commit 4e226a29668323adfc682b626b1812bc1a4d6c89
parent aadecc93d40c200e594190ed9c5863df811b5f67
Author: triesap <tyson@radroots.org>
Date: Tue, 30 Jun 2026 06:24:10 +0000
trade: route studio actions through sdk facade
- replace order enqueue requests with product trade facade commands
- ingest signed trade evidence into SDK projection before mutations
- rename SDK migration receipts to workflow receipts and remove audit scaffolding
- update sync payload contracts to trade operation names
Diffstat:
14 files changed, 1217 insertions(+), 2711 deletions(-)
diff --git a/Cargo.lock b/Cargo.lock
@@ -5327,6 +5327,7 @@ name = "radroots_sdk"
version = "0.1.0"
dependencies = [
"hex",
+ "nostr",
"radroots_authority",
"radroots_event_store",
"radroots_events",
@@ -5339,7 +5340,6 @@ dependencies = [
"radroots_outbox",
"radroots_relay_transport",
"radroots_runtime_paths",
- "radroots_sp1_guest_trade",
"radroots_trade",
"serde",
"serde_json",
@@ -5354,16 +5354,6 @@ name = "radroots_secret_vault"
version = "0.1.0-alpha.2"
[[package]]
-name = "radroots_sp1_guest_trade"
-version = "0.1.0-alpha.2"
-dependencies = [
- "serde",
- "serde_json",
- "sha2",
- "thiserror 1.0.69",
-]
-
-[[package]]
name = "radroots_sql_core"
version = "0.1.0-alpha.2"
dependencies = [
@@ -5425,6 +5415,7 @@ dependencies = [
"radroots_runtime_paths",
"radroots_sdk",
"radroots_studio_app_view",
+ "radroots_trade",
"serde",
"serde_json",
"thiserror 2.0.18",
@@ -5490,6 +5481,7 @@ dependencies = [
name = "radroots_studio_app_sync"
version = "0.1.0"
dependencies = [
+ "radroots_events",
"radroots_sdk",
"radroots_studio_app_view",
"serde",
diff --git a/crates/desktop/src/runtime.rs b/crates/desktop/src/runtime.rs
@@ -20,8 +20,17 @@ use radroots_events::{
KIND_ORDER_REQUEST, KIND_ORDER_REVISION_DECISION, KIND_ORDER_REVISION_PROPOSAL,
KIND_PROFILE,
},
+ order::{
+ RadrootsOrderDecision, RadrootsOrderEconomics, RadrootsOrderEventType,
+ RadrootsOrderInventoryCommitment, RadrootsOrderItem, RadrootsOrderRequest,
+ RadrootsOrderRevisionDecision, RadrootsOrderRevisionOutcome, RadrootsOrderRevisionProposal,
+ },
+};
+use radroots_events_codec::order::{
+ order_cancellation_from_event, order_decision_from_event, order_event_context_from_tags,
+ order_request_from_event, order_revision_decision_from_event,
+ order_revision_proposal_from_event,
};
-use radroots_events_codec::order::order_event_context_from_tags;
use radroots_identity::{RadrootsIdentity, RadrootsIdentityId};
use radroots_local_events::{
BUYER_ORDER_REQUEST_ACTOR_SOURCE_RESOLVED_ACCOUNT,
@@ -45,27 +54,21 @@ use radroots_sdk::protocol::listing::{
RadrootsListingDeliveryMethod, RadrootsListingProduct, RadrootsListingPublicLocation,
RadrootsListingStatus,
};
-use radroots_sdk::protocol::order::{
- RadrootsOrderCancellation, RadrootsOrderDecision, RadrootsOrderDecisionOutcome,
- RadrootsOrderEconomics, RadrootsOrderInventoryCommitment, RadrootsOrderItem,
- RadrootsOrderRequest, RadrootsOrderRevisionDecision, RadrootsOrderRevisionOutcome,
- RadrootsOrderRevisionProposal,
-};
use radroots_sdk::{
- FARM_PUBLISH_OPERATION_KIND, LISTING_PUBLISH_OPERATION_KIND, ORDER_CANCELLATION_OPERATION_KIND,
- ORDER_DECISION_OPERATION_KIND, ORDER_REVISION_DECISION_OPERATION_KIND,
- ORDER_REVISION_PROPOSAL_OPERATION_KIND, ORDER_SUBMIT_OPERATION_KIND,
+ FARM_PUBLISH_OPERATION_KIND, LISTING_PUBLISH_OPERATION_KIND, TRADE_CANCELLATION_OPERATION_KIND,
+ TRADE_DECISION_OPERATION_KIND, TRADE_REVISION_DECISION_OPERATION_KIND,
+ TRADE_REVISION_PROPOSAL_OPERATION_KIND, TRADE_SUBMIT_OPERATION_KIND,
};
use radroots_sql_core::SqliteExecutor;
use radroots_studio_app_core::{
AppBuildIdentity, AppDesktopRuntimePaths, AppRuntimeCapture, AppRuntimeMode,
AppRuntimePathsError, AppRuntimeSnapshot, AppSdkConfig, AppSdkDiagnostics,
AppSdkFarmPublicLocationRequest, AppSdkFarmPublishRequest, AppSdkLifecycleState,
- AppSdkListingPublishRequest, AppSdkOrderCancellationRequest, AppSdkOrderDecisionRequest,
- AppSdkOrderRevisionDecisionRequest, AppSdkOrderRevisionProposalRequest,
- AppSdkOrderSubmitRequest, AppSdkProjectionLifecycleState, AppSdkPublicFarmLocation,
+ AppSdkListingPublishRequest, AppSdkProjectionLifecycleState, AppSdkPublicFarmLocation,
AppSdkRelayUrlPolicy, AppSdkRuntime, AppSdkRuntimeError, AppSdkRuntimeIssue,
- AppSdkRuntimeStatus, AppSdkStoragePaths, AppSdkWorkflowReceipt, AppSharedAccountsPaths,
+ AppSdkRuntimeStatus, AppSdkStoragePaths, AppSdkTradeCancellationRequest, AppSdkTradeDecision,
+ AppSdkTradeDecisionRequest, AppSdkTradeProposeRequest, AppSdkTradeRevisionDecisionRequest,
+ AppSdkTradeRevisionProposalRequest, AppSdkWorkflowReceipt, AppSharedAccountsPaths,
PackDayExportWriteError, prepare_pack_day_export_bundle_at_data_root,
shared_local_events_database_path_from_shared_accounts, write_prepared_pack_day_export_bundle,
};
@@ -73,8 +76,8 @@ use radroots_studio_app_remote_signer::{
RadrootsAppRemoteSignerApprovedSession, RadrootsAppRemoteSignerPendingSession,
};
use radroots_studio_app_sqlite::{
- APP_ACTIVITY_CONTEXT_LIMIT, AppLocalInteropImportReport, AppSdkMigrationReceiptInput,
- AppSdkMigrationReceiptSourceKind, AppSdkMigrationState, AppSqliteError, AppSqliteStore,
+ APP_ACTIVITY_CONTEXT_LIMIT, AppLocalInteropImportReport, AppSdkWorkflowReceiptInput,
+ AppSdkWorkflowReceiptSourceKind, AppSdkWorkflowReceiptState, AppSqliteError, AppSqliteStore,
BuyerOrderLocalEventExport, BuyerOrderLocalEventLine, BuyerRepeatDemandApplyOutcome,
DatabaseTarget, SelectedBuyerOrderScope, SellerOrderDecisionExport, StoredPendingSyncOperation,
StoredRelayIngestCursor, StoredSyncConflict, derive_farm_rules_readiness,
@@ -119,6 +122,7 @@ use radroots_studio_app_view::{
ReminderUrgency, SettingsAccountProjection, SettingsPreference, SettingsSection, ShellSection,
TodayAgendaProjection,
};
+use radroots_trade::identity::RadrootsTradeLocator;
use radroots_trade::listing::parse_public_listing_address;
use radroots_trade::order::{
RadrootsOrderCancellationRecord, RadrootsOrderDecisionRecord, RadrootsOrderReductionInputs,
@@ -2678,7 +2682,7 @@ impl DesktopAppRuntimeState {
let source_record_id = order_decision_sdk_source_record_id(&payload);
self.enqueue_order_decision_payload_via_sdk(
&payload,
- AppSdkMigrationReceiptSourceKind::LocalOutbox,
+ AppSdkWorkflowReceiptSourceKind::LocalOutbox,
source_record_id.as_str(),
)?;
let _ = self.refresh_selected_account_sync()?;
@@ -2784,7 +2788,6 @@ impl DesktopAppRuntimeState {
farm_id,
trade_order_id: request.payload.order_id.to_string(),
request_event_id: request.request_event_id,
- prev_event_id: lifecycle.request_event_id,
revision_id: format!("app-revision-{}", d_tag_from_uuid(Uuid::now_v7())),
listing_addr: request.payload.listing_addr.to_string(),
buyer_pubkey: request.payload.buyer_pubkey.to_string(),
@@ -2813,7 +2816,7 @@ impl DesktopAppRuntimeState {
let source_record_id = order_revision_proposal_sdk_source_record_id(&payload);
self.enqueue_order_revision_proposal_payload_via_sdk(
&payload,
- AppSdkMigrationReceiptSourceKind::LocalOutbox,
+ AppSdkWorkflowReceiptSourceKind::LocalOutbox,
source_record_id.as_str(),
)?;
let _ = self.refresh_selected_account_sync()?;
@@ -2908,7 +2911,6 @@ impl DesktopAppRuntimeState {
farm_id: detail.farm_id,
trade_order_id: request.payload.order_id.to_string(),
request_event_id: request.request_event_id,
- prev_event_id: proposal.event_id.clone(),
revision_id: proposal.payload.revision_id.to_string(),
listing_addr: request.payload.listing_addr.to_string(),
buyer_pubkey: request.payload.buyer_pubkey.to_string(),
@@ -2951,7 +2953,7 @@ impl DesktopAppRuntimeState {
let source_record_id = order_revision_decision_sdk_source_record_id(&payload);
self.enqueue_order_revision_decision_payload_via_sdk(
&payload,
- AppSdkMigrationReceiptSourceKind::LocalOutbox,
+ AppSdkWorkflowReceiptSourceKind::LocalOutbox,
source_record_id.as_str(),
)?;
let _ = self.refresh_selected_account_sync()?;
@@ -3025,8 +3027,8 @@ impl DesktopAppRuntimeState {
reason: "buyer order cancellation requires no pending seller proposal",
});
}
- let prev_event_id = match lifecycle.status {
- RadrootsTradeWorkflowState::Requested => request.request_event_id.clone(),
+ match lifecycle.status {
+ RadrootsTradeWorkflowState::Requested => {}
RadrootsTradeWorkflowState::RevisionProposed
| RadrootsTradeWorkflowState::AgreedPendingRhi
| RadrootsTradeWorkflowState::Committed => {
@@ -3049,7 +3051,6 @@ impl DesktopAppRuntimeState {
farm_id: detail.farm_id,
trade_order_id: request.payload.order_id.to_string(),
request_event_id: request.request_event_id,
- prev_event_id,
listing_addr: request.payload.listing_addr.to_string(),
buyer_pubkey: request.payload.buyer_pubkey.to_string(),
seller_pubkey: request.payload.seller_pubkey.to_string(),
@@ -3071,7 +3072,7 @@ impl DesktopAppRuntimeState {
let source_record_id = order_cancellation_sdk_source_record_id(&payload);
self.enqueue_order_cancellation_payload_via_sdk(
&payload,
- AppSdkMigrationReceiptSourceKind::LocalOutbox,
+ AppSdkWorkflowReceiptSourceKind::LocalOutbox,
source_record_id.as_str(),
)?;
let _ = self.refresh_selected_account_sync()?;
@@ -4437,7 +4438,7 @@ impl DesktopAppRuntimeState {
.unwrap_or_else(|| format!("app:order_request:{}", payload.order_id));
self.enqueue_order_request_payload_via_sdk(
&payload,
- AppSdkMigrationReceiptSourceKind::SharedLocalEvent,
+ AppSdkWorkflowReceiptSourceKind::SharedLocalEvent,
source_record_id.as_str(),
)?;
self.refresh_selected_account_sync()
@@ -4680,7 +4681,7 @@ impl DesktopAppRuntimeState {
fn enqueue_farm_profile_payload_via_sdk(
&self,
payload: &AppFarmProfilePublishPayload,
- source_kind: AppSdkMigrationReceiptSourceKind,
+ source_kind: AppSdkWorkflowReceiptSourceKind,
source_record_id: &str,
) -> Result<(), AppSqliteError> {
let operation_kind = FARM_PUBLISH_OPERATION_KIND;
@@ -4709,14 +4710,14 @@ impl DesktopAppRuntimeState {
.map_err(sync_transport_error_from_sdk_runtime_error)
});
match actor_pubkey {
- Ok((actor_pubkey, receipt)) => self.record_app_sdk_migration_success(
+ Ok((actor_pubkey, receipt)) => self.record_app_sdk_workflow_success(
source_kind,
source_record_id,
operation_kind,
actor_pubkey.as_str(),
&receipt,
),
- Err(error) => self.record_app_sdk_migration_failure(
+ Err(error) => self.record_app_sdk_workflow_failure(
source_kind,
source_record_id,
operation_kind,
@@ -4729,7 +4730,7 @@ impl DesktopAppRuntimeState {
fn enqueue_listing_payload_via_sdk(
&self,
payload: &AppListingPublishPayload,
- source_kind: AppSdkMigrationReceiptSourceKind,
+ source_kind: AppSdkWorkflowReceiptSourceKind,
source_record_id: &str,
) -> Result<(), DesktopAppRuntimeProductPublishError> {
let operation_kind = LISTING_PUBLISH_OPERATION_KIND;
@@ -4763,7 +4764,7 @@ impl DesktopAppRuntimeState {
});
match actor_pubkey {
Ok((actor_pubkey, receipt)) => self
- .record_app_sdk_migration_success(
+ .record_app_sdk_workflow_success(
source_kind,
source_record_id,
operation_kind,
@@ -4772,7 +4773,7 @@ impl DesktopAppRuntimeState {
)
.map_err(DesktopAppRuntimeProductPublishError::from),
Err(error) => {
- self.record_app_sdk_migration_failure(
+ self.record_app_sdk_workflow_failure(
source_kind,
source_record_id,
operation_kind,
@@ -4787,41 +4788,37 @@ impl DesktopAppRuntimeState {
fn enqueue_order_request_payload_via_sdk(
&self,
payload: &AppOrderRequestPublishPayload,
- source_kind: AppSdkMigrationReceiptSourceKind,
+ source_kind: AppSdkWorkflowReceiptSourceKind,
source_record_id: &str,
) -> Result<(), AppSqliteError> {
- let operation_kind = ORDER_SUBMIT_OPERATION_KIND;
+ let operation_kind = TRADE_SUBMIT_OPERATION_KIND;
let actor_pubkey = self
.local_signing_identity_for_publish_payload(&AppPublishPayload::OrderRequest(
payload.clone(),
))
.and_then(|identity| {
let actor_pubkey = identity.public_key_hex();
- let target_relays =
- order_request_sdk_target_relays(payload, self.nostr_relay_urls.as_slice())?;
- let request = AppSdkOrderSubmitRequest {
+ let request = AppSdkTradeProposeRequest {
actor_account_id: payload.context.account_id.clone(),
actor_pubkey: actor_pubkey.clone(),
signer_keys: identity.into_keys(),
listing_event: order_request_sdk_listing_event_ptr(payload)?,
order: order_request_publish_payload_to_sdk_order(payload)?,
- relay_url_policy: sdk_relay_url_policy_for_targets(target_relays.as_slice()),
- target_relays,
idempotency_key: Some(sdk_idempotency_key(source_record_id)),
};
- self.enqueue_app_sdk_order_submit(request)
+ self.enqueue_app_sdk_trade_propose(request)
.map(|receipt| (actor_pubkey, receipt))
.map_err(sync_transport_error_from_sdk_runtime_error)
});
match actor_pubkey {
- Ok((actor_pubkey, receipt)) => self.record_app_sdk_migration_success(
+ Ok((actor_pubkey, receipt)) => self.record_app_sdk_workflow_success(
source_kind,
source_record_id,
operation_kind,
actor_pubkey.as_str(),
&receipt,
),
- Err(error) => self.record_app_sdk_migration_failure(
+ Err(error) => self.record_app_sdk_workflow_failure(
source_kind,
source_record_id,
operation_kind,
@@ -4834,45 +4831,39 @@ impl DesktopAppRuntimeState {
fn enqueue_order_decision_payload_via_sdk(
&self,
payload: &AppOrderDecisionPublishPayload,
- source_kind: AppSdkMigrationReceiptSourceKind,
+ source_kind: AppSdkWorkflowReceiptSourceKind,
source_record_id: &str,
) -> Result<(), AppSqliteError> {
- let operation_kind = ORDER_DECISION_OPERATION_KIND;
- let request_evidence = self.resolve_seller_order_request_evidence(payload.app_order_id)?;
+ let evidence_events = self.sdk_trade_request_evidence_events(payload.app_order_id)?;
+ let operation_kind = TRADE_DECISION_OPERATION_KIND;
let actor_pubkey = self
.local_signing_identity_for_publish_payload(&AppPublishPayload::OrderDecision(
payload.clone(),
))
.and_then(|identity| {
let actor_pubkey = identity.public_key_hex();
- let target_relays = normalized_app_sync_relay_urls(&self.nostr_relay_urls)?;
- let request = AppSdkOrderDecisionRequest {
+ let request = AppSdkTradeDecisionRequest {
actor_account_id: payload.context.account_id.clone(),
actor_pubkey: actor_pubkey.clone(),
signer_keys: identity.into_keys(),
- request_event: request_evidence.request_event,
- request_event_ptr: order_decision_sdk_request_event_ptr(
- payload,
- target_relays.as_slice(),
- )?,
- decision: order_decision_publish_payload_to_sdk_decision(payload)?,
- relay_url_policy: sdk_relay_url_policy_for_targets(target_relays.as_slice()),
- target_relays,
+ evidence_events: evidence_events.clone(),
+ locator: trade_locator_from_decision_payload(payload)?,
+ decision: trade_decision_from_publish_payload(payload)?,
idempotency_key: Some(sdk_idempotency_key(source_record_id)),
};
- self.enqueue_app_sdk_order_decision(request)
+ self.enqueue_app_sdk_trade_decision(request)
.map(|receipt| (actor_pubkey, receipt))
.map_err(sync_transport_error_from_sdk_runtime_error)
});
match actor_pubkey {
- Ok((actor_pubkey, receipt)) => self.record_app_sdk_migration_success(
+ Ok((actor_pubkey, receipt)) => self.record_app_sdk_workflow_success(
source_kind,
source_record_id,
operation_kind,
actor_pubkey.as_str(),
&receipt,
),
- Err(error) => self.record_app_sdk_migration_failure(
+ Err(error) => self.record_app_sdk_workflow_failure(
source_kind,
source_record_id,
operation_kind,
@@ -4885,52 +4876,42 @@ impl DesktopAppRuntimeState {
fn enqueue_order_revision_proposal_payload_via_sdk(
&self,
payload: &AppOrderRevisionProposalPublishPayload,
- source_kind: AppSdkMigrationReceiptSourceKind,
+ source_kind: AppSdkWorkflowReceiptSourceKind,
source_record_id: &str,
) -> Result<(), AppSqliteError> {
- let operation_kind = ORDER_REVISION_PROPOSAL_OPERATION_KIND;
- let request_evidence = self.resolve_seller_order_request_evidence(payload.app_order_id)?;
- let lifecycle = self.resolve_order_lifecycle_evidence(&request_evidence)?;
+ let evidence_events = self.sdk_trade_lifecycle_evidence_events(payload.app_order_id)?;
+ let operation_kind = TRADE_REVISION_PROPOSAL_OPERATION_KIND;
let actor_pubkey = self
.local_signing_identity_for_publish_payload(&AppPublishPayload::OrderRevisionProposal(
payload.clone(),
))
.and_then(|identity| {
let actor_pubkey = identity.public_key_hex();
- let target_relays = normalized_app_sync_relay_urls(&self.nostr_relay_urls)?;
- let request = AppSdkOrderRevisionProposalRequest {
+ let request = AppSdkTradeRevisionProposalRequest {
actor_account_id: payload.context.account_id.clone(),
actor_pubkey: actor_pubkey.clone(),
signer_keys: identity.into_keys(),
- evidence_events: lifecycle.evidence_events,
- root_event: order_lifecycle_sdk_event_ptr(
- payload.request_event_id.as_str(),
- target_relays.as_slice(),
- "order revision proposal requires request event id",
- )?,
- previous_event: order_lifecycle_sdk_event_ptr(
- payload.prev_event_id.as_str(),
- target_relays.as_slice(),
- "order revision proposal requires previous event id",
- )?,
- proposal: order_revision_proposal_publish_payload_to_sdk_revision(payload)?,
- relay_url_policy: sdk_relay_url_policy_for_targets(target_relays.as_slice()),
- target_relays,
+ evidence_events: evidence_events.clone(),
+ locator: trade_locator_from_revision_proposal_payload(payload)?,
+ revision_id: publish_revision_id(payload.revision_id.as_str())?,
+ items: payload.items.clone(),
+ economics: payload.economics.clone(),
+ reason: payload.reason.clone(),
idempotency_key: Some(sdk_idempotency_key(source_record_id)),
};
- self.enqueue_app_sdk_order_revision_proposal(request)
+ self.enqueue_app_sdk_trade_revision_proposal(request)
.map(|receipt| (actor_pubkey, receipt))
.map_err(sync_transport_error_from_sdk_runtime_error)
});
match actor_pubkey {
- Ok((actor_pubkey, receipt)) => self.record_app_sdk_migration_success(
+ Ok((actor_pubkey, receipt)) => self.record_app_sdk_workflow_success(
source_kind,
source_record_id,
operation_kind,
actor_pubkey.as_str(),
&receipt,
),
- Err(error) => self.record_app_sdk_migration_failure(
+ Err(error) => self.record_app_sdk_workflow_failure(
source_kind,
source_record_id,
operation_kind,
@@ -4943,54 +4924,40 @@ impl DesktopAppRuntimeState {
fn enqueue_order_revision_decision_payload_via_sdk(
&self,
payload: &AppOrderRevisionDecisionPublishPayload,
- source_kind: AppSdkMigrationReceiptSourceKind,
+ source_kind: AppSdkWorkflowReceiptSourceKind,
source_record_id: &str,
) -> Result<(), AppSqliteError> {
- let operation_kind = ORDER_REVISION_DECISION_OPERATION_KIND;
- let request_evidence = self.resolve_seller_order_request_evidence(payload.app_order_id)?;
- let lifecycle = self.resolve_order_lifecycle_evidence(&request_evidence)?;
+ let evidence_events = self.sdk_trade_lifecycle_evidence_events(payload.app_order_id)?;
+ let operation_kind = TRADE_REVISION_DECISION_OPERATION_KIND;
let actor_pubkey = self
.local_signing_identity_for_publish_payload(&AppPublishPayload::OrderRevisionDecision(
payload.clone(),
))
.and_then(|identity| {
let actor_pubkey = identity.public_key_hex();
- let target_relays = normalized_app_sync_relay_urls(&self.nostr_relay_urls)?;
- let request = AppSdkOrderRevisionDecisionRequest {
+ let request = AppSdkTradeRevisionDecisionRequest {
actor_account_id: payload.context.account_id.clone(),
actor_pubkey: actor_pubkey.clone(),
signer_keys: identity.into_keys(),
- evidence_events: lifecycle.evidence_events,
- root_event: order_lifecycle_sdk_event_ptr(
- payload.request_event_id.as_str(),
- target_relays.as_slice(),
- "order revision decision requires request event id",
- )?,
- previous_event: order_lifecycle_sdk_event_ptr(
- payload.prev_event_id.as_str(),
- target_relays.as_slice(),
- "order revision decision requires previous event id",
- )?,
- decision: order_revision_decision_publish_payload_to_sdk_revision_decision(
- payload,
- )?,
- relay_url_policy: sdk_relay_url_policy_for_targets(target_relays.as_slice()),
- target_relays,
+ evidence_events: evidence_events.clone(),
+ locator: trade_locator_from_revision_decision_payload(payload)?,
+ revision_id: publish_revision_id(payload.revision_id.as_str())?,
+ decision: payload.decision.clone(),
idempotency_key: Some(sdk_idempotency_key(source_record_id)),
};
- self.enqueue_app_sdk_order_revision_decision(request)
+ self.enqueue_app_sdk_trade_revision_decision(request)
.map(|receipt| (actor_pubkey, receipt))
.map_err(sync_transport_error_from_sdk_runtime_error)
});
match actor_pubkey {
- Ok((actor_pubkey, receipt)) => self.record_app_sdk_migration_success(
+ Ok((actor_pubkey, receipt)) => self.record_app_sdk_workflow_success(
source_kind,
source_record_id,
operation_kind,
actor_pubkey.as_str(),
&receipt,
),
- Err(error) => self.record_app_sdk_migration_failure(
+ Err(error) => self.record_app_sdk_workflow_failure(
source_kind,
source_record_id,
operation_kind,
@@ -5003,52 +4970,39 @@ impl DesktopAppRuntimeState {
fn enqueue_order_cancellation_payload_via_sdk(
&self,
payload: &AppOrderCancellationPublishPayload,
- source_kind: AppSdkMigrationReceiptSourceKind,
+ source_kind: AppSdkWorkflowReceiptSourceKind,
source_record_id: &str,
) -> Result<(), AppSqliteError> {
- let operation_kind = ORDER_CANCELLATION_OPERATION_KIND;
- let request_evidence = self.resolve_seller_order_request_evidence(payload.app_order_id)?;
- let lifecycle = self.resolve_order_lifecycle_evidence(&request_evidence)?;
+ let evidence_events = self.sdk_trade_lifecycle_evidence_events(payload.app_order_id)?;
+ let operation_kind = TRADE_CANCELLATION_OPERATION_KIND;
let actor_pubkey = self
.local_signing_identity_for_publish_payload(&AppPublishPayload::OrderCancellation(
payload.clone(),
))
.and_then(|identity| {
let actor_pubkey = identity.public_key_hex();
- let target_relays = normalized_app_sync_relay_urls(&self.nostr_relay_urls)?;
- let request = AppSdkOrderCancellationRequest {
+ let request = AppSdkTradeCancellationRequest {
actor_account_id: payload.context.account_id.clone(),
actor_pubkey: actor_pubkey.clone(),
signer_keys: identity.into_keys(),
- evidence_events: lifecycle.evidence_events,
- root_event: order_lifecycle_sdk_event_ptr(
- payload.request_event_id.as_str(),
- target_relays.as_slice(),
- "order cancellation requires request event id",
- )?,
- previous_event: order_lifecycle_sdk_event_ptr(
- payload.prev_event_id.as_str(),
- target_relays.as_slice(),
- "order cancellation requires previous event id",
- )?,
- cancellation: order_cancellation_publish_payload_to_sdk_cancellation(payload)?,
- relay_url_policy: sdk_relay_url_policy_for_targets(target_relays.as_slice()),
- target_relays,
+ evidence_events: evidence_events.clone(),
+ locator: trade_locator_from_cancellation_payload(payload)?,
+ reason: payload.reason.clone(),
idempotency_key: Some(sdk_idempotency_key(source_record_id)),
};
- self.enqueue_app_sdk_order_cancellation(request)
+ self.enqueue_app_sdk_trade_cancellation(request)
.map(|receipt| (actor_pubkey, receipt))
.map_err(sync_transport_error_from_sdk_runtime_error)
});
match actor_pubkey {
- Ok((actor_pubkey, receipt)) => self.record_app_sdk_migration_success(
+ Ok((actor_pubkey, receipt)) => self.record_app_sdk_workflow_success(
source_kind,
source_record_id,
operation_kind,
actor_pubkey.as_str(),
&receipt,
),
- Err(error) => self.record_app_sdk_migration_failure(
+ Err(error) => self.record_app_sdk_workflow_failure(
source_kind,
source_record_id,
operation_kind,
@@ -5095,39 +5049,39 @@ impl DesktopAppRuntimeState {
self.with_app_sdk_runtime(|runtime| runtime.enqueue_listing_publish(request))
}
- fn enqueue_app_sdk_order_submit(
+ fn enqueue_app_sdk_trade_propose(
&self,
- request: AppSdkOrderSubmitRequest,
+ request: AppSdkTradeProposeRequest,
) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeError> {
- self.with_app_sdk_runtime(|runtime| runtime.enqueue_order_submit(request))
+ self.with_app_sdk_runtime(|runtime| runtime.trade_propose(request))
}
- fn enqueue_app_sdk_order_decision(
+ fn enqueue_app_sdk_trade_decision(
&self,
- request: AppSdkOrderDecisionRequest,
+ request: AppSdkTradeDecisionRequest,
) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeError> {
- self.with_app_sdk_runtime(|runtime| runtime.enqueue_order_decision(request))
+ self.with_app_sdk_runtime(|runtime| runtime.trade_decide(request))
}
- fn enqueue_app_sdk_order_revision_proposal(
+ fn enqueue_app_sdk_trade_revision_proposal(
&self,
- request: AppSdkOrderRevisionProposalRequest,
+ request: AppSdkTradeRevisionProposalRequest,
) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeError> {
- self.with_app_sdk_runtime(|runtime| runtime.enqueue_order_revision_proposal(request))
+ self.with_app_sdk_runtime(|runtime| runtime.trade_revision_propose(request))
}
- fn enqueue_app_sdk_order_revision_decision(
+ fn enqueue_app_sdk_trade_revision_decision(
&self,
- request: AppSdkOrderRevisionDecisionRequest,
+ request: AppSdkTradeRevisionDecisionRequest,
) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeError> {
- self.with_app_sdk_runtime(|runtime| runtime.enqueue_order_revision_decision(request))
+ self.with_app_sdk_runtime(|runtime| runtime.trade_revision_decide(request))
}
- fn enqueue_app_sdk_order_cancellation(
+ fn enqueue_app_sdk_trade_cancellation(
&self,
- request: AppSdkOrderCancellationRequest,
+ request: AppSdkTradeCancellationRequest,
) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeError> {
- self.with_app_sdk_runtime(|runtime| runtime.enqueue_order_cancellation(request))
+ self.with_app_sdk_runtime(|runtime| runtime.trade_cancel(request))
}
fn with_app_sdk_runtime<T>(
@@ -5144,9 +5098,9 @@ impl DesktopAppRuntimeState {
command(runtime)
}
- fn record_app_sdk_migration_success(
+ fn record_app_sdk_workflow_success(
&self,
- source_kind: AppSdkMigrationReceiptSourceKind,
+ source_kind: AppSdkWorkflowReceiptSourceKind,
source_record_id: &str,
operation_kind: &str,
actor_pubkey: &str,
@@ -5160,7 +5114,7 @@ impl DesktopAppRuntimeState {
"outbox_event_id": receipt.outbox_event_id,
"state": receipt.state,
});
- self.record_app_sdk_migration_receipt(AppSdkMigrationReceiptInput {
+ self.record_app_sdk_workflow_receipt(AppSdkWorkflowReceiptInput {
source_kind,
source_record_id: source_record_id.to_owned(),
sdk_operation_kind: operation_kind.to_owned(),
@@ -5168,21 +5122,21 @@ impl DesktopAppRuntimeState {
expected_event_id: Some(receipt.expected_event_id.clone()),
actor_pubkey: Some(actor_pubkey.to_owned()),
idempotency_digest_prefix: receipt.idempotency_digest_prefix.clone(),
- migration_state: AppSdkMigrationState::Enqueued,
+ workflow_state: AppSdkWorkflowReceiptState::Enqueued,
recorded_at: current_utc_timestamp(),
detail_json,
})
}
- fn record_app_sdk_migration_failure(
+ fn record_app_sdk_workflow_failure(
&self,
- source_kind: AppSdkMigrationReceiptSourceKind,
+ source_kind: AppSdkWorkflowReceiptSourceKind,
source_record_id: &str,
operation_kind: &str,
actor_pubkey: Option<&str>,
detail_json: serde_json::Value,
) -> Result<(), AppSqliteError> {
- self.record_app_sdk_migration_receipt(AppSdkMigrationReceiptInput {
+ self.record_app_sdk_workflow_receipt(AppSdkWorkflowReceiptInput {
source_kind,
source_record_id: source_record_id.to_owned(),
sdk_operation_kind: operation_kind.to_owned(),
@@ -5190,21 +5144,21 @@ impl DesktopAppRuntimeState {
expected_event_id: None,
actor_pubkey: actor_pubkey.map(str::to_owned),
idempotency_digest_prefix: None,
- migration_state: AppSdkMigrationState::Failed,
+ workflow_state: AppSdkWorkflowReceiptState::Failed,
recorded_at: current_utc_timestamp(),
detail_json,
})
}
- fn record_app_sdk_migration_receipt(
+ fn record_app_sdk_workflow_receipt(
&self,
- input: AppSdkMigrationReceiptInput,
+ input: AppSdkWorkflowReceiptInput,
) -> Result<(), AppSqliteError> {
let Some(sqlite_store) = self.sqlite_store.as_ref() else {
return Ok(());
};
let _ = sqlite_store
- .sdk_migration_receipt_repository()
+ .sdk_workflow_receipt_repository()
.record_receipt(&input)?;
Ok(())
}
@@ -5563,10 +5517,7 @@ impl DesktopAppRuntimeState {
for record in records {
if record.family != LocalRecordFamily::SignedEvent
|| record.event_kind
- != Some(i64::from(
- radroots_sdk::protocol::order::RadrootsOrderEventType::OrderRequested
- .kind(),
- ))
+ != Some(i64::from(RadrootsOrderEventType::OrderRequested.kind()))
|| !signed_order_request_evidence_record_is_usable(&record)
{
continue;
@@ -5574,8 +5525,7 @@ impl DesktopAppRuntimeState {
let Some(event) = signed_event_from_local_record(&record)? else {
continue;
};
- let Ok(envelope) = radroots_sdk::protocol::order::parse_order_request(&event)
- else {
+ let Ok(envelope) = order_request_from_event(&event) else {
continue;
};
insert_seller_order_request_evidence(
@@ -5603,11 +5553,11 @@ impl DesktopAppRuntimeState {
return Ok(());
};
let events = sqlite_store.load_local_interop_signed_events_by_kind(i64::from(
- radroots_sdk::protocol::order::RadrootsOrderEventType::OrderRequested.kind(),
+ RadrootsOrderEventType::OrderRequested.kind(),
))?;
for event in events {
- let Ok(envelope) = radroots_sdk::protocol::order::parse_order_request(&event) else {
+ let Ok(envelope) = order_request_from_event(&event) else {
continue;
};
insert_seller_order_request_evidence(
@@ -5656,10 +5606,11 @@ impl DesktopAppRuntimeState {
let author_pubkey = active_order_pubkey(event.author.as_str(), "author_pubkey")?;
match event.kind {
KIND_ORDER_DECISION => {
- let envelope = radroots_sdk::protocol::order::parse_order_decision(&event)
- .map_err(|_| AppSqliteError::InvalidProjection {
+ let envelope = order_decision_from_event(&event).map_err(|_| {
+ AppSqliteError::InvalidProjection {
reason: "order lifecycle evidence is invalid",
- })?;
+ }
+ })?;
let context = active_order_event_record_context(&event, envelope.message_type)?;
buckets.decisions.push(RadrootsOrderDecisionRecord {
event_id,
@@ -5671,9 +5622,7 @@ impl DesktopAppRuntimeState {
});
}
KIND_ORDER_REVISION_PROPOSAL => {
- let Ok(envelope) =
- radroots_sdk::protocol::order::parse_order_revision_proposal(&event)
- else {
+ let Ok(envelope) = order_revision_proposal_from_event(&event) else {
return Err(AppSqliteError::InvalidProjection {
reason: "order lifecycle evidence is invalid",
});
@@ -5691,9 +5640,7 @@ impl DesktopAppRuntimeState {
});
}
KIND_ORDER_REVISION_DECISION => {
- let Ok(envelope) =
- radroots_sdk::protocol::order::parse_order_revision_decision(&event)
- else {
+ let Ok(envelope) = order_revision_decision_from_event(&event) else {
return Err(AppSqliteError::InvalidProjection {
reason: "order lifecycle evidence is invalid",
});
@@ -5711,9 +5658,7 @@ impl DesktopAppRuntimeState {
});
}
KIND_ORDER_CANCELLATION => {
- let Ok(envelope) =
- radroots_sdk::protocol::order::parse_order_cancellation(&event)
- else {
+ let Ok(envelope) = order_cancellation_from_event(&event) else {
return Err(AppSqliteError::InvalidProjection {
reason: "order lifecycle evidence is invalid",
});
@@ -5804,6 +5749,23 @@ impl DesktopAppRuntimeState {
})
}
+ fn sdk_trade_request_evidence_events(
+ &self,
+ order_id: OrderId,
+ ) -> Result<Vec<SdkRadrootsNostrEvent>, AppSqliteError> {
+ let request = self.resolve_seller_order_request_evidence(order_id)?;
+ Ok(vec![request.request_event])
+ }
+
+ fn sdk_trade_lifecycle_evidence_events(
+ &self,
+ order_id: OrderId,
+ ) -> Result<Vec<SdkRadrootsNostrEvent>, AppSqliteError> {
+ let request = self.resolve_seller_order_request_evidence(order_id)?;
+ let lifecycle = self.resolve_order_lifecycle_evidence(&request)?;
+ Ok(lifecycle.evidence_events)
+ }
+
fn collect_order_lifecycle_signed_events(
&self,
) -> Result<Vec<SdkRadrootsNostrEvent>, AppSqliteError> {
@@ -7364,17 +7326,17 @@ fn farm_publish_source_record(
farm_id: FarmId,
source: &str,
source_local_event_id: Option<&str>,
-) -> (AppSdkMigrationReceiptSourceKind, String) {
+) -> (AppSdkWorkflowReceiptSourceKind, String) {
source_local_event_id
.map(|record_id| {
(
- AppSdkMigrationReceiptSourceKind::SharedLocalEvent,
+ AppSdkWorkflowReceiptSourceKind::SharedLocalEvent,
record_id.to_owned(),
)
})
.unwrap_or_else(|| {
(
- AppSdkMigrationReceiptSourceKind::LocalOutbox,
+ AppSdkWorkflowReceiptSourceKind::LocalOutbox,
format!("app:farm_publish:{farm_id}:{source}"),
)
})
@@ -7384,17 +7346,17 @@ fn listing_publish_source_record(
product_id: ProductId,
source: &str,
source_local_event_id: Option<&str>,
-) -> (AppSdkMigrationReceiptSourceKind, String) {
+) -> (AppSdkWorkflowReceiptSourceKind, String) {
source_local_event_id
.map(|record_id| {
(
- AppSdkMigrationReceiptSourceKind::SharedLocalEvent,
+ AppSdkWorkflowReceiptSourceKind::SharedLocalEvent,
record_id.to_owned(),
)
})
.unwrap_or_else(|| {
(
- AppSdkMigrationReceiptSourceKind::LocalOutbox,
+ AppSdkWorkflowReceiptSourceKind::LocalOutbox,
format!("app:listing_publish:{product_id}:{source}"),
)
})
@@ -7759,62 +7721,6 @@ fn order_request_sdk_listing_event_ptr(
})
}
-fn order_request_sdk_target_relays(
- payload: &AppOrderRequestPublishPayload,
- configured_relay_urls: &[String],
-) -> Result<Vec<String>, AppSyncTransportError> {
- let known_relays = normalized_listing_relays(payload.listing_relays.as_slice())?;
- let configured_relays = configured_relay_urls
- .iter()
- .map(|relay| relay.trim())
- .filter(|relay| !relay.is_empty())
- .map(str::to_owned)
- .collect::<BTreeSet<_>>();
- if configured_relays.is_empty() {
- return Ok(known_relays);
- }
- let selected_relays = known_relays
- .iter()
- .filter(|relay| configured_relays.contains(*relay))
- .cloned()
- .collect::<Vec<_>>();
- if selected_relays.is_empty() {
- return Ok(known_relays);
- }
- Ok(selected_relays)
-}
-
-fn order_decision_sdk_request_event_ptr(
- payload: &AppOrderDecisionPublishPayload,
- target_relays: &[String],
-) -> Result<RadrootsNostrEventPtr, AppSyncTransportError> {
- let request_event_id = payload.request_event_id.trim();
- if request_event_id.is_empty() {
- return Err(AppSyncTransportError::failed(
- "order decision publish requires request event id",
- ));
- }
- Ok(RadrootsNostrEventPtr {
- id: request_event_id.to_owned(),
- relays: target_relays.first().cloned(),
- })
-}
-
-fn order_lifecycle_sdk_event_ptr(
- event_id: &str,
- target_relays: &[String],
- missing_message: &'static str,
-) -> Result<RadrootsNostrEventPtr, AppSyncTransportError> {
- let event_id = event_id.trim();
- if event_id.is_empty() {
- return Err(AppSyncTransportError::failed(missing_message));
- }
- Ok(RadrootsNostrEventPtr {
- id: event_id.to_owned(),
- relays: target_relays.first().cloned(),
- })
-}
-
#[cfg(test)]
fn selected_listing_relay(
listing_relays: &[String],
@@ -9513,7 +9419,7 @@ fn active_order_pubkey(
fn active_order_event_record_context(
event: &SdkRadrootsNostrEvent,
- message_type: radroots_sdk::protocol::order::RadrootsOrderEventType,
+ message_type: RadrootsOrderEventType,
) -> Result<(RadrootsPublicKey, RadrootsEventId, RadrootsEventId), AppSqliteError> {
let context = order_event_context_from_tags(message_type, &event.tags).map_err(|_| {
AppSqliteError::InvalidProjection {
@@ -9672,79 +9578,91 @@ fn publish_bin_id(value: &str) -> Result<RadrootsInventoryBinId, AppSyncTranspor
.map_err(|error| AppSyncTransportError::failed(error.to_string()))
}
-fn order_decision_publish_payload_to_sdk_decision(
+fn trade_locator_from_parts(
+ trade_order_id: &str,
+ request_event_id: &str,
+ listing_addr: &str,
+ buyer_pubkey: &str,
+ seller_pubkey: &str,
+) -> Result<RadrootsTradeLocator, AppSyncTransportError> {
+ Ok(
+ RadrootsTradeLocator::from_order_id(publish_order_id(trade_order_id)?)
+ .with_root_event_id(publish_event_id(request_event_id)?)
+ .with_listing_addr(publish_listing_addr(listing_addr)?)
+ .with_buyer_pubkey(publish_pubkey(buyer_pubkey)?)
+ .with_seller_pubkey(publish_pubkey(seller_pubkey)?),
+ )
+}
+
+fn trade_locator_from_decision_payload(
payload: &AppOrderDecisionPublishPayload,
-) -> Result<RadrootsOrderDecision, AppSyncTransportError> {
- Ok(RadrootsOrderDecision {
- order_id: publish_order_id(payload.trade_order_id.as_str())?,
- listing_addr: publish_listing_addr(payload.listing_addr.as_str())?,
- buyer_pubkey: publish_pubkey(payload.buyer_pubkey.as_str())?,
- seller_pubkey: publish_pubkey(payload.seller_pubkey.as_str())?,
- decision: match &payload.decision {
- AppOrderDecisionPayload::Accepted {
- inventory_commitments,
- } => RadrootsOrderDecisionOutcome::Accepted {
- inventory_commitments: inventory_commitments
- .iter()
- .map(|commitment| {
- Ok(RadrootsOrderInventoryCommitment {
- bin_id: publish_bin_id(commitment.bin_id.as_str())?,
- bin_count: commitment.bin_count,
- })
- })
- .collect::<Result<Vec<_>, AppSyncTransportError>>()?,
- },
- AppOrderDecisionPayload::Declined { reason } => {
- RadrootsOrderDecisionOutcome::Declined {
- reason: reason.clone(),
- }
- }
- },
- })
+) -> Result<RadrootsTradeLocator, AppSyncTransportError> {
+ trade_locator_from_parts(
+ payload.trade_order_id.as_str(),
+ payload.request_event_id.as_str(),
+ payload.listing_addr.as_str(),
+ payload.buyer_pubkey.as_str(),
+ payload.seller_pubkey.as_str(),
+ )
}
-fn order_revision_proposal_publish_payload_to_sdk_revision(
+fn trade_locator_from_revision_proposal_payload(
payload: &AppOrderRevisionProposalPublishPayload,
-) -> Result<RadrootsOrderRevisionProposal, AppSyncTransportError> {
- Ok(RadrootsOrderRevisionProposal {
- revision_id: publish_revision_id(payload.revision_id.as_str())?,
- order_id: publish_order_id(payload.trade_order_id.as_str())?,
- listing_addr: publish_listing_addr(payload.listing_addr.as_str())?,
- buyer_pubkey: publish_pubkey(payload.buyer_pubkey.as_str())?,
- seller_pubkey: publish_pubkey(payload.seller_pubkey.as_str())?,
- root_event_id: publish_event_id(payload.request_event_id.as_str())?,
- prev_event_id: publish_event_id(payload.prev_event_id.as_str())?,
- items: payload.items.clone(),
- economics: payload.economics.clone(),
- reason: payload.reason.clone(),
- })
+) -> Result<RadrootsTradeLocator, AppSyncTransportError> {
+ trade_locator_from_parts(
+ payload.trade_order_id.as_str(),
+ payload.request_event_id.as_str(),
+ payload.listing_addr.as_str(),
+ payload.buyer_pubkey.as_str(),
+ payload.seller_pubkey.as_str(),
+ )
}
-fn order_revision_decision_publish_payload_to_sdk_revision_decision(
+fn trade_locator_from_revision_decision_payload(
payload: &AppOrderRevisionDecisionPublishPayload,
-) -> Result<RadrootsOrderRevisionDecision, AppSyncTransportError> {
- Ok(RadrootsOrderRevisionDecision {
- revision_id: publish_revision_id(payload.revision_id.as_str())?,
- order_id: publish_order_id(payload.trade_order_id.as_str())?,
- listing_addr: publish_listing_addr(payload.listing_addr.as_str())?,
- buyer_pubkey: publish_pubkey(payload.buyer_pubkey.as_str())?,
- seller_pubkey: publish_pubkey(payload.seller_pubkey.as_str())?,
- root_event_id: publish_event_id(payload.request_event_id.as_str())?,
- prev_event_id: publish_event_id(payload.prev_event_id.as_str())?,
- decision: payload.decision.clone(),
- })
+) -> Result<RadrootsTradeLocator, AppSyncTransportError> {
+ trade_locator_from_parts(
+ payload.trade_order_id.as_str(),
+ payload.request_event_id.as_str(),
+ payload.listing_addr.as_str(),
+ payload.buyer_pubkey.as_str(),
+ payload.seller_pubkey.as_str(),
+ )
}
-fn order_cancellation_publish_payload_to_sdk_cancellation(
+fn trade_locator_from_cancellation_payload(
payload: &AppOrderCancellationPublishPayload,
-) -> Result<RadrootsOrderCancellation, AppSyncTransportError> {
- Ok(RadrootsOrderCancellation {
- order_id: publish_order_id(payload.trade_order_id.as_str())?,
- listing_addr: publish_listing_addr(payload.listing_addr.as_str())?,
- buyer_pubkey: publish_pubkey(payload.buyer_pubkey.as_str())?,
- seller_pubkey: publish_pubkey(payload.seller_pubkey.as_str())?,
- reason: payload.reason.clone(),
- })
+) -> Result<RadrootsTradeLocator, AppSyncTransportError> {
+ trade_locator_from_parts(
+ payload.trade_order_id.as_str(),
+ payload.request_event_id.as_str(),
+ payload.listing_addr.as_str(),
+ payload.buyer_pubkey.as_str(),
+ payload.seller_pubkey.as_str(),
+ )
+}
+
+fn trade_decision_from_publish_payload(
+ payload: &AppOrderDecisionPublishPayload,
+) -> Result<AppSdkTradeDecision, AppSyncTransportError> {
+ match &payload.decision {
+ AppOrderDecisionPayload::Accepted {
+ inventory_commitments,
+ } => Ok(AppSdkTradeDecision::Accept {
+ inventory_commitments: inventory_commitments
+ .iter()
+ .map(|commitment| {
+ Ok(RadrootsOrderInventoryCommitment {
+ bin_id: publish_bin_id(commitment.bin_id.as_str())?,
+ bin_count: commitment.bin_count,
+ })
+ })
+ .collect::<Result<Vec<_>, AppSyncTransportError>>()?,
+ }),
+ AppOrderDecisionPayload::Declined { reason } => Ok(AppSdkTradeDecision::Decline {
+ reason: reason.clone(),
+ }),
+ }
}
#[cfg(test)]
@@ -9818,7 +9736,19 @@ mod tests {
RadrootsEventId, RadrootsInventoryBinId, RadrootsListingAddress, RadrootsOrderId,
RadrootsOrderQuoteId, RadrootsOrderRevisionId, RadrootsPublicKey,
};
- use radroots_events_codec::wire::WireEventParts;
+ use radroots_events::order::{
+ RadrootsOrderCancellation, RadrootsOrderDecision, RadrootsOrderDecisionOutcome,
+ RadrootsOrderEconomicItem, RadrootsOrderEconomics, RadrootsOrderInventoryCommitment,
+ RadrootsOrderItem, RadrootsOrderPricingBasis, RadrootsOrderRequest,
+ RadrootsOrderRevisionOutcome, RadrootsOrderRevisionProposal,
+ };
+ use radroots_events_codec::{
+ order::{
+ order_cancellation_event_build, order_decision_event_build, order_request_event_build,
+ order_revision_proposal_event_build,
+ },
+ wire::WireEventParts,
+ };
use radroots_identity::{RadrootsIdentity, RadrootsIdentityId};
use radroots_local_events::{
BUYER_ORDER_REQUEST_LOCAL_WORK_RECORD_KIND, LocalEventRecord, LocalEventRecordInput,
@@ -9837,29 +9767,24 @@ mod tests {
use radroots_sdk::protocol::events::{
RadrootsNostrEvent as SdkRadrootsNostrEvent, RadrootsNostrEventPtr,
};
- use radroots_sdk::protocol::order::{
- RadrootsOrderCancellation, RadrootsOrderDecision, RadrootsOrderDecisionOutcome,
- RadrootsOrderEconomicItem, RadrootsOrderEconomics, RadrootsOrderInventoryCommitment,
- RadrootsOrderItem, RadrootsOrderPricingBasis, RadrootsOrderRequest,
- RadrootsOrderRevisionOutcome, RadrootsOrderRevisionProposal,
- };
use radroots_sdk::{
- LISTING_PUBLISH_OPERATION_KIND, ORDER_CANCELLATION_OPERATION_KIND,
- ORDER_DECISION_OPERATION_KIND, ORDER_REVISION_DECISION_OPERATION_KIND,
- ORDER_SUBMIT_OPERATION_KIND,
+ LISTING_PUBLISH_OPERATION_KIND, TRADE_CANCELLATION_OPERATION_KIND,
+ TRADE_DECISION_OPERATION_KIND, TRADE_REVISION_DECISION_OPERATION_KIND,
+ TRADE_SUBMIT_OPERATION_KIND,
};
use radroots_sql_core::{SqlExecutor, SqliteExecutor};
use radroots_studio_app_core::{
AppDesktopRuntimePaths, AppRuntimeHostEnvironment, AppRuntimePlatform,
AppSdkLifecycleState, AppSdkProjectionLifecycleState, AppSdkPublicFarmLocation,
- AppSharedAccountsPaths, SHARED_ACCOUNTS_STORE_FILE_NAME, SHARED_IDENTITY_FILE_NAME,
+ AppSdkTradeDecision, AppSharedAccountsPaths, SHARED_ACCOUNTS_STORE_FILE_NAME,
+ SHARED_IDENTITY_FILE_NAME,
};
use radroots_studio_app_remote_signer::{
RadrootsAppRemoteSignerPendingSession, RadrootsAppRemoteSignerSessionRecord,
};
use radroots_studio_app_sqlite::{
- AppSdkMigrationReceiptSourceKind, AppSdkMigrationState, AppSqliteError, AppSqliteStore,
- BuyerOrderCoordinationState, DatabaseTarget, latest_schema_version,
+ AppSdkWorkflowReceiptSourceKind, AppSdkWorkflowReceiptState, AppSqliteError,
+ AppSqliteStore, BuyerOrderCoordinationState, DatabaseTarget, latest_schema_version,
projected_order_id_from_trade_request,
};
use radroots_studio_app_state::{
@@ -9921,9 +9846,8 @@ mod tests {
DesktopAppRuntimeCommandError, DesktopAppRuntimeMetadataSummary, DesktopAppRuntimeState,
DesktopAppSdkDiagnosticsState, DesktopAppSyncStatusSummary, DesktopRemoteSignerPaths,
SYNC_TRANSPORT_UNAVAILABLE_MESSAGE, TokioRuntimeBuilder, default_sync_transport,
- direct_relay_event_source_runtime, farm_sync_payload, is_hex_64,
- order_decision_publish_payload_to_sdk_decision, pending_sync_upsert,
- signed_event_from_local_record,
+ direct_relay_event_source_runtime, farm_sync_payload, is_hex_64, pending_sync_upsert,
+ signed_event_from_local_record, trade_decision_from_publish_payload,
};
use crate::pack_day_host_handoff::PackDayHostHandoffError;
use crate::pack_day_print::{
@@ -10319,6 +10243,35 @@ mod tests {
RadrootsListingAddress::parse(value).expect("listing address")
}
+ fn test_order_decision_publish_payload_to_event_payload(
+ payload: &AppOrderDecisionPublishPayload,
+ ) -> RadrootsOrderDecision {
+ RadrootsOrderDecision {
+ order_id: test_order_id(payload.trade_order_id.as_str()),
+ listing_addr: test_listing_addr(payload.listing_addr.as_str()),
+ buyer_pubkey: test_pubkey(payload.buyer_pubkey.as_str()),
+ seller_pubkey: test_pubkey(payload.seller_pubkey.as_str()),
+ decision: match &payload.decision {
+ AppOrderDecisionPayload::Accepted {
+ inventory_commitments,
+ } => RadrootsOrderDecisionOutcome::Accepted {
+ inventory_commitments: inventory_commitments
+ .iter()
+ .map(|commitment| RadrootsOrderInventoryCommitment {
+ bin_id: test_bin_id(commitment.bin_id.as_str()),
+ bin_count: commitment.bin_count,
+ })
+ .collect(),
+ },
+ AppOrderDecisionPayload::Declined { reason } => {
+ RadrootsOrderDecisionOutcome::Declined {
+ reason: reason.clone(),
+ }
+ }
+ },
+ }
+ }
+
fn install_recorded_sync_transport(
runtime: &DesktopAppRuntime,
transport: RecordedAppSyncTransport,
@@ -10625,7 +10578,6 @@ mod tests {
farm_id: common.1,
trade_order_id: common.2.clone(),
request_event_id: common.3.clone(),
- prev_event_id: test_event_id_seed("order-decision-event-1"),
revision_id: "revision-1".to_owned(),
listing_addr: common.4.clone(),
buyer_pubkey: common.5.clone(),
@@ -10647,7 +10599,6 @@ mod tests {
farm_id: common.1,
trade_order_id: common.2.clone(),
request_event_id: common.3.clone(),
- prev_event_id: test_event_id_seed("order-revision-proposal-event-1"),
revision_id: "revision-1".to_owned(),
listing_addr: common.4.clone(),
buyer_pubkey: common.5.clone(),
@@ -10664,7 +10615,6 @@ mod tests {
farm_id: common.1,
trade_order_id: common.2.clone(),
request_event_id: common.3.clone(),
- prev_event_id: common.3.clone(),
listing_addr: common.4.clone(),
buyer_pubkey: common.5.clone(),
seller_pubkey: common.6.clone(),
@@ -11818,16 +11768,16 @@ mod tests {
.sqlite_store
.as_ref()
.expect("sqlite store")
- .sdk_migration_receipt_repository()
+ .sdk_workflow_receipt_repository()
.load_receipt(
- AppSdkMigrationReceiptSourceKind::SharedLocalEvent,
+ AppSdkWorkflowReceiptSourceKind::SharedLocalEvent,
listing_record.record_id.as_str(),
)
- .expect("listing SDK migration receipt should load")
- .expect("listing SDK migration receipt should exist");
+ .expect("listing SDK workflow receipt should load")
+ .expect("listing SDK workflow receipt should exist");
assert_eq!(receipt.source_record_id, listing_record.record_id);
assert_eq!(receipt.sdk_operation_kind, LISTING_PUBLISH_OPERATION_KIND);
- assert_eq!(receipt.migration_state, AppSdkMigrationState::Enqueued);
+ assert_eq!(receipt.workflow_state, AppSdkWorkflowReceiptState::Enqueued);
assert!(receipt.expected_event_id.is_some());
assert!(
receipt
@@ -11953,16 +11903,16 @@ mod tests {
.sqlite_store
.as_ref()
.expect("sqlite store")
- .sdk_migration_receipt_repository()
+ .sdk_workflow_receipt_repository()
.load_receipt(
- AppSdkMigrationReceiptSourceKind::SharedLocalEvent,
+ AppSdkWorkflowReceiptSourceKind::SharedLocalEvent,
listing_record.record_id.as_str(),
)
- .expect("failed listing SDK migration receipt should load")
- .expect("failed listing SDK migration receipt should exist");
+ .expect("failed listing SDK workflow receipt should load")
+ .expect("failed listing SDK workflow receipt should exist");
assert_eq!(receipt.source_record_id, listing_record.record_id);
assert_eq!(receipt.sdk_operation_kind, LISTING_PUBLISH_OPERATION_KIND);
- assert_eq!(receipt.migration_state, AppSdkMigrationState::Failed);
+ assert_eq!(receipt.workflow_state, AppSdkWorkflowReceiptState::Failed);
assert!(receipt.sdk_outbox_event_ids.is_empty());
assert!(receipt.expected_event_id.is_none());
assert!(receipt.actor_pubkey.is_none());
@@ -11984,7 +11934,7 @@ mod tests {
.sqlite_store
.as_ref()
.expect("sqlite store")
- .sdk_migration_receipt_repository();
+ .sdk_workflow_receipt_repository();
retry_records
.iter()
.filter(|record| {
@@ -11997,12 +11947,12 @@ mod tests {
.filter_map(|record| {
repository
.load_receipt(
- AppSdkMigrationReceiptSourceKind::SharedLocalEvent,
+ AppSdkWorkflowReceiptSourceKind::SharedLocalEvent,
record.record_id.as_str(),
)
- .expect("retry listing SDK migration receipt should load")
+ .expect("retry listing SDK workflow receipt should load")
})
- .filter(|receipt| receipt.migration_state == AppSdkMigrationState::Enqueued)
+ .filter(|receipt| receipt.workflow_state == AppSdkWorkflowReceiptState::Enqueued)
.count()
};
assert!(enqueued_listing_receipts >= 1);
@@ -12120,11 +12070,14 @@ mod tests {
.sqlite_store
.as_ref()
.expect("sqlite store")
- .sdk_migration_receipt_repository()
+ .sdk_workflow_receipt_repository()
.load_receipt(source_kind, source_record_id.as_str())
- .expect("failed stock listing SDK migration receipt should load")
- .expect("failed stock listing SDK migration receipt should exist");
- assert_eq!(failed_receipt.migration_state, AppSdkMigrationState::Failed);
+ .expect("failed stock listing SDK workflow receipt should load")
+ .expect("failed stock listing SDK workflow receipt should exist");
+ assert_eq!(
+ failed_receipt.workflow_state,
+ AppSdkWorkflowReceiptState::Failed
+ );
assert_eq!(
failed_receipt.detail_json["code"],
"sdk_runtime_not_available"
@@ -12141,13 +12094,13 @@ mod tests {
.sqlite_store
.as_ref()
.expect("sqlite store")
- .sdk_migration_receipt_repository()
+ .sdk_workflow_receipt_repository()
.load_receipt(source_kind, source_record_id.as_str())
- .expect("retry stock listing SDK migration receipt should load")
- .expect("retry stock listing SDK migration receipt should exist");
+ .expect("retry stock listing SDK workflow receipt should load")
+ .expect("retry stock listing SDK workflow receipt should exist");
assert_eq!(
- retry_receipt.migration_state,
- AppSdkMigrationState::Enqueued
+ retry_receipt.workflow_state,
+ AppSdkWorkflowReceiptState::Enqueued
);
assert!(retry_receipt.expected_event_id.is_some());
assert!(
@@ -15331,8 +15284,8 @@ mod tests {
let payload = runtime
.prepare_order_accept(order_id)
.expect("seller order accept payload should prepare");
- let decision = order_decision_publish_payload_to_sdk_decision(&payload)
- .expect("order accept payload should convert to SDK decision");
+ let decision = trade_decision_from_publish_payload(&payload)
+ .expect("order accept payload should convert to SDK trade decision");
assert_eq!(payload.app_order_id, order_id);
assert_eq!(payload.trade_order_id, "seller-order-decision-1");
@@ -15346,10 +15299,9 @@ mod tests {
);
assert_eq!(payload.buyer_pubkey, buyer_pubkey);
assert_eq!(payload.seller_pubkey, seller_pubkey);
- assert_eq!(decision.order_id, "seller-order-decision-1");
- let RadrootsOrderDecisionOutcome::Accepted {
+ let AppSdkTradeDecision::Accept {
inventory_commitments,
- } = decision.decision
+ } = decision
else {
panic!("expected accepted decision");
};
@@ -15370,8 +15322,8 @@ mod tests {
let payload = runtime
.prepare_order_decline(order_id, " out of stock ")
.expect("seller order decline payload should prepare");
- let decision = order_decision_publish_payload_to_sdk_decision(&payload)
- .expect("order decline payload should convert to SDK decision");
+ let decision = trade_decision_from_publish_payload(&payload)
+ .expect("order decline payload should convert to SDK trade decision");
assert_eq!(payload.buyer_pubkey, buyer_pubkey);
assert_eq!(payload.seller_pubkey, seller_pubkey);
@@ -15381,7 +15333,7 @@ mod tests {
reason: "out of stock".to_owned()
}
);
- let RadrootsOrderDecisionOutcome::Declined { reason } = decision.decision else {
+ let AppSdkTradeDecision::Decline { reason } = decision else {
panic!("expected declined decision");
};
assert_eq!(reason, "out of stock");
@@ -15546,10 +15498,10 @@ mod tests {
&& record.event_kind == Some(3423)
&& record.event_pubkey.as_deref() == Some(seller_pubkey.as_str())
}));
- assert_order_decision_sdk_migration_receipt(
+ assert_order_decision_sdk_workflow_receipt(
&runtime,
order_id,
- AppSdkMigrationState::Enqueued,
+ AppSdkWorkflowReceiptState::Enqueued,
);
cleanup_bootstrapped_runtime_paths(&paths);
@@ -15575,10 +15527,10 @@ mod tests {
&& record.event_kind == Some(3423)
&& record.event_pubkey.as_deref() == Some(seller_pubkey.as_str())
}));
- assert_order_decision_sdk_migration_receipt(
+ assert_order_decision_sdk_workflow_receipt(
&runtime,
order_id,
- AppSdkMigrationState::Enqueued,
+ AppSdkWorkflowReceiptState::Enqueued,
);
cleanup_bootstrapped_runtime_paths(&paths);
@@ -15705,10 +15657,10 @@ mod tests {
buyer_account_id.as_str(),
order_id,
);
- assert_order_request_sdk_migration_receipt(
+ assert_order_request_sdk_workflow_receipt(
&runtime,
order_id,
- AppSdkMigrationState::Enqueued,
+ AppSdkWorkflowReceiptState::Enqueued,
);
{
@@ -15870,10 +15822,10 @@ mod tests {
buyer_account_id.as_str(),
order_id,
);
- assert_order_request_sdk_migration_receipt(
+ assert_order_request_sdk_workflow_receipt(
&runtime,
order_id,
- AppSdkMigrationState::Enqueued,
+ AppSdkWorkflowReceiptState::Enqueued,
);
assert_eq!(
summary_after_retry
@@ -15932,10 +15884,10 @@ mod tests {
buyer_account_id.as_str(),
order_id,
);
- assert_order_request_sdk_migration_receipt(
+ assert_order_request_sdk_workflow_receipt(
&runtime,
order_id,
- AppSdkMigrationState::Enqueued,
+ AppSdkWorkflowReceiptState::Enqueued,
);
cleanup_bootstrapped_runtime_paths(&paths);
@@ -16001,10 +15953,10 @@ mod tests {
buyer_account_id.as_str(),
order_id,
);
- assert_order_request_sdk_migration_receipt(
+ assert_order_request_sdk_workflow_receipt(
&restarted_runtime,
order_id,
- AppSdkMigrationState::Enqueued,
+ AppSdkWorkflowReceiptState::Enqueued,
);
cleanup_bootstrapped_runtime_paths(&paths);
@@ -16037,10 +15989,10 @@ mod tests {
buyer_account_id.as_str(),
order_id,
);
- assert_order_request_sdk_migration_receipt(
+ assert_order_request_sdk_workflow_receipt(
&runtime,
order_id,
- AppSdkMigrationState::Enqueued,
+ AppSdkWorkflowReceiptState::Enqueued,
);
{
let state = runtime.lock_state_mut();
@@ -16233,10 +16185,10 @@ mod tests {
let cancellation_events =
shared_order_events_by_kind(&fixture.paths, 3432, fixture.buyer_pubkey.as_str());
assert!(cancellation_events.is_empty());
- assert_order_cancellation_sdk_migration_receipt(
+ assert_order_cancellation_sdk_workflow_receipt(
&fixture.runtime,
fixture.order_id,
- AppSdkMigrationState::Enqueued,
+ AppSdkWorkflowReceiptState::Enqueued,
);
cleanup_bootstrapped_runtime_paths(&fixture.paths);
@@ -16404,11 +16356,11 @@ mod tests {
let revision_decision_events =
shared_order_events_by_kind(&fixture.paths, 3425, fixture.buyer_pubkey.as_str());
assert!(revision_decision_events.is_empty());
- assert_order_revision_decision_sdk_migration_receipt(
+ assert_order_revision_decision_sdk_workflow_receipt(
&fixture.runtime,
fixture.order_id,
revision_id.as_str(),
- AppSdkMigrationState::Enqueued,
+ AppSdkWorkflowReceiptState::Enqueued,
);
cleanup_bootstrapped_runtime_paths(&fixture.paths);
@@ -19986,15 +19938,9 @@ mod tests {
.expect("seller signer lookup should succeed")
.expect("seller account should have local signer");
let request_event_id = test_event_id(payload.request_event_id.as_str());
- let decision = order_decision_publish_payload_to_sdk_decision(&payload)
- .expect("order decision payload should convert to SDK decision");
- let parts = radroots_sdk::protocol::order::build_order_decision_draft(
- &request_event_id,
- &request_event_id,
- &decision,
- )
- .expect("order decision draft should build")
- .into_wire_parts();
+ let decision = test_order_decision_publish_payload_to_event_payload(&payload);
+ let parts = order_decision_event_build(&request_event_id, &request_event_id, &decision)
+ .expect("order decision draft should build");
let event = radroots_nostr_build_event(parts.kind, parts.content, parts.tags)
.expect("order decision event builder should build")
.sign_with_keys(identity.keys())
@@ -20139,15 +20085,14 @@ mod tests {
}],
economics: signed_order_request_economics(trade_order_id, order_quantity),
};
- let parts = radroots_sdk::protocol::order::build_order_request_draft(
+ let parts = order_request_event_build(
&RadrootsNostrEventPtr {
id: test_event_id_seed(listing_event_id),
relays: Some("wss://relay.example".to_owned()),
},
&order,
)
- .expect("order request draft should build")
- .into_wire_parts();
+ .expect("order request draft should build");
let record_id = format!("app:signed_event:order-request:{trade_order_id}");
let event_id = signed_event_id(record_id.as_str());
let event = test_event_from_parts(
@@ -20234,15 +20179,14 @@ mod tests {
}],
economics: signed_order_request_economics(trade_order_id, order_quantity),
};
- let parts = radroots_sdk::protocol::order::build_order_request_draft(
+ let parts = order_request_event_build(
&RadrootsNostrEventPtr {
id: test_event_id_seed(listing_event_id),
relays: Some("wss://relay.example".to_owned()),
},
&order,
)
- .expect("order request draft should build")
- .into_wire_parts();
+ .expect("order request draft should build");
let secret_key = RadrootsNostrSecretKey::from_hex(SDK_TEST_BUYER_SECRET_KEY_HEX)
.expect("SDK test buyer secret key should parse");
let keys = RadrootsNostrKeys::new(secret_key);
@@ -20347,13 +20291,8 @@ mod tests {
},
};
let request_event_id = test_event_id(request_event_id);
- let parts = radroots_sdk::protocol::order::build_order_decision_draft(
- &request_event_id,
- &request_event_id,
- &payload,
- )
- .expect("order decision draft should build")
- .into_wire_parts();
+ let parts = order_decision_event_build(&request_event_id, &request_event_id, &payload)
+ .expect("order decision draft should build");
let record_id = format!("app:signed_event:order-decision:{trade_order_id}");
append_trade_signed_event_record(
paths,
@@ -20383,13 +20322,8 @@ mod tests {
seller_pubkey: test_pubkey(seller_pubkey),
reason: "buyer cancelled order".to_owned(),
};
- let parts = radroots_sdk::protocol::order::build_order_cancellation_draft(
- &request_event_id,
- &prev_event_id,
- &payload,
- )
- .expect("order cancellation draft should build")
- .into_wire_parts();
+ let parts = order_cancellation_event_build(&request_event_id, &prev_event_id, &payload)
+ .expect("order cancellation draft should build");
let record_id = format!("app:signed_event:cancellation:{event_key}");
append_trade_signed_event_record(
paths,
@@ -20424,13 +20358,9 @@ mod tests {
economics: revision_test_order_economics(),
reason: "harvest count updated".to_owned(),
};
- let parts = radroots_sdk::protocol::order::build_order_revision_proposal_draft(
- &request_event_id,
- &prev_event_id,
- &payload,
- )
- .expect("order revision proposal draft should build")
- .into_wire_parts();
+ let parts =
+ order_revision_proposal_event_build(&request_event_id, &prev_event_id, &payload)
+ .expect("order revision proposal draft should build");
let record_id = format!("app:signed_event:revision-proposal:{event_key}");
append_trade_signed_event_record(
paths,
@@ -20876,10 +20806,10 @@ mod tests {
assert!(pending_order_sync_payloads(runtime, account_id, order_id).is_empty());
}
- fn assert_order_request_sdk_migration_receipt(
+ fn assert_order_request_sdk_workflow_receipt(
runtime: &DesktopAppRuntime,
order_id: OrderId,
- expected_state: AppSdkMigrationState,
+ expected_state: AppSdkWorkflowReceiptState,
) {
let source_record_id = format!("app:local_work:order_request:{order_id}");
let receipt = runtime
@@ -20887,31 +20817,31 @@ mod tests {
.sqlite_store
.as_ref()
.expect("sqlite store")
- .sdk_migration_receipt_repository()
+ .sdk_workflow_receipt_repository()
.load_receipt(
- AppSdkMigrationReceiptSourceKind::SharedLocalEvent,
+ AppSdkWorkflowReceiptSourceKind::SharedLocalEvent,
source_record_id.as_str(),
)
- .expect("SDK migration receipt should load")
- .expect("SDK migration receipt should exist");
+ .expect("SDK workflow receipt should load")
+ .expect("SDK workflow receipt should exist");
assert_eq!(receipt.source_record_id, source_record_id);
- assert_eq!(receipt.sdk_operation_kind, ORDER_SUBMIT_OPERATION_KIND);
+ assert_eq!(receipt.sdk_operation_kind, TRADE_SUBMIT_OPERATION_KIND);
assert_eq!(
- receipt.migration_state, expected_state,
+ receipt.workflow_state, expected_state,
"receipt detail: {}",
receipt.detail_json
);
- if expected_state == AppSdkMigrationState::Enqueued {
+ if expected_state == AppSdkWorkflowReceiptState::Enqueued {
assert!(receipt.expected_event_id.is_some());
assert!(receipt.actor_pubkey.as_deref().is_some_and(is_hex_64));
assert!(!receipt.sdk_outbox_event_ids.is_empty());
}
}
- fn assert_order_decision_sdk_migration_receipt(
+ fn assert_order_decision_sdk_workflow_receipt(
runtime: &DesktopAppRuntime,
order_id: OrderId,
- expected_state: AppSdkMigrationState,
+ expected_state: AppSdkWorkflowReceiptState,
) {
let source_record_id = format!("app:order_decision:{order_id}");
let receipt = runtime
@@ -20919,80 +20849,80 @@ mod tests {
.sqlite_store
.as_ref()
.expect("sqlite store")
- .sdk_migration_receipt_repository()
+ .sdk_workflow_receipt_repository()
.load_receipt(
- AppSdkMigrationReceiptSourceKind::LocalOutbox,
+ AppSdkWorkflowReceiptSourceKind::LocalOutbox,
source_record_id.as_str(),
)
- .expect("SDK migration receipt should load")
- .expect("SDK migration receipt should exist");
+ .expect("SDK workflow receipt should load")
+ .expect("SDK workflow receipt should exist");
assert_eq!(receipt.source_record_id, source_record_id);
- assert_eq!(receipt.sdk_operation_kind, ORDER_DECISION_OPERATION_KIND);
+ assert_eq!(receipt.sdk_operation_kind, TRADE_DECISION_OPERATION_KIND);
assert_eq!(
- receipt.migration_state, expected_state,
+ receipt.workflow_state, expected_state,
"receipt detail: {}",
receipt.detail_json
);
- if expected_state == AppSdkMigrationState::Enqueued {
+ if expected_state == AppSdkWorkflowReceiptState::Enqueued {
assert!(receipt.expected_event_id.is_some());
assert!(receipt.actor_pubkey.as_deref().is_some_and(is_hex_64));
assert!(!receipt.sdk_outbox_event_ids.is_empty());
}
}
- fn assert_order_revision_decision_sdk_migration_receipt(
+ fn assert_order_revision_decision_sdk_workflow_receipt(
runtime: &DesktopAppRuntime,
order_id: OrderId,
revision_id: &str,
- expected_state: AppSdkMigrationState,
+ expected_state: AppSdkWorkflowReceiptState,
) {
- assert_order_sdk_migration_receipt(
+ assert_order_sdk_workflow_receipt(
runtime,
format!("app:order_revision_decision:{order_id}:{revision_id}").as_str(),
- ORDER_REVISION_DECISION_OPERATION_KIND,
+ TRADE_REVISION_DECISION_OPERATION_KIND,
expected_state,
);
}
- fn assert_order_cancellation_sdk_migration_receipt(
+ fn assert_order_cancellation_sdk_workflow_receipt(
runtime: &DesktopAppRuntime,
order_id: OrderId,
- expected_state: AppSdkMigrationState,
+ expected_state: AppSdkWorkflowReceiptState,
) {
- assert_order_sdk_migration_receipt(
+ assert_order_sdk_workflow_receipt(
runtime,
format!("app:order_cancellation:{order_id}").as_str(),
- ORDER_CANCELLATION_OPERATION_KIND,
+ TRADE_CANCELLATION_OPERATION_KIND,
expected_state,
);
}
- fn assert_order_sdk_migration_receipt(
+ fn assert_order_sdk_workflow_receipt(
runtime: &DesktopAppRuntime,
source_record_id: &str,
operation_kind: &str,
- expected_state: AppSdkMigrationState,
+ expected_state: AppSdkWorkflowReceiptState,
) {
let receipt = runtime
.lock_state()
.sqlite_store
.as_ref()
.expect("sqlite store")
- .sdk_migration_receipt_repository()
+ .sdk_workflow_receipt_repository()
.load_receipt(
- AppSdkMigrationReceiptSourceKind::LocalOutbox,
+ AppSdkWorkflowReceiptSourceKind::LocalOutbox,
source_record_id,
)
- .expect("SDK migration receipt should load")
- .expect("SDK migration receipt should exist");
+ .expect("SDK workflow receipt should load")
+ .expect("SDK workflow receipt should exist");
assert_eq!(receipt.source_record_id, source_record_id);
assert_eq!(receipt.sdk_operation_kind, operation_kind);
assert_eq!(
- receipt.migration_state, expected_state,
+ receipt.workflow_state, expected_state,
"receipt detail: {}",
receipt.detail_json
);
- if expected_state == AppSdkMigrationState::Enqueued {
+ if expected_state == AppSdkWorkflowReceiptState::Enqueued {
assert!(receipt.expected_event_id.is_some());
assert!(receipt.actor_pubkey.as_deref().is_some_and(is_hex_64));
assert!(!receipt.sdk_outbox_event_ids.is_empty());
diff --git a/crates/runtime/Cargo.toml b/crates/runtime/Cargo.toml
@@ -17,6 +17,7 @@ radroots_local_events.workspace = true
radroots_nostr.workspace = true
radroots_runtime_paths.workspace = true
radroots_sdk.workspace = true
+radroots_trade.workspace = true
serde.workspace = true
serde_json.workspace = true
thiserror.workspace = true
diff --git a/crates/runtime/src/lib.rs b/crates/runtime/src/lib.rs
@@ -39,14 +39,15 @@ pub use sdk::{
APP_SDK_DEFAULT_COMMAND_QUEUE_CAPACITY, APP_SDK_STORAGE_DIR_NAME, AppSdkConfig,
AppSdkDiagnostics, AppSdkEventStoreDiagnostics, AppSdkFarmPublicLocationRequest,
AppSdkFarmPublishRequest, AppSdkIntegrityDiagnostics, AppSdkLifecycleState,
- AppSdkListingPublishRequest, AppSdkOrderCancellationRequest, AppSdkOrderDecisionRequest,
- AppSdkOrderRevisionDecisionRequest, AppSdkOrderRevisionProposalRequest,
- AppSdkOrderSubmitRequest, AppSdkOutboxDiagnostics, AppSdkProjectionLifecycleState,
+ AppSdkListingPublishRequest, AppSdkOutboxDiagnostics, AppSdkProjectionLifecycleState,
AppSdkProjectionLifecycleStatus, AppSdkPublicFarmLocation, AppSdkRelayUrlPolicy,
AppSdkRestorePreflightReceipt, AppSdkRestorePreflightRequest, AppSdkRuntime,
AppSdkRuntimeError, AppSdkRuntimeIssue, AppSdkRuntimeStatus, AppSdkSqliteStoreDiagnostics,
AppSdkStorageDiagnostics, AppSdkStoragePaths, AppSdkSyncDiagnostics,
AppSdkSyncEventStoreDiagnostics, AppSdkSyncOutboxDiagnostics, AppSdkSyncRelayTargetDiagnostics,
- AppSdkWorkflowReceipt, app_sdk_storage_root_from_data_root,
+ AppSdkTradeCancellationRequest, AppSdkTradeDecision, AppSdkTradeDecisionRequest,
+ AppSdkTradeProposeRequest, AppSdkTradeResyncRequest, AppSdkTradeRevisionDecisionRequest,
+ AppSdkTradeRevisionProposalRequest, AppSdkTradeStatusRequest, AppSdkWorkflowReceipt,
+ app_sdk_storage_root_from_data_root,
};
pub use startup::{AppStartupEvent, AppStartupEventMetadata, launch_startup_event};
diff --git a/crates/runtime/src/sdk.rs b/crates/runtime/src/sdk.rs
@@ -15,30 +15,35 @@ use radroots_events::{
RadrootsNostrEvent, RadrootsNostrEventPtr,
contract::RadrootsActorRole,
farm::RadrootsFarm,
- ids::RadrootsAddressableCoordinate,
+ ids::{RadrootsAddressableCoordinate, RadrootsOrderRevisionId},
kinds::KIND_FARM,
listing::RadrootsListing,
order::{
- RadrootsOrderCancellation, RadrootsOrderDecision, RadrootsOrderRequest,
- RadrootsOrderRevisionDecision, RadrootsOrderRevisionProposal,
+ RadrootsOrderEconomics, RadrootsOrderInventoryCommitment, RadrootsOrderItem,
+ RadrootsOrderRequest, RadrootsOrderRevisionOutcome,
},
};
use radroots_nostr::prelude::RadrootsNostrKeys;
use radroots_sdk::{
- FARM_PUBLISH_OPERATION_KIND, FarmEnqueuePublishRequest, FarmEnqueueReceipt, IntegrityReceipt,
- IntegrityRequest, LISTING_PUBLISH_OPERATION_KIND, ListingEnqueuePublishRequest,
- ListingEnqueueReceipt, ORDER_CANCELLATION_OPERATION_KIND, ORDER_DECISION_OPERATION_KIND,
- ORDER_REVISION_DECISION_OPERATION_KIND, ORDER_REVISION_PROPOSAL_OPERATION_KIND,
- ORDER_SUBMIT_OPERATION_KIND, OrderCancellationEnqueueRequest, OrderCancellationReceipt,
- OrderDecisionEnqueueRequest, OrderDecisionReceipt, OrderEvidenceIngestRequest,
- OrderRequestEvidenceIngestRequest, OrderRevisionDecisionEnqueueRequest,
- OrderRevisionDecisionReceipt, OrderRevisionProposalEnqueueRequest,
- OrderRevisionProposalReceipt, OrderSubmitEnqueueRequest, OrderSubmitReceipt, RadrootsClient,
- RadrootsSdkError, RadrootsSdkStoragePaths, RestoreReceipt, RestoreRequest,
- SdkBackupVerification, SdkPublicLocality, SdkRelayUrlPolicy as SdkRuntimeRelayUrlPolicy,
- StorageStatusReceipt, StorageStatusRequest, SyncStatusReceipt, SyncStatusRequest,
+ AckPolicy, FARM_PUBLISH_OPERATION_KIND, FarmEnqueuePublishRequest, FarmEnqueueReceipt,
+ IntegrityReceipt, IntegrityRequest, LISTING_PUBLISH_OPERATION_KIND,
+ ListingEnqueuePublishRequest, ListingEnqueueReceipt, PrivacyPreflightConfirmation,
+ ProductSensitivityField, PublishMode, RadrootsClient, RadrootsSdkError,
+ RadrootsSdkLocalKeySigner, RadrootsSdkSignerProvider, RadrootsSdkStoragePaths,
+ RelayResolutionPolicy, RestoreReceipt, RestoreRequest, SdkBackupVerification,
+ SdkPublicLocality, SdkRelayUrlPolicy as SdkRuntimeRelayUrlPolicy, StorageStatusReceipt,
+ StorageStatusRequest, SyncStatusReceipt, SyncStatusRequest, TRADE_CANCELLATION_OPERATION_KIND,
+ TRADE_DECISION_OPERATION_KIND, TRADE_REVISION_DECISION_OPERATION_KIND,
+ TRADE_REVISION_PROPOSAL_OPERATION_KIND, TRADE_SUBMIT_OPERATION_KIND, TradeAcceptRequest,
+ TradeCancelRequest, TradeCancellationPlan, TradeCancellationReceipt, TradeDecisionPlan,
+ TradeDecisionReceipt, TradeDeclineRequest, TradeEvidenceIngestRequest, TradeMutationOutcome,
+ TradeProposeRequest, TradeResyncReceipt, TradeResyncRequest, TradeRevisionDecisionPlan,
+ TradeRevisionDecisionReceipt, TradeRevisionDecisionRequest, TradeRevisionProposalPlan,
+ TradeRevisionProposalReceipt, TradeRevisionProposalRequest, TradeStatusReceipt,
+ TradeStatusRequest, TradeSubmitPlan, TradeSubmitReceipt,
};
use radroots_sdk::{SdkMutationState, SdkRelayTargetPolicy};
+use radroots_trade::identity::RadrootsTradeLocator;
use serde::Serialize;
use serde_json::{Value, json};
use thiserror::Error;
@@ -240,68 +245,77 @@ pub struct AppSdkListingPublishRequest {
pub idempotency_key: Option<String>,
}
-pub struct AppSdkOrderSubmitRequest {
+pub struct AppSdkTradeProposeRequest {
pub actor_account_id: String,
pub actor_pubkey: String,
pub signer_keys: RadrootsNostrKeys,
pub listing_event: RadrootsNostrEventPtr,
pub order: RadrootsOrderRequest,
- pub target_relays: Vec<String>,
- pub relay_url_policy: AppSdkRelayUrlPolicy,
pub idempotency_key: Option<String>,
}
-pub struct AppSdkOrderDecisionRequest {
+pub enum AppSdkTradeDecision {
+ Accept {
+ inventory_commitments: Vec<RadrootsOrderInventoryCommitment>,
+ },
+ Decline {
+ reason: String,
+ },
+}
+
+pub struct AppSdkTradeDecisionRequest {
pub actor_account_id: String,
pub actor_pubkey: String,
pub signer_keys: RadrootsNostrKeys,
- pub request_event: RadrootsNostrEvent,
- pub request_event_ptr: RadrootsNostrEventPtr,
- pub decision: RadrootsOrderDecision,
- pub target_relays: Vec<String>,
- pub relay_url_policy: AppSdkRelayUrlPolicy,
+ pub evidence_events: Vec<RadrootsNostrEvent>,
+ pub locator: RadrootsTradeLocator,
+ pub decision: AppSdkTradeDecision,
pub idempotency_key: Option<String>,
}
-pub struct AppSdkOrderRevisionProposalRequest {
+pub struct AppSdkTradeRevisionProposalRequest {
pub actor_account_id: String,
pub actor_pubkey: String,
pub signer_keys: RadrootsNostrKeys,
pub evidence_events: Vec<RadrootsNostrEvent>,
- pub root_event: RadrootsNostrEventPtr,
- pub previous_event: RadrootsNostrEventPtr,
- pub proposal: RadrootsOrderRevisionProposal,
- pub target_relays: Vec<String>,
- pub relay_url_policy: AppSdkRelayUrlPolicy,
+ pub locator: RadrootsTradeLocator,
+ pub revision_id: RadrootsOrderRevisionId,
+ pub items: Vec<RadrootsOrderItem>,
+ pub economics: RadrootsOrderEconomics,
+ pub reason: String,
pub idempotency_key: Option<String>,
}
-pub struct AppSdkOrderRevisionDecisionRequest {
+pub struct AppSdkTradeRevisionDecisionRequest {
pub actor_account_id: String,
pub actor_pubkey: String,
pub signer_keys: RadrootsNostrKeys,
pub evidence_events: Vec<RadrootsNostrEvent>,
- pub root_event: RadrootsNostrEventPtr,
- pub previous_event: RadrootsNostrEventPtr,
- pub decision: RadrootsOrderRevisionDecision,
- pub target_relays: Vec<String>,
- pub relay_url_policy: AppSdkRelayUrlPolicy,
+ pub locator: RadrootsTradeLocator,
+ pub revision_id: RadrootsOrderRevisionId,
+ pub decision: RadrootsOrderRevisionOutcome,
pub idempotency_key: Option<String>,
}
-pub struct AppSdkOrderCancellationRequest {
+pub struct AppSdkTradeCancellationRequest {
pub actor_account_id: String,
pub actor_pubkey: String,
pub signer_keys: RadrootsNostrKeys,
pub evidence_events: Vec<RadrootsNostrEvent>,
- pub root_event: RadrootsNostrEventPtr,
- pub previous_event: RadrootsNostrEventPtr,
- pub cancellation: RadrootsOrderCancellation,
- pub target_relays: Vec<String>,
- pub relay_url_policy: AppSdkRelayUrlPolicy,
+ pub locator: RadrootsTradeLocator,
+ pub reason: String,
pub idempotency_key: Option<String>,
}
+pub struct AppSdkTradeStatusRequest {
+ pub locator: RadrootsTradeLocator,
+}
+
+pub struct AppSdkTradeResyncRequest {
+ pub locator: RadrootsTradeLocator,
+ pub limit: u32,
+}
+
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct AppSdkWorkflowReceipt {
pub operation_kind: String,
@@ -406,26 +420,34 @@ enum AppSdkWorkerCommand {
AppSdkListingPublishRequest,
mpsc::Sender<Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue>>,
),
- EnqueueOrderSubmit(
- AppSdkOrderSubmitRequest,
+ TradePropose(
+ AppSdkTradeProposeRequest,
mpsc::Sender<Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue>>,
),
- EnqueueOrderDecision(
- AppSdkOrderDecisionRequest,
+ TradeDecision(
+ AppSdkTradeDecisionRequest,
mpsc::Sender<Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue>>,
),
- EnqueueOrderRevisionProposal(
- AppSdkOrderRevisionProposalRequest,
+ TradeRevisionProposal(
+ AppSdkTradeRevisionProposalRequest,
mpsc::Sender<Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue>>,
),
- EnqueueOrderRevisionDecision(
- AppSdkOrderRevisionDecisionRequest,
+ TradeRevisionDecision(
+ AppSdkTradeRevisionDecisionRequest,
mpsc::Sender<Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue>>,
),
- EnqueueOrderCancellation(
- AppSdkOrderCancellationRequest,
+ TradeCancellation(
+ AppSdkTradeCancellationRequest,
mpsc::Sender<Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue>>,
),
+ TradeStatus(
+ AppSdkTradeStatusRequest,
+ mpsc::Sender<Result<TradeStatusReceipt, AppSdkRuntimeIssue>>,
+ ),
+ TradeResync(
+ AppSdkTradeResyncRequest,
+ mpsc::Sender<Result<TradeResyncReceipt, AppSdkRuntimeIssue>>,
+ ),
BeginProjectionRebuild(
mpsc::Sender<Result<AppSdkProjectionLifecycleStatus, AppSdkRuntimeIssue>>,
),
@@ -445,15 +467,13 @@ impl fmt::Debug for AppSdkWorkerCommand {
Self::FarmPublicLocation(_, _) => formatter.write_str("FarmPublicLocation"),
Self::EnqueueFarmPublish(_, _) => formatter.write_str("EnqueueFarmPublish"),
Self::EnqueueListingPublish(_, _) => formatter.write_str("EnqueueListingPublish"),
- Self::EnqueueOrderSubmit(_, _) => formatter.write_str("EnqueueOrderSubmit"),
- Self::EnqueueOrderDecision(_, _) => formatter.write_str("EnqueueOrderDecision"),
- Self::EnqueueOrderRevisionProposal(_, _) => {
- formatter.write_str("EnqueueOrderRevisionProposal")
- }
- Self::EnqueueOrderRevisionDecision(_, _) => {
- formatter.write_str("EnqueueOrderRevisionDecision")
- }
- Self::EnqueueOrderCancellation(_, _) => formatter.write_str("EnqueueOrderCancellation"),
+ Self::TradePropose(_, _) => formatter.write_str("TradePropose"),
+ Self::TradeDecision(_, _) => formatter.write_str("TradeDecision"),
+ Self::TradeRevisionProposal(_, _) => formatter.write_str("TradeRevisionProposal"),
+ Self::TradeRevisionDecision(_, _) => formatter.write_str("TradeRevisionDecision"),
+ Self::TradeCancellation(_, _) => formatter.write_str("TradeCancellation"),
+ Self::TradeStatus(_, _) => formatter.write_str("TradeStatus"),
+ Self::TradeResync(_, _) => formatter.write_str("TradeResync"),
Self::BeginProjectionRebuild(_) => formatter.write_str("BeginProjectionRebuild"),
Self::CompleteProjectionRebuild(_) => formatter.write_str("CompleteProjectionRebuild"),
}
@@ -602,48 +622,66 @@ impl AppSdkRuntime {
})
}
- pub fn enqueue_order_submit(
+ pub fn trade_propose(
&self,
- request: AppSdkOrderSubmitRequest,
+ request: AppSdkTradeProposeRequest,
) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeError> {
self.run_command(|response_sender| {
- AppSdkWorkerCommand::EnqueueOrderSubmit(request, response_sender)
+ AppSdkWorkerCommand::TradePropose(request, response_sender)
})
}
- pub fn enqueue_order_decision(
+ pub fn trade_decide(
&self,
- request: AppSdkOrderDecisionRequest,
+ request: AppSdkTradeDecisionRequest,
) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeError> {
self.run_command(|response_sender| {
- AppSdkWorkerCommand::EnqueueOrderDecision(request, response_sender)
+ AppSdkWorkerCommand::TradeDecision(request, response_sender)
})
}
- pub fn enqueue_order_revision_proposal(
+ pub fn trade_revision_propose(
&self,
- request: AppSdkOrderRevisionProposalRequest,
+ request: AppSdkTradeRevisionProposalRequest,
) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeError> {
self.run_command(|response_sender| {
- AppSdkWorkerCommand::EnqueueOrderRevisionProposal(request, response_sender)
+ AppSdkWorkerCommand::TradeRevisionProposal(request, response_sender)
})
}
- pub fn enqueue_order_revision_decision(
+ pub fn trade_revision_decide(
&self,
- request: AppSdkOrderRevisionDecisionRequest,
+ request: AppSdkTradeRevisionDecisionRequest,
) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeError> {
self.run_command(|response_sender| {
- AppSdkWorkerCommand::EnqueueOrderRevisionDecision(request, response_sender)
+ AppSdkWorkerCommand::TradeRevisionDecision(request, response_sender)
})
}
- pub fn enqueue_order_cancellation(
+ pub fn trade_cancel(
&self,
- request: AppSdkOrderCancellationRequest,
+ request: AppSdkTradeCancellationRequest,
) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeError> {
self.run_command(|response_sender| {
- AppSdkWorkerCommand::EnqueueOrderCancellation(request, response_sender)
+ AppSdkWorkerCommand::TradeCancellation(request, response_sender)
+ })
+ }
+
+ pub fn trade_status(
+ &self,
+ request: AppSdkTradeStatusRequest,
+ ) -> Result<TradeStatusReceipt, AppSdkRuntimeError> {
+ self.run_command(|response_sender| {
+ AppSdkWorkerCommand::TradeStatus(request, response_sender)
+ })
+ }
+
+ pub fn trade_resync(
+ &self,
+ request: AppSdkTradeResyncRequest,
+ ) -> Result<TradeResyncReceipt, AppSdkRuntimeError> {
+ self.run_command(|response_sender| {
+ AppSdkWorkerCommand::TradeResync(request, response_sender)
})
}
@@ -1157,60 +1195,88 @@ fn run_app_sdk_worker(
};
send_worker_result(&shared, response_sender, result);
}
- AppSdkWorkerCommand::EnqueueOrderSubmit(request, response_sender) => {
+ AppSdkWorkerCommand::TradePropose(request, response_sender) => {
let result = if let Some(issue) = lifecycle_busy_issue(&shared) {
Err(issue)
} else {
match sdk.as_ref() {
- Some(sdk) => enqueue_order_submit_with_sdk(&runtime, sdk, request),
+ Some(_) => trade_propose_with_sdk(&runtime, &config, request),
None => Err(runtime_unavailable_issue(&shared)),
}
};
send_worker_result(&shared, response_sender, result);
}
- AppSdkWorkerCommand::EnqueueOrderDecision(request, response_sender) => {
+ AppSdkWorkerCommand::TradeDecision(request, response_sender) => {
let result = if let Some(issue) = lifecycle_busy_issue(&shared) {
Err(issue)
} else {
match sdk.as_ref() {
- Some(sdk) => enqueue_order_decision_with_sdk(&runtime, sdk, request),
+ Some(_) => trade_decision_with_sdk(&runtime, &config, request),
None => Err(runtime_unavailable_issue(&shared)),
}
};
send_worker_result(&shared, response_sender, result);
}
- AppSdkWorkerCommand::EnqueueOrderRevisionProposal(request, response_sender) => {
+ AppSdkWorkerCommand::TradeRevisionProposal(request, response_sender) => {
let result = if let Some(issue) = lifecycle_busy_issue(&shared) {
Err(issue)
} else {
match sdk.as_ref() {
- Some(sdk) => {
- enqueue_order_revision_proposal_with_sdk(&runtime, sdk, request)
- }
+ Some(_) => trade_revision_propose_with_sdk(&runtime, &config, request),
None => Err(runtime_unavailable_issue(&shared)),
}
};
send_worker_result(&shared, response_sender, result);
}
- AppSdkWorkerCommand::EnqueueOrderRevisionDecision(request, response_sender) => {
+ AppSdkWorkerCommand::TradeRevisionDecision(request, response_sender) => {
let result = if let Some(issue) = lifecycle_busy_issue(&shared) {
Err(issue)
} else {
match sdk.as_ref() {
- Some(sdk) => {
- enqueue_order_revision_decision_with_sdk(&runtime, sdk, request)
- }
+ Some(_) => trade_revision_decide_with_sdk(&runtime, &config, request),
None => Err(runtime_unavailable_issue(&shared)),
}
};
send_worker_result(&shared, response_sender, result);
}
- AppSdkWorkerCommand::EnqueueOrderCancellation(request, response_sender) => {
+ AppSdkWorkerCommand::TradeCancellation(request, response_sender) => {
let result = if let Some(issue) = lifecycle_busy_issue(&shared) {
Err(issue)
} else {
match sdk.as_ref() {
- Some(sdk) => enqueue_order_cancellation_with_sdk(&runtime, sdk, request),
+ Some(_) => trade_cancel_with_sdk(&runtime, &config, request),
+ None => Err(runtime_unavailable_issue(&shared)),
+ }
+ };
+ send_worker_result(&shared, response_sender, result);
+ }
+ AppSdkWorkerCommand::TradeStatus(request, response_sender) => {
+ let result = if let Some(issue) = lifecycle_busy_issue(&shared) {
+ Err(issue)
+ } else {
+ match sdk.as_ref() {
+ Some(sdk) => runtime
+ .block_on(
+ sdk.trades()
+ .status_client()
+ .status(TradeStatusRequest::new(request.locator)),
+ )
+ .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error)),
+ None => Err(runtime_unavailable_issue(&shared)),
+ }
+ };
+ send_worker_result(&shared, response_sender, result);
+ }
+ AppSdkWorkerCommand::TradeResync(request, response_sender) => {
+ let result = if let Some(issue) = lifecycle_busy_issue(&shared) {
+ Err(issue)
+ } else {
+ match sdk.as_ref() {
+ Some(sdk) => runtime
+ .block_on(sdk.trades().resync().resync(
+ TradeResyncRequest::new(request.locator).with_limit(request.limit),
+ ))
+ .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error)),
None => Err(runtime_unavailable_issue(&shared)),
}
};
@@ -1304,35 +1370,49 @@ fn run_degraded_worker(
Err(runtime_unavailable_issue(&shared)),
);
}
- AppSdkWorkerCommand::EnqueueOrderSubmit(_, response_sender) => {
+ AppSdkWorkerCommand::TradePropose(_, response_sender) => {
+ send_worker_result(
+ &shared,
+ response_sender,
+ Err(runtime_unavailable_issue(&shared)),
+ );
+ }
+ AppSdkWorkerCommand::TradeDecision(_, response_sender) => {
send_worker_result(
&shared,
response_sender,
Err(runtime_unavailable_issue(&shared)),
);
}
- AppSdkWorkerCommand::EnqueueOrderDecision(_, response_sender) => {
+ AppSdkWorkerCommand::TradeRevisionProposal(_, response_sender) => {
send_worker_result(
&shared,
response_sender,
Err(runtime_unavailable_issue(&shared)),
);
}
- AppSdkWorkerCommand::EnqueueOrderRevisionProposal(_, response_sender) => {
+ AppSdkWorkerCommand::TradeRevisionDecision(_, response_sender) => {
send_worker_result(
&shared,
response_sender,
Err(runtime_unavailable_issue(&shared)),
);
}
- AppSdkWorkerCommand::EnqueueOrderRevisionDecision(_, response_sender) => {
+ AppSdkWorkerCommand::TradeCancellation(_, response_sender) => {
send_worker_result(
&shared,
response_sender,
Err(runtime_unavailable_issue(&shared)),
);
}
- AppSdkWorkerCommand::EnqueueOrderCancellation(_, response_sender) => {
+ AppSdkWorkerCommand::TradeStatus(_, response_sender) => {
+ send_worker_result(
+ &shared,
+ response_sender,
+ Err(runtime_unavailable_issue(&shared)),
+ );
+ }
+ AppSdkWorkerCommand::TradeResync(_, response_sender) => {
send_worker_result(
&shared,
response_sender,
@@ -1373,6 +1453,41 @@ async fn build_sdk_runtime(config: &AppSdkConfig) -> Result<RadrootsClient, Radr
builder.build().await
}
+async fn build_sdk_runtime_with_signer(
+ config: &AppSdkConfig,
+ keys: RadrootsNostrKeys,
+) -> Result<RadrootsClient, AppSdkRuntimeIssue> {
+ let signer = RadrootsSdkLocalKeySigner::new(keys)
+ .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))?;
+ let mut builder = RadrootsClient::builder()
+ .directory_storage(config.storage_root.clone())
+ .relay_url_policy(config.relay_url_policy.into())
+ .signer_provider(RadrootsSdkSignerProvider::LocalKey(signer));
+ for relay_url in &config.relay_urls {
+ builder = builder.relay_url(relay_url.clone());
+ }
+ builder
+ .build()
+ .await
+ .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))
+}
+
+fn app_trade_publish_mode() -> PublishMode {
+ PublishMode::EnqueueOnly
+}
+
+fn app_trade_ack_policy() -> AckPolicy {
+ AckPolicy::NoWait
+}
+
+fn app_trade_relay_resolution_policy() -> RelayResolutionPolicy {
+ RelayResolutionPolicy::configured_relays()
+}
+
+fn app_trade_privacy_confirmation() -> PrivacyPreflightConfirmation {
+ PrivacyPreflightConfirmation::new().confirm(ProductSensitivityField::PublicButSensitiveNotes)
+}
+
fn run_restore_preflight(
runtime: &tokio::runtime::Runtime,
shared: &AppSdkRuntimeShared,
@@ -1500,202 +1615,220 @@ fn enqueue_listing_publish_with_sdk(
Ok(app_sdk_listing_receipt(receipt, request.actor_pubkey))
}
-fn enqueue_order_submit_with_sdk(
+fn trade_propose_with_sdk(
runtime: &tokio::runtime::Runtime,
- sdk: &RadrootsClient,
- request: AppSdkOrderSubmitRequest,
+ config: &AppSdkConfig,
+ request: AppSdkTradeProposeRequest,
) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue> {
let actor = sdk_actor_context(
request.actor_pubkey.as_str(),
request.actor_account_id.as_str(),
RadrootsActorRole::Buyer,
)?;
- let signer = sdk_local_signer(request.signer_keys)?;
- let target_relays = sdk_relay_targets(request.target_relays, request.relay_url_policy)?;
- let mut enqueue =
- OrderSubmitEnqueueRequest::new(actor, request.listing_event, request.order, target_relays);
+ let sdk = runtime.block_on(build_sdk_runtime_with_signer(config, request.signer_keys))?;
+ let mut sdk_request = TradeProposeRequest::new(
+ actor,
+ request.listing_event,
+ request.order,
+ app_trade_relay_resolution_policy(),
+ app_trade_publish_mode(),
+ app_trade_ack_policy(),
+ )
+ .with_privacy_confirmation(app_trade_privacy_confirmation());
if let Some(idempotency_key) = request.idempotency_key.as_deref() {
- enqueue = enqueue
+ sdk_request = sdk_request
.try_with_idempotency_key(idempotency_key)
.map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))?;
}
- let receipt = runtime
- .block_on(
- sdk.trades()
- .enqueue_submit_with_explicit_signer(enqueue, &signer),
- )
+ let outcome = runtime
+ .block_on(sdk.trades().buyer().propose_trade(sdk_request))
.map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))?;
- Ok(app_sdk_order_submit_ack(receipt, request.actor_pubkey))
+ app_sdk_trade_propose_receipt(outcome, request.actor_pubkey)
}
-fn enqueue_order_decision_with_sdk(
+fn trade_decision_with_sdk(
runtime: &tokio::runtime::Runtime,
- sdk: &RadrootsClient,
- request: AppSdkOrderDecisionRequest,
+ config: &AppSdkConfig,
+ request: AppSdkTradeDecisionRequest,
) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue> {
+ let evidence_events = request.evidence_events.clone();
let actor = sdk_actor_context(
request.actor_pubkey.as_str(),
request.actor_account_id.as_str(),
RadrootsActorRole::Seller,
)?;
- let signer = sdk_local_signer(request.signer_keys)?;
- let target_relays = sdk_relay_targets(request.target_relays, request.relay_url_policy)?;
- runtime
- .block_on(
- sdk.trades()
- .ingest_request_evidence(OrderRequestEvidenceIngestRequest::new(
- request.request_event,
- )),
- )
- .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))?;
- let mut enqueue = OrderDecisionEnqueueRequest::new(
- actor,
- request.request_event_ptr,
- request.decision,
- target_relays,
- );
- if let Some(idempotency_key) = request.idempotency_key.as_deref() {
- enqueue = enqueue
- .try_with_idempotency_key(idempotency_key)
- .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))?;
- }
- let receipt = runtime
- .block_on(
- sdk.trades()
- .enqueue_decision_with_explicit_signer(enqueue, &signer),
- )
- .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))?;
- Ok(app_sdk_order_decision_receipt(
- receipt,
- request.actor_pubkey,
- ))
+ let sdk = runtime.block_on(build_sdk_runtime_with_signer(config, request.signer_keys))?;
+ ingest_trade_evidence_events(runtime, &sdk, evidence_events)?;
+ let publish_mode = app_trade_publish_mode();
+ let ack_policy = app_trade_ack_policy();
+ let outcome = match request.decision {
+ AppSdkTradeDecision::Accept {
+ inventory_commitments,
+ } => {
+ let mut sdk_request = TradeAcceptRequest::new(
+ actor,
+ request.locator,
+ inventory_commitments,
+ app_trade_relay_resolution_policy(),
+ publish_mode,
+ ack_policy,
+ )
+ .with_privacy_confirmation(app_trade_privacy_confirmation());
+ if let Some(idempotency_key) = request.idempotency_key.as_deref() {
+ sdk_request = sdk_request
+ .try_with_idempotency_key(idempotency_key)
+ .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))?;
+ }
+ runtime
+ .block_on(sdk.trades().seller().accept_trade(sdk_request))
+ .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))?
+ }
+ AppSdkTradeDecision::Decline { reason } => {
+ let mut sdk_request = TradeDeclineRequest::new(
+ actor,
+ request.locator,
+ reason,
+ app_trade_relay_resolution_policy(),
+ publish_mode,
+ ack_policy,
+ )
+ .with_privacy_confirmation(app_trade_privacy_confirmation());
+ if let Some(idempotency_key) = request.idempotency_key.as_deref() {
+ sdk_request = sdk_request
+ .try_with_idempotency_key(idempotency_key)
+ .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))?;
+ }
+ runtime
+ .block_on(sdk.trades().seller().decline_trade(sdk_request))
+ .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))?
+ }
+ };
+ app_sdk_trade_decision_receipt(outcome, request.actor_pubkey)
}
-fn enqueue_order_revision_proposal_with_sdk(
+fn trade_revision_propose_with_sdk(
runtime: &tokio::runtime::Runtime,
- sdk: &RadrootsClient,
- request: AppSdkOrderRevisionProposalRequest,
+ config: &AppSdkConfig,
+ request: AppSdkTradeRevisionProposalRequest,
) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue> {
+ let evidence_events = request.evidence_events.clone();
let actor = sdk_actor_context(
request.actor_pubkey.as_str(),
request.actor_account_id.as_str(),
RadrootsActorRole::Seller,
)?;
- let signer = sdk_local_signer(request.signer_keys)?;
- let target_relays = sdk_relay_targets(request.target_relays, request.relay_url_policy)?;
- ingest_order_evidence_with_sdk(runtime, sdk, request.evidence_events)?;
- let mut enqueue = OrderRevisionProposalEnqueueRequest::new(
+ let sdk = runtime.block_on(build_sdk_runtime_with_signer(config, request.signer_keys))?;
+ ingest_trade_evidence_events(runtime, &sdk, evidence_events)?;
+ let mut sdk_request = TradeRevisionProposalRequest::new(
actor,
- request.root_event,
- request.previous_event,
- request.proposal,
- target_relays,
- );
+ request.locator,
+ request.revision_id,
+ request.items,
+ request.economics,
+ request.reason,
+ app_trade_relay_resolution_policy(),
+ app_trade_publish_mode(),
+ app_trade_ack_policy(),
+ )
+ .with_privacy_confirmation(app_trade_privacy_confirmation());
if let Some(idempotency_key) = request.idempotency_key.as_deref() {
- enqueue = enqueue
+ sdk_request = sdk_request
.try_with_idempotency_key(idempotency_key)
.map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))?;
}
- let receipt = runtime
- .block_on(
- sdk.trades()
- .enqueue_revision_proposal_with_explicit_signer(enqueue, &signer),
- )
+ let outcome = runtime
+ .block_on(sdk.trades().seller().propose_revision(sdk_request))
.map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))?;
- Ok(app_sdk_order_revision_proposal_receipt(
- receipt,
- request.actor_pubkey,
- ))
+ app_sdk_trade_revision_proposal_receipt(outcome, request.actor_pubkey)
}
-fn enqueue_order_revision_decision_with_sdk(
+fn trade_revision_decide_with_sdk(
runtime: &tokio::runtime::Runtime,
- sdk: &RadrootsClient,
- request: AppSdkOrderRevisionDecisionRequest,
+ config: &AppSdkConfig,
+ request: AppSdkTradeRevisionDecisionRequest,
) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue> {
+ let evidence_events = request.evidence_events.clone();
let actor = sdk_actor_context(
request.actor_pubkey.as_str(),
request.actor_account_id.as_str(),
RadrootsActorRole::Buyer,
)?;
- let signer = sdk_local_signer(request.signer_keys)?;
- let target_relays = sdk_relay_targets(request.target_relays, request.relay_url_policy)?;
- ingest_order_evidence_with_sdk(runtime, sdk, request.evidence_events)?;
- let mut enqueue = OrderRevisionDecisionEnqueueRequest::new(
+ let sdk = runtime.block_on(build_sdk_runtime_with_signer(config, request.signer_keys))?;
+ ingest_trade_evidence_events(runtime, &sdk, evidence_events)?;
+ let mut sdk_request = TradeRevisionDecisionRequest::new(
actor,
- request.root_event,
- request.previous_event,
- request.decision,
- target_relays,
- );
+ request.locator,
+ request.revision_id,
+ request.decision.clone(),
+ app_trade_relay_resolution_policy(),
+ app_trade_publish_mode(),
+ app_trade_ack_policy(),
+ )
+ .with_privacy_confirmation(app_trade_privacy_confirmation());
if let Some(idempotency_key) = request.idempotency_key.as_deref() {
- enqueue = enqueue
+ sdk_request = sdk_request
.try_with_idempotency_key(idempotency_key)
.map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))?;
}
- let receipt = runtime
- .block_on(
- sdk.trades()
- .enqueue_revision_decision_with_explicit_signer(enqueue, &signer),
- )
- .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))?;
- Ok(app_sdk_order_revision_decision_receipt(
- receipt,
- request.actor_pubkey,
- ))
+ let outcome = match request.decision {
+ RadrootsOrderRevisionOutcome::Accepted => runtime
+ .block_on(sdk.trades().buyer().accept_revision(sdk_request))
+ .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))?,
+ RadrootsOrderRevisionOutcome::Declined { .. } => runtime
+ .block_on(sdk.trades().buyer().decline_revision(sdk_request))
+ .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))?,
+ };
+ app_sdk_trade_revision_decision_receipt(outcome, request.actor_pubkey)
}
-fn enqueue_order_cancellation_with_sdk(
+fn trade_cancel_with_sdk(
runtime: &tokio::runtime::Runtime,
- sdk: &RadrootsClient,
- request: AppSdkOrderCancellationRequest,
+ config: &AppSdkConfig,
+ request: AppSdkTradeCancellationRequest,
) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue> {
+ let evidence_events = request.evidence_events.clone();
let actor = sdk_actor_context(
request.actor_pubkey.as_str(),
request.actor_account_id.as_str(),
RadrootsActorRole::Buyer,
)?;
- let signer = sdk_local_signer(request.signer_keys)?;
- let target_relays = sdk_relay_targets(request.target_relays, request.relay_url_policy)?;
- ingest_order_evidence_with_sdk(runtime, sdk, request.evidence_events)?;
- let mut enqueue = OrderCancellationEnqueueRequest::new(
+ let sdk = runtime.block_on(build_sdk_runtime_with_signer(config, request.signer_keys))?;
+ ingest_trade_evidence_events(runtime, &sdk, evidence_events)?;
+ let mut sdk_request = TradeCancelRequest::new(
actor,
- request.root_event,
- request.previous_event,
- request.cancellation,
- target_relays,
- );
+ request.locator,
+ request.reason,
+ app_trade_relay_resolution_policy(),
+ app_trade_publish_mode(),
+ app_trade_ack_policy(),
+ )
+ .with_privacy_confirmation(app_trade_privacy_confirmation());
if let Some(idempotency_key) = request.idempotency_key.as_deref() {
- enqueue = enqueue
+ sdk_request = sdk_request
.try_with_idempotency_key(idempotency_key)
.map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))?;
}
- let receipt = runtime
- .block_on(
- sdk.trades()
- .enqueue_cancellation_with_explicit_signer(enqueue, &signer),
- )
+ let outcome = runtime
+ .block_on(sdk.trades().buyer().cancel_trade(sdk_request))
.map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))?;
- Ok(app_sdk_order_cancellation_receipt(
- receipt,
- request.actor_pubkey,
- ))
+ app_sdk_trade_cancellation_receipt(outcome, request.actor_pubkey)
}
-fn ingest_order_evidence_with_sdk(
+fn ingest_trade_evidence_events(
runtime: &tokio::runtime::Runtime,
sdk: &RadrootsClient,
evidence_events: Vec<RadrootsNostrEvent>,
) -> Result<(), AppSdkRuntimeIssue> {
- for event in evidence_events {
- runtime
- .block_on(
+ runtime
+ .block_on(async {
+ for event in evidence_events {
sdk.trades()
- .ingest_evidence(OrderEvidenceIngestRequest::new(event)),
- )
- .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))?;
- }
- Ok(())
+ .ingest_evidence(TradeEvidenceIngestRequest::new(event))
+ .await?;
+ }
+ Ok::<(), RadrootsSdkError>(())
+ })
+ .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))
}
fn sdk_actor_context(
@@ -1756,84 +1889,115 @@ fn app_sdk_listing_receipt(
}
}
-fn app_sdk_order_submit_ack(
- receipt: OrderSubmitReceipt,
+fn app_sdk_trade_propose_receipt(
+ outcome: TradeMutationOutcome<TradeSubmitPlan, TradeSubmitReceipt>,
actor_pubkey: String,
-) -> AppSdkWorkflowReceipt {
- AppSdkWorkflowReceipt {
- operation_kind: ORDER_SUBMIT_OPERATION_KIND.to_owned(),
- expected_event_id: receipt.expected_event_id.as_str().to_owned(),
- signed_event_id: receipt.signed_event_id.as_str().to_owned(),
- outbox_operation_id: receipt.outbox_operation_id,
- outbox_event_id: receipt.outbox_event_id,
- state: sdk_mutation_state_key(receipt.state).to_owned(),
- idempotency_digest_prefix: receipt.idempotency_digest_prefix,
- actor_pubkey,
- }
-}
-
-fn app_sdk_order_decision_receipt(
- receipt: OrderDecisionReceipt,
+) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue> {
+ match outcome {
+ TradeMutationOutcome::Enqueued { receipt }
+ | TradeMutationOutcome::Published { receipt, .. } => Ok(AppSdkWorkflowReceipt {
+ operation_kind: TRADE_SUBMIT_OPERATION_KIND.to_owned(),
+ expected_event_id: receipt.expected_event_id.as_str().to_owned(),
+ signed_event_id: receipt.signed_event_id.as_str().to_owned(),
+ outbox_operation_id: receipt.outbox_operation_id,
+ outbox_event_id: receipt.outbox_event_id,
+ state: sdk_mutation_state_key(receipt.state).to_owned(),
+ idempotency_digest_prefix: receipt.idempotency_digest_prefix,
+ actor_pubkey,
+ }),
+ TradeMutationOutcome::DryRun { .. } => Err(unexpected_trade_dry_run_issue("trade.propose")),
+ }
+}
+
+fn app_sdk_trade_decision_receipt(
+ outcome: TradeMutationOutcome<TradeDecisionPlan, TradeDecisionReceipt>,
actor_pubkey: String,
-) -> AppSdkWorkflowReceipt {
- AppSdkWorkflowReceipt {
- operation_kind: ORDER_DECISION_OPERATION_KIND.to_owned(),
- expected_event_id: receipt.expected_event_id.as_str().to_owned(),
- signed_event_id: receipt.signed_event_id.as_str().to_owned(),
- outbox_operation_id: receipt.outbox_operation_id,
- outbox_event_id: receipt.outbox_event_id,
- state: sdk_mutation_state_key(receipt.state).to_owned(),
- idempotency_digest_prefix: receipt.idempotency_digest_prefix,
- actor_pubkey,
- }
-}
-
-fn app_sdk_order_revision_proposal_receipt(
- receipt: OrderRevisionProposalReceipt,
+) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue> {
+ match outcome {
+ TradeMutationOutcome::Enqueued { receipt }
+ | TradeMutationOutcome::Published { receipt, .. } => Ok(AppSdkWorkflowReceipt {
+ operation_kind: TRADE_DECISION_OPERATION_KIND.to_owned(),
+ expected_event_id: receipt.expected_event_id.as_str().to_owned(),
+ signed_event_id: receipt.signed_event_id.as_str().to_owned(),
+ outbox_operation_id: receipt.outbox_operation_id,
+ outbox_event_id: receipt.outbox_event_id,
+ state: sdk_mutation_state_key(receipt.state).to_owned(),
+ idempotency_digest_prefix: receipt.idempotency_digest_prefix,
+ actor_pubkey,
+ }),
+ TradeMutationOutcome::DryRun { .. } => Err(unexpected_trade_dry_run_issue("trade.decide")),
+ }
+}
+
+fn app_sdk_trade_revision_proposal_receipt(
+ outcome: TradeMutationOutcome<TradeRevisionProposalPlan, TradeRevisionProposalReceipt>,
actor_pubkey: String,
-) -> AppSdkWorkflowReceipt {
- AppSdkWorkflowReceipt {
- operation_kind: ORDER_REVISION_PROPOSAL_OPERATION_KIND.to_owned(),
- expected_event_id: receipt.expected_event_id.as_str().to_owned(),
- signed_event_id: receipt.signed_event_id.as_str().to_owned(),
- outbox_operation_id: receipt.outbox_operation_id,
- outbox_event_id: receipt.outbox_event_id,
- state: sdk_mutation_state_key(receipt.state).to_owned(),
- idempotency_digest_prefix: receipt.idempotency_digest_prefix,
- actor_pubkey,
+) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue> {
+ match outcome {
+ TradeMutationOutcome::Enqueued { receipt }
+ | TradeMutationOutcome::Published { receipt, .. } => Ok(AppSdkWorkflowReceipt {
+ operation_kind: TRADE_REVISION_PROPOSAL_OPERATION_KIND.to_owned(),
+ expected_event_id: receipt.expected_event_id.as_str().to_owned(),
+ signed_event_id: receipt.signed_event_id.as_str().to_owned(),
+ outbox_operation_id: receipt.outbox_operation_id,
+ outbox_event_id: receipt.outbox_event_id,
+ state: sdk_mutation_state_key(receipt.state).to_owned(),
+ idempotency_digest_prefix: receipt.idempotency_digest_prefix,
+ actor_pubkey,
+ }),
+ TradeMutationOutcome::DryRun { .. } => {
+ Err(unexpected_trade_dry_run_issue("trade.revision.propose"))
+ }
}
}
-fn app_sdk_order_revision_decision_receipt(
- receipt: OrderRevisionDecisionReceipt,
+fn app_sdk_trade_revision_decision_receipt(
+ outcome: TradeMutationOutcome<TradeRevisionDecisionPlan, TradeRevisionDecisionReceipt>,
actor_pubkey: String,
-) -> AppSdkWorkflowReceipt {
- AppSdkWorkflowReceipt {
- operation_kind: ORDER_REVISION_DECISION_OPERATION_KIND.to_owned(),
- expected_event_id: receipt.expected_event_id.as_str().to_owned(),
- signed_event_id: receipt.signed_event_id.as_str().to_owned(),
- outbox_operation_id: receipt.outbox_operation_id,
- outbox_event_id: receipt.outbox_event_id,
- state: sdk_mutation_state_key(receipt.state).to_owned(),
- idempotency_digest_prefix: receipt.idempotency_digest_prefix,
- actor_pubkey,
+) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue> {
+ match outcome {
+ TradeMutationOutcome::Enqueued { receipt }
+ | TradeMutationOutcome::Published { receipt, .. } => Ok(AppSdkWorkflowReceipt {
+ operation_kind: TRADE_REVISION_DECISION_OPERATION_KIND.to_owned(),
+ expected_event_id: receipt.expected_event_id.as_str().to_owned(),
+ signed_event_id: receipt.signed_event_id.as_str().to_owned(),
+ outbox_operation_id: receipt.outbox_operation_id,
+ outbox_event_id: receipt.outbox_event_id,
+ state: sdk_mutation_state_key(receipt.state).to_owned(),
+ idempotency_digest_prefix: receipt.idempotency_digest_prefix,
+ actor_pubkey,
+ }),
+ TradeMutationOutcome::DryRun { .. } => {
+ Err(unexpected_trade_dry_run_issue("trade.revision.decide"))
+ }
}
}
-fn app_sdk_order_cancellation_receipt(
- receipt: OrderCancellationReceipt,
+fn app_sdk_trade_cancellation_receipt(
+ outcome: TradeMutationOutcome<TradeCancellationPlan, TradeCancellationReceipt>,
actor_pubkey: String,
-) -> AppSdkWorkflowReceipt {
- AppSdkWorkflowReceipt {
- operation_kind: ORDER_CANCELLATION_OPERATION_KIND.to_owned(),
- expected_event_id: receipt.expected_event_id.as_str().to_owned(),
- signed_event_id: receipt.signed_event_id.as_str().to_owned(),
- outbox_operation_id: receipt.outbox_operation_id,
- outbox_event_id: receipt.outbox_event_id,
- state: sdk_mutation_state_key(receipt.state).to_owned(),
- idempotency_digest_prefix: receipt.idempotency_digest_prefix,
- actor_pubkey,
- }
+) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue> {
+ match outcome {
+ TradeMutationOutcome::Enqueued { receipt }
+ | TradeMutationOutcome::Published { receipt, .. } => Ok(AppSdkWorkflowReceipt {
+ operation_kind: TRADE_CANCELLATION_OPERATION_KIND.to_owned(),
+ expected_event_id: receipt.expected_event_id.as_str().to_owned(),
+ signed_event_id: receipt.signed_event_id.as_str().to_owned(),
+ outbox_operation_id: receipt.outbox_operation_id,
+ outbox_event_id: receipt.outbox_event_id,
+ state: sdk_mutation_state_key(receipt.state).to_owned(),
+ idempotency_digest_prefix: receipt.idempotency_digest_prefix,
+ actor_pubkey,
+ }),
+ TradeMutationOutcome::DryRun { .. } => Err(unexpected_trade_dry_run_issue("trade.cancel")),
+ }
+}
+
+fn unexpected_trade_dry_run_issue(operation: &'static str) -> AppSdkRuntimeIssue {
+ AppSdkRuntimeIssue::runtime_error(
+ "sdk_trade_unexpected_dry_run",
+ format!("{operation} returned a dry-run plan for an enqueue-only Studio command"),
+ )
}
fn sdk_mutation_state_key(state: SdkMutationState) -> &'static str {
diff --git a/crates/store/migrations/0028_sdk_migration_receipts.sql b/crates/store/migrations/0028_sdk_migration_receipts.sql
@@ -1,38 +0,0 @@
-CREATE TABLE app_sdk_migration_receipts (
- id TEXT PRIMARY KEY NOT NULL,
- source_kind TEXT NOT NULL CHECK (
- source_kind IN ('local_outbox', 'shared_local_event')
- ),
- source_record_id TEXT NOT NULL,
- sdk_operation_kind TEXT NOT NULL,
- sdk_outbox_event_ids_json TEXT NOT NULL,
- expected_event_id TEXT,
- actor_pubkey TEXT,
- idempotency_digest_prefix TEXT,
- migration_state TEXT NOT NULL CHECK (
- migration_state IN (
- 'pending',
- 'prepared',
- 'enqueued',
- 'pushed',
- 'failed',
- 'blocked',
- 'skipped',
- 'unsupported',
- 'manual_review',
- 'unknown'
- )
- ),
- created_at TEXT NOT NULL,
- updated_at TEXT NOT NULL,
- detail_json TEXT NOT NULL,
- UNIQUE(source_kind, source_record_id)
-);
-
-CREATE INDEX idx_app_sdk_migration_receipts_source_record ON app_sdk_migration_receipts(
- source_record_id
-);
-CREATE INDEX idx_app_sdk_migration_receipts_state ON app_sdk_migration_receipts(
- migration_state,
- updated_at
-);
diff --git a/crates/store/migrations/0028_sdk_workflow_receipts.sql b/crates/store/migrations/0028_sdk_workflow_receipts.sql
@@ -0,0 +1,38 @@
+CREATE TABLE app_sdk_workflow_receipts (
+ id TEXT PRIMARY KEY NOT NULL,
+ source_kind TEXT NOT NULL CHECK (
+ source_kind IN ('local_outbox', 'shared_local_event')
+ ),
+ source_record_id TEXT NOT NULL,
+ sdk_operation_kind TEXT NOT NULL,
+ sdk_outbox_event_ids_json TEXT NOT NULL,
+ expected_event_id TEXT,
+ actor_pubkey TEXT,
+ idempotency_digest_prefix TEXT,
+ workflow_state TEXT NOT NULL CHECK (
+ workflow_state IN (
+ 'pending',
+ 'prepared',
+ 'enqueued',
+ 'pushed',
+ 'failed',
+ 'blocked',
+ 'skipped',
+ 'unsupported',
+ 'manual_review',
+ 'unknown'
+ )
+ ),
+ created_at TEXT NOT NULL,
+ updated_at TEXT NOT NULL,
+ detail_json TEXT NOT NULL,
+ UNIQUE(source_kind, source_record_id)
+);
+
+CREATE INDEX idx_app_sdk_workflow_receipts_source_record ON app_sdk_workflow_receipts(
+ source_record_id
+);
+CREATE INDEX idx_app_sdk_workflow_receipts_state ON app_sdk_workflow_receipts(
+ workflow_state,
+ updated_at
+);
diff --git a/crates/store/src/lib.rs b/crates/store/src/lib.rs
@@ -2,10 +2,9 @@
mod error;
mod interop;
-mod migration_audit;
mod migrations;
mod repo;
-mod sdk_migration_receipts;
+mod sdk_workflow_receipts;
mod sync;
use std::{collections::BTreeSet, fs, path::PathBuf, time::Duration};
@@ -32,12 +31,6 @@ pub use interop::{
AppLocalInteropImportReport, AppLocalInteropRepository, StoredLocalInteropRecord,
projected_order_id_from_trade_request,
};
-pub use migration_audit::{
- APP_SDK_MIGRATION_AUDIT_DEFAULT_BATCH_SIZE, APP_SDK_MIGRATION_AUDIT_MAX_BATCH_SIZE,
- AppSdkMigrationAuditClassification, AppSdkMigrationAuditCount,
- AppSdkMigrationAuditDuplicateCandidate, AppSdkMigrationAuditIssue, AppSdkMigrationAuditReport,
- AppSdkMigrationAuditRequest, AppSdkMigrationAuditSource, AppSdkMigrationAuditSourceReport,
-};
pub use migrations::latest_schema_version;
pub use repo::{
APP_ACTIVITY_CONTEXT_LIMIT, APP_ACTIVITY_RETENTION_LIMIT, AppActivationRepository,
@@ -48,9 +41,9 @@ pub use repo::{
SellerOrderDecisionExport, SellerOrderDecisionLineExport, TODAY_AGENDA_LIST_LIMIT,
TODAY_AGENDA_LOW_STOCK_THRESHOLD, derive_farm_rules_readiness,
};
-pub use sdk_migration_receipts::{
- AppSdkMigrationReceipt, AppSdkMigrationReceiptInput, AppSdkMigrationReceiptRepository,
- AppSdkMigrationReceiptSourceKind, AppSdkMigrationState,
+pub use sdk_workflow_receipts::{
+ AppSdkStoredWorkflowReceipt, AppSdkWorkflowReceiptInput, AppSdkWorkflowReceiptRepository,
+ AppSdkWorkflowReceiptSourceKind, AppSdkWorkflowReceiptState,
};
pub use sync::{
AppSyncRepository, StoredPendingSyncOperation, StoredRelayIngestCursor, StoredSyncConflict,
@@ -124,8 +117,8 @@ impl AppSqliteStore {
AppSyncRepository::new(&self.connection)
}
- pub fn sdk_migration_receipt_repository(&self) -> AppSdkMigrationReceiptRepository<'_> {
- AppSdkMigrationReceiptRepository::new(&self.connection)
+ pub fn sdk_workflow_receipt_repository(&self) -> AppSdkWorkflowReceiptRepository<'_> {
+ AppSdkWorkflowReceiptRepository::new(&self.connection)
}
pub fn reminders_repository(&self) -> AppRemindersRepository<'_> {
@@ -837,7 +830,7 @@ mod tests {
assert!(table_exists(connection, "reminder_log_entries"));
assert!(table_exists(connection, "buyer_order_coordination_records"));
assert!(table_exists(connection, "order_validation_receipts"));
- assert!(table_exists(connection, "app_sdk_migration_receipts"));
+ assert!(table_exists(connection, "app_sdk_workflow_receipts"));
assert!(column_exists(connection, "farms", "timezone"));
assert!(column_exists(connection, "farms", "currency_code"));
assert!(column_exists(connection, "local_outbox", "account_id"));
@@ -1005,47 +998,47 @@ mod tests {
));
assert!(column_exists(
connection,
- "app_sdk_migration_receipts",
+ "app_sdk_workflow_receipts",
"source_record_id"
));
assert!(column_exists(
connection,
- "app_sdk_migration_receipts",
+ "app_sdk_workflow_receipts",
"source_kind"
));
assert!(column_exists(
connection,
- "app_sdk_migration_receipts",
+ "app_sdk_workflow_receipts",
"sdk_operation_kind"
));
assert!(column_exists(
connection,
- "app_sdk_migration_receipts",
+ "app_sdk_workflow_receipts",
"sdk_outbox_event_ids_json"
));
assert!(column_exists(
connection,
- "app_sdk_migration_receipts",
+ "app_sdk_workflow_receipts",
"expected_event_id"
));
assert!(column_exists(
connection,
- "app_sdk_migration_receipts",
+ "app_sdk_workflow_receipts",
"actor_pubkey"
));
assert!(column_exists(
connection,
- "app_sdk_migration_receipts",
+ "app_sdk_workflow_receipts",
"idempotency_digest_prefix"
));
assert!(column_exists(
connection,
- "app_sdk_migration_receipts",
- "migration_state"
+ "app_sdk_workflow_receipts",
+ "workflow_state"
));
assert!(column_exists(
connection,
- "app_sdk_migration_receipts",
+ "app_sdk_workflow_receipts",
"detail_json"
));
connection
diff --git a/crates/store/src/migration_audit.rs b/crates/store/src/migration_audit.rs
@@ -1,1555 +0,0 @@
-use std::collections::{BTreeMap, BTreeSet};
-
-use radroots_events::kinds::{
- KIND_FARM, KIND_LISTING, KIND_LISTING_DRAFT, KIND_ORDER_CANCELLATION, KIND_ORDER_DECISION,
- KIND_ORDER_REQUEST, KIND_ORDER_REVISION_DECISION, KIND_ORDER_REVISION_PROPOSAL,
- KIND_TRADE_VALIDATION_RECEIPT,
-};
-use radroots_local_events::{
- LocalEventRecord, LocalEventsStore, LocalRecordFamily, LocalRecordStatus, PublishOutboxStatus,
-};
-use radroots_sql_core::SqlExecutor;
-use radroots_studio_app_sync::{AppPublishPayload, SyncOperationKind};
-use rusqlite::params;
-use serde_json::Value;
-
-use crate::{
- AppSdkMigrationReceipt, AppSdkMigrationReceiptSourceKind, AppSdkMigrationState, AppSqliteError,
- AppSqliteStore,
-};
-
-pub const APP_SDK_MIGRATION_AUDIT_DEFAULT_BATCH_SIZE: u32 = 500;
-pub const APP_SDK_MIGRATION_AUDIT_MAX_BATCH_SIZE: u32 = 1_000;
-
-#[derive(Clone, Copy, Debug, Eq, PartialEq)]
-pub struct AppSdkMigrationAuditRequest {
- pub batch_size: u32,
-}
-
-impl Default for AppSdkMigrationAuditRequest {
- fn default() -> Self {
- Self {
- batch_size: APP_SDK_MIGRATION_AUDIT_DEFAULT_BATCH_SIZE,
- }
- }
-}
-
-impl AppSdkMigrationAuditRequest {
- pub fn normalized_batch_size(self) -> u32 {
- if self.batch_size == 0 {
- APP_SDK_MIGRATION_AUDIT_DEFAULT_BATCH_SIZE
- } else {
- self.batch_size.min(APP_SDK_MIGRATION_AUDIT_MAX_BATCH_SIZE)
- }
- }
-}
-
-#[derive(Clone, Debug, Eq, PartialEq)]
-pub struct AppSdkMigrationAuditReport {
- pub local_outbox: AppSdkMigrationAuditSourceReport,
- pub shared_local_events: AppSdkMigrationAuditSourceReport,
- pub issues: Vec<AppSdkMigrationAuditIssue>,
-}
-
-#[derive(Clone, Debug, Eq, PartialEq)]
-pub struct AppSdkMigrationAuditSourceReport {
- pub source: AppSdkMigrationAuditSource,
- pub batch_size: u32,
- pub batch_count: u64,
- pub scanned_records: u64,
- pub kind_counts: Vec<AppSdkMigrationAuditCount>,
- pub status_counts: Vec<AppSdkMigrationAuditCount>,
- pub classification_counts: Vec<AppSdkMigrationAuditCount>,
- pub duplicate_candidates: Vec<AppSdkMigrationAuditDuplicateCandidate>,
- pub issues: Vec<AppSdkMigrationAuditIssue>,
-}
-
-#[derive(Clone, Debug, Eq, PartialEq)]
-pub struct AppSdkMigrationAuditCount {
- pub key: String,
- pub count: u64,
-}
-
-#[derive(Clone, Debug, Eq, PartialEq)]
-pub struct AppSdkMigrationAuditDuplicateCandidate {
- pub identity_kind: String,
- pub identity_key: String,
- pub record_count: u64,
- pub record_ids: Vec<String>,
-}
-
-#[derive(Clone, Debug, Eq, PartialEq)]
-pub struct AppSdkMigrationAuditIssue {
- pub source: AppSdkMigrationAuditSource,
- pub code: String,
- pub record_id: Option<String>,
- pub message: String,
-}
-
-#[derive(Clone, Copy, Debug, Eq, PartialEq)]
-pub enum AppSdkMigrationAuditSource {
- LocalOutbox,
- SharedLocalEvents,
-}
-
-#[derive(Clone, Copy, Debug, Eq, PartialEq)]
-pub enum AppSdkMigrationAuditClassification {
- PublishableCandidate,
- AlreadyRepresentedCandidate,
- RepresentedRecord,
- SkippedRecord,
- FailedRecord,
- LocalWorkDeferred,
- ManualReviewRequired,
- ValidationReceiptDeferred,
- Unsupported,
- Unknown,
-}
-
-impl AppSdkMigrationAuditClassification {
- pub const fn storage_key(self) -> &'static str {
- match self {
- Self::PublishableCandidate => "publishable_candidate",
- Self::AlreadyRepresentedCandidate => "already_represented_candidate",
- Self::RepresentedRecord => "represented_record",
- Self::SkippedRecord => "skipped_record",
- Self::FailedRecord => "failed_record",
- Self::LocalWorkDeferred => "local_work_deferred",
- Self::ManualReviewRequired => "manual_review_required",
- Self::ValidationReceiptDeferred => "validation_receipt_deferred",
- Self::Unsupported => "unsupported",
- Self::Unknown => "unknown",
- }
- }
-}
-
-impl AppSqliteStore {
- pub fn audit_sdk_migration<E>(
- &self,
- shared_local_events: &LocalEventsStore<E>,
- request: AppSdkMigrationAuditRequest,
- ) -> Result<AppSdkMigrationAuditReport, AppSqliteError>
- where
- E: SqlExecutor,
- {
- let local_outbox = self.audit_sdk_migration_local_outbox(request)?;
- let shared_local_events =
- self.audit_sdk_migration_shared_local_events(shared_local_events, request)?;
- let issues = local_outbox
- .issues
- .iter()
- .chain(shared_local_events.issues.iter())
- .cloned()
- .collect();
-
- Ok(AppSdkMigrationAuditReport {
- local_outbox,
- shared_local_events,
- issues,
- })
- }
-
- pub fn audit_sdk_migration_local_outbox(
- &self,
- request: AppSdkMigrationAuditRequest,
- ) -> Result<AppSdkMigrationAuditSourceReport, AppSqliteError> {
- let batch_size = request.normalized_batch_size();
- let mut report = AppSdkMigrationAuditSourceBuilder::new(
- AppSdkMigrationAuditSource::LocalOutbox,
- batch_size,
- );
- let mut last_rowid = 0_i64;
-
- loop {
- let rows = self.load_local_outbox_audit_batch(last_rowid, batch_size)?;
- if rows.is_empty() {
- break;
- }
- report.batch_count += 1;
- for row in &rows {
- last_rowid = row.rowid;
- let receipt = self.sdk_migration_receipt_repository().load_receipt(
- AppSdkMigrationReceiptSourceKind::LocalOutbox,
- row.id.as_str(),
- )?;
- audit_local_outbox_row(row, receipt.as_ref(), &mut report);
- }
- if rows.len() < batch_size as usize {
- break;
- }
- }
-
- Ok(report.finish())
- }
-
- pub fn audit_sdk_migration_shared_local_events<E>(
- &self,
- store: &LocalEventsStore<E>,
- request: AppSdkMigrationAuditRequest,
- ) -> Result<AppSdkMigrationAuditSourceReport, AppSqliteError>
- where
- E: SqlExecutor,
- {
- audit_sdk_migration_shared_local_events_with_receipts(store, request, |record_id| {
- self.sdk_migration_receipt_repository().load_receipt(
- AppSdkMigrationReceiptSourceKind::SharedLocalEvent,
- record_id,
- )
- })
- }
-}
-
-fn audit_sdk_migration_shared_local_events_with_receipts<E>(
- store: &LocalEventsStore<E>,
- request: AppSdkMigrationAuditRequest,
- mut load_receipt: impl FnMut(&str) -> Result<Option<AppSdkMigrationReceipt>, AppSqliteError>,
-) -> Result<AppSdkMigrationAuditSourceReport, AppSqliteError>
-where
- E: SqlExecutor,
-{
- let batch_size = request.normalized_batch_size();
- let mut report = AppSdkMigrationAuditSourceBuilder::new(
- AppSdkMigrationAuditSource::SharedLocalEvents,
- batch_size,
- );
- let mut after_change_seq = 0_i64;
-
- loop {
- let records = store
- .list_records_changed_after(after_change_seq, batch_size)
- .map_err(|source| AppSqliteError::LocalEvents {
- operation: "audit shared local event records",
- source,
- })?;
- if records.is_empty() {
- break;
- }
- report.batch_count += 1;
- for record in &records {
- after_change_seq = record.change_seq;
- let receipt = load_receipt(record.record_id.as_str())?;
- audit_shared_local_event_record(record, receipt.as_ref(), &mut report);
- }
- if records.len() < batch_size as usize {
- break;
- }
- }
-
- Ok(report.finish())
-}
-
-impl AppSqliteStore {
- fn load_local_outbox_audit_batch(
- &self,
- after_rowid: i64,
- limit: u32,
- ) -> Result<Vec<LocalOutboxAuditRow>, AppSqliteError> {
- let mut statement = self
- .connection()
- .prepare(
- "SELECT
- rowid,
- id,
- account_id,
- operation_key,
- aggregate_kind,
- aggregate_id,
- operation_kind,
- payload_json,
- state
- FROM local_outbox
- WHERE rowid > ?1
- ORDER BY rowid ASC
- LIMIT ?2",
- )
- .map_err(|source| AppSqliteError::Query {
- operation: "prepare SDK migration local outbox audit query",
- source,
- })?;
- let rows = statement
- .query_map(params![after_rowid, i64::from(limit)], |row| {
- Ok(LocalOutboxAuditRow {
- rowid: row.get(0)?,
- id: row.get(1)?,
- account_id: row.get(2)?,
- operation_key: row.get(3)?,
- aggregate_kind: row.get(4)?,
- aggregate_id: row.get(5)?,
- operation_kind: row.get(6)?,
- payload_json: row.get(7)?,
- state: row.get(8)?,
- })
- })
- .map_err(|source| AppSqliteError::Query {
- operation: "query SDK migration local outbox audit rows",
- source,
- })?;
-
- rows.map(|row| {
- row.map_err(|source| AppSqliteError::Query {
- operation: "read SDK migration local outbox audit row",
- source,
- })
- })
- .collect()
- }
-}
-
-#[derive(Clone, Debug, Eq, PartialEq)]
-struct LocalOutboxAuditRow {
- rowid: i64,
- id: String,
- account_id: String,
- operation_key: String,
- aggregate_kind: String,
- aggregate_id: String,
- operation_kind: String,
- payload_json: String,
- state: String,
-}
-
-struct AppSdkMigrationAuditSourceBuilder {
- source: AppSdkMigrationAuditSource,
- batch_size: u32,
- batch_count: u64,
- scanned_records: u64,
- kind_counts: BTreeMap<String, u64>,
- status_counts: BTreeMap<String, u64>,
- classification_counts: BTreeMap<String, u64>,
- duplicate_records: BTreeMap<DuplicateIdentity, BTreeSet<String>>,
- issues: Vec<AppSdkMigrationAuditIssue>,
-}
-
-impl AppSdkMigrationAuditSourceBuilder {
- fn new(source: AppSdkMigrationAuditSource, batch_size: u32) -> Self {
- Self {
- source,
- batch_size,
- batch_count: 0,
- scanned_records: 0,
- kind_counts: BTreeMap::new(),
- status_counts: BTreeMap::new(),
- classification_counts: BTreeMap::new(),
- duplicate_records: BTreeMap::new(),
- issues: Vec::new(),
- }
- }
-
- fn record(
- &mut self,
- record_id: &str,
- kind: String,
- status: String,
- classification: AppSdkMigrationAuditClassification,
- duplicate_identities: Vec<DuplicateIdentity>,
- ) {
- self.scanned_records += 1;
- increment_count(&mut self.kind_counts, kind);
- increment_count(&mut self.status_counts, status);
- increment_count(
- &mut self.classification_counts,
- classification.storage_key().to_owned(),
- );
- for identity in duplicate_identities {
- self.duplicate_records
- .entry(identity)
- .or_default()
- .insert(record_id.to_owned());
- }
- }
-
- fn issue(&mut self, code: &str, record_id: Option<&str>, message: impl Into<String>) {
- self.issues.push(AppSdkMigrationAuditIssue {
- source: self.source,
- code: code.to_owned(),
- record_id: record_id.map(ToOwned::to_owned),
- message: message.into(),
- });
- }
-
- fn finish(self) -> AppSdkMigrationAuditSourceReport {
- AppSdkMigrationAuditSourceReport {
- source: self.source,
- batch_size: self.batch_size,
- batch_count: self.batch_count,
- scanned_records: self.scanned_records,
- kind_counts: counts_from_map(self.kind_counts),
- status_counts: counts_from_map(self.status_counts),
- classification_counts: counts_from_map(self.classification_counts),
- duplicate_candidates: duplicate_candidates_from_map(self.duplicate_records),
- issues: self.issues,
- }
- }
-}
-
-#[derive(Clone, Debug, Eq, Ord, PartialEq, PartialOrd)]
-struct DuplicateIdentity {
- kind: String,
- key: String,
-}
-
-fn audit_local_outbox_row(
- row: &LocalOutboxAuditRow,
- receipt: Option<&AppSdkMigrationReceipt>,
- report: &mut AppSdkMigrationAuditSourceBuilder,
-) {
- let payload = serde_json::from_str::<AppPublishPayload>(row.payload_json.as_str());
- let (kind, source_classification) = match payload {
- Ok(payload) => {
- if row.operation_kind == SyncOperationKind::Delete.storage_key() {
- report.issue(
- "unsupported_local_outbox_operation",
- Some(row.id.as_str()),
- format!(
- "local outbox delete operation `{}` is not a SDK publish migration candidate",
- row.operation_key
- ),
- );
- (
- payload.work_kind().storage_key().to_owned(),
- AppSdkMigrationAuditClassification::Unsupported,
- )
- } else {
- (
- payload.work_kind().storage_key().to_owned(),
- classify_local_outbox_state(row, report),
- )
- }
- }
- Err(source) => {
- report.issue(
- "unknown_local_outbox_payload",
- Some(row.id.as_str()),
- format!(
- "local outbox payload for operation `{}` could not be decoded: {source}",
- row.operation_key
- ),
- );
- (
- format!("{}:{}", row.aggregate_kind, row.operation_kind),
- AppSdkMigrationAuditClassification::Unknown,
- )
- }
- };
- let classification =
- classify_receipt_overlay(row.id.as_str(), source_classification, receipt, report);
- let identities = vec![
- DuplicateIdentity {
- kind: "operation".to_owned(),
- key: format!("{}:{}", row.account_id, row.operation_key),
- },
- DuplicateIdentity {
- kind: "aggregate".to_owned(),
- key: format!(
- "{}:{}:{}:{}",
- row.account_id, row.aggregate_kind, row.aggregate_id, row.operation_kind
- ),
- },
- ];
- report.record(
- row.id.as_str(),
- kind,
- row.state.clone(),
- classification,
- identities,
- );
-}
-
-fn classify_local_outbox_state(
- row: &LocalOutboxAuditRow,
- report: &mut AppSdkMigrationAuditSourceBuilder,
-) -> AppSdkMigrationAuditClassification {
- match row.state.as_str() {
- "pending" | "in_progress" | "retryable" => {
- AppSdkMigrationAuditClassification::PublishableCandidate
- }
- "succeeded" => AppSdkMigrationAuditClassification::AlreadyRepresentedCandidate,
- "failed" | "blocked" => {
- report.issue(
- "manual_review_local_outbox_state",
- Some(row.id.as_str()),
- format!(
- "local outbox operation `{}` is in `{}` state and requires migration review",
- row.operation_key, row.state
- ),
- );
- AppSdkMigrationAuditClassification::ManualReviewRequired
- }
- _ => {
- report.issue(
- "unknown_local_outbox_state",
- Some(row.id.as_str()),
- format!(
- "local outbox operation `{}` has unknown state `{}`",
- row.operation_key, row.state
- ),
- );
- AppSdkMigrationAuditClassification::Unknown
- }
- }
-}
-
-fn audit_shared_local_event_record(
- record: &LocalEventRecord,
- receipt: Option<&AppSdkMigrationReceipt>,
- report: &mut AppSdkMigrationAuditSourceBuilder,
-) {
- let kind = shared_local_event_kind(record);
- let source_classification = shared_local_event_classification(record, report);
- let classification = classify_receipt_overlay(
- record.record_id.as_str(),
- source_classification,
- receipt,
- report,
- );
- report.record(
- record.record_id.as_str(),
- kind,
- format!(
- "{}:{}",
- record.status.as_str(),
- record.outbox_status.as_str()
- ),
- classification,
- shared_local_event_duplicate_identities(record),
- );
-}
-
-fn classify_receipt_overlay(
- record_id: &str,
- source_classification: AppSdkMigrationAuditClassification,
- receipt: Option<&AppSdkMigrationReceipt>,
- report: &mut AppSdkMigrationAuditSourceBuilder,
-) -> AppSdkMigrationAuditClassification {
- let Some(receipt) = receipt else {
- return source_classification;
- };
- if !receipt_allowed_for_source_classification(source_classification) {
- report.issue(
- "sdk_migration_receipt_for_non_migratable_source",
- Some(record_id),
- format!(
- "SDK migration receipt `{}` for operation `{}` cannot override source classification `{}`",
- receipt.id,
- receipt.sdk_operation_kind,
- source_classification.storage_key()
- ),
- );
- return source_classification;
- }
-
- match receipt.migration_state {
- AppSdkMigrationState::Pending | AppSdkMigrationState::Prepared => source_classification,
- AppSdkMigrationState::Enqueued | AppSdkMigrationState::Pushed => {
- AppSdkMigrationAuditClassification::RepresentedRecord
- }
- AppSdkMigrationState::Skipped => AppSdkMigrationAuditClassification::SkippedRecord,
- AppSdkMigrationState::Failed => {
- report.issue(
- "sdk_migration_receipt_failed",
- Some(record_id),
- format!(
- "SDK migration receipt `{}` for operation `{}` is failed",
- receipt.id, receipt.sdk_operation_kind
- ),
- );
- AppSdkMigrationAuditClassification::FailedRecord
- }
- AppSdkMigrationState::Blocked | AppSdkMigrationState::ManualReview => {
- report.issue(
- "sdk_migration_receipt_manual_review",
- Some(record_id),
- format!(
- "SDK migration receipt `{}` for operation `{}` requires manual review",
- receipt.id, receipt.sdk_operation_kind
- ),
- );
- AppSdkMigrationAuditClassification::ManualReviewRequired
- }
- AppSdkMigrationState::Unsupported => {
- report.issue(
- "sdk_migration_receipt_unsupported",
- Some(record_id),
- format!(
- "SDK migration receipt `{}` for operation `{}` is unsupported",
- receipt.id, receipt.sdk_operation_kind
- ),
- );
- AppSdkMigrationAuditClassification::Unsupported
- }
- AppSdkMigrationState::Unknown => {
- report.issue(
- "sdk_migration_receipt_unknown",
- Some(record_id),
- format!(
- "SDK migration receipt `{}` for operation `{}` is unknown",
- receipt.id, receipt.sdk_operation_kind
- ),
- );
- AppSdkMigrationAuditClassification::Unknown
- }
- }
-}
-
-fn receipt_allowed_for_source_classification(
- classification: AppSdkMigrationAuditClassification,
-) -> bool {
- matches!(
- classification,
- AppSdkMigrationAuditClassification::PublishableCandidate
- | AppSdkMigrationAuditClassification::AlreadyRepresentedCandidate
- )
-}
-
-fn shared_local_event_kind(record: &LocalEventRecord) -> String {
- match record.family {
- LocalRecordFamily::LocalWork => record
- .local_work_json
- .as_ref()
- .and_then(local_work_record_kind)
- .map(|kind| format!("local_work:{kind}"))
- .unwrap_or_else(|| "local_work:unknown".to_owned()),
- LocalRecordFamily::SignedEvent => record
- .event_kind
- .map(shared_signed_event_kind)
- .unwrap_or_else(|| "signed_event:unknown".to_owned()),
- }
-}
-
-fn shared_local_event_classification(
- record: &LocalEventRecord,
- report: &mut AppSdkMigrationAuditSourceBuilder,
-) -> AppSdkMigrationAuditClassification {
- match record.family {
- LocalRecordFamily::LocalWork => classify_shared_local_work(record, report),
- LocalRecordFamily::SignedEvent => classify_shared_signed_event(record, report),
- }
-}
-
-fn classify_shared_local_work(
- record: &LocalEventRecord,
- report: &mut AppSdkMigrationAuditSourceBuilder,
-) -> AppSdkMigrationAuditClassification {
- match record
- .local_work_json
- .as_ref()
- .and_then(local_work_record_kind)
- {
- Some("farm_config_v1" | "listing_draft_v1") => classify_shared_local_work_status(record),
- Some(record_kind) => {
- report.issue(
- "unsupported_shared_local_work_kind",
- Some(record.record_id.as_str()),
- format!("shared local work kind `{record_kind}` is not a SDK migration candidate"),
- );
- AppSdkMigrationAuditClassification::Unsupported
- }
- None => {
- report.issue(
- "unknown_shared_local_work_kind",
- Some(record.record_id.as_str()),
- "shared local work record does not expose a record_kind",
- );
- AppSdkMigrationAuditClassification::Unknown
- }
- }
-}
-
-fn classify_shared_local_work_status(
- record: &LocalEventRecord,
-) -> AppSdkMigrationAuditClassification {
- if matches!(record.outbox_status, PublishOutboxStatus::Acknowledged)
- || matches!(record.status, LocalRecordStatus::Published)
- {
- AppSdkMigrationAuditClassification::AlreadyRepresentedCandidate
- } else if matches!(record.outbox_status, PublishOutboxStatus::Failed)
- || matches!(
- record.status,
- LocalRecordStatus::Failed | LocalRecordStatus::Conflict
- )
- {
- AppSdkMigrationAuditClassification::ManualReviewRequired
- } else if matches!(record.status, LocalRecordStatus::PendingPublish) {
- AppSdkMigrationAuditClassification::PublishableCandidate
- } else {
- AppSdkMigrationAuditClassification::LocalWorkDeferred
- }
-}
-
-fn classify_shared_signed_event(
- record: &LocalEventRecord,
- report: &mut AppSdkMigrationAuditSourceBuilder,
-) -> AppSdkMigrationAuditClassification {
- match record.event_kind {
- Some(kind) if kind == KIND_TRADE_VALIDATION_RECEIPT as i64 => {
- AppSdkMigrationAuditClassification::ValidationReceiptDeferred
- }
- Some(kind) if supported_signed_event_kind(kind) => {
- if signed_event_is_already_represented(record.status, record.outbox_status) {
- AppSdkMigrationAuditClassification::AlreadyRepresentedCandidate
- } else {
- AppSdkMigrationAuditClassification::PublishableCandidate
- }
- }
- Some(kind) => {
- report.issue(
- "unsupported_shared_signed_event_kind",
- Some(record.record_id.as_str()),
- format!("shared signed event kind `{kind}` is not a SDK migration candidate"),
- );
- AppSdkMigrationAuditClassification::Unsupported
- }
- None => {
- report.issue(
- "unknown_shared_signed_event_kind",
- Some(record.record_id.as_str()),
- "shared signed event record does not expose an event_kind",
- );
- AppSdkMigrationAuditClassification::Unknown
- }
- }
-}
-
-fn signed_event_is_already_represented(
- status: LocalRecordStatus,
- outbox_status: PublishOutboxStatus,
-) -> bool {
- matches!(status, LocalRecordStatus::Published)
- || matches!(outbox_status, PublishOutboxStatus::Acknowledged)
-}
-
-fn shared_local_event_duplicate_identities(record: &LocalEventRecord) -> Vec<DuplicateIdentity> {
- let mut identities = Vec::new();
- if let (Some(event_kind), Some(event_id)) = (
- record.event_kind,
- non_empty_value(record.event_id.as_deref()),
- ) {
- identities.push(DuplicateIdentity {
- kind: "event".to_owned(),
- key: format!("{event_kind}:{event_id}"),
- });
- }
- if let Some(key) = shared_local_event_aggregate_key(record) {
- identities.push(DuplicateIdentity {
- kind: "aggregate".to_owned(),
- key,
- });
- }
- identities
-}
-
-fn shared_local_event_aggregate_key(record: &LocalEventRecord) -> Option<String> {
- match record.family {
- LocalRecordFamily::LocalWork => {
- let record_kind = record
- .local_work_json
- .as_ref()
- .and_then(local_work_record_kind)?;
- non_empty_value(record.farm_id.as_deref())
- .map(|farm_id| format!("local_work:{record_kind}:farm:{farm_id}"))
- .or_else(|| {
- non_empty_value(record.listing_addr.as_deref()).map(|listing_addr| {
- format!("local_work:{record_kind}:listing:{listing_addr}")
- })
- })
- }
- LocalRecordFamily::SignedEvent => {
- let event_kind = record.event_kind?;
- non_empty_value(record.listing_addr.as_deref())
- .map(|listing_addr| format!("signed_event:{event_kind}:listing:{listing_addr}"))
- .or_else(|| {
- non_empty_value(record.farm_id.as_deref())
- .map(|farm_id| format!("signed_event:{event_kind}:farm:{farm_id}"))
- })
- }
- }
-}
-
-fn local_work_record_kind(payload: &Value) -> Option<&str> {
- payload
- .get("record_kind")
- .and_then(Value::as_str)
- .map(str::trim)
- .filter(|value| !value.is_empty())
-}
-
-fn supported_signed_event_kind(kind: i64) -> bool {
- matches!(
- kind,
- value if value == KIND_FARM as i64
- || value == KIND_LISTING as i64
- || value == KIND_LISTING_DRAFT as i64
- || value == KIND_ORDER_REQUEST as i64
- || value == KIND_ORDER_DECISION as i64
- || value == KIND_ORDER_REVISION_PROPOSAL as i64
- || value == KIND_ORDER_REVISION_DECISION as i64
- || value == KIND_ORDER_CANCELLATION as i64
- )
-}
-
-fn shared_signed_event_kind(kind: i64) -> String {
- let name = match kind {
- value if value == KIND_FARM as i64 => "farm",
- value if value == KIND_LISTING as i64 => "listing",
- value if value == KIND_LISTING_DRAFT as i64 => "listing_draft",
- value if value == KIND_ORDER_REQUEST as i64 => "order_request",
- value if value == KIND_ORDER_DECISION as i64 => "order_decision",
- value if value == KIND_ORDER_REVISION_PROPOSAL as i64 => "order_revision_proposal",
- value if value == KIND_ORDER_REVISION_DECISION as i64 => "order_revision_decision",
- value if value == KIND_ORDER_CANCELLATION as i64 => "order_cancellation",
- value if value == KIND_TRADE_VALIDATION_RECEIPT as i64 => "trade_validation_receipt",
- _ => "unsupported",
- };
- format!("signed_event:{name}:{kind}")
-}
-
-fn non_empty_value(value: Option<&str>) -> Option<&str> {
- value.map(str::trim).filter(|value| !value.is_empty())
-}
-
-fn increment_count(counts: &mut BTreeMap<String, u64>, key: String) {
- *counts.entry(key).or_default() += 1;
-}
-
-fn counts_from_map(counts: BTreeMap<String, u64>) -> Vec<AppSdkMigrationAuditCount> {
- counts
- .into_iter()
- .map(|(key, count)| AppSdkMigrationAuditCount { key, count })
- .collect()
-}
-
-fn duplicate_candidates_from_map(
- duplicate_records: BTreeMap<DuplicateIdentity, BTreeSet<String>>,
-) -> Vec<AppSdkMigrationAuditDuplicateCandidate> {
- duplicate_records
- .into_iter()
- .filter_map(|(identity, records)| {
- if records.len() < 2 {
- return None;
- }
- Some(AppSdkMigrationAuditDuplicateCandidate {
- identity_kind: identity.kind,
- identity_key: identity.key,
- record_count: records.len() as u64,
- record_ids: records.into_iter().collect(),
- })
- })
- .collect()
-}
-
-#[cfg(test)]
-mod tests {
- use radroots_events::kinds::{KIND_LISTING, KIND_ORDER_REQUEST, KIND_TRADE_VALIDATION_RECEIPT};
- use radroots_local_events::{
- LocalEventRecord, LocalEventRecordInput, LocalEventsStore, LocalRecordFamily,
- LocalRecordStatus, PublishOutboxStatus, SourceRuntime,
- };
- use radroots_sql_core::SqliteExecutor;
- use radroots_studio_app_sync::{
- AppFarmProfilePublishPayload, AppPublishContext, AppPublishPayload, PendingSyncOperation,
- };
- use radroots_studio_app_view::{FarmId, FarmReadiness};
- use rusqlite::params;
- use serde_json::json;
-
- use crate::{
- AppSdkMigrationAuditClassification, AppSdkMigrationAuditRequest,
- AppSdkMigrationReceiptInput, AppSdkMigrationReceiptSourceKind, AppSdkMigrationState,
- AppSqliteStore, DatabaseTarget,
- };
-
- fn local_events_store() -> LocalEventsStore<SqliteExecutor> {
- let executor = SqliteExecutor::open_memory().expect("open local events memory db");
- let store = LocalEventsStore::new(executor);
- store.migrate_up().expect("migrate local events store");
- store
- }
-
- fn count_named(counts: &[crate::AppSdkMigrationAuditCount], key: &str) -> u64 {
- counts
- .iter()
- .find(|count| count.key == key)
- .map(|count| count.count)
- .unwrap_or_default()
- }
-
- #[test]
- fn local_outbox_audit_reads_batches_without_mutating_rows() {
- let store = AppSqliteStore::open(DatabaseTarget::InMemory).expect("open app store");
- let shared_events = local_events_store();
- let farm_id = FarmId::new();
- let operation = PendingSyncOperation::from_publish_payload(
- AppPublishPayload::FarmProfile(AppFarmProfilePublishPayload {
- context: AppPublishContext::new("acct_a", "farm_setup"),
- farm_id,
- display_name: "Green Loop Farm".to_owned(),
- readiness: Some(FarmReadiness::Ready),
- }),
- "2026-06-18T12:00:00Z",
- )
- .expect("build publish operation");
- store
- .sync_repository()
- .enqueue_pending_operation("acct_a", &operation)
- .expect("enqueue operation");
- store
- .connection()
- .execute(
- "INSERT INTO local_outbox (
- id,
- account_id,
- operation_key,
- aggregate_kind,
- aggregate_id,
- operation_kind,
- payload_json,
- created_at,
- available_at,
- attempt_count,
- state,
- last_error_message
- ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, NULL)",
- params![
- "succeeded-duplicate",
- "acct_a",
- operation.operation_key,
- operation.aggregate.aggregate_kind(),
- operation.aggregate.aggregate_id(),
- operation.operation.storage_key(),
- operation.payload_json,
- "2026-06-18T11:00:00Z",
- "2026-06-18T11:00:00Z",
- 0_i64,
- "succeeded",
- ],
- )
- .expect("insert succeeded duplicate");
- let before_count = local_outbox_row_count(&store);
-
- let report = store
- .audit_sdk_migration(
- &shared_events,
- AppSdkMigrationAuditRequest { batch_size: 1 },
- )
- .expect("audit should run");
-
- assert_eq!(local_outbox_row_count(&store), before_count);
- assert_eq!(report.local_outbox.batch_size, 1);
- assert_eq!(report.local_outbox.batch_count, 2);
- assert_eq!(report.local_outbox.scanned_records, 2);
- assert_eq!(
- count_named(&report.local_outbox.kind_counts, "farm_profile"),
- 2
- );
- assert_eq!(
- count_named(&report.local_outbox.status_counts, "pending"),
- 1
- );
- assert_eq!(
- count_named(&report.local_outbox.status_counts, "succeeded"),
- 1
- );
- assert_eq!(
- count_named(
- &report.local_outbox.classification_counts,
- AppSdkMigrationAuditClassification::PublishableCandidate.storage_key()
- ),
- 1
- );
- assert_eq!(
- count_named(
- &report.local_outbox.classification_counts,
- AppSdkMigrationAuditClassification::AlreadyRepresentedCandidate.storage_key()
- ),
- 1
- );
- assert!(
- report
- .local_outbox
- .duplicate_candidates
- .iter()
- .any(|candidate| candidate.identity_kind == "operation"
- && candidate.record_count == 2)
- );
- assert_eq!(report.shared_local_events.scanned_records, 0);
- }
-
- #[test]
- fn local_outbox_audit_classifies_status_matrix() {
- let store = AppSqliteStore::open(DatabaseTarget::InMemory).expect("open app store");
- let shared_events = local_events_store();
- let operation = farm_profile_operation("acct_seed", "status_matrix");
-
- for (index, state) in [
- "pending",
- "in_progress",
- "retryable",
- "failed",
- "blocked",
- "succeeded",
- ]
- .iter()
- .enumerate()
- {
- insert_local_outbox_audit_row(
- &store,
- &format!("local-outbox-{state}"),
- &format!("acct_{index}"),
- state,
- &operation,
- );
- }
-
- let report = store
- .audit_sdk_migration(
- &shared_events,
- AppSdkMigrationAuditRequest { batch_size: 2 },
- )
- .expect("audit should run");
-
- assert_eq!(report.local_outbox.scanned_records, 6);
- assert_eq!(
- count_named(
- &report.local_outbox.classification_counts,
- AppSdkMigrationAuditClassification::PublishableCandidate.storage_key()
- ),
- 3
- );
- assert_eq!(
- count_named(
- &report.local_outbox.classification_counts,
- AppSdkMigrationAuditClassification::ManualReviewRequired.storage_key()
- ),
- 2
- );
- assert_eq!(
- count_named(
- &report.local_outbox.classification_counts,
- AppSdkMigrationAuditClassification::AlreadyRepresentedCandidate.storage_key()
- ),
- 1
- );
- assert_eq!(
- report
- .local_outbox
- .issues
- .iter()
- .filter(|issue| issue.code == "manual_review_local_outbox_state")
- .count(),
- 2
- );
- }
-
- #[test]
- fn local_outbox_audit_uses_migration_receipts_for_migratable_records() {
- let store = AppSqliteStore::open(DatabaseTarget::InMemory).expect("open app store");
- let shared_events = local_events_store();
- let operation = farm_profile_operation("acct_seed", "receipt_matrix");
-
- for (id, state) in [
- ("represented-source", AppSdkMigrationState::Enqueued),
- ("skipped-source", AppSdkMigrationState::Skipped),
- ("failed-source", AppSdkMigrationState::Failed),
- ] {
- insert_local_outbox_audit_row(&store, id, id, "pending", &operation);
- record_local_outbox_receipt(&store, id, state);
- }
-
- let report = store
- .audit_sdk_migration(
- &shared_events,
- AppSdkMigrationAuditRequest { batch_size: 10 },
- )
- .expect("audit should run");
-
- assert_eq!(report.local_outbox.scanned_records, 3);
- assert_eq!(
- count_named(
- &report.local_outbox.classification_counts,
- AppSdkMigrationAuditClassification::RepresentedRecord.storage_key()
- ),
- 1
- );
- assert_eq!(
- count_named(
- &report.local_outbox.classification_counts,
- AppSdkMigrationAuditClassification::SkippedRecord.storage_key()
- ),
- 1
- );
- assert_eq!(
- count_named(
- &report.local_outbox.classification_counts,
- AppSdkMigrationAuditClassification::FailedRecord.storage_key()
- ),
- 1
- );
- assert!(
- report
- .local_outbox
- .issues
- .iter()
- .any(|issue| issue.code == "sdk_migration_receipt_failed")
- );
- }
-
- #[test]
- fn local_outbox_audit_does_not_let_receipts_hide_non_migratable_rows() {
- let store = AppSqliteStore::open(DatabaseTarget::InMemory).expect("open app store");
- let shared_events = local_events_store();
- let operation = farm_profile_operation("acct_seed", "non_migratable");
-
- insert_local_outbox_audit_row(&store, "failed-source", "acct_failed", "failed", &operation);
- record_local_outbox_receipt(&store, "failed-source", AppSdkMigrationState::Enqueued);
- insert_local_outbox_audit_row(
- &store,
- "unsupported-source",
- "acct_unsupported",
- "pending",
- &PendingSyncOperation {
- operation: radroots_studio_app_sync::SyncOperationKind::Delete,
- ..operation.clone()
- },
- );
- record_local_outbox_receipt(&store, "unsupported-source", AppSdkMigrationState::Enqueued);
- store
- .connection()
- .execute_batch("PRAGMA ignore_check_constraints = ON")
- .expect("disable sqlite checks for defensive unknown state row");
- insert_local_outbox_audit_row(
- &store,
- "unknown-source",
- "acct_unknown",
- "mystery",
- &operation,
- );
- store
- .connection()
- .execute_batch("PRAGMA ignore_check_constraints = OFF")
- .expect("restore sqlite checks");
- record_local_outbox_receipt(&store, "unknown-source", AppSdkMigrationState::Enqueued);
-
- let report = store
- .audit_sdk_migration(
- &shared_events,
- AppSdkMigrationAuditRequest { batch_size: 10 },
- )
- .expect("audit should run");
-
- assert_eq!(
- count_named(
- &report.local_outbox.classification_counts,
- AppSdkMigrationAuditClassification::RepresentedRecord.storage_key()
- ),
- 0
- );
- assert_eq!(
- count_named(
- &report.local_outbox.classification_counts,
- AppSdkMigrationAuditClassification::ManualReviewRequired.storage_key()
- ),
- 1
- );
- assert_eq!(
- count_named(
- &report.local_outbox.classification_counts,
- AppSdkMigrationAuditClassification::Unsupported.storage_key()
- ),
- 1
- );
- assert_eq!(
- count_named(
- &report.local_outbox.classification_counts,
- AppSdkMigrationAuditClassification::Unknown.storage_key()
- ),
- 1
- );
- assert_eq!(
- report
- .local_outbox
- .issues
- .iter()
- .filter(|issue| issue.code == "sdk_migration_receipt_for_non_migratable_source")
- .count(),
- 3
- );
- }
-
- #[test]
- fn shared_local_events_audit_classifies_supported_events_and_validation_receipts() {
- let store = AppSqliteStore::open(DatabaseTarget::InMemory).expect("open app store");
- let shared_events = local_events_store();
- shared_events
- .append_record(&signed_event_record(
- "listing-a",
- "duplicate-listing-event",
- KIND_LISTING as i64,
- ))
- .expect("append listing a");
- shared_events
- .append_record(&signed_event_record(
- "listing-b",
- "duplicate-listing-event",
- KIND_LISTING as i64,
- ))
- .expect("append listing b");
- shared_events
- .append_record(&signed_event_record(
- "request",
- "request-event",
- KIND_ORDER_REQUEST as i64,
- ))
- .expect("append request");
- shared_events
- .append_record(&signed_event_record(
- "validation-receipt",
- "validation-receipt-event",
- KIND_TRADE_VALIDATION_RECEIPT as i64,
- ))
- .expect("append validation receipt");
- let before_records = shared_events
- .list_records_changed_after(0, 10)
- .expect("list records before audit")
- .len();
-
- let report = store
- .audit_sdk_migration(
- &shared_events,
- AppSdkMigrationAuditRequest { batch_size: 1 },
- )
- .expect("audit should run");
-
- assert_eq!(
- shared_events
- .list_records_changed_after(0, 10)
- .expect("list records after audit")
- .len(),
- before_records
- );
- assert_eq!(report.shared_local_events.batch_count, 4);
- assert_eq!(report.shared_local_events.scanned_records, 4);
- assert_eq!(
- count_named(
- &report.shared_local_events.classification_counts,
- AppSdkMigrationAuditClassification::AlreadyRepresentedCandidate.storage_key()
- ),
- 3
- );
- assert_eq!(
- count_named(
- &report.shared_local_events.classification_counts,
- AppSdkMigrationAuditClassification::ValidationReceiptDeferred.storage_key()
- ),
- 1
- );
- assert_eq!(
- count_named(
- &report.shared_local_events.classification_counts,
- AppSdkMigrationAuditClassification::PublishableCandidate.storage_key()
- ),
- 0
- );
- assert!(
- report
- .shared_local_events
- .duplicate_candidates
- .iter()
- .any(|candidate| candidate.identity_kind == "event" && candidate.record_count == 2)
- );
- }
-
- #[test]
- fn shared_local_work_audit_classifies_status_matrix() {
- let store = AppSqliteStore::open(DatabaseTarget::InMemory).expect("open app store");
- let shared_events = local_events_store();
-
- for (record_id, record_kind, status, outbox_status) in [
- (
- "local-draft",
- "farm_config_v1",
- LocalRecordStatus::LocalDraft,
- PublishOutboxStatus::None,
- ),
- (
- "local-saved",
- "listing_draft_v1",
- LocalRecordStatus::LocalSaved,
- PublishOutboxStatus::None,
- ),
- (
- "pending-publish",
- "listing_draft_v1",
- LocalRecordStatus::PendingPublish,
- PublishOutboxStatus::None,
- ),
- (
- "published",
- "farm_config_v1",
- LocalRecordStatus::Published,
- PublishOutboxStatus::None,
- ),
- (
- "failed",
- "farm_config_v1",
- LocalRecordStatus::Failed,
- PublishOutboxStatus::None,
- ),
- (
- "conflict",
- "listing_draft_v1",
- LocalRecordStatus::Conflict,
- PublishOutboxStatus::None,
- ),
- ] {
- shared_events
- .append_record(&local_work_record(
- record_id,
- record_kind,
- status,
- outbox_status,
- ))
- .expect("append local work record");
- }
-
- let report = store
- .audit_sdk_migration(
- &shared_events,
- AppSdkMigrationAuditRequest { batch_size: 3 },
- )
- .expect("audit should run");
-
- assert_eq!(report.shared_local_events.scanned_records, 6);
- assert_eq!(
- count_named(
- &report.shared_local_events.classification_counts,
- AppSdkMigrationAuditClassification::LocalWorkDeferred.storage_key()
- ),
- 2
- );
- assert_eq!(
- count_named(
- &report.shared_local_events.classification_counts,
- AppSdkMigrationAuditClassification::PublishableCandidate.storage_key()
- ),
- 1
- );
- assert_eq!(
- count_named(
- &report.shared_local_events.classification_counts,
- AppSdkMigrationAuditClassification::AlreadyRepresentedCandidate.storage_key()
- ),
- 1
- );
- assert_eq!(
- count_named(
- &report.shared_local_events.classification_counts,
- AppSdkMigrationAuditClassification::ManualReviewRequired.storage_key()
- ),
- 2
- );
- }
-
- #[test]
- fn shared_local_work_status_classifier_handles_defensive_outbox_states() {
- assert_eq!(
- super::classify_shared_local_work_status(&local_work_model_record(
- LocalRecordStatus::PendingPublish,
- PublishOutboxStatus::Acknowledged,
- )),
- AppSdkMigrationAuditClassification::AlreadyRepresentedCandidate
- );
- assert_eq!(
- super::classify_shared_local_work_status(&local_work_model_record(
- LocalRecordStatus::PendingPublish,
- PublishOutboxStatus::Failed,
- )),
- AppSdkMigrationAuditClassification::ManualReviewRequired
- );
- }
-
- fn local_outbox_row_count(store: &AppSqliteStore) -> i64 {
- store
- .connection()
- .query_row("SELECT count(*) FROM local_outbox", [], |row| row.get(0))
- .expect("count local outbox rows")
- }
-
- fn farm_profile_operation(account_id: &str, source: &str) -> PendingSyncOperation {
- PendingSyncOperation::from_publish_payload(
- AppPublishPayload::FarmProfile(AppFarmProfilePublishPayload {
- context: AppPublishContext::new(account_id, source),
- farm_id: FarmId::new(),
- display_name: "Green Loop Farm".to_owned(),
- readiness: Some(FarmReadiness::Ready),
- }),
- "2026-06-18T12:00:00Z",
- )
- .expect("build publish operation")
- }
-
- fn insert_local_outbox_audit_row(
- store: &AppSqliteStore,
- id: &str,
- account_id: &str,
- state: &str,
- operation: &PendingSyncOperation,
- ) {
- store
- .connection()
- .execute(
- "INSERT INTO local_outbox (
- id,
- account_id,
- operation_key,
- aggregate_kind,
- aggregate_id,
- operation_kind,
- payload_json,
- created_at,
- available_at,
- attempt_count,
- state,
- last_error_message
- ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, NULL)",
- params![
- id,
- account_id,
- operation.operation_key.as_str(),
- operation.aggregate.aggregate_kind(),
- operation.aggregate.aggregate_id(),
- operation.operation.storage_key(),
- operation.payload_json.as_str(),
- operation.created_at.as_str(),
- operation.available_at.as_str(),
- i64::from(operation.attempt_count),
- state,
- ],
- )
- .expect("insert local outbox audit row");
- }
-
- fn record_local_outbox_receipt(
- store: &AppSqliteStore,
- source_record_id: &str,
- migration_state: AppSdkMigrationState,
- ) {
- store
- .sdk_migration_receipt_repository()
- .record_receipt(&AppSdkMigrationReceiptInput {
- source_kind: AppSdkMigrationReceiptSourceKind::LocalOutbox,
- source_record_id: source_record_id.to_owned(),
- sdk_operation_kind: "farm.publish".to_owned(),
- sdk_outbox_event_ids: vec![format!("sdk-outbox-{source_record_id}")],
- expected_event_id: Some(format!("event-{source_record_id}")),
- actor_pubkey: Some("actor-pubkey".to_owned()),
- idempotency_digest_prefix: Some("digest-prefix".to_owned()),
- migration_state,
- recorded_at: "2026-06-18T12:00:00Z".to_owned(),
- detail_json: json!({"source": source_record_id}),
- })
- .expect("record local outbox receipt");
- }
-
- fn local_work_record(
- record_id: &str,
- record_kind: &str,
- status: LocalRecordStatus,
- outbox_status: PublishOutboxStatus,
- ) -> LocalEventRecordInput {
- LocalEventRecordInput {
- record_id: record_id.to_owned(),
- family: LocalRecordFamily::LocalWork,
- status,
- source_runtime: SourceRuntime::App,
- created_at_ms: 1000,
- inserted_at_ms: 1001,
- owner_account_id: Some("acct_a".to_owned()),
- owner_pubkey: Some("seller-pubkey".to_owned()),
- farm_id: Some("farm-key".to_owned()),
- listing_addr: Some("30402:seller-pubkey:listing-key".to_owned()),
- local_work_json: Some(json!({"record_kind": record_kind})),
- event_id: None,
- event_kind: None,
- event_pubkey: None,
- event_created_at: None,
- event_tags_json: None,
- event_content: None,
- event_sig: None,
- raw_event_json: None,
- outbox_status,
- relay_set_fingerprint: None,
- relay_delivery_json: None,
- }
- }
-
- fn local_work_model_record(
- status: LocalRecordStatus,
- outbox_status: PublishOutboxStatus,
- ) -> LocalEventRecord {
- LocalEventRecord {
- seq: 1,
- change_seq: 1,
- record_id: "defensive-local-work".to_owned(),
- family: LocalRecordFamily::LocalWork,
- status,
- source_runtime: SourceRuntime::App,
- created_at_ms: 1000,
- inserted_at_ms: 1001,
- updated_at_ms: 1002,
- owner_account_id: Some("acct_a".to_owned()),
- owner_pubkey: Some("seller-pubkey".to_owned()),
- farm_id: Some("farm-key".to_owned()),
- listing_addr: Some("30402:seller-pubkey:listing-key".to_owned()),
- local_work_json: Some(json!({"record_kind": "listing_draft_v1"})),
- event_id: None,
- event_kind: None,
- event_pubkey: None,
- event_created_at: None,
- event_tags_json: None,
- event_content: None,
- event_sig: None,
- raw_event_json: None,
- outbox_status,
- relay_set_fingerprint: None,
- relay_delivery_json: None,
- }
- }
-
- fn signed_event_record(
- record_id: &str,
- event_id: &str,
- event_kind: i64,
- ) -> LocalEventRecordInput {
- LocalEventRecordInput {
- record_id: record_id.to_owned(),
- family: LocalRecordFamily::SignedEvent,
- status: LocalRecordStatus::Published,
- source_runtime: SourceRuntime::App,
- created_at_ms: 1000,
- inserted_at_ms: 1001,
- owner_account_id: Some("acct_a".to_owned()),
- owner_pubkey: Some("seller-pubkey".to_owned()),
- farm_id: Some("farm-key".to_owned()),
- listing_addr: Some("30402:seller-pubkey:listing-key".to_owned()),
- local_work_json: None,
- event_id: Some(event_id.to_owned()),
- event_kind: Some(event_kind),
- event_pubkey: Some("seller-pubkey".to_owned()),
- event_created_at: Some(1000),
- event_tags_json: Some(json!([["d", "listing-key"]])),
- event_content: Some("{}".to_owned()),
- event_sig: Some("signature".to_owned()),
- raw_event_json: Some(json!({
- "id": event_id,
- "kind": event_kind,
- "pubkey": "seller-pubkey"
- })),
- outbox_status: PublishOutboxStatus::Acknowledged,
- relay_set_fingerprint: Some("relay-set".to_owned()),
- relay_delivery_json: Some(json!({
- "state": "acknowledged",
- "acknowledged_relays": ["wss://relay.example"]
- })),
- }
- }
-}
diff --git a/crates/store/src/migrations.rs b/crates/store/src/migrations.rs
@@ -116,7 +116,7 @@ const MIGRATIONS: &[Migration] = &[
},
Migration {
version: 28,
- sql: include_str!("../migrations/0028_sdk_migration_receipts.sql"),
+ sql: include_str!("../migrations/0028_sdk_workflow_receipts.sql"),
},
Migration {
version: 29,
diff --git a/crates/store/src/sdk_migration_receipts.rs b/crates/store/src/sdk_migration_receipts.rs
@@ -1,316 +0,0 @@
-use rusqlite::{Connection, OptionalExtension, params};
-use serde_json::Value;
-use uuid::Uuid;
-
-use crate::AppSqliteError;
-
-#[derive(Clone, Copy, Debug, Eq, PartialEq)]
-pub enum AppSdkMigrationReceiptSourceKind {
- LocalOutbox,
- SharedLocalEvent,
-}
-
-impl AppSdkMigrationReceiptSourceKind {
- pub const fn storage_key(self) -> &'static str {
- match self {
- Self::LocalOutbox => "local_outbox",
- Self::SharedLocalEvent => "shared_local_event",
- }
- }
-
- pub fn parse(value: &str) -> Result<Self, AppSqliteError> {
- match value {
- "local_outbox" => Ok(Self::LocalOutbox),
- "shared_local_event" => Ok(Self::SharedLocalEvent),
- _ => Err(AppSqliteError::DecodeEnum {
- field: "app_sdk_migration_receipts.source_kind",
- value: value.to_owned(),
- }),
- }
- }
-}
-
-#[derive(Clone, Copy, Debug, Eq, PartialEq)]
-pub enum AppSdkMigrationState {
- Pending,
- Prepared,
- Enqueued,
- Pushed,
- Failed,
- Blocked,
- Skipped,
- Unsupported,
- ManualReview,
- Unknown,
-}
-
-impl AppSdkMigrationState {
- pub const fn storage_key(self) -> &'static str {
- match self {
- Self::Pending => "pending",
- Self::Prepared => "prepared",
- Self::Enqueued => "enqueued",
- Self::Pushed => "pushed",
- Self::Failed => "failed",
- Self::Blocked => "blocked",
- Self::Skipped => "skipped",
- Self::Unsupported => "unsupported",
- Self::ManualReview => "manual_review",
- Self::Unknown => "unknown",
- }
- }
-
- pub fn parse(value: &str) -> Result<Self, AppSqliteError> {
- match value {
- "pending" => Ok(Self::Pending),
- "prepared" => Ok(Self::Prepared),
- "enqueued" => Ok(Self::Enqueued),
- "pushed" => Ok(Self::Pushed),
- "failed" => Ok(Self::Failed),
- "blocked" => Ok(Self::Blocked),
- "skipped" => Ok(Self::Skipped),
- "unsupported" => Ok(Self::Unsupported),
- "manual_review" => Ok(Self::ManualReview),
- "unknown" => Ok(Self::Unknown),
- _ => Err(AppSqliteError::DecodeEnum {
- field: "app_sdk_migration_receipts.migration_state",
- value: value.to_owned(),
- }),
- }
- }
-}
-
-#[derive(Clone, Debug, Eq, PartialEq)]
-pub struct AppSdkMigrationReceiptInput {
- pub source_kind: AppSdkMigrationReceiptSourceKind,
- pub source_record_id: String,
- pub sdk_operation_kind: String,
- pub sdk_outbox_event_ids: Vec<String>,
- pub expected_event_id: Option<String>,
- pub actor_pubkey: Option<String>,
- pub idempotency_digest_prefix: Option<String>,
- pub migration_state: AppSdkMigrationState,
- pub recorded_at: String,
- pub detail_json: Value,
-}
-
-#[derive(Clone, Debug, Eq, PartialEq)]
-pub struct AppSdkMigrationReceipt {
- pub id: String,
- pub source_kind: AppSdkMigrationReceiptSourceKind,
- pub source_record_id: String,
- pub sdk_operation_kind: String,
- pub sdk_outbox_event_ids: Vec<String>,
- pub expected_event_id: Option<String>,
- pub actor_pubkey: Option<String>,
- pub idempotency_digest_prefix: Option<String>,
- pub migration_state: AppSdkMigrationState,
- pub created_at: String,
- pub updated_at: String,
- pub detail_json: Value,
-}
-
-pub struct AppSdkMigrationReceiptRepository<'a> {
- connection: &'a Connection,
-}
-
-impl<'a> AppSdkMigrationReceiptRepository<'a> {
- pub const fn new(connection: &'a Connection) -> Self {
- Self { connection }
- }
-
- pub fn record_receipt(
- &self,
- input: &AppSdkMigrationReceiptInput,
- ) -> Result<AppSdkMigrationReceipt, AppSqliteError> {
- let receipt_id = Uuid::now_v7().to_string();
- let outbox_ids_json =
- serde_json::to_string(&input.sdk_outbox_event_ids).map_err(|source| {
- AppSqliteError::EncodeJson {
- field: "app_sdk_migration_receipts.sdk_outbox_event_ids_json",
- source,
- }
- })?;
- let detail_json = serde_json::to_string(&input.detail_json).map_err(|source| {
- AppSqliteError::EncodeJson {
- field: "app_sdk_migration_receipts.detail_json",
- source,
- }
- })?;
-
- self.connection
- .execute(
- "INSERT INTO app_sdk_migration_receipts (
- id,
- source_kind,
- source_record_id,
- sdk_operation_kind,
- sdk_outbox_event_ids_json,
- expected_event_id,
- actor_pubkey,
- idempotency_digest_prefix,
- migration_state,
- created_at,
- updated_at,
- detail_json
- ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?10, ?11)
- ON CONFLICT(source_kind, source_record_id)
- DO UPDATE SET
- sdk_operation_kind = excluded.sdk_operation_kind,
- sdk_outbox_event_ids_json = excluded.sdk_outbox_event_ids_json,
- expected_event_id = excluded.expected_event_id,
- actor_pubkey = excluded.actor_pubkey,
- idempotency_digest_prefix = excluded.idempotency_digest_prefix,
- migration_state = excluded.migration_state,
- updated_at = excluded.updated_at,
- detail_json = excluded.detail_json",
- params![
- receipt_id,
- input.source_kind.storage_key(),
- input.source_record_id.as_str(),
- input.sdk_operation_kind.as_str(),
- outbox_ids_json.as_str(),
- input.expected_event_id.as_deref(),
- input.actor_pubkey.as_deref(),
- input.idempotency_digest_prefix.as_deref(),
- input.migration_state.storage_key(),
- input.recorded_at.as_str(),
- detail_json.as_str(),
- ],
- )
- .map_err(|source| AppSqliteError::Query {
- operation: "record app SDK migration receipt",
- source,
- })?;
-
- self.load_receipt(input.source_kind, input.source_record_id.as_str())?
- .ok_or(AppSqliteError::MissingColumn {
- field: "app_sdk_migration_receipts.id",
- })
- }
-
- pub fn load_receipt(
- &self,
- source_kind: AppSdkMigrationReceiptSourceKind,
- source_record_id: &str,
- ) -> Result<Option<AppSdkMigrationReceipt>, AppSqliteError> {
- self.connection
- .query_row(
- "SELECT
- id,
- source_kind,
- source_record_id,
- sdk_operation_kind,
- sdk_outbox_event_ids_json,
- expected_event_id,
- actor_pubkey,
- idempotency_digest_prefix,
- migration_state,
- created_at,
- updated_at,
- detail_json
- FROM app_sdk_migration_receipts
- WHERE source_kind = ?1
- AND source_record_id = ?2
- LIMIT 1",
- params![source_kind.storage_key(), source_record_id],
- decode_receipt_row,
- )
- .optional()
- .map_err(|source| AppSqliteError::Query {
- operation: "load app SDK migration receipt",
- source,
- })
- }
-}
-
-fn decode_receipt_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<AppSdkMigrationReceipt> {
- let source_kind: String = row.get(1)?;
- let outbox_ids_json: String = row.get(4)?;
- let migration_state: String = row.get(8)?;
- let detail_json: String = row.get(11)?;
- Ok(AppSdkMigrationReceipt {
- id: row.get(0)?,
- source_kind: AppSdkMigrationReceiptSourceKind::parse(source_kind.as_str())
- .map_err(decode_app_error)?,
- source_record_id: row.get(2)?,
- sdk_operation_kind: row.get(3)?,
- sdk_outbox_event_ids: serde_json::from_str(outbox_ids_json.as_str()).map_err(|source| {
- decode_app_error(AppSqliteError::DecodeJson {
- field: "app_sdk_migration_receipts.sdk_outbox_event_ids_json",
- source,
- })
- })?,
- expected_event_id: row.get(5)?,
- actor_pubkey: row.get(6)?,
- idempotency_digest_prefix: row.get(7)?,
- migration_state: AppSdkMigrationState::parse(migration_state.as_str())
- .map_err(decode_app_error)?,
- created_at: row.get(9)?,
- updated_at: row.get(10)?,
- detail_json: serde_json::from_str(detail_json.as_str()).map_err(|source| {
- decode_app_error(AppSqliteError::DecodeJson {
- field: "app_sdk_migration_receipts.detail_json",
- source,
- })
- })?,
- })
-}
-
-fn decode_app_error(error: AppSqliteError) -> rusqlite::Error {
- rusqlite::Error::FromSqlConversionFailure(0, rusqlite::types::Type::Text, Box::new(error))
-}
-
-#[cfg(test)]
-mod tests {
- use serde_json::json;
-
- use crate::{
- AppSdkMigrationReceiptInput, AppSdkMigrationReceiptSourceKind, AppSdkMigrationState,
- AppSqliteStore, DatabaseTarget,
- };
-
- #[test]
- fn migration_receipts_are_idempotent_by_source_record() {
- let store = AppSqliteStore::open(DatabaseTarget::InMemory).expect("open app store");
- let first = store
- .sdk_migration_receipt_repository()
- .record_receipt(&AppSdkMigrationReceiptInput {
- source_kind: AppSdkMigrationReceiptSourceKind::LocalOutbox,
- source_record_id: "source-record-a".to_owned(),
- sdk_operation_kind: "farm.publish".to_owned(),
- sdk_outbox_event_ids: vec!["outbox-a".to_owned()],
- expected_event_id: Some("expected-a".to_owned()),
- actor_pubkey: Some("actor-a".to_owned()),
- idempotency_digest_prefix: Some("digest-a".to_owned()),
- migration_state: AppSdkMigrationState::Enqueued,
- recorded_at: "2026-06-18T12:00:00Z".to_owned(),
- detail_json: json!({"attempt": 1}),
- })
- .expect("record first receipt");
- let second = store
- .sdk_migration_receipt_repository()
- .record_receipt(&AppSdkMigrationReceiptInput {
- source_kind: AppSdkMigrationReceiptSourceKind::LocalOutbox,
- source_record_id: "source-record-a".to_owned(),
- sdk_operation_kind: "farm.publish".to_owned(),
- sdk_outbox_event_ids: vec!["outbox-b".to_owned()],
- expected_event_id: Some("expected-b".to_owned()),
- actor_pubkey: Some("actor-b".to_owned()),
- idempotency_digest_prefix: Some("digest-b".to_owned()),
- migration_state: AppSdkMigrationState::Pushed,
- recorded_at: "2026-06-18T12:05:00Z".to_owned(),
- detail_json: json!({"attempt": 2}),
- })
- .expect("record second receipt");
-
- assert_eq!(first.id, second.id);
- assert_eq!(second.created_at, "2026-06-18T12:00:00Z");
- assert_eq!(second.updated_at, "2026-06-18T12:05:00Z");
- assert_eq!(second.sdk_outbox_event_ids, vec!["outbox-b".to_owned()]);
- assert_eq!(second.expected_event_id.as_deref(), Some("expected-b"));
- assert_eq!(second.actor_pubkey.as_deref(), Some("actor-b"));
- assert_eq!(second.migration_state, AppSdkMigrationState::Pushed);
- assert_eq!(second.detail_json, json!({"attempt": 2}));
- }
-}
diff --git a/crates/store/src/sdk_workflow_receipts.rs b/crates/store/src/sdk_workflow_receipts.rs
@@ -0,0 +1,316 @@
+use rusqlite::{Connection, OptionalExtension, params};
+use serde_json::Value;
+use uuid::Uuid;
+
+use crate::AppSqliteError;
+
+#[derive(Clone, Copy, Debug, Eq, PartialEq)]
+pub enum AppSdkWorkflowReceiptSourceKind {
+ LocalOutbox,
+ SharedLocalEvent,
+}
+
+impl AppSdkWorkflowReceiptSourceKind {
+ pub const fn storage_key(self) -> &'static str {
+ match self {
+ Self::LocalOutbox => "local_outbox",
+ Self::SharedLocalEvent => "shared_local_event",
+ }
+ }
+
+ pub fn parse(value: &str) -> Result<Self, AppSqliteError> {
+ match value {
+ "local_outbox" => Ok(Self::LocalOutbox),
+ "shared_local_event" => Ok(Self::SharedLocalEvent),
+ _ => Err(AppSqliteError::DecodeEnum {
+ field: "app_sdk_workflow_receipts.source_kind",
+ value: value.to_owned(),
+ }),
+ }
+ }
+}
+
+#[derive(Clone, Copy, Debug, Eq, PartialEq)]
+pub enum AppSdkWorkflowReceiptState {
+ Pending,
+ Prepared,
+ Enqueued,
+ Pushed,
+ Failed,
+ Blocked,
+ Skipped,
+ Unsupported,
+ ManualReview,
+ Unknown,
+}
+
+impl AppSdkWorkflowReceiptState {
+ pub const fn storage_key(self) -> &'static str {
+ match self {
+ Self::Pending => "pending",
+ Self::Prepared => "prepared",
+ Self::Enqueued => "enqueued",
+ Self::Pushed => "pushed",
+ Self::Failed => "failed",
+ Self::Blocked => "blocked",
+ Self::Skipped => "skipped",
+ Self::Unsupported => "unsupported",
+ Self::ManualReview => "manual_review",
+ Self::Unknown => "unknown",
+ }
+ }
+
+ pub fn parse(value: &str) -> Result<Self, AppSqliteError> {
+ match value {
+ "pending" => Ok(Self::Pending),
+ "prepared" => Ok(Self::Prepared),
+ "enqueued" => Ok(Self::Enqueued),
+ "pushed" => Ok(Self::Pushed),
+ "failed" => Ok(Self::Failed),
+ "blocked" => Ok(Self::Blocked),
+ "skipped" => Ok(Self::Skipped),
+ "unsupported" => Ok(Self::Unsupported),
+ "manual_review" => Ok(Self::ManualReview),
+ "unknown" => Ok(Self::Unknown),
+ _ => Err(AppSqliteError::DecodeEnum {
+ field: "app_sdk_workflow_receipts.workflow_state",
+ value: value.to_owned(),
+ }),
+ }
+ }
+}
+
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct AppSdkWorkflowReceiptInput {
+ pub source_kind: AppSdkWorkflowReceiptSourceKind,
+ pub source_record_id: String,
+ pub sdk_operation_kind: String,
+ pub sdk_outbox_event_ids: Vec<String>,
+ pub expected_event_id: Option<String>,
+ pub actor_pubkey: Option<String>,
+ pub idempotency_digest_prefix: Option<String>,
+ pub workflow_state: AppSdkWorkflowReceiptState,
+ pub recorded_at: String,
+ pub detail_json: Value,
+}
+
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct AppSdkStoredWorkflowReceipt {
+ pub id: String,
+ pub source_kind: AppSdkWorkflowReceiptSourceKind,
+ pub source_record_id: String,
+ pub sdk_operation_kind: String,
+ pub sdk_outbox_event_ids: Vec<String>,
+ pub expected_event_id: Option<String>,
+ pub actor_pubkey: Option<String>,
+ pub idempotency_digest_prefix: Option<String>,
+ pub workflow_state: AppSdkWorkflowReceiptState,
+ pub created_at: String,
+ pub updated_at: String,
+ pub detail_json: Value,
+}
+
+pub struct AppSdkWorkflowReceiptRepository<'a> {
+ connection: &'a Connection,
+}
+
+impl<'a> AppSdkWorkflowReceiptRepository<'a> {
+ pub const fn new(connection: &'a Connection) -> Self {
+ Self { connection }
+ }
+
+ pub fn record_receipt(
+ &self,
+ input: &AppSdkWorkflowReceiptInput,
+ ) -> Result<AppSdkStoredWorkflowReceipt, AppSqliteError> {
+ let receipt_id = Uuid::now_v7().to_string();
+ let outbox_ids_json =
+ serde_json::to_string(&input.sdk_outbox_event_ids).map_err(|source| {
+ AppSqliteError::EncodeJson {
+ field: "app_sdk_workflow_receipts.sdk_outbox_event_ids_json",
+ source,
+ }
+ })?;
+ let detail_json = serde_json::to_string(&input.detail_json).map_err(|source| {
+ AppSqliteError::EncodeJson {
+ field: "app_sdk_workflow_receipts.detail_json",
+ source,
+ }
+ })?;
+
+ self.connection
+ .execute(
+ "INSERT INTO app_sdk_workflow_receipts (
+ id,
+ source_kind,
+ source_record_id,
+ sdk_operation_kind,
+ sdk_outbox_event_ids_json,
+ expected_event_id,
+ actor_pubkey,
+ idempotency_digest_prefix,
+ workflow_state,
+ created_at,
+ updated_at,
+ detail_json
+ ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?10, ?11)
+ ON CONFLICT(source_kind, source_record_id)
+ DO UPDATE SET
+ sdk_operation_kind = excluded.sdk_operation_kind,
+ sdk_outbox_event_ids_json = excluded.sdk_outbox_event_ids_json,
+ expected_event_id = excluded.expected_event_id,
+ actor_pubkey = excluded.actor_pubkey,
+ idempotency_digest_prefix = excluded.idempotency_digest_prefix,
+ workflow_state = excluded.workflow_state,
+ updated_at = excluded.updated_at,
+ detail_json = excluded.detail_json",
+ params![
+ receipt_id,
+ input.source_kind.storage_key(),
+ input.source_record_id.as_str(),
+ input.sdk_operation_kind.as_str(),
+ outbox_ids_json.as_str(),
+ input.expected_event_id.as_deref(),
+ input.actor_pubkey.as_deref(),
+ input.idempotency_digest_prefix.as_deref(),
+ input.workflow_state.storage_key(),
+ input.recorded_at.as_str(),
+ detail_json.as_str(),
+ ],
+ )
+ .map_err(|source| AppSqliteError::Query {
+ operation: "record app SDK workflow receipt",
+ source,
+ })?;
+
+ self.load_receipt(input.source_kind, input.source_record_id.as_str())?
+ .ok_or(AppSqliteError::MissingColumn {
+ field: "app_sdk_workflow_receipts.id",
+ })
+ }
+
+ pub fn load_receipt(
+ &self,
+ source_kind: AppSdkWorkflowReceiptSourceKind,
+ source_record_id: &str,
+ ) -> Result<Option<AppSdkStoredWorkflowReceipt>, AppSqliteError> {
+ self.connection
+ .query_row(
+ "SELECT
+ id,
+ source_kind,
+ source_record_id,
+ sdk_operation_kind,
+ sdk_outbox_event_ids_json,
+ expected_event_id,
+ actor_pubkey,
+ idempotency_digest_prefix,
+ workflow_state,
+ created_at,
+ updated_at,
+ detail_json
+ FROM app_sdk_workflow_receipts
+ WHERE source_kind = ?1
+ AND source_record_id = ?2
+ LIMIT 1",
+ params![source_kind.storage_key(), source_record_id],
+ decode_receipt_row,
+ )
+ .optional()
+ .map_err(|source| AppSqliteError::Query {
+ operation: "load app SDK workflow receipt",
+ source,
+ })
+ }
+}
+
+fn decode_receipt_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<AppSdkStoredWorkflowReceipt> {
+ let source_kind: String = row.get(1)?;
+ let outbox_ids_json: String = row.get(4)?;
+ let workflow_state: String = row.get(8)?;
+ let detail_json: String = row.get(11)?;
+ Ok(AppSdkStoredWorkflowReceipt {
+ id: row.get(0)?,
+ source_kind: AppSdkWorkflowReceiptSourceKind::parse(source_kind.as_str())
+ .map_err(decode_app_error)?,
+ source_record_id: row.get(2)?,
+ sdk_operation_kind: row.get(3)?,
+ sdk_outbox_event_ids: serde_json::from_str(outbox_ids_json.as_str()).map_err(|source| {
+ decode_app_error(AppSqliteError::DecodeJson {
+ field: "app_sdk_workflow_receipts.sdk_outbox_event_ids_json",
+ source,
+ })
+ })?,
+ expected_event_id: row.get(5)?,
+ actor_pubkey: row.get(6)?,
+ idempotency_digest_prefix: row.get(7)?,
+ workflow_state: AppSdkWorkflowReceiptState::parse(workflow_state.as_str())
+ .map_err(decode_app_error)?,
+ created_at: row.get(9)?,
+ updated_at: row.get(10)?,
+ detail_json: serde_json::from_str(detail_json.as_str()).map_err(|source| {
+ decode_app_error(AppSqliteError::DecodeJson {
+ field: "app_sdk_workflow_receipts.detail_json",
+ source,
+ })
+ })?,
+ })
+}
+
+fn decode_app_error(error: AppSqliteError) -> rusqlite::Error {
+ rusqlite::Error::FromSqlConversionFailure(0, rusqlite::types::Type::Text, Box::new(error))
+}
+
+#[cfg(test)]
+mod tests {
+ use serde_json::json;
+
+ use crate::{
+ AppSdkWorkflowReceiptInput, AppSdkWorkflowReceiptSourceKind, AppSdkWorkflowReceiptState,
+ AppSqliteStore, DatabaseTarget,
+ };
+
+ #[test]
+ fn workflow_receipts_are_idempotent_by_source_record() {
+ let store = AppSqliteStore::open(DatabaseTarget::InMemory).expect("open app store");
+ let first = store
+ .sdk_workflow_receipt_repository()
+ .record_receipt(&AppSdkWorkflowReceiptInput {
+ source_kind: AppSdkWorkflowReceiptSourceKind::LocalOutbox,
+ source_record_id: "source-record-a".to_owned(),
+ sdk_operation_kind: "farm.publish".to_owned(),
+ sdk_outbox_event_ids: vec!["outbox-a".to_owned()],
+ expected_event_id: Some("expected-a".to_owned()),
+ actor_pubkey: Some("actor-a".to_owned()),
+ idempotency_digest_prefix: Some("digest-a".to_owned()),
+ workflow_state: AppSdkWorkflowReceiptState::Enqueued,
+ recorded_at: "2026-06-18T12:00:00Z".to_owned(),
+ detail_json: json!({"attempt": 1}),
+ })
+ .expect("record first receipt");
+ let second = store
+ .sdk_workflow_receipt_repository()
+ .record_receipt(&AppSdkWorkflowReceiptInput {
+ source_kind: AppSdkWorkflowReceiptSourceKind::LocalOutbox,
+ source_record_id: "source-record-a".to_owned(),
+ sdk_operation_kind: "farm.publish".to_owned(),
+ sdk_outbox_event_ids: vec!["outbox-b".to_owned()],
+ expected_event_id: Some("expected-b".to_owned()),
+ actor_pubkey: Some("actor-b".to_owned()),
+ idempotency_digest_prefix: Some("digest-b".to_owned()),
+ workflow_state: AppSdkWorkflowReceiptState::Pushed,
+ recorded_at: "2026-06-18T12:05:00Z".to_owned(),
+ detail_json: json!({"attempt": 2}),
+ })
+ .expect("record second receipt");
+
+ assert_eq!(first.id, second.id);
+ assert_eq!(second.created_at, "2026-06-18T12:00:00Z");
+ assert_eq!(second.updated_at, "2026-06-18T12:05:00Z");
+ assert_eq!(second.sdk_outbox_event_ids, vec!["outbox-b".to_owned()]);
+ assert_eq!(second.expected_event_id.as_deref(), Some("expected-b"));
+ assert_eq!(second.actor_pubkey.as_deref(), Some("actor-b"));
+ assert_eq!(second.workflow_state, AppSdkWorkflowReceiptState::Pushed);
+ assert_eq!(second.detail_json, json!({"attempt": 2}));
+ }
+}
diff --git a/crates/sync/Cargo.toml b/crates/sync/Cargo.toml
@@ -9,6 +9,7 @@ license.workspace = true
publish = false
[dependencies]
+radroots_events.workspace = true
radroots_studio_app_view.workspace = true
radroots_sdk.workspace = true
serde.workspace = true
diff --git a/crates/sync/src/publish.rs b/crates/sync/src/publish.rs
@@ -1,10 +1,10 @@
-use radroots_sdk::protocol::order::{
+use radroots_events::order::{
RadrootsOrderEconomics, RadrootsOrderItem, RadrootsOrderRevisionOutcome,
};
use radroots_sdk::{
- FARM_PUBLISH_OPERATION_KIND, LISTING_PUBLISH_OPERATION_KIND, ORDER_CANCELLATION_OPERATION_KIND,
- ORDER_DECISION_OPERATION_KIND, ORDER_REVISION_DECISION_OPERATION_KIND,
- ORDER_REVISION_PROPOSAL_OPERATION_KIND, ORDER_SUBMIT_OPERATION_KIND,
+ FARM_PUBLISH_OPERATION_KIND, LISTING_PUBLISH_OPERATION_KIND, TRADE_CANCELLATION_OPERATION_KIND,
+ TRADE_DECISION_OPERATION_KIND, TRADE_REVISION_DECISION_OPERATION_KIND,
+ TRADE_REVISION_PROPOSAL_OPERATION_KIND, TRADE_SUBMIT_OPERATION_KIND,
};
use radroots_studio_app_view::{
FarmId, FarmReadiness, FulfillmentWindowId, OrderId, ProductId, ProductStatus,
@@ -43,11 +43,11 @@ impl AppPublishWorkKind {
match self {
Self::FarmProfile => FARM_PUBLISH_OPERATION_KIND,
Self::Listing => LISTING_PUBLISH_OPERATION_KIND,
- Self::OrderRequest => ORDER_SUBMIT_OPERATION_KIND,
- Self::OrderDecision => ORDER_DECISION_OPERATION_KIND,
- Self::OrderRevisionProposal => ORDER_REVISION_PROPOSAL_OPERATION_KIND,
- Self::OrderRevisionDecision => ORDER_REVISION_DECISION_OPERATION_KIND,
- Self::OrderCancellation => ORDER_CANCELLATION_OPERATION_KIND,
+ Self::OrderRequest => TRADE_SUBMIT_OPERATION_KIND,
+ Self::OrderDecision => TRADE_DECISION_OPERATION_KIND,
+ Self::OrderRevisionProposal => TRADE_REVISION_PROPOSAL_OPERATION_KIND,
+ Self::OrderRevisionDecision => TRADE_REVISION_DECISION_OPERATION_KIND,
+ Self::OrderCancellation => TRADE_CANCELLATION_OPERATION_KIND,
}
}
}
@@ -186,7 +186,6 @@ pub struct AppOrderRevisionProposalPublishPayload {
pub farm_id: FarmId,
pub trade_order_id: String,
pub request_event_id: String,
- pub prev_event_id: String,
pub revision_id: String,
pub listing_addr: String,
pub buyer_pubkey: String,
@@ -203,7 +202,6 @@ pub struct AppOrderRevisionDecisionPublishPayload {
pub farm_id: FarmId,
pub trade_order_id: String,
pub request_event_id: String,
- pub prev_event_id: String,
pub revision_id: String,
pub listing_addr: String,
pub buyer_pubkey: String,
@@ -218,7 +216,6 @@ pub struct AppOrderCancellationPublishPayload {
pub farm_id: FarmId,
pub trade_order_id: String,
pub request_event_id: String,
- pub prev_event_id: String,
pub listing_addr: String,
pub buyer_pubkey: String,
pub seller_pubkey: String,
@@ -433,7 +430,6 @@ impl AppPublishPayload {
&payload.context,
payload.trade_order_id.as_str(),
payload.request_event_id.as_str(),
- payload.prev_event_id.as_str(),
payload.listing_addr.as_str(),
payload.buyer_pubkey.as_str(),
payload.seller_pubkey.as_str(),
@@ -462,7 +458,6 @@ impl AppPublishPayload {
&payload.context,
payload.trade_order_id.as_str(),
payload.request_event_id.as_str(),
- payload.prev_event_id.as_str(),
payload.listing_addr.as_str(),
payload.buyer_pubkey.as_str(),
payload.seller_pubkey.as_str(),
@@ -480,7 +475,6 @@ impl AppPublishPayload {
&payload.context,
payload.trade_order_id.as_str(),
payload.request_event_id.as_str(),
- payload.prev_event_id.as_str(),
payload.listing_addr.as_str(),
payload.buyer_pubkey.as_str(),
payload.seller_pubkey.as_str(),
@@ -515,7 +509,6 @@ fn validate_lifecycle_order_fields(
context: &AppPublishContext,
trade_order_id: &str,
request_event_id: &str,
- prev_event_id: &str,
listing_addr: &str,
buyer_pubkey: &str,
seller_pubkey: &str,
@@ -528,9 +521,6 @@ fn validate_lifecycle_order_fields(
if request_event_id.trim().is_empty() {
failures.push(AppPublishValidationFailure::MissingOrderRequestEventId);
}
- if prev_event_id.trim().is_empty() {
- failures.push(AppPublishValidationFailure::MissingOrderPreviousEventId);
- }
if listing_addr.trim().is_empty() {
failures.push(AppPublishValidationFailure::MissingOrderListingAddress);
}
@@ -570,7 +560,6 @@ pub enum AppPublishValidationFailure {
MissingOrderTotal,
MissingOrderTradeOrderId,
MissingOrderRequestEventId,
- MissingOrderPreviousEventId,
MissingOrderDecisionInventory,
MissingOrderDeclineReason,
MissingOrderRevisionId,
@@ -609,7 +598,6 @@ impl AppPublishValidationFailure {
Self::MissingOrderTotal => "missing_order_total",
Self::MissingOrderTradeOrderId => "missing_order_trade_order_id",
Self::MissingOrderRequestEventId => "missing_order_request_event_id",
- Self::MissingOrderPreviousEventId => "missing_order_previous_event_id",
Self::MissingOrderDecisionInventory => "missing_order_decision_inventory",
Self::MissingOrderDeclineReason => "missing_order_decline_reason",
Self::MissingOrderRevisionId => "missing_order_revision_id",
@@ -674,13 +662,13 @@ mod tests {
AppOrderRequestPublishPayload, AppOrderRevisionDecisionPublishPayload,
AppOrderRevisionProposalPublishPayload, AppPublishContext, AppPublishPayload,
AppPublishValidationFailure, AppPublishWorkKind, FARM_PUBLISH_OPERATION_KIND,
- ORDER_CANCELLATION_OPERATION_KIND, ORDER_DECISION_OPERATION_KIND,
- ORDER_REVISION_DECISION_OPERATION_KIND, ORDER_REVISION_PROPOSAL_OPERATION_KIND,
+ TRADE_CANCELLATION_OPERATION_KIND, TRADE_DECISION_OPERATION_KIND,
+ TRADE_REVISION_DECISION_OPERATION_KIND, TRADE_REVISION_PROPOSAL_OPERATION_KIND,
};
use crate::{
PendingSyncOperation, PendingSyncOperationState, SyncAggregateRef, SyncOperationKind,
};
- use radroots_sdk::protocol::order::{
+ use radroots_events::order::{
RadrootsOrderEconomics, RadrootsOrderItem, RadrootsOrderRevisionOutcome,
};
use radroots_studio_app_view::{FarmId, FarmReadiness, OrderId, ProductId, ProductStatus};
@@ -856,7 +844,7 @@ mod tests {
assert_eq!(payload.work_kind().storage_key(), "order_decision");
assert_eq!(
payload.work_kind().sdk_operation(),
- ORDER_DECISION_OPERATION_KIND
+ TRADE_DECISION_OPERATION_KIND
);
let reason_codes: Vec<&str> = payload
.validation_failures()
@@ -890,7 +878,6 @@ mod tests {
farm_id,
trade_order_id: " ".to_owned(),
request_event_id: String::new(),
- prev_event_id: String::new(),
listing_addr: String::new(),
buyer_pubkey: String::new(),
seller_pubkey: String::new(),
@@ -899,7 +886,7 @@ mod tests {
assert_eq!(
cancellation.work_kind().sdk_operation(),
- ORDER_CANCELLATION_OPERATION_KIND
+ TRADE_CANCELLATION_OPERATION_KIND
);
let cancellation_reason_codes: Vec<&str> = cancellation
@@ -915,7 +902,6 @@ mod tests {
"missing_source",
"missing_order_trade_order_id",
"missing_order_request_event_id",
- "missing_order_previous_event_id",
"missing_order_listing_address",
"missing_order_buyer_pubkey",
"missing_order_seller_pubkey",
@@ -930,7 +916,6 @@ mod tests {
farm_id,
trade_order_id: "order-1".to_owned(),
request_event_id: "request-event-1".to_owned(),
- prev_event_id: "decision-event-1".to_owned(),
listing_addr: "30402:seller:listing".to_owned(),
buyer_pubkey: "buyer".to_owned(),
seller_pubkey: "seller".to_owned(),
@@ -951,7 +936,6 @@ mod tests {
farm_id,
trade_order_id: "order-1".to_owned(),
request_event_id: "request-event-1".to_owned(),
- prev_event_id: "decision-event-1".to_owned(),
listing_addr: "30402:seller:listing".to_owned(),
buyer_pubkey: "buyer".to_owned(),
seller_pubkey: "seller".to_owned(),
@@ -972,7 +956,6 @@ mod tests {
farm_id,
trade_order_id: "order-1".to_owned(),
request_event_id: "request-event-1".to_owned(),
- prev_event_id: "decision-event-1".to_owned(),
revision_id: "revision-1".to_owned(),
listing_addr: "30402:seller:listing".to_owned(),
buyer_pubkey: "buyer".to_owned(),
@@ -991,7 +974,6 @@ mod tests {
farm_id,
trade_order_id: " ".to_owned(),
request_event_id: String::new(),
- prev_event_id: String::new(),
revision_id: String::new(),
listing_addr: String::new(),
buyer_pubkey: String::new(),
@@ -1007,7 +989,6 @@ mod tests {
farm_id,
trade_order_id: " ".to_owned(),
request_event_id: String::new(),
- prev_event_id: String::new(),
revision_id: String::new(),
listing_addr: String::new(),
buyer_pubkey: String::new(),
@@ -1019,12 +1000,12 @@ mod tests {
assert_eq!(
valid_proposal.work_kind().sdk_operation(),
- ORDER_REVISION_PROPOSAL_OPERATION_KIND
+ TRADE_REVISION_PROPOSAL_OPERATION_KIND
);
assert_eq!(valid_proposal.validation_failures(), Vec::new());
assert_eq!(
invalid_decision.work_kind().sdk_operation(),
- ORDER_REVISION_DECISION_OPERATION_KIND
+ TRADE_REVISION_DECISION_OPERATION_KIND
);
let proposal_reason_codes: Vec<&str> = invalid_proposal
@@ -1045,7 +1026,6 @@ mod tests {
"missing_source",
"missing_order_trade_order_id",
"missing_order_request_event_id",
- "missing_order_previous_event_id",
"missing_order_listing_address",
"missing_order_buyer_pubkey",
"missing_order_seller_pubkey",
@@ -1061,7 +1041,6 @@ mod tests {
"missing_source",
"missing_order_trade_order_id",
"missing_order_request_event_id",
- "missing_order_previous_event_id",
"missing_order_listing_address",
"missing_order_buyer_pubkey",
"missing_order_seller_pubkey",