1mod 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
76pub const MAX_LISTED_EVENTS: usize = 1024;
79
80pub 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#[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#[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#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
126#[error("no crash boundary is named `{name}`")]
127pub struct UnknownBoundary {
128 pub name: String,
130}
131
132#[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 #[must_use]
148 pub fn new() -> Self {
149 Self::at(DateTime::<Utc>::UNIX_EPOCH)
150 }
151
152 #[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 .unwrap_or_else(|_| Stores::from_memory(memory.clone()));
171 Self {
172 memory,
173 counter,
174 clock,
175 stores,
176 }
177 }
178
179 #[must_use]
181 pub fn stores(&self) -> &Stores {
182 &self.stores
183 }
184
185 #[must_use]
187 pub fn memory(&self) -> &Arc<MemoryStores> {
188 &self.memory
189 }
190
191 #[must_use]
193 pub fn clock(&self) -> &Arc<ManualClock> {
194 &self.clock
195 }
196
197 #[must_use]
199 pub fn now(&self) -> DateTime<Utc> {
200 self.clock.now()
201 }
202
203 pub fn advance(&self, delta: TimeDelta) {
205 self.clock.advance(delta);
206 }
207
208 pub fn set_time(&self, instant: DateTime<Utc>) {
210 self.clock.set(instant);
211 }
212
213 pub fn fail_at(&self, point: FailurePoint, error: StoreError) -> Result<(), StoreError> {
223 self.memory.fail_next(point, error)
224 }
225
226 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 let _ = self.memory.fail_next(point, error);
238 Ok(())
239 }
240
241 #[must_use]
243 pub fn armed_failures(&self) -> Vec<FailurePoint> {
244 self.memory.armed_failures().unwrap_or_default()
245 }
246
247 pub fn clear_failures(&self) {
249 let _ = self.memory.clear_failures();
250 }
251
252 #[must_use]
256 pub fn call_count(&self, method: &str) -> usize {
257 self.counter.count(method)
258 }
259
260 #[must_use]
262 pub fn calls(&self) -> Vec<(&'static str, usize)> {
263 self.counter.counts()
264 }
265
266 #[must_use]
268 pub fn total_calls(&self) -> usize {
269 self.counter.total()
270 }
271
272 pub fn reset_calls(&self) {
275 self.counter.reset();
276 }
277
278 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 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 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 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 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 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 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 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 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 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 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 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 offers: Vec::new(),
494 }
495 }
496
497 async fn arranged() -> (FakeStores, ConversationId, TurnId) {
499 let fake = FakeStores::new();
500 let conversation = ConversationId::nil();
501 let turn_id = TurnId::nil();
502 fake.stores()
503 .conversations()
504 .create_conversation(ConversationRecord::new(conversation, account(), fake.now()))
505 .await
506 .unwrap();
507 fake.stores()
508 .conversations()
509 .append_user_turn(StoredUserTurn::new(
510 turn_input(conversation, turn_id),
511 fake.now(),
512 ))
513 .await
514 .unwrap();
515 fake.reset_calls();
516 (fake, conversation, turn_id)
517 }
518
519 #[tokio::test]
520 async fn every_call_through_the_set_is_counted_by_role_and_method() {
521 let (fake, _conversation, turn_id) = arranged().await;
522 fake.stores()
523 .conversations()
524 .turn_phase(&account(), &turn_id)
525 .await
526 .unwrap();
527 fake.stores()
528 .conversations()
529 .turn_phase(&account(), &turn_id)
530 .await
531 .unwrap();
532
533 assert_eq!(fake.call_count("conversations.turn_phase"), 2);
534 assert_eq!(fake.total_calls(), 2);
535 assert_eq!(fake.calls(), vec![("conversations.turn_phase", 2)]);
536 }
537
538 #[tokio::test]
539 async fn reading_what_was_persisted_does_not_move_the_tally() {
540 let (fake, conversation, turn_id) = arranged().await;
541 fake.stores()
542 .conversations()
543 .append_assistant_turn(&account(), assistant(turn_id, conversation))
544 .await
545 .unwrap();
546
547 let stored = fake.turn(&account(), &turn_id).await.unwrap();
548 assert!(stored.assistant.is_some());
549 assert_eq!(
550 fake.calls(),
551 vec![("conversations.append_assistant_turn", 1)],
552 "the accessor bypasses the counting layer"
553 );
554 }
555
556 #[tokio::test]
557 async fn a_named_boundary_fails_the_next_call_that_reaches_it() {
558 let (fake, conversation, turn_id) = arranged().await;
559 fake.fail_at_boundary("before_response_persistence", StoreError::Unavailable)
560 .unwrap();
561 assert_eq!(
562 fake.armed_failures(),
563 vec![FailurePoint::BeforeResponsePersistence]
564 );
565
566 let refused = fake
567 .stores()
568 .conversations()
569 .append_assistant_turn(&account(), assistant(turn_id, conversation))
570 .await;
571
572 assert_eq!(refused, Err(StoreError::Unavailable));
573 assert!(
575 fake.turn(&account(), &turn_id)
576 .await
577 .unwrap()
578 .assistant
579 .is_none()
580 );
581 assert!(fake.armed_failures().is_empty());
582 }
583
584 #[test]
585 fn the_boundary_table_covers_the_specification_and_round_trips() {
586 assert_eq!(CRASH_BOUNDARIES.len(), FailurePoint::ALL.len());
587 for point in FailurePoint::ALL {
588 let name = boundary_name(point);
589 assert_ne!(name, "unknown", "{point:?} is missing from the table");
590 assert_eq!(crash_boundary(name), Some(point));
591 }
592 assert_eq!(crash_boundary("nope"), None);
593
594 let fake = FakeStores::new();
595 assert_eq!(
596 fake.fail_at_boundary("nope", StoreError::Timeout)
597 .unwrap_err(),
598 UnknownBoundary {
599 name: "nope".to_owned()
600 }
601 );
602 }
603
604 #[tokio::test]
605 async fn the_clock_only_moves_when_the_test_moves_it() {
606 let (fake, _conversation, turn_id) = arranged().await;
607 assert_eq!(fake.now(), DateTime::<Utc>::UNIX_EPOCH);
608 let first = fake.turn_phase(&account(), &turn_id).await.unwrap();
609 assert_eq!(first.updated_at, DateTime::<Utc>::UNIX_EPOCH);
610
611 fake.advance(TimeDelta::minutes(5));
612 fake.stores()
613 .conversations()
614 .set_turn_phase(&account(), &turn_id, TurnPhase::Committed)
615 .await
616 .unwrap();
617
618 let second = fake.turn_phase(&account(), &turn_id).await.unwrap();
619 assert_eq!(second.phase, TurnPhase::Committed);
620 assert_eq!(
621 second.updated_at,
622 DateTime::<Utc>::UNIX_EPOCH + TimeDelta::minutes(5)
623 );
624 fake.set_time(DateTime::<Utc>::UNIX_EPOCH);
625 assert_eq!(fake.now(), DateTime::<Utc>::UNIX_EPOCH);
626 }
627
628 #[tokio::test]
629 async fn armed_failures_can_be_disarmed_again() {
630 let fake = FakeStores::default();
631 fake.fail_at(FailurePoint::BeforeJournalInsert, StoreError::Timeout)
632 .unwrap();
633 assert_eq!(fake.armed_failures().len(), 1);
634 fake.clear_failures();
635 assert!(fake.armed_failures().is_empty());
636 }
637
638 #[tokio::test]
639 async fn an_empty_case_simply_has_no_events() {
640 let fake = FakeStores::new();
641 let case_key = CaseKey {
642 workflow: turnframe_core::ids::WorkflowKey::from("trip"),
643 case_id: turnframe_core::ids::CaseId::from("trip-1"),
644 };
645 assert!(
646 fake.events_for(&account(), &case_key)
647 .await
648 .unwrap()
649 .is_empty()
650 );
651 assert!(
652 fake.event_types(&account(), &case_key)
653 .await
654 .unwrap()
655 .is_empty()
656 );
657 }
658}