From 804ea44a365ea04cd5fd59cbc1dbcbae2db8dc94 Mon Sep 17 00:00:00 2001 From: Rudy Fraser Date: Sat, 15 Aug 2026 14:44:37 -0400 Subject: [PATCH] fix(repo): blob refs parsed from CAR imports keep their blobs ipld_to_lex only recognized blob shapes on the Json variant, which DAG-CBOR decoding never produces, so every blob ref in an imported CAR dissolved into a plain map: no record_blob association, expectedBlobs 0, and uploaded blobs stranded in temp storage forever. Blob-shaped CBOR maps (typed and legacy) now become Lex::Blob, import-time blob processing tolerates blobs that have not been uploaded yet, and an upload whose cid is already associated is promoted immediately. --- Cargo.lock | 4 +- Cargo.toml | 2 +- rsky-pds/Cargo.toml | 2 +- rsky-pds/src/actor_store/blob/mod.rs | 51 +++++++++++++++++++- rsky-pds/src/actor_store/mod.rs | 2 +- rsky-repo/Cargo.toml | 2 +- rsky-repo/src/util.rs | 70 ++++++++++++++++++++++++++++ 7 files changed, 126 insertions(+), 7 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index d5af0974..6457601c 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -8211,7 +8211,7 @@ dependencies = [ [[package]] name = "rsky-pds" -version = "0.13.11" +version = "0.13.12" dependencies = [ "anyhow", "argon2", @@ -8339,7 +8339,7 @@ dependencies = [ [[package]] name = "rsky-repo" -version = "0.0.3" +version = "0.0.4" dependencies = [ "anyhow", "async-recursion", diff --git a/Cargo.toml b/Cargo.toml index ffa740f0..a8e0561d 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -44,7 +44,7 @@ rsky-identity = {path = "rsky-identity", version = "0.2.0"} rsky-crypto = {path = "rsky-crypto", version = "0.2.0"} rsky-syntax = {path = "rsky-syntax", version = "0.1.0"} rsky-common = {path = "rsky-common", version = "0.1.3"} -rsky-repo = {path = "rsky-repo", version = "0.0.3"} +rsky-repo = {path = "rsky-repo", version = "0.0.4"} rsky-firehose = {path = "rsky-firehose", version = "0.2.1"} [profile.release] diff --git a/rsky-pds/Cargo.toml b/rsky-pds/Cargo.toml index 9e9f397a..7620116a 100644 --- a/rsky-pds/Cargo.toml +++ b/rsky-pds/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "rsky-pds" -version = "0.13.11" +version = "0.13.12" authors = ["Rudy Fraser "] description = "Rust reference implementation of an atproto PDS." license = "Apache-2.0" diff --git a/rsky-pds/src/actor_store/blob/mod.rs b/rsky-pds/src/actor_store/blob/mod.rs index ddae32c7..05340e7b 100644 --- a/rsky-pds/src/actor_store/blob/mod.rs +++ b/rsky-pds/src/actor_store/blob/mod.rs @@ -12,7 +12,7 @@ use rsky_common::now; use rsky_lexicon::blob_refs::BlobRef; use rsky_lexicon::com::atproto::admin::StatusAttr; use rsky_lexicon::com::atproto::repo::ListMissingBlobsRefRecordBlob; -use rsky_repo::types::{PreparedBlobRef, PreparedWrite}; +use rsky_repo::types::{BlobConstraint, PreparedBlobRef, PreparedWrite}; use rusqlite::OptionalExtension; use sha2::{Digest, Sha256}; use std::str::FromStr; @@ -194,9 +194,58 @@ impl BlobReader { Ok(()) }) .await?; + // A blob uploaded after the record referencing it was imported is + // already associated; promote it now or it stays temp forever. + let cid_str = cid.to_string(); + let associated: bool = self + .db + .run(move |conn| { + Ok(conn + .query_row( + "SELECT 1 FROM record_blob WHERE \"blobCid\" = ?1 LIMIT 1", + [cid_str.clone()], + |_| Ok(()), + ) + .optional()? + .is_some()) + }) + .await?; + if associated { + self.verify_blob_and_make_permanent(PreparedBlobRef { + cid, + mime_type: mime_type.clone(), + constraints: BlobConstraint { + max_size: None, + accept: None, + }, + }) + .await?; + } Ok(BlobRef::new(cid, mime_type, size, None)) } + /// Blob processing for CAR imports: the referenced blobs typically + /// arrive by uploadBlob *after* the import, so a missing blob records + /// the association and moves on instead of failing the import; the + /// upload promotes it (see `track_untethered_blob`). + pub async fn process_import_blobs(&self, writes: Vec) -> Result<()> { + for write in writes { + let (blobs, uri) = match write { + PreparedWrite::Create(w) => (w.blobs, w.uri), + PreparedWrite::Update(w) => (w.blobs, w.uri), + _ => continue, + }; + for blob in blobs { + self.associate_blob(blob.clone(), uri.clone()).await?; + if let Err(error) = self.verify_blob_and_make_permanent(blob.clone()).await { + tracing::debug!(cid = %blob.cid, %error, + "imported record references a blob not yet uploaded"); + } + } + } + Ok(()) + } + pub async fn process_write_blobs(&self, writes: Vec) -> Result<()> { self.delete_dereferenced_blobs(writes.clone()).await?; for write in writes { diff --git a/rsky-pds/src/actor_store/mod.rs b/rsky-pds/src/actor_store/mod.rs index 4d2dd5c9..7ea01bec 100644 --- a/rsky-pds/src/actor_store/mod.rs +++ b/rsky-pds/src/actor_store/mod.rs @@ -491,7 +491,7 @@ impl ActorStoreTransactor { storage_guard.apply_commit(commit.clone(), None).await?; } // process blobs - self.blob.process_write_blobs(writes).await?; + self.blob.process_import_blobs(writes).await?; Ok(()) } diff --git a/rsky-repo/Cargo.toml b/rsky-repo/Cargo.toml index 3dac34b1..bf39a75d 100644 --- a/rsky-repo/Cargo.toml +++ b/rsky-repo/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "rsky-repo" -version = "0.0.3" +version = "0.0.4" authors = ["Rudy Fraser "] description = "Rust crate for atproto repositories, an in particular the Merkle Search Tree (MST) data structure." license = "Apache-2.0" diff --git a/rsky-repo/src/util.rs b/rsky-repo/src/util.rs index ec7ab442..6bb828a1 100644 --- a/rsky-repo/src/util.rs +++ b/rsky-repo/src/util.rs @@ -9,6 +9,7 @@ use futures::{stream, Stream, StreamExt, TryStreamExt}; use lexicon_cid::Cid; use rsky_common::sign::sign_without_indexmap; use rsky_common::tid::Ticker; +use rsky_lexicon::blob_refs::{BlobRef, JsonBlobRef}; use secp256k1::Keypair; use serde_json::Value as JsonValue; use sha2::{Digest, Sha256}; @@ -70,10 +71,38 @@ pub fn lex_to_ipld(val: Lex) -> Ipld { } } +/// A DAG-CBOR blob reference decodes as an `Ipld::Map` whose `ref` is an +/// `Ipld::Link`; recognize that shape here or every blob ref parsed from a +/// CAR dissolves into a plain map and record-blob associations are lost. +fn ipld_map_as_blob(map: &BTreeMap) -> Option { + let typed = matches!(map.get("$type"), Some(Ipld::String(t)) if t == "blob"); + let legacy = matches!(map.get("cid"), Some(Ipld::String(_))) + && matches!(map.get("mimeType"), Some(Ipld::String(_))); + if !typed && !legacy { + return None; + } + let mut obj = serde_json::Map::new(); + for (key, value) in map { + let json = match value { + Ipld::Link(cid) => serde_json::json!({ "$link": cid.to_string() }), + Ipld::String(s) => JsonValue::String(s.clone()), + Ipld::Json(value) => value.clone(), + _ => return None, + }; + obj.insert(key.clone(), json); + } + serde_json::from_value::(JsonValue::Object(obj)) + .ok() + .map(|original| BlobRef { original }) +} + pub fn ipld_to_lex(val: Ipld) -> Lex { match val { Ipld::List(list) => Lex::List(list.into_iter().map(ipld_to_lex).collect::>()), Ipld::Map(map) => { + if let Some(blob) = ipld_map_as_blob(&map) { + return Lex::Blob(blob); + } let mut to_return: BTreeMap = BTreeMap::new(); for key in map.keys() { to_return.insert(key.to_owned(), ipld_to_lex(map.get(key).unwrap().clone())); @@ -219,3 +248,44 @@ where } Ok(buffer) } + +#[cfg(test)] +mod tests { + use super::*; + + /// A blob ref decoded from real DAG-CBOR (map with a tag-42 link) must + /// come out as `Lex::Blob`, or record-blob associations are silently + /// lost on every CAR import. + #[test] + fn cbor_blob_ref_parses_as_lex_blob() { + let cid: Cid = "bafkreiey6e2xp4ncufvsyfmubucbsz5xujbc7lguospuziohtgfdik3pr4" + .parse() + .unwrap(); + let record = Ipld::Map(BTreeMap::from([ + ( + "$type".to_string(), + Ipld::String("app.bsky.actor.profile".to_string()), + ), + ( + "avatar".to_string(), + Ipld::Map(BTreeMap::from([ + ("$type".to_string(), Ipld::String("blob".to_string())), + ("ref".to_string(), Ipld::Link(cid)), + ( + "mimeType".to_string(), + Ipld::String("image/jpeg".to_string()), + ), + ("size".to_string(), Ipld::Json(JsonValue::from(924586))), + ])), + ), + ])); + let bytes = serde_ipld_dagcbor::to_vec(&record).unwrap(); + let parsed = cbor_to_lex_record(bytes).unwrap(); + match parsed.get("avatar") { + Some(Lex::Blob(blob)) => { + assert_eq!(blob.get_cid().unwrap(), cid); + } + other => panic!("avatar should parse as Lex::Blob, got {other:?}"), + } + } +} -- 2.51.2