1use chrono::Utc;
9use std::sync::Arc;
10
11use super::consolidate::ConsolidationOutput;
12use super::proposal::{MemorySelector, ResolutionKind, ResolvedMutation};
13use super::resolver::Resolver;
14use crate::core::{
15 CommitReceipt, MemoryError, MemoryEvent, MemoryStatus, SessionEventWriter, UserId, stable_hash,
16};
17use crate::okf::{MemoryRepository, MemoryTransaction, ReconciliationSelector};
18
19pub struct MemoryCommitter {
21 repository: Arc<dyn MemoryRepository>,
22 events: Option<SessionEventWriter>,
23}
24
25#[derive(Debug, Clone, Default, PartialEq)]
27pub struct ReconciliationReport {
28 pub creates: usize,
30 pub reinforces: usize,
32 pub refines: usize,
34 pub supersedes: usize,
36 pub stages: usize,
38 pub deletes: usize,
40 pub discards: usize,
42 pub revision: u64,
44}
45
46impl ReconciliationReport {
47 fn record(&mut self, kind: ResolutionKind) {
48 match kind {
49 ResolutionKind::Create => self.creates += 1,
50 ResolutionKind::Reinforce => self.reinforces += 1,
51 ResolutionKind::Refine => self.refines += 1,
52 ResolutionKind::Supersede => self.supersedes += 1,
53 ResolutionKind::Coexist => self.creates += 1,
54 ResolutionKind::Stage => self.stages += 1,
55 ResolutionKind::Delete => self.deletes += 1,
56 ResolutionKind::Discard => self.discards += 1,
57 }
58 }
59
60 pub fn is_empty(&self) -> bool {
62 self.creates + self.reinforces + self.refines + self.supersedes + self.stages + self.deletes
63 == 0
64 }
65}
66
67impl MemoryCommitter {
68 pub fn new(repository: Arc<dyn MemoryRepository>) -> Self {
70 Self {
71 repository,
72 events: None,
73 }
74 }
75
76 pub fn with_events(mut self, events: SessionEventWriter) -> Self {
78 self.events = Some(events);
79 self
80 }
81
82 pub async fn reconcile(
84 &self,
85 owner: &UserId,
86 output: ConsolidationOutput,
87 idempotency_key: &str,
88 ) -> Result<ReconciliationReport, MemoryError> {
89 let now = Utc::now();
90 let resolver = Resolver::new(owner.clone());
91 let mut report = ReconciliationReport::default();
92 let mut transaction = MemoryTransaction::new(owner.clone(), idempotency_key.to_string());
93
94 for proposal in output.proposals {
95 let mut window = self
99 .repository
100 .find_candidates(
101 owner,
102 &ReconciliationSelector::by_subject_predicate(
103 proposal.fingerprint.subject_predicate().to_string(),
104 ),
105 )
106 .await?;
107
108 let wider = self
115 .repository
116 .find_candidates(
117 owner,
118 &ReconciliationSelector::by_subject(proposal.fingerprint.subject().to_string()),
119 )
120 .await?;
121 for candidate in wider {
122 if !window.iter().any(|m| m.id == candidate.id) {
123 window.push(candidate);
124 }
125 }
126
127 let resolved = resolver.resolve(proposal, &window, now);
128 report.record(resolved.kind);
129 self.emit(&resolved).await;
130 transaction = apply(transaction, resolved);
131 }
132
133 for selector in &output.deletions {
134 let targets = self.resolve_deletion(owner, selector).await?;
135 for id in targets {
136 report.deletes += 1;
137 if let Some(events) = &self.events {
138 let _ = events
139 .append(
140 None,
141 MemoryEvent::MemoryDeleted {
142 memory_id: id.clone(),
143 },
144 )
145 .await;
146 }
147 transaction = transaction.delete(id);
148 }
149 }
150
151 if transaction.is_empty() {
152 return Ok(report);
153 }
154
155 let receipt = self.repository.commit(transaction).await?;
156 report.revision = receipt.revision;
157 Ok(report)
158 }
159
160 pub async fn commit_one(
162 &self,
163 owner: &UserId,
164 resolved: ResolvedMutation,
165 ) -> Result<CommitReceipt, MemoryError> {
166 let key = stable_hash(&format!("{}|{:?}", resolved.fingerprint, resolved.kind));
167 self.emit(&resolved).await;
168 let transaction = apply(MemoryTransaction::new(owner.clone(), key), resolved);
169 self.repository.commit(transaction).await
170 }
171
172 async fn resolve_deletion(
174 &self,
175 owner: &UserId,
176 selector: &MemorySelector,
177 ) -> Result<Vec<crate::core::MemoryId>, MemoryError> {
178 let all = self.repository.all(owner).await?;
181 Ok(all
182 .iter()
183 .filter(|m| selector.matches(m))
184 .map(|m| m.id.clone())
185 .collect())
186 }
187
188 async fn emit(&self, resolved: &ResolvedMutation) {
189 let Some(events) = &self.events else { return };
190 let event = match resolved.kind {
191 ResolutionKind::Discard => MemoryEvent::MutationRejected {
192 fingerprint: resolved.fingerprint.clone(),
193 reason: resolved
194 .discard_reason
195 .unwrap_or(crate::core::DiscardReason::InsufficientEvidence),
196 },
197 kind => match resolved.writes.first() {
198 Some(memory) => MemoryEvent::MutationCommitted {
199 memory_id: memory.id.clone(),
200 kind: kind.label().to_string(),
201 },
202 None => MemoryEvent::MutationProposed {
203 fingerprint: resolved.fingerprint.clone(),
204 kind: kind.label().to_string(),
205 },
206 },
207 };
208 let _ = events.append(None, event).await;
209
210 if (resolved.kind == ResolutionKind::Supersede || resolved.kind == ResolutionKind::Refine)
211 && let (Some(new), Some(old)) = (resolved.writes.first(), resolved.writes.get(1))
212 {
213 let _ = events
214 .append(
215 None,
216 MemoryEvent::MemorySuperseded {
217 old: old.id.clone(),
218 new: new.id.clone(),
219 },
220 )
221 .await;
222 }
223 }
224}
225
226fn apply(mut transaction: MemoryTransaction, resolved: ResolvedMutation) -> MemoryTransaction {
227 for memory in resolved.writes {
228 transaction = transaction.put(memory);
229 }
230 for id in resolved.deletes {
231 transaction = transaction.delete(id);
232 }
233 transaction
234}
235
236pub async fn commit_promotions(
238 repository: &Arc<dyn MemoryRepository>,
239 owner: &UserId,
240 outcomes: Vec<super::promotion::PromotionOutcome>,
241 idempotency_key: &str,
242) -> Result<CommitReceipt, MemoryError> {
243 use super::promotion::PromotionOutcome;
244
245 let mut transaction = MemoryTransaction::new(owner.clone(), idempotency_key.to_string());
246 for outcome in outcomes {
247 match outcome {
248 PromotionOutcome::Promote(memory) => transaction = transaction.put(*memory),
249 PromotionOutcome::Expire(memory) => {
250 let mut expired = *memory;
251 expired.status = MemoryStatus::Expired;
252 expired.temporal.updated_at = Utc::now();
253 transaction = transaction.put(expired);
254 }
255 PromotionOutcome::Hold { .. } => {}
256 }
257 }
258 repository.commit(transaction).await
259}
260
261#[cfg(test)]
262mod tests {
263 use super::*;
264 use crate::core::{
265 Explicitness, InMemoryEventLog, IngestionConfig, MemoryEventLog, MemoryKind,
266 MemoryObservation, MemoryValue, ObservationId, ProposedPersistence, SensitivityClass,
267 SessionId, SpeakerAttribution, TemporalScope, TranscriptEvidence, TurnId,
268 };
269 use crate::ingestion::{InMemorySessionLedger, SessionLedger};
270 use crate::okf::OkfRepository;
271 use crate::reconcile::consolidate::consolidate;
272
273 fn observation(
274 predicate: &str,
275 value: &str,
276 turn: u64,
277 intent: Option<crate::core::MutationIntent>,
278 ) -> MemoryObservation {
279 MemoryObservation {
280 observation_id: ObservationId::generate(),
281 session_id: SessionId::new("ses_1"),
282 turn_id: TurnId(turn),
283 subject: crate::core::EntityRef::user(),
284 predicate: crate::core::CanonicalPredicate::new(predicate),
285 value: MemoryValue::Text(value.to_string()),
286 canonical_statement: format!("The user is {value}."),
287 kind: MemoryKind::Preference,
288 explicitness: Explicitness::ExplicitStatement,
289 confidence: 0.9,
290 persistence: ProposedPersistence::Durable,
291 temporal_scope: TemporalScope::Persistent,
292 valid_from: None,
293 expected_expiry: None,
294 transcript_evidence: TranscriptEvidence::new(format!("I am {value}")),
295 speaker_attribution: SpeakerAttribution::User,
296 sensitivity: SensitivityClass::Normal,
297 mutation_intent: intent,
298 search_terms: Vec::new(),
299 }
300 }
301
302 async fn run_session(
303 committer: &MemoryCommitter,
304 owner: &UserId,
305 session: &str,
306 observations: Vec<MemoryObservation>,
307 ) -> ReconciliationReport {
308 let ledger =
309 InMemorySessionLedger::new(SessionId::new(session), IngestionConfig::default());
310 for obs in observations {
311 ledger.append_observation(obs).await.unwrap();
312 }
313 ledger.micro_reconcile();
314 let sealed = ledger.seal().await.unwrap();
315 committer
316 .reconcile(owner, consolidate(&sealed), session)
317 .await
318 .unwrap()
319 }
320
321 fn setup() -> (Arc<dyn MemoryRepository>, MemoryCommitter, UserId) {
322 let repository: Arc<dyn MemoryRepository> = Arc::new(OkfRepository::in_memory());
323 let committer = MemoryCommitter::new(repository.clone());
324 (repository, committer, UserId::new("usr_1"))
325 }
326
327 #[tokio::test]
328 async fn a_first_session_creates_the_record() {
329 let (repository, committer, owner) = setup();
330 let report = run_session(
331 &committer,
332 &owner,
333 "ses_1",
334 vec![observation("dietary_identity", "pescatarian", 1, None)],
335 )
336 .await;
337
338 assert_eq!(report.creates, 1);
339 let stored = repository.all(&owner).await.unwrap();
340 assert_eq!(stored.len(), 1);
341 assert_eq!(stored[0].statement, "The user is pescatarian.");
342 }
343
344 #[tokio::test]
345 async fn the_same_fact_in_a_later_session_reinforces_rather_than_duplicating() {
346 let (repository, committer, owner) = setup();
347 run_session(
348 &committer,
349 &owner,
350 "ses_1",
351 vec![observation("dietary_identity", "pescatarian", 1, None)],
352 )
353 .await;
354 let report = run_session(
355 &committer,
356 &owner,
357 "ses_2",
358 vec![observation("dietary_identity", "pescatarian", 1, None)],
359 )
360 .await;
361
362 assert_eq!(report.reinforces, 1);
363 assert_eq!(report.creates, 0);
364 let stored = repository.all(&owner).await.unwrap();
365 assert_eq!(stored.len(), 1, "no duplicate record");
366 assert!(stored[0].evidence.count >= 2);
367 assert_eq!(stored[0].evidence.distinct_sessions, 2);
368 }
369
370 #[tokio::test]
371 async fn a_correction_in_a_later_session_supersedes_the_old_record() {
372 let (repository, committer, owner) = setup();
373 run_session(
374 &committer,
375 &owner,
376 "ses_1",
377 vec![observation("dietary_identity", "vegetarian", 1, None)],
378 )
379 .await;
380 let report = run_session(
381 &committer,
382 &owner,
383 "ses_2",
384 vec![observation("dietary_identity", "pescatarian", 1, None)],
385 )
386 .await;
387
388 assert_eq!(report.supersedes, 1);
389 let stored = repository.all(&owner).await.unwrap();
390 let active: Vec<_> = stored
391 .iter()
392 .filter(|m| m.status == MemoryStatus::Active)
393 .collect();
394 assert_eq!(active.len(), 1);
395 assert_eq!(active[0].statement, "The user is pescatarian.");
396
397 let retired: Vec<_> = stored
398 .iter()
399 .filter(|m| m.status == MemoryStatus::Superseded)
400 .collect();
401 assert_eq!(retired.len(), 1);
402 assert_eq!(retired[0].superseded_by.as_ref(), Some(&active[0].id));
403 }
404
405 #[tokio::test]
406 async fn a_forget_command_removes_the_record_and_its_superseded_copies() {
407 let (repository, committer, owner) = setup();
408 run_session(
409 &committer,
410 &owner,
411 "ses_1",
412 vec![observation("dietary_identity", "vegetarian", 1, None)],
413 )
414 .await;
415 run_session(
416 &committer,
417 &owner,
418 "ses_2",
419 vec![observation("dietary_identity", "pescatarian", 1, None)],
420 )
421 .await;
422 assert_eq!(repository.all(&owner).await.unwrap().len(), 2);
423
424 let report = run_session(
425 &committer,
426 &owner,
427 "ses_3",
428 vec![observation(
429 "memory_removal",
430 "pescatarian",
431 1,
432 Some(crate::core::MutationIntent::Forget),
433 )],
434 )
435 .await;
436
437 assert_eq!(report.deletes, 1, "the superseded copy mentions vegetarian");
438 let remaining = repository.all(&owner).await.unwrap();
439 assert!(
440 remaining
441 .iter()
442 .all(|m| !m.statement.contains("pescatarian")),
443 "deleted content is gone: {remaining:?}"
444 );
445 }
446
447 #[tokio::test]
448 async fn a_retried_reconciliation_does_not_double_write() {
449 let (repository, committer, owner) = setup();
450 let ledger =
451 InMemorySessionLedger::new(SessionId::new("ses_1"), IngestionConfig::default());
452 ledger
453 .append_observation(observation("dietary_identity", "pescatarian", 1, None))
454 .await
455 .unwrap();
456 let sealed = ledger.seal().await.unwrap();
457
458 committer
459 .reconcile(&owner, consolidate(&sealed), "same-key")
460 .await
461 .unwrap();
462 committer
463 .reconcile(&owner, consolidate(&sealed), "same-key")
464 .await
465 .unwrap();
466
467 assert_eq!(repository.all(&owner).await.unwrap().len(), 1);
468 }
469
470 #[tokio::test]
471 async fn a_session_that_proposed_nothing_commits_nothing() {
472 let (repository, committer, owner) = setup();
473 let report = run_session(&committer, &owner, "ses_1", Vec::new()).await;
474 assert!(report.is_empty());
475 assert_eq!(repository.revision(&owner).await.unwrap(), 0);
476 }
477
478 #[tokio::test]
479 async fn commits_are_audited() {
480 let log = Arc::new(InMemoryEventLog::new());
481 let repository: Arc<dyn MemoryRepository> = Arc::new(OkfRepository::in_memory());
482 let owner = UserId::new("usr_1");
483 let committer = MemoryCommitter::new(repository).with_events(SessionEventWriter::new(
484 log.clone() as Arc<dyn MemoryEventLog>,
485 owner.clone(),
486 SessionId::new("ses_1"),
487 ));
488
489 run_session(
490 &committer,
491 &owner,
492 "ses_1",
493 vec![observation("dietary_identity", "vegetarian", 1, None)],
494 )
495 .await;
496 run_session(
497 &committer,
498 &owner,
499 "ses_2",
500 vec![observation("dietary_identity", "pescatarian", 1, None)],
501 )
502 .await;
503
504 assert!(log.count_label("mutation_committed") >= 2);
505 assert_eq!(log.count_label("memory_superseded"), 1);
506 }
507}