1use 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 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 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 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 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 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 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 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 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}