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
22pub struct DirtyInput {
34 pub original_range: Range<u64>,
35 pub reader: Pin<Box<dyn AsyncRead + Send>>,
36 pub new_length: u64,
37}
38
39const STREAM_BLOCK_SIZE: usize = 4 * 1024 * 1024; struct UploadedWindow {
46 start: u64,
47 end: u64,
48 chunks: ChunkHashList,
49 mdb: MDBFileInfo,
50}
51
52pub 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 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 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 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 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 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 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 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 let mut hash_ranges = response.hash_ranges;
275 let trailing_gap = hash_ranges.pop().flatten();
276 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 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
336fn 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 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
422fn 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
455fn 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
470async 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
497async 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
514fn 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
521fn 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(all(test, feature = "simulation"))]
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 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 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 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 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 #[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 let cas_client: Arc<dyn Client> = Arc::new(server);
644
645 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 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 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 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 #[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 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 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 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 #[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 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 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 #[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 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 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 #[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 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 let mut modified_data = original_data.clone();
826 modified_data[10_000..12_000].fill(0xBB); modified_data[200_000..202_000].fill(0xCC); 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 #[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 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 #[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 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 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 {
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 {
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 {
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 {
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 {
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 #[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 #[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 #[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 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 #[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 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 #[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 #[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 #[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 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 fn empty_reader() -> Pin<Box<dyn AsyncRead + Send>> {
1461 Box::pin(Cursor::new(Vec::<u8>::new()))
1462 }
1463
1464 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 let mut inputs = make_dirty_inputs(dirty_ranges, expected);
1499
1500 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 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 #[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 #[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 #[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 #[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 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 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 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 #[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 #[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 #[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 #[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 #[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 #[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 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}