commit 85e78fea1feca5bd030b0156cae9fe570402c8a0
parent ac54364e484fb565e65e58581ec9bb2bb750e680
Author: triesap <tyson@radroots.org>
Date: Mon, 3 Aug 2026 07:52:22 +0000
sync: define the sync engine composition boundary
- Compose storage with independently optional source, sink, and signer SPIs.
- Inject bounded deadline, clock, and operation identity policies.
- Reject empty capability sets and signer-without-sink configurations.
- Prove source-only, sink-only, and full compositions without polling I/O.
Diffstat:
5 files changed, 482 insertions(+), 0 deletions(-)
diff --git a/crates/sync/Cargo.toml b/crates/sync/Cargo.toml
@@ -37,5 +37,8 @@ radroots_trade = { workspace = true, default-features = false }
radroots_transport = { workspace = true, default-features = false }
serde = { workspace = true, optional = true }
+[dev-dependencies]
+radroots_storage = { workspace = true, features = ["memory"] }
+
[lints]
workspace = true
diff --git a/crates/sync/src/engine.rs b/crates/sync/src/engine.rs
@@ -0,0 +1,78 @@
+use std::sync::Arc;
+
+use radroots_signing::Signer;
+use radroots_storage::Storage;
+use radroots_transport::{EventSink, EventSource};
+
+use crate::policy::{Clock, DeadlinePolicy, EngineBuilder, IdSource};
+
+/// Injected, executor-neutral synchronization composition boundary.
+#[derive(Clone)]
+pub struct Engine {
+ pub(crate) storage: Arc<dyn Storage>,
+ pub(crate) source: Option<Arc<dyn EventSource>>,
+ pub(crate) sink: Option<Arc<dyn EventSink>>,
+ pub(crate) signer: Option<Arc<dyn Signer>>,
+ pub(crate) clock: Arc<dyn Clock>,
+ pub(crate) ids: Arc<dyn IdSource>,
+ pub(crate) deadlines: DeadlinePolicy,
+}
+
+impl Engine {
+ /// Starts an explicit capability builder around required host policies.
+ pub fn builder(
+ storage: Arc<dyn Storage>,
+ clock: Arc<dyn Clock>,
+ ids: Arc<dyn IdSource>,
+ deadlines: DeadlinePolicy,
+ ) -> EngineBuilder {
+ EngineBuilder::new(storage, clock, ids, deadlines)
+ }
+
+ /// Returns the canonical storage capability.
+ pub fn storage(&self) -> &dyn Storage {
+ self.storage.as_ref()
+ }
+
+ /// Returns the configured event source, when pull is enabled.
+ pub fn source(&self) -> Option<&dyn EventSource> {
+ self.source.as_deref()
+ }
+
+ /// Returns the configured event sink, when delivery is enabled.
+ pub fn sink(&self) -> Option<&dyn EventSink> {
+ self.sink.as_deref()
+ }
+
+ /// Returns the configured signer, when outbound authoring is enabled.
+ pub fn signer(&self) -> Option<&dyn Signer> {
+ self.signer.as_deref()
+ }
+
+ /// Returns the injected clock policy.
+ pub fn clock(&self) -> &dyn Clock {
+ self.clock.as_ref()
+ }
+
+ /// Returns the injected operation identity source.
+ pub fn ids(&self) -> &dyn IdSource {
+ self.ids.as_ref()
+ }
+
+ /// Returns the bounded deadline policy.
+ pub const fn deadlines(&self) -> DeadlinePolicy {
+ self.deadlines
+ }
+}
+
+impl core::fmt::Debug for Engine {
+ fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
+ formatter
+ .debug_struct("Engine")
+ .field("source", &self.source.is_some())
+ .field("sink", &self.sink.is_some())
+ .field("signer", &self.signer.is_some())
+ .field("deadlines", &self.deadlines)
+ .finish_non_exhaustive()
+ }
+}
diff --git a/crates/sync/src/lib.rs b/crates/sync/src/lib.rs
@@ -1,8 +1,13 @@
//! Executor-neutral local-first synchronization orchestration.
+mod engine;
+
pub mod ingest;
pub mod policy;
pub mod projection;
pub mod pull;
pub mod push;
pub mod status;
+
+pub use engine::Engine;
+pub use policy::Error;
diff --git a/crates/sync/src/policy.rs b/crates/sync/src/policy.rs
@@ -1 +1,236 @@
//! Explicit clocks, identifiers, deadlines, and retry decisions.
+
+use std::sync::Arc;
+
+use radroots_signing::Signer;
+use radroots_storage::Storage;
+use radroots_transport::{EventSink, EventSource};
+
+use crate::Engine;
+
+const MAX_OPERATION_TIMEOUT_MS: u64 = 86_400_000;
+
+/// Sync operation class used for identity and deadline policy.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
+#[derive(Clone, Copy, Debug, Eq, PartialEq)]
+#[non_exhaustive]
+pub enum OperationKind {
+ Pull,
+ Sign,
+ Deliver,
+}
+
+/// Opaque host-generated identity for one synchronization operation.
+#[cfg_attr(feature = "serde", derive(serde::Serialize))]
+#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
+pub struct SyncId([u8; 16]);
+
+impl SyncId {
+ pub const fn new(bytes: [u8; 16]) -> Result<Self, Error> {
+ let mut index = 0;
+ while index < bytes.len() {
+ if bytes[index] != 0 {
+ return Ok(Self(bytes));
+ }
+ index += 1;
+ }
+ Err(Error::InvalidSyncId)
+ }
+
+ pub const fn as_bytes(&self) -> &[u8; 16] {
+ &self.0
+ }
+}
+
+/// Host clock used instead of reading ambient time inside orchestration.
+pub trait Clock: Send + Sync {
+ fn now_unix_ms(&self) -> Result<u64, Error>;
+}
+
+/// Host identity source used instead of ambient randomness or global counters.
+pub trait IdSource: Send + Sync {
+ fn next_id(&self, operation: OperationKind) -> Result<SyncId, Error>;
+}
+
+/// Bounded time budgets applied to individual orchestration calls.
+#[cfg_attr(feature = "serde", derive(serde::Serialize))]
+#[derive(Clone, Copy, Debug, Eq, PartialEq)]
+pub struct DeadlinePolicy {
+ pull_timeout_ms: u64,
+ sign_timeout_ms: u64,
+ delivery_timeout_ms: u64,
+}
+
+impl DeadlinePolicy {
+ pub const fn new(
+ pull_timeout_ms: u64,
+ sign_timeout_ms: u64,
+ delivery_timeout_ms: u64,
+ ) -> Result<Self, Error> {
+ if !valid_timeout(pull_timeout_ms)
+ || !valid_timeout(sign_timeout_ms)
+ || !valid_timeout(delivery_timeout_ms)
+ {
+ return Err(Error::InvalidDeadlinePolicy);
+ }
+ Ok(Self {
+ pull_timeout_ms,
+ sign_timeout_ms,
+ delivery_timeout_ms,
+ })
+ }
+
+ pub const fn timeout_ms(self, operation: OperationKind) -> u64 {
+ match operation {
+ OperationKind::Pull => self.pull_timeout_ms,
+ OperationKind::Sign => self.sign_timeout_ms,
+ OperationKind::Deliver => self.delivery_timeout_ms,
+ }
+ }
+
+ pub fn deadline_unix_ms(
+ self,
+ operation: OperationKind,
+ now_unix_ms: u64,
+ ) -> Result<u64, Error> {
+ if now_unix_ms == 0 {
+ return Err(Error::ClockUnavailable);
+ }
+ now_unix_ms
+ .checked_add(self.timeout_ms(operation))
+ .ok_or(Error::DeadlineOverflow)
+ }
+}
+
+const fn valid_timeout(value: u64) -> bool {
+ value != 0 && value <= MAX_OPERATION_TIMEOUT_MS
+}
+
+/// Builder for an [`Engine`] with explicit optional transport capabilities.
+pub struct EngineBuilder {
+ storage: Arc<dyn Storage>,
+ source: Option<Arc<dyn EventSource>>,
+ sink: Option<Arc<dyn EventSink>>,
+ signer: Option<Arc<dyn Signer>>,
+ clock: Arc<dyn Clock>,
+ ids: Arc<dyn IdSource>,
+ deadlines: DeadlinePolicy,
+}
+
+impl EngineBuilder {
+ pub(crate) fn new(
+ storage: Arc<dyn Storage>,
+ clock: Arc<dyn Clock>,
+ ids: Arc<dyn IdSource>,
+ deadlines: DeadlinePolicy,
+ ) -> Self {
+ Self {
+ storage,
+ source: None,
+ sink: None,
+ signer: None,
+ clock,
+ ids,
+ deadlines,
+ }
+ }
+
+ #[must_use]
+ pub fn source(mut self, source: Arc<dyn EventSource>) -> Self {
+ self.source = Some(source);
+ self
+ }
+
+ #[must_use]
+ pub fn sink(mut self, sink: Arc<dyn EventSink>) -> Self {
+ self.sink = Some(sink);
+ self
+ }
+
+ #[must_use]
+ pub fn signer(mut self, signer: Arc<dyn Signer>) -> Self {
+ self.signer = Some(signer);
+ self
+ }
+
+ pub fn build(self) -> Result<Engine, Error> {
+ if self.signer.is_some() && self.sink.is_none() {
+ return Err(Error::SignerWithoutSink);
+ }
+ if self.source.is_none() && self.sink.is_none() {
+ return Err(Error::MissingTransportCapability);
+ }
+ Ok(Engine {
+ storage: self.storage,
+ source: self.source,
+ sink: self.sink,
+ signer: self.signer,
+ clock: self.clock,
+ ids: self.ids,
+ deadlines: self.deadlines,
+ })
+ }
+}
+
+/// Sync composition and host-policy error.
+#[derive(Clone, Copy, Debug, Eq, PartialEq)]
+#[non_exhaustive]
+pub enum Error {
+ InvalidSyncId,
+ InvalidDeadlinePolicy,
+ ClockUnavailable,
+ DeadlineOverflow,
+ MissingTransportCapability,
+ SignerWithoutSink,
+}
+
+impl core::fmt::Display for Error {
+ fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
+ formatter.write_str(match self {
+ Self::InvalidSyncId => "sync identity must not be all zero",
+ Self::InvalidDeadlinePolicy => "sync deadline policy is outside its bounds",
+ Self::ClockUnavailable => "sync clock did not provide a valid timestamp",
+ Self::DeadlineOverflow => "sync deadline overflowed",
+ Self::MissingTransportCapability => "sync engine requires a source or sink",
+ Self::SignerWithoutSink => "sync signer requires a sink",
+ })
+ }
+}
+
+impl std::error::Error for Error {}
+
+#[cfg(feature = "serde")]
+impl<'de> serde::Deserialize<'de> for SyncId {
+ fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
+ where
+ D: serde::Deserializer<'de>,
+ {
+ let bytes = <[u8; 16] as serde::Deserialize>::deserialize(deserializer)?;
+ Self::new(bytes).map_err(serde::de::Error::custom)
+ }
+}
+
+#[cfg(feature = "serde")]
+impl<'de> serde::Deserialize<'de> for DeadlinePolicy {
+ fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
+ where
+ D: serde::Deserializer<'de>,
+ {
+ #[derive(serde::Deserialize)]
+ #[serde(deny_unknown_fields)]
+ struct Wire {
+ pull_timeout_ms: u64,
+ sign_timeout_ms: u64,
+ delivery_timeout_ms: u64,
+ }
+
+ let wire = Wire::deserialize(deserializer)?;
+ Self::new(
+ wire.pull_timeout_ms,
+ wire.sign_timeout_ms,
+ wire.delivery_timeout_ms,
+ )
+ .map_err(serde::de::Error::custom)
+ }
+}
diff --git a/crates/sync/tests/engine_composition.rs b/crates/sync/tests/engine_composition.rs
@@ -0,0 +1,161 @@
+use std::sync::Arc;
+
+use radroots_signing::{Error as SigningError, SignReceipt, SignRequest, Signer, SignerStatus};
+use radroots_storage::{Storage, event::SourceGeneration, memory::MemoryStorage};
+use radroots_sync::{
+ Engine,
+ policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId},
+};
+use radroots_transport::{
+ DeliveryReceipt, DeliveryRequest, Error as TransportError, EventSink, EventSource, FetchPage,
+ FetchRequest, SinkStatus, SourceStatus,
+};
+
+struct MockSource;
+struct MockSink;
+struct MockSigner;
+struct FixedClock;
+struct FixedIds;
+
+type TestDependencies = (
+ Arc<dyn Storage>,
+ Arc<dyn Clock>,
+ Arc<dyn IdSource>,
+ DeadlinePolicy,
+);
+
+impl EventSource for MockSource {
+ fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SourceStatus, TransportError>> {
+ Box::pin(async { unreachable!("composition does not poll source") })
+ }
+
+ fn fetch(
+ &self,
+ _request: FetchRequest,
+ ) -> radroots_transport::BoxFuture<'_, Result<FetchPage, TransportError>> {
+ Box::pin(async { unreachable!("composition 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") })
+ }
+
+ fn deliver(
+ &self,
+ _request: DeliveryRequest,
+ ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, TransportError>> {
+ Box::pin(async { unreachable!("composition does not deliver") })
+ }
+}
+
+impl Signer for MockSigner {
+ fn status(
+ &self,
+ ) -> radroots_signing::signer::BoxFuture<'_, Result<SignerStatus, SigningError>> {
+ Box::pin(async { unreachable!("composition does not poll signer") })
+ }
+
+ fn sign(
+ &self,
+ _request: SignRequest,
+ ) -> radroots_signing::signer::BoxFuture<'_, Result<SignReceipt, SigningError>> {
+ Box::pin(async { unreachable!("composition does not sign") })
+ }
+}
+
+impl Clock for FixedClock {
+ fn now_unix_ms(&self) -> Result<u64, Error> {
+ Ok(1_700_000_000_000)
+ }
+}
+
+impl IdSource for FixedIds {
+ fn next_id(&self, operation: OperationKind) -> Result<SyncId, Error> {
+ let byte = match operation {
+ OperationKind::Pull => 1,
+ OperationKind::Sign => 2,
+ OperationKind::Deliver => 3,
+ _ => 4,
+ };
+ SyncId::new([byte; 16])
+ }
+}
+
+fn dependencies() -> TestDependencies {
+ let generation = SourceGeneration::new([7; 32]).expect("generation");
+ (
+ Arc::new(MemoryStorage::new(generation)),
+ Arc::new(FixedClock),
+ Arc::new(FixedIds),
+ DeadlinePolicy::new(10_000, 20_000, 30_000).expect("deadlines"),
+ )
+}
+
+#[test]
+fn source_only_sink_only_and_full_compositions_are_explicit() {
+ let (storage, clock, ids, deadlines) = dependencies();
+ let source_only = Engine::builder(storage, clock, ids, deadlines)
+ .source(Arc::new(MockSource))
+ .build()
+ .expect("source engine");
+ assert!(source_only.source().is_some());
+ assert!(source_only.sink().is_none());
+ assert!(source_only.signer().is_none());
+
+ let (storage, clock, ids, deadlines) = dependencies();
+ let sink_only = Engine::builder(storage, clock, ids, deadlines)
+ .sink(Arc::new(MockSink))
+ .build()
+ .expect("sink engine");
+ assert!(sink_only.source().is_none());
+ assert!(sink_only.sink().is_some());
+ assert!(sink_only.signer().is_none());
+
+ 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");
+ assert!(full.source().is_some());
+ assert!(full.sink().is_some());
+ assert!(full.signer().is_some());
+ assert_eq!(
+ full.clock().now_unix_ms().expect("clock"),
+ 1_700_000_000_000
+ );
+ assert_eq!(
+ full.deadlines()
+ .deadline_unix_ms(OperationKind::Deliver, 1_000)
+ .expect("deadline"),
+ 31_000
+ );
+}
+
+#[test]
+fn invalid_compositions_and_ambient_policy_inputs_fail_closed() {
+ let (storage, clock, ids, deadlines) = dependencies();
+ assert_eq!(
+ Engine::builder(storage, clock, ids, deadlines)
+ .build()
+ .expect_err("missing transport"),
+ Error::MissingTransportCapability
+ );
+
+ let (storage, clock, ids, deadlines) = dependencies();
+ assert_eq!(
+ Engine::builder(storage, clock, ids, deadlines)
+ .signer(Arc::new(MockSigner))
+ .build()
+ .expect_err("signer without sink"),
+ Error::SignerWithoutSink
+ );
+ assert_eq!(
+ DeadlinePolicy::new(0, 1, 1),
+ Err(Error::InvalidDeadlinePolicy)
+ );
+ assert_eq!(SyncId::new([0; 16]), Err(Error::InvalidSyncId));
+}