From 3cda23007c0b8f1de58aa26cf3272a146424ae11 Mon Sep 17 00:00:00 2001 From: M3gA-Mind Date: Wed, 7 Oct 2026 00:29:33 +0530 Subject: [PATCH 1/9] Follow up #205: survive dropped packs, validate pieces, budget beliefs - Recall builds the chosen scope's pack again when /v1/answer answers 404 for use_pack_id, up to three answers. CortexDB 0.10.4 drops every pack it holds on any successful forget, even in another scope, so a concurrent forget made recall fail (CortexDB live at 9d40d5d7). - get and list return a chunked document whole only when every piece agrees on one positive count, each index is below it, and all are present; an unchunked envelope of the same id is the whole body and wins over pieces. - The beliefs read sends max_tokens for each belief it asks for. - The live long-document test is two pieces, so its forget stays inside the request timeout, and after forget it polls fetch until no piece of the item is left. - The spec states the chunking threshold on the envelope as a piece. --- .../src/cortex/README.md | 4 +- .../src/cortex/engine/beliefs.rs | 3 +- .../src/cortex/engine/beliefs_tests.rs | 8 +++ .../src/cortex/engine/mod_tests.rs | 30 ++++++++++ .../src/cortex/engine/recall.rs | 42 ++++++++++---- .../src/cortex/envelope/mod_tests.rs | 48 ++++++++++++++++ .../src/cortex/envelope/rebuild.rs | 57 ++++++++++++------- .../src/cortex/testing/mod.rs | 3 + .../src/cortex/testing/routes.rs | 3 + .../tests/live_cortexdb.rs | 27 +++++++-- docs/architecture/cortex-flows.md | 10 +++- docs/architecture/cortex-wire.md | 7 ++- docs/specs/memory-v2.md | 2 +- 13 files changed, 201 insertions(+), 43 deletions(-) diff --git a/crates/tinymemory-integrations/src/cortex/README.md b/crates/tinymemory-integrations/src/cortex/README.md index d1a3e4fa..d1a51c03 100644 --- a/crates/tinymemory-integrations/src/cortex/README.md +++ b/crates/tinymemory-integrations/src/cortex/README.md @@ -178,7 +178,9 @@ as prefixes, so they cannot be labelled and are filtered only client-side. scopes, or, for an unscoped read, every kind scope the engine holds. With nothing to read the answer is empty and nothing is sent. The answer comes from the pack holding the most admitted events. The answer route is - called **once** with `use_pack_id`. Hosted omits a null + called **once** with `use_pack_id` (again, from a rebuilt pack, when + that pack was dropped in between, up to three answers: CortexDB drops + every pack on any forget). Hosted omits a null `answer_instructions`, because its schema is strict; Direct sends `null`. Citations come from the packs' decoded events, filtered (reach included), merged rank by rank (the most specific node's first), one per item, capped diff --git a/crates/tinymemory-integrations/src/cortex/engine/beliefs.rs b/crates/tinymemory-integrations/src/cortex/engine/beliefs.rs index e837ce41..39b3410f 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/beliefs.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/beliefs.rs @@ -34,7 +34,7 @@ use tinymemory_api::{ }; use super::CortexEngine; -use super::fetch::recall_body; +use super::fetch::{recall_body, whole_items_budget}; use super::items::hit; use crate::cortex::descriptor::{CortexWire, Route}; use crate::cortex::envelope::parse_scope; @@ -174,6 +174,7 @@ impl CortexEngine { body["budgets"]["per_layer_limits"] = json!({ "events": 0, "facts": 0, "episodes": 0, "understanding": 0, "beliefs": limit, }); + body["budgets"]["max_tokens"] = json!(whole_items_budget(limit)); let pack = self.log.recall(&body).await?; Ok(beliefs_in(&pack, "/layers/beliefs")) } diff --git a/crates/tinymemory-integrations/src/cortex/engine/beliefs_tests.rs b/crates/tinymemory-integrations/src/cortex/engine/beliefs_tests.rs index 352d8307..bd303a88 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/beliefs_tests.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/beliefs_tests.rs @@ -98,6 +98,14 @@ async fn a_built_scope_s_beliefs_are_read_with_and_without_a_query() { let belief_pack = state.seen.lock().unwrap().recalls.last().cloned().unwrap(); assert_eq!(belief_pack["include"], json!(["beliefs"])); assert_eq!(belief_pack["view"], "granular"); + let asked = belief_pack["budgets"]["per_layer_limits"]["beliefs"] + .as_u64() + .unwrap(); + assert_eq!( + belief_pack["budgets"]["max_tokens"].as_u64(), + Some(asked * crate::cortex::envelope::chunks::MAX_EVENT_TEXT_BYTES as u64), + "room for each belief asked for" + ); assert_eq!(ranked[0].meta.namespace, node); let listed = engine .beliefs(BeliefsRequest::new(Reach::exact(node), 5)) diff --git a/crates/tinymemory-integrations/src/cortex/engine/mod_tests.rs b/crates/tinymemory-integrations/src/cortex/engine/mod_tests.rs index ca0374d0..81825c63 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/mod_tests.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/mod_tests.rs @@ -320,3 +320,33 @@ async fn a_store_succeeds_when_ranked_recall_is_down() { assert_eq!(listed.items.len(), 1, "the settle probe is best-effort"); } } + +#[tokio::test] +async fn an_answer_whose_pack_expired_recalls_that_scope_again() { + for (engine, state) in both().await { + for item in items() { + engine.store(item).await.unwrap(); + } + let mut req = RecallRequest::new("which editor helix", 2); + req.filter = MetaFilter::kinds([ItemKind::Learning]); + state.seen.lock().unwrap().recalls.clear(); + state.expire_packs.store(1, Ordering::SeqCst); + let answer = engine.recall(req.clone()).await.unwrap(); + assert_eq!(answer.answer, "grounded answer for which editor helix"); + { + let seen = state.seen.lock().unwrap(); + assert_eq!(seen.recalls.len(), 2, "the one scope, packed again"); + assert_eq!(seen.recalls[0], seen.recalls[1], "the same pack request"); + assert_eq!(seen.answers.len(), 2, "answered from the fresh pack"); + } + + state.seen.lock().unwrap().answers.clear(); + state.expire_packs.store(3, Ordering::SeqCst); + assert!( + matches!(engine.recall(req).await, Err(Error::NotFound(_))), + "a bounded retry, not forever" + ); + assert_eq!(state.seen.lock().unwrap().answers.len(), 3); + state.expire_packs.store(0, Ordering::SeqCst); + } +} diff --git a/crates/tinymemory-integrations/src/cortex/engine/recall.rs b/crates/tinymemory-integrations/src/cortex/engine/recall.rs index 3adcee31..d5adb0d5 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/recall.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/recall.rs @@ -9,7 +9,9 @@ //! order, not by relevance, so it would answer from an arbitrary sample. //! With no scope to read, the answer is empty and nothing is sent. //! -//! The answer route is asked once with `use_pack_id`, so it answers from +//! The answer route is asked once with `use_pack_id` (again, from that +//! scope's pack built again, when CortexDB dropped the pack in between, up +//! to three answers: it drops every pack on any forget), so it answers from //! exactly the evidence that pack holds: with several packs, the one holding //! the most admitted events, the most specific node on a tie. **The answer //! text is grounded on that one pack, while the citations come from every @@ -43,6 +45,10 @@ use crate::cortex::error::{Error, Result}; /// Recall packs built at once when a reach spans several scopes. const PACKS_AT_ONCE: usize = 4; +/// Answers asked at most per recall: the first, and one after each pack +/// CortexDB dropped before it was used. +const ANSWER_ATTEMPTS: usize = 3; + /// The derived layers a pack also draws on, besides events. const DERIVED_LAYERS: [&str; 4] = ["facts", "beliefs", "episodes", "understanding"]; @@ -71,6 +77,13 @@ fn pack_budgets(limit: usize) -> Value { Value::Object(layers) } +/// The id a recall pack names. +fn pack_id_of(pack: &Value) -> Result<&str> { + pack.get("pack_id") + .and_then(Value::as_str) + .ok_or_else(|| Error::Engine("CortexDB recall omitted pack_id".to_string())) +} + /// The answer request body. /// /// The hosted route's schema is strict (an unknown key, or a `null` @@ -139,20 +152,29 @@ impl CortexEngine { .max_by_key(|(index, events)| (events.len(), std::cmp::Reverse(*index))) .map_or(0, |(index, _)| index); let (scope, pack) = &packs[chosen]; - let pack_id = pack - .get("pack_id") - .and_then(Value::as_str) - .ok_or_else(|| Error::Engine("CortexDB recall omitted pack_id".to_string()))?; - let answered = self - .log - .answer(&answer_body( + let ask = |pack_id: &str| { + answer_body( self.wire(), scope, &req.question, pack_id, req.instructions.as_deref(), - )) - .await?; + ) + }; + // A pack lives 60 s, and CortexDB drops every pack it holds once + // anything is forgotten (measured on 0.10.4: a forget in another + // scope turns the next `use_pack_id` into a 404). On a 404, build + // this scope's pack again and answer from it, up to + // [`ANSWER_ATTEMPTS`] answers in all. + let mut answered = self.log.answer(&ask(pack_id_of(pack)?)).await; + for _ in 1..ANSWER_ATTEMPTS { + if !matches!(answered, Err(Error::NotFound(_))) { + break; + } + let fresh = self.pack(req, scope).await?; + answered = self.log.answer(&ask(pack_id_of(&fresh)?)).await; + } + let answered = answered?; let answer = answered .get("answer") .and_then(Value::as_str) diff --git a/crates/tinymemory-integrations/src/cortex/envelope/mod_tests.rs b/crates/tinymemory-integrations/src/cortex/envelope/mod_tests.rs index f5665fc5..5837e491 100644 --- a/crates/tinymemory-integrations/src/cortex/envelope/mod_tests.rs +++ b/crates/tinymemory-integrations/src/cortex/envelope/mod_tests.rs @@ -287,3 +287,51 @@ fn a_whitespace_only_document_keeps_its_text() { assert_eq!(rebuild(&envelopes).unwrap(), item, "{body:?}"); } } + +#[test] +fn a_whole_read_needs_every_piece_of_one_agreed_layout() { + let item = long_document(); + let pieces = Envelope::for_item(&item, "id").unwrap(); + assert!(pieces.len() > 1); + assert_eq!(rebuild_whole(&pieces), Some(item.clone())); + assert_eq!(rebuild_whole(&pieces[1..]), None, "a piece missing"); + + let relaid = |edit: &dyn Fn(&mut ChunkInfo)| { + let mut changed = pieces.clone(); + edit(changed[0].chunk.as_mut().unwrap()); + rebuild_whole(&changed) + }; + assert_eq!(relaid(&|chunk| chunk.count += 1), None, "counts disagree"); + assert_eq!( + relaid(&|chunk| chunk.index = 99), + None, + "an index past the count" + ); + let mut zero = pieces[..1].to_vec(); + zero[0].chunk = Some(ChunkInfo { + index: 0, + count: 0, + pages: None, + section: None, + }); + assert_eq!(rebuild_whole(&zero), None, "a zero count"); + + let mut whole = Envelope::for_item(&StoreItem::document("x", meta()), "id") + .unwrap() + .remove(0); + whole.text = match &item { + StoreItem::Document { + body: DocumentBody::Text(text), + .. + } => text.clone(), + _ => unreachable!("a document"), + }; + whole.title.clone_from(&pieces[0].title); + whole.meta = pieces[0].meta.clone(); + let mixed = vec![pieces[1].clone(), whole]; + assert_eq!( + rebuild_whole(&mixed), + Some(item), + "the same item written whole before chunking reads as that body" + ); +} diff --git a/crates/tinymemory-integrations/src/cortex/envelope/rebuild.rs b/crates/tinymemory-integrations/src/cortex/envelope/rebuild.rs index 90a2d7f0..7ae2e849 100644 --- a/crates/tinymemory-integrations/src/cortex/envelope/rebuild.rs +++ b/crates/tinymemory-integrations/src/cortex/envelope/rebuild.rs @@ -83,35 +83,52 @@ pub(crate) fn rebuild(envelopes: &[Envelope]) -> Option { /// does not list yet. A `fetch` hit or a `recall` citation is one piece and /// uses [`rebuild`]. pub(crate) fn rebuild_whole(envelopes: &[Envelope]) -> Option { - if let Some(count) = envelopes.iter().find_map(|e| Some(e.chunk.as_ref()?.count)) { - let held: HashSet = envelopes - .iter() - .filter_map(|envelope| Some(envelope.chunk.as_ref()?.index)) - .collect(); - if (0..count).any(|index| !held.contains(&index)) { - log::debug!( - "[cortex] chunked document {:?} is missing pieces; not returned whole", - envelopes[0].id - ); - return None; - } + if !pieces_complete(envelopes) { + log::debug!( + "[cortex] chunked document {:?} is missing pieces or disagrees on its layout; \ + not returned whole", + envelopes.first().map(|envelope| &envelope.id) + ); + return None; } rebuild(envelopes) } -/// A document's text from its envelopes: the first one's, or the pieces of -/// a chunked document in index order, each once. +/// Whether `envelopes` hold a whole item: one written as a single event +/// (a document's unchunked envelope is its whole body), or every piece of a +/// chunked document, all agreeing on one positive count with every index +/// below it and each index present. +fn pieces_complete(envelopes: &[Envelope]) -> bool { + if envelopes.iter().any(|envelope| envelope.chunk.is_none()) { + return true; + } + let Some(count) = envelopes.iter().find_map(|e| Some(e.chunk.as_ref()?.count)) else { + return true; + }; + let mut held = HashSet::new(); + for chunk in envelopes + .iter() + .filter_map(|envelope| envelope.chunk.as_ref()) + { + if chunk.count != count || chunk.index >= count { + return false; + } + held.insert(chunk.index); + } + count > 0 && held.len() == count as usize +} + +/// A document's text from its envelopes: an unchunked envelope's (the whole +/// body: the same item written before chunking, since an id is a content +/// digest), else the pieces of a chunked document in index order, each once. fn document_text(envelopes: &[Envelope]) -> String { + if let Some(whole) = envelopes.iter().find(|envelope| envelope.chunk.is_none()) { + return whole.text.clone(); + } let mut pieces: Vec<(u32, &str)> = envelopes .iter() .filter_map(|envelope| Some((envelope.chunk.as_ref()?.index, envelope.text.as_str()))) .collect(); - if pieces.is_empty() { - return envelopes - .first() - .map(|envelope| envelope.text.clone()) - .unwrap_or_default(); - } pieces.sort_by_key(|(index, _)| *index); pieces.dedup_by_key(|(index, _)| *index); pieces.into_iter().map(|(_, text)| text).collect() diff --git a/crates/tinymemory-integrations/src/cortex/testing/mod.rs b/crates/tinymemory-integrations/src/cortex/testing/mod.rs index d80a39be..a714c4b7 100644 --- a/crates/tinymemory-integrations/src/cortex/testing/mod.rs +++ b/crates/tinymemory-integrations/src/cortex/testing/mod.rs @@ -77,6 +77,9 @@ pub(crate) struct Double { pub(crate) rate_limit_forget: AtomicUsize, /// Recall answers 500. pub(crate) recall_down: AtomicBool, + /// Answers refuse their `use_pack_id` as expired (404) this many times, + /// as CortexDB does once anything is forgotten after the pack was built. + pub(crate) expire_packs: AtomicUsize, /// Once the next write is applied, rate limit this many listings and /// hide the listing this many more times: (429s, hidden). Lets a test /// aim at the reads a write makes after it is sent, not the replay diff --git a/crates/tinymemory-integrations/src/cortex/testing/routes.rs b/crates/tinymemory-integrations/src/cortex/testing/routes.rs index 74d0675f..3451f082 100644 --- a/crates/tinymemory-integrations/src/cortex/testing/routes.rs +++ b/crates/tinymemory-integrations/src/cortex/testing/routes.rs @@ -297,6 +297,9 @@ async fn answer( if body["use_pack_id"].as_str() != Some("pack_test") { return fail(&state, 400, "MISSING_PACK"); } + if take_one(&state.expire_packs) { + return fail(&state, 404, "NOT_FOUND"); + } ok( &state, 200, diff --git a/crates/tinymemory-integrations/tests/live_cortexdb.rs b/crates/tinymemory-integrations/tests/live_cortexdb.rs index dfd6990e..150483c2 100644 --- a/crates/tinymemory-integrations/tests/live_cortexdb.rs +++ b/crates/tinymemory-integrations/tests/live_cortexdb.rs @@ -267,7 +267,7 @@ async fn a_long_document_round_trips_in_pieces() { for (wire, engine) in live_engines() { eprintln!("long document on {wire}"); let workspace = format!("ws-long-{}", run_id()); - let pages: Vec = (1..=24) + let pages: Vec = (1..=12) .map(|page| { let filler = format!("Clause {page} covers refunds and delivery terms. ").repeat(600); @@ -276,10 +276,14 @@ async fn a_long_document_round_trips_in_pieces() { .collect(); let body = pages.join("\u{c}"); // A document splits once its envelope passes the 256 KiB chunk - // target (not the 1 MiB event limit): this one is several pieces. + // target (not the 1 MiB event limit): this one is two pieces. It is + // kept that small because CortexDB 0.10.4 takes seconds to forget a + // ~240 KiB event (8 to 13 s measured), and a bigger document's + // forget outlasts the 60 s request timeout while the other live + // tests load the server. assert!( - body.len() > 2 * 256 * 1024, - "over twice the chunk target: {} bytes", + body.len() > 256 * 1024, + "over the chunk target: {} bytes", body.len() ); let document = StoreItem::Document { @@ -351,7 +355,7 @@ async fn a_long_document_round_trips_in_pieces() { ); let report = engine - .forget(ForgetTarget::Ids(vec![receipt.id])) + .forget(ForgetTarget::Ids(vec![receipt.id.clone()])) .await .expect("forget"); assert_eq!(report.forgotten, 1); @@ -362,7 +366,18 @@ async fn a_long_document_round_trips_in_pieces() { .expect("list after forget") .items .is_empty(), - "no piece is left behind" + "the document is gone" ); + // `list` hides a document missing a piece, so ask ranked recall, + // which hits single pieces: none may be left behind. + let deadline = Instant::now() + VISIBILITY; + loop { + let page = engine.fetch(fetch.clone()).await.expect("fetch"); + if !page.hits.iter().any(|hit| hit.id == receipt.id) { + break; + } + assert!(Instant::now() < deadline, "a piece outlived forget"); + tokio::time::sleep(Duration::from_millis(500)).await; + } } } diff --git a/docs/architecture/cortex-flows.md b/docs/architecture/cortex-flows.md index eef2b0d5..171727e3 100644 --- a/docs/architecture/cortex-flows.md +++ b/docs/architecture/cortex-flows.md @@ -152,8 +152,8 @@ Only `Hybrid`. Other modes fail `Error::Unsupported` before any request. ## Recall -Recall builds a pack, asks the answer route **once** with `use_pack_id`, and -cites from the pack. +Recall builds a pack, asks the answer route **once** with `use_pack_id` +(again from a rebuilt pack when the pack was dropped in between, step 5), and cites from the pack. 1. Resolve the scopes: a reach's kind scopes (its node and, when it inherits, every ancestor; below it too for a subtree reach), or, with **no reach** (an @@ -178,7 +178,11 @@ cites from the pack. by default. 5. Ask the answer route with that pack's scope and `use_pack_id`. A response without `answer` text is `Error::Engine`. `model` is - `diagnostics.answer_model`. + `diagnostics.answer_model`. A pack lives 60 s, and CortexDB drops every + pack it holds once anything is forgotten (measured on 0.10.4: a forget in + an unrelated scope turns the next `use_pack_id` into a 404). On that 404 + the chosen scope's pack is built again and the answer asked again, up to + three answers in all; a third 404 is returned as `Error::NotFound`. 6. **Citations** come from the packs' decoded events, merged rank by rank (each pack's best first), one per item, the most specific node's first, capped at `limit`, with `score: None` and the diff --git a/docs/architecture/cortex-wire.md b/docs/architecture/cortex-wire.md index cc3f216c..e0936234 100644 --- a/docs/architecture/cortex-wire.md +++ b/docs/architecture/cortex-wire.md @@ -199,6 +199,11 @@ ones). Response fields read: `answer` (required, string) and `diagnostics.answer_model` (optional, becomes `RecallAnswer.model`). +A 404 means the pack is gone: packs live 60 s, and CortexDB drops every +pack it holds once anything is forgotten. The scope's pack is then built +again and the answer asked again, up to three answers in all (see +[recall](cortex-flows.md#recall), step 5). + `answer_instructions` is the request's instructions when set. When unset, Direct sends `null` and TinyHumans **omits the key**: its answer schema is strict (an unknown key, or a `null` instructions, is a 400). @@ -266,7 +271,7 @@ The beliefs land in a derived layer, read two ways: ```json { "scope": "…", "query": "…", - "budgets": { "max_tokens": 786432, + "budgets": { "max_tokens": 6291456, "per_layer_limits": { "events": 0, "facts": 0, "episodes": 0, "understanding": 0, "beliefs": 8 } } } ``` diff --git a/docs/specs/memory-v2.md b/docs/specs/memory-v2.md index 87ead2e7..4659d95f 100644 --- a/docs/specs/memory-v2.md +++ b/docs/specs/memory-v2.md @@ -237,7 +237,7 @@ credentialed cleartext non-loopback endpoint, all as `Error::Config`. - **Wires.** `Direct` (`v1/experience`, `v1/events`, `v1/recall`, `v1/forget`, `v1/answer`) and `TinyHumans` (`memory/*` with `{success,data}` envelopes), as in the v1 adapter. - **Store.** - - Each item becomes one experience: a conversation becomes a bulk append of its turns, and a document whose envelope would pass 256 KiB becomes an ordered append of its pieces, cut at page breaks and headings, and a page or section still over the target cut again at blank lines, then line ends, then characters (`chunk {index, count, pages?, section?}` in the envelope; `pages` an inclusive `[first, last]` range counted from 1, only when the text marks pages; `section` only under a heading). No event over 768 KiB of encoded envelope is sent (CortexDB refuses one over 1 MiB): an item that cannot fit even split (a learning or a turn that long, or a document whose title and metadata leave no room for a piece and which does not fit whole) is `InvalidRequest`, and nothing of its batch is sent. `get`/`list` reassemble a chunked document, returning it only when every piece is present, and a ranked read gives one hit per document, its best-ranked piece, with `page:`/`section:` tags. + - Each item becomes one experience: a conversation becomes a bulk append of its turns, and a document whose encoded envelope as a chunked piece would pass 256 KiB becomes an ordered append of its pieces, cut at page breaks and headings, and a page or section still over the target cut again at blank lines, then line ends, then characters (`chunk {index, count, pages?, section?}` in the envelope; `pages` an inclusive `[first, last]` range counted from 1, only when the text marks pages; `section` only under a heading). No event over 768 KiB of encoded envelope is sent (CortexDB refuses one over 1 MiB): an item that cannot fit even split (a learning or a turn that long, or a document whose title and metadata leave no room for a piece and which does not fit whole) is `InvalidRequest`, and nothing of its batch is sent. `get`/`list` reassemble a chunked document, returning it only when every piece is present, and a ranked read gives one hit per document, its best-ranked piece, with `page:`/`section:` tags. - The envelope carries `{v:2, kind, meta, title?, learning_kind?, confidence?}`, and `meta` maps to scope labels where CortexDB can filter. - Writes wait for the indexed barrier, keeping the v1 `await_readable` behaviour. - **Scope.** One scope per item kind *per namespace node*, under the TinyMemory root `app:tinymemory` (which the hosted backend further roots under the tenant): the root node keeps `app:tinymemory/app:{documents,conversations,learnings}`, and a node adds its segments in between, e.g. `app:tinymemory/team:acme/agent:writer/app:learnings`. Namespace segments map to CortexDB's built-in `agent`, `team`, `user`, `ws` and `project` types and the kind leaf uses `app`, because CortexDB v0.10+ refuses scope types outside the deployment's `allowed_scope_types` (`422 UNREGISTERED_SCOPE_TYPE`); every shipped preset allows all of them. A `MetaFilter`'s `kinds` and `reach` pick the scopes read, each admitted kind at each node: an ordinary reach (`Reach::of`, `inherit: true`) reads its node **and every ancestor**, an exact reach its node alone, and a subtree reach its node and every node below; those nodes are known, and only a subtree reach or an unscoped read discovers nodes, from the registered scopes (`v1/scopes/list` / `memory/scopes`). Reads are always exact (`view: "granular"`, sent explicitly because public recall defaults to `holistic`), never server-side traversal. From d344303d59b1186eeb58f2b76beb76f51ebfe21b Mon Sep 17 00:00:00 2001 From: M3gA-Mind Date: Wed, 7 Oct 2026 00:49:19 +0530 Subject: [PATCH 2/9] Rebuild every pack when one is dropped; split out chunked documents - A 404 for use_pack_id now repeats the whole round: every scope's pack is built again and the answer asked from the new chosen pack, so the citations also come from packs read after the drop and cannot cite an item forgotten in between. At most three rounds. - docs/architecture/cortex-wire.md was over the 500-line limit; its "Chunked documents" section is now docs/architecture/cortex-chunks.md. --- .../src/cortex/README.md | 6 +- .../src/cortex/engine/mod_tests.rs | 13 +++ .../src/cortex/engine/recall.rs | 97 ++++++++++--------- docs/architecture/README.md | 2 +- docs/architecture/cortex-chunks.md | 51 ++++++++++ docs/architecture/cortex-flows.md | 8 +- docs/architecture/cortex-wire.md | 54 +---------- 7 files changed, 130 insertions(+), 101 deletions(-) create mode 100644 docs/architecture/cortex-chunks.md diff --git a/crates/tinymemory-integrations/src/cortex/README.md b/crates/tinymemory-integrations/src/cortex/README.md index d1a51c03..82ee4dbc 100644 --- a/crates/tinymemory-integrations/src/cortex/README.md +++ b/crates/tinymemory-integrations/src/cortex/README.md @@ -178,9 +178,9 @@ as prefixes, so they cannot be labelled and are filtered only client-side. scopes, or, for an unscoped read, every kind scope the engine holds. With nothing to read the answer is empty and nothing is sent. The answer comes from the pack holding the most admitted events. The answer route is - called **once** with `use_pack_id` (again, from a rebuilt pack, when - that pack was dropped in between, up to three answers: CortexDB drops - every pack on any forget). Hosted omits a null + called **once** with `use_pack_id` (again, after every pack is built + anew, when the packs were dropped in between, up to three rounds: + CortexDB drops every pack on any forget). Hosted omits a null `answer_instructions`, because its schema is strict; Direct sends `null`. Citations come from the packs' decoded events, filtered (reach included), merged rank by rank (the most specific node's first), one per item, capped diff --git a/crates/tinymemory-integrations/src/cortex/engine/mod_tests.rs b/crates/tinymemory-integrations/src/cortex/engine/mod_tests.rs index 81825c63..eae5723d 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/mod_tests.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/mod_tests.rs @@ -340,6 +340,19 @@ async fn an_answer_whose_pack_expired_recalls_that_scope_again() { assert_eq!(seen.answers.len(), 2, "answered from the fresh pack"); } + // With several scopes every pack is read again, so the citations + // come from packs read after whatever dropped them, not before. + let mut wide = RecallRequest::new("which editor helix", 2); + wide.filter = MetaFilter::kinds([ItemKind::Learning, ItemKind::Conversation]); + state.seen.lock().unwrap().recalls.clear(); + state.expire_packs.store(1, Ordering::SeqCst); + engine.recall(wide).await.unwrap(); + assert_eq!( + state.seen.lock().unwrap().recalls.len(), + 4, + "two scopes, each packed in both rounds" + ); + state.seen.lock().unwrap().answers.clear(); state.expire_packs.store(3, Ordering::SeqCst); assert!( diff --git a/crates/tinymemory-integrations/src/cortex/engine/recall.rs b/crates/tinymemory-integrations/src/cortex/engine/recall.rs index d5adb0d5..bc0abfbb 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/recall.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/recall.rs @@ -9,9 +9,9 @@ //! order, not by relevance, so it would answer from an arbitrary sample. //! With no scope to read, the answer is empty and nothing is sent. //! -//! The answer route is asked once with `use_pack_id` (again, from that -//! scope's pack built again, when CortexDB dropped the pack in between, up -//! to three answers: it drops every pack on any forget), so it answers from +//! The answer route is asked once with `use_pack_id` (again, after every +//! pack is built anew, when CortexDB dropped the packs in between, up to +//! three rounds: it drops every pack on any forget), so it answers from //! exactly the evidence that pack holds: with several packs, the one holding //! the most admitted events, the most specific node on a tie. **The answer //! text is grounded on that one pack, while the citations come from every @@ -45,8 +45,8 @@ use crate::cortex::error::{Error, Result}; /// Recall packs built at once when a reach spans several scopes. const PACKS_AT_ONCE: usize = 4; -/// Answers asked at most per recall: the first, and one after each pack -/// CortexDB dropped before it was used. +/// Rounds (packs, then one answer) at most per recall: the first, and one +/// after each time CortexDB dropped the packs before the answer used one. const ANSWER_ATTEMPTS: usize = 3; /// The derived layers a pack also draws on, besides events. @@ -133,48 +133,20 @@ impl CortexEngine { let mut ordered: Vec<&KindScope> = scopes.iter().collect(); ordered.sort_by_key(|scope| std::cmp::Reverse(scope.namespace.depth())); let paths: Vec = ordered.into_iter().map(|s| s.path.clone()).collect(); - let req = &req; - let packs: Vec<(String, Value)> = stream::iter(paths) - .map(|path| async move { - let pack = self.pack(req, &path).await?; - Ok::<_, Error>((path, pack)) - }) - .buffered(PACKS_AT_ONCE) - .try_collect() - .await?; - let per_pack: Vec> = packs - .iter() - .map(|(_, pack)| ranked(pack, None, &req.filter)) - .collect(); - let chosen = per_pack - .iter() - .enumerate() - .max_by_key(|(index, events)| (events.len(), std::cmp::Reverse(*index))) - .map_or(0, |(index, _)| index); - let (scope, pack) = &packs[chosen]; - let ask = |pack_id: &str| { - answer_body( - self.wire(), - scope, - &req.question, - pack_id, - req.instructions.as_deref(), - ) - }; // A pack lives 60 s, and CortexDB drops every pack it holds once // anything is forgotten (measured on 0.10.4: a forget in another - // scope turns the next `use_pack_id` into a 404). On a 404, build - // this scope's pack again and answer from it, up to - // [`ANSWER_ATTEMPTS`] answers in all. - let mut answered = self.log.answer(&ask(pack_id_of(pack)?)).await; - for _ in 1..ANSWER_ATTEMPTS { - if !matches!(answered, Err(Error::NotFound(_))) { - break; + // scope turns the next `use_pack_id` into a 404). On a 404 every + // pack is built again, so the answer and the citations both come + // from packs read after that forget, up to [`ANSWER_ATTEMPTS`] + // rounds in all. + let mut round = 1; + let (per_pack, answered) = loop { + let (per_pack, answered) = self.answer_round(&req, &paths).await?; + match answered { + Err(Error::NotFound(_)) if round < ANSWER_ATTEMPTS => round += 1, + answered => break (per_pack, answered?), } - let fresh = self.pack(req, scope).await?; - answered = self.log.answer(&ask(pack_id_of(&fresh)?)).await; - } - let answered = answered?; + }; let answer = answered .get("answer") .and_then(Value::as_str) @@ -205,6 +177,43 @@ impl CortexEngine { } impl CortexEngine { + /// One round of recall: a pack per path (four at a time), decoded and + /// filtered, and the answer route asked once with the pack holding the + /// most admitted events, the most specific node on a tie. The answer's + /// own error is handed back, so the caller can tell a dropped pack. + async fn answer_round( + &self, + req: &RecallRequest, + paths: &[String], + ) -> Result<(Vec>, Result)> { + let packs: Vec<(String, Value)> = stream::iter(paths.to_vec()) + .map(|path| async move { + let pack = self.pack(req, &path).await?; + Ok::<_, Error>((path, pack)) + }) + .buffered(PACKS_AT_ONCE) + .try_collect() + .await?; + let per_pack: Vec> = packs + .iter() + .map(|(_, pack)| ranked(pack, None, &req.filter)) + .collect(); + let chosen = per_pack + .iter() + .enumerate() + .max_by_key(|(index, events)| (events.len(), std::cmp::Reverse(*index))) + .map_or(0, |(index, _)| index); + let (scope, pack) = &packs[chosen]; + let body = answer_body( + self.wire(), + scope, + &req.question, + pack_id_of(pack)?, + req.instructions.as_deref(), + ); + Ok((per_pack, self.log.answer(&body).await)) + } + /// A recall pack for `req` over exactly `scope`, sized for `req.limit` /// citations. It includes the derived layers the answer route reads, /// after the events the citations come from. diff --git a/docs/architecture/README.md b/docs/architecture/README.md index 411c6dc3..32c30114 100644 --- a/docs/architecture/README.md +++ b/docs/architecture/README.md @@ -14,7 +14,7 @@ rustdoc next to the code. | [operations.md](operations.md) | Step-by-step semantics of store, store_many, fetch, recall, list, forget, explore and get | | [namespaces.md](namespaces.md) | The memory tree: `Namespace`, `Segment`, `Reach`, and what each operation does with them | | [cortex.md](cortex.md) | The CortexDB engine: wires, scopes, envelopes, recall | -| [cortex-wire.md](cortex-wire.md), [cortex-flows.md](cortex-flows.md) | The CortexDB wire formats and the step-by-step request flows | +| [cortex-wire.md](cortex-wire.md), [cortex-flows.md](cortex-flows.md), [cortex-chunks.md](cortex-chunks.md) | The CortexDB wire formats, the step-by-step request flows, and chunked documents | | [tools.md](tools.md) | `tinymemory-tools`: the seven agent tools, host-fixed scoping, `context.md` | | [lifecycle.md](lifecycle.md) | The agent memory lifecycle: the standard layout (brain, conversations, learnings), holistic recall, pre- and post-turn, compaction, background belief builds | | [integrations.md](integrations.md) | `tinymemory-integrations`: registry and config, documents, sources, safety, legacy import | diff --git a/docs/architecture/cortex-chunks.md b/docs/architecture/cortex-chunks.md new file mode 100644 index 00000000..b7de6bb3 --- /dev/null +++ b/docs/architecture/cortex-chunks.md @@ -0,0 +1,51 @@ +# CortexDB: chunked documents + +How the CortexDB engine writes a document too long for one event, and how it +reads one back. The event layout and labels are in +[cortex-wire.md](cortex-wire.md); the request flows in +[cortex-flows.md](cortex-flows.md). + +CortexDB refuses an experience whose flattened text is over 1 MiB +(`422 INVALID_ENVELOPE`, 0.10.4 API §6.10), and a document's event text is +its whole envelope. So (`envelope/chunks.rs`): + +- **When.** Measured as a piece would be written (the envelope with its + `chunk` field): a document whose encoded envelope fits in + `DOCUMENT_CHUNK_TARGET_BYTES` (256 KiB) is one event, byte-identical to an + unchunked one (no `chunk` field). A longer one is split. The target is the + only granularity knob: at `0` every page and every section becomes its own + event, except that a page or section too big for one event is still cut + into several (see Where). +- **Where.** First at page breaks (the form feed the PDF converter puts + between pages), then before markdown heading lines; a stretch of only + whitespace (or page breaks) never becomes a piece of its own but joins the + unit next to it, except in a text that is nothing else, which is then one + piece, so every byte is kept. These units are packed greedily, in order, + up to the target. A unit over the target is cut at blank lines, then line + ends, then characters. Sizes are JSON-escaped bytes plus the envelope + around the piece (its metadata and the `chunk` field at full width, with + a page range reserved only when the document marks pages). +- **Limit.** No event over `MAX_EVENT_TEXT_BYTES` (768 KiB of encoded + envelope, a quarter under the server's limit) is ever sent: every event of + a batch is encoded and checked before the first write, and an item over it + (a learning or a conversation turn that long) is `Error::InvalidRequest`. +- **Identity and replay.** Every piece carries the item's id and label, so + replay detection, `forget` by id or filter, and `get` see all of them; a + store that failed part way writes only the missing pieces. Until then + `get` and `list` do not return the document (never a truncated body); + `fetch` and `recall` still hit the pieces that are there. +- **Reads.** `get` and `list` give the whole document (pieces in index + order). `fetch` and `recall` give one hit or citation per document, as for + every item: its best-ranked piece, with the + item's id and its metadata plus, when known, a `page:` (or + `page:-`) tag and a `section:` tag: a document without + page breaks gets no page tag, a piece before the first heading no section + tag. These tags are read-side metadata, not part of the item's identity. + `pages` is an inclusive range `[first, last]`, counted from 1, so a piece + on one page has `first == last`. Readable CortexDB labels for page and + section are not written yet. +- **Why one piece per target rather than per page.** CortexDB 0.10.4 + already fragments every event over about 500 bytes for retrieval + (`matched_fragments`) and serves an over-budget event as an excerpt, and a + hosted write is billed per event, so splitting is used only to stay under + the limit, along the document's structure. diff --git a/docs/architecture/cortex-flows.md b/docs/architecture/cortex-flows.md index 171727e3..f7713e21 100644 --- a/docs/architecture/cortex-flows.md +++ b/docs/architecture/cortex-flows.md @@ -153,7 +153,7 @@ Only `Hybrid`. Other modes fail `Error::Unsupported` before any request. ## Recall Recall builds a pack, asks the answer route **once** with `use_pack_id` -(again from a rebuilt pack when the pack was dropped in between, step 5), and cites from the pack. +(again, from every pack built anew, when the packs were dropped in between, step 5), and cites from the packs. 1. Resolve the scopes: a reach's kind scopes (its node and, when it inherits, every ancestor; below it too for a subtree reach), or, with **no reach** (an @@ -181,8 +181,10 @@ Recall builds a pack, asks the answer route **once** with `use_pack_id` `diagnostics.answer_model`. A pack lives 60 s, and CortexDB drops every pack it holds once anything is forgotten (measured on 0.10.4: a forget in an unrelated scope turns the next `use_pack_id` into a 404). On that 404 - the chosen scope's pack is built again and the answer asked again, up to - three answers in all; a third 404 is returned as `Error::NotFound`. + the round is repeated from step 1's packs: every scope's pack is built + again and steps 3 to 5 run on the new ones, so the answer and the + citations both come from packs read after the drop. At most three rounds; + a third 404 is returned as `Error::NotFound`. 6. **Citations** come from the packs' decoded events, merged rank by rank (each pack's best first), one per item, the most specific node's first, capped at `limit`, with `score: None` and the diff --git a/docs/architecture/cortex-wire.md b/docs/architecture/cortex-wire.md index e0936234..2fde8b86 100644 --- a/docs/architecture/cortex-wire.md +++ b/docs/architecture/cortex-wire.md @@ -200,8 +200,9 @@ Response fields read: `answer` (required, string) and `diagnostics.answer_model` (optional, becomes `RecallAnswer.model`). A 404 means the pack is gone: packs live 60 s, and CortexDB drops every -pack it holds once anything is forgotten. The scope's pack is then built -again and the answer asked again, up to three answers in all (see +pack it holds once anything is forgotten. Every scope's pack is then built +again and the answer asked from the new chosen pack, up to three rounds in +all (see [recall](cortex-flows.md#recall), step 5). `answer_instructions` is the request's instructions when set. When unset, @@ -364,58 +365,11 @@ double enforces). parent-scope sample. - A filter whose `kinds` admits nothing reads no scopes. -## Chunked documents - -CortexDB refuses an experience whose flattened text is over 1 MiB -(`422 INVALID_ENVELOPE`, 0.10.4 API §6.10), and a document's event text is -its whole envelope. So (`envelope/chunks.rs`): - -- **When.** Measured as a piece would be written (the envelope with its - `chunk` field): a document whose encoded envelope fits in - `DOCUMENT_CHUNK_TARGET_BYTES` (256 KiB) is one event, byte-identical to an - unchunked one (no `chunk` field). A longer one is split. The target is the - only granularity knob: at `0` every page and every section becomes its own - event, except that a page or section too big for one event is still cut - into several (see Where). -- **Where.** First at page breaks (the form feed the PDF converter puts - between pages), then before markdown heading lines; a stretch of only - whitespace (or page breaks) never becomes a piece of its own but joins the - unit next to it, except in a text that is nothing else, which is then one - piece, so every byte is kept. These units are packed greedily, in order, - up to the target. A unit over the target is cut at blank lines, then line - ends, then characters. Sizes are JSON-escaped bytes plus the envelope - around the piece (its metadata and the `chunk` field at full width, with - a page range reserved only when the document marks pages). -- **Limit.** No event over `MAX_EVENT_TEXT_BYTES` (768 KiB of encoded - envelope, a quarter under the server's limit) is ever sent: every event of - a batch is encoded and checked before the first write, and an item over it - (a learning or a conversation turn that long) is `Error::InvalidRequest`. -- **Identity and replay.** Every piece carries the item's id and label, so - replay detection, `forget` by id or filter, and `get` see all of them; a - store that failed part way writes only the missing pieces. Until then - `get` and `list` do not return the document (never a truncated body); - `fetch` and `recall` still hit the pieces that are there. -- **Reads.** `get` and `list` give the whole document (pieces in index - order). `fetch` and `recall` give one hit or citation per document, as for - every item: its best-ranked piece, with the - item's id and its metadata plus, when known, a `page:<n>` (or - `page:<first>-<last>`) tag and a `section:<title>` tag: a document without - page breaks gets no page tag, a piece before the first heading no section - tag. These tags are read-side metadata, not part of the item's identity. - `pages` is an inclusive range `[first, last]`, counted from 1, so a piece - on one page has `first == last`. Readable CortexDB labels for page and - section are not written yet. -- **Why one piece per target rather than per page.** CortexDB 0.10.4 - already fragments every event over about 500 bytes for retrieval - (`matched_fragments`) and serves an over-budget event as an excerpt, and a - hosted write is billed per event, so splitting is used only to stay under - the limit, along the document's structure. - ## The v2 envelope A learning is one event; a conversation is one event per turn, appended in order. A document is one event unless its envelope would pass 256 KiB; then it -is one event per piece (see [Chunked documents](#chunked-documents)). +is one event per piece (see [cortex-chunks.md](cortex-chunks.md)). CortexDB's experience schema is closed (an unknown field is a 422), so the structured data rides in the one free-form field: the event's `content.text` is a JSON **envelope**: From cb42b6c5d53a43becd15f26045272e746b54c395 Mon Sep 17 00:00:00 2001 From: M3gA-Mind <elvin@tinyhumans.ai> Date: Wed, 7 Oct 2026 01:30:54 +0530 Subject: [PATCH 3/9] Run the live tests one at a time; keep the long document at three pieces Run together, one live test's forgets drop the packs another is about to answer from, and load the server enough that forgetting a ~700 KiB document (seconds per ~240 KiB event on CortexDB 0.10.4) outlasted the request timeout. Each live test now holds one async lock for its run, and the long document is back to 24 pages (three pieces, two boundaries). --- .../tests/live_cortexdb.rs | 26 +++++++++++++------ 1 file changed, 18 insertions(+), 8 deletions(-) diff --git a/crates/tinymemory-integrations/tests/live_cortexdb.rs b/crates/tinymemory-integrations/tests/live_cortexdb.rs index 150483c2..e3fe9f20 100644 --- a/crates/tinymemory-integrations/tests/live_cortexdb.rs +++ b/crates/tinymemory-integrations/tests/live_cortexdb.rs @@ -102,8 +102,16 @@ async fn list_until(engine: &CortexEngine, filter: &MetaFilter, want: usize) -> } } +/// Held by each live test for its whole run, so they reach the server one +/// at a time. Run together, one test's forgets drop the packs another is +/// about to answer from, and load the server enough that forgetting a long +/// document (seconds per ~240 KiB event on 0.10.4) outlasts the request +/// timeout. +static ONE_AT_A_TIME: futures::lock::Mutex<()> = futures::lock::Mutex::new(()); + #[tokio::test] async fn the_live_server_upholds_the_contract() { + let _alone = ONE_AT_A_TIME.lock().await; for (wire, engine) in live_engines() { eprintln!("conformance on {wire}"); tinymemory_api::conformance::run(&engine) @@ -114,6 +122,7 @@ async fn the_live_server_upholds_the_contract() { #[tokio::test] async fn documents_conversations_and_learnings_round_trip_into_context() { + let _alone = ONE_AT_A_TIME.lock().await; for (wire, engine) in live_engines() { eprintln!("round trip on {wire}"); round_trip(&engine).await; @@ -264,10 +273,11 @@ async fn round_trip(engine: &CortexEngine) { /// section, and `forget` removes every piece. #[tokio::test] async fn a_long_document_round_trips_in_pieces() { + let _alone = ONE_AT_A_TIME.lock().await; for (wire, engine) in live_engines() { eprintln!("long document on {wire}"); let workspace = format!("ws-long-{}", run_id()); - let pages: Vec<String> = (1..=12) + let pages: Vec<String> = (1..=24) .map(|page| { let filler = format!("Clause {page} covers refunds and delivery terms. ").repeat(600); @@ -276,14 +286,14 @@ async fn a_long_document_round_trips_in_pieces() { .collect(); let body = pages.join("\u{c}"); // A document splits once its envelope passes the 256 KiB chunk - // target (not the 1 MiB event limit): this one is two pieces. It is - // kept that small because CortexDB 0.10.4 takes seconds to forget a - // ~240 KiB event (8 to 13 s measured), and a bigger document's - // forget outlasts the 60 s request timeout while the other live - // tests load the server. + // target (not the 1 MiB event limit): this one is three pieces, so + // two boundaries are crossed. It stays under 1 MiB because CortexDB + // 0.10.4 takes seconds to forget a ~240 KiB event (8 to 13 s + // measured), and a longer document's forget can outlast the 60 s + // request timeout. assert!( - body.len() > 256 * 1024, - "over the chunk target: {} bytes", + body.len() > 2 * 256 * 1024, + "over twice the chunk target: {} bytes", body.len() ); let document = StoreItem::Document { From 6fbd6d4ed2f413f8e3985d627571d8a1cdeeb2ca Mon Sep 17 00:00:00 2001 From: M3gA-Mind <elvin@tinyhumans.ai> Date: Wed, 7 Oct 2026 00:35:28 +0530 Subject: [PATCH 4/9] Write events as prose with the envelope in labels; let items opt out of derivation CortexDB extracts from an event's text and splits it for search at sentence boundaries, which a JSON text lacks, and counts that text toward its 1 MiB limit; labels are its app-metadata extension point. So events are now written as v3: - content.text is the item's own text: the body or piece, the turn's text, or the learning's statement; - context.labels hold the lookup labels, readable kind:/file:/page:/ section: labels (at most 256 bytes each, never lang:), and the rest of the envelope as compact JSON in tm:e:<NN>: parts of at most 240 bytes. An event with empty text, or whose labels would pass 64, is written as v2 (the whole envelope as JSON text) as every event was before. Readers take both, so existing stores need no rewrite. A recovered hosted write is matched on its labels as well as its text, since two v3 turns can say the same words. MemoryMeta gains derive: Option<bool>. Some(false) stores and indexes an item but asks CortexDB to derive nothing from it (directives.extract: []); every tool turn is sent the same way. Unset, it is not serialized, so fingerprints are unchanged. The test double's recall no longer prefixes a pack event's text with [role]: CortexDB 0.10.3 and 0.10.4 return the stored text there. --- crates/tinymemory-api/src/item/mod_tests.rs | 13 + .../tinymemory-api/src/meta/filter_tests.rs | 1 + crates/tinymemory-api/src/meta/mod.rs | 9 + .../src/cortex/README.md | 20 +- .../src/cortex/engine/mod_chunk_tests.rs | 14 +- .../src/cortex/engine/mod_hosted_tests.rs | 27 +++ .../src/cortex/engine/store.rs | 4 +- .../src/cortex/envelope/mod.rs | 225 ++++++++++++++---- .../src/cortex/envelope/mod_tests.rs | 210 +++++++++++++++- .../src/cortex/envelope/rebuild.rs | 13 +- .../src/cortex/log/write.rs | 5 + .../tinymemory-integrations/src/cortex/mod.rs | 8 +- .../src/cortex/testing/log.rs | 10 +- docs/architecture/cortex-chunks.md | 11 +- docs/architecture/cortex-wire.md | 42 +++- docs/architecture/cortex.md | 2 +- docs/specs/memory-v2.md | 3 +- 17 files changed, 529 insertions(+), 88 deletions(-) diff --git a/crates/tinymemory-api/src/item/mod_tests.rs b/crates/tinymemory-api/src/item/mod_tests.rs index 5b14055a..b0ba261e 100644 --- a/crates/tinymemory-api/src/item/mod_tests.rs +++ b/crates/tinymemory-api/src/item/mod_tests.rs @@ -173,3 +173,16 @@ fn a_turn_renders_its_tool_calls() { "assistant: done [tools: shell (call-1), grep]" ); } + +#[test] +fn an_unset_derive_flag_leaves_the_fingerprint_as_it_was() { + let item = StoreItem::document("Refunds take five days.", MemoryMeta::default()); + let json = serde_json::to_string(item.meta()).unwrap(); + assert!(!json.contains("derive"), "{json}"); + let mut opted_out = item.clone(); + opted_out.meta_mut().derive = Some(false); + assert_ne!(opted_out.fingerprint(), item.fingerprint()); + let back: MemoryMeta = + serde_json::from_str(&serde_json::to_string(opted_out.meta()).unwrap()).unwrap(); + assert_eq!(back.derive, Some(false)); +} diff --git a/crates/tinymemory-api/src/meta/filter_tests.rs b/crates/tinymemory-api/src/meta/filter_tests.rs index c1c3a6cf..6275ab93 100644 --- a/crates/tinymemory-api/src/meta/filter_tests.rs +++ b/crates/tinymemory-api/src/meta/filter_tests.rs @@ -27,6 +27,7 @@ fn meta() -> MemoryMeta { }, tags: vec!["x".into(), "y".into()], observed_at: Some(Utc.with_ymd_and_hms(2026, 1, 2, 3, 4, 5).unwrap()), + derive: None, } } diff --git a/crates/tinymemory-api/src/meta/mod.rs b/crates/tinymemory-api/src/meta/mod.rs index 8b7ef77f..3a4328f4 100644 --- a/crates/tinymemory-api/src/meta/mod.rs +++ b/crates/tinymemory-api/src/meta/mod.rs @@ -68,6 +68,15 @@ pub struct MemoryMeta { /// When the underlying fact was observed, as opposed to when it was stored. #[serde(skip_serializing_if = "Option::is_none")] pub observed_at: Option<DateTime<Utc>>, + /// Whether an engine that derives layers from what it stores (facts, + /// beliefs, concepts) may derive them from this item. `Some(false)` + /// stores and indexes the item, so it stays searchable, but derives + /// nothing from it: runtime state, tool output, machine-written + /// summaries. `None`, the default, leaves it to the engine; engines that + /// derive nothing ignore it. Unset, it is not serialized, so an item's + /// fingerprint is the same as before the field existed. + #[serde(skip_serializing_if = "Option::is_none")] + pub derive: Option<bool>, } impl MemoryMeta { diff --git a/crates/tinymemory-integrations/src/cortex/README.md b/crates/tinymemory-integrations/src/cortex/README.md index 82ee4dbc..4a4f0445 100644 --- a/crates/tinymemory-integrations/src/cortex/README.md +++ b/crates/tinymemory-integrations/src/cortex/README.md @@ -21,7 +21,7 @@ This README is the short in-tree summary. The full reference is under - [`cortex.md`](../../../../docs/architecture/cortex.md): surface, credentials, transport, failure mapping, endpoint security, the registry and `MemoryConfig`; - [`cortex-wire.md`](../../../../docs/architecture/cortex-wire.md): every endpoint - and its shapes, scope layout, the v2 envelope, lookup labels; + and its shapes, scope layout, the envelope (v3 and v2), lookup labels; - [`cortex-flows.md`](../../../../docs/architecture/cortex-flows.md): step-by-step store, list, fetch, recall, forget, get, discovery; - [`testing.md`](../../../../docs/architecture/testing.md): the doubles, the @@ -122,18 +122,26 @@ reassemble a chunked document, and return it only when every piece is present; ` its best-ranked piece, with a `page:<n>` (or `page:<first>-<last>`) tag when the document marks its pages and a `section:<title>` tag when the piece starts under a heading; a piece with neither carries no extra tag. Each event's -`content.text` is a JSON envelope: +`content.text` is the item's own text (the body or piece, the turn's text, or +the statement), and the rest of its envelope rides in `context.labels` as +`tm:e:<NN>:` parts of at most 240 bytes of JSON (v3): ```json -{ "v": 2, "id": "<40-hex fingerprint>", "kind": "conversation", - "text": "<body | turn text | statement>", "meta": { ... MemoryMeta ... }, +{ "v": 3, "id": "<40-hex fingerprint>", "kind": "conversation", "text": "", + "meta": { ... MemoryMeta ... }, "title": "...", "mime": "...", "learning_kind": "...", "confidence": 0.8, "evidence": "...", "turn": { "index": 0, "count": 3, "role": "user", "at": "...", "tool_calls": [] } } ``` -Kind-specific fields appear only when set. Text that is not a v2 envelope is -someone else's event and is ignored. `context.observed_at` carries the turn's +Kind-specific fields appear only when set. Readable labels (`kind:`, `file:`, +`page:`, `section:`) sit beside the parts. An event with empty text, or whose +labels would pass 64, is written as v2 (the whole envelope as JSON text, as +every event was before), and both layouts read. An event that is neither is +someone else's and is ignored. An item that opts out of derivation +(`MemoryMeta::derive == Some(false)`), and every tool turn, is sent with +`directives.extract: []`: indexed and searchable, but no facts, beliefs or +concepts are derived from it. `context.observed_at` carries the turn's `at` or the item's `meta.observed_at`. **Labels.** Each event carries up to eight `context.labels`, each a 16-hex diff --git a/crates/tinymemory-integrations/src/cortex/engine/mod_chunk_tests.rs b/crates/tinymemory-integrations/src/cortex/engine/mod_chunk_tests.rs index 24ef5ae2..b6b047de 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/mod_chunk_tests.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/mod_chunk_tests.rs @@ -189,7 +189,7 @@ async fn a_store_that_lost_a_piece_writes_only_that_piece_again() { } #[tokio::test] -async fn a_short_document_is_one_event_exactly_as_before_and_reads_the_same() { +async fn a_short_document_is_one_event_and_reads_the_same() { for (engine, state) in both().await { let item = StoreItem::Document { title: Some("Note".into()), @@ -200,9 +200,15 @@ async fn a_short_document_is_one_event_exactly_as_before_and_reads_the_same() { engine.store(item.clone()).await.unwrap(); let written = events(&state, SCOPE); assert_eq!(written.len(), 1); - let text = written[0]["content"]["text"].as_str().unwrap(); - let envelope: Value = serde_json::from_str(text).unwrap(); - assert!(envelope.get("chunk").is_none(), "no piece info: {envelope}"); + assert_eq!( + written[0]["content"]["text"], + format!("Page one.{PAGE_BREAK}Page two."), + "the body itself" + ); + let envelope = crate::cortex::envelope::decode_event(&written[0]) + .unwrap() + .envelope; + assert!(envelope.chunk.is_none(), "no piece info: {envelope:?}"); let listed = engine .list(ListRequest::new(MetaFilter::default(), 10)) .await diff --git a/crates/tinymemory-integrations/src/cortex/engine/mod_hosted_tests.rs b/crates/tinymemory-integrations/src/cortex/engine/mod_hosted_tests.rs index 8ac93e5c..358a2654 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/mod_hosted_tests.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/mod_hosted_tests.rs @@ -395,3 +395,30 @@ async fn the_answer_body_holds_only_keys_the_strict_schema_allows() { with.instructions = Some("be brief".into()); engine.recall(with).await.unwrap(); } + +#[tokio::test] +async fn recovery_does_not_take_another_turn_with_the_same_words_for_the_lost_one() { + use tinymemory_api::{MemoryMeta, Role, Turn}; + let (endpoint, state) = hosted_double().await; + let engine = hosted_engine(&endpoint).with_test_timing(std::time::Duration::from_millis(200)); + let item = StoreItem::Conversation { + turns: vec![ + Turn::new(Role::User, "ok"), + Turn::new(Role::Assistant, "ok"), + ], + meta: MemoryMeta::default(), + }; + engine.store(item.clone()).await.unwrap(); + { + // The second turn's event is lost. + let mut log = state.log.lock().unwrap(); + let last = log.events.len() - 1; + log.events.remove(last); + } + // Its re-write is claimed but never applied, so recovery must look for + // it, and must not take the first turn, which says the same words. + state.claim_then_fail.store(1, Ordering::SeqCst); + let error = engine.store(item).await.unwrap_err(); + assert!(matches!(error, Error::Unavailable(_)), "{error:?}"); + assert_eq!(state.event_count(), 1, "the lost turn is still missing"); +} diff --git a/crates/tinymemory-integrations/src/cortex/engine/store.rs b/crates/tinymemory-integrations/src/cortex/engine/store.rs index 8f558dd4..52ae5534 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/store.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/store.rs @@ -26,7 +26,7 @@ use tinymemory_api::{ItemId, StoreItem, StoreReceipt, WaitFor, validate_many}; use super::CortexEngine; use super::scopes::KindScope; -use crate::cortex::envelope::Envelope; +use crate::cortex::envelope::{Encoded, Envelope}; use crate::cortex::error::Result; use crate::cortex::log::Written; @@ -53,7 +53,7 @@ impl CortexEngine { let ids: Vec<String> = items.iter().map(StoreItem::fingerprint).collect(); // Every event of the batch is laid out and size-checked before any is // sent, so an item CortexDB would refuse leaves nothing half-written. - let mut planned: Vec<Vec<(Option<u32>, Envelope, String)>> = + let mut planned: Vec<Vec<(Option<u32>, Envelope, Encoded)>> = Vec::with_capacity(items.len()); for (item, id) in items.iter().zip(&ids) { let mut events = Vec::new(); diff --git a/crates/tinymemory-integrations/src/cortex/envelope/mod.rs b/crates/tinymemory-integrations/src/cortex/envelope/mod.rs index 1485a973..9167bb7c 100644 --- a/crates/tinymemory-integrations/src/cortex/envelope/mod.rs +++ b/crates/tinymemory-integrations/src/cortex/envelope/mod.rs @@ -32,13 +32,18 @@ //! No event is sent whose encoded envelope is over //! [`chunks::MAX_EVENT_TEXT_BYTES`] ([`Envelope::encode_checked`]): CortexDB //! refuses an experience over 1 MiB of text, and a conversation turn or a -//! learning cannot be split. Each event's `content.text` is a JSON -//! [`Envelope`] (`"v": 2`) carrying the item id, kind, the event's own text -//! (the body, the turn's text, or the learning's statement), the item's full -//! [`MemoryMeta`], and the kind's extra fields. CortexDB's experience schema -//! is closed (an unknown field is a 422), so the envelope rides in the one -//! free-form field there is; anything that does not parse as a v2 envelope is -//! somebody else's event and is ignored. +//! learning cannot be split. An [`Envelope`] carries the item id, kind, the +//! event's own text (the body, the turn's text, or the learning's +//! statement), the item's full [`MemoryMeta`], and the kind's extra fields. +//! CortexDB's experience schema is closed (an unknown field is a 422), and +//! `context.labels` is its app-metadata extension point. So an event is +//! written as v3: `content.text` is the event's own text, which CortexDB +//! extracts from and splits for search at sentence boundaries, and the rest +//! of the envelope rides in `tm:e:<NN>:` labels ([`Envelope::encode_checked`]). +//! An event with empty text, or too much envelope for its labels, is written +//! as v2: the whole envelope as JSON text, as every event was before v3. +//! Both read back ([`decode_event`]); anything else is somebody else's event +//! and is ignored. //! //! Each event also carries lookup labels (see [`labels`]) and, when the item //! has one, `context.observed_at`. @@ -67,8 +72,37 @@ pub(crate) use rebuild::{Decoded, decode_event, rebuild, rebuild_whole}; /// The TinyMemory root every kind scope sits under. pub(crate) const ROOT_SCOPE: &str = "app:tinymemory"; -/// The envelope version this crate writes and reads. -const VERSION: u8 = 2; +/// The envelope version of an event whose text is the whole JSON envelope: +/// every event before v3, and a v3-era event whose envelope does not fit in +/// its labels. +const V2: u8 = 2; + +/// The envelope version of an event whose text is the item's own text and +/// whose envelope (all but the text) rides in its labels. +const V3: u8 = 3; + +/// The label prefix of a v3 envelope part: `tm:e:<NN>:<JSON slice>`. +const PART_PREFIX: &str = "tm:e:"; + +/// The most bytes of envelope JSON one part label carries, so a label stays +/// within the 256 bytes CortexDB asks labels to keep to. +const PART_BYTES: usize = 240; + +/// The most labels one event carries, as CortexDB asks. +const MAX_LABELS: usize = 64; + +/// The longest readable label written; a longer one is left out. +const MAX_LABEL_BYTES: usize = 256; + +/// One event as it is sent: its `content.text` and `context.labels`. +#[derive(Debug, Clone, PartialEq)] +pub(crate) struct Encoded { + /// The item's own text (v3), or the whole JSON envelope (v2). + pub(crate) text: String, + /// The lookup labels, then (v3) the readable labels and the envelope + /// parts. + pub(crate) labels: Vec<String>, +} /// The leaf segment of `kind`'s scope inside a namespace node. pub(crate) fn kind_leaf(kind: ItemKind) -> &'static str { @@ -109,7 +143,7 @@ pub(crate) fn parse_scope(path: &str) -> Option<(Namespace, ItemKind)> { /// One event's payload: the item it belongs to and the event's share of it. #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] pub(crate) struct Envelope { - /// Always [`VERSION`]. + /// [`V2`] or [`V3`] as stored; [`V3`] for an envelope laid out here. pub(crate) v: u8, /// The item id ([`StoreItem::fingerprint`]). pub(crate) id: String, @@ -184,7 +218,7 @@ impl Envelope { /// A bare envelope for item `id` of `kind`. fn new(id: &str, kind: ItemKind, text: String, meta: &MemoryMeta) -> Self { Self { - v: VERSION, + v: V3, id: id.to_string(), kind, text, @@ -306,54 +340,135 @@ impl Envelope { } } - /// Reads an envelope from an event's text, whichever read path it came - /// from. + /// Reads a v2 envelope from an event's text, whichever read path it + /// came from. /// - /// The two read paths disagree on the bytes: `/v1/events` returns the - /// text as stored, while `/v1/recall` renders it for a reader and - /// prefixes the speaker (`[user] {...}`). The prefix is stripped only - /// when the text does not parse without it. Anything that is not a v2 - /// envelope is `None`. + /// The two read paths may disagree on the bytes: `/v1/events` returns + /// the text as stored, while an older `/v1/recall` rendered it for a + /// reader and prefixed the speaker (`[user] {...}`). The prefix is + /// stripped only when the text does not parse without it. Anything that + /// is not a v2 envelope is `None`. pub(crate) fn decode(text: &str) -> Option<Self> { let parsed = serde_json::from_str::<Self>(text).ok().or_else(|| { let rendered = text.strip_prefix('[')?; let (_role, rest) = rendered.split_once("] ")?; serde_json::from_str::<Self>(rest).ok() })?; - (parsed.v == VERSION).then_some(parsed) + (parsed.v == V2).then_some(parsed) } - /// The stored text. + /// Reads a v3 envelope from an event's text and labels: the envelope + /// parts (`tm:e:00:`, `tm:e:01:`, …) joined in order, with `text` as its + /// text. `None` when the parts are absent, not numbered `0..n`, or not a + /// v3 envelope. + pub(crate) fn from_labels<'a>( + text: &str, + labels: impl IntoIterator<Item = &'a str>, + ) -> Option<Self> { + let mut parts: Vec<(usize, &str)> = Vec::new(); + for label in labels { + if let Some(rest) = label.strip_prefix(PART_PREFIX) { + let (index, json) = rest.split_once(':')?; + parts.push((index.parse().ok()?, json)); + } + } + parts.sort_by_key(|(index, _)| *index); + if parts.is_empty() + || parts + .iter() + .enumerate() + .any(|(at, (index, _))| at != *index) + { + return None; + } + let json: String = parts.into_iter().map(|(_, json)| json).collect(); + let mut envelope = serde_json::from_str::<Self>(&json).ok()?; + if envelope.v != V3 { + return None; + } + envelope.text = text.to_string(); + Some(envelope) + } + + /// The whole envelope as v2 JSON: what a v2 event's text holds, and the + /// size every limit here is checked against. /// /// # Errors /// /// [`Error::Engine`] if serialisation fails, which plain data cannot. pub(crate) fn encode(&self) -> Result<String> { - serde_json::to_string(self) - .map_err(|_| Error::Engine("an item envelope could not be serialised".to_string())) + let mut v2 = self.clone(); + v2.v = V2; + json_of(&v2) } - /// The stored text, refused when it is over - /// [`chunks::MAX_EVENT_TEXT_BYTES`]: CortexDB would refuse the event, - /// so nothing of the item is sent. + /// The event as sent, refused when its v2 JSON is over + /// [`chunks::MAX_EVENT_TEXT_BYTES`]: CortexDB would refuse such an event + /// as v2, so nothing of the item is sent. (Checking the v2 size keeps one + /// limit for both layouts, and a v3 text is always shorter.) + /// + /// The event is v3 (the item's own text, the envelope in labels) unless + /// its text is empty or its labels would be more than [`MAX_LABELS`]; + /// then it is v2, as every event was before. /// /// # Errors /// /// [`Error::InvalidRequest`] for an envelope over the limit (a /// conversation turn or a learning that long, or metadata too large to /// leave room for a document piece); as [`Envelope::encode`] otherwise. - pub(crate) fn encode_checked(&self) -> Result<String> { - let encoded = self.encode()?; - if encoded.len() > chunks::MAX_EVENT_TEXT_BYTES { + pub(crate) fn encode_checked(&self) -> Result<Encoded> { + let v2 = self.encode()?; + if v2.len() > chunks::MAX_EVENT_TEXT_BYTES { return Err(Error::InvalidRequest(format!( "a {:?} event would be {} bytes; CortexDB refuses an event over 1 MiB, so at \ most {} are sent", self.kind, - encoded.len(), + v2.len(), chunks::MAX_EVENT_TEXT_BYTES ))); } - Ok(encoded) + let lookup = labels::for_item(&self.id, &self.meta); + let mut head = self.clone(); + head.v = V3; + head.text = String::new(); + let parts = part_labels(&json_of(&head)?); + let readable = self.readable_labels(); + if self.text.is_empty() || lookup.len() + readable.len() + parts.len() > MAX_LABELS { + return Ok(Encoded { + text: v2, + labels: lookup, + }); + } + let mut labels = lookup; + labels.extend(readable); + labels.extend(parts); + Ok(Encoded { + text: self.text.clone(), + labels, + }) + } + + /// Labels a person reading the events can make sense of: the item kind, + /// and for a document its file, and a piece's pages and section. Never + /// filtered on (the lookup labels are), and left out when longer than + /// [`MAX_LABEL_BYTES`]. + fn readable_labels(&self) -> Vec<String> { + let mut out = vec![format!("kind:{}", self.kind.as_str())]; + if let Some(path) = &self.meta.file_path { + out.push(format!("file:{path}")); + } + if let Some(chunk) = &self.chunk { + match chunk.pages { + Some([first, last]) if first == last => out.push(format!("page:{first}")), + Some([first, last]) => out.push(format!("page:{first}-{last}")), + None => {} + } + if let Some(section) = &chunk.section { + out.push(format!("section:{section}")); + } + } + out.retain(|label| label.len() <= MAX_LABEL_BYTES); + out } /// The encoded size of this envelope as a document piece with an empty @@ -371,9 +486,9 @@ impl Envelope { Ok(probe.encode()?.len() + SECTION_RESERVE) } - /// The experience request appending this envelope, with a fresh body - /// idempotency key. - pub(crate) fn request(&self, text: &str) -> Value { + /// The experience request appending this envelope as `encoded`, with a + /// fresh body idempotency key. + pub(crate) fn request(&self, encoded: &Encoded) -> Value { let (modality, role) = match (&self.kind, &self.turn) { (ItemKind::Conversation, Some(turn)) => ("conversation", role_of(turn.role)), (ItemKind::Document, _) => ("document", "user"), @@ -385,21 +500,51 @@ impl Envelope { .and_then(|turn| turn.at) .or(self.meta.observed_at); let mut context = serde_json::Map::new(); - context.insert( - "labels".to_string(), - json!(labels::for_item(&self.id, &self.meta)), - ); + context.insert("labels".to_string(), json!(encoded.labels)); if let Some(at) = observed_at { context.insert("observed_at".to_string(), json!(at.to_rfc3339())); } - json!({ + let mut request = json!({ "scope": scope_path(&self.meta.namespace, self.kind), "modality": modality, "idempotency_key": crate::cortex::transport::fresh_idempotency_key(), - "content": { "kind": "message", "role": role, "text": text }, + "content": { "kind": "message", "role": role, "text": encoded.text }, "context": Value::Object(context), - }) + }); + // Index and embed, derive nothing (no facts, beliefs or concepts) + // for an item that opts out, and for a tool's output: raw tool + // results are searchable but are not memory about the person. + let tool_turn = self + .turn + .as_ref() + .is_some_and(|turn| turn.role == Role::Tool); + if self.meta.derive == Some(false) || tool_turn { + request["directives"] = json!({ "extract": [] }); + } + request + } +} + +/// `value` as compact JSON. +fn json_of(value: &Envelope) -> Result<String> { + serde_json::to_string(value) + .map_err(|_| Error::Engine("an item envelope could not be serialised".to_string())) +} + +/// `json` cut into numbered part labels of at most [`PART_BYTES`] bytes +/// each, at char boundaries. +fn part_labels(json: &str) -> Vec<String> { + let mut out = Vec::new(); + let mut rest = json; + while !rest.is_empty() { + let mut end = rest.len().min(PART_BYTES); + while !rest.is_char_boundary(end) { + end -= 1; + } + out.push(format!("{PART_PREFIX}{:02}:{}", out.len(), &rest[..end])); + rest = &rest[end..]; } + out } /// CortexDB's four-value message role for a turn's speaker. diff --git a/crates/tinymemory-integrations/src/cortex/envelope/mod_tests.rs b/crates/tinymemory-integrations/src/cortex/envelope/mod_tests.rs index 5837e491..9b8febdd 100644 --- a/crates/tinymemory-integrations/src/cortex/envelope/mod_tests.rs +++ b/crates/tinymemory-integrations/src/cortex/envelope/mod_tests.rs @@ -72,7 +72,7 @@ fn a_conversation_is_one_event_per_turn_in_order() { .map(|e| e.turn.as_ref().map(|t| (t.index, t.count))) .collect(); assert_eq!(turns, vec![Some((0, 2)), Some((1, 2))]); - let request = envelopes[1].request("x"); + let request = envelopes[1].request(&envelopes[1].encode_checked().unwrap()); assert_eq!(request["content"]["role"], "assistant"); assert_eq!(request["scope"], "app:tinymemory/app:conversations"); assert_eq!(request["modality"], "conversation"); @@ -80,13 +80,14 @@ fn a_conversation_is_one_event_per_turn_in_order() { #[test] fn a_recall_rendering_is_read_as_well_as_the_stored_text() { - let envelope = &Envelope::for_item(&conversation(), "id").unwrap()[0]; + let mut envelope = Envelope::for_item(&conversation(), "id").unwrap().remove(0); let stored = envelope.encode().unwrap(); + envelope.v = 2; assert_eq!( - Envelope::decode(&format!("[user] {stored}")).as_ref(), - Some(envelope) + Envelope::decode(&format!("[user] {stored}")), + Some(envelope.clone()) ); - assert_eq!(Envelope::decode(&stored).as_ref(), Some(envelope)); + assert_eq!(Envelope::decode(&stored), Some(envelope)); } #[test] @@ -95,7 +96,7 @@ fn events_this_crate_did_not_write_are_ignored() { assert!(Envelope::decode(r#"{"k":"v1-key","c":"v1 content"}"#).is_none()); let mut old = Envelope::for_item(&conversation(), "id").unwrap().remove(0); old.v = 1; - assert!(Envelope::decode(&old.encode().unwrap()).is_none()); + assert!(Envelope::decode(&json_of(&old).unwrap()).is_none()); assert!(decode_event(&json!({ "id": "e", "content": { "text": "plain" } })).is_none()); } @@ -116,7 +117,7 @@ fn observed_at_and_labels_reach_the_event_context() { meta.observed_at = Some("2026-01-02T03:04:05Z".parse().unwrap()); let item = StoreItem::document("text", meta); let envelope = &Envelope::for_item(&item, "id").unwrap()[0]; - let request = envelope.request("payload"); + let request = envelope.request(&envelope.encode_checked().unwrap()); assert_eq!( request["context"]["observed_at"], "2026-01-02T03:04:05+00:00" @@ -125,7 +126,7 @@ fn observed_at_and_labels_reach_the_event_context() { assert_eq!(request["scope"], "app:tinymemory/app:documents"); assert_ne!( request["idempotency_key"], - envelope.request("payload")["idempotency_key"], + envelope.request(&envelope.encode_checked().unwrap())["idempotency_key"], "every write mints a fresh key" ); } @@ -183,7 +184,7 @@ fn an_item_is_written_to_its_namespace_scope() { let item = StoreItem::document("notes", meta); let id = item.fingerprint(); let envelope = Envelope::for_item(&item, &id).unwrap().remove(0); - let request = envelope.request(&envelope.encode().unwrap()); + let request = envelope.request(&envelope.encode_checked().unwrap()); assert_eq!( request["scope"], "app:tinymemory/agent:researcher/app:documents" @@ -335,3 +336,194 @@ fn a_whole_read_needs_every_piece_of_one_agreed_layout() { "the same item written whole before chunking reads as that body" ); } + +#[test] +fn an_item_that_opts_out_of_derivation_asks_to_extract_nothing() { + let mut opted_out = meta(); + opted_out.derive = Some(false); + let item = StoreItem::document("Run digest: 3 items sent.", opted_out); + let envelope = Envelope::for_item(&item, &item.fingerprint()) + .unwrap() + .remove(0); + let request = envelope.request(&envelope.encode_checked().unwrap()); + assert_eq!(request["directives"], serde_json::json!({ "extract": [] })); + + let plain = StoreItem::document("A note.", meta()); + let envelope = Envelope::for_item(&plain, &plain.fingerprint()) + .unwrap() + .remove(0); + let request = envelope.request(&envelope.encode_checked().unwrap()); + assert!(request.get("directives").is_none(), "{request}"); +} + +/// The event CortexDB would hand back for `request`: its id, content and +/// context as sent. +fn stored(request: &Value) -> Value { + json!({ + "id": "evt_1", + "content": request["content"], + "context": request["context"], + }) +} + +#[test] +fn an_event_is_written_as_its_own_text_with_the_envelope_in_labels() { + let mut meta = meta(); + meta.file_path = Some("/docs/handbook.pdf".into()); + let item = StoreItem::document("Refunds take five days.", meta); + let envelope = Envelope::for_item(&item, &item.fingerprint()) + .unwrap() + .remove(0); + let encoded = envelope.encode_checked().unwrap(); + assert_eq!(encoded.text, "Refunds take five days.", "prose, not JSON"); + let request = envelope.request(&encoded); + let labels: Vec<&str> = request["context"]["labels"] + .as_array() + .unwrap() + .iter() + .map(|label| label.as_str().unwrap()) + .collect(); + assert_eq!(labels[0], labels::item(&item.fingerprint()), "lookup first"); + assert!(labels.contains(&"kind:document"), "{labels:?}"); + assert!(labels.contains(&"file:/docs/handbook.pdf"), "{labels:?}"); + assert!(labels.iter().any(|label| label.starts_with("tm:e:00:"))); + assert!(labels.len() <= 64); + assert!(labels.iter().all(|label| label.len() <= 256), "{labels:?}"); + assert!(!labels.iter().any(|label| label.starts_with("lang:"))); + + let decoded = decode_event(&stored(&request)).unwrap(); + assert_eq!(decoded.envelope, envelope); + assert_eq!(rebuild(&[decoded.envelope]).unwrap(), item); +} + +#[test] +fn every_piece_and_turn_round_trips_through_its_labels() { + for item in [conversation(), long_document()] { + let envelopes = Envelope::for_item(&item, &item.fingerprint()).unwrap(); + let read: Vec<Envelope> = envelopes + .iter() + .map(|envelope| { + let request = envelope.request(&envelope.encode_checked().unwrap()); + assert!( + !request["content"]["text"] + .as_str() + .unwrap() + .starts_with('{') + ); + decode_event(&stored(&request)).unwrap().envelope + }) + .collect(); + assert_eq!(read, envelopes); + assert_eq!(rebuild(&read).unwrap(), item); + } +} + +#[test] +fn a_piece_names_its_pages_and_section_in_readable_labels() { + let mut envelope = Envelope::for_item(&StoreItem::document("piece", meta()), "id") + .unwrap() + .remove(0); + envelope.chunk = Some(ChunkInfo { + index: 1, + count: 3, + pages: Some([3, 5]), + section: Some("Billing".into()), + }); + let labels = envelope.encode_checked().unwrap().labels; + assert!(labels.contains(&"page:3-5".to_string()), "{labels:?}"); + assert!( + labels.contains(&"section:Billing".to_string()), + "{labels:?}" + ); +} + +#[test] +fn an_envelope_too_big_for_its_labels_or_with_no_text_is_written_as_v2() { + let mut big = meta(); + big.tags = (0..2000).map(|n| format!("tag-number-{n:04}")).collect(); + let item = StoreItem::document("Some text.", big); + let envelope = Envelope::for_item(&item, &item.fingerprint()) + .unwrap() + .remove(0); + let encoded = envelope.encode_checked().unwrap(); + assert!( + encoded.text.starts_with('{'), + "the whole envelope: {}", + encoded.text + ); + assert!( + !encoded + .labels + .iter() + .any(|label| label.starts_with("tm:e:")) + ); + let decoded = decode_event(&stored(&envelope.request(&encoded))).unwrap(); + assert_eq!(rebuild(&[decoded.envelope]).unwrap(), item); + + let mut silent = Envelope::for_item(&conversation(), "id").unwrap().remove(0); + silent.text = String::new(); + let encoded = silent.encode_checked().unwrap(); + assert!(encoded.text.starts_with('{'), "never an empty message text"); + assert_eq!( + decode_event(&stored(&silent.request(&encoded))) + .unwrap() + .envelope + .text, + "" + ); +} + +#[test] +fn incomplete_or_foreign_part_labels_are_not_an_envelope() { + let mut tagged = meta(); + tagged.tags = (0..40).map(|n| format!("tag-{n}")).collect(); + let item = StoreItem::document("Some text.", tagged); + let envelope = Envelope::for_item(&item, "id").unwrap().remove(0); + let labels = envelope.encode_checked().unwrap().labels; + let parts: Vec<&str> = labels + .iter() + .map(String::as_str) + .filter(|label| label.starts_with("tm:e:")) + .collect(); + assert!(parts.len() > 1, "{parts:?}"); + let mut shuffled = parts.clone(); + shuffled.reverse(); + assert_eq!( + Envelope::from_labels("t", shuffled).unwrap().text, + "t", + "order is by number" + ); + assert!( + Envelope::from_labels("t", parts[1..].iter().copied()).is_none(), + "a part missing" + ); + assert!( + Envelope::from_labels("t", ["tm:e:00:{\"v\":2}"]).is_none(), + "not v3" + ); + assert!( + Envelope::from_labels("t", ["owner:someone"]).is_none(), + "no parts" + ); +} + +#[test] +fn a_tool_turn_asks_to_extract_nothing_and_other_turns_do_not() { + let item = StoreItem::Conversation { + turns: vec![ + Turn::new(Role::User, "Look up the weather."), + Turn::new(Role::Tool, r#"{"temp_c": 21}"#), + Turn::new(Role::Assistant, "It is 21 degrees."), + ], + meta: meta(), + }; + let extract: Vec<Value> = Envelope::for_item(&item, "id") + .unwrap() + .iter() + .map(|envelope| envelope.request(&envelope.encode_checked().unwrap())["directives"].clone()) + .collect(); + assert_eq!( + extract, + [Value::Null, json!({ "extract": [] }), Value::Null] + ); +} diff --git a/crates/tinymemory-integrations/src/cortex/envelope/rebuild.rs b/crates/tinymemory-integrations/src/cortex/envelope/rebuild.rs index 7ae2e849..08aab583 100644 --- a/crates/tinymemory-integrations/src/cortex/envelope/rebuild.rs +++ b/crates/tinymemory-integrations/src/cortex/envelope/rebuild.rs @@ -16,12 +16,19 @@ pub(crate) struct Decoded { pub(crate) envelope: Envelope, } -/// Decodes one event from a listing or a recall pack. `None` for an event -/// with no id or text, or one this crate did not write. +/// Decodes one event from a listing or a recall pack: a v3 event from its +/// text and labels, else a v2 event from its text. `None` for an event with +/// no id or text, or one this crate did not write. pub(crate) fn decode_event(event: &Value) -> Option<Decoded> { let event_id = event.get("id").and_then(Value::as_str)?.to_string(); let text = event.pointer("/content/text").and_then(Value::as_str)?; - let envelope = Envelope::decode(text)?; + let labels = event + .pointer("/context/labels") + .and_then(Value::as_array) + .into_iter() + .flatten() + .filter_map(Value::as_str); + let envelope = Envelope::from_labels(text, labels).or_else(|| Envelope::decode(text))?; Some(Decoded { event_id, envelope }) } diff --git a/crates/tinymemory-integrations/src/cortex/log/write.rs b/crates/tinymemory-integrations/src/cortex/log/write.rs index 91f6b54d..e2fb8996 100644 --- a/crates/tinymemory-integrations/src/cortex/log/write.rs +++ b/crates/tinymemory-integrations/src/cortex/log/write.rs @@ -188,8 +188,13 @@ impl Log { loop { match self.page(scope, Some(&labels), None, PAGE_SIZE).await { Ok(page) => { + // The text alone does not tell two turns of one item + // apart (two v3 turns can both say "ok"); their labels + // (the envelope parts, with the turn index) do. let found = page.items.iter().find(|event| { event.pointer("/content/text").and_then(Value::as_str) == Some(text) + && event.pointer("/context/labels") + == request.pointer("/context/labels") }); if let Some(event) = found { return receipt(&json!({ "event_id": event.get("id") })); diff --git a/crates/tinymemory-integrations/src/cortex/mod.rs b/crates/tinymemory-integrations/src/cortex/mod.rs index 0be8dd9f..da0e2017 100644 --- a/crates/tinymemory-integrations/src/cortex/mod.rs +++ b/crates/tinymemory-integrations/src/cortex/mod.rs @@ -22,9 +22,11 @@ //! Items live in one scope per kind under the TinyMemory root: //! `app:tinymemory/app:documents`, `app:tinymemory/app:conversations`, //! `app:tinymemory/app:learnings`. A document or learning is one event; a -//! conversation is one event per turn. Each event's text is a JSON envelope -//! (`"v": 2`) carrying the item id ([`tinymemory_api::StoreItem::fingerprint`]), -//! kind, text and full metadata, and each event carries lookup labels (digests +//! conversation is one event per turn. Each event's text is the item's own +//! text, and its labels carry the rest of a v3 envelope: the item id +//! ([`tinymemory_api::StoreItem::fingerprint`]), kind and full metadata +//! (events written before v3, and a few that cannot fit their labels, carry +//! the whole envelope as JSON text, v2). Each event also carries lookup labels (digests //! of the item id and of the exact-match metadata fields) so reads can narrow //! server-side before the full [`tinymemory_api::MetaFilter`] is applied //! client-side. This module's `README.md` summarises the layout and every diff --git a/crates/tinymemory-integrations/src/cortex/testing/log.rs b/crates/tinymemory-integrations/src/cortex/testing/log.rs index 583db306..26341b08 100644 --- a/crates/tinymemory-integrations/src/cortex/testing/log.rs +++ b/crates/tinymemory-integrations/src/cortex/testing/log.rs @@ -233,15 +233,13 @@ impl CortexLog { }) .collect(); scored.sort_by_key(|(score, _)| std::cmp::Reverse(*score)); + // A pack's event carries its stored text: the `[role] ` marker is + // only in `context_block` and the index copy (0.10.3/0.10.4 API + // §9.4), which this double does not render. let events: Vec<Value> = scored .into_iter() .take(budget) - .map(|(_, mut hit)| { - let role = str_of(&hit, "/content/role").to_string(); - let text = str_of(&hit, "/content/text").to_string(); - hit["content"]["text"] = json!(format!("[{role}] {text}")); - hit - }) + .map(|(_, hit)| hit) .collect(); let wanted = body .pointer("/budgets/per_layer_limits/beliefs") diff --git a/docs/architecture/cortex-chunks.md b/docs/architecture/cortex-chunks.md index b7de6bb3..1f1155d6 100644 --- a/docs/architecture/cortex-chunks.md +++ b/docs/architecture/cortex-chunks.md @@ -6,8 +6,10 @@ reads one back. The event layout and labels are in [cortex-flows.md](cortex-flows.md). CortexDB refuses an experience whose flattened text is over 1 MiB -(`422 INVALID_ENVELOPE`, 0.10.4 API §6.10), and a document's event text is -its whole envelope. So (`envelope/chunks.rs`): +(`422 INVALID_ENVELOPE`, 0.10.4 API §6.10). Every limit here is measured on +the event's whole envelope as v2 JSON, which is never shorter than a v3 +event's text, so a piece fits whichever layout it is written in +(`envelope/chunks.rs`): - **When.** Measured as a piece would be written (the envelope with its `chunk` field): a document whose encoded envelope fits in @@ -42,8 +44,9 @@ its whole envelope. So (`envelope/chunks.rs`): page breaks gets no page tag, a piece before the first heading no section tag. These tags are read-side metadata, not part of the item's identity. `pages` is an inclusive range `[first, last]`, counted from 1, so a piece - on one page has `first == last`. Readable CortexDB labels for page and - section are not written yet. + on one page has `first == last`. A v3 piece also carries readable + `page:` and `section:` CortexDB labels (see + [cortex-wire.md](cortex-wire.md)). - **Why one piece per target rather than per page.** CortexDB 0.10.4 already fragments every event over about 500 bytes for retrieval (`matched_fragments`) and serves an over-budget event as an excerpt, and a diff --git a/docs/architecture/cortex-wire.md b/docs/architecture/cortex-wire.md index 2fde8b86..75c4e447 100644 --- a/docs/architecture/cortex-wire.md +++ b/docs/architecture/cortex-wire.md @@ -365,14 +365,38 @@ double enforces). parent-scope sample. - A filter whose `kinds` admits nothing reads no scopes. -## The v2 envelope +## The envelope (v3, and v2) A learning is one event; a conversation is one event per turn, appended in order. A document is one event unless its envelope would pass 256 KiB; then it is one event per piece (see [cortex-chunks.md](cortex-chunks.md)). -CortexDB's experience schema is closed (an unknown field is -a 422), so the structured data rides in the one free-form field: the event's -`content.text` is a JSON **envelope**: +CortexDB's experience schema is closed (an unknown field is a 422), and its +`context.labels` are the app-metadata extension point: stored verbatim, +returned on listings and on recall's `layers.events`, and not counted toward +the 1 MiB text limit. + +**v3 (written now).** An event's `content.text` is the item's own text: the +document body or piece, the turn's text, or the learning's statement. This is +what CortexDB extracts from and splits for search at sentence boundaries, +which a JSON text lacks. Its `context.labels` are, in order: + +- the lookup labels (below), `tm:i:` first; +- readable labels, each at most 256 bytes or left out: `kind:<kind>`, + `file:<path>`, and a piece's `page:<n>` or `page:<first>-<last>` and + `section:<title>`. Never filtered on; no `lang:` label is ever written + (CortexDB reserves it); +- the **envelope parts**: the envelope below with `"v": 3` and an empty + `text`, as compact JSON cut at char boundaries into slices of at most 240 + bytes, each written `tm:e:<NN>:<slice>` with `NN` counting from `00`. + +An event whose text is empty (CortexDB requires message text), or whose +labels would be more than 64, is written as v2 instead. + +**v2 (every event before v3).** `content.text` is the whole envelope as JSON, +with `"v": 2`. Both layouts can live in one scope, and every reader takes +either: an event whose labels hold envelope parts numbered `0..n` that join to +a v3 envelope is v3, with its text as the envelope's `text`; otherwise its +text is tried as a v2 envelope. The envelope: ```json { "v": 2, "id": "<40-hex fingerprint>", "kind": "conversation", @@ -385,7 +409,7 @@ a 422), so the structured data rides in the one free-form field: the event's | Field | Present on | Meaning | | --- | --- | --- | -| `v` | every event | always `2`; any other value is ignored | +| `v` | every event | `3` in envelope parts, `2` in a v2 text; any other value is ignored | | `id` | every event | the item id, `StoreItem::fingerprint()` (a content digest) | | `kind` | every event | `document`, `conversation` or `learning` | | `text` | every event | body, the turn's text, or the learning's statement | @@ -395,10 +419,10 @@ a 422), so the structured data rides in the one free-form field: the event's | `turn` | conversation turns | `index` (0-based), `count`, `role`, `at`, `tool_calls` | | `chunk` | pieces of a chunked document only (a document written whole has none) | `index` (0-based), `count`; optional `pages` (`[first, last]`, only when the text marks pages) and `section` (only when the piece starts under a heading) | -Text that is not a v2 envelope is someone else's event and is ignored by -every reader. Decoding first tries the text as written (`/events` returns it -as stored), then, failing that, strips a `[role] ` prefix (`/recall` renders -text for a reader). A document whose body is still an unresolved URI is +An event that is neither is someone else's and is ignored by every reader. A +v2 text is first tried as written, then without a leading `[role] ` (an older +recall rendered one; 0.10.3 and 0.10.4 return the stored text in +`layers.events`, the marker only in `context_block`). A document whose body is still an unresolved URI is refused at write time as `Error::InvalidRequest`. Rebuilding an item from envelopes: a learning, or a document written whole, diff --git a/docs/architecture/cortex.md b/docs/architecture/cortex.md index a53f910a..cca304a2 100644 --- a/docs/architecture/cortex.md +++ b/docs/architecture/cortex.md @@ -11,7 +11,7 @@ This is the overview. The detail is split into focused pages: - **this page**: surface, credentials, transport, failure mapping, endpoint security, the registry and `MemoryConfig`; - [cortex-wire.md](cortex-wire.md): the two wires, every endpoint and its - request and response shape, the scope layout, the v2 envelope and the lookup + request and response shape, the scope layout, the envelope (v3 and v2) and the lookup labels; - [cortex-flows.md](cortex-flows.md): step-by-step store, list, fetch, recall, forget, get, explore, scope discovery and health; diff --git a/docs/specs/memory-v2.md b/docs/specs/memory-v2.md index 4659d95f..7b1015b2 100644 --- a/docs/specs/memory-v2.md +++ b/docs/specs/memory-v2.md @@ -238,7 +238,8 @@ credentialed cleartext non-loopback endpoint, all as `Error::Config`. - **Wires.** `Direct` (`v1/experience`, `v1/events`, `v1/recall`, `v1/forget`, `v1/answer`) and `TinyHumans` (`memory/*` with `{success,data}` envelopes), as in the v1 adapter. - **Store.** - Each item becomes one experience: a conversation becomes a bulk append of its turns, and a document whose encoded envelope as a chunked piece would pass 256 KiB becomes an ordered append of its pieces, cut at page breaks and headings, and a page or section still over the target cut again at blank lines, then line ends, then characters (`chunk {index, count, pages?, section?}` in the envelope; `pages` an inclusive `[first, last]` range counted from 1, only when the text marks pages; `section` only under a heading). No event over 768 KiB of encoded envelope is sent (CortexDB refuses one over 1 MiB): an item that cannot fit even split (a learning or a turn that long, or a document whose title and metadata leave no room for a piece and which does not fit whole) is `InvalidRequest`, and nothing of its batch is sent. `get`/`list` reassemble a chunked document, returning it only when every piece is present, and a ranked read gives one hit per document, its best-ranked piece, with `page:`/`section:` tags. - - The envelope carries `{v:2, kind, meta, title?, learning_kind?, confidence?}`, and `meta` maps to scope labels where CortexDB can filter. + - An event's `content.text` is the item's own text (body or piece, turn text, statement), and its envelope `{v:3, id, kind, meta, title?, learning_kind?, confidence?, turn?, chunk?}` rides in `tm:e:` labels beside readable `kind:`/`file:`/`page:`/`section:` labels (v3); an event with empty text or too many labels for 64 carries the whole envelope as JSON text (v2), as every event before v3 did, and both read. `meta` maps to lookup labels where CortexDB can filter. + - `MemoryMeta.derive: Some(false)`, and every tool turn, is written with `directives.extract: []`: indexed and searchable, nothing derived. - Writes wait for the indexed barrier, keeping the v1 `await_readable` behaviour. - **Scope.** One scope per item kind *per namespace node*, under the TinyMemory root `app:tinymemory` (which the hosted backend further roots under the tenant): the root node keeps `app:tinymemory/app:{documents,conversations,learnings}`, and a node adds its segments in between, e.g. `app:tinymemory/team:acme/agent:writer/app:learnings`. Namespace segments map to CortexDB's built-in `agent`, `team`, `user`, `ws` and `project` types and the kind leaf uses `app`, because CortexDB v0.10+ refuses scope types outside the deployment's `allowed_scope_types` (`422 UNREGISTERED_SCOPE_TYPE`); every shipped preset allows all of them. A `MetaFilter`'s `kinds` and `reach` pick the scopes read, each admitted kind at each node: an ordinary reach (`Reach::of`, `inherit: true`) reads its node **and every ancestor**, an exact reach its node alone, and a subtree reach its node and every node below; those nodes are known, and only a subtree reach or an unscoped read discovers nodes, from the registered scopes (`v1/scopes/list` / `memory/scopes`). Reads are always exact (`view: "granular"`, sent explicitly because public recall defaults to `holistic`), never server-side traversal. - **Fetch.** From 59e57be39ec6cd83343f22f33589336feb2d74dc Mon Sep 17 00:00:00 2001 From: M3gA-Mind <elvin@tinyhumans.ai> Date: Wed, 7 Oct 2026 01:07:32 +0530 Subject: [PATCH 5/9] Say the v3 envelope parts leave the text to content.text --- crates/tinymemory-integrations/src/cortex/README.md | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/crates/tinymemory-integrations/src/cortex/README.md b/crates/tinymemory-integrations/src/cortex/README.md index 4a4f0445..5bb5768b 100644 --- a/crates/tinymemory-integrations/src/cortex/README.md +++ b/crates/tinymemory-integrations/src/cortex/README.md @@ -124,7 +124,9 @@ document marks its pages and a `section:<title>` tag when the piece starts under a heading; a piece with neither carries no extra tag. Each event's `content.text` is the item's own text (the body or piece, the turn's text, or the statement), and the rest of its envelope rides in `context.labels` as -`tm:e:<NN>:` parts of at most 240 bytes of JSON (v3): +`tm:e:<NN>:` parts of at most 240 bytes of JSON (v3). The parts carry the +envelope with `text` left empty, because the text is the event's +`content.text`; joined, they read: ```json { "v": 3, "id": "<40-hex fingerprint>", "kind": "conversation", "text": "", From 007ffdb4b1aaf9f70256fe8fa4462c9660e918b9 Mon Sep 17 00:00:00 2001 From: M3gA-Mind <elvin@tinyhumans.ai> Date: Wed, 7 Oct 2026 01:13:04 +0530 Subject: [PATCH 6/9] Say how to build MemoryMeta so a field addition does not break callers --- crates/tinymemory-api/src/meta/mod.rs | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/crates/tinymemory-api/src/meta/mod.rs b/crates/tinymemory-api/src/meta/mod.rs index 3a4328f4..a97acd62 100644 --- a/crates/tinymemory-api/src/meta/mod.rs +++ b/crates/tinymemory-api/src/meta/mod.rs @@ -20,6 +20,11 @@ use serde::{Deserialize, Serialize}; use crate::namespace::Namespace; /// Where an item came from and what it is about. +/// +/// Fields are added as the contract grows, so build one with +/// `..MemoryMeta::default()` for the fields you do not set. A literal that +/// names every field stops compiling when a field is added, which is a +/// major release. #[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)] #[serde(default)] pub struct MemoryMeta { From 06f27635357317aef07e5b17ba3cf018cb5b3349 Mon Sep 17 00:00:00 2001 From: M3gA-Mind <elvin@tinyhumans.ai> Date: Wed, 7 Oct 2026 01:31:35 +0530 Subject: [PATCH 7/9] Define how readers tell v3 from v2 and order envelope parts --- docs/architecture/cortex-wire.md | 5 ++++- docs/specs/memory-v2.md | 2 +- 2 files changed, 5 insertions(+), 2 deletions(-) diff --git a/docs/architecture/cortex-wire.md b/docs/architecture/cortex-wire.md index 75c4e447..f8463fdd 100644 --- a/docs/architecture/cortex-wire.md +++ b/docs/architecture/cortex-wire.md @@ -387,7 +387,10 @@ which a JSON text lacks. Its `context.labels` are, in order: (CortexDB reserves it); - the **envelope parts**: the envelope below with `"v": 3` and an empty `text`, as compact JSON cut at char boundaries into slices of at most 240 - bytes, each written `tm:e:<NN>:<slice>` with `NN` counting from `00`. + bytes, each written `tm:e:<NN>:<slice>`. `NN` is the part's number as a + decimal integer, zero-padded to two digits (`00`, `01`, …); an event has + at most 64 labels, so there are at most 63 parts. A reader orders parts by + that number, never by the label's text, and needs every number from 0 up. An event whose text is empty (CortexDB requires message text), or whose labels would be more than 64, is written as v2 instead. diff --git a/docs/specs/memory-v2.md b/docs/specs/memory-v2.md index 7b1015b2..2f7065d5 100644 --- a/docs/specs/memory-v2.md +++ b/docs/specs/memory-v2.md @@ -238,7 +238,7 @@ credentialed cleartext non-loopback endpoint, all as `Error::Config`. - **Wires.** `Direct` (`v1/experience`, `v1/events`, `v1/recall`, `v1/forget`, `v1/answer`) and `TinyHumans` (`memory/*` with `{success,data}` envelopes), as in the v1 adapter. - **Store.** - Each item becomes one experience: a conversation becomes a bulk append of its turns, and a document whose encoded envelope as a chunked piece would pass 256 KiB becomes an ordered append of its pieces, cut at page breaks and headings, and a page or section still over the target cut again at blank lines, then line ends, then characters (`chunk {index, count, pages?, section?}` in the envelope; `pages` an inclusive `[first, last]` range counted from 1, only when the text marks pages; `section` only under a heading). No event over 768 KiB of encoded envelope is sent (CortexDB refuses one over 1 MiB): an item that cannot fit even split (a learning or a turn that long, or a document whose title and metadata leave no room for a piece and which does not fit whole) is `InvalidRequest`, and nothing of its batch is sent. `get`/`list` reassemble a chunked document, returning it only when every piece is present, and a ranked read gives one hit per document, its best-ranked piece, with `page:`/`section:` tags. - - An event's `content.text` is the item's own text (body or piece, turn text, statement), and its envelope `{v:3, id, kind, meta, title?, learning_kind?, confidence?, turn?, chunk?}` rides in `tm:e:` labels beside readable `kind:`/`file:`/`page:`/`section:` labels (v3); an event with empty text or too many labels for 64 carries the whole envelope as JSON text (v2), as every event before v3 did, and both read. `meta` maps to lookup labels where CortexDB can filter. + - An event's `content.text` is the item's own text (body or piece, turn text, statement), and its envelope `{v:3, id, kind, meta, title?, learning_kind?, confidence?, turn?, chunk?}` rides in `tm:e:` labels beside readable `kind:`/`file:`/`page:`/`section:` labels (v3); an event with empty text or too many labels for 64 carries the whole envelope as JSON text (v2), as every event before v3 did, and both read. A reader decides by the labels: an event whose `tm:e:` parts, numbered from 0 with none missing, join to a `v:3` envelope is v3 (its `content.text` is the item's text, empty or not); any other event's `content.text` is tried as a `v:2` envelope; an event that is neither is not this crate's. `meta` maps to lookup labels where CortexDB can filter. - `MemoryMeta.derive: Some(false)`, and every tool turn, is written with `directives.extract: []`: indexed and searchable, nothing derived. - Writes wait for the indexed barrier, keeping the v1 `await_readable` behaviour. - **Scope.** One scope per item kind *per namespace node*, under the TinyMemory root `app:tinymemory` (which the hosted backend further roots under the tenant): the root node keeps `app:tinymemory/app:{documents,conversations,learnings}`, and a node adds its segments in between, e.g. `app:tinymemory/team:acme/agent:writer/app:learnings`. Namespace segments map to CortexDB's built-in `agent`, `team`, `user`, `ws` and `project` types and the kind leaf uses `app`, because CortexDB v0.10+ refuses scope types outside the deployment's `allowed_scope_types` (`422 UNREGISTERED_SCOPE_TYPE`); every shipped preset allows all of them. A `MetaFilter`'s `kinds` and `reach` pick the scopes read, each admitted kind at each node: an ordinary reach (`Reach::of`, `inherit: true`) reads its node **and every ancestor**, an exact reach its node alone, and a subtree reach its node and every node below; those nodes are known, and only a subtree reach or an unscoped read discovers nodes, from the registered scopes (`v1/scopes/list` / `memory/scopes`). Reads are always exact (`view: "granular"`, sent explicitly because public recall defaults to `holistic`), never server-side traversal. From 89ecf9e1f931925e083284a2964679bc7b6bbbca Mon Sep 17 00:00:00 2001 From: M3gA-Mind <elvin@tinyhumans.ai> Date: Wed, 7 Oct 2026 01:09:22 +0530 Subject: [PATCH 8/9] Key each event by its own body; skip the lookup on the turn hot path CortexDB 0.10.4 (measured): a reused body idempotency_key with the same body is a replay answered with the first event's id and replayed_from_idempotency: true; with another body it is a 409; and /v1/forget by memory_ids releases the key. The old reason for fresh keys (a forgotten item's key is never released) no longer holds. - Each event's idempotency_key is tm3: and 56 hex of the SHA-256 of its request without the key (60 chars, under CortexDB's 64). An identical retry replays; any change to the body is a new key, so a 409 for a reused key cannot happen. - Writes read replayed_from_idempotency (single, bulk, hosted); a receipt is a replay when every written event was replayed. - The pre-write lookup is skipped only for a Direct, accepted-only store of one single-turn conversation (the agent lifecycle's two writes per turn). Keys last 24 hours and change with observed_at, so documents, batches, waited-for writes and every hosted write still look up. - The hosted Idempotency-Key claim stays fresh per write: the hosted API answers every replay of a claim with 409. - The test double releases a key on forget, as 0.10.4 does. --- .../src/cortex/README.md | 6 +- .../src/cortex/engine/mod_chunk_tests.rs | 10 +-- .../src/cortex/engine/mod_hosted_tests.rs | 12 ++-- .../src/cortex/engine/store.rs | 34 ++++++++-- .../src/cortex/engine/store_tests.rs | 63 +++++++++++++++++++ .../src/cortex/envelope/mod.rs | 8 ++- .../src/cortex/envelope/mod_tests.rs | 13 +++- .../src/cortex/log/mod.rs | 7 ++- .../src/cortex/log/write.rs | 43 +++++++++---- .../src/cortex/testing/log.rs | 24 ++++++- .../src/cortex/transport/mod.rs | 38 ++++++++--- docs/architecture/cortex-flows.md | 31 ++++++--- docs/architecture/cortex-wire.md | 14 +++-- docs/architecture/testing.md | 6 +- 14 files changed, 238 insertions(+), 71 deletions(-) diff --git a/crates/tinymemory-integrations/src/cortex/README.md b/crates/tinymemory-integrations/src/cortex/README.md index 5bb5768b..90961af6 100644 --- a/crates/tinymemory-integrations/src/cortex/README.md +++ b/crates/tinymemory-integrations/src/cortex/README.md @@ -214,8 +214,10 @@ as prefixes, so they cannot be labelled and are filtered only client-side. These were measured against a live CortexDB by the v1 adapter. The doubles in `testing/` reproduce all of them. -- **Append-only.** There is no update route. Forget removes events but not - their idempotency records. +- **Append-only.** There is no update route. A body `idempotency_key` replays + the same body for 24 hours and refuses another body (409); forget by + `memory_ids` releases it. Keys are derived from the body, so a retry is a + replay. - **Accepted is not readable.** A write first waits until the label-narrowed listing carries its event (fatal after 30s). It then waits until ranked recall returns it (best-effort, 10s); a recall that is down or slow does not diff --git a/crates/tinymemory-integrations/src/cortex/engine/mod_chunk_tests.rs b/crates/tinymemory-integrations/src/cortex/engine/mod_chunk_tests.rs index b6b047de..c0b54761 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/mod_chunk_tests.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/mod_chunk_tests.rs @@ -156,15 +156,7 @@ async fn a_store_that_lost_a_piece_writes_only_that_piece_again() { let item = handbook(30); engine.store(item.clone()).await.unwrap(); let before = events(&state, SCOPE).len(); - { - let mut log = state.log.lock().unwrap(); - let last = log - .events - .iter() - .rposition(|event| event["scope"] == SCOPE) - .unwrap(); - log.events.remove(last); - } + state.log.lock().unwrap().lose_last(SCOPE); let got = engine .get(GetRequest { ids: vec![ItemId::new(item.fingerprint())], diff --git a/crates/tinymemory-integrations/src/cortex/engine/mod_hosted_tests.rs b/crates/tinymemory-integrations/src/cortex/engine/mod_hosted_tests.rs index 358a2654..d1005a33 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/mod_hosted_tests.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/mod_hosted_tests.rs @@ -409,12 +409,12 @@ async fn recovery_does_not_take_another_turn_with_the_same_words_for_the_lost_on meta: MemoryMeta::default(), }; engine.store(item.clone()).await.unwrap(); - { - // The second turn's event is lost. - let mut log = state.log.lock().unwrap(); - let last = log.events.len() - 1; - log.events.remove(last); - } + // The second turn's event is lost. + state + .log + .lock() + .unwrap() + .lose_last("app:tinymemory/app:conversations"); // Its re-write is claimed but never applied, so recovery must look for // it, and must not take the first turn, which says the same words. state.claim_then_fail.store(1, Ordering::SeqCst); diff --git a/crates/tinymemory-integrations/src/cortex/engine/store.rs b/crates/tinymemory-integrations/src/cortex/engine/store.rs index 52ae5534..b88b89d9 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/store.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/store.rs @@ -15,10 +15,15 @@ //! are written, in order; //! - nothing present — every event is written. //! -//! Writes use fresh idempotency keys (see `transport::fresh_idempotency_key`) -//! rather than ones derived from content: CortexDB never releases a key on -//! forget, so a content key would make re-storing a forgotten item a silent -//! no-op. +//! Each event is keyed by its own body (`transport::body_idempotency_key`), +//! so an identical retry is a replay CortexDB answers without writing +//! (`replayed_from_idempotency`), and forgetting an item releases its keys. +//! A key lasts 24 hours and changes with any byte of the body, so it does +//! not replace the lookup: an unchanged file synced again a day later, or +//! with a new `observed_at`, would be written again without it. The lookup +//! is skipped only on the turn-logging hot path ([`looks_up`]): a Direct, +//! accepted-only store of one single-turn conversation, which the agent +//! lifecycle makes twice per turn, and where a retry is the replay to catch. use std::collections::{BTreeMap, HashMap, HashSet}; @@ -26,6 +31,7 @@ use tinymemory_api::{ItemId, StoreItem, StoreReceipt, WaitFor, validate_many}; use super::CortexEngine; use super::scopes::KindScope; +use crate::cortex::descriptor::CortexWire; use crate::cortex::envelope::{Encoded, Envelope}; use crate::cortex::error::Result; use crate::cortex::log::Written; @@ -71,7 +77,8 @@ impl CortexEngine { .or_default() .push(id.clone()); } - for (scope, of_scope) in &by_scope { + let lookup = looks_up(self.wire(), &items, wait); + for (scope, of_scope) in by_scope.iter().filter(|_| lookup) { for (id, events) in self.item_events(scope, of_scope).await? { held.entry(id) .or_default() @@ -92,8 +99,9 @@ impl CortexEngine { requests.push(envelope.request(&encoded)); } } - let replayed = requests.is_empty(); + let mut replayed = requests.is_empty(); if let Some(written) = self.log.write(&requests, wait).await? { + replayed = written.replayed; last_per_scope.retain(|w| w.scope != written.scope); last_per_scope.push(written); } @@ -116,6 +124,20 @@ impl CortexEngine { } } +/// Whether a store looks its items up before writing. Always, except on the +/// turn-logging hot path: a Direct store of one single-turn conversation +/// that waits only for acceptance. There the body key catches a retry +/// (Direct answers it as a replay), and a listing per turn would cost the +/// turn latency. The hosted wire always looks up: its answer need not say +/// it replayed. +fn looks_up(wire: CortexWire, items: &[StoreItem], wait: WaitFor) -> bool { + let hot = matches!( + items, + [StoreItem::Conversation { turns, .. }] if turns.len() == 1 + ); + !(wire == CortexWire::Direct && wait == WaitFor::Accepted && hot) +} + #[cfg(test)] #[path = "store_tests.rs"] mod tests; diff --git a/crates/tinymemory-integrations/src/cortex/engine/store_tests.rs b/crates/tinymemory-integrations/src/cortex/engine/store_tests.rs index dbd23978..d0db8e5a 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/store_tests.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/store_tests.rs @@ -141,3 +141,66 @@ async fn a_single_store_is_listed_and_settled_on_return_like_a_batch_of_one() { ); } } + +fn turn(text: &str) -> StoreItem { + StoreItem::Conversation { + turns: vec![tinymemory_api::Turn::new(tinymemory_api::Role::User, text)], + meta: MemoryMeta { + thread_id: Some("t-hot".into()), + ..MemoryMeta::default() + }, + } +} + +#[tokio::test] +async fn a_logged_turn_skips_the_lookup_and_a_retry_is_a_replay() { + use tinymemory_api::WriteOptions; + let (endpoint, state) = crate::cortex::testing::direct_double().await; + let engine = crate::cortex::testing::direct_engine(&endpoint); + let first = engine + .store_with(turn("Ship it on Friday."), WriteOptions::accepted()) + .await + .unwrap(); + assert!(!first.replayed); + assert_eq!( + state.count("GET /v1/events"), + 0, + "no lookup on the hot path" + ); + let retry = engine + .store_with(turn("Ship it on Friday."), WriteOptions::accepted()) + .await + .unwrap(); + assert!(retry.replayed, "CortexDB answered the retry as a replay"); + assert_eq!(retry.id, first.id); + assert_eq!(state.event_count(), 1, "written once"); + assert_eq!(state.count("GET /v1/events"), 0); +} + +#[tokio::test] +async fn every_other_store_still_looks_its_items_up_first() { + use tinymemory_api::WriteOptions; + for (engine, state) in both().await { + let listing = match engine.wire() { + crate::cortex::CortexWire::Direct => "GET /v1/events", + crate::cortex::CortexWire::TinyHumans => "GET /memory/events", + }; + let cases: Vec<(StoreItem, WriteOptions)> = vec![ + (doc("A synced file."), WriteOptions::accepted()), + (turn("Waited for."), WriteOptions::visible()), + ]; + for (item, options) in cases { + let before = state.count(listing); + engine.store_with(item.clone(), options).await.unwrap(); + assert!(state.count(listing) > before, "{:?} looked up", item.kind()); + } + if engine.wire() == crate::cortex::CortexWire::TinyHumans { + let before = state.count(listing); + engine + .store_with(turn("Hosted turn."), WriteOptions::accepted()) + .await + .unwrap(); + assert!(state.count(listing) > before, "hosted always looks up"); + } + } +} diff --git a/crates/tinymemory-integrations/src/cortex/envelope/mod.rs b/crates/tinymemory-integrations/src/cortex/envelope/mod.rs index 9167bb7c..72850811 100644 --- a/crates/tinymemory-integrations/src/cortex/envelope/mod.rs +++ b/crates/tinymemory-integrations/src/cortex/envelope/mod.rs @@ -486,8 +486,9 @@ impl Envelope { Ok(probe.encode()?.len() + SECTION_RESERVE) } - /// The experience request appending this envelope as `encoded`, with a - /// fresh body idempotency key. + /// The experience request appending this envelope as `encoded`, keyed by + /// its own body (`transport::body_idempotency_key`), so an identical + /// retry is a replay. pub(crate) fn request(&self, encoded: &Encoded) -> Value { let (modality, role) = match (&self.kind, &self.turn) { (ItemKind::Conversation, Some(turn)) => ("conversation", role_of(turn.role)), @@ -507,7 +508,6 @@ impl Envelope { let mut request = json!({ "scope": scope_path(&self.meta.namespace, self.kind), "modality": modality, - "idempotency_key": crate::cortex::transport::fresh_idempotency_key(), "content": { "kind": "message", "role": role, "text": encoded.text }, "context": Value::Object(context), }); @@ -521,6 +521,8 @@ impl Envelope { if self.meta.derive == Some(false) || tool_turn { request["directives"] = json!({ "extract": [] }); } + request["idempotency_key"] = + json!(crate::cortex::transport::body_idempotency_key(&request)); request } } diff --git a/crates/tinymemory-integrations/src/cortex/envelope/mod_tests.rs b/crates/tinymemory-integrations/src/cortex/envelope/mod_tests.rs index 9b8febdd..6c822034 100644 --- a/crates/tinymemory-integrations/src/cortex/envelope/mod_tests.rs +++ b/crates/tinymemory-integrations/src/cortex/envelope/mod_tests.rs @@ -124,10 +124,19 @@ fn observed_at_and_labels_reach_the_event_context() { ); assert_eq!(request["context"]["labels"][0], labels::item("id")); assert_eq!(request["scope"], "app:tinymemory/app:documents"); - assert_ne!( + let key = request["idempotency_key"].as_str().unwrap(); + assert!(key.starts_with("tm3:") && key.len() <= 64, "{key}"); + assert_eq!( request["idempotency_key"], envelope.request(&envelope.encode_checked().unwrap())["idempotency_key"], - "every write mints a fresh key" + "the same body, the same key: a retry is a replay" + ); + let mut later = envelope.clone(); + later.meta.observed_at = Some("2026-01-03T00:00:00Z".parse().unwrap()); + assert_ne!( + request["idempotency_key"], + later.request(&later.encode_checked().unwrap())["idempotency_key"], + "any change to the body is a new key, never a reused key's 409" ); } diff --git a/crates/tinymemory-integrations/src/cortex/log/mod.rs b/crates/tinymemory-integrations/src/cortex/log/mod.rs index ff23b451..717558ad 100644 --- a/crates/tinymemory-integrations/src/cortex/log/mod.rs +++ b/crates/tinymemory-integrations/src/cortex/log/mod.rs @@ -7,9 +7,10 @@ //! was wrong in that adapter first. The test doubles reproduce all of them. //! //! - **It is append-only.** `/v1/experience` only appends; there is no -//! update. `/v1/forget` removes events but **not** their idempotency -//! records, so a body `idempotency_key` reused after a forget is swallowed -//! as a replay. Every write therefore gets a fresh key. +//! update. A body `idempotency_key` replays the same body for 24 hours +//! (`replayed_from_idempotency`) and refuses another body (409); forget by +//! `memory_ids` releases it (0.10.4). Keys are derived from the body, so a +//! retry is a replay (`transport::body_idempotency_key`). //! - **Accepted is not readable.** `/v1/experience` answers `202 captured` //! and indexes afterwards. Writes wait (see `visibility`): first until the //! listing carries the event, which is fatal on timeout, then until ranked diff --git a/crates/tinymemory-integrations/src/cortex/log/write.rs b/crates/tinymemory-integrations/src/cortex/log/write.rs index e2fb8996..b65dabca 100644 --- a/crates/tinymemory-integrations/src/cortex/log/write.rs +++ b/crates/tinymemory-integrations/src/cortex/log/write.rs @@ -43,6 +43,9 @@ pub(crate) struct Written { pub(crate) label: String, pub(crate) text: String, pub(crate) event_id: String, + /// Whether CortexDB answered every event of the write as a replay of + /// one it already holds (`replayed_from_idempotency`). + pub(crate) replayed: bool, } /// How many times a hosted write is sent before a transient fault surfaces. @@ -61,6 +64,16 @@ fn parts(request: &Value) -> Result<(&str, &str, &str)> { } } +/// Whether a write receipt says CortexDB replayed the event rather than +/// writing it. Absent (an older server, or a hosted answer without it) is +/// not a replay. +fn replayed(answer: &Value) -> bool { + answer + .get("replayed_from_idempotency") + .and_then(Value::as_bool) + == Some(true) +} + /// The event id a write receipt names. fn receipt(answer: &Value) -> Result<String> { answer @@ -83,14 +96,15 @@ impl Log { let Some(last) = requests.last() else { return Ok(None); }; - let event_id = match self.client.wire() { + let (event_id, replayed) = match self.client.wire() { CortexWire::Direct => self.append_direct(requests, wait).await?, CortexWire::TinyHumans => { - let mut event_id = String::new(); + let mut last = (String::new(), true); for request in requests { - event_id = self.send_hosted_write(request).await?; + let (event_id, replayed) = self.send_hosted_write(request).await?; + last = (event_id, last.1 && replayed); } - event_id + last } }; let (scope, text, label) = parts(last)?; @@ -99,6 +113,7 @@ impl Log { label: label.to_string(), text: text.to_string(), event_id, + replayed, })) } @@ -119,8 +134,8 @@ impl Log { } /// One Direct write of one event or one ordered batch; the last event's - /// id. - async fn append_direct(&self, requests: &[Value], wait: WaitFor) -> Result<String> { + /// id, and whether every event was a replay. + async fn append_direct(&self, requests: &[Value], wait: WaitFor) -> Result<(String, bool)> { let wire = self.client.wire(); let query = match wait { WaitFor::Visible => "?wait=indexed", @@ -132,7 +147,7 @@ impl Log { .client .json(reqwest::Method::POST, &path, Some(single), Attempts::Once) .await?; - return receipt(&answer); + return Ok((receipt(&answer)?, replayed(&answer))); } let path = format!("{}{query}", wire.path(Route::Bulk)); let body = json!({ "items": requests, "ordering": "strict_temporal" }); @@ -151,23 +166,25 @@ impl Log { requests.len() ))); } - results.last().map_or_else( + let last = results.last().map_or_else( || Err(Error::Engine("empty bulk results".to_string())), receipt, - ) + )?; + Ok((last, results.iter().all(replayed))) } - /// One hosted write under one claim, with the outcome-unknown recovery. - async fn send_hosted_write(&self, request: &Value) -> Result<String> { + /// One hosted write under one claim, with the outcome-unknown recovery; + /// the event id, and whether it was a replay. + async fn send_hosted_write(&self, request: &Value) -> Result<(String, bool)> { let path = self.client.wire().path(Route::Experience); let claim = fresh_idempotency_key(); let mut attempt = 0; loop { attempt += 1; match self.client.json_keyed(path, request, &claim).await { - Ok(answer) => return receipt(&answer), + Ok(answer) => return Ok((receipt(&answer)?, replayed(&answer))), Err(Error::Conflict(_)) if attempt > 1 => { - return self.recover_unknown_write(request).await; + return Ok((self.recover_unknown_write(request).await?, false)); } Err(error) if attempt < HOSTED_WRITE_ATTEMPTS && error.is_transient() => { tokio::time::sleep(self.timing.poll * 2_u32.pow(attempt - 1)).await; diff --git a/crates/tinymemory-integrations/src/cortex/testing/log.rs b/crates/tinymemory-integrations/src/cortex/testing/log.rs index 26341b08..e6d7c5f0 100644 --- a/crates/tinymemory-integrations/src/cortex/testing/log.rs +++ b/crates/tinymemory-integrations/src/cortex/testing/log.rs @@ -2,9 +2,10 @@ //! //! Deliberately unaccommodating, because a tidy double proves nothing: //! -//! - append-only, with the body `idempotency_key` remembered for ever (a -//! reused key with a different body is `409 IDEMPOTENCY_CONFLICT`, the -//! same body is a replay, and forgetting an event does not release it); +//! - append-only, with the body `idempotency_key` remembered (a reused key +//! with a different body is `409 IDEMPOTENCY_CONFLICT`, the same body is +//! a replay); forgetting an event by `memory_ids` releases its key, as +//! CortexDB 0.10.4 does; //! - the listing is newest first and emits **every event twice**, with //! `limit` counting the copies; //! - unknown query parameters are ignored; @@ -185,6 +186,8 @@ impl CortexLog { }); self.forgotten .extend(gone.iter().map(|e| str_of(e, "/id").to_string())); + self.idempotency + .retain(|_, (_, id)| !gone.iter().any(|e| str_of(e, "/id") == id)); self.events = kept; } else { self.events.retain(|e| str_of(e, "/scope") != scope); @@ -196,6 +199,21 @@ impl CortexLog { ) } + /// Removes the last event of `scope` and its idempotency record, as if + /// it had never been written (a store that failed part way). + pub(crate) fn lose_last(&mut self, scope: &str) { + let Some(at) = self + .events + .iter() + .rposition(|e| str_of(e, "/scope") == scope) + else { + return; + }; + let lost = self.events.remove(at); + let id = str_of(&lost, "/id").to_string(); + self.idempotency.retain(|_, (_, held)| *held != id); + } + /// `POST /v1/recall`: events ranked by how many query words they hold. pub(crate) fn recall(&self, body: &Value) -> Value { let scope = str_of(body, "/scope"); diff --git a/crates/tinymemory-integrations/src/cortex/transport/mod.rs b/crates/tinymemory-integrations/src/cortex/transport/mod.rs index 2e3f0770..5aa31b7f 100644 --- a/crates/tinymemory-integrations/src/cortex/transport/mod.rs +++ b/crates/tinymemory-integrations/src/cortex/transport/mod.rs @@ -388,15 +388,37 @@ pub(crate) fn credential_header(token: &str) -> Result<HeaderValue> { Ok(header) } -/// A fresh key for every write: the body's `idempotency_key` and the hosted -/// `Idempotency-Key` claim. +/// The body `idempotency_key` of an experience request: `tm3:` and the +/// first 56 hex digits of the SHA-256 of the request without its key, 60 +/// characters (CortexDB refuses one over 64). /// -/// Never derived from content. CortexDB keeps a forgotten event's -/// idempotency record, so a content-derived key would make re-storing an item -/// after forgetting it a silent no-op; and the hosted memory API answers -/// every replay of a claim with 409, so a content-derived claim would refuse -/// an identical re-store. Store detects a replay itself, by looking the item -/// up, before it writes. +/// Derived from the whole body, so an identical retry replays: CortexDB +/// 0.10.4 answers it with the first event's id and +/// `replayed_from_idempotency: true`, and writes nothing, for 24 hours. +/// Any change to the body (a new `observed_at` on an unchanged item) is a +/// new key and a new event, never a 409 for a reused key with another body. +/// `/v1/forget` by `memory_ids` releases a key (measured on 0.10.4), so an +/// item stored again after it was forgotten is written again. `tm3` names +/// the event layout the body is in. +pub(crate) fn body_idempotency_key(request: &Value) -> String { + use sha2::{Digest, Sha256}; + let mut body = request.clone(); + if let Some(object) = body.as_object_mut() { + object.remove("idempotency_key"); + } + let bytes = serde_json::to_vec(&body).unwrap_or_default(); + let mut key = String::from("tm3:"); + for byte in Sha256::digest(&bytes).iter().take(28) { + key.push_str(&format!("{byte:02x}")); + } + key +} + +/// A fresh key for every hosted write's `Idempotency-Key` claim. +/// +/// Never derived from content: the hosted memory API answers every replay +/// of a claim with 409, so a content-derived claim would refuse an identical +/// re-store. One claim is reused across the retries of one write. /// /// Three parts: a per-process salt from the OS-seeded `RandomState`, the /// wall-clock nanoseconds, and a counter, so neither two writes in one diff --git a/docs/architecture/cortex-flows.md b/docs/architecture/cortex-flows.md index f7713e21..361c7016 100644 --- a/docs/architecture/cortex-flows.md +++ b/docs/architecture/cortex-flows.md @@ -24,7 +24,10 @@ valid items; an empty or oversized batch is `Error::InvalidRequest`). Then: the same text at two nodes is two items and a re-sync that only restamps `observed_at` is a replay. 2. **Group** the items by scope (kind at namespace node). -3. **Replay detection.** One id lookup per scope: the listing narrowed by the +3. **Replay detection.** Skipped only on the turn-logging hot path: a Direct + store of one single-turn conversation with `WaitFor::Accepted` (the agent + lifecycle makes two per turn), where the body key catches a retry (below). + Otherwise one id lookup per scope: the listing narrowed by the items' `tm:i:` labels (batches of up to 50 labels), each hit re-checked against the envelope's real id. The result is, per id, which turn indexes are already held (`None` for a document or learning). @@ -38,7 +41,9 @@ valid items; an empty or oversized batch is `Error::InvalidRequest`). Then: - nothing present: every event is written. An item repeated inside the batch is a replay of its first copy. Each - write uses a fresh idempotency key (see below). The wire call is made per + event is keyed by its own body (see below), and a receipt is also a + replay when CortexDB answered every written event as one + (`replayed_from_idempotency`). The wire call is made per item: Direct sends one experience, or one ordered bulk when two or more events are due; TinyHumans sends the events one at a time. 5. **Wait, once per scope.** For the **last event written** in each scope the @@ -50,11 +55,23 @@ valid items; an empty or oversized batch is `Error::InvalidRequest`). Then: Receipts come back in item order, each `{id, replayed}`. On an error the items before the failing one are stored, and storing them again is a replay. -**Fresh idempotency keys.** Writes use a fresh `tm-<salt>-<nanos>-<seq>` key, -never one derived from content. CortexDB never releases a key on forget, so a -content key would make re-storing a forgotten item a silent no-op; and the -hosted memory API answers every replay of a claim with 409. Replay detection is -done by the engine, by looking the item up, before it writes. +**Body idempotency keys.** Each event's `idempotency_key` is `tm3:` and 56 hex +digits of the SHA-256 of its request body without the key (60 characters; +CortexDB refuses one over 64). Measured on CortexDB 0.10.4: + +- the same key and body is a replay, answered with the first event's id and + `replayed_from_idempotency: true`, and nothing is written; +- the same key with another body is `409 IDEMPOTENCY_CONFLICT`, which a body + key never produces (any change to the body is a new key); +- `/v1/forget` by `memory_ids` releases the key, so an item stored again + after it was forgotten is written again. + +A key lasts 24 hours, and an unchanged item restamped with a new +`observed_at` has a new body, so the key does not replace the lookup. It +catches retries, and on the hot path it is the only replay check. The hosted +memory API's `Idempotency-Key` **claim** stays fresh per write +(`tm-<salt>-<nanos>-<seq>`), reused only across the retries of that write: it +answers every replay of a claim with 409. The hosted wire always looks up. ### Waiting for a write to be readable diff --git a/docs/architecture/cortex-wire.md b/docs/architecture/cortex-wire.md index f8463fdd..e1dc331e 100644 --- a/docs/architecture/cortex-wire.md +++ b/docs/architecture/cortex-wire.md @@ -63,7 +63,7 @@ Request body (one event). It is the same on both wires: { "scope": "app:tinymemory/agent:researcher/app:documents", "modality": "document", - "idempotency_key": "tm-<salt>-<nanos>-<seq>", + "idempotency_key": "tm3:<56 hex: SHA-256 of this body without the key>", "content": { "kind": "message", "role": "user", "text": "<envelope JSON, see below>" }, "context": { "labels": ["tm:i:<16 hex>", "tm:k:<16 hex>"], @@ -76,8 +76,9 @@ Request body (one event). It is the same on both wires: `conversation` for a turn. `content.role` is `user` for documents and learnings and the turn's speaker (`user`, `assistant`, `system`, `tool`) for a conversation turn. -- `idempotency_key` is a fresh value on every write, never derived from - content (see [flows: store](cortex-flows.md#store-and-store_many)). +- `idempotency_key` is derived from the body, so an identical retry is a + replay (see [flows: store](cortex-flows.md#store-and-store_many)). The + answer's `replayed_from_idempotency` is read; absent counts as `false`. - `context.observed_at` is the turn's `at`, else the item's `meta.observed_at`; it is omitted when neither is set. - `context.labels[0]` is always the item label; the writer relies on that. @@ -478,9 +479,10 @@ digest cannot, so they are never labelled and are filtered client-side only. Each was measured against a live CortexDB and was wrong in the first adapter. The loopback doubles reproduce all of them (see [testing](testing.md)). -- **Append-only.** There is no update route. Forget removes events but **not** - their idempotency records, so a reused body `idempotency_key` after a forget - is swallowed as a replay. +- **Append-only.** There is no update route. A reused body `idempotency_key` + with the same body is a replay for 24 hours, and with another body a 409. + Forget by `memory_ids` releases the keys of what it removes (measured on + 0.10.4; the v1 adapter's notes said the opposite, on an unrecorded build). - **Accepted is not readable.** An append answers `202` and indexes afterwards. The status route and the lifecycle stream are not readiness signals, so the engine waits on the listing and on recall itself (see diff --git a/docs/architecture/testing.md b/docs/architecture/testing.md index b9a13b29..07340215 100644 --- a/docs/architecture/testing.md +++ b/docs/architecture/testing.md @@ -196,9 +196,9 @@ Both doubles serve the same in-memory `CortexLog`, which is deliberately **unaccommodating**, because a tidy double proves nothing. It reproduces every behaviour in [the wire page](cortex-wire.md#cortexdb-behaviours-the-engine-is-shaped-around): -- append-only, a body `idempotency_key` remembered for ever (same key, same - body is a replay; same key, different body is `409 IDEMPOTENCY_CONFLICT`; - forget does not release it); +- append-only, a body `idempotency_key` remembered (same key, same body is a + replay; same key, different body is `409 IDEMPOTENCY_CONFLICT`); forget by + `memory_ids` releases it, as CortexDB 0.10.4 does; - the listing is newest first, emits **every event twice**, counts the copies in `limit`, ignores unknown query parameters, and pages by offset cursor; - the forget selector reads only `memory_ids`; an empty selector without From 6e1e0f9628543f2ac9a62f09e6dfee2dc668a096 Mon Sep 17 00:00:00 2001 From: M3gA-Mind <elvin@tinyhumans.ai> Date: Wed, 7 Oct 2026 01:35:14 +0530 Subject: [PATCH 9/9] Live-test the turn replay; pin down the key and forget wording - live_cortexdb: a turn logged twice on the hot path is written once, the retry answered by CortexDB as a replay of the first event. - The key is the first 56 hex digits (224 bits) of the SHA-256 of the compact JSON body; a body-derived key never produces the 409 for a reused key. - The double's scope-wide forget keeps keys held, as CortexDB's redact-only scope forget does; a forget by memory_ids releases them. --- .../src/cortex/testing/log.rs | 3 ++ .../tests/live_cortexdb.rs | 48 +++++++++++++++++++ docs/architecture/cortex-flows.md | 5 +- docs/architecture/cortex-wire.md | 5 +- 4 files changed, 57 insertions(+), 4 deletions(-) diff --git a/crates/tinymemory-integrations/src/cortex/testing/log.rs b/crates/tinymemory-integrations/src/cortex/testing/log.rs index e6d7c5f0..87dbc235 100644 --- a/crates/tinymemory-integrations/src/cortex/testing/log.rs +++ b/crates/tinymemory-integrations/src/cortex/testing/log.rs @@ -190,6 +190,9 @@ impl CortexLog { .retain(|_, (_, id)| !gone.iter().any(|e| str_of(e, "/id") == id)); self.events = kept; } else { + // A scope-wide forget only redacts, so its keys stay held for + // their 24 hours (CortexDB's answer for 0.10.3/0.10.4); only a + // forget by `memory_ids` releases them. self.events.retain(|e| str_of(e, "/scope") != scope); } let deleted = before - self.events.len(); diff --git a/crates/tinymemory-integrations/tests/live_cortexdb.rs b/crates/tinymemory-integrations/tests/live_cortexdb.rs index e3fe9f20..59b0663c 100644 --- a/crates/tinymemory-integrations/tests/live_cortexdb.rs +++ b/crates/tinymemory-integrations/tests/live_cortexdb.rs @@ -391,3 +391,51 @@ async fn a_long_document_round_trips_in_pieces() { } } } + +/// A turn logged on the hot path (one single-turn conversation, accepted +/// only) skips the lookup, so a retry is caught by CortexDB itself: the +/// same body is the same idempotency key, answered as a replay of the +/// first event. The hosted wire always looks up, so this is Direct only. +#[tokio::test] +async fn a_logged_turn_sent_twice_is_written_once() { + let _alone = ONE_AT_A_TIME.lock().await; + for (wire, engine) in live_engines() { + if wire != "cortexdb" { + continue; + } + let thread = run_id(); + let turn = StoreItem::Conversation { + turns: vec![Turn::new( + Role::User, + format!("Ship the Aurora build on Friday ({thread})."), + )], + meta: MemoryMeta { + thread_id: Some(thread.clone()), + ..MemoryMeta::default() + }, + }; + let first = engine + .store_with(turn.clone(), tinymemory_api::WriteOptions::accepted()) + .await + .expect("log the turn"); + assert!(!first.replayed); + let retry = engine + .store_with(turn, tinymemory_api::WriteOptions::accepted()) + .await + .expect("log it again"); + assert!(retry.replayed, "CortexDB answered the retry as a replay"); + assert_eq!(retry.id, first.id); + + let filter = MetaFilter { + thread_id: Some(thread), + ..MetaFilter::default() + }; + let listed = list_until(&engine, &filter, 1).await; + assert_eq!(listed.len(), 1, "one conversation"); + let report = engine + .forget(ForgetTarget::Ids(vec![first.id])) + .await + .expect("forget"); + assert_eq!(report.forgotten, 1); + } +} diff --git a/docs/architecture/cortex-flows.md b/docs/architecture/cortex-flows.md index 361c7016..14a22618 100644 --- a/docs/architecture/cortex-flows.md +++ b/docs/architecture/cortex-flows.md @@ -55,8 +55,9 @@ valid items; an empty or oversized batch is `Error::InvalidRequest`). Then: Receipts come back in item order, each `{id, replayed}`. On an error the items before the failing one are stored, and storing them again is a replay. -**Body idempotency keys.** Each event's `idempotency_key` is `tm3:` and 56 hex -digits of the SHA-256 of its request body without the key (60 characters; +**Body idempotency keys.** Each event's `idempotency_key` is `tm3:` and the +first 56 hex digits (224 bits) of the SHA-256 of its compact JSON request +body without the key (60 characters; CortexDB refuses one over 64). Measured on CortexDB 0.10.4: - the same key and body is a replay, answered with the first event's id and diff --git a/docs/architecture/cortex-wire.md b/docs/architecture/cortex-wire.md index e1dc331e..8034179b 100644 --- a/docs/architecture/cortex-wire.md +++ b/docs/architecture/cortex-wire.md @@ -63,7 +63,7 @@ Request body (one event). It is the same on both wires: { "scope": "app:tinymemory/agent:researcher/app:documents", "modality": "document", - "idempotency_key": "tm3:<56 hex: SHA-256 of this body without the key>", + "idempotency_key": "tm3:<the first 56 hex digits (224 bits) of the SHA-256 of this body without the key>", "content": { "kind": "message", "role": "user", "text": "<envelope JSON, see below>" }, "context": { "labels": ["tm:i:<16 hex>", "tm:k:<16 hex>"], @@ -480,7 +480,8 @@ Each was measured against a live CortexDB and was wrong in the first adapter. The loopback doubles reproduce all of them (see [testing](testing.md)). - **Append-only.** There is no update route. A reused body `idempotency_key` - with the same body is a replay for 24 hours, and with another body a 409. + with the same body is a replay for 24 hours, and with another body a 409 + (which this crate's keys, derived from the body, never produce). Forget by `memory_ids` releases the keys of what it removes (measured on 0.10.4; the v1 adapter's notes said the opposite, on an unrecorded build). - **Accepted is not readable.** An append answers `202` and indexes afterwards.