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::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";
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(());
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 by_hash = index_fragments_by_handle_hash(&augmented.fragments);
let ttl_seconds = positive_ttl(ttl_seconds);
for source in &out.sources {
self.store_one_source(&source.handle, &by_hash, ttl_seconds)?;
}
Ok(())
}
fn store_one_source(
&self,
handle: &str,
by_hash: &BTreeMap<u64, &ContextFragment>,
ttl_seconds: Option<u64>,
) -> Result<(), MemoryError> {
let Some(hash) = provenance::parse_handle(handle) else {
return Ok(());
};
let Some(fragment) = by_hash.get(&hash) else {
return Ok(());
};
let slot = source_id(hash);
if !self.prepare_source_slot(slot, ttl_seconds)? {
return Ok(());
}
let (embedding, media_meta) = self.source_vector(fragment, hash)?;
let mut extra: Vec<(&str, Value)> = vec![(CTX_SOURCE_FIELD, Value::Bool(true))];
if let Some(media) = media_meta {
extra.push((CTX_SOURCE_MEDIA_FIELD, media));
}
self.store_fact(
slot,
fragment.content.as_str(),
&embedding,
Some(&system_meta(&extra)),
ttl_seconds,
)
}
fn prepare_source_slot(
&self,
slot: u64,
ttl_seconds: Option<u64>,
) -> Result<bool, MemoryError> {
if !self.should_store_source(slot, ttl_seconds)? {
return Ok(false);
}
if ttl_seconds.is_none() && self.store.get(slot)?.is_some() {
self.store.delete(slot)?;
}
Ok(true)
}
fn source_vector(
&self,
fragment: &ContextFragment,
hash: u64,
) -> Result<(Vec<f32>, Option<Value>), MemoryError> {
let Some(media_ref) = &fragment.media else {
let embeddable = super::embeddable_prefix(fragment.content.as_str());
return Ok((self.embedder.embed(embeddable)?, None));
};
let descriptor = serde_json::to_value(media_ref).unwrap_or(Value::Null);
Ok((self.media_placeholder_embedding(hash), Some(descriptor)))
}
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),
})
}
pub fn explain_compilation(
&self,
request: &CompileRequest,
fragment_id: u64,
fragment_index: Option<usize>,
) -> Result<ContextDecision, MemoryError> {
if let Some(index) = fragment_index {
let len = request.fragments.len();
if index >= len {
return Err(MemoryError::FragmentIndexOutOfBounds { index, len });
}
}
let mut request = request.clone();
let mut policy = request.policy.take().unwrap_or_default();
policy.record_events = false;
policy.store_sources = false;
policy.slim_response = false;
request.policy = Some(policy);
let compiled =
self.compile_context(&ContextCompiler::new(CompilePolicy::default()), &request)?;
let decision = if let Some(index) = fragment_index {
compiled.decisions.into_iter().nth(index)
} else {
compiled
.decisions
.into_iter()
.find(|decision| decision.fragment_id == fragment_id)
};
decision.ok_or(MemoryError::FragmentNotFound(fragment_id))
}
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> {
if working.is_empty() {
return Err(MemoryError::EmptyWorkingContext);
}
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);
}
let Some((content, _)) = self.store.get(slot)? else {
return Err(MemoryError::WorkingContextCodec(format!(
"working context for project '{project}', session '{session}' is corrupt: \
the reserved marker is present but the stored body is gone"
)));
};
serde_json::from_str(&content)
.map(Some)
.map_err(|err| MemoryError::WorkingContextCodec(err.to_string()))
}
pub fn resume_working_context(
&self,
project: &str,
session: &str,
) -> Result<LoadedWorkingContext, MemoryError> {
let working = self.load_working_context(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(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(format!(
"storage returned {} metadata rows for {} working-context ids",
payloads.len(),
ids.len()
)));
}
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 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(err.to_string())),
None => Err(MemoryError::WorkingContextCodec(format!(
"working-context index for project '{project}' is corrupt: the index \
marker is present but the stored body is gone"
))),
}
}
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(err.to_string()))?;
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;