use async_trait::async_trait;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use turnframe_core::event::OutboxEntry;
use turnframe_core::ids::{CommandId, OutboxId};
use crate::error::StoreError;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct OutboxClaim {
pub worker_id: String,
pub claimed_at: DateTime<Utc>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct OutboxRecord {
pub entry: OutboxEntry,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub claim: Option<OutboxClaim>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_failure: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub remote_ref: Option<String>,
}
impl OutboxRecord {
#[must_use]
pub fn new(entry: OutboxEntry) -> Self {
Self {
entry,
claim: None,
last_failure: None,
remote_ref: None,
}
}
}
#[async_trait]
pub trait OutboxReader: Send + Sync {
async fn get(&self, outbox_id: &OutboxId) -> Result<OutboxRecord, StoreError>;
async fn list_for_command(
&self,
command_id: &CommandId,
) -> Result<Vec<OutboxRecord>, StoreError>;
}
#[async_trait]
pub trait OutboxWriter: Send + Sync {
async fn enqueue(&self, entry: OutboxEntry) -> Result<(), StoreError>;
async fn claim_due(
&self,
now: DateTime<Utc>,
limit: usize,
worker_id: &str,
) -> Result<Vec<OutboxEntry>, StoreError>;
async fn mark_completed(&self, outbox_id: &OutboxId) -> Result<(), StoreError>;
async fn mark_failed(
&self,
outbox_id: &OutboxId,
reason: String,
retry_at: Option<DateTime<Utc>>,
) -> Result<(), StoreError>;
async fn mark_outcome_unknown(
&self,
outbox_id: &OutboxId,
remote_ref: Option<String>,
) -> Result<(), StoreError>;
async fn reschedule(
&self,
outbox_id: &OutboxId,
next_attempt_at: DateTime<Utc>,
) -> Result<(), StoreError>;
async fn release_expired_claims(
&self,
claimed_before: DateTime<Utc>,
) -> Result<Vec<OutboxId>, StoreError>;
}
pub trait OutboxStore: OutboxReader + OutboxWriter {}
impl<T: OutboxReader + OutboxWriter + ?Sized> OutboxStore for T {}
#[cfg(test)]
mod tests {
use super::*;
use turnframe_core::command::IdempotencyKey;
use turnframe_core::event::OutboxStatus;
#[test]
fn record_round_trips() {
let record = OutboxRecord::new(OutboxEntry {
outbox_id: OutboxId::nil(),
command_id: CommandId::nil(),
destination: "airline".into(),
payload: serde_json::Value::Null,
idempotency_key: IdempotencyKey::new("k"),
status: OutboxStatus::Pending,
attempt_count: 0,
next_attempt_at: None,
created_at: DateTime::<Utc>::UNIX_EPOCH,
completed_at: None,
});
let json = serde_json::to_string(&record).unwrap();
assert_eq!(serde_json::from_str::<OutboxRecord>(&json).unwrap(), record);
}
}