lib

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

commit 4caff310a8c4814bf0551b09ec9733410e464157
parent eff18cb4c3d5ed9c8bee656f20d19e4dcadfb103
Author: triesap <tyson@radroots.org>
Date:   Sat, 11 Jul 2026 23:24:28 +0000

crate-surface: rename runtime store crate

- rename radroots_local_events to radroots_runtime_store
- move shared runtime paths and SQLite filename to runtime_store
- preserve relay delivery metadata fields and migrations
- update runtime-store README, tests, and destructive namespace contract

Diffstat:
MCargo.lock | 2+-
MCargo.toml | 4++--
Mcontracts/coverage.toml | 2+-
Dcrates/local_events/Cargo.toml | 29-----------------------------
Dcrates/local_events/README | 8--------
Dcrates/local_events/migrations/0000_local_events.down.sql | 2--
Dcrates/local_events/migrations/0000_local_events.up.sql | 41-----------------------------------------
Dcrates/local_events/migrations/0001_change_tracking.down.sql | 108-------------------------------------------------------------------------------
Dcrates/local_events/migrations/0001_change_tracking.up.sql | 113-------------------------------------------------------------------------------
Dcrates/local_events/migrations/0002_network_source_runtime.up.sql | 91-------------------------------------------------------------------------------
Dcrates/local_events/src/error.rs | 14--------------
Dcrates/local_events/src/lib.rs | 24------------------------
Dcrates/local_events/src/migrations.rs | 55-------------------------------------------------------
Dcrates/local_events/src/models.rs | 487-------------------------------------------------------------------------------
Dcrates/local_events/src/order_work.rs | 806-------------------------------------------------------------------------------
Dcrates/local_events/src/store.rs | 933-------------------------------------------------------------------------------
Dcrates/local_events/tests/order_work.rs | 294-------------------------------------------------------------------------------
Dcrates/local_events/tests/store.rs | 499-------------------------------------------------------------------------------
Mcrates/runtime_paths/src/conventions.rs | 77++++++++++++++++++++++++++++++++++++++---------------------------------------
Mcrates/runtime_paths/src/lib.rs | 13++++++-------
Acrates/runtime_store/Cargo.toml | 29+++++++++++++++++++++++++++++
Acrates/runtime_store/README | 8++++++++
Acrates/runtime_store/migrations/0000_runtime_store.down.sql | 2++
Acrates/runtime_store/migrations/0000_runtime_store.up.sql | 43+++++++++++++++++++++++++++++++++++++++++++
Acrates/runtime_store/migrations/0001_change_tracking.down.sql | 114+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/runtime_store/migrations/0001_change_tracking.up.sql | 119+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Rcrates/local_events/migrations/0002_network_source_runtime.down.sql -> crates/runtime_store/migrations/0002_network_source_runtime.down.sql | 0
Acrates/runtime_store/migrations/0002_network_source_runtime.up.sql | 97+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/runtime_store/src/error.rs | 14++++++++++++++
Acrates/runtime_store/src/lib.rs | 25+++++++++++++++++++++++++
Acrates/runtime_store/src/migrations.rs | 55+++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/runtime_store/src/models.rs | 757+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/runtime_store/src/order_work.rs | 806+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/runtime_store/src/store.rs | 970+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/runtime_store/tests/order_work.rs | 294+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/runtime_store/tests/store.rs | 545+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
36 files changed, 3926 insertions(+), 3554 deletions(-)

diff --git a/Cargo.lock b/Cargo.lock @@ -4367,7 +4367,7 @@ dependencies = [ ] [[package]] -name = "radroots_local_events" +name = "radroots_runtime_store" version = "0.1.0-alpha.2" dependencies = [ "radroots_sql_core", diff --git a/Cargo.toml b/Cargo.toml @@ -8,7 +8,7 @@ members = [ "crates/authority", "crates/geocoder", "crates/identity", - "crates/local_events", + "crates/runtime_store", "crates/log", "crates/net", "crates/nostr", @@ -72,7 +72,7 @@ radroots_event_index = { path = "crates/event_index", version = "0.1.0-alpha.2", radroots_authority = { path = "crates/authority", version = "0.1.0-alpha.2", default-features = false } radroots_geocoder = { path = "crates/geocoder", version = "0.1.0-alpha.2" } radroots_identity = { path = "crates/identity", version = "0.1.0-alpha.2", default-features = false } -radroots_local_events = { path = "crates/local_events", version = "0.1.0-alpha.2", default-features = false } +radroots_runtime_store = { path = "crates/runtime_store", version = "0.1.0-alpha.2", default-features = false } radroots_nostr = { path = "crates/nostr", version = "0.1.0-alpha.2", default-features = false } radroots_nostr_accounts = { path = "crates/nostr_accounts", version = "0.1.0-alpha.2", default-features = false } radroots_nostr_connect = { path = "crates/nostr_connect", version = "0.1.0-alpha.2", default-features = false } diff --git a/contracts/coverage.toml b/contracts/coverage.toml @@ -39,7 +39,7 @@ crates = [ "radroots_event_index", "radroots_geocoder", "radroots_identity", - "radroots_local_events", + "radroots_runtime_store", "radroots_log", "radroots_mesh", "radroots_mesh_agent_client", diff --git a/crates/local_events/Cargo.toml b/crates/local_events/Cargo.toml @@ -1,29 +0,0 @@ -[package] -name = "radroots_local_events" -publish = ["crates-io"] -version = "0.1.0-alpha.2" -edition.workspace = true -authors = ["Tyson Lupul <tyson@radroots.org>"] -rust-version.workspace = true -license.workspace = true -description = "Local event workspace for Radroots" -repository.workspace = true -homepage.workspace = true -documentation = "https://docs.rs/radroots_local_events" -readme = "README" - -[lib] -crate-type = ["rlib"] - -[features] -default = [] -native = ["radroots_sql_core/native"] - -[dependencies] -radroots_sql_core = { workspace = true } -serde = { workspace = true } -serde_json = { workspace = true } -thiserror = { workspace = true } - -[dev-dependencies] -radroots_sql_core = { workspace = true, features = ["native"] } diff --git a/crates/local_events/README b/crates/local_events/README @@ -1,8 +0,0 @@ -# radroots_local_events - -Shared local event and work-store primitives for same-host Rad Roots runtimes. - -This crate owns the SQLite schema and typed API for the `shared/local_events` -namespace. It is an interop store for local work records, signed event records, -publish outbox status, relay delivery metadata, and projection cursors. It is -not an application primary database. diff --git a/crates/local_events/migrations/0000_local_events.down.sql b/crates/local_events/migrations/0000_local_events.down.sql @@ -1,2 +0,0 @@ -drop table if exists local_event_projection_cursor; -drop table if exists local_event_record; diff --git a/crates/local_events/migrations/0000_local_events.up.sql b/crates/local_events/migrations/0000_local_events.up.sql @@ -1,41 +0,0 @@ -create table if not exists local_event_record ( - seq integer primary key autoincrement, - record_id text not null unique, - family text not null check (family in ('local_work', 'signed_event')), - status text not null check (status in ('local_draft', 'local_saved', 'pending_publish', 'published', 'failed', 'conflict')), - source_runtime text not null check (source_runtime in ('cli', 'app', 'network', 'service', 'worker', 'test')), - created_at_ms integer not null, - inserted_at_ms integer not null, - updated_at_ms integer not null, - owner_account_id text, - owner_pubkey text, - farm_id text, - listing_addr text, - local_work_json text, - event_id text, - event_kind integer, - event_pubkey text, - event_created_at integer, - event_tags_json text, - event_content text, - event_sig text, - raw_event_json text, - outbox_status text not null check (outbox_status in ('none', 'pending', 'acknowledged', 'failed')), - check (trim(record_id) <> ''), - check (family <> 'local_work' or local_work_json is not null), - check (family <> 'local_work' or outbox_status = 'none'), - check (family <> 'signed_event' or (event_id is not null and event_kind is not null and event_pubkey is not null and event_sig is not null and raw_event_json is not null)) -); - -create index if not exists local_event_record_event_id_idx on local_event_record(event_id); -create index if not exists local_event_record_listing_addr_idx on local_event_record(listing_addr); -create index if not exists local_event_record_owner_pubkey_idx on local_event_record(owner_pubkey); -create index if not exists local_event_record_status_idx on local_event_record(status); - -create table if not exists local_event_projection_cursor ( - consumer_id text primary key, - last_seq integer not null, - updated_at_ms integer not null, - check (trim(consumer_id) <> ''), - check (last_seq >= 0) -); diff --git a/crates/local_events/migrations/0001_change_tracking.down.sql b/crates/local_events/migrations/0001_change_tracking.down.sql @@ -1,108 +0,0 @@ -create table local_event_projection_cursor_previous ( - consumer_id text primary key, - last_seq integer not null, - updated_at_ms integer not null, - check (trim(consumer_id) <> ''), - check (last_seq >= 0) -); - -insert into local_event_projection_cursor_previous( - consumer_id, - last_seq, - updated_at_ms -) -select - consumer_id, - last_change_seq, - updated_at_ms -from local_event_projection_cursor; - -drop table local_event_projection_cursor; -alter table local_event_projection_cursor_previous rename to local_event_projection_cursor; - -create table local_event_record_previous ( - seq integer primary key autoincrement, - record_id text not null unique, - family text not null check (family in ('local_work', 'signed_event')), - status text not null check (status in ('local_draft', 'local_saved', 'pending_publish', 'published', 'failed', 'conflict')), - source_runtime text not null check (source_runtime in ('cli', 'app', 'network', 'service', 'worker', 'test')), - created_at_ms integer not null, - inserted_at_ms integer not null, - updated_at_ms integer not null, - owner_account_id text, - owner_pubkey text, - farm_id text, - listing_addr text, - local_work_json text, - event_id text, - event_kind integer, - event_pubkey text, - event_created_at integer, - event_tags_json text, - event_content text, - event_sig text, - raw_event_json text, - outbox_status text not null check (outbox_status in ('none', 'pending', 'acknowledged', 'failed')), - check (trim(record_id) <> ''), - check (family <> 'local_work' or local_work_json is not null), - check (family <> 'local_work' or outbox_status = 'none'), - check (family <> 'signed_event' or (event_id is not null and event_kind is not null and event_pubkey is not null and event_sig is not null and raw_event_json is not null)) -); - -insert into local_event_record_previous( - seq, - record_id, - family, - status, - source_runtime, - created_at_ms, - inserted_at_ms, - updated_at_ms, - owner_account_id, - owner_pubkey, - farm_id, - listing_addr, - local_work_json, - event_id, - event_kind, - event_pubkey, - event_created_at, - event_tags_json, - event_content, - event_sig, - raw_event_json, - outbox_status -) -select - seq, - record_id, - family, - status, - source_runtime, - created_at_ms, - inserted_at_ms, - updated_at_ms, - owner_account_id, - owner_pubkey, - farm_id, - listing_addr, - local_work_json, - event_id, - event_kind, - event_pubkey, - event_created_at, - event_tags_json, - event_content, - event_sig, - raw_event_json, - outbox_status -from local_event_record -order by seq asc; - -drop table local_event_record; -alter table local_event_record_previous rename to local_event_record; - -create index local_event_record_event_id_idx on local_event_record(event_id); -create index local_event_record_listing_addr_idx on local_event_record(listing_addr); -create index local_event_record_owner_pubkey_idx on local_event_record(owner_pubkey); -create index local_event_record_status_idx on local_event_record(status); diff --git a/crates/local_events/migrations/0001_change_tracking.up.sql b/crates/local_events/migrations/0001_change_tracking.up.sql @@ -1,113 +0,0 @@ -create table local_event_record_next ( - seq integer primary key autoincrement, - change_seq integer not null unique, - record_id text not null unique, - family text not null check (family in ('local_work', 'signed_event')), - status text not null check (status in ('local_draft', 'local_saved', 'pending_publish', 'published', 'failed', 'conflict')), - source_runtime text not null check (source_runtime in ('cli', 'app', 'network', 'service', 'worker', 'test')), - created_at_ms integer not null, - inserted_at_ms integer not null, - updated_at_ms integer not null, - owner_account_id text, - owner_pubkey text, - farm_id text, - listing_addr text, - local_work_json text, - event_id text, - event_kind integer, - event_pubkey text, - event_created_at integer, - event_tags_json text, - event_content text, - event_sig text, - raw_event_json text, - outbox_status text not null check (outbox_status in ('none', 'pending', 'acknowledged', 'failed')), - check (change_seq >= 1), - check (trim(record_id) <> ''), - check (family <> 'local_work' or local_work_json is not null), - check (family <> 'local_work' or outbox_status = 'none'), - check (family <> 'signed_event' or (event_id is not null and event_kind is not null and event_pubkey is not null and event_sig is not null and raw_event_json is not null)) -); - -insert into local_event_record_next( - seq, - change_seq, - record_id, - family, - status, - source_runtime, - created_at_ms, - inserted_at_ms, - updated_at_ms, - owner_account_id, - owner_pubkey, - farm_id, - listing_addr, - local_work_json, - event_id, - event_kind, - event_pubkey, - event_created_at, - event_tags_json, - event_content, - event_sig, - raw_event_json, - outbox_status -) -select - seq, - seq, - record_id, - family, - status, - source_runtime, - created_at_ms, - inserted_at_ms, - updated_at_ms, - owner_account_id, - owner_pubkey, - farm_id, - listing_addr, - local_work_json, - event_id, - event_kind, - event_pubkey, - event_created_at, - event_tags_json, - event_content, - event_sig, - raw_event_json, - outbox_status -from local_event_record -order by seq asc; - -drop table local_event_record; -alter table local_event_record_next rename to local_event_record; - -create index local_event_record_change_seq_idx on local_event_record(change_seq); -create index local_event_record_event_id_idx on local_event_record(event_id); -create index local_event_record_listing_addr_idx on local_event_record(listing_addr); -create index local_event_record_owner_pubkey_idx on local_event_record(owner_pubkey); -create index local_event_record_status_idx on local_event_record(status); - -create table local_event_projection_cursor_next ( - consumer_id text primary key, - last_change_seq integer not null, - updated_at_ms integer not null, - check (trim(consumer_id) <> ''), - check (last_change_seq >= 0) -); - -insert into local_event_projection_cursor_next( - consumer_id, - last_change_seq, - updated_at_ms -) -select - consumer_id, - last_seq, - updated_at_ms -from local_event_projection_cursor; - -drop table local_event_projection_cursor; -alter table local_event_projection_cursor_next rename to local_event_projection_cursor; diff --git a/crates/local_events/migrations/0002_network_source_runtime.up.sql b/crates/local_events/migrations/0002_network_source_runtime.up.sql @@ -1,91 +0,0 @@ -create table local_event_record_network_source_next ( - seq integer primary key autoincrement, - change_seq integer not null unique, - record_id text not null unique, - family text not null check (family in ('local_work', 'signed_event')), - status text not null check (status in ('local_draft', 'local_saved', 'pending_publish', 'published', 'failed', 'conflict')), - source_runtime text not null check (source_runtime in ('cli', 'app', 'network', 'service', 'worker', 'test')), - created_at_ms integer not null, - inserted_at_ms integer not null, - updated_at_ms integer not null, - owner_account_id text, - owner_pubkey text, - farm_id text, - listing_addr text, - local_work_json text, - event_id text, - event_kind integer, - event_pubkey text, - event_created_at integer, - event_tags_json text, - event_content text, - event_sig text, - raw_event_json text, - outbox_status text not null check (outbox_status in ('none', 'pending', 'acknowledged', 'failed')), - check (change_seq >= 1), - check (trim(record_id) <> ''), - check (family <> 'local_work' or local_work_json is not null), - check (family <> 'local_work' or outbox_status = 'none'), - check (family <> 'signed_event' or (event_id is not null and event_kind is not null and event_pubkey is not null and event_sig is not null and raw_event_json is not null)) -); - -insert into local_event_record_network_source_next( - seq, - change_seq, - record_id, - family, - status, - source_runtime, - created_at_ms, - inserted_at_ms, - updated_at_ms, - owner_account_id, - owner_pubkey, - farm_id, - listing_addr, - local_work_json, - event_id, - event_kind, - event_pubkey, - event_created_at, - event_tags_json, - event_content, - event_sig, - raw_event_json, - outbox_status -) -select - seq, - change_seq, - record_id, - family, - status, - source_runtime, - created_at_ms, - inserted_at_ms, - updated_at_ms, - owner_account_id, - owner_pubkey, - farm_id, - listing_addr, - local_work_json, - event_id, - event_kind, - event_pubkey, - event_created_at, - event_tags_json, - event_content, - event_sig, - raw_event_json, - outbox_status -from local_event_record -order by seq asc; - -drop table local_event_record; -alter table local_event_record_network_source_next rename to local_event_record; - -create index local_event_record_change_seq_idx on local_event_record(change_seq); -create index local_event_record_event_id_idx on local_event_record(event_id); -create index local_event_record_listing_addr_idx on local_event_record(listing_addr); -create index local_event_record_owner_pubkey_idx on local_event_record(owner_pubkey); -create index local_event_record_status_idx on local_event_record(status); diff --git a/crates/local_events/src/error.rs b/crates/local_events/src/error.rs @@ -1,14 +0,0 @@ -#![forbid(unsafe_code)] - -use radroots_sql_core::error::SqlError; -use thiserror::Error; - -#[derive(Debug, Error)] -pub enum LocalEventsError { - #[error("invalid local event record: {0}")] - InvalidRecord(String), - #[error("sql error: {0}")] - Sql(#[from] SqlError), - #[error("serialization error: {0}")] - Serialization(#[from] serde_json::Error), -} diff --git a/crates/local_events/src/lib.rs b/crates/local_events/src/lib.rs @@ -1,24 +0,0 @@ -#![forbid(unsafe_code)] - -mod error; -mod migrations; -mod models; -mod order_work; -mod store; - -pub use error::LocalEventsError; -pub use migrations::{MIGRATIONS, run_all_down, run_all_up}; -pub use models::{ - LocalEventRecord, LocalEventRecordInput, LocalEventRecordUpdate, LocalEventsCursor, - LocalRecordFamily, LocalRecordStatus, PublishOutboxStatus, SourceRuntime, -}; -pub use order_work::{ - BUYER_ORDER_REQUEST_ACTOR_SOURCE_RESOLVED_ACCOUNT, - BUYER_ORDER_REQUEST_ACTOR_SOURCE_UNRESOLVED_APP, BUYER_ORDER_REQUEST_DOCUMENT_KIND, - BUYER_ORDER_REQUEST_LOCAL_WORK_RECORD_KIND, BuyerOrderRequestLocalWorkValidation, - BuyerOrderRequestSupportState, buyer_order_request_local_work_record_id, - validate_buyer_order_request_local_work_payload, - validate_supported_buyer_order_request_local_work_payload, - validate_unsupported_buyer_order_request_local_work_payload, -}; -pub use store::LocalEventsStore; diff --git a/crates/local_events/src/migrations.rs b/crates/local_events/src/migrations.rs @@ -1,55 +0,0 @@ -#![forbid(unsafe_code)] - -use radroots_sql_core::SqlExecutor; -use radroots_sql_core::error::SqlError; -use radroots_sql_core::migrations::{Migration, migrations_run_all_down, migrations_run_all_up}; - -pub static MIGRATIONS: &[Migration] = &[ - Migration { - name: "0000_local_events", - up_sql: include_str!("../migrations/0000_local_events.up.sql"), - down_sql: include_str!("../migrations/0000_local_events.down.sql"), - }, - Migration { - name: "0001_change_tracking", - up_sql: include_str!("../migrations/0001_change_tracking.up.sql"), - down_sql: include_str!("../migrations/0001_change_tracking.down.sql"), - }, - Migration { - name: "0002_network_source_runtime", - up_sql: include_str!("../migrations/0002_network_source_runtime.up.sql"), - down_sql: include_str!("../migrations/0002_network_source_runtime.down.sql"), - }, -]; - -pub fn run_all_up<E>(executor: &E) -> Result<(), SqlError> -where - E: SqlExecutor, -{ - migrations_run_all_up(executor, MIGRATIONS) -} - -pub fn run_all_down<E>(executor: &E) -> Result<(), SqlError> -where - E: SqlExecutor, -{ - migrations_run_all_down(executor, MIGRATIONS) -} - -#[cfg(test)] -mod tests { - use radroots_sql_core::SqliteExecutor; - - use super::*; - - #[test] - fn migration_entrypoints_apply_and_reverse_schema() { - let executor = SqliteExecutor::open_memory().expect("open memory sqlite"); - - run_all_up(&executor).expect("migrate up"); - executor - .query_raw("select name from __migrations order by name", "[]") - .expect("query migrations"); - run_all_down(&executor).expect("migrate down"); - } -} diff --git a/crates/local_events/src/models.rs b/crates/local_events/src/models.rs @@ -1,487 +0,0 @@ -#![forbid(unsafe_code)] - -use serde::{Deserialize, Serialize}; -use serde_json::Value; - -use crate::LocalEventsError; - -#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] -#[serde(rename_all = "snake_case")] -pub enum LocalRecordFamily { - LocalWork, - SignedEvent, -} - -impl LocalRecordFamily { - pub fn as_str(self) -> &'static str { - match self { - Self::LocalWork => "local_work", - Self::SignedEvent => "signed_event", - } - } - - pub fn parse(value: &str) -> Result<Self, LocalEventsError> { - match value { - "local_work" => Ok(Self::LocalWork), - "signed_event" => Ok(Self::SignedEvent), - other => Err(LocalEventsError::InvalidRecord(format!( - "unknown record family `{other}`" - ))), - } - } -} - -#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] -#[serde(rename_all = "snake_case")] -pub enum LocalRecordStatus { - LocalDraft, - LocalSaved, - PendingPublish, - Published, - Failed, - Conflict, -} - -impl LocalRecordStatus { - pub fn as_str(self) -> &'static str { - match self { - Self::LocalDraft => "local_draft", - Self::LocalSaved => "local_saved", - Self::PendingPublish => "pending_publish", - Self::Published => "published", - Self::Failed => "failed", - Self::Conflict => "conflict", - } - } - - pub fn parse(value: &str) -> Result<Self, LocalEventsError> { - match value { - "local_draft" => Ok(Self::LocalDraft), - "local_saved" => Ok(Self::LocalSaved), - "pending_publish" => Ok(Self::PendingPublish), - "published" => Ok(Self::Published), - "failed" => Ok(Self::Failed), - "conflict" => Ok(Self::Conflict), - other => Err(LocalEventsError::InvalidRecord(format!( - "unknown record status `{other}`" - ))), - } - } -} - -#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] -#[serde(rename_all = "snake_case")] -pub enum PublishOutboxStatus { - None, - Pending, - Acknowledged, - Failed, -} - -impl PublishOutboxStatus { - pub fn as_str(self) -> &'static str { - match self { - Self::None => "none", - Self::Pending => "pending", - Self::Acknowledged => "acknowledged", - Self::Failed => "failed", - } - } - - pub fn parse(value: &str) -> Result<Self, LocalEventsError> { - match value { - "none" => Ok(Self::None), - "pending" => Ok(Self::Pending), - "acknowledged" => Ok(Self::Acknowledged), - "failed" => Ok(Self::Failed), - other => Err(LocalEventsError::InvalidRecord(format!( - "unknown outbox status `{other}`" - ))), - } - } -} - -#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] -#[serde(rename_all = "snake_case")] -pub enum SourceRuntime { - Cli, - App, - Network, - Service, - Worker, - Test, -} - -impl SourceRuntime { - pub fn as_str(self) -> &'static str { - match self { - Self::Cli => "cli", - Self::App => "app", - Self::Network => "network", - Self::Service => "service", - Self::Worker => "worker", - Self::Test => "test", - } - } - - pub fn parse(value: &str) -> Result<Self, LocalEventsError> { - match value { - "cli" => Ok(Self::Cli), - "app" => Ok(Self::App), - "network" => Ok(Self::Network), - "service" => Ok(Self::Service), - "worker" => Ok(Self::Worker), - "test" => Ok(Self::Test), - other => Err(LocalEventsError::InvalidRecord(format!( - "unknown source runtime `{other}`" - ))), - } - } -} - -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct LocalEventRecordInput { - pub record_id: String, - pub family: LocalRecordFamily, - pub status: LocalRecordStatus, - pub source_runtime: SourceRuntime, - pub created_at_ms: i64, - pub inserted_at_ms: i64, - pub owner_account_id: Option<String>, - pub owner_pubkey: Option<String>, - pub farm_id: Option<String>, - pub listing_addr: Option<String>, - pub local_work_json: Option<Value>, - pub event_id: Option<String>, - pub event_kind: Option<i64>, - pub event_pubkey: Option<String>, - pub event_created_at: Option<i64>, - pub event_tags_json: Option<Value>, - pub event_content: Option<String>, - pub event_sig: Option<String>, - pub raw_event_json: Option<Value>, - pub outbox_status: PublishOutboxStatus, -} - -impl LocalEventRecordInput { - pub fn validate(&self) -> Result<(), LocalEventsError> { - validate_non_empty("record_id", &self.record_id)?; - if let Some(value) = self.owner_account_id.as_deref() { - validate_non_empty("owner_account_id", value)?; - } - if let Some(value) = self.owner_pubkey.as_deref() { - validate_non_empty("owner_pubkey", value)?; - } - if let Some(value) = self.farm_id.as_deref() { - validate_non_empty("farm_id", value)?; - } - if let Some(value) = self.listing_addr.as_deref() { - validate_non_empty("listing_addr", value)?; - } - match self.family { - LocalRecordFamily::LocalWork => { - if self.local_work_json.is_none() { - return Err(LocalEventsError::InvalidRecord( - "local work records require local_work_json".to_owned(), - )); - } - if self.outbox_status != PublishOutboxStatus::None { - return Err(LocalEventsError::InvalidRecord( - "local work records must use outbox status none".to_owned(), - )); - } - } - LocalRecordFamily::SignedEvent => { - validate_required("event_id", self.event_id.as_deref())?; - validate_required("event_pubkey", self.event_pubkey.as_deref())?; - validate_required("event_sig", self.event_sig.as_deref())?; - if self.event_kind.is_none() { - return Err(LocalEventsError::InvalidRecord( - "signed event records require event_kind".to_owned(), - )); - } - if self.raw_event_json.is_none() { - return Err(LocalEventsError::InvalidRecord( - "signed event records require raw_event_json".to_owned(), - )); - } - } - } - Ok(()) - } -} - -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct LocalEventRecord { - pub seq: i64, - pub change_seq: i64, - pub record_id: String, - pub family: LocalRecordFamily, - pub status: LocalRecordStatus, - pub source_runtime: SourceRuntime, - pub created_at_ms: i64, - pub inserted_at_ms: i64, - pub updated_at_ms: i64, - pub owner_account_id: Option<String>, - pub owner_pubkey: Option<String>, - pub farm_id: Option<String>, - pub listing_addr: Option<String>, - pub local_work_json: Option<Value>, - pub event_id: Option<String>, - pub event_kind: Option<i64>, - pub event_pubkey: Option<String>, - pub event_created_at: Option<i64>, - pub event_tags_json: Option<Value>, - pub event_content: Option<String>, - pub event_sig: Option<String>, - pub raw_event_json: Option<Value>, - pub outbox_status: PublishOutboxStatus, -} - -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct LocalEventRecordUpdate { - pub record_id: String, - pub status: LocalRecordStatus, - pub outbox_status: PublishOutboxStatus, - pub updated_at_ms: i64, -} - -#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] -pub struct LocalEventsCursor { - pub consumer_id: String, - pub last_change_seq: i64, - pub updated_at_ms: i64, -} - -pub(crate) fn validate_non_empty(field: &str, value: &str) -> Result<(), LocalEventsError> { - if value.trim().is_empty() { - return Err(LocalEventsError::InvalidRecord(format!( - "{field} must not be empty" - ))); - } - Ok(()) -} - -fn validate_required(field: &str, value: Option<&str>) -> Result<(), LocalEventsError> { - match value { - Some(value) => validate_non_empty(field, value), - None => Err(LocalEventsError::InvalidRecord(format!( - "{field} is required" - ))), - } -} - -#[cfg(test)] -mod tests { - use serde_json::json; - - use super::*; - - #[test] - fn enum_strings_and_parse_errors_cover_all_model_variants() { - for (variant, value) in [ - (LocalRecordFamily::LocalWork, "local_work"), - (LocalRecordFamily::SignedEvent, "signed_event"), - ] { - assert_eq!(variant.as_str(), value); - assert_eq!( - LocalRecordFamily::parse(value).expect("record family"), - variant - ); - } - - for (variant, value) in [ - (LocalRecordStatus::LocalDraft, "local_draft"), - (LocalRecordStatus::LocalSaved, "local_saved"), - (LocalRecordStatus::PendingPublish, "pending_publish"), - (LocalRecordStatus::Published, "published"), - (LocalRecordStatus::Failed, "failed"), - (LocalRecordStatus::Conflict, "conflict"), - ] { - assert_eq!(variant.as_str(), value); - assert_eq!( - LocalRecordStatus::parse(value).expect("record status"), - variant - ); - } - - for (variant, value) in [ - (PublishOutboxStatus::None, "none"), - (PublishOutboxStatus::Pending, "pending"), - (PublishOutboxStatus::Acknowledged, "acknowledged"), - (PublishOutboxStatus::Failed, "failed"), - ] { - assert_eq!(variant.as_str(), value); - assert_eq!( - PublishOutboxStatus::parse(value).expect("outbox status"), - variant - ); - } - - for (variant, value) in [ - (SourceRuntime::Cli, "cli"), - (SourceRuntime::App, "app"), - (SourceRuntime::Network, "network"), - (SourceRuntime::Service, "service"), - (SourceRuntime::Worker, "worker"), - (SourceRuntime::Test, "test"), - ] { - assert_eq!(variant.as_str(), value); - assert_eq!( - SourceRuntime::parse(value).expect("source runtime"), - variant - ); - } - - assert!(LocalRecordFamily::parse("other").is_err()); - assert!(LocalRecordStatus::parse("other").is_err()); - assert!(PublishOutboxStatus::parse("other").is_err()); - assert!(SourceRuntime::parse("other").is_err()); - } - - #[test] - fn local_record_input_validation_covers_success_and_error_paths() { - let mut local_work = local_work_input(); - local_work.validate().expect("valid local work"); - - for (field, update) in [ - ( - "owner_account_id", - Box::new(|input: &mut LocalEventRecordInput| { - input.owner_account_id = Some(" ".to_owned()); - }) as Box<dyn Fn(&mut LocalEventRecordInput)>, - ), - ( - "owner_pubkey", - Box::new(|input: &mut LocalEventRecordInput| { - input.owner_pubkey = Some(" ".to_owned()); - }), - ), - ( - "farm_id", - Box::new(|input: &mut LocalEventRecordInput| { - input.farm_id = Some(" ".to_owned()); - }), - ), - ( - "listing_addr", - Box::new(|input: &mut LocalEventRecordInput| { - input.listing_addr = Some(" ".to_owned()); - }), - ), - ] { - let mut input = local_work_input(); - update(&mut input); - assert_error_contains(input.validate(), field); - } - - local_work.record_id = " ".to_owned(); - assert_error_contains(local_work.validate(), "record_id"); - - let mut missing_work = local_work_input(); - missing_work.local_work_json = None; - assert_error_contains(missing_work.validate(), "local_work_json"); - - let mut queued_work = local_work_input(); - queued_work.outbox_status = PublishOutboxStatus::Pending; - assert_error_contains(queued_work.validate(), "outbox status none"); - - let signed_event = signed_event_input(); - signed_event.validate().expect("valid signed event"); - - for (field, update) in [ - ( - "event_id", - Box::new(|input: &mut LocalEventRecordInput| { - input.event_id = Some(" ".to_owned()); - }) as Box<dyn Fn(&mut LocalEventRecordInput)>, - ), - ( - "event_pubkey", - Box::new(|input: &mut LocalEventRecordInput| { - input.event_pubkey = None; - }), - ), - ( - "event_sig", - Box::new(|input: &mut LocalEventRecordInput| { - input.event_sig = None; - }), - ), - ( - "event_kind", - Box::new(|input: &mut LocalEventRecordInput| { - input.event_kind = None; - }), - ), - ( - "raw_event_json", - Box::new(|input: &mut LocalEventRecordInput| { - input.raw_event_json = None; - }), - ), - ] { - let mut input = signed_event_input(); - update(&mut input); - assert_error_contains(input.validate(), field); - } - } - - fn local_work_input() -> LocalEventRecordInput { - LocalEventRecordInput { - record_id: "local-work-a".to_owned(), - family: LocalRecordFamily::LocalWork, - status: LocalRecordStatus::LocalSaved, - source_runtime: SourceRuntime::App, - created_at_ms: 10, - inserted_at_ms: 11, - owner_account_id: Some("account-a".to_owned()), - owner_pubkey: Some("pubkey-a".to_owned()), - farm_id: Some("farm-a".to_owned()), - listing_addr: Some("listing-a".to_owned()), - local_work_json: Some(json!({"kind":"buyer_order_request_v1"})), - event_id: None, - event_kind: None, - event_pubkey: None, - event_created_at: None, - event_tags_json: None, - event_content: None, - event_sig: None, - raw_event_json: None, - outbox_status: PublishOutboxStatus::None, - } - } - - fn signed_event_input() -> LocalEventRecordInput { - LocalEventRecordInput { - record_id: "signed-event-a".to_owned(), - family: LocalRecordFamily::SignedEvent, - status: LocalRecordStatus::PendingPublish, - source_runtime: SourceRuntime::Service, - created_at_ms: 20, - inserted_at_ms: 21, - owner_account_id: None, - owner_pubkey: None, - farm_id: None, - listing_addr: None, - local_work_json: None, - event_id: Some("event-a".to_owned()), - event_kind: Some(30402), - event_pubkey: Some("pubkey-a".to_owned()), - event_created_at: Some(20), - event_tags_json: Some(json!([["d", "listing-a"]])), - event_content: Some("{}".to_owned()), - event_sig: Some("sig-a".to_owned()), - raw_event_json: Some(json!({"id":"event-a"})), - outbox_status: PublishOutboxStatus::Pending, - } - } - - fn assert_error_contains(result: Result<(), LocalEventsError>, expected: &str) { - let err = result.expect_err("validation error"); - assert!( - err.to_string().contains(expected), - "expected error to contain {expected}, got {err}" - ); - } -} diff --git a/crates/local_events/src/order_work.rs b/crates/local_events/src/order_work.rs @@ -1,806 +0,0 @@ -use serde_json::Value; - -use crate::LocalEventsError; -use crate::models::validate_non_empty; - -pub const BUYER_ORDER_REQUEST_LOCAL_WORK_RECORD_KIND: &str = "buyer_order_request_v1"; -pub const BUYER_ORDER_REQUEST_DOCUMENT_KIND: &str = "order_draft_v1"; -pub const BUYER_ORDER_REQUEST_ACTOR_SOURCE_RESOLVED_ACCOUNT: &str = "resolved_account"; -pub const BUYER_ORDER_REQUEST_ACTOR_SOURCE_UNRESOLVED_APP: &str = "app_unresolved"; - -#[derive(Clone, Copy, Debug, Eq, PartialEq)] -pub enum BuyerOrderRequestSupportState { - Supported, - Unsupported, -} - -impl BuyerOrderRequestSupportState { - pub fn as_str(self) -> &'static str { - match self { - Self::Supported => "supported", - Self::Unsupported => "unsupported", - } - } -} - -#[derive(Clone, Debug, Eq, PartialEq)] -pub struct BuyerOrderRequestLocalWorkValidation { - pub order_id: String, - pub support_state: BuyerOrderRequestSupportState, - pub support_issues: Vec<String>, -} - -pub fn buyer_order_request_local_work_record_id( - order_id: &str, -) -> Result<String, LocalEventsError> { - let order_id = order_id.trim(); - validate_non_empty("order_id", order_id)?; - Ok(format!("app:local_work:order_request:{order_id}")) -} - -pub fn validate_buyer_order_request_local_work_payload( - payload: &Value, -) -> Result<BuyerOrderRequestLocalWorkValidation, LocalEventsError> { - validate_string_field( - payload, - &["record_kind"], - BUYER_ORDER_REQUEST_LOCAL_WORK_RECORD_KIND, - )?; - validate_string_field(payload, &["scope"], "app")?; - validate_string_field( - payload, - &["document", "kind"], - BUYER_ORDER_REQUEST_DOCUMENT_KIND, - )?; - validate_bool_field(payload, &["currentness", "current"], true)?; - validate_string_field(payload, &["currentness", "source"], "app_sqlite_order")?; - - let order_id = validate_required_string(payload, &["document", "order", "order_id"])?; - let currentness_order_id = validate_required_string(payload, &["currentness", "order_id"])?; - if currentness_order_id != order_id { - return Err(invalid_field( - "currentness.order_id", - "must match document.order.order_id", - )); - } - validate_required_string(payload, &["currentness", "record_id"])?; - validate_positive_i64(payload, &["currentness", "created_at_ms"])?; - validate_required_string(payload, &["currentness", "order_updated_at"])?; - - let (support_state, support_issues) = validate_support_status(payload)?; - validate_exportability(payload, support_state)?; - validate_order_identity(payload, support_state)?; - validate_order_items(payload)?; - validate_order_economics(payload)?; - - Ok(BuyerOrderRequestLocalWorkValidation { - order_id: order_id.to_owned(), - support_state, - support_issues, - }) -} - -pub fn validate_supported_buyer_order_request_local_work_payload( - payload: &Value, -) -> Result<BuyerOrderRequestLocalWorkValidation, LocalEventsError> { - let validation = validate_buyer_order_request_local_work_payload(payload)?; - if validation.support_state != BuyerOrderRequestSupportState::Supported { - return Err(invalid_field( - "support_status.state", - "must be supported for exportable app order work", - )); - } - Ok(validation) -} - -pub fn validate_unsupported_buyer_order_request_local_work_payload( - payload: &Value, -) -> Result<BuyerOrderRequestLocalWorkValidation, LocalEventsError> { - let validation = validate_buyer_order_request_local_work_payload(payload)?; - if validation.support_state != BuyerOrderRequestSupportState::Unsupported { - return Err(invalid_field( - "support_status.state", - "must be unsupported for unsupported app order work", - )); - } - Ok(validation) -} - -fn validate_support_status( - payload: &Value, -) -> Result<(BuyerOrderRequestSupportState, Vec<String>), LocalEventsError> { - let state = validate_required_string(payload, &["support_status", "state"])?; - let issues = support_issues(payload)?; - match state { - "supported" => { - if !issues.is_empty() { - return Err(invalid_field( - "support_status.issues", - "must be empty when support_status.state is supported", - )); - } - Ok((BuyerOrderRequestSupportState::Supported, issues)) - } - "unsupported" => { - if issues.is_empty() { - return Err(invalid_field( - "support_status.issues", - "must contain at least one issue when support_status.state is unsupported", - )); - } - Ok((BuyerOrderRequestSupportState::Unsupported, issues)) - } - _ => Err(invalid_field( - "support_status.state", - "must be supported or unsupported", - )), - } -} - -fn validate_exportability( - payload: &Value, - support_state: BuyerOrderRequestSupportState, -) -> Result<(), LocalEventsError> { - let state = validate_required_string(payload, &["exportability", "state"])?; - match state { - "exportable" => { - validate_string_field( - payload, - &["document", "buyer_actor", "source"], - BUYER_ORDER_REQUEST_ACTOR_SOURCE_RESOLVED_ACCOUNT, - )?; - validate_buyer_pubkey(payload)?; - } - "identity_unresolved" => { - validate_required_string(payload, &["exportability", "reason"])?; - validate_string_field( - payload, - &["document", "buyer_actor", "source"], - BUYER_ORDER_REQUEST_ACTOR_SOURCE_UNRESOLVED_APP, - )?; - if support_state == BuyerOrderRequestSupportState::Supported { - return Err(invalid_field( - "exportability.state", - "supported app order work must be exportable", - )); - } - } - _ => { - return Err(invalid_field( - "exportability.state", - "must be exportable or identity_unresolved", - )); - } - } - Ok(()) -} - -fn validate_order_identity( - payload: &Value, - support_state: BuyerOrderRequestSupportState, -) -> Result<(), LocalEventsError> { - validate_required_string(payload, &["document", "order", "listing_addr"])?; - validate_required_string(payload, &["document", "order", "listing_event_id"])?; - validate_required_string(payload, &["document", "order", "seller_pubkey"])?; - if support_state == BuyerOrderRequestSupportState::Supported { - validate_buyer_pubkey(payload)?; - } - Ok(()) -} - -fn validate_buyer_pubkey(payload: &Value) -> Result<(), LocalEventsError> { - let order_buyer_pubkey = - validate_required_string(payload, &["document", "order", "buyer_pubkey"])?; - let actor_buyer_pubkey = - validate_required_string(payload, &["document", "buyer_actor", "pubkey"])?; - if order_buyer_pubkey != actor_buyer_pubkey { - return Err(invalid_field( - "document.buyer_actor.pubkey", - "must match document.order.buyer_pubkey", - )); - } - Ok(()) -} - -fn validate_order_items(payload: &Value) -> Result<(), LocalEventsError> { - let items = required_array(payload, &["document", "order", "items"])?; - if items.is_empty() { - return Err(invalid_field( - "document.order.items", - "must contain at least one item", - )); - } - for (index, item) in items.iter().enumerate() { - validate_required_string(item, &["bin_id"]).map_err(|_| { - invalid_field_at( - format!("document.order.items[{index}].bin_id"), - "is required", - ) - })?; - validate_positive_u64(item, &["bin_count"]).map_err(|_| { - invalid_field_at( - format!("document.order.items[{index}].bin_count"), - "must be positive", - ) - })?; - } - Ok(()) -} - -fn validate_order_economics(payload: &Value) -> Result<(), LocalEventsError> { - let economics = value_at(payload, &["document", "order", "economics"]).ok_or_else(|| { - invalid_field("document.order.economics", "is required for app order work") - })?; - if !economics.is_object() { - return Err(invalid_field( - "document.order.economics", - "must be an object", - )); - } - validate_string_field(economics, &["pricing_basis"], "listing_event")?; - let currency = validate_required_string(economics, &["currency"])?; - validate_currency("document.order.economics.currency", currency)?; - let economics_items = required_array(economics, &["items"])?; - let order_items = required_array(payload, &["document", "order", "items"])?; - if economics_items.is_empty() { - return Err(invalid_field( - "document.order.economics.items", - "must contain at least one item", - )); - } - if economics_items.len() != order_items.len() { - return Err(invalid_field( - "document.order.economics.items", - "must match document.order.items length", - )); - } - for (index, item) in economics_items.iter().enumerate() { - let order_item = &order_items[index]; - let economics_bin_id = validate_required_string(item, &["bin_id"]).map_err(|_| { - invalid_field_at( - format!("document.order.economics.items[{index}].bin_id"), - "is required", - ) - })?; - let order_bin_id = validate_required_string(order_item, &["bin_id"])?; - if economics_bin_id != order_bin_id { - return Err(invalid_field_at( - format!("document.order.economics.items[{index}].bin_id"), - "must match document.order.items bin_id", - )); - } - let economics_bin_count = validate_positive_u64(item, &["bin_count"]).map_err(|_| { - invalid_field_at( - format!("document.order.economics.items[{index}].bin_count"), - "must be positive", - ) - })?; - let order_bin_count = validate_positive_u64(order_item, &["bin_count"])?; - if economics_bin_count != order_bin_count { - return Err(invalid_field_at( - format!("document.order.economics.items[{index}].bin_count"), - "must match document.order.items bin_count", - )); - } - validate_required_string(item, &["quantity_amount"]).map_err(|_| { - invalid_field_at( - format!("document.order.economics.items[{index}].quantity_amount"), - "is required", - ) - })?; - validate_required_string(item, &["quantity_unit"]).map_err(|_| { - invalid_field_at( - format!("document.order.economics.items[{index}].quantity_unit"), - "is required", - ) - })?; - validate_required_string(item, &["unit_price_amount"]).map_err(|_| { - invalid_field_at( - format!("document.order.economics.items[{index}].unit_price_amount"), - "is required", - ) - })?; - let unit_price_currency = validate_required_string(item, &["unit_price_currency"])?; - if unit_price_currency != currency { - return Err(invalid_field_at( - format!("document.order.economics.items[{index}].unit_price_currency"), - "must match document.order.economics.currency", - )); - } - validate_money(item, &["line_subtotal"], currency)?; - } - validate_money(economics, &["subtotal"], currency)?; - validate_money(economics, &["discount_total"], currency)?; - validate_money(economics, &["adjustment_total"], currency)?; - validate_money(economics, &["total"], currency)?; - Ok(()) -} - -fn validate_money(payload: &Value, path: &[&str], currency: &str) -> Result<(), LocalEventsError> { - let Some(money) = value_at(payload, path) else { - return Err(missing_field(path)); - }; - validate_required_string(money, &["amount"])?; - let money_currency = validate_required_string(money, &["currency"])?; - if money_currency != currency { - return Err(invalid_field( - &format!("{}.currency", path.join(".")), - "must match currency", - )); - } - Ok(()) -} - -fn validate_string_field( - payload: &Value, - path: &[&str], - expected: &str, -) -> Result<(), LocalEventsError> { - let Some(value) = value_at(payload, path).and_then(Value::as_str) else { - return Err(missing_field(path)); - }; - if value != expected { - return Err(invalid_field( - &path.join("."), - &format!("must be `{expected}`"), - )); - } - Ok(()) -} - -fn validate_required_string<'a>( - payload: &'a Value, - path: &[&str], -) -> Result<&'a str, LocalEventsError> { - let Some(value) = value_at(payload, path).and_then(Value::as_str) else { - return Err(missing_field(path)); - }; - validate_non_empty(&path.join("."), value)?; - Ok(value.trim()) -} - -fn validate_bool_field( - payload: &Value, - path: &[&str], - expected: bool, -) -> Result<(), LocalEventsError> { - let Some(value) = value_at(payload, path).and_then(Value::as_bool) else { - return Err(missing_field(path)); - }; - if value != expected { - return Err(invalid_field( - &path.join("."), - &format!("must be `{expected}`"), - )); - } - Ok(()) -} - -fn validate_positive_i64(payload: &Value, path: &[&str]) -> Result<(), LocalEventsError> { - match value_at(payload, path).and_then(Value::as_i64) { - Some(value) if value > 0 => Ok(()), - _ => Err(invalid_field(&path.join("."), "must be positive")), - } -} - -fn validate_positive_u64(payload: &Value, path: &[&str]) -> Result<u64, LocalEventsError> { - match value_at(payload, path).and_then(Value::as_u64) { - Some(value) if value > 0 => Ok(value), - _ => Err(invalid_field(&path.join("."), "must be positive")), - } -} - -fn validate_currency(field: &str, value: &str) -> Result<(), LocalEventsError> { - if value.len() != 3 || !value.bytes().all(|byte| byte.is_ascii_uppercase()) { - return Err(invalid_field( - field, - "must be an uppercase ISO currency code", - )); - } - Ok(()) -} - -fn required_array<'a>( - payload: &'a Value, - path: &[&str], -) -> Result<&'a Vec<Value>, LocalEventsError> { - let Some(value) = value_at(payload, path).and_then(Value::as_array) else { - return Err(missing_field(path)); - }; - Ok(value) -} - -fn support_issues(payload: &Value) -> Result<Vec<String>, LocalEventsError> { - let issues = required_array(payload, &["support_status", "issues"])?; - let mut parsed = Vec::with_capacity(issues.len()); - for (index, issue) in issues.iter().enumerate() { - let Some(issue) = issue.as_str() else { - return Err(invalid_field_at( - format!("support_status.issues[{index}]"), - "must be a string", - )); - }; - validate_non_empty("support_status.issues", issue)?; - parsed.push(issue.trim().to_owned()); - } - Ok(parsed) -} - -fn value_at<'a>(payload: &'a Value, path: &[&str]) -> Option<&'a Value> { - let mut current = payload; - for part in path { - current = current.get(*part)?; - } - Some(current) -} - -fn missing_field(path: &[&str]) -> LocalEventsError { - invalid_field(&path.join("."), "is required") -} - -fn invalid_field(field: &str, requirement: &str) -> LocalEventsError { - LocalEventsError::InvalidRecord(format!("local order field `{field}` {requirement}")) -} - -fn invalid_field_at(field: String, requirement: &str) -> LocalEventsError { - LocalEventsError::InvalidRecord(format!("local order field `{field}` {requirement}")) -} - -#[cfg(test)] -mod tests { - use serde_json::{Value, json}; - - use super::*; - - #[test] - fn support_state_labels_and_record_id_validation_are_stable() { - assert_eq!( - BuyerOrderRequestSupportState::Supported.as_str(), - "supported" - ); - assert_eq!( - BuyerOrderRequestSupportState::Unsupported.as_str(), - "unsupported" - ); - assert_eq!( - buyer_order_request_local_work_record_id(" ord-a ").expect("record id"), - "app:local_work:order_request:ord-a" - ); - assert_error_contains( - buyer_order_request_local_work_record_id(" "), - "order_id must not be empty", - ); - } - - #[test] - fn private_validation_helpers_cover_successful_payload() { - let payload = supported_payload(); - - assert_eq!( - validate_support_status(&payload).expect("support status"), - ( - BuyerOrderRequestSupportState::Supported, - Vec::<String>::new() - ) - ); - validate_supported_buyer_order_request_local_work_payload(&payload) - .expect("supported payload"); - validate_exportability(&payload, BuyerOrderRequestSupportState::Supported) - .expect("exportability"); - validate_order_identity(&payload, BuyerOrderRequestSupportState::Supported) - .expect("identity"); - validate_order_items(&payload).expect("items"); - validate_order_economics(&payload).expect("economics"); - assert_eq!( - validate_required_string(&payload, &["document", "order", "order_id"]) - .expect("order id"), - "ord_1" - ); - validate_bool_field(&payload, &["currentness", "current"], true).expect("bool"); - assert_eq!( - support_issues(&payload).expect("support issues"), - Vec::<String>::new() - ); - assert!(value_at(&payload, &["document", "order"]).is_some()); - } - - #[test] - fn payload_validation_rejects_top_level_contract_drift() { - let mut wrong_kind = supported_payload(); - wrong_kind["record_kind"] = json!("other"); - assert_invalid(wrong_kind, "record_kind"); - - let mut missing_scope = supported_payload(); - missing_scope["scope"] = Value::Null; - assert_invalid(missing_scope, "scope"); - - let mut wrong_document_kind = supported_payload(); - wrong_document_kind["document"]["kind"] = json!("other"); - assert_invalid(wrong_document_kind, "document.kind"); - - let mut wrong_currentness_source = supported_payload(); - wrong_currentness_source["currentness"]["source"] = json!("other"); - assert_invalid(wrong_currentness_source, "currentness.source"); - - let mut missing_order_updated = supported_payload(); - missing_order_updated["currentness"]["order_updated_at"] = Value::Null; - assert_invalid(missing_order_updated, "order_updated_at"); - - let mut bad_created_at = supported_payload(); - bad_created_at["currentness"]["created_at_ms"] = json!(0); - assert_invalid(bad_created_at, "created_at_ms"); - } - - #[test] - fn support_and_exportability_rejections_cover_private_branches() { - let mut invalid_state = supported_payload(); - invalid_state["support_status"]["state"] = json!("partial"); - assert_invalid(invalid_state, "support_status.state"); - - let mut issue_not_string = supported_payload(); - issue_not_string["support_status"] = json!({ - "state": "unsupported", - "issues": [42] - }); - assert_invalid(issue_not_string, "support_status.issues[0]"); - - let mut issue_empty = supported_payload(); - issue_empty["support_status"] = json!({ - "state": "unsupported", - "issues": [" "] - }); - assert_invalid(issue_empty, "support_status.issues"); - - let mut supported_but_unresolved = unsupported_payload(); - supported_but_unresolved["support_status"] = json!({ - "state": "supported", - "issues": [] - }); - assert_invalid(supported_but_unresolved, "exportability.state"); - - let mut unknown_exportability = supported_payload(); - unknown_exportability["exportability"]["state"] = json!("queued"); - assert_invalid(unknown_exportability, "exportability.state"); - - let mut missing_reason = unsupported_payload(); - missing_reason["exportability"]["reason"] = Value::Null; - assert_invalid(missing_reason, "exportability.reason"); - - let mut wrong_actor_source = unsupported_payload(); - wrong_actor_source["document"]["buyer_actor"]["source"] = - json!(BUYER_ORDER_REQUEST_ACTOR_SOURCE_RESOLVED_ACCOUNT); - assert_invalid(wrong_actor_source, "buyer_actor.source"); - - let mut mismatched_buyer = supported_payload(); - mismatched_buyer["document"]["buyer_actor"]["pubkey"] = json!("other"); - assert_invalid(mismatched_buyer, "buyer_actor.pubkey"); - - let supported_error = - validate_unsupported_buyer_order_request_local_work_payload(&supported_payload()) - .expect_err("supported payload is not unsupported"); - assert!(supported_error.to_string().contains("support_status.state")); - } - - #[test] - fn item_and_economics_rejections_cover_private_branches() { - let mut economics_not_object = supported_payload(); - economics_not_object["document"]["order"]["economics"] = json!("bad"); - assert_invalid(economics_not_object, "economics"); - - let mut bad_pricing_basis = supported_payload(); - bad_pricing_basis["document"]["order"]["economics"]["pricing_basis"] = json!("manual"); - assert_invalid(bad_pricing_basis, "pricing_basis"); - - let mut bad_currency = supported_payload(); - bad_currency["document"]["order"]["economics"]["currency"] = json!("usd"); - assert_invalid(bad_currency, "currency"); - - let mut bad_currency_length = supported_payload(); - bad_currency_length["document"]["order"]["economics"]["currency"] = json!("US"); - assert_invalid(bad_currency_length, "currency"); - - let mut missing_economics = supported_payload(); - missing_economics["document"]["order"] - .as_object_mut() - .expect("order object") - .remove("economics"); - assert_invalid(missing_economics, "economics"); - - let mut economics_items_missing = supported_payload(); - economics_items_missing["document"]["order"]["economics"]["items"] = Value::Null; - assert_invalid(economics_items_missing, "items"); - - let mut economics_items_short = supported_payload(); - economics_items_short["document"]["order"]["economics"]["items"] = json!([]); - assert_invalid(economics_items_short, "economics.items"); - - let mut economics_items_long = supported_payload(); - economics_items_long["document"]["order"]["economics"]["items"] = json!([ - { - "bin_id": "dozen-eggs", - "bin_count": 2, - "quantity_amount": "1", - "quantity_unit": "dozen", - "unit_price_amount": "8.00", - "unit_price_currency": "USD", - "line_subtotal": { - "amount": "16.00", - "currency": "USD" - } - }, - { - "bin_id": "half-dozen-eggs", - "bin_count": 1 - } - ]); - assert_invalid(economics_items_long, "economics.items"); - - let mut economics_bin_missing = supported_payload(); - economics_bin_missing["document"]["order"]["economics"]["items"][0]["bin_id"] = Value::Null; - assert_invalid(economics_bin_missing, "economics.items[0].bin_id"); - - let mut economics_count_bad = supported_payload(); - economics_count_bad["document"]["order"]["economics"]["items"][0]["bin_count"] = json!(0); - assert_invalid(economics_count_bad, "economics.items[0].bin_count"); - - let mut order_count_mismatch = supported_payload(); - order_count_mismatch["document"]["order"]["economics"]["items"][0]["bin_count"] = json!(3); - assert_invalid(order_count_mismatch, "economics.items[0].bin_count"); - - let mut quantity_amount_missing = supported_payload(); - quantity_amount_missing["document"]["order"]["economics"]["items"][0]["quantity_amount"] = - Value::Null; - assert_invalid(quantity_amount_missing, "quantity_amount"); - - let mut quantity_unit_missing = supported_payload(); - quantity_unit_missing["document"]["order"]["economics"]["items"][0]["quantity_unit"] = - Value::Null; - assert_invalid(quantity_unit_missing, "quantity_unit"); - - let mut unit_price_amount_missing = supported_payload(); - unit_price_amount_missing["document"]["order"]["economics"]["items"][0]["unit_price_amount"] = - Value::Null; - assert_invalid(unit_price_amount_missing, "unit_price_amount"); - - let mut line_subtotal_missing = supported_payload(); - line_subtotal_missing["document"]["order"]["economics"]["items"][0]["line_subtotal"] = - Value::Null; - assert_invalid(line_subtotal_missing, "amount"); - - let mut missing_line_subtotal = supported_payload(); - missing_line_subtotal["document"]["order"]["economics"]["items"][0] - .as_object_mut() - .expect("economics item") - .remove("line_subtotal"); - assert_invalid(missing_line_subtotal, "line_subtotal"); - - let mut line_subtotal_currency = supported_payload(); - line_subtotal_currency["document"]["order"]["economics"]["items"][0]["line_subtotal"]["currency"] = - json!("CAD"); - assert_invalid(line_subtotal_currency, "line_subtotal.currency"); - - let mut subtotal_currency = supported_payload(); - subtotal_currency["document"]["order"]["economics"]["subtotal"]["currency"] = json!("CAD"); - assert_invalid(subtotal_currency, "subtotal.currency"); - - let mut order_item_missing = supported_payload(); - order_item_missing["document"]["order"]["items"] = Value::Null; - assert_invalid(order_item_missing, "document.order.items"); - - let mut missing_order_bin = supported_payload(); - missing_order_bin["document"]["order"]["items"][0]["bin_id"] = Value::Null; - assert_error_contains(validate_order_items(&missing_order_bin), "items[0].bin_id"); - } - - fn supported_payload() -> Value { - json!({ - "record_kind": BUYER_ORDER_REQUEST_LOCAL_WORK_RECORD_KIND, - "scope": "app", - "exportability": { - "state": "exportable" - }, - "support_status": { - "state": "supported", - "issues": [] - }, - "currentness": { - "current": true, - "source": "app_sqlite_order", - "record_id": "app:local_work:order_request:ord_1", - "order_id": "ord_1", - "order_updated_at": "2026-05-24T12:00:00Z", - "created_at_ms": 1777777777000_i64 - }, - "document": { - "kind": BUYER_ORDER_REQUEST_DOCUMENT_KIND, - "order": { - "order_id": "ord_1", - "listing_addr": "30402:seller_pubkey:listing_key", - "listing_event_id": "event-listing-1", - "buyer_pubkey": "buyer_pubkey", - "seller_pubkey": "seller_pubkey", - "items": [ - { - "bin_id": "dozen-eggs", - "bin_count": 2 - } - ], - "economics": { - "pricing_basis": "listing_event", - "currency": "USD", - "items": [ - { - "bin_id": "dozen-eggs", - "bin_count": 2, - "quantity_amount": "1", - "quantity_unit": "dozen", - "unit_price_amount": "8.00", - "unit_price_currency": "USD", - "line_subtotal": { - "amount": "16.00", - "currency": "USD" - } - } - ], - "subtotal": { - "amount": "16.00", - "currency": "USD" - }, - "discount_total": { - "amount": "0", - "currency": "USD" - }, - "adjustment_total": { - "amount": "0", - "currency": "USD" - }, - "total": { - "amount": "16.00", - "currency": "USD" - } - } - }, - "buyer_actor": { - "account_id": "buyer-account", - "pubkey": "buyer_pubkey", - "source": BUYER_ORDER_REQUEST_ACTOR_SOURCE_RESOLVED_ACCOUNT - } - } - }) - } - - fn unsupported_payload() -> Value { - let mut payload = supported_payload(); - payload["exportability"] = json!({ - "state": "identity_unresolved", - "reason": "canonical_hex_pubkey_required" - }); - payload["support_status"] = json!({ - "state": "unsupported", - "issues": ["buyer_pubkey_required"] - }); - payload["document"]["order"]["buyer_pubkey"] = json!(""); - payload["document"]["buyer_actor"]["pubkey"] = json!(""); - payload["document"]["buyer_actor"]["source"] = - json!(BUYER_ORDER_REQUEST_ACTOR_SOURCE_UNRESOLVED_APP); - payload - } - - fn assert_invalid(payload: Value, expected: &str) { - assert_error_contains( - validate_buyer_order_request_local_work_payload(&payload), - expected, - ); - } - - fn assert_error_contains<T: std::fmt::Debug>( - result: Result<T, LocalEventsError>, - expected: &str, - ) { - let error = result.expect_err("expected validation error"); - assert!( - error.to_string().contains(expected), - "expected error to contain {expected}, got {error}" - ); - } -} diff --git a/crates/local_events/src/store.rs b/crates/local_events/src/store.rs @@ -1,933 +0,0 @@ -#![forbid(unsafe_code)] - -use radroots_sql_core::SqlExecutor; -use radroots_sql_core::error::SqlError; -use serde::Deserialize; -use serde_json::{Value, json}; - -use crate::migrations; -use crate::models::validate_non_empty; -use crate::{ - LocalEventRecord, LocalEventRecordInput, LocalEventRecordUpdate, LocalEventsCursor, - LocalEventsError, LocalRecordFamily, LocalRecordStatus, PublishOutboxStatus, SourceRuntime, -}; - -pub struct LocalEventsStore<E: SqlExecutor> { - executor: E, -} - -impl<E: SqlExecutor> LocalEventsStore<E> { - pub fn new(executor: E) -> Self { - Self { executor } - } - - pub fn executor(&self) -> &E { - &self.executor - } - - pub fn migrate_up(&self) -> Result<(), SqlError> { - migrations::run_all_up(self.executor()) - } - - pub fn migrate_down(&self) -> Result<(), SqlError> { - migrations::run_all_down(self.executor()) - } - - pub fn append_record( - &self, - input: &LocalEventRecordInput, - ) -> Result<LocalEventRecord, LocalEventsError> { - input.validate()?; - self.executor.begin()?; - let result = (|| -> Result<(), LocalEventsError> { - let change_seq = self.next_change_seq()?; - let params = json!([ - change_seq, - input.record_id, - input.family.as_str(), - input.status.as_str(), - input.source_runtime.as_str(), - input.created_at_ms, - input.inserted_at_ms, - input.inserted_at_ms, - input.owner_account_id, - input.owner_pubkey, - input.farm_id, - input.listing_addr, - encode_json(input.local_work_json.as_ref()), - input.event_id, - input.event_kind, - input.event_pubkey, - input.event_created_at, - encode_json(input.event_tags_json.as_ref()), - input.event_content, - input.event_sig, - encode_json(input.raw_event_json.as_ref()), - input.outbox_status.as_str() - ]) - .to_string(); - let sql = "insert or ignore into local_event_record( - change_seq, - record_id, - family, - status, - source_runtime, - created_at_ms, - inserted_at_ms, - updated_at_ms, - owner_account_id, - owner_pubkey, - farm_id, - listing_addr, - local_work_json, - event_id, - event_kind, - event_pubkey, - event_created_at, - event_tags_json, - event_content, - event_sig, - raw_event_json, - outbox_status - ) values(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)"; - let _ = self.executor.exec(sql, &params)?; - Ok(()) - })(); - match result { - Ok(()) => self.executor.commit()?, - Err(err) => { - let _ = self.executor.rollback(); - return Err(err); - } - } - self.get_record(&input.record_id)? - .ok_or_else(|| LocalEventsError::InvalidRecord("record append failed".to_owned())) - } - - pub fn get_record( - &self, - record_id: &str, - ) -> Result<Option<LocalEventRecord>, LocalEventsError> { - validate_non_empty("record_id", record_id)?; - let params = json!([record_id]).to_string(); - let rows = self.query_records( - "select * from local_event_record where record_id = ? limit 1", - &params, - )?; - Ok(rows.into_iter().next()) - } - - pub fn list_records_after_seq( - &self, - after_seq: i64, - limit: u32, - ) -> Result<Vec<LocalEventRecord>, LocalEventsError> { - let params = json!([after_seq, i64::from(limit)]).to_string(); - self.query_records( - "select * from local_event_record where seq > ? order by seq asc limit ?", - &params, - ) - } - - pub fn list_records_changed_after( - &self, - after_change_seq: i64, - limit: u32, - ) -> Result<Vec<LocalEventRecord>, LocalEventsError> { - let params = json!([after_change_seq, i64::from(limit)]).to_string(); - self.query_records( - "select * from local_event_record where change_seq > ? order by change_seq asc, seq asc limit ?", - &params, - ) - } - - pub fn list_records_changed_latest( - &self, - limit: u32, - ) -> Result<Vec<LocalEventRecord>, LocalEventsError> { - let params = json!([i64::from(limit)]).to_string(); - self.query_records( - "select * from local_event_record order by change_seq desc, seq desc, record_id asc limit ?", - &params, - ) - } - - pub fn list_records_changed_before( - &self, - before_change_seq: i64, - before_seq: i64, - limit: u32, - ) -> Result<Vec<LocalEventRecord>, LocalEventsError> { - let params = json!([ - before_change_seq, - before_change_seq, - before_seq, - i64::from(limit) - ]) - .to_string(); - self.query_records( - "select * from local_event_record - where change_seq < ? or (change_seq = ? and seq < ?) - order by change_seq desc, seq desc, record_id asc - limit ?", - &params, - ) - } - - pub fn update_outbox( - &self, - update: &LocalEventRecordUpdate, - ) -> Result<LocalEventRecord, LocalEventsError> { - validate_non_empty("record_id", &update.record_id)?; - self.executor.begin()?; - let result = (|| -> Result<i64, LocalEventsError> { - let change_seq = self.next_change_seq()?; - let params = json!([ - change_seq, - update.status.as_str(), - update.outbox_status.as_str(), - update.updated_at_ms, - update.record_id - ]) - .to_string(); - let outcome = self.executor.exec( - "update local_event_record - set change_seq = ?, - status = ?, - outbox_status = ?, - updated_at_ms = ? - where record_id = ?", - &params, - )?; - Ok(outcome.changes) - })(); - let changes = match result { - Ok(changes) => { - self.executor.commit()?; - changes - } - Err(err) => { - let _ = self.executor.rollback(); - return Err(err); - } - }; - if changes == 0 { - return Err(LocalEventsError::Sql(SqlError::NotFound( - update.record_id.clone(), - ))); - } - self.get_record(&update.record_id)? - .ok_or_else(|| LocalEventsError::Sql(SqlError::NotFound(update.record_id.clone()))) - } - - pub fn get_cursor( - &self, - consumer_id: &str, - ) -> Result<Option<LocalEventsCursor>, LocalEventsError> { - validate_non_empty("consumer_id", consumer_id)?; - let params = json!([consumer_id]).to_string(); - let raw = self.executor.query_raw( - "select consumer_id, last_change_seq, updated_at_ms from local_event_projection_cursor where consumer_id = ? limit 1", - &params, - )?; - let rows: Vec<CursorRow> = serde_json::from_str(&raw)?; - Ok(rows.into_iter().next().map(Into::into)) - } - - pub fn advance_cursor( - &self, - consumer_id: &str, - last_change_seq: i64, - updated_at_ms: i64, - ) -> Result<LocalEventsCursor, LocalEventsError> { - validate_non_empty("consumer_id", consumer_id)?; - let params = json!([consumer_id, last_change_seq, updated_at_ms]).to_string(); - self.executor.exec( - "insert into local_event_projection_cursor(consumer_id, last_change_seq, updated_at_ms) - values(?,?,?) - on conflict(consumer_id) do update set - last_change_seq = max(local_event_projection_cursor.last_change_seq, excluded.last_change_seq), - updated_at_ms = excluded.updated_at_ms", - &params, - )?; - self.get_cursor(consumer_id)? - .ok_or_else(|| LocalEventsError::InvalidRecord("cursor advance failed".to_owned())) - } - - fn query_records( - &self, - sql: &str, - params: &str, - ) -> Result<Vec<LocalEventRecord>, LocalEventsError> { - let raw = self.executor.query_raw(sql, params)?; - let rows: Vec<RecordRow> = serde_json::from_str(&raw)?; - rows.into_iter().map(TryInto::try_into).collect() - } - - fn next_change_seq(&self) -> Result<i64, LocalEventsError> { - let raw = self.executor.query_raw( - "select coalesce(max(change_seq), 0) + 1 as change_seq from local_event_record", - "[]", - )?; - let rows: Vec<ChangeSeqRow> = serde_json::from_str(&raw)?; - rows.into_iter() - .next() - .map(|row| row.change_seq) - .ok_or_else(|| { - LocalEventsError::InvalidRecord("change sequence unavailable".to_owned()) - }) - } -} - -#[derive(Debug, Deserialize)] -struct RecordRow { - seq: i64, - change_seq: i64, - record_id: String, - family: String, - status: String, - source_runtime: String, - created_at_ms: i64, - inserted_at_ms: i64, - updated_at_ms: i64, - owner_account_id: Option<String>, - owner_pubkey: Option<String>, - farm_id: Option<String>, - listing_addr: Option<String>, - local_work_json: Option<String>, - event_id: Option<String>, - event_kind: Option<i64>, - event_pubkey: Option<String>, - event_created_at: Option<i64>, - event_tags_json: Option<String>, - event_content: Option<String>, - event_sig: Option<String>, - raw_event_json: Option<String>, - outbox_status: String, -} - -impl TryFrom<RecordRow> for LocalEventRecord { - type Error = LocalEventsError; - - fn try_from(row: RecordRow) -> Result<Self, Self::Error> { - Ok(Self { - seq: row.seq, - change_seq: row.change_seq, - record_id: row.record_id, - family: LocalRecordFamily::parse(&row.family)?, - status: LocalRecordStatus::parse(&row.status)?, - source_runtime: SourceRuntime::parse(&row.source_runtime)?, - created_at_ms: row.created_at_ms, - inserted_at_ms: row.inserted_at_ms, - updated_at_ms: row.updated_at_ms, - owner_account_id: row.owner_account_id, - owner_pubkey: row.owner_pubkey, - farm_id: row.farm_id, - listing_addr: row.listing_addr, - local_work_json: decode_json(row.local_work_json)?, - event_id: row.event_id, - event_kind: row.event_kind, - event_pubkey: row.event_pubkey, - event_created_at: row.event_created_at, - event_tags_json: decode_json(row.event_tags_json)?, - event_content: row.event_content, - event_sig: row.event_sig, - raw_event_json: decode_json(row.raw_event_json)?, - outbox_status: PublishOutboxStatus::parse(&row.outbox_status)?, - }) - } -} - -#[derive(Debug, Deserialize)] -struct CursorRow { - consumer_id: String, - last_change_seq: i64, - updated_at_ms: i64, -} - -impl From<CursorRow> for LocalEventsCursor { - fn from(row: CursorRow) -> Self { - Self { - consumer_id: row.consumer_id, - last_change_seq: row.last_change_seq, - updated_at_ms: row.updated_at_ms, - } - } -} - -#[derive(Debug, Deserialize)] -struct ChangeSeqRow { - change_seq: i64, -} - -fn encode_json(value: Option<&Value>) -> Option<String> { - value.map(Value::to_string) -} - -fn decode_json(value: Option<String>) -> Result<Option<Value>, LocalEventsError> { - value - .map(|value| serde_json::from_str(&value)) - .transpose() - .map_err(Into::into) -} - -#[cfg(test)] -mod tests { - use std::collections::VecDeque; - use std::sync::Mutex; - use std::sync::atomic::{AtomicUsize, Ordering}; - - use radroots_sql_core::{ExecOutcome, SqlExecutor, SqliteExecutor}; - use serde_json::json; - - use super::*; - - fn store() -> LocalEventsStore<SqliteExecutor> { - let executor = SqliteExecutor::open_memory().expect("open memory sqlite"); - let store = LocalEventsStore::new(executor); - store.migrate_up().expect("migrate up"); - store - } - - fn local_work(record_id: &str) -> LocalEventRecordInput { - LocalEventRecordInput { - record_id: record_id.to_owned(), - family: LocalRecordFamily::LocalWork, - status: LocalRecordStatus::LocalSaved, - source_runtime: SourceRuntime::Cli, - created_at_ms: 1000, - inserted_at_ms: 1001, - owner_account_id: Some("seller-account".to_owned()), - owner_pubkey: Some("seller-pubkey".to_owned()), - farm_id: Some("farm-a".to_owned()), - listing_addr: Some("listing-a".to_owned()), - local_work_json: Some(json!({"kind":"listing","title":"Eggs"})), - event_id: None, - event_kind: None, - event_pubkey: None, - event_created_at: None, - event_tags_json: None, - event_content: None, - event_sig: None, - raw_event_json: None, - outbox_status: PublishOutboxStatus::None, - } - } - - fn signed_event(record_id: &str) -> LocalEventRecordInput { - LocalEventRecordInput { - record_id: record_id.to_owned(), - family: LocalRecordFamily::SignedEvent, - status: LocalRecordStatus::PendingPublish, - source_runtime: SourceRuntime::Cli, - created_at_ms: 2000, - inserted_at_ms: 2001, - owner_account_id: Some("seller-account".to_owned()), - owner_pubkey: Some("seller-pubkey".to_owned()), - farm_id: Some("farm-a".to_owned()), - listing_addr: Some("listing-a".to_owned()), - local_work_json: None, - event_id: Some(record_id.to_owned()), - event_kind: Some(3421), - event_pubkey: Some("seller-pubkey".to_owned()), - event_created_at: Some(2000), - event_tags_json: Some(json!([["d", "listing-a"]])), - event_content: Some("{\"title\":\"Eggs\"}".to_owned()), - event_sig: Some("sig-a".to_owned()), - raw_event_json: Some(json!({"id":record_id,"kind":3421})), - outbox_status: PublishOutboxStatus::Pending, - } - } - - #[derive(Debug)] - struct ScriptedExecutor { - begin_result: Mutex<Result<(), SqlError>>, - commit_result: Mutex<Result<(), SqlError>>, - exec_results: Mutex<VecDeque<Result<ExecOutcome, SqlError>>>, - query_results: Mutex<VecDeque<Result<String, SqlError>>>, - rollbacks: AtomicUsize, - } - - impl ScriptedExecutor { - fn new( - exec_results: Vec<Result<ExecOutcome, SqlError>>, - query_results: Vec<Result<String, SqlError>>, - ) -> Self { - Self { - begin_result: Mutex::new(Ok(())), - commit_result: Mutex::new(Ok(())), - exec_results: Mutex::new(exec_results.into()), - query_results: Mutex::new(query_results.into()), - rollbacks: AtomicUsize::new(0), - } - } - - fn with_begin_error(error: SqlError) -> Self { - let executor = Self::new(Vec::new(), Vec::new()); - *executor.begin_result.lock().expect("begin result") = Err(error); - executor - } - - fn with_commit_error(error: SqlError) -> Self { - let executor = Self::new( - vec![Ok(ExecOutcome { - changes: 1, - last_insert_id: 0, - })], - vec![Ok(r#"[{"change_seq":1}]"#.to_owned())], - ); - *executor.commit_result.lock().expect("commit result") = Err(error); - executor - } - } - - impl SqlExecutor for ScriptedExecutor { - fn exec(&self, _sql: &str, _params_json: &str) -> Result<ExecOutcome, SqlError> { - self.exec_results - .lock() - .expect("exec results") - .pop_front() - .unwrap_or(Ok(ExecOutcome { - changes: 1, - last_insert_id: 0, - })) - } - - fn query_raw(&self, _sql: &str, _params_json: &str) -> Result<String, SqlError> { - self.query_results - .lock() - .expect("query results") - .pop_front() - .unwrap_or_else(|| Ok("[]".to_owned())) - } - - fn begin(&self) -> Result<(), SqlError> { - self.begin_result.lock().expect("begin result").clone() - } - - fn commit(&self) -> Result<(), SqlError> { - self.commit_result.lock().expect("commit result").clone() - } - - fn rollback(&self) -> Result<(), SqlError> { - self.rollbacks.fetch_add(1, Ordering::SeqCst); - Ok(()) - } - } - - fn record_row_with(field: &str, value: serde_json::Value) -> String { - let mut row = json!({ - "seq": 1, - "change_seq": 1, - "record_id": "record-a", - "family": "signed_event", - "status": "pending_publish", - "source_runtime": "cli", - "created_at_ms": 1000, - "inserted_at_ms": 1001, - "updated_at_ms": 1001, - "owner_account_id": "seller-account", - "owner_pubkey": "seller-pubkey", - "farm_id": "farm-a", - "listing_addr": "listing-a", - "local_work_json": null, - "event_id": "event-a", - "event_kind": 3421, - "event_pubkey": "seller-pubkey", - "event_created_at": 1000, - "event_tags_json": "[[\"d\",\"listing-a\"]]", - "event_content": "{}", - "event_sig": "sig-a", - "raw_event_json": "{\"id\":\"event-a\",\"kind\":3421}", - "outbox_status": "pending" - }); - row[field] = value; - json!([row]).to_string() - } - - #[test] - fn store_methods_round_trip_records_and_cursors() { - let store = store(); - - assert!( - store - .executor() - .query_raw("select 1 as value", "[]") - .is_ok() - ); - assert!(store.get_record("missing").expect("get missing").is_none()); - assert!(store.get_cursor("app").expect("cursor missing").is_none()); - - let local = store - .append_record(&local_work("local-a")) - .expect("append local work"); - let event = store - .append_record(&signed_event("event-a")) - .expect("append signed event"); - - assert_eq!( - store - .get_record("local-a") - .expect("get local") - .expect("local record") - .record_id, - local.record_id - ); - assert_eq!( - store - .list_records_after_seq(0, 10) - .expect("list after seq") - .len(), - 2 - ); - assert_eq!( - store - .list_records_changed_after(local.change_seq, 10) - .expect("list changed after")[0] - .record_id, - event.record_id - ); - assert_eq!( - store.list_records_changed_latest(1).expect("list latest")[0].record_id, - event.record_id - ); - assert_eq!( - store - .list_records_changed_before(event.change_seq, event.seq, 10) - .expect("list before")[0] - .record_id, - local.record_id - ); - - let cursor = store - .advance_cursor("app", event.change_seq, 3000) - .expect("advance cursor"); - assert_eq!(cursor.consumer_id, "app"); - assert_eq!( - store - .get_cursor("app") - .expect("get cursor") - .expect("cursor") - .last_change_seq, - event.change_seq - ); - - let updated = store - .update_outbox(&LocalEventRecordUpdate { - record_id: "event-a".to_owned(), - status: LocalRecordStatus::Published, - outbox_status: PublishOutboxStatus::Acknowledged, - updated_at_ms: 4000, - }) - .expect("update outbox"); - - assert_eq!(updated.status, LocalRecordStatus::Published); - assert_eq!(updated.outbox_status, PublishOutboxStatus::Acknowledged); - store.migrate_down().expect("migrate down"); - } - - #[test] - fn store_reports_missing_updates_and_decode_errors() { - let store = store(); - assert!( - store - .get_record(" ") - .expect_err("empty record id") - .to_string() - .contains("record_id") - ); - assert!( - store - .get_cursor(" ") - .expect_err("empty consumer id") - .to_string() - .contains("consumer_id") - ); - assert!( - store - .advance_cursor(" ", 1, 1000) - .expect_err("empty cursor consumer") - .to_string() - .contains("consumer_id") - ); - assert!( - store - .update_outbox(&LocalEventRecordUpdate { - record_id: " ".to_owned(), - status: LocalRecordStatus::Published, - outbox_status: PublishOutboxStatus::Acknowledged, - updated_at_ms: 4000, - }) - .expect_err("empty update record id") - .to_string() - .contains("record_id") - ); - - let missing_update = store - .update_outbox(&LocalEventRecordUpdate { - record_id: "missing-event".to_owned(), - status: LocalRecordStatus::Published, - outbox_status: PublishOutboxStatus::Acknowledged, - updated_at_ms: 4000, - }) - .expect_err("missing record update"); - - assert!(missing_update.to_string().contains("missing-event")); - - store - .append_record(&local_work("local-a")) - .expect("append local"); - let params = json!(["{", "local-a"]).to_string(); - store - .executor() - .exec( - "update local_event_record set local_work_json = ? where record_id = ?", - &params, - ) - .expect("corrupt local work json"); - let decode_error = store.get_record("local-a").expect_err("decode error"); - - assert!(decode_error.to_string().contains("EOF")); - } - - #[test] - fn store_rolls_back_when_change_sequence_is_unavailable() { - let append_store = - LocalEventsStore::new(ScriptedExecutor::new(Vec::new(), vec![Ok("[]".to_owned())])); - let append_error = append_store - .append_record(&local_work("local-a")) - .expect_err("append error"); - - assert!(append_error.to_string().contains("change sequence")); - assert_eq!(append_store.executor().rollbacks.load(Ordering::SeqCst), 1); - - let update_store = - LocalEventsStore::new(ScriptedExecutor::new(Vec::new(), vec![Ok("[]".to_owned())])); - let update_error = update_store - .update_outbox(&LocalEventRecordUpdate { - record_id: "event-a".to_owned(), - status: LocalRecordStatus::Published, - outbox_status: PublishOutboxStatus::Acknowledged, - updated_at_ms: 4000, - }) - .expect_err("update error"); - - assert!(update_error.to_string().contains("change sequence")); - assert_eq!(update_store.executor().rollbacks.load(Ordering::SeqCst), 1); - } - - #[test] - fn store_reports_cursor_advance_without_returned_cursor() { - let store = LocalEventsStore::new(ScriptedExecutor::new(Vec::new(), Vec::new())); - - assert!(store.get_cursor("app").expect("missing cursor").is_none()); - let cursor_error = store - .advance_cursor("app", 1, 1000) - .expect_err("cursor advance error"); - - assert!(cursor_error.to_string().contains("cursor advance failed")); - } - - #[test] - fn store_reports_executor_and_decode_failures() { - let begin_store = LocalEventsStore::new(ScriptedExecutor::with_begin_error( - SqlError::InvalidQuery("begin failed".to_owned()), - )); - assert!( - begin_store - .append_record(&local_work("local-a")) - .expect_err("begin failure") - .to_string() - .contains("begin failed") - ); - - let exec_store = LocalEventsStore::new(ScriptedExecutor::new( - vec![Err(SqlError::InvalidQuery("insert failed".to_owned()))], - vec![Ok(r#"[{"change_seq":1}]"#.to_owned())], - )); - assert!( - exec_store - .append_record(&local_work("local-a")) - .expect_err("exec failure") - .to_string() - .contains("insert failed") - ); - assert_eq!(exec_store.executor().rollbacks.load(Ordering::SeqCst), 1); - - let commit_store = LocalEventsStore::new(ScriptedExecutor::with_commit_error( - SqlError::InvalidQuery("commit failed".to_owned()), - )); - assert!( - commit_store - .append_record(&local_work("local-a")) - .expect_err("commit failure") - .to_string() - .contains("commit failed") - ); - - let query_error_store = LocalEventsStore::new(ScriptedExecutor::new( - Vec::new(), - vec![Err(SqlError::InvalidQuery("query failed".to_owned()))], - )); - assert!( - query_error_store - .get_record("record-a") - .expect_err("query failure") - .to_string() - .contains("query failed") - ); - - let invalid_rows_store = - LocalEventsStore::new(ScriptedExecutor::new(Vec::new(), vec![Ok("{".to_owned())])); - let _ = invalid_rows_store - .get_record("record-a") - .expect_err("invalid rows"); - - let cursor_rows_store = - LocalEventsStore::new(ScriptedExecutor::new(Vec::new(), vec![Ok("{".to_owned())])); - let _ = cursor_rows_store - .get_cursor("app") - .expect_err("invalid cursor rows"); - - let change_rows_store = - LocalEventsStore::new(ScriptedExecutor::new(Vec::new(), vec![Ok("{".to_owned())])); - let _ = change_rows_store - .append_record(&local_work("local-a")) - .expect_err("invalid change rows"); - - let cursor_exec_store = LocalEventsStore::new(ScriptedExecutor::new( - vec![Err(SqlError::InvalidQuery("cursor failed".to_owned()))], - Vec::new(), - )); - assert!( - cursor_exec_store - .advance_cursor("app", 1, 1000) - .expect_err("cursor exec failure") - .to_string() - .contains("cursor failed") - ); - - let append_lookup_store = LocalEventsStore::new(ScriptedExecutor::new( - vec![Ok(ExecOutcome { - changes: 1, - last_insert_id: 0, - })], - vec![Ok(r#"[{"change_seq":1}]"#.to_owned()), Ok("[]".to_owned())], - )); - assert!( - append_lookup_store - .append_record(&local_work("local-a")) - .expect_err("append lookup failure") - .to_string() - .contains("record append failed") - ); - - let update_lookup_store = LocalEventsStore::new(ScriptedExecutor::new( - vec![Ok(ExecOutcome { - changes: 1, - last_insert_id: 0, - })], - vec![Ok(r#"[{"change_seq":1}]"#.to_owned()), Ok("[]".to_owned())], - )); - assert!( - update_lookup_store - .update_outbox(&LocalEventRecordUpdate { - record_id: "event-a".to_owned(), - status: LocalRecordStatus::Published, - outbox_status: PublishOutboxStatus::Acknowledged, - updated_at_ms: 4000, - }) - .expect_err("update lookup failure") - .to_string() - .contains("event-a") - ); - - let cursor_query_store = LocalEventsStore::new(ScriptedExecutor::new( - Vec::new(), - vec![Err(SqlError::InvalidQuery( - "cursor query failed".to_owned(), - ))], - )); - assert!( - cursor_query_store - .get_cursor("app") - .expect_err("cursor query failure") - .to_string() - .contains("cursor query failed") - ); - - let advance_cursor_query_store = LocalEventsStore::new(ScriptedExecutor::new( - vec![Ok(ExecOutcome { - changes: 1, - last_insert_id: 0, - })], - vec![Err(SqlError::InvalidQuery( - "advanced cursor query failed".to_owned(), - ))], - )); - assert!( - advance_cursor_query_store - .advance_cursor("app", 1, 1000) - .expect_err("advance cursor query failure") - .to_string() - .contains("advanced cursor query failed") - ); - - let change_query_store = LocalEventsStore::new(ScriptedExecutor::new( - Vec::new(), - vec![Err(SqlError::InvalidQuery( - "change query failed".to_owned(), - ))], - )); - assert!( - change_query_store - .append_record(&local_work("local-a")) - .expect_err("change query failure") - .to_string() - .contains("change query failed") - ); - } - - #[test] - fn store_reports_record_row_conversion_failures() { - for (field, value, expected) in [ - ("family", json!("bad_family"), "family"), - ("status", json!("bad_status"), "status"), - ("source_runtime", json!("bad_runtime"), "runtime"), - ("event_tags_json", json!("{"), "EOF"), - ("raw_event_json", json!("{"), "EOF"), - ("outbox_status", json!("bad_outbox"), "outbox"), - ] { - let store = LocalEventsStore::new(ScriptedExecutor::new( - Vec::new(), - vec![Ok(record_row_with(field, value))], - )); - let error = store.get_record("record-a").expect_err("conversion error"); - - assert!( - error.to_string().contains(expected), - "expected error to contain {expected}, got {error}" - ); - } - - let store = LocalEventsStore::new(ScriptedExecutor::new( - vec![Ok(ExecOutcome { - changes: 1, - last_insert_id: 0, - })], - vec![ - Ok(r#"[{"change_seq":1}]"#.to_owned()), - Ok(record_row_with("status", json!("published"))), - ], - )); - let updated = store - .update_outbox(&LocalEventRecordUpdate { - record_id: "record-a".to_owned(), - status: LocalRecordStatus::Published, - outbox_status: PublishOutboxStatus::Acknowledged, - updated_at_ms: 4000, - }) - .expect("scripted update"); - assert_eq!(updated.status, LocalRecordStatus::Published); - } -} diff --git a/crates/local_events/tests/order_work.rs b/crates/local_events/tests/order_work.rs @@ -1,294 +0,0 @@ -use radroots_local_events::{ - BUYER_ORDER_REQUEST_ACTOR_SOURCE_RESOLVED_ACCOUNT, - BUYER_ORDER_REQUEST_ACTOR_SOURCE_UNRESOLVED_APP, BUYER_ORDER_REQUEST_DOCUMENT_KIND, - BUYER_ORDER_REQUEST_LOCAL_WORK_RECORD_KIND, BuyerOrderRequestSupportState, - buyer_order_request_local_work_record_id, validate_buyer_order_request_local_work_payload, - validate_supported_buyer_order_request_local_work_payload, - validate_unsupported_buyer_order_request_local_work_payload, -}; -use serde_json::{Value, json}; - -#[test] -fn buyer_order_request_record_id_is_deterministic_for_app_orders() { - assert_eq!( - buyer_order_request_local_work_record_id(" order-1 ").expect("record id"), - "app:local_work:order_request:order-1" - ); -} - -#[test] -fn buyer_order_request_payload_accepts_supported_exportable_work() { - let payload = supported_payload(); - - let validation = - validate_buyer_order_request_local_work_payload(&payload).expect("valid payload"); - let supported = validate_supported_buyer_order_request_local_work_payload(&payload) - .expect("supported payload"); - - assert_eq!(validation.order_id, "ord_1"); - assert_eq!( - validation.support_state, - BuyerOrderRequestSupportState::Supported - ); - assert_eq!(validation.support_state.as_str(), "supported"); - assert!(validation.support_issues.is_empty()); - assert_eq!(supported, validation); -} - -#[test] -fn buyer_order_request_payload_accepts_explicit_unsupported_work() { - let mut payload = supported_payload(); - payload["exportability"] = json!({ - "state": "identity_unresolved", - "reason": "canonical_hex_pubkey_required" - }); - payload["support_status"] = json!({ - "state": "unsupported", - "issues": ["buyer_pubkey_required"] - }); - payload["document"]["order"]["buyer_pubkey"] = json!(""); - payload["document"]["buyer_actor"]["pubkey"] = json!(""); - payload["document"]["buyer_actor"]["source"] = - json!(BUYER_ORDER_REQUEST_ACTOR_SOURCE_UNRESOLVED_APP); - - let validation = - validate_buyer_order_request_local_work_payload(&payload).expect("valid payload"); - let unsupported = validate_unsupported_buyer_order_request_local_work_payload(&payload) - .expect("unsupported payload"); - let supported_error = validate_supported_buyer_order_request_local_work_payload(&payload) - .expect_err("unsupported payload should not validate as supported"); - - assert_eq!( - validation.support_state, - BuyerOrderRequestSupportState::Unsupported - ); - assert_eq!(validation.support_state.as_str(), "unsupported"); - assert_eq!(validation.support_issues, vec!["buyer_pubkey_required"]); - assert_eq!(unsupported, validation); - assert!(supported_error.to_string().contains("support_status.state")); -} - -#[test] -fn buyer_order_request_payload_rejects_missing_identity() { - for (path, expected) in [ - (vec!["document", "order", "listing_addr"], "listing_addr"), - ( - vec!["document", "order", "listing_event_id"], - "listing_event_id", - ), - (vec!["document", "order", "seller_pubkey"], "seller_pubkey"), - (vec!["document", "order", "buyer_pubkey"], "buyer_pubkey"), - ] { - let mut payload = supported_payload(); - set_path(&mut payload, &path, json!("")); - - assert_invalid(payload, expected); - } -} - -#[test] -fn buyer_order_request_payload_rejects_missing_items() { - let mut payload = supported_payload(); - payload["document"]["order"]["items"] = json!([]); - - assert_invalid(payload, "items"); -} - -#[test] -fn buyer_order_request_payload_rejects_invalid_item_identity() { - let mut missing_bin = supported_payload(); - missing_bin["document"]["order"]["items"][0]["bin_id"] = json!(""); - assert_invalid(missing_bin, "items[0].bin_id"); - - let mut zero_count = supported_payload(); - zero_count["document"]["order"]["items"][0]["bin_count"] = json!(0); - assert_invalid(zero_count, "items[0].bin_count"); -} - -#[test] -fn buyer_order_request_payload_rejects_invalid_economics() { - let mut missing_economics = supported_payload(); - missing_economics["document"]["order"] - .as_object_mut() - .expect("order object") - .remove("economics"); - assert_invalid(missing_economics, "economics"); - - let mut non_object_economics = supported_payload(); - non_object_economics["document"]["order"]["economics"] = Value::Null; - assert_invalid(non_object_economics, "economics"); - - let mut mismatched_currency = supported_payload(); - mismatched_currency["document"]["order"]["economics"]["items"][0]["unit_price_currency"] = - json!("CAD"); - assert_invalid(mismatched_currency, "unit_price_currency"); - - let mut mismatched_items = supported_payload(); - mismatched_items["document"]["order"]["economics"]["items"] = json!([ - { - "bin_id": "dozen-eggs", - "bin_count": 2, - "quantity_amount": "1", - "quantity_unit": "dozen", - "unit_price_amount": "8.00", - "unit_price_currency": "USD", - "line_subtotal": { - "amount": "16.00", - "currency": "USD" - } - }, - { - "bin_id": "half-dozen-eggs", - "bin_count": 1 - } - ]); - assert_invalid(mismatched_items, "economics.items"); - - let mut empty_economics_items = supported_payload(); - empty_economics_items["document"]["order"]["economics"]["items"] = json!([]); - assert_invalid(empty_economics_items, "economics.items"); - - let mut mismatched_bin = supported_payload(); - mismatched_bin["document"]["order"]["economics"]["items"][0]["bin_id"] = json!("other-bin"); - assert_invalid(mismatched_bin, "economics.items[0].bin_id"); - - let mut bad_currency_length = supported_payload(); - bad_currency_length["document"]["order"]["economics"]["currency"] = json!("US"); - assert_invalid(bad_currency_length, "currency"); - - let mut missing_line_subtotal = supported_payload(); - missing_line_subtotal["document"]["order"]["economics"]["items"][0] - .as_object_mut() - .expect("economics item") - .remove("line_subtotal"); - assert_invalid(missing_line_subtotal, "line_subtotal"); -} - -#[test] -fn buyer_order_request_payload_rejects_stale_or_conflicting_currentness() { - let mut stale = supported_payload(); - stale["currentness"]["current"] = json!(false); - assert_invalid(stale, "currentness.current"); - - let mut missing_current = supported_payload(); - missing_current["currentness"]["current"] = Value::Null; - assert_invalid(missing_current, "currentness.current"); - - let mut wrong_order = supported_payload(); - wrong_order["currentness"]["order_id"] = json!("ord_other"); - assert_invalid(wrong_order, "currentness.order_id"); -} - -#[test] -fn buyer_order_request_payload_rejects_malformed_support_status() { - let mut supported_with_issue = supported_payload(); - supported_with_issue["support_status"]["issues"] = json!(["unit_price_required"]); - assert_invalid(supported_with_issue, "support_status.issues"); - - let mut unsupported_without_issue = supported_payload(); - unsupported_without_issue["support_status"] = json!({ - "state": "unsupported", - "issues": [] - }); - assert_invalid(unsupported_without_issue, "support_status.issues"); -} - -fn supported_payload() -> Value { - json!({ - "record_kind": BUYER_ORDER_REQUEST_LOCAL_WORK_RECORD_KIND, - "scope": "app", - "exportability": { - "state": "exportable" - }, - "support_status": { - "state": "supported", - "issues": [] - }, - "currentness": { - "current": true, - "source": "app_sqlite_order", - "record_id": "app:local_work:order_request:ord_1", - "order_id": "ord_1", - "order_updated_at": "2026-05-24T12:00:00Z", - "created_at_ms": 1777777777000_i64 - }, - "document": { - "version": 1, - "kind": BUYER_ORDER_REQUEST_DOCUMENT_KIND, - "order": { - "order_id": "ord_1", - "listing_addr": "30402:seller_pubkey:listing_key", - "listing_event_id": "event-listing-1", - "buyer_pubkey": "buyer_pubkey", - "seller_pubkey": "seller_pubkey", - "items": [ - { - "bin_id": "dozen-eggs", - "bin_count": 2 - } - ], - "economics": { - "quote_id": "app-order:ord_1", - "quote_version": 1, - "pricing_basis": "listing_event", - "currency": "USD", - "items": [ - { - "bin_id": "dozen-eggs", - "bin_count": 2, - "quantity_amount": "1", - "quantity_unit": "dozen", - "unit_price_amount": "8.00", - "unit_price_currency": "USD", - "line_subtotal": { - "amount": "16.00", - "currency": "USD" - } - } - ], - "discounts": [], - "adjustments": [], - "subtotal": { - "amount": "16.00", - "currency": "USD" - }, - "discount_total": { - "amount": "0", - "currency": "USD" - }, - "adjustment_total": { - "amount": "0", - "currency": "USD" - }, - "total": { - "amount": "16.00", - "currency": "USD" - } - } - }, - "buyer_actor": { - "account_id": "buyer-account", - "pubkey": "buyer_pubkey", - "source": BUYER_ORDER_REQUEST_ACTOR_SOURCE_RESOLVED_ACCOUNT - }, - "listing_lookup": "30402:seller_pubkey:listing_key" - } - }) -} - -fn assert_invalid(payload: Value, expected: &str) { - let error = - validate_buyer_order_request_local_work_payload(&payload).expect_err("invalid payload"); - assert!( - error.to_string().contains(expected), - "expected error to contain {expected}, got {error}" - ); -} - -fn set_path(payload: &mut Value, path: &[&str], value: Value) { - let mut current = payload; - for segment in &path[..path.len() - 1] { - current = current.get_mut(*segment).expect("path segment"); - } - current[path[path.len() - 1]] = value; -} diff --git a/crates/local_events/tests/store.rs b/crates/local_events/tests/store.rs @@ -1,499 +0,0 @@ -use radroots_local_events::{ - LocalEventRecordInput, LocalEventRecordUpdate, LocalEventsStore, LocalRecordFamily, - LocalRecordStatus, MIGRATIONS, PublishOutboxStatus, SourceRuntime, -}; -use radroots_sql_core::migrations::migrations_run_all_up; -use radroots_sql_core::{SqlExecutor, SqliteExecutor}; -use serde_json::json; - -fn store() -> LocalEventsStore<SqliteExecutor> { - let executor = SqliteExecutor::open_memory().expect("open memory sqlite"); - let store = LocalEventsStore::new(executor); - store.migrate_up().expect("migrate local events"); - store -} - -fn local_work(record_id: &str) -> LocalEventRecordInput { - LocalEventRecordInput { - record_id: record_id.to_owned(), - family: LocalRecordFamily::LocalWork, - status: LocalRecordStatus::LocalSaved, - source_runtime: SourceRuntime::Cli, - created_at_ms: 1000, - inserted_at_ms: 1001, - owner_account_id: Some("seller-account".to_owned()), - owner_pubkey: Some("seller-pubkey".to_owned()), - farm_id: Some("farm-a".to_owned()), - listing_addr: Some("listing-a".to_owned()), - local_work_json: Some(json!({"kind":"listing","title":"Eggs"})), - event_id: None, - event_kind: None, - event_pubkey: None, - event_created_at: None, - event_tags_json: None, - event_content: None, - event_sig: None, - raw_event_json: None, - outbox_status: PublishOutboxStatus::None, - } -} - -fn signed_event(record_id: &str) -> LocalEventRecordInput { - LocalEventRecordInput { - record_id: record_id.to_owned(), - family: LocalRecordFamily::SignedEvent, - status: LocalRecordStatus::PendingPublish, - source_runtime: SourceRuntime::Cli, - created_at_ms: 2000, - inserted_at_ms: 2001, - owner_account_id: Some("seller-account".to_owned()), - owner_pubkey: Some("seller-pubkey".to_owned()), - farm_id: Some("farm-a".to_owned()), - listing_addr: Some("listing-a".to_owned()), - local_work_json: None, - event_id: Some("event-a".to_owned()), - event_kind: Some(3421), - event_pubkey: Some("seller-pubkey".to_owned()), - event_created_at: Some(2000), - event_tags_json: Some(json!([["d", "listing-a"]])), - event_content: Some("{\"title\":\"Eggs\"}".to_owned()), - event_sig: Some("sig-a".to_owned()), - raw_event_json: Some(json!({"id":"event-a","kind":3421})), - outbox_status: PublishOutboxStatus::Pending, - } -} - -#[test] -fn append_rejects_malformed_local_work_records() { - let store = store(); - let mut input = local_work("local-a"); - input.local_work_json = None; - - let err = store.append_record(&input).expect_err("invalid record"); - - assert!(err.to_string().contains("local_work_json")); -} - -#[test] -fn append_is_idempotent_by_record_id() { - let store = store(); - let input = local_work("local-a"); - - let first = store.append_record(&input).expect("append first"); - let second = store.append_record(&input).expect("append second"); - let rows = store.list_records_after_seq(0, 10).expect("list records"); - - assert_eq!(first.seq, second.seq); - assert_eq!(first.change_seq, second.change_seq); - assert_eq!(rows.len(), 1); - assert_eq!(rows[0].record_id, "local-a"); - assert_eq!( - rows[0].local_work_json, - Some(json!({"kind":"listing","title":"Eggs"})) - ); -} - -#[test] -fn source_runtime_network_round_trips() { - let store = store(); - let mut input = signed_event("event-network-a"); - input.source_runtime = SourceRuntime::Network; - - let inserted = store.append_record(&input).expect("append network event"); - let rows = store - .list_records_after_seq(0, 10) - .expect("list network event"); - - assert_eq!(SourceRuntime::Network.as_str(), "network"); - assert_eq!( - SourceRuntime::parse("network").expect("parse network runtime"), - SourceRuntime::Network - ); - assert_eq!(inserted.source_runtime, SourceRuntime::Network); - assert_eq!(rows.len(), 1); - assert_eq!(rows[0].source_runtime, SourceRuntime::Network); -} - -#[test] -fn projection_cursor_advances_without_rewinding() { - let store = store(); - - let first = store - .advance_cursor("app", 10, 100) - .expect("advance cursor"); - let second = store.advance_cursor("app", 5, 200).expect("ignore rewind"); - let third = store.advance_cursor("app", 12, 300).expect("advance again"); - - assert_eq!(first.last_change_seq, 10); - assert_eq!(second.last_change_seq, 10); - assert_eq!(third.last_change_seq, 12); -} - -#[test] -fn outbox_status_updates_signed_event_records() { - let store = store(); - let input = signed_event("event-a"); - store.append_record(&input).expect("append signed event"); - - let updated = store - .update_outbox(&LocalEventRecordUpdate { - record_id: "event-a".to_owned(), - status: LocalRecordStatus::Published, - outbox_status: PublishOutboxStatus::Acknowledged, - updated_at_ms: 3000, - }) - .expect("update outbox"); - - assert_eq!(updated.status, LocalRecordStatus::Published); - assert_eq!(updated.outbox_status, PublishOutboxStatus::Acknowledged); -} - -#[test] -fn changed_after_uses_change_seq_for_appends_and_outbox_updates() { - let store = store(); - let input = signed_event("event-a"); - let appended = store.append_record(&input).expect("append signed event"); - let initial_rows = store - .list_records_changed_after(0, 10) - .expect("list initial changes"); - - assert_eq!(initial_rows.len(), 1); - assert_eq!(initial_rows[0].record_id, "event-a"); - assert_eq!(initial_rows[0].seq, appended.seq); - assert_eq!(initial_rows[0].change_seq, appended.change_seq); - - let updated = store - .update_outbox(&LocalEventRecordUpdate { - record_id: "event-a".to_owned(), - status: LocalRecordStatus::Published, - outbox_status: PublishOutboxStatus::Acknowledged, - updated_at_ms: 3000, - }) - .expect("update outbox"); - let changed_rows = store - .list_records_changed_after(appended.change_seq, 10) - .expect("list changed rows"); - - assert_eq!(updated.seq, appended.seq); - assert!(updated.change_seq > appended.change_seq); - assert_eq!(changed_rows.len(), 1); - assert_eq!(changed_rows[0].record_id, "event-a"); - assert_eq!(changed_rows[0].change_seq, updated.change_seq); -} - -#[test] -fn changed_latest_lists_newest_records_first() { - let store = store(); - let first = store - .append_record(&local_work("local-a")) - .expect("append first"); - let second = store - .append_record(&local_work("local-b")) - .expect("append second"); - let third = store - .append_record(&local_work("local-c")) - .expect("append third"); - - let rows = store - .list_records_changed_latest(2) - .expect("list latest changed rows"); - - assert_eq!(rows.len(), 2); - assert_eq!(rows[0].record_id, "local-c"); - assert_eq!(rows[0].change_seq, third.change_seq); - assert_eq!(rows[1].record_id, "local-b"); - assert_eq!(rows[1].change_seq, second.change_seq); - assert!(rows[1].change_seq > first.change_seq); -} - -#[test] -fn changed_before_pages_newest_first_by_cursor() { - let store = store(); - let _first = store - .append_record(&local_work("local-a")) - .expect("append first"); - let second = store - .append_record(&local_work("local-b")) - .expect("append second"); - let third = store - .append_record(&local_work("local-c")) - .expect("append third"); - let fourth = store - .append_record(&local_work("local-d")) - .expect("append fourth"); - - let first_page = store - .list_records_changed_latest(2) - .expect("list first page"); - let cursor = first_page.last().expect("last first page"); - let second_page = store - .list_records_changed_before(cursor.change_seq, cursor.seq, 2) - .expect("list second page"); - - assert_eq!(first_page.len(), 2); - assert_eq!(first_page[0].record_id, "local-d"); - assert_eq!(first_page[0].change_seq, fourth.change_seq); - assert_eq!(first_page[1].record_id, "local-c"); - assert_eq!(first_page[1].change_seq, third.change_seq); - assert_eq!(second_page.len(), 2); - assert_eq!(second_page[0].record_id, "local-b"); - assert_eq!(second_page[0].change_seq, second.change_seq); - assert_eq!(second_page[1].record_id, "local-a"); -} - -#[test] -fn changed_latest_is_not_blocked_by_older_record_volume() { - let store = store(); - for index in 0..505 { - store - .append_record(&local_work(&format!("older-{index:03}"))) - .expect("append older record"); - } - let current = store - .append_record(&local_work("current-record")) - .expect("append current record"); - - let rows = store - .list_records_changed_latest(1) - .expect("list latest record"); - - assert_eq!(rows.len(), 1); - assert_eq!(rows[0].record_id, "current-record"); - assert_eq!(rows[0].change_seq, current.change_seq); -} - -#[test] -fn migration_assigns_existing_records_change_seq_from_insert_order() { - let executor = SqliteExecutor::open_memory().expect("open memory sqlite"); - migrations_run_all_up(&executor, &MIGRATIONS[..1]).expect("apply initial migration"); - let first = insert_pre_change_tracking_record(&executor, "local-a"); - let second = insert_pre_change_tracking_record(&executor, "local-b"); - let store = LocalEventsStore::new(executor); - - store.migrate_up().expect("apply change tracking migration"); - let rows = store - .list_records_changed_after(0, 10) - .expect("list changed rows after migration"); - - assert_eq!(rows.len(), 2); - assert_eq!(rows[0].seq, first); - assert_eq!(rows[0].change_seq, first); - assert_eq!(rows[1].seq, second); - assert_eq!(rows[1].change_seq, second); -} - -#[test] -fn migration_repairs_pre_network_source_runtime_constraint() { - let executor = SqliteExecutor::open_memory().expect("open memory sqlite"); - create_pre_network_change_tracking_schema(&executor); - let legacy_seq = insert_pre_network_change_tracking_record(&executor, "legacy-cli", 1); - let store = LocalEventsStore::new(executor); - - store - .migrate_up() - .expect("apply network source repair migration"); - let mut input = signed_event("event-network-repaired"); - input.source_runtime = SourceRuntime::Network; - input.event_id = Some("event-network-repaired".to_owned()); - input.raw_event_json = Some(json!({"id":"event-network-repaired","kind":3421})); - let inserted = store - .append_record(&input) - .expect("append repaired network event"); - let rows = store - .list_records_changed_after(0, 10) - .expect("list changed rows after repair"); - - assert_eq!(legacy_seq, 1); - assert_eq!(rows.len(), 2); - assert_eq!(rows[0].record_id, "legacy-cli"); - assert_eq!(rows[0].change_seq, 1); - assert_eq!(rows[0].source_runtime, SourceRuntime::Cli); - assert_eq!(rows[1].record_id, "event-network-repaired"); - assert_eq!(rows[1].seq, inserted.seq); - assert_eq!(rows[1].source_runtime, SourceRuntime::Network); -} - -fn insert_pre_change_tracking_record(executor: &SqliteExecutor, record_id: &str) -> i64 { - let input = local_work(record_id); - let params = json!([ - input.record_id, - input.family.as_str(), - input.status.as_str(), - input.source_runtime.as_str(), - input.created_at_ms, - input.inserted_at_ms, - input.inserted_at_ms, - input.owner_account_id, - input.owner_pubkey, - input.farm_id, - input.listing_addr, - serde_json::to_string(&input.local_work_json).expect("encode local work"), - input.event_id, - input.event_kind, - input.event_pubkey, - input.event_created_at, - input - .event_tags_json - .map(|value| serde_json::to_string(&value).expect("encode tags")), - input.event_content, - input.event_sig, - input - .raw_event_json - .map(|value| serde_json::to_string(&value).expect("encode raw event")), - input.outbox_status.as_str(), - ]) - .to_string(); - let outcome = executor - .exec( - "insert into local_event_record( - record_id, - family, - status, - source_runtime, - created_at_ms, - inserted_at_ms, - updated_at_ms, - owner_account_id, - owner_pubkey, - farm_id, - listing_addr, - local_work_json, - event_id, - event_kind, - event_pubkey, - event_created_at, - event_tags_json, - event_content, - event_sig, - raw_event_json, - outbox_status - ) values(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)", - &params, - ) - .expect("insert old local event record"); - outcome.last_insert_id -} - -fn create_pre_network_change_tracking_schema(executor: &SqliteExecutor) { - let schema = [ - "create table __migrations(id integer primary key, name text not null unique, applied_at text not null default (datetime('now')))", - "create table local_event_record ( - seq integer primary key autoincrement, - change_seq integer not null unique, - record_id text not null unique, - family text not null check (family in ('local_work', 'signed_event')), - status text not null check (status in ('local_draft', 'local_saved', 'pending_publish', 'published', 'failed', 'conflict')), - source_runtime text not null check (source_runtime in ('cli', 'app', 'service', 'worker', 'test')), - created_at_ms integer not null, - inserted_at_ms integer not null, - updated_at_ms integer not null, - owner_account_id text, - owner_pubkey text, - farm_id text, - listing_addr text, - local_work_json text, - event_id text, - event_kind integer, - event_pubkey text, - event_created_at integer, - event_tags_json text, - event_content text, - event_sig text, - raw_event_json text, - outbox_status text not null check (outbox_status in ('none', 'pending', 'acknowledged', 'failed')), - check (change_seq >= 1), - check (trim(record_id) <> ''), - check (family <> 'local_work' or local_work_json is not null), - check (family <> 'local_work' or outbox_status = 'none'), - check (family <> 'signed_event' or (event_id is not null and event_kind is not null and event_pubkey is not null and event_sig is not null and raw_event_json is not null)) - )", - "create index local_event_record_change_seq_idx on local_event_record(change_seq)", - "create index local_event_record_event_id_idx on local_event_record(event_id)", - "create index local_event_record_listing_addr_idx on local_event_record(listing_addr)", - "create index local_event_record_owner_pubkey_idx on local_event_record(owner_pubkey)", - "create index local_event_record_status_idx on local_event_record(status)", - "create table local_event_projection_cursor ( - consumer_id text primary key, - last_change_seq integer not null, - updated_at_ms integer not null, - check (trim(consumer_id) <> ''), - check (last_change_seq >= 0) - )", - ]; - for sql in schema { - executor.exec(sql, "[]").expect("schema statement"); - } - for name in ["0000_local_events", "0001_change_tracking"] { - let params = json!([name]).to_string(); - executor - .exec("insert into __migrations(name) values(?)", &params) - .expect("migration marker"); - } -} - -fn insert_pre_network_change_tracking_record( - executor: &SqliteExecutor, - record_id: &str, - change_seq: i64, -) -> i64 { - let input = local_work(record_id); - let params = json!([ - change_seq, - input.record_id, - input.family.as_str(), - input.status.as_str(), - input.source_runtime.as_str(), - input.created_at_ms, - input.inserted_at_ms, - input.inserted_at_ms, - input.owner_account_id, - input.owner_pubkey, - input.farm_id, - input.listing_addr, - serde_json::to_string(&input.local_work_json).expect("encode local work"), - input.event_id, - input.event_kind, - input.event_pubkey, - input.event_created_at, - input - .event_tags_json - .map(|value| serde_json::to_string(&value).expect("encode tags")), - input.event_content, - input.event_sig, - input - .raw_event_json - .map(|value| serde_json::to_string(&value).expect("encode raw event")), - input.outbox_status.as_str(), - ]) - .to_string(); - let outcome = executor - .exec( - "insert into local_event_record( - change_seq, - record_id, - family, - status, - source_runtime, - created_at_ms, - inserted_at_ms, - updated_at_ms, - owner_account_id, - owner_pubkey, - farm_id, - listing_addr, - local_work_json, - event_id, - event_kind, - event_pubkey, - event_created_at, - event_tags_json, - event_content, - event_sig, - raw_event_json, - outbox_status - ) values(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)", - &params, - ) - .expect("insert pre-network local event record"); - outcome.last_insert_id -} diff --git a/crates/runtime_paths/src/conventions.rs b/crates/runtime_paths/src/conventions.rs @@ -11,10 +11,10 @@ pub const DEFAULT_SHARED_IDENTITY_FILE_NAME: &str = "default.json"; pub const DEFAULT_SHARED_GEONAMES_NAMESPACE_KIND: &str = "shared"; pub const DEFAULT_SHARED_GEONAMES_NAMESPACE_VALUE: &str = "geonames"; pub const DEFAULT_SHARED_GEONAMES_NAMESPACE: &str = "shared/geonames"; -pub const DEFAULT_SHARED_LOCAL_EVENTS_NAMESPACE_KIND: &str = "shared"; -pub const DEFAULT_SHARED_LOCAL_EVENTS_NAMESPACE_VALUE: &str = "local_events"; -pub const DEFAULT_SHARED_LOCAL_EVENTS_NAMESPACE: &str = "shared/local_events"; -pub const DEFAULT_SHARED_LOCAL_EVENTS_DB_FILE_NAME: &str = "local_events.sqlite"; +pub const DEFAULT_SHARED_RUNTIME_STORE_NAMESPACE_KIND: &str = "shared"; +pub const DEFAULT_SHARED_RUNTIME_STORE_NAMESPACE_VALUE: &str = "runtime_store"; +pub const DEFAULT_SHARED_RUNTIME_STORE_NAMESPACE: &str = "shared/runtime_store"; +pub const DEFAULT_SHARED_RUNTIME_STORE_DB_FILE_NAME: &str = "runtime_store.sqlite"; #[derive(Debug, Clone, PartialEq, Eq)] pub struct RadrootsBootstrapPaths { @@ -59,19 +59,19 @@ pub fn default_shared_runtime_logs_dir( } #[must_use] -pub fn default_shared_local_events_root_from_data_root(data_root: impl AsRef<Path>) -> PathBuf { +pub fn default_shared_runtime_store_root_from_data_root(data_root: impl AsRef<Path>) -> PathBuf { data_root .as_ref() - .join(DEFAULT_SHARED_LOCAL_EVENTS_NAMESPACE_KIND) - .join(DEFAULT_SHARED_LOCAL_EVENTS_NAMESPACE_VALUE) + .join(DEFAULT_SHARED_RUNTIME_STORE_NAMESPACE_KIND) + .join(DEFAULT_SHARED_RUNTIME_STORE_NAMESPACE_VALUE) } #[must_use] -pub fn default_shared_local_events_database_path_from_data_root( +pub fn default_shared_runtime_store_database_path_from_data_root( data_root: impl AsRef<Path>, ) -> PathBuf { - default_shared_local_events_root_from_data_root(data_root) - .join(DEFAULT_SHARED_LOCAL_EVENTS_DB_FILE_NAME) + default_shared_runtime_store_root_from_data_root(data_root) + .join(DEFAULT_SHARED_RUNTIME_STORE_DB_FILE_NAME) } #[must_use] @@ -96,7 +96,7 @@ pub fn default_shared_geonames_database_path_from_cache_root( .join(default_shared_geonames_database_file_name(version)) } -pub fn default_shared_local_events_root_from_shared_accounts_data_root( +pub fn default_shared_runtime_store_root_from_shared_accounts_data_root( shared_accounts_data_root: impl AsRef<Path>, ) -> Result<PathBuf, RadrootsRuntimePathsError> { let shared_accounts_data_root = shared_accounts_data_root.as_ref(); @@ -105,14 +105,14 @@ pub fn default_shared_local_events_root_from_shared_accounts_data_root( path: shared_accounts_data_root.to_path_buf(), } })?; - Ok(shared_data_root.join(DEFAULT_SHARED_LOCAL_EVENTS_NAMESPACE_VALUE)) + Ok(shared_data_root.join(DEFAULT_SHARED_RUNTIME_STORE_NAMESPACE_VALUE)) } -pub fn default_shared_local_events_database_path_from_shared_accounts_data_root( +pub fn default_shared_runtime_store_database_path_from_shared_accounts_data_root( shared_accounts_data_root: impl AsRef<Path>, ) -> Result<PathBuf, RadrootsRuntimePathsError> { - default_shared_local_events_root_from_shared_accounts_data_root(shared_accounts_data_root) - .map(|root| root.join(DEFAULT_SHARED_LOCAL_EVENTS_DB_FILE_NAME)) + default_shared_runtime_store_root_from_shared_accounts_data_root(shared_accounts_data_root) + .map(|root| root.join(DEFAULT_SHARED_RUNTIME_STORE_DB_FILE_NAME)) } #[cfg(test)] @@ -126,16 +126,15 @@ mod tests { use super::{ DEFAULT_SERVICE_IDENTITY_FILE_NAME, DEFAULT_SHARED_GEONAMES_NAMESPACE, - DEFAULT_SHARED_IDENTITY_FILE_NAME, DEFAULT_SHARED_LOCAL_EVENTS_DB_FILE_NAME, - DEFAULT_SHARED_LOCAL_EVENTS_NAMESPACE, default_namespaced_bootstrap_paths, + DEFAULT_SHARED_IDENTITY_FILE_NAME, DEFAULT_SHARED_RUNTIME_STORE_DB_FILE_NAME, + DEFAULT_SHARED_RUNTIME_STORE_NAMESPACE, default_namespaced_bootstrap_paths, default_shared_geonames_database_file_name, default_shared_geonames_database_path_from_cache_root, default_shared_geonames_root_from_cache_root, default_shared_identity_path, - default_shared_local_events_database_path_from_data_root, - default_shared_local_events_database_path_from_shared_accounts_data_root, - default_shared_local_events_root_from_data_root, - default_shared_local_events_root_from_shared_accounts_data_root, - default_shared_runtime_logs_dir, + default_shared_runtime_logs_dir, default_shared_runtime_store_database_path_from_data_root, + default_shared_runtime_store_database_path_from_shared_accounts_data_root, + default_shared_runtime_store_root_from_data_root, + default_shared_runtime_store_root_from_shared_accounts_data_root, }; #[test] @@ -210,18 +209,18 @@ mod tests { } #[test] - fn shared_local_events_paths_use_canonical_shared_namespace() { + fn shared_runtime_store_paths_use_canonical_shared_namespace() { let data_root = PathBuf::from("/repo/infra/local/runtime/radroots/data"); assert_eq!( - default_shared_local_events_root_from_data_root(&data_root), - data_root.join(DEFAULT_SHARED_LOCAL_EVENTS_NAMESPACE) + default_shared_runtime_store_root_from_data_root(&data_root), + data_root.join(DEFAULT_SHARED_RUNTIME_STORE_NAMESPACE) ); assert_eq!( - default_shared_local_events_database_path_from_data_root(&data_root), + default_shared_runtime_store_database_path_from_data_root(&data_root), data_root - .join(DEFAULT_SHARED_LOCAL_EVENTS_NAMESPACE) - .join(DEFAULT_SHARED_LOCAL_EVENTS_DB_FILE_NAME) + .join(DEFAULT_SHARED_RUNTIME_STORE_NAMESPACE) + .join(DEFAULT_SHARED_RUNTIME_STORE_DB_FILE_NAME) ); } @@ -246,36 +245,36 @@ mod tests { } #[test] - fn shared_local_events_paths_derive_from_shared_accounts_data_root() { + fn shared_runtime_store_paths_derive_from_shared_accounts_data_root() { let shared_accounts_data_root = PathBuf::from("/repo/infra/local/runtime/radroots/data/shared/accounts"); assert_eq!( - default_shared_local_events_root_from_shared_accounts_data_root( + default_shared_runtime_store_root_from_shared_accounts_data_root( &shared_accounts_data_root ) - .expect("shared local-events root"), - PathBuf::from("/repo/infra/local/runtime/radroots/data/shared/local_events") + .expect("shared runtime-store root"), + PathBuf::from("/repo/infra/local/runtime/radroots/data/shared/runtime_store") ); assert_eq!( - default_shared_local_events_root_from_shared_accounts_data_root( + default_shared_runtime_store_root_from_shared_accounts_data_root( shared_accounts_data_root.clone() ) - .expect("shared local-events root from owned path"), - PathBuf::from("/repo/infra/local/runtime/radroots/data/shared/local_events") + .expect("shared runtime-store root from owned path"), + PathBuf::from("/repo/infra/local/runtime/radroots/data/shared/runtime_store") ); assert_eq!( - default_shared_local_events_database_path_from_shared_accounts_data_root( + default_shared_runtime_store_database_path_from_shared_accounts_data_root( &shared_accounts_data_root ) - .expect("shared local-events database path"), + .expect("shared runtime-store database path"), PathBuf::from( - "/repo/infra/local/runtime/radroots/data/shared/local_events/local_events.sqlite" + "/repo/infra/local/runtime/radroots/data/shared/runtime_store/runtime_store.sqlite" ) ); let err = - default_shared_local_events_root_from_shared_accounts_data_root(PathBuf::from("/")) + default_shared_runtime_store_root_from_shared_accounts_data_root(PathBuf::from("/")) .expect_err("root path has no parent shared data root"); assert_eq!( err, diff --git a/crates/runtime_paths/src/lib.rs b/crates/runtime_paths/src/lib.rs @@ -11,17 +11,16 @@ pub use conventions::{ DEFAULT_CONFIG_FILE_NAME, DEFAULT_SERVICE_IDENTITY_FILE_NAME, DEFAULT_SHARED_GEONAMES_NAMESPACE, DEFAULT_SHARED_GEONAMES_NAMESPACE_KIND, DEFAULT_SHARED_GEONAMES_NAMESPACE_VALUE, DEFAULT_SHARED_IDENTITY_FILE_NAME, - DEFAULT_SHARED_LOCAL_EVENTS_DB_FILE_NAME, DEFAULT_SHARED_LOCAL_EVENTS_NAMESPACE, - DEFAULT_SHARED_LOCAL_EVENTS_NAMESPACE_KIND, DEFAULT_SHARED_LOCAL_EVENTS_NAMESPACE_VALUE, + DEFAULT_SHARED_RUNTIME_STORE_DB_FILE_NAME, DEFAULT_SHARED_RUNTIME_STORE_NAMESPACE, + DEFAULT_SHARED_RUNTIME_STORE_NAMESPACE_KIND, DEFAULT_SHARED_RUNTIME_STORE_NAMESPACE_VALUE, RadrootsBootstrapPaths, default_namespaced_bootstrap_paths, default_shared_geonames_database_file_name, default_shared_geonames_database_path_from_cache_root, default_shared_geonames_root_from_cache_root, default_shared_identity_path, - default_shared_local_events_database_path_from_data_root, - default_shared_local_events_database_path_from_shared_accounts_data_root, - default_shared_local_events_root_from_data_root, - default_shared_local_events_root_from_shared_accounts_data_root, - default_shared_runtime_logs_dir, + default_shared_runtime_logs_dir, default_shared_runtime_store_database_path_from_data_root, + default_shared_runtime_store_database_path_from_shared_accounts_data_root, + default_shared_runtime_store_root_from_data_root, + default_shared_runtime_store_root_from_shared_accounts_data_root, }; pub use error::RadrootsRuntimePathsError; pub use namespace::{RadrootsRuntimeNamespace, RadrootsRuntimeNamespaceKind}; diff --git a/crates/runtime_store/Cargo.toml b/crates/runtime_store/Cargo.toml @@ -0,0 +1,29 @@ +[package] +name = "radroots_runtime_store" +publish = ["crates-io"] +version = "0.1.0-alpha.2" +edition.workspace = true +authors = ["Tyson Lupul <tyson@radroots.org>"] +rust-version.workspace = true +license.workspace = true +description = "Runtime store workspace for Radroots" +repository.workspace = true +homepage.workspace = true +documentation = "https://docs.rs/radroots_runtime_store" +readme = "README" + +[lib] +crate-type = ["rlib"] + +[features] +default = [] +native = ["radroots_sql_core/native"] + +[dependencies] +radroots_sql_core = { workspace = true } +serde = { workspace = true } +serde_json = { workspace = true } +thiserror = { workspace = true } + +[dev-dependencies] +radroots_sql_core = { workspace = true, features = ["native"] } diff --git a/crates/runtime_store/README b/crates/runtime_store/README @@ -0,0 +1,8 @@ +# radroots_runtime_store + +Shared runtime-store primitives for same-host Rad Roots runtimes. + +This crate owns the SQLite schema and typed API for the `shared/runtime_store` +namespace. It is an interop store for local work records, signed event records, +publish outbox status, relay delivery metadata, and projection cursors. It is +not an application primary database. diff --git a/crates/runtime_store/migrations/0000_runtime_store.down.sql b/crates/runtime_store/migrations/0000_runtime_store.down.sql @@ -0,0 +1,2 @@ +drop table if exists runtime_store_projection_cursor; +drop table if exists runtime_store_record; diff --git a/crates/runtime_store/migrations/0000_runtime_store.up.sql b/crates/runtime_store/migrations/0000_runtime_store.up.sql @@ -0,0 +1,43 @@ +create table if not exists runtime_store_record ( + seq integer primary key autoincrement, + record_id text not null unique, + family text not null check (family in ('local_work', 'signed_event')), + status text not null check (status in ('local_draft', 'local_saved', 'pending_publish', 'published', 'failed', 'conflict')), + source_runtime text not null check (source_runtime in ('cli', 'app', 'network', 'service', 'worker', 'test')), + created_at_ms integer not null, + inserted_at_ms integer not null, + updated_at_ms integer not null, + owner_account_id text, + owner_pubkey text, + farm_id text, + listing_addr text, + local_work_json text, + event_id text, + event_kind integer, + event_pubkey text, + event_created_at integer, + event_tags_json text, + event_content text, + event_sig text, + raw_event_json text, + outbox_status text not null check (outbox_status in ('none', 'pending', 'acknowledged', 'failed')), + relay_set_fingerprint text, + relay_delivery_json text, + check (trim(record_id) <> ''), + check (family <> 'local_work' or local_work_json is not null), + check (family <> 'local_work' or outbox_status = 'none'), + check (family <> 'signed_event' or (event_id is not null and event_kind is not null and event_pubkey is not null and event_sig is not null and raw_event_json is not null)) +); + +create index if not exists runtime_store_record_event_id_idx on runtime_store_record(event_id); +create index if not exists runtime_store_record_listing_addr_idx on runtime_store_record(listing_addr); +create index if not exists runtime_store_record_owner_pubkey_idx on runtime_store_record(owner_pubkey); +create index if not exists runtime_store_record_status_idx on runtime_store_record(status); + +create table if not exists runtime_store_projection_cursor ( + consumer_id text primary key, + last_seq integer not null, + updated_at_ms integer not null, + check (trim(consumer_id) <> ''), + check (last_seq >= 0) +); diff --git a/crates/runtime_store/migrations/0001_change_tracking.down.sql b/crates/runtime_store/migrations/0001_change_tracking.down.sql @@ -0,0 +1,114 @@ +create table runtime_store_projection_cursor_previous ( + consumer_id text primary key, + last_seq integer not null, + updated_at_ms integer not null, + check (trim(consumer_id) <> ''), + check (last_seq >= 0) +); + +insert into runtime_store_projection_cursor_previous( + consumer_id, + last_seq, + updated_at_ms +) +select + consumer_id, + last_change_seq, + updated_at_ms +from runtime_store_projection_cursor; + +drop table runtime_store_projection_cursor; +alter table runtime_store_projection_cursor_previous rename to runtime_store_projection_cursor; + +create table runtime_store_record_previous ( + seq integer primary key autoincrement, + record_id text not null unique, + family text not null check (family in ('local_work', 'signed_event')), + status text not null check (status in ('local_draft', 'local_saved', 'pending_publish', 'published', 'failed', 'conflict')), + source_runtime text not null check (source_runtime in ('cli', 'app', 'network', 'service', 'worker', 'test')), + created_at_ms integer not null, + inserted_at_ms integer not null, + updated_at_ms integer not null, + owner_account_id text, + owner_pubkey text, + farm_id text, + listing_addr text, + local_work_json text, + event_id text, + event_kind integer, + event_pubkey text, + event_created_at integer, + event_tags_json text, + event_content text, + event_sig text, + raw_event_json text, + outbox_status text not null check (outbox_status in ('none', 'pending', 'acknowledged', 'failed')), + relay_set_fingerprint text, + relay_delivery_json text, + check (trim(record_id) <> ''), + check (family <> 'local_work' or local_work_json is not null), + check (family <> 'local_work' or outbox_status = 'none'), + check (family <> 'signed_event' or (event_id is not null and event_kind is not null and event_pubkey is not null and event_sig is not null and raw_event_json is not null)) +); + +insert into runtime_store_record_previous( + seq, + record_id, + family, + status, + source_runtime, + created_at_ms, + inserted_at_ms, + updated_at_ms, + owner_account_id, + owner_pubkey, + farm_id, + listing_addr, + local_work_json, + event_id, + event_kind, + event_pubkey, + event_created_at, + event_tags_json, + event_content, + event_sig, + raw_event_json, + outbox_status, + relay_set_fingerprint, + relay_delivery_json +) +select + seq, + record_id, + family, + status, + source_runtime, + created_at_ms, + inserted_at_ms, + updated_at_ms, + owner_account_id, + owner_pubkey, + farm_id, + listing_addr, + local_work_json, + event_id, + event_kind, + event_pubkey, + event_created_at, + event_tags_json, + event_content, + event_sig, + raw_event_json, + outbox_status, + relay_set_fingerprint, + relay_delivery_json +from runtime_store_record +order by seq asc; + +drop table runtime_store_record; +alter table runtime_store_record_previous rename to runtime_store_record; + +create index runtime_store_record_event_id_idx on runtime_store_record(event_id); +create index runtime_store_record_listing_addr_idx on runtime_store_record(listing_addr); +create index runtime_store_record_owner_pubkey_idx on runtime_store_record(owner_pubkey); +create index runtime_store_record_status_idx on runtime_store_record(status); diff --git a/crates/runtime_store/migrations/0001_change_tracking.up.sql b/crates/runtime_store/migrations/0001_change_tracking.up.sql @@ -0,0 +1,119 @@ +create table runtime_store_record_next ( + seq integer primary key autoincrement, + change_seq integer not null unique, + record_id text not null unique, + family text not null check (family in ('local_work', 'signed_event')), + status text not null check (status in ('local_draft', 'local_saved', 'pending_publish', 'published', 'failed', 'conflict')), + source_runtime text not null check (source_runtime in ('cli', 'app', 'network', 'service', 'worker', 'test')), + created_at_ms integer not null, + inserted_at_ms integer not null, + updated_at_ms integer not null, + owner_account_id text, + owner_pubkey text, + farm_id text, + listing_addr text, + local_work_json text, + event_id text, + event_kind integer, + event_pubkey text, + event_created_at integer, + event_tags_json text, + event_content text, + event_sig text, + raw_event_json text, + outbox_status text not null check (outbox_status in ('none', 'pending', 'acknowledged', 'failed')), + relay_set_fingerprint text, + relay_delivery_json text, + check (change_seq >= 1), + check (trim(record_id) <> ''), + check (family <> 'local_work' or local_work_json is not null), + check (family <> 'local_work' or outbox_status = 'none'), + check (family <> 'signed_event' or (event_id is not null and event_kind is not null and event_pubkey is not null and event_sig is not null and raw_event_json is not null)) +); + +insert into runtime_store_record_next( + seq, + change_seq, + record_id, + family, + status, + source_runtime, + created_at_ms, + inserted_at_ms, + updated_at_ms, + owner_account_id, + owner_pubkey, + farm_id, + listing_addr, + local_work_json, + event_id, + event_kind, + event_pubkey, + event_created_at, + event_tags_json, + event_content, + event_sig, + raw_event_json, + outbox_status, + relay_set_fingerprint, + relay_delivery_json +) +select + seq, + seq, + record_id, + family, + status, + source_runtime, + created_at_ms, + inserted_at_ms, + updated_at_ms, + owner_account_id, + owner_pubkey, + farm_id, + listing_addr, + local_work_json, + event_id, + event_kind, + event_pubkey, + event_created_at, + event_tags_json, + event_content, + event_sig, + raw_event_json, + outbox_status, + relay_set_fingerprint, + relay_delivery_json +from runtime_store_record +order by seq asc; + +drop table runtime_store_record; +alter table runtime_store_record_next rename to runtime_store_record; + +create index runtime_store_record_change_seq_idx on runtime_store_record(change_seq); +create index runtime_store_record_event_id_idx on runtime_store_record(event_id); +create index runtime_store_record_listing_addr_idx on runtime_store_record(listing_addr); +create index runtime_store_record_owner_pubkey_idx on runtime_store_record(owner_pubkey); +create index runtime_store_record_status_idx on runtime_store_record(status); + +create table runtime_store_projection_cursor_next ( + consumer_id text primary key, + last_change_seq integer not null, + updated_at_ms integer not null, + check (trim(consumer_id) <> ''), + check (last_change_seq >= 0) +); + +insert into runtime_store_projection_cursor_next( + consumer_id, + last_change_seq, + updated_at_ms +) +select + consumer_id, + last_seq, + updated_at_ms +from runtime_store_projection_cursor; + +drop table runtime_store_projection_cursor; +alter table runtime_store_projection_cursor_next rename to runtime_store_projection_cursor; diff --git a/crates/local_events/migrations/0002_network_source_runtime.down.sql b/crates/runtime_store/migrations/0002_network_source_runtime.down.sql diff --git a/crates/runtime_store/migrations/0002_network_source_runtime.up.sql b/crates/runtime_store/migrations/0002_network_source_runtime.up.sql @@ -0,0 +1,97 @@ +create table runtime_store_record_network_source_next ( + seq integer primary key autoincrement, + change_seq integer not null unique, + record_id text not null unique, + family text not null check (family in ('local_work', 'signed_event')), + status text not null check (status in ('local_draft', 'local_saved', 'pending_publish', 'published', 'failed', 'conflict')), + source_runtime text not null check (source_runtime in ('cli', 'app', 'network', 'service', 'worker', 'test')), + created_at_ms integer not null, + inserted_at_ms integer not null, + updated_at_ms integer not null, + owner_account_id text, + owner_pubkey text, + farm_id text, + listing_addr text, + local_work_json text, + event_id text, + event_kind integer, + event_pubkey text, + event_created_at integer, + event_tags_json text, + event_content text, + event_sig text, + raw_event_json text, + outbox_status text not null check (outbox_status in ('none', 'pending', 'acknowledged', 'failed')), + relay_set_fingerprint text, + relay_delivery_json text, + check (change_seq >= 1), + check (trim(record_id) <> ''), + check (family <> 'local_work' or local_work_json is not null), + check (family <> 'local_work' or outbox_status = 'none'), + check (family <> 'signed_event' or (event_id is not null and event_kind is not null and event_pubkey is not null and event_sig is not null and raw_event_json is not null)) +); + +insert into runtime_store_record_network_source_next( + seq, + change_seq, + record_id, + family, + status, + source_runtime, + created_at_ms, + inserted_at_ms, + updated_at_ms, + owner_account_id, + owner_pubkey, + farm_id, + listing_addr, + local_work_json, + event_id, + event_kind, + event_pubkey, + event_created_at, + event_tags_json, + event_content, + event_sig, + raw_event_json, + outbox_status, + relay_set_fingerprint, + relay_delivery_json +) +select + seq, + change_seq, + record_id, + family, + status, + source_runtime, + created_at_ms, + inserted_at_ms, + updated_at_ms, + owner_account_id, + owner_pubkey, + farm_id, + listing_addr, + local_work_json, + event_id, + event_kind, + event_pubkey, + event_created_at, + event_tags_json, + event_content, + event_sig, + raw_event_json, + outbox_status, + relay_set_fingerprint, + relay_delivery_json +from runtime_store_record +order by seq asc; + +drop table runtime_store_record; +alter table runtime_store_record_network_source_next rename to runtime_store_record; + +create index runtime_store_record_change_seq_idx on runtime_store_record(change_seq); +create index runtime_store_record_event_id_idx on runtime_store_record(event_id); +create index runtime_store_record_listing_addr_idx on runtime_store_record(listing_addr); +create index runtime_store_record_owner_pubkey_idx on runtime_store_record(owner_pubkey); +create index runtime_store_record_status_idx on runtime_store_record(status); diff --git a/crates/runtime_store/src/error.rs b/crates/runtime_store/src/error.rs @@ -0,0 +1,14 @@ +#![forbid(unsafe_code)] + +use radroots_sql_core::error::SqlError; +use thiserror::Error; + +#[derive(Debug, Error)] +pub enum RuntimeStoreError { + #[error("invalid runtime store record: {0}")] + InvalidRecord(String), + #[error("sql error: {0}")] + Sql(#[from] SqlError), + #[error("serialization error: {0}")] + Serialization(#[from] serde_json::Error), +} diff --git a/crates/runtime_store/src/lib.rs b/crates/runtime_store/src/lib.rs @@ -0,0 +1,25 @@ +#![forbid(unsafe_code)] + +mod error; +mod migrations; +mod models; +mod order_work; +mod store; + +pub use error::RuntimeStoreError; +pub use migrations::{MIGRATIONS, run_all_down, run_all_up}; +pub use models::{ + PublishOutboxStatus, RelayDeliveryEvidence, RelayDeliveryFailure, RelayDeliveryState, + RuntimeStoreCursor, RuntimeStoreRecord, RuntimeStoreRecordFamily, RuntimeStoreRecordInput, + RuntimeStoreRecordStatus, RuntimeStoreRecordUpdate, SourceRuntime, +}; +pub use order_work::{ + BUYER_ORDER_REQUEST_ACTOR_SOURCE_RESOLVED_ACCOUNT, + BUYER_ORDER_REQUEST_ACTOR_SOURCE_UNRESOLVED_APP, BUYER_ORDER_REQUEST_DOCUMENT_KIND, + BUYER_ORDER_REQUEST_LOCAL_WORK_RECORD_KIND, BuyerOrderRequestLocalWorkValidation, + BuyerOrderRequestSupportState, buyer_order_request_local_work_record_id, + validate_buyer_order_request_local_work_payload, + validate_supported_buyer_order_request_local_work_payload, + validate_unsupported_buyer_order_request_local_work_payload, +}; +pub use store::RuntimeStore; diff --git a/crates/runtime_store/src/migrations.rs b/crates/runtime_store/src/migrations.rs @@ -0,0 +1,55 @@ +#![forbid(unsafe_code)] + +use radroots_sql_core::SqlExecutor; +use radroots_sql_core::error::SqlError; +use radroots_sql_core::migrations::{Migration, migrations_run_all_down, migrations_run_all_up}; + +pub static MIGRATIONS: &[Migration] = &[ + Migration { + name: "0000_runtime_store", + up_sql: include_str!("../migrations/0000_runtime_store.up.sql"), + down_sql: include_str!("../migrations/0000_runtime_store.down.sql"), + }, + Migration { + name: "0001_change_tracking", + up_sql: include_str!("../migrations/0001_change_tracking.up.sql"), + down_sql: include_str!("../migrations/0001_change_tracking.down.sql"), + }, + Migration { + name: "0002_network_source_runtime", + up_sql: include_str!("../migrations/0002_network_source_runtime.up.sql"), + down_sql: include_str!("../migrations/0002_network_source_runtime.down.sql"), + }, +]; + +pub fn run_all_up<E>(executor: &E) -> Result<(), SqlError> +where + E: SqlExecutor, +{ + migrations_run_all_up(executor, MIGRATIONS) +} + +pub fn run_all_down<E>(executor: &E) -> Result<(), SqlError> +where + E: SqlExecutor, +{ + migrations_run_all_down(executor, MIGRATIONS) +} + +#[cfg(test)] +mod tests { + use radroots_sql_core::SqliteExecutor; + + use super::*; + + #[test] + fn migration_entrypoints_apply_and_reverse_schema() { + let executor = SqliteExecutor::open_memory().expect("open memory sqlite"); + + run_all_up(&executor).expect("migrate up"); + executor + .query_raw("select name from __migrations order by name", "[]") + .expect("query migrations"); + run_all_down(&executor).expect("migrate down"); + } +} diff --git a/crates/runtime_store/src/models.rs b/crates/runtime_store/src/models.rs @@ -0,0 +1,757 @@ +#![forbid(unsafe_code)] + +use std::collections::BTreeSet; + +use serde::{Deserialize, Serialize}; +use serde_json::Value; + +use crate::RuntimeStoreError; + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum RuntimeStoreRecordFamily { + LocalWork, + SignedEvent, +} + +impl RuntimeStoreRecordFamily { + pub fn as_str(self) -> &'static str { + match self { + Self::LocalWork => "local_work", + Self::SignedEvent => "signed_event", + } + } + + pub fn parse(value: &str) -> Result<Self, RuntimeStoreError> { + match value { + "local_work" => Ok(Self::LocalWork), + "signed_event" => Ok(Self::SignedEvent), + other => Err(RuntimeStoreError::InvalidRecord(format!( + "unknown record family `{other}`" + ))), + } + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum RuntimeStoreRecordStatus { + LocalDraft, + LocalSaved, + PendingPublish, + Published, + Failed, + Conflict, +} + +impl RuntimeStoreRecordStatus { + pub fn as_str(self) -> &'static str { + match self { + Self::LocalDraft => "local_draft", + Self::LocalSaved => "local_saved", + Self::PendingPublish => "pending_publish", + Self::Published => "published", + Self::Failed => "failed", + Self::Conflict => "conflict", + } + } + + pub fn parse(value: &str) -> Result<Self, RuntimeStoreError> { + match value { + "local_draft" => Ok(Self::LocalDraft), + "local_saved" => Ok(Self::LocalSaved), + "pending_publish" => Ok(Self::PendingPublish), + "published" => Ok(Self::Published), + "failed" => Ok(Self::Failed), + "conflict" => Ok(Self::Conflict), + other => Err(RuntimeStoreError::InvalidRecord(format!( + "unknown record status `{other}`" + ))), + } + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum PublishOutboxStatus { + None, + Pending, + Acknowledged, + Failed, +} + +impl PublishOutboxStatus { + pub fn as_str(self) -> &'static str { + match self { + Self::None => "none", + Self::Pending => "pending", + Self::Acknowledged => "acknowledged", + Self::Failed => "failed", + } + } + + pub fn parse(value: &str) -> Result<Self, RuntimeStoreError> { + match value { + "none" => Ok(Self::None), + "pending" => Ok(Self::Pending), + "acknowledged" => Ok(Self::Acknowledged), + "failed" => Ok(Self::Failed), + other => Err(RuntimeStoreError::InvalidRecord(format!( + "unknown outbox status `{other}`" + ))), + } + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum RelayDeliveryState { + Pending, + Acknowledged, + Observed, + Failed, +} + +impl RelayDeliveryState { + pub fn as_str(self) -> &'static str { + match self { + Self::Pending => "pending", + Self::Acknowledged => "acknowledged", + Self::Observed => "observed", + Self::Failed => "failed", + } + } +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct RelayDeliveryFailure { + pub relay_url: String, + pub error: String, +} + +impl RelayDeliveryFailure { + pub fn new( + relay_url: impl Into<String>, + error: impl Into<String>, + ) -> Result<Self, RuntimeStoreError> { + let relay_url = relay_url.into(); + let error = error.into(); + validate_non_empty("relay_url", &relay_url)?; + validate_non_empty("relay_delivery_error", &error)?; + Ok(Self { relay_url, error }) + } +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct RelayDeliveryEvidence { + pub state: RelayDeliveryState, + #[serde(default)] + pub target_relays: Vec<String>, + #[serde(default)] + pub connected_relays: Vec<String>, + #[serde(default)] + pub acknowledged_relays: Vec<String>, + #[serde(default)] + pub observed_relays: Vec<String>, + #[serde(default)] + pub failed_relays: Vec<RelayDeliveryFailure>, +} + +impl RelayDeliveryEvidence { + pub fn pending<I, S>(target_relays: I) -> Result<Self, RuntimeStoreError> + where + I: IntoIterator<Item = S>, + S: AsRef<str>, + { + let target_relays = normalized_relay_values(target_relays)?; + if target_relays.is_empty() { + return Err(RuntimeStoreError::InvalidRecord( + "pending relay delivery evidence requires target_relays".to_owned(), + )); + } + Ok(Self { + state: RelayDeliveryState::Pending, + target_relays, + connected_relays: Vec::new(), + acknowledged_relays: Vec::new(), + observed_relays: Vec::new(), + failed_relays: Vec::new(), + }) + } + + pub fn acknowledged<T, C, A, TS, CS, AS>( + target_relays: T, + connected_relays: C, + acknowledged_relays: A, + failed_relays: Vec<RelayDeliveryFailure>, + ) -> Result<Self, RuntimeStoreError> + where + T: IntoIterator<Item = TS>, + C: IntoIterator<Item = CS>, + A: IntoIterator<Item = AS>, + TS: AsRef<str>, + CS: AsRef<str>, + AS: AsRef<str>, + { + let target_relays = normalized_relay_values(target_relays)?; + let connected_relays = normalized_relay_values(connected_relays)?; + let acknowledged_relays = normalized_relay_values(acknowledged_relays)?; + if target_relays.is_empty() { + return Err(RuntimeStoreError::InvalidRecord( + "acknowledged relay delivery evidence requires target_relays".to_owned(), + )); + } + if acknowledged_relays.is_empty() { + return Err(RuntimeStoreError::InvalidRecord( + "acknowledged relay delivery evidence requires acknowledged_relays".to_owned(), + )); + } + Ok(Self { + state: RelayDeliveryState::Acknowledged, + target_relays, + connected_relays, + acknowledged_relays, + observed_relays: Vec::new(), + failed_relays, + }) + } + + pub fn observed<T, C, O, TS, CS, OS>( + target_relays: T, + connected_relays: C, + observed_relays: O, + failed_relays: Vec<RelayDeliveryFailure>, + ) -> Result<Self, RuntimeStoreError> + where + T: IntoIterator<Item = TS>, + C: IntoIterator<Item = CS>, + O: IntoIterator<Item = OS>, + TS: AsRef<str>, + CS: AsRef<str>, + OS: AsRef<str>, + { + let target_relays = normalized_relay_values(target_relays)?; + let connected_relays = normalized_relay_values(connected_relays)?; + let observed_relays = normalized_relay_values(observed_relays)?; + if target_relays.is_empty() { + return Err(RuntimeStoreError::InvalidRecord( + "observed relay delivery evidence requires target_relays".to_owned(), + )); + } + if observed_relays.is_empty() { + return Err(RuntimeStoreError::InvalidRecord( + "observed relay delivery evidence requires observed_relays".to_owned(), + )); + } + Ok(Self { + state: RelayDeliveryState::Observed, + target_relays, + connected_relays, + acknowledged_relays: Vec::new(), + observed_relays, + failed_relays, + }) + } + + pub fn from_json_value(value: &Value) -> Result<Self, RuntimeStoreError> { + let evidence: Self = serde_json::from_value(value.clone())?; + evidence.validate()?; + Ok(evidence) + } + + pub fn to_json_value(&self) -> Result<Value, RuntimeStoreError> { + self.validate()?; + Ok(serde_json::to_value(self)?) + } + + pub fn relay_set_fingerprint(&self) -> Option<String> { + let relays = self + .target_relays + .iter() + .chain(self.connected_relays.iter()) + .chain(self.acknowledged_relays.iter()) + .chain(self.observed_relays.iter()) + .map(String::as_str) + .chain( + self.failed_relays + .iter() + .map(|failure| failure.relay_url.as_str()), + ) + .map(str::trim) + .filter(|relay| !relay.is_empty()) + .collect::<BTreeSet<_>>(); + if relays.is_empty() { + None + } else { + Some(relays.into_iter().collect::<Vec<_>>().join("\n")) + } + } + + fn validate(&self) -> Result<(), RuntimeStoreError> { + validate_relays("target_relays", &self.target_relays)?; + validate_relays("connected_relays", &self.connected_relays)?; + validate_relays("acknowledged_relays", &self.acknowledged_relays)?; + validate_relays("observed_relays", &self.observed_relays)?; + match self.state { + RelayDeliveryState::Pending => { + if self.target_relays.is_empty() { + return Err(RuntimeStoreError::InvalidRecord( + "pending relay delivery evidence requires target_relays".to_owned(), + )); + } + } + RelayDeliveryState::Acknowledged => { + if self.acknowledged_relays.is_empty() { + return Err(RuntimeStoreError::InvalidRecord( + "acknowledged relay delivery evidence requires acknowledged_relays" + .to_owned(), + )); + } + } + RelayDeliveryState::Observed => { + if self.observed_relays.is_empty() { + return Err(RuntimeStoreError::InvalidRecord( + "observed relay delivery evidence requires observed_relays".to_owned(), + )); + } + } + RelayDeliveryState::Failed => { + if self.failed_relays.is_empty() { + return Err(RuntimeStoreError::InvalidRecord( + "failed relay delivery evidence requires failed_relays".to_owned(), + )); + } + } + } + for failure in &self.failed_relays { + validate_non_empty("failed_relay_url", &failure.relay_url)?; + validate_non_empty("failed_relay_error", &failure.error)?; + } + Ok(()) + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum SourceRuntime { + Cli, + App, + Network, + Service, + Worker, + Test, +} + +impl SourceRuntime { + pub fn as_str(self) -> &'static str { + match self { + Self::Cli => "cli", + Self::App => "app", + Self::Network => "network", + Self::Service => "service", + Self::Worker => "worker", + Self::Test => "test", + } + } + + pub fn parse(value: &str) -> Result<Self, RuntimeStoreError> { + match value { + "cli" => Ok(Self::Cli), + "app" => Ok(Self::App), + "network" => Ok(Self::Network), + "service" => Ok(Self::Service), + "worker" => Ok(Self::Worker), + "test" => Ok(Self::Test), + other => Err(RuntimeStoreError::InvalidRecord(format!( + "unknown source runtime `{other}`" + ))), + } + } +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct RuntimeStoreRecordInput { + pub record_id: String, + pub family: RuntimeStoreRecordFamily, + pub status: RuntimeStoreRecordStatus, + pub source_runtime: SourceRuntime, + pub created_at_ms: i64, + pub inserted_at_ms: i64, + pub owner_account_id: Option<String>, + pub owner_pubkey: Option<String>, + pub farm_id: Option<String>, + pub listing_addr: Option<String>, + pub local_work_json: Option<Value>, + pub event_id: Option<String>, + pub event_kind: Option<i64>, + pub event_pubkey: Option<String>, + pub event_created_at: Option<i64>, + pub event_tags_json: Option<Value>, + pub event_content: Option<String>, + pub event_sig: Option<String>, + pub raw_event_json: Option<Value>, + pub outbox_status: PublishOutboxStatus, + pub relay_set_fingerprint: Option<String>, + pub relay_delivery_json: Option<Value>, +} + +impl RuntimeStoreRecordInput { + pub fn validate(&self) -> Result<(), RuntimeStoreError> { + validate_non_empty("record_id", &self.record_id)?; + if let Some(value) = self.owner_account_id.as_deref() { + validate_non_empty("owner_account_id", value)?; + } + if let Some(value) = self.owner_pubkey.as_deref() { + validate_non_empty("owner_pubkey", value)?; + } + if let Some(value) = self.farm_id.as_deref() { + validate_non_empty("farm_id", value)?; + } + if let Some(value) = self.listing_addr.as_deref() { + validate_non_empty("listing_addr", value)?; + } + if let Some(value) = self.relay_set_fingerprint.as_deref() { + validate_non_empty("relay_set_fingerprint", value)?; + } + if let Some(value) = self.relay_delivery_json.as_ref() { + RelayDeliveryEvidence::from_json_value(value)?; + } + match self.family { + RuntimeStoreRecordFamily::LocalWork => { + if self.local_work_json.is_none() { + return Err(RuntimeStoreError::InvalidRecord( + "local work records require local_work_json".to_owned(), + )); + } + if self.outbox_status != PublishOutboxStatus::None { + return Err(RuntimeStoreError::InvalidRecord( + "local work records must use outbox status none".to_owned(), + )); + } + } + RuntimeStoreRecordFamily::SignedEvent => { + validate_required("event_id", self.event_id.as_deref())?; + validate_required("event_pubkey", self.event_pubkey.as_deref())?; + validate_required("event_sig", self.event_sig.as_deref())?; + if self.event_kind.is_none() { + return Err(RuntimeStoreError::InvalidRecord( + "signed event records require event_kind".to_owned(), + )); + } + if self.raw_event_json.is_none() { + return Err(RuntimeStoreError::InvalidRecord( + "signed event records require raw_event_json".to_owned(), + )); + } + } + } + Ok(()) + } +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct RuntimeStoreRecord { + pub seq: i64, + pub change_seq: i64, + pub record_id: String, + pub family: RuntimeStoreRecordFamily, + pub status: RuntimeStoreRecordStatus, + pub source_runtime: SourceRuntime, + pub created_at_ms: i64, + pub inserted_at_ms: i64, + pub updated_at_ms: i64, + pub owner_account_id: Option<String>, + pub owner_pubkey: Option<String>, + pub farm_id: Option<String>, + pub listing_addr: Option<String>, + pub local_work_json: Option<Value>, + pub event_id: Option<String>, + pub event_kind: Option<i64>, + pub event_pubkey: Option<String>, + pub event_created_at: Option<i64>, + pub event_tags_json: Option<Value>, + pub event_content: Option<String>, + pub event_sig: Option<String>, + pub raw_event_json: Option<Value>, + pub outbox_status: PublishOutboxStatus, + pub relay_set_fingerprint: Option<String>, + pub relay_delivery_json: Option<Value>, +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct RuntimeStoreRecordUpdate { + pub record_id: String, + pub status: RuntimeStoreRecordStatus, + pub outbox_status: PublishOutboxStatus, + pub relay_set_fingerprint: Option<String>, + pub relay_delivery_json: Option<Value>, + pub updated_at_ms: i64, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct RuntimeStoreCursor { + pub consumer_id: String, + pub last_change_seq: i64, + pub updated_at_ms: i64, +} + +pub(crate) fn validate_non_empty(field: &str, value: &str) -> Result<(), RuntimeStoreError> { + if value.trim().is_empty() { + return Err(RuntimeStoreError::InvalidRecord(format!( + "{field} must not be empty" + ))); + } + Ok(()) +} + +fn normalized_relay_values<I, S>(values: I) -> Result<Vec<String>, RuntimeStoreError> +where + I: IntoIterator<Item = S>, + S: AsRef<str>, +{ + let mut seen = BTreeSet::new(); + let mut normalized = Vec::new(); + for value in values { + let value = value.as_ref().trim(); + validate_non_empty("relay_url", value)?; + if seen.insert(value.to_owned()) { + normalized.push(value.to_owned()); + } + } + Ok(normalized) +} + +fn validate_relays(field: &str, relays: &[String]) -> Result<(), RuntimeStoreError> { + for relay in relays { + validate_non_empty(field, relay)?; + } + Ok(()) +} + +fn validate_required(field: &str, value: Option<&str>) -> Result<(), RuntimeStoreError> { + match value { + Some(value) => validate_non_empty(field, value), + None => Err(RuntimeStoreError::InvalidRecord(format!( + "{field} is required" + ))), + } +} + +#[cfg(test)] +mod tests { + use serde_json::json; + + use super::*; + + #[test] + fn enum_strings_and_parse_errors_cover_all_model_variants() { + for (variant, value) in [ + (RuntimeStoreRecordFamily::LocalWork, "local_work"), + (RuntimeStoreRecordFamily::SignedEvent, "signed_event"), + ] { + assert_eq!(variant.as_str(), value); + assert_eq!( + RuntimeStoreRecordFamily::parse(value).expect("record family"), + variant + ); + } + + for (variant, value) in [ + (RuntimeStoreRecordStatus::LocalDraft, "local_draft"), + (RuntimeStoreRecordStatus::LocalSaved, "local_saved"), + (RuntimeStoreRecordStatus::PendingPublish, "pending_publish"), + (RuntimeStoreRecordStatus::Published, "published"), + (RuntimeStoreRecordStatus::Failed, "failed"), + (RuntimeStoreRecordStatus::Conflict, "conflict"), + ] { + assert_eq!(variant.as_str(), value); + assert_eq!( + RuntimeStoreRecordStatus::parse(value).expect("record status"), + variant + ); + } + + for (variant, value) in [ + (PublishOutboxStatus::None, "none"), + (PublishOutboxStatus::Pending, "pending"), + (PublishOutboxStatus::Acknowledged, "acknowledged"), + (PublishOutboxStatus::Failed, "failed"), + ] { + assert_eq!(variant.as_str(), value); + assert_eq!( + PublishOutboxStatus::parse(value).expect("outbox status"), + variant + ); + } + + for (variant, value) in [ + (SourceRuntime::Cli, "cli"), + (SourceRuntime::App, "app"), + (SourceRuntime::Network, "network"), + (SourceRuntime::Service, "service"), + (SourceRuntime::Worker, "worker"), + (SourceRuntime::Test, "test"), + ] { + assert_eq!(variant.as_str(), value); + assert_eq!( + SourceRuntime::parse(value).expect("source runtime"), + variant + ); + } + + assert!(RuntimeStoreRecordFamily::parse("other").is_err()); + assert!(RuntimeStoreRecordStatus::parse("other").is_err()); + assert!(PublishOutboxStatus::parse("other").is_err()); + assert!(SourceRuntime::parse("other").is_err()); + } + + #[test] + fn local_record_input_validation_covers_success_and_error_paths() { + let mut local_work = local_work_input(); + local_work.validate().expect("valid local work"); + + for (field, update) in [ + ( + "owner_account_id", + Box::new(|input: &mut RuntimeStoreRecordInput| { + input.owner_account_id = Some(" ".to_owned()); + }) as Box<dyn Fn(&mut RuntimeStoreRecordInput)>, + ), + ( + "owner_pubkey", + Box::new(|input: &mut RuntimeStoreRecordInput| { + input.owner_pubkey = Some(" ".to_owned()); + }), + ), + ( + "farm_id", + Box::new(|input: &mut RuntimeStoreRecordInput| { + input.farm_id = Some(" ".to_owned()); + }), + ), + ( + "listing_addr", + Box::new(|input: &mut RuntimeStoreRecordInput| { + input.listing_addr = Some(" ".to_owned()); + }), + ), + ] { + let mut input = local_work_input(); + update(&mut input); + assert_error_contains(input.validate(), field); + } + + local_work.record_id = " ".to_owned(); + assert_error_contains(local_work.validate(), "record_id"); + + let mut missing_work = local_work_input(); + missing_work.local_work_json = None; + assert_error_contains(missing_work.validate(), "local_work_json"); + + let mut queued_work = local_work_input(); + queued_work.outbox_status = PublishOutboxStatus::Pending; + assert_error_contains(queued_work.validate(), "outbox status none"); + + let signed_event = signed_event_input(); + signed_event.validate().expect("valid signed event"); + + for (field, update) in [ + ( + "event_id", + Box::new(|input: &mut RuntimeStoreRecordInput| { + input.event_id = Some(" ".to_owned()); + }) as Box<dyn Fn(&mut RuntimeStoreRecordInput)>, + ), + ( + "event_pubkey", + Box::new(|input: &mut RuntimeStoreRecordInput| { + input.event_pubkey = None; + }), + ), + ( + "event_sig", + Box::new(|input: &mut RuntimeStoreRecordInput| { + input.event_sig = None; + }), + ), + ( + "event_kind", + Box::new(|input: &mut RuntimeStoreRecordInput| { + input.event_kind = None; + }), + ), + ( + "raw_event_json", + Box::new(|input: &mut RuntimeStoreRecordInput| { + input.raw_event_json = None; + }), + ), + ] { + let mut input = signed_event_input(); + update(&mut input); + assert_error_contains(input.validate(), field); + } + } + + fn local_work_input() -> RuntimeStoreRecordInput { + RuntimeStoreRecordInput { + record_id: "local-work-a".to_owned(), + family: RuntimeStoreRecordFamily::LocalWork, + status: RuntimeStoreRecordStatus::LocalSaved, + source_runtime: SourceRuntime::App, + created_at_ms: 10, + inserted_at_ms: 11, + owner_account_id: Some("account-a".to_owned()), + owner_pubkey: Some("pubkey-a".to_owned()), + farm_id: Some("farm-a".to_owned()), + listing_addr: Some("listing-a".to_owned()), + local_work_json: Some(json!({"kind":"buyer_order_request_v1"})), + event_id: None, + event_kind: None, + event_pubkey: None, + event_created_at: None, + event_tags_json: None, + event_content: None, + event_sig: None, + raw_event_json: None, + outbox_status: PublishOutboxStatus::None, + relay_set_fingerprint: None, + relay_delivery_json: None, + } + } + + fn signed_event_input() -> RuntimeStoreRecordInput { + RuntimeStoreRecordInput { + record_id: "signed-event-a".to_owned(), + family: RuntimeStoreRecordFamily::SignedEvent, + status: RuntimeStoreRecordStatus::PendingPublish, + source_runtime: SourceRuntime::Service, + created_at_ms: 20, + inserted_at_ms: 21, + owner_account_id: None, + owner_pubkey: None, + farm_id: None, + listing_addr: None, + local_work_json: None, + event_id: Some("event-a".to_owned()), + event_kind: Some(30402), + event_pubkey: Some("pubkey-a".to_owned()), + event_created_at: Some(20), + event_tags_json: Some(json!([["d", "listing-a"]])), + event_content: Some("{}".to_owned()), + event_sig: Some("sig-a".to_owned()), + raw_event_json: Some(json!({"id":"event-a"})), + outbox_status: PublishOutboxStatus::Pending, + relay_set_fingerprint: None, + relay_delivery_json: None, + } + } + + fn assert_error_contains(result: Result<(), RuntimeStoreError>, expected: &str) { + let err = result.expect_err("validation error"); + assert!( + err.to_string().contains(expected), + "expected error to contain {expected}, got {err}" + ); + } +} diff --git a/crates/runtime_store/src/order_work.rs b/crates/runtime_store/src/order_work.rs @@ -0,0 +1,806 @@ +use serde_json::Value; + +use crate::RuntimeStoreError; +use crate::models::validate_non_empty; + +pub const BUYER_ORDER_REQUEST_LOCAL_WORK_RECORD_KIND: &str = "buyer_order_request_v1"; +pub const BUYER_ORDER_REQUEST_DOCUMENT_KIND: &str = "order_draft_v1"; +pub const BUYER_ORDER_REQUEST_ACTOR_SOURCE_RESOLVED_ACCOUNT: &str = "resolved_account"; +pub const BUYER_ORDER_REQUEST_ACTOR_SOURCE_UNRESOLVED_APP: &str = "app_unresolved"; + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum BuyerOrderRequestSupportState { + Supported, + Unsupported, +} + +impl BuyerOrderRequestSupportState { + pub fn as_str(self) -> &'static str { + match self { + Self::Supported => "supported", + Self::Unsupported => "unsupported", + } + } +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct BuyerOrderRequestLocalWorkValidation { + pub order_id: String, + pub support_state: BuyerOrderRequestSupportState, + pub support_issues: Vec<String>, +} + +pub fn buyer_order_request_local_work_record_id( + order_id: &str, +) -> Result<String, RuntimeStoreError> { + let order_id = order_id.trim(); + validate_non_empty("order_id", order_id)?; + Ok(format!("app:local_work:order_request:{order_id}")) +} + +pub fn validate_buyer_order_request_local_work_payload( + payload: &Value, +) -> Result<BuyerOrderRequestLocalWorkValidation, RuntimeStoreError> { + validate_string_field( + payload, + &["record_kind"], + BUYER_ORDER_REQUEST_LOCAL_WORK_RECORD_KIND, + )?; + validate_string_field(payload, &["scope"], "app")?; + validate_string_field( + payload, + &["document", "kind"], + BUYER_ORDER_REQUEST_DOCUMENT_KIND, + )?; + validate_bool_field(payload, &["currentness", "current"], true)?; + validate_string_field(payload, &["currentness", "source"], "app_sqlite_order")?; + + let order_id = validate_required_string(payload, &["document", "order", "order_id"])?; + let currentness_order_id = validate_required_string(payload, &["currentness", "order_id"])?; + if currentness_order_id != order_id { + return Err(invalid_field( + "currentness.order_id", + "must match document.order.order_id", + )); + } + validate_required_string(payload, &["currentness", "record_id"])?; + validate_positive_i64(payload, &["currentness", "created_at_ms"])?; + validate_required_string(payload, &["currentness", "order_updated_at"])?; + + let (support_state, support_issues) = validate_support_status(payload)?; + validate_exportability(payload, support_state)?; + validate_order_identity(payload, support_state)?; + validate_order_items(payload)?; + validate_order_economics(payload)?; + + Ok(BuyerOrderRequestLocalWorkValidation { + order_id: order_id.to_owned(), + support_state, + support_issues, + }) +} + +pub fn validate_supported_buyer_order_request_local_work_payload( + payload: &Value, +) -> Result<BuyerOrderRequestLocalWorkValidation, RuntimeStoreError> { + let validation = validate_buyer_order_request_local_work_payload(payload)?; + if validation.support_state != BuyerOrderRequestSupportState::Supported { + return Err(invalid_field( + "support_status.state", + "must be supported for exportable app order work", + )); + } + Ok(validation) +} + +pub fn validate_unsupported_buyer_order_request_local_work_payload( + payload: &Value, +) -> Result<BuyerOrderRequestLocalWorkValidation, RuntimeStoreError> { + let validation = validate_buyer_order_request_local_work_payload(payload)?; + if validation.support_state != BuyerOrderRequestSupportState::Unsupported { + return Err(invalid_field( + "support_status.state", + "must be unsupported for unsupported app order work", + )); + } + Ok(validation) +} + +fn validate_support_status( + payload: &Value, +) -> Result<(BuyerOrderRequestSupportState, Vec<String>), RuntimeStoreError> { + let state = validate_required_string(payload, &["support_status", "state"])?; + let issues = support_issues(payload)?; + match state { + "supported" => { + if !issues.is_empty() { + return Err(invalid_field( + "support_status.issues", + "must be empty when support_status.state is supported", + )); + } + Ok((BuyerOrderRequestSupportState::Supported, issues)) + } + "unsupported" => { + if issues.is_empty() { + return Err(invalid_field( + "support_status.issues", + "must contain at least one issue when support_status.state is unsupported", + )); + } + Ok((BuyerOrderRequestSupportState::Unsupported, issues)) + } + _ => Err(invalid_field( + "support_status.state", + "must be supported or unsupported", + )), + } +} + +fn validate_exportability( + payload: &Value, + support_state: BuyerOrderRequestSupportState, +) -> Result<(), RuntimeStoreError> { + let state = validate_required_string(payload, &["exportability", "state"])?; + match state { + "exportable" => { + validate_string_field( + payload, + &["document", "buyer_actor", "source"], + BUYER_ORDER_REQUEST_ACTOR_SOURCE_RESOLVED_ACCOUNT, + )?; + validate_buyer_pubkey(payload)?; + } + "identity_unresolved" => { + validate_required_string(payload, &["exportability", "reason"])?; + validate_string_field( + payload, + &["document", "buyer_actor", "source"], + BUYER_ORDER_REQUEST_ACTOR_SOURCE_UNRESOLVED_APP, + )?; + if support_state == BuyerOrderRequestSupportState::Supported { + return Err(invalid_field( + "exportability.state", + "supported app order work must be exportable", + )); + } + } + _ => { + return Err(invalid_field( + "exportability.state", + "must be exportable or identity_unresolved", + )); + } + } + Ok(()) +} + +fn validate_order_identity( + payload: &Value, + support_state: BuyerOrderRequestSupportState, +) -> Result<(), RuntimeStoreError> { + validate_required_string(payload, &["document", "order", "listing_addr"])?; + validate_required_string(payload, &["document", "order", "listing_event_id"])?; + validate_required_string(payload, &["document", "order", "seller_pubkey"])?; + if support_state == BuyerOrderRequestSupportState::Supported { + validate_buyer_pubkey(payload)?; + } + Ok(()) +} + +fn validate_buyer_pubkey(payload: &Value) -> Result<(), RuntimeStoreError> { + let order_buyer_pubkey = + validate_required_string(payload, &["document", "order", "buyer_pubkey"])?; + let actor_buyer_pubkey = + validate_required_string(payload, &["document", "buyer_actor", "pubkey"])?; + if order_buyer_pubkey != actor_buyer_pubkey { + return Err(invalid_field( + "document.buyer_actor.pubkey", + "must match document.order.buyer_pubkey", + )); + } + Ok(()) +} + +fn validate_order_items(payload: &Value) -> Result<(), RuntimeStoreError> { + let items = required_array(payload, &["document", "order", "items"])?; + if items.is_empty() { + return Err(invalid_field( + "document.order.items", + "must contain at least one item", + )); + } + for (index, item) in items.iter().enumerate() { + validate_required_string(item, &["bin_id"]).map_err(|_| { + invalid_field_at( + format!("document.order.items[{index}].bin_id"), + "is required", + ) + })?; + validate_positive_u64(item, &["bin_count"]).map_err(|_| { + invalid_field_at( + format!("document.order.items[{index}].bin_count"), + "must be positive", + ) + })?; + } + Ok(()) +} + +fn validate_order_economics(payload: &Value) -> Result<(), RuntimeStoreError> { + let economics = value_at(payload, &["document", "order", "economics"]).ok_or_else(|| { + invalid_field("document.order.economics", "is required for app order work") + })?; + if !economics.is_object() { + return Err(invalid_field( + "document.order.economics", + "must be an object", + )); + } + validate_string_field(economics, &["pricing_basis"], "listing_event")?; + let currency = validate_required_string(economics, &["currency"])?; + validate_currency("document.order.economics.currency", currency)?; + let economics_items = required_array(economics, &["items"])?; + let order_items = required_array(payload, &["document", "order", "items"])?; + if economics_items.is_empty() { + return Err(invalid_field( + "document.order.economics.items", + "must contain at least one item", + )); + } + if economics_items.len() != order_items.len() { + return Err(invalid_field( + "document.order.economics.items", + "must match document.order.items length", + )); + } + for (index, item) in economics_items.iter().enumerate() { + let order_item = &order_items[index]; + let economics_bin_id = validate_required_string(item, &["bin_id"]).map_err(|_| { + invalid_field_at( + format!("document.order.economics.items[{index}].bin_id"), + "is required", + ) + })?; + let order_bin_id = validate_required_string(order_item, &["bin_id"])?; + if economics_bin_id != order_bin_id { + return Err(invalid_field_at( + format!("document.order.economics.items[{index}].bin_id"), + "must match document.order.items bin_id", + )); + } + let economics_bin_count = validate_positive_u64(item, &["bin_count"]).map_err(|_| { + invalid_field_at( + format!("document.order.economics.items[{index}].bin_count"), + "must be positive", + ) + })?; + let order_bin_count = validate_positive_u64(order_item, &["bin_count"])?; + if economics_bin_count != order_bin_count { + return Err(invalid_field_at( + format!("document.order.economics.items[{index}].bin_count"), + "must match document.order.items bin_count", + )); + } + validate_required_string(item, &["quantity_amount"]).map_err(|_| { + invalid_field_at( + format!("document.order.economics.items[{index}].quantity_amount"), + "is required", + ) + })?; + validate_required_string(item, &["quantity_unit"]).map_err(|_| { + invalid_field_at( + format!("document.order.economics.items[{index}].quantity_unit"), + "is required", + ) + })?; + validate_required_string(item, &["unit_price_amount"]).map_err(|_| { + invalid_field_at( + format!("document.order.economics.items[{index}].unit_price_amount"), + "is required", + ) + })?; + let unit_price_currency = validate_required_string(item, &["unit_price_currency"])?; + if unit_price_currency != currency { + return Err(invalid_field_at( + format!("document.order.economics.items[{index}].unit_price_currency"), + "must match document.order.economics.currency", + )); + } + validate_money(item, &["line_subtotal"], currency)?; + } + validate_money(economics, &["subtotal"], currency)?; + validate_money(economics, &["discount_total"], currency)?; + validate_money(economics, &["adjustment_total"], currency)?; + validate_money(economics, &["total"], currency)?; + Ok(()) +} + +fn validate_money(payload: &Value, path: &[&str], currency: &str) -> Result<(), RuntimeStoreError> { + let Some(money) = value_at(payload, path) else { + return Err(missing_field(path)); + }; + validate_required_string(money, &["amount"])?; + let money_currency = validate_required_string(money, &["currency"])?; + if money_currency != currency { + return Err(invalid_field( + &format!("{}.currency", path.join(".")), + "must match currency", + )); + } + Ok(()) +} + +fn validate_string_field( + payload: &Value, + path: &[&str], + expected: &str, +) -> Result<(), RuntimeStoreError> { + let Some(value) = value_at(payload, path).and_then(Value::as_str) else { + return Err(missing_field(path)); + }; + if value != expected { + return Err(invalid_field( + &path.join("."), + &format!("must be `{expected}`"), + )); + } + Ok(()) +} + +fn validate_required_string<'a>( + payload: &'a Value, + path: &[&str], +) -> Result<&'a str, RuntimeStoreError> { + let Some(value) = value_at(payload, path).and_then(Value::as_str) else { + return Err(missing_field(path)); + }; + validate_non_empty(&path.join("."), value)?; + Ok(value.trim()) +} + +fn validate_bool_field( + payload: &Value, + path: &[&str], + expected: bool, +) -> Result<(), RuntimeStoreError> { + let Some(value) = value_at(payload, path).and_then(Value::as_bool) else { + return Err(missing_field(path)); + }; + if value != expected { + return Err(invalid_field( + &path.join("."), + &format!("must be `{expected}`"), + )); + } + Ok(()) +} + +fn validate_positive_i64(payload: &Value, path: &[&str]) -> Result<(), RuntimeStoreError> { + match value_at(payload, path).and_then(Value::as_i64) { + Some(value) if value > 0 => Ok(()), + _ => Err(invalid_field(&path.join("."), "must be positive")), + } +} + +fn validate_positive_u64(payload: &Value, path: &[&str]) -> Result<u64, RuntimeStoreError> { + match value_at(payload, path).and_then(Value::as_u64) { + Some(value) if value > 0 => Ok(value), + _ => Err(invalid_field(&path.join("."), "must be positive")), + } +} + +fn validate_currency(field: &str, value: &str) -> Result<(), RuntimeStoreError> { + if value.len() != 3 || !value.bytes().all(|byte| byte.is_ascii_uppercase()) { + return Err(invalid_field( + field, + "must be an uppercase ISO currency code", + )); + } + Ok(()) +} + +fn required_array<'a>( + payload: &'a Value, + path: &[&str], +) -> Result<&'a Vec<Value>, RuntimeStoreError> { + let Some(value) = value_at(payload, path).and_then(Value::as_array) else { + return Err(missing_field(path)); + }; + Ok(value) +} + +fn support_issues(payload: &Value) -> Result<Vec<String>, RuntimeStoreError> { + let issues = required_array(payload, &["support_status", "issues"])?; + let mut parsed = Vec::with_capacity(issues.len()); + for (index, issue) in issues.iter().enumerate() { + let Some(issue) = issue.as_str() else { + return Err(invalid_field_at( + format!("support_status.issues[{index}]"), + "must be a string", + )); + }; + validate_non_empty("support_status.issues", issue)?; + parsed.push(issue.trim().to_owned()); + } + Ok(parsed) +} + +fn value_at<'a>(payload: &'a Value, path: &[&str]) -> Option<&'a Value> { + let mut current = payload; + for part in path { + current = current.get(*part)?; + } + Some(current) +} + +fn missing_field(path: &[&str]) -> RuntimeStoreError { + invalid_field(&path.join("."), "is required") +} + +fn invalid_field(field: &str, requirement: &str) -> RuntimeStoreError { + RuntimeStoreError::InvalidRecord(format!("local order field `{field}` {requirement}")) +} + +fn invalid_field_at(field: String, requirement: &str) -> RuntimeStoreError { + RuntimeStoreError::InvalidRecord(format!("local order field `{field}` {requirement}")) +} + +#[cfg(test)] +mod tests { + use serde_json::{Value, json}; + + use super::*; + + #[test] + fn support_state_labels_and_record_id_validation_are_stable() { + assert_eq!( + BuyerOrderRequestSupportState::Supported.as_str(), + "supported" + ); + assert_eq!( + BuyerOrderRequestSupportState::Unsupported.as_str(), + "unsupported" + ); + assert_eq!( + buyer_order_request_local_work_record_id(" ord-a ").expect("record id"), + "app:local_work:order_request:ord-a" + ); + assert_error_contains( + buyer_order_request_local_work_record_id(" "), + "order_id must not be empty", + ); + } + + #[test] + fn private_validation_helpers_cover_successful_payload() { + let payload = supported_payload(); + + assert_eq!( + validate_support_status(&payload).expect("support status"), + ( + BuyerOrderRequestSupportState::Supported, + Vec::<String>::new() + ) + ); + validate_supported_buyer_order_request_local_work_payload(&payload) + .expect("supported payload"); + validate_exportability(&payload, BuyerOrderRequestSupportState::Supported) + .expect("exportability"); + validate_order_identity(&payload, BuyerOrderRequestSupportState::Supported) + .expect("identity"); + validate_order_items(&payload).expect("items"); + validate_order_economics(&payload).expect("economics"); + assert_eq!( + validate_required_string(&payload, &["document", "order", "order_id"]) + .expect("order id"), + "ord_1" + ); + validate_bool_field(&payload, &["currentness", "current"], true).expect("bool"); + assert_eq!( + support_issues(&payload).expect("support issues"), + Vec::<String>::new() + ); + assert!(value_at(&payload, &["document", "order"]).is_some()); + } + + #[test] + fn payload_validation_rejects_top_level_contract_drift() { + let mut wrong_kind = supported_payload(); + wrong_kind["record_kind"] = json!("other"); + assert_invalid(wrong_kind, "record_kind"); + + let mut missing_scope = supported_payload(); + missing_scope["scope"] = Value::Null; + assert_invalid(missing_scope, "scope"); + + let mut wrong_document_kind = supported_payload(); + wrong_document_kind["document"]["kind"] = json!("other"); + assert_invalid(wrong_document_kind, "document.kind"); + + let mut wrong_currentness_source = supported_payload(); + wrong_currentness_source["currentness"]["source"] = json!("other"); + assert_invalid(wrong_currentness_source, "currentness.source"); + + let mut missing_order_updated = supported_payload(); + missing_order_updated["currentness"]["order_updated_at"] = Value::Null; + assert_invalid(missing_order_updated, "order_updated_at"); + + let mut bad_created_at = supported_payload(); + bad_created_at["currentness"]["created_at_ms"] = json!(0); + assert_invalid(bad_created_at, "created_at_ms"); + } + + #[test] + fn support_and_exportability_rejections_cover_private_branches() { + let mut invalid_state = supported_payload(); + invalid_state["support_status"]["state"] = json!("partial"); + assert_invalid(invalid_state, "support_status.state"); + + let mut issue_not_string = supported_payload(); + issue_not_string["support_status"] = json!({ + "state": "unsupported", + "issues": [42] + }); + assert_invalid(issue_not_string, "support_status.issues[0]"); + + let mut issue_empty = supported_payload(); + issue_empty["support_status"] = json!({ + "state": "unsupported", + "issues": [" "] + }); + assert_invalid(issue_empty, "support_status.issues"); + + let mut supported_but_unresolved = unsupported_payload(); + supported_but_unresolved["support_status"] = json!({ + "state": "supported", + "issues": [] + }); + assert_invalid(supported_but_unresolved, "exportability.state"); + + let mut unknown_exportability = supported_payload(); + unknown_exportability["exportability"]["state"] = json!("queued"); + assert_invalid(unknown_exportability, "exportability.state"); + + let mut missing_reason = unsupported_payload(); + missing_reason["exportability"]["reason"] = Value::Null; + assert_invalid(missing_reason, "exportability.reason"); + + let mut wrong_actor_source = unsupported_payload(); + wrong_actor_source["document"]["buyer_actor"]["source"] = + json!(BUYER_ORDER_REQUEST_ACTOR_SOURCE_RESOLVED_ACCOUNT); + assert_invalid(wrong_actor_source, "buyer_actor.source"); + + let mut mismatched_buyer = supported_payload(); + mismatched_buyer["document"]["buyer_actor"]["pubkey"] = json!("other"); + assert_invalid(mismatched_buyer, "buyer_actor.pubkey"); + + let supported_error = + validate_unsupported_buyer_order_request_local_work_payload(&supported_payload()) + .expect_err("supported payload is not unsupported"); + assert!(supported_error.to_string().contains("support_status.state")); + } + + #[test] + fn item_and_economics_rejections_cover_private_branches() { + let mut economics_not_object = supported_payload(); + economics_not_object["document"]["order"]["economics"] = json!("bad"); + assert_invalid(economics_not_object, "economics"); + + let mut bad_pricing_basis = supported_payload(); + bad_pricing_basis["document"]["order"]["economics"]["pricing_basis"] = json!("manual"); + assert_invalid(bad_pricing_basis, "pricing_basis"); + + let mut bad_currency = supported_payload(); + bad_currency["document"]["order"]["economics"]["currency"] = json!("usd"); + assert_invalid(bad_currency, "currency"); + + let mut bad_currency_length = supported_payload(); + bad_currency_length["document"]["order"]["economics"]["currency"] = json!("US"); + assert_invalid(bad_currency_length, "currency"); + + let mut missing_economics = supported_payload(); + missing_economics["document"]["order"] + .as_object_mut() + .expect("order object") + .remove("economics"); + assert_invalid(missing_economics, "economics"); + + let mut economics_items_missing = supported_payload(); + economics_items_missing["document"]["order"]["economics"]["items"] = Value::Null; + assert_invalid(economics_items_missing, "items"); + + let mut economics_items_short = supported_payload(); + economics_items_short["document"]["order"]["economics"]["items"] = json!([]); + assert_invalid(economics_items_short, "economics.items"); + + let mut economics_items_long = supported_payload(); + economics_items_long["document"]["order"]["economics"]["items"] = json!([ + { + "bin_id": "dozen-eggs", + "bin_count": 2, + "quantity_amount": "1", + "quantity_unit": "dozen", + "unit_price_amount": "8.00", + "unit_price_currency": "USD", + "line_subtotal": { + "amount": "16.00", + "currency": "USD" + } + }, + { + "bin_id": "half-dozen-eggs", + "bin_count": 1 + } + ]); + assert_invalid(economics_items_long, "economics.items"); + + let mut economics_bin_missing = supported_payload(); + economics_bin_missing["document"]["order"]["economics"]["items"][0]["bin_id"] = Value::Null; + assert_invalid(economics_bin_missing, "economics.items[0].bin_id"); + + let mut economics_count_bad = supported_payload(); + economics_count_bad["document"]["order"]["economics"]["items"][0]["bin_count"] = json!(0); + assert_invalid(economics_count_bad, "economics.items[0].bin_count"); + + let mut order_count_mismatch = supported_payload(); + order_count_mismatch["document"]["order"]["economics"]["items"][0]["bin_count"] = json!(3); + assert_invalid(order_count_mismatch, "economics.items[0].bin_count"); + + let mut quantity_amount_missing = supported_payload(); + quantity_amount_missing["document"]["order"]["economics"]["items"][0]["quantity_amount"] = + Value::Null; + assert_invalid(quantity_amount_missing, "quantity_amount"); + + let mut quantity_unit_missing = supported_payload(); + quantity_unit_missing["document"]["order"]["economics"]["items"][0]["quantity_unit"] = + Value::Null; + assert_invalid(quantity_unit_missing, "quantity_unit"); + + let mut unit_price_amount_missing = supported_payload(); + unit_price_amount_missing["document"]["order"]["economics"]["items"][0]["unit_price_amount"] = + Value::Null; + assert_invalid(unit_price_amount_missing, "unit_price_amount"); + + let mut line_subtotal_missing = supported_payload(); + line_subtotal_missing["document"]["order"]["economics"]["items"][0]["line_subtotal"] = + Value::Null; + assert_invalid(line_subtotal_missing, "amount"); + + let mut missing_line_subtotal = supported_payload(); + missing_line_subtotal["document"]["order"]["economics"]["items"][0] + .as_object_mut() + .expect("economics item") + .remove("line_subtotal"); + assert_invalid(missing_line_subtotal, "line_subtotal"); + + let mut line_subtotal_currency = supported_payload(); + line_subtotal_currency["document"]["order"]["economics"]["items"][0]["line_subtotal"]["currency"] = + json!("CAD"); + assert_invalid(line_subtotal_currency, "line_subtotal.currency"); + + let mut subtotal_currency = supported_payload(); + subtotal_currency["document"]["order"]["economics"]["subtotal"]["currency"] = json!("CAD"); + assert_invalid(subtotal_currency, "subtotal.currency"); + + let mut order_item_missing = supported_payload(); + order_item_missing["document"]["order"]["items"] = Value::Null; + assert_invalid(order_item_missing, "document.order.items"); + + let mut missing_order_bin = supported_payload(); + missing_order_bin["document"]["order"]["items"][0]["bin_id"] = Value::Null; + assert_error_contains(validate_order_items(&missing_order_bin), "items[0].bin_id"); + } + + fn supported_payload() -> Value { + json!({ + "record_kind": BUYER_ORDER_REQUEST_LOCAL_WORK_RECORD_KIND, + "scope": "app", + "exportability": { + "state": "exportable" + }, + "support_status": { + "state": "supported", + "issues": [] + }, + "currentness": { + "current": true, + "source": "app_sqlite_order", + "record_id": "app:local_work:order_request:ord_1", + "order_id": "ord_1", + "order_updated_at": "2026-05-24T12:00:00Z", + "created_at_ms": 1777777777000_i64 + }, + "document": { + "kind": BUYER_ORDER_REQUEST_DOCUMENT_KIND, + "order": { + "order_id": "ord_1", + "listing_addr": "30402:seller_pubkey:listing_key", + "listing_event_id": "event-listing-1", + "buyer_pubkey": "buyer_pubkey", + "seller_pubkey": "seller_pubkey", + "items": [ + { + "bin_id": "dozen-eggs", + "bin_count": 2 + } + ], + "economics": { + "pricing_basis": "listing_event", + "currency": "USD", + "items": [ + { + "bin_id": "dozen-eggs", + "bin_count": 2, + "quantity_amount": "1", + "quantity_unit": "dozen", + "unit_price_amount": "8.00", + "unit_price_currency": "USD", + "line_subtotal": { + "amount": "16.00", + "currency": "USD" + } + } + ], + "subtotal": { + "amount": "16.00", + "currency": "USD" + }, + "discount_total": { + "amount": "0", + "currency": "USD" + }, + "adjustment_total": { + "amount": "0", + "currency": "USD" + }, + "total": { + "amount": "16.00", + "currency": "USD" + } + } + }, + "buyer_actor": { + "account_id": "buyer-account", + "pubkey": "buyer_pubkey", + "source": BUYER_ORDER_REQUEST_ACTOR_SOURCE_RESOLVED_ACCOUNT + } + } + }) + } + + fn unsupported_payload() -> Value { + let mut payload = supported_payload(); + payload["exportability"] = json!({ + "state": "identity_unresolved", + "reason": "canonical_hex_pubkey_required" + }); + payload["support_status"] = json!({ + "state": "unsupported", + "issues": ["buyer_pubkey_required"] + }); + payload["document"]["order"]["buyer_pubkey"] = json!(""); + payload["document"]["buyer_actor"]["pubkey"] = json!(""); + payload["document"]["buyer_actor"]["source"] = + json!(BUYER_ORDER_REQUEST_ACTOR_SOURCE_UNRESOLVED_APP); + payload + } + + fn assert_invalid(payload: Value, expected: &str) { + assert_error_contains( + validate_buyer_order_request_local_work_payload(&payload), + expected, + ); + } + + fn assert_error_contains<T: std::fmt::Debug>( + result: Result<T, RuntimeStoreError>, + expected: &str, + ) { + let error = result.expect_err("expected validation error"); + assert!( + error.to_string().contains(expected), + "expected error to contain {expected}, got {error}" + ); + } +} diff --git a/crates/runtime_store/src/store.rs b/crates/runtime_store/src/store.rs @@ -0,0 +1,970 @@ +#![forbid(unsafe_code)] + +use radroots_sql_core::SqlExecutor; +use radroots_sql_core::error::SqlError; +use serde::Deserialize; +use serde_json::{Value, json}; + +use crate::migrations; +use crate::models::validate_non_empty; +use crate::{ + PublishOutboxStatus, RuntimeStoreCursor, RuntimeStoreError, RuntimeStoreRecord, + RuntimeStoreRecordFamily, RuntimeStoreRecordInput, RuntimeStoreRecordStatus, + RuntimeStoreRecordUpdate, SourceRuntime, +}; + +pub struct RuntimeStore<E: SqlExecutor> { + executor: E, +} + +impl<E: SqlExecutor> RuntimeStore<E> { + pub fn new(executor: E) -> Self { + Self { executor } + } + + pub fn executor(&self) -> &E { + &self.executor + } + + pub fn migrate_up(&self) -> Result<(), SqlError> { + migrations::run_all_up(self.executor()) + } + + pub fn migrate_down(&self) -> Result<(), SqlError> { + migrations::run_all_down(self.executor()) + } + + pub fn append_record( + &self, + input: &RuntimeStoreRecordInput, + ) -> Result<RuntimeStoreRecord, RuntimeStoreError> { + input.validate()?; + self.executor.begin()?; + let result = (|| -> Result<(), RuntimeStoreError> { + let change_seq = self.next_change_seq()?; + let params = json!([ + change_seq, + input.record_id, + input.family.as_str(), + input.status.as_str(), + input.source_runtime.as_str(), + input.created_at_ms, + input.inserted_at_ms, + input.inserted_at_ms, + input.owner_account_id, + input.owner_pubkey, + input.farm_id, + input.listing_addr, + encode_json(input.local_work_json.as_ref()), + input.event_id, + input.event_kind, + input.event_pubkey, + input.event_created_at, + encode_json(input.event_tags_json.as_ref()), + input.event_content, + input.event_sig, + encode_json(input.raw_event_json.as_ref()), + input.outbox_status.as_str(), + input.relay_set_fingerprint, + encode_json(input.relay_delivery_json.as_ref()) + ]) + .to_string(); + let sql = "insert or ignore into runtime_store_record( + change_seq, + record_id, + family, + status, + source_runtime, + created_at_ms, + inserted_at_ms, + updated_at_ms, + owner_account_id, + owner_pubkey, + farm_id, + listing_addr, + local_work_json, + event_id, + event_kind, + event_pubkey, + event_created_at, + event_tags_json, + event_content, + event_sig, + raw_event_json, + outbox_status, + relay_set_fingerprint, + relay_delivery_json + ) values(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)"; + let _ = self.executor.exec(sql, &params)?; + Ok(()) + })(); + match result { + Ok(()) => self.executor.commit()?, + Err(err) => { + let _ = self.executor.rollback(); + return Err(err); + } + } + self.get_record(&input.record_id)? + .ok_or_else(|| RuntimeStoreError::InvalidRecord("record append failed".to_owned())) + } + + pub fn get_record( + &self, + record_id: &str, + ) -> Result<Option<RuntimeStoreRecord>, RuntimeStoreError> { + validate_non_empty("record_id", record_id)?; + let params = json!([record_id]).to_string(); + let rows = self.query_records( + "select * from runtime_store_record where record_id = ? limit 1", + &params, + )?; + Ok(rows.into_iter().next()) + } + + pub fn list_records_after_seq( + &self, + after_seq: i64, + limit: u32, + ) -> Result<Vec<RuntimeStoreRecord>, RuntimeStoreError> { + let params = json!([after_seq, i64::from(limit)]).to_string(); + self.query_records( + "select * from runtime_store_record where seq > ? order by seq asc limit ?", + &params, + ) + } + + pub fn list_records_changed_after( + &self, + after_change_seq: i64, + limit: u32, + ) -> Result<Vec<RuntimeStoreRecord>, RuntimeStoreError> { + let params = json!([after_change_seq, i64::from(limit)]).to_string(); + self.query_records( + "select * from runtime_store_record where change_seq > ? order by change_seq asc, seq asc limit ?", + &params, + ) + } + + pub fn list_records_changed_latest( + &self, + limit: u32, + ) -> Result<Vec<RuntimeStoreRecord>, RuntimeStoreError> { + let params = json!([i64::from(limit)]).to_string(); + self.query_records( + "select * from runtime_store_record order by change_seq desc, seq desc, record_id asc limit ?", + &params, + ) + } + + pub fn list_records_changed_before( + &self, + before_change_seq: i64, + before_seq: i64, + limit: u32, + ) -> Result<Vec<RuntimeStoreRecord>, RuntimeStoreError> { + let params = json!([ + before_change_seq, + before_change_seq, + before_seq, + i64::from(limit) + ]) + .to_string(); + self.query_records( + "select * from runtime_store_record + where change_seq < ? or (change_seq = ? and seq < ?) + order by change_seq desc, seq desc, record_id asc + limit ?", + &params, + ) + } + + pub fn update_outbox( + &self, + update: &RuntimeStoreRecordUpdate, + ) -> Result<RuntimeStoreRecord, RuntimeStoreError> { + validate_non_empty("record_id", &update.record_id)?; + self.executor.begin()?; + let result = (|| -> Result<i64, RuntimeStoreError> { + let change_seq = self.next_change_seq()?; + let params = json!([ + change_seq, + update.status.as_str(), + update.outbox_status.as_str(), + update.relay_set_fingerprint, + encode_json(update.relay_delivery_json.as_ref()), + update.updated_at_ms, + update.record_id + ]) + .to_string(); + let outcome = self.executor.exec( + "update runtime_store_record + set change_seq = ?, + status = ?, + outbox_status = ?, + relay_set_fingerprint = ?, + relay_delivery_json = ?, + updated_at_ms = ? + where record_id = ?", + &params, + )?; + Ok(outcome.changes) + })(); + let changes = match result { + Ok(changes) => { + self.executor.commit()?; + changes + } + Err(err) => { + let _ = self.executor.rollback(); + return Err(err); + } + }; + if changes == 0 { + return Err(RuntimeStoreError::Sql(SqlError::NotFound( + update.record_id.clone(), + ))); + } + self.get_record(&update.record_id)? + .ok_or_else(|| RuntimeStoreError::Sql(SqlError::NotFound(update.record_id.clone()))) + } + + pub fn get_cursor( + &self, + consumer_id: &str, + ) -> Result<Option<RuntimeStoreCursor>, RuntimeStoreError> { + validate_non_empty("consumer_id", consumer_id)?; + let params = json!([consumer_id]).to_string(); + let raw = self.executor.query_raw( + "select consumer_id, last_change_seq, updated_at_ms from runtime_store_projection_cursor where consumer_id = ? limit 1", + &params, + )?; + let rows: Vec<CursorRow> = serde_json::from_str(&raw)?; + Ok(rows.into_iter().next().map(Into::into)) + } + + pub fn advance_cursor( + &self, + consumer_id: &str, + last_change_seq: i64, + updated_at_ms: i64, + ) -> Result<RuntimeStoreCursor, RuntimeStoreError> { + validate_non_empty("consumer_id", consumer_id)?; + let params = json!([consumer_id, last_change_seq, updated_at_ms]).to_string(); + self.executor.exec( + "insert into runtime_store_projection_cursor(consumer_id, last_change_seq, updated_at_ms) + values(?,?,?) + on conflict(consumer_id) do update set + last_change_seq = max(runtime_store_projection_cursor.last_change_seq, excluded.last_change_seq), + updated_at_ms = excluded.updated_at_ms", + &params, + )?; + self.get_cursor(consumer_id)? + .ok_or_else(|| RuntimeStoreError::InvalidRecord("cursor advance failed".to_owned())) + } + + fn query_records( + &self, + sql: &str, + params: &str, + ) -> Result<Vec<RuntimeStoreRecord>, RuntimeStoreError> { + let raw = self.executor.query_raw(sql, params)?; + let rows: Vec<RecordRow> = serde_json::from_str(&raw)?; + rows.into_iter().map(TryInto::try_into).collect() + } + + fn next_change_seq(&self) -> Result<i64, RuntimeStoreError> { + let raw = self.executor.query_raw( + "select coalesce(max(change_seq), 0) + 1 as change_seq from runtime_store_record", + "[]", + )?; + let rows: Vec<ChangeSeqRow> = serde_json::from_str(&raw)?; + rows.into_iter() + .next() + .map(|row| row.change_seq) + .ok_or_else(|| { + RuntimeStoreError::InvalidRecord("change sequence unavailable".to_owned()) + }) + } +} + +#[derive(Debug, Deserialize)] +struct RecordRow { + seq: i64, + change_seq: i64, + record_id: String, + family: String, + status: String, + source_runtime: String, + created_at_ms: i64, + inserted_at_ms: i64, + updated_at_ms: i64, + owner_account_id: Option<String>, + owner_pubkey: Option<String>, + farm_id: Option<String>, + listing_addr: Option<String>, + local_work_json: Option<String>, + event_id: Option<String>, + event_kind: Option<i64>, + event_pubkey: Option<String>, + event_created_at: Option<i64>, + event_tags_json: Option<String>, + event_content: Option<String>, + event_sig: Option<String>, + raw_event_json: Option<String>, + outbox_status: String, + relay_set_fingerprint: Option<String>, + relay_delivery_json: Option<String>, +} + +impl TryFrom<RecordRow> for RuntimeStoreRecord { + type Error = RuntimeStoreError; + + fn try_from(row: RecordRow) -> Result<Self, Self::Error> { + Ok(Self { + seq: row.seq, + change_seq: row.change_seq, + record_id: row.record_id, + family: RuntimeStoreRecordFamily::parse(&row.family)?, + status: RuntimeStoreRecordStatus::parse(&row.status)?, + source_runtime: SourceRuntime::parse(&row.source_runtime)?, + created_at_ms: row.created_at_ms, + inserted_at_ms: row.inserted_at_ms, + updated_at_ms: row.updated_at_ms, + owner_account_id: row.owner_account_id, + owner_pubkey: row.owner_pubkey, + farm_id: row.farm_id, + listing_addr: row.listing_addr, + local_work_json: decode_json(row.local_work_json)?, + event_id: row.event_id, + event_kind: row.event_kind, + event_pubkey: row.event_pubkey, + event_created_at: row.event_created_at, + event_tags_json: decode_json(row.event_tags_json)?, + event_content: row.event_content, + event_sig: row.event_sig, + raw_event_json: decode_json(row.raw_event_json)?, + outbox_status: PublishOutboxStatus::parse(&row.outbox_status)?, + relay_set_fingerprint: row.relay_set_fingerprint, + relay_delivery_json: decode_json(row.relay_delivery_json)?, + }) + } +} + +#[derive(Debug, Deserialize)] +struct CursorRow { + consumer_id: String, + last_change_seq: i64, + updated_at_ms: i64, +} + +impl From<CursorRow> for RuntimeStoreCursor { + fn from(row: CursorRow) -> Self { + Self { + consumer_id: row.consumer_id, + last_change_seq: row.last_change_seq, + updated_at_ms: row.updated_at_ms, + } + } +} + +#[derive(Debug, Deserialize)] +struct ChangeSeqRow { + change_seq: i64, +} + +fn encode_json(value: Option<&Value>) -> Option<String> { + value.map(Value::to_string) +} + +fn decode_json(value: Option<String>) -> Result<Option<Value>, RuntimeStoreError> { + value + .map(|value| serde_json::from_str(&value)) + .transpose() + .map_err(Into::into) +} + +#[cfg(test)] +mod tests { + use std::collections::VecDeque; + use std::sync::Mutex; + use std::sync::atomic::{AtomicUsize, Ordering}; + + use radroots_sql_core::{ExecOutcome, SqlExecutor, SqliteExecutor}; + use serde_json::json; + + use super::*; + + fn store() -> RuntimeStore<SqliteExecutor> { + let executor = SqliteExecutor::open_memory().expect("open memory sqlite"); + let store = RuntimeStore::new(executor); + store.migrate_up().expect("migrate up"); + store + } + + fn local_work(record_id: &str) -> RuntimeStoreRecordInput { + RuntimeStoreRecordInput { + record_id: record_id.to_owned(), + family: RuntimeStoreRecordFamily::LocalWork, + status: RuntimeStoreRecordStatus::LocalSaved, + source_runtime: SourceRuntime::Cli, + created_at_ms: 1000, + inserted_at_ms: 1001, + owner_account_id: Some("seller-account".to_owned()), + owner_pubkey: Some("seller-pubkey".to_owned()), + farm_id: Some("farm-a".to_owned()), + listing_addr: Some("listing-a".to_owned()), + local_work_json: Some(json!({"kind":"listing","title":"Eggs"})), + event_id: None, + event_kind: None, + event_pubkey: None, + event_created_at: None, + event_tags_json: None, + event_content: None, + event_sig: None, + raw_event_json: None, + outbox_status: PublishOutboxStatus::None, + relay_set_fingerprint: None, + relay_delivery_json: None, + } + } + + fn signed_event(record_id: &str) -> RuntimeStoreRecordInput { + RuntimeStoreRecordInput { + record_id: record_id.to_owned(), + family: RuntimeStoreRecordFamily::SignedEvent, + status: RuntimeStoreRecordStatus::PendingPublish, + source_runtime: SourceRuntime::Cli, + created_at_ms: 2000, + inserted_at_ms: 2001, + owner_account_id: Some("seller-account".to_owned()), + owner_pubkey: Some("seller-pubkey".to_owned()), + farm_id: Some("farm-a".to_owned()), + listing_addr: Some("listing-a".to_owned()), + local_work_json: None, + event_id: Some(record_id.to_owned()), + event_kind: Some(3421), + event_pubkey: Some("seller-pubkey".to_owned()), + event_created_at: Some(2000), + event_tags_json: Some(json!([["d", "listing-a"]])), + event_content: Some("{\"title\":\"Eggs\"}".to_owned()), + event_sig: Some("sig-a".to_owned()), + raw_event_json: Some(json!({"id":record_id,"kind":3421})), + outbox_status: PublishOutboxStatus::Pending, + relay_set_fingerprint: None, + relay_delivery_json: None, + } + } + + #[derive(Debug)] + struct ScriptedExecutor { + begin_result: Mutex<Result<(), SqlError>>, + commit_result: Mutex<Result<(), SqlError>>, + exec_results: Mutex<VecDeque<Result<ExecOutcome, SqlError>>>, + query_results: Mutex<VecDeque<Result<String, SqlError>>>, + rollbacks: AtomicUsize, + } + + impl ScriptedExecutor { + fn new( + exec_results: Vec<Result<ExecOutcome, SqlError>>, + query_results: Vec<Result<String, SqlError>>, + ) -> Self { + Self { + begin_result: Mutex::new(Ok(())), + commit_result: Mutex::new(Ok(())), + exec_results: Mutex::new(exec_results.into()), + query_results: Mutex::new(query_results.into()), + rollbacks: AtomicUsize::new(0), + } + } + + fn with_begin_error(error: SqlError) -> Self { + let executor = Self::new(Vec::new(), Vec::new()); + *executor.begin_result.lock().expect("begin result") = Err(error); + executor + } + + fn with_commit_error(error: SqlError) -> Self { + let executor = Self::new( + vec![Ok(ExecOutcome { + changes: 1, + last_insert_id: 0, + })], + vec![Ok(r#"[{"change_seq":1}]"#.to_owned())], + ); + *executor.commit_result.lock().expect("commit result") = Err(error); + executor + } + } + + impl SqlExecutor for ScriptedExecutor { + fn exec(&self, _sql: &str, _params_json: &str) -> Result<ExecOutcome, SqlError> { + self.exec_results + .lock() + .expect("exec results") + .pop_front() + .unwrap_or(Ok(ExecOutcome { + changes: 1, + last_insert_id: 0, + })) + } + + fn query_raw(&self, _sql: &str, _params_json: &str) -> Result<String, SqlError> { + self.query_results + .lock() + .expect("query results") + .pop_front() + .unwrap_or_else(|| Ok("[]".to_owned())) + } + + fn begin(&self) -> Result<(), SqlError> { + self.begin_result.lock().expect("begin result").clone() + } + + fn commit(&self) -> Result<(), SqlError> { + self.commit_result.lock().expect("commit result").clone() + } + + fn rollback(&self) -> Result<(), SqlError> { + self.rollbacks.fetch_add(1, Ordering::SeqCst); + Ok(()) + } + } + + fn record_row_with(field: &str, value: serde_json::Value) -> String { + let mut row = json!({ + "seq": 1, + "change_seq": 1, + "record_id": "record-a", + "family": "signed_event", + "status": "pending_publish", + "source_runtime": "cli", + "created_at_ms": 1000, + "inserted_at_ms": 1001, + "updated_at_ms": 1001, + "owner_account_id": "seller-account", + "owner_pubkey": "seller-pubkey", + "farm_id": "farm-a", + "listing_addr": "listing-a", + "local_work_json": null, + "event_id": "event-a", + "event_kind": 3421, + "event_pubkey": "seller-pubkey", + "event_created_at": 1000, + "event_tags_json": "[[\"d\",\"listing-a\"]]", + "event_content": "{}", + "event_sig": "sig-a", + "raw_event_json": "{\"id\":\"event-a\",\"kind\":3421}", + "outbox_status": "pending", + "relay_set_fingerprint": null, + "relay_delivery_json": null + }); + row[field] = value; + json!([row]).to_string() + } + + #[test] + fn store_methods_round_trip_records_and_cursors() { + let store = store(); + + assert!( + store + .executor() + .query_raw("select 1 as value", "[]") + .is_ok() + ); + assert!(store.get_record("missing").expect("get missing").is_none()); + assert!(store.get_cursor("app").expect("cursor missing").is_none()); + + let local = store + .append_record(&local_work("local-a")) + .expect("append local work"); + let event = store + .append_record(&signed_event("event-a")) + .expect("append signed event"); + + assert_eq!( + store + .get_record("local-a") + .expect("get local") + .expect("local record") + .record_id, + local.record_id + ); + assert_eq!( + store + .list_records_after_seq(0, 10) + .expect("list after seq") + .len(), + 2 + ); + assert_eq!( + store + .list_records_changed_after(local.change_seq, 10) + .expect("list changed after")[0] + .record_id, + event.record_id + ); + assert_eq!( + store.list_records_changed_latest(1).expect("list latest")[0].record_id, + event.record_id + ); + assert_eq!( + store + .list_records_changed_before(event.change_seq, event.seq, 10) + .expect("list before")[0] + .record_id, + local.record_id + ); + + let cursor = store + .advance_cursor("app", event.change_seq, 3000) + .expect("advance cursor"); + assert_eq!(cursor.consumer_id, "app"); + assert_eq!( + store + .get_cursor("app") + .expect("get cursor") + .expect("cursor") + .last_change_seq, + event.change_seq + ); + + let updated = store + .update_outbox(&RuntimeStoreRecordUpdate { + record_id: "event-a".to_owned(), + status: RuntimeStoreRecordStatus::Published, + outbox_status: PublishOutboxStatus::Acknowledged, + relay_set_fingerprint: Some("relay-a".to_owned()), + relay_delivery_json: Some(json!({ + "state": "acknowledged", + "target_relays": ["wss://relay.example"], + "connected_relays": ["wss://relay.example"], + "acknowledged_relays": ["wss://relay.example"] + })), + updated_at_ms: 4000, + }) + .expect("update outbox"); + + assert_eq!(updated.status, RuntimeStoreRecordStatus::Published); + assert_eq!(updated.outbox_status, PublishOutboxStatus::Acknowledged); + assert_eq!(updated.relay_set_fingerprint.as_deref(), Some("relay-a")); + store.migrate_down().expect("migrate down"); + } + + #[test] + fn store_reports_missing_updates_and_decode_errors() { + let store = store(); + assert!( + store + .get_record(" ") + .expect_err("empty record id") + .to_string() + .contains("record_id") + ); + assert!( + store + .get_cursor(" ") + .expect_err("empty consumer id") + .to_string() + .contains("consumer_id") + ); + assert!( + store + .advance_cursor(" ", 1, 1000) + .expect_err("empty cursor consumer") + .to_string() + .contains("consumer_id") + ); + assert!( + store + .update_outbox(&RuntimeStoreRecordUpdate { + record_id: " ".to_owned(), + status: RuntimeStoreRecordStatus::Published, + outbox_status: PublishOutboxStatus::Acknowledged, + relay_set_fingerprint: None, + relay_delivery_json: None, + updated_at_ms: 4000, + }) + .expect_err("empty update record id") + .to_string() + .contains("record_id") + ); + + let missing_update = store + .update_outbox(&RuntimeStoreRecordUpdate { + record_id: "missing-event".to_owned(), + status: RuntimeStoreRecordStatus::Published, + outbox_status: PublishOutboxStatus::Acknowledged, + relay_set_fingerprint: None, + relay_delivery_json: None, + updated_at_ms: 4000, + }) + .expect_err("missing record update"); + + assert!(missing_update.to_string().contains("missing-event")); + + store + .append_record(&local_work("local-a")) + .expect("append local"); + let params = json!(["{", "local-a"]).to_string(); + store + .executor() + .exec( + "update runtime_store_record set local_work_json = ? where record_id = ?", + &params, + ) + .expect("corrupt local work json"); + let decode_error = store.get_record("local-a").expect_err("decode error"); + + assert!(decode_error.to_string().contains("EOF")); + } + + #[test] + fn store_rolls_back_when_change_sequence_is_unavailable() { + let append_store = + RuntimeStore::new(ScriptedExecutor::new(Vec::new(), vec![Ok("[]".to_owned())])); + let append_error = append_store + .append_record(&local_work("local-a")) + .expect_err("append error"); + + assert!(append_error.to_string().contains("change sequence")); + assert_eq!(append_store.executor().rollbacks.load(Ordering::SeqCst), 1); + + let update_store = + RuntimeStore::new(ScriptedExecutor::new(Vec::new(), vec![Ok("[]".to_owned())])); + let update_error = update_store + .update_outbox(&RuntimeStoreRecordUpdate { + record_id: "event-a".to_owned(), + status: RuntimeStoreRecordStatus::Published, + outbox_status: PublishOutboxStatus::Acknowledged, + relay_set_fingerprint: None, + relay_delivery_json: None, + updated_at_ms: 4000, + }) + .expect_err("update error"); + + assert!(update_error.to_string().contains("change sequence")); + assert_eq!(update_store.executor().rollbacks.load(Ordering::SeqCst), 1); + } + + #[test] + fn store_reports_cursor_advance_without_returned_cursor() { + let store = RuntimeStore::new(ScriptedExecutor::new(Vec::new(), Vec::new())); + + assert!(store.get_cursor("app").expect("missing cursor").is_none()); + let cursor_error = store + .advance_cursor("app", 1, 1000) + .expect_err("cursor advance error"); + + assert!(cursor_error.to_string().contains("cursor advance failed")); + } + + #[test] + fn store_reports_executor_and_decode_failures() { + let begin_store = RuntimeStore::new(ScriptedExecutor::with_begin_error( + SqlError::InvalidQuery("begin failed".to_owned()), + )); + assert!( + begin_store + .append_record(&local_work("local-a")) + .expect_err("begin failure") + .to_string() + .contains("begin failed") + ); + + let exec_store = RuntimeStore::new(ScriptedExecutor::new( + vec![Err(SqlError::InvalidQuery("insert failed".to_owned()))], + vec![Ok(r#"[{"change_seq":1}]"#.to_owned())], + )); + assert!( + exec_store + .append_record(&local_work("local-a")) + .expect_err("exec failure") + .to_string() + .contains("insert failed") + ); + assert_eq!(exec_store.executor().rollbacks.load(Ordering::SeqCst), 1); + + let commit_store = RuntimeStore::new(ScriptedExecutor::with_commit_error( + SqlError::InvalidQuery("commit failed".to_owned()), + )); + assert!( + commit_store + .append_record(&local_work("local-a")) + .expect_err("commit failure") + .to_string() + .contains("commit failed") + ); + + let query_error_store = RuntimeStore::new(ScriptedExecutor::new( + Vec::new(), + vec![Err(SqlError::InvalidQuery("query failed".to_owned()))], + )); + assert!( + query_error_store + .get_record("record-a") + .expect_err("query failure") + .to_string() + .contains("query failed") + ); + + let invalid_rows_store = + RuntimeStore::new(ScriptedExecutor::new(Vec::new(), vec![Ok("{".to_owned())])); + let _ = invalid_rows_store + .get_record("record-a") + .expect_err("invalid rows"); + + let cursor_rows_store = + RuntimeStore::new(ScriptedExecutor::new(Vec::new(), vec![Ok("{".to_owned())])); + let _ = cursor_rows_store + .get_cursor("app") + .expect_err("invalid cursor rows"); + + let change_rows_store = + RuntimeStore::new(ScriptedExecutor::new(Vec::new(), vec![Ok("{".to_owned())])); + let _ = change_rows_store + .append_record(&local_work("local-a")) + .expect_err("invalid change rows"); + + let cursor_exec_store = RuntimeStore::new(ScriptedExecutor::new( + vec![Err(SqlError::InvalidQuery("cursor failed".to_owned()))], + Vec::new(), + )); + assert!( + cursor_exec_store + .advance_cursor("app", 1, 1000) + .expect_err("cursor exec failure") + .to_string() + .contains("cursor failed") + ); + + let append_lookup_store = RuntimeStore::new(ScriptedExecutor::new( + vec![Ok(ExecOutcome { + changes: 1, + last_insert_id: 0, + })], + vec![Ok(r#"[{"change_seq":1}]"#.to_owned()), Ok("[]".to_owned())], + )); + assert!( + append_lookup_store + .append_record(&local_work("local-a")) + .expect_err("append lookup failure") + .to_string() + .contains("record append failed") + ); + + let update_lookup_store = RuntimeStore::new(ScriptedExecutor::new( + vec![Ok(ExecOutcome { + changes: 1, + last_insert_id: 0, + })], + vec![Ok(r#"[{"change_seq":1}]"#.to_owned()), Ok("[]".to_owned())], + )); + assert!( + update_lookup_store + .update_outbox(&RuntimeStoreRecordUpdate { + record_id: "event-a".to_owned(), + status: RuntimeStoreRecordStatus::Published, + outbox_status: PublishOutboxStatus::Acknowledged, + relay_set_fingerprint: None, + relay_delivery_json: None, + updated_at_ms: 4000, + }) + .expect_err("update lookup failure") + .to_string() + .contains("event-a") + ); + + let cursor_query_store = RuntimeStore::new(ScriptedExecutor::new( + Vec::new(), + vec![Err(SqlError::InvalidQuery( + "cursor query failed".to_owned(), + ))], + )); + assert!( + cursor_query_store + .get_cursor("app") + .expect_err("cursor query failure") + .to_string() + .contains("cursor query failed") + ); + + let advance_cursor_query_store = RuntimeStore::new(ScriptedExecutor::new( + vec![Ok(ExecOutcome { + changes: 1, + last_insert_id: 0, + })], + vec![Err(SqlError::InvalidQuery( + "advanced cursor query failed".to_owned(), + ))], + )); + assert!( + advance_cursor_query_store + .advance_cursor("app", 1, 1000) + .expect_err("advance cursor query failure") + .to_string() + .contains("advanced cursor query failed") + ); + + let change_query_store = RuntimeStore::new(ScriptedExecutor::new( + Vec::new(), + vec![Err(SqlError::InvalidQuery( + "change query failed".to_owned(), + ))], + )); + assert!( + change_query_store + .append_record(&local_work("local-a")) + .expect_err("change query failure") + .to_string() + .contains("change query failed") + ); + } + + #[test] + fn store_reports_record_row_conversion_failures() { + for (field, value, expected) in [ + ("family", json!("bad_family"), "family"), + ("status", json!("bad_status"), "status"), + ("source_runtime", json!("bad_runtime"), "runtime"), + ("event_tags_json", json!("{"), "EOF"), + ("raw_event_json", json!("{"), "EOF"), + ("outbox_status", json!("bad_outbox"), "outbox"), + ] { + let store = RuntimeStore::new(ScriptedExecutor::new( + Vec::new(), + vec![Ok(record_row_with(field, value))], + )); + let error = store.get_record("record-a").expect_err("conversion error"); + + assert!( + error.to_string().contains(expected), + "expected error to contain {expected}, got {error}" + ); + } + + let store = RuntimeStore::new(ScriptedExecutor::new( + vec![Ok(ExecOutcome { + changes: 1, + last_insert_id: 0, + })], + vec![ + Ok(r#"[{"change_seq":1}]"#.to_owned()), + Ok(record_row_with("status", json!("published"))), + ], + )); + let updated = store + .update_outbox(&RuntimeStoreRecordUpdate { + record_id: "record-a".to_owned(), + status: RuntimeStoreRecordStatus::Published, + outbox_status: PublishOutboxStatus::Acknowledged, + relay_set_fingerprint: None, + relay_delivery_json: None, + updated_at_ms: 4000, + }) + .expect("scripted update"); + assert_eq!(updated.status, RuntimeStoreRecordStatus::Published); + } +} diff --git a/crates/runtime_store/tests/order_work.rs b/crates/runtime_store/tests/order_work.rs @@ -0,0 +1,294 @@ +use radroots_runtime_store::{ + BUYER_ORDER_REQUEST_ACTOR_SOURCE_RESOLVED_ACCOUNT, + BUYER_ORDER_REQUEST_ACTOR_SOURCE_UNRESOLVED_APP, BUYER_ORDER_REQUEST_DOCUMENT_KIND, + BUYER_ORDER_REQUEST_LOCAL_WORK_RECORD_KIND, BuyerOrderRequestSupportState, + buyer_order_request_local_work_record_id, validate_buyer_order_request_local_work_payload, + validate_supported_buyer_order_request_local_work_payload, + validate_unsupported_buyer_order_request_local_work_payload, +}; +use serde_json::{Value, json}; + +#[test] +fn buyer_order_request_record_id_is_deterministic_for_app_orders() { + assert_eq!( + buyer_order_request_local_work_record_id(" order-1 ").expect("record id"), + "app:local_work:order_request:order-1" + ); +} + +#[test] +fn buyer_order_request_payload_accepts_supported_exportable_work() { + let payload = supported_payload(); + + let validation = + validate_buyer_order_request_local_work_payload(&payload).expect("valid payload"); + let supported = validate_supported_buyer_order_request_local_work_payload(&payload) + .expect("supported payload"); + + assert_eq!(validation.order_id, "ord_1"); + assert_eq!( + validation.support_state, + BuyerOrderRequestSupportState::Supported + ); + assert_eq!(validation.support_state.as_str(), "supported"); + assert!(validation.support_issues.is_empty()); + assert_eq!(supported, validation); +} + +#[test] +fn buyer_order_request_payload_accepts_explicit_unsupported_work() { + let mut payload = supported_payload(); + payload["exportability"] = json!({ + "state": "identity_unresolved", + "reason": "canonical_hex_pubkey_required" + }); + payload["support_status"] = json!({ + "state": "unsupported", + "issues": ["buyer_pubkey_required"] + }); + payload["document"]["order"]["buyer_pubkey"] = json!(""); + payload["document"]["buyer_actor"]["pubkey"] = json!(""); + payload["document"]["buyer_actor"]["source"] = + json!(BUYER_ORDER_REQUEST_ACTOR_SOURCE_UNRESOLVED_APP); + + let validation = + validate_buyer_order_request_local_work_payload(&payload).expect("valid payload"); + let unsupported = validate_unsupported_buyer_order_request_local_work_payload(&payload) + .expect("unsupported payload"); + let supported_error = validate_supported_buyer_order_request_local_work_payload(&payload) + .expect_err("unsupported payload should not validate as supported"); + + assert_eq!( + validation.support_state, + BuyerOrderRequestSupportState::Unsupported + ); + assert_eq!(validation.support_state.as_str(), "unsupported"); + assert_eq!(validation.support_issues, vec!["buyer_pubkey_required"]); + assert_eq!(unsupported, validation); + assert!(supported_error.to_string().contains("support_status.state")); +} + +#[test] +fn buyer_order_request_payload_rejects_missing_identity() { + for (path, expected) in [ + (vec!["document", "order", "listing_addr"], "listing_addr"), + ( + vec!["document", "order", "listing_event_id"], + "listing_event_id", + ), + (vec!["document", "order", "seller_pubkey"], "seller_pubkey"), + (vec!["document", "order", "buyer_pubkey"], "buyer_pubkey"), + ] { + let mut payload = supported_payload(); + set_path(&mut payload, &path, json!("")); + + assert_invalid(payload, expected); + } +} + +#[test] +fn buyer_order_request_payload_rejects_missing_items() { + let mut payload = supported_payload(); + payload["document"]["order"]["items"] = json!([]); + + assert_invalid(payload, "items"); +} + +#[test] +fn buyer_order_request_payload_rejects_invalid_item_identity() { + let mut missing_bin = supported_payload(); + missing_bin["document"]["order"]["items"][0]["bin_id"] = json!(""); + assert_invalid(missing_bin, "items[0].bin_id"); + + let mut zero_count = supported_payload(); + zero_count["document"]["order"]["items"][0]["bin_count"] = json!(0); + assert_invalid(zero_count, "items[0].bin_count"); +} + +#[test] +fn buyer_order_request_payload_rejects_invalid_economics() { + let mut missing_economics = supported_payload(); + missing_economics["document"]["order"] + .as_object_mut() + .expect("order object") + .remove("economics"); + assert_invalid(missing_economics, "economics"); + + let mut non_object_economics = supported_payload(); + non_object_economics["document"]["order"]["economics"] = Value::Null; + assert_invalid(non_object_economics, "economics"); + + let mut mismatched_currency = supported_payload(); + mismatched_currency["document"]["order"]["economics"]["items"][0]["unit_price_currency"] = + json!("CAD"); + assert_invalid(mismatched_currency, "unit_price_currency"); + + let mut mismatched_items = supported_payload(); + mismatched_items["document"]["order"]["economics"]["items"] = json!([ + { + "bin_id": "dozen-eggs", + "bin_count": 2, + "quantity_amount": "1", + "quantity_unit": "dozen", + "unit_price_amount": "8.00", + "unit_price_currency": "USD", + "line_subtotal": { + "amount": "16.00", + "currency": "USD" + } + }, + { + "bin_id": "half-dozen-eggs", + "bin_count": 1 + } + ]); + assert_invalid(mismatched_items, "economics.items"); + + let mut empty_economics_items = supported_payload(); + empty_economics_items["document"]["order"]["economics"]["items"] = json!([]); + assert_invalid(empty_economics_items, "economics.items"); + + let mut mismatched_bin = supported_payload(); + mismatched_bin["document"]["order"]["economics"]["items"][0]["bin_id"] = json!("other-bin"); + assert_invalid(mismatched_bin, "economics.items[0].bin_id"); + + let mut bad_currency_length = supported_payload(); + bad_currency_length["document"]["order"]["economics"]["currency"] = json!("US"); + assert_invalid(bad_currency_length, "currency"); + + let mut missing_line_subtotal = supported_payload(); + missing_line_subtotal["document"]["order"]["economics"]["items"][0] + .as_object_mut() + .expect("economics item") + .remove("line_subtotal"); + assert_invalid(missing_line_subtotal, "line_subtotal"); +} + +#[test] +fn buyer_order_request_payload_rejects_stale_or_conflicting_currentness() { + let mut stale = supported_payload(); + stale["currentness"]["current"] = json!(false); + assert_invalid(stale, "currentness.current"); + + let mut missing_current = supported_payload(); + missing_current["currentness"]["current"] = Value::Null; + assert_invalid(missing_current, "currentness.current"); + + let mut wrong_order = supported_payload(); + wrong_order["currentness"]["order_id"] = json!("ord_other"); + assert_invalid(wrong_order, "currentness.order_id"); +} + +#[test] +fn buyer_order_request_payload_rejects_malformed_support_status() { + let mut supported_with_issue = supported_payload(); + supported_with_issue["support_status"]["issues"] = json!(["unit_price_required"]); + assert_invalid(supported_with_issue, "support_status.issues"); + + let mut unsupported_without_issue = supported_payload(); + unsupported_without_issue["support_status"] = json!({ + "state": "unsupported", + "issues": [] + }); + assert_invalid(unsupported_without_issue, "support_status.issues"); +} + +fn supported_payload() -> Value { + json!({ + "record_kind": BUYER_ORDER_REQUEST_LOCAL_WORK_RECORD_KIND, + "scope": "app", + "exportability": { + "state": "exportable" + }, + "support_status": { + "state": "supported", + "issues": [] + }, + "currentness": { + "current": true, + "source": "app_sqlite_order", + "record_id": "app:local_work:order_request:ord_1", + "order_id": "ord_1", + "order_updated_at": "2026-05-24T12:00:00Z", + "created_at_ms": 1777777777000_i64 + }, + "document": { + "version": 1, + "kind": BUYER_ORDER_REQUEST_DOCUMENT_KIND, + "order": { + "order_id": "ord_1", + "listing_addr": "30402:seller_pubkey:listing_key", + "listing_event_id": "event-listing-1", + "buyer_pubkey": "buyer_pubkey", + "seller_pubkey": "seller_pubkey", + "items": [ + { + "bin_id": "dozen-eggs", + "bin_count": 2 + } + ], + "economics": { + "quote_id": "app-order:ord_1", + "quote_version": 1, + "pricing_basis": "listing_event", + "currency": "USD", + "items": [ + { + "bin_id": "dozen-eggs", + "bin_count": 2, + "quantity_amount": "1", + "quantity_unit": "dozen", + "unit_price_amount": "8.00", + "unit_price_currency": "USD", + "line_subtotal": { + "amount": "16.00", + "currency": "USD" + } + } + ], + "discounts": [], + "adjustments": [], + "subtotal": { + "amount": "16.00", + "currency": "USD" + }, + "discount_total": { + "amount": "0", + "currency": "USD" + }, + "adjustment_total": { + "amount": "0", + "currency": "USD" + }, + "total": { + "amount": "16.00", + "currency": "USD" + } + } + }, + "buyer_actor": { + "account_id": "buyer-account", + "pubkey": "buyer_pubkey", + "source": BUYER_ORDER_REQUEST_ACTOR_SOURCE_RESOLVED_ACCOUNT + }, + "listing_lookup": "30402:seller_pubkey:listing_key" + } + }) +} + +fn assert_invalid(payload: Value, expected: &str) { + let error = + validate_buyer_order_request_local_work_payload(&payload).expect_err("invalid payload"); + assert!( + error.to_string().contains(expected), + "expected error to contain {expected}, got {error}" + ); +} + +fn set_path(payload: &mut Value, path: &[&str], value: Value) { + let mut current = payload; + for segment in &path[..path.len() - 1] { + current = current.get_mut(*segment).expect("path segment"); + } + current[path[path.len() - 1]] = value; +} diff --git a/crates/runtime_store/tests/store.rs b/crates/runtime_store/tests/store.rs @@ -0,0 +1,545 @@ +use radroots_runtime_store::{ + MIGRATIONS, PublishOutboxStatus, RelayDeliveryEvidence, RuntimeStore, RuntimeStoreRecordFamily, + RuntimeStoreRecordInput, RuntimeStoreRecordStatus, RuntimeStoreRecordUpdate, SourceRuntime, +}; +use radroots_sql_core::migrations::migrations_run_all_up; +use radroots_sql_core::{SqlExecutor, SqliteExecutor}; +use serde_json::json; + +fn store() -> RuntimeStore<SqliteExecutor> { + let executor = SqliteExecutor::open_memory().expect("open memory sqlite"); + let store = RuntimeStore::new(executor); + store.migrate_up().expect("migrate runtime store"); + store +} + +fn local_work(record_id: &str) -> RuntimeStoreRecordInput { + RuntimeStoreRecordInput { + record_id: record_id.to_owned(), + family: RuntimeStoreRecordFamily::LocalWork, + status: RuntimeStoreRecordStatus::LocalSaved, + source_runtime: SourceRuntime::Cli, + created_at_ms: 1000, + inserted_at_ms: 1001, + owner_account_id: Some("seller-account".to_owned()), + owner_pubkey: Some("seller-pubkey".to_owned()), + farm_id: Some("farm-a".to_owned()), + listing_addr: Some("listing-a".to_owned()), + local_work_json: Some(json!({"kind":"listing","title":"Eggs"})), + event_id: None, + event_kind: None, + event_pubkey: None, + event_created_at: None, + event_tags_json: None, + event_content: None, + event_sig: None, + raw_event_json: None, + outbox_status: PublishOutboxStatus::None, + relay_set_fingerprint: None, + relay_delivery_json: None, + } +} + +fn signed_event(record_id: &str) -> RuntimeStoreRecordInput { + RuntimeStoreRecordInput { + record_id: record_id.to_owned(), + family: RuntimeStoreRecordFamily::SignedEvent, + status: RuntimeStoreRecordStatus::PendingPublish, + source_runtime: SourceRuntime::Cli, + created_at_ms: 2000, + inserted_at_ms: 2001, + owner_account_id: Some("seller-account".to_owned()), + owner_pubkey: Some("seller-pubkey".to_owned()), + farm_id: Some("farm-a".to_owned()), + listing_addr: Some("listing-a".to_owned()), + local_work_json: None, + event_id: Some("event-a".to_owned()), + event_kind: Some(3421), + event_pubkey: Some("seller-pubkey".to_owned()), + event_created_at: Some(2000), + event_tags_json: Some(json!([["d", "listing-a"]])), + event_content: Some("{\"title\":\"Eggs\"}".to_owned()), + event_sig: Some("sig-a".to_owned()), + raw_event_json: Some(json!({"id":"event-a","kind":3421})), + outbox_status: PublishOutboxStatus::Pending, + relay_set_fingerprint: None, + relay_delivery_json: None, + } +} + +#[test] +fn append_rejects_malformed_local_work_records() { + let store = store(); + let mut input = local_work("local-a"); + input.local_work_json = None; + + let err = store.append_record(&input).expect_err("invalid record"); + + assert!(err.to_string().contains("local_work_json")); +} + +#[test] +fn append_is_idempotent_by_record_id() { + let store = store(); + let input = local_work("local-a"); + + let first = store.append_record(&input).expect("append first"); + let second = store.append_record(&input).expect("append second"); + let rows = store.list_records_after_seq(0, 10).expect("list records"); + + assert_eq!(first.seq, second.seq); + assert_eq!(first.change_seq, second.change_seq); + assert_eq!(rows.len(), 1); + assert_eq!(rows[0].record_id, "local-a"); + assert_eq!( + rows[0].local_work_json, + Some(json!({"kind":"listing","title":"Eggs"})) + ); +} + +#[test] +fn source_runtime_network_round_trips() { + let store = store(); + let mut input = signed_event("event-network-a"); + input.source_runtime = SourceRuntime::Network; + + let inserted = store.append_record(&input).expect("append network event"); + let rows = store + .list_records_after_seq(0, 10) + .expect("list network event"); + + assert_eq!(SourceRuntime::Network.as_str(), "network"); + assert_eq!( + SourceRuntime::parse("network").expect("parse network runtime"), + SourceRuntime::Network + ); + assert_eq!(inserted.source_runtime, SourceRuntime::Network); + assert_eq!(rows.len(), 1); + assert_eq!(rows[0].source_runtime, SourceRuntime::Network); +} + +#[test] +fn projection_cursor_advances_without_rewinding() { + let store = store(); + + let first = store + .advance_cursor("app", 10, 100) + .expect("advance cursor"); + let second = store.advance_cursor("app", 5, 200).expect("ignore rewind"); + let third = store.advance_cursor("app", 12, 300).expect("advance again"); + + assert_eq!(first.last_change_seq, 10); + assert_eq!(second.last_change_seq, 10); + assert_eq!(third.last_change_seq, 12); +} + +#[test] +fn outbox_status_updates_signed_event_records() { + let store = store(); + let input = signed_event("event-a"); + store.append_record(&input).expect("append signed event"); + + let updated = store + .update_outbox(&RuntimeStoreRecordUpdate { + record_id: "event-a".to_owned(), + status: RuntimeStoreRecordStatus::Published, + outbox_status: PublishOutboxStatus::Acknowledged, + relay_set_fingerprint: None, + relay_delivery_json: None, + updated_at_ms: 3000, + }) + .expect("update outbox"); + + assert_eq!(updated.status, RuntimeStoreRecordStatus::Published); + assert_eq!(updated.outbox_status, PublishOutboxStatus::Acknowledged); +} + +#[test] +fn relay_delivery_evidence_round_trips_on_signed_event_records() { + let store = store(); + let mut input = signed_event("event-a"); + let evidence = RelayDeliveryEvidence::acknowledged( + ["wss://relay.example"], + ["wss://relay.example"], + ["wss://relay.example"], + Vec::new(), + ) + .expect("relay delivery evidence"); + input.relay_set_fingerprint = evidence.relay_set_fingerprint(); + input.relay_delivery_json = Some(evidence.to_json_value().expect("relay evidence json")); + + let inserted = store.append_record(&input).expect("append signed event"); + let loaded = store + .get_record("event-a") + .expect("load signed event") + .expect("signed event"); + + assert_eq!(inserted.relay_set_fingerprint, loaded.relay_set_fingerprint); + assert_eq!(inserted.relay_delivery_json, loaded.relay_delivery_json); +} + +#[test] +fn changed_after_uses_change_seq_for_appends_and_outbox_updates() { + let store = store(); + let input = signed_event("event-a"); + let appended = store.append_record(&input).expect("append signed event"); + let initial_rows = store + .list_records_changed_after(0, 10) + .expect("list initial changes"); + + assert_eq!(initial_rows.len(), 1); + assert_eq!(initial_rows[0].record_id, "event-a"); + assert_eq!(initial_rows[0].seq, appended.seq); + assert_eq!(initial_rows[0].change_seq, appended.change_seq); + + let updated = store + .update_outbox(&RuntimeStoreRecordUpdate { + record_id: "event-a".to_owned(), + status: RuntimeStoreRecordStatus::Published, + outbox_status: PublishOutboxStatus::Acknowledged, + relay_set_fingerprint: None, + relay_delivery_json: None, + updated_at_ms: 3000, + }) + .expect("update outbox"); + let changed_rows = store + .list_records_changed_after(appended.change_seq, 10) + .expect("list changed rows"); + + assert_eq!(updated.seq, appended.seq); + assert!(updated.change_seq > appended.change_seq); + assert_eq!(changed_rows.len(), 1); + assert_eq!(changed_rows[0].record_id, "event-a"); + assert_eq!(changed_rows[0].change_seq, updated.change_seq); +} + +#[test] +fn changed_latest_lists_newest_records_first() { + let store = store(); + let first = store + .append_record(&local_work("local-a")) + .expect("append first"); + let second = store + .append_record(&local_work("local-b")) + .expect("append second"); + let third = store + .append_record(&local_work("local-c")) + .expect("append third"); + + let rows = store + .list_records_changed_latest(2) + .expect("list latest changed rows"); + + assert_eq!(rows.len(), 2); + assert_eq!(rows[0].record_id, "local-c"); + assert_eq!(rows[0].change_seq, third.change_seq); + assert_eq!(rows[1].record_id, "local-b"); + assert_eq!(rows[1].change_seq, second.change_seq); + assert!(rows[1].change_seq > first.change_seq); +} + +#[test] +fn changed_before_pages_newest_first_by_cursor() { + let store = store(); + let _first = store + .append_record(&local_work("local-a")) + .expect("append first"); + let second = store + .append_record(&local_work("local-b")) + .expect("append second"); + let third = store + .append_record(&local_work("local-c")) + .expect("append third"); + let fourth = store + .append_record(&local_work("local-d")) + .expect("append fourth"); + + let first_page = store + .list_records_changed_latest(2) + .expect("list first page"); + let cursor = first_page.last().expect("last first page"); + let second_page = store + .list_records_changed_before(cursor.change_seq, cursor.seq, 2) + .expect("list second page"); + + assert_eq!(first_page.len(), 2); + assert_eq!(first_page[0].record_id, "local-d"); + assert_eq!(first_page[0].change_seq, fourth.change_seq); + assert_eq!(first_page[1].record_id, "local-c"); + assert_eq!(first_page[1].change_seq, third.change_seq); + assert_eq!(second_page.len(), 2); + assert_eq!(second_page[0].record_id, "local-b"); + assert_eq!(second_page[0].change_seq, second.change_seq); + assert_eq!(second_page[1].record_id, "local-a"); +} + +#[test] +fn changed_latest_is_not_blocked_by_older_record_volume() { + let store = store(); + for index in 0..505 { + store + .append_record(&local_work(&format!("older-{index:03}"))) + .expect("append older record"); + } + let current = store + .append_record(&local_work("current-record")) + .expect("append current record"); + + let rows = store + .list_records_changed_latest(1) + .expect("list latest record"); + + assert_eq!(rows.len(), 1); + assert_eq!(rows[0].record_id, "current-record"); + assert_eq!(rows[0].change_seq, current.change_seq); +} + +#[test] +fn migration_assigns_existing_records_change_seq_from_insert_order() { + let executor = SqliteExecutor::open_memory().expect("open memory sqlite"); + migrations_run_all_up(&executor, &MIGRATIONS[..1]).expect("apply initial migration"); + let first = insert_pre_change_tracking_record(&executor, "local-a"); + let second = insert_pre_change_tracking_record(&executor, "local-b"); + let store = RuntimeStore::new(executor); + + store.migrate_up().expect("apply change tracking migration"); + let rows = store + .list_records_changed_after(0, 10) + .expect("list changed rows after migration"); + + assert_eq!(rows.len(), 2); + assert_eq!(rows[0].seq, first); + assert_eq!(rows[0].change_seq, first); + assert_eq!(rows[1].seq, second); + assert_eq!(rows[1].change_seq, second); +} + +#[test] +fn migration_repairs_pre_network_source_runtime_constraint() { + let executor = SqliteExecutor::open_memory().expect("open memory sqlite"); + create_pre_network_change_tracking_schema(&executor); + let legacy_seq = insert_pre_network_change_tracking_record(&executor, "legacy-cli", 1); + let store = RuntimeStore::new(executor); + + store + .migrate_up() + .expect("apply network source repair migration"); + let mut input = signed_event("event-network-repaired"); + input.source_runtime = SourceRuntime::Network; + input.event_id = Some("event-network-repaired".to_owned()); + input.raw_event_json = Some(json!({"id":"event-network-repaired","kind":3421})); + let inserted = store + .append_record(&input) + .expect("append repaired network event"); + let rows = store + .list_records_changed_after(0, 10) + .expect("list changed rows after repair"); + + assert_eq!(legacy_seq, 1); + assert_eq!(rows.len(), 2); + assert_eq!(rows[0].record_id, "legacy-cli"); + assert_eq!(rows[0].change_seq, 1); + assert_eq!(rows[0].source_runtime, SourceRuntime::Cli); + assert_eq!(rows[1].record_id, "event-network-repaired"); + assert_eq!(rows[1].seq, inserted.seq); + assert_eq!(rows[1].source_runtime, SourceRuntime::Network); +} + +fn insert_pre_change_tracking_record(executor: &SqliteExecutor, record_id: &str) -> i64 { + let input = local_work(record_id); + let params = json!([ + input.record_id, + input.family.as_str(), + input.status.as_str(), + input.source_runtime.as_str(), + input.created_at_ms, + input.inserted_at_ms, + input.inserted_at_ms, + input.owner_account_id, + input.owner_pubkey, + input.farm_id, + input.listing_addr, + serde_json::to_string(&input.local_work_json).expect("encode local work"), + input.event_id, + input.event_kind, + input.event_pubkey, + input.event_created_at, + input + .event_tags_json + .map(|value| serde_json::to_string(&value).expect("encode tags")), + input.event_content, + input.event_sig, + input + .raw_event_json + .map(|value| serde_json::to_string(&value).expect("encode raw event")), + input.outbox_status.as_str(), + input.relay_set_fingerprint, + input + .relay_delivery_json + .map(|value| serde_json::to_string(&value).expect("encode relay delivery")), + ]) + .to_string(); + let outcome = executor + .exec( + "insert into runtime_store_record( + record_id, + family, + status, + source_runtime, + created_at_ms, + inserted_at_ms, + updated_at_ms, + owner_account_id, + owner_pubkey, + farm_id, + listing_addr, + local_work_json, + event_id, + event_kind, + event_pubkey, + event_created_at, + event_tags_json, + event_content, + event_sig, + raw_event_json, + outbox_status, + relay_set_fingerprint, + relay_delivery_json + ) values(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)", + &params, + ) + .expect("insert pre-change-tracking runtime store record"); + outcome.last_insert_id +} + +fn create_pre_network_change_tracking_schema(executor: &SqliteExecutor) { + let schema = [ + "create table __migrations(id integer primary key, name text not null unique, applied_at text not null default (datetime('now')))", + "create table runtime_store_record ( + seq integer primary key autoincrement, + change_seq integer not null unique, + record_id text not null unique, + family text not null check (family in ('local_work', 'signed_event')), + status text not null check (status in ('local_draft', 'local_saved', 'pending_publish', 'published', 'failed', 'conflict')), + source_runtime text not null check (source_runtime in ('cli', 'app', 'service', 'worker', 'test')), + created_at_ms integer not null, + inserted_at_ms integer not null, + updated_at_ms integer not null, + owner_account_id text, + owner_pubkey text, + farm_id text, + listing_addr text, + local_work_json text, + event_id text, + event_kind integer, + event_pubkey text, + event_created_at integer, + event_tags_json text, + event_content text, + event_sig text, + raw_event_json text, + outbox_status text not null check (outbox_status in ('none', 'pending', 'acknowledged', 'failed')), + relay_set_fingerprint text, + relay_delivery_json text, + check (change_seq >= 1), + check (trim(record_id) <> ''), + check (family <> 'local_work' or local_work_json is not null), + check (family <> 'local_work' or outbox_status = 'none'), + check (family <> 'signed_event' or (event_id is not null and event_kind is not null and event_pubkey is not null and event_sig is not null and raw_event_json is not null)) + )", + "create index runtime_store_record_change_seq_idx on runtime_store_record(change_seq)", + "create index runtime_store_record_event_id_idx on runtime_store_record(event_id)", + "create index runtime_store_record_listing_addr_idx on runtime_store_record(listing_addr)", + "create index runtime_store_record_owner_pubkey_idx on runtime_store_record(owner_pubkey)", + "create index runtime_store_record_status_idx on runtime_store_record(status)", + "create table runtime_store_projection_cursor ( + consumer_id text primary key, + last_change_seq integer not null, + updated_at_ms integer not null, + check (trim(consumer_id) <> ''), + check (last_change_seq >= 0) + )", + ]; + for sql in schema { + executor.exec(sql, "[]").expect("schema statement"); + } + for name in ["0000_runtime_store", "0001_change_tracking"] { + let params = json!([name]).to_string(); + executor + .exec("insert into __migrations(name) values(?)", &params) + .expect("migration marker"); + } +} + +fn insert_pre_network_change_tracking_record( + executor: &SqliteExecutor, + record_id: &str, + change_seq: i64, +) -> i64 { + let input = local_work(record_id); + let params = json!([ + change_seq, + input.record_id, + input.family.as_str(), + input.status.as_str(), + input.source_runtime.as_str(), + input.created_at_ms, + input.inserted_at_ms, + input.inserted_at_ms, + input.owner_account_id, + input.owner_pubkey, + input.farm_id, + input.listing_addr, + serde_json::to_string(&input.local_work_json).expect("encode local work"), + input.event_id, + input.event_kind, + input.event_pubkey, + input.event_created_at, + input + .event_tags_json + .map(|value| serde_json::to_string(&value).expect("encode tags")), + input.event_content, + input.event_sig, + input + .raw_event_json + .map(|value| serde_json::to_string(&value).expect("encode raw event")), + input.outbox_status.as_str(), + input.relay_set_fingerprint, + input + .relay_delivery_json + .map(|value| serde_json::to_string(&value).expect("encode relay delivery")), + ]) + .to_string(); + let outcome = executor + .exec( + "insert into runtime_store_record( + change_seq, + record_id, + family, + status, + source_runtime, + created_at_ms, + inserted_at_ms, + updated_at_ms, + owner_account_id, + owner_pubkey, + farm_id, + listing_addr, + local_work_json, + event_id, + event_kind, + event_pubkey, + event_created_at, + event_tags_json, + event_content, + event_sig, + raw_event_json, + outbox_status, + relay_set_fingerprint, + relay_delivery_json + ) values(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)", + &params, + ) + .expect("insert pre-network runtime store record"); + outcome.last_insert_id +}