farm.rs (18045B)
1 //! Side-effect-free farm planning and canonical durable enqueue operations. 2 3 use std::{error, fmt}; 4 5 use radroots_event::{ 6 GenericEventDraft, 7 contract::AuthorRole, 8 envelope::kind::KIND_FARM, 9 farm::Farm, 10 id::{AddressableCoordinate, ParseError}, 11 }; 12 use radroots_event_codec::{ 13 authoring::AuthoredEventPlan, encode::EventEncodeError, encode::farm::to_wire_parts, 14 }; 15 use radroots_signing::Actor; 16 17 const FARM_PROFILE_CONTRACT_ID: &str = "radroots.farm.profile.v1"; 18 19 /// Pure inputs for one frozen farm profile plan. 20 #[derive(Clone, Debug)] 21 pub struct PrepareRequest { 22 actor: Actor, 23 farm: Farm, 24 created_at_unix: u64, 25 } 26 27 impl PrepareRequest { 28 /// Creates explicit canonical planning inputs. 29 #[must_use] 30 pub const fn new(actor: Actor, farm: Farm, created_at_unix: u64) -> Self { 31 Self { 32 actor, 33 farm, 34 created_at_unix, 35 } 36 } 37 } 38 39 /// Frozen, replay-stable farm publication plan. 40 #[derive(Clone, Debug, Eq, PartialEq)] 41 pub struct Plan { 42 actor: Actor, 43 coordinate: AddressableCoordinate, 44 authored_event: AuthoredEventPlan, 45 } 46 47 impl Plan { 48 /// Returns the exact authorized actor carried into signing. 49 #[must_use] 50 pub const fn actor(&self) -> &Actor { 51 &self.actor 52 } 53 54 /// Returns the canonical addressable farm coordinate. 55 #[must_use] 56 pub const fn coordinate(&self) -> &AddressableCoordinate { 57 &self.coordinate 58 } 59 60 /// Returns the immutable canonical authored event plan. 61 #[must_use] 62 pub const fn authored_event(&self) -> &AuthoredEventPlan { 63 &self.authored_event 64 } 65 } 66 67 /// Farm plan validation stage. 68 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 69 #[non_exhaustive] 70 pub enum PrepareErrorKind { 71 /// The actor does not claim the required farm author role. 72 UnauthorizedActor, 73 /// The canonical farm codec rejected the native farm model. 74 Encode, 75 /// The canonical addressable coordinate rejected the farm identity. 76 Coordinate, 77 /// The canonical event draft rejected the encoded parts. 78 Draft, 79 } 80 81 /// One secret-safe farm planning failure retaining its lower source. 82 pub struct PrepareError { 83 kind: PrepareErrorKind, 84 source: Option<Box<dyn error::Error + Send + Sync>>, 85 } 86 87 impl PrepareError { 88 /// Returns the stable client-level planning stage. 89 #[must_use] 90 pub const fn kind(&self) -> PrepareErrorKind { 91 self.kind 92 } 93 94 fn unauthorized_actor() -> Self { 95 Self { 96 kind: PrepareErrorKind::UnauthorizedActor, 97 source: None, 98 } 99 } 100 101 fn encode(source: EventEncodeError) -> Self { 102 Self::with_source(PrepareErrorKind::Encode, source) 103 } 104 105 fn coordinate(source: ParseError) -> Self { 106 Self::with_source(PrepareErrorKind::Coordinate, source) 107 } 108 109 fn draft(source: radroots_event::draft::DraftError) -> Self { 110 Self::with_source(PrepareErrorKind::Draft, source) 111 } 112 113 fn with_source( 114 kind: PrepareErrorKind, 115 source: impl error::Error + Send + Sync + 'static, 116 ) -> Self { 117 Self { 118 kind, 119 source: Some(Box::new(source)), 120 } 121 } 122 } 123 124 impl fmt::Display for PrepareError { 125 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 126 formatter.write_str(match self.kind { 127 PrepareErrorKind::UnauthorizedActor => "farm actor is not authorized", 128 PrepareErrorKind::Encode => "farm model is invalid", 129 PrepareErrorKind::Coordinate => "farm coordinate is invalid", 130 PrepareErrorKind::Draft => "farm event draft is invalid", 131 }) 132 } 133 } 134 135 impl fmt::Debug for PrepareError { 136 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 137 formatter 138 .debug_struct("PrepareError") 139 .field("kind", &self.kind) 140 .finish_non_exhaustive() 141 } 142 } 143 144 impl error::Error for PrepareError { 145 fn source(&self) -> Option<&(dyn error::Error + 'static)> { 146 self.source 147 .as_deref() 148 .map(|source| source as &(dyn error::Error + 'static)) 149 } 150 } 151 152 /// Validates and freezes one farm profile without storage, signing, or network work. 153 /// 154 /// `Farm` contains only public profile locality. Exact coordinates and other 155 /// private farm artifacts are deliberately not accepted here; hosts persist 156 /// their typed references and metadata through 157 /// `radroots_storage::private_artifact::PrivateArtifactStore`. 158 pub fn prepare(request: PrepareRequest) -> Result<Plan, PrepareError> { 159 if !request.actor.satisfies(AuthorRole::Farmer) { 160 return Err(PrepareError::unauthorized_actor()); 161 } 162 let parts = to_wire_parts(&request.farm).map_err(PrepareError::encode)?; 163 let coordinate = AddressableCoordinate::parse(format!( 164 "{KIND_FARM}:{}:{}", 165 request.actor.public_key(), 166 request.farm.d_tag 167 )) 168 .map_err(PrepareError::coordinate)?; 169 let authored_event = AuthoredEventPlan::from_generic( 170 GenericEventDraft::new( 171 FARM_PROFILE_CONTRACT_ID, 172 parts.kind, 173 request.created_at_unix, 174 parts.tags, 175 parts.content, 176 request.actor.public_key().to_hex(), 177 ) 178 .map_err(PrepareError::draft)?, 179 ) 180 .map_err(PrepareError::draft)?; 181 Ok(Plan { 182 actor: request.actor, 183 coordinate, 184 authored_event, 185 }) 186 } 187 188 #[cfg(feature = "sync")] 189 use radroots_signing::request::CancellationPolicy; 190 #[cfg(feature = "sync")] 191 use radroots_storage::journal::IdempotencyKey; 192 #[cfg(feature = "sync")] 193 use radroots_sync::{ 194 policy::{Error as SyncError, SyncId}, 195 push::PushStatus, 196 }; 197 198 /// Explicit commit inputs for one prepared farm publication. 199 #[cfg(feature = "sync")] 200 #[derive(Clone, Debug)] 201 pub struct EnqueueRequest { 202 operation_id: SyncId, 203 idempotency_key: IdempotencyKey, 204 plan: Plan, 205 profile: crate::transport::Profile, 206 delivery_deadline_unix_ms: u64, 207 cancellation: CancellationPolicy, 208 } 209 210 #[cfg(feature = "sync")] 211 impl EnqueueRequest { 212 /// Creates an enqueue request whose transport selection has no fallback. 213 #[must_use] 214 pub const fn new( 215 operation_id: SyncId, 216 idempotency_key: IdempotencyKey, 217 plan: Plan, 218 profile: crate::transport::Profile, 219 delivery_deadline_unix_ms: u64, 220 cancellation: CancellationPolicy, 221 ) -> Self { 222 Self { 223 operation_id, 224 idempotency_key, 225 plan, 226 profile, 227 delivery_deadline_unix_ms, 228 cancellation, 229 } 230 } 231 } 232 233 /// Borrowed farm commit operations over the canonical sync engine. 234 #[cfg(feature = "sync")] 235 #[derive(Clone, Copy, Debug)] 236 pub struct Operations<'a> { 237 sync: crate::sync::Operations<'a>, 238 } 239 240 #[cfg(feature = "sync")] 241 impl<'a> Operations<'a> { 242 pub(crate) const fn new(sync: crate::sync::Operations<'a>) -> Self { 243 Self { sync } 244 } 245 246 /// Durably prepares, signs, and locally admits a farm publication. 247 /// 248 /// Once preparation commits, cancellation cannot erase the authored 249 /// intent. Replay with identical operation and idempotency inputs resumes 250 /// from its durable signing or admission phase. 251 pub async fn enqueue(&self, request: EnqueueRequest) -> Result<PushStatus, SyncError> { 252 let targets = request 253 .profile 254 .targets() 255 .cloned() 256 .ok_or(SyncError::InvalidPushRequest)?; 257 let satisfaction = request 258 .profile 259 .satisfaction() 260 .cloned() 261 .ok_or(SyncError::InvalidPushRequest)?; 262 self.sync 263 .submit_push(radroots_sync::PushRequest::new( 264 request.operation_id, 265 request.idempotency_key, 266 request.plan.actor, 267 request.plan.authored_event, 268 targets, 269 satisfaction, 270 request.delivery_deadline_unix_ms, 271 request.cancellation, 272 )?) 273 .await 274 } 275 } 276 277 #[cfg(test)] 278 mod tests { 279 use radroots_event::farm::FarmPublicLocation; 280 use radroots_signing::actor::ActorSource; 281 282 use super::*; 283 284 const PUBLIC_KEY: &str = "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df"; 285 286 fn actor(role: AuthorRole) -> Actor { 287 Actor::from_public_key_hex(PUBLIC_KEY, ActorSource::ExplicitPublicKey, [role]) 288 .expect("actor") 289 } 290 291 fn farm() -> Farm { 292 Farm { 293 d_tag: "AAAAAAAAAAAAAAAAAAAAAA".to_owned(), 294 name: "Moss Street Farm".to_owned(), 295 about: Some("seasonal vegetables".to_owned()), 296 website: None, 297 picture: None, 298 banner: None, 299 location: Some(FarmPublicLocation { 300 primary: "Santa Cruz, California".to_owned(), 301 city: Some("Santa Cruz".to_owned()), 302 region: Some("California".to_owned()), 303 country: Some("US".to_owned()), 304 geohash: "9q8yy".to_owned(), 305 }), 306 tags: Some(vec!["vegetables".to_owned()]), 307 } 308 } 309 310 #[test] 311 fn prepare_is_pure_deterministic_and_uses_canonical_public_types() { 312 let request = PrepareRequest::new(actor(AuthorRole::Farmer), farm(), 1_800_000_000); 313 let first = prepare(request.clone()).expect("first plan"); 314 let second = prepare(request).expect("second plan"); 315 316 assert_eq!(first, second); 317 assert_eq!( 318 first 319 .authored_event() 320 .body() 321 .contract() 322 .contract_id() 323 .as_str(), 324 FARM_PROFILE_CONTRACT_ID 325 ); 326 assert_eq!(first.authored_event().body().kind(), KIND_FARM); 327 assert_eq!(first.authored_event().created_at(), 1_800_000_000); 328 assert_eq!( 329 first.authored_event().author(), 330 &actor(AuthorRole::Farmer).public_key() 331 ); 332 assert_eq!( 333 first.coordinate().as_str(), 334 format!("{KIND_FARM}:{PUBLIC_KEY}:AAAAAAAAAAAAAAAAAAAAAA") 335 ); 336 assert!( 337 first 338 .authored_event() 339 .body() 340 .content() 341 .contains("Moss Street Farm") 342 ); 343 assert!(!first.authored_event().body().content().contains("latitude")); 344 assert!( 345 !first 346 .authored_event() 347 .body() 348 .content() 349 .contains("longitude") 350 ); 351 } 352 353 #[test] 354 fn prepare_maps_authorization_and_lower_validation_once() { 355 let unauthorized = prepare(PrepareRequest::new( 356 actor(AuthorRole::Buyer), 357 farm(), 358 1_800_000_000, 359 )) 360 .expect_err("unauthorized"); 361 assert_eq!(unauthorized.kind(), PrepareErrorKind::UnauthorizedActor); 362 assert!(std::error::Error::source(&unauthorized).is_none()); 363 364 let mut invalid = farm(); 365 invalid.name.clear(); 366 let encoded = prepare(PrepareRequest::new( 367 actor(AuthorRole::Farmer), 368 invalid, 369 1_800_000_000, 370 )) 371 .expect_err("invalid model"); 372 assert_eq!(encoded.kind(), PrepareErrorKind::Encode); 373 assert!(std::error::Error::source(&encoded).is_some()); 374 assert_eq!(encoded.to_string(), "farm model is invalid"); 375 assert!(!format!("{encoded:?}").contains("name")); 376 } 377 378 #[cfg(all(feature = "sync", feature = "memory", feature = "local-signing"))] 379 mod enqueue { 380 use std::sync::{ 381 Arc, 382 atomic::{AtomicU8, Ordering}, 383 }; 384 385 use radroots_storage::{ 386 Outbox, event::SourceGeneration, journal::IdempotencyKey, memory::MemoryStorage, 387 }; 388 use radroots_sync::{ 389 Engine, 390 policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId, SyncStorage}, 391 }; 392 use radroots_transport::{ 393 DeliveryReceipt, DeliveryRequest, Error as TransportError, EventSink, SinkFailure, 394 SinkStatus, Target, TargetSet, TransportId, 395 capability::{Availability, Maturity, SinkCapabilities}, 396 outcome::Retryability, 397 policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy}, 398 }; 399 400 use crate::{ClientBuilder, transport::Profile}; 401 402 use super::*; 403 404 struct HostClock; 405 struct SequenceIds(AtomicU8); 406 struct NoopSink; 407 408 impl Clock for HostClock { 409 fn now_unix_ms(&self) -> Result<u64, Error> { 410 std::time::SystemTime::now() 411 .duration_since(std::time::UNIX_EPOCH) 412 .ok() 413 .and_then(|duration| u64::try_from(duration.as_millis()).ok()) 414 .filter(|value| *value != 0) 415 .ok_or(Error::ClockUnavailable) 416 } 417 } 418 419 impl IdSource for SequenceIds { 420 fn next_id(&self, _operation: OperationKind) -> Result<SyncId, Error> { 421 SyncId::new([self.0.fetch_add(1, Ordering::Relaxed); 16]) 422 } 423 } 424 425 impl EventSink for NoopSink { 426 fn status( 427 &self, 428 ) -> radroots_transport::BoxFuture<'_, Result<SinkStatus, TransportError>> { 429 Box::pin(async { 430 Ok(SinkStatus::new( 431 TransportId::NOSTR, 432 true, 433 Maturity::Stable, 434 Availability::Available, 435 SinkCapabilities::DELIVER, 436 "ready", 437 )) 438 }) 439 } 440 441 fn deliver( 442 &self, 443 request: DeliveryRequest, 444 ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> 445 { 446 Box::pin(async move { 447 Err(SinkFailure::for_request( 448 &request, 449 "test_sink_unavailable", 450 Retryability::Terminal, 451 None, 452 None, 453 Vec::new(), 454 ) 455 .expect("test sink failure")) 456 }) 457 } 458 } 459 460 #[tokio::test] 461 async fn enqueue_preserves_commit_cancellation_idempotency_and_delivery_policy() { 462 let storage = Arc::new(MemoryStorage::new( 463 SourceGeneration::new([4; 32]).expect("generation"), 464 )); 465 let signer = 466 Arc::new(radroots_nostr::signing::LocalSigner::generate().expect("local signer")); 467 let farm_actor = Actor::new( 468 signer.public_key(), 469 ActorSource::ExplicitPublicKey, 470 [AuthorRole::Farmer], 471 ) 472 .expect("actor"); 473 let plan = 474 prepare(PrepareRequest::new(farm_actor, farm(), 1_800_000_000)).expect("plan"); 475 let targets = TargetSet::new(vec![ 476 Target::nostr_relay("wss://farm.example").expect("target"), 477 ]) 478 .expect("targets"); 479 let satisfaction = 480 SatisfactionPolicy::new(SatisfactionClass::Delivered, TargetPolicy::all()); 481 let profile = Profile::delivery(targets.clone(), satisfaction.clone()) 482 .expect("transport profile"); 483 let capability: Arc<dyn SyncStorage> = storage.clone(); 484 let engine = Engine::builder( 485 capability, 486 Arc::new(HostClock), 487 Arc::new(SequenceIds(AtomicU8::new(1))), 488 DeadlinePolicy::new(30_000, 30_000, 30_000).expect("deadlines"), 489 ) 490 .sink(Arc::new(NoopSink)) 491 .signer(signer) 492 .build() 493 .expect("engine"); 494 let client = ClientBuilder::new() 495 .storage(storage.clone()) 496 .sync_engine(engine) 497 .build() 498 .expect("client"); 499 let operations = client.farm().expect("open").expect("farm operations"); 500 let request = EnqueueRequest::new( 501 SyncId::new([7; 16]).expect("operation id"), 502 IdempotencyKey::parse("farm-publish-a").expect("idempotency key"), 503 plan, 504 profile, 505 2_000_000_001_000, 506 CancellationPolicy::PreservePublishedRequest, 507 ); 508 509 drop(operations.enqueue(request.clone())); 510 assert_eq!( 511 Outbox::status(storage.as_ref()) 512 .await 513 .expect("outbox status") 514 .pending, 515 0 516 ); 517 518 let committed = operations 519 .enqueue(request.clone()) 520 .await 521 .expect("committed enqueue"); 522 assert_eq!(committed.delivery_plan().intent().target_set(), &targets); 523 assert_eq!( 524 committed.delivery_plan().intent().satisfaction(), 525 &satisfaction 526 ); 527 let replay = operations.enqueue(request).await.expect("replay"); 528 assert_eq!(replay, committed); 529 530 let unavailable = EnqueueRequest::new( 531 SyncId::new([8; 16]).expect("operation id"), 532 IdempotencyKey::parse("farm-publish-preview").expect("idempotency key"), 533 prepare(PrepareRequest::new( 534 actor(AuthorRole::Farmer), 535 farm(), 536 1_800_000_000, 537 )) 538 .expect("plan"), 539 Profile::unavailable_preview(TransportId::RETICULUM), 540 2_000_000_001_000, 541 CancellationPolicy::PreservePublishedRequest, 542 ); 543 assert_eq!( 544 operations.enqueue(unavailable).await, 545 Err(Error::InvalidPushRequest) 546 ); 547 } 548 } 549 }