commit dd0fa0f4befe1e36509369a3aed7ebd6dd095ef4
parent 2689256de9af66aef2a92e1a32377a3bf24562f8
Author: triesap <tyson@radroots.org>
Date: Mon, 24 Aug 2026 06:50:00 +0000
rhi: bind overlap-safe reconciliation replay
- Scope cursor evidence to exact durable authority.
- Canonicalize bounded replay and first provenance.
- Reject conflicting mutation and signed-event reuse.
- Freeze contracts, tests, documentation, and API surface.
Diffstat:
12 files changed, 1390 insertions(+), 8 deletions(-)
diff --git a/AGENTS.md b/AGENTS.md
@@ -125,6 +125,11 @@
- Query every configured source outside write transactions. A timeout,
unsupported adapter, partial result, raw upstream error, or unknown
completion never masquerades as success.
+- Bind every resumed request to its exact prior authored-time/event-ID cursor,
+ subtract the configured overlap for an inclusive query start, canonicalize
+ distinct signed-event identities, retain the earliest injected provenance,
+ and reject conflicting mutation or signed-event identity reuse before the
+ Step 190 commit boundary.
- Freeze the exact accepted mutation, signed-event, provenance, and per-source
completion inventory in an immutable canonical manifest. Reducers consume a
canonically ordered immutable set with explicit policy digest, reducer
diff --git a/README b/README
@@ -162,6 +162,32 @@ 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).
+## Overlap-safe reconciliation replay
+
+`RhiReconciliationSourceReplayPlan` binds one exact source request to optional
+sealed prior-cursor evidence, configured overlap, and an inclusive query start
+under a domain-separated identity. The evidence is sealed to the source,
+trade, evidence-policy digest, and selector digest, so callers cannot forge or
+relabel progress. Step 189 provides no evidence-minting path: Step 190 alone
+may construct it after the replay and cursor commit atomically. Initial
+queries use the configured lookback; resumed queries subtract the configured
+overlap from the prior authored-time
+cursor so equal-timestamp events remain discoverable. A cursor becomes
+eligible only after exact `complete` source evidence, but remains an in-memory
+candidate that this pure model cannot promote to resumption authority.
+
+The replay inventory consumes at most the request event limit plus one and
+bounds all original event bytes before deduplication. It accepts only the
+request's trade, orders admitted signed events canonically, collapses exact
+event replay, retains the earliest injected source observation, and rejects
+conflicting mutation or signed-event identity reuse. Result counts and bytes
+are derived from the distinct canonical inventory rather than caller claims.
+It performs no SQLite, source, relay, network, filesystem, task, clock, or
+entropy operation. Durable scope revalidation, result/completion persistence,
+and checkpoint advancement remain Step 190 authority. The exact machine
+contract is
+[`reconciliation_replay.v1.json`](contracts/services_hardening/reconciliation_replay.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
@@ -168,6 +168,15 @@ pub rhi::RhiReconciliationJobState::Ready
pub rhi::RhiReconciliationJobState::Superseded
impl rhi::RhiReconciliationJobState
pub const fn rhi::RhiReconciliationJobState::code(self) -> &'static str
+pub enum rhi::RhiReconciliationReplayErrorKind
+pub rhi::RhiReconciliationReplayErrorKind::InvalidConfiguration
+pub rhi::RhiReconciliationReplayErrorKind::InvalidInput
+pub rhi::RhiReconciliationReplayErrorKind::MutationConflict
+pub rhi::RhiReconciliationReplayErrorKind::PolicyMismatch
+pub rhi::RhiReconciliationReplayErrorKind::ResourceLimit
+pub rhi::RhiReconciliationReplayErrorKind::SignedEventConflict
+impl rhi::RhiReconciliationReplayErrorKind
+pub const fn rhi::RhiReconciliationReplayErrorKind::code(self) -> &'static str
pub enum rhi::RhiRuntimeAdapterErrorKind
pub rhi::RhiRuntimeAdapterErrorKind::CredentialAccess
pub rhi::RhiRuntimeAdapterErrorKind::EntropyUnavailable
@@ -681,6 +690,15 @@ impl rhi::RhiReconciliationLeaseOwner
pub fn rhi::RhiReconciliationLeaseOwner::from_bytes([u8; 16]) -> core::result::Result<Self, rhi::RhiReconciliationJobError>
impl core::fmt::Debug for rhi::RhiReconciliationLeaseOwner
pub fn rhi::RhiReconciliationLeaseOwner::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct rhi::RhiReconciliationReplayError
+impl rhi::RhiReconciliationReplayError
+pub const fn rhi::RhiReconciliationReplayError::code(self) -> &'static str
+pub const fn rhi::RhiReconciliationReplayError::kind(self) -> rhi::RhiReconciliationReplayErrorKind
+impl core::error::Error for rhi::RhiReconciliationReplayError
+impl core::fmt::Debug for rhi::RhiReconciliationReplayError
+pub fn rhi::RhiReconciliationReplayError::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+impl core::fmt::Display for rhi::RhiReconciliationReplayError
+pub fn rhi::RhiReconciliationReplayError::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
pub struct rhi::RhiReconciliationRetryDelayMilliseconds(_)
impl rhi::RhiReconciliationRetryDelayMilliseconds
pub const fn rhi::RhiReconciliationRetryDelayMilliseconds::get(self) -> u64
@@ -689,6 +707,39 @@ 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::RhiReconciliationSourceCursorEvidence
+impl rhi::RhiReconciliationSourceCursorEvidence
+pub const fn rhi::RhiReconciliationSourceCursorEvidence::cursor(&self) -> rhi::RhiTradeSourceCursor
+impl core::fmt::Debug for rhi::RhiReconciliationSourceCursorEvidence
+pub fn rhi::RhiReconciliationSourceCursorEvidence::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct rhi::RhiReconciliationSourceReplay
+impl rhi::RhiReconciliationSourceReplay
+pub fn rhi::RhiReconciliationSourceReplay::accepted_event_count(&self) -> usize
+pub const fn rhi::RhiReconciliationSourceReplay::accepted_original_event_bytes(&self) -> u64
+pub const fn rhi::RhiReconciliationSourceReplay::cursor_candidate(&self) -> core::option::Option<rhi::RhiTradeSourceCursor>
+pub const fn rhi::RhiReconciliationSourceReplay::duplicate_observation_count(&self) -> u32
+pub fn rhi::RhiReconciliationSourceReplay::eligible_cursor(&self) -> core::option::Option<rhi::RhiTradeSourceCursor>
+pub const fn rhi::RhiReconciliationSourceReplay::first_observed_at(&self) -> core::option::Option<rhi::RhiTradeMutationObservedAtUnixSeconds>
+pub const fn rhi::RhiReconciliationSourceReplay::id(&self) -> rhi::RhiReconciliationSourceReplayId
+pub const fn rhi::RhiReconciliationSourceReplay::result(&self) -> rhi::RhiReconciliationSourceResult
+impl core::fmt::Debug for rhi::RhiReconciliationSourceReplay
+pub fn rhi::RhiReconciliationSourceReplay::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct rhi::RhiReconciliationSourceReplayId(_)
+impl rhi::RhiReconciliationSourceReplayId
+pub const fn rhi::RhiReconciliationSourceReplayId::as_bytes(&self) -> &[u8; 32]
+impl core::fmt::Debug for rhi::RhiReconciliationSourceReplayId
+pub fn rhi::RhiReconciliationSourceReplayId::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct rhi::RhiReconciliationSourceReplayPlan
+impl rhi::RhiReconciliationSourceReplayPlan
+pub fn rhi::RhiReconciliationSourceReplayPlan::finish<I>(self, &rhi::RhiReconciliationSourceRequest, rhi::RhiTradeSourceCompletion, rhi::RhiReconciliationUnixMilliseconds, rhi::RhiReconciliationUnixMilliseconds, I) -> core::result::Result<rhi::RhiReconciliationSourceReplay, rhi::RhiReconciliationReplayError> where I: core::iter::traits::collect::IntoIterator<Item = rhi::RhiAdmittedTradeMutationEvent>
+pub fn rhi::RhiReconciliationSourceReplayPlan::from_request(&rhi::RhiReconciliationAttemptPlan, &rhi::RhiReconciliationSourceRequest, &rhi::RhiConfigDocumentV1, core::option::Option<rhi::RhiReconciliationSourceCursorEvidence>) -> core::result::Result<Self, rhi::RhiReconciliationReplayError>
+pub const fn rhi::RhiReconciliationSourceReplayPlan::id(&self) -> rhi::RhiReconciliationSourceReplayId
+pub const fn rhi::RhiReconciliationSourceReplayPlan::overlap_seconds(&self) -> u64
+pub fn rhi::RhiReconciliationSourceReplayPlan::prior_cursor(&self) -> core::option::Option<rhi::RhiTradeSourceCursor>
+pub const fn rhi::RhiReconciliationSourceReplayPlan::request_id(&self) -> rhi::RhiReconciliationSourceRequestId
+pub const fn rhi::RhiReconciliationSourceReplayPlan::since_unix_seconds(&self) -> u64
+impl core::fmt::Debug for rhi::RhiReconciliationSourceReplayPlan
+pub fn rhi::RhiReconciliationSourceReplayPlan::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
pub struct rhi::RhiReconciliationSourceRequest
impl rhi::RhiReconciliationSourceRequest
pub const fn rhi::RhiReconciliationSourceRequest::attempt_started_at(&self) -> rhi::RhiReconciliationUnixMilliseconds
@@ -1109,6 +1160,7 @@ 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_RECONCILIATION_REPLAY_CONTRACT_VERSION: u32
pub const rhi::RHI_RUNTIME_ADAPTER_CONTRACT_VERSION: u32
pub const rhi::RHI_RUNTIME_FOUNDATION_CONTRACT_VERSION: u32
pub const rhi::RHI_RUNTIME_JITTER_MAX_ENTROPY_DRAWS: usize
diff --git a/contracts/services_hardening/reconciliation_attempts.v1.json b/contracts/services_hardening/reconciliation_attempts.v1.json
@@ -22,7 +22,7 @@
"selector_id": "trade_mutation_lineage_v1",
"event_kinds": [3470, 3471, 3472, 3473, 3474],
"exact_tag": "#d",
- "cursor_binding": "deferred_to_step_189"
+ "cursor_binding": "contracts/services_hardening/reconciliation_replay.v1.json"
},
"request_identity": {
"algorithm": "sha256",
@@ -106,7 +106,6 @@
"source_query_inside_transaction"
],
"deferred": [
- "cursor_and_overlap_derivation",
"source_execution",
"durable_source_result_commit",
"checkpoint_advance",
diff --git a/contracts/services_hardening/reconciliation_replay.v1.json b/contracts/services_hardening/reconciliation_replay.v1.json
@@ -0,0 +1,98 @@
+{
+ "schema": "radroots.rhi.reconciliation-replay",
+ "schema_version": 1,
+ "contract_version": 1,
+ "state_schema_version": 5,
+ "configuration_authority": "contracts/services_hardening/evidence_policy.v1.json",
+ "request_authority": "contracts/services_hardening/reconciliation_attempts.v1.json",
+ "replay_identity": {
+ "algorithm": "sha256",
+ "domain": "radroots.rhi.reconciliation_source_replay.v1\\0",
+ "preimage": [
+ "request_id_32_bytes",
+ "prior_cursor_present_u8",
+ "prior_cursor_created_at_unix_seconds_u64_be_when_present",
+ "prior_cursor_verified_event_id_32_bytes_when_present",
+ "overlap_seconds_u64_be",
+ "inclusive_since_unix_seconds_u64_be"
+ ],
+ "exact_vector": {
+ "request_id_hex": "ef58f9e8a61f7964734cf0c2aabe0bdb2cbdcec529b16f40ae18ec224d0c896d",
+ "prior_cursor": null,
+ "overlap_seconds": 300,
+ "inclusive_since_unix_seconds": 0,
+ "replay_id_hex": "2490a2e9a6e85051e92f6c2fc2ff7e98a1367afd6c26f8625eb421cd2ab30c68"
+ }
+ },
+ "cursor": {
+ "tuple": ["event_authored_unix_seconds", "verified_event_id"],
+ "input": "sealed_step_190_committed_cursor_evidence",
+ "scope": ["source_id", "trade_id", "evidence_policy_digest", "selector_digest"],
+ "scope_mismatch": "reject_before_resume",
+ "ordering": "lexicographic_ascending",
+ "initial_since": "max_zero_floor_attempt_start_unix_ms_to_seconds_minus_lookback_seconds",
+ "resume_since": "max_zero_prior_cursor_authored_seconds_minus_configured_overlap_seconds",
+ "equal_timestamp_safe": true,
+ "candidate": "greatest_canonical_admitted_cursor",
+ "eligible_only_for": "complete_and_strictly_after_prior_cursor",
+ "scope_revalidation_before_commit": "step_190"
+ },
+ "inventory": {
+ "input_ingestion_bound": "request_maximum_events_plus_one",
+ "input_original_bytes_bound": "request_maximum_original_event_bytes_before_deduplication",
+ "trade_scope": "exact_request_trade_id",
+ "canonical_order": [
+ "event_authored_unix_seconds_ascending",
+ "verified_event_id_bytes_ascending",
+ "verified_signature_bytes_ascending",
+ "source_observed_unix_seconds_ascending",
+ "original_event_byte_length_ascending"
+ ],
+ "signed_event_identity": ["verified_event_id", "verified_signature"],
+ "exact_replay": "one_canonical_fact",
+ "observation_deduplication": "one_fact_per_signed_event_identity",
+ "first_provenance": "earliest_injected_observation_time_retained",
+ "mutation_identity_reuse": "exact_canonical_mutation_required",
+ "signed_event_identity_reuse": "exact_canonical_signed_event_required",
+ "accepted_count_and_bytes": "derived_from_canonical_distinct_inventory"
+ },
+ "timing": {
+ "observation": "injected_whole_unix_second_interval_intersects_source_result_start_and_finish",
+ "event_authored_time": "distinct_untrusted_previously_admitted_input"
+ },
+ "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": [
+ "unscoped_trade",
+ "cursor_only_global_progress",
+ "caller_forged_or_relabelled_cursor_scope",
+ "in_memory_candidate_used_as_committed_cursor",
+ "exclusive_since_cursor",
+ "arrival_order_authority",
+ "duplicate_authoritative_observation",
+ "later_provenance_overwrite",
+ "conflicting_mutation_identity_reuse",
+ "conflicting_signed_event_identity_reuse",
+ "unbounded_event_iterator",
+ "unbounded_original_event_bytes",
+ "checkpoint_advance_before_complete_commit"
+ ],
+ "deferred": [
+ "source_execution",
+ "durable_source_result_commit",
+ "durable_scope_revalidation",
+ "checkpoint_advance",
+ "manifest",
+ "reducer",
+ "attestation",
+ "publication"
+ ]
+}
diff --git a/src/lib.rs b/src/lib.rs
@@ -10,6 +10,7 @@ mod identity_credential;
mod identity_envelope;
mod reconciliation_attempt;
mod reconciliation_job;
+mod reconciliation_replay;
mod runtime_adapters;
mod runtime_context;
mod runtime_foundation;
@@ -82,6 +83,12 @@ pub use reconciliation_job::{
RhiReconciliationLease, RhiReconciliationLeaseOwner, RhiReconciliationRetryDelayMilliseconds,
RhiReconciliationScheduleOutcome, RhiReconciliationUnixMilliseconds,
};
+pub use reconciliation_replay::{
+ RHI_RECONCILIATION_REPLAY_CONTRACT_VERSION, RhiReconciliationReplayError,
+ RhiReconciliationReplayErrorKind, RhiReconciliationSourceCursorEvidence,
+ RhiReconciliationSourceReplay, RhiReconciliationSourceReplayId,
+ RhiReconciliationSourceReplayPlan,
+};
pub use runtime_adapters::{
CanonicalRhiCredentialAccess, CanonicalRhiIdentityAccess, RHI_RUNTIME_ADAPTER_CONTRACT_VERSION,
RHI_RUNTIME_JITTER_MAX_ENTROPY_DRAWS, RHI_RUNTIME_JITTER_MAX_MILLISECONDS, RhiCredentialAccess,
diff --git a/src/reconciliation_replay.rs b/src/reconciliation_replay.rs
@@ -0,0 +1,857 @@
+//! Pure overlap-safe reconciliation replay and provenance canonicalization.
+
+use core::{cmp::Ordering, fmt};
+use std::{collections::BTreeMap, error::Error};
+
+use sha2::{Digest, Sha256};
+
+use crate::{
+ RhiAdmittedTradeMutationEvent, RhiConfigDocumentV1, RhiReconciliationAttemptPlan,
+ RhiReconciliationSourceRequest, RhiReconciliationSourceRequestId,
+ RhiReconciliationSourceResult, RhiReconciliationUnixMilliseconds,
+ RhiTradeMutationObservedAtUnixSeconds, RhiTradeSourceCompletion, RhiTradeSourceCursor,
+ state_metadata, state_trade::PersistenceRecord,
+};
+
+/// Exact version of the reconciliation replay contract.
+pub const RHI_RECONCILIATION_REPLAY_CONTRACT_VERSION: u32 = 1;
+
+const REPLAY_ID_DOMAIN: &[u8] = b"radroots.rhi.reconciliation_source_replay.v1\0";
+const SOURCE_KIND: &str = "nostr_relay";
+const SOURCE_SELECTOR: &str = "trade_mutation_lineage_v1";
+const MAX_OVERLAP_SECONDS: u64 = 86_400;
+
+/// Stable source-free replay validation class.
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+pub enum RhiReconciliationReplayErrorKind {
+ InvalidInput,
+ InvalidConfiguration,
+ PolicyMismatch,
+ ResourceLimit,
+ MutationConflict,
+ SignedEventConflict,
+}
+
+impl RhiReconciliationReplayErrorKind {
+ /// Returns the stable machine-readable failure code.
+ #[must_use]
+ pub const fn code(self) -> &'static str {
+ match self {
+ Self::InvalidInput => "reconciliation_replay_input_invalid",
+ Self::InvalidConfiguration => "reconciliation_replay_configuration_invalid",
+ Self::PolicyMismatch => "reconciliation_replay_policy_mismatch",
+ Self::ResourceLimit => "reconciliation_replay_resource_limit",
+ Self::MutationConflict => "reconciliation_replay_mutation_conflict",
+ Self::SignedEventConflict => "reconciliation_replay_signed_event_conflict",
+ }
+ }
+}
+
+/// Redacted source-free replay validation failure.
+#[derive(Clone, Copy, PartialEq, Eq)]
+pub struct RhiReconciliationReplayError {
+ kind: RhiReconciliationReplayErrorKind,
+}
+
+impl RhiReconciliationReplayError {
+ /// Returns the stable failure class.
+ #[must_use]
+ pub const fn kind(self) -> RhiReconciliationReplayErrorKind {
+ 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 RhiReconciliationReplayError {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter.write_str(match self.kind {
+ RhiReconciliationReplayErrorKind::InvalidInput => {
+ "RHI reconciliation replay input is invalid"
+ }
+ RhiReconciliationReplayErrorKind::InvalidConfiguration => {
+ "RHI reconciliation replay configuration is invalid"
+ }
+ RhiReconciliationReplayErrorKind::PolicyMismatch => {
+ "RHI reconciliation replay policy does not match the attempt"
+ }
+ RhiReconciliationReplayErrorKind::ResourceLimit => {
+ "RHI reconciliation replay exceeds its resource limit"
+ }
+ RhiReconciliationReplayErrorKind::MutationConflict => {
+ "RHI reconciliation replay conflicts with canonical mutation evidence"
+ }
+ RhiReconciliationReplayErrorKind::SignedEventConflict => {
+ "RHI reconciliation replay conflicts with signed event evidence"
+ }
+ })
+ }
+}
+
+impl fmt::Debug for RhiReconciliationReplayError {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter
+ .debug_struct("RhiReconciliationReplayError")
+ .field("kind", &self.kind)
+ .finish()
+ }
+}
+
+impl Error for RhiReconciliationReplayError {}
+
+/// Domain-separated identity of one request's exact cursor/overlap binding.
+#[derive(Clone, Copy, PartialEq, Eq, Hash)]
+pub struct RhiReconciliationSourceReplayId([u8; 32]);
+
+impl RhiReconciliationSourceReplayId {
+ /// Returns the exact identity bytes.
+ #[must_use]
+ pub const fn as_bytes(&self) -> &[u8; 32] {
+ &self.0
+ }
+}
+
+impl fmt::Debug for RhiReconciliationSourceReplayId {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter.write_str("RhiReconciliationSourceReplayId([redacted])")
+ }
+}
+
+/// Sealed source/trade/policy/selector-scoped evidence for a committed cursor.
+///
+/// Step 189 defines and consumes this non-forgeable capability but provides no
+/// minting path. Step 190 alone may construct it after the replay and cursor
+/// are committed atomically under their exact durable scope.
+#[derive(Clone, PartialEq, Eq)]
+pub struct RhiReconciliationSourceCursorEvidence {
+ source_id: Box<str>,
+ trade_id: radroots_event::id::TradeId,
+ policy_digest: [u8; 32],
+ selector_digest: [u8; 32],
+ cursor: RhiTradeSourceCursor,
+}
+
+impl RhiReconciliationSourceCursorEvidence {
+ /// Returns the retained exact cursor tuple.
+ #[must_use]
+ pub const fn cursor(&self) -> RhiTradeSourceCursor {
+ self.cursor
+ }
+}
+
+impl fmt::Debug for RhiReconciliationSourceCursorEvidence {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter
+ .debug_struct("RhiReconciliationSourceCursorEvidence")
+ .field("scope", &"[redacted]")
+ .field("cursor", &self.cursor)
+ .finish()
+ }
+}
+
+/// Pure source-request binding to an optional prior cursor and overlap window.
+#[derive(Clone, PartialEq, Eq)]
+pub struct RhiReconciliationSourceReplayPlan {
+ id: RhiReconciliationSourceReplayId,
+ request_id: RhiReconciliationSourceRequestId,
+ prior_cursor: Option<RhiReconciliationSourceCursorEvidence>,
+ overlap_seconds: u64,
+ since_unix_seconds: u64,
+}
+
+impl RhiReconciliationSourceReplayPlan {
+ /// Derives the exact overlap-safe cursor binding for one attempt request.
+ pub fn from_request(
+ attempt: &RhiReconciliationAttemptPlan,
+ request: &RhiReconciliationSourceRequest,
+ configuration: &RhiConfigDocumentV1,
+ prior_cursor: Option<RhiReconciliationSourceCursorEvidence>,
+ ) -> Result<Self, RhiReconciliationReplayError> {
+ if !attempt
+ .requests()
+ .iter()
+ .any(|candidate| candidate.id() == request.id())
+ {
+ return Err(failure(RhiReconciliationReplayErrorKind::InvalidInput));
+ }
+ let normalized = configuration.normalized();
+ let policy = state_metadata::evidence_policy_digest(normalized)
+ .map_err(|_| failure(RhiReconciliationReplayErrorKind::InvalidConfiguration))?;
+ if policy != attempt.evidence_policy_digest() {
+ return Err(failure(RhiReconciliationReplayErrorKind::PolicyMismatch));
+ }
+ if prior_cursor.as_ref().is_some_and(|evidence| {
+ !cursor_scope_matches(
+ evidence,
+ request.source_id(),
+ request.trade_id(),
+ policy.as_bytes(),
+ request.selector_digest().as_bytes(),
+ )
+ }) {
+ return Err(failure(RhiReconciliationReplayErrorKind::PolicyMismatch));
+ }
+ let source = normalized
+ .pointer("/evidence/sources")
+ .and_then(serde_json::Value::as_array)
+ .and_then(|sources| {
+ sources.iter().find(|source| {
+ source
+ .pointer("/source_id")
+ .and_then(serde_json::Value::as_str)
+ == Some(request.source_id())
+ })
+ })
+ .filter(|source| {
+ 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)
+ })
+ .ok_or_else(|| failure(RhiReconciliationReplayErrorKind::InvalidConfiguration))?;
+ let overlap_seconds = source
+ .pointer("/overlap_seconds")
+ .and_then(serde_json::Value::as_u64)
+ .filter(|value| (1..=MAX_OVERLAP_SECONDS).contains(value))
+ .ok_or_else(|| failure(RhiReconciliationReplayErrorKind::InvalidConfiguration))?;
+ let lookback_seconds = source
+ .pointer("/lookback_seconds")
+ .and_then(serde_json::Value::as_u64)
+ .filter(|value| *value == request.lookback_seconds() && overlap_seconds <= *value)
+ .ok_or_else(|| failure(RhiReconciliationReplayErrorKind::InvalidConfiguration))?;
+ let raw_prior_cursor = prior_cursor.as_ref().map(|evidence| evidence.cursor);
+ let since_unix_seconds = raw_prior_cursor.map_or_else(
+ || {
+ request
+ .attempt_started_at()
+ .get()
+ .checked_div(1_000)
+ .unwrap_or(0)
+ .saturating_sub(lookback_seconds)
+ },
+ |cursor| resume_since(cursor, overlap_seconds),
+ );
+ Ok(Self {
+ id: replay_id(
+ request.id(),
+ raw_prior_cursor,
+ overlap_seconds,
+ since_unix_seconds,
+ ),
+ request_id: request.id(),
+ prior_cursor,
+ overlap_seconds,
+ since_unix_seconds,
+ })
+ }
+
+ /// Returns the exact cursor-binding identity.
+ #[must_use]
+ pub const fn id(&self) -> RhiReconciliationSourceReplayId {
+ self.id
+ }
+
+ /// Returns the exact source request identity.
+ #[must_use]
+ pub const fn request_id(&self) -> RhiReconciliationSourceRequestId {
+ self.request_id
+ }
+
+ /// Returns the retained prior cursor, when one exists.
+ #[must_use]
+ pub fn prior_cursor(&self) -> Option<RhiTradeSourceCursor> {
+ self.prior_cursor.as_ref().map(|evidence| evidence.cursor)
+ }
+
+ /// Returns the configured overlap in whole seconds.
+ #[must_use]
+ pub const fn overlap_seconds(&self) -> u64 {
+ self.overlap_seconds
+ }
+
+ /// Returns the inclusive overlap-safe query start in Unix seconds.
+ #[must_use]
+ pub const fn since_unix_seconds(&self) -> u64 {
+ self.since_unix_seconds
+ }
+
+ /// Canonicalizes a bounded admitted-event result for a later atomic commit.
+ pub fn finish<I>(
+ self,
+ request: &RhiReconciliationSourceRequest,
+ outcome: RhiTradeSourceCompletion,
+ started_at: RhiReconciliationUnixMilliseconds,
+ finished_at: RhiReconciliationUnixMilliseconds,
+ events: I,
+ ) -> Result<RhiReconciliationSourceReplay, RhiReconciliationReplayError>
+ where
+ I: IntoIterator<Item = RhiAdmittedTradeMutationEvent>,
+ {
+ if request.id() != self.request_id {
+ return Err(failure(RhiReconciliationReplayErrorKind::InvalidInput));
+ }
+ RhiReconciliationSourceResult::new(request, outcome, started_at, finished_at, 0, 0)
+ .map_err(|_| failure(RhiReconciliationReplayErrorKind::InvalidInput))?;
+ let maximum_events = usize::try_from(request.maximum_events())
+ .map_err(|_| failure(RhiReconciliationReplayErrorKind::InvalidConfiguration))?;
+ let candidates = ingest_bounded_facts(
+ events.into_iter().map(|event| {
+ if event.mutation().trade_id != request.trade_id() {
+ return Err(failure(RhiReconciliationReplayErrorKind::InvalidInput));
+ }
+ let event_bytes = u64::try_from(event.original_bytes().len())
+ .map_err(|_| failure(RhiReconciliationReplayErrorKind::ResourceLimit))?;
+ let observed_at = event.observed_at_unix_seconds();
+ let observed_at_milliseconds = observed_at
+ .get()
+ .checked_mul(1_000)
+ .ok_or_else(|| failure(RhiReconciliationReplayErrorKind::InvalidInput))?;
+ let observed_interval_end = observed_at_milliseconds
+ .checked_add(999)
+ .ok_or_else(|| failure(RhiReconciliationReplayErrorKind::InvalidInput))?;
+ if observed_interval_end < started_at.get()
+ || observed_at_milliseconds > finished_at.get()
+ {
+ return Err(failure(RhiReconciliationReplayErrorKind::InvalidInput));
+ }
+ Ok(ReplayFact {
+ original_bytes: event_bytes,
+ observed_at,
+ record: PersistenceRecord::from_admitted(event)
+ .map_err(|_| failure(RhiReconciliationReplayErrorKind::InvalidInput))?,
+ })
+ }),
+ maximum_events,
+ request.maximum_bytes(),
+ )?;
+ let CanonicalReplay {
+ facts,
+ duplicate_observations,
+ accepted_original_bytes,
+ cursor_candidate,
+ first_observed_at,
+ } = canonicalize(candidates)?;
+ let accepted_event_count = u32::try_from(facts.len())
+ .map_err(|_| failure(RhiReconciliationReplayErrorKind::ResourceLimit))?;
+ let result = RhiReconciliationSourceResult::new(
+ request,
+ outcome,
+ started_at,
+ finished_at,
+ accepted_event_count,
+ accepted_original_bytes,
+ )
+ .map_err(|_| failure(RhiReconciliationReplayErrorKind::InvalidInput))?;
+ Ok(RhiReconciliationSourceReplay {
+ plan: self,
+ result,
+ duplicate_observations,
+ accepted_original_bytes,
+ cursor_candidate,
+ first_observed_at,
+ facts: facts.into_boxed_slice(),
+ })
+ }
+}
+
+impl fmt::Debug for RhiReconciliationSourceReplayPlan {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter
+ .debug_struct("RhiReconciliationSourceReplayPlan")
+ .field("identity", &"[redacted]")
+ .field("has_prior_cursor", &self.prior_cursor.is_some())
+ .field("overlap_seconds", &self.overlap_seconds)
+ .field("since_unix_seconds", &self.since_unix_seconds)
+ .finish()
+ }
+}
+
+/// Canonical bounded replay result prepared for the Step 190 commit boundary.
+pub struct RhiReconciliationSourceReplay {
+ plan: RhiReconciliationSourceReplayPlan,
+ result: RhiReconciliationSourceResult,
+ duplicate_observations: u32,
+ accepted_original_bytes: u64,
+ cursor_candidate: Option<RhiTradeSourceCursor>,
+ first_observed_at: Option<RhiTradeMutationObservedAtUnixSeconds>,
+ facts: Box<[ReplayFact]>,
+}
+
+impl RhiReconciliationSourceReplay {
+ /// Returns the exact cursor-binding identity.
+ #[must_use]
+ pub const fn id(&self) -> RhiReconciliationSourceReplayId {
+ self.plan.id()
+ }
+
+ /// Returns the validated source result derived from the canonical inventory.
+ #[must_use]
+ pub const fn result(&self) -> RhiReconciliationSourceResult {
+ self.result
+ }
+
+ /// Returns the number of distinct accepted signed-event identities.
+ #[must_use]
+ pub fn accepted_event_count(&self) -> usize {
+ self.facts.len()
+ }
+
+ /// Returns the retained first-provenance original-byte total.
+ #[must_use]
+ pub const fn accepted_original_event_bytes(&self) -> u64 {
+ self.accepted_original_bytes
+ }
+
+ /// Returns the exact replay observations removed by canonical deduplication.
+ #[must_use]
+ pub const fn duplicate_observation_count(&self) -> u32 {
+ self.duplicate_observations
+ }
+
+ /// Returns the earliest retained source observation, when evidence exists.
+ #[must_use]
+ pub const fn first_observed_at(&self) -> Option<RhiTradeMutationObservedAtUnixSeconds> {
+ self.first_observed_at
+ }
+
+ /// Returns the greatest admitted cursor candidate regardless of completion.
+ #[must_use]
+ pub const fn cursor_candidate(&self) -> Option<RhiTradeSourceCursor> {
+ self.cursor_candidate
+ }
+
+ /// Returns a cursor only when exact source completion makes it eligible.
+ #[must_use]
+ pub fn eligible_cursor(&self) -> Option<RhiTradeSourceCursor> {
+ eligible_cursor(
+ self.result.outcome(),
+ self.cursor_candidate,
+ self.plan
+ .prior_cursor
+ .as_ref()
+ .map(|evidence| evidence.cursor),
+ )
+ }
+}
+
+fn cursor_scope_matches(
+ evidence: &RhiReconciliationSourceCursorEvidence,
+ source_id: &str,
+ trade_id: radroots_event::id::TradeId,
+ policy_digest: &[u8; 32],
+ selector_digest: &[u8; 32],
+) -> bool {
+ evidence.source_id.as_ref() == source_id
+ && evidence.trade_id == trade_id
+ && &evidence.policy_digest == policy_digest
+ && &evidence.selector_digest == selector_digest
+}
+
+impl fmt::Debug for RhiReconciliationSourceReplay {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter
+ .debug_struct("RhiReconciliationSourceReplay")
+ .field("identity", &"[redacted]")
+ .field("outcome", &self.result.outcome())
+ .field("accepted_event_count", &self.facts.len())
+ .field(
+ "accepted_original_event_bytes",
+ &self.accepted_original_bytes,
+ )
+ .field("duplicate_observations", &self.duplicate_observations)
+ .field("has_cursor_candidate", &self.cursor_candidate.is_some())
+ .finish()
+ }
+}
+
+struct ReplayFact {
+ record: PersistenceRecord,
+ original_bytes: u64,
+ observed_at: RhiTradeMutationObservedAtUnixSeconds,
+}
+
+struct CanonicalReplay {
+ facts: Vec<ReplayFact>,
+ duplicate_observations: u32,
+ accepted_original_bytes: u64,
+ cursor_candidate: Option<RhiTradeSourceCursor>,
+ first_observed_at: Option<RhiTradeMutationObservedAtUnixSeconds>,
+}
+
+fn ingest_bounded_facts<I>(
+ facts: I,
+ maximum_events: usize,
+ maximum_bytes: u64,
+) -> Result<Vec<ReplayFact>, RhiReconciliationReplayError>
+where
+ I: IntoIterator<Item = Result<ReplayFact, RhiReconciliationReplayError>>,
+{
+ let mut accepted_bytes = 0_u64;
+ let mut bounded = Vec::with_capacity(maximum_events);
+ for fact in facts.into_iter().take(maximum_events.saturating_add(1)) {
+ if bounded.len() == maximum_events {
+ return Err(failure(RhiReconciliationReplayErrorKind::ResourceLimit));
+ }
+ let fact = fact?;
+ accepted_bytes = accepted_bytes
+ .checked_add(fact.original_bytes)
+ .filter(|value| *value <= maximum_bytes)
+ .ok_or_else(|| failure(RhiReconciliationReplayErrorKind::ResourceLimit))?;
+ bounded.push(fact);
+ }
+ Ok(bounded)
+}
+
+fn canonicalize(
+ mut candidates: Vec<ReplayFact>,
+) -> Result<CanonicalReplay, RhiReconciliationReplayError> {
+ candidates.sort_by(compare_fact);
+ let mut event_indexes = BTreeMap::<([u8; 32], [u8; 64]), usize>::new();
+ let mut mutation_indexes = BTreeMap::<[u8; 32], usize>::new();
+ let mut facts = Vec::<ReplayFact>::with_capacity(candidates.len());
+ let mut duplicate_observations = 0_u32;
+ let mut accepted_original_bytes = 0_u64;
+ let mut cursor_candidate = None;
+ let mut first_observed_at: Option<RhiTradeMutationObservedAtUnixSeconds> = None;
+ for candidate in candidates {
+ let event_key = (candidate.record.event_id, candidate.record.event_signature);
+ if let Some(index) = event_indexes.get(&event_key).copied() {
+ if !same_event(&facts[index].record, &candidate.record) {
+ return Err(failure(
+ RhiReconciliationReplayErrorKind::SignedEventConflict,
+ ));
+ }
+ duplicate_observations = duplicate_observations
+ .checked_add(1)
+ .ok_or_else(|| failure(RhiReconciliationReplayErrorKind::ResourceLimit))?;
+ continue;
+ }
+ if let Some(index) = mutation_indexes.get(&candidate.record.mutation_id).copied()
+ && !same_mutation(&facts[index].record, &candidate.record)
+ {
+ return Err(failure(RhiReconciliationReplayErrorKind::MutationConflict));
+ }
+ accepted_original_bytes = accepted_original_bytes
+ .checked_add(candidate.original_bytes)
+ .ok_or_else(|| failure(RhiReconciliationReplayErrorKind::ResourceLimit))?;
+ let cursor = RhiTradeSourceCursor::from_verified_parts(
+ candidate.record.authored_at_unix_s,
+ candidate.record.event_id,
+ );
+ cursor_candidate = Some(cursor_candidate.map_or(cursor, |current| {
+ if compare_cursor(current, cursor).is_lt() {
+ cursor
+ } else {
+ current
+ }
+ }));
+ first_observed_at = Some(first_observed_at.map_or(candidate.observed_at, |current| {
+ if candidate.observed_at.get() < current.get() {
+ candidate.observed_at
+ } else {
+ current
+ }
+ }));
+ let index = facts.len();
+ event_indexes.insert(event_key, index);
+ mutation_indexes
+ .entry(candidate.record.mutation_id)
+ .or_insert(index);
+ facts.push(candidate);
+ }
+ Ok(CanonicalReplay {
+ facts,
+ duplicate_observations,
+ accepted_original_bytes,
+ cursor_candidate,
+ first_observed_at,
+ })
+}
+
+fn compare_fact(left: &ReplayFact, right: &ReplayFact) -> Ordering {
+ left.record
+ .authored_at_unix_s
+ .cmp(&right.record.authored_at_unix_s)
+ .then_with(|| left.record.event_id.cmp(&right.record.event_id))
+ .then_with(|| {
+ left.record
+ .event_signature
+ .cmp(&right.record.event_signature)
+ })
+ .then_with(|| left.observed_at.get().cmp(&right.observed_at.get()))
+ .then_with(|| left.original_bytes.cmp(&right.original_bytes))
+}
+
+fn same_mutation(left: &PersistenceRecord, right: &PersistenceRecord) -> bool {
+ left.mutation_id == right.mutation_id
+ && left.trade_id == right.trade_id
+ && left.contract_id == right.contract_id
+ && left.schema_version == right.schema_version
+ && left.event_kind == right.event_kind
+ && left.author_pubkey == right.author_pubkey
+ && left.canonical_content == right.canonical_content
+}
+
+fn same_event(left: &PersistenceRecord, right: &PersistenceRecord) -> bool {
+ same_mutation(left, right)
+ && left.event_id == right.event_id
+ && left.event_signature == right.event_signature
+ && left.event_kind == right.event_kind
+ && left.authored_at_unix_s == right.authored_at_unix_s
+ && left.canonical_event_json == right.canonical_event_json
+}
+
+fn compare_cursor(left: RhiTradeSourceCursor, right: RhiTradeSourceCursor) -> Ordering {
+ (left.created_at_unix_seconds(), left.event_id())
+ .cmp(&(right.created_at_unix_seconds(), right.event_id()))
+}
+
+fn resume_since(cursor: RhiTradeSourceCursor, overlap_seconds: u64) -> u64 {
+ cursor
+ .created_at_unix_seconds()
+ .saturating_sub(overlap_seconds)
+}
+
+fn eligible_cursor(
+ outcome: RhiTradeSourceCompletion,
+ candidate: Option<RhiTradeSourceCursor>,
+ prior: Option<RhiTradeSourceCursor>,
+) -> Option<RhiTradeSourceCursor> {
+ if !matches!(outcome, RhiTradeSourceCompletion::Complete) {
+ return None;
+ }
+ candidate
+ .filter(|candidate| prior.is_none_or(|prior| compare_cursor(prior, *candidate).is_lt()))
+}
+
+fn replay_id(
+ request_id: RhiReconciliationSourceRequestId,
+ prior_cursor: Option<RhiTradeSourceCursor>,
+ overlap_seconds: u64,
+ since_unix_seconds: u64,
+) -> RhiReconciliationSourceReplayId {
+ let mut digest = Sha256::new();
+ digest.update(REPLAY_ID_DOMAIN);
+ digest.update(request_id.as_bytes());
+ match prior_cursor {
+ Some(cursor) => {
+ digest.update([1]);
+ digest.update(cursor.created_at_unix_seconds().to_be_bytes());
+ digest.update(cursor.event_id());
+ }
+ None => digest.update([0]),
+ }
+ digest.update(overlap_seconds.to_be_bytes());
+ digest.update(since_unix_seconds.to_be_bytes());
+ RhiReconciliationSourceReplayId(digest.finalize().into())
+}
+
+const fn failure(kind: RhiReconciliationReplayErrorKind) -> RhiReconciliationReplayError {
+ RhiReconciliationReplayError { kind }
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+
+ fn observed(value: u64) -> RhiTradeMutationObservedAtUnixSeconds {
+ RhiTradeMutationObservedAtUnixSeconds::new(value).expect("observation")
+ }
+
+ fn record(event: u8, signature: u8, mutation: u8, content: &[u8]) -> PersistenceRecord {
+ PersistenceRecord {
+ mutation_id: [mutation; 32],
+ trade_id: [0x11; 16],
+ contract_id: "radroots.trade.proposal.v1",
+ schema_version: 1,
+ event_id: [event; 32],
+ event_signature: [signature; 64],
+ author_pubkey: [0x22; 32],
+ event_kind: 3470,
+ authored_at_unix_s: 1_784_347_200,
+ canonical_content: content.into(),
+ canonical_event_json: [b"event:".as_slice(), content].concat().into_boxed_slice(),
+ }
+ }
+
+ fn fact(
+ event: u8,
+ signature: u8,
+ mutation: u8,
+ content: &[u8],
+ observation: u64,
+ bytes: u64,
+ ) -> ReplayFact {
+ ReplayFact {
+ record: record(event, signature, mutation, content),
+ original_bytes: bytes,
+ observed_at: observed(observation),
+ }
+ }
+
+ #[test]
+ fn canonicalization_deduplicates_replay_and_retains_first_provenance() {
+ let canonical = canonicalize(vec![
+ fact(1, 2, 3, b"same", 1_784_347_202, 12),
+ fact(1, 2, 3, b"same", 1_784_347_200, 10),
+ fact(1, 4, 3, b"same", 1_784_347_201, 11),
+ ])
+ .expect("canonical replay");
+ assert_eq!(canonical.facts.len(), 2);
+ assert_eq!(canonical.duplicate_observations, 1);
+ assert_eq!(canonical.accepted_original_bytes, 21);
+ assert_eq!(
+ canonical.first_observed_at.expect("first").get(),
+ 1_784_347_200
+ );
+ }
+
+ #[test]
+ fn conflicting_mutation_and_event_identity_reuse_fail_closed() {
+ let mutation = canonicalize(vec![
+ fact(1, 2, 3, b"first", 1_784_347_200, 10),
+ fact(4, 5, 3, b"second", 1_784_347_201, 10),
+ ])
+ .err()
+ .expect("mutation conflict");
+ assert_eq!(
+ mutation.kind(),
+ RhiReconciliationReplayErrorKind::MutationConflict
+ );
+
+ let event = canonicalize(vec![
+ fact(1, 2, 3, b"first", 1_784_347_200, 10),
+ fact(1, 2, 4, b"second", 1_784_347_201, 10),
+ ])
+ .err()
+ .expect("event conflict");
+ assert_eq!(
+ event.kind(),
+ RhiReconciliationReplayErrorKind::SignedEventConflict
+ );
+ }
+
+ #[test]
+ fn pre_dedup_bound_is_exact_and_infinite_iterators_terminate() {
+ let exact = ingest_bounded_facts([Ok(fact(1, 2, 3, b"one", 1_784_347_200, 10))], 1, 10)
+ .expect("exact maximum");
+ assert_eq!(exact.len(), 1);
+
+ let over_bytes =
+ ingest_bounded_facts([Ok(fact(1, 2, 3, b"one", 1_784_347_200, 11))], 1, 10)
+ .err()
+ .expect("byte maximum plus one");
+ assert_eq!(
+ over_bytes.kind(),
+ RhiReconciliationReplayErrorKind::ResourceLimit
+ );
+
+ let mut event = 0_u8;
+ let over_count = ingest_bounded_facts(
+ std::iter::repeat_with(|| {
+ event = event.wrapping_add(1);
+ Ok(fact(event, 2, event, b"one", 1_784_347_200, 1))
+ }),
+ 1,
+ 10,
+ )
+ .err()
+ .expect("count maximum plus one");
+ assert_eq!(
+ over_count.kind(),
+ RhiReconciliationReplayErrorKind::ResourceLimit
+ );
+ }
+
+ #[test]
+ fn errors_are_source_free_and_redacted() {
+ let error = failure(RhiReconciliationReplayErrorKind::SignedEventConflict);
+ assert!(Error::source(&error).is_none());
+ assert_eq!(error.code(), "reconciliation_replay_signed_event_conflict");
+ assert_eq!(
+ format!("{error:?}"),
+ "RhiReconciliationReplayError { kind: SignedEventConflict }"
+ );
+ }
+
+ #[test]
+ fn cursor_evidence_scope_requires_every_exact_dimension() {
+ let exact = RhiReconciliationSourceCursorEvidence {
+ source_id: "trade-primary".into(),
+ trade_id: radroots_event::id::TradeId::from_bytes([0x11; 16]),
+ policy_digest: [0x22; 32],
+ selector_digest: [0x33; 32],
+ cursor: RhiTradeSourceCursor::from_verified_parts(42, [0x44; 32]),
+ };
+ assert!(cursor_scope_matches(
+ &exact,
+ "trade-primary",
+ radroots_event::id::TradeId::from_bytes([0x11; 16]),
+ &[0x22; 32],
+ &[0x33; 32],
+ ));
+ assert!(!cursor_scope_matches(
+ &exact,
+ "trade-secondary",
+ radroots_event::id::TradeId::from_bytes([0x11; 16]),
+ &[0x22; 32],
+ &[0x33; 32],
+ ));
+ assert!(!cursor_scope_matches(
+ &exact,
+ "trade-primary",
+ radroots_event::id::TradeId::from_bytes([0x12; 16]),
+ &[0x22; 32],
+ &[0x33; 32],
+ ));
+ assert!(!cursor_scope_matches(
+ &exact,
+ "trade-primary",
+ radroots_event::id::TradeId::from_bytes([0x11; 16]),
+ &[0x23; 32],
+ &[0x33; 32],
+ ));
+ assert!(!cursor_scope_matches(
+ &exact,
+ "trade-primary",
+ radroots_event::id::TradeId::from_bytes([0x11; 16]),
+ &[0x22; 32],
+ &[0x34; 32],
+ ));
+ }
+
+ #[test]
+ fn resume_overlap_and_cursor_eligibility_are_exact() {
+ let prior = RhiTradeSourceCursor::from_verified_parts(500, [0x11; 32]);
+ let equal = RhiTradeSourceCursor::from_verified_parts(500, [0x11; 32]);
+ let older = RhiTradeSourceCursor::from_verified_parts(499, [0xff; 32]);
+ let newer = RhiTradeSourceCursor::from_verified_parts(500, [0x12; 32]);
+ assert_eq!(resume_since(prior, 300), 200);
+ assert_eq!(resume_since(prior, 600), 0);
+ assert_eq!(
+ eligible_cursor(RhiTradeSourceCompletion::Complete, Some(newer), Some(prior)),
+ Some(newer)
+ );
+ assert_eq!(
+ eligible_cursor(RhiTradeSourceCompletion::Complete, Some(equal), Some(prior)),
+ None
+ );
+ assert_eq!(
+ eligible_cursor(RhiTradeSourceCompletion::Complete, Some(older), Some(prior)),
+ None
+ );
+ assert_eq!(
+ eligible_cursor(
+ RhiTradeSourceCompletion::IncompleteUnavailable,
+ Some(newer),
+ Some(prior),
+ ),
+ None
+ );
+ }
+}
diff --git a/src/source_ingest.rs b/src/source_ingest.rs
@@ -121,6 +121,16 @@ pub struct RhiTradeSourceCursor {
}
impl RhiTradeSourceCursor {
+ pub(crate) const fn from_verified_parts(
+ created_at_unix_seconds: u64,
+ event_id: [u8; 32],
+ ) -> Self {
+ Self {
+ created_at_unix_seconds,
+ event_id,
+ }
+ }
+
/// Returns the inclusive event-authored UTC second.
#[must_use]
pub const fn created_at_unix_seconds(self) -> u64 {
diff --git a/tests/package_boundary.rs b/tests/package_boundary.rs
@@ -13,8 +13,11 @@ const RUNTIME_ADAPTER_CONTRACT: &str =
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_REPLAY: &str = include_str!("../src/reconciliation_replay.rs");
const RECONCILIATION_ATTEMPT_CONTRACT: &str =
include_str!("../contracts/services_hardening/reconciliation_attempts.v1.json");
+const RECONCILIATION_REPLAY_CONTRACT: &str =
+ include_str!("../contracts/services_hardening/reconciliation_replay.v1.json");
const RUNTIME_FOUNDATION_CONTRACT: &str =
include_str!("../contracts/services_hardening/runtime_foundation.v1.json");
const TRADE_INGEST_CONTRACT: &str =
@@ -33,6 +36,7 @@ const SOURCES: &[&str] = &[
include_str!("../src/identity_envelope.rs"),
include_str!("../src/reconciliation_attempt.rs"),
include_str!("../src/reconciliation_job.rs"),
+ include_str!("../src/reconciliation_replay.rs"),
include_str!("../src/runtime_context.rs"),
include_str!("../src/runtime_adapters.rs"),
include_str!("../src/runtime_foundation.rs"),
@@ -101,6 +105,7 @@ fn state_catalog_module_is_private_and_root_api_is_curated() {
"identity_envelope",
"reconciliation_attempt",
"reconciliation_job",
+ "reconciliation_replay",
"runtime_context",
"runtime_adapters",
"runtime_foundation",
@@ -140,6 +145,8 @@ fn state_catalog_module_is_private_and_root_api_is_curated() {
"RhiReconciliationSourceRequest",
"RhiReconciliationSourceResult",
"RhiReconciliationAttemptResults",
+ "RhiReconciliationSourceReplayPlan",
+ "RhiReconciliationSourceReplay",
"RhiReconciliationJobPolicy",
"RhiReconciliationLease",
"RhiTimeEntropyAdapters",
@@ -227,7 +234,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, 18);
+ assert_eq!(public_error_count, 19);
}
#[test]
@@ -280,6 +287,63 @@ fn reconciliation_attempts_are_exact_bounded_and_effect_free() {
}
#[test]
+fn reconciliation_replay_is_overlap_safe_bounded_and_effect_free() {
+ let contract: serde_json::Value = serde_json::from_str(RECONCILIATION_REPLAY_CONTRACT)
+ .expect("reconciliation-replay contract");
+ assert_eq!(contract["schema"], "radroots.rhi.reconciliation-replay");
+ assert_eq!(contract["contract_version"], 1);
+ assert_eq!(contract["state_schema_version"], 5);
+ assert_eq!(contract["cursor"]["equal_timestamp_safe"], true);
+ assert_eq!(
+ contract["cursor"]["input"],
+ "sealed_step_190_committed_cursor_evidence"
+ );
+ assert_eq!(
+ contract["cursor"]["eligible_only_for"],
+ "complete_and_strictly_after_prior_cursor"
+ );
+ assert_eq!(
+ contract["inventory"]["input_ingestion_bound"],
+ "request_maximum_events_plus_one"
+ );
+ assert_eq!(
+ contract["inventory"]["first_provenance"],
+ "earliest_injected_observation_time_retained"
+ );
+ assert_eq!(contract["effects"]["sqlite"], false);
+ assert_eq!(contract["effects"]["source_or_relay"], false);
+ for required in [
+ ".take(maximum_events.saturating_add(1))",
+ "saturating_sub(overlap_seconds)",
+ "RhiReconciliationReplayErrorKind::MutationConflict",
+ "RhiReconciliationReplayErrorKind::SignedEventConflict",
+ "RhiReconciliationSourceCursorEvidence",
+ "cursor_scope_matches(",
+ ] {
+ assert!(
+ RECONCILIATION_REPLAY.contains(required),
+ "replay boundary is missing {required}"
+ );
+ }
+ for forbidden in [
+ "sqlx::",
+ "std::fs",
+ "std::net",
+ "tokio::",
+ "SystemTime",
+ "thread_rng",
+ "OsRng",
+ ] {
+ assert!(
+ !RECONCILIATION_REPLAY.contains(forbidden),
+ "replay boundary gained forbidden authority {forbidden}"
+ );
+ }
+ assert!(!ROOT.contains("pub mod reconciliation_replay"));
+ assert!(!PUBLIC_API.contains("rhi::reconciliation_replay::"));
+}
+
+#[test]
fn trade_source_ingest_is_exact_bounded_generation_fenced_and_sealed() {
let contract: serde_json::Value =
serde_json::from_str(TRADE_SOURCE_INGEST_CONTRACT).expect("trade-source ingest contract");
@@ -584,6 +648,8 @@ fn readme_freezes_the_root_only_boundary_and_exact_baseline() {
"[`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)",
+ "## Overlap-safe reconciliation replay",
+ "[`reconciliation_replay.v1.json`](contracts/services_hardening/reconciliation_replay.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",
diff --git a/tests/services_hardening_reconciliation_attempt_contract.rs b/tests/services_hardening_reconciliation_attempt_contract.rs
@@ -38,7 +38,7 @@ fn machine_contract_freezes_the_complete_step_188_boundary() {
assert_eq!(contract["selector_identity"]["exact_tag"], "#d");
assert_eq!(
contract["selector_identity"]["cursor_binding"],
- "deferred_to_step_189"
+ "contracts/services_hardening/reconciliation_replay.v1.json"
);
assert_eq!(
contract["plan"]["source_count_maximum"],
@@ -78,7 +78,6 @@ fn machine_contract_freezes_the_complete_step_188_boundary() {
assert_eq!(
contract["deferred"],
json!([
- "cursor_and_overlap_derivation",
"source_execution",
"durable_source_result_commit",
"checkpoint_advance",
diff --git a/tests/services_hardening_reconciliation_jobs.rs b/tests/services_hardening_reconciliation_jobs.rs
@@ -10,15 +10,19 @@ use rhi::{
RhiReconciliationAttemptErrorKind, RhiReconciliationAttemptPlan,
RhiReconciliationAttemptResults, RhiReconciliationJobErrorKind, RhiReconciliationJobPolicy,
RhiReconciliationJobState, RhiReconciliationLeaseOwner,
- RhiReconciliationRetryDelayMilliseconds, RhiReconciliationSourceResult,
- RhiReconciliationUnixMilliseconds, RhiRuntimeContext, RhiStateMetadata,
- RhiTradeSourceCompletion, TradeId, initialize_rhi_state, open_rhi_state_inspection,
+ RhiReconciliationRetryDelayMilliseconds, RhiReconciliationSourceReplayPlan,
+ RhiReconciliationSourceResult, RhiReconciliationUnixMilliseconds, RhiRuntimeContext,
+ RhiStateMetadata, RhiTradeMutationAdmissionLimits, RhiTradeMutationAuthoredTimePolicy,
+ RhiTradeMutationObservedAtUnixSeconds, RhiTradeSourceCompletion, TradeId,
+ admit_rhi_trade_mutation_event, initialize_rhi_state, open_rhi_state_inspection,
open_rhi_state_read_write, parse_rhi_cli_v1_from, parse_rhi_config_v1,
resolve_rhi_runtime_context,
};
use sqlx::{Connection, SqliteConnection, sqlite::SqliteConnectOptions};
const EXAMPLE: &str = include_str!("../contracts/services_hardening/config.v1.example.toml");
+const TRADE_VECTOR: &str =
+ include_str!("../contracts/conformance/vectors/trade_ingest_proposal.v1.json");
fn runtime(root: &Path, instance: &str) -> RhiRuntimeContext {
let invocation = parse_rhi_cli_v1_from([
@@ -522,6 +526,15 @@ async fn attempt_plan_binds_the_exact_claim_policy_sources_and_frozen_identities
lower_hex(request.id().as_bytes()),
"ef58f9e8a61f7964734cf0c2aabe0bdb2cbdcec529b16f40ae18ec224d0c896d"
);
+ let replay =
+ RhiReconciliationSourceReplayPlan::from_request(&plan, request, &configuration, None)
+ .expect("replay plan");
+ assert_eq!(replay.overlap_seconds(), 300);
+ assert_eq!(replay.since_unix_seconds(), 0);
+ assert_eq!(
+ lower_hex(replay.id().as_bytes()),
+ "2490a2e9a6e85051e92f6c2fc2ff7e98a1367afd6c26f8625eb421cd2ab30c68"
+ );
host.close().await.expect("close");
}
@@ -900,6 +913,68 @@ async fn attempt_diagnostics_are_redacted_and_source_free() {
drop((configuration, metadata, runtime, root));
}
+#[tokio::test]
+async fn replay_plan_binds_overlap_deduplicates_and_retains_first_provenance() {
+ let started_ms = 1_784_347_200_000;
+ let (root, runtime, metadata, configuration, host, plan) =
+ replay_fixture("replay-plan", EXAMPLE, started_ms).await;
+ let request = &plan.requests()[0];
+ let cursor_plan =
+ RhiReconciliationSourceReplayPlan::from_request(&plan, request, &configuration, None)
+ .expect("initial cursor plan");
+ assert_eq!(cursor_plan.overlap_seconds(), 300);
+ assert_eq!(cursor_plan.since_unix_seconds(), 1_784_260_800);
+ assert!(cursor_plan.prior_cursor().is_none());
+
+ let wire = replay_wire();
+ let replay = cursor_plan
+ .finish(
+ request,
+ RhiTradeSourceCompletion::Complete,
+ now(started_ms),
+ now(started_ms + 2_000),
+ [
+ admitted_replay(&configuration, &wire, 1_784_347_201),
+ admitted_replay(&configuration, &wire, 1_784_347_200),
+ ],
+ )
+ .expect("canonical replay");
+ assert_eq!(replay.accepted_event_count(), 1);
+ assert_eq!(
+ replay.accepted_original_event_bytes(),
+ u64::try_from(wire.len()).expect("wire bytes")
+ );
+ assert_eq!(replay.duplicate_observation_count(), 1);
+ assert_eq!(
+ replay.first_observed_at().expect("first provenance").get(),
+ 1_784_347_200
+ );
+ assert_eq!(replay.result().accepted_event_count(), 1);
+ assert_eq!(
+ replay.result().accepted_event_bytes(),
+ u64::try_from(wire.len()).expect("wire bytes")
+ );
+ let cursor = replay.eligible_cursor().expect("eligible cursor");
+ assert_eq!(cursor.created_at_unix_seconds(), 1_784_347_200);
+
+ let incomplete =
+ RhiReconciliationSourceReplayPlan::from_request(&plan, request, &configuration, None)
+ .expect("incomplete plan")
+ .finish(
+ request,
+ RhiTradeSourceCompletion::IncompleteUnavailable,
+ now(started_ms),
+ now(started_ms + 1_000),
+ [admitted_replay(&configuration, &wire, 1_784_347_200)],
+ )
+ .expect("incomplete replay");
+ assert!(incomplete.cursor_candidate().is_some());
+ assert!(incomplete.eligible_cursor().is_none());
+ assert!(!format!("{replay:?} {incomplete:?}").contains("trade-primary"));
+ host.close().await.expect("close");
+ drop((configuration, metadata, runtime, root));
+}
+
async fn attempt_fixture(
instance: &str,
) -> (
@@ -940,6 +1015,70 @@ async fn attempt_fixture(
(root, runtime, metadata, configuration, host, plan)
}
+async fn replay_fixture(
+ instance: &str,
+ source: &str,
+ started_ms: u64,
+) -> (
+ 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(source.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(),
+ started_ms / 1_000,
+ )
+ .await;
+ let host = open_writer(&runtime, &metadata).await;
+ let jobs = host.repositories().reconciliation_jobs();
+ jobs.schedule_trade(trade, configured_policy(&configuration), now(started_ms))
+ .await
+ .expect("schedule");
+ let lease = jobs
+ .claim_next(owner(0x72), now(started_ms))
+ .await
+ .expect("claim")
+ .expect("job");
+ let plan = RhiReconciliationAttemptPlan::from_claim(lease, &configuration, now(started_ms))
+ .expect("plan");
+ (root, runtime, metadata, configuration, host, plan)
+}
+
+fn replay_wire() -> Vec<u8> {
+ serde_json::from_str::<serde_json::Value>(TRADE_VECTOR).expect("trade vector")["raw_json"]
+ .as_str()
+ .expect("raw event")
+ .as_bytes()
+ .to_vec()
+}
+
+fn admitted_replay(
+ configuration: &rhi::RhiConfigDocumentV1,
+ wire: &[u8],
+ observed_at: u64,
+) -> rhi::RhiAdmittedTradeMutationEvent {
+ admit_rhi_trade_mutation_event(
+ RhiTradeMutationAdmissionLimits::from_config(configuration).expect("admission limits"),
+ wire,
+ RhiTradeMutationObservedAtUnixSeconds::new(observed_at).expect("observed time"),
+ RhiTradeMutationAuthoredTimePolicy::new(0).expect("authored-time policy"),
+ )
+ .expect("admitted event")
+}
+
#[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)
diff --git a/tests/services_hardening_reconciliation_replay_contract.rs b/tests/services_hardening_reconciliation_replay_contract.rs
@@ -0,0 +1,124 @@
+#![forbid(unsafe_code)]
+
+use rhi::RHI_RECONCILIATION_REPLAY_CONTRACT_VERSION;
+use serde_json::json;
+
+const CONTRACT: &str =
+ include_str!("../contracts/services_hardening/reconciliation_replay.v1.json");
+const ROOT: &str = include_str!("../src/lib.rs");
+const SOURCE: &str = include_str!("../src/reconciliation_replay.rs");
+const README: &str = include_str!("../README");
+
+#[test]
+fn machine_contract_freezes_the_complete_step_189_boundary() {
+ let contract: serde_json::Value = serde_json::from_str(CONTRACT).expect("contract");
+ assert_eq!(contract["schema"], "radroots.rhi.reconciliation-replay");
+ assert_eq!(contract["schema_version"], 1);
+ assert_eq!(
+ contract["contract_version"],
+ RHI_RECONCILIATION_REPLAY_CONTRACT_VERSION
+ );
+ assert_eq!(contract["state_schema_version"], 5);
+ assert_eq!(
+ contract["replay_identity"]["domain"],
+ "radroots.rhi.reconciliation_source_replay.v1\\0"
+ );
+ assert_eq!(
+ contract["replay_identity"]["exact_vector"]["replay_id_hex"],
+ "2490a2e9a6e85051e92f6c2fc2ff7e98a1367afd6c26f8625eb421cd2ab30c68"
+ );
+ assert_eq!(
+ contract["cursor"]["tuple"],
+ json!(["event_authored_unix_seconds", "verified_event_id"])
+ );
+ assert_eq!(contract["cursor"]["equal_timestamp_safe"], true);
+ assert_eq!(
+ contract["cursor"]["input"],
+ "sealed_step_190_committed_cursor_evidence"
+ );
+ assert_eq!(
+ contract["cursor"]["scope"],
+ json!([
+ "source_id",
+ "trade_id",
+ "evidence_policy_digest",
+ "selector_digest"
+ ])
+ );
+ assert_eq!(
+ contract["cursor"]["eligible_only_for"],
+ "complete_and_strictly_after_prior_cursor"
+ );
+ assert_eq!(
+ contract["inventory"]["signed_event_identity"],
+ json!(["verified_event_id", "verified_signature"])
+ );
+ assert_eq!(
+ contract["inventory"]["first_provenance"],
+ "earliest_injected_observation_time_retained"
+ );
+ 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!([
+ "source_execution",
+ "durable_source_result_commit",
+ "durable_scope_revalidation",
+ "checkpoint_advance",
+ "manifest",
+ "reducer",
+ "attestation",
+ "publication"
+ ])
+ );
+}
+
+#[test]
+fn replay_model_is_private_bounded_pure_and_documented() {
+ assert!(ROOT.contains("mod reconciliation_replay;"));
+ assert!(!ROOT.contains("pub mod reconciliation_replay;"));
+ for required in [
+ "RhiReconciliationSourceReplayPlan",
+ "RhiReconciliationSourceReplay",
+ "RhiReconciliationSourceCursorEvidence",
+ "RhiReconciliationReplayError",
+ ] {
+ assert!(ROOT.contains(required), "root API is missing {required}");
+ }
+ for required in [
+ ".take(maximum_events.saturating_add(1))",
+ "saturating_sub(overlap_seconds)",
+ "cursor_scope_matches(",
+ "fn eligible_cursor(",
+ "RhiTradeSourceCompletion::Complete",
+ "RhiReconciliationReplayErrorKind::MutationConflict",
+ "RhiReconciliationReplayErrorKind::SignedEventConflict",
+ ] {
+ assert!(
+ SOURCE.contains(required),
+ "replay boundary is missing {required}"
+ );
+ }
+ for forbidden in [
+ "sqlx::",
+ "std::fs",
+ "std::net",
+ "tokio::",
+ "SystemTime",
+ "thread_rng",
+ "OsRng",
+ "pub fn cursor_evidence",
+ ] {
+ assert!(
+ !SOURCE.contains(forbidden),
+ "replay model gained deferred authority {forbidden}"
+ );
+ }
+ assert!(README.contains("## Overlap-safe reconciliation replay"));
+ assert!(README.contains(
+ "[`reconciliation_replay.v1.json`](contracts/services_hardening/reconciliation_replay.v1.json)"
+ ));
+}