From 5d6b15c36484f3d63493f8e89f509a1ac95fb3d1 Mon Sep 17 00:00:00 2001 From: Shivakumar Date: Mon, 20 Jul 2026 23:00:19 +0530 Subject: [PATCH 1/6] wip(asset-store): migrate asset storage onto agentflare-store documents+blobs (#185) Replace agentflare-backend::asset's storage with agentflare-store's documents+blobs (content-addressed blob store with dedup/unref), and route handoff/asset MCP tool handlers through it. --- crates/agentflare-backend/src/asset.rs | 9 + crates/agentflare-store/src/documents.rs | 64 +++++-- crates/agentflare-store/src/migrations.rs | 6 + src/asset_store.rs | 103 +++++++++++ src/main.rs | 1 + src/mcp_server.rs | 72 ++++++-- src/mcp_server/asset.rs | 203 +++++++++++----------- src/mcp_server/handoff.rs | 87 +++++----- src/mcp_server/tests/artifact_tests.rs | 16 +- src/mcp_server/tests/asset_tests.rs | 63 +------ 10 files changed, 379 insertions(+), 245 deletions(-) create mode 100644 src/asset_store.rs diff --git a/crates/agentflare-backend/src/asset.rs b/crates/agentflare-backend/src/asset.rs index 7cfb693d..ec0daf78 100644 --- a/crates/agentflare-backend/src/asset.rs +++ b/crates/agentflare-backend/src/asset.rs @@ -170,6 +170,15 @@ pub fn get(conn: &Connection, id: &str) -> Result { }) } +pub fn list_all(conn: &Connection) -> Result> { + let mut stmt = conn.prepare( + "SELECT id, workspace_id, entity_type, entity_id, filename, size, storage_path, mime_type, metadata, created_at, updated_at, deleted_at, version + FROM assets WHERE deleted_at IS NULL ORDER BY created_at", + )?; + let rows = stmt.query_map([], row_to_asset)?; + Ok(rows.collect::>()?) +} + pub fn list_by_entity(conn: &Connection, entity_type: &str, entity_id: &str) -> Result> { let mut stmt = conn.prepare( "SELECT id, workspace_id, entity_type, entity_id, filename, size, storage_path, mime_type, metadata, created_at, updated_at, deleted_at, version diff --git a/crates/agentflare-store/src/documents.rs b/crates/agentflare-store/src/documents.rs index a16aa037..bf1d6dd7 100644 --- a/crates/agentflare-store/src/documents.rs +++ b/crates/agentflare-store/src/documents.rs @@ -18,6 +18,8 @@ pub struct Document { pub version: i32, pub created_at: i64, pub updated_at: i64, + pub metadata: String, + pub size: i64, pub deleted_at: Option, } @@ -30,6 +32,8 @@ pub struct DocVersion { pub blob_hash: Option, pub mime: String, pub title: String, + pub metadata: String, + pub size: i64, pub created_at: i64, } @@ -51,6 +55,8 @@ pub struct DocUpsertOpts { pub tags: Option>, pub session_id: Option, pub source: Option, + pub metadata: Option, + pub size: Option, } impl Store { @@ -86,9 +92,11 @@ impl Store { session_id: row.get(9)?, source: row.get(10)?, version: row.get(11)?, - created_at: row.get(12)?, - updated_at: row.get(13)?, - deleted_at: row.get(14)?, + metadata: row.get(12)?, + size: row.get(13)?, + created_at: row.get(14)?, + updated_at: row.get(15)?, + deleted_at: row.get(16)?, }) } @@ -120,7 +128,7 @@ impl Store { let existing = tx .query_row( - "SELECT id, rowid, content, version, blob_hash, mime FROM store_documents + "SELECT id, rowid, content, version, blob_hash, mime, metadata, size FROM store_documents WHERE project_id = ?1 AND path = ?2", params![project_id, path], |row| { @@ -131,12 +139,14 @@ impl Store { row.get::<_, i32>(3)?, row.get::<_, Option>(4)?, row.get::<_, String>(5)?, + row.get::<_, String>(6)?, + row.get::<_, i64>(7)?, )) }, ) .optional()?; - if let Some((existing_id, rowid, old_content, old_version, old_blob_hash, old_mime)) = + if let Some((existing_id, rowid, old_content, old_version, old_blob_hash, old_mime, old_metadata, old_size)) = existing { let new_version = old_version + 1; @@ -144,9 +154,9 @@ impl Store { // Snapshot current version to history tx.execute( - "INSERT INTO store_doc_history (id, doc_id, version, content, blob_hash, mime, title, created_at) - VALUES (?1, ?2, ?3, ?4, ?5, ?6, (SELECT title FROM store_documents WHERE id = ?2), ?7)", - params![history_id, existing_id, old_version, old_content, old_blob_hash, old_mime, now], + "INSERT INTO store_doc_history (id, doc_id, version, content, blob_hash, mime, title, metadata, size, created_at) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, (SELECT title FROM store_documents WHERE id = ?2), ?7, ?8, ?9)", + params![history_id, existing_id, old_version, old_content, old_blob_hash, old_mime, old_metadata, old_size, now], )?; tx.execute( @@ -201,6 +211,18 @@ impl Store { params![source, existing_id], )?; } + if let Some(metadata) = &opts.metadata { + tx.execute( + "UPDATE store_documents SET metadata = ?1 WHERE id = ?2", + params![metadata, existing_id], + )?; + } + if let Some(size) = opts.size { + tx.execute( + "UPDATE store_documents SET size = ?1 WHERE id = ?2", + params![size, existing_id], + )?; + } Self::doc_sync_fts(&tx, rowid, content)?; tx.commit()?; @@ -214,16 +236,18 @@ impl Store { let tags_val = opts.tags.unwrap_or_default(); let tags_json = serde_json::to_string(&tags_val).unwrap_or_else(|_| "[]".to_string()); let source = opts.source.unwrap_or_default(); + let metadata = opts.metadata.unwrap_or_else(|| "{}".to_string()); + let size = opts.size.unwrap_or(0); // Insert + FTS sync share this transaction so a failure between // the two can't leave a document without its search index row. tx.execute( "INSERT INTO store_documents - (id, project_id, path, content, title, doc_type, blob_hash, mime, tags, session_id, source, version, created_at, updated_at) - VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, 1, ?12, ?12)", + (id, project_id, path, content, title, doc_type, blob_hash, mime, tags, session_id, source, metadata, size, version, created_at, updated_at) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, 1, ?14, ?14)", params![ id, project_id, path, content, title, doc_type, opts.blob_hash, - mime, tags_json, opts.session_id, source, now + mime, tags_json, opts.session_id, source, metadata, size, now ], )?; let rowid = tx.last_insert_rowid(); @@ -241,6 +265,8 @@ impl Store { tags: tags_val, session_id: opts.session_id, source, + metadata, + size, version: 1, created_at: now, updated_at: now, @@ -253,7 +279,7 @@ impl Store { let conn = self.conn(); conn.query_row( "SELECT id, project_id, path, content, title, doc_type, blob_hash, mime, tags, - session_id, source, version, created_at, updated_at, deleted_at + session_id, source, version, metadata, size, created_at, updated_at, deleted_at FROM store_documents WHERE id = ?1", params![id], Self::row_to_document, @@ -319,7 +345,7 @@ impl Store { pub fn doc_history(&self, doc_id: &str) -> rusqlite::Result> { let conn = self.conn(); let mut stmt = conn.prepare( - "SELECT id, doc_id, version, content, blob_hash, mime, title, created_at + "SELECT id, doc_id, version, content, blob_hash, mime, title, metadata, size, created_at FROM store_doc_history WHERE doc_id = ?1 ORDER BY version DESC", @@ -333,7 +359,9 @@ impl Store { blob_hash: row.get(4)?, mime: row.get(5)?, title: row.get(6)?, - created_at: row.get(7)?, + metadata: row.get(7)?, + size: row.get(8)?, + created_at: row.get(9)?, }) })?; rows.collect() @@ -346,7 +374,7 @@ impl Store { ) -> rusqlite::Result> { let conn = self.conn(); conn.query_row( - "SELECT id, doc_id, version, content, blob_hash, mime, title, created_at + "SELECT id, doc_id, version, content, blob_hash, mime, title, metadata, size, created_at FROM store_doc_history WHERE doc_id = ?1 AND version = ?2", params![doc_id, version], |row| { @@ -358,7 +386,9 @@ impl Store { blob_hash: row.get(4)?, mime: row.get(5)?, title: row.get(6)?, - created_at: row.get(7)?, + metadata: row.get(7)?, + size: row.get(8)?, + created_at: row.get(9)?, }) }, ) @@ -526,7 +556,7 @@ impl Store { let conn = self.conn(); let mut stmt = conn.prepare( "SELECT id, project_id, path, content, title, doc_type, blob_hash, mime, tags, - session_id, source, version, created_at, updated_at, deleted_at + session_id, source, version, metadata, size, created_at, updated_at, deleted_at FROM store_documents WHERE project_id = ?1 AND deleted_at IS NULL ORDER BY path", diff --git a/crates/agentflare-store/src/migrations.rs b/crates/agentflare-store/src/migrations.rs index 032b6334..79128d9c 100644 --- a/crates/agentflare-store/src/migrations.rs +++ b/crates/agentflare-store/src/migrations.rs @@ -80,5 +80,11 @@ pub fn migrations() -> Migrations<'static> { ); CREATE INDEX IF NOT EXISTS idx_doc_history_doc ON store_doc_history(doc_id);", ), + M::up( + "ALTER TABLE store_documents ADD COLUMN metadata TEXT NOT NULL DEFAULT '{}'; + ALTER TABLE store_documents ADD COLUMN size INTEGER NOT NULL DEFAULT 0; + ALTER TABLE store_doc_history ADD COLUMN metadata TEXT NOT NULL DEFAULT '{}'; + ALTER TABLE store_doc_history ADD COLUMN size INTEGER NOT NULL DEFAULT 0;", + ), ]) } diff --git a/src/asset_store.rs b/src/asset_store.rs new file mode 100644 index 00000000..0d5522a9 --- /dev/null +++ b/src/asset_store.rs @@ -0,0 +1,103 @@ +use agentflare_store::Store; + +pub fn entity_path(entity_type: &str, entity_id: &str, filename: &str) -> String { + format!("{entity_type}/{entity_id}/{filename}") +} + +pub fn parse_entity_path(path: &str) -> Option<(&str, &str, &str)> { + let mut parts = path.splitn(3, '/'); + let entity_type = parts.next()?; + let entity_id = parts.next()?; + let filename = parts.next()?; + if entity_type.is_empty() || entity_id.is_empty() || filename.is_empty() { + return None; + } + Some((entity_type, entity_id, filename)) +} + +pub fn document_to_asset_json( + doc: &agentflare_store::documents::Document, +) -> serde_json::Value { + let (entity_type, entity_id, filename) = parse_entity_path(&doc.path) + .unwrap_or(("unknown", "unknown", &doc.path)); + let meta: serde_json::Value = + serde_json::from_str(&doc.metadata).unwrap_or(serde_json::Value::Object(Default::default())); + serde_json::json!({ + "id": doc.id, + "workspace_id": doc.project_id, + "entity_type": entity_type, + "entity_id": entity_id, + "filename": filename, + "size": doc.size, + "mime_type": doc.mime, + "metadata": meta, + "created_at": doc.created_at, + "updated_at": doc.updated_at, + "deleted_at": doc.deleted_at, + "version": doc.version, + }) +} + +pub fn backfill_legacy_assets( + store: &Store, + backend_conn: &rusqlite::Connection, + asset_base_path: &std::path::Path, +) -> Result> { + if let Some(marker) = store.kv_get("_asset_backfill_done")? { + let ts: i64 = serde_json::from_slice(&marker.value)?; + return Err(format!("backfill already ran at {ts}").into()); + } + + let assets = agentflare_backend::asset::list_all(backend_conn)?; + if assets.is_empty() { + let now = db_kit::ids::now(); + store.kv_set( + "_asset_backfill_done", + &serde_json::to_vec(&now)?, + )?; + return Ok(0); + } + + for asset in &assets { + let path = entity_path(&asset.entity_type, &asset.entity_id, &asset.filename); + let bytes = agentflare_backend::asset::read_file(asset_base_path, &asset.storage_path)?; + let blob_hash = store.blob_store(&bytes)?; + + store.doc_upsert_with_opts( + &asset.workspace_id.clone().unwrap_or_default(), + &path, + "", + agentflare_store::documents::DocUpsertOpts { + title: Some(asset.filename.clone()), + doc_type: Some("asset".into()), + blob_hash: Some(blob_hash), + mime: Some(asset.mime_type.clone().unwrap_or_default()), + source: Some("backfill".into()), + metadata: Some(asset.metadata.clone()), + size: Some(asset.size), + ..Default::default() + }, + )?; + } + + let now = db_kit::ids::now(); + store.kv_set( + "_asset_backfill_done", + &serde_json::to_vec(&now)?, + )?; + + Ok(assets.len()) +} + +pub fn get_blob_content( + store: &Store, + doc: &agentflare_store::documents::Document, +) -> Result>, Box> { + match &doc.blob_hash { + Some(hash) => Ok(store.blob_get(hash)?), + None => { + let content = doc.content.as_bytes().to_vec(); + Ok(Some(content)) + } + } +} diff --git a/src/main.rs b/src/main.rs index 828e14b8..d278f78d 100644 --- a/src/main.rs +++ b/src/main.rs @@ -3,6 +3,7 @@ mod agent_install; mod agent_launch; mod agents; mod alias; +mod asset_store; mod artifacts; mod atomic_fs; mod auth; diff --git a/src/mcp_server.rs b/src/mcp_server.rs index 136cb69b..b55336a9 100644 --- a/src/mcp_server.rs +++ b/src/mcp_server.rs @@ -90,6 +90,13 @@ pub struct AgentflareMcp { /// never runs real git worktree/branch operations against this actual /// repository (worktree add, force-remove, branch -D). worktree_repo_root_override: Option, + /// Lazily-opened agentflare-store (documents + blobs), replacing the + /// hand-rolled `assets` table. Persisted across calls so migrations + /// and the one-time backfill run only once per process lifetime. + store: std::sync::Mutex>, + /// Tests inject either a temp path (file-based) or `":memory:"` so + /// they never touch ~/.agentflare/store.db. + store_override: Option, } #[tool_router] @@ -248,12 +255,6 @@ impl AgentflareMcp { } } - fn content_hash(bytes: &[u8]) -> String { - use sha2::Digest; - let digest = sha2::Sha256::digest(bytes); - hex::encode(&digest[..]) - } - fn infer_mime_type(ext: &str) -> String { match ext.to_lowercase().as_str() { "pdf" => "application/pdf".into(), @@ -321,14 +322,6 @@ impl AgentflareMcp { .unwrap_or(1024 * 1024) } - fn strip_storage_path(asset: &agentflare_backend::asset::Asset) -> serde_json::Value { - let mut v = serde_json::to_value(asset).unwrap_or_default(); - if let serde_json::Value::Object(ref mut map) = v { - map.remove("storage_path"); - } - v - } - /// Lock the artifact store + backend pair, resolving the backend on /// first use: reuse an already-running artifact server (flared's /// /artifacts routes, or another session), else bind the fixed port @@ -501,10 +494,13 @@ impl AgentflareMcp { worktree_repo_root: std::path::PathBuf, project_link: std::path::PathBuf, ) -> Self { + let store_dir = backend_db.parent().unwrap().join("store"); + std::fs::create_dir_all(&store_dir).unwrap(); Self { backend_db_override: Some(backend_db), backend_project_link_override: Some(project_link), worktree_repo_root_override: Some(worktree_repo_root), + store_override: Some(store_dir.join("store.db")), ..Default::default() } } @@ -609,6 +605,54 @@ impl AgentflareMcp { Ok(f(guard.as_ref().expect("just initialized above"))) } + /// Open the store (create + migrate) if not yet open. + fn ensure_store(&self) -> Result<(), ErrorData> { + let mut guard = self + .store + .lock() + .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; + if guard.is_some() { + return Ok(()); + } + let path = self + .store_override + .clone() + .unwrap_or_else(crate::store::store_path); + let store = agentflare_store::Store::open_file(&path) + .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; + *guard = Some(store); + Ok(()) + } + + /// Lock the agentflare-store, lazily opening it on first use. + /// After opening, runs the one-time asset backfill (best-effort, + /// skipped if backend_db is already locked by this thread). + fn with_store( + &self, + f: impl FnOnce(&agentflare_store::Store) -> T, + ) -> Result { + self.ensure_store()?; + + // One-time backfill: try_lock to avoid deadlock when called from + // within with_backend_db (attach handler nests store inside backend). + // If backend isn't open yet, skip — with_backend_db will handle it. + if let Ok(bg) = self.backend_db.try_lock() { + if let Some(ref conn) = *bg { + let base_path = crate::paths::home().join(".agentflare"); + let s = self.store.lock().unwrap(); + if let Some(ref store) = *s { + let _ = crate::asset_store::backfill_legacy_assets(store, conn, &base_path); + } + } + } + + let guard = self + .store + .lock() + .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; + Ok(f(guard.as_ref().expect("ensure_store just initialized it"))) + } + /// The one and only workspace on this system: reused if it already /// exists, auto-created (named "default") on first use. Never exposed /// as an MCP parameter. diff --git a/src/mcp_server/asset.rs b/src/mcp_server/asset.rs index af196d47..0d89b53d 100644 --- a/src/mcp_server/asset.rs +++ b/src/mcp_server/asset.rs @@ -1,5 +1,3 @@ -//! `asset` MCP tool handler body -- split out of mcp_server.rs (item #168). - use super::*; impl AgentflareMcp { @@ -27,7 +25,6 @@ impl AgentflareMcp { let fn_val = filename.ok_or_else(|| { ErrorData::invalid_params("filename is required for attach", None) })?; - // path traversal guard: reject filename with .. or absolute components let staged_rel = std::path::Path::new(&fn_val); if staged_rel .components() @@ -66,8 +63,8 @@ impl AgentflareMcp { } let bytes = std::fs::read(&staged) .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; - let hash = Self::content_hash(&bytes); let meta = metadata.unwrap_or_else(|| "{}".to_string()); + self.with_backend_db(|conn| { let ws_id = Self::resolve_workspace_id(conn)?; let (entity_type, entity_id) = if has_item { @@ -84,69 +81,65 @@ impl AgentflareMcp { .and_then(|e| e.to_str()) .unwrap_or(""); let mime = Self::infer_mime_type(ext); - let stem = std::path::Path::new(&fn_val) - .file_stem() - .and_then(|s| s.to_str()) - .unwrap_or(&fn_val); - let safe_stem: String = { - let s: String = stem - .chars() - .filter(|c| c.is_ascii_alphanumeric() || *c == '-' || *c == '_') - .collect(); - if s.is_empty() { "file".to_string() } else { s } - }; - let full_storage = if ext.is_empty() { - format!("{}/assets/{}-{}", ws_id, safe_stem, hash) - } else { - format!("{}/assets/{}-{}.{}", ws_id, safe_stem, hash, ext) - }; - let base_path = crate::paths::home().join(".agentflare"); - // only write if file doesn't already exist (same content already stored) - let target = base_path.join(&full_storage); - if !target.exists() { - agentflare_backend::asset::write_file(&base_path, &full_storage, &bytes) + + let path = crate::asset_store::entity_path(&entity_type, &entity_id, &fn_val); + + let result = self.with_store(|store| -> Result { + let blob_hash = store + .blob_store(&bytes) .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; - } - let asset = agentflare_backend::asset::create( - conn, - agentflare_backend::asset::CreateAsset { - workspace_id: Some(ws_id.clone()), - entity_type: entity_type.into(), - entity_id, - filename: fn_val.clone(), - size: size as i64, - mime_type: Some(mime), - metadata: Some(meta), - storage_path: Some(full_storage), - }, - ) - .map_err(map_backend_err)?; - // remove staging file only after the DB insert succeeds - let _ = std::fs::remove_file(&staged); - Ok( - serde_json::to_string_pretty(&Self::strip_storage_path(&asset)) - .unwrap_or_default(), - ) + + let doc = store + .doc_upsert_with_opts( + &ws_id, + &path, + "", + agentflare_store::documents::DocUpsertOpts { + title: Some(fn_val.clone()), + doc_type: Some("asset".into()), + blob_hash: Some(blob_hash), + mime: Some(mime), + source: Some("attach".into()), + metadata: Some(meta.clone()), + size: Some(size as i64), + ..Default::default() + }, + ) + .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; + + let _ = std::fs::remove_file(&staged); + Ok(crate::asset_store::document_to_asset_json(&doc)) + })??; + + Ok(serde_json::to_string_pretty(&result).unwrap_or_default()) })? } "get" => { let id = id.ok_or_else(|| ErrorData::invalid_params("id is required for get", None))?; - self.with_backend_db(|conn| { - let asset = agentflare_backend::asset::get(conn, &id) - .map_err(map_backend_err)?; - let base_path = crate::paths::home().join(".agentflare"); + self.with_store(|store| -> Result { + let doc = store + .doc_get(&id) + .map_err(|e| ErrorData::internal_error(e.to_string(), None))? + .ok_or_else(|| ErrorData::invalid_params("asset not found", None))?; + let max_inline = Self::asset_max_inline_bytes(); - let meta = Self::strip_storage_path(&asset); - let size = asset.size as u64; + let meta = crate::asset_store::document_to_asset_json(&doc); + let size = doc.size as u64; + if size <= max_inline { - match agentflare_backend::asset::read_file(&base_path, &asset.storage_path) { - Ok(bytes) => { - // Textual MIME + valid UTF-8 => return readable text so - // callers don't decode every text asset; everything else - // (binary MIME, or invalid UTF-8) => Base64. + let content = + crate::asset_store::get_blob_content(store, &doc).map_err(|e| { + ErrorData::internal_error(e.to_string(), None) + })?; + match content { + Some(bytes) => { let (content, encoding) = match std::str::from_utf8(&bytes) { - Ok(text) if Self::mime_is_textual(asset.mime_type.as_deref()) => (text.to_string(), "utf8"), + Ok(text) + if Self::mime_is_textual(Some(&doc.mime)) => + { + (text.to_string(), "utf8") + } _ => (base64_encode(&bytes), "base64"), }; let result = serde_json::json!({ @@ -156,11 +149,11 @@ impl AgentflareMcp { }); Ok(serde_json::to_string_pretty(&result).unwrap_or_default()) } - Err(e) => { + None => { let result = serde_json::json!({ "asset": meta, "content": null, - "content_omitted_reason": format!("could not read file: {}", e), + "content_omitted_reason": "blob content not found in store", }); Ok(serde_json::to_string_pretty(&result).unwrap_or_default()) } @@ -169,66 +162,68 @@ impl AgentflareMcp { let result = serde_json::json!({ "asset": meta, "content": null, - "content_omitted_reason": format!("file is {} bytes, exceeds the {} byte inline limit", size, max_inline), + "content_omitted_reason": format!( + "file is {} bytes, exceeds the {} byte inline limit", + size, max_inline + ), }); Ok(serde_json::to_string_pretty(&result).unwrap_or_default()) } })? } - "list" => self.with_backend_db(|conn| { - let ws_id = Self::resolve_workspace_id(conn)?; - let assets: Vec = match (item_id, project_id) { - (Some(iid), None) => { - agentflare_backend::asset::list_by_entity(conn, "item_attachment", &iid) - .map_err(map_backend_err)? - } - (None, Some(pid)) => { - agentflare_backend::asset::list_by_entity(conn, "project_attachment", &pid) - .map_err(map_backend_err)? - } + "list" => { + let ws_id = match self.with_backend_db(|conn| Self::resolve_workspace_id(conn)) { + Ok(Ok(id)) => id, + Ok(Err(e)) => return Err(ErrorData::internal_error(e.to_string(), None)), + Err(e) => return Err(e), + }; + let prefix = match (item_id, project_id) { + (Some(iid), None) => format!("item_attachment/{iid}"), + (None, Some(pid)) => format!("project_attachment/{pid}"), (Some(_), Some(_)) => { return Err(ErrorData::invalid_params( "only one of item_id or project_id allowed for list, not both", None, )); } - (None, None) => { - let mut assets: Vec = Vec::new(); - for a in agentflare_backend::asset::list_by_workspace(conn, &ws_id) - .map_err(map_backend_err)? - { - assets.push(Self::strip_storage_path(&a)); - } - return Ok(serde_json::to_string_pretty(&assets).unwrap_or_default()); - } + (None, None) => String::new(), }; - let mut stripped: Vec = Vec::new(); - for a in assets { - stripped.push(Self::strip_storage_path(&a)); - } - Ok(serde_json::to_string_pretty(&stripped).unwrap_or_default()) - })?, + + self.with_store(|store| -> Result { + let docs = store + .doc_list(&ws_id) + .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; + + let filtered: Vec = docs + .into_iter() + .filter(|d| { + d.doc_type == "asset" + && (prefix.is_empty() || d.path.starts_with(&prefix)) + }) + .map(|d| crate::asset_store::document_to_asset_json(&d)) + .collect(); + + Ok(serde_json::to_string_pretty(&filtered).unwrap_or_default()) + })? + } "delete" => { let id = id.ok_or_else(|| ErrorData::invalid_params("id is required for delete", None))?; - self.with_backend_db(|conn| { - let asset = agentflare_backend::asset::get(conn, &id) - .map_err(map_backend_err)?; - // soft-delete the row - agentflare_backend::asset::delete(conn, &id) - .map_err(map_backend_err)?; - // only unlink from disk if no other live row references the same storage_path - let remaining: i64 = conn - .query_row( - "SELECT count(*) FROM assets WHERE storage_path = ?1 AND deleted_at IS NULL", - rusqlite::params![&asset.storage_path], - |r| r.get(0), - ) - .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; - if remaining == 0 { - let base_path = crate::paths::home().join(".agentflare"); - let _ = agentflare_backend::asset::delete_file(&base_path, &asset.storage_path); + self.with_store(|store| -> Result { + let doc = store + .doc_get(&id) + .map_err(|e| ErrorData::internal_error(e.to_string(), None))? + .ok_or_else(|| ErrorData::invalid_params("asset not found", None))?; + + if let Some(ref hash) = doc.blob_hash { + store + .blob_unref(hash) + .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; } + store + .doc_delete(&id) + .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; + Ok(serde_json::json!({"deleted": true, "id": id}).to_string()) })? } diff --git a/src/mcp_server/handoff.rs b/src/mcp_server/handoff.rs index a861c0c9..e9812229 100644 --- a/src/mcp_server/handoff.rs +++ b/src/mcp_server/handoff.rs @@ -1,5 +1,3 @@ -//! `handoff` MCP tool handler body -- split out of mcp_server.rs (item #168). - use super::*; impl AgentflareMcp { @@ -88,21 +86,10 @@ impl AgentflareMcp { }; let bytes = content.as_bytes(); - let hash = Self::content_hash(bytes); - // Keyed on item.id, not name — name is the per-call brief and - // can legitimately differ between messages on the same item - // (e.g. a reply's brief vs. the original ask); keying on it - // would silently reset versioning to 1 instead of continuing - // the chain. let safe_stem = Self::slugify(&item.id); - let filename = format!("{safe_stem}.{ext}"); - let full_storage = format!("{ws_id}/assets/{safe_stem}-{hash}.{ext}"); - let base_path = crate::paths::home().join(".agentflare"); - let target = base_path.join(&full_storage); - if !target.exists() { - agentflare_backend::asset::write_file(&base_path, &full_storage, bytes) - .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; - } + let asset_id = db_kit::ids::new_id(); + let filename = format!("{safe_stem}-{asset_id}.{ext}"); + let entity_path = crate::asset_store::entity_path("item_attachment", &item.id, &filename); let mut meta = serde_json::json!({ "sender": self.agent, "recipient": recipient }); if let Some(t) = &thread_id { meta["thread_id"] = serde_json::json!(t); @@ -110,7 +97,6 @@ impl AgentflareMcp { if let Some(r) = &reply_to { meta["reply_to"] = serde_json::json!(r); } - // Session snapshot — lets the recipient see context without extra round-trips. if let Some(s) = summary { meta["session_summary"] = serde_json::json!(s); } @@ -126,24 +112,49 @@ impl AgentflareMcp { if let Some(e) = evidence { meta["evidence"] = serde_json::json!(e); } - let asset = agentflare_backend::asset::create( - conn, - agentflare_backend::asset::CreateAsset { - workspace_id: Some(ws_id), - entity_type: "item_attachment".into(), - entity_id: item.id.clone(), - filename, - size: bytes.len() as i64, - mime_type: Some(Self::infer_mime_type(ext)), - metadata: Some(meta.to_string()), - storage_path: Some(full_storage), - }, - ) - .map_err(map_backend_err)?; - // Knowledge fact import: persist each fact into the recipient's memory, - // scoped to this handoff's project. Runs only after the item/asset are - // committed, so a failed handoff never leaves orphaned facts behind. + let mime_type = Self::infer_mime_type(ext); + + let result = self.with_store(|store| -> Result { + let prefix = format!("item_attachment/{}", item.id); + let existing = store + .doc_list(&ws_id) + .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; + let version = existing.iter().filter(|d| d.path.starts_with(&prefix)).count() as i32 + 1; + + let blob_hash = store + .blob_store(bytes) + .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; + + let doc = store + .doc_upsert_with_opts( + &ws_id, + &entity_path, + "", + agentflare_store::documents::DocUpsertOpts { + title: Some(filename.clone()), + doc_type: Some("asset".into()), + blob_hash: Some(blob_hash), + mime: Some(mime_type.clone()), + source: Some("handoff".into()), + metadata: Some(meta.to_string()), + size: Some(bytes.len() as i64), + + ..Default::default() + }, + ) + .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; + + Ok(serde_json::json!({ + "item_id": item.id, + "item_sequence_id": item.sequence_id, + "asset_id": doc.id, + "asset_version": version, + "recipient": recipient, + })) + })??; + + // Knowledge fact import: persist each fact into the recipient's memory if let Some(ref facts) = facts { let sender = self.agent.as_deref().unwrap_or("unknown"); for fact in facts { @@ -172,19 +183,11 @@ impl AgentflareMcp { scope: None, }; if let Err(e) = crate::memory::mcp::handle_remember(input) { - // Fact import is best-effort — never fail the handoff for a memory write eprintln!("[handoff] fact import failed: {e}"); } } } - let result = serde_json::json!({ - "item_id": item.id, - "item_sequence_id": item.sequence_id, - "asset_id": asset.id, - "asset_version": asset.version, - "recipient": recipient, - }); Ok(serde_json::to_string_pretty(&result).unwrap_or_default()) })? } diff --git a/src/mcp_server/tests/artifact_tests.rs b/src/mcp_server/tests/artifact_tests.rs index ee2fbcab..3e6d32d9 100644 --- a/src/mcp_server/tests/artifact_tests.rs +++ b/src/mcp_server/tests/artifact_tests.rs @@ -339,13 +339,15 @@ fn handoff_tool_requires_recipient_and_assigns_item() { let assets = item_assets(&s, &item_id); assert_eq!(assets.as_array().unwrap().len(), 1); - // Match production's actual transform (AgentflareMcp::slugify), not - // just .to_lowercase() -- nanoid ids can contain `_`, which slugify - // collapses to `-` but to_lowercase() leaves untouched, making the - // naive comparison flaky whenever a generated id contains one. - assert_eq!( - assets[0]["filename"], - format!("{}.md", AgentflareMcp::slugify(&item_id)) + // Filename includes a UUID suffix for uniqueness per handoff + let slug = AgentflareMcp::slugify(&item_id); + assert!( + assets[0]["filename"] + .as_str() + .map(|f| f.starts_with(&format!("{slug}-")) && f.ends_with(".md")) + .unwrap_or(false), + "expected filename to start with '{slug}-*.md' but got {:?}", + assets[0]["filename"] ); }); } diff --git a/src/mcp_server/tests/asset_tests.rs b/src/mcp_server/tests/asset_tests.rs index 1e18e811..a292a4e1 100644 --- a/src/mcp_server/tests/asset_tests.rs +++ b/src/mcp_server/tests/asset_tests.rs @@ -258,67 +258,8 @@ fn asset_shared_storage_delete_safety() { }); } -#[test] -fn asset_content_dedup() { - crate::paths::test_support::with_temp_home(|| { - let (tmp, s) = harness(); - let home = crate::paths::home(); - let staging = home.join(".agentflare").join("staging"); - std::fs::create_dir_all(&staging).unwrap(); - - let item_a: serde_json::Value = - serde_json::from_str(&s.item(Parameters(empty_item_create("dedup-a"))).unwrap()) - .unwrap(); - let item_b: serde_json::Value = - serde_json::from_str(&s.item(Parameters(empty_item_create("dedup-b"))).unwrap()) - .unwrap(); - let id_a = item_a["id"].as_str().unwrap().to_string(); - let id_b = item_b["id"].as_str().unwrap().to_string(); - - let content = b"dedup me please"; - std::fs::write(staging.join("dedup.txt"), content).unwrap(); - s.asset(Parameters(AssetRequest { - action: "attach".into(), - id: None, - item_id: Some(id_a.clone()), - project_id: None, - filename: Some("dedup.txt".into()), - metadata: None, - })) - .unwrap(); - - std::fs::write(staging.join("dedup.txt"), content).unwrap(); - s.asset(Parameters(AssetRequest { - action: "attach".into(), - id: None, - item_id: Some(id_b.clone()), - project_id: None, - filename: Some("dedup.txt".into()), - metadata: None, - })) - .unwrap(); - - // two rows, one file on disk: count unique storage_path values - let conn = backend_conn(&tmp); - let unique_paths: i64 = conn - .query_row( - "SELECT count(DISTINCT storage_path) FROM assets WHERE deleted_at IS NULL", - [], - |r| r.get(0), - ) - .unwrap(); - assert_eq!(unique_paths, 1, "same content must share one storage_path"); - - let total_rows: i64 = conn - .query_row( - "SELECT count(*) FROM assets WHERE deleted_at IS NULL", - [], - |r| r.get(0), - ) - .unwrap(); - assert_eq!(total_rows, 2, "two rows despite one file on disk"); - }); -} +// asset_content_dedup was removed in #185: content dedup via store_blobs +// (blob_store/blob_unref) is tested by the agentflare-store crate itself. #[test] fn asset_attach_to_project() { From dbd386d6a24c6b4a4b23e982784c18861c7d1d1d Mon Sep 17 00:00:00 2001 From: Shivakumar Date: Tue, 21 Jul 2026 00:11:17 +0530 Subject: [PATCH 2/6] fix(asset-store): deleted assets stay gettable, backfill can be skipped forever - Store::doc_get() didn't filter deleted_at, so asset get/delete on an already-deleted id returned 200 instead of not-found (regression from agentflare_backend::asset::get, worse now that delete also physically purges the blob once ref_count hits 0 -- soft-deleted docs pointed at already-gone content). - with_store()'s one-time legacy backfill only fires when backend_db is unlocked at call time; attach always nests with_store inside with_backend_db so it never got a chance there. Trigger it once, unlocked, at the top of attach so first-run migration isn't dependent on call order. - replaced a lock().unwrap() in the backfill path with a graceful skip on a poisoned mutex, matching the rest of with_store/with_backend_db. --- crates/agentflare-store/src/documents.rs | 6 ++- src/mcp_server.rs | 7 +-- src/mcp_server/asset.rs | 5 ++ src/mcp_server/tests/asset_tests.rs | 63 ++++++++++++++++++++++++ 4 files changed, 77 insertions(+), 4 deletions(-) diff --git a/crates/agentflare-store/src/documents.rs b/crates/agentflare-store/src/documents.rs index bf1d6dd7..796588c8 100644 --- a/crates/agentflare-store/src/documents.rs +++ b/crates/agentflare-store/src/documents.rs @@ -280,7 +280,7 @@ impl Store { conn.query_row( "SELECT id, project_id, path, content, title, doc_type, blob_hash, mime, tags, session_id, source, version, metadata, size, created_at, updated_at, deleted_at - FROM store_documents WHERE id = ?1", + FROM store_documents WHERE id = ?1 AND deleted_at IS NULL", params![id], Self::row_to_document, ) @@ -606,6 +606,10 @@ mod tests { let list = s.doc_list("p").unwrap(); assert_eq!(list.len(), 1); assert_eq!(list[0].path, "/a.md"); + + // doc_get must not resurrect a soft-deleted row -- callers (e.g. the + // asset MCP tool) rely on this to report "not found" post-delete. + assert!(s.doc_get(&b.id).unwrap().is_none()); } #[test] diff --git a/src/mcp_server.rs b/src/mcp_server.rs index b55336a9..d6d2298a 100644 --- a/src/mcp_server.rs +++ b/src/mcp_server.rs @@ -639,9 +639,10 @@ impl AgentflareMcp { if let Ok(bg) = self.backend_db.try_lock() { if let Some(ref conn) = *bg { let base_path = crate::paths::home().join(".agentflare"); - let s = self.store.lock().unwrap(); - if let Some(ref store) = *s { - let _ = crate::asset_store::backfill_legacy_assets(store, conn, &base_path); + if let Ok(s) = self.store.lock() { + if let Some(ref store) = *s { + let _ = crate::asset_store::backfill_legacy_assets(store, conn, &base_path); + } } } } diff --git a/src/mcp_server/asset.rs b/src/mcp_server/asset.rs index 0d89b53d..88ee9bdf 100644 --- a/src/mcp_server/asset.rs +++ b/src/mcp_server/asset.rs @@ -65,6 +65,11 @@ impl AgentflareMcp { .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; let meta = metadata.unwrap_or_else(|| "{}".to_string()); + // attach always nests with_store inside with_backend_db below, which + // skips the one-time legacy backfill to avoid a non-reentrant-mutex + // deadlock -- so trigger it here first, while backend_db is unlocked. + self.with_store(|_| ())?; + self.with_backend_db(|conn| { let ws_id = Self::resolve_workspace_id(conn)?; let (entity_type, entity_id) = if has_item { diff --git a/src/mcp_server/tests/asset_tests.rs b/src/mcp_server/tests/asset_tests.rs index a292a4e1..713a6535 100644 --- a/src/mcp_server/tests/asset_tests.rs +++ b/src/mcp_server/tests/asset_tests.rs @@ -161,6 +161,69 @@ fn asset_get_rejects_missing_id() { assert_eq!(err.code, rmcp::model::ErrorCode::INVALID_PARAMS); } +#[test] +fn asset_get_and_delete_after_delete_return_not_found() { + crate::paths::test_support::with_temp_home(|| { + let (_tmp, s) = harness(); + let home = crate::paths::home(); + let staging = home.join(".agentflare").join("staging"); + std::fs::create_dir_all(&staging).unwrap(); + + let item: serde_json::Value = + serde_json::from_str(&s.item(Parameters(empty_item_create("gone"))).unwrap()).unwrap(); + let item_id = item["id"].as_str().unwrap().to_string(); + + std::fs::write(staging.join("gone.txt"), b"bye").unwrap(); + let attached: serde_json::Value = serde_json::from_str( + &s.asset(Parameters(AssetRequest { + action: "attach".into(), + id: None, + item_id: Some(item_id), + project_id: None, + filename: Some("gone.txt".into()), + metadata: None, + })) + .unwrap(), + ) + .unwrap(); + let asset_id = attached["id"].as_str().unwrap().to_string(); + + s.asset(Parameters(AssetRequest { + action: "delete".into(), + id: Some(asset_id.clone()), + item_id: None, + project_id: None, + filename: None, + metadata: None, + })) + .unwrap(); + + // A deleted asset must not be gettable, matching the pre-#185 + // agentflare_backend::asset::get contract (deleted_at IS NULL). + s.asset(Parameters(AssetRequest { + action: "get".into(), + id: Some(asset_id.clone()), + item_id: None, + project_id: None, + filename: None, + metadata: None, + })) + .unwrap_err(); + + // Deleting an already-deleted asset must also report not-found, + // not silently double-unref an already-purged blob. + s.asset(Parameters(AssetRequest { + action: "delete".into(), + id: Some(asset_id), + item_id: None, + project_id: None, + filename: None, + metadata: None, + })) + .unwrap_err(); + }); +} + #[test] fn asset_shared_storage_delete_safety() { crate::paths::test_support::with_temp_home(|| { From c9f5121a61a037d0313fbc347058b455d29ecc8d Mon Sep 17 00:00:00 2001 From: Shivakumar Date: Tue, 21 Jul 2026 00:26:09 +0530 Subject: [PATCH 3/6] fix(asset-store): satisfy CI clippy -D warnings (collapsible-if, redundant-closure) --- src/mcp_server.rs | 16 ++++++++-------- src/mcp_server/asset.rs | 2 +- 2 files changed, 9 insertions(+), 9 deletions(-) diff --git a/src/mcp_server.rs b/src/mcp_server.rs index d6d2298a..db4b7af2 100644 --- a/src/mcp_server.rs +++ b/src/mcp_server.rs @@ -636,14 +636,14 @@ impl AgentflareMcp { // One-time backfill: try_lock to avoid deadlock when called from // within with_backend_db (attach handler nests store inside backend). // If backend isn't open yet, skip — with_backend_db will handle it. - if let Ok(bg) = self.backend_db.try_lock() { - if let Some(ref conn) = *bg { - let base_path = crate::paths::home().join(".agentflare"); - if let Ok(s) = self.store.lock() { - if let Some(ref store) = *s { - let _ = crate::asset_store::backfill_legacy_assets(store, conn, &base_path); - } - } + if let Ok(bg) = self.backend_db.try_lock() + && let Some(ref conn) = *bg + { + let base_path = crate::paths::home().join(".agentflare"); + if let Ok(s) = self.store.lock() + && let Some(ref store) = *s + { + let _ = crate::asset_store::backfill_legacy_assets(store, conn, &base_path); } } diff --git a/src/mcp_server/asset.rs b/src/mcp_server/asset.rs index 88ee9bdf..52771125 100644 --- a/src/mcp_server/asset.rs +++ b/src/mcp_server/asset.rs @@ -177,7 +177,7 @@ impl AgentflareMcp { })? } "list" => { - let ws_id = match self.with_backend_db(|conn| Self::resolve_workspace_id(conn)) { + let ws_id = match self.with_backend_db(Self::resolve_workspace_id) { Ok(Ok(id)) => id, Ok(Err(e)) => return Err(ErrorData::internal_error(e.to_string(), None)), Err(e) => return Err(e), From 0710e83b3c7adcf73a7ce0c832711c44040983d7 Mon Sep 17 00:00:00 2001 From: Shivakumar Date: Tue, 21 Jul 2026 00:30:12 +0530 Subject: [PATCH 4/6] fix(asset-store): remove needless borrow flagged by clippy -D warnings --- src/mcp_server/asset.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/mcp_server/asset.rs b/src/mcp_server/asset.rs index 52771125..38620c3d 100644 --- a/src/mcp_server/asset.rs +++ b/src/mcp_server/asset.rs @@ -87,7 +87,7 @@ impl AgentflareMcp { .unwrap_or(""); let mime = Self::infer_mime_type(ext); - let path = crate::asset_store::entity_path(&entity_type, &entity_id, &fn_val); + let path = crate::asset_store::entity_path(entity_type, &entity_id, &fn_val); let result = self.with_store(|store| -> Result { let blob_hash = store From bd22efcb1fc6c71db1f1f3a4e204fb35bd2e7a0e Mon Sep 17 00:00:00 2001 From: Shivakumar Date: Tue, 21 Jul 2026 00:31:08 +0530 Subject: [PATCH 5/6] style(asset-store): cargo fmt --- crates/agentflare-store/src/documents.rs | 12 ++++- src/asset_store.rs | 22 +++------ src/main.rs | 2 +- src/mcp_server.rs | 5 +- src/mcp_server/asset.rs | 59 +++++++++++------------- src/mcp_server/handoff.rs | 11 +++-- 6 files changed, 55 insertions(+), 56 deletions(-) diff --git a/crates/agentflare-store/src/documents.rs b/crates/agentflare-store/src/documents.rs index 796588c8..5d7ac468 100644 --- a/crates/agentflare-store/src/documents.rs +++ b/crates/agentflare-store/src/documents.rs @@ -146,8 +146,16 @@ impl Store { ) .optional()?; - if let Some((existing_id, rowid, old_content, old_version, old_blob_hash, old_mime, old_metadata, old_size)) = - existing + if let Some(( + existing_id, + rowid, + old_content, + old_version, + old_blob_hash, + old_mime, + old_metadata, + old_size, + )) = existing { let new_version = old_version + 1; let history_id = db_kit::ids::new_id(); diff --git a/src/asset_store.rs b/src/asset_store.rs index 0d5522a9..888e05e5 100644 --- a/src/asset_store.rs +++ b/src/asset_store.rs @@ -15,13 +15,11 @@ pub fn parse_entity_path(path: &str) -> Option<(&str, &str, &str)> { Some((entity_type, entity_id, filename)) } -pub fn document_to_asset_json( - doc: &agentflare_store::documents::Document, -) -> serde_json::Value { - let (entity_type, entity_id, filename) = parse_entity_path(&doc.path) - .unwrap_or(("unknown", "unknown", &doc.path)); - let meta: serde_json::Value = - serde_json::from_str(&doc.metadata).unwrap_or(serde_json::Value::Object(Default::default())); +pub fn document_to_asset_json(doc: &agentflare_store::documents::Document) -> serde_json::Value { + let (entity_type, entity_id, filename) = + parse_entity_path(&doc.path).unwrap_or(("unknown", "unknown", &doc.path)); + let meta: serde_json::Value = serde_json::from_str(&doc.metadata) + .unwrap_or(serde_json::Value::Object(Default::default())); serde_json::json!({ "id": doc.id, "workspace_id": doc.project_id, @@ -51,10 +49,7 @@ pub fn backfill_legacy_assets( let assets = agentflare_backend::asset::list_all(backend_conn)?; if assets.is_empty() { let now = db_kit::ids::now(); - store.kv_set( - "_asset_backfill_done", - &serde_json::to_vec(&now)?, - )?; + store.kv_set("_asset_backfill_done", &serde_json::to_vec(&now)?)?; return Ok(0); } @@ -81,10 +76,7 @@ pub fn backfill_legacy_assets( } let now = db_kit::ids::now(); - store.kv_set( - "_asset_backfill_done", - &serde_json::to_vec(&now)?, - )?; + store.kv_set("_asset_backfill_done", &serde_json::to_vec(&now)?)?; Ok(assets.len()) } diff --git a/src/main.rs b/src/main.rs index d278f78d..140e5612 100644 --- a/src/main.rs +++ b/src/main.rs @@ -3,8 +3,8 @@ mod agent_install; mod agent_launch; mod agents; mod alias; -mod asset_store; mod artifacts; +mod asset_store; mod atomic_fs; mod auth; mod auth_crypt; diff --git a/src/mcp_server.rs b/src/mcp_server.rs index db4b7af2..ba164785 100644 --- a/src/mcp_server.rs +++ b/src/mcp_server.rs @@ -627,10 +627,7 @@ impl AgentflareMcp { /// Lock the agentflare-store, lazily opening it on first use. /// After opening, runs the one-time asset backfill (best-effort, /// skipped if backend_db is already locked by this thread). - fn with_store( - &self, - f: impl FnOnce(&agentflare_store::Store) -> T, - ) -> Result { + fn with_store(&self, f: impl FnOnce(&agentflare_store::Store) -> T) -> Result { self.ensure_store()?; // One-time backfill: try_lock to avoid deadlock when called from diff --git a/src/mcp_server/asset.rs b/src/mcp_server/asset.rs index 38620c3d..da664393 100644 --- a/src/mcp_server/asset.rs +++ b/src/mcp_server/asset.rs @@ -89,32 +89,33 @@ impl AgentflareMcp { let path = crate::asset_store::entity_path(entity_type, &entity_id, &fn_val); - let result = self.with_store(|store| -> Result { - let blob_hash = store - .blob_store(&bytes) - .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; + let result = + self.with_store(|store| -> Result { + let blob_hash = store + .blob_store(&bytes) + .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; - let doc = store - .doc_upsert_with_opts( - &ws_id, - &path, - "", - agentflare_store::documents::DocUpsertOpts { - title: Some(fn_val.clone()), - doc_type: Some("asset".into()), - blob_hash: Some(blob_hash), - mime: Some(mime), - source: Some("attach".into()), - metadata: Some(meta.clone()), - size: Some(size as i64), - ..Default::default() - }, - ) - .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; + let doc = store + .doc_upsert_with_opts( + &ws_id, + &path, + "", + agentflare_store::documents::DocUpsertOpts { + title: Some(fn_val.clone()), + doc_type: Some("asset".into()), + blob_hash: Some(blob_hash), + mime: Some(mime), + source: Some("attach".into()), + metadata: Some(meta.clone()), + size: Some(size as i64), + ..Default::default() + }, + ) + .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; - let _ = std::fs::remove_file(&staged); - Ok(crate::asset_store::document_to_asset_json(&doc)) - })??; + let _ = std::fs::remove_file(&staged); + Ok(crate::asset_store::document_to_asset_json(&doc)) + })??; Ok(serde_json::to_string_pretty(&result).unwrap_or_default()) })? @@ -133,16 +134,12 @@ impl AgentflareMcp { let size = doc.size as u64; if size <= max_inline { - let content = - crate::asset_store::get_blob_content(store, &doc).map_err(|e| { - ErrorData::internal_error(e.to_string(), None) - })?; + let content = crate::asset_store::get_blob_content(store, &doc) + .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; match content { Some(bytes) => { let (content, encoding) = match std::str::from_utf8(&bytes) { - Ok(text) - if Self::mime_is_textual(Some(&doc.mime)) => - { + Ok(text) if Self::mime_is_textual(Some(&doc.mime)) => { (text.to_string(), "utf8") } _ => (base64_encode(&bytes), "base64"), diff --git a/src/mcp_server/handoff.rs b/src/mcp_server/handoff.rs index e9812229..b795982c 100644 --- a/src/mcp_server/handoff.rs +++ b/src/mcp_server/handoff.rs @@ -89,7 +89,8 @@ impl AgentflareMcp { let safe_stem = Self::slugify(&item.id); let asset_id = db_kit::ids::new_id(); let filename = format!("{safe_stem}-{asset_id}.{ext}"); - let entity_path = crate::asset_store::entity_path("item_attachment", &item.id, &filename); + let entity_path = + crate::asset_store::entity_path("item_attachment", &item.id, &filename); let mut meta = serde_json::json!({ "sender": self.agent, "recipient": recipient }); if let Some(t) = &thread_id { meta["thread_id"] = serde_json::json!(t); @@ -120,7 +121,11 @@ impl AgentflareMcp { let existing = store .doc_list(&ws_id) .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; - let version = existing.iter().filter(|d| d.path.starts_with(&prefix)).count() as i32 + 1; + let version = existing + .iter() + .filter(|d| d.path.starts_with(&prefix)) + .count() as i32 + + 1; let blob_hash = store .blob_store(bytes) @@ -139,7 +144,7 @@ impl AgentflareMcp { source: Some("handoff".into()), metadata: Some(meta.to_string()), size: Some(bytes.len() as i64), - + ..Default::default() }, ) From 4a19ec27d2cc6508067c0ca1b74ec4663bc2454a Mon Sep 17 00:00:00 2001 From: Shivakumar Date: Tue, 21 Jul 2026 00:50:30 +0530 Subject: [PATCH 6/6] fix(asset-store): address CodeRabbit review findings on PR #282 - backfill_legacy_assets: don't abort the whole batch on one bad legacy asset; skip and log it, migrate the rest, and always write the completion marker. Previously a single unreadable asset meant the marker never got written, so every later with_store() call replayed the entire batch, re-upserting already-migrated docs and bumping their version/history each time. - attach: open backend_db (and release it) before triggering the backfill check, so a totally fresh instance's very first attach call can actually backfill, not just the second one onward. - delete: soft-delete the document before unref'ing its blob, so a failure between the two steps can't leave the last reference's content purged while the asset still reads as live. Agentflare-Agent: claude-code_2-1-215_agent Agentflare-Branch: item-185-asset-store-migration --- src/asset_store.rs | 121 +++++++++++++++++++++++++++++++--------- src/mcp_server/asset.rs | 14 ++++- 2 files changed, 107 insertions(+), 28 deletions(-) diff --git a/src/asset_store.rs b/src/asset_store.rs index 888e05e5..3236358f 100644 --- a/src/asset_store.rs +++ b/src/asset_store.rs @@ -47,38 +47,50 @@ pub fn backfill_legacy_assets( } let assets = agentflare_backend::asset::list_all(backend_conn)?; - if assets.is_empty() { - let now = db_kit::ids::now(); - store.kv_set("_asset_backfill_done", &serde_json::to_vec(&now)?)?; - return Ok(0); - } + let mut migrated = 0usize; + // A single unreadable/corrupt legacy asset must not abort the whole + // batch: that would leave the marker unwritten, so every later + // with_store() call would replay it -- re-upserting the assets that + // already migrated fine and bumping their version/history each time. for asset in &assets { - let path = entity_path(&asset.entity_type, &asset.entity_id, &asset.filename); - let bytes = agentflare_backend::asset::read_file(asset_base_path, &asset.storage_path)?; - let blob_hash = store.blob_store(&bytes)?; - - store.doc_upsert_with_opts( - &asset.workspace_id.clone().unwrap_or_default(), - &path, - "", - agentflare_store::documents::DocUpsertOpts { - title: Some(asset.filename.clone()), - doc_type: Some("asset".into()), - blob_hash: Some(blob_hash), - mime: Some(asset.mime_type.clone().unwrap_or_default()), - source: Some("backfill".into()), - metadata: Some(asset.metadata.clone()), - size: Some(asset.size), - ..Default::default() - }, - )?; + match backfill_one(store, asset_base_path, asset) { + Ok(()) => migrated += 1, + Err(e) => eprintln!("[asset-store backfill] skipping asset {}: {e}", asset.id), + } } let now = db_kit::ids::now(); store.kv_set("_asset_backfill_done", &serde_json::to_vec(&now)?)?; - Ok(assets.len()) + Ok(migrated) +} + +fn backfill_one( + store: &Store, + asset_base_path: &std::path::Path, + asset: &agentflare_backend::asset::Asset, +) -> Result<(), Box> { + let path = entity_path(&asset.entity_type, &asset.entity_id, &asset.filename); + let bytes = agentflare_backend::asset::read_file(asset_base_path, &asset.storage_path)?; + let blob_hash = store.blob_store(&bytes)?; + + store.doc_upsert_with_opts( + &asset.workspace_id.clone().unwrap_or_default(), + &path, + "", + agentflare_store::documents::DocUpsertOpts { + title: Some(asset.filename.clone()), + doc_type: Some("asset".into()), + blob_hash: Some(blob_hash), + mime: Some(asset.mime_type.clone().unwrap_or_default()), + source: Some("backfill".into()), + metadata: Some(asset.metadata.clone()), + size: Some(asset.size), + ..Default::default() + }, + )?; + Ok(()) } pub fn get_blob_content( @@ -93,3 +105,62 @@ pub fn get_blob_content( } } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn backfill_skips_bad_asset_and_still_marks_done() { + let tmp = tempfile::tempdir().unwrap(); + let base_path = tmp.path().join("home"); + std::fs::create_dir_all(&base_path).unwrap(); + + let conn = agentflare_backend::db::open_db(&tmp.path().join("backend.db")).unwrap(); + let good_storage = "ws/assets/good-hash.txt"; + agentflare_backend::asset::write_file(&base_path, good_storage, b"hello").unwrap(); + agentflare_backend::asset::create( + &conn, + agentflare_backend::asset::CreateAsset { + workspace_id: Some("ws".into()), + entity_type: "item_attachment".into(), + entity_id: "item-good".into(), + filename: "good.txt".into(), + size: 5, + mime_type: Some("text/plain".into()), + metadata: None, + storage_path: Some(good_storage.into()), + }, + ) + .unwrap(); + // No file written for this one -- read_file will fail, simulating a + // corrupt/missing legacy asset. + agentflare_backend::asset::create( + &conn, + agentflare_backend::asset::CreateAsset { + workspace_id: Some("ws".into()), + entity_type: "item_attachment".into(), + entity_id: "item-bad".into(), + filename: "bad.txt".into(), + size: 5, + mime_type: Some("text/plain".into()), + metadata: None, + storage_path: Some("ws/assets/missing-hash.txt".into()), + }, + ) + .unwrap(); + + let store = Store::open_memory().unwrap(); + let migrated = backfill_legacy_assets(&store, &conn, &base_path).unwrap(); + assert_eq!(migrated, 1, "only the readable asset should migrate"); + + let docs = store.doc_list("ws").unwrap(); + assert_eq!(docs.len(), 1); + assert_eq!(docs[0].path, "item_attachment/item-good/good.txt"); + + // The completion marker must be written even though one asset + // failed, so a later with_store() call doesn't replay this batch. + let err = backfill_legacy_assets(&store, &conn, &base_path).unwrap_err(); + assert!(err.to_string().contains("already ran")); + } +} diff --git a/src/mcp_server/asset.rs b/src/mcp_server/asset.rs index da664393..6ca89a36 100644 --- a/src/mcp_server/asset.rs +++ b/src/mcp_server/asset.rs @@ -68,6 +68,10 @@ impl AgentflareMcp { // attach always nests with_store inside with_backend_db below, which // skips the one-time legacy backfill to avoid a non-reentrant-mutex // deadlock -- so trigger it here first, while backend_db is unlocked. + // Open backend_db first (lock is released immediately after) so the + // backfill's try_lock below sees a real connection on a totally fresh + // instance's very first attach call, not just on the second+. + self.with_backend_db(|_| ())?; self.with_store(|_| ())?; self.with_backend_db(|conn| { @@ -217,14 +221,18 @@ impl AgentflareMcp { .map_err(|e| ErrorData::internal_error(e.to_string(), None))? .ok_or_else(|| ErrorData::invalid_params("asset not found", None))?; + // Delete the document row before releasing the blob ref: if + // blob_unref ran first and doc_delete then failed, the last + // reference's content would already be gone while the asset + // still showed up as live (not soft-deleted). + store + .doc_delete(&id) + .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; if let Some(ref hash) = doc.blob_hash { store .blob_unref(hash) .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; } - store - .doc_delete(&id) - .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; Ok(serde_json::json!({"deleted": true, "id": id}).to_string()) })?