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 }
494 }
495
496 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 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}