1use async_trait::async_trait;
8use chrono::{DateTime, Utc};
9use serde::{Deserialize, Serialize};
10use std::sync::Arc;
11
12use super::domain::{FactFingerprint, MemoryObservation, MutationIntent, stable_hash};
13use super::error::MemoryError;
14use super::ids::{EventId, MemoryId, SessionId, TurnId, UserId};
15use super::policy::DiscardReason;
16
17#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
19#[serde(tag = "event", rename_all = "snake_case")]
20pub enum MemoryEvent {
21 FinalTranscriptRecorded {
23 text: String,
25 },
26 ObservationExtracted {
28 observation: Box<MemoryObservation>,
30 },
31 ExtractionFailed {
37 stage: String,
39 reason: String,
41 },
42 ObservationRejected {
44 fingerprint: FactFingerprint,
46 reason: DiscardReason,
48 },
49 SessionCandidateMerged {
51 fingerprint: FactFingerprint,
53 evidence_count: u32,
55 },
56 SessionOverlayUpdated {
58 revision: u64,
60 },
61 ExplicitMutationRequested {
63 intent: MutationIntent,
65 statement: String,
67 },
68 SessionCheckpointed {
70 turns: u64,
72 },
73 SessionSealed {
75 candidate_count: usize,
77 },
78 MutationProposed {
80 fingerprint: FactFingerprint,
82 kind: String,
84 },
85 MutationCommitted {
87 memory_id: MemoryId,
89 kind: String,
91 },
92 MutationRejected {
94 fingerprint: FactFingerprint,
96 reason: DiscardReason,
98 },
99 MemorySuperseded {
101 old: MemoryId,
103 new: MemoryId,
105 },
106 MemoryDeleted {
108 memory_id: MemoryId,
110 },
111 IndexRevisionPublished {
113 revision: u64,
115 },
116}
117
118impl MemoryEvent {
119 pub fn label(&self) -> &'static str {
121 match self {
122 Self::FinalTranscriptRecorded { .. } => "final_transcript_recorded",
123 Self::ObservationExtracted { .. } => "observation_extracted",
124 Self::ExtractionFailed { .. } => "extraction_failed",
125 Self::ObservationRejected { .. } => "observation_rejected",
126 Self::SessionCandidateMerged { .. } => "session_candidate_merged",
127 Self::SessionOverlayUpdated { .. } => "session_overlay_updated",
128 Self::ExplicitMutationRequested { .. } => "explicit_mutation_requested",
129 Self::SessionCheckpointed { .. } => "session_checkpointed",
130 Self::SessionSealed { .. } => "session_sealed",
131 Self::MutationProposed { .. } => "mutation_proposed",
132 Self::MutationCommitted { .. } => "mutation_committed",
133 Self::MutationRejected { .. } => "mutation_rejected",
134 Self::MemorySuperseded { .. } => "memory_superseded",
135 Self::MemoryDeleted { .. } => "memory_deleted",
136 Self::IndexRevisionPublished { .. } => "index_revision_published",
137 }
138 }
139}
140
141pub const EVENT_SCHEMA_VERSION: u32 = 1;
144
145#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
147pub struct MemoryEventEnvelope {
148 pub event_id: EventId,
150 pub occurred_at: DateTime<Utc>,
152 pub user_id: UserId,
154 pub logical_session_id: SessionId,
156 pub turn_id: Option<TurnId>,
158 pub idempotency_key: String,
160 pub schema_version: u32,
162 pub payload: MemoryEvent,
164}
165
166impl MemoryEventEnvelope {
167 pub fn new(
172 user_id: UserId,
173 logical_session_id: SessionId,
174 turn_id: Option<TurnId>,
175 payload: MemoryEvent,
176 now: DateTime<Utc>,
177 ) -> Self {
178 let payload_repr = serde_json::to_string(&payload).unwrap_or_default();
179 let idempotency_key = stable_hash(&format!(
180 "{user_id}|{logical_session_id}|{}|{}",
181 turn_id.map(|t| t.0).unwrap_or_default(),
182 payload_repr
183 ));
184 Self {
185 event_id: EventId::generate(),
186 occurred_at: now,
187 user_id,
188 logical_session_id,
189 turn_id,
190 idempotency_key,
191 schema_version: EVENT_SCHEMA_VERSION,
192 payload,
193 }
194 }
195}
196
197#[async_trait]
203pub trait MemoryEventLog: Send + Sync {
204 async fn append(&self, envelope: MemoryEventEnvelope) -> Result<(), MemoryError>;
206
207 async fn replay_session(
209 &self,
210 session_id: &SessionId,
211 ) -> Result<Vec<MemoryEventEnvelope>, MemoryError>;
212}
213
214#[derive(Debug, Default)]
218pub struct InMemoryEventLog {
219 entries: parking_lot::RwLock<Vec<MemoryEventEnvelope>>,
220 seen: parking_lot::RwLock<std::collections::HashSet<String>>,
221}
222
223impl InMemoryEventLog {
224 pub fn new() -> Self {
226 Self::default()
227 }
228
229 pub fn entries(&self) -> Vec<MemoryEventEnvelope> {
231 self.entries.read().clone()
232 }
233
234 pub fn len(&self) -> usize {
236 self.entries.read().len()
237 }
238
239 pub fn is_empty(&self) -> bool {
241 self.entries.read().is_empty()
242 }
243
244 pub fn count_label(&self, label: &str) -> usize {
246 self.entries
247 .read()
248 .iter()
249 .filter(|e| e.payload.label() == label)
250 .count()
251 }
252}
253
254#[async_trait]
255impl MemoryEventLog for InMemoryEventLog {
256 async fn append(&self, envelope: MemoryEventEnvelope) -> Result<(), MemoryError> {
257 let mut seen = self.seen.write();
258 if !seen.insert(envelope.idempotency_key.clone()) {
259 return Ok(());
260 }
261 drop(seen);
262 self.entries.write().push(envelope);
263 Ok(())
264 }
265
266 async fn replay_session(
267 &self,
268 session_id: &SessionId,
269 ) -> Result<Vec<MemoryEventEnvelope>, MemoryError> {
270 Ok(self
271 .entries
272 .read()
273 .iter()
274 .filter(|e| &e.logical_session_id == session_id)
275 .cloned()
276 .collect())
277 }
278}
279
280#[derive(Clone)]
282pub struct SessionEventWriter {
283 log: Arc<dyn MemoryEventLog>,
284 user_id: UserId,
285 session_id: SessionId,
286}
287
288impl SessionEventWriter {
289 pub fn new(log: Arc<dyn MemoryEventLog>, user_id: UserId, session_id: SessionId) -> Self {
291 Self {
292 log,
293 user_id,
294 session_id,
295 }
296 }
297
298 pub async fn append(
300 &self,
301 turn_id: Option<TurnId>,
302 payload: MemoryEvent,
303 ) -> Result<(), MemoryError> {
304 let envelope = MemoryEventEnvelope::new(
305 self.user_id.clone(),
306 self.session_id.clone(),
307 turn_id,
308 payload,
309 Utc::now(),
310 );
311 self.log.append(envelope).await
312 }
313
314 pub fn user_id(&self) -> &UserId {
316 &self.user_id
317 }
318
319 pub fn session_id(&self) -> &SessionId {
321 &self.session_id
322 }
323}
324
325#[derive(Debug, Clone, Default, PartialEq, Eq)]
327pub struct CommitReceipt {
328 pub revision: u64,
330 pub written: Vec<MemoryId>,
332 pub deleted: Vec<MemoryId>,
334}
335
336#[cfg(test)]
337mod tests {
338 use super::*;
339
340 fn envelope(payload: MemoryEvent) -> MemoryEventEnvelope {
341 MemoryEventEnvelope::new(
342 UserId::new("usr_1"),
343 SessionId::new("ses_1"),
344 Some(TurnId(1)),
345 payload,
346 Utc::now(),
347 )
348 }
349
350 #[tokio::test]
351 async fn appends_are_idempotent_on_replay() {
352 let log = InMemoryEventLog::new();
353 let e = envelope(MemoryEvent::FinalTranscriptRecorded {
354 text: "hello".into(),
355 });
356 log.append(e.clone()).await.unwrap();
357 log.append(e).await.unwrap();
358 assert_eq!(log.len(), 1);
359 }
360
361 #[tokio::test]
362 async fn distinct_payloads_are_distinct_events() {
363 let log = InMemoryEventLog::new();
364 log.append(envelope(MemoryEvent::FinalTranscriptRecorded {
365 text: "a".into(),
366 }))
367 .await
368 .unwrap();
369 log.append(envelope(MemoryEvent::FinalTranscriptRecorded {
370 text: "b".into(),
371 }))
372 .await
373 .unwrap();
374 assert_eq!(log.len(), 2);
375 }
376
377 #[tokio::test]
378 async fn replay_is_scoped_to_one_session() {
379 let log = InMemoryEventLog::new();
380 log.append(envelope(MemoryEvent::SessionSealed { candidate_count: 1 }))
381 .await
382 .unwrap();
383 log.append(MemoryEventEnvelope::new(
384 UserId::new("usr_1"),
385 SessionId::new("ses_other"),
386 None,
387 MemoryEvent::SessionSealed { candidate_count: 2 },
388 Utc::now(),
389 ))
390 .await
391 .unwrap();
392
393 let replayed = log.replay_session(&SessionId::new("ses_1")).await.unwrap();
394 assert_eq!(replayed.len(), 1);
395 }
396}