lib

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

commit e4e8ae87f6c230c4f26c317bd3107a6710e936a0
parent 2a3090443e763a208a84ab8e2b0ce8ca0bf323bd
Author: triesap <tyson@radroots.org>
Date:   Thu, 16 Jul 2026 09:44:37 +0000

store: add semantic trade mutation storage

- replace order projection storage with immutable trade mutation rows
- persist transport envelopes missing parents and seller reservations
- add semantic trade mutation outbox metadata and idempotency
- cover storage outbox and reservation contracts with tests

Diffstat:
MCargo.lock | 2++
Mcrates/event_store/Cargo.toml | 2++
Mcrates/event_store/migrations/0001_event_store.down.sql | 9++++++++-
Mcrates/event_store/migrations/0001_event_store.up.sql | 181++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------------
Mcrates/event_store/src/lib.rs | 6+++++-
Mcrates/event_store/src/model.rs | 104+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/event_store/src/store.rs | 1183+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------
Mcrates/outbox/migrations/0001_outbox.up.sql | 20++++++++++++++++----
Mcrates/outbox/src/error.rs | 11+++++++++++
Mcrates/outbox/src/lib.rs | 3++-
Mcrates/outbox/src/model.rs | 115++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------------
Mcrates/outbox/src/store.rs | 989+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------
12 files changed, 2415 insertions(+), 210 deletions(-)

diff --git a/Cargo.lock b/Cargo.lock @@ -4326,11 +4326,13 @@ dependencies = [ name = "radroots_event_store" version = "0.1.0-alpha.2" dependencies = [ + "hex", "radroots_event", "radroots_nostr", "radroots_transport", "serde", "serde_json", + "sha2", "sqlx", "tempfile", "thiserror 1.0.69", diff --git a/crates/event_store/Cargo.toml b/crates/event_store/Cargo.toml @@ -26,8 +26,10 @@ radroots_nostr = { workspace = true, default-features = false, features = [ "events", ] } radroots_transport = { workspace = true, default-features = false } +hex = { workspace = true } serde = { workspace = true, features = ["std"] } serde_json = { workspace = true, features = ["std"] } +sha2 = { workspace = true, features = ["std"] } sqlx = { workspace = true, optional = true, features = ["derive"] } thiserror = { workspace = true } diff --git a/crates/event_store/migrations/0001_event_store.down.sql b/crates/event_store/migrations/0001_event_store.down.sql @@ -1,5 +1,12 @@ DROP TABLE listing_search_fts; -DROP TABLE trade_projection; +DROP TABLE trade_projection_quarantine; +DROP TABLE trade_projection_checkpoint; +DROP TABLE seller_inventory_reservation_line; +DROP TABLE seller_inventory_reservation; +DROP TABLE trade_transport_envelope; +DROP TABLE trade_missing_parent; +DROP TABLE trade_mutation_parent; +DROP TABLE trade_mutation; DROP TABLE listing_projection; DROP TABLE projection_cursor; DROP TABLE event_envelope_head; diff --git a/crates/event_store/migrations/0001_event_store.up.sql b/crates/event_store/migrations/0001_event_store.up.sql @@ -130,45 +130,150 @@ CREATE VIRTUAL TABLE IF NOT EXISTS listing_search_fts USING fts5( tokenize = 'unicode61' ); -CREATE TABLE IF NOT EXISTS trade_projection ( - order_id TEXT NOT NULL, - root_event_id TEXT NOT NULL REFERENCES event_envelopes(event_id) ON DELETE CASCADE, - projection_version INTEGER NOT NULL, - status TEXT NOT NULL, - lifecycle_terminal INTEGER NOT NULL, - rhi_state TEXT NOT NULL, - listing_addr TEXT, - buyer_pubkey TEXT, - seller_pubkey TEXT, - request_event_id TEXT REFERENCES event_envelopes(event_id) ON DELETE SET NULL, - decision_event_id TEXT REFERENCES event_envelopes(event_id) ON DELETE SET NULL, - agreement_event_id TEXT REFERENCES event_envelopes(event_id) ON DELETE SET NULL, - cancellation_event_id TEXT REFERENCES event_envelopes(event_id) ON DELETE SET NULL, - validation_receipt_event_id TEXT REFERENCES event_envelopes(event_id) ON DELETE SET NULL, - last_event_id TEXT REFERENCES event_envelopes(event_id) ON DELETE SET NULL, - expected_listing_event_id TEXT, - current_listing_event_id TEXT REFERENCES event_envelopes(event_id) ON DELETE SET NULL, - economics_json TEXT, - pending_inventory_json TEXT NOT NULL, - committed_inventory_json TEXT NOT NULL, - issues_json TEXT NOT NULL, - issue_count INTEGER NOT NULL, - source_event_count INTEGER NOT NULL, - transport_observation_count INTEGER NOT NULL, - evidence_hash TEXT NOT NULL, - last_source_event_seq INTEGER, - updated_at_ms INTEGER NOT NULL, - PRIMARY KEY(order_id, root_event_id, projection_version) -); +CREATE TABLE IF NOT EXISTS trade_mutation ( + mutation_id TEXT PRIMARY KEY NOT NULL, + trade_id TEXT NOT NULL, + root_mutation_id TEXT, + contract_id TEXT NOT NULL, + mutation_kind TEXT NOT NULL CHECK (mutation_kind IN ('proposal', 'decision', 'revision_proposal', 'revision_decision', 'cancellation')), + schema_version INTEGER NOT NULL, + candidate_id TEXT, + proposal_mutation_id TEXT, + target_claim_mutation_id TEXT, + author_pubkey TEXT NOT NULL, + counterparty_pubkey TEXT NOT NULL, + buyer_pubkey TEXT NOT NULL, + seller_pubkey TEXT NOT NULL, + farm_id TEXT NOT NULL, + authored_at_unix_s INTEGER NOT NULL, + canonical_payload_bytes BLOB NOT NULL, + payload_sha256 TEXT NOT NULL CHECK (length(payload_sha256) = 64), + first_event_seq INTEGER NOT NULL REFERENCES event_envelopes(seq) ON DELETE RESTRICT, + first_transport_event_id TEXT NOT NULL REFERENCES event_envelopes(event_id) ON DELETE RESTRICT, + inserted_at_ms INTEGER NOT NULL +) STRICT; + +CREATE INDEX IF NOT EXISTS trade_mutation_trade_idx +ON trade_mutation(trade_id, authored_at_unix_s, mutation_id); + +CREATE INDEX IF NOT EXISTS trade_mutation_candidate_idx +ON trade_mutation(trade_id, candidate_id, mutation_id) +WHERE candidate_id IS NOT NULL; + +CREATE INDEX IF NOT EXISTS trade_mutation_actor_idx +ON trade_mutation(buyer_pubkey, seller_pubkey, authored_at_unix_s, mutation_id); + +CREATE TABLE IF NOT EXISTS trade_mutation_parent ( + mutation_id TEXT NOT NULL REFERENCES trade_mutation(mutation_id) ON DELETE CASCADE, + parent_mutation_id TEXT NOT NULL, + parent_index INTEGER NOT NULL, + PRIMARY KEY(mutation_id, parent_mutation_id) +) STRICT; + +CREATE INDEX IF NOT EXISTS trade_mutation_parent_lookup_idx +ON trade_mutation_parent(parent_mutation_id, mutation_id); + +CREATE TABLE IF NOT EXISTS trade_missing_parent ( + trade_id TEXT NOT NULL, + mutation_id TEXT NOT NULL REFERENCES trade_mutation(mutation_id) ON DELETE CASCADE, + missing_parent_mutation_id TEXT NOT NULL, + first_transport_event_id TEXT NOT NULL REFERENCES event_envelopes(event_id) ON DELETE CASCADE, + first_seen_at_ms INTEGER NOT NULL, + PRIMARY KEY(trade_id, mutation_id, missing_parent_mutation_id) +) STRICT; + +CREATE INDEX IF NOT EXISTS trade_missing_parent_lookup_idx +ON trade_missing_parent(missing_parent_mutation_id, trade_id, mutation_id); + +CREATE TABLE IF NOT EXISTS trade_transport_envelope ( + transport_event_id TEXT PRIMARY KEY NOT NULL REFERENCES event_envelopes(event_id) ON DELETE CASCADE, + mutation_id TEXT NOT NULL REFERENCES trade_mutation(mutation_id) ON DELETE CASCADE, + trade_id TEXT NOT NULL, + transport_kind TEXT NOT NULL, + pubkey TEXT NOT NULL, + created_at INTEGER NOT NULL, + event_seq INTEGER NOT NULL REFERENCES event_envelopes(seq) ON DELETE CASCADE, + payload_sha256 TEXT NOT NULL CHECK (length(payload_sha256) = 64), + observed_at_ms INTEGER NOT NULL +) STRICT; + +CREATE INDEX IF NOT EXISTS trade_transport_envelope_mutation_idx +ON trade_transport_envelope(mutation_id, observed_at_ms, transport_event_id); + +CREATE INDEX IF NOT EXISTS trade_transport_envelope_trade_idx +ON trade_transport_envelope(trade_id, event_seq, transport_event_id); + +CREATE TABLE IF NOT EXISTS seller_inventory_reservation ( + reservation_id TEXT PRIMARY KEY NOT NULL, + trade_id TEXT NOT NULL, + candidate_id TEXT NOT NULL, + claim_mutation_id TEXT NOT NULL REFERENCES trade_mutation(mutation_id) ON DELETE CASCADE, + inventory_authority_pubkey TEXT NOT NULL, + inventory_epoch INTEGER NOT NULL, + assertion_commitment TEXT NOT NULL CHECK (length(assertion_commitment) = 64), + reservation_expires_at_unix_s INTEGER NOT NULL, + reservation_json TEXT NOT NULL, + inserted_at_ms INTEGER NOT NULL, + UNIQUE(candidate_id, assertion_commitment) +) STRICT; + +CREATE INDEX IF NOT EXISTS seller_inventory_reservation_trade_idx +ON seller_inventory_reservation(trade_id, candidate_id, reservation_expires_at_unix_s); + +CREATE INDEX IF NOT EXISTS seller_inventory_reservation_authority_idx +ON seller_inventory_reservation(inventory_authority_pubkey, inventory_epoch, reservation_id); + +CREATE TABLE IF NOT EXISTS seller_inventory_reservation_line ( + reservation_id TEXT NOT NULL REFERENCES seller_inventory_reservation(reservation_id) ON DELETE CASCADE, + line_id TEXT NOT NULL, + bin_id TEXT NOT NULL, + quantity_mantissa TEXT NOT NULL, + quantity_scale INTEGER NOT NULL, + unit_code TEXT NOT NULL, + line_index INTEGER NOT NULL, + PRIMARY KEY(reservation_id, line_id) +) STRICT; + +CREATE INDEX IF NOT EXISTS seller_inventory_reservation_line_bin_idx +ON seller_inventory_reservation_line(bin_id, reservation_id, line_id); + +CREATE TABLE IF NOT EXISTS trade_projection_checkpoint ( + trade_id TEXT PRIMARY KEY NOT NULL, + reducer_contract_id TEXT NOT NULL, + reducer_version INTEGER NOT NULL, + projection_digest TEXT NOT NULL CHECK (length(projection_digest) = 64), + root_mutation_id TEXT, + negotiation_state TEXT NOT NULL, + agreement_state TEXT NOT NULL, + evidence_state TEXT NOT NULL, + conflict_state TEXT NOT NULL, + private_terms_state TEXT NOT NULL, + attestation_state TEXT NOT NULL, + fulfillment_state TEXT NOT NULL, + payment_state TEXT NOT NULL, + projection_json TEXT NOT NULL, + last_mutation_id TEXT, + last_transport_event_seq INTEGER, + updated_at_ms INTEGER NOT NULL +) STRICT; + +CREATE INDEX IF NOT EXISTS trade_projection_checkpoint_agreement_idx +ON trade_projection_checkpoint(agreement_state, updated_at_ms, trade_id); -CREATE INDEX IF NOT EXISTS trade_projection_status_idx -ON trade_projection(status, updated_at_ms, order_id, root_event_id); +CREATE INDEX IF NOT EXISTS trade_projection_checkpoint_actor_idx +ON trade_projection_checkpoint(root_mutation_id, updated_at_ms, trade_id); -CREATE INDEX IF NOT EXISTS trade_projection_listing_idx -ON trade_projection(listing_addr, updated_at_ms, order_id, root_event_id); +CREATE TABLE IF NOT EXISTS trade_projection_quarantine ( + quarantine_id INTEGER PRIMARY KEY AUTOINCREMENT, + trade_id TEXT, + mutation_id TEXT, + transport_event_id TEXT REFERENCES event_envelopes(event_id) ON DELETE CASCADE, + reason TEXT NOT NULL, + observed_at_ms INTEGER NOT NULL +) STRICT; -CREATE INDEX IF NOT EXISTS trade_projection_actor_idx -ON trade_projection(buyer_pubkey, seller_pubkey, updated_at_ms, order_id, root_event_id); +CREATE INDEX IF NOT EXISTS trade_projection_quarantine_trade_idx +ON trade_projection_quarantine(trade_id, observed_at_ms, quarantine_id); -CREATE INDEX IF NOT EXISTS trade_projection_root_idx -ON trade_projection(root_event_id, projection_version); +CREATE INDEX IF NOT EXISTS trade_projection_quarantine_mutation_idx +ON trade_projection_quarantine(mutation_id, observed_at_ms, quarantine_id); diff --git a/crates/event_store/src/lib.rs b/crates/event_store/src/lib.rs @@ -19,7 +19,11 @@ pub use model::{ RadrootsEventContractStatus, RadrootsEventHeadStoreDecision, RadrootsEventIngest, RadrootsEventIngestReceipt, RadrootsEventStoreStatusSummary, RadrootsEventVerificationStatus, RadrootsProjectionCursor, RadrootsStoredEvent, RadrootsStoredEventHead, RadrootsStoredEventTag, - RadrootsTransportObservation, RadrootsTransportObservationType, StoredEventClass, + RadrootsStoredSellerReservation, RadrootsStoredSellerReservationLine, + RadrootsStoredTradeMissingParent, RadrootsStoredTradeMutation, + RadrootsStoredTradeMutationParent, RadrootsStoredTradeTransportEnvelope, + RadrootsTradeProjectionCheckpoint, RadrootsTransportObservation, + RadrootsTransportObservationType, StoredEventClass, }; #[cfg(feature = "sqlite")] pub use store::{ diff --git a/crates/event_store/src/model.rs b/crates/event_store/src/model.rs @@ -5,6 +5,11 @@ use radroots_event::contract::{ }; use radroots_event::draft::RadrootsSignedEvent; use radroots_event::event_head::RadrootsEventHeadDecision; +use radroots_event::ids::{ + RadrootsDTag, RadrootsEventId, RadrootsInventoryBinId, RadrootsPublicKey, + RadrootsTradeCandidateId, RadrootsTradeId, RadrootsTradeMutationId, +}; +use radroots_event::trade::RadrootsTradeMutationKindV1; use radroots_event::wire::RadrootsNip01EventWire; use radroots_transport::{ RadrootsTransportKind, RadrootsTransportTargetFingerprint, RadrootsTransportTargetUri, @@ -354,6 +359,105 @@ pub struct RadrootsProjectionCursor { pub updated_at_ms: i64, } +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct RadrootsStoredTradeMutation { + pub mutation_id: RadrootsTradeMutationId, + pub trade_id: RadrootsTradeId, + pub root_mutation_id: Option<RadrootsTradeMutationId>, + pub contract_id: String, + pub mutation_kind: RadrootsTradeMutationKindV1, + pub schema_version: u16, + pub candidate_id: Option<RadrootsTradeCandidateId>, + pub proposal_mutation_id: Option<RadrootsTradeMutationId>, + pub target_claim_mutation_id: Option<RadrootsTradeMutationId>, + pub author_pubkey: RadrootsPublicKey, + pub counterparty_pubkey: RadrootsPublicKey, + pub buyer_pubkey: RadrootsPublicKey, + pub seller_pubkey: RadrootsPublicKey, + pub farm_id: RadrootsDTag, + pub authored_at_unix_s: u64, + pub canonical_payload_bytes: Vec<u8>, + pub payload_sha256: String, + pub first_event_seq: i64, + pub first_transport_event_id: RadrootsEventId, + pub inserted_at_ms: i64, +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct RadrootsStoredTradeMutationParent { + pub mutation_id: RadrootsTradeMutationId, + pub parent_mutation_id: RadrootsTradeMutationId, + pub parent_index: u32, +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct RadrootsStoredTradeMissingParent { + pub trade_id: RadrootsTradeId, + pub mutation_id: RadrootsTradeMutationId, + pub missing_parent_mutation_id: RadrootsTradeMutationId, + pub first_transport_event_id: RadrootsEventId, + pub first_seen_at_ms: i64, +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct RadrootsStoredTradeTransportEnvelope { + pub transport_event_id: RadrootsEventId, + pub mutation_id: RadrootsTradeMutationId, + pub trade_id: RadrootsTradeId, + pub transport_kind: String, + pub pubkey: RadrootsPublicKey, + pub created_at: u64, + pub event_seq: i64, + pub payload_sha256: String, + pub observed_at_ms: i64, +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct RadrootsStoredSellerReservation { + pub reservation_id: RadrootsDTag, + pub trade_id: RadrootsTradeId, + pub candidate_id: RadrootsTradeCandidateId, + pub claim_mutation_id: RadrootsTradeMutationId, + pub inventory_authority_pubkey: RadrootsPublicKey, + pub inventory_epoch: u64, + pub assertion_commitment: String, + pub reservation_expires_at_unix_s: u64, + pub reservation_json: String, + pub inserted_at_ms: i64, +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct RadrootsStoredSellerReservationLine { + pub reservation_id: RadrootsDTag, + pub line_id: RadrootsDTag, + pub bin_id: RadrootsInventoryBinId, + pub quantity_mantissa: String, + pub quantity_scale: u8, + pub unit_code: String, + pub line_index: u32, +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct RadrootsTradeProjectionCheckpoint { + pub trade_id: RadrootsTradeId, + pub reducer_contract_id: String, + pub reducer_version: u16, + pub projection_digest: String, + pub root_mutation_id: Option<RadrootsTradeMutationId>, + pub negotiation_state: String, + pub agreement_state: String, + pub evidence_state: String, + pub conflict_state: String, + pub private_terms_state: String, + pub attestation_state: String, + pub fulfillment_state: String, + pub payment_state: String, + pub projection_json: String, + pub last_mutation_id: Option<RadrootsTradeMutationId>, + pub last_transport_event_seq: Option<i64>, + pub updated_at_ms: i64, +} + pub fn tag_semantic_name(value: RadrootsTagSemantic) -> &'static str { match value { RadrootsTagSemantic::AddressableCoordinate => "addressable_coordinate", diff --git a/crates/event_store/src/store.rs b/crates/event_store/src/store.rs @@ -4,8 +4,11 @@ use crate::model::{ RadrootsEventContractStatus, RadrootsEventHeadStoreDecision, RadrootsEventIngest, RadrootsEventIngestReceipt, RadrootsEventStoreStatusSummary, RadrootsEventVerificationStatus, RadrootsProjectionCursor, RadrootsStoredEvent, RadrootsStoredEventHead, RadrootsStoredEventTag, - RadrootsTransportObservation, RadrootsTransportObservationType, StoredEventClass, - tag_semantic_name, tag_value_type_name, + RadrootsStoredSellerReservation, RadrootsStoredSellerReservationLine, + RadrootsStoredTradeMissingParent, RadrootsStoredTradeMutation, + RadrootsStoredTradeMutationParent, RadrootsStoredTradeTransportEnvelope, + RadrootsTradeProjectionCheckpoint, RadrootsTransportObservation, + RadrootsTransportObservationType, StoredEventClass, tag_semantic_name, tag_value_type_name, }; use radroots_event::RadrootsEventEnvelope; use radroots_event::contract::{ @@ -16,11 +19,20 @@ use radroots_event::event_head::{ RadrootsEventHeadCoordinate, RadrootsEventHeadDecision, event_head_candidate_for_contract, select_event_head, }; -use radroots_event::ids::{RadrootsEventId, RadrootsPublicKey}; +use radroots_event::ids::{ + RadrootsDTag, RadrootsEventId, RadrootsPublicKey, RadrootsTradeCandidateId, RadrootsTradeId, + RadrootsTradeMutationId, +}; +use radroots_event::trade::{ + RADROOTS_TRADE_MUTATION_CONTRACT_IDS, RadrootsSellerReservationAssertionV1, + RadrootsTradeDecisionV1, RadrootsTradeMutationBodyV1, RadrootsTradeMutationEnvelopeV1, + RadrootsTradeMutationKindV1, trade_mutation_from_canonical_content, +}; use radroots_nostr::prelude::{RadrootsNostrEventVerification, radroots_nostr_verify_event}; use radroots_transport::{ RadrootsTransportKind, RadrootsTransportTargetFingerprint, RadrootsTransportTargetUri, }; +use sha2::{Digest, Sha256}; use sqlx::sqlite::{SqliteConnectOptions, SqlitePoolOptions}; use sqlx::{Row, SqlitePool}; use std::path::Path; @@ -317,6 +329,150 @@ impl RadrootsEventStore { .await?; rows.into_iter().map(stored_event_from_row).collect() } + + pub async fn get_trade_mutation( + &self, + mutation_id: &RadrootsTradeMutationId, + ) -> Result<Option<RadrootsStoredTradeMutation>, RadrootsEventStoreError> { + let row = sqlx::query( + "SELECT mutation_id, trade_id, root_mutation_id, contract_id, mutation_kind, schema_version, candidate_id, proposal_mutation_id, target_claim_mutation_id, author_pubkey, counterparty_pubkey, buyer_pubkey, seller_pubkey, farm_id, authored_at_unix_s, canonical_payload_bytes, payload_sha256, first_event_seq, first_transport_event_id, inserted_at_ms FROM trade_mutation WHERE mutation_id = ?", + ) + .bind(mutation_id.as_str()) + .fetch_optional(&self.pool) + .await?; + row.map(trade_mutation_from_row).transpose() + } + + pub async fn trade_mutations_for_trade( + &self, + trade_id: &RadrootsTradeId, + limit: u32, + ) -> Result<Vec<RadrootsStoredTradeMutation>, RadrootsEventStoreError> { + validate_trade_query_limit(limit)?; + let rows = sqlx::query( + "SELECT mutation_id, trade_id, root_mutation_id, contract_id, mutation_kind, schema_version, candidate_id, proposal_mutation_id, target_claim_mutation_id, author_pubkey, counterparty_pubkey, buyer_pubkey, seller_pubkey, farm_id, authored_at_unix_s, canonical_payload_bytes, payload_sha256, first_event_seq, first_transport_event_id, inserted_at_ms FROM trade_mutation WHERE trade_id = ? ORDER BY authored_at_unix_s, mutation_id LIMIT ?", + ) + .bind(trade_id.as_str()) + .bind(i64::from(limit)) + .fetch_all(&self.pool) + .await?; + rows.into_iter().map(trade_mutation_from_row).collect() + } + + pub async fn trade_mutation_parents( + &self, + mutation_id: &RadrootsTradeMutationId, + ) -> Result<Vec<RadrootsStoredTradeMutationParent>, RadrootsEventStoreError> { + let rows = sqlx::query( + "SELECT mutation_id, parent_mutation_id, parent_index FROM trade_mutation_parent WHERE mutation_id = ? ORDER BY parent_index", + ) + .bind(mutation_id.as_str()) + .fetch_all(&self.pool) + .await?; + rows.into_iter() + .map(trade_mutation_parent_from_row) + .collect() + } + + pub async fn trade_transport_envelopes_for_mutation( + &self, + mutation_id: &RadrootsTradeMutationId, + ) -> Result<Vec<RadrootsStoredTradeTransportEnvelope>, RadrootsEventStoreError> { + let rows = sqlx::query( + "SELECT transport_event_id, mutation_id, trade_id, transport_kind, pubkey, created_at, event_seq, payload_sha256, observed_at_ms FROM trade_transport_envelope WHERE mutation_id = ? ORDER BY observed_at_ms, transport_event_id", + ) + .bind(mutation_id.as_str()) + .fetch_all(&self.pool) + .await?; + rows.into_iter() + .map(trade_transport_envelope_from_row) + .collect() + } + + pub async fn missing_trade_parents( + &self, + trade_id: &RadrootsTradeId, + ) -> Result<Vec<RadrootsStoredTradeMissingParent>, RadrootsEventStoreError> { + let rows = sqlx::query( + "SELECT trade_id, mutation_id, missing_parent_mutation_id, first_transport_event_id, first_seen_at_ms FROM trade_missing_parent WHERE trade_id = ? ORDER BY first_seen_at_ms, mutation_id, missing_parent_mutation_id", + ) + .bind(trade_id.as_str()) + .fetch_all(&self.pool) + .await?; + rows.into_iter() + .map(trade_missing_parent_from_row) + .collect() + } + + pub async fn seller_reservation( + &self, + reservation_id: &RadrootsDTag, + ) -> Result<Option<RadrootsStoredSellerReservation>, RadrootsEventStoreError> { + let row = sqlx::query( + "SELECT reservation_id, trade_id, candidate_id, claim_mutation_id, inventory_authority_pubkey, inventory_epoch, assertion_commitment, reservation_expires_at_unix_s, reservation_json, inserted_at_ms FROM seller_inventory_reservation WHERE reservation_id = ?", + ) + .bind(reservation_id.as_str()) + .fetch_optional(&self.pool) + .await?; + row.map(seller_reservation_from_row).transpose() + } + + pub async fn seller_reservation_lines( + &self, + reservation_id: &RadrootsDTag, + ) -> Result<Vec<RadrootsStoredSellerReservationLine>, RadrootsEventStoreError> { + let rows = sqlx::query( + "SELECT reservation_id, line_id, bin_id, quantity_mantissa, quantity_scale, unit_code, line_index FROM seller_inventory_reservation_line WHERE reservation_id = ? ORDER BY line_index", + ) + .bind(reservation_id.as_str()) + .fetch_all(&self.pool) + .await?; + rows.into_iter() + .map(seller_reservation_line_from_row) + .collect() + } + + pub async fn update_trade_projection_checkpoint( + &self, + checkpoint: &RadrootsTradeProjectionCheckpoint, + ) -> Result<(), RadrootsEventStoreError> { + sqlx::query( + "INSERT INTO trade_projection_checkpoint(trade_id, reducer_contract_id, reducer_version, projection_digest, root_mutation_id, negotiation_state, agreement_state, evidence_state, conflict_state, private_terms_state, attestation_state, fulfillment_state, payment_state, projection_json, last_mutation_id, last_transport_event_seq, updated_at_ms) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT(trade_id) DO UPDATE SET reducer_contract_id = excluded.reducer_contract_id, reducer_version = excluded.reducer_version, projection_digest = excluded.projection_digest, root_mutation_id = excluded.root_mutation_id, negotiation_state = excluded.negotiation_state, agreement_state = excluded.agreement_state, evidence_state = excluded.evidence_state, conflict_state = excluded.conflict_state, private_terms_state = excluded.private_terms_state, attestation_state = excluded.attestation_state, fulfillment_state = excluded.fulfillment_state, payment_state = excluded.payment_state, projection_json = excluded.projection_json, last_mutation_id = excluded.last_mutation_id, last_transport_event_seq = excluded.last_transport_event_seq, updated_at_ms = excluded.updated_at_ms", + ) + .bind(checkpoint.trade_id.as_str()) + .bind(checkpoint.reducer_contract_id.as_str()) + .bind(i64::from(checkpoint.reducer_version)) + .bind(checkpoint.projection_digest.as_str()) + .bind(checkpoint.root_mutation_id.as_ref().map(RadrootsTradeMutationId::as_str)) + .bind(checkpoint.negotiation_state.as_str()) + .bind(checkpoint.agreement_state.as_str()) + .bind(checkpoint.evidence_state.as_str()) + .bind(checkpoint.conflict_state.as_str()) + .bind(checkpoint.private_terms_state.as_str()) + .bind(checkpoint.attestation_state.as_str()) + .bind(checkpoint.fulfillment_state.as_str()) + .bind(checkpoint.payment_state.as_str()) + .bind(checkpoint.projection_json.as_str()) + .bind(checkpoint.last_mutation_id.as_ref().map(RadrootsTradeMutationId::as_str)) + .bind(checkpoint.last_transport_event_seq) + .bind(checkpoint.updated_at_ms) + .execute(&self.pool) + .await?; + Ok(()) + } + + pub async fn trade_projection_checkpoint( + &self, + trade_id: &RadrootsTradeId, + ) -> Result<Option<RadrootsTradeProjectionCheckpoint>, RadrootsEventStoreError> { + let row = sqlx::query( + "SELECT trade_id, reducer_contract_id, reducer_version, projection_digest, root_mutation_id, negotiation_state, agreement_state, evidence_state, conflict_state, private_terms_state, attestation_state, fulfillment_state, payment_state, projection_json, last_mutation_id, last_transport_event_seq, updated_at_ms FROM trade_projection_checkpoint WHERE trade_id = ?", + ) + .bind(trade_id.as_str()) + .fetch_optional(&self.pool) + .await?; + row.map(trade_projection_checkpoint_from_row).transpose() + } } #[derive(Clone, Debug, PartialEq, Eq)] @@ -426,6 +582,373 @@ fn classify_event(event: &RadrootsEventEnvelope) -> EventClassification { } } +fn is_trade_mutation_contract_id(contract_id: &str) -> bool { + RADROOTS_TRADE_MUTATION_CONTRACT_IDS.contains(&contract_id) +} + +async fn store_trade_mutation_event( + tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>, + ingest: &RadrootsEventIngest, + contract: &RadrootsEventContract, + event_seq: i64, +) -> Result<bool, RadrootsEventStoreError> { + let event = ingest.event(); + let payload_sha256 = sha256_hex(event.content().as_bytes()); + let parsed = match trade_mutation_from_canonical_content(event.content()) { + Ok(envelope) => envelope, + Err(error) => { + insert_trade_quarantine( + tx, + None, + None, + Some(event.id_str()), + format!("{error}").as_str(), + ingest.observed_at_ms, + ) + .await?; + return Ok(false); + } + }; + let Some(mutation_id) = parsed.mutation_id.clone() else { + insert_trade_quarantine( + tx, + Some(parsed.trade_id.as_str()), + None, + Some(event.id_str()), + "canonical trade mutation content is missing mutation_id", + ingest.observed_at_ms, + ) + .await?; + return Ok(false); + }; + if parsed.author_pubkey.as_str() != event.author_str() { + insert_trade_quarantine( + tx, + Some(parsed.trade_id.as_str()), + Some(mutation_id.as_str()), + Some(event.id_str()), + "trade mutation author_pubkey does not match transport event pubkey", + ingest.observed_at_ms, + ) + .await?; + return Ok(false); + } + let mutation_kind = parsed.mutation_kind(); + let candidate_id = candidate_id_for_mutation(&parsed); + let proposal_mutation_id = proposal_mutation_id_for_mutation(&parsed); + let target_claim_mutation_id = target_claim_mutation_id_for_mutation(&parsed); + sqlx::query( + "INSERT OR IGNORE INTO trade_mutation(mutation_id, trade_id, root_mutation_id, contract_id, mutation_kind, schema_version, candidate_id, proposal_mutation_id, target_claim_mutation_id, author_pubkey, counterparty_pubkey, buyer_pubkey, seller_pubkey, farm_id, authored_at_unix_s, canonical_payload_bytes, payload_sha256, first_event_seq, first_transport_event_id, inserted_at_ms) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", + ) + .bind(mutation_id.as_str()) + .bind(parsed.trade_id.as_str()) + .bind(parsed.root_mutation_id.as_ref().map(RadrootsTradeMutationId::as_str)) + .bind(parsed.contract_id.as_str()) + .bind(trade_mutation_kind_storage_value(mutation_kind)) + .bind(i64::from(parsed.schema_version)) + .bind(candidate_id.as_ref().map(RadrootsTradeCandidateId::as_str)) + .bind(proposal_mutation_id.as_ref().map(RadrootsTradeMutationId::as_str)) + .bind(target_claim_mutation_id.as_ref().map(RadrootsTradeMutationId::as_str)) + .bind(parsed.author_pubkey.as_str()) + .bind(parsed.counterparty_pubkey.as_str()) + .bind(parsed.buyer_pubkey.as_str()) + .bind(parsed.seller_pubkey.as_str()) + .bind(parsed.farm_id.as_str()) + .bind(i64_from_u64("authored_at_unix_s", parsed.authored_at_unix_s)?) + .bind(event.content().as_bytes()) + .bind(payload_sha256.as_str()) + .bind(event_seq) + .bind(event.id_str()) + .bind(ingest.observed_at_ms) + .execute(&mut **tx) + .await?; + insert_trade_mutation_parents(tx, &mutation_id, &parsed.parent_mutation_ids).await?; + insert_trade_transport_envelope( + tx, + event, + event_seq, + &parsed, + &mutation_id, + &payload_sha256, + ingest.observed_at_ms, + ) + .await?; + insert_missing_parent_records( + tx, + &parsed, + &mutation_id, + event.id_str(), + ingest.observed_at_ms, + ) + .await?; + delete_resolved_missing_parent_records(tx, &mutation_id).await?; + if let Some(reservation) = seller_reservation_for_mutation(&parsed) { + insert_seller_reservation( + tx, + &parsed, + &mutation_id, + reservation, + ingest.observed_at_ms, + ) + .await?; + } + let _ = contract; + Ok(true) +} + +async fn insert_trade_quarantine( + tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>, + trade_id: Option<&str>, + mutation_id: Option<&str>, + transport_event_id: Option<&str>, + reason: &str, + observed_at_ms: i64, +) -> Result<(), RadrootsEventStoreError> { + sqlx::query( + "INSERT INTO trade_projection_quarantine(trade_id, mutation_id, transport_event_id, reason, observed_at_ms) VALUES (?, ?, ?, ?, ?)", + ) + .bind(trade_id) + .bind(mutation_id) + .bind(transport_event_id) + .bind(reason) + .bind(observed_at_ms) + .execute(&mut **tx) + .await?; + Ok(()) +} + +async fn insert_trade_mutation_parents( + tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>, + mutation_id: &RadrootsTradeMutationId, + parents: &[RadrootsTradeMutationId], +) -> Result<(), RadrootsEventStoreError> { + for (index, parent) in parents.iter().enumerate() { + sqlx::query( + "INSERT OR IGNORE INTO trade_mutation_parent(mutation_id, parent_mutation_id, parent_index) VALUES (?, ?, ?)", + ) + .bind(mutation_id.as_str()) + .bind(parent.as_str()) + .bind(i64::try_from(index).map_err(|_| RadrootsEventStoreError::IntegerRange { + field: "parent_index", + value: i64::MAX, + })?) + .execute(&mut **tx) + .await?; + } + Ok(()) +} + +async fn insert_trade_transport_envelope( + tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>, + event: &RadrootsEventEnvelope, + event_seq: i64, + mutation: &RadrootsTradeMutationEnvelopeV1, + mutation_id: &RadrootsTradeMutationId, + payload_sha256: &str, + observed_at_ms: i64, +) -> Result<(), RadrootsEventStoreError> { + sqlx::query( + "INSERT OR IGNORE INTO trade_transport_envelope(transport_event_id, mutation_id, trade_id, transport_kind, pubkey, created_at, event_seq, payload_sha256, observed_at_ms) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)", + ) + .bind(event.id_str()) + .bind(mutation_id.as_str()) + .bind(mutation.trade_id.as_str()) + .bind(RadrootsTransportKind::Nostr.canonical_label()) + .bind(event.author_str()) + .bind(i64_from_u64("created_at", event.created_at_u64())?) + .bind(event_seq) + .bind(payload_sha256) + .bind(observed_at_ms) + .execute(&mut **tx) + .await?; + Ok(()) +} + +async fn insert_missing_parent_records( + tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>, + mutation: &RadrootsTradeMutationEnvelopeV1, + mutation_id: &RadrootsTradeMutationId, + transport_event_id: &str, + observed_at_ms: i64, +) -> Result<(), RadrootsEventStoreError> { + for parent in &mutation.parent_mutation_ids { + let exists: Option<i64> = + sqlx::query_scalar("SELECT 1 FROM trade_mutation WHERE mutation_id = ? LIMIT 1") + .bind(parent.as_str()) + .fetch_optional(&mut **tx) + .await?; + if exists.is_none() { + sqlx::query( + "INSERT OR IGNORE INTO trade_missing_parent(trade_id, mutation_id, missing_parent_mutation_id, first_transport_event_id, first_seen_at_ms) VALUES (?, ?, ?, ?, ?)", + ) + .bind(mutation.trade_id.as_str()) + .bind(mutation_id.as_str()) + .bind(parent.as_str()) + .bind(transport_event_id) + .bind(observed_at_ms) + .execute(&mut **tx) + .await?; + } + } + Ok(()) +} + +async fn delete_resolved_missing_parent_records( + tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>, + mutation_id: &RadrootsTradeMutationId, +) -> Result<(), RadrootsEventStoreError> { + sqlx::query("DELETE FROM trade_missing_parent WHERE missing_parent_mutation_id = ?") + .bind(mutation_id.as_str()) + .execute(&mut **tx) + .await?; + Ok(()) +} + +async fn insert_seller_reservation( + tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>, + mutation: &RadrootsTradeMutationEnvelopeV1, + claim_mutation_id: &RadrootsTradeMutationId, + reservation: &RadrootsSellerReservationAssertionV1, + inserted_at_ms: i64, +) -> Result<(), RadrootsEventStoreError> { + let reservation_json = serde_json::to_string(reservation)?; + sqlx::query( + "INSERT OR IGNORE INTO seller_inventory_reservation(reservation_id, trade_id, candidate_id, claim_mutation_id, inventory_authority_pubkey, inventory_epoch, assertion_commitment, reservation_expires_at_unix_s, reservation_json, inserted_at_ms) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", + ) + .bind(reservation.reservation_id.as_str()) + .bind(mutation.trade_id.as_str()) + .bind(reservation.candidate_id.as_str()) + .bind(claim_mutation_id.as_str()) + .bind(reservation.inventory_authority_id.as_str()) + .bind(i64_from_u64("inventory_epoch", reservation.inventory_epoch)?) + .bind(reservation.assertion_commitment.as_str()) + .bind(i64_from_u64( + "reservation_expires_at_unix_s", + reservation.reservation_expires_at_unix_s, + )?) + .bind(reservation_json.as_str()) + .bind(inserted_at_ms) + .execute(&mut **tx) + .await?; + for (index, line) in reservation.commitments.iter().enumerate() { + sqlx::query( + "INSERT OR IGNORE INTO seller_inventory_reservation_line(reservation_id, line_id, bin_id, quantity_mantissa, quantity_scale, unit_code, line_index) VALUES (?, ?, ?, ?, ?, ?, ?)", + ) + .bind(reservation.reservation_id.as_str()) + .bind(line.line_id.as_str()) + .bind(line.bin_id.as_str()) + .bind(line.quantity_mantissa.as_str()) + .bind(i64::from(line.quantity_scale)) + .bind(line.unit_code.as_str()) + .bind(i64::try_from(index).map_err(|_| RadrootsEventStoreError::IntegerRange { + field: "reservation.line_index", + value: i64::MAX, + })?) + .execute(&mut **tx) + .await?; + } + Ok(()) +} + +fn candidate_id_for_mutation( + mutation: &RadrootsTradeMutationEnvelopeV1, +) -> Option<RadrootsTradeCandidateId> { + match &mutation.body { + RadrootsTradeMutationBodyV1::Proposal { candidate } + | RadrootsTradeMutationBodyV1::RevisionProposal { candidate } => { + candidate.candidate_id.clone() + } + RadrootsTradeMutationBodyV1::Decision { candidate_id, .. } + | RadrootsTradeMutationBodyV1::RevisionDecision { candidate_id, .. } => { + Some(candidate_id.clone()) + } + RadrootsTradeMutationBodyV1::Cancellation { + target_candidate_id, + .. + } => target_candidate_id.clone(), + } +} + +fn proposal_mutation_id_for_mutation( + mutation: &RadrootsTradeMutationEnvelopeV1, +) -> Option<RadrootsTradeMutationId> { + match &mutation.body { + RadrootsTradeMutationBodyV1::Decision { + proposal_mutation_id, + .. + } + | RadrootsTradeMutationBodyV1::RevisionDecision { + proposal_mutation_id, + .. + } => Some(proposal_mutation_id.clone()), + _ => None, + } +} + +fn target_claim_mutation_id_for_mutation( + mutation: &RadrootsTradeMutationEnvelopeV1, +) -> Option<RadrootsTradeMutationId> { + match &mutation.body { + RadrootsTradeMutationBodyV1::Cancellation { + target_claim_mutation_id, + .. + } => target_claim_mutation_id.clone(), + _ => None, + } +} + +fn seller_reservation_for_mutation( + mutation: &RadrootsTradeMutationEnvelopeV1, +) -> Option<&RadrootsSellerReservationAssertionV1> { + match &mutation.body { + RadrootsTradeMutationBodyV1::Decision { + decision: + RadrootsTradeDecisionV1::Accepted { + reservation_assertion: Some(reservation), + }, + .. + } + | RadrootsTradeMutationBodyV1::RevisionDecision { + decision: + RadrootsTradeDecisionV1::Accepted { + reservation_assertion: Some(reservation), + }, + .. + } => Some(reservation), + _ => None, + } +} + +fn trade_mutation_kind_storage_value(kind: RadrootsTradeMutationKindV1) -> &'static str { + match kind { + RadrootsTradeMutationKindV1::Proposal => "proposal", + RadrootsTradeMutationKindV1::Decision => "decision", + RadrootsTradeMutationKindV1::RevisionProposal => "revision_proposal", + RadrootsTradeMutationKindV1::RevisionDecision => "revision_decision", + RadrootsTradeMutationKindV1::Cancellation => "cancellation", + } +} + +fn parse_trade_mutation_kind( + value: &str, +) -> Result<RadrootsTradeMutationKindV1, RadrootsEventStoreError> { + match value { + "proposal" => Ok(RadrootsTradeMutationKindV1::Proposal), + "decision" => Ok(RadrootsTradeMutationKindV1::Decision), + "revision_proposal" => Ok(RadrootsTradeMutationKindV1::RevisionProposal), + "revision_decision" => Ok(RadrootsTradeMutationKindV1::RevisionDecision), + "cancellation" => Ok(RadrootsTradeMutationKindV1::Cancellation), + _ => Err(RadrootsEventStoreError::InvalidStoredEnum { + field: "trade_mutation.mutation_kind", + value: value.to_owned(), + }), + } +} + +fn sha256_hex(bytes: &[u8]) -> String { + hex::encode(Sha256::digest(bytes)) +} + fn verify_event(event: &RadrootsEventEnvelope) -> RadrootsEventVerificationStatus { verification_status_from_nostr(radroots_nostr_verify_event(event)) } @@ -474,9 +997,19 @@ async fn ingest_event_in_transaction( insert_tags(tx, event, classification.contract).await?; if let Some(contract) = classification.contract { if projection_eligible { - let head = apply_event_head(tx, event, contract, ingest.observed_at_ms).await?; - projection_eligible = head.projection_eligible; - head_decision = head.decision; + if is_trade_mutation_contract_id(contract.id) { + projection_eligible = + store_trade_mutation_event(tx, &ingest, contract, insert.seq).await?; + head_decision = if projection_eligible { + RadrootsEventHeadStoreDecision::NotHeadSelected + } else { + RadrootsEventHeadStoreDecision::Malformed + }; + } else { + let head = apply_event_head(tx, event, contract, ingest.observed_at_ms).await?; + projection_eligible = head.projection_eligible; + head_decision = head.decision; + } sqlx::query( "UPDATE event_envelopes SET projection_eligible = ?, updated_at_ms = ? WHERE event_id = ?", ) @@ -840,6 +1373,142 @@ fn projection_cursor_from_row( } #[cfg_attr(coverage_nightly, coverage(off))] +fn trade_mutation_from_row( + row: sqlx::sqlite::SqliteRow, +) -> Result<RadrootsStoredTradeMutation, RadrootsEventStoreError> { + Ok(RadrootsStoredTradeMutation { + mutation_id: parse_id(row.try_get::<String, _>("mutation_id")?)?, + trade_id: parse_id(row.try_get::<String, _>("trade_id")?)?, + root_mutation_id: parse_optional_id(row.try_get("root_mutation_id")?)?, + contract_id: row.try_get("contract_id")?, + mutation_kind: parse_trade_mutation_kind( + row.try_get::<String, _>("mutation_kind")?.as_str(), + )?, + schema_version: u16_from_i64("schema_version", row.try_get("schema_version")?)?, + candidate_id: parse_optional_id(row.try_get("candidate_id")?)?, + proposal_mutation_id: parse_optional_id(row.try_get("proposal_mutation_id")?)?, + target_claim_mutation_id: parse_optional_id(row.try_get("target_claim_mutation_id")?)?, + author_pubkey: parse_id(row.try_get::<String, _>("author_pubkey")?)?, + counterparty_pubkey: parse_id(row.try_get::<String, _>("counterparty_pubkey")?)?, + buyer_pubkey: parse_id(row.try_get::<String, _>("buyer_pubkey")?)?, + seller_pubkey: parse_id(row.try_get::<String, _>("seller_pubkey")?)?, + farm_id: parse_id(row.try_get::<String, _>("farm_id")?)?, + authored_at_unix_s: u64_from_i64("authored_at_unix_s", row.try_get("authored_at_unix_s")?)?, + canonical_payload_bytes: row.try_get("canonical_payload_bytes")?, + payload_sha256: row.try_get("payload_sha256")?, + first_event_seq: row.try_get("first_event_seq")?, + first_transport_event_id: parse_id(row.try_get::<String, _>("first_transport_event_id")?)?, + inserted_at_ms: row.try_get("inserted_at_ms")?, + }) +} + +#[cfg_attr(coverage_nightly, coverage(off))] +fn trade_mutation_parent_from_row( + row: sqlx::sqlite::SqliteRow, +) -> Result<RadrootsStoredTradeMutationParent, RadrootsEventStoreError> { + Ok(RadrootsStoredTradeMutationParent { + mutation_id: parse_id(row.try_get::<String, _>("mutation_id")?)?, + parent_mutation_id: parse_id(row.try_get::<String, _>("parent_mutation_id")?)?, + parent_index: u32_from_i64("parent_index", row.try_get("parent_index")?)?, + }) +} + +#[cfg_attr(coverage_nightly, coverage(off))] +fn trade_missing_parent_from_row( + row: sqlx::sqlite::SqliteRow, +) -> Result<RadrootsStoredTradeMissingParent, RadrootsEventStoreError> { + Ok(RadrootsStoredTradeMissingParent { + trade_id: parse_id(row.try_get::<String, _>("trade_id")?)?, + mutation_id: parse_id(row.try_get::<String, _>("mutation_id")?)?, + missing_parent_mutation_id: parse_id( + row.try_get::<String, _>("missing_parent_mutation_id")?, + )?, + first_transport_event_id: parse_id(row.try_get::<String, _>("first_transport_event_id")?)?, + first_seen_at_ms: row.try_get("first_seen_at_ms")?, + }) +} + +#[cfg_attr(coverage_nightly, coverage(off))] +fn trade_transport_envelope_from_row( + row: sqlx::sqlite::SqliteRow, +) -> Result<RadrootsStoredTradeTransportEnvelope, RadrootsEventStoreError> { + Ok(RadrootsStoredTradeTransportEnvelope { + transport_event_id: parse_id(row.try_get::<String, _>("transport_event_id")?)?, + mutation_id: parse_id(row.try_get::<String, _>("mutation_id")?)?, + trade_id: parse_id(row.try_get::<String, _>("trade_id")?)?, + transport_kind: row.try_get("transport_kind")?, + pubkey: parse_id(row.try_get::<String, _>("pubkey")?)?, + created_at: u64_from_i64("created_at", row.try_get("created_at")?)?, + event_seq: row.try_get("event_seq")?, + payload_sha256: row.try_get("payload_sha256")?, + observed_at_ms: row.try_get("observed_at_ms")?, + }) +} + +#[cfg_attr(coverage_nightly, coverage(off))] +fn seller_reservation_from_row( + row: sqlx::sqlite::SqliteRow, +) -> Result<RadrootsStoredSellerReservation, RadrootsEventStoreError> { + Ok(RadrootsStoredSellerReservation { + reservation_id: parse_id(row.try_get::<String, _>("reservation_id")?)?, + trade_id: parse_id(row.try_get::<String, _>("trade_id")?)?, + candidate_id: parse_id(row.try_get::<String, _>("candidate_id")?)?, + claim_mutation_id: parse_id(row.try_get::<String, _>("claim_mutation_id")?)?, + inventory_authority_pubkey: parse_id( + row.try_get::<String, _>("inventory_authority_pubkey")?, + )?, + inventory_epoch: u64_from_i64("inventory_epoch", row.try_get("inventory_epoch")?)?, + assertion_commitment: row.try_get("assertion_commitment")?, + reservation_expires_at_unix_s: u64_from_i64( + "reservation_expires_at_unix_s", + row.try_get("reservation_expires_at_unix_s")?, + )?, + reservation_json: row.try_get("reservation_json")?, + inserted_at_ms: row.try_get("inserted_at_ms")?, + }) +} + +#[cfg_attr(coverage_nightly, coverage(off))] +fn seller_reservation_line_from_row( + row: sqlx::sqlite::SqliteRow, +) -> Result<RadrootsStoredSellerReservationLine, RadrootsEventStoreError> { + Ok(RadrootsStoredSellerReservationLine { + reservation_id: parse_id(row.try_get::<String, _>("reservation_id")?)?, + line_id: parse_id(row.try_get::<String, _>("line_id")?)?, + bin_id: parse_id(row.try_get::<String, _>("bin_id")?)?, + quantity_mantissa: row.try_get("quantity_mantissa")?, + quantity_scale: u8_from_i64("quantity_scale", row.try_get("quantity_scale")?)?, + unit_code: row.try_get("unit_code")?, + line_index: u32_from_i64("line_index", row.try_get("line_index")?)?, + }) +} + +#[cfg_attr(coverage_nightly, coverage(off))] +fn trade_projection_checkpoint_from_row( + row: sqlx::sqlite::SqliteRow, +) -> Result<RadrootsTradeProjectionCheckpoint, RadrootsEventStoreError> { + Ok(RadrootsTradeProjectionCheckpoint { + trade_id: parse_id(row.try_get::<String, _>("trade_id")?)?, + reducer_contract_id: row.try_get("reducer_contract_id")?, + reducer_version: u16_from_i64("reducer_version", row.try_get("reducer_version")?)?, + projection_digest: row.try_get("projection_digest")?, + root_mutation_id: parse_optional_id(row.try_get("root_mutation_id")?)?, + negotiation_state: row.try_get("negotiation_state")?, + agreement_state: row.try_get("agreement_state")?, + evidence_state: row.try_get("evidence_state")?, + conflict_state: row.try_get("conflict_state")?, + private_terms_state: row.try_get("private_terms_state")?, + attestation_state: row.try_get("attestation_state")?, + fulfillment_state: row.try_get("fulfillment_state")?, + payment_state: row.try_get("payment_state")?, + projection_json: row.try_get("projection_json")?, + last_mutation_id: parse_optional_id(row.try_get("last_mutation_id")?)?, + last_transport_event_seq: row.try_get("last_transport_event_seq")?, + updated_at_ms: row.try_get("updated_at_ms")?, + }) +} + +#[cfg_attr(coverage_nightly, coverage(off))] fn transport_observation_from_row( row: sqlx::sqlite::SqliteRow, ) -> Result<RadrootsTransportObservationRow, RadrootsEventStoreError> { @@ -884,6 +1553,16 @@ fn u32_from_i64(field: &'static str, value: i64) -> Result<u32, RadrootsEventSto } #[cfg_attr(coverage_nightly, coverage(off))] +fn u16_from_i64(field: &'static str, value: i64) -> Result<u16, RadrootsEventStoreError> { + u16::try_from(value).map_err(|_| RadrootsEventStoreError::IntegerRange { field, value }) +} + +#[cfg_attr(coverage_nightly, coverage(off))] +fn u8_from_i64(field: &'static str, value: i64) -> Result<u8, RadrootsEventStoreError> { + u8::try_from(value).map_err(|_| RadrootsEventStoreError::IntegerRange { field, value }) +} + +#[cfg_attr(coverage_nightly, coverage(off))] fn u64_from_i64(field: &'static str, value: i64) -> Result<u64, RadrootsEventStoreError> { u64::try_from(value).map_err(|_| RadrootsEventStoreError::IntegerRange { field, value }) } @@ -897,6 +1576,20 @@ fn bool_i64(value: bool) -> i64 { if value { 1 } else { 0 } } +fn parse_id<T>(value: String) -> Result<T, RadrootsEventStoreError> +where + T: TryFrom<String, Error = radroots_event::ids::RadrootsIdParseError>, +{ + T::try_from(value).map_err(Into::into) +} + +fn parse_optional_id<T>(value: Option<String>) -> Result<Option<T>, RadrootsEventStoreError> +where + T: TryFrom<String, Error = radroots_event::ids::RadrootsIdParseError>, +{ + value.map(parse_id).transpose() +} + fn validate_tag_query(tag_name: &str, limit: u32) -> Result<(), RadrootsEventStoreError> { if tag_name.is_empty() { return Err(RadrootsEventStoreError::EmptyTagName); @@ -931,13 +1624,33 @@ where validate_tag_query(tag_name, limit) } +fn validate_trade_query_limit(limit: u32) -> Result<(), RadrootsEventStoreError> { + if !(1..=RADROOTS_EVENT_STORE_QUERY_LIMIT_MAX).contains(&limit) { + return Err(RadrootsEventStoreError::QueryLimitOutOfRange { + min: 1, + max: RADROOTS_EVENT_STORE_QUERY_LIMIT_MAX, + actual: limit, + }); + } + Ok(()) +} + #[cfg(test)] mod tests { use super::*; use radroots_event::draft::RadrootsSignedEvent; use radroots_event::event_head::event_head_candidate_for_event; - use radroots_event::kinds::{ - KIND_GEOCHAT, KIND_LISTING, KIND_ORDER_REQUEST, KIND_POST, KIND_PROFILE, + use radroots_event::ids::{RadrootsAddressableCoordinate, RadrootsInventoryBinId}; + use radroots_event::kinds::{KIND_GEOCHAT, KIND_LISTING, KIND_POST, KIND_PROFILE}; + use radroots_event::trade::{ + RADROOTS_TRADE_DECISION_CONTRACT_ID, RADROOTS_TRADE_PROPOSAL_CONTRACT_ID, + RADROOTS_TRADE_SCHEMA_VERSION, RadrootsFulfillmentProfileV1, + RadrootsSellerReservationAssertionV1, RadrootsSellerReservationLineV1, + RadrootsTradeCancellationProfileV1, RadrootsTradeCandidateLineV1, + RadrootsTradeCandidateTermsV1, RadrootsTradeCanonicalMutationV1, RadrootsTradeDecisionV1, + RadrootsTradeEconomicAdjustmentV1, RadrootsTradeEconomicsProfileV1, + RadrootsTradeMutationBodyV1, RadrootsTradeMutationEnvelopeV1, + canonical_trade_mutation_content, }; use radroots_event::wire::{RadrootsNip01EventWire, compute_canonical_nip01_event_id}; use radroots_nostr::prelude::{ @@ -960,6 +1673,172 @@ mod tests { core::iter::repeat_n(character, 64).collect() } + fn trade_id() -> RadrootsTradeId { + RadrootsTradeId::parse("1".repeat(32)).expect("trade id") + } + + fn public_key(character: char) -> RadrootsPublicKey { + if character == 'a' { + return RadrootsPublicKey::parse(FIXTURE_ALICE_PUBLIC_KEY_HEX).expect("alice pubkey"); + } + RadrootsPublicKey::parse(event_id(character)).expect("pubkey") + } + + fn candidate_terms() -> RadrootsTradeCandidateTermsV1 { + RadrootsTradeCandidateTermsV1 { + candidate_id: None, + schema_version: RADROOTS_TRADE_SCHEMA_VERSION, + base_candidate_id: None, + supersession_intent: None, + buyer_pubkey: public_key('a'), + seller_pubkey: public_key('a'), + farm_id: RadrootsDTag::parse("farm-1").expect("farm id"), + lines: vec![RadrootsTradeCandidateLineV1 { + line_id: RadrootsDTag::parse("line-1").expect("line id"), + listing_addr: RadrootsAddressableCoordinate::parse(format!( + "{KIND_LISTING}:{}:listing-1", + FIXTURE_ALICE_PUBLIC_KEY_HEX + )) + .expect("listing address"), + listing_event_id: RadrootsEventId::parse(event_id('c')).expect("listing event id"), + listing_snapshot_sha256: event_id('d'), + product_id: "carrots".to_owned(), + option_id: None, + bin_id: RadrootsInventoryBinId::parse("bin-1").expect("bin id"), + quantity_mantissa: "2".to_owned(), + quantity_scale: 0, + unit_code: "count".to_owned(), + unit_profile: "mvp-count".to_owned(), + unit_price_mantissa: "300".to_owned(), + currency_code: "USD".to_owned(), + line_subtotal_mantissa: "600".to_owned(), + replaces_line_id: None, + }], + line_tombstones: Vec::new(), + economics: RadrootsTradeEconomicsProfileV1 { + profile_id: "mvp-no-payment".to_owned(), + currency_code: "USD".to_owned(), + currency_exponent: 2, + rounding_profile: "half-up".to_owned(), + subtotal_mantissa: "600".to_owned(), + discount_total_mantissa: "0".to_owned(), + adjustment_total_mantissa: "0".to_owned(), + total_mantissa: "600".to_owned(), + adjustments: Vec::<RadrootsTradeEconomicAdjustmentV1>::new(), + }, + fulfillment: RadrootsFulfillmentProfileV1 { + profile_id: "market-pickup".to_owned(), + method: "pickup".to_owned(), + starts_at_unix_s: 1_800_000_000, + ends_at_unix_s: 1_800_003_600, + timezone: "America/New_York".to_owned(), + utc_offset_seconds: -18_000, + fold: 0, + location_class: "private_after_agreement".to_owned(), + requires_private_terms: false, + }, + cancellation: RadrootsTradeCancellationProfileV1 { + profile_id: "mvp".to_owned(), + buyer_pre_agreement: true, + post_agreement_cutoff_unix_s: Some(1_799_999_000), + }, + private_terms: None, + proposal_expires_at_unix_s: 1_799_999_000, + } + } + + fn proposal_envelope() -> RadrootsTradeMutationEnvelopeV1 { + RadrootsTradeMutationEnvelopeV1 { + mutation_id: None, + contract_id: RADROOTS_TRADE_PROPOSAL_CONTRACT_ID.to_owned(), + schema_version: RADROOTS_TRADE_SCHEMA_VERSION, + trade_id: trade_id(), + root_mutation_id: None, + buyer_pubkey: public_key('a'), + seller_pubkey: public_key('a'), + farm_id: RadrootsDTag::parse("farm-1").expect("farm id"), + parent_mutation_ids: Vec::new(), + author_pubkey: public_key('a'), + counterparty_pubkey: public_key('a'), + authored_at_unix_s: 1_799_000_000, + body: RadrootsTradeMutationBodyV1::Proposal { + candidate: candidate_terms(), + }, + } + } + + fn decision_envelope( + proposal: &RadrootsTradeCanonicalMutationV1, + ) -> RadrootsTradeMutationEnvelopeV1 { + let RadrootsTradeMutationBodyV1::Proposal { candidate } = &proposal.envelope.body else { + panic!("proposal"); + }; + let candidate_id = candidate.candidate_id.clone().expect("candidate id"); + let line = candidate.lines.first().expect("candidate line"); + RadrootsTradeMutationEnvelopeV1 { + mutation_id: None, + contract_id: RADROOTS_TRADE_DECISION_CONTRACT_ID.to_owned(), + schema_version: RADROOTS_TRADE_SCHEMA_VERSION, + trade_id: proposal.envelope.trade_id.clone(), + root_mutation_id: Some(proposal.mutation_id.clone()), + buyer_pubkey: public_key('a'), + seller_pubkey: public_key('a'), + farm_id: RadrootsDTag::parse("farm-1").expect("farm id"), + parent_mutation_ids: vec![proposal.mutation_id.clone()], + author_pubkey: public_key('a'), + counterparty_pubkey: public_key('a'), + authored_at_unix_s: 1_799_000_060, + body: RadrootsTradeMutationBodyV1::Decision { + proposal_mutation_id: proposal.mutation_id.clone(), + candidate_id: candidate_id.clone(), + decision: RadrootsTradeDecisionV1::Accepted { + reservation_assertion: Some(RadrootsSellerReservationAssertionV1 { + reservation_id: RadrootsDTag::parse("reservation-1") + .expect("reservation id"), + inventory_authority_id: public_key('a'), + inventory_epoch: 7, + candidate_id, + commitments: vec![RadrootsSellerReservationLineV1 { + line_id: line.line_id.clone(), + bin_id: line.bin_id.clone(), + quantity_mantissa: line.quantity_mantissa.clone(), + quantity_scale: line.quantity_scale, + unit_code: line.unit_code.clone(), + }], + reservation_expires_at_unix_s: 1_799_000_600, + assertion_commitment: event_id('e'), + }), + }, + }, + } + } + + fn signed_trade_mutation(canonical: &RadrootsTradeCanonicalMutationV1) -> RadrootsSignedEvent { + let counterparty = canonical.envelope.counterparty_pubkey.as_str().to_owned(); + let mut tags = vec![ + vec![ + "contract".to_owned(), + canonical.envelope.contract_id.clone(), + ], + vec!["d".to_owned(), canonical.mutation_id.to_string()], + vec!["p".to_owned(), counterparty], + ]; + for parent in &canonical.envelope.parent_mutation_ids { + tags.push(vec!["e".to_owned(), parent.to_string()]); + } + let draft = radroots_event::draft::RadrootsEventDraft::new( + canonical.envelope.contract_id.clone(), + canonical.envelope.mutation_kind().nostr_kind(), + canonical.envelope.authored_at_unix_s, + tags, + canonical.content.clone(), + FIXTURE_ALICE_PUBLIC_KEY_HEX, + ) + .expect("trade draft"); + radroots_nostr::prelude::radroots_nostr_sign_frozen_draft(&fixture_keys(), &draft) + .expect("signed trade mutation") + } + fn signed_event( kind: u32, created_at: u32, @@ -1044,6 +1923,18 @@ mod tests { } } + async fn explain_query_plan(store: &RadrootsEventStore, sql: &str, bind: &str) -> String { + let rows = sqlx::query(sqlx::AssertSqlSafe(sql.to_owned())) + .bind(bind) + .fetch_all(store.pool()) + .await + .expect("query plan"); + rows.into_iter() + .map(|row| row.try_get::<String, _>("detail").expect("detail")) + .collect::<Vec<_>>() + .join("\n") + } + #[test] fn verification_status_values_round_trip() { for status in [ @@ -1170,14 +2061,26 @@ mod tests { .collect::<Vec<_>>(); assert!(names.iter().any(|name| name == "listing_projection")); - assert!(names.iter().any(|name| name == "trade_projection")); + for table in [ + "trade_mutation", + "trade_mutation_parent", + "trade_missing_parent", + "trade_transport_envelope", + "seller_inventory_reservation", + "seller_inventory_reservation_line", + "trade_projection_checkpoint", + "trade_projection_quarantine", + ] { + assert!(names.iter().any(|name| name == table), "{table}"); + } + assert!(!names.iter().any(|name| name == "trade_projection")); assert!(names.iter().any(|name| name == "listing_search_fts")); } #[tokio::test] - async fn migration_installs_root_aware_trade_projection_key() { + async fn migration_installs_semantic_trade_mutation_keys() { let store = RadrootsEventStore::open_memory().await.expect("open"); - let rows = sqlx::query("PRAGMA table_info(trade_projection)") + let rows = sqlx::query("PRAGMA table_info(trade_mutation)") .fetch_all(store.pool()) .await .expect("table info"); @@ -1197,18 +2100,16 @@ mod tests { .collect::<Vec<_>>(); primary_key.sort_by_key(|(_, pk)| *pk); - assert_eq!( - primary_key, - vec![ - ("order_id", 1), - ("root_event_id", 2), - ("projection_version", 3) - ] + assert_eq!(primary_key, vec![("mutation_id", 1)]); + assert!( + columns + .iter() + .any(|(name, notnull, _)| name == "canonical_payload_bytes" && *notnull == 1) ); assert!( columns .iter() - .any(|(name, notnull, _)| name == "evidence_hash" && *notnull == 1) + .any(|(name, notnull, _)| name == "payload_sha256" && *notnull == 1) ); } @@ -1277,6 +2178,92 @@ mod tests { } #[tokio::test] + async fn trade_mutation_ingest_stores_semantic_rows_missing_parents_and_reservations() { + let store = RadrootsEventStore::open_memory().await.expect("open"); + let proposal = canonical_trade_mutation_content(proposal_envelope()).expect("proposal"); + let decision = + canonical_trade_mutation_content(decision_envelope(&proposal)).expect("decision"); + let decision_event = signed_trade_mutation(&decision); + + let decision_receipt = store + .ingest_event(RadrootsEventIngest::new(decision_event.clone(), 2_000)) + .await + .expect("decision ingest"); + assert!(decision_receipt.projection_eligible); + + let stored_decision = store + .get_trade_mutation(&decision.mutation_id) + .await + .expect("decision query") + .expect("decision mutation"); + assert_eq!(stored_decision.mutation_id, decision.mutation_id); + assert_eq!(stored_decision.trade_id, trade_id()); + assert_eq!( + stored_decision.canonical_payload_bytes, + decision.content.as_bytes() + ); + assert_eq!( + stored_decision.payload_sha256, + sha256_hex(decision.content.as_bytes()) + ); + assert_eq!( + stored_decision.first_transport_event_id.as_str(), + decision_event.id_str() + ); + assert_eq!( + stored_decision.mutation_kind, + RadrootsTradeMutationKindV1::Decision + ); + + let missing = store + .missing_trade_parents(&trade_id()) + .await + .expect("missing parents"); + assert_eq!(missing.len(), 1); + assert_eq!(missing[0].mutation_id, decision.mutation_id); + assert_eq!(missing[0].missing_parent_mutation_id, proposal.mutation_id); + + let reservation_id = RadrootsDTag::parse("reservation-1").expect("reservation id"); + let reservation = store + .seller_reservation(&reservation_id) + .await + .expect("reservation query") + .expect("reservation"); + assert_eq!(reservation.claim_mutation_id, decision.mutation_id); + assert_eq!(reservation.trade_id, trade_id()); + assert_eq!(reservation.assertion_commitment, event_id('e')); + let lines = store + .seller_reservation_lines(&reservation_id) + .await + .expect("reservation lines"); + assert_eq!(lines.len(), 1); + assert_eq!(lines[0].bin_id.as_str(), "bin-1"); + + let proposal_event = signed_trade_mutation(&proposal); + store + .ingest_event(RadrootsEventIngest::new(proposal_event, 2_100)) + .await + .expect("proposal ingest"); + let missing = store + .missing_trade_parents(&trade_id()) + .await + .expect("missing parents resolved"); + assert!(missing.is_empty()); + let parents = store + .trade_mutation_parents(&decision.mutation_id) + .await + .expect("parents"); + assert_eq!(parents.len(), 1); + assert_eq!(parents[0].parent_mutation_id, proposal.mutation_id); + let transport_envelopes = store + .trade_transport_envelopes_for_mutation(&decision.mutation_id) + .await + .expect("transport envelopes"); + assert_eq!(transport_envelopes.len(), 1); + assert_eq!(transport_envelopes[0].transport_kind, "nostr"); + } + + #[tokio::test] async fn wrapper_json_is_rejected_as_event_authority() { let event = signed_event(KIND_POST, 10, Vec::new(), "hello"); let wrapper_json = serde_json::to_string(&event).expect("wrapper json"); @@ -1723,7 +2710,7 @@ mod tests { } #[tokio::test] - async fn events_by_contract_and_tag_enforces_contract_tag_and_projection_filters() { + async fn events_by_contract_and_tag_enforces_trade_contract_tag_and_projection_filters() { let store = RadrootsEventStore::open_memory().await.expect("store"); assert!(matches!( @@ -1732,8 +2719,10 @@ mod tests { .await, Err(RadrootsEventStoreError::EmptyContractList) )); - let too_many_contracts = - vec!["radroots.order.request.v1"; RADROOTS_EVENT_STORE_CONTRACT_QUERY_LIMIT_MAX + 1]; + let too_many_contracts = vec![ + RADROOTS_TRADE_PROPOSAL_CONTRACT_ID; + RADROOTS_EVENT_STORE_CONTRACT_QUERY_LIMIT_MAX + 1 + ]; assert!(matches!( store .events_by_contract_and_tag( @@ -1746,24 +2735,8 @@ mod tests { Err(RadrootsEventStoreError::ContractListTooLarge { .. }) )); - let matching_order = signed_event( - KIND_ORDER_REQUEST, - 70, - vec![ - vec!["d".to_owned(), "order-1".to_owned()], - vec!["p".to_owned(), FIXTURE_ALICE_PUBLIC_KEY_HEX.to_owned()], - ], - "{}", - ); - let wrong_tag_order = signed_event( - KIND_ORDER_REQUEST, - 71, - vec![ - vec!["d".to_owned(), "order-2".to_owned()], - vec!["p".to_owned(), event_id('b')], - ], - "{}", - ); + let proposal = canonical_trade_mutation_content(proposal_envelope()).expect("proposal"); + let matching_trade = signed_trade_mutation(&proposal); let same_tag_wrong_contract = signed_event( KIND_POST, 72, @@ -1784,8 +2757,7 @@ mod tests { ); for (event, observed_at_ms) in [ - (matching_order.clone(), 3_600), - (wrong_tag_order, 3_700), + (matching_trade.clone(), 3_600), (same_tag_wrong_contract, 3_800), (unsupported_same_tag, 3_900), ] { @@ -1797,7 +2769,7 @@ mod tests { let events = store .events_by_contract_and_tag( - &["radroots.order.request.v1"], + &[RADROOTS_TRADE_PROPOSAL_CONTRACT_ID], "p", FIXTURE_ALICE_PUBLIC_KEY_HEX, 10, @@ -1805,10 +2777,10 @@ mod tests { .await .expect("contract tag query"); assert_eq!(events.len(), 1); - assert_eq!(events[0].event_id, matching_order.id_str()); + assert_eq!(events[0].event_id, matching_trade.id_str()); assert_eq!( events[0].contract_id.as_deref(), - Some("radroots.order.request.v1") + Some(RADROOTS_TRADE_PROPOSAL_CONTRACT_ID) ); assert!(events[0].projection_eligible); } @@ -1841,51 +2813,44 @@ mod tests { } #[tokio::test] - async fn listing_event_tag_persists_event_id_contract_metadata() { + async fn trade_mutation_tags_persist_contract_and_semantic_metadata() { let store = RadrootsEventStore::open_memory().await.expect("open"); - let listing_event_id = event_id('f'); - let event = signed_event( - KIND_ORDER_REQUEST, - 16, - vec![ - vec!["d".to_owned(), "order-1".to_owned()], - vec!["p".to_owned(), FIXTURE_ALICE_PUBLIC_KEY_HEX.to_owned()], - vec![ - "a".to_owned(), - format!( - "{KIND_LISTING}:{}:AAAAAAAAAAAAAAAAAAAAAg", - FIXTURE_ALICE_PUBLIC_KEY_HEX - ), - ], - vec![ - "listing_event".to_owned(), - listing_event_id.clone(), - "wss://relay.example.com".to_owned(), - ], - ], - "{}", - ); + let proposal = canonical_trade_mutation_content(proposal_envelope()).expect("proposal"); + let event = signed_trade_mutation(&proposal); store .ingest_event(RadrootsEventIngest::new(event.clone(), 3_100)) .await .expect("ingest"); let tags = store.tags_for_event(event.id_str()).await.expect("tags"); - let listing_tag = tags + let contract_tag = tags .iter() - .find(|tag| tag.tag_name == "listing_event") - .expect("listing event tag"); + .find(|tag| tag.tag_name == "contract") + .expect("contract tag"); + let mutation_tag = tags + .iter() + .find(|tag| tag.tag_name == "d") + .expect("mutation d tag"); assert_eq!( - listing_tag.tag_value.as_deref(), - Some(listing_event_id.as_str()) + contract_tag.tag_value.as_deref(), + Some(RADROOTS_TRADE_PROPOSAL_CONTRACT_ID) + ); + assert_eq!(contract_tag.contract_semantic.as_deref(), Some("contract")); + assert_eq!( + contract_tag.contract_value_type.as_deref(), + Some("contract_id") ); + assert!(!contract_tag.relay_indexed); assert_eq!( - listing_tag.contract_semantic.as_deref(), - Some("listing_snapshot") + mutation_tag.tag_value.as_deref(), + Some(proposal.mutation_id.as_str()) ); - assert_eq!(listing_tag.contract_value_type.as_deref(), Some("event_id")); - assert!(!listing_tag.relay_indexed); + assert_eq!( + mutation_tag.contract_semantic.as_deref(), + Some("identifier") + ); + assert_eq!(mutation_tag.contract_value_type.as_deref(), Some("d_tag")); } #[tokio::test] @@ -2090,6 +3055,70 @@ mod tests { } #[tokio::test] + async fn trade_projection_checkpoint_and_list_queries_use_semantic_indexes() { + let store = RadrootsEventStore::open_memory().await.expect("open"); + let proposal = canonical_trade_mutation_content(proposal_envelope()).expect("proposal"); + let proposal_event = signed_trade_mutation(&proposal); + let proposal_receipt = store + .ingest_event(RadrootsEventIngest::new(proposal_event, 7_000)) + .await + .expect("proposal ingest"); + store + .update_trade_projection_checkpoint(&RadrootsTradeProjectionCheckpoint { + trade_id: trade_id(), + reducer_contract_id: "radroots.trade.reducer.v1".to_owned(), + reducer_version: 1, + projection_digest: event_id('f'), + root_mutation_id: Some(proposal.mutation_id.clone()), + negotiation_state: "open".to_owned(), + agreement_state: "none".to_owned(), + evidence_state: "complete".to_owned(), + conflict_state: "none".to_owned(), + private_terms_state: "not_required".to_owned(), + attestation_state: "none".to_owned(), + fulfillment_state: "not_started".to_owned(), + payment_state: "not_tracked".to_owned(), + projection_json: "{\"trade_id\":\"fixture\"}".to_owned(), + last_mutation_id: Some(proposal.mutation_id.clone()), + last_transport_event_seq: Some(proposal_receipt.seq), + updated_at_ms: 7_100, + }) + .await + .expect("checkpoint"); + let checkpoint = store + .trade_projection_checkpoint(&trade_id()) + .await + .expect("checkpoint query") + .expect("checkpoint"); + assert_eq!( + checkpoint.root_mutation_id, + Some(proposal.mutation_id.clone()) + ); + assert_eq!(checkpoint.agreement_state, "none"); + + let mutation_plan = explain_query_plan( + &store, + "EXPLAIN QUERY PLAN SELECT mutation_id FROM trade_mutation WHERE trade_id = ? ORDER BY authored_at_unix_s, mutation_id LIMIT 10", + trade_id().as_str(), + ) + .await; + assert!( + mutation_plan.contains("trade_mutation_trade_idx"), + "{mutation_plan}" + ); + let checkpoint_plan = explain_query_plan( + &store, + "EXPLAIN QUERY PLAN SELECT trade_id FROM trade_projection_checkpoint WHERE agreement_state = ? ORDER BY updated_at_ms, trade_id LIMIT 10", + "none", + ) + .await; + assert!( + checkpoint_plan.contains("trade_projection_checkpoint_agreement_idx"), + "{checkpoint_plan}" + ); + } + + #[tokio::test] async fn smoke_event_store_ingests_and_replays_ten_thousand_events() { let store = RadrootsEventStore::open_memory().await.expect("open"); for index in 0..10_000u32 { diff --git a/crates/outbox/migrations/0001_outbox.up.sql b/crates/outbox/migrations/0001_outbox.up.sql @@ -2,11 +2,19 @@ CREATE TABLE IF NOT EXISTS outbox_operations ( operation_id INTEGER PRIMARY KEY AUTOINCREMENT, operation_kind TEXT NOT NULL, expected_pubkey TEXT NOT NULL, + semantic_scope TEXT NOT NULL CHECK (semantic_scope IN ('generic_event', 'trade_mutation')), + trade_id TEXT, + mutation_id TEXT, + canonical_payload_sha256 TEXT CHECK (canonical_payload_sha256 IS NULL OR length(canonical_payload_sha256) = 64), idempotency_key TEXT, operation_idempotency_digest TEXT NOT NULL, - status TEXT NOT NULL CHECK (status IN ('queued', 'complete', 'deferred_until_implemented', 'failed_terminal', 'cancelled')), + status TEXT NOT NULL CHECK (status IN ('queued', 'complete', 'failed_terminal', 'cancelled')), created_at_ms INTEGER NOT NULL, - updated_at_ms INTEGER NOT NULL + updated_at_ms INTEGER NOT NULL, + CHECK ( + (semantic_scope = 'generic_event' AND trade_id IS NULL AND mutation_id IS NULL AND canonical_payload_sha256 IS NULL) + OR (semantic_scope = 'trade_mutation' AND trade_id IS NOT NULL AND mutation_id IS NOT NULL AND canonical_payload_sha256 IS NOT NULL) + ) ); CREATE UNIQUE INDEX IF NOT EXISTS outbox_operation_idempotency_idx @@ -16,6 +24,10 @@ WHERE idempotency_key IS NOT NULL; CREATE INDEX IF NOT EXISTS outbox_operation_status_idx ON outbox_operations(status, created_at_ms, operation_id); +CREATE UNIQUE INDEX IF NOT EXISTS outbox_operation_trade_mutation_idx +ON outbox_operations(operation_kind, expected_pubkey, mutation_id) +WHERE semantic_scope = 'trade_mutation'; + CREATE TABLE IF NOT EXISTS outbox_event ( outbox_event_id INTEGER PRIMARY KEY AUTOINCREMENT, operation_id INTEGER NOT NULL REFERENCES outbox_operations(operation_id) ON DELETE CASCADE, @@ -24,7 +36,7 @@ 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', 'deferred_until_implemented', 'failed_terminal', 'cancelled')), + state TEXT NOT NULL CHECK (state IN ('draft_queued', 'signing', 'signed', 'publishing', 'published', 'sign_retryable', 'publish_retryable', 'failed_terminal', 'cancelled')), attempt_count INTEGER NOT NULL, claim_token TEXT, claim_owner TEXT, @@ -54,7 +66,7 @@ CREATE TABLE IF NOT EXISTS outbox_delivery_plan ( satisfaction_policy TEXT NOT NULL, required_success_count INTEGER NOT NULL, delivery_plan_idempotency_digest TEXT NOT NULL, - status TEXT NOT NULL CHECK (status IN ('queued', 'complete', 'deferred_until_implemented', 'failed_terminal', 'cancelled')), + status TEXT NOT NULL CHECK (status IN ('queued', 'complete', 'failed_terminal', 'cancelled')), satisfied_at_ms INTEGER, created_at_ms INTEGER NOT NULL, updated_at_ms INTEGER NOT NULL, diff --git a/crates/outbox/src/error.rs b/crates/outbox/src/error.rs @@ -29,12 +29,23 @@ pub enum RadrootsOutboxError { #[error("transport profile id cannot be empty")] EmptyTransportProfileId, + #[error("trade mutation drafts require the semantic trade mutation outbox API")] + TradeMutationRequiresSemanticOutbox, + + #[error( + "trade mutation outbox metadata does not match the canonical mutation content: {field}" + )] + TradeMutationMetadataMismatch { field: &'static str }, + #[error("transport contract error: {0}")] Transport(RadrootsTransportError), #[error("Invalid stored enum for {field}: {value}")] InvalidStoredEnum { field: &'static str, value: String }, + #[error("Invalid stored identifier for {field}: {value}")] + InvalidStoredIdentifier { field: &'static str, value: String }, + #[error("stored integer for {field} is outside the supported range: {value}")] IntegerRange { field: &'static str, value: i64 }, diff --git a/crates/outbox/src/lib.rs b/crates/outbox/src/lib.rs @@ -17,6 +17,7 @@ pub use model::{ RadrootsOutboxIdempotencyPreflight, RadrootsOutboxOperationInput, RadrootsOutboxOperationRecord, RadrootsOutboxOperationStatus, RadrootsOutboxReticulumBehavior, RadrootsOutboxReticulumEventRecord, RadrootsOutboxSignedOperationInput, - RadrootsOutboxStatusSummary, + RadrootsOutboxSignedTradeMutationInput, RadrootsOutboxStatusSummary, + RadrootsOutboxTradeMutationInput, }; pub use store::RadrootsOutbox; diff --git a/crates/outbox/src/model.rs b/crates/outbox/src/model.rs @@ -2,6 +2,7 @@ use crate::RadrootsOutboxError; use radroots_event::draft::{RadrootsEventDraft, RadrootsSignedEvent}; +use radroots_event::ids::{RadrootsTradeId, RadrootsTradeMutationId}; use radroots_transport::{ RadrootsTransportKind, RadrootsTransportMeshScopeId, RadrootsTransportOutcomeKind, RadrootsTransportSatisfactionClass, RadrootsTransportSatisfactionPolicy, @@ -13,7 +14,6 @@ use radroots_transport::{ pub enum RadrootsOutboxOperationStatus { Queued, Complete, - DeferredUntilImplemented, FailedTerminal, Cancelled, } @@ -23,7 +23,6 @@ impl RadrootsOutboxOperationStatus { match self { Self::Queued => "queued", Self::Complete => "complete", - Self::DeferredUntilImplemented => "deferred_until_implemented", Self::FailedTerminal => "failed_terminal", Self::Cancelled => "cancelled", } @@ -33,7 +32,6 @@ impl RadrootsOutboxOperationStatus { match value { "queued" => Ok(Self::Queued), "complete" => Ok(Self::Complete), - "deferred_until_implemented" => Ok(Self::DeferredUntilImplemented), "failed_terminal" => Ok(Self::FailedTerminal), "cancelled" => Ok(Self::Cancelled), _ => Err(RadrootsOutboxError::InvalidStoredEnum { @@ -53,7 +51,6 @@ pub enum RadrootsOutboxEventState { Published, SignRetryable, PublishRetryable, - DeferredUntilImplemented, FailedTerminal, Cancelled, } @@ -68,7 +65,6 @@ impl RadrootsOutboxEventState { Self::Published => "published", Self::SignRetryable => "sign_retryable", Self::PublishRetryable => "publish_retryable", - Self::DeferredUntilImplemented => "deferred_until_implemented", Self::FailedTerminal => "failed_terminal", Self::Cancelled => "cancelled", } @@ -83,7 +79,6 @@ impl RadrootsOutboxEventState { "published" => Ok(Self::Published), "sign_retryable" => Ok(Self::SignRetryable), "publish_retryable" => Ok(Self::PublishRetryable), - "deferred_until_implemented" => Ok(Self::DeferredUntilImplemented), "failed_terminal" => Ok(Self::FailedTerminal), "cancelled" => Ok(Self::Cancelled), _ => Err(RadrootsOutboxError::InvalidStoredEnum { @@ -105,7 +100,6 @@ impl RadrootsOutboxEventState { pub enum RadrootsOutboxDeliveryPlanStatus { Queued, Complete, - DeferredUntilImplemented, FailedTerminal, Cancelled, } @@ -115,7 +109,6 @@ impl RadrootsOutboxDeliveryPlanStatus { match self { Self::Queued => "queued", Self::Complete => "complete", - Self::DeferredUntilImplemented => "deferred_until_implemented", Self::FailedTerminal => "failed_terminal", Self::Cancelled => "cancelled", } @@ -125,7 +118,6 @@ impl RadrootsOutboxDeliveryPlanStatus { match value { "queued" => Ok(Self::Queued), "complete" => Ok(Self::Complete), - "deferred_until_implemented" => Ok(Self::DeferredUntilImplemented), "failed_terminal" => Ok(Self::FailedTerminal), "cancelled" => Ok(Self::Cancelled), _ => Err(RadrootsOutboxError::InvalidStoredEnum { @@ -353,6 +345,95 @@ impl RadrootsOutboxSignedOperationInput { } } +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct RadrootsOutboxTradeMutationInput { + pub operation_kind: String, + pub trade_id: RadrootsTradeId, + pub mutation_id: RadrootsTradeMutationId, + pub canonical_payload_sha256: String, + pub draft: RadrootsEventDraft, + pub delivery_plan: RadrootsOutboxDeliveryPlanInput, + pub idempotency_key: Option<String>, + pub created_at_ms: i64, +} + +impl RadrootsOutboxTradeMutationInput { + pub fn new( + operation_kind: impl Into<String>, + trade_id: RadrootsTradeId, + mutation_id: RadrootsTradeMutationId, + canonical_payload_sha256: impl Into<String>, + draft: RadrootsEventDraft, + delivery_plan: RadrootsOutboxDeliveryPlanInput, + created_at_ms: i64, + ) -> Self { + Self { + operation_kind: operation_kind.into(), + trade_id, + mutation_id, + canonical_payload_sha256: canonical_payload_sha256.into(), + draft, + delivery_plan, + idempotency_key: None, + created_at_ms, + } + } + + pub fn with_idempotency_key(mut self, idempotency_key: impl Into<String>) -> Self { + self.idempotency_key = Some(idempotency_key.into()); + self + } +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct RadrootsOutboxSignedTradeMutationInput { + pub operation_kind: String, + pub trade_id: RadrootsTradeId, + pub mutation_id: RadrootsTradeMutationId, + pub canonical_payload_sha256: String, + pub draft: RadrootsEventDraft, + pub signed_event: RadrootsSignedEvent, + pub delivery_plan: RadrootsOutboxDeliveryPlanInput, + pub idempotency_key: Option<String>, + pub event_store_inserted: bool, + pub event_store_ingested_at_ms: i64, + pub created_at_ms: i64, +} + +impl RadrootsOutboxSignedTradeMutationInput { + pub fn new( + operation_kind: impl Into<String>, + trade_id: RadrootsTradeId, + mutation_id: RadrootsTradeMutationId, + canonical_payload_sha256: impl Into<String>, + draft: RadrootsEventDraft, + signed_event: RadrootsSignedEvent, + delivery_plan: RadrootsOutboxDeliveryPlanInput, + event_store_inserted: bool, + event_store_ingested_at_ms: i64, + created_at_ms: i64, + ) -> Self { + Self { + operation_kind: operation_kind.into(), + trade_id, + mutation_id, + canonical_payload_sha256: canonical_payload_sha256.into(), + draft, + signed_event, + delivery_plan, + idempotency_key: None, + event_store_inserted, + event_store_ingested_at_ms, + created_at_ms, + } + } + + pub fn with_idempotency_key(mut self, idempotency_key: impl Into<String>) -> Self { + self.idempotency_key = Some(idempotency_key.into()); + self + } +} + #[derive(Clone, Copy, Debug, PartialEq, Eq)] pub enum RadrootsOutboxEnqueueStatus { Inserted, @@ -381,6 +462,10 @@ pub struct RadrootsOutboxOperationRecord { pub operation_id: i64, pub operation_kind: String, pub expected_pubkey: String, + pub semantic_scope: String, + pub trade_id: Option<RadrootsTradeId>, + pub mutation_id: Option<RadrootsTradeMutationId>, + pub canonical_payload_sha256: Option<String>, pub idempotency_key: Option<String>, pub operation_idempotency_digest: String, pub status: RadrootsOutboxOperationStatus, @@ -508,10 +593,6 @@ mod tests { (RadrootsOutboxOperationStatus::Queued, "queued"), (RadrootsOutboxOperationStatus::Complete, "complete"), ( - RadrootsOutboxOperationStatus::DeferredUntilImplemented, - "deferred_until_implemented", - ), - ( RadrootsOutboxOperationStatus::FailedTerminal, "failed_terminal", ), @@ -535,10 +616,6 @@ mod tests { RadrootsOutboxEventState::PublishRetryable, "publish_retryable", ), - ( - RadrootsOutboxEventState::DeferredUntilImplemented, - "deferred_until_implemented", - ), (RadrootsOutboxEventState::FailedTerminal, "failed_terminal"), (RadrootsOutboxEventState::Cancelled, "cancelled"), ] { @@ -553,10 +630,6 @@ mod tests { (RadrootsOutboxDeliveryPlanStatus::Queued, "queued"), (RadrootsOutboxDeliveryPlanStatus::Complete, "complete"), ( - RadrootsOutboxDeliveryPlanStatus::DeferredUntilImplemented, - "deferred_until_implemented", - ), - ( RadrootsOutboxDeliveryPlanStatus::FailedTerminal, "failed_terminal", ), diff --git a/crates/outbox/src/store.rs b/crates/outbox/src/store.rs @@ -11,11 +11,15 @@ use crate::model::{ RadrootsOutboxIdempotencyPreflight, RadrootsOutboxOperationInput, RadrootsOutboxOperationRecord, RadrootsOutboxOperationStatus, RadrootsOutboxReticulumBehavior, RadrootsOutboxReticulumEventRecord, RadrootsOutboxSignedOperationInput, - RadrootsOutboxStatusSummary, + RadrootsOutboxSignedTradeMutationInput, RadrootsOutboxStatusSummary, + RadrootsOutboxTradeMutationInput, }; use radroots_event::draft::{ RadrootsEventDraft, RadrootsSignedEvent, validate_signed_nostr_event_matches_draft, }; +use radroots_event::ids::{RadrootsTradeId, RadrootsTradeMutationId}; +use radroots_event::kinds::TRADE_MUTATION_EVENT_KINDS; +use radroots_event::trade::trade_mutation_from_canonical_content; use radroots_event::wire::RadrootsNip01EventWire; use radroots_event_store::{ RadrootsEventIngest, RadrootsEventStore, RadrootsTransportObservation, @@ -99,7 +103,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 = '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", + "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", ) .fetch_one(&self.pool) .await?; @@ -109,6 +113,12 @@ impl RadrootsOutbox { .bind(now_ms) .bind(now_ms) .fetch_one(&self.pool) + .await? + .try_get(0)?; + let deferred_until_implemented_events = sqlx::query( + "SELECT COUNT(DISTINCT plan.outbox_event_id) FROM outbox_delivery_plan AS plan JOIN outbox_delivery_target AS target ON target.delivery_plan_id = plan.delivery_plan_id WHERE target.status = 'deferred_until_implemented'", + ) + .fetch_one(&self.pool) .await? .try_get(0)?; let last_attempt_at_ms = @@ -129,7 +139,7 @@ impl RadrootsOutbox { retryable_events: row.try_get("retryable_events")?, terminal_events: row.try_get("terminal_events")?, failed_terminal_events: row.try_get("failed_terminal_events")?, - deferred_until_implemented_events: row.try_get("deferred_until_implemented_events")?, + deferred_until_implemented_events, ready_signed_events, publishing_events: row.try_get("publishing_events")?, last_attempt_at_ms, @@ -141,6 +151,7 @@ impl RadrootsOutbox { &self, input: &RadrootsOutboxSignedOperationInput, ) -> Result<RadrootsOutboxIdempotencyPreflight, RadrootsOutboxError> { + ensure_not_trade_mutation_draft(&input.draft)?; validate_signed_nostr_event_matches_draft(&input.signed_event, &input.draft)?; let prepared = prepare_delivery_plan(input.draft.expected_event_id_str(), &input.delivery_plan)?; @@ -175,10 +186,73 @@ impl RadrootsOutbox { }) } + pub async fn preflight_signed_trade_mutation_idempotency( + &self, + input: &RadrootsOutboxSignedTradeMutationInput, + ) -> Result<RadrootsOutboxIdempotencyPreflight, RadrootsOutboxError> { + validate_signed_nostr_event_matches_draft(&input.signed_event, &input.draft)?; + let semantic = validate_trade_mutation_input( + input.trade_id.as_str(), + input.mutation_id.as_str(), + input.canonical_payload_sha256.as_str(), + &input.draft, + )?; + let prepared = + prepare_delivery_plan(input.draft.expected_event_id_str(), &input.delivery_plan)?; + let operation_digest = trade_mutation_operation_idempotency_digest( + input.operation_kind.as_str(), + input.draft.expected_pubkey_str(), + &semantic, + ); + + if let Some(existing) = existing_trade_mutation_operation_for_pool( + &self.pool, + input.operation_kind.as_str(), + input.draft.expected_pubkey_str(), + input.mutation_id.as_str(), + ) + .await? + && existing.operation_idempotency_digest != operation_digest + { + return Err(RadrootsOutboxError::IdempotencyConflict { + operation_kind: input.operation_kind.clone(), + expected_pubkey: input.draft.expected_pubkey_str().to_owned(), + idempotency_key: input.mutation_id.to_string(), + existing_digest: existing.operation_idempotency_digest, + new_digest: operation_digest, + }); + } + + if let Some(idempotency_key) = input.idempotency_key.as_deref() + && let Some(existing) = existing_idempotent_operation_for_pool( + &self.pool, + input.operation_kind.as_str(), + input.draft.expected_pubkey_str(), + idempotency_key, + ) + .await? + && existing.operation_idempotency_digest != operation_digest + { + return Err(RadrootsOutboxError::IdempotencyConflict { + operation_kind: input.operation_kind.clone(), + expected_pubkey: input.draft.expected_pubkey_str().to_owned(), + idempotency_key: idempotency_key.to_owned(), + existing_digest: existing.operation_idempotency_digest, + new_digest: operation_digest, + }); + } + + Ok(RadrootsOutboxIdempotencyPreflight { + operation_idempotency_digest: operation_digest, + delivery_plan_idempotency_digest: prepared.delivery_plan_idempotency_digest, + }) + } + pub async fn enqueue_operation( &self, input: RadrootsOutboxOperationInput, ) -> Result<RadrootsOutboxEnqueueReceipt, RadrootsOutboxError> { + ensure_not_trade_mutation_draft(&input.draft)?; let prepared = prepare_delivery_plan(input.draft.expected_event_id_str(), &input.delivery_plan)?; let operation_digest = operation_idempotency_digest( @@ -230,7 +304,7 @@ impl RadrootsOutbox { } let operation = sqlx::query( - "INSERT INTO outbox_operations(operation_kind, expected_pubkey, idempotency_key, operation_idempotency_digest, status, created_at_ms, updated_at_ms) VALUES (?, ?, ?, ?, ?, ?, ?)", + "INSERT INTO outbox_operations(operation_kind, expected_pubkey, semantic_scope, trade_id, mutation_id, canonical_payload_sha256, idempotency_key, operation_idempotency_digest, status, created_at_ms, updated_at_ms) VALUES (?, ?, 'generic_event', NULL, NULL, NULL, ?, ?, ?, ?, ?)", ) .bind(input.operation_kind.as_str()) .bind(input.draft.expected_pubkey_str()) @@ -276,6 +350,7 @@ impl RadrootsOutbox { &self, input: RadrootsOutboxSignedOperationInput, ) -> Result<RadrootsOutboxEnqueueReceipt, RadrootsOutboxError> { + ensure_not_trade_mutation_draft(&input.draft)?; validate_signed_nostr_event_matches_draft(&input.signed_event, &input.draft)?; let prepared = prepare_delivery_plan(input.draft.expected_event_id_str(), &input.delivery_plan)?; @@ -334,10 +409,321 @@ impl RadrootsOutbox { } let operation = sqlx::query( - "INSERT INTO outbox_operations(operation_kind, expected_pubkey, idempotency_key, operation_idempotency_digest, status, created_at_ms, updated_at_ms) VALUES (?, ?, ?, ?, ?, ?, ?)", + "INSERT INTO outbox_operations(operation_kind, expected_pubkey, semantic_scope, trade_id, mutation_id, canonical_payload_sha256, idempotency_key, operation_idempotency_digest, status, created_at_ms, updated_at_ms) VALUES (?, ?, 'generic_event', NULL, NULL, NULL, ?, ?, ?, ?, ?)", + ) + .bind(input.operation_kind.as_str()) + .bind(input.draft.expected_pubkey_str()) + .bind(input.idempotency_key.as_deref()) + .bind(operation_digest.as_str()) + .bind(RadrootsOutboxOperationStatus::Queued.as_str()) + .bind(input.created_at_ms) + .bind(input.created_at_ms) + .execute(&mut *tx) + .await?; + let operation_id = operation.last_insert_rowid(); + let draft_json = serde_json::to_string(&input.draft)?; + let signed_event_json = signed_event_wire_json(&input.signed_event)?; + let event = sqlx::query( + "INSERT INTO outbox_event(operation_id, event_id, expected_pubkey, draft_json, signed_event_json, raw_event_json, state, attempt_count, next_attempt_after_ms, event_store_ingested, event_store_inserted, event_store_ingested_at_ms, created_at_ms, updated_at_ms) VALUES (?, ?, ?, ?, ?, ?, ?, 0, ?, 1, ?, ?, ?, ?)", + ) + .bind(operation_id) + .bind(input.draft.expected_event_id_str()) + .bind(input.draft.expected_pubkey_str()) + .bind(draft_json.as_str()) + .bind(signed_event_json.as_str()) + .bind(input.signed_event.raw_json()) + .bind(RadrootsOutboxEventState::Signed.as_str()) + .bind(input.created_at_ms) + .bind(bool_i64(input.event_store_inserted)) + .bind(input.event_store_ingested_at_ms) + .bind(input.created_at_ms) + .bind(input.created_at_ms) + .execute(&mut *tx) + .await?; + let outbox_event_id = event.last_insert_rowid(); + 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, + operation_id, + outbox_event_id, + delivery_plan_id: plan.delivery_plan_id, + expected_event_id: input.draft.expected_event_id_str().to_owned(), + operation_idempotency_digest: operation_digest, + delivery_plan_idempotency_digest: prepared.delivery_plan_idempotency_digest, + }) + } + + pub async fn enqueue_trade_mutation_operation( + &self, + input: RadrootsOutboxTradeMutationInput, + ) -> Result<RadrootsOutboxEnqueueReceipt, RadrootsOutboxError> { + let semantic = validate_trade_mutation_input( + input.trade_id.as_str(), + input.mutation_id.as_str(), + input.canonical_payload_sha256.as_str(), + &input.draft, + )?; + let prepared = + prepare_delivery_plan(input.draft.expected_event_id_str(), &input.delivery_plan)?; + let operation_digest = trade_mutation_operation_idempotency_digest( + input.operation_kind.as_str(), + input.draft.expected_pubkey_str(), + &semantic, + ); + let mut tx = self.pool.begin().await?; + + if let Some(existing) = existing_trade_mutation_operation( + &mut tx, + input.operation_kind.as_str(), + input.draft.expected_pubkey_str(), + input.mutation_id.as_str(), + ) + .await? + { + if existing.operation_idempotency_digest != operation_digest { + return Err(RadrootsOutboxError::IdempotencyConflict { + operation_kind: input.operation_kind, + expected_pubkey: input.draft.expected_pubkey_str().to_owned(), + idempotency_key: input.mutation_id.to_string(), + existing_digest: existing.operation_idempotency_digest, + new_digest: operation_digest, + }); + } + let plan = insert_or_get_delivery_plan( + &mut tx, + existing.outbox_event_id, + &prepared, + input.created_at_ms, + ) + .await?; + if plan.status == RadrootsOutboxEnqueueStatus::Inserted { + sync_signed_event_lifecycle(&mut tx, existing.outbox_event_id, input.created_at_ms) + .await?; + } + tx.commit().await?; + return Ok(RadrootsOutboxEnqueueReceipt { + status: plan.status, + operation_id: existing.operation_id, + outbox_event_id: existing.outbox_event_id, + delivery_plan_id: plan.delivery_plan_id, + expected_event_id: existing.event_id, + operation_idempotency_digest: operation_digest, + delivery_plan_idempotency_digest: prepared.delivery_plan_idempotency_digest, + }); + } + + if let Some(idempotency_key) = input.idempotency_key.as_deref() + && let Some(existing) = existing_idempotent_operation( + &mut tx, + input.operation_kind.as_str(), + input.draft.expected_pubkey_str(), + idempotency_key, + ) + .await? + { + if existing.operation_idempotency_digest != operation_digest { + return Err(RadrootsOutboxError::IdempotencyConflict { + operation_kind: input.operation_kind, + expected_pubkey: input.draft.expected_pubkey_str().to_owned(), + idempotency_key: idempotency_key.to_owned(), + existing_digest: existing.operation_idempotency_digest, + new_digest: operation_digest, + }); + } + let plan = insert_or_get_delivery_plan( + &mut tx, + existing.outbox_event_id, + &prepared, + input.created_at_ms, + ) + .await?; + if plan.status == RadrootsOutboxEnqueueStatus::Inserted { + sync_signed_event_lifecycle(&mut tx, existing.outbox_event_id, input.created_at_ms) + .await?; + } + tx.commit().await?; + return Ok(RadrootsOutboxEnqueueReceipt { + status: plan.status, + operation_id: existing.operation_id, + outbox_event_id: existing.outbox_event_id, + delivery_plan_id: plan.delivery_plan_id, + expected_event_id: existing.event_id, + operation_idempotency_digest: operation_digest, + delivery_plan_idempotency_digest: prepared.delivery_plan_idempotency_digest, + }); + } + + let operation = sqlx::query( + "INSERT INTO outbox_operations(operation_kind, expected_pubkey, semantic_scope, trade_id, mutation_id, canonical_payload_sha256, idempotency_key, operation_idempotency_digest, status, created_at_ms, updated_at_ms) VALUES (?, ?, 'trade_mutation', ?, ?, ?, ?, ?, ?, ?, ?)", + ) + .bind(input.operation_kind.as_str()) + .bind(input.draft.expected_pubkey_str()) + .bind(input.trade_id.as_str()) + .bind(input.mutation_id.as_str()) + .bind(input.canonical_payload_sha256.as_str()) + .bind(input.idempotency_key.as_deref()) + .bind(operation_digest.as_str()) + .bind(RadrootsOutboxOperationStatus::Queued.as_str()) + .bind(input.created_at_ms) + .bind(input.created_at_ms) + .execute(&mut *tx) + .await?; + let operation_id = operation.last_insert_rowid(); + let draft_json = serde_json::to_string(&input.draft)?; + let event = sqlx::query( + "INSERT INTO outbox_event(operation_id, event_id, expected_pubkey, draft_json, state, attempt_count, next_attempt_after_ms, event_store_ingested, event_store_inserted, created_at_ms, updated_at_ms) VALUES (?, ?, ?, ?, ?, 0, ?, 0, 0, ?, ?)", + ) + .bind(operation_id) + .bind(input.draft.expected_event_id_str()) + .bind(input.draft.expected_pubkey_str()) + .bind(draft_json.as_str()) + .bind(RadrootsOutboxEventState::DraftQueued.as_str()) + .bind(input.created_at_ms) + .bind(input.created_at_ms) + .bind(input.created_at_ms) + .execute(&mut *tx) + .await?; + let outbox_event_id = event.last_insert_rowid(); + let plan = + insert_or_get_delivery_plan(&mut tx, outbox_event_id, &prepared, input.created_at_ms) + .await?; + tx.commit().await?; + Ok(RadrootsOutboxEnqueueReceipt { + status: RadrootsOutboxEnqueueStatus::Inserted, + operation_id, + outbox_event_id, + delivery_plan_id: plan.delivery_plan_id, + expected_event_id: input.draft.expected_event_id_str().to_owned(), + operation_idempotency_digest: operation_digest, + delivery_plan_idempotency_digest: prepared.delivery_plan_idempotency_digest, + }) + } + + pub async fn enqueue_signed_trade_mutation_operation( + &self, + input: RadrootsOutboxSignedTradeMutationInput, + ) -> Result<RadrootsOutboxEnqueueReceipt, RadrootsOutboxError> { + validate_signed_nostr_event_matches_draft(&input.signed_event, &input.draft)?; + let semantic = validate_trade_mutation_input( + input.trade_id.as_str(), + input.mutation_id.as_str(), + input.canonical_payload_sha256.as_str(), + &input.draft, + )?; + let prepared = + prepare_delivery_plan(input.draft.expected_event_id_str(), &input.delivery_plan)?; + let operation_digest = trade_mutation_operation_idempotency_digest( + input.operation_kind.as_str(), + input.draft.expected_pubkey_str(), + &semantic, + ); + let mut tx = self.pool.begin().await?; + + if let Some(existing) = existing_trade_mutation_operation( + &mut tx, + input.operation_kind.as_str(), + input.draft.expected_pubkey_str(), + input.mutation_id.as_str(), + ) + .await? + { + if existing.operation_idempotency_digest != operation_digest { + return Err(RadrootsOutboxError::IdempotencyConflict { + operation_kind: input.operation_kind, + expected_pubkey: input.draft.expected_pubkey_str().to_owned(), + idempotency_key: input.mutation_id.to_string(), + existing_digest: existing.operation_idempotency_digest, + new_digest: operation_digest, + }); + } + ensure_event_signed( + &mut tx, + existing.outbox_event_id, + &input.signed_event, + input.event_store_inserted, + input.event_store_ingested_at_ms, + ) + .await?; + let plan = insert_or_get_delivery_plan( + &mut tx, + existing.outbox_event_id, + &prepared, + 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 { + status: plan.status, + operation_id: existing.operation_id, + outbox_event_id: existing.outbox_event_id, + delivery_plan_id: plan.delivery_plan_id, + expected_event_id: existing.event_id, + operation_idempotency_digest: operation_digest, + delivery_plan_idempotency_digest: prepared.delivery_plan_idempotency_digest, + }); + } + + if let Some(idempotency_key) = input.idempotency_key.as_deref() + && let Some(existing) = existing_idempotent_operation( + &mut tx, + input.operation_kind.as_str(), + input.draft.expected_pubkey_str(), + idempotency_key, + ) + .await? + { + if existing.operation_idempotency_digest != operation_digest { + return Err(RadrootsOutboxError::IdempotencyConflict { + operation_kind: input.operation_kind, + expected_pubkey: input.draft.expected_pubkey_str().to_owned(), + idempotency_key: idempotency_key.to_owned(), + existing_digest: existing.operation_idempotency_digest, + new_digest: operation_digest, + }); + } + ensure_event_signed( + &mut tx, + existing.outbox_event_id, + &input.signed_event, + input.event_store_inserted, + input.event_store_ingested_at_ms, + ) + .await?; + let plan = insert_or_get_delivery_plan( + &mut tx, + existing.outbox_event_id, + &prepared, + 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 { + status: plan.status, + operation_id: existing.operation_id, + outbox_event_id: existing.outbox_event_id, + delivery_plan_id: plan.delivery_plan_id, + expected_event_id: existing.event_id, + operation_idempotency_digest: operation_digest, + delivery_plan_idempotency_digest: prepared.delivery_plan_idempotency_digest, + }); + } + + let operation = sqlx::query( + "INSERT INTO outbox_operations(operation_kind, expected_pubkey, semantic_scope, trade_id, mutation_id, canonical_payload_sha256, idempotency_key, operation_idempotency_digest, status, created_at_ms, updated_at_ms) VALUES (?, ?, 'trade_mutation', ?, ?, ?, ?, ?, ?, ?, ?)", ) .bind(input.operation_kind.as_str()) .bind(input.draft.expected_pubkey_str()) + .bind(input.trade_id.as_str()) + .bind(input.mutation_id.as_str()) + .bind(input.canonical_payload_sha256.as_str()) .bind(input.idempotency_key.as_deref()) .bind(operation_digest.as_str()) .bind(RadrootsOutboxOperationStatus::Queued.as_str()) @@ -387,6 +773,7 @@ impl RadrootsOutbox { tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>, input: RadrootsOutboxSignedOperationInput, ) -> Result<RadrootsOutboxEnqueueReceipt, RadrootsOutboxError> { + ensure_not_trade_mutation_draft(&input.draft)?; validate_signed_nostr_event_matches_draft(&input.signed_event, &input.draft)?; let prepared = prepare_delivery_plan(input.draft.expected_event_id_str(), &input.delivery_plan)?; @@ -442,7 +829,7 @@ impl RadrootsOutbox { } let operation = sqlx::query( - "INSERT INTO outbox_operations(operation_kind, expected_pubkey, idempotency_key, operation_idempotency_digest, status, created_at_ms, updated_at_ms) VALUES (?, ?, ?, ?, ?, ?, ?)", + "INSERT INTO outbox_operations(operation_kind, expected_pubkey, semantic_scope, trade_id, mutation_id, canonical_payload_sha256, idempotency_key, operation_idempotency_digest, status, created_at_ms, updated_at_ms) VALUES (?, ?, 'generic_event', NULL, NULL, NULL, ?, ?, ?, ?, ?)", ) .bind(input.operation_kind.as_str()) .bind(input.draft.expected_pubkey_str()) @@ -493,7 +880,7 @@ impl RadrootsOutbox { operation_id: i64, ) -> Result<Option<RadrootsOutboxOperationRecord>, RadrootsOutboxError> { let row = sqlx::query( - "SELECT operation_id, operation_kind, expected_pubkey, idempotency_key, operation_idempotency_digest, status, created_at_ms, updated_at_ms FROM outbox_operations WHERE operation_id = ?", + "SELECT operation_id, operation_kind, expected_pubkey, semantic_scope, trade_id, mutation_id, canonical_payload_sha256, idempotency_key, operation_idempotency_digest, status, created_at_ms, updated_at_ms FROM outbox_operations WHERE operation_id = ?", ) .bind(operation_id) .fetch_optional(&self.pool) @@ -813,7 +1200,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, 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', 'deferred_until_implemented')", + "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')", ) .bind(now_ms) .bind(now_ms) @@ -1378,7 +1765,6 @@ struct PlanEvaluation { all_complete: bool, any_failed_terminal: bool, any_ready: bool, - any_deferred_until_implemented: bool, } fn publish_lifecycle_from_plan_evaluation<'a>( @@ -1414,13 +1800,6 @@ fn publish_lifecycle_from_plan_evaluation<'a>( Some(terminal_error), now_ms, ) - } else if evaluation.any_deferred_until_implemented { - ( - RadrootsOutboxEventState::DeferredUntilImplemented, - Some(RadrootsOutboxOperationStatus::DeferredUntilImplemented), - None, - now_ms, - ) } else { ( RadrootsOutboxEventState::FailedTerminal, @@ -1631,7 +2010,7 @@ fn initial_delivery_plan_status( if prepared_targets.iter().all(|target| { target.initial_status == RadrootsOutboxDeliveryTargetStatus::DeferredUntilImplemented }) { - return RadrootsOutboxDeliveryPlanStatus::DeferredUntilImplemented; + return RadrootsOutboxDeliveryPlanStatus::FailedTerminal; } RadrootsOutboxDeliveryPlanStatus::Queued } @@ -1653,6 +2032,40 @@ async fn existing_idempotent_operation( row.map(existing_operation_from_row).transpose() } +async fn existing_trade_mutation_operation( + tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>, + operation_kind: &str, + expected_pubkey: &str, + mutation_id: &str, +) -> Result<Option<ExistingOperation>, RadrootsOutboxError> { + let row = sqlx::query( + "SELECT o.operation_id, o.operation_idempotency_digest, e.outbox_event_id, e.event_id FROM outbox_operations o JOIN outbox_event e ON e.operation_id = o.operation_id WHERE o.operation_kind = ? AND o.expected_pubkey = ? AND o.semantic_scope = 'trade_mutation' AND o.mutation_id = ? ORDER BY e.outbox_event_id LIMIT 1", + ) + .bind(operation_kind) + .bind(expected_pubkey) + .bind(mutation_id) + .fetch_optional(&mut **tx) + .await?; + row.map(existing_operation_from_row).transpose() +} + +async fn existing_trade_mutation_operation_for_pool( + pool: &SqlitePool, + operation_kind: &str, + expected_pubkey: &str, + mutation_id: &str, +) -> Result<Option<ExistingOperation>, RadrootsOutboxError> { + let row = sqlx::query( + "SELECT o.operation_id, o.operation_idempotency_digest, e.outbox_event_id, e.event_id FROM outbox_operations o JOIN outbox_event e ON e.operation_id = o.operation_id WHERE o.operation_kind = ? AND o.expected_pubkey = ? AND o.semantic_scope = 'trade_mutation' AND o.mutation_id = ? ORDER BY e.outbox_event_id LIMIT 1", + ) + .bind(operation_kind) + .bind(expected_pubkey) + .bind(mutation_id) + .fetch_optional(pool) + .await?; + row.map(existing_operation_from_row).transpose() +} + async fn existing_idempotent_operation_for_pool( pool: &SqlitePool, operation_kind: &str, @@ -1829,15 +2242,6 @@ async fn signed_event_lifecycle_for_plans( } 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(( @@ -2058,7 +2462,7 @@ async fn reticulum_event_ids_pool( value: i64::MAX, })?; let mut query = String::from( - "SELECT event.outbox_event_id FROM outbox_event AS event WHERE event.signed_event_json IS NOT NULL 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 plan.status IN ('queued', 'deferred_until_implemented') AND target.transport_kind = 'reticulum' AND target.status IN ('pending', 'failed_retryable', 'deferred_until_implemented'))", + "SELECT event.outbox_event_id FROM outbox_event AS event WHERE event.signed_event_json IS NOT NULL 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 plan.status IN ('queued', 'failed_terminal') AND target.transport_kind = 'reticulum' AND target.status IN ('pending', 'failed_retryable', 'deferred_until_implemented'))", ); if outbox_event_id.is_some() { query.push_str(" AND event.outbox_event_id = ?"); @@ -2080,7 +2484,7 @@ async fn reticulum_targets_for_event_pool( outbox_event_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.target_scope, target.target_label, target.endpoint_fingerprint, target.status, target.last_outcome_kind, 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 plan.status IN ('queued', 'deferred_until_implemented') AND target.transport_kind = 'reticulum' AND target.status IN ('pending', 'failed_retryable', 'deferred_until_implemented') ORDER BY target.delivery_plan_id, target.delivery_target_id", + "SELECT target.delivery_target_id, target.delivery_plan_id, target.transport_kind, target.endpoint_uri, target.target_scope, target.target_label, target.endpoint_fingerprint, target.status, target.last_outcome_kind, 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 plan.status IN ('queued', 'failed_terminal') AND target.transport_kind = 'reticulum' AND target.status IN ('pending', 'failed_retryable', 'deferred_until_implemented') ORDER BY target.delivery_plan_id, target.delivery_target_id", ) .bind(outbox_event_id) .fetch_all(pool) @@ -2162,7 +2566,6 @@ 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_deferred_until_implemented = false; for plan in plans { let targets = delivery_targets_for_plan_tx(tx, plan.delivery_plan_id).await?; let satisfied_count = outbox_satisfied_target_count(&plan.satisfaction_policy, &targets); @@ -2183,10 +2586,8 @@ async fn evaluate_delivery_plans( RadrootsOutboxDeliveryPlanStatus::Complete } else if ready_count > 0 { RadrootsOutboxDeliveryPlanStatus::Queued - } else if terminal_failure_count > 0 { + } else if terminal_failure_count > 0 || deferred_count > 0 { RadrootsOutboxDeliveryPlanStatus::FailedTerminal - } else if deferred_count > 0 { - RadrootsOutboxDeliveryPlanStatus::DeferredUntilImplemented } else { RadrootsOutboxDeliveryPlanStatus::FailedTerminal }; @@ -2199,9 +2600,6 @@ async fn evaluate_delivery_plans( if ready_count > 0 { any_ready = 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 = ?", ) @@ -2217,7 +2615,6 @@ async fn evaluate_delivery_plans( all_complete, any_failed_terminal, any_ready, - any_deferred_until_implemented, }) } @@ -2276,6 +2673,16 @@ fn operation_from_row( operation_id: row.try_get("operation_id")?, operation_kind: row.try_get("operation_kind")?, expected_pubkey: row.try_get("expected_pubkey")?, + semantic_scope: row.try_get("semantic_scope")?, + trade_id: parse_optional_stored_trade_id( + "outbox_operations.trade_id", + row.try_get("trade_id")?, + )?, + mutation_id: parse_optional_stored_mutation_id( + "outbox_operations.mutation_id", + row.try_get("mutation_id")?, + )?, + canonical_payload_sha256: row.try_get("canonical_payload_sha256")?, idempotency_key: row.try_get("idempotency_key")?, operation_idempotency_digest: row.try_get("operation_idempotency_digest")?, status, @@ -2531,6 +2938,21 @@ struct OperationDigestInput<'a> { draft: &'a RadrootsEventDraft, } +struct TradeMutationSemantic { + trade_id: String, + mutation_id: String, + canonical_payload_sha256: String, +} + +#[derive(Serialize)] +struct TradeMutationOperationDigestInput<'a> { + operation_kind: &'a str, + expected_pubkey: &'a str, + trade_id: &'a str, + mutation_id: &'a str, + canonical_payload_sha256: &'a str, +} + fn operation_idempotency_digest( operation_kind: &str, expected_pubkey: &str, @@ -2544,6 +2966,20 @@ fn operation_idempotency_digest( sha256_json(&input) } +fn trade_mutation_operation_idempotency_digest( + operation_kind: &str, + expected_pubkey: &str, + semantic: &TradeMutationSemantic, +) -> String { + sha256_json(&TradeMutationOperationDigestInput { + operation_kind, + expected_pubkey, + trade_id: semantic.trade_id.as_str(), + mutation_id: semantic.mutation_id.as_str(), + canonical_payload_sha256: semantic.canonical_payload_sha256.as_str(), + }) +} + #[derive(Serialize)] struct TargetPolicyDigestInput<'a> { satisfaction_policy: String, @@ -2611,6 +3047,100 @@ fn sha256_json<T: Serialize>(value: &T) -> String { hex::encode(Sha256::digest(bytes)) } +fn validate_trade_mutation_input( + trade_id: &str, + mutation_id: &str, + canonical_payload_sha256: &str, + draft: &RadrootsEventDraft, +) -> Result<TradeMutationSemantic, RadrootsOutboxError> { + if !TRADE_MUTATION_EVENT_KINDS.contains(&draft.kind_u32()) { + return Err(RadrootsOutboxError::TradeMutationMetadataMismatch { field: "kind" }); + } + if canonical_payload_sha256.len() != 64 + || !canonical_payload_sha256 + .bytes() + .all(|byte| byte.is_ascii_hexdigit()) + { + return Err(RadrootsOutboxError::TradeMutationMetadataMismatch { + field: "canonical_payload_sha256", + }); + } + let parsed = trade_mutation_from_canonical_content(draft.content()) + .map_err(|_| RadrootsOutboxError::TradeMutationMetadataMismatch { field: "content" })?; + if parsed.contract_id != draft.contract_id() { + return Err(RadrootsOutboxError::TradeMutationMetadataMismatch { + field: "contract_id", + }); + } + if parsed.mutation_kind().nostr_kind() != draft.kind_u32() { + return Err(RadrootsOutboxError::TradeMutationMetadataMismatch { field: "kind" }); + } + if parsed.author_pubkey.as_str() != draft.expected_pubkey_str() { + return Err(RadrootsOutboxError::TradeMutationMetadataMismatch { + field: "author_pubkey", + }); + } + if parsed.trade_id.as_str() != trade_id { + return Err(RadrootsOutboxError::TradeMutationMetadataMismatch { field: "trade_id" }); + } + let Some(parsed_mutation_id) = parsed.mutation_id.as_ref() else { + return Err(RadrootsOutboxError::TradeMutationMetadataMismatch { + field: "mutation_id", + }); + }; + if parsed_mutation_id.as_str() != mutation_id { + return Err(RadrootsOutboxError::TradeMutationMetadataMismatch { + field: "mutation_id", + }); + } + let payload_sha256 = sha256_hex(draft.content().as_bytes()); + if canonical_payload_sha256 != payload_sha256 { + return Err(RadrootsOutboxError::TradeMutationMetadataMismatch { + field: "canonical_payload_sha256", + }); + } + Ok(TradeMutationSemantic { + trade_id: trade_id.to_owned(), + mutation_id: mutation_id.to_owned(), + canonical_payload_sha256: canonical_payload_sha256.to_owned(), + }) +} + +fn sha256_hex(bytes: &[u8]) -> String { + hex::encode(Sha256::digest(bytes)) +} + +fn ensure_not_trade_mutation_draft(draft: &RadrootsEventDraft) -> Result<(), RadrootsOutboxError> { + if TRADE_MUTATION_EVENT_KINDS.contains(&draft.kind_u32()) { + return Err(RadrootsOutboxError::TradeMutationRequiresSemanticOutbox); + } + Ok(()) +} + +fn parse_optional_stored_trade_id( + field: &'static str, + value: Option<String>, +) -> Result<Option<RadrootsTradeId>, RadrootsOutboxError> { + value + .map(|value| { + RadrootsTradeId::parse(value.as_str()) + .map_err(|_| RadrootsOutboxError::InvalidStoredIdentifier { field, value }) + }) + .transpose() +} + +fn parse_optional_stored_mutation_id( + field: &'static str, + value: Option<String>, +) -> Result<Option<RadrootsTradeMutationId>, RadrootsOutboxError> { + value + .map(|value| { + RadrootsTradeMutationId::parse(value.as_str()) + .map_err(|_| RadrootsOutboxError::InvalidStoredIdentifier { field, value }) + }) + .transpose() +} + fn satisfaction_policy_storage_value(policy: &RadrootsTransportSatisfactionPolicy) -> String { match policy { RadrootsTransportSatisfactionPolicy::NoWait => "no_wait".to_owned(), @@ -2753,7 +3283,19 @@ fn u32_from_i64(field: &'static str, value: i64) -> Result<u32, RadrootsOutboxEr #[cfg(test)] mod tests { use super::*; - use radroots_event::kinds::KIND_POST; + use radroots_event::ids::{ + RadrootsAddressableCoordinate, RadrootsDTag, RadrootsEventId, RadrootsInventoryBinId, + RadrootsPublicKey, RadrootsTradeId, + }; + use radroots_event::kinds::{KIND_LISTING, KIND_POST}; + use radroots_event::trade::{ + RADROOTS_TRADE_PROPOSAL_CONTRACT_ID, RADROOTS_TRADE_SCHEMA_VERSION, + RadrootsFulfillmentProfileV1, RadrootsTradeCancellationProfileV1, + RadrootsTradeCandidateLineV1, RadrootsTradeCandidateTermsV1, + RadrootsTradeCanonicalMutationV1, RadrootsTradeEconomicAdjustmentV1, + RadrootsTradeEconomicsProfileV1, RadrootsTradeMutationBodyV1, + RadrootsTradeMutationEnvelopeV1, canonical_trade_mutation_content, + }; use radroots_nostr::prelude::{ RadrootsNostrKeys, RadrootsNostrSecretKey, radroots_nostr_sign_frozen_draft, }; @@ -2769,6 +3311,17 @@ mod tests { std::iter::repeat_n(character, 64).collect() } + fn hex_32(character: char) -> String { + std::iter::repeat_n(character, 32).collect() + } + + fn public_key(character: char) -> RadrootsPublicKey { + if character == 'a' { + return RadrootsPublicKey::parse(FIXTURE_ALICE_PUBLIC_KEY_HEX).expect("alice pubkey"); + } + RadrootsPublicKey::parse(hex_64(character)).expect("pubkey") + } + fn post_draft(expected_pubkey: &str, content: &str) -> RadrootsEventDraft { RadrootsEventDraft::new( "radroots.social.post.v1", @@ -2781,6 +3334,128 @@ mod tests { .expect("post draft") } + fn candidate_terms() -> RadrootsTradeCandidateTermsV1 { + RadrootsTradeCandidateTermsV1 { + candidate_id: None, + schema_version: RADROOTS_TRADE_SCHEMA_VERSION, + base_candidate_id: None, + supersession_intent: None, + buyer_pubkey: public_key('a'), + seller_pubkey: public_key('a'), + farm_id: RadrootsDTag::parse("farm-1").expect("farm id"), + lines: vec![RadrootsTradeCandidateLineV1 { + line_id: RadrootsDTag::parse("line-1").expect("line id"), + listing_addr: RadrootsAddressableCoordinate::parse(format!( + "{KIND_LISTING}:{}:listing-1", + FIXTURE_ALICE_PUBLIC_KEY_HEX + )) + .expect("listing address"), + listing_event_id: RadrootsEventId::parse(hex_64('c')).expect("listing event id"), + listing_snapshot_sha256: hex_64('d'), + product_id: "carrots".to_owned(), + option_id: None, + bin_id: RadrootsInventoryBinId::parse("bin-1").expect("bin id"), + quantity_mantissa: "2".to_owned(), + quantity_scale: 0, + unit_code: "count".to_owned(), + unit_profile: "mvp-count".to_owned(), + unit_price_mantissa: "300".to_owned(), + currency_code: "USD".to_owned(), + line_subtotal_mantissa: "600".to_owned(), + replaces_line_id: None, + }], + line_tombstones: Vec::new(), + economics: RadrootsTradeEconomicsProfileV1 { + profile_id: "mvp-no-payment".to_owned(), + currency_code: "USD".to_owned(), + currency_exponent: 2, + rounding_profile: "half-up".to_owned(), + subtotal_mantissa: "600".to_owned(), + discount_total_mantissa: "0".to_owned(), + adjustment_total_mantissa: "0".to_owned(), + total_mantissa: "600".to_owned(), + adjustments: Vec::<RadrootsTradeEconomicAdjustmentV1>::new(), + }, + fulfillment: RadrootsFulfillmentProfileV1 { + profile_id: "market-pickup".to_owned(), + method: "pickup".to_owned(), + starts_at_unix_s: 1_800_000_000, + ends_at_unix_s: 1_800_003_600, + timezone: "America/New_York".to_owned(), + utc_offset_seconds: -18_000, + fold: 0, + location_class: "private_after_agreement".to_owned(), + requires_private_terms: false, + }, + cancellation: RadrootsTradeCancellationProfileV1 { + profile_id: "mvp".to_owned(), + buyer_pre_agreement: true, + post_agreement_cutoff_unix_s: Some(1_799_999_000), + }, + private_terms: None, + proposal_expires_at_unix_s: 1_799_999_000, + } + } + + fn proposal_envelope() -> RadrootsTradeMutationEnvelopeV1 { + RadrootsTradeMutationEnvelopeV1 { + mutation_id: None, + contract_id: RADROOTS_TRADE_PROPOSAL_CONTRACT_ID.to_owned(), + schema_version: RADROOTS_TRADE_SCHEMA_VERSION, + trade_id: RadrootsTradeId::parse(hex_32('1')).expect("trade id"), + root_mutation_id: None, + buyer_pubkey: public_key('a'), + seller_pubkey: public_key('a'), + farm_id: RadrootsDTag::parse("farm-1").expect("farm id"), + parent_mutation_ids: Vec::new(), + author_pubkey: public_key('a'), + counterparty_pubkey: public_key('a'), + authored_at_unix_s: 1_799_000_000, + body: RadrootsTradeMutationBodyV1::Proposal { + candidate: candidate_terms(), + }, + } + } + + fn canonical_trade_proposal() -> RadrootsTradeCanonicalMutationV1 { + canonical_trade_mutation_content(proposal_envelope()).expect("canonical trade proposal") + } + + fn trade_mutation_draft(canonical: &RadrootsTradeCanonicalMutationV1) -> RadrootsEventDraft { + let mut tags = vec![ + vec![ + "contract".to_owned(), + canonical.envelope.contract_id.clone(), + ], + vec!["d".to_owned(), canonical.mutation_id.to_string()], + vec![ + "p".to_owned(), + canonical.envelope.counterparty_pubkey.to_string(), + ], + ]; + for parent in &canonical.envelope.parent_mutation_ids { + tags.push(vec!["e".to_owned(), parent.to_string()]); + } + RadrootsEventDraft::new( + canonical.envelope.contract_id.clone(), + canonical.envelope.mutation_kind().nostr_kind(), + canonical.envelope.authored_at_unix_s, + tags, + canonical.content.clone(), + canonical.envelope.author_pubkey.as_str(), + ) + .expect("trade mutation draft") + } + + fn signed_trade_mutation( + canonical: &RadrootsTradeCanonicalMutationV1, + ) -> (RadrootsEventDraft, RadrootsSignedEvent) { + let draft = trade_mutation_draft(canonical); + let signed_event = + radroots_nostr_sign_frozen_draft(&fixture_keys(), &draft).expect("signed trade event"); + (draft, signed_event) + } + fn nostr_target(uri: &str) -> RadrootsTransportTarget { RadrootsTransportTarget::nostr_relay(uri).expect("nostr target") } @@ -2985,6 +3660,27 @@ mod tests { ) } + fn signed_trade_mutation_input( + canonical: &RadrootsTradeCanonicalMutationV1, + draft: RadrootsEventDraft, + signed_event: RadrootsSignedEvent, + targets: Vec<RadrootsTransportTarget>, + created_at_ms: i64, + ) -> RadrootsOutboxSignedTradeMutationInput { + RadrootsOutboxSignedTradeMutationInput::new( + "publish_trade_mutation", + canonical.envelope.trade_id.clone(), + canonical.mutation_id.clone(), + sha256_hex(canonical.content.as_bytes()), + draft, + signed_event, + delivery_plan(targets), + true, + created_at_ms + 7, + created_at_ms, + ) + } + fn fixture_keys() -> RadrootsNostrKeys { let secret_key = RadrootsNostrSecretKey::from_hex(FIXTURE_ALICE_SECRET_KEY_HEX).expect("secret key"); @@ -3128,6 +3824,18 @@ mod tests { .await .expect("old table query"); assert!(old.is_none()); + let operation_columns = table_columns(&outbox, "outbox_operations").await; + for column in [ + "semantic_scope", + "trade_id", + "mutation_id", + "canonical_payload_sha256", + ] { + assert!( + operation_columns.iter().any(|name| name == column), + "{column}" + ); + } let target_columns = table_columns(&outbox, "outbox_delivery_target").await; for column in ["target_scope", "target_label", "last_outcome_kind"] { assert!(target_columns.iter().any(|name| name == column), "{column}"); @@ -3205,6 +3913,158 @@ mod tests { } #[tokio::test] + async fn generic_enqueue_rejects_trade_mutation_drafts_before_persistence() { + let outbox = RadrootsOutbox::open_memory().await.expect("open"); + let canonical = canonical_trade_proposal(); + let (draft, signed_event) = signed_trade_mutation(&canonical); + + let unsigned_err = outbox + .enqueue_operation(RadrootsOutboxOperationInput::new( + "publish_trade_mutation", + draft.clone(), + delivery_plan(vec![nostr_target(NOSTR_PRIMARY_WSS)]), + 1_000, + )) + .await + .expect_err("generic unsigned trade rejection"); + assert!(matches!( + unsigned_err, + RadrootsOutboxError::TradeMutationRequiresSemanticOutbox + )); + + let signed_err = outbox + .enqueue_signed_operation(RadrootsOutboxSignedOperationInput::new( + "publish_trade_mutation", + draft, + signed_event, + delivery_plan(vec![nostr_target(NOSTR_PRIMARY_WSS)]), + true, + 1_007, + 1_000, + )) + .await + .expect_err("generic signed trade rejection"); + assert!(matches!( + signed_err, + RadrootsOutboxError::TradeMutationRequiresSemanticOutbox + )); + 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); + } + + #[tokio::test] + async fn semantic_trade_mutation_enqueue_persists_metadata_and_deduplicates_by_mutation() { + let outbox = RadrootsOutbox::open_memory().await.expect("open"); + let canonical = canonical_trade_proposal(); + let payload_sha256 = sha256_hex(canonical.content.as_bytes()); + let (draft, signed_event) = signed_trade_mutation(&canonical); + let input = signed_trade_mutation_input( + &canonical, + draft.clone(), + signed_event.clone(), + vec![nostr_target(NOSTR_PRIMARY_WSS)], + 1_000, + ) + .with_idempotency_key("trade-first"); + + let preflight = outbox + .preflight_signed_trade_mutation_idempotency(&input) + .await + .expect("preflight"); + let first = outbox + .enqueue_signed_trade_mutation_operation(input) + .await + .expect("first enqueue"); + assert_eq!(first.status, RadrootsOutboxEnqueueStatus::Inserted); + assert_eq!( + first.operation_idempotency_digest, + preflight.operation_idempotency_digest + ); + assert_eq!( + first.delivery_plan_idempotency_digest, + preflight.delivery_plan_idempotency_digest + ); + + let operation = outbox + .get_operation(first.operation_id) + .await + .expect("operation") + .expect("operation"); + assert_eq!(operation.semantic_scope, "trade_mutation"); + assert_eq!( + operation.trade_id.as_ref(), + Some(&canonical.envelope.trade_id) + ); + assert_eq!(operation.mutation_id.as_ref(), Some(&canonical.mutation_id)); + assert_eq!( + operation.canonical_payload_sha256.as_deref(), + Some(payload_sha256.as_str()) + ); + assert_eq!(operation.status, RadrootsOutboxOperationStatus::Queued); + + let event = outbox + .get_event(first.outbox_event_id) + .await + .expect("event") + .expect("event"); + assert_eq!(event.state, RadrootsOutboxEventState::Signed); + assert_eq!(event.signed_event, Some(signed_event)); + + let second = outbox + .enqueue_signed_trade_mutation_operation( + signed_trade_mutation_input( + &canonical, + draft, + event.signed_event.expect("stored signed event"), + vec![nostr_target("wss://relay-3.example.com")], + 1_100, + ) + .with_idempotency_key("trade-second"), + ) + .await + .expect("second enqueue"); + assert_eq!(second.status, RadrootsOutboxEnqueueStatus::Inserted); + assert_eq!(second.operation_id, first.operation_id); + assert_eq!(second.outbox_event_id, first.outbox_event_id); + assert_eq!(table_count(&outbox, "outbox_operations").await, 1); + assert_eq!(table_count(&outbox, "outbox_event").await, 1); + assert_eq!(table_count(&outbox, "outbox_delivery_plan").await, 2); + } + + #[tokio::test] + async fn semantic_trade_mutation_enqueue_rejects_payload_hash_mismatch() { + let outbox = RadrootsOutbox::open_memory().await.expect("open"); + let canonical = canonical_trade_proposal(); + let (draft, signed_event) = signed_trade_mutation(&canonical); + + let err = outbox + .enqueue_signed_trade_mutation_operation(RadrootsOutboxSignedTradeMutationInput::new( + "publish_trade_mutation", + canonical.envelope.trade_id, + canonical.mutation_id, + "0".repeat(64), + draft, + signed_event, + delivery_plan(vec![nostr_target(NOSTR_PRIMARY_WSS)]), + true, + 1_007, + 1_000, + )) + .await + .expect_err("hash mismatch"); + match err { + RadrootsOutboxError::TradeMutationMetadataMismatch { field } => { + assert_eq!(field, "canonical_payload_sha256"); + } + other => panic!("unexpected error: {other}"), + } + 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); + } + + #[tokio::test] async fn enqueue_rejects_empty_delivery_targets_before_persistence() { let outbox = RadrootsOutbox::open_memory().await.expect("open"); let draft = post_draft(hex_64('a').as_str(), "hello"); @@ -4417,17 +5277,14 @@ mod tests { .expect("plans"); assert_eq!( plans[0].status, - RadrootsOutboxDeliveryPlanStatus::DeferredUntilImplemented + RadrootsOutboxDeliveryPlanStatus::FailedTerminal ); let event = outbox .get_event(receipt.outbox_event_id) .await .expect("event") .expect("event"); - assert_eq!( - event.state, - RadrootsOutboxEventState::DeferredUntilImplemented - ); + assert_eq!(event.state, RadrootsOutboxEventState::FailedTerminal); let operation = outbox .get_operation(receipt.operation_id) .await @@ -4435,10 +5292,12 @@ mod tests { .expect("operation"); assert_eq!( operation.status, - RadrootsOutboxOperationStatus::DeferredUntilImplemented + RadrootsOutboxOperationStatus::FailedTerminal ); let summary = outbox.status_summary(1_000).await.expect("summary"); assert_eq!(summary.pending_events, 0); + assert_eq!(summary.terminal_events, 1); + assert_eq!(summary.failed_terminal_events, 1); assert_eq!(summary.ready_signed_events, 0); assert_eq!(summary.deferred_until_implemented_events, 1); assert!( @@ -4489,17 +5348,14 @@ mod tests { .expect("plans"); assert_eq!( plans[0].status, - RadrootsOutboxDeliveryPlanStatus::DeferredUntilImplemented + RadrootsOutboxDeliveryPlanStatus::FailedTerminal ); let event = outbox .get_event(receipt.outbox_event_id) .await .expect("event") .expect("event"); - assert_eq!( - event.state, - RadrootsOutboxEventState::DeferredUntilImplemented - ); + assert_eq!(event.state, RadrootsOutboxEventState::FailedTerminal); let operation = outbox .get_operation(receipt.operation_id) .await @@ -4507,10 +5363,12 @@ mod tests { .expect("operation"); assert_eq!( operation.status, - RadrootsOutboxOperationStatus::DeferredUntilImplemented + RadrootsOutboxOperationStatus::FailedTerminal ); let summary = outbox.status_summary(1_000).await.expect("summary"); assert_eq!(summary.pending_events, 0); + assert_eq!(summary.terminal_events, 1); + assert_eq!(summary.failed_terminal_events, 1); assert_eq!(summary.ready_signed_events, 0); assert_eq!(summary.deferred_until_implemented_events, 1); assert!( @@ -4667,7 +5525,7 @@ mod tests { assert_eq!(records[0].event.outbox_event_id, reject.outbox_event_id); assert_eq!( records[0].event.state, - RadrootsOutboxEventState::DeferredUntilImplemented + RadrootsOutboxEventState::FailedTerminal ); assert_eq!(records[0].targets.len(), 1); assert_eq!( @@ -4696,7 +5554,7 @@ mod tests { assert_eq!(records[1].event.outbox_event_id, deferred.outbox_event_id); assert_eq!( records[1].event.state, - RadrootsOutboxEventState::DeferredUntilImplemented + RadrootsOutboxEventState::FailedTerminal ); assert_eq!(records[1].targets.len(), 1); assert_eq!( @@ -4750,7 +5608,7 @@ mod tests { } #[tokio::test] - async fn complete_signing_sets_deferred_until_implemented_lifecycle_for_reticulum_only_plans() { + async fn complete_signing_sets_failed_terminal_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 @@ -4786,10 +5644,7 @@ mod tests { .await .expect("event") .expect("event"); - assert_eq!( - event.state, - RadrootsOutboxEventState::DeferredUntilImplemented - ); + assert_eq!(event.state, RadrootsOutboxEventState::FailedTerminal); assert_eq!(event.claim_token, None); let operation = outbox .get_operation(receipt.operation_id) @@ -4798,7 +5653,7 @@ mod tests { .expect("operation"); assert_eq!( operation.status, - RadrootsOutboxOperationStatus::DeferredUntilImplemented + RadrootsOutboxOperationStatus::FailedTerminal ); assert!( outbox @@ -5002,7 +5857,7 @@ mod tests { ); let summary = outbox.status_summary(1_000).await.expect("summary"); assert_eq!(summary.pending_events, 1); - assert_eq!(summary.deferred_until_implemented_events, 0); + assert_eq!(summary.deferred_until_implemented_events, 1); } #[tokio::test] @@ -5246,17 +6101,17 @@ mod tests { "transport deferred", "deferred", RadrootsOutboxDeliveryTargetStatus::DeferredUntilImplemented, - RadrootsOutboxDeliveryPlanStatus::DeferredUntilImplemented, - RadrootsOutboxEventState::DeferredUntilImplemented, - RadrootsOutboxOperationStatus::DeferredUntilImplemented, + RadrootsOutboxDeliveryPlanStatus::FailedTerminal, + RadrootsOutboxEventState::FailedTerminal, + RadrootsOutboxOperationStatus::FailedTerminal, ), ( "transport unavailable", "unavailable", RadrootsOutboxDeliveryTargetStatus::DeferredUntilImplemented, - RadrootsOutboxDeliveryPlanStatus::DeferredUntilImplemented, - RadrootsOutboxEventState::DeferredUntilImplemented, - RadrootsOutboxOperationStatus::DeferredUntilImplemented, + RadrootsOutboxDeliveryPlanStatus::FailedTerminal, + RadrootsOutboxEventState::FailedTerminal, + RadrootsOutboxOperationStatus::FailedTerminal, ), ] { let outbox = RadrootsOutbox::open_memory().await.expect("open"); @@ -5361,8 +6216,8 @@ mod tests { total_events: 1, pending_events: 0, retryable_events: 0, - terminal_events: 0, - failed_terminal_events: 0, + terminal_events: 1, + failed_terminal_events: 1, deferred_until_implemented_events: if expected_status == RadrootsOutboxDeliveryTargetStatus::DeferredUntilImplemented { @@ -5373,7 +6228,7 @@ mod tests { ready_signed_events: 0, publishing_events: 0, last_attempt_at_ms: Some(1_100), - last_error: None, + last_error: Some("terminal".to_owned()), } ); }