1use std::sync::atomic::{AtomicU64, Ordering};
9use std::sync::{Arc, Mutex};
10use std::time::{Duration, Instant};
11
12use reqwest::Client;
13use tracing::{debug, warn};
14
15use crate::error::TransferError;
16use crate::file::DestFile;
17use crate::limit::RateLimiter;
18use crate::mirror::SourceSet;
19use crate::progress::{Event, Reporter};
20use crate::retry::{Decision, RetryPolicy};
21use crate::scheduler::{Lease, Scheduler};
22use crate::shutdown::Cancel;
23use crate::storage::OPEN_END;
24
25pub struct WorkerCtx {
27 pub client: Client,
28 pub sources: Arc<SourceSet>,
29 pub file: Arc<DestFile>,
30 pub scheduler: Arc<Scheduler>,
31 pub reporter: Reporter,
32 pub retry: RetryPolicy,
33 pub limiter: Option<Arc<RateLimiter>>,
34 pub read_timeout: Duration,
36 pub expected_total: Option<u64>,
37 pub ranges_supported: bool,
41 pub discovered_size: Arc<AtomicU64>,
44 pub cancel: Cancel,
45 pub primed: Mutex<Option<PrimedBody>>,
49}
50
51pub struct PrimedBody {
54 pub url: url::Url,
55 pub response: reqwest::Response,
56}
57
58impl WorkerCtx {
59 fn take_primed(&self, url: &url::Url, start: u64) -> Option<reqwest::Response> {
62 if start != 0 {
63 return None;
64 }
65 let mut slot = self.primed.lock().unwrap_or_else(|e| e.into_inner());
66 match slot.as_ref() {
67 Some(primed) if &primed.url == url => slot.take().map(|p| p.response),
68 _ => None,
69 }
70 }
71}
72
73pub enum WorkerOutcome {
75 Finished,
77 Fatal(TransferError),
79}
80
81pub async fn run(ctx: Arc<WorkerCtx>, worker_id: usize) -> WorkerOutcome {
84 loop {
85 if ctx.cancel.is_cancelled() {
86 return WorkerOutcome::Finished;
87 }
88
89 let Some((lease, split)) = ctx.scheduler.acquire() else {
90 if ctx.scheduler.is_finished() || ctx.scheduler.has_failure() {
91 return WorkerOutcome::Finished;
92 }
93 tokio::select! {
96 _ = ctx.scheduler.wait_for_change(Duration::from_millis(250)) => {}
97 _ = ctx.cancel.cancelled() => return WorkerOutcome::Finished,
98 }
99 continue;
100 };
101
102 if let Some(split) = split {
103 ctx.reporter.emit(Event::RangeSplit {
104 index: split.shrunk.idx,
105 new_index: split.added.idx,
106 at: split.added.start,
107 });
108 debug!(
109 worker = worker_id,
110 victim = split.shrunk.idx,
111 new_range = split.added.idx,
112 at = split.added.start,
113 "split a slow range"
114 );
115 }
116
117 ctx.reporter.emit(Event::RangeStarted {
118 index: lease.idx,
119 start: lease.cursor(),
120 end: lease.end(),
121 });
122
123 match transfer(&ctx, &lease, worker_id).await {
124 Ok(()) => {
125 ctx.scheduler.complete(lease.idx);
126 ctx.reporter
127 .emit(Event::RangeCompleted { index: lease.idx });
128 let (done, total) = ctx.scheduler.counts();
129 ctx.reporter.stats.set_ranges_complete(done);
130 ctx.reporter.stats.set_ranges_total(total);
131 }
132 Err(TransferError::Cancelled) => {
133 ctx.scheduler.release(lease.idx);
134 return WorkerOutcome::Finished;
135 }
136 Err(err) => {
137 ctx.scheduler.fail(lease.idx);
139 warn!(worker = worker_id, range = lease.idx, %err, "range failed permanently");
140 return WorkerOutcome::Fatal(err);
141 }
142 }
143 }
144}
145
146async fn transfer(ctx: &WorkerCtx, lease: &Lease, worker_id: usize) -> Result<(), TransferError> {
149 let mut attempts = 0u32;
150 loop {
151 if ctx.cancel.is_cancelled() {
152 return Err(TransferError::Cancelled);
153 }
154 if lease.remaining() == 0 {
155 return Ok(());
156 }
157
158 let (source_idx, source) = ctx.sources.pick();
159 let url = source.url.clone();
160 let validator = source.validator.clone();
162 let before = lease.progress();
163 match attempt(ctx, lease, &url, validator.as_deref()).await {
164 Ok(()) => {
165 ctx.sources.reward(source_idx);
166 return Ok(());
167 }
168 Err(err) => {
169 if lease.progress() > before {
173 attempts = 0;
174 }
175 attempts += 1;
176 ctx.sources.penalise(source_idx);
177 match ctx.retry.decide(&err, attempts) {
178 Decision::Retry { delay, attempt } => {
179 ctx.reporter.stats.record_retry();
180 ctx.reporter.emit(Event::RetryScheduled {
181 index: Some(lease.idx),
182 attempt,
183 delay_ms: delay.as_millis() as u64,
184 reason: err.to_string(),
185 });
186 debug!(
187 worker = worker_id,
188 range = lease.idx,
189 attempt,
190 ?delay,
191 %err,
192 "retrying range"
193 );
194 tokio::select! {
195 _ = tokio::time::sleep(delay) => {}
196 _ = ctx.cancel.cancelled() => return Err(TransferError::Cancelled),
197 }
198 }
199 Decision::GiveUp => return Err(err),
200 }
201 }
202 }
203 }
204}
205
206async fn attempt(
208 ctx: &WorkerCtx,
209 lease: &Lease,
210 url: &url::Url,
211 validator: Option<&str>,
212) -> Result<(), TransferError> {
213 let start_cursor = lease.cursor();
214 let open_ended = lease.is_open_ended();
215
216 let (req_start, req_end) = if ctx.ranges_supported {
217 (
218 start_cursor,
219 if open_ended { None } else { Some(lease.end()) },
220 )
221 } else {
222 if start_cursor != 0 {
223 return Err(TransferError::Protocol(
224 "cannot resume from the middle: this server does not support range requests".into(),
225 ));
226 }
227 (0, None)
228 };
229
230 let mut resp = match ctx.take_primed(url, req_start) {
233 Some(primed) => primed,
234 None => {
235 crate::http::get_range(
236 &ctx.client,
237 url,
238 req_start,
239 req_end,
240 validator,
241 ctx.expected_total,
242 )
243 .await?
244 }
245 };
246
247 ctx.reporter.stats.connection_opened();
248 let _conn = ConnectionGuard(&ctx.reporter);
250
251 let mut cursor = start_cursor;
252 let mut pending_event_bytes = 0u64;
253 let mut last_event = Instant::now();
254
255 loop {
256 if ctx.cancel.is_cancelled() {
257 return Err(TransferError::Cancelled);
258 }
259
260 let chunk = tokio::select! {
263 biased;
264 _ = ctx.cancel.cancelled() => return Err(TransferError::Cancelled),
265 read = tokio::time::timeout(ctx.read_timeout, resp.chunk()) => match read {
266 Err(_) => return Err(TransferError::Timeout(ctx.read_timeout)),
267 Ok(Err(e)) => return Err(TransferError::from_reqwest(&e)),
268 Ok(Ok(None)) => break,
269 Ok(Ok(Some(chunk))) => chunk,
270 },
271 };
272 if chunk.is_empty() {
273 continue;
274 }
275
276 let end = lease.end();
279 let writable = if open_ended && end >= OPEN_END {
280 chunk.len() as u64
281 } else {
282 let room = end.saturating_sub(cursor).saturating_add(1);
283 (chunk.len() as u64).min(room)
284 };
285 if writable == 0 {
286 break;
288 }
289
290 if let Some(limiter) = &ctx.limiter {
291 limiter.acquire(writable).await;
292 }
293
294 let slice = &chunk[..writable as usize];
295 ctx.file
296 .write_at(slice, cursor)
297 .map_err(|e| TransferError::Io(e.to_string()))?;
298
299 cursor += writable;
300 lease.publish_progress(cursor - lease.start);
303 ctx.reporter.stats.add_downloaded(writable);
304
305 pending_event_bytes += writable;
306 if last_event.elapsed() >= Duration::from_millis(100) {
307 ctx.reporter.emit(Event::BytesWritten {
308 index: lease.idx,
309 bytes: pending_event_bytes,
310 });
311 pending_event_bytes = 0;
312 last_event = Instant::now();
313 }
314
315 if !open_ended && cursor > end {
316 break;
317 }
318 }
319
320 if pending_event_bytes > 0 {
321 ctx.reporter.emit(Event::BytesWritten {
322 index: lease.idx,
323 bytes: pending_event_bytes,
324 });
325 }
326
327 if open_ended && lease.end() >= OPEN_END {
328 ctx.discovered_size.store(cursor, Ordering::Release);
330 ctx.scheduler
331 .set_end(lease.idx, cursor.saturating_sub(1).max(lease.start));
332 return Ok(());
333 }
334
335 if cursor > lease.end() {
336 return Ok(());
337 }
338
339 Err(TransferError::Network(format!(
342 "connection closed with {} bytes of the range still missing",
343 lease.end() + 1 - cursor
344 )))
345}
346
347struct ConnectionGuard<'a>(&'a Reporter);
349
350impl Drop for ConnectionGuard<'_> {
351 fn drop(&mut self) {
352 self.0.stats.connection_closed();
353 }
354}