rhi

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

commit 2e4f73e4b9367cb012b55641d180f6824f95493f
parent 7e457dd7496b1ff5375acd60de51c58f8b569c2e
Author: triesap <tyson@radroots.org>
Date:   Mon, 24 Aug 2026 11:48:24 +0000

refactor(rhi): commit reconciliation finalization atomically

- bind finalization to the exact config, lease, generation, source inventory, and supersession\n- persist exact manifest, report, signed bytes, publication targets, and completed job in one transaction\n- add idempotency, rollback, disabled-publication, contract, documentation, and API evidence

Diffstat:
MAGENTS.md | 8++++++++
MREADME | 22++++++++++++++++++++++
Mcontracts/api_baselines/rhi.txt | 34++++++++++++++++++++++++++++++++++
Acontracts/services_hardening/reconciliation_finalization_commit.v1.json | 89+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Msrc/lib.rs | 6++++++
Msrc/publication.rs | 33+++++++++++++++++++++++++++++----
Msrc/reconciliation_attestation.rs | 22++++++++++++++++++++++
Msrc/reconciliation_finalization.rs | 20+++++++++++++-------
Asrc/reconciliation_finalization_commit.rs | 1097+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Msrc/reconciliation_job.rs | 8++++++++
Msrc/state_metadata.rs | 2+-
Mtests/package_boundary.rs | 52+++++++++++++++++++++++++++++++++++++++++++++++++++-
Atests/services_hardening_reconciliation_finalization_commit_contract.rs | 72++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mtests/services_hardening_reconciliation_jobs.rs | 300+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
14 files changed, 1748 insertions(+), 17 deletions(-)

diff --git a/AGENTS.md b/AGENTS.md @@ -167,6 +167,14 @@ sealed manifest, projection, and evaluation chain. A finalization preflight is not commit authority: rerun its exact lease, generation, policy, and attempt validator inside the Step 199 atomic transaction before any write. +- Commit a signed finalization only through the sealed attempt repository. + Bind publication authority to the same normalized configuration as the open + state host, reconcile exact prior success before testing the consumed lease, + and write manifest, projection, report, exact signed bytes, explicit + supersession, required outbox/targets, and completed job in one short SQLx + transaction. Disabled publication must create no outbox, and no source, + relay, network, filesystem, task, clock, or entropy operation may occur in + this boundary. ## 6. Report, attestation, and publication invariants diff --git a/README b/README @@ -343,6 +343,28 @@ outcomes, retry, recovery, and wave qualification. The exact machine contract is [`publication_outbox.v1.json`](contracts/services_hardening/publication_outbox.v1.json). +## Atomic reconciliation finalization commit + +`RhiReconciliationAttemptRepository::commit_finalization` borrows one sealed, +independently verified signed attestation and the publication authority derived +from the same normalized configuration as the open state host. One short +SQLx-owned transaction first reconciles an exact prior success, then reruns the +live lease, dirty-generation, policy, committed-attempt, source-inventory, +checkpoint, supersession, and publication-capacity fences before its first +write. + +The same transaction retains the exact canonical manifest, projection, +canonical report, and independently verified signed-event bytes; records the +explicit supersession; creates the immutable outbox and ordered target set only +when publication is required; and completes the exact reconciliation job. +Disabled publication creates no outbox or target row. An exact retry returns +`created = false` even after the lease was consumed or expired, while any +mismatched durable footprint fails closed. This boundary performs no source or +relay I/O, network access, filesystem access, task spawn, ambient clock read, +or ambient entropy read. Publication claims and relay submission remain later +steps. The exact machine contract is +[`reconciliation_finalization_commit.v1.json`](contracts/services_hardening/reconciliation_finalization_commit.v1.json). + ## Existing-state runtime foundation `open_rhi_runtime_foundation` opens only an already initialized database from diff --git a/contracts/api_baselines/rhi.txt b/contracts/api_baselines/rhi.txt @@ -180,6 +180,18 @@ pub rhi::RhiReconciliationCommitErrorKind::LeaseLost pub rhi::RhiReconciliationCommitErrorKind::Storage impl rhi::RhiReconciliationCommitErrorKind pub const fn rhi::RhiReconciliationCommitErrorKind::code(self) -> &'static str +pub enum rhi::RhiReconciliationFinalizationCommitErrorKind +pub rhi::RhiReconciliationFinalizationCommitErrorKind::AttemptUnavailable +pub rhi::RhiReconciliationFinalizationCommitErrorKind::CommitOutcomeUnknown +pub rhi::RhiReconciliationFinalizationCommitErrorKind::Conflict +pub rhi::RhiReconciliationFinalizationCommitErrorKind::GenerationConflict +pub rhi::RhiReconciliationFinalizationCommitErrorKind::InvalidInput +pub rhi::RhiReconciliationFinalizationCommitErrorKind::InvalidMode +pub rhi::RhiReconciliationFinalizationCommitErrorKind::LeaseLost +pub rhi::RhiReconciliationFinalizationCommitErrorKind::PublicationQueueFull +pub rhi::RhiReconciliationFinalizationCommitErrorKind::Storage +impl rhi::RhiReconciliationFinalizationCommitErrorKind +pub const fn rhi::RhiReconciliationFinalizationCommitErrorKind::code(self) -> &'static str pub enum rhi::RhiReconciliationFinalizationErrorKind pub rhi::RhiReconciliationFinalizationErrorKind::AttemptUnavailable pub rhi::RhiReconciliationFinalizationErrorKind::CommitOutcomeUnknown @@ -658,6 +670,9 @@ pub const fn rhi::RhiPublicationAuthority::queue_capacity(&self) -> u32 pub const fn rhi::RhiPublicationAuthority::retry_policy(&self) -> core::option::Option<rhi::RhiPublicationRetryPolicy> pub const fn rhi::RhiPublicationAuthority::target_set_sha256(&self) -> &[u8; 32] pub fn rhi::RhiPublicationAuthority::targets(&self) -> &[rhi::RhiPublicationTarget] +impl core::cmp::Eq for rhi::RhiPublicationAuthority +impl core::cmp::PartialEq for rhi::RhiPublicationAuthority +pub fn rhi::RhiPublicationAuthority::eq(&self, &Self) -> bool impl core::fmt::Debug for rhi::RhiPublicationAuthority pub fn rhi::RhiPublicationAuthority::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result pub struct rhi::RhiPublicationError @@ -722,6 +737,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_finalization(&self, &rhi::RhiSignedEvidenceAttestation, &rhi::RhiPublicationAuthority, rhi::RhiReconciliationUnixMilliseconds) -> core::result::Result<rhi::RhiReconciliationFinalizationCommitOutcome, rhi::RhiReconciliationFinalizationCommitError> +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 @@ -765,6 +782,22 @@ pub const fn rhi::RhiReconciliationEvaluation::projection(&self) -> &rhi::RhiRec pub const fn rhi::RhiReconciliationEvaluation::reason_codes(&self) -> &[rhi::RhiReconciliationReasonCode] impl core::fmt::Debug for rhi::RhiReconciliationEvaluation pub fn rhi::RhiReconciliationEvaluation::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiReconciliationFinalizationCommitError +impl rhi::RhiReconciliationFinalizationCommitError +pub const fn rhi::RhiReconciliationFinalizationCommitError::code(self) -> &'static str +pub const fn rhi::RhiReconciliationFinalizationCommitError::kind(self) -> rhi::RhiReconciliationFinalizationCommitErrorKind +impl core::error::Error for rhi::RhiReconciliationFinalizationCommitError +impl core::fmt::Debug for rhi::RhiReconciliationFinalizationCommitError +pub fn rhi::RhiReconciliationFinalizationCommitError::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +impl core::fmt::Display for rhi::RhiReconciliationFinalizationCommitError +pub fn rhi::RhiReconciliationFinalizationCommitError::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiReconciliationFinalizationCommitOutcome +impl rhi::RhiReconciliationFinalizationCommitOutcome +pub const fn rhi::RhiReconciliationFinalizationCommitOutcome::created(self) -> bool +pub const fn rhi::RhiReconciliationFinalizationCommitOutcome::publication_mode(self) -> rhi::RhiPublicationMode +pub const fn rhi::RhiReconciliationFinalizationCommitOutcome::target_count(self) -> u8 +impl core::fmt::Debug for rhi::RhiReconciliationFinalizationCommitOutcome +pub fn rhi::RhiReconciliationFinalizationCommitOutcome::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result pub struct rhi::RhiReconciliationFinalizationError impl rhi::RhiReconciliationFinalizationError pub const fn rhi::RhiReconciliationFinalizationError::code(self) -> &'static str @@ -1386,6 +1419,7 @@ pub const rhi::RHI_RECONCILIATION_ATTEMPT_CONTRACT_VERSION: u32 pub const rhi::RHI_RECONCILIATION_ATTEMPT_MAX_SOURCES: usize pub const rhi::RHI_RECONCILIATION_ATTESTATION_CONTRACT_VERSION: u32 pub const rhi::RHI_RECONCILIATION_COMMIT_CONTRACT_VERSION: u32 +pub const rhi::RHI_RECONCILIATION_FINALIZATION_COMMIT_CONTRACT_VERSION: u32 pub const rhi::RHI_RECONCILIATION_FINALIZATION_CONTRACT_VERSION: u32 pub const rhi::RHI_RECONCILIATION_JOB_CONTRACT_VERSION: u32 pub const rhi::RHI_RECONCILIATION_JOB_MAX_ACTIVE: u32 diff --git a/contracts/services_hardening/reconciliation_finalization_commit.v1.json b/contracts/services_hardening/reconciliation_finalization_commit.v1.json @@ -0,0 +1,89 @@ +{ + "schema": "radroots.rhi.reconciliation-finalization-commit", + "schema_version": 1, + "contract_version": 1, + "authority": "sealed_signed_attestation_plus_exact_config_derived_publication_authority", + "input": { + "attestation": "borrowed_sealed_step_196_signed_attestation", + "publication": "same_normalized_configuration_as_open_state_host", + "time": "caller_injected_integer_utc_milliseconds", + "network": false + }, + "transaction": { + "count": 1, + "boundary": "one_short_service_sqlite_transaction", + "reconcile_exact_prior_success_before_consumed_lease_validation": true, + "fences_before_first_write": [ + "exact_live_lease", + "exact_dirty_generation", + "exact_evidence_policy_digest", + "exact_committed_attempt", + "exact_manifest_source_inventory", + "every_advanced_checkpoint", + "exact_superseded_statement_and_event_pair", + "configured_publication_queue_capacity" + ], + "atomic_writes": [ + "canonical_manifest", + "projection", + "canonical_report", + "exact_verified_signed_event_bytes", + "explicit_supersession", + "required_publication_outbox", + "immutable_ordered_publication_targets", + "completed_reconciliation_job" + ] + }, + "publication": { + "required": "one_pending_outbox_plus_exact_configured_target_inventory", + "disabled": "no_outbox_and_no_target_rows", + "outbox_identity": "sha256_domain_plus_event_id_plus_authority_digest_plus_target_set_digest", + "relay_io": false, + "claim_or_submission": false + }, + "idempotency": { + "identity": "exact_attestation_and_publication_immutable_footprint", + "exact_retry": "created_false_even_after_lease_consumption_or_expiry", + "later_checkpoint_or_outbox_progress": "does_not_change_immutable_success_reconciliation", + "mismatch": "conflict", + "commit_acknowledgement_loss": "commit_outcome_unknown_then_exact_retry" + }, + "result": { + "created": "true_for_first_commit_false_for_exact_reconciliation", + "publication_mode": [ + "required", + "disabled" + ], + "target_count": "exact_u8_or_zero_when_disabled", + "debug": "mode_and_counts_only", + "errors": "crate_owned_source_free_redacted" + }, + "effects": { + "sqlite_read": true, + "sqlite_write": true, + "filesystem": false, + "source_or_relay": false, + "network": false, + "task_spawn": false, + "ambient_clock": false, + "ambient_entropy": false + }, + "forbidden": [ + "manifest_or_projection_recomputation", + "report_or_signed_event_reserialization", + "stale_lease_or_generation_commit", + "partial_finalization_inventory", + "publication_from_a_different_configuration", + "implicit_publication_when_disabled", + "relay_io_inside_transaction", + "checkpoint_advance", + "raw_sqlx_authority_escape" + ], + "deferred": [ + "publication_claims", + "relay_submission", + "publication_result_recording", + "retry_and_terminal_outbox_state", + "integration_wave_qualification" + ] +} diff --git a/src/lib.rs b/src/lib.rs @@ -13,6 +13,7 @@ mod reconciliation_attempt; mod reconciliation_attestation; mod reconciliation_commit; mod reconciliation_finalization; +mod reconciliation_finalization_commit; mod reconciliation_job; mod reconciliation_manifest; mod reconciliation_reducer; @@ -101,6 +102,11 @@ pub use reconciliation_finalization::{ RHI_RECONCILIATION_FINALIZATION_CONTRACT_VERSION, RhiReconciliationFinalizationError, RhiReconciliationFinalizationErrorKind, RhiReconciliationFinalizationFence, }; +pub use reconciliation_finalization_commit::{ + RHI_RECONCILIATION_FINALIZATION_COMMIT_CONTRACT_VERSION, + RhiReconciliationFinalizationCommitError, RhiReconciliationFinalizationCommitErrorKind, + RhiReconciliationFinalizationCommitOutcome, +}; pub use reconciliation_job::{ RHI_RECONCILIATION_JOB_CONTRACT_VERSION, RHI_RECONCILIATION_JOB_MAX_ACTIVE, RhiReconciliationJob, RhiReconciliationJobError, RhiReconciliationJobErrorKind, diff --git a/src/publication.rs b/src/publication.rs @@ -6,7 +6,7 @@ use std::error::Error; use serde_json::Value; use sha2::{Digest, Sha256}; -use crate::RhiConfigDocumentV1; +use crate::{RhiConfigDocumentV1, state_metadata::normalized_config_digest}; /// Exact version of the RHI publication-authority contract. pub const RHI_PUBLICATION_CONTRACT_VERSION: u32 = 1; @@ -124,8 +124,9 @@ impl fmt::Debug for RhiPublicationTarget { /// /// let _forged = RhiPublicationAuthority { mode: todo!() }; /// ``` -#[derive(Clone, PartialEq, Eq)] +#[derive(Clone)] pub struct RhiPublicationAuthority { + configuration_sha256: [u8; 32], mode: RhiPublicationMode, targets: Box<[RhiPublicationTarget]>, retry: Option<RhiPublicationRetryPolicy>, @@ -134,10 +135,23 @@ pub struct RhiPublicationAuthority { authority_sha256: [u8; 32], } +impl PartialEq for RhiPublicationAuthority { + fn eq(&self, other: &Self) -> bool { + self.mode == other.mode + && self.targets == other.targets + && self.retry == other.retry + && self.queue_capacity == other.queue_capacity + && self.target_set_sha256 == other.target_set_sha256 + && self.authority_sha256 == other.authority_sha256 + } +} + +impl Eq for RhiPublicationAuthority {} + impl RhiPublicationAuthority { /// Derives the only publication authority from one complete admitted config. pub fn from_config(config: &RhiConfigDocumentV1) -> Result<Self, RhiPublicationError> { - derive_authority(config.normalized()) + derive_authority(config.normalized(), config.profile()) } /// Returns the explicit configured publication mode. @@ -175,6 +189,10 @@ impl RhiPublicationAuthority { pub const fn authority_sha256(&self) -> &[u8; 32] { &self.authority_sha256 } + + pub(crate) const fn configuration_sha256(&self) -> &[u8; 32] { + &self.configuration_sha256 + } } impl fmt::Debug for RhiPublicationAuthority { @@ -250,7 +268,13 @@ impl fmt::Debug for RhiPublicationError { impl Error for RhiPublicationError {} -fn derive_authority(document: &Value) -> Result<RhiPublicationAuthority, RhiPublicationError> { +fn derive_authority( + document: &Value, + profile: crate::RhiConfigProfile, +) -> Result<RhiPublicationAuthority, RhiPublicationError> { + let configuration_sha256 = *normalized_config_digest(profile, document) + .map_err(|_| failure(RhiPublicationErrorKind::InvalidConfiguration))? + .as_bytes(); let mode = match string(document, "/publication/mode")? { "required" => RhiPublicationMode::Required, "disabled" => RhiPublicationMode::Disabled, @@ -356,6 +380,7 @@ fn derive_authority(document: &Value) -> Result<RhiPublicationAuthority, RhiPubl queue_capacity, )?; Ok(RhiPublicationAuthority { + configuration_sha256, mode, targets: targets.into_boxed_slice(), retry, diff --git a/src/reconciliation_attestation.rs b/src/reconciliation_attestation.rs @@ -252,6 +252,28 @@ impl RhiSignedEvidenceAttestation { pub const fn has_supersession(&self) -> bool { self.report.supersession().is_some() } + + pub(crate) const fn fence(&self) -> &RhiReconciliationFinalizationFence { + &self.fence + } + + pub(crate) const fn issuer_public_key_bytes(&self) -> [u8; 32] { + self.report.issuer_public_key().into_bytes() + } + + pub(crate) const fn claim_mutation_id_bytes(&self) -> [u8; 32] { + *self.report.claim_mutation_id().as_bytes() + } + + pub(crate) const fn report_observed_at_unix_seconds(&self) -> u64 { + self.report.observed_at_unix_s() + } + + pub(crate) fn supersession_bytes(&self) -> Option<([u8; 32], [u8; 32])> { + self.report + .supersession() + .map(|value| (*value.report_id().as_bytes(), *value.event_id().as_bytes())) + } } impl fmt::Debug for RhiSignedEvidenceAttestation { diff --git a/src/reconciliation_finalization.rs b/src/reconciliation_finalization.rs @@ -160,12 +160,12 @@ impl fmt::Debug for RhiReconciliationFinalizationFence { } #[derive(Clone, Copy)] -struct FinalizationIdentity { - attempt_id: RhiReconciliationAttemptId, - job_id: RhiReconciliationJobId, - trade_id: TradeId, - generation: u64, - policy_digest: RhiEvidencePolicyDigest, +pub(crate) struct FinalizationIdentity { + pub(crate) attempt_id: RhiReconciliationAttemptId, + pub(crate) job_id: RhiReconciliationJobId, + pub(crate) trade_id: TradeId, + pub(crate) generation: u64, + pub(crate) policy_digest: RhiEvidencePolicyDigest, } #[derive(Clone, Copy, Debug, PartialEq, Eq)] @@ -209,6 +209,12 @@ impl RhiReconciliationAttemptRepository<'_> { } } +impl RhiReconciliationFinalizationFence { + pub(crate) const fn validation_parts(&self) -> (RhiReconciliationLease, FinalizationIdentity) { + (self.lease, self.identity) + } +} + pub(crate) async fn validate_finalization_fence( transaction: &mut ServiceSqliteTransaction<'_>, fence: &RhiReconciliationFinalizationFence, @@ -254,7 +260,7 @@ fn finalization_identity( }) } -async fn validate_finalization_identity( +pub(crate) async fn validate_finalization_identity( transaction: &mut ServiceSqliteTransaction<'_>, lease: RhiReconciliationLease, identity: FinalizationIdentity, diff --git a/src/reconciliation_finalization_commit.rs b/src/reconciliation_finalization_commit.rs @@ -0,0 +1,1097 @@ +//! Atomic durable finalization of one independently verified attestation. + +use core::fmt; +use std::error::Error; + +use radroots_service_sqlite::{ + ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind, +}; +use radroots_trade::evidence::{ + RadrootsTradeEvidenceOutcomeV1, RadrootsTradeEvidenceSourceCompletionV1, + RadrootsTradeEvidenceSourceRequirementV1, +}; +use sha2::{Digest, Sha256}; + +use crate::{ + RhiPublicationAuthority, RhiPublicationMode, RhiReconciliationAttemptRepository, + RhiReconciliationLease, RhiReconciliationOutcome, RhiReconciliationUnixMilliseconds, + RhiSignedEvidenceAttestation, RhiStateHostMode, + reconciliation_finalization::{FinalizationIdentity, validate_finalization_identity}, +}; + +/// Exact version of the atomic reconciliation-finalization commit contract. +pub const RHI_RECONCILIATION_FINALIZATION_COMMIT_CONTRACT_VERSION: u32 = 1; + +const OUTBOX_ID_DOMAIN: &[u8] = b"radroots.rhi.publication_outbox.v1\0"; + +const INSERT_MANIFEST_SQL: &str = r#"INSERT INTO evidence_manifests ( + manifest_sha256, attempt_id, trade_id, trade_generation, + evidence_policy_sha256, canonical_manifest, observed_at_unix_s, + source_count, observation_count +) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)"#; +const INSERT_PROJECTION_SQL: &str = r#"INSERT INTO trade_projections ( + projection_sha256, manifest_sha256, shared_projection_sha256, + reducer_contract, reducer_contract_version, issue_count +) VALUES (?, ?, ?, 'radroots.trade.reducer.v1', 1, ?)"#; +const INSERT_REPORT_SQL: &str = r#"INSERT INTO attestation_reports ( + statement_sha256, manifest_sha256, projection_sha256, trade_id, + claim_mutation_id, issuer_public_key, outcome, canonical_report, + observed_at_unix_s, supersedes_statement_sha256, supersedes_event_id +) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)"#; +const INSERT_SIGNED_EVENT_SQL: &str = r#"INSERT INTO signed_attestation_events ( + event_id, statement_sha256, event_sha256, issuer_public_key, + authored_at_unix_s, canonical_event_json +) VALUES (?, ?, ?, ?, ?, ?)"#; +const INSERT_OUTBOX_SQL: &str = r#"INSERT INTO publication_outbox ( + outbox_id, event_id, event_sha256, publication_authority_sha256, + target_set_sha256, target_count, required_target_count, max_attempts, + initial_backoff_ms, maximum_backoff_ms, attempt_deadline_ms, + state, revision, next_attempt_unix_ms, lease_owner, + lease_expires_unix_ms, created_at_unix_ms, updated_at_unix_ms +) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, + 'pending', 1, ?, NULL, NULL, ?, ?)"#; +const INSERT_TARGET_SQL: &str = r#"INSERT INTO publication_targets ( + outbox_id, target_ordinal, relay_id, required, state, revision, + attempt_count, next_attempt_unix_ms, last_attempt_id, updated_at_unix_ms +) VALUES (?, ?, ?, ?, 'pending', 1, 0, ?, NULL, ?)"#; +const COMPLETE_JOB_SQL: &str = r#"UPDATE reconciliation_jobs +SET state = 'completed', revision = revision + 1, next_attempt_unix_ms = NULL, + lease_owner = NULL, lease_expires_unix_ms = NULL, updated_at_unix_ms = ? +WHERE job_id = ? AND state = 'leased' AND revision = ? + AND lease_owner = ? AND lease_expires_unix_ms = ? + AND updated_at_unix_ms <= ?"#; + +/// Stable source-free atomic-finalization failure class. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum RhiReconciliationFinalizationCommitErrorKind { + InvalidMode, + InvalidInput, + LeaseLost, + GenerationConflict, + AttemptUnavailable, + PublicationQueueFull, + Conflict, + Storage, + CommitOutcomeUnknown, +} + +impl RhiReconciliationFinalizationCommitErrorKind { + /// Returns the stable machine-readable failure code. + #[must_use] + pub const fn code(self) -> &'static str { + match self { + Self::InvalidMode => "reconciliation_finalization_commit_mode_invalid", + Self::InvalidInput => "reconciliation_finalization_commit_input_invalid", + Self::LeaseLost => "reconciliation_finalization_commit_lease_lost", + Self::GenerationConflict => "reconciliation_finalization_commit_generation_conflict", + Self::AttemptUnavailable => "reconciliation_finalization_commit_attempt_unavailable", + Self::PublicationQueueFull => "reconciliation_publication_queue_full", + Self::Conflict => "reconciliation_finalization_commit_conflict", + Self::Storage => "reconciliation_finalization_commit_storage_failed", + Self::CommitOutcomeUnknown => "reconciliation_finalization_commit_outcome_unknown", + } + } +} + +/// Redacted source-free atomic-finalization failure. +#[derive(Clone, Copy, PartialEq, Eq)] +pub struct RhiReconciliationFinalizationCommitError { + kind: RhiReconciliationFinalizationCommitErrorKind, +} + +impl RhiReconciliationFinalizationCommitError { + /// Returns the stable failure class. + #[must_use] + pub const fn kind(self) -> RhiReconciliationFinalizationCommitErrorKind { + 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 RhiReconciliationFinalizationCommitError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str(match self.kind { + RhiReconciliationFinalizationCommitErrorKind::InvalidMode => { + "RHI reconciliation finalization requires writable state" + } + RhiReconciliationFinalizationCommitErrorKind::InvalidInput => { + "RHI reconciliation finalization commit input is invalid" + } + RhiReconciliationFinalizationCommitErrorKind::LeaseLost => { + "RHI reconciliation finalization lease is no longer authoritative" + } + RhiReconciliationFinalizationCommitErrorKind::GenerationConflict => { + "RHI reconciliation finalization generation changed" + } + RhiReconciliationFinalizationCommitErrorKind::AttemptUnavailable => { + "RHI reconciliation finalization attempt is unavailable" + } + RhiReconciliationFinalizationCommitErrorKind::PublicationQueueFull => { + "RHI publication queue capacity is exhausted" + } + RhiReconciliationFinalizationCommitErrorKind::Conflict => { + "RHI reconciliation finalization conflicts with durable state" + } + RhiReconciliationFinalizationCommitErrorKind::Storage => { + "RHI reconciliation finalization transaction failed" + } + RhiReconciliationFinalizationCommitErrorKind::CommitOutcomeUnknown => { + "RHI reconciliation finalization outcome is unknown" + } + }) + } +} + +impl fmt::Debug for RhiReconciliationFinalizationCommitError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("RhiReconciliationFinalizationCommitError") + .field("kind", &self.kind) + .finish() + } +} + +impl Error for RhiReconciliationFinalizationCommitError {} + +/// Durable result of one exact atomic finalization commit. +#[derive(Clone, Copy, PartialEq, Eq)] +pub struct RhiReconciliationFinalizationCommitOutcome { + created: bool, + publication_mode: RhiPublicationMode, + target_count: u8, +} + +impl RhiReconciliationFinalizationCommitOutcome { + /// Reports whether this call created the immutable finalization inventory. + #[must_use] + pub const fn created(self) -> bool { + self.created + } + + /// Returns the explicit configured publication mode committed with the result. + #[must_use] + pub const fn publication_mode(self) -> RhiPublicationMode { + self.publication_mode + } + + /// Returns the exact immutable target count, or zero when publication is disabled. + #[must_use] + pub const fn target_count(self) -> u8 { + self.target_count + } +} + +impl fmt::Debug for RhiReconciliationFinalizationCommitOutcome { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("RhiReconciliationFinalizationCommitOutcome") + .field("created", &self.created) + .field("publication_mode", &self.publication_mode) + .field("target_count", &self.target_count) + .finish() + } +} + +impl RhiReconciliationAttemptRepository<'_> { + /// Atomically persists one sealed signed result and finalizes its exact job. + /// + /// The attestation is borrowed so a caller can reconcile a lost commit + /// acknowledgement by retrying the same sealed value. Exact prior success + /// is recognized before the consumed lease is checked. No relay, source, + /// clock, entropy, filesystem, or task operation occurs in this boundary. + pub async fn commit_finalization( + &self, + attestation: &RhiSignedEvidenceAttestation, + publication: &RhiPublicationAuthority, + now: RhiReconciliationUnixMilliseconds, + ) -> Result<RhiReconciliationFinalizationCommitOutcome, RhiReconciliationFinalizationCommitError> + { + if self.host().mode() != RhiStateHostMode::ReadWriteExisting { + return Err(failure( + RhiReconciliationFinalizationCommitErrorKind::InvalidMode, + )); + } + if publication.configuration_sha256() + != self.host().metadata().configuration_digest().as_bytes() + { + return Err(failure( + RhiReconciliationFinalizationCommitErrorKind::InvalidInput, + )); + } + let record = FinalizationRecord::from_inputs(attestation, publication, now) + .ok_or_else(|| failure(RhiReconciliationFinalizationCommitErrorKind::InvalidInput))?; + self.host() + .sqlite_host() + .transaction(move |transaction| { + Box::pin(async move { commit(transaction, &record).await }) + }) + .await + .map_err(map_transaction_error) + } +} + +struct FinalizationTarget { + ordinal: u8, + relay_id: Box<str>, + required: bool, +} + +struct FinalizationSource { + source_id: Box<str>, + required: bool, + completion: RadrootsTradeEvidenceSourceCompletionV1, + admitted_event_count: u32, +} + +struct RequiredPublication { + outbox_id: [u8; 32], + authority_sha256: [u8; 32], + target_set_sha256: [u8; 32], + queue_capacity: u32, + maximum_attempts: u16, + initial_backoff_ms: u64, + maximum_backoff_ms: u64, + attempt_deadline_ms: u64, + target_count: u8, + targets: Box<[FinalizationTarget]>, +} + +enum FinalizationPublication { + Disabled, + Required(RequiredPublication), +} + +struct FinalizationRecord { + lease: RhiReconciliationLease, + identity: FinalizationIdentity, + now: RhiReconciliationUnixMilliseconds, + manifest_sha256: [u8; 32], + canonical_manifest: Box<[u8]>, + observed_at_unix_s: u64, + source_count: u8, + sources: Box<[FinalizationSource]>, + observation_count: u32, + projection_sha256: [u8; 32], + shared_projection_sha256: [u8; 32], + issue_count: u32, + statement_sha256: [u8; 32], + claim_mutation_id: [u8; 32], + issuer_public_key: [u8; 32], + outcome: &'static str, + canonical_report: Box<[u8]>, + supersession: Option<([u8; 32], [u8; 32])>, + event_id: [u8; 32], + event_sha256: [u8; 32], + authored_at_unix_s: u64, + canonical_event_json: Box<[u8]>, + publication: FinalizationPublication, +} + +impl FinalizationRecord { + fn from_inputs( + attestation: &RhiSignedEvidenceAttestation, + publication: &RhiPublicationAuthority, + now: RhiReconciliationUnixMilliseconds, + ) -> Option<Self> { + let fence = attestation.fence(); + let (lease, identity) = fence.validation_parts(); + let evaluation = fence.evaluation(); + let projection = evaluation.projection(); + let manifest = projection.manifest(); + let projection_sha256 = projection.digest()?; + let shared_projection_sha256 = projection.shared_projection_digest()?; + let source_count = u8::try_from(manifest.source_count()).ok()?; + let sources = manifest + .inner() + .sources() + .iter() + .map(|source| { + let result = source.result(); + FinalizationSource { + source_id: source.source_id().as_str().into(), + required: matches!( + result.requirement(), + RadrootsTradeEvidenceSourceRequirementV1::Required + ), + completion: result.completion(), + admitted_event_count: result.admitted_event_count(), + } + }) + .collect::<Vec<_>>() + .into_boxed_slice(); + let observation_count = u32::try_from(manifest.observation_count()).ok()?; + let issue_count = u32::try_from(projection.issue_count()).ok()?; + let observed_at_unix_s = manifest.observed_at_unix_seconds(); + if attestation.report_observed_at_unix_seconds() != observed_at_unix_s + || [ + now.get(), + observed_at_unix_s, + attestation.created_at_unix_seconds(), + manifest.trade_generation(), + ] + .iter() + .any(|value| i64::try_from(*value).is_err()) + { + return None; + } + let publication = match publication.mode() { + RhiPublicationMode::Disabled => { + if !publication.targets().is_empty() || publication.retry_policy().is_some() { + return None; + } + FinalizationPublication::Disabled + } + RhiPublicationMode::Required => { + let retry = publication.retry_policy()?; + let targets = publication + .targets() + .iter() + .map(|target| FinalizationTarget { + ordinal: target.ordinal(), + relay_id: target.relay_id().into(), + required: target.required(), + }) + .collect::<Vec<_>>(); + if targets.is_empty() + || targets.len() > crate::RHI_PUBLICATION_MAX_TARGETS + || targets + .iter() + .enumerate() + .any(|(ordinal, target)| usize::from(target.ordinal) != ordinal) + { + return None; + } + let target_count = u8::try_from(targets.len()).ok()?; + FinalizationPublication::Required(RequiredPublication { + outbox_id: outbox_id( + attestation.event_id(), + publication.authority_sha256(), + publication.target_set_sha256(), + ), + authority_sha256: *publication.authority_sha256(), + target_set_sha256: *publication.target_set_sha256(), + queue_capacity: publication.queue_capacity(), + maximum_attempts: retry.maximum_attempts(), + initial_backoff_ms: retry.initial_backoff_milliseconds(), + maximum_backoff_ms: retry.maximum_backoff_milliseconds(), + attempt_deadline_ms: retry.attempt_deadline_milliseconds(), + target_count, + targets: targets.into_boxed_slice(), + }) + } + }; + Some(Self { + lease, + identity, + now, + manifest_sha256: manifest.digest(), + canonical_manifest: manifest.canonical_bytes().into(), + observed_at_unix_s, + source_count, + sources, + observation_count, + projection_sha256, + shared_projection_sha256, + issue_count, + statement_sha256: attestation.statement_digest(), + claim_mutation_id: attestation.claim_mutation_id_bytes(), + issuer_public_key: attestation.issuer_public_key_bytes(), + outcome: outcome_code(attestation.outcome()), + canonical_report: attestation.canonical_report_bytes().into(), + supersession: attestation.supersession_bytes(), + event_id: *attestation.event_id(), + event_sha256: *attestation.signed_event_sha256(), + authored_at_unix_s: attestation.created_at_unix_seconds(), + canonical_event_json: attestation.signed_event_bytes().into(), + publication, + }) + } + + fn outcome(&self, created: bool) -> RhiReconciliationFinalizationCommitOutcome { + let (publication_mode, target_count) = match &self.publication { + FinalizationPublication::Disabled => (RhiPublicationMode::Disabled, 0), + FinalizationPublication::Required(required) => { + (RhiPublicationMode::Required, required.target_count) + } + }; + RhiReconciliationFinalizationCommitOutcome { + created, + publication_mode, + target_count, + } + } +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum OperationError { + LeaseLost, + GenerationConflict, + AttemptUnavailable, + QueueFull, + Conflict, + Storage, +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum ExistingState { + Absent, + Exact, +} + +async fn commit( + transaction: &mut ServiceSqliteTransaction<'_>, + record: &FinalizationRecord, +) -> Result<RhiReconciliationFinalizationCommitOutcome, OperationError> { + match reconcile_existing(transaction, record).await? { + ExistingState::Exact => return Ok(record.outcome(false)), + ExistingState::Absent => {} + } + validate_finalization_identity(transaction, record.lease, record.identity, record.now) + .await + .map_err(|error| match error { + crate::reconciliation_finalization::FinalizationOperationError::LeaseLost => { + OperationError::LeaseLost + } + crate::reconciliation_finalization::FinalizationOperationError::GenerationConflict => { + OperationError::GenerationConflict + } + crate::reconciliation_finalization::FinalizationOperationError::AttemptUnavailable => { + OperationError::AttemptUnavailable + } + crate::reconciliation_finalization::FinalizationOperationError::Storage => { + OperationError::Storage + } + })?; + validate_source_inventory(transaction, record).await?; + validate_advanced_checkpoints(transaction, record).await?; + validate_supersession(transaction, record).await?; + if let FinalizationPublication::Required(required) = &record.publication { + let active: i64 = + sqlx::query_scalar("SELECT COUNT(*) FROM publication_outbox WHERE state != 'complete'") + .fetch_one(&mut *transaction) + .await + .map_err(|_| OperationError::Storage)?; + let active = u64::try_from(active).map_err(|_| OperationError::Storage)?; + if active >= u64::from(required.queue_capacity) { + return Err(OperationError::QueueFull); + } + } + + insert_finalization(transaction, record).await?; + let lease_job = record.lease.job(); + let result = sqlx::query(COMPLETE_JOB_SQL) + .bind(i64_value(record.now.get())?) + .bind(record.identity.job_id.as_bytes().as_slice()) + .bind(i64_value(lease_job.revision())?) + .bind(record.lease.owner_bytes().as_slice()) + .bind(i64_value(record.lease.lease_expires().get())?) + .bind(i64_value(record.now.get())?) + .execute(&mut *transaction) + .await + .map_err(|_| OperationError::Storage)?; + if result.rows_affected() != 1 { + return Err(OperationError::LeaseLost); + } + Ok(record.outcome(true)) +} + +async fn validate_source_inventory( + transaction: &mut ServiceSqliteTransaction<'_>, + record: &FinalizationRecord, +) -> Result<(), OperationError> { + let attempt_count: i64 = sqlx::query_scalar( + "SELECT COUNT(*) FROM evidence_reconciliations WHERE attempt_id = ? AND source_count = ?", + ) + .bind(record.identity.attempt_id.as_bytes().as_slice()) + .bind(i64::from(record.source_count)) + .fetch_one(&mut *transaction) + .await + .map_err(|_| OperationError::Storage)?; + if attempt_count != 1 { + return Err(OperationError::AttemptUnavailable); + } + let count: i64 = sqlx::query_scalar( + "SELECT COUNT(*) FROM evidence_reconciliation_sources WHERE attempt_id = ?", + ) + .bind(record.identity.attempt_id.as_bytes().as_slice()) + .fetch_one(&mut *transaction) + .await + .map_err(|_| OperationError::Storage)?; + if count != i64::from(record.source_count) { + return Err(OperationError::AttemptUnavailable); + } + for (ordinal, source) in record.sources.iter().enumerate() { + let completion = match source.completion { + RadrootsTradeEvidenceSourceCompletionV1::Complete => "complete", + RadrootsTradeEvidenceSourceCompletionV1::Unsupported => "unsupported", + RadrootsTradeEvidenceSourceCompletionV1::Incomplete => "incomplete", + }; + let source_match: i64 = sqlx::query_scalar( + r#"SELECT COUNT(*) FROM evidence_reconciliation_sources +WHERE attempt_id = ? AND source_ordinal = ? AND source_id = ? + AND required = ? AND accepted_event_count = ? + AND CASE ? + WHEN 'complete' THEN completion = 'complete' + WHEN 'unsupported' THEN completion = 'unsupported' + ELSE completion IN ( + 'incomplete_timeout', 'incomplete_unavailable', + 'incomplete_resource_limit', 'incomplete_unknown' + ) + END"#, + ) + .bind(record.identity.attempt_id.as_bytes().as_slice()) + .bind(i64::try_from(ordinal).map_err(|_| OperationError::Storage)?) + .bind(source.source_id.as_ref()) + .bind(i64::from(source.required)) + .bind(i64::from(source.admitted_event_count)) + .bind(completion) + .fetch_one(&mut *transaction) + .await + .map_err(|_| OperationError::Storage)?; + if source_match != 1 { + return Err(OperationError::AttemptUnavailable); + } + } + Ok(()) +} + +async fn validate_advanced_checkpoints( + transaction: &mut ServiceSqliteTransaction<'_>, + record: &FinalizationRecord, +) -> Result<(), OperationError> { + let bad_checkpoint: i64 = sqlx::query_scalar( + r#"SELECT COUNT(*) +FROM evidence_reconciliation_sources AS source +LEFT JOIN relay_checkpoints AS checkpoint + ON checkpoint.source_id = source.source_id + AND checkpoint.evidence_policy_sha256 = ? + AND checkpoint.trade_id = source.trade_id +WHERE source.attempt_id = ? AND source.checkpoint_advanced = 1 + AND (checkpoint.cursor_created_at_unix_s IS NULL + OR checkpoint.cursor_event_id IS NULL + OR checkpoint.cursor_created_at_unix_s != source.candidate_created_at_unix_s + OR checkpoint.cursor_event_id != source.candidate_event_id)"#, + ) + .bind(record.identity.policy_digest.as_bytes().as_slice()) + .bind(record.identity.attempt_id.as_bytes().as_slice()) + .fetch_one(&mut *transaction) + .await + .map_err(|_| OperationError::Storage)?; + if bad_checkpoint == 0 { + Ok(()) + } else { + Err(OperationError::GenerationConflict) + } +} + +async fn insert_finalization( + transaction: &mut ServiceSqliteTransaction<'_>, + record: &FinalizationRecord, +) -> Result<(), OperationError> { + execute_one( + sqlx::query(INSERT_MANIFEST_SQL) + .bind(record.manifest_sha256.as_slice()) + .bind(record.identity.attempt_id.as_bytes().as_slice()) + .bind(record.identity.trade_id.as_bytes().as_slice()) + .bind(i64_value(record.identity.generation)?) + .bind(record.identity.policy_digest.as_bytes().as_slice()) + .bind(record.canonical_manifest.as_ref()) + .bind(i64_value(record.observed_at_unix_s)?) + .bind(i64::from(record.source_count)) + .bind(i64::from(record.observation_count)), + transaction, + ) + .await?; + execute_one( + sqlx::query(INSERT_PROJECTION_SQL) + .bind(record.projection_sha256.as_slice()) + .bind(record.manifest_sha256.as_slice()) + .bind(record.shared_projection_sha256.as_slice()) + .bind(i64::from(record.issue_count)), + transaction, + ) + .await?; + let supersedes_statement = record.supersession.as_ref().map(|value| value.0.as_slice()); + let supersedes_event = record.supersession.as_ref().map(|value| value.1.as_slice()); + execute_one( + sqlx::query(INSERT_REPORT_SQL) + .bind(record.statement_sha256.as_slice()) + .bind(record.manifest_sha256.as_slice()) + .bind(record.projection_sha256.as_slice()) + .bind(record.identity.trade_id.as_bytes().as_slice()) + .bind(record.claim_mutation_id.as_slice()) + .bind(record.issuer_public_key.as_slice()) + .bind(record.outcome) + .bind(record.canonical_report.as_ref()) + .bind(i64_value(record.observed_at_unix_s)?) + .bind(supersedes_statement) + .bind(supersedes_event), + transaction, + ) + .await?; + execute_one( + sqlx::query(INSERT_SIGNED_EVENT_SQL) + .bind(record.event_id.as_slice()) + .bind(record.statement_sha256.as_slice()) + .bind(record.event_sha256.as_slice()) + .bind(record.issuer_public_key.as_slice()) + .bind(i64_value(record.authored_at_unix_s)?) + .bind(record.canonical_event_json.as_ref()), + transaction, + ) + .await?; + if let FinalizationPublication::Required(required) = &record.publication { + let required_count = required + .targets + .iter() + .filter(|target| target.required) + .count(); + execute_one( + sqlx::query(INSERT_OUTBOX_SQL) + .bind(required.outbox_id.as_slice()) + .bind(record.event_id.as_slice()) + .bind(record.event_sha256.as_slice()) + .bind(required.authority_sha256.as_slice()) + .bind(required.target_set_sha256.as_slice()) + .bind(i64::try_from(required.targets.len()).map_err(|_| OperationError::Storage)?) + .bind(i64::try_from(required_count).map_err(|_| OperationError::Storage)?) + .bind(i64::from(required.maximum_attempts)) + .bind(i64_value(required.initial_backoff_ms)?) + .bind(i64_value(required.maximum_backoff_ms)?) + .bind(i64_value(required.attempt_deadline_ms)?) + .bind(i64_value(record.now.get())?) + .bind(i64_value(record.now.get())?) + .bind(i64_value(record.now.get())?), + transaction, + ) + .await?; + for target in &required.targets { + execute_one( + sqlx::query(INSERT_TARGET_SQL) + .bind(required.outbox_id.as_slice()) + .bind(i64::from(target.ordinal)) + .bind(target.relay_id.as_ref()) + .bind(i64::from(target.required)) + .bind(i64_value(record.now.get())?) + .bind(i64_value(record.now.get())?), + transaction, + ) + .await?; + } + } + Ok(()) +} + +async fn reconcile_existing( + transaction: &mut ServiceSqliteTransaction<'_>, + record: &FinalizationRecord, +) -> Result<ExistingState, OperationError> { + let footprint: i64 = sqlx::query_scalar( + r#"SELECT + (SELECT COUNT(*) FROM evidence_manifests + WHERE manifest_sha256 = ? OR attempt_id = ?) + + (SELECT COUNT(*) FROM trade_projections + WHERE projection_sha256 = ? OR manifest_sha256 = ?) + + (SELECT COUNT(*) FROM attestation_reports WHERE statement_sha256 = ?) + + (SELECT COUNT(*) FROM signed_attestation_events + WHERE event_id = ? OR statement_sha256 = ? OR event_sha256 = ?) + + (SELECT COUNT(*) FROM publication_outbox WHERE event_id = ?) + + (SELECT COUNT(*) FROM reconciliation_jobs + WHERE job_id = ? AND state = 'completed')"#, + ) + .bind(record.manifest_sha256.as_slice()) + .bind(record.identity.attempt_id.as_bytes().as_slice()) + .bind(record.projection_sha256.as_slice()) + .bind(record.manifest_sha256.as_slice()) + .bind(record.statement_sha256.as_slice()) + .bind(record.event_id.as_slice()) + .bind(record.statement_sha256.as_slice()) + .bind(record.event_sha256.as_slice()) + .bind(record.event_id.as_slice()) + .bind(record.identity.job_id.as_bytes().as_slice()) + .fetch_one(&mut *transaction) + .await + .map_err(|_| OperationError::Storage)?; + if footprint == 0 { + return Ok(ExistingState::Absent); + } + if manifest_matches(transaction, record).await? + && projection_matches(transaction, record).await? + && report_matches(transaction, record).await? + && event_matches(transaction, record).await? + && publication_matches(transaction, record).await? + && completed_job_matches(transaction, record).await? + { + validate_source_inventory(transaction, record).await?; + validate_supersession(transaction, record).await?; + Ok(ExistingState::Exact) + } else { + Err(OperationError::Conflict) + } +} + +async fn validate_supersession( + transaction: &mut ServiceSqliteTransaction<'_>, + record: &FinalizationRecord, +) -> Result<(), OperationError> { + let Some((statement_sha256, event_id)) = record.supersession else { + return Ok(()); + }; + let count: i64 = sqlx::query_scalar( + r#"SELECT COUNT(*) +FROM attestation_reports AS report +JOIN signed_attestation_events AS event + ON event.statement_sha256 = report.statement_sha256 +WHERE report.statement_sha256 = ? AND event.event_id = ?"#, + ) + .bind(statement_sha256.as_slice()) + .bind(event_id.as_slice()) + .fetch_one(&mut *transaction) + .await + .map_err(|_| OperationError::Storage)?; + if count == 1 { + Ok(()) + } else { + Err(OperationError::Conflict) + } +} + +async fn manifest_matches( + transaction: &mut ServiceSqliteTransaction<'_>, + record: &FinalizationRecord, +) -> Result<bool, OperationError> { + match_count( + sqlx::query_scalar( + r#"SELECT COUNT(*) FROM evidence_manifests +WHERE manifest_sha256 = ? AND attempt_id = ? AND trade_id = ? + AND trade_generation = ? AND evidence_policy_sha256 = ? + AND canonical_manifest = ? AND observed_at_unix_s = ? + AND source_count = ? AND observation_count = ?"#, + ) + .bind(record.manifest_sha256.as_slice()) + .bind(record.identity.attempt_id.as_bytes().as_slice()) + .bind(record.identity.trade_id.as_bytes().as_slice()) + .bind(i64_value(record.identity.generation)?) + .bind(record.identity.policy_digest.as_bytes().as_slice()) + .bind(record.canonical_manifest.as_ref()) + .bind(i64_value(record.observed_at_unix_s)?) + .bind(i64::from(record.source_count)) + .bind(i64::from(record.observation_count)), + transaction, + ) + .await +} + +async fn projection_matches( + transaction: &mut ServiceSqliteTransaction<'_>, + record: &FinalizationRecord, +) -> Result<bool, OperationError> { + match_count( + sqlx::query_scalar( + r#"SELECT COUNT(*) FROM trade_projections +WHERE projection_sha256 = ? AND manifest_sha256 = ? + AND shared_projection_sha256 = ? + AND reducer_contract = 'radroots.trade.reducer.v1' + AND reducer_contract_version = 1 AND issue_count = ?"#, + ) + .bind(record.projection_sha256.as_slice()) + .bind(record.manifest_sha256.as_slice()) + .bind(record.shared_projection_sha256.as_slice()) + .bind(i64::from(record.issue_count)), + transaction, + ) + .await +} + +async fn report_matches( + transaction: &mut ServiceSqliteTransaction<'_>, + record: &FinalizationRecord, +) -> Result<bool, OperationError> { + let supersedes_statement = record.supersession.as_ref().map(|value| value.0.as_slice()); + let supersedes_event = record.supersession.as_ref().map(|value| value.1.as_slice()); + match_count( + sqlx::query_scalar( + r#"SELECT COUNT(*) FROM attestation_reports +WHERE statement_sha256 = ? AND manifest_sha256 = ? AND projection_sha256 = ? + AND trade_id = ? AND claim_mutation_id = ? AND issuer_public_key = ? + AND outcome = ? AND canonical_report = ? AND observed_at_unix_s = ? + AND supersedes_statement_sha256 IS ? AND supersedes_event_id IS ?"#, + ) + .bind(record.statement_sha256.as_slice()) + .bind(record.manifest_sha256.as_slice()) + .bind(record.projection_sha256.as_slice()) + .bind(record.identity.trade_id.as_bytes().as_slice()) + .bind(record.claim_mutation_id.as_slice()) + .bind(record.issuer_public_key.as_slice()) + .bind(record.outcome) + .bind(record.canonical_report.as_ref()) + .bind(i64_value(record.observed_at_unix_s)?) + .bind(supersedes_statement) + .bind(supersedes_event), + transaction, + ) + .await +} + +async fn event_matches( + transaction: &mut ServiceSqliteTransaction<'_>, + record: &FinalizationRecord, +) -> Result<bool, OperationError> { + match_count( + sqlx::query_scalar( + r#"SELECT COUNT(*) FROM signed_attestation_events +WHERE event_id = ? AND statement_sha256 = ? AND event_sha256 = ? + AND issuer_public_key = ? AND authored_at_unix_s = ? + AND canonical_event_json = ?"#, + ) + .bind(record.event_id.as_slice()) + .bind(record.statement_sha256.as_slice()) + .bind(record.event_sha256.as_slice()) + .bind(record.issuer_public_key.as_slice()) + .bind(i64_value(record.authored_at_unix_s)?) + .bind(record.canonical_event_json.as_ref()), + transaction, + ) + .await +} + +async fn publication_matches( + transaction: &mut ServiceSqliteTransaction<'_>, + record: &FinalizationRecord, +) -> Result<bool, OperationError> { + match &record.publication { + FinalizationPublication::Disabled => { + let count: i64 = + sqlx::query_scalar("SELECT COUNT(*) FROM publication_outbox WHERE event_id = ?") + .bind(record.event_id.as_slice()) + .fetch_one(&mut *transaction) + .await + .map_err(|_| OperationError::Storage)?; + Ok(count == 0) + } + FinalizationPublication::Required(required) => { + let required_count = required + .targets + .iter() + .filter(|target| target.required) + .count(); + let outbox = match_count( + sqlx::query_scalar( + r#"SELECT COUNT(*) FROM publication_outbox +WHERE outbox_id = ? AND event_id = ? AND event_sha256 = ? + AND publication_authority_sha256 = ? AND target_set_sha256 = ? + AND target_count = ? AND required_target_count = ? AND max_attempts = ? + AND initial_backoff_ms = ? AND maximum_backoff_ms = ? + AND attempt_deadline_ms = ?"#, + ) + .bind(required.outbox_id.as_slice()) + .bind(record.event_id.as_slice()) + .bind(record.event_sha256.as_slice()) + .bind(required.authority_sha256.as_slice()) + .bind(required.target_set_sha256.as_slice()) + .bind(i64::try_from(required.targets.len()).map_err(|_| OperationError::Storage)?) + .bind(i64::try_from(required_count).map_err(|_| OperationError::Storage)?) + .bind(i64::from(required.maximum_attempts)) + .bind(i64_value(required.initial_backoff_ms)?) + .bind(i64_value(required.maximum_backoff_ms)?) + .bind(i64_value(required.attempt_deadline_ms)?), + transaction, + ) + .await?; + if !outbox { + return Ok(false); + } + let count: i64 = + sqlx::query_scalar("SELECT COUNT(*) FROM publication_targets WHERE outbox_id = ?") + .bind(required.outbox_id.as_slice()) + .fetch_one(&mut *transaction) + .await + .map_err(|_| OperationError::Storage)?; + if count + != i64::try_from(required.targets.len()).map_err(|_| OperationError::Storage)? + { + return Ok(false); + } + for target in &required.targets { + if !match_count( + sqlx::query_scalar( + r#"SELECT COUNT(*) FROM publication_targets +WHERE outbox_id = ? AND target_ordinal = ? AND relay_id = ? AND required = ?"#, + ) + .bind(required.outbox_id.as_slice()) + .bind(i64::from(target.ordinal)) + .bind(target.relay_id.as_ref()) + .bind(i64::from(target.required)), + transaction, + ) + .await? + { + return Ok(false); + } + } + Ok(true) + } + } +} + +async fn completed_job_matches( + transaction: &mut ServiceSqliteTransaction<'_>, + record: &FinalizationRecord, +) -> Result<bool, OperationError> { + let job = record.lease.job(); + match_count( + sqlx::query_scalar( + r#"SELECT COUNT(*) FROM reconciliation_jobs +WHERE job_id = ? AND trade_id = ? AND input_generation = ? + AND evidence_policy_sha256 = ? AND state = 'completed' + AND revision = ? AND attempt_count = ? AND failure_count = ? + AND max_attempts = ? AND next_attempt_unix_ms IS NULL + AND lease_owner IS NULL AND lease_expires_unix_ms IS NULL"#, + ) + .bind(record.identity.job_id.as_bytes().as_slice()) + .bind(record.identity.trade_id.as_bytes().as_slice()) + .bind(i64_value(record.identity.generation)?) + .bind(record.identity.policy_digest.as_bytes().as_slice()) + .bind(i64_value( + job.revision() + .checked_add(1) + .ok_or(OperationError::Storage)?, + )?) + .bind(i64::from(job.attempt_count())) + .bind(i64::from(job.failure_count())) + .bind(i64::from(job.policy().max_attempts())), + transaction, + ) + .await +} + +async fn execute_one<'query>( + query: sqlx::query::Query<'query, sqlx::Sqlite, sqlx::sqlite::SqliteArguments>, + transaction: &mut ServiceSqliteTransaction<'_>, +) -> Result<(), OperationError> { + let result = query + .execute(&mut *transaction) + .await + .map_err(|_| OperationError::Storage)?; + if result.rows_affected() == 1 { + Ok(()) + } else { + Err(OperationError::Storage) + } +} + +async fn match_count<'query>( + query: sqlx::query::QueryScalar<'query, sqlx::Sqlite, i64, sqlx::sqlite::SqliteArguments>, + transaction: &mut ServiceSqliteTransaction<'_>, +) -> Result<bool, OperationError> { + query + .fetch_one(&mut *transaction) + .await + .map(|count| count == 1) + .map_err(|_| OperationError::Storage) +} + +const fn outcome_code(outcome: RhiReconciliationOutcome) -> &'static str { + match outcome { + RadrootsTradeEvidenceOutcomeV1::Valid => "valid", + RadrootsTradeEvidenceOutcomeV1::Invalid => "invalid", + RadrootsTradeEvidenceOutcomeV1::Indeterminate => "indeterminate", + } +} + +fn outbox_id( + event_id: &[u8; 32], + authority_sha256: &[u8; 32], + target_set_sha256: &[u8; 32], +) -> [u8; 32] { + let mut digest = Sha256::new(); + digest.update(OUTBOX_ID_DOMAIN); + digest.update(event_id); + digest.update(authority_sha256); + digest.update(target_set_sha256); + digest.finalize().into() +} + +fn i64_value(value: u64) -> Result<i64, OperationError> { + i64::try_from(value).map_err(|_| OperationError::Storage) +} + +fn map_transaction_error( + error: ServiceSqliteTransactionError<OperationError>, +) -> RhiReconciliationFinalizationCommitError { + if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown { + return failure(RhiReconciliationFinalizationCommitErrorKind::CommitOutcomeUnknown); + } + failure(match error.operation_error().copied() { + Some(OperationError::LeaseLost) => RhiReconciliationFinalizationCommitErrorKind::LeaseLost, + Some(OperationError::GenerationConflict) => { + RhiReconciliationFinalizationCommitErrorKind::GenerationConflict + } + Some(OperationError::AttemptUnavailable) => { + RhiReconciliationFinalizationCommitErrorKind::AttemptUnavailable + } + Some(OperationError::QueueFull) => { + RhiReconciliationFinalizationCommitErrorKind::PublicationQueueFull + } + Some(OperationError::Conflict) => RhiReconciliationFinalizationCommitErrorKind::Conflict, + Some(OperationError::Storage) | None => { + RhiReconciliationFinalizationCommitErrorKind::Storage + } + }) +} + +const fn failure( + kind: RhiReconciliationFinalizationCommitErrorKind, +) -> RhiReconciliationFinalizationCommitError { + RhiReconciliationFinalizationCommitError { kind } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn outbox_identity_is_exact_and_domain_separated() { + assert_eq!( + outbox_id(&[0x11; 32], &[0x22; 32], &[0x33; 32]), + [ + 0xf7, 0x80, 0x86, 0x80, 0xd5, 0xa6, 0x84, 0x1f, 0x2e, 0xcd, 0xaf, 0xb0, 0xc3, 0x1c, + 0x71, 0xb7, 0x4f, 0x0f, 0xf4, 0x02, 0x7f, 0x85, 0xb2, 0x49, 0x1f, 0x00, 0xc8, 0xb9, + 0xa4, 0xf5, 0x9e, 0x37, + ] + ); + assert_ne!( + outbox_id(&[0x11; 32], &[0x22; 32], &[0x33; 32]), + outbox_id(&[0x11; 32], &[0x22; 32], &[0x34; 32]) + ); + } + + #[test] + fn errors_are_closed_source_free_and_redacted() { + for kind in [ + RhiReconciliationFinalizationCommitErrorKind::InvalidMode, + RhiReconciliationFinalizationCommitErrorKind::InvalidInput, + RhiReconciliationFinalizationCommitErrorKind::LeaseLost, + RhiReconciliationFinalizationCommitErrorKind::GenerationConflict, + RhiReconciliationFinalizationCommitErrorKind::AttemptUnavailable, + RhiReconciliationFinalizationCommitErrorKind::PublicationQueueFull, + RhiReconciliationFinalizationCommitErrorKind::Conflict, + RhiReconciliationFinalizationCommitErrorKind::Storage, + RhiReconciliationFinalizationCommitErrorKind::CommitOutcomeUnknown, + ] { + let error = failure(kind); + assert_eq!(error.kind(), kind); + assert!(error.code().starts_with("reconciliation_")); + assert!(Error::source(&error).is_none()); + let rendered = format!("{error} {error:?}"); + assert!(!rendered.contains("relay-primary")); + assert!(!rendered.contains("11111111")); + assert!(!rendered.contains("SELECT")); + } + } +} diff --git a/src/reconciliation_job.rs b/src/reconciliation_job.rs @@ -462,6 +462,10 @@ pub struct RhiReconciliationJob { } impl RhiReconciliationJob { + pub(crate) const fn policy(self) -> RhiReconciliationJobPolicy { + self.policy + } + pub(crate) const fn attempt_policy_matches(self, expected: RhiReconciliationJobPolicy) -> bool { self.policy.lease_duration_ms == expected.lease_duration_ms && self.policy.lease_renewal_ms == expected.lease_renewal_ms @@ -586,6 +590,10 @@ impl RhiReconciliationLease { self.lease_expires } + pub(crate) const fn owner_bytes(self) -> [u8; LEASE_OWNER_BYTES] { + self.owner.0 + } + #[must_use] pub fn renewal_due(self) -> RhiReconciliationUnixMilliseconds { RhiReconciliationUnixMilliseconds( diff --git a/src/state_metadata.rs b/src/state_metadata.rs @@ -405,7 +405,7 @@ fn require_profile_binding( .ok_or_else(|| RhiStateMetadataError::new(RhiStateMetadataErrorKind::Profile)) } -fn normalized_config_digest( +pub(crate) fn normalized_config_digest( profile: RhiConfigProfile, normalized: &Value, ) -> Result<RhiNormalizedConfigDigest, RhiStateMetadataError> { diff --git a/tests/package_boundary.rs b/tests/package_boundary.rs @@ -18,6 +18,8 @@ const RECONCILIATION_ATTEMPTS: &str = include_str!("../src/reconciliation_attemp const RECONCILIATION_COMMIT: &str = include_str!("../src/reconciliation_commit.rs"); const RECONCILIATION_ATTESTATION: &str = include_str!("../src/reconciliation_attestation.rs"); const RECONCILIATION_FINALIZATION: &str = include_str!("../src/reconciliation_finalization.rs"); +const RECONCILIATION_FINALIZATION_COMMIT: &str = + include_str!("../src/reconciliation_finalization_commit.rs"); const RECONCILIATION_MANIFEST: &str = include_str!("../src/reconciliation_manifest.rs"); const RECONCILIATION_REDUCER: &str = include_str!("../src/reconciliation_reducer.rs"); const RECONCILIATION_JOBS: &str = include_str!("../src/reconciliation_job.rs"); @@ -32,6 +34,8 @@ const RECONCILIATION_ATTESTATION_CONTRACT: &str = include_str!("../contracts/services_hardening/reconciliation_attestation.v1.json"); const RECONCILIATION_FINALIZATION_CONTRACT: &str = include_str!("../contracts/services_hardening/reconciliation_finalization.v1.json"); +const RECONCILIATION_FINALIZATION_COMMIT_CONTRACT: &str = + include_str!("../contracts/services_hardening/reconciliation_finalization_commit.v1.json"); const RECONCILIATION_MANIFEST_CONTRACT: &str = include_str!("../contracts/services_hardening/reconciliation_manifest.v1.json"); const RECONCILIATION_REDUCER_CONTRACT: &str = @@ -59,6 +63,7 @@ const SOURCES: &[&str] = &[ include_str!("../src/reconciliation_attestation.rs"), include_str!("../src/reconciliation_commit.rs"), include_str!("../src/reconciliation_finalization.rs"), + include_str!("../src/reconciliation_finalization_commit.rs"), include_str!("../src/reconciliation_job.rs"), include_str!("../src/reconciliation_manifest.rs"), include_str!("../src/reconciliation_reducer.rs"), @@ -134,6 +139,7 @@ fn state_catalog_module_is_private_and_root_api_is_curated() { "reconciliation_attestation", "reconciliation_commit", "reconciliation_finalization", + "reconciliation_finalization_commit", "reconciliation_job", "reconciliation_manifest", "reconciliation_reducer", @@ -194,6 +200,9 @@ fn state_catalog_module_is_private_and_root_api_is_curated() { "RhiReconciliationFinalizationFence", "RhiReconciliationFinalizationErrorKind", "RHI_RECONCILIATION_FINALIZATION_CONTRACT_VERSION", + "RhiReconciliationFinalizationCommitOutcome", + "RhiReconciliationFinalizationCommitErrorKind", + "RHI_RECONCILIATION_FINALIZATION_COMMIT_CONTRACT_VERSION", "RhiReconciliationManifest", "RhiReconciliationManifestErrorKind", "RhiReconciliationScopePrerequisites", @@ -411,6 +420,42 @@ fn reconciliation_finalization_is_attempt_bound_nonmutating_and_revalidated() { } #[test] +fn reconciliation_finalization_commit_is_atomic_idempotent_and_effect_bounded() { + let contract: serde_json::Value = + serde_json::from_str(RECONCILIATION_FINALIZATION_COMMIT_CONTRACT) + .expect("reconciliation-finalization-commit contract"); + assert_eq!( + contract["schema"], + "radroots.rhi.reconciliation-finalization-commit" + ); + assert_eq!(contract["contract_version"], 1); + assert_eq!(contract["transaction"]["count"], 1); + assert_eq!(contract["effects"]["network"], false); + for required in [ + "pub async fn commit_finalization(", + "reconcile_existing(transaction, record)", + "validate_finalization_identity(transaction, record.lease, record.identity, record.now)", + "validate_source_inventory(transaction, record)", + "validate_advanced_checkpoints(transaction, record)", + "validate_supersession(transaction, record)", + "COMPLETE_JOB_SQL", + ] { + assert!( + RECONCILIATION_FINALIZATION_COMMIT.contains(required), + "atomic finalization is missing {required}" + ); + } + for forbidden in ["tokio::spawn", "SystemTime", "std::fs", "reqwest"] { + assert!( + !RECONCILIATION_FINALIZATION_COMMIT.contains(forbidden), + "atomic finalization gained forbidden authority {forbidden}" + ); + } + assert!(!ROOT.contains("pub mod reconciliation_finalization_commit")); + assert!(!PUBLIC_API.contains("rhi::reconciliation_finalization_commit::")); +} + +#[test] fn reconciliation_manifest_is_canonical_sealed_and_effect_free() { let contract: serde_json::Value = serde_json::from_str(RECONCILIATION_MANIFEST_CONTRACT) .expect("reconciliation-manifest contract"); @@ -552,7 +597,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, 25); + assert_eq!(public_error_count, 26); } #[test] @@ -979,6 +1024,10 @@ fn readme_freezes_the_root_only_boundary_and_exact_baseline() { "## Generation-fenced finalization preflight", "[`reconciliation_finalization.v1.json`](contracts/services_hardening/reconciliation_finalization.v1.json)", "Step 199 must rerun the same validator inside the final", + "## Atomic reconciliation finalization commit", + "[`reconciliation_finalization_commit.v1.json`](contracts/services_hardening/reconciliation_finalization_commit.v1.json)", + "An exact retry returns", + "Disabled publication creates no outbox or target row", "## Canonical signed reconciliation attestation", "[`reconciliation_attestation.v1.json`](contracts/services_hardening/reconciliation_attestation.v1.json)", "## Explicit publication authority and durable schema", @@ -1045,6 +1094,7 @@ fn readme_freezes_the_root_only_boundary_and_exact_baseline() { "admit_rhi_trade_mutation_event", "Persist each canonical", "independently signed Nostr event", + "Commit a signed finalization only through the sealed attempt repository", ] { assert!(AGENTS.contains(required), "AGENTS is missing {required}"); } diff --git a/tests/services_hardening_reconciliation_finalization_commit_contract.rs b/tests/services_hardening_reconciliation_finalization_commit_contract.rs @@ -0,0 +1,72 @@ +#![forbid(unsafe_code)] + +const CONTRACT: &str = + include_str!("../contracts/services_hardening/reconciliation_finalization_commit.v1.json"); +const SOURCE: &str = include_str!("../src/reconciliation_finalization_commit.rs"); +const ROOT: &str = include_str!("../src/lib.rs"); + +#[test] +fn contract_freezes_one_atomic_effect_boundary() { + let contract: serde_json::Value = serde_json::from_str(CONTRACT).expect("contract JSON"); + assert_eq!( + contract["schema"], + "radroots.rhi.reconciliation-finalization-commit" + ); + assert_eq!(contract["contract_version"], 1); + assert_eq!(contract["transaction"]["count"], 1); + assert_eq!( + contract["transaction"]["reconcile_exact_prior_success_before_consumed_lease_validation"], + true + ); + assert_eq!(contract["effects"]["sqlite_write"], true); + assert_eq!(contract["effects"]["network"], false); + assert_eq!(contract["effects"]["task_spawn"], false); + assert_eq!( + contract["publication"]["disabled"], + "no_outbox_and_no_target_rows" + ); +} + +#[test] +fn implementation_retains_exact_bytes_and_has_no_external_effect_authority() { + for required in [ + "pub async fn commit_finalization(", + "reconcile_existing(transaction, record)", + "validate_finalization_identity(transaction, record.lease, record.identity, record.now)", + "validate_source_inventory(transaction, record)", + "validate_advanced_checkpoints(transaction, record)", + "validate_supersession(transaction, record)", + "INSERT_MANIFEST_SQL", + "INSERT_PROJECTION_SQL", + "INSERT_REPORT_SQL", + "INSERT_SIGNED_EVENT_SQL", + "INSERT_OUTBOX_SQL", + "INSERT_TARGET_SQL", + "COMPLETE_JOB_SQL", + "record.canonical_manifest.as_ref()", + "record.canonical_report.as_ref()", + "record.canonical_event_json.as_ref()", + ] { + assert!( + SOURCE.contains(required), + "missing finalization guard {required}" + ); + } + for forbidden in [ + "tokio::spawn", + "spawn_blocking", + "SystemTime", + "std::fs", + "reqwest", + "EventSink", + "publish(", + "send(", + ] { + assert!( + !SOURCE.contains(forbidden), + "finalization gained forbidden effect authority {forbidden}" + ); + } + assert!(ROOT.contains("mod reconciliation_finalization_commit;")); + assert!(!ROOT.contains("pub mod reconciliation_finalization_commit;")); +} diff --git a/tests/services_hardening_reconciliation_jobs.rs b/tests/services_hardening_reconciliation_jobs.rs @@ -10,11 +10,13 @@ use radroots_storage::event::SourceGeneration; use rhi::{ RadrootsHostEnvironment, RadrootsPathResolver, RadrootsPlatform, RhiDecryptedIdentity, RhiEncryptedIdentityProvisioningMaterial, RhiEvidenceAttestationSupersession, - RhiIdentityEnvelopeBinding, RhiReconciliationAttemptErrorKind, RhiReconciliationAttemptPlan, + RhiIdentityEnvelopeBinding, RhiPublicationAuthority, RhiPublicationMode, + RhiReconciliationAttemptErrorKind, RhiReconciliationAttemptPlan, RhiReconciliationAttemptResults, RhiReconciliationAttestationErrorKind, - RhiReconciliationCommitErrorKind, RhiReconciliationFinalizationErrorKind, - RhiReconciliationJobErrorKind, RhiReconciliationJobPolicy, RhiReconciliationJobState, - RhiReconciliationLease, RhiReconciliationLeaseOwner, RhiReconciliationRetryDelayMilliseconds, + RhiReconciliationCommitErrorKind, RhiReconciliationFinalizationCommitErrorKind, + RhiReconciliationFinalizationErrorKind, RhiReconciliationJobErrorKind, + RhiReconciliationJobPolicy, RhiReconciliationJobState, RhiReconciliationLease, + RhiReconciliationLeaseOwner, RhiReconciliationRetryDelayMilliseconds, RhiReconciliationScopePrerequisites, RhiReconciliationSourceReplayPlan, RhiReconciliationSourceResult, RhiReconciliationUnixMilliseconds, RhiRuntimeContext, RhiStateMetadata, RhiTradeMutationAdmissionLimits, RhiTradeMutationAuthoredTimePolicy, @@ -1915,6 +1917,296 @@ async fn signed_attestation_supersession_is_verified_ordered_and_explicit() { )); } +#[tokio::test] +async fn atomic_finalization_commits_exact_required_inventory_and_reconciles_retry() { + let ( + root, + runtime, + metadata, + configuration, + host, + lease, + signed, + publication, + manifest, + _trade, + ) = signed_finalization_fixture("finalization-commit-required", false).await; + let repositories = host.repositories(); + let first = repositories + .reconciliation_attempts() + .commit_finalization(&signed, &publication, now(1_784_347_208_000)) + .await + .expect("atomic finalization"); + assert!(first.created()); + assert_eq!(first.publication_mode(), RhiPublicationMode::Required); + assert_eq!(first.target_count(), 2); + + let mut progressed = fixture_connection(&runtime).await; + let checkpoint = sqlx::query( + r#"UPDATE relay_checkpoints +SET cursor_created_at_unix_s = cursor_created_at_unix_s + 1, + cursor_event_id = ?, revision = revision + 1, + completed_at_unix_s = completed_at_unix_s + 1"#, + ) + .bind([0xfe; 32].as_slice()) + .execute(&mut progressed) + .await + .expect("later checkpoint progression"); + assert_eq!(checkpoint.rows_affected(), 1); + progressed.close().await.expect("progression close"); + + let retry = repositories + .reconciliation_attempts() + .commit_finalization(&signed, &publication, lease.lease_expires()) + .await + .expect("exact retry after consumed lease"); + assert!(!retry.created()); + assert_eq!(retry.publication_mode(), RhiPublicationMode::Required); + assert_eq!(retry.target_count(), 2); + + let mut connection = fixture_connection(&runtime).await; + let counts: (i64, i64, i64, i64, i64, i64, i64) = sqlx::query_as( + r#"SELECT + (SELECT COUNT(*) FROM evidence_manifests), + (SELECT COUNT(*) FROM trade_projections), + (SELECT COUNT(*) FROM attestation_reports), + (SELECT COUNT(*) FROM signed_attestation_events), + (SELECT COUNT(*) FROM publication_outbox), + (SELECT COUNT(*) FROM publication_targets), + (SELECT COUNT(*) FROM reconciliation_jobs WHERE state = 'completed')"#, + ) + .fetch_one(&mut connection) + .await + .expect("final inventory counts"); + assert_eq!(counts, (1, 1, 1, 1, 1, 2, 1)); + let exact: (Vec<u8>, Vec<u8>, Vec<u8>, String, i64) = sqlx::query_as( + r#"SELECT manifest.canonical_manifest, report.canonical_report, + event.canonical_event_json, outbox.state, outbox.target_count +FROM evidence_manifests AS manifest +JOIN attestation_reports AS report + ON report.manifest_sha256 = manifest.manifest_sha256 +JOIN signed_attestation_events AS event + ON event.statement_sha256 = report.statement_sha256 +JOIN publication_outbox AS outbox ON outbox.event_id = event.event_id"#, + ) + .fetch_one(&mut connection) + .await + .expect("exact final inventory"); + assert_eq!(exact.0.as_slice(), manifest.as_ref()); + assert_eq!(exact.1.as_slice(), signed.canonical_report_bytes()); + assert_eq!(exact.2.as_slice(), signed.signed_event_bytes()); + assert_eq!(exact.3, "pending"); + assert_eq!(exact.4, 2); + connection.close().await.expect("fixture close"); + + host.close().await.expect("host close"); + drop(( + signed, + publication, + manifest, + configuration, + metadata, + runtime, + root, + )); +} + +#[tokio::test] +async fn atomic_finalization_disabled_mode_creates_no_publication_rows() { + let ( + root, + runtime, + metadata, + configuration, + host, + _lease, + signed, + publication, + manifest, + _trade, + ) = signed_finalization_fixture("finalization-commit-disabled", true).await; + let outcome = host + .repositories() + .reconciliation_attempts() + .commit_finalization(&signed, &publication, now(1_784_347_208_000)) + .await + .expect("disabled finalization"); + assert!(outcome.created()); + assert_eq!(outcome.publication_mode(), RhiPublicationMode::Disabled); + assert_eq!(outcome.target_count(), 0); + let mut connection = fixture_connection(&runtime).await; + let counts: (i64, i64, i64) = sqlx::query_as( + r#"SELECT + (SELECT COUNT(*) FROM signed_attestation_events), + (SELECT COUNT(*) FROM publication_outbox), + (SELECT COUNT(*) FROM publication_targets)"#, + ) + .fetch_one(&mut connection) + .await + .expect("disabled counts"); + assert_eq!(counts, (1, 0, 0)); + connection.close().await.expect("fixture close"); + host.close().await.expect("host close"); + drop(( + signed, + publication, + manifest, + configuration, + metadata, + runtime, + root, + )); +} + +#[tokio::test] +async fn atomic_finalization_rejects_configuration_and_generation_drift_without_partial_rows() { + let ( + root, + runtime, + metadata, + configuration, + host, + _lease, + signed, + publication, + manifest, + trade, + ) = signed_finalization_fixture("finalization-commit-drift", false).await; + let changed = attestation_configuration( + &runtime, + &Keys::new(SecretKey::from_slice(&attestation_secret()).expect("identity secret")) + .public_key() + .to_hex(), + ) + .replace("samples = 512", "samples = 511"); + let changed = parse_rhi_config_v1(changed.as_bytes(), rhi::RhiConfigProfile::RepoLocal) + .expect("changed configuration"); + let mismatched = RhiPublicationAuthority::from_config(&changed).expect("changed authority"); + assert_eq!(mismatched, publication); + let error = host + .repositories() + .reconciliation_attempts() + .commit_finalization(&signed, &mismatched, now(1_784_347_208_000)) + .await + .expect_err("configuration mismatch"); + assert_eq!( + error.kind(), + RhiReconciliationFinalizationCommitErrorKind::InvalidInput + ); + + write_dirty( + &runtime, + trade, + 3, + *metadata.evidence_policy_digest().as_bytes(), + 1_784_347_208, + ) + .await; + let error = host + .repositories() + .reconciliation_attempts() + .commit_finalization(&signed, &publication, now(1_784_347_208_001)) + .await + .expect_err("generation drift"); + assert_eq!( + error.kind(), + RhiReconciliationFinalizationCommitErrorKind::GenerationConflict + ); + let mut connection = fixture_connection(&runtime).await; + let counts: (i64, i64, i64, i64, i64) = sqlx::query_as( + r#"SELECT + (SELECT COUNT(*) FROM evidence_manifests), + (SELECT COUNT(*) FROM trade_projections), + (SELECT COUNT(*) FROM attestation_reports), + (SELECT COUNT(*) FROM signed_attestation_events), + (SELECT COUNT(*) FROM publication_outbox)"#, + ) + .fetch_one(&mut connection) + .await + .expect("rolled-back inventory"); + assert_eq!(counts, (0, 0, 0, 0, 0)); + connection.close().await.expect("fixture close"); + host.close().await.expect("host close"); + drop(( + signed, + publication, + manifest, + configuration, + metadata, + runtime, + root, + )); +} + +async fn signed_finalization_fixture( + instance: &str, + publication_disabled: bool, +) -> ( + tempfile::TempDir, + RhiRuntimeContext, + RhiStateMetadata, + rhi::RhiConfigDocumentV1, + rhi::RhiStateHost, + RhiReconciliationLease, + rhi::RhiSignedEvidenceAttestation, + RhiPublicationAuthority, + Box<[u8]>, + TradeId, +) { + let expected_public_key = + Keys::new(SecretKey::from_slice(&attestation_secret()).expect("identity secret")) + .public_key() + .to_hex(); + let (root, runtime, metadata, configuration, host, lease, evaluation) = + finalization_fixture_with_source(instance, |runtime| { + let source = attestation_configuration(runtime, &expected_public_key); + if publication_disabled { + let publication = source.find("[publication]").expect("publication section"); + let presence = source.find("[presence]").expect("presence section"); + format!( + "{}[publication]\nmode = \"disabled\"\n\n{}", + &source[..publication], + &source[presence..] + ) + } else { + source + } + }) + .await; + let manifest = evaluation.projection().manifest().canonical_bytes().into(); + let trade = *evaluation.projection().manifest().trade_id(); + let identity = provision_attestation_identity(&runtime, &configuration, &metadata); + let fence = host + .repositories() + .reconciliation_attempts() + .prepare_finalization(lease, evaluation, now(1_784_347_206_000)) + .await + .expect("finalization fence"); + let signed = build_rhi_signed_evidence_attestation( + fence, + &identity, + UnixTimeSeconds::new(1_784_347_207), + &FixedAttestationEntropy(0xb1), + None, + ) + .expect("signed attestation"); + let publication = + RhiPublicationAuthority::from_config(&configuration).expect("publication authority"); + drop(identity); + ( + root, + runtime, + metadata, + configuration, + host, + lease, + signed, + publication, + manifest, + trade, + ) +} + async fn fixture_connection(runtime: &RhiRuntimeContext) -> SqliteConnection { let options = SqliteConnectOptions::new() .filename(runtime.artifacts().state_database())