aster_forge_cloud_files_core/
hydration.rs

1//! Runtime-neutral hydration work sharing, revision fencing, and waiter-scoped cancellation.
2
3use 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
24/// Result returned by hydration coordination and waiter completion.
25pub type HydrationResult<T> = std::result::Result<T, HydrationError>;
26
27/// Product-neutral hydration coordination failure.
28#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
29pub enum HydrationError {
30    /// The individual waiter was cancelled before it won the completion race.
31    #[error("hydration waiter was cancelled")]
32    Cancelled,
33    /// The product-owned content backend failed the revision-bound read.
34    #[error(transparent)]
35    Backend(#[from] CloudBackendError),
36    /// Backend bytes violated the requested revision, range, or size contract.
37    #[error(transparent)]
38    Contract(#[from] CloudFilesCoreError),
39    /// The coordinator reached its configured global or per-content work limit.
40    #[error("hydration in-flight work limit exceeded ({scope})")]
41    InFlightLimitExceeded {
42        /// Whether the global coordinator or one content key reached its limit.
43        scope: &'static str,
44    },
45}
46
47/// Bounded hydration work limits.
48#[derive(Debug, Clone, Copy, PartialEq, Eq)]
49pub struct HydrationLimits {
50    /// Maximum number of backend work items retained by one coordinator.
51    pub max_in_flight_work: usize,
52    /// Maximum number of backend work items retained for one content key.
53    pub max_in_flight_per_content: usize,
54}
55
56impl HydrationLimits {
57    /// Conservative defaults suitable for a provider process.
58    pub const DEFAULT: Self = Self {
59        max_in_flight_work: 1_024,
60        max_in_flight_per_content: 128,
61    };
62}
63
64/// One logical hydration request plus the physical alignment required by its adapter boundary.
65#[derive(Debug, Clone, PartialEq, Eq)]
66pub struct HydrationRequest {
67    read: ContentReadRequest,
68    alignment: Alignment,
69    session_generation: SessionGeneration,
70}
71
72impl HydrationRequest {
73    /// Creates a hydration request from an exact content read and physical transfer alignment.
74    #[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    /// Creates a complete-file hydration request.
88    #[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    /// Creates a range hydration request with adapter-owned physical alignment.
103    #[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    /// Returns the exact logical content request.
120    #[must_use]
121    pub const fn read(&self) -> &ContentReadRequest {
122        &self.read
123    }
124
125    /// Returns the physical range alignment owned by the active adapter boundary.
126    #[must_use]
127    pub const fn alignment(&self) -> Alignment {
128        self.alignment
129    }
130
131    /// Returns the platform session generation that owns this waiter and its completion.
132    #[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/// Runtime-neutral coordinator for deduplicated revision-bound content reads.
193///
194/// The coordinator does not spawn tasks and does not own physical materialization. Returned
195/// waiters drive shared backend futures on the caller's executor. When the last waiter releases a
196/// work item, dropping the backend future provides the narrow cancellation boundary available to
197/// the backend adapter.
198#[derive(Clone)]
199pub struct HydrationCoordinator {
200    backend: Arc<dyn CloudContentBackend>,
201    state: Arc<Mutex<CoordinatorState>>,
202    limits: HydrationLimits,
203}
204
205impl HydrationCoordinator {
206    /// Creates an empty coordinator scoped to one product-owned content backend adapter.
207    pub fn new(backend: Arc<dyn CloudContentBackend>) -> Self {
208        Self::with_limits(backend, HydrationLimits::DEFAULT)
209    }
210
211    /// Creates a coordinator with explicit global and per-content backpressure limits.
212    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    /// Returns the configured in-flight work limits.
221    #[must_use]
222    pub const fn limits(&self) -> HydrationLimits {
223        self.limits
224    }
225
226    /// Registers one waiter and returns immediately without spawning backend work.
227    ///
228    /// Exact requests share one backend future. Overlapping range requests reuse existing work and
229    /// add backend reads only for uncovered physical gaps. All sharing keys include item scope,
230    /// content revision, and expected size.
231    /// # Errors
232    ///
233    /// Returns an error when validation fails or an underlying backend, store, or platform
234    /// operation fails.
235    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/// Result of attempting to cancel one hydration waiter.
623#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
624pub enum HydrationCancellationOutcome {
625    /// Cancellation won the terminal race for this waiter.
626    Cancelled,
627    /// This waiter had already been cancelled.
628    AlreadyCancelled,
629    /// Completion won before cancellation was applied.
630    AlreadyCompleted,
631}
632
633/// Cloneable cancellation ownership for exactly one hydration waiter.
634#[derive(Clone)]
635pub struct HydrationCancellationHandle {
636    inner: Arc<WaiterInner>,
637}
638
639impl HydrationCancellationHandle {
640    /// Attempts to cancel this waiter without cancelling work still owned by other waiters.
641    #[must_use]
642    pub fn cancel(&self) -> HydrationCancellationOutcome {
643        self.inner.cancel()
644    }
645}
646
647/// One registered hydration waiter that resolves to exactly its requested logical bytes.
648pub 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    /// Returns an independently cloneable cancellation handle for this waiter.
660    #[must_use]
661    pub fn cancellation_handle(&self) -> HydrationCancellationHandle {
662        HydrationCancellationHandle {
663            inner: Arc::clone(&self.inner),
664        }
665    }
666
667    /// Drives shared backend work and returns only the exact logical bytes requested by this waiter.
668    /// # Errors
669    ///
670    /// Returns an error when validation fails or an underlying backend, store, or platform
671    /// operation fails.
672    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            // Both offsets were checked against `Bytes::len()`, so conversion cannot truncate.
741            #[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}