diff --git a/core/src/tinycortex/sync.rs b/core/src/tinycortex/sync.rs index 1531740..54aa871 100644 --- a/core/src/tinycortex/sync.rs +++ b/core/src/tinycortex/sync.rs @@ -61,6 +61,68 @@ impl HostSyncAdapter { config: Some(config), } } + + /// Reconnect a synced Composio document to the memory tree (#5473). + /// + /// The TinyCortex migration (#4794) dropped the per-provider tree-ingest + /// half of the connector sync: synced items reached the `skill-` + /// document store but never `mem_tree_chunks`, so connector memories fell + /// out of tree-backed recall. This routes each synced item through the + /// engine's document ingest — the same L0-chunk path local folder sources + /// use via [`LocalDocumentSink`] — additively alongside the skill store. + /// + /// Scope naming matches the tree retrieval contract: the tree scope + /// (`path_scope`) is `"{toolkit}:{connection_id}"` so `query_source` resolves + /// it by platform prefix (`gmail:` → email, `slack:` → chat, …), while the + /// per-item `source_id` carries the document id so each message admits + /// independently rather than colliding on one dedup key. + /// + /// `ingest_document` writes the L0 chunk rows synchronously and enqueues the + /// summary seal on the async extract worker. Retrieval (`query_source`) reads + /// sealed summaries, so an item becomes retrievable once its buffer seals — + /// on the token threshold or the time-based `flush_stale_buffers` — and the + /// seal degrades to a fallback summary when no LLM is available. + async fn ingest_document_into_memory_tree( + &self, + config: &Config, + document: &SkillDocument, + ) -> anyhow::Result<()> { + let toolkit = document.toolkit.trim().to_ascii_lowercase(); + let connection_id = document.connection_id.trim(); + // A blank toolkit/connection would yield a scope with no platform prefix + // (`":conn"`), which no retrieval kind matches; skip rather than write an + // unreachable tree. The skill store still holds the item. + if toolkit.is_empty() || connection_id.is_empty() { + tracing::debug!( + document_id = %document.document_id, + "[tinycortex:sync] skipping memory-tree ingest: item has no toolkit/connection scope" + ); + return Ok(()); + } + let tree_scope = format!("{toolkit}:{connection_id}"); + let source_id = format!("{tree_scope}:{}", document.document_id); + let owner = format!("{toolkit}-sync:{connection_id}"); + let input = tinycortex::memory::ingest::canonicalize::document::DocumentInput { + provider: format!("composio:{toolkit}"), + title: document.title.clone(), + body: document.content.clone(), + modified_at: chrono::Utc::now(), + source_ref: Some(document.document_id.clone()), + }; + crate::ingest_pipeline::ingest_document_with_scope( + config, + &source_id, + &owner, + vec![toolkit], + input, + Some(tree_scope), + ) + .await + .map(|_| ()) + .map_err(|error| { + anyhow::anyhow!("memory-tree ingest failed for source `{source_id}`: {error}") + }) + } } /// Append one host sync audit record, logging failures without exposing source identifiers. @@ -579,14 +641,41 @@ impl SkillDocSink for HostSyncAdapter { &document.title, &document.content, Some("tinycortex-sync".into()), - Some(document.metadata), + Some(document.metadata.clone()), Some("medium".into()), None, None, - Some(document.document_id), + Some(document.document_id.clone()), ) .await - .map_err(anyhow::Error::msg) + .map_err(anyhow::Error::msg)?; + + // #5473: additively reconnect the synced item to the memory tree. This + // is a best-effort secondary index over the skill store, which is the + // source of truth and has already committed above. A failure here must + // NOT abort the connector sync: most providers do not tolerate scope + // errors, so the orchestrator turns a `store` error into a run-aborting + // `Err` — propagating would let one deterministically-poisonous item + // stall the whole connection and re-fetch the page (Composio spend) on + // every retry. Log and continue; the per-item source gate re-attempts + // the item on a later sync, and an operator rebuild can backfill. + // The config-less adapter (`sync_context`) has no ingest pipeline and is + // not on the connector sync path, so it skips tree ingest entirely. + if let Some(config) = self.config.as_deref() { + if let Err(error) = self + .ingest_document_into_memory_tree(config, &document) + .await + { + tracing::warn!( + toolkit = %document.toolkit, + connection_id = %document.connection_id, + document_id = %document.document_id, + %error, + "[tinycortex:sync] memory-tree ingest failed; skill store retained" + ); + } + } + Ok(()) } async fn delete(&self, namespace_skill_id: &str, document_id: &str) -> anyhow::Result<()> { @@ -834,4 +923,339 @@ mod tests { "expected the audit I/O error to remain distinguishable: {error:#}" ); } + + /// Regression for #5473: a Composio connector sync must feed the memory tree, + /// not just the `skill-` document store. The TinyCortex migration + /// (#4794) dropped the tree-ingest half, so synced items stopped producing + /// `mem_tree_chunks` rows and fell out of tree-backed recall. This fails if + /// the `SkillDocSink` store path ever stops writing tree chunks again. + #[tokio::test] + async fn composio_sync_document_reaches_memory_tree() { + use crate::store::{MemoryClient, MemoryClientRef}; + use std::sync::Arc; + use tinycortex::memory::sync::{SkillDocSink, SkillDocument}; + use tinymemory_api::host::test_support::TestHostConfig; + use tinymemory_api::host::MemoryHostConfig; + + crate::test_seams::init(); + let workspace = tempfile::tempdir().expect("workspace"); + let workspace_dir = workspace.path().join("workspace"); + + let mut host = TestHostConfig::default(); + host.workspace_dir = workspace_dir.clone(); + let config = host.to_arc(); + + let client: MemoryClientRef = Arc::new( + MemoryClient::from_workspace_dir(workspace_dir) + .expect("memory client initialises against a fresh workspace"), + ); + let adapter = super::HostSyncAdapter::with_config(client, config.clone()); + + // Precondition: a fresh tree is empty, so a post-store non-zero count is + // attributable to the sync path rather than to pre-existing state. + assert_eq!( + crate::store::chunks::store::count_chunks(&*config).expect("count chunks"), + 0, + "fresh workspace must start with an empty memory tree" + ); + + adapter + .store(SkillDocument { + namespace_skill_id: "gmail".into(), + connection_id: "conn-1".into(), + document_id: "gmail:msg-1".into(), + title: "Quarterly planning".into(), + content: "Let's finalise the Q3 roadmap and align on the launch date.".into(), + toolkit: "gmail".into(), + metadata: serde_json::json!({ "source": "composio-provider-incremental" }), + }) + .await + .expect("storing a synced document must also ingest it into the memory tree"); + + let chunks = crate::store::chunks::store::count_chunks(&*config).expect("count chunks"); + assert!( + chunks > 0, + "a Composio sync must add mem_tree_chunks rows for the ingested item (#5473)" + ); + + // The chunk must carry the deterministic per-item source id + // `{toolkit}:{connection_id}:{document_id}`; its `path_scope` + // (`gmail:conn-1`) is what tree retrieval resolves by platform prefix. + // A drift here is the silent "ingests but is never retrievable" trap. + let scoped = crate::store::chunks::store::list_chunks( + &*config, + &tinycortex::memory::chunks::ListChunksQuery { + source_id: Some("gmail:conn-1:gmail:msg-1".into()), + limit: Some(8), + ..Default::default() + }, + ) + .expect("list chunks by source id"); + assert!( + !scoped.is_empty(), + "ingested chunks must be keyed by the deterministic connector source id" + ); + assert!( + scoped + .iter() + .all(|chunk| chunk.metadata.path_scope.as_deref() == Some("gmail:conn-1")), + "connector chunks must carry the `{{toolkit}}:{{connection_id}}` tree scope so \ + query_source resolves them (gmail → email)" + ); + + // Retrievability is the real goal, and L0 chunks alone do NOT imply it: + // `query_source` reads sealed summaries and skips unsealed trees, so + // before a seal the freshly-ingested item is not yet retrievable. + let before = crate::tree::retrieval::query_source( + &*config, + Some("gmail:conn-1"), + None, + None, + None, + 10, + ) + .await + .expect("query_source before seal"); + assert!( + before.hits.is_empty(), + "an unsealed connector tree must not yet be retrievable" + ); + + // Drive the async extract worker to append the leaf, then force-seal the + // buffer (the time-based flush path) so a level-1 summary exists. + crate::queue::drain_until_idle(&*config) + .await + .expect("drain tree jobs"); + crate::tree::tree::flush::flush_stale_buffers( + &*config, + chrono::Duration::zero(), + &crate::tree::tree::bucket_seal::LabelStrategy::Empty, + ) + .await + .expect("force-seal stale buffers"); + + // Now the connector item is retrievable through the same path the + // product uses for tree-backed recall — the property #5473 restores. + let after = crate::tree::retrieval::query_source( + &*config, + Some("gmail:conn-1"), + None, + None, + None, + 10, + ) + .await + .expect("query_source after seal"); + assert!( + !after.hits.is_empty(), + "a sealed connector tree must be retrievable via query_source (#5473)" + ); + } + + /// The tree-ingest half of `store` is best-effort: when + /// `ingest_document_with_scope` fails, `store` must log and still return + /// `Ok(())`, so one deterministically-poisonous item cannot abort the whole + /// connector run and re-fetch the page (Composio spend) on every retry — the + /// #4947 stall that propagating the error re-created (sanil-23's review + /// blocker #2). The skill store runs first and is the source of truth, so it + /// must remain committed. This forces a real ingest failure by pointing the + /// adapter's tree-ingest `config.workspace_dir` under a regular file (so the + /// tree store cannot be created) while the skill-store client keeps a healthy + /// workspace — isolating the failure to the tree half. If `store` ever + /// propagates the ingest error again, the `.expect` on the store call fails. + #[tokio::test] + async fn tree_ingest_failure_is_tolerated_and_skill_store_is_retained() { + use crate::store::{MemoryClient, MemoryClientRef}; + use std::sync::Arc; + use tinycortex::memory::sync::{SkillDocSink, SkillDocument}; + use tinymemory_api::host::test_support::TestHostConfig; + use tinymemory_api::host::MemoryHostConfig; + + crate::test_seams::init(); + let workspace = tempfile::tempdir().expect("workspace"); + + // The skill store (source of truth) gets a healthy workspace … + let client: MemoryClientRef = Arc::new( + MemoryClient::from_workspace_dir(workspace.path().join("skill-store")) + .expect("memory client initialises against a fresh workspace"), + ); + + // … but the tree-ingest config points at a workspace *under* a regular + // file, so `ingest_document_with_scope` cannot create its store and + // returns `Err` (same failure shape as the `fallible_audit_read` guard). + let blocker = workspace.path().join("blocker"); + std::fs::write(&blocker, b"not a directory").expect("write blocker file"); + let mut host = TestHostConfig::default(); + host.workspace_dir = blocker.join("workspace"); + let config = host.to_arc(); + + let adapter = super::HostSyncAdapter::with_config(client.clone(), config.clone()); + let document = SkillDocument { + namespace_skill_id: "gmail".into(), + connection_id: "conn-1".into(), + document_id: "gmail:msg-1".into(), + title: "Quarterly planning".into(), + content: "Let's finalise the Q3 roadmap.".into(), + toolkit: "gmail".into(), + metadata: serde_json::json!({ "source": "composio-provider-incremental" }), + }; + + // Guard against a vacuous test: the tree-ingest half must *genuinely* + // fail under the broken config. If the lever ever stops failing (e.g. + // ingest resolves its store path elsewhere), this fires rather than the + // test silently passing without exercising the tolerance path. + assert!( + adapter + .ingest_document_into_memory_tree(&*config, &document) + .await + .is_err(), + "the broken tree-ingest workspace must make ingest fail" + ); + + // `store` must swallow that tree-ingest failure and still succeed. + adapter + .store(document) + .await + .expect("store must tolerate a memory-tree ingest failure (best-effort tree)"); + + // The skill store, committed before the tree half, still holds the item — + // best-effort tree ingest must never cost the durable skill write. + let skill_docs = client + .list_documents(Some("skill-gmail")) + .await + .expect("list skill-gmail documents"); + let documents = skill_docs + .get("documents") + .and_then(|value| value.as_array()) + .cloned() + .unwrap_or_default(); + assert_eq!( + documents.len(), + 1, + "the skill store must retain the synced document even when tree ingest fails" + ); + let persisted = serde_json::to_string(&documents).expect("serialise skill documents"); + assert!( + persisted.contains("gmail:msg-1"), + "the retained skill document must carry the synced id" + ); + } + + /// The config-less adapter (`sync_context`) has no ingest pipeline and is not + /// on the connector sync path, so it stores the skill document without + /// touching the memory tree. Guards the `None` branch of `store` from + /// regressing into a panic or an accidental (workspace-less) ingest. + #[tokio::test] + async fn config_less_adapter_skips_memory_tree_ingest() { + use crate::store::{MemoryClient, MemoryClientRef}; + use std::sync::Arc; + use tinycortex::memory::sync::{SkillDocSink, SkillDocument}; + use tinymemory_api::host::test_support::TestHostConfig; + use tinymemory_api::host::MemoryHostConfig; + + crate::test_seams::init(); + let workspace = tempfile::tempdir().expect("workspace"); + let workspace_dir = workspace.path().join("workspace"); + + let mut host = TestHostConfig::default(); + host.workspace_dir = workspace_dir.clone(); + let config = host.to_arc(); + + let client: MemoryClientRef = Arc::new( + MemoryClient::from_workspace_dir(workspace_dir) + .expect("memory client initialises against a fresh workspace"), + ); + // `new` leaves `config: None` — the config-less variant. Keep a handle + // to the shared client so we can read the skill store back afterwards. + let store_client = client.clone(); + let adapter = super::HostSyncAdapter::new(client); + + adapter + .store(SkillDocument { + namespace_skill_id: "gmail".into(), + connection_id: "conn-1".into(), + document_id: "gmail:msg-1".into(), + title: "Quarterly planning".into(), + content: "Let's finalise the Q3 roadmap and align on the launch date.".into(), + toolkit: "gmail".into(), + metadata: serde_json::json!({ "source": "composio-provider-incremental" }), + }) + .await + .expect("config-less store must still persist the skill document"); + + // The skill store still receives the document (the always-on half of + // `store`), keyed by its stable document id under `skill-gmail`. + let skill_docs = store_client + .list_documents(Some("skill-gmail")) + .await + .expect("list skill-gmail documents"); + let documents = skill_docs + .get("documents") + .and_then(|value| value.as_array()) + .cloned() + .unwrap_or_default(); + assert_eq!( + documents.len(), + 1, + "config-less store must persist exactly the one synced skill document" + ); + let persisted = serde_json::to_string(&documents).expect("serialise skill documents"); + assert!( + persisted.contains("gmail:msg-1") && persisted.contains("Quarterly planning"), + "the persisted skill document must carry the synced id and title" + ); + + // …but the tree is untouched, because the config-less adapter has no + // ingest pipeline. + assert_eq!( + crate::store::chunks::store::count_chunks(&*config).expect("count chunks"), + 0, + "a config-less adapter must not ingest into the memory tree" + ); + } + + /// The blank-scope guard: an item whose toolkit is empty would form an + /// unreachable `":conn"` tree scope, so `ingest_document_into_memory_tree` + /// skips it — the skill store still receives it, the tree does not. Covers + /// the early-return branch (a valid toolkit yields chunks, as the retrieval + /// test proves; a blank one must not). + #[tokio::test] + async fn blank_scope_item_is_skipped_for_memory_tree_ingest() { + use crate::store::{MemoryClient, MemoryClientRef}; + use std::sync::Arc; + use tinycortex::memory::sync::{SkillDocSink, SkillDocument}; + use tinymemory_api::host::test_support::TestHostConfig; + use tinymemory_api::host::MemoryHostConfig; + + crate::test_seams::init(); + let workspace = tempfile::tempdir().expect("workspace"); + let workspace_dir = workspace.path().join("workspace"); + let mut host = TestHostConfig::default(); + host.workspace_dir = workspace_dir.clone(); + let config = host.to_arc(); + let client: MemoryClientRef = Arc::new( + MemoryClient::from_workspace_dir(workspace_dir).expect("memory client initialises"), + ); + let adapter = super::HostSyncAdapter::with_config(client, config.clone()); + + adapter + .store(SkillDocument { + namespace_skill_id: "gmail".into(), + connection_id: "conn-1".into(), + document_id: "gmail:msg-1".into(), + title: "Quarterly planning".into(), + content: "Let's finalise the Q3 roadmap.".into(), + // Blank after trim — no platform scope can be formed. + toolkit: " ".into(), + metadata: serde_json::json!({}), + }) + .await + .expect("store must still succeed for an item without a tree scope"); + + assert_eq!( + crate::store::chunks::store::count_chunks(&*config).expect("count chunks"), + 0, + "an item without a toolkit/connection scope must be skipped for tree ingest" + ); + } }