From 653e6e1436e1ba9674cae1b6a10fb9e2c059bdca Mon Sep 17 00:00:00 2001 From: YellowSnnowmann <167776381+YellowSnnowmann@users.noreply.github.com> Date: Thu, 9 Jul 2026 07:23:01 +0530 Subject: [PATCH] fix(memory_store): self-heal drifted content-sha tokens + prune vanished folder items (#4700) Co-authored-by: Steven Enamakel --- src/openhuman/memory_sources/sync.rs | 312 ++++++++++++++++++ src/openhuman/memory_store/chunks/store.rs | 149 +++++++++ src/openhuman/memory_store/content/atomic.rs | 183 ++++++++-- .../memory_store/content/compose/yaml.rs | 53 +++ src/openhuman/memory_store/content/mod.rs | 66 ++-- src/openhuman/memory_store/content/read.rs | 132 ++++++-- src/openhuman/memory_store/content/tags.rs | 59 ++++ 7 files changed, 878 insertions(+), 76 deletions(-) diff --git a/src/openhuman/memory_sources/sync.rs b/src/openhuman/memory_sources/sync.rs index 377b9765a..5df2b6cf2 100644 --- a/src/openhuman/memory_sources/sync.rs +++ b/src/openhuman/memory_sources/sync.rs @@ -318,6 +318,42 @@ async fn sync_items_individually( let items = reader.list_items(source, config).await?; let total = items.len(); + // Reconcile before re-ingesting: for LOCAL FOLDER sources only, drop chunks + // (rows + on-disk bodies) for items previously ingested under this source + // that no longer exist on disk — e.g. a renamed or deleted file. Each item + // ingests under the content-addressed composite id + // `mem_src:{source_id}:{item_id}` as a Document, so a rename mints a fresh id + // and orphans the old chunk + its on-disk body; without this the stale body + // lingers forever and can only ever be served as a ≤500-char preview (#4689). + // + // Runs BEFORE the empty-listing early return so an emptied folder (every file + // deleted → total == 0) still reconciles instead of leaving all its chunks + // behind. This is safe on an empty or partial listing because + // `prune_vanished_items` re-checks each candidate on disk and only deletes + // files that are provably absent, so a transient listing miss (EACCES / + // EMFILE / stat stall) is never mistaken for a deletion. + // + // Restricted to Folder: for feed / web / conversation sources, absence from + // the current listing means "rolled off / not re-fetched", NOT "deleted", so + // pruning them would irrecoverably delete valid archived items. + if source_supports_prune(&source.kind) { + if let Some(base_path) = source.path.clone() { + let config = config.clone(); + let source_id = source.id.clone(); + let live: HashSet = items + .iter() + .map(|item| format!("mem_src:{source_id}:{}", item.id)) + .collect(); + if let Err(e) = tokio::task::spawn_blocking(move || { + prune_vanished_items(&config, &source_id, std::path::Path::new(&base_path), &live) + }) + .await + { + tracing::warn!(error = %e, "[memory_sources:sync] prune join error"); + } + } + } + if total == 0 { return Ok(0); } @@ -416,6 +452,107 @@ async fn sync_items_individually( Ok(ingested.load(Ordering::Relaxed)) } +/// Whether a source kind has authoritative present/absent semantics on the +/// local disk, so that "absent from the current listing" genuinely means +/// "deleted" and can drive a prune (#4689). +/// +/// Only `Folder` qualifies: for `RssFeed` / `WebPage`, `list_items` returns a +/// rolling, `max_items`-truncated window, so an item missing from it merely +/// rolled off the feed and must never be deleted; `Conversation` threads have +/// no delete-follows-listing contract. Restricting prune here prevents turning +/// an append-only archive into a destructive mirror of the latest window. +fn source_supports_prune(kind: &SourceKind) -> bool { + matches!(kind, SourceKind::Folder) +} + +/// Delete chunks (rows + on-disk bodies) for items previously ingested under +/// `source_id` whose backing file no longer exists under `base_path` — the +/// reconcile step that keeps a folder resync from orphaning renamed or deleted +/// files (#4689). +/// +/// Safety: a candidate is deleted ONLY when its file is provably absent +/// (`symlink_metadata` returns `NotFound`). A transient listing miss (the reader +/// dropped a still-present file on an `EACCES` / `EMFILE` / stat stall) leaves +/// the file on disk, so the re-check keeps it — absence from `live` alone is +/// never sufficient to delete. +/// +/// Blocking DB + FS work; call from `spawn_blocking`. Chunks land under the +/// content (`Document`) `SourceKind`, not the outer `memory_sources` source kind. +fn prune_vanished_items( + config: &Config, + source_id: &str, + base_path: &std::path::Path, + live: &HashSet, +) { + use crate::openhuman::memory_store::chunks::store as chunk_store; + use crate::openhuman::memory_store::chunks::types::SourceKind as ChunkSourceKind; + + let prefix = format!("mem_src:{source_id}:"); + let previously = match chunk_store::list_source_ids_with_prefix( + config, + ChunkSourceKind::Document, + &prefix, + ) { + Ok(ids) => ids, + Err(e) => { + tracing::warn!( + source_id = %source_id, + error = %format!("{e:#}"), + "[memory_sources:sync] prune: failed to list previously-ingested items" + ); + return; + } + }; + + let mut removed_chunks = 0usize; + let mut removed_items = 0usize; + for stale in previously.into_iter().filter(|sid| !live.contains(sid)) { + // Recover the item's relative path from the composite id and confirm the + // file is genuinely gone before deleting. Anything other than a definite + // NotFound (present file, or an ambiguous EACCES/IO error) is treated as + // "keep" so a transient listing miss can never delete live data. + let Some(rel) = stale.strip_prefix(&prefix) else { + continue; + }; + // Defense-in-depth: `rel` comes from a stored composite id. If it were + // ever empty, absolute, or contained `..`, `base_path.join(rel)` could + // resolve outside the source folder (on Unix an absolute `rel` silently + // discards `base_path`), and a `NotFound` there would delete real chunk + // rows. The current folder reader can't produce such ids, so keep any + // such candidate rather than risk deleting on a path we can't vouch for. + if rel.is_empty() || rel.contains("..") || std::path::Path::new(rel).is_absolute() { + tracing::warn!( + source_id = %source_id, + "[memory_sources:sync] prune: skipping candidate with unsafe relative path" + ); + continue; + } + match std::fs::symlink_metadata(base_path.join(rel)) { + Err(e) if e.kind() == std::io::ErrorKind::NotFound => { /* absent → prune */ } + _ => continue, + } + match chunk_store::delete_chunks_by_source(config, ChunkSourceKind::Document, &stale) { + Ok(n) => { + removed_chunks += n; + removed_items += 1; + } + Err(e) => tracing::warn!( + source_id = %source_id, + error = %format!("{e:#}"), + "[memory_sources:sync] prune: delete failed for a vanished item" + ), + } + } + if removed_items > 0 { + tracing::info!( + source_id = %source_id, + items = removed_items, + chunks = removed_chunks, + "[memory_sources:sync] pruned chunks for vanished items" + ); + } +} + /// Derive the tree scope(s) for a source and reconcile any raw files that /// are not yet covered by tree summaries (incremental — see /// `memory_sync::sources::rebuild`). @@ -539,3 +676,178 @@ pub(crate) fn derive_scopes(source: &MemorySourceEntry, config: &Config) -> Vec< _ => Vec::new(), } } + +#[cfg(test)] +mod tests { + use super::*; + use crate::openhuman::memory_store::chunks::store as chunk_store; + use crate::openhuman::memory_store::chunks::types::{ + chunk_id, Chunk, Metadata, SourceKind as ChunkSourceKind, + }; + use crate::openhuman::memory_store::content::stage_chunks; + use chrono::TimeZone; + use tempfile::TempDir; + + fn doc_chunk(source_id: &str) -> Chunk { + let ts = chrono::Utc.timestamp_millis_opt(1_700_000_000_000).unwrap(); + Chunk { + id: chunk_id(ChunkSourceKind::Document, source_id, 0, "body"), + content: format!("body of {source_id}"), + metadata: Metadata { + source_kind: ChunkSourceKind::Document, + source_id: source_id.into(), + owner: "user".into(), + timestamp: ts, + time_range: (ts, ts), + tags: vec![], + source_ref: None, + path_scope: None, + }, + token_count: 2, + seq_in_source: 0, + created_at: ts, + partial_message: false, + } + } + + fn seed(cfg: &Config, chunk: &Chunk) { + let staged = + stage_chunks(&cfg.memory_tree_content_root(), std::slice::from_ref(chunk)).unwrap(); + chunk_store::with_connection(cfg, |conn| { + let tx = conn.unchecked_transaction()?; + chunk_store::upsert_staged_chunks_tx(&tx, &staged)?; + tx.commit()?; + Ok(()) + }) + .unwrap(); + } + + #[test] + fn source_supports_prune_only_for_folder() { + assert!(source_supports_prune(&SourceKind::Folder)); + assert!(!source_supports_prune(&SourceKind::RssFeed)); + assert!(!source_supports_prune(&SourceKind::WebPage)); + assert!(!source_supports_prune(&SourceKind::Conversation)); + } + + #[test] + fn prune_vanished_items_removes_only_files_absent_from_disk() { + let tmp = TempDir::new().unwrap(); + let mut cfg = Config::default(); + cfg.workspace_dir = tmp.path().to_path_buf(); + // Folder base holds only b.md on disk; a.md was renamed/deleted. + let base = tmp.path().join("folder"); + std::fs::create_dir_all(&base).unwrap(); + std::fs::write(base.join("b.md"), b"b").unwrap(); + + seed(&cfg, &doc_chunk("mem_src:src1:a.md")); + seed(&cfg, &doc_chunk("mem_src:src1:b.md")); + assert_eq!(chunk_store::count_chunks(&cfg).unwrap(), 2); + + let live: HashSet = ["mem_src:src1:b.md".to_string()].into_iter().collect(); + prune_vanished_items(&cfg, "src1", &base, &live); + + let remaining = chunk_store::list_source_ids_with_prefix( + &cfg, + ChunkSourceKind::Document, + "mem_src:src1:", + ) + .unwrap(); + assert_eq!(remaining, vec!["mem_src:src1:b.md".to_string()]); + assert_eq!(chunk_store::count_chunks(&cfg).unwrap(), 1); + } + + #[test] + fn prune_vanished_items_keeps_still_present_file_missed_by_listing() { + // Safety guard (#4689 review): a transient listing miss must not delete a + // file that is still on disk. 'a.md' is absent from `live` but present on + // disk → it must be kept. + let tmp = TempDir::new().unwrap(); + let mut cfg = Config::default(); + cfg.workspace_dir = tmp.path().to_path_buf(); + let base = tmp.path().join("folder"); + std::fs::create_dir_all(&base).unwrap(); + std::fs::write(base.join("a.md"), b"still here").unwrap(); + + seed(&cfg, &doc_chunk("mem_src:src3:a.md")); + // 'a.md' dropped from the listing (e.g. EACCES/EMFILE) though it exists. + let live: HashSet = HashSet::new(); + prune_vanished_items(&cfg, "src3", &base, &live); + + // Not deleted — the on-disk re-check kept it. + assert_eq!(chunk_store::count_chunks(&cfg).unwrap(), 1); + } + + #[test] + fn prune_vanished_items_prunes_when_folder_emptied() { + // Emptied folder: listing is empty (live = {}) and no files remain on + // disk, so the previously-ingested chunk must be pruned (the reason prune + // now runs before the total == 0 early return). + let tmp = TempDir::new().unwrap(); + let mut cfg = Config::default(); + cfg.workspace_dir = tmp.path().to_path_buf(); + let base = tmp.path().join("empty_folder"); + std::fs::create_dir_all(&base).unwrap(); + + seed(&cfg, &doc_chunk("mem_src:src4:gone.md")); + let live: HashSet = HashSet::new(); + prune_vanished_items(&cfg, "src4", &base, &live); + assert_eq!(chunk_store::count_chunks(&cfg).unwrap(), 0); + } + + #[test] + fn prune_vanished_items_keeps_candidate_with_unsafe_relative_path() { + // Defense-in-depth: a stored id whose relative path is absolute must not + // be pruned even when its file is "absent" (join would escape the base). + let tmp = TempDir::new().unwrap(); + let mut cfg = Config::default(); + cfg.workspace_dir = tmp.path().to_path_buf(); + let base = tmp.path().join("folder"); + std::fs::create_dir_all(&base).unwrap(); + + seed(&cfg, &doc_chunk("mem_src:src5:/etc/hostname")); + let live: HashSet = HashSet::new(); + prune_vanished_items(&cfg, "src5", &base, &live); + // Not deleted — the unsafe-path guard skipped it. + assert_eq!(chunk_store::count_chunks(&cfg).unwrap(), 1); + } + + #[test] + fn list_source_ids_with_prefix_isolates_sibling_prefixes() { + // `mem_src:src1:` must not match `mem_src:src10:` items. + let tmp = TempDir::new().unwrap(); + let mut cfg = Config::default(); + cfg.workspace_dir = tmp.path().to_path_buf(); + + seed(&cfg, &doc_chunk("mem_src:src1:a.md")); + seed(&cfg, &doc_chunk("mem_src:src1:c.md")); + seed(&cfg, &doc_chunk("mem_src:src10:b.md")); + + let mut got = chunk_store::list_source_ids_with_prefix( + &cfg, + ChunkSourceKind::Document, + "mem_src:src1:", + ) + .unwrap(); + got.sort(); + assert_eq!( + got, + vec![ + "mem_src:src1:a.md".to_string(), + "mem_src:src1:c.md".to_string() + ] + ); + } + + #[test] + fn prune_vanished_items_is_noop_when_all_live() { + let tmp = TempDir::new().unwrap(); + let mut cfg = Config::default(); + cfg.workspace_dir = tmp.path().to_path_buf(); + + seed(&cfg, &doc_chunk("mem_src:src2:a.md")); + let live: HashSet = ["mem_src:src2:a.md".to_string()].into_iter().collect(); + prune_vanished_items(&cfg, "src2", tmp.path(), &live); + assert_eq!(chunk_store::count_chunks(&cfg).unwrap(), 1); + } +} diff --git a/src/openhuman/memory_store/chunks/store.rs b/src/openhuman/memory_store/chunks/store.rs index ec139fe14..d174769c9 100644 --- a/src/openhuman/memory_store/chunks/store.rs +++ b/src/openhuman/memory_store/chunks/store.rs @@ -493,6 +493,155 @@ pub(crate) fn upsert_staged_chunks_tx( Ok(staged.len()) } +/// Repair the stored body-sha token for one chunk (#4689). +/// +/// The chunk content file is content-addressed and atomically written, so it is +/// the source of truth for its body. When a read detects that the on-disk body +/// no longer hashes to the recorded `content_sha256` (e.g. an external editor +/// rewrote a synced file after ingest), the reader serves the full on-disk body +/// and calls this to re-point the stale token at the disk bytes so the next read +/// verifies cleanly instead of falling back to the ≤500-char preview. +pub fn update_chunk_content_sha256( + config: &Config, + chunk_id: &str, + new_sha256: &str, +) -> Result<()> { + with_connection(config, |conn| { + conn.execute( + "UPDATE mem_tree_chunks SET content_sha256 = ?1 WHERE id = ?2", + params![new_sha256, chunk_id], + )?; + Ok(()) + }) +} + +/// Summary counterpart of [`update_chunk_content_sha256`] (#4689). +pub fn update_summary_content_sha256( + config: &Config, + summary_id: &str, + new_sha256: &str, +) -> Result<()> { + with_connection(config, |conn| { + conn.execute( + "UPDATE mem_tree_summaries SET content_sha256 = ?1 WHERE id = ?2", + params![new_sha256, summary_id], + )?; + Ok(()) + }) +} + +/// List the distinct `source_id`s of chunk rows for `source_kind` whose id +/// starts with `source_id_prefix` (#4689). +/// +/// Used by folder/list-based resync to discover previously-ingested items that +/// have since vanished (e.g. a renamed or deleted file) so their stale rows + +/// on-disk bodies can be cleaned. The prefix is applied Rust-side (literal, not +/// a SQL `LIKE`) so ids containing `_`/`%` are matched verbatim — matching the +/// convention in [`delete_chunks_by_source_prefix`]. +pub fn list_source_ids_with_prefix( + config: &Config, + source_kind: SourceKind, + source_id_prefix: &str, +) -> Result> { + with_connection(config, |conn| { + // Bounded range scan on the `idx_mem_tree_chunks_source (source_kind, + // source_id)` index instead of scanning every source_id for the kind and + // filtering in Rust: `source_id >= prefix AND source_id < upper`, where + // `upper` is the prefix with its last byte incremented. With SQLite's + // default BINARY collation this is exactly the literal byte-prefix set + // (so `_`/`%` stay literal, matching `delete_chunks_by_source_prefix`). + let out = match prefix_upper_bound(source_id_prefix) { + Some(upper) => { + let mut stmt = conn.prepare( + "SELECT DISTINCT source_id FROM mem_tree_chunks \ + WHERE source_kind = ?1 AND source_id >= ?2 AND source_id < ?3", + )?; + let rows = stmt.query_map( + params![source_kind.as_str(), source_id_prefix, upper], + |row| row.get::<_, String>(0), + )?; + rows.collect::>>() + } + None => { + // Empty prefix (or all-0xFF): no finite upper bound. The lower + // bound still uses the index; every row of the kind qualifies — + // identical to the previous `starts_with("")` behaviour. + let mut stmt = conn.prepare( + "SELECT DISTINCT source_id FROM mem_tree_chunks \ + WHERE source_kind = ?1 AND source_id >= ?2", + )?; + let rows = stmt + .query_map(params![source_kind.as_str(), source_id_prefix], |row| { + row.get::<_, String>(0) + })?; + rows.collect::>>() + } + } + .context("Failed to list mem_tree chunk source_ids by prefix")?; + Ok(out) + }) +} + +/// Exclusive upper bound for a literal byte-prefix range scan: the least string +/// strictly greater than every string starting with `prefix`. `None` when no +/// finite bound exists (empty prefix, or every byte is `0xFF`), in which case +/// the caller applies only the lower bound. +fn prefix_upper_bound(prefix: &str) -> Option { + let mut bytes = prefix.as_bytes().to_vec(); + while let Some(&last) = bytes.last() { + if last < 0xFF { + *bytes.last_mut().unwrap() = last + 1; + // Our prefixes are ASCII (`mem_src::`), so incrementing the last + // byte yields valid UTF-8. If a future caller passes a prefix that + // ends mid-codepoint, fall back to the lower-bound-only path rather + // than emit invalid UTF-8. + return String::from_utf8(bytes).ok(); + } + bytes.pop(); + } + None +} + +#[cfg(test)] +mod prefix_bound_tests { + use super::prefix_upper_bound; + + #[test] + fn upper_bound_increments_last_byte() { + assert_eq!( + prefix_upper_bound("mem_src:s1:").as_deref(), + Some("mem_src:s1;") + ); + assert_eq!(prefix_upper_bound("a").as_deref(), Some("b")); + } + + #[test] + fn upper_bound_none_for_empty_prefix() { + // Empty prefix is the only reachable `None` for a valid UTF-8 `&str` (no + // valid string ends in 0xFF); the all-0xFF `pop` loop is defensive. + assert_eq!(prefix_upper_bound(""), None); + } + + #[test] + fn upper_bound_handles_multibyte_prefix() { + // A prefix ending in a multibyte codepoint still yields a valid bound. + assert_eq!(prefix_upper_bound("café").as_deref(), Some("cafê")); + } + + #[test] + fn range_excludes_sibling_prefix() { + // The trailing delimiter in the prefix is what keeps `mem_src:s1:` from + // matching `mem_src:s10:...`: `s10:` sorts below `s1:` (`'0' < ':'`), so + // it falls below the lower bound and is excluded — same as `starts_with`. + let prefix = "mem_src:s1:"; + let upper = prefix_upper_bound(prefix).unwrap(); + let sib = "mem_src:s10:x"; + assert!(!(sib >= prefix && *sib < *upper.as_str())); + let child = "mem_src:s1:a"; + assert!(child >= prefix && *child < *upper.as_str()); + } +} + fn upsert_chunks_with_statement( stmt: &mut rusqlite::Statement<'_>, chunks: &[Chunk], diff --git a/src/openhuman/memory_store/content/atomic.rs b/src/openhuman/memory_store/content/atomic.rs index bf1b1d460..a54b54131 100644 --- a/src/openhuman/memory_store/content/atomic.rs +++ b/src/openhuman/memory_store/content/atomic.rs @@ -42,10 +42,13 @@ pub fn write_if_new(abs_path: &Path, bytes: &[u8]) -> anyhow::Result { { let mut f = std::fs::File::create(&tmp_path) .map_err(|e| anyhow::anyhow!("create tempfile {:?}: {e}", tmp_path))?; - f.write_all(bytes) - .map_err(|e| anyhow::anyhow!("write tempfile {:?}: {e}", tmp_path))?; - f.sync_all() - .map_err(|e| anyhow::anyhow!("fsync tempfile {:?}: {e}", tmp_path))?; + // Remove the temp file on write/fsync failure so a leaked + // `.tmp_.md` never accumulates in the user-facing vault. + if let Err(e) = f.write_all(bytes).and_then(|_| f.sync_all()) { + drop(f); + let _ = std::fs::remove_file(&tmp_path); + return Err(anyhow::anyhow!("write/fsync tempfile {:?}: {e}", tmp_path)); + } } // Rename: if the target appeared concurrently (another thread/process beat @@ -97,6 +100,100 @@ pub fn write_if_new(abs_path: &Path, bytes: &[u8]) -> anyhow::Result { } } +/// Ensure the file at `abs_path` contains exactly `full_bytes`, whose body +/// (front-matter excluded) hashes to `body_sha256`. +/// +/// This is the write-side half of the content store's integrity contract +/// (#4689): the recorded `content_sha256` must always match the bytes actually +/// on disk. `write_if_new` alone cannot guarantee that — it skips an existing +/// file unconditionally, so a stale/partial pre-existing file at the target +/// path would be left in place while the caller records the *new* body's sha, +/// permanently diverging the DB token from disk. +/// +/// Behaviour when the file already exists: +/// - on-disk body sha matches `body_sha256` → idempotent no-op. +/// - mismatch → atomically replaced (temp file + rename over the destination) +/// from `full_bytes`, so the file is never observed missing or partial. At +/// ingest the freshly-composed input is authoritative, so overwriting a +/// drifted on-disk file is the correct reconciliation. +pub fn write_or_replace_body( + abs_path: &Path, + full_bytes: &[u8], + body_sha256: &str, +) -> anyhow::Result<()> { + if abs_path.exists() { + let disk_sha = read_body_sha256(abs_path).unwrap_or_default(); + if disk_sha == body_sha256 { + log::debug!( + "[content_store::atomic] file already on disk with matching body sha: {}", + abs_path.display() + ); + return Ok(()); + } + log::debug!( + "[content_store::atomic] on-disk body sha mismatch for {} (disk={disk_sha} new={body_sha256}) — re-staging", + abs_path.display() + ); + } + // Write the replacement to a sibling temp file and atomically rename it over + // the destination. rename() replaces an existing file atomically on POSIX and + // NTFS, so — unlike unlink-then-write — the destination is never observed + // missing or partially written even if the process crashes or the write + // fails mid-way (a committed DB row can never end up pointing at an absent + // body file). Post-condition: on success `abs_path` holds exactly + // `full_bytes`, whose body hashes to `body_sha256`. + write_via_temp_rename(abs_path, full_bytes) +} + +/// Write `bytes` to `abs_path` via a sibling tempfile + fsync + atomic rename, +/// **replacing** any existing file. The rename is atomic on POSIX and NTFS, so +/// the destination is never seen missing or half-written. Parent directories are +/// created on demand and the parent dir entry is fsync'd on Unix for durability. +fn write_via_temp_rename(abs_path: &Path, bytes: &[u8]) -> anyhow::Result<()> { + let parent = abs_path.parent().unwrap_or_else(|| Path::new(".")); + std::fs::create_dir_all(parent) + .map_err(|e| anyhow::anyhow!("create_dir_all {:?}: {e}", parent))?; + + let tmp_path = parent.join(format!(".tmp_{}.md", uuid_v4_hex())); + { + let mut f = std::fs::File::create(&tmp_path) + .map_err(|e| anyhow::anyhow!("create tempfile {:?}: {e}", tmp_path))?; + // Clean up the temp file on a write/fsync failure too — not just on the + // rename failure below. content_root is the user-facing vault, so a + // leaked `.tmp_.md` is visible in Obsidian and would accumulate. + if let Err(e) = f.write_all(bytes).and_then(|_| f.sync_all()) { + drop(f); + let _ = std::fs::remove_file(&tmp_path); + return Err(anyhow::anyhow!("write/fsync tempfile {:?}: {e}", tmp_path)); + } + } + + if let Err(e) = std::fs::rename(&tmp_path, abs_path) { + // Keep the old destination intact; only our temp is cleaned up. + let _ = std::fs::remove_file(&tmp_path); + return Err(anyhow::anyhow!( + "rename {:?} -> {:?}: {e}", + tmp_path, + abs_path + )); + } + + #[cfg(unix)] + if let Ok(dir) = std::fs::File::open(parent) { + if let Err(e) = dir.sync_all() { + log::warn!( + "[content_store::atomic] parent dir fsync failed for {:?}: {e}", + parent + ); + } + } + log::debug!( + "[content_store::atomic] wrote (replace) {}", + abs_path.display() + ); + Ok(()) +} + /// A summary that has been written to disk and is ready for SQLite upsert. #[derive(Debug, Clone)] pub struct StagedSummary { @@ -160,38 +257,11 @@ pub fn stage_summary_with_layout( let body_bytes = composed.body.as_bytes(); let sha256 = sha256_hex(body_bytes); - // Idempotent re-stage: if the file already exists, read and hash its - // body bytes. If the on-disk hash matches the new body's hash, return - // the StagedSummary unchanged (true idempotency). If the hashes differ - // the on-disk file is stale/corrupted — re-write it atomically with the - // new content so the db row and disk file are always consistent. - // - // Not re-writing would leave SQLite storing a content_sha256 that - // doesn't match the actual on-disk bytes, breaking integrity checks. - if abs_path.exists() { - let disk_sha = read_body_sha256(&abs_path).unwrap_or_default(); - if disk_sha == sha256 { - log::debug!( - "[content_store::atomic] summary already on disk with matching sha: {}", - input.summary_id - ); - return Ok(StagedSummary { - summary_id: input.summary_id.to_string(), - content_path: rel_path, - content_sha256: sha256, - }); - } - // Hash mismatch — overwrite atomically. - log::debug!( - "[content_store::atomic] summary on-disk sha mismatch for {} — re-staging", - input.summary_id - ); - // Remove the stale file first; write_if_new's fast-path would skip it. - let _ = std::fs::remove_file(&abs_path); - } - - let full_bytes = composed.full.as_bytes(); - write_if_new(&abs_path, full_bytes)?; + // Idempotent re-stage that self-heals a stale/corrupt on-disk file: matching + // body sha is a no-op, a mismatch is atomically rewritten. Not re-writing + // would leave SQLite storing a content_sha256 that doesn't match the actual + // on-disk bytes, breaking integrity checks. See [`write_or_replace_body`]. + write_or_replace_body(&abs_path, composed.full.as_bytes(), &sha256)?; log::debug!( "[content_store::atomic] staged summary {} → {}", @@ -277,6 +347,47 @@ mod tests { assert_eq!(a.len(), 64); // 32 bytes → 64 hex chars } + #[test] + fn write_or_replace_body_writes_new_and_is_idempotent() { + let dir = TempDir::new().unwrap(); + let path = dir.path().join("c.md"); + let full = b"---\nk: v\n---\nBODY"; + let body_sha = sha256_hex(b"BODY"); + // Fresh write. + write_or_replace_body(&path, full, &body_sha).unwrap(); + assert_eq!(std::fs::read(&path).unwrap(), full); + // Idempotent: a matching on-disk body sha leaves the file untouched. + write_or_replace_body(&path, full, &body_sha).unwrap(); + assert_eq!(std::fs::read(&path).unwrap(), full); + } + + #[test] + fn write_or_replace_body_overwrites_stale_file() { + let dir = TempDir::new().unwrap(); + let path = dir.path().join("c.md"); + // A stale pre-existing file with a different body would be skipped by + // write_if_new; write_or_replace_body must reconcile it (#4689). + write_if_new(&path, b"---\nk: v\n---\nOLD").unwrap(); + let new_full = b"---\nk: v\n---\nNEW"; + write_or_replace_body(&path, new_full, &sha256_hex(b"NEW")).unwrap(); + assert_eq!(std::fs::read(&path).unwrap(), new_full); + } + + #[test] + fn write_or_replace_body_errors_when_stale_target_cannot_be_replaced() { + // If the stale target can't be removed/overwritten (here: a directory + // sits at the path), the function must refuse rather than let the caller + // record a content_sha256 the on-disk bytes do not match (#4689). + let dir = TempDir::new().unwrap(); + let path = dir.path().join("c.md"); + std::fs::create_dir_all(&path).unwrap(); + let res = write_or_replace_body(&path, b"---\nk: v\n---\nNEW", &sha256_hex(b"NEW")); + assert!( + res.is_err(), + "must not record a sha the target does not hold" + ); + } + fn mk_summary_input<'a>( tree_kind: SummaryTreeKind, scope: &'a str, diff --git a/src/openhuman/memory_store/content/compose/yaml.rs b/src/openhuman/memory_store/content/compose/yaml.rs index 0a92da93c..05236454f 100644 --- a/src/openhuman/memory_store/content/compose/yaml.rs +++ b/src/openhuman/memory_store/content/compose/yaml.rs @@ -76,7 +76,27 @@ pub fn split_front_matter(content: &str) -> Option<(&str, &str)> { /// We conservatively quote strings containing `:`, `#`, `[`, `]`, `{`, `}`, /// `"`, `'`, `\`, leading/trailing whitespace, or that start with special /// YAML indicator characters. +/// +/// Newlines, carriage returns and tabs are collapsed to a single space first: +/// front-matter scalars are single-line identifiers/paths, and a raw newline in +/// a value would inject a spurious `\n---\n` into the outer front matter, which +/// the reader's `split_front_matter` would then mistake for the closing +/// delimiter — corrupting the body boundary and its sha (#4689). Collapsing is +/// lossless for the identifier fields that flow through here (they never +/// legitimately contain control whitespace). pub fn yaml_scalar(s: &str) -> String { + let sanitized: String = s + .chars() + .map(|c| { + if matches!(c, '\n' | '\r' | '\t') { + ' ' + } else { + c + } + }) + .collect(); + let s = sanitized.as_str(); + let needs_quoting = s.is_empty() || s.trim() != s || s.starts_with(|c: char| { @@ -94,3 +114,36 @@ pub fn yaml_scalar(s: &str) -> String { s.to_string() } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn yaml_scalar_plain_value_is_unquoted() { + assert_eq!(yaml_scalar("hello"), "hello"); + } + + #[test] + fn yaml_scalar_collapses_newlines_and_blocks_fm_injection() { + // A field value carrying an embedded `\n---\n` must not emit a raw + // newline, or it would inject a spurious front-matter closer (#4689). + let out = yaml_scalar("vault:evil\n---\ninjected"); + assert!(!out.contains('\n'), "no raw newline: {out}"); + assert!(!out.contains('\r')); + + // A composed front matter using the sanitized value still splits at the + // real closer, leaving the body intact. + let fm = format!("---\nsource_id: {out}\n---\nBODY"); + let (_, body) = split_front_matter(&fm).expect("front matter splits"); + assert_eq!(body, "BODY"); + } + + #[test] + fn yaml_scalar_collapses_tabs_and_carriage_returns() { + let out = yaml_scalar("a\tb\r\nc"); + assert!(!out.contains('\t')); + assert!(!out.contains('\r')); + assert!(!out.contains('\n')); + } +} diff --git a/src/openhuman/memory_store/content/mod.rs b/src/openhuman/memory_store/content/mod.rs index 302a76787..e9014340d 100644 --- a/src/openhuman/memory_store/content/mod.rs +++ b/src/openhuman/memory_store/content/mod.rs @@ -99,27 +99,20 @@ pub fn stage_chunks(content_root: &Path, chunks: &[Chunk]) -> anyhow::Result { - if written { - log::debug!("[content_store] wrote chunk {} → {}", chunk.id, rel_path); - } else { - log::debug!( - "[content_store] chunk {} already on disk at {}", - chunk.id, - rel_path - ); - } - } - Err(e) => { - log::error!( - "[content_store] failed to write chunk {} to {}: {e}", - chunk.id, - rel_path - ); - return Err(e); - } + // Self-heal a stale/drifted on-disk file so the recorded content_sha256 + // always matches the bytes actually on disk (#4689). A plain write_if_new + // would skip a pre-existing file at this path while we still record the + // freshly-composed sha, permanently diverging the DB token from disk and + // forcing read_chunk_body to serve the ≤500-char preview. + if let Err(e) = atomic::write_or_replace_body(&abs_path, &full_bytes, &sha256) { + log::error!( + "[content_store] failed to write chunk {} to {}: {e}", + chunk.id, + rel_path + ); + return Err(e); } + log::debug!("[content_store] staged chunk {} → {}", chunk.id, rel_path); staged.push(StagedChunk { chunk: chunk.clone(), @@ -194,4 +187,37 @@ mod tests { assert_eq!(first[0].content_sha256, second[0].content_sha256); assert_eq!(first[0].content_path, second[0].content_path); } + + #[test] + fn stage_chunks_overwrites_stale_on_disk_body() { + let dir = TempDir::new().unwrap(); + let chunk = sample_chunk(0); + + // Pre-write a stale file at the chunk's content path with a different + // body, so write_if_new would otherwise skip and leave it in place while + // stage_chunks records the fresh body's sha — the #4689 divergence. + let abs = paths::chunk_abs_path( + dir.path(), + chunk.metadata.source_kind.as_str(), + &chunk.metadata.source_id, + &chunk.id, + ); + std::fs::create_dir_all(abs.parent().unwrap()).unwrap(); + std::fs::write(&abs, b"---\nstale: 1\n---\nSTALE BODY").unwrap(); + + let staged = stage_chunks(dir.path(), std::slice::from_ref(&chunk)).unwrap(); + + // The file now holds the freshly-composed body, and the recorded sha + // matches the bytes actually on disk. + let on_disk = std::fs::read_to_string(&abs).unwrap(); + assert!( + on_disk.ends_with(&chunk.content), + "body rewritten: {on_disk}" + ); + let (_, body) = super::compose::split_front_matter(&on_disk).unwrap(); + assert_eq!( + staged[0].content_sha256, + atomic::sha256_hex(body.as_bytes()) + ); + } } diff --git a/src/openhuman/memory_store/content/read.rs b/src/openhuman/memory_store/content/read.rs index 01c4c225c..608d858df 100644 --- a/src/openhuman/memory_store/content/read.rs +++ b/src/openhuman/memory_store/content/read.rs @@ -256,20 +256,33 @@ pub fn read_chunk_body( ) })?; - // Verify the on-disk body matches the SHA stored at write time. A mismatch - // means the file was tampered with, the tx that committed the pointer - // raced with a separate writer, or the disk corrupted — all unsafe to - // hand back to a consumer. Fail loudly rather than serve stale/corrupt - // bytes into the LLM extractor / summariser pipeline. + // The content file is content-addressed and atomically written, so the file + // on disk is authoritative for this chunk's body. A sha mismatch means the + // stored token drifted from disk — e.g. an external editor rewrote a synced + // file after ingest (#4689). Serve the full on-disk body and repair the + // stale token so the next read verifies cleanly, instead of returning an Err + // that every caller converts into the ≤500-char preview (silent truncation). if result.sha256 != expected_sha256 { - return Err(anyhow::anyhow!( - "[content_store::read] sha256 mismatch for chunk_id={} \ - expected={} actual={} path_hash={}", + log::warn!( + "[content_store::read] stale sha token for chunk_id={} disk={} db={} path_hash={} \ + — serving on-disk body and repairing token", chunk_id, - expected_sha256, result.sha256, + expected_sha256, redact(&rel_path), - )); + ); + if let Err(e) = crate::openhuman::memory_store::chunks::store::update_chunk_content_sha256( + config, + chunk_id, + &result.sha256, + ) { + // Best-effort: the correct body is already in hand; a failed repair + // just means the next read re-heals. Never fail the read on this. + log::warn!( + "[content_store::read] failed to repair sha token for chunk_id={}: {e:#}", + chunk_id, + ); + } } Ok(result.body) @@ -384,17 +397,28 @@ pub fn read_summary_body( ) })?; - // Verify the on-disk body matches the SHA stored at seal time. See the - // matching guard in `read_chunk_body` for rationale. + // Self-heal a drifted sha token by trusting the on-disk file and repairing + // the stored token, rather than returning an Err that callers convert into + // the ≤500-char preview. See the matching guard in `read_chunk_body` (#4689). if result.sha256 != expected_sha256 { - return Err(anyhow::anyhow!( - "[content_store::read] sha256 mismatch for summary_id={} \ - expected={} actual={} path_hash={}", + log::warn!( + "[content_store::read] stale sha token for summary_id={} disk={} db={} path_hash={} \ + — serving on-disk body and repairing token", summary_id, - expected_sha256, result.sha256, + expected_sha256, redact(&rel_path), - )); + ); + if let Err(e) = crate::openhuman::memory_store::chunks::store::update_summary_content_sha256( + config, + summary_id, + &result.sha256, + ) { + log::warn!( + "[content_store::read] failed to repair sha token for summary_id={}: {e:#}", + summary_id, + ); + } } Ok(result.body) @@ -748,7 +772,7 @@ mod tests { } #[test] - fn read_chunk_body_errors_on_sha_mismatch() { + fn read_chunk_body_self_heals_on_sha_mismatch() { let tmp = TempDir::new().unwrap(); let cfg = test_config(&tmp); let chunk = sample_chunk(); @@ -766,6 +790,8 @@ mod tests { }) .unwrap(); + // Simulate an external editor rewriting the synced file after ingest: + // the on-disk body drifts from the recorded content_sha256 (#4689). let rel = crate::openhuman::memory_store::chunks::store::get_chunk_content_path(&cfg, &chunk.id) .unwrap() @@ -776,8 +802,74 @@ mod tests { } std::fs::write(&abs, b"---\nsource_kind: chat\n---\nmutated body").unwrap(); - let err = read_chunk_body(&cfg, &chunk.id).unwrap_err(); - assert!(err.to_string().contains("sha256 mismatch")); + // Self-heal: serve the full on-disk body instead of erroring into the + // ≤500-char preview, and repair the stale token. + let body = read_chunk_body(&cfg, &chunk.id).unwrap(); + assert_eq!(body, "mutated body"); + + let (_, sha) = crate::openhuman::memory_store::chunks::store::get_chunk_content_pointers( + &cfg, &chunk.id, + ) + .unwrap() + .unwrap(); + assert_eq!(sha, sha256_hex(b"mutated body")); + // A second read now verifies cleanly against the repaired token. + assert_eq!(read_chunk_body(&cfg, &chunk.id).unwrap(), "mutated body"); + } + + #[test] + fn read_summary_body_self_heals_on_sha_mismatch() { + let tmp = TempDir::new().unwrap(); + let cfg = test_config(&tmp); + let tree = sample_tree(); + let node = sample_summary_node(); + insert_tree(&cfg, &tree).unwrap(); + let staged = stage_summary( + &cfg.memory_tree_content_root(), + &SummaryComposeInput { + summary_id: &node.id, + tree_kind: SummaryTreeKind::Source, + tree_id: &tree.id, + tree_scope: &tree.scope, + level: node.level, + child_ids: &node.child_ids, + child_basenames: None, + child_count: node.child_ids.len(), + time_range_start: node.time_range_start, + time_range_end: node.time_range_end, + sealed_at: node.sealed_at, + body: &node.content, + }, + "slack-eng", + ) + .unwrap(); + with_connection(&cfg, |conn| { + let tx = conn.unchecked_transaction()?; + insert_summary_tx(&tx, &node, Some(&staged), "test")?; + tx.commit()?; + Ok(()) + }) + .unwrap(); + + let (rel, _) = crate::openhuman::memory_store::chunks::store::get_summary_content_pointers( + &cfg, &node.id, + ) + .unwrap() + .unwrap(); + let mut abs = cfg.memory_tree_content_root(); + for part in rel.split('/') { + abs.push(part); + } + std::fs::write(&abs, b"---\ntree_kind: source\n---\nmutated summary").unwrap(); + + let body = read_summary_body(&cfg, &node.id).unwrap(); + assert_eq!(body, "mutated summary"); + let (_, sha) = crate::openhuman::memory_store::chunks::store::get_summary_content_pointers( + &cfg, &node.id, + ) + .unwrap() + .unwrap(); + assert_eq!(sha, sha256_hex(b"mutated summary")); } #[test] diff --git a/src/openhuman/memory_store/content/tags.rs b/src/openhuman/memory_store/content/tags.rs index af342904f..f18d21705 100644 --- a/src/openhuman/memory_store/content/tags.rs +++ b/src/openhuman/memory_store/content/tags.rs @@ -48,6 +48,13 @@ pub fn update_chunk_tags(abs_path: &Path, tags: &[String]) -> anyhow::Result<()> let new_bytes = rewrite_tags(&old_bytes, &augmented) .map_err(|e| anyhow::anyhow!("rewrite_tags {:?}: {e}", abs_path))?; + // The tag rewrite must only ever touch front-matter. Verify the body is + // byte-identical before committing so a front-matter parse regression (or a + // newline-injected field) fails loud here instead of silently drifting the + // on-disk body from the DB content_sha256 and truncating retrieval (#4689). + // Mirrors the post-rewrite guard in `update_summary_tags`. + ensure_tag_rewrite_preserves_body(&old_bytes, &new_bytes, abs_path)?; + // Write the new content atomically via a sibling temp file. let parent = abs_path.parent().unwrap_or_else(|| Path::new(".")); let tmp_name = format!(".tmp_tags_{}.md", crate_temp_id()); @@ -191,6 +198,33 @@ pub fn update_summary_tags(config: &Config, summary_id: &str) -> anyhow::Result< Ok(()) } +/// Guard for [`update_chunk_tags`]: the body (front-matter excluded) must be +/// byte-identical before and after a tag rewrite. A drift here would desync the +/// on-disk body from the DB `content_sha256` and silently truncate retrieval +/// (#4689), so surface it as a loud error rather than committing the rewrite. +fn ensure_tag_rewrite_preserves_body( + old_bytes: &[u8], + new_bytes: &[u8], + abs_path: &Path, +) -> anyhow::Result<()> { + let body = |bytes: &[u8]| -> Option { + std::str::from_utf8(bytes) + .ok() + .and_then(split_front_matter) + .map(|(_, body)| body.to_string()) + }; + // Require BOTH sides to parse AND match. Comparing `Option`s directly would + // let two un-parseable sides (`None == None`) pass — the exact silent-drift + // case this guard exists to catch — so treat an unparseable body as a failure. + match (body(old_bytes), body(new_bytes)) { + (Some(a), Some(b)) if a == b => Ok(()), + _ => Err(anyhow::anyhow!( + "[content_store::tags] update_chunk_tags would mutate or invalidate the body for {:?} — aborting rewrite", + abs_path + )), + } +} + /// Slugify an entity kind string for use in an Obsidian hierarchical tag. /// /// Output: lowercase, spaces and non-alphanumeric chars replaced with `-`, @@ -424,6 +458,31 @@ mod tests { assert!(update_chunk_tags(&path, &["p/X".into()]).is_ok()); } + #[test] + fn ensure_tag_rewrite_preserves_body_accepts_equal_bodies() { + let p = std::path::Path::new("x.md"); + // Same body, different front-matter → allowed. + let old = b"---\nk: v\n---\nBODY"; + let new = b"---\nk: other\ntags:\n - t\n---\nBODY"; + assert!(ensure_tag_rewrite_preserves_body(old, new, p).is_ok()); + } + + #[test] + fn ensure_tag_rewrite_preserves_body_rejects_body_drift() { + let p = std::path::Path::new("x.md"); + let old = b"---\nk: v\n---\nBODY"; + let drifted = b"---\nk: v\n---\nDIFFERENT BODY"; + assert!(ensure_tag_rewrite_preserves_body(old, drifted, p).is_err()); + } + + #[test] + fn ensure_tag_rewrite_preserves_body_rejects_unparseable_bodies() { + // Both sides lack front-matter → both parse to None. The guard must still + // fail rather than let `None == None` pass silently. + let p = std::path::Path::new("x.md"); + assert!(ensure_tag_rewrite_preserves_body(b"no front matter", b"still none", p).is_err()); + } + #[test] fn slugify_tag_kind_examples() { assert_eq!(slugify_tag_kind("Person"), "person");