Skip to content
15 changes: 13 additions & 2 deletions crates/tinymemory-integrations/src/cortex/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -110,8 +110,19 @@ 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 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`
Comment thread
M3gA-Mind marked this conversation as resolved.
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
`content.text` is a JSON envelope:

```json
{ "v": 2, "id": "<40-hex fingerprint>", "kind": "conversation",
Expand Down
54 changes: 43 additions & 11 deletions crates/tinymemory-integrations/src/cortex/engine/fetch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,8 +33,9 @@ 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::chunks::MAX_EVENT_TEXT_BYTES;
use crate::cortex::envelope::{Envelope, decode_event, labels};
use crate::cortex::error::{Error, Result};

/// The cursor tag of a fetch.
Expand All @@ -61,23 +62,53 @@ 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 } });
}
body
}

/// 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
/// 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.
///
/// 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 {
Comment thread
M3gA-Mind marked this conversation as resolved.
Comment thread
M3gA-Mind marked this conversation as resolved.
items
.max(1)
.saturating_mul(MAX_EVENT_TEXT_BYTES)
.min(MAX_PACK_TOKENS)
Comment thread
M3gA-Mind marked this conversation as resolved.
}

/// 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> {
Expand Down Expand Up @@ -159,11 +190,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 {
Expand Down
15 changes: 13 additions & 2 deletions crates/tinymemory-integrations/src/cortex/engine/fetch_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
52 changes: 47 additions & 5 deletions crates/tinymemory-integrations/src/cortex/engine/items.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,11 +4,13 @@
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;
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
Expand All @@ -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 {
Comment thread
M3gA-Mind marked this conversation as resolved.
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 {
Expand Down Expand Up @@ -82,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));
}
}
Expand All @@ -95,17 +126,28 @@ 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 {
by_node.entry(namespace).or_default().push(id.clone());
}
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) {
if let Some(item) = rebuild_whole(&envelopes) {
out.insert(id, item);
}
}
Expand Down
56 changes: 33 additions & 23 deletions crates/tinymemory-integrations/src/cortex/engine/list.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
//!
Expand All @@ -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};
Expand All @@ -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 {
Expand Down Expand Up @@ -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)),
Comment thread
M3gA-Mind marked this conversation as resolved.
})
.collect())
}
Expand Down
4 changes: 4 additions & 0 deletions crates/tinymemory-integrations/src/cortex/engine/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -263,6 +263,10 @@ mod tests;
#[path = "mod_list_tests.rs"]
mod list_tests;

#[cfg(test)]
Comment thread
M3gA-Mind marked this conversation as resolved.
#[path = "mod_chunk_tests.rs"]
mod chunk_tests;
Comment thread
M3gA-Mind marked this conversation as resolved.
Comment thread
M3gA-Mind marked this conversation as resolved.

#[cfg(test)]
#[path = "mod_direct_tests.rs"]
mod direct_tests;
Expand Down
Loading
Loading