lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

listing.rs (22381B)


      1 //! Side-effect-free listing planning and canonical durable enqueue operations.
      2 
      3 use std::{error, fmt};
      4 
      5 use radroots_event::{contract::AuthorRole, id::ClassifiedListingAddress};
      6 use radroots_event_codec::authoring::AuthoredEventPlan;
      7 use radroots_signing::Actor;
      8 
      9 mod operational_listing;
     10 
     11 pub use operational_listing::{
     12     BinPricingTryExt, RadrootsClassifiedListingAddressParts,
     13     RadrootsOperationalListingCanonicalEdit, RadrootsOperationalListingEditDocumentV1,
     14     RadrootsOperationalListingEditError, RadrootsOperationalListingLifecycleState,
     15     RadrootsOperationalListingMutation, RadrootsOperationalListingMutationError,
     16     RadrootsOperationalListingSubtotal, RadrootsOperationalListingTotal,
     17     RadrootsOperationalListingTradeProjection, RadrootsPublicClassifiedListingAddress,
     18     parse_classified_listing_address, parse_operational_listing_event,
     19     parse_public_classified_listing_address, validate_operational_listing_event,
     20     validate_operational_listing_model,
     21 };
     22 use operational_listing::{
     23     build_operational_listing_mutation_plan, canonicalize_operational_listing_edit,
     24 };
     25 
     26 /// Supported public listing mutation intent.
     27 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
     28 #[non_exhaustive]
     29 pub enum Action {
     30     /// Publish the first public version of a listing.
     31     Publish,
     32     /// Replace an existing addressable listing with a new public version.
     33     Update,
     34 }
     35 
     36 /// Pure inputs for one frozen public listing plan.
     37 #[derive(Clone, Debug)]
     38 pub struct PrepareRequest {
     39     actor: Actor,
     40     document: RadrootsOperationalListingEditDocumentV1,
     41     action: Action,
     42     created_at_unix: u64,
     43 }
     44 
     45 impl PrepareRequest {
     46     /// Creates a public listing publication request.
     47     #[must_use]
     48     pub const fn publish(
     49         actor: Actor,
     50         document: RadrootsOperationalListingEditDocumentV1,
     51         created_at_unix: u64,
     52     ) -> Self {
     53         Self {
     54             actor,
     55             document,
     56             action: Action::Publish,
     57             created_at_unix,
     58         }
     59     }
     60 
     61     /// Creates a public listing replacement request.
     62     #[must_use]
     63     pub const fn update(
     64         actor: Actor,
     65         document: RadrootsOperationalListingEditDocumentV1,
     66         created_at_unix: u64,
     67     ) -> Self {
     68         Self {
     69             actor,
     70             document,
     71             action: Action::Update,
     72             created_at_unix,
     73         }
     74     }
     75 }
     76 
     77 /// Frozen, replay-stable public listing mutation plan.
     78 #[derive(Clone, Debug)]
     79 pub struct Plan {
     80     actor: Actor,
     81     action: Action,
     82     address: ClassifiedListingAddress,
     83     lifecycle: RadrootsOperationalListingLifecycleState,
     84     authored_event: AuthoredEventPlan,
     85 }
     86 
     87 impl Plan {
     88     /// Returns the exact authorized actor carried into signing.
     89     #[must_use]
     90     pub const fn actor(&self) -> &Actor {
     91         &self.actor
     92     }
     93 
     94     /// Returns the requested listing mutation intent.
     95     #[must_use]
     96     pub const fn action(&self) -> Action {
     97         self.action
     98     }
     99 
    100     /// Returns the canonical addressable listing identity.
    101     #[must_use]
    102     pub const fn address(&self) -> &ClassifiedListingAddress {
    103         &self.address
    104     }
    105 
    106     /// Returns the lower-owned lifecycle state resulting from the mutation.
    107     #[must_use]
    108     pub const fn lifecycle(&self) -> RadrootsOperationalListingLifecycleState {
    109         self.lifecycle
    110     }
    111 
    112     /// Returns the immutable canonical authored event plan.
    113     #[must_use]
    114     pub const fn authored_event(&self) -> &AuthoredEventPlan {
    115         &self.authored_event
    116     }
    117 }
    118 
    119 /// Listing plan validation stage.
    120 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
    121 #[non_exhaustive]
    122 pub enum PrepareErrorKind {
    123     /// The actor does not claim the seller author role.
    124     UnauthorizedActor,
    125     /// The lower trade boundary rejected the untrusted edit document.
    126     Edit,
    127     /// The lower trade/event boundary rejected mutation preparation.
    128     Mutation,
    129 }
    130 
    131 /// One secret-safe listing planning failure retaining its lower source.
    132 pub struct PrepareError {
    133     kind: PrepareErrorKind,
    134     source: Option<Box<dyn error::Error + Send + Sync>>,
    135 }
    136 
    137 impl PrepareError {
    138     /// Returns the stable client-level planning stage.
    139     #[must_use]
    140     pub const fn kind(&self) -> PrepareErrorKind {
    141         self.kind
    142     }
    143 
    144     fn unauthorized_actor() -> Self {
    145         Self {
    146             kind: PrepareErrorKind::UnauthorizedActor,
    147             source: None,
    148         }
    149     }
    150 
    151     fn edit(source: RadrootsOperationalListingEditError) -> Self {
    152         Self::with_source(PrepareErrorKind::Edit, source)
    153     }
    154 
    155     fn mutation(source: RadrootsOperationalListingMutationError) -> Self {
    156         Self::with_source(PrepareErrorKind::Mutation, source)
    157     }
    158 
    159     fn with_source(
    160         kind: PrepareErrorKind,
    161         source: impl error::Error + Send + Sync + 'static,
    162     ) -> Self {
    163         Self {
    164             kind,
    165             source: Some(Box::new(source)),
    166         }
    167     }
    168 }
    169 
    170 impl fmt::Display for PrepareError {
    171     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    172         formatter.write_str(match self.kind {
    173             PrepareErrorKind::UnauthorizedActor => "listing actor is not authorized",
    174             PrepareErrorKind::Edit => "listing edit is invalid",
    175             PrepareErrorKind::Mutation => "listing mutation is invalid",
    176         })
    177     }
    178 }
    179 
    180 impl fmt::Debug for PrepareError {
    181     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    182         formatter
    183             .debug_struct("PrepareError")
    184             .field("kind", &self.kind)
    185             .finish_non_exhaustive()
    186     }
    187 }
    188 
    189 impl error::Error for PrepareError {
    190     fn source(&self) -> Option<&(dyn error::Error + 'static)> {
    191         self.source
    192             .as_deref()
    193             .map(|source| source as &(dyn error::Error + 'static))
    194     }
    195 }
    196 
    197 /// Validates and freezes one listing mutation without storage, signing, or network work.
    198 ///
    199 /// The canonical public model permits only coarse public locality. Exact
    200 /// coordinates and other private artifacts are deliberately absent; hosts
    201 /// persist those through `radroots_storage::private_artifact::PrivateArtifactStore`.
    202 pub fn prepare(request: PrepareRequest) -> Result<Plan, PrepareError> {
    203     if !request.actor.satisfies(AuthorRole::Seller) {
    204         return Err(PrepareError::unauthorized_actor());
    205     }
    206     let canonical =
    207         canonicalize_operational_listing_edit(request.actor.public_key(), request.document)
    208             .map_err(PrepareError::edit)?;
    209     let address = canonical.public_listing_addr().clone();
    210     let mutation = match request.action {
    211         Action::Publish => RadrootsOperationalListingMutation::publish(canonical),
    212         Action::Update => RadrootsOperationalListingMutation::update(canonical),
    213     };
    214     let lifecycle = mutation.lifecycle_state().map_err(PrepareError::mutation)?;
    215     let authored_event =
    216         build_operational_listing_mutation_plan(&mutation, request.created_at_unix)
    217             .map_err(PrepareError::mutation)?;
    218     Ok(Plan {
    219         actor: request.actor,
    220         action: request.action,
    221         address,
    222         lifecycle,
    223         authored_event,
    224     })
    225 }
    226 
    227 #[cfg(feature = "sync")]
    228 use radroots_signing::request::CancellationPolicy;
    229 #[cfg(feature = "sync")]
    230 use radroots_storage::journal::IdempotencyKey;
    231 #[cfg(feature = "sync")]
    232 use radroots_sync::{
    233     policy::{Error as SyncError, SyncId},
    234     push::PushStatus,
    235 };
    236 
    237 /// Explicit commit inputs for one prepared public listing mutation.
    238 #[cfg(feature = "sync")]
    239 #[derive(Clone, Debug)]
    240 pub struct EnqueueRequest {
    241     operation_id: SyncId,
    242     idempotency_key: IdempotencyKey,
    243     plan: Plan,
    244     profile: crate::transport::Profile,
    245     delivery_deadline_unix_ms: u64,
    246     cancellation: CancellationPolicy,
    247 }
    248 
    249 #[cfg(feature = "sync")]
    250 impl EnqueueRequest {
    251     /// Creates an enqueue request whose transport selection has no fallback.
    252     #[must_use]
    253     pub const fn new(
    254         operation_id: SyncId,
    255         idempotency_key: IdempotencyKey,
    256         plan: Plan,
    257         profile: crate::transport::Profile,
    258         delivery_deadline_unix_ms: u64,
    259         cancellation: CancellationPolicy,
    260     ) -> Self {
    261         Self {
    262             operation_id,
    263             idempotency_key,
    264             plan,
    265             profile,
    266             delivery_deadline_unix_ms,
    267             cancellation,
    268         }
    269     }
    270 }
    271 
    272 /// Borrowed listing commit operations over the canonical sync engine.
    273 #[cfg(feature = "sync")]
    274 #[derive(Clone, Copy, Debug)]
    275 pub struct Operations<'a> {
    276     sync: crate::sync::Operations<'a>,
    277 }
    278 
    279 #[cfg(feature = "sync")]
    280 impl<'a> Operations<'a> {
    281     pub(crate) const fn new(sync: crate::sync::Operations<'a>) -> Self {
    282         Self { sync }
    283     }
    284 
    285     /// Durably prepares, signs, and locally admits a public listing mutation.
    286     pub async fn enqueue(&self, request: EnqueueRequest) -> Result<PushStatus, SyncError> {
    287         let targets = request
    288             .profile
    289             .targets()
    290             .cloned()
    291             .ok_or(SyncError::InvalidPushRequest)?;
    292         let satisfaction = request
    293             .profile
    294             .satisfaction()
    295             .cloned()
    296             .ok_or(SyncError::InvalidPushRequest)?;
    297         self.sync
    298             .submit_push(radroots_sync::PushRequest::new(
    299                 request.operation_id,
    300                 request.idempotency_key,
    301                 request.plan.actor,
    302                 request.plan.authored_event,
    303                 targets,
    304                 satisfaction,
    305                 request.delivery_deadline_unix_ms,
    306                 request.cancellation,
    307             )?)
    308             .await
    309     }
    310 }
    311 
    312 #[cfg(test)]
    313 mod tests {
    314     use radroots_core::{Currency, Decimal, Money, Quantity, QuantityPrice, Unit};
    315     use radroots_event::{
    316         envelope::kind::KIND_CLASSIFIED_LISTING,
    317         farm::FarmRef,
    318         id::{DTag, InventoryBinId},
    319         listing::operational::{
    320             OperationalListing, OperationalListingAvailability, OperationalListingBin,
    321             OperationalListingDeliveryMethod, OperationalListingProduct,
    322             OperationalListingPublicLocation, OperationalListingStatus,
    323         },
    324     };
    325     use radroots_signing::actor::ActorSource;
    326 
    327     use super::*;
    328 
    329     const PUBLIC_KEY: &str = "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df";
    330 
    331     fn actor(role: AuthorRole) -> Actor {
    332         Actor::from_public_key_hex(PUBLIC_KEY, ActorSource::ExplicitPublicKey, [role])
    333             .expect("actor")
    334     }
    335 
    336     fn listing(seller: &str) -> OperationalListing {
    337         OperationalListing {
    338             d_tag: DTag::parse("AAAAAAAAAAAAAAAAAAAAAg").expect("d tag"),
    339             published_at: None,
    340             farm: FarmRef {
    341                 pubkey: seller.to_owned(),
    342                 d_tag: "AAAAAAAAAAAAAAAAAAAAAA".to_owned(),
    343             },
    344             product: OperationalListingProduct {
    345                 key: "coffee".to_owned(),
    346                 title: "Coffee".to_owned(),
    347                 category: "coffee".to_owned(),
    348                 summary: Some("Single origin coffee".to_owned()),
    349                 process: None,
    350                 lot: None,
    351                 location: None,
    352                 profile: None,
    353                 year: None,
    354             },
    355             primary_bin_id: InventoryBinId::parse("bin-1").expect("bin id"),
    356             bins: vec![OperationalListingBin {
    357                 bin_id: InventoryBinId::parse("bin-1").expect("bin id"),
    358                 quantity: Quantity::try_new(Decimal::from(1000_u32), Unit::MassG)
    359                     .expect("quantity"),
    360                 price_per_canonical_unit: QuantityPrice::try_new(
    361                     Money::try_new(Decimal::from(20_u32), Currency::USD).expect("money"),
    362                     Quantity::try_new(Decimal::from(1_u32), Unit::MassG).expect("unit"),
    363                 )
    364                 .expect("price"),
    365                 display_amount: None,
    366                 display_unit: None,
    367                 display_label: None,
    368                 display_price: None,
    369                 display_price_unit: None,
    370             }],
    371             resource_area: None,
    372             plot: None,
    373             discounts: None,
    374             inventory_available: Some(Decimal::from(5_u32)),
    375             availability: Some(OperationalListingAvailability::Status {
    376                 status: OperationalListingStatus::Active,
    377             }),
    378             delivery_method: Some(OperationalListingDeliveryMethod::Pickup),
    379             location: Some(OperationalListingPublicLocation {
    380                 primary: "Santa Cruz, California".to_owned(),
    381                 city: Some("Santa Cruz".to_owned()),
    382                 region: Some("California".to_owned()),
    383                 country: Some("US".to_owned()),
    384                 geohash: "9q8yy".to_owned(),
    385             }),
    386             images: None,
    387         }
    388     }
    389 
    390     fn document(seller: &str) -> RadrootsOperationalListingEditDocumentV1 {
    391         RadrootsOperationalListingEditDocumentV1::new(listing(seller))
    392     }
    393 
    394     #[test]
    395     fn prepare_publish_and_update_are_pure_canonical_and_privacy_bounded() {
    396         let publish_request = PrepareRequest::publish(
    397             actor(AuthorRole::Seller),
    398             document(PUBLIC_KEY),
    399             1_800_000_000,
    400         );
    401         let first = prepare(publish_request.clone()).expect("publish plan");
    402         let replay = prepare(publish_request).expect("replayed plan");
    403         let update = prepare(PrepareRequest::update(
    404             actor(AuthorRole::Seller),
    405             document(PUBLIC_KEY),
    406             1_800_000_001,
    407         ))
    408         .expect("update plan");
    409 
    410         assert_eq!(first.action(), Action::Publish);
    411         assert_eq!(update.action(), Action::Update);
    412         assert_eq!(
    413             first.lifecycle(),
    414             RadrootsOperationalListingLifecycleState::Published
    415         );
    416         assert_eq!(
    417             update.lifecycle(),
    418             RadrootsOperationalListingLifecycleState::Published
    419         );
    420         assert_eq!(first.address(), replay.address());
    421         assert_eq!(first.authored_event(), replay.authored_event());
    422         assert_eq!(first.address(), update.address());
    423         assert_eq!(
    424             first.authored_event().body().kind(),
    425             KIND_CLASSIFIED_LISTING
    426         );
    427         assert_eq!(first.authored_event().created_at(), 1_800_000_000);
    428         assert!(!first.authored_event().body().content().contains("latitude"));
    429         assert!(
    430             !first
    431                 .authored_event()
    432                 .body()
    433                 .content()
    434                 .contains("longitude")
    435         );
    436     }
    437 
    438     #[test]
    439     fn prepare_maps_authorization_and_lower_validation_classes_once() {
    440         let unauthorized = prepare(PrepareRequest::publish(
    441             actor(AuthorRole::Buyer),
    442             document(PUBLIC_KEY),
    443             1_800_000_000,
    444         ))
    445         .expect_err("unauthorized");
    446         assert_eq!(unauthorized.kind(), PrepareErrorKind::UnauthorizedActor);
    447         assert!(std::error::Error::source(&unauthorized).is_none());
    448 
    449         let mut invalid = listing(PUBLIC_KEY);
    450         invalid.product.title.clear();
    451         let edit = prepare(PrepareRequest::publish(
    452             actor(AuthorRole::Seller),
    453             RadrootsOperationalListingEditDocumentV1::new(invalid),
    454             1_800_000_000,
    455         ))
    456         .expect_err("invalid listing");
    457         assert_eq!(edit.kind(), PrepareErrorKind::Edit);
    458         assert!(std::error::Error::source(&edit).is_some());
    459         assert_eq!(edit.to_string(), "listing edit is invalid");
    460         assert!(!format!("{edit:?}").contains("title"));
    461 
    462         let mismatch = prepare(PrepareRequest::update(
    463             actor(AuthorRole::Seller),
    464             document("8f"),
    465             1_800_000_000,
    466         ))
    467         .expect_err("invalid seller identity");
    468         assert_eq!(mismatch.kind(), PrepareErrorKind::Edit);
    469     }
    470 
    471     #[cfg(all(feature = "sync", feature = "memory", feature = "local-signing"))]
    472     mod enqueue {
    473         use std::sync::{
    474             Arc,
    475             atomic::{AtomicU8, Ordering},
    476         };
    477 
    478         use radroots_storage::{
    479             Outbox, event::SourceGeneration, journal::IdempotencyKey, memory::MemoryStorage,
    480         };
    481         use radroots_sync::{
    482             Engine,
    483             policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId, SyncStorage},
    484         };
    485         use radroots_transport::{
    486             DeliveryReceipt, DeliveryRequest, Error as TransportError, EventSink, SinkFailure,
    487             SinkStatus, Target, TargetSet, TransportId,
    488             capability::{Availability, Maturity, SinkCapabilities},
    489             outcome::Retryability,
    490             policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy},
    491         };
    492 
    493         use super::*;
    494         use crate::{ClientBuilder, transport::Profile};
    495 
    496         struct HostClock;
    497         struct SequenceIds(AtomicU8);
    498         struct NoopSink;
    499 
    500         impl Clock for HostClock {
    501             fn now_unix_ms(&self) -> Result<u64, Error> {
    502                 std::time::SystemTime::now()
    503                     .duration_since(std::time::UNIX_EPOCH)
    504                     .ok()
    505                     .and_then(|duration| u64::try_from(duration.as_millis()).ok())
    506                     .filter(|value| *value != 0)
    507                     .ok_or(Error::ClockUnavailable)
    508             }
    509         }
    510         impl IdSource for SequenceIds {
    511             fn next_id(&self, _operation: OperationKind) -> Result<SyncId, Error> {
    512                 SyncId::new([self.0.fetch_add(1, Ordering::Relaxed); 16])
    513             }
    514         }
    515         impl EventSink for NoopSink {
    516             fn status(
    517                 &self,
    518             ) -> radroots_transport::BoxFuture<'_, Result<SinkStatus, TransportError>> {
    519                 Box::pin(async {
    520                     Ok(SinkStatus::new(
    521                         TransportId::NOSTR,
    522                         true,
    523                         Maturity::Stable,
    524                         Availability::Available,
    525                         SinkCapabilities::DELIVER,
    526                         "ready",
    527                     ))
    528                 })
    529             }
    530             fn deliver(
    531                 &self,
    532                 request: DeliveryRequest,
    533             ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>>
    534             {
    535                 Box::pin(async move {
    536                     Err(SinkFailure::for_request(
    537                         &request,
    538                         "test_sink_unavailable",
    539                         Retryability::Terminal,
    540                         None,
    541                         None,
    542                         Vec::new(),
    543                     )
    544                     .expect("test sink failure"))
    545                 })
    546             }
    547         }
    548 
    549         #[tokio::test]
    550         async fn enqueue_preserves_commit_cancellation_idempotency_and_transport_outcomes() {
    551             let storage = Arc::new(MemoryStorage::new(
    552                 SourceGeneration::new([5; 32]).expect("generation"),
    553             ));
    554             let signer =
    555                 Arc::new(radroots_nostr::signing::LocalSigner::generate().expect("local signer"));
    556             let seller = Actor::new(
    557                 signer.public_key(),
    558                 ActorSource::ExplicitPublicKey,
    559                 [AuthorRole::Seller],
    560             )
    561             .expect("actor");
    562             let plan = prepare(PrepareRequest::publish(
    563                 seller,
    564                 document(&signer.public_key().to_hex()),
    565                 1_800_000_000,
    566             ))
    567             .expect("plan");
    568             let targets = TargetSet::new(vec![
    569                 Target::nostr_relay("wss://listing.example").expect("target"),
    570             ])
    571             .expect("targets");
    572             let satisfaction =
    573                 SatisfactionPolicy::new(SatisfactionClass::Delivered, TargetPolicy::all());
    574             let profile = Profile::delivery(targets.clone(), satisfaction.clone())
    575                 .expect("transport profile");
    576             let capability: Arc<dyn SyncStorage> = storage.clone();
    577             let engine = Engine::builder(
    578                 capability,
    579                 Arc::new(HostClock),
    580                 Arc::new(SequenceIds(AtomicU8::new(1))),
    581                 DeadlinePolicy::new(30_000, 30_000, 30_000).expect("deadlines"),
    582             )
    583             .sink(Arc::new(NoopSink))
    584             .signer(signer)
    585             .build()
    586             .expect("engine");
    587             let client = ClientBuilder::new()
    588                 .storage(storage.clone())
    589                 .sync_engine(engine)
    590                 .build()
    591                 .expect("client");
    592             let operations = client.listing().expect("open").expect("listing operations");
    593             let request = EnqueueRequest::new(
    594                 SyncId::new([9; 16]).expect("operation id"),
    595                 IdempotencyKey::parse("listing-publish-a").expect("idempotency key"),
    596                 plan,
    597                 profile,
    598                 2_000_000_001_000,
    599                 CancellationPolicy::PreservePublishedRequest,
    600             );
    601 
    602             drop(operations.enqueue(request.clone()));
    603             assert_eq!(
    604                 Outbox::status(storage.as_ref())
    605                     .await
    606                     .expect("outbox status")
    607                     .pending,
    608                 0
    609             );
    610             let committed = operations
    611                 .enqueue(request.clone())
    612                 .await
    613                 .expect("committed enqueue");
    614             assert_eq!(committed.delivery_plan().intent().target_set(), &targets);
    615             assert_eq!(
    616                 committed.delivery_plan().intent().satisfaction(),
    617                 &satisfaction
    618             );
    619             let replay = operations.enqueue(request).await.expect("replay");
    620             assert_eq!(replay, committed);
    621 
    622             let unavailable = EnqueueRequest::new(
    623                 SyncId::new([10; 16]).expect("operation id"),
    624                 IdempotencyKey::parse("listing-update-preview").expect("idempotency key"),
    625                 prepare(PrepareRequest::update(
    626                     actor(AuthorRole::Seller),
    627                     document(PUBLIC_KEY),
    628                     1_800_000_001,
    629                 ))
    630                 .expect("plan"),
    631                 Profile::unavailable_preview(TransportId::RETICULUM),
    632                 2_000_000_001_000,
    633                 CancellationPolicy::PreservePublishedRequest,
    634             );
    635             assert_eq!(
    636                 operations.enqueue(unavailable).await,
    637                 Err(Error::InvalidPushRequest)
    638             );
    639         }
    640     }
    641 }