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 }