commit 28efdc18b283cc4ddaf65203fcfc7e40a88ce498
parent dd0fa0f4befe1e36509369a3aed7ebd6dd095ef4
Author: triesap <tyson@radroots.org>
Date: Mon, 24 Aug 2026 07:33:22 +0000
rhi: commit reconciliation source results
- fence replay commits by the exact lease, generation, policy, and prior cursors
- persist immutable attempt results and full accepted-inventory digests in schema v6
- advance dirty state and eligible checkpoints atomically with canonical evidence
- qualify idempotent retry, incomplete outcomes, stale leases, catalogs, and API scope
Diffstat:
17 files changed, 1757 insertions(+), 52 deletions(-)
diff --git a/AGENTS.md b/AGENTS.md
@@ -217,8 +217,9 @@
- Create-new state begins at the shared schema-v1 baseline and applies the
governed RHI schema-v2 configuration-binding migration, schema-v3
immutable trade-evidence migration, schema-v4 source-checkpoint and
- dirty-generation migration, and schema-v5 bounded reconciliation-job
- migration. Retain at most 1,024 consecutive
+ dirty-generation migration, schema-v5 bounded reconciliation-job migration,
+ and schema-v6 immutable reconciliation-attempt/source-result 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
@@ -234,6 +235,13 @@
use intent-open, discover source generation under retained authority, and
match the latest durable binding; configuration apply is an exclusive offline
operation.
+- Commit one exact reconciliation-attempt replay inventory only through the
+ typed attempt repository. Revalidate the exact live lease, dirty generation,
+ evidence policy, and every scoped prior checkpoint before mutation; persist
+ evidence, immutable results, the exact ordered fact/provenance inventory
+ digest, at most one dirty advance, and eligible checkpoints atomically.
+ Incomplete or unsupported results never advance, and committed cursor
+ evidence is minted only after durable commit confirmation.
- Never hold a database transaction while waiting for a source, relay, DNS,
identity provider, clock, entropy, signing, reduction, or backoff.
- Never prune active jobs/outboxes, migration history, current identity/policy
diff --git a/README b/README
@@ -168,8 +168,9 @@ their later ordered checkpoints. The exact machine contract is
sealed prior-cursor evidence, configured overlap, and an inclusive query start
under a domain-separated identity. The evidence is sealed to the source,
trade, evidence-policy digest, and selector digest, so callers cannot forge or
-relabel progress. Step 189 provides no evidence-minting path: Step 190 alone
-may construct it after the replay and cursor commit atomically. Initial
+relabel progress. This pure replay model provides no evidence-minting path; the
+atomic commit boundary below alone constructs it after replay and cursor commit.
+Initial
queries use the configured lookback; resumed queries subtract the configured
overlap from the prior authored-time
cursor so equal-timestamp events remain discoverable. A cursor becomes
@@ -184,10 +185,34 @@ conflicting mutation or signed-event identity reuse. Result counts and bytes
are derived from the distinct canonical inventory rather than caller claims.
It performs no SQLite, source, relay, network, filesystem, task, clock, or
entropy operation. Durable scope revalidation, result/completion persistence,
-and checkpoint advancement remain Step 190 authority. The exact machine
-contract is
+and checkpoint advancement belong to the atomic commit boundary below. The
+exact machine contract is
[`reconciliation_replay.v1.json`](contracts/services_hardening/reconciliation_replay.v1.json).
+## Atomic reconciliation result commit
+
+`RhiReconciliationAttemptRepository::commit_source_replays` accepts exactly one
+bounded replay result for every request in a claim-derived attempt. One short
+SQLx transaction first revalidates the exact live lease, trade dirty generation,
+evidence-policy digest, and each source-scoped prior checkpoint. It then commits
+canonical mutations, signed events, first-source observations, the immutable
+attempt and source-result inventory, a domain-separated digest of every exact
+ordered persisted fact and provenance value, at most one dirty-generation
+advance, and only eligible cursor checkpoints as one atomic unit. No source,
+relay, network, clock, entropy, or task operation occurs inside that
+transaction.
+
+Only exact `complete` evidence with a strictly newer candidate advances a
+checkpoint. Timeout, unavailable, resource-limited, unknown, and unsupported
+results remain durable but never become progress. Newly inserted canonical
+mutation or signed-event evidence advances the dirty generation once for the
+whole attempt; observation-only replay does not. Exact retry of an already
+committed attempt reconciles idempotently before stale lease or generation
+fences, while any inventory mismatch fails closed. Sealed source/trade/policy/
+selector cursor evidence is minted only after durable commit confirmation. The
+exact machine contract is
+[`reconciliation_commit.v1.json`](contracts/services_hardening/reconciliation_commit.v1.json).
+
## Existing-state runtime foundation
`open_rhi_runtime_foundation` opens only an already initialized database from
@@ -285,13 +310,16 @@ 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, schema-v3 immutable trade-evidence
-migration, schema-v4 source-checkpoint and dirty-generation migration, and
-schema-v5 bounded reconciliation-job migration.
-Version five contains the six shared immutable service-metadata and
+migration, schema-v4 source-checkpoint and dirty-generation migration,
+schema-v5 bounded reconciliation-job migration, and schema-v6 immutable
+reconciliation-result migration.
+Version six contains the six shared immutable service-metadata and
migration-ledger objects, the bounded append-only `rhi_config_bindings` table,
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.
+It also retains immutable attempt and source-result rows under no-update and
+no-delete triggers.
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.
@@ -308,9 +336,8 @@ 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.
-Durable source-attempt/completion inventory, reconciliation jobs, reports,
-attestations, and publication state remain reserved for their later owning
-checkpoints.
+Evidence manifests, reports, attestations, publication state, and final job
+completion remain reserved for their later owning checkpoints.
The sealed RHI state-host lifecycle now reserves and initializes a missing
canonical database only through an explicit create-new operation. Ordinary
diff --git a/contracts/api_baselines/rhi.txt b/contracts/api_baselines/rhi.txt
@@ -149,6 +149,16 @@ pub enum rhi::RhiReconciliationCommandV1
pub rhi::RhiReconciliationCommandV1::Jobs
pub rhi::RhiReconciliationCommandV1::Refresh
pub rhi::RhiReconciliationCommandV1::Status
+pub enum rhi::RhiReconciliationCommitErrorKind
+pub rhi::RhiReconciliationCommitErrorKind::CommitOutcomeUnknown
+pub rhi::RhiReconciliationCommitErrorKind::Conflict
+pub rhi::RhiReconciliationCommitErrorKind::GenerationConflict
+pub rhi::RhiReconciliationCommitErrorKind::InvalidInput
+pub rhi::RhiReconciliationCommitErrorKind::InvalidMode
+pub rhi::RhiReconciliationCommitErrorKind::LeaseLost
+pub rhi::RhiReconciliationCommitErrorKind::Storage
+impl rhi::RhiReconciliationCommitErrorKind
+pub const fn rhi::RhiReconciliationCommitErrorKind::code(self) -> &'static str
pub enum rhi::RhiReconciliationJobErrorKind
pub rhi::RhiReconciliationJobErrorKind::CommitOutcomeUnknown
pub rhi::RhiReconciliationJobErrorKind::DirtyGenerationConflict
@@ -616,6 +626,8 @@ impl core::fmt::Debug for rhi::RhiReconciliationAttemptPlan
pub fn rhi::RhiReconciliationAttemptPlan::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
pub struct rhi::RhiReconciliationAttemptRepository<'host>
impl rhi::RhiReconciliationAttemptRepository<'_>
+pub async fn rhi::RhiReconciliationAttemptRepository<'_>::commit_source_replays<I>(&self, rhi::RhiReconciliationLease, rhi::RhiReconciliationAttemptPlan, I) -> core::result::Result<rhi::RhiReconciliationSourceCommitOutcome, rhi::RhiReconciliationCommitError> where I: core::iter::traits::collect::IntoIterator<Item = rhi::RhiReconciliationSourceReplay>
+impl rhi::RhiReconciliationAttemptRepository<'_>
pub const fn rhi::RhiReconciliationAttemptRepository<'_>::descriptor(&self) -> rhi::RhiStateRepositoryDescriptor
pub const fn rhi::RhiReconciliationAttemptRepository<'_>::kind(&self) -> rhi::RhiStateRepositoryKind
impl core::fmt::Debug for rhi::RhiReconciliationAttemptRepository<'_>
@@ -627,6 +639,15 @@ pub fn rhi::RhiReconciliationAttemptResults::new<I>(&rhi::RhiReconciliationAttem
pub fn rhi::RhiReconciliationAttemptResults::results(&self) -> &[rhi::RhiReconciliationSourceResult]
impl core::fmt::Debug for rhi::RhiReconciliationAttemptResults
pub fn rhi::RhiReconciliationAttemptResults::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct rhi::RhiReconciliationCommitError
+impl rhi::RhiReconciliationCommitError
+pub const fn rhi::RhiReconciliationCommitError::code(self) -> &'static str
+pub const fn rhi::RhiReconciliationCommitError::kind(self) -> rhi::RhiReconciliationCommitErrorKind
+impl core::error::Error for rhi::RhiReconciliationCommitError
+impl core::fmt::Debug for rhi::RhiReconciliationCommitError
+pub fn rhi::RhiReconciliationCommitError::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+impl core::fmt::Display for rhi::RhiReconciliationCommitError
+pub fn rhi::RhiReconciliationCommitError::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
pub struct rhi::RhiReconciliationJob
impl rhi::RhiReconciliationJob
pub const fn rhi::RhiReconciliationJob::attempt_count(self) -> u16
@@ -707,6 +728,15 @@ pub struct rhi::RhiReconciliationScheduleOutcome
impl rhi::RhiReconciliationScheduleOutcome
pub const fn rhi::RhiReconciliationScheduleOutcome::created(self) -> bool
pub const fn rhi::RhiReconciliationScheduleOutcome::job(self) -> rhi::RhiReconciliationJob
+pub struct rhi::RhiReconciliationSourceCommitOutcome
+impl rhi::RhiReconciliationSourceCommitOutcome
+pub const fn rhi::RhiReconciliationSourceCommitOutcome::checkpoint_advance_count(&self) -> u32
+pub fn rhi::RhiReconciliationSourceCommitOutcome::committed_cursors(&self) -> &[rhi::RhiReconciliationSourceCursorEvidence]
+pub const fn rhi::RhiReconciliationSourceCommitOutcome::created(&self) -> bool
+pub const fn rhi::RhiReconciliationSourceCommitOutcome::dirty_generation_advanced(&self) -> bool
+pub const fn rhi::RhiReconciliationSourceCommitOutcome::source_result_count(&self) -> u32
+impl core::fmt::Debug for rhi::RhiReconciliationSourceCommitOutcome
+pub fn rhi::RhiReconciliationSourceCommitOutcome::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
pub struct rhi::RhiReconciliationSourceCursorEvidence
impl rhi::RhiReconciliationSourceCursorEvidence
pub const fn rhi::RhiReconciliationSourceCursorEvidence::cursor(&self) -> rhi::RhiTradeSourceCursor
@@ -1158,6 +1188,7 @@ pub const rhi::RHI_MIGRATION_CATALOG_SHA256: [u8; 32]
pub const rhi::RHI_PROVIDER_CONTRACT_VERSION: u32
pub const rhi::RHI_RECONCILIATION_ATTEMPT_CONTRACT_VERSION: u32
pub const rhi::RHI_RECONCILIATION_ATTEMPT_MAX_SOURCES: usize
+pub const rhi::RHI_RECONCILIATION_COMMIT_CONTRACT_VERSION: u32
pub const rhi::RHI_RECONCILIATION_JOB_CONTRACT_VERSION: u32
pub const rhi::RHI_RECONCILIATION_JOB_MAX_ACTIVE: u32
pub const rhi::RHI_RECONCILIATION_REPLAY_CONTRACT_VERSION: u32
@@ -1185,6 +1216,9 @@ pub const rhi::RHI_STATE_SCHEMA_VERSION_4_SHA256: [u8; 32]
pub const rhi::RHI_STATE_SCHEMA_VERSION_5_MIGRATION_SHA256: [u8; 32]
pub const rhi::RHI_STATE_SCHEMA_VERSION_5_OBJECT_COUNT: u32
pub const rhi::RHI_STATE_SCHEMA_VERSION_5_SHA256: [u8; 32]
+pub const rhi::RHI_STATE_SCHEMA_VERSION_6_MIGRATION_SHA256: [u8; 32]
+pub const rhi::RHI_STATE_SCHEMA_VERSION_6_OBJECT_COUNT: u32
+pub const rhi::RHI_STATE_SCHEMA_VERSION_6_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
diff --git a/contracts/services_hardening/reconciliation_commit.v1.json b/contracts/services_hardening/reconciliation_commit.v1.json
@@ -0,0 +1,101 @@
+{
+ "schema": "radroots.rhi.reconciliation-commit",
+ "schema_version": 1,
+ "contract_version": 1,
+ "state_schema_version": 6,
+ "authority": "sealed_reconciliation_attempt_repository",
+ "input": {
+ "lease": "exact_live_compare_and_swap_lease",
+ "attempt": "exact_claim_derived_attempt_plan",
+ "source_inventory": "exact_canonical_request_order_bounded_to_configured_count_plus_one",
+ "events": "signature_verified_admitted_replay_facts_only",
+ "time": "injected_source_result_milliseconds_only"
+ },
+ "transaction": {
+ "mode": "one_short_sqlx_transaction",
+ "fences_before_mutation": [
+ "exact_live_lease_row",
+ "exact_trade_dirty_generation",
+ "exact_evidence_policy_digest",
+ "exact_per_source_prior_checkpoint"
+ ],
+ "atomic_writes": [
+ "canonical_mutations",
+ "signed_events",
+ "source_observations",
+ "immutable_attempt",
+ "immutable_source_results",
+ "exact_source_inventory_digests",
+ "at_most_one_dirty_generation_advance",
+ "eligible_source_checkpoints"
+ ],
+ "source_or_network_wait": false
+ },
+ "checkpoint": {
+ "scope": [
+ "source_id",
+ "trade_id",
+ "evidence_policy_digest",
+ "selector_digest"
+ ],
+ "advance_only_when": "completion_is_complete_and_candidate_is_strictly_newer_than_prior",
+ "forbidden_completion_classes": [
+ "incomplete_timeout",
+ "incomplete_unavailable",
+ "incomplete_resource_limit",
+ "incomplete_unknown",
+ "unsupported"
+ ],
+ "public_evidence": "sealed_and_minted_only_after_durable_commit_confirmation"
+ },
+ "dirty_generation": {
+ "advance_only_for": "newly_inserted_canonical_mutation_or_signed_event",
+ "advance_count_maximum_per_attempt": 1,
+ "observation_only_or_exact_replay": "no_advance"
+ },
+ "idempotency": {
+ "key": "attempt_id",
+ "source_inventory_identity": "domain_separated_sha256_over_exact_ordered_persisted_facts_and_provenance",
+ "exact_committed_inventory": "reconciled_as_created_false",
+ "mismatch": "conflict",
+ "lost_success_retry": "checked_before_live_lease_and_generation_fences"
+ },
+ "schema_objects": {
+ "tables": [
+ "evidence_reconciliations",
+ "evidence_reconciliation_sources"
+ ],
+ "immutability": "no_update_and_no_delete_triggers",
+ "source_count_maximum": 16
+ },
+ "diagnostics": "crate_owned_source_free_redacted",
+ "effects": {
+ "sqlite": true,
+ "filesystem": false,
+ "source_or_relay": false,
+ "network": false,
+ "task_spawn": false,
+ "ambient_clock": false,
+ "ambient_entropy": false
+ },
+ "forbidden": [
+ "caller_forged_cursor_evidence",
+ "raw_or_unverified_event_persistence",
+ "partial_attempt_commit",
+ "cursor_advance_for_incomplete_or_unsupported_result",
+ "cursor_advance_before_evidence_commit",
+ "dirty_advance_for_observation_only",
+ "stale_lease_commit",
+ "stale_generation_commit",
+ "stale_checkpoint_commit",
+ "unbounded_result_iterator",
+ "raw_sqlx_authority_escape"
+ ],
+ "deferred": [
+ "evidence_manifest",
+ "lineage_reducer",
+ "attestation",
+ "publication",
+ "job_finalization"
+ ]
+}
diff --git a/src/lib.rs b/src/lib.rs
@@ -9,6 +9,7 @@ mod features;
mod identity_credential;
mod identity_envelope;
mod reconciliation_attempt;
+mod reconciliation_commit;
mod reconciliation_job;
mod reconciliation_replay;
mod runtime_adapters;
@@ -76,6 +77,10 @@ pub use reconciliation_attempt::{
RhiReconciliationSourceRequestId, RhiReconciliationSourceResult,
RhiReconciliationSourceSelectorDigest,
};
+pub use reconciliation_commit::{
+ RHI_RECONCILIATION_COMMIT_CONTRACT_VERSION, RhiReconciliationCommitError,
+ RhiReconciliationCommitErrorKind, RhiReconciliationSourceCommitOutcome,
+};
pub use reconciliation_job::{
RHI_RECONCILIATION_JOB_CONTRACT_VERSION, RHI_RECONCILIATION_JOB_MAX_ACTIVE,
RhiReconciliationJob, RhiReconciliationJobError, RhiReconciliationJobErrorKind,
@@ -120,8 +125,10 @@ pub use state_catalog::{
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,
RHI_STATE_SCHEMA_VERSION_5_MIGRATION_SHA256, RHI_STATE_SCHEMA_VERSION_5_OBJECT_COUNT,
- RHI_STATE_SCHEMA_VERSION_5_SHA256, RhiStateCatalogError, RhiStateCatalogErrorKind,
- rhi_migration_catalog, rhi_schema_catalog, validate_rhi_state_catalogs,
+ RHI_STATE_SCHEMA_VERSION_5_SHA256, RHI_STATE_SCHEMA_VERSION_6_MIGRATION_SHA256,
+ RHI_STATE_SCHEMA_VERSION_6_OBJECT_COUNT, RHI_STATE_SCHEMA_VERSION_6_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/reconciliation_commit.rs b/src/reconciliation_commit.rs
@@ -0,0 +1,803 @@
+//! Atomic durable reconciliation-source result commit.
+
+use core::{cmp::Ordering, fmt};
+use std::error::Error;
+
+use radroots_event::id::TradeId;
+use radroots_service_sqlite::{
+ ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind,
+};
+use sha2::{Digest, Sha256};
+use sqlx::Row;
+
+use crate::{
+ RhiReconciliationAttemptPlan, RhiReconciliationAttemptRepository, RhiReconciliationLease,
+ RhiReconciliationSourceCursorEvidence, RhiReconciliationSourceReplay, RhiStateHostMode,
+ RhiTradeSourceCursor,
+ reconciliation_job::{LeaseValidationError, validate_exact_lease},
+ reconciliation_replay::{
+ RhiReconciliationReplayCommitFact, RhiReconciliationReplayCommitParts,
+ committed_cursor_evidence,
+ },
+ source_ingest::{
+ SourceOperationError, advance_dirty_generation, compare_cursor, read_checkpoint,
+ read_dirty, write_checkpoint,
+ },
+ state_trade::{PersistenceOperationError, RhiTradeSourceObservation, persist},
+};
+
+/// Exact version of the atomic reconciliation-source commit contract.
+pub const RHI_RECONCILIATION_COMMIT_CONTRACT_VERSION: u32 = 1;
+
+const SOURCE_INVENTORY_DIGEST_DOMAIN: &[u8] = b"radroots.rhi.reconciliation_source_inventory.v1\0";
+
+const COUNT_ATTEMPT_SQL: &str =
+ "SELECT COUNT(*) AS row_count FROM evidence_reconciliations WHERE attempt_id = ?";
+const MATCH_ATTEMPT_SQL: &str = r#"SELECT COUNT(*) AS row_count
+FROM evidence_reconciliations
+WHERE attempt_id = ? AND job_id = ? AND trade_id = ? AND input_generation = ?
+ AND evidence_policy_sha256 = ? AND attempt_started_unix_ms = ?
+ AND deadline_unix_ms = ? AND source_count = ?"#;
+const COUNT_SOURCE_RESULTS_SQL: &str = r#"SELECT COUNT(*) AS row_count
+FROM evidence_reconciliation_sources WHERE attempt_id = ?"#;
+const MATCH_SOURCE_RESULT_SQL: &str = r#"SELECT COUNT(*) AS row_count
+FROM evidence_reconciliation_sources
+WHERE attempt_id = ? AND request_id = ? AND source_ordinal = ?
+ AND source_id = ? AND trade_id = ? AND required = ? AND selector_sha256 = ?
+ AND replay_id = ? AND completion = ? AND started_unix_ms = ? AND finished_unix_ms = ?
+ AND accepted_event_count = ? AND accepted_event_bytes = ?
+ AND accepted_inventory_sha256 = ?
+ AND duplicate_observation_count = ? AND first_observed_unix_s IS ?
+ AND prior_cursor_created_at_unix_s IS ? AND prior_cursor_event_id IS ?
+ AND overlap_seconds = ? AND inclusive_since_unix_s = ?
+ AND candidate_created_at_unix_s IS ? AND candidate_event_id IS ?
+ AND checkpoint_advanced = ?"#;
+const INSERT_ATTEMPT_SQL: &str = r#"INSERT INTO evidence_reconciliations (
+ attempt_id, job_id, trade_id, input_generation, evidence_policy_sha256,
+ attempt_started_unix_ms, deadline_unix_ms, source_count
+) VALUES (?, ?, ?, ?, ?, ?, ?, ?)"#;
+const INSERT_SOURCE_RESULT_SQL: &str = r#"INSERT INTO evidence_reconciliation_sources (
+ attempt_id, request_id, source_ordinal, source_id, trade_id, required,
+ selector_sha256, replay_id, completion, started_unix_ms, finished_unix_ms,
+ accepted_event_count, accepted_event_bytes, accepted_inventory_sha256,
+ duplicate_observation_count,
+ first_observed_unix_s, prior_cursor_created_at_unix_s, prior_cursor_event_id,
+ overlap_seconds, inclusive_since_unix_s, candidate_created_at_unix_s,
+ candidate_event_id, checkpoint_advanced
+) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)"#;
+
+/// Stable source-free failure classification for one atomic result commit.
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+pub enum RhiReconciliationCommitErrorKind {
+ InvalidMode,
+ InvalidInput,
+ LeaseLost,
+ GenerationConflict,
+ Conflict,
+ Storage,
+ CommitOutcomeUnknown,
+}
+
+impl RhiReconciliationCommitErrorKind {
+ /// Returns the stable machine-readable failure code.
+ #[must_use]
+ pub const fn code(self) -> &'static str {
+ match self {
+ Self::InvalidMode => "reconciliation_commit_mode_invalid",
+ Self::InvalidInput => "reconciliation_commit_input_invalid",
+ Self::LeaseLost => "reconciliation_commit_lease_lost",
+ Self::GenerationConflict => "reconciliation_commit_generation_conflict",
+ Self::Conflict => "reconciliation_commit_conflict",
+ Self::Storage => "reconciliation_commit_storage_failed",
+ Self::CommitOutcomeUnknown => "reconciliation_commit_outcome_unknown",
+ }
+ }
+}
+
+/// Redacted source-free atomic reconciliation commit failure.
+#[derive(Clone, Copy, PartialEq, Eq)]
+pub struct RhiReconciliationCommitError {
+ kind: RhiReconciliationCommitErrorKind,
+}
+
+impl RhiReconciliationCommitError {
+ /// Returns the stable failure class.
+ #[must_use]
+ pub const fn kind(self) -> RhiReconciliationCommitErrorKind {
+ self.kind
+ }
+
+ /// Returns the stable machine-readable failure code.
+ #[must_use]
+ pub const fn code(self) -> &'static str {
+ self.kind.code()
+ }
+}
+
+impl fmt::Display for RhiReconciliationCommitError {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter.write_str(match self.kind {
+ RhiReconciliationCommitErrorKind::InvalidMode => {
+ "RHI reconciliation commit requires writable state"
+ }
+ RhiReconciliationCommitErrorKind::InvalidInput => {
+ "RHI reconciliation commit input is invalid"
+ }
+ RhiReconciliationCommitErrorKind::LeaseLost => {
+ "RHI reconciliation lease is no longer authoritative"
+ }
+ RhiReconciliationCommitErrorKind::GenerationConflict => {
+ "RHI reconciliation generation or cursor changed"
+ }
+ RhiReconciliationCommitErrorKind::Conflict => {
+ "RHI reconciliation evidence conflicts with durable state"
+ }
+ RhiReconciliationCommitErrorKind::Storage => {
+ "RHI reconciliation commit transaction failed"
+ }
+ RhiReconciliationCommitErrorKind::CommitOutcomeUnknown => {
+ "RHI reconciliation commit outcome is unknown"
+ }
+ })
+ }
+}
+
+impl fmt::Debug for RhiReconciliationCommitError {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter
+ .debug_struct("RhiReconciliationCommitError")
+ .field("kind", &self.kind)
+ .finish()
+ }
+}
+
+impl Error for RhiReconciliationCommitError {}
+
+/// Durable result of one exact bounded reconciliation-attempt commit.
+pub struct RhiReconciliationSourceCommitOutcome {
+ created: bool,
+ source_result_count: u32,
+ checkpoint_advance_count: u32,
+ dirty_generation_advanced: bool,
+ committed_cursors: Box<[RhiReconciliationSourceCursorEvidence]>,
+}
+
+impl RhiReconciliationSourceCommitOutcome {
+ /// Reports whether this call created the immutable attempt inventory.
+ #[must_use]
+ pub const fn created(&self) -> bool {
+ self.created
+ }
+
+ /// Returns the exact committed source-result count.
+ #[must_use]
+ pub const fn source_result_count(&self) -> u32 {
+ self.source_result_count
+ }
+
+ /// Returns the exact cursor checkpoint-advance count.
+ #[must_use]
+ pub const fn checkpoint_advance_count(&self) -> u32 {
+ self.checkpoint_advance_count
+ }
+
+ /// Reports whether newly durable mutation/event evidence advanced dirty state once.
+ #[must_use]
+ pub const fn dirty_generation_advanced(&self) -> bool {
+ self.dirty_generation_advanced
+ }
+
+ /// Returns sealed cursor evidence minted only after durable commit confirmation.
+ #[must_use]
+ pub fn committed_cursors(&self) -> &[RhiReconciliationSourceCursorEvidence] {
+ &self.committed_cursors
+ }
+}
+
+impl fmt::Debug for RhiReconciliationSourceCommitOutcome {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter
+ .debug_struct("RhiReconciliationSourceCommitOutcome")
+ .field("created", &self.created)
+ .field("source_result_count", &self.source_result_count)
+ .field("checkpoint_advance_count", &self.checkpoint_advance_count)
+ .field("dirty_generation_advanced", &self.dirty_generation_advanced)
+ .field("committed_cursors", &self.committed_cursors.len())
+ .finish()
+ }
+}
+
+impl RhiReconciliationAttemptRepository<'_> {
+ /// Atomically commits one exact bounded source-result inventory.
+ pub async fn commit_source_replays<I>(
+ &self,
+ lease: RhiReconciliationLease,
+ plan: RhiReconciliationAttemptPlan,
+ replays: I,
+ ) -> Result<RhiReconciliationSourceCommitOutcome, RhiReconciliationCommitError>
+ where
+ I: IntoIterator<Item = RhiReconciliationSourceReplay>,
+ {
+ if self.host().mode() != RhiStateHostMode::ReadWriteExisting {
+ return Err(failure(RhiReconciliationCommitErrorKind::InvalidMode));
+ }
+ let replays = replays
+ .into_iter()
+ .take(plan.requests().len().saturating_add(1))
+ .collect::<Vec<_>>();
+ if replays.len() != plan.requests().len()
+ || plan.job_id() != lease.job().id()
+ || plan.input_generation() != lease.job().input_generation()
+ || plan.evidence_policy_digest() != lease.job().evidence_policy_digest()
+ || plan.deadline() > lease.lease_expires()
+ {
+ return Err(failure(RhiReconciliationCommitErrorKind::InvalidInput));
+ }
+ let parts = replays
+ .into_iter()
+ .map(RhiReconciliationSourceReplay::into_commit_parts)
+ .collect::<Vec<_>>();
+ if parts.iter().zip(plan.requests()).any(|(replay, request)| {
+ replay.request_id != request.id()
+ || replay.result.request_id() != request.id()
+ || replay.source_id.as_ref() != request.source_id()
+ || replay.trade_id != request.trade_id()
+ || replay.required != request.required()
+ || replay.trade_id != lease.job().trade_id()
+ || replay.policy_digest != *plan.evidence_policy_digest().as_bytes()
+ || replay.selector_digest != *request.selector_digest().as_bytes()
+ || replay.result.finished_at() >= lease.lease_expires()
+ || replay.facts.len() != replay.result.accepted_event_count() as usize
+ || replay
+ .eligible_cursor
+ .is_some_and(|_| replay.result.finished_at().get() / 1_000 == 0)
+ }) {
+ return Err(failure(RhiReconciliationCommitErrorKind::InvalidInput));
+ }
+
+ let raw = self
+ .host()
+ .sqlite_host()
+ .transaction(move |transaction| {
+ Box::pin(async move { commit(transaction, lease, plan, parts).await })
+ })
+ .await
+ .map_err(map_transaction_error)?;
+ let committed_cursors = raw
+ .cursor_scopes
+ .into_vec()
+ .into_iter()
+ .map(|scope| {
+ committed_cursor_evidence(
+ scope.source_id,
+ scope.trade_id,
+ scope.policy_digest,
+ scope.selector_digest,
+ scope.cursor,
+ )
+ })
+ .collect::<Vec<_>>()
+ .into_boxed_slice();
+ Ok(RhiReconciliationSourceCommitOutcome {
+ created: raw.created,
+ source_result_count: raw.source_result_count,
+ checkpoint_advance_count: raw.checkpoint_advance_count,
+ dirty_generation_advanced: raw.dirty_generation_advanced,
+ committed_cursors,
+ })
+ }
+}
+
+struct RawCommitOutcome {
+ created: bool,
+ source_result_count: u32,
+ checkpoint_advance_count: u32,
+ dirty_generation_advanced: bool,
+ cursor_scopes: Box<[CommittedCursorScope]>,
+}
+
+struct CommittedCursorScope {
+ source_id: Box<str>,
+ trade_id: TradeId,
+ policy_digest: [u8; 32],
+ selector_digest: [u8; 32],
+ cursor: RhiTradeSourceCursor,
+}
+
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+enum CommitOperationError {
+ LeaseLost,
+ GenerationConflict,
+ Conflict,
+ Storage,
+}
+
+async fn commit(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ lease: RhiReconciliationLease,
+ plan: RhiReconciliationAttemptPlan,
+ parts: Vec<RhiReconciliationReplayCommitParts>,
+) -> Result<RawCommitOutcome, CommitOperationError> {
+ if reconcile_existing(transaction, &plan, &parts).await? {
+ return raw_outcome(false, false, &parts);
+ }
+ validate_exact_lease(transaction, lease)
+ .await
+ .map_err(|error| match error {
+ LeaseValidationError::LeaseLost => CommitOperationError::LeaseLost,
+ LeaseValidationError::Storage => CommitOperationError::Storage,
+ })?;
+ let policy = plan.evidence_policy_digest();
+ let dirty = read_dirty(transaction, lease.job().trade_id())
+ .await
+ .map_err(map_source_error)?
+ .filter(|state| state.generation.get() == plan.input_generation() && state.policy == policy)
+ .ok_or(CommitOperationError::GenerationConflict)?;
+ let mut checkpoints = Vec::with_capacity(parts.len());
+ for part in &parts {
+ let checkpoint =
+ read_checkpoint(transaction, part.source_id.as_ref(), policy, part.trade_id)
+ .await
+ .map_err(map_source_error)?;
+ if checkpoint.map(|value| value.cursor) != part.prior_cursor {
+ return Err(CommitOperationError::GenerationConflict);
+ }
+ checkpoints.push(checkpoint);
+ }
+
+ insert_attempt(transaction, &plan).await?;
+ let mut new_relevant_evidence = false;
+ for part in &parts {
+ for fact in &part.facts {
+ let observation = RhiTradeSourceObservation::from_persistence_parts(
+ part.source_id.clone(),
+ policy,
+ &fact.record,
+ fact.observed_at,
+ );
+ let outcome = persist(transaction, &fact.record, &observation)
+ .await
+ .map_err(map_persistence_error)?;
+ new_relevant_evidence |= outcome.mutation_inserted() || outcome.signed_event_inserted();
+ }
+ }
+
+ if new_relevant_evidence {
+ let updated_at = parts
+ .iter()
+ .map(|part| part.result.finished_at().get() / 1_000)
+ .max()
+ .ok_or(CommitOperationError::Storage)?;
+ advance_dirty_generation(
+ transaction,
+ lease.job().trade_id(),
+ policy,
+ updated_at,
+ Some(dirty),
+ )
+ .await
+ .map_err(map_source_error)?;
+ }
+
+ for ((ordinal, part), checkpoint) in parts.iter().enumerate().zip(checkpoints) {
+ if let Some(candidate) = part.eligible_cursor {
+ write_checkpoint(
+ transaction,
+ part.source_id.as_ref(),
+ policy,
+ part.trade_id,
+ checkpoint,
+ candidate,
+ part.result.finished_at().get() / 1_000,
+ )
+ .await
+ .map_err(map_source_error)?;
+ }
+ insert_source_result(transaction, plan.id().as_bytes(), ordinal, part).await?;
+ }
+ raw_outcome(true, new_relevant_evidence, &parts)
+}
+
+fn raw_outcome(
+ created: bool,
+ dirty_generation_advanced: bool,
+ parts: &[RhiReconciliationReplayCommitParts],
+) -> Result<RawCommitOutcome, CommitOperationError> {
+ let cursor_scopes = parts
+ .iter()
+ .filter_map(|part| {
+ part.eligible_cursor.map(|cursor| CommittedCursorScope {
+ source_id: part.source_id.clone(),
+ trade_id: part.trade_id,
+ policy_digest: part.policy_digest,
+ selector_digest: part.selector_digest,
+ cursor,
+ })
+ })
+ .collect::<Vec<_>>()
+ .into_boxed_slice();
+ Ok(RawCommitOutcome {
+ created,
+ source_result_count: u32::try_from(parts.len())
+ .map_err(|_| CommitOperationError::Storage)?,
+ checkpoint_advance_count: u32::try_from(cursor_scopes.len())
+ .map_err(|_| CommitOperationError::Storage)?,
+ dirty_generation_advanced,
+ cursor_scopes,
+ })
+}
+
+async fn reconcile_existing(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ plan: &RhiReconciliationAttemptPlan,
+ parts: &[RhiReconciliationReplayCommitParts],
+) -> Result<bool, CommitOperationError> {
+ let count = count_query(transaction, COUNT_ATTEMPT_SQL, plan.id().as_bytes()).await?;
+ if count == 0 {
+ return Ok(false);
+ }
+ if count != 1 || match_attempt(transaction, plan).await? != 1 {
+ return Err(CommitOperationError::Conflict);
+ }
+ let source_count =
+ count_query(transaction, COUNT_SOURCE_RESULTS_SQL, plan.id().as_bytes()).await?;
+ if source_count != parts.len() as u64 {
+ return Err(CommitOperationError::Conflict);
+ }
+ for (ordinal, part) in parts.iter().enumerate() {
+ if match_source_result(transaction, plan.id().as_bytes(), ordinal, part).await? != 1 {
+ return Err(CommitOperationError::Conflict);
+ }
+ if let Some(candidate) = part.eligible_cursor {
+ let checkpoint = read_checkpoint(
+ transaction,
+ part.source_id.as_ref(),
+ plan.evidence_policy_digest(),
+ part.trade_id,
+ )
+ .await
+ .map_err(map_source_error)?
+ .ok_or(CommitOperationError::Conflict)?;
+ if compare_cursor(checkpoint.cursor, candidate) == Ordering::Less {
+ return Err(CommitOperationError::Conflict);
+ }
+ }
+ }
+ Ok(true)
+}
+
+async fn count_query(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ sql: &'static str,
+ id: &[u8; 32],
+) -> Result<u64, CommitOperationError> {
+ sqlx::query(sql)
+ .bind(id.as_slice())
+ .fetch_one(&mut *transaction)
+ .await
+ .map_err(|_| CommitOperationError::Storage)?
+ .try_get::<i64, _>("row_count")
+ .ok()
+ .and_then(|value| u64::try_from(value).ok())
+ .ok_or(CommitOperationError::Storage)
+}
+
+async fn match_attempt(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ plan: &RhiReconciliationAttemptPlan,
+) -> Result<u64, CommitOperationError> {
+ let trade = plan
+ .requests()
+ .first()
+ .ok_or(CommitOperationError::Storage)?
+ .trade_id();
+ count_row(
+ sqlx::query(MATCH_ATTEMPT_SQL)
+ .bind(plan.id().as_bytes().as_slice())
+ .bind(plan.job_id().as_bytes().as_slice())
+ .bind(trade.as_bytes().as_slice())
+ .bind(i64_value(plan.input_generation())?)
+ .bind(plan.evidence_policy_digest().as_bytes().as_slice())
+ .bind(i64_value(plan.attempt_started_at().get())?)
+ .bind(i64_value(plan.deadline().get())?)
+ .bind(i64::try_from(plan.requests().len()).map_err(|_| CommitOperationError::Storage)?),
+ transaction,
+ )
+ .await
+}
+
+async fn match_source_result(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ attempt_id: &[u8; 32],
+ ordinal: usize,
+ part: &RhiReconciliationReplayCommitParts,
+) -> Result<u64, CommitOperationError> {
+ let prior_created = part
+ .prior_cursor
+ .map(|cursor| i64_value(cursor.created_at_unix_seconds()))
+ .transpose()?;
+ let prior_event = part.prior_cursor.map(|cursor| cursor.event_id());
+ let candidate_created = part
+ .cursor_candidate
+ .map(|cursor| i64_value(cursor.created_at_unix_seconds()))
+ .transpose()?;
+ let candidate_event = part.cursor_candidate.map(|cursor| cursor.event_id());
+ let inventory_digest = accepted_inventory_digest(part)?;
+ count_row(
+ sqlx::query(MATCH_SOURCE_RESULT_SQL)
+ .bind(attempt_id.as_slice())
+ .bind(part.request_id.as_bytes().as_slice())
+ .bind(i64::try_from(ordinal).map_err(|_| CommitOperationError::Storage)?)
+ .bind(part.source_id.as_ref())
+ .bind(part.trade_id.as_bytes().as_slice())
+ .bind(i64::from(part.required))
+ .bind(part.selector_digest.as_slice())
+ .bind(part.replay_id.as_bytes().as_slice())
+ .bind(part.result.outcome().code())
+ .bind(i64_value(part.result.started_at().get())?)
+ .bind(i64_value(part.result.finished_at().get())?)
+ .bind(i64::from(part.result.accepted_event_count()))
+ .bind(i64_value(part.result.accepted_event_bytes())?)
+ .bind(inventory_digest.as_slice())
+ .bind(i64::from(part.duplicate_observations))
+ .bind(
+ part.first_observed_at
+ .map(|value| i64_value(value.get()))
+ .transpose()?,
+ )
+ .bind(prior_created)
+ .bind(prior_event.as_ref().map(<[u8; 32]>::as_slice))
+ .bind(i64_value(part.overlap_seconds)?)
+ .bind(i64_value(part.since_unix_seconds)?)
+ .bind(candidate_created)
+ .bind(candidate_event.as_ref().map(<[u8; 32]>::as_slice))
+ .bind(i64::from(part.eligible_cursor.is_some())),
+ transaction,
+ )
+ .await
+}
+
+async fn insert_attempt(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ plan: &RhiReconciliationAttemptPlan,
+) -> Result<(), CommitOperationError> {
+ let trade = plan
+ .requests()
+ .first()
+ .ok_or(CommitOperationError::Storage)?
+ .trade_id();
+ let result = sqlx::query(INSERT_ATTEMPT_SQL)
+ .bind(plan.id().as_bytes().as_slice())
+ .bind(plan.job_id().as_bytes().as_slice())
+ .bind(trade.as_bytes().as_slice())
+ .bind(i64_value(plan.input_generation())?)
+ .bind(plan.evidence_policy_digest().as_bytes().as_slice())
+ .bind(i64_value(plan.attempt_started_at().get())?)
+ .bind(i64_value(plan.deadline().get())?)
+ .bind(i64::try_from(plan.requests().len()).map_err(|_| CommitOperationError::Storage)?)
+ .execute(&mut *transaction)
+ .await
+ .map_err(|_| CommitOperationError::Storage)?;
+ if result.rows_affected() == 1 {
+ Ok(())
+ } else {
+ Err(CommitOperationError::Storage)
+ }
+}
+
+async fn insert_source_result(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ attempt_id: &[u8; 32],
+ ordinal: usize,
+ part: &RhiReconciliationReplayCommitParts,
+) -> Result<(), CommitOperationError> {
+ let prior_created = part
+ .prior_cursor
+ .map(|cursor| i64_value(cursor.created_at_unix_seconds()))
+ .transpose()?;
+ let prior_event = part.prior_cursor.map(|cursor| cursor.event_id());
+ let candidate_created = part
+ .cursor_candidate
+ .map(|cursor| i64_value(cursor.created_at_unix_seconds()))
+ .transpose()?;
+ let candidate_event = part.cursor_candidate.map(|cursor| cursor.event_id());
+ let inventory_digest = accepted_inventory_digest(part)?;
+ let result = sqlx::query(INSERT_SOURCE_RESULT_SQL)
+ .bind(attempt_id.as_slice())
+ .bind(part.request_id.as_bytes().as_slice())
+ .bind(i64::try_from(ordinal).map_err(|_| CommitOperationError::Storage)?)
+ .bind(part.source_id.as_ref())
+ .bind(part.trade_id.as_bytes().as_slice())
+ .bind(i64::from(part.required))
+ .bind(part.selector_digest.as_slice())
+ .bind(part.replay_id.as_bytes().as_slice())
+ .bind(part.result.outcome().code())
+ .bind(i64_value(part.result.started_at().get())?)
+ .bind(i64_value(part.result.finished_at().get())?)
+ .bind(i64::from(part.result.accepted_event_count()))
+ .bind(i64_value(part.result.accepted_event_bytes())?)
+ .bind(inventory_digest.as_slice())
+ .bind(i64::from(part.duplicate_observations))
+ .bind(
+ part.first_observed_at
+ .map(|value| i64_value(value.get()))
+ .transpose()?,
+ )
+ .bind(prior_created)
+ .bind(prior_event.as_ref().map(<[u8; 32]>::as_slice))
+ .bind(i64_value(part.overlap_seconds)?)
+ .bind(i64_value(part.since_unix_seconds)?)
+ .bind(candidate_created)
+ .bind(candidate_event.as_ref().map(<[u8; 32]>::as_slice))
+ .bind(i64::from(part.eligible_cursor.is_some()))
+ .execute(&mut *transaction)
+ .await
+ .map_err(|_| CommitOperationError::Storage)?;
+ if result.rows_affected() == 1 {
+ Ok(())
+ } else {
+ Err(CommitOperationError::Storage)
+ }
+}
+
+fn accepted_inventory_digest(
+ part: &RhiReconciliationReplayCommitParts,
+) -> Result<[u8; 32], CommitOperationError> {
+ accepted_inventory_digest_for_facts(&part.facts)
+}
+
+fn accepted_inventory_digest_for_facts(
+ facts: &[RhiReconciliationReplayCommitFact],
+) -> Result<[u8; 32], CommitOperationError> {
+ let mut digest = Sha256::new();
+ digest.update(SOURCE_INVENTORY_DIGEST_DOMAIN);
+ digest.update(
+ u32::try_from(facts.len())
+ .map_err(|_| CommitOperationError::Storage)?
+ .to_be_bytes(),
+ );
+ for fact in facts {
+ let record = &fact.record;
+ digest.update(record.mutation_id);
+ digest.update(record.trade_id);
+ update_framed(&mut digest, record.contract_id.as_bytes())?;
+ digest.update(record.schema_version.to_be_bytes());
+ digest.update(record.event_id);
+ digest.update(record.event_signature);
+ digest.update(record.author_pubkey);
+ digest.update(record.event_kind.to_be_bytes());
+ digest.update(record.authored_at_unix_s.to_be_bytes());
+ update_framed(&mut digest, &record.canonical_content)?;
+ update_framed(&mut digest, &record.canonical_event_json)?;
+ digest.update(fact.observed_at.get().to_be_bytes());
+ }
+ Ok(digest.finalize().into())
+}
+
+fn update_framed(digest: &mut Sha256, bytes: &[u8]) -> Result<(), CommitOperationError> {
+ digest.update(
+ u64::try_from(bytes.len())
+ .map_err(|_| CommitOperationError::Storage)?
+ .to_be_bytes(),
+ );
+ digest.update(bytes);
+ Ok(())
+}
+
+async fn count_row<'query>(
+ query: sqlx::query::Query<'query, sqlx::Sqlite, sqlx::sqlite::SqliteArguments>,
+ transaction: &mut ServiceSqliteTransaction<'_>,
+) -> Result<u64, CommitOperationError> {
+ query
+ .fetch_one(&mut *transaction)
+ .await
+ .map_err(|_| CommitOperationError::Storage)?
+ .try_get::<i64, _>("row_count")
+ .ok()
+ .and_then(|value| u64::try_from(value).ok())
+ .ok_or(CommitOperationError::Storage)
+}
+
+fn i64_value(value: u64) -> Result<i64, CommitOperationError> {
+ i64::try_from(value).map_err(|_| CommitOperationError::Storage)
+}
+
+fn map_persistence_error(error: PersistenceOperationError) -> CommitOperationError {
+ match error {
+ PersistenceOperationError::MutationConflict
+ | PersistenceOperationError::SignedEventConflict => CommitOperationError::Conflict,
+ PersistenceOperationError::Storage => CommitOperationError::Storage,
+ }
+}
+
+fn map_source_error(error: SourceOperationError) -> CommitOperationError {
+ match error {
+ SourceOperationError::GenerationConflict => CommitOperationError::GenerationConflict,
+ SourceOperationError::Persistence(error) => map_persistence_error(error),
+ SourceOperationError::Storage => CommitOperationError::Storage,
+ }
+}
+
+fn map_transaction_error(
+ error: ServiceSqliteTransactionError<CommitOperationError>,
+) -> RhiReconciliationCommitError {
+ if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown {
+ return failure(RhiReconciliationCommitErrorKind::CommitOutcomeUnknown);
+ }
+ failure(match error.operation_error().copied() {
+ Some(CommitOperationError::LeaseLost) => RhiReconciliationCommitErrorKind::LeaseLost,
+ Some(CommitOperationError::GenerationConflict) => {
+ RhiReconciliationCommitErrorKind::GenerationConflict
+ }
+ Some(CommitOperationError::Conflict) => RhiReconciliationCommitErrorKind::Conflict,
+ Some(CommitOperationError::Storage) | None => RhiReconciliationCommitErrorKind::Storage,
+ })
+}
+
+const fn failure(kind: RhiReconciliationCommitErrorKind) -> RhiReconciliationCommitError {
+ RhiReconciliationCommitError { kind }
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+
+ fn fact(observed_at: u64, content: &[u8]) -> RhiReconciliationReplayCommitFact {
+ RhiReconciliationReplayCommitFact {
+ record: crate::state_trade::PersistenceRecord {
+ mutation_id: [0x11; 32],
+ trade_id: [0x22; 16],
+ contract_id: "radroots.trade.proposal.v1",
+ schema_version: 1,
+ event_id: [0x33; 32],
+ event_signature: [0x44; 64],
+ author_pubkey: [0x55; 32],
+ event_kind: 3_470,
+ authored_at_unix_s: 1_784_347_200,
+ canonical_content: content.into(),
+ canonical_event_json: br#"{"id":"example"}"#.as_slice().into(),
+ },
+ observed_at: crate::RhiTradeMutationObservedAtUnixSeconds::new(observed_at)
+ .expect("observation"),
+ }
+ }
+
+ #[test]
+ fn source_inventory_digest_is_exact_framed_and_provenance_sensitive() {
+ assert_eq!(
+ accepted_inventory_digest_for_facts(&[]).expect("empty digest"),
+ [
+ 0xc6, 0x77, 0x8e, 0xb5, 0x38, 0xe6, 0x1a, 0x6d, 0xa4, 0xe5, 0xba, 0xe1, 0x9c, 0xb1,
+ 0xb3, 0xf9, 0xcd, 0x14, 0x39, 0x8b, 0x7d, 0xc2, 0xfb, 0x26, 0x3c, 0xa0, 0x86, 0xf2,
+ 0xf5, 0xf3, 0x0c, 0xdb,
+ ]
+ );
+ let first = accepted_inventory_digest_for_facts(&[fact(1_784_347_201, b"alpha")])
+ .expect("first digest");
+ let changed_provenance =
+ accepted_inventory_digest_for_facts(&[fact(1_784_347_202, b"alpha")])
+ .expect("provenance digest");
+ let changed_content = accepted_inventory_digest_for_facts(&[fact(1_784_347_201, b"bravo")])
+ .expect("content digest");
+ assert_ne!(first, changed_provenance);
+ assert_ne!(first, changed_content);
+ }
+
+ #[test]
+ fn errors_are_closed_redacted_and_source_free() {
+ for kind in [
+ RhiReconciliationCommitErrorKind::InvalidMode,
+ RhiReconciliationCommitErrorKind::InvalidInput,
+ RhiReconciliationCommitErrorKind::LeaseLost,
+ RhiReconciliationCommitErrorKind::GenerationConflict,
+ RhiReconciliationCommitErrorKind::Conflict,
+ RhiReconciliationCommitErrorKind::Storage,
+ RhiReconciliationCommitErrorKind::CommitOutcomeUnknown,
+ ] {
+ let error = failure(kind);
+ assert_eq!(error.kind(), kind);
+ assert!(Error::source(&error).is_none());
+ assert!(!format!("{error} {error:?}").contains("trade-primary"));
+ }
+ }
+}
diff --git a/src/reconciliation_job.rs b/src/reconciliation_job.rs
@@ -601,6 +601,29 @@ impl RhiReconciliationLease {
}
}
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+pub(crate) enum LeaseValidationError {
+ LeaseLost,
+ Storage,
+}
+
+pub(crate) async fn validate_exact_lease(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ lease: RhiReconciliationLease,
+) -> Result<(), LeaseValidationError> {
+ let current = read_job(transaction, lease.job.id)
+ .await
+ .map_err(|_| LeaseValidationError::Storage)?;
+ if current == Some(lease.job)
+ && lease.job.state == RhiReconciliationJobState::Leased
+ && lease.job.lease_expires == Some(lease.lease_expires)
+ {
+ Ok(())
+ } else {
+ Err(LeaseValidationError::LeaseLost)
+ }
+}
+
impl fmt::Debug for RhiReconciliationLease {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
diff --git a/src/reconciliation_replay.rs b/src/reconciliation_replay.rs
@@ -158,6 +158,11 @@ impl fmt::Debug for RhiReconciliationSourceCursorEvidence {
pub struct RhiReconciliationSourceReplayPlan {
id: RhiReconciliationSourceReplayId,
request_id: RhiReconciliationSourceRequestId,
+ source_id: Box<str>,
+ trade_id: radroots_event::id::TradeId,
+ required: bool,
+ policy_digest: [u8; 32],
+ selector_digest: [u8; 32],
prior_cursor: Option<RhiReconciliationSourceCursorEvidence>,
overlap_seconds: u64,
since_unix_seconds: u64,
@@ -244,6 +249,11 @@ impl RhiReconciliationSourceReplayPlan {
since_unix_seconds,
),
request_id: request.id(),
+ source_id: request.source_id().into(),
+ trade_id: request.trade_id(),
+ required: request.required(),
+ policy_digest: *policy.as_bytes(),
+ selector_digest: *request.selector_digest().as_bytes(),
prior_cursor,
overlap_seconds,
since_unix_seconds,
@@ -437,6 +447,77 @@ impl RhiReconciliationSourceReplay {
.map(|evidence| evidence.cursor),
)
}
+
+ pub(crate) fn into_commit_parts(self) -> RhiReconciliationReplayCommitParts {
+ let eligible_cursor = self.eligible_cursor();
+ RhiReconciliationReplayCommitParts {
+ replay_id: self.plan.id,
+ request_id: self.plan.request_id,
+ source_id: self.plan.source_id,
+ trade_id: self.plan.trade_id,
+ required: self.plan.required,
+ policy_digest: self.plan.policy_digest,
+ selector_digest: self.plan.selector_digest,
+ prior_cursor: self.plan.prior_cursor.map(|evidence| evidence.cursor),
+ overlap_seconds: self.plan.overlap_seconds,
+ since_unix_seconds: self.plan.since_unix_seconds,
+ result: self.result,
+ duplicate_observations: self.duplicate_observations,
+ cursor_candidate: self.cursor_candidate,
+ eligible_cursor,
+ first_observed_at: self.first_observed_at,
+ facts: self
+ .facts
+ .into_vec()
+ .into_iter()
+ .map(|fact| RhiReconciliationReplayCommitFact {
+ record: fact.record,
+ observed_at: fact.observed_at,
+ })
+ .collect::<Vec<_>>()
+ .into_boxed_slice(),
+ }
+ }
+}
+
+pub(crate) struct RhiReconciliationReplayCommitParts {
+ pub(crate) replay_id: RhiReconciliationSourceReplayId,
+ pub(crate) request_id: RhiReconciliationSourceRequestId,
+ pub(crate) source_id: Box<str>,
+ pub(crate) trade_id: radroots_event::id::TradeId,
+ pub(crate) required: bool,
+ pub(crate) policy_digest: [u8; 32],
+ pub(crate) selector_digest: [u8; 32],
+ pub(crate) prior_cursor: Option<RhiTradeSourceCursor>,
+ pub(crate) overlap_seconds: u64,
+ pub(crate) since_unix_seconds: u64,
+ pub(crate) result: RhiReconciliationSourceResult,
+ pub(crate) duplicate_observations: u32,
+ pub(crate) cursor_candidate: Option<RhiTradeSourceCursor>,
+ pub(crate) eligible_cursor: Option<RhiTradeSourceCursor>,
+ pub(crate) first_observed_at: Option<RhiTradeMutationObservedAtUnixSeconds>,
+ pub(crate) facts: Box<[RhiReconciliationReplayCommitFact]>,
+}
+
+pub(crate) struct RhiReconciliationReplayCommitFact {
+ pub(crate) record: PersistenceRecord,
+ pub(crate) observed_at: RhiTradeMutationObservedAtUnixSeconds,
+}
+
+pub(crate) fn committed_cursor_evidence(
+ source_id: Box<str>,
+ trade_id: radroots_event::id::TradeId,
+ policy_digest: [u8; 32],
+ selector_digest: [u8; 32],
+ cursor: RhiTradeSourceCursor,
+) -> RhiReconciliationSourceCursorEvidence {
+ RhiReconciliationSourceCursorEvidence {
+ source_id,
+ trade_id,
+ policy_digest,
+ selector_digest,
+ cursor,
+ }
}
fn cursor_scope_matches(
diff --git a/src/source_ingest.rs b/src/source_ingest.rs
@@ -96,7 +96,7 @@ impl RhiTradeSourceCompletion {
}
#[must_use]
- const fn allows_checkpoint(self) -> bool {
+ pub(crate) const fn allows_checkpoint(self) -> bool {
matches!(self, Self::Complete)
}
}
@@ -486,10 +486,10 @@ impl ConfiguredSource {
}
#[derive(Clone, Copy, PartialEq, Eq)]
-struct Checkpoint {
- cursor: RhiTradeSourceCursor,
- revision: u64,
- completed_at_unix_s: u64,
+pub(crate) struct Checkpoint {
+ pub(crate) cursor: RhiTradeSourceCursor,
+ pub(crate) revision: u64,
+ pub(crate) completed_at_unix_s: u64,
}
#[derive(Clone, Copy, PartialEq, Eq)]
@@ -887,7 +887,7 @@ pub(crate) async fn read_dirty(
}))
}
-async fn read_checkpoint(
+pub(crate) async fn read_checkpoint(
transaction: &mut ServiceSqliteTransaction<'_>,
source_id: &str,
policy: RhiEvidencePolicyDigest,
@@ -914,7 +914,7 @@ async fn read_checkpoint(
}))
}
-async fn write_checkpoint(
+pub(crate) async fn write_checkpoint(
transaction: &mut ServiceSqliteTransaction<'_>,
source_id: &str,
policy: RhiEvidencePolicyDigest,
@@ -975,7 +975,7 @@ fn compare_candidate(left: &Candidate, right: &Candidate) -> Ordering {
})
}
-fn compare_cursor(left: RhiTradeSourceCursor, right: RhiTradeSourceCursor) -> Ordering {
+pub(crate) 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))
}
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 = 5;
+pub const RHI_STATE_SCHEMA_VERSION: u32 = 6;
/// The shared metadata and migration-ledger objects present at schema v1.
pub const RHI_STATE_SCHEMA_VERSION_1_OBJECT_COUNT: u32 = 6;
@@ -29,10 +29,13 @@ pub const RHI_STATE_SCHEMA_VERSION_4_OBJECT_COUNT: u32 = 28;
/// The shared objects, canonical evidence, cursors, generations, and durable jobs.
pub const RHI_STATE_SCHEMA_VERSION_5_OBJECT_COUNT: u32 = 33;
+/// The shared objects plus immutable reconciliation attempt/source results.
+pub const RHI_STATE_SCHEMA_VERSION_6_OBJECT_COUNT: u32 = 39;
+
/// SHA-256 identity of the ordered migration catalog rooted at schema v1.
pub const RHI_MIGRATION_CATALOG_SHA256: [u8; 32] = [
- 0x4e, 0xd3, 0x2a, 0x5a, 0x71, 0xf5, 0xbf, 0x45, 0x42, 0x62, 0xa9, 0x08, 0x2a, 0xd7, 0x4d, 0x0d,
- 0x6a, 0x13, 0xa5, 0x56, 0x05, 0x46, 0x59, 0xac, 0x7b, 0x9d, 0x87, 0x72, 0x59, 0xb8, 0x16, 0xd8,
+ 0x32, 0xf4, 0x9f, 0x1c, 0x50, 0xf5, 0x4c, 0x72, 0x2d, 0x94, 0xd9, 0x34, 0x3a, 0x05, 0x2e, 0x67,
+ 0x14, 0x2d, 0x00, 0x74, 0x7a, 0x0b, 0xb1, 0xcc, 0x1c, 0xb5, 0xca, 0xba, 0xfa, 0x30, 0xfe, 0x4c,
];
/// SHA-256 identity of the exact schema-v1 object snapshot.
@@ -89,10 +92,22 @@ pub const RHI_STATE_SCHEMA_VERSION_5_SHA256: [u8; 32] = [
0x99, 0x83, 0x15, 0x0b, 0x88, 0x7f, 0x26, 0x42, 0x50, 0xf9, 0xc8, 0x7a, 0x1a, 0x31, 0x3e, 0x60,
];
+/// SHA-256 identity of the schema-v6 reconciliation-result migration.
+pub const RHI_STATE_SCHEMA_VERSION_6_MIGRATION_SHA256: [u8; 32] = [
+ 0x48, 0xa1, 0x4b, 0x27, 0x44, 0xc4, 0x11, 0x86, 0x49, 0x6d, 0x59, 0x7e, 0xc7, 0x80, 0xad, 0x99,
+ 0x8f, 0xa9, 0x30, 0x8a, 0x9c, 0x9d, 0x75, 0x40, 0xd4, 0x75, 0xdb, 0x3f, 0x4d, 0xcf, 0xc4, 0xbb,
+];
+
+/// SHA-256 identity of the exact schema-v6 object snapshot.
+pub const RHI_STATE_SCHEMA_VERSION_6_SHA256: [u8; 32] = [
+ 0x5d, 0x1f, 0xa9, 0x50, 0x8b, 0x5b, 0x8a, 0x0c, 0x80, 0x65, 0xfd, 0x37, 0xf1, 0x1c, 0x68, 0x66,
+ 0xd5, 0x1f, 0x92, 0x94, 0xc7, 0x52, 0x72, 0x04, 0x1f, 0x34, 0x9e, 0x95, 0xce, 0xd5, 0x0f, 0xb4,
+];
+
/// SHA-256 identity of the schema catalog bound to the migration catalog.
pub const RHI_STATE_SCHEMA_CATALOG_SHA256: [u8; 32] = [
- 0x0d, 0x2c, 0x59, 0x5a, 0x3c, 0xa3, 0xfb, 0xc9, 0xb3, 0xbf, 0x6b, 0xfc, 0x61, 0x68, 0xa1, 0x94,
- 0xac, 0x78, 0x8d, 0x49, 0x31, 0xc4, 0x62, 0x00, 0xe2, 0xf1, 0x37, 0x9e, 0xbb, 0x52, 0x3e, 0xce,
+ 0x8c, 0x19, 0x82, 0xb2, 0x95, 0xd6, 0x54, 0x4b, 0xbe, 0x6d, 0xf6, 0x79, 0x42, 0x4e, 0xe6, 0xca,
+ 0xf3, 0x47, 0x68, 0x31, 0xb3, 0x8d, 0xe9, 0x94, 0x6b, 0x11, 0x5f, 0xbc, 0x5b, 0x24, 0x6b, 0x03,
];
macro_rules! rhi_config_bindings_table_sql {
@@ -607,6 +622,157 @@ const CREATE_RECONCILIATION_JOBS_MIGRATION_SQL: &str = concat!(
";",
);
+macro_rules! evidence_reconciliations_table_sql {
+ () => {
+ r#"CREATE TABLE evidence_reconciliations (
+ attempt_id BLOB NOT NULL PRIMARY KEY CHECK (length(attempt_id) = 32),
+ job_id BLOB NOT NULL CHECK (length(job_id) = 32),
+ trade_id BLOB NOT NULL CHECK (length(trade_id) = 16),
+ input_generation INTEGER NOT NULL
+ CHECK (input_generation BETWEEN 1 AND 9223372036854775807),
+ evidence_policy_sha256 BLOB NOT NULL CHECK (length(evidence_policy_sha256) = 32),
+ attempt_started_unix_ms INTEGER NOT NULL
+ CHECK (attempt_started_unix_ms BETWEEN 0 AND 9223372036854775807),
+ deadline_unix_ms INTEGER NOT NULL
+ CHECK (deadline_unix_ms BETWEEN attempt_started_unix_ms + 1 AND 9223372036854775807),
+ source_count INTEGER NOT NULL CHECK (source_count BETWEEN 1 AND 16),
+ FOREIGN KEY (job_id) REFERENCES reconciliation_jobs (job_id)
+) STRICT"#
+ };
+}
+
+macro_rules! evidence_reconciliation_sources_table_sql {
+ () => {
+ r#"CREATE TABLE evidence_reconciliation_sources (
+ attempt_id BLOB NOT NULL CHECK (length(attempt_id) = 32),
+ request_id BLOB NOT NULL CHECK (length(request_id) = 32),
+ source_ordinal INTEGER NOT NULL CHECK (source_ordinal BETWEEN 0 AND 15),
+ 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]'),
+ trade_id BLOB NOT NULL CHECK (length(trade_id) = 16),
+ required INTEGER NOT NULL CHECK (required IN (0, 1)),
+ selector_sha256 BLOB NOT NULL CHECK (length(selector_sha256) = 32),
+ replay_id BLOB NOT NULL CHECK (length(replay_id) = 32),
+ completion TEXT NOT NULL CHECK (completion IN (
+ 'complete', 'incomplete_timeout', 'incomplete_unavailable',
+ 'incomplete_resource_limit', 'incomplete_unknown', 'unsupported'
+ )),
+ started_unix_ms INTEGER NOT NULL
+ CHECK (started_unix_ms BETWEEN 0 AND 9223372036854775807),
+ finished_unix_ms INTEGER NOT NULL
+ CHECK (finished_unix_ms BETWEEN started_unix_ms AND 9223372036854775807),
+ accepted_event_count INTEGER NOT NULL CHECK (accepted_event_count BETWEEN 0 AND 4096),
+ accepted_event_bytes INTEGER NOT NULL CHECK (accepted_event_bytes BETWEEN 0 AND 8388608),
+ accepted_inventory_sha256 BLOB NOT NULL CHECK (length(accepted_inventory_sha256) = 32),
+ duplicate_observation_count INTEGER NOT NULL
+ CHECK (duplicate_observation_count BETWEEN 0 AND 4096),
+ first_observed_unix_s INTEGER
+ CHECK (first_observed_unix_s BETWEEN 0 AND 9223372036854775807),
+ prior_cursor_created_at_unix_s INTEGER
+ CHECK (prior_cursor_created_at_unix_s BETWEEN 0 AND 9223372036854775807),
+ prior_cursor_event_id BLOB CHECK (length(prior_cursor_event_id) = 32),
+ overlap_seconds INTEGER NOT NULL CHECK (overlap_seconds BETWEEN 1 AND 86400),
+ inclusive_since_unix_s INTEGER NOT NULL
+ CHECK (inclusive_since_unix_s BETWEEN 0 AND 9223372036854775807),
+ candidate_created_at_unix_s INTEGER
+ CHECK (candidate_created_at_unix_s BETWEEN 0 AND 9223372036854775807),
+ candidate_event_id BLOB CHECK (length(candidate_event_id) = 32),
+ checkpoint_advanced INTEGER NOT NULL CHECK (checkpoint_advanced IN (0, 1)),
+ PRIMARY KEY (attempt_id, request_id),
+ UNIQUE (attempt_id, source_ordinal),
+ FOREIGN KEY (attempt_id) REFERENCES evidence_reconciliations (attempt_id),
+ CHECK ((accepted_event_count = 0) = (accepted_event_bytes = 0)),
+ CHECK ((accepted_event_count = 0) = (first_observed_unix_s IS NULL)),
+ CHECK ((accepted_event_count = 0) = (candidate_created_at_unix_s IS NULL)),
+ CHECK ((candidate_created_at_unix_s IS NULL) = (candidate_event_id IS NULL)),
+ CHECK ((prior_cursor_created_at_unix_s IS NULL) = (prior_cursor_event_id IS NULL)),
+ CHECK (checkpoint_advanced = 0 OR completion = 'complete'),
+ CHECK (checkpoint_advanced = 0 OR candidate_event_id IS NOT NULL)
+) STRICT"#
+ };
+}
+
+const CREATE_EVIDENCE_RECONCILIATIONS_TABLE_SQL: &str = evidence_reconciliations_table_sql!();
+const CREATE_EVIDENCE_RECONCILIATIONS_NO_UPDATE_SQL: &str = immutable_no_update_sql!(
+ "evidence_reconciliations_no_update",
+ "evidence_reconciliations",
+ "reconciliation attempts are immutable"
+);
+const CREATE_EVIDENCE_RECONCILIATIONS_NO_DELETE_SQL: &str = immutable_no_delete_sql!(
+ "evidence_reconciliations_no_delete",
+ "evidence_reconciliations",
+ "reconciliation attempts are retained"
+);
+const CREATE_EVIDENCE_RECONCILIATION_SOURCES_TABLE_SQL: &str =
+ evidence_reconciliation_sources_table_sql!();
+const CREATE_EVIDENCE_RECONCILIATION_SOURCES_NO_UPDATE_SQL: &str = immutable_no_update_sql!(
+ "evidence_reconciliation_sources_no_update",
+ "evidence_reconciliation_sources",
+ "reconciliation source results are immutable"
+);
+const CREATE_EVIDENCE_RECONCILIATION_SOURCES_NO_DELETE_SQL: &str = immutable_no_delete_sql!(
+ "evidence_reconciliation_sources_no_delete",
+ "evidence_reconciliation_sources",
+ "reconciliation source results are retained"
+);
+const CREATE_RECONCILIATION_RESULTS_MIGRATION_SQL: &str = concat!(
+ evidence_reconciliations_table_sql!(),
+ ";\n",
+ immutable_no_update_sql!(
+ "evidence_reconciliations_no_update",
+ "evidence_reconciliations",
+ "reconciliation attempts are immutable"
+ ),
+ ";\n",
+ immutable_no_delete_sql!(
+ "evidence_reconciliations_no_delete",
+ "evidence_reconciliations",
+ "reconciliation attempts are retained"
+ ),
+ ";\n",
+ evidence_reconciliation_sources_table_sql!(),
+ ";\n",
+ immutable_no_update_sql!(
+ "evidence_reconciliation_sources_no_update",
+ "evidence_reconciliation_sources",
+ "reconciliation source results are immutable"
+ ),
+ ";\n",
+ immutable_no_delete_sql!(
+ "evidence_reconciliation_sources_no_delete",
+ "evidence_reconciliation_sources",
+ "reconciliation source results are retained"
+ ),
+ ";",
+);
+
+const EVIDENCE_RECONCILIATIONS_TABLE_SHA256: [u8; 32] = [
+ 0xaf, 0xed, 0xa3, 0xfb, 0xfb, 0x32, 0x6b, 0xd5, 0x27, 0x26, 0x42, 0x1d, 0x6c, 0x41, 0x5e, 0xff,
+ 0xf9, 0x71, 0xf7, 0x47, 0x72, 0x5e, 0xbf, 0x8b, 0x34, 0x26, 0x8f, 0x89, 0x6e, 0xa1, 0x2d, 0x9b,
+];
+const EVIDENCE_RECONCILIATIONS_NO_UPDATE_SHA256: [u8; 32] = [
+ 0xf4, 0x8f, 0xe6, 0x07, 0x7f, 0x32, 0x9c, 0x96, 0xc7, 0x8a, 0xad, 0xd4, 0x41, 0x53, 0x1a, 0xa3,
+ 0x24, 0x2b, 0x5d, 0x12, 0x63, 0x18, 0xdb, 0x39, 0xe7, 0x73, 0xbb, 0x0c, 0xa5, 0x40, 0xcd, 0xe9,
+];
+const EVIDENCE_RECONCILIATIONS_NO_DELETE_SHA256: [u8; 32] = [
+ 0x9a, 0x92, 0xf6, 0x89, 0x01, 0x75, 0xff, 0x89, 0x30, 0xee, 0x47, 0x8b, 0x8a, 0x6e, 0x10, 0x4f,
+ 0xc3, 0x88, 0x4f, 0x83, 0xc5, 0x01, 0x95, 0x9b, 0x72, 0x3e, 0xbd, 0xed, 0x87, 0x2b, 0x13, 0xdf,
+];
+const EVIDENCE_RECONCILIATION_SOURCES_TABLE_SHA256: [u8; 32] = [
+ 0xd1, 0x56, 0x84, 0x76, 0xc7, 0x52, 0xf8, 0x22, 0x5c, 0xa0, 0xce, 0x77, 0xa0, 0x1f, 0xc4, 0xc6,
+ 0x78, 0x22, 0x62, 0x61, 0x02, 0xaa, 0x92, 0x5c, 0xe5, 0x6d, 0x5c, 0x91, 0x4c, 0x8e, 0x1c, 0x5a,
+];
+const EVIDENCE_RECONCILIATION_SOURCES_NO_UPDATE_SHA256: [u8; 32] = [
+ 0x85, 0x8a, 0x6d, 0x3b, 0x97, 0x35, 0x91, 0x0b, 0x93, 0xe1, 0x9a, 0x73, 0x97, 0x6e, 0x47, 0x0e,
+ 0xe6, 0xb6, 0x38, 0x68, 0x31, 0x35, 0xf2, 0x1d, 0x3a, 0x30, 0xe1, 0xa6, 0x3f, 0xae, 0x07, 0x6b,
+];
+const EVIDENCE_RECONCILIATION_SOURCES_NO_DELETE_SHA256: [u8; 32] = [
+ 0x55, 0x75, 0xfc, 0xaf, 0xab, 0x07, 0xaf, 0x83, 0x0b, 0x69, 0x8d, 0xfa, 0x86, 0x56, 0xe9, 0xc6,
+ 0x98, 0x0b, 0x4c, 0xdc, 0xa2, 0xde, 0xce, 0xfc, 0x06, 0x41, 0x86, 0x69, 0x2e, 0xac, 0x00, 0x39,
+];
+
const RECONCILIATION_JOBS_TABLE_SHA256: [u8; 32] = [
0xe7, 0x0b, 0x5c, 0xb7, 0x26, 0x91, 0x9d, 0x02, 0xef, 0xb3, 0xa6, 0x21, 0x58, 0x48, 0xce, 0x92,
0x30, 0x88, 0x17, 0x2b, 0x3f, 0xe0, 0xc1, 0x40, 0xee, 0x4a, 0xc7, 0x1b, 0xd0, 0x3c, 0x66, 0xa7,
@@ -816,15 +982,23 @@ pub fn rhi_migration_catalog() -> Result<MigrationCatalog, RhiStateCatalogError>
MigrationChecksum::from_bytes(RHI_STATE_SCHEMA_VERSION_5_MIGRATION_SHA256),
)
.map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::MigrationCatalog))?;
+ let reconciliation_results = MigrationDescriptor::sql(
+ 6,
+ "create_reconciliation_source_results",
+ CREATE_RECONCILIATION_RESULTS_MIGRATION_SQL,
+ MigrationChecksum::from_bytes(RHI_STATE_SCHEMA_VERSION_6_MIGRATION_SHA256),
+ )
+ .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::MigrationCatalog))?;
let catalog = MigrationCatalog::new([
configuration,
trade_evidence,
source_checkpoints,
reconciliation_jobs,
+ reconciliation_results,
])
.map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::MigrationCatalog))?;
if catalog.current_version() != RHI_STATE_SCHEMA_VERSION
- || catalog.descriptors().len() != 4
+ || catalog.descriptors().len() != 5
|| catalog.digest().as_bytes() != &RHI_MIGRATION_CATALOG_SHA256
{
return Err(RhiStateCatalogError::new(
@@ -862,11 +1036,17 @@ pub fn rhi_schema_catalog() -> Result<SchemaCatalog, RhiStateCatalogError> {
)
.map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog))?;
let version_five = SchemaVersionCatalog::new(
- RHI_STATE_SCHEMA_VERSION,
+ 5,
rhi_schema_version_five_objects()?,
SchemaDigest::from_bytes(RHI_STATE_SCHEMA_VERSION_5_SHA256),
)
.map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog))?;
+ let version_six = SchemaVersionCatalog::new(
+ RHI_STATE_SCHEMA_VERSION,
+ rhi_schema_version_six_objects()?,
+ SchemaDigest::from_bytes(RHI_STATE_SCHEMA_VERSION_6_SHA256),
+ )
+ .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog))?;
let catalog = SchemaCatalog::new(
&migrations,
[
@@ -875,6 +1055,7 @@ pub fn rhi_schema_catalog() -> Result<SchemaCatalog, RhiStateCatalogError> {
version_three,
version_four,
version_five,
+ version_six,
],
)
.map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog))?;
@@ -889,10 +1070,10 @@ pub fn validate_rhi_state_catalogs(
) -> Result<(), RhiStateCatalogError> {
let versions = schema.versions();
let valid = migrations.current_version() == RHI_STATE_SCHEMA_VERSION
- && migrations.descriptors().len() == 4
+ && migrations.descriptors().len() == 5
&& migrations.digest().as_bytes() == &RHI_MIGRATION_CATALOG_SHA256
&& schema.migration_catalog_digest() == migrations.digest()
- && versions.len() == 5
+ && versions.len() == 6
&& 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
@@ -905,9 +1086,12 @@ pub fn validate_rhi_state_catalogs(
&& versions[3].version() == 4
&& versions[3].object_count() == RHI_STATE_SCHEMA_VERSION_4_OBJECT_COUNT
&& versions[3].digest().as_bytes() == &RHI_STATE_SCHEMA_VERSION_4_SHA256
- && versions[4].version() == RHI_STATE_SCHEMA_VERSION
+ && versions[4].version() == 5
&& versions[4].object_count() == RHI_STATE_SCHEMA_VERSION_5_OBJECT_COUNT
&& versions[4].digest().as_bytes() == &RHI_STATE_SCHEMA_VERSION_5_SHA256
+ && versions[5].version() == RHI_STATE_SCHEMA_VERSION
+ && versions[5].object_count() == RHI_STATE_SCHEMA_VERSION_6_OBJECT_COUNT
+ && versions[5].digest().as_bytes() == &RHI_STATE_SCHEMA_VERSION_6_SHA256
&& schema.digest().as_bytes() == &RHI_STATE_SCHEMA_CATALOG_SHA256;
if valid {
Ok(())
@@ -973,6 +1157,69 @@ fn rhi_schema_version_five_objects() -> Result<Vec<SchemaObject>, RhiStateCatalo
Ok(objects)
}
+fn rhi_schema_version_six_objects() -> Result<Vec<SchemaObject>, RhiStateCatalogError> {
+ let mut objects = rhi_schema_version_five_objects()?;
+ objects.extend(rhi_reconciliation_result_objects()?);
+ Ok(objects)
+}
+
+fn rhi_reconciliation_result_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,
+ "evidence_reconciliations",
+ "evidence_reconciliations",
+ CREATE_EVIDENCE_RECONCILIATIONS_TABLE_SQL,
+ EVIDENCE_RECONCILIATIONS_TABLE_SHA256,
+ )?,
+ object(
+ SchemaObjectKind::Trigger,
+ "evidence_reconciliations_no_update",
+ "evidence_reconciliations",
+ CREATE_EVIDENCE_RECONCILIATIONS_NO_UPDATE_SQL,
+ EVIDENCE_RECONCILIATIONS_NO_UPDATE_SHA256,
+ )?,
+ object(
+ SchemaObjectKind::Trigger,
+ "evidence_reconciliations_no_delete",
+ "evidence_reconciliations",
+ CREATE_EVIDENCE_RECONCILIATIONS_NO_DELETE_SQL,
+ EVIDENCE_RECONCILIATIONS_NO_DELETE_SHA256,
+ )?,
+ object(
+ SchemaObjectKind::Table,
+ "evidence_reconciliation_sources",
+ "evidence_reconciliation_sources",
+ CREATE_EVIDENCE_RECONCILIATION_SOURCES_TABLE_SQL,
+ EVIDENCE_RECONCILIATION_SOURCES_TABLE_SHA256,
+ )?,
+ object(
+ SchemaObjectKind::Trigger,
+ "evidence_reconciliation_sources_no_update",
+ "evidence_reconciliation_sources",
+ CREATE_EVIDENCE_RECONCILIATION_SOURCES_NO_UPDATE_SQL,
+ EVIDENCE_RECONCILIATION_SOURCES_NO_UPDATE_SHA256,
+ )?,
+ object(
+ SchemaObjectKind::Trigger,
+ "evidence_reconciliation_sources_no_delete",
+ "evidence_reconciliation_sources",
+ CREATE_EVIDENCE_RECONCILIATION_SOURCES_NO_DELETE_SQL,
+ EVIDENCE_RECONCILIATION_SOURCES_NO_DELETE_SHA256,
+ )?,
+ ])
+}
+
fn rhi_reconciliation_job_objects() -> Result<[SchemaObject; 5], RhiStateCatalogError> {
let object = |kind, name, table_name, sql, digest| {
SchemaObject::new(
diff --git a/src/state_repository.rs b/src/state_repository.rs
@@ -467,6 +467,12 @@ impl<'host> RhiReconciliationJobRepository<'host> {
}
}
+impl<'host> RhiReconciliationAttemptRepository<'host> {
+ pub(crate) const fn host(&self) -> &'host RhiStateHost {
+ self.host
+ }
+}
+
#[cfg(test)]
mod tests {
use super::*;
diff --git a/src/state_trade.rs b/src/state_trade.rs
@@ -121,6 +121,21 @@ impl RhiTradeSourceObservation {
observed_at: event.observed_at_unix_seconds(),
}
}
+
+ pub(crate) fn from_persistence_parts(
+ source_id: Box<str>,
+ policy: RhiEvidencePolicyDigest,
+ record: &PersistenceRecord,
+ observed_at: RhiTradeMutationObservedAtUnixSeconds,
+ ) -> Self {
+ Self {
+ source_id,
+ policy,
+ event_id: record.event_id,
+ event_signature: record.event_signature,
+ observed_at,
+ }
+ }
}
impl fmt::Debug for RhiTradeSourceObservation {
@@ -290,6 +305,7 @@ impl RhiStateRepositories<'_> {
}
}
+#[derive(Clone)]
pub(crate) struct PersistenceRecord {
pub(crate) mutation_id: [u8; 32],
pub(crate) trade_id: [u8; 16],
diff --git a/tests/package_boundary.rs b/tests/package_boundary.rs
@@ -12,12 +12,15 @@ const RUNTIME_ADAPTER_CONTRACT: &str =
include_str!("../contracts/services_hardening/runtime_adapters.v1.json");
const RUNTIME_FOUNDATION: &str = include_str!("../src/runtime_foundation.rs");
const RECONCILIATION_ATTEMPTS: &str = include_str!("../src/reconciliation_attempt.rs");
+const RECONCILIATION_COMMIT: &str = include_str!("../src/reconciliation_commit.rs");
const RECONCILIATION_JOBS: &str = include_str!("../src/reconciliation_job.rs");
const RECONCILIATION_REPLAY: &str = include_str!("../src/reconciliation_replay.rs");
const RECONCILIATION_ATTEMPT_CONTRACT: &str =
include_str!("../contracts/services_hardening/reconciliation_attempts.v1.json");
const RECONCILIATION_REPLAY_CONTRACT: &str =
include_str!("../contracts/services_hardening/reconciliation_replay.v1.json");
+const RECONCILIATION_COMMIT_CONTRACT: &str =
+ include_str!("../contracts/services_hardening/reconciliation_commit.v1.json");
const RUNTIME_FOUNDATION_CONTRACT: &str =
include_str!("../contracts/services_hardening/runtime_foundation.v1.json");
const TRADE_INGEST_CONTRACT: &str =
@@ -35,6 +38,7 @@ const SOURCES: &[&str] = &[
include_str!("../src/identity_credential.rs"),
include_str!("../src/identity_envelope.rs"),
include_str!("../src/reconciliation_attempt.rs"),
+ include_str!("../src/reconciliation_commit.rs"),
include_str!("../src/reconciliation_job.rs"),
include_str!("../src/reconciliation_replay.rs"),
include_str!("../src/runtime_context.rs"),
@@ -104,6 +108,7 @@ fn state_catalog_module_is_private_and_root_api_is_curated() {
"identity_credential",
"identity_envelope",
"reconciliation_attempt",
+ "reconciliation_commit",
"reconciliation_job",
"reconciliation_replay",
"runtime_context",
@@ -145,6 +150,8 @@ fn state_catalog_module_is_private_and_root_api_is_curated() {
"RhiReconciliationSourceRequest",
"RhiReconciliationSourceResult",
"RhiReconciliationAttemptResults",
+ "RhiReconciliationSourceCommitOutcome",
+ "RhiReconciliationCommitErrorKind",
"RhiReconciliationSourceReplayPlan",
"RhiReconciliationSourceReplay",
"RhiReconciliationJobPolicy",
@@ -203,6 +210,33 @@ fn state_catalog_module_is_private_and_root_api_is_curated() {
}
#[test]
+fn reconciliation_commit_is_atomic_bounded_and_sealed() {
+ let contract: serde_json::Value = serde_json::from_str(RECONCILIATION_COMMIT_CONTRACT)
+ .expect("reconciliation-commit contract");
+ assert_eq!(contract["schema"], "radroots.rhi.reconciliation-commit");
+ assert_eq!(contract["contract_version"], 1);
+ assert_eq!(contract["state_schema_version"], 6);
+ assert_eq!(contract["effects"]["source_or_relay"], false);
+ for required in [
+ ".take(plan.requests().len().saturating_add(1))",
+ "validate_exact_lease(transaction, lease)",
+ "reconcile_existing(transaction, &plan, &parts)",
+ "SOURCE_INVENTORY_DIGEST_DOMAIN",
+ "accepted_inventory_sha256",
+ "advance_dirty_generation(",
+ "write_checkpoint(",
+ "committed_cursor_evidence(",
+ ] {
+ assert!(
+ RECONCILIATION_COMMIT.contains(required),
+ "reconciliation commit is missing {required}"
+ );
+ }
+ assert!(!ROOT.contains("pub mod reconciliation_commit"));
+ assert!(!PUBLIC_API.contains("rhi::reconciliation_commit::"));
+}
+
+#[test]
fn public_errors_are_crate_owned_redacted_and_source_free() {
let production = SOURCES.join("\n");
assert!(!production.contains("fn source("));
@@ -234,7 +268,7 @@ fn public_errors_are_crate_owned_redacted_and_source_free() {
.lines()
.filter(|line| line.starts_with("pub struct rhi::") && line.ends_with("Error"))
.count();
- assert_eq!(public_error_count, 19);
+ assert_eq!(public_error_count, 20);
}
#[test]
@@ -650,6 +684,8 @@ fn readme_freezes_the_root_only_boundary_and_exact_baseline() {
"[`reconciliation_attempts.v1.json`](contracts/services_hardening/reconciliation_attempts.v1.json)",
"## Overlap-safe reconciliation replay",
"[`reconciliation_replay.v1.json`](contracts/services_hardening/reconciliation_replay.v1.json)",
+ "## Atomic reconciliation result commit",
+ "[`reconciliation_commit.v1.json`](contracts/services_hardening/reconciliation_commit.v1.json)",
"configured queue capacity is enforced beneath a fixed 65,536-job",
"from an unexpired claimed job lease",
"can be omitted, duplicated, reordered, or appended beyond",
@@ -659,7 +695,7 @@ fn readme_freezes_the_root_only_boundary_and_exact_baseline() {
"Only exact-target EOSE before the deadline is complete",
"4,096 distinct signed-event identities (event ID plus signature) and 8 MiB",
"observation, and operational retry do not",
- "schema-v5 bounded reconciliation-job migration",
+ "schema-v6 immutable",
"one canonical mutation",
"every distinct valid signed",
"does not advance reconciliation checkpoints or dirty generation",
diff --git a/tests/services_hardening_reconciliation_commit_contract.rs b/tests/services_hardening_reconciliation_commit_contract.rs
@@ -0,0 +1,102 @@
+#![forbid(unsafe_code)]
+
+use rhi::RHI_RECONCILIATION_COMMIT_CONTRACT_VERSION;
+use serde_json::json;
+
+const CONTRACT: &str =
+ include_str!("../contracts/services_hardening/reconciliation_commit.v1.json");
+const ROOT: &str = include_str!("../src/lib.rs");
+const SOURCE: &str = include_str!("../src/reconciliation_commit.rs");
+const README: &str = include_str!("../README");
+
+#[test]
+fn machine_contract_freezes_the_complete_step_190_boundary() {
+ let contract: serde_json::Value = serde_json::from_str(CONTRACT).expect("contract");
+ assert_eq!(contract["schema"], "radroots.rhi.reconciliation-commit");
+ assert_eq!(contract["schema_version"], 1);
+ assert_eq!(
+ contract["contract_version"],
+ RHI_RECONCILIATION_COMMIT_CONTRACT_VERSION
+ );
+ assert_eq!(contract["state_schema_version"], 6);
+ assert_eq!(
+ contract["transaction"]["fences_before_mutation"],
+ json!([
+ "exact_live_lease_row",
+ "exact_trade_dirty_generation",
+ "exact_evidence_policy_digest",
+ "exact_per_source_prior_checkpoint"
+ ])
+ );
+ assert_eq!(
+ contract["checkpoint"]["advance_only_when"],
+ "completion_is_complete_and_candidate_is_strictly_newer_than_prior"
+ );
+ assert_eq!(
+ contract["checkpoint"]["public_evidence"],
+ "sealed_and_minted_only_after_durable_commit_confirmation"
+ );
+ assert_eq!(
+ contract["dirty_generation"]["advance_count_maximum_per_attempt"],
+ 1
+ );
+ assert_eq!(
+ contract["idempotency"]["lost_success_retry"],
+ "checked_before_live_lease_and_generation_fences"
+ );
+ assert_eq!(
+ contract["idempotency"]["source_inventory_identity"],
+ "domain_separated_sha256_over_exact_ordered_persisted_facts_and_provenance"
+ );
+ assert_eq!(contract["effects"]["source_or_relay"], false);
+ assert_eq!(contract["effects"]["ambient_clock"], false);
+}
+
+#[test]
+fn commit_boundary_is_private_bounded_transactional_and_documented() {
+ assert!(ROOT.contains("mod reconciliation_commit;"));
+ assert!(!ROOT.contains("pub mod reconciliation_commit;"));
+ for required in [
+ "RhiReconciliationCommitError",
+ "RhiReconciliationCommitErrorKind",
+ "RhiReconciliationSourceCommitOutcome",
+ "RHI_RECONCILIATION_COMMIT_CONTRACT_VERSION",
+ ] {
+ assert!(ROOT.contains(required), "root API is missing {required}");
+ }
+ for required in [
+ ".take(plan.requests().len().saturating_add(1))",
+ "validate_exact_lease(transaction, lease)",
+ "reconcile_existing(transaction, &plan, &parts)",
+ "SOURCE_INVENTORY_DIGEST_DOMAIN",
+ "accepted_inventory_sha256",
+ "advance_dirty_generation(",
+ "write_checkpoint(",
+ "committed_cursor_evidence(",
+ "ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown",
+ ] {
+ assert!(
+ SOURCE.contains(required),
+ "commit boundary is missing {required}"
+ );
+ }
+ for forbidden in [
+ "std::fs",
+ "std::net",
+ "tokio::spawn",
+ "SystemTime",
+ "thread_rng",
+ "OsRng",
+ "pub fn transaction",
+ "pub fn sqlite_host",
+ ] {
+ assert!(
+ !SOURCE.contains(forbidden),
+ "commit boundary gained forbidden authority {forbidden}"
+ );
+ }
+ assert!(README.contains("## Atomic reconciliation result commit"));
+ assert!(README.contains(
+ "[`reconciliation_commit.v1.json`](contracts/services_hardening/reconciliation_commit.v1.json)"
+ ));
+}
diff --git a/tests/services_hardening_reconciliation_jobs.rs b/tests/services_hardening_reconciliation_jobs.rs
@@ -8,11 +8,12 @@ use radroots_storage::event::SourceGeneration;
use rhi::{
RadrootsHostEnvironment, RadrootsPathResolver, RadrootsPlatform,
RhiReconciliationAttemptErrorKind, RhiReconciliationAttemptPlan,
- RhiReconciliationAttemptResults, RhiReconciliationJobErrorKind, RhiReconciliationJobPolicy,
- RhiReconciliationJobState, RhiReconciliationLeaseOwner,
- RhiReconciliationRetryDelayMilliseconds, RhiReconciliationSourceReplayPlan,
- RhiReconciliationSourceResult, RhiReconciliationUnixMilliseconds, RhiRuntimeContext,
- RhiStateMetadata, RhiTradeMutationAdmissionLimits, RhiTradeMutationAuthoredTimePolicy,
+ RhiReconciliationAttemptResults, RhiReconciliationCommitErrorKind,
+ RhiReconciliationJobErrorKind, RhiReconciliationJobPolicy, RhiReconciliationJobState,
+ RhiReconciliationLease, RhiReconciliationLeaseOwner, RhiReconciliationRetryDelayMilliseconds,
+ RhiReconciliationSourceReplayPlan, RhiReconciliationSourceResult,
+ RhiReconciliationUnixMilliseconds, RhiRuntimeContext, RhiStateMetadata,
+ RhiTradeMutationAdmissionLimits, RhiTradeMutationAuthoredTimePolicy,
RhiTradeMutationObservedAtUnixSeconds, RhiTradeSourceCompletion, TradeId,
admit_rhi_trade_mutation_event, initialize_rhi_state, open_rhi_state_inspection,
open_rhi_state_read_write, parse_rhi_cli_v1_from, parse_rhi_config_v1,
@@ -916,7 +917,7 @@ async fn attempt_diagnostics_are_redacted_and_source_free() {
#[tokio::test]
async fn replay_plan_binds_overlap_deduplicates_and_retains_first_provenance() {
let started_ms = 1_784_347_200_000;
- let (root, runtime, metadata, configuration, host, plan) =
+ let (root, runtime, metadata, configuration, host, _lease, plan) =
replay_fixture("replay-plan", EXAMPLE, started_ms).await;
let request = &plan.requests()[0];
let cursor_plan =
@@ -975,6 +976,178 @@ async fn replay_plan_binds_overlap_deduplicates_and_retains_first_provenance() {
drop((configuration, metadata, runtime, root));
}
+#[tokio::test]
+async fn source_replay_commit_is_atomic_idempotent_and_mints_durable_cursor_evidence() {
+ let started_ms = 1_784_347_200_000;
+ let (root, runtime, metadata, configuration, host, lease, plan) =
+ replay_fixture("replay-commit", EXAMPLE, started_ms).await;
+ let request = &plan.requests()[0];
+ let wire = replay_wire();
+ let make_replay = || {
+ RhiReconciliationSourceReplayPlan::from_request(&plan, request, &configuration, None)
+ .expect("replay plan")
+ .finish(
+ request,
+ RhiTradeSourceCompletion::Complete,
+ now(started_ms),
+ now(started_ms + 2_000),
+ [admitted_replay(&configuration, &wire, 1_784_347_200)],
+ )
+ .expect("replay")
+ };
+ let first_replay = make_replay();
+ let retry_replay = make_replay();
+ let resume_plan = plan.clone();
+ let attempts = host.repositories().reconciliation_attempts();
+ let committed = attempts
+ .commit_source_replays(lease, plan.clone(), [first_replay])
+ .await
+ .expect("commit");
+ assert!(committed.created());
+ assert_eq!(committed.source_result_count(), 1);
+ assert_eq!(committed.checkpoint_advance_count(), 1);
+ assert!(committed.dirty_generation_advanced());
+ assert_eq!(committed.committed_cursors().len(), 1);
+ assert_eq!(
+ committed.committed_cursors()[0]
+ .cursor()
+ .created_at_unix_seconds(),
+ 1_784_347_200
+ );
+ let resumed = RhiReconciliationSourceReplayPlan::from_request(
+ &resume_plan,
+ &resume_plan.requests()[0],
+ &configuration,
+ Some(committed.committed_cursors()[0].clone()),
+ )
+ .expect("committed cursor resumes exact scope");
+ assert_eq!(resumed.since_unix_seconds(), 1_784_346_900);
+
+ let reconciled = attempts
+ .commit_source_replays(lease, plan, [retry_replay])
+ .await
+ .expect("idempotent reconcile");
+ assert!(!reconciled.created());
+ assert_eq!(reconciled.source_result_count(), 1);
+ assert_eq!(reconciled.checkpoint_advance_count(), 1);
+ assert!(!reconciled.dirty_generation_advanced());
+ host.close().await.expect("close");
+
+ let options = SqliteConnectOptions::new()
+ .filename(runtime.artifacts().state_database())
+ .create_if_missing(false)
+ .foreign_keys(true);
+ let mut connection = SqliteConnection::connect_with(&options)
+ .await
+ .expect("offline fixture connection");
+ let attempt_count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM evidence_reconciliations")
+ .fetch_one(&mut connection)
+ .await
+ .expect("attempt count");
+ let source_count: i64 =
+ sqlx::query_scalar("SELECT COUNT(*) FROM evidence_reconciliation_sources")
+ .fetch_one(&mut connection)
+ .await
+ .expect("source count");
+ let checkpoint_count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM relay_checkpoints")
+ .fetch_one(&mut connection)
+ .await
+ .expect("checkpoint count");
+ let generation: i64 = sqlx::query_scalar("SELECT generation FROM trade_dirty_generations")
+ .fetch_one(&mut connection)
+ .await
+ .expect("generation");
+ assert_eq!(
+ (attempt_count, source_count, checkpoint_count, generation),
+ (1, 1, 1, 2)
+ );
+ connection.close().await.expect("fixture close");
+ drop((configuration, metadata, runtime, root));
+}
+
+#[tokio::test]
+async fn incomplete_results_never_advance_and_stale_leases_fail_closed() {
+ let started_ms = 1_784_347_200_000;
+ let (root, runtime, metadata, configuration, host, lease, plan) =
+ replay_fixture("replay-incomplete", EXAMPLE, started_ms).await;
+ let request = &plan.requests()[0];
+ let wire = replay_wire();
+ let incomplete =
+ RhiReconciliationSourceReplayPlan::from_request(&plan, request, &configuration, None)
+ .expect("replay plan")
+ .finish(
+ request,
+ RhiTradeSourceCompletion::IncompleteUnavailable,
+ now(started_ms),
+ now(started_ms + 2_000),
+ [admitted_replay(&configuration, &wire, 1_784_347_200)],
+ )
+ .expect("incomplete replay");
+ let committed = host
+ .repositories()
+ .reconciliation_attempts()
+ .commit_source_replays(lease, plan, [incomplete])
+ .await
+ .expect("incomplete commit");
+ assert_eq!(committed.checkpoint_advance_count(), 0);
+ assert!(committed.committed_cursors().is_empty());
+ assert!(committed.dirty_generation_advanced());
+ host.close().await.expect("close");
+ drop((configuration, metadata, runtime, root));
+
+ let started_ms = 1_784_347_300_000;
+ let (root, runtime, metadata, configuration, host, lease, plan) =
+ replay_fixture("replay-stale-lease", EXAMPLE, started_ms).await;
+ let request = &plan.requests()[0];
+ let replay =
+ RhiReconciliationSourceReplayPlan::from_request(&plan, request, &configuration, None)
+ .expect("replay plan")
+ .finish(
+ request,
+ RhiTradeSourceCompletion::Complete,
+ now(started_ms),
+ now(started_ms + 2_000),
+ [],
+ )
+ .expect("empty replay");
+ host.repositories()
+ .reconciliation_jobs()
+ .renew(lease, now(started_ms + 20_000))
+ .await
+ .expect("renewed lease");
+ let error = host
+ .repositories()
+ .reconciliation_attempts()
+ .commit_source_replays(lease, plan, [replay])
+ .await
+ .expect_err("stale lease");
+ assert_eq!(error.kind(), RhiReconciliationCommitErrorKind::LeaseLost);
+ assert!(Error::source(&error).is_none());
+ host.close().await.expect("close");
+ let options = SqliteConnectOptions::new()
+ .filename(runtime.artifacts().state_database())
+ .create_if_missing(false)
+ .foreign_keys(true);
+ let mut connection = SqliteConnection::connect_with(&options)
+ .await
+ .expect("offline fixture connection");
+ let attempt_count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM evidence_reconciliations")
+ .fetch_one(&mut connection)
+ .await
+ .expect("attempt count");
+ let checkpoint_count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM relay_checkpoints")
+ .fetch_one(&mut connection)
+ .await
+ .expect("checkpoint count");
+ let generation: i64 = sqlx::query_scalar("SELECT generation FROM trade_dirty_generations")
+ .fetch_one(&mut connection)
+ .await
+ .expect("generation");
+ assert_eq!((attempt_count, checkpoint_count, generation), (0, 0, 1));
+ connection.close().await.expect("fixture close");
+ drop((configuration, metadata, runtime, root));
+}
+
async fn attempt_fixture(
instance: &str,
) -> (
@@ -1025,6 +1198,7 @@ async fn replay_fixture(
RhiStateMetadata,
rhi::RhiConfigDocumentV1,
rhi::RhiStateHost,
+ RhiReconciliationLease,
RhiReconciliationAttemptPlan,
) {
let root = tempfile::tempdir().expect("root");
@@ -1054,7 +1228,7 @@ async fn replay_fixture(
.expect("job");
let plan = RhiReconciliationAttemptPlan::from_claim(lease, &configuration, now(started_ms))
.expect("plan");
- (root, runtime, metadata, configuration, host, plan)
+ (root, runtime, metadata, configuration, host, lease, plan)
}
fn replay_wire() -> Vec<u8> {
diff --git a/tests/services_hardening_state_catalog.rs b/tests/services_hardening_state_catalog.rs
@@ -15,8 +15,10 @@ use rhi::{
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,
RHI_STATE_SCHEMA_VERSION_5_MIGRATION_SHA256, RHI_STATE_SCHEMA_VERSION_5_OBJECT_COUNT,
- RHI_STATE_SCHEMA_VERSION_5_SHA256, RhiStateCatalogErrorKind, rhi_migration_catalog,
- rhi_schema_catalog, validate_rhi_state_catalogs,
+ RHI_STATE_SCHEMA_VERSION_5_SHA256, RHI_STATE_SCHEMA_VERSION_6_MIGRATION_SHA256,
+ RHI_STATE_SCHEMA_VERSION_6_OBJECT_COUNT, RHI_STATE_SCHEMA_VERSION_6_SHA256,
+ RhiStateCatalogErrorKind, rhi_migration_catalog, rhi_schema_catalog,
+ validate_rhi_state_catalogs,
};
const CATALOG_SOURCE: &str = include_str!("../src/state_catalog.rs");
@@ -24,14 +26,14 @@ const LIB_SOURCE: &str = include_str!("../src/lib.rs");
const MANIFEST: &str = include_str!("../Cargo.toml");
#[test]
-fn schema_v1_through_v5_catalogs_have_exact_literal_identities() {
+fn schema_v1_through_v6_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, 5);
- assert_eq!(migrations.descriptors().len(), 4);
- assert_eq!(migrations.current_version(), 5);
+ assert_eq!(RHI_STATE_SCHEMA_VERSION, 6);
+ assert_eq!(migrations.descriptors().len(), 5);
+ assert_eq!(migrations.current_version(), 6);
assert_eq!(migrations.descriptors()[0].target_version(), 2);
assert_eq!(
migrations.descriptors()[0].name().as_str(),
@@ -68,12 +70,21 @@ fn schema_v1_through_v5_catalogs_have_exact_literal_identities() {
migrations.descriptors()[3].checksum().as_bytes(),
&RHI_STATE_SCHEMA_VERSION_5_MIGRATION_SHA256
);
+ assert_eq!(migrations.descriptors()[4].target_version(), 6);
+ assert_eq!(
+ migrations.descriptors()[4].name().as_str(),
+ "create_reconciliation_source_results"
+ );
+ assert_eq!(
+ migrations.descriptors()[4].checksum().as_bytes(),
+ &RHI_STATE_SCHEMA_VERSION_6_MIGRATION_SHA256
+ );
assert_eq!(
migrations.digest().as_bytes(),
&RHI_MIGRATION_CATALOG_SHA256
);
- assert_eq!(schema.versions().len(), 5);
+ assert_eq!(schema.versions().len(), 6);
let version = schema.versions()[0];
assert_eq!(version.version(), 1);
assert_eq!(
@@ -129,13 +140,24 @@ fn schema_v1_through_v5_catalogs_have_exact_literal_identities() {
version.digest().as_bytes(),
&RHI_STATE_SCHEMA_VERSION_5_SHA256
);
+ let version = schema.versions()[5];
+ assert_eq!(version.version(), 6);
+ assert_eq!(
+ version.object_count(),
+ RHI_STATE_SCHEMA_VERSION_6_OBJECT_COUNT
+ );
+ assert_eq!(version.object_count(), 39);
+ assert_eq!(
+ version.digest().as_bytes(),
+ &RHI_STATE_SCHEMA_VERSION_6_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),
- "4ed32a5a71f5bf454262a9082ad74d0d6a13a556054659ac7b9d877259b816d8"
+ "32f49f1c50f54c722d94d9343a052e67142d00747a0bb1cc1cb5cabafa30fe4c"
);
assert_eq!(
lower_hex(&RHI_STATE_SCHEMA_VERSION_1_SHA256),
@@ -174,8 +196,16 @@ fn schema_v1_through_v5_catalogs_have_exact_literal_identities() {
"eef67b48400dee234f6c36d422551c7e9983150b887f264250f9c87a1a313e60"
);
assert_eq!(
+ lower_hex(&RHI_STATE_SCHEMA_VERSION_6_MIGRATION_SHA256),
+ "48a14b2744c41186496d597ec780ad998fa9308a9c9d7540d475db3f4dcfc4bb"
+ );
+ assert_eq!(
+ lower_hex(&RHI_STATE_SCHEMA_VERSION_6_SHA256),
+ "5d1fa9508b5b8a0c8065fd37f11c6866d51f9294c75272041f349e95ced50fb4"
+ );
+ assert_eq!(
lower_hex(&RHI_STATE_SCHEMA_CATALOG_SHA256),
- "0d2c595a3ca3fbc9b3bf6bfc6168a194ac788d4931c46200e2f1379ebb523ece"
+ "8c1982b295d6544bbe6df679424ee6caf3476831b38de9946b115fbc5b246b03"
);
}
@@ -227,6 +257,10 @@ fn independent_validator_rejects_migration_or_schema_drift() {
SchemaVersionCatalog::computed_digest(5, [version_two_object()]).expect("v5 digest");
let version_five = SchemaVersionCatalog::new(5, [version_two_object()], snapshot_digest)
.expect("version five");
+ let snapshot_digest =
+ SchemaVersionCatalog::computed_digest(6, [version_two_object()]).expect("v6 digest");
+ let version_six =
+ SchemaVersionCatalog::new(6, [version_two_object()], snapshot_digest).expect("version six");
let schema = SchemaCatalog::new(
&exact_migrations,
[
@@ -235,6 +269,7 @@ fn independent_validator_rejects_migration_or_schema_drift() {
version_three,
version_four,
version_five,
+ version_six,
],
)
.expect("drift schema catalog");
@@ -285,6 +320,10 @@ fn catalog_errors_are_stable_source_free_and_redacted() {
SchemaVersionCatalog::computed_digest(5, [secret_object()]).expect("v5 digest");
let version_five =
SchemaVersionCatalog::new(5, [secret_object()], version_five_digest).expect("version five");
+ let version_six_digest =
+ SchemaVersionCatalog::computed_digest(6, [secret_object()]).expect("v6 digest");
+ let version_six =
+ SchemaVersionCatalog::new(6, [secret_object()], version_six_digest).expect("version six");
let schema = SchemaCatalog::new(
&migrations,
[
@@ -293,6 +332,7 @@ fn catalog_errors_are_stable_source_free_and_redacted() {
version_three,
version_four,
version_five,
+ version_six,
],
)
.expect("schema catalog");
diff --git a/tests/services_hardening_state_resilience.rs b/tests/services_hardening_state_resilience.rs
@@ -347,8 +347,8 @@ async fn exact_open_rejects_unexpected_migration_history_without_repair() {
service_version, service_commit, lib_revision, rust_version, target,
feature_profile, config_contract_version, state_contract_version,
admin_contract_version, status_contract_version, provider_contract_version
- ) VALUES (5, 'unexpected_schema', ?, 1725000000, '0.1.0', ?, ?,
- 'rustc-test', 'test-target', 'service-host', 1, 5, 1, 1, 1)",
+ ) VALUES (7, 'unexpected_schema', ?, 1725000000, '0.1.0', ?, ?,
+ 'rustc-test', 'test-target', 'service-host', 1, 7, 1, 1, 1)",
)
.bind([0x44_u8; 32].as_slice())
.bind("1111111111111111111111111111111111111111")