lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

commit c70f70acbdb34ea3e6e91d602d36dcd119f3589c
parent 061bf15d7fcab4834a0c87740f43f2fef6589228
Author: triesap <tyson@radroots.org>
Date:   Wed,  8 Jul 2026 00:25:32 +0000

outbox: remove event-wide publish failure helper

- delete the public terminal publish failure API
- keep cancellation as the only event-wide claimed-event path
- prove active-plan failure and cancellation semantics in tests
- add source-boundary coverage for removed helper names

Diffstat:
Mcrates/outbox/src/store.rs | 212++++++++++++++++++++++++++++++++++++++++++++++++-------------------------------
1 file changed, 129 insertions(+), 83 deletions(-)

diff --git a/crates/outbox/src/store.rs b/crates/outbox/src/store.rs @@ -965,39 +965,48 @@ impl RadrootsOutbox { Ok(event_state) } - pub async fn mark_publish_failed_terminal( + pub async fn cancel_claimed_event( &self, outbox_event_id: i64, claim_token: &str, - error: impl AsRef<str>, now_ms: i64, ) -> Result<(), RadrootsOutboxError> { - self.finish_claimed_event( - outbox_event_id, - claim_token, - RadrootsOutboxEventState::FailedTerminal, - RadrootsOutboxOperationStatus::FailedTerminal, - Some(error.as_ref()), - now_ms, + let mut tx = self.pool.begin().await?; + let row = claimed_event_identity_tx(&mut tx, outbox_event_id, claim_token).await?; + sqlx::query( + "UPDATE outbox_delivery_plan SET status = ?, updated_at_ms = ? WHERE outbox_event_id = ?", ) - .await - } + .bind(RadrootsOutboxDeliveryPlanStatus::Cancelled.as_str()) + .bind(now_ms) + .bind(outbox_event_id) + .execute(&mut *tx) + .await?; + let changed = sqlx::query( + "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::Cancelled.as_str()) + .bind(Option::<&str>::None) + .bind(now_ms) + .bind(now_ms) + .bind(outbox_event_id) + .bind(claim_token) + .execute(&mut *tx) + .await?; + if changed.rows_affected() == 0 { + return Err(RadrootsOutboxError::ClaimTokenMismatch { outbox_event_id }); + } - pub async fn cancel_claimed_event( - &self, - outbox_event_id: i64, - claim_token: &str, - now_ms: i64, - ) -> Result<(), RadrootsOutboxError> { - self.finish_claimed_event( - outbox_event_id, - claim_token, - RadrootsOutboxEventState::Cancelled, - RadrootsOutboxOperationStatus::Cancelled, - None, - now_ms, + sqlx::query( + "UPDATE outbox_operations SET status = ?, updated_at_ms = ? WHERE operation_id = ?", ) - .await + .bind(RadrootsOutboxOperationStatus::Cancelled.as_str()) + .bind(now_ms) + .bind(row.operation_id) + .execute(&mut *tx) + .await?; + + tx.commit().await?; + Ok(()) } async fn claimed_event( @@ -1045,64 +1054,6 @@ impl RadrootsOutbox { Err(RadrootsOutboxError::ClaimTokenMismatch { outbox_event_id }) } - async fn finish_claimed_event( - &self, - outbox_event_id: i64, - claim_token: &str, - event_state: RadrootsOutboxEventState, - operation_status: RadrootsOutboxOperationStatus, - last_error: Option<&str>, - now_ms: i64, - ) -> Result<(), RadrootsOutboxError> { - let mut tx = self.pool.begin().await?; - let row = claimed_event_identity_tx(&mut tx, outbox_event_id, claim_token).await?; - let plan_status = match event_state { - RadrootsOutboxEventState::Cancelled => { - Some(RadrootsOutboxDeliveryPlanStatus::Cancelled) - } - RadrootsOutboxEventState::FailedTerminal => { - Some(RadrootsOutboxDeliveryPlanStatus::FailedTerminal) - } - _ => None, - }; - if let Some(plan_status) = plan_status { - sqlx::query( - "UPDATE outbox_delivery_plan SET status = ?, updated_at_ms = ? WHERE outbox_event_id = ?", - ) - .bind(plan_status.as_str()) - .bind(now_ms) - .bind(outbox_event_id) - .execute(&mut *tx) - .await?; - } - let changed = sqlx::query( - "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) - .bind(now_ms) - .bind(now_ms) - .bind(outbox_event_id) - .bind(claim_token) - .execute(&mut *tx) - .await?; - if changed.rows_affected() == 0 { - return Err(RadrootsOutboxError::ClaimTokenMismatch { outbox_event_id }); - } - - sqlx::query( - "UPDATE outbox_operations SET status = ?, updated_at_ms = ? WHERE operation_id = ?", - ) - .bind(operation_status.as_str()) - .bind(now_ms) - .bind(row.operation_id) - .execute(&mut *tx) - .await?; - - tx.commit().await?; - Ok(()) - } - async fn mark_delivery_target_status( &self, outbox_event_id: i64, @@ -2355,6 +2306,17 @@ mod tests { .expect("table count") } + #[test] + fn event_wide_publish_failure_helper_names_stay_removed() { + let source = include_str!("store.rs"); + for removed in [ + ["mark_publish", "_failed_terminal"].concat(), + ["finish_claimed", "_event"].concat(), + ] { + assert!(!source.contains(removed.as_str()), "{removed}"); + } + } + #[tokio::test] async fn migration_applies_delivery_plan_schema_and_migrates_down() { let outbox = RadrootsOutbox::open_memory().await.expect("open"); @@ -2908,6 +2870,90 @@ mod tests { } #[tokio::test] + async fn cancel_claimed_event_marks_every_delivery_plan_cancelled() { + let outbox = RadrootsOutbox::open_memory().await.expect("open"); + let draft = post_draft(FIXTURE_ALICE_PUBLIC_KEY_HEX, "cancel all 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.cancel", + 1, + RadrootsTransportSatisfactionPolicy::all_accepted(), + vec![nostr_target(NOSTR_PRIMARY_WSS)], + ), + true, + 1_007, + 1_000, + ) + .with_idempotency_key("idem-cancel-all-plans"), + ) + .await + .expect("first plan"); + let second = outbox + .enqueue_signed_operation( + RadrootsOutboxSignedOperationInput::new( + "publish_post", + draft, + signed_event, + RadrootsOutboxDeliveryPlanInput::new( + "transport.nostr.cancel.sibling", + 1, + RadrootsTransportSatisfactionPolicy::all_accepted(), + vec![nostr_target(NOSTR_SECONDARY_WSS)], + ), + true, + 1_017, + 1_010, + ) + .with_idempotency_key("idem-cancel-all-plans"), + ) + .await + .expect("second plan"); + assert_eq!(first.outbox_event_id, second.outbox_event_id); + + outbox + .claim_next_ready_signed_event("publisher", "claim-a", 2_000, 1_020) + .await + .expect("claim") + .expect("claimed"); + outbox + .cancel_claimed_event(first.outbox_event_id, "claim-a", 1_100) + .await + .expect("cancel"); + + let event = outbox + .get_event(first.outbox_event_id) + .await + .expect("event") + .expect("event"); + assert_eq!(event.state, RadrootsOutboxEventState::Cancelled); + assert_eq!(event.claim_token, None); + assert_eq!(event.active_delivery_plan_id, None); + let operation = outbox + .get_operation(first.operation_id) + .await + .expect("operation") + .expect("operation"); + assert_eq!(operation.status, RadrootsOutboxOperationStatus::Cancelled); + let plans = outbox + .delivery_plans(first.outbox_event_id) + .await + .expect("plans"); + assert_eq!(plans.len(), 2); + assert!( + plans + .iter() + .all(|plan| plan.status == RadrootsOutboxDeliveryPlanStatus::Cancelled) + ); + } + + #[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");