From 692b862ff683ff7689006549883f0a360b1ff80a Mon Sep 17 00:00:00 2001 From: Lewis Date: Wed, 20 May 2026 18:23:45 +0300 Subject: [PATCH] fix(ingest): try to reconcile more strange records Lewis: May this revision serve well! --- crates/ingest/src/lib.rs | 2 +- crates/resolver/src/legacy_upgrade.rs | 374 ++++++++++++++++++++++++-- crates/resolver/src/lib.rs | 2 +- crates/types/src/legacy.rs | 89 +++++- crates/xrpc/src/lib.rs | 21 +- crates/xrpc/tests/cold_start.rs | 35 +++ 6 files changed, 487 insertions(+), 36 deletions(-) diff --git a/crates/ingest/src/lib.rs b/crates/ingest/src/lib.rs index bc0102b..8f6dc5c 100644 --- a/crates/ingest/src/lib.rs +++ b/crates/ingest/src/lib.rs @@ -1021,7 +1021,7 @@ async fn prepare_record( return PendingOp::ClearCache { source }; } Err(e) => { - warn!(?e, "record decode failed, clearing cache"); + warn!(?e, collection = %record.collection, "record decode failed, clearing cache"); return PendingOp::ClearCache { source }; } }; diff --git a/crates/resolver/src/legacy_upgrade.rs b/crates/resolver/src/legacy_upgrade.rs index 17e363b..73fed97 100644 --- a/crates/resolver/src/legacy_upgrade.rs +++ b/crates/resolver/src/legacy_upgrade.rs @@ -1,13 +1,16 @@ use bobbin_types::edges::{ExtractError, Record}; use bobbin_types::legacy::{ - LegacyCollaborator, LegacyIssue, LegacyPull, LegacyRecord, LegacyRefUpdate, LegacySource, - LegacyStar, LegacyTarget, + LegacyCollaborator, LegacyIssue, LegacyKnotMember, LegacyPublicKey, LegacyPull, LegacyRecord, + LegacyRefUpdate, LegacyRepo, LegacySource, LegacyStar, LegacyTarget, }; use bobbin_types::sh_tangled::feed::star::{Repo as StarRepo, Star, StarString, StarSubject}; use bobbin_types::sh_tangled::git::ref_update::RefUpdate; +use bobbin_types::sh_tangled::knot::member::Member as KnotMember; +use bobbin_types::sh_tangled::public_key::PublicKey; +use bobbin_types::sh_tangled::repo::Repo; use bobbin_types::sh_tangled::repo::collaborator::Collaborator; use bobbin_types::sh_tangled::repo::issue::Issue; -use bobbin_types::sh_tangled::repo::pull::{Pull, Source, Target}; +use bobbin_types::sh_tangled::repo::pull::{Pull, Round, Source, Target}; use jacquard_common::types::did::Did; use jacquard_common::types::nsid::Nsid; use jacquard_common::types::string::AtUri; @@ -32,14 +35,77 @@ impl DecodedRecord { ) -> Result { match Record::from_json_bytes(nsid, bytes) { Ok(record) => Ok(Self::Canon(record)), - Err(canon_err) => match LegacyRecord::from_json_bytes(nsid, bytes) { - Ok(legacy) => Ok(Self::Legacy(legacy)), - Err(_) => Err(canon_err), - }, + Err(canon_err) => { + if let Some(scrubbed) = scrub_record_bytes(nsid, bytes) { + if let Ok(record) = Record::from_json_bytes(nsid, &scrubbed) { + return Ok(Self::Canon(record)); + } + if let Ok(legacy) = LegacyRecord::from_json_bytes(nsid, &scrubbed) { + return Ok(Self::Legacy(legacy)); + } + } + match LegacyRecord::from_json_bytes(nsid, bytes) { + Ok(legacy) => Ok(Self::Legacy(legacy)), + Err(_) => Err(canon_err), + } + } } } } +#[derive(Clone, Copy, Debug)] +enum FieldRule { + DropIfEmptyString, + NullToEmptyArray, +} + +fn scrub_rules(nsid: &str) -> &'static [(&'static str, FieldRule)] { + match nsid { + "sh.tangled.actor.profile" => &[("preferredHandle", FieldRule::DropIfEmptyString)], + "sh.tangled.label.op" => &[ + ("add", FieldRule::NullToEmptyArray), + ("delete", FieldRule::NullToEmptyArray), + ], + "sh.tangled.repo.pull" => &[("rounds", FieldRule::NullToEmptyArray)], + _ => &[], + } +} + +pub fn scrub_record_bytes>( + nsid: &Nsid, + bytes: &[u8], +) -> Option> { + let rules = scrub_rules(nsid.as_ref()); + if rules.is_empty() { + return None; + } + let value: serde_json::Value = serde_json::from_slice(bytes).ok()?; + let mut obj = value.as_object()?.clone(); + let touched: alloc::vec::Vec<(&str, FieldRule)> = rules + .iter() + .filter_map(|(field, rule)| match (rule, obj.get(*field)) { + (FieldRule::DropIfEmptyString, Some(serde_json::Value::String(s))) if s.is_empty() => { + Some((*field, *rule)) + } + (FieldRule::NullToEmptyArray, Some(serde_json::Value::Null)) => Some((*field, *rule)), + _ => None, + }) + .collect(); + if touched.is_empty() { + return None; + } + touched.iter().for_each(|(field, rule)| match rule { + FieldRule::DropIfEmptyString => { + obj.remove(*field); + } + FieldRule::NullToEmptyArray => { + obj.insert((*field).to_owned(), serde_json::Value::Array(alloc::vec::Vec::new())); + } + }); + tracing::debug!(nsid = %nsid.as_ref(), ?touched, "scrubbing fields before record retry"); + serde_json::to_vec(&serde_json::Value::Object(obj)).ok() +} + async fn upgrade_repo_did( resolver: &RepoIdResolver, at_uri: Option>, @@ -95,7 +161,12 @@ fn serialize_canon_variant(record: &Record) -> Result, serde Record::Collaborator(r) => serde_json::to_vec(r), Record::RefUpdate(r) => serde_json::to_vec(r), Record::Star(r) => serde_json::to_vec(r), - _ => unreachable!("upgrade only produces Issue/Pull/Collaborator/RefUpdate/Star"), + Record::PublicKey(r) => serde_json::to_vec(r), + Record::Repo(r) => serde_json::to_vec(r), + Record::KnotMember(r) => serde_json::to_vec(r), + _ => unreachable!( + "upgrade only produces Issue/Pull/Collaborator/RefUpdate/Star/PublicKey/Repo/KnotMember" + ), } } @@ -108,6 +179,9 @@ pub async fn upgrade(legacy: LegacyRecord, resolver: &RepoIdResolver) -> Option< .map(Record::Collaborator), LegacyRecord::RefUpdate(l) => Some(Record::RefUpdate(upgrade_ref_update(l))), LegacyRecord::Star(l) => upgrade_star(l, resolver).await.map(Record::Star), + LegacyRecord::PublicKey(l) => Some(Record::PublicKey(upgrade_public_key(l))), + LegacyRecord::Repo(l) => Some(Record::Repo(upgrade_repo(l))), + LegacyRecord::KnotMember(l) => Some(Record::KnotMember(upgrade_knot_member(l))), } } @@ -160,13 +234,26 @@ async fn upgrade_pull( Some(s) => Some(upgrade_source(s, resolver).await), None => None, }; + let rounds = if l.rounds.is_empty() { + l.patch_blob + .map(|patch_blob| { + alloc::vec![Round { + created_at: l.created_at.clone(), + patch_blob, + extra_data: None, + }] + }) + .unwrap_or_default() + } else { + l.rounds + }; Some(Pull { created_at: l.created_at, body: l.body, dependent_on: l.dependent_on, mentions: l.mentions, references: l.references, - rounds: l.rounds, + rounds, source, target, title: l.title, @@ -174,6 +261,41 @@ async fn upgrade_pull( }) } +fn upgrade_public_key(l: LegacyPublicKey) -> PublicKey { + PublicKey { + created_at: l.created, + key: l.key, + name: l.name, + extra_data: l.extra_data, + } +} + +fn upgrade_repo(l: LegacyRepo) -> Repo { + let _ = l.owner; + Repo { + created_at: l.added_at, + description: l.description, + knot: l.knot, + labels: None, + name: l.name, + repo_did: None, + source: None, + spindle: None, + topics: None, + website: None, + extra_data: l.extra_data, + } +} + +fn upgrade_knot_member(l: LegacyKnotMember) -> KnotMember { + KnotMember { + created_at: l.added_at, + domain: l.domain, + subject: l.member, + extra_data: l.extra_data, + } +} + async fn upgrade_collaborator( l: LegacyCollaborator, resolver: &RepoIdResolver, @@ -297,6 +419,94 @@ mod tests { } } + #[test] + fn scrub_returns_none_for_unknown_nsid() { + let json = br#"{"$type":"sh.tangled.repo.issue","preferredHandle":""}"#; + assert!(scrub_record_bytes(&nsid("sh.tangled.repo.issue"), json).is_none()); + } + + #[test] + fn scrub_returns_none_when_target_field_is_non_empty() { + let json = br#"{"$type":"sh.tangled.actor.profile","bluesky":true,"preferredHandle":"nel.pet"}"#; + assert!(scrub_record_bytes(&nsid("sh.tangled.actor.profile"), json).is_none()); + } + + #[test] + fn scrub_returns_none_when_target_field_is_absent() { + let json = br#"{"$type":"sh.tangled.actor.profile","bluesky":true}"#; + assert!(scrub_record_bytes(&nsid("sh.tangled.actor.profile"), json).is_none()); + } + + #[test] + fn scrub_returns_none_for_non_string_value() { + let json = br#"{"$type":"sh.tangled.actor.profile","bluesky":true,"preferredHandle":42}"#; + assert!(scrub_record_bytes(&nsid("sh.tangled.actor.profile"), json).is_none()); + } + + #[test] + fn scrub_returns_none_for_non_object_json() { + assert!(scrub_record_bytes(&nsid("sh.tangled.actor.profile"), b"[]").is_none()); + assert!(scrub_record_bytes(&nsid("sh.tangled.actor.profile"), b"null").is_none()); + assert!(scrub_record_bytes(&nsid("sh.tangled.actor.profile"), b"123").is_none()); + } + + #[test] + fn scrub_returns_none_for_invalid_json() { + assert!(scrub_record_bytes(&nsid("sh.tangled.actor.profile"), b"{not json").is_none()); + } + + #[test] + fn scrub_drops_empty_preferred_handle_and_preserves_other_fields() { + let json = br#"{"$type":"sh.tangled.actor.profile","bluesky":true,"preferredHandle":"","description":"hi"}"#; + let scrubbed = scrub_record_bytes(&nsid("sh.tangled.actor.profile"), json) + .expect("empty preferredHandle must trigger scrub"); + let value: serde_json::Value = serde_json::from_slice(&scrubbed).expect("valid json"); + let obj = value.as_object().expect("object"); + assert!(!obj.contains_key("preferredHandle")); + assert_eq!(obj.get("bluesky"), Some(&serde_json::json!(true))); + assert_eq!(obj.get("description"), Some(&serde_json::json!("hi"))); + assert_eq!(obj.get("$type"), Some(&serde_json::json!("sh.tangled.actor.profile"))); + } + + #[test] + fn scrub_replaces_null_arrays_with_empty_for_label_op() { + let json = br#"{"$type":"sh.tangled.label.op","add":[{"key":"at://did:plc:limpet/sh.tangled.label.definition/k","value":"v"}],"delete":null,"performedAt":"2026-05-01T00:00:00Z","subject":"at://did:plc:limpet/sh.tangled.repo.issue/3aaa"}"#; + let scrubbed = scrub_record_bytes(&nsid("sh.tangled.label.op"), json) + .expect("null delete must trigger scrub"); + let value: serde_json::Value = serde_json::from_slice(&scrubbed).expect("valid json"); + let obj = value.as_object().expect("object"); + assert_eq!(obj.get("delete"), Some(&serde_json::json!([]))); + assert!(obj.get("add").is_some_and(|v| v.is_array())); + } + + #[test] + fn scrub_passes_through_when_label_op_arrays_are_non_null() { + let json = br#"{"$type":"sh.tangled.label.op","add":[],"delete":[],"performedAt":"2026-05-01T00:00:00Z","subject":"at://did:plc:limpet/sh.tangled.repo.issue/3aaa"}"#; + assert!(scrub_record_bytes(&nsid("sh.tangled.label.op"), json).is_none()); + } + + #[test] + fn try_decode_recovers_label_op_with_null_delete() { + let json = br#"{"$type":"sh.tangled.label.op","add":[{"key":"at://did:plc:limpet/sh.tangled.label.definition/k","value":"v"}],"delete":null,"performedAt":"2026-05-01T00:00:00Z","subject":"at://did:plc:limpet/sh.tangled.repo.issue/3aaa"}"#; + let decoded = DecodedRecord::try_decode(&nsid("sh.tangled.label.op"), json) + .expect("label.op with null delete must scrub-recover"); + assert!(matches!(decoded, DecodedRecord::Canon(Record::LabelOp(_)))); + } + + #[test] + fn try_decode_recovers_profile_with_empty_preferred_handle() { + let json = br#"{"$type":"sh.tangled.actor.profile","bluesky":true,"preferredHandle":"","description":"hi"}"#; + let decoded = DecodedRecord::try_decode(&nsid("sh.tangled.actor.profile"), json) + .expect("profile with empty preferredHandle must scrub-recover"); + match decoded { + DecodedRecord::Canon(Record::Profile(p)) => { + assert!(p.preferred_handle.is_none()); + assert_eq!(p.description.as_deref(), Some("hi")); + } + other => panic!("expected canon profile, got {other:?}"), + } + } + #[test] fn legacy_decode_passes_through_for_unaffected_nsids() { let json = br#"{"$type":"sh.tangled.graph.follow","subject":"did:plc:bailey","createdAt":"2026-05-01T00:00:00Z"}"#; @@ -308,11 +518,11 @@ mod tests { #[tokio::test] async fn upgrade_issue_uses_repo_did_directly() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); - let json = br#"{"$type":"sh.tangled.repo.issue","repoDid":"did:plc:abalone","title":"t","createdAt":"2026-05-01T00:00:00Z"}"#; + let json = br#"{"$type":"sh.tangled.repo.issue","repoDid":"did:plc:scallop","title":"t","createdAt":"2026-05-01T00:00:00Z"}"#; let legacy = LegacyRecord::from_json_bytes(&nsid("sh.tangled.repo.issue"), json).expect("decode"); let canon = upgrade(legacy, &resolver).await.expect("upgrade"); match canon { - Record::Issue(i) => assert_eq!(i.repo, did("did:plc:abalone")), + Record::Issue(i) => assert_eq!(i.repo, did("did:plc:scallop")), other => panic!("expected canon issue, got {other:?}"), } } @@ -323,13 +533,13 @@ mod tests { let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); resolver - .observe(owner.clone(), key.clone(), Some(did("did:plc:abalone"))) + .observe(owner.clone(), key.clone(), Some(did("did:plc:scallop"))) .await; let json = br#"{"$type":"sh.tangled.repo.issue","repo":"at://did:plc:nel/sh.tangled.repo/abcabcabcabcz","title":"t","createdAt":"2026-05-01T00:00:00Z"}"#; let legacy = LegacyRecord::from_json_bytes(&nsid("sh.tangled.repo.issue"), json).expect("decode"); let canon = upgrade(legacy, &resolver).await.expect("upgrade"); match canon { - Record::Issue(i) => assert_eq!(i.repo, did("did:plc:abalone")), + Record::Issue(i) => assert_eq!(i.repo, did("did:plc:scallop")), other => panic!("expected canon issue, got {other:?}"), } } @@ -348,27 +558,141 @@ mod tests { #[tokio::test] async fn upgrade_pull_propagates_target_resolution() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); - let json = br#"{"$type":"sh.tangled.repo.pull","title":"t","createdAt":"2026-05-01T00:00:00Z","rounds":[],"target":{"branch":"main","repoDid":"did:plc:abalone"}}"#; + let json = br#"{"$type":"sh.tangled.repo.pull","title":"t","createdAt":"2026-05-01T00:00:00Z","rounds":[],"target":{"branch":"main","repoDid":"did:plc:scallop"}}"#; let legacy = LegacyRecord::from_json_bytes(&nsid("sh.tangled.repo.pull"), json).expect("decode"); let canon = upgrade(legacy, &resolver).await.expect("upgrade"); match canon { Record::Pull(p) => { - assert_eq!(p.target.repo, did("did:plc:abalone")); + assert_eq!(p.target.repo, did("did:plc:scallop")); assert!(p.source.is_none()); } other => panic!("expected canon pull, got {other:?}"), } } + #[tokio::test] + async fn upgrade_pull_pre_rounds_synthesizes_round_from_top_level_patch_blob() { + let resolver = RepoIdResolver::detached(RuntimeHasher::default()); + let json = br#"{"$type":"sh.tangled.repo.pull","title":"t","createdAt":"2026-05-01T00:00:00Z","target":{"branch":"main","repoDid":"did:plc:scallop"},"patchBlob":{"$type":"blob","mimeType":"application/gzip","ref":{"$link":"bafkreibpatvbeajtwzlr4jwr4s2hnwo5l7sgdbfnqu6n7ctd2bcbtluw4a"},"size":920}}"#; + let legacy = LegacyRecord::from_json_bytes(&nsid("sh.tangled.repo.pull"), json).expect("decode"); + let canon = upgrade(legacy, &resolver).await.expect("upgrade"); + match canon { + Record::Pull(p) => { + assert_eq!(p.rounds.len(), 1, "pre-rounds wire must yield exactly one synthesized round"); + assert_eq!(p.rounds[0].patch_blob.blob().mime_type.as_ref(), "application/gzip"); + assert_eq!(p.rounds[0].created_at, p.created_at); + } + other => panic!("expected canon pull, got {other:?}"), + } + } + + #[tokio::test] + async fn upgrade_pull_omits_round_when_neither_rounds_nor_patch_blob_present() { + let resolver = RepoIdResolver::detached(RuntimeHasher::default()); + let json = br#"{"$type":"sh.tangled.repo.pull","title":"t","createdAt":"2026-05-01T00:00:00Z","target":{"branch":"main","repoDid":"did:plc:scallop"}}"#; + let legacy = LegacyRecord::from_json_bytes(&nsid("sh.tangled.repo.pull"), json).expect("decode"); + let canon = upgrade(legacy, &resolver).await.expect("upgrade"); + match canon { + Record::Pull(p) => assert!(p.rounds.is_empty()), + other => panic!("expected canon pull, got {other:?}"), + } + } + + #[tokio::test] + async fn upgrade_public_key_renames_created_to_created_at() { + let resolver = RepoIdResolver::detached(RuntimeHasher::default()); + let json = br#"{"$type":"sh.tangled.publicKey","created":"2025-04-15T18:35:38Z","key":"ssh-ed25519 AAAA","name":"laptop"}"#; + let legacy = LegacyRecord::from_json_bytes(&nsid("sh.tangled.publicKey"), json).expect("decode"); + let canon = upgrade(legacy, &resolver).await.expect("upgrade"); + match canon { + Record::PublicKey(k) => { + assert_eq!(k.created_at.as_str(), "2025-04-15T18:35:38Z"); + assert_eq!(k.key.as_str(), "ssh-ed25519 AAAA"); + assert_eq!(k.name.as_str(), "laptop"); + } + other => panic!("expected canon publicKey, got {other:?}"), + } + } + + #[tokio::test] + async fn upgrade_repo_renames_added_at_and_drops_owner() { + let resolver = RepoIdResolver::detached(RuntimeHasher::default()); + let json = br#"{"$type":"sh.tangled.repo","addedAt":"2025-03-21T10:18:58Z","description":"hi","knot":"knot1.tangled.sh","name":"site","owner":"did:plc:nel"}"#; + let legacy = LegacyRecord::from_json_bytes(&nsid("sh.tangled.repo"), json).expect("decode"); + let canon = upgrade(legacy, &resolver).await.expect("upgrade"); + match canon { + Record::Repo(r) => { + assert_eq!(r.created_at.as_str(), "2025-03-21T10:18:58Z"); + assert_eq!(r.description.as_deref(), Some("hi")); + assert_eq!(r.knot.as_str(), "knot1.tangled.sh"); + assert_eq!(r.name.as_deref(), Some("site")); + assert!(r.repo_did.is_none(), "legacy repos have no repo_did"); + } + other => panic!("expected canon repo, got {other:?}"), + } + } + + #[tokio::test] + async fn upgrade_knot_member_renames_added_at_and_member() { + let resolver = RepoIdResolver::detached(RuntimeHasher::default()); + let json = br#"{"$type":"sh.tangled.knot.member","addedAt":"2025-03-31T05:14:09Z","domain":"knot.example","member":"did:plc:nel"}"#; + let legacy = LegacyRecord::from_json_bytes(&nsid("sh.tangled.knot.member"), json).expect("decode"); + let canon = upgrade(legacy, &resolver).await.expect("upgrade"); + match canon { + Record::KnotMember(m) => { + assert_eq!(m.created_at.as_str(), "2025-03-31T05:14:09Z"); + assert_eq!(m.domain.as_str(), "knot.example"); + assert_eq!(m.subject, did("did:plc:nel")); + } + other => panic!("expected canon knot.member, got {other:?}"), + } + } + + #[tokio::test] + async fn legacy_pull_target_with_empty_repo_did_treats_as_none() { + let resolver = RepoIdResolver::detached(RuntimeHasher::default()); + resolver.observe( + did("did:plc:nel"), + rkey("abcabcabcabcz"), + Some(did("did:plc:scallop")), + ).await; + let json = br#"{"$type":"sh.tangled.repo.pull","title":"t","createdAt":"2026-05-01T00:00:00Z","rounds":[],"target":{"branch":"main","repo":"at://did:plc:nel/sh.tangled.repo/abcabcabcabcz","repoDid":""}}"#; + let legacy = LegacyRecord::from_json_bytes(&nsid("sh.tangled.repo.pull"), json).expect("decode"); + let canon = upgrade(legacy, &resolver).await.expect("upgrade"); + match canon { + Record::Pull(p) => assert_eq!(p.target.repo, did("did:plc:scallop")), + other => panic!("expected canon pull, got {other:?}"), + } + } + + #[tokio::test] + async fn try_decode_recovers_publickey_with_created_field() { + let json = br#"{"$type":"sh.tangled.publicKey","created":"2025-04-15T18:35:38Z","key":"k","name":"n"}"#; + let decoded = DecodedRecord::try_decode(&nsid("sh.tangled.publicKey"), json) + .expect("legacy publicKey must decode"); + assert!(matches!( + decoded, + DecodedRecord::Legacy(LegacyRecord::PublicKey(_)) + )); + } + + #[tokio::test] + async fn try_decode_recovers_repo_with_added_at_field() { + let json = br#"{"$type":"sh.tangled.repo","addedAt":"2025-03-21T10:18:58Z","knot":"knot1.tangled.sh","owner":"did:plc:nel"}"#; + let decoded = DecodedRecord::try_decode(&nsid("sh.tangled.repo"), json) + .expect("legacy repo must decode"); + assert!(matches!(decoded, DecodedRecord::Legacy(LegacyRecord::Repo(_)))); + } + #[tokio::test] async fn upgrade_pull_source_repo_resolution_is_independent_of_target() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); - let json = br#"{"$type":"sh.tangled.repo.pull","title":"t","createdAt":"2026-05-01T00:00:00Z","rounds":[],"target":{"branch":"main","repoDid":"did:plc:abalone"},"source":{"branch":"feat","repo":"at://did:plc:nel/sh.tangled.repo/missing"}}"#; + let json = br#"{"$type":"sh.tangled.repo.pull","title":"t","createdAt":"2026-05-01T00:00:00Z","rounds":[],"target":{"branch":"main","repoDid":"did:plc:scallop"},"source":{"branch":"feat","repo":"at://did:plc:nel/sh.tangled.repo/missing"}}"#; let legacy = LegacyRecord::from_json_bytes(&nsid("sh.tangled.repo.pull"), json).expect("decode"); let canon = upgrade(legacy, &resolver).await.expect("upgrade"); match canon { Record::Pull(p) => { - assert_eq!(p.target.repo, did("did:plc:abalone")); + assert_eq!(p.target.repo, did("did:plc:scallop")); let source = p.source.expect("source struct retained"); assert_eq!(source.branch.as_str(), "feat"); assert!( @@ -383,12 +707,12 @@ mod tests { #[tokio::test] async fn upgrade_ref_update_renames_repo_did_to_repo() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); - let json = br#"{"$type":"sh.tangled.git.refUpdate","ref":"refs/heads/main","committerDid":"did:plc:olaren","repoDid":"did:plc:abalone","oldSha":"0000000000000000000000000000000000000000","newSha":"1111111111111111111111111111111111111111","meta":{"isDefaultRef":true,"commitCount":{}}}"#; + let json = br#"{"$type":"sh.tangled.git.refUpdate","ref":"refs/heads/main","committerDid":"did:plc:olaren","repoDid":"did:plc:scallop","oldSha":"0000000000000000000000000000000000000000","newSha":"1111111111111111111111111111111111111111","meta":{"isDefaultRef":true,"commitCount":{}}}"#; let legacy = LegacyRecord::from_json_bytes(&nsid("sh.tangled.git.refUpdate"), json).expect("decode"); let canon = upgrade(legacy, &resolver).await.expect("upgrade"); match canon { - Record::RefUpdate(r) => assert_eq!(r.repo, did("did:plc:abalone")), + Record::RefUpdate(r) => assert_eq!(r.repo, did("did:plc:scallop")), other => panic!("expected canon ref update, got {other:?}"), } } @@ -396,12 +720,12 @@ mod tests { #[tokio::test] async fn upgrade_star_prefers_subject_did_over_subject_uri() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); - let json = br#"{"$type":"sh.tangled.feed.star","createdAt":"2026-05-01T00:00:00Z","subject":"at://did:plc:nel/sh.tangled.string/k1","subjectDid":"did:plc:abalone"}"#; + let json = br#"{"$type":"sh.tangled.feed.star","createdAt":"2026-05-01T00:00:00Z","subject":"at://did:plc:nel/sh.tangled.string/k1","subjectDid":"did:plc:scallop"}"#; let legacy = LegacyRecord::from_json_bytes(&nsid("sh.tangled.feed.star"), json).expect("decode"); let canon = upgrade(legacy, &resolver).await.expect("upgrade"); match canon { Record::Star(s) => match s.subject { - StarSubject::Repo(r) => assert_eq!(r.did, did("did:plc:abalone")), + StarSubject::Repo(r) => assert_eq!(r.did, did("did:plc:scallop")), StarSubject::String(_) => panic!("subjectDid must win"), }, other => panic!("expected canon star, got {other:?}"), @@ -433,14 +757,14 @@ mod tests { let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); resolver - .observe(owner.clone(), key.clone(), Some(did("did:plc:abalone"))) + .observe(owner.clone(), key.clone(), Some(did("did:plc:scallop"))) .await; let json = br#"{"$type":"sh.tangled.feed.star","createdAt":"2026-05-01T00:00:00Z","subject":"at://did:plc:nel/sh.tangled.repo/abcabcabcabcz"}"#; let legacy = LegacyRecord::from_json_bytes(&nsid("sh.tangled.feed.star"), json).expect("decode"); let canon = upgrade(legacy, &resolver).await.expect("upgrade"); match canon { Record::Star(s) => match s.subject { - StarSubject::Repo(r) => assert_eq!(r.did, did("did:plc:abalone")), + StarSubject::Repo(r) => assert_eq!(r.did, did("did:plc:scallop")), StarSubject::String(_) => panic!("observed cache must upgrade to Repo variant"), }, other => panic!("expected canon star, got {other:?}"), @@ -450,7 +774,7 @@ mod tests { #[tokio::test] async fn upgrade_collaborator_requires_repo_did() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); - let with_did = br#"{"$type":"sh.tangled.repo.collaborator","createdAt":"2026-05-01T00:00:00Z","subject":"did:plc:lyna","repoDid":"did:plc:abalone"}"#; + let with_did = br#"{"$type":"sh.tangled.repo.collaborator","createdAt":"2026-05-01T00:00:00Z","subject":"did:plc:lyna","repoDid":"did:plc:scallop"}"#; let canon = upgrade( LegacyRecord::from_json_bytes(&nsid("sh.tangled.repo.collaborator"), with_did) .expect("decode"), @@ -459,7 +783,7 @@ mod tests { .await .expect("upgrade"); match canon { - Record::Collaborator(c) => assert_eq!(c.repo, did("did:plc:abalone")), + Record::Collaborator(c) => assert_eq!(c.repo, did("did:plc:scallop")), other => panic!("expected canon collaborator, got {other:?}"), } diff --git a/crates/resolver/src/lib.rs b/crates/resolver/src/lib.rs index 537c438..3a91145 100644 --- a/crates/resolver/src/lib.rs +++ b/crates/resolver/src/lib.rs @@ -2,7 +2,7 @@ mod legacy_upgrade; mod normalize; pub use legacy_upgrade::{ - DecodedRecord, decode_canon_or_upgrade_bytes, upgrade, upgrade_wire_bytes, + DecodedRecord, decode_canon_or_upgrade_bytes, scrub_record_bytes, upgrade, upgrade_wire_bytes, }; pub use normalize::NormalizeRepoRefs; diff --git a/crates/types/src/legacy.rs b/crates/types/src/legacy.rs index 440e210..f642919 100644 --- a/crates/types/src/legacy.rs +++ b/crates/types/src/legacy.rs @@ -2,15 +2,31 @@ use alloc::collections::BTreeMap; use alloc::vec::Vec; use jacquard_common::deps::smol_str::SmolStr; +use jacquard_common::types::blob::BlobRef; use jacquard_common::types::nsid::Nsid; use jacquard_common::types::string::{AtUri, Datetime, Did}; use jacquard_common::types::value::Data; use jacquard_common::{BosStr, DefaultStr}; -use serde::Deserialize; +use serde::{Deserialize, Deserializer}; use crate::edges::ExtractError; use crate::sh_tangled::repo::pull::Round as CanonRound; +fn empty_string_as_none<'de, D, T>(d: D) -> Result, D::Error> +where + D: Deserializer<'de>, + T: Deserialize<'de>, +{ + let value: Option = Deserialize::deserialize(d)?; + match value { + None | Some(serde_json::Value::Null) => Ok(None), + Some(serde_json::Value::String(s)) if s.is_empty() => Ok(None), + Some(other) => T::deserialize(other) + .map(Some) + .map_err(serde::de::Error::custom), + } +} + #[derive(Debug, Deserialize)] #[serde( rename_all = "camelCase", @@ -44,7 +60,11 @@ pub struct LegacyTarget { pub branch: S, #[serde(default, skip_serializing_if = "Option::is_none")] pub repo: Option>, - #[serde(default, skip_serializing_if = "Option::is_none")] + #[serde( + default, + skip_serializing_if = "Option::is_none", + deserialize_with = "empty_string_as_none" + )] pub repo_did: Option>, } @@ -57,7 +77,11 @@ pub struct LegacySource { pub branch: S, #[serde(default, skip_serializing_if = "Option::is_none")] pub repo: Option>, - #[serde(default, skip_serializing_if = "Option::is_none")] + #[serde( + default, + skip_serializing_if = "Option::is_none", + deserialize_with = "empty_string_as_none" + )] pub repo_did: Option>, } @@ -71,6 +95,7 @@ pub struct LegacySource { pub struct LegacyPull { pub created_at: Datetime, pub title: S, + #[serde(default)] pub rounds: Vec>, pub target: LegacyTarget, #[serde(default, skip_serializing_if = "Option::is_none")] @@ -83,6 +108,8 @@ pub struct LegacyPull { pub references: Option>>, #[serde(default, skip_serializing_if = "Option::is_none")] pub dependent_on: Option>, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub patch_blob: Option>, #[serde(flatten, default, skip_serializing_if = "Option::is_none")] pub extra_data: Option>>, } @@ -144,6 +171,56 @@ pub struct LegacyStar { pub extra_data: Option>>, } +#[derive(Debug, Deserialize)] +#[serde( + rename_all = "camelCase", + rename = "sh.tangled.publicKey", + tag = "$type", + bound(deserialize = "S: Deserialize<'de> + BosStr") +)] +pub struct LegacyPublicKey { + pub created: Datetime, + pub key: S, + pub name: S, + #[serde(flatten, default, skip_serializing_if = "Option::is_none")] + pub extra_data: Option>>, +} + +#[derive(Debug, Deserialize)] +#[serde( + rename_all = "camelCase", + rename = "sh.tangled.repo", + tag = "$type", + bound(deserialize = "S: Deserialize<'de> + BosStr") +)] +pub struct LegacyRepo { + pub added_at: Datetime, + pub knot: S, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub description: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub name: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub owner: Option>, + #[serde(flatten, default, skip_serializing_if = "Option::is_none")] + pub extra_data: Option>>, +} + +#[derive(Debug, Deserialize)] +#[serde( + rename_all = "camelCase", + rename = "sh.tangled.knot.member", + tag = "$type", + bound(deserialize = "S: Deserialize<'de> + BosStr") +)] +pub struct LegacyKnotMember { + pub added_at: Datetime, + pub domain: S, + pub member: Did, + #[serde(flatten, default, skip_serializing_if = "Option::is_none")] + pub extra_data: Option>>, +} + #[derive(Debug)] pub enum LegacyRecord { Issue(LegacyIssue), @@ -151,6 +228,9 @@ pub enum LegacyRecord { Collaborator(LegacyCollaborator), RefUpdate(LegacyRefUpdate), Star(LegacyStar), + PublicKey(LegacyPublicKey), + Repo(LegacyRepo), + KnotMember(LegacyKnotMember), } impl LegacyRecord { @@ -166,6 +246,9 @@ impl LegacyRecord { } "sh.tangled.git.refUpdate" => Ok(Self::RefUpdate(serde_json::from_slice(bytes)?)), "sh.tangled.feed.star" => Ok(Self::Star(serde_json::from_slice(bytes)?)), + "sh.tangled.publicKey" => Ok(Self::PublicKey(serde_json::from_slice(bytes)?)), + "sh.tangled.repo" => Ok(Self::Repo(serde_json::from_slice(bytes)?)), + "sh.tangled.knot.member" => Ok(Self::KnotMember(serde_json::from_slice(bytes)?)), other => Err(ExtractError::UnknownCollection(other.into())), } } diff --git a/crates/xrpc/src/lib.rs b/crates/xrpc/src/lib.rs index 37f210b..973bd53 100644 --- a/crates/xrpc/src/lib.rs +++ b/crates/xrpc/src/lib.rs @@ -1,7 +1,7 @@ use std::future::Future; use std::sync::Arc; -use bobbin_resolver::{NormalizeRepoRefs, upgrade_wire_bytes}; +use bobbin_resolver::{NormalizeRepoRefs, scrub_record_bytes, upgrade_wire_bytes}; use axum::{ Router, @@ -1248,11 +1248,20 @@ where { match serde_json::from_slice::(bytes) { Ok(v) => Ok(v), - Err(canon_err) => match upgrade_wire_bytes(nsid, bytes, &state.resolver).await { - Ok(canon_bytes) => serde_json::from_slice(&canon_bytes) - .map_err(|e| XrpcError::InvalidRecord(e.to_string())), - Err(_) => Err(XrpcError::InvalidRecord(canon_err.to_string())), - }, + Err(canon_err) => { + let scrubbed = scrub_record_bytes(nsid, bytes); + let retry_bytes: &[u8] = scrubbed.as_deref().unwrap_or(bytes); + if scrubbed.is_some() { + if let Ok(v) = serde_json::from_slice::(retry_bytes) { + return Ok(v); + } + } + match upgrade_wire_bytes(nsid, retry_bytes, &state.resolver).await { + Ok(canon_bytes) => serde_json::from_slice(&canon_bytes) + .map_err(|e| XrpcError::InvalidRecord(e.to_string())), + Err(_) => Err(XrpcError::InvalidRecord(canon_err.to_string())), + } + } } } diff --git a/crates/xrpc/tests/cold_start.rs b/crates/xrpc/tests/cold_start.rs index 526350b..1d00bfa 100644 --- a/crates/xrpc/tests/cold_start.rs +++ b/crates/xrpc/tests/cold_start.rs @@ -373,6 +373,41 @@ async fn wrong_type_does_not_poison_cache() { drop(mock); } +#[tokio::test] +async fn profile_with_empty_preferred_handle_is_tolerated() { + let server = MockServer::start().await; + let nel = did("did:plc:nel"); + mount_record( + &server, + &nel, + &nsid("sh.tangled.actor.profile"), + &rkey("self"), + json!({ + "$type": "sh.tangled.actor.profile", + "bluesky": true, + "preferredHandle": "", + "description": "empty handle, valid profile" + }), + ) + .await; + let state = fresh_app(&Url::parse(&server.uri()).unwrap()).await; + let app = router(state); + let at_uri = format!("at://{}/sh.tangled.actor.profile/self", nel.as_ref()); + let resp = app + .oneshot(xrpc_request( + "sh.tangled.actor.getProfile", + "actor", + &at_uri, + )) + .await + .unwrap(); + let (status, body) = json_response(resp).await; + assert_eq!(status, StatusCode::OK, "status: {body}"); + assert_eq!(body["uri"], at_uri); + assert_eq!(body["value"]["description"], "empty handle, valid profile"); + assert!(body["value"]["preferredHandle"].is_null()); +} + #[tokio::test] async fn missing_uri_param_returns_json_envelope() { let server = MockServer::start().await; -- 2.51.2