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
35 changes: 34 additions & 1 deletion crates/tinymemory-api/src/conformance/reference/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,27 @@ impl ReferenceEngine {
}
}

/// 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]
pub fn with_consolidation(mut self, consolidation: Consolidation) -> Self {
self.descriptor.consolidation = consolidation;
Comment thread
M3gA-Mind marked this conversation as resolved.
self
}

/// How many items the engine holds.
#[must_use]
pub fn len(&self) -> usize {
Expand Down Expand Up @@ -206,9 +227,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<ConsolidateReceipt> {
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 => {}
Comment thread
M3gA-Mind marked this conversation as resolved.
Comment thread
M3gA-Mind marked this conversation as resolved.
}
let mut items = self.items()?;
let beliefs = distil::distil(&items, &req);
let mut nodes: Vec<&crate::Namespace> = Vec::new();
Expand Down
13 changes: 9 additions & 4 deletions crates/tinymemory-api/src/conformance/suite/lifecycle.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down Expand Up @@ -119,13 +120,13 @@ fn answers_as_promised(
check: CHECK,
source,
}),
(Consolidation::OnDemand, Ok(receipt)) => ensure(
(Consolidation::OnDemand | Consolidation::Automatic, Ok(receipt)) => ensure(
Comment thread
M3gA-Mind marked this conversation as resolved.
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,
Expand Down Expand Up @@ -174,3 +175,7 @@ async fn beliefs(ctx: &Ctx<'_>, node: Namespace) -> Result<()> {
}
Ok(())
}

#[cfg(test)]
#[path = "lifecycle_tests.rs"]
mod tests;
30 changes: 30 additions & 0 deletions crates/tinymemory-api/src/conformance/suite/lifecycle_tests.rs
Original file line number Diff line number Diff line change
@@ -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();
Comment thread
M3gA-Mind marked this conversation as resolved.
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:?}");
}
8 changes: 8 additions & 0 deletions crates/tinymemory-api/src/consolidate/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Comment thread
M3gA-Mind marked this conversation as resolved.
/// [`crate::MemoryEngine::consolidate`] still runs a build at once.
Automatic,
}

/// What to consolidate.
Expand Down
8 changes: 8 additions & 0 deletions crates/tinymemory-api/src/consolidate/mod_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
4 changes: 2 additions & 2 deletions crates/tinymemory-api/src/engine/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
///
Expand Down
16 changes: 16 additions & 0 deletions crates/tinymemory-api/tests/conformance_reference.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,22 @@ async fn the_reference_engine_passes_and_cleans_up() {
assert!(engine.is_empty());
}

#[tokio::test]
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]
async fn the_suite_leaves_foreign_items_alone() {
let engine = ReferenceEngine::new();
Expand Down
2 changes: 1 addition & 1 deletion crates/tinymemory-integrations/examples/cortex_agent.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
10 changes: 9 additions & 1 deletion crates/tinymemory-integrations/src/config/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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};

Expand Down Expand Up @@ -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<String, String>,
/// 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")]
Comment thread
M3gA-Mind marked this conversation as resolved.
pub consolidation: Option<Consolidation>,
}

#[cfg(test)]
Expand Down
19 changes: 19 additions & 0 deletions crates/tinymemory-integrations/src/config/mod_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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() {
Comment thread
M3gA-Mind marked this conversation as resolved.
let config: MemoryConfig = toml::from_str(
r#"
engine = "cortexdb"

[engines.cortexdb]
endpoint = "http://127.0.0.1:3141"
consolidation = "automatic"
"#,
)
.unwrap();
assert_eq!(
Comment thread
M3gA-Mind marked this conversation as resolved.
config.settings().consolidation,
Some(Consolidation::Automatic)
);
assert_eq!(EngineSettings::default().consolidation, None);
}
8 changes: 6 additions & 2 deletions crates/tinymemory-integrations/src/cortex/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
47 changes: 41 additions & 6 deletions crates/tinymemory-integrations/src/cortex/descriptor/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand All @@ -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 {
Expand All @@ -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),
Comment thread
M3gA-Mind marked this conversation as resolved.
}
}

Expand All @@ -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 {
Expand Down
29 changes: 29 additions & 0 deletions crates/tinymemory-integrations/src/cortex/descriptor/mod_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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() {
Comment thread
M3gA-Mind marked this conversation as resolved.
Comment thread
M3gA-Mind marked this conversation as resolved.
assert_eq!(
Comment thread
M3gA-Mind marked this conversation as resolved.
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}"
);
}
}
Loading
Loading