commit 991d63b9906070a827e235d3e7824c80508d494e
parent 85e78fea1feca5bd030b0156cae9fe570402c8a0
Author: triesap <tyson@radroots.org>
Date: Mon, 3 Aug 2026 08:03:08 +0000
sync: implement canonical ingest orchestration
- add concrete BIP-340 verification and explicit host admission decisions
- atomically persist canonical event state with exact transport provenance
- normalize durable receipts, duplicate handling, conflicts, and partial batches
- verify feature matrices, focused tests, strict linting, and architecture gates
Diffstat:
6 files changed, 549 insertions(+), 0 deletions(-)
diff --git a/Cargo.lock b/Cargo.lock
@@ -5155,6 +5155,7 @@ dependencies = [
name = "radroots_sync"
version = "0.1.0-alpha"
dependencies = [
+ "futures-executor",
"radroots_event",
"radroots_event_codec",
"radroots_protocol",
@@ -5163,6 +5164,7 @@ dependencies = [
"radroots_trade",
"radroots_transport",
"serde",
+ "sha2",
]
[[package]]
diff --git a/crates/event_codec/src/verify.rs b/crates/event_codec/src/verify.rs
@@ -12,6 +12,29 @@ pub use radroots_event::admission::{
// namespace while consumers move away from the superseded crate-root exports.
pub use crate::verification::*;
+/// Deterministic BIP-340 verifier for canonical NIP-01 event envelopes.
+///
+/// This capability performs cryptographic verification only. Contract
+/// validation, host admission, and visibility remain later explicit stages.
+#[derive(Clone, Copy, Debug, Default)]
+pub struct Nip01SignatureVerifier;
+
+impl SignatureVerifier for Nip01SignatureVerifier {
+ fn verify_signature(
+ &self,
+ event: &radroots_event::envelope::EventEnvelope,
+ ) -> Result<(), VerificationError> {
+ let public_key = secp256k1::XOnlyPublicKey::from_slice(event.author().as_bytes())
+ .map_err(|_| VerificationError::MalformedEnvelope)?;
+ let signature = secp256k1::schnorr::Signature::from_slice(event.sig().as_bytes())
+ .map_err(|_| VerificationError::MalformedEnvelope)?;
+ let message = secp256k1::Message::from_digest(*event.id().as_bytes());
+ secp256k1::Secp256k1::verification_only()
+ .verify_schnorr(&signature, &message, &public_key)
+ .map_err(|_| VerificationError::SignatureInvalid)
+ }
+}
+
/// Verifies that an event's declared identifier matches its canonical bytes.
pub fn id(event: RawEvent) -> Result<IdVerifiedEvent, VerificationError> {
event.verify_id()
diff --git a/crates/sync/Cargo.toml b/crates/sync/Cargo.toml
@@ -36,8 +36,10 @@ radroots_storage = { workspace = true, default-features = false }
radroots_trade = { workspace = true, default-features = false }
radroots_transport = { workspace = true, default-features = false }
serde = { workspace = true, optional = true }
+sha2 = { workspace = true, default-features = false }
[dev-dependencies]
+futures-executor = { workspace = true }
radroots_storage = { workspace = true, features = ["memory"] }
[lints]
diff --git a/crates/sync/src/ingest.rs b/crates/sync/src/ingest.rs
@@ -1 +1,282 @@
//! Canonical event verification and admission orchestration.
+
+use radroots_event::admission::{
+ AdmissionPolicy as EventAdmissionPolicy, AdmittedEvent, ContractValidatedEvent, RawEvent,
+ VisibilityPolicy,
+};
+use radroots_event_codec::verify::{self, Nip01SignatureVerifier};
+use radroots_storage::{
+ Error as StorageError,
+ atomic::{
+ AtomicCommit, AtomicCommitDigest, AtomicCommitDisposition, AtomicCommitId,
+ AtomicCommitOutcome, AtomicWorkflow, CommitIngested,
+ },
+ event::{AdmissionReceipt, EventAdmission},
+};
+use radroots_transport::source::ObservedEvent;
+use sha2::{Digest, Sha256};
+
+use crate::{
+ Engine,
+ policy::{Error, OperationKind, SyncId},
+};
+
+/// Host decision applied after cryptographic and contract verification.
+#[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 AdmissionDecision {
+ /// Reject the event without mutating canonical storage.
+ Reject,
+ /// Persist the signature-verified event without making it visible.
+ Verified,
+ /// Persist the event as admitted and visible.
+ Visible,
+}
+
+/// Deterministic host policy for canonical admission and visibility.
+///
+/// Implementations must be side-effect free. The engine may evaluate a policy
+/// once per observed event and owns all durable effects after that decision.
+pub trait AdmissionPolicy: Send + Sync {
+ /// Stable policy identity retained in event typestate evidence.
+ fn policy_id(&self) -> &'static str;
+
+ /// Decides whether a contract-valid event is rejected, verified-only, or visible.
+ fn decide(&self, event: &ContractValidatedEvent) -> AdmissionDecision;
+}
+
+/// Registry-valid policy useful for hosts that do not add product authorization.
+#[derive(Clone, Copy, Debug, Eq, PartialEq)]
+pub struct RegistryPolicy {
+ decision: AdmissionDecision,
+}
+
+impl RegistryPolicy {
+ /// Admits contract-valid events without authorizing visibility.
+ pub const fn verified() -> Self {
+ Self {
+ decision: AdmissionDecision::Verified,
+ }
+ }
+
+ /// Admits contract-valid events and authorizes visibility.
+ pub const fn visible() -> Self {
+ Self {
+ decision: AdmissionDecision::Visible,
+ }
+ }
+}
+
+impl AdmissionPolicy for RegistryPolicy {
+ fn policy_id(&self) -> &'static str {
+ "radroots.registry_v7"
+ }
+
+ fn decide(&self, _event: &ContractValidatedEvent) -> AdmissionDecision {
+ self.decision
+ }
+}
+
+/// Normalized durable result for one observed event.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct IngestReceipt {
+ sync_id: SyncId,
+ commit_disposition: AtomicCommitDisposition,
+ admission: AdmissionReceipt,
+ committed_at_unix_ms: u64,
+}
+
+impl IngestReceipt {
+ pub const fn sync_id(&self) -> SyncId {
+ self.sync_id
+ }
+
+ pub const fn commit_disposition(&self) -> AtomicCommitDisposition {
+ self.commit_disposition
+ }
+
+ pub const fn admission(&self) -> &AdmissionReceipt {
+ &self.admission
+ }
+
+ pub const fn committed_at_unix_ms(&self) -> u64 {
+ self.committed_at_unix_ms
+ }
+}
+
+/// Ordered independent outcomes for one bounded caller-supplied batch.
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct IngestBatchReceipt {
+ outcomes: Vec<Result<IngestReceipt, Error>>,
+}
+
+impl IngestBatchReceipt {
+ pub fn outcomes(&self) -> &[Result<IngestReceipt, Error>] {
+ self.outcomes.as_slice()
+ }
+
+ pub fn accepted(&self) -> usize {
+ self.outcomes
+ .iter()
+ .filter(|outcome| outcome.is_ok())
+ .count()
+ }
+
+ pub fn rejected(&self) -> usize {
+ self.outcomes.len() - self.accepted()
+ }
+}
+
+impl Engine {
+ /// Verifies and atomically persists one exact transport observation.
+ pub async fn ingest(
+ &self,
+ observed: ObservedEvent,
+ policy: &dyn AdmissionPolicy,
+ ) -> Result<IngestReceipt, Error> {
+ let verified = verify::signature(
+ verify::id(RawEvent::new(observed.event().envelope().clone()))
+ .map_err(|_| Error::VerificationFailed)?,
+ &Nip01SignatureVerifier,
+ )
+ .map_err(|_| Error::VerificationFailed)?;
+ let validated =
+ verify::contract(verified.clone()).map_err(|_| Error::VerificationFailed)?;
+ let decision = policy.decide(&validated);
+ if decision == AdmissionDecision::Reject {
+ return Err(Error::PolicyRejected);
+ }
+
+ let admission = match decision {
+ AdmissionDecision::Verified => EventAdmission::verified(observed.clone(), verified),
+ AdmissionDecision::Visible => {
+ let evidence = DecisionEvidence {
+ policy_id: policy.policy_id(),
+ };
+ let visible = validated
+ .admit_with(&evidence)
+ .and_then(|event| event.make_visible_with(&evidence))
+ .map_err(|never| match never {})?;
+ EventAdmission::visible(observed.clone(), visible)
+ }
+ AdmissionDecision::Reject => unreachable!("rejection returned before admission"),
+ }
+ .map_err(map_storage_error)?;
+
+ let sync_id = self.ids.next_id(OperationKind::Ingest)?;
+ let requested_at_unix_ms = self.clock.now_unix_ms()?;
+ let digest = ingest_digest(&observed, policy.policy_id(), decision);
+ let request = AtomicCommit::new(
+ AtomicCommitId::new(*sync_id.as_bytes()).map_err(map_storage_error)?,
+ digest,
+ requested_at_unix_ms,
+ AtomicWorkflow::Ingested(Box::new(CommitIngested::new(admission, None))),
+ )
+ .map_err(map_storage_error)?;
+ let receipt = self
+ .storage
+ .commit(request)
+ .await
+ .map_err(map_storage_error)?;
+ let AtomicCommitOutcome::Ingested { admission, .. } = receipt.outcome() else {
+ return Err(Error::InvalidIngestReceipt);
+ };
+ if admission.event_id() != observed.event().id() {
+ return Err(Error::InvalidIngestReceipt);
+ }
+ Ok(IngestReceipt {
+ sync_id,
+ commit_disposition: receipt.disposition(),
+ admission: admission.clone(),
+ committed_at_unix_ms: receipt.committed_at_unix_ms(),
+ })
+ }
+
+ /// Ingests each supplied observation independently and preserves input order.
+ ///
+ /// A rejected or failed item does not suppress later items and no hidden
+ /// retry, polling, or scheduling loop is created.
+ pub async fn ingest_batch(
+ &self,
+ observed: Vec<ObservedEvent>,
+ policy: &dyn AdmissionPolicy,
+ ) -> IngestBatchReceipt {
+ let mut outcomes = Vec::with_capacity(observed.len());
+ for event in observed {
+ outcomes.push(self.ingest(event, policy).await);
+ }
+ IngestBatchReceipt { outcomes }
+ }
+}
+
+struct DecisionEvidence {
+ policy_id: &'static str,
+}
+
+impl EventAdmissionPolicy for DecisionEvidence {
+ type Error = core::convert::Infallible;
+
+ fn policy_id(&self) -> &'static str {
+ self.policy_id
+ }
+
+ fn admit(&self, _event: &ContractValidatedEvent) -> Result<(), Self::Error> {
+ Ok(())
+ }
+}
+
+impl VisibilityPolicy for DecisionEvidence {
+ type Error = core::convert::Infallible;
+
+ fn policy_id(&self) -> &'static str {
+ self.policy_id
+ }
+
+ fn make_visible(&self, _event: &AdmittedEvent) -> Result<(), Self::Error> {
+ Ok(())
+ }
+}
+
+fn ingest_digest(
+ observed: &ObservedEvent,
+ policy_id: &str,
+ decision: AdmissionDecision,
+) -> AtomicCommitDigest {
+ let provenance = observed.provenance();
+ let mut hasher = Sha256::new();
+ hash_field(&mut hasher, b"radroots.sync.ingest.v1");
+ hash_field(&mut hasher, observed.event().raw_json().as_bytes());
+ hash_field(&mut hasher, provenance.transport_id().as_str().as_bytes());
+ hash_field(&mut hasher, provenance.target().as_str().as_bytes());
+ hash_field(
+ &mut hasher,
+ provenance.observed_at_unix_ms().to_be_bytes().as_slice(),
+ );
+ hash_field(
+ &mut hasher,
+ provenance
+ .cursor()
+ .map_or(&[][..], |cursor| cursor.as_str().as_bytes()),
+ );
+ hash_field(&mut hasher, policy_id.as_bytes());
+ hasher.update([match decision {
+ AdmissionDecision::Reject => 0,
+ AdmissionDecision::Verified => 1,
+ AdmissionDecision::Visible => 2,
+ }]);
+ AtomicCommitDigest::new(hasher.finalize().into())
+}
+
+fn hash_field(hasher: &mut Sha256, value: &[u8]) {
+ hasher.update(u64::try_from(value.len()).unwrap_or(u64::MAX).to_be_bytes());
+ hasher.update(value);
+}
+
+fn map_storage_error(error: StorageError) -> Error {
+ match error {
+ StorageError::EventConflict | StorageError::AtomicCommitConflict => Error::StorageConflict,
+ _ => Error::StorageFailed,
+ }
+}
diff --git a/crates/sync/src/policy.rs b/crates/sync/src/policy.rs
@@ -16,6 +16,7 @@ const MAX_OPERATION_TIMEOUT_MS: u64 = 86_400_000;
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum OperationKind {
+ Ingest,
Pull,
Sign,
Deliver,
@@ -83,6 +84,9 @@ impl DeadlinePolicy {
pub const fn timeout_ms(self, operation: OperationKind) -> u64 {
match operation {
+ // Ingest performs local verification and one atomic commit. It
+ // shares the inbound operation budget with pull orchestration.
+ OperationKind::Ingest => self.pull_timeout_ms,
OperationKind::Pull => self.pull_timeout_ms,
OperationKind::Sign => self.sign_timeout_ms,
OperationKind::Deliver => self.delivery_timeout_ms,
@@ -183,6 +187,11 @@ pub enum Error {
DeadlineOverflow,
MissingTransportCapability,
SignerWithoutSink,
+ VerificationFailed,
+ PolicyRejected,
+ StorageConflict,
+ StorageFailed,
+ InvalidIngestReceipt,
}
impl core::fmt::Display for Error {
@@ -194,6 +203,11 @@ impl core::fmt::Display for Error {
Self::DeadlineOverflow => "sync deadline overflowed",
Self::MissingTransportCapability => "sync engine requires a source or sink",
Self::SignerWithoutSink => "sync signer requires a sink",
+ Self::VerificationFailed => "sync event verification failed",
+ Self::PolicyRejected => "sync admission policy rejected the event",
+ Self::StorageConflict => "sync input conflicts with durable storage state",
+ Self::StorageFailed => "sync storage operation failed",
+ Self::InvalidIngestReceipt => "sync storage returned an invalid ingest receipt",
})
}
}
diff --git a/crates/sync/tests/ingest.rs b/crates/sync/tests/ingest.rs
@@ -0,0 +1,227 @@
+use std::sync::{
+ Arc,
+ atomic::{AtomicU8, Ordering},
+};
+
+use futures_executor::block_on;
+use radroots_event::{SignedEvent, draft::SignedEventParts};
+use radroots_storage::{
+ EventStore, Storage,
+ event::{AdmissionDisposition, AdmissionStage, EventQuery, EventQueryBounds, SourceGeneration},
+ memory::MemoryStorage,
+};
+use radroots_sync::{
+ Engine,
+ ingest::{AdmissionDecision, AdmissionPolicy, RegistryPolicy},
+ policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId},
+};
+use radroots_transport::{
+ Error as TransportError, EventSource, FetchPage, FetchRequest, SourceStatus, Target,
+ TransportId,
+ source::{EventProvenance, 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}";
+
+struct MockSource;
+
+impl EventSource for MockSource {
+ fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SourceStatus, TransportError>> {
+ Box::pin(async { unreachable!("ingest does not inspect source status") })
+ }
+
+ fn fetch(
+ &self,
+ _request: FetchRequest,
+ ) -> radroots_transport::BoxFuture<'_, Result<FetchPage, TransportError>> {
+ Box::pin(async { unreachable!("ingest does not fetch") })
+ }
+}
+
+struct FixedClock;
+
+impl Clock for FixedClock {
+ fn now_unix_ms(&self) -> Result<u64, Error> {
+ Ok(1_800_000_200_000)
+ }
+}
+
+struct SequenceIds(AtomicU8);
+
+impl IdSource for SequenceIds {
+ fn next_id(&self, operation: OperationKind) -> Result<SyncId, Error> {
+ assert_eq!(operation, OperationKind::Ingest);
+ let byte = self.0.fetch_add(1, Ordering::Relaxed);
+ SyncId::new([byte; 16])
+ }
+}
+
+struct ConstantIds(u8);
+
+impl IdSource for ConstantIds {
+ fn next_id(&self, operation: OperationKind) -> Result<SyncId, Error> {
+ assert_eq!(operation, OperationKind::Ingest);
+ SyncId::new([self.0; 16])
+ }
+}
+
+struct Reject;
+
+impl AdmissionPolicy for Reject {
+ fn policy_id(&self) -> &'static str {
+ "test.reject.v1"
+ }
+
+ fn decide(
+ &self,
+ _event: &radroots_event::admission::ContractValidatedEvent,
+ ) -> AdmissionDecision {
+ AdmissionDecision::Reject
+ }
+}
+
+fn setup_engine(first_id: u8) -> (Engine, Arc<MemoryStorage>) {
+ setup_engine_with_ids(Arc::new(SequenceIds(AtomicU8::new(first_id))))
+}
+
+fn setup_engine_with_ids(ids: Arc<dyn IdSource>) -> (Engine, Arc<MemoryStorage>) {
+ let storage = Arc::new(MemoryStorage::new(
+ SourceGeneration::new([7; 32]).expect("source generation"),
+ ));
+ let storage_capability: Arc<dyn Storage> = storage.clone();
+ let engine = Engine::builder(
+ storage_capability,
+ Arc::new(FixedClock),
+ ids,
+ DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"),
+ )
+ .source(Arc::new(MockSource))
+ .build()
+ .expect("engine");
+ (engine, storage)
+}
+
+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 signed event")
+}
+
+fn observed(signature: &str, observed_at: u64) -> ObservedEvent {
+ let target = Target::new(TransportId::NOSTR, "wss://relay.example").expect("target");
+ let provenance = EventProvenance::new(
+ TransportId::NOSTR,
+ target.fingerprint().clone(),
+ observed_at,
+ )
+ .expect("provenance");
+ ObservedEvent::new(signed_event(signature), provenance)
+}
+
+#[test]
+fn valid_visible_ingest_is_atomic_and_preserves_provenance() {
+ let (engine, storage) = setup_engine(1);
+ let receipt = block_on(engine.ingest(
+ observed(SIGNATURE, 1_800_000_100_000),
+ &RegistryPolicy::visible(),
+ ))
+ .expect("visible ingest");
+ assert_eq!(receipt.admission().stage(), AdmissionStage::Visible);
+ assert_eq!(
+ receipt.admission().disposition(),
+ AdmissionDisposition::Inserted
+ );
+
+ let bounds = EventQueryBounds::first(10).expect("bounds");
+ let visible = block_on(storage.query_visible(EventQuery::all(bounds))).expect("visible query");
+ assert_eq!(visible.items().len(), 1);
+ let provenance = block_on(storage.query_provenance(*receipt.admission().event_id(), bounds))
+ .expect("provenance query");
+ assert_eq!(provenance.items().len(), 1);
+ assert_eq!(
+ provenance.items()[0].provenance().observed_at_unix_ms(),
+ 1_800_000_100_000
+ );
+}
+
+#[test]
+fn invalid_policy_rejected_and_verified_only_inputs_fail_closed() {
+ let (engine, storage) = setup_engine(1);
+ let invalid_signature = format!("0{}", &SIGNATURE[1..]);
+ assert_eq!(
+ block_on(engine.ingest(observed(&invalid_signature, 1), &RegistryPolicy::visible())),
+ Err(Error::VerificationFailed)
+ );
+ assert_eq!(
+ block_on(engine.ingest(observed(SIGNATURE, 2), &Reject)),
+ Err(Error::PolicyRejected)
+ );
+
+ let receipt = block_on(engine.ingest(observed(SIGNATURE, 3), &RegistryPolicy::verified()))
+ .expect("verified ingest");
+ assert_eq!(receipt.admission().stage(), AdmissionStage::Verified);
+ let page = block_on(storage.query_visible(EventQuery::all(
+ EventQueryBounds::first(10).expect("bounds"),
+ )))
+ .expect("visible query");
+ assert!(page.items().is_empty());
+}
+
+#[test]
+fn duplicate_conflict_and_partial_batch_outcomes_are_normalized() {
+ let (engine, storage) = setup_engine(1);
+ let inserted = block_on(engine.ingest(observed(SIGNATURE, 10), &RegistryPolicy::visible()))
+ .expect("insert");
+ let duplicate = block_on(engine.ingest(observed(SIGNATURE, 11), &RegistryPolicy::visible()))
+ .expect("duplicate");
+ assert_eq!(
+ inserted.admission().position(),
+ duplicate.admission().position()
+ );
+ assert_eq!(
+ duplicate.admission().disposition(),
+ AdmissionDisposition::Duplicate
+ );
+ let provenance = block_on(storage.query_provenance(
+ *inserted.admission().event_id(),
+ EventQueryBounds::first(10).expect("bounds"),
+ ))
+ .expect("provenance query");
+ assert_eq!(provenance.items().len(), 2);
+
+ let (collision_engine, _) = setup_engine_with_ids(Arc::new(ConstantIds(9)));
+ block_on(collision_engine.ingest(observed(SIGNATURE, 20), &RegistryPolicy::visible()))
+ .expect("first identity use");
+ assert_eq!(
+ block_on(collision_engine.ingest(observed(SIGNATURE, 21), &RegistryPolicy::visible())),
+ Err(Error::StorageConflict)
+ );
+
+ let invalid_signature = format!("0{}", &SIGNATURE[1..]);
+ let batch = block_on(engine.ingest_batch(
+ vec![
+ observed(SIGNATURE, 30),
+ observed(&invalid_signature, 31),
+ observed(SIGNATURE, 32),
+ ],
+ &RegistryPolicy::visible(),
+ ));
+ assert_eq!(batch.accepted(), 2);
+ assert_eq!(batch.rejected(), 1);
+ assert_eq!(batch.outcomes()[1], Err(Error::VerificationFailed));
+}