From 82054f3957a88190e4ddd816ee77ff6e5929ce26 Mon Sep 17 00:00:00 2001 From: Ghost Scripter Date: Fri, 9 Oct 2026 12:14:06 +0530 Subject: [PATCH] Tag synced chunk sources as external on v1 import v1's chunk tier kept no taint, so the importer tagged only memory_docs rows with taint:external_sync. But the chunk store holds most synced content: a Gmail sync files every message there under the owner gmail-sync:. In v2 those chunks share one store with the user's own memory, and arrived indistinguishable from it. Decide from each chunk's owner instead: a source is external unless every chunk's owner is one the host writes itself (cron, cron:, or the archivist's JSON session key with a thread_id). Anything else, including a blank owner, fails closed as the memory_docs decode does. A store without the owner column gets no tag. --- .../src/import/README.md | 14 ++- .../tinymemory-integrations/src/import/mod.rs | 3 +- .../src/import/sections/chunks.rs | 44 +++++++- .../src/import/sections/memory_docs.rs | 5 +- .../src/import/workspace/schema.rs | 3 + .../tests/legacy_import.rs | 102 +++++++++++++++++- .../tests/support/mod.rs | 23 +++- 7 files changed, 183 insertions(+), 11 deletions(-) diff --git a/crates/tinymemory-integrations/src/import/README.md b/crates/tinymemory-integrations/src/import/README.md index a6be723b..3547a82d 100644 --- a/crates/tinymemory-integrations/src/import/README.md +++ b/crates/tinymemory-integrations/src/import/README.md @@ -119,8 +119,18 @@ row can land in (documents, learnings, `global`). The tag rides in the item's metadata, which a CortexDB engine stores whole, so the host can read it back on recall. The decode fails closed like v1's: an unknown or empty value is external. A store from before the `taint` column is read as all `internal`, -which is how v1 read it. `episodic_log`, `user_profile` and the chunk store -have no taint in v1 and get no tag. +which is how v1 read it. `episodic_log` and `user_profile` have no taint in v1 +and get no tag. + +The chunk store has no taint column either (v1's chunk tier refused +`ExternalSync`), yet it holds most synced content: a Gmail sync files every +message there under the owner `gmail-sync:`. In v2 those chunks +sit beside the user's own memory, so their `owner` decides instead: a chunk +source is tagged `taint:external_sync` unless every chunk's owner is one the +host writes itself, `cron` (or `cron:`) or the archivist's session key (a +JSON object with a `thread_id`). Connector owners, an agent's own label and a +blank owner are all external, failing closed. A chunk store without the +`owner` column gets no tag. ### Documents diff --git a/crates/tinymemory-integrations/src/import/mod.rs b/crates/tinymemory-integrations/src/import/mod.rs index 5c4338aa..7bcf0c2c 100644 --- a/crates/tinymemory-integrations/src/import/mod.rs +++ b/crates/tinymemory-integrations/src/import/mod.rs @@ -27,7 +27,8 @@ //! `episodic_log:lesson:`, `graph_global:`, `graph_namespace:`, //! `file:`), and //! `meta.workspace` is the workspace path. A `memory_docs` row v1 marked as -//! synced from an external service also carries [`EXTERNAL_SYNC_TAG`]. The +//! synced from an external service, and a chunk source whose owner is not +//! the host's own, also carries [`EXTERNAL_SYNC_TAG`]. The //! module's `README.md` details every mapping decision. //! //! Import is resumable: each [`ImportedItem`] carries the [`Checkpoint`] to diff --git a/crates/tinymemory-integrations/src/import/sections/chunks.rs b/crates/tinymemory-integrations/src/import/sections/chunks.rs index 07bd3ec2..11da9305 100644 --- a/crates/tinymemory-integrations/src/import/sections/chunks.rs +++ b/crates/tinymemory-integrations/src/import/sections/chunks.rs @@ -10,6 +10,15 @@ //! people rather than the assistant. Every other source kind (`document`, //! `email`) becomes one document whose body is its chunks joined by blank //! lines. +//! +//! v1 kept no taint on chunks (its chunk tier refused `ExternalSync`), but a +//! chunk's `owner` says where it came from. A source is tagged +//! [`EXTERNAL_SYNC_TAG`] unless every chunk's owner is one the host itself +//! writes: `cron` (or `cron:`) and the archivist's session key, a JSON +//! object with a `thread_id`. Anything else, a connector sync such as +//! `gmail-sync:`, an agent's own label, or a blank owner, is +//! external, failing closed as the `memory_docs` decode does. A store without +//! the `owner` column has nothing to read and gets no tag. use std::io::ErrorKind; use std::path::{Component, Path}; @@ -17,7 +26,7 @@ use std::path::{Component, Path}; use rusqlite::params; use tinymemory_api::{DocumentBody, Role, StoreItem, Turn, TurnRange}; -use super::{Mark, Scanned, import_meta, push_unique, sql_limit}; +use super::{EXTERNAL_SYNC_TAG, Mark, Scanned, import_meta, push_unique, sql_limit}; use crate::import::checkpoint::ChunkCursor; use crate::import::convert; use crate::import::error::{Error, Result}; @@ -29,6 +38,8 @@ struct Chunk { text: String, timestamp_ms: i64, tags_json: String, + /// `owner`, or `None` when the store has no such column. + owner: Option, } /// The next page of sources after `after`; empty when the workspace has no @@ -117,6 +128,13 @@ fn source_item( } } push_unique(&mut tags, format!("source_kind:{}", source.source_kind)); + if store.owner + && !chunks + .iter() + .all(|chunk| is_host_owner(chunk.owner.as_deref())) + { + push_unique(&mut tags, EXTERNAL_SYNC_TAG.to_string()); + } meta.tags = tags; meta.observed_at = chunks .iter() @@ -158,8 +176,9 @@ fn chunks(store: &ChunkStore, source: &ChunkCursor) -> Result> { } else { "NULL" }; + let owner = if store.owner { "owner" } else { "NULL" }; let sql = format!( - "SELECT content, {content_path}, timestamp_ms, tags_json FROM mem_tree_chunks \ + "SELECT content, {content_path}, timestamp_ms, tags_json, {owner} FROM mem_tree_chunks \ WHERE source_kind = ?1 AND source_id = ?2 ORDER BY seq_in_source, id" ); let mut stmt = store.conn.prepare(&sql)?; @@ -169,11 +188,12 @@ fn chunks(store: &ChunkStore, source: &ChunkCursor) -> Result> { row.get::<_, Option>(1)?, row.get::<_, Option>(2)?.unwrap_or_default(), row.get::<_, Option>(3)?.unwrap_or_default(), + row.get::<_, Option>(4)?, )) })?; let mut chunks = Vec::new(); for row in rows { - let (preview, path, timestamp_ms, tags_json) = row?; + let (preview, path, timestamp_ms, tags_json, owner) = row?; let full = match path.as_deref() { Some(path) => full_body(&store.content_dir, path)?, None => None, @@ -186,11 +206,29 @@ fn chunks(store: &ChunkStore, source: &ChunkCursor) -> Result> { text, timestamp_ms, tags_json, + owner, }); } Ok(chunks) } +/// Whether `owner` is one the host writes for its own content: `cron` (or +/// `cron:`), or the archivist's session key, a JSON object carrying a +/// `thread_id`. A missing or blank owner is not. +fn is_host_owner(owner: Option<&str>) -> bool { + let Some(owner) = owner.map(str::trim).filter(|owner| !owner.is_empty()) else { + return false; + }; + if owner == "cron" || owner.starts_with("cron:") { + return true; + } + serde_json::from_str::(owner).is_ok_and(|value| { + value + .get("thread_id") + .is_some_and(serde_json::Value::is_string) + }) +} + /// Reads `content_dir/`; `None` when the path is not a plain /// relative path inside the content directory or the file does not exist. fn full_body(content_dir: &Path, relative: &str) -> Result> { diff --git a/crates/tinymemory-integrations/src/import/sections/memory_docs.rs b/crates/tinymemory-integrations/src/import/sections/memory_docs.rs index 5315e031..6b485cde 100644 --- a/crates/tinymemory-integrations/src/import/sections/memory_docs.rs +++ b/crates/tinymemory-integrations/src/import/sections/memory_docs.rs @@ -44,8 +44,9 @@ const SECTION_PREFIXES: [&str; 9] = [ ]; /// The tag on every item from a `memory_docs` row v1 marked as synced from an -/// external service (Gmail, Slack, Notion, Composio, MCP, ...): content the -/// user did not write, which v1 kept out of external-effect tool decisions. +/// external service (Gmail, Slack, Notion, Composio, MCP, ...), and on every +/// chunk source whose owner is not the host's own: content the user did not +/// write, which v1 kept out of external-effect tool decisions. pub const EXTERNAL_SYNC_TAG: &str = "taint:external_sync"; /// One `memory_docs` row. diff --git a/crates/tinymemory-integrations/src/import/workspace/schema.rs b/crates/tinymemory-integrations/src/import/workspace/schema.rs index 020772bd..508cb3a0 100644 --- a/crates/tinymemory-integrations/src/import/workspace/schema.rs +++ b/crates/tinymemory-integrations/src/import/workspace/schema.rs @@ -150,6 +150,8 @@ pub(crate) struct ChunkStore { pub(crate) content_dir: PathBuf, /// Whether `mem_tree_chunks.content_path` exists. pub(crate) content_path: bool, + /// Whether `mem_tree_chunks.owner` exists. + pub(crate) owner: bool, } impl ChunkStore { @@ -172,6 +174,7 @@ impl ChunkStore { } Ok(Some(Self { content_path: present.contains("content_path"), + owner: present.contains("owner"), content_dir: tree.join("content"), conn, })) diff --git a/crates/tinymemory-integrations/tests/legacy_import.rs b/crates/tinymemory-integrations/tests/legacy_import.rs index a1dbe6a3..7036af0a 100644 --- a/crates/tinymemory-integrations/tests/legacy_import.rs +++ b/crates/tinymemory-integrations/tests/legacy_import.rs @@ -8,7 +8,7 @@ mod support; use std::path::Path; -use support::{OLD_MEMORY_DDL, chunk, chunk_store, doc, facet, turn, workspace}; +use support::{OLD_MEMORY_DDL, chunk, chunk_store, doc, facet, owned_chunk, turn, workspace}; use tinymemory_api::{ DocumentBody, LearningKind, Role, SourceKind, StoreItem, ToolCallRef, TurnRange, }; @@ -866,6 +866,106 @@ fn a_store_without_the_taint_column_reads_as_internal() { })); } +/// Whether the item imported from `legacy_id` carries [`EXTERNAL_SYNC_TAG`]. +fn is_external(items: &[ImportedItem], legacy_id: &str) -> bool { + find(items, legacy_id) + .meta() + .tags + .iter() + .any(|t| t == EXTERNAL_SYNC_TAG) +} + +#[test] +fn chunk_sources_a_connector_synced_are_tagged_external() { + let dir = tempfile::tempdir().unwrap(); + let chunks = chunk_store(dir.path()); + owned_chunk( + &chunks, + "k1", + "email", + "gmail:me|them", + 0, + "gmail-sync:ca_1", + ); + owned_chunk(&chunks, "k2", "document", "gmail:msg", 0, "gmail-sync:ca_1"); + owned_chunk(&chunks, "k3", "chat", "slack:c1", 0, "slack:conn"); + owned_chunk(&chunks, "k4", "chat", "cron-out", 0, "cron"); + owned_chunk(&chunks, "k5", "chat", "cron-job", 0, "cron:job-7"); + owned_chunk( + &chunks, + "k6", + "chat", + "conversations:agent", + 0, + r#"{"client_id":"c","thread_id":"thread-1"}"#, + ); + drop(chunks); + let ws = LegacyWorkspace::open(dir.path()).unwrap(); + let items = all(&ws); + assert!(is_external(&items, "mem_tree_chunks:email:gmail:me|them")); + assert!(is_external(&items, "mem_tree_chunks:document:gmail:msg")); + assert!(is_external(&items, "mem_tree_chunks:chat:slack:c1")); + assert!(!is_external(&items, "mem_tree_chunks:chat:cron-out")); + assert!(!is_external(&items, "mem_tree_chunks:chat:cron-job")); + assert!(!is_external( + &items, + "mem_tree_chunks:chat:conversations:agent" + )); + let StoreItem::Document { meta, .. } = find(&items, "mem_tree_chunks:email:gmail:me|them") + else { + panic!("email is a document"); + }; + assert_eq!(meta.tags, ["source_kind:email", EXTERNAL_SYNC_TAG]); +} + +#[test] +fn a_chunk_owner_the_host_does_not_write_fails_closed() { + let dir = tempfile::tempdir().unwrap(); + let chunks = chunk_store(dir.path()); + // One synced chunk taints the whole source. + owned_chunk(&chunks, "k1", "chat", "mixed", 0, "cron"); + owned_chunk(&chunks, "k2", "chat", "mixed", 1, "slack:conn"); + owned_chunk(&chunks, "k3", "document", "blank", 0, " "); + owned_chunk(&chunks, "k4", "document", "label", 0, "agent-notes"); + owned_chunk(&chunks, "k5", "document", "crony", 0, "cronjob"); + owned_chunk(&chunks, "k6", "document", "json", 0, r#"{"client_id":"c"}"#); + drop(chunks); + let ws = LegacyWorkspace::open(dir.path()).unwrap(); + let items = all(&ws); + for id in [ + "chat:mixed", + "document:blank", + "document:label", + "document:crony", + "document:json", + ] { + assert!( + is_external(&items, &format!("mem_tree_chunks:{id}")), + "{id}" + ); + } +} + +#[test] +fn a_chunk_store_without_the_owner_column_gets_no_taint() { + let dir = tempfile::tempdir().unwrap(); + std::fs::create_dir_all(dir.path().join("memory_tree")).unwrap(); + let chunks = rusqlite::Connection::open(dir.path().join("memory_tree/chunks.db")).unwrap(); + chunks + .execute_batch( + "CREATE TABLE mem_tree_chunks (id TEXT PRIMARY KEY, source_kind TEXT NOT NULL, + source_id TEXT NOT NULL, timestamp_ms INTEGER NOT NULL, + tags_json TEXT NOT NULL DEFAULT '[]', content TEXT NOT NULL, + seq_in_source INTEGER NOT NULL); + INSERT INTO mem_tree_chunks VALUES ('k1', 'email', 'e1', 1000, '[]', 'hello', 0);", + ) + .unwrap(); + drop(chunks); + let ws = LegacyWorkspace::open(dir.path()).unwrap(); + let items = all(&ws); + assert!(!is_external(&items, "mem_tree_chunks:email:e1")); +} + #[test] fn imports_an_early_v1_store_without_optional_columns() { let (dir, conn) = workspace(OLD_MEMORY_DDL); diff --git a/crates/tinymemory-integrations/tests/support/mod.rs b/crates/tinymemory-integrations/tests/support/mod.rs index 2a272fbe..0d86a585 100644 --- a/crates/tinymemory-integrations/tests/support/mod.rs +++ b/crates/tinymemory-integrations/tests/support/mod.rs @@ -158,7 +158,7 @@ pub(crate) fn chunk_store(root: &Path) -> Connection { conn } -/// Inserts a chunk. +/// Inserts a chunk owned by the host's own `cron`. #[allow(clippy::too_many_arguments, reason = "mirrors the table's columns")] pub(crate) fn chunk( conn: &Connection, @@ -175,7 +175,7 @@ pub(crate) fn chunk( "INSERT INTO mem_tree_chunks (id, source_kind, source_id, owner, timestamp_ms, time_range_start_ms, time_range_end_ms, tags_json, content, token_count, seq_in_source, created_at_ms, content_path) - VALUES (?1, ?2, ?3, 'me', ?4, ?4, ?4, ?5, ?6, 1, ?7, ?4, ?8)", + VALUES (?1, ?2, ?3, 'cron', ?4, ?4, ?4, ?5, ?6, 1, ?7, ?4, ?8)", params![ id, kind, @@ -189,3 +189,22 @@ pub(crate) fn chunk( ) .expect("insert chunk"); } + +/// Inserts a one-line chunk at `seq` of `source`, owned by `owner`. +pub(crate) fn owned_chunk( + conn: &Connection, + id: &str, + kind: &str, + source: &str, + seq: i64, + owner: &str, +) { + conn.execute( + "INSERT INTO mem_tree_chunks (id, source_kind, source_id, owner, timestamp_ms, + time_range_start_ms, time_range_end_ms, content, token_count, seq_in_source, + created_at_ms) + VALUES (?1, ?2, ?3, ?4, 1000, 1000, 1000, ?1, 1, ?5, 1000)", + params![id, kind, source, owner, seq], + ) + .expect("insert owned chunk"); +}