commit 2689256de9af66aef2a92e1a32377a3bf24562f8
parent ba0e5a17937b62df80eefa46d47e2202ff79ce09
Author: triesap <tyson@radroots.org>
Date: Mon, 24 Aug 2026 06:09:32 +0000
refactor(rhi): bind reconciliation source attempts
- derive attempt selector and request identities from claimed configuration
- enforce lease deadlines policy parity and bounded outcome evidence
- freeze pure result inventory contract with redacted public diagnostics
- add fixed vectors boundary tests documentation and API baseline
Diffstat:
10 files changed, 1638 insertions(+), 9 deletions(-)
diff --git a/AGENTS.md b/AGENTS.md
@@ -116,6 +116,12 @@
- Record each source result with exact selector digest, required/optional
authority, stable request identity, deadline, lookback/cursor, completion
evidence, accepted-event count, safe outcome code, and bounded timing.
+- Derive the bounded canonical per-source request inventory only from an
+ unexpired claimed reconciliation lease and the exact normalized evidence
+ policy. Attempt, selector, and request identities are domain-separated;
+ every deadline is absolute and capped by both configuration and lease
+ expiry; result ingestion consumes at most the configured source count plus
+ one before rejecting missing, duplicate, reordered, or excess evidence.
- Query every configured source outside write transactions. A timeout,
unsupported adapter, partial result, raw upstream error, or unknown
completion never masquerades as success.
diff --git a/README b/README
@@ -134,11 +134,34 @@ tokens and integer UTC milliseconds; expired leases can be reclaimed, while
an expired final attempt becomes exhausted. Failed attempts persist their next
eligible time using caller-injected full jitter bounded by the job's governed
exponential backoff. No ambient clock or entropy is read, and no source or
-network work occurs inside a database transaction. Per-source attempt
-inventory and source execution remain with the next ordered checkpoint. The
-exact machine contract is
+network work occurs inside a database transaction. The exact machine contract
+is
[`reconciliation_jobs.v1.json`](contracts/services_hardening/reconciliation_jobs.v1.json).
+## Bounded reconciliation source attempts
+
+`RhiReconciliationAttemptPlan` derives one canonical source-request inventory
+from an unexpired claimed job lease and the exact normalized evidence policy.
+The attempt, selector, and each source request have domain-separated SHA-256
+identities. Every request binds the source ID, required/optional authority,
+trade selector, absolute deadline, lookback, and configured result limits.
+The claimed job's retained lease, retry, and attempt policy must also equal the
+exact normalized configuration; queue capacity remains admission authority for
+new jobs rather than per-attempt identity. The attempt deadline and every
+source deadline are capped by the retained lease expiry; no implicit clock or
+deadline is consulted.
+
+`RhiReconciliationSourceResult` records one stable completion code, explicit
+start and finish times, and bounded accepted-event count and bytes. Exact EOSE
+completion must occur before the source deadline, timeout begins at the
+deadline, unsupported sources cannot report accepted events, and no outcome
+can be omitted, duplicated, reordered, or appended beyond the configured
+source inventory. Planning and result validation are pure and perform no
+SQLite, source, relay, network, filesystem, task, clock, or entropy operation.
+Cursor derivation, source execution, and durable completion commit remain with
+their later ordered checkpoints. The exact machine contract is
+[`reconciliation_attempts.v1.json`](contracts/services_hardening/reconciliation_attempts.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
@@ -137,6 +137,14 @@ pub enum rhi::RhiPublicationCommandV1
pub rhi::RhiPublicationCommandV1::Backlog
pub rhi::RhiPublicationCommandV1::Retry
pub rhi::RhiPublicationCommandV1::Targets
+pub enum rhi::RhiReconciliationAttemptErrorKind
+pub rhi::RhiReconciliationAttemptErrorKind::InvalidConfiguration
+pub rhi::RhiReconciliationAttemptErrorKind::InvalidInput
+pub rhi::RhiReconciliationAttemptErrorKind::LeaseExpired
+pub rhi::RhiReconciliationAttemptErrorKind::PolicyMismatch
+pub rhi::RhiReconciliationAttemptErrorKind::ResultInventory
+impl rhi::RhiReconciliationAttemptErrorKind
+pub const fn rhi::RhiReconciliationAttemptErrorKind::code(self) -> &'static str
pub enum rhi::RhiReconciliationCommandV1
pub rhi::RhiReconciliationCommandV1::Jobs
pub rhi::RhiReconciliationCommandV1::Refresh
@@ -571,12 +579,45 @@ pub const fn rhi::RhiPublicationTargetRepository<'_>::descriptor(&self) -> rhi::
pub const fn rhi::RhiPublicationTargetRepository<'_>::kind(&self) -> rhi::RhiStateRepositoryKind
impl core::fmt::Debug for rhi::RhiPublicationTargetRepository<'_>
pub fn rhi::RhiPublicationTargetRepository<'_>::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct rhi::RhiReconciliationAttemptError
+impl rhi::RhiReconciliationAttemptError
+pub const fn rhi::RhiReconciliationAttemptError::code(self) -> &'static str
+pub const fn rhi::RhiReconciliationAttemptError::kind(self) -> rhi::RhiReconciliationAttemptErrorKind
+impl core::error::Error for rhi::RhiReconciliationAttemptError
+impl core::fmt::Debug for rhi::RhiReconciliationAttemptError
+pub fn rhi::RhiReconciliationAttemptError::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+impl core::fmt::Display for rhi::RhiReconciliationAttemptError
+pub fn rhi::RhiReconciliationAttemptError::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct rhi::RhiReconciliationAttemptId(_)
+impl rhi::RhiReconciliationAttemptId
+pub const fn rhi::RhiReconciliationAttemptId::as_bytes(&self) -> &[u8; 32]
+impl core::fmt::Debug for rhi::RhiReconciliationAttemptId
+pub fn rhi::RhiReconciliationAttemptId::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct rhi::RhiReconciliationAttemptPlan
+impl rhi::RhiReconciliationAttemptPlan
+pub const fn rhi::RhiReconciliationAttemptPlan::attempt_started_at(&self) -> rhi::RhiReconciliationUnixMilliseconds
+pub const fn rhi::RhiReconciliationAttemptPlan::deadline(&self) -> rhi::RhiReconciliationUnixMilliseconds
+pub const fn rhi::RhiReconciliationAttemptPlan::evidence_policy_digest(&self) -> rhi::RhiEvidencePolicyDigest
+pub fn rhi::RhiReconciliationAttemptPlan::from_claim(rhi::RhiReconciliationLease, &rhi::RhiConfigDocumentV1, rhi::RhiReconciliationUnixMilliseconds) -> core::result::Result<Self, rhi::RhiReconciliationAttemptError>
+pub const fn rhi::RhiReconciliationAttemptPlan::id(&self) -> rhi::RhiReconciliationAttemptId
+pub const fn rhi::RhiReconciliationAttemptPlan::input_generation(&self) -> u64
+pub const fn rhi::RhiReconciliationAttemptPlan::job_id(&self) -> rhi::RhiReconciliationJobId
+pub fn rhi::RhiReconciliationAttemptPlan::requests(&self) -> &[rhi::RhiReconciliationSourceRequest]
+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 const fn rhi::RhiReconciliationAttemptRepository<'_>::descriptor(&self) -> rhi::RhiStateRepositoryDescriptor
pub const fn rhi::RhiReconciliationAttemptRepository<'_>::kind(&self) -> rhi::RhiStateRepositoryKind
impl core::fmt::Debug for rhi::RhiReconciliationAttemptRepository<'_>
pub fn rhi::RhiReconciliationAttemptRepository<'_>::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct rhi::RhiReconciliationAttemptResults
+impl rhi::RhiReconciliationAttemptResults
+pub const fn rhi::RhiReconciliationAttemptResults::attempt_id(&self) -> rhi::RhiReconciliationAttemptId
+pub fn rhi::RhiReconciliationAttemptResults::new<I>(&rhi::RhiReconciliationAttemptPlan, I) -> core::result::Result<Self, rhi::RhiReconciliationAttemptError> where I: core::iter::traits::collect::IntoIterator<Item = rhi::RhiReconciliationSourceResult>
+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::RhiReconciliationJob
impl rhi::RhiReconciliationJob
pub const fn rhi::RhiReconciliationJob::attempt_count(self) -> u16
@@ -648,6 +689,41 @@ 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::RhiReconciliationSourceRequest
+impl rhi::RhiReconciliationSourceRequest
+pub const fn rhi::RhiReconciliationSourceRequest::attempt_started_at(&self) -> rhi::RhiReconciliationUnixMilliseconds
+pub const fn rhi::RhiReconciliationSourceRequest::deadline(&self) -> rhi::RhiReconciliationUnixMilliseconds
+pub const fn rhi::RhiReconciliationSourceRequest::id(&self) -> rhi::RhiReconciliationSourceRequestId
+pub const fn rhi::RhiReconciliationSourceRequest::lookback_seconds(&self) -> u64
+pub const fn rhi::RhiReconciliationSourceRequest::maximum_bytes(&self) -> u64
+pub const fn rhi::RhiReconciliationSourceRequest::maximum_events(&self) -> u32
+pub const fn rhi::RhiReconciliationSourceRequest::required(&self) -> bool
+pub const fn rhi::RhiReconciliationSourceRequest::selector_digest(&self) -> rhi::RhiReconciliationSourceSelectorDigest
+pub fn rhi::RhiReconciliationSourceRequest::source_id(&self) -> &str
+pub const fn rhi::RhiReconciliationSourceRequest::trade_id(&self) -> radroots_event::id::TradeId
+impl core::fmt::Debug for rhi::RhiReconciliationSourceRequest
+pub fn rhi::RhiReconciliationSourceRequest::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct rhi::RhiReconciliationSourceRequestId(_)
+impl rhi::RhiReconciliationSourceRequestId
+pub const fn rhi::RhiReconciliationSourceRequestId::as_bytes(&self) -> &[u8; 32]
+impl core::fmt::Debug for rhi::RhiReconciliationSourceRequestId
+pub fn rhi::RhiReconciliationSourceRequestId::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct rhi::RhiReconciliationSourceResult
+impl rhi::RhiReconciliationSourceResult
+pub const fn rhi::RhiReconciliationSourceResult::accepted_event_bytes(self) -> u64
+pub const fn rhi::RhiReconciliationSourceResult::accepted_event_count(self) -> u32
+pub const fn rhi::RhiReconciliationSourceResult::finished_at(self) -> rhi::RhiReconciliationUnixMilliseconds
+pub fn rhi::RhiReconciliationSourceResult::new(&rhi::RhiReconciliationSourceRequest, rhi::RhiTradeSourceCompletion, rhi::RhiReconciliationUnixMilliseconds, rhi::RhiReconciliationUnixMilliseconds, u32, u64) -> core::result::Result<Self, rhi::RhiReconciliationAttemptError>
+pub const fn rhi::RhiReconciliationSourceResult::outcome(self) -> rhi::RhiTradeSourceCompletion
+pub const fn rhi::RhiReconciliationSourceResult::request_id(self) -> rhi::RhiReconciliationSourceRequestId
+pub const fn rhi::RhiReconciliationSourceResult::started_at(self) -> rhi::RhiReconciliationUnixMilliseconds
+impl core::fmt::Debug for rhi::RhiReconciliationSourceResult
+pub fn rhi::RhiReconciliationSourceResult::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct rhi::RhiReconciliationSourceSelectorDigest(_)
+impl rhi::RhiReconciliationSourceSelectorDigest
+pub const fn rhi::RhiReconciliationSourceSelectorDigest::as_bytes(&self) -> &[u8; 32]
+impl core::fmt::Debug for rhi::RhiReconciliationSourceSelectorDigest
+pub fn rhi::RhiReconciliationSourceSelectorDigest::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
pub struct rhi::RhiReconciliationUnixMilliseconds(_)
impl rhi::RhiReconciliationUnixMilliseconds
pub const fn rhi::RhiReconciliationUnixMilliseconds::get(self) -> u64
@@ -1029,6 +1105,8 @@ pub const rhi::RHI_ENCRYPTED_IDENTITY_ENVELOPE_CONTRACT_VERSION: u32
pub const rhi::RHI_ENCRYPTED_IDENTITY_ENVELOPE_MAX_BYTES: usize
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_JOB_CONTRACT_VERSION: u32
pub const rhi::RHI_RECONCILIATION_JOB_MAX_ACTIVE: u32
pub const rhi::RHI_RUNTIME_ADAPTER_CONTRACT_VERSION: u32
diff --git a/contracts/services_hardening/reconciliation_attempts.v1.json b/contracts/services_hardening/reconciliation_attempts.v1.json
@@ -0,0 +1,118 @@
+{
+ "schema": "radroots.rhi.reconciliation-attempts",
+ "schema_version": 1,
+ "contract_version": 1,
+ "state_schema_version": 5,
+ "configuration_authority": "contracts/services_hardening/evidence_policy.v1.json",
+ "attempt_identity": {
+ "algorithm": "sha256",
+ "domain": "radroots.rhi.reconciliation_attempt.v1\\0",
+ "preimage": ["job_id_32_bytes", "attempt_count_u16_be"]
+ },
+ "selector_identity": {
+ "algorithm": "sha256",
+ "domain": "radroots.rhi.reconciliation_selector.v1\\0",
+ "preimage": [
+ "selector_id_u64_length_framed_utf8",
+ "event_kind_count_u32_be",
+ "event_kinds_u32_be_in_canonical_order",
+ "exact_tag_name_byte_d",
+ "trade_id_16_bytes"
+ ],
+ "selector_id": "trade_mutation_lineage_v1",
+ "event_kinds": [3470, 3471, 3472, 3473, 3474],
+ "exact_tag": "#d",
+ "cursor_binding": "deferred_to_step_189"
+ },
+ "request_identity": {
+ "algorithm": "sha256",
+ "domain": "radroots.rhi.reconciliation_source_request.v1\\0",
+ "preimage": [
+ "attempt_id_32_bytes",
+ "source_id_u64_length_framed_utf8",
+ "required_u8",
+ "selector_digest_32_bytes",
+ "attempt_started_at_unix_ms_u64_be",
+ "absolute_deadline_unix_ms_u64_be",
+ "lookback_seconds_u64_be",
+ "maximum_events_u32_be",
+ "maximum_original_event_bytes_u64_be"
+ ]
+ },
+ "plan": {
+ "admission": "unexpired_claimed_job_lease_exact_normalized_evidence_policy_and_exact_normalized_lease_retry_attempt_policy",
+ "queue_capacity_authority": "job_admission_only_not_per_attempt_identity",
+ "source_order": "source_id_ascending_utf8",
+ "source_count_minimum": 1,
+ "source_count_maximum": 16,
+ "source_kind": "nostr_relay",
+ "attempt_deadline": "minimum_of_attempt_start_plus_configured_attempt_deadline_and_lease_expiry",
+ "source_deadline": "minimum_of_attempt_start_plus_configured_source_deadline_attempt_deadline_and_lease_expiry",
+ "result_event_maximum": 4096,
+ "result_original_event_bytes_maximum": 8388608
+ },
+ "source_request": [
+ "request_id",
+ "source_id",
+ "trade_id",
+ "required",
+ "selector_digest",
+ "attempt_started_at_unix_ms",
+ "absolute_deadline_unix_ms",
+ "lookback_seconds",
+ "maximum_events",
+ "maximum_original_event_bytes"
+ ],
+ "completion_codes": [
+ "complete",
+ "incomplete_timeout",
+ "incomplete_unavailable",
+ "incomplete_resource_limit",
+ "incomplete_unknown",
+ "unsupported"
+ ],
+ "result": {
+ "identity": "exact_request_id",
+ "timing": ["started_at_unix_ms", "finished_at_unix_ms"],
+ "complete": "finished_before_absolute_source_deadline",
+ "timeout": "finished_at_or_after_absolute_source_deadline",
+ "other_incomplete_or_unsupported": "finished_before_absolute_source_deadline",
+ "unsupported_accepted_inventory": "exactly_zero",
+ "accepted_inventory": ["accepted_event_count", "accepted_original_event_bytes"],
+ "inventory_order": "exact_canonical_request_order",
+ "inventory_ingestion_bound": "configured_source_count_plus_one"
+ },
+ "effects": {
+ "sqlite": false,
+ "source_or_relay": false,
+ "network": false,
+ "filesystem": false,
+ "task_spawn": false,
+ "ambient_clock": false,
+ "ambient_entropy": false
+ },
+ "diagnostics": "crate_owned_source_free_redacted",
+ "forbidden": [
+ "caller_selected_attempt_id",
+ "caller_selected_request_id",
+ "implicit_source",
+ "implicit_deadline",
+ "complete_at_or_after_deadline",
+ "timeout_before_deadline",
+ "missing_source_result",
+ "duplicate_source_result",
+ "reordered_source_result",
+ "unbounded_result_inventory",
+ "source_query_inside_transaction"
+ ],
+ "deferred": [
+ "cursor_and_overlap_derivation",
+ "source_execution",
+ "durable_source_result_commit",
+ "checkpoint_advance",
+ "manifest",
+ "reducer",
+ "attestation",
+ "publication"
+ ]
+}
diff --git a/src/lib.rs b/src/lib.rs
@@ -8,6 +8,7 @@ mod config_v1;
mod features;
mod identity_credential;
mod identity_envelope;
+mod reconciliation_attempt;
mod reconciliation_job;
mod runtime_adapters;
mod runtime_context;
@@ -67,6 +68,13 @@ pub use radroots_service_host::{
EntropyError, EntropySource, MonotonicClock, MonotonicClockError, MonotonicDeadline,
MonotonicTime, UnixTimeSeconds, WallClock, WallClockError,
};
+pub use reconciliation_attempt::{
+ RHI_RECONCILIATION_ATTEMPT_CONTRACT_VERSION, RHI_RECONCILIATION_ATTEMPT_MAX_SOURCES,
+ RhiReconciliationAttemptError, RhiReconciliationAttemptErrorKind, RhiReconciliationAttemptId,
+ RhiReconciliationAttemptPlan, RhiReconciliationAttemptResults, RhiReconciliationSourceRequest,
+ RhiReconciliationSourceRequestId, RhiReconciliationSourceResult,
+ RhiReconciliationSourceSelectorDigest,
+};
pub use reconciliation_job::{
RHI_RECONCILIATION_JOB_CONTRACT_VERSION, RHI_RECONCILIATION_JOB_MAX_ACTIVE,
RhiReconciliationJob, RhiReconciliationJobError, RhiReconciliationJobErrorKind,
diff --git a/src/reconciliation_attempt.rs b/src/reconciliation_attempt.rs
@@ -0,0 +1,689 @@
+//! Pure bounded per-source reconciliation-attempt planning and result inventory.
+
+use core::fmt;
+use std::error::Error;
+
+use sha2::{Digest, Sha256};
+
+use crate::{
+ RHI_TRADE_SOURCE_RESULT_MAX_BYTES, RHI_TRADE_SOURCE_RESULT_MAX_EVENTS, RhiConfigDocumentV1,
+ RhiEvidencePolicyDigest, RhiReconciliationJobId, RhiReconciliationJobPolicy,
+ RhiReconciliationJobState, RhiReconciliationLease, RhiReconciliationUnixMilliseconds,
+ RhiTradeSourceCompletion, state_metadata,
+};
+
+/// Exact version of the per-source reconciliation-attempt contract.
+pub const RHI_RECONCILIATION_ATTEMPT_CONTRACT_VERSION: u32 = 1;
+
+/// Maximum configured sources represented by one attempt.
+pub const RHI_RECONCILIATION_ATTEMPT_MAX_SOURCES: usize = 16;
+
+const ATTEMPT_ID_DOMAIN: &[u8] = b"radroots.rhi.reconciliation_attempt.v1\0";
+const SELECTOR_DIGEST_DOMAIN: &[u8] = b"radroots.rhi.reconciliation_selector.v1\0";
+const REQUEST_ID_DOMAIN: &[u8] = b"radroots.rhi.reconciliation_source_request.v1\0";
+const SOURCE_SELECTOR: &str = "trade_mutation_lineage_v1";
+const SOURCE_KIND: &str = "nostr_relay";
+const EVENT_KIND_COUNT: u32 = 5;
+const EVENT_KINDS: [u32; 5] = [3470, 3471, 3472, 3473, 3474];
+const _: [(); EVENT_KIND_COUNT as usize] = [(); EVENT_KINDS.len()];
+const MAX_UNIX_MILLISECONDS: u64 = i64::MAX as u64;
+
+/// Stable source-free attempt-model failure class.
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+pub enum RhiReconciliationAttemptErrorKind {
+ InvalidInput,
+ InvalidConfiguration,
+ PolicyMismatch,
+ LeaseExpired,
+ ResultInventory,
+}
+
+impl RhiReconciliationAttemptErrorKind {
+ /// Returns the stable machine-readable failure code.
+ #[must_use]
+ pub const fn code(self) -> &'static str {
+ match self {
+ Self::InvalidInput => "reconciliation_attempt_input_invalid",
+ Self::InvalidConfiguration => "reconciliation_attempt_configuration_invalid",
+ Self::PolicyMismatch => "reconciliation_attempt_policy_mismatch",
+ Self::LeaseExpired => "reconciliation_attempt_lease_expired",
+ Self::ResultInventory => "reconciliation_attempt_result_inventory_invalid",
+ }
+ }
+}
+
+/// Redacted source-free attempt-model failure.
+#[derive(Clone, Copy, PartialEq, Eq)]
+pub struct RhiReconciliationAttemptError {
+ kind: RhiReconciliationAttemptErrorKind,
+}
+
+impl RhiReconciliationAttemptError {
+ const fn new(kind: RhiReconciliationAttemptErrorKind) -> Self {
+ Self { kind }
+ }
+
+ /// Returns the stable failure class.
+ #[must_use]
+ pub const fn kind(self) -> RhiReconciliationAttemptErrorKind {
+ 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 RhiReconciliationAttemptError {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter.write_str(match self.kind {
+ RhiReconciliationAttemptErrorKind::InvalidInput => {
+ "RHI reconciliation attempt input is invalid"
+ }
+ RhiReconciliationAttemptErrorKind::InvalidConfiguration => {
+ "RHI reconciliation attempt configuration is invalid"
+ }
+ RhiReconciliationAttemptErrorKind::PolicyMismatch => {
+ "RHI reconciliation attempt policy does not match the claimed job"
+ }
+ RhiReconciliationAttemptErrorKind::LeaseExpired => {
+ "RHI reconciliation attempt lease has expired"
+ }
+ RhiReconciliationAttemptErrorKind::ResultInventory => {
+ "RHI reconciliation attempt result inventory is invalid"
+ }
+ })
+ }
+}
+
+impl fmt::Debug for RhiReconciliationAttemptError {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter
+ .debug_struct("RhiReconciliationAttemptError")
+ .field("kind", &self.kind)
+ .finish()
+ }
+}
+
+impl Error for RhiReconciliationAttemptError {}
+
+macro_rules! redacted_digest {
+ ($name:ident, $documentation:literal) => {
+ #[doc = $documentation]
+ #[derive(Clone, Copy, PartialEq, Eq, Hash)]
+ pub struct $name([u8; 32]);
+
+ impl $name {
+ /// Returns the exact identity bytes.
+ #[must_use]
+ pub const fn as_bytes(&self) -> &[u8; 32] {
+ &self.0
+ }
+ }
+
+ impl fmt::Debug for $name {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter.write_str(concat!(stringify!($name), "([redacted])"))
+ }
+ }
+ };
+}
+
+redacted_digest!(
+ RhiReconciliationAttemptId,
+ "Domain-separated identity of one claimed reconciliation attempt."
+);
+redacted_digest!(
+ RhiReconciliationSourceRequestId,
+ "Domain-separated identity of one exact source request."
+);
+redacted_digest!(
+ RhiReconciliationSourceSelectorDigest,
+ "Domain-separated digest of the exact trade-mutation selector."
+);
+
+/// One exact configured source request within a claimed job attempt.
+#[derive(Clone, PartialEq, Eq)]
+pub struct RhiReconciliationSourceRequest {
+ id: RhiReconciliationSourceRequestId,
+ source_id: Box<str>,
+ trade_id: radroots_event::id::TradeId,
+ required: bool,
+ selector_digest: RhiReconciliationSourceSelectorDigest,
+ attempt_started_at: RhiReconciliationUnixMilliseconds,
+ deadline: RhiReconciliationUnixMilliseconds,
+ lookback_seconds: u64,
+ maximum_events: u32,
+ maximum_bytes: u64,
+}
+
+impl RhiReconciliationSourceRequest {
+ /// Returns the exact derived request identity.
+ #[must_use]
+ pub const fn id(&self) -> RhiReconciliationSourceRequestId {
+ self.id
+ }
+
+ /// Returns the validated configured source ID.
+ #[must_use]
+ pub fn source_id(&self) -> &str {
+ &self.source_id
+ }
+
+ /// Returns the exact trade selected by this request.
+ #[must_use]
+ pub const fn trade_id(&self) -> radroots_event::id::TradeId {
+ self.trade_id
+ }
+
+ /// Reports whether this source is required by the evidence policy.
+ #[must_use]
+ pub const fn required(&self) -> bool {
+ self.required
+ }
+
+ /// Returns the exact base-selector digest.
+ #[must_use]
+ pub const fn selector_digest(&self) -> RhiReconciliationSourceSelectorDigest {
+ self.selector_digest
+ }
+
+ /// Returns the injected attempt start in integer UTC milliseconds.
+ #[must_use]
+ pub const fn attempt_started_at(&self) -> RhiReconciliationUnixMilliseconds {
+ self.attempt_started_at
+ }
+
+ /// Returns the absolute source deadline capped by attempt and lease expiry.
+ #[must_use]
+ pub const fn deadline(&self) -> RhiReconciliationUnixMilliseconds {
+ self.deadline
+ }
+
+ /// Returns the configured initial lookback in whole seconds.
+ #[must_use]
+ pub const fn lookback_seconds(&self) -> u64 {
+ self.lookback_seconds
+ }
+
+ /// Returns the configured maximum accepted event count.
+ #[must_use]
+ pub const fn maximum_events(&self) -> u32 {
+ self.maximum_events
+ }
+
+ /// Returns the configured maximum accepted original-event bytes.
+ #[must_use]
+ pub const fn maximum_bytes(&self) -> u64 {
+ self.maximum_bytes
+ }
+}
+
+impl fmt::Debug for RhiReconciliationSourceRequest {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter
+ .debug_struct("RhiReconciliationSourceRequest")
+ .field("identity", &"[redacted]")
+ .field("required", &self.required)
+ .field("attempt_started_at", &self.attempt_started_at)
+ .field("deadline", &self.deadline)
+ .field("lookback_seconds", &self.lookback_seconds)
+ .field("maximum_events", &self.maximum_events)
+ .field("maximum_bytes", &self.maximum_bytes)
+ .finish()
+ }
+}
+
+/// Canonically ordered per-source plan for one claimed durable job attempt.
+#[derive(Clone, PartialEq, Eq)]
+pub struct RhiReconciliationAttemptPlan {
+ id: RhiReconciliationAttemptId,
+ job_id: RhiReconciliationJobId,
+ input_generation: u64,
+ policy_digest: RhiEvidencePolicyDigest,
+ attempt_started_at: RhiReconciliationUnixMilliseconds,
+ deadline: RhiReconciliationUnixMilliseconds,
+ requests: Box<[RhiReconciliationSourceRequest]>,
+}
+
+impl RhiReconciliationAttemptPlan {
+ /// Derives the complete bounded source plan from one unexpired claimed lease.
+ pub fn from_claim(
+ lease: RhiReconciliationLease,
+ configuration: &RhiConfigDocumentV1,
+ attempt_started_at: RhiReconciliationUnixMilliseconds,
+ ) -> Result<Self, RhiReconciliationAttemptError> {
+ let job = lease.job();
+ if job.state() != RhiReconciliationJobState::Leased || job.attempt_count() == 0 {
+ return Err(error(RhiReconciliationAttemptErrorKind::InvalidInput));
+ }
+ if attempt_started_at >= lease.lease_expires() {
+ return Err(error(RhiReconciliationAttemptErrorKind::LeaseExpired));
+ }
+ let normalized = configuration.normalized();
+ let policy_digest = state_metadata::evidence_policy_digest(normalized)
+ .map_err(|_| error(RhiReconciliationAttemptErrorKind::InvalidConfiguration))?;
+ if policy_digest != job.evidence_policy_digest() {
+ return Err(error(RhiReconciliationAttemptErrorKind::PolicyMismatch));
+ }
+ let configured_job_policy =
+ RhiReconciliationJobPolicy::from_configuration(configuration)
+ .map_err(|_| error(RhiReconciliationAttemptErrorKind::InvalidConfiguration))?;
+ if !job.attempt_policy_matches(configured_job_policy) {
+ return Err(error(RhiReconciliationAttemptErrorKind::PolicyMismatch));
+ }
+ let attempt_deadline_ms = bounded_integer(
+ normalized,
+ "/reconciliation/attempt_deadline_ms",
+ 100,
+ 30_000,
+ )?;
+ let deadline = absolute_deadline(
+ attempt_started_at,
+ attempt_deadline_ms,
+ lease.lease_expires(),
+ )?;
+ let maximum_events = bounded_integer(
+ normalized,
+ "/resource_limits/source_results/events",
+ 1,
+ RHI_TRADE_SOURCE_RESULT_MAX_EVENTS as u64,
+ )
+ .and_then(|value| {
+ u32::try_from(value)
+ .map_err(|_| error(RhiReconciliationAttemptErrorKind::InvalidConfiguration))
+ })?;
+ let maximum_bytes = bounded_integer(
+ normalized,
+ "/resource_limits/source_results/bytes",
+ 1,
+ RHI_TRADE_SOURCE_RESULT_MAX_BYTES as u64,
+ )?;
+ let sources = normalized
+ .pointer("/evidence/sources")
+ .and_then(serde_json::Value::as_array)
+ .filter(|sources| {
+ !sources.is_empty() && sources.len() <= RHI_RECONCILIATION_ATTEMPT_MAX_SOURCES
+ })
+ .ok_or_else(|| error(RhiReconciliationAttemptErrorKind::InvalidConfiguration))?;
+ let id = attempt_id(job.id(), job.attempt_count());
+ let selector_digest = selector_digest(job.trade_id());
+ let mut requests = Vec::with_capacity(sources.len());
+ let mut prior_source_id: Option<&str> = None;
+ for source in sources {
+ let source_id = bounded_source_id(source)?;
+ if prior_source_id.is_some_and(|prior| prior >= source_id)
+ || source.pointer("/kind").and_then(serde_json::Value::as_str) != Some(SOURCE_KIND)
+ || source
+ .pointer("/selector")
+ .and_then(serde_json::Value::as_str)
+ != Some(SOURCE_SELECTOR)
+ {
+ return Err(error(
+ RhiReconciliationAttemptErrorKind::InvalidConfiguration,
+ ));
+ }
+ prior_source_id = Some(source_id);
+ let required = source
+ .pointer("/required")
+ .and_then(serde_json::Value::as_bool)
+ .ok_or_else(|| error(RhiReconciliationAttemptErrorKind::InvalidConfiguration))?;
+ let source_deadline_ms = bounded_integer(source, "/deadline_ms", 100, 30_000)?;
+ if source_deadline_ms > attempt_deadline_ms {
+ return Err(error(
+ RhiReconciliationAttemptErrorKind::InvalidConfiguration,
+ ));
+ }
+ let source_deadline =
+ absolute_deadline(attempt_started_at, source_deadline_ms, deadline)?;
+ let lookback_seconds = bounded_integer(source, "/lookback_seconds", 60, 2_678_400)?;
+ requests.push(RhiReconciliationSourceRequest {
+ id: request_id(RequestIdentityMaterial {
+ attempt_id: id,
+ source_id,
+ required,
+ selector_digest,
+ attempt_started_at,
+ deadline: source_deadline,
+ lookback_seconds,
+ maximum_events,
+ maximum_bytes,
+ }),
+ source_id: source_id.into(),
+ trade_id: job.trade_id(),
+ required,
+ selector_digest,
+ attempt_started_at,
+ deadline: source_deadline,
+ lookback_seconds,
+ maximum_events,
+ maximum_bytes,
+ });
+ }
+ Ok(Self {
+ id,
+ job_id: job.id(),
+ input_generation: job.input_generation(),
+ policy_digest,
+ attempt_started_at,
+ deadline,
+ requests: requests.into_boxed_slice(),
+ })
+ }
+
+ /// Returns the exact derived attempt identity.
+ #[must_use]
+ pub const fn id(&self) -> RhiReconciliationAttemptId {
+ self.id
+ }
+
+ /// Returns the claimed durable job identity.
+ #[must_use]
+ pub const fn job_id(&self) -> RhiReconciliationJobId {
+ self.job_id
+ }
+
+ /// Returns the exact dirty generation fenced by the claimed job.
+ #[must_use]
+ pub const fn input_generation(&self) -> u64 {
+ self.input_generation
+ }
+
+ /// Returns the exact normalized evidence-policy digest.
+ #[must_use]
+ pub const fn evidence_policy_digest(&self) -> RhiEvidencePolicyDigest {
+ self.policy_digest
+ }
+
+ /// Returns the injected attempt start in integer UTC milliseconds.
+ #[must_use]
+ pub const fn attempt_started_at(&self) -> RhiReconciliationUnixMilliseconds {
+ self.attempt_started_at
+ }
+
+ /// Returns the absolute attempt deadline capped by lease expiry.
+ #[must_use]
+ pub const fn deadline(&self) -> RhiReconciliationUnixMilliseconds {
+ self.deadline
+ }
+
+ /// Returns the canonical configured-source request inventory.
+ #[must_use]
+ pub fn requests(&self) -> &[RhiReconciliationSourceRequest] {
+ &self.requests
+ }
+}
+
+impl fmt::Debug for RhiReconciliationAttemptPlan {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter
+ .debug_struct("RhiReconciliationAttemptPlan")
+ .field("identity", &"[redacted]")
+ .field("input_generation", &self.input_generation)
+ .field("attempt_started_at", &self.attempt_started_at)
+ .field("deadline", &self.deadline)
+ .field("source_count", &self.requests.len())
+ .finish()
+ }
+}
+
+/// Bounded source result bound to one exact request identity.
+#[derive(Clone, Copy, PartialEq, Eq)]
+pub struct RhiReconciliationSourceResult {
+ request_id: RhiReconciliationSourceRequestId,
+ outcome: RhiTradeSourceCompletion,
+ started_at: RhiReconciliationUnixMilliseconds,
+ finished_at: RhiReconciliationUnixMilliseconds,
+ accepted_event_count: u32,
+ accepted_event_bytes: u64,
+}
+
+impl RhiReconciliationSourceResult {
+ /// Validates exact timing and result bounds for one source request.
+ pub fn new(
+ request: &RhiReconciliationSourceRequest,
+ outcome: RhiTradeSourceCompletion,
+ started_at: RhiReconciliationUnixMilliseconds,
+ finished_at: RhiReconciliationUnixMilliseconds,
+ accepted_event_count: u32,
+ accepted_event_bytes: u64,
+ ) -> Result<Self, RhiReconciliationAttemptError> {
+ let before_deadline = finished_at < request.deadline();
+ if started_at < request.attempt_started_at()
+ || started_at > finished_at
+ || started_at >= request.deadline()
+ || (outcome == RhiTradeSourceCompletion::IncompleteTimeout && before_deadline)
+ || (outcome != RhiTradeSourceCompletion::IncompleteTimeout && !before_deadline)
+ || accepted_event_count > request.maximum_events()
+ || accepted_event_bytes > request.maximum_bytes()
+ || (accepted_event_count == 0) != (accepted_event_bytes == 0)
+ || (outcome == RhiTradeSourceCompletion::Unsupported
+ && (accepted_event_count != 0 || accepted_event_bytes != 0))
+ {
+ return Err(error(RhiReconciliationAttemptErrorKind::InvalidInput));
+ }
+ Ok(Self {
+ request_id: request.id(),
+ outcome,
+ started_at,
+ finished_at,
+ accepted_event_count,
+ accepted_event_bytes,
+ })
+ }
+
+ /// Returns the exact request identity this result satisfies.
+ #[must_use]
+ pub const fn request_id(self) -> RhiReconciliationSourceRequestId {
+ self.request_id
+ }
+
+ /// Returns the stable terminal source-completion classification.
+ #[must_use]
+ pub const fn outcome(self) -> RhiTradeSourceCompletion {
+ self.outcome
+ }
+
+ /// Returns the injected source-operation start in UTC milliseconds.
+ #[must_use]
+ pub const fn started_at(self) -> RhiReconciliationUnixMilliseconds {
+ self.started_at
+ }
+
+ /// Returns the injected source-operation finish in UTC milliseconds.
+ #[must_use]
+ pub const fn finished_at(self) -> RhiReconciliationUnixMilliseconds {
+ self.finished_at
+ }
+
+ /// Returns the bounded accepted event count.
+ #[must_use]
+ pub const fn accepted_event_count(self) -> u32 {
+ self.accepted_event_count
+ }
+
+ /// Returns the bounded accepted original-event bytes.
+ #[must_use]
+ pub const fn accepted_event_bytes(self) -> u64 {
+ self.accepted_event_bytes
+ }
+}
+
+impl fmt::Debug for RhiReconciliationSourceResult {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter
+ .debug_struct("RhiReconciliationSourceResult")
+ .field("request_id", &"[redacted]")
+ .field("outcome", &self.outcome)
+ .field("started_at", &self.started_at)
+ .field("finished_at", &self.finished_at)
+ .field("accepted_event_count", &self.accepted_event_count)
+ .field("accepted_event_bytes", &self.accepted_event_bytes)
+ .finish()
+ }
+}
+
+/// Exact canonically ordered result inventory for one attempt plan.
+#[derive(Clone, PartialEq, Eq)]
+pub struct RhiReconciliationAttemptResults {
+ attempt_id: RhiReconciliationAttemptId,
+ results: Box<[RhiReconciliationSourceResult]>,
+}
+
+impl RhiReconciliationAttemptResults {
+ /// Boundedly ingests exactly one result for each request in canonical order.
+ pub fn new<I>(
+ plan: &RhiReconciliationAttemptPlan,
+ results: I,
+ ) -> Result<Self, RhiReconciliationAttemptError>
+ where
+ I: IntoIterator<Item = RhiReconciliationSourceResult>,
+ {
+ let results = results
+ .into_iter()
+ .take(plan.requests.len().saturating_add(1))
+ .collect::<Vec<_>>();
+ if results.len() != plan.requests.len()
+ || results
+ .iter()
+ .zip(plan.requests.iter())
+ .any(|(result, request)| result.request_id != request.id)
+ {
+ return Err(error(RhiReconciliationAttemptErrorKind::ResultInventory));
+ }
+ Ok(Self {
+ attempt_id: plan.id,
+ results: results.into_boxed_slice(),
+ })
+ }
+
+ /// Returns the exact attempt identity satisfied by this inventory.
+ #[must_use]
+ pub const fn attempt_id(&self) -> RhiReconciliationAttemptId {
+ self.attempt_id
+ }
+
+ /// Returns the canonical exact result inventory.
+ #[must_use]
+ pub fn results(&self) -> &[RhiReconciliationSourceResult] {
+ &self.results
+ }
+}
+
+impl fmt::Debug for RhiReconciliationAttemptResults {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter
+ .debug_struct("RhiReconciliationAttemptResults")
+ .field("attempt_id", &"[redacted]")
+ .field("result_count", &self.results.len())
+ .finish()
+ }
+}
+
+fn bounded_source_id(source: &serde_json::Value) -> Result<&str, RhiReconciliationAttemptError> {
+ source
+ .pointer("/source_id")
+ .and_then(serde_json::Value::as_str)
+ .filter(|value| {
+ !value.is_empty()
+ && value.len() <= 64
+ && value.bytes().enumerate().all(|(index, byte)| {
+ if index == 0 {
+ byte.is_ascii_lowercase()
+ } else {
+ byte.is_ascii_lowercase()
+ || byte.is_ascii_digit()
+ || matches!(byte, b'_' | b'-')
+ }
+ })
+ })
+ .ok_or_else(|| error(RhiReconciliationAttemptErrorKind::InvalidConfiguration))
+}
+
+fn bounded_integer(
+ value: &serde_json::Value,
+ pointer: &str,
+ minimum: u64,
+ maximum: u64,
+) -> Result<u64, RhiReconciliationAttemptError> {
+ value
+ .pointer(pointer)
+ .and_then(serde_json::Value::as_u64)
+ .filter(|value| (minimum..=maximum).contains(value))
+ .ok_or_else(|| error(RhiReconciliationAttemptErrorKind::InvalidConfiguration))
+}
+
+fn absolute_deadline(
+ started_at: RhiReconciliationUnixMilliseconds,
+ duration_ms: u64,
+ ceiling: RhiReconciliationUnixMilliseconds,
+) -> Result<RhiReconciliationUnixMilliseconds, RhiReconciliationAttemptError> {
+ let deadline = started_at
+ .get()
+ .checked_add(duration_ms)
+ .filter(|value| *value <= MAX_UNIX_MILLISECONDS)
+ .map(|value| value.min(ceiling.get()))
+ .filter(|value| *value > started_at.get())
+ .ok_or_else(|| error(RhiReconciliationAttemptErrorKind::InvalidInput))?;
+ RhiReconciliationUnixMilliseconds::new(deadline)
+ .map_err(|_| error(RhiReconciliationAttemptErrorKind::InvalidInput))
+}
+
+fn attempt_id(job_id: RhiReconciliationJobId, attempt_count: u16) -> RhiReconciliationAttemptId {
+ let mut hasher = Sha256::new();
+ hasher.update(ATTEMPT_ID_DOMAIN);
+ hasher.update(job_id.as_bytes());
+ hasher.update(attempt_count.to_be_bytes());
+ RhiReconciliationAttemptId(hasher.finalize().into())
+}
+
+fn selector_digest(trade_id: radroots_event::id::TradeId) -> RhiReconciliationSourceSelectorDigest {
+ let mut hasher = Sha256::new();
+ hasher.update(SELECTOR_DIGEST_DOMAIN);
+ hash_framed(&mut hasher, SOURCE_SELECTOR.as_bytes());
+ hasher.update(EVENT_KIND_COUNT.to_be_bytes());
+ for kind in EVENT_KINDS {
+ hasher.update(kind.to_be_bytes());
+ }
+ hasher.update(b"d");
+ hasher.update(trade_id.as_bytes());
+ RhiReconciliationSourceSelectorDigest(hasher.finalize().into())
+}
+
+struct RequestIdentityMaterial<'source> {
+ attempt_id: RhiReconciliationAttemptId,
+ source_id: &'source str,
+ required: bool,
+ selector_digest: RhiReconciliationSourceSelectorDigest,
+ attempt_started_at: RhiReconciliationUnixMilliseconds,
+ deadline: RhiReconciliationUnixMilliseconds,
+ lookback_seconds: u64,
+ maximum_events: u32,
+ maximum_bytes: u64,
+}
+
+fn request_id(material: RequestIdentityMaterial<'_>) -> RhiReconciliationSourceRequestId {
+ let mut hasher = Sha256::new();
+ hasher.update(REQUEST_ID_DOMAIN);
+ hasher.update(material.attempt_id.as_bytes());
+ hash_framed(&mut hasher, material.source_id.as_bytes());
+ hasher.update([u8::from(material.required)]);
+ hasher.update(material.selector_digest.as_bytes());
+ hasher.update(material.attempt_started_at.get().to_be_bytes());
+ hasher.update(material.deadline.get().to_be_bytes());
+ hasher.update(material.lookback_seconds.to_be_bytes());
+ hasher.update(material.maximum_events.to_be_bytes());
+ hasher.update(material.maximum_bytes.to_be_bytes());
+ RhiReconciliationSourceRequestId(hasher.finalize().into())
+}
+
+fn hash_framed(hasher: &mut Sha256, bytes: &[u8]) {
+ hasher.update(u64::try_from(bytes.len()).unwrap_or(u64::MAX).to_be_bytes());
+ hasher.update(bytes);
+}
+
+const fn error(kind: RhiReconciliationAttemptErrorKind) -> RhiReconciliationAttemptError {
+ RhiReconciliationAttemptError::new(kind)
+}
diff --git a/src/reconciliation_job.rs b/src/reconciliation_job.rs
@@ -462,6 +462,14 @@ pub struct RhiReconciliationJob {
}
impl RhiReconciliationJob {
+ 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
+ && self.policy.max_attempts == expected.max_attempts
+ && self.policy.initial_backoff_ms == expected.initial_backoff_ms
+ && self.policy.maximum_backoff_ms == expected.maximum_backoff_ms
+ }
+
#[must_use]
pub const fn id(self) -> RhiReconciliationJobId {
self.id
diff --git a/tests/package_boundary.rs b/tests/package_boundary.rs
@@ -11,7 +11,10 @@ const RUNTIME_ADAPTERS: &str = include_str!("../src/runtime_adapters.rs");
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_JOBS: &str = include_str!("../src/reconciliation_job.rs");
+const RECONCILIATION_ATTEMPT_CONTRACT: &str =
+ include_str!("../contracts/services_hardening/reconciliation_attempts.v1.json");
const RUNTIME_FOUNDATION_CONTRACT: &str =
include_str!("../contracts/services_hardening/runtime_foundation.v1.json");
const TRADE_INGEST_CONTRACT: &str =
@@ -28,6 +31,7 @@ const SOURCES: &[&str] = &[
include_str!("../src/features/trade_agreement_attestation.rs"),
include_str!("../src/identity_credential.rs"),
include_str!("../src/identity_envelope.rs"),
+ include_str!("../src/reconciliation_attempt.rs"),
include_str!("../src/reconciliation_job.rs"),
include_str!("../src/runtime_context.rs"),
include_str!("../src/runtime_adapters.rs"),
@@ -95,6 +99,7 @@ fn state_catalog_module_is_private_and_root_api_is_curated() {
"features",
"identity_credential",
"identity_envelope",
+ "reconciliation_attempt",
"reconciliation_job",
"runtime_context",
"runtime_adapters",
@@ -131,6 +136,10 @@ fn state_catalog_module_is_private_and_root_api_is_curated() {
"validate_rhi_state_catalogs",
"RhiStateCatalogError",
"RhiRuntimeAdapters",
+ "RhiReconciliationAttemptPlan",
+ "RhiReconciliationSourceRequest",
+ "RhiReconciliationSourceResult",
+ "RhiReconciliationAttemptResults",
"RhiReconciliationJobPolicy",
"RhiReconciliationLease",
"RhiTimeEntropyAdapters",
@@ -218,7 +227,56 @@ 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, 17);
+ assert_eq!(public_error_count, 18);
+}
+
+#[test]
+fn reconciliation_attempts_are_exact_bounded_and_effect_free() {
+ let contract: serde_json::Value = serde_json::from_str(RECONCILIATION_ATTEMPT_CONTRACT)
+ .expect("reconciliation-attempt contract");
+ assert_eq!(contract["schema"], "radroots.rhi.reconciliation-attempts");
+ assert_eq!(contract["contract_version"], 1);
+ assert_eq!(contract["state_schema_version"], 5);
+ assert_eq!(contract["plan"]["source_count_maximum"], 16);
+ assert_eq!(contract["plan"]["result_event_maximum"], 4_096);
+ assert_eq!(
+ contract["plan"]["result_original_event_bytes_maximum"],
+ 8_388_608
+ );
+ assert_eq!(
+ contract["result"]["inventory_ingestion_bound"],
+ "configured_source_count_plus_one"
+ );
+ assert_eq!(contract["effects"]["sqlite"], false);
+ assert_eq!(contract["effects"]["source_or_relay"], false);
+ for required in [
+ "state_metadata::evidence_policy_digest(normalized)",
+ ".take(plan.requests.len().saturating_add(1))",
+ "RhiTradeSourceCompletion::IncompleteTimeout",
+ "RhiReconciliationSourceSelectorDigest",
+ ] {
+ assert!(
+ RECONCILIATION_ATTEMPTS.contains(required),
+ "attempt boundary is missing {required}"
+ );
+ }
+ for forbidden in [
+ "sqlx::",
+ "std::fs",
+ "std::net",
+ "tokio::",
+ "SystemTime",
+ "thread_rng",
+ "OsRng",
+ "RhiTradeSourceCursor",
+ ] {
+ assert!(
+ !RECONCILIATION_ATTEMPTS.contains(forbidden),
+ "attempt boundary gained forbidden authority {forbidden}"
+ );
+ }
+ assert!(!ROOT.contains("pub mod reconciliation_attempt"));
+ assert!(!PUBLIC_API.contains("rhi::reconciliation_attempt::"));
}
#[test]
@@ -524,7 +582,13 @@ fn readme_freezes_the_root_only_boundary_and_exact_baseline() {
"[`trade_source_ingest.v1.json`](contracts/services_hardening/trade_source_ingest.v1.json)",
"## Durable reconciliation jobs",
"[`reconciliation_jobs.v1.json`](contracts/services_hardening/reconciliation_jobs.v1.json)",
+ "## Bounded reconciliation source attempts",
+ "[`reconciliation_attempts.v1.json`](contracts/services_hardening/reconciliation_attempts.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",
+ "Planning and result validation are pure and perform no",
+ "SQLite, source, relay, network, filesystem, task, clock, or entropy operation",
"No ambient clock or entropy is read",
"Only exact-target EOSE before the deadline is complete",
"4,096 distinct signed-event identities (event ID plus signature) and 8 MiB",
diff --git a/tests/services_hardening_reconciliation_attempt_contract.rs b/tests/services_hardening_reconciliation_attempt_contract.rs
@@ -0,0 +1,138 @@
+#![forbid(unsafe_code)]
+
+use rhi::{RHI_RECONCILIATION_ATTEMPT_CONTRACT_VERSION, RHI_RECONCILIATION_ATTEMPT_MAX_SOURCES};
+use serde_json::json;
+
+const CONTRACT: &str =
+ include_str!("../contracts/services_hardening/reconciliation_attempts.v1.json");
+const LIB_SOURCE: &str = include_str!("../src/lib.rs");
+const ATTEMPT_SOURCE: &str = include_str!("../src/reconciliation_attempt.rs");
+const README: &str = include_str!("../README");
+
+#[test]
+fn machine_contract_freezes_the_complete_step_188_boundary() {
+ let contract: serde_json::Value = serde_json::from_str(CONTRACT).expect("contract");
+ assert_eq!(contract["schema"], "radroots.rhi.reconciliation-attempts");
+ assert_eq!(contract["schema_version"], 1);
+ assert_eq!(
+ contract["contract_version"],
+ RHI_RECONCILIATION_ATTEMPT_CONTRACT_VERSION
+ );
+ assert_eq!(contract["state_schema_version"], 5);
+ assert_eq!(
+ contract["configuration_authority"],
+ "contracts/services_hardening/evidence_policy.v1.json"
+ );
+ assert_eq!(
+ contract["attempt_identity"],
+ json!({
+ "algorithm": "sha256",
+ "domain": "radroots.rhi.reconciliation_attempt.v1\\0",
+ "preimage": ["job_id_32_bytes", "attempt_count_u16_be"]
+ })
+ );
+ assert_eq!(
+ contract["selector_identity"]["event_kinds"],
+ json!([3470, 3471, 3472, 3473, 3474])
+ );
+ assert_eq!(contract["selector_identity"]["exact_tag"], "#d");
+ assert_eq!(
+ contract["selector_identity"]["cursor_binding"],
+ "deferred_to_step_189"
+ );
+ assert_eq!(
+ contract["plan"]["source_count_maximum"],
+ RHI_RECONCILIATION_ATTEMPT_MAX_SOURCES
+ );
+ assert_eq!(contract["plan"]["result_event_maximum"], 4_096);
+ assert_eq!(
+ contract["plan"]["result_original_event_bytes_maximum"],
+ 8_388_608
+ );
+ assert!(
+ contract["source_request"]
+ .as_array()
+ .expect("source-request fields")
+ .iter()
+ .any(|field| field == "trade_id")
+ );
+ assert_eq!(
+ contract["completion_codes"],
+ json!([
+ "complete",
+ "incomplete_timeout",
+ "incomplete_unavailable",
+ "incomplete_resource_limit",
+ "incomplete_unknown",
+ "unsupported"
+ ])
+ );
+ assert_eq!(
+ contract["result"]["inventory_ingestion_bound"],
+ "configured_source_count_plus_one"
+ );
+ assert_eq!(contract["effects"]["sqlite"], false);
+ assert_eq!(contract["effects"]["source_or_relay"], false);
+ assert_eq!(contract["effects"]["ambient_clock"], false);
+ assert_eq!(contract["effects"]["ambient_entropy"], false);
+ assert_eq!(
+ contract["deferred"],
+ json!([
+ "cursor_and_overlap_derivation",
+ "source_execution",
+ "durable_source_result_commit",
+ "checkpoint_advance",
+ "manifest",
+ "reducer",
+ "attestation",
+ "publication"
+ ])
+ );
+}
+
+#[test]
+fn model_is_private_bounded_pure_and_documented() {
+ assert!(LIB_SOURCE.contains("mod reconciliation_attempt;"));
+ assert!(!LIB_SOURCE.contains("pub mod reconciliation_attempt;"));
+ for required in [
+ "RhiReconciliationAttemptPlan",
+ "RhiReconciliationSourceRequest",
+ "RhiReconciliationSourceResult",
+ "RhiReconciliationAttemptResults",
+ ] {
+ assert!(
+ LIB_SOURCE.contains(required),
+ "root API is missing {required}"
+ );
+ }
+ assert!(README.contains("## Bounded reconciliation source attempts"));
+ assert!(README.contains(
+ "[`reconciliation_attempts.v1.json`](contracts/services_hardening/reconciliation_attempts.v1.json)"
+ ));
+ for required in [
+ ".take(plan.requests.len().saturating_add(1))",
+ "attempt_started_at >= lease.lease_expires()",
+ "state_metadata::evidence_policy_digest(normalized)",
+ "outcome == RhiTradeSourceCompletion::IncompleteTimeout",
+ ] {
+ assert!(
+ ATTEMPT_SOURCE.contains(required),
+ "attempt boundary is missing {required}"
+ );
+ }
+ for forbidden in [
+ "sqlx::",
+ "std::fs",
+ "std::net",
+ "tokio::",
+ "SystemTime",
+ "thread_rng",
+ "OsRng",
+ "RhiTradeSourceCursor",
+ ] {
+ assert!(
+ !ATTEMPT_SOURCE.contains(forbidden),
+ "attempt model gained deferred authority {forbidden}"
+ );
+ }
+}
diff --git a/tests/services_hardening_reconciliation_jobs.rs b/tests/services_hardening_reconciliation_jobs.rs
@@ -6,10 +6,13 @@ use std::{error::Error, fs, os::unix::fs::PermissionsExt, path::Path};
use radroots_service_sqlite::{MigrationAppliedAtUnixSeconds, MigrationBuildIdentity};
use radroots_storage::event::SourceGeneration;
use rhi::{
- RadrootsHostEnvironment, RadrootsPathResolver, RadrootsPlatform, RhiReconciliationJobErrorKind,
- RhiReconciliationJobPolicy, RhiReconciliationJobState, RhiReconciliationLeaseOwner,
- RhiReconciliationRetryDelayMilliseconds, RhiReconciliationUnixMilliseconds, RhiRuntimeContext,
- RhiStateMetadata, TradeId, initialize_rhi_state, open_rhi_state_inspection,
+ RadrootsHostEnvironment, RadrootsPathResolver, RadrootsPlatform,
+ RhiReconciliationAttemptErrorKind, RhiReconciliationAttemptPlan,
+ RhiReconciliationAttemptResults, RhiReconciliationJobErrorKind, RhiReconciliationJobPolicy,
+ RhiReconciliationJobState, RhiReconciliationLeaseOwner,
+ RhiReconciliationRetryDelayMilliseconds, RhiReconciliationSourceResult,
+ RhiReconciliationUnixMilliseconds, RhiRuntimeContext, RhiStateMetadata,
+ RhiTradeSourceCompletion, TradeId, initialize_rhi_state, open_rhi_state_inspection,
open_rhi_state_read_write, parse_rhi_cli_v1_from, parse_rhi_config_v1,
resolve_rhi_runtime_context,
};
@@ -39,9 +42,16 @@ fn runtime(root: &Path, instance: &str) -> RhiRuntimeContext {
fn metadata(runtime: &RhiRuntimeContext) -> RhiStateMetadata {
let config = parse_rhi_config_v1(EXAMPLE.as_bytes(), rhi::RhiConfigProfile::RepoLocal)
.expect("configuration");
+ metadata_from_config(runtime, &config)
+}
+
+fn metadata_from_config(
+ runtime: &RhiRuntimeContext,
+ config: &rhi::RhiConfigDocumentV1,
+) -> RhiStateMetadata {
RhiStateMetadata::new(
runtime,
- &config,
+ config,
SourceGeneration::new([0x5a; 32]).expect("generation"),
1_725_000_000_000,
)
@@ -143,6 +153,10 @@ fn policy(
.expect("policy")
}
+fn configured_policy(configuration: &rhi::RhiConfigDocumentV1) -> RhiReconciliationJobPolicy {
+ RhiReconciliationJobPolicy::from_configuration(configuration).expect("configured job policy")
+}
+
fn now(value: u64) -> RhiReconciliationUnixMilliseconds {
RhiReconciliationUnixMilliseconds::new(value).expect("time")
}
@@ -443,6 +457,489 @@ async fn concurrent_claims_have_one_winner_and_read_only_state_cannot_mutate() {
inspection.close().await.expect("inspection close");
}
+#[tokio::test]
+async fn attempt_plan_binds_the_exact_claim_policy_sources_and_frozen_identities() {
+ let root = tempfile::tempdir().expect("root");
+ let runtime = runtime(root.path(), "attempt-plan");
+ let configuration = parse_rhi_config_v1(EXAMPLE.as_bytes(), rhi::RhiConfigProfile::RepoLocal)
+ .expect("configuration");
+ let metadata = metadata_from_config(&runtime, &configuration);
+ initialize(&runtime, &metadata).await;
+ let trade = TradeId::from_bytes([0x11; 16]);
+ write_dirty(
+ &runtime,
+ trade,
+ 1,
+ *metadata.evidence_policy_digest().as_bytes(),
+ 1_000,
+ )
+ .await;
+ let host = open_writer(&runtime, &metadata).await;
+ let jobs = host.repositories().reconciliation_jobs();
+ jobs.schedule_trade(trade, configured_policy(&configuration), now(1_000))
+ .await
+ .expect("schedule");
+ let lease = jobs
+ .claim_next(owner(0x71), now(1_000))
+ .await
+ .expect("claim")
+ .expect("job");
+
+ let plan = RhiReconciliationAttemptPlan::from_claim(lease, &configuration, now(1_000))
+ .expect("attempt plan");
+ assert_eq!(plan.job_id(), lease.job().id());
+ assert_eq!(plan.input_generation(), 1);
+ assert_eq!(
+ plan.evidence_policy_digest(),
+ metadata.evidence_policy_digest()
+ );
+ assert_eq!(plan.attempt_started_at(), now(1_000));
+ assert_eq!(plan.deadline(), now(31_000));
+ assert_eq!(
+ lower_hex(plan.job_id().as_bytes()),
+ "4e5ecfbee585698c6a67202b30487249291b446909b16a95fd8ae729c7d51e85"
+ );
+ assert_eq!(
+ lower_hex(plan.id().as_bytes()),
+ "89b61ce985d6f11ed963a0d96a05b80e8961a16749a122010d33e1aaee04fdeb"
+ );
+ let [request] = plan.requests() else {
+ panic!("exact source inventory")
+ };
+ assert_eq!(request.source_id(), "trade-primary");
+ assert_eq!(request.trade_id(), trade);
+ assert!(request.required());
+ assert_eq!(request.attempt_started_at(), now(1_000));
+ assert_eq!(request.deadline(), now(11_000));
+ assert_eq!(request.lookback_seconds(), 86_400);
+ assert_eq!(request.maximum_events(), 4_096);
+ assert_eq!(request.maximum_bytes(), 8_388_608);
+ assert_eq!(
+ lower_hex(request.selector_digest().as_bytes()),
+ "2c489c22515b4db784f1be9ab2c224b578c8d28e95d3921ade420b6aa78345bd"
+ );
+ assert_eq!(
+ lower_hex(request.id().as_bytes()),
+ "ef58f9e8a61f7964734cf0c2aabe0bdb2cbdcec529b16f40ae18ec224d0c896d"
+ );
+ host.close().await.expect("close");
+}
+
+#[tokio::test]
+async fn attempt_plan_rejects_policy_mismatch_and_expired_claim_time() {
+ let root = tempfile::tempdir().expect("root");
+ let runtime = runtime(root.path(), "attempt-policy");
+ let configuration = parse_rhi_config_v1(EXAMPLE.as_bytes(), rhi::RhiConfigProfile::RepoLocal)
+ .expect("configuration");
+ let metadata = metadata_from_config(&runtime, &configuration);
+ initialize(&runtime, &metadata).await;
+ let host = open_writer(&runtime, &metadata).await;
+ let jobs = host.repositories().reconciliation_jobs();
+ let governed_policy = configured_policy(&configuration);
+
+ let mismatch_trade = TradeId::from_bytes([0x61; 16]);
+ write_dirty(&runtime, mismatch_trade, 1, [0x62; 32], 1_000).await;
+ jobs.schedule_trade(mismatch_trade, governed_policy, now(1_000))
+ .await
+ .expect("schedule mismatch");
+ let mismatch = jobs
+ .claim_next(owner(0x63), now(1_000))
+ .await
+ .expect("claim mismatch")
+ .expect("job");
+ assert_eq!(
+ RhiReconciliationAttemptPlan::from_claim(mismatch, &configuration, now(1_000))
+ .expect_err("policy mismatch")
+ .kind(),
+ RhiReconciliationAttemptErrorKind::PolicyMismatch
+ );
+
+ let scheduling_mismatch_trade = TradeId::from_bytes([0x69; 16]);
+ write_dirty(
+ &runtime,
+ scheduling_mismatch_trade,
+ 1,
+ *metadata.evidence_policy_digest().as_bytes(),
+ 1_001,
+ )
+ .await;
+ jobs.schedule_trade(
+ scheduling_mismatch_trade,
+ policy(8, 30_000, 10_000, 3, 250, 30_000),
+ now(1_001),
+ )
+ .await
+ .expect("schedule with mismatched job policy");
+ let scheduling_mismatch = jobs
+ .claim_next(owner(0x6a), now(1_001))
+ .await
+ .expect("claim scheduling mismatch")
+ .expect("job");
+ assert_eq!(
+ RhiReconciliationAttemptPlan::from_claim(scheduling_mismatch, &configuration, now(1_001),)
+ .expect_err("job policy mismatch")
+ .kind(),
+ RhiReconciliationAttemptErrorKind::PolicyMismatch
+ );
+
+ let expired_trade = TradeId::from_bytes([0x64; 16]);
+ write_dirty(
+ &runtime,
+ expired_trade,
+ 1,
+ *metadata.evidence_policy_digest().as_bytes(),
+ 1_002,
+ )
+ .await;
+ jobs.schedule_trade(expired_trade, governed_policy, now(1_002))
+ .await
+ .expect("schedule expired");
+ let expired = jobs
+ .claim_next(owner(0x65), now(1_002))
+ .await
+ .expect("claim expired")
+ .expect("job");
+ assert_eq!(
+ RhiReconciliationAttemptPlan::from_claim(expired, &configuration, expired.lease_expires(),)
+ .expect_err("expired attempt start")
+ .kind(),
+ RhiReconciliationAttemptErrorKind::LeaseExpired
+ );
+ host.close().await.expect("close");
+}
+
+#[tokio::test]
+async fn reclaimed_attempt_changes_identity_and_caps_deadlines_to_each_lease() {
+ let short_lease = EXAMPLE
+ .replace("lease_ms = 30000", "lease_ms = 1000")
+ .replace("lease_renewal_ms = 10000", "lease_renewal_ms = 100");
+ let root = tempfile::tempdir().expect("root");
+ let runtime = runtime(root.path(), "attempt-reclaim");
+ let configuration =
+ parse_rhi_config_v1(short_lease.as_bytes(), rhi::RhiConfigProfile::RepoLocal)
+ .expect("short-lease configuration");
+ let metadata = metadata_from_config(&runtime, &configuration);
+ initialize(&runtime, &metadata).await;
+ let trade = TradeId::from_bytes([0x66; 16]);
+ write_dirty(
+ &runtime,
+ trade,
+ 1,
+ *metadata.evidence_policy_digest().as_bytes(),
+ 1_000,
+ )
+ .await;
+ let host = open_writer(&runtime, &metadata).await;
+ let jobs = host.repositories().reconciliation_jobs();
+ jobs.schedule_trade(trade, configured_policy(&configuration), now(1_000))
+ .await
+ .expect("schedule");
+ let first_lease = jobs
+ .claim_next(owner(0x67), now(1_000))
+ .await
+ .expect("first claim")
+ .expect("job");
+ let first = RhiReconciliationAttemptPlan::from_claim(first_lease, &configuration, now(1_000))
+ .expect("first plan");
+ assert_eq!(first.deadline(), now(2_000));
+ assert_eq!(first.requests()[0].deadline(), now(2_000));
+ let later_start =
+ RhiReconciliationAttemptPlan::from_claim(first_lease, &configuration, now(1_500))
+ .expect("same claim with a later explicit start");
+ assert_eq!(later_start.id(), first.id());
+ assert_eq!(later_start.requests()[0].deadline(), now(2_000));
+ assert_ne!(later_start.requests()[0].id(), first.requests()[0].id());
+
+ let second_lease = jobs
+ .claim_next(owner(0x68), now(2_000))
+ .await
+ .expect("reclaim")
+ .expect("job");
+ let second = RhiReconciliationAttemptPlan::from_claim(second_lease, &configuration, now(2_000))
+ .expect("second plan");
+ assert_eq!(second_lease.job().attempt_count(), 2);
+ assert_eq!(second.deadline(), now(3_000));
+ assert_eq!(second.requests()[0].deadline(), now(3_000));
+ assert_ne!(first.id(), second.id());
+ assert_ne!(first.requests()[0].id(), second.requests()[0].id());
+ assert_eq!(
+ first.requests()[0].selector_digest(),
+ second.requests()[0].selector_digest()
+ );
+ host.close().await.expect("close");
+}
+
+#[tokio::test]
+async fn source_results_enforce_deadline_outcome_and_exact_resource_bounds() {
+ let (root, runtime, metadata, configuration, host, plan) =
+ attempt_fixture("attempt-results").await;
+ let request = &plan.requests()[0];
+ let complete = RhiReconciliationSourceResult::new(
+ request,
+ RhiTradeSourceCompletion::Complete,
+ now(1_000),
+ now(10_999),
+ 4_096,
+ 8_388_608,
+ )
+ .expect("exact maximum result");
+ assert_eq!(complete.request_id(), request.id());
+ assert_eq!(complete.outcome().code(), "complete");
+ assert_eq!(complete.accepted_event_count(), 4_096);
+ assert_eq!(complete.accepted_event_bytes(), 8_388_608);
+ assert_eq!(complete.started_at(), now(1_000));
+ assert_eq!(complete.finished_at(), now(10_999));
+
+ let timeout = RhiReconciliationSourceResult::new(
+ request,
+ RhiTradeSourceCompletion::IncompleteTimeout,
+ now(1_000),
+ now(11_000),
+ 0,
+ 0,
+ )
+ .expect("deadline timeout");
+ assert_eq!(timeout.outcome().code(), "incomplete_timeout");
+ for outcome in [
+ RhiTradeSourceCompletion::IncompleteUnavailable,
+ RhiTradeSourceCompletion::IncompleteResourceLimit,
+ RhiTradeSourceCompletion::IncompleteUnknown,
+ RhiTradeSourceCompletion::Unsupported,
+ ] {
+ let result =
+ RhiReconciliationSourceResult::new(request, outcome, now(1_000), now(1_001), 0, 0)
+ .expect("safe incomplete outcome");
+ assert_eq!(result.outcome(), outcome);
+ }
+ for invalid in [
+ RhiReconciliationSourceResult::new(
+ request,
+ RhiTradeSourceCompletion::Complete,
+ now(1_000),
+ now(11_000),
+ 0,
+ 0,
+ ),
+ RhiReconciliationSourceResult::new(
+ request,
+ RhiTradeSourceCompletion::IncompleteTimeout,
+ now(1_000),
+ now(10_999),
+ 0,
+ 0,
+ ),
+ RhiReconciliationSourceResult::new(
+ request,
+ RhiTradeSourceCompletion::IncompleteTimeout,
+ now(11_000),
+ now(11_000),
+ 0,
+ 0,
+ ),
+ RhiReconciliationSourceResult::new(
+ request,
+ RhiTradeSourceCompletion::Complete,
+ now(1_000),
+ now(1_001),
+ 4_097,
+ 8_388_608,
+ ),
+ RhiReconciliationSourceResult::new(
+ request,
+ RhiTradeSourceCompletion::Complete,
+ now(1_000),
+ now(1_001),
+ 4_096,
+ 8_388_609,
+ ),
+ RhiReconciliationSourceResult::new(
+ request,
+ RhiTradeSourceCompletion::Complete,
+ now(1_000),
+ now(1_001),
+ 1,
+ 0,
+ ),
+ RhiReconciliationSourceResult::new(
+ request,
+ RhiTradeSourceCompletion::Unsupported,
+ now(1_000),
+ now(1_001),
+ 1,
+ 1,
+ ),
+ ] {
+ assert_eq!(
+ invalid.expect_err("invalid result").kind(),
+ RhiReconciliationAttemptErrorKind::InvalidInput
+ );
+ }
+
+ let exact = RhiReconciliationAttemptResults::new(&plan, [complete]).expect("inventory");
+ assert_eq!(exact.attempt_id(), plan.id());
+ assert_eq!(exact.results(), [complete]);
+ assert_eq!(
+ RhiReconciliationAttemptResults::new(&plan, [])
+ .expect_err("missing result")
+ .kind(),
+ RhiReconciliationAttemptErrorKind::ResultInventory
+ );
+ assert_eq!(
+ RhiReconciliationAttemptResults::new(&plan, std::iter::repeat(complete))
+ .expect_err("bounded infinite excess")
+ .kind(),
+ RhiReconciliationAttemptErrorKind::ResultInventory
+ );
+ host.close().await.expect("close");
+ drop((configuration, metadata, runtime, root));
+}
+
+#[tokio::test]
+async fn result_inventory_rejects_reordered_configured_sources() {
+ let multi_source = EXAMPLE
+ .replace(
+ "read = false\nwrite = true\nrequired = false",
+ "read = true\nwrite = true\nrequired = false",
+ )
+ .replace(
+ "[[evidence.sources]]\nsource_id = \"trade-primary\"",
+ "[[evidence.sources]]\nsource_id = \"a-secondary\"\nkind = \"nostr_relay\"\nrelay_id = \"relay-secondary\"\nrequired = false\nselector = \"trade_mutation_lineage_v1\"\ndeadline_ms = 5000\nlookback_seconds = 3600\noverlap_seconds = 60\n\n[[evidence.sources]]\nsource_id = \"trade-primary\"",
+ );
+ let root = tempfile::tempdir().expect("root");
+ let runtime = runtime(root.path(), "attempt-order");
+ let configuration =
+ parse_rhi_config_v1(multi_source.as_bytes(), rhi::RhiConfigProfile::RepoLocal)
+ .expect("multi-source configuration");
+ let metadata = metadata_from_config(&runtime, &configuration);
+ initialize(&runtime, &metadata).await;
+ let trade = TradeId::from_bytes([0x51; 16]);
+ write_dirty(
+ &runtime,
+ trade,
+ 1,
+ *metadata.evidence_policy_digest().as_bytes(),
+ 1_000,
+ )
+ .await;
+ let host = open_writer(&runtime, &metadata).await;
+ let jobs = host.repositories().reconciliation_jobs();
+ jobs.schedule_trade(trade, configured_policy(&configuration), now(1_000))
+ .await
+ .expect("schedule");
+ let lease = jobs
+ .claim_next(owner(0x52), now(1_000))
+ .await
+ .expect("claim")
+ .expect("job");
+ let plan =
+ RhiReconciliationAttemptPlan::from_claim(lease, &configuration, now(1_000)).expect("plan");
+ assert_eq!(
+ plan.requests()
+ .iter()
+ .map(|request| request.source_id())
+ .collect::<Vec<_>>(),
+ ["a-secondary", "trade-primary"]
+ );
+ let mut results = plan
+ .requests()
+ .iter()
+ .map(|request| {
+ RhiReconciliationSourceResult::new(
+ request,
+ RhiTradeSourceCompletion::Complete,
+ now(1_000),
+ now(1_001),
+ 0,
+ 0,
+ )
+ .expect("result")
+ })
+ .collect::<Vec<_>>();
+ assert!(RhiReconciliationAttemptResults::new(&plan, results.clone()).is_ok());
+ results.reverse();
+ assert_eq!(
+ RhiReconciliationAttemptResults::new(&plan, results)
+ .expect_err("reordered")
+ .kind(),
+ RhiReconciliationAttemptErrorKind::ResultInventory
+ );
+ host.close().await.expect("close");
+}
+
+#[tokio::test]
+async fn attempt_diagnostics_are_redacted_and_source_free() {
+ let (root, runtime, metadata, configuration, host, plan) =
+ attempt_fixture("attempt-debug").await;
+ let request = &plan.requests()[0];
+ let result = RhiReconciliationSourceResult::new(
+ request,
+ RhiTradeSourceCompletion::Complete,
+ now(1_000),
+ now(1_001),
+ 1,
+ 16,
+ )
+ .expect("result");
+ let inventory = RhiReconciliationAttemptResults::new(&plan, [result]).expect("inventory");
+ let rendered = format!("{plan:?} {request:?} {result:?} {inventory:?}");
+ for secret in [
+ &lower_hex(plan.id().as_bytes()),
+ &lower_hex(request.id().as_bytes()),
+ &lower_hex(request.selector_digest().as_bytes()),
+ ] {
+ assert!(!rendered.contains(secret));
+ }
+ let error = RhiReconciliationAttemptResults::new(&plan, []).expect_err("error");
+ assert!(Error::source(&error).is_none());
+ assert_eq!(
+ error.code(),
+ "reconciliation_attempt_result_inventory_invalid"
+ );
+ assert!(!format!("{error} {error:?}").contains("trade-primary"));
+ host.close().await.expect("close");
+ drop((configuration, metadata, runtime, root));
+}
+
+async fn attempt_fixture(
+ instance: &str,
+) -> (
+ tempfile::TempDir,
+ RhiRuntimeContext,
+ RhiStateMetadata,
+ rhi::RhiConfigDocumentV1,
+ rhi::RhiStateHost,
+ RhiReconciliationAttemptPlan,
+) {
+ let root = tempfile::tempdir().expect("root");
+ let runtime = runtime(root.path(), instance);
+ let configuration = parse_rhi_config_v1(EXAMPLE.as_bytes(), rhi::RhiConfigProfile::RepoLocal)
+ .expect("configuration");
+ let metadata = metadata_from_config(&runtime, &configuration);
+ initialize(&runtime, &metadata).await;
+ let trade = TradeId::from_bytes([0x41; 16]);
+ write_dirty(
+ &runtime,
+ trade,
+ 1,
+ *metadata.evidence_policy_digest().as_bytes(),
+ 1_000,
+ )
+ .await;
+ let host = open_writer(&runtime, &metadata).await;
+ let jobs = host.repositories().reconciliation_jobs();
+ jobs.schedule_trade(trade, configured_policy(&configuration), now(1_000))
+ .await
+ .expect("schedule");
+ let lease = jobs
+ .claim_next(owner(0x42), now(1_000))
+ .await
+ .expect("claim")
+ .expect("job");
+ let plan =
+ RhiReconciliationAttemptPlan::from_claim(lease, &configuration, now(1_000)).expect("plan");
+ (root, runtime, metadata, configuration, host, plan)
+}
+
#[test]
fn public_inputs_have_exact_bounds_and_diagnostics_are_redacted() {
let maximum = RhiReconciliationJobPolicy::new(65_536, 300_000, 150_000, 100, 60_000, 3_600_000)