use std::collections::BTreeMap;
use std::sync::atomic::{AtomicU64, Ordering};
#[cfg(not(target_arch = "wasm32"))]
use std::time::{SystemTime, UNIX_EPOCH};
fn now_nanos() -> u128 {
#[cfg(target_arch = "wasm32")]
{
0
}
#[cfg(not(target_arch = "wasm32"))]
{
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|elapsed| elapsed.as_nanos())
.unwrap_or(0)
}
}
fn now_unix_secs() -> u64 {
#[cfg(target_arch = "wasm32")]
{
0
}
#[cfg(not(target_arch = "wasm32"))]
{
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|elapsed| elapsed.as_secs())
.unwrap_or(0)
}
}
use serde_json::{Map, Number, Value};
use super::{positive_ttl, MemoryService, Metadata, HUB_FIELD};
use crate::context::model::{
CompilePolicy, CompileRequest, CompiledContext, ContextDecision, ContextFragment,
ContextSavings, ContextSource, ImportanceWeights, LoadedWorkingContext, MediaRef, MemoryScope,
WorkingContext, WorkingContextIndex, WorkingContextSession,
};
use crate::context::{media, provenance, ContextCompiler};
use crate::embedder::Embedder;
use crate::error::MemoryError;
use crate::id::stable_id;
use crate::model::FusionOptions;
use crate::storage::{FactStore, GraphStore, RecallStore};
const SOURCE_ID_SALT: &str = "veles-ctx-source:";
const EVENT_ID_SALT: &str = "veles-ctx-event:";
const WORKING_ID_SALT: &str = "veles-ctx-working:";
const WORKING_INDEX_ID_SALT: &str = "veles-ctx-working-index:";
const EVENT_ANCHOR: &str = "veles context compilation event";
use crate::storage::{
CTX_EVENT_FIELD, CTX_SOURCE_FIELD, CTX_WORKING_FIELD, CTX_WORKING_INDEX_FIELD,
};
const CTX_PROJECT_FIELD: &str = "_veles_ctx_project";
const CTX_MODEL_FIELD: &str = "_veles_ctx_model";
const CTX_SOURCE_MEDIA_FIELD: &str = "_veles_ctx_source_media";
const EXPIRES_AT_FIELD: &str = "_veles_expires_at";
const CTX_SESSION_FIELD: &str = "_veles_ctx_session";
const CTX_TOKENS_IN_FIELD: &str = "_veles_ctx_tokens_in";
const CTX_TOKENS_OUT_FIELD: &str = "_veles_ctx_tokens_out";
const CTX_TOKENS_SAVED_FIELD: &str = "_veles_ctx_tokens_saved";
const CTX_COST_FIELD: &str = "_veles_ctx_cost_micros";
const CTX_CURRENCY_FIELD: &str = "_veles_ctx_currency";
const CTX_AT_FIELD: &str = "_veles_ctx_at";
static EVENT_SEQ: AtomicU64 = AtomicU64::new(0);
static WORKING_INDEX_WRITE: parking_lot::Mutex<()> = parking_lot::Mutex::new(());
#[path = "memory_bridge_compile.rs"]
mod compile;
impl<E: Embedder, S: FactStore> MemoryService<E, S> {
pub fn save_working_context(
&self,
project: &str,
session: &str,
working: &WorkingContext,
) -> Result<u64, MemoryError> {
let _generation = self.enter_generation();
if working.is_empty() {
return Err(MemoryError::EmptyWorkingContext);
}
let content =
serde_json::to_string(working).map_err(|err| MemoryError::WorkingContextCodec {
detail: "encoding the working context for storage".to_owned(),
source: Some(Box::new(err)),
})?;
if content.len() > crate::limits::MAX_FACT_BYTES {
return Err(MemoryError::ContextOverLimit(format!(
"working context of {} bytes exceeds the cap of {} bytes",
content.len(),
crate::limits::MAX_FACT_BYTES
)));
}
let id = working_id(project, session);
let embedding = self
.embedder
.embed(&format!("working context {project} {session}"))?;
let meta = system_meta(&[
(CTX_WORKING_FIELD, Value::Bool(true)),
(CTX_PROJECT_FIELD, Value::String(project.to_owned())),
(CTX_SESSION_FIELD, Value::String(session.to_owned())),
]);
self.store_fact(id, &content, &embedding, Some(&meta), None)?;
self.update_working_index(project, session)?;
Ok(id)
}
pub fn load_working_context(
&self,
project: &str,
session: &str,
) -> Result<Option<WorkingContext>, MemoryError> {
let _generation = self.enter_generation();
self.load_working_context_inner(project, session)
}
fn load_working_context_inner(
&self,
project: &str,
session: &str,
) -> Result<Option<WorkingContext>, MemoryError> {
let slot = working_id(project, session);
let payloads = self.store.get_metadata_batch(&[slot])?;
let marked = payloads
.into_iter()
.next()
.flatten()
.is_some_and(|meta| meta.get(CTX_WORKING_FIELD) == Some(&Value::Bool(true)));
if !marked {
return Ok(None);
}
let Some((content, _)) = self.store.get(slot)? else {
return Err(MemoryError::WorkingContextCodec {
detail: format!(
"working context for project '{project}', session '{session}' is corrupt: \
the reserved marker is present but the stored body is gone"
),
source: None,
});
};
serde_json::from_str(&content)
.map(Some)
.map_err(|err| MemoryError::WorkingContextCodec {
detail: format!(
"decoding the stored working context for project '{project}', \
session '{session}'"
),
source: Some(Box::new(err)),
})
}
pub fn resume_working_context(
&self,
project: &str,
session: &str,
) -> Result<LoadedWorkingContext, MemoryError> {
let _generation = self.enter_generation();
let working = self.load_working_context_inner(project, session)?;
let other_sessions = self.other_sessions_for(project, session, working.is_some())?;
Ok(LoadedWorkingContext {
found: working.is_some(),
working,
other_sessions,
})
}
fn other_sessions_for(
&self,
project: &str,
session: &str,
found: bool,
) -> Result<Vec<String>, MemoryError> {
let listed = match self.list_working_contexts_inner(project) {
Ok(listed) => listed,
Err(_) if found => return Ok(Vec::new()),
Err(err) => return Err(err),
};
Ok(listed
.into_iter()
.map(|entry| entry.session)
.filter(|candidate| candidate != session)
.collect())
}
fn live_sessions(
&self,
project: &str,
sessions: Vec<WorkingContextSession>,
) -> Result<Vec<WorkingContextSession>, MemoryError> {
if sessions.is_empty() {
return Ok(sessions);
}
let ids: Vec<u64> = sessions
.iter()
.map(|entry| working_id(project, &entry.session))
.collect();
let payloads = self.store.get_metadata_batch(&ids)?;
if payloads.len() != ids.len() {
return Err(MemoryError::WorkingContextCodec {
detail: format!(
"storage returned {} metadata rows for {} working-context ids",
payloads.len(),
ids.len()
),
source: None,
});
}
Ok(sessions
.into_iter()
.zip(payloads)
.filter(|(_, meta)| {
meta.as_ref()
.is_some_and(|meta| meta.get(CTX_WORKING_FIELD) == Some(&Value::Bool(true)))
})
.map(|(entry, _)| entry)
.collect())
}
pub fn list_working_contexts(
&self,
project: &str,
) -> Result<Vec<WorkingContextSession>, MemoryError> {
let _generation = self.enter_generation();
self.list_working_contexts_inner(project)
}
fn list_working_contexts_inner(
&self,
project: &str,
) -> Result<Vec<WorkingContextSession>, MemoryError> {
let Some(index) = self.working_index(project)? else {
return Ok(Vec::new());
};
let mut sessions = self.live_sessions(project, index.sessions)?;
sessions.sort_by(|a, b| {
b.saved_at
.cmp(&a.saved_at)
.then_with(|| a.session.cmp(&b.session))
});
Ok(sessions)
}
fn working_index(&self, project: &str) -> Result<Option<WorkingContextIndex>, MemoryError> {
let slot = working_index_id(project);
let payloads = self.store.get_metadata_batch(&[slot])?;
let marked = payloads
.into_iter()
.next()
.flatten()
.is_some_and(|meta| meta.get(CTX_WORKING_INDEX_FIELD) == Some(&Value::Bool(true)));
if !marked {
return Ok(None);
}
match self.store.get(slot)? {
Some((content, _)) => serde_json::from_str(&content).map(Some).map_err(|err| {
MemoryError::WorkingContextCodec {
detail: format!("decoding the working-context index for project '{project}'"),
source: Some(Box::new(err)),
}
}),
None => Err(MemoryError::WorkingContextCodec {
detail: format!(
"working-context index for project '{project}' is corrupt: the index \
marker is present but the stored body is gone"
),
source: None,
}),
}
}
fn update_working_index(&self, project: &str, session: &str) -> Result<(), MemoryError> {
let embedding = self
.embedder
.embed(&format!("working context index {project}"))?;
let _guard = WORKING_INDEX_WRITE.lock();
let mut index = match self.working_index(project) {
Ok(index) => index.unwrap_or_default(),
Err(MemoryError::WorkingContextCodec { .. }) => WorkingContextIndex::default(),
Err(err) => return Err(err),
};
let now = now_unix_secs();
if let Some(entry) = index.sessions.iter_mut().find(|s| s.session == session) {
entry.saved_at = now;
} else {
index.sessions.push(WorkingContextSession {
session: session.to_owned(),
saved_at: now,
});
}
index.sessions = self.live_sessions(project, index.sessions)?;
let content =
serde_json::to_string(&index).map_err(|err| MemoryError::WorkingContextCodec {
detail: format!("encoding the working-context index for project '{project}'"),
source: Some(Box::new(err)),
})?;
self.write_working_index(project, &content, &embedding)
}
fn write_working_index(
&self,
project: &str,
content: &str,
embedding: &[f32],
) -> Result<(), MemoryError> {
let slot = working_index_id(project);
let meta = system_meta(&[
(CTX_WORKING_INDEX_FIELD, Value::Bool(true)),
(CTX_PROJECT_FIELD, Value::String(project.to_owned())),
]);
self.store_fact(slot, content, embedding, Some(&meta), None)?;
Ok(())
}
}
const DEFAULT_MEMORY_K: usize = 5;
fn scope_and_k(request: &CompileRequest) -> Option<(&MemoryScope, usize)> {
let scope = request.memory_scope.as_ref()?;
let room = crate::limits::MAX_FRAGMENTS.saturating_sub(request.fragments.len());
let k = crate::limits::clamp_recall_limit(scope.k.unwrap_or(DEFAULT_MEMORY_K)).min(room);
(k > 0).then_some((scope, k))
}
fn scope_filter(scope: &MemoryScope) -> Option<Metadata> {
scope.project.as_ref().map(|project| {
let mut meta = Map::new();
meta.insert("project".to_owned(), Value::String(project.clone()));
meta
})
}
struct PulledMemory {
fragment: ContextFragment,
memory_id: u64,
relevance: f32,
vector_norm: f64,
graph_weight: f64,
confidence: f64,
recency: f64,
ventilated: bool,
}
struct MemoryCandidate {
memory_id: u64,
base: f64,
vector_norm: f64,
graph_weight: f64,
metadata: Option<Metadata>,
content: String,
}
impl MemoryCandidate {
fn into_pulled(self) -> PulledMemory {
#[allow(clippy::cast_possible_truncation)] let relevance = self.base as f32;
PulledMemory {
fragment: ContextFragment {
id: None,
content: self.content,
path: None,
kind: Some("memory".to_owned()),
priority: None,
metadata: None,
media: None,
},
memory_id: self.memory_id,
relevance,
vector_norm: self.vector_norm,
graph_weight: self.graph_weight,
confidence: NEUTRAL_CONFIDENCE,
recency: 0.0,
ventilated: false,
}
}
}
const NEUTRAL_CONFIDENCE: f64 = 0.5;
#[cfg(feature = "persistence")]
fn payload_confidence(payload: Option<&Metadata>) -> f64 {
f64::from(payload.map_or(
super::reinforce::RL_NEUTRAL_CONFIDENCE,
super::reinforce::read_confidence,
))
}
#[cfg(not(feature = "persistence"))]
fn payload_confidence(_payload: Option<&Metadata>) -> f64 {
NEUTRAL_CONFIDENCE
}
#[allow(
clippy::float_cmp,
reason = "an exact zero weight is the documented off switch; any non-zero weight, however small, is active"
)]
fn importance_active(weights: &ImportanceWeights) -> bool {
weights.confidence != 0.0 || (weights.recency != 0.0 && weights.recency_field.is_some())
}
#[allow(
clippy::float_cmp,
reason = "an exact zero weight is the documented off switch for the recency term"
)]
fn recency_norms(candidates: &[MemoryCandidate], weights: &ImportanceWeights) -> Vec<f64> {
let field = weights
.recency_field
.as_ref()
.filter(|_| weights.recency != 0.0);
let Some(field) = field else {
return vec![0.0; candidates.len()];
};
let values: Vec<Option<f64>> = candidates
.iter()
.map(|candidate| {
candidate
.metadata
.as_ref()
.and_then(|meta| meta.get(field.as_str()))
.and_then(Value::as_f64)
.filter(|value| value.is_finite())
})
.collect();
let (min, max) = values
.iter()
.flatten()
.fold((f64::INFINITY, f64::NEG_INFINITY), |(lo, hi), &v| {
(lo.min(v), hi.max(v))
});
if max <= min {
return vec![0.0; candidates.len()];
}
values
.into_iter()
.map(|value| value.map_or(0.0, |v| ((v - min) / (max - min)).clamp(0.0, 1.0)))
.collect()
}
fn annotate_memory_provenance(out: &mut CompiledContext, pulled: &BTreeMap<u64, PulledMemory>) {
for decision in &mut out.decisions {
if let Some(memory) = pulled.get(&decision.content_hash) {
decision.memory_id = Some(memory.memory_id);
decision.relevance = memory.relevance;
decision.reason = if memory.ventilated {
format!(
"{} — pulled from memory {} (vector {:.2}, graph {:.2}, confidence {:.2}, recency {:.2})",
decision.reason,
memory.memory_id,
memory.vector_norm,
memory.graph_weight,
memory.confidence,
memory.recency
)
} else {
format!(
"{} — pulled from memory {} (vector {:.2}, graph {:.2})",
decision.reason, memory.memory_id, memory.vector_norm, memory.graph_weight
)
};
}
}
for source in &mut out.sources {
if let Some(hash) = provenance::parse_handle(&source.handle) {
if let Some(memory) = pulled.get(&hash) {
source.memory_id = Some(memory.memory_id);
}
}
}
}
fn system_meta(extra: &[(&str, Value)]) -> Metadata {
let mut meta = Map::new();
meta.insert(HUB_FIELD.to_owned(), Value::Bool(true));
for (key, value) in extra {
meta.insert((*key).to_owned(), value.clone());
}
meta
}
fn event_meta(request: &CompileRequest, out: &CompiledContext, nanos: u128) -> Metadata {
let mut extra: Vec<(&str, Value)> = vec![
(CTX_EVENT_FIELD, Value::Bool(true)),
(
CTX_TOKENS_IN_FIELD,
Value::Number(out.insights.tokens_in.into()),
),
(
CTX_TOKENS_OUT_FIELD,
Value::Number(out.insights.tokens_out.into()),
),
(
CTX_TOKENS_SAVED_FIELD,
Value::Number(out.insights.tokens_saved.into()),
),
(
CTX_AT_FIELD,
Value::Number(Number::from(
u64::try_from(nanos / 1_000_000_000).unwrap_or(u64::MAX),
)),
),
];
if let Some(project) = &request.project {
extra.push((CTX_PROJECT_FIELD, Value::String(project.clone())));
}
if let Some(model) = &request.target_model {
extra.push((CTX_MODEL_FIELD, Value::String(model.clone())));
}
if let (Some(micros), Some(currency)) = (
out.insights.estimated_cost_saved_micros,
out.insights.currency.as_ref(),
) {
extra.push((CTX_COST_FIELD, Value::Number(micros.into())));
extra.push((CTX_CURRENCY_FIELD, Value::String(currency.clone())));
}
system_meta(&extra)
}
fn aggregate_events(payloads: &[Option<Metadata>]) -> ContextSavings {
let mut savings = ContextSavings {
events: payloads.len() as u64,
truncated: payloads.len() >= crate::limits::MAX_RECALL_LIMIT,
..ContextSavings::default()
};
for payload in payloads {
let Some(meta) = payload else { continue };
savings.tokens_in = savings
.tokens_in
.saturating_add(meta_u64(meta, CTX_TOKENS_IN_FIELD));
savings.tokens_out = savings
.tokens_out
.saturating_add(meta_u64(meta, CTX_TOKENS_OUT_FIELD));
savings.tokens_saved = savings
.tokens_saved
.saturating_add(meta_u64(meta, CTX_TOKENS_SAVED_FIELD));
if let (Some(Value::String(currency)), micros) =
(meta.get(CTX_CURRENCY_FIELD), meta_u64(meta, CTX_COST_FIELD))
{
if micros > 0 {
let entry = savings
.cost_saved_micros_by_currency
.entry(currency.clone())
.or_insert(0);
*entry = entry.saturating_add(micros);
}
}
}
savings
}
fn meta_u64(meta: &Metadata, key: &str) -> u64 {
meta.get(key).and_then(Value::as_u64).unwrap_or(0)
}
fn source_id(content_hash: u64) -> u64 {
stable_id(&format!("{SOURCE_ID_SALT}{content_hash}"))
}
fn fragment_handle_hash(fragment: &ContextFragment) -> u64 {
fragment.media.as_ref().map_or_else(
|| stable_id(&fragment.content),
|media_ref| media::analyze(media_ref).raw_hash,
)
}
fn index_fragments_by_handle_hash(
fragments: &[ContextFragment],
) -> BTreeMap<u64, &ContextFragment> {
let mut by_hash: BTreeMap<u64, &ContextFragment> = BTreeMap::new();
for fragment in fragments {
by_hash
.entry(fragment_handle_hash(fragment))
.or_insert(fragment);
}
by_hash
}
fn source_media(meta: &Metadata) -> Option<MediaRef> {
meta.get(CTX_SOURCE_MEDIA_FIELD)
.cloned()
.and_then(|value| serde_json::from_value(value).ok())
}
fn working_id(project: &str, session: &str) -> u64 {
stable_id(&format!("{WORKING_ID_SALT}{project}\u{1f}{session}"))
}
fn working_index_id(project: &str) -> u64 {
stable_id(&format!("{WORKING_INDEX_ID_SALT}{project}"))
}
#[cfg(all(test, feature = "persistence"))]
#[path = "memory_bridge_tests.rs"]
mod tests;