From cc9e15b16f15cef09f3f0d39cf250de7ee42ee4e Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 01:52:14 +0530 Subject: [PATCH 1/2] refactor(import): split workspace import into per-section modules The workspace import logic was reorganised into a sections module with separate files for graph, memory docs, profile, and connector handling, plus a dedicated schema module. This keeps each section's parsing and validation self-contained and makes the connector tests easier to locate. Auto-committed-on: macbook Co-authored-by: Medulla --- .../src/import/sections/connector.rs | 61 +++++++++++++++++++ .../src/import/sections/connector_tests.rs | 31 ++++++++++ .../src/import/sections/graph.rs | 20 ++++-- .../src/import/sections/memory_docs.rs | 15 ++++- .../src/import/sections/mod.rs | 1 + .../src/import/sections/profile.rs | 4 ++ .../src/import/workspace/mod.rs | 17 ++++++ .../src/import/workspace/schema.rs | 3 + 8 files changed, 147 insertions(+), 5 deletions(-) create mode 100644 crates/tinymemory-integrations/src/import/sections/connector.rs create mode 100644 crates/tinymemory-integrations/src/import/sections/connector_tests.rs diff --git a/crates/tinymemory-integrations/src/import/sections/connector.rs b/crates/tinymemory-integrations/src/import/sections/connector.rs new file mode 100644 index 00000000..01e1d4e5 --- /dev/null +++ b/crates/tinymemory-integrations/src/import/sections/connector.rs @@ -0,0 +1,61 @@ +//! What counts as a connector sync, for `LegacyWorkspace::skip_connector_syncs`. +//! +//! v1 synced outside services (Gmail, Slack, Notion, Linear, GitHub, ClickUp, +//! ... through Composio and the connector path) into its stores. That data is +//! re-synced by the connectors themselves, so an opt-in import leaves it out. +//! Every predicate here is used by both the scan and `counts()`, so the counts +//! stay exactly what `items()` yields. +//! +//! - `memory_docs`: logical namespace `skill-*` (Composio `SkillDoc` sync, +//! `skill-{toolkit}`) or `source:*` (connector path, +//! `source:{toolkit}:{conn}`). Taint is deliberately not consulted: v1 also +//! marked the agent's own `global` and flow notes `external_sync`. +//! - chunks: see [`chunk_source_by_identity`] and [`OWNER_SYNC_PATTERN`]. +//! - `user_profile`: `facet_id` starting `skill-` ([`PROFILE_SKILL_PREFIX`]). +//! - `graph_namespace`: namespace starting `skill-`, `source:` or `source_`. + +/// Connector toolkits whose chunk `source_id` is `{toolkit}:{conn}[:{item}]` +/// (Composio before July: slack chat `slack:{conn}`, docs +/// `{toolkit}:{conn}:{id}`; current: `{toolkit}:{conn}:{item}`). +pub(crate) const CONNECTOR_TOOLKIT_PREFIXES: [&str; 6] = [ + "gmail:", "slack:", "notion:", "linear:", "github:", "clickup:", +]; + +/// The `source_kind` v1 gave Gmail threads. +pub(crate) const EMAIL_SOURCE_KIND: &str = "email"; + +/// SQL `LIKE` pattern for a chunk `owner` written by a connector +/// (`{toolkit}-sync:{conn}`). +pub(crate) const OWNER_SYNC_PATTERN: &str = "%-sync:%"; + +/// Prefix of a `user_profile.facet_id` written by Composio identity sync +/// (`skill-{toolkit}-{conn}-{kind}`). +pub(crate) const PROFILE_SKILL_PREFIX: &str = "skill-"; + +/// Whether a `memory_docs` logical namespace is a connector sync. +pub(crate) fn is_connector_namespace(logical: &str) -> bool { + logical.starts_with("skill-") || logical.starts_with("source:") +} + +/// Whether a chunk source is a connector sync judging by its kind and id +/// alone (the `owner` column is checked by the chunk reader). +pub(crate) fn chunk_source_by_identity(source_kind: &str, source_id: &str) -> bool { + source_kind == EMAIL_SOURCE_KIND + || CONNECTOR_TOOLKIT_PREFIXES + .iter() + .any(|prefix| source_id.starts_with(prefix)) +} + +/// SQL condition true for a `graph_namespace` row that is NOT a connector +/// sync (`skill-*`, `source:*`, or the sanitised `source_*`). +pub(crate) const GRAPH_NAMESPACE_KEPT: &str = "NOT (substr(COALESCE(namespace, ''), 1, 6) = 'skill-' \ + OR substr(COALESCE(namespace, ''), 1, 7) = 'source:' \ + OR substr(COALESCE(namespace, ''), 1, 7) = 'source_')"; + +/// SQL condition true for a `user_profile` row that is NOT a Composio +/// identity facet. +pub(crate) const PROFILE_KEPT: &str = "substr(facet_id, 1, 6) != 'skill-'"; + +#[cfg(test)] +#[path = "connector_tests.rs"] +mod tests; diff --git a/crates/tinymemory-integrations/src/import/sections/connector_tests.rs b/crates/tinymemory-integrations/src/import/sections/connector_tests.rs new file mode 100644 index 00000000..0695270e --- /dev/null +++ b/crates/tinymemory-integrations/src/import/sections/connector_tests.rs @@ -0,0 +1,31 @@ +//! Tests for the connector-sync predicates. + +use super::*; + +#[test] +fn namespaces() { + assert!(is_connector_namespace("skill-gmail")); + assert!(is_connector_namespace("source:gmail:conn1")); + assert!(!is_connector_namespace("source_gmail")); + assert!(!is_connector_namespace("global")); + assert!(!is_connector_namespace("document:notes")); + assert!(!is_connector_namespace("skills")); +} + +#[test] +fn chunk_identity() { + assert!(chunk_source_by_identity("email", "anything")); + assert!(chunk_source_by_identity("chat", "slack:conn1")); + assert!(chunk_source_by_identity("document", "notion:c:page")); + for toolkit in CONNECTOR_TOOLKIT_PREFIXES { + assert!(chunk_source_by_identity("document", &format!("{toolkit}x"))); + } + assert!(!chunk_source_by_identity("document", "mem_src:folder")); + assert!(!chunk_source_by_identity("chat", "conversations:agent")); + assert!(!chunk_source_by_identity("document", "slackish:x")); +} + +#[test] +fn profile_prefix_matches_constant() { + assert!(PROFILE_KEPT.contains(PROFILE_SKILL_PREFIX)); +} diff --git a/crates/tinymemory-integrations/src/import/sections/graph.rs b/crates/tinymemory-integrations/src/import/sections/graph.rs index 43d88554..e4cef04c 100644 --- a/crates/tinymemory-integrations/src/import/sections/graph.rs +++ b/crates/tinymemory-integrations/src/import/sections/graph.rs @@ -49,6 +49,16 @@ impl Table { } } +/// The `AND …` filter dropping connector-synced namespaces, shared by `page` +/// and `count`. `graph_global` has no namespace and is never filtered. +fn connector_filter(ws: &LegacyWorkspace, table: Table) -> String { + if ws.skip_connector_syncs && table == Table::Namespace { + format!(" AND {}", super::connector::GRAPH_NAMESPACE_KEPT) + } else { + String::new() + } +} + /// The next page of `table`'s relations after the row `after`. pub(super) fn page( ws: &LegacyWorkspace, @@ -65,10 +75,11 @@ pub(super) fn page( }; let sql = format!( "SELECT rowid, subject, predicate, object, updated_at, {namespace} FROM {} \ - WHERE (?1 IS NULL OR rowid > ?1) AND {} AND {} ORDER BY rowid LIMIT ?2", + WHERE (?1 IS NULL OR rowid > ?1) AND {} AND {}{} ORDER BY rowid LIMIT ?2", table.name(), has_text("subject"), - has_text("object") + has_text("object"), + connector_filter(ws, table) ); let mut stmt = memory.prepare(&sql)?; let rows = stmt.query_map(params![after, sql_limit(limit)], |row| { @@ -118,10 +129,11 @@ pub(super) fn count(ws: &LegacyWorkspace, table: Table) -> Result { return Ok(0); }; let sql = format!( - "SELECT COUNT(*) FROM {} WHERE {} AND {}", + "SELECT COUNT(*) FROM {} WHERE {} AND {}{}", table.name(), has_text("subject"), - has_text("object") + has_text("object"), + connector_filter(ws, table) ); Ok(count_of(memory.query_row(&sql, [], |row| row.get(0))?)) } diff --git a/crates/tinymemory-integrations/src/import/sections/memory_docs.rs b/crates/tinymemory-integrations/src/import/sections/memory_docs.rs index 5315e031..9ded5c43 100644 --- a/crates/tinymemory-integrations/src/import/sections/memory_docs.rs +++ b/crates/tinymemory-integrations/src/import/sections/memory_docs.rs @@ -83,6 +83,7 @@ pub(super) fn documents( .map(|row| { let logical = logical_namespace(&row); let item = match classify(&logical) { + RowClass::Document if skipped(ws, &logical) => None, RowClass::Document => document(ws, &row, logical), _ => None, }; @@ -144,7 +145,9 @@ pub(super) fn count(ws: &LegacyWorkspace, learnings: bool) -> Result { let mut total = 0; for group in groups { let (namespace, logical_namespace, rows) = group?; - let wanted = match classify(&resolve_logical(&namespace, logical_namespace.as_deref())) { + let logical = resolve_logical(&namespace, logical_namespace.as_deref()); + let wanted = match classify(&logical) { + RowClass::Document if skipped(ws, &logical) => false, RowClass::Document => !learnings, RowClass::Learning(_) | RowClass::Global => learnings, RowClass::Event => false, @@ -206,6 +209,16 @@ fn mark_taint(row: &DocRow, tags: &mut Vec) { } } +/// Whether the workspace leaves out a row of this logical namespace as a +/// connector sync; the one test the documents scan and [`count`] share. +fn skipped(ws: &LegacyWorkspace, logical: &str) -> bool { + let skip = ws.skip_connector_syncs && super::connector::is_connector_namespace(logical); + if skip { + tracing_free_note(); + } + skip +} + fn logical_namespace(row: &DocRow) -> String { resolve_logical(&row.namespace, row.logical_namespace.as_deref()) } diff --git a/crates/tinymemory-integrations/src/import/sections/mod.rs b/crates/tinymemory-integrations/src/import/sections/mod.rs index 677bc713..fa1141cf 100644 --- a/crates/tinymemory-integrations/src/import/sections/mod.rs +++ b/crates/tinymemory-integrations/src/import/sections/mod.rs @@ -8,6 +8,7 @@ //! rows never stalls the iterator. mod chunks; +pub(crate) mod connector; mod episodic; mod events; mod files; diff --git a/crates/tinymemory-integrations/src/import/sections/profile.rs b/crates/tinymemory-integrations/src/import/sections/profile.rs index d1831e87..b0a44036 100644 --- a/crates/tinymemory-integrations/src/import/sections/profile.rs +++ b/crates/tinymemory-integrations/src/import/sections/profile.rs @@ -43,6 +43,10 @@ pub(super) fn count(ws: &LegacyWorkspace) -> Result { /// columns this store has. fn live_filters(ws: &LegacyWorkspace) -> String { let mut filters = String::new(); + if ws.skip_connector_syncs { + filters.push_str(" AND "); + filters.push_str(super::connector::PROFILE_KEPT); + } if ws.schema.profile_state { filters.push_str(" AND state IS NOT 'dropped'"); } diff --git a/crates/tinymemory-integrations/src/import/workspace/mod.rs b/crates/tinymemory-integrations/src/import/workspace/mod.rs index a17da881..b400848f 100644 --- a/crates/tinymemory-integrations/src/import/workspace/mod.rs +++ b/crates/tinymemory-integrations/src/import/workspace/mod.rs @@ -49,6 +49,9 @@ pub struct LegacyWorkspace { pub(crate) schema: MemorySchema, /// `memory_tree/chunks.db`, when present and usable. pub(crate) chunks: Option, + /// Whether rows synced from connectors are left out; see + /// [`Self::skip_connector_syncs`]. + pub(crate) skip_connector_syncs: bool, } impl LegacyWorkspace { @@ -146,9 +149,23 @@ impl LegacyWorkspace { memory, schema, chunks, + skip_connector_syncs: false, }) } + /// Leaves out everything v1 synced from connectors (Composio and the + /// connector path: Gmail, Slack, Notion, Linear, GitHub, ClickUp, ...): + /// `memory_docs` in `skill-*` / `source:*` namespaces, chunk sources of + /// kind `email`, with a connector toolkit prefix in `source_id`, or with a + /// `*-sync:*` owner, `skill-*` profile facets, and `skill-*` / `source:*` + /// / `source_*` graph namespaces. [`Self::counts`] excludes the same rows. + /// Off by default. See the module README for the exact rules. + #[must_use] + pub fn skip_connector_syncs(mut self, skip: bool) -> Self { + self.skip_connector_syncs = skip; + self + } + /// The suffixes of the per-profile v1 stores beside the main one in the /// workspace at `path` (`-1`, `-2`, …), sorted: every `memory` or /// `memory_tree` directory with a valid suffix. Each may still 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, })) From 439cce2a7cfa42916852b24978a17c2567f45a2f Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 01:54:07 +0530 Subject: [PATCH 2/2] feat(import): add skip_connector_syncs to omit connector-synced rows A new `LegacyWorkspace::skip_connector_syncs(bool)` option, off by default, leaves out everything v1 synced from outside services through Composio and the older connector path, since those connectors re-sync on their own. The skipped rows yield no item while the checkpoint still advances over them, and counts exclude exactly the same rows, so resumption stays exact. Auto-committed-on: macbook Co-authored-by: Medulla --- .../src/import/README.md | 28 ++- .../src/import/sections/chunks.rs | 37 +++- .../src/import/sections/connector.rs | 9 +- .../src/import/sections/connector_tests.rs | 5 - .../src/import/sections/memory_docs.rs | 6 +- .../tests/legacy_import.rs | 180 ++++++++++++++++++ 6 files changed, 244 insertions(+), 21 deletions(-) diff --git a/crates/tinymemory-integrations/src/import/README.md b/crates/tinymemory-integrations/src/import/README.md index a6be723b..98fa512e 100644 --- a/crates/tinymemory-integrations/src/import/README.md +++ b/crates/tinymemory-integrations/src/import/README.md @@ -18,6 +18,7 @@ Architecture overview: | `LegacyWorkspace::open(path)` | Detects a v1 store or refuses with a typed error. | | `LegacyWorkspace::store_suffixes(path)` / `open_store(path, suffix)` / `store_suffix()` | Lists and opens the per-profile stores (`memory-1`, `memory_tree-1`, …). | | `LegacyWorkspace::counts()` | `LegacyCounts` per section, exactly what `items()` yields: one aggregate query per `memory.db` section, the chunk store through the chunk reader, no item decoded; `total()`, `is_empty()`. Non-exhaustive. | +| `LegacyWorkspace::skip_connector_syncs(bool)` | Opt in to leaving out everything v1 synced from connectors (Composio and the connector path); see [Connector syncs](#connector-syncs). Off by default. | | `LegacyWorkspace::has_memory_db()` / `has_chunks()` | Which of the two v1 databases the workspace has. | | `LegacyWorkspace::items()` / `items_from(&Checkpoint)` | Streams `Result` from the start or after a checkpoint. | | `Items::with_page_size(n)` | Keys fetched per query (default `DEFAULT_PAGE_SIZE`, 256). Does not affect output. | @@ -44,7 +45,7 @@ the reason. Columns that later v1 migrations added are probed with `pragma_table_info` and used when present: `memory_docs.logical_namespace`, `memory_docs.taint`, `episodic_log.tool_calls_json`, `user_profile.state` / `user_state` / `class` -/ `evidence_refs_json`, and `mem_tree_chunks.content_path`. +/ `evidence_refs_json`, `mem_tree_chunks.content_path` and `mem_tree_chunks.owner`. Beside a `memory.db`, `memory_tree/chunks.db` is optional: if it is absent, not SQLite, or has no usable `mem_tree_chunks` table, the chunk section is @@ -57,7 +58,8 @@ migration must not report itself complete without it. there is anything to import and show progress against a total. Every section counts with the very predicate its scan filters by (a SQL function over Rust's `str::trim`, and for the chunk store the chunk reader itself, which reads -bodies from their files), so the counts are exactly what `items()` yields. +bodies from their files), so the counts are exactly what `items()` yields, +with or without `skip_connector_syncs`. ### Per-profile stores @@ -134,7 +136,8 @@ Chunks of one source are ordered by `(seq_in_source, id)`. A chunk's text is the file `memory_tree/content/` when the column is set, the path is a plain relative path, and the file exists; otherwise the stored preview. A `chat` source becomes a conversation of one `User` turn per chunk (chat -chunks are transcripts of host channels, whose speakers are people), with +chunks are transcripts of host channels, whose speakers are people, and, from +Composio's Slack sync, channel messages under `slack:{conn}`), with `thread_id` = `source_id` and `turns` = `0..=n-1`. Every other kind becomes a document whose body is the chunks joined by blank lines. Tags are the union of the chunks' `tags_json` plus `source_kind:`; `observed_at` is the latest @@ -231,6 +234,25 @@ modification time (none when the file system cannot report one). A missing file not UTF-8, is `Error::Io`. At most 256 KiB of a file is read: a longer one is cut there, at a character boundary, and also tagged `truncated`. +## Connector syncs + +v1 synced outside services (Gmail, Slack, Notion, Linear, GitHub, ClickUp, ...) +into its stores, through Composio and the older connector path. Those +connectors re-sync, so a host can leave them out with +`LegacyWorkspace::skip_connector_syncs(true)` (default `false`, which imports +everything). The skipped rows yield no item, so the checkpoint still advances +over them, and `counts()` excludes exactly the same rows. Conversations, +folder and file memory sources (`mem_src:*`), `conversations:agent`, meetings, +`global`, learnings, events, lessons, goals and persona files are never +skipped. The rules: + +| Section | Skipped when | +| --- | --- | +| documents (`memory_docs`) | the logical namespace starts with `skill-` (Composio `SkillDoc` sync, `skill-{toolkit}`) or `source:` (connector path, `source:{toolkit}:{conn}`, stored as `source_...`; the resolver maps it back). Taint is not consulted: v1 also marked the agent's own `global` and flow notes `external_sync`. | +| chunks | `source_kind = 'email'`; or `source_id` starts with a toolkit prefix in `CONNECTOR_TOOLKIT_PREFIXES` (`gmail:`, `slack:`, `notion:`, `linear:`, `github:`, `clickup:`); or the store has an `owner` column and any chunk of the source has `owner LIKE '%-sync:%'` (`{toolkit}-sync:{conn}`) | +| profile | `facet_id` starts with `skill-` (Composio identity facets `skill-{toolkit}-{conn}-{kind}`) | +| graph | `graph_namespace.namespace` starts with `skill-`, `source:` or `source_`; `graph_global` is unaffected | + ## Ordering and resumption Sections run in the fixed order above; within a section keys ascend in SQLite diff --git a/crates/tinymemory-integrations/src/import/sections/chunks.rs b/crates/tinymemory-integrations/src/import/sections/chunks.rs index 07bd3ec2..12992775 100644 --- a/crates/tinymemory-integrations/src/import/sections/chunks.rs +++ b/crates/tinymemory-integrations/src/import/sections/chunks.rs @@ -17,6 +17,7 @@ use std::path::{Component, Path}; use rusqlite::params; use tinymemory_api::{DocumentBody, Role, StoreItem, Turn, TurnRange}; +use super::connector; use super::{Mark, Scanned, import_meta, push_unique, sql_limit}; use crate::import::checkpoint::ChunkCursor; use crate::import::convert; @@ -59,8 +60,13 @@ pub(super) fn page( sources .into_iter() .map(|source| { + let item = if skipped(ws, store, &source)? { + None + } else { + source_item(ws, store, &source)? + }; Ok(Scanned { - item: source_item(ws, store, &source)?, + item, mark: Mark::Chunk(source), }) }) @@ -87,13 +93,40 @@ pub(super) fn count(ws: &LegacyWorkspace) -> Result { .collect::>>()?; let mut total = 0; for source in sources { - if !chunks(store, &source)?.is_empty() { + if !skipped(ws, store, &source)? && !chunks(store, &source)?.is_empty() { total = u64::saturating_add(total, 1); } } Ok(total) } +/// Whether the workspace leaves this source out as a connector sync: its +/// kind or id says so ([`connector::chunk_source_by_identity`]), or, when the +/// store has an `owner` column, any of its chunks has an owner matching +/// [`connector::OWNER_SYNC_PATTERN`]. The one test `page` and [`count`] share. +fn skipped(ws: &LegacyWorkspace, store: &ChunkStore, source: &ChunkCursor) -> Result { + if !ws.skip_connector_syncs { + return Ok(false); + } + if connector::chunk_source_by_identity(&source.source_kind, &source.source_id) { + return Ok(true); + } + if !store.owner { + return Ok(false); + } + let by_owner = store.conn.query_row( + "SELECT EXISTS(SELECT 1 FROM mem_tree_chunks \ + WHERE source_kind = ?1 AND source_id = ?2 AND owner LIKE ?3)", + params![ + source.source_kind, + source.source_id, + connector::OWNER_SYNC_PATTERN + ], + |row| row.get::<_, bool>(0), + )?; + Ok(by_owner) +} + fn source_item( ws: &LegacyWorkspace, store: &ChunkStore, diff --git a/crates/tinymemory-integrations/src/import/sections/connector.rs b/crates/tinymemory-integrations/src/import/sections/connector.rs index 01e1d4e5..002b4c91 100644 --- a/crates/tinymemory-integrations/src/import/sections/connector.rs +++ b/crates/tinymemory-integrations/src/import/sections/connector.rs @@ -11,7 +11,8 @@ //! `source:{toolkit}:{conn}`). Taint is deliberately not consulted: v1 also //! marked the agent's own `global` and flow notes `external_sync`. //! - chunks: see [`chunk_source_by_identity`] and [`OWNER_SYNC_PATTERN`]. -//! - `user_profile`: `facet_id` starting `skill-` ([`PROFILE_SKILL_PREFIX`]). +//! - `user_profile`: `facet_id` starting `skill-` (Composio identity facets +//! `skill-{toolkit}-{conn}-{kind}`). //! - `graph_namespace`: namespace starting `skill-`, `source:` or `source_`. /// Connector toolkits whose chunk `source_id` is `{toolkit}:{conn}[:{item}]` @@ -28,10 +29,6 @@ pub(crate) const EMAIL_SOURCE_KIND: &str = "email"; /// (`{toolkit}-sync:{conn}`). pub(crate) const OWNER_SYNC_PATTERN: &str = "%-sync:%"; -/// Prefix of a `user_profile.facet_id` written by Composio identity sync -/// (`skill-{toolkit}-{conn}-{kind}`). -pub(crate) const PROFILE_SKILL_PREFIX: &str = "skill-"; - /// Whether a `memory_docs` logical namespace is a connector sync. pub(crate) fn is_connector_namespace(logical: &str) -> bool { logical.starts_with("skill-") || logical.starts_with("source:") @@ -53,7 +50,7 @@ pub(crate) const GRAPH_NAMESPACE_KEPT: &str = "NOT (substr(COALESCE(namespace, ' OR substr(COALESCE(namespace, ''), 1, 7) = 'source_')"; /// SQL condition true for a `user_profile` row that is NOT a Composio -/// identity facet. +/// identity facet (`facet_id` starting `skill-`). pub(crate) const PROFILE_KEPT: &str = "substr(facet_id, 1, 6) != 'skill-'"; #[cfg(test)] diff --git a/crates/tinymemory-integrations/src/import/sections/connector_tests.rs b/crates/tinymemory-integrations/src/import/sections/connector_tests.rs index 0695270e..220d4798 100644 --- a/crates/tinymemory-integrations/src/import/sections/connector_tests.rs +++ b/crates/tinymemory-integrations/src/import/sections/connector_tests.rs @@ -24,8 +24,3 @@ fn chunk_identity() { assert!(!chunk_source_by_identity("chat", "conversations:agent")); assert!(!chunk_source_by_identity("document", "slackish:x")); } - -#[test] -fn profile_prefix_matches_constant() { - assert!(PROFILE_KEPT.contains(PROFILE_SKILL_PREFIX)); -} diff --git a/crates/tinymemory-integrations/src/import/sections/memory_docs.rs b/crates/tinymemory-integrations/src/import/sections/memory_docs.rs index 9ded5c43..8b19d8f3 100644 --- a/crates/tinymemory-integrations/src/import/sections/memory_docs.rs +++ b/crates/tinymemory-integrations/src/import/sections/memory_docs.rs @@ -212,11 +212,7 @@ fn mark_taint(row: &DocRow, tags: &mut Vec) { /// Whether the workspace leaves out a row of this logical namespace as a /// connector sync; the one test the documents scan and [`count`] share. fn skipped(ws: &LegacyWorkspace, logical: &str) -> bool { - let skip = ws.skip_connector_syncs && super::connector::is_connector_namespace(logical); - if skip { - tracing_free_note(); - } - skip + ws.skip_connector_syncs && super::connector::is_connector_namespace(logical) } fn logical_namespace(row: &DocRow) -> String { diff --git a/crates/tinymemory-integrations/tests/legacy_import.rs b/crates/tinymemory-integrations/tests/legacy_import.rs index a1dbe6a3..630608d7 100644 --- a/crates/tinymemory-integrations/tests/legacy_import.rs +++ b/crates/tinymemory-integrations/tests/legacy_import.rs @@ -1889,3 +1889,183 @@ fn an_empty_profile_tree_or_a_bad_suffix_is_not_a_store() { assert!(matches!(err, Error::NotLegacy { .. }), "{bad}: {err:?}"); } } + +/// A workspace holding both connector syncs and the data that must stay. +fn with_connector_syncs() -> tempfile::TempDir { + let (dir, conn) = workspace(&format!("{}{}", support::MEMORY_DDL, support::GRAPH_DDL)); + let add = |id: &str, ns: &str, logical: &str, content: &str| { + doc(&conn, id, ns, Some(logical), "t", content, "[]", "{}", T0); + }; + add("c01", "skill-gmail", "skill-gmail", "composio mail"); + add("c02", "source_gmail_c1", "source:gmail:c1", "connector doc"); + add("c03", "document_notes", "document:notes", "my notes"); + add("c04", "global", "global", "agent flow note"); + conn.execute( + "UPDATE memory_docs SET taint = 'external_sync' WHERE document_id IN ('c01', 'c04')", + [], + ) + .unwrap(); + facet( + &conn, + "skill-gmail-c1-name", + "identity", + "skill:gmail:name", + "Ann", + 0.9, + T0, + "active", + "auto", + None, + ); + facet( + &conn, + "f-normal", + "preference", + "tone", + "terse", + 0.9, + T0, + "active", + "auto", + None, + ); + for (ns, object) in [ + ("skill-slack", "a"), + ("source:notion:c", "b"), + ("source_notion_c", "c"), + ("document:notes", "d"), + ] { + conn.execute( + "INSERT INTO graph_namespace (namespace, subject, predicate, object, attrs_json, updated_at) \ + VALUES (?1, 's', 'p', ?2, '{}', 1.0)", + rusqlite::params![ns, object], + ) + .unwrap(); + } + conn.execute( + "INSERT INTO graph_global (subject, predicate, object, attrs_json, updated_at) \ + VALUES ('s', 'p', 'g', '{}', 1.0)", + [], + ) + .unwrap(); + let chunks = chunk_store(dir.path()); + for (id, kind, source) in [ + ("k01", "email", "thread-1"), + ("k02", "chat", "slack:conn1"), + ("k03", "document", "notion:conn1:page"), + ("k04", "document", "github:c:issue"), + ("k05", "document", "linear:c:i"), + ("k06", "document", "clickup:c:t"), + ("k07", "document", "gmail:c:m"), + ("k08", "document", "other-owner-sync"), + ("k09", "document", "mem_src:folder"), + ("k10", "chat", "conversations:agent"), + ] { + chunk(&chunks, id, kind, source, 0, 1_000, "text", "[]", None); + } + chunks + .execute( + "UPDATE mem_tree_chunks SET owner = 'jira-sync:conn1' WHERE id = 'k08'", + [], + ) + .unwrap(); + dir +} + +fn legacy_ids(ws: &LegacyWorkspace) -> Vec { + all(ws).iter().map(|i| source_id(&i.item)).collect() +} + +#[test] +fn connector_syncs_are_imported_unless_skipped() { + let dir = with_connector_syncs(); + let ws = LegacyWorkspace::open(dir.path()).unwrap(); + let ids = legacy_ids(&ws); + for id in [ + "memory_docs:c01", + "memory_docs:c02", + "user_profile:skill-gmail-c1-name", + "mem_tree_chunks:email:thread-1", + "mem_tree_chunks:document:other-owner-sync", + ] { + assert!( + ids.iter().any(|i| i == id), + "{id} missing by default: {ids:?}" + ); + } + assert_eq!(ws.counts().unwrap(), counted(&all(&ws))); +} + +#[test] +fn skip_connector_syncs_drops_exactly_the_connector_rows() { + let dir = with_connector_syncs(); + let ws = LegacyWorkspace::open(dir.path()) + .unwrap() + .skip_connector_syncs(true); + let ids = legacy_ids(&ws); + let mut expected = vec![ + "memory_docs:c03", + "memory_docs:c04", // global with external_sync taint stays + "user_profile:f-normal", + "graph_namespace:4", + "graph_global:1", + "mem_tree_chunks:chat:conversations:agent", + "mem_tree_chunks:document:mem_src:folder", + ]; + let mut got: Vec<&str> = ids.iter().map(String::as_str).collect(); + expected.sort_unstable(); + got.sort_unstable(); + assert_eq!(got, expected); + let counts = ws.counts().unwrap(); + assert_eq!(counts, counted(&all(&ws))); + assert_eq!(counts.total(), expected.len() as u64); +} + +#[test] +fn skipping_connector_syncs_keeps_resumption_exact() { + let dir = with_connector_syncs(); + let ws = LegacyWorkspace::open(dir.path()) + .unwrap() + .skip_connector_syncs(true); + let all_items = all(&ws); + for (n, imported) in all_items.iter().enumerate() { + let rest: Vec = ws + .items_from(&imported.checkpoint) + .map(|i| source_id(&i.unwrap().item)) + .collect(); + let want: Vec = all_items[n + 1..] + .iter() + .map(|i| source_id(&i.item)) + .collect(); + assert_eq!(rest, want); + } + let small: Vec = ws + .items() + .with_page_size(1) + .map(|i| source_id(&i.unwrap().item)) + .collect(); + assert_eq!(small, legacy_ids(&ws)); +} + +#[test] +fn the_owner_rule_needs_the_owner_column() { + // A chunk store without `owner` cannot match on it; the other rules hold. + let dir = tempfile::tempdir().unwrap(); + std::fs::create_dir_all(dir.path().join("memory_tree/content")).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, + content TEXT NOT NULL, seq_in_source INTEGER NOT NULL); + INSERT INTO mem_tree_chunks VALUES ('a', 'document', 'x-sync', 1, '[]', 'one', 0); + INSERT INTO mem_tree_chunks VALUES ('b', 'email', 'y', 1, '[]', 'two', 0);", + ) + .unwrap(); + drop(chunks); + let ws = LegacyWorkspace::open(dir.path()) + .unwrap() + .skip_connector_syncs(true); + assert_eq!(legacy_ids(&ws), ["mem_tree_chunks:document:x-sync"]); + assert_eq!(ws.counts().unwrap().chunks, 1); +}