commit db9845f621613817827051d3d3d863597879d567
parent 89406244f80d695a8b0d7d24f8fdebdd164ad12f
Author: triesap <tyson@radroots.org>
Date: Mon, 3 Aug 2026 11:16:26 +0000
sdk: migrate synchronization to radroots_sync
- expose client-scoped canonical sync operations
- return native pull ingest projection push and status types
- remove SDK-owned outbox projection and string mappings
- verify cancellation partial outcomes lifecycle and status
Diffstat:
5 files changed, 409 insertions(+), 3736 deletions(-)
diff --git a/crates/sdk/src/client.rs b/crates/sdk/src/client.rs
@@ -232,11 +232,11 @@ impl Client {
Ok(self.inner.sink.as_deref())
}
- /// Returns the explicit synchronization engine, when configured.
+ /// Returns client-scoped canonical synchronization operations, when configured.
#[cfg(feature = "sync")]
- pub fn sync_engine(&self) -> Result<Option<&radroots_sync::Engine>> {
+ pub fn sync(&self) -> Result<Option<crate::sync::Operations<'_>>> {
self.require_open()?;
- Ok(self.inner.sync.as_ref())
+ Ok(self.inner.sync.as_ref().map(crate::sync::Operations::new))
}
/// Returns whether explicit close completed successfully or reached the
diff --git a/crates/sdk/src/sync.rs b/crates/sdk/src/sync.rs
@@ -1 +1,369 @@
-//! Synchronization capability composition.
+//! Client-scoped access to the canonical synchronization engine.
+
+#[cfg(feature = "sync")]
+use radroots_storage::{outbox::OutboxRecord, projection::ProjectionId};
+#[cfg(feature = "sync")]
+use radroots_sync::{
+ Engine, PullReceipt, PullRequest, PushReceipt, PushRequest, SyncStatus,
+ ingest::{AdmissionPolicy, IngestBatchReceipt, IngestReceipt},
+ policy::Error,
+ projection::{Reducer, RefreshReceipt, RefreshRequest},
+ push::{DeliveryRunReceipt, DeliveryRunRequest},
+};
+#[cfg(feature = "sync")]
+use radroots_transport::source::ObservedEvent;
+
+/// Borrowed client operations over one explicitly composed sync engine.
+///
+/// This type owns no scheduling, retries, status strings, outbox state, or
+/// projection state. Every method delegates once to the canonical engine and
+/// returns its native receipt or error.
+#[cfg(feature = "sync")]
+#[derive(Clone, Copy)]
+pub struct Operations<'a> {
+ engine: &'a Engine,
+}
+
+#[cfg(feature = "sync")]
+impl<'a> Operations<'a> {
+ pub(crate) const fn new(engine: &'a Engine) -> Self {
+ Self { engine }
+ }
+
+ /// Runs one caller-bounded pull and canonical ingest sequence.
+ pub async fn pull(
+ &self,
+ request: PullRequest,
+ admission: &dyn AdmissionPolicy,
+ ) -> Result<PullReceipt, Error> {
+ self.engine.pull(request, admission).await
+ }
+
+ /// Verifies and atomically ingests one observed event.
+ pub async fn ingest(
+ &self,
+ observed: ObservedEvent,
+ admission: &dyn AdmissionPolicy,
+ ) -> Result<IngestReceipt, Error> {
+ self.engine.ingest(observed, admission).await
+ }
+
+ /// Ingests a bounded caller-owned batch while preserving partial outcomes.
+ pub async fn ingest_batch(
+ &self,
+ observed: Vec<ObservedEvent>,
+ admission: &dyn AdmissionPolicy,
+ ) -> IngestBatchReceipt {
+ self.engine.ingest_batch(observed, admission).await
+ }
+
+ /// Runs one bounded projection refresh through its owning reducer.
+ pub async fn refresh_projection(
+ &self,
+ request: RefreshRequest,
+ reducer: &dyn Reducer,
+ ) -> Result<RefreshReceipt, Error> {
+ self.engine.refresh_projection(request, reducer).await
+ }
+
+ /// Signs, verifies, and durably enqueues one outbound operation.
+ pub async fn sign_and_enqueue(&self, request: PushRequest) -> Result<PushReceipt, Error> {
+ self.engine.sign_and_enqueue(request).await
+ }
+
+ /// Runs one bounded delivery pass and retains every independent outcome.
+ pub async fn deliver_pending(
+ &self,
+ request: DeliveryRunRequest,
+ ) -> Result<DeliveryRunReceipt, Error> {
+ self.engine.deliver_pending(request).await
+ }
+
+ /// Returns the native passive sync status without starting recovery work.
+ pub async fn status(&self, projections: &[ProjectionId]) -> Result<SyncStatus, Error> {
+ self.engine.status(projections).await
+ }
+
+ /// Returns the native host scheduling decision for one durable plan.
+ pub fn retry_decision(
+ &self,
+ record: &OutboxRecord,
+ now_unix_ms: u64,
+ ) -> Result<radroots_protocol::runtime::v1::SyncRetryDecision, Error> {
+ self.engine.retry_decision(record, now_unix_ms)
+ }
+}
+
+#[cfg(feature = "sync")]
+impl std::fmt::Debug for Operations<'_> {
+ fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
+ formatter
+ .debug_struct("Operations")
+ .field("engine", &"<borrowed canonical engine>")
+ .finish()
+ }
+}
+
+#[cfg(all(test, feature = "sync", feature = "memory"))]
+mod tests {
+ use std::sync::{
+ Arc,
+ atomic::{AtomicU8, Ordering},
+ };
+
+ use radroots_event::{EventDraft, SignedEvent, contract::AuthorRole, wire::Nip01EventWire};
+ use radroots_identity::PublicKey;
+ use radroots_protocol::runtime::v1::SyncCapabilityState;
+ use radroots_signing::{Actor, actor::ActorSource, request::CancellationPolicy};
+ use radroots_storage::{
+ event::{SourceGeneration, StoredVisibleEvent},
+ journal::IdempotencyKey,
+ memory::MemoryStorage,
+ outbox::LeaseOwner,
+ projection::{ProjectionGeneration, ProjectionId},
+ };
+ use radroots_sync::{
+ Engine, PullRequest, PushRequest,
+ ingest::RegistryPolicy,
+ policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId, SyncStorage},
+ projection::{Reducer, ReducerError, RefreshRequest, RefreshState},
+ pull::PullTermination,
+ push::DeliveryRunRequest,
+ };
+ use radroots_transport::{
+ Error as TransportError, EventSource, FetchPage, FetchRequest, SourceStatus, Target,
+ TargetSet, TransportId,
+ capability::{Availability, Maturity, SourceCapabilities},
+ policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy},
+ source::{EventProvenance, NextPage, ObservedEvent},
+ };
+
+ use crate::{ClientBuilder, error::ErrorKind};
+
+ const PUBLIC_KEY: &str = "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df";
+
+ struct FixedClock;
+ struct SequenceIds(AtomicU8);
+ struct CancelledSource;
+
+ impl Clock for FixedClock {
+ fn now_unix_ms(&self) -> Result<u64, Error> {
+ Ok(1_700_000_000_000)
+ }
+ }
+
+ impl IdSource for SequenceIds {
+ fn next_id(&self, _operation: OperationKind) -> Result<SyncId, Error> {
+ let next = self.0.fetch_add(1, Ordering::Relaxed);
+ SyncId::new([next; 16])
+ }
+ }
+
+ impl EventSource for CancelledSource {
+ fn status(
+ &self,
+ ) -> radroots_transport::BoxFuture<'_, Result<SourceStatus, TransportError>> {
+ Box::pin(async {
+ Ok(SourceStatus::new(
+ TransportId::NOSTR,
+ true,
+ Maturity::Stable,
+ Availability::Available,
+ SourceCapabilities::FETCH,
+ "ready",
+ ))
+ })
+ }
+
+ fn fetch(
+ &self,
+ request: FetchRequest,
+ ) -> radroots_transport::BoxFuture<'_, Result<FetchPage, TransportError>> {
+ Box::pin(async move {
+ FetchPage::for_request(
+ &request,
+ Vec::new(),
+ Vec::new(),
+ NextPage::Cancelled { resume_from: None },
+ )
+ })
+ }
+ }
+
+ struct EmptyReducer {
+ id: ProjectionId,
+ generation: ProjectionGeneration,
+ }
+
+ impl Reducer for EmptyReducer {
+ fn projection_id(&self) -> &ProjectionId {
+ &self.id
+ }
+
+ fn generation(&self) -> ProjectionGeneration {
+ self.generation
+ }
+
+ fn reduce(
+ &self,
+ events: &[StoredVisibleEvent],
+ prior_projected_rows: u64,
+ ) -> Result<u64, ReducerError> {
+ assert!(events.is_empty());
+ Ok(prior_projected_rows)
+ }
+ }
+
+ fn target() -> Target {
+ Target::nostr_relay("wss://sync.example").expect("target")
+ }
+
+ fn engine(storage: Arc<MemoryStorage>) -> Engine {
+ let capability: Arc<dyn SyncStorage> = storage;
+ Engine::builder(
+ capability,
+ Arc::new(FixedClock),
+ Arc::new(SequenceIds(AtomicU8::new(1))),
+ DeadlinePolicy::new(1_000, 1_000, 1_000).expect("deadlines"),
+ )
+ .source(Arc::new(CancelledSource))
+ .build()
+ .expect("engine")
+ }
+
+ fn invalid_observation() -> ObservedEvent {
+ let mut wire = Nip01EventWire {
+ id: "0".repeat(64),
+ pubkey: PUBLIC_KEY.to_owned(),
+ created_at: 1_800_000_100,
+ kind: 0,
+ tags: vec![],
+ content: "invalid signature fixture".to_owned(),
+ sig: "42".repeat(64),
+ extra: Default::default(),
+ };
+ wire.id = wire.computed_event_id().expect("event id").to_hex();
+ let raw = format!(
+ "{{\"id\":\"{}\",\"pubkey\":\"{}\",\"created_at\":{},\"kind\":0,\"tags\":[],\"content\":\"invalid signature fixture\",\"sig\":\"{}\"}}",
+ wire.id, wire.pubkey, wire.created_at, wire.sig
+ );
+ let event = SignedEvent::from_wire_verified_id(wire, raw).expect("signed event");
+ let target = target();
+ let provenance = EventProvenance::new(
+ TransportId::NOSTR,
+ target.fingerprint().clone(),
+ 1_700_000_000_000,
+ )
+ .expect("provenance");
+ ObservedEvent::new(event, provenance)
+ }
+
+ fn push_request() -> PushRequest {
+ let actor = Actor::new(
+ PublicKey::from_hex(PUBLIC_KEY).expect("public key"),
+ ActorSource::ExplicitPublicKey,
+ [AuthorRole::Any],
+ )
+ .expect("actor");
+ let draft = EventDraft::new(
+ "radroots.social.geochat.v1",
+ 20_000,
+ 1_700_000_000,
+ Vec::new(),
+ "content",
+ PUBLIC_KEY,
+ )
+ .expect("draft");
+ PushRequest::new(
+ SyncId::new([8; 16]).expect("operation id"),
+ IdempotencyKey::parse("sdk-sync-wrapper").expect("idempotency key"),
+ actor,
+ draft,
+ TargetSet::new(vec![target()]).expect("targets"),
+ SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()),
+ CancellationPolicy::PreservePublishedRequest,
+ )
+ .expect("push request")
+ }
+
+ #[tokio::test]
+ async fn operations_delegate_native_pull_ingest_projection_push_delivery_and_status() {
+ let storage = Arc::new(MemoryStorage::new(
+ SourceGeneration::new([3; 32]).expect("generation"),
+ ));
+ let client = ClientBuilder::new()
+ .storage(storage.clone())
+ .sync_engine(engine(storage))
+ .build()
+ .expect("client");
+ let operations = client.sync().expect("open client").expect("sync");
+ let policy = RegistryPolicy::verified();
+
+ let pull = operations
+ .pull(
+ PullRequest::new(TargetSet::new(vec![target()]).expect("targets"), 10, 1)
+ .expect("pull request"),
+ &policy,
+ )
+ .await
+ .expect("pull");
+ assert_eq!(pull.termination(), PullTermination::Cancelled);
+
+ let ingest = operations
+ .ingest_batch(vec![invalid_observation(), invalid_observation()], &policy)
+ .await;
+ assert_eq!(ingest.accepted(), 0);
+ assert_eq!(ingest.rejected(), 2);
+ assert!(
+ ingest
+ .outcomes()
+ .iter()
+ .all(|outcome| matches!(outcome, Err(Error::VerificationFailed)))
+ );
+
+ let projection_id = ProjectionId::parse("sdk.sync.test").expect("projection id");
+ let generation = ProjectionGeneration::new([5; 32]).expect("generation");
+ let projection = operations
+ .refresh_projection(
+ RefreshRequest::new(projection_id.clone(), generation, 10, 1)
+ .expect("refresh request"),
+ &EmptyReducer {
+ id: projection_id.clone(),
+ generation,
+ },
+ )
+ .await
+ .expect("projection");
+ assert_eq!(projection.state(), RefreshState::Complete);
+
+ assert_eq!(
+ operations.sign_and_enqueue(push_request()).await,
+ Err(Error::MissingSigner)
+ );
+ let delivery = DeliveryRunRequest::new(
+ LeaseOwner::parse("sdk-sync-test").expect("owner"),
+ SyncId::new([9; 16]).expect("lease seed"),
+ 100,
+ 1,
+ )
+ .expect("delivery request");
+ assert_eq!(
+ operations.deliver_pending(delivery).await,
+ Err(Error::MissingSink)
+ );
+
+ let status = operations
+ .status(std::slice::from_ref(&projection_id))
+ .await
+ .expect("status");
+ assert_eq!(status.source().state(), SyncCapabilityState::Available);
+ assert_eq!(status.sink().state(), SyncCapabilityState::Unsupported);
+ assert_eq!(status.projections().len(), 1);
+
+ client.close().await.expect("close");
+ assert_eq!(
+ client.sync().expect_err("closed").kind(),
+ ErrorKind::ClientClosed
+ );
+ }
+}
diff --git a/crates/sdk/src/sync_runtime.rs b/crates/sdk/src/sync_runtime.rs
@@ -1,1947 +0,0 @@
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-use crate::adapters::radrootsd::{
- RadrootsdAuth, RadrootsdError, RadrootsdPublishAdapter, RadrootsdPublishConfig,
- RadrootsdPublishRequest,
-};
-#[cfg(feature = "runtime")]
-use crate::{
- NostrRelayUrlPolicy, RadrootsSdkError, SyncClient,
- runtime::{RadrootsClient, sdk_now_ms},
- transport::{ReticulumProfile, TransportProfile},
-};
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-use crate::{RadrootsdExecutionAuth, RadrootsdExecutionProfile};
-#[cfg(feature = "runtime")]
-use radroots_event::id::EventId;
-#[cfg(feature = "runtime")]
-use radroots_event_store::{RADROOTS_EVENT_STORE_QUERY_LIMIT_MAX, RadrootsEventStoreStatusSummary};
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-use radroots_outbox::RadrootsOutboxClaimedEvent;
-#[cfg(feature = "runtime")]
-use radroots_outbox::{
- RadrootsOutboxDeliveryTargetRecord, RadrootsOutboxDeliveryTargetStatus,
- RadrootsOutboxEventState, RadrootsOutboxReticulumEventRecord, RadrootsOutboxStatusSummary,
-};
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-use radroots_protocol::radrootsd::transport_publish::v5::{
- DeliveryPolicy as TransportPublishDeliveryPolicy, Job as TransportPublishJobView,
- JobStatus as TransportPublishJobStatus, OutcomeKind as TransportPublishOutcomeKind,
- Target as TransportPublishTarget, TargetFingerprint as TransportPublishTargetFingerprint,
- TargetOutcome as TransportPublishTargetOutcome, TargetPolicy as TransportPublishTargetPolicy,
-};
-#[cfg(feature = "runtime")]
-use radroots_trade::reducer::{RADROOTS_TRADE_REDUCER_CONTRACT_ID, RADROOTS_TRADE_REDUCER_VERSION};
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-use radroots_transport::target::TargetFingerprint;
-#[cfg(feature = "runtime")]
-use radroots_transport::{
- RadrootsTransportCapabilityAvailability, RadrootsTransportCapabilityMaturity,
- RadrootsTransportImplementationState, RadrootsTransportOutcomeKind, RadrootsTransportStatus,
- Target, TransportId,
-};
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-use radroots_transport::{RadrootsTransportSatisfactionClass, RadrootsTransportSatisfactionPolicy};
-#[cfg(all(feature = "runtime", feature = "transport-nostr-runtime"))]
-use radroots_transport_nostr::{
- RadrootsNostrClient, RadrootsNostrClientPublishAdapter, RadrootsNostrTransport,
-};
-#[cfg(feature = "runtime")]
-use radroots_transport_nostr::{
- RadrootsOutboxPublishPolicy, RadrootsOutboxPublishReceipt, RadrootsOutboxPublishTargetReceipt,
- RadrootsRelayOutcomeKind, publish_claimed_outbox_event_with_transport,
-};
-use radroots_transport_reticulum::RADROOTS_RETICULUM_UNAVAILABLE_MESSAGE;
-#[cfg(feature = "runtime")]
-use sha2::{Digest, Sha256};
-
-#[cfg(feature = "runtime")]
-pub const PUSH_OUTBOX_DEFAULT_LIMIT: usize = 20;
-#[cfg(feature = "runtime")]
-pub const PUSH_OUTBOX_MAX_LIMIT: usize = 100;
-#[cfg(feature = "runtime")]
-pub const PUSH_OUTBOX_DEFAULT_CLAIM_TTL_MS: i64 = 30_000;
-#[cfg(feature = "runtime")]
-pub const PUSH_OUTBOX_DEFAULT_NEXT_ATTEMPT_DELAY_MS: i64 = 60_000;
-#[cfg(feature = "runtime")]
-pub const SYNC_PROJECTION_REFRESH_DEFAULT_LIMIT: u32 = RADROOTS_EVENT_STORE_QUERY_LIMIT_MAX;
-#[cfg(feature = "runtime")]
-pub const SYNC_PROJECTION_REFRESH_MAX_LIMIT: u32 = RADROOTS_EVENT_STORE_QUERY_LIMIT_MAX;
-
-#[cfg(feature = "runtime")]
-const CLAIM_OWNER: &str = "radroots_sdk.sync.push_outbox";
-#[cfg(feature = "runtime")]
-const SDK_RELEASE_PRODUCT_PROJECTION_ID: &str = "radroots.release_product.trade_projection.v1";
-#[cfg(feature = "runtime")]
-const SDK_RELEASE_PRODUCT_PROJECTION_VERSION: u32 = 1;
-
-#[cfg(feature = "runtime")]
-#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, serde::Serialize)]
-#[non_exhaustive]
-pub struct SyncStatusRequest {}
-
-#[cfg(feature = "runtime")]
-impl SyncStatusRequest {
- pub fn new() -> Self {
- Self::default()
- }
-}
-
-#[cfg(feature = "runtime")]
-#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, serde::Serialize)]
-#[non_exhaustive]
-pub struct ReticulumTryNowRequest {}
-
-#[cfg(feature = "runtime")]
-impl ReticulumTryNowRequest {
- pub fn new() -> Self {
- Self::default()
- }
-}
-
-#[cfg(feature = "runtime")]
-#[derive(Clone, Debug, PartialEq, Eq, serde::Serialize)]
-pub struct SyncStatusReceipt {
- pub source: SyncStatusSource,
- pub observed_at_ms: i64,
- pub event_store: SyncEventStoreStatus,
- pub outbox: SyncOutboxStatus,
- pub transport_profile: SyncTransportProfileSummary,
-}
-
-#[cfg(feature = "runtime")]
-#[derive(Clone, Copy, Debug, PartialEq, Eq, serde::Serialize)]
-#[serde(rename_all = "snake_case")]
-#[non_exhaustive]
-pub enum SyncStatusSource {
- SdkCanonicalStores,
-}
-
-#[cfg(feature = "runtime")]
-#[derive(Clone, Debug, PartialEq, Eq, serde::Serialize)]
-pub struct SyncEventStoreStatus {
- pub total_events: i64,
- pub valid_stream_events: i64,
- pub transport_observations: i64,
- pub last_event_seq: Option<i64>,
- pub last_event_updated_at_ms: Option<i64>,
-}
-
-#[cfg(feature = "runtime")]
-impl From<RadrootsEventStoreStatusSummary> for SyncEventStoreStatus {
- fn from(summary: RadrootsEventStoreStatusSummary) -> Self {
- Self {
- total_events: summary.total_events,
- valid_stream_events: summary.valid_stream_events,
- transport_observations: summary.transport_observations,
- last_event_seq: summary.last_event_seq,
- last_event_updated_at_ms: summary.last_event_updated_at_ms,
- }
- }
-}
-
-#[cfg(feature = "runtime")]
-#[derive(Clone, Debug, PartialEq, Eq, serde::Serialize)]
-pub struct SyncOutboxStatus {
- pub total_events: i64,
- pub pending_events: i64,
- pub retryable_events: i64,
- pub terminal_events: i64,
- pub failed_terminal_events: i64,
- pub deferred_until_implemented_events: i64,
- pub ready_signed_events: i64,
- pub publishing_events: i64,
- pub last_attempt_at_ms: Option<i64>,
- pub last_error: Option<String>,
-}
-
-#[cfg(feature = "runtime")]
-impl From<RadrootsOutboxStatusSummary> for SyncOutboxStatus {
- fn from(summary: RadrootsOutboxStatusSummary) -> Self {
- Self {
- total_events: summary.total_events,
- pending_events: summary.pending_events,
- retryable_events: summary.retryable_events,
- terminal_events: summary.terminal_events,
- failed_terminal_events: summary.failed_terminal_events,
- deferred_until_implemented_events: summary.deferred_until_implemented_events,
- ready_signed_events: summary.ready_signed_events,
- publishing_events: summary.publishing_events,
- last_attempt_at_ms: summary.last_attempt_at_ms,
- last_error: summary.last_error,
- }
- }
-}
-
-#[cfg(feature = "runtime")]
-#[derive(Clone, Debug, PartialEq, Eq, serde::Serialize)]
-pub struct SyncTransportProfileSummary {
- pub transport_profile_id: String,
- pub configured_transport_target_count: usize,
- pub configured_transport_targets: Vec<SyncTransportTargetSummary>,
- pub transport_statuses: Vec<SyncTransportStatusSummary>,
-}
-
-#[cfg(feature = "runtime")]
-impl SyncTransportProfileSummary {
- fn from_transport_profile(profile: &TransportProfile) -> Result<Self, RadrootsSdkError> {
- let configured_transport_targets = profile
- .configured_transport_targets()?
- .iter()
- .map(SyncTransportTargetSummary::from_transport_target)
- .collect::<Vec<_>>();
- Ok(Self {
- transport_profile_id: profile.transport_profile_id().to_owned(),
- configured_transport_target_count: configured_transport_targets.len(),
- configured_transport_targets,
- transport_statuses: profile
- .transport_statuses()
- .into_iter()
- .map(SyncTransportStatusSummary::from_transport_status)
- .collect(),
- })
- }
-}
-
-#[cfg(feature = "runtime")]
-#[derive(Clone, Debug, PartialEq, Eq, serde::Serialize)]
-pub struct SyncTransportTargetSummary {
- pub transport_kind: String,
- pub endpoint_uri: String,
- pub target_scope: Option<String>,
- pub target_label: Option<String>,
- pub endpoint_fingerprint: String,
-}
-
-#[cfg(feature = "runtime")]
-impl SyncTransportTargetSummary {
- fn from_transport_target(target: &Target) -> Self {
- Self {
- transport_kind: target.kind().canonical_label(),
- endpoint_uri: target.uri().as_str().to_owned(),
- target_scope: target.scope().map(|scope| scope.as_str().to_owned()),
- target_label: target.label().map(|label| label.as_str().to_owned()),
- endpoint_fingerprint: target.fingerprint().as_str().to_owned(),
- }
- }
-}
-
-#[cfg(feature = "runtime")]
-#[derive(Clone, Debug, PartialEq, Eq, serde::Serialize)]
-pub struct SyncTransportStatusSummary {
- pub transport: String,
- pub profile_id: Option<String>,
- pub endpoint_uri: Option<String>,
- pub configured: bool,
- pub implementation: String,
- pub maturity: String,
- pub availability: String,
- pub usable_for_delivery: bool,
- pub capabilities: SyncTransportOperationCapabilitiesSummary,
- pub message: String,
-}
-
-#[cfg(feature = "runtime")]
-#[derive(Clone, Debug, PartialEq, Eq, serde::Serialize)]
-pub struct SyncTransportOperationCapabilitiesSummary {
- pub deliver: bool,
- pub fetch: bool,
-}
-
-#[cfg(feature = "runtime")]
-impl SyncTransportStatusSummary {
- fn from_transport_status(status: RadrootsTransportStatus) -> Self {
- Self {
- transport: status.kind.canonical_label(),
- profile_id: status.profile_id,
- endpoint_uri: status.endpoint_uri,
- configured: status.configured,
- implementation: transport_implementation_label(status.implementation).to_owned(),
- maturity: transport_maturity_label(status.maturity).to_owned(),
- availability: transport_availability_label(status.availability).to_owned(),
- usable_for_delivery: status.usable_for_delivery,
- capabilities: SyncTransportOperationCapabilitiesSummary {
- deliver: status.capabilities.deliver,
- fetch: status.capabilities.fetch,
- },
- message: status.message,
- }
- }
-}
-
-#[cfg(feature = "runtime")]
-fn transport_implementation_label(state: RadrootsTransportImplementationState) -> &'static str {
- match state {
- RadrootsTransportImplementationState::Real => "real",
- RadrootsTransportImplementationState::Mock => "mock",
- }
-}
-
-#[cfg(feature = "runtime")]
-fn transport_maturity_label(maturity: RadrootsTransportCapabilityMaturity) -> &'static str {
- match maturity {
- RadrootsTransportCapabilityMaturity::Experimental => "experimental",
- RadrootsTransportCapabilityMaturity::Preview => "preview",
- RadrootsTransportCapabilityMaturity::Stable => "stable",
- }
-}
-
-#[cfg(feature = "runtime")]
-fn transport_availability_label(
- availability: RadrootsTransportCapabilityAvailability,
-) -> &'static str {
- match availability {
- RadrootsTransportCapabilityAvailability::Available => "available",
- RadrootsTransportCapabilityAvailability::Degraded => "degraded",
- RadrootsTransportCapabilityAvailability::Unavailable => "unavailable",
- }
-}
-
-#[cfg(feature = "runtime")]
-#[derive(Clone, Copy, Debug, PartialEq, Eq, serde::Serialize)]
-#[non_exhaustive]
-pub struct SyncProjectionRefreshRequest {
- pub limit: u32,
-}
-
-#[cfg(feature = "runtime")]
-impl Default for SyncProjectionRefreshRequest {
- fn default() -> Self {
- Self {
- limit: SYNC_PROJECTION_REFRESH_DEFAULT_LIMIT,
- }
- }
-}
-
-#[cfg(feature = "runtime")]
-impl SyncProjectionRefreshRequest {
- pub fn new() -> Self {
- Self::default()
- }
-
- pub fn with_limit(mut self, limit: u32) -> Self {
- self.limit = limit;
- self
- }
-}
-
-#[cfg(feature = "runtime")]
-#[derive(Clone, Debug, Default, PartialEq, Eq, serde::Serialize)]
-pub struct SyncProjectionRefreshReceipt {
- pub projection_id: &'static str,
- pub projection_version: u32,
- pub refreshed_at_ms: i64,
- pub scanned_events: usize,
- pub listing_upserts: usize,
- pub trade_upserts: usize,
- pub agreement_attestations: usize,
- pub transport_observations: i64,
- pub last_event_seq: Option<i64>,
-}
-
-#[cfg(feature = "runtime")]
-impl SyncProjectionRefreshReceipt {
- fn from_event_store_snapshot(
- summary: RadrootsEventStoreStatusSummary,
- scanned_events: usize,
- trade_upserts: usize,
- refreshed_at_ms: i64,
- ) -> Self {
- Self {
- projection_id: SDK_RELEASE_PRODUCT_PROJECTION_ID,
- projection_version: SDK_RELEASE_PRODUCT_PROJECTION_VERSION,
- refreshed_at_ms,
- scanned_events,
- listing_upserts: 0,
- trade_upserts,
- agreement_attestations: 0,
- transport_observations: summary.transport_observations,
- last_event_seq: summary.last_event_seq,
- }
- }
-}
-
-#[cfg(feature = "runtime")]
-#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, serde::Serialize)]
-#[serde(rename_all = "snake_case")]
-#[non_exhaustive]
-pub enum SdkRelayAuthPolicy {
- #[default]
- DetectOnly,
-}
-
-#[cfg(feature = "runtime")]
-#[derive(Clone, Debug, PartialEq, Eq, serde::Serialize)]
-#[non_exhaustive]
-pub struct PushOutboxRequest {
- pub limit: usize,
- pub outbox_event_id: Option<i64>,
- pub republish_accepted_targets: bool,
- pub nostr_relay_url_policy: NostrRelayUrlPolicy,
- pub auth_policy: SdkRelayAuthPolicy,
- pub claim_ttl_ms: i64,
- pub next_attempt_delay_ms: i64,
-}
-
-#[cfg(feature = "runtime")]
-impl Default for PushOutboxRequest {
- fn default() -> Self {
- Self {
- limit: PUSH_OUTBOX_DEFAULT_LIMIT,
- outbox_event_id: None,
- republish_accepted_targets: false,
- nostr_relay_url_policy: NostrRelayUrlPolicy::Public,
- auth_policy: SdkRelayAuthPolicy::DetectOnly,
- claim_ttl_ms: PUSH_OUTBOX_DEFAULT_CLAIM_TTL_MS,
- next_attempt_delay_ms: PUSH_OUTBOX_DEFAULT_NEXT_ATTEMPT_DELAY_MS,
- }
- }
-}
-
-#[cfg(feature = "runtime")]
-impl PushOutboxRequest {
- pub fn new() -> Self {
- Self::default()
- }
-
- pub fn with_limit(mut self, limit: usize) -> Self {
- self.limit = limit;
- self
- }
-
- pub fn with_outbox_event_id(mut self, outbox_event_id: i64) -> Self {
- self.outbox_event_id = Some(outbox_event_id);
- self.limit = 1;
- self
- }
-
- pub fn republish_accepted_targets(mut self, enabled: bool) -> Self {
- self.republish_accepted_targets = enabled;
- self
- }
-
- pub fn with_nostr_relay_url_policy(mut self, policy: NostrRelayUrlPolicy) -> Self {
- self.nostr_relay_url_policy = policy;
- self
- }
-
- pub fn with_auth_policy(mut self, policy: SdkRelayAuthPolicy) -> Self {
- self.auth_policy = policy;
- self
- }
-
- pub fn with_claim_ttl_ms(mut self, claim_ttl_ms: i64) -> Self {
- self.claim_ttl_ms = claim_ttl_ms;
- self
- }
-
- pub fn with_next_attempt_delay_ms(mut self, next_attempt_delay_ms: i64) -> Self {
- self.next_attempt_delay_ms = next_attempt_delay_ms;
- self
- }
-
- fn validate(&self) -> Result<(), RadrootsSdkError> {
- if self.limit == 0 {
- return Err(RadrootsSdkError::InvalidRequest {
- message: format!("push_outbox limit must be between 1 and {PUSH_OUTBOX_MAX_LIMIT}"),
- });
- }
- if self.limit > PUSH_OUTBOX_MAX_LIMIT {
- return Err(RadrootsSdkError::InvalidRequest {
- message: format!("push_outbox limit must be between 1 and {PUSH_OUTBOX_MAX_LIMIT}"),
- });
- }
- if let Some(outbox_event_id) = self.outbox_event_id
- && outbox_event_id <= 0
- {
- return Err(RadrootsSdkError::InvalidRequest {
- message: "push_outbox outbox event id must be positive".to_owned(),
- });
- }
- if self.claim_ttl_ms <= 0 {
- return Err(RadrootsSdkError::InvalidRequest {
- message: "push_outbox claim TTL must be positive".to_owned(),
- });
- }
- if self.next_attempt_delay_ms <= 0 {
- return Err(RadrootsSdkError::InvalidRequest {
- message: "push_outbox next attempt delay must be positive".to_owned(),
- });
- }
- Ok(())
- }
-}
-
-#[cfg(feature = "runtime")]
-#[derive(Clone, Debug, Default, PartialEq, Eq, serde::Serialize)]
-pub struct PushOutboxReceipt {
- pub attempted_events: usize,
- pub published_events: usize,
- pub retryable_events: usize,
- pub terminal_events: usize,
- pub events: Vec<PushOutboxEventReceipt>,
-}
-
-#[cfg(feature = "runtime")]
-impl PushOutboxReceipt {
- fn push_attempted_event(&mut self, event: PushOutboxEventReceipt) {
- self.attempted_events += 1;
- self.push_reported_event(event);
- }
-
- fn push_reported_event(&mut self, event: PushOutboxEventReceipt) {
- match event.final_state {
- PushOutboxEventState::Published => self.published_events += 1,
- PushOutboxEventState::PublishRetryable => self.retryable_events += 1,
- PushOutboxEventState::FailedTerminal => self.terminal_events += 1,
- _ => {}
- }
- self.events.push(event);
- }
-}
-
-#[cfg(feature = "runtime")]
-#[derive(Clone, Debug, PartialEq, Eq, serde::Serialize)]
-pub struct PushOutboxEventReceipt {
- pub event_id: EventId,
- pub outbox_event_id: i64,
- pub final_state: PushOutboxEventState,
- pub attempted_count: usize,
- pub accepted_count: usize,
- pub retryable_count: usize,
- pub terminal_count: usize,
- pub quorum: usize,
- pub quorum_met: bool,
- pub targets: Vec<PushOutboxTargetReceipt>,
-}
-
-#[cfg(feature = "runtime")]
-#[derive(Clone, Debug, PartialEq, Eq, serde::Serialize)]
-pub struct PushOutboxTargetReceipt {
- pub transport_kind: String,
- pub endpoint_uri: String,
- pub target_scope: Option<String>,
- pub target_label: Option<String>,
- pub outcome_kind: PushOutboxTargetOutcomeKind,
- pub transport_outcome_kind: Option<PushOutboxTransportOutcomeKind>,
- pub attempted: bool,
- pub message: Option<String>,
-}
-
-#[cfg(feature = "runtime")]
-#[derive(Clone, Copy, Debug, PartialEq, Eq, serde::Serialize)]
-#[serde(rename_all = "snake_case")]
-#[non_exhaustive]
-pub enum PushOutboxEventState {
- DraftQueued,
- Signing,
- Signed,
- Publishing,
- Published,
- SignRetryable,
- PublishRetryable,
- DeferredUntilImplemented,
- FailedTerminal,
- Cancelled,
-}
-
-#[cfg(feature = "runtime")]
-impl From<RadrootsOutboxEventState> for PushOutboxEventState {
- fn from(state: RadrootsOutboxEventState) -> Self {
- match state {
- RadrootsOutboxEventState::DraftQueued => Self::DraftQueued,
- RadrootsOutboxEventState::Signing => Self::Signing,
- RadrootsOutboxEventState::Signed => Self::Signed,
- RadrootsOutboxEventState::Publishing => Self::Publishing,
- RadrootsOutboxEventState::Published => Self::Published,
- RadrootsOutboxEventState::SignRetryable => Self::SignRetryable,
- RadrootsOutboxEventState::PublishRetryable => Self::PublishRetryable,
- RadrootsOutboxEventState::FailedTerminal => Self::FailedTerminal,
- RadrootsOutboxEventState::Cancelled => Self::Cancelled,
- }
- }
-}
-
-#[cfg(feature = "runtime")]
-#[derive(Clone, Copy, Debug, PartialEq, Eq, serde::Serialize)]
-#[serde(rename_all = "snake_case")]
-#[non_exhaustive]
-pub enum PushOutboxTargetOutcomeKind {
- Accepted,
- DuplicateAccepted,
- Blocked,
- RateLimited,
- Invalid,
- PowRequired,
- Restricted,
- AuthRequired,
- Muted,
- Unsupported,
- PaymentRequired,
- Error,
- Timeout,
- ConnectionFailed,
- TargetUriRejected,
- SkippedAlreadyAccepted,
- DeferredUntilImplemented,
- Unknown,
-}
-
-#[cfg(feature = "runtime")]
-impl PushOutboxTargetOutcomeKind {
- pub fn as_str(self) -> &'static str {
- match self {
- Self::Accepted => "accepted",
- Self::DuplicateAccepted => "duplicate_accepted",
- Self::Blocked => "blocked",
- Self::RateLimited => "rate_limited",
- Self::Invalid => "invalid",
- Self::PowRequired => "pow_required",
- Self::Restricted => "restricted",
- Self::AuthRequired => "auth_required",
- Self::Muted => "muted",
- Self::Unsupported => "unsupported",
- Self::PaymentRequired => "payment_required",
- Self::Error => "error",
- Self::Timeout => "timeout",
- Self::ConnectionFailed => "connection_failed",
- Self::TargetUriRejected => "target_uri_rejected",
- Self::SkippedAlreadyAccepted => "skipped_already_accepted",
- Self::DeferredUntilImplemented => "deferred_until_implemented",
- Self::Unknown => "unknown",
- }
- }
-}
-
-#[cfg(feature = "runtime")]
-#[derive(Clone, Copy, Debug, PartialEq, Eq, serde::Serialize)]
-#[serde(rename_all = "snake_case")]
-#[non_exhaustive]
-pub enum PushOutboxTransportOutcomeKind {
- Accepted,
- DuplicateAccepted,
- Delivered,
- Forwarded,
- StoredByGateway,
- Seen,
- DeferredUntilImplemented,
- Rejected,
- RouteUnavailable,
- PayloadTooLarge,
- PolicyDenied,
- Timeout,
- ConnectionFailed,
- TransportUnavailable,
-}
-
-#[cfg(feature = "runtime")]
-impl PushOutboxTransportOutcomeKind {
- pub fn as_str(self) -> &'static str {
- match self {
- Self::Accepted => "accepted",
- Self::DuplicateAccepted => "duplicate_accepted",
- Self::Delivered => "delivered",
- Self::Forwarded => "forwarded",
- Self::StoredByGateway => "stored_by_gateway",
- Self::Seen => "seen",
- Self::DeferredUntilImplemented => "deferred_until_implemented",
- Self::Rejected => "rejected",
- Self::RouteUnavailable => "route_unavailable",
- Self::PayloadTooLarge => "payload_too_large",
- Self::PolicyDenied => "policy_denied",
- Self::Timeout => "timeout",
- Self::ConnectionFailed => "connection_failed",
- Self::TransportUnavailable => "transport_unavailable",
- }
- }
-}
-
-#[cfg(feature = "runtime")]
-impl From<RadrootsTransportOutcomeKind> for PushOutboxTransportOutcomeKind {
- fn from(kind: RadrootsTransportOutcomeKind) -> Self {
- match kind {
- RadrootsTransportOutcomeKind::Accepted => Self::Accepted,
- RadrootsTransportOutcomeKind::DuplicateAccepted => Self::DuplicateAccepted,
- RadrootsTransportOutcomeKind::Delivered => Self::Delivered,
- RadrootsTransportOutcomeKind::Forwarded => Self::Forwarded,
- RadrootsTransportOutcomeKind::StoredByGateway => Self::StoredByGateway,
- RadrootsTransportOutcomeKind::Seen => Self::Seen,
- RadrootsTransportOutcomeKind::DeferredUntilImplemented => {
- Self::DeferredUntilImplemented
- }
- RadrootsTransportOutcomeKind::Rejected => Self::Rejected,
- RadrootsTransportOutcomeKind::RouteUnavailable => Self::RouteUnavailable,
- RadrootsTransportOutcomeKind::PayloadTooLarge => Self::PayloadTooLarge,
- RadrootsTransportOutcomeKind::PolicyDenied => Self::PolicyDenied,
- RadrootsTransportOutcomeKind::Timeout => Self::Timeout,
- RadrootsTransportOutcomeKind::ConnectionFailed => Self::ConnectionFailed,
- RadrootsTransportOutcomeKind::TransportUnavailable => Self::TransportUnavailable,
- }
- }
-}
-
-#[cfg(feature = "runtime")]
-impl From<RadrootsRelayOutcomeKind> for PushOutboxTargetOutcomeKind {
- fn from(kind: RadrootsRelayOutcomeKind) -> Self {
- match kind {
- RadrootsRelayOutcomeKind::Accepted => Self::Accepted,
- RadrootsRelayOutcomeKind::DuplicateAccepted => Self::DuplicateAccepted,
- RadrootsRelayOutcomeKind::Blocked => Self::Blocked,
- RadrootsRelayOutcomeKind::RateLimited => Self::RateLimited,
- RadrootsRelayOutcomeKind::Invalid => Self::Invalid,
- RadrootsRelayOutcomeKind::PowRequired => Self::PowRequired,
- RadrootsRelayOutcomeKind::Restricted => Self::Restricted,
- RadrootsRelayOutcomeKind::AuthRequired => Self::AuthRequired,
- RadrootsRelayOutcomeKind::Muted => Self::Muted,
- RadrootsRelayOutcomeKind::Unsupported => Self::Unsupported,
- RadrootsRelayOutcomeKind::PaymentRequired => Self::PaymentRequired,
- RadrootsRelayOutcomeKind::Error => Self::Error,
- RadrootsRelayOutcomeKind::Timeout => Self::Timeout,
- RadrootsRelayOutcomeKind::ConnectionFailed => Self::ConnectionFailed,
- RadrootsRelayOutcomeKind::RelayUrlRejected => Self::TargetUriRejected,
- RadrootsRelayOutcomeKind::SkippedAlreadyAccepted => Self::SkippedAlreadyAccepted,
- RadrootsRelayOutcomeKind::Unknown => Self::Unknown,
- }
- }
-}
-
-#[cfg(feature = "runtime")]
-impl<'sdk> SyncClient<'sdk> {
- pub async fn refresh_projections(
- &self,
- request: SyncProjectionRefreshRequest,
- ) -> Result<SyncProjectionRefreshReceipt, RadrootsSdkError> {
- refresh_product_projections_for_sdk(self.sdk, request).await
- }
-
- pub async fn status(
- &self,
- _request: SyncStatusRequest,
- ) -> Result<SyncStatusReceipt, RadrootsSdkError> {
- let observed_at_ms = sdk_now_ms(self.sdk)?;
- let event_store = self.sdk._event_store.status_summary().await?;
- let outbox = self.sdk._outbox.status_summary(observed_at_ms).await?;
- Ok(SyncStatusReceipt {
- source: SyncStatusSource::SdkCanonicalStores,
- observed_at_ms,
- event_store: event_store.into(),
- outbox: outbox.into(),
- transport_profile: SyncTransportProfileSummary::from_transport_profile(
- self.sdk.transport_profile(),
- )?,
- })
- }
-
- pub async fn push_outbox(
- &self,
- request: PushOutboxRequest,
- ) -> Result<PushOutboxReceipt, RadrootsSdkError> {
- #[cfg(feature = "radrootsd-execution")]
- if let Some(profile) = self.sdk.radrootsd_execution_profile() {
- let adapter =
- RadrootsdPublishAdapter::new(radrootsd_publish_config_from_profile(profile));
- return self
- .push_outbox_with_radrootsd_adapter(&adapter, request)
- .await;
- }
-
- #[cfg(not(feature = "radrootsd-execution"))]
- if self.sdk.radrootsd_execution_profile().is_some() {
- return Err(RadrootsSdkError::ProductSyncUnsupported {
- operation: "sync.push_outbox",
- required_feature: "radrootsd-execution",
- });
- }
-
- match self.sdk.transport_profile() {
- TransportProfile::Nostr { .. } | TransportProfile::MultiTarget { .. } => {
- #[cfg(feature = "transport-nostr-runtime")]
- {
- let adapter = RadrootsNostrClientPublishAdapter::new(
- RadrootsNostrClient::new_signerless(),
- );
- let transport = RadrootsNostrTransport::new(adapter);
- self.push_outbox_with_transport(&transport, request).await
- }
-
- #[cfg(not(feature = "transport-nostr-runtime"))]
- {
- let _ = request;
- Err(RadrootsSdkError::ProductSyncUnsupported {
- operation: "sync.push_outbox",
- required_feature: "transport-nostr-runtime",
- })
- }
- }
- TransportProfile::LocalOnly => {
- if self.push_outbox_has_no_ready_signed_work(&request).await? {
- return Ok(PushOutboxReceipt::default());
- }
- Err(RadrootsSdkError::ProductSyncUnsupported {
- operation: "sync.push_outbox",
- required_feature: "delivery-capable transport profile",
- })
- }
- TransportProfile::Reticulum { .. } => self.reticulum_push_receipt(request).await,
- }
- }
-
- pub async fn try_reticulum_now(
- &self,
- _request: ReticulumTryNowRequest,
- ) -> Result<(), RadrootsSdkError> {
- let profile = active_reticulum_profile(self.sdk.transport_profile()).ok_or_else(|| {
- RadrootsSdkError::InvalidRequest {
- message: "sync.try_reticulum_now requires a Reticulum transport profile".to_owned(),
- }
- })?;
- Err(RadrootsSdkError::ReticulumTransportUnavailable {
- operation: "sync.try_reticulum_now".to_owned(),
- endpoint_uri: profile.endpoint_uri().to_owned(),
- behavior: profile.behavior(),
- })
- }
-
- async fn push_outbox_has_no_ready_signed_work(
- &self,
- request: &PushOutboxRequest,
- ) -> Result<bool, RadrootsSdkError> {
- request.validate()?;
- let now_ms = sdk_now_ms(self.sdk)?;
- let summary = self.sdk._outbox.status_summary(now_ms).await?;
- Ok(summary.ready_signed_events == 0)
- }
-
- pub async fn push_outbox_with_transport<T>(
- &self,
- transport: &T,
- request: PushOutboxRequest,
- ) -> Result<PushOutboxReceipt, RadrootsSdkError>
- where
- T: radroots_transport::EventSink + ?Sized,
- {
- request.validate()?;
- let recovery_now_ms = sdk_now_ms(self.sdk)?;
- recover_expired_outbox_claims_for_push(self.sdk, recovery_now_ms).await?;
- let mut receipt = PushOutboxReceipt::default();
- for index in 0..request.limit {
- let claim_now_ms = if index == 0 {
- recovery_now_ms
- } else {
- sdk_now_ms(self.sdk)?
- };
- let claim_token = push_outbox_claim_token();
- let Some(claimed) = claim_ready_signed_event_for_push(
- self.sdk,
- &request,
- claim_token.as_str(),
- claim_now_ms,
- )
- .await?
- else {
- break;
- };
- let publish_now_ms = claim_now_ms;
- let policy = RadrootsOutboxPublishPolicy::new(
- publish_now_ms.saturating_add(request.next_attempt_delay_ms),
- )
- .republish_accepted_relays(request.republish_accepted_targets)
- .relay_url_policy(request.nostr_relay_url_policy.nostr_transport_policy());
- let publish = publish_claimed_outbox_event_with_transport(
- &self.sdk._outbox,
- &self.sdk._event_store,
- transport,
- &claimed,
- policy,
- publish_now_ms,
- )
- .await?;
- let final_state = push_event_final_state(&publish);
- receipt.push_attempted_event(push_event_receipt(
- claimed.outbox_event_id,
- final_state,
- publish,
- )?);
- }
- Ok(receipt)
- }
-
- #[cfg(feature = "radrootsd-execution")]
- pub async fn push_outbox_with_radrootsd_adapter(
- &self,
- adapter: &RadrootsdPublishAdapter,
- request: PushOutboxRequest,
- ) -> Result<PushOutboxReceipt, RadrootsSdkError> {
- request.validate()?;
- let recovery_now_ms = sdk_now_ms(self.sdk)?;
- recover_expired_outbox_claims_for_push(self.sdk, recovery_now_ms).await?;
- let mut receipt = PushOutboxReceipt::default();
- for index in 0..request.limit {
- let claim_now_ms = if index == 0 {
- recovery_now_ms
- } else {
- sdk_now_ms(self.sdk)?
- };
- let claim_token = push_outbox_claim_token();
- let Some(claimed) = claim_ready_signed_event_for_push(
- self.sdk,
- &request,
- claim_token.as_str(),
- claim_now_ms,
- )
- .await?
- else {
- break;
- };
- let publish_now_ms = claim_now_ms;
- let publish = push_radrootsd_claimed_outbox_event(
- self,
- adapter,
- &claimed,
- request.next_attempt_delay_ms,
- publish_now_ms,
- )
- .await?;
- receipt.push_attempted_event(publish);
- }
- Ok(receipt)
- }
-
- async fn reticulum_push_receipt(
- &self,
- request: PushOutboxRequest,
- ) -> Result<PushOutboxReceipt, RadrootsSdkError> {
- request.validate()?;
- let records = self
- .sdk
- ._outbox
- .reticulum_events(request.outbox_event_id, request.limit)
- .await?;
- let mut receipt = PushOutboxReceipt::default();
- for record in records {
- receipt.push_reported_event(reticulum_event_receipt(record)?);
- }
- Ok(receipt)
- }
-}
-
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-fn radrootsd_publish_config_from_profile(
- profile: &RadrootsdExecutionProfile,
-) -> RadrootsdPublishConfig {
- let config = RadrootsdPublishConfig::new(profile.endpoint_url().to_owned());
- match profile.auth() {
- RadrootsdExecutionAuth::None => config,
- RadrootsdExecutionAuth::BearerToken(token) => {
- config.with_auth(RadrootsdAuth::BearerToken(token.to_owned()))
- }
- }
-}
-
-#[cfg(feature = "runtime")]
-async fn recover_expired_outbox_claims_for_push(
- sdk: &RadrootsClient,
- now_ms: i64,
-) -> Result<(), RadrootsSdkError> {
- sdk._outbox.recover_expired_claims(now_ms).await?;
- Ok(())
-}
-
-#[cfg(feature = "runtime")]
-async fn claim_ready_signed_event_for_push(
- sdk: &RadrootsClient,
- request: &PushOutboxRequest,
- claim_token: &str,
- claim_now_ms: i64,
-) -> Result<Option<radroots_outbox::RadrootsOutboxClaimedEvent>, RadrootsSdkError> {
- let claim_expires_at_ms = claim_now_ms.saturating_add(request.claim_ttl_ms);
- match request.outbox_event_id {
- Some(outbox_event_id) => Ok(sdk
- ._outbox
- .claim_ready_signed_event(
- outbox_event_id,
- CLAIM_OWNER,
- claim_token,
- claim_expires_at_ms,
- claim_now_ms,
- )
- .await?),
- None => Ok(sdk
- ._outbox
- .claim_next_ready_signed_event(
- CLAIM_OWNER,
- claim_token,
- claim_expires_at_ms,
- claim_now_ms,
- )
- .await?),
- }
-}
-
-#[cfg(feature = "runtime")]
-fn reticulum_event_receipt(
- record: RadrootsOutboxReticulumEventRecord,
-) -> Result<PushOutboxEventReceipt, RadrootsSdkError> {
- let event_id =
- push_receipt_event_id(record.event.event_id.as_str(), "Reticulum outbox event id")?;
- let final_state = reticulum_event_final_state(record.event.state, &record.targets);
- let quorum = record.targets.len();
- Ok(PushOutboxEventReceipt {
- event_id,
- outbox_event_id: record.event.outbox_event_id,
- final_state,
- attempted_count: 0,
- accepted_count: 0,
- retryable_count: 0,
- terminal_count: 0,
- quorum,
- quorum_met: false,
- targets: record
- .targets
- .into_iter()
- .map(reticulum_target_receipt)
- .collect(),
- })
-}
-
-#[cfg(feature = "runtime")]
-fn reticulum_event_final_state(
- _event_state: RadrootsOutboxEventState,
- _targets: &[RadrootsOutboxDeliveryTargetRecord],
-) -> PushOutboxEventState {
- PushOutboxEventState::DeferredUntilImplemented
-}
-
-#[cfg(feature = "runtime")]
-fn reticulum_target_receipt(target: RadrootsOutboxDeliveryTargetRecord) -> PushOutboxTargetReceipt {
- PushOutboxTargetReceipt {
- transport_kind: target.transport_kind.canonical_label(),
- endpoint_uri: target.endpoint_uri.as_str().to_owned(),
- target_scope: target
- .target_scope
- .as_ref()
- .map(|scope| scope.as_str().to_owned()),
- target_label: target
- .target_label
- .as_ref()
- .map(|label| label.as_str().to_owned()),
- outcome_kind: reticulum_target_outcome_kind(target.status),
- transport_outcome_kind: target.last_outcome_kind.map(Into::into),
- attempted: false,
- message: Some(
- target
- .last_error
- .unwrap_or_else(|| RADROOTS_RETICULUM_UNAVAILABLE_MESSAGE.to_owned()),
- ),
- }
-}
-
-#[cfg(feature = "runtime")]
-fn active_reticulum_profile(profile: &TransportProfile) -> Option<&ReticulumProfile> {
- match profile {
- TransportProfile::Reticulum { profile } => Some(profile),
- TransportProfile::MultiTarget { profile } => Some(profile.reticulum()),
- TransportProfile::LocalOnly | TransportProfile::Nostr { .. } => None,
- }
-}
-
-#[cfg(feature = "runtime")]
-fn reticulum_target_outcome_kind(
- status: RadrootsOutboxDeliveryTargetStatus,
-) -> PushOutboxTargetOutcomeKind {
- match status {
- RadrootsOutboxDeliveryTargetStatus::DeferredUntilImplemented => {
- PushOutboxTargetOutcomeKind::DeferredUntilImplemented
- }
- RadrootsOutboxDeliveryTargetStatus::Pending
- | RadrootsOutboxDeliveryTargetStatus::FailedRetryable => {
- PushOutboxTargetOutcomeKind::DeferredUntilImplemented
- }
- RadrootsOutboxDeliveryTargetStatus::Accepted
- | RadrootsOutboxDeliveryTargetStatus::Delivered
- | RadrootsOutboxDeliveryTargetStatus::Forwarded
- | RadrootsOutboxDeliveryTargetStatus::StoredByGateway
- | RadrootsOutboxDeliveryTargetStatus::Seen
- | RadrootsOutboxDeliveryTargetStatus::SkippedPolicyDenied
- | RadrootsOutboxDeliveryTargetStatus::FailedTerminal => {
- PushOutboxTargetOutcomeKind::Unknown
- }
- }
-}
-
-#[cfg(feature = "runtime")]
-pub(crate) async fn refresh_product_projections_for_sdk(
- sdk: &RadrootsClient,
- request: SyncProjectionRefreshRequest,
-) -> Result<SyncProjectionRefreshReceipt, RadrootsSdkError> {
- if request.limit == 0 || request.limit > SYNC_PROJECTION_REFRESH_MAX_LIMIT {
- return Err(RadrootsSdkError::InvalidRequest {
- message: format!(
- "projection refresh limit must be between 1 and {SYNC_PROJECTION_REFRESH_MAX_LIMIT}"
- ),
- });
- }
- let refreshed_at_ms = sdk_now_ms(sdk)?;
- let summary = sdk._event_store.status_summary().await?;
- let scanned_events = usize::try_from(summary.valid_stream_events.max(0))
- .unwrap_or(usize::MAX)
- .min(request.limit as usize);
- let trade_upserts = trade_mutation_count(sdk).await?;
- let source_digest = projection_source_digest(&summary, trade_upserts);
- sqlx::query(
- "INSERT INTO sdk_runtime_trade_projection_checkpoint(projection_name, reducer_contract_id, reducer_version, last_ingest_seq, source_digest, projection_digest, completeness_state, rebuilt_at_ms, updated_at_ms) VALUES (?, ?, ?, ?, ?, ?, 'current', ?, ?) ON CONFLICT(projection_name) DO UPDATE SET reducer_contract_id = excluded.reducer_contract_id, reducer_version = excluded.reducer_version, last_ingest_seq = excluded.last_ingest_seq, source_digest = excluded.source_digest, projection_digest = excluded.projection_digest, completeness_state = excluded.completeness_state, rebuilt_at_ms = excluded.rebuilt_at_ms, updated_at_ms = excluded.updated_at_ms",
- )
- .bind(SDK_RELEASE_PRODUCT_PROJECTION_ID)
- .bind(RADROOTS_TRADE_REDUCER_CONTRACT_ID)
- .bind(i64::from(RADROOTS_TRADE_REDUCER_VERSION))
- .bind(summary.last_event_seq.unwrap_or_default())
- .bind(source_digest.as_str())
- .bind(source_digest.as_str())
- .bind(refreshed_at_ms)
- .bind(refreshed_at_ms)
- .execute(sdk._event_store.pool())
- .await
- .map_err(|error| RadrootsSdkError::Projection {
- message: error.to_string(),
- })?;
- Ok(SyncProjectionRefreshReceipt::from_event_store_snapshot(
- summary,
- scanned_events,
- trade_upserts,
- refreshed_at_ms,
- ))
-}
-
-#[cfg(feature = "runtime")]
-async fn trade_mutation_count(sdk: &RadrootsClient) -> Result<usize, RadrootsSdkError> {
- let count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM trade_mutation")
- .fetch_one(sdk._event_store.pool())
- .await
- .map_err(|error| RadrootsSdkError::Projection {
- message: error.to_string(),
- })?;
- Ok(usize::try_from(count.max(0)).unwrap_or(usize::MAX))
-}
-
-#[cfg(feature = "runtime")]
-fn projection_source_digest(
- summary: &RadrootsEventStoreStatusSummary,
- trade_upserts: usize,
-) -> String {
- let mut hasher = Sha256::new();
- hasher.update(SDK_RELEASE_PRODUCT_PROJECTION_ID.as_bytes());
- hasher.update(b"\0");
- hasher.update(
- SDK_RELEASE_PRODUCT_PROJECTION_VERSION
- .to_string()
- .as_bytes(),
- );
- hasher.update(b"\0");
- hasher.update(summary.total_events.to_string().as_bytes());
- hasher.update(b"\0");
- hasher.update(summary.valid_stream_events.to_string().as_bytes());
- hasher.update(b"\0");
- hasher.update(
- summary
- .last_event_seq
- .unwrap_or_default()
- .to_string()
- .as_bytes(),
- );
- hasher.update(b"\0");
- hasher.update(summary.transport_observations.to_string().as_bytes());
- hasher.update(b"\0");
- hasher.update(trade_upserts.to_string().as_bytes());
- hex::encode(hasher.finalize())
-}
-
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-async fn push_radrootsd_claimed_outbox_event(
- sync: &SyncClient<'_>,
- adapter: &RadrootsdPublishAdapter,
- claimed: &RadrootsOutboxClaimedEvent,
- next_attempt_delay_ms: i64,
- now_ms: i64,
-) -> Result<PushOutboxEventReceipt, RadrootsSdkError> {
- let signed_event = claimed.signed_event.clone().ok_or(
- radroots_transport_nostr::RadrootsRelayTransportError::MissingSignedOutboxEvent(
- claimed.outbox_event_id,
- ),
- )?;
- let target_policy = match radrootsd_transport_publish_target_policy(claimed) {
- Ok(target_policy) => target_policy,
- Err(error) => {
- return fail_radrootsd_local_validation(sync, claimed, error, now_ms).await;
- }
- };
- let delivery_policy = match radrootsd_delivery_policy(sync, claimed).await {
- Ok(delivery_policy) => delivery_policy,
- Err(error) => {
- return fail_radrootsd_local_validation(sync, claimed, error, now_ms).await;
- }
- };
- sync.sdk
- ._outbox
- .ingest_signed_event_local(
- &sync.sdk._event_store,
- claimed.outbox_event_id,
- claimed.claim_token.as_str(),
- now_ms,
- )
- .await?;
- let request = RadrootsdPublishRequest {
- signed_event: signed_event.clone(),
- delivery_policy: delivery_policy.clone(),
- target_policy,
- idempotency_key: Some(radrootsd_outbox_idempotency_key(
- claimed.outbox_event_id,
- claimed.attempt_count,
- signed_event.id_str(),
- active_delivery_plan_id(claimed, "radrootsd publish")?,
- )),
- timeout_ms: adapter.config().request_timeout_ms,
- };
- let publish = match adapter.publish_signed_event(request).await {
- Ok(response) => response.job,
- Err(error) => {
- let message = radrootsd_error_message(&error);
- sync.sdk
- ._outbox
- .mark_publish_retryable(
- claimed.outbox_event_id,
- claimed.claim_token.as_str(),
- message.as_str(),
- now_ms.saturating_add(next_attempt_delay_ms),
- now_ms,
- )
- .await?;
- return radrootsd_transport_error_receipt(
- claimed,
- &signed_event,
- &delivery_policy,
- message,
- );
- }
- };
- complete_radrootsd_publish_attempt(sync, claimed, &publish, next_attempt_delay_ms, now_ms)
- .await?;
- push_radrootsd_event_receipt(claimed.outbox_event_id, publish)
-}
-
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-async fn fail_radrootsd_local_validation(
- sync: &SyncClient<'_>,
- claimed: &RadrootsOutboxClaimedEvent,
- error: RadrootsSdkError,
- now_ms: i64,
-) -> Result<PushOutboxEventReceipt, RadrootsSdkError> {
- let message = error.to_string();
- sync.sdk
- ._outbox
- .mark_active_delivery_plan_failed_terminal(
- claimed.outbox_event_id,
- claimed.claim_token.as_str(),
- message.as_str(),
- now_ms,
- )
- .await?;
- Err(error)
-}
-
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-async fn radrootsd_delivery_policy(
- sync: &SyncClient<'_>,
- claimed: &RadrootsOutboxClaimedEvent,
-) -> Result<TransportPublishDeliveryPolicy, RadrootsSdkError> {
- let active_delivery_plan_id = active_delivery_plan_id(claimed, "radrootsd publish")?;
- let plans = sync
- .sdk
- ._outbox
- .delivery_plans(claimed.outbox_event_id)
- .await?;
- let plan = plans
- .iter()
- .find(|plan| plan.delivery_plan_id == active_delivery_plan_id)
- .ok_or_else(|| RadrootsSdkError::InvalidRequest {
- message: format!(
- "outbox event {} active delivery plan {} was not found for radrootsd publish",
- claimed.outbox_event_id, active_delivery_plan_id
- ),
- })?;
- let targets = sync
- .sdk
- ._outbox
- .delivery_targets(claimed.outbox_event_id)
- .await?;
- let active_targets = targets
- .iter()
- .filter(|target| target.delivery_plan_id == active_delivery_plan_id)
- .collect::<Vec<_>>();
- reject_non_accepted_radrootsd_satisfaction(&plan.satisfaction_policy)?;
- let ready_target_count = active_targets
- .iter()
- .filter(|target| target.status.is_ready_for_attempt())
- .count();
- let required_remaining_targets =
- radrootsd_required_remaining_targets(&plan.satisfaction_policy, &active_targets)?;
- let required_remaining = if let Some(targets) = required_remaining_targets.as_ref() {
- targets.len()
- } else {
- let satisfied_count = plan
- .satisfaction_policy
- .target_satisfaction_class()
- .map(|satisfaction_class| {
- active_targets
- .iter()
- .filter(|target| {
- target
- .status
- .counts_as_transport_satisfaction(satisfaction_class)
- })
- .count()
- })
- .unwrap_or(0);
- (plan.required_success_count as usize).saturating_sub(satisfied_count)
- };
- radrootsd_delivery_policy_from_remaining(
- ready_target_count,
- required_remaining,
- required_remaining_targets.as_deref(),
- &plan.satisfaction_policy,
- )
-}
-
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-fn radrootsd_delivery_policy_from_remaining(
- ready_target_count: usize,
- required_remaining: usize,
- required_remaining_targets: Option<&[TargetFingerprint]>,
- satisfaction_policy: &RadrootsTransportSatisfactionPolicy,
-) -> Result<TransportPublishDeliveryPolicy, RadrootsSdkError> {
- reject_non_accepted_radrootsd_satisfaction(satisfaction_policy)?;
- if required_remaining == 0 {
- return Ok(TransportPublishDeliveryPolicy::Any);
- }
- if ready_target_count == 0 {
- if matches!(
- satisfaction_policy,
- RadrootsTransportSatisfactionPolicy::RequiredTargets { .. }
- ) {
- return Err(RadrootsSdkError::InvalidRequest {
- message: "radrootsd publish has unsatisfied required targets but no ready target"
- .to_owned(),
- });
- }
- return Ok(TransportPublishDeliveryPolicy::Any);
- }
- Ok(match satisfaction_policy {
- RadrootsTransportSatisfactionPolicy::NoWait => TransportPublishDeliveryPolicy::Any,
- RadrootsTransportSatisfactionPolicy::Any { .. } => TransportPublishDeliveryPolicy::Any,
- RadrootsTransportSatisfactionPolicy::All { .. } => TransportPublishDeliveryPolicy::All,
- RadrootsTransportSatisfactionPolicy::RequiredTargets { .. } => {
- let targets = required_remaining_targets
- .ok_or_else(|| RadrootsSdkError::InvalidRequest {
- message: "radrootsd publish missing required target fingerprints".to_owned(),
- })?
- .iter()
- .map(|fingerprint| {
- TransportPublishTargetFingerprint::parse(fingerprint.as_str().to_owned())
- .map_err(|error| RadrootsSdkError::InvalidRequest {
- message: error.to_string(),
- })
- })
- .collect::<Result<Vec<_>, _>>()?;
- TransportPublishDeliveryPolicy::required_targets(targets).map_err(|error| {
- RadrootsSdkError::InvalidRequest {
- message: error.to_string(),
- }
- })?
- }
- RadrootsTransportSatisfactionPolicy::Quorum { .. } => {
- if required_remaining >= ready_target_count {
- TransportPublishDeliveryPolicy::All
- } else if required_remaining == 1 {
- TransportPublishDeliveryPolicy::Any
- } else {
- TransportPublishDeliveryPolicy::Quorum {
- quorum: required_remaining,
- }
- }
- }
- })
-}
-
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-fn radrootsd_required_remaining_targets(
- satisfaction_policy: &RadrootsTransportSatisfactionPolicy,
- active_targets: &[&RadrootsOutboxDeliveryTargetRecord],
-) -> Result<Option<Vec<TargetFingerprint>>, RadrootsSdkError> {
- let RadrootsTransportSatisfactionPolicy::RequiredTargets { class, targets } =
- satisfaction_policy
- else {
- return Ok(None);
- };
- let mut remaining = Vec::new();
- for required in targets {
- let target = active_targets
- .iter()
- .find(|target| target.endpoint_fingerprint == *required)
- .ok_or_else(|| RadrootsSdkError::InvalidRequest {
- message: format!(
- "radrootsd publish required target {required} is not present in active delivery plan"
- ),
- })?;
- if target.status.counts_as_transport_satisfaction(*class) {
- continue;
- }
- if !target.status.is_ready_for_attempt() {
- return Err(RadrootsSdkError::InvalidRequest {
- message: format!(
- "radrootsd publish required target {required} is not ready for publish"
- ),
- });
- }
- remaining.push(required.clone());
- }
- Ok(Some(remaining))
-}
-
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-fn reject_non_accepted_radrootsd_satisfaction(
- satisfaction_policy: &RadrootsTransportSatisfactionPolicy,
-) -> Result<(), RadrootsSdkError> {
- if satisfaction_policy
- .target_satisfaction_class()
- .is_some_and(|class| class != RadrootsTransportSatisfactionClass::Accepted)
- {
- return Err(RadrootsSdkError::InvalidRequest {
- message: "radrootsd publish only supports accepted-class satisfaction policies"
- .to_owned(),
- });
- }
- Ok(())
-}
-
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-fn radrootsd_outbox_idempotency_key(
- outbox_event_id: i64,
- attempt_count: i64,
- event_id: &str,
- active_delivery_plan_id: i64,
-) -> String {
- format!(
- "radroots-sdk-outbox-{outbox_event_id}-{attempt_count}-{event_id}-{active_delivery_plan_id}"
- )
-}
-
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-fn active_delivery_plan_id(
- claimed: &RadrootsOutboxClaimedEvent,
- operation: &'static str,
-) -> Result<i64, RadrootsSdkError> {
- claimed
- .active_delivery_plan_id
- .ok_or_else(|| RadrootsSdkError::InvalidRequest {
- message: format!(
- "outbox event {} has no active delivery plan for {operation}",
- claimed.outbox_event_id
- ),
- })
-}
-
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-async fn complete_radrootsd_publish_attempt(
- sync: &SyncClient<'_>,
- claimed: &RadrootsOutboxClaimedEvent,
- publish: &TransportPublishJobView,
- next_attempt_delay_ms: i64,
- now_ms: i64,
-) -> Result<(), RadrootsSdkError> {
- let mut completed_target_ids = std::collections::BTreeSet::new();
- let mut matched_outcomes = Vec::new();
- for outcome in &publish.targets {
- let matched_targets = claimed
- .delivery_targets
- .iter()
- .filter(|target| target.status.is_ready_for_attempt())
- .filter(|target| radrootsd_target_matches_outcome(target, outcome))
- .collect::<Vec<_>>();
- if matched_targets.is_empty() {
- continue;
- }
- if matched_targets.len() > 1 {
- return Err(RadrootsSdkError::InvalidRequest {
- message: format!(
- "radrootsd publish outcome for {} {} matched multiple ready delivery targets on outbox event {}",
- outcome.transport_kind, outcome.endpoint_uri, claimed.outbox_event_id
- ),
- });
- }
- let target = matched_targets[0];
- if !completed_target_ids.insert(target.delivery_target_id) {
- return Err(RadrootsSdkError::InvalidRequest {
- message: format!(
- "radrootsd publish outcome for {} {} matched delivery target {} more than once on outbox event {}",
- outcome.transport_kind,
- outcome.endpoint_uri,
- target.delivery_target_id,
- claimed.outbox_event_id
- ),
- });
- }
- matched_outcomes.push((target, outcome));
- }
- for (target, outcome) in matched_outcomes {
- complete_radrootsd_delivery_target(sync, claimed, target, outcome, now_ms).await?;
- }
- for target in claimed
- .delivery_targets
- .iter()
- .filter(|target| target.status.is_ready_for_attempt())
- .filter(|target| !completed_target_ids.contains(&target.delivery_target_id))
- {
- complete_missing_radrootsd_delivery_target(sync, claimed, target, publish, now_ms).await?;
- }
- sync.sdk
- ._outbox
- .complete_publish_attempt(
- claimed.outbox_event_id,
- claimed.claim_token.as_str(),
- "radrootsd publish incomplete",
- "radrootsd publish terminal",
- now_ms.saturating_add(next_attempt_delay_ms),
- now_ms,
- )
- .await?;
- Ok(())
-}
-
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-fn radrootsd_transport_publish_target_policy(
- claimed: &RadrootsOutboxClaimedEvent,
-) -> Result<TransportPublishTargetPolicy, RadrootsSdkError> {
- let ready_targets = claimed
- .delivery_targets
- .iter()
- .filter(|target| target.status.is_ready_for_attempt())
- .collect::<Vec<_>>();
- Ok(TransportPublishTargetPolicy::explicit_targets(
- ready_targets
- .into_iter()
- .map(transport_publish_target_from_outbox_target)
- .collect::<Result<Vec<_>, _>>()?,
- ))
-}
-
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-fn transport_publish_target_from_outbox_target(
- target: &RadrootsOutboxDeliveryTargetRecord,
-) -> Result<TransportPublishTarget, RadrootsSdkError> {
- if target.transport_kind != TransportId::NOSTR {
- return Err(RadrootsSdkError::InvalidRequest {
- message: format!(
- "radrootsd execution explicit targets are Nostr-only and cannot publish {} target {}",
- target.transport_kind.canonical_label(),
- target.endpoint_uri.as_str()
- ),
- });
- }
- Ok(TransportPublishTarget {
- transport_kind: target.transport_kind.canonical_label(),
- endpoint_uri: target.endpoint_uri.as_str().to_owned(),
- target_scope: target
- .target_scope
- .as_ref()
- .map(|scope| scope.as_str().to_owned()),
- target_label: target
- .target_label
- .as_ref()
- .map(|label| label.as_str().to_owned()),
- reticulum_behavior: None,
- })
-}
-
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-fn radrootsd_target_matches_outcome(
- target: &RadrootsOutboxDeliveryTargetRecord,
- outcome: &TransportPublishTargetOutcome,
-) -> bool {
- target.transport_kind.canonical_label() == outcome.transport_kind
- && target.endpoint_uri.as_str() == outcome.endpoint_uri
- && target.target_scope.as_ref().map(|scope| scope.as_str())
- == outcome.target_scope.as_deref()
-}
-
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-async fn complete_radrootsd_delivery_target(
- sync: &SyncClient<'_>,
- claimed: &RadrootsOutboxClaimedEvent,
- target: &RadrootsOutboxDeliveryTargetRecord,
- outcome: &TransportPublishTargetOutcome,
- now_ms: i64,
-) -> Result<(), RadrootsSdkError> {
- if outcome.outcome_kind.counts_toward_accepted_delivery() {
- sync.sdk
- ._outbox
- .mark_delivery_target_accepted(
- claimed.outbox_event_id,
- claimed.claim_token.as_str(),
- target.delivery_target_id,
- now_ms,
- )
- .await?;
- } else if outcome.outcome_kind.is_retryable() {
- sync.sdk
- ._outbox
- .mark_delivery_target_failed_retryable(
- claimed.outbox_event_id,
- claimed.claim_token.as_str(),
- target.delivery_target_id,
- outcome
- .message
- .as_deref()
- .unwrap_or("radrootsd publish retryable"),
- now_ms,
- )
- .await?;
- } else if outcome.outcome_kind == TransportPublishOutcomeKind::DeferredUntilImplemented {
- sync.sdk
- ._outbox
- .mark_delivery_target_deferred_until_implemented(
- claimed.outbox_event_id,
- claimed.claim_token.as_str(),
- target.delivery_target_id,
- outcome
- .message
- .as_deref()
- .unwrap_or("radrootsd publish deferred until implemented"),
- now_ms,
- )
- .await?;
- } else {
- sync.sdk
- ._outbox
- .mark_delivery_target_failed_terminal(
- claimed.outbox_event_id,
- claimed.claim_token.as_str(),
- target.delivery_target_id,
- outcome
- .message
- .as_deref()
- .unwrap_or("radrootsd publish terminal"),
- now_ms,
- )
- .await?;
- }
- Ok(())
-}
-
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-async fn complete_missing_radrootsd_delivery_target(
- sync: &SyncClient<'_>,
- claimed: &RadrootsOutboxClaimedEvent,
- target: &RadrootsOutboxDeliveryTargetRecord,
- publish: &TransportPublishJobView,
- now_ms: i64,
-) -> Result<(), RadrootsSdkError> {
- if matches!(
- publish.status,
- TransportPublishJobStatus::DeliveryDeferred
- | TransportPublishJobStatus::DeliveryDeferredUntilImplemented
- ) {
- sync.sdk
- ._outbox
- .mark_delivery_target_deferred_until_implemented(
- claimed.outbox_event_id,
- claimed.claim_token.as_str(),
- target.delivery_target_id,
- "radrootsd publish deferred until implemented",
- now_ms,
- )
- .await?;
- } else if publish.retryable_count > 0
- || !publish.terminal
- || target.status == RadrootsOutboxDeliveryTargetStatus::FailedRetryable
- {
- sync.sdk
- ._outbox
- .mark_delivery_target_failed_retryable(
- claimed.outbox_event_id,
- claimed.claim_token.as_str(),
- target.delivery_target_id,
- "radrootsd publish incomplete",
- now_ms,
- )
- .await?;
- } else {
- sync.sdk
- ._outbox
- .mark_delivery_target_failed_terminal(
- claimed.outbox_event_id,
- claimed.claim_token.as_str(),
- target.delivery_target_id,
- "radrootsd publish terminal",
- now_ms,
- )
- .await?;
- }
- Ok(())
-}
-
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-fn radrootsd_error_message(error: &RadrootsdError) -> String {
- format!("radrootsd publish failed: {error}")
-}
-
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-fn radrootsd_transport_error_receipt(
- claimed: &RadrootsOutboxClaimedEvent,
- event: &radroots_event::draft::SignedEvent,
- delivery_policy: &TransportPublishDeliveryPolicy,
- message: String,
-) -> Result<PushOutboxEventReceipt, RadrootsSdkError> {
- let ready_targets = claimed
- .delivery_targets
- .iter()
- .filter(|target| target.status.is_ready_for_attempt())
- .collect::<Vec<_>>();
- let target_count = ready_targets.len();
- let event_id = push_receipt_event_id(event.id_str(), "radrootsd transport failure event id")?;
- Ok(PushOutboxEventReceipt {
- event_id,
- outbox_event_id: claimed.outbox_event_id,
- final_state: PushOutboxEventState::PublishRetryable,
- attempted_count: 0,
- accepted_count: 0,
- retryable_count: target_count,
- terminal_count: 0,
- quorum: delivery_policy.required_target_count(target_count),
- quorum_met: false,
- targets: ready_targets
- .into_iter()
- .map(|target| PushOutboxTargetReceipt {
- transport_kind: target.transport_kind.canonical_label(),
- endpoint_uri: target.endpoint_uri.as_str().to_owned(),
- target_scope: target
- .target_scope
- .as_ref()
- .map(|scope| scope.as_str().to_owned()),
- target_label: target
- .target_label
- .as_ref()
- .map(|label| label.as_str().to_owned()),
- outcome_kind: PushOutboxTargetOutcomeKind::ConnectionFailed,
- transport_outcome_kind: Some(PushOutboxTransportOutcomeKind::ConnectionFailed),
- attempted: false,
- message: Some(message.clone()),
- })
- .collect(),
- })
-}
-
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-fn radrootsd_push_event_final_state(publish: &TransportPublishJobView) -> PushOutboxEventState {
- if publish.delivery_satisfied {
- PushOutboxEventState::Published
- } else if matches!(
- publish.status,
- TransportPublishJobStatus::DeliveryDeferred
- | TransportPublishJobStatus::DeliveryDeferredUntilImplemented
- ) {
- PushOutboxEventState::DeferredUntilImplemented
- } else if publish.retryable_count > 0 || !publish.terminal {
- PushOutboxEventState::PublishRetryable
- } else {
- PushOutboxEventState::FailedTerminal
- }
-}
-
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-fn push_radrootsd_event_receipt(
- outbox_event_id: i64,
- publish: TransportPublishJobView,
-) -> Result<PushOutboxEventReceipt, RadrootsSdkError> {
- let event_id = push_receipt_event_id(
- publish.event_id.as_str(),
- "transport publish daemon job event id",
- )?;
- let quorum = publish
- .delivery_policy
- .required_target_count(publish.target_count);
- Ok(PushOutboxEventReceipt {
- event_id,
- outbox_event_id,
- final_state: radrootsd_push_event_final_state(&publish),
- attempted_count: publish
- .targets
- .iter()
- .filter(|target| target.attempted)
- .count(),
- accepted_count: publish.acknowledged_count,
- retryable_count: publish.retryable_count,
- terminal_count: publish.terminal_count,
- quorum,
- quorum_met: publish.delivery_satisfied,
- targets: publish
- .targets
- .into_iter()
- .map(push_radrootsd_target_receipt)
- .collect(),
- })
-}
-
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-fn push_radrootsd_target_receipt(
- outcome: TransportPublishTargetOutcome,
-) -> PushOutboxTargetReceipt {
- PushOutboxTargetReceipt {
- transport_kind: outcome.transport_kind,
- endpoint_uri: outcome.endpoint_uri,
- target_scope: outcome.target_scope,
- target_label: outcome.target_label,
- outcome_kind: push_radrootsd_target_outcome_kind(outcome.outcome_kind),
- transport_outcome_kind: Some(push_radrootsd_transport_outcome_kind(outcome.outcome_kind)),
- attempted: outcome.attempted,
- message: outcome.message,
- }
-}
-
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-fn push_radrootsd_target_outcome_kind(
- outcome_kind: TransportPublishOutcomeKind,
-) -> PushOutboxTargetOutcomeKind {
- match outcome_kind {
- TransportPublishOutcomeKind::Accepted => PushOutboxTargetOutcomeKind::Accepted,
- TransportPublishOutcomeKind::DuplicateAccepted => {
- PushOutboxTargetOutcomeKind::DuplicateAccepted
- }
- TransportPublishOutcomeKind::Blocked => PushOutboxTargetOutcomeKind::Blocked,
- TransportPublishOutcomeKind::RateLimited => PushOutboxTargetOutcomeKind::RateLimited,
- TransportPublishOutcomeKind::Invalid => PushOutboxTargetOutcomeKind::Invalid,
- TransportPublishOutcomeKind::PowRequired => PushOutboxTargetOutcomeKind::PowRequired,
- TransportPublishOutcomeKind::Restricted => PushOutboxTargetOutcomeKind::Restricted,
- TransportPublishOutcomeKind::AuthRequired => PushOutboxTargetOutcomeKind::AuthRequired,
- TransportPublishOutcomeKind::Muted => PushOutboxTargetOutcomeKind::Muted,
- TransportPublishOutcomeKind::Unsupported => PushOutboxTargetOutcomeKind::Unsupported,
- TransportPublishOutcomeKind::PaymentRequired => {
- PushOutboxTargetOutcomeKind::PaymentRequired
- }
- TransportPublishOutcomeKind::Error => PushOutboxTargetOutcomeKind::Error,
- TransportPublishOutcomeKind::Timeout => PushOutboxTargetOutcomeKind::Timeout,
- TransportPublishOutcomeKind::ConnectionFailed => {
- PushOutboxTargetOutcomeKind::ConnectionFailed
- }
- TransportPublishOutcomeKind::TargetRejected => {
- PushOutboxTargetOutcomeKind::TargetUriRejected
- }
- TransportPublishOutcomeKind::SkippedAlreadyAccepted => {
- PushOutboxTargetOutcomeKind::SkippedAlreadyAccepted
- }
- TransportPublishOutcomeKind::DeferredUntilImplemented => {
- PushOutboxTargetOutcomeKind::DeferredUntilImplemented
- }
- TransportPublishOutcomeKind::Unknown => PushOutboxTargetOutcomeKind::Unknown,
- }
-}
-
-#[cfg(all(feature = "runtime", feature = "radrootsd-execution"))]
-fn push_radrootsd_transport_outcome_kind(
- outcome_kind: TransportPublishOutcomeKind,
-) -> PushOutboxTransportOutcomeKind {
- match outcome_kind {
- TransportPublishOutcomeKind::Accepted => PushOutboxTransportOutcomeKind::Accepted,
- TransportPublishOutcomeKind::DuplicateAccepted
- | TransportPublishOutcomeKind::SkippedAlreadyAccepted => {
- PushOutboxTransportOutcomeKind::DuplicateAccepted
- }
- TransportPublishOutcomeKind::DeferredUntilImplemented => {
- PushOutboxTransportOutcomeKind::DeferredUntilImplemented
- }
- TransportPublishOutcomeKind::Blocked
- | TransportPublishOutcomeKind::Invalid
- | TransportPublishOutcomeKind::Restricted
- | TransportPublishOutcomeKind::Muted
- | TransportPublishOutcomeKind::Unsupported
- | TransportPublishOutcomeKind::TargetRejected => PushOutboxTransportOutcomeKind::Rejected,
- TransportPublishOutcomeKind::PaymentRequired
- | TransportPublishOutcomeKind::PowRequired
- | TransportPublishOutcomeKind::AuthRequired => PushOutboxTransportOutcomeKind::PolicyDenied,
- TransportPublishOutcomeKind::Timeout => PushOutboxTransportOutcomeKind::Timeout,
- TransportPublishOutcomeKind::ConnectionFailed => {
- PushOutboxTransportOutcomeKind::ConnectionFailed
- }
- TransportPublishOutcomeKind::RateLimited
- | TransportPublishOutcomeKind::Error
- | TransportPublishOutcomeKind::Unknown => {
- PushOutboxTransportOutcomeKind::TransportUnavailable
- }
- }
-}
-
-#[cfg(feature = "runtime")]
-fn push_outbox_claim_token() -> String {
- format!("radroots-sdk-sync-{}", uuid::Uuid::now_v7())
-}
-
-#[cfg(feature = "runtime")]
-fn push_event_final_state(publish: &RadrootsOutboxPublishReceipt) -> PushOutboxEventState {
- if publish.quorum_met {
- PushOutboxEventState::Published
- } else if publish.retryable_count > 0 {
- PushOutboxEventState::PublishRetryable
- } else {
- PushOutboxEventState::FailedTerminal
- }
-}
-
-#[cfg(feature = "runtime")]
-fn push_event_receipt(
- outbox_event_id: i64,
- final_state: PushOutboxEventState,
- publish: RadrootsOutboxPublishReceipt,
-) -> Result<PushOutboxEventReceipt, RadrootsSdkError> {
- let event_id = push_receipt_event_id(
- publish.event_id.as_str(),
- "direct Nostr outbox publish receipt event id",
- )?;
- Ok(PushOutboxEventReceipt {
- event_id,
- outbox_event_id,
- final_state,
- attempted_count: publish.attempted_count,
- accepted_count: publish.accepted_count,
- retryable_count: publish.retryable_count,
- terminal_count: publish.terminal_count,
- quorum: publish.quorum,
- quorum_met: publish.quorum_met,
- targets: publish
- .target_receipts
- .into_iter()
- .map(push_target_receipt)
- .collect(),
- })
-}
-
-#[cfg(feature = "runtime")]
-fn push_receipt_event_id(value: &str, field: &str) -> Result<EventId, RadrootsSdkError> {
- EventId::parse(value).map_err(|error| RadrootsSdkError::InvalidRequest {
- message: format!("{field} is invalid: {error}"),
- })
-}
-
-#[cfg(feature = "runtime")]
-fn push_target_receipt(target: RadrootsOutboxPublishTargetReceipt) -> PushOutboxTargetReceipt {
- PushOutboxTargetReceipt {
- transport_kind: TransportId::NOSTR.canonical_label(),
- endpoint_uri: target.endpoint_uri,
- target_scope: target.target_scope,
- target_label: target.target_label,
- outcome_kind: target.outcome.kind.into(),
- transport_outcome_kind: Some(target.outcome.kind.transport_outcome_kind().into()),
- attempted: target.attempted,
- message: target.outcome.message,
- }
-}
-
-#[cfg(all(test, feature = "runtime"))]
-#[path = "../tests/unit/sync_runtime_tests.rs"]
-mod tests;
diff --git a/crates/sdk/tests/package_boundary.rs b/crates/sdk/tests/package_boundary.rs
@@ -3,6 +3,7 @@ use std::collections::BTreeSet;
const MANIFEST: &str = include_str!("../Cargo.toml");
const ROOT: &str = include_str!("../src/lib.rs");
const CLIENT: &str = include_str!("../src/client.rs");
+const SYNC: &str = include_str!("../src/sync.rs");
const TRANSPORT: &str = include_str!("../src/transport.rs");
#[test]
@@ -158,6 +159,42 @@ fn transport_profiles_reuse_canonical_types_and_forbid_fallback() {
}
}
+#[test]
+fn sync_operations_only_delegate_to_the_canonical_engine() {
+ for delegation in [
+ "self.engine.pull(request, admission).await",
+ "self.engine.ingest(observed, admission).await",
+ "self.engine.ingest_batch(observed, admission).await",
+ "self.engine.refresh_projection(request, reducer).await",
+ "self.engine.sign_and_enqueue(request).await",
+ "self.engine.deliver_pending(request).await",
+ "self.engine.status(projections).await",
+ "self.engine.retry_decision(record, now_unix_ms)",
+ ] {
+ assert!(
+ SYNC.contains(delegation),
+ "missing sync delegation `{delegation}`"
+ );
+ }
+ for duplicate in [
+ "pub struct SyncTransport",
+ "pub struct SyncOutbox",
+ "pub struct SyncProjection",
+ "pub struct SyncStatus",
+ "pub struct PushOutbox",
+ ] {
+ assert!(
+ !SYNC.contains(duplicate),
+ "SDK sync source contains duplicate lower model `{duplicate}`"
+ );
+ }
+ assert!(
+ !std::path::Path::new(env!("CARGO_MANIFEST_DIR"))
+ .join("src/sync_runtime.rs")
+ .exists()
+ );
+}
+
fn dependency_names(manifest: &str) -> BTreeSet<&str> {
let dependencies = manifest
.split_once("[dependencies]")
diff --git a/crates/sdk/tests/unit/sync_runtime_tests.rs b/crates/sdk/tests/unit/sync_runtime_tests.rs
@@ -1,1785 +0,0 @@
-#[cfg(feature = "radrootsd-execution")]
-use super::{
- CLAIM_OWNER, complete_radrootsd_publish_attempt, push_radrootsd_claimed_outbox_event,
- push_radrootsd_event_receipt, radrootsd_delivery_policy_from_remaining,
- radrootsd_error_message, radrootsd_outbox_idempotency_key,
- radrootsd_required_remaining_targets, radrootsd_transport_error_receipt,
- transport_publish_target_from_outbox_target,
-};
-use super::{
- PushOutboxEventReceipt, PushOutboxEventState, PushOutboxReceipt, PushOutboxTargetOutcomeKind,
- PushOutboxTransportOutcomeKind, SdkRelayAuthPolicy, SyncEventStoreStatus, SyncOutboxStatus,
- push_event_final_state, push_event_receipt, push_outbox_claim_token,
-};
-use crate::RadrootsSdkError;
-#[cfg(feature = "radrootsd-execution")]
-use crate::adapters::radrootsd::{RadrootsdError, RadrootsdPublishAdapter, RadrootsdPublishConfig};
-#[cfg(feature = "radrootsd-execution")]
-use crate::workflow_runtime::{SdkWorkflowEnqueueRequest, enqueue_signed_workflow};
-use futures::future::BoxFuture;
-use nostr::Keys as RadrootsNostrKeys;
-#[cfg(feature = "radrootsd-execution")]
-use radroots_event::contract::AuthorRole;
-#[cfg(feature = "radrootsd-execution")]
-use radroots_event::draft::EventDraft;
-#[cfg(feature = "radrootsd-execution")]
-use radroots_event::envelope::kind::KIND_FARM;
-use radroots_event::id::EventId;
-use radroots_event_store::RadrootsEventStoreStatusSummary;
-#[cfg(feature = "radrootsd-execution")]
-use radroots_nostr::signing::sign_frozen_draft;
-#[cfg(feature = "radrootsd-execution")]
-use radroots_outbox::{
- RadrootsOutboxClaimedEvent, RadrootsOutboxDeliveryPlanInput, RadrootsOutboxDeliveryPlanStatus,
- RadrootsOutboxDeliveryTargetRecord, RadrootsOutboxDeliveryTargetStatus,
- RadrootsOutboxOperationInput, RadrootsOutboxSignedOperationInput,
-};
-use radroots_outbox::{
- RadrootsOutboxEventState, RadrootsOutboxEventStoreIngestReceipt, RadrootsOutboxStatusSummary,
-};
-#[cfg(feature = "radrootsd-execution")]
-use radroots_protocol::radrootsd::transport_publish::v5::{
- DeliveryPolicy as TransportPublishDeliveryPolicy, Job as TransportPublishJobView,
- JobStatus as TransportPublishJobStatus,
- NostrTargetSourcePolicy as NostrPublishTargetSourcePolicy,
- OutcomeKind as TransportPublishOutcomeKind, Target as TransportPublishTarget,
- TargetFingerprint as TransportPublishTargetFingerprint,
- TargetOutcome as TransportPublishTargetOutcome, TargetPolicy as TransportPublishTargetPolicy,
- TargetSource as TransportPublishTargetSource,
-};
-#[cfg(feature = "radrootsd-execution")]
-use radroots_signing::{
- Actor, Error as SigningError, SignReceipt, SignRequest, Signer, SignerStatus,
- actor::ActorSource, error::Kind as SigningErrorKind, signer::BoxFuture as SigningFuture,
-};
-use radroots_transport::target::{TargetLabel, TargetScope};
-use radroots_transport::{RadrootsTransportDeliveryTargetStatus, Target, TransportId};
-#[cfg(feature = "radrootsd-execution")]
-use radroots_transport::{RadrootsTransportSatisfactionClass, RadrootsTransportSatisfactionPolicy};
-use radroots_transport_nostr::{
- RadrootsNostrTransport, RadrootsOutboxPublishReceipt, RadrootsOutboxPublishTargetReceipt,
- RadrootsRelayOutcomeKind, RadrootsRelayPublishAdapter, RadrootsRelayPublishRelayReceipt,
- RadrootsRelayPublishRequest, RadrootsRelayTransportError,
-};
-use std::collections::BTreeSet;
-#[cfg(feature = "radrootsd-execution")]
-use std::io::ErrorKind;
-#[cfg(feature = "radrootsd-execution")]
-use std::net::TcpListener;
-#[cfg(feature = "radrootsd-execution")]
-use std::sync::LazyLock;
-#[cfg(feature = "radrootsd-execution")]
-use std::time::Duration;
-
-#[cfg(feature = "radrootsd-execution")]
-static RADROOTSD_FIXTURE_SIGNER_KEYS: LazyLock<RadrootsNostrKeys> =
- LazyLock::new(RadrootsNostrKeys::generate);
-#[cfg(feature = "radrootsd-execution")]
-static RADROOTSD_FIXTURE_SIGNER_PUBLIC_KEY: LazyLock<String> =
- LazyLock::new(|| RADROOTSD_FIXTURE_SIGNER_KEYS.public_key().to_hex());
-
-#[cfg(feature = "radrootsd-execution")]
-fn radrootsd_fixture_signer_pubkey() -> &'static str {
- RADROOTSD_FIXTURE_SIGNER_PUBLIC_KEY.as_str()
-}
-
-struct UnusedPublishAdapter;
-
-#[cfg(feature = "radrootsd-execution")]
-struct RadrootsdFixtureSigner {
- keys: RadrootsNostrKeys,
-}
-
-#[cfg(feature = "radrootsd-execution")]
-impl RadrootsdFixtureSigner {
- fn new() -> Self {
- Self {
- keys: RADROOTSD_FIXTURE_SIGNER_KEYS.clone(),
- }
- }
-
- fn sign_frozen_draft(
- &self,
- draft: &EventDraft,
- ) -> Result<radroots_event::SignedEvent, radroots_nostr::Error> {
- sign_frozen_draft(&self.keys, draft)
- }
-}
-
-#[cfg(feature = "radrootsd-execution")]
-impl Signer for RadrootsdFixtureSigner {
- fn status(&self) -> SigningFuture<'_, Result<SignerStatus, SigningError>> {
- Box::pin(async { Ok(SignerStatus::unavailable()) })
- }
-
- fn sign(&self, request: SignRequest) -> SigningFuture<'_, Result<SignReceipt, SigningError>> {
- Box::pin(async move {
- let signed_event =
- sign_frozen_draft(&self.keys, request.draft()).map_err(|source| {
- SigningError::with_source(SigningErrorKind::InternalError, source)
- })?;
- SignReceipt::from_signed_event(&request, signed_event, 1_700_000_001)
- })
- }
-}
-
-impl RadrootsRelayPublishAdapter for UnusedPublishAdapter {
- fn publish<'a>(
- &'a self,
- _request: RadrootsRelayPublishRequest,
- ) -> BoxFuture<'a, Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>>
- {
- Box::pin(async { Ok(Vec::new()) })
- }
-}
-
-#[cfg(feature = "radrootsd-execution")]
-fn radrootsd_actor() -> Actor {
- Actor::from_public_key_hex(
- radrootsd_fixture_signer_pubkey(),
- ActorSource::ExplicitPublicKey,
- [AuthorRole::Farmer],
- )
- .expect("actor")
-}
-
-#[cfg(feature = "radrootsd-execution")]
-fn radrootsd_frozen_draft(d_tag: &str) -> EventDraft {
- EventDraft::new(
- "radroots.farm.profile.v1",
- KIND_FARM,
- 1_700_000_000,
- vec![vec!["d".to_owned(), d_tag.to_owned()]],
- "{}",
- radrootsd_fixture_signer_pubkey(),
- )
- .expect("frozen draft")
-}
-
-#[cfg(feature = "radrootsd-execution")]
-async fn claimed_radrootsd_event(
- d_tag: &str,
-) -> (crate::RadrootsClient, RadrootsOutboxClaimedEvent) {
- let sdk = crate::RadrootsClient::builder()
- .fixed_clock(crate::RadrootsSdkTimestamp::from_unix_seconds(
- 1_700_000_000,
- ))
- .build()
- .await
- .expect("sdk");
- let actor = radrootsd_actor();
- let draft = radrootsd_frozen_draft(d_tag);
- enqueue_signed_workflow(
- &sdk,
- SdkWorkflowEnqueueRequest {
- operation_kind: "sync.push.v1",
- actor: &actor,
- frozen_draft: &draft,
- target_policy: crate::TargetPolicy::try_nostr_relays(
- ["wss://relay.example.com"],
- crate::NostrRelayUrlPolicy::Public,
- )
- .expect("target relays"),
- satisfaction_policy: crate::SatisfactionPolicy::AllAccepted,
- idempotency_key: Some(
- crate::SdkIdempotencyKey::new("01890f0e-6c00-7000-8000-00000000025a")
- .expect("idempotency key"),
- ),
- },
- &RadrootsdFixtureSigner::new(),
- )
- .await
- .expect("enqueue signed workflow");
- let claimed = sdk
- ._outbox
- .claim_next_ready_signed_event(
- CLAIM_OWNER,
- "radrootsd-unit-claim",
- 1_700_000_060_000,
- 1_700_000_000_000,
- )
- .await
- .expect("claim")
- .expect("claimed event");
- (sdk, claimed)
-}
-
-#[cfg(feature = "radrootsd-execution")]
-async fn claimed_uningested_radrootsd_event(
- d_tag: &str,
- radrootsd_endpoint: &str,
-) -> (crate::RadrootsClient, RadrootsOutboxClaimedEvent) {
- claimed_uningested_radrootsd_event_with_satisfaction(
- d_tag,
- radrootsd_endpoint,
- RadrootsTransportSatisfactionPolicy::all_accepted(),
- )
- .await
-}
-
-#[cfg(feature = "radrootsd-execution")]
-async fn claimed_uningested_radrootsd_event_with_satisfaction(
- d_tag: &str,
- _radrootsd_endpoint: &str,
- satisfaction_policy: RadrootsTransportSatisfactionPolicy,
-) -> (crate::RadrootsClient, RadrootsOutboxClaimedEvent) {
- let sdk = crate::RadrootsClient::builder()
- .fixed_clock(crate::RadrootsSdkTimestamp::from_unix_seconds(
- 1_700_000_000,
- ))
- .build()
- .await
- .expect("sdk");
- let draft = radrootsd_frozen_draft(d_tag);
- let radrootsd_target =
- Target::new(TransportId::NOSTR, "wss://relay.example.com").expect("Nostr target");
- let enqueue = sdk
- ._outbox
- .enqueue_operation(
- RadrootsOutboxOperationInput::new(
- "sync.radrootsd.unit.v1",
- draft,
- RadrootsOutboxDeliveryPlanInput::new(
- "radrootsd",
- 1,
- satisfaction_policy,
- vec![radrootsd_target],
- ),
- 1_700_000_000_000,
- )
- .with_idempotency_key(format!("radrootsd-uningested-{d_tag}")),
- )
- .await
- .expect("enqueue");
- let signing_claim = sdk
- ._outbox
- .claim_next_ready_event(
- CLAIM_OWNER,
- "radrootsd-unit-sign",
- 1_700_000_000_500,
- 1_700_000_000_000,
- )
- .await
- .expect("signing claim")
- .expect("signing claim");
- let signed_event = RadrootsdFixtureSigner::new()
- .sign_frozen_draft(&signing_claim.draft)
- .expect("signed event");
- sdk._outbox
- .complete_signing(
- enqueue.outbox_event_id,
- "radrootsd-unit-sign",
- signed_event,
- 1_700_000_000_100,
- )
- .await
- .expect("complete signing");
- let stored = sdk
- ._outbox
- .get_event(enqueue.outbox_event_id)
- .await
- .expect("stored")
- .expect("stored");
- assert_eq!(stored.state, RadrootsOutboxEventState::Signed);
- assert!(!stored.event_store_ingested);
- assert_eq!(stored.event_store_ingested_at_ms, None);
- assert_eq!(stored.claim_token, None);
- let claimed = sdk
- ._outbox
- .claim_next_ready_signed_event(
- CLAIM_OWNER,
- "radrootsd-unit-publish",
- 1_700_000_060_000,
- 1_700_000_000_100,
- )
- .await
- .expect("publish claim")
- .expect("publish claim");
- (sdk, claimed)
-}
-
-#[cfg(feature = "radrootsd-execution")]
-fn assert_no_transport_publish_request(listener: &TcpListener) {
- listener.set_nonblocking(true).expect("nonblocking");
- match listener.accept() {
- Err(error) if error.kind() == ErrorKind::WouldBlock => {}
- Err(error) => panic!("transport publish listener failed before request check: {error}"),
- Ok(_) => panic!("transport publish listener received a request after local validation"),
- }
-}
-
-#[cfg(feature = "radrootsd-execution")]
-fn delivery_target_record(
- delivery_target_id: i64,
- delivery_plan_id: i64,
- target: &Target,
-) -> RadrootsOutboxDeliveryTargetRecord {
- RadrootsOutboxDeliveryTargetRecord {
- delivery_target_id,
- delivery_plan_id,
- transport_kind: *target.kind(),
- endpoint_uri: target.uri().clone(),
- target_scope: target.scope().cloned(),
- target_label: target.label().cloned(),
- endpoint_fingerprint: target.fingerprint().clone(),
- status: RadrootsOutboxDeliveryTargetStatus::Pending,
- last_outcome_kind: None,
- attempt_count: 0,
- last_attempt_at_ms: None,
- completed_at_ms: None,
- last_error: None,
- }
-}
-
-#[cfg(feature = "radrootsd-execution")]
-fn radrootsd_job(
- event_id: &str,
- outcome_kind: TransportPublishOutcomeKind,
-) -> TransportPublishJobView {
- let delivery_satisfied = outcome_kind.counts_toward_accepted_delivery();
- let retryable = outcome_kind.is_retryable();
- let terminal_failure = outcome_kind.is_terminal_failure();
- let status = if delivery_satisfied {
- TransportPublishJobStatus::DeliverySatisfied
- } else if retryable {
- TransportPublishJobStatus::DeliveryUnsatisfiedRetryable
- } else if outcome_kind == TransportPublishOutcomeKind::DeferredUntilImplemented {
- TransportPublishJobStatus::DeliveryDeferred
- } else {
- TransportPublishJobStatus::DeliveryUnsatisfiedTerminal
- };
- TransportPublishJobView {
- job_id: "radrootsd-unit-job".to_owned(),
- status,
- terminal: !retryable,
- delivery_satisfied,
- event_id: event_id.to_owned(),
- pubkey: radrootsd_fixture_signer_pubkey().to_owned(),
- event_kind: KIND_FARM,
- target_policy: TransportPublishTargetPolicy::nostr(
- NostrPublishTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault,
- vec!["wss://relay.example.com".to_owned()],
- ),
- delivery_policy: TransportPublishDeliveryPolicy::Any,
- target_count: 1,
- acknowledged_count: usize::from(delivery_satisfied),
- retryable_count: usize::from(retryable),
- terminal_count: usize::from(terminal_failure),
- requested_at_ms: 1_700_000_000_000,
- completed_at_ms: Some(1_700_000_000_100),
- last_error: None,
- targets: vec![TransportPublishTargetOutcome {
- transport_kind: "nostr".to_owned(),
- endpoint_uri: "wss://relay.example.com".to_owned(),
- target_scope: None,
- target_label: None,
- source: TransportPublishTargetSource::Request,
- attempted: true,
- outcome_kind,
- message: Some("daemon outcome".to_owned()),
- latency_ms: Some(4),
- }],
- }
-}
-
-#[test]
-fn push_outbox_claim_tokens_are_unique_under_immediate_generation() {
- let mut tokens = BTreeSet::new();
- for _ in 0..1_024 {
- let token = push_outbox_claim_token();
- assert!(token.starts_with("radroots-sdk-sync-"));
- assert!(tokens.insert(token));
- }
-}
-
-#[test]
-fn push_event_receipt_parses_typed_event_id() {
- let event_id = "a".repeat(64);
- let receipt = push_event_receipt(
- 1,
- PushOutboxEventState::Published,
- outbox_publish_receipt(event_id.as_str()).with_target(),
- )
- .expect("receipt");
-
- assert_eq!(
- receipt.event_id,
- EventId::parse(event_id).expect("event id")
- );
- assert_eq!(receipt.targets.len(), 1);
- assert_eq!(receipt.targets[0].transport_kind, "nostr");
- assert_eq!(receipt.targets[0].endpoint_uri, "wss://relay.example.com");
- assert_eq!(
- receipt.targets[0].target_scope.as_deref(),
- Some("farm.local")
- );
- assert_eq!(
- receipt.targets[0].target_label.as_deref(),
- Some("Farm relay")
- );
- assert!(receipt.targets[0].attempted);
-}
-
-#[test]
-fn push_event_receipt_returns_typed_error_for_invalid_internal_event_id() {
- let error = push_event_receipt(
- 1,
- PushOutboxEventState::Published,
- outbox_publish_receipt("not-a-valid-event-id"),
- )
- .expect_err("invalid event id");
- assert!(matches!(
- error,
- RadrootsSdkError::InvalidRequest { message }
- if message.contains("direct Nostr outbox publish receipt event id is invalid")
- ));
-}
-
-#[test]
-fn push_event_final_state_follows_publish_quorum_and_retryability() {
- let published = outbox_publish_receipt("a".repeat(64).as_str())
- .with_quorum_met(true)
- .with_retryable_count(1);
- assert_eq!(
- push_event_final_state(&published),
- PushOutboxEventState::Published
- );
-
- let retryable = outbox_publish_receipt("b".repeat(64).as_str()).with_retryable_count(1);
- assert_eq!(
- push_event_final_state(&retryable),
- PushOutboxEventState::PublishRetryable
- );
-
- let terminal = outbox_publish_receipt("c".repeat(64).as_str());
- assert_eq!(
- push_event_final_state(&terminal),
- PushOutboxEventState::FailedTerminal
- );
-}
-
-#[test]
-fn push_relay_outcome_mapping_covers_daemon_radrootsd_results() {
- assert_eq!(
- PushOutboxTargetOutcomeKind::from(RadrootsRelayOutcomeKind::Muted),
- PushOutboxTargetOutcomeKind::Muted
- );
- assert_eq!(
- PushOutboxTargetOutcomeKind::from(RadrootsRelayOutcomeKind::Unsupported),
- PushOutboxTargetOutcomeKind::Unsupported
- );
- assert_eq!(
- PushOutboxTargetOutcomeKind::from(RadrootsRelayOutcomeKind::PaymentRequired),
- PushOutboxTargetOutcomeKind::PaymentRequired
- );
- assert_eq!(
- PushOutboxTargetOutcomeKind::from(RadrootsRelayOutcomeKind::RelayUrlRejected),
- PushOutboxTargetOutcomeKind::TargetUriRejected
- );
- assert_eq!(
- PushOutboxTargetOutcomeKind::from(RadrootsRelayOutcomeKind::SkippedAlreadyAccepted),
- PushOutboxTargetOutcomeKind::SkippedAlreadyAccepted
- );
-}
-
-#[test]
-fn auth_policy_defaults_and_outbox_state_mappings_cover_all_public_states() {
- assert_eq!(
- SdkRelayAuthPolicy::default(),
- SdkRelayAuthPolicy::DetectOnly
- );
-
- assert_eq!(
- PushOutboxEventState::from(RadrootsOutboxEventState::DraftQueued),
- PushOutboxEventState::DraftQueued
- );
- assert_eq!(
- PushOutboxEventState::from(RadrootsOutboxEventState::Signing),
- PushOutboxEventState::Signing
- );
- assert_eq!(
- PushOutboxEventState::from(RadrootsOutboxEventState::Signed),
- PushOutboxEventState::Signed
- );
- assert_eq!(
- PushOutboxEventState::from(RadrootsOutboxEventState::Publishing),
- PushOutboxEventState::Publishing
- );
- assert_eq!(
- PushOutboxEventState::from(RadrootsOutboxEventState::Published),
- PushOutboxEventState::Published
- );
- assert_eq!(
- PushOutboxEventState::from(RadrootsOutboxEventState::SignRetryable),
- PushOutboxEventState::SignRetryable
- );
- assert_eq!(
- PushOutboxEventState::from(RadrootsOutboxEventState::PublishRetryable),
- PushOutboxEventState::PublishRetryable
- );
- assert_eq!(
- PushOutboxEventState::from(RadrootsOutboxEventState::FailedTerminal),
- PushOutboxEventState::FailedTerminal
- );
- assert_eq!(
- PushOutboxEventState::from(RadrootsOutboxEventState::Cancelled),
- PushOutboxEventState::Cancelled
- );
-
- let mut receipt = PushOutboxReceipt::default();
- receipt.push_attempted_event(push_receipt(PushOutboxEventState::Published));
- receipt.push_attempted_event(push_receipt(PushOutboxEventState::PublishRetryable));
- receipt.push_attempted_event(push_receipt(PushOutboxEventState::FailedTerminal));
- receipt.push_attempted_event(push_receipt(PushOutboxEventState::Cancelled));
- assert_eq!(receipt.attempted_events, 4);
- assert_eq!(receipt.terminal_events, 1);
- assert_eq!(receipt.published_events, 1);
- assert_eq!(receipt.retryable_events, 1);
-}
-
-#[test]
-fn sync_status_summary_conversions_preserve_all_fields() {
- let event_summary = RadrootsEventStoreStatusSummary {
- total_events: 11,
- valid_stream_events: 7,
- transport_observations: 3,
- last_event_seq: Some(9),
- last_event_updated_at_ms: Some(1_700_000_000_000),
- };
- let event_status = SyncEventStoreStatus::from(event_summary);
- assert_eq!(event_status.total_events, 11);
- assert_eq!(event_status.valid_stream_events, 7);
- assert_eq!(event_status.transport_observations, 3);
- assert_eq!(event_status.last_event_seq, Some(9));
- assert_eq!(
- event_status.last_event_updated_at_ms,
- Some(1_700_000_000_000)
- );
-
- let outbox_summary = RadrootsOutboxStatusSummary {
- total_events: 13,
- pending_events: 5,
- retryable_events: 4,
- terminal_events: 2,
- failed_terminal_events: 1,
- deferred_until_implemented_events: 10,
- ready_signed_events: 6,
- publishing_events: 8,
- last_attempt_at_ms: Some(1_700_000_000_001),
- last_error: Some("relay offline".to_owned()),
- };
- let outbox_status = SyncOutboxStatus::from(outbox_summary);
- assert_eq!(outbox_status.total_events, 13);
- assert_eq!(outbox_status.pending_events, 5);
- assert_eq!(outbox_status.retryable_events, 4);
- assert_eq!(outbox_status.terminal_events, 2);
- assert_eq!(outbox_status.failed_terminal_events, 1);
- assert_eq!(outbox_status.deferred_until_implemented_events, 10);
- assert_eq!(outbox_status.ready_signed_events, 6);
- assert_eq!(outbox_status.publishing_events, 8);
- assert_eq!(outbox_status.last_attempt_at_ms, Some(1_700_000_000_001));
- assert_eq!(outbox_status.last_error.as_deref(), Some("relay offline"));
-}
-
-#[test]
-fn push_outbox_request_builders_validate_all_bounds() {
- let request = super::PushOutboxRequest::new()
- .with_limit(2)
- .with_outbox_event_id(9)
- .republish_accepted_targets(true)
- .with_nostr_relay_url_policy(crate::NostrRelayUrlPolicy::Localhost)
- .with_auth_policy(SdkRelayAuthPolicy::DetectOnly)
- .with_claim_ttl_ms(7)
- .with_next_attempt_delay_ms(11);
- assert_eq!(request.limit, 1);
- assert_eq!(request.outbox_event_id, Some(9));
- assert!(request.republish_accepted_targets);
- request.validate().expect("valid request");
-
- assert!(matches!(
- super::PushOutboxRequest::new().with_limit(0).validate(),
- Err(RadrootsSdkError::InvalidRequest { message }) if message.contains("limit")
- ));
- assert!(matches!(
- super::PushOutboxRequest::new()
- .with_limit(super::PUSH_OUTBOX_MAX_LIMIT + 1)
- .validate(),
- Err(RadrootsSdkError::InvalidRequest { message }) if message.contains("limit")
- ));
- assert!(matches!(
- super::PushOutboxRequest::new()
- .with_outbox_event_id(0)
- .validate(),
- Err(RadrootsSdkError::InvalidRequest { message }) if message.contains("outbox event id")
- ));
- assert!(matches!(
- super::PushOutboxRequest::new()
- .with_claim_ttl_ms(0)
- .validate(),
- Err(RadrootsSdkError::InvalidRequest { message }) if message.contains("TTL")
- ));
- assert!(matches!(
- super::PushOutboxRequest::new()
- .with_next_attempt_delay_ms(0)
- .validate(),
- Err(RadrootsSdkError::InvalidRequest { message }) if message.contains("next attempt")
- ));
-}
-
-#[test]
-fn relay_outcome_kind_mapping_covers_all_transport_outcomes() {
- assert_eq!(
- PushOutboxTargetOutcomeKind::from(RadrootsRelayOutcomeKind::Accepted),
- PushOutboxTargetOutcomeKind::Accepted
- );
- assert_eq!(
- PushOutboxTargetOutcomeKind::from(RadrootsRelayOutcomeKind::DuplicateAccepted),
- PushOutboxTargetOutcomeKind::DuplicateAccepted
- );
- assert_eq!(
- PushOutboxTargetOutcomeKind::from(RadrootsRelayOutcomeKind::Blocked),
- PushOutboxTargetOutcomeKind::Blocked
- );
- assert_eq!(
- PushOutboxTargetOutcomeKind::from(RadrootsRelayOutcomeKind::RateLimited),
- PushOutboxTargetOutcomeKind::RateLimited
- );
- assert_eq!(
- PushOutboxTargetOutcomeKind::from(RadrootsRelayOutcomeKind::Invalid),
- PushOutboxTargetOutcomeKind::Invalid
- );
- assert_eq!(
- PushOutboxTargetOutcomeKind::from(RadrootsRelayOutcomeKind::PowRequired),
- PushOutboxTargetOutcomeKind::PowRequired
- );
- assert_eq!(
- PushOutboxTargetOutcomeKind::from(RadrootsRelayOutcomeKind::Restricted),
- PushOutboxTargetOutcomeKind::Restricted
- );
- assert_eq!(
- PushOutboxTargetOutcomeKind::from(RadrootsRelayOutcomeKind::AuthRequired),
- PushOutboxTargetOutcomeKind::AuthRequired
- );
- assert_eq!(
- PushOutboxTargetOutcomeKind::from(RadrootsRelayOutcomeKind::Error),
- PushOutboxTargetOutcomeKind::Error
- );
- assert_eq!(
- PushOutboxTargetOutcomeKind::from(RadrootsRelayOutcomeKind::Timeout),
- PushOutboxTargetOutcomeKind::Timeout
- );
- assert_eq!(
- PushOutboxTargetOutcomeKind::from(RadrootsRelayOutcomeKind::ConnectionFailed),
- PushOutboxTargetOutcomeKind::ConnectionFailed
- );
- assert_eq!(
- PushOutboxTargetOutcomeKind::from(RadrootsRelayOutcomeKind::Unknown),
- PushOutboxTargetOutcomeKind::Unknown
- );
-}
-
-#[test]
-fn push_outbox_outcome_kind_labels_cover_all_public_variants() {
- for (kind, label) in [
- (PushOutboxTargetOutcomeKind::Accepted, "accepted"),
- (
- PushOutboxTargetOutcomeKind::DuplicateAccepted,
- "duplicate_accepted",
- ),
- (PushOutboxTargetOutcomeKind::Blocked, "blocked"),
- (PushOutboxTargetOutcomeKind::RateLimited, "rate_limited"),
- (PushOutboxTargetOutcomeKind::Invalid, "invalid"),
- (PushOutboxTargetOutcomeKind::PowRequired, "pow_required"),
- (PushOutboxTargetOutcomeKind::Restricted, "restricted"),
- (PushOutboxTargetOutcomeKind::AuthRequired, "auth_required"),
- (PushOutboxTargetOutcomeKind::Muted, "muted"),
- (PushOutboxTargetOutcomeKind::Unsupported, "unsupported"),
- (
- PushOutboxTargetOutcomeKind::PaymentRequired,
- "payment_required",
- ),
- (PushOutboxTargetOutcomeKind::Error, "error"),
- (PushOutboxTargetOutcomeKind::Timeout, "timeout"),
- (
- PushOutboxTargetOutcomeKind::ConnectionFailed,
- "connection_failed",
- ),
- (
- PushOutboxTargetOutcomeKind::TargetUriRejected,
- "target_uri_rejected",
- ),
- (
- PushOutboxTargetOutcomeKind::SkippedAlreadyAccepted,
- "skipped_already_accepted",
- ),
- (
- PushOutboxTargetOutcomeKind::DeferredUntilImplemented,
- "deferred_until_implemented",
- ),
- (
- PushOutboxTargetOutcomeKind::DeferredUntilImplemented,
- "deferred_until_implemented",
- ),
- (PushOutboxTargetOutcomeKind::Unknown, "unknown"),
- ] {
- assert_eq!(kind.as_str(), label);
- }
-
- for (kind, label) in [
- (PushOutboxTransportOutcomeKind::Accepted, "accepted"),
- (
- PushOutboxTransportOutcomeKind::DuplicateAccepted,
- "duplicate_accepted",
- ),
- (PushOutboxTransportOutcomeKind::Delivered, "delivered"),
- (PushOutboxTransportOutcomeKind::Forwarded, "forwarded"),
- (
- PushOutboxTransportOutcomeKind::StoredByGateway,
- "stored_by_gateway",
- ),
- (PushOutboxTransportOutcomeKind::Seen, "seen"),
- (
- PushOutboxTransportOutcomeKind::DeferredUntilImplemented,
- "deferred_until_implemented",
- ),
- (PushOutboxTransportOutcomeKind::Rejected, "rejected"),
- (
- PushOutboxTransportOutcomeKind::RouteUnavailable,
- "route_unavailable",
- ),
- (
- PushOutboxTransportOutcomeKind::PayloadTooLarge,
- "payload_too_large",
- ),
- (
- PushOutboxTransportOutcomeKind::PolicyDenied,
- "policy_denied",
- ),
- (PushOutboxTransportOutcomeKind::Timeout, "timeout"),
- (
- PushOutboxTransportOutcomeKind::ConnectionFailed,
- "connection_failed",
- ),
- (
- PushOutboxTransportOutcomeKind::TransportUnavailable,
- "transport_unavailable",
- ),
- ] {
- assert_eq!(kind.as_str(), label);
- }
-}
-
-#[tokio::test]
-async fn sync_status_maps_closed_store_errors() {
- let event_store_closed = crate::RadrootsClient::builder().build().await.expect("sdk");
- event_store_closed._event_store.pool().close().await;
- assert!(matches!(
- event_store_closed
- .sync()
- .status(super::SyncStatusRequest::new())
- .await,
- Err(RadrootsSdkError::EventStore { .. })
- ));
-
- let outbox_closed = crate::RadrootsClient::builder().build().await.expect("sdk");
- outbox_closed._outbox.pool().close().await;
- assert!(matches!(
- outbox_closed
- .sync()
- .status(super::SyncStatusRequest::new())
- .await,
- Err(RadrootsSdkError::EventStore { .. })
- ));
-}
-
-#[tokio::test]
-async fn sync_runtime_reports_clock_errors_before_store_or_relay_work() {
- let sdk = crate::RadrootsClient::builder()
- .clock(crate::RadrootsSdkClock::BeforeUnixEpoch)
- .build()
- .await
- .expect("sdk");
- assert!(matches!(
- sdk.sync().status(super::SyncStatusRequest::new()).await,
- Err(RadrootsSdkError::ClockBeforeUnixEpoch)
- ));
- assert!(matches!(
- sdk.sync()
- .push_outbox_with_transport(
- &RadrootsNostrTransport::new(UnusedPublishAdapter),
- super::PushOutboxRequest::new()
- )
- .await,
- Err(RadrootsSdkError::ClockBeforeUnixEpoch)
- ));
-}
-
-#[cfg(feature = "radrootsd-execution")]
-#[tokio::test]
-async fn radrootsd_push_empty_queue_and_private_helpers_are_deterministic() {
- let sdk = crate::RadrootsClient::builder().build().await.expect("sdk");
- let adapter =
- RadrootsdPublishAdapter::new(RadrootsdPublishConfig::new("http://127.0.0.1:9/rpc"));
-
- let receipt = sdk
- .sync()
- .push_outbox_with_radrootsd_adapter(&adapter, super::PushOutboxRequest::new())
- .await
- .expect("empty transport publish push");
-
- assert_eq!(receipt.attempted_events, 0);
- assert_eq!(
- radrootsd_delivery_policy_from_remaining(
- 0,
- 0,
- None,
- &RadrootsTransportSatisfactionPolicy::no_wait()
- )
- .expect("no-wait radrootsd policy"),
- TransportPublishDeliveryPolicy::Any
- );
- assert_eq!(
- radrootsd_delivery_policy_from_remaining(
- 0,
- 0,
- None,
- &RadrootsTransportSatisfactionPolicy::all_accepted()
- )
- .expect("zero-target radrootsd policy"),
- TransportPublishDeliveryPolicy::Any
- );
- assert_eq!(
- radrootsd_delivery_policy_from_remaining(
- 2,
- 2,
- None,
- &RadrootsTransportSatisfactionPolicy::all_accepted()
- )
- .expect("all-target radrootsd policy"),
- TransportPublishDeliveryPolicy::All
- );
- assert_eq!(
- radrootsd_delivery_policy_from_remaining(
- 2,
- 1,
- None,
- &RadrootsTransportSatisfactionPolicy::any_accepted()
- )
- .expect("any-target radrootsd policy"),
- TransportPublishDeliveryPolicy::Any
- );
- let first_required = Target::new(TransportId::NOSTR, "wss://required-a.example.com")
- .expect("first required target");
- let second_required = Target::new(TransportId::NOSTR, "wss://required-b.example.com")
- .expect("second required target");
- let optional =
- Target::new(TransportId::NOSTR, "wss://optional.example.com").expect("optional target");
- let policy = RadrootsTransportSatisfactionPolicy::required_targets(
- RadrootsTransportSatisfactionClass::Accepted,
- vec![
- first_required.fingerprint().clone(),
- second_required.fingerprint().clone(),
- ],
- )
- .expect("required target policy");
- let mut first_record = delivery_target_record(1, 7, &first_required);
- first_record.status = RadrootsOutboxDeliveryTargetStatus::Accepted;
- let second_record = delivery_target_record(2, 7, &second_required);
- let mut optional_record = delivery_target_record(3, 7, &optional);
- optional_record.status = RadrootsOutboxDeliveryTargetStatus::Accepted;
- let active_targets = vec![&first_record, &second_record, &optional_record];
- let remaining = radrootsd_required_remaining_targets(&policy, &active_targets)
- .expect("required remaining targets")
- .expect("required target policy");
- assert_eq!(remaining, vec![second_required.fingerprint().clone()]);
- let protocol_remaining = remaining
- .iter()
- .map(|fingerprint| {
- TransportPublishTargetFingerprint::parse(fingerprint.as_str().to_owned())
- .expect("protocol target fingerprint")
- })
- .collect();
- assert_eq!(
- radrootsd_delivery_policy_from_remaining(2, remaining.len(), Some(&remaining), &policy)
- .expect("required target radrootsd policy"),
- TransportPublishDeliveryPolicy::RequiredTargets {
- targets: protocol_remaining
- }
- );
- assert!(matches!(
- radrootsd_delivery_policy_from_remaining(0, 1, Some(&[]), &policy),
- Err(RadrootsSdkError::InvalidRequest { message })
- if message.contains("unsatisfied required targets")
- ));
- assert_eq!(
- radrootsd_outbox_idempotency_key(7, 3, "event-id", 5),
- "radroots-sdk-outbox-7-3-event-id-5"
- );
-
- let (_sdk, claimed) = claimed_radrootsd_event("radrootsd-transport-error-receipt").await;
- let signed_event = claimed.signed_event.as_ref().expect("signed event");
- let message = radrootsd_error_message(&RadrootsdError::Http("connection refused".to_owned()));
- let receipt = radrootsd_transport_error_receipt(
- &claimed,
- signed_event,
- &TransportPublishDeliveryPolicy::All,
- message.clone(),
- )
- .expect("radrootsd transport error receipt");
- assert_eq!(receipt.event_id, *signed_event.id());
- assert_eq!(receipt.final_state, PushOutboxEventState::PublishRetryable);
- assert_eq!(receipt.retryable_count, 1);
- assert!(!receipt.quorum_met);
- assert_eq!(receipt.targets.len(), 1);
- assert_eq!(
- receipt.targets[0].outcome_kind,
- PushOutboxTargetOutcomeKind::ConnectionFailed
- );
- assert!(!receipt.targets[0].attempted);
- assert_eq!(receipt.targets[0].message.as_ref(), Some(&message));
- assert_eq!(
- radrootsd_error_message(&RadrootsdError::Http("connection refused".to_owned())),
- "radrootsd publish failed: connection refused"
- );
-}
-
-#[cfg(feature = "radrootsd-execution")]
-#[tokio::test]
-async fn radrootsd_delivery_policy_rejects_non_accepted_satisfaction_before_daemon_publish() {
- let listener = TcpListener::bind("127.0.0.1:0").expect("bind radrootsd listener");
- let endpoint = format!("http://{}/rpc", listener.local_addr().expect("addr"));
- for (index, satisfaction_policy) in [
- RadrootsTransportSatisfactionPolicy::all_forwarded(),
- RadrootsTransportSatisfactionPolicy::all_stored(),
- RadrootsTransportSatisfactionPolicy::all_seen(),
- RadrootsTransportSatisfactionPolicy::all_delivered(),
- RadrootsTransportSatisfactionPolicy::all_durable_or_observed(),
- ]
- .into_iter()
- .enumerate()
- {
- let d_tag = format!("radrootsd-non-accepted-rejected-{index}");
- let (sdk, claimed) = claimed_uningested_radrootsd_event_with_satisfaction(
- d_tag.as_str(),
- endpoint.as_str(),
- satisfaction_policy,
- )
- .await;
- let sync = sdk.sync();
- let adapter = RadrootsdPublishAdapter::new(
- RadrootsdPublishConfig::new(endpoint.clone()).with_timeout(Duration::from_millis(50)),
- );
-
- let error = push_radrootsd_claimed_outbox_event(
- &sync,
- &adapter,
- &claimed,
- 60_000,
- 1_700_000_000_000,
- )
- .await
- .expect_err("non-accepted-class radrootsd satisfaction rejected");
- assert_no_transport_publish_request(&listener);
-
- assert!(matches!(
- error,
- RadrootsSdkError::InvalidRequest { message }
- if message.contains("radrootsd publish")
- && message.contains("accepted-class satisfaction")
- ));
- let stored = sdk
- ._outbox
- .get_event(claimed.outbox_event_id)
- .await
- .expect("stored")
- .expect("stored");
- assert_eq!(stored.state, RadrootsOutboxEventState::FailedTerminal);
- assert!(stored.claim_token.is_none());
- assert!(!stored.event_store_ingested);
- assert_eq!(stored.event_store_ingested_at_ms, None);
- assert!(
- stored
- .last_error
- .as_deref()
- .expect("last error")
- .contains("accepted-class satisfaction")
- );
- }
-}
-
-#[cfg(feature = "radrootsd-execution")]
-#[test]
-fn radrootsd_outbox_target_conversion_rejects_reticulum_targets_before_behavior_loss() {
- let target = Target::new(TransportId::RETICULUM, "reticulum:local").expect("Reticulum target");
- let record = RadrootsOutboxDeliveryTargetRecord {
- delivery_target_id: 1,
- delivery_plan_id: 1,
- transport_kind: *target.kind(),
- endpoint_uri: target.uri().clone(),
- target_scope: target.scope().cloned(),
- target_label: target.label().cloned(),
- endpoint_fingerprint: target.fingerprint().clone(),
- status: RadrootsOutboxDeliveryTargetStatus::Pending,
- last_outcome_kind: None,
- attempt_count: 0,
- last_attempt_at_ms: None,
- completed_at_ms: None,
- last_error: None,
- };
-
- let error =
- transport_publish_target_from_outbox_target(&record).expect_err("Reticulum rejected");
-
- assert!(matches!(
- error,
- RadrootsSdkError::InvalidRequest { message }
- if message.contains("radrootsd execution")
- && message.contains("Nostr-only")
- && message.contains("reticulum target")
- ));
-}
-
-#[cfg(feature = "radrootsd-execution")]
-#[test]
-fn radrootsd_outbox_target_conversion_preserves_nostr_scope_and_label() {
- let target = Target::new_with_metadata(
- TransportId::NOSTR,
- "wss://relay.example.com",
- Some(TargetScope::parse("farm.local").expect("scope")),
- Some(TargetLabel::parse("Farm relay").expect("label")),
- )
- .expect("scoped Nostr target");
- let record = delivery_target_record(1, 1, &target);
-
- let converted = transport_publish_target_from_outbox_target(&record).expect("converted target");
-
- assert_eq!(converted.transport_kind, "nostr");
- assert_eq!(converted.endpoint_uri, "wss://relay.example.com");
- assert_eq!(converted.target_scope.as_deref(), Some("farm.local"));
- assert_eq!(converted.target_label.as_deref(), Some("Farm relay"));
- assert_eq!(converted.reticulum_behavior, None);
-}
-
-#[cfg(feature = "radrootsd-execution")]
-#[tokio::test]
-async fn radrootsd_push_entrypoints_report_request_clock_and_claim_errors() {
- let adapter =
- RadrootsdPublishAdapter::new(RadrootsdPublishConfig::new("http://127.0.0.1:9/rpc"));
- let sdk = crate::RadrootsClient::builder().build().await.expect("sdk");
- assert!(matches!(
- sdk.sync()
- .push_outbox_with_radrootsd_adapter(
- &adapter,
- super::PushOutboxRequest::new().with_limit(0)
- )
- .await,
- Err(RadrootsSdkError::InvalidRequest { .. })
- ));
-
- let clock_sdk = crate::RadrootsClient::builder()
- .clock(crate::RadrootsSdkClock::BeforeUnixEpoch)
- .build()
- .await
- .expect("clock sdk");
- assert!(matches!(
- clock_sdk
- .sync()
- .push_outbox_with_radrootsd_adapter(&adapter, super::PushOutboxRequest::new())
- .await,
- Err(RadrootsSdkError::ClockBeforeUnixEpoch)
- ));
-
- let closed_outbox_sdk = crate::RadrootsClient::builder()
- .build()
- .await
- .expect("closed sdk");
- closed_outbox_sdk._outbox.pool().close().await;
- assert!(matches!(
- closed_outbox_sdk
- .sync()
- .push_outbox_with_radrootsd_adapter(&adapter, super::PushOutboxRequest::new())
- .await,
- Err(RadrootsSdkError::Outbox { .. })
- ));
-}
-
-#[cfg(feature = "radrootsd-execution")]
-#[tokio::test]
-async fn radrootsd_push_reports_missing_signed_claim_before_daemon_publish() {
- let sdk = crate::RadrootsClient::builder().build().await.expect("sdk");
- let sync = sdk.sync();
- let adapter =
- RadrootsdPublishAdapter::new(RadrootsdPublishConfig::new("http://127.0.0.1:9/rpc"));
- let claimed = RadrootsOutboxClaimedEvent {
- outbox_event_id: 41,
- operation_id: 42,
- expected_event_id: "b".repeat(64),
- attempt_count: 3,
- state: RadrootsOutboxEventState::Signed,
- claim_token: "claim-token".to_owned(),
- active_delivery_plan_id: Some(1),
- draft: EventDraft::new(
- "radroots.farm.profile.v1",
- KIND_FARM,
- 1_700_000_000,
- vec![vec!["d".to_owned(), "missing-signed-event".to_owned()]],
- "{}",
- radrootsd_fixture_signer_pubkey(),
- )
- .expect("draft"),
- signed_event: None,
- delivery_targets: Vec::new(),
- };
-
- assert!(matches!(
- push_radrootsd_claimed_outbox_event(&sync, &adapter, &claimed, 60_000, 1_700_000_000_000)
- .await,
- Err(RadrootsSdkError::Transport { message })
- if message.contains("Outbox claim 41 does not contain a signed event")
- ));
-}
-
-#[cfg(feature = "radrootsd-execution")]
-#[tokio::test]
-async fn radrootsd_claim_publish_marks_retryable_transport_errors() {
- let (sdk, claimed) = claimed_radrootsd_event("radrootsd-transport-error").await;
- let sync = sdk.sync();
- let adapter =
- RadrootsdPublishAdapter::new(RadrootsdPublishConfig::new("http://127.0.0.1:9/rpc"));
- let receipt =
- push_radrootsd_claimed_outbox_event(&sync, &adapter, &claimed, 60_000, 1_700_000_000_000)
- .await
- .expect("transport error job");
-
- assert_eq!(receipt.retryable_count, 1);
- assert_eq!(receipt.final_state, PushOutboxEventState::PublishRetryable);
- assert_eq!(receipt.targets.len(), 1);
- assert_eq!(
- receipt.targets[0].outcome_kind,
- PushOutboxTargetOutcomeKind::ConnectionFailed
- );
- assert!(!receipt.targets[0].attempted);
- assert!(
- receipt.targets[0]
- .message
- .as_deref()
- .is_some_and(|message| message.contains("radrootsd publish failed"))
- );
- let stored = sdk
- ._outbox
- .get_event(claimed.outbox_event_id)
- .await
- .expect("stored")
- .expect("stored");
- assert_eq!(stored.state, RadrootsOutboxEventState::PublishRetryable);
- assert!(stored.claim_token.is_none());
-}
-
-#[cfg(feature = "radrootsd-execution")]
-#[tokio::test]
-async fn radrootsd_local_validation_errors_release_claim_before_daemon_publish() {
- let listener = TcpListener::bind("127.0.0.1:0").expect("bind radrootsd listener");
- let endpoint = format!("http://{}/rpc", listener.local_addr().expect("addr"));
- let (sdk, mut claimed) =
- claimed_uningested_radrootsd_event("radrootsd-local-validation-error", endpoint.as_str())
- .await;
- let stored_before = sdk
- ._outbox
- .get_event(claimed.outbox_event_id)
- .await
- .expect("stored before")
- .expect("stored before");
- assert!(!stored_before.event_store_ingested);
- assert_eq!(stored_before.event_store_ingested_at_ms, None);
- let reticulum_target =
- Target::new(TransportId::RETICULUM, "reticulum:local").expect("Reticulum target");
- claimed.delivery_targets[0].transport_kind = *reticulum_target.kind();
- claimed.delivery_targets[0].endpoint_uri = reticulum_target.uri().clone();
- claimed.delivery_targets[0].endpoint_fingerprint = reticulum_target.fingerprint().clone();
- let sync = sdk.sync();
- let adapter = RadrootsdPublishAdapter::new(
- RadrootsdPublishConfig::new(endpoint).with_timeout(Duration::from_millis(50)),
- );
- let error =
- push_radrootsd_claimed_outbox_event(&sync, &adapter, &claimed, 60_000, 1_700_000_000_000)
- .await
- .expect_err("local radrootsd validation error");
- assert_no_transport_publish_request(&listener);
-
- assert!(matches!(
- error,
- RadrootsSdkError::InvalidRequest { message }
- if message.contains("radrootsd execution")
- && message.contains("Nostr-only")
- && message.contains("reticulum target")
- ));
- let stored = sdk
- ._outbox
- .get_event(claimed.outbox_event_id)
- .await
- .expect("stored")
- .expect("stored");
- assert_eq!(stored.state, RadrootsOutboxEventState::FailedTerminal);
- assert!(stored.claim_token.is_none());
- assert!(!stored.event_store_ingested);
- assert_eq!(stored.event_store_ingested_at_ms, None);
- assert!(
- stored
- .last_error
- .as_deref()
- .expect("last error")
- .contains("reticulum target")
- );
-}
-
-#[cfg(feature = "radrootsd-execution")]
-#[tokio::test]
-async fn radrootsd_local_validation_failure_keeps_sibling_plan_ready_and_claimable() {
- let listener = TcpListener::bind("127.0.0.1:0").expect("bind radrootsd listener");
- let endpoint = format!("http://{}/rpc", listener.local_addr().expect("addr"));
- let sdk = crate::RadrootsClient::builder()
- .fixed_clock(crate::RadrootsSdkTimestamp::from_unix_seconds(
- 1_700_000_000,
- ))
- .build()
- .await
- .expect("sdk");
- let draft = radrootsd_frozen_draft("radrootsd-local-validation-sibling");
- let signed_event = RadrootsdFixtureSigner::new()
- .sign_frozen_draft(&draft)
- .expect("signed event");
- let first = sdk
- ._outbox
- .enqueue_signed_operation(
- RadrootsOutboxSignedOperationInput::new(
- "sync.radrootsd.unit.v1",
- draft.clone(),
- signed_event.clone(),
- RadrootsOutboxDeliveryPlanInput::new(
- "radrootsd.validation.active",
- 1,
- RadrootsTransportSatisfactionPolicy::all_accepted(),
- vec![
- Target::new(TransportId::NOSTR, "wss://active.example.com")
- .expect("active target"),
- ],
- ),
- true,
- 1_700_000_000_000,
- 1_700_000_000_000,
- )
- .with_idempotency_key("radrootsd-local-validation-sibling"),
- )
- .await
- .expect("first plan");
- let second = sdk
- ._outbox
- .enqueue_signed_operation(
- RadrootsOutboxSignedOperationInput::new(
- "sync.radrootsd.unit.v1",
- draft,
- signed_event,
- RadrootsOutboxDeliveryPlanInput::new(
- "radrootsd.validation.sibling",
- 1,
- RadrootsTransportSatisfactionPolicy::all_accepted(),
- vec![
- Target::new(TransportId::NOSTR, "wss://sibling.example.com")
- .expect("sibling target"),
- ],
- ),
- true,
- 1_700_000_000_000,
- 1_700_000_000_000,
- )
- .with_idempotency_key("radrootsd-local-validation-sibling"),
- )
- .await
- .expect("second plan");
- assert_eq!(first.outbox_event_id, second.outbox_event_id);
- let mut claimed = sdk
- ._outbox
- .claim_next_ready_signed_event(
- CLAIM_OWNER,
- "radrootsd-sibling-claim-a",
- 1_700_000_060_000,
- 1_700_000_000_000,
- )
- .await
- .expect("claim")
- .expect("claim");
- let active_plan_id = claimed.active_delivery_plan_id.expect("active plan");
- let sibling_plan_id = if first.delivery_plan_id == active_plan_id {
- second.delivery_plan_id
- } else {
- first.delivery_plan_id
- };
- let stored_before = sdk
- ._outbox
- .get_event(claimed.outbox_event_id)
- .await
- .expect("stored before")
- .expect("stored before");
- let ingested_before = stored_before.event_store_ingested;
- let ingested_at_before = stored_before.event_store_ingested_at_ms;
- let reticulum_target =
- Target::new(TransportId::RETICULUM, "reticulum:local").expect("Reticulum target");
- claimed.delivery_targets[0].transport_kind = *reticulum_target.kind();
- claimed.delivery_targets[0].endpoint_uri = reticulum_target.uri().clone();
- claimed.delivery_targets[0].endpoint_fingerprint = reticulum_target.fingerprint().clone();
- let sync = sdk.sync();
- let adapter = RadrootsdPublishAdapter::new(
- RadrootsdPublishConfig::new(endpoint).with_timeout(Duration::from_millis(50)),
- );
-
- let error =
- push_radrootsd_claimed_outbox_event(&sync, &adapter, &claimed, 60_000, 1_700_000_000_000)
- .await
- .expect_err("local radrootsd validation error");
- assert_no_transport_publish_request(&listener);
-
- assert!(matches!(
- error,
- RadrootsSdkError::InvalidRequest { message }
- if message.contains("radrootsd execution")
- && message.contains("Nostr-only")
- && message.contains("reticulum target")
- ));
- let stored = sdk
- ._outbox
- .get_event(claimed.outbox_event_id)
- .await
- .expect("stored")
- .expect("stored");
- assert_eq!(stored.state, RadrootsOutboxEventState::PublishRetryable);
- assert_eq!(stored.claim_token, None);
- assert_eq!(stored.event_store_ingested, ingested_before);
- assert_eq!(stored.event_store_ingested_at_ms, ingested_at_before);
- let targets = sdk
- ._outbox
- .delivery_targets(claimed.outbox_event_id)
- .await
- .expect("targets");
- assert!(
- targets
- .iter()
- .filter(|target| target.delivery_plan_id == active_plan_id)
- .all(|target| target.status == RadrootsOutboxDeliveryTargetStatus::FailedTerminal)
- );
- assert!(
- targets
- .iter()
- .filter(|target| target.delivery_plan_id == sibling_plan_id)
- .all(|target| target.status == RadrootsOutboxDeliveryTargetStatus::Pending)
- );
- let plans = sdk
- ._outbox
- .delivery_plans(claimed.outbox_event_id)
- .await
- .expect("plans");
- assert_eq!(
- plans
- .iter()
- .find(|plan| plan.delivery_plan_id == active_plan_id)
- .expect("active plan")
- .status,
- RadrootsOutboxDeliveryPlanStatus::FailedTerminal
- );
- assert_eq!(
- plans
- .iter()
- .find(|plan| plan.delivery_plan_id == sibling_plan_id)
- .expect("sibling plan")
- .status,
- RadrootsOutboxDeliveryPlanStatus::Queued
- );
- let sibling_claim = sdk
- ._outbox
- .claim_next_ready_signed_event(
- CLAIM_OWNER,
- "radrootsd-sibling-claim-b",
- 1_700_000_060_000,
- 1_700_000_000_000,
- )
- .await
- .expect("sibling claim")
- .expect("sibling claim");
- assert_eq!(sibling_claim.active_delivery_plan_id, Some(sibling_plan_id));
-}
-
-#[cfg(feature = "radrootsd-execution")]
-#[tokio::test]
-async fn radrootsd_completion_updates_outbox_for_success_retryable_and_terminal_receipts() {
- let cases = [
- (
- "radrootsd-complete-success",
- PushOutboxEventState::Published,
- PushOutboxEventState::Published,
- RadrootsOutboxDeliveryTargetStatus::Accepted,
- TransportPublishOutcomeKind::Accepted,
- ),
- (
- "radrootsd-complete-retryable",
- PushOutboxEventState::PublishRetryable,
- PushOutboxEventState::PublishRetryable,
- RadrootsOutboxDeliveryTargetStatus::FailedRetryable,
- TransportPublishOutcomeKind::Timeout,
- ),
- (
- "radrootsd-complete-terminal",
- PushOutboxEventState::FailedTerminal,
- PushOutboxEventState::FailedTerminal,
- RadrootsOutboxDeliveryTargetStatus::FailedTerminal,
- TransportPublishOutcomeKind::Blocked,
- ),
- (
- "radrootsd-complete-deferred",
- PushOutboxEventState::DeferredUntilImplemented,
- PushOutboxEventState::FailedTerminal,
- RadrootsOutboxDeliveryTargetStatus::DeferredUntilImplemented,
- TransportPublishOutcomeKind::DeferredUntilImplemented,
- ),
- (
- "radrootsd-complete-deferred-until-implemented",
- PushOutboxEventState::DeferredUntilImplemented,
- PushOutboxEventState::FailedTerminal,
- RadrootsOutboxDeliveryTargetStatus::DeferredUntilImplemented,
- TransportPublishOutcomeKind::DeferredUntilImplemented,
- ),
- ];
-
- for (
- d_tag,
- expected_receipt_state,
- expected_stored_state,
- expected_target_status,
- outcome_kind,
- ) in cases
- {
- let (sdk, claimed) = claimed_radrootsd_event(d_tag).await;
- let publish = radrootsd_job(
- claimed
- .signed_event
- .as_ref()
- .expect("signed event")
- .id_str(),
- outcome_kind,
- );
- assert_eq!(
- publish.event_id,
- claimed
- .signed_event
- .as_ref()
- .expect("signed event")
- .id_str()
- );
- assert_eq!(
- publish.pubkey,
- claimed
- .signed_event
- .as_ref()
- .expect("signed event")
- .pubkey()
- .to_hex()
- );
- assert_eq!(
- publish.event_kind,
- claimed.signed_event.as_ref().expect("signed event").kind()
- );
- let radrootsd_receipt =
- push_radrootsd_event_receipt(claimed.outbox_event_id, publish.clone())
- .expect("receipt");
- assert_eq!(radrootsd_receipt.final_state, expected_receipt_state);
- let sync = sdk.sync();
- complete_radrootsd_publish_attempt(&sync, &claimed, &publish, 60_000, 1_700_000_000_000)
- .await
- .expect("complete radrootsd attempt");
- let stored = sdk
- ._outbox
- .get_event(claimed.outbox_event_id)
- .await
- .expect("stored")
- .expect("stored");
- assert_eq!(
- PushOutboxEventState::from(stored.state),
- expected_stored_state
- );
- assert!(stored.claim_token.is_none());
- let targets = sdk
- ._outbox
- .delivery_targets(claimed.outbox_event_id)
- .await
- .expect("targets");
- assert_eq!(targets[0].status, expected_target_status);
- }
-}
-
-#[cfg(feature = "radrootsd-execution")]
-#[tokio::test]
-async fn radrootsd_completion_matches_duplicate_endpoint_targets_by_scope() {
- let sdk = crate::RadrootsClient::builder()
- .fixed_clock(crate::RadrootsSdkTimestamp::from_unix_seconds(
- 1_700_000_000,
- ))
- .build()
- .await
- .expect("sdk");
- let draft = radrootsd_frozen_draft("radrootsd-complete-scoped-targets");
- let signed_event = RadrootsdFixtureSigner::new()
- .sign_frozen_draft(&draft)
- .expect("signed event");
- let farm_a = Target::new_with_metadata(
- TransportId::NOSTR,
- "wss://relay.example.com",
- Some(TargetScope::parse("farm.a").expect("farm a scope")),
- Some(TargetLabel::parse("Farm A").expect("farm a label")),
- )
- .expect("farm a target");
- let farm_b = Target::new_with_metadata(
- TransportId::NOSTR,
- "wss://relay.example.com",
- Some(TargetScope::parse("farm.b").expect("farm b scope")),
- Some(TargetLabel::parse("Farm B").expect("farm b label")),
- )
- .expect("farm b target");
- let enqueue = sdk
- ._outbox
- .enqueue_signed_operation(
- RadrootsOutboxSignedOperationInput::new(
- "sync.radrootsd.unit.v1",
- draft,
- signed_event.clone(),
- RadrootsOutboxDeliveryPlanInput::new(
- "radrootsd.scoped",
- 2,
- RadrootsTransportSatisfactionPolicy::all_accepted(),
- vec![farm_a, farm_b],
- ),
- true,
- 1_700_000_000_000,
- 1_700_000_000_000,
- )
- .with_idempotency_key("radrootsd-complete-scoped-targets"),
- )
- .await
- .expect("scoped radrootsd event");
- let claimed = sdk
- ._outbox
- .claim_next_ready_signed_event(
- CLAIM_OWNER,
- "radrootsd-scoped-target-claim",
- 1_700_000_060_000,
- 1_700_000_000_000,
- )
- .await
- .expect("claim")
- .expect("claim");
- assert_eq!(claimed.outbox_event_id, enqueue.outbox_event_id);
- let mut publish = radrootsd_job(signed_event.id_str(), TransportPublishOutcomeKind::Accepted);
- publish.target_policy = TransportPublishTargetPolicy::explicit_targets(vec![
- TransportPublishTarget::nostr("wss://relay.example.com")
- .with_scope("farm.a")
- .with_label("Farm A"),
- TransportPublishTarget::nostr("wss://relay.example.com")
- .with_scope("farm.b")
- .with_label("Farm B"),
- ]);
- publish.delivery_policy = TransportPublishDeliveryPolicy::All;
- publish.delivery_satisfied = false;
- publish.status = TransportPublishJobStatus::DeliveryUnsatisfiedRetryable;
- publish.terminal = false;
- publish.target_count = 2;
- publish.acknowledged_count = 1;
- publish.retryable_count = 1;
- publish.terminal_count = 0;
- publish.completed_at_ms = None;
- publish.targets[0].target_scope = Some("farm.a".to_owned());
- publish.targets[0].target_label = Some("Farm A".to_owned());
- let mut farm_b_outcome = publish.targets[0].clone();
- farm_b_outcome.target_scope = Some("farm.b".to_owned());
- farm_b_outcome.target_label = Some("Farm B".to_owned());
- farm_b_outcome.outcome_kind = TransportPublishOutcomeKind::Timeout;
- farm_b_outcome.message = Some("daemon timeout".to_owned());
- publish.targets.push(farm_b_outcome);
-
- let sync = sdk.sync();
- complete_radrootsd_publish_attempt(&sync, &claimed, &publish, 60_000, 1_700_000_000_000)
- .await
- .expect("complete scoped radrootsd attempt");
- let targets = sdk
- ._outbox
- .delivery_targets(claimed.outbox_event_id)
- .await
- .expect("targets");
- let farm_a_target = targets
- .iter()
- .find(|target| target.target_scope.as_ref().map(|scope| scope.as_str()) == Some("farm.a"))
- .expect("farm a target");
- let farm_b_target = targets
- .iter()
- .find(|target| target.target_scope.as_ref().map(|scope| scope.as_str()) == Some("farm.b"))
- .expect("farm b target");
-
- assert_eq!(
- farm_a_target.status,
- RadrootsOutboxDeliveryTargetStatus::Accepted
- );
- assert_eq!(
- farm_b_target.status,
- RadrootsOutboxDeliveryTargetStatus::FailedRetryable
- );
-}
-
-#[cfg(feature = "radrootsd-execution")]
-#[tokio::test]
-async fn radrootsd_completion_rejects_duplicate_daemon_outcome_before_local_mutation() {
- let (sdk, claimed) = claimed_radrootsd_event("radrootsd-complete-duplicate-outcome").await;
- let mut publish = radrootsd_job(
- claimed
- .signed_event
- .as_ref()
- .expect("signed event")
- .id_str(),
- TransportPublishOutcomeKind::Accepted,
- );
- publish.targets.push(publish.targets[0].clone());
-
- let sync = sdk.sync();
- let error =
- complete_radrootsd_publish_attempt(&sync, &claimed, &publish, 60_000, 1_700_000_000_000)
- .await
- .expect_err("duplicate daemon outcome must fail closed");
-
- assert!(matches!(
- error,
- RadrootsSdkError::InvalidRequest { message }
- if message.contains("matched delivery target")
- && message.contains("more than once")
- ));
- let targets = sdk
- ._outbox
- .delivery_targets(claimed.outbox_event_id)
- .await
- .expect("targets");
- assert_ne!(
- targets[0].status,
- RadrootsOutboxDeliveryTargetStatus::Accepted
- );
-}
-
-#[cfg(feature = "radrootsd-execution")]
-#[test]
-fn push_radrootsd_event_receipt_preserves_daemon_target_metadata() {
- let mut publish = radrootsd_job(
- "a".repeat(64).as_str(),
- TransportPublishOutcomeKind::Accepted,
- );
- publish.targets[0].target_scope = Some("farm.local".to_owned());
- publish.targets[0].target_label = Some("Farm relay".to_owned());
-
- let receipt = push_radrootsd_event_receipt(1, publish).expect("receipt");
-
- assert_eq!(receipt.targets.len(), 1);
- assert_eq!(
- receipt.targets[0].target_scope.as_deref(),
- Some("farm.local")
- );
- assert_eq!(
- receipt.targets[0].target_label.as_deref(),
- Some("Farm relay")
- );
-}
-
-#[cfg(feature = "radrootsd-execution")]
-#[test]
-fn push_radrootsd_event_receipt_returns_typed_error_for_invalid_daemon_event_id() {
- let error = push_radrootsd_event_receipt(
- 1,
- radrootsd_job(
- "not-a-valid-event-id",
- TransportPublishOutcomeKind::Accepted,
- ),
- )
- .expect_err("invalid daemon event id");
- assert!(matches!(
- error,
- RadrootsSdkError::InvalidRequest { message }
- if message.contains("transport publish daemon job event id is invalid")
- ));
-}
-
-fn outbox_publish_receipt(event_id: &str) -> RadrootsOutboxPublishReceipt {
- RadrootsOutboxPublishReceipt {
- local_ingest: RadrootsOutboxEventStoreIngestReceipt {
- outbox_event_id: 1,
- event_id: event_id.to_owned(),
- already_ingested: false,
- event_store_inserted: true,
- },
- event_id: event_id.to_owned(),
- attempted_count: 0,
- accepted_count: 0,
- retryable_count: 0,
- terminal_count: 0,
- quorum: 0,
- quorum_met: false,
- target_receipts: Vec::new(),
- relay_receipts: Vec::new(),
- }
-}
-
-trait OutboxPublishReceiptFixture {
- fn with_target(self) -> Self;
- fn with_quorum_met(self, quorum_met: bool) -> Self;
- fn with_retryable_count(self, retryable_count: usize) -> Self;
-}
-
-impl OutboxPublishReceiptFixture for RadrootsOutboxPublishReceipt {
- fn with_target(mut self) -> Self {
- self.target_receipts
- .push(RadrootsOutboxPublishTargetReceipt {
- delivery_target_id: 10,
- endpoint_uri: "wss://relay.example.com".to_owned(),
- endpoint_fingerprint: Target::new_with_metadata(
- TransportId::NOSTR,
- "wss://relay.example.com",
- Some(TargetScope::parse("farm.local").expect("scope")),
- Some(TargetLabel::parse("Farm relay").expect("label")),
- )
- .expect("target")
- .fingerprint()
- .clone(),
- target_scope: Some("farm.local".to_owned()),
- target_label: Some("Farm relay".to_owned()),
- attempted: true,
- transport_status: RadrootsTransportDeliveryTargetStatus::Accepted,
- outcome: radroots_transport_nostr::RadrootsRelayOutcome::accepted(),
- });
- self
- }
-
- fn with_quorum_met(mut self, quorum_met: bool) -> Self {
- self.quorum_met = quorum_met;
- self
- }
-
- fn with_retryable_count(mut self, retryable_count: usize) -> Self {
- self.retryable_count = retryable_count;
- self
- }
-}
-
-fn push_receipt(final_state: PushOutboxEventState) -> PushOutboxEventReceipt {
- PushOutboxEventReceipt {
- event_id: EventId::parse("a".repeat(64)).expect("event id"),
- outbox_event_id: 1,
- final_state,
- attempted_count: 0,
- accepted_count: 0,
- retryable_count: 0,
- terminal_count: 0,
- quorum: 0,
- quorum_met: false,
- targets: Vec::new(),
- }
-}