rhi

Coordinated trade for connected markets
git clone https://radroots.dev/git/rhi.git
Log | Files | Refs | README | LICENSE

commit 5bb6539b8cd510d79c1b2381052b1142cd86210a
parent 336f708bfbf4d56b1edfb5dbfded12663aa44cac
Author: triesap <tyson@radroots.org>
Date:   Mon, 24 Aug 2026 04:21:09 +0000

refactor(rhi): persist canonical evidence provenance

- add bounded EOSE-gated relay-source ingestion and typed outcomes
- bind checkpoint and dirty state to schema v4 with generation fencing
- preserve replay and rejection semantics across immutable evidence writes
- advance exact Lib source lock and freeze the reviewed public contract

Diffstat:
MAGENTS.md | 12++++++++----
MCargo.lock | 28++++++++++++++--------------
MCargo.toml | 24++++++++++++------------
MREADME | 35++++++++++++++++++++++++++++++-----
Mcontracts/api_baselines/rhi.txt | 65+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcontracts/services_hardening/trade_evidence_persistence.v1.json | 10++--------
Mcontracts/services_hardening/trade_ingest.v1.json | 3++-
Acontracts/services_hardening/trade_source_ingest.v1.json | 61+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mradroots.service.source-lock.v2.toml | 8++++----
Msrc/lib.rs | 14++++++++++++--
Msrc/runtime_adapters.rs | 8++++++--
Asrc/source_ingest.rs | 1135+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Msrc/state_catalog.rs | 253++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----
Msrc/state_config.rs | 84++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Msrc/state_metadata.rs | 4++++
Msrc/state_trade.rs | 44+++++++++++++++++++++++++++++---------------
Mtests/build_policy.rs | 16++++++++--------
Mtests/package_boundary.rs | 102++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mtests/services_hardening_config_lifecycle.rs | 4++--
Mtests/services_hardening_native_release.rs | 2+-
Mtests/services_hardening_runtime_foundation.rs | 2+-
Atests/services_hardening_source_ingest.rs | 522+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mtests/services_hardening_state_catalog.rs | 72++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------------
Mtests/services_hardening_state_host.rs | 4++--
Mtests/services_hardening_state_resilience.rs | 4++--
Mtests/services_hardening_trade_persistence.rs | 2+-
Mtests/services_hardening_wave_100_b.rs | 2+-
Mtests/source_guards.rs | 2+-
28 files changed, 2407 insertions(+), 115 deletions(-)

diff --git a/AGENTS.md b/AGENTS.md @@ -204,8 +204,9 @@ 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 and schema-v3 - immutable trade-evidence migration. Retain at most 1,024 consecutive + governed RHI schema-v2 configuration-binding migration, schema-v3 + immutable trade-evidence migration, and schema-v4 source-checkpoint and + dirty-generation 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. Persist each canonical @@ -213,8 +214,11 @@ 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 + checkpoints or dirty generation. The composed relay-source ingest path may + advance only its exact scoped checkpoint after complete EOSE evidence and + may dirty a trade only for newly inserted mutation or signed-event evidence; + replayed source observation alone does neither. 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. diff --git a/Cargo.lock b/Cargo.lock @@ -1576,7 +1576,7 @@ checksum = "f8dcc9c7d52a811697d2151c701e0d08956f92b0e24136cf4cf27b57a6a0d9bf" [[package]] name = "radroots_blossom" version = "0.1.0-alpha" -source = "git+https://github.com/radrootslabs/lib?rev=79d7818c8fe22a425f9524b884ddf59d25f0ef89#79d7818c8fe22a425f9524b884ddf59d25f0ef89" +source = "git+https://github.com/radrootslabs/lib?rev=21b11e7a5120ea949f7ad0838c746873fc73aac2#21b11e7a5120ea949f7ad0838c746873fc73aac2" dependencies = [ "mediatype", "serde", @@ -1588,7 +1588,7 @@ dependencies = [ [[package]] name = "radroots_core" version = "0.1.0-alpha" -source = "git+https://github.com/radrootslabs/lib?rev=79d7818c8fe22a425f9524b884ddf59d25f0ef89#79d7818c8fe22a425f9524b884ddf59d25f0ef89" +source = "git+https://github.com/radrootslabs/lib?rev=21b11e7a5120ea949f7ad0838c746873fc73aac2#21b11e7a5120ea949f7ad0838c746873fc73aac2" dependencies = [ "rust_decimal", "serde", @@ -1597,7 +1597,7 @@ dependencies = [ [[package]] name = "radroots_event" version = "0.1.0-alpha" -source = "git+https://github.com/radrootslabs/lib?rev=79d7818c8fe22a425f9524b884ddf59d25f0ef89#79d7818c8fe22a425f9524b884ddf59d25f0ef89" +source = "git+https://github.com/radrootslabs/lib?rev=21b11e7a5120ea949f7ad0838c746873fc73aac2#21b11e7a5120ea949f7ad0838c746873fc73aac2" dependencies = [ "hex", "jiff-tzdb", @@ -1615,7 +1615,7 @@ dependencies = [ [[package]] name = "radroots_event_codec" version = "0.1.0-alpha" -source = "git+https://github.com/radrootslabs/lib?rev=79d7818c8fe22a425f9524b884ddf59d25f0ef89#79d7818c8fe22a425f9524b884ddf59d25f0ef89" +source = "git+https://github.com/radrootslabs/lib?rev=21b11e7a5120ea949f7ad0838c746873fc73aac2#21b11e7a5120ea949f7ad0838c746873fc73aac2" dependencies = [ "hex", "radroots_blossom", @@ -1632,7 +1632,7 @@ dependencies = [ [[package]] name = "radroots_identity" version = "0.1.0-alpha" -source = "git+https://github.com/radrootslabs/lib?rev=79d7818c8fe22a425f9524b884ddf59d25f0ef89#79d7818c8fe22a425f9524b884ddf59d25f0ef89" +source = "git+https://github.com/radrootslabs/lib?rev=21b11e7a5120ea949f7ad0838c746873fc73aac2#21b11e7a5120ea949f7ad0838c746873fc73aac2" dependencies = [ "k256", "serde", @@ -1642,7 +1642,7 @@ dependencies = [ [[package]] name = "radroots_nostr" version = "0.1.0-alpha" -source = "git+https://github.com/radrootslabs/lib?rev=79d7818c8fe22a425f9524b884ddf59d25f0ef89#79d7818c8fe22a425f9524b884ddf59d25f0ef89" +source = "git+https://github.com/radrootslabs/lib?rev=21b11e7a5120ea949f7ad0838c746873fc73aac2#21b11e7a5120ea949f7ad0838c746873fc73aac2" dependencies = [ "nostr", "radroots_event", @@ -1656,7 +1656,7 @@ dependencies = [ [[package]] name = "radroots_protocol" version = "0.1.0-alpha" -source = "git+https://github.com/radrootslabs/lib?rev=79d7818c8fe22a425f9524b884ddf59d25f0ef89#79d7818c8fe22a425f9524b884ddf59d25f0ef89" +source = "git+https://github.com/radrootslabs/lib?rev=21b11e7a5120ea949f7ad0838c746873fc73aac2#21b11e7a5120ea949f7ad0838c746873fc73aac2" dependencies = [ "serde", ] @@ -1664,7 +1664,7 @@ dependencies = [ [[package]] name = "radroots_runtime_paths" version = "0.1.0-alpha" -source = "git+https://github.com/radrootslabs/lib?rev=79d7818c8fe22a425f9524b884ddf59d25f0ef89#79d7818c8fe22a425f9524b884ddf59d25f0ef89" +source = "git+https://github.com/radrootslabs/lib?rev=21b11e7a5120ea949f7ad0838c746873fc73aac2#21b11e7a5120ea949f7ad0838c746873fc73aac2" dependencies = [ "serde", "thiserror 1.0.69", @@ -1673,7 +1673,7 @@ dependencies = [ [[package]] name = "radroots_secrets" version = "0.1.0-alpha" -source = "git+https://github.com/radrootslabs/lib?rev=79d7818c8fe22a425f9524b884ddf59d25f0ef89#79d7818c8fe22a425f9524b884ddf59d25f0ef89" +source = "git+https://github.com/radrootslabs/lib?rev=21b11e7a5120ea949f7ad0838c746873fc73aac2#21b11e7a5120ea949f7ad0838c746873fc73aac2" dependencies = [ "chacha20poly1305", "serde", @@ -1685,7 +1685,7 @@ dependencies = [ [[package]] name = "radroots_service_host" version = "0.1.0-alpha" -source = "git+https://github.com/radrootslabs/lib?rev=79d7818c8fe22a425f9524b884ddf59d25f0ef89#79d7818c8fe22a425f9524b884ddf59d25f0ef89" +source = "git+https://github.com/radrootslabs/lib?rev=21b11e7a5120ea949f7ad0838c746873fc73aac2#21b11e7a5120ea949f7ad0838c746873fc73aac2" dependencies = [ "bytes", "fs2", @@ -1706,7 +1706,7 @@ dependencies = [ [[package]] name = "radroots_service_sqlite" version = "0.1.0-alpha" -source = "git+https://github.com/radrootslabs/lib?rev=79d7818c8fe22a425f9524b884ddf59d25f0ef89#79d7818c8fe22a425f9524b884ddf59d25f0ef89" +source = "git+https://github.com/radrootslabs/lib?rev=21b11e7a5120ea949f7ad0838c746873fc73aac2#21b11e7a5120ea949f7ad0838c746873fc73aac2" dependencies = [ "fs2", "futures", @@ -1724,7 +1724,7 @@ dependencies = [ [[package]] name = "radroots_storage" version = "0.1.0-alpha" -source = "git+https://github.com/radrootslabs/lib?rev=79d7818c8fe22a425f9524b884ddf59d25f0ef89#79d7818c8fe22a425f9524b884ddf59d25f0ef89" +source = "git+https://github.com/radrootslabs/lib?rev=21b11e7a5120ea949f7ad0838c746873fc73aac2#21b11e7a5120ea949f7ad0838c746873fc73aac2" dependencies = [ "radroots_event", "radroots_event_codec", @@ -1737,7 +1737,7 @@ dependencies = [ [[package]] name = "radroots_trade" version = "0.1.0-alpha" -source = "git+https://github.com/radrootslabs/lib?rev=79d7818c8fe22a425f9524b884ddf59d25f0ef89#79d7818c8fe22a425f9524b884ddf59d25f0ef89" +source = "git+https://github.com/radrootslabs/lib?rev=21b11e7a5120ea949f7ad0838c746873fc73aac2#21b11e7a5120ea949f7ad0838c746873fc73aac2" dependencies = [ "hex", "radroots_core", @@ -1751,7 +1751,7 @@ dependencies = [ [[package]] name = "radroots_transport" version = "0.1.0-alpha" -source = "git+https://github.com/radrootslabs/lib?rev=79d7818c8fe22a425f9524b884ddf59d25f0ef89#79d7818c8fe22a425f9524b884ddf59d25f0ef89" +source = "git+https://github.com/radrootslabs/lib?rev=21b11e7a5120ea949f7ad0838c746873fc73aac2#21b11e7a5120ea949f7ad0838c746873fc73aac2" dependencies = [ "radroots_event", "radroots_identity", 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 = 3 +state_contract_version = 4 admin_contract_version = 1 status_contract_version = 1 provider_contract_version = 1 @@ -51,17 +51,17 @@ service-host = [] workspace = true [dependencies] -radroots_event = { git = "https://github.com/radrootslabs/lib", rev = "79d7818c8fe22a425f9524b884ddf59d25f0ef89", version = "=0.1.0-alpha", features = ["serde"] } -radroots_event_codec = { git = "https://github.com/radrootslabs/lib", rev = "79d7818c8fe22a425f9524b884ddf59d25f0ef89", version = "=0.1.0-alpha", features = ["json"] } -radroots_nostr = { git = "https://github.com/radrootslabs/lib", rev = "79d7818c8fe22a425f9524b884ddf59d25f0ef89", version = "=0.1.0-alpha", features = ["events"] } -radroots_protocol = { git = "https://github.com/radrootslabs/lib", rev = "79d7818c8fe22a425f9524b884ddf59d25f0ef89", version = "=0.1.0-alpha" } -radroots_runtime_paths = { git = "https://github.com/radrootslabs/lib", rev = "79d7818c8fe22a425f9524b884ddf59d25f0ef89", version = "=0.1.0-alpha" } -radroots_service_host = { git = "https://github.com/radrootslabs/lib", rev = "79d7818c8fe22a425f9524b884ddf59d25f0ef89", version = "=0.1.0-alpha" } -radroots_service_sqlite = { git = "https://github.com/radrootslabs/lib", rev = "79d7818c8fe22a425f9524b884ddf59d25f0ef89", version = "=0.1.0-alpha" } -radroots_storage = { git = "https://github.com/radrootslabs/lib", rev = "79d7818c8fe22a425f9524b884ddf59d25f0ef89", version = "=0.1.0-alpha", default-features = false } -radroots_transport = { git = "https://github.com/radrootslabs/lib", rev = "79d7818c8fe22a425f9524b884ddf59d25f0ef89", version = "=0.1.0-alpha", default-features = false, features = ["std"] } -radroots_secrets = { git = "https://github.com/radrootslabs/lib", rev = "79d7818c8fe22a425f9524b884ddf59d25f0ef89", version = "=0.1.0-alpha" } -radroots_trade = { git = "https://github.com/radrootslabs/lib", rev = "79d7818c8fe22a425f9524b884ddf59d25f0ef89", version = "=0.1.0-alpha" } +radroots_event = { git = "https://github.com/radrootslabs/lib", rev = "21b11e7a5120ea949f7ad0838c746873fc73aac2", version = "=0.1.0-alpha", features = ["serde"] } +radroots_event_codec = { git = "https://github.com/radrootslabs/lib", rev = "21b11e7a5120ea949f7ad0838c746873fc73aac2", version = "=0.1.0-alpha", features = ["json"] } +radroots_nostr = { git = "https://github.com/radrootslabs/lib", rev = "21b11e7a5120ea949f7ad0838c746873fc73aac2", version = "=0.1.0-alpha", features = ["events"] } +radroots_protocol = { git = "https://github.com/radrootslabs/lib", rev = "21b11e7a5120ea949f7ad0838c746873fc73aac2", version = "=0.1.0-alpha" } +radroots_runtime_paths = { git = "https://github.com/radrootslabs/lib", rev = "21b11e7a5120ea949f7ad0838c746873fc73aac2", version = "=0.1.0-alpha" } +radroots_service_host = { git = "https://github.com/radrootslabs/lib", rev = "21b11e7a5120ea949f7ad0838c746873fc73aac2", version = "=0.1.0-alpha" } +radroots_service_sqlite = { git = "https://github.com/radrootslabs/lib", rev = "21b11e7a5120ea949f7ad0838c746873fc73aac2", version = "=0.1.0-alpha" } +radroots_storage = { git = "https://github.com/radrootslabs/lib", rev = "21b11e7a5120ea949f7ad0838c746873fc73aac2", version = "=0.1.0-alpha", default-features = false } +radroots_transport = { git = "https://github.com/radrootslabs/lib", rev = "21b11e7a5120ea949f7ad0838c746873fc73aac2", version = "=0.1.0-alpha", default-features = false, features = ["std"] } +radroots_secrets = { git = "https://github.com/radrootslabs/lib", rev = "21b11e7a5120ea949f7ad0838c746873fc73aac2", version = "=0.1.0-alpha" } +radroots_trade = { git = "https://github.com/radrootslabs/lib", rev = "21b11e7a5120ea949f7ad0838c746873fc73aac2", version = "=0.1.0-alpha" } chacha20poly1305 = { version = "0.10" } clap = { version = "4", features = ["derive"] } diff --git a/README b/README @@ -95,6 +95,29 @@ 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). +## Bounded relay-source ingestion + +`ingest_rhi_trade_source` resolves one explicit configured `nostr_relay` +source and its read-authorized relay, then fetches the exact five trade-mutation +kinds under the requested trade's indexed `#d` tag. It uses one absolute +configured deadline, pages at no more than 1,000 events, and admits at most +4,096 distinct event identities and 8 MiB of original event bytes per attempt. +Every returned event still passes the sealed signature, identifier, mutation, +tag, trade, and authored-time boundary before it can affect durable state. + +Only exact-target EOSE before the deadline is complete. Cancellation, +unavailability, bounded-result exhaustion, partial or malformed results, and +unsupported operation remain distinct incomplete outcomes and cannot advance a +checkpoint. Checkpoints are scoped by source, selector, evidence-policy digest, +and trade; their authored-time plus verified-event-ID cursor uses configured +overlap so equal timestamps remain discoverable. One short generation-fenced +SQLx transaction persists admitted evidence and an eligible checkpoint. +Relevant newly inserted mutation or signed-event evidence advances the +per-trade dirty generation; rejection, verified-event replay, repeated source +observation, and operational retry do not. Source fetching never occurs inside +a database transaction. The exact machine contract is +[`trade_source_ingest.v1.json`](contracts/services_hardening/trade_source_ingest.v1.json). + ## Existing-state runtime foundation `open_rhi_runtime_foundation` opens only an already initialized database from @@ -191,11 +214,13 @@ 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 configuration migration and schema-v3 immutable trade-evidence -migration. Version three contains the six shared immutable service-metadata and +schema-v2 configuration migration, schema-v3 immutable trade-evidence +migration, and schema-v4 source-checkpoint and dirty-generation migration. +Version four 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. +separate immutable tables for canonical mutations, signed Nostr events, and +accepted source observations, and generation-guarded relay checkpoints and +per-trade dirty generations 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. @@ -212,7 +237,7 @@ 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 migrations and typed state transactions. -Reconciliation checkpoints, completion, dirty generation, reports, +Durable source-attempt/completion inventory, reconciliation jobs, reports, attestations, and publication state remain reserved for their later owning checkpoints. diff --git a/contracts/api_baselines/rhi.txt b/contracts/api_baselines/rhi.txt @@ -15,6 +15,7 @@ pub use rhi::RadrootsServiceInstanceArtifacts pub use rhi::RuntimeContext pub use rhi::RuntimeContextSource pub use rhi::ServiceId +pub use rhi::TradeId pub use rhi::UnixTimeSeconds pub use rhi::WallClock pub use rhi::WallClockError @@ -305,6 +306,24 @@ pub rhi::RhiTradeMutationAdmissionErrorKind::TooManyExtraFields pub rhi::RhiTradeMutationAdmissionErrorKind::TooManyTagElements pub rhi::RhiTradeMutationAdmissionErrorKind::TooManyTags pub rhi::RhiTradeMutationAdmissionErrorKind::UnsupportedKind +pub enum rhi::RhiTradeSourceCompletion +pub rhi::RhiTradeSourceCompletion::Complete +pub rhi::RhiTradeSourceCompletion::IncompleteResourceLimit +pub rhi::RhiTradeSourceCompletion::IncompleteTimeout +pub rhi::RhiTradeSourceCompletion::IncompleteUnavailable +pub rhi::RhiTradeSourceCompletion::IncompleteUnknown +pub rhi::RhiTradeSourceCompletion::Unsupported +impl rhi::RhiTradeSourceCompletion +pub const fn rhi::RhiTradeSourceCompletion::code(self) -> &'static str +pub enum rhi::RhiTradeSourceIngestErrorKind +pub rhi::RhiTradeSourceIngestErrorKind::CommitOutcomeUnknown +pub rhi::RhiTradeSourceIngestErrorKind::GenerationConflict +pub rhi::RhiTradeSourceIngestErrorKind::InvalidConfiguration +pub rhi::RhiTradeSourceIngestErrorKind::InvalidInput +pub rhi::RhiTradeSourceIngestErrorKind::InvalidMode +pub rhi::RhiTradeSourceIngestErrorKind::Storage +impl rhi::RhiTradeSourceIngestErrorKind +pub const fn rhi::RhiTradeSourceIngestErrorKind::code(self) -> &'static str pub enum rhi::TradeAgreementAttestationBackend pub rhi::TradeAgreementAttestationBackend::LocalStatementHash impl rhi::TradeAgreementAttestationBackend @@ -760,6 +779,9 @@ 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::RhiTradeDirtyGeneration(_) +impl rhi::RhiTradeDirtyGeneration +pub const fn rhi::RhiTradeDirtyGeneration::get(self) -> u64 pub struct rhi::RhiTradeEvidencePersistenceError impl rhi::RhiTradeEvidencePersistenceError pub const fn rhi::RhiTradeEvidencePersistenceError::code(self) -> &'static str @@ -803,6 +825,42 @@ 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::RhiTradeSourceAttempt +impl rhi::RhiTradeSourceAttempt +pub fn rhi::RhiTradeSourceAttempt::new(impl core::convert::AsRef<str>, radroots_service_host::time::UnixTimeSeconds, rhi::RhiTradeMutationObservedAtUnixSeconds, rhi::RhiTradeMutationAuthoredTimePolicy) -> core::result::Result<Self, rhi::RhiTradeSourceIngestError> +impl core::fmt::Debug for rhi::RhiTradeSourceAttempt +pub fn rhi::RhiTradeSourceAttempt::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiTradeSourceCursor +impl rhi::RhiTradeSourceCursor +pub const fn rhi::RhiTradeSourceCursor::created_at_unix_seconds(self) -> u64 +pub const fn rhi::RhiTradeSourceCursor::event_id(self) -> [u8; 32] +impl core::fmt::Debug for rhi::RhiTradeSourceCursor +pub fn rhi::RhiTradeSourceCursor::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiTradeSourceIngestError +impl rhi::RhiTradeSourceIngestError +pub const fn rhi::RhiTradeSourceIngestError::code(self) -> &'static str +pub const fn rhi::RhiTradeSourceIngestError::kind(self) -> rhi::RhiTradeSourceIngestErrorKind +impl core::error::Error for rhi::RhiTradeSourceIngestError +impl core::fmt::Debug for rhi::RhiTradeSourceIngestError +pub fn rhi::RhiTradeSourceIngestError::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +impl core::fmt::Display for rhi::RhiTradeSourceIngestError +pub fn rhi::RhiTradeSourceIngestError::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiTradeSourceIngestOutcome +impl rhi::RhiTradeSourceIngestOutcome +pub const fn rhi::RhiTradeSourceIngestOutcome::admitted_events(self) -> u32 +pub const fn rhi::RhiTradeSourceIngestOutcome::checkpoint(self) -> core::option::Option<rhi::RhiTradeSourceCursor> +pub const fn rhi::RhiTradeSourceIngestOutcome::checkpoint_advanced(self) -> bool +pub const fn rhi::RhiTradeSourceIngestOutcome::completion(self) -> rhi::RhiTradeSourceCompletion +pub const fn rhi::RhiTradeSourceIngestOutcome::dirty_generation(self) -> core::option::Option<rhi::RhiTradeDirtyGeneration> +pub const fn rhi::RhiTradeSourceIngestOutcome::dirty_generation_advanced(self) -> bool +pub const fn rhi::RhiTradeSourceIngestOutcome::duplicate_events(self) -> u32 +pub const fn rhi::RhiTradeSourceIngestOutcome::inserted_mutations(self) -> u32 +pub const fn rhi::RhiTradeSourceIngestOutcome::inserted_observations(self) -> u32 +pub const fn rhi::RhiTradeSourceIngestOutcome::inserted_signed_events(self) -> u32 +pub const fn rhi::RhiTradeSourceIngestOutcome::received_events(self) -> u32 +pub const fn rhi::RhiTradeSourceIngestOutcome::rejected_events(self) -> u32 +impl core::fmt::Debug for rhi::RhiTradeSourceIngestOutcome +pub fn rhi::RhiTradeSourceIngestOutcome::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result pub struct rhi::RhiTradeSourceObservation impl rhi::RhiTradeSourceObservation pub fn rhi::RhiTradeSourceObservation::from_config(&rhi::RhiConfigDocumentV1, &str, &rhi::RhiAdmittedTradeMutationEvent) -> core::result::Result<Self, rhi::RhiTradeEvidencePersistenceError> @@ -901,6 +959,9 @@ 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_STATE_SCHEMA_VERSION_4_MIGRATION_SHA256: [u8; 32] +pub const rhi::RHI_STATE_SCHEMA_VERSION_4_OBJECT_COUNT: u32 +pub const rhi::RHI_STATE_SCHEMA_VERSION_4_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 @@ -909,6 +970,9 @@ 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_TRADE_SOURCE_INGEST_CONTRACT_VERSION: u32 +pub const rhi::RHI_TRADE_SOURCE_RESULT_MAX_BYTES: usize +pub const rhi::RHI_TRADE_SOURCE_RESULT_MAX_EVENTS: usize pub const rhi::RHI_WRAPPING_CREDENTIAL_ARTIFACT_BYTES: usize pub const rhi::RHI_WRAPPING_CREDENTIAL_CONTRACT_VERSION: u32 pub trait rhi::RhiCredentialAccess: core::marker::Send + core::marker::Sync @@ -923,6 +987,7 @@ pub fn rhi::admit_rhi_trade_mutation_event(rhi::RhiTradeMutationAdmissionLimits, pub async fn rhi::apply_rhi_configuration(&rhi::RhiRuntimeContext, &rhi::RhiConfigDocumentV1, &rhi::RhiConfigDocumentV1, radroots_service_sqlite::migration::MigrationAppliedAtUnixSeconds, &radroots_service_sqlite::migration::MigrationBuildIdentity) -> core::result::Result<rhi::RhiConfigApplyOutcome, rhi::RhiConfigApplyError> pub fn rhi::attest_projection_claim(&radroots_trade::trade_contract_v1::RadrootsTradeProjectionV1, &radroots_event::id::MutationId, &rhi::TradeAgreementAttestationPolicy) -> core::result::Result<rhi::TradeAgreementAttestationReportV1, rhi::TradeAgreementAttestationError> pub async fn rhi::finalize_rhi_state_restore(rhi::RhiStagedStateRestore) -> core::result::Result<(), rhi::RhiStateMaintenanceError> +pub async fn rhi::ingest_rhi_trade_source(&rhi::RhiStateRepositories<'_>, &rhi::RhiTransportAdapters, &rhi::RhiConfigDocumentV1, &str, radroots_event::id::TradeId, rhi::RhiTradeSourceAttempt) -> core::result::Result<rhi::RhiTradeSourceIngestOutcome, rhi::RhiTradeSourceIngestError> pub async fn rhi::initialize_rhi_state(&rhi::RhiRuntimeContext, &rhi::RhiStateMetadata, radroots_service_sqlite::migration::MigrationAppliedAtUnixSeconds, &radroots_service_sqlite::migration::MigrationBuildIdentity) -> core::result::Result<(), rhi::RhiStateHostError> pub fn rhi::open_rhi_encrypted_identity(&rhi::RhiIdentityEnvelopeBinding, &rhi::RhiWrappingCredential) -> core::result::Result<rhi::RhiDecryptedIdentity, rhi::RhiEncryptedIdentityEnvelopeError> pub async fn rhi::open_rhi_runtime_foundation(rhi::RhiRuntimeContext, rhi::RhiConfigDocumentV1, rhi::RhiRuntimeAdapters, radroots_service_sqlite::migration::MigrationAppliedAtUnixSeconds, &radroots_service_sqlite::migration::MigrationBuildIdentity) -> core::result::Result<rhi::RhiRuntimeFoundation, rhi::RhiRuntimeFoundationError> diff --git a/contracts/services_hardening/trade_evidence_persistence.v1.json b/contracts/services_hardening/trade_evidence_persistence.v1.json @@ -45,12 +45,6 @@ "dirty_generation": false }, "errors": "crate_owned_source_free_redacted", - "deferred": [ - "nostr_relay_source_adapter", - "source_completion", - "checkpoint_advance", - "dirty_generation", - "reconciliation", - "publication" - ] + "source_ingest_contract": "contracts/services_hardening/trade_source_ingest.v1.json", + "deferred": ["reconciliation", "publication"] } diff --git a/contracts/services_hardening/trade_ingest.v1.json b/contracts/services_hardening/trade_ingest.v1.json @@ -53,5 +53,6 @@ "errors": "crate_owned_source_free_redacted", "effects": { "filesystem": false, "sqlite": false, "clock_read": false, "network": false, "checkpoint": false, "dirty_generation": false }, "persistence_contract": "contracts/services_hardening/trade_evidence_persistence.v1.json", - "deferred": ["source_adapter_and_completion", "checkpoint_and_dirty_generation", "reconciliation", "publication"] + "source_ingest_contract": "contracts/services_hardening/trade_source_ingest.v1.json", + "deferred": ["reconciliation", "publication"] } diff --git a/contracts/services_hardening/trade_source_ingest.v1.json b/contracts/services_hardening/trade_source_ingest.v1.json @@ -0,0 +1,61 @@ +{ + "schema": "radroots.rhi.trade-source-ingest.v1", + "contract_version": 1, + "configuration_authority": "contracts/services_hardening/evidence_policy.v1.json", + "source": { + "kind": "nostr_relay", + "configured_source_required": true, + "configured_read_relay_required": true, + "implicit_source": false, + "selector_id": "trade_mutation_lineage_v1", + "selector": { + "kinds": [3470, 3471, 3472, 3473, 3474], + "exact_tag": "#d", + "exact_tag_value": "requested_trade_id", + "initial_since": "max_zero_attempt_started_at_utc_minus_lookback_seconds", + "resume_since": "max_zero_cursor_created_at_minus_overlap_seconds" + }, + "deadline": "single_absolute_configured_deadline", + "page_events_maximum": 1000, + "result_events_maximum": 4096, + "result_original_event_bytes_maximum": 8388608 + }, + "completion": { + "complete": "exact_target_eose_before_deadline", + "cancelled": "incomplete_timeout", + "unavailable_or_retryable_failure": "incomplete_unavailable", + "bounded_result_exceeded": "incomplete_resource_limit", + "terminal_failure_partial_or_malformed_result": "incomplete_unknown", + "unsupported_operation": "unsupported", + "checkpoint_allowed_only_for": "complete" + }, + "admission": { + "boundary": "admit_rhi_trade_mutation_event", + "trade_id": "exact_requested_trade_id", + "deduplicate_by": "verified_event_id", + "ordering": ["created_at_unix_seconds_ascending", "event_id_bytes_ascending", "signature_bytes_ascending"], + "rejected_event": "cannot_advance_checkpoint_or_dirty_generation" + }, + "checkpoint": { + "scope": ["source_id", "selector_id", "evidence_policy_sha256", "trade_id"], + "cursor": ["created_at_unix_seconds", "verified_event_id"], + "ordering": "lexicographic_ascending", + "equal_timestamp_safe": true, + "advance": "greatest_new_admitted_cursor_after_complete_only", + "write": "same_generation_fenced_sqlx_transaction_as_evidence" + }, + "dirty_generation": { + "scope": "trade_id", + "advance_when": ["new_canonical_mutation", "new_verified_signed_event", "governing_policy_change"], + "do_not_advance_when": ["rejected_event", "duplicate_event", "repeated_source_observation", "operational_retry"], + "write": "same_generation_fenced_sqlx_transaction_as_evidence" + }, + "transaction": { + "source_fetch_inside_transaction": false, + "initial_fence": ["checkpoint", "dirty_generation"], + "commit_effects": ["canonical_mutations", "verified_signed_events", "source_observations", "eligible_checkpoint", "eligible_dirty_generation"], + "commit_outcome_unknown": "explicit_error" + }, + "diagnostics": "crate_owned_source_free_redacted", + "deferred": ["durable_source_attempts_and_completion_inventory", "reconciliation_jobs_and_leases", "coverage_and_outcome", "reports_and_attestations", "publication"] +} diff --git a/radroots.service.source-lock.v2.toml b/radroots.service.source-lock.v2.toml @@ -2,12 +2,12 @@ schema = "radroots.service.source-lock.v2" contract_version = 2 service = "rhi" repository = "https://github.com/radrootslabs/lib" -revision = "79d7818c8fe22a425f9524b884ddf59d25f0ef89" +revision = "21b11e7a5120ea949f7ad0838c746873fc73aac2" architecture = "radroots.crates.release.v2" workspace_catalog_sha256 = "deca0c080deae187ff8186c0708903e42f41ea57f77c5f91581e23aa561164a4" version = "0.1.0-alpha" -source_archive_sha256 = "6f95d80d5ecb9d026ba6f9145b60fff67e370b2c175c438537319278e6974fa6" -cargo_lock_sha256 = "99d8cfcf6e50ec54e3557219ba5ffd0c999df3a17155ed609f2919efe21e2979" +source_archive_sha256 = "7e584a4b679264620d7bb6cf0a7028cc7651b33977b263c213f4e7b29c0e5a19" +cargo_lock_sha256 = "b85bee310965fc4c4da8f6641007840f0d736c194d3d0fb74ea3db757e7f1433" rust_version = "1.97.1" host_feature_profile = "service-host" @@ -16,7 +16,7 @@ material = "absent" [contract_versions] config = 1 -state = 3 +state = 4 admin = 1 status = 1 provider = 1 diff --git a/src/lib.rs b/src/lib.rs @@ -11,6 +11,7 @@ mod identity_envelope; mod runtime_adapters; mod runtime_context; mod runtime_foundation; +mod source_ingest; mod state_catalog; mod state_config; mod state_host; @@ -55,6 +56,7 @@ pub use identity_envelope::{ RhiIdentityRole, RhiWrappingCredential, open_rhi_encrypted_identity, provision_rhi_encrypted_identity, }; +pub use radroots_event::id::TradeId; pub use radroots_runtime_paths::{ INSTANCE_ID_MAX_BYTES, InstanceId, RadrootsHostEnvironment, RadrootsPathProfile, RadrootsPathResolver, RadrootsPlatform, RadrootsServiceInstanceArtifacts, RuntimeContext, @@ -80,14 +82,22 @@ pub use runtime_foundation::{ RhiRuntimeFoundationErrorKind, RhiRuntimePrerequisite, RhiRuntimeReadiness, RhiRuntimeReadinessReason, open_rhi_runtime_foundation, }; +pub use source_ingest::{ + RHI_TRADE_SOURCE_INGEST_CONTRACT_VERSION, RHI_TRADE_SOURCE_RESULT_MAX_BYTES, + RHI_TRADE_SOURCE_RESULT_MAX_EVENTS, RhiTradeDirtyGeneration, RhiTradeSourceAttempt, + RhiTradeSourceCompletion, RhiTradeSourceCursor, RhiTradeSourceIngestError, + RhiTradeSourceIngestErrorKind, RhiTradeSourceIngestOutcome, ingest_rhi_trade_source, +}; pub use state_catalog::{ RHI_MIGRATION_CATALOG_SHA256, RHI_STATE_BASE_SCHEMA_VERSION, RHI_STATE_SCHEMA_CATALOG_SHA256, 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, 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, + RHI_STATE_SCHEMA_VERSION_3_SHA256, RHI_STATE_SCHEMA_VERSION_4_MIGRATION_SHA256, + RHI_STATE_SCHEMA_VERSION_4_OBJECT_COUNT, RHI_STATE_SCHEMA_VERSION_4_SHA256, + RhiStateCatalogError, RhiStateCatalogErrorKind, rhi_migration_catalog, rhi_schema_catalog, + validate_rhi_state_catalogs, }; pub use state_config::{ RHI_CONFIG_BINDING_MAX_GENERATIONS, RhiConfigApplyError, RhiConfigApplyErrorKind, diff --git a/src/runtime_adapters.rs b/src/runtime_adapters.rs @@ -240,7 +240,7 @@ impl fmt::Debug for RhiTimeEntropyAdapters { /// remain sealed inside RHI so concrete transports and detachable I/O handles /// do not become public runtime authority. pub struct RhiTransportAdapters { - _evidence_source: Arc<dyn EventSource>, + evidence_source: Arc<dyn EventSource>, _evidence_subscriber: Arc<dyn EventSubscriber>, _publication_sink: Arc<dyn EventSink>, } @@ -254,11 +254,15 @@ impl RhiTransportAdapters { publication_sink: Arc<dyn EventSink>, ) -> Self { Self { - _evidence_source: evidence_source, + evidence_source, _evidence_subscriber: evidence_subscriber, _publication_sink: publication_sink, } } + + pub(crate) fn evidence_source(&self) -> &dyn EventSource { + self.evidence_source.as_ref() + } } impl fmt::Debug for RhiTransportAdapters { diff --git a/src/source_ingest.rs b/src/source_ingest.rs @@ -0,0 +1,1135 @@ +//! Bounded relay-source ingestion with generation-fenced checkpoint commit. + +use core::{cmp::Ordering, fmt}; +use std::{collections::BTreeSet, error::Error}; + +use radroots_event::{SignedEvent, id::TradeId}; +use radroots_service_host::UnixTimeSeconds; +use radroots_service_sqlite::{ + ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind, +}; +use radroots_transport::{ + FetchRequest, Target, TargetSet, + outcome::FetchTargetState, + source::{FetchBounds, FetchCursor, FetchSelector, NextPage}, +}; +use serde_json::Value; +use sqlx::Row; + +use crate::{ + RhiAdmittedTradeMutationEvent, RhiConfigDocumentV1, RhiEvidencePolicyDigest, RhiStateHostMode, + RhiStateRepositories, RhiTradeMutationAdmissionLimits, RhiTradeMutationAuthoredTimePolicy, + RhiTradeMutationObservedAtUnixSeconds, RhiTradeSourceObservation, RhiTransportAdapters, + admit_rhi_trade_mutation_event, state_metadata, + state_trade::{PersistenceOperationError, PersistenceRecord, persist}, +}; + +/// Exact version of the RHI relay-source ingestion contract. +pub const RHI_TRADE_SOURCE_INGEST_CONTRACT_VERSION: u32 = 1; + +/// Maximum distinct event identities admitted from one source attempt. +pub const RHI_TRADE_SOURCE_RESULT_MAX_EVENTS: usize = 4_096; + +/// Maximum aggregate original event bytes admitted from one source attempt. +pub const RHI_TRADE_SOURCE_RESULT_MAX_BYTES: usize = 8 * 1024 * 1024; + +const SOURCE_SELECTOR: &str = "trade_mutation_lineage_v1"; +const FETCH_REQUEST_ID_MAX_BYTES: usize = 256; +const FETCH_PAGE_MAX_EVENTS: u16 = 1_000; +const EVENT_KINDS: [u32; 5] = [3470, 3471, 3472, 3473, 3474]; + +const READ_CHECKPOINT_SQL: &str = r#"SELECT + cursor_created_at_unix_s, + length(cursor_event_id) AS cursor_event_id_bytes, + substr(cursor_event_id, 1, 33) AS cursor_event_id, + revision, + completed_at_unix_s +FROM relay_checkpoints +WHERE source_id = ? AND selector_id = ? AND evidence_policy_sha256 = ? AND trade_id = ? +LIMIT 1"#; +const INSERT_CHECKPOINT_SQL: &str = r#"INSERT INTO relay_checkpoints ( + source_id, selector_id, evidence_policy_sha256, trade_id, + cursor_created_at_unix_s, cursor_event_id, revision, completed_at_unix_s +) VALUES (?, ?, ?, ?, ?, ?, 1, ?)"#; +const UPDATE_CHECKPOINT_SQL: &str = r#"UPDATE relay_checkpoints +SET cursor_created_at_unix_s = ?, cursor_event_id = ?, + revision = revision + 1, completed_at_unix_s = ? +WHERE source_id = ? AND selector_id = ? AND evidence_policy_sha256 = ? AND trade_id = ? + AND revision = ?"#; +const READ_DIRTY_SQL: &str = r#"SELECT generation, + length(evidence_policy_sha256) AS evidence_policy_bytes, + substr(evidence_policy_sha256, 1, 33) AS evidence_policy_sha256, + updated_at_unix_s +FROM trade_dirty_generations +WHERE trade_id = ? +LIMIT 1"#; +const INSERT_DIRTY_SQL: &str = r#"INSERT INTO trade_dirty_generations ( + trade_id, generation, evidence_policy_sha256, updated_at_unix_s +) VALUES (?, 1, ?, ?)"#; +const UPDATE_DIRTY_SQL: &str = r#"UPDATE trade_dirty_generations +SET generation = generation + 1, evidence_policy_sha256 = ?, updated_at_unix_s = ? +WHERE trade_id = ? AND generation = ?"#; + +/// Stable terminal classification for one exact relay-source attempt. +#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)] +pub enum RhiTradeSourceCompletion { + Complete, + IncompleteTimeout, + IncompleteUnavailable, + IncompleteResourceLimit, + IncompleteUnknown, + Unsupported, +} + +impl RhiTradeSourceCompletion { + /// Returns the exact machine-contract spelling. + #[must_use] + pub const fn code(self) -> &'static str { + match self { + Self::Complete => "complete", + Self::IncompleteTimeout => "incomplete_timeout", + Self::IncompleteUnavailable => "incomplete_unavailable", + Self::IncompleteResourceLimit => "incomplete_resource_limit", + Self::IncompleteUnknown => "incomplete_unknown", + Self::Unsupported => "unsupported", + } + } + + #[must_use] + const fn allows_checkpoint(self) -> bool { + matches!(self, Self::Complete) + } +} + +/// Monotonic per-trade invalidation generation. +#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)] +pub struct RhiTradeDirtyGeneration(u64); + +impl RhiTradeDirtyGeneration { + /// Returns the positive durable generation. + #[must_use] + pub const fn get(self) -> u64 { + self.0 + } +} + +/// Exact canonically admitted cursor tuple for one source scope. +#[derive(Clone, Copy, PartialEq, Eq)] +pub struct RhiTradeSourceCursor { + created_at_unix_seconds: u64, + event_id: [u8; 32], +} + +impl RhiTradeSourceCursor { + /// Returns the inclusive event-authored UTC second. + #[must_use] + pub const fn created_at_unix_seconds(self) -> u64 { + self.created_at_unix_seconds + } + + /// Returns the exact verified Nostr event identifier bytes. + #[must_use] + pub const fn event_id(self) -> [u8; 32] { + self.event_id + } +} + +impl fmt::Debug for RhiTradeSourceCursor { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("RhiTradeSourceCursor") + .field("created_at_unix_seconds", &self.created_at_unix_seconds) + .field("event_id", &"[redacted]") + .finish() + } +} + +/// Caller-owned, injected timing and request evidence for one fetch attempt. +pub struct RhiTradeSourceAttempt { + request_id: Box<str>, + attempt_started_at: UnixTimeSeconds, + observed_at: RhiTradeMutationObservedAtUnixSeconds, + authored_time_policy: RhiTradeMutationAuthoredTimePolicy, +} + +impl RhiTradeSourceAttempt { + /// Validates a bounded request identity and explicit, ordered timestamps. + pub fn new( + request_id: impl AsRef<str>, + attempt_started_at: UnixTimeSeconds, + observed_at: RhiTradeMutationObservedAtUnixSeconds, + authored_time_policy: RhiTradeMutationAuthoredTimePolicy, + ) -> Result<Self, RhiTradeSourceIngestError> { + let request_id = request_id.as_ref(); + if request_id.is_empty() + || request_id.len() > FETCH_REQUEST_ID_MAX_BYTES + || request_id != request_id.trim() + || request_id.chars().any(char::is_control) + || attempt_started_at.get() == 0 + || i64::try_from(attempt_started_at.get()).is_err() + || observed_at.get() < attempt_started_at.get() + { + return Err(failure(RhiTradeSourceIngestErrorKind::InvalidInput)); + } + Ok(Self { + request_id: request_id.into(), + attempt_started_at, + observed_at, + authored_time_policy, + }) + } +} + +impl fmt::Debug for RhiTradeSourceAttempt { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("RhiTradeSourceAttempt") + .field("request_id", &"[redacted]") + .field("attempt_started_at", &self.attempt_started_at.get()) + .field("observed_at", &self.observed_at.get()) + .field("authored_time_policy", &self.authored_time_policy) + .finish() + } +} + +/// Stable source-free relay-ingest failure classification. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum RhiTradeSourceIngestErrorKind { + InvalidMode, + InvalidInput, + InvalidConfiguration, + GenerationConflict, + Storage, + CommitOutcomeUnknown, +} + +impl RhiTradeSourceIngestErrorKind { + /// Returns the stable machine-readable code. + #[must_use] + pub const fn code(self) -> &'static str { + match self { + Self::InvalidMode => "trade_source_mode_invalid", + Self::InvalidInput => "trade_source_input_invalid", + Self::InvalidConfiguration => "trade_source_configuration_invalid", + Self::GenerationConflict => "trade_source_generation_conflict", + Self::Storage => "trade_source_storage_failed", + Self::CommitOutcomeUnknown => "trade_source_commit_outcome_unknown", + } + } +} + +/// Redacted source-free relay-ingest failure. +#[derive(Clone, Copy, PartialEq, Eq)] +pub struct RhiTradeSourceIngestError { + kind: RhiTradeSourceIngestErrorKind, +} + +impl RhiTradeSourceIngestError { + /// Returns the stable failure kind. + #[must_use] + pub const fn kind(self) -> RhiTradeSourceIngestErrorKind { + self.kind + } + + /// Returns the stable machine-readable code. + #[must_use] + pub const fn code(self) -> &'static str { + self.kind.code() + } +} + +impl fmt::Debug for RhiTradeSourceIngestError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("RhiTradeSourceIngestError") + .field("kind", &self.kind) + .finish() + } +} + +impl fmt::Display for RhiTradeSourceIngestError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str(match self.kind { + RhiTradeSourceIngestErrorKind::InvalidMode => { + "RHI trade-source ingestion requires writable state" + } + RhiTradeSourceIngestErrorKind::InvalidInput => { + "RHI trade-source attempt input is invalid" + } + RhiTradeSourceIngestErrorKind::InvalidConfiguration => { + "RHI trade-source configuration is invalid" + } + RhiTradeSourceIngestErrorKind::GenerationConflict => { + "RHI trade-source generation changed during the attempt" + } + RhiTradeSourceIngestErrorKind::Storage => "RHI trade-source state transaction failed", + RhiTradeSourceIngestErrorKind::CommitOutcomeUnknown => { + "RHI trade-source commit outcome is unknown" + } + }) + } +} + +impl Error for RhiTradeSourceIngestError {} + +/// Durable outcome of one bounded source fetch and atomic evidence commit. +#[derive(Clone, Copy, PartialEq, Eq)] +pub struct RhiTradeSourceIngestOutcome { + completion: RhiTradeSourceCompletion, + received_events: u32, + admitted_events: u32, + rejected_events: u32, + duplicate_events: u32, + inserted_mutations: u32, + inserted_signed_events: u32, + inserted_observations: u32, + checkpoint: Option<RhiTradeSourceCursor>, + checkpoint_advanced: bool, + dirty_generation: Option<RhiTradeDirtyGeneration>, + dirty_generation_advanced: bool, +} + +impl RhiTradeSourceIngestOutcome { + #[must_use] + pub const fn completion(self) -> RhiTradeSourceCompletion { + self.completion + } + + #[must_use] + pub const fn received_events(self) -> u32 { + self.received_events + } + + #[must_use] + pub const fn admitted_events(self) -> u32 { + self.admitted_events + } + + #[must_use] + pub const fn rejected_events(self) -> u32 { + self.rejected_events + } + + #[must_use] + pub const fn duplicate_events(self) -> u32 { + self.duplicate_events + } + + #[must_use] + pub const fn inserted_mutations(self) -> u32 { + self.inserted_mutations + } + + #[must_use] + pub const fn inserted_signed_events(self) -> u32 { + self.inserted_signed_events + } + + #[must_use] + pub const fn inserted_observations(self) -> u32 { + self.inserted_observations + } + + #[must_use] + pub const fn checkpoint(self) -> Option<RhiTradeSourceCursor> { + self.checkpoint + } + + #[must_use] + pub const fn checkpoint_advanced(self) -> bool { + self.checkpoint_advanced + } + + #[must_use] + pub const fn dirty_generation(self) -> Option<RhiTradeDirtyGeneration> { + self.dirty_generation + } + + #[must_use] + pub const fn dirty_generation_advanced(self) -> bool { + self.dirty_generation_advanced + } +} + +impl fmt::Debug for RhiTradeSourceIngestOutcome { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("RhiTradeSourceIngestOutcome") + .field("completion", &self.completion) + .field("received_events", &self.received_events) + .field("admitted_events", &self.admitted_events) + .field("rejected_events", &self.rejected_events) + .field("duplicate_events", &self.duplicate_events) + .field("inserted_mutations", &self.inserted_mutations) + .field("inserted_signed_events", &self.inserted_signed_events) + .field("inserted_observations", &self.inserted_observations) + .field("checkpoint", &self.checkpoint) + .field("checkpoint_advanced", &self.checkpoint_advanced) + .field("dirty_generation", &self.dirty_generation) + .field("dirty_generation_advanced", &self.dirty_generation_advanced) + .finish() + } +} + +/// Fetches one exact configured relay source and atomically commits admitted evidence. +pub async fn ingest_rhi_trade_source( + repositories: &RhiStateRepositories<'_>, + transports: &RhiTransportAdapters, + configuration: &RhiConfigDocumentV1, + source_id: &str, + trade_id: TradeId, + attempt: RhiTradeSourceAttempt, +) -> Result<RhiTradeSourceIngestOutcome, RhiTradeSourceIngestError> { + let host = repositories.host(); + if host.mode() != RhiStateHostMode::ReadWriteExisting { + return Err(failure(RhiTradeSourceIngestErrorKind::InvalidMode)); + } + let source = ConfiguredSource::new(host, configuration, source_id)?; + let initial = read_initial_state(repositories, &source, trade_id).await?; + let fetched = fetch_source(transports, &source, trade_id, &attempt, initial.checkpoint).await?; + commit_source_result(repositories, source, trade_id, attempt, initial, fetched).await +} + +struct ConfiguredSource { + source_id: Box<str>, + relay_url: Box<str>, + policy: RhiEvidencePolicyDigest, + deadline_ms: u64, + lookback_seconds: u64, + overlap_seconds: u64, + maximum_events: usize, + maximum_bytes: usize, + admission_limits: RhiTradeMutationAdmissionLimits, +} + +impl ConfiguredSource { + fn new( + host: &crate::RhiStateHost, + configuration: &RhiConfigDocumentV1, + source_id: &str, + ) -> Result<Self, RhiTradeSourceIngestError> { + let normalized = configuration.normalized(); + let source = configured_source(normalized, source_id) + .ok_or_else(|| failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration))?; + 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(RhiTradeSourceIngestErrorKind::InvalidConfiguration)); + } + let relay_id = source + .pointer("/relay_id") + .and_then(Value::as_str) + .ok_or_else(|| failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration))?; + let relay = normalized + .pointer("/relays") + .and_then(Value::as_array) + .and_then(|relays| { + relays + .iter() + .find(|relay| relay.pointer("/id").and_then(Value::as_str) == Some(relay_id)) + }) + .filter(|relay| relay.pointer("/read").and_then(Value::as_bool) == Some(true)) + .ok_or_else(|| failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration))?; + let relay_url = relay + .pointer("/url") + .and_then(Value::as_str) + .ok_or_else(|| failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration))?; + let policy = state_metadata::evidence_policy_digest(normalized) + .map_err(|_| failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration))?; + if policy != host.metadata().evidence_policy_digest() { + return Err(failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration)); + } + let deadline_ms = exact_u64(source, "/deadline_ms", 100, 30_000)?; + let lookback_seconds = exact_u64(source, "/lookback_seconds", 60, 2_678_400)?; + let overlap_seconds = exact_u64(source, "/overlap_seconds", 1, 86_400)?; + if overlap_seconds > lookback_seconds { + return Err(failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration)); + } + let maximum_events = exact_usize( + normalized, + "/resource_limits/source_results/events", + 1, + RHI_TRADE_SOURCE_RESULT_MAX_EVENTS, + )?; + let maximum_bytes = exact_usize( + normalized, + "/resource_limits/source_results/bytes", + 1, + RHI_TRADE_SOURCE_RESULT_MAX_BYTES, + )?; + Target::nostr_relay(relay_url) + .map_err(|_| failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration))?; + let admission_limits = RhiTradeMutationAdmissionLimits::from_config(configuration) + .map_err(|_| failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration))?; + Ok(Self { + source_id: source_id.into(), + relay_url: relay_url.into(), + policy, + deadline_ms, + lookback_seconds, + overlap_seconds, + maximum_events, + maximum_bytes, + admission_limits, + }) + } +} + +#[derive(Clone, Copy, PartialEq, Eq)] +struct Checkpoint { + cursor: RhiTradeSourceCursor, + revision: u64, + completed_at_unix_s: u64, +} + +#[derive(Clone, Copy, PartialEq, Eq)] +pub(crate) struct DirtyState { + pub(crate) generation: RhiTradeDirtyGeneration, + pub(crate) policy: RhiEvidencePolicyDigest, + pub(crate) updated_at_unix_s: u64, +} + +#[derive(Clone, Copy)] +struct InitialState { + checkpoint: Option<Checkpoint>, + dirty: Option<DirtyState>, +} + +struct Candidate { + event: SignedEvent, +} + +struct FetchedSource { + completion: RhiTradeSourceCompletion, + received_events: usize, + duplicate_events: usize, + admitted: Vec<RhiAdmittedTradeMutationEvent>, + rejected_events: usize, + cursor_candidate: Option<RhiTradeSourceCursor>, +} + +async fn fetch_source( + transports: &RhiTransportAdapters, + source: &ConfiguredSource, + trade_id: TradeId, + attempt: &RhiTradeSourceAttempt, + checkpoint: Option<Checkpoint>, +) -> Result<FetchedSource, RhiTradeSourceIngestError> { + let deadline_unix_ms = attempt + .attempt_started_at + .get() + .checked_mul(1_000) + .and_then(|value| value.checked_add(source.deadline_ms)) + .ok_or_else(|| failure(RhiTradeSourceIngestErrorKind::InvalidInput))?; + let since = checkpoint.map_or_else( + || { + attempt + .attempt_started_at + .get() + .saturating_sub(source.lookback_seconds) + }, + |value| { + value + .cursor + .created_at_unix_seconds + .saturating_sub(source.overlap_seconds) + }, + ); + let selector = FetchSelector::all() + .with_kinds(EVENT_KINDS.to_vec()) + .and_then(|selector| selector.with_exact_tag_value('d', trade_id.to_hex())) + .and_then(|selector| selector.with_since_unix_seconds(since)) + .map_err(|_| failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration))?; + let target = Target::nostr_relay(source.relay_url.as_ref()) + .map_err(|_| failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration))?; + let fingerprint = target.fingerprint().clone(); + let targets = TargetSet::new(vec![target]) + .map_err(|_| failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration))?; + let mut adapter_cursor = None::<FetchCursor>; + let mut seen_adapter_cursors = BTreeSet::new(); + let mut candidates = Vec::new(); + let mut received_events = 0_usize; + let mut received_bytes = 0_usize; + let completion = loop { + let remaining = source.maximum_events.saturating_sub(received_events); + let request_limit = usize::min(usize::from(FETCH_PAGE_MAX_EVENTS), remaining + 1); + let bounds = FetchBounds::new( + u16::try_from(request_limit) + .map_err(|_| failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration))?, + deadline_unix_ms, + ) + .map_err(|_| failure(RhiTradeSourceIngestErrorKind::InvalidInput))?; + let mut request = FetchRequest::new(attempt.request_id.as_ref(), targets.clone(), bounds) + .map_err(|_| failure(RhiTradeSourceIngestErrorKind::InvalidInput))? + .with_selector(selector.clone()); + if let Some(cursor) = adapter_cursor.take() { + request = request.with_cursor(cursor); + } + let page = match transports.evidence_source().fetch(request.clone()).await { + Ok(page) => page, + Err(radroots_transport::Error::UnsupportedOperation) => { + break RhiTradeSourceCompletion::Unsupported; + } + Err(_) => break RhiTradeSourceCompletion::IncompleteUnknown, + }; + if page.validate_for_request(&request).is_err() { + break RhiTradeSourceCompletion::IncompleteUnknown; + } + let outcome = page + .target_outcomes() + .iter() + .find(|outcome| outcome.target() == &fingerprint); + let Some(outcome) = outcome.filter(|_| page.target_outcomes().len() == 1) else { + break RhiTradeSourceCompletion::IncompleteUnknown; + }; + let target_state = outcome.state(); + match target_state { + FetchTargetState::Complete | FetchTargetState::Partial => {} + FetchTargetState::Unavailable | FetchTargetState::FailedRetryable => { + break RhiTradeSourceCompletion::IncompleteUnavailable; + } + FetchTargetState::FailedTerminal => { + break RhiTradeSourceCompletion::IncompleteUnknown; + } + FetchTargetState::Cancelled => break RhiTradeSourceCompletion::IncompleteTimeout, + } + for observed in page.events() { + received_events = received_events.saturating_add(1); + if received_events > source.maximum_events { + break; + } + received_bytes = match received_bytes.checked_add(observed.event().raw_json().len()) { + Some(value) if value <= source.maximum_bytes => value, + _ => { + received_events = source.maximum_events.saturating_add(1); + break; + } + }; + candidates.push(Candidate { + event: observed.event().clone(), + }); + } + if received_events > source.maximum_events { + break RhiTradeSourceCompletion::IncompleteResourceLimit; + } + if target_state == FetchTargetState::Partial { + break RhiTradeSourceCompletion::IncompleteUnknown; + } + match page.next_page() { + NextPage::Complete => break RhiTradeSourceCompletion::Complete, + NextPage::Cancelled { .. } => break RhiTradeSourceCompletion::IncompleteTimeout, + NextPage::Cursor(cursor) => { + if page.events().is_empty() + || !seen_adapter_cursors.insert(cursor.as_str().to_owned()) + { + break RhiTradeSourceCompletion::IncompleteUnknown; + } + adapter_cursor = Some(cursor.clone()); + } + } + }; + + candidates.sort_by(compare_candidate); + let mut duplicate_events = 0_usize; + let mut admitted_event_ids = BTreeSet::new(); + let mut admitted = Vec::with_capacity(candidates.len()); + let mut rejected_events = 0_usize; + let mut cursor_candidate = None; + for candidate in candidates { + match admit_rhi_trade_mutation_event( + source.admission_limits, + candidate.event.raw_json().as_bytes(), + attempt.observed_at, + attempt.authored_time_policy, + ) { + Ok(event) if event.mutation().trade_id == trade_id => { + if !admitted_event_ids.insert(*event.event_id().as_bytes()) { + duplicate_events = duplicate_events.saturating_add(1); + continue; + } + let cursor = RhiTradeSourceCursor { + created_at_unix_seconds: event.authored_at_unix_seconds(), + event_id: *event.event_id().as_bytes(), + }; + cursor_candidate = Some(cursor_candidate.map_or(cursor, |current| { + if compare_cursor(current, cursor).is_lt() { + cursor + } else { + current + } + })); + admitted.push(event); + } + Ok(_) | Err(_) => rejected_events = rejected_events.saturating_add(1), + } + } + Ok(FetchedSource { + completion, + received_events: received_events.min(source.maximum_events), + duplicate_events, + admitted, + rejected_events, + cursor_candidate, + }) +} + +async fn read_initial_state( + repositories: &RhiStateRepositories<'_>, + source: &ConfiguredSource, + trade_id: TradeId, +) -> Result<InitialState, RhiTradeSourceIngestError> { + let source_id = source.source_id.clone(); + let policy = source.policy; + repositories + .host() + .sqlite_host() + .transaction(move |transaction| { + Box::pin(async move { + Ok(InitialState { + checkpoint: read_checkpoint(transaction, source_id.as_ref(), policy, trade_id) + .await?, + dirty: read_dirty(transaction, trade_id).await?, + }) + }) + }) + .await + .map_err(map_transaction_error) +} + +async fn commit_source_result( + repositories: &RhiStateRepositories<'_>, + source: ConfiguredSource, + trade_id: TradeId, + attempt: RhiTradeSourceAttempt, + initial: InitialState, + fetched: FetchedSource, +) -> Result<RhiTradeSourceIngestOutcome, RhiTradeSourceIngestError> { + let source_id = source.source_id; + let policy = source.policy; + repositories + .host() + .sqlite_host() + .transaction(move |transaction| { + Box::pin(async move { + if read_checkpoint(transaction, source_id.as_ref(), policy, trade_id).await? + != initial.checkpoint + || read_dirty(transaction, trade_id).await? != initial.dirty + { + return Err(SourceOperationError::GenerationConflict); + } + let mut inserted_mutations = 0_u32; + let mut inserted_signed_events = 0_u32; + let mut inserted_observations = 0_u32; + let admitted_events = u32::try_from(fetched.admitted.len()) + .map_err(|_| SourceOperationError::Storage)?; + for admitted in fetched.admitted { + let observation = + RhiTradeSourceObservation::from_parts(source_id.clone(), policy, &admitted); + let record = PersistenceRecord::from_admitted(admitted) + .map_err(|_| SourceOperationError::Storage)?; + let persisted = persist(transaction, &record, &observation) + .await + .map_err(SourceOperationError::Persistence)?; + inserted_mutations = inserted_mutations + .checked_add(u32::from(persisted.mutation_inserted())) + .ok_or(SourceOperationError::Storage)?; + inserted_signed_events = inserted_signed_events + .checked_add(u32::from(persisted.signed_event_inserted())) + .ok_or(SourceOperationError::Storage)?; + inserted_observations = inserted_observations + .checked_add(u32::from(persisted.observation_inserted())) + .ok_or(SourceOperationError::Storage)?; + } + let new_relevant_evidence = inserted_mutations != 0 || inserted_signed_events != 0; + let (dirty_generation, dirty_generation_advanced) = if new_relevant_evidence { + let generation = advance_dirty_generation( + transaction, + trade_id, + policy, + attempt.observed_at.get(), + initial.dirty, + ) + .await?; + (Some(generation), true) + } else { + (initial.dirty.map(|dirty| dirty.generation), false) + }; + let mut checkpoint = initial.checkpoint.map(|value| value.cursor); + let mut checkpoint_advanced = false; + if fetched.completion.allows_checkpoint() + && fetched.cursor_candidate.is_some_and(|candidate| { + initial + .checkpoint + .is_none_or(|current| compare_cursor(current.cursor, candidate).is_lt()) + }) + { + let Some(candidate) = fetched.cursor_candidate else { + return Err(SourceOperationError::Storage); + }; + write_checkpoint( + transaction, + source_id.as_ref(), + policy, + trade_id, + initial.checkpoint, + candidate, + attempt.observed_at.get(), + ) + .await?; + checkpoint = Some(candidate); + checkpoint_advanced = true; + } + Ok(RhiTradeSourceIngestOutcome { + completion: fetched.completion, + received_events: u32::try_from(fetched.received_events) + .map_err(|_| SourceOperationError::Storage)?, + admitted_events, + rejected_events: u32::try_from(fetched.rejected_events) + .map_err(|_| SourceOperationError::Storage)?, + duplicate_events: u32::try_from(fetched.duplicate_events) + .map_err(|_| SourceOperationError::Storage)?, + inserted_mutations, + inserted_signed_events, + inserted_observations, + checkpoint, + checkpoint_advanced, + dirty_generation, + dirty_generation_advanced, + }) + }) + }) + .await + .map_err(map_transaction_error) +} + +pub(crate) async fn advance_dirty_generation( + transaction: &mut ServiceSqliteTransaction<'_>, + trade_id: TradeId, + policy: RhiEvidencePolicyDigest, + updated_at_unix_s: u64, + expected: Option<DirtyState>, +) -> Result<RhiTradeDirtyGeneration, SourceOperationError> { + let updated_at = i64::try_from(updated_at_unix_s).map_err(|_| SourceOperationError::Storage)?; + match expected { + None => { + let result = sqlx::query(INSERT_DIRTY_SQL) + .bind(trade_id.as_bytes().as_slice()) + .bind(policy.as_bytes().as_slice()) + .bind(updated_at) + .execute(&mut *transaction) + .await + .map_err(|_| SourceOperationError::Storage)?; + if result.rows_affected() != 1 { + return Err(SourceOperationError::GenerationConflict); + } + Ok(RhiTradeDirtyGeneration(1)) + } + Some(current) + if current.policy == policy && updated_at_unix_s >= current.updated_at_unix_s => + { + let next = current + .generation + .get() + .checked_add(1) + .filter(|value| i64::try_from(*value).is_ok()) + .ok_or(SourceOperationError::Storage)?; + let result = sqlx::query(UPDATE_DIRTY_SQL) + .bind(policy.as_bytes().as_slice()) + .bind(updated_at) + .bind(trade_id.as_bytes().as_slice()) + .bind( + i64::try_from(current.generation.get()) + .map_err(|_| SourceOperationError::Storage)?, + ) + .execute(&mut *transaction) + .await + .map_err(|_| SourceOperationError::Storage)?; + if result.rows_affected() != 1 { + return Err(SourceOperationError::GenerationConflict); + } + Ok(RhiTradeDirtyGeneration(next)) + } + Some(_) => Err(SourceOperationError::GenerationConflict), + } +} + +pub(crate) async fn read_dirty( + transaction: &mut ServiceSqliteTransaction<'_>, + trade_id: TradeId, +) -> Result<Option<DirtyState>, SourceOperationError> { + let Some(row) = sqlx::query(READ_DIRTY_SQL) + .bind(trade_id.as_bytes().as_slice()) + .fetch_optional(&mut *transaction) + .await + .map_err(|_| SourceOperationError::Storage)? + else { + return Ok(None); + }; + let generation = positive_i64_u64(&row, "generation")?; + let policy = exact_digest(&row, "evidence_policy_sha256", "evidence_policy_bytes")?; + let updated_at_unix_s = nonnegative_i64_u64(&row, "updated_at_unix_s")?; + Ok(Some(DirtyState { + generation: RhiTradeDirtyGeneration(generation), + policy: RhiEvidencePolicyDigest::from_bytes(policy), + updated_at_unix_s, + })) +} + +async fn read_checkpoint( + transaction: &mut ServiceSqliteTransaction<'_>, + source_id: &str, + policy: RhiEvidencePolicyDigest, + trade_id: TradeId, +) -> Result<Option<Checkpoint>, SourceOperationError> { + let Some(row) = sqlx::query(READ_CHECKPOINT_SQL) + .bind(source_id) + .bind(SOURCE_SELECTOR) + .bind(policy.as_bytes().as_slice()) + .bind(trade_id.as_bytes().as_slice()) + .fetch_optional(&mut *transaction) + .await + .map_err(|_| SourceOperationError::Storage)? + else { + return Ok(None); + }; + Ok(Some(Checkpoint { + cursor: RhiTradeSourceCursor { + created_at_unix_seconds: nonnegative_i64_u64(&row, "cursor_created_at_unix_s")?, + event_id: exact_digest(&row, "cursor_event_id", "cursor_event_id_bytes")?, + }, + revision: positive_i64_u64(&row, "revision")?, + completed_at_unix_s: positive_i64_u64(&row, "completed_at_unix_s")?, + })) +} + +async fn write_checkpoint( + transaction: &mut ServiceSqliteTransaction<'_>, + source_id: &str, + policy: RhiEvidencePolicyDigest, + trade_id: TradeId, + current: Option<Checkpoint>, + next: RhiTradeSourceCursor, + completed_at_unix_s: u64, +) -> Result<(), SourceOperationError> { + let created_at = + i64::try_from(next.created_at_unix_seconds).map_err(|_| SourceOperationError::Storage)?; + let completed_at = + i64::try_from(completed_at_unix_s).map_err(|_| SourceOperationError::Storage)?; + let result = match current { + None => { + sqlx::query(INSERT_CHECKPOINT_SQL) + .bind(source_id) + .bind(SOURCE_SELECTOR) + .bind(policy.as_bytes().as_slice()) + .bind(trade_id.as_bytes().as_slice()) + .bind(created_at) + .bind(next.event_id.as_slice()) + .bind(completed_at) + .execute(&mut *transaction) + .await + } + Some(current) => { + sqlx::query(UPDATE_CHECKPOINT_SQL) + .bind(created_at) + .bind(next.event_id.as_slice()) + .bind(completed_at) + .bind(source_id) + .bind(SOURCE_SELECTOR) + .bind(policy.as_bytes().as_slice()) + .bind(trade_id.as_bytes().as_slice()) + .bind(i64::try_from(current.revision).map_err(|_| SourceOperationError::Storage)?) + .execute(&mut *transaction) + .await + } + } + .map_err(|_| SourceOperationError::Storage)?; + if result.rows_affected() == 1 { + Ok(()) + } else { + Err(SourceOperationError::GenerationConflict) + } +} + +fn compare_candidate(left: &Candidate, right: &Candidate) -> Ordering { + left.event + .created_at() + .cmp(&right.event.created_at()) + .then_with(|| left.event.id().as_bytes().cmp(right.event.id().as_bytes())) + .then_with(|| { + left.event + .sig() + .as_bytes() + .cmp(right.event.sig().as_bytes()) + }) +} + +fn compare_cursor(left: RhiTradeSourceCursor, right: RhiTradeSourceCursor) -> Ordering { + (left.created_at_unix_seconds, left.event_id) + .cmp(&(right.created_at_unix_seconds, right.event_id)) +} + +fn configured_source<'a>(configuration: &'a Value, source_id: &str) -> Option<&'a Value> { + if source_id.is_empty() + || source_id.len() > 64 + || !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_u64( + value: &Value, + pointer: &str, + minimum: u64, + maximum: u64, +) -> Result<u64, RhiTradeSourceIngestError> { + value + .pointer(pointer) + .and_then(Value::as_u64) + .filter(|value| (minimum..=maximum).contains(value)) + .ok_or_else(|| failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration)) +} + +fn exact_usize( + value: &Value, + pointer: &str, + minimum: usize, + maximum: usize, +) -> Result<usize, RhiTradeSourceIngestError> { + exact_u64( + value, + pointer, + u64::try_from(minimum).unwrap_or(u64::MAX), + u64::try_from(maximum).unwrap_or(u64::MAX), + ) + .and_then(|value| { + usize::try_from(value) + .map_err(|_| failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration)) + }) +} + +fn exact_digest( + row: &sqlx::sqlite::SqliteRow, + field: &str, + length_field: &str, +) -> Result<[u8; 32], SourceOperationError> { + if row.try_get::<i64, _>(length_field).ok() != Some(32) { + return Err(SourceOperationError::Storage); + } + row.try_get::<Vec<u8>, _>(field) + .map_err(|_| SourceOperationError::Storage)? + .try_into() + .map_err(|_| SourceOperationError::Storage) +} + +fn positive_i64_u64( + row: &sqlx::sqlite::SqliteRow, + field: &str, +) -> Result<u64, SourceOperationError> { + row.try_get::<i64, _>(field) + .ok() + .filter(|value| *value > 0) + .and_then(|value| u64::try_from(value).ok()) + .ok_or(SourceOperationError::Storage) +} + +fn nonnegative_i64_u64( + row: &sqlx::sqlite::SqliteRow, + field: &str, +) -> Result<u64, SourceOperationError> { + row.try_get::<i64, _>(field) + .ok() + .filter(|value| *value >= 0) + .and_then(|value| u64::try_from(value).ok()) + .ok_or(SourceOperationError::Storage) +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(crate) enum SourceOperationError { + GenerationConflict, + Persistence(PersistenceOperationError), + Storage, +} + +fn map_transaction_error( + error: ServiceSqliteTransactionError<SourceOperationError>, +) -> RhiTradeSourceIngestError { + if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown { + return failure(RhiTradeSourceIngestErrorKind::CommitOutcomeUnknown); + } + failure(match error.operation_error().copied() { + Some(SourceOperationError::GenerationConflict) => { + RhiTradeSourceIngestErrorKind::GenerationConflict + } + Some(SourceOperationError::Persistence(_)) | Some(SourceOperationError::Storage) | None => { + RhiTradeSourceIngestErrorKind::Storage + } + }) +} + +const fn failure(kind: RhiTradeSourceIngestErrorKind) -> RhiTradeSourceIngestError { + RhiTradeSourceIngestError { kind } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn equal_timestamp_cursor_order_uses_verified_event_id() { + let lower = RhiTradeSourceCursor { + created_at_unix_seconds: 100, + event_id: [0x11; 32], + }; + let higher = RhiTradeSourceCursor { + created_at_unix_seconds: 100, + event_id: [0x22; 32], + }; + assert_eq!(compare_cursor(lower, higher), Ordering::Less); + assert_eq!(compare_cursor(higher, lower), Ordering::Greater); + assert_eq!(compare_cursor(lower, lower), Ordering::Equal); + } + + #[test] + fn completion_codes_and_checkpoint_policy_are_closed() { + let vectors = [ + (RhiTradeSourceCompletion::Complete, "complete", true), + ( + RhiTradeSourceCompletion::IncompleteTimeout, + "incomplete_timeout", + false, + ), + ( + RhiTradeSourceCompletion::IncompleteUnavailable, + "incomplete_unavailable", + false, + ), + ( + RhiTradeSourceCompletion::IncompleteResourceLimit, + "incomplete_resource_limit", + false, + ), + ( + RhiTradeSourceCompletion::IncompleteUnknown, + "incomplete_unknown", + false, + ), + (RhiTradeSourceCompletion::Unsupported, "unsupported", false), + ]; + for (completion, code, checkpoint) in vectors { + assert_eq!(completion.code(), code); + assert_eq!(completion.allows_checkpoint(), checkpoint); + } + } +} 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 = 3; +pub const RHI_STATE_SCHEMA_VERSION: u32 = 4; /// The shared metadata and migration-ledger objects present at schema v1. pub const RHI_STATE_SCHEMA_VERSION_1_OBJECT_COUNT: u32 = 6; @@ -23,10 +23,13 @@ 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; +/// The shared objects, immutable trade evidence, source cursors, and dirty generations. +pub const RHI_STATE_SCHEMA_VERSION_4_OBJECT_COUNT: u32 = 28; + /// SHA-256 identity of the ordered migration catalog rooted at schema v1. pub const RHI_MIGRATION_CATALOG_SHA256: [u8; 32] = [ - 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, + 0x29, 0x5f, 0xf5, 0xe9, 0xac, 0x0e, 0xad, 0xc1, 0xf2, 0xc2, 0xd8, 0xc2, 0x58, 0xc2, 0xea, 0xa8, + 0x16, 0xf9, 0x46, 0x56, 0x0b, 0x68, 0x9d, 0xcb, 0x43, 0x2f, 0xdc, 0x99, 0xe7, 0x9e, 0xc3, 0x53, ]; /// SHA-256 identity of the exact schema-v1 object snapshot. @@ -59,10 +62,22 @@ pub const RHI_STATE_SCHEMA_VERSION_3_SHA256: [u8; 32] = [ 0x7a, 0xab, 0x2d, 0xbd, 0x23, 0xfe, 0xad, 0xac, 0x16, 0x53, 0x33, 0x49, 0x0d, 0x6f, 0x0a, 0xd5, ]; +/// SHA-256 identity of the schema-v4 source-checkpoint migration. +pub const RHI_STATE_SCHEMA_VERSION_4_MIGRATION_SHA256: [u8; 32] = [ + 0x24, 0x41, 0xf7, 0xc4, 0xc1, 0x15, 0x94, 0xc6, 0xfd, 0xe7, 0x87, 0xdb, 0xb9, 0x4e, 0x44, 0xba, + 0x00, 0xb7, 0x2d, 0x2c, 0x97, 0x9b, 0x0e, 0x7f, 0xd3, 0xc7, 0x33, 0xea, 0xe9, 0x14, 0x4d, 0x78, +]; + +/// SHA-256 identity of the exact schema-v4 object snapshot. +pub const RHI_STATE_SCHEMA_VERSION_4_SHA256: [u8; 32] = [ + 0x9c, 0xa7, 0x8a, 0x54, 0xb0, 0xea, 0x20, 0x13, 0xe7, 0xaa, 0x70, 0xd0, 0x9b, 0xea, 0xdc, 0xfb, + 0x6c, 0xdd, 0xe0, 0x7c, 0x40, 0x42, 0x9d, 0xff, 0xe9, 0x1c, 0xb8, 0x3c, 0x2a, 0x42, 0x5c, 0xea, +]; + /// SHA-256 identity of the schema catalog bound to the migration catalog. pub const RHI_STATE_SCHEMA_CATALOG_SHA256: [u8; 32] = [ - 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, + 0x58, 0x78, 0x9c, 0x1d, 0x39, 0x65, 0x14, 0xce, 0x19, 0xa5, 0x60, 0xb6, 0x7b, 0xf7, 0x60, 0x65, + 0x0b, 0x9d, 0x6f, 0xa3, 0x78, 0xbc, 0x82, 0x34, 0xc1, 0xf6, 0xa3, 0x9a, 0xd2, 0xf5, 0x56, 0xe6, ]; macro_rules! rhi_config_bindings_table_sql { @@ -351,6 +366,112 @@ const CREATE_TRADE_EVIDENCE_MIGRATION_SQL: &str = concat!( ";", ); +macro_rules! relay_checkpoints_table_sql { + () => { + r#"CREATE TABLE relay_checkpoints ( + 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), + trade_id BLOB NOT NULL CHECK (length(trade_id) = 16), + cursor_created_at_unix_s INTEGER NOT NULL + CHECK (cursor_created_at_unix_s BETWEEN 0 AND 9223372036854775807), + cursor_event_id BLOB NOT NULL CHECK (length(cursor_event_id) = 32), + revision INTEGER NOT NULL CHECK (revision BETWEEN 1 AND 9223372036854775807), + completed_at_unix_s INTEGER NOT NULL + CHECK (completed_at_unix_s BETWEEN 1 AND 9223372036854775807), + PRIMARY KEY (source_id, selector_id, evidence_policy_sha256, trade_id) +) STRICT"# + }; +} + +macro_rules! relay_checkpoints_guard_update_sql { + () => { + r#"CREATE TRIGGER relay_checkpoints_guard_update +BEFORE UPDATE ON relay_checkpoints +WHEN NEW.source_id != OLD.source_id + OR NEW.selector_id != OLD.selector_id + OR NEW.evidence_policy_sha256 != OLD.evidence_policy_sha256 + OR NEW.trade_id != OLD.trade_id + OR NEW.revision != OLD.revision + 1 + OR NEW.completed_at_unix_s < OLD.completed_at_unix_s + OR NEW.cursor_created_at_unix_s < OLD.cursor_created_at_unix_s + OR ( + NEW.cursor_created_at_unix_s = OLD.cursor_created_at_unix_s + AND NEW.cursor_event_id <= OLD.cursor_event_id + ) +BEGIN + SELECT RAISE(ABORT, 'relay checkpoint transition is invalid'); +END"# + }; +} + +macro_rules! trade_dirty_generations_table_sql { + () => { + r#"CREATE TABLE trade_dirty_generations ( + trade_id BLOB NOT NULL PRIMARY KEY CHECK (length(trade_id) = 16), + generation INTEGER NOT NULL CHECK (generation BETWEEN 1 AND 9223372036854775807), + evidence_policy_sha256 BLOB NOT NULL CHECK (length(evidence_policy_sha256) = 32), + updated_at_unix_s INTEGER NOT NULL + CHECK (updated_at_unix_s BETWEEN 0 AND 9223372036854775807) +) STRICT"# + }; +} + +macro_rules! trade_dirty_generations_guard_update_sql { + () => { + r#"CREATE TRIGGER trade_dirty_generations_guard_update +BEFORE UPDATE ON trade_dirty_generations +WHEN NEW.trade_id != OLD.trade_id + OR NEW.generation != OLD.generation + 1 + OR NEW.updated_at_unix_s < OLD.updated_at_unix_s +BEGIN + SELECT RAISE(ABORT, 'trade dirty generation transition is invalid'); +END"# + }; +} + +pub(crate) const CREATE_RELAY_CHECKPOINTS_TABLE_SQL: &str = relay_checkpoints_table_sql!(); +const CREATE_RELAY_CHECKPOINTS_GUARD_UPDATE_SQL: &str = relay_checkpoints_guard_update_sql!(); +const CREATE_RELAY_CHECKPOINTS_NO_DELETE_SQL: &str = immutable_no_delete_sql!( + "relay_checkpoints_no_delete", + "relay_checkpoints", + "relay checkpoints are retained" +); +pub(crate) const CREATE_TRADE_DIRTY_GENERATIONS_TABLE_SQL: &str = + trade_dirty_generations_table_sql!(); +const CREATE_TRADE_DIRTY_GENERATIONS_GUARD_UPDATE_SQL: &str = + trade_dirty_generations_guard_update_sql!(); +const CREATE_TRADE_DIRTY_GENERATIONS_NO_DELETE_SQL: &str = immutable_no_delete_sql!( + "trade_dirty_generations_no_delete", + "trade_dirty_generations", + "trade dirty generations are retained" +); +const CREATE_SOURCE_CHECKPOINT_MIGRATION_SQL: &str = concat!( + relay_checkpoints_table_sql!(), + ";\n", + relay_checkpoints_guard_update_sql!(), + ";\n", + immutable_no_delete_sql!( + "relay_checkpoints_no_delete", + "relay_checkpoints", + "relay checkpoints are retained" + ), + ";\n", + trade_dirty_generations_table_sql!(), + ";\n", + trade_dirty_generations_guard_update_sql!(), + ";\n", + immutable_no_delete_sql!( + "trade_dirty_generations_no_delete", + "trade_dirty_generations", + "trade dirty generations are 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, @@ -415,6 +536,30 @@ 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, ]; +const RELAY_CHECKPOINTS_TABLE_SHA256: [u8; 32] = [ + 0x5f, 0xf7, 0xc3, 0x41, 0xa3, 0x48, 0x37, 0x3e, 0x92, 0xf5, 0x46, 0xe1, 0xff, 0xfd, 0x45, 0x4a, + 0x86, 0x6e, 0x23, 0x2c, 0xff, 0x60, 0x83, 0xbf, 0x3f, 0xfe, 0xc3, 0xfc, 0x65, 0x53, 0x2c, 0xb6, +]; +const RELAY_CHECKPOINTS_GUARD_UPDATE_SHA256: [u8; 32] = [ + 0x5c, 0xd7, 0x2a, 0xaf, 0x31, 0x2c, 0x17, 0xcf, 0xfa, 0xe4, 0xa8, 0x37, 0x91, 0x3c, 0x5e, 0x92, + 0xf2, 0xaa, 0x31, 0xea, 0x71, 0x3e, 0x86, 0x0e, 0x3a, 0xaa, 0x9f, 0xf7, 0x05, 0xd0, 0xc4, 0x92, +]; +const RELAY_CHECKPOINTS_NO_DELETE_SHA256: [u8; 32] = [ + 0x8d, 0xb0, 0x54, 0xf7, 0x75, 0x97, 0xd5, 0xe5, 0x42, 0xd0, 0x0c, 0x2e, 0xdc, 0xf2, 0xd9, 0x26, + 0xa9, 0x2b, 0x8e, 0xf5, 0x5a, 0x69, 0x54, 0xb2, 0x01, 0x6b, 0x93, 0x36, 0x48, 0x84, 0xfb, 0x1f, +]; +const TRADE_DIRTY_GENERATIONS_TABLE_SHA256: [u8; 32] = [ + 0xe9, 0x93, 0x53, 0x2c, 0x08, 0x36, 0xa2, 0x99, 0x40, 0xd2, 0xe3, 0x51, 0x5b, 0x11, 0x34, 0xea, + 0x68, 0xf7, 0xd3, 0x50, 0xe5, 0xbe, 0xdb, 0x3b, 0xc1, 0xb1, 0xa5, 0x6f, 0x7f, 0xa6, 0xe4, 0xd1, +]; +const TRADE_DIRTY_GENERATIONS_GUARD_UPDATE_SHA256: [u8; 32] = [ + 0x50, 0x31, 0xbc, 0x10, 0x4d, 0x51, 0x66, 0xae, 0xf2, 0xdc, 0x9a, 0x1b, 0xaf, 0x42, 0x0e, 0x47, + 0xc1, 0x6b, 0xa9, 0x70, 0x91, 0x6f, 0x5d, 0xc5, 0x6d, 0xb3, 0xc1, 0x70, 0x4e, 0x76, 0x16, 0xcb, +]; +const TRADE_DIRTY_GENERATIONS_NO_DELETE_SHA256: [u8; 32] = [ + 0x36, 0x66, 0x3a, 0x1a, 0xd4, 0x65, 0x41, 0x9f, 0x23, 0x1e, 0xe6, 0xcd, 0x1e, 0xd1, 0x6d, 0xa4, + 0x26, 0xac, 0x7d, 0x70, 0xde, 0xc3, 0x67, 0x75, 0x55, 0x62, 0xa8, 0x45, 0x52, 0x86, 0x54, 0x23, +]; /// Stable classes for invalid embedded RHI catalog definitions. #[derive(Clone, Copy, Debug, PartialEq, Eq)] @@ -501,10 +646,17 @@ pub fn rhi_migration_catalog() -> Result<MigrationCatalog, RhiStateCatalogError> MigrationChecksum::from_bytes(RHI_STATE_SCHEMA_VERSION_3_MIGRATION_SHA256), ) .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::MigrationCatalog))?; - let catalog = MigrationCatalog::new([configuration, trade_evidence]) + let source_checkpoints = MigrationDescriptor::sql( + 4, + "create_source_checkpoints_and_dirty_generations", + CREATE_SOURCE_CHECKPOINT_MIGRATION_SQL, + MigrationChecksum::from_bytes(RHI_STATE_SCHEMA_VERSION_4_MIGRATION_SHA256), + ) + .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::MigrationCatalog))?; + let catalog = MigrationCatalog::new([configuration, trade_evidence, source_checkpoints]) .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::MigrationCatalog))?; if catalog.current_version() != RHI_STATE_SCHEMA_VERSION - || catalog.descriptors().len() != 2 + || catalog.descriptors().len() != 3 || catalog.digest().as_bytes() != &RHI_MIGRATION_CATALOG_SHA256 { return Err(RhiStateCatalogError::new( @@ -530,13 +682,22 @@ pub fn rhi_schema_catalog() -> Result<SchemaCatalog, RhiStateCatalogError> { ) .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog))?; let version_three = SchemaVersionCatalog::new( - RHI_STATE_SCHEMA_VERSION, + 3, 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))?; + let version_four = SchemaVersionCatalog::new( + RHI_STATE_SCHEMA_VERSION, + rhi_schema_version_four_objects()?, + SchemaDigest::from_bytes(RHI_STATE_SCHEMA_VERSION_4_SHA256), + ) + .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog))?; + let catalog = SchemaCatalog::new( + &migrations, + [version_one, version_two, version_three, version_four], + ) + .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog))?; validate_rhi_state_catalogs(&migrations, &catalog)?; Ok(catalog) } @@ -548,19 +709,22 @@ pub fn validate_rhi_state_catalogs( ) -> Result<(), RhiStateCatalogError> { let versions = schema.versions(); let valid = migrations.current_version() == RHI_STATE_SCHEMA_VERSION - && migrations.descriptors().len() == 2 + && migrations.descriptors().len() == 3 && migrations.digest().as_bytes() == &RHI_MIGRATION_CATALOG_SHA256 && schema.migration_catalog_digest() == migrations.digest() - && versions.len() == 3 + && versions.len() == 4 && 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() == 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].version() == 3 && versions[2].object_count() == RHI_STATE_SCHEMA_VERSION_3_OBJECT_COUNT && versions[2].digest().as_bytes() == &RHI_STATE_SCHEMA_VERSION_3_SHA256 + && versions[3].version() == RHI_STATE_SCHEMA_VERSION + && versions[3].object_count() == RHI_STATE_SCHEMA_VERSION_4_OBJECT_COUNT + && versions[3].digest().as_bytes() == &RHI_STATE_SCHEMA_VERSION_4_SHA256 && schema.digest().as_bytes() == &RHI_STATE_SCHEMA_CATALOG_SHA256; if valid { Ok(()) @@ -614,6 +778,69 @@ fn rhi_schema_version_three_objects() -> Result<Vec<SchemaObject>, RhiStateCatal Ok(objects) } +fn rhi_schema_version_four_objects() -> Result<Vec<SchemaObject>, RhiStateCatalogError> { + let mut objects = rhi_schema_version_three_objects()?; + objects.extend(rhi_source_checkpoint_objects()?); + Ok(objects) +} + +fn rhi_source_checkpoint_objects() -> Result<[SchemaObject; 6], 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, + "relay_checkpoints", + "relay_checkpoints", + CREATE_RELAY_CHECKPOINTS_TABLE_SQL, + RELAY_CHECKPOINTS_TABLE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "relay_checkpoints_guard_update", + "relay_checkpoints", + CREATE_RELAY_CHECKPOINTS_GUARD_UPDATE_SQL, + RELAY_CHECKPOINTS_GUARD_UPDATE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "relay_checkpoints_no_delete", + "relay_checkpoints", + CREATE_RELAY_CHECKPOINTS_NO_DELETE_SQL, + RELAY_CHECKPOINTS_NO_DELETE_SHA256, + )?, + object( + SchemaObjectKind::Table, + "trade_dirty_generations", + "trade_dirty_generations", + CREATE_TRADE_DIRTY_GENERATIONS_TABLE_SQL, + TRADE_DIRTY_GENERATIONS_TABLE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "trade_dirty_generations_guard_update", + "trade_dirty_generations", + CREATE_TRADE_DIRTY_GENERATIONS_GUARD_UPDATE_SQL, + TRADE_DIRTY_GENERATIONS_GUARD_UPDATE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "trade_dirty_generations_no_delete", + "trade_dirty_generations", + CREATE_TRADE_DIRTY_GENERATIONS_NO_DELETE_SQL, + TRADE_DIRTY_GENERATIONS_NO_DELETE_SHA256, + )?, + ]) +} + fn rhi_trade_evidence_objects() -> Result<[SchemaObject; 12], RhiStateCatalogError> { let object = |kind, name, table_name, sql, digest| { SchemaObject::new( diff --git a/src/state_config.rs b/src/state_config.rs @@ -43,6 +43,14 @@ const INSERT_BINDING_SQL: &str = r#"INSERT INTO rhi_config_bindings ( applied_at_unix_s, service_version, service_commit, lib_revision, rust_version, target, feature_profile ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)"#; +const READ_DIRTY_POLICY_BOUNDS_SQL: &str = r#"SELECT + COUNT(*) FILTER (WHERE generation >= 9223372036854775807) AS exhausted, + COALESCE(MAX(updated_at_unix_s), 0) AS latest_updated_at +FROM trade_dirty_generations +WHERE evidence_policy_sha256 != ?"#; +const ADVANCE_DIRTY_POLICY_SQL: &str = r#"UPDATE trade_dirty_generations +SET generation = generation + 1, evidence_policy_sha256 = ?, updated_at_unix_s = ? +WHERE evidence_policy_sha256 != ?"#; /// Stable source-free offline configuration-application failure classes. #[derive(Clone, Copy, Debug, PartialEq, Eq)] @@ -168,11 +176,23 @@ struct ConfigBinding { impl ConfigBinding { fn is_governed(&self) -> bool { self.config_contract_version == crate::RHI_CONFIG_SCHEMA_VERSION - && self.state_contract_version == crate::RHI_STATE_SCHEMA_VERSION + && (crate::RHI_STATE_BASE_SCHEMA_VERSION..=crate::RHI_STATE_SCHEMA_VERSION) + .contains(&self.state_contract_version) && self.admin_contract_version == crate::RHI_ADMIN_CONTRACT_VERSION && self.status_contract_version == crate::RHI_STATUS_CONTRACT_VERSION && self.provider_contract_version == crate::RHI_PROVIDER_CONTRACT_VERSION } + + fn is_same_identity_and_policy_except_state_version(&self, other: &Self) -> bool { + self.normalized_config_sha256 == other.normalized_config_sha256 + && self.evidence_policy_sha256 == other.evidence_policy_sha256 + && self.service_public_key == other.service_public_key + && self.config_contract_version == other.config_contract_version + && self.admin_contract_version == other.admin_contract_version + && self.status_contract_version == other.status_contract_version + && self.provider_contract_version == other.provider_contract_version + && self.state_contract_version < other.state_contract_version + } } impl From<&RhiStateMetadata> for ConfigBinding { @@ -232,6 +252,26 @@ pub(crate) async fn bind_or_verify( validate_history(&history)?; match history.last() { Some(actual) if actual.binding == expected => Ok(()), + Some(actual) + if actual + .binding + .is_same_identity_and_policy_except_state_version(&expected) => + { + if actual.generation >= RHI_CONFIG_BINDING_MAX_GENERATIONS { + return Err(ConfigOperationError::ResourceExhausted); + } + if applied_at.get() < actual.applied_at_unix_s { + return Err(ConfigOperationError::InvalidInput); + } + insert_binding( + transaction, + actual.generation + 1, + &expected, + applied_at.get(), + &build, + ) + .await + } Some(_) | None => Err(ConfigOperationError::Binding), } }) @@ -292,6 +332,14 @@ pub(crate) async fn append_configuration( return Err(ConfigOperationError::InvalidInput); } let generation = latest.generation + 1; + if latest.binding.evidence_policy_sha256 != candidate.evidence_policy_sha256 { + advance_dirty_policy( + transaction, + candidate.evidence_policy_sha256, + applied_at.get(), + ) + .await?; + } insert_binding( transaction, generation, @@ -316,6 +364,40 @@ pub(crate) async fn append_configuration( .map_err(map_transaction_error) } +async fn advance_dirty_policy( + transaction: &mut ServiceSqliteTransaction<'_>, + policy: [u8; 32], + updated_at_unix_s: u64, +) -> Result<(), ConfigOperationError> { + let row = sqlx::query(READ_DIRTY_POLICY_BOUNDS_SQL) + .bind(policy.as_slice()) + .fetch_one(&mut *transaction) + .await + .map_err(|_| ConfigOperationError::Storage)?; + let exhausted = row + .try_get::<i64, _>("exhausted") + .map_err(|_| ConfigOperationError::Storage)?; + let latest_updated_at = row + .try_get::<i64, _>("latest_updated_at") + .ok() + .and_then(|value| u64::try_from(value).ok()) + .ok_or(ConfigOperationError::Storage)?; + if exhausted != 0 { + return Err(ConfigOperationError::ResourceExhausted); + } + if updated_at_unix_s < latest_updated_at { + return Err(ConfigOperationError::InvalidInput); + } + sqlx::query(ADVANCE_DIRTY_POLICY_SQL) + .bind(policy.as_slice()) + .bind(i64::try_from(updated_at_unix_s).map_err(|_| ConfigOperationError::InvalidInput)?) + .bind(policy.as_slice()) + .execute(&mut *transaction) + .await + .map_err(|_| ConfigOperationError::Storage)?; + Ok(()) +} + async fn read_history( transaction: &mut ServiceSqliteTransaction<'_>, ) -> Result<Vec<HistoryEntry>, ConfigOperationError> { diff --git a/src/state_metadata.rs b/src/state_metadata.rs @@ -55,6 +55,10 @@ impl fmt::Debug for RhiNormalizedConfigDigest { pub struct RhiEvidencePolicyDigest([u8; 32]); impl RhiEvidencePolicyDigest { + pub(crate) const fn from_bytes(bytes: [u8; 32]) -> Self { + Self(bytes) + } + /// Returns the exact digest bytes. #[must_use] pub const fn as_bytes(&self) -> &[u8; 32] { diff --git a/src/state_trade.rs b/src/state_trade.rs @@ -107,6 +107,20 @@ impl RhiTradeSourceObservation { observed_at: event.observed_at_unix_seconds(), }) } + + pub(crate) fn from_parts( + source_id: Box<str>, + policy: RhiEvidencePolicyDigest, + event: &RhiAdmittedTradeMutationEvent, + ) -> Self { + Self { + source_id, + 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 { @@ -276,22 +290,22 @@ impl RhiStateRepositories<'_> { } } -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]>, +pub(crate) struct PersistenceRecord { + pub(crate) mutation_id: [u8; 32], + pub(crate) trade_id: [u8; 16], + pub(crate) contract_id: &'static str, + pub(crate) schema_version: u16, + pub(crate) event_id: [u8; 32], + pub(crate) event_signature: [u8; 64], + pub(crate) author_pubkey: [u8; 32], + pub(crate) event_kind: u32, + pub(crate) authored_at_unix_s: u64, + pub(crate) canonical_content: Box<[u8]>, + pub(crate) canonical_event_json: Box<[u8]>, } impl PersistenceRecord { - fn from_admitted( + pub(crate) fn from_admitted( admitted: RhiAdmittedTradeMutationEvent, ) -> Result<Self, RhiTradeEvidencePersistenceError> { let (_original, event, mutation, mutation_id, _) = admitted.into_parts(); @@ -322,13 +336,13 @@ impl PersistenceRecord { } #[derive(Clone, Copy, Debug, PartialEq, Eq)] -enum PersistenceOperationError { +pub(crate) enum PersistenceOperationError { MutationConflict, SignedEventConflict, Storage, } -async fn persist( +pub(crate) async fn persist( transaction: &mut ServiceSqliteTransaction<'_>, record: &PersistenceRecord, observation: &RhiTradeSourceObservation, 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 = 3", + "state_contract_version = 4", "admin_contract_version = 1", "status_contract_version = 1", "provider_contract_version = 1", @@ -43,7 +43,7 @@ fn source_lock_metadata_is_exact_and_nix_is_absent() { fn shared_host_packages_are_exactly_source_locked() { for dependency in ["radroots_service_host", "radroots_service_sqlite"] { assert!(MANIFEST.contains(&format!( - "{dependency} = {{ git = \"https://github.com/radrootslabs/lib\", rev = \"79d7818c8fe22a425f9524b884ddf59d25f0ef89\", version = \"=0.1.0-alpha\" }}" + "{dependency} = {{ git = \"https://github.com/radrootslabs/lib\", rev = \"21b11e7a5120ea949f7ad0838c746873fc73aac2\", version = \"=0.1.0-alpha\" }}" ))); } } @@ -51,7 +51,7 @@ fn shared_host_packages_are_exactly_source_locked() { #[test] fn shared_service_sqlite_is_the_only_catalog_authority() { assert!(MANIFEST.contains( - "radroots_service_sqlite = { git = \"https://github.com/radrootslabs/lib\", rev = \"79d7818c8fe22a425f9524b884ddf59d25f0ef89\", version = \"=0.1.0-alpha\" }" + "radroots_service_sqlite = { git = \"https://github.com/radrootslabs/lib\", rev = \"21b11e7a5120ea949f7ad0838c746873fc73aac2\", version = \"=0.1.0-alpha\" }" )); for forbidden in ["rusqlite", "libsqlite3-sys"] { assert!( @@ -64,14 +64,14 @@ fn shared_service_sqlite_is_the_only_catalog_authority() { #[test] fn shared_storage_generation_type_is_exactly_source_locked() { assert!(MANIFEST.contains( - "radroots_storage = { git = \"https://github.com/radrootslabs/lib\", rev = \"79d7818c8fe22a425f9524b884ddf59d25f0ef89\", version = \"=0.1.0-alpha\", default-features = false }" + "radroots_storage = { git = \"https://github.com/radrootslabs/lib\", rev = \"21b11e7a5120ea949f7ad0838c746873fc73aac2\", version = \"=0.1.0-alpha\", default-features = false }" )); } #[test] fn shared_transport_spi_is_exactly_source_locked_without_serde() { assert!(MANIFEST.contains( - "radroots_transport = { git = \"https://github.com/radrootslabs/lib\", rev = \"79d7818c8fe22a425f9524b884ddf59d25f0ef89\", version = \"=0.1.0-alpha\", default-features = false, features = [\"std\"] }" + "radroots_transport = { git = \"https://github.com/radrootslabs/lib\", rev = \"21b11e7a5120ea949f7ad0838c746873fc73aac2\", version = \"=0.1.0-alpha\", default-features = false, features = [\"std\"] }" )); } @@ -82,18 +82,18 @@ fn source_lock_binds_the_current_cargo_lock() { "schema = \"radroots.service.source-lock.v2\"\ncontract_version = 2\nservice = \"rhi\"\n" )); assert!(SOURCE_LOCK.contains(&format!("cargo_lock_sha256 = \"{digest}\""))); - assert!(SOURCE_LOCK.contains("revision = \"79d7818c8fe22a425f9524b884ddf59d25f0ef89\"")); + assert!(SOURCE_LOCK.contains("revision = \"21b11e7a5120ea949f7ad0838c746873fc73aac2\"")); assert!(SOURCE_LOCK.contains( "workspace_catalog_sha256 = \"deca0c080deae187ff8186c0708903e42f41ea57f77c5f91581e23aa561164a4\"" )); assert!(SOURCE_LOCK.contains( - "source_archive_sha256 = \"6f95d80d5ecb9d026ba6f9145b60fff67e370b2c175c438537319278e6974fa6\"" + "source_archive_sha256 = \"7e584a4b679264620d7bb6cf0a7028cc7651b33977b263c213f4e7b29c0e5a19\"" )); assert!(SOURCE_LOCK.contains("\n[nix]\nmaterial = \"absent\"\n")); assert!(!SOURCE_LOCK.contains("flake_lock_sha256")); assert!(!SOURCE_LOCK.contains("lib_revision =")); assert!(SOURCE_LOCK.ends_with( - "[contract_versions]\nconfig = 1\nstate = 3\nadmin = 1\nstatus = 1\nprovider = 1\n" + "[contract_versions]\nconfig = 1\nstate = 4\nadmin = 1\nstatus = 1\nprovider = 1\n" )); } diff --git a/tests/package_boundary.rs b/tests/package_boundary.rs @@ -17,6 +17,8 @@ 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 TRADE_SOURCE_INGEST_CONTRACT: &str = + include_str!("../contracts/services_hardening/trade_source_ingest.v1.json"); const PUBLIC_API: &str = include_str!("../contracts/api_baselines/rhi.txt"); const SOURCES: &[&str] = &[ include_str!("../src/adapters/nostr/event.rs"), @@ -28,6 +30,7 @@ const SOURCES: &[&str] = &[ include_str!("../src/runtime_context.rs"), include_str!("../src/runtime_adapters.rs"), include_str!("../src/runtime_foundation.rs"), + include_str!("../src/source_ingest.rs"), include_str!("../src/state_catalog.rs"), include_str!("../src/state_config.rs"), include_str!("../src/state_host.rs"), @@ -93,6 +96,7 @@ fn state_catalog_module_is_private_and_root_api_is_curated() { "runtime_context", "runtime_adapters", "runtime_foundation", + "source_ingest", "state_catalog", "state_config", "state_host", @@ -147,6 +151,17 @@ fn state_catalog_module_is_private_and_root_api_is_curated() { "RhiTradeEvidencePersistenceOutcome", "RhiTradeSourceObservation", "RHI_TRADE_EVIDENCE_PERSISTENCE_CONTRACT_VERSION", + "ingest_rhi_trade_source", + "RhiTradeDirtyGeneration", + "RhiTradeSourceAttempt", + "RhiTradeSourceCompletion", + "RhiTradeSourceCursor", + "RhiTradeSourceIngestError", + "RhiTradeSourceIngestErrorKind", + "RhiTradeSourceIngestOutcome", + "RHI_TRADE_SOURCE_INGEST_CONTRACT_VERSION", + "RHI_TRADE_SOURCE_RESULT_MAX_BYTES", + "RHI_TRADE_SOURCE_RESULT_MAX_EVENTS", ] { assert!( ROOT.contains(required), @@ -198,7 +213,86 @@ 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, 15); + assert_eq!(public_error_count, 16); +} + +#[test] +fn trade_source_ingest_is_exact_bounded_generation_fenced_and_sealed() { + let contract: serde_json::Value = + serde_json::from_str(TRADE_SOURCE_INGEST_CONTRACT).expect("trade-source ingest contract"); + assert_eq!(contract["schema"], "radroots.rhi.trade-source-ingest.v1"); + assert_eq!(contract["contract_version"], 1); + assert_eq!(contract["source"]["kind"], "nostr_relay"); + assert_eq!( + contract["source"]["selector"]["kinds"], + serde_json::json!([3470, 3471, 3472, 3473, 3474]) + ); + assert_eq!(contract["source"]["selector"]["exact_tag"], "#d"); + assert_eq!(contract["source"]["page_events_maximum"], 1_000); + assert_eq!(contract["source"]["result_events_maximum"], 4_096); + assert_eq!( + contract["source"]["result_original_event_bytes_maximum"], + 8_388_608 + ); + assert_eq!( + contract["completion"]["complete"], + "exact_target_eose_before_deadline" + ); + assert_eq!( + contract["checkpoint"]["scope"], + serde_json::json!([ + "source_id", + "selector_id", + "evidence_policy_sha256", + "trade_id" + ]) + ); + assert_eq!(contract["checkpoint"]["equal_timestamp_safe"], true); + assert_eq!( + contract["dirty_generation"]["do_not_advance_when"], + serde_json::json!([ + "rejected_event", + "duplicate_event", + "repeated_source_observation", + "operational_retry" + ]) + ); + assert_eq!( + contract["transaction"]["source_fetch_inside_transaction"], + false + ); + + let source = include_str!("../src/source_ingest.rs"); + for required in [ + "with_kinds(EVENT_KINDS.to_vec())", + "with_exact_tag_value('d', trade_id.to_hex())", + "with_since_unix_seconds(since)", + "admit_rhi_trade_mutation_event(", + "read_checkpoint(transaction", + "read_dirty(transaction", + "compare_cursor(current.cursor, candidate).is_lt()", + "new_relevant_evidence", + ] { + assert!( + source.contains(required), + "source ingest is missing {required}" + ); + } + for forbidden in [ + "pub fn host(", + "pub fn sqlite_host(", + "SystemTime", + "std::fs", + "std::net", + "tokio::spawn", + ] { + assert!( + !source.contains(forbidden), + "source ingest gained forbidden authority {forbidden}" + ); + } + assert!(!ROOT.contains("pub mod source_ingest")); + assert!(!PUBLIC_API.contains("rhi::source_ingest::")); } #[test] @@ -417,6 +511,12 @@ fn readme_freezes_the_root_only_boundary_and_exact_baseline() { "[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)", + "## Bounded relay-source ingestion", + "[`trade_source_ingest.v1.json`](contracts/services_hardening/trade_source_ingest.v1.json)", + "Only exact-target EOSE before the deadline is complete", + "4,096 distinct event identities and 8 MiB", + "observation, and operational retry do not", + "schema-v4 source-checkpoint and dirty-generation migration", "one canonical mutation", "every distinct valid signed", "does not advance reconciliation checkpoints or dirty generation", diff --git a/tests/services_hardening_config_lifecycle.rs b/tests/services_hardening_config_lifecycle.rs @@ -53,7 +53,7 @@ fn evidence(at: u64) -> (MigrationAppliedAtUnixSeconds, MigrationBuildIdentity) let build = MigrationBuildIdentity::new( env!("CARGO_PKG_VERSION"), "1111111111111111111111111111111111111111", - "79d7818c8fe22a425f9524b884ddf59d25f0ef89", + "21b11e7a5120ea949f7ad0838c746873fc73aac2", "rustc-test", "test-target", "service-host", @@ -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(), 3); + assert_eq!(row.try_get::<i64, _>("state_contract_version").unwrap(), 4); assert_eq!( row.try_get::<String, _>("service_public_key") .unwrap() diff --git a/tests/services_hardening_native_release.rs b/tests/services_hardening_native_release.rs @@ -13,7 +13,7 @@ const SYSTEMD_UNIT: &str = include_str!("../packaging/systemd/rhi@.service"); const RELEASE_ACCEPTANCE: &str = include_str!("../scripts/release-acceptance.sh"); const XTASK_MANIFEST: &str = include_str!("../tools/xtask/Cargo.toml"); -const LIB_REVISION: &str = "79d7818c8fe22a425f9524b884ddf59d25f0ef89"; +const LIB_REVISION: &str = "21b11e7a5120ea949f7ad0838c746873fc73aac2"; const LIB_REPOSITORY: &str = "https://github.com/radrootslabs/lib"; #[test] diff --git a/tests/services_hardening_runtime_foundation.rs b/tests/services_hardening_runtime_foundation.rs @@ -81,7 +81,7 @@ fn evidence() -> (MigrationAppliedAtUnixSeconds, MigrationBuildIdentity) { let build = MigrationBuildIdentity::new( env!("CARGO_PKG_VERSION"), "1111111111111111111111111111111111111111", - "79d7818c8fe22a425f9524b884ddf59d25f0ef89", + "21b11e7a5120ea949f7ad0838c746873fc73aac2", "rustc-test", "test-target", "service-host", diff --git a/tests/services_hardening_source_ingest.rs b/tests/services_hardening_source_ingest.rs @@ -0,0 +1,522 @@ +#![forbid(unsafe_code)] +#![cfg(any(target_os = "linux", target_os = "macos"))] + +use std::{ + collections::VecDeque, + fs, + os::unix::fs::PermissionsExt, + path::Path, + sync::{Arc, Mutex}, +}; + +use nostr::secp256k1::{Keypair, Message}; +use nostr::{EventId, Keys, SECP256K1}; +use radroots_service_sqlite::{MigrationAppliedAtUnixSeconds, MigrationBuildIdentity}; +use radroots_storage::event::SourceGeneration; +use radroots_transport::{ + BoxFuture, DeliveryReceipt, DeliveryRequest, EventSink, EventSource, EventSubscriber, + EventSubscription, FetchPage, FetchRequest, SinkFailure, SinkStatus, SourceStatus, + SubscriptionRequest, + outcome::{FetchTargetOutcome, FetchTargetState}, + source::{EventProvenance, FetchCursor, NextPage, ObservedEvent}, +}; +use rhi::{ + RadrootsHostEnvironment, RadrootsPathResolver, RadrootsPlatform, RhiConfigProfile, + RhiStateMetadata, RhiTradeMutationAuthoredTimePolicy, RhiTradeMutationObservedAtUnixSeconds, + RhiTradeSourceAttempt, RhiTradeSourceCompletion, RhiTransportAdapters, TradeId, + UnixTimeSeconds, ingest_rhi_trade_source, initialize_rhi_state, open_rhi_state_read_write, + parse_rhi_cli_v1_from, parse_rhi_config_v1, resolve_rhi_runtime_context, +}; +use serde_json::Value; + +const CONFIG: &str = include_str!("../contracts/services_hardening/config.v1.example.toml"); +const VECTOR: &str = include_str!("../contracts/conformance/vectors/trade_ingest_proposal.v1.json"); + +#[derive(Clone)] +struct PageSpec { + events: Vec<String>, + state: FetchTargetState, + next: Option<&'static str>, +} + +#[derive(Clone)] +struct ScriptedTransport { + pages: Arc<Mutex<VecDeque<PageSpec>>>, + requests: Arc<Mutex<Vec<FetchRequest>>>, +} + +impl ScriptedTransport { + fn new(pages: impl IntoIterator<Item = PageSpec>) -> Self { + Self { + pages: Arc::new(Mutex::new(pages.into_iter().collect())), + requests: Arc::new(Mutex::new(Vec::new())), + } + } + + fn requests(&self) -> Vec<FetchRequest> { + self.requests.lock().expect("requests").clone() + } +} + +impl EventSource for ScriptedTransport { + fn status(&self) -> BoxFuture<'_, Result<SourceStatus, radroots_transport::Error>> { + Box::pin(async { Err(radroots_transport::Error::UnsupportedOperation) }) + } + + fn fetch( + &self, + request: FetchRequest, + ) -> BoxFuture<'_, Result<FetchPage, radroots_transport::Error>> { + self.requests + .lock() + .expect("requests") + .push(request.clone()); + let page = self.pages.lock().expect("pages").pop_front(); + Box::pin(async move { + let spec = page.ok_or(radroots_transport::Error::UnsupportedOperation)?; + let target = request.target_set().targets().first().expect("one target"); + let events = spec + .events + .into_iter() + .map(|raw| { + let event = radroots_event_codec::decode::signed_event(raw.as_str()) + .expect("signed event"); + let mut provenance = EventProvenance::new( + radroots_transport::TransportId::NOSTR, + target.fingerprint().clone(), + 1, + ) + .expect("provenance"); + if let Some(cursor) = request.cursor().cloned() { + provenance = provenance.with_cursor(cursor); + } + ObservedEvent::new(event, provenance) + }) + .collect(); + let outcome = FetchTargetOutcome::new(target.fingerprint().clone(), spec.state); + let next = match spec.next { + Some(cursor) => NextPage::Cursor(FetchCursor::parse(cursor).expect("cursor")), + None => NextPage::Complete, + }; + FetchPage::for_request(&request, events, vec![outcome], next) + }) + } +} + +impl EventSubscriber for ScriptedTransport { + fn subscribe( + &self, + _request: SubscriptionRequest, + ) -> BoxFuture<'_, Result<Box<dyn EventSubscription>, radroots_transport::Error>> { + Box::pin(async { Err(radroots_transport::Error::UnsupportedOperation) }) + } +} + +impl EventSink for ScriptedTransport { + fn status(&self) -> BoxFuture<'_, Result<SinkStatus, radroots_transport::Error>> { + Box::pin(async { Err(radroots_transport::Error::UnsupportedOperation) }) + } + + fn deliver( + &self, + request: DeliveryRequest, + ) -> BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> { + Box::pin(async move { Err(SinkFailure::invalid_contract(&request)) }) + } +} + +fn adapters(source: &ScriptedTransport) -> RhiTransportAdapters { + RhiTransportAdapters::new( + Arc::new(source.clone()), + Arc::new(source.clone()), + Arc::new(source.clone()), + ) +} + +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", + "21b11e7a5120ea949f7ad0838c746873fc73aac2", + "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 vector_wire() -> String { + let vector: Value = serde_json::from_str(VECTOR).expect("vector"); + vector["raw_json"].as_str().expect("raw event").to_owned() +} + +fn resign_wire(raw: &str, auxiliary: u8) -> String { + let mut event: Value = serde_json::from_str(raw).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()); + event["sig"] = SECP256K1 + .sign_schnorr_with_aux_rand(&message, &keypair, &[auxiliary; 32]) + .to_string() + .into(); + serde_json::to_string(&event).expect("event JSON") +} + +fn attempt(request_id: &str, started_at: u64, observed_at: u64) -> RhiTradeSourceAttempt { + RhiTradeSourceAttempt::new( + request_id, + UnixTimeSeconds::new(started_at), + RhiTradeMutationObservedAtUnixSeconds::new(observed_at).expect("observed"), + RhiTradeMutationAuthoredTimePolicy::new(0).expect("authored policy"), + ) + .expect("attempt") +} + +#[tokio::test] +async fn exact_selector_fetch_replay_and_rejection_preserve_cursor_and_generation() { + 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 trade_id = TradeId::parse("11111111111111111111111111111111").expect("trade id"); + let wire = vector_wire(); + + let source = ScriptedTransport::new([PageSpec { + events: vec![wire.clone()], + state: FetchTargetState::Complete, + next: None, + }]); + let outcome = ingest_rhi_trade_source( + &host.repositories(), + &adapters(&source), + &configuration, + "trade-primary", + trade_id, + attempt("attempt-1", 1_784_347_200, 1_784_347_200), + ) + .await + .expect("ingest"); + assert_eq!(outcome.completion(), RhiTradeSourceCompletion::Complete); + assert_eq!(outcome.received_events(), 1); + assert_eq!(outcome.admitted_events(), 1); + assert_eq!(outcome.rejected_events(), 0); + assert_eq!(outcome.inserted_mutations(), 1); + assert_eq!(outcome.inserted_signed_events(), 1); + assert_eq!(outcome.inserted_observations(), 1); + assert!(outcome.checkpoint_advanced()); + assert_eq!(outcome.dirty_generation().expect("generation").get(), 1); + assert!(outcome.dirty_generation_advanced()); + + let requests = source.requests(); + assert_eq!(requests.len(), 1); + assert_eq!( + requests[0].selector().kinds(), + &[3470, 3471, 3472, 3473, 3474] + ); + assert_eq!( + requests[0] + .selector() + .exact_tag_filters() + .collect::<Vec<_>>(), + vec![('d', &["11111111111111111111111111111111".to_owned()][..])] + ); + assert_eq!( + requests[0].selector().since_unix_seconds(), + Some(1_784_260_800) + ); + assert_eq!(requests[0].bounds().deadline_unix_ms(), 1_784_347_210_000); + + let replay_source = ScriptedTransport::new([PageSpec { + events: vec![wire.clone(), wire.clone()], + state: FetchTargetState::Complete, + next: None, + }]); + let replay = ingest_rhi_trade_source( + &host.repositories(), + &adapters(&replay_source), + &configuration, + "trade-primary", + trade_id, + attempt("attempt-2", 1_784_347_201, 1_784_347_201), + ) + .await + .expect("replay"); + assert_eq!(replay.admitted_events(), 1); + assert_eq!(replay.duplicate_events(), 1); + assert_eq!(replay.inserted_mutations(), 0); + assert_eq!(replay.inserted_signed_events(), 0); + assert_eq!(replay.inserted_observations(), 1); + assert!(!replay.checkpoint_advanced()); + assert_eq!(replay.dirty_generation().expect("generation").get(), 1); + assert!(!replay.dirty_generation_advanced()); + + let rejected_source = ScriptedTransport::new([PageSpec { + events: vec![resign_wire(&wire, 9)], + state: FetchTargetState::Complete, + next: None, + }]); + let rejected = ingest_rhi_trade_source( + &host.repositories(), + &adapters(&rejected_source), + &configuration, + "trade-primary", + trade_id, + attempt("attempt-3", 1_784_347_000, 1_784_347_199), + ) + .await + .expect("rejected attempt"); + assert_eq!(rejected.completion(), RhiTradeSourceCompletion::Complete); + assert_eq!(rejected.rejected_events(), 1); + assert_eq!(rejected.inserted_signed_events(), 0); + assert!(!rejected.checkpoint_advanced()); + assert_eq!(rejected.dirty_generation().expect("generation").get(), 1); + assert!(!rejected.dirty_generation_advanced()); + + host.close().await.expect("close"); +} + +#[tokio::test] +async fn incomplete_and_unsupported_results_never_advance_checkpoint_or_dirty_generation() { + 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 trade_id = TradeId::parse("11111111111111111111111111111111").expect("trade id"); + let vectors = [ + ( + FetchTargetState::Cancelled, + RhiTradeSourceCompletion::IncompleteTimeout, + ), + ( + FetchTargetState::Unavailable, + RhiTradeSourceCompletion::IncompleteUnavailable, + ), + ( + FetchTargetState::FailedRetryable, + RhiTradeSourceCompletion::IncompleteUnavailable, + ), + ( + FetchTargetState::FailedTerminal, + RhiTradeSourceCompletion::IncompleteUnknown, + ), + ( + FetchTargetState::Partial, + RhiTradeSourceCompletion::IncompleteUnknown, + ), + ]; + for (index, (state, expected)) in vectors.into_iter().enumerate() { + let source = ScriptedTransport::new([PageSpec { + events: Vec::new(), + state, + next: None, + }]); + let outcome = ingest_rhi_trade_source( + &host.repositories(), + &adapters(&source), + &configuration, + "trade-primary", + trade_id, + attempt( + &format!("incomplete-{index}"), + 1_784_347_300 + index as u64, + 1_784_347_300 + index as u64, + ), + ) + .await + .expect("classified incomplete result"); + assert_eq!(outcome.completion(), expected); + assert_eq!(outcome.checkpoint(), None); + assert_eq!(outcome.dirty_generation(), None); + assert!(!outcome.checkpoint_advanced()); + assert!(!outcome.dirty_generation_advanced()); + } + + let unsupported = ScriptedTransport::new([]); + let outcome = ingest_rhi_trade_source( + &host.repositories(), + &adapters(&unsupported), + &configuration, + "trade-primary", + trade_id, + attempt("unsupported", 1_784_347_400, 1_784_347_400), + ) + .await + .expect("unsupported result"); + assert_eq!(outcome.completion(), RhiTradeSourceCompletion::Unsupported); + assert_eq!(outcome.checkpoint(), None); + assert_eq!(outcome.dirty_generation(), None); + + host.close().await.expect("close"); +} + +#[tokio::test] +async fn configured_result_bound_retains_admitted_evidence_without_checkpoint() { + 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 trade_id = TradeId::parse("11111111111111111111111111111111").expect("trade id"); + let wire = vector_wire(); + let source = ScriptedTransport::new([ + PageSpec { + events: vec![wire.clone(); 1_000], + state: FetchTargetState::Complete, + next: Some("page-1"), + }, + PageSpec { + events: vec![wire.clone(); 1_000], + state: FetchTargetState::Complete, + next: Some("page-2"), + }, + PageSpec { + events: vec![wire.clone(); 1_000], + state: FetchTargetState::Complete, + next: Some("page-3"), + }, + PageSpec { + events: vec![wire.clone(); 1_000], + state: FetchTargetState::Complete, + next: Some("page-4"), + }, + PageSpec { + events: vec![wire; 97], + state: FetchTargetState::Complete, + next: None, + }, + ]); + let outcome = ingest_rhi_trade_source( + &host.repositories(), + &adapters(&source), + &configuration, + "trade-primary", + trade_id, + attempt("resource-limit", 1_784_347_500, 1_784_347_500), + ) + .await + .expect("resource-limited result"); + assert_eq!( + outcome.completion(), + RhiTradeSourceCompletion::IncompleteResourceLimit + ); + assert_eq!(outcome.received_events(), 4_096); + assert_eq!(outcome.admitted_events(), 1); + assert_eq!(outcome.inserted_mutations(), 1); + assert_eq!(outcome.inserted_signed_events(), 1); + assert_eq!(outcome.inserted_observations(), 1); + assert_eq!(outcome.checkpoint(), None); + assert!(!outcome.checkpoint_advanced()); + assert_eq!(outcome.dirty_generation().expect("dirty").get(), 1); + assert!(outcome.dirty_generation_advanced()); + + host.close().await.expect("close"); +} + +#[test] +fn attempt_and_public_diagnostics_are_bounded_and_redacted() { + let observed = RhiTradeMutationObservedAtUnixSeconds::new(100).expect("observed"); + let policy = RhiTradeMutationAuthoredTimePolicy::new(0).expect("policy"); + assert!(RhiTradeSourceAttempt::new("", UnixTimeSeconds::new(1), observed, policy).is_err()); + assert!( + RhiTradeSourceAttempt::new("x".repeat(257), UnixTimeSeconds::new(1), observed, policy,) + .is_err() + ); + let attempt = RhiTradeSourceAttempt::new( + "secret-attempt-id", + UnixTimeSeconds::new(99), + observed, + policy, + ) + .expect("attempt"); + assert!(!format!("{attempt:?}").contains("secret-attempt-id")); + + let error = RhiTradeSourceAttempt::new( + "\nsecret-request", + UnixTimeSeconds::new(99), + observed, + policy, + ) + .expect_err("control character"); + assert_eq!(error.code(), "trade_source_input_invalid"); + assert!(!format!("{error}").contains("secret-request")); + assert!(!format!("{error:?}").contains("secret-request")); + assert!(std::error::Error::source(&error).is_none()); +} diff --git a/tests/services_hardening_state_catalog.rs b/tests/services_hardening_state_catalog.rs @@ -12,8 +12,10 @@ use rhi::{ 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, 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, + RHI_STATE_SCHEMA_VERSION_3_SHA256, RHI_STATE_SCHEMA_VERSION_4_MIGRATION_SHA256, + RHI_STATE_SCHEMA_VERSION_4_OBJECT_COUNT, RHI_STATE_SCHEMA_VERSION_4_SHA256, + RhiStateCatalogErrorKind, rhi_migration_catalog, rhi_schema_catalog, + validate_rhi_state_catalogs, }; const CATALOG_SOURCE: &str = include_str!("../src/state_catalog.rs"); @@ -21,14 +23,14 @@ const LIB_SOURCE: &str = include_str!("../src/lib.rs"); const MANIFEST: &str = include_str!("../Cargo.toml"); #[test] -fn schema_v1_through_v3_catalogs_have_exact_literal_identities() { +fn schema_v1_through_v4_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, 3); - assert_eq!(migrations.descriptors().len(), 2); - assert_eq!(migrations.current_version(), 3); + assert_eq!(RHI_STATE_SCHEMA_VERSION, 4); + assert_eq!(migrations.descriptors().len(), 3); + assert_eq!(migrations.current_version(), 4); assert_eq!(migrations.descriptors()[0].target_version(), 2); assert_eq!( migrations.descriptors()[0].name().as_str(), @@ -47,12 +49,21 @@ fn schema_v1_through_v3_catalogs_have_exact_literal_identities() { migrations.descriptors()[1].checksum().as_bytes(), &RHI_STATE_SCHEMA_VERSION_3_MIGRATION_SHA256 ); + assert_eq!(migrations.descriptors()[2].target_version(), 4); + assert_eq!( + migrations.descriptors()[2].name().as_str(), + "create_source_checkpoints_and_dirty_generations" + ); + assert_eq!( + migrations.descriptors()[2].checksum().as_bytes(), + &RHI_STATE_SCHEMA_VERSION_4_MIGRATION_SHA256 + ); assert_eq!( migrations.digest().as_bytes(), &RHI_MIGRATION_CATALOG_SHA256 ); - assert_eq!(schema.versions().len(), 3); + assert_eq!(schema.versions().len(), 4); let version = schema.versions()[0]; assert_eq!(version.version(), 1); assert_eq!( @@ -86,13 +97,24 @@ fn schema_v1_through_v3_catalogs_have_exact_literal_identities() { version.digest().as_bytes(), &RHI_STATE_SCHEMA_VERSION_3_SHA256 ); + let version = schema.versions()[3]; + assert_eq!(version.version(), 4); + assert_eq!( + version.object_count(), + RHI_STATE_SCHEMA_VERSION_4_OBJECT_COUNT + ); + assert_eq!(version.object_count(), 28); + assert_eq!( + version.digest().as_bytes(), + &RHI_STATE_SCHEMA_VERSION_4_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), - "14046048b468836f2602ec51e538f2a98b71c945f83d933dccb8603b84f865f0" + "295ff5e9ac0eadc1f2c2d8c258c2eaa816f946560b689dcb432fdc99e79ec353" ); assert_eq!( lower_hex(&RHI_STATE_SCHEMA_VERSION_1_SHA256), @@ -115,8 +137,16 @@ fn schema_v1_through_v3_catalogs_have_exact_literal_identities() { "fd96226405ab68655aee00f683f6023c7aab2dbd23feadac165333490d6f0ad5" ); assert_eq!( + lower_hex(&RHI_STATE_SCHEMA_VERSION_4_MIGRATION_SHA256), + "2441f7c4c11594c6fde787dbb94e44ba00b72d2c979b0e7fd3c733eae9144d78" + ); + assert_eq!( + lower_hex(&RHI_STATE_SCHEMA_VERSION_4_SHA256), + "9ca78a54b0ea2013e7aa70d09beadcfb6cdde07c40429dffe91cb83c2a425cea" + ); + assert_eq!( lower_hex(&RHI_STATE_SCHEMA_CATALOG_SHA256), - "1325a41b90abfc7d523bbfe83500a1b1a73e492cc930d6ead6c71e29d5447436" + "58789c1d396514ce19a560b67bf760650b9d6fa378bc8234c1f6a39ad2f556e6" ); } @@ -160,8 +190,15 @@ fn independent_validator_rejects_migration_or_schema_drift() { 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"); + let snapshot_digest = + SchemaVersionCatalog::computed_digest(4, [version_two_object()]).expect("v4 digest"); + let version_four = SchemaVersionCatalog::new(4, [version_two_object()], snapshot_digest) + .expect("version four"); + let schema = SchemaCatalog::new( + &exact_migrations, + [version_one, version_two, version_three, version_four], + ) + .expect("drift schema catalog"); assert_eq!( validate_rhi_state_catalogs(&exact_migrations, &schema) .expect_err("schema drift") @@ -201,8 +238,15 @@ fn catalog_errors_are_stable_source_free_and_redacted() { 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 version_four_digest = + SchemaVersionCatalog::computed_digest(4, [secret_object()]).expect("v4 digest"); + let version_four = + SchemaVersionCatalog::new(4, [secret_object()], version_four_digest).expect("version four"); + let schema = SchemaCatalog::new( + &migrations, + [version_one, version_two, version_three, version_four], + ) + .expect("schema catalog"); let error = validate_rhi_state_catalogs(&migrations, &schema).expect_err("mismatch"); assert_eq!(error.kind(), RhiStateCatalogErrorKind::CatalogMismatch); @@ -249,7 +293,7 @@ fn secret_object() -> SchemaObject { #[test] fn catalog_source_is_pure_pinned_and_uses_only_the_shared_authority() { assert!(MANIFEST.contains( - "radroots_service_sqlite = { git = \"https://github.com/radrootslabs/lib\", rev = \"79d7818c8fe22a425f9524b884ddf59d25f0ef89\", version = \"=0.1.0-alpha\" }" + "radroots_service_sqlite = { git = \"https://github.com/radrootslabs/lib\", rev = \"21b11e7a5120ea949f7ad0838c746873fc73aac2\", version = \"=0.1.0-alpha\" }" )); assert!(LIB_SOURCE.contains("mod state_catalog;")); assert!(!LIB_SOURCE.contains("pub mod state_catalog;")); diff --git a/tests/services_hardening_state_host.rs b/tests/services_hardening_state_host.rs @@ -60,7 +60,7 @@ fn migration_evidence() -> (MigrationAppliedAtUnixSeconds, MigrationBuildIdentit let build = MigrationBuildIdentity::new( env!("CARGO_PKG_VERSION"), "1111111111111111111111111111111111111111", - "79d7818c8fe22a425f9524b884ddf59d25f0ef89", + "21b11e7a5120ea949f7ad0838c746873fc73aac2", "rustc-test", "test-target", "service-host", @@ -227,7 +227,7 @@ async fn missing_state_and_mismatched_evidence_fail_before_database_creation() { let invalid_build = MigrationBuildIdentity::new( env!("CARGO_PKG_VERSION"), "1111111111111111111111111111111111111111", - "79d7818c8fe22a425f9524b884ddf59d25f0ef89", + "21b11e7a5120ea949f7ad0838c746873fc73aac2", "rustc-test", "test-target", "service-host", diff --git a/tests/services_hardening_state_resilience.rs b/tests/services_hardening_state_resilience.rs @@ -73,7 +73,7 @@ fn migration_evidence() -> (MigrationAppliedAtUnixSeconds, MigrationBuildIdentit let build = MigrationBuildIdentity::new( env!("CARGO_PKG_VERSION"), "1111111111111111111111111111111111111111", - "79d7818c8fe22a425f9524b884ddf59d25f0ef89", + "21b11e7a5120ea949f7ad0838c746873fc73aac2", "rustc-test", "test-target", "service-host", @@ -352,7 +352,7 @@ async fn exact_open_rejects_unexpected_migration_history_without_repair() { ) .bind([0x44_u8; 32].as_slice()) .bind("1111111111111111111111111111111111111111") - .bind("79d7818c8fe22a425f9524b884ddf59d25f0ef89") + .bind("21b11e7a5120ea949f7ad0838c746873fc73aac2") .execute(&mut connection) .await .expect("insert unexpected ledger row"); diff --git a/tests/services_hardening_trade_persistence.rs b/tests/services_hardening_trade_persistence.rs @@ -66,7 +66,7 @@ fn migration_evidence() -> (MigrationAppliedAtUnixSeconds, MigrationBuildIdentit MigrationBuildIdentity::new( env!("CARGO_PKG_VERSION"), "1111111111111111111111111111111111111111", - "79d7818c8fe22a425f9524b884ddf59d25f0ef89", + "21b11e7a5120ea949f7ad0838c746873fc73aac2", "rustc-test", "test-target", "service-host", diff --git a/tests/services_hardening_wave_100_b.rs b/tests/services_hardening_wave_100_b.rs @@ -92,7 +92,7 @@ fn migration_evidence() -> (MigrationAppliedAtUnixSeconds, MigrationBuildIdentit let build = MigrationBuildIdentity::new( env!("CARGO_PKG_VERSION"), "1111111111111111111111111111111111111111", - "79d7818c8fe22a425f9524b884ddf59d25f0ef89", + "21b11e7a5120ea949f7ad0838c746873fc73aac2", "rustc-test", "test-target", "service-host", diff --git a/tests/source_guards.rs b/tests/source_guards.rs @@ -43,7 +43,7 @@ fn rhi_manifest_exact_pins_radroots_contract() { ); assert_eq!( dependency.get("rev").and_then(toml::Value::as_str), - Some("79d7818c8fe22a425f9524b884ddf59d25f0ef89"), + Some("21b11e7a5120ea949f7ad0838c746873fc73aac2"), "RHI must source-lock {name} to the exact promoted Lib revision" ); assert_eq!(