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::{
CompileRequest, CompiledContext, ContextFragment, ContextSavings, ContextSource,
ImportanceWeights, 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::MemoryStore;
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";
const CTX_EVENT_FIELD: &str = "_veles_ctx_event";
const CTX_PROJECT_FIELD: &str = "_veles_ctx_project";
const CTX_MODEL_FIELD: &str = "_veles_ctx_model";
const CTX_SOURCE_FIELD: &str = "_veles_ctx_source";
const CTX_SOURCE_MEDIA_FIELD: &str = "_veles_ctx_source_media";
const EXPIRES_AT_FIELD: &str = "_veles_expires_at";
const CTX_WORKING_FIELD: &str = "_veles_ctx_working";
const CTX_WORKING_INDEX_FIELD: &str = "_veles_ctx_working_index";
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);
impl<E: Embedder, S: MemoryStore> MemoryService<E, S> {
pub fn compile_context(
&self,
compiler: &ContextCompiler,
request: &CompileRequest,
) -> Result<CompiledContext, MemoryError> {
let importance = compiler.effective_policy(request).importance.clone();
let memories = self.context_memories(request, &importance)?;
self.compile_with_memories(compiler, request, memories)
}
pub fn compile_context_reranked<R: crate::Reranker>(
&self,
compiler: &ContextCompiler,
request: &CompileRequest,
reranker: &R,
) -> Result<CompiledContext, MemoryError> {
let importance = compiler.effective_policy(request).importance.clone();
let memories = self.context_memories_reranked(request, reranker, &importance)?;
self.compile_with_memories(compiler, request, memories)
}
fn compile_with_memories(
&self,
compiler: &ContextCompiler,
request: &CompileRequest,
memories: Vec<PulledMemory>,
) -> Result<CompiledContext, MemoryError> {
let mut augmented = request.clone();
let mut pulled: BTreeMap<u64, PulledMemory> = BTreeMap::new();
for memory in memories {
augmented.fragments.push(memory.fragment.clone());
pulled.insert(stable_id(&memory.fragment.content), memory);
}
let mut out = compiler.compile_raw(&augmented)?;
annotate_memory_provenance(&mut out, &pulled);
out.warnings = crate::context::warnings_for(&out.decisions);
let policy = compiler.effective_policy(request);
if policy.store_sources {
self.store_context_sources(&augmented, &out, policy.source_ttl_seconds)?;
}
if policy.record_events {
self.record_context_event(request, &out, policy.event_ttl_seconds)?;
}
Ok(crate::context::apply_slim(out, policy))
}
fn context_memories(
&self,
request: &CompileRequest,
importance: &ImportanceWeights,
) -> Result<Vec<PulledMemory>, MemoryError> {
let Some((scope, k)) = scope_and_k(request) else {
return Ok(Vec::new());
};
let filter = scope_filter(scope);
let opts = FusionOptions::from_knobs(scope.hops, scope.graph_boost, None);
let scored = self.recall_fused_scored(&request.query, k, filter.as_ref(), opts)?;
let max_fused = scored
.iter()
.map(|s| s.fused)
.fold(f64::MIN, f64::max)
.max(f64::EPSILON);
let candidates = scored
.into_iter()
.map(|scored| {
let fused = if scored.fused.is_finite() {
scored.fused
} else {
0.0
};
MemoryCandidate {
memory_id: scored.recollection.id,
base: (fused / max_fused).clamp(0.0, 1.0),
vector_norm: scored.vector_norm,
graph_weight: scored.graph_weight,
metadata: scored.recollection.metadata,
content: scored.recollection.content,
}
})
.collect();
self.blend_importance(candidates, importance)
}
fn context_memories_reranked<R: crate::Reranker>(
&self,
request: &CompileRequest,
reranker: &R,
importance: &ImportanceWeights,
) -> Result<Vec<PulledMemory>, MemoryError> {
let Some((scope, k)) = scope_and_k(request) else {
return Ok(Vec::new());
};
let filter = scope_filter(scope);
let opts = FusionOptions::from_knobs(scope.hops, scope.graph_boost, None);
let ranked =
self.recall_fused_reranked(&request.query, k, filter.as_ref(), opts, reranker)?;
let count = ranked.len().max(1);
let candidates = ranked
.into_iter()
.enumerate()
.map(|(rank, recollection)| {
#[allow(clippy::cast_precision_loss)] let relevance = 1.0 - (rank as f32 / count as f32);
MemoryCandidate {
memory_id: recollection.id,
base: f64::from(relevance),
vector_norm: 0.0,
graph_weight: 0.0,
metadata: recollection.metadata,
content: recollection.content,
}
})
.collect();
self.blend_importance(candidates, importance)
}
fn blend_importance(
&self,
candidates: Vec<MemoryCandidate>,
weights: &ImportanceWeights,
) -> Result<Vec<PulledMemory>, MemoryError> {
if !importance_active(weights) {
return Ok(candidates
.into_iter()
.map(MemoryCandidate::into_pulled)
.collect());
}
let ids: Vec<u64> = candidates.iter().map(|c| c.memory_id).collect();
let raw = self.store.get_metadata_batch(&ids)?;
let recencies = recency_norms(&candidates, weights);
let mut blended: Vec<(f64, PulledMemory)> = candidates
.into_iter()
.zip(raw)
.zip(recencies)
.map(|((candidate, payload), recency)| {
let confidence = payload_confidence(payload.as_ref());
let score = candidate.base
+ weights.confidence * (confidence - NEUTRAL_CONFIDENCE) * 2.0
+ weights.recency * recency;
let mut pulled = candidate.into_pulled();
#[allow(clippy::cast_possible_truncation)] {
pulled.relevance = score.clamp(0.0, 1.0) as f32;
}
pulled.confidence = confidence;
pulled.recency = recency;
pulled.ventilated = true;
(score, pulled)
})
.collect();
blended.sort_by(|a, b| b.0.total_cmp(&a.0));
Ok(blended.into_iter().map(|(_, pulled)| pulled).collect())
}
fn store_context_sources(
&self,
augmented: &CompileRequest,
out: &CompiledContext,
ttl_seconds: Option<u64>,
) -> Result<(), MemoryError> {
let mut by_hash: BTreeMap<u64, &ContextFragment> = BTreeMap::new();
for fragment in &augmented.fragments {
by_hash
.entry(fragment_handle_hash(fragment))
.or_insert(fragment);
}
let ttl_seconds = positive_ttl(ttl_seconds);
for source in &out.sources {
let Some(hash) = provenance::parse_handle(&source.handle) else {
continue;
};
let Some(fragment) = by_hash.get(&hash) else {
continue;
};
let slot = source_id(hash);
if !self.should_store_source(slot, ttl_seconds)? {
continue;
}
if ttl_seconds.is_none() && self.store.get(slot)?.is_some() {
self.store.delete(slot)?;
}
let content = fragment.content.as_str();
let mut extra: Vec<(&str, Value)> = vec![(CTX_SOURCE_FIELD, Value::Bool(true))];
let embedding = if let Some(media_ref) = &fragment.media {
extra.push((
CTX_SOURCE_MEDIA_FIELD,
serde_json::to_value(media_ref).unwrap_or(Value::Null),
));
self.media_placeholder_embedding(hash)
} else {
self.embedder.embed(content)?
};
self.store_fact(
slot,
content,
&embedding,
Some(&system_meta(&extra)),
ttl_seconds,
)?;
}
Ok(())
}
fn should_store_source(
&self,
slot: u64,
requested_ttl: Option<u64>,
) -> Result<bool, MemoryError> {
match self.context_source_metadata(slot)? {
Some(existing) => Ok(Self::should_upgrade_ttl(&existing, requested_ttl)),
None => Ok(self.store.get(slot)?.is_none()),
}
}
fn should_upgrade_ttl(existing: &Metadata, requested_ttl: Option<u64>) -> bool {
let existing_expiry = existing.get(EXPIRES_AT_FIELD).and_then(Value::as_u64);
match (requested_ttl, existing_expiry) {
(None, Some(_)) => true,
(None | Some(_), None) => false,
(Some(ttl), Some(existing_exp)) => now_unix_secs().saturating_add(ttl) > existing_exp,
}
}
fn media_placeholder_embedding(&self, raw_hash: u64) -> Vec<f32> {
let dim = self.embedder.dimension();
let mut vector = vec![0.0_f32; dim];
let Ok(dim_u64) = u64::try_from(dim) else {
return vector;
};
if dim_u64 == 0 {
return vector;
}
let bucket = usize::try_from(raw_hash % dim_u64).unwrap_or(0);
vector[bucket] = 1.0;
velesdb_core::simd_native::normalize_inplace_native(&mut vector);
vector
}
fn context_source_metadata(&self, slot: u64) -> Result<Option<Metadata>, MemoryError> {
let payloads = self.store.get_metadata_batch(&[slot])?;
Ok(payloads
.into_iter()
.next()
.flatten()
.filter(|meta| meta.get(CTX_SOURCE_FIELD) == Some(&Value::Bool(true))))
}
pub fn retrieve_context_source(&self, handle: &str) -> Result<ContextSource, MemoryError> {
let unknown = || MemoryError::UnknownHandle(handle.to_owned());
let hash = provenance::parse_handle(handle).ok_or_else(unknown)?;
let slot = source_id(hash);
let meta = self.context_source_metadata(slot)?.ok_or_else(unknown)?;
let content = self
.store
.get(slot)?
.map(|(content, _embedding)| content)
.ok_or_else(unknown)?;
Ok(ContextSource {
content,
media: source_media(&meta),
})
}
fn record_context_event(
&self,
request: &CompileRequest,
out: &CompiledContext,
ttl_seconds: Option<u64>,
) -> Result<(), MemoryError> {
let occurred_at_nanos = now_nanos();
let seq = EVENT_SEQ.fetch_add(1, Ordering::Relaxed);
let content = format!("{EVENT_ANCHOR} {occurred_at_nanos}-{seq}");
let id = stable_id(&format!("{EVENT_ID_SALT}{occurred_at_nanos}:{seq}"));
let embedding = self.embedder.embed(&content)?;
let meta = event_meta(request, out, occurred_at_nanos);
self.store_fact(
id,
&content,
&embedding,
Some(&meta),
positive_ttl(ttl_seconds),
)?;
Ok(())
}
pub fn context_savings(&self, project: Option<&str>) -> Result<ContextSavings, MemoryError> {
let mut filter = Map::new();
filter.insert(CTX_EVENT_FIELD.to_owned(), Value::Bool(true));
if let Some(project) = project {
filter.insert(
CTX_PROJECT_FIELD.to_owned(),
Value::String(project.to_owned()),
);
}
let embedding = self.embedder.embed(EVENT_ANCHOR)?;
let hits =
self.store
.query_filtered(&embedding, crate::limits::MAX_RECALL_LIMIT, &filter, 0)?;
let ids: Vec<u64> = hits.iter().map(|(id, _, _)| *id).collect();
let payloads = self.store.get_metadata_batch(&ids)?;
Ok(aggregate_events(&payloads))
}
pub fn save_working_context(
&self,
project: &str,
session: &str,
working: &WorkingContext,
) -> Result<u64, MemoryError> {
let content = serde_json::to_string(working)
.map_err(|err| MemoryError::WorkingContextCodec(err.to_string()))?;
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 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);
}
match self.store.get(slot)? {
Some((content, _)) => serde_json::from_str(&content)
.map(Some)
.map_err(|err| MemoryError::WorkingContextCodec(err.to_string())),
None => Ok(None),
}
}
pub fn list_working_contexts(
&self,
project: &str,
) -> Result<Vec<WorkingContextSession>, MemoryError> {
let mut sessions = self
.working_index(project)?
.map(|index| index.sessions)
.unwrap_or_default();
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(err.to_string())),
None => Ok(None),
}
}
fn update_working_index(&self, project: &str, session: &str) -> Result<(), MemoryError> {
let mut index = self.working_index(project)?.unwrap_or_default();
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,
});
}
let content = serde_json::to_string(&index)
.map_err(|err| MemoryError::WorkingContextCodec(err.to_string()))?;
let slot = working_index_id(project);
let embedding = self
.embedder
.embed(&format!("working context index {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,
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 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;