1use async_trait::async_trait;
27use chrono::{DateTime, Utc};
28use sqlx::PgConnection;
29use sqlx::postgres::PgRow;
30use turnframe_core::case::CaseKey;
31use turnframe_core::ids::{
32 AccountId, CaseRevision, ConversationId, EventId, InteractionId, OptionId, TurnId,
33};
34use turnframe_core::interaction::{Interaction, InteractionStatus};
35use turnframe_store::error::{StoreError, invalid_record};
36use turnframe_store::interaction::{
37 InteractionReader, InteractionRecord, InteractionWriter, InvalidationReason,
38 InvalidationRecord, ResolutionOutcome,
39};
40use uuid::Uuid;
41
42use crate::codec::{column, from_json, from_label, label, now, revision_to_sql, to_json};
43use crate::error::store_error;
44use crate::store::{PgStores, commit};
45
46const RECORD_COLUMNS: &str = "interaction, status, resolved_at, resolved_option_id, \
48 resolved_by_turn, resolution_event_ids, failure_code, invalidation";
49
50const OPEN_STATUSES: &str = "('active', 'resolving')";
53const ANSWERED_STATUSES: &str = "('resolved', 'declined', 'dismissed')";
55
56pub(crate) async fn insert_interaction(
61 conn: &mut PgConnection,
62 interaction: Interaction,
63 replace_blocking: bool,
64 at: DateTime<Utc>,
65) -> Result<Vec<InteractionId>, StoreError> {
66 if interaction.status != InteractionStatus::Active {
67 return Err(invalid_record());
68 }
69 let mut invalidated = Vec::new();
70 if interaction.blocking
71 && let Some((occupant, status)) = lock_blocking_slot(
72 &mut *conn,
73 &interaction.account_id,
74 &interaction.case_ref.key(),
75 )
76 .await?
77 {
78 if !replace_blocking || status != InteractionStatus::Active {
81 return Err(StoreError::Conflict);
82 }
83 invalidate(
84 &mut *conn,
85 &interaction.account_id,
86 &occupant,
87 &InvalidationRecord {
88 reason: InvalidationReason::Superseded { by: interaction.id },
89 new_revision: None,
90 at,
91 },
92 )
93 .await?;
94 invalidated.push(occupant);
95 }
96 write_row(conn, &interaction).await?;
97 Ok(invalidated)
98}
99
100async fn lock_blocking_slot(
103 conn: &mut PgConnection,
104 account: &AccountId,
105 case: &CaseKey,
106) -> Result<Option<(InteractionId, InteractionStatus)>, StoreError> {
107 let statement = format!(
108 "SELECT interaction_id, status FROM tf_interaction
109 WHERE account_id = $1 AND workflow_key = $2 AND case_id = $3
110 AND blocking AND status IN {OPEN_STATUSES}
111 FOR UPDATE"
112 );
113 let row = sqlx::query(&statement)
114 .bind(account.as_str())
115 .bind(case.workflow.as_str())
116 .bind(case.case_id.as_str())
117 .fetch_optional(conn)
118 .await
119 .map_err(|error| store_error(&error))?;
120 let Some(row) = row else {
121 return Ok(None);
122 };
123 let id: Uuid = column(&row, "interaction_id")?;
124 let status: String = column(&row, "status")?;
125 Ok(Some((InteractionId::from(id), from_label(&status)?)))
126}
127
128async fn write_row(conn: &mut PgConnection, interaction: &Interaction) -> Result<(), StoreError> {
130 sqlx::query(
131 "INSERT INTO tf_interaction (
132 account_id, interaction_id, conversation_id, workflow_key, case_id, case_revision,
133 kind, blocking, revision_independent, payload_hash, interaction, status,
134 created_at, expires_at, resolved_at, resolved_option_id
135 ) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16)",
136 )
137 .bind(interaction.account_id.as_str())
138 .bind(interaction.id.as_uuid())
139 .bind(interaction.conversation_id.as_uuid())
140 .bind(interaction.case_ref.workflow.as_str())
141 .bind(interaction.case_ref.case_id.as_str())
142 .bind(revision_to_sql(interaction.case_ref.expected_revision)?)
143 .bind(label(&interaction.kind)?)
144 .bind(interaction.blocking)
145 .bind(interaction.revision_independent)
146 .bind(interaction.payload_hash.as_str())
147 .bind(to_json(interaction)?)
148 .bind(label(&interaction.status)?)
149 .bind(interaction.created_at)
150 .bind(interaction.expires_at)
151 .bind(interaction.resolved_at)
152 .bind(
153 interaction
154 .resolved_option_id
155 .as_ref()
156 .map(OptionId::as_str),
157 )
158 .execute(conn)
159 .await
160 .map_err(|error| store_error(&error))?;
161 Ok(())
162}
163
164async fn invalidate(
166 conn: &mut PgConnection,
167 account: &AccountId,
168 id: &InteractionId,
169 record: &InvalidationRecord,
170) -> Result<(), StoreError> {
171 let updated = sqlx::query(
172 "UPDATE tf_interaction SET status = 'invalidated', invalidation = $3
173 WHERE account_id = $1 AND interaction_id = $2 AND status = 'active'",
174 )
175 .bind(account.as_str())
176 .bind(id.as_uuid())
177 .bind(to_json(record)?)
178 .execute(conn)
179 .await
180 .map_err(|error| store_error(&error))?;
181 if updated.rows_affected() == 0 {
182 return Err(StoreError::Conflict);
183 }
184 Ok(())
185}
186
187pub(crate) async fn get_interaction(
189 conn: &mut PgConnection,
190 account: &AccountId,
191 id: &InteractionId,
192) -> Result<InteractionRecord, StoreError> {
193 read_record(conn, account, id)
194 .await?
195 .ok_or(StoreError::NotFound)
196}
197
198async fn read_record(
200 conn: &mut PgConnection,
201 account: &AccountId,
202 id: &InteractionId,
203) -> Result<Option<InteractionRecord>, StoreError> {
204 let statement = format!(
205 "SELECT {RECORD_COLUMNS} FROM tf_interaction
206 WHERE account_id = $1 AND interaction_id = $2"
207 );
208 let row = sqlx::query(&statement)
209 .bind(account.as_str())
210 .bind(id.as_uuid())
211 .fetch_optional(conn)
212 .await
213 .map_err(|error| store_error(&error))?;
214 row.as_ref().map(decode_record).transpose()
215}
216
217pub(crate) async fn list_open_for_conversation(
219 conn: &mut PgConnection,
220 account: &AccountId,
221 conversation: &ConversationId,
222) -> Result<Vec<Interaction>, StoreError> {
223 let statement = format!(
224 "SELECT {RECORD_COLUMNS} FROM tf_interaction
225 WHERE account_id = $1 AND conversation_id = $2 AND status IN {OPEN_STATUSES}
226 ORDER BY created_at, interaction_id"
227 );
228 let rows = sqlx::query(&statement)
229 .bind(account.as_str())
230 .bind(conversation.as_uuid())
231 .fetch_all(conn)
232 .await
233 .map_err(|error| store_error(&error))?;
234 decode_interactions(&rows)
235}
236
237pub(crate) async fn list_open_for_case(
239 conn: &mut PgConnection,
240 account: &AccountId,
241 case_key: &CaseKey,
242) -> Result<Vec<Interaction>, StoreError> {
243 let statement = format!(
244 "SELECT {RECORD_COLUMNS} FROM tf_interaction
245 WHERE account_id = $1 AND workflow_key = $2 AND case_id = $3
246 AND status IN {OPEN_STATUSES}
247 ORDER BY created_at, interaction_id"
248 );
249 let rows = sqlx::query(&statement)
250 .bind(account.as_str())
251 .bind(case_key.workflow.as_str())
252 .bind(case_key.case_id.as_str())
253 .fetch_all(conn)
254 .await
255 .map_err(|error| store_error(&error))?;
256 decode_interactions(&rows)
257}
258
259pub(crate) async fn blocking_answered_at(
264 conn: &mut PgConnection,
265 account: &AccountId,
266 case_key: &CaseKey,
267 revision: CaseRevision,
268) -> Result<bool, StoreError> {
269 let statement = format!(
270 "SELECT 1 FROM tf_interaction
271 WHERE account_id = $1 AND workflow_key = $2 AND case_id = $3
272 AND case_revision = $4 AND blocking AND status IN {ANSWERED_STATUSES}
273 LIMIT 1"
274 );
275 let found = sqlx::query(&statement)
276 .bind(account.as_str())
277 .bind(case_key.workflow.as_str())
278 .bind(case_key.case_id.as_str())
279 .bind(revision_to_sql(revision)?)
280 .fetch_optional(conn)
281 .await
282 .map_err(|error| store_error(&error))?;
283 Ok(found.is_some())
284}
285
286pub(crate) async fn begin_resolution(
288 conn: &mut PgConnection,
289 account: &AccountId,
290 id: &InteractionId,
291 expected_status: InteractionStatus,
292 option_id: OptionId,
293 resolved_by: TurnId,
294 at: DateTime<Utc>,
295) -> Result<InteractionRecord, StoreError> {
296 let updated = if expected_status == InteractionStatus::Active {
299 let statement = format!(
300 "UPDATE tf_interaction
301 SET status = 'resolving', resolved_option_id = $3, resolved_at = $4,
302 resolved_by_turn = $5
303 WHERE account_id = $1 AND interaction_id = $2 AND status = 'active'
304 RETURNING {RECORD_COLUMNS}"
305 );
306 sqlx::query(&statement)
307 .bind(account.as_str())
308 .bind(id.as_uuid())
309 .bind(option_id.as_str())
310 .bind(at)
311 .bind(resolved_by.as_uuid())
312 .fetch_optional(&mut *conn)
313 .await
314 .map_err(|error| store_error(&error))?
315 } else {
316 None
317 };
318 match updated {
319 Some(row) => decode_record(&row),
320 None => Err(missing_or_conflict(conn, account, id).await?),
321 }
322}
323
324pub(crate) async fn finish_resolution(
329 conn: &mut PgConnection,
330 account: &AccountId,
331 id: &InteractionId,
332 outcome: ResolutionOutcome,
333) -> Result<InteractionRecord, StoreError> {
334 let updated = settle(&mut *conn, account, id, &outcome).await?;
335 if let Some(row) = updated {
336 return decode_record(&row);
337 }
338 let Some(record) = read_record(conn, account, id).await? else {
339 return Err(StoreError::NotFound);
340 };
341 if record.status() == outcome.target_status() && already_settled(&record, &outcome) {
342 return Ok(record);
343 }
344 Err(StoreError::Conflict)
345}
346
347async fn settle(
349 conn: &mut PgConnection,
350 account: &AccountId,
351 id: &InteractionId,
352 outcome: &ResolutionOutcome,
353) -> Result<Option<PgRow>, StoreError> {
354 let statement = match outcome {
355 ResolutionOutcome::Resolved { .. } => format!(
356 "UPDATE tf_interaction
357 SET status = 'resolved', resolution_event_ids = $3, failure_code = NULL
358 WHERE account_id = $1 AND interaction_id = $2 AND status = 'resolving'
359 RETURNING {RECORD_COLUMNS}"
360 ),
361 ResolutionOutcome::Failed { .. } => format!(
362 "UPDATE tf_interaction
363 SET status = 'failed', failure_code = $3, resolution_event_ids = '{{}}'
364 WHERE account_id = $1 AND interaction_id = $2 AND status = 'resolving'
365 RETURNING {RECORD_COLUMNS}"
366 ),
367 ResolutionOutcome::RestoreActive => format!(
370 "UPDATE tf_interaction
371 SET status = 'active', resolution_event_ids = '{{}}', failure_code = NULL,
372 resolved_by_turn = NULL, resolved_option_id = NULL, resolved_at = NULL
373 WHERE account_id = $1 AND interaction_id = $2 AND status = 'resolving'
374 RETURNING {RECORD_COLUMNS}"
375 ),
376 };
377 let query = sqlx::query(&statement)
378 .bind(account.as_str())
379 .bind(id.as_uuid());
380 let query = match outcome {
381 ResolutionOutcome::Resolved { event_ids } => query.bind(
382 event_ids
383 .iter()
384 .map(|id| *id.as_uuid())
385 .collect::<Vec<Uuid>>(),
386 ),
387 ResolutionOutcome::Failed { code } => query.bind(code.clone()),
388 ResolutionOutcome::RestoreActive => query,
389 };
390 query
391 .fetch_optional(conn)
392 .await
393 .map_err(|error| store_error(&error))
394}
395
396fn already_settled(record: &InteractionRecord, outcome: &ResolutionOutcome) -> bool {
399 match outcome {
400 ResolutionOutcome::Resolved { event_ids } => &record.resolution_event_ids == event_ids,
401 ResolutionOutcome::Failed { code } => record.failure_code.as_deref() == Some(code.as_str()),
402 ResolutionOutcome::RestoreActive => {
403 record.resolved_by_turn.is_none() && record.interaction.resolved_option_id.is_none()
404 }
405 }
406}
407
408pub(crate) async fn invalidate_case_cards(
420 conn: &mut PgConnection,
421 account: &AccountId,
422 case_key: &CaseKey,
423 reason: InvalidationReason,
424 at: DateTime<Utc>,
425) -> Result<Vec<InteractionId>, StoreError> {
426 let record = InvalidationRecord {
427 reason,
428 new_revision: None,
429 at,
430 };
431 let rows = sqlx::query(
432 "UPDATE tf_interaction SET status = 'invalidated', invalidation = $4
433 WHERE account_id = $1 AND workflow_key = $2 AND case_id = $3
434 AND status = 'active'
435 RETURNING interaction_id, created_at",
436 )
437 .bind(account.as_str())
438 .bind(case_key.workflow.as_str())
439 .bind(case_key.case_id.as_str())
440 .bind(to_json(&record)?)
441 .fetch_all(conn)
442 .await
443 .map_err(|error| store_error(&error))?;
444 let mut invalidated = rows
446 .iter()
447 .map(|row| {
448 let id: Uuid = column(row, "interaction_id")?;
449 let created_at: DateTime<Utc> = column(row, "created_at")?;
450 Ok((created_at, InteractionId::from(id)))
451 })
452 .collect::<Result<Vec<_>, StoreError>>()?;
453 invalidated.sort_unstable();
454 Ok(invalidated.into_iter().map(|(_, id)| id).collect())
455}
456
457pub(crate) async fn invalidate_for_case(
458 conn: &mut PgConnection,
459 account: &AccountId,
460 case_key: &CaseKey,
461 new_revision: CaseRevision,
462 reason: InvalidationReason,
463 at: DateTime<Utc>,
464) -> Result<Vec<InteractionId>, StoreError> {
465 let record = InvalidationRecord {
466 reason,
467 new_revision: Some(new_revision),
468 at,
469 };
470 let rows = sqlx::query(
471 "UPDATE tf_interaction SET status = 'invalidated', invalidation = $5
472 WHERE account_id = $1 AND workflow_key = $2 AND case_id = $3
473 AND status = 'active'
474 AND NOT revision_independent
475 AND case_revision <> $4
476 RETURNING interaction_id, created_at",
477 )
478 .bind(account.as_str())
479 .bind(case_key.workflow.as_str())
480 .bind(case_key.case_id.as_str())
481 .bind(revision_to_sql(new_revision)?)
482 .bind(to_json(&record)?)
483 .fetch_all(conn)
484 .await
485 .map_err(|error| store_error(&error))?;
486 let mut invalidated = rows
488 .iter()
489 .map(|row| {
490 let id: Uuid = column(row, "interaction_id")?;
491 let created_at: DateTime<Utc> = column(row, "created_at")?;
492 Ok((created_at, InteractionId::from(id)))
493 })
494 .collect::<Result<Vec<_>, StoreError>>()?;
495 invalidated.sort_unstable();
496 Ok(invalidated.into_iter().map(|(_, id)| id).collect())
497}
498
499pub(crate) async fn expire_due(
501 conn: &mut PgConnection,
502 at: DateTime<Utc>,
503) -> Result<Vec<InteractionId>, StoreError> {
504 let rows = sqlx::query(
505 "UPDATE tf_interaction SET status = 'expired'
506 WHERE status = 'active' AND expires_at IS NOT NULL AND expires_at <= $1
507 RETURNING account_id, interaction_id",
508 )
509 .bind(at)
510 .fetch_all(conn)
511 .await
512 .map_err(|error| store_error(&error))?;
513 let mut expired = rows
514 .iter()
515 .map(|row| {
516 let account: String = column(row, "account_id")?;
517 let id: Uuid = column(row, "interaction_id")?;
518 Ok((account, InteractionId::from(id)))
519 })
520 .collect::<Result<Vec<_>, StoreError>>()?;
521 expired.sort_unstable();
522 Ok(expired.into_iter().map(|(_, id)| id).collect())
523}
524
525async fn missing_or_conflict(
528 conn: &mut PgConnection,
529 account: &AccountId,
530 id: &InteractionId,
531) -> Result<StoreError, StoreError> {
532 let exists =
533 sqlx::query("SELECT 1 FROM tf_interaction WHERE account_id = $1 AND interaction_id = $2")
534 .bind(account.as_str())
535 .bind(id.as_uuid())
536 .fetch_optional(conn)
537 .await
538 .map_err(|error| store_error(&error))?
539 .is_some();
540 Ok(if exists {
541 StoreError::Conflict
542 } else {
543 StoreError::NotFound
544 })
545}
546
547fn decode_record(row: &PgRow) -> Result<InteractionRecord, StoreError> {
549 let mut interaction: Interaction = from_json(column(row, "interaction")?)?;
550 let status: String = column(row, "status")?;
551 interaction.status = from_label(&status)?;
553 interaction.resolved_at = column(row, "resolved_at")?;
554 interaction.resolved_option_id =
555 column::<Option<String>>(row, "resolved_option_id")?.map(OptionId::new);
556 let resolved_by: Option<Uuid> = column(row, "resolved_by_turn")?;
557 let event_ids: Vec<Uuid> = column(row, "resolution_event_ids")?;
558 let invalidation: Option<serde_json::Value> = column(row, "invalidation")?;
559 Ok(InteractionRecord {
560 interaction,
561 resolved_by_turn: resolved_by.map(TurnId::from),
562 resolution_event_ids: event_ids.into_iter().map(EventId::from).collect(),
563 failure_code: column(row, "failure_code")?,
564 invalidation: invalidation.map(from_json).transpose()?,
565 })
566}
567
568fn decode_interactions(rows: &[PgRow]) -> Result<Vec<Interaction>, StoreError> {
570 rows.iter()
571 .map(|row| decode_record(row).map(|record| record.interaction))
572 .collect()
573}
574
575#[async_trait]
576impl InteractionReader for PgStores {
577 async fn get(
578 &self,
579 account: &AccountId,
580 id: &InteractionId,
581 ) -> Result<InteractionRecord, StoreError> {
582 let mut conn = self.connection().await?;
583 get_interaction(&mut conn, account, id).await
584 }
585
586 async fn list_open_for_conversation(
587 &self,
588 account: &AccountId,
589 conversation: &ConversationId,
590 ) -> Result<Vec<Interaction>, StoreError> {
591 let mut conn = self.connection().await?;
592 list_open_for_conversation(&mut conn, account, conversation).await
593 }
594
595 async fn list_open_for_case(
596 &self,
597 account: &AccountId,
598 case_key: &CaseKey,
599 ) -> Result<Vec<Interaction>, StoreError> {
600 let mut conn = self.connection().await?;
601 list_open_for_case(&mut conn, account, case_key).await
602 }
603
604 async fn blocking_answered_at(
605 &self,
606 account: &AccountId,
607 case_key: &CaseKey,
608 revision: CaseRevision,
609 ) -> Result<bool, StoreError> {
610 let mut conn = self.connection().await?;
611 blocking_answered_at(&mut conn, account, case_key, revision).await
612 }
613}
614
615#[async_trait]
616impl InteractionWriter for PgStores {
617 async fn insert(&self, interaction: Interaction) -> Result<(), StoreError> {
618 let mut transaction = self.transaction().await?;
619 insert_interaction(&mut transaction, interaction, false, now()).await?;
620 commit(transaction).await
621 }
622
623 async fn insert_replacing_blocking(
624 &self,
625 interaction: Interaction,
626 ) -> Result<Vec<InteractionId>, StoreError> {
627 let mut transaction = self.transaction().await?;
628 let invalidated = insert_interaction(&mut transaction, interaction, true, now()).await?;
629 commit(transaction).await?;
630 Ok(invalidated)
631 }
632
633 async fn begin_resolution(
634 &self,
635 account: &AccountId,
636 id: &InteractionId,
637 expected_status: InteractionStatus,
638 option_id: OptionId,
639 resolved_by: TurnId,
640 ) -> Result<InteractionRecord, StoreError> {
641 let mut transaction = self.transaction().await?;
642 let record = begin_resolution(
643 &mut transaction,
644 account,
645 id,
646 expected_status,
647 option_id,
648 resolved_by,
649 now(),
650 )
651 .await?;
652 commit(transaction).await?;
653 Ok(record)
654 }
655
656 async fn finish_resolution(
657 &self,
658 account: &AccountId,
659 id: &InteractionId,
660 outcome: ResolutionOutcome,
661 ) -> Result<InteractionRecord, StoreError> {
662 let mut transaction = self.transaction().await?;
663 let record = finish_resolution(&mut transaction, account, id, outcome).await?;
664 commit(transaction).await?;
665 Ok(record)
666 }
667
668 async fn invalidate_case_cards(
669 &self,
670 account: &AccountId,
671 case_key: &CaseKey,
672 reason: InvalidationReason,
673 ) -> Result<Vec<InteractionId>, StoreError> {
674 let mut transaction = self.transaction().await?;
675 let invalidated =
676 invalidate_case_cards(&mut transaction, account, case_key, reason, now()).await?;
677 commit(transaction).await?;
678 Ok(invalidated)
679 }
680
681 async fn invalidate_for_case(
682 &self,
683 account: &AccountId,
684 case_key: &CaseKey,
685 new_revision: CaseRevision,
686 reason: InvalidationReason,
687 ) -> Result<Vec<InteractionId>, StoreError> {
688 let mut transaction = self.transaction().await?;
689 let invalidated = invalidate_for_case(
690 &mut transaction,
691 account,
692 case_key,
693 new_revision,
694 reason,
695 now(),
696 )
697 .await?;
698 commit(transaction).await?;
699 Ok(invalidated)
700 }
701
702 async fn expire_due(&self, at: DateTime<Utc>) -> Result<Vec<InteractionId>, StoreError> {
703 let mut transaction = self.transaction().await?;
704 let expired = expire_due(&mut transaction, at).await?;
705 commit(transaction).await?;
706 Ok(expired)
707 }
708}