Skip to main content

turnframe_test/stores/
mod.rs

1//! Fake stores: the in-memory persistence layer with the three things a
2//! runtime test always ends up needing (spec §27.4, §27.7).
3//!
4//! `turnframe-store` already ships the deterministic implementation and the
5//! conformance suite that proves an adapter right. What it does not ship is the
6//! ergonomics of *driving* it from a turn test, and those are the same three
7//! every time:
8//!
9//! * **crash where I say.** Recovery code is never exercised by the happy path,
10//!   so [`FakeStores::fail_at`] arms a failure at one of the boundaries of
11//!   spec §27.7, and [`crash_boundary`] looks one up by name so a table-driven
12//!   chaos test can walk all seven.
13//! * **count what was called.** "The turn wrote the journal once" and "the
14//!   recovery path did not re-append the events" are assertions about call
15//!   counts, not about state. Every trait method is tallied as
16//!   `"<role>.<method>"`.
17//! * **stop the clock.** Expiry, claim staleness and phase markers are
18//!   timestamps the store stamps itself; a frozen clock turns "the card
19//!   expired" into something the test makes true instead of waits for.
20//!
21//! And then: **assert what was persisted**. The accessors on [`FakeStores`]
22//! read through the underlying store *without* counting, so checking the
23//! outcome of a turn never pollutes the tally the same test is asserting on.
24//!
25//! # The suite comes with it
26//!
27//! [`conformance`] is the store crate's own suite, re-exported so an adopter
28//! building a store over their own database reaches fixtures, doubles and the
29//! contract through this one crate.
30//!
31//! ```
32//! use turnframe_core::ids::{AccountId, ConversationId};
33//! use turnframe_store::conversation::{ConversationRecord, ConversationStore};
34//! use turnframe_test::stores::FakeStores;
35//!
36//! # fn main() -> Result<(), Box<dyn std::error::Error>> {
37//! # tokio::runtime::Runtime::new()?.block_on(async {
38//! let fake = FakeStores::new();
39//! let account = AccountId::from("aurora");
40//! let conversation = ConversationId::nil();
41//!
42//! fake.stores()
43//!     .conversations()
44//!     .create_conversation(ConversationRecord::new(conversation, account.clone(), fake.now()))
45//!     .await?;
46//!
47//! assert_eq!(fake.call_count("conversations.create_conversation"), 1);
48//! # Ok::<(), turnframe_store::error::StoreError>(())
49//! # })?;
50//! # Ok(())
51//! # }
52//! ```
53
54mod counting;
55
56use std::sync::Arc;
57
58use chrono::{DateTime, TimeDelta, Utc};
59use turnframe_core::case::CaseKey;
60use turnframe_core::ids::{AccountId, CommandId, ConversationId, EventId, InteractionId, TurnId};
61use turnframe_core::interaction::Interaction;
62use turnframe_core::replay::ReplayRecord;
63use turnframe_store::conversation::{ConversationReader, StoredTurn, TurnPhaseMarker};
64use turnframe_store::error::StoreError;
65use turnframe_store::events::{EventJournalReader, StoredEvent};
66use turnframe_store::interaction::{InteractionReader, InteractionRecord};
67use turnframe_store::journal::{CommandJournalEntry, CommandJournalReader};
68use turnframe_store::outbox::{OutboxReader, OutboxRecord};
69use turnframe_store::replay::ReplayReader;
70
71pub use counting::{CountingStores, StoreCallCounter};
72pub use turnframe_store::conformance;
73pub use turnframe_store::memory::{Clock, FailurePoint, ManualClock, MemoryStores, SystemClock};
74pub use turnframe_store::stores::Stores;
75
76/// How many events [`FakeStores::events_for`] reads at most. High enough that a
77/// test never has to think about it, low enough to stay a bounded read.
78pub const MAX_LISTED_EVENTS: usize = 1024;
79
80/// The crash boundaries of spec §27.7, by name.
81///
82/// The names are the ones the specification uses, so a chaos test can be
83/// written as a table and read next to the paragraph it implements.
84pub const CRASH_BOUNDARIES: [(&str, FailurePoint); 7] = [
85    (
86        "after_interaction_persistence",
87        FailurePoint::AfterInteractionPersistence,
88    ),
89    ("before_journal_insert", FailurePoint::BeforeJournalInsert),
90    (
91        "after_journal_insert_before_commit",
92        FailurePoint::AfterJournalInsertBeforeCommit,
93    ),
94    (
95        "after_commit_before_event_readback",
96        FailurePoint::AfterCommitBeforeEventReadback,
97    ),
98    ("before_outbox_dispatch", FailurePoint::BeforeOutboxDispatch),
99    ("after_outbox_dispatch", FailurePoint::AfterOutboxDispatch),
100    (
101        "before_response_persistence",
102        FailurePoint::BeforeResponsePersistence,
103    ),
104];
105
106/// The boundary with this name, or `None`.
107#[must_use]
108pub fn crash_boundary(name: &str) -> Option<FailurePoint> {
109    CRASH_BOUNDARIES
110        .iter()
111        .find(|(candidate, _)| *candidate == name)
112        .map(|(_, point)| *point)
113}
114
115/// The name of a boundary, or `"unknown"` for one added after this table.
116#[must_use]
117pub fn boundary_name(point: FailurePoint) -> &'static str {
118    CRASH_BOUNDARIES
119        .iter()
120        .find(|(_, candidate)| *candidate == point)
121        .map_or("unknown", |(name, _)| *name)
122}
123
124/// A name that is not one of the [`CRASH_BOUNDARIES`].
125#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
126#[error("no crash boundary is named `{name}`")]
127pub struct UnknownBoundary {
128    /// The name that was not found.
129    pub name: String,
130}
131
132/// The in-memory stores, wrapped for turn tests.
133///
134/// Cheap to build and self-contained: one [`MemoryStores`] on a frozen
135/// [`ManualClock`], a [`CountingStores`] in front of it, and a [`Stores`] built
136/// from the counting layer — which is the value the runtime takes.
137#[derive(Debug)]
138pub struct FakeStores {
139    memory: Arc<MemoryStores>,
140    counter: Arc<StoreCallCounter>,
141    clock: Arc<ManualClock>,
142    stores: Stores,
143}
144
145impl FakeStores {
146    /// Empty stores on a clock frozen at the Unix epoch.
147    #[must_use]
148    pub fn new() -> Self {
149        Self::at(DateTime::<Utc>::UNIX_EPOCH)
150    }
151
152    /// Empty stores on a clock frozen at `start`.
153    #[must_use]
154    pub fn at(start: DateTime<Utc>) -> Self {
155        let clock = Arc::new(ManualClock::new(start));
156        let memory = Arc::new(MemoryStores::with_clock(clock.clone()));
157        let counter = Arc::new(StoreCallCounter::new());
158        let counting = Arc::new(CountingStores::new(memory.clone(), counter.clone()));
159        let stores = Stores::builder()
160            .conversations(counting.clone())
161            .interactions(counting.clone())
162            .journal(counting.clone())
163            .events(counting.clone())
164            .outbox(counting.clone())
165            .replay(counting.clone())
166            .commit(counting)
167            .build()
168            // Every role was just supplied, so the builder cannot refuse. The
169            // fallback keeps this constructor infallible without an unwrap.
170            .unwrap_or_else(|_| Stores::from_memory(memory.clone()));
171        Self {
172            memory,
173            counter,
174            clock,
175            stores,
176        }
177    }
178
179    /// The set of stores to hand to the code under test.
180    #[must_use]
181    pub fn stores(&self) -> &Stores {
182        &self.stores
183    }
184
185    /// The underlying store, for the few things only it can do.
186    #[must_use]
187    pub fn memory(&self) -> &Arc<MemoryStores> {
188        &self.memory
189    }
190
191    /// The clock the store stamps its own columns from.
192    #[must_use]
193    pub fn clock(&self) -> &Arc<ManualClock> {
194        &self.clock
195    }
196
197    /// The instant the store would stamp a write with right now.
198    #[must_use]
199    pub fn now(&self) -> DateTime<Utc> {
200        self.clock.now()
201    }
202
203    /// Moves the clock forward.
204    pub fn advance(&self, delta: TimeDelta) {
205        self.clock.advance(delta);
206    }
207
208    /// Moves the clock to `instant`.
209    pub fn set_time(&self, instant: DateTime<Utc>) {
210        self.clock.set(instant);
211    }
212
213    // ---- failure injection -------------------------------------------------
214
215    /// Arms one failure at `point`: the next call reaching it returns `error`.
216    ///
217    /// Whether the write before the boundary survives is part of the point's
218    /// meaning; see [`FailurePoint`].
219    ///
220    /// # Errors
221    /// * `Corrupt` when the store's internal lock is poisoned.
222    pub fn fail_at(&self, point: FailurePoint, error: StoreError) -> Result<(), StoreError> {
223        self.memory.fail_next(point, error)
224    }
225
226    /// Arms one failure at the boundary with this name (see
227    /// [`CRASH_BOUNDARIES`]).
228    ///
229    /// # Errors
230    /// * [`UnknownBoundary`] when no boundary carries that name.
231    pub fn fail_at_boundary(&self, name: &str, error: StoreError) -> Result<(), UnknownBoundary> {
232        let point = crash_boundary(name).ok_or_else(|| UnknownBoundary {
233            name: name.to_owned(),
234        })?;
235        // A poisoned lock cannot be reported through this signature and is not
236        // what the caller is testing; the arming is simply dropped.
237        let _ = self.memory.fail_next(point, error);
238        Ok(())
239    }
240
241    /// The boundaries still armed, in arming order.
242    #[must_use]
243    pub fn armed_failures(&self) -> Vec<FailurePoint> {
244        self.memory.armed_failures().unwrap_or_default()
245    }
246
247    /// Disarms every pending failure.
248    pub fn clear_failures(&self) {
249        let _ = self.memory.clear_failures();
250    }
251
252    // ---- call counting -----------------------------------------------------
253
254    /// How many times `method` was called, as `"<role>.<method>"`.
255    #[must_use]
256    pub fn call_count(&self, method: &str) -> usize {
257        self.counter.count(method)
258    }
259
260    /// Every method called at least once, in key order.
261    #[must_use]
262    pub fn calls(&self) -> Vec<(&'static str, usize)> {
263        self.counter.counts()
264    }
265
266    /// Total calls across every store method.
267    #[must_use]
268    pub fn total_calls(&self) -> usize {
269        self.counter.total()
270    }
271
272    /// Forgets the tally, keeping the data. Use it between the arrange and act
273    /// phases of a test.
274    pub fn reset_calls(&self) {
275        self.counter.reset();
276    }
277
278    // ---- what was persisted ------------------------------------------------
279    //
280    // These read the store directly rather than through the counting layer: an
281    // assertion about the outcome must not change the tally the same test is
282    // asserting on.
283
284    /// Every event of a case, in append order.
285    ///
286    /// # Errors
287    /// * Whatever the store returns; `NotFound` is never one of them, an
288    ///   unknown case simply has no events.
289    pub async fn events_for(
290        &self,
291        account: &AccountId,
292        case_key: &CaseKey,
293    ) -> Result<Vec<StoredEvent>, StoreError> {
294        EventJournalReader::list_since(
295            self.memory.as_ref(),
296            account,
297            case_key,
298            turnframe_core::ids::CaseRevision::ZERO,
299            MAX_LISTED_EVENTS,
300        )
301        .await
302    }
303
304    /// The event types of a case, in append order: the shortest assertion about
305    /// what a turn actually committed.
306    ///
307    /// # Errors
308    /// * Whatever the store returns.
309    pub async fn event_types(
310        &self,
311        account: &AccountId,
312        case_key: &CaseKey,
313    ) -> Result<Vec<String>, StoreError> {
314        Ok(self
315            .events_for(account, case_key)
316            .await?
317            .into_iter()
318            .map(|event| event.event_type)
319            .collect())
320    }
321
322    /// The events with these identifiers, for checking a receipt's citations.
323    ///
324    /// # Errors
325    /// * Whatever the store returns.
326    pub async fn events_by_ids(
327        &self,
328        account: &AccountId,
329        ids: &[EventId],
330    ) -> Result<Vec<StoredEvent>, StoreError> {
331        EventJournalReader::get_by_ids(self.memory.as_ref(), account, ids).await
332    }
333
334    /// The journal entry of a command.
335    ///
336    /// # Errors
337    /// * `NotFound` when the command was never admitted for this account.
338    pub async fn journal_entry(
339        &self,
340        account: &AccountId,
341        command_id: &CommandId,
342    ) -> Result<CommandJournalEntry, StoreError> {
343        CommandJournalReader::get(self.memory.as_ref(), account, command_id).await
344    }
345
346    /// Every journal entry of a turn, in admission order.
347    ///
348    /// # Errors
349    /// * Whatever the store returns.
350    pub async fn journal_for_turn(
351        &self,
352        account: &AccountId,
353        turn_id: &TurnId,
354    ) -> Result<Vec<CommandJournalEntry>, StoreError> {
355        CommandJournalReader::for_turn(self.memory.as_ref(), account, turn_id).await
356    }
357
358    /// The open cards of a case.
359    ///
360    /// # Errors
361    /// * Whatever the store returns.
362    pub async fn open_interactions(
363        &self,
364        account: &AccountId,
365        case_key: &CaseKey,
366    ) -> Result<Vec<Interaction>, StoreError> {
367        InteractionReader::list_open_for_case(self.memory.as_ref(), account, case_key).await
368    }
369
370    /// One card, whatever its status.
371    ///
372    /// # Errors
373    /// * `NotFound` when it does not exist for this account.
374    pub async fn interaction(
375        &self,
376        account: &AccountId,
377        id: &InteractionId,
378    ) -> Result<InteractionRecord, StoreError> {
379        InteractionReader::get(self.memory.as_ref(), account, id).await
380    }
381
382    /// A stored turn, user side and assistant side.
383    ///
384    /// # Errors
385    /// * `NotFound` when the turn does not exist for this account.
386    pub async fn turn(
387        &self,
388        account: &AccountId,
389        turn_id: &TurnId,
390    ) -> Result<StoredTurn, StoreError> {
391        ConversationReader::load_turn(self.memory.as_ref(), account, turn_id).await
392    }
393
394    /// The phase marker of a turn.
395    ///
396    /// # Errors
397    /// * `NotFound` when the turn does not exist for this account.
398    pub async fn turn_phase(
399        &self,
400        account: &AccountId,
401        turn_id: &TurnId,
402    ) -> Result<TurnPhaseMarker, StoreError> {
403        ConversationReader::turn_phase(self.memory.as_ref(), account, turn_id).await
404    }
405
406    /// The replay record of a turn.
407    ///
408    /// # Errors
409    /// * `NotFound` when no record was written for this account and turn.
410    pub async fn replay_record(
411        &self,
412        account: &AccountId,
413        turn_id: &TurnId,
414    ) -> Result<ReplayRecord, StoreError> {
415        ReplayReader::get(self.memory.as_ref(), account, turn_id).await
416    }
417
418    /// The outbox rows a command produced.
419    ///
420    /// # Errors
421    /// * Whatever the store returns.
422    pub async fn outbox_for_command(
423        &self,
424        command_id: &CommandId,
425    ) -> Result<Vec<OutboxRecord>, StoreError> {
426        OutboxReader::list_for_command(self.memory.as_ref(), command_id).await
427    }
428
429    /// Every conversation-scoped replay record, most recent first.
430    ///
431    /// # Errors
432    /// * Whatever the store returns.
433    pub async fn replay_records_for(
434        &self,
435        account: &AccountId,
436        conversation: &ConversationId,
437        limit: usize,
438    ) -> Result<Vec<ReplayRecord>, StoreError> {
439        ReplayReader::list_for_conversation(self.memory.as_ref(), account, conversation, limit)
440            .await
441    }
442}
443
444impl Default for FakeStores {
445    fn default() -> Self {
446        Self::new()
447    }
448}
449
450#[cfg(test)]
451mod tests {
452    use super::*;
453    use turnframe_core::ids::{BlockId, ConversationId, TurnId};
454    use turnframe_core::locale::Locale;
455    use turnframe_core::replay::TurnPhase;
456    use turnframe_core::response::{
457        AssistantTurn, GeneratedTransition, ReplayToken, ResponseBlock,
458    };
459    use turnframe_core::turn::{ActorContext, TurnInput};
460    use turnframe_store::conversation::{ConversationRecord, StoredUserTurn};
461
462    fn account() -> AccountId {
463        AccountId::from("aurora")
464    }
465
466    fn turn_input(conversation: ConversationId, turn_id: TurnId) -> TurnInput {
467        TurnInput {
468            turn_id,
469            conversation_id: conversation,
470            actor: ActorContext::new("aurora", "user-1"),
471            text: Some("ciao".to_owned()),
472            interaction_response: None,
473            attachments: Vec::new(),
474            origin: None,
475            locale: Locale::from("it"),
476            effort: None,
477        }
478    }
479
480    fn assistant(turn_id: TurnId, conversation: ConversationId) -> AssistantTurn {
481        AssistantTurn {
482            turn_id,
483            conversation_id: conversation,
484            blocks: vec![ResponseBlock::Transition(GeneratedTransition {
485                block_id: BlockId::from("t1"),
486                text: "fatto".to_owned(),
487                facts_used: Vec::new(),
488            })],
489            replay_token: ReplayToken::from("rt"),
490            subjects: Vec::new(),
491            expectations: Vec::new(),
492            done: Vec::new(),
493        }
494    }
495
496    /// Creates a conversation and a user turn, then forgets the tally.
497    async fn arranged() -> (FakeStores, ConversationId, TurnId) {
498        let fake = FakeStores::new();
499        let conversation = ConversationId::nil();
500        let turn_id = TurnId::nil();
501        fake.stores()
502            .conversations()
503            .create_conversation(ConversationRecord::new(conversation, account(), fake.now()))
504            .await
505            .unwrap();
506        fake.stores()
507            .conversations()
508            .append_user_turn(StoredUserTurn::new(
509                turn_input(conversation, turn_id),
510                fake.now(),
511            ))
512            .await
513            .unwrap();
514        fake.reset_calls();
515        (fake, conversation, turn_id)
516    }
517
518    #[tokio::test]
519    async fn every_call_through_the_set_is_counted_by_role_and_method() {
520        let (fake, _conversation, turn_id) = arranged().await;
521        fake.stores()
522            .conversations()
523            .turn_phase(&account(), &turn_id)
524            .await
525            .unwrap();
526        fake.stores()
527            .conversations()
528            .turn_phase(&account(), &turn_id)
529            .await
530            .unwrap();
531
532        assert_eq!(fake.call_count("conversations.turn_phase"), 2);
533        assert_eq!(fake.total_calls(), 2);
534        assert_eq!(fake.calls(), vec![("conversations.turn_phase", 2)]);
535    }
536
537    #[tokio::test]
538    async fn reading_what_was_persisted_does_not_move_the_tally() {
539        let (fake, conversation, turn_id) = arranged().await;
540        fake.stores()
541            .conversations()
542            .append_assistant_turn(&account(), assistant(turn_id, conversation))
543            .await
544            .unwrap();
545
546        let stored = fake.turn(&account(), &turn_id).await.unwrap();
547        assert!(stored.assistant.is_some());
548        assert_eq!(
549            fake.calls(),
550            vec![("conversations.append_assistant_turn", 1)],
551            "the accessor bypasses the counting layer"
552        );
553    }
554
555    #[tokio::test]
556    async fn a_named_boundary_fails_the_next_call_that_reaches_it() {
557        let (fake, conversation, turn_id) = arranged().await;
558        fake.fail_at_boundary("before_response_persistence", StoreError::Unavailable)
559            .unwrap();
560        assert_eq!(
561            fake.armed_failures(),
562            vec![FailurePoint::BeforeResponsePersistence]
563        );
564
565        let refused = fake
566            .stores()
567            .conversations()
568            .append_assistant_turn(&account(), assistant(turn_id, conversation))
569            .await;
570
571        assert_eq!(refused, Err(StoreError::Unavailable));
572        // The point is named `Before…`, so nothing was written.
573        assert!(
574            fake.turn(&account(), &turn_id)
575                .await
576                .unwrap()
577                .assistant
578                .is_none()
579        );
580        assert!(fake.armed_failures().is_empty());
581    }
582
583    #[test]
584    fn the_boundary_table_covers_the_specification_and_round_trips() {
585        assert_eq!(CRASH_BOUNDARIES.len(), FailurePoint::ALL.len());
586        for point in FailurePoint::ALL {
587            let name = boundary_name(point);
588            assert_ne!(name, "unknown", "{point:?} is missing from the table");
589            assert_eq!(crash_boundary(name), Some(point));
590        }
591        assert_eq!(crash_boundary("nope"), None);
592
593        let fake = FakeStores::new();
594        assert_eq!(
595            fake.fail_at_boundary("nope", StoreError::Timeout)
596                .unwrap_err(),
597            UnknownBoundary {
598                name: "nope".to_owned()
599            }
600        );
601    }
602
603    #[tokio::test]
604    async fn the_clock_only_moves_when_the_test_moves_it() {
605        let (fake, _conversation, turn_id) = arranged().await;
606        assert_eq!(fake.now(), DateTime::<Utc>::UNIX_EPOCH);
607        let first = fake.turn_phase(&account(), &turn_id).await.unwrap();
608        assert_eq!(first.updated_at, DateTime::<Utc>::UNIX_EPOCH);
609
610        fake.advance(TimeDelta::minutes(5));
611        fake.stores()
612            .conversations()
613            .set_turn_phase(&account(), &turn_id, TurnPhase::Committed)
614            .await
615            .unwrap();
616
617        let second = fake.turn_phase(&account(), &turn_id).await.unwrap();
618        assert_eq!(second.phase, TurnPhase::Committed);
619        assert_eq!(
620            second.updated_at,
621            DateTime::<Utc>::UNIX_EPOCH + TimeDelta::minutes(5)
622        );
623        fake.set_time(DateTime::<Utc>::UNIX_EPOCH);
624        assert_eq!(fake.now(), DateTime::<Utc>::UNIX_EPOCH);
625    }
626
627    #[tokio::test]
628    async fn armed_failures_can_be_disarmed_again() {
629        let fake = FakeStores::default();
630        fake.fail_at(FailurePoint::BeforeJournalInsert, StoreError::Timeout)
631            .unwrap();
632        assert_eq!(fake.armed_failures().len(), 1);
633        fake.clear_failures();
634        assert!(fake.armed_failures().is_empty());
635    }
636
637    #[tokio::test]
638    async fn an_empty_case_simply_has_no_events() {
639        let fake = FakeStores::new();
640        let case_key = CaseKey {
641            workflow: turnframe_core::ids::WorkflowKey::from("trip"),
642            case_id: turnframe_core::ids::CaseId::from("trip-1"),
643        };
644        assert!(
645            fake.events_for(&account(), &case_key)
646                .await
647                .unwrap()
648                .is_empty()
649        );
650        assert!(
651            fake.event_types(&account(), &case_key)
652                .await
653                .unwrap()
654                .is_empty()
655        );
656    }
657}