gemini_memory_rs/reconcile/
commit.rs

1//! Committing resolved mutations to canonical memory.
2//!
3//! One transaction per session. A contradiction resolution writes both the new
4//! active record and the retired one; committing only half would leave the
5//! corpus asserting two incompatible facts, so the unit of commit is the whole
6//! reconciliation, not the individual record.
7
8use 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
19/// Turns proposals into durable memory.
20pub struct MemoryCommitter {
21    repository: Arc<dyn MemoryRepository>,
22    events: Option<SessionEventWriter>,
23}
24
25/// What one reconciliation did.
26#[derive(Debug, Clone, Default, PartialEq)]
27pub struct ReconciliationReport {
28    /// Count of each resolution kind, for metrics.
29    pub creates: usize,
30    /// Records reinforced.
31    pub reinforces: usize,
32    /// Records refined.
33    pub refines: usize,
34    /// Records superseded.
35    pub supersedes: usize,
36    /// Records staged.
37    pub stages: usize,
38    /// Records deleted.
39    pub deletes: usize,
40    /// Proposals refused.
41    pub discards: usize,
42    /// Repository revision after the commit.
43    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    /// Whether anything at all changed.
61    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    /// A committer writing to `repository`.
69    pub fn new(repository: Arc<dyn MemoryRepository>) -> Self {
70        Self {
71            repository,
72            events: None,
73        }
74    }
75
76    /// Emit audit events alongside every commit.
77    pub fn with_events(mut self, events: SessionEventWriter) -> Self {
78        self.events = Some(events);
79        self
80    }
81
82    /// Resolve and commit a session's consolidation output.
83    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            // The candidate window is the subject-and-predicate neighbourhood,
96            // not the whole corpus: reconciliation compares a proposal against
97            // what could plausibly be the same fact, and nothing else.
98            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            // ยง22.1 step 3: widen before concluding a proposal is novel.
109            // Restricting the window to an exact predicate match means a
110            // renamed predicate looks like a brand-new fact, and the corpus
111            // grows a duplicate every session. Widening by subject is cheaper
112            // and more exact than a lexical search at this corpus size, and
113            // the resolver does the discrimination.
114            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    /// Commit a single resolved mutation, for explicit in-session commands.
161    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    /// Expand a deletion selector into concrete record ids.
173    async fn resolve_deletion(
174        &self,
175        owner: &UserId,
176        selector: &MemorySelector,
177    ) -> Result<Vec<crate::core::MemoryId>, MemoryError> {
178        // Deletion reaches every status, not only active records: a user who
179        // asks to forget something means the superseded copies too.
180        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
236/// Apply a promotion sweep's outcomes to the repository.
237pub 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}