From a7e58a8a71f274622c35b047d17c6bac16d8c0d1 Mon Sep 17 00:00:00 2001 From: M3gA-Mind Date: Tue, 6 Oct 2026 17:08:30 +0530 Subject: [PATCH 1/4] Skip post-turn and ingest belief builds on self-building engines CortexDB's managed API extracts facts and rebuilds a scope's beliefs on its own shortly after each write, so the build jobs the lifecycle queued after every N turns and every brain ingest only re-ran a full scope rebuild. - tinymemory-api: add `Consolidation::Automatic`. An explicit `consolidate` still builds and answers as `OnDemand` does; the conformance suite holds it to that, and `ReferenceEngine::with_consolidation` lets a host test it. - tinymemory-integrations: a direct engine declares `Automatic` when its endpoint is the managed API and `OnDemand` anywhere else (self-hosted CortexDB builds in the background only when its operator turns the layer scheduler on). `CortexEngine::with_consolidation` and `EngineSettings::consolidation` override it; the TinyHumans wire stays `Scheduled`. The choice is made in one place, `direct_consolidation`, so a server-reported signal can replace the endpoint check later. - tinymemory-tools: `post_turn` and `Brain::ingest`/`ingest_many` hand back no build on an `Automatic` engine. `Brain::build` (new) and `AgentMemory::history_build` still ask for one at any time. Breaking: `Ingested::job` is now `Option`, and `Consolidation` gains a variant. --- .../src/conformance/reference/mod.rs | 10 +++ .../src/conformance/suite/lifecycle.rs | 13 ++- .../src/conformance/suite/lifecycle_tests.rs | 30 +++++++ crates/tinymemory-api/src/consolidate/mod.rs | 8 ++ .../src/consolidate/mod_tests.rs | 8 ++ crates/tinymemory-api/src/engine/mod.rs | 4 +- .../tests/conformance_reference.rs | 9 ++ .../examples/cortex_agent.rs | 2 +- .../examples/memory_eval/main.rs | 2 +- .../tinymemory-integrations/src/config/mod.rs | 10 ++- .../src/config/mod_tests.rs | 19 +++++ .../src/cortex/README.md | 8 +- .../src/cortex/descriptor/mod.rs | 47 ++++++++-- .../src/cortex/descriptor/mod_tests.rs | 29 +++++++ .../src/cortex/engine/consolidate_tests.rs | 85 ++++++++++++++++++- .../src/cortex/engine/mod.rs | 58 +++++++++++-- .../src/cortex/lifecycle_tests.rs | 2 +- .../src/registry/mod.rs | 10 ++- .../src/registry/mod_tests.rs | 34 ++++++++ .../tests/live_cortex_lifecycle.rs | 14 +-- crates/tinymemory-tools/examples/brain.rs | 4 +- crates/tinymemory-tools/src/background/mod.rs | 9 ++ crates/tinymemory-tools/src/brain/mod.rs | 40 ++++++--- .../tinymemory-tools/src/brain/mod_tests.rs | 44 +++++++++- crates/tinymemory-tools/src/brain/types.rs | 8 +- crates/tinymemory-tools/src/lifecycle/mod.rs | 17 ++-- .../src/lifecycle/mod_tests.rs | 32 ++++++- .../tinymemory-tools/src/lifecycle/types.rs | 4 +- docs/architecture/lifecycle.md | 6 +- docs/integration.md | 7 +- docs/specs/agent-memory.md | 9 +- 31 files changed, 511 insertions(+), 71 deletions(-) create mode 100644 crates/tinymemory-api/src/conformance/suite/lifecycle_tests.rs diff --git a/crates/tinymemory-api/src/conformance/reference/mod.rs b/crates/tinymemory-api/src/conformance/reference/mod.rs index d7700e81..a1b72d87 100644 --- a/crates/tinymemory-api/src/conformance/reference/mod.rs +++ b/crates/tinymemory-api/src/conformance/reference/mod.rs @@ -59,6 +59,16 @@ impl ReferenceEngine { } } + /// The same engine, declaring `consolidation`. Its builds run either + /// way, so it serves a host test of how a lifecycle treats an engine + /// that declares [`Consolidation::Automatic`]; the conformance suite + /// holds it to whatever it declares. + #[must_use] + pub fn with_consolidation(mut self, consolidation: Consolidation) -> Self { + self.descriptor.consolidation = consolidation; + self + } + /// How many items the engine holds. #[must_use] pub fn len(&self) -> usize { diff --git a/crates/tinymemory-api/src/conformance/suite/lifecycle.rs b/crates/tinymemory-api/src/conformance/suite/lifecycle.rs index 94ece066..08165436 100644 --- a/crates/tinymemory-api/src/conformance/suite/lifecycle.rs +++ b/crates/tinymemory-api/src/conformance/suite/lifecycle.rs @@ -5,8 +5,9 @@ //! id; and a repeat of a visible item is a replay. //! - `consolidate` — an invalid request is refused, and a valid one answers //! as the descriptor promises: [`Consolidation::None`] refuses with -//! `Unsupported`, [`Consolidation::OnDemand`] starts or completes a build, -//! and [`Consolidation::Scheduled`] acknowledges without one. Then +//! `Unsupported`, [`Consolidation::OnDemand`] and +//! [`Consolidation::Automatic`] start or complete a build, and +//! [`Consolidation::Scheduled`] acknowledges without one. Then //! `beliefs` refuses a zero limit, and every belief it returns is a //! learning tagged [`BELIEF_TAG`] within the reach asked for. @@ -119,13 +120,13 @@ fn answers_as_promised( check: CHECK, source, }), - (Consolidation::OnDemand, Ok(receipt)) => ensure( + (Consolidation::OnDemand | Consolidation::Automatic, Ok(receipt)) => ensure( CHECK, matches!( receipt.status, ConsolidateStatus::Started | ConsolidateStatus::Completed ), - || format!("an on-demand engine answered {:?}", receipt.status), + || format!("a {promised:?} engine answered {:?}", receipt.status), ), (Consolidation::Scheduled, Ok(receipt)) => ensure( CHECK, @@ -174,3 +175,7 @@ async fn beliefs(ctx: &Ctx<'_>, node: Namespace) -> Result<()> { } Ok(()) } + +#[cfg(test)] +#[path = "lifecycle_tests.rs"] +mod tests; diff --git a/crates/tinymemory-api/src/conformance/suite/lifecycle_tests.rs b/crates/tinymemory-api/src/conformance/suite/lifecycle_tests.rs new file mode 100644 index 00000000..ccf19904 --- /dev/null +++ b/crates/tinymemory-api/src/conformance/suite/lifecycle_tests.rs @@ -0,0 +1,30 @@ +//! How a consolidation's answer is held to the descriptor's promise. + +use super::*; +use crate::ConsolidateReceipt; + +fn receipt(status: ConsolidateStatus) -> ConsolidateReceipt { + ConsolidateReceipt { + status, + jobs: Vec::new(), + scopes: 1, + built: None, + } +} + +#[test] +fn an_automatic_engine_answers_an_explicit_build_like_an_on_demand_one() { + for status in [ConsolidateStatus::Started, ConsolidateStatus::Completed] { + answers_as_promised(Consolidation::Automatic, Ok(receipt(status))).unwrap(); + answers_as_promised(Consolidation::OnDemand, Ok(receipt(status))).unwrap(); + } +} + +#[test] +fn an_automatic_engine_may_not_merely_acknowledge_an_explicit_build() { + let outcome = answers_as_promised( + Consolidation::Automatic, + Ok(ConsolidateReceipt::scheduled()), + ); + assert!(matches!(outcome, Err(Error::Check { .. })), "{outcome:?}"); +} diff --git a/crates/tinymemory-api/src/consolidate/mod.rs b/crates/tinymemory-api/src/consolidate/mod.rs index 44f8bc56..47225e5b 100644 --- a/crates/tinymemory-api/src/consolidate/mod.rs +++ b/crates/tinymemory-api/src/consolidate/mod.rs @@ -21,6 +21,11 @@ //! - [`Consolidation::Scheduled`] — the engine consolidates on its own //! schedule; `consolidate` acknowledges with //! [`ConsolidateStatus::Scheduled`] and does nothing more. +//! - [`Consolidation::Automatic`] — the engine rebuilds beliefs on its own +//! shortly after each write, so a host need not ask after its writes; +//! `consolidate` still runs a build at once, as a refresh before a read +//! that must see the latest, and answers as [`Consolidation::OnDemand`] +//! does. //! //! **Reading beliefs.** An engine that keeps what it builds apart from its //! stored items (CortexDB's belief layer) serves it through @@ -49,6 +54,9 @@ pub enum Consolidation { OnDemand, /// The engine consolidates on its own schedule. Scheduled, + /// The engine rebuilds beliefs on its own after writes; an explicit + /// [`crate::MemoryEngine::consolidate`] still runs a build at once. + Automatic, } /// What to consolidate. diff --git a/crates/tinymemory-api/src/consolidate/mod_tests.rs b/crates/tinymemory-api/src/consolidate/mod_tests.rs index 53c99c20..5bff771f 100644 --- a/crates/tinymemory-api/src/consolidate/mod_tests.rs +++ b/crates/tinymemory-api/src/consolidate/mod_tests.rs @@ -51,3 +51,11 @@ fn round_trips_through_json() { ); assert_eq!(Consolidation::default(), Consolidation::None); } + +#[test] +fn automatic_consolidation_round_trips_by_its_snake_case_name() { + let json = serde_json::to_value(Consolidation::Automatic).unwrap(); + assert_eq!(json, serde_json::json!("automatic")); + let back: Consolidation = serde_json::from_value(json).unwrap(); + assert_eq!(back, Consolidation::Automatic); +} diff --git a/crates/tinymemory-api/src/engine/mod.rs b/crates/tinymemory-api/src/engine/mod.rs index 6fa51509..71665034 100644 --- a/crates/tinymemory-api/src/engine/mod.rs +++ b/crates/tinymemory-api/src/engine/mod.rs @@ -138,8 +138,8 @@ pub trait MemoryEngine: Send + Sync { /// builds surfaces through ordinary reads. /// /// The default refuses: an engine declaring - /// [`Consolidation::OnDemand`] or [`Consolidation::Scheduled`] overrides - /// it. + /// [`Consolidation::OnDemand`], [`Consolidation::Scheduled`] or + /// [`Consolidation::Automatic`] overrides it. /// /// # Errors /// diff --git a/crates/tinymemory-api/tests/conformance_reference.rs b/crates/tinymemory-api/tests/conformance_reference.rs index 53e804ae..a67e5ade 100644 --- a/crates/tinymemory-api/tests/conformance_reference.rs +++ b/crates/tinymemory-api/tests/conformance_reference.rs @@ -20,6 +20,15 @@ async fn the_reference_engine_passes_and_cleans_up() { assert!(engine.is_empty()); } +#[tokio::test] +async fn an_engine_that_builds_on_its_own_passes_by_still_building_on_request() { + let engine = ReferenceEngine::new().with_consolidation(Consolidation::Automatic); + run(&engine) + .await + .expect("an automatic engine conforms when an explicit build runs"); + assert!(engine.is_empty()); +} + #[tokio::test] async fn the_suite_leaves_foreign_items_alone() { let engine = ReferenceEngine::new(); diff --git a/crates/tinymemory-integrations/examples/cortex_agent.rs b/crates/tinymemory-integrations/examples/cortex_agent.rs index d3c7850c..e88b6c89 100644 --- a/crates/tinymemory-integrations/examples/cortex_agent.rs +++ b/crates/tinymemory-integrations/examples/cortex_agent.rs @@ -85,7 +85,7 @@ async fn main() -> Result<(), Error> { let source = document.source.clone(); let ingested = brain.ingest(document).await?; took(&format!("ingest {file} -> source:{source}"), started); - jobs.push(ingested.job); + jobs.extend(ingested.job); } let policy = RecallPolicy { diff --git a/crates/tinymemory-integrations/examples/memory_eval/main.rs b/crates/tinymemory-integrations/examples/memory_eval/main.rs index e714edaa..b68cd689 100644 --- a/crates/tinymemory-integrations/examples/memory_eval/main.rs +++ b/crates/tinymemory-integrations/examples/memory_eval/main.rs @@ -390,7 +390,7 @@ impl Eval { .ingest(BrainDocument::new(source.clone(), *text).titled(*title)) .await?; timings.add("brain ingest (visible)", ms(started)); - jobs.push(ingested.job); + jobs.extend(ingested.job); *writes.entry(tenant).or_default() += 1; } Step::Learning { diff --git a/crates/tinymemory-integrations/src/config/mod.rs b/crates/tinymemory-integrations/src/config/mod.rs index e8ccf69b..bf031c3b 100644 --- a/crates/tinymemory-integrations/src/config/mod.rs +++ b/crates/tinymemory-integrations/src/config/mod.rs @@ -12,6 +12,7 @@ //! //! [engines.cortexdb] //! endpoint = "https://cortex.example.com" +//! consolidation = "automatic" # optional: this server builds beliefs itself //! ``` //! //! `engines` is optional, an engine with no entry uses its defaults, and a @@ -22,7 +23,7 @@ use std::collections::BTreeMap; use std::sync::Arc; use serde::{Deserialize, Serialize}; -use tinymemory_api::{MemoryEngine, Result}; +use tinymemory_api::{Consolidation, MemoryEngine, Result}; use crate::registry::{EngineCredential, build_engine}; @@ -77,6 +78,13 @@ pub struct EngineSettings { /// `Authorization` and the other headers it sets itself. #[serde(default, skip_serializing_if = "BTreeMap::is_empty")] pub headers: BTreeMap, + /// How the engine consolidates, overriding its default: a `cortexdb` + /// engine is `automatic` on CortexDB's managed API and `on_demand` + /// anywhere else, and `tinyhumans` is always `scheduled`. Set + /// `automatic` for a self-hosted CortexDB that runs its own layer + /// scheduler. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub consolidation: Option, } #[cfg(test)] diff --git a/crates/tinymemory-integrations/src/config/mod_tests.rs b/crates/tinymemory-integrations/src/config/mod_tests.rs index 21d7cd1b..4e787355 100644 --- a/crates/tinymemory-integrations/src/config/mod_tests.rs +++ b/crates/tinymemory-integrations/src/config/mod_tests.rs @@ -41,3 +41,22 @@ fn building_from_config_applies_the_registry_rules() { Err(tinymemory_api::Error::Config(_)) )); } + +#[test] +fn a_consolidation_override_is_read_by_its_snake_case_name() { + let config: MemoryConfig = toml::from_str( + r#" + engine = "cortexdb" + + [engines.cortexdb] + endpoint = "http://127.0.0.1:3141" + consolidation = "automatic" + "#, + ) + .unwrap(); + assert_eq!( + config.settings().consolidation, + Some(Consolidation::Automatic) + ); + assert_eq!(EngineSettings::default().consolidation, None); +} diff --git a/crates/tinymemory-integrations/src/cortex/README.md b/crates/tinymemory-integrations/src/cortex/README.md index 8464e5bb..b798bfc3 100644 --- a/crates/tinymemory-integrations/src/cortex/README.md +++ b/crates/tinymemory-integrations/src/cortex/README.md @@ -41,8 +41,12 @@ From `tinymemory_integrations::cortex`: Beyond the contract's reads and writes, the engine consolidates: Direct posts `v1/beliefs/build` once per held scope a `ConsolidateRequest` admits -(`engine/consolidate.rs`, declared `Consolidation::OnDemand`) and reports the -beliefs built, since the server builds within the request; hosted +(`engine/consolidate.rs`) and reports the beliefs built, since the server +builds within the request. It declares `Consolidation::Automatic` on the +managed API, which rebuilds beliefs on its own after writes, so the lifecycle +queues no build there, and `Consolidation::OnDemand` on any other endpoint +(`descriptor::direct_consolidation`, overridden by +`CortexEngine::with_consolidation`); hosted declares `Consolidation::Scheduled` and sends nothing. What was built is read back by `beliefs` (`engine/beliefs.rs`): a `beliefs`-only recall per held scope for a query, or the `v1/beliefs` listing without one, each belief diff --git a/crates/tinymemory-integrations/src/cortex/descriptor/mod.rs b/crates/tinymemory-integrations/src/cortex/descriptor/mod.rs index 561a03cc..db5323d4 100644 --- a/crates/tinymemory-integrations/src/cortex/descriptor/mod.rs +++ b/crates/tinymemory-integrations/src/cortex/descriptor/mod.rs @@ -19,10 +19,25 @@ //! //! # Consolidation //! -//! CortexDB builds beliefs on demand at `v1/beliefs/build`, one scope per -//! call, so the direct descriptor declares [`Consolidation::OnDemand`]. The -//! TinyHumans backend exposes no build route; CortexDB's own scheduler -//! consolidates behind it, so the hosted descriptor declares +//! CortexDB builds a scope's beliefs at `v1/beliefs/build`, one scope per +//! call. Whether a host also has to ask after its writes depends on the +//! deployment, not the wire ([`direct_consolidation`]): +//! +//! - CortexDB's managed API ([`CORTEX_API_ENDPOINT`]) extracts facts and +//! rebuilds a scope's beliefs on its own shortly after each write, so a +//! direct engine there declares [`Consolidation::Automatic`]: no build is +//! queued after a turn or an ingest, and an explicit build is a refresh. +//! - A self-hosted CortexDB builds beliefs in the background only when its +//! operator turns the layer scheduler on (`CORTEX_V1_LAYERS_AUTO`, off by +//! default), so a direct engine anywhere else declares +//! [`Consolidation::OnDemand`]. +//! +//! [`cortexdb_descriptor`] describes the engine at its default endpoint, the +//! managed API. A host that knows better overrides the choice with +//! [`crate::cortex::CortexEngine::with_consolidation`]. +//! +//! The TinyHumans backend exposes no build route; CortexDB consolidates +//! behind it on its own, so the hosted descriptor declares //! [`Consolidation::Scheduled`]. use tinymemory_api::{Consolidation, EngineDescriptor, FetchMode}; @@ -44,7 +59,8 @@ const FETCH_MODES: [FetchMode; 1] = [FetchMode::Hybrid]; /// The descriptor of CortexDB reached directly: not hosted by a third party, /// an endpoint is optional ([`CORTEX_API_ENDPOINT`] by default), an API key -/// is required, fetch is hybrid only, and beliefs build on demand. +/// is required, fetch is hybrid only, and beliefs build as they do at the +/// default endpoint: on their own ([`Consolidation::Automatic`]). #[must_use] pub fn cortexdb_descriptor() -> EngineDescriptor { EngineDescriptor { @@ -57,7 +73,7 @@ pub fn cortexdb_descriptor() -> EngineDescriptor { needs_key: true, default_endpoint: Some(CORTEX_API_ENDPOINT), fetch_modes: FETCH_MODES.to_vec(), - consolidation: Consolidation::OnDemand, + consolidation: direct_consolidation(CORTEX_API_ENDPOINT), } } @@ -80,6 +96,25 @@ pub fn tinyhumans_descriptor() -> EngineDescriptor { } } +/// How a direct engine whose endpoint has `origin` (scheme, host and port, +/// with no trailing slash) consolidates by +/// default: [`Consolidation::Automatic`] on CortexDB's managed API, +/// [`Consolidation::OnDemand`] anywhere else. +/// +/// This is the one place the choice is made. It goes by the endpoint alone +/// because a deployment's layer settings are not readable through the public +/// API today; if CortexDB confirms a readiness signal it can be asked here +/// instead (for example whether `v1/derivation/status` reports a running +/// layer scheduler), and every caller follows. +#[must_use] +pub(crate) fn direct_consolidation(origin: &str) -> Consolidation { + if origin == CORTEX_API_ENDPOINT { + Consolidation::Automatic + } else { + Consolidation::OnDemand + } +} + /// Which HTTP surface an engine talks to. #[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)] pub enum CortexWire { diff --git a/crates/tinymemory-integrations/src/cortex/descriptor/mod_tests.rs b/crates/tinymemory-integrations/src/cortex/descriptor/mod_tests.rs index a0c51a7d..f202009f 100644 --- a/crates/tinymemory-integrations/src/cortex/descriptor/mod_tests.rs +++ b/crates/tinymemory-integrations/src/cortex/descriptor/mod_tests.rs @@ -45,3 +45,32 @@ fn every_hosted_route_is_under_memory_and_every_direct_one_under_v1() { assert_eq!(CortexWire::TinyHumans.descriptor().id, TINYHUMANS_ENGINE_ID); assert_eq!(CortexWire::Direct.descriptor().id, CORTEXDB_ENGINE_ID); } + +#[test] +fn only_the_managed_api_consolidates_on_its_own_by_default() { + assert_eq!( + cortexdb_descriptor().consolidation, + Consolidation::Automatic, + "the direct descriptor describes its default endpoint, the managed API" + ); + assert_eq!( + tinyhumans_descriptor().consolidation, + Consolidation::Scheduled + ); + assert_eq!( + direct_consolidation("https://api-v1.cortexdb.ai"), + Consolidation::Automatic + ); + for self_hosted in [ + "http://127.0.0.1:3141", + "https://cortex.example.test", + "https://api-v1.cortexdb.ai:8443", + "http://api-v1.cortexdb.ai", + ] { + assert_eq!( + direct_consolidation(self_hosted), + Consolidation::OnDemand, + "{self_hosted}" + ); + } +} diff --git a/crates/tinymemory-integrations/src/cortex/engine/consolidate_tests.rs b/crates/tinymemory-integrations/src/cortex/engine/consolidate_tests.rs index 311cd35e..eb017444 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/consolidate_tests.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/consolidate_tests.rs @@ -1,10 +1,12 @@ //! Belief builds: one `v1/beliefs/build` per held scope on Direct, none on -//! the hosted wire. +//! the hosted wire; and which consolidation each engine declares. use super::*; +use crate::cortex::CortexCredential; use crate::cortex::testing::{both, direct_double, direct_engine, sample_items}; use tinymemory_api::{ - ConsolidateRequest, ItemKind, MemoryEngine, MemoryMeta, Namespace, Reach, StoreItem, + ConsolidateRequest, Consolidation, ItemKind, MemoryEngine, MemoryMeta, Namespace, Reach, + StoreItem, }; fn at(namespace: Namespace) -> MemoryMeta { @@ -159,3 +161,82 @@ async fn a_malformed_request_is_refused_before_any_request() { ); assert!(state.requests().is_empty()); } + +#[test] +fn a_direct_engine_consolidates_as_its_endpoint_does_unless_told_otherwise() { + let key = || CortexCredential::api_key("key"); + for managed in ["https://api-v1.cortexdb.ai", "https://api-v1.cortexdb.ai/"] { + let engine = CortexEngine::direct(managed, key()).unwrap(); + assert_eq!( + engine.descriptor().consolidation, + Consolidation::Automatic, + "{managed}" + ); + } + let self_hosted = CortexEngine::direct("http://127.0.0.1:3141", key()).unwrap(); + assert_eq!( + self_hosted.descriptor().consolidation, + Consolidation::OnDemand + ); + let told = self_hosted + .with_consolidation(Consolidation::Automatic) + .unwrap(); + assert_eq!(told.descriptor().consolidation, Consolidation::Automatic); + let hosted = CortexEngine::tinyhumans( + "https://api.example.test", + std::sync::Arc::new(crate::cortex::StaticBearer::new("jwt")), + ) + .unwrap(); + assert_eq!(hosted.descriptor().consolidation, Consolidation::Scheduled); +} + +#[test] +fn a_consolidation_the_wire_cannot_serve_is_refused() { + let direct = || CortexEngine::direct("http://127.0.0.1:3141", CortexCredential::api_key("k")); + for refused in [Consolidation::None, Consolidation::Scheduled] { + let error = direct().unwrap().with_consolidation(refused).unwrap_err(); + assert!( + matches!(error, crate::cortex::Error::Config(_)), + "{refused:?}: {error:?}" + ); + } + let hosted = CortexEngine::tinyhumans( + "https://api.example.test", + std::sync::Arc::new(crate::cortex::StaticBearer::new("jwt")), + ) + .unwrap(); + for refused in [ + Consolidation::None, + Consolidation::OnDemand, + Consolidation::Automatic, + ] { + assert!( + hosted.clone().with_consolidation(refused).is_err(), + "{refused:?}" + ); + } + hosted.with_consolidation(Consolidation::Scheduled).unwrap(); +} + +#[tokio::test] +async fn an_explicit_build_still_runs_on_an_automatic_engine() { + let (endpoint, state) = direct_double().await; + let engine = direct_engine(&endpoint) + .with_consolidation(Consolidation::Automatic) + .unwrap(); + engine + .store(StoreItem::document( + "Refunds take five days.", + at(Namespace::source("pdf")), + )) + .await + .unwrap(); + let receipt = engine + .consolidate(ConsolidateRequest::new(Reach::exact(Namespace::source( + "pdf", + )))) + .await + .unwrap(); + assert_eq!(receipt.status, ConsolidateStatus::Completed); + assert_eq!(state.count("POST /v1/beliefs/build"), 1); +} diff --git a/crates/tinymemory-integrations/src/cortex/engine/mod.rs b/crates/tinymemory-integrations/src/cortex/engine/mod.rs index 88b20c4c..b68b1b10 100644 --- a/crates/tinymemory-integrations/src/cortex/engine/mod.rs +++ b/crates/tinymemory-integrations/src/cortex/engine/mod.rs @@ -27,13 +27,14 @@ use std::sync::Arc; use async_trait::async_trait; use tinymemory_api::{ - BeliefsRequest, ConsolidateReceipt, ConsolidateRequest, EngineDescriptor, EngineHealth, - FetchPage, FetchRequest, ForgetReport, ForgetTarget, GetRequest, Hit, ListPage, ListRequest, - MemoryEngine, RecallAnswer, RecallRequest, StoreItem, StoreReceipt, WaitFor, WriteOptions, + BeliefsRequest, ConsolidateReceipt, ConsolidateRequest, Consolidation, EngineDescriptor, + EngineHealth, FetchPage, FetchRequest, ForgetReport, ForgetTarget, GetRequest, Hit, ListPage, + ListRequest, MemoryEngine, RecallAnswer, RecallRequest, StoreItem, StoreReceipt, WaitFor, + WriteOptions, }; use crate::cortex::credential::{BearerSource, CortexCredential}; -use crate::cortex::descriptor::{CortexWire, Route}; +use crate::cortex::descriptor::{CortexWire, Route, direct_consolidation}; use crate::cortex::error::{Error, Result}; use crate::cortex::log::Log; use crate::cortex::transport::{HttpClient, health_reason, urlencode}; @@ -72,13 +73,52 @@ impl CortexEngine { /// [`Error::Config`] for an invalid or non-HTTP(S) endpoint, a cleartext /// endpoint that is not loopback (the credential would cross the network /// in the clear), or a blank static credential. + /// + /// A direct engine consolidates as its endpoint does by default: + /// [`Consolidation::Automatic`] on CortexDB's managed API, + /// [`Consolidation::OnDemand`] on any other (see + /// [`CortexEngine::with_consolidation`]). pub fn new(wire: CortexWire, endpoint: &str, credential: CortexCredential) -> Result { + let client = HttpClient::new(wire, endpoint, credential)?; + let mut descriptor = wire.descriptor(); + if wire == CortexWire::Direct { + descriptor.consolidation = direct_consolidation(&client.origin()); + } Ok(Self { - descriptor: wire.descriptor(), - log: Log::new(HttpClient::new(wire, endpoint, credential)?), + descriptor, + log: Log::new(client), }) } + /// The same engine, declaring `consolidation` instead of the endpoint's + /// default: for a self-hosted CortexDB whose operator runs the layer + /// scheduler ([`Consolidation::Automatic`]), or a managed one a host + /// wants to build by hand ([`Consolidation::OnDemand`]). Either way an + /// explicit [`MemoryEngine::consolidate`] still builds. + /// + /// # Errors + /// + /// [`Error::Config`] for a mode the wire cannot serve: a direct engine + /// is [`Consolidation::OnDemand`] or [`Consolidation::Automatic`], and + /// the TinyHumans backend only [`Consolidation::Scheduled`]. + pub fn with_consolidation(mut self, consolidation: Consolidation) -> Result { + let served = match self.wire() { + CortexWire::Direct => matches!( + consolidation, + Consolidation::OnDemand | Consolidation::Automatic + ), + CortexWire::TinyHumans => consolidation == Consolidation::Scheduled, + }; + if !served { + return Err(Error::Config(format!( + "the `{}` engine cannot consolidate as {consolidation:?}", + self.descriptor.id + ))); + } + self.descriptor.consolidation = consolidation; + Ok(self) + } + /// CortexDB's own `/v1/*` API at `endpoint` (for example /// [`crate::cortex::CORTEX_API_ENDPOINT`]), registered as `cortexdb`. /// @@ -202,8 +242,10 @@ impl MemoryEngine for CortexEngine { self.get_items(req).await } - /// Direct: one `v1/beliefs/build` per held scope in reach. Hosted: - /// acknowledged as scheduled, with no request. + /// Direct: one `v1/beliefs/build` per held scope in reach, whether the + /// engine declares [`Consolidation::OnDemand`] or + /// [`Consolidation::Automatic`]. Hosted: acknowledged as scheduled, with + /// no request. async fn consolidate(&self, req: ConsolidateRequest) -> Result { self.build_beliefs(req).await } diff --git a/crates/tinymemory-integrations/src/cortex/lifecycle_tests.rs b/crates/tinymemory-integrations/src/cortex/lifecycle_tests.rs index f70d5664..c24a4c3c 100644 --- a/crates/tinymemory-integrations/src/cortex/lifecycle_tests.rs +++ b/crates/tinymemory-integrations/src/cortex/lifecycle_tests.rs @@ -96,7 +96,7 @@ async fn an_agent_loop_runs_the_same_on_either_wire() { resumed.markdown ); - for job in report.jobs.into_iter().chain([ingested.job]) { + for job in report.jobs.into_iter().chain(ingested.job) { let ran = support.run_background(job).await.unwrap(); match wire { CortexWire::Direct => assert_eq!(ran.outcome, JobOutcome::Done), diff --git a/crates/tinymemory-integrations/src/registry/mod.rs b/crates/tinymemory-integrations/src/registry/mod.rs index d12726a7..5b747419 100644 --- a/crates/tinymemory-integrations/src/registry/mod.rs +++ b/crates/tinymemory-integrations/src/registry/mod.rs @@ -63,8 +63,9 @@ pub fn list_engines() -> Vec { /// # Errors /// /// [`Error::Config`] for an unknown id, a missing endpoint or credential, an -/// endpoint that is not an HTTP(S) URL, or a credentialed cleartext -/// (`http://`) endpoint that is not loopback. +/// endpoint that is not an HTTP(S) URL, a credentialed cleartext +/// (`http://`) endpoint that is not loopback, or a consolidation the engine +/// cannot serve ([`CortexEngine::with_consolidation`]). pub fn build_engine( id: &str, settings: &EngineSettings, @@ -93,8 +94,11 @@ pub fn build_engine( ))); } }; - let engine = + let mut engine = CortexEngine::new(wire, endpoint, credential)?.with_default_headers(&settings.headers)?; + if let Some(consolidation) = settings.consolidation { + engine = engine.with_consolidation(consolidation)?; + } Ok(Arc::new(engine)) } diff --git a/crates/tinymemory-integrations/src/registry/mod_tests.rs b/crates/tinymemory-integrations/src/registry/mod_tests.rs index 70d97f98..da9804dd 100644 --- a/crates/tinymemory-integrations/src/registry/mod_tests.rs +++ b/crates/tinymemory-integrations/src/registry/mod_tests.rs @@ -172,3 +172,37 @@ fn fixed_headers_are_applied_and_a_credential_header_is_refused() { assert!(matches!(refused, Error::Config(_)), "{refused:?}"); assert!(!refused.to_string().contains("smuggled")); } + +#[test] +fn a_configured_consolidation_overrides_the_endpoint_default() { + let self_hosted = settings(Some("https://cortex.example.test")); + let key = || EngineCredential::Static("key".into()); + let default = build_engine("cortexdb", &self_hosted, key()).unwrap(); + assert_eq!( + default.descriptor().consolidation, + tinymemory_api::Consolidation::OnDemand + ); + let managed = build_engine("cortexdb", &settings(None), key()).unwrap(); + assert_eq!( + managed.descriptor().consolidation, + tinymemory_api::Consolidation::Automatic + ); + let told = EngineSettings { + consolidation: Some(tinymemory_api::Consolidation::Automatic), + ..self_hosted + }; + let automatic = build_engine("cortexdb", &told, key()).unwrap(); + assert_eq!( + automatic.descriptor().consolidation, + tinymemory_api::Consolidation::Automatic + ); + let message = config_error(build_engine( + "tinyhumans", + &EngineSettings { + consolidation: Some(tinymemory_api::Consolidation::OnDemand), + ..settings(None) + }, + EngineCredential::Dynamic(Arc::new(Session)), + )); + assert!(message.contains("cannot consolidate"), "{message}"); +} diff --git a/crates/tinymemory-integrations/tests/live_cortex_lifecycle.rs b/crates/tinymemory-integrations/tests/live_cortex_lifecycle.rs index 0656d4a1..bdbb423e 100644 --- a/crates/tinymemory-integrations/tests/live_cortex_lifecycle.rs +++ b/crates/tinymemory-integrations/tests/live_cortex_lifecycle.rs @@ -83,10 +83,9 @@ async fn live_an_agent_loop_runs_against_cortexdb() { ) .await .expect("convert"); - let ingested = Brain::new(engine.clone(), layout.clone()) - .ingest(document) - .await - .expect("ingest"); + let source = document.source.clone(); + let brain = Brain::new(engine.clone(), layout.clone()); + let ingested = brain.ingest(document).await.expect("ingest"); let support = AgentMemory::new(engine.clone(), layout.clone(), "support-01").expect("agent"); let coder = AgentMemory::new(engine.clone(), layout.clone(), "coder-42").expect("agent"); @@ -134,7 +133,12 @@ async fn live_an_agent_loop_runs_against_cortexdb() { ); let built = support - .run_background(ingested.job) + // The managed API builds on its own and hands back no job; ask for + // one explicitly, as a refresh, so the build route is still proven. + .run_background(match ingested.job { + Some(job) => job, + None => brain.build(&source).expect("build job"), + }) .await .expect("a belief build is accepted"); eprintln!("belief build: {:?}", built.outcome); diff --git a/crates/tinymemory-tools/examples/brain.rs b/crates/tinymemory-tools/examples/brain.rs index 311b6496..d78cb5dd 100644 --- a/crates/tinymemory-tools/examples/brain.rs +++ b/crates/tinymemory-tools/examples/brain.rs @@ -10,7 +10,7 @@ use std::sync::Arc; use tinymemory_api::conformance::ReferenceEngine; -use tinymemory_tools::{Brain, BrainDocument, BrainSource, MemoryLayout}; +use tinymemory_tools::{BackgroundJob, Brain, BrainDocument, BrainSource, MemoryLayout}; #[tokio::main(flavor = "current_thread")] async fn main() -> Result<(), Box> { @@ -47,7 +47,7 @@ async fn main() -> Result<(), Box> { "{:<9} -> {} (then: {})", source.to_string(), layout.brain(&source)?, - ingested.job.name() + ingested.job.as_ref().map_or("nothing", BackgroundJob::name) ); } diff --git a/crates/tinymemory-tools/src/background/mod.rs b/crates/tinymemory-tools/src/background/mod.rs index 43a08eb7..0f4667a0 100644 --- a/crates/tinymemory-tools/src/background/mod.rs +++ b/crates/tinymemory-tools/src/background/mod.rs @@ -22,9 +22,18 @@ use tinymemory_api::{ StoreReceipt, }; +use tinymemory_api::Consolidation; + use crate::brain::{Brain, BrainDocument}; use crate::layout::MemoryLayout; +/// Whether `engine` rebuilds beliefs on its own after writes +/// ([`Consolidation::Automatic`]), so the lifecycle and the brain hand back +/// no belief build after a turn or an ingest. +pub(crate) fn builds_on_its_own(engine: &dyn MemoryEngine) -> bool { + engine.descriptor().consolidation == Consolidation::Automatic +} + /// One unit of deferred work. #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] #[serde(tag = "job", rename_all = "snake_case")] diff --git a/crates/tinymemory-tools/src/brain/mod.rs b/crates/tinymemory-tools/src/brain/mod.rs index c5e3486d..77e82171 100644 --- a/crates/tinymemory-tools/src/brain/mod.rs +++ b/crates/tinymemory-tools/src/brain/mod.rs @@ -9,7 +9,10 @@ //! [`Brain::ingest`] stores one document (by default waiting until it is //! readable, since ingestion is not on a live turn) and returns the //! [`BackgroundJob`] that would build beliefs from its source — the host runs -//! it whenever suits, through [`crate::BackgroundRunner`]. Converting a +//! it whenever suits, through [`crate::BackgroundRunner`]. An engine that +//! rebuilds beliefs on its own after writes +//! ([`tinymemory_api::Consolidation::Automatic`]) gets no such job; +//! [`Brain::build`] still asks for one at any time. Converting a //! file's bytes to text is the integrations crate's job; the brain takes //! text. //! @@ -47,7 +50,7 @@ use tinymemory_api::{ MAX_STORE_MANY, MemoryEngine, Reach, Result, StoreItem, WriteOptions, }; -use crate::background::BackgroundJob; +use crate::background::{BackgroundJob, builds_on_its_own}; use crate::layout::{BrainSource, MemoryLayout}; pub use types::{BrainBatch, BrainDocument, Ingested}; @@ -104,10 +107,12 @@ impl Brain { let source = document.source.clone(); let item = document.into_item(&self.layout)?; let receipt = self.engine.store_with(item, options).await?; - Ok(Ingested { - receipt, - job: self.build_job(&source)?, - }) + let job = if builds_on_its_own(self.engine.as_ref()) { + None + } else { + Some(self.build(&source)?) + }; + Ok(Ingested { receipt, job }) } /// Stores many documents, in order, in batches of at most @@ -132,10 +137,14 @@ impl Brain { receipts.extend(self.engine.store_many(rest).await?); rest = tail; } - let jobs = sources - .iter() - .map(|source| self.build_job(source)) - .collect::>>()?; + let jobs = if builds_on_its_own(self.engine.as_ref()) { + Vec::new() + } else { + sources + .iter() + .map(|source| self.build(source)) + .collect::>>()? + }; Ok(BrainBatch { receipts, jobs }) } @@ -179,8 +188,15 @@ impl Brain { self.engine.forget(ForgetTarget::Filter(filter)).await } - /// The belief build for `source`'s documents. - fn build_job(&self, source: &BrainSource) -> Result { + /// A belief build of `source`'s documents, for the host to run now or + /// queue: what an ingest hands back, and a refresh on an engine that + /// builds on its own. + /// + /// # Errors + /// + /// Never for a layout built by [`MemoryLayout::new`] (see + /// [`MemoryLayout::brain`]). + pub fn build(&self, source: &BrainSource) -> Result { Ok(BackgroundJob::BuildBeliefs { request: ConsolidateRequest::new(Reach::exact(self.layout.brain(source)?)) .kinds([ItemKind::Document]), diff --git a/crates/tinymemory-tools/src/brain/mod_tests.rs b/crates/tinymemory-tools/src/brain/mod_tests.rs index 7cc1f27c..c7f59ffe 100644 --- a/crates/tinymemory-tools/src/brain/mod_tests.rs +++ b/crates/tinymemory-tools/src/brain/mod_tests.rs @@ -2,7 +2,8 @@ use tinymemory_api::conformance::ReferenceEngine; use tinymemory_api::{ - ConsolidateRequest, Error, ListRequest, MemoryMeta, Namespace, SourceKind, SourceRef, + ConsolidateRequest, Consolidation, Error, ListRequest, MemoryMeta, Namespace, SourceKind, + SourceRef, }; use super::*; @@ -31,10 +32,10 @@ async fn a_document_lands_at_its_source_node_without_an_agent() { .unwrap(); assert_eq!( ingested.job, - BackgroundJob::BuildBeliefs { + Some(BackgroundJob::BuildBeliefs { request: ConsolidateRequest::new(Reach::exact(Namespace::source("pdf"))) .kinds([ItemKind::Document]), - } + }) ); let listed = engine .list(ListRequest::new(Default::default(), 10)) @@ -137,3 +138,40 @@ async fn search_and_forget_stay_inside_one_source() { .is_empty() ); } + +#[tokio::test] +async fn an_engine_that_builds_on_its_own_gets_no_build_from_an_ingest() { + let engine = Arc::new(ReferenceEngine::new().with_consolidation(Consolidation::Automatic)); + let brain = Brain::new(engine.clone(), MemoryLayout::default()); + let ingested = brain + .ingest(BrainDocument::new( + BrainSource::Pdf, + "Refunds take five days.", + )) + .await + .unwrap(); + assert_eq!(ingested.job, None); + let batch = brain + .ingest_many(vec![ + BrainDocument::new(BrainSource::Notion, "Deploys run on Fridays."), + BrainDocument::new(BrainSource::Github, "CI runs on every push."), + ]) + .await + .unwrap(); + assert!(batch.jobs.is_empty(), "{:?}", batch.jobs); + assert_eq!(engine.len(), 3, "the documents are still stored"); + + let refresh = brain.build(&BrainSource::Pdf).unwrap(); + assert_eq!( + refresh, + BackgroundJob::BuildBeliefs { + request: ConsolidateRequest::new(Reach::exact(Namespace::source("pdf"))) + .kinds([ItemKind::Document]), + } + ); + let report = crate::background::BackgroundRunner::new(engine, MemoryLayout::default()) + .run(refresh) + .await + .unwrap(); + assert_eq!(report.outcome, crate::background::JobOutcome::Done); +} diff --git a/crates/tinymemory-tools/src/brain/types.rs b/crates/tinymemory-tools/src/brain/types.rs index ce13ae2c..786f9c58 100644 --- a/crates/tinymemory-tools/src/brain/types.rs +++ b/crates/tinymemory-tools/src/brain/types.rs @@ -90,8 +90,9 @@ pub struct Ingested { /// The engine's receipt. pub receipt: StoreReceipt, /// The belief build for the document's source, for the host to run when - /// it suits. - pub job: BackgroundJob, + /// it suits; `None` on an engine that rebuilds beliefs on its own + /// ([`tinymemory_api::Consolidation::Automatic`]). + pub job: Option, } /// What [`crate::Brain::ingest_many`] did. @@ -99,6 +100,7 @@ pub struct Ingested { pub struct BrainBatch { /// One receipt per document, in order. pub receipts: Vec, - /// One belief build per source touched. + /// One belief build per source touched; none on an engine that rebuilds + /// beliefs on its own ([`tinymemory_api::Consolidation::Automatic`]). pub jobs: Vec, } diff --git a/crates/tinymemory-tools/src/lifecycle/mod.rs b/crates/tinymemory-tools/src/lifecycle/mod.rs index 13db7698..50d3f9fa 100644 --- a/crates/tinymemory-tools/src/lifecycle/mod.rs +++ b/crates/tinymemory-tools/src/lifecycle/mod.rs @@ -9,7 +9,7 @@ //! | --- | --- | --- | //! | session start or resume | [`AgentMemory::start_session`] | reads only | //! | user turn, before the model | [`AgentMemory::pre_turn`] | logs the turn (accepted, not indexed) while fetching the pack | -//! | after the reply | [`AgentMemory::post_turn`] | logs the reply; may return a belief build | +//! | after the reply | [`AgentMemory::post_turn`] | logs the reply; may return a belief build (never on an engine that builds on its own) | //! | prompt truncated | [`AgentMemory::recall_for_compaction`] | an answered summary of the thread, plus related memory | //! | any time | [`AgentMemory::recall`] | a pre-turn pack without logging | //! | off the turn | [`AgentMemory::run_background`] | the job | @@ -79,7 +79,7 @@ use tinymemory_api::{ Result, Role, SourceKind, SourceRef, StoreItem, StoreReceipt, Turn, TurnRange, WriteOptions, }; -use crate::background::{BackgroundJob, BackgroundRunner, JobReport}; +use crate::background::{BackgroundJob, BackgroundRunner, JobReport, builds_on_its_own}; use crate::brain::Brain; use crate::layout::{CoreScope, MemoryLayout}; use crate::recall::{ @@ -397,7 +397,9 @@ impl AgentMemory { } /// Logs the assistant's reply, and returns the belief build the policy - /// asks for at this turn, if any. + /// asks for at this turn, if any. An engine that rebuilds beliefs on its + /// own ([`tinymemory_api::Consolidation::Automatic`]) gets none; + /// [`AgentMemory::history_build`] still asks for one at any time. /// /// # Errors /// @@ -419,10 +421,11 @@ impl AgentMemory { .engine .store_with(item, WriteOptions::accepted()) .await?; - let due = self - .policy - .build_beliefs_every - .is_some_and(|every| every > 0 && (turn.turn_index + 1).is_multiple_of(every)); + let due = !builds_on_its_own(self.engine.as_ref()) + && self + .policy + .build_beliefs_every + .is_some_and(|every| every > 0 && (turn.turn_index + 1).is_multiple_of(every)); let jobs = if due { vec![self.history_build()] } else { diff --git a/crates/tinymemory-tools/src/lifecycle/mod_tests.rs b/crates/tinymemory-tools/src/lifecycle/mod_tests.rs index 1d22c61d..4acaf855 100644 --- a/crates/tinymemory-tools/src/lifecycle/mod_tests.rs +++ b/crates/tinymemory-tools/src/lifecycle/mod_tests.rs @@ -3,8 +3,9 @@ use async_trait::async_trait; use tinymemory_api::conformance::ReferenceEngine; use tinymemory_api::{ - EngineDescriptor, EngineHealth, FetchPage, FetchRequest, ForgetReport, ForgetTarget, - LearningKind, ListPage, ListRequest, RecallAnswer, RecallRequest, StoreReceipt, ToolCallRef, + Consolidation, EngineDescriptor, EngineHealth, FetchPage, FetchRequest, ForgetReport, + ForgetTarget, LearningKind, ListPage, ListRequest, RecallAnswer, RecallRequest, StoreReceipt, + ToolCallRef, }; use super::*; @@ -137,6 +138,33 @@ async fn other_agents_turns_appear_once_under_the_team() { assert_eq!(md.matches("refund delayed").count(), 1, "shown once: {md}"); } +#[tokio::test] +async fn post_turn_asks_no_build_of_an_engine_that_builds_on_its_own() { + let engine = Arc::new(ReferenceEngine::new().with_consolidation(Consolidation::Automatic)); + let support = memory(&engine, "support-01").with_policy(RecallPolicy { + build_beliefs_every: Some(1), + ..RecallPolicy::default() + }); + for index in 0..3 { + let report = support + .post_turn(PostTurn::new("t1", index, format!("reply {index}"))) + .await + .unwrap(); + assert!(report.jobs.is_empty(), "turn {index}: {:?}", report.jobs); + } + assert_eq!(engine.len(), 3, "every reply is still logged"); + + let refresh = support + .run_background(support.history_build()) + .await + .unwrap(); + assert_eq!(refresh.job, "build_beliefs"); + assert!( + refresh.consolidation.is_some(), + "an explicit build still runs" + ); +} + #[tokio::test] async fn post_turn_asks_for_a_belief_build_on_the_policy_s_cadence() { let engine = Arc::new(ReferenceEngine::new()); diff --git a/crates/tinymemory-tools/src/lifecycle/types.rs b/crates/tinymemory-tools/src/lifecycle/types.rs index a39e2b9e..50d17a3b 100644 --- a/crates/tinymemory-tools/src/lifecycle/types.rs +++ b/crates/tinymemory-tools/src/lifecycle/types.rs @@ -27,7 +27,9 @@ pub struct RecallPolicy { /// conversations out. pub team_limit: usize, /// Ask for a belief build of this agent's conversations after every this - /// many turns (by `turn_index + 1`); `None` never asks. + /// many turns (by `turn_index + 1`); `None` never asks, and neither does + /// an engine that rebuilds beliefs on its own + /// ([`tinymemory_api::Consolidation::Automatic`]). pub build_beliefs_every: Option, } diff --git a/docs/architecture/lifecycle.md b/docs/architecture/lifecycle.md index 874d5cad..161d248d 100644 --- a/docs/architecture/lifecycle.md +++ b/docs/architecture/lifecycle.md @@ -99,6 +99,9 @@ duplicates. - `Brain::ingest` returns a `BuildBeliefs` job for the source's scope. - `post_turn` returns one every `build_beliefs_every` turns, for the agent's conversations. +- Neither does on an engine that declares `Automatic`: it rebuilds beliefs + on its own after writes. `Brain::build` and `AgentMemory::history_build` + still hand back a build to run at once, as a refresh. - `BackgroundJob::IngestBrain` defers a whole ingestion. Its report hands back the follow-up builds. @@ -107,7 +110,8 @@ What a build does depends on the engine's `consolidation`: | Engine | `consolidation` | `run_background(BuildBeliefs)` | | --- | --- | --- | | Reference | `OnDemand` | distils one `Fact` per item at once → `Done` | -| CortexDB direct | `OnDemand` | `POST v1/beliefs/build` per held scope, built within the request → `Done` | +| CortexDB direct, managed API | `Automatic` | none queued; an explicit build posts `v1/beliefs/build` per held scope → `Done` | +| CortexDB direct, elsewhere | `OnDemand` | `POST v1/beliefs/build` per held scope, built within the request → `Done` | | CortexDB via TinyHumans | `Scheduled` | nothing sent → `Scheduled` | | an engine without it | `None` | `Skipped { reason }` | diff --git a/docs/integration.md b/docs/integration.md index fbe6e9c8..8df617be 100644 --- a/docs/integration.md +++ b/docs/integration.md @@ -86,7 +86,7 @@ the `EngineCredential`'s alone. | Engine id | Where it runs | Consolidation (belief builds) | | --- | --- | --- | -| `cortexdb` | a CortexDB server, `v1/*` routes | on demand: `v1/beliefs/build`, built within the request | +| `cortexdb` | a CortexDB server, `v1/*` routes | automatic on the managed API (`api-v1.cortexdb.ai`), on demand elsewhere: `v1/beliefs/build`, built within the request; `EngineSettings::consolidation` overrides | | `tinyhumans` | the hosted TinyHumans backend, `memory/*` routes | on the server's own schedule | | `ReferenceEngine` | in process (tests) | on demand, a deterministic toy | @@ -206,7 +206,10 @@ for job in queue.drain(..) { - `post_turn` hands back a `BuildBeliefs` job every `RecallPolicy::build_beliefs_every` turns (10 by default). `Brain::ingest` - hands one back per ingest. + hands one back per ingest. Neither does on an engine that declares + `Consolidation::Automatic` (CortexDB's managed API rebuilds beliefs on its + own within about a minute of a write); `Brain::build` and + `AgentMemory::history_build` still ask for a refresh. - **On CortexDB, a build is only as good as the extraction before it.** Beliefs are built from facts, which CortexDB extracts from every event with a model, in the background. With a hosted model that takes minutes. diff --git a/docs/specs/agent-memory.md b/docs/specs/agent-memory.md index a35a84d5..9a80a24c 100644 --- a/docs/specs/agent-memory.md +++ b/docs/specs/agent-memory.md @@ -81,7 +81,8 @@ by default, or a host node such as `team:acme`. the receipt names any job handles, the number of scopes covered and, for a completed build that reports it, the beliefs `built`. - **`EngineDescriptor::consolidation`** declares how the engine consolidates: - `None`, `OnDemand` or `Scheduled`. + `None`, `OnDemand`, `Scheduled` or `Automatic` (it rebuilds beliefs on its + own after writes, and an explicit `consolidate` still builds at once). - **`MemoryEngine::beliefs(BeliefsRequest { reach, query, limit })`** reads the beliefs an engine built and keeps apart from its stored items. With a query they are ranked for it; without one, the most confident come first, @@ -174,7 +175,8 @@ skipped, engine }`. `TurnContext::log_error` and the pack is still returned. - **`post_turn` reports belief builds.** It returns a `BuildBeliefs` job for the agent's conversations every `RecallPolicy::build_beliefs_every` turns, - counted as `turn_index + 1`. + counted as `turn_index + 1`, unless the engine declares `Automatic`; then + it returns none, and `history_build` still asks for one. ### Brain and background @@ -202,6 +204,9 @@ skipped, engine }`. that is held in reach and admitted. The server builds within the request, so the receipt is `Completed` with the beliefs `built`; an answer that names a job instead makes it `Started` with the handles. + - It declares `Automatic` when its endpoint is CortexDB's managed API, + which rebuilds beliefs on its own after writes, and `OnDemand` + anywhere else. `EngineSettings::consolidation` overrides that. - **CortexDB, TinyHumans wire:** declares `Scheduled` and sends nothing. ## Invariants From f0f88a635cd57de55b0d5c871a016d5f968b9ac0 Mon Sep 17 00:00:00 2001 From: M3gA-Mind Date: Tue, 6 Oct 2026 17:43:16 +0530 Subject: [PATCH 2/4] Make the reference engine consolidate as it declares; cover Automatic live - `ReferenceEngine::consolidate` honours `with_consolidation`: `None` refuses as `Unsupported`, `Scheduled` only acknowledges, `OnDemand` and `Automatic` build. The conformance suite now passes for all four. - live_cortex_lifecycle: an engine configured with `consolidation = "automatic"` through the registry queues no build after an ingest or a turn, while `Brain::build` and `history_build` still reach `v1/beliefs/build` on the real server. - docs: an ingest's job is optional (`integration.md` example, `agent-memory.md` spec). --- .../src/conformance/reference/mod.rs | 24 +++-- .../tests/conformance_reference.rs | 19 ++-- .../tests/live_cortex_lifecycle.rs | 88 ++++++++++++++++++- docs/integration.md | 2 +- docs/specs/agent-memory.md | 5 +- 5 files changed, 123 insertions(+), 15 deletions(-) diff --git a/crates/tinymemory-api/src/conformance/reference/mod.rs b/crates/tinymemory-api/src/conformance/reference/mod.rs index a1b72d87..a55c3799 100644 --- a/crates/tinymemory-api/src/conformance/reference/mod.rs +++ b/crates/tinymemory-api/src/conformance/reference/mod.rs @@ -59,10 +59,12 @@ impl ReferenceEngine { } } - /// The same engine, declaring `consolidation`. Its builds run either - /// way, so it serves a host test of how a lifecycle treats an engine - /// that declares [`Consolidation::Automatic`]; the conformance suite - /// holds it to whatever it declares. + /// The same engine, declaring `consolidation`, and consolidating as it + /// declares: [`Consolidation::OnDemand`] and [`Consolidation::Automatic`] + /// build at once, [`Consolidation::Scheduled`] only acknowledges, and + /// [`Consolidation::None`] refuses. For a host test of how a lifecycle + /// treats each kind of engine; the conformance suite holds it to what + /// it declares. #[must_use] pub fn with_consolidation(mut self, consolidation: Consolidation) -> Self { self.descriptor.consolidation = consolidation; @@ -216,9 +218,21 @@ impl MemoryEngine for ReferenceEngine { } /// Distils one belief per admitted document or conversation, at once: - /// the build is [`ConsolidateStatus::Completed`] on return. + /// the build is [`ConsolidateStatus::Completed`] on return. It answers as + /// it declares ([`ReferenceEngine::with_consolidation`]): no + /// consolidation refuses, a scheduled one only acknowledges. async fn consolidate(&self, req: ConsolidateRequest) -> Result { req.validate()?; + match self.descriptor.consolidation { + Consolidation::None => { + return Err(Error::Unsupported(format!( + "engine `{}` does not consolidate memory", + self.descriptor.id + ))); + } + Consolidation::Scheduled => return Ok(ConsolidateReceipt::scheduled()), + Consolidation::OnDemand | Consolidation::Automatic => {} + } let mut items = self.items()?; let beliefs = distil::distil(&items, &req); let mut nodes: Vec<&crate::Namespace> = Vec::new(); diff --git a/crates/tinymemory-api/tests/conformance_reference.rs b/crates/tinymemory-api/tests/conformance_reference.rs index a67e5ade..cfc4d2e1 100644 --- a/crates/tinymemory-api/tests/conformance_reference.rs +++ b/crates/tinymemory-api/tests/conformance_reference.rs @@ -21,12 +21,19 @@ async fn the_reference_engine_passes_and_cleans_up() { } #[tokio::test] -async fn an_engine_that_builds_on_its_own_passes_by_still_building_on_request() { - let engine = ReferenceEngine::new().with_consolidation(Consolidation::Automatic); - run(&engine) - .await - .expect("an automatic engine conforms when an explicit build runs"); - assert!(engine.is_empty()); +async fn the_reference_engine_conforms_whatever_consolidation_it_declares() { + for declared in [ + Consolidation::None, + Consolidation::OnDemand, + Consolidation::Scheduled, + Consolidation::Automatic, + ] { + let engine = ReferenceEngine::new().with_consolidation(declared); + run(&engine) + .await + .unwrap_or_else(|error| panic!("declaring {declared:?}: {error}")); + assert!(engine.is_empty(), "{declared:?}"); + } } #[tokio::test] diff --git a/crates/tinymemory-integrations/tests/live_cortex_lifecycle.rs b/crates/tinymemory-integrations/tests/live_cortex_lifecycle.rs index bdbb423e..f93ac9ae 100644 --- a/crates/tinymemory-integrations/tests/live_cortex_lifecycle.rs +++ b/crates/tinymemory-integrations/tests/live_cortex_lifecycle.rs @@ -19,13 +19,16 @@ use std::sync::Arc; use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; use tinymemory_api::{ - ForgetTarget, LearningKind, MemoryEngine, MemoryMeta, MetaFilter, Namespace, Reach, StoreItem, + Consolidation, ForgetTarget, LearningKind, MemoryEngine, MemoryMeta, MetaFilter, Namespace, + Reach, StoreItem, }; use tinymemory_integrations::brain::brain_document; use tinymemory_integrations::cortex::{CortexCredential, CortexEngine}; use tinymemory_integrations::documents::{ConverterChain, RawDocument}; +use tinymemory_integrations::{EngineCredential, EngineSettings, build_engine}; use tinymemory_tools::{ - AgentMemory, Brain, ContextPack, CoreScope, MemoryLayout, PostTurn, PreTurn, + AgentMemory, Brain, BrainDocument, BrainSource, ContextPack, CoreScope, JobOutcome, + MemoryLayout, PostTurn, PreTurn, RecallPolicy, }; const DEFAULT_KEY: &str = "tinymemory-cortex-test"; @@ -236,3 +239,84 @@ async fn live_core_scope_recall_and_promotion_respect_tenant_boundaries() { assert_eq!(forgotten.forgotten, 1, "{forgotten:?}"); } } + +#[tokio::test] +async fn live_an_automatic_engine_queues_no_builds_but_still_builds_on_request() { + let Ok(url) = std::env::var("TINYMEMORY_LIVE_CORTEXDB_URL") else { + eprintln!("TINYMEMORY_LIVE_CORTEXDB_URL unset; skipping"); + return; + }; + let key = std::env::var("TINYMEMORY_TEST_CORTEX_KEY").unwrap_or_else(|_| DEFAULT_KEY.into()); + // A self-hosted server that runs its own layer scheduler, declared so in + // the config, as a host would. + let settings = EngineSettings { + endpoint: Some(url), + consolidation: Some(Consolidation::Automatic), + ..EngineSettings::default() + }; + let engine = build_engine("cortexdb", &settings, EngineCredential::Static(key)) + .expect("a valid live engine"); + assert_eq!(engine.descriptor().consolidation, Consolidation::Automatic); + + let nanos = SystemTime::now() + .duration_since(UNIX_EPOCH) + .expect("clock after the epoch") + .as_nanos(); + let layout = MemoryLayout::new( + format!("project:live-auto-{nanos}") + .parse() + .expect("a valid root"), + ) + .expect("a valid layout"); + + let brain = Brain::new(engine.clone(), layout.clone()); + let ingested = brain + .ingest(BrainDocument::new( + BrainSource::Markdown, + "Expense reports are due on the fifth of each month.", + )) + .await + .expect("ingest"); + assert_eq!( + ingested.job, None, + "an automatic engine gets no ingest build" + ); + + let support = AgentMemory::new(engine.clone(), layout.clone(), "support-01") + .expect("agent") + .with_policy(RecallPolicy { + build_beliefs_every: Some(1), + ..RecallPolicy::default() + }); + let report = support + .post_turn(PostTurn::new("auto-1", 1, "They are due on the fifth.")) + .await + .expect("post_turn"); + assert!(report.jobs.is_empty(), "no turn build: {:?}", report.jobs); + + // An explicit refresh still reaches `v1/beliefs/build` on the real server. + let built = support + .run_background(brain.build(&BrainSource::Markdown).expect("build job")) + .await + .expect("a belief build is accepted"); + assert!( + matches!(built.outcome, JobOutcome::Done | JobOutcome::Started), + "{:?}", + built.outcome + ); + let history = support + .run_background(support.history_build()) + .await + .expect("a history build is accepted"); + assert!( + matches!(history.outcome, JobOutcome::Done | JobOutcome::Started), + "{:?}", + history.outcome + ); + + let forgotten = engine + .forget(ForgetTarget::Filter(layout.holistic_filter())) + .await + .expect("forget"); + assert!(forgotten.forgotten >= 2, "{forgotten:?}"); +} diff --git a/docs/integration.md b/docs/integration.md index 8df617be..b4f7b22b 100644 --- a/docs/integration.md +++ b/docs/integration.md @@ -180,7 +180,7 @@ let brain = Brain::new(engine.clone(), layout.clone()); let ingested = brain .ingest(BrainDocument::new(BrainSource::Notion, text).titled("Refund policy")) .await?; -queue.push(ingested.job); // build this source's beliefs later +queue.extend(ingested.job); // build this source's beliefs later (none on an `Automatic` engine) // A file: converted, its source picked from the format (PDF → pdf, md → markdown). let document = brain_document(&converters, &raw, None, MemoryMeta::default()).await?; diff --git a/docs/specs/agent-memory.md b/docs/specs/agent-memory.md index 9a80a24c..b75f8529 100644 --- a/docs/specs/agent-memory.md +++ b/docs/specs/agent-memory.md @@ -184,7 +184,10 @@ skipped, engine }`. - `ingest` and `ingest_with(WaitFor)` store a `BrainDocument` at its source's node. The source kind defaults from the `BrainSource`. - `ingest_many` batches by `MAX_STORE_MANY`. - - Each ingest returns the `BuildBeliefs` job for its source scope. + - Each ingest returns the `BuildBeliefs` job for its source scope, and + `ingest_many` one per source touched; on an engine that declares + `Automatic` neither returns any (`Ingested::job` is `None`, `jobs` is + empty), and `Brain::build` asks for one explicitly. - **`Brain::search`** fetches within one source or across the whole brain, and **`Brain::forget`** erases one source. - **`BackgroundJob`** is `BuildBeliefs { request }` or From c1f2f310d05c1f2112717dd35b25caba44e95eb0 Mon Sep 17 00:00:00 2001 From: M3gA-Mind Date: Tue, 6 Oct 2026 20:50:13 +0530 Subject: [PATCH 3/4] Prove server-side automatic builds live; document the reference engine - live_cortex_lifecycle: after an ingest on an Automatic engine, with no build requested, poll CortexDB's v1/derivation/status until the scope's last build is no older than its last write: the harness server's layer scheduler rebuilt it on its own. - ReferenceEngine::with_consolidation: say that the in-memory engine has no background builder and builds only when asked, and why it does not build on every write (its beliefs are ordinary items, which every read would then return). --- .../src/conformance/reference/mod.rs | 17 ++++-- .../tests/live_cortex_lifecycle.rs | 53 ++++++++++++++++++- 2 files changed, 64 insertions(+), 6 deletions(-) diff --git a/crates/tinymemory-api/src/conformance/reference/mod.rs b/crates/tinymemory-api/src/conformance/reference/mod.rs index a55c3799..b3867448 100644 --- a/crates/tinymemory-api/src/conformance/reference/mod.rs +++ b/crates/tinymemory-api/src/conformance/reference/mod.rs @@ -59,10 +59,19 @@ impl ReferenceEngine { } } - /// The same engine, declaring `consolidation`, and consolidating as it - /// declares: [`Consolidation::OnDemand`] and [`Consolidation::Automatic`] - /// build at once, [`Consolidation::Scheduled`] only acknowledges, and - /// [`Consolidation::None`] refuses. For a host test of how a lifecycle + /// The same engine, declaring `consolidation`, and answering + /// [`MemoryEngine::consolidate`] as it declares: + /// [`Consolidation::OnDemand`] and [`Consolidation::Automatic`] build at + /// once, [`Consolidation::Scheduled`] only acknowledges, and + /// [`Consolidation::None`] refuses. + /// + /// This in-memory engine has no background builder: declared + /// [`Consolidation::Automatic`], it builds only when asked. That is what + /// a host test of the lifecycle needs (an `Automatic` engine is handed no + /// build jobs, and an explicit build still runs). Building on every write + /// would not model CortexDB either: the reference engine keeps beliefs as + /// ordinary learning items, so they would appear in every read, where + /// CortexDB keeps them in a separate layer. For a host test of how a lifecycle /// treats each kind of engine; the conformance suite holds it to what /// it declares. #[must_use] diff --git a/crates/tinymemory-integrations/tests/live_cortex_lifecycle.rs b/crates/tinymemory-integrations/tests/live_cortex_lifecycle.rs index f93ac9ae..bd32573d 100644 --- a/crates/tinymemory-integrations/tests/live_cortex_lifecycle.rs +++ b/crates/tinymemory-integrations/tests/live_cortex_lifecycle.rs @@ -240,6 +240,47 @@ async fn live_core_scope_recall_and_promotion_respect_tenant_boundaries() { } } +/// Polls CortexDB's `v1/derivation/status` for `scope` until its last build +/// is no older than its last write: the server's own scheduler caught up +/// with the writes, with no build requested. `false` if that never happens +/// within a few minutes. +async fn built_after_last_write(url: &str, key: &str, scope: &str) -> bool { + use tinymemory_api::chrono::{DateTime, FixedOffset}; + let parse = |value: &serde_json::Value| -> Option> { + DateTime::parse_from_rfc3339(value.as_str()?).ok() + }; + let client = reqwest::Client::new(); + let endpoint = format!("{}/v1/derivation/status", url.trim_end_matches('/')); + let deadline = Instant::now() + Duration::from_secs(240); + loop { + let status: Option = match client + .get(&endpoint) + .query(&[("scope", scope)]) + .bearer_auth(key) + .send() + .await + { + Ok(response) => response.json().await.ok(), + Err(_) => None, + }; + if let Some(status) = &status { + let wrote = parse(&status["last_write_at"]); + let built = parse(&status["last_built_at"]); + if let (Some(wrote), Some(built)) = (wrote, built) + && built >= wrote + { + eprintln!("derivation status: {status}"); + return true; + } + } + if Instant::now() >= deadline { + eprintln!("derivation status at the deadline: {status:?}"); + return false; + } + tokio::time::sleep(Duration::from_secs(2)).await; + } +} + #[tokio::test] async fn live_an_automatic_engine_queues_no_builds_but_still_builds_on_request() { let Ok(url) = std::env::var("TINYMEMORY_LIVE_CORTEXDB_URL") else { @@ -250,11 +291,11 @@ async fn live_an_automatic_engine_queues_no_builds_but_still_builds_on_request() // A self-hosted server that runs its own layer scheduler, declared so in // the config, as a host would. let settings = EngineSettings { - endpoint: Some(url), + endpoint: Some(url.clone()), consolidation: Some(Consolidation::Automatic), ..EngineSettings::default() }; - let engine = build_engine("cortexdb", &settings, EngineCredential::Static(key)) + let engine = build_engine("cortexdb", &settings, EngineCredential::Static(key.clone())) .expect("a valid live engine"); assert_eq!(engine.descriptor().consolidation, Consolidation::Automatic); @@ -282,6 +323,14 @@ async fn live_an_automatic_engine_queues_no_builds_but_still_builds_on_request() "an automatic engine gets no ingest build" ); + // Nothing asked for a build, yet the server rebuilds the written scope on + // its own (the harness runs the layer scheduler: `CORTEX_V1_LAYERS_AUTO`). + let scope = format!("app:tinymemory/project:live-auto-{nanos}/source:markdown/app:documents"); + assert!( + built_after_last_write(&url, &key, &scope).await, + "the server never rebuilt `{scope}` after the write on its own" + ); + let support = AgentMemory::new(engine.clone(), layout.clone(), "support-01") .expect("agent") .with_policy(RecallPolicy { From 05d60bd9b1dce40ff186b4dc4907104103f3b66a Mon Sep 17 00:00:00 2001 From: M3gA-Mind Date: Tue, 6 Oct 2026 21:04:07 +0530 Subject: [PATCH 4/4] Live test: require a build strictly after the write; bound each poll The derivation-status poll now needs last_built_at strictly after last_write_at (the scope is unique to the run, so that write is the test's own), and every request is bounded by the time left, so a stalled server cannot hold the poll past its deadline. --- .../tests/live_cortex_lifecycle.rs | 30 +++++++++++-------- 1 file changed, 17 insertions(+), 13 deletions(-) diff --git a/crates/tinymemory-integrations/tests/live_cortex_lifecycle.rs b/crates/tinymemory-integrations/tests/live_cortex_lifecycle.rs index bd32573d..4a7b98ad 100644 --- a/crates/tinymemory-integrations/tests/live_cortex_lifecycle.rs +++ b/crates/tinymemory-integrations/tests/live_cortex_lifecycle.rs @@ -241,9 +241,11 @@ async fn live_core_scope_recall_and_promotion_respect_tenant_boundaries() { } /// Polls CortexDB's `v1/derivation/status` for `scope` until its last build -/// is no older than its last write: the server's own scheduler caught up -/// with the writes, with no build requested. `false` if that never happens -/// within a few minutes. +/// is strictly newer than its last write: the server's own scheduler rebuilt +/// after the write, with no build requested. The scope is unique to the run, +/// so its last write is this test's. Every request is bounded by the time +/// left, so the poll ends within its deadline even if the server stalls. +/// `false` if no such build happens within a few minutes. async fn built_after_last_write(url: &str, key: &str, scope: &str) -> bool { use tinymemory_api::chrono::{DateTime, FixedOffset}; let parse = |value: &serde_json::Value| -> Option> { @@ -253,21 +255,23 @@ async fn built_after_last_write(url: &str, key: &str, scope: &str) -> bool { let endpoint = format!("{}/v1/derivation/status", url.trim_end_matches('/')); let deadline = Instant::now() + Duration::from_secs(240); loop { - let status: Option = match client - .get(&endpoint) - .query(&[("scope", scope)]) - .bearer_auth(key) - .send() - .await - { - Ok(response) => response.json().await.ok(), - Err(_) => None, + let left = deadline.saturating_duration_since(Instant::now()); + let request = async { + let response = client + .get(&endpoint) + .query(&[("scope", scope)]) + .bearer_auth(key) + .send() + .await + .ok()?; + response.json::().await.ok() }; + let status = tokio::time::timeout(left, request).await.ok().flatten(); if let Some(status) = &status { let wrote = parse(&status["last_write_at"]); let built = parse(&status["last_built_at"]); if let (Some(wrote), Some(built)) = (wrote, built) - && built >= wrote + && built > wrote { eprintln!("derivation status: {status}"); return true;