use aion_store::{
ClaimScope, DEFAULT_OUTBOX_ROUTE, OutboxRow, OutboxStatus, Payload, RunId, StoreError,
WorkflowId,
};
use chrono::{DateTime, SecondsFormat, Utc};
use libsql::{Connection, Row, Transaction, TransactionBehavior, params};
mod transitions;
pub(crate) use transitions::{
complete_outbox_row, fail_outbox_row, retry_outbox_row, settle_outbox_row_cancelled,
};
const INSERT_OUTBOX_SQL: &str = "
INSERT OR IGNORE INTO outbox
(dispatch_key, workflow_id, ordinal, activity_type, input, status, attempt, visible_after, claimed_at, run_id, namespace, task_queue, node)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13)";
const REARM_OUTBOX_SQL: &str = "
INSERT INTO outbox
(dispatch_key, workflow_id, ordinal, activity_type, input, status, attempt, visible_after, claimed_at, namespace, task_queue, node)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, NULL, ?9, ?10, ?11)
ON CONFLICT(dispatch_key) DO UPDATE SET status = 'pending', visible_after = ?8, claimed_at = NULL";
const SELECT_CLAIMABLE_SQL: &str = "
SELECT dispatch_key, workflow_id, ordinal, activity_type, input, status, attempt, visible_after, claimed_at, run_id, namespace, task_queue, node
FROM outbox
WHERE status = 'pending' AND visible_after <= ?1
ORDER BY visible_after ASC, dispatch_key ASC
LIMIT ?2";
const SELECT_CLAIMABLE_UNBOUNDED_SQL: &str = "
SELECT dispatch_key, workflow_id, ordinal, activity_type, input, status, attempt, visible_after, claimed_at, run_id, namespace, task_queue, node
FROM outbox
WHERE status = 'pending' AND visible_after <= ?1
ORDER BY visible_after ASC, dispatch_key ASC";
const SELECT_CLAIMABLE_SCOPED_PREFIX_SQL: &str = "
SELECT dispatch_key, workflow_id, ordinal, activity_type, input, status, attempt, visible_after, claimed_at, run_id, namespace, task_queue, node
FROM outbox
WHERE status = 'pending' AND visible_after <= ?1 AND namespace = ?3 AND task_queue = ?4";
const SELECT_CLAIMABLE_SCOPED_SUFFIX_SQL: &str = "
ORDER BY visible_after ASC, dispatch_key ASC
LIMIT ?2";
const SELECT_CLAIMABLE_SCOPED_SUFFIX_UNBOUNDED_SQL: &str = "
ORDER BY visible_after ASC, dispatch_key ASC";
const NODE_CLAUSE_PINNED_OR_UNPINNED: &str = " AND (node IS NULL OR node = ?5)";
const NODE_CLAUSE_UNPINNED_ONLY: &str = " AND node IS NULL";
const CLAIM_ROW_SQL: &str = "
UPDATE outbox SET status = 'claimed', claimed_at = ?2 WHERE dispatch_key = ?1 AND status = 'pending'";
const COUNT_INFLIGHT_OUTBOX_SQL: &str = "
SELECT COUNT(*) FROM outbox WHERE namespace = ?1 AND status IN ('pending', 'claimed')";
const COUNT_CLAIMED_OUTBOX_SQL: &str = "
SELECT COUNT(*) FROM outbox WHERE namespace = ?1 AND status = 'claimed'";
const COUNT_CLAIMED_OUTBOX_BY_NAMESPACE_SQL: &str = "
SELECT namespace, COUNT(*) FROM outbox WHERE status = 'claimed' GROUP BY namespace";
const PENDING_OUTBOX_ROUTES_SQL: &str = "
SELECT DISTINCT namespace, task_queue, node
FROM outbox
WHERE status = 'pending' AND visible_after <= ?1";
const SELECT_STALE_CLAIMED_SQL: &str = "
SELECT dispatch_key, workflow_id, ordinal, activity_type, input, status, attempt, visible_after, claimed_at, run_id, namespace, task_queue, node
FROM outbox
WHERE status = 'claimed' AND claimed_at IS NOT NULL AND claimed_at < ?1
ORDER BY claimed_at ASC, dispatch_key ASC
LIMIT ?2";
const SELECT_LIVE_KEYS_FOR_WORKFLOW_SQL: &str = "
SELECT dispatch_key FROM outbox
WHERE workflow_id = ?1 AND status IN ('pending', 'claimed')
ORDER BY dispatch_key ASC";
const SETTLE_WORKFLOW_LIVE_ROWS_SQL: &str = "
UPDATE outbox SET status = 'cancelled', claimed_at = NULL
WHERE workflow_id = ?1 AND status IN ('pending', 'claimed')";
const SELECT_UNSETTLED_WORKFLOW_IDS_SQL: &str = "
SELECT DISTINCT workflow_id FROM outbox
WHERE status IN ('pending', 'claimed')
ORDER BY workflow_id ASC";
const REARM_STALE_CLAIMED_ROW_SQL: &str = "
UPDATE outbox SET status = 'pending', visible_after = ?2, claimed_at = NULL
WHERE dispatch_key = ?1 AND status = 'claimed'";
pub(crate) async fn insert_outbox_row(tx: &Transaction, row: &OutboxRow) -> Result<(), StoreError> {
let input = encode_payload(&row.input)?;
tx.execute(
INSERT_OUTBOX_SQL,
params![
row.dispatch_key.clone(),
row.workflow_id.to_string(),
i64::try_from(row.ordinal).map_err(|_| StoreError::Backend(format!(
"outbox ordinal overflow: {}",
row.ordinal
)))?,
row.activity_type.clone(),
input,
row.status.as_str(),
i64::from(row.attempt),
encode_instant(row.visible_after),
row.claimed_at.map(encode_instant),
row.run_id.as_ref().map(ToString::to_string),
row.namespace.clone(),
row.task_queue.clone(),
row.node.clone()
],
)
.await
.map(|_| ())
.map_err(|error| crate::error::libsql_error(&error))
}
pub(crate) async fn append_outbox_batch(
conn: &Connection,
rows: &[OutboxRow],
) -> Result<(), StoreError> {
if rows.is_empty() {
return Ok(());
}
let tx = conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.await
.map_err(|error| crate::error::libsql_error(&error))?;
for row in rows {
if let Err(error) = insert_outbox_row(&tx, row).await {
rollback(tx).await?;
return Err(error);
}
}
tx.commit()
.await
.map_err(|error| crate::error::libsql_error(&error))
}
pub(crate) async fn rearm_outbox_pending(
conn: &Connection,
rows: &[OutboxRow],
) -> Result<(), StoreError> {
if rows.is_empty() {
return Ok(());
}
let tx = conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.await
.map_err(|error| crate::error::libsql_error(&error))?;
for row in rows {
if let Err(error) = rearm_outbox_row(&tx, row).await {
rollback(tx).await?;
return Err(error);
}
}
tx.commit()
.await
.map_err(|error| crate::error::libsql_error(&error))
}
async fn rearm_outbox_row(tx: &Transaction, row: &OutboxRow) -> Result<(), StoreError> {
let input = encode_payload(&row.input)?;
tx.execute(
REARM_OUTBOX_SQL,
params![
row.dispatch_key.clone(),
row.workflow_id.to_string(),
i64::try_from(row.ordinal).map_err(|_| StoreError::Backend(format!(
"outbox ordinal overflow: {}",
row.ordinal
)))?,
row.activity_type.clone(),
input,
OutboxStatus::Pending.as_str(),
i64::from(row.attempt),
encode_instant(row.visible_after),
row.namespace.clone(),
row.task_queue.clone(),
row.node.clone()
],
)
.await
.map(|_| ())
.map_err(|error| crate::error::libsql_error(&error))
}
pub(crate) async fn claim_outbox_rows(
conn: &Connection,
limit: u32,
) -> Result<Vec<OutboxRow>, StoreError> {
if limit == 0 {
return Ok(Vec::new());
}
let tx = conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.await
.map_err(|error| crate::error::libsql_error(&error))?;
let claimed = match select_and_claim(&tx, limit).await {
Ok(claimed) => claimed,
Err(error) => {
rollback(tx).await?;
return Err(error);
}
};
tx.commit()
.await
.map_err(|error| crate::error::libsql_error(&error))?;
Ok(claimed)
}
pub(crate) async fn claim_outbox_rows_excluding(
conn: &Connection,
limit: u32,
held: &std::collections::HashSet<aion_core::WorkflowId>,
) -> Result<Vec<OutboxRow>, StoreError> {
if limit == 0 {
return Ok(Vec::new());
}
if held.is_empty() {
return claim_outbox_rows(conn, limit).await;
}
let tx = conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.await
.map_err(|error| crate::error::libsql_error(&error))?;
let claimed = match select_and_claim_excluding(&tx, limit, held).await {
Ok(claimed) => claimed,
Err(error) => {
rollback(tx).await?;
return Err(error);
}
};
tx.commit()
.await
.map_err(|error| crate::error::libsql_error(&error))?;
Ok(claimed)
}
async fn select_and_claim_excluding(
tx: &Transaction,
limit: u32,
held: &std::collections::HashSet<aion_core::WorkflowId>,
) -> Result<Vec<OutboxRow>, StoreError> {
let claimed_at = Utc::now();
let now = encode_instant(claimed_at);
let mut rows = tx
.query(SELECT_CLAIMABLE_UNBOUNDED_SQL, params![now.clone()])
.await
.map_err(|error| crate::error::libsql_error(&error))?;
let mut claimed = Vec::new();
while let Some(row) = rows
.next()
.await
.map_err(|error| crate::error::libsql_error(&error))?
{
if claimed.len() >= limit as usize {
break;
}
let decoded = decode_row(&row)?;
if held.contains(&decoded.workflow_id) {
continue;
}
tx.execute(
CLAIM_ROW_SQL,
params![decoded.dispatch_key.clone(), now.clone()],
)
.await
.map_err(|error| crate::error::libsql_error(&error))?;
claimed.push(OutboxRow {
status: OutboxStatus::Claimed,
claimed_at: Some(claimed_at),
..decoded
});
}
Ok(claimed)
}
async fn select_and_claim(tx: &Transaction, limit: u32) -> Result<Vec<OutboxRow>, StoreError> {
let claimed_at = Utc::now();
let now = encode_instant(claimed_at);
let mut rows = tx
.query(SELECT_CLAIMABLE_SQL, params![now.clone(), i64::from(limit)])
.await
.map_err(|error| crate::error::libsql_error(&error))?;
let mut claimed = Vec::new();
while let Some(row) = rows
.next()
.await
.map_err(|error| crate::error::libsql_error(&error))?
{
let decoded = decode_row(&row)?;
tx.execute(
CLAIM_ROW_SQL,
params![decoded.dispatch_key.clone(), now.clone()],
)
.await
.map_err(|error| crate::error::libsql_error(&error))?;
claimed.push(OutboxRow {
status: OutboxStatus::Claimed,
claimed_at: Some(claimed_at),
..decoded
});
}
Ok(claimed)
}
pub(crate) async fn claim_outbox_rows_scoped(
conn: &Connection,
scope: &ClaimScope,
limit: u32,
) -> Result<Vec<OutboxRow>, StoreError> {
claim_outbox_rows_scoped_excluding(conn, scope, limit, &std::collections::HashSet::new()).await
}
pub(crate) async fn claim_outbox_rows_scoped_excluding(
conn: &Connection,
scope: &ClaimScope,
limit: u32,
held: &std::collections::HashSet<aion_core::WorkflowId>,
) -> Result<Vec<OutboxRow>, StoreError> {
if limit == 0 {
return Ok(Vec::new());
}
let tx = conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.await
.map_err(|error| crate::error::libsql_error(&error))?;
let claimed = match select_and_claim_scoped(&tx, scope, limit, held).await {
Ok(claimed) => claimed,
Err(error) => {
rollback(tx).await?;
return Err(error);
}
};
tx.commit()
.await
.map_err(|error| crate::error::libsql_error(&error))?;
Ok(claimed)
}
async fn select_and_claim_scoped(
tx: &Transaction,
scope: &ClaimScope,
limit: u32,
held: &std::collections::HashSet<aion_core::WorkflowId>,
) -> Result<Vec<OutboxRow>, StoreError> {
let claimed_at = Utc::now();
let now = encode_instant(claimed_at);
let node_clause = if scope.node.is_some() {
NODE_CLAUSE_PINNED_OR_UNPINNED
} else {
NODE_CLAUSE_UNPINNED_ONLY
};
let suffix = if held.is_empty() {
SELECT_CLAIMABLE_SCOPED_SUFFIX_SQL
} else {
SELECT_CLAIMABLE_SCOPED_SUFFIX_UNBOUNDED_SQL
};
let sql = format!("{SELECT_CLAIMABLE_SCOPED_PREFIX_SQL}{node_clause}{suffix}");
let mut rows = if let Some(node) = scope.node.as_deref() {
tx.query(
&sql,
params![
now.clone(),
i64::from(limit),
scope.namespace.clone(),
scope.task_queue.clone(),
node.to_string()
],
)
.await
} else {
tx.query(
&sql,
params![
now.clone(),
i64::from(limit),
scope.namespace.clone(),
scope.task_queue.clone()
],
)
.await
}
.map_err(|error| crate::error::libsql_error(&error))?;
let mut claimed = Vec::new();
while let Some(row) = rows
.next()
.await
.map_err(|error| crate::error::libsql_error(&error))?
{
if claimed.len() >= limit as usize {
break;
}
let decoded = decode_row(&row)?;
if held.contains(&decoded.workflow_id) {
continue;
}
tx.execute(
CLAIM_ROW_SQL,
params![decoded.dispatch_key.clone(), now.clone()],
)
.await
.map_err(|error| crate::error::libsql_error(&error))?;
claimed.push(OutboxRow {
status: OutboxStatus::Claimed,
claimed_at: Some(claimed_at),
..decoded
});
}
Ok(claimed)
}
pub(crate) async fn rearm_stale_claimed_outbox_rows(
conn: &Connection,
older_than: DateTime<Utc>,
visible_after: DateTime<Utc>,
limit: u32,
) -> Result<Vec<OutboxRow>, StoreError> {
if limit == 0 {
return Ok(Vec::new());
}
let tx = conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.await
.map_err(|error| crate::error::libsql_error(&error))?;
let rows = match select_and_rearm_stale_claimed(&tx, older_than, visible_after, limit).await {
Ok(rows) => rows,
Err(error) => {
rollback(tx).await?;
return Err(error);
}
};
tx.commit()
.await
.map_err(|error| crate::error::libsql_error(&error))?;
Ok(rows)
}
async fn select_and_rearm_stale_claimed(
tx: &Transaction,
older_than: DateTime<Utc>,
visible_after: DateTime<Utc>,
limit: u32,
) -> Result<Vec<OutboxRow>, StoreError> {
let mut rows = tx
.query(
SELECT_STALE_CLAIMED_SQL,
params![encode_instant(older_than), i64::from(limit)],
)
.await
.map_err(|error| crate::error::libsql_error(&error))?;
let mut rearmed = Vec::new();
while let Some(row) = rows
.next()
.await
.map_err(|error| crate::error::libsql_error(&error))?
{
let decoded = decode_row(&row)?;
tx.execute(
REARM_STALE_CLAIMED_ROW_SQL,
params![decoded.dispatch_key.clone(), encode_instant(visible_after)],
)
.await
.map_err(|error| crate::error::libsql_error(&error))?;
rearmed.push(OutboxRow {
status: OutboxStatus::Pending,
visible_after,
claimed_at: None,
..decoded
});
}
Ok(rearmed)
}
pub(crate) async fn list_stale_claimed_outbox_rows(
conn: &Connection,
older_than: DateTime<Utc>,
limit: u32,
) -> Result<Vec<OutboxRow>, StoreError> {
if limit == 0 {
return Ok(Vec::new());
}
let mut rows = conn
.query(
SELECT_STALE_CLAIMED_SQL,
params![encode_instant(older_than), i64::from(limit)],
)
.await
.map_err(|error| crate::error::libsql_error(&error))?;
let mut stale = Vec::new();
while let Some(row) = rows
.next()
.await
.map_err(|error| crate::error::libsql_error(&error))?
{
stale.push(decode_row(&row)?);
}
Ok(stale)
}
pub(crate) async fn list_unsettled_outbox_workflow_ids(
conn: &Connection,
) -> Result<Vec<WorkflowId>, StoreError> {
let mut rows = conn
.query(SELECT_UNSETTLED_WORKFLOW_IDS_SQL, ())
.await
.map_err(|error| crate::error::libsql_error(&error))?;
let mut workflow_ids = Vec::new();
while let Some(row) = rows
.next()
.await
.map_err(|error| crate::error::libsql_error(&error))?
{
let workflow_id: String = row
.get(0)
.map_err(|error| crate::error::libsql_error(&error))?;
workflow_ids.push(decode_workflow_id(&workflow_id)?);
}
Ok(workflow_ids)
}
pub(crate) async fn settle_workflow_outbox_rows_cancelled(
conn: &Connection,
workflow_id: &WorkflowId,
) -> Result<Vec<String>, StoreError> {
let tx = conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.await
.map_err(|error| crate::error::libsql_error(&error))?;
let settled = match select_and_settle_workflow_rows(&tx, workflow_id).await {
Ok(settled) => settled,
Err(error) => {
rollback(tx).await?;
return Err(error);
}
};
tx.commit()
.await
.map_err(|error| crate::error::libsql_error(&error))?;
Ok(settled)
}
async fn select_and_settle_workflow_rows(
tx: &Transaction,
workflow_id: &WorkflowId,
) -> Result<Vec<String>, StoreError> {
let mut rows = tx
.query(
SELECT_LIVE_KEYS_FOR_WORKFLOW_SQL,
params![workflow_id.to_string()],
)
.await
.map_err(|error| crate::error::libsql_error(&error))?;
let mut settled = Vec::new();
while let Some(row) = rows
.next()
.await
.map_err(|error| crate::error::libsql_error(&error))?
{
let dispatch_key: String = row
.get(0)
.map_err(|error| crate::error::libsql_error(&error))?;
settled.push(dispatch_key);
}
if !settled.is_empty() {
tx.execute(
SETTLE_WORKFLOW_LIVE_ROWS_SQL,
params![workflow_id.to_string()],
)
.await
.map_err(|error| crate::error::libsql_error(&error))?;
}
Ok(settled)
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct OutboxRowState {
pub status: OutboxStatus,
pub attempt: u32,
pub visible_after: DateTime<Utc>,
}
const SELECT_ROW_STATE_SQL: &str = "
SELECT status, attempt, visible_after FROM outbox WHERE dispatch_key = ?1";
pub(crate) async fn outbox_row_state(
conn: &Connection,
dispatch_key: &str,
) -> Result<Option<OutboxRowState>, StoreError> {
let mut rows = conn
.query(SELECT_ROW_STATE_SQL, params![dispatch_key.to_string()])
.await
.map_err(|error| crate::error::libsql_error(&error))?;
let Some(row) = rows
.next()
.await
.map_err(|error| crate::error::libsql_error(&error))?
else {
return Ok(None);
};
let status: String = row
.get(0)
.map_err(|error| crate::error::libsql_error(&error))?;
let attempt: i64 = row
.get(1)
.map_err(|error| crate::error::libsql_error(&error))?;
let visible_after: String = row
.get(2)
.map_err(|error| crate::error::libsql_error(&error))?;
Ok(Some(OutboxRowState {
status: OutboxStatus::parse_token(&status)?,
attempt: u32::try_from(attempt)
.map_err(|_| StoreError::Backend(format!("outbox attempt out of range: {attempt}")))?,
visible_after: decode_instant(&visible_after)?,
}))
}
pub(crate) async fn count_inflight_outbox_rows(
conn: &Connection,
namespace: &str,
) -> Result<u64, StoreError> {
let mut rows = conn
.query(COUNT_INFLIGHT_OUTBOX_SQL, params![namespace.to_string()])
.await
.map_err(|error| crate::error::libsql_error(&error))?;
let Some(row) = rows
.next()
.await
.map_err(|error| crate::error::libsql_error(&error))?
else {
return Ok(0);
};
let count: i64 = row
.get(0)
.map_err(|error| crate::error::libsql_error(&error))?;
u64::try_from(count)
.map_err(|_| StoreError::Serialization(format!("outbox in-flight count negative: {count}")))
}
pub(crate) async fn count_claimed_outbox_rows(
conn: &Connection,
namespace: &str,
) -> Result<u64, StoreError> {
let mut rows = conn
.query(COUNT_CLAIMED_OUTBOX_SQL, params![namespace.to_string()])
.await
.map_err(|error| crate::error::libsql_error(&error))?;
let Some(row) = rows
.next()
.await
.map_err(|error| crate::error::libsql_error(&error))?
else {
return Ok(0);
};
let count: i64 = row
.get(0)
.map_err(|error| crate::error::libsql_error(&error))?;
u64::try_from(count)
.map_err(|_| StoreError::Serialization(format!("outbox claimed count negative: {count}")))
}
pub(crate) async fn count_claimed_outbox_rows_by_namespace(
conn: &Connection,
namespaces: &[&str],
) -> Result<std::collections::BTreeMap<String, u64>, StoreError> {
let requested: std::collections::BTreeSet<&str> = namespaces.iter().copied().collect();
let mut counts: std::collections::BTreeMap<String, u64> =
requested.iter().map(|ns| ((*ns).to_owned(), 0)).collect();
let mut rows = conn
.query(COUNT_CLAIMED_OUTBOX_BY_NAMESPACE_SQL, ())
.await
.map_err(|error| crate::error::libsql_error(&error))?;
while let Some(row) = rows
.next()
.await
.map_err(|error| crate::error::libsql_error(&error))?
{
let namespace: String = row
.get(0)
.map_err(|error| crate::error::libsql_error(&error))?;
if !requested.contains(namespace.as_str()) {
continue;
}
let count: i64 = row
.get(1)
.map_err(|error| crate::error::libsql_error(&error))?;
let count = u64::try_from(count).map_err(|_| {
StoreError::Serialization(format!("outbox claimed count negative: {count}"))
})?;
counts.insert(namespace, count);
}
Ok(counts)
}
pub(crate) async fn pending_outbox_routes(
conn: &Connection,
) -> Result<Vec<ClaimScope>, StoreError> {
let now = encode_instant(Utc::now());
let mut rows = conn
.query(PENDING_OUTBOX_ROUTES_SQL, params![now])
.await
.map_err(|error| crate::error::libsql_error(&error))?;
let mut routes = Vec::new();
while let Some(row) = rows
.next()
.await
.map_err(|error| crate::error::libsql_error(&error))?
{
let namespace: Option<String> = row
.get(0)
.map_err(|error| crate::error::libsql_error(&error))?;
let task_queue: Option<String> = row
.get(1)
.map_err(|error| crate::error::libsql_error(&error))?;
let node: Option<String> = row
.get(2)
.map_err(|error| crate::error::libsql_error(&error))?;
routes.push(
ClaimScope::new(
namespace.unwrap_or_else(|| String::from(DEFAULT_OUTBOX_ROUTE)),
task_queue.unwrap_or_else(|| String::from(DEFAULT_OUTBOX_ROUTE)),
)
.with_node(node),
);
}
Ok(routes)
}
fn decode_row(row: &Row) -> Result<OutboxRow, StoreError> {
let dispatch_key: String = row
.get(0)
.map_err(|error| crate::error::libsql_error(&error))?;
let workflow_id: String = row
.get(1)
.map_err(|error| crate::error::libsql_error(&error))?;
let ordinal: i64 = row
.get(2)
.map_err(|error| crate::error::libsql_error(&error))?;
let activity_type: String = row
.get(3)
.map_err(|error| crate::error::libsql_error(&error))?;
let input: Vec<u8> = row
.get(4)
.map_err(|error| crate::error::libsql_error(&error))?;
let status: String = row
.get(5)
.map_err(|error| crate::error::libsql_error(&error))?;
let attempt: i64 = row
.get(6)
.map_err(|error| crate::error::libsql_error(&error))?;
let visible_after: String = row
.get(7)
.map_err(|error| crate::error::libsql_error(&error))?;
let claimed_at: Option<String> = row
.get(8)
.map_err(|error| crate::error::libsql_error(&error))?;
let run_id: Option<String> = row
.get(9)
.map_err(|error| crate::error::libsql_error(&error))?;
let namespace: Option<String> = row
.get(10)
.map_err(|error| crate::error::libsql_error(&error))?;
let task_queue: Option<String> = row
.get(11)
.map_err(|error| crate::error::libsql_error(&error))?;
let node: Option<String> = row
.get(12)
.map_err(|error| crate::error::libsql_error(&error))?;
Ok(OutboxRow {
dispatch_key,
workflow_id: decode_workflow_id(&workflow_id)?,
ordinal: u64::try_from(ordinal)
.map_err(|_| StoreError::Backend(format!("outbox ordinal was negative: {ordinal}")))?,
activity_type,
input: decode_payload(&input)?,
status: OutboxStatus::parse_token(&status)?,
attempt: u32::try_from(attempt)
.map_err(|_| StoreError::Backend(format!("outbox attempt out of range: {attempt}")))?,
visible_after: decode_instant(&visible_after)?,
claimed_at: claimed_at.as_deref().map(decode_instant).transpose()?,
run_id: run_id.as_deref().map(decode_run_id).transpose()?,
namespace: namespace.unwrap_or_else(|| String::from(DEFAULT_OUTBOX_ROUTE)),
task_queue: task_queue.unwrap_or_else(|| String::from(DEFAULT_OUTBOX_ROUTE)),
node,
})
}
fn encode_payload(payload: &Payload) -> Result<Vec<u8>, StoreError> {
serde_json::to_vec(payload).map_err(|error| crate::error::serde_json_error(&error))
}
fn decode_payload(bytes: &[u8]) -> Result<Payload, StoreError> {
serde_json::from_slice(bytes).map_err(|error| crate::error::serde_json_error(&error))
}
fn decode_workflow_id(value: &str) -> Result<WorkflowId, StoreError> {
uuid::Uuid::parse_str(value)
.map(WorkflowId::new)
.map_err(|error| StoreError::Serialization(format!("invalid outbox workflow id: {error}")))
}
fn decode_run_id(value: &str) -> Result<RunId, StoreError> {
uuid::Uuid::parse_str(value)
.map(RunId::new)
.map_err(|error| StoreError::Serialization(format!("invalid outbox run id: {error}")))
}
fn encode_instant(instant: DateTime<Utc>) -> String {
instant.to_rfc3339_opts(SecondsFormat::Nanos, true)
}
fn decode_instant(value: &str) -> Result<DateTime<Utc>, StoreError> {
DateTime::parse_from_rfc3339(value)
.map(|date_time| date_time.with_timezone(&Utc))
.map_err(|error| StoreError::Serialization(error.to_string()))
}
async fn rollback(tx: Transaction) -> Result<(), StoreError> {
tx.rollback()
.await
.map_err(|error| crate::error::libsql_error(&error))
}