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 }