lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

commit acd3477ba8314e3cd88b47d977e362aa8c115292
parent 78e2c5ac569b4ac57da6cb78b522bc82ec3d5c12
Author: triesap <tyson@radroots.org>
Date:   Mon,  3 Aug 2026 09:01:04 +0000

sync: implement status and retry-decision reports

- aggregate storage, transport, signer, projection, and outbox health
- classify optional capabilities with typed readiness states
- expose passive deadline-aware retry decisions without scheduling
- generate and verify the versioned sync status protocol receipt

Diffstat:
Mcontracts/codegen/protocol_v1.inventory.json | 26+++++++++++++++++++++++++-
Mcontracts/codegen/protocol_v1.inventory.sha256 | 2+-
Mcrates/protocol/src/runtime/v1.rs | 163+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/sync/src/lib.rs | 1+
Mcrates/sync/src/policy.rs | 2++
Mcrates/sync/src/status.rs | 360++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcrates/sync/tests/engine_composition.rs | 104++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----
Mcrates/sync/tests/push_enqueue.rs | 83++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
8 files changed, 733 insertions(+), 8 deletions(-)

diff --git a/contracts/codegen/protocol_v1.inventory.json b/contracts/codegen/protocol_v1.inventory.json @@ -203,7 +203,7 @@ { "module": "runtime::v1", "path": "crates/protocol/src/runtime/v1.rs", - "sha256": "3ef493bd3d2e18f52f2efa0c508fe1bcdaced26d7f53674f130ccf38dd1b1e36", + "sha256": "15b0d104b7f1c01dffd69892fc842eee6de214fe11495b3b1267bfe2aed66a82", "types": [ { "rust_path": "radroots_protocol::runtime::v1::ApprovalRequirement", @@ -250,6 +250,30 @@ "kind": "enum" }, { + "rust_path": "radroots_protocol::runtime::v1::SyncCapabilityState", + "kind": "enum" + }, + { + "rust_path": "radroots_protocol::runtime::v1::SyncHealth", + "kind": "enum" + }, + { + "rust_path": "radroots_protocol::runtime::v1::SyncOutboxStatus", + "kind": "struct" + }, + { + "rust_path": "radroots_protocol::runtime::v1::SyncProjectionStatus", + "kind": "struct" + }, + { + "rust_path": "radroots_protocol::runtime::v1::SyncRetryDecision", + "kind": "enum" + }, + { + "rust_path": "radroots_protocol::runtime::v1::SyncStatusReceipt", + "kind": "struct" + }, + { "rust_path": "radroots_protocol::runtime::v1::TransportRoute", "kind": "struct" } diff --git a/contracts/codegen/protocol_v1.inventory.sha256 b/contracts/codegen/protocol_v1.inventory.sha256 @@ -1 +1 @@ -7b9f51e93a819f2d345006dc7ea6f4d9be3be3ab322d0eb791886cdee084d019 +f197c1b3afaa8eb1420f4c9d27168a681b22023eede2eb8fb7e0e0fe8d920e82 diff --git a/crates/protocol/src/runtime/v1.rs b/crates/protocol/src/runtime/v1.rs @@ -20,6 +20,127 @@ use crate::{ /// Schema generation shared by every runtime operation request and receipt. pub const OPERATION_SCHEMA_VERSION: u16 = 1; +/// Typed readiness of one optional synchronization capability. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum SyncCapabilityState { + Unsupported, + Compiled, + Configured, + Available, + Degraded, +} + +/// Aggregate synchronization health for the passive status operation. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum SyncHealth { + Healthy, + Degraded, + Unavailable, +} + +/// Host-action classification for one durable outbox record. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum SyncRetryDecision { + Ready, + DeferredUntil { unix_ms: u64 }, + InFlightUntil { unix_ms: u64 }, + Satisfied, + Exhausted, + Expired, +} + +/// Passive durable outbox cardinalities for `sync.status` generation 1. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[cfg_attr(feature = "serde", serde(deny_unknown_fields))] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct SyncOutboxStatus { + pub pending: u64, + pub leased: u64, + pub retryable: u64, + pub satisfied: u64, + pub exhausted: u64, +} + +/// Passive projection cardinalities for `sync.status` generation 1. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[cfg_attr(feature = "serde", serde(deny_unknown_fields))] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct SyncProjectionStatus { + pub ready: u32, + pub invalidated: u32, + pub rebuilding: u32, + pub failed: u32, + pub untracked: u32, +} + +/// Versioned passive receipt for the `sync.status` operation. +#[cfg_attr(feature = "serde", derive(serde::Serialize))] +#[cfg_attr(feature = "serde", serde(deny_unknown_fields))] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct SyncStatusReceipt { + pub schema_version: u16, + pub health: SyncHealth, + pub storage: SyncCapabilityState, + pub source: SyncCapabilityState, + pub sink: SyncCapabilityState, + pub signer: SyncCapabilityState, + pub outbox: SyncOutboxStatus, + pub projections: SyncProjectionStatus, +} + +impl SyncStatusReceipt { + /// Rejects status receipts from an unsupported operation generation. + pub const fn validate(&self) -> Result<(), Error> { + if self.schema_version != OPERATION_SCHEMA_VERSION { + return Err(Error::UnsupportedOperationSchemaVersion { + operation_id: OperationId::SyncStatus, + version: self.schema_version, + }); + } + Ok(()) + } +} + +#[cfg(feature = "serde")] +impl<'de> serde::Deserialize<'de> for SyncStatusReceipt { + fn deserialize<D>(deserializer: D) -> Result<Self, D::Error> + where + D: serde::Deserializer<'de>, + { + #[derive(serde::Deserialize)] + #[serde(deny_unknown_fields)] + struct Wire { + schema_version: u16, + health: SyncHealth, + storage: SyncCapabilityState, + source: SyncCapabilityState, + sink: SyncCapabilityState, + signer: SyncCapabilityState, + outbox: SyncOutboxStatus, + projections: SyncProjectionStatus, + } + let wire = Wire::deserialize(deserializer)?; + let receipt = Self { + schema_version: wire.schema_version, + health: wire.health, + storage: wire.storage, + source: wire.source, + sink: wire.sink, + signer: wire.signer, + outbox: wire.outbox, + projections: wire.projections, + }; + receipt.validate().map_err(serde::de::Error::custom)?; + Ok(receipt) + } +} + macro_rules! operation_ids { ($( $variant:ident => $value:literal ),+ $(,)?) => { #[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)] @@ -1178,4 +1299,46 @@ mod tests { "\"profile.inspect\"" ); } + + #[cfg(feature = "serde")] + #[test] + fn sync_status_receipt_is_typed_and_rejects_wire_drift() { + let receipt = SyncStatusReceipt { + schema_version: OPERATION_SCHEMA_VERSION, + health: SyncHealth::Degraded, + storage: SyncCapabilityState::Available, + source: SyncCapabilityState::Available, + sink: SyncCapabilityState::Degraded, + signer: SyncCapabilityState::Configured, + outbox: SyncOutboxStatus { + pending: 1, + leased: 2, + retryable: 3, + satisfied: 4, + exhausted: 5, + }, + projections: SyncProjectionStatus { + ready: 6, + invalidated: 7, + rebuilding: 8, + failed: 9, + untracked: 10, + }, + }; + receipt.validate().expect("status receipt"); + let value = serde_json::to_value(receipt).expect("status JSON"); + assert_eq!(value["source"], "available"); + assert_eq!(value["outbox"]["retryable"], 3); + assert_eq!(value["projections"]["untracked"], 10); + assert_eq!( + serde_json::from_value::<SyncStatusReceipt>(value.clone()).expect("status decode"), + receipt + ); + let mut invalid_version = value.clone(); + invalid_version["schema_version"] = 2.into(); + assert!(serde_json::from_value::<SyncStatusReceipt>(invalid_version).is_err()); + let mut unknown = value; + unknown["unknown"] = true.into(); + assert!(serde_json::from_value::<SyncStatusReceipt>(unknown).is_err()); + } } diff --git a/crates/sync/src/lib.rs b/crates/sync/src/lib.rs @@ -13,3 +13,4 @@ pub use engine::Engine; pub use policy::Error; pub use pull::{PullReceipt, PullRequest}; pub use push::{PushReceipt, PushRequest}; +pub use status::SyncStatus; diff --git a/crates/sync/src/policy.rs b/crates/sync/src/policy.rs @@ -209,6 +209,7 @@ pub enum Error { InvalidSignerOutput, InvalidDeliveryRequest, MissingSink, + InvalidStatusRequest, } impl core::fmt::Display for Error { @@ -238,6 +239,7 @@ impl core::fmt::Display for Error { Self::InvalidSignerOutput => "sync signer output failed canonical verification", Self::InvalidDeliveryRequest => "sync delivery request is invalid", Self::MissingSink => "sync engine has no event sink", + Self::InvalidStatusRequest => "sync status request is invalid", }) } } diff --git a/crates/sync/src/status.rs b/crates/sync/src/status.rs @@ -1 +1,359 @@ -//! Passive synchronization status aggregation. +//! Passive synchronization status aggregation and host retry decisions. + +use std::collections::BTreeSet; + +use radroots_protocol::runtime::v1::{ + OPERATION_SCHEMA_VERSION, SyncCapabilityState, SyncHealth, SyncOutboxStatus, + SyncProjectionStatus, SyncRetryDecision, SyncStatusReceipt, +}; +use radroots_signing::{SignerStatus, status::SignerAvailability}; +use radroots_storage::{ + BackupSource, EventStore, Outbox, ProjectionStore, + outbox::{OutboxRecord, OutboxStage, OutboxStatus}, + projection::{ProjectionHealth, ProjectionId, ProjectionStatus}, + status::{EventStoreHealth, EventStoreStatus, IntegrityHealth, ShutdownState, StorageStatus}, +}; +use radroots_transport::{SinkStatus, SourceStatus, capability::Availability}; + +use crate::{Engine, policy::Error}; + +const STATUS_PROJECTION_LIMIT: usize = 256; + +/// Typed report for one optional injected host capability. +#[cfg_attr(feature = "serde", derive(serde::Serialize))] +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct CapabilityReport<T> { + state: SyncCapabilityState, + status: Option<T>, +} + +impl<T> CapabilityReport<T> { + pub const fn state(&self) -> SyncCapabilityState { + self.state + } + + pub const fn status(&self) -> Option<&T> { + self.status.as_ref() + } + + const fn unsupported() -> Self { + Self { + state: SyncCapabilityState::Unsupported, + status: None, + } + } + + const fn compiled(status: Option<T>) -> Self { + Self { + state: SyncCapabilityState::Compiled, + status, + } + } + + const fn reported(state: SyncCapabilityState, status: T) -> Self { + Self { + state, + status: Some(status), + } + } +} + +/// Requested projection and its optional durable status record. +#[cfg_attr(feature = "serde", derive(serde::Serialize))] +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct ProjectionReport { + projection_id: ProjectionId, + status: Option<ProjectionStatus>, +} + +impl ProjectionReport { + pub const fn projection_id(&self) -> &ProjectionId { + &self.projection_id + } + + pub const fn status(&self) -> Option<&ProjectionStatus> { + self.status.as_ref() + } +} + +/// One passive, side-effect-free synchronization health snapshot. +#[cfg_attr(feature = "serde", derive(serde::Serialize))] +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct SyncStatus { + health: SyncHealth, + storage: StorageStatus, + events: EventStoreStatus, + outbox: OutboxStatus, + source: CapabilityReport<SourceStatus>, + sink: CapabilityReport<SinkStatus>, + signer: CapabilityReport<SignerStatus>, + projections: Vec<ProjectionReport>, +} + +impl SyncStatus { + pub const fn health(&self) -> SyncHealth { + self.health + } + + pub const fn storage(&self) -> StorageStatus { + self.storage + } + + pub const fn events(&self) -> &EventStoreStatus { + &self.events + } + + pub const fn outbox(&self) -> OutboxStatus { + self.outbox + } + + pub const fn source(&self) -> &CapabilityReport<SourceStatus> { + &self.source + } + + pub const fn sink(&self) -> &CapabilityReport<SinkStatus> { + &self.sink + } + + pub const fn signer(&self) -> &CapabilityReport<SignerStatus> { + &self.signer + } + + pub fn projections(&self) -> &[ProjectionReport] { + self.projections.as_slice() + } + + /// Converts native reports into the versioned passive protocol receipt. + pub fn to_protocol(&self) -> SyncStatusReceipt { + let mut projection = SyncProjectionStatus { + ready: 0, + invalidated: 0, + rebuilding: 0, + failed: 0, + untracked: 0, + }; + for report in &self.projections { + match report.status.as_ref().map(ProjectionStatus::health) { + Some(ProjectionHealth::Ready) => projection.ready += 1, + Some(ProjectionHealth::Invalidated) => projection.invalidated += 1, + Some(ProjectionHealth::Rebuilding) => projection.rebuilding += 1, + Some(ProjectionHealth::Failed) => projection.failed += 1, + None => projection.untracked += 1, + } + } + SyncStatusReceipt { + schema_version: OPERATION_SCHEMA_VERSION, + health: self.health, + storage: storage_state(self.storage, self.events.health()), + source: self.source.state, + sink: self.sink.state, + signer: self.signer.state, + outbox: SyncOutboxStatus { + pending: self.outbox.pending, + leased: self.outbox.leased, + retryable: self.outbox.retryable, + satisfied: self.outbox.satisfied, + exhausted: self.outbox.exhausted, + }, + projections: projection, + } + } +} + +impl Engine { + /// Aggregates passive status without spawning work or initiating recovery. + pub async fn status(&self, projection_ids: &[ProjectionId]) -> Result<SyncStatus, Error> { + if projection_ids.len() > STATUS_PROJECTION_LIMIT + || projection_ids.iter().collect::<BTreeSet<_>>().len() != projection_ids.len() + { + return Err(Error::InvalidStatusRequest); + } + let storage = BackupSource::status(self.storage.as_ref()) + .await + .map_err(|_| Error::StorageFailed)?; + let events = EventStore::status(self.storage.as_ref()) + .await + .map_err(|_| Error::StorageFailed)?; + let outbox = Outbox::status(self.storage.as_ref()) + .await + .map_err(|_| Error::StorageFailed)?; + if outbox.total().is_none() { + return Err(Error::StorageFailed); + } + let source = source_report(self).await; + let sink = sink_report(self).await; + let signer = signer_report(self).await; + let mut projections = Vec::with_capacity(projection_ids.len()); + for projection_id in projection_ids { + let status = ProjectionStore::status(self.storage.as_ref(), projection_id.clone()) + .await + .map_err(|_| Error::StorageFailed)?; + projections.push(ProjectionReport { + projection_id: projection_id.clone(), + status, + }); + } + let health = aggregate_health(storage, &events, &source, &sink, &signer, &projections); + Ok(SyncStatus { + health, + storage, + events, + outbox, + source, + sink, + signer, + projections, + }) + } + + /// Classifies host action for one durable plan without mutating it. + pub fn retry_decision( + &self, + record: &OutboxRecord, + now_unix_ms: u64, + ) -> Result<SyncRetryDecision, Error> { + if now_unix_ms == 0 { + return Err(Error::ClockUnavailable); + } + match record.stage() { + OutboxStage::Satisfied => return Ok(SyncRetryDecision::Satisfied), + OutboxStage::Exhausted => return Ok(SyncRetryDecision::Exhausted), + OutboxStage::Pending | OutboxStage::Leased | OutboxStage::Retryable => {} + } + if now_unix_ms >= record.request().deadline_unix_ms() { + return Ok(SyncRetryDecision::Expired); + } + if let Some(lease) = record.lease() + && lease.is_active_at(now_unix_ms) + { + return Ok(SyncRetryDecision::InFlightUntil { + unix_ms: lease.expires_at_unix_ms(), + }); + } + if let Some(unix_ms) = record.retry_not_before_unix_ms() + && now_unix_ms < unix_ms + { + return Ok(SyncRetryDecision::DeferredUntil { unix_ms }); + } + Ok(SyncRetryDecision::Ready) + } +} + +async fn source_report(engine: &Engine) -> CapabilityReport<SourceStatus> { + let Some(source) = engine.source.as_deref() else { + return CapabilityReport::unsupported(); + }; + match source.status().await { + Ok(status) if !status.is_configured() => CapabilityReport::compiled(Some(status)), + Ok(status) => { + let state = availability_state(status.availability()); + CapabilityReport::reported(state, status) + } + Err(_) => CapabilityReport::compiled(None), + } +} + +async fn sink_report(engine: &Engine) -> CapabilityReport<SinkStatus> { + let Some(sink) = engine.sink.as_deref() else { + return CapabilityReport::unsupported(); + }; + match sink.status().await { + Ok(status) if !status.is_configured() => CapabilityReport::compiled(Some(status)), + Ok(status) => { + let state = availability_state(status.availability()); + CapabilityReport::reported(state, status) + } + Err(_) => CapabilityReport::compiled(None), + } +} + +async fn signer_report(engine: &Engine) -> CapabilityReport<SignerStatus> { + let Some(signer) = engine.signer.as_deref() else { + return CapabilityReport::unsupported(); + }; + match signer.status().await { + Ok(status) => { + let state = match status.availability() { + SignerAvailability::Ready => SyncCapabilityState::Available, + SignerAvailability::Busy | SignerAvailability::AwaitingAuthentication => { + SyncCapabilityState::Degraded + } + SignerAvailability::Unavailable => SyncCapabilityState::Configured, + _ => SyncCapabilityState::Degraded, + }; + CapabilityReport::reported(state, status) + } + Err(_) => CapabilityReport::compiled(None), + } +} + +const fn availability_state(availability: Availability) -> SyncCapabilityState { + match availability { + Availability::Available => SyncCapabilityState::Available, + Availability::Degraded => SyncCapabilityState::Degraded, + Availability::Unavailable => SyncCapabilityState::Configured, + } +} + +fn aggregate_health( + storage: StorageStatus, + events: &EventStoreStatus, + source: &CapabilityReport<SourceStatus>, + sink: &CapabilityReport<SinkStatus>, + signer: &CapabilityReport<SignerStatus>, + projections: &[ProjectionReport], +) -> SyncHealth { + if matches!( + storage.shutdown(), + ShutdownState::Closing | ShutdownState::Closed + ) || storage.integrity().health() == IntegrityHealth::Corrupt + || events.health() == EventStoreHealth::Unavailable + { + return SyncHealth::Unavailable; + } + let capability_degraded = [source.state, sink.state, signer.state] + .into_iter() + .any(|state| { + matches!( + state, + SyncCapabilityState::Compiled + | SyncCapabilityState::Configured + | SyncCapabilityState::Degraded + ) + }); + let projection_degraded = projections.iter().any(|projection| { + !matches!( + projection.status.as_ref().map(ProjectionStatus::health), + Some(ProjectionHealth::Ready) + ) + }); + if storage.integrity().health() != IntegrityHealth::Healthy + || events.health() == EventStoreHealth::Degraded + || capability_degraded + || projection_degraded + { + SyncHealth::Degraded + } else { + SyncHealth::Healthy + } +} + +const fn storage_state( + storage: StorageStatus, + event_health: EventStoreHealth, +) -> SyncCapabilityState { + if matches!( + storage.shutdown(), + ShutdownState::Closing | ShutdownState::Closed + ) || matches!(storage.integrity().health(), IntegrityHealth::Corrupt) + || matches!(event_health, EventStoreHealth::Unavailable) + { + SyncCapabilityState::Configured + } else if matches!(storage.integrity().health(), IntegrityHealth::Healthy) + && matches!(event_health, EventStoreHealth::Available) + { + SyncCapabilityState::Available + } else { + SyncCapabilityState::Degraded + } +} diff --git a/crates/sync/tests/engine_composition.rs b/crates/sync/tests/engine_composition.rs @@ -1,7 +1,11 @@ use std::sync::Arc; +use futures_executor::block_on; +use radroots_protocol::runtime::v1::{OPERATION_SCHEMA_VERSION, SyncCapabilityState, SyncHealth}; use radroots_signing::{Error as SigningError, SignReceipt, SignRequest, Signer, SignerStatus}; -use radroots_storage::{Storage, event::SourceGeneration, memory::MemoryStorage}; +use radroots_storage::{ + Storage, event::SourceGeneration, memory::MemoryStorage, projection::ProjectionId, +}; use radroots_sync::{ Engine, policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId}, @@ -9,11 +13,13 @@ use radroots_sync::{ use radroots_transport::{ DeliveryReceipt, DeliveryRequest, Error as TransportError, EventSink, EventSource, FetchPage, FetchRequest, SinkStatus, SourceStatus, + capability::{Availability, Maturity, SinkCapabilities, SourceCapabilities}, }; struct MockSource; struct MockSink; struct MockSigner; +struct UnconfiguredSource; struct FixedClock; struct FixedIds; @@ -26,7 +32,16 @@ type TestDependencies = ( impl EventSource for MockSource { fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SourceStatus, TransportError>> { - Box::pin(async { unreachable!("composition does not poll source") }) + Box::pin(async { + Ok(SourceStatus::new( + radroots_transport::TransportId::NOSTR, + true, + Maturity::Stable, + Availability::Available, + SourceCapabilities::FETCH, + "ready", + )) + }) } fn fetch( @@ -37,9 +52,40 @@ impl EventSource for MockSource { } } +impl EventSource for UnconfiguredSource { + fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SourceStatus, TransportError>> { + Box::pin(async { + Ok(SourceStatus::new( + radroots_transport::TransportId::NOSTR, + false, + Maturity::Preview, + Availability::Available, + SourceCapabilities::FETCH, + "not configured", + )) + }) + } + + fn fetch( + &self, + _request: FetchRequest, + ) -> radroots_transport::BoxFuture<'_, Result<FetchPage, TransportError>> { + Box::pin(async { unreachable!("status does not fetch") }) + } +} + impl EventSink for MockSink { fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SinkStatus, TransportError>> { - Box::pin(async { unreachable!("composition does not poll sink") }) + Box::pin(async { + Ok(SinkStatus::new( + radroots_transport::TransportId::NOSTR, + true, + Maturity::Preview, + Availability::Degraded, + SinkCapabilities::DELIVER, + "degraded", + )) + }) } fn deliver( @@ -54,7 +100,7 @@ impl Signer for MockSigner { fn status( &self, ) -> radroots_signing::signer::BoxFuture<'_, Result<SignerStatus, SigningError>> { - Box::pin(async { unreachable!("composition does not poll signer") }) + Box::pin(async { Ok(SignerStatus::unavailable()) }) } fn sign( @@ -159,3 +205,53 @@ fn invalid_compositions_and_ambient_policy_inputs_fail_closed() { ); assert_eq!(SyncId::new([0; 16]), Err(Error::InvalidSyncId)); } + +#[test] +fn status_aggregates_typed_capability_and_protocol_reports() { + let (storage, clock, ids, deadlines) = dependencies(); + let full = Engine::builder(storage, clock, ids, deadlines) + .source(Arc::new(MockSource)) + .sink(Arc::new(MockSink)) + .signer(Arc::new(MockSigner)) + .build() + .expect("full engine"); + let projection = ProjectionId::parse("market-listings").expect("projection id"); + let status = block_on(full.status(std::slice::from_ref(&projection))).expect("sync status"); + assert_eq!(status.health(), SyncHealth::Degraded); + assert_eq!(status.source().state(), SyncCapabilityState::Available); + assert_eq!(status.sink().state(), SyncCapabilityState::Degraded); + assert_eq!(status.signer().state(), SyncCapabilityState::Configured); + assert_eq!(status.projections()[0].projection_id(), &projection); + assert!(status.projections()[0].status().is_none()); + let protocol = status.to_protocol(); + assert_eq!(protocol.schema_version, OPERATION_SCHEMA_VERSION); + assert_eq!(protocol.health, SyncHealth::Degraded); + assert_eq!(protocol.source, SyncCapabilityState::Available); + assert_eq!(protocol.sink, SyncCapabilityState::Degraded); + assert_eq!(protocol.signer, SyncCapabilityState::Configured); + assert_eq!(protocol.projections.untracked, 1); + + let (storage, clock, ids, deadlines) = dependencies(); + let sink_only = Engine::builder(storage, clock, ids, deadlines) + .sink(Arc::new(MockSink)) + .build() + .expect("sink engine"); + let status = block_on(sink_only.status(&[])).expect("sink-only status"); + assert_eq!(status.source().state(), SyncCapabilityState::Unsupported); + assert_eq!(status.signer().state(), SyncCapabilityState::Unsupported); + + let (storage, clock, ids, deadlines) = dependencies(); + let compiled = Engine::builder(storage, clock, ids, deadlines) + .source(Arc::new(UnconfiguredSource)) + .build() + .expect("unconfigured source engine"); + let status = block_on(compiled.status(&[])).expect("compiled status"); + assert_eq!(status.source().state(), SyncCapabilityState::Compiled); + assert!(status.source().status().is_some()); + assert_eq!(status.sink().state(), SyncCapabilityState::Unsupported); + + assert_eq!( + block_on(full.status(&[projection.clone(), projection])), + Err(Error::InvalidStatusRequest) + ); +} diff --git a/crates/sync/tests/push_enqueue.rs b/crates/sync/tests/push_enqueue.rs @@ -9,6 +9,7 @@ use std::{ use futures::{FutureExt, task::noop_waker_ref}; use futures_executor::block_on; use radroots_event::{EventDraft, SignedEvent, contract::AuthorRole, draft::SignedEventParts}; +use radroots_protocol::runtime::v1::SyncRetryDecision; use radroots_signing::{ Actor, Error as SigningError, SignReceipt, SignRequest, Signer, SignerStatus, actor::ActorSource, error::Kind as SigningErrorKind, request::CancellationPolicy, @@ -18,7 +19,7 @@ use radroots_storage::{ event::{EventQuery, EventQueryBounds, SourceGeneration}, journal::{IdempotencyKey, JournalStage, OperationInstanceId}, memory::MemoryStorage, - outbox::{LeaseOwner, OutboxStage, SatisfactionResult}, + outbox::{ClaimOutboxItems, LeaseId, LeaseOwner, OutboxStage, SatisfactionResult}, }; use radroots_sync::{ Engine, PushRequest, @@ -659,3 +660,83 @@ fn delivery_run_rejects_unbounded_claims() { Err(Error::InvalidDeliveryRequest) ); } + +#[test] +fn retry_decisions_are_passive_typed_and_deadline_aware() { + let signer = Arc::new(MockSigner::new(SignBehavior::Success { + completed_at_unix: 1_800_000_200, + })); + let (engine, storage) = setup_engine(signer); + let enqueued = block_on(engine.sign_and_enqueue(request(61, "wss://one.example"))) + .expect("enqueue retry decision plan"); + let now = enqueued.outbox().created_at_unix_ms() + 1; + assert_eq!( + engine.retry_decision(enqueued.outbox(), now), + Ok(SyncRetryDecision::Ready) + ); + let claimed = block_on(Outbox::claim( + &*storage, + ClaimOutboxItems::new( + LeaseOwner::parse("retry-decision-test").expect("owner"), + LeaseId::new([81; 16]).expect("lease seed"), + now, + now + 100, + 1, + ) + .expect("claim request"), + )) + .expect("claim") + .pop() + .expect("claimed plan"); + assert_eq!( + engine.retry_decision(claimed.record(), now + 1), + Ok(SyncRetryDecision::InFlightUntil { unix_ms: now + 100 }) + ); + let deferred = block_on(Outbox::release( + &*storage, + claimed.record().item_id(), + claimed.lease().id(), + claimed.record().revision(), + now + 2, + Some(now + 50), + )) + .expect("release with deferral"); + assert_eq!( + engine.retry_decision(&deferred, now + 3), + Ok(SyncRetryDecision::DeferredUntil { unix_ms: now + 50 }) + ); + assert_eq!( + engine.retry_decision(&deferred, now + 50), + Ok(SyncRetryDecision::Ready) + ); + assert_eq!( + engine.retry_decision(&deferred, deferred.request().deadline_unix_ms()), + Ok(SyncRetryDecision::Expired) + ); + assert_eq!( + engine.retry_decision(&deferred, 0), + Err(Error::ClockUnavailable) + ); + + let sink = Arc::new(ScriptedSink::new([ + DeliveryBehavior::Outcomes(vec![DeliveryOutcome::accepted()]), + DeliveryBehavior::Outcomes(vec![DeliveryOutcome::rejected()]), + ])); + let signer = Arc::new(MockSigner::new(SignBehavior::Success { + completed_at_unix: 1_800_000_200, + })); + let ((engine, _), _) = setup_engine_with_sink(signer, sink); + block_on(engine.sign_and_enqueue(request(62, "wss://one.example"))) + .expect("enqueue satisfied plan"); + block_on(engine.sign_and_enqueue(request(63, "wss://two.example"))) + .expect("enqueue exhausted plan"); + let delivered = block_on(engine.deliver_pending(delivery_run(82, 2))).expect("deliver plans"); + assert_eq!( + engine.retry_decision(delivered.outcomes()[0].as_ref().expect("satisfied"), now), + Ok(SyncRetryDecision::Satisfied) + ); + assert_eq!( + engine.retry_decision(delivered.outcomes()[1].as_ref().expect("exhausted"), now), + Ok(SyncRetryDecision::Exhausted) + ); +}