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..5d7ac468 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,22 +139,32 @@ 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)) = - 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(); // 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 +219,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 +244,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 +273,8 @@ impl Store { tags: tags_val, session_id: opts.session_id, source, + metadata, + size, version: 1, created_at: now, updated_at: now, @@ -253,8 +287,8 @@ 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 - FROM store_documents WHERE id = ?1", + session_id, source, version, metadata, size, created_at, updated_at, deleted_at + FROM store_documents WHERE id = ?1 AND deleted_at IS NULL", params![id], Self::row_to_document, ) @@ -319,7 +353,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 +367,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 +382,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 +394,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 +564,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", @@ -576,6 +614,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/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..3236358f --- /dev/null +++ b/src/asset_store.rs @@ -0,0 +1,166 @@ +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)?; + 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 { + 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(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( + 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)) + } + } +} + +#[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/main.rs b/src/main.rs index 9e719f7d..6e11def4 100644 --- a/src/main.rs +++ b/src/main.rs @@ -4,6 +4,7 @@ mod agent_launch; mod agents; mod alias; 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 5a4c9dd8..35483e77 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,52 @@ 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() + && 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); + } + } + + 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..6ca89a36 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,17 @@ 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()); + + // 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| { let ws_id = Self::resolve_workspace_id(conn)?; let (entity_type, entity_id) = if has_item { @@ -84,69 +90,62 @@ 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) - .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 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 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 +155,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 +168,72 @@ 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(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), + }; + 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), - ) + 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))?; + + // 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 remaining == 0 { - let base_path = crate::paths::home().join(".agentflare"); - let _ = agentflare_backend::asset::delete_file(&base_path, &asset.storage_path); + if let Some(ref hash) = doc.blob_hash { + store + .blob_unref(hash) + .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..b795982c 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,11 @@ 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 +98,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 +113,53 @@ 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 +188,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..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(|| { @@ -258,67 +321,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() {