Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 25 additions & 3 deletions crates/tinymemory-integrations/src/import/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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<ImportedItem>` 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. |
Expand All @@ -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
Expand All @@ -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

Expand Down Expand Up @@ -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/<content_path>` 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:<kind>`; `observed_at` is the latest
Expand Down Expand Up @@ -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
Expand Down
37 changes: 35 additions & 2 deletions crates/tinymemory-integrations/src/import/sections/chunks.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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),
})
})
Expand All @@ -87,13 +93,40 @@ pub(super) fn count(ws: &LegacyWorkspace) -> Result<u64> {
.collect::<rusqlite::Result<Vec<_>>>()?;
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<bool> {
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,
Expand Down
58 changes: 58 additions & 0 deletions crates/tinymemory-integrations/src/import/sections/connector.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,58 @@
//! 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-` (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}]`
/// (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:%";

/// 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:")

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Preserve non-connector source documents

When skip_connector_syncs(true) is enabled for a workspace containing folder, file, RSS, web, or GitHub-repository sources, this predicate drops their memory_docs too. The legacy MemorySourceSink::accept_source_items used source:{source_id} for every source kind, not only Composio, so the namespace prefix alone cannot identify connector content; inspect the persisted source kind/metadata before skipping these rows or the migration silently omits user source documents.

Useful? React with 👍 / 👎.

}

/// 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_')";
Comment on lines +49 to +50

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Retain graphs extracted from non-connector sources

When the option is enabled after a non-Composio source was ingested, this filter also removes that source's graph relations. The legacy ingestion path extracted graphs under the same sanitized namespace as its document, so generic source:{source_id} namespaces appear here as source_<source_id> for folders, files, RSS, web pages, and repository sources as well as connectors; blanket-filtering both prefixes contradicts the promise that those sources remain migrated.

Useful? React with 👍 / 👎.


/// SQL condition true for a `user_profile` row that is NOT a Composio
/// identity facet (`facet_id` starting `skill-`).
pub(crate) const PROFILE_KEPT: &str = "substr(facet_id, 1, 6) != 'skill-'";

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

priority medium critique uncertain

Keep profiles whose facet ID is NULL

In SQL, substr(NULL, 1, 6) != 'skill-' evaluates to NULL, not TRUE, so a WHERE clause using this condition drops every user_profile row with a NULL facet_id even though it is not a Composio identity facet. I could not verify the schema nullability from the supplied context; if facet_id is nullable, this loses profiles during import. Use COALESCE (or an explicit IS NULL branch) so only IDs beginning with skill- are excluded.

Suggested change
pub(crate) const PROFILE_KEPT: &str = "substr(facet_id, 1, 6) != 'skill-'";
pub(crate) const PROFILE_KEPT: &str = "substr(COALESCE(facet_id, ''), 1, 6) != 'skill-'";

[RULE] null-filtering ·


#[cfg(test)]
#[path = "connector_tests.rs"]
mod tests;
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
//! 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"));
}
20 changes: 16 additions & 4 deletions crates/tinymemory-integrations/src/import/sections/graph.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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| {
Expand Down Expand Up @@ -118,10 +129,11 @@ pub(super) fn count(ws: &LegacyWorkspace, table: Table) -> Result<u64> {
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))?))
}
Original file line number Diff line number Diff line change
Expand Up @@ -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,
};
Expand Down Expand Up @@ -144,7 +145,9 @@ pub(super) fn count(ws: &LegacyWorkspace, learnings: bool) -> Result<u64> {
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,
Expand Down Expand Up @@ -206,6 +209,12 @@ fn mark_taint(row: &DocRow, tags: &mut Vec<String>) {
}
}

/// 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 {
ws.skip_connector_syncs && super::connector::is_connector_namespace(logical)
}

fn logical_namespace(row: &DocRow) -> String {
resolve_logical(&row.namespace, row.logical_namespace.as_deref())
}
Expand Down
1 change: 1 addition & 0 deletions crates/tinymemory-integrations/src/import/sections/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
//! rows never stalls the iterator.

mod chunks;
pub(crate) mod connector;
mod episodic;
mod events;
mod files;
Expand Down
4 changes: 4 additions & 0 deletions crates/tinymemory-integrations/src/import/sections/profile.rs
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,10 @@ pub(super) fn count(ws: &LegacyWorkspace) -> Result<u64> {
/// 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'");
}
Expand Down
17 changes: 17 additions & 0 deletions crates/tinymemory-integrations/src/import/workspace/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,9 @@ pub struct LegacyWorkspace {
pub(crate) schema: MemorySchema,
/// `memory_tree/chunks.db`, when present and usable.
pub(crate) chunks: Option<ChunkStore>,
/// Whether rows synced from connectors are left out; see
/// [`Self::skip_connector_syncs`].
pub(crate) skip_connector_syncs: bool,
}

impl LegacyWorkspace {
Expand Down Expand Up @@ -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 {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

priority medium e2e uncertain

No end-to-end test drives the skip_connector_syncs import behaviour

skip_connector_syncs is a new public flag on LegacyWorkspace that changes what a legacy import yields across documents, chunks, profile and graph sections. No end-to-end harness reaches it: the repository's e2e harness (integration/cortexdb) runs the cortexdb service with flag environments and a mock inference server, and never opens a legacy v1 workspace or runs an import; the candidate coverage lines above are lexical matches on unrelated words (object, document, items()). An end-to-end test would have to boot the host with a real v1 workspace containing connector-synced rows, run the import with the flag set, and observe the resulting store contents/counts — none does. Coverage currently rests entirely on crates/tinymemory-integrations/tests/legacy_import.rs, which is an integration test of the library API, not the e2e lane.

[RULE] e2e-uncovered ·

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<suffix>` or
/// `memory_tree<suffix>` directory with a valid suffix. Each may still
Expand Down
3 changes: 3 additions & 0 deletions crates/tinymemory-integrations/src/import/workspace/schema.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -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,
}))
Expand Down
Loading
Loading