app

Local-first trade for farms and co-ops
git clone https://radroots.dev/git/app.git
Log | Files | Refs | README | LICENSE

commit 5fae5fe5bef662d4b41551a2a35b70d60a251b1f
parent 3470bf906340ff522d0972180bf82f14c7434550
Author: triesap <tyson@radroots.org>
Date:   Wed, 15 Jul 2026 20:30:11 +0000

runtime: adopt desktop runtime supervisor

- replace AppSdkRuntime with nonblocking DesktopRuntimeSupervisor effects

- move workflow receipt storage to desktop runtime receipt tables

- align desktop summaries, source guards, and i18n metadata with V1 stores

Diffstat:
Mcrates/desktop/src/runtime.rs | 762+++++++++++++++++++++++++++++++++++++------------------------------------------
Mcrates/desktop/src/source_guards.rs | 86+++++++++++++++++++++++++++++--------------------------------------------------
Mcrates/desktop/src/window.rs | 124+++++++++++++++++++++++++++++++++++++++++++------------------------------------
Mcrates/runtime/src/lib.rs | 30++++++++++++++++++------------
Mcrates/runtime/src/sdk.rs | 2069++++++++++++++++++++++++++++++++++++++++++++++---------------------------------
Acrates/store/migrations/0028_runtime_workflow_receipts.sql | 38++++++++++++++++++++++++++++++++++++++
Dcrates/store/migrations/0028_sdk_workflow_receipts.sql | 38--------------------------------------
Mcrates/store/src/lib.rs | 40+++++++++++++++++++++++-----------------
Mcrates/store/src/migrations.rs | 2+-
Acrates/store/src/runtime_workflow_receipts.rs | 323+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Dcrates/store/src/sdk_workflow_receipts.rs | 320-------------------------------------------------------------------------------
Mi18n/locales/en/messages.json | 5+++--
12 files changed, 2071 insertions(+), 1766 deletions(-)

diff --git a/crates/desktop/src/runtime.rs b/crates/desktop/src/runtime.rs @@ -57,13 +57,15 @@ use radroots_sdk::{ use radroots_sql_core::SqlxSqliteExecutor; use radroots_studio_app_core::{ AppBuildIdentity, AppDesktopRuntimePaths, AppRuntimeCapture, AppRuntimeMode, - AppRuntimePathsError, AppRuntimeSnapshot, AppSdkConfig, AppSdkDiagnostics, - AppSdkFarmPublicLocationRequest, AppSdkFarmPublishRequest, AppSdkLifecycleState, - AppSdkListingPublishRequest, AppSdkProjectionLifecycleState, AppSdkPublicFarmLocation, - AppSdkRelayUrlPolicy, AppSdkRuntime, AppSdkRuntimeError, AppSdkRuntimeIssue, - AppSdkRuntimeStatus, AppSdkStoragePaths, AppSdkTradeCancellationRequest, AppSdkTradeDecision, - AppSdkTradeDecisionRequest, AppSdkTradeProposeRequest, AppSdkWorkflowReceipt, - AppSharedAccountsPaths, PackDayExportWriteError, prepare_pack_day_export_bundle_at_data_root, + AppRuntimePathsError, AppRuntimeSnapshot, AppSharedAccountsPaths, DesktopRuntimeDiagnostics, + DesktopRuntimeEffectReceipt, DesktopRuntimeFarmPublishRequest, DesktopRuntimeIssue, + DesktopRuntimeLifecycleState, DesktopRuntimeListingPublishRequest, DesktopRuntimeLocalSigner, + DesktopRuntimeProjectionLifecycleState, DesktopRuntimePublicFarmLocation, + DesktopRuntimeRelayUrlPolicy, DesktopRuntimeSnapshot, DesktopRuntimeSupervisor, + DesktopRuntimeSupervisorConfig, DesktopRuntimeSupervisorError, + DesktopRuntimeTradeCancellationRequest, DesktopRuntimeTradeDecision, + DesktopRuntimeTradeDecisionRequest, DesktopRuntimeTradeProposeRequest, PackDayExportWriteError, + prepare_pack_day_export_bundle_at_data_root, shared_runtime_store_database_path_from_shared_accounts, write_prepared_pack_day_export_bundle, }; use radroots_studio_app_remote_signer::{ @@ -71,12 +73,12 @@ use radroots_studio_app_remote_signer::{ }; use radroots_studio_app_sqlite::{ APP_ACTIVITY_CONTEXT_LIMIT, AppLocalInteropImportReport, AppRelayIngestFailureInput, - AppRelayIngestSuccessInput, AppSdkWorkflowReceiptInput, AppSdkWorkflowReceiptSourceKind, - AppSdkWorkflowReceiptState, AppSqliteError, AppSqliteStore, BuyerOrderRuntimeStoreExport, + AppRelayIngestSuccessInput, AppSqliteError, AppSqliteStore, BuyerOrderRuntimeStoreExport, BuyerOrderRuntimeStoreLine, BuyerRepeatDemandApplyOutcome, DatabaseTarget, - SelectedBuyerOrderScope, SellerOrderDecisionExport, StoredPendingSyncOperation, - StoredRelayIngestCursor, StoredSyncConflict, derive_farm_rules_readiness, - projected_order_id_from_trade_request, + DesktopRuntimeWorkflowReceiptInput, DesktopRuntimeWorkflowReceiptSourceKind, + DesktopRuntimeWorkflowReceiptState, SelectedBuyerOrderScope, SellerOrderDecisionExport, + StoredPendingSyncOperation, StoredRelayIngestCursor, StoredSyncConflict, + derive_farm_rules_readiness, projected_order_id_from_trade_request, }; use radroots_studio_app_state::{ APP_STATE_FILE_NAME, AppShellProjection, AppStateCommand, AppStatePersistenceRepository, @@ -151,7 +153,8 @@ use crate::remote_signer::{ const APP_DATABASE_FILE_NAME: &str = "app.sqlite3"; const SYNC_TRANSPORT_UNAVAILABLE_MESSAGE: &str = "remote sync transport is not configured"; -const APP_SYNC_PUBLISH_USES_SDK_RUNTIME_MESSAGE: &str = "app sync publish work uses AppSdkRuntime"; +const APP_SYNC_PUBLISH_USES_SDK_RUNTIME_MESSAGE: &str = + "app sync publish work uses DesktopRuntimeSupervisor"; const APP_DIRECT_RELAY_SYNC_TIMEOUT_MS: u64 = 2_000; const APP_DIRECT_RELAY_CONNECT_TIMEOUT: StdDuration = StdDuration::from_secs(10); const APP_DIRECT_RELAY_INGEST_LIMIT: usize = 1_000; @@ -361,7 +364,7 @@ impl AppSyncTransport for ConfiguredRelayAppSyncTransport { #[derive(Clone, Debug)] pub struct DesktopAppRuntime { state: Arc<Mutex<DesktopAppRuntimeState>>, - sdk_runtime: Arc<Mutex<Option<AppSdkRuntime>>>, + sdk_runtime: Arc<Mutex<Option<DesktopRuntimeSupervisor>>>, } impl DesktopAppRuntime { @@ -409,7 +412,6 @@ impl DesktopAppRuntime { match start_desktop_sdk_runtime(&paths, nostr_relay_urls) { Ok(sdk_runtime) => { let runtime = Self::from_state_with_sdk_runtime(state, sdk_runtime); - let _ = runtime.wait_for_sdk_startup(StdDuration::from_secs(5)); if let Err(error) = runtime.retry_pending_personal_order_coordination() { error!( target: "buyer_order", @@ -469,72 +471,69 @@ impl DesktopAppRuntime { self.lock_state().nostr_relay_urls.clone() } - pub fn sdk_status(&self) -> Option<AppSdkRuntimeStatus> { + pub fn sdk_status(&self) -> Option<DesktopRuntimeSnapshot> { self.sdk_runtime .lock() .unwrap_or_else(|poisoned| poisoned.into_inner()) .as_ref() - .map(AppSdkRuntime::status) + .map(DesktopRuntimeSupervisor::snapshot) } - pub fn sdk_status_summary(&self) -> Option<DesktopAppSdkStatusSummary> { + pub fn sdk_status_summary(&self) -> Option<DesktopRuntimeSupervisorStatusSummary> { self.sdk_status() .as_ref() - .map(DesktopAppSdkStatusSummary::from_status) + .map(DesktopRuntimeSupervisorStatusSummary::from_status) } - pub fn wait_for_sdk_startup(&self, timeout: StdDuration) -> Option<AppSdkRuntimeStatus> { - self.sdk_runtime - .lock() - .unwrap_or_else(|poisoned| poisoned.into_inner()) - .as_ref() - .map(|runtime| runtime.wait_for_startup(timeout)) - } - - pub fn shutdown_sdk_runtime(&self) -> Result<bool, AppSdkRuntimeError> { + pub fn shutdown_sdk_runtime(&self) -> bool { let mut sdk_runtime = self .sdk_runtime .lock() .unwrap_or_else(|poisoned| poisoned.into_inner()); let Some(runtime) = sdk_runtime.take() else { - return Ok(false); + return false; }; - runtime.shutdown()?; - Ok(true) + runtime.request_shutdown() } - pub fn sdk_diagnostics(&self) -> Result<Option<AppSdkDiagnostics>, AppSdkRuntimeError> { + pub fn sdk_diagnostics(&self) -> Option<DesktopRuntimeDiagnostics> { let sdk_runtime = self .sdk_runtime .lock() .unwrap_or_else(|poisoned| poisoned.into_inner()); let Some(runtime) = sdk_runtime.as_ref() else { - return Ok(None); + return None; }; - runtime.diagnostics().map(Some) + DesktopRuntimeDiagnostics::from_snapshot(runtime.snapshot()) } - pub fn sdk_diagnostics_summary(&self) -> Option<DesktopAppSdkDiagnosticsSummary> { + pub fn sdk_diagnostics_summary(&self) -> Option<DesktopRuntimeSupervisorDiagnosticsSummary> { let sdk_runtime = self .sdk_runtime .lock() .unwrap_or_else(|poisoned| poisoned.into_inner()); let runtime = sdk_runtime.as_ref()?; - let status = runtime.status(); - match runtime.diagnostics() { - Ok(diagnostics) => Some(DesktopAppSdkDiagnosticsSummary { - status: DesktopAppSdkStatusSummary::from_status(&diagnostics.runtime), - state: DesktopAppSdkDiagnosticsState::Ready( - DesktopAppSdkReadyDiagnosticsSummary::from_diagnostics(&diagnostics), + let status = runtime.snapshot(); + match DesktopRuntimeDiagnostics::from_snapshot(status.clone()) { + Some(diagnostics) => Some(DesktopRuntimeSupervisorDiagnosticsSummary { + status: DesktopRuntimeSupervisorStatusSummary::from_status(&diagnostics.runtime), + state: DesktopRuntimeSupervisorDiagnosticsState::Ready( + DesktopRuntimeSupervisorReadyDiagnosticsSummary::from_diagnostics(&diagnostics), ), }), - Err(error) => { - let issue = desktop_app_sdk_issue_from_runtime_error(&error); - let mut status = DesktopAppSdkStatusSummary::from_status(&status); + None => { + let mut status = DesktopRuntimeSupervisorStatusSummary::from_status(&status); + let issue = status.last_issue.clone().unwrap_or_else(|| { + DesktopRuntimeSupervisorIssueSummary::runtime( + "desktop_runtime_diagnostics_pending", + true, + ["wait_for_runtime_snapshot"], + ) + }); status.last_issue = Some(issue.clone()); - Some(DesktopAppSdkDiagnosticsSummary { + Some(DesktopRuntimeSupervisorDiagnosticsSummary { status, - state: DesktopAppSdkDiagnosticsState::Blocked(issue), + state: DesktopRuntimeSupervisorDiagnosticsState::Blocked(issue), }) } } @@ -1087,7 +1086,7 @@ impl DesktopAppRuntime { fn from_state_with_sdk_runtime( mut state: DesktopAppRuntimeState, - sdk_runtime: AppSdkRuntime, + sdk_runtime: DesktopRuntimeSupervisor, ) -> Self { let sdk_runtime = Arc::new(Mutex::new(Some(sdk_runtime))); state.sdk_runtime = Some(Arc::clone(&sdk_runtime)); @@ -1144,43 +1143,11 @@ fn default_runtime_snapshot() -> AppRuntimeSnapshot { fn start_desktop_sdk_runtime( paths: &AppDesktopRuntimePaths, nostr_relay_urls: Vec<String>, -) -> Result<AppSdkRuntime, AppSdkRuntimeError> { - AppSdkRuntime::start(AppSdkConfig::from_desktop_paths(paths, nostr_relay_urls)) -} - -fn desktop_app_sdk_issue_from_runtime_error( - error: &AppSdkRuntimeError, -) -> DesktopAppSdkIssueSummary { - match error { - AppSdkRuntimeError::CommandFailed(issue) => DesktopAppSdkIssueSummary::from_issue(issue), - AppSdkRuntimeError::CommandQueueCapacityZero => DesktopAppSdkIssueSummary::runtime( - "sdk_command_queue_capacity_zero", - false, - ["review_runtime_configuration"], - ), - AppSdkRuntimeError::WorkerSpawn(_) => { - DesktopAppSdkIssueSummary::runtime("sdk_worker_spawn_failed", true, ["retry_startup"]) - } - AppSdkRuntimeError::CommandQueueFull => DesktopAppSdkIssueSummary::runtime( - "sdk_command_queue_full", - true, - ["retry_status_refresh"], - ), - AppSdkRuntimeError::CommandQueueClosed => { - DesktopAppSdkIssueSummary::runtime("sdk_command_queue_closed", true, ["retry_startup"]) - } - AppSdkRuntimeError::CommandResponseClosed => DesktopAppSdkIssueSummary::runtime( - "sdk_command_response_closed", - true, - ["retry_status_refresh"], - ), - AppSdkRuntimeError::ShutdownAck => { - DesktopAppSdkIssueSummary::runtime("sdk_shutdown_ack_failed", true, ["retry_startup"]) - } - AppSdkRuntimeError::WorkerJoin => { - DesktopAppSdkIssueSummary::runtime("sdk_worker_join_failed", true, ["retry_startup"]) - } - } +) -> Result<DesktopRuntimeSupervisor, DesktopRuntimeSupervisorError> { + DesktopRuntimeSupervisor::start(DesktopRuntimeSupervisorConfig::from_desktop_paths( + paths, + nostr_relay_urls, + )) } #[derive(Clone, Debug, Default, Eq, PartialEq)] @@ -1204,21 +1171,21 @@ pub struct DesktopAppSyncConflictSummary { } #[derive(Clone, Debug, Eq, PartialEq)] -pub struct DesktopAppSdkStatusSummary { - pub lifecycle_state: AppSdkLifecycleState, - pub projection_lifecycle_state: AppSdkProjectionLifecycleState, +pub struct DesktopRuntimeSupervisorStatusSummary { + pub lifecycle_state: DesktopRuntimeLifecycleState, + pub projection_lifecycle_state: DesktopRuntimeProjectionLifecycleState, pub projection_lifecycle_reason: Option<String>, pub storage_root: PathBuf, pub runtime_path: Option<PathBuf>, pub private_path: Option<PathBuf>, pub studio_path: Option<PathBuf>, pub relay_target_count: usize, - pub relay_url_policy: AppSdkRelayUrlPolicy, - pub last_issue: Option<DesktopAppSdkIssueSummary>, + pub relay_url_policy: DesktopRuntimeRelayUrlPolicy, + pub last_issue: Option<DesktopRuntimeSupervisorIssueSummary>, } -impl DesktopAppSdkStatusSummary { - fn from_status(status: &AppSdkRuntimeStatus) -> Self { +impl DesktopRuntimeSupervisorStatusSummary { + fn from_status(status: &DesktopRuntimeSnapshot) -> Self { let runtime_path = status .storage_paths .as_ref() @@ -1244,25 +1211,25 @@ impl DesktopAppSdkStatusSummary { last_issue: status .last_issue .as_ref() - .map(DesktopAppSdkIssueSummary::from_issue), + .map(DesktopRuntimeSupervisorIssueSummary::from_issue), } } } #[derive(Clone, Debug, Eq, PartialEq)] -pub struct DesktopAppSdkDiagnosticsSummary { - pub status: DesktopAppSdkStatusSummary, - pub state: DesktopAppSdkDiagnosticsState, +pub struct DesktopRuntimeSupervisorDiagnosticsSummary { + pub status: DesktopRuntimeSupervisorStatusSummary, + pub state: DesktopRuntimeSupervisorDiagnosticsState, } #[derive(Clone, Debug, Eq, PartialEq)] -pub enum DesktopAppSdkDiagnosticsState { - Ready(DesktopAppSdkReadyDiagnosticsSummary), - Blocked(DesktopAppSdkIssueSummary), +pub enum DesktopRuntimeSupervisorDiagnosticsState { + Ready(DesktopRuntimeSupervisorReadyDiagnosticsSummary), + Blocked(DesktopRuntimeSupervisorIssueSummary), } #[derive(Clone, Debug, Eq, PartialEq)] -pub struct DesktopAppSdkReadyDiagnosticsSummary { +pub struct DesktopRuntimeSupervisorReadyDiagnosticsSummary { pub storage_kind: String, pub event_store_total_events: i64, pub outbox_total_events: i64, @@ -1275,8 +1242,8 @@ pub struct DesktopAppSdkReadyDiagnosticsSummary { pub sync_relay_target_count: usize, } -impl DesktopAppSdkReadyDiagnosticsSummary { - fn from_diagnostics(diagnostics: &AppSdkDiagnostics) -> Self { +impl DesktopRuntimeSupervisorReadyDiagnosticsSummary { + fn from_diagnostics(diagnostics: &DesktopRuntimeDiagnostics) -> Self { Self { storage_kind: diagnostics.storage.storage_kind.clone(), event_store_total_events: diagnostics.storage.event_store.total_events, @@ -1297,15 +1264,15 @@ impl DesktopAppSdkReadyDiagnosticsSummary { } #[derive(Clone, Debug, Eq, PartialEq)] -pub struct DesktopAppSdkIssueSummary { +pub struct DesktopRuntimeSupervisorIssueSummary { pub code: String, pub class: String, pub retryable: bool, pub recovery_actions: Vec<String>, } -impl DesktopAppSdkIssueSummary { - fn from_issue(issue: &AppSdkRuntimeIssue) -> Self { +impl DesktopRuntimeSupervisorIssueSummary { + fn from_issue(issue: &DesktopRuntimeIssue) -> Self { Self { code: issue.code.clone(), class: issue.class.clone(), @@ -1389,7 +1356,7 @@ pub struct DesktopAppRuntimeSummary { pub runtime_metadata: DesktopAppRuntimeMetadataSummary, pub sync_status: DesktopAppSyncStatusSummary, pub startup_issue: Option<String>, - pub sdk_status: Option<DesktopAppSdkStatusSummary>, + pub sdk_status: Option<DesktopRuntimeSupervisorStatusSummary>, } #[derive(Debug, Error)] @@ -1458,7 +1425,7 @@ struct DesktopAppRuntimeState { remote_signer_paths: Option<DesktopRemoteSignerPaths>, accounts_manager: Option<RadrootsNostrAccountsManager>, sqlite_store: Option<AppSqliteStore>, - sdk_runtime: Option<Arc<Mutex<Option<AppSdkRuntime>>>>, + sdk_runtime: Option<Arc<Mutex<Option<DesktopRuntimeSupervisor>>>>, sync_transport: Box<dyn AppSyncTransport + Send>, runtime_metadata: DesktopAppRuntimeMetadataSummary, selected_account_pending_sync_write_count: usize, @@ -2645,7 +2612,7 @@ impl DesktopAppRuntimeState { let source_record_id = order_decision_sdk_source_record_id(&payload); self.enqueue_order_decision_payload_via_sdk( &payload, - AppSdkWorkflowReceiptSourceKind::LocalOutbox, + DesktopRuntimeWorkflowReceiptSourceKind::AppWorkflow, source_record_id.as_str(), )?; let _ = self.refresh_selected_account_sync()?; @@ -2762,7 +2729,7 @@ impl DesktopAppRuntimeState { let source_record_id = order_cancellation_sdk_source_record_id(&payload); self.enqueue_order_cancellation_payload_via_sdk( &payload, - AppSdkWorkflowReceiptSourceKind::LocalOutbox, + DesktopRuntimeWorkflowReceiptSourceKind::AppWorkflow, source_record_id.as_str(), )?; let _ = self.refresh_selected_account_sync()?; @@ -3264,7 +3231,7 @@ impl DesktopAppRuntimeState { farm_id: account .farmer_activation .farm_id - .unwrap_or_else(FarmId::new), + .unwrap_or_else(FarmId::generate), display_name: draft.farm_name.trim().to_owned(), readiness: FarmReadiness::Incomplete, }; @@ -4126,7 +4093,7 @@ impl DesktopAppRuntimeState { .unwrap_or_else(|| format!("app:order_request:{}", payload.order_id)); self.enqueue_order_request_payload_via_sdk( &payload, - AppSdkWorkflowReceiptSourceKind::SharedRuntimeStore, + DesktopRuntimeWorkflowReceiptSourceKind::SharedRuntimeStore, source_record_id.as_str(), )?; self.refresh_selected_account_sync() @@ -4370,7 +4337,7 @@ impl DesktopAppRuntimeState { fn enqueue_farm_profile_payload_via_sdk( &self, payload: &AppFarmProfilePublishPayload, - source_kind: AppSdkWorkflowReceiptSourceKind, + source_kind: DesktopRuntimeWorkflowReceiptSourceKind, source_record_id: &str, ) -> Result<(), AppSqliteError> { let operation_kind = FARM_PUBLISH_OPERATION_KIND; @@ -4380,16 +4347,13 @@ impl DesktopAppRuntimeState { )) .and_then(|identity| { let actor_pubkey = identity.public_key_hex(); - let public_location = - self.sdk_public_farm_location(actor_pubkey.as_str(), payload.farm_id)?; - let request = AppSdkFarmPublishRequest { + let request = DesktopRuntimeFarmPublishRequest { actor_account_id: payload.context.account_id.clone(), actor_pubkey: actor_pubkey.clone(), - signer_keys: identity.into_keys(), - farm: farm_profile_publish_payload_to_sdk_farm( - payload, - public_location.as_ref(), + signer: DesktopRuntimeLocalSigner::from_local_identity_keys( + identity.into_keys(), ), + farm: farm_profile_publish_payload_to_sdk_farm(payload, None), target_relays: normalized_app_sync_relay_urls(&self.nostr_relay_urls)?, relay_url_policy: sdk_relay_url_policy_for_targets(&self.nostr_relay_urls), idempotency_key: Some(sdk_idempotency_key(source_record_id)), @@ -4399,14 +4363,14 @@ impl DesktopAppRuntimeState { .map_err(sync_transport_error_from_sdk_runtime_error) }); match actor_pubkey { - Ok((actor_pubkey, receipt)) => self.record_app_sdk_workflow_success( + Ok((actor_pubkey, receipt)) => self.record_desktop_runtime_workflow_success( source_kind, source_record_id, operation_kind, actor_pubkey.as_str(), &receipt, ), - Err(error) => self.record_app_sdk_workflow_failure( + Err(error) => self.record_desktop_runtime_workflow_failure( source_kind, source_record_id, operation_kind, @@ -4419,7 +4383,7 @@ impl DesktopAppRuntimeState { fn enqueue_listing_payload_via_sdk( &self, payload: &AppListingPublishPayload, - source_kind: AppSdkWorkflowReceiptSourceKind, + source_kind: DesktopRuntimeWorkflowReceiptSourceKind, source_record_id: &str, ) -> Result<(), DesktopAppRuntimeProductPublishError> { let operation_kind = LISTING_PUBLISH_OPERATION_KIND; @@ -4429,20 +4393,13 @@ impl DesktopAppRuntimeState { )) .and_then(|identity| { let actor_pubkey = identity.public_key_hex(); - let public_location = match payload.farm_id { - Some(farm_id) => { - self.sdk_public_farm_location(actor_pubkey.as_str(), farm_id)? - } - None => None, - }; - let request = AppSdkListingPublishRequest { + let request = DesktopRuntimeListingPublishRequest { actor_account_id: payload.context.account_id.clone(), actor_pubkey: actor_pubkey.clone(), - signer_keys: identity.into_keys(), - listing: listing_publish_payload_to_sdk_listing( - payload, - public_location.as_ref(), - )?, + signer: DesktopRuntimeLocalSigner::from_local_identity_keys( + identity.into_keys(), + ), + listing: listing_publish_payload_to_sdk_listing(payload, None)?, target_relays: normalized_app_sync_relay_urls(&self.nostr_relay_urls)?, relay_url_policy: sdk_relay_url_policy_for_targets(&self.nostr_relay_urls), idempotency_key: Some(sdk_idempotency_key(source_record_id)), @@ -4453,7 +4410,7 @@ impl DesktopAppRuntimeState { }); match actor_pubkey { Ok((actor_pubkey, receipt)) => self - .record_app_sdk_workflow_success( + .record_desktop_runtime_workflow_success( source_kind, source_record_id, operation_kind, @@ -4462,7 +4419,7 @@ impl DesktopAppRuntimeState { ) .map_err(DesktopAppRuntimeProductPublishError::from), Err(error) => { - self.record_app_sdk_workflow_failure( + self.record_desktop_runtime_workflow_failure( source_kind, source_record_id, operation_kind, @@ -4477,7 +4434,7 @@ impl DesktopAppRuntimeState { fn enqueue_order_request_payload_via_sdk( &self, payload: &AppOrderRequestPublishPayload, - source_kind: AppSdkWorkflowReceiptSourceKind, + source_kind: DesktopRuntimeWorkflowReceiptSourceKind, source_record_id: &str, ) -> Result<(), AppSqliteError> { let operation_kind = TRADE_SUBMIT_OPERATION_KIND; @@ -4488,10 +4445,12 @@ impl DesktopAppRuntimeState { .and_then(|identity| { let actor_pubkey = identity.public_key_hex(); let order_parts = order_request_publish_payload_to_sdk_product_parts(payload)?; - let request = AppSdkTradeProposeRequest { + let request = DesktopRuntimeTradeProposeRequest { actor_account_id: payload.context.account_id.clone(), actor_pubkey: actor_pubkey.clone(), - signer_keys: identity.into_keys(), + signer: DesktopRuntimeLocalSigner::from_local_identity_keys( + identity.into_keys(), + ), listing_event: order_request_sdk_listing_event_ptr(payload)?, order_id: order_parts.order_id, listing_addr: order_parts.listing_addr, @@ -4507,14 +4466,14 @@ impl DesktopAppRuntimeState { .map_err(sync_transport_error_from_sdk_runtime_error) }); match actor_pubkey { - Ok((actor_pubkey, receipt)) => self.record_app_sdk_workflow_success( + Ok((actor_pubkey, receipt)) => self.record_desktop_runtime_workflow_success( source_kind, source_record_id, operation_kind, actor_pubkey.as_str(), &receipt, ), - Err(error) => self.record_app_sdk_workflow_failure( + Err(error) => self.record_desktop_runtime_workflow_failure( source_kind, source_record_id, operation_kind, @@ -4527,7 +4486,7 @@ impl DesktopAppRuntimeState { fn enqueue_order_decision_payload_via_sdk( &self, payload: &AppOrderDecisionPublishPayload, - source_kind: AppSdkWorkflowReceiptSourceKind, + source_kind: DesktopRuntimeWorkflowReceiptSourceKind, source_record_id: &str, ) -> Result<(), AppSqliteError> { let operation_kind = TRADE_DECISION_OPERATION_KIND; @@ -4537,10 +4496,12 @@ impl DesktopAppRuntimeState { )) .and_then(|identity| { let actor_pubkey = identity.public_key_hex(); - let request = AppSdkTradeDecisionRequest { + let request = DesktopRuntimeTradeDecisionRequest { actor_account_id: payload.context.account_id.clone(), actor_pubkey: actor_pubkey.clone(), - signer_keys: identity.into_keys(), + signer: DesktopRuntimeLocalSigner::from_local_identity_keys( + identity.into_keys(), + ), locator: trade_locator_from_decision_payload(payload)?, decision: trade_decision_from_publish_payload(payload)?, confirm_public_note: payload.confirm_public_note, @@ -4551,14 +4512,14 @@ impl DesktopAppRuntimeState { .map_err(sync_transport_error_from_sdk_runtime_error) }); match actor_pubkey { - Ok((actor_pubkey, receipt)) => self.record_app_sdk_workflow_success( + Ok((actor_pubkey, receipt)) => self.record_desktop_runtime_workflow_success( source_kind, source_record_id, operation_kind, actor_pubkey.as_str(), &receipt, ), - Err(error) => self.record_app_sdk_workflow_failure( + Err(error) => self.record_desktop_runtime_workflow_failure( source_kind, source_record_id, operation_kind, @@ -4571,7 +4532,7 @@ impl DesktopAppRuntimeState { fn enqueue_order_cancellation_payload_via_sdk( &self, payload: &AppOrderCancellationPublishPayload, - source_kind: AppSdkWorkflowReceiptSourceKind, + source_kind: DesktopRuntimeWorkflowReceiptSourceKind, source_record_id: &str, ) -> Result<(), AppSqliteError> { let operation_kind = TRADE_CANCELLATION_OPERATION_KIND; @@ -4581,10 +4542,12 @@ impl DesktopAppRuntimeState { )) .and_then(|identity| { let actor_pubkey = identity.public_key_hex(); - let request = AppSdkTradeCancellationRequest { + let request = DesktopRuntimeTradeCancellationRequest { actor_account_id: payload.context.account_id.clone(), actor_pubkey: actor_pubkey.clone(), - signer_keys: identity.into_keys(), + signer: DesktopRuntimeLocalSigner::from_local_identity_keys( + identity.into_keys(), + ), locator: trade_locator_from_cancellation_payload(payload)?, reason: payload.reason.clone(), confirm_public_note: payload.confirm_public_note, @@ -4595,14 +4558,14 @@ impl DesktopAppRuntimeState { .map_err(sync_transport_error_from_sdk_runtime_error) }); match actor_pubkey { - Ok((actor_pubkey, receipt)) => self.record_app_sdk_workflow_success( + Ok((actor_pubkey, receipt)) => self.record_desktop_runtime_workflow_success( source_kind, source_record_id, operation_kind, actor_pubkey.as_str(), &receipt, ), - Err(error) => self.record_app_sdk_workflow_failure( + Err(error) => self.record_desktop_runtime_workflow_failure( source_kind, source_record_id, operation_kind, @@ -4622,58 +4585,45 @@ impl DesktopAppRuntimeState { signing_identity_for_publish_payload(accounts_manager, payload) } - fn sdk_public_farm_location( - &self, - actor_pubkey: &str, - farm_id: FarmId, - ) -> Result<Option<AppSdkPublicFarmLocation>, AppSyncTransportError> { - let request = AppSdkFarmPublicLocationRequest { - actor_pubkey: actor_pubkey.to_owned(), - farm_d_tag: d_tag_from_uuid(farm_id.as_uuid()), - }; - self.with_app_sdk_runtime(|runtime| runtime.farm_public_location(request)) - .map_err(sync_transport_error_from_sdk_runtime_error) - } - fn enqueue_app_sdk_farm_publish( &self, - request: AppSdkFarmPublishRequest, - ) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeError> { + request: DesktopRuntimeFarmPublishRequest, + ) -> Result<DesktopRuntimeEffectReceipt, DesktopRuntimeSupervisorError> { self.with_app_sdk_runtime(|runtime| runtime.enqueue_farm_publish(request)) } fn enqueue_app_sdk_listing_publish( &self, - request: AppSdkListingPublishRequest, - ) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeError> { + request: DesktopRuntimeListingPublishRequest, + ) -> Result<DesktopRuntimeEffectReceipt, DesktopRuntimeSupervisorError> { self.with_app_sdk_runtime(|runtime| runtime.enqueue_listing_publish(request)) } fn enqueue_app_sdk_trade_propose( &self, - request: AppSdkTradeProposeRequest, - ) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeError> { + request: DesktopRuntimeTradeProposeRequest, + ) -> Result<DesktopRuntimeEffectReceipt, DesktopRuntimeSupervisorError> { self.with_app_sdk_runtime(|runtime| runtime.trade_propose(request)) } fn enqueue_app_sdk_trade_decision( &self, - request: AppSdkTradeDecisionRequest, - ) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeError> { + request: DesktopRuntimeTradeDecisionRequest, + ) -> Result<DesktopRuntimeEffectReceipt, DesktopRuntimeSupervisorError> { self.with_app_sdk_runtime(|runtime| runtime.trade_decide(request)) } fn enqueue_app_sdk_trade_cancellation( &self, - request: AppSdkTradeCancellationRequest, - ) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeError> { + request: DesktopRuntimeTradeCancellationRequest, + ) -> Result<DesktopRuntimeEffectReceipt, DesktopRuntimeSupervisorError> { self.with_app_sdk_runtime(|runtime| runtime.trade_cancel(request)) } fn with_app_sdk_runtime<T>( &self, - command: impl FnOnce(&AppSdkRuntime) -> Result<T, AppSdkRuntimeError>, - ) -> Result<T, AppSdkRuntimeError> { + command: impl FnOnce(&DesktopRuntimeSupervisor) -> Result<T, DesktopRuntimeSupervisorError>, + ) -> Result<T, DesktopRuntimeSupervisorError> { let Some(handle) = self.sdk_runtime.as_ref() else { return Err(sdk_runtime_unavailable_error()); }; @@ -4684,67 +4634,66 @@ impl DesktopAppRuntimeState { command(runtime) } - fn record_app_sdk_workflow_success( + fn record_desktop_runtime_workflow_success( &self, - source_kind: AppSdkWorkflowReceiptSourceKind, + source_kind: DesktopRuntimeWorkflowReceiptSourceKind, source_record_id: &str, operation_kind: &str, actor_pubkey: &str, - receipt: &AppSdkWorkflowReceipt, + receipt: &DesktopRuntimeEffectReceipt, ) -> Result<(), AppSqliteError> { let detail_json = json!({ + "effect_id": receipt.effect_id, + "effect_kind": format!("{:?}", receipt.effect_kind), "operation_kind": receipt.operation_kind, - "expected_event_id": receipt.expected_event_id, - "signed_event_id": receipt.signed_event_id, - "outbox_operation_id": receipt.outbox_operation_id, - "outbox_event_id": receipt.outbox_event_id, - "state": receipt.state, + "actor_pubkey": receipt.actor_pubkey, + "state": "accepted", }); - self.record_app_sdk_workflow_receipt(AppSdkWorkflowReceiptInput { + self.record_desktop_runtime_workflow_receipt(DesktopRuntimeWorkflowReceiptInput { source_kind, source_record_id: source_record_id.to_owned(), sdk_operation_kind: operation_kind.to_owned(), - sdk_outbox_event_ids: vec![receipt.outbox_event_id.to_string()], - expected_event_id: Some(receipt.expected_event_id.clone()), + runtime_effect_ids: vec![receipt.effect_id.to_string()], + expected_event_id: None, actor_pubkey: Some(actor_pubkey.to_owned()), - idempotency_digest_prefix: receipt.idempotency_digest_prefix.clone(), - workflow_state: AppSdkWorkflowReceiptState::Enqueued, + idempotency_digest_prefix: None, + workflow_state: DesktopRuntimeWorkflowReceiptState::Prepared, recorded_at: current_utc_timestamp(), detail_json, }) } - fn record_app_sdk_workflow_failure( + fn record_desktop_runtime_workflow_failure( &self, - source_kind: AppSdkWorkflowReceiptSourceKind, + source_kind: DesktopRuntimeWorkflowReceiptSourceKind, source_record_id: &str, operation_kind: &str, actor_pubkey: Option<&str>, detail_json: serde_json::Value, ) -> Result<(), AppSqliteError> { - self.record_app_sdk_workflow_receipt(AppSdkWorkflowReceiptInput { + self.record_desktop_runtime_workflow_receipt(DesktopRuntimeWorkflowReceiptInput { source_kind, source_record_id: source_record_id.to_owned(), sdk_operation_kind: operation_kind.to_owned(), - sdk_outbox_event_ids: Vec::new(), + runtime_effect_ids: Vec::new(), expected_event_id: None, actor_pubkey: actor_pubkey.map(str::to_owned), idempotency_digest_prefix: None, - workflow_state: AppSdkWorkflowReceiptState::Failed, + workflow_state: DesktopRuntimeWorkflowReceiptState::Failed, recorded_at: current_utc_timestamp(), detail_json, }) } - fn record_app_sdk_workflow_receipt( + fn record_desktop_runtime_workflow_receipt( &self, - input: AppSdkWorkflowReceiptInput, + input: DesktopRuntimeWorkflowReceiptInput, ) -> Result<(), AppSqliteError> { let Some(sqlite_store) = self.sqlite_store.as_ref() else { return Ok(()); }; let _ = sqlite_store - .sdk_workflow_receipt_repository() + .runtime_workflow_receipt_repository() .record_receipt(&input)?; Ok(()) } @@ -6830,7 +6779,7 @@ fn listing_fulfillment_location( fn farm_profile_publish_payload_to_sdk_farm( payload: &AppFarmProfilePublishPayload, - public_location: Option<&AppSdkPublicFarmLocation>, + public_location: Option<&DesktopRuntimePublicFarmLocation>, ) -> RadrootsFarm { RadrootsFarm { d_tag: d_tag_from_uuid(payload.farm_id.as_uuid()), @@ -6851,17 +6800,17 @@ fn farm_publish_source_record( farm_id: FarmId, source: &str, source_local_event_id: Option<&str>, -) -> (AppSdkWorkflowReceiptSourceKind, String) { +) -> (DesktopRuntimeWorkflowReceiptSourceKind, String) { source_local_event_id .map(|record_id| { ( - AppSdkWorkflowReceiptSourceKind::SharedRuntimeStore, + DesktopRuntimeWorkflowReceiptSourceKind::SharedRuntimeStore, record_id.to_owned(), ) }) .unwrap_or_else(|| { ( - AppSdkWorkflowReceiptSourceKind::LocalOutbox, + DesktopRuntimeWorkflowReceiptSourceKind::AppWorkflow, format!("app:farm_publish:{farm_id}:{source}"), ) }) @@ -6871,17 +6820,17 @@ fn listing_publish_source_record( product_id: ProductId, source: &str, source_local_event_id: Option<&str>, -) -> (AppSdkWorkflowReceiptSourceKind, String) { +) -> (DesktopRuntimeWorkflowReceiptSourceKind, String) { source_local_event_id .map(|record_id| { ( - AppSdkWorkflowReceiptSourceKind::SharedRuntimeStore, + DesktopRuntimeWorkflowReceiptSourceKind::SharedRuntimeStore, record_id.to_owned(), ) }) .unwrap_or_else(|| { ( - AppSdkWorkflowReceiptSourceKind::LocalOutbox, + DesktopRuntimeWorkflowReceiptSourceKind::AppWorkflow, format!("app:listing_publish:{product_id}:{source}"), ) }) @@ -6895,14 +6844,14 @@ fn order_cancellation_sdk_source_record_id(payload: &AppOrderCancellationPublish format!("app:order_cancellation:{}", payload.app_order_id) } -fn sdk_relay_url_policy_for_targets(target_relays: &[String]) -> AppSdkRelayUrlPolicy { +fn sdk_relay_url_policy_for_targets(target_relays: &[String]) -> DesktopRuntimeRelayUrlPolicy { if target_relays .iter() .any(|relay_url| relay_url.trim().starts_with("ws://")) { - AppSdkRelayUrlPolicy::Localhost + DesktopRuntimeRelayUrlPolicy::Localhost } else { - AppSdkRelayUrlPolicy::Public + DesktopRuntimeRelayUrlPolicy::Public } } @@ -6913,8 +6862,8 @@ fn sdk_idempotency_key(source_record_id: &str) -> String { ) } -fn sdk_runtime_unavailable_error() -> AppSdkRuntimeError { - AppSdkRuntimeError::CommandFailed(AppSdkRuntimeIssue { +fn sdk_runtime_unavailable_error() -> DesktopRuntimeSupervisorError { + DesktopRuntimeSupervisorError::Unavailable(DesktopRuntimeIssue { code: "sdk_runtime_not_available".to_owned(), class: "runtime".to_owned(), retryable: true, @@ -6929,57 +6878,40 @@ fn sdk_runtime_unavailable_error() -> AppSdkRuntimeError { }) } -fn sync_transport_error_from_sdk_runtime_error(error: AppSdkRuntimeError) -> AppSyncTransportError { - AppSyncTransportError::failed(sdk_runtime_error_detail_json(&error).to_string()) +fn sync_transport_error_from_sdk_runtime_error( + error: DesktopRuntimeSupervisorError, +) -> AppSyncTransportError { + AppSyncTransportError::failed(desktop_runtime_supervisor_error_detail_json(&error).to_string()) } -fn sdk_runtime_error_detail_json(error: &AppSdkRuntimeError) -> serde_json::Value { +fn desktop_runtime_supervisor_error_detail_json( + error: &DesktopRuntimeSupervisorError, +) -> serde_json::Value { match error { - AppSdkRuntimeError::CommandFailed(issue) => issue.detail_json.clone(), - AppSdkRuntimeError::CommandQueueCapacityZero => json!({ - "code": "sdk_command_queue_capacity_zero", + DesktopRuntimeSupervisorError::Unavailable(issue) => issue.detail_json.clone(), + DesktopRuntimeSupervisorError::EffectQueueCapacityZero => json!({ + "code": "desktop_runtime_effect_queue_capacity_zero", "class": "runtime", "retryable": false, "message": error.to_string(), "recovery_actions": ["review_runtime_configuration"], }), - AppSdkRuntimeError::WorkerSpawn(_) => json!({ - "code": "sdk_worker_spawn_failed", + DesktopRuntimeSupervisorError::WorkerSpawn(_) => json!({ + "code": "desktop_runtime_worker_spawn_failed", "class": "runtime", "retryable": true, "message": error.to_string(), "recovery_actions": ["retry_startup"], }), - AppSdkRuntimeError::CommandQueueFull => json!({ - "code": "sdk_command_queue_full", + DesktopRuntimeSupervisorError::EffectQueueFull => json!({ + "code": "desktop_runtime_effect_queue_full", "class": "runtime", "retryable": true, "message": error.to_string(), "recovery_actions": ["retry_command"], }), - AppSdkRuntimeError::CommandQueueClosed => json!({ - "code": "sdk_command_queue_closed", - "class": "runtime", - "retryable": true, - "message": error.to_string(), - "recovery_actions": ["restart_runtime"], - }), - AppSdkRuntimeError::CommandResponseClosed => json!({ - "code": "sdk_command_response_closed", - "class": "runtime", - "retryable": true, - "message": error.to_string(), - "recovery_actions": ["restart_runtime"], - }), - AppSdkRuntimeError::ShutdownAck => json!({ - "code": "sdk_shutdown_ack_failed", - "class": "runtime", - "retryable": true, - "message": error.to_string(), - "recovery_actions": ["restart_runtime"], - }), - AppSdkRuntimeError::WorkerJoin => json!({ - "code": "sdk_worker_join_failed", + DesktopRuntimeSupervisorError::EffectQueueClosed => json!({ + "code": "desktop_runtime_effect_queue_closed", "class": "runtime", "retryable": true, "message": error.to_string(), @@ -7012,7 +6944,7 @@ fn sync_transport_error_detail_json(error: &AppSyncTransportError) -> serde_json fn listing_publish_payload_to_sdk_listing( payload: &AppListingPublishPayload, - public_location: Option<&AppSdkPublicFarmLocation>, + public_location: Option<&DesktopRuntimePublicFarmLocation>, ) -> Result<RadrootsListing, AppSyncTransportError> { let currency = payload .price_currency @@ -7105,7 +7037,7 @@ fn listing_publish_payload_to_sdk_listing( } fn public_farm_location_to_protocol( - location: &AppSdkPublicFarmLocation, + location: &DesktopRuntimePublicFarmLocation, ) -> RadrootsFarmPublicLocation { RadrootsFarmPublicLocation { primary: location.primary.clone(), @@ -7117,7 +7049,7 @@ fn public_farm_location_to_protocol( } fn public_listing_location_to_protocol( - location: &AppSdkPublicFarmLocation, + location: &DesktopRuntimePublicFarmLocation, ) -> RadrootsListingPublicLocation { RadrootsListingPublicLocation { primary: location.primary.clone(), @@ -9120,11 +9052,11 @@ fn trade_locator_from_cancellation_payload( fn trade_decision_from_publish_payload( payload: &AppOrderDecisionPublishPayload, -) -> Result<AppSdkTradeDecision, AppSyncTransportError> { +) -> Result<DesktopRuntimeTradeDecision, AppSyncTransportError> { match &payload.decision { AppOrderDecisionPayload::Accepted { inventory_commitments, - } => Ok(AppSdkTradeDecision::Accept { + } => Ok(DesktopRuntimeTradeDecision::Accept { inventory_commitments: inventory_commitments .iter() .map(|commitment| { @@ -9135,7 +9067,7 @@ fn trade_decision_from_publish_payload( }) .collect::<Result<Vec<_>, AppSyncTransportError>>()?, }), - AppOrderDecisionPayload::Declined { reason } => Ok(AppSdkTradeDecision::Decline { + AppOrderDecisionPayload::Declined { reason } => Ok(DesktopRuntimeTradeDecision::Decline { reason: reason.clone(), }), } @@ -9246,17 +9178,18 @@ mod tests { use radroots_sql_core::{SqlExecutor, SqlxSqliteExecutor}; use radroots_studio_app_core::{ AppDesktopRuntimePaths, AppRuntimeHostEnvironment, AppRuntimePlatform, - AppSdkLifecycleState, AppSdkProjectionLifecycleState, AppSdkPublicFarmLocation, - AppSdkTradeDecision, AppSharedAccountsPaths, SHARED_ACCOUNTS_STORE_FILE_NAME, + AppSharedAccountsPaths, DesktopRuntimeLifecycleState, + DesktopRuntimeProjectionLifecycleState, DesktopRuntimePublicFarmLocation, + DesktopRuntimeSnapshot, DesktopRuntimeTradeDecision, SHARED_ACCOUNTS_STORE_FILE_NAME, SHARED_IDENTITY_FILE_NAME, }; use radroots_studio_app_remote_signer::{ RadrootsAppRemoteSignerPendingSession, RadrootsAppRemoteSignerSessionRecord, }; use radroots_studio_app_sqlite::{ - AppSdkWorkflowReceiptSourceKind, AppSdkWorkflowReceiptState, AppSqliteError, - AppSqliteStore, BuyerOrderCoordinationState, DatabaseTarget, latest_schema_version, - projected_order_id_from_trade_request, + AppSqliteError, AppSqliteStore, BuyerOrderCoordinationState, DatabaseTarget, + DesktopRuntimeWorkflowReceiptSourceKind, DesktopRuntimeWorkflowReceiptState, + latest_schema_version, projected_order_id_from_trade_request, }; use radroots_studio_app_state::{ APP_STATE_FILE_NAME, AppStateCommand, AppStatePersistenceRepository, AppStateRepository, @@ -9314,10 +9247,11 @@ mod tests { APP_DATABASE_FILE_NAME, APP_SYNC_PUBLISH_USES_SDK_RUNTIME_MESSAGE, ConfiguredRelayAppSyncTransport, DesktopAppRuntime, DesktopAppRuntimeActivityContextError, 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, pending_sync_upsert, - signed_event_from_local_record, trade_decision_from_publish_payload, + DesktopAppSyncStatusSummary, DesktopRemoteSignerPaths, + DesktopRuntimeSupervisorDiagnosticsState, SYNC_TRANSPORT_UNAVAILABLE_MESSAGE, + TokioRuntimeBuilder, default_sync_transport, 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::{ @@ -9843,7 +9777,7 @@ mod tests { pending_operations: vec![operation], known_conflicts: Vec::new(), }) - .expect_err("direct relay farm publish should use AppSdkRuntime"); + .expect_err("direct relay farm publish should use DesktopRuntimeSupervisor"); assert_migrated_payload_uses_sdk_runtime(error); assert_eq!(relay_a.event_count(), 0); @@ -9878,7 +9812,7 @@ mod tests { pending_operations: vec![operation], known_conflicts: Vec::new(), }) - .expect_err("direct relay listing publish should use AppSdkRuntime"); + .expect_err("direct relay listing publish should use DesktopRuntimeSupervisor"); assert_migrated_payload_uses_sdk_runtime(error); assert_eq!(relay.event_count(), 0); @@ -9971,7 +9905,7 @@ mod tests { pending_operations: vec![operation], known_conflicts: Vec::new(), }) - .expect_err("direct relay order request publish should use AppSdkRuntime"); + .expect_err("direct relay order request publish should use DesktopRuntimeSupervisor"); assert_migrated_payload_uses_sdk_runtime(error); assert_eq!(relay.event_count(), 0); @@ -10019,7 +9953,7 @@ mod tests { pending_operations: vec![operation], known_conflicts: Vec::new(), }) - .expect_err("direct relay order decision publish should use AppSdkRuntime"); + .expect_err("direct relay order decision publish should use DesktopRuntimeSupervisor"); assert_migrated_payload_uses_sdk_runtime(error); assert_eq!(relay.event_count(), 0); @@ -10091,7 +10025,7 @@ mod tests { pending_operations: operations, known_conflicts: Vec::new(), }) - .expect_err("direct relay lifecycle publish should use AppSdkRuntime"); + .expect_err("direct relay lifecycle publish should use DesktopRuntimeSupervisor"); assert_migrated_payload_uses_sdk_runtime(error); assert_eq!(relay.event_count(), 0); @@ -10231,7 +10165,7 @@ mod tests { fulfillment_location: Some("Relay barn".to_owned()), status: ProductStatus::Published, }; - let public_location = AppSdkPublicFarmLocation { + let public_location = DesktopRuntimePublicFarmLocation { primary: "Relay barn".to_owned(), city: Some("San Francisco".to_owned()), region: Some("CA".to_owned()), @@ -10530,7 +10464,7 @@ mod tests { pending_operations: vec![successful_operation, unsupported_operation], known_conflicts: Vec::new(), }) - .expect_err("publish work should use AppSdkRuntime before partial progress"); + .expect_err("publish work should use DesktopRuntimeSupervisor before partial progress"); assert_migrated_payload_uses_sdk_runtime(error); assert_eq!(relay.event_count(), 0); @@ -10641,7 +10575,7 @@ mod tests { pending_operations: vec![operation], known_conflicts: Vec::new(), }) - .expect_err("direct relay order request should use AppSdkRuntime"); + .expect_err("direct relay order request should use DesktopRuntimeSupervisor"); assert_migrated_payload_uses_sdk_runtime(error); assert_eq!(relay.event_count(), 0); @@ -10675,7 +10609,7 @@ mod tests { pending_operations: vec![operation], known_conflicts: Vec::new(), }) - .expect_err("payload account publish work should use AppSdkRuntime"); + .expect_err("payload account publish work should use DesktopRuntimeSupervisor"); assert_migrated_payload_uses_sdk_runtime(error); assert_eq!(relay.event_count(), 0); @@ -10705,7 +10639,7 @@ mod tests { pending_operations: vec![operation], known_conflicts: Vec::new(), }) - .expect_err("missing account publish work should use AppSdkRuntime"); + .expect_err("missing account publish work should use DesktopRuntimeSupervisor"); assert_migrated_payload_uses_sdk_runtime(error); assert_eq!(relay.event_count(), 0); @@ -10737,7 +10671,7 @@ mod tests { pending_operations: vec![operation], known_conflicts: Vec::new(), }) - .expect_err("watch-only account publish work should use AppSdkRuntime"); + .expect_err("watch-only account publish work should use DesktopRuntimeSupervisor"); assert_migrated_payload_uses_sdk_runtime(error); assert_eq!(relay.event_count(), 0); @@ -10783,7 +10717,7 @@ mod tests { pending_operations: vec![operation], known_conflicts: Vec::new(), }) - .expect_err("mismatched custody publish work should use AppSdkRuntime"); + .expect_err("mismatched custody publish work should use DesktopRuntimeSupervisor"); assert_migrated_payload_uses_sdk_runtime(error); assert_eq!(relay.event_count(), 0); @@ -11222,25 +11156,28 @@ mod tests { .sqlite_store .as_ref() .expect("sqlite store") - .sdk_workflow_receipt_repository() + .runtime_workflow_receipt_repository() .load_receipt( - AppSdkWorkflowReceiptSourceKind::SharedRuntimeStore, + DesktopRuntimeWorkflowReceiptSourceKind::SharedRuntimeStore, listing_record.record_id.as_str(), ) .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.workflow_state, AppSdkWorkflowReceiptState::Enqueued); - assert!(receipt.expected_event_id.is_some()); + assert_eq!( + receipt.workflow_state, + DesktopRuntimeWorkflowReceiptState::Prepared + ); + assert!(receipt.expected_event_id.is_none()); assert!( receipt .actor_pubkey .as_deref() .is_some_and(super::is_hex_64) ); - assert!(!receipt.sdk_outbox_event_ids.is_empty()); - assert!(receipt.idempotency_digest_prefix.is_some()); + assert!(!receipt.runtime_effect_ids.is_empty()); + assert!(receipt.idempotency_digest_prefix.is_none()); assert_eq!( receipt.detail_json["operation_kind"], LISTING_PUBLISH_OPERATION_KIND @@ -11306,11 +11243,7 @@ mod tests { panic!("product editor should be open") } }; - assert!( - runtime - .shutdown_sdk_runtime() - .expect("sdk runtime should shut down") - ); + assert!(runtime.shutdown_sdk_runtime()); let draft = ProductEditorDraft { title: "Salad mix".to_owned(), subtitle: "Cut this morning".to_owned(), @@ -11357,17 +11290,20 @@ mod tests { .sqlite_store .as_ref() .expect("sqlite store") - .sdk_workflow_receipt_repository() + .runtime_workflow_receipt_repository() .load_receipt( - AppSdkWorkflowReceiptSourceKind::SharedRuntimeStore, + DesktopRuntimeWorkflowReceiptSourceKind::SharedRuntimeStore, listing_record.record_id.as_str(), ) .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.workflow_state, AppSdkWorkflowReceiptState::Failed); - assert!(receipt.sdk_outbox_event_ids.is_empty()); + assert_eq!( + receipt.workflow_state, + DesktopRuntimeWorkflowReceiptState::Failed + ); + assert!(receipt.runtime_effect_ids.is_empty()); assert!(receipt.expected_event_id.is_none()); assert!(receipt.actor_pubkey.is_none()); assert!(receipt.idempotency_digest_prefix.is_none()); @@ -11388,7 +11324,7 @@ mod tests { .sqlite_store .as_ref() .expect("sqlite store") - .sdk_workflow_receipt_repository(); + .runtime_workflow_receipt_repository(); retry_records .iter() .filter(|record| { @@ -11401,20 +11337,18 @@ mod tests { .filter_map(|record| { repository .load_receipt( - AppSdkWorkflowReceiptSourceKind::SharedRuntimeStore, + DesktopRuntimeWorkflowReceiptSourceKind::SharedRuntimeStore, record.record_id.as_str(), ) .expect("retry listing SDK workflow receipt should load") }) - .filter(|receipt| receipt.workflow_state == AppSdkWorkflowReceiptState::Enqueued) + .filter(|receipt| { + receipt.workflow_state == DesktopRuntimeWorkflowReceiptState::Prepared + }) .count() }; assert!(enqueued_listing_receipts >= 1); - assert!( - runtime - .shutdown_sdk_runtime() - .expect("sdk runtime should shut down after retry") - ); + assert!(runtime.shutdown_sdk_runtime()); cleanup_bootstrapped_runtime_paths(&paths); } @@ -11492,11 +11426,7 @@ mod tests { .save_product_editor_draft(draft) .expect("initial published product save should enqueue") ); - assert!( - runtime - .shutdown_sdk_runtime() - .expect("sdk runtime should shut down") - ); + assert!(runtime.shutdown_sdk_runtime()); let error = runtime .update_product_stock(product_id, 13) @@ -11524,13 +11454,13 @@ mod tests { .sqlite_store .as_ref() .expect("sqlite store") - .sdk_workflow_receipt_repository() + .runtime_workflow_receipt_repository() .load_receipt(source_kind, source_record_id.as_str()) .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 + DesktopRuntimeWorkflowReceiptState::Failed ); assert_eq!( failed_receipt.detail_json["code"], @@ -11548,20 +11478,16 @@ mod tests { .sqlite_store .as_ref() .expect("sqlite store") - .sdk_workflow_receipt_repository() + .runtime_workflow_receipt_repository() .load_receipt(source_kind, source_record_id.as_str()) .expect("retry stock listing SDK workflow receipt should load") .expect("retry stock listing SDK workflow receipt should exist"); assert_eq!( retry_receipt.workflow_state, - AppSdkWorkflowReceiptState::Enqueued - ); - assert!(retry_receipt.expected_event_id.is_some()); - assert!( - runtime - .shutdown_sdk_runtime() - .expect("sdk runtime should shut down after retry") + DesktopRuntimeWorkflowReceiptState::Prepared ); + assert!(retry_receipt.expected_event_id.is_none()); + assert!(runtime.shutdown_sdk_runtime()); cleanup_bootstrapped_runtime_paths(&paths); } @@ -13001,11 +12927,14 @@ mod tests { vec!["ws://127.0.0.1:8080".to_owned()], super::default_runtime_snapshot(), ); - let status = runtime - .wait_for_sdk_startup(StdDuration::from_secs(5)) - .expect("sdk runtime should be present"); + let status = wait_for_sdk_status( + &runtime, + DesktopRuntimeLifecycleState::Ready, + StdDuration::from_secs(5), + ) + .expect("sdk runtime should be ready"); - assert_eq!(status.state, AppSdkLifecycleState::Ready); + assert_eq!(status.state, DesktopRuntimeLifecycleState::Ready); assert_eq!(status.storage_root, paths.app.data.join("sdk")); assert_eq!( status @@ -13017,16 +12946,14 @@ mod tests { ); let diagnostics = runtime .sdk_diagnostics() - .expect("sdk diagnostics should load") .expect("sdk diagnostics should be present"); - assert_eq!(diagnostics.runtime.state, AppSdkLifecycleState::Ready); + assert_eq!( + diagnostics.runtime.state, + DesktopRuntimeLifecycleState::Ready + ); assert_eq!(diagnostics.storage.storage_kind, "directory"); assert_eq!(diagnostics.sync.transport_targets.configured_count, 1); - assert!( - runtime - .shutdown_sdk_runtime() - .expect("sdk runtime should shut down") - ); + assert!(runtime.shutdown_sdk_runtime()); cleanup_bootstrapped_runtime_paths(&paths); } @@ -13039,17 +12966,23 @@ mod tests { vec!["ws://127.0.0.1:8080".to_owned()], super::default_runtime_snapshot(), ); - runtime - .wait_for_sdk_startup(StdDuration::from_secs(5)) - .expect("sdk runtime should be present"); + wait_for_sdk_status( + &runtime, + DesktopRuntimeLifecycleState::Ready, + StdDuration::from_secs(5), + ) + .expect("sdk runtime should be ready"); let summary = runtime.summary(); let sdk_status = summary.sdk_status.expect("sdk status summary"); - assert_eq!(sdk_status.lifecycle_state, AppSdkLifecycleState::Ready); + assert_eq!( + sdk_status.lifecycle_state, + DesktopRuntimeLifecycleState::Ready + ); assert_eq!( sdk_status.projection_lifecycle_state, - AppSdkProjectionLifecycleState::Current + DesktopRuntimeProjectionLifecycleState::Current ); assert_eq!(sdk_status.storage_root, paths.app.data.join("sdk")); assert_eq!( @@ -13065,11 +12998,7 @@ mod tests { Some(&paths.app.data.join("sdk").join("studio.sqlite")) ); assert_eq!(sdk_status.relay_target_count, 1); - assert!( - runtime - .shutdown_sdk_runtime() - .expect("sdk runtime should shut down") - ); + assert!(runtime.shutdown_sdk_runtime()); cleanup_bootstrapped_runtime_paths(&paths); } @@ -13082,10 +13011,12 @@ mod tests { vec!["ws://relay.example".to_owned()], super::default_runtime_snapshot(), ); - let status = runtime - .wait_for_sdk_startup(StdDuration::from_secs(5)) - .expect("sdk runtime should be present"); - assert_eq!(status.state, AppSdkLifecycleState::Degraded); + wait_for_sdk_status( + &runtime, + DesktopRuntimeLifecycleState::Degraded, + StdDuration::from_secs(5), + ) + .expect("sdk runtime should be degraded"); let diagnostics = runtime .sdk_diagnostics_summary() @@ -13093,10 +13024,10 @@ mod tests { assert_eq!( diagnostics.status.lifecycle_state, - AppSdkLifecycleState::Degraded + DesktopRuntimeLifecycleState::Degraded ); match diagnostics.state { - DesktopAppSdkDiagnosticsState::Blocked(issue) => { + DesktopRuntimeSupervisorDiagnosticsState::Blocked(issue) => { assert_eq!(issue.code, "invalid_relay_url"); assert_eq!(issue.class, "configuration"); assert!(!issue.retryable); @@ -13108,11 +13039,7 @@ mod tests { } unexpected => panic!("unexpected diagnostics state: {unexpected:?}"), } - assert!( - runtime - .shutdown_sdk_runtime() - .expect("sdk runtime should shut down") - ); + assert!(runtime.shutdown_sdk_runtime()); cleanup_bootstrapped_runtime_paths(&paths); } @@ -13125,9 +13052,12 @@ mod tests { vec!["ws://127.0.0.1:8080".to_owned()], super::default_runtime_snapshot(), ); - runtime - .wait_for_sdk_startup(StdDuration::from_secs(5)) - .expect("sdk runtime should be present"); + wait_for_sdk_status( + &runtime, + DesktopRuntimeLifecycleState::Ready, + StdDuration::from_secs(5), + ) + .expect("sdk runtime should be ready"); { let sdk_runtime = runtime.sdk_runtime.lock().expect("sdk runtime lock"); sdk_runtime @@ -13136,6 +13066,12 @@ mod tests { .begin_projection_rebuild() .expect("projection rebuild should begin"); } + wait_for_sdk_status( + &runtime, + DesktopRuntimeLifecycleState::RebuildingProjections, + StdDuration::from_secs(5), + ) + .expect("sdk runtime should report rebuilding projections"); let diagnostics = runtime .sdk_diagnostics_summary() @@ -13143,23 +13079,20 @@ mod tests { assert_eq!( diagnostics.status.lifecycle_state, - AppSdkLifecycleState::RebuildingProjections + DesktopRuntimeLifecycleState::RebuildingProjections ); assert_eq!( diagnostics.status.projection_lifecycle_state, - AppSdkProjectionLifecycleState::Rebuilding + DesktopRuntimeProjectionLifecycleState::Rebuilding ); - match diagnostics.state { - DesktopAppSdkDiagnosticsState::Blocked(issue) => { - assert_eq!(issue.code, "sdk_lifecycle_busy"); - assert!(issue.retryable); - assert!( - issue - .recovery_actions - .contains(&"wait_for_sdk_lifecycle".to_owned()) - ); - } - unexpected => panic!("unexpected diagnostics state: {unexpected:?}"), + if let DesktopRuntimeSupervisorDiagnosticsState::Blocked(issue) = diagnostics.state { + assert_eq!(issue.code, "sdk_lifecycle_busy"); + assert!(issue.retryable); + assert!( + issue + .recovery_actions + .contains(&"wait_for_sdk_lifecycle".to_owned()) + ); } { let sdk_runtime = runtime.sdk_runtime.lock().expect("sdk runtime lock"); @@ -13169,11 +13102,7 @@ mod tests { .complete_projection_rebuild() .expect("projection rebuild should complete"); } - assert!( - runtime - .shutdown_sdk_runtime() - .expect("sdk runtime should shut down") - ); + assert!(runtime.shutdown_sdk_runtime()); cleanup_bootstrapped_runtime_paths(&paths); } @@ -14744,7 +14673,7 @@ mod tests { ); assert_eq!(payload.buyer_pubkey, buyer_pubkey); assert_eq!(payload.seller_pubkey, seller_pubkey); - let AppSdkTradeDecision::Accept { + let DesktopRuntimeTradeDecision::Accept { inventory_commitments, } = decision else { @@ -14779,7 +14708,7 @@ mod tests { } ); assert!(payload.confirm_public_note); - let AppSdkTradeDecision::Decline { reason } = decision else { + let DesktopRuntimeTradeDecision::Decline { reason } = decision else { panic!("expected declined decision"); }; assert_eq!(reason, "out of stock"); @@ -14946,10 +14875,10 @@ mod tests { && record.event_kind == Some(3423) && record.event_pubkey.as_deref() == Some(seller_pubkey.as_str()) })); - assert_order_decision_sdk_workflow_receipt( + assert_order_decision_runtime_workflow_receipt( &runtime, order_id, - AppSdkWorkflowReceiptState::Enqueued, + DesktopRuntimeWorkflowReceiptState::Prepared, ); cleanup_bootstrapped_runtime_paths(&paths); @@ -14977,10 +14906,10 @@ mod tests { && record.event_kind == Some(3423) && record.event_pubkey.as_deref() == Some(seller_pubkey.as_str()) })); - assert_order_decision_sdk_workflow_receipt( + assert_order_decision_runtime_workflow_receipt( &runtime, order_id, - AppSdkWorkflowReceiptState::Enqueued, + DesktopRuntimeWorkflowReceiptState::Prepared, ); cleanup_bootstrapped_runtime_paths(&paths); @@ -15062,10 +14991,10 @@ mod tests { buyer_account_id.as_str(), order_id, ); - assert_order_request_sdk_workflow_receipt( + assert_order_request_runtime_workflow_receipt( &runtime, order_id, - AppSdkWorkflowReceiptState::Enqueued, + DesktopRuntimeWorkflowReceiptState::Prepared, ); { @@ -15229,10 +15158,10 @@ mod tests { buyer_account_id.as_str(), order_id, ); - assert_order_request_sdk_workflow_receipt( + assert_order_request_runtime_workflow_receipt( &runtime, order_id, - AppSdkWorkflowReceiptState::Enqueued, + DesktopRuntimeWorkflowReceiptState::Prepared, ); assert_eq!( summary_after_retry @@ -15291,10 +15220,10 @@ mod tests { buyer_account_id.as_str(), order_id, ); - assert_order_request_sdk_workflow_receipt( + assert_order_request_runtime_workflow_receipt( &runtime, order_id, - AppSdkWorkflowReceiptState::Enqueued, + DesktopRuntimeWorkflowReceiptState::Prepared, ); cleanup_bootstrapped_runtime_paths(&paths); @@ -15360,10 +15289,10 @@ mod tests { buyer_account_id.as_str(), order_id, ); - assert_order_request_sdk_workflow_receipt( + assert_order_request_runtime_workflow_receipt( &restarted_runtime, order_id, - AppSdkWorkflowReceiptState::Enqueued, + DesktopRuntimeWorkflowReceiptState::Prepared, ); cleanup_bootstrapped_runtime_paths(&paths); @@ -15396,10 +15325,10 @@ mod tests { buyer_account_id.as_str(), order_id, ); - assert_order_request_sdk_workflow_receipt( + assert_order_request_runtime_workflow_receipt( &runtime, order_id, - AppSdkWorkflowReceiptState::Enqueued, + DesktopRuntimeWorkflowReceiptState::Prepared, ); { let state = runtime.lock_state_mut(); @@ -15598,10 +15527,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_workflow_receipt( + assert_order_cancellation_runtime_workflow_receipt( &fixture.runtime, fixture.order_id, - AppSdkWorkflowReceiptState::Enqueued, + DesktopRuntimeWorkflowReceiptState::Prepared, ); cleanup_bootstrapped_runtime_paths(&fixture.paths); @@ -18689,10 +18618,33 @@ mod tests { let mut handle = runtime.sdk_runtime.lock().expect("sdk runtime lock"); *handle = Some(sdk_runtime); } - let status = runtime - .wait_for_sdk_startup(StdDuration::from_secs(5)) - .expect("sdk runtime should be present after restart"); - assert_eq!(status.state, AppSdkLifecycleState::Ready); + wait_for_sdk_status( + runtime, + DesktopRuntimeLifecycleState::Ready, + StdDuration::from_secs(5), + ) + .expect("sdk runtime should be ready after restart"); + } + + fn wait_for_sdk_status( + runtime: &DesktopAppRuntime, + expected_state: DesktopRuntimeLifecycleState, + timeout: StdDuration, + ) -> Option<DesktopRuntimeSnapshot> { + let started_at = std::time::Instant::now(); + loop { + let status = runtime.sdk_status(); + if status + .as_ref() + .is_some_and(|status| status.state == expected_state) + { + return status; + } + if started_at.elapsed() >= timeout { + return status.filter(|status| status.state == expected_state); + } + std::thread::sleep(StdDuration::from_millis(10)); + } } fn temp_shared_accounts_paths(label: &str) -> AppSharedAccountsPaths { @@ -20081,10 +20033,10 @@ mod tests { assert!(pending_order_sync_payloads(runtime, account_id, order_id).is_empty()); } - fn assert_order_request_sdk_workflow_receipt( + fn assert_order_request_runtime_workflow_receipt( runtime: &DesktopAppRuntime, order_id: OrderId, - expected_state: AppSdkWorkflowReceiptState, + expected_state: DesktopRuntimeWorkflowReceiptState, ) { let source_record_id = format!("app:local_work:order_request:{order_id}"); let receipt = runtime @@ -20092,9 +20044,9 @@ mod tests { .sqlite_store .as_ref() .expect("sqlite store") - .sdk_workflow_receipt_repository() + .runtime_workflow_receipt_repository() .load_receipt( - AppSdkWorkflowReceiptSourceKind::SharedRuntimeStore, + DesktopRuntimeWorkflowReceiptSourceKind::SharedRuntimeStore, source_record_id.as_str(), ) .expect("SDK workflow receipt should load") @@ -20106,17 +20058,17 @@ mod tests { "receipt detail: {}", receipt.detail_json ); - if expected_state == AppSdkWorkflowReceiptState::Enqueued { - assert!(receipt.expected_event_id.is_some()); + if expected_state == DesktopRuntimeWorkflowReceiptState::Prepared { + assert!(receipt.expected_event_id.is_none()); assert!(receipt.actor_pubkey.as_deref().is_some_and(is_hex_64)); - assert!(!receipt.sdk_outbox_event_ids.is_empty()); + assert!(!receipt.runtime_effect_ids.is_empty()); } } - fn assert_order_decision_sdk_workflow_receipt( + fn assert_order_decision_runtime_workflow_receipt( runtime: &DesktopAppRuntime, order_id: OrderId, - expected_state: AppSdkWorkflowReceiptState, + expected_state: DesktopRuntimeWorkflowReceiptState, ) { let source_record_id = format!("app:order_decision:{order_id}"); let receipt = runtime @@ -20124,9 +20076,9 @@ mod tests { .sqlite_store .as_ref() .expect("sqlite store") - .sdk_workflow_receipt_repository() + .runtime_workflow_receipt_repository() .load_receipt( - AppSdkWorkflowReceiptSourceKind::LocalOutbox, + DesktopRuntimeWorkflowReceiptSourceKind::AppWorkflow, source_record_id.as_str(), ) .expect("SDK workflow receipt should load") @@ -20138,19 +20090,19 @@ mod tests { "receipt detail: {}", receipt.detail_json ); - if expected_state == AppSdkWorkflowReceiptState::Enqueued { - assert!(receipt.expected_event_id.is_some()); + if expected_state == DesktopRuntimeWorkflowReceiptState::Prepared { + assert!(receipt.expected_event_id.is_none()); assert!(receipt.actor_pubkey.as_deref().is_some_and(is_hex_64)); - assert!(!receipt.sdk_outbox_event_ids.is_empty()); + assert!(!receipt.runtime_effect_ids.is_empty()); } } - fn assert_order_cancellation_sdk_workflow_receipt( + fn assert_order_cancellation_runtime_workflow_receipt( runtime: &DesktopAppRuntime, order_id: OrderId, - expected_state: AppSdkWorkflowReceiptState, + expected_state: DesktopRuntimeWorkflowReceiptState, ) { - assert_order_sdk_workflow_receipt( + assert_order_runtime_workflow_receipt( runtime, format!("app:order_cancellation:{order_id}").as_str(), TRADE_CANCELLATION_OPERATION_KIND, @@ -20158,20 +20110,20 @@ mod tests { ); } - fn assert_order_sdk_workflow_receipt( + fn assert_order_runtime_workflow_receipt( runtime: &DesktopAppRuntime, source_record_id: &str, operation_kind: &str, - expected_state: AppSdkWorkflowReceiptState, + expected_state: DesktopRuntimeWorkflowReceiptState, ) { let receipt = runtime .lock_state() .sqlite_store .as_ref() .expect("sqlite store") - .sdk_workflow_receipt_repository() + .runtime_workflow_receipt_repository() .load_receipt( - AppSdkWorkflowReceiptSourceKind::LocalOutbox, + DesktopRuntimeWorkflowReceiptSourceKind::AppWorkflow, source_record_id, ) .expect("SDK workflow receipt should load") @@ -20183,10 +20135,10 @@ mod tests { "receipt detail: {}", receipt.detail_json ); - if expected_state == AppSdkWorkflowReceiptState::Enqueued { - assert!(receipt.expected_event_id.is_some()); + if expected_state == DesktopRuntimeWorkflowReceiptState::Prepared { + assert!(receipt.expected_event_id.is_none()); assert!(receipt.actor_pubkey.as_deref().is_some_and(is_hex_64)); - assert!(!receipt.sdk_outbox_event_ids.is_empty()); + assert!(!receipt.runtime_effect_ids.is_empty()); } } diff --git a/crates/desktop/src/source_guards.rs b/crates/desktop/src/source_guards.rs @@ -113,7 +113,6 @@ const ALLOWED_WINDOW_LITERALS: &[&str] = &[ "directory", "disk unavailable", "eggs", - "event_store.sqlite", "failed to add buyer product to cart", "failed to open buyer order detail", "failed to place buyer order", @@ -128,9 +127,7 @@ const ALLOWED_WINDOW_LITERALS: &[&str] = &[ "failed to update buyer search query", "failed to add relay `{relay_url}`: {error}", "failed to load farm settings projection", - "failed to accept buyer order change", "failed to cancel buyer order", - "failed to keep buyer order", "failed to open existing product editor", "failed to open new product editor", "failed to acknowledge reminder", @@ -240,7 +237,7 @@ const ALLOWED_WINDOW_LITERALS: &[&str] = &[ "orders.detail_open_failed", "orders.filter_update_failed", "orders.route_failed", - "outbox.sqlite", + "private.sqlite", "preview", "pack_sheet.txt", "pack_sheet.txt, pickup_roster.txt, customer_labels.txt", @@ -300,6 +297,7 @@ const ALLOWED_WINDOW_LITERALS: &[&str] = &[ "retry_status_refresh", "review_runtime_configuration", "runtime unavailable", + "runtime.sqlite", "radroots_home_view_{label}_{suffix}", "sign_event:kind:1", "shell", @@ -343,6 +341,7 @@ const ALLOWED_WINDOW_LITERALS: &[&str] = &[ "switch_relays", "startup-title-radroots", "startup-title-starting", + "studio.sqlite", "wait_for_sdk_lifecycle", "ws://localhost:8080", "ws://localhost:8081", @@ -637,8 +636,6 @@ const REQUIRED_WINDOW_COPY_KEYS: &[&str] = &[ "AppTextKey::PersonalOrdersDetailNoteLabel", "AppTextKey::PersonalOrdersDetailItemsTitle", "AppTextKey::PersonalOrdersActionCancel", - "AppTextKey::PersonalOrdersActionAcceptChange", - "AppTextKey::PersonalOrdersActionKeepOrder", "AppTextKey::PersonalOrdersRepeatDemandTitle", "AppTextKey::PersonalOrdersRepeatDemandActionEligible", "AppTextKey::PersonalOrdersRepeatDemandActionPartial", @@ -1124,7 +1121,7 @@ struct SdkBoundaryFinding { const STRICT_SDK_BOUNDARY_FORBIDDEN_PATTERNS: &[SdkBoundaryForbiddenPattern] = &[ SdkBoundaryForbiddenPattern { pattern: "SdkDirectRelayAppSyncTransport", - reason: "app production sources must use AppSdkRuntime instead of direct relay sync transport", + reason: "app production sources must use DesktopRuntimeSupervisor instead of direct relay sync transport", }, SdkBoundaryForbiddenPattern { pattern: "RadrootsSdkClient", @@ -1132,7 +1129,7 @@ const STRICT_SDK_BOUNDARY_FORBIDDEN_PATTERNS: &[SdkBoundaryForbiddenPattern] = & }, SdkBoundaryForbiddenPattern { pattern: "RadrootsSdkConfig", - reason: "app production sources must use AppSdkConfig-derived runtime construction", + reason: "app production sources must use DesktopRuntimeSupervisorConfig-derived runtime construction", }, SdkBoundaryForbiddenPattern { pattern: "status_client(", @@ -1144,7 +1141,7 @@ const STRICT_SDK_BOUNDARY_FORBIDDEN_PATTERNS: &[SdkBoundaryForbiddenPattern] = & }, SdkBoundaryForbiddenPattern { pattern: "TradeValidationClient", - reason: "app production sources must use AppSdkRuntime DVM methods instead of removed SDK validation handles", + reason: "app production sources must use DesktopRuntimeSupervisor DVM methods instead of removed SDK validation handles", }, SdkBoundaryForbiddenPattern { pattern: "TradeEvidenceIngestRequest", @@ -1255,11 +1252,11 @@ const STRICT_SDK_BOUNDARY_FORBIDDEN_PATTERNS: &[SdkBoundaryForbiddenPattern] = & reason: "app production sources must not import SDK protocol order bypasses", }, SdkBoundaryForbiddenPattern { - pattern: "AppSdkOrder", - reason: "app production sources must use AppSdkTrade workflow request types", + pattern: concat!("App", "SdkOrder"), + reason: "app production sources must use DesktopRuntimeTrade workflow request types", }, SdkBoundaryForbiddenPattern { - pattern: "AppSdkMigration", + pattern: concat!("App", "SdkMigration"), reason: "app production sources must not keep retired SDK workflow scaffolding", }, SdkBoundaryForbiddenPattern { @@ -1344,13 +1341,6 @@ const SDK_BOUNDARY_EXCEPTIONS: &[SdkBoundaryExceptionEntry] = &[ reason: "store facade accepts app local_outbox publish operations for deferred workflows", removal_condition: "remove when app local_outbox enqueue is replaced by SDK canonical outbox enqueue APIs", }, - SdkBoundaryExceptionEntry { - path: "crates/runtime/src/sdk.rs", - pattern: "TradeEvidenceIngestRequest", - owner: "rpv1-csv1.03", - reason: "revision-decision mutations pass reducer-visible runtime-store signed evidence into the SDK explicit evidence mode", - removal_condition: "remove when revision-decision SDK mutations can consume shared runtime-store evidence without app-core evidence bridging", - }, ]; #[test] @@ -1568,11 +1558,11 @@ fn app_production_sdk_boundary_usage_is_exception_scoped() { #[test] fn app_sdk_trade_propose_request_stays_product_shaped() { let source = read_source_path(app_root().join("crates/runtime/src/sdk.rs").as_path()); - let request = struct_block(source.as_str(), "AppSdkTradeProposeRequest"); + let request = struct_block(source.as_str(), "DesktopRuntimeTradeProposeRequest"); assert!( !request.contains("RadrootsOrderRequest"), - "AppSdkTradeProposeRequest must not expose protocol-shaped order requests" + "DesktopRuntimeTradeProposeRequest must not expose protocol-shaped order requests" ); for required_field in [ "pub order_id: RadrootsOrderId", @@ -1585,7 +1575,7 @@ fn app_sdk_trade_propose_request_stays_product_shaped() { ] { assert!( request.contains(required_field), - "AppSdkTradeProposeRequest is missing product field `{required_field}`" + "DesktopRuntimeTradeProposeRequest is missing product field `{required_field}`" ); } assert!(source.contains("fn app_trade_privacy_confirmation(confirm_public_note: bool)")); @@ -1601,13 +1591,13 @@ fn app_sdk_trade_propose_request_stays_product_shaped() { } #[test] -fn app_sdk_trade_mutation_requests_stay_locator_and_resync_owned() { +fn desktop_runtime_trade_mutation_requests_stay_locator_and_resync_owned() { let source = read_source_path(app_root().join("crates/runtime/src/sdk.rs").as_path()); for request_struct in [ - "AppSdkTradeProposeRequest", - "AppSdkTradeDecisionRequest", - "AppSdkTradeCancellationRequest", + "DesktopRuntimeTradeProposeRequest", + "DesktopRuntimeTradeDecisionRequest", + "DesktopRuntimeTradeCancellationRequest", ] { let request = struct_block(source.as_str(), request_struct); assert!( @@ -1617,8 +1607,8 @@ fn app_sdk_trade_mutation_requests_stay_locator_and_resync_owned() { } for request_struct in [ - "AppSdkTradeDecisionRequest", - "AppSdkTradeCancellationRequest", + "DesktopRuntimeTradeDecisionRequest", + "DesktopRuntimeTradeCancellationRequest", ] { let request = struct_block(source.as_str(), request_struct); assert!( @@ -1630,22 +1620,16 @@ fn app_sdk_trade_mutation_requests_stay_locator_and_resync_owned() { assert!(source.contains("sdk.trades().buyer().propose_trade(sdk_request)")); for (label, start, end, required_count) in [ ( - "trade decision", + "trade proposal", + "fn trade_propose_with_sdk(", "fn trade_decision_with_sdk(", - "fn trade_revision_propose_with_sdk(", - 2usize, - ), - ( - "trade revision proposal", - "fn trade_revision_propose_with_sdk(", - "fn trade_revision_decide_with_sdk(", - 1usize, + 0usize, ), ( - "trade revision decision", - "fn trade_revision_decide_with_sdk(", + "trade decision", + "fn trade_decision_with_sdk(", "fn trade_cancel_with_sdk(", - 0usize, + 2usize, ), ( "trade cancellation", @@ -1667,19 +1651,6 @@ fn app_sdk_trade_mutation_requests_stay_locator_and_resync_owned() { "{label} must not use local-only evidence for app mutation authority" ); } - let revision_decision = source_segment( - source.as_str(), - "fn trade_revision_decide_with_sdk(", - "fn trade_cancel_with_sdk(", - ); - assert!( - revision_decision.contains("TradeEvidenceMode::require_explicit_evidence("), - "trade revision decision must pass reducer-visible runtime-store evidence into SDK explicit evidence mode" - ); - assert!( - revision_decision.contains("TradeEvidenceIngestRequest::new"), - "trade revision decision explicit evidence must be converted through SDK evidence ingest requests" - ); } #[test] @@ -1800,7 +1771,7 @@ const APP_STORE_RETIRED_SCHEMA_TERMS: &[&str] = &[ "dual_write", "dual-read", "dual-write", - "AppSdkMigration", + concat!("App", "SdkMigration"), "sdk_migration", "migration_receipt", "migration_audit", @@ -1841,12 +1812,17 @@ fn strict_sdk_boundary_scanner_rejects_unexcepted_new_production_paths() { "crates/runtime/src/sdk.rs", "fn mutate() { sdk.trades().ingest_evidence(TradeEvidenceIngestRequest::new(event)); }", ); - assert_eq!(evidence_findings.len(), 1); + assert_eq!(evidence_findings.len(), 2); assert!( evidence_findings .iter() .any(|finding| finding.pattern == ".ingest_evidence(") ); + assert!( + evidence_findings + .iter() + .any(|finding| finding.pattern == "TradeEvidenceIngestRequest") + ); let unexcepted_evidence_type_findings = unexcepted_sdk_boundary_patterns( "crates/desktop/src/runtime.rs", "fn mutate() { let _ = TradeEvidenceIngestRequest::new(event); }", diff --git a/crates/desktop/src/window.rs b/crates/desktop/src/window.rs @@ -13,7 +13,8 @@ use gpui_component::{ }; use radroots_nostr::prelude::RadrootsNostrClient; use radroots_studio_app_core::{ - AppSdkLifecycleState, AppSdkProjectionLifecycleState, AppSdkRelayUrlPolicy, + DesktopRuntimeLifecycleState, DesktopRuntimeProjectionLifecycleState, + DesktopRuntimeRelayUrlPolicy, }; use radroots_studio_app_i18n::{AppTextKey, app_text}; use radroots_studio_app_remote_signer::{ @@ -103,9 +104,10 @@ use crate::pack_day_print::{ use crate::runtime::{ DesktopAppRuntime, DesktopAppRuntimeProductEditorSaveError, DesktopAppRuntimeProductStockUpdateError, DesktopAppRuntimeSummary, - DesktopAppSdkDiagnosticsState, DesktopAppSdkDiagnosticsSummary, DesktopAppSdkIssueSummary, - DesktopAppSdkReadyDiagnosticsSummary, DesktopAppSdkStatusSummary, DesktopAppSyncConflictSummary, DesktopAppSyncStatusSummary, + DesktopRuntimeSupervisorDiagnosticsState, DesktopRuntimeSupervisorDiagnosticsSummary, + DesktopRuntimeSupervisorIssueSummary, DesktopRuntimeSupervisorReadyDiagnosticsSummary, + DesktopRuntimeSupervisorStatusSummary, }; const HOME_WINDOW_MIN_WIDTH_PX: f32 = 1080.0; @@ -6231,7 +6233,7 @@ impl SettingsFarmPanelState { .farm_profile .as_ref() .map(|farm_profile| farm_profile.farm_id) - .unwrap_or_else(FarmId::new); + .unwrap_or_else(FarmId::generate); let initial_draft = SettingsFarmRulesDraft::from_projection(farm_id, &projection); let farm_name_input = cx.new(|cx| { InputState::new(window, cx) @@ -8184,7 +8186,7 @@ fn settings_account_activation_key( fn about_status_rows( runtime: &DesktopAppRuntimeSummary, - sdk_diagnostics: Option<&DesktopAppSdkDiagnosticsSummary>, + sdk_diagnostics: Option<&DesktopRuntimeSupervisorDiagnosticsSummary>, ) -> Vec<LabelValueRow> { let mut rows = vec![ LabelValueRow::new( @@ -8248,8 +8250,8 @@ fn about_status_rows( fn append_sdk_status_rows( rows: &mut Vec<LabelValueRow>, - sdk_status: Option<&DesktopAppSdkStatusSummary>, - sdk_diagnostics: Option<&DesktopAppSdkDiagnosticsSummary>, + sdk_status: Option<&DesktopRuntimeSupervisorStatusSummary>, + sdk_diagnostics: Option<&DesktopRuntimeSupervisorDiagnosticsSummary>, ) { let status = sdk_diagnostics .map(|diagnostics| &diagnostics.status) @@ -8280,11 +8282,11 @@ fn append_sdk_status_rows( )); match sdk_diagnostics.map(|diagnostics| &diagnostics.state) { - Some(DesktopAppSdkDiagnosticsState::Ready(ready)) => { + Some(DesktopRuntimeSupervisorDiagnosticsState::Ready(ready)) => { append_ready_sdk_rows(rows, ready); append_sdk_issue_rows(rows, status.last_issue.as_ref()); } - Some(DesktopAppSdkDiagnosticsState::Blocked(issue)) => { + Some(DesktopRuntimeSupervisorDiagnosticsState::Blocked(issue)) => { rows.push(LabelValueRow::new( app_shared_text(AppTextKey::MetadataSdkDiagnosticState), app_text(AppTextKey::ValueSdkDiagnosticsBlocked), @@ -8303,7 +8305,7 @@ fn append_sdk_status_rows( fn append_ready_sdk_rows( rows: &mut Vec<LabelValueRow>, - ready: &DesktopAppSdkReadyDiagnosticsSummary, + ready: &DesktopRuntimeSupervisorReadyDiagnosticsSummary, ) { rows.push(LabelValueRow::new( app_shared_text(AppTextKey::MetadataSdkDiagnosticState), @@ -8339,7 +8341,10 @@ fn append_ready_sdk_rows( )); } -fn append_sdk_issue_rows(rows: &mut Vec<LabelValueRow>, issue: Option<&DesktopAppSdkIssueSummary>) { +fn append_sdk_issue_rows( + rows: &mut Vec<LabelValueRow>, + issue: Option<&DesktopRuntimeSupervisorIssueSummary>, +) { rows.push(LabelValueRow::new( app_shared_text(AppTextKey::MetadataSdkLastIssueCode), issue @@ -8593,34 +8598,36 @@ fn path_or_none(path: Option<&PathBuf>) -> String { .unwrap_or_else(|| app_text(AppTextKey::ValueNone)) } -fn sdk_lifecycle_state_text(state: AppSdkLifecycleState) -> String { +fn sdk_lifecycle_state_text(state: DesktopRuntimeLifecycleState) -> String { app_text(match state { - AppSdkLifecycleState::Starting => AppTextKey::ValueSdkLifecycleStarting, - AppSdkLifecycleState::Ready => AppTextKey::ValueSdkLifecycleReady, - AppSdkLifecycleState::Degraded => AppTextKey::ValueSdkLifecycleDegraded, - AppSdkLifecycleState::Pausing => AppTextKey::ValueSdkLifecyclePausing, - AppSdkLifecycleState::Paused => AppTextKey::ValueSdkLifecyclePaused, - AppSdkLifecycleState::Restoring => AppTextKey::ValueSdkLifecycleRestoring, - AppSdkLifecycleState::RebuildingProjections => { + DesktopRuntimeLifecycleState::Starting => AppTextKey::ValueSdkLifecycleStarting, + DesktopRuntimeLifecycleState::Ready => AppTextKey::ValueSdkLifecycleReady, + DesktopRuntimeLifecycleState::Degraded => AppTextKey::ValueSdkLifecycleDegraded, + DesktopRuntimeLifecycleState::Pausing => AppTextKey::ValueSdkLifecyclePausing, + DesktopRuntimeLifecycleState::Paused => AppTextKey::ValueSdkLifecyclePaused, + DesktopRuntimeLifecycleState::Restoring => AppTextKey::ValueSdkLifecycleRestoring, + DesktopRuntimeLifecycleState::RebuildingProjections => { AppTextKey::ValueSdkLifecycleRebuildingProjections } - AppSdkLifecycleState::ShuttingDown => AppTextKey::ValueSdkLifecycleShuttingDown, - AppSdkLifecycleState::Stopped => AppTextKey::ValueSdkLifecycleStopped, + DesktopRuntimeLifecycleState::ShuttingDown => AppTextKey::ValueSdkLifecycleShuttingDown, + DesktopRuntimeLifecycleState::Stopped => AppTextKey::ValueSdkLifecycleStopped, }) } -fn sdk_projection_lifecycle_state_text(state: AppSdkProjectionLifecycleState) -> String { +fn sdk_projection_lifecycle_state_text(state: DesktopRuntimeProjectionLifecycleState) -> String { app_text(match state { - AppSdkProjectionLifecycleState::Current => AppTextKey::ValueSdkProjectionCurrent, - AppSdkProjectionLifecycleState::Stale => AppTextKey::ValueSdkProjectionStale, - AppSdkProjectionLifecycleState::Rebuilding => AppTextKey::ValueSdkProjectionRebuilding, + DesktopRuntimeProjectionLifecycleState::Current => AppTextKey::ValueSdkProjectionCurrent, + DesktopRuntimeProjectionLifecycleState::Stale => AppTextKey::ValueSdkProjectionStale, + DesktopRuntimeProjectionLifecycleState::Rebuilding => { + AppTextKey::ValueSdkProjectionRebuilding + } }) } -fn sdk_relay_url_policy_text(policy: AppSdkRelayUrlPolicy) -> String { +fn sdk_relay_url_policy_text(policy: DesktopRuntimeRelayUrlPolicy) -> String { app_text(match policy { - AppSdkRelayUrlPolicy::Public => AppTextKey::ValueSdkRelayPolicyPublic, - AppSdkRelayUrlPolicy::Localhost => AppTextKey::ValueSdkRelayPolicyLocalhost, + DesktopRuntimeRelayUrlPolicy::Public => AppTextKey::ValueSdkRelayPolicyPublic, + DesktopRuntimeRelayUrlPolicy::Localhost => AppTextKey::ValueSdkRelayPolicyLocalhost, }) } @@ -16271,15 +16278,16 @@ mod tests { trade_workflow_source_key, }; use crate::runtime::{ - DesktopAppRuntimeMetadataSummary, DesktopAppRuntimeSummary, DesktopAppSdkDiagnosticsState, - DesktopAppSdkDiagnosticsSummary, DesktopAppSdkIssueSummary, - DesktopAppSdkReadyDiagnosticsSummary, DesktopAppSdkStatusSummary, - DesktopAppSyncConflictSummary, DesktopAppSyncStatusSummary, + DesktopAppRuntimeMetadataSummary, DesktopAppRuntimeSummary, DesktopAppSyncConflictSummary, + DesktopAppSyncStatusSummary, DesktopRuntimeSupervisorDiagnosticsState, + DesktopRuntimeSupervisorDiagnosticsSummary, DesktopRuntimeSupervisorIssueSummary, + DesktopRuntimeSupervisorReadyDiagnosticsSummary, DesktopRuntimeSupervisorStatusSummary, }; use radroots_identity::RadrootsIdentity; use radroots_studio_app_core::{ AppDesktopRuntimePaths, AppRuntimeHostEnvironment, AppRuntimePlatform, - AppSdkLifecycleState, AppSdkProjectionLifecycleState, AppSdkRelayUrlPolicy, + DesktopRuntimeLifecycleState, DesktopRuntimeProjectionLifecycleState, + DesktopRuntimeRelayUrlPolicy, }; use radroots_studio_app_remote_signer::{ RadrootsAppRemoteSignerApprovedSession, RadrootsAppRemoteSignerPendingSession, @@ -18594,21 +18602,23 @@ mod tests { #[test] fn about_status_rows_surface_ready_sdk_diagnostics() { - let sdk_status = fixture_sdk_status(AppSdkLifecycleState::Ready); - let sdk_diagnostics = DesktopAppSdkDiagnosticsSummary { + let sdk_status = fixture_sdk_status(DesktopRuntimeLifecycleState::Ready); + let sdk_diagnostics = DesktopRuntimeSupervisorDiagnosticsSummary { status: sdk_status.clone(), - state: DesktopAppSdkDiagnosticsState::Ready(DesktopAppSdkReadyDiagnosticsSummary { - storage_kind: "directory".to_owned(), - event_store_total_events: 7, - outbox_total_events: 3, - outbox_pending_events: 2, - outbox_failed_terminal_events: 0, - integrity_event_store_ok: true, - integrity_outbox_ok: true, - sync_source: "sdk_canonical_stores".to_owned(), - sync_observed_at_ms: 42, - sync_relay_target_count: 2, - }), + state: DesktopRuntimeSupervisorDiagnosticsState::Ready( + DesktopRuntimeSupervisorReadyDiagnosticsSummary { + storage_kind: "directory".to_owned(), + event_store_total_events: 7, + outbox_total_events: 3, + outbox_pending_events: 2, + outbox_failed_terminal_events: 0, + integrity_event_store_ok: true, + integrity_outbox_ok: true, + sync_source: "sdk_canonical_stores".to_owned(), + sync_observed_at_ms: 42, + sync_relay_target_count: 2, + }, + ), }; let mut runtime = summary( HomeRoute::Today, @@ -18646,17 +18656,17 @@ mod tests { #[test] fn about_status_rows_surface_blocked_sdk_issue_metadata() { - let issue = DesktopAppSdkIssueSummary { + let issue = DesktopRuntimeSupervisorIssueSummary { code: "invalid_relay_url".to_owned(), class: "configuration".to_owned(), retryable: false, recovery_actions: vec!["configure_transport_targets".to_owned()], }; - let mut sdk_status = fixture_sdk_status(AppSdkLifecycleState::Degraded); + let mut sdk_status = fixture_sdk_status(DesktopRuntimeLifecycleState::Degraded); sdk_status.last_issue = Some(issue.clone()); - let sdk_diagnostics = DesktopAppSdkDiagnosticsSummary { + let sdk_diagnostics = DesktopRuntimeSupervisorDiagnosticsSummary { status: sdk_status.clone(), - state: DesktopAppSdkDiagnosticsState::Blocked(issue), + state: DesktopRuntimeSupervisorDiagnosticsState::Blocked(issue), }; let mut runtime = summary( HomeRoute::Today, @@ -18708,7 +18718,7 @@ mod tests { database_schema_version: Some(7), ..DesktopAppRuntimeMetadataSummary::default() }; - runtime.sdk_status = Some(fixture_sdk_status(AppSdkLifecycleState::Ready)); + runtime.sdk_status = Some(fixture_sdk_status(DesktopRuntimeLifecycleState::Ready)); let rows = about_runtime_rows(&runtime); @@ -18780,18 +18790,20 @@ mod tests { } } - fn fixture_sdk_status(lifecycle_state: AppSdkLifecycleState) -> DesktopAppSdkStatusSummary { + fn fixture_sdk_status( + lifecycle_state: DesktopRuntimeLifecycleState, + ) -> DesktopRuntimeSupervisorStatusSummary { let storage_root = PathBuf::from("/tmp/radroots/data/apps/app/sdk"); - DesktopAppSdkStatusSummary { + DesktopRuntimeSupervisorStatusSummary { lifecycle_state, - projection_lifecycle_state: AppSdkProjectionLifecycleState::Current, + projection_lifecycle_state: DesktopRuntimeProjectionLifecycleState::Current, projection_lifecycle_reason: None, storage_root: storage_root.clone(), runtime_path: Some(storage_root.join("runtime.sqlite")), private_path: Some(storage_root.join("private.sqlite")), studio_path: Some(storage_root.join("studio.sqlite")), relay_target_count: 2, - relay_url_policy: AppSdkRelayUrlPolicy::Localhost, + relay_url_policy: DesktopRuntimeRelayUrlPolicy::Localhost, last_issue: None, } } diff --git a/crates/runtime/src/lib.rs b/crates/runtime/src/lib.rs @@ -37,17 +37,23 @@ pub use runtime::{ runtime_mode_label, }; pub use sdk::{ - APP_SDK_DEFAULT_COMMAND_QUEUE_CAPACITY, APP_SDK_STORAGE_DIR_NAME, AppSdkConfig, - AppSdkDiagnostics, AppSdkEventStoreDiagnostics, AppSdkFarmPublicLocationRequest, - AppSdkFarmPublishRequest, AppSdkIntegrityDiagnostics, AppSdkLifecycleState, - AppSdkListingPublishRequest, AppSdkOutboxDiagnostics, AppSdkProjectionLifecycleState, - AppSdkProjectionLifecycleStatus, AppSdkPublicFarmLocation, AppSdkRelayUrlPolicy, - AppSdkRestorePreflightReceipt, AppSdkRestorePreflightRequest, AppSdkRuntime, - AppSdkRuntimeError, AppSdkRuntimeIssue, AppSdkRuntimeStatus, AppSdkSqliteStoreDiagnostics, - AppSdkStorageDiagnostics, AppSdkStoragePaths, AppSdkSyncDiagnostics, - AppSdkSyncEventStoreDiagnostics, AppSdkSyncOutboxDiagnostics, - AppSdkSyncTransportTargetDiagnostics, AppSdkTradeCancellationRequest, AppSdkTradeDecision, - AppSdkTradeDecisionRequest, AppSdkTradeProposeRequest, AppSdkWorkflowReceipt, - app_sdk_storage_root_from_data_root, + DESKTOP_RUNTIME_DEFAULT_EFFECT_QUEUE_CAPACITY, DESKTOP_RUNTIME_STORAGE_DIR_NAME, + DesktopRuntimeDiagnostics, DesktopRuntimeEffectKind, DesktopRuntimeEffectReceipt, + DesktopRuntimeEffectState, DesktopRuntimeEffectStatus, DesktopRuntimeEventStoreDiagnostics, + DesktopRuntimeFarmPublicLocationRequest, DesktopRuntimeFarmPublishRequest, + DesktopRuntimeIntegrityDiagnostics, DesktopRuntimeIssue, DesktopRuntimeLifecycleState, + DesktopRuntimeListingPublishRequest, DesktopRuntimeLocalSigner, + DesktopRuntimeOutboxDiagnostics, DesktopRuntimeProjectionLifecycleState, + DesktopRuntimeProjectionLifecycleStatus, DesktopRuntimePublicFarmLocation, + DesktopRuntimeRelayUrlPolicy, DesktopRuntimeRestorePreflightReceipt, + DesktopRuntimeRestorePreflightRequest, DesktopRuntimeSnapshot, + DesktopRuntimeSqliteStoreDiagnostics, DesktopRuntimeStartupMilestone, + DesktopRuntimeStorageDiagnostics, DesktopRuntimeStoragePaths, DesktopRuntimeSupervisor, + DesktopRuntimeSupervisorConfig, DesktopRuntimeSupervisorError, DesktopRuntimeSyncDiagnostics, + DesktopRuntimeSyncEventStoreDiagnostics, DesktopRuntimeSyncOutboxDiagnostics, + DesktopRuntimeSyncTransportTargetDiagnostics, DesktopRuntimeTradeCancellationRequest, + DesktopRuntimeTradeDecision, DesktopRuntimeTradeDecisionRequest, + DesktopRuntimeTradeProposeRequest, DesktopRuntimeWorkflowReceipt, + desktop_runtime_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 @@ -1,25 +1,27 @@ use std::{ - fmt, io, + fmt, + future::Future, + io, path::{Path, PathBuf}, sync::{ - Arc, Condvar, Mutex, MutexGuard, - atomic::{AtomicBool, Ordering}, - mpsc::{self, Receiver, SyncSender, TrySendError}, + Arc, Mutex, MutexGuard, + atomic::{AtomicBool, AtomicU64, Ordering}, + mpsc::{self, Receiver, SyncSender, TryRecvError, TrySendError}, }, - thread::{self, JoinHandle}, - time::{Duration, Instant}, + thread, + time::Duration, }; use radroots_authority::{RadrootsActorContext, RadrootsLocalEventSigner}; use radroots_event::{ RadrootsEventPtr, contract::RadrootsActorRole, - farm::RadrootsFarm, + farm::{RadrootsFarm, RadrootsFarmPublicLocation}, ids::{ RadrootsAddressableCoordinate, RadrootsListingAddress, RadrootsOrderId, RadrootsPublicKey, }, kinds::KIND_FARM, - listing::RadrootsListing, + listing::{RadrootsListing, RadrootsListingPublicLocation}, order::{ RadrootsOrderEconomics, RadrootsOrderInventoryCommitment, RadrootsOrderItem, RadrootsOrderRequest, @@ -49,17 +51,18 @@ use tokio::runtime::Builder as TokioRuntimeBuilder; use crate::AppDesktopRuntimePaths; -pub const APP_SDK_STORAGE_DIR_NAME: &str = "sdk"; -pub const APP_SDK_DEFAULT_COMMAND_QUEUE_CAPACITY: usize = 32; +pub const DESKTOP_RUNTIME_STORAGE_DIR_NAME: &str = "sdk"; +pub const DESKTOP_RUNTIME_DEFAULT_EFFECT_QUEUE_CAPACITY: usize = 32; +const DESKTOP_RUNTIME_SDK_EFFECT_TIMEOUT_MS: u64 = 2_000; #[derive(Clone, Copy, Debug, Eq, PartialEq)] -pub enum AppSdkRelayUrlPolicy { +pub enum DesktopRuntimeRelayUrlPolicy { Public, Localhost, } #[derive(Clone, Copy, Debug, Eq, PartialEq)] -pub enum AppSdkLifecycleState { +pub enum DesktopRuntimeLifecycleState { Starting, Ready, Degraded, @@ -71,23 +74,33 @@ pub enum AppSdkLifecycleState { Stopped, } +#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)] +pub enum DesktopRuntimeStartupMilestone { + ShellReady, + RuntimeStoreReady, + PrivateStoreReady, + SignerReady, + ProjectionsReady, + NetworkObserved, +} + #[derive(Clone, Debug, Eq, PartialEq)] -pub struct AppSdkConfig { +pub struct DesktopRuntimeSupervisorConfig { pub storage_root: PathBuf, pub relay_urls: Vec<String>, - pub relay_url_policy: AppSdkRelayUrlPolicy, - pub command_queue_capacity: usize, + pub relay_url_policy: DesktopRuntimeRelayUrlPolicy, + pub effect_queue_capacity: usize, } #[derive(Clone, Debug, Eq, PartialEq)] -pub struct AppSdkStoragePaths { +pub struct DesktopRuntimeStoragePaths { pub runtime_path: PathBuf, pub private_path: PathBuf, pub studio_path: PathBuf, } #[derive(Clone, Debug, PartialEq)] -pub struct AppSdkRuntimeIssue { +pub struct DesktopRuntimeIssue { pub code: String, pub class: String, pub retryable: bool, @@ -97,34 +110,39 @@ pub struct AppSdkRuntimeIssue { } #[derive(Clone, Debug, PartialEq)] -pub struct AppSdkRuntimeStatus { - pub state: AppSdkLifecycleState, +pub struct DesktopRuntimeSnapshot { + pub state: DesktopRuntimeLifecycleState, + pub startup_milestones: Vec<DesktopRuntimeStartupMilestone>, pub storage_root: PathBuf, pub relay_urls: Vec<String>, - pub relay_url_policy: AppSdkRelayUrlPolicy, - pub storage_paths: Option<AppSdkStoragePaths>, - pub last_issue: Option<AppSdkRuntimeIssue>, - pub projection_lifecycle: AppSdkProjectionLifecycleStatus, + pub relay_url_policy: DesktopRuntimeRelayUrlPolicy, + pub storage_paths: Option<DesktopRuntimeStoragePaths>, + pub last_issue: Option<DesktopRuntimeIssue>, + pub last_effect: Option<DesktopRuntimeEffectStatus>, + pub storage_diagnostics: Option<DesktopRuntimeStorageDiagnostics>, + pub integrity_diagnostics: Option<DesktopRuntimeIntegrityDiagnostics>, + pub sync_diagnostics: Option<DesktopRuntimeSyncDiagnostics>, + pub projection_lifecycle: DesktopRuntimeProjectionLifecycleStatus, } #[derive(Clone, Debug, PartialEq)] -pub struct AppSdkDiagnostics { - pub runtime: AppSdkRuntimeStatus, - pub storage: AppSdkStorageDiagnostics, - pub integrity: AppSdkIntegrityDiagnostics, - pub sync: AppSdkSyncDiagnostics, +pub struct DesktopRuntimeDiagnostics { + pub runtime: DesktopRuntimeSnapshot, + pub storage: DesktopRuntimeStorageDiagnostics, + pub integrity: DesktopRuntimeIntegrityDiagnostics, + pub sync: DesktopRuntimeSyncDiagnostics, } #[derive(Clone, Debug, PartialEq)] -pub struct AppSdkStorageDiagnostics { +pub struct DesktopRuntimeStorageDiagnostics { pub storage_kind: String, - pub paths: Option<AppSdkStoragePaths>, - pub event_store: AppSdkEventStoreDiagnostics, - pub outbox: AppSdkOutboxDiagnostics, + pub paths: Option<DesktopRuntimeStoragePaths>, + pub event_store: DesktopRuntimeEventStoreDiagnostics, + pub outbox: DesktopRuntimeOutboxDiagnostics, } #[derive(Clone, Debug, PartialEq)] -pub struct AppSdkSqliteStoreDiagnostics { +pub struct DesktopRuntimeSqliteStoreDiagnostics { pub schema_version: i64, pub journal_mode: String, pub foreign_keys_enabled: bool, @@ -134,8 +152,8 @@ pub struct AppSdkSqliteStoreDiagnostics { } #[derive(Clone, Debug, PartialEq)] -pub struct AppSdkEventStoreDiagnostics { - pub store: AppSdkSqliteStoreDiagnostics, +pub struct DesktopRuntimeEventStoreDiagnostics { + pub store: DesktopRuntimeSqliteStoreDiagnostics, pub total_events: i64, pub projection_eligible_events: i64, pub transport_observations: i64, @@ -144,8 +162,8 @@ pub struct AppSdkEventStoreDiagnostics { } #[derive(Clone, Debug, PartialEq)] -pub struct AppSdkOutboxDiagnostics { - pub store: AppSdkSqliteStoreDiagnostics, +pub struct DesktopRuntimeOutboxDiagnostics { + pub store: DesktopRuntimeSqliteStoreDiagnostics, pub total_events: i64, pub pending_events: i64, pub retryable_events: i64, @@ -158,7 +176,7 @@ pub struct AppSdkOutboxDiagnostics { } #[derive(Clone, Debug, PartialEq)] -pub struct AppSdkIntegrityDiagnostics { +pub struct DesktopRuntimeIntegrityDiagnostics { pub checked_paths: Vec<PathBuf>, pub event_store_ok: bool, pub outbox_ok: bool, @@ -167,16 +185,16 @@ pub struct AppSdkIntegrityDiagnostics { } #[derive(Clone, Debug, PartialEq)] -pub struct AppSdkSyncDiagnostics { +pub struct DesktopRuntimeSyncDiagnostics { pub source: String, pub observed_at_ms: i64, - pub event_store: AppSdkSyncEventStoreDiagnostics, - pub outbox: AppSdkSyncOutboxDiagnostics, - pub transport_targets: AppSdkSyncTransportTargetDiagnostics, + pub event_store: DesktopRuntimeSyncEventStoreDiagnostics, + pub outbox: DesktopRuntimeSyncOutboxDiagnostics, + pub transport_targets: DesktopRuntimeSyncTransportTargetDiagnostics, } #[derive(Clone, Debug, PartialEq)] -pub struct AppSdkSyncEventStoreDiagnostics { +pub struct DesktopRuntimeSyncEventStoreDiagnostics { pub total_events: i64, pub projection_eligible_events: i64, pub transport_observations: i64, @@ -185,7 +203,7 @@ pub struct AppSdkSyncEventStoreDiagnostics { } #[derive(Clone, Debug, PartialEq)] -pub struct AppSdkSyncOutboxDiagnostics { +pub struct DesktopRuntimeSyncOutboxDiagnostics { pub total_events: i64, pub pending_events: i64, pub retryable_events: i64, @@ -198,25 +216,25 @@ pub struct AppSdkSyncOutboxDiagnostics { } #[derive(Clone, Debug, PartialEq)] -pub struct AppSdkSyncTransportTargetDiagnostics { +pub struct DesktopRuntimeSyncTransportTargetDiagnostics { pub configured_count: usize, pub configured_targets: Vec<String>, } #[derive(Clone, Debug, PartialEq)] -pub struct AppSdkRestorePreflightRequest { +pub struct DesktopRuntimeRestorePreflightRequest { pub source: PathBuf, pub overwrite_existing_sdk_storage: bool, } #[derive(Clone, Debug, PartialEq, Eq)] -pub struct AppSdkFarmPublicLocationRequest { +pub struct DesktopRuntimeFarmPublicLocationRequest { pub actor_pubkey: String, pub farm_d_tag: String, } #[derive(Clone, Debug, PartialEq, Eq)] -pub struct AppSdkPublicFarmLocation { +pub struct DesktopRuntimePublicFarmLocation { pub primary: String, pub city: Option<String>, pub region: Option<String>, @@ -224,30 +242,44 @@ pub struct AppSdkPublicFarmLocation { pub geohash5: String, } -pub struct AppSdkFarmPublishRequest { +pub struct DesktopRuntimeLocalSigner { + keys: RadrootsNostrKeys, +} + +impl DesktopRuntimeLocalSigner { + pub fn from_local_identity_keys(keys: RadrootsNostrKeys) -> Self { + Self { keys } + } + + fn into_keys(self) -> RadrootsNostrKeys { + self.keys + } +} + +pub struct DesktopRuntimeFarmPublishRequest { pub actor_account_id: String, pub actor_pubkey: String, - pub signer_keys: RadrootsNostrKeys, + pub signer: DesktopRuntimeLocalSigner, pub farm: RadrootsFarm, pub target_relays: Vec<String>, - pub relay_url_policy: AppSdkRelayUrlPolicy, + pub relay_url_policy: DesktopRuntimeRelayUrlPolicy, pub idempotency_key: Option<String>, } -pub struct AppSdkListingPublishRequest { +pub struct DesktopRuntimeListingPublishRequest { pub actor_account_id: String, pub actor_pubkey: String, - pub signer_keys: RadrootsNostrKeys, + pub signer: DesktopRuntimeLocalSigner, pub listing: RadrootsListing, pub target_relays: Vec<String>, - pub relay_url_policy: AppSdkRelayUrlPolicy, + pub relay_url_policy: DesktopRuntimeRelayUrlPolicy, pub idempotency_key: Option<String>, } -pub struct AppSdkTradeProposeRequest { +pub struct DesktopRuntimeTradeProposeRequest { pub actor_account_id: String, pub actor_pubkey: String, - pub signer_keys: RadrootsNostrKeys, + pub signer: DesktopRuntimeLocalSigner, pub listing_event: RadrootsEventPtr, pub order_id: RadrootsOrderId, pub listing_addr: RadrootsListingAddress, @@ -259,7 +291,7 @@ pub struct AppSdkTradeProposeRequest { pub idempotency_key: Option<String>, } -pub enum AppSdkTradeDecision { +pub enum DesktopRuntimeTradeDecision { Accept { inventory_commitments: Vec<RadrootsOrderInventoryCommitment>, }, @@ -268,20 +300,20 @@ pub enum AppSdkTradeDecision { }, } -pub struct AppSdkTradeDecisionRequest { +pub struct DesktopRuntimeTradeDecisionRequest { pub actor_account_id: String, pub actor_pubkey: String, - pub signer_keys: RadrootsNostrKeys, + pub signer: DesktopRuntimeLocalSigner, pub locator: RadrootsTradeLocator, - pub decision: AppSdkTradeDecision, + pub decision: DesktopRuntimeTradeDecision, pub confirm_public_note: bool, pub idempotency_key: Option<String>, } -pub struct AppSdkTradeCancellationRequest { +pub struct DesktopRuntimeTradeCancellationRequest { pub actor_account_id: String, pub actor_pubkey: String, - pub signer_keys: RadrootsNostrKeys, + pub signer: DesktopRuntimeLocalSigner, pub locator: RadrootsTradeLocator, pub reason: String, pub confirm_public_note: bool, @@ -289,7 +321,7 @@ pub struct AppSdkTradeCancellationRequest { } #[derive(Clone, Debug, Eq, PartialEq)] -pub struct AppSdkWorkflowReceipt { +pub struct DesktopRuntimeWorkflowReceipt { pub operation_kind: String, pub expected_event_id: String, pub signed_event_id: String, @@ -301,23 +333,23 @@ pub struct AppSdkWorkflowReceipt { } #[derive(Clone, Debug, PartialEq)] -pub struct AppSdkRestorePreflightReceipt { +pub struct DesktopRuntimeRestorePreflightReceipt { pub source: PathBuf, pub destination: PathBuf, pub state: String, - pub destination_paths: Option<AppSdkStoragePaths>, - pub restored_paths: Option<AppSdkStoragePaths>, + pub destination_paths: Option<DesktopRuntimeStoragePaths>, + pub restored_paths: Option<DesktopRuntimeStoragePaths>, pub runtime_path: PathBuf, pub private_path: PathBuf, pub studio_path: PathBuf, pub manifest_path: PathBuf, - pub verification: AppSdkBackupVerificationDiagnostics, - pub source_storage: AppSdkStorageDiagnostics, - pub projection_lifecycle: AppSdkProjectionLifecycleStatus, + pub verification: DesktopRuntimeBackupVerificationDiagnostics, + pub source_storage: DesktopRuntimeStorageDiagnostics, + pub projection_lifecycle: DesktopRuntimeProjectionLifecycleStatus, } #[derive(Clone, Debug, Eq, PartialEq)] -pub struct AppSdkBackupVerificationDiagnostics { +pub struct DesktopRuntimeBackupVerificationDiagnostics { pub event_store_ok: bool, pub outbox_ok: bool, pub event_store_events: i64, @@ -325,103 +357,118 @@ pub struct AppSdkBackupVerificationDiagnostics { } #[derive(Clone, Debug, Eq, PartialEq)] -pub struct AppSdkProjectionLifecycleStatus { - pub state: AppSdkProjectionLifecycleState, +pub struct DesktopRuntimeProjectionLifecycleStatus { + pub state: DesktopRuntimeProjectionLifecycleState, pub reason: Option<String>, pub restore_source: Option<PathBuf>, } #[derive(Clone, Copy, Debug, Eq, PartialEq)] -pub enum AppSdkProjectionLifecycleState { +pub enum DesktopRuntimeProjectionLifecycleState { Current, Stale, Rebuilding, } +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum DesktopRuntimeEffectKind { + RefreshDiagnostics, + RestorePreflight, + FarmPublish, + ListingPublish, + TradePropose, + TradeDecision, + TradeCancellation, + BeginProjectionRebuild, + CompleteProjectionRebuild, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum DesktopRuntimeEffectState { + Accepted, + Completed, + Failed, +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct DesktopRuntimeEffectReceipt { + pub effect_id: u64, + pub effect_kind: DesktopRuntimeEffectKind, + pub operation_kind: Option<String>, + pub actor_pubkey: Option<String>, +} + +#[derive(Clone, Debug, PartialEq)] +pub struct DesktopRuntimeEffectStatus { + pub receipt: DesktopRuntimeEffectReceipt, + pub state: DesktopRuntimeEffectState, + pub issue: Option<DesktopRuntimeIssue>, + pub workflow_receipt: Option<DesktopRuntimeWorkflowReceipt>, + pub restore_preflight: Option<DesktopRuntimeRestorePreflightReceipt>, +} + #[derive(Debug, Error)] -pub enum AppSdkRuntimeError { - #[error("app sdk command queue capacity must be greater than zero")] - CommandQueueCapacityZero, - #[error("failed to start app sdk worker: {0}")] +pub enum DesktopRuntimeSupervisorError { + #[error("desktop runtime supervisor effect queue capacity must be greater than zero")] + EffectQueueCapacityZero, + #[error("failed to start desktop runtime supervisor worker: {0}")] WorkerSpawn(#[from] io::Error), - #[error("app sdk command queue is full")] - CommandQueueFull, - #[error("app sdk command queue is closed")] - CommandQueueClosed, - #[error("app sdk command response channel is closed")] - CommandResponseClosed, - #[error("app sdk command failed: {0}")] - CommandFailed(AppSdkRuntimeIssue), - #[error("app sdk shutdown acknowledgement failed")] - ShutdownAck, - #[error("app sdk worker failed to join")] - WorkerJoin, + #[error("desktop runtime supervisor effect queue is full")] + EffectQueueFull, + #[error("desktop runtime supervisor effect queue is closed")] + EffectQueueClosed, + #[error("desktop runtime supervisor is unavailable: {0}")] + Unavailable(DesktopRuntimeIssue), } #[derive(Debug)] -pub struct AppSdkRuntime { - command_sender: Mutex<Option<SyncSender<AppSdkWorkerCommand>>>, - shared: Arc<AppSdkRuntimeShared>, - worker: Mutex<Option<JoinHandle<()>>>, +pub struct DesktopRuntimeSupervisor { + command_sender: Mutex<Option<SyncSender<DesktopRuntimeEffect>>>, + shared: Arc<DesktopRuntimeSupervisorShared>, + next_effect_id: AtomicU64, } #[derive(Debug)] -struct AppSdkRuntimeShared { - status: Mutex<AppSdkRuntimeStatus>, - status_changed: Condvar, +struct DesktopRuntimeSupervisorShared { + status: Mutex<DesktopRuntimeSnapshot>, shutdown_requested: AtomicBool, } -enum AppSdkWorkerCommand { - StorageStatus(mpsc::Sender<Result<AppSdkStorageDiagnostics, AppSdkRuntimeIssue>>), - IntegrityStatus(mpsc::Sender<Result<AppSdkIntegrityDiagnostics, AppSdkRuntimeIssue>>), - SyncStatus(mpsc::Sender<Result<AppSdkSyncDiagnostics, AppSdkRuntimeIssue>>), - Diagnostics(mpsc::Sender<Result<AppSdkDiagnostics, AppSdkRuntimeIssue>>), +enum DesktopRuntimeEffect { + RefreshDiagnostics(DesktopRuntimeEffectReceipt), RestorePreflight( - AppSdkRestorePreflightRequest, - mpsc::Sender<Result<AppSdkRestorePreflightReceipt, AppSdkRuntimeIssue>>, - ), - FarmPublicLocation( - AppSdkFarmPublicLocationRequest, - mpsc::Sender<Result<Option<AppSdkPublicFarmLocation>, AppSdkRuntimeIssue>>, + DesktopRuntimeEffectReceipt, + DesktopRuntimeRestorePreflightRequest, ), EnqueueFarmPublish( - AppSdkFarmPublishRequest, - mpsc::Sender<Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue>>, + DesktopRuntimeEffectReceipt, + DesktopRuntimeFarmPublishRequest, ), EnqueueListingPublish( - AppSdkListingPublishRequest, - mpsc::Sender<Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue>>, + DesktopRuntimeEffectReceipt, + DesktopRuntimeListingPublishRequest, ), TradePropose( - AppSdkTradeProposeRequest, - mpsc::Sender<Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue>>, + DesktopRuntimeEffectReceipt, + DesktopRuntimeTradeProposeRequest, ), TradeDecision( - AppSdkTradeDecisionRequest, - mpsc::Sender<Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue>>, + DesktopRuntimeEffectReceipt, + DesktopRuntimeTradeDecisionRequest, ), TradeCancellation( - AppSdkTradeCancellationRequest, - mpsc::Sender<Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue>>, - ), - BeginProjectionRebuild( - mpsc::Sender<Result<AppSdkProjectionLifecycleStatus, AppSdkRuntimeIssue>>, - ), - CompleteProjectionRebuild( - mpsc::Sender<Result<AppSdkProjectionLifecycleStatus, AppSdkRuntimeIssue>>, + DesktopRuntimeEffectReceipt, + DesktopRuntimeTradeCancellationRequest, ), + BeginProjectionRebuild(DesktopRuntimeEffectReceipt), + CompleteProjectionRebuild(DesktopRuntimeEffectReceipt), } -impl fmt::Debug for AppSdkWorkerCommand { +impl fmt::Debug for DesktopRuntimeEffect { fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { match self { - Self::StorageStatus(_) => formatter.write_str("StorageStatus"), - Self::IntegrityStatus(_) => formatter.write_str("IntegrityStatus"), - Self::SyncStatus(_) => formatter.write_str("SyncStatus"), - Self::Diagnostics(_) => formatter.write_str("Diagnostics"), + Self::RefreshDiagnostics(_) => formatter.write_str("RefreshDiagnostics"), Self::RestorePreflight(_, _) => formatter.write_str("RestorePreflight"), - Self::FarmPublicLocation(_, _) => formatter.write_str("FarmPublicLocation"), Self::EnqueueFarmPublish(_, _) => formatter.write_str("EnqueueFarmPublish"), Self::EnqueueListingPublish(_, _) => formatter.write_str("EnqueueListingPublish"), Self::TradePropose(_, _) => formatter.write_str("TradePropose"), @@ -433,27 +480,27 @@ impl fmt::Debug for AppSdkWorkerCommand { } } -impl AppSdkConfig { +impl DesktopRuntimeSupervisorConfig { pub fn from_desktop_paths(paths: &AppDesktopRuntimePaths, relay_urls: Vec<String>) -> Self { Self::from_app_data_root(paths.app.data.as_path(), relay_urls) } pub fn from_app_data_root(data_root: &Path, relay_urls: Vec<String>) -> Self { Self { - storage_root: app_sdk_storage_root_from_data_root(data_root), - relay_url_policy: app_sdk_relay_url_policy(relay_urls.as_slice()), + storage_root: desktop_runtime_storage_root_from_data_root(data_root), + relay_url_policy: desktop_runtime_relay_url_policy(relay_urls.as_slice()), relay_urls, - command_queue_capacity: APP_SDK_DEFAULT_COMMAND_QUEUE_CAPACITY, + effect_queue_capacity: DESKTOP_RUNTIME_DEFAULT_EFFECT_QUEUE_CAPACITY, } } - pub fn with_command_queue_capacity(mut self, capacity: usize) -> Self { - self.command_queue_capacity = capacity; + pub fn with_effect_queue_capacity(mut self, capacity: usize) -> Self { + self.effect_queue_capacity = capacity; self } } -impl AppSdkRestorePreflightRequest { +impl DesktopRuntimeRestorePreflightRequest { pub fn new(source: impl Into<PathBuf>) -> Self { Self { source: source.into(), @@ -467,10 +514,10 @@ impl AppSdkRestorePreflightRequest { } } -impl AppSdkProjectionLifecycleStatus { +impl DesktopRuntimeProjectionLifecycleStatus { pub fn current() -> Self { Self { - state: AppSdkProjectionLifecycleState::Current, + state: DesktopRuntimeProjectionLifecycleState::Current, reason: None, restore_source: None, } @@ -478,7 +525,7 @@ impl AppSdkProjectionLifecycleStatus { fn stale(reason: impl Into<String>, restore_source: Option<PathBuf>) -> Self { Self { - state: AppSdkProjectionLifecycleState::Stale, + state: DesktopRuntimeProjectionLifecycleState::Stale, reason: Some(reason.into()), restore_source, } @@ -486,235 +533,235 @@ impl AppSdkProjectionLifecycleStatus { fn rebuilding(reason: impl Into<String>, restore_source: Option<PathBuf>) -> Self { Self { - state: AppSdkProjectionLifecycleState::Rebuilding, + state: DesktopRuntimeProjectionLifecycleState::Rebuilding, reason: Some(reason.into()), restore_source, } } } -impl AppSdkRuntime { - pub fn start(config: AppSdkConfig) -> Result<Self, AppSdkRuntimeError> { - if config.command_queue_capacity == 0 { - return Err(AppSdkRuntimeError::CommandQueueCapacityZero); +impl DesktopRuntimeSupervisor { + pub fn start( + config: DesktopRuntimeSupervisorConfig, + ) -> Result<Self, DesktopRuntimeSupervisorError> { + if config.effect_queue_capacity == 0 { + return Err(DesktopRuntimeSupervisorError::EffectQueueCapacityZero); } - let initial_status = - AppSdkRuntimeStatus::from_config(&config, AppSdkLifecycleState::Starting, None, None); - let shared = Arc::new(AppSdkRuntimeShared { + let initial_status = DesktopRuntimeSnapshot::from_config( + &config, + DesktopRuntimeLifecycleState::Starting, + Some(vec![DesktopRuntimeStartupMilestone::ShellReady]), + None, + None, + ); + let shared = Arc::new(DesktopRuntimeSupervisorShared { status: Mutex::new(initial_status), - status_changed: Condvar::new(), shutdown_requested: AtomicBool::new(false), }); - let (command_sender, command_receiver) = mpsc::sync_channel(config.command_queue_capacity); + let (command_sender, command_receiver) = mpsc::sync_channel(config.effect_queue_capacity); let worker_shared = Arc::clone(&shared); - let worker = thread::Builder::new() - .name("radroots-app-sdk-runtime".to_owned()) - .spawn(move || run_app_sdk_worker(config, worker_shared, command_receiver))?; + let _worker = thread::Builder::new() + .name("radroots-desktop-runtime-supervisor".to_owned()) + .spawn(move || run_desktop_runtime_worker(config, worker_shared, command_receiver))?; Ok(Self { command_sender: Mutex::new(Some(command_sender)), shared, - worker: Mutex::new(Some(worker)), + next_effect_id: AtomicU64::new(1), }) } - pub fn status(&self) -> AppSdkRuntimeStatus { + pub fn snapshot(&self) -> DesktopRuntimeSnapshot { lock_status(&self.shared).clone() } - pub fn storage_status(&self) -> Result<AppSdkStorageDiagnostics, AppSdkRuntimeError> { - self.run_command(AppSdkWorkerCommand::StorageStatus) - } - - pub fn integrity_status(&self) -> Result<AppSdkIntegrityDiagnostics, AppSdkRuntimeError> { - self.run_command(AppSdkWorkerCommand::IntegrityStatus) - } - - pub fn sync_status(&self) -> Result<AppSdkSyncDiagnostics, AppSdkRuntimeError> { - self.run_command(AppSdkWorkerCommand::SyncStatus) - } - - pub fn diagnostics(&self) -> Result<AppSdkDiagnostics, AppSdkRuntimeError> { - self.run_command(AppSdkWorkerCommand::Diagnostics) - } - - pub fn restore_preflight( + pub fn request_diagnostics_refresh( &self, - request: AppSdkRestorePreflightRequest, - ) -> Result<AppSdkRestorePreflightReceipt, AppSdkRuntimeError> { - self.run_command(|response_sender| { - AppSdkWorkerCommand::RestorePreflight(request, response_sender) - }) + ) -> Result<DesktopRuntimeEffectReceipt, DesktopRuntimeSupervisorError> { + self.submit_effect( + DesktopRuntimeEffectKind::RefreshDiagnostics, + None, + None, + DesktopRuntimeEffect::RefreshDiagnostics, + ) } - pub fn farm_public_location( + pub fn request_restore_preflight( &self, - request: AppSdkFarmPublicLocationRequest, - ) -> Result<Option<AppSdkPublicFarmLocation>, AppSdkRuntimeError> { - self.run_command(|response_sender| { - AppSdkWorkerCommand::FarmPublicLocation(request, response_sender) - }) + request: DesktopRuntimeRestorePreflightRequest, + ) -> Result<DesktopRuntimeEffectReceipt, DesktopRuntimeSupervisorError> { + self.submit_effect( + DesktopRuntimeEffectKind::RestorePreflight, + None, + None, + |receipt| DesktopRuntimeEffect::RestorePreflight(receipt, request), + ) } pub fn enqueue_farm_publish( &self, - request: AppSdkFarmPublishRequest, - ) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeError> { - self.run_command(|response_sender| { - AppSdkWorkerCommand::EnqueueFarmPublish(request, response_sender) - }) + request: DesktopRuntimeFarmPublishRequest, + ) -> Result<DesktopRuntimeEffectReceipt, DesktopRuntimeSupervisorError> { + let actor_pubkey = Some(request.actor_pubkey.clone()); + self.submit_effect( + DesktopRuntimeEffectKind::FarmPublish, + Some(FARM_PUBLISH_OPERATION_KIND.to_owned()), + actor_pubkey, + |receipt| DesktopRuntimeEffect::EnqueueFarmPublish(receipt, request), + ) } pub fn enqueue_listing_publish( &self, - request: AppSdkListingPublishRequest, - ) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeError> { - self.run_command(|response_sender| { - AppSdkWorkerCommand::EnqueueListingPublish(request, response_sender) - }) + request: DesktopRuntimeListingPublishRequest, + ) -> Result<DesktopRuntimeEffectReceipt, DesktopRuntimeSupervisorError> { + let actor_pubkey = Some(request.actor_pubkey.clone()); + self.submit_effect( + DesktopRuntimeEffectKind::ListingPublish, + Some(LISTING_PUBLISH_OPERATION_KIND.to_owned()), + actor_pubkey, + |receipt| DesktopRuntimeEffect::EnqueueListingPublish(receipt, request), + ) } pub fn trade_propose( &self, - request: AppSdkTradeProposeRequest, - ) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeError> { - self.run_command(|response_sender| { - AppSdkWorkerCommand::TradePropose(request, response_sender) - }) + request: DesktopRuntimeTradeProposeRequest, + ) -> Result<DesktopRuntimeEffectReceipt, DesktopRuntimeSupervisorError> { + let actor_pubkey = Some(request.actor_pubkey.clone()); + self.submit_effect( + DesktopRuntimeEffectKind::TradePropose, + Some(TRADE_SUBMIT_OPERATION_KIND.to_owned()), + actor_pubkey, + |receipt| DesktopRuntimeEffect::TradePropose(receipt, request), + ) } pub fn trade_decide( &self, - request: AppSdkTradeDecisionRequest, - ) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeError> { - self.run_command(|response_sender| { - AppSdkWorkerCommand::TradeDecision(request, response_sender) - }) + request: DesktopRuntimeTradeDecisionRequest, + ) -> Result<DesktopRuntimeEffectReceipt, DesktopRuntimeSupervisorError> { + let actor_pubkey = Some(request.actor_pubkey.clone()); + self.submit_effect( + DesktopRuntimeEffectKind::TradeDecision, + Some(TRADE_DECISION_OPERATION_KIND.to_owned()), + actor_pubkey, + |receipt| DesktopRuntimeEffect::TradeDecision(receipt, request), + ) } pub fn trade_cancel( &self, - request: AppSdkTradeCancellationRequest, - ) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeError> { - self.run_command(|response_sender| { - AppSdkWorkerCommand::TradeCancellation(request, response_sender) - }) + request: DesktopRuntimeTradeCancellationRequest, + ) -> Result<DesktopRuntimeEffectReceipt, DesktopRuntimeSupervisorError> { + let actor_pubkey = Some(request.actor_pubkey.clone()); + self.submit_effect( + DesktopRuntimeEffectKind::TradeCancellation, + Some(TRADE_CANCELLATION_OPERATION_KIND.to_owned()), + actor_pubkey, + |receipt| DesktopRuntimeEffect::TradeCancellation(receipt, request), + ) } pub fn begin_projection_rebuild( &self, - ) -> Result<AppSdkProjectionLifecycleStatus, AppSdkRuntimeError> { - self.run_command(AppSdkWorkerCommand::BeginProjectionRebuild) + ) -> Result<DesktopRuntimeEffectReceipt, DesktopRuntimeSupervisorError> { + self.submit_effect( + DesktopRuntimeEffectKind::BeginProjectionRebuild, + None, + None, + DesktopRuntimeEffect::BeginProjectionRebuild, + ) } pub fn complete_projection_rebuild( &self, - ) -> Result<AppSdkProjectionLifecycleStatus, AppSdkRuntimeError> { - self.run_command(AppSdkWorkerCommand::CompleteProjectionRebuild) - } - - pub fn wait_for_startup(&self, timeout: Duration) -> AppSdkRuntimeStatus { - let deadline = Instant::now() - .checked_add(timeout) - .unwrap_or_else(Instant::now); - let mut status = lock_status(&self.shared); - loop { - if !matches!(status.state, AppSdkLifecycleState::Starting) { - return status.clone(); - } - let now = Instant::now(); - if now >= deadline { - return status.clone(); - } - let remaining = deadline.saturating_duration_since(now); - let wait_result = self.shared.status_changed.wait_timeout(status, remaining); - let (next_status, timeout_result) = wait_result.unwrap_or_else(|poisoned| { - let (guard, timeout_result) = poisoned.into_inner(); - (guard, timeout_result) - }); - status = next_status; - if timeout_result.timed_out() { - return status.clone(); - } - } + ) -> Result<DesktopRuntimeEffectReceipt, DesktopRuntimeSupervisorError> { + self.submit_effect( + DesktopRuntimeEffectKind::CompleteProjectionRebuild, + None, + None, + DesktopRuntimeEffect::CompleteProjectionRebuild, + ) } - pub fn shutdown(&self) -> Result<(), AppSdkRuntimeError> { - if matches!(self.status().state, AppSdkLifecycleState::Stopped) { - return self.join_worker(); + pub fn request_shutdown(&self) -> bool { + if matches!(self.snapshot().state, DesktopRuntimeLifecycleState::Stopped) { + return false; } self.shared.shutdown_requested.store(true, Ordering::SeqCst); - transition_status_state(&self.shared, AppSdkLifecycleState::ShuttingDown); + transition_status_state(&self.shared, DesktopRuntimeLifecycleState::ShuttingDown); let command_sender = self .command_sender .lock() .unwrap_or_else(|poisoned| poisoned.into_inner()) .take(); drop(command_sender); - self.join_worker() - } - - fn join_worker(&self) -> Result<(), AppSdkRuntimeError> { - let mut worker = self - .worker - .lock() - .unwrap_or_else(|poisoned| poisoned.into_inner()); - let Some(worker) = worker.take() else { - return Ok(()); - }; - worker.join().map_err(|_| AppSdkRuntimeError::WorkerJoin) + true } - fn run_command<T>( + fn submit_effect( &self, - command: impl FnOnce(mpsc::Sender<Result<T, AppSdkRuntimeIssue>>) -> AppSdkWorkerCommand, - ) -> Result<T, AppSdkRuntimeError> { - let (response_sender, response_receiver) = mpsc::channel(); + effect_kind: DesktopRuntimeEffectKind, + operation_kind: Option<String>, + actor_pubkey: Option<String>, + effect: impl FnOnce(DesktopRuntimeEffectReceipt) -> DesktopRuntimeEffect, + ) -> Result<DesktopRuntimeEffectReceipt, DesktopRuntimeSupervisorError> { + let receipt = DesktopRuntimeEffectReceipt { + effect_id: self.next_effect_id.fetch_add(1, Ordering::SeqCst), + effect_kind, + operation_kind, + actor_pubkey, + }; let command_sender = { let command_sender = self .command_sender .lock() .unwrap_or_else(|poisoned| poisoned.into_inner()); if self.shared.shutdown_requested.load(Ordering::SeqCst) { - return Err(AppSdkRuntimeError::CommandQueueClosed); + return Err(DesktopRuntimeSupervisorError::EffectQueueClosed); } command_sender .as_ref() .cloned() - .ok_or(AppSdkRuntimeError::CommandQueueClosed)? + .ok_or(DesktopRuntimeSupervisorError::EffectQueueClosed)? }; - match command_sender.try_send(command(response_sender)) { - Ok(()) => {} - Err(TrySendError::Full(_)) => return Err(AppSdkRuntimeError::CommandQueueFull), + set_last_effect( + &self.shared, + DesktopRuntimeEffectStatus::accepted(receipt.clone()), + ); + match command_sender.try_send(effect(receipt.clone())) { + Ok(()) => Ok(receipt), + Err(TrySendError::Full(_)) => { + clear_last_effect(&self.shared, receipt.effect_id); + Err(DesktopRuntimeSupervisorError::EffectQueueFull) + } Err(TrySendError::Disconnected(_)) => { - return Err(AppSdkRuntimeError::CommandQueueClosed); + clear_last_effect(&self.shared, receipt.effect_id); + Err(DesktopRuntimeSupervisorError::EffectQueueClosed) } } - response_receiver - .recv() - .map_err(|_| AppSdkRuntimeError::CommandResponseClosed)? - .map_err(AppSdkRuntimeError::CommandFailed) } } -impl Drop for AppSdkRuntime { +impl Drop for DesktopRuntimeSupervisor { fn drop(&mut self) { - let _ = self.shutdown(); + let _ = self.request_shutdown(); } } -impl From<AppSdkRelayUrlPolicy> for NostrRelayUrlPolicy { - fn from(policy: AppSdkRelayUrlPolicy) -> Self { +impl From<DesktopRuntimeRelayUrlPolicy> for NostrRelayUrlPolicy { + fn from(policy: DesktopRuntimeRelayUrlPolicy) -> Self { match policy { - AppSdkRelayUrlPolicy::Public => Self::Public, - AppSdkRelayUrlPolicy::Localhost => Self::Localhost, + DesktopRuntimeRelayUrlPolicy::Public => Self::Public, + DesktopRuntimeRelayUrlPolicy::Localhost => Self::Localhost, } } } -impl From<&RadrootsSdkStoragePaths> for AppSdkStoragePaths { +impl From<&RadrootsSdkStoragePaths> for DesktopRuntimeStoragePaths { fn from(paths: &RadrootsSdkStoragePaths) -> Self { Self { runtime_path: paths.runtime_path.clone(), @@ -724,7 +771,7 @@ impl From<&RadrootsSdkStoragePaths> for AppSdkStoragePaths { } } -impl AppSdkRuntimeIssue { +impl DesktopRuntimeIssue { fn from_sdk_error(error: &RadrootsSdkError) -> Self { Self { code: error.code().to_owned(), @@ -759,7 +806,7 @@ impl AppSdkRuntimeIssue { } } - fn lifecycle_blocked(state: AppSdkLifecycleState) -> Self { + fn lifecycle_blocked(state: DesktopRuntimeLifecycleState) -> Self { Self { code: "sdk_lifecycle_busy".to_owned(), class: "runtime".to_owned(), @@ -777,7 +824,7 @@ impl AppSdkRuntimeIssue { } } -impl From<SdkPublicLocality> for AppSdkPublicFarmLocation { +impl From<SdkPublicLocality> for DesktopRuntimePublicFarmLocation { fn from(value: SdkPublicLocality) -> Self { Self { primary: value.primary, @@ -789,37 +836,90 @@ impl From<SdkPublicLocality> for AppSdkPublicFarmLocation { } } -impl fmt::Display for AppSdkRuntimeIssue { +impl fmt::Display for DesktopRuntimeIssue { fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { write!(formatter, "{}: {}", self.code, self.message) } } -impl AppSdkRuntimeStatus { +impl DesktopRuntimeSnapshot { fn from_config( - config: &AppSdkConfig, - state: AppSdkLifecycleState, - storage_paths: Option<AppSdkStoragePaths>, - last_issue: Option<AppSdkRuntimeIssue>, + config: &DesktopRuntimeSupervisorConfig, + state: DesktopRuntimeLifecycleState, + startup_milestones: Option<Vec<DesktopRuntimeStartupMilestone>>, + storage_paths: Option<DesktopRuntimeStoragePaths>, + last_issue: Option<DesktopRuntimeIssue>, ) -> Self { Self { state, + startup_milestones: startup_milestones.unwrap_or_default(), storage_root: config.storage_root.clone(), relay_urls: config.relay_urls.clone(), relay_url_policy: config.relay_url_policy, storage_paths, last_issue, - projection_lifecycle: AppSdkProjectionLifecycleStatus::current(), + last_effect: None, + storage_diagnostics: None, + integrity_diagnostics: None, + sync_diagnostics: None, + projection_lifecycle: DesktopRuntimeProjectionLifecycleStatus::current(), } } } -impl From<StorageStatusReceipt> for AppSdkStorageDiagnostics { +impl DesktopRuntimeDiagnostics { + pub fn from_snapshot(snapshot: DesktopRuntimeSnapshot) -> Option<Self> { + Some(Self { + storage: snapshot.storage_diagnostics.clone()?, + integrity: snapshot.integrity_diagnostics.clone()?, + sync: snapshot.sync_diagnostics.clone()?, + runtime: snapshot, + }) + } +} + +impl DesktopRuntimeEffectStatus { + fn accepted(receipt: DesktopRuntimeEffectReceipt) -> Self { + Self { + receipt, + state: DesktopRuntimeEffectState::Accepted, + issue: None, + workflow_receipt: None, + restore_preflight: None, + } + } + + fn completed( + receipt: DesktopRuntimeEffectReceipt, + workflow_receipt: Option<DesktopRuntimeWorkflowReceipt>, + restore_preflight: Option<DesktopRuntimeRestorePreflightReceipt>, + ) -> Self { + Self { + receipt, + state: DesktopRuntimeEffectState::Completed, + issue: None, + workflow_receipt, + restore_preflight, + } + } + + fn failed(receipt: DesktopRuntimeEffectReceipt, issue: DesktopRuntimeIssue) -> Self { + Self { + receipt, + state: DesktopRuntimeEffectState::Failed, + issue: Some(issue), + workflow_receipt: None, + restore_preflight: None, + } + } +} + +impl From<StorageStatusReceipt> for DesktopRuntimeStorageDiagnostics { fn from(receipt: StorageStatusReceipt) -> Self { Self { storage_kind: serialized_label(&receipt.storage), - paths: receipt.paths.as_ref().map(AppSdkStoragePaths::from), - event_store: AppSdkEventStoreDiagnostics { + paths: receipt.paths.as_ref().map(DesktopRuntimeStoragePaths::from), + event_store: DesktopRuntimeEventStoreDiagnostics { store: receipt.event_store.store.into(), total_events: receipt.event_store.total_events, projection_eligible_events: receipt.event_store.projection_eligible_events, @@ -827,7 +927,7 @@ impl From<StorageStatusReceipt> for AppSdkStorageDiagnostics { last_event_seq: receipt.event_store.last_event_seq, last_event_updated_at_ms: receipt.event_store.last_event_updated_at_ms, }, - outbox: AppSdkOutboxDiagnostics { + outbox: DesktopRuntimeOutboxDiagnostics { store: receipt.outbox.store.into(), total_events: receipt.outbox.total_events, pending_events: receipt.outbox.pending_events, @@ -843,7 +943,7 @@ impl From<StorageStatusReceipt> for AppSdkStorageDiagnostics { } } -impl From<radroots_sdk::SdkSqliteStoreStatus> for AppSdkSqliteStoreDiagnostics { +impl From<radroots_sdk::SdkSqliteStoreStatus> for DesktopRuntimeSqliteStoreDiagnostics { fn from(status: radroots_sdk::SdkSqliteStoreStatus) -> Self { Self { schema_version: status.schema_version, @@ -856,7 +956,7 @@ impl From<radroots_sdk::SdkSqliteStoreStatus> for AppSdkSqliteStoreDiagnostics { } } -impl From<IntegrityReceipt> for AppSdkIntegrityDiagnostics { +impl From<IntegrityReceipt> for DesktopRuntimeIntegrityDiagnostics { fn from(receipt: IntegrityReceipt) -> Self { Self { checked_paths: receipt.checked_paths, @@ -868,19 +968,19 @@ impl From<IntegrityReceipt> for AppSdkIntegrityDiagnostics { } } -impl From<SyncStatusReceipt> for AppSdkSyncDiagnostics { +impl From<SyncStatusReceipt> for DesktopRuntimeSyncDiagnostics { fn from(receipt: SyncStatusReceipt) -> Self { Self { source: serialized_label(&receipt.source), observed_at_ms: receipt.observed_at_ms, - event_store: AppSdkSyncEventStoreDiagnostics { + event_store: DesktopRuntimeSyncEventStoreDiagnostics { total_events: receipt.event_store.total_events, projection_eligible_events: receipt.event_store.projection_eligible_events, transport_observations: receipt.event_store.transport_observations, last_event_seq: receipt.event_store.last_event_seq, last_event_updated_at_ms: receipt.event_store.last_event_updated_at_ms, }, - outbox: AppSdkSyncOutboxDiagnostics { + outbox: DesktopRuntimeSyncOutboxDiagnostics { total_events: receipt.outbox.total_events, pending_events: receipt.outbox.pending_events, retryable_events: receipt.outbox.retryable_events, @@ -891,7 +991,7 @@ impl From<SyncStatusReceipt> for AppSdkSyncDiagnostics { last_attempt_at_ms: receipt.outbox.last_attempt_at_ms, last_error: receipt.outbox.last_error, }, - transport_targets: AppSdkSyncTransportTargetDiagnostics { + transport_targets: DesktopRuntimeSyncTransportTargetDiagnostics { configured_count: receipt.transport_profile.configured_transport_target_count, configured_targets: receipt .transport_profile @@ -904,7 +1004,7 @@ impl From<SyncStatusReceipt> for AppSdkSyncDiagnostics { } } -impl From<SdkBackupVerification> for AppSdkBackupVerificationDiagnostics { +impl From<SdkBackupVerification> for DesktopRuntimeBackupVerificationDiagnostics { fn from(verification: SdkBackupVerification) -> Self { Self { event_store_ok: verification.event_store_ok, @@ -915,11 +1015,11 @@ impl From<SdkBackupVerification> for AppSdkBackupVerificationDiagnostics { } } -impl AppSdkRestorePreflightReceipt { +impl DesktopRuntimeRestorePreflightReceipt { fn from_restore_receipt( receipt: RestoreReceipt, destination: PathBuf, - projection_lifecycle: AppSdkProjectionLifecycleStatus, + projection_lifecycle: DesktopRuntimeProjectionLifecycleStatus, ) -> Self { Self { source: receipt.source, @@ -928,11 +1028,11 @@ impl AppSdkRestorePreflightReceipt { destination_paths: receipt .destination_paths .as_ref() - .map(AppSdkStoragePaths::from), + .map(DesktopRuntimeStoragePaths::from), restored_paths: receipt .restored_paths .as_ref() - .map(AppSdkStoragePaths::from), + .map(DesktopRuntimeStoragePaths::from), runtime_path: receipt.runtime_path, private_path: receipt.private_path, studio_path: receipt.studio_path, @@ -944,27 +1044,29 @@ impl AppSdkRestorePreflightReceipt { } } -pub fn app_sdk_storage_root_from_data_root(data_root: &Path) -> PathBuf { - data_root.join(APP_SDK_STORAGE_DIR_NAME) +pub fn desktop_runtime_storage_root_from_data_root(data_root: &Path) -> PathBuf { + data_root.join(DESKTOP_RUNTIME_STORAGE_DIR_NAME) } -fn app_sdk_relay_url_policy(relay_urls: &[String]) -> AppSdkRelayUrlPolicy { +fn desktop_runtime_relay_url_policy(relay_urls: &[String]) -> DesktopRuntimeRelayUrlPolicy { if relay_urls .iter() .any(|relay_url| relay_url.trim().to_ascii_lowercase().starts_with("ws://")) { - AppSdkRelayUrlPolicy::Localhost + DesktopRuntimeRelayUrlPolicy::Localhost } else { - AppSdkRelayUrlPolicy::Public + DesktopRuntimeRelayUrlPolicy::Public } } -fn run_app_sdk_worker( - config: AppSdkConfig, - shared: Arc<AppSdkRuntimeShared>, - command_receiver: Receiver<AppSdkWorkerCommand>, +fn run_desktop_runtime_worker( + config: DesktopRuntimeSupervisorConfig, + shared: Arc<DesktopRuntimeSupervisorShared>, + command_receiver: Receiver<DesktopRuntimeEffect>, ) { - let runtime = match TokioRuntimeBuilder::new_current_thread() + let runtime = match TokioRuntimeBuilder::new_multi_thread() + .worker_threads(2) + .thread_name("radroots-desktop-runtime-supervisor-async") .enable_all() .build() { @@ -972,11 +1074,12 @@ fn run_app_sdk_worker( Err(error) => { replace_status( &shared, - AppSdkRuntimeStatus::from_config( + DesktopRuntimeSnapshot::from_config( &config, - AppSdkLifecycleState::Degraded, + DesktopRuntimeLifecycleState::Degraded, + Some(vec![DesktopRuntimeStartupMilestone::ShellReady]), None, - Some(AppSdkRuntimeIssue::runtime_error( + Some(DesktopRuntimeIssue::runtime_error( "tokio_runtime_init", error.to_string(), )), @@ -989,115 +1092,81 @@ fn run_app_sdk_worker( let mut sdk = match runtime.block_on(build_sdk_runtime(&config)) { Ok(sdk) => { - replace_status( - &shared, - AppSdkRuntimeStatus::from_config( - &config, - AppSdkLifecycleState::Ready, - sdk.storage_paths().map(AppSdkStoragePaths::from), - None, - ), + let storage_paths = sdk.storage_paths().map(DesktopRuntimeStoragePaths::from); + let mut ready_status = DesktopRuntimeSnapshot::from_config( + &config, + DesktopRuntimeLifecycleState::Ready, + Some(ready_startup_milestones(&config, storage_paths.as_ref())), + storage_paths, + None, ); + match block_on_sdk_result( + &runtime, + "desktop_runtime_initial_diagnostics", + collect_sdk_diagnostics(&sdk, ready_status.clone()), + ) { + Ok(diagnostics) => { + ready_status.storage_diagnostics = Some(diagnostics.storage); + ready_status.integrity_diagnostics = Some(diagnostics.integrity); + ready_status.sync_diagnostics = Some(diagnostics.sync); + } + Err(error) => { + ready_status.state = DesktopRuntimeLifecycleState::Degraded; + ready_status.last_issue = Some(error); + } + } + replace_status(&shared, ready_status); Some(sdk) } Err(error) => { replace_status( &shared, - AppSdkRuntimeStatus::from_config( + DesktopRuntimeSnapshot::from_config( &config, - AppSdkLifecycleState::Degraded, + DesktopRuntimeLifecycleState::Degraded, + Some(vec![DesktopRuntimeStartupMilestone::ShellReady]), None, - Some(AppSdkRuntimeIssue::from_sdk_error(&error)), + Some(DesktopRuntimeIssue::from_sdk_error(&error)), ), ); None } }; - while let Ok(command) = command_receiver.recv() { + loop { if shared.shutdown_requested.load(Ordering::SeqCst) { break; } - match command { - AppSdkWorkerCommand::StorageStatus(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.storage_status(StorageStatusRequest::new())) - .map(AppSdkStorageDiagnostics::from) - .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error)), - None => Err(runtime_unavailable_issue(&shared)), - } - }; - send_worker_result(&shared, response_sender, result); - } - AppSdkWorkerCommand::IntegrityStatus(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.integrity(IntegrityRequest::new())) - .map(AppSdkIntegrityDiagnostics::from) - .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error)), - None => Err(runtime_unavailable_issue(&shared)), - } - }; - send_worker_result(&shared, response_sender, result); - } - AppSdkWorkerCommand::SyncStatus(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.sync().status(SyncStatusRequest::new())) - .map(AppSdkSyncDiagnostics::from) - .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error)), - None => Err(runtime_unavailable_issue(&shared)), - } - }; - send_worker_result(&shared, response_sender, result); + let command = match command_receiver.try_recv() { + Ok(command) => command, + Err(TryRecvError::Empty) => { + thread::sleep(Duration::from_millis(10)); + continue; } - AppSdkWorkerCommand::Diagnostics(response_sender) => { + Err(TryRecvError::Disconnected) => break, + }; + + match command { + DesktopRuntimeEffect::RefreshDiagnostics(receipt) => { let result = if let Some(issue) = lifecycle_busy_issue(&shared) { Err(issue) } else { match sdk.as_ref() { - Some(sdk) => { - let mut runtime_status = lock_status(&shared).clone(); - runtime_status.last_issue = None; - runtime - .block_on(collect_sdk_diagnostics(sdk, runtime_status)) - .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error)) - } + Some(sdk) => refresh_supervisor_diagnostics(&runtime, &shared, sdk), None => Err(runtime_unavailable_issue(&shared)), } }; - send_worker_result(&shared, response_sender, result); + finish_effect_result(&shared, receipt, result.map(|()| None)); } - AppSdkWorkerCommand::RestorePreflight(request, response_sender) => { + DesktopRuntimeEffect::RestorePreflight(receipt, request) => { let result = match sdk.as_ref() { Some(_) => run_restore_preflight(&runtime, &shared, &config, request), None => Err(runtime_unavailable_issue(&shared)), }; - send_worker_result(&shared, response_sender, result); + finish_restore_preflight_result(&shared, receipt, result); } - AppSdkWorkerCommand::FarmPublicLocation(request, response_sender) => { - let result = if let Some(issue) = lifecycle_busy_issue(&shared) { - Err(issue) - } else { - match sdk.as_ref() { - Some(sdk) => farm_public_location_with_sdk(&runtime, sdk, request), - None => Err(runtime_unavailable_issue(&shared)), - } - }; - send_worker_result(&shared, response_sender, result); - } - AppSdkWorkerCommand::EnqueueFarmPublish(request, response_sender) => { + DesktopRuntimeEffect::EnqueueFarmPublish(receipt, request) => { let result = if let Some(issue) = lifecycle_busy_issue(&shared) { Err(issue) } else { @@ -1106,9 +1175,9 @@ fn run_app_sdk_worker( None => Err(runtime_unavailable_issue(&shared)), } }; - send_worker_result(&shared, response_sender, result); + finish_workflow_result(&shared, receipt, result); } - AppSdkWorkerCommand::EnqueueListingPublish(request, response_sender) => { + DesktopRuntimeEffect::EnqueueListingPublish(receipt, request) => { let result = if let Some(issue) = lifecycle_busy_issue(&shared) { Err(issue) } else { @@ -1117,9 +1186,9 @@ fn run_app_sdk_worker( None => Err(runtime_unavailable_issue(&shared)), } }; - send_worker_result(&shared, response_sender, result); + finish_workflow_result(&shared, receipt, result); } - AppSdkWorkerCommand::TradePropose(request, response_sender) => { + DesktopRuntimeEffect::TradePropose(receipt, request) => { let result = if let Some(issue) = lifecycle_busy_issue(&shared) { Err(issue) } else { @@ -1128,9 +1197,9 @@ fn run_app_sdk_worker( None => Err(runtime_unavailable_issue(&shared)), } }; - send_worker_result(&shared, response_sender, result); + finish_workflow_result(&shared, receipt, result); } - AppSdkWorkerCommand::TradeDecision(request, response_sender) => { + DesktopRuntimeEffect::TradeDecision(receipt, request) => { let result = if let Some(issue) = lifecycle_busy_issue(&shared) { Err(issue) } else { @@ -1139,9 +1208,9 @@ fn run_app_sdk_worker( None => Err(runtime_unavailable_issue(&shared)), } }; - send_worker_result(&shared, response_sender, result); + finish_workflow_result(&shared, receipt, result); } - AppSdkWorkerCommand::TradeCancellation(request, response_sender) => { + DesktopRuntimeEffect::TradeCancellation(receipt, request) => { let result = if let Some(issue) = lifecycle_busy_issue(&shared) { Err(issue) } else { @@ -1150,142 +1219,70 @@ fn run_app_sdk_worker( None => Err(runtime_unavailable_issue(&shared)), } }; - send_worker_result(&shared, response_sender, result); + finish_workflow_result(&shared, receipt, result); } - AppSdkWorkerCommand::BeginProjectionRebuild(response_sender) => { + DesktopRuntimeEffect::BeginProjectionRebuild(receipt) => { let result = match sdk.as_ref() { Some(_) => Ok(begin_projection_rebuild(&shared)), None => Err(runtime_unavailable_issue(&shared)), }; - send_worker_result(&shared, response_sender, result); + finish_projection_result(&shared, receipt, result); } - AppSdkWorkerCommand::CompleteProjectionRebuild(response_sender) => { + DesktopRuntimeEffect::CompleteProjectionRebuild(receipt) => { let result = match sdk.as_ref() { Some(_) => complete_projection_rebuild(&shared), None => Err(runtime_unavailable_issue(&shared)), }; - send_worker_result(&shared, response_sender, result); + finish_projection_result(&shared, receipt, result); } } } drop(sdk.take()); - transition_status_state(&shared, AppSdkLifecycleState::Stopped); + transition_status_state(&shared, DesktopRuntimeLifecycleState::Stopped); } fn run_degraded_worker( - config: AppSdkConfig, - shared: Arc<AppSdkRuntimeShared>, - command_receiver: Receiver<AppSdkWorkerCommand>, + config: DesktopRuntimeSupervisorConfig, + shared: Arc<DesktopRuntimeSupervisorShared>, + command_receiver: Receiver<DesktopRuntimeEffect>, ) { - while let Ok(command) = command_receiver.recv() { + loop { if shared.shutdown_requested.load(Ordering::SeqCst) { break; } - match command { - AppSdkWorkerCommand::StorageStatus(response_sender) => { - send_worker_result( - &shared, - response_sender, - Err(runtime_unavailable_issue(&shared)), - ); - } - AppSdkWorkerCommand::IntegrityStatus(response_sender) => { - send_worker_result( - &shared, - response_sender, - Err(runtime_unavailable_issue(&shared)), - ); - } - AppSdkWorkerCommand::SyncStatus(response_sender) => { - send_worker_result( - &shared, - response_sender, - Err(runtime_unavailable_issue(&shared)), - ); - } - AppSdkWorkerCommand::Diagnostics(response_sender) => { - send_worker_result( - &shared, - response_sender, - Err(runtime_unavailable_issue(&shared)), - ); - } - AppSdkWorkerCommand::RestorePreflight(_, response_sender) => { - send_worker_result( - &shared, - response_sender, - Err(runtime_unavailable_issue(&shared)), - ); - } - AppSdkWorkerCommand::FarmPublicLocation(_, response_sender) => { - send_worker_result( - &shared, - response_sender, - Err(runtime_unavailable_issue(&shared)), - ); - } - AppSdkWorkerCommand::EnqueueFarmPublish(_, response_sender) => { - send_worker_result( - &shared, - response_sender, - Err(runtime_unavailable_issue(&shared)), - ); - } - AppSdkWorkerCommand::EnqueueListingPublish(_, response_sender) => { - send_worker_result( - &shared, - response_sender, - Err(runtime_unavailable_issue(&shared)), - ); - } - AppSdkWorkerCommand::TradePropose(_, response_sender) => { - send_worker_result( - &shared, - response_sender, - Err(runtime_unavailable_issue(&shared)), - ); + let command = match command_receiver.try_recv() { + Ok(command) => command, + Err(TryRecvError::Empty) => { + thread::sleep(Duration::from_millis(10)); + continue; } - AppSdkWorkerCommand::TradeDecision(_, response_sender) => { - send_worker_result( - &shared, - response_sender, - Err(runtime_unavailable_issue(&shared)), - ); - } - AppSdkWorkerCommand::TradeCancellation(_, response_sender) => { - send_worker_result( - &shared, - response_sender, - Err(runtime_unavailable_issue(&shared)), - ); - } - AppSdkWorkerCommand::BeginProjectionRebuild(response_sender) => { - send_worker_result( - &shared, - response_sender, - Err(runtime_unavailable_issue(&shared)), - ); - } - AppSdkWorkerCommand::CompleteProjectionRebuild(response_sender) => { - send_worker_result( - &shared, - response_sender, - Err(runtime_unavailable_issue(&shared)), - ); - } - } + Err(TryRecvError::Disconnected) => break, + }; + set_effect_failed( + &shared, + effect_receipt(&command), + runtime_unavailable_issue(&shared), + ); } let last_issue = lock_status(&shared).last_issue.clone(); replace_status( &shared, - AppSdkRuntimeStatus::from_config(&config, AppSdkLifecycleState::Stopped, None, last_issue), + DesktopRuntimeSnapshot::from_config( + &config, + DesktopRuntimeLifecycleState::Stopped, + Some(vec![DesktopRuntimeStartupMilestone::ShellReady]), + None, + last_issue, + ), ); } -async fn build_sdk_runtime(config: &AppSdkConfig) -> Result<RadrootsClient, RadrootsSdkError> { +async fn build_sdk_runtime( + config: &DesktopRuntimeSupervisorConfig, +) -> Result<RadrootsClient, RadrootsSdkError> { RadrootsClient::builder() .directory_storage(config.storage_root.clone()) .transport_profile(app_transport_profile(config)?) @@ -1294,26 +1291,28 @@ async fn build_sdk_runtime(config: &AppSdkConfig) -> Result<RadrootsClient, Radr } async fn build_sdk_runtime_with_signer( - config: &AppSdkConfig, - keys: RadrootsNostrKeys, -) -> Result<RadrootsClient, AppSdkRuntimeIssue> { - let local_signer = RadrootsLocalEventSigner::new(keys).map_err(|error| { - AppSdkRuntimeIssue::runtime_error("sdk_signer_init_failed", error.to_string()) + config: &DesktopRuntimeSupervisorConfig, + signer: DesktopRuntimeLocalSigner, +) -> Result<RadrootsClient, DesktopRuntimeIssue> { + let local_signer = RadrootsLocalEventSigner::new(signer.into_keys()).map_err(|error| { + DesktopRuntimeIssue::runtime_error("sdk_signer_init_failed", error.to_string()) })?; let signer = RadrootsSdkLocalKeySigner::from_event_signer(local_signer) - .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))?; + .map_err(|error| DesktopRuntimeIssue::from_sdk_error(&error))?; let transport_profile = app_transport_profile(config) - .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))?; + .map_err(|error| DesktopRuntimeIssue::from_sdk_error(&error))?; RadrootsClient::builder() .directory_storage(config.storage_root.clone()) .transport_profile(transport_profile) .signer_provider(RadrootsSdkSignerProvider::LocalKey(signer)) .build() .await - .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error)) + .map_err(|error| DesktopRuntimeIssue::from_sdk_error(&error)) } -fn app_transport_profile(config: &AppSdkConfig) -> Result<TransportProfile, RadrootsSdkError> { +fn app_transport_profile( + config: &DesktopRuntimeSupervisorConfig, +) -> Result<TransportProfile, RadrootsSdkError> { if config.relay_urls.is_empty() { return Ok(TransportProfile::local_only()); } @@ -1344,52 +1343,96 @@ fn app_trade_privacy_confirmation(confirm_public_note: bool) -> PrivacyPreflight } } +fn block_on_sdk_result<T>( + runtime: &tokio::runtime::Runtime, + operation: &'static str, + future: impl Future<Output = Result<T, RadrootsSdkError>>, +) -> Result<T, DesktopRuntimeIssue> { + match runtime.block_on(async { + tokio::time::timeout( + Duration::from_millis(DESKTOP_RUNTIME_SDK_EFFECT_TIMEOUT_MS), + future, + ) + .await + }) { + Ok(Ok(value)) => Ok(value), + Ok(Err(error)) => Err(DesktopRuntimeIssue::from_sdk_error(&error)), + Err(_) => Err(sdk_effect_timeout_issue(operation)), + } +} + +fn block_on_desktop_result<T>( + runtime: &tokio::runtime::Runtime, + operation: &'static str, + future: impl Future<Output = Result<T, DesktopRuntimeIssue>>, +) -> Result<T, DesktopRuntimeIssue> { + match runtime.block_on(async { + tokio::time::timeout( + Duration::from_millis(DESKTOP_RUNTIME_SDK_EFFECT_TIMEOUT_MS), + future, + ) + .await + }) { + Ok(result) => result, + Err(_) => Err(sdk_effect_timeout_issue(operation)), + } +} + +fn sdk_effect_timeout_issue(operation: &'static str) -> DesktopRuntimeIssue { + DesktopRuntimeIssue::runtime_error( + "desktop_runtime_sdk_effect_timeout", + format!("{operation} did not complete within {DESKTOP_RUNTIME_SDK_EFFECT_TIMEOUT_MS} ms"), + ) +} + fn run_restore_preflight( runtime: &tokio::runtime::Runtime, - shared: &AppSdkRuntimeShared, - config: &AppSdkConfig, - request: AppSdkRestorePreflightRequest, -) -> Result<AppSdkRestorePreflightReceipt, AppSdkRuntimeIssue> { + shared: &DesktopRuntimeSupervisorShared, + config: &DesktopRuntimeSupervisorConfig, + request: DesktopRuntimeRestorePreflightRequest, +) -> Result<DesktopRuntimeRestorePreflightReceipt, DesktopRuntimeIssue> { if let Some(issue) = lifecycle_busy_issue(shared) { return Err(issue); } - transition_status_state(shared, AppSdkLifecycleState::Pausing); - transition_status_state(shared, AppSdkLifecycleState::Paused); - transition_status_state(shared, AppSdkLifecycleState::Restoring); + transition_status_state(shared, DesktopRuntimeLifecycleState::Pausing); + transition_status_state(shared, DesktopRuntimeLifecycleState::Paused); + transition_status_state(shared, DesktopRuntimeLifecycleState::Restoring); let restore_request = RestoreRequest::new(request.source.clone()) .with_destination(config.storage_root.clone()) .with_overwrite(request.overwrite_existing_sdk_storage) .dry_run(); - let result = runtime - .block_on(RadrootsClient::restore(restore_request)) - .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error)) - .map(|receipt| { - let projection_lifecycle = mark_projections_stale( - shared, - "sdk_restore_preflight", - Some(request.source.clone()), - ); - AppSdkRestorePreflightReceipt::from_restore_receipt( - receipt, - config.storage_root.clone(), - projection_lifecycle, - ) - }); + let result = block_on_sdk_result( + runtime, + "desktop_runtime_restore_preflight", + RadrootsClient::restore(restore_request), + ) + .map(|receipt| { + let projection_lifecycle = mark_projections_stale( + shared, + "sdk_restore_preflight", + Some(request.source.clone()), + ); + DesktopRuntimeRestorePreflightReceipt::from_restore_receipt( + receipt, + config.storage_root.clone(), + projection_lifecycle, + ) + }); if result.is_err() { - transition_status_state(shared, AppSdkLifecycleState::Ready); + transition_status_state(shared, DesktopRuntimeLifecycleState::Ready); } result } async fn collect_sdk_diagnostics( sdk: &RadrootsClient, - runtime: AppSdkRuntimeStatus, -) -> Result<AppSdkDiagnostics, RadrootsSdkError> { + runtime: DesktopRuntimeSnapshot, +) -> Result<DesktopRuntimeDiagnostics, RadrootsSdkError> { let storage = sdk.storage_status(StorageStatusRequest::new()).await?; let integrity = sdk.integrity(IntegrityRequest::new()).await?; let sync = sdk.sync().status(SyncStatusRequest::new()).await?; - Ok(AppSdkDiagnostics { + Ok(DesktopRuntimeDiagnostics { runtime, storage: storage.into(), integrity: integrity.into(), @@ -1400,91 +1443,124 @@ async fn collect_sdk_diagnostics( fn farm_public_location_with_sdk( runtime: &tokio::runtime::Runtime, sdk: &RadrootsClient, - request: AppSdkFarmPublicLocationRequest, -) -> Result<Option<AppSdkPublicFarmLocation>, AppSdkRuntimeIssue> { + request: DesktopRuntimeFarmPublicLocationRequest, +) -> Result<Option<DesktopRuntimePublicFarmLocation>, DesktopRuntimeIssue> { let farm_addr = RadrootsAddressableCoordinate::parse(format!( "{KIND_FARM}:{}:{}", request.actor_pubkey, request.farm_d_tag )) .map_err(|error| { - AppSdkRuntimeIssue::from_sdk_error(&RadrootsSdkError::InvalidRequest { + DesktopRuntimeIssue::from_sdk_error(&RadrootsSdkError::InvalidRequest { message: format!("farm public location address is invalid: {error}"), }) })?; - runtime - .block_on(sdk.farms().private_location(&farm_addr)) - .map(|location| location.map(|receipt| receipt.public_locality.into())) - .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error)) + block_on_sdk_result( + runtime, + "desktop_runtime_farm_public_location", + sdk.farms().private_location(&farm_addr), + ) + .map(|location| location.map(|receipt| receipt.public_locality.into())) } fn enqueue_farm_publish_with_sdk( runtime: &tokio::runtime::Runtime, sdk: &RadrootsClient, - request: AppSdkFarmPublishRequest, -) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue> { + request: DesktopRuntimeFarmPublishRequest, +) -> Result<DesktopRuntimeWorkflowReceipt, DesktopRuntimeIssue> { let actor = sdk_actor_context( request.actor_pubkey.as_str(), request.actor_account_id.as_str(), RadrootsActorRole::Farmer, )?; - let signer = sdk_local_signer(request.signer_keys)?; + let signer = sdk_local_signer(request.signer)?; let target_relays = sdk_transport_targets(request.target_relays, request.relay_url_policy)?; - let mut enqueue = FarmEnqueuePublishRequest::new(actor, request.farm, target_relays); + let mut farm = request.farm; + if farm.location.is_none() { + let public_location = farm_public_location_with_sdk( + runtime, + sdk, + DesktopRuntimeFarmPublicLocationRequest { + actor_pubkey: request.actor_pubkey.clone(), + farm_d_tag: farm.d_tag.clone(), + }, + )?; + farm.location = public_location.map(desktop_runtime_public_farm_location_to_protocol); + } + let mut enqueue = FarmEnqueuePublishRequest::new(actor, farm, 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))?; + .map_err(|error| DesktopRuntimeIssue::from_sdk_error(&error))?; } - let receipt = runtime - .block_on( - sdk.farms() - .enqueue_publish_with_explicit_signer(enqueue, &signer), - ) - .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))?; - Ok(app_sdk_farm_receipt(receipt, request.actor_pubkey)) + let receipt = block_on_sdk_result( + runtime, + "desktop_runtime_farm_publish", + sdk.farms() + .enqueue_publish_with_explicit_signer(enqueue, &signer), + )?; + Ok(desktop_runtime_farm_receipt(receipt, request.actor_pubkey)) } fn enqueue_listing_publish_with_sdk( runtime: &tokio::runtime::Runtime, sdk: &RadrootsClient, - request: AppSdkListingPublishRequest, -) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue> { + request: DesktopRuntimeListingPublishRequest, +) -> Result<DesktopRuntimeWorkflowReceipt, DesktopRuntimeIssue> { 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 signer = sdk_local_signer(request.signer)?; let target_relays = sdk_transport_targets(request.target_relays, request.relay_url_policy)?; - let mut enqueue = ListingEnqueuePublishRequest::new(actor, request.listing, target_relays); + let mut listing = request.listing; + if listing.location.is_none() { + let public_location = farm_public_location_with_sdk( + runtime, + sdk, + DesktopRuntimeFarmPublicLocationRequest { + actor_pubkey: listing.farm.pubkey.clone(), + farm_d_tag: listing.farm.d_tag.clone(), + }, + )?; + listing.location = public_location.map(desktop_runtime_public_listing_location_to_protocol); + } + let mut enqueue = ListingEnqueuePublishRequest::new(actor, listing, 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))?; + .map_err(|error| DesktopRuntimeIssue::from_sdk_error(&error))?; } - let receipt = runtime - .block_on( - sdk.listings() - .enqueue_publish_with_explicit_signer(enqueue, &signer), - ) - .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))?; - Ok(app_sdk_listing_receipt(receipt, request.actor_pubkey)) + let receipt = block_on_sdk_result( + runtime, + "desktop_runtime_listing_publish", + sdk.listings() + .enqueue_publish_with_explicit_signer(enqueue, &signer), + )?; + Ok(desktop_runtime_listing_receipt( + receipt, + request.actor_pubkey, + )) } fn trade_propose_with_sdk( runtime: &tokio::runtime::Runtime, - config: &AppSdkConfig, - request: AppSdkTradeProposeRequest, -) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue> { + config: &DesktopRuntimeSupervisorConfig, + request: DesktopRuntimeTradeProposeRequest, +) -> Result<DesktopRuntimeWorkflowReceipt, DesktopRuntimeIssue> { let actor = sdk_actor_context( request.actor_pubkey.as_str(), request.actor_account_id.as_str(), RadrootsActorRole::Buyer, )?; - let sdk = runtime.block_on(build_sdk_runtime_with_signer(config, request.signer_keys))?; + let sdk = block_on_desktop_result( + runtime, + "desktop_runtime_trade_propose_build", + build_sdk_runtime_with_signer(config, request.signer), + )?; let buyer_pubkey = RadrootsPublicKey::parse(request.actor_pubkey.as_str()).map_err(|error| { - AppSdkRuntimeIssue::from_sdk_error(&RadrootsSdkError::InvalidRequest { + DesktopRuntimeIssue::from_sdk_error(&RadrootsSdkError::InvalidRequest { message: format!("trade proposal buyer public key is invalid: {error}"), }) })?; @@ -1509,29 +1585,35 @@ fn trade_propose_with_sdk( 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))?; + .map_err(|error| DesktopRuntimeIssue::from_sdk_error(&error))?; } - let outcome = runtime - .block_on(sdk.trades().buyer().propose_trade(sdk_request)) - .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))?; - app_sdk_trade_propose_receipt(outcome, request.actor_pubkey) + let outcome = block_on_sdk_result( + runtime, + "desktop_runtime_trade_propose", + sdk.trades().buyer().propose_trade(sdk_request), + )?; + desktop_runtime_trade_propose_receipt(outcome, request.actor_pubkey) } fn trade_decision_with_sdk( runtime: &tokio::runtime::Runtime, - config: &AppSdkConfig, - request: AppSdkTradeDecisionRequest, -) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue> { + config: &DesktopRuntimeSupervisorConfig, + request: DesktopRuntimeTradeDecisionRequest, +) -> Result<DesktopRuntimeWorkflowReceipt, DesktopRuntimeIssue> { let actor = sdk_actor_context( request.actor_pubkey.as_str(), request.actor_account_id.as_str(), RadrootsActorRole::Seller, )?; - let sdk = runtime.block_on(build_sdk_runtime_with_signer(config, request.signer_keys))?; + let sdk = block_on_desktop_result( + runtime, + "desktop_runtime_trade_decision_build", + build_sdk_runtime_with_signer(config, request.signer), + )?; let publish_mode = app_trade_publish_mode(); let satisfaction_policy = app_trade_satisfaction_policy(); let outcome = match request.decision { - AppSdkTradeDecision::Accept { + DesktopRuntimeTradeDecision::Accept { inventory_commitments, } => { let mut sdk_request = TradeAcceptRequest::new( @@ -1547,13 +1629,15 @@ fn trade_decision_with_sdk( 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))?; + .map_err(|error| DesktopRuntimeIssue::from_sdk_error(&error))?; } - runtime - .block_on(sdk.trades().seller().accept_trade(sdk_request)) - .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))? + block_on_sdk_result( + runtime, + "desktop_runtime_trade_accept", + sdk.trades().seller().accept_trade(sdk_request), + )? } - AppSdkTradeDecision::Decline { reason } => { + DesktopRuntimeTradeDecision::Decline { reason } => { let mut sdk_request = TradeDeclineRequest::new( actor, request.locator, @@ -1567,27 +1651,33 @@ fn trade_decision_with_sdk( 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))?; + .map_err(|error| DesktopRuntimeIssue::from_sdk_error(&error))?; } - runtime - .block_on(sdk.trades().seller().decline_trade(sdk_request)) - .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))? + block_on_sdk_result( + runtime, + "desktop_runtime_trade_decline", + sdk.trades().seller().decline_trade(sdk_request), + )? } }; - app_sdk_trade_decision_receipt(outcome, request.actor_pubkey) + desktop_runtime_trade_decision_receipt(outcome, request.actor_pubkey) } fn trade_cancel_with_sdk( runtime: &tokio::runtime::Runtime, - config: &AppSdkConfig, - request: AppSdkTradeCancellationRequest, -) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue> { + config: &DesktopRuntimeSupervisorConfig, + request: DesktopRuntimeTradeCancellationRequest, +) -> Result<DesktopRuntimeWorkflowReceipt, DesktopRuntimeIssue> { let actor = sdk_actor_context( request.actor_pubkey.as_str(), request.actor_account_id.as_str(), RadrootsActorRole::Buyer, )?; - let sdk = runtime.block_on(build_sdk_runtime_with_signer(config, request.signer_keys))?; + let sdk = block_on_desktop_result( + runtime, + "desktop_runtime_trade_cancel_build", + build_sdk_runtime_with_signer(config, request.signer), + )?; let mut sdk_request = TradeCancelRequest::new( actor, request.locator, @@ -1601,45 +1691,71 @@ fn trade_cancel_with_sdk( 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))?; + .map_err(|error| DesktopRuntimeIssue::from_sdk_error(&error))?; } - let outcome = runtime - .block_on(sdk.trades().buyer().cancel_trade(sdk_request)) - .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error))?; - app_sdk_trade_cancellation_receipt(outcome, request.actor_pubkey) + let outcome = block_on_sdk_result( + runtime, + "desktop_runtime_trade_cancel", + sdk.trades().buyer().cancel_trade(sdk_request), + )?; + desktop_runtime_trade_cancellation_receipt(outcome, request.actor_pubkey) } fn sdk_actor_context( actor_pubkey: &str, actor_account_id: &str, role: RadrootsActorRole, -) -> Result<RadrootsActorContext, AppSdkRuntimeIssue> { +) -> Result<RadrootsActorContext, DesktopRuntimeIssue> { RadrootsActorContext::local_account(actor_pubkey, actor_account_id.to_owned(), [role]).map_err( - |error| AppSdkRuntimeIssue::runtime_error("sdk_actor_context_invalid", error.to_string()), + |error| DesktopRuntimeIssue::runtime_error("sdk_actor_context_invalid", error.to_string()), ) } fn sdk_local_signer( - keys: RadrootsNostrKeys, -) -> Result<RadrootsLocalEventSigner, AppSdkRuntimeIssue> { - RadrootsLocalEventSigner::new(keys).map_err(|error| { - AppSdkRuntimeIssue::runtime_error("sdk_signer_init_failed", error.to_string()) + signer: DesktopRuntimeLocalSigner, +) -> Result<RadrootsLocalEventSigner, DesktopRuntimeIssue> { + RadrootsLocalEventSigner::new(signer.into_keys()).map_err(|error| { + DesktopRuntimeIssue::runtime_error("sdk_signer_init_failed", error.to_string()) }) } +fn desktop_runtime_public_farm_location_to_protocol( + location: DesktopRuntimePublicFarmLocation, +) -> RadrootsFarmPublicLocation { + RadrootsFarmPublicLocation { + primary: location.primary, + city: location.city, + region: location.region, + country: location.country, + geohash: location.geohash5, + } +} + +fn desktop_runtime_public_listing_location_to_protocol( + location: DesktopRuntimePublicFarmLocation, +) -> RadrootsListingPublicLocation { + RadrootsListingPublicLocation { + primary: location.primary, + city: location.city, + region: location.region, + country: location.country, + geohash: location.geohash5, + } +} + fn sdk_transport_targets( relays: Vec<String>, - policy: AppSdkRelayUrlPolicy, -) -> Result<TargetPolicy, AppSdkRuntimeIssue> { + policy: DesktopRuntimeRelayUrlPolicy, +) -> Result<TargetPolicy, DesktopRuntimeIssue> { TargetPolicy::try_nostr_relays(relays, policy.into()) - .map_err(|error| AppSdkRuntimeIssue::from_sdk_error(&error)) + .map_err(|error| DesktopRuntimeIssue::from_sdk_error(&error)) } -fn app_sdk_farm_receipt( +fn desktop_runtime_farm_receipt( receipt: FarmEnqueueReceipt, actor_pubkey: String, -) -> AppSdkWorkflowReceipt { - AppSdkWorkflowReceipt { +) -> DesktopRuntimeWorkflowReceipt { + DesktopRuntimeWorkflowReceipt { operation_kind: FARM_PUBLISH_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(), @@ -1651,11 +1767,11 @@ fn app_sdk_farm_receipt( } } -fn app_sdk_listing_receipt( +fn desktop_runtime_listing_receipt( receipt: ListingEnqueueReceipt, actor_pubkey: String, -) -> AppSdkWorkflowReceipt { - AppSdkWorkflowReceipt { +) -> DesktopRuntimeWorkflowReceipt { + DesktopRuntimeWorkflowReceipt { operation_kind: LISTING_PUBLISH_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(), @@ -1667,13 +1783,13 @@ fn app_sdk_listing_receipt( } } -fn app_sdk_trade_propose_receipt( +fn desktop_runtime_trade_propose_receipt( outcome: TradeMutationOutcome<TradeSubmitPlan, TradeSubmitReceipt>, actor_pubkey: String, -) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue> { +) -> Result<DesktopRuntimeWorkflowReceipt, DesktopRuntimeIssue> { match outcome { TradeMutationOutcome::Enqueued { receipt } - | TradeMutationOutcome::Published { receipt, .. } => Ok(AppSdkWorkflowReceipt { + | TradeMutationOutcome::Published { receipt, .. } => Ok(DesktopRuntimeWorkflowReceipt { 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(), @@ -1687,13 +1803,13 @@ fn app_sdk_trade_propose_receipt( } } -fn app_sdk_trade_decision_receipt( +fn desktop_runtime_trade_decision_receipt( outcome: TradeMutationOutcome<TradeDecisionPlan, TradeDecisionReceipt>, actor_pubkey: String, -) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue> { +) -> Result<DesktopRuntimeWorkflowReceipt, DesktopRuntimeIssue> { match outcome { TradeMutationOutcome::Enqueued { receipt } - | TradeMutationOutcome::Published { receipt, .. } => Ok(AppSdkWorkflowReceipt { + | TradeMutationOutcome::Published { receipt, .. } => Ok(DesktopRuntimeWorkflowReceipt { 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(), @@ -1707,13 +1823,13 @@ fn app_sdk_trade_decision_receipt( } } -fn app_sdk_trade_cancellation_receipt( +fn desktop_runtime_trade_cancellation_receipt( outcome: TradeMutationOutcome<TradeCancellationPlan, TradeCancellationReceipt>, actor_pubkey: String, -) -> Result<AppSdkWorkflowReceipt, AppSdkRuntimeIssue> { +) -> Result<DesktopRuntimeWorkflowReceipt, DesktopRuntimeIssue> { match outcome { TradeMutationOutcome::Enqueued { receipt } - | TradeMutationOutcome::Published { receipt, .. } => Ok(AppSdkWorkflowReceipt { + | TradeMutationOutcome::Published { receipt, .. } => Ok(DesktopRuntimeWorkflowReceipt { 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(), @@ -1727,8 +1843,8 @@ fn app_sdk_trade_cancellation_receipt( } } -fn unexpected_trade_dry_run_issue(operation: &'static str) -> AppSdkRuntimeIssue { - AppSdkRuntimeIssue::runtime_error( +fn unexpected_trade_dry_run_issue(operation: &'static str) -> DesktopRuntimeIssue { + DesktopRuntimeIssue::runtime_error( "sdk_trade_unexpected_dry_run", format!("{operation} returned a dry-run plan for an enqueue-only Studio command"), ) @@ -1742,106 +1858,258 @@ fn sdk_mutation_state_key(state: SdkMutationState) -> &'static str { } } -fn send_worker_result<T>( - shared: &AppSdkRuntimeShared, - response_sender: mpsc::Sender<Result<T, AppSdkRuntimeIssue>>, - result: Result<T, AppSdkRuntimeIssue>, +fn ready_startup_milestones( + config: &DesktopRuntimeSupervisorConfig, + storage_paths: Option<&DesktopRuntimeStoragePaths>, +) -> Vec<DesktopRuntimeStartupMilestone> { + let mut milestones = vec![DesktopRuntimeStartupMilestone::ShellReady]; + if storage_paths.is_some() { + milestones.push(DesktopRuntimeStartupMilestone::RuntimeStoreReady); + milestones.push(DesktopRuntimeStartupMilestone::PrivateStoreReady); + } + milestones.push(DesktopRuntimeStartupMilestone::SignerReady); + milestones.push(DesktopRuntimeStartupMilestone::ProjectionsReady); + if !config.relay_urls.is_empty() { + milestones.push(DesktopRuntimeStartupMilestone::NetworkObserved); + } + milestones +} + +fn refresh_supervisor_diagnostics( + runtime: &tokio::runtime::Runtime, + shared: &DesktopRuntimeSupervisorShared, + sdk: &RadrootsClient, +) -> Result<(), DesktopRuntimeIssue> { + let mut runtime_status = lock_status(shared).clone(); + runtime_status.last_issue = None; + let diagnostics = block_on_sdk_result( + runtime, + "desktop_runtime_refresh_diagnostics", + collect_sdk_diagnostics(sdk, runtime_status), + )?; + let mut status = lock_status(shared); + status.storage_diagnostics = Some(diagnostics.storage); + status.integrity_diagnostics = Some(diagnostics.integrity); + status.sync_diagnostics = Some(diagnostics.sync); + status.last_issue = None; + Ok(()) +} + +fn finish_effect_result( + shared: &DesktopRuntimeSupervisorShared, + receipt: DesktopRuntimeEffectReceipt, + result: Result<Option<DesktopRuntimeWorkflowReceipt>, DesktopRuntimeIssue>, +) { + match result { + Ok(workflow_receipt) => set_effect_completed(shared, receipt, workflow_receipt, None), + Err(issue) => set_effect_failed(shared, receipt, issue), + } +} + +fn finish_workflow_result( + shared: &DesktopRuntimeSupervisorShared, + receipt: DesktopRuntimeEffectReceipt, + result: Result<DesktopRuntimeWorkflowReceipt, DesktopRuntimeIssue>, ) { - set_last_issue( - shared, - match &result { - Ok(_) => None, - Err(issue) => Some(issue.clone()), - }, - ); - let _ = response_sender.send(result); + finish_effect_result(shared, receipt, result.map(Some)); } -fn lifecycle_busy_issue(shared: &AppSdkRuntimeShared) -> Option<AppSdkRuntimeIssue> { +fn finish_restore_preflight_result( + shared: &DesktopRuntimeSupervisorShared, + receipt: DesktopRuntimeEffectReceipt, + result: Result<DesktopRuntimeRestorePreflightReceipt, DesktopRuntimeIssue>, +) { + match result { + Ok(restore_preflight) => { + set_effect_completed(shared, receipt, None, Some(restore_preflight)); + } + Err(issue) => set_effect_failed(shared, receipt, issue), + } +} + +fn finish_projection_result( + shared: &DesktopRuntimeSupervisorShared, + receipt: DesktopRuntimeEffectReceipt, + result: Result<DesktopRuntimeProjectionLifecycleStatus, DesktopRuntimeIssue>, +) { + match result { + Ok(_) => set_effect_completed(shared, receipt, None, None), + Err(issue) => set_effect_failed(shared, receipt, issue), + } +} + +fn set_effect_completed( + shared: &DesktopRuntimeSupervisorShared, + receipt: DesktopRuntimeEffectReceipt, + workflow_receipt: Option<DesktopRuntimeWorkflowReceipt>, + restore_preflight: Option<DesktopRuntimeRestorePreflightReceipt>, +) { + let mut status = lock_status(shared); + status.last_issue = None; + status.last_effect = Some(DesktopRuntimeEffectStatus::completed( + receipt, + workflow_receipt, + restore_preflight, + )); +} + +fn set_effect_failed( + shared: &DesktopRuntimeSupervisorShared, + receipt: DesktopRuntimeEffectReceipt, + issue: DesktopRuntimeIssue, +) { + let mut status = lock_status(shared); + status.last_issue = Some(issue.clone()); + status.last_effect = Some(DesktopRuntimeEffectStatus::failed(receipt, issue)); +} + +fn set_last_effect(shared: &DesktopRuntimeSupervisorShared, status: DesktopRuntimeEffectStatus) { + lock_status(shared).last_effect = Some(status); +} + +fn clear_last_effect(shared: &DesktopRuntimeSupervisorShared, effect_id: u64) { + let mut status = lock_status(shared); + if status + .last_effect + .as_ref() + .is_some_and(|effect| effect.receipt.effect_id == effect_id) + { + status.last_effect = None; + } +} + +fn set_startup_milestone( + status: &mut DesktopRuntimeSnapshot, + milestone: DesktopRuntimeStartupMilestone, + enabled: bool, +) { + if enabled { + if !status.startup_milestones.contains(&milestone) { + status.startup_milestones.push(milestone); + status.startup_milestones.sort(); + } + } else { + status + .startup_milestones + .retain(|candidate| *candidate != milestone); + } +} + +fn effect_receipt(effect: &DesktopRuntimeEffect) -> DesktopRuntimeEffectReceipt { + match effect { + DesktopRuntimeEffect::RefreshDiagnostics(receipt) + | DesktopRuntimeEffect::RestorePreflight(receipt, _) + | DesktopRuntimeEffect::EnqueueFarmPublish(receipt, _) + | DesktopRuntimeEffect::EnqueueListingPublish(receipt, _) + | DesktopRuntimeEffect::TradePropose(receipt, _) + | DesktopRuntimeEffect::TradeDecision(receipt, _) + | DesktopRuntimeEffect::TradeCancellation(receipt, _) + | DesktopRuntimeEffect::BeginProjectionRebuild(receipt) + | DesktopRuntimeEffect::CompleteProjectionRebuild(receipt) => receipt.clone(), + } +} + +fn lifecycle_busy_issue(shared: &DesktopRuntimeSupervisorShared) -> Option<DesktopRuntimeIssue> { let state = lock_status(shared).state; if matches!( state, - AppSdkLifecycleState::Pausing - | AppSdkLifecycleState::Paused - | AppSdkLifecycleState::Restoring - | AppSdkLifecycleState::RebuildingProjections - | AppSdkLifecycleState::ShuttingDown + DesktopRuntimeLifecycleState::Pausing + | DesktopRuntimeLifecycleState::Paused + | DesktopRuntimeLifecycleState::Restoring + | DesktopRuntimeLifecycleState::RebuildingProjections + | DesktopRuntimeLifecycleState::ShuttingDown ) { - Some(AppSdkRuntimeIssue::lifecycle_blocked(state)) + Some(DesktopRuntimeIssue::lifecycle_blocked(state)) } else { None } } -fn runtime_unavailable_issue(shared: &AppSdkRuntimeShared) -> AppSdkRuntimeIssue { +fn runtime_unavailable_issue(shared: &DesktopRuntimeSupervisorShared) -> DesktopRuntimeIssue { let status = lock_status(shared).clone(); if let Some(issue) = status.last_issue { issue } else { - AppSdkRuntimeIssue::runtime_error( + DesktopRuntimeIssue::runtime_error( "sdk_runtime_not_ready", format!("app sdk runtime is {:?}", status.state), ) } } -fn replace_status(shared: &AppSdkRuntimeShared, status: AppSdkRuntimeStatus) { +fn replace_status(shared: &DesktopRuntimeSupervisorShared, status: DesktopRuntimeSnapshot) { *lock_status(shared) = status; - shared.status_changed.notify_all(); } -fn set_last_issue(shared: &AppSdkRuntimeShared, issue: Option<AppSdkRuntimeIssue>) { - lock_status(shared).last_issue = issue; - shared.status_changed.notify_all(); -} - -fn transition_status_state(shared: &AppSdkRuntimeShared, state: AppSdkLifecycleState) { +fn transition_status_state( + shared: &DesktopRuntimeSupervisorShared, + state: DesktopRuntimeLifecycleState, +) { lock_status(shared).state = state; - shared.status_changed.notify_all(); } fn mark_projections_stale( - shared: &AppSdkRuntimeShared, + shared: &DesktopRuntimeSupervisorShared, reason: impl Into<String>, restore_source: Option<PathBuf>, -) -> AppSdkProjectionLifecycleStatus { +) -> DesktopRuntimeProjectionLifecycleStatus { let mut status = lock_status(shared); - status.projection_lifecycle = AppSdkProjectionLifecycleStatus::stale(reason, restore_source); - status.state = AppSdkLifecycleState::Ready; + status.projection_lifecycle = + DesktopRuntimeProjectionLifecycleStatus::stale(reason, restore_source); + status.state = DesktopRuntimeLifecycleState::Ready; + set_startup_milestone( + &mut status, + DesktopRuntimeStartupMilestone::ProjectionsReady, + false, + ); let projection_lifecycle = status.projection_lifecycle.clone(); - shared.status_changed.notify_all(); projection_lifecycle } -fn begin_projection_rebuild(shared: &AppSdkRuntimeShared) -> AppSdkProjectionLifecycleStatus { +fn begin_projection_rebuild( + shared: &DesktopRuntimeSupervisorShared, +) -> DesktopRuntimeProjectionLifecycleStatus { let restore_source = lock_status(shared) .projection_lifecycle .restore_source .clone(); let mut status = lock_status(shared); - status.state = AppSdkLifecycleState::RebuildingProjections; - status.projection_lifecycle = - AppSdkProjectionLifecycleStatus::rebuilding("sdk_projection_rebuild", restore_source); + status.state = DesktopRuntimeLifecycleState::RebuildingProjections; + status.projection_lifecycle = DesktopRuntimeProjectionLifecycleStatus::rebuilding( + "sdk_projection_rebuild", + restore_source, + ); + set_startup_milestone( + &mut status, + DesktopRuntimeStartupMilestone::ProjectionsReady, + false, + ); let projection_lifecycle = status.projection_lifecycle.clone(); - shared.status_changed.notify_all(); projection_lifecycle } fn complete_projection_rebuild( - shared: &AppSdkRuntimeShared, -) -> Result<AppSdkProjectionLifecycleStatus, AppSdkRuntimeIssue> { + shared: &DesktopRuntimeSupervisorShared, +) -> Result<DesktopRuntimeProjectionLifecycleStatus, DesktopRuntimeIssue> { let mut status = lock_status(shared); - if !matches!(status.state, AppSdkLifecycleState::RebuildingProjections) { - return Err(AppSdkRuntimeIssue::lifecycle_blocked(status.state)); - } - status.state = AppSdkLifecycleState::Ready; - status.projection_lifecycle = AppSdkProjectionLifecycleStatus::current(); + if !matches!( + status.state, + DesktopRuntimeLifecycleState::RebuildingProjections + ) { + return Err(DesktopRuntimeIssue::lifecycle_blocked(status.state)); + } + status.state = DesktopRuntimeLifecycleState::Ready; + status.projection_lifecycle = DesktopRuntimeProjectionLifecycleStatus::current(); + set_startup_milestone( + &mut status, + DesktopRuntimeStartupMilestone::ProjectionsReady, + true, + ); let projection_lifecycle = status.projection_lifecycle.clone(); - shared.status_changed.notify_all(); Ok(projection_lifecycle) } -fn lock_status(shared: &AppSdkRuntimeShared) -> MutexGuard<'_, AppSdkRuntimeStatus> { +fn lock_status(shared: &DesktopRuntimeSupervisorShared) -> MutexGuard<'_, DesktopRuntimeSnapshot> { shared .status .lock() @@ -1865,13 +2133,7 @@ fn serialized_label(value: &(impl Serialize + fmt::Debug)) -> String { #[cfg(test)] mod tests { use std::{ - fs, - sync::{ - Arc, Condvar, Mutex, - atomic::{AtomicBool, Ordering}, - mpsc, - }, - thread, + fs, thread, time::{Duration, SystemTime, UNIX_EPOCH}, }; @@ -1900,10 +2162,12 @@ mod tests { }; use super::{ - APP_SDK_STORAGE_DIR_NAME, AppSdkConfig, AppSdkLifecycleState, AppSdkListingPublishRequest, - AppSdkProjectionLifecycleState, AppSdkRelayUrlPolicy, AppSdkRestorePreflightRequest, - AppSdkRuntime, AppSdkRuntimeError, AppSdkRuntimeShared, AppSdkRuntimeStatus, - AppSdkWorkerCommand, app_sdk_storage_root_from_data_root, transition_status_state, + DESKTOP_RUNTIME_STORAGE_DIR_NAME, DesktopRuntimeEffectKind, DesktopRuntimeEffectState, + DesktopRuntimeLifecycleState, DesktopRuntimeListingPublishRequest, + DesktopRuntimeLocalSigner, DesktopRuntimeProjectionLifecycleState, + DesktopRuntimeRelayUrlPolicy, DesktopRuntimeRestorePreflightRequest, + DesktopRuntimeSnapshot, DesktopRuntimeStartupMilestone, DesktopRuntimeSupervisor, + DesktopRuntimeSupervisorConfig, desktop_runtime_storage_root_from_data_root, }; const SDK_TEST_SELLER_SECRET_KEY_HEX: &str = @@ -1919,25 +2183,30 @@ mod tests { }, ) .expect("desktop paths should resolve"); - let config = - AppSdkConfig::from_desktop_paths(&paths, vec!["wss://relay.example".to_owned()]); + let config = DesktopRuntimeSupervisorConfig::from_desktop_paths( + &paths, + vec!["wss://relay.example".to_owned()], + ); assert_eq!( config.storage_root, - paths.app.data.join(APP_SDK_STORAGE_DIR_NAME) + paths.app.data.join(DESKTOP_RUNTIME_STORAGE_DIR_NAME) ); assert_eq!( config.storage_root, - app_sdk_storage_root_from_data_root(paths.app.data.as_path()) + desktop_runtime_storage_root_from_data_root(paths.app.data.as_path()) ); assert_eq!(config.storage_root.parent(), Some(paths.app.data.as_path())); assert!(paths.app.data.ends_with(APP_RUNTIME_NAMESPACE)); - assert_eq!(config.relay_url_policy, AppSdkRelayUrlPolicy::Public); + assert_eq!( + config.relay_url_policy, + DesktopRuntimeRelayUrlPolicy::Public + ); } #[test] fn sdk_config_uses_localhost_policy_for_ws_relay_urls() { - let config = AppSdkConfig::from_app_data_root( + let config = DesktopRuntimeSupervisorConfig::from_app_data_root( "/tmp/radroots-app-data".as_ref(), vec![ "wss://relay.example".to_owned(), @@ -1945,27 +2214,69 @@ mod tests { ], ); - assert_eq!(config.relay_url_policy, AppSdkRelayUrlPolicy::Localhost); + assert_eq!( + config.relay_url_policy, + DesktopRuntimeRelayUrlPolicy::Localhost + ); } #[test] - fn sdk_runtime_reaches_ready_with_directory_storage() { + fn desktop_runtime_supervisor_reaches_ready_with_snapshot_diagnostics() { let storage_root = temp_storage_root("ready"); - let config = AppSdkConfig::from_app_data_root( + let config = DesktopRuntimeSupervisorConfig::from_app_data_root( storage_root .parent() .expect("storage root should have parent"), vec!["ws://127.0.0.1:8080".to_owned()], ); - let runtime = AppSdkRuntime::start(config).expect("sdk runtime should start"); + let runtime = DesktopRuntimeSupervisor::start(config).expect("supervisor should start"); - let status = runtime.wait_for_startup(Duration::from_secs(5)); + let status = poll_snapshot(&runtime, |snapshot| { + matches!(snapshot.state, DesktopRuntimeLifecycleState::Ready) + && snapshot.storage_diagnostics.is_some() + && snapshot.integrity_diagnostics.is_some() + && snapshot.sync_diagnostics.is_some() + }); - assert_eq!(status.state, AppSdkLifecycleState::Ready); + assert_eq!(status.state, DesktopRuntimeLifecycleState::Ready); assert_eq!(status.storage_root, storage_root); - assert_eq!(status.relay_url_policy, AppSdkRelayUrlPolicy::Localhost); + assert_eq!( + status.relay_url_policy, + DesktopRuntimeRelayUrlPolicy::Localhost + ); + assert!( + status + .startup_milestones + .contains(&DesktopRuntimeStartupMilestone::ShellReady) + ); + assert!( + status + .startup_milestones + .contains(&DesktopRuntimeStartupMilestone::RuntimeStoreReady) + ); + assert!( + status + .startup_milestones + .contains(&DesktopRuntimeStartupMilestone::PrivateStoreReady) + ); + assert!( + status + .startup_milestones + .contains(&DesktopRuntimeStartupMilestone::SignerReady) + ); + assert!( + status + .startup_milestones + .contains(&DesktopRuntimeStartupMilestone::ProjectionsReady) + ); + assert!( + status + .startup_milestones + .contains(&DesktopRuntimeStartupMilestone::NetworkObserved) + ); let storage_paths = status .storage_paths + .as_ref() .expect("storage paths should be present"); assert_eq!( storage_paths.runtime_path, @@ -1979,88 +2290,124 @@ mod tests { storage_paths.studio_path, storage_root.join("studio.sqlite") ); - let storage = runtime - .storage_status() + let storage = status + .storage_diagnostics + .as_ref() .expect("storage diagnostics should load"); assert_eq!(storage.storage_kind, "directory"); assert!(storage.event_store.store.integrity_ok); assert!(storage.outbox.store.integrity_ok); - let integrity = runtime - .integrity_status() + let integrity = status + .integrity_diagnostics + .as_ref() .expect("integrity diagnostics should load"); assert!(integrity.event_store_ok); assert!(integrity.outbox_ok); - let sync = runtime.sync_status().expect("sync diagnostics should load"); + let sync = status + .sync_diagnostics + .as_ref() + .expect("sync diagnostics should load"); assert_eq!(sync.source, "sdk_canonical_stores"); assert_eq!(sync.transport_targets.configured_count, 1); - let diagnostics = runtime.diagnostics().expect("diagnostics should load"); - assert_eq!(diagnostics.runtime.state, AppSdkLifecycleState::Ready); - assert_eq!(diagnostics.storage.storage_kind, "directory"); - assert_eq!(diagnostics.sync.transport_targets.configured_count, 1); - runtime.shutdown().expect("sdk runtime should shut down"); - assert_eq!(runtime.status().state, AppSdkLifecycleState::Stopped); + assert!(runtime.request_shutdown()); + let stopped = poll_snapshot(&runtime, |snapshot| { + matches!(snapshot.state, DesktopRuntimeLifecycleState::Stopped) + }); + assert_eq!(stopped.state, DesktopRuntimeLifecycleState::Stopped); let _ = fs::remove_dir_all(storage_root); } #[test] - fn sdk_runtime_enqueues_listing_publish_work() { + fn desktop_runtime_supervisor_enqueues_listing_publish_as_effect() { let storage_root = temp_storage_root("listing_enqueue"); - let config = AppSdkConfig::from_app_data_root( + let config = DesktopRuntimeSupervisorConfig::from_app_data_root( storage_root .parent() .expect("storage root should have parent"), vec!["ws://127.0.0.1:8080".to_owned()], ); - let runtime = AppSdkRuntime::start(config).expect("sdk runtime should start"); - assert_eq!( - runtime.wait_for_startup(Duration::from_secs(5)).state, - AppSdkLifecycleState::Ready - ); + let runtime = DesktopRuntimeSupervisor::start(config).expect("supervisor should start"); + poll_snapshot(&runtime, |snapshot| { + matches!(snapshot.state, DesktopRuntimeLifecycleState::Ready) + }); let secret_key = RadrootsNostrSecretKey::from_hex(SDK_TEST_SELLER_SECRET_KEY_HEX) .expect("secret key should parse"); - let signer_keys = RadrootsNostrKeys::new(secret_key); - let seller_pubkey = signer_keys.public_key().to_hex(); + let local_identity_keys = RadrootsNostrKeys::new(secret_key); + let seller_pubkey = local_identity_keys.public_key().to_hex(); let receipt = runtime - .enqueue_listing_publish(AppSdkListingPublishRequest { + .enqueue_listing_publish(DesktopRuntimeListingPublishRequest { actor_account_id: "seller-account".to_owned(), actor_pubkey: seller_pubkey.clone(), - signer_keys, + signer: DesktopRuntimeLocalSigner::from_local_identity_keys(local_identity_keys), listing: test_listing(seller_pubkey.as_str()), - target_relays: vec!["ws://127.0.0.1:8080".to_owned()], - relay_url_policy: AppSdkRelayUrlPolicy::Localhost, - idempotency_key: Some("listing-enqueue-idempotency".to_owned()), + target_relays: vec!["wss://relay.radroots.test".to_owned()], + relay_url_policy: DesktopRuntimeRelayUrlPolicy::Public, + idempotency_key: Some("01890f0e-6c00-7000-8000-000000000224".to_owned()), }) - .expect("listing publish should enqueue"); - - assert_eq!(receipt.operation_kind, LISTING_PUBLISH_OPERATION_KIND); - assert_eq!(receipt.actor_pubkey, seller_pubkey); - assert_eq!(receipt.state, "enqueued"); - assert!(!receipt.expected_event_id.is_empty()); - assert_eq!(receipt.expected_event_id, receipt.signed_event_id); - assert!(receipt.outbox_operation_id > 0); - assert!(receipt.outbox_event_id > 0); - assert!(receipt.idempotency_digest_prefix.is_some()); - let sync = runtime.sync_status().expect("sync diagnostics should load"); - assert_eq!(sync.outbox.ready_signed_events, 1); - runtime.shutdown().expect("sdk runtime should shut down"); + .expect("listing publish effect should enqueue"); + + assert_eq!( + receipt.effect_kind, + DesktopRuntimeEffectKind::ListingPublish + ); + assert_eq!( + receipt.operation_kind.as_deref(), + Some(LISTING_PUBLISH_OPERATION_KIND) + ); + assert_eq!( + receipt.actor_pubkey.as_deref(), + Some(seller_pubkey.as_str()) + ); + assert!(receipt.effect_id > 0); + + let completed = poll_snapshot(&runtime, |snapshot| { + snapshot.last_effect.as_ref().is_some_and(|effect| { + effect.receipt.effect_id == receipt.effect_id + && matches!(effect.state, DesktopRuntimeEffectState::Completed) + && effect.workflow_receipt.is_some() + }) + }); + let completed_effect = completed + .last_effect + .as_ref() + .expect("listing effect should be captured in snapshot"); + assert_eq!(completed_effect.receipt.effect_id, receipt.effect_id); + assert!(matches!( + completed_effect.state, + DesktopRuntimeEffectState::Completed + )); + let workflow = completed_effect + .workflow_receipt + .as_ref() + .expect("workflow receipt should be captured in snapshot"); + assert_eq!(workflow.operation_kind, LISTING_PUBLISH_OPERATION_KIND); + assert_eq!(workflow.actor_pubkey, seller_pubkey); + assert_eq!(workflow.state, "enqueued"); + assert!(!workflow.expected_event_id.is_empty()); + assert_eq!(workflow.expected_event_id, workflow.signed_event_id); + assert!(workflow.outbox_operation_id > 0); + assert!(workflow.outbox_event_id > 0); + assert!(workflow.idempotency_digest_prefix.is_some()); + assert!(runtime.request_shutdown()); let _ = fs::remove_dir_all(storage_root); } #[test] - fn sdk_runtime_degrades_with_structured_sdk_error() { + fn desktop_runtime_supervisor_degrades_with_structured_sdk_error() { let storage_root = temp_storage_root("invalid_relay"); - let config = AppSdkConfig::from_app_data_root( + let config = DesktopRuntimeSupervisorConfig::from_app_data_root( storage_root .parent() .expect("storage root should have parent"), vec!["ws://relay.example".to_owned()], ); - let runtime = AppSdkRuntime::start(config).expect("sdk runtime should start"); + let runtime = DesktopRuntimeSupervisor::start(config).expect("supervisor should start"); - let status = runtime.wait_for_startup(Duration::from_secs(5)); + let status = poll_snapshot(&runtime, |snapshot| { + matches!(snapshot.state, DesktopRuntimeLifecycleState::Degraded) + }); - assert_eq!(status.state, AppSdkLifecycleState::Degraded); let issue = status .last_issue .expect("degraded status should include issue"); @@ -2073,77 +2420,27 @@ mod tests { .contains(&"configure_transport_targets".to_owned()) ); assert_eq!(issue.detail_json["code"], "invalid_relay_url"); - let error = runtime - .diagnostics() - .expect_err("degraded diagnostics should fail"); - match error { - AppSdkRuntimeError::CommandFailed(issue) => { - assert_eq!(issue.code, "invalid_relay_url"); - assert_eq!(issue.class, "configuration"); - assert_eq!(issue.detail_json["code"], "invalid_relay_url"); - } - unexpected => panic!("unexpected degraded diagnostics error: {unexpected:?}"), - } - runtime.shutdown().expect("sdk runtime should shut down"); - let _ = fs::remove_dir_all(storage_root); - } - - #[test] - fn sdk_shutdown_joins_when_normal_command_queue_is_full() { - let config = AppSdkConfig::from_app_data_root( - "/tmp/radroots-app-sdk-full-queue".as_ref(), - vec!["ws://127.0.0.1:8080".to_owned()], - ) - .with_command_queue_capacity(1); - let shared = Arc::new(AppSdkRuntimeShared { - status: Mutex::new(AppSdkRuntimeStatus::from_config( - &config, - AppSdkLifecycleState::Ready, - None, - None, - )), - status_changed: Condvar::new(), - shutdown_requested: AtomicBool::new(false), - }); - let (command_sender, command_receiver) = mpsc::sync_channel(config.command_queue_capacity); - let worker_shared = Arc::clone(&shared); - let worker = thread::spawn(move || { - while !worker_shared.shutdown_requested.load(Ordering::SeqCst) { - thread::sleep(Duration::from_millis(1)); - } - drop(command_receiver); - transition_status_state(&worker_shared, AppSdkLifecycleState::Stopped); + let refresh = runtime + .request_diagnostics_refresh() + .expect("diagnostics refresh effect should be accepted"); + let failed = poll_snapshot(&runtime, |snapshot| { + snapshot.last_effect.as_ref().is_some_and(|effect| { + effect.receipt.effect_id == refresh.effect_id + && matches!(effect.state, DesktopRuntimeEffectState::Failed) + }) }); - let runtime = AppSdkRuntime { - command_sender: Mutex::new(Some(command_sender)), - shared, - worker: Mutex::new(Some(worker)), - }; - let (response_sender, _response_receiver) = mpsc::channel(); - runtime - .command_sender - .lock() - .expect("command sender lock") + let issue = failed + .last_effect .as_ref() - .expect("command sender") - .try_send(AppSdkWorkerCommand::Diagnostics(response_sender)) - .expect("normal command queue should fill"); - - assert!(matches!( - runtime.sync_status(), - Err(AppSdkRuntimeError::CommandQueueFull) - )); - assert_eq!(runtime.status().state, AppSdkLifecycleState::Ready); - - runtime - .shutdown() - .expect("shutdown should not depend on normal command queue capacity"); - - assert_eq!(runtime.status().state, AppSdkLifecycleState::Stopped); + .and_then(|effect| effect.issue.as_ref()) + .expect("failed diagnostics effect should include issue"); + assert_eq!(issue.code, "invalid_relay_url"); + assert!(runtime.request_shutdown()); + let _ = fs::remove_dir_all(storage_root); } #[test] - fn sdk_restore_preflight_marks_projections_stale_without_writing_destination() { + fn desktop_runtime_supervisor_restore_preflight_marks_projections_stale() { let backup_source_root = temp_storage_root("restore_backup_source"); let backup_archive = backup_source_root .parent() @@ -2173,43 +2470,53 @@ mod tests { .parent() .expect("app storage root should have parent") .to_path_buf(); - let config = AppSdkConfig::from_app_data_root( + let config = DesktopRuntimeSupervisorConfig::from_app_data_root( app_data_root.as_path(), vec!["ws://127.0.0.1:8080".to_owned()], ); - let runtime = AppSdkRuntime::start(config).expect("sdk runtime should start"); - assert_eq!( - runtime.wait_for_startup(Duration::from_secs(5)).state, - AppSdkLifecycleState::Ready - ); + let runtime = DesktopRuntimeSupervisor::start(config).expect("supervisor should start"); + poll_snapshot(&runtime, |snapshot| { + matches!(snapshot.state, DesktopRuntimeLifecycleState::Ready) + }); let sentinel = app_storage_root.join("restore-preflight-sentinel"); fs::write(&sentinel, "existing destination").expect("sentinel should write"); - let receipt = runtime - .restore_preflight( - AppSdkRestorePreflightRequest::new(backup_archive.clone()) + let effect = runtime + .request_restore_preflight( + DesktopRuntimeRestorePreflightRequest::new(backup_archive.clone()) .with_overwrite_existing_sdk_storage(true), ) - .expect("restore preflight should succeed"); + .expect("restore preflight effect should be accepted"); + let completed = poll_snapshot(&runtime, |snapshot| { + snapshot.last_effect.as_ref().is_some_and(|last_effect| { + last_effect.receipt.effect_id == effect.effect_id + && matches!(last_effect.state, DesktopRuntimeEffectState::Completed) + }) + }); + let receipt = completed + .last_effect + .as_ref() + .and_then(|effect| effect.restore_preflight.as_ref()) + .expect("restore preflight receipt should be captured"); assert_eq!(receipt.state, "dry_run"); assert_eq!(receipt.destination, app_storage_root); assert_eq!(receipt.restored_paths, None); assert!(sentinel.exists()); assert_eq!( receipt.projection_lifecycle.state, - AppSdkProjectionLifecycleState::Stale + DesktopRuntimeProjectionLifecycleState::Stale ); assert_eq!( - receipt.projection_lifecycle.reason.as_deref(), - Some("sdk_restore_preflight") + completed.projection_lifecycle.state, + DesktopRuntimeProjectionLifecycleState::Stale ); - assert_eq!( - runtime.status().projection_lifecycle.state, - AppSdkProjectionLifecycleState::Stale + assert!( + !completed + .startup_milestones + .contains(&DesktopRuntimeStartupMilestone::ProjectionsReady) ); - assert_eq!(runtime.status().state, AppSdkLifecycleState::Ready); - runtime.shutdown().expect("sdk runtime should shut down"); + assert!(runtime.request_shutdown()); let _ = fs::remove_dir_all( backup_source_root .parent() @@ -2219,61 +2526,103 @@ mod tests { } #[test] - fn sdk_projection_rebuild_state_rejects_conflicting_commands() { + fn desktop_runtime_supervisor_projection_rebuild_uses_effect_snapshots() { let storage_root = temp_storage_root("projection_rebuild"); - let config = AppSdkConfig::from_app_data_root( + let config = DesktopRuntimeSupervisorConfig::from_app_data_root( storage_root .parent() .expect("storage root should have parent"), vec!["ws://127.0.0.1:8080".to_owned()], ); - let runtime = AppSdkRuntime::start(config).expect("sdk runtime should start"); - assert_eq!( - runtime.wait_for_startup(Duration::from_secs(5)).state, - AppSdkLifecycleState::Ready - ); + let runtime = DesktopRuntimeSupervisor::start(config).expect("supervisor should start"); + poll_snapshot(&runtime, |snapshot| { + matches!(snapshot.state, DesktopRuntimeLifecycleState::Ready) + }); - let rebuilding = runtime + let begin = runtime .begin_projection_rebuild() - .expect("projection rebuild should start"); - - assert_eq!(rebuilding.state, AppSdkProjectionLifecycleState::Rebuilding); + .expect("projection rebuild begin effect should enqueue"); + let rebuilding = poll_snapshot(&runtime, |snapshot| { + snapshot + .last_effect + .as_ref() + .is_some_and(|effect| effect.receipt.effect_id == begin.effect_id) + && matches!( + snapshot.state, + DesktopRuntimeLifecycleState::RebuildingProjections + ) + }); assert_eq!( - runtime.status().state, - AppSdkLifecycleState::RebuildingProjections + rebuilding.projection_lifecycle.state, + DesktopRuntimeProjectionLifecycleState::Rebuilding ); - let error = runtime - .sync_status() - .expect_err("sync status should wait for rebuild completion"); - match error { - AppSdkRuntimeError::CommandFailed(issue) => { - assert_eq!(issue.code, "sdk_lifecycle_busy"); - assert_eq!(issue.detail_json["state"], "RebuildingProjections"); - } - unexpected => panic!("unexpected lifecycle error: {unexpected:?}"), - } + assert!( + !rebuilding + .startup_milestones + .contains(&DesktopRuntimeStartupMilestone::ProjectionsReady) + ); + + let refresh = runtime + .request_diagnostics_refresh() + .expect("diagnostics refresh should be accepted while rebuilding"); + let blocked = poll_snapshot(&runtime, |snapshot| { + snapshot.last_effect.as_ref().is_some_and(|effect| { + effect.receipt.effect_id == refresh.effect_id + && matches!(effect.state, DesktopRuntimeEffectState::Failed) + }) + }); + let issue = blocked + .last_effect + .as_ref() + .and_then(|effect| effect.issue.as_ref()) + .expect("blocked refresh should expose lifecycle issue"); + assert_eq!(issue.code, "sdk_lifecycle_busy"); + assert_eq!(issue.detail_json["state"], "RebuildingProjections"); let complete = runtime .complete_projection_rebuild() - .expect("projection rebuild should complete"); - - assert_eq!(complete.state, AppSdkProjectionLifecycleState::Current); - assert_eq!(runtime.status().state, AppSdkLifecycleState::Ready); - runtime - .sync_status() - .expect("sync status should work after rebuild"); - runtime.shutdown().expect("sdk runtime should shut down"); + .expect("projection rebuild complete effect should enqueue"); + let current = poll_snapshot(&runtime, |snapshot| { + snapshot.last_effect.as_ref().is_some_and(|effect| { + effect.receipt.effect_id == complete.effect_id + && matches!(effect.state, DesktopRuntimeEffectState::Completed) + }) && matches!(snapshot.state, DesktopRuntimeLifecycleState::Ready) + }); + assert_eq!( + current.projection_lifecycle.state, + DesktopRuntimeProjectionLifecycleState::Current + ); + assert!( + current + .startup_milestones + .contains(&DesktopRuntimeStartupMilestone::ProjectionsReady) + ); + assert!(runtime.request_shutdown()); let _ = fs::remove_dir_all(storage_root); } + fn poll_snapshot( + runtime: &DesktopRuntimeSupervisor, + predicate: impl Fn(&DesktopRuntimeSnapshot) -> bool, + ) -> DesktopRuntimeSnapshot { + for _ in 0..500 { + let snapshot = runtime.snapshot(); + if predicate(&snapshot) { + return snapshot; + } + thread::sleep(Duration::from_millis(10)); + } + runtime.snapshot() + } + fn temp_storage_root(label: &str) -> std::path::PathBuf { let nanos = SystemTime::now() .duration_since(UNIX_EPOCH) .expect("clock") .as_nanos(); std::env::temp_dir() - .join(format!("radroots_studio_app_sdk_runtime_{label}_{nanos}")) - .join(APP_SDK_STORAGE_DIR_NAME) + .join(format!("radroots_studio_desktop_runtime_{label}_{nanos}")) + .join(DESKTOP_RUNTIME_STORAGE_DIR_NAME) } fn test_listing(seller_pubkey: &str) -> RadrootsListing { diff --git a/crates/store/migrations/0028_runtime_workflow_receipts.sql b/crates/store/migrations/0028_runtime_workflow_receipts.sql @@ -0,0 +1,38 @@ +CREATE TABLE desktop_runtime_workflow_receipts ( + id TEXT PRIMARY KEY NOT NULL, + source_kind TEXT NOT NULL CHECK ( + source_kind IN ('app_workflow', 'shared_runtime_store') + ), + source_record_id TEXT NOT NULL, + sdk_operation_kind TEXT NOT NULL, + runtime_effect_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_desktop_runtime_workflow_receipts_source_record ON desktop_runtime_workflow_receipts( + source_record_id +); +CREATE INDEX idx_desktop_runtime_workflow_receipts_state ON desktop_runtime_workflow_receipts( + workflow_state, + updated_at +); diff --git a/crates/store/migrations/0028_sdk_workflow_receipts.sql b/crates/store/migrations/0028_sdk_workflow_receipts.sql @@ -1,38 +0,0 @@ -CREATE TABLE app_sdk_workflow_receipts ( - id TEXT PRIMARY KEY NOT NULL, - source_kind TEXT NOT NULL CHECK ( - source_kind IN ('local_outbox', 'shared_runtime_store') - ), - 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 @@ -5,7 +5,7 @@ mod error; mod interop; mod migrations; mod repo; -mod sdk_workflow_receipts; +mod runtime_workflow_receipts; #[cfg(test)] mod source_guards; mod sync; @@ -47,9 +47,10 @@ pub use repo::{ SellerOrderDecisionExport, SellerOrderDecisionLineExport, TODAY_AGENDA_LIST_LIMIT, TODAY_AGENDA_LOW_STOCK_THRESHOLD, derive_farm_rules_readiness, }; -pub use sdk_workflow_receipts::{ - AppSdkStoredWorkflowReceipt, AppSdkWorkflowReceiptInput, AppSdkWorkflowReceiptRepository, - AppSdkWorkflowReceiptSourceKind, AppSdkWorkflowReceiptState, +pub use runtime_workflow_receipts::{ + DesktopRuntimeStoredWorkflowReceipt, DesktopRuntimeWorkflowReceiptInput, + DesktopRuntimeWorkflowReceiptRepository, DesktopRuntimeWorkflowReceiptSourceKind, + DesktopRuntimeWorkflowReceiptState, }; pub use sync::{ AppRelayIngestFailureInput, AppRelayIngestSuccessInput, AppSyncRepository, @@ -162,8 +163,10 @@ impl AppSqliteStore { AppSyncRepository::new(&self.connection) } - pub fn sdk_workflow_receipt_repository(&self) -> AppSdkWorkflowReceiptRepository<'_> { - AppSdkWorkflowReceiptRepository::new(&self.connection) + pub fn runtime_workflow_receipt_repository( + &self, + ) -> DesktopRuntimeWorkflowReceiptRepository<'_> { + DesktopRuntimeWorkflowReceiptRepository::new(&self.connection) } pub fn reminders_repository(&self) -> AppRemindersRepository<'_> { @@ -852,7 +855,10 @@ 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_workflow_receipts")); + assert!(table_exists( + connection, + "desktop_runtime_workflow_receipts" + )); assert!(column_exists(connection, "farms", "timezone")); assert!(column_exists(connection, "farms", "currency_code")); assert!(column_exists(connection, "local_outbox", "account_id")); @@ -1030,47 +1036,47 @@ mod tests { )); assert!(column_exists( connection, - "app_sdk_workflow_receipts", + "desktop_runtime_workflow_receipts", "source_record_id" )); assert!(column_exists( connection, - "app_sdk_workflow_receipts", + "desktop_runtime_workflow_receipts", "source_kind" )); assert!(column_exists( connection, - "app_sdk_workflow_receipts", + "desktop_runtime_workflow_receipts", "sdk_operation_kind" )); assert!(column_exists( connection, - "app_sdk_workflow_receipts", - "sdk_outbox_event_ids_json" + "desktop_runtime_workflow_receipts", + "runtime_effect_ids_json" )); assert!(column_exists( connection, - "app_sdk_workflow_receipts", + "desktop_runtime_workflow_receipts", "expected_event_id" )); assert!(column_exists( connection, - "app_sdk_workflow_receipts", + "desktop_runtime_workflow_receipts", "actor_pubkey" )); assert!(column_exists( connection, - "app_sdk_workflow_receipts", + "desktop_runtime_workflow_receipts", "idempotency_digest_prefix" )); assert!(column_exists( connection, - "app_sdk_workflow_receipts", + "desktop_runtime_workflow_receipts", "workflow_state" )); assert!(column_exists( connection, - "app_sdk_workflow_receipts", + "desktop_runtime_workflow_receipts", "detail_json" )); connection 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_workflow_receipts.sql"), + sql: include_str!("../migrations/0028_runtime_workflow_receipts.sql"), }, Migration { version: 29, diff --git a/crates/store/src/runtime_workflow_receipts.rs b/crates/store/src/runtime_workflow_receipts.rs @@ -0,0 +1,323 @@ +use sqlx::Row; + +use crate::{AppSqliteDatabase, OptionalSqliteResult}; +use serde_json::Value; +use uuid::Uuid; + +use crate::AppSqliteError; + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum DesktopRuntimeWorkflowReceiptSourceKind { + AppWorkflow, + SharedRuntimeStore, +} + +impl DesktopRuntimeWorkflowReceiptSourceKind { + pub const fn storage_key(self) -> &'static str { + match self { + Self::AppWorkflow => "app_workflow", + Self::SharedRuntimeStore => "shared_runtime_store", + } + } + + pub fn parse(value: &str) -> Result<Self, AppSqliteError> { + match value { + "app_workflow" => Ok(Self::AppWorkflow), + "shared_runtime_store" => Ok(Self::SharedRuntimeStore), + _ => Err(AppSqliteError::DecodeEnum { + field: "desktop_runtime_workflow_receipts.source_kind", + value: value.to_owned(), + }), + } + } +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum DesktopRuntimeWorkflowReceiptState { + Pending, + Prepared, + Enqueued, + Pushed, + Failed, + Blocked, + Skipped, + Unsupported, + ManualReview, + Unknown, +} + +impl DesktopRuntimeWorkflowReceiptState { + 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: "desktop_runtime_workflow_receipts.workflow_state", + value: value.to_owned(), + }), + } + } +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct DesktopRuntimeWorkflowReceiptInput { + pub source_kind: DesktopRuntimeWorkflowReceiptSourceKind, + pub source_record_id: String, + pub sdk_operation_kind: String, + pub runtime_effect_ids: Vec<String>, + pub expected_event_id: Option<String>, + pub actor_pubkey: Option<String>, + pub idempotency_digest_prefix: Option<String>, + pub workflow_state: DesktopRuntimeWorkflowReceiptState, + pub recorded_at: String, + pub detail_json: Value, +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct DesktopRuntimeStoredWorkflowReceipt { + pub id: String, + pub source_kind: DesktopRuntimeWorkflowReceiptSourceKind, + pub source_record_id: String, + pub sdk_operation_kind: String, + pub runtime_effect_ids: Vec<String>, + pub expected_event_id: Option<String>, + pub actor_pubkey: Option<String>, + pub idempotency_digest_prefix: Option<String>, + pub workflow_state: DesktopRuntimeWorkflowReceiptState, + pub created_at: String, + pub updated_at: String, + pub detail_json: Value, +} + +pub struct DesktopRuntimeWorkflowReceiptRepository<'a> { + connection: &'a AppSqliteDatabase, +} + +impl<'a> DesktopRuntimeWorkflowReceiptRepository<'a> { + pub(crate) const fn new(connection: &'a AppSqliteDatabase) -> Self { + Self { connection } + } + + pub fn record_receipt( + &self, + input: &DesktopRuntimeWorkflowReceiptInput, + ) -> Result<DesktopRuntimeStoredWorkflowReceipt, AppSqliteError> { + let receipt_id = Uuid::now_v7().to_string(); + let effect_ids_json = + serde_json::to_string(&input.runtime_effect_ids).map_err(|source| { + AppSqliteError::EncodeJson { + field: "desktop_runtime_workflow_receipts.runtime_effect_ids_json", + source, + } + })?; + let detail_json = serde_json::to_string(&input.detail_json).map_err(|source| { + AppSqliteError::EncodeJson { + field: "desktop_runtime_workflow_receipts.detail_json", + source, + } + })?; + + self.connection + .execute( + "INSERT INTO desktop_runtime_workflow_receipts ( + id, + source_kind, + source_record_id, + sdk_operation_kind, + runtime_effect_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, + runtime_effect_ids_json = excluded.runtime_effect_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", + crate::app_sqlite_params![ + receipt_id, + input.source_kind.storage_key(), + input.source_record_id.as_str(), + input.sdk_operation_kind.as_str(), + effect_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 desktop runtime workflow receipt", + source, + })?; + + self.load_receipt(input.source_kind, input.source_record_id.as_str())? + .ok_or(AppSqliteError::MissingColumn { + field: "desktop_runtime_workflow_receipts.id", + }) + } + + pub fn load_receipt( + &self, + source_kind: DesktopRuntimeWorkflowReceiptSourceKind, + source_record_id: &str, + ) -> Result<Option<DesktopRuntimeStoredWorkflowReceipt>, AppSqliteError> { + self.connection + .query_row( + "SELECT + id, + source_kind, + source_record_id, + sdk_operation_kind, + runtime_effect_ids_json, + expected_event_id, + actor_pubkey, + idempotency_digest_prefix, + workflow_state, + created_at, + updated_at, + detail_json + FROM desktop_runtime_workflow_receipts + WHERE source_kind = ?1 + AND source_record_id = ?2 + LIMIT 1", + crate::app_sqlite_params![source_kind.storage_key(), source_record_id], + decode_receipt_row, + ) + .optional() + .map_err(|source| AppSqliteError::Query { + operation: "load desktop runtime workflow receipt", + source, + }) + } +} + +fn decode_receipt_row( + row: &sqlx::sqlite::SqliteRow, +) -> Result<DesktopRuntimeStoredWorkflowReceipt, sqlx::Error> { + let source_kind: String = row.try_get(1)?; + let effect_ids_json: String = row.try_get(4)?; + let workflow_state: String = row.try_get(8)?; + let detail_json: String = row.try_get(11)?; + Ok(DesktopRuntimeStoredWorkflowReceipt { + id: row.try_get(0)?, + source_kind: DesktopRuntimeWorkflowReceiptSourceKind::parse(source_kind.as_str()) + .map_err(decode_app_error)?, + source_record_id: row.try_get(2)?, + sdk_operation_kind: row.try_get(3)?, + runtime_effect_ids: serde_json::from_str(effect_ids_json.as_str()).map_err(|source| { + decode_app_error(AppSqliteError::DecodeJson { + field: "desktop_runtime_workflow_receipts.runtime_effect_ids_json", + source, + }) + })?, + expected_event_id: row.try_get(5)?, + actor_pubkey: row.try_get(6)?, + idempotency_digest_prefix: row.try_get(7)?, + workflow_state: DesktopRuntimeWorkflowReceiptState::parse(workflow_state.as_str()) + .map_err(decode_app_error)?, + created_at: row.try_get(9)?, + updated_at: row.try_get(10)?, + detail_json: serde_json::from_str(detail_json.as_str()).map_err(|source| { + decode_app_error(AppSqliteError::DecodeJson { + field: "desktop_runtime_workflow_receipts.detail_json", + source, + }) + })?, + }) +} + +fn decode_app_error(error: AppSqliteError) -> sqlx::Error { + sqlx::Error::Decode(Box::new(error)) +} + +#[cfg(test)] +mod tests { + use serde_json::json; + + use crate::{ + AppSqliteStore, DatabaseTarget, DesktopRuntimeWorkflowReceiptInput, + DesktopRuntimeWorkflowReceiptSourceKind, DesktopRuntimeWorkflowReceiptState, + }; + + #[test] + fn workflow_receipts_are_idempotent_by_source_record() { + let store = AppSqliteStore::open(DatabaseTarget::InMemory).expect("open app store"); + let first = store + .runtime_workflow_receipt_repository() + .record_receipt(&DesktopRuntimeWorkflowReceiptInput { + source_kind: DesktopRuntimeWorkflowReceiptSourceKind::AppWorkflow, + source_record_id: "source-record-a".to_owned(), + sdk_operation_kind: "farm.publish".to_owned(), + runtime_effect_ids: vec!["effect-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: DesktopRuntimeWorkflowReceiptState::Enqueued, + recorded_at: "2026-06-18T12:00:00Z".to_owned(), + detail_json: json!({"attempt": 1}), + }) + .expect("record first receipt"); + let second = store + .runtime_workflow_receipt_repository() + .record_receipt(&DesktopRuntimeWorkflowReceiptInput { + source_kind: DesktopRuntimeWorkflowReceiptSourceKind::AppWorkflow, + source_record_id: "source-record-a".to_owned(), + sdk_operation_kind: "farm.publish".to_owned(), + runtime_effect_ids: vec!["effect-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: DesktopRuntimeWorkflowReceiptState::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.runtime_effect_ids, vec!["effect-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, + DesktopRuntimeWorkflowReceiptState::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 @@ -1,320 +0,0 @@ -use sqlx::Row; - -use crate::{AppSqliteDatabase, OptionalSqliteResult}; -use serde_json::Value; -use uuid::Uuid; - -use crate::AppSqliteError; - -#[derive(Clone, Copy, Debug, Eq, PartialEq)] -pub enum AppSdkWorkflowReceiptSourceKind { - LocalOutbox, - SharedRuntimeStore, -} - -impl AppSdkWorkflowReceiptSourceKind { - pub const fn storage_key(self) -> &'static str { - match self { - Self::LocalOutbox => "local_outbox", - Self::SharedRuntimeStore => "shared_runtime_store", - } - } - - pub fn parse(value: &str) -> Result<Self, AppSqliteError> { - match value { - "local_outbox" => Ok(Self::LocalOutbox), - "shared_runtime_store" => Ok(Self::SharedRuntimeStore), - _ => 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 AppSqliteDatabase, -} - -impl<'a> AppSdkWorkflowReceiptRepository<'a> { - pub(crate) const fn new(connection: &'a AppSqliteDatabase) -> 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", - crate::app_sqlite_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", - crate::app_sqlite_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: &sqlx::sqlite::SqliteRow, -) -> Result<AppSdkStoredWorkflowReceipt, sqlx::Error> { - let source_kind: String = row.try_get(1)?; - let outbox_ids_json: String = row.try_get(4)?; - let workflow_state: String = row.try_get(8)?; - let detail_json: String = row.try_get(11)?; - Ok(AppSdkStoredWorkflowReceipt { - id: row.try_get(0)?, - source_kind: AppSdkWorkflowReceiptSourceKind::parse(source_kind.as_str()) - .map_err(decode_app_error)?, - source_record_id: row.try_get(2)?, - sdk_operation_kind: row.try_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.try_get(5)?, - actor_pubkey: row.try_get(6)?, - idempotency_digest_prefix: row.try_get(7)?, - workflow_state: AppSdkWorkflowReceiptState::parse(workflow_state.as_str()) - .map_err(decode_app_error)?, - created_at: row.try_get(9)?, - updated_at: row.try_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) -> sqlx::Error { - sqlx::Error::Decode(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/i18n/locales/en/messages.json b/i18n/locales/en/messages.json @@ -763,8 +763,9 @@ "metadata.sdk_issue_retryable": "SDK issue retryable", "metadata.sdk_recovery_action": "SDK recovery", "metadata.sdk_storage_root": "SDK storage root", - "metadata.sdk_event_store_path": "SDK event store", - "metadata.sdk_outbox_path": "SDK outbox", + "metadata.sdk_runtime_path": "SDK runtime store", + "metadata.sdk_private_path": "SDK private store", + "metadata.sdk_studio_path": "Studio store", "metadata.sdk_relay_url_policy": "SDK relay policy", "value.none": "none", "value.yes": "yes",