availability_ports.rs (14196B)
1 //! Trusted local availability admission, read ports and independent refresh. 2 //! 3 //! This service validates structural requests and returned evidence. Concrete 4 //! adapters must admit the local account/data at the actual I/O boundary and 5 //! bind reads atomically to the admitted storage snapshot. They also own scan, 6 //! byte, deadline and response-buffer limits. No storage, network, signer, 7 //! scheduler, native admission or runtime resource enforcement is implemented 8 //! here, and local rows remain outside the application snapshot. 9 10 use std::sync::Arc; 11 use std::time::{Duration, Instant}; 12 13 use harvestcircle_domain::error::AvailabilityFailure; 14 use harvestcircle_domain::{ 15 AvailabilityEventVersion, AvailabilityHeadView, AvailabilityListingCoordinate, 16 AvailabilityPage, AvailabilityPageContinuation, AvailabilityPageCursor, 17 AvailabilityVersionView, SafeError, 18 }; 19 20 use crate::{ 21 AvailabilityDiscoveryOutcome, AvailabilityDiscoveryRequest, AvailabilityLocalQueryScope, 22 BoxFuture, CommandContext, CommandReceipt, CommandResult, CommandSubmission, RequestId, 23 ScopedAvailabilityQuery, 24 }; 25 26 /// Trusted local account/data authority, independent of signing capability. 27 /// 28 /// A concrete runtime adapter must establish OS/account/data admission and 29 /// return the selected owner's persisted context/store/source/projection and 30 /// session bindings. A caller-created structural scope is not admission proof. 31 pub trait AvailabilityLocalAdmission: Send + Sync { 32 fn current_scope<'a>(&'a self) 33 -> BoxFuture<'a, Result<AvailabilityLocalQueryScope, SafeError>>; 34 } 35 36 /// Local storage reads without transport or signing dependencies. 37 /// 38 /// Concrete adapters must enforce admission and snapshot selection at I/O, 39 /// order/filter rows and continuations consistently, and enforce the actual 40 /// scan, byte, deadline and response-buffer budgets. Returned errors retain 41 /// their existing safe codes and messages when the scope remains current. 42 pub trait AvailabilityLocalReadPort: Send + Sync { 43 fn read_page<'a>( 44 &'a self, 45 query: &'a ScopedAvailabilityQuery, 46 ) -> BoxFuture<'a, Result<AvailabilityPage<AvailabilityHeadView>, SafeError>>; 47 48 fn read_head<'a>( 49 &'a self, 50 scope: &'a AvailabilityLocalQueryScope, 51 coordinate: &'a AvailabilityListingCoordinate, 52 ) -> BoxFuture<'a, Result<AvailabilityHeadView, SafeError>>; 53 54 /// `None` describes missing retained evidence, not exhaustive relay data. 55 fn read_version<'a>( 56 &'a self, 57 scope: &'a AvailabilityLocalQueryScope, 58 coordinate: &'a AvailabilityListingCoordinate, 59 version: AvailabilityEventVersion, 60 ) -> BoxFuture<'a, Result<Option<AvailabilityVersionView>, SafeError>>; 61 } 62 63 /// Nonblocking refresh admission using the existing bounded command protocol. 64 /// 65 /// The later runtime adapter owns execution, supervision, deadline enforcement, 66 /// join/account-switch handling and resource accounting. Dropping a caller's 67 /// ticket establishes neither cancellation nor task-permit release. 68 pub trait AvailabilityRefreshPort: Send + Sync { 69 fn submit( 70 &self, 71 context: CommandContext, 72 request: AvailabilityDiscoveryRequest, 73 ) -> CommandSubmission<AvailabilityDiscoveryOutcome<AvailabilityEventVersion>>; 74 } 75 76 /// Trusted per-instance monotonic clock in the command deadline's time domain. 77 /// 78 /// Composition must supply monotonic `std::time::Instant` values. This seam 79 /// creates no timer and does not admit caller-controlled presentation clocks. 80 pub trait AvailabilityMonotonicClock: Send + Sync { 81 fn now(&self) -> Instant; 82 } 83 84 struct SystemClock; 85 86 impl AvailabilityMonotonicClock for SystemClock { 87 fn now(&self) -> Instant { 88 Instant::now() 89 } 90 } 91 92 /// Narrow orchestration over trusted admission and independently usable ports. 93 pub struct AvailabilityQueryService { 94 admission: Arc<dyn AvailabilityLocalAdmission>, 95 reader: Arc<dyn AvailabilityLocalReadPort>, 96 refresh: Arc<dyn AvailabilityRefreshPort>, 97 clock: Arc<dyn AvailabilityMonotonicClock>, 98 } 99 100 impl AvailabilityQueryService { 101 #[must_use] 102 pub fn new( 103 admission: Arc<dyn AvailabilityLocalAdmission>, 104 reader: Arc<dyn AvailabilityLocalReadPort>, 105 refresh: Arc<dyn AvailabilityRefreshPort>, 106 ) -> Self { 107 Self::new_with_clock(admission, reader, refresh, Arc::new(SystemClock)) 108 } 109 110 #[must_use] 111 pub fn new_with_clock( 112 admission: Arc<dyn AvailabilityLocalAdmission>, 113 reader: Arc<dyn AvailabilityLocalReadPort>, 114 refresh: Arc<dyn AvailabilityRefreshPort>, 115 clock: Arc<dyn AvailabilityMonotonicClock>, 116 ) -> Self { 117 Self { 118 admission, 119 reader, 120 refresh, 121 clock, 122 } 123 } 124 125 /// Reads an admitted local page and preserves its owned rows/continuation. 126 /// 127 /// # Errors 128 /// 129 /// Returns trusted admission/scope or port errors, capacity for too many 130 /// rows, stale query for a different projection, and the existing cursor 131 /// error for a continuation belonging to a different complete request. 132 pub async fn read_page( 133 &self, 134 query: &ScopedAvailabilityQuery, 135 ) -> Result<AvailabilityPage<AvailabilityHeadView>, SafeError> { 136 validate_admission(self.admission.as_ref(), query.scope()).await?; 137 let result = self.reader.read_page(query).await; 138 validate_admission(self.admission.as_ref(), query.scope()).await?; 139 let page = result?; 140 if page.items().len() > usize::from(query.limit().rows()) { 141 return Err(AvailabilityFailure::Capacity.into()); 142 } 143 if page.projection_generation() != query.scope().context().projection_generation() { 144 return Err(AvailabilityFailure::StaleQuery.into()); 145 } 146 if let AvailabilityPageContinuation::More(cursor) = page.continuation() { 147 AvailabilityPageCursor::parse(cursor.as_str(), query.fingerprint())?; 148 } 149 Ok(page) 150 } 151 152 /// Reads the current selected head, retaining missing/suppression evidence. 153 /// 154 /// # Errors 155 /// 156 /// Returns trusted admission/scope or port errors, or scope mismatch for 157 /// a returned coordinate differing from the request, including absence. 158 pub async fn read_head( 159 &self, 160 scope: &AvailabilityLocalQueryScope, 161 coordinate: &AvailabilityListingCoordinate, 162 ) -> Result<AvailabilityHeadView, SafeError> { 163 validate_admission(self.admission.as_ref(), scope).await?; 164 let result = self.reader.read_head(scope, coordinate).await; 165 validate_admission(self.admission.as_ref(), scope).await?; 166 let head = result?; 167 if head.listing_coordinate() != Some(coordinate) { 168 return Err(AvailabilityFailure::ScopeMismatch.into()); 169 } 170 Ok(head) 171 } 172 173 /// Reads exact retained historical evidence without reconstructing a head. 174 /// 175 /// # Errors 176 /// 177 /// Returns trusted admission/scope or port errors, or scope mismatch for 178 /// a retained version with a different coordinate or event ID. 179 pub async fn read_version( 180 &self, 181 scope: &AvailabilityLocalQueryScope, 182 coordinate: &AvailabilityListingCoordinate, 183 version: AvailabilityEventVersion, 184 ) -> Result<Option<AvailabilityVersionView>, SafeError> { 185 validate_admission(self.admission.as_ref(), scope).await?; 186 let result = self.reader.read_version(scope, coordinate, version).await; 187 validate_admission(self.admission.as_ref(), scope).await?; 188 let retained = result?; 189 if retained.as_ref().is_some_and(|view| { 190 view.listing_coordinate() != Some(coordinate) || view.version() != version 191 }) { 192 return Err(AvailabilityFailure::ScopeMismatch.into()); 193 } 194 Ok(retained) 195 } 196 197 /// Admits and submits refresh independently, preserving the whole context. 198 /// 199 /// The initial remaining absolute deadline must fit the request's relative 200 /// budget. Awaited admission never renews it. Expiry before/after admission 201 /// produces a correlated timeout receipt without invoking the refresh port. 202 /// This method does not wait for accepted work or automatically execute it. 203 /// 204 /// # Errors 205 /// 206 /// Returns invalid input for ID/deadline disagreement, or the trusted 207 /// admission/scope error. Submission/receipt failures use the original ID. 208 pub async fn submit_refresh( 209 &self, 210 request: AvailabilityDiscoveryRequest, 211 context: CommandContext, 212 ) -> Result<AvailabilityRefreshOperation, SafeError> { 213 if request.request_id() != context.request_id() { 214 return Err(AvailabilityFailure::InvalidInput.into()); 215 } 216 let initial_now = self.clock.now(); 217 let expired_before = context.is_expired(initial_now); 218 if !expired_before 219 && context.deadline().duration_since(initial_now) 220 > Duration::from_millis(request.deadline_millis()) 221 { 222 return Err(AvailabilityFailure::InvalidInput.into()); 223 } 224 validate_admission(self.admission.as_ref(), request.scope()).await?; 225 let submission = if expired_before || context.is_expired(self.clock.now()) { 226 CommandSubmission::Rejected(CommandReceipt::new( 227 request.request_id(), 228 CommandResult::TimedOut, 229 )) 230 } else { 231 // Only this already bounded structural request is cloned. Local 232 // view/page datasets are never cloned by service orchestration. 233 self.refresh.submit(context, request.clone()) 234 }; 235 Ok(AvailabilityRefreshOperation::new( 236 request, 237 submission, 238 self.admission.clone(), 239 )) 240 } 241 } 242 243 /// Original request correlation and the existing independent command ticket. 244 /// 245 /// Dropping this wrapper only loses the caller's reply receiver; it does not 246 /// cancel queued/running work or establish release of runtime task capacity. 247 pub struct AvailabilityRefreshOperation { 248 request: AvailabilityDiscoveryRequest, 249 submission: CommandSubmission<AvailabilityDiscoveryOutcome<AvailabilityEventVersion>>, 250 admission: Arc<dyn AvailabilityLocalAdmission>, 251 } 252 253 impl AvailabilityRefreshOperation { 254 fn new( 255 request: AvailabilityDiscoveryRequest, 256 submission: CommandSubmission<AvailabilityDiscoveryOutcome<AvailabilityEventVersion>>, 257 admission: Arc<dyn AvailabilityLocalAdmission>, 258 ) -> Self { 259 let submission = if submission.request_id() == request.request_id() { 260 submission 261 } else { 262 // Reject both foreign accepted and rejected submissions now, 263 // without ever awaiting a foreign ticket. Dropping that ticket 264 // does not remove its envelope or cancel actual adapter work. 265 drop(submission); 266 CommandSubmission::Rejected(failed_receipt( 267 request.request_id(), 268 AvailabilityFailure::InvalidInput.into(), 269 )) 270 }; 271 Self { 272 request, 273 submission, 274 admission, 275 } 276 } 277 278 #[must_use] 279 pub const fn request_id(&self) -> RequestId { 280 self.request.request_id() 281 } 282 283 /// Collects only this command's receipt and rechecks trusted scope. 284 /// 285 /// A late authority denial/mismatch wins before publishing any completion 286 /// or failure. Foreign IDs or changed outcome bindings become fixed typed 287 /// failures under the original ID. Otherwise all command variants and valid 288 /// partial results remain intact, even after the original effect deadline: 289 /// collecting a receipt creates no work and cannot renew that deadline. 290 pub async fn receipt( 291 self, 292 ) -> CommandReceipt<AvailabilityDiscoveryOutcome<AvailabilityEventVersion>> { 293 let Self { 294 request, 295 submission, 296 admission, 297 } = self; 298 let receipt = match submission { 299 CommandSubmission::Accepted(ticket) => ticket.receipt().await, 300 CommandSubmission::Rejected(receipt) => receipt, 301 }; 302 if let Err(error) = validate_admission(admission.as_ref(), request.scope()).await { 303 return failed_receipt(request.request_id(), error); 304 } 305 if receipt.request_id() != request.request_id() { 306 return failed_receipt( 307 request.request_id(), 308 AvailabilityFailure::InvalidInput.into(), 309 ); 310 } 311 if let CommandResult::Completed(outcome) = receipt.result() 312 && let Err(error) = validate_outcome_request(&request, outcome.request()) 313 { 314 return failed_receipt(request.request_id(), error); 315 } 316 receipt 317 } 318 } 319 320 async fn validate_admission( 321 admission: &dyn AvailabilityLocalAdmission, 322 expected: &AvailabilityLocalQueryScope, 323 ) -> Result<(), SafeError> { 324 let current = admission.current_scope().await?; 325 expected.validate_current(¤t).map_err(SafeError::from) 326 } 327 328 fn validate_outcome_request( 329 expected: &AvailabilityDiscoveryRequest, 330 returned: &AvailabilityDiscoveryRequest, 331 ) -> Result<(), SafeError> { 332 if expected.request_id() != returned.request_id() 333 || expected.deadline_millis() != returned.deadline_millis() 334 { 335 return Err(AvailabilityFailure::InvalidInput.into()); 336 } 337 expected.scope().validate_current(returned.scope())?; 338 // Request/outcome constructors already retain unique bounded target sets 339 // and exact counts. Set agreement does not discard either supplied order. 340 if expected.targets().len() != returned.targets().len() 341 || returned 342 .targets() 343 .iter() 344 .any(|target| !expected.targets().contains(target)) 345 { 346 return Err(AvailabilityFailure::ScopeMismatch.into()); 347 } 348 Ok(()) 349 } 350 351 fn failed_receipt( 352 request_id: RequestId, 353 error: SafeError, 354 ) -> CommandReceipt<AvailabilityDiscoveryOutcome<AvailabilityEventVersion>> { 355 CommandReceipt::new(request_id, CommandResult::Failed(error)) 356 }