Squashed from work: - asset-ai: FastH3 4-step fast video backend; clip keyframes on the wire - asset-ui: loop video chains — text→image→video that ends where it began - h3: safetensors -> pruned-Q4_K GGUF quantizer for the 24GB DiT tiers - h3_quant_gguf verify: row-error gates calibrated to the measured Q4_K floor - asset-ai realtime: the feedback loop — the source anchors, the drifted frame inits - asset-ai realtime: a feedback loop survives a resize and travels by default - asset-ai realtime: the feedback loop frees itself from the feed handshake and pauses for its listener - asset-ai realtime: the outbound encode leaves the loop's critical path - asset-ai ocr: the ocr domain — Chandra 2 at page resolution, and the tower goes planner-owned - llm slots: a lane can hold an image span — embedding prefill and a rope cursor of its own - vision tower on CUDA: the encode leg gets its two missing kernels - llm/ocr: one M-RoPE grid encoder for both image paths, and a livelock made an error - vision tower on CUDA: the f16 GEMM keeps the precision it was throwing away - live: a feed that moves box takes its trip with it — one seed image - vision tower on CUDA: the tiled attention becomes bit-exact, and tensor cores go - llm prefill on CUDA: the MMA attention kernel gets the tile a 4-to-1 model needs - asset-ai ocr: the CUDA encode lane joins the integration — vision-parity sits beside run's three arms, and the kernels - Merge branch 'ocr-perf-integration' into work - asset-ai: the live anchor can follow the trip, and text leaves the 5090 - asset-ai: the camera moves the world, and the world starts still - asset-import: the EA strategy classics, in the one 2D contract - rtsmap: one seeded generator for tiled strategy maps - asset-ui: one card for the strategy classics, with a pack dropdown - asset-ai: music3 reference-audio path, ocr/h3 backends, registry - asset: mp4 sample index for range-streaming, chat tools, import profiles - cnc: tiberium is twelve growth frames, not twelve empty variants - platform: native file and save dialogs, in-house on all three desktops - chat: the scan holds out for a lane home - chat: a full home queues you — take the free lane - chat: the preload has a percentage, and the boundless cap stops showing - llm cuda: the 32x2 attention tile — even GQA ratios stay on MMA - sa3 gets a bake path: the sfx model's tables precomputed by a diffusion-side bin - sqlite_query: anti-join regression test - td import: HARV's second frame block is its harvesting cycle, not a turret - asset-ui: sprite enhancement runs on the 32B dev DiT — distillation, not the prompt, was the ceiling - ai-hub: makepad-asset-ai becomes makepad-ai-hub at libs/ai/hub, the chat pane becomes makepad-chat-ui, the service bin - asset-ui: test health fixtures grow the realtime field they were born without - ai-hub: one home at ~/.makepad — weights/ run/ cache/ logs/, the service cache migrates from ai_content by a single re - ai-hub: subprocess workers die with the node — process groups everywhere, PDEATHSIG on linux, one KILL_ON_JOB_CLOSE Jo - ai-hub: the hub object — AiHub::in_process, pipes vocabulary, and the local LLM engine generalized out of mpfiles (aic - strict-json: the dependency-free JSON module gets its own crate; asset-client re-exports it so nothing downstream move - ai-hub: the machine layer — node entries, the 0600 machine token, and the residency election that IS the lock (aicore - ai-hub: MPHUB1 — the fabric beacon only dedicated nodes can send (aicore §4) - ai-hub: job leases — work lives only while it is renewed (aicore §8) - asset-creator: the pipeline library is born — specs, the deps gate, and the derived-state law (aicore §9) - ai-hub: RAM residency facts — the CPU-side twin of residency.rs (aicore §3) - ai-hub: ETA placement primitives — relative GPU throughput, the four-term estimate, and an observable breakdown (aicor - ai-hub: leases go live on the wire — origin fields on submit, /job/<id>/keepalive, /bye, and the reaper that cancels w - ai-hub: the chat providers move in — fleet qwen, openai, grok, claude/codex/grok CLIs, the responses driver, and the w - asset-creator: the engine — one pipeline run against the hub, deps-gated, spliced, cancellable, resumable-by-construct - ai-hub: the machine node mode — --machine binds loopback, registers in ~/.makepad/run, and exits on its own once idle - asset-creator: makepad-creator-run — the detached client for runs that must outlive a window (aicore §9) - ai-hub: a native Claude Messages-API provider — API-key or Claude Code OAuth, bounded SSE streaming, injected tools (a - route + converse: off makepad_ai — the Agent seam moves to converse, route's cloud dispatcher rides the hub's Claude p - asset-creator: the preset tables move in — fifteen chain-policy constants shared by every creator app (aicore §9 / P6) - makepad_ai is deleted — every backend is a hub pipe, the agent seam lives with its consumers (aicore §14, decided 2026 - ai-hub: loads hold the machine residency election — set_model_state claims on Loaded and publishes the service port (a - ai-hub: chats run the machine election — route to a serving holder, wait on a loading one, claim and publish when open - ai-hub: pick_for_domain_eta — ETA-ranked placement over the shared hard-filter core (aicore §6 / P4) - asset-creator: the engine picks a provider per stage at dispatch time — a chain's later stages see fresh fleet state ( - ai-hub: the fabric secret gates the service HTTP surface — bearer on everything but /health and the ticketed peer path - vj: DREAM runs execute in the app — pipelines.rs becomes the run it used to watch (aicore §9 / F1) - asset-creator: the runner — generate one thing and put it in the catalog, one implementation for every surface (aicore - chat-ui: the session runs in the app — no broker anywhere on the chat path (aicore P8 / F5) - asset-store: assets.query is a first-class query endpoint — the bounded SQL surface outlives the broker (aicore P8 / F - asset-creator: CreatorTools — the chat tool pack for a store that only stores (aicore §9 / P8) - asset-store: the shrink — the store stores (aicore P7) - importer + asset-server host: the coordination era ends (aicore P7) - store config purge + asset-ui goes fleet-direct; the derive protocol gets its route proof (aicore P7) - client + chat dispatcher: the dead wire comes out (aicore P7/P8) - ai-hub: 0.3.0 — the health version says which era a node runs - ai-hub: the default fleet is 'gen' — apps hear the LAN without env plumbing - ai-hub: the preload note percents the prefill, not the job bar - ai-hub: conversations keep their KV — the wire mirror, the lane identity, the in-turn dynamic context (aicore §7) - ai-hub: an open-think model is thinking from its first token - libs: the zero-warning sweep — stitch casts say what they mean, xatlas keeps upstream's surface quietly - zero-warning sweep, round two — the first full-workspace pass - zero-warning sweep, round three — the model lanes and the deep examples - zero-warning sweep, round four — the last stragglers - zero-warning sweep, round five — vj and chat-ui - zero-warning sweep, round six — three cascades Co-authored-by: Claude <info@makepad.nl>
475 lines
17 KiB
Rust
475 lines
17 KiB
Rust
//! Downloader tests against a local fixture HTTP server (plain TCP): fresh
|
|
//! download, Range resume, redirect following, servers that ignore Range,
|
|
//! and sha256 verification.
|
|
|
|
use makepad_ai_hub::backend::CancelToken;
|
|
use makepad_ai_hub::download::{part_path, Downloader};
|
|
use makepad_ai_hub::error::AssetAiError;
|
|
use makepad_ai_hub::registry::FileSpec;
|
|
use makepad_ai_hub::sha256::sha256_hex;
|
|
use std::io::{Read, Write};
|
|
use std::net::{SocketAddr, TcpListener, TcpStream};
|
|
use std::path::PathBuf;
|
|
use std::sync::{Arc, Mutex};
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Fixture server: serves `data` at /repo/resolve/main/file.bin with optional
|
|
// Range support and an optional redirect hop; records request headers.
|
|
// ---------------------------------------------------------------------------
|
|
|
|
#[derive(Clone, Copy, PartialEq)]
|
|
enum RangeMode {
|
|
/// Honors Range with 206 + Content-Range.
|
|
Honor,
|
|
/// Ignores Range and always answers 200 with the full body.
|
|
Ignore,
|
|
}
|
|
|
|
struct Fixture {
|
|
addr: SocketAddr,
|
|
/// One entry per request: the raw request head.
|
|
requests: Arc<Mutex<Vec<String>>>,
|
|
}
|
|
|
|
fn spawn_fixture(data: Vec<u8>, range_mode: RangeMode, redirect_first: bool) -> Fixture {
|
|
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
|
|
let addr = listener.local_addr().unwrap();
|
|
let requests = Arc::new(Mutex::new(Vec::new()));
|
|
let requests_thread = requests.clone();
|
|
std::thread::spawn(move || {
|
|
let mut redirect_pending = redirect_first;
|
|
for stream in listener.incoming() {
|
|
let Ok(mut stream) = stream else { break };
|
|
let head = read_head(&mut stream);
|
|
requests_thread.lock().unwrap().push(head.clone());
|
|
let path = head
|
|
.lines()
|
|
.next()
|
|
.and_then(|line| line.split_whitespace().nth(1))
|
|
.unwrap_or("/")
|
|
.to_string();
|
|
if redirect_pending && !path.starts_with("/cdn") {
|
|
redirect_pending = false;
|
|
let response =
|
|
"HTTP/1.1 302 Found\r\nLocation: /cdn/file.bin\r\nContent-Length: 0\r\nConnection: close\r\n\r\n";
|
|
let _ = stream.write_all(response.as_bytes());
|
|
continue;
|
|
}
|
|
let range_from = parse_range(&head);
|
|
match (range_mode, range_from) {
|
|
(RangeMode::Honor, Some(from)) if from >= data.len() as u64 => {
|
|
let response = format!(
|
|
"HTTP/1.1 416 Range Not Satisfiable\r\nContent-Range: bytes */{}\r\nContent-Length: 0\r\nConnection: close\r\n\r\n",
|
|
data.len()
|
|
);
|
|
let _ = stream.write_all(response.as_bytes());
|
|
}
|
|
(RangeMode::Honor, Some(from)) => {
|
|
let body = &data[from as usize..];
|
|
let response = format!(
|
|
"HTTP/1.1 206 Partial Content\r\nContent-Range: bytes {}-{}/{}\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
|
|
from,
|
|
data.len() - 1,
|
|
data.len(),
|
|
body.len()
|
|
);
|
|
let _ = stream.write_all(response.as_bytes());
|
|
let _ = stream.write_all(body);
|
|
}
|
|
_ => {
|
|
let response = format!(
|
|
"HTTP/1.1 200 OK\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
|
|
data.len()
|
|
);
|
|
let _ = stream.write_all(response.as_bytes());
|
|
let _ = stream.write_all(&data);
|
|
}
|
|
}
|
|
}
|
|
});
|
|
Fixture { addr, requests }
|
|
}
|
|
|
|
fn read_head(stream: &mut TcpStream) -> String {
|
|
let mut buf = Vec::new();
|
|
let mut byte = [0u8; 1];
|
|
while !buf.ends_with(b"\r\n\r\n") {
|
|
match stream.read(&mut byte) {
|
|
Ok(1) => buf.push(byte[0]),
|
|
_ => break,
|
|
}
|
|
}
|
|
String::from_utf8_lossy(&buf).to_string()
|
|
}
|
|
|
|
fn parse_range(head: &str) -> Option<u64> {
|
|
for line in head.lines() {
|
|
let lower = line.to_ascii_lowercase();
|
|
if let Some(rest) = lower.strip_prefix("range: bytes=") {
|
|
return rest.trim_end_matches('-').trim().parse().ok();
|
|
}
|
|
}
|
|
None
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Helpers
|
|
// ---------------------------------------------------------------------------
|
|
|
|
fn test_dir(name: &str) -> PathBuf {
|
|
let dir = std::env::temp_dir().join(format!(
|
|
"makepad-asset-ai-test-{name}-{}",
|
|
std::process::id()
|
|
));
|
|
let _ = std::fs::remove_dir_all(&dir);
|
|
std::fs::create_dir_all(&dir).unwrap();
|
|
dir
|
|
}
|
|
|
|
fn file_spec(sha256: Option<String>) -> FileSpec {
|
|
FileSpec {
|
|
role: None,
|
|
repo: "test-org/test-repo".to_string(),
|
|
path: "file.bin".to_string(),
|
|
revision: None,
|
|
cache_as: "unet/file.bin".to_string(),
|
|
size: None,
|
|
sha256,
|
|
local: false,
|
|
optional: false,
|
|
converts_to: None,
|
|
conversion: None,
|
|
}
|
|
}
|
|
|
|
fn downloader_for(fixture: &Fixture) -> Downloader {
|
|
Downloader::new(&format!("http://{}", fixture.addr), None).unwrap()
|
|
}
|
|
|
|
fn test_data(len: usize) -> Vec<u8> {
|
|
(0..len).map(|i| (i % 251) as u8).collect()
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Tests
|
|
// ---------------------------------------------------------------------------
|
|
|
|
#[test]
|
|
fn fresh_download() {
|
|
let data = test_data(10_000);
|
|
let fixture = spawn_fixture(data.clone(), RangeMode::Honor, false);
|
|
let cache = test_dir("fresh");
|
|
let spec = file_spec(None);
|
|
|
|
let mut progress_reports = Vec::new();
|
|
let dest = downloader_for(&fixture)
|
|
.ensure_file(&spec, &cache, &mut |p| {
|
|
progress_reports.push((p.done, p.total))
|
|
}, &CancelToken::new())
|
|
.unwrap();
|
|
|
|
assert_eq!(std::fs::read(&dest).unwrap(), data);
|
|
assert!(!part_path(&dest).exists(), ".part must be renamed away");
|
|
// First request had no Range header.
|
|
let requests = fixture.requests.lock().unwrap();
|
|
assert!(!requests[0].to_ascii_lowercase().contains("range:"));
|
|
// Progress reached the total.
|
|
let last = progress_reports.last().unwrap();
|
|
assert_eq!(last.0, data.len() as u64);
|
|
assert_eq!(last.1, Some(data.len() as u64));
|
|
// The downloader identifies itself.
|
|
assert!(requests[0].contains("makepad-asset-ai"));
|
|
}
|
|
|
|
#[test]
|
|
fn resume_uses_range_and_appends() {
|
|
let data = test_data(50_000);
|
|
let fixture = spawn_fixture(data.clone(), RangeMode::Honor, false);
|
|
let cache = test_dir("resume");
|
|
let spec = file_spec(None);
|
|
|
|
// Pre-seed a partial download: first 12_345 bytes.
|
|
let dest = spec.dest_path(&cache);
|
|
std::fs::create_dir_all(dest.parent().unwrap()).unwrap();
|
|
std::fs::write(part_path(&dest), &data[..12_345]).unwrap();
|
|
|
|
let mut first_done = None;
|
|
let out = downloader_for(&fixture)
|
|
.ensure_file(&spec, &cache, &mut |p| {
|
|
if first_done.is_none() {
|
|
first_done = Some(p.done);
|
|
}
|
|
}, &CancelToken::new())
|
|
.unwrap();
|
|
|
|
assert_eq!(std::fs::read(&out).unwrap(), data);
|
|
// The request asked to resume exactly where the .part ended.
|
|
let requests = fixture.requests.lock().unwrap();
|
|
assert!(
|
|
requests[0].to_ascii_lowercase().contains("range: bytes=12345-"),
|
|
"expected Range header in: {}",
|
|
requests[0]
|
|
);
|
|
// Progress started from the resumed offset, not zero.
|
|
assert_eq!(first_done, Some(12_345));
|
|
}
|
|
|
|
#[test]
|
|
fn resume_restarts_when_server_ignores_range() {
|
|
let data = test_data(20_000);
|
|
let fixture = spawn_fixture(data.clone(), RangeMode::Ignore, false);
|
|
let cache = test_dir("norange");
|
|
let spec = file_spec(None);
|
|
|
|
// Pre-seed garbage that would corrupt the file if it were kept.
|
|
let dest = spec.dest_path(&cache);
|
|
std::fs::create_dir_all(dest.parent().unwrap()).unwrap();
|
|
std::fs::write(part_path(&dest), vec![0xffu8; 5000]).unwrap();
|
|
|
|
let out = downloader_for(&fixture)
|
|
.ensure_file(&spec, &cache, &mut |_| {}, &CancelToken::new())
|
|
.unwrap();
|
|
// Server answered 200: the stale prefix must have been discarded.
|
|
assert_eq!(std::fs::read(&out).unwrap(), data);
|
|
}
|
|
|
|
#[test]
|
|
fn follows_redirect() {
|
|
let data = test_data(4_000);
|
|
let fixture = spawn_fixture(data.clone(), RangeMode::Honor, true);
|
|
let cache = test_dir("redirect");
|
|
let spec = file_spec(None);
|
|
|
|
let out = downloader_for(&fixture)
|
|
.ensure_file(&spec, &cache, &mut |_| {}, &CancelToken::new())
|
|
.unwrap();
|
|
assert_eq!(std::fs::read(&out).unwrap(), data);
|
|
let requests = fixture.requests.lock().unwrap();
|
|
assert_eq!(requests.len(), 2, "redirect + follow");
|
|
assert!(requests[1].starts_with("GET /cdn/file.bin"));
|
|
}
|
|
|
|
#[test]
|
|
fn sha256_verify_pass_and_fail() {
|
|
let data = test_data(8_000);
|
|
let good = sha256_hex(&data);
|
|
|
|
// Passing case.
|
|
let fixture = spawn_fixture(data.clone(), RangeMode::Honor, false);
|
|
let cache = test_dir("shapass");
|
|
let out = downloader_for(&fixture)
|
|
.ensure_file(&file_spec(Some(good)), &cache, &mut |_| {}, &CancelToken::new())
|
|
.unwrap();
|
|
assert_eq!(std::fs::read(&out).unwrap(), data);
|
|
|
|
// Failing case: wrong hash -> error, partial data discarded, no dest.
|
|
let fixture = spawn_fixture(data, RangeMode::Honor, false);
|
|
let cache = test_dir("shafail");
|
|
let bad_spec = file_spec(Some("00".repeat(32)));
|
|
let err = downloader_for(&fixture)
|
|
.ensure_file(&bad_spec, &cache, &mut |_| {}, &CancelToken::new())
|
|
.unwrap_err();
|
|
match err {
|
|
AssetAiError::Download(message) => assert!(message.contains("sha256 mismatch")),
|
|
other => panic!("expected Download error, got {other:?}"),
|
|
}
|
|
let dest = bad_spec.dest_path(&cache);
|
|
assert!(!dest.exists());
|
|
assert!(!part_path(&dest).exists());
|
|
}
|
|
|
|
#[test]
|
|
fn existing_file_is_not_refetched() {
|
|
let data = test_data(1_000);
|
|
let fixture = spawn_fixture(data.clone(), RangeMode::Honor, false);
|
|
let cache = test_dir("cached");
|
|
let spec = file_spec(None);
|
|
|
|
let dest = spec.dest_path(&cache);
|
|
std::fs::create_dir_all(dest.parent().unwrap()).unwrap();
|
|
std::fs::write(&dest, &data).unwrap();
|
|
|
|
let out = downloader_for(&fixture)
|
|
.ensure_file(&spec, &cache, &mut |_| {}, &CancelToken::new())
|
|
.unwrap();
|
|
assert_eq!(out, dest);
|
|
assert!(
|
|
fixture.requests.lock().unwrap().is_empty(),
|
|
"no network traffic for a cached file"
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn fully_downloaded_part_finishes_via_416() {
|
|
let data = test_data(6_000);
|
|
let fixture = spawn_fixture(data.clone(), RangeMode::Honor, false);
|
|
let cache = test_dir("part416");
|
|
let spec = file_spec(None);
|
|
|
|
// The whole file is already in .part (e.g. the process died between the
|
|
// last write and the rename).
|
|
let dest = spec.dest_path(&cache);
|
|
std::fs::create_dir_all(dest.parent().unwrap()).unwrap();
|
|
std::fs::write(part_path(&dest), &data).unwrap();
|
|
|
|
let mut spec = spec;
|
|
spec.size = Some(data.len() as u64);
|
|
let out = downloader_for(&fixture)
|
|
.ensure_file(&spec, &cache, &mut |_| {}, &CancelToken::new())
|
|
.unwrap();
|
|
assert_eq!(std::fs::read(&out).unwrap(), data);
|
|
assert!(!part_path(&dest).exists());
|
|
}
|
|
|
|
#[test]
|
|
fn local_files_never_download() {
|
|
// No fixture server involved: local files must not touch the network.
|
|
let cache = test_dir("localfile");
|
|
let mut spec = file_spec(None);
|
|
spec.local = true;
|
|
let downloader = Downloader::new("http://127.0.0.1:9", None).unwrap();
|
|
|
|
// Missing -> helpful error naming the expected cache path.
|
|
let err = downloader
|
|
.ensure_file(&spec, &cache, &mut |_| {}, &CancelToken::new())
|
|
.unwrap_err();
|
|
let message = err.to_string();
|
|
assert!(message.contains("locally-converted"), "{message}");
|
|
assert!(message.contains("unet"), "{message}");
|
|
|
|
// Present -> returned as-is.
|
|
let dest = spec.dest_path(&cache);
|
|
std::fs::create_dir_all(dest.parent().unwrap()).unwrap();
|
|
std::fs::write(&dest, b"weights").unwrap();
|
|
let out = downloader
|
|
.ensure_file(&spec, &cache, &mut |_| {}, &CancelToken::new())
|
|
.unwrap();
|
|
assert_eq!(out, dest);
|
|
}
|
|
|
|
#[test]
|
|
fn pinned_revision_is_used_in_resolve_url() {
|
|
let downloader = Downloader::new("https://huggingface.co", None).unwrap();
|
|
let mut spec = file_spec(Some("11".repeat(32)));
|
|
spec.revision = Some("a".repeat(40));
|
|
spec.size = Some(1);
|
|
assert_eq!(
|
|
downloader.file_url(&spec),
|
|
format!(
|
|
"https://huggingface.co/test-org/test-repo/resolve/{}/file.bin",
|
|
"a".repeat(40)
|
|
)
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn corrupt_existing_file_is_verified_and_replaced() {
|
|
let data = test_data(16_000);
|
|
let fixture = spawn_fixture(data.clone(), RangeMode::Honor, false);
|
|
let cache = test_dir("corrupt-existing");
|
|
let mut spec = file_spec(Some(sha256_hex(&data)));
|
|
spec.size = Some(data.len() as u64);
|
|
let dest = spec.dest_path(&cache);
|
|
std::fs::create_dir_all(dest.parent().unwrap()).unwrap();
|
|
std::fs::write(&dest, vec![0xff; data.len()]).unwrap();
|
|
|
|
let out = downloader_for(&fixture)
|
|
.ensure_file(&spec, &cache, &mut |_| {}, &CancelToken::new())
|
|
.unwrap();
|
|
assert_eq!(std::fs::read(out).unwrap(), data);
|
|
assert_eq!(fixture.requests.lock().unwrap().len(), 1);
|
|
|
|
// The identity receipt makes the second verification network- and
|
|
// multi-megabyte-hash-free while still binding current len+mtime.
|
|
downloader_for(&fixture)
|
|
.ensure_file(&spec, &cache, &mut |_| {}, &CancelToken::new())
|
|
.unwrap();
|
|
assert_eq!(fixture.requests.lock().unwrap().len(), 1);
|
|
}
|
|
|
|
#[test]
|
|
fn wrong_size_partial_416_is_rejected() {
|
|
let data = test_data(6_000);
|
|
let fixture = spawn_fixture(data.clone(), RangeMode::Honor, false);
|
|
let cache = test_dir("wrong-416");
|
|
let mut spec = file_spec(None);
|
|
// Server truth is 6000, manifest claims 7000. A complete server-sized
|
|
// partial provokes 416 but must not be promoted to final.
|
|
spec.size = Some(7_000);
|
|
let dest = spec.dest_path(&cache);
|
|
std::fs::create_dir_all(dest.parent().unwrap()).unwrap();
|
|
std::fs::write(part_path(&dest), &data).unwrap();
|
|
let err = downloader_for(&fixture)
|
|
.ensure_file(&spec, &cache, &mut |_| {}, &CancelToken::new())
|
|
.unwrap_err();
|
|
assert!(err.to_string().contains("expected exact manifest size"));
|
|
assert!(!dest.exists());
|
|
assert!(part_path(&dest).exists(), "resume bytes must be retained");
|
|
}
|
|
|
|
#[test]
|
|
fn cancelled_download_keeps_resumable_part() {
|
|
let data = test_data(200_000);
|
|
let fixture = spawn_fixture(data, RangeMode::Honor, false);
|
|
let cache = test_dir("cancel-resume");
|
|
let spec = file_spec(None);
|
|
let cancel = CancelToken::new();
|
|
let cancel_progress = cancel.clone();
|
|
let err = downloader_for(&fixture)
|
|
.ensure_file(
|
|
&spec,
|
|
&cache,
|
|
&mut |progress| {
|
|
if progress.done >= 65_536 {
|
|
cancel_progress.cancel();
|
|
}
|
|
},
|
|
&cancel,
|
|
)
|
|
.unwrap_err();
|
|
assert_eq!(err, AssetAiError::Cancelled);
|
|
let dest = spec.dest_path(&cache);
|
|
assert!(!dest.exists());
|
|
let partial = part_path(&dest);
|
|
assert!(partial.exists());
|
|
let partial_len = std::fs::metadata(partial).unwrap().len();
|
|
assert!(partial_len >= 65_536 && partial_len < 200_000, "{partial_len}");
|
|
}
|
|
|
|
#[test]
|
|
fn concurrent_downloaders_share_one_artifact_transaction() {
|
|
let data = test_data(512_000);
|
|
let fixture = spawn_fixture(data.clone(), RangeMode::Honor, false);
|
|
let cache = test_dir("concurrent-lock");
|
|
let mut spec = file_spec(Some(sha256_hex(&data)));
|
|
spec.size = Some(data.len() as u64);
|
|
let downloader = downloader_for(&fixture);
|
|
let barrier = Arc::new(std::sync::Barrier::new(3));
|
|
let mut threads = Vec::new();
|
|
for _ in 0..2 {
|
|
let spec = spec.clone();
|
|
let cache = cache.clone();
|
|
let downloader = downloader.clone();
|
|
let barrier = barrier.clone();
|
|
threads.push(std::thread::spawn(move || {
|
|
barrier.wait();
|
|
downloader
|
|
.ensure_file(&spec, &cache, &mut |_| {}, &CancelToken::new())
|
|
.unwrap()
|
|
}));
|
|
}
|
|
barrier.wait();
|
|
for thread in threads {
|
|
assert_eq!(std::fs::read(thread.join().unwrap()).unwrap(), data);
|
|
}
|
|
assert_eq!(
|
|
fixture.requests.lock().unwrap().len(),
|
|
1,
|
|
"only the lock winner may hit the network"
|
|
);
|
|
let dest = spec.dest_path(&cache);
|
|
assert!(!part_path(&dest).exists());
|
|
let mut lock = dest.as_os_str().to_os_string();
|
|
lock.push(".lock");
|
|
assert!(!PathBuf::from(lock).exists());
|
|
}
|