1use std::collections::BTreeMap;
18use std::sync::atomic::{AtomicU64, Ordering};
19use std::time::{SystemTime, UNIX_EPOCH};
20
21use serde_json::{Map, Number, Value};
22
23use super::{positive_ttl, MemoryService, Metadata, HUB_FIELD};
24use crate::context::model::{
25 CompileRequest, CompiledContext, ContextFragment, ContextSavings, MemoryScope, WorkingContext,
26};
27use crate::context::{provenance, ContextCompiler};
28use crate::embedder::Embedder;
29use crate::error::MemoryError;
30use crate::id::stable_id;
31use crate::model::FusionOptions;
32use crate::storage::MemoryStore;
33
34const SOURCE_ID_SALT: &str = "veles-ctx-source:";
38const EVENT_ID_SALT: &str = "veles-ctx-event:";
40const WORKING_ID_SALT: &str = "veles-ctx-working:";
43
44const EVENT_ANCHOR: &str = "veles context compilation event";
47
48const CTX_EVENT_FIELD: &str = "_veles_ctx_event";
54const CTX_PROJECT_FIELD: &str = "_veles_ctx_project";
55const CTX_MODEL_FIELD: &str = "_veles_ctx_model";
56const CTX_SOURCE_FIELD: &str = "_veles_ctx_source";
57const CTX_WORKING_FIELD: &str = "_veles_ctx_working";
58const CTX_SESSION_FIELD: &str = "_veles_ctx_session";
59const CTX_TOKENS_IN_FIELD: &str = "_veles_ctx_tokens_in";
60const CTX_TOKENS_OUT_FIELD: &str = "_veles_ctx_tokens_out";
61const CTX_TOKENS_SAVED_FIELD: &str = "_veles_ctx_tokens_saved";
62const CTX_COST_FIELD: &str = "_veles_ctx_cost_micros";
63const CTX_CURRENCY_FIELD: &str = "_veles_ctx_currency";
64const CTX_AT_FIELD: &str = "_veles_ctx_at";
65
66static EVENT_SEQ: AtomicU64 = AtomicU64::new(0);
69
70impl<E: Embedder, S: MemoryStore> MemoryService<E, S> {
71 pub fn compile_context(
84 &self,
85 compiler: &ContextCompiler,
86 request: &CompileRequest,
87 ) -> Result<CompiledContext, MemoryError> {
88 let memories = self.context_memories(request)?;
89 self.compile_with_memories(compiler, request, memories)
90 }
91
92 pub fn compile_context_reranked<R: crate::Reranker>(
107 &self,
108 compiler: &ContextCompiler,
109 request: &CompileRequest,
110 reranker: &R,
111 ) -> Result<CompiledContext, MemoryError> {
112 let memories = self.context_memories_reranked(request, reranker)?;
113 self.compile_with_memories(compiler, request, memories)
114 }
115
116 fn compile_with_memories(
120 &self,
121 compiler: &ContextCompiler,
122 request: &CompileRequest,
123 memories: Vec<PulledMemory>,
124 ) -> Result<CompiledContext, MemoryError> {
125 let mut augmented = request.clone();
126 let mut pulled: BTreeMap<u64, PulledMemory> = BTreeMap::new();
127 for memory in memories {
128 augmented.fragments.push(memory.fragment.clone());
129 pulled.insert(stable_id(&memory.fragment.content), memory);
130 }
131 let mut out = compiler.compile(&augmented)?;
132 annotate_memory_provenance(&mut out, &pulled);
133 let policy = compiler.effective_policy(request);
134 if policy.store_sources {
135 self.store_context_sources(&augmented, &out, policy.source_ttl_seconds)?;
136 }
137 if policy.record_events {
138 self.record_context_event(request, &out, policy.event_ttl_seconds)?;
139 }
140 Ok(out)
141 }
142
143 fn context_memories(&self, request: &CompileRequest) -> Result<Vec<PulledMemory>, MemoryError> {
146 let Some((scope, k)) = scope_and_k(request) else {
147 return Ok(Vec::new());
148 };
149 let filter = scope_filter(scope);
150 let opts = FusionOptions::from_knobs(scope.hops, scope.graph_boost, None);
154 let scored = self.recall_fused_scored(&request.query, k, filter.as_ref(), opts)?;
155 let max_fused = scored
156 .iter()
157 .map(|s| s.fused)
158 .fold(f64::MIN, f64::max)
159 .max(f64::EPSILON);
160 Ok(scored
161 .into_iter()
162 .map(|scored| {
163 let memory_id = scored.recollection.id;
164 let fused = if scored.fused.is_finite() {
169 scored.fused
170 } else {
171 0.0
172 };
173 #[allow(clippy::cast_possible_truncation)] let relevance = (fused / max_fused).clamp(0.0, 1.0) as f32;
175 let fragment = ContextFragment {
176 id: None,
177 content: scored.recollection.content,
178 kind: Some("memory".to_owned()),
179 priority: None,
180 metadata: None,
181 };
182 PulledMemory {
183 fragment,
184 memory_id,
185 relevance,
186 vector_norm: scored.vector_norm,
187 graph_weight: scored.graph_weight,
188 }
189 })
190 .collect())
191 }
192
193 fn context_memories_reranked<R: crate::Reranker>(
199 &self,
200 request: &CompileRequest,
201 reranker: &R,
202 ) -> Result<Vec<PulledMemory>, MemoryError> {
203 let Some((scope, k)) = scope_and_k(request) else {
204 return Ok(Vec::new());
205 };
206 let filter = scope_filter(scope);
207 let opts = FusionOptions::from_knobs(scope.hops, scope.graph_boost, None);
208 let ranked =
209 self.recall_fused_reranked(&request.query, k, filter.as_ref(), opts, reranker)?;
210 let count = ranked.len().max(1);
211 Ok(ranked
212 .into_iter()
213 .enumerate()
214 .map(|(rank, recollection)| {
215 #[allow(clippy::cast_precision_loss)] let relevance = 1.0 - (rank as f32 / count as f32);
217 PulledMemory {
218 fragment: ContextFragment {
219 id: None,
220 content: recollection.content,
221 kind: Some("memory".to_owned()),
222 priority: None,
223 metadata: None,
224 },
225 memory_id: recollection.id,
226 relevance,
227 vector_norm: 0.0,
228 graph_weight: 0.0,
229 }
230 })
231 .collect())
232 }
233
234 fn store_context_sources(
237 &self,
238 augmented: &CompileRequest,
239 out: &CompiledContext,
240 ttl_seconds: Option<u64>,
241 ) -> Result<(), MemoryError> {
242 let by_hash: BTreeMap<u64, &str> = augmented
243 .fragments
244 .iter()
245 .map(|fragment| (stable_id(&fragment.content), fragment.content.as_str()))
246 .collect();
247 let ttl_seconds = positive_ttl(ttl_seconds);
248 for source in &out.sources {
249 let Some(hash) = provenance::parse_handle(&source.handle) else {
250 continue;
251 };
252 let Some(content) = by_hash.get(&hash) else {
253 continue;
254 };
255 let slot = source_id(hash);
256 if self.store.get(slot)?.is_some() {
263 continue;
264 }
265 let embedding = self.embedder.embed(content)?;
266 self.store_fact(
267 slot,
268 content,
269 &embedding,
270 Some(&system_meta(&[(CTX_SOURCE_FIELD, Value::Bool(true))])),
271 ttl_seconds,
272 )?;
273 }
274 Ok(())
275 }
276
277 fn slot_is_context_source(&self, slot: u64) -> Result<bool, MemoryError> {
279 let payloads = self.store.get_metadata_batch(&[slot])?;
280 Ok(payloads.first().is_some_and(|payload| {
281 payload
282 .as_ref()
283 .is_some_and(|meta| meta.get(CTX_SOURCE_FIELD) == Some(&Value::Bool(true)))
284 }))
285 }
286
287 pub fn retrieve_context_source(&self, handle: &str) -> Result<String, MemoryError> {
293 let unknown = || MemoryError::UnknownHandle(handle.to_owned());
294 let hash = provenance::parse_handle(handle).ok_or_else(unknown)?;
295 let slot = source_id(hash);
296 if !self.slot_is_context_source(slot)? {
299 return Err(unknown());
300 }
301 self.store
302 .get(slot)?
303 .map(|(content, _)| content)
304 .ok_or_else(unknown)
305 }
306
307 fn record_context_event(
311 &self,
312 request: &CompileRequest,
313 out: &CompiledContext,
314 ttl_seconds: Option<u64>,
315 ) -> Result<(), MemoryError> {
316 let occurred_at_nanos = SystemTime::now()
317 .duration_since(UNIX_EPOCH)
318 .map(|elapsed| elapsed.as_nanos())
319 .unwrap_or(0);
320 let seq = EVENT_SEQ.fetch_add(1, Ordering::Relaxed);
323 let content = format!("{EVENT_ANCHOR} {occurred_at_nanos}-{seq}");
324 let id = stable_id(&format!("{EVENT_ID_SALT}{occurred_at_nanos}:{seq}"));
325 let embedding = self.embedder.embed(&content)?;
326 let meta = event_meta(request, out, occurred_at_nanos);
327 self.store_fact(
328 id,
329 &content,
330 &embedding,
331 Some(&meta),
332 positive_ttl(ttl_seconds),
333 )?;
334 Ok(())
335 }
336
337 pub fn context_savings(&self, project: Option<&str>) -> Result<ContextSavings, MemoryError> {
346 let mut filter = Map::new();
351 filter.insert(CTX_EVENT_FIELD.to_owned(), Value::Bool(true));
352 if let Some(project) = project {
353 filter.insert(
354 CTX_PROJECT_FIELD.to_owned(),
355 Value::String(project.to_owned()),
356 );
357 }
358 let embedding = self.embedder.embed(EVENT_ANCHOR)?;
359 let hits =
360 self.store
361 .query_filtered(&embedding, crate::limits::MAX_RECALL_LIMIT, &filter, 0)?;
362 let ids: Vec<u64> = hits.iter().map(|(id, _, _)| *id).collect();
363 let payloads = self.store.get_metadata_batch(&ids)?;
364 Ok(aggregate_events(&payloads))
365 }
366
367 pub fn save_working_context(
374 &self,
375 project: &str,
376 session: &str,
377 working: &WorkingContext,
378 ) -> Result<u64, MemoryError> {
379 let content = serde_json::to_string(working)
380 .map_err(|err| MemoryError::WorkingContextCodec(err.to_string()))?;
381 let id = working_id(project, session);
382 let embedding = self
383 .embedder
384 .embed(&format!("working context {project} {session}"))?;
385 let meta = system_meta(&[
386 (CTX_WORKING_FIELD, Value::Bool(true)),
387 (CTX_PROJECT_FIELD, Value::String(project.to_owned())),
388 (CTX_SESSION_FIELD, Value::String(session.to_owned())),
389 ]);
390 self.store_fact(id, &content, &embedding, Some(&meta), None)?;
391 Ok(id)
392 }
393
394 pub fn load_working_context(
401 &self,
402 project: &str,
403 session: &str,
404 ) -> Result<Option<WorkingContext>, MemoryError> {
405 match self.store.get(working_id(project, session))? {
406 Some((content, _)) => serde_json::from_str(&content)
407 .map(Some)
408 .map_err(|err| MemoryError::WorkingContextCodec(err.to_string())),
409 None => Ok(None),
410 }
411 }
412}
413
414const DEFAULT_MEMORY_K: usize = 5;
416
417fn scope_and_k(request: &CompileRequest) -> Option<(&MemoryScope, usize)> {
423 let scope = request.memory_scope.as_ref()?;
424 let room = crate::limits::MAX_FRAGMENTS.saturating_sub(request.fragments.len());
425 let k = crate::limits::clamp_recall_limit(scope.k.unwrap_or(DEFAULT_MEMORY_K)).min(room);
426 (k > 0).then_some((scope, k))
427}
428
429fn scope_filter(scope: &MemoryScope) -> Option<Metadata> {
431 scope.project.as_ref().map(|project| {
432 let mut meta = Map::new();
433 meta.insert("project".to_owned(), Value::String(project.clone()));
434 meta
435 })
436}
437
438struct PulledMemory {
440 fragment: ContextFragment,
441 memory_id: u64,
442 relevance: f32,
444 vector_norm: f64,
446 graph_weight: f64,
448}
449
450fn annotate_memory_provenance(out: &mut CompiledContext, pulled: &BTreeMap<u64, PulledMemory>) {
456 for decision in &mut out.decisions {
457 if let Some(memory) = pulled.get(&decision.content_hash) {
458 decision.memory_id = Some(memory.memory_id);
459 decision.relevance = memory.relevance;
460 decision.reason = format!(
461 "{} — pulled from memory {} (vector {:.2}, graph {:.2})",
462 decision.reason, memory.memory_id, memory.vector_norm, memory.graph_weight
463 );
464 }
465 }
466 for source in &mut out.sources {
467 if let Some(hash) = provenance::parse_handle(&source.handle) {
468 if let Some(memory) = pulled.get(&hash) {
469 source.memory_id = Some(memory.memory_id);
470 }
471 }
472 }
473}
474
475fn system_meta(extra: &[(&str, Value)]) -> Metadata {
478 let mut meta = Map::new();
479 meta.insert(HUB_FIELD.to_owned(), Value::Bool(true));
480 for (key, value) in extra {
481 meta.insert((*key).to_owned(), value.clone());
482 }
483 meta
484}
485
486fn event_meta(request: &CompileRequest, out: &CompiledContext, nanos: u128) -> Metadata {
489 let mut extra: Vec<(&str, Value)> = vec![
490 (CTX_EVENT_FIELD, Value::Bool(true)),
491 (
492 CTX_TOKENS_IN_FIELD,
493 Value::Number(out.insights.tokens_in.into()),
494 ),
495 (
496 CTX_TOKENS_OUT_FIELD,
497 Value::Number(out.insights.tokens_out.into()),
498 ),
499 (
500 CTX_TOKENS_SAVED_FIELD,
501 Value::Number(out.insights.tokens_saved.into()),
502 ),
503 (
504 CTX_AT_FIELD,
505 Value::Number(Number::from(
506 u64::try_from(nanos / 1_000_000_000).unwrap_or(u64::MAX),
507 )),
508 ),
509 ];
510 if let Some(project) = &request.project {
511 extra.push((CTX_PROJECT_FIELD, Value::String(project.clone())));
512 }
513 if let Some(model) = &request.target_model {
514 extra.push((CTX_MODEL_FIELD, Value::String(model.clone())));
515 }
516 if let (Some(micros), Some(currency)) = (
517 out.insights.estimated_cost_saved_micros,
518 out.insights.currency.as_ref(),
519 ) {
520 extra.push((CTX_COST_FIELD, Value::Number(micros.into())));
521 extra.push((CTX_CURRENCY_FIELD, Value::String(currency.clone())));
522 }
523 system_meta(&extra)
524}
525
526fn aggregate_events(payloads: &[Option<Metadata>]) -> ContextSavings {
530 let mut savings = ContextSavings {
531 events: payloads.len() as u64,
532 truncated: payloads.len() >= crate::limits::MAX_RECALL_LIMIT,
533 ..ContextSavings::default()
534 };
535 for payload in payloads {
536 let Some(meta) = payload else { continue };
537 savings.tokens_in = savings
538 .tokens_in
539 .saturating_add(meta_u64(meta, CTX_TOKENS_IN_FIELD));
540 savings.tokens_out = savings
541 .tokens_out
542 .saturating_add(meta_u64(meta, CTX_TOKENS_OUT_FIELD));
543 savings.tokens_saved = savings
544 .tokens_saved
545 .saturating_add(meta_u64(meta, CTX_TOKENS_SAVED_FIELD));
546 if let (Some(Value::String(currency)), micros) =
547 (meta.get(CTX_CURRENCY_FIELD), meta_u64(meta, CTX_COST_FIELD))
548 {
549 if micros > 0 {
550 let entry = savings
551 .cost_saved_micros_by_currency
552 .entry(currency.clone())
553 .or_insert(0);
554 *entry = entry.saturating_add(micros);
555 }
556 }
557 }
558 savings
559}
560
561fn meta_u64(meta: &Metadata, key: &str) -> u64 {
563 meta.get(key).and_then(Value::as_u64).unwrap_or(0)
564}
565
566fn source_id(content_hash: u64) -> u64 {
568 stable_id(&format!("{SOURCE_ID_SALT}{content_hash}"))
569}
570
571fn working_id(project: &str, session: &str) -> u64 {
573 stable_id(&format!("{WORKING_ID_SALT}{project}\u{1f}{session}"))
574}