commit 925d3a35f04568501a2c1b349b14adced4471576
parent 172f6a83e8877b3afebfc015000be50845321525
Author: triesap <tyson@radroots.org>
Date: Wed, 5 Aug 2026 14:36:42 +0000
refactor: adopt authored operations v2
- pin every lib dependency to the final authored-operations revision
- expose durable preparation, signing, admission, delivery, and status
- migrate product flows, facade exports, and typed transport failures
- qualify exact Git provenance, public APIs, features, packages, and coverage
Diffstat:
16 files changed, 721 insertions(+), 337 deletions(-)
diff --git a/crates/radroots/src/client.rs b/crates/radroots/src/client.rs
@@ -84,7 +84,8 @@ mod tests {
};
use radroots_transport::{
DeliveryReceipt, DeliveryRequest, Error as TransportError, FetchPage, FetchRequest,
- SinkStatus, SourceStatus, source::BoxFuture as TransportFuture,
+ SinkFailure, SinkStatus, SourceStatus, outcome::Retryability,
+ source::BoxFuture as TransportFuture,
};
struct TestSource;
@@ -110,9 +111,19 @@ mod tests {
fn deliver(
&self,
- _request: DeliveryRequest,
- ) -> TransportFuture<'_, Result<DeliveryReceipt, TransportError>> {
- Box::pin(async { Err(TransportError::UnsupportedOperation) })
+ request: DeliveryRequest,
+ ) -> TransportFuture<'_, Result<DeliveryReceipt, SinkFailure>> {
+ Box::pin(async move {
+ Err(SinkFailure::for_request(
+ &request,
+ "test_sink_unavailable",
+ Retryability::Terminal,
+ None,
+ None,
+ Vec::new(),
+ )
+ .expect("test sink failure"))
+ })
}
}
diff --git a/crates/radroots/src/event.rs b/crates/radroots/src/event.rs
@@ -1,5 +1,5 @@
//! Curated event authoring and inspection entry points.
pub use radroots_event::{
- Error, Event, EventDraft, EventId, EventKind, EventTag, SignedEvent, VerifiedEvent,
+ Error, Event, EventId, EventKind, EventTag, GenericEventDraft, SignedEvent, VerifiedEvent,
};
diff --git a/crates/radroots/tests/public_surface.rs b/crates/radroots/tests/public_surface.rs
@@ -23,7 +23,7 @@ fn root_exports_only_the_client_boundary() {
fn canonical_domain_paths_compile() {
#[allow(unused_imports)]
use radroots::{
- event::{Event, EventDraft, EventId},
+ event::{Event, EventId, GenericEventDraft},
farm::{Farm, FarmPublicLocation, Plan as FarmPlan, PrepareRequest as FarmRequest},
identity::{AccountId, PublicKey, Username},
listing::{EditV1, OperationalListing, Plan as ListingPlan},
@@ -49,8 +49,8 @@ fn explicit_transport_composition_is_inert() {
use radroots::transport::{EventSink, EventSource};
use radroots_transport::{
- DeliveryReceipt, DeliveryRequest, Error, FetchPage, FetchRequest, SinkStatus, SourceStatus,
- source::BoxFuture,
+ DeliveryReceipt, DeliveryRequest, Error, FetchPage, FetchRequest, SinkFailure, SinkStatus,
+ SourceStatus, outcome::Retryability, source::BoxFuture,
};
struct Source;
@@ -73,9 +73,19 @@ fn explicit_transport_composition_is_inert() {
fn deliver(
&self,
- _request: DeliveryRequest,
- ) -> BoxFuture<'_, Result<DeliveryReceipt, Error>> {
- Box::pin(async { Err(Error::UnsupportedOperation) })
+ request: DeliveryRequest,
+ ) -> BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> {
+ Box::pin(async move {
+ Err(SinkFailure::for_request(
+ &request,
+ "test_sink_unavailable",
+ Retryability::Terminal,
+ None,
+ None,
+ Vec::new(),
+ )
+ .expect("test sink failure"))
+ })
}
}
diff --git a/crates/sdk/src/client.rs b/crates/sdk/src/client.rs
@@ -10,6 +10,8 @@ use std::{
use radroots_signing::Signer;
use radroots_storage::Storage;
+#[cfg(all(feature = "sync", feature = "nostr", feature = "local-signing"))]
+use radroots_storage::authored_delivery::AuthoredDeliveryState;
#[cfg(feature = "memory")]
use radroots_storage::{event::SourceGeneration, memory::MemoryStorage};
use radroots_transport::{EventSink, EventSource};
@@ -53,9 +55,9 @@ struct ClientInner {
sink: Option<Arc<dyn EventSink>>,
#[cfg(feature = "sync")]
sync: Option<radroots_sync::Engine>,
- #[cfg(feature = "local-signing")]
+ #[cfg(all(feature = "sync", feature = "nostr", feature = "local-signing"))]
signing_slot: Option<crate::signing::Slot>,
- #[cfg(feature = "nostr")]
+ #[cfg(all(feature = "sync", feature = "nostr", feature = "local-signing"))]
nostr_slot: Option<crate::transport::NostrSlot>,
capability_availability: BTreeMap<CapabilityId, Availability>,
explicitly_configured_capabilities: BTreeSet<CapabilityId>,
@@ -217,7 +219,7 @@ impl PostEvent {
pub struct PublishReceipt {
event_id: String,
replay: bool,
- delivered: bool,
+ delivery_state: AuthoredDeliveryState,
}
#[cfg(all(feature = "sync", feature = "nostr", feature = "local-signing"))]
@@ -226,17 +228,24 @@ impl PublishReceipt {
pub fn event_id(&self) -> &str {
self.event_id.as_str()
}
- /// Returns whether the durable enqueue replayed an identical operation.
+ /// Returns whether preparation replayed an identical durable operation.
pub const fn is_replay(&self) -> bool {
self.replay
}
- /// Returns whether this explicit pass recorded at least one success.
+ /// Returns the complete durable delivery state after this explicit pass.
+ pub const fn delivery_state(&self) -> AuthoredDeliveryState {
+ self.delivery_state
+ }
+ /// Returns whether this explicit pass satisfied the delivery policy.
pub const fn is_delivered(&self) -> bool {
- self.delivered
+ matches!(self.delivery_state, AuthoredDeliveryState::Satisfied)
}
- /// Returns whether durable local intent remains pending delivery.
+ /// Returns whether durable local intent remains eligible for delivery.
pub const fn is_delivery_pending(&self) -> bool {
- !self.delivered
+ matches!(
+ self.delivery_state,
+ AuthoredDeliveryState::Pending | AuthoredDeliveryState::Retryable
+ )
}
}
@@ -332,13 +341,16 @@ impl ClientBuilder {
.await
.map_err(Error::storage_open_failed)?;
let storage = Arc::new(storage);
- let mut builder = Self::new()
+ let builder = Self::new()
.storage(storage.clone())
.capability_availability(CapabilityId::PERSISTENT_STORAGE, Availability::Available);
#[cfg(feature = "sync")]
{
+ let mut builder = builder;
builder.sync_storage = Some(storage);
+ Ok(builder)
}
+ #[cfg(not(feature = "sync"))]
Ok(builder)
}
@@ -472,9 +484,9 @@ impl ClientBuilder {
sink: self.sink,
#[cfg(feature = "sync")]
sync: self.sync,
- #[cfg(feature = "local-signing")]
+ #[cfg(all(feature = "sync", feature = "nostr", feature = "local-signing"))]
signing_slot: self.signing_slot,
- #[cfg(feature = "nostr")]
+ #[cfg(all(feature = "sync", feature = "nostr", feature = "local-signing"))]
nostr_slot: self.nostr_slot,
capability_availability: self.capability_availability,
explicitly_configured_capabilities: self.explicitly_configured_capabilities,
@@ -810,9 +822,10 @@ impl SocialOperations<'_> {
async fn publish(&self, authored: SocialDraft) -> Result<PublishReceipt> {
use radroots_event::contract::AuthorRole;
+ use radroots_event_codec::authoring::AuthoredEventPlan;
use radroots_signing::{Actor, actor::ActorSource, request::CancellationPolicy};
- use radroots_storage::{journal::IdempotencyKey, outbox::LeaseOwner};
- use radroots_sync::{PushRequest, policy::SyncId, push::DeliveryRunRequest};
+ use radroots_storage::journal::IdempotencyKey;
+ use radroots_sync::{PushRequest, policy::SyncId};
use radroots_transport::policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy};
let identity = self.identity()?;
@@ -823,23 +836,18 @@ impl SocialOperations<'_> {
let operation_uuid = uuid::Uuid::new_v4();
let operation_id =
SyncId::new(*operation_uuid.as_bytes()).map_err(Error::invalid_host_configuration)?;
- let created_at = now_unix_ms()? / 1_000;
- let draft = match authored {
- SocialDraft::Profile(profile) => radroots_event::EventDraft::from_authored_profile(
- &profile,
- created_at,
- identity.public_key_hex(),
- ),
- SocialDraft::Update(update) => radroots_event::EventDraft::from_authored_update(
- &update,
- created_at,
- identity.public_key_hex(),
- ),
- SocialDraft::Reply(reply) => radroots_event::EventDraft::from_authored_reply(
- &reply,
- created_at,
- identity.public_key_hex(),
- ),
+ let now = now_unix_ms()?;
+ let created_at = now / 1_000;
+ let plan = match authored {
+ SocialDraft::Profile(profile) => {
+ AuthoredEventPlan::from_profile(&profile, created_at, identity.public_key_hex())
+ }
+ SocialDraft::Update(update) => {
+ AuthoredEventPlan::from_update(&update, created_at, identity.public_key_hex())
+ }
+ SocialDraft::Reply(reply) => {
+ AuthoredEventPlan::from_nip10_reply(&reply, created_at, identity.public_key_hex())
+ }
}
.map_err(Error::invalid_host_configuration)?;
let actor = Actor::new(
@@ -853,9 +861,11 @@ impl SocialOperations<'_> {
IdempotencyKey::parse(format!("sdk-{operation_uuid}"))
.map_err(Error::invalid_host_configuration)?,
actor,
- draft,
+ plan,
targets,
SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::any()),
+ now.checked_add(30_000)
+ .ok_or_else(Error::invalid_host_configuration_without_source)?,
CancellationPolicy::PreservePublishedRequest,
)
.map_err(Error::invalid_host_configuration)?;
@@ -863,37 +873,35 @@ impl SocialOperations<'_> {
.client
.sync()?
.ok_or_else(Error::shared_operation_unavailable)?;
- let push = sync
- .sign_and_enqueue(request)
+ let preparation = sync
+ .prepare_push(request.clone())
+ .await
+ .map_err(Error::shared_operation_failed)?;
+ sync.sign_prepared(request)
.await
.map_err(Error::shared_operation_failed)?;
- let event_id = push.outbox().request().payload().event().id_hex();
+ sync.admit_signed(operation_id)
+ .await
+ .map_err(Error::shared_operation_failed)?;
+ let status = sync
+ .push_status(operation_id)
+ .await
+ .map_err(Error::shared_operation_failed)?
+ .ok_or_else(Error::shared_operation_failed_without_source)?;
+ let event_id = status
+ .artifact()
+ .signed()
+ .ok_or_else(Error::shared_operation_failed_without_source)?
+ .event()
+ .id_hex();
let delivery = sync
- .deliver_pending(
- DeliveryRunRequest::new(
- LeaseOwner::parse("radroots-sdk-host")
- .map_err(Error::invalid_host_configuration)?,
- SyncId::new(*uuid::Uuid::new_v4().as_bytes())
- .map_err(Error::invalid_host_configuration)?,
- 30_000,
- radroots_storage::outbox::OUTBOX_CLAIM_LIMIT_MAX,
- )
- .map_err(Error::invalid_host_configuration)?,
- )
+ .deliver_push(operation_id)
.await
.map_err(Error::shared_operation_failed)?;
Ok(PublishReceipt {
event_id,
- replay: push.is_replay(),
- delivered: delivery.outcomes().iter().any(|outcome| {
- outcome.as_ref().is_ok_and(|record| {
- (
- record.item_id() == push.outbox().item_id(),
- record.satisfaction()
- != radroots_storage::outbox::SatisfactionResult::Pending,
- ) == (true, true)
- })
- }),
+ replay: preparation.is_replay(),
+ delivery_state: delivery.plan().state(),
})
}
@@ -1010,7 +1018,8 @@ mod tests {
};
use radroots_transport::{
DeliveryReceipt, DeliveryRequest, Error as TransportError, FetchPage, FetchRequest,
- SinkStatus, SourceStatus, source::BoxFuture as TransportFuture,
+ SinkFailure, SinkStatus, SourceStatus, outcome::Retryability,
+ source::BoxFuture as TransportFuture,
};
use std::{
future::Future,
@@ -1042,9 +1051,19 @@ mod tests {
fn deliver(
&self,
- _request: DeliveryRequest,
- ) -> TransportFuture<'_, std::result::Result<DeliveryReceipt, TransportError>> {
- Box::pin(async { Err(TransportError::UnsupportedOperation) })
+ request: DeliveryRequest,
+ ) -> TransportFuture<'_, std::result::Result<DeliveryReceipt, SinkFailure>> {
+ Box::pin(async move {
+ Err(SinkFailure::for_request(
+ &request,
+ "test_sink_unavailable",
+ Retryability::Terminal,
+ None,
+ None,
+ Vec::new(),
+ )
+ .expect("test sink failure"))
+ })
}
}
@@ -1168,14 +1187,22 @@ mod tests {
assert_eq!(post.created_at(), 8);
assert_eq!(post.content(), "content");
- for (delivered, pending) in [(true, false), (false, true)] {
+ for (delivery_state, delivered, pending) in [
+ (AuthoredDeliveryState::Pending, false, true),
+ (AuthoredDeliveryState::Retryable, false, true),
+ (AuthoredDeliveryState::Satisfied, true, false),
+ (AuthoredDeliveryState::Exhausted, false, false),
+ (AuthoredDeliveryState::FailedTerminal, false, false),
+ (AuthoredDeliveryState::Cancelled, false, false),
+ ] {
let receipt = PublishReceipt {
event_id: "published".to_owned(),
replay: true,
- delivered,
+ delivery_state,
};
assert_eq!(receipt.event_id(), "published");
assert!(receipt.is_replay());
+ assert_eq!(receipt.delivery_state(), delivery_state);
assert_eq!(receipt.is_delivered(), delivered);
assert_eq!(receipt.is_delivery_pending(), pending);
}
diff --git a/crates/sdk/src/farm.rs b/crates/sdk/src/farm.rs
@@ -3,13 +3,15 @@
use std::{error, fmt};
use radroots_event::{
- EventDraft,
+ GenericEventDraft,
contract::AuthorRole,
envelope::kind::KIND_FARM,
farm::Farm,
id::{AddressableCoordinate, ParseError},
};
-use radroots_event_codec::{encode::EventEncodeError, encode::farm::to_wire_parts};
+use radroots_event_codec::{
+ authoring::AuthoredEventPlan, encode::EventEncodeError, encode::farm::to_wire_parts,
+};
use radroots_signing::Actor;
const FARM_PROFILE_CONTRACT_ID: &str = "radroots.farm.profile.v1";
@@ -39,7 +41,7 @@ impl PrepareRequest {
pub struct Plan {
actor: Actor,
coordinate: AddressableCoordinate,
- draft: EventDraft,
+ authored_event: AuthoredEventPlan,
}
impl Plan {
@@ -55,10 +57,10 @@ impl Plan {
&self.coordinate
}
- /// Returns the frozen canonical event draft.
+ /// Returns the immutable canonical authored event plan.
#[must_use]
- pub const fn draft(&self) -> &EventDraft {
- &self.draft
+ pub const fn authored_event(&self) -> &AuthoredEventPlan {
+ &self.authored_event
}
}
@@ -164,19 +166,22 @@ pub fn prepare(request: PrepareRequest) -> Result<Plan, PrepareError> {
request.farm.d_tag
))
.map_err(PrepareError::coordinate)?;
- let draft = EventDraft::new(
- FARM_PROFILE_CONTRACT_ID,
- parts.kind,
- request.created_at_unix,
- parts.tags,
- parts.content,
- request.actor.public_key().to_hex(),
+ let authored_event = AuthoredEventPlan::from_generic(
+ GenericEventDraft::new(
+ FARM_PROFILE_CONTRACT_ID,
+ parts.kind,
+ request.created_at_unix,
+ parts.tags,
+ parts.content,
+ request.actor.public_key().to_hex(),
+ )
+ .map_err(PrepareError::draft)?,
)
.map_err(PrepareError::draft)?;
Ok(Plan {
actor: request.actor,
coordinate,
- draft,
+ authored_event,
})
}
@@ -186,8 +191,8 @@ use radroots_signing::request::CancellationPolicy;
use radroots_storage::journal::IdempotencyKey;
#[cfg(feature = "sync")]
use radroots_sync::{
- PushReceipt,
policy::{Error as SyncError, SyncId},
+ push::PushStatus,
};
/// Explicit commit inputs for one prepared farm publication.
@@ -198,6 +203,7 @@ pub struct EnqueueRequest {
idempotency_key: IdempotencyKey,
plan: Plan,
profile: crate::transport::Profile,
+ delivery_deadline_unix_ms: u64,
cancellation: CancellationPolicy,
}
@@ -210,6 +216,7 @@ impl EnqueueRequest {
idempotency_key: IdempotencyKey,
plan: Plan,
profile: crate::transport::Profile,
+ delivery_deadline_unix_ms: u64,
cancellation: CancellationPolicy,
) -> Self {
Self {
@@ -217,6 +224,7 @@ impl EnqueueRequest {
idempotency_key,
plan,
profile,
+ delivery_deadline_unix_ms,
cancellation,
}
}
@@ -235,13 +243,12 @@ impl<'a> Operations<'a> {
Self { sync }
}
- /// Signs and atomically enqueues a prepared farm publication.
+ /// Durably prepares, signs, and locally admits a farm publication.
///
- /// Before the lower atomic enqueue commit, cancellation may leave only
- /// recoverable prepared/signed journal state. After commit, cancellation
- /// cannot claim rollback; replay with the same idempotency input returns
- /// the durable outbox record.
- pub async fn enqueue(&self, request: EnqueueRequest) -> Result<PushReceipt, SyncError> {
+ /// Once preparation commits, cancellation cannot erase the authored
+ /// intent. Replay with identical operation and idempotency inputs resumes
+ /// from its durable signing or admission phase.
+ pub async fn enqueue(&self, request: EnqueueRequest) -> Result<PushStatus, SyncError> {
let targets = request
.profile
.targets()
@@ -253,13 +260,14 @@ impl<'a> Operations<'a> {
.cloned()
.ok_or(SyncError::InvalidPushRequest)?;
self.sync
- .sign_and_enqueue(radroots_sync::PushRequest::new(
+ .submit_push(radroots_sync::PushRequest::new(
request.operation_id,
request.idempotency_key,
request.plan.actor,
- request.plan.draft,
+ request.plan.authored_event,
targets,
satisfaction,
+ request.delivery_deadline_unix_ms,
request.cancellation,
)?)
.await
@@ -306,20 +314,40 @@ mod tests {
let second = prepare(request).expect("second plan");
assert_eq!(first, second);
- assert_eq!(first.draft().contract_id(), FARM_PROFILE_CONTRACT_ID);
- assert_eq!(first.draft().kind_u32(), KIND_FARM);
- assert_eq!(first.draft().created_at_u64(), 1_800_000_000);
assert_eq!(
- first.draft().expected_pubkey(),
+ first
+ .authored_event()
+ .body()
+ .contract()
+ .contract_id()
+ .as_str(),
+ FARM_PROFILE_CONTRACT_ID
+ );
+ assert_eq!(first.authored_event().body().kind(), KIND_FARM);
+ assert_eq!(first.authored_event().created_at(), 1_800_000_000);
+ assert_eq!(
+ first.authored_event().author(),
&actor(AuthorRole::Farmer).public_key()
);
assert_eq!(
first.coordinate().as_str(),
format!("{KIND_FARM}:{PUBLIC_KEY}:AAAAAAAAAAAAAAAAAAAAAA")
);
- assert!(first.draft().content().contains("Moss Street Farm"));
- assert!(!first.draft().content().contains("latitude"));
- assert!(!first.draft().content().contains("longitude"));
+ assert!(
+ first
+ .authored_event()
+ .body()
+ .content()
+ .contains("Moss Street Farm")
+ );
+ assert!(!first.authored_event().body().content().contains("latitude"));
+ assert!(
+ !first
+ .authored_event()
+ .body()
+ .content()
+ .contains("longitude")
+ );
}
#[test]
@@ -362,9 +390,10 @@ mod tests {
policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId, SyncStorage},
};
use radroots_transport::{
- DeliveryReceipt, DeliveryRequest, Error as TransportError, EventSink, SinkStatus,
- Target, TargetSet, TransportId,
+ DeliveryReceipt, DeliveryRequest, Error as TransportError, EventSink, SinkFailure,
+ SinkStatus, Target, TargetSet, TransportId,
capability::{Availability, Maturity, SinkCapabilities},
+ outcome::Retryability,
policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy},
};
@@ -372,13 +401,18 @@ mod tests {
use super::*;
- struct FixedClock;
+ struct HostClock;
struct SequenceIds(AtomicU8);
struct NoopSink;
- impl Clock for FixedClock {
+ impl Clock for HostClock {
fn now_unix_ms(&self) -> Result<u64, Error> {
- Ok(2_000_000_000_000)
+ std::time::SystemTime::now()
+ .duration_since(std::time::UNIX_EPOCH)
+ .ok()
+ .and_then(|duration| u64::try_from(duration.as_millis()).ok())
+ .filter(|value| *value != 0)
+ .ok_or(Error::ClockUnavailable)
}
}
@@ -406,10 +440,20 @@ mod tests {
fn deliver(
&self,
- _request: DeliveryRequest,
- ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, TransportError>>
+ request: DeliveryRequest,
+ ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>>
{
- Box::pin(async { Err(TransportError::UnsupportedOperation) })
+ Box::pin(async move {
+ Err(SinkFailure::for_request(
+ &request,
+ "test_sink_unavailable",
+ Retryability::Terminal,
+ None,
+ None,
+ Vec::new(),
+ )
+ .expect("test sink failure"))
+ })
}
}
@@ -439,9 +483,9 @@ mod tests {
let capability: Arc<dyn SyncStorage> = storage.clone();
let engine = Engine::builder(
capability,
- Arc::new(FixedClock),
+ Arc::new(HostClock),
Arc::new(SequenceIds(AtomicU8::new(1))),
- DeadlinePolicy::new(1_000, 1_000, 1_000).expect("deadlines"),
+ DeadlinePolicy::new(30_000, 30_000, 30_000).expect("deadlines"),
)
.sink(Arc::new(NoopSink))
.signer(signer)
@@ -458,6 +502,7 @@ mod tests {
IdempotencyKey::parse("farm-publish-a").expect("idempotency key"),
plan,
profile,
+ 2_000_000_001_000,
CancellationPolicy::PreservePublishedRequest,
);
@@ -474,12 +519,13 @@ mod tests {
.enqueue(request.clone())
.await
.expect("committed enqueue");
- assert!(!committed.is_replay());
- assert_eq!(committed.outbox().request().target_set(), &targets);
- assert_eq!(committed.outbox().request().satisfaction(), &satisfaction);
+ assert_eq!(committed.delivery_plan().intent().target_set(), &targets);
+ assert_eq!(
+ committed.delivery_plan().intent().satisfaction(),
+ &satisfaction
+ );
let replay = operations.enqueue(request).await.expect("replay");
- assert!(replay.is_replay());
- assert_eq!(replay.outbox().item_id(), committed.outbox().item_id());
+ assert_eq!(replay, committed);
let unavailable = EnqueueRequest::new(
SyncId::new([8; 16]).expect("operation id"),
@@ -491,6 +537,7 @@ mod tests {
))
.expect("plan"),
Profile::unavailable_preview(TransportId::RETICULUM),
+ 2_000_000_001_000,
CancellationPolicy::PreservePublishedRequest,
);
assert_eq!(
diff --git a/crates/sdk/src/listing.rs b/crates/sdk/src/listing.rs
@@ -2,7 +2,8 @@
use std::{error, fmt};
-use radroots_event::{EventDraft, contract::AuthorRole, id::ClassifiedListingAddress};
+use radroots_event::{contract::AuthorRole, id::ClassifiedListingAddress};
+use radroots_event_codec::authoring::AuthoredEventPlan;
use radroots_signing::Actor;
mod operational_listing;
@@ -19,7 +20,7 @@ pub use operational_listing::{
validate_operational_listing_model,
};
use operational_listing::{
- build_operational_listing_mutation_draft, canonicalize_operational_listing_edit,
+ build_operational_listing_mutation_plan, canonicalize_operational_listing_edit,
};
/// Supported public listing mutation intent.
@@ -80,7 +81,7 @@ pub struct Plan {
action: Action,
address: ClassifiedListingAddress,
lifecycle: RadrootsOperationalListingLifecycleState,
- draft: EventDraft,
+ authored_event: AuthoredEventPlan,
}
impl Plan {
@@ -108,10 +109,10 @@ impl Plan {
self.lifecycle
}
- /// Returns the frozen canonical event draft.
+ /// Returns the immutable canonical authored event plan.
#[must_use]
- pub const fn draft(&self) -> &EventDraft {
- &self.draft
+ pub const fn authored_event(&self) -> &AuthoredEventPlan {
+ &self.authored_event
}
}
@@ -211,14 +212,15 @@ pub fn prepare(request: PrepareRequest) -> Result<Plan, PrepareError> {
Action::Update => RadrootsOperationalListingMutation::update(canonical),
};
let lifecycle = mutation.lifecycle_state().map_err(PrepareError::mutation)?;
- let draft = build_operational_listing_mutation_draft(&mutation, request.created_at_unix)
- .map_err(PrepareError::mutation)?;
+ let authored_event =
+ build_operational_listing_mutation_plan(&mutation, request.created_at_unix)
+ .map_err(PrepareError::mutation)?;
Ok(Plan {
actor: request.actor,
action: request.action,
address,
lifecycle,
- draft,
+ authored_event,
})
}
@@ -228,8 +230,8 @@ use radroots_signing::request::CancellationPolicy;
use radroots_storage::journal::IdempotencyKey;
#[cfg(feature = "sync")]
use radroots_sync::{
- PushReceipt,
policy::{Error as SyncError, SyncId},
+ push::PushStatus,
};
/// Explicit commit inputs for one prepared public listing mutation.
@@ -240,6 +242,7 @@ pub struct EnqueueRequest {
idempotency_key: IdempotencyKey,
plan: Plan,
profile: crate::transport::Profile,
+ delivery_deadline_unix_ms: u64,
cancellation: CancellationPolicy,
}
@@ -252,6 +255,7 @@ impl EnqueueRequest {
idempotency_key: IdempotencyKey,
plan: Plan,
profile: crate::transport::Profile,
+ delivery_deadline_unix_ms: u64,
cancellation: CancellationPolicy,
) -> Self {
Self {
@@ -259,6 +263,7 @@ impl EnqueueRequest {
idempotency_key,
plan,
profile,
+ delivery_deadline_unix_ms,
cancellation,
}
}
@@ -277,8 +282,8 @@ impl<'a> Operations<'a> {
Self { sync }
}
- /// Signs and atomically enqueues a prepared public listing mutation.
- pub async fn enqueue(&self, request: EnqueueRequest) -> Result<PushReceipt, SyncError> {
+ /// Durably prepares, signs, and locally admits a public listing mutation.
+ pub async fn enqueue(&self, request: EnqueueRequest) -> Result<PushStatus, SyncError> {
let targets = request
.profile
.targets()
@@ -290,13 +295,14 @@ impl<'a> Operations<'a> {
.cloned()
.ok_or(SyncError::InvalidPushRequest)?;
self.sync
- .sign_and_enqueue(radroots_sync::PushRequest::new(
+ .submit_push(radroots_sync::PushRequest::new(
request.operation_id,
request.idempotency_key,
request.plan.actor,
- request.plan.draft,
+ request.plan.authored_event,
targets,
satisfaction,
+ request.delivery_deadline_unix_ms,
request.cancellation,
)?)
.await
@@ -412,12 +418,21 @@ mod tests {
RadrootsOperationalListingLifecycleState::Published
);
assert_eq!(first.address(), replay.address());
- assert_eq!(first.draft(), replay.draft());
+ assert_eq!(first.authored_event(), replay.authored_event());
assert_eq!(first.address(), update.address());
- assert_eq!(first.draft().kind_u32(), KIND_CLASSIFIED_LISTING);
- assert_eq!(first.draft().created_at_u64(), 1_800_000_000);
- assert!(!first.draft().content().contains("latitude"));
- assert!(!first.draft().content().contains("longitude"));
+ assert_eq!(
+ first.authored_event().body().kind(),
+ KIND_CLASSIFIED_LISTING
+ );
+ assert_eq!(first.authored_event().created_at(), 1_800_000_000);
+ assert!(!first.authored_event().body().content().contains("latitude"));
+ assert!(
+ !first
+ .authored_event()
+ .body()
+ .content()
+ .contains("longitude")
+ );
}
#[test]
@@ -468,22 +483,28 @@ mod tests {
policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId, SyncStorage},
};
use radroots_transport::{
- DeliveryReceipt, DeliveryRequest, Error as TransportError, EventSink, SinkStatus,
- Target, TargetSet, TransportId,
+ DeliveryReceipt, DeliveryRequest, Error as TransportError, EventSink, SinkFailure,
+ SinkStatus, Target, TargetSet, TransportId,
capability::{Availability, Maturity, SinkCapabilities},
+ outcome::Retryability,
policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy},
};
use super::*;
use crate::{ClientBuilder, transport::Profile};
- struct FixedClock;
+ struct HostClock;
struct SequenceIds(AtomicU8);
struct NoopSink;
- impl Clock for FixedClock {
+ impl Clock for HostClock {
fn now_unix_ms(&self) -> Result<u64, Error> {
- Ok(2_000_000_000_000)
+ std::time::SystemTime::now()
+ .duration_since(std::time::UNIX_EPOCH)
+ .ok()
+ .and_then(|duration| u64::try_from(duration.as_millis()).ok())
+ .filter(|value| *value != 0)
+ .ok_or(Error::ClockUnavailable)
}
}
impl IdSource for SequenceIds {
@@ -508,10 +529,20 @@ mod tests {
}
fn deliver(
&self,
- _request: DeliveryRequest,
- ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, TransportError>>
+ request: DeliveryRequest,
+ ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>>
{
- Box::pin(async { Err(TransportError::UnsupportedOperation) })
+ Box::pin(async move {
+ Err(SinkFailure::for_request(
+ &request,
+ "test_sink_unavailable",
+ Retryability::Terminal,
+ None,
+ None,
+ Vec::new(),
+ )
+ .expect("test sink failure"))
+ })
}
}
@@ -545,9 +576,9 @@ mod tests {
let capability: Arc<dyn SyncStorage> = storage.clone();
let engine = Engine::builder(
capability,
- Arc::new(FixedClock),
+ Arc::new(HostClock),
Arc::new(SequenceIds(AtomicU8::new(1))),
- DeadlinePolicy::new(1_000, 1_000, 1_000).expect("deadlines"),
+ DeadlinePolicy::new(30_000, 30_000, 30_000).expect("deadlines"),
)
.sink(Arc::new(NoopSink))
.signer(signer)
@@ -564,6 +595,7 @@ mod tests {
IdempotencyKey::parse("listing-publish-a").expect("idempotency key"),
plan,
profile,
+ 2_000_000_001_000,
CancellationPolicy::PreservePublishedRequest,
);
@@ -579,12 +611,13 @@ mod tests {
.enqueue(request.clone())
.await
.expect("committed enqueue");
- assert!(!committed.is_replay());
- assert_eq!(committed.outbox().request().target_set(), &targets);
- assert_eq!(committed.outbox().request().satisfaction(), &satisfaction);
+ assert_eq!(committed.delivery_plan().intent().target_set(), &targets);
+ assert_eq!(
+ committed.delivery_plan().intent().satisfaction(),
+ &satisfaction
+ );
let replay = operations.enqueue(request).await.expect("replay");
- assert!(replay.is_replay());
- assert_eq!(replay.outbox().item_id(), committed.outbox().item_id());
+ assert_eq!(replay, committed);
let unavailable = EnqueueRequest::new(
SyncId::new([10; 16]).expect("operation id"),
@@ -596,6 +629,7 @@ mod tests {
))
.expect("plan"),
Profile::unavailable_preview(TransportId::RETICULUM),
+ 2_000_000_001_000,
CancellationPolicy::PreservePublishedRequest,
);
assert_eq!(
diff --git a/crates/sdk/src/listing/operational_listing/mod.rs b/crates/sdk/src/listing/operational_listing/mod.rs
@@ -17,7 +17,7 @@ pub use self::draft::{
RadrootsOperationalListingEditError, canonicalize_operational_listing_edit,
};
pub use self::model::{RadrootsOperationalListingSubtotal, RadrootsOperationalListingTotal};
-pub use self::mutation::build_operational_listing_mutation_draft;
+pub use self::mutation::build_operational_listing_mutation_plan;
pub use self::mutation::{
RadrootsOperationalListingLifecycleState, RadrootsOperationalListingMutation,
RadrootsOperationalListingMutationError,
diff --git a/crates/sdk/src/listing/operational_listing/mutation.rs b/crates/sdk/src/listing/operational_listing/mutation.rs
@@ -1,4 +1,4 @@
-//! Mutation draft preparation for Radroots Listing v1.
+//! Mutation plan preparation for Radroots Listing v1.
#![forbid(unsafe_code)]
@@ -8,10 +8,11 @@ use std::string::{String, ToString};
use radroots_event::id::ClassifiedListingAddress;
use radroots_event::{
- draft::{DraftError, EventDraft},
- envelope::kind::KIND_CLASSIFIED_LISTING,
+ GenericEventDraft, draft::DraftError, envelope::kind::KIND_CLASSIFIED_LISTING,
+};
+use radroots_event_codec::{
+ authoring::AuthoredEventPlan, encode::operational_listing::to_wire_parts_with_kind,
};
-use radroots_event_codec::encode::operational_listing::to_wire_parts_with_kind;
use crate::listing::operational_listing::draft::RadrootsOperationalListingCanonicalEdit;
@@ -45,7 +46,7 @@ pub enum RadrootsOperationalListingLifecycleState {
pub enum RadrootsOperationalListingMutationError {
UnsupportedMutation,
EncodeListing(String),
- FrozenDraft(DraftError),
+ AuthoredPlan(DraftError),
}
impl fmt::Display for RadrootsOperationalListingMutationError {
@@ -55,8 +56,8 @@ impl fmt::Display for RadrootsOperationalListingMutationError {
Self::EncodeListing(error) => {
write!(f, "failed to encode listing mutation: {error}")
}
- Self::FrozenDraft(error) => {
- write!(f, "failed to build listing mutation draft: {error}")
+ Self::AuthoredPlan(error) => {
+ write!(f, "failed to build listing authored plan: {error}")
}
}
}
@@ -125,10 +126,10 @@ impl RadrootsOperationalListingMutation {
}
}
-pub fn build_operational_listing_mutation_draft(
+pub fn build_operational_listing_mutation_plan(
mutation: &RadrootsOperationalListingMutation,
created_at: u64,
-) -> Result<EventDraft, RadrootsOperationalListingMutationError> {
+) -> Result<AuthoredEventPlan, RadrootsOperationalListingMutationError> {
let (draft, kind, contract_id) = match mutation {
RadrootsOperationalListingMutation::Publish { draft }
| RadrootsOperationalListingMutation::Update { draft } => (
@@ -144,15 +145,18 @@ pub fn build_operational_listing_mutation_draft(
let parts = to_wire_parts_with_kind(draft.listing(), kind).map_err(|error| {
RadrootsOperationalListingMutationError::EncodeListing(error.to_string())
})?;
- EventDraft::new(
- contract_id,
- parts.kind,
- created_at,
- parts.tags,
- parts.content,
- draft.seller_pubkey().to_hex(),
+ AuthoredEventPlan::from_generic(
+ GenericEventDraft::new(
+ contract_id,
+ parts.kind,
+ created_at,
+ parts.tags,
+ parts.content,
+ draft.seller_pubkey().to_hex(),
+ )
+ .map_err(RadrootsOperationalListingMutationError::AuthoredPlan)?,
)
- .map_err(RadrootsOperationalListingMutationError::FrozenDraft)
+ .map_err(RadrootsOperationalListingMutationError::AuthoredPlan)
}
#[cfg(all(test, feature = "local-signing"))]
@@ -165,7 +169,6 @@ mod tests {
use radroots_core::{Currency, Decimal, Money, Quantity, QuantityPrice, Unit};
use radroots_event::{
contract::validate_event_contract_shape,
- draft::EventDraft,
envelope::kind::KIND_CLASSIFIED_LISTING,
farm::FarmRef,
farm::resource_area::ResourceAreaRef,
@@ -177,7 +180,7 @@ mod tests {
},
wire::Nip01EventWire,
};
- use radroots_event_codec::verify::verify_nip01_event;
+ use radroots_event_codec::{authoring::AuthoredEventPlan, verify::verify_nip01_event};
use radroots_identity::PublicKey;
use crate::listing::operational_listing::draft::RadrootsOperationalListingCanonicalEdit;
@@ -186,7 +189,7 @@ mod tests {
use super::{
OPERATIONAL_LISTING_PUBLISHED_CONTRACT_ID, RadrootsOperationalListingLifecycleState,
RadrootsOperationalListingMutation, RadrootsOperationalListingMutationError,
- build_operational_listing_mutation_draft,
+ build_operational_listing_mutation_plan,
};
const SELLER: &str = FIXTURE_ALICE_PUBLIC_KEY_HEX;
@@ -199,24 +202,26 @@ mod tests {
InventoryBinId::parse(raw).expect("bin id")
}
- fn sign_draft(draft: &EventDraft) -> radroots_event::envelope::EventEnvelope {
+ fn sign_plan(plan: &AuthoredEventPlan) -> radroots_event::envelope::EventEnvelope {
let keys = Keys::parse(FIXTURE_ALICE_SECRET_KEY_HEX).expect("fixture signing key");
- assert_eq!(keys.public_key().to_hex(), draft.expected_pubkey().to_hex());
- let tags = draft
- .tags_as_vec()
- .into_iter()
- .map(|tag| Tag::parse(tag).expect("draft tag"))
+ assert_eq!(keys.public_key().to_hex(), plan.author().to_hex());
+ let tags = plan
+ .body()
+ .tags()
+ .iter()
+ .cloned()
+ .map(|tag| Tag::parse(tag).expect("plan tag"))
.collect::<Vec<_>>();
let event = EventBuilder::new(
- Kind::Custom(u16::try_from(draft.kind_u32()).expect("NIP-01 kind")),
- draft.content(),
+ Kind::Custom(u16::try_from(plan.body().kind()).expect("NIP-01 kind")),
+ plan.body().content(),
)
.tags(tags)
.allow_self_tagging()
- .custom_created_at(Timestamp::from_secs(draft.created_at_u64()))
+ .custom_created_at(Timestamp::from_secs(plan.created_at()))
.sign_with_keys(&keys)
.expect("signed listing event");
- assert_eq!(event.id.to_hex(), draft.expected_event_id_hex());
+ assert_eq!(event.id.to_hex(), plan.expected_event_id().to_hex());
let raw_json = event.as_json();
Nip01EventWire::parse_json(raw_json.as_str())
.expect("canonical event wire")
@@ -392,43 +397,46 @@ mod tests {
}
#[test]
- fn build_operational_listing_mutation_draft_maps_publish_and_update_to_published_listing() {
+ fn build_operational_listing_mutation_plan_maps_publish_and_update_to_published_listing() {
let publish = RadrootsOperationalListingMutation::publish(canonical_draft());
let update = RadrootsOperationalListingMutation::update(canonical_draft());
let publish_draft =
- build_operational_listing_mutation_draft(&publish, 1_700_000_000).expect("draft");
+ build_operational_listing_mutation_plan(&publish, 1_700_000_000).expect("draft");
let update_draft =
- build_operational_listing_mutation_draft(&update, 1_700_000_000).expect("draft");
+ build_operational_listing_mutation_plan(&update, 1_700_000_000).expect("draft");
- assert_eq!(publish_draft.kind_u32(), KIND_CLASSIFIED_LISTING);
+ assert_eq!(publish_draft.body().kind(), KIND_CLASSIFIED_LISTING);
assert_eq!(
- publish_draft.contract_id(),
+ publish_draft.body().contract().contract_id().as_str(),
OPERATIONAL_LISTING_PUBLISHED_CONTRACT_ID
);
- assert_eq!(publish_draft.expected_pubkey().to_hex(), SELLER);
- assert_eq!(publish_draft.created_at_u64(), 1_700_000_000);
- assert_eq!(publish_draft.content(), "# Coffee\n\nSingle origin coffee");
- assert_eq!(update_draft.kind_u32(), KIND_CLASSIFIED_LISTING);
+ assert_eq!(publish_draft.author().to_hex(), SELLER);
+ assert_eq!(publish_draft.created_at(), 1_700_000_000);
+ assert_eq!(
+ publish_draft.body().content(),
+ "# Coffee\n\nSingle origin coffee"
+ );
+ assert_eq!(update_draft.body().kind(), KIND_CLASSIFIED_LISTING);
assert_eq!(
- update_draft.contract_id(),
+ update_draft.body().contract().contract_id().as_str(),
OPERATIONAL_LISTING_PUBLISHED_CONTRACT_ID
);
- assert_eq!(update_draft.expected_pubkey().to_hex(), SELLER);
+ assert_eq!(update_draft.author().to_hex(), SELLER);
}
#[test]
- fn build_operational_listing_mutation_draft_rejects_save_draft() {
+ fn build_operational_listing_mutation_plan_rejects_save_draft() {
let save_draft = RadrootsOperationalListingMutation::save_draft(canonical_draft());
assert_eq!(
- build_operational_listing_mutation_draft(&save_draft, 1_700_000_000).unwrap_err(),
+ build_operational_listing_mutation_plan(&save_draft, 1_700_000_000).unwrap_err(),
RadrootsOperationalListingMutationError::UnsupportedMutation
);
}
#[test]
- fn build_operational_listing_mutation_draft_rejects_archive() {
+ fn build_operational_listing_mutation_plan_rejects_archive() {
let archive = RadrootsOperationalListingMutation::archive(
ClassifiedListingAddress::parse(format!(
"{KIND_CLASSIFIED_LISTING}:{SELLER}:AAAAAAAAAAAAAAAAAAAAAg"
@@ -437,13 +445,13 @@ mod tests {
);
assert_eq!(
- build_operational_listing_mutation_draft(&archive, 1_700_000_000).unwrap_err(),
+ build_operational_listing_mutation_plan(&archive, 1_700_000_000).unwrap_err(),
RadrootsOperationalListingMutationError::UnsupportedMutation
);
}
#[test]
- fn build_operational_listing_mutation_draft_reports_encode_errors() {
+ fn build_operational_listing_mutation_plan_reports_encode_errors() {
let mut listing = listing();
listing.resource_area = Some(ResourceAreaRef {
pubkey: SELLER.to_string(),
@@ -456,7 +464,7 @@ mod tests {
.expect("canonical listing edit");
let publish = RadrootsOperationalListingMutation::publish(draft);
- let err = build_operational_listing_mutation_draft(&publish, 1_700_000_000).unwrap_err();
+ let err = build_operational_listing_mutation_plan(&publish, 1_700_000_000).unwrap_err();
assert!(matches!(
err,
@@ -465,30 +473,27 @@ mod tests {
}
#[test]
- fn build_operational_listing_mutation_draft_event_id_is_stable_for_fixed_input() {
+ fn build_operational_listing_mutation_plan_event_id_is_stable_for_fixed_input() {
let publish = RadrootsOperationalListingMutation::publish(canonical_draft());
let first =
- build_operational_listing_mutation_draft(&publish, 1_700_000_000).expect("draft");
+ build_operational_listing_mutation_plan(&publish, 1_700_000_000).expect("draft");
let second =
- build_operational_listing_mutation_draft(&publish, 1_700_000_000).expect("draft");
+ build_operational_listing_mutation_plan(&publish, 1_700_000_000).expect("draft");
- assert_eq!(
- first.expected_event_id_hex(),
- second.expected_event_id_hex()
- );
- assert_eq!(first.expected_event_id_hex().len(), 64);
- assert_eq!(first.tags_as_vec(), second.tags_as_vec());
- assert_eq!(first.content(), second.content());
+ assert_eq!(first.expected_event_id(), second.expected_event_id());
+ assert_eq!(first.expected_event_id().to_hex().len(), 64);
+ assert_eq!(first.body().tags(), second.body().tags());
+ assert_eq!(first.body().content(), second.body().content());
}
#[test]
- fn build_operational_listing_mutation_draft_output_validates_as_operational_listing() {
+ fn build_operational_listing_mutation_plan_output_validates_as_operational_listing() {
let publish = RadrootsOperationalListingMutation::publish(canonical_draft());
let draft =
- build_operational_listing_mutation_draft(&publish, 1_700_000_000).expect("draft");
+ build_operational_listing_mutation_plan(&publish, 1_700_000_000).expect("draft");
- let signed = sign_draft(&draft);
+ let signed = sign_plan(&draft);
validate_event_contract_shape(&signed, OPERATIONAL_LISTING_PUBLISHED_CONTRACT_ID)
.expect("operational listing contract");
let verified = verify_nip01_event(signed).expect("verified listing");
diff --git a/crates/sdk/src/signing.rs b/crates/sdk/src/signing.rs
@@ -144,6 +144,7 @@ impl LocalIdentity {
self.npub.as_str()
}
+ #[cfg(any(test, all(feature = "sync", feature = "nostr")))]
pub(crate) const fn public_key(&self) -> radroots_identity::PublicKey {
self.public_key
}
@@ -305,11 +306,13 @@ mod tests {
atomic::{AtomicUsize, Ordering},
};
- use radroots_event::{EventDraft, contract::AuthorRole};
+ use radroots_event::{GenericEventDraft, contract::AuthorRole};
+ use radroots_event_codec::authoring::AuthoredEventPlan;
use radroots_identity::PublicKey;
use radroots_protocol::runtime::v1::OperationId;
use radroots_signing::{
- Actor, Error, SignReceipt, SignRequest, SignerStatus,
+ Actor, AuthoredArtifactId, Error, SignReceipt, SignRequest, SignerStatus, SigningIntentId,
+ SigningOperationId,
actor::ActorSource,
error::Kind,
request::{CancellationPolicy, SignPolicy},
@@ -320,6 +323,7 @@ mod tests {
#[cfg(feature = "nip46")]
use radroots_signing::{
capability::{CancellationSupport, SignerCapability},
+ recovery::ReplayCapability,
status::{AuthChallenge, SignProgress},
};
@@ -356,19 +360,26 @@ mod tests {
[AuthorRole::Any],
)
.expect("actor");
- let draft = EventDraft::new(
- "radroots.social.geochat.v1",
- 20_000,
- 1_700_000_000,
- Vec::new(),
- "frozen-content",
- PUBLIC_KEY,
+ let plan = AuthoredEventPlan::from_generic(
+ GenericEventDraft::new(
+ "radroots.social.geochat.v1",
+ 20_000,
+ 1_700_000_000,
+ Vec::new(),
+ "frozen-content",
+ PUBLIC_KEY,
+ )
+ .expect("draft"),
)
- .expect("draft");
+ .expect("authored plan");
SignRequest::new(
OperationId::SyncPush,
+ SigningIntentId::new(
+ SigningOperationId::new([1; 16]).expect("operation id"),
+ AuthoredArtifactId::new([2; 16]).expect("artifact id"),
+ ),
actor,
- draft,
+ plan,
SignPolicy::new(1_700_000_100, CancellationPolicy::PreservePublishedRequest)
.expect("policy"),
)
@@ -381,6 +392,7 @@ mod tests {
SignerAvailability::AwaitingAuthentication,
vec![SignerCapability::new(
SignerKind::Remote,
+ ReplayCapability::ExactReplayByRequestId,
CancellationSupport::BeforeAndAfterPublication,
true,
true,
diff --git a/crates/sdk/src/sync.rs b/crates/sdk/src/sync.rs
@@ -4,14 +4,17 @@
use std::sync::Arc;
#[cfg(feature = "sync")]
-use radroots_storage::{outbox::OutboxRecord, projection::ProjectionId};
+use radroots_storage::{authored_delivery::AuthoredDeliveryPlan, projection::ProjectionId};
#[cfg(feature = "sync")]
use radroots_sync::{
- Engine, PullReceipt, PullRequest, PushReceipt, PushRequest, SyncStatus,
+ Engine, PullReceipt, PullRequest, PushRequest, SyncStatus,
ingest::{AdmissionPolicy, IngestBatchReceipt, IngestReceipt},
policy::Error,
projection::{Reducer, RefreshReceipt, RefreshRequest},
- push::{DeliveryRunReceipt, DeliveryRunRequest},
+ push::{
+ AdmissionRunReceipt, DeliveryExecutionReceipt, PushPreparation, PushStatus,
+ SigningRunReceipt,
+ },
};
#[cfg(feature = "sync")]
use radroots_transport::source::ObservedEvent;
@@ -148,17 +151,51 @@ impl<'a> Operations<'a> {
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
+ /// Atomically persists one complete authored operation before any side effect.
+ pub async fn prepare_push(&self, request: PushRequest) -> Result<PushPreparation, Error> {
+ self.engine.prepare_push(request).await
}
- /// Runs one bounded delivery pass and retains every independent outcome.
- pub async fn deliver_pending(
+ /// Returns the complete durable state of one prepared push operation.
+ pub async fn push_status(
&self,
- request: DeliveryRunRequest,
- ) -> Result<DeliveryRunReceipt, Error> {
- self.engine.deliver_pending(request).await
+ operation_id: radroots_sync::policy::SyncId,
+ ) -> Result<Option<PushStatus>, Error> {
+ self.engine.push_status(operation_id).await
+ }
+
+ /// Runs one bounded signing phase for an exactly prepared operation.
+ pub async fn sign_prepared(&self, request: PushRequest) -> Result<SigningRunReceipt, Error> {
+ self.engine.sign_prepared(request).await
+ }
+
+ /// Runs one bounded local-admission phase for a durably signed artifact.
+ pub async fn admit_signed(
+ &self,
+ operation_id: radroots_sync::policy::SyncId,
+ ) -> Result<AdmissionRunReceipt, Error> {
+ self.engine.admit_signed(operation_id).await
+ }
+
+ /// Prepares, signs, and locally admits one authored operation.
+ ///
+ /// Preparation commits the complete recoverable intent before the signer
+ /// can be invoked. Delivery remains an explicit caller-driven phase.
+ pub async fn submit_push(&self, request: PushRequest) -> Result<PushStatus, Error> {
+ let operation_id = request.operation_id();
+ self.sign_prepared(request).await?;
+ self.admit_signed(operation_id).await?;
+ self.push_status(operation_id)
+ .await?
+ .ok_or(Error::StorageFailed)
+ }
+
+ /// Runs one bounded delivery attempt for one durable authored plan.
+ pub async fn deliver_push(
+ &self,
+ operation_id: radroots_sync::policy::SyncId,
+ ) -> Result<DeliveryExecutionReceipt, Error> {
+ self.engine.deliver_push(operation_id).await
}
/// Returns the native passive sync status without starting recovery work.
@@ -169,10 +206,10 @@ impl<'a> Operations<'a> {
/// Returns the native host scheduling decision for one durable plan.
pub fn retry_decision(
&self,
- record: &OutboxRecord,
+ plan: &AuthoredDeliveryPlan,
now_unix_ms: u64,
) -> Result<radroots_protocol::runtime::v1::SyncRetryDecision, Error> {
- self.engine.retry_decision(record, now_unix_ms)
+ self.engine.retry_decision(plan, now_unix_ms)
}
}
@@ -193,7 +230,10 @@ mod tests {
atomic::{AtomicU8, Ordering},
};
- use radroots_event::{EventDraft, SignedEvent, contract::AuthorRole, wire::Nip01EventWire};
+ use radroots_event::{
+ GenericEventDraft, SignedEvent, contract::AuthorRole, wire::Nip01EventWire,
+ };
+ use radroots_event_codec::authoring::AuthoredEventPlan;
use radroots_identity::PublicKey;
use radroots_protocol::runtime::v1::SyncCapabilityState;
use radroots_signing::{Actor, actor::ActorSource, request::CancellationPolicy};
@@ -201,7 +241,6 @@ mod tests {
event::{SourceGeneration, StoredVisibleEvent},
journal::IdempotencyKey,
memory::MemoryStorage,
- outbox::LeaseOwner,
projection::{
ProjectionGeneration, ProjectionId, RawSourceDigest, RebuildFailure, RebuildTicketId,
},
@@ -212,7 +251,6 @@ mod tests {
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,
@@ -367,22 +405,26 @@ mod tests {
[AuthorRole::Any],
)
.expect("actor");
- let draft = EventDraft::new(
- "radroots.social.geochat.v1",
- 20_000,
- 1_700_000_000,
- Vec::new(),
- "content",
- PUBLIC_KEY,
+ let plan = AuthoredEventPlan::from_generic(
+ GenericEventDraft::new(
+ "radroots.social.geochat.v1",
+ 20_000,
+ 1_700_000_000,
+ Vec::new(),
+ "content",
+ PUBLIC_KEY,
+ )
+ .expect("draft"),
)
- .expect("draft");
+ .expect("authored plan");
PushRequest::new(
SyncId::new([8; 16]).expect("operation id"),
IdempotencyKey::parse("sdk-sync-wrapper").expect("idempotency key"),
actor,
- draft,
+ plan,
TargetSet::new(vec![target()]).expect("targets"),
SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()),
+ 1_700_000_001_000,
CancellationPolicy::PreservePublishedRequest,
)
.expect("push request")
@@ -438,19 +480,26 @@ mod tests {
.expect("projection");
assert_eq!(projection.state(), RefreshState::Complete);
+ let push = push_request();
+ let operation_id = push.operation_id();
+ let preparation = operations
+ .prepare_push(push.clone())
+ .await
+ .expect("durable preparation");
+ assert!(!preparation.is_replay());
assert_eq!(
- operations.sign_and_enqueue(push_request()).await,
+ operations.sign_prepared(push).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!(
+ operations
+ .push_status(operation_id)
+ .await
+ .expect("push status")
+ .is_some()
+ );
assert_eq!(
- operations.deliver_pending(delivery).await,
+ operations.deliver_push(operation_id).await,
Err(Error::MissingSink)
);
diff --git a/crates/sdk/src/trade.rs b/crates/sdk/src/trade.rs
@@ -3,11 +3,14 @@
use std::{error, fmt};
use radroots_event::{
- EventDraft,
+ GenericEventDraft,
contract::AuthorRole,
trade::{TradeMutationEnvelopeV1, TradeProtocolError, canonical_trade_mutation_content},
};
-use radroots_event_codec::{encode::EventEncodeError, encode::trade::trade_mutation_event_build};
+use radroots_event_codec::{
+ authoring::AuthoredEventPlan, encode::EventEncodeError,
+ encode::trade::trade_mutation_event_build,
+};
use radroots_signing::Actor;
use radroots_trade::{Projection, ReductionInput, WorkflowPlan, reducer::reduce_trade_records};
@@ -32,7 +35,7 @@ impl PrepareRequest {
pub struct Plan {
actor: Actor,
workflow: WorkflowPlan,
- draft: EventDraft,
+ authored_event: AuthoredEventPlan,
}
impl Plan {
@@ -47,10 +50,10 @@ impl Plan {
&self.workflow
}
- /// Returns the frozen canonical event draft.
+ /// Returns the immutable canonical authored event plan.
#[must_use]
- pub const fn draft(&self) -> &EventDraft {
- &self.draft
+ pub const fn authored_event(&self) -> &AuthoredEventPlan {
+ &self.authored_event
}
}
@@ -172,19 +175,22 @@ pub fn prepare(request: PrepareRequest) -> Result<Plan, PrepareError> {
}
let workflow = WorkflowPlan::prepare(canonical.clone()).map_err(PrepareError::workflow)?;
let parts = trade_mutation_event_build(canonical.clone()).map_err(PrepareError::encode)?;
- let draft = EventDraft::new(
- canonical.contract_id.clone(),
- parts.kind,
- canonical.authored_at_unix_s,
- parts.tags,
- parts.content,
- canonical.author_pubkey.to_hex(),
+ let authored_event = AuthoredEventPlan::from_generic(
+ GenericEventDraft::new(
+ canonical.contract_id.clone(),
+ parts.kind,
+ canonical.authored_at_unix_s,
+ parts.tags,
+ parts.content,
+ canonical.author_pubkey.to_hex(),
+ )
+ .map_err(PrepareError::draft)?,
)
.map_err(PrepareError::draft)?;
Ok(Plan {
actor: request.actor,
workflow,
- draft,
+ authored_event,
})
}
@@ -204,8 +210,8 @@ use radroots_storage::{
};
#[cfg(feature = "sync")]
use radroots_sync::{
- PushReceipt,
policy::{Error as SyncError, SyncId},
+ push::PushStatus,
};
/// Explicit commit inputs for one prepared trade command.
@@ -216,6 +222,7 @@ pub struct EnqueueRequest {
idempotency_key: IdempotencyKey,
plan: Plan,
profile: crate::transport::Profile,
+ delivery_deadline_unix_ms: u64,
cancellation: CancellationPolicy,
}
@@ -228,6 +235,7 @@ impl EnqueueRequest {
idempotency_key: IdempotencyKey,
plan: Plan,
profile: crate::transport::Profile,
+ delivery_deadline_unix_ms: u64,
cancellation: CancellationPolicy,
) -> Self {
Self {
@@ -235,6 +243,7 @@ impl EnqueueRequest {
idempotency_key,
plan,
profile,
+ delivery_deadline_unix_ms,
cancellation,
}
}
@@ -281,9 +290,9 @@ impl<'a> Operations<'a> {
Self { storage, sync }
}
- /// Signs and atomically enqueues a prepared command. Reusing the same
+ /// Durably prepares, signs, and locally admits a command. Reusing the same
/// idempotency input is the canonical resume/replay operation.
- pub async fn enqueue(&self, request: EnqueueRequest) -> Result<PushReceipt, SyncError> {
+ pub async fn enqueue(&self, request: EnqueueRequest) -> Result<PushStatus, SyncError> {
let targets = request
.profile
.targets()
@@ -295,13 +304,14 @@ impl<'a> Operations<'a> {
.cloned()
.ok_or(SyncError::InvalidPushRequest)?;
self.sync
- .sign_and_enqueue(radroots_sync::PushRequest::new(
+ .submit_push(radroots_sync::PushRequest::new(
request.operation_id,
request.idempotency_key,
request.plan.actor,
- request.plan.draft,
+ request.plan.authored_event,
targets,
satisfaction,
+ request.delivery_deadline_unix_ms,
request.cancellation,
)?)
.await
@@ -594,10 +604,17 @@ mod tests {
);
for plan in plans {
assert_eq!(
- plan.draft().contract_id(),
+ plan.authored_event()
+ .body()
+ .contract()
+ .contract_id()
+ .as_str(),
plan.workflow().kind().contract_id()
);
- assert_eq!(plan.draft().kind_u32(), plan.workflow().kind().nostr_kind());
+ assert_eq!(
+ plan.authored_event().body().kind(),
+ plan.workflow().kind().nostr_kind()
+ );
}
}
@@ -692,9 +709,10 @@ mod tests {
policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId, SyncStorage},
};
use radroots_transport::{
- DeliveryReceipt, DeliveryRequest, Error as TransportError, EventSink, SinkStatus,
- Target, TargetSet, TransportId,
+ DeliveryReceipt, DeliveryRequest, Error as TransportError, EventSink, SinkFailure,
+ SinkStatus, Target, TargetSet, TransportId,
capability::{Availability, Maturity, SinkCapabilities},
+ outcome::Retryability,
policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy},
};
@@ -704,13 +722,18 @@ mod tests {
const BUYER_SECRET: &str =
"10c5304d6c9ae3a1a16f7860f1cc8f5e3a76225a2663b3a989a0d775919b7df5";
- struct FixedClock;
+ struct HostClock;
struct SequenceIds(AtomicU8);
struct NoopSink;
- impl Clock for FixedClock {
+ impl Clock for HostClock {
fn now_unix_ms(&self) -> Result<u64, Error> {
- Ok(2_000_000_000_000)
+ std::time::SystemTime::now()
+ .duration_since(std::time::UNIX_EPOCH)
+ .ok()
+ .and_then(|duration| u64::try_from(duration.as_millis()).ok())
+ .filter(|value| *value != 0)
+ .ok_or(Error::ClockUnavailable)
}
}
impl IdSource for SequenceIds {
@@ -735,10 +758,20 @@ mod tests {
}
fn deliver(
&self,
- _request: DeliveryRequest,
- ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, TransportError>>
+ request: DeliveryRequest,
+ ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>>
{
- Box::pin(async { Err(TransportError::UnsupportedOperation) })
+ Box::pin(async move {
+ Err(SinkFailure::for_request(
+ &request,
+ "test_sink_unavailable",
+ Retryability::Terminal,
+ None,
+ None,
+ Vec::new(),
+ )
+ .expect("test sink failure"))
+ })
}
}
@@ -772,9 +805,9 @@ mod tests {
let capability: Arc<dyn SyncStorage> = storage.clone();
let engine = Engine::builder(
capability,
- Arc::new(FixedClock),
+ Arc::new(HostClock),
Arc::new(SequenceIds(AtomicU8::new(1))),
- DeadlinePolicy::new(1_000, 1_000, 1_000).expect("deadlines"),
+ DeadlinePolicy::new(30_000, 30_000, 30_000).expect("deadlines"),
)
.sink(Arc::new(NoopSink))
.signer(signer)
@@ -815,6 +848,7 @@ mod tests {
IdempotencyKey::parse("trade-proposal-a").expect("idempotency"),
plan,
Profile::delivery(targets.clone(), satisfaction.clone()).expect("profile"),
+ 2_000_000_001_000,
CancellationPolicy::PreservePublishedRequest,
);
drop(operations.enqueue(request.clone()));
@@ -826,11 +860,9 @@ mod tests {
0
);
let committed = operations.enqueue(request.clone()).await.expect("commit");
- assert!(!committed.is_replay());
- assert_eq!(committed.outbox().request().target_set(), &targets);
+ assert_eq!(committed.delivery_plan().intent().target_set(), &targets);
let replay = operations.enqueue(request).await.expect("resume replay");
- assert!(replay.is_replay());
- assert_eq!(replay.outbox().item_id(), committed.outbox().item_id());
+ assert_eq!(replay, committed);
}
}
}
diff --git a/crates/sdk/src/transport.rs b/crates/sdk/src/transport.rs
@@ -279,12 +279,20 @@ impl radroots_transport::EventSink for NostrSlot {
request: radroots_transport::DeliveryRequest,
) -> radroots_transport::BoxFuture<
'_,
- Result<radroots_transport::DeliveryReceipt, radroots_transport::Error>,
+ Result<radroots_transport::DeliveryReceipt, radroots_transport::SinkFailure>,
> {
Box::pin(async move {
- let state = self
- .snapshot()
- .ok_or(radroots_transport::Error::UnsupportedOperation)?;
+ let Some(state) = self.snapshot() else {
+ return Err(radroots_transport::SinkFailure::for_request(
+ &request,
+ "nostr_transport_not_configured",
+ radroots_transport::outcome::Retryability::Terminal,
+ None,
+ None,
+ Vec::new(),
+ )
+ .expect("static unconfigured sink failure is valid"));
+ };
radroots_transport::EventSink::deliver(state.transport.as_ref(), request).await
})
}
@@ -728,9 +736,7 @@ mod tests {
1,
)
.expect("delivery");
- assert_eq!(
- slot.deliver(deliver).await,
- Err(Error::UnsupportedOperation)
- );
+ let failure = slot.deliver(deliver).await.expect_err("unconfigured sink");
+ assert_eq!(failure.code(), "nostr_transport_not_configured");
}
}
diff --git a/crates/sdk/tests/package_boundary.rs b/crates/sdk/tests/package_boundary.rs
@@ -311,10 +311,13 @@ fn sync_operations_only_delegate_to_the_canonical_engine() {
"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.prepare_push(request).await",
+ "self.engine.push_status(operation_id).await",
+ "self.engine.sign_prepared(request).await",
+ "self.engine.admit_signed(operation_id).await",
+ "self.engine.deliver_push(operation_id).await",
"self.engine.status(projections).await",
- "self.engine.retry_decision(record, now_unix_ms)",
+ "self.engine.retry_decision(plan, now_unix_ms)",
] {
assert!(
SYNC.contains(delegation),
@@ -341,14 +344,31 @@ fn sync_operations_only_delegate_to_the_canonical_engine() {
}
#[test]
+fn authored_submission_is_one_protocol_neutral_boundary_for_all_typed_contracts() {
+ use radroots_event_codec::authoring::REGISTRY_V7_TYPED_AUTHORING_CONTRACT_IDS;
+
+ assert_eq!(REGISTRY_V7_TYPED_AUTHORING_CONTRACT_IDS.len(), 8);
+ assert!(SYNC.contains("pub async fn submit_push"));
+ assert!(SYNC.contains("request: PushRequest"));
+ assert!(SYNC.contains("Result<PushStatus, Error>"));
+ for contract_id in REGISTRY_V7_TYPED_AUTHORING_CONTRACT_IDS {
+ assert!(
+ !SYNC.contains(contract_id),
+ "SDK submission must not branch on typed contract `{contract_id}`"
+ );
+ }
+}
+
+#[test]
fn farm_operations_preserve_pure_planning_commit_and_privacy_boundaries() {
for required in [
"encode::farm::to_wire_parts",
"AddressableCoordinate",
- "EventDraft",
+ "GenericEventDraft",
+ "AuthoredEventPlan",
"radroots_sync::PushRequest::new",
"self.sync",
- ".sign_and_enqueue(",
+ ".submit_push(",
"PrivateArtifactStore",
] {
assert!(
@@ -380,7 +400,7 @@ fn listing_operations_reuse_trade_event_sync_and_privacy_boundaries() {
"canonicalize_operational_listing_edit",
"RadrootsOperationalListingMutation",
"RadrootsOperationalListingLifecycleState",
- "build_operational_listing_mutation_draft",
+ "build_operational_listing_mutation_plan",
"radroots_sync::PushRequest::new",
"PrivateArtifactStore",
] {
diff --git a/tools/sdk_xtask_import/src/api_qualification.rs b/tools/sdk_xtask_import/src/api_qualification.rs
@@ -6,9 +6,10 @@ use serde::Deserialize;
struct Contract {
schema_version: u16,
baseline_revision: String,
+ package_version: String,
tool: String,
minimum_tool_version: String,
- release_type: String,
+ policy: String,
feature_policy: String,
packages: Vec<String>,
}
@@ -19,15 +20,32 @@ pub fn run(root: &Path) -> Result<(), String> {
verify_tool(&contract)?;
verify_revision(root, &contract.baseline_revision)?;
for package in &contract.packages {
- let args = invocation(&contract, package);
+ let args = invocation(package);
eprintln!("cargo {}", args.join(" "));
- let status = Command::new("cargo")
+ let output = Command::new("cargo")
.args(&args)
.current_dir(root)
- .status()
- .map_err(|error| format!("failed to start cargo-semver-checks: {error}"))?;
- if !status.success() {
- return Err(format!("public API qualification failed for {package}"));
+ .output()
+ .map_err(|error| format!("failed to start cargo-public-api: {error}"))?;
+ if !output.status.success() {
+ return Err(format!(
+ "public API extraction failed for {package}: {}",
+ String::from_utf8_lossy(&output.stderr).trim()
+ ));
+ }
+ let actual = String::from_utf8(output.stdout)
+ .map_err(|error| format!("public API output for {package} was not UTF-8: {error}"))?;
+ let snapshot_path = root.join(format!(
+ "docs/api/{package}-{}.txt",
+ contract.package_version
+ ));
+ let expected = fs::read_to_string(&snapshot_path)
+ .map_err(|error| format!("read {}: {error}", snapshot_path.display()))?;
+ if actual != expected {
+ return Err(format!(
+ "reviewed public API snapshot drifted for {package}: {}",
+ snapshot_path.display()
+ ));
}
}
Ok(())
@@ -42,9 +60,10 @@ fn load(root: &Path) -> Result<Contract, String> {
fn validate(contract: &Contract) -> Result<(), String> {
let unique = contract.packages.iter().collect::<BTreeSet<_>>();
- if contract.schema_version != 1
- || contract.tool != "cargo-semver-checks"
- || contract.release_type != "minor"
+ if contract.schema_version != 2
+ || contract.package_version != "0.1.0-alpha"
+ || contract.tool != "cargo-public-api"
+ || contract.policy != "reviewed-breaking-snapshot"
|| contract.feature_policy != "all"
|| contract.baseline_revision.len() < 7
|| unique.len() != 2
@@ -57,22 +76,22 @@ fn validate(contract: &Contract) -> Result<(), String> {
fn verify_tool(contract: &Contract) -> Result<(), String> {
let output = Command::new("cargo")
- .args(["semver-checks", "--version"])
+ .args(["public-api", "--version"])
.output()
- .map_err(|error| format!("failed to start cargo-semver-checks: {error}"))?;
+ .map_err(|error| format!("failed to start cargo-public-api: {error}"))?;
if !output.status.success() {
- return Err("cargo-semver-checks is required".to_owned());
+ return Err("cargo-public-api is required".to_owned());
}
let stdout = String::from_utf8_lossy(&output.stdout);
let installed = stdout
.split_whitespace()
.find_map(|value| semver::Version::parse(value).ok())
- .ok_or_else(|| format!("could not parse cargo-semver-checks version: {stdout}"))?;
+ .ok_or_else(|| format!("could not parse cargo-public-api version: {stdout}"))?;
let minimum = semver::Version::parse(&contract.minimum_tool_version)
.map_err(|error| format!("invalid minimum tool version: {error}"))?;
if installed < minimum {
return Err(format!(
- "cargo-semver-checks {} or newer is required, found {installed}",
+ "cargo-public-api {} or newer is required, found {installed}",
contract.minimum_tool_version
));
}
@@ -92,17 +111,15 @@ fn verify_revision(root: &Path, revision: &str) -> Result<(), String> {
}
}
-fn invocation(contract: &Contract, package: &str) -> Vec<String> {
+fn invocation(package: &str) -> Vec<String> {
vec![
- "semver-checks".to_owned(),
- "check-release".to_owned(),
- "--package".to_owned(),
+ "public-api".to_owned(),
+ "-p".to_owned(),
package.to_owned(),
- "--baseline-rev".to_owned(),
- contract.baseline_revision.clone(),
"--all-features".to_owned(),
- "--release-type".to_owned(),
- contract.release_type.clone(),
+ "-sss".to_owned(),
+ "--color".to_owned(),
+ "never".to_owned(),
]
}
@@ -115,6 +132,8 @@ mod tests {
let root = crate::fs::workspace_root().expect("workspace root");
let contract = load(&root).expect("contract");
validate(&contract).expect("valid contract");
- assert!(invocation(&contract, "radroots").contains(&"--all-features".to_owned()));
+ let invocation = invocation("radroots");
+ assert!(invocation.contains(&"--all-features".to_owned()));
+ assert!(invocation.contains(&"-sss".to_owned()));
}
}
diff --git a/tools/sdk_xtask_import/src/release_graph.rs b/tools/sdk_xtask_import/src/release_graph.rs
@@ -17,10 +17,23 @@ const ALL_KINDS: &[&str] = &["build", "dev", "normal"];
struct Architecture {
spec_id: String,
package_count: usize,
+ repositories: ArchitectureRepositories,
package: Vec<ArchitecturePackage>,
}
#[derive(Debug, Deserialize)]
+struct ArchitectureRepositories {
+ lib: ArchitectureRepository,
+ sdk: ArchitectureRepository,
+}
+
+#[derive(Debug, Deserialize)]
+struct ArchitectureRepository {
+ url: String,
+ packages: Vec<String>,
+}
+
+#[derive(Debug, Deserialize)]
struct ArchitecturePackage {
name: String,
publish_order_hint: usize,
@@ -181,12 +194,7 @@ fn validate(architecture: &Architecture, metadata: &CargoMetadata) -> Result<Vec
package.version
));
}
- if package.source.is_some() {
- return Err(format!(
- "release package {name} must resolve from the staged source graph; source={:?}",
- package.source
- ));
- }
+ validate_package_source(architecture, package, &workspace_members)?;
if !nodes.contains_key(package.id.as_str()) {
return Err(format!(
"release package {name} has no Cargo resolve node; manifest={}",
@@ -285,6 +293,72 @@ fn validate(architecture: &Architecture, metadata: &CargoMetadata) -> Result<Vec
Ok(order.into_iter().map(str::to_owned).collect())
}
+fn validate_package_source(
+ architecture: &Architecture,
+ package: &CargoPackage,
+ workspace_members: &BTreeSet<&str>,
+) -> Result<(), String> {
+ if workspace_members.contains(package.id.as_str()) {
+ return if package.source.is_none() {
+ Ok(())
+ } else {
+ Err(format!(
+ "workspace release package {} must resolve locally; source={:?}",
+ package.name, package.source
+ ))
+ };
+ }
+
+ let repository = [
+ &architecture.repositories.lib,
+ &architecture.repositories.sdk,
+ ]
+ .into_iter()
+ .find(|repository| repository.packages.iter().any(|name| name == &package.name))
+ .ok_or_else(|| {
+ format!(
+ "release package {} has no repository allocation",
+ package.name
+ )
+ })?;
+ let source = package.source.as_deref().ok_or_else(|| {
+ format!(
+ "external release package {} must resolve from an exact Git revision",
+ package.name
+ )
+ })?;
+ validate_exact_git_source(source, &repository.url).map_err(|reason| {
+ format!(
+ "external release package {} has invalid source {source}: {reason}",
+ package.name
+ )
+ })
+}
+
+fn validate_exact_git_source(source: &str, repository_url: &str) -> Result<(), &'static str> {
+ let canonical_url = repository_url
+ .trim_end_matches('/')
+ .trim_end_matches(".git");
+ let prefix = format!("git+{canonical_url}.git?rev=");
+ let revision_and_commit = source
+ .strip_prefix(&prefix)
+ .ok_or("repository URL or exact-revision query does not match its allocation")?;
+ let (revision, commit) = revision_and_commit
+ .split_once('#')
+ .ok_or("resolved commit fragment is absent")?;
+ if revision.len() != 40
+ || commit.len() != 40
+ || !revision.bytes().all(|byte| byte.is_ascii_hexdigit())
+ || !commit.bytes().all(|byte| byte.is_ascii_hexdigit())
+ {
+ return Err("revision and resolved commit must be full hexadecimal object IDs");
+ }
+ if revision != commit {
+ return Err("resolved commit does not equal the requested exact revision");
+ }
+ Ok(())
+}
+
fn validate_architecture(architecture: &Architecture) -> Result<BTreeMap<&str, usize>, String> {
if architecture.spec_id != SPEC_ID {
return Err(format!("architecture spec_id must be {SPEC_ID}"));
@@ -487,7 +561,7 @@ fn target_label(target: Option<&str>) -> &str {
#[cfg(test)]
mod tests {
- use super::{Architecture, CargoMetadata, validate};
+ use super::{Architecture, CargoMetadata, validate, validate_exact_git_source};
const ARCHITECTURE: &str = include_str!("../../../docs/specs/radroots_crates_release_v1.toml");
@@ -501,6 +575,30 @@ mod tests {
}
#[test]
+ fn exact_cross_repository_git_sources_require_the_requested_commit() {
+ const REVISION: &str = "691b3c844bb8824fd16b2ff4fb37b8c09bac208d";
+ let source =
+ format!("git+https://github.com/radrootslabs/lib.git?rev={REVISION}#{REVISION}");
+ validate_exact_git_source(&source, "https://github.com/radrootslabs/lib")
+ .expect("exact source");
+
+ let drifted = format!(
+ "git+https://github.com/radrootslabs/lib.git?rev={REVISION}#{}",
+ "1".repeat(40)
+ );
+ assert!(
+ validate_exact_git_source(&drifted, "https://github.com/radrootslabs/lib").is_err()
+ );
+ assert!(
+ validate_exact_git_source(
+ "registry+https://github.com/rust-lang/crates.io-index",
+ "https://github.com/radrootslabs/lib"
+ )
+ .is_err()
+ );
+ }
+
+ #[test]
fn metadata_without_resolve_graph_fails_closed() {
let architecture = toml::from_str::<Architecture>(ARCHITECTURE).expect("architecture");
let metadata: CargoMetadata =
diff --git a/tools/sdk_xtask_import/src/supply_chain_qualification.rs b/tools/sdk_xtask_import/src/supply_chain_qualification.rs
@@ -274,6 +274,20 @@ fn validate_exceptions(root: &Path, contract: &Contract) -> Result<(), String> {
if ignored != expected_ignored {
return Err("deny.toml advisory ignores differ from the governed exceptions".to_owned());
}
+ let allowed_git = deny
+ .get("sources")
+ .and_then(|value| value.get("allow-git"))
+ .and_then(toml::Value::as_array)
+ .ok_or_else(|| "deny.toml sources.allow-git is missing".to_owned())?
+ .iter()
+ .filter_map(toml::Value::as_str)
+ .collect::<BTreeSet<_>>();
+ if allowed_git != BTreeSet::from(["https://github.com/radrootslabs/lib.git"]) {
+ return Err(
+ "deny.toml Git sources must allow only the allocated Radroots lib repository"
+ .to_owned(),
+ );
+ }
let relay_source = root.join("crates/transport_nostr/src/relay.rs");
if relay_source.exists() {