1use async_trait::async_trait;
22use chrono::{DateTime, Utc};
23use sqlx::PgConnection;
24use sqlx::postgres::PgRow;
25use turnframe_core::command::IdempotencyKey;
26use turnframe_core::event::{OutboxEntry, OutboxStatus};
27use turnframe_core::ids::{CommandId, OutboxId};
28use turnframe_store::error::{StoreError, invalid_record};
29use turnframe_store::outbox::{OutboxClaim, OutboxReader, OutboxRecord, OutboxWriter};
30use uuid::Uuid;
31
32use crate::codec::{
33 attempts_from_sql, attempts_to_sql, column, from_label, label, limit_to_sql, now, to_json,
34};
35use crate::error::store_error;
36use crate::store::{PgStores, commit};
37
38const RECORD_COLUMNS: &str = "outbox_id, command_id, destination, payload, idempotency_key, \
40 status, attempt_count, next_attempt_at, created_at, completed_at, claim_worker_id, \
41 claim_taken_at, last_failure, remote_ref";
42
43const CLAIMED_COLUMNS: &str = "o.outbox_id, o.command_id, o.destination, o.payload, \
46 o.idempotency_key, o.status, o.attempt_count, o.next_attempt_at, o.created_at, \
47 o.completed_at, o.claim_worker_id, o.claim_taken_at, o.last_failure, o.remote_ref";
48
49const TERMINAL_STATUSES: &str = "('completed', 'failed')";
51
52pub(crate) async fn enqueue(conn: &mut PgConnection, entry: OutboxEntry) -> Result<(), StoreError> {
54 if entry.status != OutboxStatus::Pending {
55 return Err(invalid_record());
56 }
57 sqlx::query(
58 "INSERT INTO tf_outbox (
59 outbox_id, command_id, destination, payload, idempotency_key, status,
60 attempt_count, next_attempt_at, created_at, completed_at
61 ) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10)",
62 )
63 .bind(entry.outbox_id.as_uuid())
64 .bind(entry.command_id.as_uuid())
65 .bind(entry.destination.as_str())
66 .bind(to_json(&entry.payload)?)
67 .bind(entry.idempotency_key.as_str())
68 .bind(label(&entry.status)?)
69 .bind(attempts_to_sql(entry.attempt_count)?)
70 .bind(entry.next_attempt_at)
71 .bind(entry.created_at)
72 .bind(entry.completed_at)
73 .execute(conn)
74 .await
75 .map_err(|error| store_error(&error))?;
76 Ok(())
77}
78
79pub(crate) async fn get(
81 conn: &mut PgConnection,
82 outbox_id: &OutboxId,
83) -> Result<OutboxRecord, StoreError> {
84 let statement = format!("SELECT {RECORD_COLUMNS} FROM tf_outbox WHERE outbox_id = $1");
85 let row = sqlx::query(&statement)
86 .bind(outbox_id.as_uuid())
87 .fetch_optional(conn)
88 .await
89 .map_err(|error| store_error(&error))?
90 .ok_or(StoreError::NotFound)?;
91 decode_record(&row)
92}
93
94pub(crate) async fn list_for_command(
96 conn: &mut PgConnection,
97 command_id: &CommandId,
98) -> Result<Vec<OutboxRecord>, StoreError> {
99 let statement = format!(
100 "SELECT {RECORD_COLUMNS} FROM tf_outbox
101 WHERE command_id = $1
102 ORDER BY created_at, outbox_id"
103 );
104 let rows = sqlx::query(&statement)
105 .bind(command_id.as_uuid())
106 .fetch_all(conn)
107 .await
108 .map_err(|error| store_error(&error))?;
109 rows.iter().map(decode_record).collect()
110}
111
112pub(crate) async fn claim_due(
117 conn: &mut PgConnection,
118 at: DateTime<Utc>,
119 limit: usize,
120 worker_id: &str,
121) -> Result<Vec<OutboxEntry>, StoreError> {
122 let statement = format!(
123 "WITH due AS (
124 SELECT outbox_id FROM tf_outbox
125 WHERE status = 'pending' AND (next_attempt_at IS NULL OR next_attempt_at <= $1)
126 ORDER BY created_at, outbox_id
127 LIMIT $2
128 FOR UPDATE SKIP LOCKED
129 )
130 UPDATE tf_outbox o
131 SET status = 'dispatching',
132 attempt_count = o.attempt_count + 1,
133 claim_worker_id = $3,
134 claim_taken_at = $1
135 FROM due
136 WHERE o.outbox_id = due.outbox_id
137 RETURNING {CLAIMED_COLUMNS}"
138 );
139 let rows = sqlx::query(&statement)
140 .bind(at)
141 .bind(limit_to_sql(limit))
142 .bind(worker_id)
143 .fetch_all(conn)
144 .await
145 .map_err(|error| store_error(&error))?;
146 let mut claimed = rows
147 .iter()
148 .map(decode_record)
149 .collect::<Result<Vec<OutboxRecord>, StoreError>>()?;
150 claimed.sort_by(|a, b| {
152 a.entry
153 .created_at
154 .cmp(&b.entry.created_at)
155 .then_with(|| a.entry.outbox_id.cmp(&b.entry.outbox_id))
156 });
157 Ok(claimed.into_iter().map(|record| record.entry).collect())
158}
159
160pub(crate) async fn mark_completed(
162 conn: &mut PgConnection,
163 outbox_id: &OutboxId,
164 at: DateTime<Utc>,
165) -> Result<(), StoreError> {
166 let updated = sqlx::query(
167 "UPDATE tf_outbox
168 SET status = 'completed', completed_at = $2, claim_worker_id = NULL, claim_taken_at = NULL
169 WHERE outbox_id = $1 AND status IN ('dispatching', 'outcome_unknown')",
170 )
171 .bind(outbox_id.as_uuid())
172 .bind(at)
173 .execute(&mut *conn)
174 .await
175 .map_err(|error| store_error(&error))?;
176 if updated.rows_affected() == 1 {
177 return Ok(());
178 }
179 accept_if(conn, outbox_id, OutboxStatus::Completed).await
180}
181
182pub(crate) async fn mark_failed(
184 conn: &mut PgConnection,
185 outbox_id: &OutboxId,
186 reason: String,
187 retry_at: Option<DateTime<Utc>>,
188 at: DateTime<Utc>,
189) -> Result<(), StoreError> {
190 match retry_at {
191 Some(retry_at) => requeue(conn, outbox_id, Some(reason), retry_at).await,
192 None => fail(conn, outbox_id, reason, at).await,
193 }
194}
195
196async fn requeue(
198 conn: &mut PgConnection,
199 outbox_id: &OutboxId,
200 reason: Option<String>,
201 retry_at: DateTime<Utc>,
202) -> Result<(), StoreError> {
203 let statement = format!(
204 "UPDATE tf_outbox
205 SET status = 'pending', next_attempt_at = $2,
206 last_failure = COALESCE($3, last_failure),
207 claim_worker_id = NULL, claim_taken_at = NULL
208 WHERE outbox_id = $1 AND status NOT IN {TERMINAL_STATUSES}"
209 );
210 let updated = sqlx::query(&statement)
211 .bind(outbox_id.as_uuid())
212 .bind(retry_at)
213 .bind(reason)
214 .execute(&mut *conn)
215 .await
216 .map_err(|error| store_error(&error))?;
217 if updated.rows_affected() == 1 {
218 return Ok(());
219 }
220 Err(missing_or_conflict(conn, outbox_id).await?)
221}
222
223async fn fail(
225 conn: &mut PgConnection,
226 outbox_id: &OutboxId,
227 reason: String,
228 at: DateTime<Utc>,
229) -> Result<(), StoreError> {
230 let updated = sqlx::query(
231 "UPDATE tf_outbox
232 SET status = 'failed', completed_at = $2, last_failure = $3,
233 claim_worker_id = NULL, claim_taken_at = NULL
234 WHERE outbox_id = $1 AND status IN ('dispatching', 'outcome_unknown')",
235 )
236 .bind(outbox_id.as_uuid())
237 .bind(at)
238 .bind(reason.as_str())
239 .execute(&mut *conn)
240 .await
241 .map_err(|error| store_error(&error))?;
242 if updated.rows_affected() == 1 {
243 return Ok(());
244 }
245 let repeated = sqlx::query(
247 "UPDATE tf_outbox SET last_failure = $2 WHERE outbox_id = $1 AND status = 'failed'",
248 )
249 .bind(outbox_id.as_uuid())
250 .bind(reason.as_str())
251 .execute(&mut *conn)
252 .await
253 .map_err(|error| store_error(&error))?;
254 if repeated.rows_affected() == 1 {
255 return Ok(());
256 }
257 Err(missing_or_conflict(conn, outbox_id).await?)
258}
259
260pub(crate) async fn mark_outcome_unknown(
262 conn: &mut PgConnection,
263 outbox_id: &OutboxId,
264 remote_ref: Option<String>,
265) -> Result<(), StoreError> {
266 let updated = sqlx::query(
267 "UPDATE tf_outbox
268 SET status = 'outcome_unknown', remote_ref = $2,
269 claim_worker_id = NULL, claim_taken_at = NULL
270 WHERE outbox_id = $1 AND status = 'dispatching'",
271 )
272 .bind(outbox_id.as_uuid())
273 .bind(remote_ref.as_deref())
274 .execute(&mut *conn)
275 .await
276 .map_err(|error| store_error(&error))?;
277 if updated.rows_affected() == 1 {
278 return Ok(());
279 }
280 let repeated = sqlx::query(
283 "UPDATE tf_outbox SET remote_ref = COALESCE(remote_ref, $2)
284 WHERE outbox_id = $1 AND status = 'outcome_unknown'",
285 )
286 .bind(outbox_id.as_uuid())
287 .bind(remote_ref.as_deref())
288 .execute(&mut *conn)
289 .await
290 .map_err(|error| store_error(&error))?;
291 if repeated.rows_affected() == 1 {
292 return Ok(());
293 }
294 Err(missing_or_conflict(conn, outbox_id).await?)
295}
296
297pub(crate) async fn reschedule(
299 conn: &mut PgConnection,
300 outbox_id: &OutboxId,
301 next_attempt_at: DateTime<Utc>,
302) -> Result<(), StoreError> {
303 requeue(conn, outbox_id, None, next_attempt_at).await
304}
305
306pub(crate) async fn release_expired_claims(
308 conn: &mut PgConnection,
309 claimed_before: DateTime<Utc>,
310) -> Result<Vec<OutboxId>, StoreError> {
311 let rows = sqlx::query(
312 "UPDATE tf_outbox
313 SET status = 'pending', next_attempt_at = NULL,
314 claim_worker_id = NULL, claim_taken_at = NULL
315 WHERE status = 'dispatching' AND claim_taken_at < $1
316 RETURNING outbox_id",
317 )
318 .bind(claimed_before)
319 .fetch_all(conn)
320 .await
321 .map_err(|error| store_error(&error))?;
322 let mut released = rows
323 .iter()
324 .map(|row| column::<Uuid>(row, "outbox_id").map(OutboxId::from))
325 .collect::<Result<Vec<OutboxId>, StoreError>>()?;
326 released.sort_unstable();
327 Ok(released)
328}
329
330async fn accept_if(
332 conn: &mut PgConnection,
333 outbox_id: &OutboxId,
334 settled: OutboxStatus,
335) -> Result<(), StoreError> {
336 match read_status(conn, outbox_id).await? {
337 None => Err(StoreError::NotFound),
338 Some(status) if status == settled => Ok(()),
339 Some(_) => Err(StoreError::Conflict),
340 }
341}
342
343async fn missing_or_conflict(
346 conn: &mut PgConnection,
347 outbox_id: &OutboxId,
348) -> Result<StoreError, StoreError> {
349 Ok(match read_status(conn, outbox_id).await? {
350 None => StoreError::NotFound,
351 Some(_) => StoreError::Conflict,
352 })
353}
354
355async fn read_status(
357 conn: &mut PgConnection,
358 outbox_id: &OutboxId,
359) -> Result<Option<OutboxStatus>, StoreError> {
360 let row = sqlx::query("SELECT status FROM tf_outbox WHERE outbox_id = $1")
361 .bind(outbox_id.as_uuid())
362 .fetch_optional(conn)
363 .await
364 .map_err(|error| store_error(&error))?;
365 match row {
366 None => Ok(None),
367 Some(row) => {
368 let status: String = column(&row, "status")?;
369 from_label(&status).map(Some)
370 }
371 }
372}
373
374fn decode_record(row: &PgRow) -> Result<OutboxRecord, StoreError> {
376 let outbox_id: Uuid = column(row, "outbox_id")?;
377 let command_id: Uuid = column(row, "command_id")?;
378 let idempotency_key: String = column(row, "idempotency_key")?;
379 let status: String = column(row, "status")?;
380 let attempts: i32 = column(row, "attempt_count")?;
381 let worker_id: Option<String> = column(row, "claim_worker_id")?;
382 let claimed_at: Option<DateTime<Utc>> = column(row, "claim_taken_at")?;
383 Ok(OutboxRecord {
384 entry: OutboxEntry {
385 outbox_id: OutboxId::from(outbox_id),
386 command_id: CommandId::from(command_id),
387 destination: column(row, "destination")?,
388 payload: column(row, "payload")?,
389 idempotency_key: IdempotencyKey::new(idempotency_key),
390 status: from_label(&status)?,
391 attempt_count: attempts_from_sql(attempts)?,
392 next_attempt_at: column(row, "next_attempt_at")?,
393 created_at: column(row, "created_at")?,
394 completed_at: column(row, "completed_at")?,
395 },
396 claim: worker_id
397 .zip(claimed_at)
398 .map(|(worker_id, claimed_at)| OutboxClaim {
399 worker_id,
400 claimed_at,
401 }),
402 last_failure: column(row, "last_failure")?,
403 remote_ref: column(row, "remote_ref")?,
404 })
405}
406
407#[async_trait]
408impl OutboxReader for PgStores {
409 async fn get(&self, outbox_id: &OutboxId) -> Result<OutboxRecord, StoreError> {
410 let mut conn = self.connection().await?;
411 get(&mut conn, outbox_id).await
412 }
413
414 async fn list_for_command(
415 &self,
416 command_id: &CommandId,
417 ) -> Result<Vec<OutboxRecord>, StoreError> {
418 let mut conn = self.connection().await?;
419 list_for_command(&mut conn, command_id).await
420 }
421}
422
423#[async_trait]
424impl OutboxWriter for PgStores {
425 async fn enqueue(&self, entry: OutboxEntry) -> Result<(), StoreError> {
426 let mut transaction = self.transaction().await?;
427 enqueue(&mut transaction, entry).await?;
428 commit(transaction).await
429 }
430
431 async fn claim_due(
432 &self,
433 at: DateTime<Utc>,
434 limit: usize,
435 worker_id: &str,
436 ) -> Result<Vec<OutboxEntry>, StoreError> {
437 let mut transaction = self.transaction().await?;
438 let claimed = claim_due(&mut transaction, at, limit, worker_id).await?;
439 commit(transaction).await?;
440 Ok(claimed)
441 }
442
443 async fn mark_completed(&self, outbox_id: &OutboxId) -> Result<(), StoreError> {
444 let mut transaction = self.transaction().await?;
445 mark_completed(&mut transaction, outbox_id, now()).await?;
446 commit(transaction).await
447 }
448
449 async fn mark_failed(
450 &self,
451 outbox_id: &OutboxId,
452 reason: String,
453 retry_at: Option<DateTime<Utc>>,
454 ) -> Result<(), StoreError> {
455 let mut transaction = self.transaction().await?;
456 mark_failed(&mut transaction, outbox_id, reason, retry_at, now()).await?;
457 commit(transaction).await
458 }
459
460 async fn mark_outcome_unknown(
461 &self,
462 outbox_id: &OutboxId,
463 remote_ref: Option<String>,
464 ) -> Result<(), StoreError> {
465 let mut transaction = self.transaction().await?;
466 mark_outcome_unknown(&mut transaction, outbox_id, remote_ref).await?;
467 commit(transaction).await
468 }
469
470 async fn reschedule(
471 &self,
472 outbox_id: &OutboxId,
473 next_attempt_at: DateTime<Utc>,
474 ) -> Result<(), StoreError> {
475 let mut transaction = self.transaction().await?;
476 reschedule(&mut transaction, outbox_id, next_attempt_at).await?;
477 commit(transaction).await
478 }
479
480 async fn release_expired_claims(
481 &self,
482 claimed_before: DateTime<Utc>,
483 ) -> Result<Vec<OutboxId>, StoreError> {
484 let mut transaction = self.transaction().await?;
485 let released = release_expired_claims(&mut transaction, claimed_before).await?;
486 commit(transaction).await?;
487 Ok(released)
488 }
489}