Skip to main content

xet_data/processing/
range_upload.rs

1use std::ops::Range;
2use std::pin::Pin;
3use std::sync::Arc;
4
5use tokio::io::{AsyncRead, AsyncReadExt};
6use tracing::{debug, info};
7use xet_client::cas_client::Client;
8use xet_client::cas_types::{FileChunkHashesResponse, FileRange, HexMerkleHash};
9use xet_core_structures::merklehash::{ChunkHashList, MerkleHash, MerkleHashSubtree};
10use xet_core_structures::metadata_shard::file_structs::{
11    FileDataSequenceEntry, FileDataSequenceHeader, FileVerificationEntry, MDBFileInfo,
12};
13use xet_runtime::core::XetContext;
14
15use super::XetFileInfo;
16use super::configurations::TranslatorConfig;
17use super::file_cleaner::Sha256Policy;
18use super::file_upload_session::FileUploadSession;
19use crate::error::{DataError, Result};
20use crate::file_reconstruction::FileReconstructor;
21
22/// A single edit applied to the original file: replace `original_range` with `new_length`
23/// bytes from `reader`.
24///
25/// All three combinations are supported:
26/// - `original_range.len() == new_length` → in-place edit (no file size change).
27/// - `original_range.len() != new_length` → resize edit (file grows or shrinks).
28/// - `original_range.start == original_range.end` → pure insert at that position.
29/// - `new_length == 0` → pure delete of `original_range`.
30///
31/// Pure append at end of file is `{ original_range: original_size..original_size, new_length: N, reader: ... }`.
32/// Pure truncate-to-N is `{ original_range: N..original_size, new_length: 0, reader: empty }`.
33pub struct DirtyInput {
34    pub original_range: Range<u64>,
35    pub reader: Pin<Box<dyn AsyncRead + Send>>,
36    pub new_length: u64,
37}
38
39/// Size of blocks read from the dirty source and fed to the cleaner.
40const STREAM_BLOCK_SIZE: usize = 4 * 1024 * 1024; // 4 MB
41
42/// A dirty window the server told us to re-upload, augmented with the bytes the cleaner
43/// produced and the resulting per-window file info. `start`/`end` are file byte offsets
44/// extended for append/truncation past `original_size` if needed.
45struct UploadedWindow {
46    start: u64,
47    end: u64,
48    chunks: ChunkHashList,
49    mdb: MDBFileInfo,
50}
51
52/// Upload an edited version of an existing file, reusing the unchanged regions from the
53/// original file's CAS segments and only re-uploading the parts the caller actually
54/// rewrites.
55///
56/// `dirty_inputs` is a list of edits expressed in the **original file's coordinates**.
57/// Each edit replaces `original_range` with `new_length` bytes from its reader; resize
58/// edits (including pure inserts and pure deletes) are supported. The output file size is
59/// derived from the inputs.
60///
61/// # When to use
62///
63/// - **In-place edit**: `original_range.len() == new_length`, file size unchanged.
64/// - **Resize edit**: any `new_length`. Replaces `original_range.len()` original bytes with `new_length` new bytes.
65/// - **Pure insert**: `original_range.start == original_range.end`, `new_length > 0`.
66/// - **Pure delete**: `original_range.start < original_range.end`, `new_length == 0`.
67/// - **Append**: `original_range == original_size..original_size`, `new_length > 0`.
68/// - **Truncate to N**: `original_range == N..original_size`, `new_length == 0`.
69/// - **No change**: empty `dirty_inputs`. Returns the original hash without any CAS call.
70///
71/// # Arguments
72///
73/// * `config` - Translator configuration for creating upload sessions.
74/// * `cas_client` - CAS client for fetching original file metadata and downloading boundary bytes.
75/// * `original_hash` - Merkle hash of the original file in CAS.
76/// * `original_size` - Size of the original file in bytes.
77/// * `dirty_inputs` - Edits to apply. Must be sorted by `original_range.start` and non-overlapping
78///   (`prev.original_range.end <= next.original_range.start`). Each reader must yield exactly `new_length` bytes. Each
79///   reader is consumed exactly once.
80///
81/// # Limitations
82///
83/// The composed file has no SHA-256 metadata (`metadata_ext = None`), since recomputing
84/// it would require reading the full file. This means `upload_ranges` is only suitable
85/// for contexts that don't require SHA-256 verification (e.g. HF buckets, xet-native
86/// repos), not for Git LFS-backed repos that verify SHA-256 on download.
87pub async fn upload_ranges(
88    config: Arc<TranslatorConfig>,
89    cas_client: Arc<dyn Client>,
90    original_hash: MerkleHash,
91    original_size: u64,
92    mut dirty_inputs: Vec<DirtyInput>,
93) -> Result<XetFileInfo> {
94    validate_dirty_ranges(&dirty_inputs, original_size)?;
95    let total_size = compute_total_size(original_size, &dirty_inputs)?;
96
97    if dirty_inputs.is_empty() {
98        debug_assert_eq!(total_size, original_size);
99        return Ok(XetFileInfo::new(original_hash.hex(), original_size));
100    }
101
102    // Empty original: nothing to compose against — upload as a fresh file (concatenation of
103    // the edits' new bytes, since every `original_range` must be `0..0`).
104    if original_size == 0 {
105        return upload_fresh_file(config, dirty_inputs, total_size).await;
106    }
107
108    let recon_result = cas_client.get_file_reconstruction_info(&original_hash).await?;
109    let original_mdb = recon_result
110        .map(|(mdb, _)| mdb)
111        .ok_or_else(|| DataError::ParameterError(format!("file {} not found in CAS", original_hash.hex())))?;
112    if original_mdb.file_size() != original_size {
113        return Err(DataError::ParameterError(format!(
114            "caller said original_size={original_size} but reconstruction info reports {}",
115            original_mdb.file_size()
116        )));
117    }
118
119    // `seg_byte_starts[i]` is the first byte of segment `i`; the trailing entry equals `original_size`.
120    let mut seg_byte_starts: Vec<u64> = Vec::with_capacity(original_mdb.segments.len() + 1);
121    seg_byte_starts.push(0);
122    let mut acc = 0u64;
123    for s in &original_mdb.segments {
124        acc += s.unpacked_segment_bytes as u64;
125        seg_byte_starts.push(acc);
126    }
127
128    // Snapping to *segment* boundaries (rather than chunk boundaries) lets us swap whole
129    // segments during composition, so the client never has to truncate a segment mid-chunk.
130    // Safe to send to the server because segment edges are chunk edges, so the server's
131    // chunk-aligned windows come back identical to our snapped ranges.
132    //
133    // Pure inserts (`original_range.start == original_range.end`) snap to the segment that
134    // owns the insert position so the cleaner has enough surrounding bytes to re-chunk
135    // around the insertion. An insert at `original_size` snaps to the last segment.
136    let n_segs = original_mdb.segments.len();
137    let mut snapped: Vec<(u64, u64)> = Vec::with_capacity(dirty_inputs.len());
138    for input in &dirty_inputs {
139        let r = &input.original_range;
140        let (s, e) = if r.start == r.end {
141            // Pure insert: pick the segment containing `r.start`. At end-of-file, fall back to
142            // the last segment.
143            if r.start == original_size {
144                (seg_byte_starts[n_segs - 1], seg_byte_starts[n_segs])
145            } else {
146                (snap_to_segment_start(&seg_byte_starts, r.start), snap_to_segment_end(&seg_byte_starts, r.start + 1))
147            }
148        } else {
149            (snap_to_segment_start(&seg_byte_starts, r.start), snap_to_segment_end(&seg_byte_starts, r.end))
150        };
151        snapped.push((s, e));
152    }
153
154    snapped.sort_by_key(|&(s, _)| s);
155    let mut coalesced: Vec<(u64, u64)> = Vec::with_capacity(snapped.len());
156    for r in snapped {
157        if let Some(last) = coalesced.last_mut()
158            && r.0 <= last.1
159        {
160            last.1 = last.1.max(r.1);
161            continue;
162        }
163        coalesced.push(r);
164    }
165
166    let server_query: Vec<FileRange> = coalesced.iter().map(|&(s, e)| FileRange::new(s, e)).collect();
167    if server_query.is_empty() {
168        return Err(DataError::InternalError("internal: non-empty dirty_inputs produced no server query".into()));
169    }
170
171    let response: FileChunkHashesResponse = cas_client.get_file_chunk_hashes(&original_hash, server_query).await?;
172    // The server may coalesce adjacent/overlapping dirty ranges after extending them to
173    // stable boundaries, so only enforce shape invariants on the returned payload itself.
174    if response.windows.is_empty() {
175        return Err(DataError::InternalError("server returned no windows".into()));
176    }
177    if response.hash_ranges.len() != response.windows.len() + 1 {
178        return Err(DataError::InternalError(format!(
179            "server returned {} hash_ranges, expected {} (n_windows + 1)",
180            response.hash_ranges.len(),
181            response.windows.len() + 1
182        )));
183    }
184    let gap_verification = response.gap_verification;
185
186    let ctx = config.ctx.clone();
187    let session = FileUploadSession::new(config.clone()).await?;
188    let mut input_idx = 0usize;
189    let mut uploaded: Vec<UploadedWindow> = Vec::with_capacity(response.windows.len());
190
191    let mut buf = vec![0u8; STREAM_BLOCK_SIZE];
192    for window in response.windows.iter() {
193        let w_start = window.dirty_byte_range[0];
194        let w_end = window.dirty_byte_range[1];
195
196        // Find the slice of edits that land in this window. A pure insert at exactly
197        // `w_end` belongs here only when `w_end == original_size` (no later window can
198        // take it); anywhere else it belongs to the next window starting at that byte.
199        // Exercised by `test_resize_insert_at_segment_boundary` and `test_mid_edit_plus_append`.
200        let edits_end = dirty_inputs[input_idx..]
201            .iter()
202            .take_while(|d| {
203                let r = &d.original_range;
204                if r.start == r.end {
205                    r.start < w_end || (r.start == w_end && w_end == original_size)
206                } else {
207                    r.end <= w_end
208                }
209            })
210            .count()
211            + input_idx;
212        let window_edits = &mut dirty_inputs[input_idx..edits_end];
213
214        let (removed, added): (u64, u64) = window_edits
215            .iter()
216            .map(|d| (d.original_range.end - d.original_range.start, d.new_length))
217            .fold((0, 0), |(rm, ad), (r, a)| (rm + r, ad + a));
218        let middle_size = (w_end - w_start) + added - removed;
219
220        let (_id, mut cleaner) = session.start_clean(None, Some(middle_size), Sha256Policy::Skip)?;
221
222        let mut cursor = w_start;
223        for input in window_edits.iter_mut() {
224            let edit_start = input.original_range.start;
225            let edit_end = input.original_range.end;
226            debug_assert!(edit_start >= w_start && edit_end <= w_end, "edit straddles window (validation bug)");
227
228            if cursor < edit_start {
229                stream_cas_range(&ctx, &cas_client, original_hash, cursor, edit_start, &mut cleaner).await?;
230            }
231
232            let mut remaining = input.new_length as usize;
233            while remaining > 0 {
234                let to_read = buf.len().min(remaining);
235                input.reader.read_exact(&mut buf[..to_read]).await.map_err(|err| {
236                    DataError::InternalError(format!(
237                        "failed to read dirty input [{}, {}): {err}",
238                        input.original_range.start, input.original_range.end
239                    ))
240                })?;
241                cleaner.add_data(&buf[..to_read]).await?;
242                remaining -= to_read;
243            }
244
245            cursor = edit_end;
246        }
247        input_idx = edits_end;
248
249        if cursor < w_end {
250            stream_cas_range(&ctx, &cas_client, original_hash, cursor, w_end, &mut cleaner).await?;
251        }
252
253        let (_info, chunks, mdb, _metrics) = cleaner.finish_with_chunks_detached().await?;
254        uploaded.push(UploadedWindow {
255            start: w_start,
256            end: w_end,
257            chunks,
258            mdb,
259        });
260    }
261
262    // Every edit must have been assigned to exactly one window. If the server
263    // returned narrower `dirty_byte_range`s than requested, leftover edits would
264    // silently drop and produce a corrupt file with no error.
265    if input_idx != dirty_inputs.len() {
266        return Err(DataError::InternalError(format!(
267            "{} dirty edits not assigned to any window (input_idx={input_idx}, total={})",
268            dirty_inputs.len() - input_idx,
269            dirty_inputs.len()
270        )));
271    }
272
273    // Merge sequence: [gap0, w0, gap1, w1, ..., gapN]. Empty gaps (`None`) are skipped.
274    let mut hash_ranges = response.hash_ranges;
275    let trailing_gap = hash_ranges.pop().flatten();
276    // Leading gap is empty (None) => first window starts at byte 0.
277    let first_window_at_start = matches!(hash_ranges.first(), Some(None));
278    let last_window_at_end = trailing_gap.is_none();
279    let last_idx = uploaded.len() - 1;
280
281    let mut merge_seq: Vec<MerkleHashSubtree> = Vec::with_capacity(2 * uploaded.len() + 1);
282    for (i, (w, gap)) in uploaded.iter().zip(hash_ranges).enumerate() {
283        if let Some(g) = gap {
284            merge_seq.push(g);
285        }
286        let at_start = i == 0 && first_window_at_start;
287        let at_end = i == last_idx && last_window_at_end;
288        merge_seq.push(MerkleHashSubtree::from_chunks(at_start, &w.chunks, at_end));
289    }
290    if let Some(g) = trailing_gap {
291        merge_seq.push(g);
292    }
293
294    let merged = MerkleHashSubtree::merge(&merge_seq)
295        .map_err(|err| DataError::InternalError(format!("MerkleHashSubtree::merge failed: {err}")))?;
296    // `final_hash()` is the aggregated chunk hash; the file hash is its HMAC with the zero
297    // salt (matching `file_hash` / the cleaner's output for files without SHA-256 metadata,
298    // the only flavor `upload_ranges` produces). Empty content is the one exception:
299    // `file_hash([])` short-circuits to `MerkleHash::default()` *without* HMAC, so we mirror
300    // that here when the result is zero-length (e.g. truncating to empty).
301    let aggregated_hash = merged.final_hash().ok_or_else(|| {
302        DataError::InternalError("merged subtree is not fully closed; cannot derive final hash".into())
303    })?;
304    let combined_hash = if total_size == 0 {
305        MerkleHash::default()
306    } else {
307        aggregated_hash.hmac(MerkleHash::default())
308    };
309
310    let composed_mdb =
311        compose_mdb(&original_mdb, &seg_byte_starts, &uploaded, gap_verification, combined_hash, original_size)?;
312
313    debug!(
314        "upload_ranges: composed hash={}, {} segments, {} windows",
315        combined_hash.hex(),
316        composed_mdb.segments.len(),
317        uploaded.len()
318    );
319
320    session.register_composed_file(composed_mdb).await?;
321    session.finalize().await?;
322
323    let total_dirty: u64 = dirty_inputs.iter().map(|d| d.new_length).sum();
324    info!(
325        "upload_ranges: hash={} size={} (original={}, {} windows, {} dirty bytes)",
326        combined_hash.hex(),
327        total_size,
328        original_size,
329        uploaded.len(),
330        total_dirty
331    );
332
333    Ok(XetFileInfo::new(combined_hash.hex(), total_size))
334}
335
336/// Assemble the composed MDB by splicing uploaded window segments into the original
337/// file's segment list, pulling verification entries from the server (for stable gaps)
338/// and from the cleaner (for re-uploaded windows).
339fn compose_mdb(
340    original_mdb: &MDBFileInfo,
341    seg_byte_starts: &[u64],
342    uploaded: &[UploadedWindow],
343    gap_verification: Vec<HexMerkleHash>,
344    combined_hash: MerkleHash,
345    original_size: u64,
346) -> Result<MDBFileInfo> {
347    let mut all_segments: Vec<FileDataSequenceEntry> = Vec::new();
348    let mut all_verification: Vec<FileVerificationEntry> = Vec::new();
349    let mut seg_idx = 0usize;
350    let n_segs = original_mdb.segments.len();
351    let mut gap_idx = 0usize;
352
353    for w in uploaded {
354        while seg_idx < n_segs && seg_byte_starts[seg_idx] < w.start {
355            if seg_byte_starts[seg_idx + 1] > w.start {
356                return Err(DataError::InternalError(format!(
357                    "server returned a window starting at {} that straddles segment {} \
358                     ({}..{}); composition requires segment-aligned windows",
359                    w.start,
360                    seg_idx,
361                    seg_byte_starts[seg_idx],
362                    seg_byte_starts[seg_idx + 1]
363                )));
364            }
365            all_segments.push(original_mdb.segments[seg_idx].clone());
366            let entry = gap_verification.get(gap_idx).ok_or_else(|| {
367                DataError::InternalError(format!(
368                    "ran out of gap_verification entries at stable segment {seg_idx}; \
369                     server response is inconsistent with the segment layout"
370                ))
371            })?;
372            all_verification.push(FileVerificationEntry::new(entry.into()));
373            gap_idx += 1;
374            seg_idx += 1;
375        }
376        // Server guarantees window ends are clamped to file_size (see xetcas
377        // `core::get_file_chunk_hashes`), so this invariant should always hold.
378        debug_assert!(w.end <= original_size, "window end {} exceeds original_size {}", w.end, original_size);
379        while seg_idx < n_segs && seg_byte_starts[seg_idx] < w.end {
380            seg_idx += 1;
381        }
382        if w.mdb.verification.len() != w.mdb.segments.len() {
383            return Err(DataError::InternalError(format!(
384                "window MDB has {} segments but {} verification entries",
385                w.mdb.segments.len(),
386                w.mdb.verification.len()
387            )));
388        }
389        all_segments.extend_from_slice(&w.mdb.segments);
390        all_verification.extend_from_slice(&w.mdb.verification);
391    }
392    while seg_idx < n_segs {
393        all_segments.push(original_mdb.segments[seg_idx].clone());
394        let entry = gap_verification.get(gap_idx).ok_or_else(|| {
395            DataError::InternalError(format!(
396                "ran out of gap_verification entries at stable segment {seg_idx}; \
397                 server response is inconsistent with the segment layout"
398            ))
399        })?;
400        all_verification.push(FileVerificationEntry::new(entry.into()));
401        gap_idx += 1;
402        seg_idx += 1;
403    }
404    if gap_idx < gap_verification.len() {
405        return Err(DataError::InternalError(format!(
406            "server returned {} gap_verification entries but only {} stable segments were emitted",
407            gap_verification.len(),
408            gap_idx
409        )));
410    }
411
412    debug_assert_eq!(all_segments.len(), all_verification.len());
413
414    Ok(MDBFileInfo {
415        metadata: FileDataSequenceHeader::new(combined_hash, all_segments.len(), true, false),
416        segments: all_segments,
417        verification: all_verification,
418        metadata_ext: None,
419    })
420}
421
422/// Validate the caller-provided dirty ranges.
423///
424/// `dirty_inputs` must be sorted by `original_range.start`, non-overlapping
425/// (`prev.original_range.end <= next.original_range.start`), and every
426/// `original_range.end <= original_size`. Empty ranges (pure inserts) are allowed at any
427/// position including `original_size` (append).
428fn validate_dirty_ranges(dirty_inputs: &[DirtyInput], original_size: u64) -> Result<()> {
429    let mut prev_end = 0u64;
430    for (i, input) in dirty_inputs.iter().enumerate() {
431        let r = &input.original_range;
432        if r.start > r.end {
433            return Err(DataError::ParameterError(format!(
434                "dirty_inputs[{i}].original_range is reversed: {}..{}",
435                r.start, r.end
436            )));
437        }
438        if r.end > original_size {
439            return Err(DataError::ParameterError(format!(
440                "dirty_inputs[{i}].original_range end ({}) exceeds original_size ({original_size})",
441                r.end
442            )));
443        }
444        if i > 0 && r.start < prev_end {
445            return Err(DataError::ParameterError(format!(
446                "dirty_inputs[{i}].original_range overlaps the previous edit (starts at {} < {prev_end})",
447                r.start
448            )));
449        }
450        prev_end = r.end;
451    }
452    Ok(())
453}
454
455/// Compute the resulting file size: `original_size` minus bytes removed plus bytes added.
456fn compute_total_size(original_size: u64, dirty_inputs: &[DirtyInput]) -> Result<u64> {
457    let (removed, added) = dirty_inputs
458        .iter()
459        .fold((0u64, 0u64), |(r, a), d| (r + (d.original_range.end - d.original_range.start), a + d.new_length));
460    original_size
461        .checked_add(added)
462        .and_then(|s| s.checked_sub(removed))
463        .ok_or_else(|| {
464            DataError::ParameterError(format!(
465                "total size overflows: original_size={original_size}, added={added}, removed={removed}"
466            ))
467        })
468}
469
470/// Upload a brand-new file from `dirty_inputs` (no original to compose against).
471/// Used when the original file is empty: every edit's `original_range` is `0..0`, so we
472/// just stream their new bytes through the cleaner.
473async fn upload_fresh_file(
474    config: Arc<TranslatorConfig>,
475    mut dirty_inputs: Vec<DirtyInput>,
476    total_size: u64,
477) -> Result<XetFileInfo> {
478    let session = FileUploadSession::new(config).await?;
479    let (_id, mut cleaner) = session.start_clean(None, Some(total_size), Sha256Policy::Skip)?;
480    for input in &mut dirty_inputs {
481        let mut remaining = input.new_length as usize;
482        let mut buf = vec![0u8; STREAM_BLOCK_SIZE.min(remaining.max(1))];
483        while remaining > 0 {
484            let to_read = buf.len().min(remaining);
485            input.reader.read_exact(&mut buf[..to_read]).await.map_err(|err| {
486                DataError::InternalError(format!("failed to read dirty input at {}: {err}", input.original_range.start))
487            })?;
488            cleaner.add_data(&buf[..to_read]).await?;
489            remaining -= to_read;
490        }
491    }
492    let (info, _metrics) = cleaner.finish().await?;
493    session.finalize().await?;
494    Ok(info)
495}
496
497/// Stream a byte range from CAS into the cleaner.
498async fn stream_cas_range(
499    ctx: &XetContext,
500    cas_client: &Arc<dyn Client>,
501    file_hash: MerkleHash,
502    start: u64,
503    end: u64,
504    cleaner: &mut super::SingleFileCleaner,
505) -> Result<()> {
506    let reconstructor = FileReconstructor::new(ctx, cas_client, file_hash).with_byte_range(FileRange::new(start, end));
507    let mut stream = reconstructor.reconstruct_to_stream();
508    while let Some(chunk) = stream.next().await? {
509        cleaner.add_data(&chunk).await?;
510    }
511    Ok(())
512}
513
514/// Largest segment-start byte that is `<= byte`. Used to snap a dirty-range start back
515/// to its enclosing segment boundary.
516fn snap_to_segment_start(seg_byte_starts: &[u64], byte: u64) -> u64 {
517    let idx = seg_byte_starts.partition_point(|&s| s <= byte);
518    seg_byte_starts[idx.saturating_sub(1)]
519}
520
521/// Smallest segment-start byte that is `>= byte`. Used to snap a dirty-range end forward
522/// to the segment boundary that fully contains it.
523fn snap_to_segment_end(seg_byte_starts: &[u64], byte: u64) -> u64 {
524    let idx = seg_byte_starts.partition_point(|&s| s < byte);
525    seg_byte_starts[idx]
526}
527
528#[cfg(test)]
529mod tests {
530    use std::io::Cursor;
531    use std::ops::Range;
532    use std::path::Path;
533    use std::sync::Arc;
534
535    use tempfile::TempDir;
536    use xet_client::cas_client::{Client, LocalTestServerBuilder};
537    use xet_core_structures::merklehash::MerkleHash;
538
539    use super::*;
540    use crate::processing::configurations::TranslatorConfig;
541    use crate::processing::file_cleaner::Sha256Policy;
542    use crate::processing::file_download_session::FileDownloadSession;
543    use crate::processing::file_upload_session::FileUploadSession;
544
545    fn test_config(endpoint: impl AsRef<str>, base_dir: impl AsRef<Path>) -> Arc<TranslatorConfig> {
546        let ctx = XetContext::default().unwrap();
547        Arc::new(TranslatorConfig::test_server_config(&ctx, endpoint, base_dir).unwrap())
548    }
549
550    /// Test helper: fetch the original file's per-segment byte sizes. Tests use these to
551    /// build dirty ranges that align with segment (== chunk-group) boundaries — handy for
552    /// scenarios that want to overwrite or truncate exactly on a boundary.
553    async fn fetch_segment_sizes(cas_client: &Arc<dyn Client>, hash: &MerkleHash) -> Vec<u64> {
554        let (mdb, _) = cas_client.get_file_reconstruction_info(hash).await.unwrap().unwrap();
555        mdb.segments.iter().map(|s| s.unpacked_segment_bytes as u64).collect()
556    }
557
558    /// Build in-place `DirtyInput`s (each edit's `new_length` matches its `original_range`
559    /// length, so the file size doesn't change) from a source buffer and range list.
560    fn make_dirty_inputs(ranges: &[(u64, u64)], data: &[u8]) -> Vec<DirtyInput> {
561        ranges
562            .iter()
563            .map(|&(start, end)| {
564                let slice = data[start as usize..end as usize].to_vec();
565                DirtyInput {
566                    original_range: start..end,
567                    new_length: end - start,
568                    reader: Box::pin(Cursor::new(slice)),
569                }
570            })
571            .collect()
572    }
573
574    /// Build `DirtyInput`s with dummy readers (for validation error tests where
575    /// the reader is never consumed).
576    fn make_dummy_inputs(ranges: &[(u64, u64)]) -> Vec<DirtyInput> {
577        ranges
578            .iter()
579            .map(|&(start, end)| DirtyInput {
580                original_range: start..end,
581                new_length: end - start,
582                reader: Box::pin(Cursor::new(Vec::new())),
583            })
584            .collect()
585    }
586
587    /// Bridge the old `(start, end)` + `total_size`-shaped test cases into the new edit
588    /// model. `(start, end)` indexes into `data` AND describes the output byte range.
589    /// Handles in-place edits, pure appends, mid-edit-plus-append spans (split into two
590    /// edits), and trailing truncations.
591    fn make_legacy_inputs(specs: &[(u64, u64)], data: &[u8], original_size: u64, total_size: u64) -> Vec<DirtyInput> {
592        let mut out: Vec<DirtyInput> = Vec::new();
593        for &(start, end) in specs {
594            let bytes = data[start as usize..end as usize].to_vec();
595            if end <= original_size {
596                out.push(DirtyInput {
597                    original_range: start..end,
598                    new_length: bytes.len() as u64,
599                    reader: Box::pin(Cursor::new(bytes)),
600                });
601            } else if start >= original_size {
602                out.push(DirtyInput {
603                    original_range: original_size..original_size,
604                    new_length: bytes.len() as u64,
605                    reader: Box::pin(Cursor::new(bytes)),
606                });
607            } else {
608                let split = (original_size - start) as usize;
609                let (head, tail) = bytes.split_at(split);
610                out.push(DirtyInput {
611                    original_range: start..original_size,
612                    new_length: head.len() as u64,
613                    reader: Box::pin(Cursor::new(head.to_vec())),
614                });
615                out.push(DirtyInput {
616                    original_range: original_size..original_size,
617                    new_length: tail.len() as u64,
618                    reader: Box::pin(Cursor::new(tail.to_vec())),
619                });
620            }
621        }
622        if total_size < original_size {
623            out.push(DirtyInput {
624                original_range: total_size..original_size,
625                new_length: 0,
626                reader: Box::pin(Cursor::new(Vec::new())),
627            });
628        }
629        out
630    }
631
632    // original: [=========================== 256 KB ===========================]
633    // dirty:                  [== 1 KB ===]
634    // result:   [===stable===][re-uploaded][============stable=================]
635    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
636    async fn test_upload_ranges_mid_file_edit() {
637        let server = LocalTestServerBuilder::new().start().await;
638        let base_dir = TempDir::new().unwrap();
639        let endpoint = server.http_endpoint().to_string();
640        let config = test_config(&endpoint, base_dir.path());
641
642        // Use the server directly as the CAS client (bypasses HTTP for get_file_chunk_hashes).
643        let cas_client: Arc<dyn Client> = Arc::new(server);
644
645        // 1. Upload an original file: 256 KB of pseudo-random bytes.
646        let original_data = random_data(42, 256 * 1024);
647        let original_hash = {
648            let upload_session = FileUploadSession::new(config.clone()).await.unwrap();
649            let (_id, mut cleaner) = upload_session
650                .start_clean(Some("original".into()), Some(original_data.len() as u64), Sha256Policy::Skip)
651                .unwrap();
652            cleaner.add_data(&original_data).await.unwrap();
653            let (xfi, _metrics) = cleaner.finish().await.unwrap();
654            upload_session.finalize().await.unwrap();
655            MerkleHash::from_hex(xfi.hash()).unwrap()
656        };
657        let original_size = original_data.len() as u64;
658
659        // 2. Build modified content: overwrite [100_000, 101_000) with 0xBB.
660        let mut modified_data = original_data.clone();
661        let dirty_start = 100_000usize;
662        let dirty_end = 101_000usize;
663        modified_data[dirty_start..dirty_end].fill(0xBB);
664        let total_size = modified_data.len() as u64;
665        let result = upload_ranges(
666            config.clone(),
667            cas_client.clone(),
668            original_hash,
669            original_size,
670            make_legacy_inputs(&[(dirty_start as u64, dirty_end as u64)], &modified_data, original_size, total_size),
671        )
672        .await
673        .unwrap();
674
675        assert_eq!(result.file_size, Some(total_size));
676
677        // 3. Download and verify the composed file.
678        let composed_hash = MerkleHash::from_hex(result.hash()).unwrap();
679        let session = FileDownloadSession::new(config.clone(), None).await.unwrap();
680        let file_info = crate::processing::XetFileInfo::new(composed_hash.hex(), total_size);
681        let out_path = base_dir.path().join("output");
682        session.download_file(&file_info, &out_path).await.unwrap();
683        let downloaded = std::fs::read(&out_path).unwrap();
684
685        assert_eq!(downloaded.len(), modified_data.len());
686        assert_eq!(downloaded, modified_data);
687
688        // Hash must match a clean upload of the same content.
689        let clean_hash = upload_file(&config, &modified_data).await;
690        assert_eq!(result.hash(), clean_hash.hex(), "hash mismatch with clean upload");
691    }
692
693    // original: [=========================== 256 KB ===========================]
694    // result:   [========= 100 KB =========]
695    //                                       ^ cut here (mid-chunk)
696    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
697    async fn test_upload_ranges_truncation() {
698        let server = LocalTestServerBuilder::new().start().await;
699        let base_dir = TempDir::new().unwrap();
700        let config = test_config(server.http_endpoint(), base_dir.path());
701        let cas_client: Arc<dyn Client> = Arc::new(server);
702
703        // Upload 256 KB file.
704        let original_data = random_data(43, 256 * 1024);
705        let original_hash = upload_file(&config, &original_data).await;
706        let original_size = original_data.len() as u64;
707
708        // Truncate to 100 KB (no dirty ranges, pure truncation).
709        let truncated_size = 100_000u64;
710        let result = upload_ranges(
711            config.clone(),
712            cas_client.clone(),
713            original_hash,
714            original_size,
715            make_legacy_inputs(&[], &[], original_size, truncated_size),
716        )
717        .await
718        .unwrap();
719
720        assert_eq!(result.file_size(), Some(truncated_size));
721
722        // Download and verify: first truncated_size bytes should match original.
723        let downloaded = download_file(&config, MerkleHash::from_hex(result.hash()).unwrap(), truncated_size).await;
724        assert_eq!(downloaded.len(), truncated_size as usize);
725        assert_eq!(downloaded, &original_data[..truncated_size as usize]);
726
727        let clean_hash = upload_file(&config, &original_data[..truncated_size as usize]).await;
728        assert_eq!(result.hash(), clean_hash.hex(), "hash mismatch with clean upload");
729    }
730
731    // original: [======== 100 KB ========]
732    // result:   [======== 100 KB ========][== 50 KB appended ==]
733    //                                     ^ last chunk re-chunked with append
734    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
735    async fn test_upload_ranges_append() {
736        let server = LocalTestServerBuilder::new().start().await;
737        let base_dir = TempDir::new().unwrap();
738        let config = test_config(server.http_endpoint(), base_dir.path());
739        let cas_client: Arc<dyn Client> = Arc::new(server);
740
741        // Upload 100 KB file.
742        let original_data = random_data(44, 100 * 1024);
743        let original_hash = upload_file(&config, &original_data).await;
744        let original_size = original_data.len() as u64;
745
746        // Append 50 KB of pseudo-random data.
747        let mut full_data = original_data.clone();
748        full_data.extend(random_data(99, 50 * 1024));
749        let total_size = full_data.len() as u64;
750        let result = upload_ranges(
751            config.clone(),
752            cas_client.clone(),
753            original_hash,
754            original_size,
755            make_legacy_inputs(&[(original_size, total_size)], &full_data, original_size, total_size),
756        )
757        .await
758        .unwrap();
759
760        assert_eq!(result.file_size(), Some(total_size));
761
762        let downloaded = download_file(&config, MerkleHash::from_hex(result.hash()).unwrap(), total_size).await;
763        assert_eq!(downloaded, full_data);
764
765        let clean_hash = upload_file(&config, &full_data).await;
766        assert_eq!(result.hash(), clean_hash.hex(), "hash mismatch with clean upload");
767    }
768
769    // original: [=========================== 256 KB ==============================]
770    // dirty:    [4K]
771    // result:   [re-uploaded][=================stable=============================]
772    //           ^ no stable prefix
773    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
774    async fn test_upload_ranges_at_file_start() {
775        let server = LocalTestServerBuilder::new().start().await;
776        let base_dir = TempDir::new().unwrap();
777        let config = test_config(server.http_endpoint(), base_dir.path());
778        let cas_client: Arc<dyn Client> = Arc::new(server);
779
780        // Upload 256 KB file.
781        let original_data = random_data(45, 256 * 1024);
782        let original_hash = upload_file(&config, &original_data).await;
783        let original_size = original_data.len() as u64;
784
785        // Overwrite [0, 4096) with 0xBB (dirty range at offset 0, no stable prefix).
786        let mut modified_data = original_data.clone();
787        modified_data[..4096].fill(0xBB);
788        let total_size = modified_data.len() as u64;
789        let result = upload_ranges(
790            config.clone(),
791            cas_client.clone(),
792            original_hash,
793            original_size,
794            make_legacy_inputs(&[(0, 4096)], &modified_data, original_size, total_size),
795        )
796        .await
797        .unwrap();
798
799        assert_eq!(result.file_size(), Some(total_size));
800
801        let downloaded = download_file(&config, MerkleHash::from_hex(result.hash()).unwrap(), total_size).await;
802        assert_eq!(downloaded.len(), modified_data.len());
803        assert_eq!(downloaded, modified_data);
804
805        let clean_hash = upload_file(&config, &modified_data).await;
806        assert_eq!(result.hash(), clean_hash.hex(), "hash mismatch with clean upload");
807    }
808
809    // original: [=========================== 256 KB ===========================]
810    // dirty:       [2K]                               [2K]
811    // result:   [s][re-up][=======stable========][re-up][========stable========]
812    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
813    async fn test_upload_ranges_multiple_regions() {
814        let server = LocalTestServerBuilder::new().start().await;
815        let base_dir = TempDir::new().unwrap();
816        let config = test_config(server.http_endpoint(), base_dir.path());
817        let cas_client: Arc<dyn Client> = Arc::new(server);
818
819        // Upload 256 KB file.
820        let original_data = random_data(46, 256 * 1024);
821        let original_hash = upload_file(&config, &original_data).await;
822        let original_size = original_data.len() as u64;
823
824        // Two non-adjacent dirty ranges with a stable gap between them.
825        let mut modified_data = original_data.clone();
826        modified_data[10_000..12_000].fill(0xBB); // first dirty range
827        modified_data[200_000..202_000].fill(0xCC); // second dirty range
828        let total_size = modified_data.len() as u64;
829        let result = upload_ranges(
830            config.clone(),
831            cas_client.clone(),
832            original_hash,
833            original_size,
834            make_legacy_inputs(&[(10_000, 12_000), (200_000, 202_000)], &modified_data, original_size, total_size),
835        )
836        .await
837        .unwrap();
838
839        assert_eq!(result.file_size(), Some(total_size));
840
841        let downloaded = download_file(&config, MerkleHash::from_hex(result.hash()).unwrap(), total_size).await;
842        assert_eq!(downloaded.len(), modified_data.len());
843        assert_eq!(downloaded, modified_data);
844
845        let clean_hash = upload_file(&config, &modified_data).await;
846        assert_eq!(result.hash(), clean_hash.hex(), "hash mismatch with clean upload");
847    }
848
849    // original: [======== 100 KB ========]
850    // inputs:   [======== 100 KB ========][000][=4K written=]
851    //                                     ^ gap (zeros from seek past EOF)
852    //
853    // With the new API, the caller provides the full append region [original_size, total_size)
854    // as a single DirtyInput, including the sparse gap.
855    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
856    async fn test_append_with_gap_before_dirty_range() {
857        let server = LocalTestServerBuilder::new().start().await;
858        let base_dir = TempDir::new().unwrap();
859        let config = test_config(server.http_endpoint(), base_dir.path());
860        let cas_client: Arc<dyn Client> = Arc::new(server);
861
862        let original_data = random_data(50, 100 * 1024);
863        let original_hash = upload_file(&config, &original_data).await;
864        let original_size = original_data.len() as u64;
865
866        let gap = 500u64;
867        let write_data = random_data(101, 4096);
868        let total_size = original_size + gap + write_data.len() as u64;
869
870        let mut full_data = original_data.clone();
871        full_data.extend(vec![0x00u8; gap as usize]);
872        full_data.extend(&write_data);
873
874        // The caller provides the entire append region (including the sparse gap).
875        let result = upload_ranges(
876            config.clone(),
877            cas_client.clone(),
878            original_hash,
879            original_size,
880            make_legacy_inputs(&[(original_size, total_size)], &full_data, original_size, total_size),
881        )
882        .await
883        .unwrap();
884
885        assert_eq!(result.file_size(), Some(total_size));
886
887        let downloaded = download_file(&config, MerkleHash::from_hex(result.hash()).unwrap(), total_size).await;
888        assert_eq!(downloaded.len(), full_data.len(), "size mismatch");
889        assert_eq!(&downloaded[..], &full_data[..], "content mismatch: gap bytes were lost");
890
891        let clean_hash = upload_file(&config, &full_data).await;
892        assert_eq!(result.hash(), clean_hash.hex(), "hash mismatch with clean upload");
893    }
894
895    // staging:  [000000 (zeros) 000000][=== appended ===]
896    // CAS:      [=== original data ==]
897    // result:   [=== original data ==][=== appended ===]
898    //           ^ boundary prefix must come from CAS, not from staging zeros
899    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
900    async fn test_append_sparse_staging_file() {
901        let server = LocalTestServerBuilder::new().start().await;
902        let base_dir = TempDir::new().unwrap();
903        let config = test_config(server.http_endpoint(), base_dir.path());
904        let cas_client: Arc<dyn Client> = Arc::new(server);
905
906        let original_data = vec![0xDDu8; 100 * 1024];
907        let original_hash = upload_file(&config, &original_data).await;
908        let original_size = original_data.len() as u64;
909
910        let append_data = vec![0xEEu8; 50 * 1024];
911        let total_size = original_size + append_data.len() as u64;
912
913        // Build a sparse staging file: zeros for [0, original_size), real data after.
914        // This is what hf-mount produces (sparse hole + appended bytes).
915        let mut sparse_staging = vec![0u8; total_size as usize];
916        sparse_staging[original_size as usize..].copy_from_slice(&append_data);
917        let result = upload_ranges(
918            config.clone(),
919            cas_client.clone(),
920            original_hash,
921            original_size,
922            make_legacy_inputs(&[(original_size, total_size)], &sparse_staging, original_size, total_size),
923        )
924        .await
925        .unwrap();
926
927        // Expected: original data + appended data (not zeros + appended data).
928        let mut expected = original_data.clone();
929        expected.extend(&append_data);
930
931        let downloaded = download_file(&config, MerkleHash::from_hex(result.hash()).unwrap(), total_size).await;
932        assert_eq!(downloaded.len(), expected.len(), "size mismatch");
933        assert_eq!(&downloaded[..], &expected[..], "content mismatch: CAS data replaced by zeros from sparse file");
934
935        let clean_hash = upload_file(&config, &expected).await;
936        assert_eq!(result.hash(), clean_hash.hex(), "hash mismatch with clean upload");
937    }
938
939    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
940    async fn test_data_integrity_scenarios() {
941        let server = LocalTestServerBuilder::new().start().await;
942        let base_dir = TempDir::new().unwrap();
943        let config = test_config(server.http_endpoint(), base_dir.path());
944        let cas_client: Arc<dyn Client> = Arc::new(server);
945
946        // original: [========================= 256 KB =========================]
947        // dirty:                         [10K]
948        // result:   [====== 90 KB ======][10K]
949        //           0           90K   100K   ^ truncate here, dirty touches cut
950        {
951            let original = vec![0xAAu8; 256 * 1024];
952            let mut expected = original[..100_000].to_vec();
953            expected[90_000..100_000].fill(0xBB);
954            assert_range_edit(&config, &cas_client, &original, &expected, &[(90_000, 100_000)], 100_000).await;
955        }
956
957        // original: [========= 128 KB =========]
958        // dirty:    [========= 128 KB =========]
959        // result:   [======= re-uploaded ======]  (no stable regions)
960        {
961            let original = vec![0xAAu8; 128 * 1024];
962            let expected = vec![0xBBu8; 128 * 1024];
963            let size = original.len() as u64;
964            assert_range_edit(&config, &cas_client, &original, &expected, &[(0, size)], size).await;
965        }
966
967        // original: [===================== 256 KB =====================]
968        // dirty:                [1K][1K][1K]
969        //                       ^-- coalesced into one region
970        {
971            let original = vec![0xAAu8; 256 * 1024];
972            let mut expected = original.clone();
973            expected[50_000..51_000].fill(0xBB);
974            expected[51_000..52_000].fill(0xCC);
975            expected[52_000..53_000].fill(0xDD);
976            let size = original.len() as u64;
977            assert_range_edit(
978                &config,
979                &cas_client,
980                &original,
981                &expected,
982                &[(50_000, 51_000), (51_000, 52_000), (52_000, 53_000)],
983                size,
984            )
985            .await;
986        }
987
988        // original: [======== 100 KB ========]
989        // result:   [======== 100 KB ========][== 50 KB ==]
990        //           no dirty_ranges, only total_size > original_size
991        {
992            let original = vec![0xAAu8; 100 * 1024];
993            let mut expected = original.clone();
994            expected.extend(vec![0xEEu8; 50 * 1024]);
995            let total = expected.len() as u64;
996            assert_range_edit(&config, &cas_client, &original, &expected, &[], total).await;
997        }
998
999        // original: [chunk0][chunk1][chunk2][...]
1000        // dirty:                    [chunk2]
1001        //                           ^      ^-- starts/ends on chunk boundary
1002        {
1003            let original: Vec<u8> = (0..256 * 1024)
1004                .map(|i: usize| {
1005                    let x = i.wrapping_mul(2654435761);
1006                    (x >> 16) as u8
1007                })
1008                .collect();
1009            let original_hash = upload_file(&config, &original).await;
1010            let seg_sizes = fetch_segment_sizes(&cas_client, &original_hash).await;
1011            if seg_sizes.len() >= 3 {
1012                let boundary: u64 = seg_sizes[0] + seg_sizes[1];
1013                let dirty_end = boundary + seg_sizes[2];
1014                let mut expected = original.clone();
1015                expected[boundary as usize..dirty_end as usize].fill(0xFF);
1016                let size = original.len() as u64;
1017                let result = upload_ranges(
1018                    config.clone(),
1019                    cas_client.clone(),
1020                    original_hash,
1021                    size,
1022                    make_legacy_inputs(&[(boundary, dirty_end)], &expected, size, size),
1023                )
1024                .await
1025                .unwrap();
1026                let downloaded = download_file(&config, MerkleHash::from_hex(result.hash()).unwrap(), size).await;
1027                assert_eq!(downloaded, expected, "chunk-boundary edit mismatch");
1028
1029                let clean_hash = upload_file(&config, &expected).await;
1030                assert_eq!(result.hash(), clean_hash.hex(), "hash mismatch with clean upload");
1031            }
1032        }
1033    }
1034
1035    // No changes: dirty_ranges=[], total_size == original_size -> early return.
1036    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1037    async fn test_noop_returns_original_hash() {
1038        let server = LocalTestServerBuilder::new().start().await;
1039        let base_dir = TempDir::new().unwrap();
1040        let config = test_config(server.http_endpoint(), base_dir.path());
1041        let cas_client: Arc<dyn Client> = Arc::new(server);
1042
1043        let data = random_data(70, 256 * 1024);
1044        let hash = upload_file(&config, &data).await;
1045        let size = data.len() as u64;
1046        let result = upload_ranges(config, cas_client, hash, size, make_legacy_inputs(&[], &[], size, size))
1047            .await
1048            .unwrap();
1049
1050        assert_eq!(result.hash(), hash.hex());
1051        assert_eq!(result.file_size(), Some(size));
1052    }
1053
1054    // dirty_range end > total_size -> rejected.
1055    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1056    async fn test_rejects_dirty_range_past_total_size() {
1057        let server = LocalTestServerBuilder::new().start().await;
1058        let base_dir = TempDir::new().unwrap();
1059        let config = test_config(server.http_endpoint(), base_dir.path());
1060        let cas_client: Arc<dyn Client> = Arc::new(server);
1061
1062        let data = random_data(71, 256 * 1024);
1063        let hash = upload_file(&config, &data).await;
1064        let size = data.len() as u64;
1065        let err = upload_ranges(config, cas_client, hash, size, make_dummy_inputs(&[(100, size + 1)])).await;
1066        assert!(err.is_err(), "dirty range past total_size should be rejected");
1067    }
1068
1069    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1070    async fn test_rejects_overlapping_dirty_ranges() {
1071        let server = LocalTestServerBuilder::new().start().await;
1072        let base_dir = TempDir::new().unwrap();
1073        let config = test_config(server.http_endpoint(), base_dir.path());
1074        let cas_client: Arc<dyn Client> = Arc::new(server);
1075
1076        let data = random_data(60, 256 * 1024);
1077        let hash = upload_file(&config, &data).await;
1078        let size = data.len() as u64;
1079        let err = upload_ranges(config, cas_client, hash, size, make_dummy_inputs(&[(100, 300), (200, 400)])).await;
1080        assert!(err.is_err(), "overlapping ranges should be rejected");
1081    }
1082
1083    // Regression: validation must run *before* the `original_size == 0` short-circuit, or
1084    // the empty-original path silently accepts ranges that exceed `original_size`.
1085    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1086    async fn test_empty_original_validates_ranges() {
1087        let server = LocalTestServerBuilder::new().start().await;
1088        let base_dir = TempDir::new().unwrap();
1089        let config = test_config(server.http_endpoint(), base_dir.path());
1090        let cas_client: Arc<dyn Client> = Arc::new(server);
1091
1092        let original_hash = upload_file(&config, &[]).await;
1093
1094        // `0..10` and `5..15` both have `end > original_size == 0`; either would corrupt
1095        // the upload if validation didn't fire.
1096        let inputs = vec![
1097            DirtyInput {
1098                original_range: 0..10,
1099                new_length: 10,
1100                reader: Box::pin(Cursor::new(vec![0xAA; 10])),
1101            },
1102            DirtyInput {
1103                original_range: 5..15,
1104                new_length: 10,
1105                reader: Box::pin(Cursor::new(vec![0xBB; 10])),
1106            },
1107        ];
1108        let err = upload_ranges(config, cas_client, original_hash, 0, inputs).await;
1109        assert!(err.is_err(), "ranges with end > original_size must be rejected for empty originals too");
1110    }
1111
1112    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1113    async fn test_rejects_unsorted_dirty_ranges() {
1114        let server = LocalTestServerBuilder::new().start().await;
1115        let base_dir = TempDir::new().unwrap();
1116        let config = test_config(server.http_endpoint(), base_dir.path());
1117        let cas_client: Arc<dyn Client> = Arc::new(server);
1118
1119        let data = random_data(62, 256 * 1024);
1120        let hash = upload_file(&config, &data).await;
1121        let size = data.len() as u64;
1122        let err = upload_ranges(config, cas_client, hash, size, make_dummy_inputs(&[(300, 400), (100, 200)])).await;
1123        assert!(err.is_err(), "unsorted ranges should be rejected");
1124    }
1125
1126    // original: [chunk0][chunk1][chunk2][chunk3][...more chunks...]
1127    // input:              [========= single large write ==========]
1128    //
1129    // A single DirtyInput that spans many segments. Verifies the reader is consumed
1130    // correctly across the multi-segment window.
1131    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1132    async fn test_single_input_spanning_many_chunks() {
1133        let server = LocalTestServerBuilder::new().start().await;
1134        let base_dir = TempDir::new().unwrap();
1135        let config = test_config(server.http_endpoint(), base_dir.path());
1136        let cas_client: Arc<dyn Client> = Arc::new(server);
1137
1138        let original_data = random_data(99, 256 * 1024);
1139        let original_hash = upload_file(&config, &original_data).await;
1140        let original_size = original_data.len() as u64;
1141
1142        // Overwrite a large middle section (likely spans many CDC chunks).
1143        let mut modified = original_data.clone();
1144        let dirty_start = 10_000u64;
1145        let dirty_end = 200_000u64;
1146        modified[dirty_start as usize..dirty_end as usize].fill(0xFF);
1147
1148        let result = upload_ranges(
1149            config.clone(),
1150            cas_client.clone(),
1151            original_hash,
1152            original_size,
1153            make_legacy_inputs(&[(dirty_start, dirty_end)], &modified, original_size, original_size),
1154        )
1155        .await
1156        .unwrap();
1157
1158        let downloaded = download_file(&config, MerkleHash::from_hex(result.hash()).unwrap(), original_size).await;
1159        assert_eq!(downloaded, modified, "large spanning input produced wrong content");
1160
1161        let clean_hash = upload_file(&config, &modified).await;
1162        assert_eq!(result.hash(), clean_hash.hex(), "hash mismatch with clean upload");
1163    }
1164
1165    // original: b"AAAA_HEADER_AAAA|" (17 bytes, single CAS chunk)
1166    // dirty:          [SPARSE]        (bytes [5, 11))
1167    // expected: b"AAAA_SPARSE_AAAA|"  (17 bytes)
1168    //
1169    // Tests mid-file edit on a very small file (single chunk, smaller than
1170    // typical CDC minimum). See test_truncate_then_mid_edit for the regression
1171    // test that reproduces the real production bug.
1172    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1173    async fn test_upload_ranges_small_file_mid_edit() {
1174        let server = LocalTestServerBuilder::new().start().await;
1175        let base_dir = TempDir::new().unwrap();
1176        let config = test_config(server.http_endpoint(), base_dir.path());
1177        let cas_client: Arc<dyn Client> = Arc::new(server);
1178
1179        let original_data = b"AAAA_HEADER_AAAA|";
1180        let original_hash = upload_file(&config, original_data).await;
1181        let original_size = original_data.len() as u64;
1182
1183        let dirty_data = b"SPARSE";
1184        let dirty_inputs = vec![DirtyInput {
1185            original_range: 5..11,
1186            new_length: dirty_data.len() as u64,
1187            reader: Box::pin(Cursor::new(dirty_data.to_vec())),
1188        }];
1189
1190        let result = upload_ranges(config.clone(), cas_client.clone(), original_hash, original_size, dirty_inputs)
1191            .await
1192            .unwrap();
1193
1194        assert_eq!(result.file_size(), Some(original_size));
1195
1196        let downloaded = download_file(&config, MerkleHash::from_hex(result.hash()).unwrap(), original_size).await;
1197        assert_eq!(downloaded.len(), original_size as usize, "reconstructed size mismatch");
1198        assert_eq!(&downloaded[..5], b"AAAA_", "prefix from CAS");
1199        assert_eq!(&downloaded[5..11], b"SPARSE", "dirty range");
1200        assert_eq!(&downloaded[11..], b"_AAAA|", "suffix from CAS");
1201
1202        let expected = b"AAAA_SPARSE_AAAA|";
1203        let clean_hash = upload_file(&config, expected).await;
1204        assert_eq!(result.hash(), clean_hash.hex(), "hash mismatch with clean upload");
1205    }
1206
1207    // original: [=========================== 256 KB ===========================]
1208    // staging:  [0000000000000000000000000000] (all zeros, file never opened for write)
1209    // result:   [====== 100 KB from CAS =====]
1210    //                                         ^ cut here (mid-chunk)
1211    //
1212    // The boundary chunk bytes must come from CAS, not from the zero-filled staging.
1213    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1214    async fn test_upload_ranges_truncation_empty_staging() {
1215        let server = LocalTestServerBuilder::new().start().await;
1216        let base_dir = TempDir::new().unwrap();
1217        let config = test_config(server.http_endpoint(), base_dir.path());
1218        let cas_client: Arc<dyn Client> = Arc::new(server);
1219
1220        let original_data = random_data(77, 256 * 1024);
1221        let original_hash = upload_file(&config, &original_data).await;
1222        let original_size = original_data.len() as u64;
1223
1224        let truncated_size = 100_000u64;
1225
1226        let result = upload_ranges(
1227            config.clone(),
1228            cas_client.clone(),
1229            original_hash,
1230            original_size,
1231            make_legacy_inputs(&[], &[], original_size, truncated_size),
1232        )
1233        .await
1234        .unwrap();
1235
1236        assert_eq!(result.file_size(), Some(truncated_size));
1237
1238        let downloaded = download_file(&config, MerkleHash::from_hex(result.hash()).unwrap(), truncated_size).await;
1239        assert_eq!(downloaded.len(), truncated_size as usize);
1240        assert_eq!(
1241            &downloaded[..],
1242            &original_data[..truncated_size as usize],
1243            "truncated content should match original CAS data, not staging zeros"
1244        );
1245
1246        let clean_hash = upload_file(&config, &original_data[..truncated_size as usize]).await;
1247        assert_eq!(result.hash(), clean_hash.hex(), "hash mismatch with clean upload");
1248    }
1249
1250    // original: [=========================== 256 KB ===========================]
1251    // staging:  [000000000000][0xBB][0000000] (zeros except the dirty range)
1252    //                          ^  ^
1253    //                       90K  95K (dirty from caller)
1254    // result:   [==CAS==][stg][===CAS===]
1255    //                                    ^ cut at 100K (mid-chunk)
1256    //
1257    // Dirty bytes [90K,95K) come from staging; boundary bytes from CAS.
1258    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1259    async fn test_upload_ranges_truncation_with_overlapping_dirty() {
1260        let server = LocalTestServerBuilder::new().start().await;
1261        let base_dir = TempDir::new().unwrap();
1262        let config = test_config(server.http_endpoint(), base_dir.path());
1263        let cas_client: Arc<dyn Client> = Arc::new(server);
1264
1265        let original_data = random_data(88, 256 * 1024);
1266        let original_hash = upload_file(&config, &original_data).await;
1267        let original_size = original_data.len() as u64;
1268
1269        let truncated_size = 100_000u64;
1270
1271        let dirty_start = 90_000u64;
1272        let dirty_end = 95_000u64;
1273
1274        let mut expected = original_data[..truncated_size as usize].to_vec();
1275        expected[dirty_start as usize..dirty_end as usize].fill(0xBB);
1276
1277        let mut staging = vec![0u8; truncated_size as usize];
1278        staging[dirty_start as usize..dirty_end as usize].fill(0xBB);
1279        let result = upload_ranges(
1280            config.clone(),
1281            cas_client.clone(),
1282            original_hash,
1283            original_size,
1284            make_legacy_inputs(&[(dirty_start, dirty_end)], &staging, original_size, truncated_size),
1285        )
1286        .await
1287        .unwrap();
1288
1289        assert_eq!(result.file_size(), Some(truncated_size));
1290
1291        let downloaded = download_file(&config, MerkleHash::from_hex(result.hash()).unwrap(), truncated_size).await;
1292        assert_eq!(downloaded.len(), expected.len());
1293        assert_eq!(&downloaded[..], &expected[..], "dirty bytes should come from staging, boundary bytes from CAS");
1294
1295        let clean_hash = upload_file(&config, &expected).await;
1296        assert_eq!(result.hash(), clean_hash.hex(), "hash mismatch with clean upload");
1297    }
1298
1299    // ── Helpers ──────────────────────────────────────────────────────
1300
1301    fn random_data(seed: u64, len: usize) -> Vec<u8> {
1302        (0..len)
1303            .map(|i| {
1304                let x = (i as u64).wrapping_add(seed).wrapping_mul(2654435761);
1305                (x >> 16) as u8
1306            })
1307            .collect()
1308    }
1309
1310    #[derive(Clone, Debug)]
1311    struct DeterministicRng {
1312        state: u64,
1313    }
1314
1315    impl DeterministicRng {
1316        fn new(seed: u64) -> Self {
1317            Self { state: seed }
1318        }
1319
1320        fn next_u64(&mut self) -> u64 {
1321            self.state = self.state.wrapping_mul(6364136223846793005).wrapping_add(1);
1322            self.state
1323        }
1324
1325        fn gen_range(&mut self, start: usize, end: usize) -> usize {
1326            if end <= start {
1327                return start;
1328            }
1329            start + (self.next_u64() as usize % (end - start))
1330        }
1331
1332        fn gen_bytes(&mut self, len: usize) -> Vec<u8> {
1333            (0..len).map(|_| (self.next_u64() >> 56) as u8).collect()
1334        }
1335    }
1336
1337    #[derive(Clone, Debug)]
1338    struct PlannedEdit {
1339        original_range: Range<usize>,
1340        replacement: Vec<u8>,
1341    }
1342
1343    fn build_random_non_overlapping_edits(
1344        rng: &mut DeterministicRng,
1345        original_len: usize,
1346        max_edits: usize,
1347    ) -> Vec<PlannedEdit> {
1348        if original_len == 0 {
1349            let replacement_len = 1 + rng.gen_range(0, 8 * 1024);
1350            return vec![PlannedEdit {
1351                original_range: 0..0,
1352                replacement: rng.gen_bytes(replacement_len),
1353            }];
1354        }
1355
1356        let target_edits = 1 + rng.gen_range(0, max_edits.max(1));
1357        let mut edits: Vec<PlannedEdit> = Vec::with_capacity(target_edits);
1358        let mut cursor = 0usize;
1359
1360        while edits.len() < target_edits && cursor <= original_len {
1361            let remaining = original_len - cursor;
1362            let max_gap = remaining.min(64 * 1024);
1363            let start = cursor + rng.gen_range(0, max_gap + 1);
1364
1365            let (end, replacement_len) = if start == original_len {
1366                (start, 1 + rng.gen_range(0, 32 * 1024))
1367            } else {
1368                let op = rng.gen_range(0, 5);
1369                let max_span = (original_len - start).clamp(1, 64 * 1024);
1370                let span = 1 + rng.gen_range(0, max_span);
1371                let end = start + span;
1372                match op {
1373                    0 => (start, 1 + rng.gen_range(0, 32 * 1024)),
1374                    1 => (end, span),
1375                    2 => (end, span + 1 + rng.gen_range(0, 16 * 1024)),
1376                    3 => (end, rng.gen_range(0, span + 1)),
1377                    _ => (end, 0),
1378                }
1379            };
1380
1381            edits.push(PlannedEdit {
1382                original_range: start..end,
1383                replacement: rng.gen_bytes(replacement_len),
1384            });
1385            cursor = if end > start { end } else { start.saturating_add(1) };
1386        }
1387
1388        if edits.is_empty() {
1389            let replacement_len = 1 + rng.gen_range(0, 32 * 1024);
1390            edits.push(PlannedEdit {
1391                original_range: original_len..original_len,
1392                replacement: rng.gen_bytes(replacement_len),
1393            });
1394        }
1395
1396        for w in edits.windows(2) {
1397            assert!(w[0].original_range.end <= w[1].original_range.start);
1398        }
1399
1400        edits
1401    }
1402
1403    fn apply_planned_edits(original: &[u8], edits: &[PlannedEdit]) -> Vec<u8> {
1404        let removed: usize = edits.iter().map(|e| e.original_range.end - e.original_range.start).sum();
1405        let added: usize = edits.iter().map(|e| e.replacement.len()).sum();
1406        let mut out: Vec<u8> = Vec::with_capacity(original.len() + added.saturating_sub(removed));
1407        let mut cursor = 0usize;
1408
1409        for edit in edits {
1410            assert!(edit.original_range.start >= cursor);
1411            out.extend_from_slice(&original[cursor..edit.original_range.start]);
1412            out.extend_from_slice(&edit.replacement);
1413            cursor = edit.original_range.end;
1414        }
1415
1416        out.extend_from_slice(&original[cursor..]);
1417        out
1418    }
1419
1420    fn edits_to_dirty_inputs(edits: &[PlannedEdit]) -> Vec<DirtyInput> {
1421        edits
1422            .iter()
1423            .map(|e| DirtyInput {
1424                original_range: e.original_range.start as u64..e.original_range.end as u64,
1425                new_length: e.replacement.len() as u64,
1426                reader: Box::pin(Cursor::new(e.replacement.clone())),
1427            })
1428            .collect()
1429    }
1430
1431    fn summarize_edits(edits: &[PlannedEdit]) -> String {
1432        edits
1433            .iter()
1434            .map(|e| format!("[{}..{}, new_len={}]", e.original_range.start, e.original_range.end, e.replacement.len()))
1435            .collect::<Vec<_>>()
1436            .join(", ")
1437    }
1438
1439    async fn upload_file(config: &Arc<TranslatorConfig>, data: &[u8]) -> MerkleHash {
1440        let session = FileUploadSession::new(config.clone()).await.unwrap();
1441        let (_id, mut cleaner) = session
1442            .start_clean(Some("test".into()), Some(data.len() as u64), Sha256Policy::Skip)
1443            .unwrap();
1444        cleaner.add_data(data).await.unwrap();
1445        let (xfi, _metrics) = cleaner.finish().await.unwrap();
1446        session.finalize().await.unwrap();
1447        MerkleHash::from_hex(xfi.hash()).unwrap()
1448    }
1449
1450    async fn download_file(config: &Arc<TranslatorConfig>, hash: MerkleHash, size: u64) -> Vec<u8> {
1451        let session = FileDownloadSession::new(config.clone(), None).await.unwrap();
1452        let xfi = crate::processing::XetFileInfo::new(hash.hex(), size);
1453        let dir = TempDir::new().unwrap();
1454        let out = dir.path().join("out");
1455        session.download_file(&xfi, &out).await.unwrap();
1456        std::fs::read(&out).unwrap()
1457    }
1458
1459    /// Empty `Box::pin(Cursor::new(Vec::new()))` for delete / dummy edits.
1460    fn empty_reader() -> Pin<Box<dyn AsyncRead + Send>> {
1461        Box::pin(Cursor::new(Vec::<u8>::new()))
1462    }
1463
1464    /// End-to-end check: upload `original`, apply `inputs` via `upload_ranges`, download,
1465    /// compare against `expected`, and verify the hash matches a clean upload of `expected`.
1466    async fn assert_edits(
1467        config: &Arc<TranslatorConfig>,
1468        cas_client: &Arc<dyn Client>,
1469        original: &[u8],
1470        inputs: Vec<DirtyInput>,
1471        expected: &[u8],
1472    ) {
1473        let original_hash = upload_file(config, original).await;
1474        let original_size = original.len() as u64;
1475        let result = upload_ranges(config.clone(), cas_client.clone(), original_hash, original_size, inputs)
1476            .await
1477            .unwrap();
1478        assert_eq!(result.file_size(), Some(expected.len() as u64), "file size mismatch");
1479        let downloaded =
1480            download_file(config, MerkleHash::from_hex(result.hash()).unwrap(), expected.len() as u64).await;
1481        assert_eq!(downloaded, expected, "content mismatch");
1482        let clean = upload_file(config, expected).await;
1483        assert_eq!(result.hash(), clean.hex(), "hash diverges from clean upload");
1484    }
1485
1486    async fn assert_range_edit(
1487        config: &Arc<TranslatorConfig>,
1488        cas_client: &Arc<dyn Client>,
1489        original_data: &[u8],
1490        expected: &[u8],
1491        dirty_ranges: &[(u64, u64)],
1492        total_size: u64,
1493    ) {
1494        let original_hash = upload_file(config, original_data).await;
1495        let original_size = original_data.len() as u64;
1496
1497        // Build dirty inputs from the caller's ranges.
1498        let mut inputs = make_dirty_inputs(dirty_ranges, expected);
1499
1500        // For appends, ensure the appended region is included as a pure-insert edit at the
1501        // end of the original.
1502        if total_size > original_size {
1503            let append_start = original_size;
1504            let already_covered = dirty_ranges.iter().any(|&(s, e)| s <= append_start && e >= total_size);
1505            if !already_covered {
1506                inputs.push(DirtyInput {
1507                    original_range: original_size..original_size,
1508                    new_length: total_size - original_size,
1509                    reader: Box::pin(Cursor::new(expected[append_start as usize..total_size as usize].to_vec())),
1510                });
1511                inputs.sort_by_key(|d| d.original_range.start);
1512            }
1513        }
1514
1515        // Truncation: drop bytes past `total_size` from the original via a pure-delete edit.
1516        if total_size < original_size {
1517            inputs.push(DirtyInput {
1518                original_range: total_size..original_size,
1519                new_length: 0,
1520                reader: empty_reader(),
1521            });
1522            inputs.sort_by_key(|d| d.original_range.start);
1523        }
1524
1525        let result = upload_ranges(config.clone(), cas_client.clone(), original_hash, original_size, inputs)
1526            .await
1527            .unwrap();
1528
1529        assert_eq!(result.file_size(), Some(total_size), "file size mismatch");
1530        let downloaded = download_file(config, MerkleHash::from_hex(result.hash()).unwrap(), total_size).await;
1531        assert_eq!(downloaded.len(), expected.len(), "downloaded length mismatch");
1532        assert_eq!(&downloaded[..], expected, "content mismatch");
1533
1534        let clean_hash = upload_file(config, expected).await;
1535        assert_eq!(result.hash(), clean_hash.hex(), "hash mismatch with clean upload");
1536    }
1537
1538    // Regression: mid-file edit + tail append in a single call. Codex caught that the last
1539    // existing window (covering the mid-file edit) was being stretched to `total_size`,
1540    // dropping the stable bytes between the edit and the appended region.
1541    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1542    async fn test_mid_edit_plus_append() {
1543        let server = LocalTestServerBuilder::new().start().await;
1544        let base_dir = TempDir::new().unwrap();
1545        let config = test_config(server.http_endpoint(), base_dir.path());
1546        let cas_client: Arc<dyn Client> = Arc::new(server);
1547
1548        let original_data = random_data(7, 256 * 1024);
1549        let original_hash = upload_file(&config, &original_data).await;
1550        let original_size = original_data.len() as u64;
1551
1552        let dirty_start = 50_000usize;
1553        let dirty_end = 51_000usize;
1554        let append_extra: Vec<u8> = (0..16 * 1024).map(|i| (i % 251) as u8).collect();
1555        let mut expected = original_data.clone();
1556        expected[dirty_start..dirty_end].fill(0xAA);
1557        expected.extend_from_slice(&append_extra);
1558        let total_size = expected.len() as u64;
1559
1560        let inputs = vec![
1561            DirtyInput {
1562                original_range: dirty_start as u64..dirty_end as u64,
1563                new_length: (dirty_end - dirty_start) as u64,
1564                reader: Box::pin(Cursor::new(expected[dirty_start..dirty_end].to_vec())),
1565            },
1566            DirtyInput {
1567                original_range: original_size..original_size,
1568                new_length: append_extra.len() as u64,
1569                reader: Box::pin(Cursor::new(append_extra)),
1570            },
1571        ];
1572        let result = upload_ranges(config.clone(), cas_client.clone(), original_hash, original_size, inputs)
1573            .await
1574            .unwrap();
1575
1576        let downloaded = download_file(&config, MerkleHash::from_hex(result.hash()).unwrap(), total_size).await;
1577        assert_eq!(downloaded, expected, "content mismatch (mid-edit + append regression)");
1578        let clean_hash = upload_file(&config, &expected).await;
1579        assert_eq!(result.hash(), clean_hash.hex(), "hash mismatch with clean upload");
1580    }
1581
1582    // Regression: empty original + append. Codex caught that we'd return an internal error
1583    // instead of treating it as a fresh upload.
1584    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1585    async fn test_empty_original_append() {
1586        let server = LocalTestServerBuilder::new().start().await;
1587        let base_dir = TempDir::new().unwrap();
1588        let config = test_config(server.http_endpoint(), base_dir.path());
1589        let cas_client: Arc<dyn Client> = Arc::new(server);
1590
1591        let original_data: &[u8] = &[];
1592        let original_hash = upload_file(&config, original_data).await;
1593        let new_data: Vec<u8> = (0..32 * 1024).map(|i| (i % 251) as u8).collect();
1594        let total_size = new_data.len() as u64;
1595
1596        let inputs = vec![DirtyInput {
1597            original_range: 0..0,
1598            new_length: total_size,
1599            reader: Box::pin(Cursor::new(new_data.clone())),
1600        }];
1601        let result = upload_ranges(config.clone(), cas_client.clone(), original_hash, 0, inputs)
1602            .await
1603            .unwrap();
1604
1605        let downloaded = download_file(&config, MerkleHash::from_hex(result.hash()).unwrap(), total_size).await;
1606        assert_eq!(downloaded, new_data);
1607        let clean_hash = upload_file(&config, &new_data).await;
1608        assert_eq!(result.hash(), clean_hash.hex(), "hash mismatch with clean upload");
1609    }
1610
1611    // Regression: truncating to empty must produce the canonical empty-file hash
1612    // (`MerkleHash::default()` without HMAC), not `default.hmac(default)`.
1613    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1614    async fn test_truncate_to_empty_matches_clean_empty() {
1615        let server = LocalTestServerBuilder::new().start().await;
1616        let base_dir = TempDir::new().unwrap();
1617        let config = test_config(server.http_endpoint(), base_dir.path());
1618        let cas_client: Arc<dyn Client> = Arc::new(server);
1619
1620        let original_data = random_data(11, 64 * 1024);
1621        let original_hash = upload_file(&config, &original_data).await;
1622        let original_size = original_data.len() as u64;
1623
1624        let result = upload_ranges(
1625            config.clone(),
1626            cas_client.clone(),
1627            original_hash,
1628            original_size,
1629            make_legacy_inputs(&[], &[], original_size, 0),
1630        )
1631        .await
1632        .unwrap();
1633
1634        let clean_empty = upload_file(&config, &[]).await;
1635        assert_eq!(result.hash(), clean_empty.hex(), "truncate-to-empty must match clean empty upload hash");
1636    }
1637
1638    // The three small examples that motivated the resize-edit API: replace, insert, delete
1639    // on a tiny 3-byte file. Each produces an output of a different length than the input.
1640    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1641    async fn test_resize_edits_abc() {
1642        let server = LocalTestServerBuilder::new().start().await;
1643        let base_dir = TempDir::new().unwrap();
1644        let config = test_config(server.http_endpoint(), base_dir.path());
1645        let cas_client: Arc<dyn Client> = Arc::new(server);
1646
1647        // abc + replace [0, 1) with "foo" => "foobc"
1648        assert_edits(
1649            &config,
1650            &cas_client,
1651            b"abc",
1652            vec![DirtyInput {
1653                original_range: 0..1,
1654                new_length: 3,
1655                reader: Box::pin(Cursor::new(b"foo".to_vec())),
1656            }],
1657            b"foobc",
1658        )
1659        .await;
1660
1661        // abc + insert "foo" at 0 => "fooabc"
1662        assert_edits(
1663            &config,
1664            &cas_client,
1665            b"abc",
1666            vec![DirtyInput {
1667                original_range: 0..0,
1668                new_length: 3,
1669                reader: Box::pin(Cursor::new(b"foo".to_vec())),
1670            }],
1671            b"fooabc",
1672        )
1673        .await;
1674
1675        // abc + delete [0, 1) => "bc"
1676        assert_edits(
1677            &config,
1678            &cas_client,
1679            b"abc",
1680            vec![DirtyInput {
1681                original_range: 0..1,
1682                new_length: 0,
1683                reader: empty_reader(),
1684            }],
1685            b"bc",
1686        )
1687        .await;
1688    }
1689
1690    // original: [============= 256 KB =============]
1691    // edit:                [== 4K stale ==]
1692    //                      ^ replaced with 32 K (resize +28 K)
1693    // result:   [== prefix ==][== 32K new ==][== suffix ==]
1694    //
1695    // Big in-place replace where new_length >> original_range.len(); exercises the
1696    // resize branch with a multi-segment-touching window.
1697    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1698    async fn test_resize_large_replace_grows_file() {
1699        let server = LocalTestServerBuilder::new().start().await;
1700        let base_dir = TempDir::new().unwrap();
1701        let config = test_config(server.http_endpoint(), base_dir.path());
1702        let cas_client: Arc<dyn Client> = Arc::new(server);
1703
1704        let original = random_data(101, 256 * 1024);
1705        let drop_start = 100_000usize;
1706        let drop_end = 104_000usize;
1707        let new_bytes: Vec<u8> = (0..32 * 1024).map(|i| (i % 251) as u8).collect();
1708        let mut expected = Vec::with_capacity(original.len() - (drop_end - drop_start) + new_bytes.len());
1709        expected.extend_from_slice(&original[..drop_start]);
1710        expected.extend_from_slice(&new_bytes);
1711        expected.extend_from_slice(&original[drop_end..]);
1712
1713        assert_edits(
1714            &config,
1715            &cas_client,
1716            &original,
1717            vec![DirtyInput {
1718                original_range: drop_start as u64..drop_end as u64,
1719                new_length: new_bytes.len() as u64,
1720                reader: Box::pin(Cursor::new(new_bytes)),
1721            }],
1722            &expected,
1723        )
1724        .await;
1725    }
1726
1727    // original: [============= 256 KB =============]
1728    // edit:                [============= 80 KB stale =============]
1729    //                      ^ replaced with 4 K (resize -76 K)
1730    // result:   [== prefix ==][4 K new][== suffix ==]
1731    //
1732    // Big in-place replace where new_length << original_range.len(); shrink case.
1733    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1734    async fn test_resize_large_replace_shrinks_file() {
1735        let server = LocalTestServerBuilder::new().start().await;
1736        let base_dir = TempDir::new().unwrap();
1737        let config = test_config(server.http_endpoint(), base_dir.path());
1738        let cas_client: Arc<dyn Client> = Arc::new(server);
1739
1740        let original = random_data(102, 256 * 1024);
1741        let drop_start = 80_000usize;
1742        let drop_end = 160_000usize;
1743        let new_bytes: Vec<u8> = (0..4 * 1024).map(|i| (0xCC ^ i) as u8).collect();
1744        let mut expected = Vec::with_capacity(original.len() - (drop_end - drop_start) + new_bytes.len());
1745        expected.extend_from_slice(&original[..drop_start]);
1746        expected.extend_from_slice(&new_bytes);
1747        expected.extend_from_slice(&original[drop_end..]);
1748
1749        assert_edits(
1750            &config,
1751            &cas_client,
1752            &original,
1753            vec![DirtyInput {
1754                original_range: drop_start as u64..drop_end as u64,
1755                new_length: new_bytes.len() as u64,
1756                reader: Box::pin(Cursor::new(new_bytes)),
1757            }],
1758            &expected,
1759        )
1760        .await;
1761    }
1762
1763    // original: [============= 100 KB =============]
1764    // edit:                  [insert 8 KB here]
1765    // result:   [== prefix ==][== 8 KB new ==][======= suffix =======]
1766    //
1767    // Pure mid-file insert (range start == end), large payload that forces re-chunking
1768    // around the insertion point.
1769    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1770    async fn test_resize_mid_file_insert() {
1771        let server = LocalTestServerBuilder::new().start().await;
1772        let base_dir = TempDir::new().unwrap();
1773        let config = test_config(server.http_endpoint(), base_dir.path());
1774        let cas_client: Arc<dyn Client> = Arc::new(server);
1775
1776        let original = random_data(103, 100 * 1024);
1777        let at = 40_000usize;
1778        let new_bytes: Vec<u8> = (0..8u32 * 1024).map(|i| (i.wrapping_mul(13) % 251) as u8).collect();
1779        let mut expected = Vec::with_capacity(original.len() + new_bytes.len());
1780        expected.extend_from_slice(&original[..at]);
1781        expected.extend_from_slice(&new_bytes);
1782        expected.extend_from_slice(&original[at..]);
1783
1784        assert_edits(
1785            &config,
1786            &cas_client,
1787            &original,
1788            vec![DirtyInput {
1789                original_range: at as u64..at as u64,
1790                new_length: new_bytes.len() as u64,
1791                reader: Box::pin(Cursor::new(new_bytes)),
1792            }],
1793            &expected,
1794        )
1795        .await;
1796    }
1797
1798    // original: [============= 256 KB =============]
1799    // edit:                [== 64 KB hole ==]
1800    //                      ^ delete this slice (no replacement)
1801    // result:   [== prefix ==][== suffix ==]
1802    //
1803    // Pure mid-file delete: shrinks the file without touching the tail boundary.
1804    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1805    async fn test_resize_mid_file_delete() {
1806        let server = LocalTestServerBuilder::new().start().await;
1807        let base_dir = TempDir::new().unwrap();
1808        let config = test_config(server.http_endpoint(), base_dir.path());
1809        let cas_client: Arc<dyn Client> = Arc::new(server);
1810
1811        let original = random_data(104, 256 * 1024);
1812        let drop_start = 80_000usize;
1813        let drop_end = 144_000usize;
1814        let mut expected = Vec::with_capacity(original.len() - (drop_end - drop_start));
1815        expected.extend_from_slice(&original[..drop_start]);
1816        expected.extend_from_slice(&original[drop_end..]);
1817
1818        assert_edits(
1819            &config,
1820            &cas_client,
1821            &original,
1822            vec![DirtyInput {
1823                original_range: drop_start as u64..drop_end as u64,
1824                new_length: 0,
1825                reader: empty_reader(),
1826            }],
1827            &expected,
1828        )
1829        .await;
1830    }
1831
1832    // original: [== seg0 ==][== seg1 ==][== seg2 ==]
1833    // edits:        [shrink][grow][   delete   ]
1834    //
1835    // Three independent edits in one call: an in-place shrink, a pure insert in seg1, and
1836    // a pure delete in seg2. Mix of resize directions, far enough apart that they don't
1837    // coalesce into one window.
1838    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1839    async fn test_resize_multi_edit_mix() {
1840        let server = LocalTestServerBuilder::new().start().await;
1841        let base_dir = TempDir::new().unwrap();
1842        let config = test_config(server.http_endpoint(), base_dir.path());
1843        let cas_client: Arc<dyn Client> = Arc::new(server);
1844
1845        let original = random_data(105, 384 * 1024);
1846        let (a_start, a_end) = (10 * 1024usize, 20 * 1024usize);
1847        let a_new: Vec<u8> = vec![0xAA; 2 * 1024];
1848        let b_at = 150 * 1024usize;
1849        let b_new: Vec<u8> = vec![0xBB; 4 * 1024];
1850        let (c_start, c_end) = (300 * 1024usize, 320 * 1024usize);
1851
1852        let mut expected = Vec::with_capacity(original.len() + b_new.len());
1853        expected.extend_from_slice(&original[..a_start]);
1854        expected.extend_from_slice(&a_new);
1855        expected.extend_from_slice(&original[a_end..b_at]);
1856        expected.extend_from_slice(&b_new);
1857        expected.extend_from_slice(&original[b_at..c_start]);
1858        expected.extend_from_slice(&original[c_end..]);
1859
1860        assert_edits(
1861            &config,
1862            &cas_client,
1863            &original,
1864            vec![
1865                DirtyInput {
1866                    original_range: a_start as u64..a_end as u64,
1867                    new_length: a_new.len() as u64,
1868                    reader: Box::pin(Cursor::new(a_new)),
1869                },
1870                DirtyInput {
1871                    original_range: b_at as u64..b_at as u64,
1872                    new_length: b_new.len() as u64,
1873                    reader: Box::pin(Cursor::new(b_new)),
1874                },
1875                DirtyInput {
1876                    original_range: c_start as u64..c_end as u64,
1877                    new_length: 0,
1878                    reader: empty_reader(),
1879                },
1880            ],
1881            &expected,
1882        )
1883        .await;
1884    }
1885
1886    // original: [== seg0 ==][== seg1 ==]
1887    //                       ^ insert here, exactly on segment boundary
1888    // result:   [== seg0 ==][== inserted ==][== seg1 ==]
1889    //
1890    // Pure insert on a segment boundary mid-file. Snap picks the segment starting at the
1891    // boundary; insert lands at `w_start` of that window.
1892    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1893    async fn test_resize_insert_at_segment_boundary() {
1894        let server = LocalTestServerBuilder::new().start().await;
1895        let base_dir = TempDir::new().unwrap();
1896        let config = test_config(server.http_endpoint(), base_dir.path());
1897        let cas_client: Arc<dyn Client> = Arc::new(server);
1898
1899        let original = random_data(106, 200 * 1024);
1900        let original_hash = upload_file(&config, &original).await;
1901        let original_size = original.len() as u64;
1902
1903        // Pick the first interior segment boundary; bail if the file lands in one segment.
1904        let seg_sizes = fetch_segment_sizes(&cas_client, &original_hash).await;
1905        let Some(boundary) = seg_sizes
1906            .iter()
1907            .scan(0u64, |acc, s| {
1908                *acc += s;
1909                Some(*acc)
1910            })
1911            .find(|&b| b > 0 && b < original_size)
1912        else {
1913            return;
1914        };
1915
1916        let new_bytes: Vec<u8> = vec![0x42; 4 * 1024];
1917        let mut expected = Vec::with_capacity(original.len() + new_bytes.len());
1918        expected.extend_from_slice(&original[..boundary as usize]);
1919        expected.extend_from_slice(&new_bytes);
1920        expected.extend_from_slice(&original[boundary as usize..]);
1921
1922        assert_edits(
1923            &config,
1924            &cas_client,
1925            &original,
1926            vec![DirtyInput {
1927                original_range: boundary..boundary,
1928                new_length: new_bytes.len() as u64,
1929                reader: Box::pin(Cursor::new(new_bytes)),
1930            }],
1931            &expected,
1932        )
1933        .await;
1934    }
1935
1936    #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
1937    #[ignore = "stress test"]
1938    async fn test_stress_random_resize_sequences() {
1939        let server = LocalTestServerBuilder::new().start().await;
1940        let base_dir = TempDir::new().unwrap();
1941        let config = test_config(server.http_endpoint(), base_dir.path());
1942        let cas_client: Arc<dyn Client> = Arc::new(server);
1943
1944        for seed in 0..6u64 {
1945            let mut rng = DeterministicRng::new(0x9E37_79B9_7F4A_7C15 ^ seed.wrapping_mul(0xD1B5_4A32_D192_ED03));
1946            let mut expected = random_data(10_000 + seed, 1_048_576 + (seed as usize * 91_117 % 262_144));
1947            let mut original_hash = upload_file(&config, &expected).await;
1948            let mut original_size = expected.len() as u64;
1949
1950            for round in 0..25usize {
1951                let edits = build_random_non_overlapping_edits(&mut rng, expected.len(), 8);
1952                let expected_next = apply_planned_edits(&expected, &edits);
1953                let inputs = edits_to_dirty_inputs(&edits);
1954                let result = upload_ranges(config.clone(), cas_client.clone(), original_hash, original_size, inputs)
1955                    .await
1956                    .unwrap();
1957                let result_hash = MerkleHash::from_hex(result.hash()).unwrap();
1958
1959                assert_eq!(
1960                    result.file_size(),
1961                    Some(expected_next.len() as u64),
1962                    "seed={seed}, round={round}: size mismatch"
1963                );
1964                let clean_hash = upload_file(&config, &expected_next).await;
1965                assert_eq!(result.hash(), clean_hash.hex(), "seed={seed}, round={round}: hash mismatch");
1966
1967                let downloaded = download_file(&config, result_hash, expected_next.len() as u64).await;
1968                assert_eq!(downloaded, expected_next, "seed={seed}, round={round}: content mismatch");
1969
1970                expected = expected_next;
1971                original_hash = result_hash;
1972                original_size = expected.len() as u64;
1973            }
1974        }
1975    }
1976
1977    #[cfg(not(feature = "smoke-test"))]
1978    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1979    async fn test_regression_hash_matches_clean_upload_seed1_round17() {
1980        let server = LocalTestServerBuilder::new().start().await;
1981        let base_dir = TempDir::new().unwrap();
1982        let config = test_config(server.http_endpoint(), base_dir.path());
1983        let cas_client: Arc<dyn Client> = Arc::new(server);
1984
1985        let seed = 1u64;
1986        let mut rng = DeterministicRng::new(0x9E37_79B9_7F4A_7C15 ^ seed.wrapping_mul(0xD1B5_4A32_D192_ED03));
1987        let mut expected = random_data(10_000 + seed, 1_048_576 + (seed as usize * 91_117 % 262_144));
1988        let mut original_hash = upload_file(&config, &expected).await;
1989        let mut original_size = expected.len() as u64;
1990
1991        for round in 0..=17usize {
1992            let edits = build_random_non_overlapping_edits(&mut rng, expected.len(), 8);
1993            let edits_summary = summarize_edits(&edits);
1994            let expected_next = apply_planned_edits(&expected, &edits);
1995            let inputs = edits_to_dirty_inputs(&edits);
1996            let result = upload_ranges(config.clone(), cas_client.clone(), original_hash, original_size, inputs)
1997                .await
1998                .unwrap();
1999            let result_hash = MerkleHash::from_hex(result.hash()).unwrap();
2000
2001            let clean_hash = upload_file(&config, &expected_next).await;
2002            assert_eq!(
2003                result.hash(),
2004                clean_hash.hex(),
2005                "seed={seed}, round={round}: hash mismatch; original_size={original_size}, expected_size={}, edits={edits_summary}",
2006                expected_next.len()
2007            );
2008
2009            let downloaded = download_file(&config, result_hash, expected_next.len() as u64).await;
2010            assert_eq!(downloaded, expected_next, "seed={seed}, round={round}: content mismatch");
2011
2012            expected = expected_next;
2013            original_hash = result_hash;
2014            original_size = expected.len() as u64;
2015        }
2016    }
2017
2018    #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
2019    #[ignore = "stress test"]
2020    async fn test_stress_many_sparse_windows_single_call() {
2021        let server = LocalTestServerBuilder::new().start().await;
2022        let base_dir = TempDir::new().unwrap();
2023        let config = test_config(server.http_endpoint(), base_dir.path());
2024        let cas_client: Arc<dyn Client> = Arc::new(server);
2025
2026        let original = random_data(13_337, 16 * 1024 * 1024);
2027        let mut rng = DeterministicRng::new(0xA5A5_5A5A_0123_4567);
2028        let mut edits: Vec<PlannedEdit> = Vec::new();
2029        let stride = original.len() / 200;
2030        let mut cursor = stride / 2;
2031
2032        while edits.len() < 128 && cursor < original.len() {
2033            let start = cursor;
2034            let max_span = (original.len() - start).clamp(1, 1536);
2035            let span = 128 + rng.gen_range(0, max_span);
2036            let end = (start + span).min(original.len());
2037            let replacement_len = match rng.gen_range(0, 4) {
2038                0 => end - start,
2039                1 => (end - start) + 64 + rng.gen_range(0, 512),
2040                2 => rng.gen_range(0, end - start + 1),
2041                _ => 0,
2042            };
2043            edits.push(PlannedEdit {
2044                original_range: start..end,
2045                replacement: rng.gen_bytes(replacement_len),
2046            });
2047            cursor = cursor.saturating_add(stride.max(1));
2048        }
2049
2050        let expected = apply_planned_edits(&original, &edits);
2051        let inputs = edits_to_dirty_inputs(&edits);
2052        assert_edits(&config, &cas_client, &original, inputs, &expected).await;
2053    }
2054
2055    #[tokio::test(flavor = "multi_thread", worker_threads = 8)]
2056    #[ignore = "stress test"]
2057    async fn test_stress_parallel_random_resize_sequences() {
2058        let server = LocalTestServerBuilder::new().start().await;
2059        let base_dir = TempDir::new().unwrap();
2060        let config = test_config(server.http_endpoint(), base_dir.path());
2061        let cas_client: Arc<dyn Client> = Arc::new(server);
2062
2063        let mut handles = Vec::new();
2064        for worker in 0..8u64 {
2065            let config = config.clone();
2066            let cas_client = cas_client.clone();
2067            handles.push(tokio::spawn(async move {
2068                let mut rng = DeterministicRng::new(0xC0FF_EE00_1234_5678 ^ worker.wrapping_mul(0x94D0_49BB_1331_11EB));
2069                let mut expected = random_data(20_000 + worker, 786_432 + worker as usize * 17_321);
2070                let mut original_hash = upload_file(&config, &expected).await;
2071                let mut original_size = expected.len() as u64;
2072
2073                for round in 0..18usize {
2074                    let edits = build_random_non_overlapping_edits(&mut rng, expected.len(), 6);
2075                    let expected_next = apply_planned_edits(&expected, &edits);
2076                    let inputs = edits_to_dirty_inputs(&edits);
2077                    let result =
2078                        upload_ranges(config.clone(), cas_client.clone(), original_hash, original_size, inputs)
2079                            .await
2080                            .unwrap();
2081                    let result_hash = MerkleHash::from_hex(result.hash()).unwrap();
2082
2083                    assert_eq!(
2084                        result.file_size(),
2085                        Some(expected_next.len() as u64),
2086                        "worker={worker}, round={round}: size mismatch"
2087                    );
2088                    let clean_hash = upload_file(&config, &expected_next).await;
2089                    assert_eq!(result.hash(), clean_hash.hex(), "worker={worker}, round={round}: hash mismatch");
2090
2091                    let downloaded = download_file(&config, result_hash, expected_next.len() as u64).await;
2092                    assert_eq!(downloaded, expected_next, "worker={worker}, round={round}: content mismatch");
2093
2094                    expected = expected_next;
2095                    original_hash = result_hash;
2096                    original_size = expected.len() as u64;
2097                }
2098            }));
2099        }
2100
2101        for handle in handles {
2102            handle.await.unwrap();
2103        }
2104    }
2105}