1use std::fmt;
39use std::sync::Arc;
40
41use chrono::{DateTime, Utc};
42use turnframe_core::case::{CaseKey, CaseRef};
43use turnframe_core::command::{CommandOrigin, ResolutionChannel};
44use turnframe_core::error::InteractionError;
45use turnframe_core::hash::derive_uuid;
46use turnframe_core::ids::{
47 AccountId, CaseRevision, ConversationId, EventId, InteractionId, TurnId,
48};
49use turnframe_core::interaction::{
50 AcceptedResponse, Interaction, InteractionRejection, InteractionSpec, InteractionStatus,
51 validate_response,
52};
53use turnframe_core::observe::{NoopObserver, Observer, Signal, SignalLabels};
54use turnframe_core::reduce::ActiveInteractionSummary;
55use turnframe_core::turn::{ActorContext, InteractionResponse};
56use turnframe_store::error::StoreError;
57use turnframe_store::interaction::{InteractionRecord, InteractionStore, ResolutionOutcome};
58
59use crate::config::InteractionConfig;
60
61const INTERACTION_ID_DOMAIN: &str = "turnframe.interaction_id.v1";
63
64#[must_use]
70pub fn derive_interaction_id(turn_id: &TurnId, key: &str) -> InteractionId {
71 InteractionId::from(derive_uuid(
72 INTERACTION_ID_DOMAIN,
73 &[&turn_id.to_string(), key],
74 ))
75}
76
77#[derive(Debug, Clone, Default)]
84#[non_exhaustive]
85pub struct PersistedInteractions {
86 pub created: Vec<Interaction>,
88 pub invalidated: Vec<InteractionId>,
90 pub failed: Option<InteractionError>,
92}
93
94impl PersistedInteractions {
95 #[must_use]
97 pub fn is_complete(&self) -> bool {
98 self.failed.is_none()
99 }
100
101 #[must_use]
103 pub fn views(&self) -> Vec<turnframe_core::interaction::InteractionView> {
104 self.created.iter().map(Interaction::view).collect()
105 }
106}
107
108#[derive(Debug, Clone)]
111#[non_exhaustive]
112pub struct AcceptedInteraction {
113 pub response: AcceptedResponse,
115 pub record: InteractionRecord,
117}
118
119impl AcceptedInteraction {
120 #[must_use]
123 pub fn origin(&self) -> Option<CommandOrigin> {
124 self.response.origin()
125 }
126
127 #[must_use]
129 pub fn case_ref(&self) -> &CaseRef {
130 &self.response.case_ref
131 }
132
133 #[must_use]
135 pub fn interaction_id(&self) -> InteractionId {
136 self.response.interaction_id
137 }
138}
139
140#[derive(Debug, Clone, Copy)]
147pub struct ResponseContext<'a> {
148 pub actor: &'a ActorContext,
150 pub conversation: &'a ConversationId,
152 pub turn_id: TurnId,
154 pub channel: ResolutionChannel,
156 pub current_revision: CaseRevision,
158 pub now: DateTime<Utc>,
160}
161
162impl<'a> ResponseContext<'a> {
163 #[must_use]
165 pub fn click(
166 actor: &'a ActorContext,
167 conversation: &'a ConversationId,
168 turn_id: TurnId,
169 current_revision: CaseRevision,
170 now: DateTime<Utc>,
171 ) -> Self {
172 Self {
173 actor,
174 conversation,
175 turn_id,
176 channel: ResolutionChannel::Click,
177 current_revision,
178 now,
179 }
180 }
181
182 #[must_use]
184 pub const fn through(mut self, channel: ResolutionChannel) -> Self {
185 self.channel = channel;
186 self
187 }
188}
189
190#[derive(Debug, Clone)]
192#[non_exhaustive]
193pub enum ResponseAdmission {
194 Accepted(Box<AcceptedInteraction>),
196 AlreadyAnswered(Box<InteractionRecord>),
202}
203
204impl ResponseAdmission {
205 #[must_use]
207 pub fn accepted(&self) -> Option<&AcceptedInteraction> {
208 match self {
209 Self::Accepted(accepted) => Some(accepted),
210 Self::AlreadyAnswered(_) => None,
211 }
212 }
213
214 #[must_use]
216 pub fn is_replay(&self) -> bool {
217 matches!(self, Self::AlreadyAnswered(_))
218 }
219}
220
221#[derive(Clone)]
223pub struct InteractionEngine {
224 store: Arc<dyn InteractionStore>,
225 config: InteractionConfig,
226 observer: Arc<dyn Observer>,
227}
228
229impl fmt::Debug for InteractionEngine {
230 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
231 f.debug_struct("InteractionEngine")
232 .field("config", &self.config)
233 .finish_non_exhaustive()
234 }
235}
236
237impl InteractionEngine {
238 #[must_use]
240 pub fn new(store: Arc<dyn InteractionStore>, config: InteractionConfig) -> Self {
241 Self {
242 store,
243 config,
244 observer: Arc::new(NoopObserver),
245 }
246 }
247
248 #[must_use]
255 pub fn with_observer(mut self, observer: Arc<dyn Observer>) -> Self {
256 self.observer = observer;
257 self
258 }
259
260 #[must_use]
262 pub const fn config(&self) -> &InteractionConfig {
263 &self.config
264 }
265
266 pub async fn create(
279 &self,
280 spec: InteractionSpec,
281 account: &AccountId,
282 conversation: ConversationId,
283 turn_id: TurnId,
284 now: DateTime<Utc>,
285 ) -> Result<(Interaction, Vec<InteractionId>), InteractionError> {
286 let id = derive_interaction_id(&turn_id, &spec.key);
287 let mut spec = spec;
288 if spec.expires_in.is_none()
289 && let Some(ttl) = self.config.default_ttl
290 {
291 spec = spec.expires_in(ttl);
292 }
293 let interaction =
294 Interaction::from_spec(spec, id, account.clone(), conversation, turn_id, now)?;
295 let invalidated = if interaction.blocking {
296 self.store
297 .insert_replacing_blocking(interaction.clone())
298 .await
299 .map_err(persistence_failure)?
300 } else {
301 self.store
302 .insert(interaction.clone())
303 .await
304 .map_err(persistence_failure)?;
305 Vec::new()
306 };
307 Ok((interaction, invalidated))
308 }
309
310 pub async fn persist(
316 &self,
317 specs: &[InteractionSpec],
318 account: &AccountId,
319 conversation: ConversationId,
320 turn_id: TurnId,
321 now: DateTime<Utc>,
322 ) -> PersistedInteractions {
323 let mut created = Vec::with_capacity(specs.len());
324 let mut invalidated = Vec::new();
325 let mut failed = None;
326 for spec in specs {
327 match self
328 .create(spec.clone(), account, conversation, turn_id, now)
329 .await
330 {
331 Ok((interaction, gone)) => {
332 created.push(interaction);
333 invalidated.extend(gone);
334 }
335 Err(error) => {
336 tracing::warn!(
337 target: "turnframe.interactions",
338 key = %spec.key,
339 "interaction persistence failed; the response may not mention this card"
340 );
341 failed = Some(error);
342 break;
343 }
344 }
345 }
346 PersistedInteractions {
347 created,
348 invalidated,
349 failed,
350 }
351 }
352
353 pub async fn accept(
366 &self,
367 context: ResponseContext<'_>,
368 response: &InteractionResponse,
369 ) -> Result<ResponseAdmission, InteractionError> {
370 let ResponseContext {
371 actor,
372 conversation,
373 turn_id,
374 channel,
375 current_revision,
376 now,
377 } = context;
378 let record = match self
379 .store
380 .get(&actor.account_id, &response.interaction_id)
381 .await
382 {
383 Ok(record) => record,
384 Err(StoreError::NotFound) => {
387 return Err(InteractionError::Rejected(InteractionRejection::NotFound));
388 }
389 Err(_) => return Err(InteractionError::NotPersisted),
390 };
391 let accepted = match validate_response(
392 &record.interaction,
393 response,
394 channel,
395 actor,
396 conversation,
397 current_revision,
398 now,
399 ) {
400 Ok(accepted) => accepted,
401 Err(InteractionRejection::AlreadyResolved { .. })
405 | Err(InteractionRejection::NotActive {
406 status: InteractionStatus::Resolving,
407 }) => {
408 return Ok(ResponseAdmission::AlreadyAnswered(Box::new(record)));
409 }
410 Err(rejection) => return Err(InteractionError::Rejected(rejection)),
411 };
412 match self
413 .store
414 .begin_resolution(
415 &actor.account_id,
416 &response.interaction_id,
417 InteractionStatus::Active,
418 accepted.option_id.clone(),
419 turn_id,
420 )
421 .await
422 {
423 Ok(record) => Ok(ResponseAdmission::Accepted(Box::new(AcceptedInteraction {
424 response: accepted,
425 record,
426 }))),
427 Err(StoreError::Conflict) => {
431 let record = self
432 .store
433 .get(&actor.account_id, &response.interaction_id)
434 .await
435 .map_err(|_| InteractionError::NotPersisted)?;
436 Ok(ResponseAdmission::AlreadyAnswered(Box::new(record)))
437 }
438 Err(StoreError::NotFound) => {
439 Err(InteractionError::Rejected(InteractionRejection::NotFound))
440 }
441 Err(_) => Err(InteractionError::NotPersisted),
442 }
443 }
444
445 pub async fn settle(
453 &self,
454 account: &AccountId,
455 id: &InteractionId,
456 outcome: ResolutionOutcome,
457 ) -> Result<InteractionRecord, InteractionError> {
458 let settled = self
459 .store
460 .finish_resolution(account, id, outcome)
461 .await
462 .map_err(|error| match error {
463 StoreError::NotFound => InteractionError::Rejected(InteractionRejection::NotFound),
464 _ => InteractionError::NotPersisted,
465 })?;
466 observe_settled(self.observer.as_ref(), &settled);
469 Ok(settled)
470 }
471
472 pub async fn mark_resolved(
479 &self,
480 account: &AccountId,
481 id: &InteractionId,
482 event_ids: Vec<EventId>,
483 ) -> Result<InteractionRecord, InteractionError> {
484 self.settle(account, id, ResolutionOutcome::Resolved { event_ids })
485 .await
486 }
487
488 pub async fn mark_failed(
494 &self,
495 account: &AccountId,
496 id: &InteractionId,
497 code: impl Into<String>,
498 ) -> Result<InteractionRecord, InteractionError> {
499 self.settle(account, id, ResolutionOutcome::Failed { code: code.into() })
500 .await
501 }
502
503 pub async fn restore(
509 &self,
510 account: &AccountId,
511 id: &InteractionId,
512 ) -> Result<InteractionRecord, InteractionError> {
513 self.settle(account, id, ResolutionOutcome::RestoreActive)
514 .await
515 }
516
517 pub async fn open_for_conversation(
525 &self,
526 account: &AccountId,
527 conversation: &ConversationId,
528 ) -> Result<Vec<Interaction>, InteractionError> {
529 self.store
530 .list_open_for_conversation(account, conversation)
531 .await
532 .map_err(|_| InteractionError::NotPersisted)
533 }
534
535 pub async fn open_for_case(
541 &self,
542 account: &AccountId,
543 case_key: &CaseKey,
544 ) -> Result<Vec<Interaction>, InteractionError> {
545 self.store
546 .list_open_for_case(account, case_key)
547 .await
548 .map_err(|_| InteractionError::NotPersisted)
549 }
550
551 pub async fn blocking_answered_at(
561 &self,
562 account: &AccountId,
563 case_key: &CaseKey,
564 revision: turnframe_core::ids::CaseRevision,
565 ) -> Result<bool, InteractionError> {
566 self.store
567 .blocking_answered_at(account, case_key, revision)
568 .await
569 .map_err(|_| InteractionError::NotPersisted)
570 }
571
572 pub async fn get(
579 &self,
580 account: &AccountId,
581 id: &InteractionId,
582 ) -> Result<InteractionRecord, InteractionError> {
583 self.store
584 .get(account, id)
585 .await
586 .map_err(|error| match error {
587 StoreError::NotFound => InteractionError::Rejected(InteractionRejection::NotFound),
588 _ => InteractionError::NotPersisted,
589 })
590 }
591}
592
593fn settlement_labels(card: &Interaction) -> SignalLabels {
596 SignalLabels::workflow(card.case_ref.workflow.clone()).with_interaction(card.kind)
597}
598
599pub(crate) fn observe_settled(observer: &dyn Observer, settled: &InteractionRecord) {
608 let card = &settled.interaction;
609 let labels = settlement_labels(card);
610 match card.status {
611 InteractionStatus::Resolved => {
612 observer.observe_labeled(&Signal::InteractionResolved, &labels);
613 }
614 InteractionStatus::Failed => {
615 let code = settled
616 .failure_code
617 .clone()
618 .unwrap_or_else(|| String::from("not_committed"));
619 observer.observe_labeled(&Signal::InteractionFailed, &labels.with_error_code(code));
620 }
621 _ => {}
622 }
623}
624
625pub(crate) fn observe_resolution(
631 observer: &dyn Observer,
632 card: &Interaction,
633 outcome: &ResolutionOutcome,
634) {
635 let labels = settlement_labels(card);
636 match outcome {
637 ResolutionOutcome::Resolved { .. } => {
638 observer.observe_labeled(&Signal::InteractionResolved, &labels);
639 }
640 ResolutionOutcome::Failed { code } => {
641 observer.observe_labeled(
642 &Signal::InteractionFailed,
643 &labels.with_error_code(code.clone()),
644 );
645 }
646 ResolutionOutcome::RestoreActive => {}
647 }
648}
649
650fn persistence_failure(_error: StoreError) -> InteractionError {
653 InteractionError::NotPersisted
654}
655
656#[must_use]
658pub fn summarize(interaction: &Interaction) -> ActiveInteractionSummary {
659 ActiveInteractionSummary {
660 interaction_id: interaction.id,
661 case_ref: interaction.case_ref.clone(),
662 kind: interaction.kind,
663 blocking: interaction.blocking,
664 option_ids: interaction.payload.option_ids(),
665 text_resolution: interaction.text_resolution.clone(),
666 confirms_risk: interaction.confirms_risk,
667 payload_hash: interaction.payload_hash.clone(),
668 }
669}
670
671#[must_use]
673pub fn summarize_all(interactions: &[Interaction]) -> Vec<ActiveInteractionSummary> {
674 interactions.iter().map(summarize).collect()
675}
676
677#[cfg(test)]
678mod tests {
679 use turnframe_core::case::CaseRef;
680 use turnframe_core::ids::{CaseRevision, OptionId};
681 use turnframe_core::interaction::{
682 InteractionKind, InteractionOption, InteractionPayload, StoredInteractionAction,
683 };
684 use turnframe_store::memory::MemoryStores;
685
686 use super::*;
687
688 fn engine() -> (InteractionEngine, Arc<MemoryStores>) {
689 let memory = Arc::new(MemoryStores::new());
690 let engine = InteractionEngine::new(
691 Arc::clone(&memory) as Arc<dyn InteractionStore>,
692 InteractionConfig::conservative(),
693 );
694 (engine, memory)
695 }
696
697 fn case() -> CaseRef {
698 CaseRef::new("trip", "trip-1", CaseRevision(3))
699 }
700
701 fn spec(key: &str) -> InteractionSpec {
702 InteractionSpec::new(
703 key,
704 case(),
705 InteractionKind::SingleSelect,
706 InteractionPayload::new("Which one?").with_option(InteractionOption::new(
707 OptionId::from("a"),
708 "The first",
709 StoredInteractionAction::ResolveClarification {
710 answer_key: "a".to_owned(),
711 },
712 )),
713 )
714 }
715
716 fn now() -> DateTime<Utc> {
717 DateTime::from_timestamp(1_700_000_000, 0).expect("a valid fixed instant")
718 }
719
720 #[tokio::test]
721 async fn a_second_blocking_card_replaces_and_invalidates_the_first() {
722 let (engine, _memory) = engine();
723 let account = AccountId::from("aurora");
724 let conversation = ConversationId::nil();
725
726 let (first, gone) = engine
727 .create(
728 spec("first"),
729 &account,
730 conversation,
731 TurnId::from(uuid::Uuid::from_u128(1)),
732 now(),
733 )
734 .await
735 .expect("the slot was free");
736 assert!(gone.is_empty());
737
738 let (second, gone) = engine
739 .create(
740 spec("second"),
741 &account,
742 conversation,
743 TurnId::from(uuid::Uuid::from_u128(2)),
744 now(),
745 )
746 .await
747 .expect("the occupant is replaceable");
748 assert_eq!(
749 gone,
750 vec![first.id],
751 "the previous blocking card was invalidated (I5, §15.6)"
752 );
753 assert_ne!(second.id, first.id, "and the replacement is a new card");
754
755 let open = engine
756 .open_for_case(&account, &case().key())
757 .await
758 .expect("the store answers");
759 assert_eq!(open.len(), 1);
760 assert_eq!(open[0].id, second.id);
761 }
762
763 #[tokio::test]
764 async fn a_card_is_resolved_only_once_its_command_committed() {
765 let (engine, _memory) = engine();
766 let account = AccountId::from("aurora");
767 let conversation = ConversationId::nil();
768 let turn = TurnId::from(uuid::Uuid::from_u128(1));
769 let (card, _) = engine
770 .create(spec("first"), &account, conversation, turn, now())
771 .await
772 .expect("the card is written");
773
774 let actor = ActorContext::new(account.clone(), "u1");
775 let response = InteractionResponse {
776 interaction_id: card.id,
777 option_id: OptionId::from("a"),
778 expected_case_revision: CaseRevision(3),
779 freeform_input: None,
780 };
781 let admitted = engine
782 .accept(
783 ResponseContext::click(&actor, &conversation, turn, CaseRevision(3), now()),
784 &response,
785 )
786 .await
787 .expect("the answer is valid");
788 let accepted = admitted.accepted().expect("it was accepted");
789 assert_eq!(
790 accepted.record.status(),
791 InteractionStatus::Resolving,
792 "answering starts a resolution; it does not finish one"
793 );
794
795 let again = engine
798 .accept(
799 ResponseContext::click(
800 &actor,
801 &conversation,
802 TurnId::from(uuid::Uuid::from_u128(2)),
803 CaseRevision(3),
804 now(),
805 ),
806 &response,
807 )
808 .await
809 .expect("the second click is answered, not accepted");
810 assert!(again.is_replay());
811
812 let settled = engine
813 .mark_resolved(&account, &card.id, Vec::new())
814 .await
815 .expect("the command committed");
816 assert_eq!(settled.status(), InteractionStatus::Resolved);
817 }
818
819 #[tokio::test]
820 async fn a_card_of_another_tenant_is_simply_not_there() {
821 let (engine, _memory) = engine();
822 let (card, _) = engine
823 .create(
824 spec("first"),
825 &AccountId::from("aurora"),
826 ConversationId::nil(),
827 TurnId::nil(),
828 now(),
829 )
830 .await
831 .expect("the card is written");
832
833 let stranger = engine.get(&AccountId::from("other"), &card.id).await;
834 let unknown = engine
835 .get(
836 &AccountId::from("other"),
837 &InteractionId::from(uuid::Uuid::from_u128(999)),
838 )
839 .await;
840 assert_eq!(
841 format!("{stranger:?}"),
842 format!("{unknown:?}"),
843 "another tenant's card and one that never existed must be one answer (§25.4)"
844 );
845 }
846
847 #[test]
848 fn interaction_ids_are_derived_and_stable() {
849 let turn = TurnId::nil();
850 assert_eq!(
851 derive_interaction_id(&turn, "confirm:acts[0]"),
852 derive_interaction_id(&turn, "confirm:acts[0]"),
853 );
854 assert_ne!(
855 derive_interaction_id(&turn, "confirm:acts[0]"),
856 derive_interaction_id(&turn, "confirm:acts[1]"),
857 );
858 }
859}