lib

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

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 }