use std::fmt;
use std::sync::Arc;
use chrono::{DateTime, Utc};
use turnframe_core::case::{CaseKey, CaseRef};
use turnframe_core::command::{CommandOrigin, ResolutionChannel};
use turnframe_core::error::InteractionError;
use turnframe_core::hash::derive_uuid;
use turnframe_core::ids::{
AccountId, CaseRevision, ConversationId, EventId, InteractionId, TurnId,
};
use turnframe_core::interaction::{
AcceptedResponse, Interaction, InteractionRejection, InteractionSpec, InteractionStatus,
validate_response,
};
use turnframe_core::observe::{NoopObserver, Observer, Signal, SignalLabels};
use turnframe_core::reduce::ActiveInteractionSummary;
use turnframe_core::turn::{ActorContext, InteractionResponse};
use turnframe_store::error::StoreError;
use turnframe_store::interaction::{InteractionRecord, InteractionStore, ResolutionOutcome};
use crate::config::InteractionConfig;
const INTERACTION_ID_DOMAIN: &str = "turnframe.interaction_id.v1";
#[must_use]
pub fn derive_interaction_id(turn_id: &TurnId, key: &str) -> InteractionId {
InteractionId::from(derive_uuid(
INTERACTION_ID_DOMAIN,
&[&turn_id.to_string(), key],
))
}
#[derive(Debug, Clone, Default)]
#[non_exhaustive]
pub struct PersistedInteractions {
pub created: Vec<Interaction>,
pub invalidated: Vec<InteractionId>,
pub failed: Option<InteractionError>,
}
impl PersistedInteractions {
#[must_use]
pub fn is_complete(&self) -> bool {
self.failed.is_none()
}
#[must_use]
pub fn views(&self) -> Vec<turnframe_core::interaction::InteractionView> {
self.created.iter().map(Interaction::view).collect()
}
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct AcceptedInteraction {
pub response: AcceptedResponse,
pub record: InteractionRecord,
}
impl AcceptedInteraction {
#[must_use]
pub fn origin(&self) -> Option<CommandOrigin> {
self.response.origin()
}
#[must_use]
pub fn case_ref(&self) -> &CaseRef {
&self.response.case_ref
}
#[must_use]
pub fn interaction_id(&self) -> InteractionId {
self.response.interaction_id
}
}
#[derive(Debug, Clone, Copy)]
pub struct ResponseContext<'a> {
pub actor: &'a ActorContext,
pub conversation: &'a ConversationId,
pub turn_id: TurnId,
pub channel: ResolutionChannel,
pub current_revision: CaseRevision,
pub now: DateTime<Utc>,
}
impl<'a> ResponseContext<'a> {
#[must_use]
pub fn click(
actor: &'a ActorContext,
conversation: &'a ConversationId,
turn_id: TurnId,
current_revision: CaseRevision,
now: DateTime<Utc>,
) -> Self {
Self {
actor,
conversation,
turn_id,
channel: ResolutionChannel::Click,
current_revision,
now,
}
}
#[must_use]
pub const fn through(mut self, channel: ResolutionChannel) -> Self {
self.channel = channel;
self
}
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub enum ResponseAdmission {
Accepted(Box<AcceptedInteraction>),
AlreadyAnswered(Box<InteractionRecord>),
}
impl ResponseAdmission {
#[must_use]
pub fn accepted(&self) -> Option<&AcceptedInteraction> {
match self {
Self::Accepted(accepted) => Some(accepted),
Self::AlreadyAnswered(_) => None,
}
}
#[must_use]
pub fn is_replay(&self) -> bool {
matches!(self, Self::AlreadyAnswered(_))
}
}
#[derive(Clone)]
pub struct InteractionEngine {
store: Arc<dyn InteractionStore>,
config: InteractionConfig,
observer: Arc<dyn Observer>,
}
impl fmt::Debug for InteractionEngine {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("InteractionEngine")
.field("config", &self.config)
.finish_non_exhaustive()
}
}
impl InteractionEngine {
#[must_use]
pub fn new(store: Arc<dyn InteractionStore>, config: InteractionConfig) -> Self {
Self {
store,
config,
observer: Arc::new(NoopObserver),
}
}
#[must_use]
pub fn with_observer(mut self, observer: Arc<dyn Observer>) -> Self {
self.observer = observer;
self
}
#[must_use]
pub const fn config(&self) -> &InteractionConfig {
&self.config
}
pub async fn create(
&self,
spec: InteractionSpec,
account: &AccountId,
conversation: ConversationId,
turn_id: TurnId,
now: DateTime<Utc>,
) -> Result<(Interaction, Vec<InteractionId>), InteractionError> {
let id = derive_interaction_id(&turn_id, &spec.key);
let mut spec = spec;
if spec.expires_in.is_none()
&& let Some(ttl) = self.config.default_ttl
{
spec = spec.expires_in(ttl);
}
let interaction =
Interaction::from_spec(spec, id, account.clone(), conversation, turn_id, now)?;
let invalidated = if interaction.blocking {
self.store
.insert_replacing_blocking(interaction.clone())
.await
.map_err(persistence_failure)?
} else {
self.store
.insert(interaction.clone())
.await
.map_err(persistence_failure)?;
Vec::new()
};
Ok((interaction, invalidated))
}
pub async fn persist(
&self,
specs: &[InteractionSpec],
account: &AccountId,
conversation: ConversationId,
turn_id: TurnId,
now: DateTime<Utc>,
) -> PersistedInteractions {
let mut created = Vec::with_capacity(specs.len());
let mut invalidated = Vec::new();
let mut failed = None;
for spec in specs {
match self
.create(spec.clone(), account, conversation, turn_id, now)
.await
{
Ok((interaction, gone)) => {
created.push(interaction);
invalidated.extend(gone);
}
Err(error) => {
tracing::warn!(
target: "turnframe.interactions",
key = %spec.key,
"interaction persistence failed; the response may not mention this card"
);
failed = Some(error);
break;
}
}
}
PersistedInteractions {
created,
invalidated,
failed,
}
}
pub async fn accept(
&self,
context: ResponseContext<'_>,
response: &InteractionResponse,
) -> Result<ResponseAdmission, InteractionError> {
let ResponseContext {
actor,
conversation,
turn_id,
channel,
current_revision,
now,
} = context;
let record = match self
.store
.get(&actor.account_id, &response.interaction_id)
.await
{
Ok(record) => record,
Err(StoreError::NotFound) => {
return Err(InteractionError::Rejected(InteractionRejection::NotFound));
}
Err(_) => return Err(InteractionError::NotPersisted),
};
let accepted = match validate_response(
&record.interaction,
response,
channel,
actor,
conversation,
current_revision,
now,
) {
Ok(accepted) => accepted,
Err(InteractionRejection::AlreadyResolved { .. })
| Err(InteractionRejection::NotActive {
status: InteractionStatus::Resolving,
}) => {
return Ok(ResponseAdmission::AlreadyAnswered(Box::new(record)));
}
Err(rejection) => return Err(InteractionError::Rejected(rejection)),
};
match self
.store
.begin_resolution(
&actor.account_id,
&response.interaction_id,
InteractionStatus::Active,
accepted.option_id.clone(),
turn_id,
)
.await
{
Ok(record) => Ok(ResponseAdmission::Accepted(Box::new(AcceptedInteraction {
response: accepted,
record,
}))),
Err(StoreError::Conflict) => {
let record = self
.store
.get(&actor.account_id, &response.interaction_id)
.await
.map_err(|_| InteractionError::NotPersisted)?;
Ok(ResponseAdmission::AlreadyAnswered(Box::new(record)))
}
Err(StoreError::NotFound) => {
Err(InteractionError::Rejected(InteractionRejection::NotFound))
}
Err(_) => Err(InteractionError::NotPersisted),
}
}
pub async fn settle(
&self,
account: &AccountId,
id: &InteractionId,
outcome: ResolutionOutcome,
) -> Result<InteractionRecord, InteractionError> {
let settled = self
.store
.finish_resolution(account, id, outcome)
.await
.map_err(|error| match error {
StoreError::NotFound => InteractionError::Rejected(InteractionRejection::NotFound),
_ => InteractionError::NotPersisted,
})?;
observe_settled(self.observer.as_ref(), &settled);
Ok(settled)
}
pub async fn mark_resolved(
&self,
account: &AccountId,
id: &InteractionId,
event_ids: Vec<EventId>,
) -> Result<InteractionRecord, InteractionError> {
self.settle(account, id, ResolutionOutcome::Resolved { event_ids })
.await
}
pub async fn mark_failed(
&self,
account: &AccountId,
id: &InteractionId,
code: impl Into<String>,
) -> Result<InteractionRecord, InteractionError> {
self.settle(account, id, ResolutionOutcome::Failed { code: code.into() })
.await
}
pub async fn restore(
&self,
account: &AccountId,
id: &InteractionId,
) -> Result<InteractionRecord, InteractionError> {
self.settle(account, id, ResolutionOutcome::RestoreActive)
.await
}
pub async fn open_for_conversation(
&self,
account: &AccountId,
conversation: &ConversationId,
) -> Result<Vec<Interaction>, InteractionError> {
self.store
.list_open_for_conversation(account, conversation)
.await
.map_err(|_| InteractionError::NotPersisted)
}
pub async fn open_for_case(
&self,
account: &AccountId,
case_key: &CaseKey,
) -> Result<Vec<Interaction>, InteractionError> {
self.store
.list_open_for_case(account, case_key)
.await
.map_err(|_| InteractionError::NotPersisted)
}
pub async fn blocking_answered_at(
&self,
account: &AccountId,
case_key: &CaseKey,
revision: turnframe_core::ids::CaseRevision,
) -> Result<bool, InteractionError> {
self.store
.blocking_answered_at(account, case_key, revision)
.await
.map_err(|_| InteractionError::NotPersisted)
}
pub async fn get(
&self,
account: &AccountId,
id: &InteractionId,
) -> Result<InteractionRecord, InteractionError> {
self.store
.get(account, id)
.await
.map_err(|error| match error {
StoreError::NotFound => InteractionError::Rejected(InteractionRejection::NotFound),
_ => InteractionError::NotPersisted,
})
}
}
fn settlement_labels(card: &Interaction) -> SignalLabels {
SignalLabels::workflow(card.case_ref.workflow.clone()).with_interaction(card.kind)
}
pub(crate) fn observe_settled(observer: &dyn Observer, settled: &InteractionRecord) {
let card = &settled.interaction;
let labels = settlement_labels(card);
match card.status {
InteractionStatus::Resolved => {
observer.observe_labeled(&Signal::InteractionResolved, &labels);
}
InteractionStatus::Failed => {
let code = settled
.failure_code
.clone()
.unwrap_or_else(|| String::from("not_committed"));
observer.observe_labeled(&Signal::InteractionFailed, &labels.with_error_code(code));
}
_ => {}
}
}
pub(crate) fn observe_resolution(
observer: &dyn Observer,
card: &Interaction,
outcome: &ResolutionOutcome,
) {
let labels = settlement_labels(card);
match outcome {
ResolutionOutcome::Resolved { .. } => {
observer.observe_labeled(&Signal::InteractionResolved, &labels);
}
ResolutionOutcome::Failed { code } => {
observer.observe_labeled(
&Signal::InteractionFailed,
&labels.with_error_code(code.clone()),
);
}
ResolutionOutcome::RestoreActive => {}
}
}
fn persistence_failure(_error: StoreError) -> InteractionError {
InteractionError::NotPersisted
}
#[must_use]
pub fn summarize(interaction: &Interaction) -> ActiveInteractionSummary {
ActiveInteractionSummary {
interaction_id: interaction.id,
case_ref: interaction.case_ref.clone(),
kind: interaction.kind,
blocking: interaction.blocking,
option_ids: interaction.payload.option_ids(),
text_resolution: interaction.text_resolution.clone(),
confirms_risk: interaction.confirms_risk,
payload_hash: interaction.payload_hash.clone(),
}
}
#[must_use]
pub fn summarize_all(interactions: &[Interaction]) -> Vec<ActiveInteractionSummary> {
interactions.iter().map(summarize).collect()
}
#[cfg(test)]
mod tests {
use turnframe_core::case::CaseRef;
use turnframe_core::ids::{CaseRevision, OptionId};
use turnframe_core::interaction::{
InteractionKind, InteractionOption, InteractionPayload, StoredInteractionAction,
};
use turnframe_store::memory::MemoryStores;
use super::*;
fn engine() -> (InteractionEngine, Arc<MemoryStores>) {
let memory = Arc::new(MemoryStores::new());
let engine = InteractionEngine::new(
Arc::clone(&memory) as Arc<dyn InteractionStore>,
InteractionConfig::conservative(),
);
(engine, memory)
}
fn case() -> CaseRef {
CaseRef::new("trip", "trip-1", CaseRevision(3))
}
fn spec(key: &str) -> InteractionSpec {
InteractionSpec::new(
key,
case(),
InteractionKind::SingleSelect,
InteractionPayload::new("Which one?").with_option(InteractionOption::new(
OptionId::from("a"),
"The first",
StoredInteractionAction::ResolveClarification {
answer_key: "a".to_owned(),
},
)),
)
}
fn now() -> DateTime<Utc> {
DateTime::from_timestamp(1_700_000_000, 0).expect("a valid fixed instant")
}
#[tokio::test]
async fn a_second_blocking_card_replaces_and_invalidates_the_first() {
let (engine, _memory) = engine();
let account = AccountId::from("aurora");
let conversation = ConversationId::nil();
let (first, gone) = engine
.create(
spec("first"),
&account,
conversation,
TurnId::from(uuid::Uuid::from_u128(1)),
now(),
)
.await
.expect("the slot was free");
assert!(gone.is_empty());
let (second, gone) = engine
.create(
spec("second"),
&account,
conversation,
TurnId::from(uuid::Uuid::from_u128(2)),
now(),
)
.await
.expect("the occupant is replaceable");
assert_eq!(
gone,
vec![first.id],
"the previous blocking card was invalidated (I5, §15.6)"
);
assert_ne!(second.id, first.id, "and the replacement is a new card");
let open = engine
.open_for_case(&account, &case().key())
.await
.expect("the store answers");
assert_eq!(open.len(), 1);
assert_eq!(open[0].id, second.id);
}
#[tokio::test]
async fn a_card_is_resolved_only_once_its_command_committed() {
let (engine, _memory) = engine();
let account = AccountId::from("aurora");
let conversation = ConversationId::nil();
let turn = TurnId::from(uuid::Uuid::from_u128(1));
let (card, _) = engine
.create(spec("first"), &account, conversation, turn, now())
.await
.expect("the card is written");
let actor = ActorContext::new(account.clone(), "u1");
let response = InteractionResponse {
interaction_id: card.id,
option_id: OptionId::from("a"),
expected_case_revision: CaseRevision(3),
freeform_input: None,
};
let admitted = engine
.accept(
ResponseContext::click(&actor, &conversation, turn, CaseRevision(3), now()),
&response,
)
.await
.expect("the answer is valid");
let accepted = admitted.accepted().expect("it was accepted");
assert_eq!(
accepted.record.status(),
InteractionStatus::Resolving,
"answering starts a resolution; it does not finish one"
);
let again = engine
.accept(
ResponseContext::click(
&actor,
&conversation,
TurnId::from(uuid::Uuid::from_u128(2)),
CaseRevision(3),
now(),
),
&response,
)
.await
.expect("the second click is answered, not accepted");
assert!(again.is_replay());
let settled = engine
.mark_resolved(&account, &card.id, Vec::new())
.await
.expect("the command committed");
assert_eq!(settled.status(), InteractionStatus::Resolved);
}
#[tokio::test]
async fn a_card_of_another_tenant_is_simply_not_there() {
let (engine, _memory) = engine();
let (card, _) = engine
.create(
spec("first"),
&AccountId::from("aurora"),
ConversationId::nil(),
TurnId::nil(),
now(),
)
.await
.expect("the card is written");
let stranger = engine.get(&AccountId::from("other"), &card.id).await;
let unknown = engine
.get(
&AccountId::from("other"),
&InteractionId::from(uuid::Uuid::from_u128(999)),
)
.await;
assert_eq!(
format!("{stranger:?}"),
format!("{unknown:?}"),
"another tenant's card and one that never existed must be one answer (§25.4)"
);
}
#[test]
fn interaction_ids_are_derived_and_stable() {
let turn = TurnId::nil();
assert_eq!(
derive_interaction_id(&turn, "confirm:acts[0]"),
derive_interaction_id(&turn, "confirm:acts[0]"),
);
assert_ne!(
derive_interaction_id(&turn, "confirm:acts[0]"),
derive_interaction_id(&turn, "confirm:acts[1]"),
);
}
}