commit 2cdf57bd4376803470f5c43e81c2bf470fcab72c
parent 2473fac05942475239202ea005c5751259a0690f
Author: triesap <tyson@radroots.org>
Date: Tue, 7 Jul 2026 21:49:40 +0000
outbox: add plan scoped publish claims
- persist the active delivery plan on publish claims
- restrict delivery target outcomes to the claimed plan
- add typed preview and deferred operation lifecycle states
- cover duplicate endpoint plans and proxy exclusivity
Diffstat:
4 files changed, 641 insertions(+), 138 deletions(-)
diff --git a/crates/outbox/migrations/0001_outbox.up.sql b/crates/outbox/migrations/0001_outbox.up.sql
@@ -4,7 +4,7 @@ CREATE TABLE IF NOT EXISTS outbox_operations (
expected_pubkey TEXT NOT NULL,
idempotency_key TEXT,
operation_idempotency_digest TEXT NOT NULL,
- status TEXT NOT NULL CHECK (status IN ('queued', 'complete', 'failed_terminal', 'cancelled')),
+ status TEXT NOT NULL CHECK (status IN ('queued', 'complete', 'deferred_until_implemented', 'preview_unavailable', 'failed_terminal', 'cancelled')),
created_at_ms INTEGER NOT NULL,
updated_at_ms INTEGER NOT NULL
);
@@ -24,11 +24,12 @@ CREATE TABLE IF NOT EXISTS outbox_event (
draft_json TEXT NOT NULL,
signed_event_json TEXT,
raw_event_json TEXT,
- state TEXT NOT NULL CHECK (state IN ('draft_queued', 'signing', 'signed', 'publishing', 'published', 'sign_retryable', 'publish_retryable', 'failed_terminal', 'cancelled')),
+ state TEXT NOT NULL CHECK (state IN ('draft_queued', 'signing', 'signed', 'publishing', 'published', 'sign_retryable', 'publish_retryable', 'deferred_until_implemented', 'preview_unavailable', 'failed_terminal', 'cancelled')),
attempt_count INTEGER NOT NULL,
claim_token TEXT,
claim_owner TEXT,
claim_expires_at_ms INTEGER,
+ active_delivery_plan_id INTEGER REFERENCES outbox_delivery_plan(delivery_plan_id) ON DELETE SET NULL,
next_attempt_after_ms INTEGER NOT NULL,
last_error TEXT,
event_store_ingested INTEGER NOT NULL,
diff --git a/crates/outbox/src/error.rs b/crates/outbox/src/error.rs
@@ -50,6 +50,9 @@ pub enum RadrootsOutboxError {
#[error("Claim token mismatch for outbox event {outbox_event_id}")]
ClaimTokenMismatch { outbox_event_id: i64 },
+ #[error("Outbox event {outbox_event_id} has no active delivery plan")]
+ MissingActiveDeliveryPlan { outbox_event_id: i64 },
+
#[error("Signed event missing for outbox event {0}")]
MissingSignedEvent(i64),
diff --git a/crates/outbox/src/model.rs b/crates/outbox/src/model.rs
@@ -11,6 +11,8 @@ use radroots_transport::{
pub enum RadrootsOutboxOperationStatus {
Queued,
Complete,
+ DeferredUntilImplemented,
+ PreviewUnavailable,
FailedTerminal,
Cancelled,
}
@@ -20,6 +22,8 @@ impl RadrootsOutboxOperationStatus {
match self {
Self::Queued => "queued",
Self::Complete => "complete",
+ Self::DeferredUntilImplemented => "deferred_until_implemented",
+ Self::PreviewUnavailable => "preview_unavailable",
Self::FailedTerminal => "failed_terminal",
Self::Cancelled => "cancelled",
}
@@ -29,6 +33,8 @@ impl RadrootsOutboxOperationStatus {
match value {
"queued" => Ok(Self::Queued),
"complete" => Ok(Self::Complete),
+ "deferred_until_implemented" => Ok(Self::DeferredUntilImplemented),
+ "preview_unavailable" => Ok(Self::PreviewUnavailable),
"failed_terminal" => Ok(Self::FailedTerminal),
"cancelled" => Ok(Self::Cancelled),
_ => Err(RadrootsOutboxError::InvalidStoredEnum {
@@ -48,6 +54,8 @@ pub enum RadrootsOutboxEventState {
Published,
SignRetryable,
PublishRetryable,
+ DeferredUntilImplemented,
+ PreviewUnavailable,
FailedTerminal,
Cancelled,
}
@@ -62,6 +70,8 @@ impl RadrootsOutboxEventState {
Self::Published => "published",
Self::SignRetryable => "sign_retryable",
Self::PublishRetryable => "publish_retryable",
+ Self::DeferredUntilImplemented => "deferred_until_implemented",
+ Self::PreviewUnavailable => "preview_unavailable",
Self::FailedTerminal => "failed_terminal",
Self::Cancelled => "cancelled",
}
@@ -76,6 +86,8 @@ impl RadrootsOutboxEventState {
"published" => Ok(Self::Published),
"sign_retryable" => Ok(Self::SignRetryable),
"publish_retryable" => Ok(Self::PublishRetryable),
+ "deferred_until_implemented" => Ok(Self::DeferredUntilImplemented),
+ "preview_unavailable" => Ok(Self::PreviewUnavailable),
"failed_terminal" => Ok(Self::FailedTerminal),
"cancelled" => Ok(Self::Cancelled),
_ => Err(RadrootsOutboxError::InvalidStoredEnum {
@@ -404,6 +416,7 @@ pub struct RadrootsOutboxEventRecord {
pub claim_token: Option<String>,
pub claim_owner: Option<String>,
pub claim_expires_at_ms: Option<i64>,
+ pub active_delivery_plan_id: Option<i64>,
pub next_attempt_after_ms: i64,
pub last_error: Option<String>,
pub event_store_ingested: bool,
@@ -461,6 +474,7 @@ pub struct RadrootsOutboxClaimedEvent {
pub attempt_count: i64,
pub state: RadrootsOutboxEventState,
pub claim_token: String,
+ pub active_delivery_plan_id: Option<i64>,
pub draft: RadrootsFrozenEventDraft,
pub signed_event: Option<RadrootsSignedNostrEvent>,
pub delivery_targets: Vec<RadrootsOutboxDeliveryTargetRecord>,
@@ -499,6 +513,14 @@ mod tests {
(RadrootsOutboxOperationStatus::Queued, "queued"),
(RadrootsOutboxOperationStatus::Complete, "complete"),
(
+ RadrootsOutboxOperationStatus::DeferredUntilImplemented,
+ "deferred_until_implemented",
+ ),
+ (
+ RadrootsOutboxOperationStatus::PreviewUnavailable,
+ "preview_unavailable",
+ ),
+ (
RadrootsOutboxOperationStatus::FailedTerminal,
"failed_terminal",
),
@@ -511,6 +533,35 @@ mod tests {
);
}
+ for (state, expected) in [
+ (RadrootsOutboxEventState::DraftQueued, "draft_queued"),
+ (RadrootsOutboxEventState::Signing, "signing"),
+ (RadrootsOutboxEventState::Signed, "signed"),
+ (RadrootsOutboxEventState::Publishing, "publishing"),
+ (RadrootsOutboxEventState::Published, "published"),
+ (RadrootsOutboxEventState::SignRetryable, "sign_retryable"),
+ (
+ RadrootsOutboxEventState::PublishRetryable,
+ "publish_retryable",
+ ),
+ (
+ RadrootsOutboxEventState::DeferredUntilImplemented,
+ "deferred_until_implemented",
+ ),
+ (
+ RadrootsOutboxEventState::PreviewUnavailable,
+ "preview_unavailable",
+ ),
+ (RadrootsOutboxEventState::FailedTerminal, "failed_terminal"),
+ (RadrootsOutboxEventState::Cancelled, "cancelled"),
+ ] {
+ assert_eq!(state.as_str(), expected);
+ assert_eq!(
+ RadrootsOutboxEventState::parse(expected).expect("event state"),
+ state
+ );
+ }
+
for (status, expected) in [
(RadrootsOutboxDeliveryPlanStatus::Queued, "queued"),
(RadrootsOutboxDeliveryPlanStatus::Complete, "complete"),
diff --git a/crates/outbox/src/store.rs b/crates/outbox/src/store.rs
@@ -19,7 +19,7 @@ use radroots_events::draft::{
RadrootsFrozenEventDraft, RadrootsSignedNostrEvent, validate_signed_nostr_event_matches_draft,
};
use radroots_transport::{
- RADROOTS_RETICULUM_PREVIEW_ENDPOINT_URI, RadrootsTransportKind,
+ RADROOTS_RETICULUM_PREVIEW_ENDPOINT_URI, RadrootsTransportError, RadrootsTransportKind,
RadrootsTransportSatisfactionClass, RadrootsTransportSatisfactionPolicy,
RadrootsTransportTarget, RadrootsTransportTargetFingerprint, RadrootsTransportTargetUri,
};
@@ -85,7 +85,7 @@ impl RadrootsOutbox {
now_ms: i64,
) -> Result<RadrootsOutboxStatusSummary, RadrootsOutboxError> {
let row = sqlx::query(
- "SELECT COUNT(*) AS total_events, COALESCE(SUM(CASE WHEN state IN ('draft_queued', 'signing', 'signed', 'publishing') THEN 1 ELSE 0 END), 0) AS pending_events, COALESCE(SUM(CASE WHEN state IN ('sign_retryable', 'publish_retryable') THEN 1 ELSE 0 END), 0) AS retryable_events, COALESCE(SUM(CASE WHEN state IN ('published', 'failed_terminal', 'cancelled') THEN 1 ELSE 0 END), 0) AS terminal_events, COALESCE(SUM(CASE WHEN state = 'failed_terminal' THEN 1 ELSE 0 END), 0) AS failed_terminal_events, COALESCE(SUM(CASE WHEN state = 'publishing' THEN 1 ELSE 0 END), 0) AS publishing_events FROM outbox_event",
+ "SELECT COUNT(*) AS total_events, COALESCE(SUM(CASE WHEN state IN ('draft_queued', 'signing', 'signed', 'publishing') THEN 1 ELSE 0 END), 0) AS pending_events, COALESCE(SUM(CASE WHEN state IN ('sign_retryable', 'publish_retryable') THEN 1 ELSE 0 END), 0) AS retryable_events, COALESCE(SUM(CASE WHEN state IN ('published', 'failed_terminal', 'cancelled') THEN 1 ELSE 0 END), 0) AS terminal_events, COALESCE(SUM(CASE WHEN state = 'failed_terminal' THEN 1 ELSE 0 END), 0) AS failed_terminal_events, COALESCE(SUM(CASE WHEN state = 'preview_unavailable' THEN 1 ELSE 0 END), 0) AS preview_unavailable_events, COALESCE(SUM(CASE WHEN state = 'deferred_until_implemented' THEN 1 ELSE 0 END), 0) AS deferred_until_implemented_events, COALESCE(SUM(CASE WHEN state = 'publishing' THEN 1 ELSE 0 END), 0) AS publishing_events FROM outbox_event",
)
.fetch_one(&self.pool)
.await?;
@@ -97,59 +97,6 @@ impl RadrootsOutbox {
.fetch_one(&self.pool)
.await?
.try_get(0)?;
- let deferred_preview_row = sqlx::query(
- r#"
- WITH signed_target_state AS (
- SELECT
- event.outbox_event_id,
- EXISTS (
- SELECT 1
- FROM outbox_delivery_plan AS plan
- JOIN outbox_delivery_target AS target
- ON target.delivery_plan_id = plan.delivery_plan_id
- WHERE plan.outbox_event_id = event.outbox_event_id
- AND target.status = 'preview_unavailable'
- ) AS has_preview_unavailable,
- EXISTS (
- SELECT 1
- FROM outbox_delivery_plan AS plan
- JOIN outbox_delivery_target AS target
- ON target.delivery_plan_id = plan.delivery_plan_id
- WHERE plan.outbox_event_id = event.outbox_event_id
- AND target.status = 'deferred_until_implemented'
- ) AS has_deferred_until_implemented,
- EXISTS (
- SELECT 1
- FROM outbox_delivery_plan AS plan
- JOIN outbox_delivery_target AS target
- ON target.delivery_plan_id = plan.delivery_plan_id
- WHERE plan.outbox_event_id = event.outbox_event_id
- AND target.status IN ('pending', 'failed_retryable')
- ) AS has_ready_target
- FROM outbox_event AS event
- WHERE event.state = 'signed'
- AND event.signed_event_json IS NOT NULL
- )
- SELECT
- COALESCE(SUM(CASE
- WHEN has_preview_unavailable AND NOT has_ready_target THEN 1
- ELSE 0
- END), 0) AS preview_unavailable_events,
- COALESCE(SUM(CASE
- WHEN NOT has_preview_unavailable
- AND has_deferred_until_implemented
- AND NOT has_ready_target THEN 1
- ELSE 0
- END), 0) AS deferred_until_implemented_events
- FROM signed_target_state
- "#,
- )
- .fetch_one(&self.pool)
- .await?;
- let preview_unavailable_events: i64 =
- deferred_preview_row.try_get("preview_unavailable_events")?;
- let deferred_until_implemented_events: i64 =
- deferred_preview_row.try_get("deferred_until_implemented_events")?;
let last_attempt_at_ms =
sqlx::query("SELECT MAX(attempted_at_ms) FROM outbox_delivery_attempt")
.fetch_one(&self.pool)
@@ -162,17 +109,14 @@ impl RadrootsOutbox {
.await?
.map(|row| row.try_get("last_error"))
.transpose()?;
- let pending_events: i64 = row.try_get("pending_events")?;
Ok(RadrootsOutboxStatusSummary {
total_events: row.try_get("total_events")?,
- pending_events: pending_events
- .saturating_sub(preview_unavailable_events)
- .saturating_sub(deferred_until_implemented_events),
+ pending_events: row.try_get("pending_events")?,
retryable_events: row.try_get("retryable_events")?,
terminal_events: row.try_get("terminal_events")?,
failed_terminal_events: row.try_get("failed_terminal_events")?,
- preview_unavailable_events,
- deferred_until_implemented_events,
+ preview_unavailable_events: row.try_get("preview_unavailable_events")?,
+ deferred_until_implemented_events: row.try_get("deferred_until_implemented_events")?,
ready_signed_events,
publishing_events: row.try_get("publishing_events")?,
last_attempt_at_ms,
@@ -255,12 +199,8 @@ impl RadrootsOutbox {
)
.await?;
if plan.status == RadrootsOutboxEnqueueStatus::Inserted {
- reactivate_event_for_new_plan(
- &mut tx,
- existing.outbox_event_id,
- input.created_at_ms,
- )
- .await?;
+ sync_signed_event_lifecycle(&mut tx, existing.outbox_event_id, input.created_at_ms)
+ .await?;
}
tx.commit().await?;
return Ok(RadrootsOutboxEnqueueReceipt {
@@ -363,14 +303,8 @@ impl RadrootsOutbox {
input.created_at_ms,
)
.await?;
- if plan.status == RadrootsOutboxEnqueueStatus::Inserted {
- reactivate_event_for_new_plan(
- &mut tx,
- existing.outbox_event_id,
- input.created_at_ms,
- )
+ sync_signed_event_lifecycle(&mut tx, existing.outbox_event_id, input.created_at_ms)
.await?;
- }
tx.commit().await?;
return Ok(RadrootsOutboxEnqueueReceipt {
status: plan.status,
@@ -419,6 +353,7 @@ impl RadrootsOutbox {
let plan =
insert_or_get_delivery_plan(&mut tx, outbox_event_id, &prepared, input.created_at_ms)
.await?;
+ sync_signed_event_lifecycle(&mut tx, outbox_event_id, input.created_at_ms).await?;
tx.commit().await?;
Ok(RadrootsOutboxEnqueueReceipt {
status: RadrootsOutboxEnqueueStatus::Inserted,
@@ -449,7 +384,7 @@ impl RadrootsOutbox {
outbox_event_id: i64,
) -> Result<Option<RadrootsOutboxEventRecord>, RadrootsOutboxError> {
let row = sqlx::query(
- "SELECT outbox_event_id, operation_id, event_id, expected_pubkey, draft_json, signed_event_json, raw_event_json, state, attempt_count, claim_token, claim_owner, claim_expires_at_ms, next_attempt_after_ms, last_error, event_store_ingested, event_store_inserted, event_store_ingested_at_ms, created_at_ms, updated_at_ms FROM outbox_event WHERE outbox_event_id = ?",
+ "SELECT outbox_event_id, operation_id, event_id, expected_pubkey, draft_json, signed_event_json, raw_event_json, state, attempt_count, claim_token, claim_owner, claim_expires_at_ms, active_delivery_plan_id, next_attempt_after_ms, last_error, event_store_ingested, event_store_inserted, event_store_ingested_at_ms, created_at_ms, updated_at_ms FROM outbox_event WHERE outbox_event_id = ?",
)
.bind(outbox_event_id)
.fetch_optional(&self.pool)
@@ -507,17 +442,40 @@ impl RadrootsOutbox {
) => RadrootsOutboxEventState::Signing,
_ => RadrootsOutboxEventState::Publishing,
};
- let changed = claim_event(
- &mut tx,
- outbox_event_id,
- claimed_state,
- claim_owner.as_ref(),
- claim_token.as_ref(),
- claim_expires_at_ms,
- now_ms,
- "AND (claim_token IS NULL OR claim_expires_at_ms <= ?)",
- )
- .await?;
+ let active_delivery_plan_id = if claimed_state == RadrootsOutboxEventState::Publishing {
+ Some(
+ ready_delivery_plan_id_for_event_tx(&mut tx, outbox_event_id)
+ .await?
+ .ok_or(RadrootsOutboxError::MissingActiveDeliveryPlan { outbox_event_id })?,
+ )
+ } else {
+ None
+ };
+ let changed = if let Some(delivery_plan_id) = active_delivery_plan_id {
+ claim_publish_event_for_plan(
+ &mut tx,
+ outbox_event_id,
+ delivery_plan_id,
+ claim_owner.as_ref(),
+ claim_token.as_ref(),
+ claim_expires_at_ms,
+ now_ms,
+ )
+ .await?
+ } else {
+ claim_event(
+ &mut tx,
+ outbox_event_id,
+ claimed_state,
+ None,
+ claim_owner.as_ref(),
+ claim_token.as_ref(),
+ claim_expires_at_ms,
+ now_ms,
+ "AND (claim_token IS NULL OR claim_expires_at_ms <= ?)",
+ )
+ .await?
+ };
if changed.rows_affected() == 0 {
tx.commit().await?;
return Ok(None);
@@ -526,6 +484,7 @@ impl RadrootsOutbox {
&mut tx,
outbox_event_id,
claimed_state,
+ active_delivery_plan_id,
claim_token.as_ref(),
)
.await?;
@@ -542,7 +501,7 @@ impl RadrootsOutbox {
) -> Result<Option<RadrootsOutboxClaimedEvent>, RadrootsOutboxError> {
let mut tx = self.pool.begin().await?;
let row = sqlx::query(
- "SELECT outbox_event_id FROM outbox_event AS event WHERE event.state IN ('signed', 'publish_retryable') AND event.signed_event_json IS NOT NULL AND event.next_attempt_after_ms <= ? AND (event.claim_token IS NULL OR event.claim_expires_at_ms <= ?) AND EXISTS (SELECT 1 FROM outbox_delivery_plan AS plan JOIN outbox_delivery_target AS target ON target.delivery_plan_id = plan.delivery_plan_id WHERE plan.outbox_event_id = event.outbox_event_id AND target.status IN ('pending', 'failed_retryable')) ORDER BY event.created_at_ms, event.outbox_event_id LIMIT 1",
+ "SELECT event.outbox_event_id, plan.delivery_plan_id FROM outbox_event AS event JOIN outbox_delivery_plan AS plan ON plan.outbox_event_id = event.outbox_event_id JOIN outbox_delivery_target AS target ON target.delivery_plan_id = plan.delivery_plan_id WHERE event.state IN ('signed', 'publish_retryable') AND event.signed_event_json IS NOT NULL AND event.next_attempt_after_ms <= ? AND (event.claim_token IS NULL OR event.claim_expires_at_ms <= ?) AND target.status IN ('pending', 'failed_retryable') GROUP BY event.outbox_event_id, plan.delivery_plan_id ORDER BY event.created_at_ms, event.outbox_event_id, plan.delivery_plan_id LIMIT 1",
)
.bind(now_ms)
.bind(now_ms)
@@ -553,15 +512,15 @@ impl RadrootsOutbox {
return Ok(None);
};
let outbox_event_id: i64 = row.try_get("outbox_event_id")?;
- let changed = claim_event(
+ let active_delivery_plan_id: i64 = row.try_get("delivery_plan_id")?;
+ let changed = claim_publish_event_for_plan(
&mut tx,
outbox_event_id,
- RadrootsOutboxEventState::Publishing,
+ active_delivery_plan_id,
claim_owner.as_ref(),
claim_token.as_ref(),
claim_expires_at_ms,
now_ms,
- "AND state IN ('signed', 'publish_retryable') AND signed_event_json IS NOT NULL AND (claim_token IS NULL OR claim_expires_at_ms <= ?)",
)
.await?;
if changed.rows_affected() == 0 {
@@ -572,6 +531,7 @@ impl RadrootsOutbox {
&mut tx,
outbox_event_id,
RadrootsOutboxEventState::Publishing,
+ Some(active_delivery_plan_id),
claim_token.as_ref(),
)
.await?;
@@ -588,18 +548,21 @@ impl RadrootsOutbox {
now_ms: i64,
) -> Result<Option<RadrootsOutboxClaimedEvent>, RadrootsOutboxError> {
let mut tx = self.pool.begin().await?;
- let changed = sqlx::query(
- "UPDATE outbox_event SET state = ?, claim_token = ?, claim_owner = ?, claim_expires_at_ms = ?, attempt_count = attempt_count + 1, updated_at_ms = ? WHERE outbox_event_id = ? AND state IN ('signed', 'publish_retryable') AND signed_event_json IS NOT NULL AND next_attempt_after_ms <= ? AND (claim_token IS NULL OR claim_expires_at_ms <= ?) AND EXISTS (SELECT 1 FROM outbox_delivery_plan AS plan JOIN outbox_delivery_target AS target ON target.delivery_plan_id = plan.delivery_plan_id WHERE plan.outbox_event_id = outbox_event.outbox_event_id AND target.status IN ('pending', 'failed_retryable'))",
+ let Some(active_delivery_plan_id) =
+ ready_delivery_plan_id_for_event_tx(&mut tx, outbox_event_id).await?
+ else {
+ tx.commit().await?;
+ return Ok(None);
+ };
+ let changed = claim_publish_event_for_plan(
+ &mut tx,
+ outbox_event_id,
+ active_delivery_plan_id,
+ claim_owner.as_ref(),
+ claim_token.as_ref(),
+ claim_expires_at_ms,
+ now_ms,
)
- .bind(RadrootsOutboxEventState::Publishing.as_str())
- .bind(claim_token.as_ref())
- .bind(claim_owner.as_ref())
- .bind(claim_expires_at_ms)
- .bind(now_ms)
- .bind(outbox_event_id)
- .bind(now_ms)
- .bind(now_ms)
- .execute(&mut *tx)
.await?;
if changed.rows_affected() == 0 {
tx.commit().await?;
@@ -609,6 +572,7 @@ impl RadrootsOutbox {
&mut tx,
outbox_event_id,
RadrootsOutboxEventState::Publishing,
+ Some(active_delivery_plan_id),
claim_token.as_ref(),
)
.await?;
@@ -623,7 +587,12 @@ impl RadrootsOutbox {
signed_event: RadrootsSignedNostrEvent,
now_ms: i64,
) -> Result<RadrootsSignedNostrEvent, RadrootsOutboxError> {
- let record = self.claimed_event(outbox_event_id, claim_token).await?;
+ let mut tx = self.pool.begin().await?;
+ let record = event_by_id_tx(&mut tx, outbox_event_id).await?;
+ let stored = record.claim_token.as_deref();
+ if stored != Some(claim_token) {
+ return Err(RadrootsOutboxError::ClaimTokenMismatch { outbox_event_id });
+ }
if signed_event.id != record.event_id {
return Err(RadrootsOutboxError::SignedEventIdMismatch {
expected_event_id: record.event_id,
@@ -640,10 +609,13 @@ impl RadrootsOutbox {
.bind(now_ms)
.bind(outbox_event_id)
.bind(claim_token)
- .execute(&self.pool)
+ .execute(&mut *tx)
.await?;
- self.ensure_claimed_update(outbox_event_id, claim_token, changed)
- .await?;
+ if changed.rows_affected() == 0 {
+ return Err(RadrootsOutboxError::ClaimTokenMismatch { outbox_event_id });
+ }
+ sync_signed_event_lifecycle(&mut tx, outbox_event_id, now_ms).await?;
+ tx.commit().await?;
Ok(signed_event)
}
@@ -656,7 +628,7 @@ impl RadrootsOutbox {
now_ms: i64,
) -> Result<(), RadrootsOutboxError> {
let changed = sqlx::query(
- "UPDATE outbox_event SET state = ?, claim_token = NULL, claim_owner = NULL, claim_expires_at_ms = NULL, last_error = ?, next_attempt_after_ms = ?, updated_at_ms = ? WHERE outbox_event_id = ? AND claim_token = ?",
+ "UPDATE outbox_event SET state = ?, claim_token = NULL, claim_owner = NULL, claim_expires_at_ms = NULL, active_delivery_plan_id = NULL, last_error = ?, next_attempt_after_ms = ?, updated_at_ms = ? WHERE outbox_event_id = ? AND claim_token = ?",
)
.bind(RadrootsOutboxEventState::SignRetryable.as_str())
.bind(error.as_ref())
@@ -679,7 +651,7 @@ impl RadrootsOutbox {
now_ms: i64,
) -> Result<(), RadrootsOutboxError> {
let changed = sqlx::query(
- "UPDATE outbox_event SET state = ?, claim_token = NULL, claim_owner = NULL, claim_expires_at_ms = NULL, last_error = ?, next_attempt_after_ms = ?, updated_at_ms = ? WHERE outbox_event_id = ? AND claim_token = ?",
+ "UPDATE outbox_event SET state = ?, claim_token = NULL, claim_owner = NULL, claim_expires_at_ms = NULL, active_delivery_plan_id = NULL, last_error = ?, next_attempt_after_ms = ?, updated_at_ms = ? WHERE outbox_event_id = ? AND claim_token = ?",
)
.bind(RadrootsOutboxEventState::PublishRetryable.as_str())
.bind(error.as_ref())
@@ -695,7 +667,7 @@ impl RadrootsOutbox {
pub async fn recover_expired_claims(&self, now_ms: i64) -> Result<u64, RadrootsOutboxError> {
let changed = sqlx::query(
- "UPDATE outbox_event SET state = CASE WHEN state = 'signing' AND signed_event_json IS NULL THEN 'sign_retryable' WHEN state = 'signing' AND signed_event_json IS NOT NULL THEN 'signed' WHEN state = 'publishing' THEN 'publish_retryable' ELSE state END, claim_token = NULL, claim_owner = NULL, claim_expires_at_ms = NULL, updated_at_ms = ? WHERE claim_token IS NOT NULL AND claim_expires_at_ms <= ? AND state IN ('signing', 'signed', 'publishing')",
+ "UPDATE outbox_event SET state = CASE WHEN state = 'signing' AND signed_event_json IS NULL THEN 'sign_retryable' WHEN state = 'signing' AND signed_event_json IS NOT NULL THEN 'signed' WHEN state = 'publishing' THEN 'publish_retryable' ELSE state END, claim_token = NULL, claim_owner = NULL, claim_expires_at_ms = NULL, active_delivery_plan_id = NULL, updated_at_ms = ? WHERE claim_token IS NOT NULL AND claim_expires_at_ms <= ? AND state IN ('signing', 'signed', 'publishing', 'preview_unavailable', 'deferred_until_implemented')",
)
.bind(now_ms)
.bind(now_ms)
@@ -872,6 +844,9 @@ impl RadrootsOutbox {
) -> Result<RadrootsOutboxEventState, RadrootsOutboxError> {
let mut tx = self.pool.begin().await?;
let row = claimed_event_identity_tx(&mut tx, outbox_event_id, claim_token).await?;
+ if row.active_delivery_plan_id.is_none() {
+ return Err(RadrootsOutboxError::MissingActiveDeliveryPlan { outbox_event_id });
+ }
let evaluation = evaluate_delivery_plans(&mut tx, outbox_event_id, now_ms).await?;
let (event_state, operation_status, last_error, next_attempt_after_ms) =
if evaluation.all_complete {
@@ -895,17 +870,31 @@ impl RadrootsOutbox {
Some(retryable_error.as_ref()),
next_attempt_after_ms,
)
- } else {
+ } else if evaluation.any_preview_unavailable {
+ (
+ RadrootsOutboxEventState::PreviewUnavailable,
+ Some(RadrootsOutboxOperationStatus::PreviewUnavailable),
+ None,
+ now_ms,
+ )
+ } else if evaluation.any_deferred_until_implemented {
(
- RadrootsOutboxEventState::Signed,
+ RadrootsOutboxEventState::DeferredUntilImplemented,
+ Some(RadrootsOutboxOperationStatus::DeferredUntilImplemented),
None,
- Some("delivery deferred until implemented"),
+ now_ms,
+ )
+ } else {
+ (
+ RadrootsOutboxEventState::FailedTerminal,
+ Some(RadrootsOutboxOperationStatus::FailedTerminal),
+ Some(terminal_error.as_ref()),
now_ms,
)
};
let changed = sqlx::query(
- "UPDATE outbox_event SET state = ?, claim_token = NULL, claim_owner = NULL, claim_expires_at_ms = NULL, last_error = ?, next_attempt_after_ms = ?, updated_at_ms = ? WHERE outbox_event_id = ? AND claim_token = ?",
+ "UPDATE outbox_event SET state = ?, claim_token = NULL, claim_owner = NULL, claim_expires_at_ms = NULL, active_delivery_plan_id = NULL, last_error = ?, next_attempt_after_ms = ?, updated_at_ms = ? WHERE outbox_event_id = ? AND claim_token = ?",
)
.bind(event_state.as_str())
.bind(last_error)
@@ -1045,7 +1034,7 @@ impl RadrootsOutbox {
.await?;
}
let changed = sqlx::query(
- "UPDATE outbox_event SET state = ?, claim_token = NULL, claim_owner = NULL, claim_expires_at_ms = NULL, last_error = ?, next_attempt_after_ms = ?, updated_at_ms = ? WHERE outbox_event_id = ? AND claim_token = ?",
+ "UPDATE outbox_event SET state = ?, claim_token = NULL, claim_owner = NULL, claim_expires_at_ms = NULL, active_delivery_plan_id = NULL, last_error = ?, next_attempt_after_ms = ?, updated_at_ms = ? WHERE outbox_event_id = ? AND claim_token = ?",
)
.bind(event_state.as_str())
.bind(last_error)
@@ -1082,16 +1071,20 @@ impl RadrootsOutbox {
attempted_at_ms: i64,
) -> Result<(), RadrootsOutboxError> {
let mut tx = self.pool.begin().await?;
- claimed_event_identity_tx(&mut tx, outbox_event_id, claim_token).await?;
+ let identity = claimed_event_identity_tx(&mut tx, outbox_event_id, claim_token).await?;
+ let Some(active_delivery_plan_id) = identity.active_delivery_plan_id else {
+ return Err(RadrootsOutboxError::MissingActiveDeliveryPlan { outbox_event_id });
+ };
let completed_at_ms = status.is_completed().then_some(attempted_at_ms);
let changed = sqlx::query(
- "UPDATE outbox_delivery_target SET status = ?, attempt_count = attempt_count + 1, last_attempt_at_ms = ?, completed_at_ms = ?, last_error = ? WHERE delivery_target_id = ? AND delivery_plan_id IN (SELECT delivery_plan_id FROM outbox_delivery_plan WHERE outbox_event_id = ?)",
+ "UPDATE outbox_delivery_target SET status = ?, attempt_count = attempt_count + 1, last_attempt_at_ms = ?, completed_at_ms = ?, last_error = ? WHERE delivery_target_id = ? AND delivery_plan_id = ? AND delivery_plan_id IN (SELECT delivery_plan_id FROM outbox_delivery_plan WHERE outbox_event_id = ?)",
)
.bind(status.as_str())
.bind(attempted_at_ms)
.bind(completed_at_ms)
.bind(message)
.bind(delivery_target_id)
+ .bind(active_delivery_plan_id)
.bind(outbox_event_id)
.execute(&mut *tx)
.await?;
@@ -1131,6 +1124,7 @@ struct ExistingOperation {
struct ClaimedEventIdentity {
operation_id: i64,
+ active_delivery_plan_id: Option<i64>,
}
struct PreparedDeliveryPlan {
@@ -1158,6 +1152,8 @@ struct PlanEvaluation {
all_complete: bool,
any_failed_terminal: bool,
any_ready: bool,
+ any_preview_unavailable: bool,
+ any_deferred_until_implemented: bool,
}
async fn configure_connection(
@@ -1213,6 +1209,15 @@ fn prepare_delivery_plan(
if targets.is_empty() {
return Err(RadrootsOutboxError::EmptyDeliveryTargets);
}
+ if targets
+ .iter()
+ .any(|target| target.kind == RadrootsTransportKind::Proxy)
+ && targets.len() != 1
+ {
+ return Err(RadrootsOutboxError::Transport(
+ RadrootsTransportError::InvalidTransportKind,
+ ));
+ }
let required_success_count = input
.satisfaction_policy
.required_target_count(targets.len())? as i64;
@@ -1401,32 +1406,109 @@ async fn insert_or_get_delivery_plan(
})
}
-async fn reactivate_event_for_new_plan(
+async fn sync_signed_event_lifecycle(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
outbox_event_id: i64,
now_ms: i64,
) -> Result<(), RadrootsOutboxError> {
+ let row = sqlx::query(
+ "SELECT operation_id, signed_event_json FROM outbox_event WHERE outbox_event_id = ?",
+ )
+ .bind(outbox_event_id)
+ .fetch_one(&mut **tx)
+ .await?;
+ let signed_event_json: Option<String> = row.try_get("signed_event_json")?;
+ if signed_event_json.is_none() {
+ return Ok(());
+ }
+ let (event_state, operation_status) =
+ signed_event_lifecycle_for_plans(tx, outbox_event_id).await?;
sqlx::query(
- "UPDATE outbox_operations SET status = ?, updated_at_ms = ? WHERE operation_id = (SELECT operation_id FROM outbox_event WHERE outbox_event_id = ?)",
+ "UPDATE outbox_event SET state = ?, last_error = NULL, next_attempt_after_ms = ?, updated_at_ms = ? WHERE outbox_event_id = ?",
)
- .bind(RadrootsOutboxOperationStatus::Queued.as_str())
+ .bind(event_state.as_str())
+ .bind(now_ms)
.bind(now_ms)
.bind(outbox_event_id)
.execute(&mut **tx)
.await?;
sqlx::query(
- "UPDATE outbox_event SET state = CASE WHEN signed_event_json IS NULL THEN ? ELSE ? END, last_error = NULL, next_attempt_after_ms = ?, updated_at_ms = ? WHERE outbox_event_id = ? AND state IN ('published', 'failed_terminal')",
+ "UPDATE outbox_operations SET status = ?, updated_at_ms = ? WHERE operation_id = ?",
)
- .bind(RadrootsOutboxEventState::DraftQueued.as_str())
- .bind(RadrootsOutboxEventState::Signed.as_str())
- .bind(now_ms)
+ .bind(operation_status.as_str())
.bind(now_ms)
- .bind(outbox_event_id)
+ .bind(row.try_get::<i64, _>("operation_id")?)
.execute(&mut **tx)
.await?;
Ok(())
}
+async fn signed_event_lifecycle_for_plans(
+ tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
+ outbox_event_id: i64,
+) -> Result<(RadrootsOutboxEventState, RadrootsOutboxOperationStatus), RadrootsOutboxError> {
+ let plans = delivery_plans_for_tx(tx, outbox_event_id).await?;
+ if plans.is_empty()
+ || plans
+ .iter()
+ .any(|plan| plan.status == RadrootsOutboxDeliveryPlanStatus::Queued)
+ {
+ return Ok((
+ RadrootsOutboxEventState::Signed,
+ RadrootsOutboxOperationStatus::Queued,
+ ));
+ }
+ if plans
+ .iter()
+ .all(|plan| plan.status == RadrootsOutboxDeliveryPlanStatus::Complete)
+ {
+ return Ok((
+ RadrootsOutboxEventState::Published,
+ RadrootsOutboxOperationStatus::Complete,
+ ));
+ }
+ if plans
+ .iter()
+ .any(|plan| plan.status == RadrootsOutboxDeliveryPlanStatus::FailedTerminal)
+ {
+ return Ok((
+ RadrootsOutboxEventState::FailedTerminal,
+ RadrootsOutboxOperationStatus::FailedTerminal,
+ ));
+ }
+ if plans
+ .iter()
+ .any(|plan| plan.status == RadrootsOutboxDeliveryPlanStatus::PreviewUnavailable)
+ {
+ return Ok((
+ RadrootsOutboxEventState::PreviewUnavailable,
+ RadrootsOutboxOperationStatus::PreviewUnavailable,
+ ));
+ }
+ if plans
+ .iter()
+ .any(|plan| plan.status == RadrootsOutboxDeliveryPlanStatus::DeferredUntilImplemented)
+ {
+ return Ok((
+ RadrootsOutboxEventState::DeferredUntilImplemented,
+ RadrootsOutboxOperationStatus::DeferredUntilImplemented,
+ ));
+ }
+ if plans
+ .iter()
+ .all(|plan| plan.status == RadrootsOutboxDeliveryPlanStatus::Cancelled)
+ {
+ return Ok((
+ RadrootsOutboxEventState::Cancelled,
+ RadrootsOutboxOperationStatus::Cancelled,
+ ));
+ }
+ Ok((
+ RadrootsOutboxEventState::Signed,
+ RadrootsOutboxOperationStatus::Queued,
+ ))
+}
+
async fn ensure_event_signed(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
outbox_event_id: i64,
@@ -1453,6 +1535,7 @@ async fn claim_event(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
outbox_event_id: i64,
claimed_state: RadrootsOutboxEventState,
+ active_delivery_plan_id: Option<i64>,
claim_owner: &str,
claim_token: &str,
claim_expires_at_ms: i64,
@@ -1460,13 +1543,14 @@ async fn claim_event(
suffix: &str,
) -> Result<SqliteQueryResult, RadrootsOutboxError> {
let sql = format!(
- "UPDATE outbox_event SET state = ?, claim_token = ?, claim_owner = ?, claim_expires_at_ms = ?, attempt_count = attempt_count + 1, updated_at_ms = ? WHERE outbox_event_id = ? {suffix}"
+ "UPDATE outbox_event SET state = ?, claim_token = ?, claim_owner = ?, claim_expires_at_ms = ?, active_delivery_plan_id = ?, attempt_count = attempt_count + 1, updated_at_ms = ? WHERE outbox_event_id = ? {suffix}"
);
let changed = sqlx::query(sql.as_str())
.bind(claimed_state.as_str())
.bind(claim_token)
.bind(claim_owner)
.bind(claim_expires_at_ms)
+ .bind(active_delivery_plan_id)
.bind(now_ms)
.bind(outbox_event_id)
.bind(now_ms)
@@ -1475,17 +1559,49 @@ async fn claim_event(
Ok(changed)
}
+async fn claim_publish_event_for_plan(
+ tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
+ outbox_event_id: i64,
+ delivery_plan_id: i64,
+ claim_owner: &str,
+ claim_token: &str,
+ claim_expires_at_ms: i64,
+ now_ms: i64,
+) -> Result<SqliteQueryResult, RadrootsOutboxError> {
+ let changed = sqlx::query(
+ "UPDATE outbox_event SET state = ?, claim_token = ?, claim_owner = ?, claim_expires_at_ms = ?, active_delivery_plan_id = ?, attempt_count = attempt_count + 1, updated_at_ms = ? WHERE outbox_event_id = ? AND state IN ('signed', 'publish_retryable') AND signed_event_json IS NOT NULL AND next_attempt_after_ms <= ? AND (claim_token IS NULL OR claim_expires_at_ms <= ?) AND EXISTS (SELECT 1 FROM outbox_delivery_plan AS plan JOIN outbox_delivery_target AS target ON target.delivery_plan_id = plan.delivery_plan_id WHERE plan.outbox_event_id = outbox_event.outbox_event_id AND plan.delivery_plan_id = ? AND target.status IN ('pending', 'failed_retryable'))",
+ )
+ .bind(RadrootsOutboxEventState::Publishing.as_str())
+ .bind(claim_token)
+ .bind(claim_owner)
+ .bind(claim_expires_at_ms)
+ .bind(delivery_plan_id)
+ .bind(now_ms)
+ .bind(outbox_event_id)
+ .bind(now_ms)
+ .bind(now_ms)
+ .bind(delivery_plan_id)
+ .execute(&mut **tx)
+ .await?;
+ Ok(changed)
+}
+
async fn claimed_event_from_tx(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
outbox_event_id: i64,
claimed_state: RadrootsOutboxEventState,
+ active_delivery_plan_id: Option<i64>,
claim_token: &str,
) -> Result<RadrootsOutboxClaimedEvent, RadrootsOutboxError> {
let record = event_by_id_tx(tx, outbox_event_id).await?;
- let delivery_targets = if claimed_state == RadrootsOutboxEventState::Publishing {
- ready_delivery_targets_for_event_tx(tx, outbox_event_id).await?
- } else {
- delivery_targets_for_event_tx(tx, outbox_event_id).await?
+ let delivery_targets = match (claimed_state, active_delivery_plan_id) {
+ (RadrootsOutboxEventState::Publishing, Some(delivery_plan_id)) => {
+ ready_delivery_targets_for_plan_tx(tx, delivery_plan_id).await?
+ }
+ (RadrootsOutboxEventState::Publishing, None) => {
+ return Err(RadrootsOutboxError::MissingActiveDeliveryPlan { outbox_event_id });
+ }
+ _ => delivery_targets_for_event_tx(tx, outbox_event_id).await?,
};
Ok(RadrootsOutboxClaimedEvent {
outbox_event_id: record.outbox_event_id,
@@ -1494,6 +1610,7 @@ async fn claimed_event_from_tx(
attempt_count: record.attempt_count,
state: claimed_state,
claim_token: claim_token.to_owned(),
+ active_delivery_plan_id,
draft: record.draft,
signed_event: record.signed_event,
delivery_targets,
@@ -1505,7 +1622,7 @@ async fn event_by_id_tx(
outbox_event_id: i64,
) -> Result<RadrootsOutboxEventRecord, RadrootsOutboxError> {
let row = sqlx::query(
- "SELECT outbox_event_id, operation_id, event_id, expected_pubkey, draft_json, signed_event_json, raw_event_json, state, attempt_count, claim_token, claim_owner, claim_expires_at_ms, next_attempt_after_ms, last_error, event_store_ingested, event_store_inserted, event_store_ingested_at_ms, created_at_ms, updated_at_ms FROM outbox_event WHERE outbox_event_id = ?",
+ "SELECT outbox_event_id, operation_id, event_id, expected_pubkey, draft_json, signed_event_json, raw_event_json, state, attempt_count, claim_token, claim_owner, claim_expires_at_ms, active_delivery_plan_id, next_attempt_after_ms, last_error, event_store_ingested, event_store_inserted, event_store_ingested_at_ms, created_at_ms, updated_at_ms FROM outbox_event WHERE outbox_event_id = ?",
)
.bind(outbox_event_id)
.fetch_one(&mut **tx)
@@ -1519,7 +1636,7 @@ async fn claimed_event_identity_tx(
claim_token: &str,
) -> Result<ClaimedEventIdentity, RadrootsOutboxError> {
let row =
- sqlx::query("SELECT operation_id, claim_token FROM outbox_event WHERE outbox_event_id = ?")
+ sqlx::query("SELECT operation_id, claim_token, active_delivery_plan_id FROM outbox_event WHERE outbox_event_id = ?")
.bind(outbox_event_id)
.fetch_optional(&mut **tx)
.await?;
@@ -1532,6 +1649,7 @@ async fn claimed_event_identity_tx(
}
Ok(ClaimedEventIdentity {
operation_id: row.try_get("operation_id")?,
+ active_delivery_plan_id: row.try_get("active_delivery_plan_id")?,
})
}
@@ -1587,14 +1705,27 @@ async fn delivery_targets_for_event_tx(
rows.into_iter().map(delivery_target_from_row).collect()
}
-async fn ready_delivery_targets_for_event_tx(
+async fn ready_delivery_plan_id_for_event_tx(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
outbox_event_id: i64,
+) -> Result<Option<i64>, RadrootsOutboxError> {
+ let row = sqlx::query(
+ "SELECT plan.delivery_plan_id FROM outbox_delivery_plan AS plan JOIN outbox_delivery_target AS target ON target.delivery_plan_id = plan.delivery_plan_id WHERE plan.outbox_event_id = ? AND target.status IN ('pending', 'failed_retryable') GROUP BY plan.delivery_plan_id ORDER BY plan.delivery_plan_id LIMIT 1",
+ )
+ .bind(outbox_event_id)
+ .fetch_optional(&mut **tx)
+ .await?;
+ Ok(row.map(|row| row.try_get("delivery_plan_id")).transpose()?)
+}
+
+async fn ready_delivery_targets_for_plan_tx(
+ tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
+ delivery_plan_id: i64,
) -> Result<Vec<RadrootsOutboxDeliveryTargetRecord>, RadrootsOutboxError> {
let rows = sqlx::query(
- "SELECT target.delivery_target_id, target.delivery_plan_id, target.transport_kind, target.endpoint_uri, target.endpoint_fingerprint, target.status, target.attempt_count, target.last_attempt_at_ms, target.completed_at_ms, target.last_error FROM outbox_delivery_target AS target JOIN outbox_delivery_plan AS plan ON plan.delivery_plan_id = target.delivery_plan_id WHERE plan.outbox_event_id = ? AND target.status IN ('pending', 'failed_retryable') ORDER BY target.delivery_plan_id, target.delivery_target_id",
+ "SELECT delivery_target_id, delivery_plan_id, transport_kind, endpoint_uri, endpoint_fingerprint, status, attempt_count, last_attempt_at_ms, completed_at_ms, last_error FROM outbox_delivery_target WHERE delivery_plan_id = ? AND status IN ('pending', 'failed_retryable') ORDER BY delivery_target_id",
)
- .bind(outbox_event_id)
+ .bind(delivery_plan_id)
.fetch_all(&mut **tx)
.await?;
rows.into_iter().map(delivery_target_from_row).collect()
@@ -1635,6 +1766,8 @@ async fn evaluate_delivery_plans(
let mut all_complete = !plans.is_empty();
let mut any_failed_terminal = false;
let mut any_ready = false;
+ let mut any_preview_unavailable = false;
+ let mut any_deferred_until_implemented = false;
for plan in plans {
let targets = delivery_targets_for_plan_tx(tx, plan.delivery_plan_id).await?;
let satisfied_count = targets
@@ -1685,6 +1818,12 @@ async fn evaluate_delivery_plans(
if ready_count > 0 {
any_ready = true;
}
+ if plan_status == RadrootsOutboxDeliveryPlanStatus::PreviewUnavailable {
+ any_preview_unavailable = true;
+ }
+ if plan_status == RadrootsOutboxDeliveryPlanStatus::DeferredUntilImplemented {
+ any_deferred_until_implemented = true;
+ }
sqlx::query(
"UPDATE outbox_delivery_plan SET status = ?, satisfied_at_ms = CASE WHEN ? = 'complete' THEN ? ELSE satisfied_at_ms END, updated_at_ms = ? WHERE delivery_plan_id = ?",
)
@@ -1700,6 +1839,8 @@ async fn evaluate_delivery_plans(
all_complete,
any_failed_terminal,
any_ready,
+ any_preview_unavailable,
+ any_deferred_until_implemented,
})
}
@@ -1743,6 +1884,7 @@ fn event_from_row(
claim_token: row.try_get("claim_token")?,
claim_owner: row.try_get("claim_owner")?,
claim_expires_at_ms: row.try_get("claim_expires_at_ms")?,
+ active_delivery_plan_id: row.try_get("active_delivery_plan_id")?,
next_attempt_after_ms: row.try_get("next_attempt_after_ms")?,
last_error: row.try_get("last_error")?,
event_store_ingested: row.try_get::<i64, _>("event_store_ingested")? != 0,
@@ -2039,6 +2181,10 @@ mod tests {
.expect("reticulum target")
}
+ fn proxy_target(uri: &str) -> RadrootsTransportTarget {
+ RadrootsTransportTarget::new(RadrootsTransportKind::Proxy, uri).expect("proxy target")
+ }
+
fn malformed_reticulum_target(uri: &str) -> RadrootsTransportTarget {
let endpoint_uri = RadrootsTransportTargetUri::parse(uri).expect("target uri");
let endpoint_fingerprint = RadrootsTransportTargetFingerprint::from_target(
@@ -2264,6 +2410,34 @@ mod tests {
}
#[tokio::test]
+ async fn enqueue_rejects_mixed_proxy_delegate_targets_before_persistence() {
+ let outbox = RadrootsOutbox::open_memory().await.expect("open");
+ let draft = post_draft(hex_64('a').as_str(), "mixed proxy");
+
+ let err = outbox
+ .enqueue_operation(RadrootsOutboxOperationInput::new(
+ "publish_post",
+ draft,
+ delivery_plan(vec![
+ proxy_target("radrootsd-proxy:publish"),
+ nostr_target(NOSTR_PRIMARY_WSS),
+ ]),
+ 1_000,
+ ))
+ .await
+ .expect_err("mixed proxy targets");
+
+ assert!(matches!(
+ err,
+ RadrootsOutboxError::Transport(RadrootsTransportError::InvalidTransportKind)
+ ));
+ assert_eq!(table_count(&outbox, "outbox_operations").await, 0);
+ assert_eq!(table_count(&outbox, "outbox_event").await, 0);
+ assert_eq!(table_count(&outbox, "outbox_delivery_plan").await, 0);
+ assert_eq!(table_count(&outbox, "outbox_delivery_target").await, 0);
+ }
+
+ #[tokio::test]
async fn signed_enqueue_claims_ready_delivery_targets_and_records_attempts() {
let outbox = RadrootsOutbox::open_memory().await.expect("open");
let draft = post_draft(FIXTURE_ALICE_PUBLIC_KEY_HEX, "signed");
@@ -2288,7 +2462,17 @@ mod tests {
.expect("claim")
.expect("claimed");
assert_eq!(claimed.state, RadrootsOutboxEventState::Publishing);
+ assert_eq!(
+ claimed.active_delivery_plan_id,
+ Some(receipt.delivery_plan_id)
+ );
assert_eq!(claimed.delivery_targets.len(), 2);
+ assert!(
+ claimed
+ .delivery_targets
+ .iter()
+ .all(|target| target.delivery_plan_id == receipt.delivery_plan_id)
+ );
outbox
.mark_delivery_target_accepted(
@@ -2339,6 +2523,149 @@ mod tests {
}
#[tokio::test]
+ async fn publish_claims_one_delivery_plan_and_rejects_sibling_target_outcomes() {
+ let outbox = RadrootsOutbox::open_memory().await.expect("open");
+ let draft = post_draft(FIXTURE_ALICE_PUBLIC_KEY_HEX, "duplicate endpoint plans");
+ let signed_event =
+ radroots_nostr_sign_frozen_draft(&fixture_keys(), &draft).expect("signed event");
+ let first = outbox
+ .enqueue_signed_operation(
+ RadrootsOutboxSignedOperationInput::new(
+ "publish_post",
+ draft.clone(),
+ signed_event.clone(),
+ RadrootsOutboxDeliveryPlanInput::new(
+ "transport.nostr.primary-plan",
+ 1,
+ RadrootsTransportSatisfactionPolicy::all_accepted(),
+ vec![nostr_target(NOSTR_PRIMARY_WSS)],
+ ),
+ true,
+ 1_007,
+ 1_000,
+ )
+ .with_idempotency_key("idem-duplicate-endpoint"),
+ )
+ .await
+ .expect("first plan");
+ let second = outbox
+ .enqueue_signed_operation(
+ RadrootsOutboxSignedOperationInput::new(
+ "publish_post",
+ draft,
+ signed_event,
+ RadrootsOutboxDeliveryPlanInput::new(
+ "transport.nostr.secondary-plan",
+ 1,
+ RadrootsTransportSatisfactionPolicy::all_accepted(),
+ vec![nostr_target(NOSTR_PRIMARY_WSS)],
+ ),
+ true,
+ 1_017,
+ 1_010,
+ )
+ .with_idempotency_key("idem-duplicate-endpoint"),
+ )
+ .await
+ .expect("second plan");
+
+ assert_eq!(first.outbox_event_id, second.outbox_event_id);
+ assert_ne!(first.delivery_plan_id, second.delivery_plan_id);
+ let initial_targets = outbox
+ .delivery_targets(first.outbox_event_id)
+ .await
+ .expect("targets");
+ assert_eq!(initial_targets.len(), 2);
+ assert_eq!(
+ initial_targets[0].endpoint_fingerprint,
+ initial_targets[1].endpoint_fingerprint
+ );
+
+ let first_claim = outbox
+ .claim_next_ready_signed_event("publisher", "claim-a", 2_000, 1_020)
+ .await
+ .expect("claim")
+ .expect("claimed");
+ let active_plan_id = first_claim.active_delivery_plan_id.expect("active plan");
+ assert_eq!(first_claim.delivery_targets.len(), 1);
+ assert_eq!(
+ first_claim.delivery_targets[0].delivery_plan_id,
+ active_plan_id
+ );
+ let sibling_target = initial_targets
+ .iter()
+ .find(|target| target.delivery_plan_id != active_plan_id)
+ .expect("sibling target");
+ let sibling_err = outbox
+ .mark_delivery_target_accepted(
+ first.outbox_event_id,
+ "claim-a",
+ sibling_target.delivery_target_id,
+ 1_100,
+ )
+ .await
+ .expect_err("sibling target is not in active plan");
+ assert!(matches!(
+ sibling_err,
+ RadrootsOutboxError::DeliveryTargetNotFound(_)
+ ));
+
+ outbox
+ .mark_delivery_target_accepted(
+ first.outbox_event_id,
+ "claim-a",
+ first_claim.delivery_targets[0].delivery_target_id,
+ 1_110,
+ )
+ .await
+ .expect("active target accepted");
+ let retry_state = outbox
+ .complete_publish_attempt(
+ first.outbox_event_id,
+ "claim-a",
+ "retryable",
+ "terminal",
+ 2_500,
+ 1_200,
+ )
+ .await
+ .expect("complete first plan");
+ assert_eq!(retry_state, RadrootsOutboxEventState::PublishRetryable);
+
+ let retry_claim = outbox
+ .claim_next_ready_signed_event("publisher", "claim-b", 3_000, 2_500)
+ .await
+ .expect("retry claim")
+ .expect("claimed");
+ assert_eq!(retry_claim.delivery_targets.len(), 1);
+ assert_ne!(
+ retry_claim.active_delivery_plan_id,
+ first_claim.active_delivery_plan_id
+ );
+ outbox
+ .mark_delivery_target_accepted(
+ first.outbox_event_id,
+ "claim-b",
+ retry_claim.delivery_targets[0].delivery_target_id,
+ 2_600,
+ )
+ .await
+ .expect("second active target accepted");
+ let final_state = outbox
+ .complete_publish_attempt(
+ first.outbox_event_id,
+ "claim-b",
+ "retryable",
+ "terminal",
+ 3_500,
+ 2_700,
+ )
+ .await
+ .expect("complete second plan");
+ assert_eq!(final_state, RadrootsOutboxEventState::Published);
+ }
+
+ #[tokio::test]
async fn quorum_accepted_delivery_plan_round_trips_and_completes_after_required_target() {
let outbox = RadrootsOutbox::open_memory().await.expect("open");
let draft = post_draft(FIXTURE_ALICE_PUBLIC_KEY_HEX, "at least");
@@ -2488,6 +2815,21 @@ mod tests {
plans[0].status,
RadrootsOutboxDeliveryPlanStatus::PreviewUnavailable
);
+ let event = outbox
+ .get_event(receipt.outbox_event_id)
+ .await
+ .expect("event")
+ .expect("event");
+ assert_eq!(event.state, RadrootsOutboxEventState::PreviewUnavailable);
+ let operation = outbox
+ .get_operation(receipt.operation_id)
+ .await
+ .expect("operation")
+ .expect("operation");
+ assert_eq!(
+ operation.status,
+ RadrootsOutboxOperationStatus::PreviewUnavailable
+ );
let summary = outbox.status_summary(1_000).await.expect("summary");
assert_eq!(summary.pending_events, 0);
assert_eq!(summary.ready_signed_events, 0);
@@ -2545,6 +2887,24 @@ mod tests {
plans[0].status,
RadrootsOutboxDeliveryPlanStatus::DeferredUntilImplemented
);
+ let event = outbox
+ .get_event(receipt.outbox_event_id)
+ .await
+ .expect("event")
+ .expect("event");
+ assert_eq!(
+ event.state,
+ RadrootsOutboxEventState::DeferredUntilImplemented
+ );
+ let operation = outbox
+ .get_operation(receipt.operation_id)
+ .await
+ .expect("operation")
+ .expect("operation");
+ assert_eq!(
+ operation.status,
+ RadrootsOutboxOperationStatus::DeferredUntilImplemented
+ );
let summary = outbox.status_summary(1_000).await.expect("summary");
assert_eq!(summary.pending_events, 0);
assert_eq!(summary.ready_signed_events, 0);
@@ -2585,6 +2945,77 @@ mod tests {
}
#[tokio::test]
+ async fn complete_signing_sets_preview_unavailable_lifecycle_for_reticulum_only_plans() {
+ let outbox = RadrootsOutbox::open_memory().await.expect("open");
+ let draft = post_draft(FIXTURE_ALICE_PUBLIC_KEY_HEX, "reticulum after signing");
+ let receipt = outbox
+ .enqueue_operation(RadrootsOutboxOperationInput::new(
+ "publish_post",
+ draft,
+ RadrootsOutboxDeliveryPlanInput::new(
+ "transport.reticulum.preview",
+ 1,
+ RadrootsTransportSatisfactionPolicy::all_accepted(),
+ vec![reticulum_target("reticulum:preview-unavailable")],
+ ),
+ 1_000,
+ ))
+ .await
+ .expect("enqueue");
+ let claimed = outbox
+ .claim_next_ready_event("signer", "claim-a", 2_000, 1_000)
+ .await
+ .expect("claim")
+ .expect("claimed");
+ assert_eq!(claimed.active_delivery_plan_id, None);
+ let signed =
+ radroots_nostr_sign_frozen_draft(&fixture_keys(), &claimed.draft).expect("signed");
+
+ outbox
+ .complete_signing(receipt.outbox_event_id, "claim-a", signed, 1_100)
+ .await
+ .expect("complete signing");
+
+ let event = outbox
+ .get_event(receipt.outbox_event_id)
+ .await
+ .expect("event")
+ .expect("event");
+ assert_eq!(event.state, RadrootsOutboxEventState::PreviewUnavailable);
+ assert_eq!(event.claim_token.as_deref(), Some("claim-a"));
+ let operation = outbox
+ .get_operation(receipt.operation_id)
+ .await
+ .expect("operation")
+ .expect("operation");
+ assert_eq!(
+ operation.status,
+ RadrootsOutboxOperationStatus::PreviewUnavailable
+ );
+ assert_eq!(
+ outbox.recover_expired_claims(2_001).await.expect("recover"),
+ 1
+ );
+ let recovered_event = outbox
+ .get_event(receipt.outbox_event_id)
+ .await
+ .expect("event")
+ .expect("event");
+ assert_eq!(
+ recovered_event.state,
+ RadrootsOutboxEventState::PreviewUnavailable
+ );
+ assert_eq!(recovered_event.claim_token, None);
+ assert!(
+ outbox
+ .claim_next_ready_signed_event("publisher", "claim-b", 3_000, 2_100)
+ .await
+ .expect("claim")
+ .is_none()
+ );
+ }
+
+ #[tokio::test]
async fn hybrid_nostr_and_reticulum_preview_preserves_ready_nostr_work() {
let outbox = RadrootsOutbox::open_memory().await.expect("open");
let draft = post_draft(FIXTURE_ALICE_PUBLIC_KEY_HEX, "hybrid preview");
@@ -2644,18 +3075,29 @@ mod tests {
#[tokio::test]
async fn claimed_delivery_targets_can_complete_as_preview_outcomes() {
- for (content, mark, expected_status, expected_plan_status) in [
+ for (
+ content,
+ mark,
+ expected_status,
+ expected_plan_status,
+ expected_event_state,
+ expected_operation_status,
+ ) in [
(
"proxy deferred",
"deferred",
RadrootsOutboxDeliveryTargetStatus::DeferredUntilImplemented,
RadrootsOutboxDeliveryPlanStatus::DeferredUntilImplemented,
+ RadrootsOutboxEventState::DeferredUntilImplemented,
+ RadrootsOutboxOperationStatus::DeferredUntilImplemented,
),
(
"proxy preview unavailable",
"preview",
RadrootsOutboxDeliveryTargetStatus::PreviewUnavailable,
RadrootsOutboxDeliveryPlanStatus::PreviewUnavailable,
+ RadrootsOutboxEventState::PreviewUnavailable,
+ RadrootsOutboxOperationStatus::PreviewUnavailable,
),
] {
let outbox = RadrootsOutbox::open_memory().await.expect("open");
@@ -2725,7 +3167,13 @@ mod tests {
)
.await
.expect("complete");
- assert_eq!(state, RadrootsOutboxEventState::Signed);
+ assert_eq!(state, expected_event_state);
+ let operation = outbox
+ .get_operation(receipt.operation_id)
+ .await
+ .expect("operation")
+ .expect("operation");
+ assert_eq!(operation.status, expected_operation_status);
let targets = outbox
.delivery_targets(receipt.outbox_event_id)
.await
@@ -2763,7 +3211,7 @@ mod tests {
ready_signed_events: 0,
publishing_events: 0,
last_attempt_at_ms: Some(1_100),
- last_error: Some("delivery deferred until implemented".to_owned()),
+ last_error: None,
}
);
}