commit 3fcba60d5f3639e06e88999fdd4d3eb4d9d014eb
parent 65d7461bbdae768c30e2f76a119a74772eec5ee8
Author: triesap <tyson@radroots.org>
Date: Mon, 24 Aug 2026 01:54:50 +0000
rhi: persist immutable trade evidence
Add the pinned schema-v3 catalog and one typed SQLx transaction for canonical mutations, verified signed events, and configured source observations. Preserve replay idempotency, conflict rollback, bounded reads, redacted errors, and explicit checkpoint/dirty-generation deferral.
Diffstat:
19 files changed, 1778 insertions(+), 51 deletions(-)
diff --git a/AGENTS.md b/AGENTS.md
@@ -86,7 +86,8 @@
- Route accepted trade mutations only through the sealed
`admit_rhi_trade_mutation_event` boundary. Its wire limits come from the
validated configuration, its future-time tolerance is explicit with no
- default, and it must remain pure until Step 184 owns persistence.
+ default, and it remains pure; persistence belongs only to the separate typed
+ repository transaction introduced by Step 184.
- Persist canonical mutation, every distinct signed event carrying it, and
every accepted source observation as separate typed facts. Two events for one
mutation never overwrite one another, and arrival order never selects truth.
@@ -203,14 +204,20 @@
typed RHI repositories; live clients mutate only through the Unix admin
boundary and offline state operations must prove that no daemon writer exists.
- Create-new state begins at the shared schema-v1 baseline and applies the
- governed RHI schema-v2 configuration-binding migration. Retain at most 1,024
- consecutive immutable configuration generations containing only normalized
+ governed RHI schema-v2 configuration-binding migration and schema-v3
+ immutable trade-evidence migration. Retain at most 1,024 consecutive
+ immutable configuration generations containing only normalized
config/evidence-policy digests, public identity, exact contract versions,
- injected apply time, and bounded build identity. Never persist raw TOML,
- paths, URLs, credential references, or protected identity material. Ordinary
- startup must use intent-open, discover source generation under retained
- authority, and match the latest durable binding; configuration apply is an
- exclusive offline operation.
+ injected apply time, and bounded build identity. Persist each canonical
+ mutation, independently signed Nostr event, and accepted configured-source
+ observation as distinct immutable facts in one SQLx transaction. Exact replay
+ is idempotent, conflicts fail closed, observation time remains distinct from
+ authored time, and this persistence step must not advance reconciliation
+ checkpoints or dirty generation. Never persist raw TOML, paths, URLs,
+ credential references, or protected identity material. Ordinary startup must
+ use intent-open, discover source generation under retained authority, and
+ match the latest durable binding; configuration apply is an exclusive offline
+ operation.
- Never hold a database transaction while waiting for a source, relay, DNS,
identity provider, clock, entropy, signing, reduction, or backoff.
- Never prune active jobs/outboxes, migration history, current identity/policy
diff --git a/Cargo.toml b/Cargo.toml
@@ -20,7 +20,7 @@ service = "rhi"
host_feature_profile = "service-host"
nix_material = "absent"
config_contract_version = 1
-state_contract_version = 2
+state_contract_version = 3
admin_contract_version = 1
status_contract_version = 1
provider_contract_version = 1
diff --git a/README b/README
@@ -77,6 +77,24 @@ SQLite access, checkpoint update, or dirty-generation change. Its exact
machine contract is
[`trade_ingest.v1.json`](contracts/services_hardening/trade_ingest.v1.json).
+## Immutable trade-evidence persistence
+
+`RhiStateRepositories::persist_trade_evidence` atomically retains three
+separate immutable facts: one canonical mutation, every distinct valid signed
+event carrying it, and every accepted configured-source observation. Signed
+event identity is the verified Nostr event identifier plus its verified
+signature, so independent valid signatures over the same canonical event are
+not collapsed. Observation time is the injected admission time and remains
+distinct from the event-authored time.
+
+The source observation can be constructed only from a validated RHI
+configuration and an admitted signed event. Exact replay is idempotent;
+conflicting durable mutation or signed-event content fails closed; and the
+three inserts share one SQLx transaction. This operation performs no network
+I/O and does not advance reconciliation checkpoints or dirty generation. The
+machine contract is
+[`trade_evidence_persistence.v1.json`](contracts/services_hardening/trade_evidence_persistence.v1.json).
+
## Existing-state runtime foundation
`open_rhi_runtime_foundation` opens only an already initialized database from
@@ -173,11 +191,14 @@ Path resolution performs no directory creation or filesystem I/O.
RHI owns one `state.sqlite` per service instance. Create-new initialization
starts from the shared schema-v1 baseline and immediately applies the pinned
-schema-v2 migration. Version two contains the six shared immutable
-service-metadata and migration-ledger objects plus one bounded append-only
-`rhi_config_bindings` table and its three enforcement triggers. Exact literal
-SHA-256 values bind the migration, both schema snapshots, and the schema
-catalog. RHI validates every identity before it can become database authority.
+schema-v2 configuration migration and schema-v3 immutable trade-evidence
+migration. Version three contains the six shared immutable service-metadata and
+migration-ledger objects, the bounded append-only `rhi_config_bindings` table,
+and separate immutable tables for canonical mutations, signed Nostr events,
+and accepted source observations with their enforcement triggers and indexes.
+Exact literal SHA-256 values bind both migrations, every schema snapshot, and
+the schema catalog. RHI validates every identity before it can become database
+authority.
The configuration history retains at most 1,024 consecutive generations. Each
row stores only normalized configuration and evidence-policy digests, the
@@ -190,9 +211,10 @@ atomically, treats exact replay idempotently, and explicitly closes state.
Catalog construction itself performs no filesystem or SQLite I/O and owns no
pool, connection, transaction, query, or migration executor. The sealed state
-host alone executes the governed migration and configuration transactions.
-Service-owned evidence and attestation tables remain reserved for their later
-owning checkpoints.
+host alone executes the governed migrations and typed state transactions.
+Reconciliation checkpoints, completion, dirty generation, reports,
+attestations, and publication state remain reserved for their later owning
+checkpoints.
The sealed RHI state-host lifecycle now reserves and initializes a missing
canonical database only through an explicit create-new operation. Ordinary
diff --git a/contracts/api_baselines/rhi.txt b/contracts/api_baselines/rhi.txt
@@ -271,6 +271,16 @@ pub enum rhi::RhiTradeCommandV1
pub rhi::RhiTradeCommandV1::Projection
pub rhi::RhiTradeCommandV1::ReportCurrent
pub rhi::RhiTradeCommandV1::Reports
+pub enum rhi::RhiTradeEvidencePersistenceErrorKind
+pub rhi::RhiTradeEvidencePersistenceErrorKind::CommitOutcomeUnknown
+pub rhi::RhiTradeEvidencePersistenceErrorKind::Encoding
+pub rhi::RhiTradeEvidencePersistenceErrorKind::InvalidMode
+pub rhi::RhiTradeEvidencePersistenceErrorKind::InvalidObservation
+pub rhi::RhiTradeEvidencePersistenceErrorKind::MutationConflict
+pub rhi::RhiTradeEvidencePersistenceErrorKind::SignedEventConflict
+pub rhi::RhiTradeEvidencePersistenceErrorKind::Storage
+impl rhi::RhiTradeEvidencePersistenceErrorKind
+pub const fn rhi::RhiTradeEvidencePersistenceErrorKind::code(self) -> &'static str
pub enum rhi::RhiTradeMutationAdmissionErrorKind
pub rhi::RhiTradeMutationAdmissionErrorKind::AuthoredTimeRejected
pub rhi::RhiTradeMutationAdmissionErrorKind::DuplicateEventField
@@ -341,6 +351,7 @@ pub fn rhi::RhiAdmittedTradeMutationEvent::event_id(&self) -> &radroots_event::i
pub fn rhi::RhiAdmittedTradeMutationEvent::event_kind(&self) -> u32
pub const fn rhi::RhiAdmittedTradeMutationEvent::mutation(&self) -> &radroots_event::trade::TradeMutationEnvelopeV1
pub const fn rhi::RhiAdmittedTradeMutationEvent::mutation_id(&self) -> &radroots_event::id::MutationId
+pub const fn rhi::RhiAdmittedTradeMutationEvent::observed_at_unix_seconds(&self) -> rhi::RhiTradeMutationObservedAtUnixSeconds
pub fn rhi::RhiAdmittedTradeMutationEvent::original_bytes(&self) -> &[u8]
impl core::fmt::Debug for rhi::RhiAdmittedTradeMutationEvent
pub fn rhi::RhiAdmittedTradeMutationEvent::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
@@ -702,6 +713,8 @@ pub const fn rhi::RhiStatePolicyVersions::provider(self) -> u32
pub const fn rhi::RhiStatePolicyVersions::state(self) -> u32
pub const fn rhi::RhiStatePolicyVersions::status(self) -> u32
pub struct rhi::RhiStateRepositories<'host>
+impl rhi::RhiStateRepositories<'_>
+pub async fn rhi::RhiStateRepositories<'_>::persist_trade_evidence(&self, rhi::RhiAdmittedTradeMutationEvent, rhi::RhiTradeSourceObservation) -> core::result::Result<rhi::RhiTradeEvidencePersistenceOutcome, rhi::RhiTradeEvidencePersistenceError>
impl<'host> rhi::RhiStateRepositories<'host>
pub const fn rhi::RhiStateRepositories<'host>::desired_presence(&self) -> rhi::RhiDesiredPresenceRepository<'host>
pub const fn rhi::RhiStateRepositories<'host>::dirty_trades(&self) -> rhi::RhiDirtyTradeRepository<'host>
@@ -747,6 +760,22 @@ pub fn rhi::RhiTimeEntropyAdapters::sample_full_jitter(&self, rhi::RhiJitterBoun
pub fn rhi::RhiTimeEntropyAdapters::system() -> Self
impl core::fmt::Debug for rhi::RhiTimeEntropyAdapters
pub fn rhi::RhiTimeEntropyAdapters::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct rhi::RhiTradeEvidencePersistenceError
+impl rhi::RhiTradeEvidencePersistenceError
+pub const fn rhi::RhiTradeEvidencePersistenceError::code(self) -> &'static str
+pub const fn rhi::RhiTradeEvidencePersistenceError::kind(self) -> rhi::RhiTradeEvidencePersistenceErrorKind
+impl core::error::Error for rhi::RhiTradeEvidencePersistenceError
+impl core::fmt::Debug for rhi::RhiTradeEvidencePersistenceError
+pub fn rhi::RhiTradeEvidencePersistenceError::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+impl core::fmt::Display for rhi::RhiTradeEvidencePersistenceError
+pub fn rhi::RhiTradeEvidencePersistenceError::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct rhi::RhiTradeEvidencePersistenceOutcome
+impl rhi::RhiTradeEvidencePersistenceOutcome
+pub const fn rhi::RhiTradeEvidencePersistenceOutcome::mutation_inserted(self) -> bool
+pub const fn rhi::RhiTradeEvidencePersistenceOutcome::observation_inserted(self) -> bool
+pub const fn rhi::RhiTradeEvidencePersistenceOutcome::signed_event_inserted(self) -> bool
+impl core::fmt::Debug for rhi::RhiTradeEvidencePersistenceOutcome
+pub fn rhi::RhiTradeEvidencePersistenceOutcome::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
pub struct rhi::RhiTradeMutationAdmissionError
impl rhi::RhiTradeMutationAdmissionError
pub const fn rhi::RhiTradeMutationAdmissionError::kind(self) -> rhi::RhiTradeMutationAdmissionErrorKind
@@ -774,6 +803,11 @@ pub struct rhi::RhiTradeMutationObservedAtUnixSeconds(_)
impl rhi::RhiTradeMutationObservedAtUnixSeconds
pub const fn rhi::RhiTradeMutationObservedAtUnixSeconds::get(self) -> u64
pub fn rhi::RhiTradeMutationObservedAtUnixSeconds::new(u64) -> core::result::Result<Self, rhi::RhiTradeMutationAdmissionError>
+pub struct rhi::RhiTradeSourceObservation
+impl rhi::RhiTradeSourceObservation
+pub fn rhi::RhiTradeSourceObservation::from_config(&rhi::RhiConfigDocumentV1, &str, &rhi::RhiAdmittedTradeMutationEvent) -> core::result::Result<Self, rhi::RhiTradeEvidencePersistenceError>
+impl core::fmt::Debug for rhi::RhiTradeSourceObservation
+pub fn rhi::RhiTradeSourceObservation::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
pub struct rhi::RhiTransportAdapters
impl rhi::RhiTransportAdapters
pub fn rhi::RhiTransportAdapters::new(alloc::sync::Arc<dyn radroots_transport::source::EventSource>, alloc::sync::Arc<dyn radroots_transport::source::EventSubscriber>, alloc::sync::Arc<dyn radroots_transport::sink::EventSink>) -> Self
@@ -864,12 +898,16 @@ pub const rhi::RHI_STATE_SCHEMA_VERSION_1_SHA256: [u8; 32]
pub const rhi::RHI_STATE_SCHEMA_VERSION_2_MIGRATION_SHA256: [u8; 32]
pub const rhi::RHI_STATE_SCHEMA_VERSION_2_OBJECT_COUNT: u32
pub const rhi::RHI_STATE_SCHEMA_VERSION_2_SHA256: [u8; 32]
+pub const rhi::RHI_STATE_SCHEMA_VERSION_3_MIGRATION_SHA256: [u8; 32]
+pub const rhi::RHI_STATE_SCHEMA_VERSION_3_OBJECT_COUNT: u32
+pub const rhi::RHI_STATE_SCHEMA_VERSION_3_SHA256: [u8; 32]
pub const rhi::RHI_STATUS_CONTRACT_VERSION: u32
pub const rhi::RHI_TRADE_EVENT_EXTRA_FIELD_MAX_COUNT: usize
pub const rhi::RHI_TRADE_EVENT_EXTRA_JSON_MAX_BYTES: usize
pub const rhi::RHI_TRADE_EVENT_ID_MAX_BYTES: usize
pub const rhi::RHI_TRADE_EVENT_PUBLIC_KEY_MAX_BYTES: usize
pub const rhi::RHI_TRADE_EVENT_SIGNATURE_MAX_BYTES: usize
+pub const rhi::RHI_TRADE_EVIDENCE_PERSISTENCE_CONTRACT_VERSION: u32
pub const rhi::RHI_TRADE_INGEST_CONTRACT_VERSION: u32
pub const rhi::RHI_WRAPPING_CREDENTIAL_ARTIFACT_BYTES: usize
pub const rhi::RHI_WRAPPING_CREDENTIAL_CONTRACT_VERSION: u32
diff --git a/contracts/services_hardening/trade_evidence_persistence.v1.json b/contracts/services_hardening/trade_evidence_persistence.v1.json
@@ -0,0 +1,56 @@
+{
+ "schema": "radroots.rhi.trade-evidence-persistence.v1",
+ "contract_version": 1,
+ "input": {
+ "event": "sealed_signature_verified_trade_mutation_event",
+ "observation": "sealed_configured_nostr_relay_source_observation",
+ "observation_time": "injected_positive_unix_seconds",
+ "event_authored_time": "distinct_untrusted_policy_admitted_unix_seconds"
+ },
+ "transaction": "one_governed_sqlx_transaction",
+ "facts": {
+ "canonical_mutation": {
+ "table": "trade_mutations",
+ "identity": "content_derived_mutation_id",
+ "content": "exact_canonical_mutation_json",
+ "event_authored_time_column": "absent_event_fact_only",
+ "conflicting_identity_reuse": "reject"
+ },
+ "signed_event": {
+ "table": "nostr_events",
+ "identity": ["verified_event_id", "verified_event_signature"],
+ "content": "canonical_nip01_signed_event_json",
+ "multiple_signatures_for_one_event_id": "retained_separately",
+ "mutation_relationship": "many_signed_events_to_one_mutation"
+ },
+ "source_observation": {
+ "table": "relay_observations",
+ "identity": ["source_id", "selector", "policy_digest", "event_id", "event_signature", "observed_at_unix_s"],
+ "source_id": "exact_validated_configuration_member",
+ "selector": "trade_mutation_lineage_v1",
+ "policy_digest": "exact_normalized_evidence_policy_sha256"
+ }
+ },
+ "immutability": {
+ "updates": "reject",
+ "deletes": "reject",
+ "exact_replay": "idempotent",
+ "arrival_order": "non_authoritative"
+ },
+ "effects": {
+ "filesystem": false,
+ "network": false,
+ "clock_read": false,
+ "checkpoint": false,
+ "dirty_generation": false
+ },
+ "errors": "crate_owned_source_free_redacted",
+ "deferred": [
+ "nostr_relay_source_adapter",
+ "source_completion",
+ "checkpoint_advance",
+ "dirty_generation",
+ "reconciliation",
+ "publication"
+ ]
+}
diff --git a/contracts/services_hardening/trade_ingest.v1.json b/contracts/services_hardening/trade_ingest.v1.json
@@ -52,5 +52,6 @@
},
"errors": "crate_owned_source_free_redacted",
"effects": { "filesystem": false, "sqlite": false, "clock_read": false, "network": false, "checkpoint": false, "dirty_generation": false },
- "deferred": ["mutation_event_and_provenance_persistence", "idempotency_and_conflict_resolution", "source_adapter_and_completion", "checkpoint_and_dirty_generation", "reconciliation", "publication"]
+ "persistence_contract": "contracts/services_hardening/trade_evidence_persistence.v1.json",
+ "deferred": ["source_adapter_and_completion", "checkpoint_and_dirty_generation", "reconciliation", "publication"]
}
diff --git a/radroots.service.source-lock.v2.toml b/radroots.service.source-lock.v2.toml
@@ -16,7 +16,7 @@ material = "absent"
[contract_versions]
config = 1
-state = 2
+state = 3
admin = 1
status = 1
provider = 1
diff --git a/src/lib.rs b/src/lib.rs
@@ -17,6 +17,7 @@ mod state_host;
mod state_maintenance;
mod state_metadata;
mod state_repository;
+mod state_trade;
mod trade_ingest;
pub use adapters::nostr::event::NostrEventAdapter;
@@ -84,8 +85,9 @@ pub use state_catalog::{
RHI_STATE_SCHEMA_VERSION, RHI_STATE_SCHEMA_VERSION_1_OBJECT_COUNT,
RHI_STATE_SCHEMA_VERSION_1_SHA256, RHI_STATE_SCHEMA_VERSION_2_MIGRATION_SHA256,
RHI_STATE_SCHEMA_VERSION_2_OBJECT_COUNT, RHI_STATE_SCHEMA_VERSION_2_SHA256,
- RhiStateCatalogError, RhiStateCatalogErrorKind, rhi_migration_catalog, rhi_schema_catalog,
- validate_rhi_state_catalogs,
+ RHI_STATE_SCHEMA_VERSION_3_MIGRATION_SHA256, RHI_STATE_SCHEMA_VERSION_3_OBJECT_COUNT,
+ RHI_STATE_SCHEMA_VERSION_3_SHA256, RhiStateCatalogError, RhiStateCatalogErrorKind,
+ rhi_migration_catalog, rhi_schema_catalog, validate_rhi_state_catalogs,
};
pub use state_config::{
RHI_CONFIG_BINDING_MAX_GENERATIONS, RhiConfigApplyError, RhiConfigApplyErrorKind,
@@ -119,6 +121,11 @@ pub use state_repository::{
RhiStateRepositoryKind, RhiStateRepositoryWriteClass, RhiSupersessionRepository,
rhi_state_repository_descriptors,
};
+pub use state_trade::{
+ RHI_TRADE_EVIDENCE_PERSISTENCE_CONTRACT_VERSION, RhiTradeEvidencePersistenceError,
+ RhiTradeEvidencePersistenceErrorKind, RhiTradeEvidencePersistenceOutcome,
+ RhiTradeSourceObservation,
+};
pub use trade_ingest::{
RHI_TRADE_EVENT_EXTRA_FIELD_MAX_COUNT, RHI_TRADE_EVENT_EXTRA_JSON_MAX_BYTES,
RHI_TRADE_EVENT_ID_MAX_BYTES, RHI_TRADE_EVENT_PUBLIC_KEY_MAX_BYTES,
diff --git a/src/state_catalog.rs b/src/state_catalog.rs
@@ -12,7 +12,7 @@ use radroots_service_sqlite::{
pub const RHI_STATE_BASE_SCHEMA_VERSION: u32 = 1;
/// The newest governed RHI state schema understood by this binary.
-pub const RHI_STATE_SCHEMA_VERSION: u32 = 2;
+pub const RHI_STATE_SCHEMA_VERSION: u32 = 3;
/// The shared metadata and migration-ledger objects present at schema v1.
pub const RHI_STATE_SCHEMA_VERSION_1_OBJECT_COUNT: u32 = 6;
@@ -20,10 +20,13 @@ pub const RHI_STATE_SCHEMA_VERSION_1_OBJECT_COUNT: u32 = 6;
/// The shared objects plus the bounded append-only RHI configuration history.
pub const RHI_STATE_SCHEMA_VERSION_2_OBJECT_COUNT: u32 = 10;
+/// The shared objects, configuration history, and immutable trade evidence.
+pub const RHI_STATE_SCHEMA_VERSION_3_OBJECT_COUNT: u32 = 22;
+
/// SHA-256 identity of the ordered migration catalog rooted at schema v1.
pub const RHI_MIGRATION_CATALOG_SHA256: [u8; 32] = [
- 0xb6, 0x40, 0xa9, 0x09, 0x5d, 0x53, 0x18, 0xdb, 0xfd, 0x0a, 0xfb, 0xfb, 0xe6, 0xe0, 0x52, 0x81,
- 0x11, 0x3a, 0x9c, 0xef, 0x45, 0x64, 0x80, 0x5c, 0x02, 0x1c, 0x36, 0x8c, 0x3c, 0x52, 0xfc, 0x08,
+ 0x14, 0x04, 0x60, 0x48, 0xb4, 0x68, 0x83, 0x6f, 0x26, 0x02, 0xec, 0x51, 0xe5, 0x38, 0xf2, 0xa9,
+ 0x8b, 0x71, 0xc9, 0x45, 0xf8, 0x3d, 0x93, 0x3d, 0xcc, 0xb8, 0x60, 0x3b, 0x84, 0xf8, 0x65, 0xf0,
];
/// SHA-256 identity of the exact schema-v1 object snapshot.
@@ -44,10 +47,22 @@ pub const RHI_STATE_SCHEMA_VERSION_2_SHA256: [u8; 32] = [
0x8e, 0x08, 0x8b, 0x8c, 0x26, 0xd1, 0x5b, 0xa4, 0x51, 0x33, 0xca, 0x5e, 0x9b, 0x73, 0x15, 0xa9,
];
+/// SHA-256 identity of the schema-v3 trade-evidence migration.
+pub const RHI_STATE_SCHEMA_VERSION_3_MIGRATION_SHA256: [u8; 32] = [
+ 0x07, 0xb0, 0x98, 0xc3, 0x93, 0x14, 0x0a, 0xfe, 0xc2, 0x22, 0xfc, 0xe7, 0x6e, 0xc6, 0x68, 0x69,
+ 0x8d, 0x50, 0xf1, 0xd2, 0x37, 0x85, 0x85, 0x68, 0x73, 0xe3, 0x04, 0x45, 0xaf, 0x1a, 0x01, 0x3f,
+];
+
+/// SHA-256 identity of the exact schema-v3 object snapshot.
+pub const RHI_STATE_SCHEMA_VERSION_3_SHA256: [u8; 32] = [
+ 0xfd, 0x96, 0x22, 0x64, 0x05, 0xab, 0x68, 0x65, 0x5a, 0xee, 0x00, 0xf6, 0x83, 0xf6, 0x02, 0x3c,
+ 0x7a, 0xab, 0x2d, 0xbd, 0x23, 0xfe, 0xad, 0xac, 0x16, 0x53, 0x33, 0x49, 0x0d, 0x6f, 0x0a, 0xd5,
+];
+
/// SHA-256 identity of the schema catalog bound to the migration catalog.
pub const RHI_STATE_SCHEMA_CATALOG_SHA256: [u8; 32] = [
- 0x1b, 0x7f, 0x73, 0x59, 0xb6, 0x2e, 0xd7, 0xdc, 0xd4, 0x76, 0x28, 0x9e, 0x46, 0x69, 0x3d, 0x9b,
- 0x3d, 0x13, 0x69, 0xdc, 0xaf, 0x7d, 0x59, 0x52, 0xf0, 0x32, 0x9a, 0x5c, 0x86, 0xf8, 0xb8, 0x5f,
+ 0x13, 0x25, 0xa4, 0x1b, 0x90, 0xab, 0xfc, 0x7d, 0x52, 0x3b, 0xbf, 0xe8, 0x35, 0x00, 0xa1, 0xb1,
+ 0xa7, 0x3e, 0x49, 0x2c, 0xc9, 0x30, 0xd6, 0xea, 0xd6, 0xc7, 0x1e, 0x29, 0xd5, 0x44, 0x74, 0x36,
];
macro_rules! rhi_config_bindings_table_sql {
@@ -138,6 +153,204 @@ const CREATE_RHI_CONFIG_BINDINGS_MIGRATION_SQL: &str = concat!(
";",
);
+macro_rules! trade_mutations_table_sql {
+ () => {
+ r#"CREATE TABLE trade_mutations (
+ mutation_id BLOB NOT NULL PRIMARY KEY CHECK (length(mutation_id) = 32),
+ trade_id BLOB NOT NULL CHECK (length(trade_id) = 16),
+ contract_id TEXT NOT NULL CHECK (contract_id IN (
+ 'radroots.trade.proposal.v1',
+ 'radroots.trade.decision.v1',
+ 'radroots.trade.revision_proposal.v1',
+ 'radroots.trade.revision_decision.v1',
+ 'radroots.trade.cancellation.v1'
+ )),
+ schema_version INTEGER NOT NULL CHECK (schema_version = 1),
+ event_kind INTEGER NOT NULL CHECK (event_kind IN (3470, 3471, 3472, 3473, 3474)),
+ author_pubkey BLOB NOT NULL CHECK (length(author_pubkey) = 32),
+ canonical_content BLOB NOT NULL
+ CHECK (length(canonical_content) BETWEEN 1 AND 131072)
+) STRICT"#
+ };
+}
+
+macro_rules! trade_mutations_by_trade_sql {
+ () => {
+ r#"CREATE INDEX trade_mutations_by_trade
+ON trade_mutations (trade_id, mutation_id)"#
+ };
+}
+
+macro_rules! nostr_events_table_sql {
+ () => {
+ r#"CREATE TABLE nostr_events (
+ event_id BLOB NOT NULL CHECK (length(event_id) = 32),
+ event_signature BLOB NOT NULL CHECK (length(event_signature) = 64),
+ mutation_id BLOB NOT NULL CHECK (length(mutation_id) = 32)
+ REFERENCES trade_mutations (mutation_id),
+ author_pubkey BLOB NOT NULL CHECK (length(author_pubkey) = 32),
+ event_kind INTEGER NOT NULL CHECK (event_kind IN (3470, 3471, 3472, 3473, 3474)),
+ authored_at_unix_s INTEGER NOT NULL
+ CHECK (authored_at_unix_s BETWEEN 0 AND 9223372036854775807),
+ canonical_event_json BLOB NOT NULL
+ CHECK (length(canonical_event_json) BETWEEN 1 AND 524288),
+ PRIMARY KEY (event_id, event_signature)
+) STRICT"#
+ };
+}
+
+macro_rules! nostr_events_by_mutation_sql {
+ () => {
+ r#"CREATE INDEX nostr_events_by_mutation
+ON nostr_events (mutation_id, authored_at_unix_s, event_id)"#
+ };
+}
+
+macro_rules! relay_observations_table_sql {
+ () => {
+ r#"CREATE TABLE relay_observations (
+ source_id TEXT NOT NULL
+ CHECK (length(CAST(source_id AS BLOB)) BETWEEN 1 AND 64)
+ CHECK (source_id NOT GLOB '*[^a-z0-9_-]*')
+ CHECK (substr(source_id, 1, 1) GLOB '[a-z]'),
+ selector_id TEXT NOT NULL CHECK (selector_id = 'trade_mutation_lineage_v1'),
+ evidence_policy_sha256 BLOB NOT NULL CHECK (length(evidence_policy_sha256) = 32),
+ event_id BLOB NOT NULL CHECK (length(event_id) = 32),
+ event_signature BLOB NOT NULL CHECK (length(event_signature) = 64),
+ observed_at_unix_s INTEGER NOT NULL
+ CHECK (observed_at_unix_s BETWEEN 1 AND 9223372036854775807),
+ FOREIGN KEY (event_id, event_signature)
+ REFERENCES nostr_events (event_id, event_signature),
+ PRIMARY KEY (
+ source_id, selector_id, evidence_policy_sha256, event_id, event_signature,
+ observed_at_unix_s
+ )
+) STRICT"#
+ };
+}
+
+macro_rules! relay_observations_by_event_sql {
+ () => {
+ r#"CREATE INDEX relay_observations_by_event
+ON relay_observations (event_id, event_signature, observed_at_unix_s, source_id)"#
+ };
+}
+
+macro_rules! immutable_no_update_sql {
+ ($trigger:literal, $table:literal, $message:literal) => {
+ concat!(
+ "CREATE TRIGGER ",
+ $trigger,
+ "\nBEFORE UPDATE ON ",
+ $table,
+ "\nBEGIN\n SELECT RAISE(ABORT, '",
+ $message,
+ "');\nEND"
+ )
+ };
+}
+
+macro_rules! immutable_no_delete_sql {
+ ($trigger:literal, $table:literal, $message:literal) => {
+ concat!(
+ "CREATE TRIGGER ",
+ $trigger,
+ "\nBEFORE DELETE ON ",
+ $table,
+ "\nBEGIN\n SELECT RAISE(ABORT, '",
+ $message,
+ "');\nEND"
+ )
+ };
+}
+
+pub(crate) const CREATE_TRADE_MUTATIONS_TABLE_SQL: &str = trade_mutations_table_sql!();
+const CREATE_TRADE_MUTATIONS_BY_TRADE_SQL: &str = trade_mutations_by_trade_sql!();
+pub(crate) const CREATE_NOSTR_EVENTS_TABLE_SQL: &str = nostr_events_table_sql!();
+const CREATE_NOSTR_EVENTS_BY_MUTATION_SQL: &str = nostr_events_by_mutation_sql!();
+pub(crate) const CREATE_RELAY_OBSERVATIONS_TABLE_SQL: &str = relay_observations_table_sql!();
+const CREATE_RELAY_OBSERVATIONS_BY_EVENT_SQL: &str = relay_observations_by_event_sql!();
+const CREATE_TRADE_MUTATIONS_NO_UPDATE_SQL: &str = immutable_no_update_sql!(
+ "trade_mutations_no_update",
+ "trade_mutations",
+ "trade mutation evidence is immutable"
+);
+const CREATE_TRADE_MUTATIONS_NO_DELETE_SQL: &str = immutable_no_delete_sql!(
+ "trade_mutations_no_delete",
+ "trade_mutations",
+ "trade mutation evidence is retained"
+);
+const CREATE_NOSTR_EVENTS_NO_UPDATE_SQL: &str = immutable_no_update_sql!(
+ "nostr_events_no_update",
+ "nostr_events",
+ "signed event evidence is immutable"
+);
+const CREATE_NOSTR_EVENTS_NO_DELETE_SQL: &str = immutable_no_delete_sql!(
+ "nostr_events_no_delete",
+ "nostr_events",
+ "signed event evidence is retained"
+);
+const CREATE_RELAY_OBSERVATIONS_NO_UPDATE_SQL: &str = immutable_no_update_sql!(
+ "relay_observations_no_update",
+ "relay_observations",
+ "source observation evidence is immutable"
+);
+const CREATE_RELAY_OBSERVATIONS_NO_DELETE_SQL: &str = immutable_no_delete_sql!(
+ "relay_observations_no_delete",
+ "relay_observations",
+ "source observation evidence is retained"
+);
+const CREATE_TRADE_EVIDENCE_MIGRATION_SQL: &str = concat!(
+ trade_mutations_table_sql!(),
+ ";\n",
+ trade_mutations_by_trade_sql!(),
+ ";\n",
+ nostr_events_table_sql!(),
+ ";\n",
+ nostr_events_by_mutation_sql!(),
+ ";\n",
+ relay_observations_table_sql!(),
+ ";\n",
+ relay_observations_by_event_sql!(),
+ ";\n",
+ immutable_no_update_sql!(
+ "trade_mutations_no_update",
+ "trade_mutations",
+ "trade mutation evidence is immutable"
+ ),
+ ";\n",
+ immutable_no_delete_sql!(
+ "trade_mutations_no_delete",
+ "trade_mutations",
+ "trade mutation evidence is retained"
+ ),
+ ";\n",
+ immutable_no_update_sql!(
+ "nostr_events_no_update",
+ "nostr_events",
+ "signed event evidence is immutable"
+ ),
+ ";\n",
+ immutable_no_delete_sql!(
+ "nostr_events_no_delete",
+ "nostr_events",
+ "signed event evidence is retained"
+ ),
+ ";\n",
+ immutable_no_update_sql!(
+ "relay_observations_no_update",
+ "relay_observations",
+ "source observation evidence is immutable"
+ ),
+ ";\n",
+ immutable_no_delete_sql!(
+ "relay_observations_no_delete",
+ "relay_observations",
+ "source observation evidence is retained"
+ ),
+ ";",
+);
+
const RHI_CONFIG_BINDINGS_TABLE_SHA256: [u8; 32] = [
0x4d, 0x6e, 0x8f, 0xff, 0xda, 0x43, 0xe6, 0xf5, 0x3e, 0x23, 0x77, 0xd2, 0x77, 0xa4, 0x52, 0x9e,
0x63, 0x3e, 0xaf, 0xb6, 0xea, 0xa2, 0xad, 0xd7, 0x56, 0xde, 0x0d, 0xc9, 0x24, 0xc5, 0x77, 0xeb,
@@ -154,6 +367,54 @@ const RHI_CONFIG_BINDINGS_NO_DELETE_SHA256: [u8; 32] = [
0x5d, 0x26, 0x82, 0xe9, 0xf2, 0xdc, 0x84, 0x97, 0x61, 0xd1, 0xd7, 0x10, 0xdc, 0xda, 0x75, 0xea,
0x40, 0x6d, 0x10, 0x95, 0xae, 0x1b, 0xc0, 0xad, 0xdf, 0x70, 0x64, 0xec, 0x8b, 0xed, 0x0b, 0x90,
];
+const TRADE_MUTATIONS_TABLE_SHA256: [u8; 32] = [
+ 0x86, 0x4f, 0x45, 0xc0, 0x87, 0xe3, 0x87, 0x71, 0x53, 0xb8, 0x6a, 0xa1, 0x88, 0x6c, 0x73, 0x3e,
+ 0x56, 0x9f, 0x1b, 0x8b, 0xad, 0x8e, 0x9c, 0xbf, 0xab, 0xc4, 0x84, 0x7d, 0x33, 0x33, 0xea, 0xb5,
+];
+const TRADE_MUTATIONS_BY_TRADE_SHA256: [u8; 32] = [
+ 0xcf, 0xd1, 0x79, 0x90, 0xf6, 0x0a, 0x18, 0x94, 0x0a, 0xf9, 0xe6, 0xa8, 0x93, 0xda, 0x9e, 0x9a,
+ 0xc7, 0x5d, 0x2f, 0x69, 0x28, 0xd9, 0xb6, 0x7c, 0x25, 0x06, 0x64, 0x8d, 0x94, 0x9f, 0x97, 0x2e,
+];
+const NOSTR_EVENTS_TABLE_SHA256: [u8; 32] = [
+ 0xcd, 0x08, 0xad, 0x41, 0xb4, 0xd7, 0x88, 0xde, 0xc1, 0x23, 0xc9, 0x83, 0x27, 0x17, 0xb8, 0xab,
+ 0x26, 0x0d, 0x58, 0x5c, 0x71, 0x56, 0x68, 0x6c, 0xaa, 0xf0, 0x5f, 0x07, 0xba, 0xfc, 0x72, 0x12,
+];
+const NOSTR_EVENTS_BY_MUTATION_SHA256: [u8; 32] = [
+ 0x57, 0xd2, 0x8a, 0x14, 0xed, 0x74, 0x84, 0xef, 0xa4, 0x1f, 0x78, 0xe8, 0xbf, 0x6e, 0x85, 0x49,
+ 0xfa, 0x50, 0x76, 0x65, 0xc7, 0x95, 0x32, 0x34, 0xe0, 0xc5, 0xb4, 0xcb, 0x03, 0xb7, 0xa4, 0x18,
+];
+const RELAY_OBSERVATIONS_TABLE_SHA256: [u8; 32] = [
+ 0xb0, 0x4d, 0x44, 0x0b, 0xe0, 0x1b, 0x93, 0x2e, 0xea, 0x53, 0x0b, 0x6a, 0x2e, 0xb2, 0x13, 0xe5,
+ 0x19, 0xab, 0xbf, 0xc2, 0xd2, 0xff, 0xa3, 0xb8, 0x86, 0x98, 0xb1, 0xc7, 0x28, 0x17, 0x86, 0xfc,
+];
+const RELAY_OBSERVATIONS_BY_EVENT_SHA256: [u8; 32] = [
+ 0xdd, 0x69, 0x9f, 0x49, 0x98, 0x67, 0x0f, 0x78, 0xb4, 0x62, 0xcb, 0x1c, 0xd5, 0x17, 0x2a, 0xdb,
+ 0x8d, 0xa1, 0x68, 0x7c, 0x51, 0x74, 0xb3, 0x34, 0x43, 0xf1, 0xc6, 0x6f, 0x41, 0x03, 0xf7, 0x82,
+];
+const TRADE_MUTATIONS_NO_UPDATE_SHA256: [u8; 32] = [
+ 0xd7, 0xeb, 0x61, 0x74, 0x43, 0x6a, 0xe3, 0x46, 0xf8, 0x29, 0x31, 0x37, 0x33, 0x93, 0xc0, 0xe9,
+ 0x0c, 0xe7, 0xf3, 0x06, 0xd6, 0xad, 0xc9, 0xe7, 0xc1, 0xe0, 0x19, 0x3e, 0xa2, 0x46, 0x7d, 0xbe,
+];
+const TRADE_MUTATIONS_NO_DELETE_SHA256: [u8; 32] = [
+ 0xde, 0x61, 0x5b, 0xe4, 0x8f, 0xac, 0xc6, 0xd9, 0xae, 0xb1, 0x11, 0xe7, 0xbb, 0xd3, 0xc7, 0x31,
+ 0x17, 0x6e, 0xc6, 0xe4, 0x82, 0x61, 0x77, 0xba, 0x1c, 0xb7, 0x49, 0x44, 0x02, 0xb1, 0xba, 0xc4,
+];
+const NOSTR_EVENTS_NO_UPDATE_SHA256: [u8; 32] = [
+ 0x3e, 0x6e, 0xec, 0x38, 0x89, 0xd7, 0xfd, 0x25, 0xde, 0x26, 0xe2, 0x03, 0x6e, 0xa5, 0xa5, 0x2a,
+ 0xd5, 0xa0, 0xeb, 0xd6, 0x04, 0x62, 0x61, 0x08, 0xaa, 0x80, 0x72, 0x56, 0xe9, 0xb6, 0x0a, 0x89,
+];
+const NOSTR_EVENTS_NO_DELETE_SHA256: [u8; 32] = [
+ 0x17, 0x6b, 0x15, 0x16, 0xce, 0xe2, 0xf1, 0xfc, 0x62, 0xa7, 0x40, 0x90, 0x01, 0x79, 0x51, 0xae,
+ 0x15, 0x2c, 0xbe, 0x79, 0x53, 0xd0, 0x7c, 0x86, 0xa0, 0xae, 0xca, 0xd3, 0x19, 0x1e, 0xf5, 0xf6,
+];
+const RELAY_OBSERVATIONS_NO_UPDATE_SHA256: [u8; 32] = [
+ 0x6d, 0x20, 0xb3, 0xf0, 0xef, 0x1a, 0x26, 0x55, 0xff, 0xd7, 0x46, 0x5c, 0xce, 0xfa, 0x4c, 0x64,
+ 0x25, 0x26, 0x64, 0xe2, 0x1b, 0xdc, 0xd8, 0x59, 0xdc, 0xb2, 0xde, 0x15, 0xf0, 0x55, 0xf6, 0xaa,
+];
+const RELAY_OBSERVATIONS_NO_DELETE_SHA256: [u8; 32] = [
+ 0xe9, 0xaa, 0x66, 0xff, 0x29, 0xa0, 0x62, 0xc9, 0xf9, 0x97, 0x27, 0x0f, 0x59, 0xad, 0x63, 0x65,
+ 0xdc, 0x4d, 0x88, 0xc6, 0x59, 0x4b, 0xe5, 0xe5, 0xfe, 0x44, 0xf4, 0x44, 0x6d, 0x00, 0xb7, 0xe0,
+];
/// Stable classes for invalid embedded RHI catalog definitions.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
@@ -233,10 +494,17 @@ pub fn rhi_migration_catalog() -> Result<MigrationCatalog, RhiStateCatalogError>
MigrationChecksum::from_bytes(RHI_STATE_SCHEMA_VERSION_2_MIGRATION_SHA256),
)
.map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::MigrationCatalog))?;
- let catalog = MigrationCatalog::new([configuration])
+ let trade_evidence = MigrationDescriptor::sql(
+ 3,
+ "create_immutable_trade_evidence",
+ CREATE_TRADE_EVIDENCE_MIGRATION_SQL,
+ MigrationChecksum::from_bytes(RHI_STATE_SCHEMA_VERSION_3_MIGRATION_SHA256),
+ )
+ .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::MigrationCatalog))?;
+ let catalog = MigrationCatalog::new([configuration, trade_evidence])
.map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::MigrationCatalog))?;
if catalog.current_version() != RHI_STATE_SCHEMA_VERSION
- || catalog.descriptors().len() != 1
+ || catalog.descriptors().len() != 2
|| catalog.digest().as_bytes() != &RHI_MIGRATION_CATALOG_SHA256
{
return Err(RhiStateCatalogError::new(
@@ -256,12 +524,18 @@ pub fn rhi_schema_catalog() -> Result<SchemaCatalog, RhiStateCatalogError> {
)
.map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog))?;
let version_two = SchemaVersionCatalog::new(
- RHI_STATE_SCHEMA_VERSION,
+ 2,
rhi_config_binding_objects()?,
SchemaDigest::from_bytes(RHI_STATE_SCHEMA_VERSION_2_SHA256),
)
.map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog))?;
- let catalog = SchemaCatalog::new(&migrations, [version_one, version_two])
+ let version_three = SchemaVersionCatalog::new(
+ RHI_STATE_SCHEMA_VERSION,
+ rhi_schema_version_three_objects()?,
+ SchemaDigest::from_bytes(RHI_STATE_SCHEMA_VERSION_3_SHA256),
+ )
+ .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog))?;
+ let catalog = SchemaCatalog::new(&migrations, [version_one, version_two, version_three])
.map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog))?;
validate_rhi_state_catalogs(&migrations, &catalog)?;
Ok(catalog)
@@ -274,16 +548,19 @@ pub fn validate_rhi_state_catalogs(
) -> Result<(), RhiStateCatalogError> {
let versions = schema.versions();
let valid = migrations.current_version() == RHI_STATE_SCHEMA_VERSION
- && migrations.descriptors().len() == 1
+ && migrations.descriptors().len() == 2
&& migrations.digest().as_bytes() == &RHI_MIGRATION_CATALOG_SHA256
&& schema.migration_catalog_digest() == migrations.digest()
- && versions.len() == 2
+ && versions.len() == 3
&& versions[0].version() == RHI_STATE_BASE_SCHEMA_VERSION
&& versions[0].object_count() == RHI_STATE_SCHEMA_VERSION_1_OBJECT_COUNT
&& versions[0].digest().as_bytes() == &RHI_STATE_SCHEMA_VERSION_1_SHA256
- && versions[1].version() == RHI_STATE_SCHEMA_VERSION
+ && versions[1].version() == 2
&& versions[1].object_count() == RHI_STATE_SCHEMA_VERSION_2_OBJECT_COUNT
&& versions[1].digest().as_bytes() == &RHI_STATE_SCHEMA_VERSION_2_SHA256
+ && versions[2].version() == RHI_STATE_SCHEMA_VERSION
+ && versions[2].object_count() == RHI_STATE_SCHEMA_VERSION_3_OBJECT_COUNT
+ && versions[2].digest().as_bytes() == &RHI_STATE_SCHEMA_VERSION_3_SHA256
&& schema.digest().as_bytes() == &RHI_STATE_SCHEMA_CATALOG_SHA256;
if valid {
Ok(())
@@ -330,3 +607,108 @@ fn rhi_config_binding_objects() -> Result<[SchemaObject; 4], RhiStateCatalogErro
.map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog))?,
])
}
+
+fn rhi_schema_version_three_objects() -> Result<Vec<SchemaObject>, RhiStateCatalogError> {
+ let mut objects = rhi_config_binding_objects()?.to_vec();
+ objects.extend(rhi_trade_evidence_objects()?);
+ Ok(objects)
+}
+
+fn rhi_trade_evidence_objects() -> Result<[SchemaObject; 12], RhiStateCatalogError> {
+ let object = |kind, name, table_name, sql, digest| {
+ SchemaObject::new(
+ kind,
+ name,
+ table_name,
+ sql,
+ SchemaDigest::from_bytes(digest),
+ )
+ .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog))
+ };
+ Ok([
+ object(
+ SchemaObjectKind::Table,
+ "trade_mutations",
+ "trade_mutations",
+ CREATE_TRADE_MUTATIONS_TABLE_SQL,
+ TRADE_MUTATIONS_TABLE_SHA256,
+ )?,
+ object(
+ SchemaObjectKind::Index,
+ "trade_mutations_by_trade",
+ "trade_mutations",
+ CREATE_TRADE_MUTATIONS_BY_TRADE_SQL,
+ TRADE_MUTATIONS_BY_TRADE_SHA256,
+ )?,
+ object(
+ SchemaObjectKind::Table,
+ "nostr_events",
+ "nostr_events",
+ CREATE_NOSTR_EVENTS_TABLE_SQL,
+ NOSTR_EVENTS_TABLE_SHA256,
+ )?,
+ object(
+ SchemaObjectKind::Index,
+ "nostr_events_by_mutation",
+ "nostr_events",
+ CREATE_NOSTR_EVENTS_BY_MUTATION_SQL,
+ NOSTR_EVENTS_BY_MUTATION_SHA256,
+ )?,
+ object(
+ SchemaObjectKind::Table,
+ "relay_observations",
+ "relay_observations",
+ CREATE_RELAY_OBSERVATIONS_TABLE_SQL,
+ RELAY_OBSERVATIONS_TABLE_SHA256,
+ )?,
+ object(
+ SchemaObjectKind::Index,
+ "relay_observations_by_event",
+ "relay_observations",
+ CREATE_RELAY_OBSERVATIONS_BY_EVENT_SQL,
+ RELAY_OBSERVATIONS_BY_EVENT_SHA256,
+ )?,
+ object(
+ SchemaObjectKind::Trigger,
+ "trade_mutations_no_update",
+ "trade_mutations",
+ CREATE_TRADE_MUTATIONS_NO_UPDATE_SQL,
+ TRADE_MUTATIONS_NO_UPDATE_SHA256,
+ )?,
+ object(
+ SchemaObjectKind::Trigger,
+ "trade_mutations_no_delete",
+ "trade_mutations",
+ CREATE_TRADE_MUTATIONS_NO_DELETE_SQL,
+ TRADE_MUTATIONS_NO_DELETE_SHA256,
+ )?,
+ object(
+ SchemaObjectKind::Trigger,
+ "nostr_events_no_update",
+ "nostr_events",
+ CREATE_NOSTR_EVENTS_NO_UPDATE_SQL,
+ NOSTR_EVENTS_NO_UPDATE_SHA256,
+ )?,
+ object(
+ SchemaObjectKind::Trigger,
+ "nostr_events_no_delete",
+ "nostr_events",
+ CREATE_NOSTR_EVENTS_NO_DELETE_SQL,
+ NOSTR_EVENTS_NO_DELETE_SHA256,
+ )?,
+ object(
+ SchemaObjectKind::Trigger,
+ "relay_observations_no_update",
+ "relay_observations",
+ CREATE_RELAY_OBSERVATIONS_NO_UPDATE_SQL,
+ RELAY_OBSERVATIONS_NO_UPDATE_SHA256,
+ )?,
+ object(
+ SchemaObjectKind::Trigger,
+ "relay_observations_no_delete",
+ "relay_observations",
+ CREATE_RELAY_OBSERVATIONS_NO_DELETE_SQL,
+ RELAY_OBSERVATIONS_NO_DELETE_SHA256,
+ )?,
+ ])
+}
diff --git a/src/state_metadata.rs b/src/state_metadata.rs
@@ -420,7 +420,7 @@ fn normalized_config_digest(
Ok(RhiNormalizedConfigDigest(hasher.finalize().into()))
}
-fn evidence_policy_digest(
+pub(crate) fn evidence_policy_digest(
normalized: &Value,
) -> Result<RhiEvidencePolicyDigest, RhiStateMetadataError> {
let relays = normalized
diff --git a/src/state_repository.rs b/src/state_repository.rs
@@ -286,6 +286,10 @@ impl<'host> RhiStateRepositories<'host> {
Self { host }
}
+ pub(crate) const fn host(&self) -> &'host RhiStateHost {
+ self.host
+ }
+
/// Returns typed source-result access.
#[must_use]
pub const fn sources(&self) -> RhiSourceRepository<'host> {
diff --git a/src/state_trade.rs b/src/state_trade.rs
@@ -0,0 +1,569 @@
+//! Atomic immutable persistence for admitted trade-event evidence.
+
+use core::fmt;
+use std::error::Error;
+
+use radroots_service_sqlite::{
+ ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind,
+};
+use serde_json::Value;
+use sqlx::Row;
+
+use crate::{
+ RhiAdmittedTradeMutationEvent, RhiConfigDocumentV1, RhiEvidencePolicyDigest, RhiStateHostMode,
+ RhiStateRepositories, RhiTradeMutationObservedAtUnixSeconds, state_metadata,
+};
+
+/// Exact version of the immutable trade-evidence persistence contract.
+pub const RHI_TRADE_EVIDENCE_PERSISTENCE_CONTRACT_VERSION: u32 = 1;
+
+const SOURCE_SELECTOR: &str = "trade_mutation_lineage_v1";
+const MAX_SOURCE_ID_BYTES: usize = 64;
+const MAX_CONTRACT_ID_BYTES: usize = 128;
+const MAX_MUTATION_CONTENT_BYTES: usize = 131_072;
+const MAX_CANONICAL_EVENT_BYTES: usize = 524_288;
+
+const INSERT_MUTATION_SQL: &str = r#"INSERT INTO trade_mutations (
+ mutation_id, trade_id, contract_id, schema_version, event_kind,
+ author_pubkey, canonical_content
+) VALUES (?, ?, ?, ?, ?, ?, ?)
+ON CONFLICT (mutation_id) DO NOTHING"#;
+const READ_MUTATION_SQL: &str = r#"SELECT
+ length(trade_id) AS trade_id_bytes,
+ substr(trade_id, 1, 17) AS trade_id,
+ length(CAST(contract_id AS BLOB)) AS contract_id_bytes,
+ substr(contract_id, 1, 129) AS contract_id,
+ schema_version,
+ event_kind,
+ length(author_pubkey) AS author_pubkey_bytes,
+ substr(author_pubkey, 1, 33) AS author_pubkey,
+ length(canonical_content) AS canonical_content_bytes,
+ substr(canonical_content, 1, 131073) AS canonical_content
+FROM trade_mutations
+WHERE mutation_id = ?
+LIMIT 1"#;
+const INSERT_EVENT_SQL: &str = r#"INSERT INTO nostr_events (
+ event_id, event_signature, mutation_id, author_pubkey, event_kind,
+ authored_at_unix_s, canonical_event_json
+) VALUES (?, ?, ?, ?, ?, ?, ?)
+ON CONFLICT (event_id, event_signature) DO NOTHING"#;
+const READ_EVENT_SQL: &str = r#"SELECT
+ length(mutation_id) AS mutation_id_bytes,
+ substr(mutation_id, 1, 33) AS mutation_id,
+ length(author_pubkey) AS author_pubkey_bytes,
+ substr(author_pubkey, 1, 33) AS author_pubkey,
+ event_kind,
+ authored_at_unix_s,
+ length(canonical_event_json) AS canonical_event_json_bytes,
+ substr(canonical_event_json, 1, 524289) AS canonical_event_json
+FROM nostr_events
+WHERE event_id = ? AND event_signature = ?
+LIMIT 1"#;
+const INSERT_OBSERVATION_SQL: &str = r#"INSERT INTO relay_observations (
+ source_id, selector_id, evidence_policy_sha256, event_id,
+ event_signature, observed_at_unix_s
+) VALUES (?, ?, ?, ?, ?, ?)
+ON CONFLICT (
+ source_id, selector_id, evidence_policy_sha256, event_id,
+ event_signature, observed_at_unix_s
+) DO NOTHING"#;
+
+/// One accepted configured source observation bound to an admitted signed event.
+///
+/// Construction is sealed to a validated RHI configuration and an admitted
+/// event, so arbitrary source labels, selectors, policy digests, identifiers,
+/// signatures, and observation times cannot be supplied independently.
+pub struct RhiTradeSourceObservation {
+ source_id: Box<str>,
+ policy: RhiEvidencePolicyDigest,
+ event_id: [u8; 32],
+ event_signature: [u8; 64],
+ observed_at: RhiTradeMutationObservedAtUnixSeconds,
+}
+
+impl RhiTradeSourceObservation {
+ /// Binds one admitted event to an exact configured `nostr_relay` source.
+ pub fn from_config(
+ configuration: &RhiConfigDocumentV1,
+ source_id: &str,
+ event: &RhiAdmittedTradeMutationEvent,
+ ) -> Result<Self, RhiTradeEvidencePersistenceError> {
+ let source = configured_source(configuration.normalized(), source_id)
+ .ok_or_else(|| failure(RhiTradeEvidencePersistenceErrorKind::InvalidObservation))?;
+ if source.pointer("/kind").and_then(Value::as_str) != Some("nostr_relay")
+ || source.pointer("/selector").and_then(Value::as_str) != Some(SOURCE_SELECTOR)
+ {
+ return Err(failure(
+ RhiTradeEvidencePersistenceErrorKind::InvalidObservation,
+ ));
+ }
+ let policy = state_metadata::evidence_policy_digest(configuration.normalized())
+ .map_err(|_| failure(RhiTradeEvidencePersistenceErrorKind::InvalidObservation))?;
+ Ok(Self {
+ source_id: source_id.into(),
+ policy,
+ event_id: *event.event_id().as_bytes(),
+ event_signature: event.event_signature_bytes(),
+ observed_at: event.observed_at_unix_seconds(),
+ })
+ }
+}
+
+impl fmt::Debug for RhiTradeSourceObservation {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter
+ .debug_struct("RhiTradeSourceObservation")
+ .field("selector", &SOURCE_SELECTOR)
+ .field("observed_at_unix_seconds", &self.observed_at.get())
+ .field("source", &"[redacted]")
+ .field("event", &"[redacted]")
+ .finish()
+ }
+}
+
+/// Stable source-free classification for immutable evidence persistence.
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+pub enum RhiTradeEvidencePersistenceErrorKind {
+ InvalidMode,
+ InvalidObservation,
+ Encoding,
+ MutationConflict,
+ SignedEventConflict,
+ Storage,
+ CommitOutcomeUnknown,
+}
+
+impl RhiTradeEvidencePersistenceErrorKind {
+ /// Returns the stable machine-readable failure code.
+ #[must_use]
+ pub const fn code(self) -> &'static str {
+ match self {
+ Self::InvalidMode => "trade_evidence_mode_invalid",
+ Self::InvalidObservation => "trade_evidence_observation_invalid",
+ Self::Encoding => "trade_evidence_encoding_failed",
+ Self::MutationConflict => "trade_mutation_conflict",
+ Self::SignedEventConflict => "trade_signed_event_conflict",
+ Self::Storage => "trade_evidence_storage_failed",
+ Self::CommitOutcomeUnknown => "trade_evidence_commit_outcome_unknown",
+ }
+ }
+}
+
+/// Redacted source-free immutable evidence persistence failure.
+#[derive(Clone, Copy, PartialEq, Eq)]
+pub struct RhiTradeEvidencePersistenceError {
+ kind: RhiTradeEvidencePersistenceErrorKind,
+}
+
+impl RhiTradeEvidencePersistenceError {
+ /// Returns the stable failure classification.
+ #[must_use]
+ pub const fn kind(self) -> RhiTradeEvidencePersistenceErrorKind {
+ self.kind
+ }
+
+ /// Returns the stable machine-readable failure code.
+ #[must_use]
+ pub const fn code(self) -> &'static str {
+ self.kind.code()
+ }
+}
+
+impl fmt::Display for RhiTradeEvidencePersistenceError {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter.write_str(match self.kind {
+ RhiTradeEvidencePersistenceErrorKind::InvalidMode => {
+ "RHI trade evidence requires writable state"
+ }
+ RhiTradeEvidencePersistenceErrorKind::InvalidObservation => {
+ "RHI trade source observation is invalid"
+ }
+ RhiTradeEvidencePersistenceErrorKind::Encoding => {
+ "RHI signed trade event encoding failed"
+ }
+ RhiTradeEvidencePersistenceErrorKind::MutationConflict => {
+ "RHI canonical trade mutation conflicts with durable evidence"
+ }
+ RhiTradeEvidencePersistenceErrorKind::SignedEventConflict => {
+ "RHI signed trade event conflicts with durable evidence"
+ }
+ RhiTradeEvidencePersistenceErrorKind::Storage => {
+ "RHI trade evidence transaction failed"
+ }
+ RhiTradeEvidencePersistenceErrorKind::CommitOutcomeUnknown => {
+ "RHI trade evidence commit outcome is unknown"
+ }
+ })
+ }
+}
+
+impl fmt::Debug for RhiTradeEvidencePersistenceError {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter
+ .debug_struct("RhiTradeEvidencePersistenceError")
+ .field("kind", &self.kind)
+ .finish()
+ }
+}
+
+impl Error for RhiTradeEvidencePersistenceError {}
+
+/// Exact immutable facts newly inserted by one committed persistence call.
+#[derive(Clone, Copy, PartialEq, Eq)]
+pub struct RhiTradeEvidencePersistenceOutcome {
+ mutation_inserted: bool,
+ signed_event_inserted: bool,
+ observation_inserted: bool,
+}
+
+impl RhiTradeEvidencePersistenceOutcome {
+ /// Returns whether the canonical mutation fact was newly inserted.
+ #[must_use]
+ pub const fn mutation_inserted(self) -> bool {
+ self.mutation_inserted
+ }
+
+ /// Returns whether the exact signed-event fact was newly inserted.
+ #[must_use]
+ pub const fn signed_event_inserted(self) -> bool {
+ self.signed_event_inserted
+ }
+
+ /// Returns whether the configured-source observation was newly inserted.
+ #[must_use]
+ pub const fn observation_inserted(self) -> bool {
+ self.observation_inserted
+ }
+}
+
+impl fmt::Debug for RhiTradeEvidencePersistenceOutcome {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter
+ .debug_struct("RhiTradeEvidencePersistenceOutcome")
+ .field("mutation_inserted", &self.mutation_inserted)
+ .field("signed_event_inserted", &self.signed_event_inserted)
+ .field("observation_inserted", &self.observation_inserted)
+ .finish()
+ }
+}
+
+impl RhiStateRepositories<'_> {
+ /// Atomically persists one admitted mutation, signed event, and observation.
+ pub async fn persist_trade_evidence(
+ &self,
+ event: RhiAdmittedTradeMutationEvent,
+ observation: RhiTradeSourceObservation,
+ ) -> Result<RhiTradeEvidencePersistenceOutcome, RhiTradeEvidencePersistenceError> {
+ let host = self.host();
+ if host.mode() != RhiStateHostMode::ReadWriteExisting {
+ return Err(failure(RhiTradeEvidencePersistenceErrorKind::InvalidMode));
+ }
+ let record = PersistenceRecord::from_admitted(event)?;
+ if observation.policy != host.metadata().evidence_policy_digest()
+ || observation.event_id != record.event_id
+ || observation.event_signature != record.event_signature
+ {
+ return Err(failure(
+ RhiTradeEvidencePersistenceErrorKind::InvalidObservation,
+ ));
+ }
+ host.sqlite_host()
+ .transaction(move |transaction| {
+ Box::pin(async move { persist(transaction, &record, &observation).await })
+ })
+ .await
+ .map_err(map_transaction_error)
+ }
+}
+
+struct PersistenceRecord {
+ mutation_id: [u8; 32],
+ trade_id: [u8; 16],
+ contract_id: &'static str,
+ schema_version: u16,
+ event_id: [u8; 32],
+ event_signature: [u8; 64],
+ author_pubkey: [u8; 32],
+ event_kind: u32,
+ authored_at_unix_s: u64,
+ canonical_content: Box<[u8]>,
+ canonical_event_json: Box<[u8]>,
+}
+
+impl PersistenceRecord {
+ fn from_admitted(
+ admitted: RhiAdmittedTradeMutationEvent,
+ ) -> Result<Self, RhiTradeEvidencePersistenceError> {
+ let (_original, event, mutation, mutation_id, _) = admitted.into_parts();
+ let canonical_event_json = serde_json::to_vec(&event.to_nip01_wire())
+ .map_err(|_| failure(RhiTradeEvidencePersistenceErrorKind::Encoding))?;
+ if canonical_event_json.is_empty() || canonical_event_json.len() > MAX_CANONICAL_EVENT_BYTES
+ {
+ return Err(failure(RhiTradeEvidencePersistenceErrorKind::Encoding));
+ }
+ let canonical_content = event.content().as_bytes().to_vec();
+ if canonical_content.is_empty() || canonical_content.len() > MAX_MUTATION_CONTENT_BYTES {
+ return Err(failure(RhiTradeEvidencePersistenceErrorKind::Encoding));
+ }
+ Ok(Self {
+ mutation_id: *mutation_id.as_bytes(),
+ trade_id: *mutation.trade_id.as_bytes(),
+ contract_id: mutation.mutation_kind().contract_id(),
+ schema_version: mutation.schema_version,
+ event_id: *event.id().as_bytes(),
+ event_signature: *event.sig().as_bytes(),
+ author_pubkey: *event.author().as_bytes(),
+ event_kind: event.kind_u32(),
+ authored_at_unix_s: event.created_at_u64(),
+ canonical_content: canonical_content.into_boxed_slice(),
+ canonical_event_json: canonical_event_json.into_boxed_slice(),
+ })
+ }
+}
+
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+enum PersistenceOperationError {
+ MutationConflict,
+ SignedEventConflict,
+ Storage,
+}
+
+async fn persist(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ record: &PersistenceRecord,
+ observation: &RhiTradeSourceObservation,
+) -> Result<RhiTradeEvidencePersistenceOutcome, PersistenceOperationError> {
+ let mutation_inserted = insert_mutation(transaction, record).await?;
+ let signed_event_inserted = insert_event(transaction, record).await?;
+ let observation_inserted = insert_observation(transaction, observation).await?;
+ Ok(RhiTradeEvidencePersistenceOutcome {
+ mutation_inserted,
+ signed_event_inserted,
+ observation_inserted,
+ })
+}
+
+async fn insert_mutation(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ record: &PersistenceRecord,
+) -> Result<bool, PersistenceOperationError> {
+ let result = sqlx::query(INSERT_MUTATION_SQL)
+ .bind(record.mutation_id.as_slice())
+ .bind(record.trade_id.as_slice())
+ .bind(record.contract_id)
+ .bind(i64::from(record.schema_version))
+ .bind(i64::from(record.event_kind))
+ .bind(record.author_pubkey.as_slice())
+ .bind(record.canonical_content.as_ref())
+ .execute(&mut *transaction)
+ .await
+ .map_err(|_| PersistenceOperationError::Storage)?;
+ match result.rows_affected() {
+ 1 => Ok(true),
+ 0 if mutation_matches(transaction, record).await? => Ok(false),
+ 0 => Err(PersistenceOperationError::MutationConflict),
+ _ => Err(PersistenceOperationError::Storage),
+ }
+}
+
+async fn mutation_matches(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ expected: &PersistenceRecord,
+) -> Result<bool, PersistenceOperationError> {
+ let row = sqlx::query(READ_MUTATION_SQL)
+ .bind(expected.mutation_id.as_slice())
+ .fetch_optional(&mut *transaction)
+ .await
+ .map_err(|_| PersistenceOperationError::Storage)?
+ .ok_or(PersistenceOperationError::Storage)?;
+ Ok(
+ exact_blob(&row, "trade_id", "trade_id_bytes")? == expected.trade_id
+ && bounded_text(
+ &row,
+ "contract_id",
+ "contract_id_bytes",
+ MAX_CONTRACT_ID_BYTES,
+ )? == expected.contract_id
+ && row.try_get::<i64, _>("schema_version").ok()
+ == Some(i64::from(expected.schema_version))
+ && row.try_get::<i64, _>("event_kind").ok() == Some(i64::from(expected.event_kind))
+ && exact_blob(&row, "author_pubkey", "author_pubkey_bytes")? == expected.author_pubkey
+ && bounded_blob(
+ &row,
+ "canonical_content",
+ "canonical_content_bytes",
+ MAX_MUTATION_CONTENT_BYTES,
+ )? == expected.canonical_content.as_ref(),
+ )
+}
+
+async fn insert_event(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ record: &PersistenceRecord,
+) -> Result<bool, PersistenceOperationError> {
+ let result = sqlx::query(INSERT_EVENT_SQL)
+ .bind(record.event_id.as_slice())
+ .bind(record.event_signature.as_slice())
+ .bind(record.mutation_id.as_slice())
+ .bind(record.author_pubkey.as_slice())
+ .bind(i64::from(record.event_kind))
+ .bind(
+ i64::try_from(record.authored_at_unix_s)
+ .map_err(|_| PersistenceOperationError::Storage)?,
+ )
+ .bind(record.canonical_event_json.as_ref())
+ .execute(&mut *transaction)
+ .await
+ .map_err(|_| PersistenceOperationError::Storage)?;
+ match result.rows_affected() {
+ 1 => Ok(true),
+ 0 if event_matches(transaction, record).await? => Ok(false),
+ 0 => Err(PersistenceOperationError::SignedEventConflict),
+ _ => Err(PersistenceOperationError::Storage),
+ }
+}
+
+async fn event_matches(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ expected: &PersistenceRecord,
+) -> Result<bool, PersistenceOperationError> {
+ let row = sqlx::query(READ_EVENT_SQL)
+ .bind(expected.event_id.as_slice())
+ .bind(expected.event_signature.as_slice())
+ .fetch_optional(&mut *transaction)
+ .await
+ .map_err(|_| PersistenceOperationError::Storage)?
+ .ok_or(PersistenceOperationError::Storage)?;
+ Ok(
+ exact_blob(&row, "mutation_id", "mutation_id_bytes")? == expected.mutation_id
+ && exact_blob(&row, "author_pubkey", "author_pubkey_bytes")? == expected.author_pubkey
+ && row.try_get::<i64, _>("event_kind").ok() == Some(i64::from(expected.event_kind))
+ && row.try_get::<i64, _>("authored_at_unix_s").ok()
+ == i64::try_from(expected.authored_at_unix_s).ok()
+ && bounded_blob(
+ &row,
+ "canonical_event_json",
+ "canonical_event_json_bytes",
+ MAX_CANONICAL_EVENT_BYTES,
+ )? == expected.canonical_event_json.as_ref(),
+ )
+}
+
+async fn insert_observation(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ observation: &RhiTradeSourceObservation,
+) -> Result<bool, PersistenceOperationError> {
+ let result = sqlx::query(INSERT_OBSERVATION_SQL)
+ .bind(observation.source_id.as_ref())
+ .bind(SOURCE_SELECTOR)
+ .bind(observation.policy.as_bytes().as_slice())
+ .bind(observation.event_id.as_slice())
+ .bind(observation.event_signature.as_slice())
+ .bind(
+ i64::try_from(observation.observed_at.get())
+ .map_err(|_| PersistenceOperationError::Storage)?,
+ )
+ .execute(&mut *transaction)
+ .await
+ .map_err(|_| PersistenceOperationError::Storage)?;
+ match result.rows_affected() {
+ 0 => Ok(false),
+ 1 => Ok(true),
+ _ => Err(PersistenceOperationError::Storage),
+ }
+}
+
+fn configured_source<'a>(configuration: &'a Value, source_id: &str) -> Option<&'a Value> {
+ if source_id.is_empty()
+ || source_id.len() > MAX_SOURCE_ID_BYTES
+ || !source_id.bytes().enumerate().all(|(index, byte)| {
+ if index == 0 {
+ byte.is_ascii_lowercase()
+ } else {
+ byte.is_ascii_lowercase() || byte.is_ascii_digit() || matches!(byte, b'_' | b'-')
+ }
+ })
+ {
+ return None;
+ }
+ configuration
+ .pointer("/evidence/sources")?
+ .as_array()?
+ .iter()
+ .find(|source| source.pointer("/source_id").and_then(Value::as_str) == Some(source_id))
+}
+
+fn exact_blob<const N: usize>(
+ row: &sqlx::sqlite::SqliteRow,
+ field: &str,
+ length_field: &str,
+) -> Result<[u8; N], PersistenceOperationError> {
+ bounded_blob(row, field, length_field, N)?
+ .try_into()
+ .map_err(|_| PersistenceOperationError::Storage)
+}
+
+fn bounded_text(
+ row: &sqlx::sqlite::SqliteRow,
+ field: &str,
+ length_field: &str,
+ maximum: usize,
+) -> Result<String, PersistenceOperationError> {
+ let length = bounded_length(row, length_field, maximum)?;
+ let value = row
+ .try_get::<String, _>(field)
+ .map_err(|_| PersistenceOperationError::Storage)?;
+ (value.len() == length)
+ .then_some(value)
+ .ok_or(PersistenceOperationError::Storage)
+}
+
+fn bounded_blob(
+ row: &sqlx::sqlite::SqliteRow,
+ field: &str,
+ length_field: &str,
+ maximum: usize,
+) -> Result<Vec<u8>, PersistenceOperationError> {
+ let length = bounded_length(row, length_field, maximum)?;
+ let value = row
+ .try_get::<Vec<u8>, _>(field)
+ .map_err(|_| PersistenceOperationError::Storage)?;
+ (value.len() == length)
+ .then_some(value)
+ .ok_or(PersistenceOperationError::Storage)
+}
+
+fn bounded_length(
+ row: &sqlx::sqlite::SqliteRow,
+ field: &str,
+ maximum: usize,
+) -> Result<usize, PersistenceOperationError> {
+ row.try_get::<i64, _>(field)
+ .ok()
+ .and_then(|value| usize::try_from(value).ok())
+ .filter(|value| (1..=maximum).contains(value))
+ .ok_or(PersistenceOperationError::Storage)
+}
+
+fn map_transaction_error(
+ error: ServiceSqliteTransactionError<PersistenceOperationError>,
+) -> RhiTradeEvidencePersistenceError {
+ if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown {
+ return failure(RhiTradeEvidencePersistenceErrorKind::CommitOutcomeUnknown);
+ }
+ failure(match error.operation_error().copied() {
+ Some(PersistenceOperationError::MutationConflict) => {
+ RhiTradeEvidencePersistenceErrorKind::MutationConflict
+ }
+ Some(PersistenceOperationError::SignedEventConflict) => {
+ RhiTradeEvidencePersistenceErrorKind::SignedEventConflict
+ }
+ Some(PersistenceOperationError::Storage) | None => {
+ RhiTradeEvidencePersistenceErrorKind::Storage
+ }
+ })
+}
+
+const fn failure(kind: RhiTradeEvidencePersistenceErrorKind) -> RhiTradeEvidencePersistenceError {
+ RhiTradeEvidencePersistenceError { kind }
+}
diff --git a/src/trade_ingest.rs b/src/trade_ingest.rs
@@ -296,6 +296,7 @@ pub struct RhiAdmittedTradeMutationEvent {
event: EventEnvelope,
mutation: TradeMutationEnvelopeV1,
mutation_id: MutationId,
+ observed_at: RhiTradeMutationObservedAtUnixSeconds,
}
impl RhiAdmittedTradeMutationEvent {
@@ -329,11 +330,39 @@ impl RhiAdmittedTradeMutationEvent {
self.event.created_at_u64()
}
+ /// Returns the injected UTC second used to admit this source observation.
+ #[must_use]
+ pub const fn observed_at_unix_seconds(&self) -> RhiTradeMutationObservedAtUnixSeconds {
+ self.observed_at
+ }
+
/// Returns the canonical typed mutation bound to the signed event.
#[must_use]
pub const fn mutation(&self) -> &TradeMutationEnvelopeV1 {
&self.mutation
}
+
+ pub(crate) fn event_signature_bytes(&self) -> [u8; 64] {
+ *self.event.sig().as_bytes()
+ }
+
+ pub(crate) fn into_parts(
+ self,
+ ) -> (
+ Box<[u8]>,
+ EventEnvelope,
+ TradeMutationEnvelopeV1,
+ MutationId,
+ RhiTradeMutationObservedAtUnixSeconds,
+ ) {
+ (
+ self.original,
+ self.event,
+ self.mutation,
+ self.mutation_id,
+ self.observed_at,
+ )
+ }
}
impl fmt::Debug for RhiAdmittedTradeMutationEvent {
@@ -410,6 +439,7 @@ pub fn admit_rhi_trade_mutation_event(
event,
mutation,
mutation_id,
+ observed_at,
})
}
diff --git a/tests/build_policy.rs b/tests/build_policy.rs
@@ -27,7 +27,7 @@ fn source_lock_metadata_is_exact_and_nix_is_absent() {
));
for field in [
"config_contract_version = 1",
- "state_contract_version = 2",
+ "state_contract_version = 3",
"admin_contract_version = 1",
"status_contract_version = 1",
"provider_contract_version = 1",
@@ -93,7 +93,7 @@ fn source_lock_binds_the_current_cargo_lock() {
assert!(!SOURCE_LOCK.contains("flake_lock_sha256"));
assert!(!SOURCE_LOCK.contains("lib_revision ="));
assert!(SOURCE_LOCK.ends_with(
- "[contract_versions]\nconfig = 1\nstate = 2\nadmin = 1\nstatus = 1\nprovider = 1\n"
+ "[contract_versions]\nconfig = 1\nstate = 3\nadmin = 1\nstatus = 1\nprovider = 1\n"
));
}
diff --git a/tests/package_boundary.rs b/tests/package_boundary.rs
@@ -15,6 +15,8 @@ const RUNTIME_FOUNDATION_CONTRACT: &str =
include_str!("../contracts/services_hardening/runtime_foundation.v1.json");
const TRADE_INGEST_CONTRACT: &str =
include_str!("../contracts/services_hardening/trade_ingest.v1.json");
+const TRADE_EVIDENCE_PERSISTENCE_CONTRACT: &str =
+ include_str!("../contracts/services_hardening/trade_evidence_persistence.v1.json");
const PUBLIC_API: &str = include_str!("../contracts/api_baselines/rhi.txt");
const SOURCES: &[&str] = &[
include_str!("../src/adapters/nostr/event.rs"),
@@ -32,6 +34,7 @@ const SOURCES: &[&str] = &[
include_str!("../src/state_maintenance.rs"),
include_str!("../src/state_metadata.rs"),
include_str!("../src/state_repository.rs"),
+ include_str!("../src/state_trade.rs"),
include_str!("../src/trade_ingest.rs"),
];
@@ -96,6 +99,7 @@ fn state_catalog_module_is_private_and_root_api_is_curated() {
"state_maintenance",
"state_metadata",
"state_repository",
+ "state_trade",
"trade_ingest",
] {
assert!(
@@ -138,6 +142,11 @@ fn state_catalog_module_is_private_and_root_api_is_curated() {
"admit_rhi_trade_mutation_event",
"RhiTradeMutationAdmissionLimits",
"RhiAdmittedTradeMutationEvent",
+ "RhiTradeEvidencePersistenceError",
+ "RhiTradeEvidencePersistenceErrorKind",
+ "RhiTradeEvidencePersistenceOutcome",
+ "RhiTradeSourceObservation",
+ "RHI_TRADE_EVIDENCE_PERSISTENCE_CONTRACT_VERSION",
] {
assert!(
ROOT.contains(required),
@@ -189,7 +198,7 @@ fn public_errors_are_crate_owned_redacted_and_source_free() {
.lines()
.filter(|line| line.starts_with("pub struct rhi::") && line.ends_with("Error"))
.count();
- assert_eq!(public_error_count, 14);
+ assert_eq!(public_error_count, 15);
}
#[test]
@@ -209,6 +218,10 @@ fn trade_ingest_is_sealed_bounded_verified_and_effect_free() {
assert_eq!(contract["effects"]["filesystem"], false);
assert_eq!(contract["effects"]["sqlite"], false);
assert_eq!(contract["effects"]["network"], false);
+ assert_eq!(
+ contract["persistence_contract"],
+ "contracts/services_hardening/trade_evidence_persistence.v1.json"
+ );
let source = include_str!("../src/trade_ingest.rs");
for required in [
@@ -241,6 +254,56 @@ fn trade_ingest_is_sealed_bounded_verified_and_effect_free() {
}
#[test]
+fn trade_evidence_persistence_is_typed_atomic_and_sealed() {
+ let contract: serde_json::Value = serde_json::from_str(TRADE_EVIDENCE_PERSISTENCE_CONTRACT)
+ .expect("trade-evidence persistence contract");
+ assert_eq!(
+ contract["schema"],
+ "radroots.rhi.trade-evidence-persistence.v1"
+ );
+ assert_eq!(contract["contract_version"], 1);
+ assert_eq!(contract["transaction"], "one_governed_sqlx_transaction");
+ assert_eq!(contract["effects"]["checkpoint"], false);
+ assert_eq!(contract["effects"]["dirty_generation"], false);
+ assert_eq!(
+ contract["facts"]["signed_event"]["identity"],
+ serde_json::json!(["verified_event_id", "verified_event_signature"])
+ );
+
+ let source = include_str!("../src/state_trade.rs");
+ for required in [
+ "pub async fn persist_trade_evidence",
+ "RhiAdmittedTradeMutationEvent",
+ "RhiTradeSourceObservation",
+ ".transaction(move |transaction|",
+ "ON CONFLICT (event_id, event_signature) DO NOTHING",
+ ] {
+ assert!(
+ source.contains(required),
+ "trade evidence persistence is missing {required}"
+ );
+ }
+ for forbidden in [
+ "pub fn host(",
+ "pub fn into_parts(",
+ "pub fn event_signature_bytes(",
+ "SystemTime",
+ "std::fs",
+ "std::net",
+ "checkpoint",
+ "dirty_generation",
+ ] {
+ assert!(
+ !source.contains(forbidden),
+ "trade evidence persistence gained forbidden authority {forbidden}"
+ );
+ }
+ assert!(!ROOT.contains("pub mod state_trade"));
+ assert!(!PUBLIC_API.contains("rhi::state_trade::"));
+ assert!(!PUBLIC_API.contains("sqlx::"));
+}
+
+#[test]
fn runtime_foundation_is_existing_only_passive_and_process_neutral() {
let contract: serde_json::Value =
serde_json::from_str(RUNTIME_FOUNDATION_CONTRACT).expect("runtime foundation contract");
@@ -353,6 +416,10 @@ fn readme_freezes_the_root_only_boundary_and_exact_baseline() {
"```compile_fail",
"[RHI API baseline](contracts/api_baselines/rhi.txt)",
"[`trade_ingest.v1.json`](contracts/services_hardening/trade_ingest.v1.json)",
+ "[`trade_evidence_persistence.v1.json`](contracts/services_hardening/trade_evidence_persistence.v1.json)",
+ "one canonical mutation",
+ "every distinct valid signed",
+ "does not advance reconciliation checkpoints or dirty generation",
"## Injected runtime adapters",
"whole-second wall UTC",
"process-local monotonic time",
@@ -379,6 +446,8 @@ fn readme_freezes_the_root_only_boundary_and_exact_baseline() {
"Compose those dependencies only through the sealed runtime-adapter boundary",
"exposes no task handle or concrete transport handle",
"admit_rhi_trade_mutation_event",
+ "Persist each canonical",
+ "independently signed Nostr event",
] {
assert!(AGENTS.contains(required), "AGENTS is missing {required}");
}
diff --git a/tests/services_hardening_config_lifecycle.rs b/tests/services_hardening_config_lifecycle.rs
@@ -163,7 +163,7 @@ async fn existing_intent_and_offline_apply_bind_exact_append_only_evidence() {
for row in &rows {
assert_eq!(row.try_get::<i64, _>("config_bytes").unwrap(), 32);
assert_eq!(row.try_get::<i64, _>("policy_bytes").unwrap(), 32);
- assert_eq!(row.try_get::<i64, _>("state_contract_version").unwrap(), 2);
+ assert_eq!(row.try_get::<i64, _>("state_contract_version").unwrap(), 3);
assert_eq!(
row.try_get::<String, _>("service_public_key")
.unwrap()
diff --git a/tests/services_hardening_state_catalog.rs b/tests/services_hardening_state_catalog.rs
@@ -11,8 +11,9 @@ use rhi::{
RHI_STATE_SCHEMA_VERSION, RHI_STATE_SCHEMA_VERSION_1_OBJECT_COUNT,
RHI_STATE_SCHEMA_VERSION_1_SHA256, RHI_STATE_SCHEMA_VERSION_2_MIGRATION_SHA256,
RHI_STATE_SCHEMA_VERSION_2_OBJECT_COUNT, RHI_STATE_SCHEMA_VERSION_2_SHA256,
- RhiStateCatalogErrorKind, rhi_migration_catalog, rhi_schema_catalog,
- validate_rhi_state_catalogs,
+ RHI_STATE_SCHEMA_VERSION_3_MIGRATION_SHA256, RHI_STATE_SCHEMA_VERSION_3_OBJECT_COUNT,
+ RHI_STATE_SCHEMA_VERSION_3_SHA256, RhiStateCatalogErrorKind, rhi_migration_catalog,
+ rhi_schema_catalog, validate_rhi_state_catalogs,
};
const CATALOG_SOURCE: &str = include_str!("../src/state_catalog.rs");
@@ -20,14 +21,14 @@ const LIB_SOURCE: &str = include_str!("../src/lib.rs");
const MANIFEST: &str = include_str!("../Cargo.toml");
#[test]
-fn schema_v1_to_v2_catalogs_have_exact_literal_identities() {
+fn schema_v1_through_v3_catalogs_have_exact_literal_identities() {
let migrations = rhi_migration_catalog().expect("RHI migration catalog");
let schema = rhi_schema_catalog().expect("RHI schema catalog");
assert_eq!(RHI_STATE_BASE_SCHEMA_VERSION, 1);
- assert_eq!(RHI_STATE_SCHEMA_VERSION, 2);
- assert_eq!(migrations.descriptors().len(), 1);
- assert_eq!(migrations.current_version(), 2);
+ assert_eq!(RHI_STATE_SCHEMA_VERSION, 3);
+ assert_eq!(migrations.descriptors().len(), 2);
+ assert_eq!(migrations.current_version(), 3);
assert_eq!(migrations.descriptors()[0].target_version(), 2);
assert_eq!(
migrations.descriptors()[0].name().as_str(),
@@ -37,12 +38,21 @@ fn schema_v1_to_v2_catalogs_have_exact_literal_identities() {
migrations.descriptors()[0].checksum().as_bytes(),
&RHI_STATE_SCHEMA_VERSION_2_MIGRATION_SHA256
);
+ assert_eq!(migrations.descriptors()[1].target_version(), 3);
+ assert_eq!(
+ migrations.descriptors()[1].name().as_str(),
+ "create_immutable_trade_evidence"
+ );
+ assert_eq!(
+ migrations.descriptors()[1].checksum().as_bytes(),
+ &RHI_STATE_SCHEMA_VERSION_3_MIGRATION_SHA256
+ );
assert_eq!(
migrations.digest().as_bytes(),
&RHI_MIGRATION_CATALOG_SHA256
);
- assert_eq!(schema.versions().len(), 2);
+ assert_eq!(schema.versions().len(), 3);
let version = schema.versions()[0];
assert_eq!(version.version(), 1);
assert_eq!(
@@ -65,13 +75,24 @@ fn schema_v1_to_v2_catalogs_have_exact_literal_identities() {
version.digest().as_bytes(),
&RHI_STATE_SCHEMA_VERSION_2_SHA256
);
+ let version = schema.versions()[2];
+ assert_eq!(version.version(), 3);
+ assert_eq!(
+ version.object_count(),
+ RHI_STATE_SCHEMA_VERSION_3_OBJECT_COUNT
+ );
+ assert_eq!(version.object_count(), 22);
+ assert_eq!(
+ version.digest().as_bytes(),
+ &RHI_STATE_SCHEMA_VERSION_3_SHA256
+ );
assert_eq!(schema.digest().as_bytes(), &RHI_STATE_SCHEMA_CATALOG_SHA256);
assert_eq!(schema.migration_catalog_digest(), migrations.digest());
validate_rhi_state_catalogs(&migrations, &schema).expect("exact catalogs");
assert_eq!(
lower_hex(&RHI_MIGRATION_CATALOG_SHA256),
- "b640a9095d5318dbfd0afbfbe6e05281113a9cef4564805c021c368c3c52fc08"
+ "14046048b468836f2602ec51e538f2a98b71c945f83d933dccb8603b84f865f0"
);
assert_eq!(
lower_hex(&RHI_STATE_SCHEMA_VERSION_1_SHA256),
@@ -86,8 +107,16 @@ fn schema_v1_to_v2_catalogs_have_exact_literal_identities() {
"bccbf1e62fe7644c9c377205b25a92298e088b8c26d15ba45133ca5e9b7315a9"
);
assert_eq!(
+ lower_hex(&RHI_STATE_SCHEMA_VERSION_3_MIGRATION_SHA256),
+ "07b098c393140afec222fce76ec668698d50f1d23785856873e30445af1a013f"
+ );
+ assert_eq!(
+ lower_hex(&RHI_STATE_SCHEMA_VERSION_3_SHA256),
+ "fd96226405ab68655aee00f683f6023c7aab2dbd23feadac165333490d6f0ad5"
+ );
+ assert_eq!(
lower_hex(&RHI_STATE_SCHEMA_CATALOG_SHA256),
- "1b7f7359b62ed7dcd476289e46693d9b3d1369dcaf7d5952f0329a5c86f8b85f"
+ "1325a41b90abfc7d523bbfe83500a1b1a73e492cc930d6ead6c71e29d5447436"
);
}
@@ -127,7 +156,11 @@ fn independent_validator_rejects_migration_or_schema_drift() {
let snapshot_digest =
SchemaVersionCatalog::computed_digest(2, [object.clone()]).expect("snapshot digest");
let version_two = SchemaVersionCatalog::new(2, [object], snapshot_digest).expect("version two");
- let schema = SchemaCatalog::new(&exact_migrations, [version_one, version_two])
+ let snapshot_digest =
+ SchemaVersionCatalog::computed_digest(3, [version_two_object()]).expect("v3 digest");
+ let version_three = SchemaVersionCatalog::new(3, [version_two_object()], snapshot_digest)
+ .expect("version three");
+ let schema = SchemaCatalog::new(&exact_migrations, [version_one, version_two, version_three])
.expect("drift schema catalog");
assert_eq!(
validate_rhi_state_catalogs(&exact_migrations, &schema)
@@ -164,8 +197,12 @@ fn catalog_errors_are_stable_source_free_and_redacted() {
let snapshot =
SchemaVersionCatalog::computed_digest(2, [object.clone()]).expect("snapshot digest");
let version_two = SchemaVersionCatalog::new(2, [object], snapshot).expect("version two");
- let schema =
- SchemaCatalog::new(&migrations, [version_one, version_two]).expect("schema catalog");
+ let version_three_digest =
+ SchemaVersionCatalog::computed_digest(3, [secret_object()]).expect("v3 digest");
+ let version_three = SchemaVersionCatalog::new(3, [secret_object()], version_three_digest)
+ .expect("version three");
+ let schema = SchemaCatalog::new(&migrations, [version_one, version_two, version_three])
+ .expect("schema catalog");
let error = validate_rhi_state_catalogs(&migrations, &schema).expect_err("mismatch");
assert_eq!(error.kind(), RhiStateCatalogErrorKind::CatalogMismatch);
@@ -176,6 +213,39 @@ fn catalog_errors_are_stable_source_free_and_redacted() {
assert!(!rendered.contains(&lower_hex(snapshot.as_bytes())));
}
+fn version_two_object() -> SchemaObject {
+ const SQL: &str = "CREATE TABLE unexpected (value INTEGER NOT NULL) STRICT";
+ let digest =
+ SchemaObject::computed_digest(SchemaObjectKind::Table, "unexpected", "unexpected", SQL)
+ .expect("object digest");
+ SchemaObject::new(
+ SchemaObjectKind::Table,
+ "unexpected",
+ "unexpected",
+ SQL,
+ digest,
+ )
+ .expect("schema object")
+}
+
+fn secret_object() -> SchemaObject {
+ let digest = SchemaObject::computed_digest(
+ SchemaObjectKind::Table,
+ "secret_table",
+ "secret_table",
+ "secret SQL text",
+ )
+ .expect("object digest");
+ SchemaObject::new(
+ SchemaObjectKind::Table,
+ "secret_table",
+ "secret_table",
+ "secret SQL text",
+ digest,
+ )
+ .expect("object")
+}
+
#[test]
fn catalog_source_is_pure_pinned_and_uses_only_the_shared_authority() {
assert!(MANIFEST.contains(
diff --git a/tests/services_hardening_state_resilience.rs b/tests/services_hardening_state_resilience.rs
@@ -347,7 +347,7 @@ async fn exact_open_rejects_unexpected_migration_history_without_repair() {
service_version, service_commit, lib_revision, rust_version, target,
feature_profile, config_contract_version, state_contract_version,
admin_contract_version, status_contract_version, provider_contract_version
- ) VALUES (3, 'unexpected_schema', ?, 1725000000, '0.1.0', ?, ?,
+ ) VALUES (4, 'unexpected_schema', ?, 1725000000, '0.1.0', ?, ?,
'rustc-test', 'test-target', 'service-host', 1, 3, 1, 1, 1)",
)
.bind([0x44_u8; 32].as_slice())
diff --git a/tests/services_hardening_trade_persistence.rs b/tests/services_hardening_trade_persistence.rs
@@ -0,0 +1,472 @@
+#![forbid(unsafe_code)]
+#![cfg(any(target_os = "linux", target_os = "macos"))]
+
+use std::{error::Error, fs, os::unix::fs::PermissionsExt, path::Path};
+
+use nostr::secp256k1::{Keypair, Message};
+use nostr::{EventId, Keys, SECP256K1};
+use radroots_service_sqlite::{MigrationAppliedAtUnixSeconds, MigrationBuildIdentity};
+use radroots_storage::event::SourceGeneration;
+use rhi::{
+ RHI_TRADE_EVIDENCE_PERSISTENCE_CONTRACT_VERSION, RadrootsHostEnvironment, RadrootsPathResolver,
+ RadrootsPlatform, RhiConfigProfile, RhiStateMetadata, RhiTradeEvidencePersistenceErrorKind,
+ RhiTradeMutationAdmissionLimits, RhiTradeMutationAuthoredTimePolicy,
+ RhiTradeMutationObservedAtUnixSeconds, RhiTradeSourceObservation,
+ admit_rhi_trade_mutation_event, initialize_rhi_state, open_rhi_state_inspection,
+ open_rhi_state_read_write, parse_rhi_cli_v1_from, parse_rhi_config_v1,
+ resolve_rhi_runtime_context,
+};
+use serde_json::Value;
+use sqlx::{ConnectOptions, Connection, SqliteConnection, sqlite::SqliteConnectOptions};
+
+const CONFIG: &str = include_str!("../contracts/services_hardening/config.v1.example.toml");
+const CONTRACT: &str =
+ include_str!("../contracts/services_hardening/trade_evidence_persistence.v1.json");
+const VECTOR: &str = include_str!("../contracts/conformance/vectors/trade_ingest_proposal.v1.json");
+
+fn runtime(root: &Path) -> rhi::RhiRuntimeContext {
+ let invocation = parse_rhi_cli_v1_from([
+ "rhi",
+ "--profile",
+ "repo-local",
+ "--instance",
+ "primary",
+ "--repo-local-root",
+ root.to_str().expect("UTF-8 root"),
+ "run",
+ ])
+ .expect("invocation");
+ resolve_rhi_runtime_context(
+ &RadrootsPathResolver::new(RadrootsPlatform::Linux, RadrootsHostEnvironment::default()),
+ &invocation,
+ )
+ .expect("runtime")
+}
+
+fn configuration() -> rhi::RhiConfigDocumentV1 {
+ parse_rhi_config_v1(CONFIG.as_bytes(), RhiConfigProfile::RepoLocal).expect("configuration")
+}
+
+fn metadata(
+ runtime: &rhi::RhiRuntimeContext,
+ configuration: &rhi::RhiConfigDocumentV1,
+) -> RhiStateMetadata {
+ RhiStateMetadata::new(
+ runtime,
+ configuration,
+ SourceGeneration::new([0x5a; 32]).expect("generation"),
+ 1_725_000_000_000,
+ )
+ .expect("metadata")
+}
+
+fn migration_evidence() -> (MigrationAppliedAtUnixSeconds, MigrationBuildIdentity) {
+ (
+ MigrationAppliedAtUnixSeconds::new(1_725_000_000).expect("migration time"),
+ MigrationBuildIdentity::new(
+ env!("CARGO_PKG_VERSION"),
+ "1111111111111111111111111111111111111111",
+ "7d7b454b4c9ed86569671993bd03ca868b676665",
+ "rustc-test",
+ "test-target",
+ "service-host",
+ 1,
+ rhi::RHI_STATE_SCHEMA_VERSION,
+ 1,
+ 1,
+ 1,
+ )
+ .expect("build identity"),
+ )
+}
+
+fn prepare_state_directory(runtime: &rhi::RhiRuntimeContext) {
+ fs::create_dir_all(runtime.context().paths().state()).expect("state directory");
+ fs::set_permissions(
+ runtime.context().paths().state(),
+ fs::Permissions::from_mode(0o700),
+ )
+ .expect("state mode");
+}
+
+fn signed_wire(auxiliary: u8) -> Vec<u8> {
+ let vector: Value = serde_json::from_str(VECTOR).expect("vector");
+ let mut event: Value =
+ serde_json::from_str(vector["raw_json"].as_str().expect("raw event")).expect("event JSON");
+ let keys = Keys::parse("10c5304d6c9ae3a1a16f7860f1cc8f5e3a76225a2663b3a989a0d775919b7df5")
+ .expect("approved fixture keys");
+ let event_id = EventId::from_hex(event["id"].as_str().expect("event id")).expect("event id");
+ let message = Message::from_digest(event_id.to_bytes());
+ let keypair = Keypair::from_secret_key(SECP256K1, keys.secret_key());
+ let signature = SECP256K1.sign_schnorr_with_aux_rand(&message, &keypair, &[auxiliary; 32]);
+ event["sig"] = signature.to_string().into();
+ serde_json::to_vec(&event).expect("event JSON")
+}
+
+fn admitted(
+ configuration: &rhi::RhiConfigDocumentV1,
+ wire: &[u8],
+ observed_at: u64,
+) -> rhi::RhiAdmittedTradeMutationEvent {
+ admit_rhi_trade_mutation_event(
+ RhiTradeMutationAdmissionLimits::from_config(configuration).expect("limits"),
+ wire,
+ RhiTradeMutationObservedAtUnixSeconds::new(observed_at).expect("observation time"),
+ RhiTradeMutationAuthoredTimePolicy::new(0).expect("time policy"),
+ )
+ .expect("admitted event")
+}
+
+async fn offline_connection(runtime: &rhi::RhiRuntimeContext) -> SqliteConnection {
+ SqliteConnection::connect_with(
+ &SqliteConnectOptions::new()
+ .filename(runtime.artifacts().state_database())
+ .create_if_missing(false)
+ .disable_statement_logging(),
+ )
+ .await
+ .expect("offline connection")
+}
+
+async fn insert_mutation(
+ connection: &mut SqliteConnection,
+ event: &rhi::RhiAdmittedTradeMutationEvent,
+ canonical_content: &[u8],
+) {
+ let mutation = event.mutation();
+ sqlx::query(
+ r#"INSERT INTO trade_mutations (
+ mutation_id, trade_id, contract_id, schema_version, event_kind,
+ author_pubkey, canonical_content
+ ) VALUES (?, ?, ?, ?, ?, ?, ?)"#,
+ )
+ .bind(event.mutation_id().as_bytes().as_slice())
+ .bind(mutation.trade_id.as_bytes().as_slice())
+ .bind(mutation.mutation_kind().contract_id())
+ .bind(i64::from(mutation.schema_version))
+ .bind(i64::from(event.event_kind()))
+ .bind(mutation.author_pubkey.as_bytes().as_slice())
+ .bind(canonical_content)
+ .execute(connection)
+ .await
+ .expect("seed mutation");
+}
+
+fn decode_hex<const N: usize>(value: &str) -> [u8; N] {
+ assert_eq!(value.len(), N * 2);
+ let mut bytes = [0_u8; N];
+ for (index, byte) in bytes.iter_mut().enumerate() {
+ *byte = u8::from_str_radix(&value[index * 2..index * 2 + 2], 16).expect("hex byte");
+ }
+ bytes
+}
+
+#[test]
+fn machine_contract_freezes_three_separate_immutable_facts() {
+ let contract: Value = serde_json::from_str(CONTRACT).expect("contract");
+ assert_eq!(
+ contract["schema"],
+ "radroots.rhi.trade-evidence-persistence.v1"
+ );
+ assert_eq!(contract["contract_version"], 1);
+ assert_eq!(RHI_TRADE_EVIDENCE_PERSISTENCE_CONTRACT_VERSION, 1);
+ assert_eq!(
+ contract["facts"]["canonical_mutation"]["table"],
+ "trade_mutations"
+ );
+ assert_eq!(
+ contract["facts"]["canonical_mutation"]["event_authored_time_column"],
+ "absent_event_fact_only"
+ );
+ assert_eq!(contract["facts"]["signed_event"]["table"], "nostr_events");
+ assert_eq!(
+ contract["facts"]["signed_event"]["identity"],
+ serde_json::json!(["verified_event_id", "verified_event_signature"])
+ );
+ assert_eq!(
+ contract["facts"]["source_observation"]["table"],
+ "relay_observations"
+ );
+ assert_eq!(contract["effects"]["checkpoint"], false);
+ assert_eq!(contract["effects"]["dirty_generation"], false);
+}
+
+#[tokio::test]
+async fn signed_events_mutations_and_observations_are_atomic_distinct_and_idempotent() {
+ let root = tempfile::tempdir().expect("root");
+ let runtime = runtime(root.path());
+ prepare_state_directory(&runtime);
+ let configuration = configuration();
+ let metadata = metadata(&runtime, &configuration);
+ let (applied_at, build) = migration_evidence();
+ initialize_rhi_state(&runtime, &metadata, applied_at, &build)
+ .await
+ .expect("initialize");
+ let host = open_rhi_state_read_write(&runtime, &metadata, applied_at, &build)
+ .await
+ .expect("writer");
+
+ let first_wire = signed_wire(1);
+ let second_wire = signed_wire(2);
+ let first_json: Value = serde_json::from_slice(&first_wire).expect("first JSON");
+ let second_json: Value = serde_json::from_slice(&second_wire).expect("second JSON");
+ assert_eq!(first_json["id"], second_json["id"]);
+ assert_ne!(first_json["sig"], second_json["sig"]);
+
+ let first = admitted(&configuration, &first_wire, 1_784_347_200);
+ let first_observation =
+ RhiTradeSourceObservation::from_config(&configuration, "trade-primary", &first)
+ .expect("observation");
+ let outcome = host
+ .repositories()
+ .persist_trade_evidence(first, first_observation)
+ .await
+ .expect("first persistence");
+ assert!(outcome.mutation_inserted());
+ assert!(outcome.signed_event_inserted());
+ assert!(outcome.observation_inserted());
+
+ let second = admitted(&configuration, &second_wire, 1_784_347_201);
+ let second_observation =
+ RhiTradeSourceObservation::from_config(&configuration, "trade-primary", &second)
+ .expect("second observation");
+ let outcome = host
+ .repositories()
+ .persist_trade_evidence(second, second_observation)
+ .await
+ .expect("second signature");
+ assert!(!outcome.mutation_inserted());
+ assert!(outcome.signed_event_inserted());
+ assert!(outcome.observation_inserted());
+
+ let replay = admitted(&configuration, &first_wire, 1_784_347_200);
+ let replay_observation =
+ RhiTradeSourceObservation::from_config(&configuration, "trade-primary", &replay)
+ .expect("replay observation");
+ let outcome = host
+ .repositories()
+ .persist_trade_evidence(replay, replay_observation)
+ .await
+ .expect("exact replay");
+ assert!(!outcome.mutation_inserted());
+ assert!(!outcome.signed_event_inserted());
+ assert!(!outcome.observation_inserted());
+
+ let later = admitted(&configuration, &first_wire, 1_784_347_202);
+ let later_observation =
+ RhiTradeSourceObservation::from_config(&configuration, "trade-primary", &later)
+ .expect("later observation");
+ let outcome = host
+ .repositories()
+ .persist_trade_evidence(later, later_observation)
+ .await
+ .expect("later observation");
+ assert!(!outcome.mutation_inserted());
+ assert!(!outcome.signed_event_inserted());
+ assert!(outcome.observation_inserted());
+
+ host.close().await.expect("close");
+ let mut connection = offline_connection(&runtime).await;
+ for (query, table, expected) in [
+ (
+ "SELECT COUNT(*) FROM trade_mutations",
+ "trade_mutations",
+ 1_i64,
+ ),
+ ("SELECT COUNT(*) FROM nostr_events", "nostr_events", 2),
+ (
+ "SELECT COUNT(*) FROM relay_observations",
+ "relay_observations",
+ 3,
+ ),
+ ] {
+ let count = sqlx::query_scalar::<_, i64>(query)
+ .fetch_one(&mut connection)
+ .await
+ .expect("count");
+ assert_eq!(count, expected, "{table}");
+ }
+ assert!(
+ sqlx::query("UPDATE trade_mutations SET mutation_id = mutation_id")
+ .execute(&mut connection)
+ .await
+ .is_err()
+ );
+ assert!(
+ sqlx::query("DELETE FROM nostr_events")
+ .execute(&mut connection)
+ .await
+ .is_err()
+ );
+ assert!(
+ sqlx::query("DELETE FROM relay_observations")
+ .execute(&mut connection)
+ .await
+ .is_err()
+ );
+ connection.close().await.expect("offline close");
+
+ let inspection = open_rhi_state_inspection(&runtime, &metadata)
+ .await
+ .expect("inspection");
+ let rejected = admitted(&configuration, &first_wire, 1_784_347_203);
+ let observation =
+ RhiTradeSourceObservation::from_config(&configuration, "trade-primary", &rejected)
+ .expect("observation");
+ let error = inspection
+ .repositories()
+ .persist_trade_evidence(rejected, observation)
+ .await
+ .expect_err("inspection cannot persist");
+ assert_eq!(
+ error.kind(),
+ RhiTradeEvidencePersistenceErrorKind::InvalidMode
+ );
+ inspection.close().await.expect("inspection close");
+}
+
+#[tokio::test]
+async fn durable_mutation_conflict_rolls_back_the_event_and_observation() {
+ let root = tempfile::tempdir().expect("root");
+ let runtime = runtime(root.path());
+ prepare_state_directory(&runtime);
+ let configuration = configuration();
+ let metadata = metadata(&runtime, &configuration);
+ let (applied_at, build) = migration_evidence();
+ initialize_rhi_state(&runtime, &metadata, applied_at, &build)
+ .await
+ .expect("initialize");
+
+ let wire = signed_wire(1);
+ let event = admitted(&configuration, &wire, 1_784_347_200);
+ let mut connection = offline_connection(&runtime).await;
+ insert_mutation(&mut connection, &event, br#"{}"#).await;
+ connection.close().await.expect("offline close");
+
+ let host = open_rhi_state_read_write(&runtime, &metadata, applied_at, &build)
+ .await
+ .expect("writer");
+ let observation =
+ RhiTradeSourceObservation::from_config(&configuration, "trade-primary", &event)
+ .expect("observation");
+ let error = host
+ .repositories()
+ .persist_trade_evidence(event, observation)
+ .await
+ .expect_err("conflicting mutation");
+ assert_eq!(
+ error.kind(),
+ RhiTradeEvidencePersistenceErrorKind::MutationConflict
+ );
+ assert!(Error::source(&error).is_none());
+ host.close().await.expect("close");
+
+ let mut connection = offline_connection(&runtime).await;
+ assert_eq!(
+ sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM nostr_events")
+ .fetch_one(&mut connection)
+ .await
+ .expect("event count"),
+ 0
+ );
+ assert_eq!(
+ sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM relay_observations")
+ .fetch_one(&mut connection)
+ .await
+ .expect("observation count"),
+ 0
+ );
+ connection.close().await.expect("offline close");
+}
+
+#[tokio::test]
+async fn durable_signed_event_conflict_rolls_back_the_observation() {
+ let root = tempfile::tempdir().expect("root");
+ let runtime = runtime(root.path());
+ prepare_state_directory(&runtime);
+ let configuration = configuration();
+ let metadata = metadata(&runtime, &configuration);
+ let (applied_at, build) = migration_evidence();
+ initialize_rhi_state(&runtime, &metadata, applied_at, &build)
+ .await
+ .expect("initialize");
+
+ let wire = signed_wire(1);
+ let wire_json: Value = serde_json::from_slice(&wire).expect("wire JSON");
+ let event = admitted(&configuration, &wire, 1_784_347_200);
+ let mut connection = offline_connection(&runtime).await;
+ insert_mutation(
+ &mut connection,
+ &event,
+ wire_json["content"].as_str().expect("content").as_bytes(),
+ )
+ .await;
+ sqlx::query(
+ r#"INSERT INTO nostr_events (
+ event_id, event_signature, mutation_id, author_pubkey, event_kind,
+ authored_at_unix_s, canonical_event_json
+ ) VALUES (?, ?, ?, ?, ?, ?, ?)"#,
+ )
+ .bind(event.event_id().as_bytes().as_slice())
+ .bind(decode_hex::<64>(wire_json["sig"].as_str().expect("signature")).as_slice())
+ .bind(event.mutation_id().as_bytes().as_slice())
+ .bind(event.mutation().author_pubkey.as_bytes().as_slice())
+ .bind(i64::from(event.event_kind()))
+ .bind(i64::try_from(event.authored_at_unix_seconds()).expect("authored time"))
+ .bind(br#"{}"#.as_slice())
+ .execute(&mut connection)
+ .await
+ .expect("seed conflicting event");
+ connection.close().await.expect("offline close");
+
+ let host = open_rhi_state_read_write(&runtime, &metadata, applied_at, &build)
+ .await
+ .expect("writer");
+ let observation =
+ RhiTradeSourceObservation::from_config(&configuration, "trade-primary", &event)
+ .expect("observation");
+ let error = host
+ .repositories()
+ .persist_trade_evidence(event, observation)
+ .await
+ .expect_err("conflicting signed event");
+ assert_eq!(
+ error.kind(),
+ RhiTradeEvidencePersistenceErrorKind::SignedEventConflict
+ );
+ assert!(Error::source(&error).is_none());
+ host.close().await.expect("close");
+
+ let mut connection = offline_connection(&runtime).await;
+ assert_eq!(
+ sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM relay_observations")
+ .fetch_one(&mut connection)
+ .await
+ .expect("observation count"),
+ 0
+ );
+ connection.close().await.expect("offline close");
+}
+
+#[test]
+fn observation_construction_and_public_diagnostics_are_closed_and_redacted() {
+ let configuration = configuration();
+ let wire = signed_wire(1);
+ let event = admitted(&configuration, &wire, 1_784_347_200);
+ let missing = RhiTradeSourceObservation::from_config(&configuration, "missing", &event)
+ .expect_err("unknown source");
+ assert_eq!(
+ missing.kind(),
+ RhiTradeEvidencePersistenceErrorKind::InvalidObservation
+ );
+ assert!(Error::source(&missing).is_none());
+
+ let observation =
+ RhiTradeSourceObservation::from_config(&configuration, "trade-primary", &event)
+ .expect("observation");
+ let rendered = format!("{observation:?} {missing} {missing:?}");
+ let wire: Value = serde_json::from_slice(&wire).expect("wire JSON");
+ assert!(!rendered.contains("trade-primary"));
+ assert!(!rendered.contains(wire["id"].as_str().expect("id")));
+ assert!(!rendered.contains(wire["sig"].as_str().expect("signature")));
+}