1use std::{
4 collections::{BTreeSet, HashMap},
5 sync::{
6 Arc, Mutex, MutexGuard, Weak,
7 atomic::{AtomicBool, AtomicU8, Ordering},
8 },
9};
10
11use bytes::{Bytes, BytesMut};
12use futures::{
13 FutureExt,
14 channel::oneshot,
15 future::{BoxFuture, Either, Shared, select, try_join_all},
16};
17
18use crate::{
19 Alignment, ByteRange, CloudBackendError, CloudBackendErrorKind, CloudContentBackend,
20 CloudFilesCoreError, CloudItemKey, ContentReadRange, ContentReadRequest, ContentReadResponse,
21 ContentRevision, SessionGeneration,
22};
23
24pub type HydrationResult<T> = std::result::Result<T, HydrationError>;
26
27#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
29pub enum HydrationError {
30 #[error("hydration waiter was cancelled")]
32 Cancelled,
33 #[error(transparent)]
35 Backend(#[from] CloudBackendError),
36 #[error(transparent)]
38 Contract(#[from] CloudFilesCoreError),
39 #[error("hydration in-flight work limit exceeded ({scope})")]
41 InFlightLimitExceeded {
42 scope: &'static str,
44 },
45}
46
47#[derive(Debug, Clone, Copy, PartialEq, Eq)]
49pub struct HydrationLimits {
50 pub max_in_flight_work: usize,
52 pub max_in_flight_per_content: usize,
54}
55
56impl HydrationLimits {
57 pub const DEFAULT: Self = Self {
59 max_in_flight_work: 1_024,
60 max_in_flight_per_content: 128,
61 };
62}
63
64#[derive(Debug, Clone, PartialEq, Eq)]
66pub struct HydrationRequest {
67 read: ContentReadRequest,
68 alignment: Alignment,
69 session_generation: SessionGeneration,
70}
71
72impl HydrationRequest {
73 #[must_use]
75 pub const fn new(
76 read: ContentReadRequest,
77 alignment: Alignment,
78 session_generation: SessionGeneration,
79 ) -> Self {
80 Self {
81 read,
82 alignment,
83 session_generation,
84 }
85 }
86
87 #[must_use]
89 pub const fn whole(
90 key: CloudItemKey,
91 revision: ContentRevision,
92 expected_size: u64,
93 session_generation: SessionGeneration,
94 ) -> Self {
95 Self::new(
96 ContentReadRequest::whole(key, revision, expected_size),
97 Alignment::ONE,
98 session_generation,
99 )
100 }
101
102 #[must_use]
104 pub const fn range(
105 key: CloudItemKey,
106 revision: ContentRevision,
107 expected_size: u64,
108 range: ByteRange,
109 alignment: Alignment,
110 session_generation: SessionGeneration,
111 ) -> Self {
112 Self::new(
113 ContentReadRequest::range(key, revision, expected_size, range),
114 alignment,
115 session_generation,
116 )
117 }
118
119 #[must_use]
121 pub const fn read(&self) -> &ContentReadRequest {
122 &self.read
123 }
124
125 #[must_use]
127 pub const fn alignment(&self) -> Alignment {
128 self.alignment
129 }
130
131 #[must_use]
133 pub const fn session_generation(&self) -> SessionGeneration {
134 self.session_generation
135 }
136}
137
138#[derive(Debug, Clone, PartialEq, Eq, Hash)]
139struct HydrationContentKey {
140 key: CloudItemKey,
141 revision: ContentRevision,
142 expected_size: u64,
143 session_generation: SessionGeneration,
144 alignment: Alignment,
145}
146
147impl HydrationContentKey {
148 fn from_request(request: &HydrationRequest) -> Self {
149 Self {
150 key: request.read().key().clone(),
151 revision: request.read().revision().clone(),
152 expected_size: request.read().expected_size(),
153 session_generation: request.session_generation(),
154 alignment: request.alignment(),
155 }
156 }
157}
158
159type SharedRead = Shared<BoxFuture<'static, HydrationResult<ContentReadResponse>>>;
160
161#[derive(Clone)]
162struct WorkDependency {
163 id: u64,
164 start: u64,
165 end: u64,
166 future: SharedRead,
167}
168
169struct WorkEntry {
170 content: HydrationContentKey,
171 start: u64,
172 end: u64,
173 whole: bool,
174 future: SharedRead,
175 waiter_count: usize,
176}
177
178#[derive(Default)]
179struct CoordinatorState {
180 next_work_id: u64,
181 work: HashMap<u64, WorkEntry>,
182 by_content: HashMap<HydrationContentKey, BTreeSet<(u64, u64, u64)>>,
183}
184
185fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
186 match mutex.lock() {
187 Ok(guard) => guard,
188 Err(poisoned) => poisoned.into_inner(),
189 }
190}
191
192#[derive(Clone)]
199pub struct HydrationCoordinator {
200 backend: Arc<dyn CloudContentBackend>,
201 state: Arc<Mutex<CoordinatorState>>,
202 limits: HydrationLimits,
203}
204
205impl HydrationCoordinator {
206 pub fn new(backend: Arc<dyn CloudContentBackend>) -> Self {
208 Self::with_limits(backend, HydrationLimits::DEFAULT)
209 }
210
211 pub fn with_limits(backend: Arc<dyn CloudContentBackend>, limits: HydrationLimits) -> Self {
213 Self {
214 backend,
215 state: Arc::new(Mutex::new(CoordinatorState::default())),
216 limits,
217 }
218 }
219
220 #[must_use]
222 pub const fn limits(&self) -> HydrationLimits {
223 self.limits
224 }
225
226 pub fn request(&self, request: HydrationRequest) -> HydrationResult<HydrationWaiter> {
236 let content = HydrationContentKey::from_request(&request);
237 let (target_start, target_end) = logical_extent(request.read())?;
238 let dependencies = match request.read().read_range() {
239 ContentReadRange::Whole => self.register_whole(&content, request.read())?,
240 ContentReadRange::Range(_) if target_start == target_end => {
241 self.register_empty_range(&content, request.read(), target_start)?
242 }
243 ContentReadRange::Range(_) => {
244 let (physical_start, physical_end) = aligned_extent(
245 target_start,
246 target_end,
247 content.expected_size,
248 request.alignment(),
249 )?;
250 self.register_range(&content, request.read(), physical_start, physical_end)?
251 }
252 };
253
254 let work_ids = dependencies
255 .iter()
256 .map(|dependency| dependency.id)
257 .collect();
258 let (cancel_tx, cancel_rx) = oneshot::channel();
259 let inner = Arc::new(WaiterInner {
260 terminal: AtomicU8::new(TERMINAL_PENDING),
261 released: AtomicBool::new(false),
262 coordinator: Arc::downgrade(&self.state),
263 work_ids,
264 cancel_tx: Mutex::new(Some(cancel_tx)),
265 });
266 Ok(HydrationWaiter {
267 request,
268 target_start,
269 target_end,
270 dependencies,
271 cancel_rx,
272 inner,
273 finished: false,
274 })
275 }
276
277 fn register_whole(
278 &self,
279 content: &HydrationContentKey,
280 request: &ContentReadRequest,
281 ) -> HydrationResult<Vec<WorkDependency>> {
282 let mut state = lock(&self.state);
283 if let Some((id, entry)) = state
284 .work
285 .iter_mut()
286 .find(|(_, entry)| entry.content == *content && entry.whole)
287 {
288 entry.waiter_count = entry.waiter_count.saturating_add(1);
289 return Ok(vec![WorkDependency {
290 id: *id,
291 start: entry.start,
292 end: entry.end,
293 future: entry.future.clone(),
294 }]);
295 }
296
297 ensure_capacity(&state, content, 1, self.limits)?;
298 let dependency = insert_work(
299 &mut state,
300 Arc::clone(&self.backend),
301 content.clone(),
302 request,
303 0,
304 content.expected_size,
305 true,
306 )?;
307 if let Some(entry) = state.work.get_mut(&dependency.id) {
308 entry.waiter_count = 1;
309 }
310 Ok(vec![dependency])
311 }
312
313 fn register_empty_range(
314 &self,
315 content: &HydrationContentKey,
316 request: &ContentReadRequest,
317 offset: u64,
318 ) -> HydrationResult<Vec<WorkDependency>> {
319 let mut state = lock(&self.state);
320 if let Some((id, entry)) = state.work.iter_mut().find(|(_, entry)| {
321 entry.content == *content
322 && !entry.whole
323 && entry.start == offset
324 && entry.end == offset
325 }) {
326 entry.waiter_count = entry.waiter_count.saturating_add(1);
327 return Ok(vec![WorkDependency {
328 id: *id,
329 start: entry.start,
330 end: entry.end,
331 future: entry.future.clone(),
332 }]);
333 }
334
335 ensure_capacity(&state, content, 1, self.limits)?;
336 let dependency = insert_work(
337 &mut state,
338 Arc::clone(&self.backend),
339 content.clone(),
340 request,
341 offset,
342 offset,
343 false,
344 )?;
345 if let Some(entry) = state.work.get_mut(&dependency.id) {
346 entry.waiter_count = 1;
347 }
348 Ok(vec![dependency])
349 }
350
351 fn register_range(
352 &self,
353 content: &HydrationContentKey,
354 original_request: &ContentReadRequest,
355 physical_start: u64,
356 physical_end: u64,
357 ) -> HydrationResult<Vec<WorkDependency>> {
358 let mut state = lock(&self.state);
359 let mut dependencies = state
360 .by_content
361 .get(content)
362 .into_iter()
363 .flat_map(|entries| entries.iter())
364 .take_while(|(start, _, _)| *start < physical_end)
365 .filter_map(|(start, end, id)| {
366 let entry = state.work.get(id)?;
367 if entry.whole || ranges_overlap(physical_start, physical_end, *start, *end) {
368 Some(WorkDependency {
369 id: *id,
370 start: *start,
371 end: *end,
372 future: entry.future.clone(),
373 })
374 } else {
375 None
376 }
377 })
378 .collect::<Vec<_>>();
379 dependencies.sort_by_key(|dependency| (dependency.start, dependency.end));
380
381 let mut cursor = physical_start;
382 let existing = dependencies.clone();
383 let mut gaps = Vec::new();
384 for dependency in existing {
385 let gap_end = dependency.start.min(physical_end);
386 if cursor < gap_end {
387 gaps.push((cursor, gap_end));
388 }
389 cursor = cursor.max(dependency.end.min(physical_end));
390 }
391 if cursor < physical_end {
392 gaps.push((cursor, physical_end));
393 }
394 ensure_capacity(&state, content, gaps.len(), self.limits)?;
395 for (start, end) in gaps {
396 dependencies.push(insert_range_work(
397 &mut state,
398 Arc::clone(&self.backend),
399 content.clone(),
400 original_request,
401 start,
402 end,
403 )?);
404 }
405
406 dependencies.sort_by_key(|dependency| (dependency.start, dependency.end));
407 for dependency in &dependencies {
408 if let Some(entry) = state.work.get_mut(&dependency.id) {
409 entry.waiter_count = entry.waiter_count.saturating_add(1);
410 }
411 }
412 Ok(dependencies)
413 }
414}
415
416fn insert_range_work(
417 state: &mut CoordinatorState,
418 backend: Arc<dyn CloudContentBackend>,
419 content: HydrationContentKey,
420 original_request: &ContentReadRequest,
421 start: u64,
422 end: u64,
423) -> HydrationResult<WorkDependency> {
424 let range = ByteRange::new(start, end.saturating_sub(start))?;
425 let request = ContentReadRequest::range(
426 original_request.key().clone(),
427 original_request.revision().clone(),
428 original_request.expected_size(),
429 range,
430 );
431 insert_work(state, backend, content, &request, start, end, false)
432}
433
434fn insert_work(
435 state: &mut CoordinatorState,
436 backend: Arc<dyn CloudContentBackend>,
437 content: HydrationContentKey,
438 request: &ContentReadRequest,
439 start: u64,
440 end: u64,
441 whole: bool,
442) -> HydrationResult<WorkDependency> {
443 let id = state.next_work_id;
444 state.next_work_id = state.next_work_id.checked_add(1).ok_or_else(|| {
445 CloudFilesCoreError::invalid_content_response("hydration work identity exhausted")
446 })?;
447 let future_request = request.clone();
448 let future = async move {
449 let response = backend.read_content(&future_request).await?;
450 future_request.validate_response(&response)?;
451 Ok(response)
452 }
453 .boxed()
454 .shared();
455 let index_content = content.clone();
456 state.work.insert(
457 id,
458 WorkEntry {
459 content,
460 start,
461 end,
462 whole,
463 future: future.clone(),
464 waiter_count: 0,
465 },
466 );
467 state
468 .by_content
469 .entry(index_content)
470 .or_default()
471 .insert((start, end, id));
472 Ok(WorkDependency {
473 id,
474 start,
475 end,
476 future,
477 })
478}
479
480fn ensure_capacity(
481 state: &CoordinatorState,
482 content: &HydrationContentKey,
483 additional: usize,
484 limits: HydrationLimits,
485) -> HydrationResult<()> {
486 if limits.max_in_flight_work == 0
487 || state.work.len().saturating_add(additional) > limits.max_in_flight_work
488 {
489 return Err(HydrationError::InFlightLimitExceeded { scope: "global" });
490 }
491 let per_content = state.by_content.get(content).map_or(0, BTreeSet::len);
492 if limits.max_in_flight_per_content == 0
493 || per_content.saturating_add(additional) > limits.max_in_flight_per_content
494 {
495 return Err(HydrationError::InFlightLimitExceeded { scope: "content" });
496 }
497 Ok(())
498}
499
500fn logical_extent(request: &ContentReadRequest) -> HydrationResult<(u64, u64)> {
501 match request.read_range() {
502 ContentReadRange::Whole => Ok((0, request.expected_size())),
503 ContentReadRange::Range(range) => {
504 if range.offset() > request.expected_size() {
505 return Err(CloudBackendError::new(CloudBackendErrorKind::InvalidRequest).into());
506 }
507 Ok((
508 range.offset(),
509 range.end_exclusive().min(request.expected_size()),
510 ))
511 }
512 }
513}
514
515fn aligned_extent(
516 start: u64,
517 end: u64,
518 total_size: u64,
519 alignment: Alignment,
520) -> HydrationResult<(u64, u64)> {
521 let alignment = alignment.get();
522 let aligned_start = start / alignment * alignment;
523 if end == total_size {
524 return Ok((aligned_start, total_size));
525 }
526 let remainder = end % alignment;
527 let aligned_end = if remainder == 0 {
528 end
529 } else {
530 end.checked_add(alignment - remainder).ok_or_else(|| {
531 CloudFilesCoreError::invalid_byte_range("aligned hydration range exceeds u64")
532 })?
533 };
534 Ok((aligned_start, aligned_end.min(total_size)))
535}
536
537const fn ranges_overlap(left_start: u64, left_end: u64, right_start: u64, right_end: u64) -> bool {
538 left_start < right_end && right_start < left_end
539}
540
541const TERMINAL_PENDING: u8 = 0;
542const TERMINAL_COMPLETED: u8 = 1;
543const TERMINAL_CANCELLED: u8 = 2;
544
545struct WaiterInner {
546 terminal: AtomicU8,
547 released: AtomicBool,
548 coordinator: Weak<Mutex<CoordinatorState>>,
549 work_ids: Vec<u64>,
550 cancel_tx: Mutex<Option<oneshot::Sender<()>>>,
551}
552
553impl WaiterInner {
554 fn cancel(&self) -> HydrationCancellationOutcome {
555 match self.terminal.compare_exchange(
556 TERMINAL_PENDING,
557 TERMINAL_CANCELLED,
558 Ordering::AcqRel,
559 Ordering::Acquire,
560 ) {
561 Ok(_) => {
562 if let Some(sender) = lock(&self.cancel_tx).take() {
563 let _ = sender.send(());
564 }
565 self.release_work();
566 HydrationCancellationOutcome::Cancelled
567 }
568 Err(TERMINAL_CANCELLED) => HydrationCancellationOutcome::AlreadyCancelled,
569 Err(_) => HydrationCancellationOutcome::AlreadyCompleted,
570 }
571 }
572
573 fn complete(&self) -> bool {
574 if self
575 .terminal
576 .compare_exchange(
577 TERMINAL_PENDING,
578 TERMINAL_COMPLETED,
579 Ordering::AcqRel,
580 Ordering::Acquire,
581 )
582 .is_ok()
583 {
584 lock(&self.cancel_tx).take();
585 self.release_work();
586 true
587 } else {
588 false
589 }
590 }
591
592 fn release_work(&self) {
593 if self.released.swap(true, Ordering::AcqRel) {
594 return;
595 }
596 let Some(coordinator) = self.coordinator.upgrade() else {
597 return;
598 };
599 let mut state = lock(&coordinator);
600 for id in &self.work_ids {
601 let remove = if let Some(entry) = state.work.get_mut(id) {
602 entry.waiter_count = entry.waiter_count.saturating_sub(1);
603 entry.waiter_count == 0
604 } else {
605 false
606 };
607 if remove && let Some(entry) = state.work.remove(id) {
608 let empty = if let Some(index) = state.by_content.get_mut(&entry.content) {
609 index.remove(&(entry.start, entry.end, *id));
610 index.is_empty()
611 } else {
612 false
613 };
614 if empty {
615 state.by_content.remove(&entry.content);
616 }
617 }
618 }
619 }
620}
621
622#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
624pub enum HydrationCancellationOutcome {
625 Cancelled,
627 AlreadyCancelled,
629 AlreadyCompleted,
631}
632
633#[derive(Clone)]
635pub struct HydrationCancellationHandle {
636 inner: Arc<WaiterInner>,
637}
638
639impl HydrationCancellationHandle {
640 #[must_use]
642 pub fn cancel(&self) -> HydrationCancellationOutcome {
643 self.inner.cancel()
644 }
645}
646
647pub struct HydrationWaiter {
649 request: HydrationRequest,
650 target_start: u64,
651 target_end: u64,
652 dependencies: Vec<WorkDependency>,
653 cancel_rx: oneshot::Receiver<()>,
654 inner: Arc<WaiterInner>,
655 finished: bool,
656}
657
658impl HydrationWaiter {
659 #[must_use]
661 pub fn cancellation_handle(&self) -> HydrationCancellationHandle {
662 HydrationCancellationHandle {
663 inner: Arc::clone(&self.inner),
664 }
665 }
666
667 pub async fn wait(mut self) -> HydrationResult<ContentReadResponse> {
673 if self.inner.terminal.load(Ordering::Acquire) == TERMINAL_CANCELLED {
674 self.finished = true;
675 return Err(HydrationError::Cancelled);
676 }
677
678 let futures = self
679 .dependencies
680 .iter()
681 .map(|dependency| dependency.future.clone());
682 let selected = {
683 let work = try_join_all(futures).boxed();
684 let cancelled = (&mut self.cancel_rx).map(|_| ()).boxed();
685 match select(work, cancelled).await {
686 Either::Left((result, _)) => Some(result),
687 Either::Right(((), _)) => None,
688 }
689 };
690 let result = match selected {
691 Some(result) => {
692 if self.inner.complete() {
693 result.and_then(|responses| self.assemble(responses))
694 } else {
695 Err(HydrationError::Cancelled)
696 }
697 }
698 None => Err(HydrationError::Cancelled),
699 };
700 self.finished = true;
701 result
702 }
703
704 fn assemble(
705 &self,
706 responses: Vec<ContentReadResponse>,
707 ) -> HydrationResult<ContentReadResponse> {
708 if self.target_start == self.target_end {
709 return ContentReadResponse::new(
710 self.request.read().revision().clone(),
711 self.target_start,
712 Bytes::new(),
713 self.request.read().expected_size(),
714 )
715 .map_err(Into::into);
716 }
717 if responses.len() != self.dependencies.len() {
718 return Err(CloudFilesCoreError::invalid_content_response(
719 "hydration dependency result count changed",
720 )
721 .into());
722 }
723
724 let mut cursor = self.target_start;
725 let mut slices = Vec::new();
726 for (dependency, response) in self.dependencies.iter().zip(responses) {
727 if dependency.end <= cursor || dependency.start >= self.target_end {
728 continue;
729 }
730 let copy_start = cursor.max(dependency.start);
731 let copy_end = self.target_end.min(dependency.end);
732 let relative_start = copy_start.saturating_sub(response.offset());
733 let relative_end = copy_end.saturating_sub(response.offset());
734 if relative_end > response.bytes().len() as u64 || relative_start > relative_end {
735 return Err(CloudFilesCoreError::invalid_content_response(
736 "hydration dependency did not cover its claimed logical range",
737 )
738 .into());
739 }
740 #[expect(
742 clippy::cast_possible_truncation,
743 reason = "both offsets were bounded by the usize-sized response buffer"
744 )]
745 let (relative_start, relative_end) = (relative_start as usize, relative_end as usize);
746 slices.push(response.bytes().slice(relative_start..relative_end));
747 cursor = copy_end;
748 if cursor == self.target_end {
749 break;
750 }
751 }
752 if cursor != self.target_end {
753 return Err(CloudFilesCoreError::invalid_content_response(
754 "hydration dependencies left an uncovered logical range",
755 )
756 .into());
757 }
758
759 let bytes = combine_slices(slices)?;
760 let revision = self.request.read().revision().clone();
761 let offset = self.target_start;
762 let total_size = self.request.read().expected_size();
763 let response = ContentReadResponse::new(revision, offset, bytes, total_size)?;
764 self.request.read().validate_response(&response)?;
765 Ok(response)
766 }
767}
768
769impl Drop for HydrationWaiter {
770 fn drop(&mut self) {
771 if !self.finished {
772 self.inner.cancel();
773 }
774 }
775}
776
777fn combine_slices(mut slices: Vec<Bytes>) -> HydrationResult<Bytes> {
778 match slices.len() {
779 0 => Ok(Bytes::new()),
780 1 => Ok(slices.remove(0)),
781 _ => {
782 let total_len = checked_slice_total(slices.iter().map(Bytes::len))?;
783 let mut combined = BytesMut::with_capacity(total_len);
784 for bytes in slices {
785 combined.extend_from_slice(&bytes);
786 }
787 Ok(combined.freeze())
788 }
789 }
790}
791
792fn checked_slice_total(lengths: impl IntoIterator<Item = usize>) -> HydrationResult<usize> {
793 lengths.into_iter().try_fold(0usize, |total, length| {
794 total.checked_add(length).ok_or_else(|| {
795 CloudFilesCoreError::invalid_content_response("hydration result length exceeds usize")
796 .into()
797 })
798 })
799}
800
801#[cfg(test)]
802mod tests {
803 use super::*;
804
805 struct TestBackend;
806
807 #[async_trait::async_trait]
808 impl CloudContentBackend for TestBackend {
809 async fn read_content(
810 &self,
811 _request: &ContentReadRequest,
812 ) -> crate::BackendResult<ContentReadResponse> {
813 Err(CloudBackendError::new(
814 CloudBackendErrorKind::InvalidRequest,
815 ))
816 }
817 }
818
819 fn request(range: ContentReadRange) -> HydrationRequest {
820 let scope = crate::CloudScope::new(
821 crate::CloudNamespaceId::new("namespace").expect("namespace fixture should be valid"),
822 crate::CloudRootId::new("root").expect("root fixture should be valid"),
823 );
824 let key = CloudItemKey::new(
825 scope,
826 crate::CloudItemId::new("item").expect("item fixture should be valid"),
827 );
828 let revision = ContentRevision::from_slice(b"content-v1")
829 .expect("content revision fixture should be valid");
830 let read = match range {
831 ContentReadRange::Whole => ContentReadRequest::whole(key, revision, 4),
832 ContentReadRange::Range(range) => ContentReadRequest::range(key, revision, 4, range),
833 };
834 HydrationRequest::new(
835 read,
836 Alignment::ONE,
837 SessionGeneration::new(1).expect("generation fixture should be valid"),
838 )
839 }
840
841 fn response(offset: u64, bytes: &'static [u8]) -> ContentReadResponse {
842 ContentReadResponse::new(
843 ContentRevision::from_slice(b"content-v1")
844 .expect("content revision fixture should be valid"),
845 offset,
846 Bytes::from_static(bytes),
847 4,
848 )
849 .expect("response fixture should be valid")
850 }
851
852 fn dependency(id: u64, start: u64, end: u64) -> WorkDependency {
853 WorkDependency {
854 id,
855 start,
856 end,
857 future: futures::future::ready(Err(HydrationError::Cancelled))
858 .boxed()
859 .shared(),
860 }
861 }
862
863 fn waiter(
864 request: HydrationRequest,
865 target_start: u64,
866 target_end: u64,
867 dependencies: Vec<WorkDependency>,
868 ) -> HydrationWaiter {
869 let (_, cancel_rx) = oneshot::channel();
870 HydrationWaiter {
871 request,
872 target_start,
873 target_end,
874 dependencies,
875 cancel_rx,
876 inner: Arc::new(WaiterInner {
877 terminal: AtomicU8::new(TERMINAL_PENDING),
878 released: AtomicBool::new(false),
879 coordinator: Weak::new(),
880 work_ids: Vec::new(),
881 cancel_tx: Mutex::new(None),
882 }),
883 finished: false,
884 }
885 }
886
887 #[test]
888 fn poisoned_coordinator_mutex_recovers_its_inner_state() {
889 let mutex = Arc::new(Mutex::new(7usize));
890 let poison_target = Arc::clone(&mutex);
891 let _ = std::thread::spawn(move || {
892 let _guard = poison_target.lock().expect("mutex should initially lock");
893 panic!("poison coordinator fixture");
894 })
895 .join();
896
897 assert_eq!(*lock(&mutex), 7);
898 }
899
900 #[test]
901 fn work_identity_and_alignment_overflow_are_contract_failures() {
902 let mut state = CoordinatorState {
903 next_work_id: u64::MAX,
904 ..CoordinatorState::default()
905 };
906 let scope = crate::CloudScope::new(
907 crate::CloudNamespaceId::new("namespace").expect("namespace fixture should be valid"),
908 crate::CloudRootId::new("root").expect("root fixture should be valid"),
909 );
910 let key = CloudItemKey::new(
911 scope,
912 crate::CloudItemId::new("item").expect("item fixture should be valid"),
913 );
914 let revision = ContentRevision::from_slice(b"content-v1")
915 .expect("content revision fixture should be valid");
916 let request = HydrationRequest::whole(
917 key,
918 revision,
919 1,
920 SessionGeneration::new(1).expect("generation fixture should be valid"),
921 );
922 let content = HydrationContentKey::from_request(&request);
923 assert!(futures::executor::block_on(TestBackend.read_content(request.read())).is_err());
924
925 assert!(matches!(
926 insert_work(
927 &mut state,
928 Arc::new(TestBackend),
929 content,
930 request.read(),
931 0,
932 1,
933 true,
934 ),
935 Err(HydrationError::Contract(
936 CloudFilesCoreError::InvalidContentResponse { .. }
937 ))
938 ));
939 assert!(matches!(
940 aligned_extent(
941 u64::MAX - 1,
942 u64::MAX - 1,
943 u64::MAX,
944 Alignment::new(8).expect("alignment fixture should be valid"),
945 ),
946 Err(HydrationError::Contract(
947 CloudFilesCoreError::InvalidByteRange { .. }
948 ))
949 ));
950 }
951
952 #[test]
953 fn coordinator_propagates_work_identity_and_alignment_failures_from_each_request_shape() {
954 let coordinator = HydrationCoordinator {
955 backend: Arc::new(TestBackend),
956 state: Arc::new(Mutex::new(CoordinatorState {
957 next_work_id: u64::MAX,
958 ..CoordinatorState::default()
959 })),
960 limits: HydrationLimits::DEFAULT,
961 };
962
963 let whole = request(ContentReadRange::Whole);
964 assert!(matches!(
965 coordinator.request(whole),
966 Err(HydrationError::Contract(
967 CloudFilesCoreError::InvalidContentResponse { .. }
968 ))
969 ));
970
971 let empty = request(ContentReadRange::Range(
972 ByteRange::new(4, 1).expect("range fixture should be valid"),
973 ));
974 assert!(matches!(
975 coordinator.request(empty),
976 Err(HydrationError::Contract(
977 CloudFilesCoreError::InvalidContentResponse { .. }
978 ))
979 ));
980
981 let range = request(ContentReadRange::Range(
982 ByteRange::new(0, 1).expect("range fixture should be valid"),
983 ));
984 assert!(matches!(
985 coordinator.request(range),
986 Err(HydrationError::Contract(
987 CloudFilesCoreError::InvalidContentResponse { .. }
988 ))
989 ));
990
991 let overflow = HydrationRequest::range(
992 request(ContentReadRange::Whole).read().key().clone(),
993 ContentRevision::from_slice(b"content-v1")
994 .expect("content revision fixture should be valid"),
995 u64::MAX,
996 ByteRange::new(u64::MAX - 2, 1).expect("range fixture should be valid"),
997 Alignment::new(8).expect("alignment fixture should be valid"),
998 SessionGeneration::new(1).expect("generation fixture should be valid"),
999 );
1000 assert!(matches!(
1001 coordinator.request(overflow),
1002 Err(HydrationError::Contract(
1003 CloudFilesCoreError::InvalidByteRange { .. }
1004 ))
1005 ));
1006 }
1007
1008 #[test]
1009 fn release_work_is_idempotent_and_tolerates_missing_internal_indexes() {
1010 let empty_state = Arc::new(Mutex::new(CoordinatorState::default()));
1011 let repeated = WaiterInner {
1012 terminal: AtomicU8::new(TERMINAL_PENDING),
1013 released: AtomicBool::new(false),
1014 coordinator: Arc::downgrade(&empty_state),
1015 work_ids: Vec::new(),
1016 cancel_tx: Mutex::new(None),
1017 };
1018 repeated.release_work();
1019 repeated.release_work();
1020
1021 let missing_work = WaiterInner {
1022 terminal: AtomicU8::new(TERMINAL_PENDING),
1023 released: AtomicBool::new(false),
1024 coordinator: Arc::downgrade(&empty_state),
1025 work_ids: vec![7],
1026 cancel_tx: Mutex::new(None),
1027 };
1028 missing_work.release_work();
1029
1030 let request = request(ContentReadRange::Whole);
1031 let content = HydrationContentKey::from_request(&request);
1032 let mut state = CoordinatorState::default();
1033 state.work.insert(
1034 9,
1035 WorkEntry {
1036 content,
1037 start: 0,
1038 end: 4,
1039 whole: true,
1040 future: futures::future::ready(Err(HydrationError::Cancelled))
1041 .boxed()
1042 .shared(),
1043 waiter_count: 1,
1044 },
1045 );
1046 let missing_index_state = Arc::new(Mutex::new(state));
1047 let missing_index = WaiterInner {
1048 terminal: AtomicU8::new(TERMINAL_PENDING),
1049 released: AtomicBool::new(false),
1050 coordinator: Arc::downgrade(&missing_index_state),
1051 work_ids: vec![9],
1052 cancel_tx: Mutex::new(None),
1053 };
1054 missing_index.release_work();
1055 assert!(lock(&missing_index_state).work.is_empty());
1056 }
1057
1058 #[test]
1059 fn assembly_rejects_changed_missing_and_short_dependencies() {
1060 let request = request(ContentReadRange::Range(
1061 ByteRange::new(0, 4).expect("range fixture should be valid"),
1062 ));
1063
1064 let changed = waiter(request.clone(), 0, 4, vec![dependency(1, 0, 4)]);
1065 assert!(matches!(
1066 changed.assemble(Vec::new()),
1067 Err(HydrationError::Contract(
1068 CloudFilesCoreError::InvalidContentResponse { .. }
1069 ))
1070 ));
1071
1072 let missing = waiter(request.clone(), 0, 4, vec![dependency(2, 4, 4)]);
1073 assert!(matches!(
1074 missing.assemble(vec![response(4, b"")]),
1075 Err(HydrationError::Contract(
1076 CloudFilesCoreError::InvalidContentResponse { .. }
1077 ))
1078 ));
1079
1080 let short = waiter(request.clone(), 0, 4, vec![dependency(3, 0, 4)]);
1081 assert!(matches!(
1082 short.assemble(vec![response(0, b"ab")]),
1083 Err(HydrationError::Contract(
1084 CloudFilesCoreError::InvalidContentResponse { .. }
1085 ))
1086 ));
1087
1088 let complete = waiter(request, 0, 4, vec![dependency(4, 0, 4)]);
1089 assert_eq!(
1090 complete
1091 .assemble(vec![response(0, b"abcd")])
1092 .expect("complete dependency should assemble")
1093 .bytes()
1094 .as_ref(),
1095 b"abcd"
1096 );
1097 assert!(
1098 combine_slices(Vec::new())
1099 .expect("empty slices should combine")
1100 .is_empty()
1101 );
1102 }
1103
1104 #[test]
1105 fn slice_length_accumulator_rejects_usize_overflow() {
1106 assert!(matches!(
1107 checked_slice_total([usize::MAX, 1]),
1108 Err(HydrationError::Contract(
1109 CloudFilesCoreError::InvalidContentResponse { .. }
1110 ))
1111 ));
1112 }
1113}