commit ab084e6943b24e24cb7408f05fbb6241ae8f8c83
parent a64fe38236edd357c9ad4d1b055b1bf6af4bd29e
Author: triesap <tyson@radroots.org>
Date: Wed, 9 Sep 2026 23:15:49 +0000
sync: Preserve bounded target evidence across pull pages
- Retain incomplete and missing outcomes through later successful pages
- Bound target summaries and validate cumulative receipt serialization
- Preserve final outcomes and legacy unknown evidence for callers
- Verify SDK delegation, additive API, coverage and workspace preflight
Diffstat:
11 files changed, 1282 insertions(+), 3 deletions(-)
diff --git a/contracts/api_baselines/radroots_sync.txt b/contracts/api_baselines/radroots_sync.txt
@@ -0,0 +1,378 @@
+pub mod radroots_sync
+pub mod radroots_sync::ingest
+pub enum radroots_sync::ingest::AdmissionDecision
+pub radroots_sync::ingest::AdmissionDecision::Reject
+pub radroots_sync::ingest::AdmissionDecision::Verified
+pub radroots_sync::ingest::AdmissionDecision::Visible
+pub struct radroots_sync::ingest::IngestBatchReceipt
+impl radroots_sync::ingest::IngestBatchReceipt
+pub fn radroots_sync::ingest::IngestBatchReceipt::accepted(&self) -> usize
+pub fn radroots_sync::ingest::IngestBatchReceipt::outcomes(&self) -> &[core::result::Result<radroots_sync::ingest::IngestReceipt, radroots_sync::policy::Error>]
+pub fn radroots_sync::ingest::IngestBatchReceipt::rejected(&self) -> usize
+pub struct radroots_sync::ingest::IngestReceipt
+impl radroots_sync::ingest::IngestReceipt
+pub const fn radroots_sync::ingest::IngestReceipt::admission(&self) -> &radroots_storage::event::AdmissionReceipt
+pub const fn radroots_sync::ingest::IngestReceipt::commit_disposition(&self) -> radroots_storage::atomic::AtomicCommitDisposition
+pub const fn radroots_sync::ingest::IngestReceipt::committed_at_unix_ms(&self) -> u64
+pub const fn radroots_sync::ingest::IngestReceipt::sync_id(&self) -> radroots_sync::policy::SyncId
+pub struct radroots_sync::ingest::RegistryPolicy
+impl radroots_sync::ingest::RegistryPolicy
+pub const fn radroots_sync::ingest::RegistryPolicy::verified() -> Self
+pub const fn radroots_sync::ingest::RegistryPolicy::visible() -> Self
+impl radroots_sync::ingest::AdmissionPolicy for radroots_sync::ingest::RegistryPolicy
+pub fn radroots_sync::ingest::RegistryPolicy::decide(&self, &radroots_event::verification::ContractValidatedEvent) -> radroots_sync::ingest::AdmissionDecision
+pub fn radroots_sync::ingest::RegistryPolicy::policy_id(&self) -> &'static str
+pub fn radroots_sync::ingest::RegistryPolicy::select_contract(&self, &radroots_event::verification::SignatureVerifiedEvent) -> core::option::Option<&'static str>
+pub trait radroots_sync::ingest::AdmissionPolicy: core::marker::Send + core::marker::Sync
+pub fn radroots_sync::ingest::AdmissionPolicy::decide(&self, &radroots_event::verification::ContractValidatedEvent) -> radroots_sync::ingest::AdmissionDecision
+pub fn radroots_sync::ingest::AdmissionPolicy::policy_id(&self) -> &'static str
+pub fn radroots_sync::ingest::AdmissionPolicy::select_contract(&self, &radroots_event::verification::SignatureVerifiedEvent) -> core::option::Option<&'static str>
+impl radroots_sync::ingest::AdmissionPolicy for radroots_sync::ingest::RegistryPolicy
+pub fn radroots_sync::ingest::RegistryPolicy::decide(&self, &radroots_event::verification::ContractValidatedEvent) -> radroots_sync::ingest::AdmissionDecision
+pub fn radroots_sync::ingest::RegistryPolicy::policy_id(&self) -> &'static str
+pub fn radroots_sync::ingest::RegistryPolicy::select_contract(&self, &radroots_event::verification::SignatureVerifiedEvent) -> core::option::Option<&'static str>
+pub mod radroots_sync::policy
+#[non_exhaustive] pub enum radroots_sync::policy::Error
+pub radroots_sync::policy::Error::AdmissionFailed
+pub radroots_sync::policy::Error::ClockUnavailable
+pub radroots_sync::policy::Error::DeadlineOverflow
+pub radroots_sync::policy::Error::DeliveryDeferred
+pub radroots_sync::policy::Error::InvalidDeadlinePolicy
+pub radroots_sync::policy::Error::InvalidDeliveryRequest
+pub radroots_sync::policy::Error::InvalidIngestReceipt
+pub radroots_sync::policy::Error::InvalidProjectionRequest
+pub radroots_sync::policy::Error::InvalidPullRequest
+pub radroots_sync::policy::Error::InvalidPushRequest
+pub radroots_sync::policy::Error::InvalidReducerOutput
+pub radroots_sync::policy::Error::InvalidSignerOutput
+pub radroots_sync::policy::Error::InvalidSourcePage
+pub radroots_sync::policy::Error::InvalidStatusRequest
+pub radroots_sync::policy::Error::InvalidSyncId
+pub radroots_sync::policy::Error::MissingSigner
+pub radroots_sync::policy::Error::MissingSink
+pub radroots_sync::policy::Error::MissingSource
+pub radroots_sync::policy::Error::MissingTransportCapability
+pub radroots_sync::policy::Error::PolicyRejected
+pub radroots_sync::policy::Error::ReducerFailed
+pub radroots_sync::policy::Error::SignerCapabilityUnavailable
+pub radroots_sync::policy::Error::SignerDeadlineExceeded
+pub radroots_sync::policy::Error::SignerFailed
+pub radroots_sync::policy::Error::SignerWithoutSink
+pub radroots_sync::policy::Error::SigningCancelled
+pub radroots_sync::policy::Error::SigningIndeterminate
+pub radroots_sync::policy::Error::StorageConflict
+pub radroots_sync::policy::Error::StorageFailed
+pub radroots_sync::policy::Error::VerificationFailed
+pub radroots_sync::policy::Error::WorkClaimConflict
+impl core::error::Error for radroots_sync::policy::Error
+impl core::fmt::Display for radroots_sync::policy::Error
+pub fn radroots_sync::policy::Error::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+#[non_exhaustive] pub enum radroots_sync::policy::OperationKind
+pub radroots_sync::policy::OperationKind::Deliver
+pub radroots_sync::policy::OperationKind::Ingest
+pub radroots_sync::policy::OperationKind::Projection
+pub radroots_sync::policy::OperationKind::Pull
+pub radroots_sync::policy::OperationKind::Sign
+pub struct radroots_sync::policy::DeadlinePolicy
+impl radroots_sync::policy::DeadlinePolicy
+pub fn radroots_sync::policy::DeadlinePolicy::deadline_unix_ms(self, radroots_sync::policy::OperationKind, u64) -> core::result::Result<u64, radroots_sync::policy::Error>
+pub const fn radroots_sync::policy::DeadlinePolicy::new(u64, u64, u64) -> core::result::Result<Self, radroots_sync::policy::Error>
+pub const fn radroots_sync::policy::DeadlinePolicy::timeout_ms(self, radroots_sync::policy::OperationKind) -> u64
+impl<'de> serde_core::de::Deserialize<'de> for radroots_sync::policy::DeadlinePolicy
+pub fn radroots_sync::policy::DeadlinePolicy::deserialize<D>(D) -> core::result::Result<Self, <D as serde_core::de::Deserializer>::Error> where D: serde_core::de::Deserializer<'de>
+pub struct radroots_sync::policy::EngineBuilder
+impl radroots_sync::policy::EngineBuilder
+pub fn radroots_sync::policy::EngineBuilder::build(self) -> core::result::Result<radroots_sync::Engine, radroots_sync::policy::Error>
+pub fn radroots_sync::policy::EngineBuilder::signer(self, alloc::sync::Arc<dyn radroots_signing::signer::Signer>) -> Self
+pub fn radroots_sync::policy::EngineBuilder::sink(self, alloc::sync::Arc<dyn radroots_transport::sink::EventSink>) -> Self
+pub fn radroots_sync::policy::EngineBuilder::source(self, alloc::sync::Arc<dyn radroots_transport::source::EventSource>) -> Self
+pub struct radroots_sync::policy::SyncId(_)
+impl radroots_sync::policy::SyncId
+pub const fn radroots_sync::policy::SyncId::as_bytes(&self) -> &[u8; 16]
+pub const fn radroots_sync::policy::SyncId::new([u8; 16]) -> core::result::Result<Self, radroots_sync::policy::Error>
+impl<'de> serde_core::de::Deserialize<'de> for radroots_sync::policy::SyncId
+pub fn radroots_sync::policy::SyncId::deserialize<D>(D) -> core::result::Result<Self, <D as serde_core::de::Deserializer>::Error> where D: serde_core::de::Deserializer<'de>
+pub trait radroots_sync::policy::Clock: core::marker::Send + core::marker::Sync
+pub fn radroots_sync::policy::Clock::now_unix_ms(&self) -> core::result::Result<u64, radroots_sync::policy::Error>
+pub trait radroots_sync::policy::IdSource: core::marker::Send + core::marker::Sync
+pub fn radroots_sync::policy::IdSource::next_id(&self, radroots_sync::policy::OperationKind) -> core::result::Result<radroots_sync::policy::SyncId, radroots_sync::policy::Error>
+pub trait radroots_sync::policy::SyncStorage: radroots_storage::event::EventStore + radroots_storage::journal::Journal + radroots_storage::outbox::Outbox + radroots_storage::projection::ProjectionStore + radroots_storage::atomic::AtomicStorage + radroots_storage::authored_atomic::AuthoredAtomicStorage + radroots_storage::status::StorageStatusProvider
+impl<T> radroots_sync::policy::SyncStorage for T where T: radroots_storage::event::EventStore + radroots_storage::journal::Journal + radroots_storage::outbox::Outbox + radroots_storage::projection::ProjectionStore + radroots_storage::atomic::AtomicStorage + radroots_storage::authored_atomic::AuthoredAtomicStorage + radroots_storage::status::StorageStatusProvider
+pub mod radroots_sync::projection
+pub enum radroots_sync::projection::RefreshKind
+pub radroots_sync::projection::RefreshKind::Incremental
+pub radroots_sync::projection::RefreshKind::Rebuild
+pub enum radroots_sync::projection::RefreshState
+pub radroots_sync::projection::RefreshState::Complete
+pub radroots_sync::projection::RefreshState::Failed
+pub radroots_sync::projection::RefreshState::Partial
+pub struct radroots_sync::projection::ReducerError
+pub struct radroots_sync::projection::RefreshReceipt
+impl radroots_sync::projection::RefreshReceipt
+pub const fn radroots_sync::projection::RefreshReceipt::batches(&self) -> u16
+pub const fn radroots_sync::projection::RefreshReceipt::checkpoint(&self) -> core::option::Option<&radroots_storage::projection::ProjectionCheckpoint>
+pub const fn radroots_sync::projection::RefreshReceipt::events_reduced(&self) -> usize
+pub const fn radroots_sync::projection::RefreshReceipt::kind(&self) -> radroots_sync::projection::RefreshKind
+pub const fn radroots_sync::projection::RefreshReceipt::rebuild_ticket(&self) -> core::option::Option<radroots_storage::projection::RebuildTicketId>
+pub const fn radroots_sync::projection::RefreshReceipt::state(&self) -> radroots_sync::projection::RefreshState
+pub struct radroots_sync::projection::RefreshRequest
+impl radroots_sync::projection::RefreshRequest
+pub const fn radroots_sync::projection::RefreshRequest::batch_limit(&self) -> u16
+pub const fn radroots_sync::projection::RefreshRequest::generation(&self) -> radroots_storage::projection::ProjectionGeneration
+pub const fn radroots_sync::projection::RefreshRequest::max_batches(&self) -> u16
+pub fn radroots_sync::projection::RefreshRequest::new(radroots_storage::projection::ProjectionId, radroots_storage::projection::ProjectionGeneration, u16, u16) -> core::result::Result<Self, radroots_sync::policy::Error>
+pub const fn radroots_sync::projection::RefreshRequest::projection_id(&self) -> &radroots_storage::projection::ProjectionId
+pub const radroots_sync::projection::PROJECTION_RAW_SOURCE_MAX_EVENTS: u64
+pub const radroots_sync::projection::PROJECTION_REFRESH_MAX_BATCHES: u16
+pub trait radroots_sync::projection::Reducer: core::marker::Send + core::marker::Sync
+pub fn radroots_sync::projection::Reducer::abort_rebuild(&self, radroots_storage::projection::RebuildTicketId, radroots_storage::projection::RebuildFailure) -> core::result::Result<(), radroots_sync::projection::ReducerError>
+pub fn radroots_sync::projection::Reducer::begin_rebuild(&self, radroots_storage::projection::RebuildTicketId, radroots_storage::event::SourceGeneration, radroots_storage::projection::RawSourceDigest) -> core::result::Result<(), radroots_sync::projection::ReducerError>
+pub fn radroots_sync::projection::Reducer::generation(&self) -> radroots_storage::projection::ProjectionGeneration
+pub fn radroots_sync::projection::Reducer::projection_id(&self) -> &radroots_storage::projection::ProjectionId
+pub fn radroots_sync::projection::Reducer::reduce(&self, &[radroots_storage::event::StoredVisibleEvent], u64, core::option::Option<radroots_storage::projection::RebuildTicketId>) -> core::result::Result<u64, radroots_sync::projection::ReducerError>
+pub mod radroots_sync::pull
+pub enum radroots_sync::pull::PullTermination
+pub radroots_sync::pull::PullTermination::Cancelled
+pub radroots_sync::pull::PullTermination::Complete
+pub radroots_sync::pull::PullTermination::Deadline
+pub radroots_sync::pull::PullTermination::PageLimit
+pub radroots_sync::pull::PullTermination::SourceFailed
+pub struct radroots_sync::pull::PullReceipt
+impl radroots_sync::pull::PullReceipt
+pub const fn radroots_sync::pull::PullReceipt::deadline_unix_ms(&self) -> u64
+pub const fn radroots_sync::pull::PullReceipt::events_observed(&self) -> usize
+pub fn radroots_sync::pull::PullReceipt::ingest_outcomes(&self) -> &[core::result::Result<radroots_sync::ingest::IngestReceipt, radroots_sync::policy::Error>]
+pub const fn radroots_sync::pull::PullReceipt::pages_fetched(&self) -> u16
+pub const fn radroots_sync::pull::PullReceipt::resume_from(&self) -> core::option::Option<&radroots_transport::source::FetchCursor>
+pub const fn radroots_sync::pull::PullReceipt::sync_id(&self) -> radroots_sync::policy::SyncId
+pub fn radroots_sync::pull::PullReceipt::target_outcomes(&self) -> &[radroots_transport::outcome::FetchTargetOutcome]
+pub fn radroots_sync::pull::PullReceipt::target_summaries(&self) -> core::option::Option<&[radroots_sync::pull::PullTargetSummary]>
+pub const fn radroots_sync::pull::PullReceipt::termination(&self) -> radroots_sync::pull::PullTermination
+impl<'de> serde_core::de::Deserialize<'de> for radroots_sync::pull::PullReceipt
+pub fn radroots_sync::pull::PullReceipt::deserialize<D>(D) -> core::result::Result<Self, <D as serde_core::de::Deserializer>::Error> where D: serde_core::de::Deserializer<'de>
+pub struct radroots_sync::pull::PullRequest
+impl radroots_sync::pull::PullRequest
+pub const fn radroots_sync::pull::PullRequest::cursor(&self) -> core::option::Option<&radroots_transport::source::FetchCursor>
+pub const fn radroots_sync::pull::PullRequest::max_pages(&self) -> u16
+pub fn radroots_sync::pull::PullRequest::new(radroots_transport::target::TargetSet, u16, u16) -> core::result::Result<Self, radroots_sync::policy::Error>
+pub const fn radroots_sync::pull::PullRequest::page_limit(&self) -> u16
+pub const fn radroots_sync::pull::PullRequest::selector(&self) -> &radroots_transport::source::FetchSelector
+pub const fn radroots_sync::pull::PullRequest::targets(&self) -> &radroots_transport::target::TargetSet
+pub fn radroots_sync::pull::PullRequest::with_cursor(self, radroots_transport::source::FetchCursor) -> Self
+pub fn radroots_sync::pull::PullRequest::with_selector(self, radroots_transport::source::FetchSelector) -> Self
+impl<'de> serde_core::de::Deserialize<'de> for radroots_sync::pull::PullRequest
+pub fn radroots_sync::pull::PullRequest::deserialize<D>(D) -> core::result::Result<Self, <D as serde_core::de::Deserializer>::Error> where D: serde_core::de::Deserializer<'de>
+pub struct radroots_sync::pull::PullTargetSummary
+impl radroots_sync::pull::PullTargetSummary
+pub const fn radroots_sync::pull::PullTargetSummary::all_pages_complete(&self) -> bool
+pub const fn radroots_sync::pull::PullTargetSummary::incomplete_pages(&self) -> u16
+pub const fn radroots_sync::pull::PullTargetSummary::last_incomplete(&self) -> core::option::Option<radroots_transport::outcome::FetchTargetState>
+pub const fn radroots_sync::pull::PullTargetSummary::missing_outcome_pages(&self) -> u16
+pub const fn radroots_sync::pull::PullTargetSummary::pages_observed(&self) -> u16
+pub const fn radroots_sync::pull::PullTargetSummary::target(&self) -> &radroots_transport::target::TargetFingerprint
+impl<'de> serde_core::de::Deserialize<'de> for radroots_sync::pull::PullTargetSummary
+pub fn radroots_sync::pull::PullTargetSummary::deserialize<D>(D) -> core::result::Result<Self, <D as serde_core::de::Deserializer>::Error> where D: serde_core::de::Deserializer<'de>
+pub const radroots_sync::pull::PULL_MAX_PAGES: u16
+pub mod radroots_sync::push
+pub struct radroots_sync::push::AdmissionRunReceipt
+impl radroots_sync::push::AdmissionRunReceipt
+pub const fn radroots_sync::push::AdmissionRunReceipt::artifact(&self) -> &radroots_storage::authored::AuthoredArtifact
+pub const fn radroots_sync::push::AdmissionRunReceipt::is_replay(&self) -> bool
+pub struct radroots_sync::push::DeliveryExecutionReceipt
+impl radroots_sync::push::DeliveryExecutionReceipt
+pub const fn radroots_sync::push::DeliveryExecutionReceipt::is_replay(&self) -> bool
+pub const fn radroots_sync::push::DeliveryExecutionReceipt::plan(&self) -> &radroots_storage::authored_delivery::AuthoredDeliveryPlan
+pub struct radroots_sync::push::PushCancellationReceipt
+impl radroots_sync::push::PushCancellationReceipt
+pub const fn radroots_sync::push::PushCancellationReceipt::changed(&self) -> bool
+pub const fn radroots_sync::push::PushCancellationReceipt::status(&self) -> &radroots_sync::push::PushStatus
+pub struct radroots_sync::push::PushPreparation
+impl radroots_sync::push::PushPreparation
+pub const fn radroots_sync::push::PushPreparation::artifact(&self) -> &radroots_storage::authored::AuthoredArtifact
+pub const fn radroots_sync::push::PushPreparation::delivery_plan(&self) -> &radroots_storage::authored_delivery::AuthoredDeliveryPlan
+pub const fn radroots_sync::push::PushPreparation::is_replay(&self) -> bool
+pub const fn radroots_sync::push::PushPreparation::operation(&self) -> &radroots_storage::authored::AuthoredOperation
+pub struct radroots_sync::push::PushRequest
+impl radroots_sync::push::PushRequest
+pub const fn radroots_sync::push::PushRequest::actor(&self) -> &radroots_signing::actor::Actor
+pub const fn radroots_sync::push::PushRequest::cancellation(&self) -> radroots_signing::request::CancellationPolicy
+pub const fn radroots_sync::push::PushRequest::delivery_deadline_unix_ms(&self) -> u64
+pub const fn radroots_sync::push::PushRequest::idempotency_key(&self) -> &radroots_storage::journal::IdempotencyKey
+pub fn radroots_sync::push::PushRequest::new(radroots_sync::policy::SyncId, radroots_storage::journal::IdempotencyKey, radroots_signing::actor::Actor, radroots_event_codec::authoring::AuthoredEventPlan, radroots_transport::target::TargetSet, radroots_transport::policy::SatisfactionPolicy, u64, radroots_signing::request::CancellationPolicy) -> core::result::Result<Self, radroots_sync::policy::Error>
+pub const fn radroots_sync::push::PushRequest::operation_id(&self) -> radroots_sync::policy::SyncId
+pub const fn radroots_sync::push::PushRequest::plan(&self) -> &radroots_event_codec::authoring::AuthoredEventPlan
+pub const fn radroots_sync::push::PushRequest::satisfaction(&self) -> &radroots_transport::policy::SatisfactionPolicy
+pub const fn radroots_sync::push::PushRequest::targets(&self) -> &radroots_transport::target::TargetSet
+impl core::fmt::Debug for radroots_sync::push::PushRequest
+pub fn radroots_sync::push::PushRequest::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct radroots_sync::push::PushStatus
+impl radroots_sync::push::PushStatus
+pub const fn radroots_sync::push::PushStatus::artifact(&self) -> &radroots_storage::authored::AuthoredArtifact
+pub const fn radroots_sync::push::PushStatus::delivery_plan(&self) -> &radroots_storage::authored_delivery::AuthoredDeliveryPlan
+pub const fn radroots_sync::push::PushStatus::operation(&self) -> &radroots_storage::authored::AuthoredOperation
+pub const fn radroots_sync::push::PushStatus::settlement(&self) -> radroots_storage::authored::OperationSettlement
+pub struct radroots_sync::push::SigningRunReceipt
+impl radroots_sync::push::SigningRunReceipt
+pub const fn radroots_sync::push::SigningRunReceipt::artifact(&self) -> &radroots_storage::authored::AuthoredArtifact
+pub const fn radroots_sync::push::SigningRunReceipt::is_replay(&self) -> bool
+pub mod radroots_sync::status
+pub struct radroots_sync::status::CapabilityReport<T>
+impl<T> radroots_sync::status::CapabilityReport<T>
+pub const fn radroots_sync::status::CapabilityReport<T>::state(&self) -> radroots_protocol::runtime::v1::SyncCapabilityState
+pub const fn radroots_sync::status::CapabilityReport<T>::status(&self) -> core::option::Option<&T>
+pub struct radroots_sync::status::ProjectionReport
+impl radroots_sync::status::ProjectionReport
+pub const fn radroots_sync::status::ProjectionReport::projection_id(&self) -> &radroots_storage::projection::ProjectionId
+pub const fn radroots_sync::status::ProjectionReport::status(&self) -> core::option::Option<&radroots_storage::projection::ProjectionStatus>
+pub struct radroots_sync::status::SyncStatus
+impl radroots_sync::status::SyncStatus
+pub const fn radroots_sync::status::SyncStatus::events(&self) -> &radroots_storage::status::EventStoreStatus
+pub const fn radroots_sync::status::SyncStatus::health(&self) -> radroots_protocol::runtime::v1::SyncHealth
+pub const fn radroots_sync::status::SyncStatus::outbox(&self) -> radroots_storage::outbox::OutboxStatus
+pub fn radroots_sync::status::SyncStatus::projections(&self) -> &[radroots_sync::status::ProjectionReport]
+pub const fn radroots_sync::status::SyncStatus::signer(&self) -> &radroots_sync::status::CapabilityReport<radroots_signing::status::SignerStatus>
+pub const fn radroots_sync::status::SyncStatus::sink(&self) -> &radroots_sync::status::CapabilityReport<radroots_transport::status::SinkStatus>
+pub const fn radroots_sync::status::SyncStatus::source(&self) -> &radroots_sync::status::CapabilityReport<radroots_transport::status::SourceStatus>
+pub const fn radroots_sync::status::SyncStatus::storage(&self) -> radroots_storage::status::StorageStatus
+pub fn radroots_sync::status::SyncStatus::to_protocol(&self) -> radroots_protocol::runtime::v1::SyncStatusReceipt
+#[non_exhaustive] pub enum radroots_sync::Error
+pub radroots_sync::Error::AdmissionFailed
+pub radroots_sync::Error::ClockUnavailable
+pub radroots_sync::Error::DeadlineOverflow
+pub radroots_sync::Error::DeliveryDeferred
+pub radroots_sync::Error::InvalidDeadlinePolicy
+pub radroots_sync::Error::InvalidDeliveryRequest
+pub radroots_sync::Error::InvalidIngestReceipt
+pub radroots_sync::Error::InvalidProjectionRequest
+pub radroots_sync::Error::InvalidPullRequest
+pub radroots_sync::Error::InvalidPushRequest
+pub radroots_sync::Error::InvalidReducerOutput
+pub radroots_sync::Error::InvalidSignerOutput
+pub radroots_sync::Error::InvalidSourcePage
+pub radroots_sync::Error::InvalidStatusRequest
+pub radroots_sync::Error::InvalidSyncId
+pub radroots_sync::Error::MissingSigner
+pub radroots_sync::Error::MissingSink
+pub radroots_sync::Error::MissingSource
+pub radroots_sync::Error::MissingTransportCapability
+pub radroots_sync::Error::PolicyRejected
+pub radroots_sync::Error::ReducerFailed
+pub radroots_sync::Error::SignerCapabilityUnavailable
+pub radroots_sync::Error::SignerDeadlineExceeded
+pub radroots_sync::Error::SignerFailed
+pub radroots_sync::Error::SignerWithoutSink
+pub radroots_sync::Error::SigningCancelled
+pub radroots_sync::Error::SigningIndeterminate
+pub radroots_sync::Error::StorageConflict
+pub radroots_sync::Error::StorageFailed
+pub radroots_sync::Error::VerificationFailed
+pub radroots_sync::Error::WorkClaimConflict
+impl core::error::Error for radroots_sync::policy::Error
+impl core::fmt::Display for radroots_sync::policy::Error
+pub fn radroots_sync::policy::Error::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct radroots_sync::AdmissionRunReceipt
+impl radroots_sync::push::AdmissionRunReceipt
+pub const fn radroots_sync::push::AdmissionRunReceipt::artifact(&self) -> &radroots_storage::authored::AuthoredArtifact
+pub const fn radroots_sync::push::AdmissionRunReceipt::is_replay(&self) -> bool
+pub struct radroots_sync::DeliveryExecutionReceipt
+impl radroots_sync::push::DeliveryExecutionReceipt
+pub const fn radroots_sync::push::DeliveryExecutionReceipt::is_replay(&self) -> bool
+pub const fn radroots_sync::push::DeliveryExecutionReceipt::plan(&self) -> &radroots_storage::authored_delivery::AuthoredDeliveryPlan
+pub struct radroots_sync::Engine
+impl radroots_sync::Engine
+pub async fn radroots_sync::Engine::admit_signed(&self, radroots_sync::policy::SyncId) -> core::result::Result<radroots_sync::push::AdmissionRunReceipt, radroots_sync::policy::Error>
+pub async fn radroots_sync::Engine::cancel_push(&self, radroots_sync::policy::SyncId) -> core::result::Result<radroots_sync::push::PushCancellationReceipt, radroots_sync::policy::Error>
+pub async fn radroots_sync::Engine::deliver_push(&self, radroots_sync::policy::SyncId) -> core::result::Result<radroots_sync::push::DeliveryExecutionReceipt, radroots_sync::policy::Error>
+pub async fn radroots_sync::Engine::prepare_push(&self, radroots_sync::push::PushRequest) -> core::result::Result<radroots_sync::push::PushPreparation, radroots_sync::policy::Error>
+pub async fn radroots_sync::Engine::push_status(&self, radroots_sync::policy::SyncId) -> core::result::Result<core::option::Option<radroots_sync::push::PushStatus>, radroots_sync::policy::Error>
+pub async fn radroots_sync::Engine::sign_prepared(&self, radroots_sync::push::PushRequest) -> core::result::Result<radroots_sync::push::SigningRunReceipt, radroots_sync::policy::Error>
+impl radroots_sync::Engine
+pub fn radroots_sync::Engine::builder(alloc::sync::Arc<dyn radroots_sync::policy::SyncStorage>, alloc::sync::Arc<dyn radroots_sync::policy::Clock>, alloc::sync::Arc<dyn radroots_sync::policy::IdSource>, radroots_sync::policy::DeadlinePolicy) -> radroots_sync::policy::EngineBuilder
+pub fn radroots_sync::Engine::clock(&self) -> &dyn radroots_sync::policy::Clock
+pub const fn radroots_sync::Engine::deadlines(&self) -> radroots_sync::policy::DeadlinePolicy
+pub fn radroots_sync::Engine::ids(&self) -> &dyn radroots_sync::policy::IdSource
+pub fn radroots_sync::Engine::signer(&self) -> core::option::Option<&dyn radroots_signing::signer::Signer>
+pub fn radroots_sync::Engine::sink(&self) -> core::option::Option<&dyn radroots_transport::sink::EventSink>
+pub fn radroots_sync::Engine::source(&self) -> core::option::Option<&dyn radroots_transport::source::EventSource>
+pub fn radroots_sync::Engine::storage(&self) -> &dyn radroots_sync::policy::SyncStorage
+impl radroots_sync::Engine
+pub async fn radroots_sync::Engine::ingest(&self, radroots_transport::source::ObservedEvent, &dyn radroots_sync::ingest::AdmissionPolicy) -> core::result::Result<radroots_sync::ingest::IngestReceipt, radroots_sync::policy::Error>
+pub async fn radroots_sync::Engine::ingest_batch(&self, alloc::vec::Vec<radroots_transport::source::ObservedEvent>, &dyn radroots_sync::ingest::AdmissionPolicy) -> radroots_sync::ingest::IngestBatchReceipt
+impl radroots_sync::Engine
+pub async fn radroots_sync::Engine::pull(&self, radroots_sync::pull::PullRequest, &dyn radroots_sync::ingest::AdmissionPolicy) -> core::result::Result<radroots_sync::pull::PullReceipt, radroots_sync::policy::Error>
+impl radroots_sync::Engine
+pub async fn radroots_sync::Engine::refresh_projection(&self, radroots_sync::projection::RefreshRequest, &dyn radroots_sync::projection::Reducer) -> core::result::Result<radroots_sync::projection::RefreshReceipt, radroots_sync::policy::Error>
+impl radroots_sync::Engine
+pub fn radroots_sync::Engine::retry_decision(&self, &radroots_storage::authored_delivery::AuthoredDeliveryPlan, u64) -> core::result::Result<radroots_protocol::runtime::v1::SyncRetryDecision, radroots_sync::policy::Error>
+pub async fn radroots_sync::Engine::status(&self, &[radroots_storage::projection::ProjectionId]) -> core::result::Result<radroots_sync::status::SyncStatus, radroots_sync::policy::Error>
+impl core::fmt::Debug for radroots_sync::Engine
+pub fn radroots_sync::Engine::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct radroots_sync::PullReceipt
+impl radroots_sync::pull::PullReceipt
+pub const fn radroots_sync::pull::PullReceipt::deadline_unix_ms(&self) -> u64
+pub const fn radroots_sync::pull::PullReceipt::events_observed(&self) -> usize
+pub fn radroots_sync::pull::PullReceipt::ingest_outcomes(&self) -> &[core::result::Result<radroots_sync::ingest::IngestReceipt, radroots_sync::policy::Error>]
+pub const fn radroots_sync::pull::PullReceipt::pages_fetched(&self) -> u16
+pub const fn radroots_sync::pull::PullReceipt::resume_from(&self) -> core::option::Option<&radroots_transport::source::FetchCursor>
+pub const fn radroots_sync::pull::PullReceipt::sync_id(&self) -> radroots_sync::policy::SyncId
+pub fn radroots_sync::pull::PullReceipt::target_outcomes(&self) -> &[radroots_transport::outcome::FetchTargetOutcome]
+pub fn radroots_sync::pull::PullReceipt::target_summaries(&self) -> core::option::Option<&[radroots_sync::pull::PullTargetSummary]>
+pub const fn radroots_sync::pull::PullReceipt::termination(&self) -> radroots_sync::pull::PullTermination
+impl<'de> serde_core::de::Deserialize<'de> for radroots_sync::pull::PullReceipt
+pub fn radroots_sync::pull::PullReceipt::deserialize<D>(D) -> core::result::Result<Self, <D as serde_core::de::Deserializer>::Error> where D: serde_core::de::Deserializer<'de>
+pub struct radroots_sync::PullRequest
+impl radroots_sync::pull::PullRequest
+pub const fn radroots_sync::pull::PullRequest::cursor(&self) -> core::option::Option<&radroots_transport::source::FetchCursor>
+pub const fn radroots_sync::pull::PullRequest::max_pages(&self) -> u16
+pub fn radroots_sync::pull::PullRequest::new(radroots_transport::target::TargetSet, u16, u16) -> core::result::Result<Self, radroots_sync::policy::Error>
+pub const fn radroots_sync::pull::PullRequest::page_limit(&self) -> u16
+pub const fn radroots_sync::pull::PullRequest::selector(&self) -> &radroots_transport::source::FetchSelector
+pub const fn radroots_sync::pull::PullRequest::targets(&self) -> &radroots_transport::target::TargetSet
+pub fn radroots_sync::pull::PullRequest::with_cursor(self, radroots_transport::source::FetchCursor) -> Self
+pub fn radroots_sync::pull::PullRequest::with_selector(self, radroots_transport::source::FetchSelector) -> Self
+impl<'de> serde_core::de::Deserialize<'de> for radroots_sync::pull::PullRequest
+pub fn radroots_sync::pull::PullRequest::deserialize<D>(D) -> core::result::Result<Self, <D as serde_core::de::Deserializer>::Error> where D: serde_core::de::Deserializer<'de>
+pub struct radroots_sync::PushCancellationReceipt
+impl radroots_sync::push::PushCancellationReceipt
+pub const fn radroots_sync::push::PushCancellationReceipt::changed(&self) -> bool
+pub const fn radroots_sync::push::PushCancellationReceipt::status(&self) -> &radroots_sync::push::PushStatus
+pub struct radroots_sync::PushPreparation
+impl radroots_sync::push::PushPreparation
+pub const fn radroots_sync::push::PushPreparation::artifact(&self) -> &radroots_storage::authored::AuthoredArtifact
+pub const fn radroots_sync::push::PushPreparation::delivery_plan(&self) -> &radroots_storage::authored_delivery::AuthoredDeliveryPlan
+pub const fn radroots_sync::push::PushPreparation::is_replay(&self) -> bool
+pub const fn radroots_sync::push::PushPreparation::operation(&self) -> &radroots_storage::authored::AuthoredOperation
+pub struct radroots_sync::PushRequest
+impl radroots_sync::push::PushRequest
+pub const fn radroots_sync::push::PushRequest::actor(&self) -> &radroots_signing::actor::Actor
+pub const fn radroots_sync::push::PushRequest::cancellation(&self) -> radroots_signing::request::CancellationPolicy
+pub const fn radroots_sync::push::PushRequest::delivery_deadline_unix_ms(&self) -> u64
+pub const fn radroots_sync::push::PushRequest::idempotency_key(&self) -> &radroots_storage::journal::IdempotencyKey
+pub fn radroots_sync::push::PushRequest::new(radroots_sync::policy::SyncId, radroots_storage::journal::IdempotencyKey, radroots_signing::actor::Actor, radroots_event_codec::authoring::AuthoredEventPlan, radroots_transport::target::TargetSet, radroots_transport::policy::SatisfactionPolicy, u64, radroots_signing::request::CancellationPolicy) -> core::result::Result<Self, radroots_sync::policy::Error>
+pub const fn radroots_sync::push::PushRequest::operation_id(&self) -> radroots_sync::policy::SyncId
+pub const fn radroots_sync::push::PushRequest::plan(&self) -> &radroots_event_codec::authoring::AuthoredEventPlan
+pub const fn radroots_sync::push::PushRequest::satisfaction(&self) -> &radroots_transport::policy::SatisfactionPolicy
+pub const fn radroots_sync::push::PushRequest::targets(&self) -> &radroots_transport::target::TargetSet
+impl core::fmt::Debug for radroots_sync::push::PushRequest
+pub fn radroots_sync::push::PushRequest::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct radroots_sync::PushStatus
+impl radroots_sync::push::PushStatus
+pub const fn radroots_sync::push::PushStatus::artifact(&self) -> &radroots_storage::authored::AuthoredArtifact
+pub const fn radroots_sync::push::PushStatus::delivery_plan(&self) -> &radroots_storage::authored_delivery::AuthoredDeliveryPlan
+pub const fn radroots_sync::push::PushStatus::operation(&self) -> &radroots_storage::authored::AuthoredOperation
+pub const fn radroots_sync::push::PushStatus::settlement(&self) -> radroots_storage::authored::OperationSettlement
+pub struct radroots_sync::SigningRunReceipt
+impl radroots_sync::push::SigningRunReceipt
+pub const fn radroots_sync::push::SigningRunReceipt::artifact(&self) -> &radroots_storage::authored::AuthoredArtifact
+pub const fn radroots_sync::push::SigningRunReceipt::is_replay(&self) -> bool
+pub struct radroots_sync::SyncStatus
+impl radroots_sync::status::SyncStatus
+pub const fn radroots_sync::status::SyncStatus::events(&self) -> &radroots_storage::status::EventStoreStatus
+pub const fn radroots_sync::status::SyncStatus::health(&self) -> radroots_protocol::runtime::v1::SyncHealth
+pub const fn radroots_sync::status::SyncStatus::outbox(&self) -> radroots_storage::outbox::OutboxStatus
+pub fn radroots_sync::status::SyncStatus::projections(&self) -> &[radroots_sync::status::ProjectionReport]
+pub const fn radroots_sync::status::SyncStatus::signer(&self) -> &radroots_sync::status::CapabilityReport<radroots_signing::status::SignerStatus>
+pub const fn radroots_sync::status::SyncStatus::sink(&self) -> &radroots_sync::status::CapabilityReport<radroots_transport::status::SinkStatus>
+pub const fn radroots_sync::status::SyncStatus::source(&self) -> &radroots_sync::status::CapabilityReport<radroots_transport::status::SourceStatus>
+pub const fn radroots_sync::status::SyncStatus::storage(&self) -> radroots_storage::status::StorageStatus
+pub fn radroots_sync::status::SyncStatus::to_protocol(&self) -> radroots_protocol::runtime::v1::SyncStatusReceipt
diff --git a/contracts/architecture/decisions/pull_target_evidence.v1.json b/contracts/architecture/decisions/pull_target_evidence.v1.json
@@ -0,0 +1,27 @@
+{
+ "schema": "radroots.pull-target-evidence.v1",
+ "owner": "radroots_sync",
+ "status": "implemented",
+ "scope": "One explicit bounded Engine::pull call with its exact target set, selector, initial cursor and deadline.",
+ "max_targets": 64,
+ "max_pages": 1000,
+ "target_order": "Stable validated request order; every requested target is represented in a newly measured receipt, including when no page returns.",
+ "summary_fields": [
+ "target",
+ "pages_observed",
+ "incomplete_pages",
+ "missing_outcome_pages",
+ "last_incomplete"
+ ],
+ "counting": "Every returned validated page increments pages_observed for every requested target. A present non-complete typed outcome increments incomplete_pages and replaces last_incomplete. An absent outcome increments missing_outcome_pages without fabricating a transport state.",
+ "bounds": "Use the existing target-set and pull-page constants. Retain no additional messages, event bodies or per-page list. Reject an oversized serialized summary sequence before retaining its extra item; reject duplicate targets and invalid or inconsistent summary counters.",
+ "complete_evidence": "A summary reports all_pages_complete only for positive pages_observed with zero incomplete_pages and zero missing_outcome_pages. Callers must independently require appropriate pull termination and inspect the request scope; no global history is proved.",
+ "compatibility": "Keep target_outcomes as the last available outcome for each target. Add optional target_summaries. Missing or null summaries in legacy deserialization remain unknown; newly measured receipts always contain summaries. Existing source-page and transport behavior is unchanged.",
+ "source_failure": "No returned page means no new target observation. Preserve earlier summaries and the existing SourceFailed termination; never infer a target-specific failure or successful completion from a missing source response.",
+ "non_goals": [
+ "new scheduler, clock or transport policy",
+ "global-history completeness",
+ "per-page history or unbounded diagnostics",
+ "release, Nix, device or deployment qualification"
+ ]
+}
diff --git a/contracts/architecture/deviations.toml b/contracts/architecture/deviations.toml
@@ -2,6 +2,33 @@ schema_version = 1
architecture_id = "radroots.crates.release.v1"
[[deviation]]
+id = "RCRV1-DEV-015"
+date = "2026-09-09"
+status = "closed"
+approval = "Standing user approval of the reviewed refactor and necessary shared-owner prerequisites, including verified producer commits and non-force publication."
+affected_steps = ["201"]
+spec_anchors = [
+ "contracts/crates/release_v1/radroots_crates_release_v1.toml#package.radroots_sync",
+]
+source_evidence = [
+ "PullReceipt intentionally retains the final per-target outcome, so a later complete page replaces an earlier partial or failed page.",
+ "FetchPage permits missing target outcomes; a complete page marker alone cannot prove every requested target supplied complete evidence.",
+]
+replacement_action = "Add bounded cumulative per-target summaries under contracts/architecture/decisions/pull_target_evidence.v1.json, preserving final-page API semantics and legacy receipts as unknown evidence."
+verification = [
+ "Exercise incomplete/complete ordering, omitted outcomes, multiple targets, exact limits, termination and compatible serialization through the real shared engine and SDK delegation.",
+ "Qualify current package/workspace, API, feature, coverage and release-preflight surfaces before publication and consumer adoption.",
+]
+unresolved_risk = "Current shared sync, SDK, workspace, API and unchanged coverage gates pass. Summary completeness remains scoped to returned pages and must be combined with pull termination; no global-history or new release authority is claimed."
+normative_architecture_change = false
+adr_required = false
+closure_evidence = [
+ "The reproduced partial-then-complete regression passes while final outcomes remain unchanged; all actual incomplete states, missing outcomes, multiple targets, 64 targets and 1,000 pages are exercised.",
+ "Legacy and malformed serialization, source/page/deadline/cancellation cases, and actual SDK cumulative-evidence delegation pass with the supported feature profiles.",
+ "Generated shared API retains every previous declaration, SDK API is byte-identical, fresh sync/SDK coverage and all 45 reports/aggregate pass, and the full workspace and release-preflight lanes pass.",
+]
+
+[[deviation]]
id = "RCRV1-DEV-013"
date = "2026-09-09"
status = "closed"
diff --git a/crates/sdk/src/sync.rs b/crates/sdk/src/sync.rs
@@ -320,6 +320,92 @@ mod tests {
}
}
+ struct PartialThenCompleteSource(AtomicU8);
+
+ impl EventSource for PartialThenCompleteSource {
+ fn status(
+ &self,
+ ) -> radroots_transport::BoxFuture<'_, Result<SourceStatus, TransportError>> {
+ Box::pin(async { unreachable!("explicit pull only") })
+ }
+
+ fn fetch(
+ &self,
+ request: FetchRequest,
+ ) -> radroots_transport::BoxFuture<'_, Result<FetchPage, TransportError>> {
+ Box::pin(async move {
+ use radroots_transport::{
+ outcome::{FetchTargetOutcome, FetchTargetState},
+ source::FetchCursor,
+ };
+ let page = self.0.fetch_add(1, Ordering::Relaxed);
+ assert!(page < 2);
+ let (state, next) = if page == 0 {
+ assert!(request.cursor().is_none());
+ (
+ FetchTargetState::Partial,
+ NextPage::Cursor(FetchCursor::parse("second").expect("cursor")),
+ )
+ } else {
+ assert_eq!(request.cursor().map(FetchCursor::as_str), Some("second"));
+ (FetchTargetState::Complete, NextPage::Complete)
+ };
+ let outcome = FetchTargetOutcome::new(
+ request.target_set().targets()[0].fingerprint().clone(),
+ state,
+ );
+ FetchPage::for_request(&request, vec![], vec![outcome], next)
+ })
+ }
+ }
+
+ #[tokio::test]
+ async fn operations_preserve_cumulative_pull_evidence_after_a_later_complete_page() {
+ let storage = Arc::new(MemoryStorage::new(
+ SourceGeneration::new([3; 32]).expect("generation"),
+ ));
+ let source = Arc::new(PartialThenCompleteSource(AtomicU8::new(0)));
+ let engine = Engine::builder(
+ storage.clone(),
+ Arc::new(FixedClock),
+ Arc::new(SequenceIds(AtomicU8::new(1))),
+ DeadlinePolicy::new(1_000, 1_000, 1_000).expect("deadlines"),
+ )
+ .source(source.clone())
+ .build()
+ .expect("engine");
+ let client = ClientBuilder::new()
+ .storage(storage)
+ .sync_engine(engine)
+ .build()
+ .expect("client");
+ let receipt = client
+ .sync()
+ .expect("open client")
+ .expect("sync")
+ .pull(
+ PullRequest::new(TargetSet::new(vec![target()]).expect("targets"), 1, 2)
+ .expect("request"),
+ &RegistryPolicy::visible(),
+ )
+ .await
+ .expect("receipt");
+ assert_eq!(source.0.load(Ordering::Relaxed), 2);
+ assert_eq!(receipt.termination(), PullTermination::Complete);
+ assert_eq!(
+ receipt.target_outcomes()[0].state(),
+ radroots_transport::outcome::FetchTargetState::Complete
+ );
+ let summary = &receipt.target_summaries().expect("measured")[0];
+ assert_eq!(summary.pages_observed(), 2);
+ assert_eq!(summary.incomplete_pages(), 1);
+ assert_eq!(
+ summary.last_incomplete(),
+ Some(radroots_transport::outcome::FetchTargetState::Partial)
+ );
+ assert!(!summary.all_pages_complete());
+ }
+
struct EmptyReducer {
id: ProjectionId,
generation: ProjectionGeneration,
@@ -460,6 +546,12 @@ mod tests {
.await
.expect("pull");
assert_eq!(pull.termination(), PullTermination::Cancelled);
+ let summaries = pull.target_summaries().expect("measured target evidence");
+ assert_eq!(summaries.len(), 1);
+ assert_eq!(summaries[0].target(), target().fingerprint());
+ assert_eq!(summaries[0].pages_observed(), 1);
+ assert_eq!(summaries[0].missing_outcome_pages(), 1);
+ assert!(!summaries[0].all_pages_complete());
let ingest = operations
.ingest_batch(vec![invalid_observation(), invalid_observation()], &policy)
diff --git a/crates/sync/README.md b/crates/sync/README.md
@@ -6,5 +6,13 @@ The package owns the shared ingest, pull, projection, push, policy, and status
boundaries. It does not create an executor, spawn workers, install timers, own
process lifecycle, store UI state, or branch on concrete transport adapters.
+Pull receipts retain the last available outcome for each target and bounded
+cumulative target summaries across returned pages. A later complete outcome
+does not erase an earlier incomplete or missing outcome. Summaries preserve
+request order and use the existing 64-target and 1,000-page bounds. A legacy
+receipt without summaries has unknown cumulative evidence. Even when every
+page reports completion, callers must inspect pull termination and request
+scope; a receipt never proves complete global history.
+
Publication remains disabled while behavior is implemented and qualified in
the subsequent Release V1 sync checkpoints.
diff --git a/crates/sync/src/pull.rs b/crates/sync/src/pull.rs
@@ -1,5 +1,11 @@
//! Bounded source pagination and ingestion.
+mod summary;
+#[cfg(feature = "serde")]
+mod wire;
+
+pub use summary::PullTargetSummary;
+
use radroots_transport::{
FetchRequest,
outcome::FetchTargetOutcome,
@@ -92,8 +98,8 @@ pub enum PullTermination {
SourceFailed,
}
-/// Normalized receipt retaining every ingest outcome and final target state.
-#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+/// Normalized receipt retaining ingest outcomes, final states and bounded page evidence.
+#[cfg_attr(feature = "serde", derive(serde::Serialize))]
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct PullReceipt {
sync_id: SyncId,
@@ -102,6 +108,7 @@ pub struct PullReceipt {
events_observed: usize,
ingest_outcomes: Vec<Result<IngestReceipt, Error>>,
target_outcomes: Vec<FetchTargetOutcome>,
+ target_summaries: Option<summary::PullTargetSummaries>,
termination: PullTermination,
resume_from: Option<FetchCursor>,
}
@@ -131,6 +138,15 @@ impl PullReceipt {
self.target_outcomes.as_slice()
}
+ /// Cumulative evidence in request-target order, or `None` for a legacy receipt.
+ ///
+ /// A missing summary inventory is unknown evidence. Even complete summaries
+ /// must be interpreted with the termination and exact request scope; they
+ /// do not prove complete global history.
+ pub fn target_summaries(&self) -> Option<&[PullTargetSummary]> {
+ self.target_summaries.as_ref().map(|value| value.as_slice())
+ }
+
pub const fn termination(&self) -> PullTermination {
self.termination
}
@@ -162,6 +178,7 @@ impl Engine {
events_observed: 0,
ingest_outcomes: Vec::new(),
target_outcomes: Vec::new(),
+ target_summaries: Some(summary::PullTargetSummaries::new(&request.targets)),
termination: PullTermination::Complete,
resume_from: request.cursor.clone(),
};
@@ -196,6 +213,9 @@ impl Engine {
receipt.pages_fetched += 1;
receipt.events_observed += page.events().len();
merge_target_outcomes(&mut receipt.target_outcomes, page.target_outcomes());
+ if let Some(summaries) = &mut receipt.target_summaries {
+ summaries.observe(page.target_outcomes());
+ }
let outcomes = self.ingest_batch(page.events().to_vec(), admission).await;
receipt
.ingest_outcomes
diff --git a/crates/sync/src/pull/summary.rs b/crates/sync/src/pull/summary.rs
@@ -0,0 +1,180 @@
+use radroots_transport::{
+ outcome::{FetchTargetOutcome, FetchTargetState},
+ target::{TargetFingerprint, TargetSet},
+};
+
+/// Bounded evidence for one requested target across all returned pull pages.
+///
+/// Counts describe returned, validated pages only. A source failure returning no
+/// page adds no target observation; the pull termination retains that failure.
+/// No diagnostic messages, event bodies or per-page history are retained here.
+#[cfg_attr(feature = "serde", derive(serde::Serialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct PullTargetSummary {
+ target: TargetFingerprint,
+ pages_observed: u16,
+ incomplete_pages: u16,
+ missing_outcome_pages: u16,
+ last_incomplete: Option<FetchTargetState>,
+}
+
+impl PullTargetSummary {
+ /// Exact requested target identity.
+ pub const fn target(&self) -> &TargetFingerprint {
+ &self.target
+ }
+
+ /// Number of validated pages returned during this pull.
+ pub const fn pages_observed(&self) -> u16 {
+ self.pages_observed
+ }
+
+ /// Pages that supplied an explicit non-complete outcome for this target.
+ pub const fn incomplete_pages(&self) -> u16 {
+ self.incomplete_pages
+ }
+
+ /// Pages that supplied no outcome for this requested target.
+ pub const fn missing_outcome_pages(&self) -> u16 {
+ self.missing_outcome_pages
+ }
+
+ /// Last actual non-complete state, preserved across later success or omission.
+ pub const fn last_incomplete(&self) -> Option<FetchTargetState> {
+ self.last_incomplete
+ }
+
+ /// Whether every returned page positively reported this target complete.
+ ///
+ /// False when no page returned. The caller must also inspect pull termination
+ /// and request bounds; this is not a claim about complete global history.
+ pub const fn all_pages_complete(&self) -> bool {
+ self.pages_observed > 0 && self.incomplete_pages == 0 && self.missing_outcome_pages == 0
+ }
+}
+
+#[cfg_attr(feature = "serde", derive(serde::Serialize))]
+#[cfg_attr(feature = "serde", serde(transparent))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub(super) struct PullTargetSummaries(Vec<PullTargetSummary>);
+
+impl PullTargetSummaries {
+ pub(super) fn new(targets: &TargetSet) -> Self {
+ Self(
+ targets
+ .targets()
+ .iter()
+ .map(|target| PullTargetSummary {
+ target: target.fingerprint().clone(),
+ pages_observed: 0,
+ incomplete_pages: 0,
+ missing_outcome_pages: 0,
+ last_incomplete: None,
+ })
+ .collect(),
+ )
+ }
+
+ pub(super) fn as_slice(&self) -> &[PullTargetSummary] {
+ &self.0
+ }
+
+ pub(super) fn observe(&mut self, outcomes: &[FetchTargetOutcome]) {
+ // Only the validated pull loop calls this, at most PULL_MAX_PAGES times.
+ for summary in &mut self.0 {
+ summary.pages_observed += 1;
+ match outcomes
+ .iter()
+ .find(|outcome| outcome.target() == summary.target())
+ {
+ Some(outcome) if outcome.state() != FetchTargetState::Complete => {
+ summary.incomplete_pages += 1;
+ summary.last_incomplete = Some(outcome.state());
+ }
+ Some(_) => {}
+ None => summary.missing_outcome_pages += 1,
+ }
+ }
+ }
+}
+
+#[cfg(feature = "serde")]
+impl<'de> serde::Deserialize<'de> for PullTargetSummary {
+ fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
+ where
+ D: serde::Deserializer<'de>,
+ {
+ #[derive(serde::Deserialize)]
+ #[serde(deny_unknown_fields)]
+ struct Wire {
+ target: TargetFingerprint,
+ pages_observed: u16,
+ incomplete_pages: u16,
+ missing_outcome_pages: u16,
+ last_incomplete: Option<FetchTargetState>,
+ }
+ let wire = Wire::deserialize(deserializer)?;
+ if wire.pages_observed > super::PULL_MAX_PAGES
+ || u32::from(wire.incomplete_pages) + u32::from(wire.missing_outcome_pages)
+ > u32::from(wire.pages_observed)
+ || (wire.incomplete_pages > 0) != wire.last_incomplete.is_some()
+ || wire.last_incomplete == Some(FetchTargetState::Complete)
+ {
+ return Err(serde::de::Error::custom(
+ "invalid pull target summary counts or state",
+ ));
+ }
+ Ok(Self {
+ target: wire.target,
+ pages_observed: wire.pages_observed,
+ incomplete_pages: wire.incomplete_pages,
+ missing_outcome_pages: wire.missing_outcome_pages,
+ last_incomplete: wire.last_incomplete,
+ })
+ }
+}
+
+#[cfg(feature = "serde")]
+impl<'de> serde::Deserialize<'de> for PullTargetSummaries {
+ fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
+ where
+ D: serde::Deserializer<'de>,
+ {
+ struct BoundedSummaries;
+ impl<'de> serde::de::Visitor<'de> for BoundedSummaries {
+ type Value = PullTargetSummaries;
+
+ fn expecting(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
+ formatter.write_str("a nonempty bounded list of unique pull target summaries")
+ }
+
+ fn visit_seq<A>(self, mut sequence: A) -> Result<Self::Value, A::Error>
+ where
+ A: serde::de::SeqAccess<'de>,
+ {
+ let mut summaries: Vec<PullTargetSummary> = Vec::new();
+ while summaries.len() < radroots_transport::target::TARGET_SET_MAX_ITEMS {
+ let Some(summary) = sequence.next_element::<PullTargetSummary>()? else {
+ if summaries.is_empty() {
+ return Err(serde::de::Error::custom("empty pull target summaries"));
+ }
+ return Ok(PullTargetSummaries(summaries));
+ };
+ if summaries
+ .iter()
+ .any(|prior| prior.target() == summary.target())
+ {
+ return Err(serde::de::Error::custom("duplicate pull summary target"));
+ }
+ summaries.push(summary);
+ }
+ // Skip an extra item without deserializing/retaining its fields.
+ if sequence.next_element::<serde::de::IgnoredAny>()?.is_some() {
+ return Err(serde::de::Error::custom("too many pull target summaries"));
+ }
+ Ok(PullTargetSummaries(summaries))
+ }
+ }
+ deserializer.deserialize_seq(BoundedSummaries)
+ }
+}
diff --git a/crates/sync/src/pull/wire.rs b/crates/sync/src/pull/wire.rs
@@ -0,0 +1,84 @@
+use super::{PullReceipt, PullTermination, summary::PullTargetSummaries};
+use crate::{
+ ingest::IngestReceipt,
+ policy::{Error, SyncId},
+};
+use radroots_transport::{
+ outcome::{FetchTargetOutcome, FetchTargetState},
+ source::FetchCursor,
+};
+
+impl<'de> serde::Deserialize<'de> for PullReceipt {
+ fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
+ where
+ D: serde::Deserializer<'de>,
+ {
+ // Preserve the existing receipt's legacy fields and unknown-field policy.
+ #[derive(serde::Deserialize)]
+ struct Wire {
+ sync_id: SyncId,
+ deadline_unix_ms: u64,
+ pages_fetched: u16,
+ events_observed: usize,
+ ingest_outcomes: Vec<Result<IngestReceipt, Error>>,
+ target_outcomes: Vec<FetchTargetOutcome>,
+ #[serde(default)]
+ target_summaries: Option<PullTargetSummaries>,
+ termination: PullTermination,
+ resume_from: Option<FetchCursor>,
+ }
+ let wire = Wire::deserialize(deserializer)?;
+ if let Some(summaries) = &wire.target_summaries {
+ let summaries = summaries.as_slice();
+ if summaries
+ .iter()
+ .any(|summary| summary.pages_observed() != wire.pages_fetched)
+ {
+ return Err(serde::de::Error::custom("pull summary page count differs"));
+ }
+ let mut targets = std::collections::BTreeSet::new();
+ for outcome in &wire.target_outcomes {
+ let Some(summary) = summaries
+ .iter()
+ .find(|summary| summary.target() == outcome.target())
+ else {
+ return Err(serde::de::Error::custom(
+ "pull summary target inventory differs",
+ ));
+ };
+ if !targets.insert(outcome.target()) {
+ return Err(serde::de::Error::custom("duplicate pull final target"));
+ }
+ let consistent = match outcome.state() {
+ FetchTargetState::Complete => {
+ summary.pages_observed()
+ > summary.incomplete_pages() + summary.missing_outcome_pages()
+ }
+ state => summary.last_incomplete() == Some(state),
+ };
+ if !consistent {
+ return Err(serde::de::Error::custom("pull summary final state differs"));
+ }
+ }
+ for summary in summaries {
+ let observed_outcome = summary.pages_observed() > summary.missing_outcome_pages();
+ if targets.contains(summary.target()) != observed_outcome {
+ return Err(serde::de::Error::custom(
+ "pull summary outcome evidence differs",
+ ));
+ }
+ }
+ }
+ Ok(Self {
+ sync_id: wire.sync_id,
+ deadline_unix_ms: wire.deadline_unix_ms,
+ pages_fetched: wire.pages_fetched,
+ events_observed: wire.events_observed,
+ ingest_outcomes: wire.ingest_outcomes,
+ target_outcomes: wire.target_outcomes,
+ target_summaries: wire.target_summaries,
+ termination: wire.termination,
+ resume_from: wire.resume_from,
+ })
+ }
+}
diff --git a/crates/sync/tests/package_boundary.rs b/crates/sync/tests/package_boundary.rs
@@ -9,11 +9,43 @@ const SOURCES: &[(&str, &str)] = &[
("policy.rs", include_str!("../src/policy.rs")),
("projection.rs", include_str!("../src/projection.rs")),
("pull.rs", include_str!("../src/pull.rs")),
+ ("pull/summary.rs", include_str!("../src/pull/summary.rs")),
+ ("pull/wire.rs", include_str!("../src/pull/wire.rs")),
("push.rs", include_str!("../src/push.rs")),
("status.rs", include_str!("../src/status.rs")),
];
#[test]
+fn pull_summary_contract_retains_existing_bounds_and_final_outcome_compatibility() {
+ let contract: serde_json::Value = serde_json::from_str(include_str!(
+ "../../../contracts/architecture/decisions/pull_target_evidence.v1.json"
+ ))
+ .expect("pull target evidence contract");
+ assert_eq!(contract["owner"], "radroots_sync");
+ assert_eq!(
+ contract["max_targets"],
+ radroots_transport::target::TARGET_SET_MAX_ITEMS
+ );
+ assert_eq!(contract["max_pages"], radroots_sync::pull::PULL_MAX_PAGES);
+ assert_eq!(
+ contract["summary_fields"],
+ serde_json::json!([
+ "target",
+ "pages_observed",
+ "incomplete_pages",
+ "missing_outcome_pages",
+ "last_incomplete"
+ ])
+ );
+ assert!(
+ contract["compatibility"]
+ .as_str()
+ .unwrap()
+ .contains("legacy")
+ );
+}
+
+#[test]
fn sync_depends_only_on_final_orchestration_boundaries() {
for required in [
"name = \"radroots_sync\"",
diff --git a/crates/sync/tests/pull.rs b/crates/sync/tests/pull.rs
@@ -3,6 +3,9 @@ use std::{
sync::{Arc, Mutex},
};
+#[path = "pull/summary.rs"]
+mod summary;
+
use futures_executor::block_on;
use radroots_event::{SignedEvent, draft::SignedEventParts};
use radroots_storage::{event::SourceGeneration, memory::MemoryStorage};
@@ -174,7 +177,7 @@ fn observed(signature: &str, observed_at: u64) -> ObservedEvent {
)
}
-fn engine(source: Arc<ScriptedSource>, clock: Arc<dyn Clock>, timeout_ms: u64) -> Engine {
+fn engine(source: Arc<dyn EventSource>, clock: Arc<dyn Clock>, timeout_ms: u64) -> Engine {
let storage: Arc<dyn SyncStorage> = Arc::new(MemoryStorage::new(
SourceGeneration::new([8; 32]).expect("generation"),
));
@@ -247,6 +250,37 @@ fn single_and_multiple_pages_propagate_cursor_deadline_and_ingest_results() {
assert_eq!(requests[0].deadline_unix_ms, requests[1].deadline_unix_ms);
}
+#[cfg(feature = "serde")]
+#[test]
+fn later_complete_page_retains_earlier_incomplete_evidence() {
+ let source = Arc::new(ScriptedSource::new(vec![
+ Response::Page {
+ events: vec![],
+ state: FetchTargetState::Partial,
+ next: NextPage::Cursor(FetchCursor::parse("next").expect("cursor")),
+ },
+ Response::Page {
+ events: vec![],
+ state: FetchTargetState::Complete,
+ next: NextPage::Complete,
+ },
+ ]));
+ let pull = engine(source, Arc::new(FixedClock(100)), 50);
+ let receipt = block_on(pull.pull(
+ PullRequest::new(targets(), 10, 2).expect("request"),
+ &RegistryPolicy::visible(),
+ ))
+ .expect("pull");
+ assert_eq!(receipt.termination(), PullTermination::Complete);
+ assert_eq!(
+ receipt.target_outcomes()[0].state(),
+ FetchTargetState::Complete
+ );
+ let wire = serde_json::to_value(receipt).expect("receipt JSON");
+ assert_eq!(wire["target_summaries"][0]["incomplete_pages"], 1);
+ assert_eq!(wire["target_summaries"][0]["last_incomplete"], "partial");
+}
+
#[test]
fn pull_propagates_the_exact_selector_to_every_page() {
let next = FetchCursor::parse("page-2").expect("cursor");
diff --git a/crates/sync/tests/pull/summary.rs b/crates/sync/tests/pull/summary.rs
@@ -0,0 +1,397 @@
+use super::*;
+use radroots_sync::PullReceipt;
+use radroots_transport::target::TARGET_SET_MAX_ITEMS;
+
+struct Source {
+ pages: Mutex<VecDeque<Option<Vec<Option<FetchTargetState>>>>>,
+ finish: NextPage,
+}
+
+impl EventSource for Source {
+ fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SourceStatus, TransportError>> {
+ Box::pin(async { unreachable!("explicit pull only") })
+ }
+
+ fn fetch(
+ &self,
+ request: FetchRequest,
+ ) -> radroots_transport::BoxFuture<'_, Result<FetchPage, TransportError>> {
+ Box::pin(async move {
+ let mut pages = self.pages.lock().expect("pages");
+ let states = pages
+ .pop_front()
+ .expect("bounded scripted fetch")
+ .ok_or(TransportError::UnsupportedOperation)?;
+ assert_eq!(states.len(), request.target_set().len());
+ let outcomes = request
+ .target_set()
+ .targets()
+ .iter()
+ .zip(states)
+ .filter_map(|(target, state)| {
+ state.map(|state| FetchTargetOutcome::new(target.fingerprint().clone(), state))
+ })
+ .collect();
+ let next = if pages.is_empty() {
+ self.finish.clone()
+ } else {
+ NextPage::Cursor(
+ FetchCursor::parse(format!("remaining-{}", pages.len())).expect("cursor"),
+ )
+ };
+ FetchPage::for_request(&request, vec![], outcomes, next)
+ })
+ }
+}
+
+fn target_set(count: usize) -> TargetSet {
+ TargetSet::new(
+ (0..count)
+ .rev()
+ .map(|index| {
+ Target::nostr_relay(format!("wss://relay-{index}.example")).expect("target")
+ })
+ .collect(),
+ )
+ .expect("targets")
+}
+
+fn pull(
+ pages: Vec<Option<Vec<Option<FetchTargetState>>>>,
+ targets: TargetSet,
+ max_pages: u16,
+ finish: NextPage,
+ clock: Arc<dyn Clock>,
+) -> PullReceipt {
+ let source = Arc::new(Source {
+ pages: Mutex::new(pages.into()),
+ finish,
+ });
+ block_on(engine(source, clock, 50).pull(
+ PullRequest::new(targets, 1, max_pages).expect("request"),
+ &RegistryPolicy::visible(),
+ ))
+ .expect("receipt")
+}
+
+fn ordered(states: &[Option<FetchTargetState>]) -> PullReceipt {
+ pull(
+ states.iter().map(|state| Some(vec![*state])).collect(),
+ target_set(1),
+ states.len() as u16,
+ NextPage::Complete,
+ Arc::new(FixedClock(100)),
+ )
+}
+
+#[test]
+fn every_incomplete_state_survives_later_success_and_missing_outcomes() {
+ use FetchTargetState::*;
+ for state in [
+ Partial,
+ Unavailable,
+ FailedRetryable,
+ FailedTerminal,
+ Cancelled,
+ ] {
+ for states in [
+ vec![Some(state), Some(Complete)],
+ vec![Some(Complete), Some(state)],
+ vec![Some(state), None, Some(Complete)],
+ vec![Some(state), Some(Complete), None],
+ ] {
+ let receipt = ordered(&states);
+ let summaries = receipt.target_summaries().expect("measured");
+ assert_eq!(summaries.len(), 1);
+ let summary = &summaries[0];
+ assert_eq!(summary.target(), target_set(1).targets()[0].fingerprint());
+ assert_eq!(summary.pages_observed(), states.len() as u16);
+ assert_eq!(summary.incomplete_pages(), 1);
+ assert_eq!(
+ summary.missing_outcome_pages(),
+ states.iter().filter(|value| value.is_none()).count() as u16
+ );
+ assert_eq!(summary.last_incomplete(), Some(state));
+ assert!(!summary.all_pages_complete());
+ assert_eq!(
+ receipt.target_outcomes()[0].state(),
+ states.iter().rev().flatten().next().copied().unwrap()
+ );
+ #[cfg(feature = "serde")]
+ assert_eq!(
+ serde_json::from_str::<PullReceipt>(&serde_json::to_string(&receipt).unwrap())
+ .unwrap(),
+ receipt
+ );
+ }
+ }
+ let receipt = ordered(&[Some(Partial), Some(FailedTerminal), Some(Complete)]);
+ let summary = &receipt.target_summaries().unwrap()[0];
+ assert_eq!(summary.incomplete_pages(), 2);
+ assert_eq!(summary.last_incomplete(), Some(FailedTerminal));
+}
+
+#[test]
+fn omitted_outcomes_never_become_positive_evidence() {
+ use FetchTargetState::Complete;
+ for states in [
+ vec![None],
+ vec![None, Some(Complete)],
+ vec![Some(Complete), None],
+ ] {
+ let receipt = ordered(&states);
+ assert_eq!(receipt.termination(), PullTermination::Complete);
+ let summary = &receipt.target_summaries().unwrap()[0];
+ assert_eq!(summary.missing_outcome_pages(), 1);
+ assert_eq!(summary.incomplete_pages(), 0);
+ assert_eq!(summary.last_incomplete(), None);
+ assert!(!summary.all_pages_complete());
+ #[cfg(feature = "serde")]
+ assert_eq!(
+ serde_json::from_value::<PullReceipt>(serde_json::to_value(&receipt).unwrap()).unwrap(),
+ receipt
+ );
+ }
+}
+
+#[test]
+fn maximum_inventory_preserves_request_order_and_counts_without_page_history() {
+ let targets = target_set(TARGET_SET_MAX_ITEMS);
+ let receipt = pull(
+ vec![
+ Some(vec![Some(FetchTargetState::Complete); TARGET_SET_MAX_ITEMS]);
+ usize::from(PULL_MAX_PAGES)
+ ],
+ targets.clone(),
+ PULL_MAX_PAGES,
+ NextPage::Complete,
+ Arc::new(FixedClock(100)),
+ );
+ assert_eq!(receipt.pages_fetched(), PULL_MAX_PAGES);
+ let summaries = receipt.target_summaries().unwrap();
+ assert_eq!(summaries.len(), TARGET_SET_MAX_ITEMS);
+ for (summary, target) in summaries.iter().zip(targets.targets()) {
+ assert_eq!(summary.target(), target.fingerprint());
+ assert_eq!(summary.pages_observed(), PULL_MAX_PAGES);
+ assert!(summary.all_pages_complete());
+ }
+ #[cfg(feature = "serde")]
+ {
+ let encoded = serde_json::to_string(&receipt).unwrap();
+ assert!(
+ encoded.len() < 32_768,
+ "bounded evidence excludes per-page history"
+ );
+ assert_eq!(
+ serde_json::from_str::<PullReceipt>(&encoded).unwrap(),
+ receipt
+ );
+ }
+}
+
+#[test]
+fn multiple_targets_retain_independent_evidence() {
+ use FetchTargetState::*;
+ let receipt = pull(
+ vec![
+ Some(vec![Some(Partial), Some(Complete), None]),
+ Some(vec![Some(Complete), Some(FailedRetryable), Some(Complete)]),
+ ],
+ target_set(3),
+ 2,
+ NextPage::Complete,
+ Arc::new(FixedClock(100)),
+ );
+ let summaries = receipt.target_summaries().unwrap();
+ assert_eq!(
+ summaries
+ .iter()
+ .map(|s| s.incomplete_pages())
+ .collect::<Vec<_>>(),
+ [1, 1, 0]
+ );
+ assert_eq!(
+ summaries
+ .iter()
+ .map(|s| s.missing_outcome_pages())
+ .collect::<Vec<_>>(),
+ [0, 0, 1]
+ );
+ assert_eq!(
+ summaries
+ .iter()
+ .map(|s| s.last_incomplete())
+ .collect::<Vec<_>>(),
+ [Some(Partial), Some(FailedRetryable), None]
+ );
+ assert!(summaries.iter().all(|s| !s.all_pages_complete()));
+}
+
+#[test]
+fn termination_and_zero_returned_pages_remain_explicit() {
+ use FetchTargetState::*;
+ let failed = pull(
+ vec![None],
+ target_set(1),
+ 1,
+ NextPage::Complete,
+ Arc::new(FixedClock(100)),
+ );
+ assert_eq!(failed.termination(), PullTermination::SourceFailed);
+ assert_eq!(failed.target_summaries().unwrap()[0].pages_observed(), 0);
+ assert!(!failed.target_summaries().unwrap()[0].all_pages_complete());
+ let later_failure = pull(
+ vec![Some(vec![Some(Complete)]), None],
+ target_set(1),
+ 2,
+ NextPage::Complete,
+ Arc::new(FixedClock(100)),
+ );
+ assert_eq!(later_failure.termination(), PullTermination::SourceFailed);
+ assert!(later_failure.target_summaries().unwrap()[0].all_pages_complete());
+ assert_eq!(
+ later_failure.target_summaries().unwrap()[0].pages_observed(),
+ 1
+ );
+ let limited = pull(
+ vec![Some(vec![Some(Partial)]); 2],
+ target_set(1),
+ 1,
+ NextPage::Complete,
+ Arc::new(FixedClock(100)),
+ );
+ assert_eq!(limited.termination(), PullTermination::PageLimit);
+ let deadline = pull(
+ vec![Some(vec![Some(Partial)]); 2],
+ target_set(1),
+ 2,
+ NextPage::Complete,
+ Arc::new(DeadlineClock(Mutex::new(VecDeque::from([100, 150])))),
+ );
+ assert_eq!(deadline.termination(), PullTermination::Deadline);
+ let cancelled = pull(
+ vec![Some(vec![Some(Cancelled)])],
+ target_set(1),
+ 1,
+ NextPage::Cancelled { resume_from: None },
+ Arc::new(FixedClock(100)),
+ );
+ assert_eq!(cancelled.termination(), PullTermination::Cancelled);
+ for receipt in [failed, later_failure, limited, deadline, cancelled] {
+ #[cfg(feature = "serde")]
+ assert_eq!(
+ serde_json::from_value::<PullReceipt>(serde_json::to_value(&receipt).unwrap()).unwrap(),
+ receipt
+ );
+ assert_eq!(
+ receipt.target_summaries().unwrap()[0].pages_observed(),
+ receipt.pages_fetched()
+ );
+ }
+}
+
+#[cfg(feature = "serde")]
+#[test]
+fn legacy_receipts_remain_unknown_and_new_receipts_reject_inconsistent_evidence() {
+ use serde_json::json;
+ let receipt = ordered(&[Some(FetchTargetState::Complete)]);
+ let original = serde_json::to_value(&receipt).unwrap();
+ let mut legacy = original.clone();
+ legacy.as_object_mut().unwrap().remove("target_summaries");
+ assert!(
+ serde_json::from_value::<PullReceipt>(legacy.clone())
+ .unwrap()
+ .target_summaries()
+ .is_none()
+ );
+ legacy["target_summaries"] = json!(null);
+ assert!(
+ serde_json::from_value::<PullReceipt>(legacy)
+ .unwrap()
+ .target_summaries()
+ .is_none()
+ );
+ for replacement in [json!([]), json!({}), json!(42)] {
+ let mut changed = original.clone();
+ changed["target_summaries"] = replacement;
+ assert!(serde_json::from_value::<PullReceipt>(changed).is_err());
+ }
+ let mut duplicate = original.clone();
+ duplicate["target_summaries"]
+ .as_array_mut()
+ .unwrap()
+ .push(original["target_summaries"][0].clone());
+ assert!(serde_json::from_value::<PullReceipt>(duplicate).is_err());
+ for (field, value) in [
+ ("pages_observed", json!(1001)),
+ ("pages_observed", json!(2)),
+ ("incomplete_pages", json!(1)),
+ ("incomplete_pages", json!(65535)),
+ ("missing_outcome_pages", json!(2)),
+ ("missing_outcome_pages", json!(1)),
+ ("last_incomplete", json!("partial")),
+ ("target", json!("bad")),
+ ("extra", json!(true)),
+ ] {
+ let mut changed = original.clone();
+ changed["target_summaries"][0][field] = value;
+ assert!(
+ serde_json::from_value::<PullReceipt>(changed).is_err(),
+ "{field}"
+ );
+ }
+ let mut complete_as_incomplete = original.clone();
+ complete_as_incomplete["target_summaries"][0]["incomplete_pages"] = json!(1);
+ complete_as_incomplete["target_summaries"][0]["last_incomplete"] = json!("complete");
+ assert!(serde_json::from_value::<PullReceipt>(complete_as_incomplete).is_err());
+ for change in 0..5 {
+ let mut changed = original.clone();
+ match change {
+ 0 => {
+ changed["target_outcomes"] = json!([]);
+ }
+ 1 => {
+ changed["target_outcomes"]
+ .as_array_mut()
+ .unwrap()
+ .push(original["target_outcomes"][0].clone());
+ }
+ 2 => {
+ changed["target_outcomes"][0]["target"] =
+ json!(target_set(2).targets()[0].fingerprint());
+ }
+ 3 => {
+ changed["target_outcomes"][0]["state"] = json!("partial");
+ }
+ _ => {
+ changed["target_summaries"][0]["incomplete_pages"] = json!(1);
+ changed["target_summaries"][0]["last_incomplete"] = json!("partial");
+ }
+ }
+ assert!(
+ serde_json::from_value::<PullReceipt>(changed).is_err(),
+ "case {change}"
+ );
+ }
+ let mut all_missing = serde_json::to_value(ordered(&[None])).unwrap();
+ all_missing["target_summaries"][0]["missing_outcome_pages"] = json!(0);
+ assert!(serde_json::from_value::<PullReceipt>(all_missing).is_err());
+ let maximum = pull(
+ vec![Some(vec![
+ Some(FetchTargetState::Complete);
+ TARGET_SET_MAX_ITEMS
+ ])],
+ target_set(TARGET_SET_MAX_ITEMS),
+ 1,
+ NextPage::Complete,
+ Arc::new(FixedClock(100)),
+ );
+ let mut too_many = serde_json::to_value(maximum).unwrap();
+ too_many["target_summaries"]
+ .as_array_mut()
+ .unwrap()
+ .push(json!({"unparsed_extra": [1, 2, 3]}));
+ let error = serde_json::from_str::<PullReceipt>(&serde_json::to_string(&too_many).unwrap())
+ .unwrap_err();
+ assert!(error.to_string().contains("too many pull target summaries"));
+}