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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 13 additions & 0 deletions crates/tinymemory-api/src/item/mod_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -173,3 +173,16 @@ fn a_turn_renders_its_tool_calls() {
"assistant: done [tools: shell (call-1), grep]"
);
}

#[test]
fn an_unset_derive_flag_leaves_the_fingerprint_as_it_was() {
let item = StoreItem::document("Refunds take five days.", MemoryMeta::default());
let json = serde_json::to_string(item.meta()).unwrap();
assert!(!json.contains("derive"), "{json}");
let mut opted_out = item.clone();
opted_out.meta_mut().derive = Some(false);
assert_ne!(opted_out.fingerprint(), item.fingerprint());
let back: MemoryMeta =
serde_json::from_str(&serde_json::to_string(opted_out.meta()).unwrap()).unwrap();
assert_eq!(back.derive, Some(false));
}
1 change: 1 addition & 0 deletions crates/tinymemory-api/src/meta/filter_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ fn meta() -> MemoryMeta {
},
tags: vec!["x".into(), "y".into()],
observed_at: Some(Utc.with_ymd_and_hms(2026, 1, 2, 3, 4, 5).unwrap()),
derive: None,
}
}

Expand Down
14 changes: 14 additions & 0 deletions crates/tinymemory-api/src/meta/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,11 @@ use serde::{Deserialize, Serialize};
use crate::namespace::Namespace;

/// Where an item came from and what it is about.
///
/// Fields are added as the contract grows, so build one with
/// `..MemoryMeta::default()` for the fields you do not set. A literal that
/// names every field stops compiling when a field is added, which is a
/// major release.
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
#[serde(default)]
pub struct MemoryMeta {
Expand Down Expand Up @@ -68,6 +73,15 @@ pub struct MemoryMeta {
/// When the underlying fact was observed, as opposed to when it was stored.
#[serde(skip_serializing_if = "Option::is_none")]
pub observed_at: Option<DateTime<Utc>>,
/// Whether an engine that derives layers from what it stores (facts,
/// beliefs, concepts) may derive them from this item. `Some(false)`
/// stores and indexes the item, so it stays searchable, but derives
/// nothing from it: runtime state, tool output, machine-written
/// summaries. `None`, the default, leaves it to the engine; engines that
/// derive nothing ignore it. Unset, it is not serialized, so an item's
/// fingerprint is the same as before the field existed.
#[serde(skip_serializing_if = "Option::is_none")]
pub derive: Option<bool>,
}

impl MemoryMeta {
Expand Down
32 changes: 23 additions & 9 deletions crates/tinymemory-integrations/src/cortex/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@ This README is the short in-tree summary. The full reference is under
- [`cortex.md`](../../../../docs/architecture/cortex.md): surface, credentials,
transport, failure mapping, endpoint security, the registry and `MemoryConfig`;
- [`cortex-wire.md`](../../../../docs/architecture/cortex-wire.md): every endpoint
and its shapes, scope layout, the v2 envelope, lookup labels;
and its shapes, scope layout, the envelope (v3 and v2), lookup labels;
- [`cortex-flows.md`](../../../../docs/architecture/cortex-flows.md): step-by-step
store, list, fetch, recall, forget, get, discovery;
- [`testing.md`](../../../../docs/architecture/testing.md): the doubles, the
Expand Down Expand Up @@ -122,18 +122,28 @@ reassemble a chunked document, and return it only when every piece is present; `
its best-ranked piece, with a `page:<n>` (or `page:<first>-<last>`) tag when the
document marks its pages and a `section:<title>` tag when the piece starts
under a heading; a piece with neither carries no extra tag. Each event's
`content.text` is a JSON envelope:
`content.text` is the item's own text (the body or piece, the turn's text, or
the statement), and the rest of its envelope rides in `context.labels` as
`tm:e:<NN>:` parts of at most 240 bytes of JSON (v3). The parts carry the
envelope with `text` left empty, because the text is the event's
`content.text`; joined, they read:

```json
{ "v": 2, "id": "<40-hex fingerprint>", "kind": "conversation",
"text": "<body | turn text | statement>", "meta": { ... MemoryMeta ... },
{ "v": 3, "id": "<40-hex fingerprint>", "kind": "conversation", "text": "",
"meta": { ... MemoryMeta ... },
"title": "...", "mime": "...", "learning_kind": "...", "confidence": 0.8,
"evidence": "...",
"turn": { "index": 0, "count": 3, "role": "user", "at": "...", "tool_calls": [] } }
```

Kind-specific fields appear only when set. Text that is not a v2 envelope is
someone else's event and is ignored. `context.observed_at` carries the turn's
Kind-specific fields appear only when set. Readable labels (`kind:`, `file:`,
`page:`, `section:`) sit beside the parts. An event with empty text, or whose
labels would pass 64, is written as v2 (the whole envelope as JSON text, as
every event was before), and both layouts read. An event that is neither is
someone else's and is ignored. An item that opts out of derivation
(`MemoryMeta::derive == Some(false)`), and every tool turn, is sent with
`directives.extract: []`: indexed and searchable, but no facts, beliefs or
concepts are derived from it. `context.observed_at` carries the turn's
`at` or the item's `meta.observed_at`.

**Labels.** Each event carries up to eight `context.labels`, each a 16-hex
Expand Down Expand Up @@ -178,7 +188,9 @@ as prefixes, so they cannot be labelled and are filtered only client-side.
scopes, or, for an unscoped read, every kind scope the engine holds. With
nothing to read the answer is empty and nothing is sent. The answer comes
from the pack holding the most admitted events. The answer route is
called **once** with `use_pack_id`. Hosted omits a null
called **once** with `use_pack_id` (again, after every pack is built
anew, when the packs were dropped in between, up to three rounds:
CortexDB drops every pack on any forget). Hosted omits a null
`answer_instructions`, because its schema is strict; Direct sends `null`.
Citations come from the packs' decoded events, filtered (reach included),
merged rank by rank (the most specific node's first), one per item, capped
Expand All @@ -202,8 +214,10 @@ as prefixes, so they cannot be labelled and are filtered only client-side.
These were measured against a live CortexDB by the v1 adapter. The doubles in
`testing/` reproduce all of them.

- **Append-only.** There is no update route. Forget removes events but not
their idempotency records.
- **Append-only.** There is no update route. A body `idempotency_key` replays
the same body for 24 hours and refuses another body (409); forget by
`memory_ids` releases it. Keys are derived from the body, so a retry is a
replay.
- **Accepted is not readable.** A write first waits until the label-narrowed
listing carries its event (fatal after 30s). It then waits until ranked
recall returns it (best-effort, 10s); a recall that is down or slow does not
Expand Down
3 changes: 2 additions & 1 deletion crates/tinymemory-integrations/src/cortex/engine/beliefs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@ use tinymemory_api::{
};

use super::CortexEngine;
use super::fetch::recall_body;
use super::fetch::{recall_body, whole_items_budget};
use super::items::hit;
use crate::cortex::descriptor::{CortexWire, Route};
use crate::cortex::envelope::parse_scope;
Expand Down Expand Up @@ -174,6 +174,7 @@ impl CortexEngine {
body["budgets"]["per_layer_limits"] = json!({
"events": 0, "facts": 0, "episodes": 0, "understanding": 0, "beliefs": limit,
});
body["budgets"]["max_tokens"] = json!(whole_items_budget(limit));
let pack = self.log.recall(&body).await?;
Ok(beliefs_in(&pack, "/layers/beliefs"))
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -98,6 +98,14 @@ async fn a_built_scope_s_beliefs_are_read_with_and_without_a_query() {
let belief_pack = state.seen.lock().unwrap().recalls.last().cloned().unwrap();
assert_eq!(belief_pack["include"], json!(["beliefs"]));
assert_eq!(belief_pack["view"], "granular");
let asked = belief_pack["budgets"]["per_layer_limits"]["beliefs"]
.as_u64()
.unwrap();
assert_eq!(
belief_pack["budgets"]["max_tokens"].as_u64(),
Some(asked * crate::cortex::envelope::chunks::MAX_EVENT_TEXT_BYTES as u64),
"room for each belief asked for"
);
assert_eq!(ranked[0].meta.namespace, node);
let listed = engine
.beliefs(BeliefsRequest::new(Reach::exact(node), 5))
Expand Down
24 changes: 11 additions & 13 deletions crates/tinymemory-integrations/src/cortex/engine/mod_chunk_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -156,15 +156,7 @@ async fn a_store_that_lost_a_piece_writes_only_that_piece_again() {
let item = handbook(30);
engine.store(item.clone()).await.unwrap();
let before = events(&state, SCOPE).len();
{
let mut log = state.log.lock().unwrap();
let last = log
.events
.iter()
.rposition(|event| event["scope"] == SCOPE)
.unwrap();
log.events.remove(last);
}
state.log.lock().unwrap().lose_last(SCOPE);
let got = engine
.get(GetRequest {
ids: vec![ItemId::new(item.fingerprint())],
Expand All @@ -189,7 +181,7 @@ async fn a_store_that_lost_a_piece_writes_only_that_piece_again() {
}

#[tokio::test]
async fn a_short_document_is_one_event_exactly_as_before_and_reads_the_same() {
async fn a_short_document_is_one_event_and_reads_the_same() {
for (engine, state) in both().await {
let item = StoreItem::Document {
title: Some("Note".into()),
Expand All @@ -200,9 +192,15 @@ async fn a_short_document_is_one_event_exactly_as_before_and_reads_the_same() {
engine.store(item.clone()).await.unwrap();
let written = events(&state, SCOPE);
assert_eq!(written.len(), 1);
let text = written[0]["content"]["text"].as_str().unwrap();
let envelope: Value = serde_json::from_str(text).unwrap();
assert!(envelope.get("chunk").is_none(), "no piece info: {envelope}");
assert_eq!(
written[0]["content"]["text"],
format!("Page one.{PAGE_BREAK}Page two."),
"the body itself"
);
let envelope = crate::cortex::envelope::decode_event(&written[0])
.unwrap()
.envelope;
assert!(envelope.chunk.is_none(), "no piece info: {envelope:?}");
let listed = engine
.list(ListRequest::new(MetaFilter::default(), 10))
.await
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -395,3 +395,30 @@ async fn the_answer_body_holds_only_keys_the_strict_schema_allows() {
with.instructions = Some("be brief".into());
engine.recall(with).await.unwrap();
}

#[tokio::test]
async fn recovery_does_not_take_another_turn_with_the_same_words_for_the_lost_one() {
use tinymemory_api::{MemoryMeta, Role, Turn};
let (endpoint, state) = hosted_double().await;
let engine = hosted_engine(&endpoint).with_test_timing(std::time::Duration::from_millis(200));
let item = StoreItem::Conversation {
turns: vec![
Turn::new(Role::User, "ok"),
Turn::new(Role::Assistant, "ok"),
],
meta: MemoryMeta::default(),
};
engine.store(item.clone()).await.unwrap();
// The second turn's event is lost.
state
.log
.lock()
.unwrap()
.lose_last("app:tinymemory/app:conversations");
// Its re-write is claimed but never applied, so recovery must look for
// it, and must not take the first turn, which says the same words.
state.claim_then_fail.store(1, Ordering::SeqCst);
let error = engine.store(item).await.unwrap_err();
assert!(matches!(error, Error::Unavailable(_)), "{error:?}");
assert_eq!(state.event_count(), 1, "the lost turn is still missing");
}
43 changes: 43 additions & 0 deletions crates/tinymemory-integrations/src/cortex/engine/mod_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -320,3 +320,46 @@ async fn a_store_succeeds_when_ranked_recall_is_down() {
assert_eq!(listed.items.len(), 1, "the settle probe is best-effort");
}
}

#[tokio::test]
async fn an_answer_whose_pack_expired_recalls_that_scope_again() {
for (engine, state) in both().await {
for item in items() {
engine.store(item).await.unwrap();
}
let mut req = RecallRequest::new("which editor helix", 2);
req.filter = MetaFilter::kinds([ItemKind::Learning]);
state.seen.lock().unwrap().recalls.clear();
state.expire_packs.store(1, Ordering::SeqCst);
let answer = engine.recall(req.clone()).await.unwrap();
assert_eq!(answer.answer, "grounded answer for which editor helix");
{
let seen = state.seen.lock().unwrap();
assert_eq!(seen.recalls.len(), 2, "the one scope, packed again");
assert_eq!(seen.recalls[0], seen.recalls[1], "the same pack request");
assert_eq!(seen.answers.len(), 2, "answered from the fresh pack");
}

// With several scopes every pack is read again, so the citations
// come from packs read after whatever dropped them, not before.
let mut wide = RecallRequest::new("which editor helix", 2);
wide.filter = MetaFilter::kinds([ItemKind::Learning, ItemKind::Conversation]);
state.seen.lock().unwrap().recalls.clear();
state.expire_packs.store(1, Ordering::SeqCst);
engine.recall(wide).await.unwrap();
assert_eq!(
state.seen.lock().unwrap().recalls.len(),
4,
"two scopes, each packed in both rounds"
);

state.seen.lock().unwrap().answers.clear();
state.expire_packs.store(3, Ordering::SeqCst);
assert!(
matches!(engine.recall(req).await, Err(Error::NotFound(_))),
"a bounded retry, not forever"
);
assert_eq!(state.seen.lock().unwrap().answers.len(), 3);
state.expire_packs.store(0, Ordering::SeqCst);
}
}
Loading
Loading