From 66a6956a335f7b554b8faa4f2c6dca0a3716f6f4 Mon Sep 17 00:00:00 2001 From: M3gA-Mind Date: Tue, 6 Oct 2026 17:36:24 +0530 Subject: [PATCH 1/8] Write long documents as pieces under CortexDB's 1 MiB event limit CortexDB refuses an experience over 1 MiB of flattened text (422 INVALID_ENVELOPE, never truncated), and a document's event text is its whole JSON envelope, so a long file could not be stored at all. - A document whose encoded envelope fits in 256 KiB is still one event, byte-identical to before. A longer one is split into contiguous pieces: at page breaks, then before markdown headings, packed greedily up to that target, with blank-line, line and character cuts for a unit that is too big. `DOCUMENT_CHUNK_TARGET_BYTES` is the only granularity knob; at 0 every page and section is its own event. - Each piece carries `chunk {index, count, pages?, section?}` and the item's id and label, so replay, `forget` and `get` see every piece, and a store that failed part way writes only the missing ones. - No event over 768 KiB of encoded envelope is ever sent: a batch is encoded and checked before its first write, and an item that cannot fit (a learning or a turn that long) is `InvalidRequest`. - `get` and `list` reassemble a chunked document. A fetch hit or recall citation on a piece is that piece, with the item's id and `page:` and `section:` tags. Events written before chunking read unchanged. - The PDF converter extracts pages one by one and joins them with a form feed (`documents::PAGE_BREAK`), each normalized on its own, so page numbers survive conversion. - The test double refuses events over 1 MiB, as CortexDB does. This differs on purpose from "one event per page or section": CortexDB 0.10.4 already fragments each event for retrieval and serves an over-budget event as an excerpt, and hosted writes are billed per event. --- .../src/cortex/README.md | 12 +- .../src/cortex/engine/fetch.rs | 15 +- .../src/cortex/engine/items.rs | 46 ++- .../src/cortex/engine/list.rs | 56 ++-- .../src/cortex/engine/mod.rs | 4 + .../src/cortex/engine/mod_chunk_tests.rs | 255 ++++++++++++++++ .../src/cortex/engine/recall.rs | 4 +- .../src/cortex/engine/store.rs | 34 ++- .../src/cortex/envelope/chunks.rs | 286 ++++++++++++++++++ .../src/cortex/envelope/chunks_tests.rs | 143 +++++++++ .../src/cortex/envelope/mod.rs | 118 +++++++- .../src/cortex/envelope/rebuild.rs | 25 +- .../src/cortex/testing/log.rs | 6 +- .../src/documents/README.md | 4 +- .../src/documents/mod.rs | 6 + .../src/documents/office/mod.rs | 27 +- .../src/documents/office/mod_tests.rs | 54 +++- .../src/documents/office/pdf.rs | 13 +- docs/architecture/cortex-flows.md | 8 +- docs/architecture/cortex-wire.md | 49 ++- docs/specs/memory-v2.md | 2 +- 21 files changed, 1079 insertions(+), 88 deletions(-) create mode 100644 crates/tinymemory-integrations/src/cortex/engine/mod_chunk_tests.rs create mode 100644 crates/tinymemory-integrations/src/cortex/envelope/chunks.rs create mode 100644 crates/tinymemory-integrations/src/cortex/envelope/chunks_tests.rs diff --git a/crates/tinymemory-integrations/src/cortex/README.md b/crates/tinymemory-integrations/src/cortex/README.md index d5695730..97c30054 100644 --- a/crates/tinymemory-integrations/src/cortex/README.md +++ b/crates/tinymemory-integrations/src/cortex/README.md @@ -110,8 +110,16 @@ tokens and refuses a request without it (`401 ACTOR_MISMATCH`); a static operator key is served as `user:local`; a server with no `whoami` route gets no header. The hosted (TinyHumans) wire names the actor itself. -**Events.** A document or learning is one event; a conversation is one event -per turn, appended in order. Each event's `content.text` is a JSON envelope: +**Events.** A learning is one event; a conversation is one event per turn, +appended in order; a document is one event, or, when its envelope would pass +256 KiB (`envelope::chunks::DOCUMENT_CHUNK_TARGET_BYTES`, the one granularity +knob), one event per piece of its body, cut at page breaks and headings and +packed up to that size (`envelope/chunks.rs`). No event is sent over 768 KiB of +encoded envelope (CortexDB refuses an experience over 1 MiB); an item that +cannot fit is refused before anything of its batch is sent. `get` and `list` +reassemble a chunked document; a `fetch` or `recall` hit on a piece carries +that piece, tagged `page:` and `section:`. Each event's +`content.text` is a JSON envelope: ```json { "v": 2, "id": "<40-hex fingerprint>", "kind": "conversation", diff --git a/crates/tinymemory-integrations/src/cortex/engine/fetch.rs b/crates/tinymemory-integrations/src/cortex/engine/fetch.rs index 637b9507..6016e2aa 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/fetch.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/fetch.rs @@ -33,8 +33,8 @@ use tinymemory_api::{FetchPage, FetchRequest, Hit, ItemKind, MetaFilter}; use super::CortexEngine; use super::beliefs::{beliefs_in, merge}; use super::cursor::{self, FetchCursor}; -use super::items::{hit, keeps}; -use crate::cortex::envelope::{Envelope, decode_event, labels, rebuild}; +use super::items::{event_hit, hit, keeps}; +use crate::cortex::envelope::{Envelope, decode_event, labels}; use crate::cortex::error::{Error, Result}; /// The cursor tag of a fetch. @@ -159,11 +159,12 @@ impl CortexEngine { .into_iter() .filter_map(|(rank, envelope)| { let score = 1.0 / (1.0 + rank as f32); - let item = match envelope.kind { - ItemKind::Conversation => conversations.get(&envelope.id)?.clone(), - _ => rebuild(std::slice::from_ref(&envelope))?, - }; - Some(hit(&envelope.id, &item, score)) + match envelope.kind { + ItemKind::Conversation => { + Some(hit(&envelope.id, conversations.get(&envelope.id)?, score)) + } + _ => event_hit(&envelope, score), + } }) .collect(); let next_cursor = if more { diff --git a/crates/tinymemory-integrations/src/cortex/engine/items.rs b/crates/tinymemory-integrations/src/cortex/engine/items.rs index b730ede7..5f94e25d 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/items.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/items.rs @@ -4,7 +4,9 @@ use std::collections::{BTreeMap, HashMap}; use tinymemory_api::explore::in_request_order; -use tinymemory_api::{GetRequest, Hit, ItemId, ItemKind, MetaFilter, Namespace, StoreItem}; +use tinymemory_api::{ + GetRequest, Hit, ItemId, ItemKind, MemoryMeta, MetaFilter, Namespace, StoreItem, +}; use super::CortexEngine; use super::scopes::KindScope; @@ -25,6 +27,35 @@ pub(super) fn keeps(filter: &MetaFilter, kind: ItemKind, envelope: &Envelope) -> envelope.kind == kind && filter.matches(kind, &envelope.meta) } +/// An envelope's metadata, located: for a piece of a chunked document, the +/// item's metadata plus a `page:<n>` (or `page:<first>-<last>`) tag and a +/// `section:<title>` tag for what the piece covers. Read-side only: the +/// stored item carries neither. +pub(super) fn located_meta(envelope: &Envelope) -> MemoryMeta { + let mut meta = envelope.meta.clone(); + if let Some(chunk) = &envelope.chunk { + match chunk.pages { + Some([first, last]) if first == last => meta.tags.push(format!("page:{first}")), + Some([first, last]) => meta.tags.push(format!("page:{first}-{last}")), + None => {} + } + if let Some(section) = &chunk.section { + meta.tags.push(format!("section:{section}")); + } + } + meta +} + +/// A ranked hit for one event's envelope: a learning or a whole document as +/// it was stored, or, for a piece of a chunked document, that piece (the +/// item's id, the piece's text, its located metadata). +pub(super) fn event_hit(envelope: &Envelope, score: f32) -> Option<Hit> { + let item = rebuild(std::slice::from_ref(envelope))?; + let mut found = hit(&envelope.id, &item, score); + found.meta = located_meta(envelope); + Some(found) +} + /// A hit for `item`. pub(super) fn hit(id: &str, item: &StoreItem, score: f32) -> Hit { Hit { @@ -95,6 +126,17 @@ impl CortexEngine { pub(super) async fn conversations( &self, ids: &[(String, Namespace)], + ) -> Result<HashMap<String, StoreItem>> { + self.assembled(ItemKind::Conversation, ids).await + } + + /// The whole items of `kind` named by `ids`, each at its namespace, + /// rebuilt from all their events: a conversation's turns, a chunked + /// document's pieces (one lookup per namespace). + pub(super) async fn assembled( + &self, + kind: ItemKind, + ids: &[(String, Namespace)], ) -> Result<HashMap<String, StoreItem>> { let mut by_node: BTreeMap<&Namespace, Vec<String>> = BTreeMap::new(); for (id, namespace) in ids { @@ -102,7 +144,7 @@ impl CortexEngine { } let mut out = HashMap::new(); for (namespace, ids) in by_node { - let scope = KindScope::new(namespace.clone(), ItemKind::Conversation); + let scope = KindScope::new(namespace.clone(), kind); for (id, events) in self.item_events(&scope, &ids).await? { let envelopes: Vec<Envelope> = events.into_iter().map(|d| d.envelope).collect(); if let Some(item) = rebuild(&envelopes) { diff --git a/crates/tinymemory-integrations/src/cortex/engine/list.rs b/crates/tinymemory-integrations/src/cortex/engine/list.rs index e7e6e25b..d49a55d1 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/list.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/list.rs @@ -9,9 +9,10 @@ //! that one label first (see `envelope::labels`); the client-side check runs //! regardless, and the cursor stays the engine's. //! -//! **Each item once.** A document or learning is one event. A conversation -//! is emitted only on the page holding its turn-0 event, and its text is -//! assembled from all its turns by one label lookup per page. Writes are +//! **Each item once.** A learning, or a document written whole, is one +//! event. A conversation, or a chunked document, is emitted only on the page +//! holding its first event (turn 0, piece 0), and its text is assembled from +//! all its events by one label lookup per page and kind. Writes are //! ordered, so a conversation whose store failed part-way still has its //! turn 0 and lists with the turns it holds. //! @@ -20,7 +21,7 @@ //! cursor remembers the last id). A page that ends mid-way is resumed by //! re-reading the same engine page and skipping the consumed events. -use std::collections::HashSet; +use std::collections::{HashMap, HashSet}; use serde_json::Value; use tinymemory_api::{Hit, ItemKind, ListPage, ListRequest, Namespace}; @@ -36,11 +37,11 @@ use crate::cortex::log::{MAX_PAGES, PAGE_SIZE}; /// The cursor tag of a listing. const TAG: char = 'l'; -/// A hit, or a conversation whose turns are assembled before the page -/// returns. +/// A hit, or an item whose events (a conversation's turns, a chunked +/// document's pieces) are assembled before the page returns. enum Pending { Ready(Box<Hit>), - Conversation(String, Namespace), + Assembled(ItemKind, String, Namespace), } impl CortexEngine { @@ -137,36 +138,45 @@ impl CortexEngine { if !keeps(&req.filter, kind, &envelope) { return None; } - let starts = envelope.turn.as_ref().is_none_or(|turn| turn.index == 0); + let starts = envelope.part().is_none_or(|index| index == 0); if !starts || !seen.insert(envelope.id.clone()) { return None; } - if kind == ItemKind::Conversation { - return Some(Pending::Conversation(envelope.id, envelope.meta.namespace)); + if kind == ItemKind::Conversation || envelope.chunk.is_some() { + return Some(Pending::Assembled( + kind, + envelope.id, + envelope.meta.namespace, + )); } let id = envelope.id.clone(); let item = rebuild(std::slice::from_ref::<Envelope>(&envelope))?; Some(Pending::Ready(Box::new(hit(&id, &item, 0.0)))) } - /// Assembles the page's conversations (one lookup for all of them) and - /// returns the hits in listing order. + /// Assembles the page's conversations and chunked documents (one lookup + /// per kind) and returns the hits in listing order. async fn resolve(&self, pending: Vec<Pending>) -> Result<Vec<Hit>> { - let ids: Vec<(String, Namespace)> = pending - .iter() - .filter_map(|p| match p { - Pending::Conversation(id, namespace) => Some((id.clone(), namespace.clone())), - Pending::Ready(_) => None, - }) - .collect(); - let conversations = self.conversations(&ids).await?; + let mut assembled = HashMap::new(); + for kind in [ItemKind::Conversation, ItemKind::Document] { + let ids: Vec<(String, Namespace)> = pending + .iter() + .filter_map(|p| match p { + Pending::Assembled(of, id, namespace) if *of == kind => { + Some((id.clone(), namespace.clone())) + } + _ => None, + }) + .collect(); + if !ids.is_empty() { + assembled.extend(self.assembled(kind, &ids).await?); + } + } Ok(pending .into_iter() .filter_map(|p| match p { Pending::Ready(hit) => Some(*hit), - Pending::Conversation(id, _) => { - conversations.get(&id).map(|item| hit(&id, item, 0.0)) - } + Pending::Assembled(_, id, _) => assembled.get(&id).map(|item| hit(&id, item, 0.0)), }) .collect()) } diff --git a/crates/tinymemory-integrations/src/cortex/engine/mod.rs b/crates/tinymemory-integrations/src/cortex/engine/mod.rs index b68b1b10..296ef681 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/mod.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/mod.rs @@ -263,6 +263,10 @@ mod tests; #[path = "mod_list_tests.rs"] mod list_tests; +#[cfg(test)] +#[path = "mod_chunk_tests.rs"] +mod chunk_tests; + #[cfg(test)] #[path = "mod_direct_tests.rs"] mod direct_tests; diff --git a/crates/tinymemory-integrations/src/cortex/engine/mod_chunk_tests.rs b/crates/tinymemory-integrations/src/cortex/engine/mod_chunk_tests.rs new file mode 100644 index 00000000..68e07a67 --- /dev/null +++ b/crates/tinymemory-integrations/src/cortex/engine/mod_chunk_tests.rs @@ -0,0 +1,255 @@ +//! Chunked documents end to end on both wires: written in pieces under the +//! event limit, read back whole by `get` and `list`, ranked by piece, and +//! forgotten whole; documents written as one event read as before. + +use serde_json::{Value, json}; +use tinymemory_api::{ + DocumentBody, FetchMode, FetchRequest, ForgetTarget, GetRequest, ItemId, ItemKind, + LearningKind, ListRequest, MemoryMeta, MetaFilter, Namespace, Reach, StoreItem, +}; + +use super::*; +use crate::cortex::envelope::chunks::{MAX_EVENT_TEXT_BYTES, PAGE_BREAK}; +use crate::cortex::testing::{Shared, both, direct_double, direct_engine}; + +const SCOPE: &str = "app:tinymemory/source:pdf/app:documents"; + +fn meta() -> MemoryMeta { + MemoryMeta { + namespace: Namespace::source("pdf"), + file_path: Some("/docs/handbook.pdf".into()), + ..MemoryMeta::default() + } +} + +/// A document of `pages` pages, each a titled section of about 20 KiB, so +/// the whole is well over one event's target. +fn handbook(pages: usize) -> StoreItem { + let body: Vec<String> = (1..=pages) + .map(|page| { + let filler = format!("Policy text for page {page}. ").repeat(700); + format!("# Chapter {page}\n\n{filler}\n") + }) + .collect(); + StoreItem::Document { + title: Some("Handbook".into()), + body: DocumentBody::Text(body.join(&PAGE_BREAK.to_string())), + mime: Some("application/pdf".into()), + meta: meta(), + } +} + +fn body(item: &StoreItem) -> &str { + match item { + StoreItem::Document { + body: DocumentBody::Text(text), + .. + } => text, + _ => panic!("a document"), + } +} + +/// The events the double holds in `scope`. +fn events(state: &Shared, scope: &str) -> Vec<Value> { + state + .log + .lock() + .unwrap() + .events + .iter() + .filter(|event| event["scope"] == scope) + .cloned() + .collect() +} + +#[tokio::test] +async fn a_long_document_is_written_in_pieces_under_the_limit_and_read_back_whole() { + for (engine, state) in both().await { + let item = handbook(40); + let id = item.fingerprint(); + let receipt = engine.store(item.clone()).await.unwrap(); + assert_eq!(receipt.id.as_str(), id, "the item keeps its identity"); + + let written = events(&state, SCOPE); + assert!(written.len() > 2, "{} pieces", written.len()); + for event in &written { + let text = event["content"]["text"].as_str().unwrap(); + assert!(text.len() <= MAX_EVENT_TEXT_BYTES, "{}", text.len()); + assert!( + text.len() <= 300 * 1024, + "packed near the target: {}", + text.len() + ); + } + + let got = engine + .get(GetRequest { + ids: vec![ItemId::new(id.clone())], + reach: None, + }) + .await + .unwrap(); + assert_eq!(got.len(), 1); + assert_eq!(got[0].text, item.render_text(), "get reassembles the body"); + + let listed = engine + .list(ListRequest::new( + MetaFilter::kinds([ItemKind::Document]), + 10, + )) + .await + .unwrap(); + assert_eq!(listed.items.len(), 1, "one item, not one per piece"); + assert_eq!(listed.items[0].text, item.render_text()); + } +} + +#[tokio::test] +async fn a_hit_on_a_piece_is_that_piece_with_its_page_and_section() { + let (endpoint, _state) = direct_double().await; + let engine = direct_engine(&endpoint); + let item = handbook(40); + engine.store(item.clone()).await.unwrap(); + let mut request = FetchRequest::new("Policy text page", FetchMode::Hybrid, 3); + request.filter.reach = Some(Reach::exact(Namespace::source("pdf"))); + let page = engine.fetch(request).await.unwrap(); + assert!(!page.hits.is_empty()); + let hit = &page.hits[0]; + assert_eq!(hit.id.as_str(), item.fingerprint(), "the item's own id"); + assert!(hit.text.len() < body(&item).len(), "a piece, not the whole"); + assert!(body(&item).contains(hit.text.split("\n\n").nth(1).unwrap_or_default())); + assert_eq!(hit.meta.file_path.as_deref(), Some("/docs/handbook.pdf")); + assert!( + hit.meta.tags.iter().any(|tag| tag.starts_with("page:")), + "{:?}", + hit.meta.tags + ); + assert!( + hit.meta + .tags + .iter() + .any(|tag| tag.starts_with("section:Chapter ")), + "{:?}", + hit.meta.tags + ); +} + +#[tokio::test] +async fn forgetting_a_chunked_document_removes_every_piece() { + for (engine, state) in both().await { + let item = handbook(30); + let receipt = engine.store(item).await.unwrap(); + assert!(events(&state, SCOPE).len() > 2); + let report = engine + .forget(ForgetTarget::Ids(vec![receipt.id.clone()])) + .await + .unwrap(); + assert_eq!(report.forgotten, 1, "one item"); + assert!(events(&state, SCOPE).is_empty(), "no piece left behind"); + } +} + +#[tokio::test] +async fn a_store_that_lost_a_piece_writes_only_that_piece_again() { + let (endpoint, state) = direct_double().await; + let engine = direct_engine(&endpoint); + 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); + } + let again = engine.store(item.clone()).await.unwrap(); + assert!(!again.replayed, "a piece was missing"); + assert_eq!(events(&state, SCOPE).len(), before, "exactly that piece"); + let replay = engine.store(item).await.unwrap(); + assert!(replay.replayed, "every piece present: a replay"); +} + +#[tokio::test] +async fn a_short_document_is_one_event_exactly_as_before_and_reads_the_same() { + for (engine, state) in both().await { + let item = StoreItem::Document { + title: Some("Note".into()), + body: DocumentBody::Text(format!("Page one.{PAGE_BREAK}Page two.")), + mime: None, + meta: meta(), + }; + 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}"); + let listed = engine + .list(ListRequest::new(MetaFilter::default(), 10)) + .await + .unwrap(); + assert_eq!(listed.items[0].text, item.render_text()); + assert!(listed.items[0].meta.tags.is_empty()); + } +} + +#[tokio::test] +async fn an_event_written_before_chunking_still_reads() { + let (endpoint, state) = direct_double().await; + let engine = direct_engine(&endpoint); + let id = "legacy-doc-id"; + let envelope = json!({ + "v": 2, + "id": id, + "kind": "document", + "text": "Refunds take five days.", + "meta": { "namespace": "source:pdf", "source": { "kind": "file" } }, + "title": "Refunds" + }); + state.log.lock().unwrap().append(&json!({ + "scope": SCOPE, + "modality": "document", + "idempotency_key": "legacy", + "content": { "kind": "message", "role": "user", "text": envelope.to_string() }, + "context": { "labels": [crate::cortex::envelope::labels::item(id)] }, + })); + let got = engine + .get(GetRequest { + ids: vec![ItemId::new(id)], + reach: None, + }) + .await + .unwrap(); + assert_eq!(got.len(), 1); + assert!(got[0].text.contains("Refunds take five days.")); + let mut request = FetchRequest::new("refunds", FetchMode::Hybrid, 3); + request.filter.reach = Some(Reach::exact(Namespace::source("pdf"))); + let page = engine.fetch(request).await.unwrap(); + assert_eq!(page.hits.len(), 1); + assert!(page.hits[0].text.contains("Refunds take five days.")); + assert!(page.hits[0].meta.tags.is_empty()); +} + +#[tokio::test] +async fn an_event_over_the_limit_is_refused_before_anything_is_sent() { + for (engine, state) in both().await { + let huge = StoreItem::learning( + "x".repeat(MAX_EVENT_TEXT_BYTES), + LearningKind::Fact, + 0.9, + MemoryMeta::default(), + ); + let small = StoreItem::document("A small note.", MemoryMeta::default()); + let writes_before = state.count("POST"); + let error = engine.store_many(vec![small, huge]).await.unwrap_err(); + assert!( + matches!(error, crate::cortex::Error::InvalidRequest(_)), + "{error:?}" + ); + assert!(error.to_string().contains("1 MiB"), "{error}"); + assert_eq!(state.count("POST"), writes_before, "nothing was written"); + } +} diff --git a/crates/tinymemory-integrations/src/cortex/engine/recall.rs b/crates/tinymemory-integrations/src/cortex/engine/recall.rs index 720e5b58..e21529b0 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/recall.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/recall.rs @@ -34,7 +34,7 @@ use tinymemory_api::{Citation, ItemId, Namespace, Reach, RecallAnswer, RecallReq use super::CortexEngine; use super::fetch::{interleave, ranked, recall_body}; -use super::items::admitted; +use super::items::{admitted, located_meta}; use super::scopes::KindScope; use crate::cortex::descriptor::CortexWire; use crate::cortex::envelope::Envelope; @@ -164,10 +164,10 @@ impl CortexEngine { .filter(|envelope| seen.insert(envelope.id.clone())) .take(req.limit) .map(|envelope| Citation { + meta: located_meta(&envelope), id: ItemId::new(envelope.id), kind: envelope.kind, snippet: envelope.text, - meta: envelope.meta, score: None, }) .collect(); diff --git a/crates/tinymemory-integrations/src/cortex/engine/store.rs b/crates/tinymemory-integrations/src/cortex/engine/store.rs index f8e75d1e..8f558dd4 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/store.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/store.rs @@ -10,8 +10,9 @@ //! //! - every event present (the whole document, learning, or every turn) — //! a replay: nothing is written and the receipt says so; -//! - some turns of a conversation present — a previous store failed part -//! way, and only the missing turns are written, in order; +//! - some turns of a conversation, or some pieces of a chunked document, +//! present — a previous store failed part way, and only the missing ones +//! are written, in order; //! - nothing present — every event is written. //! //! Writes use fresh idempotency keys (see `transport::fresh_idempotency_key`) @@ -50,6 +51,18 @@ impl CortexEngine { ) -> Result<Vec<StoreReceipt>> { validate_many(&items)?; 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)>> = + Vec::with_capacity(items.len()); + for (item, id) in items.iter().zip(&ids) { + let mut events = Vec::new(); + for envelope in Envelope::for_item(item, id)? { + let encoded = envelope.encode_checked()?; + events.push((envelope.part(), envelope, encoded)); + } + planned.push(events); + } let mut held: HashMap<String, HashSet<Option<u32>>> = HashMap::new(); let mut by_scope: BTreeMap<KindScope, Vec<String>> = BTreeMap::new(); for (item, id) in items.iter().zip(&ids) { @@ -60,26 +73,23 @@ impl CortexEngine { } for (scope, of_scope) in &by_scope { for (id, events) in self.item_events(scope, of_scope).await? { - held.entry(id).or_default().extend( - events - .iter() - .map(|decoded| decoded.envelope.turn.as_ref().map(|turn| turn.index)), - ); + held.entry(id) + .or_default() + .extend(events.iter().map(|decoded| decoded.envelope.part())); } } let mut receipts = Vec::with_capacity(items.len()); let mut written_here: HashSet<String> = HashSet::new(); let mut last_per_scope: Vec<Written> = Vec::new(); - for (item, id) in items.iter().zip(ids) { + for (events, id) in planned.into_iter().zip(ids) { let present = held.get(&id); let mut requests = Vec::new(); if !written_here.contains(&id) { - for envelope in Envelope::for_item(item, &id)? { - let turn = envelope.turn.as_ref().map(|turn| turn.index); - if present.is_some_and(|present| present.contains(&turn)) { + for (part, envelope, encoded) in events { + if present.is_some_and(|present| present.contains(&part)) { continue; } - requests.push(envelope.request(&envelope.encode()?)); + requests.push(envelope.request(&encoded)); } } let replayed = requests.is_empty(); diff --git a/crates/tinymemory-integrations/src/cortex/envelope/chunks.rs b/crates/tinymemory-integrations/src/cortex/envelope/chunks.rs new file mode 100644 index 00000000..81423d50 --- /dev/null +++ b/crates/tinymemory-integrations/src/cortex/envelope/chunks.rs @@ -0,0 +1,286 @@ +//! Splitting a document's text into the events it is written as. +//! +//! CortexDB refuses an experience whose flattened text is over 1 MiB +//! (`422 INVALID_ENVELOPE`, never truncated). A document's event text is its +//! whole [`super::Envelope`], so a long document is written as several +//! events, each carrying one contiguous piece of the body; read back in order +//! they concatenate to the body exactly. +//! +//! **Units.** The text is first cut where its structure says to: at every +//! page break ([`PAGE_BREAK`], which the PDF converter puts between pages), +//! and within a page before every markdown heading line. A unit never spans +//! two pages. +//! +//! **Packing.** Units are packed greedily, in order, into pieces of at most +//! [`DOCUMENT_CHUNK_TARGET_BYTES`]. A unit larger than that is cut finer, at +//! blank lines, then line ends, then characters, and those parts packed the +//! same way. Every size is the piece's cost inside the encoded envelope (its +//! JSON-escaped length), so a piece and the envelope around it stay under +//! [`MAX_EVENT_TEXT_BYTES`]. +//! +//! A document that fits in one piece is one event, exactly as before +//! chunking existed. + +/// The page break the PDF converter writes between pages: a form feed, the +/// same character as `documents::PAGE_BREAK` (this module does not depend on +/// the `documents` feature). +pub(crate) const PAGE_BREAK: char = '\u{c}'; + +/// The most one event's text (its encoded envelope) may take: CortexDB's +/// 1 MiB limit with a quarter left as a safety margin, because the server +/// counts its own flattening of the text, which this crate cannot see. +pub(crate) const MAX_EVENT_TEXT_BYTES: usize = 768 * 1024; + +/// How much of a document's body one event carries at most, as encoded +/// bytes. +/// +/// The one knob for granularity. At this value a document under it stays +/// one event and a longer one is packed into as few events as fit. At `0` +/// every page and every section is its own event, and a part is cut finer +/// only when it alone would exceed [`MAX_EVENT_TEXT_BYTES`]. +pub(crate) const DOCUMENT_CHUNK_TARGET_BYTES: usize = 256 * 1024; + +/// The longest section title an event carries, in characters. +pub(crate) const MAX_SECTION_CHARS: usize = 120; + +/// One piece of a document: a contiguous slice of its text. +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) struct Piece<'a> { + /// The text, exactly as it appears in the document. + pub(crate) text: &'a str, + /// The pages it covers, first and last, counted from 1, when the + /// document marks its pages. + pub(crate) pages: Option<(u32, u32)>, + /// The title of the section it starts in, when one was seen. + pub(crate) section: Option<String>, +} + +/// A unit before packing: a byte range, its page, and its section. +struct Unit { + start: usize, + end: usize, + page: u32, + section: Option<String>, +} + +/// The pieces of `text`, each at most `target` (or, at `0`, one per unit) +/// and never over `limit`, both as encoded bytes on top of `overhead` (the +/// envelope around the piece). One piece for a text that fits. +pub(crate) fn split(text: &str, overhead: usize, target: usize, limit: usize) -> Vec<Piece<'_>> { + let room = limit.saturating_sub(overhead).max(1); + let pack = if target == 0 { + room + } else { + target.saturating_sub(overhead).clamp(1, room) + }; + let paged = text.contains(PAGE_BREAK); + if target != 0 && escaped_len(text) <= pack { + return vec![Piece { + text, + pages: paged.then(|| (1, page_count(text))), + section: None, + }]; + } + let mut pieces: Vec<Piece<'_>> = Vec::new(); + let mut open: Option<Open> = None; + for unit in units(text) { + for (start, end) in parts(text, unit.start, unit.end, pack) { + let size = escaped_len(&text[start..end]); + match open.as_mut() { + Some(current) if target != 0 && current.size + size <= pack => { + current.end = end; + current.last_page = unit.page; + current.size += size; + } + _ => { + if let Some(done) = open.take() { + pieces.push(done.piece(text, paged)); + } + open = Some(Open { + start, + end, + first_page: unit.page, + last_page: unit.page, + section: unit.section.clone(), + size, + }); + } + } + } + } + pieces.extend(open.map(|done| done.piece(text, paged))); + pieces +} + +/// The piece being packed. +struct Open { + start: usize, + end: usize, + first_page: u32, + last_page: u32, + section: Option<String>, + size: usize, +} + +impl Open { + fn piece(self, text: &str, paged: bool) -> Piece<'_> { + Piece { + text: &text[self.start..self.end], + pages: paged.then_some((self.first_page, self.last_page)), + section: self.section, + } + } +} + +/// The units of `text`: one per page, cut again before every heading line. +/// A stretch holding only whitespace and page breaks is never a unit of its +/// own; it joins the unit that follows. +fn units(text: &str) -> Vec<Unit> { + let mut units = Vec::new(); + let mut page = 1; + let mut section: Option<String> = None; + let mut start = 0; + let mut start_page = 1; + let mut line_start = 0; + for (at, ch) in text.char_indices() { + let breaks_page = ch == PAGE_BREAK; + let title = if at == line_start { + heading(&text[at..]) + } else { + None + }; + if (breaks_page || title.is_some()) && !blank(&text[start..at]) { + units.push(Unit { + start, + end: at, + page: start_page, + section: section.clone(), + }); + start = at; + start_page = page; + } + if breaks_page { + page += 1; + if blank(&text[start..at]) { + start_page = page; + } + } + if let Some(title) = title { + section = Some(title); + } + if ch == '\n' || breaks_page { + line_start = at + ch.len_utf8(); + } + } + if start < text.len() { + units.push(Unit { + start, + end: text.len(), + page: start_page, + section, + }); + } + units +} + +/// Whether `text` holds nothing but whitespace and page breaks. +fn blank(text: &str) -> bool { + text.chars().all(|c| c == PAGE_BREAK || c.is_whitespace()) +} + +/// The title of a markdown heading line starting at the head of `rest` +/// (`#` to `######`, a space, then text), trimmed and shortened. +fn heading(rest: &str) -> Option<String> { + let line = rest.split('\n').next().unwrap_or_default(); + let trimmed = line.trim_start_matches(' '); + if line.len() - trimmed.len() > 3 { + return None; + } + let hashes = trimmed.chars().take_while(|c| *c == '#').count(); + if !(1..=6).contains(&hashes) { + return None; + } + let after = &trimmed[hashes..]; + if !after.starts_with(' ') { + return None; + } + let title = after.trim().trim_end_matches('#').trim(); + (!title.is_empty()).then(|| title.chars().take(MAX_SECTION_CHARS).collect()) +} + +/// `text[start..end]` cut into ranges of at most `cap` encoded bytes: whole +/// when it fits, else at blank lines, then line ends, then characters. +fn parts(text: &str, start: usize, end: usize, cap: usize) -> Vec<(usize, usize)> { + if escaped_len(&text[start..end]) <= cap { + return vec![(start, end)]; + } + for separator in ["\n\n", "\n"] { + let cuts = cut_after(text, start, end, separator); + if cuts.len() > 1 { + return cuts + .into_iter() + .flat_map(|(s, e)| parts(text, s, e, cap)) + .collect::<Vec<_>>() + .into_iter() + .fold(Vec::new(), |mut packed: Vec<(usize, usize)>, (s, e)| { + match packed.last_mut() { + Some(last) if escaped_len(&text[last.0..e]) <= cap => last.1 = e, + _ => packed.push((s, e)), + } + packed + }); + } + } + let mut ranges = Vec::new(); + let mut from = start; + let mut used = 0; + for (at, ch) in text[start..end].char_indices() { + let size = escaped_char_len(ch); + if used + size > cap && start + at > from { + ranges.push((from, start + at)); + from = start + at; + used = 0; + } + used += size; + } + ranges.push((from, end)); + ranges +} + +/// `text[start..end]` cut just after every `separator`. +fn cut_after(text: &str, start: usize, end: usize, separator: &str) -> Vec<(usize, usize)> { + let mut cuts = Vec::new(); + let mut from = start; + for (at, _) in text[start..end].match_indices(separator) { + let to = start + at + separator.len(); + if to < end { + cuts.push((from, to)); + from = to; + } + } + cuts.push((from, end)); + cuts +} + +/// How many pages `text` marks: one more than its page breaks. +fn page_count(text: &str) -> u32 { + let breaks = text.matches(PAGE_BREAK).count(); + u32::try_from(breaks).map_or(u32::MAX, |breaks| breaks.saturating_add(1)) +} + +/// The length of `text` as a JSON string body, as `serde_json` escapes it. +pub(crate) fn escaped_len(text: &str) -> usize { + text.chars().map(escaped_char_len).sum() +} + +fn escaped_char_len(ch: char) -> usize { + match ch { + '"' | '\\' | '\n' | '\r' | '\t' | '\u{8}' | '\u{c}' => 2, + c if u32::from(c) < 0x20 => 6, + c => c.len_utf8(), + } +} + +#[cfg(test)] +#[path = "chunks_tests.rs"] +mod tests; diff --git a/crates/tinymemory-integrations/src/cortex/envelope/chunks_tests.rs b/crates/tinymemory-integrations/src/cortex/envelope/chunks_tests.rs new file mode 100644 index 00000000..0e971d6e --- /dev/null +++ b/crates/tinymemory-integrations/src/cortex/envelope/chunks_tests.rs @@ -0,0 +1,143 @@ +//! Splitting a document: exact reassembly, structure, pages and the limits. + +use super::*; + +/// The pieces of `text` under a small overhead and the given sizes, +/// checked to concatenate back to `text`. +fn pieces(text: &str, target: usize, limit: usize) -> Vec<Piece<'_>> { + let pieces = split(text, 10, target, limit); + let joined: String = pieces.iter().map(|piece| piece.text).collect(); + assert_eq!(joined, text, "the pieces concatenate to the text"); + pieces +} + +#[test] +fn a_text_that_fits_is_one_piece() { + let text = "# Title\n\nShort body.\n"; + let one = pieces(text, 1000, 10_000); + assert_eq!(one.len(), 1); + assert_eq!(one[0].text, text); + assert_eq!(one[0].pages, None); +} + +#[test] +fn sections_are_packed_up_to_the_target_and_named_by_their_heading() { + let section = |title: &str| format!("## {title}\n\n{}\n\n", "word ".repeat(40)); + let text = format!( + "{}{}{}{}", + section("Alpha"), + section("Beta"), + section("Gamma"), + section("Delta") + ); + let packed = pieces(&text, 500, 10_000); + assert!(packed.len() >= 2 && packed.len() < 4, "{}", packed.len()); + for piece in &packed { + assert!(escaped_len(piece.text) + 10 <= 500, "{}", piece.text.len()); + assert!( + piece.text.starts_with("## "), + "cut at a heading: {:?}", + piece.text + ); + } + assert_eq!(packed[0].section.as_deref(), Some("Alpha")); + let sections: Vec<_> = packed.iter().filter_map(|p| p.section.as_deref()).collect(); + assert!(sections.windows(2).all(|w| w[0] != w[1]), "{sections:?}"); +} + +#[test] +fn at_a_zero_target_every_section_is_its_own_piece() { + let text = "intro\n# One\nfirst\n# Two\nsecond\n# Three\nthird\n"; + let each = pieces(text, 0, 10_000); + let sections: Vec<_> = each.iter().map(|p| p.section.as_deref()).collect(); + assert_eq!( + sections, + [None, Some("One"), Some("Two"), Some("Three")], + "{each:?}" + ); +} + +#[test] +fn page_breaks_number_the_pieces_and_never_split_a_page_across_units() { + let page = |n: u32| format!("Page {n} text. {}\n", "x".repeat(60)); + let text = format!( + "{}{PAGE_BREAK}{}{PAGE_BREAK}{}{PAGE_BREAK}{}", + page(1), + page(2), + page(3), + page(4) + ); + let each = pieces(&text, 0, 10_000); + let pages: Vec<_> = each.iter().map(|p| p.pages).collect(); + assert_eq!( + pages, + [Some((1, 1)), Some((2, 2)), Some((3, 3)), Some((4, 4))] + ); + assert!(each[2].text.contains("Page 3"), "{:?}", each[2].text); + let packed = pieces(&text, 200, 10_000); + assert!(packed.len() > 1 && packed.len() < 4, "{packed:?}"); + assert_eq!(packed[0].pages.map(|(first, _)| first), Some(1)); + assert_eq!(packed.last().unwrap().pages.map(|(_, last)| last), Some(4)); + let one = pieces(&text, 10_000, 10_000); + assert_eq!( + one[0].pages, + Some((1, 4)), + "a paged text that fits spans all" + ); +} + +#[test] +fn a_page_break_alone_never_becomes_a_piece() { + let text = format!("First page.\n{PAGE_BREAK}# Heading\nSecond page.\n"); + let each = pieces(&text, 0, 10_000); + assert_eq!(each.len(), 2, "{each:?}"); + assert_eq!(each[1].pages, Some((2, 2))); + assert_eq!(each[1].section.as_deref(), Some("Heading")); +} + +#[test] +fn an_oversized_unit_is_cut_at_paragraphs_then_lines_then_characters() { + let paragraphs = format!("{}\n\n", "p".repeat(80)).repeat(10); + for piece in pieces(¶graphs, 0, 300) { + assert!(escaped_len(piece.text) + 10 <= 300); + assert!(piece.text.ends_with("\n\n"), "cut after a blank line"); + } + let lines = format!("{}\n", "l".repeat(80)).repeat(10); + for piece in pieces(&lines, 0, 300) { + assert!(escaped_len(piece.text) + 10 <= 300); + assert!(piece.text.ends_with('\n'), "cut after a line end"); + } + let one_line = "é".repeat(1000); + let cut = pieces(&one_line, 0, 300); + assert!(cut.len() > 1); + for piece in &cut { + assert!(escaped_len(piece.text) + 10 <= 300); + } +} + +#[test] +fn sizes_count_json_escaping() { + assert_eq!(escaped_len("ab"), 2); + assert_eq!(escaped_len("\"\\\n"), 6); + assert_eq!(escaped_len("\u{1}"), 6); + assert_eq!(escaped_len("é"), 2); + let quotes = "\"".repeat(400); + for piece in pieces("es, 0, 300) { + assert!( + escaped_len(piece.text) + 10 <= 300, + "escaped, not raw, size" + ); + } +} + +#[test] +fn heading_lines_are_markdown_atx_headings_only() { + assert_eq!(heading("# Title\nbody").as_deref(), Some("Title")); + assert_eq!(heading(" ### Deep ###").as_deref(), Some("Deep")); + assert_eq!(heading("#hashtag"), None); + assert_eq!(heading("####### seven"), None); + assert_eq!(heading(" # indented code"), None); + assert_eq!(heading("# "), None); + let long = format!("# {}", "t".repeat(500)); + assert_eq!(heading(&long).unwrap().chars().count(), MAX_SECTION_CHARS); +} diff --git a/crates/tinymemory-integrations/src/cortex/envelope/mod.rs b/crates/tinymemory-integrations/src/cortex/envelope/mod.rs index 4459bbab..23172c73 100644 --- a/crates/tinymemory-integrations/src/cortex/envelope/mod.rs +++ b/crates/tinymemory-integrations/src/cortex/envelope/mod.rs @@ -23,8 +23,16 @@ //! //! # Events //! -//! A document or a learning is one event. A conversation is one event per -//! turn, appended in order. Each event's `content.text` is a JSON +//! A learning is one event. A conversation is one event per turn, appended +//! in order. A document is one event, or, when its envelope would be too big +//! for one, one event per piece of its body ([`chunks`]), appended in order; +//! each piece carries its index, and the pages and section it covers, and +//! the pieces concatenate back to the body. +//! +//! 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 @@ -41,6 +49,7 @@ //! an identical item again resolves to the same id and is detected as a //! replay by looking that id's label up. +pub(crate) mod chunks; pub(crate) mod labels; mod rebuild; @@ -128,8 +137,32 @@ pub(crate) struct Envelope { /// Which turn of a conversation this event is. #[serde(default, skip_serializing_if = "Option::is_none")] pub(crate) turn: Option<TurnInfo>, + /// Which piece of a chunked document this event is; absent for a + /// document written as one event. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub(crate) chunk: Option<ChunkInfo>, +} + +/// A piece of a chunked document: its place, and what it covers. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub(crate) struct ChunkInfo { + /// Zero-based position. + pub(crate) index: u32, + /// How many pieces the document has. + pub(crate) count: u32, + /// The first and last page the piece covers, counted from 1, when the + /// document marks its pages. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub(crate) pages: Option<[u32; 2]>, + /// The title of the section the piece starts in, when known. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub(crate) section: Option<String>, } +/// Room left in a piece's envelope for its section title: the longest title +/// at the worst JSON escaping, plus its key. +const SECTION_RESERVE: usize = chunks::MAX_SECTION_CHARS * 6 + 32; + /// A conversation turn's place and attributes. #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] pub(crate) struct TurnInfo { @@ -162,9 +195,20 @@ impl Envelope { confidence: None, evidence: None, turn: None, + chunk: None, } } + /// The position of this event within its item: a conversation turn's + /// index, a document piece's index, or `None` for an item written as + /// one event. + pub(crate) fn part(&self) -> Option<u32> { + self.turn + .as_ref() + .map(|turn| turn.index) + .or_else(|| self.chunk.as_ref().map(|chunk| chunk.index)) + } + /// The envelopes `item` is written as, one per event, in write order. /// /// # Errors @@ -184,10 +228,36 @@ impl Envelope { "document body is an unresolved uri".to_string(), )); }; - let mut envelope = Self::new(id, ItemKind::Document, text.clone(), meta); - envelope.title.clone_from(title); - envelope.mime.clone_from(mime); - Ok(vec![envelope]) + let mut whole = Self::new(id, ItemKind::Document, String::new(), meta); + whole.title.clone_from(title); + whole.mime.clone_from(mime); + let pieces = chunks::split( + text, + whole.piece_overhead()?, + chunks::DOCUMENT_CHUNK_TARGET_BYTES, + chunks::MAX_EVENT_TEXT_BYTES, + ); + if pieces.len() <= 1 { + whole.text.clone_from(text); + return Ok(vec![whole]); + } + let count = u32::try_from(pieces.len()).map_err(|_| { + Error::InvalidRequest("document has too many pieces".to_string()) + })?; + Ok((0..count) + .zip(pieces) + .map(|(index, piece)| { + let mut envelope = whole.clone(); + envelope.text = piece.text.to_string(); + envelope.chunk = Some(ChunkInfo { + index, + count, + pages: piece.pages.map(|(first, last)| [first, last]), + section: piece.section, + }); + envelope + }) + .collect()) } StoreItem::Learning { text, @@ -251,6 +321,42 @@ impl Envelope { .map_err(|_| Error::Engine("an item envelope could not be serialised".to_string())) } + /// 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. + /// + /// # 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 { + 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(), + chunks::MAX_EVENT_TEXT_BYTES + ))); + } + Ok(encoded) + } + + /// The encoded size of this envelope as a document piece with an empty + /// text: the room every piece's own text is added to. + fn piece_overhead(&self) -> Result<usize> { + let mut probe = self.clone(); + probe.chunk = Some(ChunkInfo { + index: u32::MAX, + count: u32::MAX, + pages: Some([u32::MAX, u32::MAX]), + section: None, + }); + 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 { diff --git a/crates/tinymemory-integrations/src/cortex/envelope/rebuild.rs b/crates/tinymemory-integrations/src/cortex/envelope/rebuild.rs index c5709ae0..e270c78c 100644 --- a/crates/tinymemory-integrations/src/cortex/envelope/rebuild.rs +++ b/crates/tinymemory-integrations/src/cortex/envelope/rebuild.rs @@ -25,7 +25,10 @@ pub(crate) fn decode_event(event: &Value) -> Option<Decoded> { /// The item a set of one item's envelopes describes. /// -/// A document or learning takes the first envelope. A conversation orders +/// A learning, or a document written as one event, takes the first +/// envelope. A chunked document orders its pieces by index, keeps one per +/// index and concatenates their text, so the full set gives back the body +/// exactly (and one piece alone gives that piece). A conversation orders /// its turns by index and keeps one envelope per index, so a duplicated or /// re-written turn does not repeat; turns that were never written (a store /// that failed part-way) are simply absent. `None` for an empty set. @@ -34,7 +37,7 @@ pub(crate) fn rebuild(envelopes: &[Envelope]) -> Option<StoreItem> { Some(match first.kind { ItemKind::Document => StoreItem::Document { title: first.title.clone(), - body: DocumentBody::Text(first.text.clone()), + body: DocumentBody::Text(document_text(envelopes)), mime: first.mime.clone(), meta: first.meta.clone(), }, @@ -70,3 +73,21 @@ pub(crate) fn rebuild(envelopes: &[Envelope]) -> Option<StoreItem> { } }) } + +/// A document's text from its envelopes: the first one's, or the pieces of +/// a chunked document in index order, each once. +fn document_text(envelopes: &[Envelope]) -> String { + 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/log.rs b/crates/tinymemory-integrations/src/cortex/testing/log.rs index 17d44b5a..ff74d1af 100644 --- a/crates/tinymemory-integrations/src/cortex/testing/log.rs +++ b/crates/tinymemory-integrations/src/cortex/testing/log.rs @@ -57,10 +57,14 @@ fn str_of<'a>(value: &'a Value, pointer: &str) -> &'a str { } impl CortexLog { - /// `POST /v1/experience`: (status, body). + /// `POST /v1/experience`: (status, body). Like CortexDB 0.10.4, an + /// event over 1 MiB of text is refused (`422 INVALID_ENVELOPE`). pub(crate) fn append(&mut self, body: &Value) -> (u16, Value) { let key = str_of(body, "/idempotency_key").to_string(); let text = str_of(body, "/content/text").to_string(); + if text.len() > 1024 * 1024 { + return (422, json!({ "error_code": "INVALID_ENVELOPE" })); + } if let Some((seen, id)) = self.idempotency.get(&key) { if seen != &text { return (409, json!({ "error_code": "IDEMPOTENCY_CONFLICT" })); diff --git a/crates/tinymemory-integrations/src/documents/README.md b/crates/tinymemory-integrations/src/documents/README.md index 56ed00c1..545d3045 100644 --- a/crates/tinymemory-integrations/src/documents/README.md +++ b/crates/tinymemory-integrations/src/documents/README.md @@ -68,7 +68,9 @@ let chain = ConverterChain::default().prepend(Box::new(MyPdfConverter)); ``` The `documents-office` feature ships one such binding, `OfficeConverter`: PDF (text -layer only — a scanned PDF is refused as having no text), DOCX, PPTX (slides +layer only — a scanned PDF is refused as having no text; pages are extracted +one by one, normalized on their own and joined by `PAGE_BREAK`, a form feed, +so a reader can number them), DOCX, PPTX (slides in numeric order) and XLSX (one `sheet | cell | cell` line per row), all pure Rust. It refuses hostile input rather than allocating for it: an archive whose declared uncompressed size exceeds `MAX_DECOMPRESSED_BYTES` (64 MiB), and a diff --git a/crates/tinymemory-integrations/src/documents/mod.rs b/crates/tinymemory-integrations/src/documents/mod.rs index ac1f84a4..cb1db49f 100644 --- a/crates/tinymemory-integrations/src/documents/mod.rs +++ b/crates/tinymemory-integrations/src/documents/mod.rs @@ -54,6 +54,12 @@ pub mod language; #[cfg(feature = "documents-office")] pub mod office; +/// The character converted text puts between two pages of a paged source +/// (a PDF): a form feed, one per page boundary, so page `n` follows the +/// `n - 1`th. Readers that split a document (the CortexDB engine writes a +/// long one as several events) number pages by it. +pub const PAGE_BREAK: char = '\u{c}'; + pub use convert::{ ConvertedDocument, ConverterChain, DocumentConverter, MAX_DOCUMENT_BYTES, NativeConverter, RawDocument, check_size, markdown_from_text, diff --git a/crates/tinymemory-integrations/src/documents/office/mod.rs b/crates/tinymemory-integrations/src/documents/office/mod.rs index e529e5be..b3865b08 100644 --- a/crates/tinymemory-integrations/src/documents/office/mod.rs +++ b/crates/tinymemory-integrations/src/documents/office/mod.rs @@ -13,7 +13,7 @@ //! //! | Format | Reader | Markdown it produces | //! | --- | --- | --- | -//! | PDF | `pdf-extract`, text layer only | the page text, whitespace-normalized | +//! | PDF | `pdf-extract`, text layer only | each page's text, whitespace-normalized, pages joined by [`PAGE_BREAK`] | //! | DOCX | `zip` + `quick-xml` over `word/document.xml` | one paragraph per `w:p` | //! | PPTX | `zip` + `quick-xml` over `ppt/slides/slideN.xml` | slides in numeric order, one paragraph per `a:p` | //! | XLSX | `calamine` | one `sheet \| cell \| cell` line per non-empty row | @@ -56,6 +56,7 @@ mod xlsx; use async_trait::async_trait; +use crate::documents::PAGE_BREAK; #[cfg(test)] use crate::documents::convert::MAX_DOCUMENT_BYTES; use crate::documents::convert::{ConvertedDocument, DocumentConverter, RawDocument, check_size}; @@ -96,19 +97,18 @@ impl OfficeConverter { check_size(document)?; let format = document.format(); let bytes = document.bytes.as_slice(); - let text = match format { - DocumentFormat::Pdf => pdf::extract(bytes)?, - DocumentFormat::Docx => ooxml::docx(bytes)?, - DocumentFormat::Pptx => ooxml::pptx(bytes)?, - DocumentFormat::Xlsx => xlsx::extract(bytes)?, + let markdown = match format { + DocumentFormat::Pdf => paged(&pdf::extract(bytes)?), + DocumentFormat::Docx => normalize::normalize(&ooxml::docx(bytes)?), + DocumentFormat::Pptx => normalize::normalize(&ooxml::pptx(bytes)?), + DocumentFormat::Xlsx => normalize::normalize(&xlsx::extract(bytes)?), other => { return Err(Error::UnsupportedFormat(format!( "the office converter does not handle {other}" ))); } }; - let markdown = normalize::normalize(&text); - if markdown.is_empty() { + if markdown.trim().is_empty() { return Err(unreadable(format!("converting {format} produced no text"))); } Ok(ConvertedDocument::new(markdown, format, bytes.len()) @@ -116,6 +116,17 @@ impl OfficeConverter { } } +/// A PDF's pages as one markdown text: each page normalized on its own, +/// joined by [`PAGE_BREAK`] so a reader can still tell where each page +/// starts (an empty page keeps its place, so numbering survives). +fn paged(pages: &[String]) -> String { + pages + .iter() + .map(|page| normalize::normalize(page)) + .collect::<Vec<_>>() + .join(&PAGE_BREAK.to_string()) +} + #[async_trait] impl DocumentConverter for OfficeConverter { fn name(&self) -> &str { diff --git a/crates/tinymemory-integrations/src/documents/office/mod_tests.rs b/crates/tinymemory-integrations/src/documents/office/mod_tests.rs index b41a19c9..a740d53b 100644 --- a/crates/tinymemory-integrations/src/documents/office/mod_tests.rs +++ b/crates/tinymemory-integrations/src/documents/office/mod_tests.rs @@ -107,18 +107,35 @@ fn xlsx_with_sheet(sheet_data: &str) -> Vec<u8> { /// cross-reference table so the parser takes the normal path rather than a /// recovery one. fn pdf(content: &str) -> Vec<u8> { - let objects = [ + pdf_pages(&[content]) +} + +/// A PDF with one page per content stream, in order. +fn pdf_pages(contents: &[&str]) -> Vec<u8> { + // 1 catalog, 2 pages, 3 font, then each page and its content stream. + let kids: Vec<String> = (0..contents.len()) + .map(|index| format!("{} 0 R", 4 + 2 * index)) + .collect(); + let mut objects = vec![ "<< /Type /Catalog /Pages 2 0 R >>".to_string(), - "<< /Type /Pages /Kids [3 0 R] /Count 1 >>".to_string(), - "<< /Type /Page /Parent 2 0 R /MediaBox [0 0 612 792] /Contents 4 0 R \ - /Resources << /Font << /F1 5 0 R >> >> >>" - .to_string(), format!( - "<< /Length {} >>\nstream\n{content}\nendstream", - content.len() + "<< /Type /Pages /Kids [{}] /Count {} >>", + kids.join(" "), + contents.len() ), "<< /Type /Font /Subtype /Type1 /BaseFont /Helvetica >>".to_string(), ]; + for (index, content) in contents.iter().enumerate() { + objects.push(format!( + "<< /Type /Page /Parent 2 0 R /MediaBox [0 0 612 792] /Contents {} 0 R \ + /Resources << /Font << /F1 3 0 R >> >> >>", + 5 + 2 * index + )); + objects.push(format!( + "<< /Length {} >>\nstream\n{content}\nendstream", + content.len() + )); + } let mut out = String::from("%PDF-1.4\n"); let mut offsets = Vec::new(); for (index, object) in objects.iter().enumerate() { @@ -272,6 +289,29 @@ async fn a_pdf_yields_its_text_layer() { assert_eq!(converted.format, DocumentFormat::Pdf); } +#[tokio::test] +async fn a_pdf_keeps_its_page_boundaries_as_page_breaks() { + let bytes = pdf_pages(&[ + "BT /F1 24 Tf 72 720 Td (Refunds take five days) Tj ET", + "BT /F1 24 Tf 72 720 Td (Shipping is free over fifty) Tj ET", + "", + "BT /F1 24 Tf 72 720 Td (Returns need a receipt) Tj ET", + ]); + let converted = convert("policy.pdf", bytes).await.unwrap(); + let pages: Vec<&str> = converted + .markdown + .split(crate::documents::PAGE_BREAK) + .collect(); + assert_eq!(pages.len(), 4, "{:?}", converted.markdown); + assert!(pages[0].contains("Refunds take five days"), "{pages:?}"); + assert!( + pages[1].contains("Shipping is free over fifty"), + "{pages:?}" + ); + assert!(pages[2].trim().is_empty(), "an empty page keeps its place"); + assert!(pages[3].contains("Returns need a receipt"), "{pages:?}"); +} + #[tokio::test] async fn a_pdf_without_a_text_layer_says_so_rather_than_storing_nothing() { // A scanned PDF carries pictures of words. That is not a parse failure, diff --git a/crates/tinymemory-integrations/src/documents/office/pdf.rs b/crates/tinymemory-integrations/src/documents/office/pdf.rs index 6481b584..d5e1c64f 100644 --- a/crates/tinymemory-integrations/src/documents/office/pdf.rs +++ b/crates/tinymemory-integrations/src/documents/office/pdf.rs @@ -1,18 +1,19 @@ -//! A PDF's text layer. +//! A PDF's text layer, page by page. use super::unreadable; use crate::documents::error::Result; -/// Extracts a PDF's text layer, which may be empty. +/// Extracts a PDF's text layer, one string per page, any of which may be +/// empty. /// /// A scanned PDF has none: the file was read, it simply carries pictures of -/// words. That comes back as empty text, and the converter reports it as a +/// words. That comes back as empty pages, and the converter reports it as a /// document with no text rather than as a parse failure. -pub(super) fn extract(bytes: &[u8]) -> Result<String> { +pub(super) fn extract(bytes: &[u8]) -> Result<Vec<String>> { // `pdf-extract` panics on some malformed documents rather than erroring. // Caught so one bad file is one refused document, not a crashed task. - match std::panic::catch_unwind(|| pdf_extract::extract_text_from_mem(bytes)) { - Ok(Ok(text)) => Ok(text), + match std::panic::catch_unwind(|| pdf_extract::extract_text_from_mem_by_pages(bytes)) { + Ok(Ok(pages)) => Ok(pages), Ok(Err(error)) => Err(unreadable(format!("the PDF could not be read: {error}"))), Err(_) => Err(unreadable( "the PDF is malformed enough that the parser gave up on it".to_string(), diff --git a/docs/architecture/cortex-flows.md b/docs/architecture/cortex-flows.md index 1b403d14..eef2b0d5 100644 --- a/docs/architecture/cortex-flows.md +++ b/docs/architecture/cortex-flows.md @@ -110,10 +110,10 @@ in `ItemKind::ALL` order and then by namespace, each newest first. event id (the engine emits each event twice; the cursor remembers the last id so this works across page boundaries), decode it, and keep it when it is an envelope of the scope's kind and the **full** `MetaFilter` matches. -4. **Each item once.** A document or learning is one event. A conversation is - emitted only on the page holding its **turn 0** event; its text is - assembled from all its turns by one label lookup for all the conversations - on the page. A conversation whose store failed part way still has turn 0 and +4. **Each item once.** A learning, or a document written whole, is one + event. A conversation, or a chunked document, is emitted only on the page + holding its first event (**turn 0**, **piece 0**); its text is assembled + from all its events by one label lookup per kind for the page. A conversation whose store failed part way still has turn 0 and lists with the turns it holds. 5. Stop when `limit` items are collected and return a cursor, unless the end of the last scope was reached (then there is none). diff --git a/docs/architecture/cortex-wire.md b/docs/architecture/cortex-wire.md index 1de3bc1a..92c134c2 100644 --- a/docs/architecture/cortex-wire.md +++ b/docs/architecture/cortex-wire.md @@ -345,10 +345,47 @@ 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.** 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. +- **Where.** First at page breaks (the form feed the PDF converter puts + between pages), then before markdown heading lines; a stretch of only + whitespace never becomes a piece. 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). +- **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. +- **Reads.** `get` and `list` give the whole document (pieces in index + order). A ranked hit or citation on a piece gives that piece, with the + item's id and its metadata plus `page:<n>` (or `page:<first>-<last>`) and + `section:<title>` tags. 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 document or learning is one event; a conversation is one event per turn, -appended in order. CortexDB's experience schema is closed (an unknown field is +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)). +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**: @@ -371,6 +408,7 @@ a 422), so the structured data rides in the one free-form field: the event's | `title`, `mime` | documents, when set | | | `learning_kind`, `confidence`, `evidence` | learnings (`evidence` when set) | | | `turn` | conversation turns | `index` (0-based), `count`, `role`, `at`, `tool_calls` | +| `chunk` | pieces of a chunked document | `index` (0-based), `count`, `pages` (`[first, last]`, when the text marks pages), `section` (the heading the piece starts under) | 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 @@ -378,8 +416,11 @@ 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 refused at write time as `Error::InvalidRequest`. -Rebuilding an item from envelopes: a document or learning takes the first -envelope; a conversation orders turns by `index` and keeps one per index (so a +Rebuilding an item from envelopes: a learning, or a document written whole, +takes the first envelope; a chunked document orders its pieces by `index`, +keeps one per index and concatenates their text (the pieces are contiguous +slices, so all of them give back the body exactly, and one gives that piece); +a conversation orders turns by `index` and keeps one per index (so a duplicated or re-written turn does not repeat, and a turn that was never written is absent). A learning with no `learning_kind` reads back as `Other`. diff --git a/docs/specs/memory-v2.md b/docs/specs/memory-v2.md index dd8de9e4..5aa293e6 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. + - 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 (`chunk {index, count, pages?, section?}` in the envelope). No event over 768 KiB of encoded envelope is sent (CortexDB refuses one over 1 MiB); `get`/`list` reassemble a chunked document, and a ranked hit on a piece carries that 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 4eba98cb04d5c5d1df2d934b2cba0dd97d59c8b8 Mon Sep 17 00:00:00 2001 From: M3gA-Mind <elvin@tinyhumans.ai> Date: Tue, 6 Oct 2026 20:56:41 +0530 Subject: [PATCH 2/8] Guard the splitter's room, test reassembly, and state tag optionality - chunks::split returns None when the envelope overhead leaves less than one escaped character (6 bytes) of room, and never packs below that, so no piece can exceed the limit; a document with no room for a piece stays whole and is refused by encode_checked if it does not fit, never cut into pieces that each exceed it. - Tests: rebuilding a chunked document from shuffled, duplicated and single pieces; a document written whole has no chunk field; oversized metadata; the room boundary; a multi-page PDF with no text is refused. - The test double checks a reused idempotency key (409) before the size limit (422), as its contract says. - Docs: page and section tags are present only when the document marks pages or the piece starts under a heading, page tags may be ranges, and oversized pages or sections are cut again at blank lines, lines, then characters. --- .../src/cortex/README.md | 4 +- .../src/cortex/engine/mod_chunk_tests.rs | 22 +++++ .../src/cortex/envelope/chunks.rs | 28 ++++-- .../src/cortex/envelope/chunks_tests.rs | 28 +++++- .../src/cortex/envelope/mod.rs | 7 +- .../src/cortex/envelope/mod_tests.rs | 92 +++++++++++++++++++ .../src/cortex/envelope/rebuild.rs | 5 +- .../src/cortex/testing/log.rs | 6 +- .../src/documents/office/mod_tests.rs | 3 + docs/architecture/cortex-wire.md | 8 +- docs/specs/memory-v2.md | 2 +- 11 files changed, 187 insertions(+), 18 deletions(-) diff --git a/crates/tinymemory-integrations/src/cortex/README.md b/crates/tinymemory-integrations/src/cortex/README.md index 97c30054..ded16947 100644 --- a/crates/tinymemory-integrations/src/cortex/README.md +++ b/crates/tinymemory-integrations/src/cortex/README.md @@ -118,7 +118,9 @@ packed up to that size (`envelope/chunks.rs`). No event is sent over 768 KiB of encoded envelope (CortexDB refuses an experience over 1 MiB); an item that cannot fit is refused before anything of its batch is sent. `get` and `list` reassemble a chunked document; a `fetch` or `recall` hit on a piece carries -that piece, tagged `page:<n>` and `section:<title>`. Each event's +that 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: ```json 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 68e07a67..feea7812 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/mod_chunk_tests.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/mod_chunk_tests.rs @@ -253,3 +253,25 @@ async fn an_event_over_the_limit_is_refused_before_anything_is_sent() { assert_eq!(state.count("POST"), writes_before, "nothing was written"); } } + +#[test] +fn the_double_answers_a_reused_key_as_a_conflict_before_checking_size() { + let mut log = crate::cortex::testing::CortexLog::default(); + let event = |text: String| { + json!({ + "scope": SCOPE, + "modality": "document", + "idempotency_key": "k", + "content": { "kind": "message", "role": "user", "text": text }, + }) + }; + assert_eq!(log.append(&event("small".into())).0, 202); + let (status, body) = log.append(&event("x".repeat(2 * 1024 * 1024))); + assert_eq!(status, 409, "{body}"); + let (status, body) = log.append(&json!({ + "scope": SCOPE, + "idempotency_key": "other", + "content": { "text": "x".repeat(2 * 1024 * 1024) }, + })); + assert_eq!(status, 422, "{body}"); +} diff --git a/crates/tinymemory-integrations/src/cortex/envelope/chunks.rs b/crates/tinymemory-integrations/src/cortex/envelope/chunks.rs index 81423d50..7cbb44f7 100644 --- a/crates/tinymemory-integrations/src/cortex/envelope/chunks.rs +++ b/crates/tinymemory-integrations/src/cortex/envelope/chunks.rs @@ -63,23 +63,39 @@ struct Unit { section: Option<String>, } +/// The most bytes one character takes JSON-escaped (`\u0001`): the least +/// room a piece needs to make progress. +const MAX_ESCAPED_CHAR: usize = 6; + /// The pieces of `text`, each at most `target` (or, at `0`, one per unit) /// and never over `limit`, both as encoded bytes on top of `overhead` (the /// envelope around the piece). One piece for a text that fits. -pub(crate) fn split(text: &str, overhead: usize, target: usize, limit: usize) -> Vec<Piece<'_>> { - let room = limit.saturating_sub(overhead).max(1); +/// +/// `None` when `overhead` leaves less than one escaped character of room +/// under `limit`: no piece could be written, so none is made up. +pub(crate) fn split( + text: &str, + overhead: usize, + target: usize, + limit: usize, +) -> Option<Vec<Piece<'_>>> { + let room = limit + .checked_sub(overhead) + .filter(|room| *room >= MAX_ESCAPED_CHAR)?; let pack = if target == 0 { room } else { - target.saturating_sub(overhead).clamp(1, room) + target + .saturating_sub(overhead) + .clamp(MAX_ESCAPED_CHAR, room) }; let paged = text.contains(PAGE_BREAK); if target != 0 && escaped_len(text) <= pack { - return vec![Piece { + return Some(vec![Piece { text, pages: paged.then(|| (1, page_count(text))), section: None, - }]; + }]); } let mut pieces: Vec<Piece<'_>> = Vec::new(); let mut open: Option<Open> = None; @@ -109,7 +125,7 @@ pub(crate) fn split(text: &str, overhead: usize, target: usize, limit: usize) -> } } pieces.extend(open.map(|done| done.piece(text, paged))); - pieces + Some(pieces) } /// The piece being packed. diff --git a/crates/tinymemory-integrations/src/cortex/envelope/chunks_tests.rs b/crates/tinymemory-integrations/src/cortex/envelope/chunks_tests.rs index 0e971d6e..aa31ce50 100644 --- a/crates/tinymemory-integrations/src/cortex/envelope/chunks_tests.rs +++ b/crates/tinymemory-integrations/src/cortex/envelope/chunks_tests.rs @@ -5,7 +5,7 @@ use super::*; /// The pieces of `text` under a small overhead and the given sizes, /// checked to concatenate back to `text`. fn pieces(text: &str, target: usize, limit: usize) -> Vec<Piece<'_>> { - let pieces = split(text, 10, target, limit); + let pieces = split(text, 10, target, limit).expect("room for a piece"); let joined: String = pieces.iter().map(|piece| piece.text).collect(); assert_eq!(joined, text, "the pieces concatenate to the text"); pieces @@ -141,3 +141,29 @@ fn heading_lines_are_markdown_atx_headings_only() { let long = format!("# {}", "t".repeat(500)); assert_eq!(heading(&long).unwrap().chars().count(), MAX_SECTION_CHARS); } + +#[test] +fn no_room_for_one_escaped_character_makes_no_piece() { + assert_eq!(split("a", 10, 0, 10), None, "overhead equals the limit"); + assert_eq!(split("a", 10, 0, 15), None, "five bytes of room"); + assert_eq!(split("", 20, 0, 10), None, "overhead over the limit"); + assert!(split("a", 10, 0, 16).is_some(), "six bytes of room"); +} + +#[test] +fn every_piece_fits_even_at_the_smallest_room() { + for text in ["\u{1}\u{1}\u{1}", "\"\"\"\"", "ééé", "a\nb\nc"] { + for target in [0, 1, 7] { + let cut = split(text, 10, target, 16).expect("six bytes of room"); + let joined: String = cut.iter().map(|piece| piece.text).collect(); + assert_eq!(joined, text); + for piece in &cut { + assert!( + 10 + escaped_len(piece.text) <= 16, + "{text:?} at target {target}: {:?}", + piece.text + ); + } + } + } +} diff --git a/crates/tinymemory-integrations/src/cortex/envelope/mod.rs b/crates/tinymemory-integrations/src/cortex/envelope/mod.rs index 23172c73..8d75cfcc 100644 --- a/crates/tinymemory-integrations/src/cortex/envelope/mod.rs +++ b/crates/tinymemory-integrations/src/cortex/envelope/mod.rs @@ -231,12 +231,17 @@ impl Envelope { let mut whole = Self::new(id, ItemKind::Document, String::new(), meta); whole.title.clone_from(title); whole.mime.clone_from(mime); + // No room for a piece (metadata alone near the limit) leaves + // the document whole: it is written if it fits and refused + // by `encode_checked` if not, never cut into pieces that are + // each over the limit. let pieces = chunks::split( text, whole.piece_overhead()?, chunks::DOCUMENT_CHUNK_TARGET_BYTES, chunks::MAX_EVENT_TEXT_BYTES, - ); + ) + .unwrap_or_default(); if pieces.len() <= 1 { whole.text.clone_from(text); return Ok(vec![whole]); diff --git a/crates/tinymemory-integrations/src/cortex/envelope/mod_tests.rs b/crates/tinymemory-integrations/src/cortex/envelope/mod_tests.rs index 128b7606..418d173d 100644 --- a/crates/tinymemory-integrations/src/cortex/envelope/mod_tests.rs +++ b/crates/tinymemory-integrations/src/cortex/envelope/mod_tests.rs @@ -189,3 +189,95 @@ fn an_item_is_written_to_its_namespace_scope() { "app:tinymemory/agent:researcher/app:documents" ); } + +/// A document long enough to be written as several pieces. +fn long_document() -> StoreItem { + let body: String = (1..=30) + .map(|n| { + format!( + "# Part {n}\n\n{}\n", + "Body text for this part. ".repeat(600) + ) + }) + .collect(); + StoreItem::Document { + title: Some("Long".into()), + body: DocumentBody::Text(body), + mime: None, + meta: meta(), + } +} + +#[test] +fn a_chunked_document_rebuilds_from_pieces_in_any_order_each_once() { + let item = long_document(); + let id = item.fingerprint(); + let mut pieces = Envelope::for_item(&item, &id).unwrap(); + assert!(pieces.len() >= 2, "{} pieces", pieces.len()); + assert!(pieces.iter().all(|piece| piece.chunk.is_some())); + + pieces.reverse(); + let duplicated = pieces[1].clone(); + pieces.push(duplicated); + let rebuilt = rebuild(&pieces).unwrap(); + assert_eq!( + rebuilt, item, + "shuffled and duplicated pieces give the body" + ); + assert_eq!(rebuilt.fingerprint(), id); + + let one = rebuild(std::slice::from_ref(&pieces[0])).unwrap(); + let StoreItem::Document { + body: DocumentBody::Text(text), + .. + } = one + else { + panic!("a document"); + }; + assert_eq!(text, pieces[0].text, "one piece alone gives that piece"); +} + +#[test] +fn a_document_without_pieces_rebuilds_from_its_first_envelope() { + let item = StoreItem::document("Short note.", meta()); + let id = item.fingerprint(); + let envelopes = Envelope::for_item(&item, &id).unwrap(); + assert_eq!(envelopes.len(), 1); + assert_eq!(envelopes[0].chunk, None); + assert!(!envelopes[0].encode().unwrap().contains("\"chunk\"")); + assert_eq!(rebuild(&envelopes).unwrap(), item); +} + +#[test] +fn metadata_too_large_for_a_piece_keeps_the_document_whole_and_refused() { + let mut huge = meta(); + huge.tags = vec!["t".repeat(chunks::MAX_EVENT_TEXT_BYTES)]; + let small = StoreItem::document("Short note.", huge.clone()); + let id = small.fingerprint(); + let envelopes = Envelope::for_item(&small, &id).unwrap(); + assert_eq!(envelopes.len(), 1, "no pieces are made up"); + assert!(matches!( + envelopes[0].encode_checked(), + Err(Error::InvalidRequest(_)) + )); + let StoreItem::Document { + body, title, mime, .. + } = long_document() + else { + unreachable!("a document"); + }; + let long = StoreItem::Document { + title, + body, + mime, + meta: huge, + }; + let id = long.fingerprint(); + let envelopes = Envelope::for_item(&long, &id).unwrap(); + assert_eq!( + envelopes.len(), + 1, + "never pieces that each exceed the limit" + ); + assert!(envelopes[0].encode_checked().is_err()); +} diff --git a/crates/tinymemory-integrations/src/cortex/envelope/rebuild.rs b/crates/tinymemory-integrations/src/cortex/envelope/rebuild.rs index e270c78c..b9ffbcbd 100644 --- a/crates/tinymemory-integrations/src/cortex/envelope/rebuild.rs +++ b/crates/tinymemory-integrations/src/cortex/envelope/rebuild.rs @@ -25,8 +25,9 @@ pub(crate) fn decode_event(event: &Value) -> Option<Decoded> { /// The item a set of one item's envelopes describes. /// -/// A learning, or a document written as one event, takes the first -/// envelope. A chunked document orders its pieces by index, keeps one per +/// A learning, or a document written as one event (no `chunk` field, as +/// every document was before chunking), takes the first envelope. A chunked +/// document (every piece carries `chunk`) orders its pieces by index, keeps one per /// index and concatenates their text, so the full set gives back the body /// exactly (and one piece alone gives that piece). A conversation orders /// its turns by index and keeps one envelope per index, so a duplicated or diff --git a/crates/tinymemory-integrations/src/cortex/testing/log.rs b/crates/tinymemory-integrations/src/cortex/testing/log.rs index ff74d1af..583db306 100644 --- a/crates/tinymemory-integrations/src/cortex/testing/log.rs +++ b/crates/tinymemory-integrations/src/cortex/testing/log.rs @@ -62,9 +62,6 @@ impl CortexLog { pub(crate) fn append(&mut self, body: &Value) -> (u16, Value) { let key = str_of(body, "/idempotency_key").to_string(); let text = str_of(body, "/content/text").to_string(); - if text.len() > 1024 * 1024 { - return (422, json!({ "error_code": "INVALID_ENVELOPE" })); - } if let Some((seen, id)) = self.idempotency.get(&key) { if seen != &text { return (409, json!({ "error_code": "IDEMPOTENCY_CONFLICT" })); @@ -74,6 +71,9 @@ impl CortexLog { json!({ "event_id": id, "replayed_from_idempotency": true }), ); } + if text.len() > 1024 * 1024 { + return (422, json!({ "error_code": "INVALID_ENVELOPE" })); + } self.next_id += 1; let id = format!("evt_{}", self.next_id); self.idempotency.insert(key, (text, id.clone())); diff --git a/crates/tinymemory-integrations/src/documents/office/mod_tests.rs b/crates/tinymemory-integrations/src/documents/office/mod_tests.rs index a740d53b..b696e30e 100644 --- a/crates/tinymemory-integrations/src/documents/office/mod_tests.rs +++ b/crates/tinymemory-integrations/src/documents/office/mod_tests.rs @@ -319,6 +319,9 @@ async fn a_pdf_without_a_text_layer_says_so_rather_than_storing_nothing() { // succeeded. let reason = refusal("scan.pdf", pdf("")).await; assert!(reason.contains("no text"), "{reason}"); + // Several empty pages are joined by page breaks, which are not text. + let reason = refusal("scan.pdf", pdf_pages(&["", "", ""])).await; + assert!(reason.contains("no text"), "{reason}"); } #[tokio::test] diff --git a/docs/architecture/cortex-wire.md b/docs/architecture/cortex-wire.md index 92c134c2..108235a3 100644 --- a/docs/architecture/cortex-wire.md +++ b/docs/architecture/cortex-wire.md @@ -371,8 +371,10 @@ its whole envelope. So (`envelope/chunks.rs`): store that failed part way writes only the missing pieces. - **Reads.** `get` and `list` give the whole document (pieces in index order). A ranked hit or citation on a piece gives that piece, with the - item's id and its metadata plus `page:<n>` (or `page:<first>-<last>`) and - `section:<title>` tags. Readable CortexDB labels for page and section are + 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. 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 @@ -408,7 +410,7 @@ a 422), so the structured data rides in the one free-form field: the event's | `title`, `mime` | documents, when set | | | `learning_kind`, `confidence`, `evidence` | learnings (`evidence` when set) | | | `turn` | conversation turns | `index` (0-based), `count`, `role`, `at`, `tool_calls` | -| `chunk` | pieces of a chunked document | `index` (0-based), `count`, `pages` (`[first, last]`, when the text marks pages), `section` (the heading the piece starts under) | +| `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 diff --git a/docs/specs/memory-v2.md b/docs/specs/memory-v2.md index 5aa293e6..ceb4878d 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 (`chunk {index, count, pages?, section?}` in the envelope). No event over 768 KiB of encoded envelope is sent (CortexDB refuses one over 1 MiB); `get`/`list` reassemble a chunked document, and a ranked hit on a piece carries that piece with `page:`/`section:` tags. + - 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` 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); `get`/`list` reassemble a chunked document, and a ranked hit on a piece carries that 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 8629692de0759db66675fee39fd58dac96ea524b Mon Sep 17 00:00:00 2001 From: M3gA-Mind <elvin@tinyhumans.ai> Date: Tue, 6 Oct 2026 21:08:51 +0530 Subject: [PATCH 3/8] Refuse an unsplittable oversized document up front; size overhead exactly - A document whose metadata leaves no room for a piece is laid out whole only if it fits; otherwise `for_item` refuses it with InvalidRequest instead of returning an envelope that would be refused later. - The piece overhead reserves a page range only for a document that marks pages. - When the metadata alone uses up the 256 KiB target, pieces pack up to the room under the event limit instead of collapsing to a few bytes each (a short note with large metadata was being split). - Tests: tag branches of located_meta, an unpaged headingless document gets no extra tags, a whitespace-only document keeps its text, metadata over the target, refusal of an unsplittable document. - Docs: at a zero target, a page or section too big for one event is still cut into several. --- .../src/cortex/engine/mod_chunk_tests.rs | 51 +++++++++++++++ .../src/cortex/envelope/chunks.rs | 7 +- .../src/cortex/envelope/chunks_tests.rs | 12 ++++ .../src/cortex/envelope/mod.rs | 33 ++++++---- .../src/cortex/envelope/mod_tests.rs | 64 ++++++++++--------- docs/architecture/cortex-wire.md | 3 +- 6 files changed, 127 insertions(+), 43 deletions(-) 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 feea7812..e51ac1df 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/mod_chunk_tests.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/mod_chunk_tests.rs @@ -275,3 +275,54 @@ fn the_double_answers_a_reused_key_as_a_conflict_before_checking_size() { })); assert_eq!(status, 422, "{body}"); } + +#[test] +fn a_piece_is_tagged_only_with_the_page_and_section_it_has() { + use crate::cortex::envelope::{ChunkInfo, Envelope}; + let item = StoreItem::document("piece text", meta()); + let mut envelope = Envelope::for_item(&item, &item.fingerprint()) + .unwrap() + .remove(0); + let tags = |pages: Option<[u32; 2]>, section: Option<&str>| { + let mut piece = envelope.clone(); + piece.chunk = Some(ChunkInfo { + index: 1, + count: 3, + pages, + section: section.map(str::to_string), + }); + super::items::located_meta(&piece).tags + }; + assert_eq!( + tags(Some([3, 3]), Some("Billing")), + ["page:3", "section:Billing"] + ); + assert_eq!(tags(Some([3, 5]), None), ["page:3-5"]); + assert_eq!(tags(None, Some("Billing")), ["section:Billing"]); + assert!(tags(None, None).is_empty(), "no page, no section: no tag"); + envelope.chunk = None; + assert!(super::items::located_meta(&envelope).tags.is_empty()); +} + +#[tokio::test] +async fn a_long_document_without_pages_or_headings_gets_no_extra_tags() { + let (endpoint, state) = direct_double().await; + let engine = direct_engine(&endpoint); + let body = format!("{}\n\n", "Plain paragraph about refunds. ".repeat(40)).repeat(300); + let item = StoreItem::Document { + title: None, + body: DocumentBody::Text(body.clone()), + mime: None, + meta: meta(), + }; + engine.store(item.clone()).await.unwrap(); + assert!(events(&state, SCOPE).len() > 1, "written in pieces"); + let mut request = FetchRequest::new("refunds paragraph", FetchMode::Hybrid, 3); + request.filter.reach = Some(Reach::exact(Namespace::source("pdf"))); + let page = engine.fetch(request).await.unwrap(); + assert!(!page.hits.is_empty()); + for hit in &page.hits { + assert!(hit.meta.tags.is_empty(), "{:?}", hit.meta.tags); + assert!(body.contains(&hit.text) || hit.text.len() < body.len()); + } +} diff --git a/crates/tinymemory-integrations/src/cortex/envelope/chunks.rs b/crates/tinymemory-integrations/src/cortex/envelope/chunks.rs index 7cbb44f7..38d135c4 100644 --- a/crates/tinymemory-integrations/src/cortex/envelope/chunks.rs +++ b/crates/tinymemory-integrations/src/cortex/envelope/chunks.rs @@ -82,12 +82,15 @@ pub(crate) fn split( let room = limit .checked_sub(overhead) .filter(|room| *room >= MAX_ESCAPED_CHAR)?; + // A target the metadata alone uses up gives way to the room under the + // limit, rather than packing a few bytes per piece. let pack = if target == 0 { room } else { target - .saturating_sub(overhead) - .clamp(MAX_ESCAPED_CHAR, room) + .checked_sub(overhead) + .filter(|pack| *pack >= MAX_ESCAPED_CHAR) + .map_or(room, |pack| pack.min(room)) }; let paged = text.contains(PAGE_BREAK); if target != 0 && escaped_len(text) <= pack { diff --git a/crates/tinymemory-integrations/src/cortex/envelope/chunks_tests.rs b/crates/tinymemory-integrations/src/cortex/envelope/chunks_tests.rs index aa31ce50..ebbf4c77 100644 --- a/crates/tinymemory-integrations/src/cortex/envelope/chunks_tests.rs +++ b/crates/tinymemory-integrations/src/cortex/envelope/chunks_tests.rs @@ -167,3 +167,15 @@ fn every_piece_fits_even_at_the_smallest_room() { } } } + +#[test] +fn metadata_over_the_target_packs_up_to_the_limit_not_a_few_bytes() { + let text = "A short note that fits under the limit."; + let one = split(text, 900, 500, 10_000).expect("room under the limit"); + assert_eq!(one.len(), 1, "{one:?}"); + let long = "word ".repeat(4_000); + for piece in split(&long, 900, 500, 10_000).expect("room") { + assert!(900 + escaped_len(piece.text) <= 10_000); + assert!(escaped_len(piece.text) > 1_000, "not a few bytes per piece"); + } +} diff --git a/crates/tinymemory-integrations/src/cortex/envelope/mod.rs b/crates/tinymemory-integrations/src/cortex/envelope/mod.rs index 8d75cfcc..627d1de5 100644 --- a/crates/tinymemory-integrations/src/cortex/envelope/mod.rs +++ b/crates/tinymemory-integrations/src/cortex/envelope/mod.rs @@ -231,17 +231,26 @@ impl Envelope { let mut whole = Self::new(id, ItemKind::Document, String::new(), meta); whole.title.clone_from(title); whole.mime.clone_from(mime); - // No room for a piece (metadata alone near the limit) leaves - // the document whole: it is written if it fits and refused - // by `encode_checked` if not, never cut into pieces that are - // each over the limit. - let pieces = chunks::split( + let paged = text.contains(chunks::PAGE_BREAK); + let Some(pieces) = chunks::split( text, - whole.piece_overhead()?, + whole.piece_overhead(paged)?, chunks::DOCUMENT_CHUNK_TARGET_BYTES, chunks::MAX_EVENT_TEXT_BYTES, - ) - .unwrap_or_default(); + ) else { + // The metadata leaves no room for a piece: the document + // is written whole if it fits, and refused here if not, + // never cut into pieces that would each be over the limit. + whole.text.clone_from(text); + let size = whole.encode()?.len(); + if size > chunks::MAX_EVENT_TEXT_BYTES { + return Err(Error::InvalidRequest(format!( + "a document event would be {size} bytes and its metadata leaves no \ + room to split it; CortexDB refuses an event over 1 MiB" + ))); + } + return Ok(vec![whole]); + }; if pieces.len() <= 1 { whole.text.clone_from(text); return Ok(vec![whole]); @@ -350,13 +359,15 @@ impl Envelope { } /// The encoded size of this envelope as a document piece with an empty - /// text: the room every piece's own text is added to. - fn piece_overhead(&self) -> Result<usize> { + /// text: the room every piece's own text is added to. A page range is + /// reserved only for a document that marks pages (`paged`), since only + /// its pieces carry one. + fn piece_overhead(&self, paged: bool) -> Result<usize> { let mut probe = self.clone(); probe.chunk = Some(ChunkInfo { index: u32::MAX, count: u32::MAX, - pages: Some([u32::MAX, u32::MAX]), + pages: paged.then_some([u32::MAX, u32::MAX]), section: None, }); Ok(probe.encode()?.len() + SECTION_RESERVE) diff --git a/crates/tinymemory-integrations/src/cortex/envelope/mod_tests.rs b/crates/tinymemory-integrations/src/cortex/envelope/mod_tests.rs index 418d173d..f5665fc5 100644 --- a/crates/tinymemory-integrations/src/cortex/envelope/mod_tests.rs +++ b/crates/tinymemory-integrations/src/cortex/envelope/mod_tests.rs @@ -249,35 +249,41 @@ fn a_document_without_pieces_rebuilds_from_its_first_envelope() { } #[test] -fn metadata_too_large_for_a_piece_keeps_the_document_whole_and_refused() { +fn metadata_that_leaves_no_room_for_a_piece_keeps_a_fitting_document_whole() { + // Metadata just under the limit: no room for a piece's envelope, yet a + // short document still fits whole. + let mut near = meta(); + near.tags = vec!["t".repeat(chunks::MAX_EVENT_TEXT_BYTES - 2_000)]; + let fits = StoreItem::document("Short note.", near); + let id = fits.fingerprint(); + let envelopes = Envelope::for_item(&fits, &id).unwrap(); + assert_eq!(envelopes.len(), 1, "no pieces are made up"); + assert_eq!(envelopes[0].chunk, None); + envelopes[0].encode_checked().unwrap(); +} + +#[test] +fn a_document_that_cannot_fit_or_be_split_is_refused_when_laid_out() { let mut huge = meta(); huge.tags = vec!["t".repeat(chunks::MAX_EVENT_TEXT_BYTES)]; - let small = StoreItem::document("Short note.", huge.clone()); - let id = small.fingerprint(); - let envelopes = Envelope::for_item(&small, &id).unwrap(); - assert_eq!(envelopes.len(), 1, "no pieces are made up"); - assert!(matches!( - envelopes[0].encode_checked(), - Err(Error::InvalidRequest(_)) - )); - let StoreItem::Document { - body, title, mime, .. - } = long_document() - else { - unreachable!("a document"); - }; - let long = StoreItem::Document { - title, - body, - mime, - meta: huge, - }; - let id = long.fingerprint(); - let envelopes = Envelope::for_item(&long, &id).unwrap(); - assert_eq!( - envelopes.len(), - 1, - "never pieces that each exceed the limit" - ); - assert!(envelopes[0].encode_checked().is_err()); + for body in ["Short note.".to_string(), "x".repeat(400_000)] { + let item = StoreItem::document(body, huge.clone()); + let id = item.fingerprint(); + let refused = Envelope::for_item(&item, &id); + assert!( + matches!(&refused, Err(Error::InvalidRequest(message)) if message.contains("1 MiB")), + "{refused:?}" + ); + } +} + +#[test] +fn a_whitespace_only_document_keeps_its_text() { + for body in [" \n\n ", "\u{c}\u{c}", " \u{c} \n"] { + let item = StoreItem::document(body, meta()); + let id = item.fingerprint(); + let envelopes = Envelope::for_item(&item, &id).unwrap(); + assert_eq!(envelopes.len(), 1, "{body:?}"); + assert_eq!(rebuild(&envelopes).unwrap(), item, "{body:?}"); + } } diff --git a/docs/architecture/cortex-wire.md b/docs/architecture/cortex-wire.md index 108235a3..6fb5bef1 100644 --- a/docs/architecture/cortex-wire.md +++ b/docs/architecture/cortex-wire.md @@ -355,7 +355,8 @@ its whole envelope. So (`envelope/chunks.rs`): `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. + 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 never becomes a piece. These units are packed greedily, in order, From 390a1640c303926049fc8a2ee24598c855e0f1ba Mon Sep 17 00:00:00 2001 From: M3gA-Mind <elvin@tinyhumans.ai> Date: Tue, 6 Oct 2026 21:18:32 +0530 Subject: [PATCH 4/8] Keep blank stretches and starting sections in the splitter - A text that is only whitespace or page breaks is one unit (one piece at a zero target), so split always reassembles exactly. - A trailing page break joins the unit before it instead of becoming a piece of its own. - A text that fits in one piece keeps the section it starts in. - Every whole-document envelope for_item returns goes through encode_checked, as the pieces do. - Tests: blank-only text, whitespace before a heading kept, trailing page break, starting section; a fetch hit must be a piece of the body. --- .../src/cortex/engine/mod_chunk_tests.rs | 6 ++++- .../src/cortex/envelope/chunks.rs | 21 ++++++++++------- .../src/cortex/envelope/chunks_tests.rs | 23 +++++++++++++++++++ .../src/cortex/envelope/mod.rs | 16 ++++++------- 4 files changed, 48 insertions(+), 18 deletions(-) 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 e51ac1df..36feed16 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/mod_chunk_tests.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/mod_chunk_tests.rs @@ -323,6 +323,10 @@ async fn a_long_document_without_pages_or_headings_gets_no_extra_tags() { assert!(!page.hits.is_empty()); for hit in &page.hits { assert!(hit.meta.tags.is_empty(), "{:?}", hit.meta.tags); - assert!(body.contains(&hit.text) || hit.text.len() < body.len()); + assert!( + body.contains(&hit.text), + "a hit is a piece of the body: {:?}", + &hit.text[..hit.text.len().min(80)] + ); } } diff --git a/crates/tinymemory-integrations/src/cortex/envelope/chunks.rs b/crates/tinymemory-integrations/src/cortex/envelope/chunks.rs index 38d135c4..f951e8ec 100644 --- a/crates/tinymemory-integrations/src/cortex/envelope/chunks.rs +++ b/crates/tinymemory-integrations/src/cortex/envelope/chunks.rs @@ -97,7 +97,7 @@ pub(crate) fn split( return Some(vec![Piece { text, pages: paged.then(|| (1, page_count(text))), - section: None, + section: units(text).into_iter().next().and_then(|unit| unit.section), }]); } let mut pieces: Vec<Piece<'_>> = Vec::new(); @@ -153,7 +153,9 @@ impl Open { /// The units of `text`: one per page, cut again before every heading line. /// A stretch holding only whitespace and page breaks is never a unit of its -/// own; it joins the unit that follows. +/// own: it joins the unit that follows, or, at the end of the text, the one +/// before. A text that is nothing but such a stretch is one unit, so every +/// byte of the text is in exactly one unit. fn units(text: &str) -> Vec<Unit> { let mut units = Vec::new(); let mut page = 1; @@ -192,12 +194,15 @@ fn units(text: &str) -> Vec<Unit> { } } if start < text.len() { - units.push(Unit { - start, - end: text.len(), - page: start_page, - section, - }); + match units.last_mut() { + Some(last) if blank(&text[start..]) => last.end = text.len(), + _ => units.push(Unit { + start, + end: text.len(), + page: start_page, + section, + }), + } } units } diff --git a/crates/tinymemory-integrations/src/cortex/envelope/chunks_tests.rs b/crates/tinymemory-integrations/src/cortex/envelope/chunks_tests.rs index ebbf4c77..0f7234b2 100644 --- a/crates/tinymemory-integrations/src/cortex/envelope/chunks_tests.rs +++ b/crates/tinymemory-integrations/src/cortex/envelope/chunks_tests.rs @@ -179,3 +179,26 @@ fn metadata_over_the_target_packs_up_to_the_limit_not_a_few_bytes() { assert!(escaped_len(piece.text) > 1_000, "not a few bytes per piece"); } } + +#[test] +fn blank_stretches_are_kept_and_never_a_piece_of_their_own() { + for text in [" ", "\u{c}", "\n\u{c}\n"] { + let each = pieces(text, 0, 10_000); + assert_eq!(each.len(), 1, "a blank-only text is one piece: {text:?}"); + } + for text in [" # Heading\nbody", "\n# Heading\nbody"] { + let each = pieces(text, 0, 10_000); + assert_eq!(each.len(), 1, "{each:?}"); + assert_eq!(each[0].section.as_deref(), Some("Heading")); + } + let trailing = pieces("a\u{c}", 0, 10_000); + assert_eq!(trailing.len(), 1, "a trailing break joins the unit before"); + assert_eq!(trailing[0].text, "a\u{c}"); +} + +#[test] +fn a_text_that_fits_keeps_the_section_it_starts_in() { + let one = pieces("# Heading\nbody", 1_000, 10_000); + assert_eq!(one.len(), 1); + assert_eq!(one[0].section.as_deref(), Some("Heading")); +} diff --git a/crates/tinymemory-integrations/src/cortex/envelope/mod.rs b/crates/tinymemory-integrations/src/cortex/envelope/mod.rs index 627d1de5..8a2a478c 100644 --- a/crates/tinymemory-integrations/src/cortex/envelope/mod.rs +++ b/crates/tinymemory-integrations/src/cortex/envelope/mod.rs @@ -26,8 +26,8 @@ //! A learning is one event. A conversation is one event per turn, appended //! in order. A document is one event, or, when its envelope would be too big //! for one, one event per piece of its body ([`chunks`]), appended in order; -//! each piece carries its index, and the pages and section it covers, and -//! the pieces concatenate back to the body. +//! each piece carries its index and, when available, the page range and +//! section it covers, and the pieces concatenate back to the body. //! //! No event is sent whose encoded envelope is over //! [`chunks::MAX_EVENT_TEXT_BYTES`] ([`Envelope::encode_checked`]): CortexDB @@ -242,17 +242,15 @@ impl Envelope { // is written whole if it fits, and refused here if not, // never cut into pieces that would each be over the limit. whole.text.clone_from(text); - let size = whole.encode()?.len(); - if size > chunks::MAX_EVENT_TEXT_BYTES { - return Err(Error::InvalidRequest(format!( - "a document event would be {size} bytes and its metadata leaves no \ - room to split it; CortexDB refuses an event over 1 MiB" - ))); - } + whole.encode_checked()?; return Ok(vec![whole]); }; if pieces.len() <= 1 { + // One piece fits under the limit with a chunk field, so + // the whole envelope (which has none) does too; checked + // all the same, as every envelope this returns is. whole.text.clone_from(text); + whole.encode_checked()?; return Ok(vec![whole]); } let count = u32::try_from(pieces.len()).map_err(|_| { From a856f642f69b10a1faf8b8bb03057c464cb2b390 Mon Sep 17 00:00:00 2001 From: M3gA-Mind <elvin@tinyhumans.ai> Date: Tue, 6 Oct 2026 22:15:47 +0530 Subject: [PATCH 5/8] Docs: the piece overhead reserves a page range only for paged documents --- docs/architecture/cortex-wire.md | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/docs/architecture/cortex-wire.md b/docs/architecture/cortex-wire.md index 6fb5bef1..2661e1af 100644 --- a/docs/architecture/cortex-wire.md +++ b/docs/architecture/cortex-wire.md @@ -362,7 +362,8 @@ its whole envelope. So (`envelope/chunks.rs`): whitespace never becomes a piece. 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). + 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 From 303ed92bc50e841845910c7f64cf7b8b8f36bb92 Mon Sep 17 00:00:00 2001 From: M3gA-Mind <elvin@tinyhumans.ai> Date: Tue, 6 Oct 2026 22:53:12 +0530 Subject: [PATCH 6/8] Live-test chunked documents; make the chunking docs precise - live_cortexdb: a ~700 KiB, 24-page, sectioned document is stored through the real wire, comes back whole from list and get, a fetch hit is one piece tagged with its page and section, and forget removes every piece. - Docs: the threshold is measured on the piece envelope (with its chunk field); fetch and recall give one hit per document, its best-ranked piece; pages is an inclusive [first, last] range counted from 1; an all-whitespace text is one piece. --- .../src/cortex/README.md | 11 ++- .../tests/live_cortexdb.rs | 95 +++++++++++++++++++ docs/architecture/cortex-wire.md | 16 +++- docs/specs/memory-v2.md | 2 +- 4 files changed, 113 insertions(+), 11 deletions(-) diff --git a/crates/tinymemory-integrations/src/cortex/README.md b/crates/tinymemory-integrations/src/cortex/README.md index ded16947..26cc2961 100644 --- a/crates/tinymemory-integrations/src/cortex/README.md +++ b/crates/tinymemory-integrations/src/cortex/README.md @@ -111,14 +111,15 @@ operator key is served as `user:local`; a server with no `whoami` route gets no header. The hosted (TinyHumans) wire names the actor itself. **Events.** A learning is one event; a conversation is one event per turn, -appended in order; a document is one event, or, when its envelope would pass -256 KiB (`envelope::chunks::DOCUMENT_CHUNK_TARGET_BYTES`, the one granularity -knob), one event per piece of its body, cut at page breaks and headings and +appended in order; a document is one event, or, when its text with a piece's +envelope (metadata plus the `chunk` field) would pass 256 KiB +(`envelope::chunks::DOCUMENT_CHUNK_TARGET_BYTES`, the one granularity knob), +one event per piece of its body, cut at page breaks and headings and packed up to that size (`envelope/chunks.rs`). No event is sent over 768 KiB of encoded envelope (CortexDB refuses an experience over 1 MiB); an item that cannot fit is refused before anything of its batch is sent. `get` and `list` -reassemble a chunked document; a `fetch` or `recall` hit on a piece carries -that piece, with a `page:<n>` (or `page:<first>-<last>`) tag when the +reassemble a chunked document; `fetch` and `recall` give one hit per document, +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: diff --git a/crates/tinymemory-integrations/tests/live_cortexdb.rs b/crates/tinymemory-integrations/tests/live_cortexdb.rs index 49a0b09d..2a03e027 100644 --- a/crates/tinymemory-integrations/tests/live_cortexdb.rs +++ b/crates/tinymemory-integrations/tests/live_cortexdb.rs @@ -257,3 +257,98 @@ async fn round_trip(engine: &CortexEngine) { assert_eq!(report.forgotten, 3, "all three items are forgotten"); let _ = (conversation, learning); } + +/// A long, paged, sectioned document goes to the real server as several +/// events, each under its 1 MiB limit, and comes back whole: `get` and +/// `list` reassemble it, a `fetch` hit is one piece tagged with its page and +/// section, and `forget` removes every piece. +#[tokio::test] +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<String> = (1..=24) + .map(|page| { + let filler = + format!("Clause {page} covers refunds and delivery terms. ").repeat(600); + format!("# Section {page}\n\n{filler}\n") + }) + .collect(); + let body = pages.join("\u{c}"); + assert!(body.len() > 600 * 1024, "{} bytes", body.len()); + let document = StoreItem::Document { + title: Some("Long contract".into()), + body: tinymemory_api::DocumentBody::Text(body.clone()), + mime: Some("application/pdf".into()), + meta: MemoryMeta { + workspace: Some(workspace.clone()), + file_path: Some("/contracts/long.pdf".into()), + ..MemoryMeta::from_source(SourceKind::File, Some("contracts".into())) + }, + }; + let receipt = engine + .store(document.clone()) + .await + .expect("store in pieces"); + + let filter = MetaFilter { + workspace: Some(workspace.clone()), + ..MetaFilter::default() + }; + let listed = list_until(&engine, &filter, 1).await; + assert_eq!(listed.len(), 1, "one item, not one per piece"); + assert_eq!( + listed[0], + document.render_text(), + "list reassembles the body" + ); + let got = engine + .get(tinymemory_api::GetRequest { + ids: vec![receipt.id.clone()], + reach: None, + }) + .await + .expect("get"); + assert_eq!(got.len(), 1); + assert_eq!( + got[0].text, + document.render_text(), + "get reassembles the body" + ); + + let mut fetch = FetchRequest::new("Clause 7 refunds and delivery", FetchMode::Hybrid, 5); + fetch.filter = filter.clone(); + let page = engine.fetch(fetch).await.expect("fetch"); + let hit = page + .hits + .iter() + .find(|hit| hit.id == receipt.id) + .expect("fetch finds the document"); + assert!(hit.text.len() < body.len(), "a piece, not the whole"); + assert!( + hit.meta.tags.iter().any(|tag| tag.starts_with("page:")) + && hit + .meta + .tags + .iter() + .any(|tag| tag.starts_with("section:Section ")), + "{:?}", + hit.meta.tags + ); + + let report = engine + .forget(ForgetTarget::Ids(vec![receipt.id])) + .await + .expect("forget"); + assert_eq!(report.forgotten, 1); + assert!( + engine + .list(ListRequest::new(filter, 10)) + .await + .expect("list after forget") + .items + .is_empty(), + "no piece is left behind" + ); + } +} diff --git a/docs/architecture/cortex-wire.md b/docs/architecture/cortex-wire.md index 2661e1af..c196bf15 100644 --- a/docs/architecture/cortex-wire.md +++ b/docs/architecture/cortex-wire.md @@ -351,7 +351,8 @@ 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.** A document whose encoded envelope fits in +- **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 @@ -359,7 +360,9 @@ its whole envelope. So (`envelope/chunks.rs`): 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 never becomes a piece. These units are packed greedily, in order, + 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 @@ -372,12 +375,15 @@ its whole envelope. So (`envelope/chunks.rs`): replay detection, `forget` by id or filter, and `get` see all of them; a store that failed part way writes only the missing pieces. - **Reads.** `get` and `list` give the whole document (pieces in index - order). A ranked hit or citation on a piece gives that piece, with the + 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. Readable CortexDB labels for page and section are - not written yet. + 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 diff --git a/docs/specs/memory-v2.md b/docs/specs/memory-v2.md index ceb4878d..7d847046 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` 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); `get`/`list` reassemble a chunked document, and a ranked hit on a piece carries that piece with `page:`/`section:` tags. + - 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); `get`/`list` reassemble a chunked document, 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 51f25ba63abbb700996aa194ad59bd60d9d8f4a9 Mon Sep 17 00:00:00 2001 From: M3gA-Mind <elvin@tinyhumans.ai> Date: Tue, 6 Oct 2026 23:15:56 +0530 Subject: [PATCH 7/8] Size recall packs so every event comes back whole MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit CortexDB's default budgets.max_tokens (4000, about 14 KB) serves a longer event as a budget_excerpt (0.10.4 API §9.5): a slice of the stored JSON envelope that no longer decodes, so the event was never a fetch hit or recall citation. Every pack now sends a budget of a token per byte of the largest event this crate writes, per item asked for, capped at 8 Mi tokens. per_layer_limits still bounds what a pack holds. The live long-document test now polls fetch, as list_until polls the listing. --- .../src/cortex/engine/fetch.rs | 34 ++++++++++++++++--- .../src/cortex/engine/fetch_tests.rs | 15 ++++++-- .../src/cortex/engine/mod_tests.rs | 10 +++++- .../src/cortex/engine/recall.rs | 3 +- .../src/cortex/log/notes.rs | 4 +-- .../tests/live_cortexdb.rs | 18 ++++++---- docs/architecture/cortex-wire.md | 19 ++++++++--- 7 files changed, 83 insertions(+), 20 deletions(-) diff --git a/crates/tinymemory-integrations/src/cortex/engine/fetch.rs b/crates/tinymemory-integrations/src/cortex/engine/fetch.rs index 6016e2aa..00966e3b 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/fetch.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/fetch.rs @@ -34,6 +34,7 @@ use super::CortexEngine; use super::beliefs::{beliefs_in, merge}; use super::cursor::{self, FetchCursor}; use super::items::{event_hit, hit, keeps}; +use crate::cortex::envelope::chunks::MAX_EVENT_TEXT_BYTES; use crate::cortex::envelope::{Envelope, decode_event, labels}; use crate::cortex::error::{Error, Result}; @@ -61,16 +62,19 @@ const EVENTS_PER_HIT: usize = 3; /// - `include` lists `events` first, so the cross-layer token budget funds /// the events this crate reads before any derived layer (by default it /// evicts events first). A caller that reads more layers names them after. -/// - No `budgets.max_tokens`: the server's default only ever evicts, and an -/// evicted event is a hit this crate never sees. Its client-side token -/// budget applies after the read. +/// - `budgets.max_tokens` is [`whole_items_budget`]: room for every event +/// asked for to come back whole. Its client-side token budget applies +/// after the read. pub(super) fn recall_body(scope: &str, query: &str, events: usize, filter: &MetaFilter) -> Value { let mut body = json!({ "scope": scope, "query": query, "view": "granular", "include": ["events"], - "budgets": { "per_layer_limits": { "events": events } }, + "budgets": { + "max_tokens": whole_items_budget(events), + "per_layer_limits": { "events": events }, + }, }); if let Some(labels) = labels::narrowing(filter) { body["filters"] = json!({ "metadata": { "labels": labels } }); @@ -78,6 +82,28 @@ pub(super) fn recall_body(scope: &str, query: &str, events: usize, filter: &Meta body } +/// The most [`whole_items_budget`] asks for: 8 Mi tokens. CortexDB counts +/// about 3.5 bytes a token, so a pack then holds at most about 28 MiB of +/// event text, under the 32 MiB request cap. +pub(super) const MAX_PACK_TOKENS: usize = 8 * 1024 * 1024; + +/// A pack's `budgets.max_tokens` for `items` items: a token per byte of the +/// largest event this crate writes, for each, at most [`MAX_PACK_TOKENS`]. +/// CortexDB's default, 4000 tokens (about 14 KB), cuts a longer event to a +/// `budget_excerpt` (0.10.4 API §9.5): a slice of the stored envelope that +/// no longer decodes, so the hit is lost. +/// +/// The budget only stops cutting; `per_layer_limits` still bounds what a +/// pack holds. A pack of `n` events carries at most `n` × 768 KiB of event +/// text, and past [`MAX_PACK_TOKENS`] (about 37 events at that size) the +/// server excerpts again. +pub(super) fn whole_items_budget(items: usize) -> usize { + items + .max(1) + .saturating_mul(MAX_EVENT_TEXT_BYTES) + .min(MAX_PACK_TOKENS) +} + /// The distinct items of `kind` a pack's events decode to, best rank first, /// keeping only what `filter` matches. pub(super) fn ranked(pack: &Value, kind: Option<ItemKind>, filter: &MetaFilter) -> Vec<Envelope> { diff --git a/crates/tinymemory-integrations/src/cortex/engine/fetch_tests.rs b/crates/tinymemory-integrations/src/cortex/engine/fetch_tests.rs index 0757d23b..14ccdcd1 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/fetch_tests.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/fetch_tests.rs @@ -87,8 +87,19 @@ fn a_recall_body_reads_one_scope_exactly_and_funds_events_first() { "query": "refunds", "view": "granular", "include": ["events"], - "budgets": { "per_layer_limits": { "events": 6 } }, + "budgets": { + "max_tokens": 6 * MAX_EVENT_TEXT_BYTES, + "per_layer_limits": { "events": 6 }, + }, }), - "exact scope, events funded first, no max_tokens and no temporal" + "exact scope, events funded first, room for each event whole, no temporal" ); } + +#[test] +fn a_pack_budget_fits_each_event_whole_up_to_a_ceiling() { + assert_eq!(whole_items_budget(0), MAX_EVENT_TEXT_BYTES, "never zero"); + assert_eq!(whole_items_budget(6), 6 * MAX_EVENT_TEXT_BYTES); + assert_eq!(whole_items_budget(MAX_PACK_EVENTS), MAX_PACK_TOKENS); + assert_eq!(whole_items_budget(usize::MAX), MAX_PACK_TOKENS); +} diff --git a/crates/tinymemory-integrations/src/cortex/engine/mod_tests.rs b/crates/tinymemory-integrations/src/cortex/engine/mod_tests.rs index 619143bd..ca0374d0 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/mod_tests.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/mod_tests.rs @@ -208,7 +208,15 @@ async fn an_unscoped_recall_packs_each_held_scope_exactly_and_answers_once() { body["include"], serde_json::json!(["events", "facts", "beliefs", "episodes", "understanding"]) ); - assert!(body["budgets"].get("max_tokens").is_none(), "{body}"); + let events = body["budgets"]["per_layer_limits"]["events"] + .as_u64() + .unwrap(); + let whole = (events * crate::cortex::envelope::chunks::MAX_EVENT_TEXT_BYTES as u64) + .min(super::fetch::MAX_PACK_TOKENS as u64); + assert!( + body["budgets"]["max_tokens"].as_u64() >= Some(whole), + "room for every event whole: {body}" + ); assert!(body.get("temporal").is_none(), "{body}"); } assert_eq!(seen.answers.len(), 1, "one answer"); diff --git a/crates/tinymemory-integrations/src/cortex/engine/recall.rs b/crates/tinymemory-integrations/src/cortex/engine/recall.rs index e21529b0..3adcee31 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/recall.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/recall.rs @@ -33,7 +33,7 @@ use serde_json::{Value, json}; use tinymemory_api::{Citation, ItemId, Namespace, Reach, RecallAnswer, RecallRequest}; use super::CortexEngine; -use super::fetch::{interleave, ranked, recall_body}; +use super::fetch::{interleave, ranked, recall_body, whole_items_budget}; use super::items::{admitted, located_meta}; use super::scopes::KindScope; use crate::cortex::descriptor::CortexWire; @@ -190,6 +190,7 @@ impl CortexEngine { let mut body = recall_body(scope, &req.question, 0, &req.filter); body["include"] = pack_layers(); body["budgets"]["per_layer_limits"] = pack_budgets(req.limit); + body["budgets"]["max_tokens"] = json!(whole_items_budget(req.limit.saturating_mul(3))); self.log.recall(&body).await } } diff --git a/crates/tinymemory-integrations/src/cortex/log/notes.rs b/crates/tinymemory-integrations/src/cortex/log/notes.rs index 16810706..bcccea08 100644 --- a/crates/tinymemory-integrations/src/cortex/log/notes.rs +++ b/crates/tinymemory-integrations/src/cortex/log/notes.rs @@ -9,8 +9,8 @@ //! storage order, not ranked. This crate reads one exact scope per pack //! (`view: "granular"`), so it must never appear; if it does, a read //! strayed into a parent scope. -//! - **Knapsack evictions.** `budgets.max_tokens` (4000 by default) is a -//! cross-layer budget; items it evicts are counted in +//! - **Knapsack evictions.** `budgets.max_tokens` is a cross-layer budget +//! (this crate sizes it to fit every event whole; see `whole_items_budget`); items it evicts are counted in //! `diagnostics.knapsack_evictions` (when the caller may read //! diagnostics), and a rendered event evicted from `layers.events` is named //! by the `context_contributors` provenance entry with diff --git a/crates/tinymemory-integrations/tests/live_cortexdb.rs b/crates/tinymemory-integrations/tests/live_cortexdb.rs index 2a03e027..5917b48b 100644 --- a/crates/tinymemory-integrations/tests/live_cortexdb.rs +++ b/crates/tinymemory-integrations/tests/live_cortexdb.rs @@ -316,14 +316,20 @@ async fn a_long_document_round_trips_in_pieces() { "get reassembles the body" ); + // Ranking may lag the write; poll as `list_until` polls the listing. let mut fetch = FetchRequest::new("Clause 7 refunds and delivery", FetchMode::Hybrid, 5); fetch.filter = filter.clone(); - let page = engine.fetch(fetch).await.expect("fetch"); - let hit = page - .hits - .iter() - .find(|hit| hit.id == receipt.id) - .expect("fetch finds the document"); + let deadline = Instant::now() + VISIBILITY; + let hit = loop { + let page = engine.fetch(fetch.clone()).await.expect("fetch"); + if let Some(hit) = page.hits.into_iter().find(|hit| hit.id == receipt.id) { + break hit; + } + if Instant::now() >= deadline { + panic!("fetch never found the document"); + } + tokio::time::sleep(Duration::from_millis(500)).await; + }; assert!(hit.text.len() < body.len(), "a piece, not the whole"); assert!( hit.meta.tags.iter().any(|tag| tag.starts_with("page:")) diff --git a/docs/architecture/cortex-wire.md b/docs/architecture/cortex-wire.md index c196bf15..8fe27cdb 100644 --- a/docs/architecture/cortex-wire.md +++ b/docs/architecture/cortex-wire.md @@ -135,7 +135,7 @@ Response: ```json { "scope": "app:tinymemory/app:documents", "query": "...", "view": "granular", "include": ["events"], - "budgets": { "per_layer_limits": { "events": 30 } }, + "budgets": { "max_tokens": 23592960, "per_layer_limits": { "events": 30 } }, "filters": { "metadata": { "labels": ["tm:t:<16 hex>"] } } } ``` @@ -152,8 +152,18 @@ Response: are evicted first. A fetch that wants beliefs sends `["events", "beliefs"]`, an answer pack `["events", "facts", "beliefs", "episodes", "understanding"]`, a beliefs read `["beliefs"]`. -- `max_tokens` is not sent: the default only ever evicts, and an evicted - event is a hit the engine never sees. +- `max_tokens` is sent, sized so every event asked for comes back whole: a + token per byte of the largest event this crate writes (768 KiB), for each + event (and, in an answer pack, each derived item). The default, 4000 + tokens (about 14 KB), cuts a longer event to a `budget_excerpt` (0.10.4 + API §9.5): a slice of the stored envelope that no longer decodes, so a + document piece would never be a hit. It also evicts. + The budget only stops the cutting: `per_layer_limits` still bounds a pack, + so a pack of `n` events carries at most `n` × 768 KiB of event text (a + fetch of 5 asks 18 events: at most 13.5 MiB, typically far less). The + budget is capped at 8 Mi tokens, about 28 MiB at CortexDB's 3.5 bytes a + token; a deeper fetch page past that (about 37 events of the largest size) + gets excerpts again, which do not decode. - `temporal` is not sent. `temporal.reference_date` only anchors `temporal.natural` (a phrase such as "last 30 days", reduced to a capture-time filter) and already defaults to the request time; the field @@ -253,7 +263,8 @@ The beliefs land in a derived layer, read two ways: ```json { "scope": "…", "query": "…", - "budgets": { "per_layer_limits": { "events": 0, "facts": 0, "episodes": 0, + "budgets": { "max_tokens": 786432, + "per_layer_limits": { "events": 0, "facts": 0, "episodes": 0, "understanding": 0, "beliefs": 8 } } } ``` From 9d40d5d74c866ea5abd4acea831b4be0a92986fd Mon Sep 17 00:00:00 2001 From: M3gA-Mind <elvin@tinyhumans.ai> Date: Tue, 6 Oct 2026 23:29:55 +0530 Subject: [PATCH 8/8] Return chunked documents whole or not at all; log partial pack events - get and list return a chunked document only when every piece is present (rebuild_whole), never a truncated body. A store that failed part-way is completed by the next store of the item; fetch and recall still hit the pieces that are there. - A pack event served as a partial view (_partial, e.g. budget_excerpt) does not decode and is dropped; it is now logged at warn with the scope and the reasons. - The budget docs state the bytes-per-token ratio measured on CortexDB 0.10.4 (3 to 3.5 for English, CJK and random text). - The spec states how an item that cannot fit even split is refused. - The live test notes that list waits for every piece, and asserts its document is over twice the chunk target. --- .../src/cortex/README.md | 2 +- .../src/cortex/engine/fetch.rs | 19 ++++++---- .../src/cortex/engine/items.rs | 6 ++-- .../src/cortex/engine/mod_chunk_tests.rs | 16 +++++++++ .../src/cortex/envelope/mod.rs | 2 +- .../src/cortex/envelope/rebuild.rs | 24 +++++++++++++ .../src/cortex/log/notes.rs | 36 +++++++++++++++++-- .../src/cortex/log/notes_tests.rs | 21 +++++++++++ .../tests/live_cortexdb.rs | 10 +++++- docs/architecture/cortex-wire.md | 13 ++++--- docs/specs/memory-v2.md | 2 +- 11 files changed, 131 insertions(+), 20 deletions(-) diff --git a/crates/tinymemory-integrations/src/cortex/README.md b/crates/tinymemory-integrations/src/cortex/README.md index 26cc2961..d1a3e4fa 100644 --- a/crates/tinymemory-integrations/src/cortex/README.md +++ b/crates/tinymemory-integrations/src/cortex/README.md @@ -118,7 +118,7 @@ one event per piece of its body, cut at page breaks and headings and packed up to that size (`envelope/chunks.rs`). No event is sent over 768 KiB of encoded envelope (CortexDB refuses an experience over 1 MiB); an item that cannot fit is refused before anything of its batch is sent. `get` and `list` -reassemble a chunked document; `fetch` and `recall` give one hit per document, +reassemble a chunked document, and return it only when every piece is present; `fetch` and `recall` give one hit per document, 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 diff --git a/crates/tinymemory-integrations/src/cortex/engine/fetch.rs b/crates/tinymemory-integrations/src/cortex/engine/fetch.rs index 00966e3b..07fe27af 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/fetch.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/fetch.rs @@ -82,9 +82,12 @@ pub(super) fn recall_body(scope: &str, query: &str, events: usize, filter: &Meta body } -/// The most [`whole_items_budget`] asks for: 8 Mi tokens. CortexDB counts -/// about 3.5 bytes a token, so a pack then holds at most about 28 MiB of -/// event text, under the 32 MiB request cap. +/// The most [`whole_items_budget`] asks for: 8 Mi tokens. CortexDB 0.10.4 +/// counts 3 to 3.5 bytes a token (measured: 700,000 bytes of English text +/// come back whole at 210,000 tokens and are cut at 200,000; 900,000 bytes +/// of CJK text whole at 300,000; 300,000 random bytes whole at 100,000), so a +/// pack then holds at most about 24 to 28 MiB of event text, under the +/// 32 MiB request cap. pub(super) const MAX_PACK_TOKENS: usize = 8 * 1024 * 1024; /// A pack's `budgets.max_tokens` for `items` items: a token per byte of the @@ -93,10 +96,12 @@ pub(super) const MAX_PACK_TOKENS: usize = 8 * 1024 * 1024; /// `budget_excerpt` (0.10.4 API §9.5): a slice of the stored envelope that /// no longer decodes, so the hit is lost. /// -/// The budget only stops cutting; `per_layer_limits` still bounds what a -/// pack holds. A pack of `n` events carries at most `n` × 768 KiB of event -/// text, and past [`MAX_PACK_TOKENS`] (about 37 events at that size) the -/// server excerpts again. +/// A token per byte is at least three times the room an event needs. The +/// budget only stops cutting; `per_layer_limits` still bounds what a pack +/// holds. A pack of `n` events carries at most `n` × 768 KiB of event text, +/// and only past [`MAX_PACK_TOKENS`] (at least 32 events of that size in one +/// pack, or about 100 at the 256 KiB chunk target) does the server excerpt +/// again; the pack notes log any excerpt at warn. pub(super) fn whole_items_budget(items: usize) -> usize { items .max(1) diff --git a/crates/tinymemory-integrations/src/cortex/engine/items.rs b/crates/tinymemory-integrations/src/cortex/engine/items.rs index 5f94e25d..ab9fdc1b 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/items.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/items.rs @@ -10,7 +10,7 @@ use tinymemory_api::{ use super::CortexEngine; use super::scopes::KindScope; -use crate::cortex::envelope::{Decoded, Envelope, decode_event, labels, rebuild}; +use crate::cortex::envelope::{Decoded, Envelope, decode_event, labels, rebuild, rebuild_whole}; use crate::cortex::error::Result; /// The kinds `filter` admits, in the fixed order @@ -113,7 +113,7 @@ impl CortexEngine { } for (id, events) in self.item_events(&scope, &ids).await? { let envelopes: Vec<Envelope> = events.into_iter().map(|d| d.envelope).collect(); - if let Some(item) = rebuild(&envelopes) { + if let Some(item) = rebuild_whole(&envelopes) { found.insert(ItemId::new(id.clone()), hit(&id, &item, 0.0)); } } @@ -147,7 +147,7 @@ impl CortexEngine { let scope = KindScope::new(namespace.clone(), kind); for (id, events) in self.item_events(&scope, &ids).await? { let envelopes: Vec<Envelope> = events.into_iter().map(|d| d.envelope).collect(); - if let Some(item) = rebuild(&envelopes) { + if let Some(item) = rebuild_whole(&envelopes) { out.insert(id, item); } } 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 36feed16..24ef5ae2 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/mod_chunk_tests.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/mod_chunk_tests.rs @@ -165,6 +165,22 @@ async fn a_store_that_lost_a_piece_writes_only_that_piece_again() { .unwrap(); log.events.remove(last); } + let got = engine + .get(GetRequest { + ids: vec![ItemId::new(item.fingerprint())], + reach: None, + }) + .await + .unwrap(); + assert!(got.is_empty(), "never a truncated body: {got:?}"); + let listed = engine + .list(ListRequest::new( + MetaFilter::kinds([ItemKind::Document]), + 10, + )) + .await + .unwrap(); + assert!(listed.items.is_empty(), "never a truncated body"); let again = engine.store(item.clone()).await.unwrap(); assert!(!again.replayed, "a piece was missing"); assert_eq!(events(&state, SCOPE).len(), before, "exactly that piece"); diff --git a/crates/tinymemory-integrations/src/cortex/envelope/mod.rs b/crates/tinymemory-integrations/src/cortex/envelope/mod.rs index 8a2a478c..1485a973 100644 --- a/crates/tinymemory-integrations/src/cortex/envelope/mod.rs +++ b/crates/tinymemory-integrations/src/cortex/envelope/mod.rs @@ -62,7 +62,7 @@ use tinymemory_api::{ use crate::cortex::error::{Error, Result}; -pub(crate) use rebuild::{Decoded, decode_event, rebuild}; +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"; diff --git a/crates/tinymemory-integrations/src/cortex/envelope/rebuild.rs b/crates/tinymemory-integrations/src/cortex/envelope/rebuild.rs index b9ffbcbd..90a2d7f0 100644 --- a/crates/tinymemory-integrations/src/cortex/envelope/rebuild.rs +++ b/crates/tinymemory-integrations/src/cortex/envelope/rebuild.rs @@ -1,5 +1,7 @@ //! Reading events back into envelopes and envelopes back into items. +use std::collections::HashSet; + use serde_json::Value; use tinymemory_api::{DocumentBody, ItemKind, LearningKind, StoreItem, Turn}; @@ -75,6 +77,28 @@ pub(crate) fn rebuild(envelopes: &[Envelope]) -> Option<StoreItem> { }) } +/// [`rebuild`] for a read that returns whole items (`get`, `list`): `None` +/// for a chunked document missing a piece, from a store that failed part-way +/// (the next store of the item writes the missing ones) or pieces the engine +/// 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<StoreItem> { + if let Some(count) = envelopes.iter().find_map(|e| Some(e.chunk.as_ref()?.count)) { + let held: HashSet<u32> = 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; + } + } + 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. fn document_text(envelopes: &[Envelope]) -> String { diff --git a/crates/tinymemory-integrations/src/cortex/log/notes.rs b/crates/tinymemory-integrations/src/cortex/log/notes.rs index bcccea08..a232454d 100644 --- a/crates/tinymemory-integrations/src/cortex/log/notes.rs +++ b/crates/tinymemory-integrations/src/cortex/log/notes.rs @@ -19,6 +19,11 @@ //! renders no `context_block`, so it carries no contributor rows; there the //! diagnostics count is the only signal. Every pack lists `events` first in //! `include` so the budget funds events before derived layers. +//! - **Partial events.** An event whose `content._partial` is `true` (a +//! `budget_excerpt`, §9.5, or another recall-time view) carries a slice of +//! its stored text, which does not decode, so it is a hit the engine +//! drops. Logged at warn with the scope, the count and the reasons. The +//! budget makes this rare (see `whole_items_budget`); this keeps it loud. use serde_json::Value; @@ -40,6 +45,9 @@ pub(crate) enum Note<'a> { /// Event ids marked `evicted_from_layers: true`. events: Vec<&'a str>, }, + /// Events served as partial views: each one's `_partial_reason`, or + /// `"unknown"`. + Partial(Vec<&'a str>), } impl Note<'_> { @@ -47,12 +55,13 @@ impl Note<'_> { pub(crate) fn level(&self) -> log::Level { match self { Self::Warning(_) => log::Level::Debug, - Self::ParentSample(_) | Self::Evicted { .. } => log::Level::Warn, + Self::ParentSample(_) | Self::Evicted { .. } | Self::Partial(_) => log::Level::Warn, } } } -/// Every note `pack` carries, in a fixed order: warnings, then evictions. +/// Every note `pack` carries, in a fixed order: warnings, evictions, then +/// partial events. pub(crate) fn notes(pack: &Value) -> Vec<Note<'_>> { let mut notes: Vec<Note<'_>> = pack .get("warnings") @@ -85,6 +94,23 @@ pub(crate) fn notes(pack: &Value) -> Vec<Note<'_>> { if count.is_some() || !events.is_empty() { notes.push(Note::Evicted { count, events }); } + let partial: Vec<&str> = pack + .pointer("/layers/events") + .and_then(Value::as_array) + .into_iter() + .flatten() + .filter_map(|event| event.get("content")) + .filter(|content| content.get("_partial").and_then(Value::as_bool) == Some(true)) + .map(|content| { + content + .get("_partial_reason") + .and_then(Value::as_str) + .unwrap_or("unknown") + }) + .collect(); + if !partial.is_empty() { + notes.push(Note::Partial(partial)); + } notes } @@ -123,6 +149,12 @@ pub(crate) fn report(scope: &str, pack: &Value) { count.map_or_else(|| "unreported".to_string(), |count| count.to_string()), events.len() ), + Note::Partial(reasons) => log::log!( + level, + "[cortex] recall pack served events as partial views, which do not decode \ + scope={scope:?} partial_events={} reasons={reasons:?}", + reasons.len() + ), } } } diff --git a/crates/tinymemory-integrations/src/cortex/log/notes_tests.rs b/crates/tinymemory-integrations/src/cortex/log/notes_tests.rs index 176d27f7..0c90d9a0 100644 --- a/crates/tinymemory-integrations/src/cortex/log/notes_tests.rs +++ b/crates/tinymemory-integrations/src/cortex/log/notes_tests.rs @@ -202,3 +202,24 @@ fn a_newline_in_a_scope_or_a_warning_cannot_forge_a_log_line() { assert!(message.contains("\\n[cortex] forged"), "{message}"); assert!(message.contains("\\nERROR forged line"), "{message}"); } + +#[test] +fn a_partial_event_is_logged_at_warn_with_its_reason() { + capture(); + let scope = "app:tinymemory/agent:notes-partial/app:documents"; + let pack = json!({ "layers": { "events": [ + { "id": "a", "content": { "text": "whole" } }, + { "id": "b", "content": { "text": "[…]\nslice", "_partial": true, + "_partial_reason": "budget_excerpt" } }, + { "id": "c", "content": { "text": "view", "_partial": true } }, + ] } }); + assert_eq!( + notes(&pack), + [Note::Partial(vec!["budget_excerpt", "unknown"])] + ); + report(scope, &pack); + let logged = logged_for(scope); + assert_eq!(logged.len(), 1, "{logged:?}"); + assert_eq!(logged[0].0, log::Level::Warn); + assert!(logged[0].1.contains("partial_events=2"), "{}", logged[0].1); +} diff --git a/crates/tinymemory-integrations/tests/live_cortexdb.rs b/crates/tinymemory-integrations/tests/live_cortexdb.rs index 5917b48b..dfd6990e 100644 --- a/crates/tinymemory-integrations/tests/live_cortexdb.rs +++ b/crates/tinymemory-integrations/tests/live_cortexdb.rs @@ -275,7 +275,13 @@ async fn a_long_document_round_trips_in_pieces() { }) .collect(); let body = pages.join("\u{c}"); - assert!(body.len() > 600 * 1024, "{} bytes", body.len()); + // A document splits once its envelope passes the 256 KiB chunk + // target (not the 1 MiB event limit): this one is several pieces. + assert!( + body.len() > 2 * 256 * 1024, + "over twice the chunk target: {} bytes", + body.len() + ); let document = StoreItem::Document { title: Some("Long contract".into()), body: tinymemory_api::DocumentBody::Text(body.clone()), @@ -295,6 +301,8 @@ async fn a_long_document_round_trips_in_pieces() { workspace: Some(workspace.clone()), ..MetaFilter::default() }; + // `list` returns a chunked document only once every piece is + // listed, so this waits for all of them, not just the first. let listed = list_until(&engine, &filter, 1).await; assert_eq!(listed.len(), 1, "one item, not one per piece"); assert_eq!( diff --git a/docs/architecture/cortex-wire.md b/docs/architecture/cortex-wire.md index 8fe27cdb..cc3f216c 100644 --- a/docs/architecture/cortex-wire.md +++ b/docs/architecture/cortex-wire.md @@ -161,9 +161,12 @@ Response: The budget only stops the cutting: `per_layer_limits` still bounds a pack, so a pack of `n` events carries at most `n` × 768 KiB of event text (a fetch of 5 asks 18 events: at most 13.5 MiB, typically far less). The - budget is capped at 8 Mi tokens, about 28 MiB at CortexDB's 3.5 bytes a - token; a deeper fetch page past that (about 37 events of the largest size) - gets excerpts again, which do not decode. + budget is capped at 8 Mi tokens, about 24 to 28 MiB at the 3 to 3.5 bytes + a token CortexDB 0.10.4 counts (measured on English, CJK and random text), + so a token per byte is at least three times the room an event needs. Only + a pack holding more than that (at least 32 events of the largest size, or + about 100 at the chunk target) gets excerpts again, which do not decode + and are logged at warn (`log/notes.rs`). - `temporal` is not sent. `temporal.reference_date` only anchors `temporal.natural` (a phrase such as "last 30 days", reduced to a capture-time filter) and already defaults to the request time; the field @@ -384,7 +387,9 @@ its whole envelope. So (`envelope/chunks.rs`): (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. + 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 diff --git a/docs/specs/memory-v2.md b/docs/specs/memory-v2.md index 7d847046..87ead2e7 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); `get`/`list` reassemble a chunked document, 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 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. - 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.