commit ba0e5a17937b62df80eefa46d47e2202ff79ce09
parent 79cc730521a60d11b7c08bd391d8f839d4dd0c2c
Author: triesap <tyson@radroots.org>
Date: Mon, 24 Aug 2026 05:33:59 +0000
refactor(rhi): add durable reconciliation leases
- add the pinned schema-v5 job catalog and deterministic job identities
- enforce bounded CAS scheduling, claims, renewal, reclaim, and retry state
- freeze the curated public API, machine contract, and redacted diagnostics
- cover migration identity, limits, concurrency, restart, and lease expiry
Diffstat:
18 files changed, 2394 insertions(+), 41 deletions(-)
diff --git a/AGENTS.md b/AGENTS.md
@@ -205,8 +205,9 @@
boundary and offline state operations must prove that no daemon writer exists.
- Create-new state begins at the shared schema-v1 baseline and applies the
governed RHI schema-v2 configuration-binding migration, schema-v3
- immutable trade-evidence migration, and schema-v4 source-checkpoint and
- dirty-generation migration. Retain at most 1,024 consecutive
+ immutable trade-evidence migration, schema-v4 source-checkpoint and
+ dirty-generation migration, and schema-v5 bounded reconciliation-job
+ migration. Retain at most 1,024 consecutive
immutable configuration generations containing only normalized
config/evidence-policy digests, public identity, exact contract versions,
injected apply time, and bounded build identity. Persist each canonical
diff --git a/Cargo.toml b/Cargo.toml
@@ -20,7 +20,7 @@ service = "rhi"
host_feature_profile = "service-host"
nix_material = "absent"
config_contract_version = 1
-state_contract_version = 4
+state_contract_version = 5
admin_contract_version = 1
status_contract_version = 1
provider_contract_version = 1
diff --git a/README b/README
@@ -119,6 +119,26 @@ observation, and operational retry do not. Source fetching never occurs inside
a database transaction. The exact machine contract is
[`trade_source_ingest.v1.json`](contracts/services_hardening/trade_source_ingest.v1.json).
+## Durable reconciliation jobs
+
+`RhiReconciliationJobRepository` schedules at most one active job per trade
+against the exact durable dirty generation and evidence-policy digest. Job IDs
+are deterministic domain-separated SHA-256 identities, exact scheduling is
+idempotent, a newer dirty generation atomically supersedes the prior active
+job, and the configured queue capacity is enforced beneath a fixed 65,536-job
+hard ceiling.
+
+Claims and renewals are compare-and-swap transitions over an immutable job
+identity and monotonic revision. Leases use caller-injected 16-byte owner
+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
+[`reconciliation_jobs.v1.json`](contracts/services_hardening/reconciliation_jobs.v1.json).
+
## Existing-state runtime foundation
`open_rhi_runtime_foundation` opens only an already initialized database from
@@ -216,8 +236,9 @@ Path resolution performs no directory creation or filesystem I/O.
RHI owns one `state.sqlite` per service instance. Create-new initialization
starts from the shared schema-v1 baseline and immediately applies the pinned
schema-v2 configuration migration, schema-v3 immutable trade-evidence
-migration, and schema-v4 source-checkpoint and dirty-generation migration.
-Version four contains the six shared immutable service-metadata and
+migration, schema-v4 source-checkpoint and dirty-generation migration, and
+schema-v5 bounded reconciliation-job migration.
+Version five contains the six shared immutable service-metadata and
migration-ledger objects, the bounded append-only `rhi_config_bindings` table,
separate immutable tables for canonical mutations, signed Nostr events, and
accepted source observations, and generation-guarded relay checkpoints and
diff --git a/contracts/api_baselines/rhi.txt b/contracts/api_baselines/rhi.txt
@@ -141,6 +141,25 @@ pub enum rhi::RhiReconciliationCommandV1
pub rhi::RhiReconciliationCommandV1::Jobs
pub rhi::RhiReconciliationCommandV1::Refresh
pub rhi::RhiReconciliationCommandV1::Status
+pub enum rhi::RhiReconciliationJobErrorKind
+pub rhi::RhiReconciliationJobErrorKind::CommitOutcomeUnknown
+pub rhi::RhiReconciliationJobErrorKind::DirtyGenerationConflict
+pub rhi::RhiReconciliationJobErrorKind::InvalidInput
+pub rhi::RhiReconciliationJobErrorKind::InvalidMode
+pub rhi::RhiReconciliationJobErrorKind::LeaseLost
+pub rhi::RhiReconciliationJobErrorKind::NotReady
+pub rhi::RhiReconciliationJobErrorKind::QueueFull
+pub rhi::RhiReconciliationJobErrorKind::Storage
+impl rhi::RhiReconciliationJobErrorKind
+pub const fn rhi::RhiReconciliationJobErrorKind::code(self) -> &'static str
+pub enum rhi::RhiReconciliationJobState
+pub rhi::RhiReconciliationJobState::Completed
+pub rhi::RhiReconciliationJobState::Exhausted
+pub rhi::RhiReconciliationJobState::Leased
+pub rhi::RhiReconciliationJobState::Ready
+pub rhi::RhiReconciliationJobState::Superseded
+impl rhi::RhiReconciliationJobState
+pub const fn rhi::RhiReconciliationJobState::code(self) -> &'static str
pub enum rhi::RhiRuntimeAdapterErrorKind
pub rhi::RhiRuntimeAdapterErrorKind::CredentialAccess
pub rhi::RhiRuntimeAdapterErrorKind::EntropyUnavailable
@@ -558,12 +577,81 @@ pub const fn rhi::RhiReconciliationAttemptRepository<'_>::descriptor(&self) -> r
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::RhiReconciliationJob
+impl rhi::RhiReconciliationJob
+pub const fn rhi::RhiReconciliationJob::attempt_count(self) -> u16
+pub const fn rhi::RhiReconciliationJob::created_at(self) -> rhi::RhiReconciliationUnixMilliseconds
+pub const fn rhi::RhiReconciliationJob::evidence_policy_digest(self) -> rhi::RhiEvidencePolicyDigest
+pub const fn rhi::RhiReconciliationJob::failure_count(self) -> u16
+pub const fn rhi::RhiReconciliationJob::id(self) -> rhi::RhiReconciliationJobId
+pub const fn rhi::RhiReconciliationJob::input_generation(self) -> u64
+pub const fn rhi::RhiReconciliationJob::next_attempt(self) -> core::option::Option<rhi::RhiReconciliationUnixMilliseconds>
+pub const fn rhi::RhiReconciliationJob::revision(self) -> u64
+pub const fn rhi::RhiReconciliationJob::state(self) -> rhi::RhiReconciliationJobState
+pub const fn rhi::RhiReconciliationJob::trade_id(self) -> radroots_event::id::TradeId
+pub const fn rhi::RhiReconciliationJob::updated_at(self) -> rhi::RhiReconciliationUnixMilliseconds
+impl core::fmt::Debug for rhi::RhiReconciliationJob
+pub fn rhi::RhiReconciliationJob::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct rhi::RhiReconciliationJobError
+impl rhi::RhiReconciliationJobError
+pub const fn rhi::RhiReconciliationJobError::code(self) -> &'static str
+pub const fn rhi::RhiReconciliationJobError::kind(self) -> rhi::RhiReconciliationJobErrorKind
+impl core::error::Error for rhi::RhiReconciliationJobError
+impl core::fmt::Debug for rhi::RhiReconciliationJobError
+pub fn rhi::RhiReconciliationJobError::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+impl core::fmt::Display for rhi::RhiReconciliationJobError
+pub fn rhi::RhiReconciliationJobError::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct rhi::RhiReconciliationJobId(_)
+impl rhi::RhiReconciliationJobId
+pub const fn rhi::RhiReconciliationJobId::as_bytes(&self) -> &[u8; 32]
+impl core::fmt::Debug for rhi::RhiReconciliationJobId
+pub fn rhi::RhiReconciliationJobId::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct rhi::RhiReconciliationJobPolicy
+impl rhi::RhiReconciliationJobPolicy
+pub fn rhi::RhiReconciliationJobPolicy::from_configuration(&rhi::RhiConfigDocumentV1) -> core::result::Result<Self, rhi::RhiReconciliationJobError>
+pub const fn rhi::RhiReconciliationJobPolicy::initial_backoff_milliseconds(self) -> u64
+pub const fn rhi::RhiReconciliationJobPolicy::lease_duration_milliseconds(self) -> u64
+pub const fn rhi::RhiReconciliationJobPolicy::lease_renewal_milliseconds(self) -> u64
+pub const fn rhi::RhiReconciliationJobPolicy::max_attempts(self) -> u16
+pub const fn rhi::RhiReconciliationJobPolicy::maximum_backoff_milliseconds(self) -> u64
+pub fn rhi::RhiReconciliationJobPolicy::new(u32, u64, u64, u16, u64, u64) -> core::result::Result<Self, rhi::RhiReconciliationJobError>
+pub const fn rhi::RhiReconciliationJobPolicy::queue_capacity(self) -> u32
pub struct rhi::RhiReconciliationJobRepository<'host>
impl rhi::RhiReconciliationJobRepository<'_>
+pub async fn rhi::RhiReconciliationJobRepository<'_>::claim_next(&self, rhi::RhiReconciliationLeaseOwner, rhi::RhiReconciliationUnixMilliseconds) -> core::result::Result<core::option::Option<rhi::RhiReconciliationLease>, rhi::RhiReconciliationJobError>
+pub async fn rhi::RhiReconciliationJobRepository<'_>::record_failure(&self, rhi::RhiReconciliationLease, rhi::RhiReconciliationUnixMilliseconds, rhi::RhiReconciliationRetryDelayMilliseconds) -> core::result::Result<rhi::RhiReconciliationJob, rhi::RhiReconciliationJobError>
+pub async fn rhi::RhiReconciliationJobRepository<'_>::renew(&self, rhi::RhiReconciliationLease, rhi::RhiReconciliationUnixMilliseconds) -> core::result::Result<rhi::RhiReconciliationLease, rhi::RhiReconciliationJobError>
+pub async fn rhi::RhiReconciliationJobRepository<'_>::schedule_trade(&self, radroots_event::id::TradeId, rhi::RhiReconciliationJobPolicy, rhi::RhiReconciliationUnixMilliseconds) -> core::result::Result<rhi::RhiReconciliationScheduleOutcome, rhi::RhiReconciliationJobError>
+impl rhi::RhiReconciliationJobRepository<'_>
pub const fn rhi::RhiReconciliationJobRepository<'_>::descriptor(&self) -> rhi::RhiStateRepositoryDescriptor
pub const fn rhi::RhiReconciliationJobRepository<'_>::kind(&self) -> rhi::RhiStateRepositoryKind
impl core::fmt::Debug for rhi::RhiReconciliationJobRepository<'_>
pub fn rhi::RhiReconciliationJobRepository<'_>::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct rhi::RhiReconciliationLease
+impl rhi::RhiReconciliationLease
+pub const fn rhi::RhiReconciliationLease::job(self) -> rhi::RhiReconciliationJob
+pub const fn rhi::RhiReconciliationLease::lease_expires(self) -> rhi::RhiReconciliationUnixMilliseconds
+pub fn rhi::RhiReconciliationLease::renewal_due(self) -> rhi::RhiReconciliationUnixMilliseconds
+pub fn rhi::RhiReconciliationLease::retry_delay_upper_bound(self) -> u64
+impl core::fmt::Debug for rhi::RhiReconciliationLease
+pub fn rhi::RhiReconciliationLease::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct rhi::RhiReconciliationLeaseOwner(_)
+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::RhiReconciliationRetryDelayMilliseconds(_)
+impl rhi::RhiReconciliationRetryDelayMilliseconds
+pub const fn rhi::RhiReconciliationRetryDelayMilliseconds::get(self) -> u64
+pub fn rhi::RhiReconciliationRetryDelayMilliseconds::new(u64) -> core::result::Result<Self, rhi::RhiReconciliationJobError>
+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::RhiReconciliationUnixMilliseconds(_)
+impl rhi::RhiReconciliationUnixMilliseconds
+pub const fn rhi::RhiReconciliationUnixMilliseconds::get(self) -> u64
+pub fn rhi::RhiReconciliationUnixMilliseconds::new(u64) -> core::result::Result<Self, rhi::RhiReconciliationJobError>
pub struct rhi::RhiReportRepository<'host>
impl rhi::RhiReportRepository<'_>
pub const fn rhi::RhiReportRepository<'_>::descriptor(&self) -> rhi::RhiStateRepositoryDescriptor
@@ -941,6 +1029,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_JOB_CONTRACT_VERSION: u32
+pub const rhi::RHI_RECONCILIATION_JOB_MAX_ACTIVE: 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
@@ -962,6 +1052,9 @@ pub const rhi::RHI_STATE_SCHEMA_VERSION_3_SHA256: [u8; 32]
pub const rhi::RHI_STATE_SCHEMA_VERSION_4_MIGRATION_SHA256: [u8; 32]
pub const rhi::RHI_STATE_SCHEMA_VERSION_4_OBJECT_COUNT: u32
pub const rhi::RHI_STATE_SCHEMA_VERSION_4_SHA256: [u8; 32]
+pub const rhi::RHI_STATE_SCHEMA_VERSION_5_MIGRATION_SHA256: [u8; 32]
+pub const rhi::RHI_STATE_SCHEMA_VERSION_5_OBJECT_COUNT: u32
+pub const rhi::RHI_STATE_SCHEMA_VERSION_5_SHA256: [u8; 32]
pub const rhi::RHI_STATUS_CONTRACT_VERSION: u32
pub const rhi::RHI_TRADE_EVENT_EXTRA_FIELD_MAX_COUNT: usize
pub const rhi::RHI_TRADE_EVENT_EXTRA_JSON_MAX_BYTES: usize
diff --git a/contracts/services_hardening/reconciliation_jobs.v1.json b/contracts/services_hardening/reconciliation_jobs.v1.json
@@ -0,0 +1,66 @@
+{
+ "schema": "radroots.rhi.reconciliation-jobs",
+ "schema_version": 1,
+ "contract_version": 1,
+ "state_schema_version": 5,
+ "backing_table": "reconciliation_jobs",
+ "job_identity": {
+ "algorithm": "sha256",
+ "domain": "radroots.rhi.reconciliation_job.v1\\0",
+ "preimage": ["trade_id_16_bytes", "input_generation_u64_be", "evidence_policy_sha256_32_bytes"]
+ },
+ "active_job_bound": {
+ "absolute_maximum": 65536,
+ "configured_field": "reconciliation.queue_capacity",
+ "one_active_per_trade": true
+ },
+ "states": ["ready", "leased", "exhausted", "superseded", "completed"],
+ "claim": {
+ "ordering": ["eligible_at_unix_ms", "created_at_unix_ms", "job_id"],
+ "compare_and_swap": ["job_id", "revision", "state", "eligibility"],
+ "attempt_count_increment": "exactly_once_per_successful_claim_or_reclaim",
+ "source_io_inside_transaction": false
+ },
+ "lease": {
+ "owner_bytes": 16,
+ "duration_ms": { "minimum": 1000, "maximum": 300000 },
+ "renewal_ms": { "minimum": 100, "maximum": 150000 },
+ "renewal_requires_unexpired_current_revision": true,
+ "expired_claims_are_reclaimable": true,
+ "expired_final_attempt_becomes_exhausted": true
+ },
+ "retry": {
+ "max_attempts": { "minimum": 1, "maximum": 100 },
+ "initial_backoff_ms": { "minimum": 1, "maximum": 60000 },
+ "maximum_backoff_ms": { "minimum": 1, "maximum": 3600000 },
+ "strategy": "caller_injected_full_jitter_with_governed_exponential_upper_bound",
+ "next_attempt_persisted_as": "integer_unix_milliseconds"
+ },
+ "durability": {
+ "schedule_is_idempotent_for_exact_job_identity": true,
+ "newer_dirty_generation_supersedes_prior_active_job_atomically": true,
+ "jobs_are_not_deleted": true,
+ "restart_recovers_persisted_schedule_and_lease_expiry": true
+ },
+ "injected_authority": ["wall_time_unix_ms", "lease_owner", "retry_jitter_ms"],
+ "forbidden": [
+ "ambient_clock",
+ "ambient_entropy",
+ "network_io",
+ "source_query_inside_transaction",
+ "unbounded_queue",
+ "raw_sqlite_handle",
+ "caller_selected_job_id",
+ "caller_selected_next_attempt_timestamp",
+ "delete_job"
+ ],
+ "deferred": [
+ "per_source_attempt_inventory",
+ "source_execution",
+ "source_completion_commit",
+ "manifest",
+ "reducer",
+ "attestation",
+ "publication"
+ ]
+}
diff --git a/contracts/services_hardening/state_repository_topology.v1.json b/contracts/services_hardening/state_repository_topology.v1.json
@@ -26,8 +26,8 @@
{ "kind": "desired_presence", "backing_table": "presence_desired_state", "write_class": "compare_and_swap" }
],
"deferred_behavior": [
- "schema_migration_and_verification",
- "repository_crud",
+ "later_schema_migrations",
+ "non_job_repository_crud",
"backup_restore_and_recovery",
"network_io",
"task_supervision"
diff --git a/radroots.service.source-lock.v2.toml b/radroots.service.source-lock.v2.toml
@@ -16,7 +16,7 @@ material = "absent"
[contract_versions]
config = 1
-state = 4
+state = 5
admin = 1
status = 1
provider = 1
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_job;
mod runtime_adapters;
mod runtime_context;
mod runtime_foundation;
@@ -66,6 +67,13 @@ pub use radroots_service_host::{
EntropyError, EntropySource, MonotonicClock, MonotonicClockError, MonotonicDeadline,
MonotonicTime, UnixTimeSeconds, WallClock, WallClockError,
};
+pub use reconciliation_job::{
+ RHI_RECONCILIATION_JOB_CONTRACT_VERSION, RHI_RECONCILIATION_JOB_MAX_ACTIVE,
+ RhiReconciliationJob, RhiReconciliationJobError, RhiReconciliationJobErrorKind,
+ RhiReconciliationJobId, RhiReconciliationJobPolicy, RhiReconciliationJobState,
+ RhiReconciliationLease, RhiReconciliationLeaseOwner, RhiReconciliationRetryDelayMilliseconds,
+ RhiReconciliationScheduleOutcome, RhiReconciliationUnixMilliseconds,
+};
pub use runtime_adapters::{
CanonicalRhiCredentialAccess, CanonicalRhiIdentityAccess, RHI_RUNTIME_ADAPTER_CONTRACT_VERSION,
RHI_RUNTIME_JITTER_MAX_ENTROPY_DRAWS, RHI_RUNTIME_JITTER_MAX_MILLISECONDS, RhiCredentialAccess,
@@ -96,8 +104,9 @@ pub use state_catalog::{
RHI_STATE_SCHEMA_VERSION_3_MIGRATION_SHA256, RHI_STATE_SCHEMA_VERSION_3_OBJECT_COUNT,
RHI_STATE_SCHEMA_VERSION_3_SHA256, RHI_STATE_SCHEMA_VERSION_4_MIGRATION_SHA256,
RHI_STATE_SCHEMA_VERSION_4_OBJECT_COUNT, RHI_STATE_SCHEMA_VERSION_4_SHA256,
- RhiStateCatalogError, RhiStateCatalogErrorKind, rhi_migration_catalog, rhi_schema_catalog,
- validate_rhi_state_catalogs,
+ RHI_STATE_SCHEMA_VERSION_5_MIGRATION_SHA256, RHI_STATE_SCHEMA_VERSION_5_OBJECT_COUNT,
+ RHI_STATE_SCHEMA_VERSION_5_SHA256, RhiStateCatalogError, RhiStateCatalogErrorKind,
+ rhi_migration_catalog, rhi_schema_catalog, validate_rhi_state_catalogs,
};
pub use state_config::{
RHI_CONFIG_BINDING_MAX_GENERATIONS, RhiConfigApplyError, RhiConfigApplyErrorKind,
diff --git a/src/reconciliation_job.rs b/src/reconciliation_job.rs
@@ -0,0 +1,1238 @@
+//! Bounded durable reconciliation-job scheduling and lease state machine.
+
+use core::fmt;
+use std::error::Error;
+
+use radroots_event::id::TradeId;
+use radroots_service_sqlite::{
+ ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind,
+};
+use sha2::{Digest, Sha256};
+use sqlx::Row;
+
+use crate::{
+ RhiConfigDocumentV1, RhiEvidencePolicyDigest, RhiReconciliationJobRepository, RhiStateHostMode,
+};
+
+/// Exact version of the durable reconciliation-job contract.
+pub const RHI_RECONCILIATION_JOB_CONTRACT_VERSION: u32 = 1;
+
+/// Absolute hard ceiling for active durable reconciliation jobs.
+pub const RHI_RECONCILIATION_JOB_MAX_ACTIVE: u32 = 65_536;
+
+const JOB_ID_DOMAIN: &[u8] = b"radroots.rhi.reconciliation_job.v1\0";
+const LEASE_OWNER_BYTES: usize = 16;
+const MAX_UNIX_MILLISECONDS: u64 = i64::MAX as u64;
+const MAX_ATTEMPTS: u16 = 100;
+const MAX_LEASE_MILLISECONDS: u64 = 300_000;
+const MAX_RENEWAL_MILLISECONDS: u64 = 150_000;
+const MAX_INITIAL_BACKOFF_MILLISECONDS: u64 = 60_000;
+const MAX_BACKOFF_MILLISECONDS: u64 = 3_600_000;
+
+const READ_DIRTY_SQL: &str = r#"SELECT generation,
+ length(evidence_policy_sha256) AS evidence_policy_bytes,
+ substr(evidence_policy_sha256, 1, 33) AS evidence_policy_sha256
+FROM trade_dirty_generations
+WHERE trade_id = ?
+LIMIT 1"#;
+const READ_JOB_SQL: &str = r#"SELECT
+ length(job_id) AS job_id_bytes, substr(job_id, 1, 33) AS job_id,
+ length(trade_id) AS trade_id_bytes, substr(trade_id, 1, 17) AS trade_id,
+ input_generation,
+ length(evidence_policy_sha256) AS evidence_policy_bytes,
+ substr(evidence_policy_sha256, 1, 33) AS evidence_policy_sha256,
+ length(CAST(state AS BLOB)) AS state_bytes, substr(state, 1, 12) AS state,
+ revision, attempt_count, failure_count, max_attempts,
+ lease_duration_ms, lease_renewal_ms, initial_backoff_ms, maximum_backoff_ms,
+ next_attempt_unix_ms,
+ CASE WHEN lease_owner IS NULL THEN NULL ELSE length(lease_owner) END AS lease_owner_bytes,
+ CASE WHEN lease_owner IS NULL THEN NULL ELSE substr(lease_owner, 1, 17) END AS lease_owner,
+ lease_expires_unix_ms, created_at_unix_ms, updated_at_unix_ms
+FROM reconciliation_jobs
+WHERE job_id = ?
+LIMIT 1"#;
+const READ_ACTIVE_JOB_SQL: &str = r#"SELECT
+ length(job_id) AS job_id_bytes, substr(job_id, 1, 33) AS job_id,
+ length(trade_id) AS trade_id_bytes, substr(trade_id, 1, 17) AS trade_id,
+ input_generation,
+ length(evidence_policy_sha256) AS evidence_policy_bytes,
+ substr(evidence_policy_sha256, 1, 33) AS evidence_policy_sha256,
+ length(CAST(state AS BLOB)) AS state_bytes, substr(state, 1, 12) AS state,
+ revision, attempt_count, failure_count, max_attempts,
+ lease_duration_ms, lease_renewal_ms, initial_backoff_ms, maximum_backoff_ms,
+ next_attempt_unix_ms,
+ CASE WHEN lease_owner IS NULL THEN NULL ELSE length(lease_owner) END AS lease_owner_bytes,
+ CASE WHEN lease_owner IS NULL THEN NULL ELSE substr(lease_owner, 1, 17) END AS lease_owner,
+ lease_expires_unix_ms, created_at_unix_ms, updated_at_unix_ms
+FROM reconciliation_jobs
+WHERE trade_id = ? AND state IN ('ready', 'leased')
+LIMIT 2"#;
+const ACTIVE_JOB_COUNT_SQL: &str = r#"SELECT COUNT(*) AS active_count
+FROM reconciliation_jobs
+WHERE state IN ('ready', 'leased')"#;
+const SUPERSEDE_JOB_SQL: &str = r#"UPDATE reconciliation_jobs
+SET state = 'superseded', revision = revision + 1,
+ next_attempt_unix_ms = NULL, lease_owner = NULL, lease_expires_unix_ms = NULL,
+ updated_at_unix_ms = ?
+WHERE job_id = ? AND revision = ? AND state IN ('ready', 'leased')"#;
+const INSERT_JOB_SQL: &str = r#"INSERT INTO reconciliation_jobs (
+ job_id, trade_id, input_generation, evidence_policy_sha256, state, revision,
+ attempt_count, failure_count, max_attempts, lease_duration_ms,
+ lease_renewal_ms, initial_backoff_ms, maximum_backoff_ms,
+ next_attempt_unix_ms, lease_owner, lease_expires_unix_ms,
+ created_at_unix_ms, updated_at_unix_ms
+) VALUES (?, ?, ?, ?, 'ready', 1, 0, 0, ?, ?, ?, ?, ?, ?, NULL, NULL, ?, ?)"#;
+const EXHAUST_EXPIRED_SQL: &str = r#"UPDATE reconciliation_jobs
+SET state = 'exhausted', revision = revision + 1,
+ failure_count = attempt_count, lease_owner = NULL, lease_expires_unix_ms = NULL,
+ updated_at_unix_ms = ?
+WHERE job_id IN (
+ SELECT job_id FROM reconciliation_jobs
+ WHERE state = 'leased' AND lease_expires_unix_ms <= ?
+ AND attempt_count >= max_attempts AND updated_at_unix_ms <= ?
+ ORDER BY lease_expires_unix_ms, created_at_unix_ms, job_id
+ LIMIT 65536
+)"#;
+const READ_CLAIMABLE_SQL: &str = r#"SELECT
+ length(job_id) AS job_id_bytes, substr(job_id, 1, 33) AS job_id,
+ length(trade_id) AS trade_id_bytes, substr(trade_id, 1, 17) AS trade_id,
+ input_generation,
+ length(evidence_policy_sha256) AS evidence_policy_bytes,
+ substr(evidence_policy_sha256, 1, 33) AS evidence_policy_sha256,
+ length(CAST(state AS BLOB)) AS state_bytes, substr(state, 1, 12) AS state,
+ revision, attempt_count, failure_count, max_attempts,
+ lease_duration_ms, lease_renewal_ms, initial_backoff_ms, maximum_backoff_ms,
+ next_attempt_unix_ms,
+ CASE WHEN lease_owner IS NULL THEN NULL ELSE length(lease_owner) END AS lease_owner_bytes,
+ CASE WHEN lease_owner IS NULL THEN NULL ELSE substr(lease_owner, 1, 17) END AS lease_owner,
+ lease_expires_unix_ms, created_at_unix_ms, updated_at_unix_ms
+FROM reconciliation_jobs
+WHERE attempt_count < max_attempts AND updated_at_unix_ms <= ? AND (
+ (state = 'ready' AND next_attempt_unix_ms <= ?)
+ OR (state = 'leased' AND lease_expires_unix_ms <= ?)
+)
+ORDER BY
+ CASE state WHEN 'ready' THEN next_attempt_unix_ms ELSE lease_expires_unix_ms END,
+ created_at_unix_ms, job_id
+LIMIT 1"#;
+const CLAIM_JOB_SQL: &str = r#"UPDATE reconciliation_jobs
+SET state = 'leased', revision = revision + 1, attempt_count = attempt_count + 1,
+ next_attempt_unix_ms = NULL, lease_owner = ?, lease_expires_unix_ms = ?,
+ updated_at_unix_ms = ?
+WHERE job_id = ? AND revision = ? AND attempt_count < max_attempts
+ AND updated_at_unix_ms <= ? AND (
+ (state = 'ready' AND next_attempt_unix_ms <= ?)
+ OR (state = 'leased' AND lease_expires_unix_ms <= ?)
+ )"#;
+const RENEW_JOB_SQL: &str = r#"UPDATE reconciliation_jobs
+SET revision = revision + 1, lease_expires_unix_ms = ?, updated_at_unix_ms = ?
+WHERE job_id = ? AND revision = ? AND state = 'leased'
+ AND lease_owner = ? AND lease_expires_unix_ms = ?
+ AND lease_expires_unix_ms > ? AND updated_at_unix_ms <= ?"#;
+const RETRY_JOB_SQL: &str = r#"UPDATE reconciliation_jobs
+SET state = 'ready', revision = revision + 1, failure_count = failure_count + 1,
+ next_attempt_unix_ms = ?, lease_owner = NULL, lease_expires_unix_ms = NULL,
+ updated_at_unix_ms = ?
+WHERE job_id = ? AND revision = ? AND state = 'leased'
+ AND lease_owner = ? AND lease_expires_unix_ms = ?
+ AND lease_expires_unix_ms > ? AND attempt_count < max_attempts"#;
+const EXHAUST_JOB_SQL: &str = r#"UPDATE reconciliation_jobs
+SET state = 'exhausted', revision = revision + 1, failure_count = failure_count + 1,
+ next_attempt_unix_ms = NULL, lease_owner = NULL, lease_expires_unix_ms = NULL,
+ updated_at_unix_ms = ?
+WHERE job_id = ? AND revision = ? AND state = 'leased'
+ AND lease_owner = ? AND lease_expires_unix_ms = ?
+ AND lease_expires_unix_ms > ? AND attempt_count >= max_attempts"#;
+
+/// Stable durable reconciliation-job lifecycle state.
+#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
+pub enum RhiReconciliationJobState {
+ Ready,
+ Leased,
+ Exhausted,
+ Superseded,
+ Completed,
+}
+
+impl RhiReconciliationJobState {
+ /// Returns the exact machine-contract spelling.
+ #[must_use]
+ pub const fn code(self) -> &'static str {
+ match self {
+ Self::Ready => "ready",
+ Self::Leased => "leased",
+ Self::Exhausted => "exhausted",
+ Self::Superseded => "superseded",
+ Self::Completed => "completed",
+ }
+ }
+}
+
+/// Stable source-free durable job failure class.
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+pub enum RhiReconciliationJobErrorKind {
+ InvalidMode,
+ InvalidInput,
+ QueueFull,
+ DirtyGenerationConflict,
+ LeaseLost,
+ NotReady,
+ Storage,
+ CommitOutcomeUnknown,
+}
+
+impl RhiReconciliationJobErrorKind {
+ /// Returns the stable machine-readable failure code.
+ #[must_use]
+ pub const fn code(self) -> &'static str {
+ match self {
+ Self::InvalidMode => "reconciliation_job_mode_invalid",
+ Self::InvalidInput => "reconciliation_job_input_invalid",
+ Self::QueueFull => "reconciliation_queue_full",
+ Self::DirtyGenerationConflict => "reconciliation_dirty_generation_conflict",
+ Self::LeaseLost => "reconciliation_lease_lost",
+ Self::NotReady => "reconciliation_job_not_ready",
+ Self::Storage => "reconciliation_job_storage_failed",
+ Self::CommitOutcomeUnknown => "reconciliation_job_commit_outcome_unknown",
+ }
+ }
+}
+
+/// Redacted source-free durable job failure.
+#[derive(Clone, Copy, PartialEq, Eq)]
+pub struct RhiReconciliationJobError {
+ kind: RhiReconciliationJobErrorKind,
+}
+
+impl RhiReconciliationJobError {
+ const fn new(kind: RhiReconciliationJobErrorKind) -> Self {
+ Self { kind }
+ }
+
+ /// Returns the stable failure class.
+ #[must_use]
+ pub const fn kind(self) -> RhiReconciliationJobErrorKind {
+ 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 RhiReconciliationJobError {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter.write_str(match self.kind {
+ RhiReconciliationJobErrorKind::InvalidMode => {
+ "RHI reconciliation jobs require writable state"
+ }
+ RhiReconciliationJobErrorKind::InvalidInput => {
+ "RHI reconciliation job input is invalid"
+ }
+ RhiReconciliationJobErrorKind::QueueFull => {
+ "RHI reconciliation job capacity is exhausted"
+ }
+ RhiReconciliationJobErrorKind::DirtyGenerationConflict => {
+ "RHI reconciliation dirty generation changed"
+ }
+ RhiReconciliationJobErrorKind::LeaseLost => {
+ "RHI reconciliation lease is no longer authoritative"
+ }
+ RhiReconciliationJobErrorKind::NotReady => "RHI reconciliation job is not ready",
+ RhiReconciliationJobErrorKind::Storage => "RHI reconciliation job transaction failed",
+ RhiReconciliationJobErrorKind::CommitOutcomeUnknown => {
+ "RHI reconciliation job commit outcome is unknown"
+ }
+ })
+ }
+}
+
+impl fmt::Debug for RhiReconciliationJobError {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter
+ .debug_struct("RhiReconciliationJobError")
+ .field("kind", &self.kind)
+ .finish()
+ }
+}
+
+impl Error for RhiReconciliationJobError {}
+
+/// Validated numeric scheduling authority copied into each durable job.
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+pub struct RhiReconciliationJobPolicy {
+ queue_capacity: u32,
+ lease_duration_ms: u64,
+ lease_renewal_ms: u64,
+ max_attempts: u16,
+ initial_backoff_ms: u64,
+ maximum_backoff_ms: u64,
+}
+
+impl RhiReconciliationJobPolicy {
+ /// Constructs an explicit bounded policy with no ambient defaults.
+ pub fn new(
+ queue_capacity: u32,
+ lease_duration_ms: u64,
+ lease_renewal_ms: u64,
+ max_attempts: u16,
+ initial_backoff_ms: u64,
+ maximum_backoff_ms: u64,
+ ) -> Result<Self, RhiReconciliationJobError> {
+ if queue_capacity == 0
+ || queue_capacity > RHI_RECONCILIATION_JOB_MAX_ACTIVE
+ || !(1_000..=MAX_LEASE_MILLISECONDS).contains(&lease_duration_ms)
+ || !(100..=MAX_RENEWAL_MILLISECONDS).contains(&lease_renewal_ms)
+ || lease_renewal_ms >= lease_duration_ms
+ || max_attempts == 0
+ || max_attempts > MAX_ATTEMPTS
+ || !(1..=MAX_INITIAL_BACKOFF_MILLISECONDS).contains(&initial_backoff_ms)
+ || !(1..=MAX_BACKOFF_MILLISECONDS).contains(&maximum_backoff_ms)
+ || initial_backoff_ms > maximum_backoff_ms
+ {
+ return Err(RhiReconciliationJobError::new(
+ RhiReconciliationJobErrorKind::InvalidInput,
+ ));
+ }
+ Ok(Self {
+ queue_capacity,
+ lease_duration_ms,
+ lease_renewal_ms,
+ max_attempts,
+ initial_backoff_ms,
+ maximum_backoff_ms,
+ })
+ }
+
+ /// Extracts the exact admitted reconciliation policy from one immutable configuration.
+ pub fn from_configuration(
+ configuration: &RhiConfigDocumentV1,
+ ) -> Result<Self, RhiReconciliationJobError> {
+ let value = |pointer: &str| {
+ configuration
+ .normalized()
+ .pointer(pointer)
+ .and_then(serde_json::Value::as_u64)
+ .ok_or_else(|| {
+ RhiReconciliationJobError::new(RhiReconciliationJobErrorKind::InvalidInput)
+ })
+ };
+ Self::new(
+ u32::try_from(value("/reconciliation/queue_capacity")?).map_err(|_| {
+ RhiReconciliationJobError::new(RhiReconciliationJobErrorKind::InvalidInput)
+ })?,
+ value("/reconciliation/lease_ms")?,
+ value("/reconciliation/lease_renewal_ms")?,
+ u16::try_from(value("/reconciliation/max_attempts")?).map_err(|_| {
+ RhiReconciliationJobError::new(RhiReconciliationJobErrorKind::InvalidInput)
+ })?,
+ value("/reconciliation/initial_backoff_ms")?,
+ value("/reconciliation/maximum_backoff_ms")?,
+ )
+ }
+
+ #[must_use]
+ pub const fn queue_capacity(self) -> u32 {
+ self.queue_capacity
+ }
+
+ #[must_use]
+ pub const fn lease_duration_milliseconds(self) -> u64 {
+ self.lease_duration_ms
+ }
+
+ #[must_use]
+ pub const fn lease_renewal_milliseconds(self) -> u64 {
+ self.lease_renewal_ms
+ }
+
+ #[must_use]
+ pub const fn max_attempts(self) -> u16 {
+ self.max_attempts
+ }
+
+ #[must_use]
+ pub const fn initial_backoff_milliseconds(self) -> u64 {
+ self.initial_backoff_ms
+ }
+
+ #[must_use]
+ pub const fn maximum_backoff_milliseconds(self) -> u64 {
+ self.maximum_backoff_ms
+ }
+}
+
+/// Injected wall-clock millisecond used only for durable scheduling evidence.
+#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)]
+pub struct RhiReconciliationUnixMilliseconds(u64);
+
+impl RhiReconciliationUnixMilliseconds {
+ pub fn new(value: u64) -> Result<Self, RhiReconciliationJobError> {
+ if value > MAX_UNIX_MILLISECONDS {
+ return Err(RhiReconciliationJobError::new(
+ RhiReconciliationJobErrorKind::InvalidInput,
+ ));
+ }
+ Ok(Self(value))
+ }
+
+ #[must_use]
+ pub const fn get(self) -> u64 {
+ self.0
+ }
+}
+
+/// Injected full-jitter delay for one failed reconciliation attempt.
+#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)]
+pub struct RhiReconciliationRetryDelayMilliseconds(u64);
+
+impl RhiReconciliationRetryDelayMilliseconds {
+ pub fn new(value: u64) -> Result<Self, RhiReconciliationJobError> {
+ if value > MAX_BACKOFF_MILLISECONDS {
+ return Err(RhiReconciliationJobError::new(
+ RhiReconciliationJobErrorKind::InvalidInput,
+ ));
+ }
+ Ok(Self(value))
+ }
+
+ #[must_use]
+ pub const fn get(self) -> u64 {
+ self.0
+ }
+}
+
+/// Stable process-local owner token for compare-and-swap leases.
+#[derive(Clone, Copy, PartialEq, Eq, Hash)]
+pub struct RhiReconciliationLeaseOwner([u8; LEASE_OWNER_BYTES]);
+
+impl RhiReconciliationLeaseOwner {
+ pub fn from_bytes(bytes: [u8; LEASE_OWNER_BYTES]) -> Result<Self, RhiReconciliationJobError> {
+ if bytes.iter().all(|byte| *byte == 0) {
+ return Err(RhiReconciliationJobError::new(
+ RhiReconciliationJobErrorKind::InvalidInput,
+ ));
+ }
+ Ok(Self(bytes))
+ }
+}
+
+impl fmt::Debug for RhiReconciliationLeaseOwner {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter.write_str("RhiReconciliationLeaseOwner([redacted])")
+ }
+}
+
+/// Deterministic identity of one trade generation and evidence policy.
+#[derive(Clone, Copy, PartialEq, Eq, Hash)]
+pub struct RhiReconciliationJobId([u8; 32]);
+
+impl RhiReconciliationJobId {
+ #[must_use]
+ pub const fn as_bytes(&self) -> &[u8; 32] {
+ &self.0
+ }
+}
+
+impl fmt::Debug for RhiReconciliationJobId {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter.write_str("RhiReconciliationJobId([redacted])")
+ }
+}
+
+/// Immutable decoded view of one retained durable job row.
+#[derive(Clone, Copy, PartialEq, Eq)]
+pub struct RhiReconciliationJob {
+ id: RhiReconciliationJobId,
+ trade_id: TradeId,
+ input_generation: u64,
+ policy_digest: RhiEvidencePolicyDigest,
+ state: RhiReconciliationJobState,
+ revision: u64,
+ attempt_count: u16,
+ failure_count: u16,
+ policy: RhiReconciliationJobPolicy,
+ next_attempt: Option<RhiReconciliationUnixMilliseconds>,
+ lease_owner: Option<RhiReconciliationLeaseOwner>,
+ lease_expires: Option<RhiReconciliationUnixMilliseconds>,
+ created_at: RhiReconciliationUnixMilliseconds,
+ updated_at: RhiReconciliationUnixMilliseconds,
+}
+
+impl RhiReconciliationJob {
+ #[must_use]
+ pub const fn id(self) -> RhiReconciliationJobId {
+ self.id
+ }
+
+ #[must_use]
+ pub const fn trade_id(self) -> TradeId {
+ self.trade_id
+ }
+
+ #[must_use]
+ pub const fn input_generation(self) -> u64 {
+ self.input_generation
+ }
+
+ #[must_use]
+ pub const fn evidence_policy_digest(self) -> RhiEvidencePolicyDigest {
+ self.policy_digest
+ }
+
+ #[must_use]
+ pub const fn state(self) -> RhiReconciliationJobState {
+ self.state
+ }
+
+ #[must_use]
+ pub const fn revision(self) -> u64 {
+ self.revision
+ }
+
+ #[must_use]
+ pub const fn attempt_count(self) -> u16 {
+ self.attempt_count
+ }
+
+ #[must_use]
+ pub const fn failure_count(self) -> u16 {
+ self.failure_count
+ }
+
+ #[must_use]
+ pub const fn next_attempt(self) -> Option<RhiReconciliationUnixMilliseconds> {
+ self.next_attempt
+ }
+
+ #[must_use]
+ pub const fn created_at(self) -> RhiReconciliationUnixMilliseconds {
+ self.created_at
+ }
+
+ #[must_use]
+ pub const fn updated_at(self) -> RhiReconciliationUnixMilliseconds {
+ self.updated_at
+ }
+}
+
+impl fmt::Debug for RhiReconciliationJob {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter
+ .debug_struct("RhiReconciliationJob")
+ .field("identity", &"[redacted]")
+ .field("state", &self.state)
+ .field("revision", &self.revision)
+ .field("attempt_count", &self.attempt_count)
+ .field("failure_count", &self.failure_count)
+ .field("next_attempt", &self.next_attempt)
+ .field("lease", &self.lease_owner.map(|_| "[redacted]"))
+ .field("lease_expires", &self.lease_expires)
+ .finish()
+ }
+}
+
+/// Idempotent schedule result for one exact dirty generation.
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+pub struct RhiReconciliationScheduleOutcome {
+ job: RhiReconciliationJob,
+ created: bool,
+}
+
+impl RhiReconciliationScheduleOutcome {
+ #[must_use]
+ pub const fn job(self) -> RhiReconciliationJob {
+ self.job
+ }
+
+ #[must_use]
+ pub const fn created(self) -> bool {
+ self.created
+ }
+}
+
+/// Non-forgeable compare-and-swap lease returned by a successful claim.
+///
+/// ```compile_fail
+/// use rhi::RhiReconciliationLease;
+///
+/// let _ = RhiReconciliationLease { job: todo!(), owner: todo!(), lease_expires: todo!() };
+/// ```
+#[derive(Clone, Copy, PartialEq, Eq)]
+pub struct RhiReconciliationLease {
+ job: RhiReconciliationJob,
+ owner: RhiReconciliationLeaseOwner,
+ lease_expires: RhiReconciliationUnixMilliseconds,
+}
+
+impl RhiReconciliationLease {
+ #[must_use]
+ pub const fn job(self) -> RhiReconciliationJob {
+ self.job
+ }
+
+ #[must_use]
+ pub const fn lease_expires(self) -> RhiReconciliationUnixMilliseconds {
+ self.lease_expires
+ }
+
+ #[must_use]
+ pub fn renewal_due(self) -> RhiReconciliationUnixMilliseconds {
+ RhiReconciliationUnixMilliseconds(
+ self.lease_expires()
+ .get()
+ .saturating_sub(self.job.policy.lease_renewal_ms),
+ )
+ }
+
+ #[must_use]
+ pub fn retry_delay_upper_bound(self) -> u64 {
+ retry_delay_upper_bound(self.job.policy, self.job.failure_count)
+ }
+}
+
+impl fmt::Debug for RhiReconciliationLease {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter
+ .debug_struct("RhiReconciliationLease")
+ .field("job", &self.job)
+ .field("owner", &"[redacted]")
+ .finish()
+ }
+}
+
+impl RhiReconciliationJobRepository<'_> {
+ /// Idempotently schedules the current durable dirty generation for one trade.
+ pub async fn schedule_trade(
+ &self,
+ trade_id: TradeId,
+ policy: RhiReconciliationJobPolicy,
+ now: RhiReconciliationUnixMilliseconds,
+ ) -> Result<RhiReconciliationScheduleOutcome, RhiReconciliationJobError> {
+ require_writable(self)?;
+ self.host()
+ .sqlite_host()
+ .transaction(move |transaction| {
+ Box::pin(async move { schedule(transaction, trade_id, policy, now).await })
+ })
+ .await
+ .map_err(map_transaction_error)
+ }
+
+ /// Claims the oldest eligible job or reclaims one expired lease.
+ pub async fn claim_next(
+ &self,
+ owner: RhiReconciliationLeaseOwner,
+ now: RhiReconciliationUnixMilliseconds,
+ ) -> Result<Option<RhiReconciliationLease>, RhiReconciliationJobError> {
+ require_writable(self)?;
+ self.host()
+ .sqlite_host()
+ .transaction(move |transaction| {
+ Box::pin(async move { claim_next(transaction, owner, now).await })
+ })
+ .await
+ .map_err(map_transaction_error)
+ }
+
+ /// Renews an unexpired lease only after its configured renewal point.
+ pub async fn renew(
+ &self,
+ lease: RhiReconciliationLease,
+ now: RhiReconciliationUnixMilliseconds,
+ ) -> Result<RhiReconciliationLease, RhiReconciliationJobError> {
+ require_writable(self)?;
+ if now >= lease.lease_expires() {
+ return Err(RhiReconciliationJobError::new(
+ RhiReconciliationJobErrorKind::LeaseLost,
+ ));
+ }
+ if now < lease.renewal_due() {
+ return Err(RhiReconciliationJobError::new(
+ RhiReconciliationJobErrorKind::NotReady,
+ ));
+ }
+ self.host()
+ .sqlite_host()
+ .transaction(move |transaction| {
+ Box::pin(async move { renew(transaction, lease, now).await })
+ })
+ .await
+ .map_err(map_transaction_error)
+ }
+
+ /// Records one failed leased attempt and schedules bounded full-jitter retry.
+ pub async fn record_failure(
+ &self,
+ lease: RhiReconciliationLease,
+ now: RhiReconciliationUnixMilliseconds,
+ delay: RhiReconciliationRetryDelayMilliseconds,
+ ) -> Result<RhiReconciliationJob, RhiReconciliationJobError> {
+ require_writable(self)?;
+ if now >= lease.lease_expires() {
+ return Err(RhiReconciliationJobError::new(
+ RhiReconciliationJobErrorKind::LeaseLost,
+ ));
+ }
+ if delay.get() > lease.retry_delay_upper_bound() {
+ return Err(RhiReconciliationJobError::new(
+ RhiReconciliationJobErrorKind::InvalidInput,
+ ));
+ }
+ self.host()
+ .sqlite_host()
+ .transaction(move |transaction| {
+ Box::pin(async move { record_failure(transaction, lease, now, delay).await })
+ })
+ .await
+ .map_err(map_transaction_error)
+ }
+}
+
+fn require_writable(
+ repository: &RhiReconciliationJobRepository<'_>,
+) -> Result<(), RhiReconciliationJobError> {
+ if repository.host().mode() == RhiStateHostMode::ReadWriteExisting {
+ Ok(())
+ } else {
+ Err(RhiReconciliationJobError::new(
+ RhiReconciliationJobErrorKind::InvalidMode,
+ ))
+ }
+}
+
+async fn schedule(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ trade_id: TradeId,
+ policy: RhiReconciliationJobPolicy,
+ now: RhiReconciliationUnixMilliseconds,
+) -> Result<RhiReconciliationScheduleOutcome, OperationError> {
+ let dirty = read_dirty(transaction, trade_id).await?;
+ let id = job_id(trade_id, dirty.generation, dirty.policy);
+ if let Some(job) = read_job(transaction, id).await? {
+ if job.trade_id != trade_id
+ || job.input_generation != dirty.generation
+ || job.policy_digest != dirty.policy
+ {
+ return Err(OperationError::Storage);
+ }
+ return Ok(RhiReconciliationScheduleOutcome {
+ job,
+ created: false,
+ });
+ }
+
+ if let Some(active) = read_active_job(transaction, trade_id).await? {
+ if active.input_generation >= dirty.generation {
+ return Err(OperationError::DirtyGenerationConflict);
+ }
+ let result = sqlx::query(SUPERSEDE_JOB_SQL)
+ .bind(i64_value(now.get())?)
+ .bind(active.id.as_bytes().as_slice())
+ .bind(i64_value(active.revision)?)
+ .execute(&mut *transaction)
+ .await
+ .map_err(|_| OperationError::Storage)?;
+ if result.rows_affected() != 1 {
+ return Err(OperationError::DirtyGenerationConflict);
+ }
+ }
+
+ let active = sqlx::query(ACTIVE_JOB_COUNT_SQL)
+ .fetch_one(&mut *transaction)
+ .await
+ .map_err(|_| OperationError::Storage)?
+ .try_get::<i64, _>("active_count")
+ .ok()
+ .and_then(|value| u32::try_from(value).ok())
+ .ok_or(OperationError::Storage)?;
+ if active >= policy.queue_capacity {
+ return Err(OperationError::QueueFull);
+ }
+
+ let now = i64_value(now.get())?;
+ let result = sqlx::query(INSERT_JOB_SQL)
+ .bind(id.as_bytes().as_slice())
+ .bind(trade_id.as_bytes().as_slice())
+ .bind(i64_value(dirty.generation)?)
+ .bind(dirty.policy.as_bytes().as_slice())
+ .bind(i64::from(policy.max_attempts))
+ .bind(i64_value(policy.lease_duration_ms)?)
+ .bind(i64_value(policy.lease_renewal_ms)?)
+ .bind(i64_value(policy.initial_backoff_ms)?)
+ .bind(i64_value(policy.maximum_backoff_ms)?)
+ .bind(now)
+ .bind(now)
+ .bind(now)
+ .execute(&mut *transaction)
+ .await
+ .map_err(|_| OperationError::Storage)?;
+ if result.rows_affected() != 1 {
+ return Err(OperationError::Storage);
+ }
+ let job = read_job(transaction, id)
+ .await?
+ .ok_or(OperationError::Storage)?;
+ Ok(RhiReconciliationScheduleOutcome { job, created: true })
+}
+
+async fn claim_next(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ owner: RhiReconciliationLeaseOwner,
+ now: RhiReconciliationUnixMilliseconds,
+) -> Result<Option<RhiReconciliationLease>, OperationError> {
+ let now_value = i64_value(now.get())?;
+ sqlx::query(EXHAUST_EXPIRED_SQL)
+ .bind(now_value)
+ .bind(now_value)
+ .bind(now_value)
+ .execute(&mut *transaction)
+ .await
+ .map_err(|_| OperationError::Storage)?;
+ let Some(row) = sqlx::query(READ_CLAIMABLE_SQL)
+ .bind(now_value)
+ .bind(now_value)
+ .bind(now_value)
+ .fetch_optional(&mut *transaction)
+ .await
+ .map_err(|_| OperationError::Storage)?
+ else {
+ return Ok(None);
+ };
+ let candidate = decode_job(row)?;
+ let expires = now
+ .get()
+ .checked_add(candidate.policy.lease_duration_ms)
+ .filter(|value| *value <= MAX_UNIX_MILLISECONDS)
+ .ok_or(OperationError::InvalidInput)?;
+ let result = sqlx::query(CLAIM_JOB_SQL)
+ .bind(owner.0.as_slice())
+ .bind(i64_value(expires)?)
+ .bind(now_value)
+ .bind(candidate.id.as_bytes().as_slice())
+ .bind(i64_value(candidate.revision)?)
+ .bind(now_value)
+ .bind(now_value)
+ .bind(now_value)
+ .execute(&mut *transaction)
+ .await
+ .map_err(|_| OperationError::Storage)?;
+ if result.rows_affected() != 1 {
+ return Err(OperationError::LeaseLost);
+ }
+ let job = read_job(transaction, candidate.id)
+ .await?
+ .filter(|job| job.state == RhiReconciliationJobState::Leased)
+ .ok_or(OperationError::Storage)?;
+ let lease_expires = job.lease_expires.ok_or(OperationError::Storage)?;
+ Ok(Some(RhiReconciliationLease {
+ job,
+ owner,
+ lease_expires,
+ }))
+}
+
+async fn renew(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ lease: RhiReconciliationLease,
+ now: RhiReconciliationUnixMilliseconds,
+) -> Result<RhiReconciliationLease, OperationError> {
+ let previous_expiry = lease.lease_expires().get();
+ let expires = now
+ .get()
+ .checked_add(lease.job.policy.lease_duration_ms)
+ .filter(|value| *value > previous_expiry && *value <= MAX_UNIX_MILLISECONDS)
+ .ok_or(OperationError::NotReady)?;
+ let result = sqlx::query(RENEW_JOB_SQL)
+ .bind(i64_value(expires)?)
+ .bind(i64_value(now.get())?)
+ .bind(lease.job.id.as_bytes().as_slice())
+ .bind(i64_value(lease.job.revision)?)
+ .bind(lease.owner.0.as_slice())
+ .bind(i64_value(previous_expiry)?)
+ .bind(i64_value(now.get())?)
+ .bind(i64_value(now.get())?)
+ .execute(&mut *transaction)
+ .await
+ .map_err(|_| OperationError::Storage)?;
+ if result.rows_affected() != 1 {
+ return Err(OperationError::LeaseLost);
+ }
+ let job = read_job(transaction, lease.job.id)
+ .await?
+ .filter(|job| job.state == RhiReconciliationJobState::Leased)
+ .ok_or(OperationError::Storage)?;
+ Ok(RhiReconciliationLease {
+ job,
+ owner: lease.owner,
+ lease_expires: job.lease_expires.ok_or(OperationError::Storage)?,
+ })
+}
+
+async fn record_failure(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ lease: RhiReconciliationLease,
+ now: RhiReconciliationUnixMilliseconds,
+ delay: RhiReconciliationRetryDelayMilliseconds,
+) -> Result<RhiReconciliationJob, OperationError> {
+ let now_value = i64_value(now.get())?;
+ let query = if lease.job.attempt_count >= lease.job.policy.max_attempts {
+ sqlx::query(EXHAUST_JOB_SQL)
+ .bind(now_value)
+ .bind(lease.job.id.as_bytes().as_slice())
+ .bind(i64_value(lease.job.revision)?)
+ .bind(lease.owner.0.as_slice())
+ .bind(i64_value(lease.lease_expires().get())?)
+ .bind(now_value)
+ } else {
+ let next = now
+ .get()
+ .checked_add(delay.get())
+ .filter(|value| *value <= MAX_UNIX_MILLISECONDS)
+ .ok_or(OperationError::InvalidInput)?;
+ sqlx::query(RETRY_JOB_SQL)
+ .bind(i64_value(next)?)
+ .bind(now_value)
+ .bind(lease.job.id.as_bytes().as_slice())
+ .bind(i64_value(lease.job.revision)?)
+ .bind(lease.owner.0.as_slice())
+ .bind(i64_value(lease.lease_expires().get())?)
+ .bind(now_value)
+ };
+ let result = query
+ .execute(&mut *transaction)
+ .await
+ .map_err(|_| OperationError::Storage)?;
+ if result.rows_affected() != 1 {
+ return Err(OperationError::LeaseLost);
+ }
+ read_job(transaction, lease.job.id)
+ .await?
+ .ok_or(OperationError::Storage)
+}
+
+async fn read_dirty(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ trade_id: TradeId,
+) -> Result<DirtyGeneration, OperationError> {
+ let row = sqlx::query(READ_DIRTY_SQL)
+ .bind(trade_id.as_bytes().as_slice())
+ .fetch_optional(&mut *transaction)
+ .await
+ .map_err(|_| OperationError::Storage)?
+ .ok_or(OperationError::DirtyGenerationConflict)?;
+ Ok(DirtyGeneration {
+ generation: positive_u64(&row, "generation")?,
+ policy: RhiEvidencePolicyDigest::from_bytes(exact_bytes::<32>(
+ &row,
+ "evidence_policy_sha256",
+ "evidence_policy_bytes",
+ )?),
+ })
+}
+
+async fn read_job(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ id: RhiReconciliationJobId,
+) -> Result<Option<RhiReconciliationJob>, OperationError> {
+ sqlx::query(READ_JOB_SQL)
+ .bind(id.as_bytes().as_slice())
+ .fetch_optional(&mut *transaction)
+ .await
+ .map_err(|_| OperationError::Storage)?
+ .map(decode_job)
+ .transpose()
+}
+
+async fn read_active_job(
+ transaction: &mut ServiceSqliteTransaction<'_>,
+ trade_id: TradeId,
+) -> Result<Option<RhiReconciliationJob>, OperationError> {
+ let rows = sqlx::query(READ_ACTIVE_JOB_SQL)
+ .bind(trade_id.as_bytes().as_slice())
+ .fetch_all(&mut *transaction)
+ .await
+ .map_err(|_| OperationError::Storage)?;
+ match rows.len() {
+ 0 => Ok(None),
+ 1 => Ok(Some(decode_job(
+ rows.into_iter().next().ok_or(OperationError::Storage)?,
+ )?)),
+ _ => Err(OperationError::Storage),
+ }
+}
+
+fn decode_job(row: sqlx::sqlite::SqliteRow) -> Result<RhiReconciliationJob, OperationError> {
+ let id = RhiReconciliationJobId(exact_bytes::<32>(&row, "job_id", "job_id_bytes")?);
+ let trade_id = TradeId::from_bytes(exact_bytes::<16>(&row, "trade_id", "trade_id_bytes")?);
+ let input_generation = positive_u64(&row, "input_generation")?;
+ let policy_digest = RhiEvidencePolicyDigest::from_bytes(exact_bytes::<32>(
+ &row,
+ "evidence_policy_sha256",
+ "evidence_policy_bytes",
+ )?);
+ let state = match bounded_text(&row, "state", "state_bytes", 10)?.as_str() {
+ "ready" => RhiReconciliationJobState::Ready,
+ "leased" => RhiReconciliationJobState::Leased,
+ "exhausted" => RhiReconciliationJobState::Exhausted,
+ "superseded" => RhiReconciliationJobState::Superseded,
+ "completed" => RhiReconciliationJobState::Completed,
+ _ => return Err(OperationError::Storage),
+ };
+ let revision = positive_u64(&row, "revision")?;
+ let attempt_count = bounded_u16(&row, "attempt_count", 0, MAX_ATTEMPTS)?;
+ let failure_count = bounded_u16(&row, "failure_count", 0, attempt_count)?;
+ let max_attempts = bounded_u16(&row, "max_attempts", 1, MAX_ATTEMPTS)?;
+ let lease_duration_ms = bounded_u64(&row, "lease_duration_ms", 1_000, MAX_LEASE_MILLISECONDS)?;
+ let lease_renewal_ms = bounded_u64(&row, "lease_renewal_ms", 100, MAX_RENEWAL_MILLISECONDS)?;
+ let initial_backoff_ms = bounded_u64(
+ &row,
+ "initial_backoff_ms",
+ 1,
+ MAX_INITIAL_BACKOFF_MILLISECONDS,
+ )?;
+ let maximum_backoff_ms = bounded_u64(&row, "maximum_backoff_ms", 1, MAX_BACKOFF_MILLISECONDS)?;
+ let policy = RhiReconciliationJobPolicy::new(
+ RHI_RECONCILIATION_JOB_MAX_ACTIVE,
+ lease_duration_ms,
+ lease_renewal_ms,
+ max_attempts,
+ initial_backoff_ms,
+ maximum_backoff_ms,
+ )
+ .map_err(|_| OperationError::Storage)?;
+ let next_attempt = optional_millis(&row, "next_attempt_unix_ms")?;
+ let lease_expires = optional_millis(&row, "lease_expires_unix_ms")?;
+ let lease_owner = optional_owner(&row)?;
+ let created_at =
+ RhiReconciliationUnixMilliseconds(nonnegative_u64(&row, "created_at_unix_ms")?);
+ let updated_at =
+ RhiReconciliationUnixMilliseconds(nonnegative_u64(&row, "updated_at_unix_ms")?);
+ if updated_at < created_at
+ || failure_count > attempt_count
+ || job_id(trade_id, input_generation, policy_digest) != id
+ || !valid_state_fields(
+ state,
+ attempt_count,
+ max_attempts,
+ next_attempt,
+ lease_owner,
+ lease_expires,
+ )
+ {
+ return Err(OperationError::Storage);
+ }
+ Ok(RhiReconciliationJob {
+ id,
+ trade_id,
+ input_generation,
+ policy_digest,
+ state,
+ revision,
+ attempt_count,
+ failure_count,
+ policy,
+ next_attempt,
+ lease_owner,
+ lease_expires,
+ created_at,
+ updated_at,
+ })
+}
+
+fn valid_state_fields(
+ state: RhiReconciliationJobState,
+ attempt_count: u16,
+ max_attempts: u16,
+ next_attempt: Option<RhiReconciliationUnixMilliseconds>,
+ owner: Option<RhiReconciliationLeaseOwner>,
+ expires: Option<RhiReconciliationUnixMilliseconds>,
+) -> bool {
+ match state {
+ RhiReconciliationJobState::Ready => {
+ attempt_count < max_attempts
+ && next_attempt.is_some()
+ && owner.is_none()
+ && expires.is_none()
+ }
+ RhiReconciliationJobState::Leased => {
+ attempt_count > 0
+ && attempt_count <= max_attempts
+ && next_attempt.is_none()
+ && owner.is_some()
+ && expires.is_some()
+ }
+ RhiReconciliationJobState::Exhausted
+ | RhiReconciliationJobState::Superseded
+ | RhiReconciliationJobState::Completed => {
+ next_attempt.is_none() && owner.is_none() && expires.is_none()
+ }
+ }
+}
+
+fn job_id(
+ trade_id: TradeId,
+ generation: u64,
+ policy: RhiEvidencePolicyDigest,
+) -> RhiReconciliationJobId {
+ let mut hasher = Sha256::new();
+ hasher.update(JOB_ID_DOMAIN);
+ hasher.update(trade_id.as_bytes());
+ hasher.update(generation.to_be_bytes());
+ hasher.update(policy.as_bytes());
+ RhiReconciliationJobId(hasher.finalize().into())
+}
+
+fn retry_delay_upper_bound(policy: RhiReconciliationJobPolicy, failure_count: u16) -> u64 {
+ let exponent = u32::from(failure_count.min(63));
+ policy
+ .initial_backoff_ms
+ .saturating_mul(1_u64.checked_shl(exponent).unwrap_or(u64::MAX))
+ .min(policy.maximum_backoff_ms)
+}
+
+fn exact_bytes<const N: usize>(
+ row: &sqlx::sqlite::SqliteRow,
+ field: &str,
+ length_field: &str,
+) -> Result<[u8; N], OperationError> {
+ if row.try_get::<i64, _>(length_field).ok() != i64::try_from(N).ok() {
+ return Err(OperationError::Storage);
+ }
+ row.try_get::<Vec<u8>, _>(field)
+ .map_err(|_| OperationError::Storage)?
+ .try_into()
+ .map_err(|_| OperationError::Storage)
+}
+
+fn optional_owner(
+ row: &sqlx::sqlite::SqliteRow,
+) -> Result<Option<RhiReconciliationLeaseOwner>, OperationError> {
+ let length = row
+ .try_get::<Option<i64>, _>("lease_owner_bytes")
+ .map_err(|_| OperationError::Storage)?;
+ match length {
+ None => Ok(None),
+ Some(value) if value == LEASE_OWNER_BYTES as i64 => {
+ let bytes: [u8; LEASE_OWNER_BYTES] = row
+ .try_get::<Vec<u8>, _>("lease_owner")
+ .map_err(|_| OperationError::Storage)?
+ .try_into()
+ .map_err(|_| OperationError::Storage)?;
+ RhiReconciliationLeaseOwner::from_bytes(bytes)
+ .map(Some)
+ .map_err(|_| OperationError::Storage)
+ }
+ Some(_) => Err(OperationError::Storage),
+ }
+}
+
+fn bounded_text(
+ row: &sqlx::sqlite::SqliteRow,
+ field: &str,
+ length_field: &str,
+ maximum: usize,
+) -> Result<String, OperationError> {
+ let length = row
+ .try_get::<i64, _>(length_field)
+ .ok()
+ .and_then(|value| usize::try_from(value).ok())
+ .filter(|value| *value > 0 && *value <= maximum)
+ .ok_or(OperationError::Storage)?;
+ let value = row
+ .try_get::<String, _>(field)
+ .map_err(|_| OperationError::Storage)?;
+ if value.len() == length {
+ Ok(value)
+ } else {
+ Err(OperationError::Storage)
+ }
+}
+
+fn bounded_u16(
+ row: &sqlx::sqlite::SqliteRow,
+ field: &str,
+ minimum: u16,
+ maximum: u16,
+) -> Result<u16, OperationError> {
+ row.try_get::<i64, _>(field)
+ .ok()
+ .and_then(|value| u16::try_from(value).ok())
+ .filter(|value| (*value >= minimum) && (*value <= maximum))
+ .ok_or(OperationError::Storage)
+}
+
+fn bounded_u64(
+ row: &sqlx::sqlite::SqliteRow,
+ field: &str,
+ minimum: u64,
+ maximum: u64,
+) -> Result<u64, OperationError> {
+ row.try_get::<i64, _>(field)
+ .ok()
+ .and_then(|value| u64::try_from(value).ok())
+ .filter(|value| (*value >= minimum) && (*value <= maximum))
+ .ok_or(OperationError::Storage)
+}
+
+fn positive_u64(row: &sqlx::sqlite::SqliteRow, field: &str) -> Result<u64, OperationError> {
+ bounded_u64(row, field, 1, MAX_UNIX_MILLISECONDS)
+}
+
+fn nonnegative_u64(row: &sqlx::sqlite::SqliteRow, field: &str) -> Result<u64, OperationError> {
+ bounded_u64(row, field, 0, MAX_UNIX_MILLISECONDS)
+}
+
+fn optional_millis(
+ row: &sqlx::sqlite::SqliteRow,
+ field: &str,
+) -> Result<Option<RhiReconciliationUnixMilliseconds>, OperationError> {
+ row.try_get::<Option<i64>, _>(field)
+ .map_err(|_| OperationError::Storage)?
+ .map(|value| {
+ u64::try_from(value)
+ .map(RhiReconciliationUnixMilliseconds)
+ .map_err(|_| OperationError::Storage)
+ })
+ .transpose()
+}
+
+fn i64_value(value: u64) -> Result<i64, OperationError> {
+ i64::try_from(value).map_err(|_| OperationError::InvalidInput)
+}
+
+#[derive(Clone, Copy)]
+struct DirtyGeneration {
+ generation: u64,
+ policy: RhiEvidencePolicyDigest,
+}
+
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+enum OperationError {
+ InvalidInput,
+ QueueFull,
+ DirtyGenerationConflict,
+ LeaseLost,
+ NotReady,
+ Storage,
+}
+
+fn map_transaction_error(
+ error: ServiceSqliteTransactionError<OperationError>,
+) -> RhiReconciliationJobError {
+ if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown {
+ return RhiReconciliationJobError::new(RhiReconciliationJobErrorKind::CommitOutcomeUnknown);
+ }
+ RhiReconciliationJobError::new(match error.operation_error().copied() {
+ Some(OperationError::InvalidInput) => RhiReconciliationJobErrorKind::InvalidInput,
+ Some(OperationError::QueueFull) => RhiReconciliationJobErrorKind::QueueFull,
+ Some(OperationError::DirtyGenerationConflict) => {
+ RhiReconciliationJobErrorKind::DirtyGenerationConflict
+ }
+ Some(OperationError::LeaseLost) => RhiReconciliationJobErrorKind::LeaseLost,
+ Some(OperationError::NotReady) => RhiReconciliationJobErrorKind::NotReady,
+ Some(OperationError::Storage) | None => RhiReconciliationJobErrorKind::Storage,
+ })
+}
diff --git a/src/state_catalog.rs b/src/state_catalog.rs
@@ -12,7 +12,7 @@ use radroots_service_sqlite::{
pub const RHI_STATE_BASE_SCHEMA_VERSION: u32 = 1;
/// The newest governed RHI state schema understood by this binary.
-pub const RHI_STATE_SCHEMA_VERSION: u32 = 4;
+pub const RHI_STATE_SCHEMA_VERSION: u32 = 5;
/// The shared metadata and migration-ledger objects present at schema v1.
pub const RHI_STATE_SCHEMA_VERSION_1_OBJECT_COUNT: u32 = 6;
@@ -26,10 +26,13 @@ pub const RHI_STATE_SCHEMA_VERSION_3_OBJECT_COUNT: u32 = 22;
/// The shared objects, immutable trade evidence, source cursors, and dirty generations.
pub const RHI_STATE_SCHEMA_VERSION_4_OBJECT_COUNT: u32 = 28;
+/// The shared objects, canonical evidence, cursors, generations, and durable jobs.
+pub const RHI_STATE_SCHEMA_VERSION_5_OBJECT_COUNT: u32 = 33;
+
/// SHA-256 identity of the ordered migration catalog rooted at schema v1.
pub const RHI_MIGRATION_CATALOG_SHA256: [u8; 32] = [
- 0x29, 0x5f, 0xf5, 0xe9, 0xac, 0x0e, 0xad, 0xc1, 0xf2, 0xc2, 0xd8, 0xc2, 0x58, 0xc2, 0xea, 0xa8,
- 0x16, 0xf9, 0x46, 0x56, 0x0b, 0x68, 0x9d, 0xcb, 0x43, 0x2f, 0xdc, 0x99, 0xe7, 0x9e, 0xc3, 0x53,
+ 0x4e, 0xd3, 0x2a, 0x5a, 0x71, 0xf5, 0xbf, 0x45, 0x42, 0x62, 0xa9, 0x08, 0x2a, 0xd7, 0x4d, 0x0d,
+ 0x6a, 0x13, 0xa5, 0x56, 0x05, 0x46, 0x59, 0xac, 0x7b, 0x9d, 0x87, 0x72, 0x59, 0xb8, 0x16, 0xd8,
];
/// SHA-256 identity of the exact schema-v1 object snapshot.
@@ -74,10 +77,22 @@ pub const RHI_STATE_SCHEMA_VERSION_4_SHA256: [u8; 32] = [
0x6c, 0xdd, 0xe0, 0x7c, 0x40, 0x42, 0x9d, 0xff, 0xe9, 0x1c, 0xb8, 0x3c, 0x2a, 0x42, 0x5c, 0xea,
];
+/// SHA-256 identity of the schema-v5 reconciliation-job migration.
+pub const RHI_STATE_SCHEMA_VERSION_5_MIGRATION_SHA256: [u8; 32] = [
+ 0x82, 0x75, 0xff, 0xb3, 0xd5, 0xc9, 0xfa, 0x0c, 0x76, 0x48, 0xf8, 0x9e, 0x8b, 0xfb, 0x92, 0x87,
+ 0x7b, 0x0d, 0xcd, 0x5e, 0xf3, 0xbf, 0x0b, 0x52, 0x54, 0xb7, 0xff, 0xd2, 0x58, 0x0b, 0xd4, 0x5d,
+];
+
+/// SHA-256 identity of the exact schema-v5 object snapshot.
+pub const RHI_STATE_SCHEMA_VERSION_5_SHA256: [u8; 32] = [
+ 0xee, 0xf6, 0x7b, 0x48, 0x40, 0x0d, 0xee, 0x23, 0x4f, 0x6c, 0x36, 0xd4, 0x22, 0x55, 0x1c, 0x7e,
+ 0x99, 0x83, 0x15, 0x0b, 0x88, 0x7f, 0x26, 0x42, 0x50, 0xf9, 0xc8, 0x7a, 0x1a, 0x31, 0x3e, 0x60,
+];
+
/// SHA-256 identity of the schema catalog bound to the migration catalog.
pub const RHI_STATE_SCHEMA_CATALOG_SHA256: [u8; 32] = [
- 0x58, 0x78, 0x9c, 0x1d, 0x39, 0x65, 0x14, 0xce, 0x19, 0xa5, 0x60, 0xb6, 0x7b, 0xf7, 0x60, 0x65,
- 0x0b, 0x9d, 0x6f, 0xa3, 0x78, 0xbc, 0x82, 0x34, 0xc1, 0xf6, 0xa3, 0x9a, 0xd2, 0xf5, 0x56, 0xe6,
+ 0x0d, 0x2c, 0x59, 0x5a, 0x3c, 0xa3, 0xfb, 0xc9, 0xb3, 0xbf, 0x6b, 0xfc, 0x61, 0x68, 0xa1, 0x94,
+ 0xac, 0x78, 0x8d, 0x49, 0x31, 0xc4, 0x62, 0x00, 0xe2, 0xf1, 0x37, 0x9e, 0xbb, 0x52, 0x3e, 0xce,
];
macro_rules! rhi_config_bindings_table_sql {
@@ -472,6 +487,147 @@ const CREATE_SOURCE_CHECKPOINT_MIGRATION_SQL: &str = concat!(
";",
);
+macro_rules! reconciliation_jobs_table_sql {
+ () => {
+ r#"CREATE TABLE reconciliation_jobs (
+ job_id BLOB NOT NULL PRIMARY KEY CHECK (length(job_id) = 32),
+ trade_id BLOB NOT NULL CHECK (length(trade_id) = 16),
+ input_generation INTEGER NOT NULL
+ CHECK (input_generation BETWEEN 1 AND 9223372036854775807),
+ evidence_policy_sha256 BLOB NOT NULL CHECK (length(evidence_policy_sha256) = 32),
+ state TEXT NOT NULL
+ CHECK (state IN ('ready', 'leased', 'exhausted', 'superseded', 'completed')),
+ revision INTEGER NOT NULL CHECK (revision BETWEEN 1 AND 9223372036854775807),
+ attempt_count INTEGER NOT NULL CHECK (attempt_count BETWEEN 0 AND 100),
+ failure_count INTEGER NOT NULL CHECK (failure_count BETWEEN 0 AND attempt_count),
+ max_attempts INTEGER NOT NULL CHECK (max_attempts BETWEEN 1 AND 100),
+ lease_duration_ms INTEGER NOT NULL CHECK (lease_duration_ms BETWEEN 1000 AND 300000),
+ lease_renewal_ms INTEGER NOT NULL CHECK (lease_renewal_ms BETWEEN 100 AND 150000),
+ initial_backoff_ms INTEGER NOT NULL CHECK (initial_backoff_ms BETWEEN 1 AND 60000),
+ maximum_backoff_ms INTEGER NOT NULL CHECK (maximum_backoff_ms BETWEEN 1 AND 3600000),
+ next_attempt_unix_ms INTEGER,
+ lease_owner BLOB,
+ lease_expires_unix_ms INTEGER,
+ created_at_unix_ms INTEGER NOT NULL
+ CHECK (created_at_unix_ms BETWEEN 0 AND 9223372036854775807),
+ updated_at_unix_ms INTEGER NOT NULL
+ CHECK (updated_at_unix_ms BETWEEN created_at_unix_ms AND 9223372036854775807),
+ FOREIGN KEY (trade_id) REFERENCES trade_dirty_generations (trade_id),
+ CHECK (lease_renewal_ms < lease_duration_ms),
+ CHECK (initial_backoff_ms <= maximum_backoff_ms),
+ CHECK (
+ (state = 'ready'
+ AND attempt_count < max_attempts
+ AND next_attempt_unix_ms BETWEEN 0 AND 9223372036854775807
+ AND lease_owner IS NULL
+ AND lease_expires_unix_ms IS NULL)
+ OR (state = 'leased'
+ AND attempt_count BETWEEN 1 AND max_attempts
+ AND next_attempt_unix_ms IS NULL
+ AND length(lease_owner) = 16
+ AND lease_expires_unix_ms BETWEEN 1 AND 9223372036854775807)
+ OR (state IN ('exhausted', 'superseded', 'completed')
+ AND next_attempt_unix_ms IS NULL
+ AND lease_owner IS NULL
+ AND lease_expires_unix_ms IS NULL)
+ )
+) STRICT"#
+ };
+}
+
+macro_rules! reconciliation_jobs_one_active_sql {
+ () => {
+ r#"CREATE UNIQUE INDEX reconciliation_jobs_one_active_per_trade
+ON reconciliation_jobs (trade_id)
+WHERE state IN ('ready', 'leased')"#
+ };
+}
+
+macro_rules! reconciliation_jobs_schedule_sql {
+ () => {
+ r#"CREATE INDEX reconciliation_jobs_by_schedule
+ON reconciliation_jobs (
+ state, next_attempt_unix_ms, lease_expires_unix_ms,
+ created_at_unix_ms, job_id
+)"#
+ };
+}
+
+macro_rules! reconciliation_jobs_guard_update_sql {
+ () => {
+ r#"CREATE TRIGGER reconciliation_jobs_guard_update
+BEFORE UPDATE ON reconciliation_jobs
+WHEN NEW.job_id != OLD.job_id
+ OR NEW.trade_id != OLD.trade_id
+ OR NEW.input_generation != OLD.input_generation
+ OR NEW.evidence_policy_sha256 != OLD.evidence_policy_sha256
+ OR NEW.max_attempts != OLD.max_attempts
+ OR NEW.lease_duration_ms != OLD.lease_duration_ms
+ OR NEW.lease_renewal_ms != OLD.lease_renewal_ms
+ OR NEW.initial_backoff_ms != OLD.initial_backoff_ms
+ OR NEW.maximum_backoff_ms != OLD.maximum_backoff_ms
+ OR NEW.created_at_unix_ms != OLD.created_at_unix_ms
+ OR NEW.revision != OLD.revision + 1
+ OR NEW.updated_at_unix_ms < OLD.updated_at_unix_ms
+ OR NOT (
+ (OLD.state = 'ready' AND NEW.state IN ('leased', 'superseded'))
+ OR (OLD.state = 'leased'
+ AND NEW.state IN ('leased', 'ready', 'exhausted', 'superseded', 'completed'))
+ OR (OLD.state = 'exhausted' AND NEW.state = 'superseded')
+ )
+BEGIN
+ SELECT RAISE(ABORT, 'reconciliation job transition is invalid');
+END"#
+ };
+}
+
+pub(crate) const CREATE_RECONCILIATION_JOBS_TABLE_SQL: &str = reconciliation_jobs_table_sql!();
+const CREATE_RECONCILIATION_JOBS_ONE_ACTIVE_SQL: &str = reconciliation_jobs_one_active_sql!();
+const CREATE_RECONCILIATION_JOBS_SCHEDULE_SQL: &str = reconciliation_jobs_schedule_sql!();
+const CREATE_RECONCILIATION_JOBS_GUARD_UPDATE_SQL: &str = reconciliation_jobs_guard_update_sql!();
+const CREATE_RECONCILIATION_JOBS_NO_DELETE_SQL: &str = immutable_no_delete_sql!(
+ "reconciliation_jobs_no_delete",
+ "reconciliation_jobs",
+ "reconciliation jobs are retained"
+);
+const CREATE_RECONCILIATION_JOBS_MIGRATION_SQL: &str = concat!(
+ reconciliation_jobs_table_sql!(),
+ ";\n",
+ reconciliation_jobs_one_active_sql!(),
+ ";\n",
+ reconciliation_jobs_schedule_sql!(),
+ ";\n",
+ reconciliation_jobs_guard_update_sql!(),
+ ";\n",
+ immutable_no_delete_sql!(
+ "reconciliation_jobs_no_delete",
+ "reconciliation_jobs",
+ "reconciliation jobs are retained"
+ ),
+ ";",
+);
+
+const RECONCILIATION_JOBS_TABLE_SHA256: [u8; 32] = [
+ 0xe7, 0x0b, 0x5c, 0xb7, 0x26, 0x91, 0x9d, 0x02, 0xef, 0xb3, 0xa6, 0x21, 0x58, 0x48, 0xce, 0x92,
+ 0x30, 0x88, 0x17, 0x2b, 0x3f, 0xe0, 0xc1, 0x40, 0xee, 0x4a, 0xc7, 0x1b, 0xd0, 0x3c, 0x66, 0xa7,
+];
+const RECONCILIATION_JOBS_ONE_ACTIVE_SHA256: [u8; 32] = [
+ 0xf9, 0xbe, 0x78, 0xc6, 0x46, 0x7a, 0x76, 0x2d, 0x38, 0xf4, 0xe0, 0xb4, 0x80, 0xd4, 0x1b, 0x37,
+ 0x6f, 0x66, 0x8c, 0x21, 0xd2, 0x96, 0xfd, 0x72, 0x12, 0x11, 0x9d, 0x2b, 0x9e, 0x3a, 0x14, 0x3f,
+];
+const RECONCILIATION_JOBS_SCHEDULE_SHA256: [u8; 32] = [
+ 0x31, 0xd9, 0xf7, 0x38, 0x3b, 0x87, 0x8d, 0xd8, 0x00, 0xda, 0xe6, 0x1f, 0x40, 0x54, 0xd8, 0xd2,
+ 0x99, 0xca, 0x93, 0xef, 0x38, 0xbe, 0x5b, 0x63, 0x63, 0x0e, 0x1c, 0x3f, 0x57, 0x2a, 0x1e, 0x75,
+];
+const RECONCILIATION_JOBS_GUARD_UPDATE_SHA256: [u8; 32] = [
+ 0xe4, 0x6a, 0x3a, 0xf2, 0x7e, 0xe1, 0xce, 0xe2, 0x4e, 0x29, 0x13, 0x56, 0x51, 0x6c, 0x72, 0x6d,
+ 0xf0, 0xd1, 0x11, 0x29, 0x01, 0x22, 0x97, 0x5e, 0xe1, 0x79, 0x76, 0xe3, 0x6a, 0x74, 0x13, 0x1b,
+];
+const RECONCILIATION_JOBS_NO_DELETE_SHA256: [u8; 32] = [
+ 0x1a, 0x86, 0xb1, 0x70, 0x05, 0x05, 0x4c, 0x27, 0x1b, 0xa4, 0x9a, 0xf2, 0x1a, 0x44, 0x9c, 0x9d,
+ 0xdc, 0xfb, 0xbc, 0xd9, 0x13, 0xfa, 0x2a, 0x8e, 0x64, 0x6d, 0x5a, 0x6d, 0x22, 0x69, 0x59, 0xa2,
+];
+
const RHI_CONFIG_BINDINGS_TABLE_SHA256: [u8; 32] = [
0x4d, 0x6e, 0x8f, 0xff, 0xda, 0x43, 0xe6, 0xf5, 0x3e, 0x23, 0x77, 0xd2, 0x77, 0xa4, 0x52, 0x9e,
0x63, 0x3e, 0xaf, 0xb6, 0xea, 0xa2, 0xad, 0xd7, 0x56, 0xde, 0x0d, 0xc9, 0x24, 0xc5, 0x77, 0xeb,
@@ -653,10 +809,22 @@ pub fn rhi_migration_catalog() -> Result<MigrationCatalog, RhiStateCatalogError>
MigrationChecksum::from_bytes(RHI_STATE_SCHEMA_VERSION_4_MIGRATION_SHA256),
)
.map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::MigrationCatalog))?;
- let catalog = MigrationCatalog::new([configuration, trade_evidence, source_checkpoints])
- .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::MigrationCatalog))?;
+ let reconciliation_jobs = MigrationDescriptor::sql(
+ 5,
+ "create_reconciliation_jobs",
+ CREATE_RECONCILIATION_JOBS_MIGRATION_SQL,
+ MigrationChecksum::from_bytes(RHI_STATE_SCHEMA_VERSION_5_MIGRATION_SHA256),
+ )
+ .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::MigrationCatalog))?;
+ let catalog = MigrationCatalog::new([
+ configuration,
+ trade_evidence,
+ source_checkpoints,
+ reconciliation_jobs,
+ ])
+ .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::MigrationCatalog))?;
if catalog.current_version() != RHI_STATE_SCHEMA_VERSION
- || catalog.descriptors().len() != 3
+ || catalog.descriptors().len() != 4
|| catalog.digest().as_bytes() != &RHI_MIGRATION_CATALOG_SHA256
{
return Err(RhiStateCatalogError::new(
@@ -688,14 +856,26 @@ pub fn rhi_schema_catalog() -> Result<SchemaCatalog, RhiStateCatalogError> {
)
.map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog))?;
let version_four = SchemaVersionCatalog::new(
- RHI_STATE_SCHEMA_VERSION,
+ 4,
rhi_schema_version_four_objects()?,
SchemaDigest::from_bytes(RHI_STATE_SCHEMA_VERSION_4_SHA256),
)
.map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog))?;
+ let version_five = SchemaVersionCatalog::new(
+ RHI_STATE_SCHEMA_VERSION,
+ rhi_schema_version_five_objects()?,
+ SchemaDigest::from_bytes(RHI_STATE_SCHEMA_VERSION_5_SHA256),
+ )
+ .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog))?;
let catalog = SchemaCatalog::new(
&migrations,
- [version_one, version_two, version_three, version_four],
+ [
+ version_one,
+ version_two,
+ version_three,
+ version_four,
+ version_five,
+ ],
)
.map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog))?;
validate_rhi_state_catalogs(&migrations, &catalog)?;
@@ -709,10 +889,10 @@ pub fn validate_rhi_state_catalogs(
) -> Result<(), RhiStateCatalogError> {
let versions = schema.versions();
let valid = migrations.current_version() == RHI_STATE_SCHEMA_VERSION
- && migrations.descriptors().len() == 3
+ && migrations.descriptors().len() == 4
&& migrations.digest().as_bytes() == &RHI_MIGRATION_CATALOG_SHA256
&& schema.migration_catalog_digest() == migrations.digest()
- && versions.len() == 4
+ && versions.len() == 5
&& versions[0].version() == RHI_STATE_BASE_SCHEMA_VERSION
&& versions[0].object_count() == RHI_STATE_SCHEMA_VERSION_1_OBJECT_COUNT
&& versions[0].digest().as_bytes() == &RHI_STATE_SCHEMA_VERSION_1_SHA256
@@ -722,9 +902,12 @@ pub fn validate_rhi_state_catalogs(
&& versions[2].version() == 3
&& versions[2].object_count() == RHI_STATE_SCHEMA_VERSION_3_OBJECT_COUNT
&& versions[2].digest().as_bytes() == &RHI_STATE_SCHEMA_VERSION_3_SHA256
- && versions[3].version() == RHI_STATE_SCHEMA_VERSION
+ && versions[3].version() == 4
&& versions[3].object_count() == RHI_STATE_SCHEMA_VERSION_4_OBJECT_COUNT
&& versions[3].digest().as_bytes() == &RHI_STATE_SCHEMA_VERSION_4_SHA256
+ && versions[4].version() == RHI_STATE_SCHEMA_VERSION
+ && versions[4].object_count() == RHI_STATE_SCHEMA_VERSION_5_OBJECT_COUNT
+ && versions[4].digest().as_bytes() == &RHI_STATE_SCHEMA_VERSION_5_SHA256
&& schema.digest().as_bytes() == &RHI_STATE_SCHEMA_CATALOG_SHA256;
if valid {
Ok(())
@@ -784,6 +967,62 @@ fn rhi_schema_version_four_objects() -> Result<Vec<SchemaObject>, RhiStateCatalo
Ok(objects)
}
+fn rhi_schema_version_five_objects() -> Result<Vec<SchemaObject>, RhiStateCatalogError> {
+ let mut objects = rhi_schema_version_four_objects()?;
+ objects.extend(rhi_reconciliation_job_objects()?);
+ Ok(objects)
+}
+
+fn rhi_reconciliation_job_objects() -> Result<[SchemaObject; 5], RhiStateCatalogError> {
+ let object = |kind, name, table_name, sql, digest| {
+ SchemaObject::new(
+ kind,
+ name,
+ table_name,
+ sql,
+ SchemaDigest::from_bytes(digest),
+ )
+ .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog))
+ };
+ Ok([
+ object(
+ SchemaObjectKind::Table,
+ "reconciliation_jobs",
+ "reconciliation_jobs",
+ CREATE_RECONCILIATION_JOBS_TABLE_SQL,
+ RECONCILIATION_JOBS_TABLE_SHA256,
+ )?,
+ object(
+ SchemaObjectKind::Index,
+ "reconciliation_jobs_one_active_per_trade",
+ "reconciliation_jobs",
+ CREATE_RECONCILIATION_JOBS_ONE_ACTIVE_SQL,
+ RECONCILIATION_JOBS_ONE_ACTIVE_SHA256,
+ )?,
+ object(
+ SchemaObjectKind::Index,
+ "reconciliation_jobs_by_schedule",
+ "reconciliation_jobs",
+ CREATE_RECONCILIATION_JOBS_SCHEDULE_SQL,
+ RECONCILIATION_JOBS_SCHEDULE_SHA256,
+ )?,
+ object(
+ SchemaObjectKind::Trigger,
+ "reconciliation_jobs_guard_update",
+ "reconciliation_jobs",
+ CREATE_RECONCILIATION_JOBS_GUARD_UPDATE_SQL,
+ RECONCILIATION_JOBS_GUARD_UPDATE_SHA256,
+ )?,
+ object(
+ SchemaObjectKind::Trigger,
+ "reconciliation_jobs_no_delete",
+ "reconciliation_jobs",
+ CREATE_RECONCILIATION_JOBS_NO_DELETE_SQL,
+ RECONCILIATION_JOBS_NO_DELETE_SHA256,
+ )?,
+ ])
+}
+
fn rhi_source_checkpoint_objects() -> Result<[SchemaObject; 6], RhiStateCatalogError> {
let object = |kind, name, table_name, sql, digest| {
SchemaObject::new(
diff --git a/src/state_repository.rs b/src/state_repository.rs
@@ -461,6 +461,12 @@ repository_handle!(RhiPublicationTargetRepository, PublicationTarget);
repository_handle!(RhiPublicationAttemptRepository, PublicationAttempt);
repository_handle!(RhiDesiredPresenceRepository, DesiredPresence);
+impl<'host> RhiReconciliationJobRepository<'host> {
+ pub(crate) const fn host(&self) -> &'host RhiStateHost {
+ self.host
+ }
+}
+
#[cfg(test)]
mod tests {
use super::*;
diff --git a/tests/build_policy.rs b/tests/build_policy.rs
@@ -27,7 +27,7 @@ fn source_lock_metadata_is_exact_and_nix_is_absent() {
));
for field in [
"config_contract_version = 1",
- "state_contract_version = 4",
+ "state_contract_version = 5",
"admin_contract_version = 1",
"status_contract_version = 1",
"provider_contract_version = 1",
@@ -93,7 +93,7 @@ fn source_lock_binds_the_current_cargo_lock() {
assert!(!SOURCE_LOCK.contains("flake_lock_sha256"));
assert!(!SOURCE_LOCK.contains("lib_revision ="));
assert!(SOURCE_LOCK.ends_with(
- "[contract_versions]\nconfig = 1\nstate = 4\nadmin = 1\nstatus = 1\nprovider = 1\n"
+ "[contract_versions]\nconfig = 1\nstate = 5\nadmin = 1\nstatus = 1\nprovider = 1\n"
));
}
diff --git a/tests/package_boundary.rs b/tests/package_boundary.rs
@@ -11,6 +11,7 @@ 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_JOBS: &str = include_str!("../src/reconciliation_job.rs");
const RUNTIME_FOUNDATION_CONTRACT: &str =
include_str!("../contracts/services_hardening/runtime_foundation.v1.json");
const TRADE_INGEST_CONTRACT: &str =
@@ -27,6 +28,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_job.rs"),
include_str!("../src/runtime_context.rs"),
include_str!("../src/runtime_adapters.rs"),
include_str!("../src/runtime_foundation.rs"),
@@ -93,6 +95,7 @@ fn state_catalog_module_is_private_and_root_api_is_curated() {
"features",
"identity_credential",
"identity_envelope",
+ "reconciliation_job",
"runtime_context",
"runtime_adapters",
"runtime_foundation",
@@ -128,6 +131,8 @@ fn state_catalog_module_is_private_and_root_api_is_curated() {
"validate_rhi_state_catalogs",
"RhiStateCatalogError",
"RhiRuntimeAdapters",
+ "RhiReconciliationJobPolicy",
+ "RhiReconciliationLease",
"RhiTimeEntropyAdapters",
"RhiTransportAdapters",
"RhiCredentialAccess",
@@ -213,7 +218,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, 16);
+ assert_eq!(public_error_count, 17);
}
#[test]
@@ -517,10 +522,14 @@ fn readme_freezes_the_root_only_boundary_and_exact_baseline() {
"[`trade_evidence_persistence.v1.json`](contracts/services_hardening/trade_evidence_persistence.v1.json)",
"## Bounded relay-source ingestion",
"[`trade_source_ingest.v1.json`](contracts/services_hardening/trade_source_ingest.v1.json)",
+ "## Durable reconciliation jobs",
+ "[`reconciliation_jobs.v1.json`](contracts/services_hardening/reconciliation_jobs.v1.json)",
+ "configured queue capacity is enforced beneath a fixed 65,536-job",
+ "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",
"observation, and operational retry do not",
- "schema-v4 source-checkpoint and dirty-generation migration",
+ "schema-v5 bounded reconciliation-job migration",
"one canonical mutation",
"every distinct valid signed",
"does not advance reconciliation checkpoints or dirty generation",
@@ -542,6 +551,17 @@ fn readme_freezes_the_root_only_boundary_and_exact_baseline() {
] {
assert!(README.contains(required), "README is missing {required}");
}
+ for forbidden in [
+ "SystemTime",
+ "thread_rng",
+ "OsRng",
+ "pub mod reconciliation_job",
+ ] {
+ assert!(
+ !RECONCILIATION_JOBS.contains(forbidden),
+ "reconciliation jobs expose forbidden authority {forbidden}"
+ );
+ }
for required in [
"Keep every implementation module private",
"contracts/api_baselines/rhi.txt",
diff --git a/tests/services_hardening_config_lifecycle.rs b/tests/services_hardening_config_lifecycle.rs
@@ -163,7 +163,10 @@ async fn existing_intent_and_offline_apply_bind_exact_append_only_evidence() {
for row in &rows {
assert_eq!(row.try_get::<i64, _>("config_bytes").unwrap(), 32);
assert_eq!(row.try_get::<i64, _>("policy_bytes").unwrap(), 32);
- assert_eq!(row.try_get::<i64, _>("state_contract_version").unwrap(), 4);
+ assert_eq!(
+ row.try_get::<i64, _>("state_contract_version").unwrap(),
+ i64::from(rhi::RHI_STATE_SCHEMA_VERSION)
+ );
assert_eq!(
row.try_get::<String, _>("service_public_key")
.unwrap()
diff --git a/tests/services_hardening_reconciliation_job_contract.rs b/tests/services_hardening_reconciliation_job_contract.rs
@@ -0,0 +1,108 @@
+#![forbid(unsafe_code)]
+
+use rhi::{RHI_RECONCILIATION_JOB_CONTRACT_VERSION, RHI_RECONCILIATION_JOB_MAX_ACTIVE};
+use serde_json::json;
+
+const CONTRACT: &str = include_str!("../contracts/services_hardening/reconciliation_jobs.v1.json");
+const LIB_SOURCE: &str = include_str!("../src/lib.rs");
+const JOB_SOURCE: &str = include_str!("../src/reconciliation_job.rs");
+const README: &str = include_str!("../README");
+
+#[test]
+fn machine_contract_freezes_the_complete_step_187_boundary() {
+ let contract: serde_json::Value = serde_json::from_str(CONTRACT).expect("contract");
+ assert_eq!(contract["schema"], "radroots.rhi.reconciliation-jobs");
+ assert_eq!(contract["schema_version"], 1);
+ assert_eq!(
+ contract["contract_version"],
+ RHI_RECONCILIATION_JOB_CONTRACT_VERSION
+ );
+ assert_eq!(contract["state_schema_version"], 5);
+ assert_eq!(contract["backing_table"], "reconciliation_jobs");
+ assert_eq!(
+ contract["job_identity"],
+ json!({
+ "algorithm": "sha256",
+ "domain": "radroots.rhi.reconciliation_job.v1\\0",
+ "preimage": [
+ "trade_id_16_bytes",
+ "input_generation_u64_be",
+ "evidence_policy_sha256_32_bytes"
+ ]
+ })
+ );
+ assert_eq!(
+ contract["active_job_bound"]["absolute_maximum"],
+ RHI_RECONCILIATION_JOB_MAX_ACTIVE
+ );
+ assert_eq!(
+ contract["states"],
+ json!(["ready", "leased", "exhausted", "superseded", "completed"])
+ );
+ assert_eq!(
+ contract["injected_authority"],
+ json!(["wall_time_unix_ms", "lease_owner", "retry_jitter_ms"])
+ );
+ for forbidden in [
+ "ambient_clock",
+ "ambient_entropy",
+ "network_io",
+ "source_query_inside_transaction",
+ "unbounded_queue",
+ "raw_sqlite_handle",
+ "caller_selected_job_id",
+ "caller_selected_next_attempt_timestamp",
+ "delete_job",
+ ] {
+ assert!(
+ contract["forbidden"]
+ .as_array()
+ .expect("forbidden")
+ .iter()
+ .any(|value| value == forbidden),
+ "missing {forbidden}"
+ );
+ }
+ assert_eq!(
+ contract["deferred"],
+ json!([
+ "per_source_attempt_inventory",
+ "source_execution",
+ "source_completion_commit",
+ "manifest",
+ "reducer",
+ "attestation",
+ "publication"
+ ])
+ );
+}
+
+#[test]
+fn module_and_sqlite_authority_remain_private() {
+ assert!(LIB_SOURCE.contains("mod reconciliation_job;"));
+ assert!(!LIB_SOURCE.contains("pub mod reconciliation_job;"));
+ assert!(LIB_SOURCE.contains("RhiReconciliationJobPolicy"));
+ assert!(README.contains("## Durable reconciliation jobs"));
+ assert!(README.contains(
+ "[`reconciliation_jobs.v1.json`](contracts/services_hardening/reconciliation_jobs.v1.json)"
+ ));
+ for forbidden in [
+ "pub fn sqlite",
+ "pub fn transaction",
+ "pub fn connection",
+ "pub fn into_inner",
+ "pub fn owner",
+ "pub fn from_job_id",
+ "std::time::",
+ "SystemTime",
+ "thread_rng",
+ "OsRng",
+ "radroots_transport",
+ ] {
+ assert!(
+ !JOB_SOURCE.contains(forbidden),
+ "forbidden job authority {forbidden}"
+ );
+ }
+ assert!(JOB_SOURCE.contains("ServiceSqliteTransaction"));
+}
diff --git a/tests/services_hardening_reconciliation_jobs.rs b/tests/services_hardening_reconciliation_jobs.rs
@@ -0,0 +1,500 @@
+#![forbid(unsafe_code)]
+#![cfg(any(target_os = "linux", target_os = "macos"))]
+
+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,
+ 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");
+
+fn runtime(root: &Path, instance: &str) -> RhiRuntimeContext {
+ let invocation = parse_rhi_cli_v1_from([
+ "rhi",
+ "--profile",
+ "repo-local",
+ "--instance",
+ instance,
+ "--repo-local-root",
+ root.to_str().expect("UTF-8 temporary root"),
+ "run",
+ ])
+ .expect("runtime invocation");
+ resolve_rhi_runtime_context(
+ &RadrootsPathResolver::new(RadrootsPlatform::Linux, RadrootsHostEnvironment::default()),
+ &invocation,
+ )
+ .expect("runtime context")
+}
+
+fn metadata(runtime: &RhiRuntimeContext) -> RhiStateMetadata {
+ let config = parse_rhi_config_v1(EXAMPLE.as_bytes(), rhi::RhiConfigProfile::RepoLocal)
+ .expect("configuration");
+ RhiStateMetadata::new(
+ runtime,
+ &config,
+ SourceGeneration::new([0x5a; 32]).expect("generation"),
+ 1_725_000_000_000,
+ )
+ .expect("metadata")
+}
+
+fn migration_evidence() -> (MigrationAppliedAtUnixSeconds, MigrationBuildIdentity) {
+ (
+ MigrationAppliedAtUnixSeconds::new(1_725_000_000).expect("time"),
+ MigrationBuildIdentity::new(
+ env!("CARGO_PKG_VERSION"),
+ "1111111111111111111111111111111111111111",
+ "21b11e7a5120ea949f7ad0838c746873fc73aac2",
+ "rustc-test",
+ "test-target",
+ "service-host",
+ 1,
+ rhi::RHI_STATE_SCHEMA_VERSION,
+ 1,
+ 1,
+ 1,
+ )
+ .expect("build"),
+ )
+}
+
+async fn initialize(runtime: &RhiRuntimeContext, metadata: &RhiStateMetadata) {
+ fs::create_dir_all(runtime.context().paths().state()).expect("state directory");
+ fs::set_permissions(
+ runtime.context().paths().state(),
+ fs::Permissions::from_mode(0o700),
+ )
+ .expect("state permissions");
+ let (applied_at, build) = migration_evidence();
+ initialize_rhi_state(runtime, metadata, applied_at, &build)
+ .await
+ .expect("initialize");
+}
+
+async fn open_writer(
+ runtime: &RhiRuntimeContext,
+ metadata: &RhiStateMetadata,
+) -> rhi::RhiStateHost {
+ let (applied_at, build) = migration_evidence();
+ open_rhi_state_read_write(runtime, metadata, applied_at, &build)
+ .await
+ .expect("writer")
+}
+
+async fn write_dirty(
+ runtime: &RhiRuntimeContext,
+ trade: TradeId,
+ generation: u64,
+ policy: [u8; 32],
+ updated_at_unix_s: u64,
+) {
+ let options = SqliteConnectOptions::new()
+ .filename(runtime.artifacts().state_database())
+ .create_if_missing(false)
+ .foreign_keys(true);
+ let mut connection = SqliteConnection::connect_with(&options)
+ .await
+ .expect("offline fixture connection");
+ sqlx::query(
+ r#"INSERT INTO trade_dirty_generations (
+ trade_id, generation, evidence_policy_sha256, updated_at_unix_s
+ ) VALUES (?, ?, ?, ?)
+ ON CONFLICT(trade_id) DO UPDATE SET
+ generation = excluded.generation,
+ evidence_policy_sha256 = excluded.evidence_policy_sha256,
+ updated_at_unix_s = excluded.updated_at_unix_s"#,
+ )
+ .bind(trade.as_bytes().as_slice())
+ .bind(i64::try_from(generation).expect("generation"))
+ .bind(policy.as_slice())
+ .bind(i64::try_from(updated_at_unix_s).expect("time"))
+ .execute(&mut connection)
+ .await
+ .expect("dirty fixture");
+ connection.close().await.expect("fixture close");
+}
+
+fn policy(
+ queue_capacity: u32,
+ lease_ms: u64,
+ renewal_ms: u64,
+ max_attempts: u16,
+ initial_backoff_ms: u64,
+ maximum_backoff_ms: u64,
+) -> RhiReconciliationJobPolicy {
+ RhiReconciliationJobPolicy::new(
+ queue_capacity,
+ lease_ms,
+ renewal_ms,
+ max_attempts,
+ initial_backoff_ms,
+ maximum_backoff_ms,
+ )
+ .expect("policy")
+}
+
+fn now(value: u64) -> RhiReconciliationUnixMilliseconds {
+ RhiReconciliationUnixMilliseconds::new(value).expect("time")
+}
+
+fn owner(byte: u8) -> RhiReconciliationLeaseOwner {
+ RhiReconciliationLeaseOwner::from_bytes([byte; 16]).expect("owner")
+}
+
+fn delay(value: u64) -> RhiReconciliationRetryDelayMilliseconds {
+ RhiReconciliationRetryDelayMilliseconds::new(value).expect("delay")
+}
+
+#[tokio::test]
+async fn schedule_is_bounded_idempotent_and_uses_the_frozen_job_identity() {
+ let root = tempfile::tempdir().expect("root");
+ let runtime = runtime(root.path(), "primary");
+ let metadata = metadata(&runtime);
+ initialize(&runtime, &metadata).await;
+ let first_trade = TradeId::from_bytes([0x11; 16]);
+ let second_trade = TradeId::from_bytes([0x33; 16]);
+ write_dirty(&runtime, first_trade, 1, [0x22; 32], 1_000).await;
+ write_dirty(&runtime, second_trade, 1, [0x44; 32], 1_000).await;
+ let host = open_writer(&runtime, &metadata).await;
+ let jobs = host.repositories().reconciliation_jobs();
+ let bounded = policy(1, 1_000, 100, 3, 100, 1_000);
+
+ let first = jobs
+ .schedule_trade(first_trade, bounded, now(1_000))
+ .await
+ .expect("first schedule");
+ assert!(first.created());
+ assert_eq!(first.job().state(), RhiReconciliationJobState::Ready);
+ assert_eq!(first.job().input_generation(), 1);
+ assert_eq!(first.job().revision(), 1);
+ assert_eq!(first.job().next_attempt(), Some(now(1_000)));
+ assert_eq!(
+ lower_hex(first.job().id().as_bytes()),
+ "dc7b98b36fb8e839cba83dd7274f25fc27d021d2b46c2f5ab089ddde591fad02"
+ );
+ let replay = jobs
+ .schedule_trade(first_trade, bounded, now(1_001))
+ .await
+ .expect("idempotent schedule");
+ assert!(!replay.created());
+ assert_eq!(replay.job(), first.job());
+
+ let full = jobs
+ .schedule_trade(second_trade, bounded, now(1_001))
+ .await
+ .expect_err("configured queue bound");
+ assert_eq!(full.kind(), RhiReconciliationJobErrorKind::QueueFull);
+ host.close().await.expect("close");
+}
+
+#[tokio::test]
+async fn claims_renew_only_when_due_and_expired_leases_are_reclaimed() {
+ let root = tempfile::tempdir().expect("root");
+ let runtime = runtime(root.path(), "leases");
+ let metadata = metadata(&runtime);
+ initialize(&runtime, &metadata).await;
+ let trade = TradeId::from_bytes([0x12; 16]);
+ write_dirty(&runtime, trade, 1, [0x23; 32], 1_000).await;
+ let host = open_writer(&runtime, &metadata).await;
+ let jobs = host.repositories().reconciliation_jobs();
+ jobs.schedule_trade(trade, policy(8, 1_000, 100, 4, 100, 1_000), now(1_000))
+ .await
+ .expect("schedule");
+
+ let first = jobs
+ .claim_next(owner(1), now(1_000))
+ .await
+ .expect("claim")
+ .expect("job");
+ assert_eq!(first.job().attempt_count(), 1);
+ assert_eq!(first.lease_expires(), now(2_000));
+ assert_eq!(first.renewal_due(), now(1_900));
+ assert_eq!(
+ jobs.renew(first, now(1_899))
+ .await
+ .expect_err("early renewal")
+ .kind(),
+ RhiReconciliationJobErrorKind::NotReady
+ );
+ let renewed = jobs.renew(first, now(1_900)).await.expect("renew");
+ assert_eq!(renewed.job().revision(), 3);
+ assert_eq!(renewed.lease_expires(), now(2_900));
+ assert!(
+ jobs.claim_next(owner(2), now(2_899))
+ .await
+ .expect("no early reclaim")
+ .is_none()
+ );
+ let reclaimed = jobs
+ .claim_next(owner(2), now(2_900))
+ .await
+ .expect("reclaim")
+ .expect("expired lease");
+ assert_eq!(reclaimed.job().attempt_count(), 2);
+ assert_eq!(reclaimed.job().revision(), 4);
+ assert_eq!(
+ jobs.record_failure(renewed, now(2_900), delay(0))
+ .await
+ .expect_err("expired lease")
+ .kind(),
+ RhiReconciliationJobErrorKind::LeaseLost
+ );
+ assert_eq!(
+ jobs.renew(first, now(1_950))
+ .await
+ .expect_err("stale revision")
+ .kind(),
+ RhiReconciliationJobErrorKind::LeaseLost
+ );
+ host.close().await.expect("close");
+}
+
+#[tokio::test]
+async fn failures_schedule_exact_jitter_and_the_final_attempt_exhausts() {
+ let root = tempfile::tempdir().expect("root");
+ let runtime = runtime(root.path(), "retry");
+ let metadata = metadata(&runtime);
+ initialize(&runtime, &metadata).await;
+ let trade = TradeId::from_bytes([0x13; 16]);
+ write_dirty(&runtime, trade, 1, [0x24; 32], 1_000).await;
+ let host = open_writer(&runtime, &metadata).await;
+ let jobs = host.repositories().reconciliation_jobs();
+ jobs.schedule_trade(trade, policy(8, 1_000, 100, 2, 100, 1_000), now(1_000))
+ .await
+ .expect("schedule");
+ let first = jobs
+ .claim_next(owner(3), now(1_000))
+ .await
+ .expect("claim")
+ .expect("job");
+ assert_eq!(first.retry_delay_upper_bound(), 100);
+ assert_eq!(
+ jobs.record_failure(first, now(1_100), delay(101))
+ .await
+ .expect_err("delay above governed cap")
+ .kind(),
+ RhiReconciliationJobErrorKind::InvalidInput
+ );
+ let retry = jobs
+ .record_failure(first, now(1_100), delay(100))
+ .await
+ .expect("retry schedule");
+ assert_eq!(retry.state(), RhiReconciliationJobState::Ready);
+ assert_eq!(retry.failure_count(), 1);
+ assert_eq!(retry.next_attempt(), Some(now(1_200)));
+ assert!(
+ jobs.claim_next(owner(4), now(1_199))
+ .await
+ .expect("not due")
+ .is_none()
+ );
+ let second = jobs
+ .claim_next(owner(4), now(1_200))
+ .await
+ .expect("second claim")
+ .expect("job");
+ assert_eq!(second.retry_delay_upper_bound(), 200);
+ let exhausted = jobs
+ .record_failure(second, now(1_300), delay(200))
+ .await
+ .expect("final failure");
+ assert_eq!(exhausted.state(), RhiReconciliationJobState::Exhausted);
+ assert_eq!(exhausted.attempt_count(), 2);
+ assert_eq!(exhausted.failure_count(), 2);
+ assert_eq!(exhausted.next_attempt(), None);
+ assert!(
+ jobs.claim_next(owner(5), now(9_000))
+ .await
+ .expect("terminal queue")
+ .is_none()
+ );
+ host.close().await.expect("close");
+}
+
+#[tokio::test]
+async fn an_expired_final_attempt_is_exhausted_instead_of_reclaimed() {
+ let root = tempfile::tempdir().expect("root");
+ let runtime = runtime(root.path(), "expired-final");
+ let metadata = metadata(&runtime);
+ initialize(&runtime, &metadata).await;
+ let trade = TradeId::from_bytes([0x18; 16]);
+ write_dirty(&runtime, trade, 1, [0x29; 32], 1_000).await;
+ let host = open_writer(&runtime, &metadata).await;
+ let jobs = host.repositories().reconciliation_jobs();
+ let job_policy = policy(8, 1_000, 100, 1, 100, 1_000);
+ jobs.schedule_trade(trade, job_policy, now(1_000))
+ .await
+ .expect("schedule");
+ jobs.claim_next(owner(10), now(1_000))
+ .await
+ .expect("claim")
+ .expect("job");
+
+ assert!(
+ jobs.claim_next(owner(11), now(2_000))
+ .await
+ .expect("expired final attempt")
+ .is_none()
+ );
+ let retained = jobs
+ .schedule_trade(trade, job_policy, now(2_000))
+ .await
+ .expect("idempotent retained job");
+ assert!(!retained.created());
+ assert_eq!(retained.job().state(), RhiReconciliationJobState::Exhausted);
+ assert_eq!(retained.job().attempt_count(), 1);
+ assert_eq!(retained.job().failure_count(), 1);
+ host.close().await.expect("close");
+}
+
+#[tokio::test]
+async fn lease_and_schedule_state_survive_reopen_and_new_generation_supersedes() {
+ let root = tempfile::tempdir().expect("root");
+ let runtime = runtime(root.path(), "reopen");
+ let metadata = metadata(&runtime);
+ initialize(&runtime, &metadata).await;
+ let trade = TradeId::from_bytes([0x14; 16]);
+ write_dirty(&runtime, trade, 1, [0x25; 32], 1_000).await;
+ let job_policy = policy(8, 1_000, 100, 3, 100, 1_000);
+ let host = open_writer(&runtime, &metadata).await;
+ let jobs = host.repositories().reconciliation_jobs();
+ let first_job = jobs
+ .schedule_trade(trade, job_policy, now(1_000))
+ .await
+ .expect("schedule")
+ .job();
+ let stale = jobs
+ .claim_next(owner(6), now(1_000))
+ .await
+ .expect("claim")
+ .expect("job");
+ host.close().await.expect("crash boundary close");
+
+ write_dirty(&runtime, trade, 2, [0x26; 32], 1_001).await;
+ let host = open_writer(&runtime, &metadata).await;
+ let jobs = host.repositories().reconciliation_jobs();
+ let newer = jobs
+ .schedule_trade(trade, job_policy, now(1_900))
+ .await
+ .expect("new generation");
+ assert!(newer.created());
+ assert_eq!(newer.job().input_generation(), 2);
+ assert_ne!(newer.job().id(), first_job.id());
+ assert_eq!(
+ jobs.renew(stale, now(1_900))
+ .await
+ .expect_err("superseded lease")
+ .kind(),
+ RhiReconciliationJobErrorKind::LeaseLost
+ );
+ let claimed = jobs
+ .claim_next(owner(7), now(1_900))
+ .await
+ .expect("claim new")
+ .expect("new job");
+ assert_eq!(claimed.job().id(), newer.job().id());
+ host.close().await.expect("close");
+}
+
+#[tokio::test]
+async fn concurrent_claims_have_one_winner_and_read_only_state_cannot_mutate() {
+ let root = tempfile::tempdir().expect("root");
+ let runtime = runtime(root.path(), "concurrent");
+ let metadata = metadata(&runtime);
+ initialize(&runtime, &metadata).await;
+ let trade = TradeId::from_bytes([0x15; 16]);
+ write_dirty(&runtime, trade, 1, [0x27; 32], 1_000).await;
+ let host = open_writer(&runtime, &metadata).await;
+ let jobs = host.repositories().reconciliation_jobs();
+ jobs.schedule_trade(trade, policy(8, 1_000, 100, 3, 100, 1_000), now(1_000))
+ .await
+ .expect("schedule");
+ let (left, right) = tokio::join!(
+ jobs.claim_next(owner(8), now(1_000)),
+ jobs.claim_next(owner(9), now(1_000))
+ );
+ let winners = [left, right]
+ .into_iter()
+ .filter_map(|result| result.expect("claim result"))
+ .count();
+ assert_eq!(winners, 1);
+ host.close().await.expect("close");
+
+ let inspection = open_rhi_state_inspection(&runtime, &metadata)
+ .await
+ .expect("inspection");
+ let error = inspection
+ .repositories()
+ .reconciliation_jobs()
+ .schedule_trade(trade, policy(8, 1_000, 100, 3, 100, 1_000), now(2_000))
+ .await
+ .expect_err("read-only mutation");
+ assert_eq!(error.kind(), RhiReconciliationJobErrorKind::InvalidMode);
+ inspection.close().await.expect("inspection close");
+}
+
+#[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)
+ .expect("exact maximum policy");
+ assert_eq!(maximum.queue_capacity(), 65_536);
+ assert_eq!(maximum.lease_duration_milliseconds(), 300_000);
+ assert_eq!(maximum.lease_renewal_milliseconds(), 150_000);
+ assert_eq!(maximum.max_attempts(), 100);
+ assert_eq!(maximum.initial_backoff_milliseconds(), 60_000);
+ assert_eq!(maximum.maximum_backoff_milliseconds(), 3_600_000);
+ let configuration = parse_rhi_config_v1(EXAMPLE.as_bytes(), rhi::RhiConfigProfile::RepoLocal)
+ .expect("configuration");
+ let configured = RhiReconciliationJobPolicy::from_configuration(&configuration)
+ .expect("configured reconciliation policy");
+ assert_eq!(configured.queue_capacity(), 4_096);
+ assert_eq!(configured.lease_duration_milliseconds(), 30_000);
+ assert_eq!(configured.lease_renewal_milliseconds(), 10_000);
+ assert_eq!(configured.max_attempts(), 10);
+ assert_eq!(configured.initial_backoff_milliseconds(), 250);
+ assert_eq!(configured.maximum_backoff_milliseconds(), 30_000);
+ assert!(RhiReconciliationJobPolicy::new(0, 1_000, 100, 1, 1, 1).is_err());
+ assert!(RhiReconciliationJobPolicy::new(65_537, 1_000, 100, 1, 1, 1).is_err());
+ assert!(RhiReconciliationJobPolicy::new(1, 999, 100, 1, 1, 1).is_err());
+ assert!(RhiReconciliationJobPolicy::new(1, 300_001, 100, 1, 1, 1).is_err());
+ assert!(RhiReconciliationJobPolicy::new(1, 1_000, 1_000, 1, 1, 1).is_err());
+ assert!(RhiReconciliationJobPolicy::new(1, 1_000, 100, 0, 1, 1).is_err());
+ assert!(RhiReconciliationJobPolicy::new(1, 1_000, 100, 101, 1, 1).is_err());
+ assert!(RhiReconciliationJobPolicy::new(1, 1_000, 100, 1, 0, 1).is_err());
+ assert!(RhiReconciliationJobPolicy::new(1, 1_000, 100, 1, 2, 1).is_err());
+ assert!(RhiReconciliationUnixMilliseconds::new(i64::MAX as u64).is_ok());
+ assert!(RhiReconciliationUnixMilliseconds::new(i64::MAX as u64 + 1).is_err());
+ assert!(RhiReconciliationRetryDelayMilliseconds::new(3_600_000).is_ok());
+ assert!(RhiReconciliationRetryDelayMilliseconds::new(3_600_001).is_err());
+ assert!(RhiReconciliationLeaseOwner::from_bytes([0; 16]).is_err());
+ let secret = RhiReconciliationLeaseOwner::from_bytes([0xab; 16]).expect("owner");
+ assert_eq!(
+ format!("{secret:?}"),
+ "RhiReconciliationLeaseOwner([redacted])"
+ );
+ let error =
+ RhiReconciliationJobPolicy::new(0, 1_000, 100, 1, 1, 1).expect_err("invalid policy");
+ assert!(Error::source(&error).is_none());
+ assert_eq!(error.code(), "reconciliation_job_input_invalid");
+ assert!(!format!("{error} {error:?}").contains("65536"));
+}
+
+fn lower_hex(bytes: &[u8]) -> String {
+ const DIGITS: &[u8; 16] = b"0123456789abcdef";
+ let mut output = String::with_capacity(bytes.len() * 2);
+ for byte in bytes {
+ output.push(char::from(DIGITS[usize::from(byte >> 4)]));
+ output.push(char::from(DIGITS[usize::from(byte & 0x0f)]));
+ }
+ output
+}
diff --git a/tests/services_hardening_state_catalog.rs b/tests/services_hardening_state_catalog.rs
@@ -14,8 +14,9 @@ use rhi::{
RHI_STATE_SCHEMA_VERSION_3_MIGRATION_SHA256, RHI_STATE_SCHEMA_VERSION_3_OBJECT_COUNT,
RHI_STATE_SCHEMA_VERSION_3_SHA256, RHI_STATE_SCHEMA_VERSION_4_MIGRATION_SHA256,
RHI_STATE_SCHEMA_VERSION_4_OBJECT_COUNT, RHI_STATE_SCHEMA_VERSION_4_SHA256,
- RhiStateCatalogErrorKind, rhi_migration_catalog, rhi_schema_catalog,
- validate_rhi_state_catalogs,
+ RHI_STATE_SCHEMA_VERSION_5_MIGRATION_SHA256, RHI_STATE_SCHEMA_VERSION_5_OBJECT_COUNT,
+ RHI_STATE_SCHEMA_VERSION_5_SHA256, RhiStateCatalogErrorKind, rhi_migration_catalog,
+ rhi_schema_catalog, validate_rhi_state_catalogs,
};
const CATALOG_SOURCE: &str = include_str!("../src/state_catalog.rs");
@@ -23,14 +24,14 @@ const LIB_SOURCE: &str = include_str!("../src/lib.rs");
const MANIFEST: &str = include_str!("../Cargo.toml");
#[test]
-fn schema_v1_through_v4_catalogs_have_exact_literal_identities() {
+fn schema_v1_through_v5_catalogs_have_exact_literal_identities() {
let migrations = rhi_migration_catalog().expect("RHI migration catalog");
let schema = rhi_schema_catalog().expect("RHI schema catalog");
assert_eq!(RHI_STATE_BASE_SCHEMA_VERSION, 1);
- assert_eq!(RHI_STATE_SCHEMA_VERSION, 4);
- assert_eq!(migrations.descriptors().len(), 3);
- assert_eq!(migrations.current_version(), 4);
+ assert_eq!(RHI_STATE_SCHEMA_VERSION, 5);
+ assert_eq!(migrations.descriptors().len(), 4);
+ assert_eq!(migrations.current_version(), 5);
assert_eq!(migrations.descriptors()[0].target_version(), 2);
assert_eq!(
migrations.descriptors()[0].name().as_str(),
@@ -58,12 +59,21 @@ fn schema_v1_through_v4_catalogs_have_exact_literal_identities() {
migrations.descriptors()[2].checksum().as_bytes(),
&RHI_STATE_SCHEMA_VERSION_4_MIGRATION_SHA256
);
+ assert_eq!(migrations.descriptors()[3].target_version(), 5);
+ assert_eq!(
+ migrations.descriptors()[3].name().as_str(),
+ "create_reconciliation_jobs"
+ );
+ assert_eq!(
+ migrations.descriptors()[3].checksum().as_bytes(),
+ &RHI_STATE_SCHEMA_VERSION_5_MIGRATION_SHA256
+ );
assert_eq!(
migrations.digest().as_bytes(),
&RHI_MIGRATION_CATALOG_SHA256
);
- assert_eq!(schema.versions().len(), 4);
+ assert_eq!(schema.versions().len(), 5);
let version = schema.versions()[0];
assert_eq!(version.version(), 1);
assert_eq!(
@@ -108,13 +118,24 @@ fn schema_v1_through_v4_catalogs_have_exact_literal_identities() {
version.digest().as_bytes(),
&RHI_STATE_SCHEMA_VERSION_4_SHA256
);
+ let version = schema.versions()[4];
+ assert_eq!(version.version(), 5);
+ assert_eq!(
+ version.object_count(),
+ RHI_STATE_SCHEMA_VERSION_5_OBJECT_COUNT
+ );
+ assert_eq!(version.object_count(), 33);
+ assert_eq!(
+ version.digest().as_bytes(),
+ &RHI_STATE_SCHEMA_VERSION_5_SHA256
+ );
assert_eq!(schema.digest().as_bytes(), &RHI_STATE_SCHEMA_CATALOG_SHA256);
assert_eq!(schema.migration_catalog_digest(), migrations.digest());
validate_rhi_state_catalogs(&migrations, &schema).expect("exact catalogs");
assert_eq!(
lower_hex(&RHI_MIGRATION_CATALOG_SHA256),
- "295ff5e9ac0eadc1f2c2d8c258c2eaa816f946560b689dcb432fdc99e79ec353"
+ "4ed32a5a71f5bf454262a9082ad74d0d6a13a556054659ac7b9d877259b816d8"
);
assert_eq!(
lower_hex(&RHI_STATE_SCHEMA_VERSION_1_SHA256),
@@ -145,8 +166,16 @@ fn schema_v1_through_v4_catalogs_have_exact_literal_identities() {
"9ca78a54b0ea2013e7aa70d09beadcfb6cdde07c40429dffe91cb83c2a425cea"
);
assert_eq!(
+ lower_hex(&RHI_STATE_SCHEMA_VERSION_5_MIGRATION_SHA256),
+ "8275ffb3d5c9fa0c7648f89e8bfb92877b0dcd5ef3bf0b5254b7ffd2580bd45d"
+ );
+ assert_eq!(
+ lower_hex(&RHI_STATE_SCHEMA_VERSION_5_SHA256),
+ "eef67b48400dee234f6c36d422551c7e9983150b887f264250f9c87a1a313e60"
+ );
+ assert_eq!(
lower_hex(&RHI_STATE_SCHEMA_CATALOG_SHA256),
- "58789c1d396514ce19a560b67bf760650b9d6fa378bc8234c1f6a39ad2f556e6"
+ "0d2c595a3ca3fbc9b3bf6bfc6168a194ac788d4931c46200e2f1379ebb523ece"
);
}
@@ -194,9 +223,19 @@ fn independent_validator_rejects_migration_or_schema_drift() {
SchemaVersionCatalog::computed_digest(4, [version_two_object()]).expect("v4 digest");
let version_four = SchemaVersionCatalog::new(4, [version_two_object()], snapshot_digest)
.expect("version four");
+ let snapshot_digest =
+ SchemaVersionCatalog::computed_digest(5, [version_two_object()]).expect("v5 digest");
+ let version_five = SchemaVersionCatalog::new(5, [version_two_object()], snapshot_digest)
+ .expect("version five");
let schema = SchemaCatalog::new(
&exact_migrations,
- [version_one, version_two, version_three, version_four],
+ [
+ version_one,
+ version_two,
+ version_three,
+ version_four,
+ version_five,
+ ],
)
.expect("drift schema catalog");
assert_eq!(
@@ -242,9 +281,19 @@ fn catalog_errors_are_stable_source_free_and_redacted() {
SchemaVersionCatalog::computed_digest(4, [secret_object()]).expect("v4 digest");
let version_four =
SchemaVersionCatalog::new(4, [secret_object()], version_four_digest).expect("version four");
+ let version_five_digest =
+ SchemaVersionCatalog::computed_digest(5, [secret_object()]).expect("v5 digest");
+ let version_five =
+ SchemaVersionCatalog::new(5, [secret_object()], version_five_digest).expect("version five");
let schema = SchemaCatalog::new(
&migrations,
- [version_one, version_two, version_three, version_four],
+ [
+ version_one,
+ version_two,
+ version_three,
+ version_four,
+ version_five,
+ ],
)
.expect("schema catalog");
let error = validate_rhi_state_catalogs(&migrations, &schema).expect_err("mismatch");
diff --git a/tests/services_hardening_state_repository_topology.rs b/tests/services_hardening_state_repository_topology.rs
@@ -27,8 +27,8 @@ fn machine_contract_and_typed_descriptor_inventory_are_exact() {
assert_eq!(
contract["deferred_behavior"],
json!([
- "schema_migration_and_verification",
- "repository_crud",
+ "later_schema_migrations",
+ "non_job_repository_crud",
"backup_restore_and_recovery",
"network_io",
"task_supervision"