#![deny(unsafe_code)]
#![warn(missing_docs, rust_2018_idioms)]
pub mod fire;
pub mod mcr2;
pub mod trigger;
use std::sync::Arc;
use std::time::Duration;
use dashmap::DashMap;
use smol_str::SmolStr;
use tokio::sync::mpsc;
use tracing::{info, instrument, warn};
use exocortex_kernel::{Memory, MemoryId, Provenance, Relationship, RelationshipId};
use exocortex_storage::{FencedBatchCommit, FencedRestore, LeaseKey, RegionKey, Storage};
use mcr2::{
compute_sparsity, effective_strength, GraphSparsity, MCR2Engine, MCR2Value, MemoryWithEmbedding,
};
use trigger::{DreamsTrigger, RegionWriteCounters};
pub const SIMILAR_TO_THRESHOLD: f32 = 0.85;
pub const MAX_DISCOVERIES_PER_CYCLE: usize = 16;
pub const MAX_DISCOVERY_PATH_INSPECTIONS: usize = 50_000;
const MAX_REGION_MEMORIES: usize = 50_000;
const MAX_REGION_RELATIONSHIPS: usize = 50_000;
type DiscoveryEdge = (MemoryId, MemoryId, u32, bool);
struct RegionWorkingSet {
memories: std::collections::HashMap<MemoryId, Memory>,
relationships: std::collections::HashMap<RelationshipId, Relationship>,
}
impl RegionWorkingSet {
fn apply_relationship_writes(&mut self, writes: &[Relationship]) {
let ontology = dreams_ontology();
for relationship in writes {
self.relationships
.insert(relationship.id, relationship.clone());
if let Some(inverse) = exocortex_kernel::materialize_inverse(ontology, relationship) {
self.relationships.insert(inverse.id, inverse);
}
}
}
}
#[derive(Clone, Debug)]
pub struct ConsolidationResult {
pub session_id: SmolStr,
pub user_id: Option<SmolStr>,
pub started_at: chrono::DateTime<chrono::Utc>,
pub completed_at: chrono::DateTime<chrono::Utc>,
pub region: RegionKey,
pub memories_input: u32,
pub memories_output: u32,
pub mcr2_before: MCR2Value,
pub mcr2_after: MCR2Value,
pub sparsity_before: GraphSparsity,
pub sparsity_after: GraphSparsity,
pub merged: Vec<MemoryId>,
pub abstracted: Vec<MemoryId>,
pub pruned: Vec<(MemoryId, PruneReason)>,
pub strengthened: Vec<RelationshipId>,
pub rewired: Vec<RelationshipId>,
pub similar_edges: Vec<RelationshipId>,
pub owner_node_id: SmolStr,
pub lease_epoch: u64,
pub regression: bool,
pub hairball_regression: bool,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum PruneReason {
Redundant,
Superseded,
Stale,
LowValue,
}
struct CycleJournal {
cycle_id: SmolStr,
memories: std::collections::BTreeMap<MemoryId, Memory>,
relationships: std::collections::BTreeMap<RelationshipId, Relationship>,
created_memories: std::collections::BTreeMap<MemoryId, Memory>,
created_relationships: std::collections::BTreeMap<RelationshipId, Relationship>,
owned_memory_lsns: std::collections::BTreeMap<MemoryId, std::collections::BTreeSet<u64>>,
owned_relationship_lsns:
std::collections::BTreeMap<RelationshipId, std::collections::BTreeSet<u64>>,
}
impl CycleJournal {
fn new(cycle_id: impl Into<SmolStr>) -> Self {
Self {
cycle_id: cycle_id.into(),
memories: Default::default(),
relationships: Default::default(),
created_memories: Default::default(),
created_relationships: Default::default(),
owned_memory_lsns: Default::default(),
owned_relationship_lsns: Default::default(),
}
}
fn record_memory(&mut self, memory: &Memory) {
if !self.created_memories.contains_key(&memory.id) {
self.memories
.entry(memory.id)
.or_insert_with(|| memory.clone());
}
}
fn record_relationship(&mut self, relationship: &Relationship) {
if !self.created_relationships.contains_key(&relationship.id) {
self.relationships
.entry(relationship.id)
.or_insert_with(|| relationship.clone());
}
}
fn create_relationship(&mut self, relationship: &Relationship) {
if !self.relationships.contains_key(&relationship.id) {
self.created_relationships
.entry(relationship.id)
.or_insert_with(|| relationship.clone());
}
}
fn prepare_relationship_writes(
&mut self,
writes: &[Relationship],
current: &std::collections::HashMap<RelationshipId, Relationship>,
) {
let ontology = dreams_ontology();
let mut seen: std::collections::HashSet<RelationshipId> =
writes.iter().map(|relationship| relationship.id).collect();
for relationship in writes {
self.prepare_relationship(relationship, current);
if let Some(inverse) = exocortex_kernel::materialize_inverse(ontology, relationship) {
if seen.insert(inverse.id) {
self.prepare_relationship(&inverse, current);
}
}
}
}
fn prepare_relationship(
&mut self,
relationship: &Relationship,
current: &std::collections::HashMap<RelationshipId, Relationship>,
) {
match current.get(&relationship.id) {
Some(preimage) => self.record_relationship(preimage),
None => self.create_relationship(relationship),
}
}
fn record_commit(&mut self, commit: &FencedBatchCommit) {
for (id, lsns) in &commit.memory_lsns {
self.owned_memory_lsns.entry(*id).or_default().extend(lsns);
}
for (id, lsns) in &commit.relationship_lsns {
self.owned_relationship_lsns
.entry(*id)
.or_default()
.extend(lsns);
}
}
fn restore(&self) -> FencedRestore {
FencedRestore {
memories: self.memories.values().cloned().collect(),
relationships: self.relationships.values().cloned().collect(),
created_memories: self.created_memories.values().cloned().collect(),
created_relationships: self.created_relationships.values().cloned().collect(),
owned_memory_lsns: self.owned_memory_lsns.clone(),
owned_relationship_lsns: self.owned_relationship_lsns.clone(),
}
}
fn is_empty(&self) -> bool {
self.owned_memory_lsns.is_empty() && self.owned_relationship_lsns.is_empty()
}
}
fn dreams_ontology() -> &'static exocortex_kernel::Ontology {
static ONTOLOGY: std::sync::OnceLock<exocortex_kernel::Ontology> = std::sync::OnceLock::new();
ONTOLOGY.get_or_init(|| {
exocortex_kernel::Ontology::from_packs(vec![exocortex_pack_dev_v1::pack_def()])
.expect("the compiled development ontology must be valid")
})
}
pub struct DreamsEngine<S: Storage> {
pub storage: Arc<S>,
pub leader_gate: Option<std::sync::Arc<std::sync::atomic::AtomicBool>>,
pub counters: DashMap<RegionKey, RegionWriteCounters>,
pending_regions: DashMap<RegionKey, ()>,
distributed_fire: Option<Arc<tokio::sync::Mutex<fire::RedisFireQueue>>>,
last_local_write_event: DashMap<RegionKey, SmolStr>,
distributed_notifications: DashMap<RegionKey, fire::FireMessage>,
pub last_cycle_at: DashMap<RegionKey, chrono::DateTime<chrono::Utc>>,
pub dreams_trigger: DreamsTrigger,
pub tolerance: f32,
pub hairball_tolerance: f32,
pub rollback_on_regression: bool,
pub tx_fire: mpsc::Sender<(RegionKey, RegionWriteCounters)>,
pub rx_fire: tokio::sync::Mutex<mpsc::Receiver<(RegionKey, RegionWriteCounters)>>,
pub node_id: SmolStr,
pub discoveries: DashMap<uuid::Uuid, Discovery>,
pub last_result: tokio::sync::RwLock<Option<ConsolidationResult>>,
lease_ttl: Duration,
#[cfg(feature = "testing")]
cycle_fault_after: Option<usize>,
#[cfg(feature = "testing")]
cycle_crash_after: Option<usize>,
#[cfg(feature = "testing")]
cycle_pause_after: Option<(usize, Duration)>,
#[cfg(feature = "testing")]
rollback_pause: Option<Duration>,
#[cfg(feature = "testing")]
rollback_concurrent_memories: Vec<Memory>,
#[cfg(feature = "testing")]
renewal_failure_after: Option<usize>,
}
impl<S: Storage + 'static> DreamsEngine<S> {
pub fn with_leader_gate(mut self, gate: std::sync::Arc<std::sync::atomic::AtomicBool>) -> Self {
self.leader_gate = Some(gate);
self
}
pub fn new(
storage: Arc<S>,
dreams_trigger: DreamsTrigger,
tolerance: f32,
hairball_tolerance: f32,
rollback_on_regression: bool,
node_id: SmolStr,
) -> Self {
let (tx_fire, rx_fire) = mpsc::channel(1000);
Self {
leader_gate: None,
storage,
counters: DashMap::new(),
pending_regions: DashMap::new(),
distributed_fire: None,
last_local_write_event: DashMap::new(),
distributed_notifications: DashMap::new(),
last_cycle_at: DashMap::new(),
dreams_trigger,
tolerance,
hairball_tolerance,
rollback_on_regression,
tx_fire,
rx_fire: tokio::sync::Mutex::new(rx_fire),
node_id,
discoveries: DashMap::new(),
last_result: tokio::sync::RwLock::new(None),
lease_ttl: Duration::from_secs(60),
#[cfg(feature = "testing")]
cycle_fault_after: None,
#[cfg(feature = "testing")]
cycle_crash_after: None,
#[cfg(feature = "testing")]
cycle_pause_after: None,
#[cfg(feature = "testing")]
rollback_pause: None,
#[cfg(feature = "testing")]
rollback_concurrent_memories: Vec::new(),
#[cfg(feature = "testing")]
renewal_failure_after: None,
}
}
pub fn with_distributed_fire(
mut self,
queue: Arc<tokio::sync::Mutex<fire::RedisFireQueue>>,
) -> Self {
self.distributed_fire = Some(queue);
self
}
#[doc(hidden)]
#[cfg(feature = "testing")]
pub fn with_cycle_fault_after(mut self, mutation: usize) -> Self {
self.cycle_fault_after = Some(mutation);
self
}
#[doc(hidden)]
#[cfg(feature = "testing")]
pub fn with_cycle_crash_after(mut self, mutation: usize) -> Self {
self.cycle_crash_after = Some(mutation);
self
}
#[doc(hidden)]
#[cfg(feature = "testing")]
pub async fn recover_active_cycle_for_test(&self, region: &RegionKey) -> anyhow::Result<()> {
let lease_key = LeaseKey::Dreams {
org: region.org.clone(),
region: format!("{}:{}", region.project, region.memory_type).into(),
};
let lease = self
.storage
.acquire_lease(&lease_key, self.lease_ttl)
.await
.map_err(|error| anyhow::anyhow!("lease: {error}"))?;
let recovery = self.recover_active_cycle(&lease_key, &lease).await;
let release = self.storage.release_lease(lease).await;
recovery?;
release.map_err(|error| anyhow::anyhow!("release recovery lease: {error}"))
}
#[doc(hidden)]
#[cfg(feature = "testing")]
pub fn with_lease_ttl(mut self, ttl: Duration) -> Self {
self.lease_ttl = ttl;
self
}
#[doc(hidden)]
#[cfg(feature = "testing")]
pub fn with_cycle_pause_after(mut self, mutation: usize, pause: Duration) -> Self {
self.cycle_pause_after = Some((mutation, pause));
self
}
#[doc(hidden)]
#[cfg(feature = "testing")]
pub fn with_rollback_pause(mut self, pause: Duration) -> Self {
self.rollback_pause = Some(pause);
self
}
#[doc(hidden)]
#[cfg(feature = "testing")]
pub fn with_rollback_concurrent_memories(mut self, memories: Vec<Memory>) -> Self {
self.rollback_concurrent_memories = memories;
self
}
#[doc(hidden)]
#[cfg(feature = "testing")]
pub fn with_renewal_failure_after(mut self, attempt: usize) -> Self {
self.renewal_failure_after = Some(attempt);
self
}
pub async fn on_write(&self, region: RegionKey) -> anyhow::Result<()> {
self.on_writes(region, 1, 0).await
}
pub async fn on_writes(
&self,
region: RegionKey,
memories: u32,
edges: u32,
) -> anyhow::Result<()> {
if let Some(queue) = &self.distributed_fire {
queue
.lock()
.await
.record_write(
®ion,
memories,
edges,
self.dreams_trigger,
self.node_id.as_str(),
)
.await?;
return Ok(());
}
let now = chrono::Utc::now();
let anchor = *self
.last_cycle_at
.entry(region.clone())
.or_insert(now)
.value();
let mut e = self.counters.entry(region.clone()).or_default();
e.memories_since_last_cycle = e.memories_since_last_cycle.saturating_add(memories);
e.edges_since_last_cycle = e.edges_since_last_cycle.saturating_add(edges);
e.seconds_since_last_cycle = (now - anchor).num_seconds().max(0) as u64;
if self.is_leader() && self.dreams_trigger.should_fire(&e) {
let snap = *e;
drop(e);
self.schedule_region(region, snap);
}
Ok(())
}
pub async fn on_writes_once(
&self,
event_id: &str,
delivery_generation: u64,
region: RegionKey,
memories: u32,
edges: u32,
) -> anyhow::Result<()> {
anyhow::ensure!(
!event_id.is_empty(),
"Dreams write event id must not be empty"
);
if let Some(queue) = &self.distributed_fire {
queue
.lock()
.await
.record_write_once(
®ion,
memories,
edges,
self.dreams_trigger,
self.node_id.as_str(),
fire::DeliveryFence {
event_id,
generation: delivery_generation,
},
)
.await?;
return Ok(());
}
if self
.last_local_write_event
.get(®ion)
.is_some_and(|seen| seen.as_str() == event_id)
{
return Ok(());
}
self.on_writes(region.clone(), memories, edges).await?;
self.last_local_write_event.insert(region, event_id.into());
Ok(())
}
pub async fn settle_writes_once(
&self,
event_id: &str,
delivery_generation: u64,
retain_legacy_identity: bool,
regions: Vec<RegionKey>,
) -> anyhow::Result<()> {
anyhow::ensure!(
!event_id.is_empty(),
"Dreams write event id must not be empty"
);
if let Some(queue) = &self.distributed_fire {
let mut queue = queue.lock().await;
for region in regions {
queue
.forget_write_once(
®ion,
event_id,
delivery_generation,
retain_legacy_identity,
)
.await?;
}
return Ok(());
}
for region in regions {
if self
.last_local_write_event
.get(®ion)
.is_some_and(|seen| seen.as_str() == event_id)
{
self.last_local_write_event.remove(®ion);
}
}
Ok(())
}
pub fn notify(&self, region: RegionKey) {
if !self.is_leader() {
return;
}
let snap = self.counters.get(®ion).map(|e| *e).unwrap_or_default();
self.schedule_region(region, snap);
}
pub fn notify_distributed(&self, notification: fire::FireMessage) {
if !self.is_leader() {
return;
}
let fired_at = notification.fired_at.unwrap_or_default();
let region = notification.region.clone();
self.distributed_notifications
.insert(region.clone(), notification);
self.schedule_region(region, fired_at);
}
fn is_leader(&self) -> bool {
self.leader_gate
.as_ref()
.is_none_or(|gate| gate.load(std::sync::atomic::Ordering::SeqCst))
}
fn schedule_region(&self, region: RegionKey, snapshot: RegionWriteCounters) {
use dashmap::mapref::entry::Entry;
match self.pending_regions.entry(region.clone()) {
Entry::Occupied(_) => {}
Entry::Vacant(entry) => {
entry.insert(());
if self.tx_fire.try_send((region.clone(), snapshot)).is_err() {
self.pending_regions.remove(®ion);
metrics::counter!("exocortex_dreams_queue_dropped_total").increment(1);
}
}
}
}
fn complete_region(&self, region: &RegionKey, fired_at: RegionWriteCounters, success: bool) {
if success {
if let Some(mut current) = self.counters.get_mut(region) {
current.memories_since_last_cycle = current
.memories_since_last_cycle
.saturating_sub(fired_at.memories_since_last_cycle);
current.edges_since_last_cycle = current
.edges_since_last_cycle
.saturating_sub(fired_at.edges_since_last_cycle);
current.seconds_since_last_cycle = 0;
}
self.last_cycle_at
.insert(region.clone(), chrono::Utc::now());
}
self.pending_regions.remove(region);
if success && self.is_leader() {
if let Some(current) = self.counters.get(region).map(|entry| *entry) {
if self.dreams_trigger.should_fire(¤t) {
self.schedule_region(region.clone(), current);
}
}
}
}
pub async fn run(self: Arc<Self>) {
while let Some((region, fired_at)) = { self.rx_fire.lock().await.recv().await } {
let fire_cycle_id = self
.distributed_notifications
.get(®ion)
.and_then(|notification| notification.fire_id.clone())
.map(|fire_id| format!("dream-fire:{fire_id}"));
let success = match self
.try_consolidate_inner(®ion, fire_cycle_id.as_deref())
.await
{
Ok(Some(res)) => {
*self.last_result.write().await = Some(res.clone());
info!(?res, "Dreams cycle ok");
true
}
Ok(None) => {
info!(?region, "Dreams fire already settled; acknowledging replay");
true
}
Err(e) => {
warn!(?e, "Dreams cycle failed");
false
}
};
self.complete_region(®ion, fired_at, success);
if let Some((_, notification)) = self.distributed_notifications.remove(®ion) {
if let Some(queue) = &self.distributed_fire {
let mut delay = Duration::from_millis(50);
loop {
match queue
.lock()
.await
.acknowledge(
¬ification,
success,
self.dreams_trigger,
self.node_id.as_str(),
)
.await
{
Ok(_) => break,
Err(error) => {
warn!(?error, "Dreams distributed acknowledgement retrying");
tokio::time::sleep(delay).await;
delay = (delay * 2).min(Duration::from_secs(5));
}
}
}
}
}
}
}
#[instrument(skip(self))]
pub async fn try_consolidate(&self, region: &RegionKey) -> anyhow::Result<ConsolidationResult> {
self.try_consolidate_inner(region, None)
.await?
.ok_or_else(|| anyhow::anyhow!("fresh Dreams cycle unexpectedly already settled"))
}
#[doc(hidden)]
#[cfg(feature = "testing")]
pub async fn try_consolidate_once_for_testing(
&self,
region: &RegionKey,
fire_id: &str,
) -> anyhow::Result<Option<ConsolidationResult>> {
self.try_consolidate_inner(region, Some(&format!("dream-fire:{fire_id}")))
.await
}
async fn try_consolidate_inner(
&self,
region: &RegionKey,
stable_cycle_id: Option<&str>,
) -> anyhow::Result<Option<ConsolidationResult>> {
if let Some(gate) = &self.leader_gate {
if !gate.load(std::sync::atomic::Ordering::SeqCst) {
return Err(anyhow::anyhow!("not the elected leader"));
}
}
let lease_key = LeaseKey::Dreams {
org: region.org.clone(),
region: format!("{}:{}", region.project, region.memory_type).into(),
};
let lease = self
.storage
.acquire_lease(&lease_key, self.lease_ttl)
.await
.map_err(|e| anyhow::anyhow!("lease: {e}"))?;
let cycle_id = stable_cycle_id
.map(SmolStr::new)
.unwrap_or_else(|| format!("dream:{}", uuid::Uuid::new_v4()).into());
if self
.storage
.cycle_succeeded_fenced(cycle_id.as_str(), &lease)
.await
.map_err(|error| anyhow::anyhow!("load Dreams cycle settlement: {error}"))?
{
self.storage
.release_lease(lease)
.await
.map_err(|error| anyhow::anyhow!("release settled Dreams lease: {error}"))?;
return Ok(None);
}
if let Err(error) = self.recover_active_cycle(&lease_key, &lease).await {
let _ = self.storage.release_lease(lease).await;
return Err(error);
}
let working_set = match self.load_region_working_set(region).await {
Ok(working_set) => working_set,
Err(error) => {
let _ = self.storage.release_lease(lease).await;
return Err(error);
}
};
if let Err(error) = self.validate_loaded_region(region, &working_set) {
let _ = self.storage.release_lease(lease).await;
return Err(error);
}
let mut renewal = self.spawn_lease_renewal(lease.clone());
let mut journal = CycleJournal::new(cycle_id);
let mut renewal_stopped = false;
let outcome = {
let consolidation =
self.consolidate_under_tracked(&lease, region, &mut journal, working_set);
tokio::pin!(consolidation);
tokio::select! {
biased;
renewal_outcome = &mut renewal => {
renewal_stopped = true;
Err(Self::renewal_task_error(renewal_outcome))
}
outcome = &mut consolidation => outcome,
}
};
#[cfg(feature = "testing")]
let crashed = self.cycle_crash_after.is_some() && outcome.is_err();
#[cfg(not(feature = "testing"))]
let crashed = false;
let outcome = if crashed {
outcome
} else if renewal_stopped {
self.finish_cycle(outcome, &journal, &lease).await
} else {
let finishing = self.finish_cycle(outcome, &journal, &lease);
tokio::pin!(finishing);
tokio::select! {
renewal_outcome = &mut renewal => {
renewal_stopped = true;
let renewal_error = Self::renewal_task_error(renewal_outcome);
match finishing.await {
Ok(_) => Err(renewal_error),
Err(cycle_error) => Err(anyhow::anyhow!(
"cycle failed ({cycle_error}); lease renewal also failed ({renewal_error})"
)),
}
}
outcome = &mut finishing => outcome,
}
};
let outcome = match outcome {
Ok(result) if !renewal_stopped && !crashed => {
let discovery =
self.run_discovery_fenced(region, &lease, Some(journal.cycle_id.as_str()));
tokio::pin!(discovery);
tokio::select! {
biased;
renewal_outcome = &mut renewal => {
renewal_stopped = true;
Err(Self::renewal_task_error(renewal_outcome))
}
outcome = &mut discovery => outcome.map(|proposals| {
info!(count = proposals.len(), "discovery ok");
result
}),
}
}
other => other,
};
if !renewal_stopped {
renewal.abort();
let _ = renewal.await;
}
let release = self.storage.release_lease(lease).await;
match (outcome, release) {
(Ok(result), Ok(())) => Ok(Some(result)),
(Err(error), Ok(())) => Err(error),
(Ok(_), Err(release)) => Err(anyhow::anyhow!("release Dreams lease: {release}")),
(Err(error), Err(release)) => Err(anyhow::anyhow!(
"cycle failed ({error}); lease release also failed ({release})"
)),
}
}
async fn recover_active_cycle(
&self,
lease_key: &LeaseKey,
lease: &exocortex_storage::OwnerLease,
) -> anyhow::Result<()> {
let Some(active) = self
.storage
.get_active_cycle_journal(lease_key)
.await
.map_err(|error| anyhow::anyhow!("load active Dreams journal: {error}"))?
else {
return Ok(());
};
self.storage
.restore_fenced(&active.restore, lease)
.await
.map_err(|error| anyhow::anyhow!("recover active Dreams cycle: {error}"))?;
self.storage
.complete_cycle_journal_fenced(active.cycle_id.as_str(), lease)
.await
.map_err(|error| anyhow::anyhow!("complete recovered Dreams journal: {error}"))?;
Ok(())
}
fn renewal_task_error(
outcome: Result<anyhow::Result<()>, tokio::task::JoinError>,
) -> anyhow::Error {
match outcome {
Ok(Err(error)) => error,
Ok(Ok(())) => anyhow::anyhow!("Dreams lease renewal stopped unexpectedly"),
Err(error) => anyhow::anyhow!("Dreams lease renewal task failed: {error}"),
}
}
fn spawn_lease_renewal(
&self,
lease: exocortex_storage::OwnerLease,
) -> tokio::task::JoinHandle<anyhow::Result<()>> {
let storage = self.storage.clone();
let ttl = (lease.expires_at - lease.acquired_at)
.to_std()
.unwrap_or(self.lease_ttl);
let interval = (ttl / 3).max(Duration::from_millis(1));
let retry = (ttl / 12).max(Duration::from_millis(1));
let reserve = interval;
#[cfg(feature = "testing")]
let renewal_failure_after = self.renewal_failure_after;
tokio::spawn(async move {
let mut current = lease;
let mut delay = interval;
#[cfg(feature = "testing")]
let mut attempt = 0usize;
loop {
tokio::time::sleep(delay).await;
#[cfg(feature = "testing")]
{
attempt += 1;
}
#[cfg(feature = "testing")]
let renewed = if renewal_failure_after.is_some_and(|start| attempt >= start) {
Err(exocortex_storage::StorageError::Backend(
"injected lease renewal failure".into(),
))
} else {
storage.renew_lease(¤t).await
};
#[cfg(not(feature = "testing"))]
let renewed = storage.renew_lease(¤t).await;
match renewed {
Ok(renewed) => {
current = renewed;
delay = interval;
}
Err(error) => {
let remaining = (current.expires_at - chrono::Utc::now())
.to_std()
.unwrap_or_default();
if remaining <= reserve {
anyhow::bail!(
"Dreams owner lease renewal could not be confirmed with rollback reserve: {error}"
);
}
warn!(?error, ?remaining, "Dreams owner lease renewal retrying");
delay = retry.min(remaining.saturating_sub(reserve));
}
}
}
})
}
async fn finish_cycle(
&self,
outcome: anyhow::Result<ConsolidationResult>,
journal: &CycleJournal,
lease: &exocortex_storage::OwnerLease,
) -> anyhow::Result<ConsolidationResult> {
match outcome {
Ok(result) => Ok(result),
Err(error) => {
if !journal.is_empty() {
self.rollback(journal, lease).await.map_err(|rollback| {
anyhow::anyhow!(
"cycle failed ({error}); atomic restore failed ({rollback})"
)
})?;
self.storage
.complete_cycle_journal_fenced(journal.cycle_id.as_str(), lease)
.await
.map_err(|complete| {
anyhow::anyhow!(
"cycle failed ({error}); rollback succeeded but journal completion failed ({complete})"
)
})?;
}
Err(error)
}
}
}
async fn consolidate_under_tracked(
&self,
lease: &exocortex_storage::OwnerLease,
region: &RegionKey,
journal: &mut CycleJournal,
mut working_set: RegionWorkingSet,
) -> anyhow::Result<ConsolidationResult> {
let mut mutations = 0usize;
let anchors = self.select_anchors(&working_set);
let mcr2_before = match self.score_with(&anchors) {
Ok(v) => v,
Err(e)
if matches!(
e.downcast_ref::<crate::mcr2::MCR2Error>(),
Some(crate::mcr2::MCR2Error::TooFew(_))
) =>
{
return self
.empty_result(region, lease, anchors.len() as u32, &working_set)
.await;
}
Err(e) => return Err(e),
};
let sparsity_before = self.sparsity(&working_set);
let mut res = ConsolidationResult {
session_id: journal.cycle_id.clone(),
user_id: None,
started_at: chrono::Utc::now(),
completed_at: chrono::Utc::now(),
region: region.clone(),
memories_input: anchors.len() as u32,
memories_output: anchors.len() as u32,
mcr2_before: mcr2_before.clone(),
mcr2_after: mcr2_before,
sparsity_before: sparsity_before.clone(),
sparsity_after: sparsity_before,
merged: vec![],
abstracted: vec![],
pruned: vec![],
strengthened: vec![],
rewired: vec![],
similar_edges: vec![],
owner_node_id: self.node_id.clone(),
lease_epoch: lease.epoch,
regression: false,
hairball_regression: false,
};
let engine = MCR2Engine::default();
let candidates = engine.identify_merge_candidates(&anchors);
for c in candidates.iter().take(32) {
if c.cosine_similarity >= 0.92 {
let merged_before = res.merged.len();
self.merge(&mut res, c, lease, &mut working_set, journal)
.await?;
if res.merged.len() > merged_before {
self.mutation_checkpoint(&mut mutations).await?;
}
}
}
{
let mut by_class: std::collections::BTreeMap<u32, Vec<MemoryId>> = Default::default();
for a in anchors.iter().filter(|a| !res.merged.contains(&a.id)) {
by_class.entry(a.class as u32).or_default().push(a.id);
}
for (_class, members) in by_class {
if members.len() >= 3 {
res.abstracted.push(*members.last().expect("non-empty"));
}
}
}
self.strengthen(&mut res, lease, &mut working_set, journal)
.await?;
if !res.strengthened.is_empty() {
self.mutation_checkpoint(&mut mutations).await?;
}
self.prune(&mut res, &working_set);
res.memories_output = (res.memories_input as usize - res.merged.len()).max(0) as u32;
let survivors: Vec<MemoryWithEmbedding> = anchors
.iter()
.filter(|a| !res.merged.contains(&a.id))
.cloned()
.collect();
self.write_similar_edges(&mut res, &survivors, lease, &mut working_set, journal)
.await?;
if !res.similar_edges.is_empty() {
self.mutation_checkpoint(&mut mutations).await?;
}
let remaining: Vec<MemoryWithEmbedding> = anchors
.into_iter()
.filter(|a| !res.merged.contains(&a.id))
.collect();
res.mcr2_after = match self.score_with(&remaining) {
Ok(score) => score,
Err(e)
if matches!(
e.downcast_ref::<crate::mcr2::MCR2Error>(),
Some(crate::mcr2::MCR2Error::TooFew(_))
) =>
{
res.mcr2_before.clone()
}
Err(e) => return Err(e),
};
res.sparsity_after = self.sparsity(&working_set);
res.completed_at = chrono::Utc::now();
if res.mcr2_after.delta_r < res.mcr2_before.delta_r - self.tolerance {
res.regression = true; if self.rollback_on_regression {
warn!(
"MCR2 degraded {} -> {} - rolling back",
res.mcr2_before.delta_r, res.mcr2_after.delta_r
);
self.rollback(journal, lease).await?;
}
}
if res.sparsity_after.hairball_fraction
> res.sparsity_before.hairball_fraction + self.hairball_tolerance
{
res.hairball_regression = true; }
self.write_audit(&res).await?;
Ok(res)
}
async fn mutation_checkpoint(&self, mutations: &mut usize) -> anyhow::Result<()> {
*mutations += 1;
#[cfg(feature = "testing")]
if self
.cycle_pause_after
.is_some_and(|(mutation, _)| mutation == *mutations)
{
let pause = self
.cycle_pause_after
.expect("pause configuration checked")
.1;
tokio::time::sleep(pause).await;
}
#[cfg(feature = "testing")]
if self.cycle_fault_after == Some(*mutations) || self.cycle_crash_after == Some(*mutations)
{
anyhow::bail!("injected cycle failure after mutation {mutations}");
}
Ok(())
}
async fn load_region_working_set(
&self,
region: &RegionKey,
) -> anyhow::Result<RegionWorkingSet> {
let memories = self
.storage
.memories_in_region(region, MAX_REGION_MEMORIES as u32)
.await?
.into_iter()
.map(|memory| (memory.id, memory))
.collect();
let relationships = self
.storage
.current_relationships_in_region(region, MAX_REGION_RELATIONSHIPS as u32)
.await?
.into_iter()
.map(|relationship| (relationship.id, relationship))
.collect();
Ok(RegionWorkingSet {
memories,
relationships,
})
}
fn select_anchors(&self, working_set: &RegionWorkingSet) -> Vec<MemoryWithEmbedding> {
let mut rows: Vec<(chrono::DateTime<chrono::Utc>, MemoryWithEmbedding)> = Vec::new();
for memory in working_set.memories.values() {
if memory.valid_until.is_some() {
continue;
}
if let Some(embedding) = &memory.embedding {
rows.push((
memory.recorded_at,
MemoryWithEmbedding {
id: memory.id,
class: memory.memory_type,
visibility: memory.visibility,
embedding: embedding.clone(),
},
));
}
}
let now = chrono::Utc::now();
rows.sort_by(|a, b| {
let score = |t: chrono::DateTime<chrono::Utc>| -(now - t).num_days();
score(b.0).cmp(&score(a.0)).then(b.1.id.cmp(&a.1.id))
});
rows.into_iter()
.take(32)
.map(|(_, memory)| memory)
.collect()
}
fn score_with(&self, anchors: &[MemoryWithEmbedding]) -> anyhow::Result<MCR2Value> {
Ok(MCR2Engine::default().compute(anchors)?)
}
fn sparsity(&self, working_set: &RegionWorkingSet) -> GraphSparsity {
let nodes: Vec<_> = working_set
.memories
.values()
.filter(|memory| memory.valid_until.is_none())
.map(|memory| (memory.id, memory.memory_type))
.collect();
let members: std::collections::HashSet<exocortex_kernel::MemoryId> =
nodes.iter().map(|(id, _)| *id).collect();
let edges: Vec<_> = working_set
.relationships
.values()
.filter(|relationship| {
relationship.valid_until.is_none()
&& members.contains(&relationship.from)
&& members.contains(&relationship.to)
})
.map(|relationship| {
(
relationship.from,
relationship.to,
relationship.kind.0,
0u64,
relationship.properties.confidence,
)
})
.collect();
compute_sparsity(&nodes, &edges, 32, similar_to_kind())
}
async fn merge(
&self,
res: &mut ConsolidationResult,
c: &mcr2::MergeCandidate,
lease: &exocortex_storage::OwnerLease,
working_set: &mut RegionWorkingSet,
journal: &mut CycleJournal,
) -> anyhow::Result<()> {
let survivor_live = working_set
.memories
.get(&c.a)
.is_some_and(|memory| memory.valid_until.is_none());
let newer = working_set
.memories
.get(&c.b)
.filter(|memory| memory.valid_until.is_none())
.cloned();
if survivor_live {
if let Some(mut m) = newer {
let now = chrono::Utc::now();
journal.record_memory(&m);
let mut relationship_updates = std::collections::BTreeMap::new();
for relationship in working_set.relationships.values().filter(|relationship| {
relationship.valid_until.is_none()
&& (relationship.from == c.b || relationship.to == c.b)
}) {
journal.record_relationship(relationship);
let mut closed = relationship.clone();
closed.valid_until = Some(now);
relationship_updates.insert(closed.id, closed);
let mut rewired = relationship.clone();
if rewired.from == c.b {
rewired.from = c.a;
}
if rewired.to == c.b {
rewired.to = c.a;
}
if rewired.from == rewired.to {
continue;
}
rewired.id =
RelationshipId::derive(rewired.from, rewired.kind, rewired.to, None);
rewired.recorded_at = now;
rewired.valid_from = now;
rewired.valid_until = None;
rewired.invalidated_by = None;
if let Some(existing) = working_set.relationships.get(&rewired.id) {
if existing.valid_until.is_none() {
continue;
}
journal.record_relationship(existing);
} else {
journal.create_relationship(&rewired);
}
res.rewired.push(rewired.id);
relationship_updates.insert(rewired.id, rewired);
}
m.valid_until = Some(now);
m.invalidated_by = Some(c.a);
let relationship_updates: Vec<_> = relationship_updates.into_values().collect();
journal
.prepare_relationship_writes(&relationship_updates, &working_set.relationships);
let commit = self
.storage
.upsert_batch_fenced_journaled(
&[m.clone()],
&relationship_updates,
&journal.restore(),
journal.cycle_id.as_str(),
lease,
)
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
journal.record_commit(&commit);
working_set.memories.insert(m.id, m);
working_set.apply_relationship_writes(&relationship_updates);
res.merged.push(c.b);
}
}
Ok(())
}
async fn strengthen(
&self,
res: &mut ConsolidationResult,
lease: &exocortex_storage::OwnerLease,
working_set: &mut RegionWorkingSet,
journal: &mut CycleJournal,
) -> anyhow::Result<()> {
let now = chrono::Utc::now();
let mut updates = Vec::new();
let member_ids: std::collections::HashSet<exocortex_kernel::MemoryId> = working_set
.memories
.values()
.filter(|memory| memory.valid_until.is_none())
.map(|memory| memory.id)
.collect();
for row in working_set.relationships.values() {
let mut r = row.clone();
if !member_ids.contains(&r.from) || !member_ids.contains(&r.to) {
continue;
}
if r.valid_until.is_some() {
continue;
}
if matches!(r.provenance, Provenance::Computed { .. }) {
continue;
}
if res.strengthened.contains(&r.id) {
continue; }
journal.record_relationship(&r);
r.properties.evidence_count += 1;
let age_days = (now - r.recorded_at).num_days().max(0) as f32;
r.properties.strength = effective_strength(
r.properties.strength / decay_factor(age_days).max(f32::EPSILON),
r.properties.evidence_count,
r.properties.success_rate.unwrap_or(1.0),
age_days,
);
updates.push(r);
}
updates.sort_by_key(|relationship| relationship.id);
journal.prepare_relationship_writes(&updates, &working_set.relationships);
for r in &updates {
res.strengthened.push(r.id);
}
if !updates.is_empty() {
let commit = self
.storage
.upsert_batch_fenced_journaled(
&[],
&updates,
&journal.restore(),
journal.cycle_id.as_str(),
lease,
)
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
journal.record_commit(&commit);
working_set.apply_relationship_writes(&updates);
}
Ok(())
}
async fn write_similar_edges(
&self,
res: &mut ConsolidationResult,
survivors: &[MemoryWithEmbedding],
lease: &exocortex_storage::OwnerLease,
working_set: &mut RegionWorkingSet,
journal: &mut CycleJournal,
) -> anyhow::Result<()> {
let Some(similar_kind) = similar_to_kind() else {
return Ok(());
};
let now = chrono::Utc::now();
let mut edges = Vec::new();
for i in 0..survivors.len() {
for j in (i + 1)..survivors.len() {
let (a, b) = (&survivors[i], &survivors[j]);
if a.class != b.class || a.id == b.id {
continue;
}
let sim = mcr2::cosine(&a.embedding.vector, &b.embedding.vector);
if sim < SIMILAR_TO_THRESHOLD {
continue;
}
use exocortex_kernel::{
relationship_visibility, Relationship, RelationshipProperties, LSN,
};
edges.push(Relationship {
id: exocortex_kernel::RelationshipId::derive(a.id, similar_kind, b.id, None),
kind: similar_kind,
from: a.id,
to: b.id,
visibility: relationship_visibility(a.visibility, b.visibility),
provenance: Provenance::Computed {
producer: exocortex_kernel::provenance::ComputedProducer::SimilarityHnsw,
threshold: SIMILAR_TO_THRESHOLD,
},
properties: RelationshipProperties {
strength: sim.clamp(0.0, 1.0),
confidence: sim.clamp(0.0, 1.0),
context: None,
evidence_count: 1,
success_rate: None,
validation_count: 0,
counter_evidence_count: 0,
last_validated: now,
},
description: None,
bidirectional: true,
valid_from: now,
valid_until: None,
recorded_at: now,
invalidated_by: None,
lsn: LSN::new_local(0),
});
}
}
let mut fresh = Vec::new();
for edge in edges {
match working_set.relationships.get(&edge.id) {
Some(row) if row.valid_until.is_none() => continue,
Some(row) => journal.record_relationship(row),
None => journal.create_relationship(&edge),
}
fresh.push(edge);
}
if !fresh.is_empty() {
journal.prepare_relationship_writes(&fresh, &working_set.relationships);
let commit = self
.storage
.upsert_batch_fenced_journaled(
&[],
&fresh,
&journal.restore(),
journal.cycle_id.as_str(),
lease,
)
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
journal.record_commit(&commit);
res.similar_edges.extend(fresh.iter().map(|e| e.id));
working_set.apply_relationship_writes(&fresh);
metrics::counter!("exocortex_dreams_similar_edges_total").increment(fresh.len() as u64);
}
Ok(())
}
fn prune(&self, res: &mut ConsolidationResult, working_set: &RegionWorkingSet) {
for memory in working_set.memories.values() {
if memory.valid_until.is_some() {
res.pruned.push((memory.id, PruneReason::Redundant));
}
}
res.pruned.sort_by_key(|(id, _)| *id);
}
async fn rollback(
&self,
journal: &CycleJournal,
lease: &exocortex_storage::OwnerLease,
) -> anyhow::Result<()> {
#[cfg(feature = "testing")]
if let Some(pause) = self.rollback_pause {
tokio::time::sleep(pause).await;
}
#[cfg(feature = "testing")]
if !self.rollback_concurrent_memories.is_empty() {
self.storage
.upsert_batch(&self.rollback_concurrent_memories, &[])
.await
.map_err(|error| anyhow::anyhow!("inject concurrent rollback write: {error}"))?;
}
self.storage
.restore_fenced(&journal.restore(), lease)
.await
.map_err(|error| anyhow::anyhow!("{error}"))?;
Ok(())
}
async fn write_audit(&self, res: &ConsolidationResult) -> anyhow::Result<()> {
metrics::counter!("exocortex_dreams_discoveries_total", "quality" => "consolidation")
.increment((res.merged.len() + res.pruned.len()) as u64);
Ok(())
}
}
impl<S: Storage + 'static> DreamsEngine<S> {
fn validate_loaded_region(
&self,
region: &RegionKey,
working_set: &RegionWorkingSet,
) -> anyhow::Result<()> {
if region.org == "*" || region.project == "*" {
return Ok(());
}
if working_set
.memories
.values()
.any(|memory| memory.valid_until.is_none())
{
return Ok(());
}
anyhow::bail!("unknown project region {}/{}", region.org, region.project)
}
async fn validate_region(&self, region: &RegionKey) -> anyhow::Result<()> {
if region.org == "*" || region.project == "*" {
return Ok(());
}
let memories = self
.storage
.memories_in_region(region, MAX_REGION_MEMORIES as u32)
.await?;
let working_set = RegionWorkingSet {
memories: memories
.into_iter()
.map(|memory| (memory.id, memory))
.collect(),
relationships: Default::default(),
};
self.validate_loaded_region(region, &working_set)
}
pub async fn run_discovery(&self, region: &RegionKey) -> anyhow::Result<Vec<Discovery>> {
let lease_key = LeaseKey::Dreams {
org: region.org.clone(),
region: format!("{}:{}", region.project, region.memory_type).into(),
};
let lease = self
.storage
.acquire_lease(&lease_key, self.lease_ttl)
.await
.map_err(|error| anyhow::anyhow!("lease: {error}"))?;
let mut renewal = self.spawn_lease_renewal(lease.clone());
let outcome = {
let discovery = async {
self.validate_region(region).await?;
self.run_discovery_fenced(region, &lease, None).await
};
tokio::pin!(discovery);
tokio::select! {
biased;
renewal_outcome = &mut renewal => Err(Self::renewal_task_error(renewal_outcome)),
outcome = &mut discovery => outcome,
}
};
renewal.abort();
let _ = renewal.await;
let release = self.storage.release_lease(lease).await;
match (outcome, release) {
(Ok(discoveries), Ok(())) => Ok(discoveries),
(Err(error), Ok(())) => Err(error),
(Ok(_), Err(release)) => Err(anyhow::anyhow!("release Dreams lease: {release}")),
(Err(error), Err(release)) => Err(anyhow::anyhow!(
"discovery failed ({error}); lease release also failed ({release})"
)),
}
}
async fn run_discovery_fenced(
&self,
region: &RegionKey,
lease: &exocortex_storage::OwnerLease,
cycle_id: Option<&str>,
) -> anyhow::Result<Vec<Discovery>> {
let relationships = self
.storage
.relationships_in_region(region, MAX_DISCOVERY_PATH_INSPECTIONS as u32)
.await?;
let edges: Vec<DiscoveryEdge> = relationships
.into_iter()
.map(|relationship| {
(
relationship.from,
relationship.to,
relationship.kind.0,
matches!(relationship.provenance, Provenance::Derived { .. }),
)
})
.collect();
let (candidates, _) = transitive_candidates(&edges, MAX_DISCOVERY_PATH_INSPECTIONS);
let cycle: SmolStr = cycle_id
.map(SmolStr::new)
.unwrap_or_else(|| format!("dream:{}", uuid::Uuid::new_v4()).into());
let out: Vec<_> = candidates
.into_iter()
.map(|(a, c, k1, k2)| Discovery {
id: uuid::Uuid::new_v4(),
kind: DiscoveryKind::Transitive,
endpoints: (a, c),
quality: DiscoveryKind::Transitive.default_quality(),
via_types: (k1, k2),
discovery_cycle_id: cycle.clone(),
discovered_at: chrono::Utc::now(),
})
.collect();
let records = out
.iter()
.map(|discovery| exocortex_storage::DiscoveryRecord {
discovery_id: discovery.id.to_string().into(),
region: region.clone(),
from: discovery.endpoints.0,
to: discovery.endpoints.1,
discovery_type: "transitive".into(),
quality: discovery.quality,
via_types: [discovery.via_types.0, discovery.via_types.1],
discovery_cycle_id: discovery.discovery_cycle_id.clone(),
discovered_at: discovery.discovered_at,
})
.collect::<Vec<_>>();
if let Some(cycle_id) = cycle_id {
self.storage
.settle_dreams_cycle_fenced(cycle_id, &records, lease)
.await?;
} else {
for record in &records {
self.storage.store_discovery_fenced(record, lease).await?;
}
}
self.discoveries.clear();
for discovery in &out {
emit_discovery_metric(discovery);
self.discoveries.insert(discovery.id, discovery.clone());
}
Ok(out)
}
pub async fn issue_discovery_proposal(
&self,
discovery: &Discovery,
region: &RegionKey,
relationship_kind: exocortex_kernel::RelKindId,
proposed_visibility: exocortex_kernel::Visibility,
caller_scope: exocortex_storage::VisibilityContext,
) -> anyhow::Result<exocortex_storage::DiscoveryProposal> {
let proposal = exocortex_storage::DiscoveryProposal {
discovery_id: discovery.id.to_string().into(),
region: region.clone(),
from: discovery.endpoints.0,
to: discovery.endpoints.1,
kind: relationship_kind,
proposed_visibility,
caller_scope,
issued_at: discovery.discovered_at,
};
self.storage.create_discovery_proposal(&proposal).await?;
self.discoveries.remove(&discovery.id);
Ok(proposal)
}
pub fn pending_discoveries(&self) -> Vec<Discovery> {
self.discoveries.iter().map(|e| e.value().clone()).collect()
}
}
fn transitive_candidates(
edges: &[DiscoveryEdge],
inspection_budget: usize,
) -> (Vec<(MemoryId, MemoryId, u32, u32)>, usize) {
let direct: std::collections::HashSet<(MemoryId, MemoryId)> =
edges.iter().map(|(from, to, _, _)| (*from, *to)).collect();
let mut outgoing: std::collections::BTreeMap<MemoryId, Vec<&DiscoveryEdge>> =
std::collections::BTreeMap::new();
for edge in edges {
outgoing.entry(edge.0).or_default().push(edge);
}
let mut seen = std::collections::BTreeSet::new();
let mut out = Vec::new();
let mut inspected = 0usize;
'paths: for (a, b, k1, first_derived) in edges {
if *first_derived {
continue;
}
let Some(next_edges) = outgoing.get(b) else {
continue;
};
for (_, c, k2, second_derived) in next_edges {
if inspected >= inspection_budget {
break 'paths;
}
inspected += 1;
let candidate = (*a, *c, *k1, *k2);
if !*second_derived && c != a && !direct.contains(&(*a, *c)) && seen.insert(candidate) {
out.push(candidate);
if out.len() >= MAX_DISCOVERIES_PER_CYCLE {
break 'paths;
}
}
}
}
(out, inspected)
}
#[cfg(test)]
mod cycle_scheduling_tests {
use super::*;
use exocortex_storage::InMemoryStorage;
fn region() -> RegionKey {
RegionKey {
org: "org".into(),
project: "project".into(),
memory_type: 3,
}
}
fn engine(storage: &InMemoryStorage) -> DreamsEngine<InMemoryStorage> {
DreamsEngine::new(
Arc::new(storage.clone_dyn()),
DreamsTrigger {
memory_threshold: 1,
edge_threshold: u32::MAX,
age_floor_days: u32::MAX,
min_interval_hours: 0,
},
0.01,
0.05,
false,
"cycle-test".into(),
)
}
fn storage() -> InMemoryStorage {
InMemoryStorage::new(Arc::new(
exocortex_kernel::Ontology::from_packs(vec![exocortex_pack_dev_v1::pack_def()])
.expect("development ontology"),
))
}
#[tokio::test]
async fn region_has_one_pending_cycle_and_retains_post_fire_writes() {
let storage = storage();
let engine = engine(&storage);
let region = region();
engine.on_write(region.clone()).await.unwrap();
engine.on_write(region.clone()).await.unwrap();
engine.notify(region.clone());
assert_eq!(engine.pending_regions.len(), 1);
let fired_at = {
let mut receiver = engine.rx_fire.lock().await;
let (_, fired_at) = receiver.try_recv().expect("one scheduled cycle");
assert!(receiver.try_recv().is_err(), "the region must be coalesced");
fired_at
};
assert_eq!(fired_at.memories_since_last_cycle, 1);
engine.complete_region(®ion, fired_at, true);
assert_eq!(
engine
.counters
.get(®ion)
.unwrap()
.memories_since_last_cycle,
1,
"the write after fire remains pending"
);
assert_eq!(engine.pending_regions.len(), 1);
let second = {
let mut receiver = engine.rx_fire.lock().await;
receiver
.try_recv()
.expect("retained write schedules next cycle")
};
assert_eq!(second.1.memories_since_last_cycle, 1);
engine.complete_region(®ion, second.1, true);
assert_eq!(
engine
.counters
.get(®ion)
.unwrap()
.memories_since_last_cycle,
0
);
assert!(engine.pending_regions.is_empty());
}
#[tokio::test]
async fn stable_local_write_event_is_counted_once() {
let storage = storage();
let engine = engine(&storage);
let region = region();
engine
.on_writes_once("batch:region", 1, region.clone(), 2, 3)
.await
.unwrap();
engine
.on_writes_once("batch:region", 1, region.clone(), 2, 3)
.await
.unwrap();
let counters = *engine.counters.get(®ion).unwrap();
assert_eq!(counters.memories_since_last_cycle, 2);
assert_eq!(counters.edges_since_last_cycle, 3);
}
#[tokio::test]
async fn failed_cycle_keeps_counters_for_a_later_retry() {
let storage = storage();
let engine = engine(&storage);
let region = region();
engine.on_write(region.clone()).await.unwrap();
let fired_at = {
let mut receiver = engine.rx_fire.lock().await;
receiver.try_recv().expect("scheduled cycle").1
};
engine.complete_region(®ion, fired_at, false);
assert_eq!(
engine
.counters
.get(®ion)
.unwrap()
.memories_since_last_cycle,
1
);
assert!(engine.pending_regions.is_empty());
engine.notify(region.clone());
let retry = engine.rx_fire.lock().await.try_recv().expect("retry cycle");
assert_eq!(retry.1.memories_since_last_cycle, 1);
}
#[tokio::test]
async fn follower_rejects_cycle_before_any_region_scan() {
let storage = storage();
let leader = Arc::new(std::sync::atomic::AtomicBool::new(false));
let engine = engine(&storage).with_leader_gate(leader);
let error = engine.try_consolidate(®ion()).await.unwrap_err();
assert!(error.to_string().contains("not the elected leader"));
assert_eq!(
storage.reasoning_query_counts(),
(0, 0, 0, 0),
"a follower must not bulk-scan storage"
);
}
}
#[cfg(test)]
mod discovery_scaling_tests {
use super::*;
#[test]
fn indexed_discovery_is_linear_for_disjoint_edges_and_budgeted_for_dense_paths() {
let disjoint: Vec<_> = (0..10_000u32)
.map(|i| {
let mut from = [0u8; 16];
from[..4].copy_from_slice(&(i * 2).to_be_bytes());
let mut to = [0u8; 16];
to[..4].copy_from_slice(&(i * 2 + 1).to_be_bytes());
(MemoryId(from), MemoryId(to), 1, false)
})
.collect();
let (candidates, inspected) = transitive_candidates(&disjoint, 50_000);
assert!(candidates.is_empty());
assert_eq!(
inspected, 0,
"unconnected edges require no global pair scan"
);
let hub = MemoryId([0x7f; 16]);
let mut dense = Vec::new();
for i in 0..400u32 {
let mut endpoint = [0u8; 16];
endpoint[..4].copy_from_slice(&i.to_be_bytes());
dense.push((MemoryId(endpoint), hub, 1, false));
dense.push((hub, MemoryId(endpoint), 2, true));
}
dense.sort_by_key(|edge| (edge.0, edge.1, edge.2, edge.3));
let (_, inspected) = transitive_candidates(&dense, 97);
assert_eq!(inspected, 97, "dense candidate work stops at its budget");
}
}
#[derive(Clone, Debug)]
pub struct Discovery {
pub id: uuid::Uuid,
pub kind: DiscoveryKind,
pub endpoints: (MemoryId, MemoryId),
pub quality: f32,
pub via_types: (u32, u32),
pub discovery_cycle_id: SmolStr,
pub discovered_at: chrono::DateTime<chrono::Utc>,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum DiscoveryKind {
CrossDomain,
TemporalEcho,
Orphan,
Transitive,
}
impl DiscoveryKind {
pub fn default_quality(self) -> f32 {
match self {
DiscoveryKind::CrossDomain => 0.9,
DiscoveryKind::TemporalEcho => 0.7,
DiscoveryKind::Orphan => 0.4,
DiscoveryKind::Transitive => 0.6,
}
}
}
impl Discovery {
pub fn rate_quality(&self) -> f32 {
self.quality
}
}
fn emit_discovery_metric(discovery: &Discovery) {
match (discovery.kind, discovery.quality.to_bits()) {
(DiscoveryKind::CrossDomain, bits) if bits == 0.9_f32.to_bits() => {
metrics::counter!("exocortex_dreams_discoveries_total", "type" => "cross_domain", "quality" => "0.9").increment(1);
}
(DiscoveryKind::TemporalEcho, bits) if bits == 0.7_f32.to_bits() => {
metrics::counter!("exocortex_dreams_discoveries_total", "type" => "temporal_echo", "quality" => "0.7").increment(1);
}
(DiscoveryKind::Orphan, bits) if bits == 0.4_f32.to_bits() => {
metrics::counter!("exocortex_dreams_discoveries_total", "type" => "orphan", "quality" => "0.4").increment(1);
}
(DiscoveryKind::Transitive, bits) if bits == 0.6_f32.to_bits() => {
metrics::counter!("exocortex_dreams_discoveries_total", "type" => "transitive", "quality" => "0.6").increment(1);
}
_ => metrics::counter!("exocortex_dreams_discoveries_total", "type" => "invalid", "quality" => "invalid").increment(1),
}
}
#[cfg(test)]
mod discovery_metric_tests {
use super::*;
use std::sync::Mutex;
#[derive(Default)]
struct CaptureRecorder {
labels: Mutex<Vec<(String, String)>>,
}
impl metrics::Recorder for CaptureRecorder {
fn describe_counter(
&self,
_: metrics::KeyName,
_: Option<metrics::Unit>,
_: metrics::SharedString,
) {
}
fn describe_gauge(
&self,
_: metrics::KeyName,
_: Option<metrics::Unit>,
_: metrics::SharedString,
) {
}
fn describe_histogram(
&self,
_: metrics::KeyName,
_: Option<metrics::Unit>,
_: metrics::SharedString,
) {
}
fn register_counter(
&self,
key: &metrics::Key,
_: &metrics::Metadata<'_>,
) -> metrics::Counter {
self.labels.lock().unwrap().extend(
key.labels()
.map(|label| (label.key().to_owned(), label.value().to_owned())),
);
metrics::Counter::noop()
}
fn register_gauge(&self, _: &metrics::Key, _: &metrics::Metadata<'_>) -> metrics::Gauge {
metrics::Gauge::noop()
}
fn register_histogram(
&self,
_: &metrics::Key,
_: &metrics::Metadata<'_>,
) -> metrics::Histogram {
metrics::Histogram::noop()
}
}
#[test]
fn discovery_metric_reports_the_stamped_quality() {
let recorder = CaptureRecorder::default();
let _guard = metrics::set_default_local_recorder(&recorder);
let discovery = Discovery {
id: uuid::Uuid::nil(),
kind: DiscoveryKind::Transitive,
endpoints: (MemoryId([0; 16]), MemoryId([1; 16])),
quality: DiscoveryKind::Transitive.default_quality(),
via_types: (1, 2),
discovery_cycle_id: "cycle".into(),
discovered_at: chrono::Utc::now(),
};
emit_discovery_metric(&discovery);
let labels = recorder.labels.lock().unwrap();
assert!(labels.contains(&("quality".into(), discovery.rate_quality().to_string())));
}
}
pub fn discovery_provenance(score: f32) -> Provenance {
Provenance::Proposed {
discovery_id: uuid::Uuid::new_v4(),
score,
}
}
fn similar_to_kind() -> Option<exocortex_kernel::RelKindId> {
static KIND: std::sync::OnceLock<Option<exocortex_kernel::RelKindId>> =
std::sync::OnceLock::new();
*KIND.get_or_init(|| {
exocortex_kernel::Ontology::from_packs(vec![exocortex_pack_dev_v1::pack_def()])
.ok()
.and_then(|onto| onto.kind_id("SimilarTo"))
})
}
fn decay_factor(age_days: f32) -> f32 {
(1.0 - 0.01 * age_days).max(0.5)
}
impl<S: Storage + 'static> DreamsEngine<S> {
async fn empty_result(
&self,
region: &RegionKey,
lease: &exocortex_storage::OwnerLease,
n: u32,
working_set: &RegionWorkingSet,
) -> anyhow::Result<ConsolidationResult> {
let sparsity = self.sparsity(working_set);
let res = ConsolidationResult {
session_id: format!("dream:{}", uuid::Uuid::new_v4()).into(),
user_id: None,
started_at: chrono::Utc::now(),
completed_at: chrono::Utc::now(),
region: region.clone(),
memories_input: n,
memories_output: n,
mcr2_before: crate::mcr2::MCR2Value {
delta_r: 0.0,
total_rate: 0.0,
class_rates: Default::default(),
compact_rate: 0.0,
n_memories: 0,
embedding_model: crate::mcr2::EmbeddingModelId::bge_small(),
computed_at: chrono::Utc::now(),
},
mcr2_after: crate::mcr2::MCR2Value {
delta_r: 0.0,
total_rate: 0.0,
class_rates: Default::default(),
compact_rate: 0.0,
n_memories: 0,
embedding_model: crate::mcr2::EmbeddingModelId::bge_small(),
computed_at: chrono::Utc::now(),
},
sparsity_before: sparsity.clone(),
sparsity_after: sparsity,
merged: vec![],
abstracted: vec![],
pruned: vec![],
strengthened: vec![],
rewired: vec![],
similar_edges: vec![],
owner_node_id: self.node_id.clone(),
lease_epoch: lease.epoch,
regression: false,
hairball_regression: false,
};
self.write_audit(&res).await?;
Ok(res)
}
}