commit a82210b769e6405ee8e74c4a31221cd89f36352e
parent 991d63b9906070a827e235d3e7824c80508d494e
Author: triesap <tyson@radroots.org>
Date: Mon, 3 Aug 2026 08:09:54 +0000
sync: implement bounded pull
- add validated pull bounds, cursors, deadlines, and root receipt exports
- fetch deterministic pages and feed observations through canonical ingest
- aggregate final target states and resumable termination evidence
- verify pagination, failure, cancellation, and hard-limit behavior
Diffstat:
5 files changed, 603 insertions(+), 0 deletions(-)
diff --git a/crates/sync/src/ingest.rs b/crates/sync/src/ingest.rs
@@ -107,6 +107,7 @@ impl IngestReceipt {
}
/// Ordered independent outcomes for one bounded caller-supplied batch.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct IngestBatchReceipt {
outcomes: Vec<Result<IngestReceipt, Error>>,
diff --git a/crates/sync/src/lib.rs b/crates/sync/src/lib.rs
@@ -11,3 +11,4 @@ pub mod status;
pub use engine::Engine;
pub use policy::Error;
+pub use pull::{PullReceipt, PullRequest};
diff --git a/crates/sync/src/policy.rs b/crates/sync/src/policy.rs
@@ -178,6 +178,8 @@ impl EngineBuilder {
}
/// Sync composition and host-policy error.
+#[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 Error {
@@ -192,6 +194,9 @@ pub enum Error {
StorageConflict,
StorageFailed,
InvalidIngestReceipt,
+ InvalidPullRequest,
+ MissingSource,
+ InvalidSourcePage,
}
impl core::fmt::Display for Error {
@@ -208,6 +213,9 @@ impl core::fmt::Display for Error {
Self::StorageConflict => "sync input conflicts with durable storage state",
Self::StorageFailed => "sync storage operation failed",
Self::InvalidIngestReceipt => "sync storage returned an invalid ingest receipt",
+ Self::InvalidPullRequest => "sync pull request is outside its bounds",
+ Self::MissingSource => "sync engine has no event source",
+ Self::InvalidSourcePage => "sync source returned an invalid page",
})
}
}
diff --git a/crates/sync/src/pull.rs b/crates/sync/src/pull.rs
@@ -1 +1,263 @@
//! Bounded source pagination and ingestion.
+
+use radroots_transport::{
+ FetchRequest,
+ outcome::FetchTargetOutcome,
+ source::{FETCH_PAGE_MAX_EVENTS, FetchBounds, FetchCursor, NextPage},
+ target::TargetSet,
+};
+
+use crate::{
+ Engine,
+ ingest::{AdmissionPolicy, IngestReceipt},
+ policy::{Error, OperationKind, SyncId},
+};
+
+/// Hard maximum number of source pages in one explicit pull call.
+pub const PULL_MAX_PAGES: u16 = 1_000;
+
+/// Caller-owned bounds and continuation state for one pull operation.
+#[cfg_attr(feature = "serde", derive(serde::Serialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct PullRequest {
+ targets: TargetSet,
+ page_limit: u16,
+ max_pages: u16,
+ cursor: Option<FetchCursor>,
+}
+
+impl PullRequest {
+ pub fn new(targets: TargetSet, page_limit: u16, max_pages: u16) -> Result<Self, Error> {
+ if page_limit == 0
+ || page_limit > FETCH_PAGE_MAX_EVENTS
+ || max_pages == 0
+ || max_pages > PULL_MAX_PAGES
+ {
+ return Err(Error::InvalidPullRequest);
+ }
+ Ok(Self {
+ targets,
+ page_limit,
+ max_pages,
+ cursor: None,
+ })
+ }
+
+ #[must_use]
+ pub fn with_cursor(mut self, cursor: FetchCursor) -> Self {
+ self.cursor = Some(cursor);
+ self
+ }
+
+ pub const fn targets(&self) -> &TargetSet {
+ &self.targets
+ }
+
+ pub const fn page_limit(&self) -> u16 {
+ self.page_limit
+ }
+
+ pub const fn max_pages(&self) -> u16 {
+ self.max_pages
+ }
+
+ pub const fn cursor(&self) -> Option<&FetchCursor> {
+ self.cursor.as_ref()
+ }
+}
+
+/// Deterministic reason why a bounded pull returned control to its caller.
+#[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 PullTermination {
+ Complete,
+ PageLimit,
+ Deadline,
+ Cancelled,
+ SourceFailed,
+}
+
+/// Normalized receipt retaining every ingest outcome and final target state.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct PullReceipt {
+ sync_id: SyncId,
+ deadline_unix_ms: u64,
+ pages_fetched: u16,
+ events_observed: usize,
+ ingest_outcomes: Vec<Result<IngestReceipt, Error>>,
+ target_outcomes: Vec<FetchTargetOutcome>,
+ termination: PullTermination,
+ resume_from: Option<FetchCursor>,
+}
+
+impl PullReceipt {
+ pub const fn sync_id(&self) -> SyncId {
+ self.sync_id
+ }
+
+ pub const fn deadline_unix_ms(&self) -> u64 {
+ self.deadline_unix_ms
+ }
+
+ pub const fn pages_fetched(&self) -> u16 {
+ self.pages_fetched
+ }
+
+ pub const fn events_observed(&self) -> usize {
+ self.events_observed
+ }
+
+ pub fn ingest_outcomes(&self) -> &[Result<IngestReceipt, Error>] {
+ self.ingest_outcomes.as_slice()
+ }
+
+ pub fn target_outcomes(&self) -> &[FetchTargetOutcome] {
+ self.target_outcomes.as_slice()
+ }
+
+ pub const fn termination(&self) -> PullTermination {
+ self.termination
+ }
+
+ pub const fn resume_from(&self) -> Option<&FetchCursor> {
+ self.resume_from.as_ref()
+ }
+}
+
+impl Engine {
+ /// Fetches and ingests at most the caller-bounded number of source pages.
+ pub async fn pull(
+ &self,
+ request: PullRequest,
+ admission: &dyn AdmissionPolicy,
+ ) -> Result<PullReceipt, Error> {
+ let source = self.source.as_deref().ok_or(Error::MissingSource)?;
+ let sync_id = self.ids.next_id(OperationKind::Pull)?;
+ let started_at = self.clock.now_unix_ms()?;
+ let deadline_unix_ms = self
+ .deadlines
+ .deadline_unix_ms(OperationKind::Pull, started_at)?;
+ let bounds = FetchBounds::new(request.page_limit, deadline_unix_ms)
+ .map_err(|_| Error::InvalidPullRequest)?;
+ let mut receipt = PullReceipt {
+ sync_id,
+ deadline_unix_ms,
+ pages_fetched: 0,
+ events_observed: 0,
+ ingest_outcomes: Vec::new(),
+ target_outcomes: Vec::new(),
+ termination: PullTermination::Complete,
+ resume_from: request.cursor.clone(),
+ };
+ let mut cursor = request.cursor;
+
+ for page_index in 0..request.max_pages {
+ if page_index != 0 && self.clock.now_unix_ms()? >= deadline_unix_ms {
+ receipt.termination = PullTermination::Deadline;
+ receipt.resume_from = cursor;
+ return Ok(receipt);
+ }
+ let mut fetch = FetchRequest::new(
+ fetch_request_id(sync_id, page_index),
+ request.targets.clone(),
+ bounds,
+ )
+ .map_err(|_| Error::InvalidPullRequest)?;
+ if let Some(current) = cursor.clone() {
+ fetch = fetch.with_cursor(current);
+ }
+ let page = match source.fetch(fetch.clone()).await {
+ Ok(page) => page,
+ Err(_) => {
+ receipt.termination = PullTermination::SourceFailed;
+ receipt.resume_from = cursor;
+ return Ok(receipt);
+ }
+ };
+ page.validate_for_request(&fetch)
+ .map_err(|_| Error::InvalidSourcePage)?;
+ receipt.pages_fetched += 1;
+ receipt.events_observed += page.events().len();
+ merge_target_outcomes(&mut receipt.target_outcomes, page.target_outcomes());
+ let outcomes = self.ingest_batch(page.events().to_vec(), admission).await;
+ receipt
+ .ingest_outcomes
+ .extend_from_slice(outcomes.outcomes());
+
+ match page.next_page() {
+ NextPage::Complete => {
+ receipt.termination = PullTermination::Complete;
+ receipt.resume_from = None;
+ return Ok(receipt);
+ }
+ NextPage::Cancelled { resume_from } => {
+ receipt.termination = PullTermination::Cancelled;
+ receipt.resume_from = resume_from.clone();
+ return Ok(receipt);
+ }
+ NextPage::Cursor(next) => {
+ cursor = Some(next.clone());
+ if page_index + 1 == request.max_pages {
+ receipt.termination = PullTermination::PageLimit;
+ receipt.resume_from = cursor;
+ return Ok(receipt);
+ }
+ }
+ }
+ }
+ unreachable!("validated pull requests execute at least one page")
+ }
+}
+
+fn merge_target_outcomes(current: &mut Vec<FetchTargetOutcome>, page: &[FetchTargetOutcome]) {
+ for outcome in page {
+ if let Some(existing) = current
+ .iter_mut()
+ .find(|existing| existing.target() == outcome.target())
+ {
+ *existing = outcome.clone();
+ } else {
+ current.push(outcome.clone());
+ }
+ }
+}
+
+fn fetch_request_id(sync_id: SyncId, page_index: u16) -> String {
+ const HEX: &[u8; 16] = b"0123456789abcdef";
+ let mut value = String::with_capacity(48);
+ value.push_str("sync-");
+ for byte in sync_id.as_bytes() {
+ value.push(HEX[(byte >> 4) as usize] as char);
+ value.push(HEX[(byte & 0x0f) as usize] as char);
+ }
+ value.push('-');
+ value.push_str(page_index.to_string().as_str());
+ value
+}
+
+#[cfg(feature = "serde")]
+impl<'de> serde::Deserialize<'de> for PullRequest {
+ fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
+ where
+ D: serde::Deserializer<'de>,
+ {
+ #[derive(serde::Deserialize)]
+ #[serde(deny_unknown_fields)]
+ struct Wire {
+ targets: TargetSet,
+ page_limit: u16,
+ max_pages: u16,
+ cursor: Option<FetchCursor>,
+ }
+
+ let wire = Wire::deserialize(deserializer)?;
+ let mut request = Self::new(wire.targets, wire.page_limit, wire.max_pages)
+ .map_err(serde::de::Error::custom)?;
+ if let Some(cursor) = wire.cursor {
+ request = request.with_cursor(cursor);
+ }
+ Ok(request)
+ }
+}
diff --git a/crates/sync/tests/pull.rs b/crates/sync/tests/pull.rs
@@ -0,0 +1,331 @@
+use std::{
+ collections::VecDeque,
+ sync::{Arc, Mutex},
+};
+
+use futures_executor::block_on;
+use radroots_event::{SignedEvent, draft::SignedEventParts};
+use radroots_storage::{Storage, event::SourceGeneration, memory::MemoryStorage};
+use radroots_sync::{
+ Engine, PullRequest,
+ ingest::RegistryPolicy,
+ policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId},
+ pull::{PULL_MAX_PAGES, PullTermination},
+};
+use radroots_transport::{
+ Error as TransportError, EventSource, FetchPage, FetchRequest, SourceStatus, Target, TargetSet,
+ TransportId,
+ outcome::{FetchTargetOutcome, FetchTargetState},
+ source::{EventProvenance, FetchCursor, NextPage, ObservedEvent},
+};
+
+const EVENT_ID: &str = "762bee187e9e645b81ec26ade05a69b5e8398caf527be8de0d9a45311ed0c7a0";
+const PUBKEY: &str = "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df";
+const SIGNATURE: &str = "4290da0bb6422986647bc8cd5f63bd52d49f41e7b665d3b47105b8109183e8d596f322c531d4061df53e1d2b70fda12d5d1c14f3720d7a56d9d0a03746af5109";
+const CONTENT: &str = "{\"display_name\":\"Moss Street Farm\",\"bot\":false,\"website\":\"https://mossstreet.example\",\"picture\":42}";
+
+enum Response {
+ Page {
+ events: Vec<ObservedEvent>,
+ state: FetchTargetState,
+ next: NextPage,
+ },
+ Failure,
+}
+
+#[derive(Clone, Debug, Eq, PartialEq)]
+struct RequestEvidence {
+ cursor: Option<String>,
+ limit: u16,
+ deadline_unix_ms: u64,
+}
+
+struct ScriptedSource {
+ responses: Mutex<VecDeque<Response>>,
+ requests: Mutex<Vec<RequestEvidence>>,
+}
+
+impl ScriptedSource {
+ fn new(responses: Vec<Response>) -> Self {
+ Self {
+ responses: Mutex::new(responses.into()),
+ requests: Mutex::new(Vec::new()),
+ }
+ }
+
+ fn requests(&self) -> Vec<RequestEvidence> {
+ self.requests.lock().expect("requests").clone()
+ }
+}
+
+impl EventSource for ScriptedSource {
+ fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SourceStatus, TransportError>> {
+ Box::pin(async { unreachable!("pull does not inspect source status") })
+ }
+
+ fn fetch(
+ &self,
+ request: FetchRequest,
+ ) -> radroots_transport::BoxFuture<'_, Result<FetchPage, TransportError>> {
+ Box::pin(async move {
+ self.requests
+ .lock()
+ .expect("requests")
+ .push(RequestEvidence {
+ cursor: request.cursor().map(|cursor| cursor.as_str().to_owned()),
+ limit: request.bounds().limit(),
+ deadline_unix_ms: request.bounds().deadline_unix_ms(),
+ });
+ match self
+ .responses
+ .lock()
+ .expect("responses")
+ .pop_front()
+ .expect("scripted response")
+ {
+ Response::Failure => Err(TransportError::UnsupportedOperation),
+ Response::Page {
+ events,
+ state,
+ next,
+ } => {
+ let target = request.target_set().targets()[0].fingerprint().clone();
+ FetchPage::for_request(
+ &request,
+ events,
+ vec![FetchTargetOutcome::new(target, state)],
+ next,
+ )
+ }
+ }
+ })
+ }
+}
+
+struct FixedClock(u64);
+
+impl Clock for FixedClock {
+ fn now_unix_ms(&self) -> Result<u64, Error> {
+ Ok(self.0)
+ }
+}
+
+struct DeadlineClock(Mutex<VecDeque<u64>>);
+
+impl Clock for DeadlineClock {
+ fn now_unix_ms(&self) -> Result<u64, Error> {
+ self.0
+ .lock()
+ .expect("clock")
+ .pop_front()
+ .ok_or(Error::ClockUnavailable)
+ }
+}
+
+struct SequenceIds(Mutex<u8>);
+
+impl IdSource for SequenceIds {
+ fn next_id(&self, _operation: OperationKind) -> Result<SyncId, Error> {
+ let mut next = self.0.lock().expect("ids");
+ let value = *next;
+ *next = next.checked_add(1).ok_or(Error::InvalidSyncId)?;
+ SyncId::new([value; 16])
+ }
+}
+
+fn target() -> Target {
+ Target::new(TransportId::NOSTR, "wss://relay.example").expect("target")
+}
+
+fn targets() -> TargetSet {
+ TargetSet::new(vec![target()]).expect("target set")
+}
+
+fn signed_event(signature: &str) -> SignedEvent {
+ let raw_json = format!(
+ "{{\"id\":\"{EVENT_ID}\",\"pubkey\":\"{PUBKEY}\",\"created_at\":1800000100,\"kind\":0,\"tags\":[],\"content\":{content:?},\"sig\":\"{signature}\"}}",
+ content = CONTENT,
+ );
+ SignedEvent::new(SignedEventParts {
+ id: EVENT_ID.to_owned(),
+ pubkey: PUBKEY.to_owned(),
+ created_at: 1_800_000_100,
+ kind: 0,
+ tags: vec![],
+ content: CONTENT.to_owned(),
+ sig: signature.to_owned(),
+ raw_json,
+ })
+ .expect("ID-valid event")
+}
+
+fn observed(signature: &str, observed_at: u64) -> ObservedEvent {
+ let target = target();
+ ObservedEvent::new(
+ signed_event(signature),
+ EventProvenance::new(
+ TransportId::NOSTR,
+ target.fingerprint().clone(),
+ observed_at,
+ )
+ .expect("provenance"),
+ )
+}
+
+fn engine(source: Arc<ScriptedSource>, clock: Arc<dyn Clock>, timeout_ms: u64) -> Engine {
+ let storage: Arc<dyn Storage> = Arc::new(MemoryStorage::new(
+ SourceGeneration::new([8; 32]).expect("generation"),
+ ));
+ Engine::builder(
+ storage,
+ clock,
+ Arc::new(SequenceIds(Mutex::new(1))),
+ DeadlinePolicy::new(timeout_ms, 10, 10).expect("deadlines"),
+ )
+ .source(source)
+ .build()
+ .expect("engine")
+}
+
+#[test]
+fn single_and_multiple_pages_propagate_cursor_deadline_and_ingest_results() {
+ let single_source = Arc::new(ScriptedSource::new(vec![Response::Page {
+ events: vec![observed(SIGNATURE, 1)],
+ state: FetchTargetState::Complete,
+ next: NextPage::Complete,
+ }]));
+ let single = engine(single_source.clone(), Arc::new(FixedClock(100)), 50);
+ let receipt = block_on(single.pull(
+ PullRequest::new(targets(), 20, 1).expect("request"),
+ &RegistryPolicy::visible(),
+ ))
+ .expect("pull");
+ assert_eq!(receipt.termination(), PullTermination::Complete);
+ assert_eq!(receipt.pages_fetched(), 1);
+ assert_eq!(receipt.events_observed(), 1);
+ assert!(receipt.ingest_outcomes()[0].is_ok());
+ assert_eq!(single_source.requests()[0].deadline_unix_ms, 150);
+
+ let next = FetchCursor::parse("page-2").expect("cursor");
+ let multiple_source = Arc::new(ScriptedSource::new(vec![
+ Response::Page {
+ events: vec![],
+ state: FetchTargetState::Partial,
+ next: NextPage::Cursor(next.clone()),
+ },
+ Response::Page {
+ events: vec![observed(SIGNATURE, 2)],
+ state: FetchTargetState::Complete,
+ next: NextPage::Complete,
+ },
+ ]));
+ let multiple = engine(multiple_source.clone(), Arc::new(FixedClock(200)), 50);
+ let receipt = block_on(
+ multiple.pull(
+ PullRequest::new(targets(), 10, 2)
+ .expect("request")
+ .with_cursor(FetchCursor::parse("starting").expect("initial cursor")),
+ &RegistryPolicy::visible(),
+ ),
+ )
+ .expect("pull");
+ assert_eq!(receipt.pages_fetched(), 2);
+ assert_eq!(receipt.termination(), PullTermination::Complete);
+ assert_eq!(
+ receipt.target_outcomes()[0].state(),
+ FetchTargetState::Complete
+ );
+ let requests = multiple_source.requests();
+ assert_eq!(requests[0].cursor.as_deref(), Some("starting"));
+ assert_eq!(requests[1].cursor.as_deref(), Some(next.as_str()));
+ assert_eq!(requests[0].deadline_unix_ms, requests[1].deadline_unix_ms);
+}
+
+#[test]
+fn source_failure_and_cancelled_page_return_resumable_partial_receipts() {
+ let cursor = FetchCursor::parse("resume").expect("cursor");
+ let source = Arc::new(ScriptedSource::new(vec![
+ Response::Page {
+ events: vec![observed(SIGNATURE, 1)],
+ state: FetchTargetState::Partial,
+ next: NextPage::Cursor(cursor.clone()),
+ },
+ Response::Failure,
+ ]));
+ let pull = engine(source, Arc::new(FixedClock(100)), 50);
+ let receipt = block_on(pull.pull(
+ PullRequest::new(targets(), 10, 3).expect("request"),
+ &RegistryPolicy::visible(),
+ ))
+ .expect("partial receipt");
+ assert_eq!(receipt.termination(), PullTermination::SourceFailed);
+ assert_eq!(receipt.pages_fetched(), 1);
+ assert_eq!(
+ receipt.resume_from().map(FetchCursor::as_str),
+ Some("resume")
+ );
+ assert!(receipt.ingest_outcomes()[0].is_ok());
+
+ let cancelled_from = FetchCursor::parse("cancelled-at").expect("cursor");
+ let source = Arc::new(ScriptedSource::new(vec![Response::Page {
+ events: vec![],
+ state: FetchTargetState::Cancelled,
+ next: NextPage::Cancelled {
+ resume_from: Some(cancelled_from.clone()),
+ },
+ }]));
+ let pull = engine(source, Arc::new(FixedClock(100)), 50);
+ let receipt = block_on(pull.pull(
+ PullRequest::new(targets(), 10, 3).expect("request"),
+ &RegistryPolicy::visible(),
+ ))
+ .expect("cancelled receipt");
+ assert_eq!(receipt.termination(), PullTermination::Cancelled);
+ assert_eq!(
+ receipt.resume_from().map(FetchCursor::as_str),
+ Some(cancelled_from.as_str())
+ );
+}
+
+#[test]
+fn page_and_deadline_limits_stop_without_hidden_fetches() {
+ assert_eq!(
+ PullRequest::new(targets(), 0, 1),
+ Err(Error::InvalidPullRequest)
+ );
+ assert_eq!(
+ PullRequest::new(targets(), 1, PULL_MAX_PAGES + 1),
+ Err(Error::InvalidPullRequest)
+ );
+
+ let cursor = FetchCursor::parse("more").expect("cursor");
+ let page_limited_source = Arc::new(ScriptedSource::new(vec![Response::Page {
+ events: vec![],
+ state: FetchTargetState::Partial,
+ next: NextPage::Cursor(cursor.clone()),
+ }]));
+ let pull = engine(page_limited_source.clone(), Arc::new(FixedClock(100)), 10);
+ let receipt = block_on(pull.pull(
+ PullRequest::new(targets(), 1, 1).expect("request"),
+ &RegistryPolicy::visible(),
+ ))
+ .expect("limited receipt");
+ assert_eq!(receipt.termination(), PullTermination::PageLimit);
+ assert_eq!(page_limited_source.requests().len(), 1);
+
+ let deadline_source = Arc::new(ScriptedSource::new(vec![Response::Page {
+ events: vec![],
+ state: FetchTargetState::Partial,
+ next: NextPage::Cursor(cursor),
+ }]));
+ let clock = Arc::new(DeadlineClock(Mutex::new(VecDeque::from([100, 110]))));
+ let pull = engine(deadline_source.clone(), clock, 10);
+ let receipt = block_on(pull.pull(
+ PullRequest::new(targets(), 1, 2).expect("request"),
+ &RegistryPolicy::visible(),
+ ))
+ .expect("deadline receipt");
+ assert_eq!(receipt.termination(), PullTermination::Deadline);
+ assert_eq!(receipt.deadline_unix_ms(), 110);
+ assert_eq!(deadline_source.requests().len(), 1);
+}