Skip to main content

turnframe_store/memory/
impls.rs

1//! The trait implementations of [`MemoryStores`].
2//!
3//! Each method does the same three things: check whether a failure is armed at
4//! the boundary it crosses, take the state lock for exactly one synchronous
5//! rule call from [`super::state`], and release it. No `.await` happens while
6//! the lock is held — there is no `.await` in this file at all — which is what
7//! keeps the returned futures `Send`.
8
9use async_trait::async_trait;
10use chrono::{DateTime, Utc};
11use turnframe_core::case::CaseKey;
12use turnframe_core::event::{EventRedaction, OutboxEntry};
13use turnframe_core::ids::{
14    AccountId, CaseRevision, CommandId, ConversationId, EventId, InteractionId, OptionId, OutboxId,
15    RedactionAuthority, TurnId,
16};
17use turnframe_core::interaction::{Interaction, InteractionStatus};
18use turnframe_core::replay::{ReplayRecord, TurnPhase};
19use turnframe_core::response::AssistantTurn;
20
21use super::MemoryStores;
22use super::fault::FailurePoint;
23use crate::commit::{CommitBundle, CommitReceipt, CommitStore};
24use crate::conversation::{
25    ConversationReader, ConversationRecord, ConversationWriter, RecoveryScope, StoredTurn,
26    StoredUserTurn, TurnPhaseMarker,
27};
28use crate::error::StoreError;
29use crate::events::{
30    EventBatch, EventCursor, EventJournalReader, EventJournalWriter, EventPage, StoredEvent,
31};
32use crate::interaction::{
33    InteractionReader, InteractionRecord, InteractionWriter, InvalidationReason, ResolutionOutcome,
34};
35use crate::journal::{
36    CommandJournalEntry, CommandJournalReader, CommandJournalWriter, JournalAdmission,
37    JournalOutcome,
38};
39use crate::outbox::{OutboxReader, OutboxRecord, OutboxWriter};
40use crate::replay::{ReplayReader, ReplayWriter};
41
42#[async_trait]
43impl ConversationReader for MemoryStores {
44    async fn load_conversation(
45        &self,
46        account: &AccountId,
47        id: &ConversationId,
48    ) -> Result<ConversationRecord, StoreError> {
49        self.inner()?.load_conversation(account, id)
50    }
51
52    async fn load_recent_turns(
53        &self,
54        account: &AccountId,
55        conversation: &ConversationId,
56        limit: usize,
57    ) -> Result<Vec<StoredTurn>, StoreError> {
58        self.inner()?
59            .load_recent_turns(account, conversation, limit)
60    }
61
62    async fn load_turn(
63        &self,
64        account: &AccountId,
65        turn_id: &TurnId,
66    ) -> Result<StoredTurn, StoreError> {
67        self.inner()?.load_turn(account, turn_id)
68    }
69
70    async fn turn_phase(
71        &self,
72        account: &AccountId,
73        turn_id: &TurnId,
74    ) -> Result<TurnPhaseMarker, StoreError> {
75        self.inner()?.turn_phase(account, turn_id)
76    }
77
78    async fn list_unfinished_turns(
79        &self,
80        scope: RecoveryScope,
81        limit: usize,
82    ) -> Result<Vec<TurnPhaseMarker>, StoreError> {
83        Ok(self.inner()?.list_unfinished_turns(&scope, limit))
84    }
85}
86
87#[async_trait]
88impl ConversationWriter for MemoryStores {
89    async fn create_conversation(&self, record: ConversationRecord) -> Result<(), StoreError> {
90        self.inner()?.create_conversation(record)
91    }
92
93    async fn append_user_turn(&self, turn: StoredUserTurn) -> Result<(), StoreError> {
94        let now = self.now();
95        self.inner()?.append_user_turn(turn, now)
96    }
97
98    async fn append_assistant_turn(
99        &self,
100        account: &AccountId,
101        turn: AssistantTurn,
102    ) -> Result<(), StoreError> {
103        self.fire(FailurePoint::BeforeResponsePersistence)?;
104        self.inner()?.append_assistant_turn(account, turn)
105    }
106
107    async fn set_turn_phase(
108        &self,
109        account: &AccountId,
110        turn_id: &TurnId,
111        phase: TurnPhase,
112    ) -> Result<TurnPhaseMarker, StoreError> {
113        let now = self.now();
114        self.inner()?.set_turn_phase(account, turn_id, phase, now)
115    }
116}
117
118#[async_trait]
119impl InteractionReader for MemoryStores {
120    async fn get(
121        &self,
122        account: &AccountId,
123        id: &InteractionId,
124    ) -> Result<InteractionRecord, StoreError> {
125        self.inner()?.get_interaction(account, id)
126    }
127
128    async fn list_open_for_conversation(
129        &self,
130        account: &AccountId,
131        conversation: &ConversationId,
132    ) -> Result<Vec<Interaction>, StoreError> {
133        Ok(self
134            .inner()?
135            .list_open_for_conversation(account, conversation))
136    }
137
138    async fn list_open_for_case(
139        &self,
140        account: &AccountId,
141        case_key: &CaseKey,
142    ) -> Result<Vec<Interaction>, StoreError> {
143        Ok(self.inner()?.list_open_for_case(account, case_key))
144    }
145
146    async fn blocking_answered_at(
147        &self,
148        account: &AccountId,
149        case_key: &CaseKey,
150        revision: CaseRevision,
151    ) -> Result<bool, StoreError> {
152        Ok(self
153            .inner()?
154            .blocking_answered_at(account, case_key, revision))
155    }
156}
157
158#[async_trait]
159impl InteractionWriter for MemoryStores {
160    async fn insert(&self, interaction: Interaction) -> Result<(), StoreError> {
161        let now = self.now();
162        self.inner()?.insert_interaction(interaction, false, now)?;
163        self.fire(FailurePoint::AfterInteractionPersistence)
164    }
165
166    async fn insert_replacing_blocking(
167        &self,
168        interaction: Interaction,
169    ) -> Result<Vec<InteractionId>, StoreError> {
170        let now = self.now();
171        let invalidated = self.inner()?.insert_interaction(interaction, true, now)?;
172        self.fire(FailurePoint::AfterInteractionPersistence)?;
173        Ok(invalidated)
174    }
175
176    async fn begin_resolution(
177        &self,
178        account: &AccountId,
179        id: &InteractionId,
180        expected_status: InteractionStatus,
181        option_id: OptionId,
182        resolved_by: TurnId,
183    ) -> Result<InteractionRecord, StoreError> {
184        let now = self.now();
185        self.inner()?
186            .begin_resolution(account, id, expected_status, option_id, resolved_by, now)
187    }
188
189    async fn finish_resolution(
190        &self,
191        account: &AccountId,
192        id: &InteractionId,
193        outcome: ResolutionOutcome,
194    ) -> Result<InteractionRecord, StoreError> {
195        self.inner()?.finish_resolution(account, id, outcome)
196    }
197
198    async fn invalidate_for_case(
199        &self,
200        account: &AccountId,
201        case_key: &CaseKey,
202        new_revision: CaseRevision,
203        reason: InvalidationReason,
204    ) -> Result<Vec<InteractionId>, StoreError> {
205        let now = self.now();
206        self.inner()?
207            .invalidate_for_case(account, case_key, new_revision, reason, now)
208    }
209
210    async fn invalidate_case_cards(
211        &self,
212        account: &AccountId,
213        case_key: &CaseKey,
214        reason: InvalidationReason,
215    ) -> Result<Vec<InteractionId>, StoreError> {
216        let now = self.now();
217        self.inner()?
218            .invalidate_case_cards(account, case_key, reason, now)
219    }
220
221    async fn expire_due(&self, now: DateTime<Utc>) -> Result<Vec<InteractionId>, StoreError> {
222        self.inner()?.expire_due(now)
223    }
224}
225
226#[async_trait]
227impl CommandJournalReader for MemoryStores {
228    async fn get(
229        &self,
230        account: &AccountId,
231        command_id: &CommandId,
232    ) -> Result<CommandJournalEntry, StoreError> {
233        self.inner()?.journal_get(account, command_id)
234    }
235
236    async fn for_turn(
237        &self,
238        account: &AccountId,
239        turn_id: &TurnId,
240    ) -> Result<Vec<CommandJournalEntry>, StoreError> {
241        Ok(self.inner()?.journal_for_turn(account, turn_id, false))
242    }
243
244    async fn pending_for_turn(
245        &self,
246        account: &AccountId,
247        turn_id: &TurnId,
248    ) -> Result<Vec<CommandJournalEntry>, StoreError> {
249        Ok(self.inner()?.journal_for_turn(account, turn_id, true))
250    }
251}
252
253#[async_trait]
254impl CommandJournalWriter for MemoryStores {
255    async fn begin(&self, entry: CommandJournalEntry) -> Result<JournalAdmission, StoreError> {
256        self.fire(FailurePoint::BeforeJournalInsert)?;
257        let admission = self.inner()?.journal_begin(entry)?;
258        if admission.is_fresh() {
259            self.fire(FailurePoint::AfterJournalInsertBeforeCommit)?;
260        }
261        Ok(admission)
262    }
263
264    async fn mark_executing(
265        &self,
266        account: &AccountId,
267        command_id: &CommandId,
268    ) -> Result<(), StoreError> {
269        self.inner()?.journal_mark_executing(account, command_id)
270    }
271
272    async fn complete(
273        &self,
274        account: &AccountId,
275        command_id: &CommandId,
276        outcome: JournalOutcome,
277    ) -> Result<(), StoreError> {
278        let now = self.now();
279        self.inner()?
280            .journal_complete(account, command_id, outcome, now)
281    }
282}
283
284#[async_trait]
285impl EventJournalReader for MemoryStores {
286    async fn list_since(
287        &self,
288        account: &AccountId,
289        case_key: &CaseKey,
290        since: CaseRevision,
291        limit: usize,
292    ) -> Result<Vec<StoredEvent>, StoreError> {
293        Ok(self
294            .inner()?
295            .list_events_since(account, case_key, since, limit))
296    }
297
298    async fn read_from(
299        &self,
300        account: &AccountId,
301        after: EventCursor,
302        limit: usize,
303    ) -> Result<EventPage, StoreError> {
304        Ok(self.inner()?.read_events_from(account, after, limit))
305    }
306
307    async fn get_by_ids(
308        &self,
309        account: &AccountId,
310        ids: &[EventId],
311    ) -> Result<Vec<StoredEvent>, StoreError> {
312        Ok(self.inner()?.events_by_ids(account, ids))
313    }
314
315    async fn count(&self, account: &AccountId, case_key: &CaseKey) -> Result<u64, StoreError> {
316        Ok(self.inner()?.count_events(account, case_key))
317    }
318}
319
320#[async_trait]
321impl EventJournalWriter for MemoryStores {
322    async fn append(&self, batch: EventBatch) -> Result<Vec<EventId>, StoreError> {
323        let appended = self.inner()?.append_events(batch)?;
324        self.fire(FailurePoint::AfterCommitBeforeEventReadback)?;
325        Ok(appended)
326    }
327
328    async fn redact_payload(
329        &self,
330        account: &AccountId,
331        event_id: &EventId,
332        authority: &RedactionAuthority,
333    ) -> Result<EventRedaction, StoreError> {
334        let now = self.now();
335        self.inner()?
336            .redact_event_payload(account, event_id, authority, now)
337    }
338}
339
340#[async_trait]
341impl OutboxReader for MemoryStores {
342    async fn get(&self, outbox_id: &OutboxId) -> Result<OutboxRecord, StoreError> {
343        self.inner()?.get_outbox(outbox_id)
344    }
345
346    async fn list_for_command(
347        &self,
348        command_id: &CommandId,
349    ) -> Result<Vec<OutboxRecord>, StoreError> {
350        Ok(self.inner()?.outbox_for_command(command_id))
351    }
352}
353
354#[async_trait]
355impl OutboxWriter for MemoryStores {
356    async fn enqueue(&self, entry: OutboxEntry) -> Result<(), StoreError> {
357        self.inner()?.enqueue_outbox(entry)
358    }
359
360    async fn claim_due(
361        &self,
362        now: DateTime<Utc>,
363        limit: usize,
364        worker_id: &str,
365    ) -> Result<Vec<OutboxEntry>, StoreError> {
366        self.fire(FailurePoint::BeforeOutboxDispatch)?;
367        Ok(self.inner()?.claim_due(now, limit, worker_id))
368    }
369
370    async fn mark_completed(&self, outbox_id: &OutboxId) -> Result<(), StoreError> {
371        self.fire(FailurePoint::AfterOutboxDispatch)?;
372        let now = self.now();
373        self.inner()?.mark_outbox_completed(outbox_id, now)
374    }
375
376    async fn mark_failed(
377        &self,
378        outbox_id: &OutboxId,
379        reason: String,
380        retry_at: Option<DateTime<Utc>>,
381    ) -> Result<(), StoreError> {
382        self.fire(FailurePoint::AfterOutboxDispatch)?;
383        let now = self.now();
384        self.inner()?
385            .mark_outbox_failed(outbox_id, reason, retry_at, now)
386    }
387
388    async fn mark_outcome_unknown(
389        &self,
390        outbox_id: &OutboxId,
391        remote_ref: Option<String>,
392    ) -> Result<(), StoreError> {
393        self.fire(FailurePoint::AfterOutboxDispatch)?;
394        self.inner()?
395            .mark_outbox_outcome_unknown(outbox_id, remote_ref)
396    }
397
398    async fn reschedule(
399        &self,
400        outbox_id: &OutboxId,
401        next_attempt_at: DateTime<Utc>,
402    ) -> Result<(), StoreError> {
403        self.inner()?.reschedule_outbox(outbox_id, next_attempt_at)
404    }
405
406    async fn release_expired_claims(
407        &self,
408        claimed_before: DateTime<Utc>,
409    ) -> Result<Vec<OutboxId>, StoreError> {
410        Ok(self.inner()?.release_expired_claims(claimed_before))
411    }
412}
413
414#[async_trait]
415impl ReplayReader for MemoryStores {
416    async fn get(&self, account: &AccountId, turn_id: &TurnId) -> Result<ReplayRecord, StoreError> {
417        self.inner()?.get_replay(account, turn_id)
418    }
419
420    async fn list_for_conversation(
421        &self,
422        account: &AccountId,
423        conversation: &ConversationId,
424        limit: usize,
425    ) -> Result<Vec<ReplayRecord>, StoreError> {
426        Ok(self
427            .inner()?
428            .replays_for_conversation(account, conversation, limit))
429    }
430}
431
432#[async_trait]
433impl ReplayWriter for MemoryStores {
434    async fn put(&self, record: ReplayRecord) -> Result<(), StoreError> {
435        self.inner()?.put_replay(record);
436        Ok(())
437    }
438}
439
440#[async_trait]
441impl CommitStore for MemoryStores {
442    async fn commit(
443        &self,
444        account: &AccountId,
445        bundle: CommitBundle,
446    ) -> Result<CommitReceipt, StoreError> {
447        let now = self.now();
448        let faults = &self.faults;
449        let mut probe = |point: FailurePoint| match faults.lock() {
450            Ok(mut queue) => queue.take(point),
451            Err(_) => Some(StoreError::Corrupt),
452        };
453        // The bundle mutates the state directly and records how to undo each
454        // mutation; an error rolls the recorded reversals back before returning,
455        // so nothing of a failed bundle is ever visible. That is the
456        // all-or-nothing of spec §16.3; see `Inner::apply_bundle` for why it is
457        // done this way rather than by copying the state.
458        self.inner()?.apply_bundle(account, bundle, now, &mut probe)
459    }
460}
461
462#[cfg(test)]
463mod tests {
464    use std::sync::Arc;
465
466    use turnframe_core::case::CaseRef;
467    use turnframe_core::command::{CommandOrigin, IdempotencyKey};
468    use turnframe_core::event::CommittedEvent;
469    use turnframe_core::ids::{BlockId, ConversationId, UserId};
470    use turnframe_core::interaction::{
471        InteractionKind, InteractionOption, InteractionPayload, InteractionSpec,
472        StoredInteractionAction,
473    };
474    use turnframe_core::locale::Locale;
475    use turnframe_core::response::{NoticeSeverity, ReplayToken, ResponseBlock, ServerNotice};
476    use turnframe_core::turn::{ActorContext, TurnInput};
477
478    use super::*;
479    use crate::journal::{CommandJournalStatus, JournalAdmission};
480    use crate::memory::{ManualClock, MemoryStores};
481
482    fn account() -> AccountId {
483        AccountId::from("chaos")
484    }
485
486    fn epoch() -> DateTime<Utc> {
487        DateTime::<Utc>::UNIX_EPOCH
488    }
489
490    fn case() -> CaseRef {
491        CaseRef::new("chaos", "case-1", CaseRevision(1))
492    }
493
494    fn stores() -> (Arc<MemoryStores>, Arc<ManualClock>) {
495        let clock = Arc::new(ManualClock::epoch());
496        (Arc::new(MemoryStores::with_clock(clock.clone())), clock)
497    }
498
499    fn card(id: InteractionId, blocking: bool) -> Interaction {
500        let mut spec = InteractionSpec::new(
501            "chaos",
502            case(),
503            InteractionKind::SingleSelect,
504            InteractionPayload::new("chaos card").with_option(InteractionOption::new(
505                "ack",
506                "Got it",
507                StoredInteractionAction::Dismiss,
508            )),
509        );
510        if !blocking {
511            spec = spec.non_blocking();
512        }
513        Interaction::from_spec(
514            spec,
515            id,
516            account(),
517            ConversationId::nil(),
518            TurnId::nil(),
519            epoch(),
520        )
521        .expect("a valid fixture card")
522    }
523
524    fn entry(command_id: CommandId, turn: TurnId, key: &str) -> CommandJournalEntry {
525        CommandJournalEntry {
526            command_id,
527            account_id: account(),
528            idempotency_key: IdempotencyKey::new(key),
529            turn_id: turn,
530            case_ref: case(),
531            command_type: "chaos.noop".to_owned(),
532            command_payload: serde_json::json!({ "key": key }),
533            origin: CommandOrigin::InternalPolicy {
534                policy_key: "chaos".to_owned(),
535            },
536            status: CommandJournalStatus::Pending,
537            result: None,
538            created_at: epoch(),
539            completed_at: None,
540        }
541    }
542
543    fn batch(command_id: CommandId, ids: &[EventId]) -> EventBatch {
544        EventBatch::new(
545            account(),
546            case().key(),
547            command_id,
548            CaseRevision(2),
549            ids.iter()
550                .map(|id| CommittedEvent {
551                    event_id: *id,
552                    event_type: "chaos.happened".to_owned(),
553                    occurred_at: epoch(),
554                    payload: serde_json::Value::Null,
555                })
556                .collect(),
557        )
558    }
559
560    async fn conversation_with_turn(store: &MemoryStores, turn: TurnId) -> ConversationId {
561        let conversation = ConversationId::new();
562        store
563            .create_conversation(ConversationRecord::new(conversation, account(), epoch()))
564            .await
565            .expect("the conversation is created");
566        store
567            .append_user_turn(StoredUserTurn::new(
568                TurnInput {
569                    turn_id: turn,
570                    conversation_id: conversation,
571                    actor: ActorContext::new(account(), UserId::from("u")),
572                    text: Some("hello".to_owned()),
573                    interaction_response: None,
574                    attachments: Vec::new(),
575                    origin: None,
576                    locale: Locale::from("en"),
577                    effort: None,
578                },
579                epoch(),
580            ))
581            .await
582            .expect("the turn is appended");
583        conversation
584    }
585
586    #[tokio::test]
587    async fn a_failure_after_interaction_persistence_leaves_the_card_written() {
588        let (store, _clock) = stores();
589        let id = InteractionId::new();
590        store
591            .fail_next(
592                FailurePoint::AfterInteractionPersistence,
593                StoreError::Unavailable,
594            )
595            .unwrap();
596        assert_eq!(
597            InteractionWriter::insert(store.as_ref(), card(id, true)).await,
598            Err(StoreError::Unavailable)
599        );
600        // Spec §15.5: the response must not mention a card whose persistence
601        // was reported as failed — but the row is there, and recovery sees it.
602        let record = InteractionReader::get(store.as_ref(), &account(), &id)
603            .await
604            .expect("the card survived the failure");
605        assert_eq!(record.status(), InteractionStatus::Active);
606        assert!(store.armed_failures().unwrap().is_empty());
607    }
608
609    #[tokio::test]
610    async fn a_failure_before_the_journal_insert_leaves_the_key_free() {
611        let (store, _clock) = stores();
612        let turn = TurnId::new();
613        let command = CommandId::new();
614        store
615            .fail_next(FailurePoint::BeforeJournalInsert, StoreError::Unavailable)
616            .unwrap();
617        assert_eq!(
618            CommandJournalWriter::begin(store.as_ref(), entry(command, turn, "k")).await,
619            Err(StoreError::Unavailable)
620        );
621        assert_eq!(
622            CommandJournalReader::get(store.as_ref(), &account(), &command).await,
623            Err(StoreError::NotFound)
624        );
625        // Nothing was written, so the command may be admitted from scratch.
626        assert_eq!(
627            CommandJournalWriter::begin(store.as_ref(), entry(command, turn, "k")).await,
628            Ok(JournalAdmission::Fresh)
629        );
630    }
631
632    #[tokio::test]
633    async fn a_failure_after_the_journal_insert_leaves_a_pending_entry_for_recovery() {
634        let (store, _clock) = stores();
635        let turn = TurnId::new();
636        let command = CommandId::new();
637        store
638            .fail_next(
639                FailurePoint::AfterJournalInsertBeforeCommit,
640                StoreError::Timeout,
641            )
642            .unwrap();
643        assert_eq!(
644            CommandJournalWriter::begin(store.as_ref(), entry(command, turn, "k")).await,
645            Err(StoreError::Timeout)
646        );
647        let pending = CommandJournalReader::pending_for_turn(store.as_ref(), &account(), &turn)
648            .await
649            .expect("recovery lists the turn");
650        assert_eq!(pending.len(), 1, "spec §23.1: resume by idempotency key");
651        assert_eq!(pending[0].command_id, command);
652        // And the key is now taken, so a retry replays instead of duplicating.
653        let admission =
654            CommandJournalWriter::begin(store.as_ref(), entry(CommandId::new(), turn, "k"))
655                .await
656                .expect("the retry is admitted");
657        assert_eq!(
658            admission.replayed().map(|entry| entry.command_id),
659            Some(command)
660        );
661    }
662
663    #[tokio::test]
664    async fn a_failure_after_the_event_append_leaves_the_events_readable() {
665        let (store, _clock) = stores();
666        let command = CommandId::new();
667        let ids = vec![EventId::new()];
668        store
669            .fail_next(
670                FailurePoint::AfterCommitBeforeEventReadback,
671                StoreError::Timeout,
672            )
673            .unwrap();
674        assert_eq!(
675            EventJournalWriter::append(store.as_ref(), batch(command, &ids)).await,
676            Err(StoreError::Timeout)
677        );
678        let readback = EventJournalReader::get_by_ids(store.as_ref(), &account(), &ids)
679            .await
680            .expect("the ledger answers");
681        assert_eq!(readback.len(), 1, "the events landed before the failure");
682        // Appending them again must not duplicate the claim ledger.
683        assert_eq!(
684            EventJournalWriter::append(store.as_ref(), batch(command, &ids)).await,
685            Err(StoreError::Conflict)
686        );
687    }
688
689    #[tokio::test]
690    async fn outbox_failures_bracket_the_dispatch() {
691        let (store, _clock) = stores();
692        let outbox_id = OutboxId::new();
693        let entry = OutboxEntry {
694            outbox_id,
695            command_id: CommandId::new(),
696            destination: "chaos".to_owned(),
697            payload: serde_json::Value::Null,
698            idempotency_key: IdempotencyKey::new("k"),
699            status: turnframe_core::event::OutboxStatus::Pending,
700            attempt_count: 0,
701            next_attempt_at: None,
702            created_at: epoch(),
703            completed_at: None,
704        };
705        OutboxWriter::enqueue(store.as_ref(), entry)
706            .await
707            .expect("the row is enqueued");
708
709        store
710            .fail_next(FailurePoint::BeforeOutboxDispatch, StoreError::Unavailable)
711            .unwrap();
712        assert_eq!(
713            OutboxWriter::claim_due(store.as_ref(), epoch(), 10, "w").await,
714            Err(StoreError::Unavailable)
715        );
716        let untouched = OutboxReader::get(store.as_ref(), &outbox_id)
717            .await
718            .expect("the row is readable");
719        assert_eq!(untouched.entry.attempt_count, 0, "nothing was claimed");
720
721        let claimed = OutboxWriter::claim_due(store.as_ref(), epoch(), 10, "w")
722            .await
723            .expect("the row is claimed");
724        assert_eq!(claimed.len(), 1);
725        store
726            .fail_next(FailurePoint::AfterOutboxDispatch, StoreError::Timeout)
727            .unwrap();
728        assert_eq!(
729            OutboxWriter::mark_completed(store.as_ref(), &outbox_id).await,
730            Err(StoreError::Timeout)
731        );
732        // The worst case of spec §16.5: the call may have happened and nothing
733        // local says so. Only the reaper gets the row moving again.
734        let stranded = OutboxReader::get(store.as_ref(), &outbox_id)
735            .await
736            .expect("the row is readable");
737        assert_eq!(
738            stranded.entry.status,
739            turnframe_core::event::OutboxStatus::Dispatching
740        );
741        assert!(stranded.claim.is_some());
742        let released = OutboxWriter::release_expired_claims(
743            store.as_ref(),
744            epoch() + chrono::TimeDelta::minutes(5),
745        )
746        .await
747        .expect("the reaper runs");
748        assert_eq!(released, vec![outbox_id]);
749    }
750
751    #[tokio::test]
752    async fn a_failure_before_response_persistence_keeps_the_turn_answerable() {
753        let (store, _clock) = stores();
754        let turn = TurnId::new();
755        let conversation = conversation_with_turn(store.as_ref(), turn).await;
756        let assistant = AssistantTurn {
757            turn_id: turn,
758            conversation_id: conversation,
759            blocks: vec![ResponseBlock::Notice(ServerNotice {
760                block_id: BlockId::from("b1"),
761                code: "turnframe.notice.chaos".to_owned(),
762                severity: NoticeSeverity::Info,
763                text: "hello".into(),
764            })],
765            subjects: Vec::new(),
766            expectations: Vec::new(),
767            replay_token: ReplayToken::new("t"),
768            done: Vec::new(),
769        };
770        store
771            .fail_next(
772                FailurePoint::BeforeResponsePersistence,
773                StoreError::Unavailable,
774            )
775            .unwrap();
776        assert_eq!(
777            ConversationWriter::append_assistant_turn(
778                store.as_ref(),
779                &account(),
780                assistant.clone()
781            )
782            .await,
783            Err(StoreError::Unavailable)
784        );
785        let stored = ConversationReader::load_turn(store.as_ref(), &account(), &turn)
786            .await
787            .expect("the turn is still there");
788        assert!(stored.assistant.is_none(), "nothing was written");
789        // Recovery regenerates and persists it; the slot is still free.
790        ConversationWriter::append_assistant_turn(store.as_ref(), &account(), assistant.clone())
791            .await
792            .expect("the retry succeeds");
793        let recovered = ConversationReader::load_turn(store.as_ref(), &account(), &turn)
794            .await
795            .expect("the turn is readable");
796        assert_eq!(recovered.assistant, Some(assistant));
797    }
798
799    #[tokio::test]
800    async fn a_failure_mid_bundle_makes_nothing_of_it_visible() {
801        for point in [
802            FailurePoint::AfterJournalInsertBeforeCommit,
803            FailurePoint::AfterCommitBeforeEventReadback,
804            FailurePoint::AfterInteractionPersistence,
805        ] {
806            let (store, _clock) = stores();
807            let turn = TurnId::new();
808            let conversation = conversation_with_turn(store.as_ref(), turn).await;
809            let command = CommandId::new();
810            CommandJournalWriter::begin(store.as_ref(), entry(command, turn, "k"))
811                .await
812                .expect("the command is admitted");
813            let event_ids = vec![EventId::new()];
814            let inserted = InteractionId::new();
815            let bundle = CommitBundle::new()
816                .with_journal_completion(
817                    command,
818                    JournalOutcome::Committed {
819                        new_revision: CaseRevision(2),
820                        event_ids: event_ids.clone(),
821                    },
822                )
823                .with_events(batch(command, &event_ids))
824                .with_interaction_insert(card(inserted, true), false)
825                .with_replay_record(turnframe_core::replay::ReplayRecord::received(
826                    turn,
827                    conversation,
828                    account(),
829                    epoch(),
830                ))
831                .with_turn_phase(turn, TurnPhase::Committed);
832
833            store.fail_next(point, StoreError::Unavailable).unwrap();
834            assert_eq!(
835                CommitStore::commit(store.as_ref(), &account(), bundle).await,
836                Err(StoreError::Unavailable),
837                "{point:?}"
838            );
839
840            // Not one stage of the bundle may have survived, including the
841            // stages that ran before the injected failure.
842            let entry = CommandJournalReader::get(store.as_ref(), &account(), &command)
843                .await
844                .expect("the entry is readable");
845            assert_eq!(entry.status, CommandJournalStatus::Pending, "{point:?}");
846            assert!(entry.result.is_none(), "{point:?}");
847            assert!(
848                EventJournalReader::get_by_ids(store.as_ref(), &account(), &event_ids)
849                    .await
850                    .expect("the ledger answers")
851                    .is_empty(),
852                "{point:?}"
853            );
854            assert_eq!(
855                InteractionReader::get(store.as_ref(), &account(), &inserted).await,
856                Err(StoreError::NotFound),
857                "{point:?}"
858            );
859            assert_eq!(
860                ReplayReader::get(store.as_ref(), &account(), &turn).await,
861                Err(StoreError::NotFound),
862                "{point:?}"
863            );
864            assert_eq!(
865                ConversationReader::turn_phase(store.as_ref(), &account(), &turn)
866                    .await
867                    .expect("the marker is readable")
868                    .phase,
869                TurnPhase::Received,
870                "{point:?}"
871            );
872        }
873    }
874
875    #[tokio::test]
876    async fn the_clock_stamps_what_the_store_owns() {
877        let (store, clock) = stores();
878        let turn = TurnId::new();
879        conversation_with_turn(store.as_ref(), turn).await;
880        clock.advance(chrono::TimeDelta::seconds(42));
881        let marker = ConversationWriter::set_turn_phase(
882            store.as_ref(),
883            &account(),
884            &turn,
885            TurnPhase::Reduced,
886        )
887        .await
888        .expect("the phase moves");
889        assert_eq!(marker.updated_at, epoch() + chrono::TimeDelta::seconds(42));
890    }
891
892    #[tokio::test]
893    async fn non_blocking_cards_never_take_the_blocking_slot() {
894        let (store, _clock) = stores();
895        for _ in 0..3 {
896            InteractionWriter::insert(store.as_ref(), card(InteractionId::new(), false))
897                .await
898                .expect("a non-blocking card always fits");
899        }
900        InteractionWriter::insert(store.as_ref(), card(InteractionId::new(), true))
901            .await
902            .expect("the slot is still free");
903        assert_eq!(
904            InteractionWriter::insert(store.as_ref(), card(InteractionId::new(), true)).await,
905            Err(StoreError::Conflict)
906        );
907    }
908}