app

Local-first trade for farms and co-ops
git clone https://radroots.dev/git/app.git
Log | Files | Refs | README | LICENSE

availability_outcomes.rs (18192B)


      1 //! Bounded caller-supplied discovery requests, progress and owned valid results.
      2 //!
      3 //! These immutable values perform no I/O, admission, event validation, deadline
      4 //! enforcement or resource reservation. Their counters describe supplied work;
      5 //! they do not establish exhaustive network knowledge or current local storage.
      6 
      7 use std::fmt::{self, Write};
      8 
      9 use harvestcircle_domain::SafeError;
     10 use harvestcircle_domain::error::AvailabilityFailure;
     11 pub use radroots_transport::outcome::FetchTargetState;
     12 pub use radroots_transport::target::TargetFingerprint;
     13 
     14 use crate::{AvailabilityLocalQueryScope, RequestId};
     15 
     16 pub const MAX_DISCOVERY_TARGETS: usize = 16;
     17 pub const MAX_DISCOVERY_FETCH_CALLS: u8 = 4;
     18 pub const DISCOVERY_FETCH_RAW_RESERVATION_BYTES: u64 = 8_388_608;
     19 pub const MAX_DISCOVERY_RETURNED_PER_TARGET: u16 = 64;
     20 pub const MAX_DISCOVERY_RETURNED_EVENTS: u16 = 1_024;
     21 pub const MAX_DISCOVERY_METADATA_BYTES: usize = 16_384;
     22 
     23 // Fixed keys, finite labels and bounded decimals fit in 384 envelope bytes.
     24 // Each target's 64 hex bytes and three progress objects fit in 384 bytes,
     25 // including its separator. The complete structural bound is therefore 6,528.
     26 const METADATA_ENVELOPE_MAX_BYTES: usize = 384;
     27 const METADATA_TARGET_MAX_BYTES: usize = 384;
     28 const _: () = assert!(
     29     METADATA_ENVELOPE_MAX_BYTES + MAX_DISCOVERY_TARGETS * METADATA_TARGET_MAX_BYTES
     30         <= MAX_DISCOVERY_METADATA_BYTES
     31 );
     32 
     33 /// Structural source selection and local bindings without effect authority.
     34 #[derive(Clone, Eq, PartialEq)]
     35 pub struct AvailabilityDiscoveryRequest {
     36     request_id: RequestId,
     37     scope: AvailabilityLocalQueryScope,
     38     targets: Vec<TargetFingerprint>,
     39     deadline_millis: u64,
     40 }
     41 
     42 impl AvailabilityDiscoveryRequest {
     43     /// Retains the original bounded unique target vector and structural scope.
     44     ///
     45     /// # Errors
     46     ///
     47     /// Returns capacity for more than sixteen targets, or invalid input for
     48     /// duplicate canonical fingerprints or a deadline outside 1..=30,000 ms.
     49     pub fn new(
     50         request_id: RequestId,
     51         scope: AvailabilityLocalQueryScope,
     52         targets: Vec<TargetFingerprint>,
     53         deadline_millis: u64,
     54     ) -> Result<Self, SafeError> {
     55         if targets.len() > MAX_DISCOVERY_TARGETS {
     56             return Err(AvailabilityFailure::Capacity.into());
     57         }
     58         if targets
     59             .iter()
     60             .enumerate()
     61             .any(|(index, target)| targets[..index].contains(target))
     62             || !(1..=30_000).contains(&deadline_millis)
     63         {
     64             return Err(AvailabilityFailure::InvalidInput.into());
     65         }
     66         Ok(Self {
     67             request_id,
     68             scope,
     69             targets,
     70             deadline_millis,
     71         })
     72     }
     73 
     74     #[must_use]
     75     pub const fn request_id(&self) -> RequestId {
     76         self.request_id
     77     }
     78 
     79     #[must_use]
     80     pub const fn scope(&self) -> &AvailabilityLocalQueryScope {
     81         &self.scope
     82     }
     83 
     84     #[must_use]
     85     pub fn targets(&self) -> &[TargetFingerprint] {
     86         &self.targets
     87     }
     88 
     89     #[must_use]
     90     pub const fn deadline_millis(&self) -> u64 {
     91         self.deadline_millis
     92     }
     93 }
     94 
     95 impl fmt::Debug for AvailabilityDiscoveryRequest {
     96     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
     97         formatter
     98             .debug_struct("AvailabilityDiscoveryRequest")
     99             .field("selected_targets", &self.targets.len())
    100             .field("deadline_millis", &self.deadline_millis)
    101             .finish_non_exhaustive()
    102     }
    103 }
    104 
    105 /// Exact shared target state, with an explicit absence of requested work.
    106 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
    107 pub enum AvailabilityDiscoveryState {
    108     NotRequested,
    109     Requested(FetchTargetState),
    110 }
    111 
    112 impl AvailabilityDiscoveryState {
    113     const fn metadata_label(self) -> &'static str {
    114         match self {
    115             Self::NotRequested => "not_requested",
    116             Self::Requested(FetchTargetState::Complete) => "complete",
    117             Self::Requested(FetchTargetState::Partial) => "partial",
    118             Self::Requested(FetchTargetState::Unavailable) => "unavailable",
    119             Self::Requested(FetchTargetState::FailedRetryable) => "failed_retryable",
    120             Self::Requested(FetchTargetState::FailedTerminal) => "failed_terminal",
    121             Self::Requested(FetchTargetState::Cancelled) => "cancelled",
    122         }
    123     }
    124 }
    125 
    126 /// Finite target-local stopping information without arbitrary diagnostic text.
    127 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
    128 pub enum AvailabilityDiscoveryStopReason {
    129     None,
    130     BudgetExhausted,
    131     DeadlineExpired,
    132     Stopped,
    133 }
    134 
    135 impl AvailabilityDiscoveryStopReason {
    136     const fn metadata_label(self) -> &'static str {
    137         match self {
    138             Self::None => "none",
    139             Self::BudgetExhausted => "budget_exhausted",
    140             Self::DeadlineExpired => "deadline_expired",
    141             Self::Stopped => "stopped",
    142         }
    143     }
    144 }
    145 
    146 /// One stream's consistent target-local facts, including earlier results.
    147 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
    148 pub struct AvailabilityDiscoveryProgress {
    149     state: AvailabilityDiscoveryState,
    150     reason: AvailabilityDiscoveryStopReason,
    151     returned: u16,
    152 }
    153 
    154 impl AvailabilityDiscoveryProgress {
    155     /// Preserves complete-empty, interrupted and unrequested distinctions.
    156     ///
    157     /// # Errors
    158     ///
    159     /// Returns capacity for more than sixty-four returned events, or invalid
    160     /// input for a reason/count inconsistent with the supplied stream state.
    161     pub fn new(
    162         state: AvailabilityDiscoveryState,
    163         reason: AvailabilityDiscoveryStopReason,
    164         returned: u16,
    165     ) -> Result<Self, SafeError> {
    166         if returned > MAX_DISCOVERY_RETURNED_PER_TARGET {
    167             return Err(AvailabilityFailure::Capacity.into());
    168         }
    169         let valid = match state {
    170             AvailabilityDiscoveryState::NotRequested
    171             | AvailabilityDiscoveryState::Requested(FetchTargetState::Unavailable) => {
    172                 reason == AvailabilityDiscoveryStopReason::None && returned == 0
    173             }
    174             AvailabilityDiscoveryState::Requested(
    175                 FetchTargetState::Complete
    176                 | FetchTargetState::FailedRetryable
    177                 | FetchTargetState::FailedTerminal,
    178             ) => reason == AvailabilityDiscoveryStopReason::None,
    179             AvailabilityDiscoveryState::Requested(FetchTargetState::Partial) => {
    180                 matches!(
    181                     reason,
    182                     AvailabilityDiscoveryStopReason::None
    183                         | AvailabilityDiscoveryStopReason::BudgetExhausted
    184                         | AvailabilityDiscoveryStopReason::DeadlineExpired
    185                 )
    186             }
    187             AvailabilityDiscoveryState::Requested(FetchTargetState::Cancelled) => {
    188                 reason == AvailabilityDiscoveryStopReason::Stopped
    189             }
    190         };
    191         if !valid {
    192             return Err(AvailabilityFailure::InvalidInput.into());
    193         }
    194         Ok(Self {
    195             state,
    196             reason,
    197             returned,
    198         })
    199     }
    200 
    201     #[must_use]
    202     pub const fn state(&self) -> AvailabilityDiscoveryState {
    203         self.state
    204     }
    205 
    206     #[must_use]
    207     pub const fn reason(&self) -> AvailabilityDiscoveryStopReason {
    208         self.reason
    209     }
    210 
    211     #[must_use]
    212     pub const fn returned(&self) -> u16 {
    213         self.returned
    214     }
    215 
    216     fn append_metadata(self, json: &mut String) {
    217         write!(
    218             json,
    219             "{{\"state\":\"{}\",\"reason\":\"{}\",\"returned\":{}}}",
    220             self.state.metadata_label(),
    221             self.reason.metadata_label(),
    222             self.returned
    223         )
    224         .expect("writing finite metadata to String cannot fail");
    225     }
    226 }
    227 
    228 /// Three independent streams sharing one sixty-four-event target allowance.
    229 #[derive(Clone, Eq, PartialEq)]
    230 pub struct AvailabilityDiscoveryTargetOutcome {
    231     target: TargetFingerprint,
    232     listings: AvailabilityDiscoveryProgress,
    233     profiles: AvailabilityDiscoveryProgress,
    234     deletions: AvailabilityDiscoveryProgress,
    235 }
    236 
    237 impl AvailabilityDiscoveryTargetOutcome {
    238     /// Retains admitted progress only after checking the shared target sum.
    239     ///
    240     /// # Errors
    241     ///
    242     /// Returns capacity if the three returned counts together exceed sixty-four.
    243     pub fn new(
    244         target: TargetFingerprint,
    245         listings: AvailabilityDiscoveryProgress,
    246         profiles: AvailabilityDiscoveryProgress,
    247         deletions: AvailabilityDiscoveryProgress,
    248     ) -> Result<Self, SafeError> {
    249         if listings.returned() + profiles.returned() + deletions.returned()
    250             > MAX_DISCOVERY_RETURNED_PER_TARGET
    251         {
    252             return Err(AvailabilityFailure::Capacity.into());
    253         }
    254         Ok(Self {
    255             target,
    256             listings,
    257             profiles,
    258             deletions,
    259         })
    260     }
    261 
    262     #[must_use]
    263     pub const fn target(&self) -> &TargetFingerprint {
    264         &self.target
    265     }
    266 
    267     #[must_use]
    268     pub const fn listings(&self) -> AvailabilityDiscoveryProgress {
    269         self.listings
    270     }
    271 
    272     #[must_use]
    273     pub const fn profiles(&self) -> AvailabilityDiscoveryProgress {
    274         self.profiles
    275     }
    276 
    277     #[must_use]
    278     pub const fn deletions(&self) -> AvailabilityDiscoveryProgress {
    279         self.deletions
    280     }
    281 
    282     const fn returned(&self) -> u16 {
    283         self.listings.returned() + self.profiles.returned() + self.deletions.returned()
    284     }
    285 
    286     fn has_complete_stream(&self) -> bool {
    287         [self.listings, self.profiles, self.deletions]
    288             .iter()
    289             .any(|progress| {
    290                 progress.state()
    291                     == AvailabilityDiscoveryState::Requested(FetchTargetState::Complete)
    292             })
    293     }
    294 }
    295 
    296 impl fmt::Debug for AvailabilityDiscoveryTargetOutcome {
    297     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    298         formatter
    299             .debug_struct("AvailabilityDiscoveryTargetOutcome")
    300             .field("listings", &self.listings)
    301             .field("profiles", &self.profiles)
    302             .field("deletions", &self.deletions)
    303             .finish_non_exhaustive()
    304     }
    305 }
    306 
    307 /// Bounded supplied operation accounting without actual transport reservation.
    308 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
    309 pub struct AvailabilityDiscoveryUsage {
    310     fetch_calls: u8,
    311     returned_events: u16,
    312 }
    313 
    314 impl AvailabilityDiscoveryUsage {
    315     /// Admits at most four fetch calls and 1,024 aggregate returned events.
    316     ///
    317     /// # Errors
    318     ///
    319     /// Returns capacity for either exceeded bound, or invalid input for returned
    320     /// events with zero supplied calls. Reported reservations are never refunded.
    321     pub fn new(fetch_calls: u8, returned_events: u16) -> Result<Self, SafeError> {
    322         if fetch_calls > MAX_DISCOVERY_FETCH_CALLS
    323             || returned_events > MAX_DISCOVERY_RETURNED_EVENTS
    324         {
    325             return Err(AvailabilityFailure::Capacity.into());
    326         }
    327         if fetch_calls == 0 && returned_events != 0 {
    328             return Err(AvailabilityFailure::InvalidInput.into());
    329         }
    330         Ok(Self {
    331             fetch_calls,
    332             returned_events,
    333         })
    334     }
    335 
    336     #[must_use]
    337     pub const fn fetch_calls(&self) -> u8 {
    338         self.fetch_calls
    339     }
    340 
    341     #[must_use]
    342     pub const fn returned_events(&self) -> u16 {
    343         self.returned_events
    344     }
    345 
    346     #[must_use]
    347     pub const fn reserved_raw_bytes(&self) -> u64 {
    348         self.fetch_calls as u64 * DISCOVERY_FETCH_RAW_RESERVATION_BYTES
    349     }
    350 }
    351 
    352 /// Original owned valid items retained across independently interrupted targets.
    353 ///
    354 /// Item admission remains the caller's responsibility. No item serialization,
    355 /// event evidence, freshness or persistence is inferred from these values.
    356 pub struct AvailabilityDiscoveryOutcome<T> {
    357     request: AvailabilityDiscoveryRequest,
    358     usage: AvailabilityDiscoveryUsage,
    359     targets: Vec<AvailabilityDiscoveryTargetOutcome>,
    360     items: Vec<T>,
    361 }
    362 
    363 impl<T> AvailabilityDiscoveryOutcome<T> {
    364     /// Binds exactly one outcome to each selected target without cloning items.
    365     ///
    366     /// # Errors
    367     ///
    368     /// Returns capacity for oversized vectors, scope mismatch for a different
    369     /// target set, or invalid input for duplicate outcomes, inconsistent counts,
    370     /// completion without a call, or fabricated work for an empty selection.
    371     pub fn new(
    372         request: AvailabilityDiscoveryRequest,
    373         usage: AvailabilityDiscoveryUsage,
    374         targets: Vec<AvailabilityDiscoveryTargetOutcome>,
    375         items: Vec<T>,
    376     ) -> Result<Self, SafeError> {
    377         if targets.len() > MAX_DISCOVERY_TARGETS
    378             || items.len() > usize::from(MAX_DISCOVERY_RETURNED_EVENTS)
    379         {
    380             return Err(AvailabilityFailure::Capacity.into());
    381         }
    382         if targets.iter().enumerate().any(|(index, target)| {
    383             targets[..index]
    384                 .iter()
    385                 .any(|previous| previous.target() == target.target())
    386         }) {
    387             return Err(AvailabilityFailure::InvalidInput.into());
    388         }
    389         if targets.len() != request.targets().len()
    390             || targets
    391                 .iter()
    392                 .any(|target| !request.targets().contains(target.target()))
    393         {
    394             return Err(AvailabilityFailure::ScopeMismatch.into());
    395         }
    396         // At most sixteen independently admitted sums of at most sixty-four.
    397         let returned: u16 = targets
    398             .iter()
    399             .map(AvailabilityDiscoveryTargetOutcome::returned)
    400             .sum();
    401         if returned != usage.returned_events()
    402             || items.len() > usize::from(returned)
    403             || (usage.fetch_calls() == 0
    404                 && targets
    405                     .iter()
    406                     .any(AvailabilityDiscoveryTargetOutcome::has_complete_stream))
    407             || (request.targets().is_empty() && usage.fetch_calls() != 0)
    408         {
    409             return Err(AvailabilityFailure::InvalidInput.into());
    410         }
    411         Ok(Self {
    412             request,
    413             usage,
    414             targets,
    415             items,
    416         })
    417     }
    418 
    419     #[must_use]
    420     pub const fn request(&self) -> &AvailabilityDiscoveryRequest {
    421         &self.request
    422     }
    423 
    424     #[must_use]
    425     pub const fn usage(&self) -> AvailabilityDiscoveryUsage {
    426         self.usage
    427     }
    428 
    429     #[must_use]
    430     pub fn targets(&self) -> &[AvailabilityDiscoveryTargetOutcome] {
    431         &self.targets
    432     }
    433 
    434     #[must_use]
    435     pub fn items(&self) -> &[T] {
    436         &self.items
    437     }
    438 
    439     #[must_use]
    440     pub fn into_items(self) -> Vec<T> {
    441         self.items
    442     }
    443 
    444     #[must_use]
    445     pub fn listings_state(&self) -> AvailabilityDiscoveryState {
    446         self.stream_state(AvailabilityDiscoveryTargetOutcome::listings)
    447     }
    448 
    449     #[must_use]
    450     pub fn profiles_state(&self) -> AvailabilityDiscoveryState {
    451         self.stream_state(AvailabilityDiscoveryTargetOutcome::profiles)
    452     }
    453 
    454     #[must_use]
    455     pub fn deletions_state(&self) -> AvailabilityDiscoveryState {
    456         self.stream_state(AvailabilityDiscoveryTargetOutcome::deletions)
    457     }
    458 
    459     fn stream_state(
    460         &self,
    461         progress: fn(&AvailabilityDiscoveryTargetOutcome) -> AvailabilityDiscoveryProgress,
    462     ) -> AvailabilityDiscoveryState {
    463         let Some(first) = self.targets.first() else {
    464             return AvailabilityDiscoveryState::NotRequested;
    465         };
    466         let state = progress(first).state();
    467         if self
    468             .targets
    469             .iter()
    470             .all(|target| progress(target).state() == state)
    471         {
    472             state
    473         } else {
    474             AvailabilityDiscoveryState::Requested(FetchTargetState::Partial)
    475         }
    476     }
    477 
    478     /// Serializes only finite metadata, preserving the supplied target order.
    479     ///
    480     /// The fixed envelope and sixteen bounded target objects fit within 6,528
    481     /// ASCII bytes, below the 16 KiB contract. Only shared admitted fingerprints,
    482     /// finite labels and bounded decimals enter JSON; items and scope do not.
    483     #[must_use]
    484     pub fn metadata_json(&self) -> String {
    485         let mut json = String::with_capacity(
    486             METADATA_ENVELOPE_MAX_BYTES + self.targets.len() * METADATA_TARGET_MAX_BYTES,
    487         );
    488         write!(
    489             json,
    490             concat!(
    491                 "{{\"version\":1,\"fetch_calls\":{},\"reserved_raw_bytes\":{},",
    492                 "\"returned_events\":{},\"retained_items\":{},\"listings\":\"{}\",",
    493                 "\"profiles\":\"{}\",\"deletions\":\"{}\",\"targets\":["
    494             ),
    495             self.usage.fetch_calls(),
    496             self.usage.reserved_raw_bytes(),
    497             self.usage.returned_events(),
    498             self.items.len(),
    499             self.listings_state().metadata_label(),
    500             self.profiles_state().metadata_label(),
    501             self.deletions_state().metadata_label()
    502         )
    503         .expect("writing finite metadata to String cannot fail");
    504         for (index, target) in self.targets.iter().enumerate() {
    505             if index != 0 {
    506                 json.push(',');
    507             }
    508             json.push_str("{\"target\":\"");
    509             json.push_str(target.target().as_str());
    510             json.push_str("\",\"listings\":");
    511             target.listings().append_metadata(&mut json);
    512             json.push_str(",\"profiles\":");
    513             target.profiles().append_metadata(&mut json);
    514             json.push_str(",\"deletions\":");
    515             target.deletions().append_metadata(&mut json);
    516             json.push('}');
    517         }
    518         json.push_str("]}");
    519         json
    520     }
    521 }
    522 
    523 impl<T> fmt::Debug for AvailabilityDiscoveryOutcome<T> {
    524     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    525         formatter
    526             .debug_struct("AvailabilityDiscoveryOutcome")
    527             .field("fetch_calls", &self.usage.fetch_calls())
    528             .field("reserved_raw_bytes", &self.usage.reserved_raw_bytes())
    529             .field("returned_events", &self.usage.returned_events())
    530             .field("retained_items", &self.items.len())
    531             .field("selected_targets", &self.targets.len())
    532             .field("listings", &self.listings_state())
    533             .field("profiles", &self.profiles_state())
    534             .field("deletions", &self.deletions_state())
    535             .finish_non_exhaustive()
    536     }
    537 }