use crate::transfer_status::TransferStatus;
use crate::inner::scheduler_state::SchedulerState;
use crate::inner::worker_event::WorkerEvent;
pub(crate) async fn handle_worker_event(event: WorkerEvent, state: &mut SchedulerState) {
match event {
WorkerEvent::Progress {
key,
next_offset,
total_size,
} => {
crate::meow_trace_log!(
"worker_event",
"progress: key={} next_offset={} total_size={}",
crate::inner::safe_key(&key),
next_offset,
total_size
);
if state.groups().contains_key(&key) {
state.offsets_mut().insert(key.clone(), next_offset);
if total_size > 0 {
state.known_totals_mut().insert(key.clone(), total_size);
}
if let Some(group) = state.groups().get(&key) {
let total = total_size.max(crate::inner::exec_impl::emit::effective_total(
state,
&key,
group.entry().inner(),
));
crate::inner::exec_impl::emit::emit_status(
state,
group.entry(),
TransferStatus::Transmission,
next_offset,
total,
);
}
} else {
crate::log::emit_lazy(|| {
crate::log::Log::warn(
"worker_event",
format!(
"stray Progress for dead group; ignored next_offset={} total_size={}",
next_offset, total_size
),
)
.with_key(key.1.as_str())
.with_offset(next_offset)
});
}
}
WorkerEvent::Completed {
key,
total_size,
completion_payload,
} => {
crate::meow_key_log!(
"worker_event",
"completed: key={} total_size={}",
crate::inner::safe_key(&key),
total_size
);
state.active_mut().remove(&key);
state.paused_set_mut().remove(&key);
state.offsets_mut().insert(key.clone(), total_size);
if let Some(group) = state.groups_mut().remove(&key) {
state
.task_id_to_dedupe_mut()
.remove(&group.leader_inner().task_id());
let task_id = group.entry().inner().task_id();
let total = total_size.max(crate::inner::exec_impl::emit::effective_total(
state,
&key,
group.entry().inner(),
));
crate::inner::exec_impl::emit::emit_status(
state,
group.entry(),
TransferStatus::Complete,
total,
total,
);
if let Some(cb) = group.entry().callbacks().complete_cb() {
crate::inner::exec_impl::emit::invoke_complete_cb(
state,
cb,
task_id,
completion_payload,
);
}
} else {
crate::log::emit_lazy(|| {
crate::log::Log::warn(
"worker_event",
format!("Completed for unknown group; nothing to finalize total_size={}", total_size),
)
.with_key(key.1.as_str())
});
}
state.offsets_mut().remove(&key);
state.known_totals_mut().remove(&key);
}
WorkerEvent::Failed {
key,
error,
total_size,
} => {
state.active_mut().remove(&key);
state.paused_set_mut().remove(&key);
if let Some(group) = state.groups_mut().remove(&key) {
state
.task_id_to_dedupe_mut()
.remove(&group.leader_inner().task_id());
let current = state.offsets().get(&key).copied().unwrap_or(0);
crate::log::emit_lazy(|| {
let mut log = crate::log::Log::error(
"worker_event",
format!("task failed: {}", crate::log::redact_secrets(&error.to_string())),
)
.with_key(key.1.as_str())
.with_offset(current)
.with_error_code(error.code());
if let Some(status) = error.http_status() {
log = log.with_http_status(status);
}
log
});
let total = total_size.max(crate::inner::exec_impl::emit::effective_total(
state,
&key,
group.entry().inner(),
));
crate::inner::exec_impl::emit::emit_status(
state,
group.entry(),
TransferStatus::Failed(error),
current,
total,
);
} else {
crate::log::emit_lazy(|| {
crate::log::Log::warn(
"worker_event",
format!("Failed for unknown group; nothing to finalize error={}", crate::log::redact_secrets(&error.to_string())),
)
.with_key(key.1.as_str())
.with_error_code(error.code())
});
}
state.offsets_mut().remove(&key);
state.known_totals_mut().remove(&key);
}
WorkerEvent::Canceled { key, total_size } => {
crate::meow_key_log!("worker_event", "canceled: key={}", crate::inner::safe_key(&key));
state.active_mut().remove(&key);
if state.paused_set().contains(&key) {
crate::meow_key_log!(
"worker_event",
"canceled from pause flow; keep group for resume: key={}",
crate::inner::safe_key(&key)
);
return;
}
if let Some(group) = state.groups_mut().remove(&key) {
state
.task_id_to_dedupe_mut()
.remove(&group.leader_inner().task_id());
let current = state.offsets().get(&key).copied().unwrap_or(0);
let total = total_size.max(crate::inner::exec_impl::emit::effective_total(
state,
&key,
group.entry().inner(),
));
crate::inner::exec_impl::emit::emit_status(
state,
group.entry(),
TransferStatus::Canceled,
current,
total,
);
} else {
crate::log::emit_lazy(|| {
crate::log::Log::warn(
"worker_event",
"Canceled for unknown group; nothing to finalize".to_string(),
)
.with_key(key.1.as_str())
});
}
state.offsets_mut().remove(&key);
state.known_totals_mut().remove(&key);
}
}
}
#[cfg(test)]
mod tests {
use std::sync::{Arc, RwLock};
use super::handle_worker_event;
use crate::dflt::default_http_transfer::default_breakpoint_arcs;
use crate::error::{InnerErrorCode, MeowError};
use crate::http_breakpoint::BreakpointDownloadHttpConfig;
use crate::inner::cb_dispatcher;
use crate::inner::group_state::{GroupState, RecordEntry};
use crate::inner::inner_task::InnerTask;
use crate::inner::scheduler_state::SchedulerState;
use crate::inner::task_callbacks::TaskCallbacks;
use crate::inner::test_support::{attach_capture, live_download_state, wait_for_record};
use crate::inner::worker_event::WorkerEvent;
use crate::inner::UniqueId;
use crate::transfer_status::TransferStatus;
use crate::up_pounce_builder::UploadPounceBuilder;
async fn live_upload_state() -> (SchedulerState, UniqueId, u64) {
let mut path = std::env::temp_dir();
let ts = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.expect("clock")
.as_nanos();
path.push(format!("rusty_cat_offsets_teardown_{ts}.bin"));
std::fs::write(&path, vec![7u8; 4096]).expect("write fixture");
let pounce = UploadPounceBuilder::new("teardown.bin", &path, 1024)
.with_url("https://placeholder/upload")
.build()
.expect("build pounce");
let (def_up, def_down) = default_breakpoint_arcs();
let inner = InnerTask::from_pounce(
pounce,
BreakpointDownloadHttpConfig::default(),
None,
def_up,
def_down,
)
.await
.expect("from_pounce");
let _ = std::fs::remove_file(&path);
let key = inner.dedupe_key();
let total = inner.total_size();
let (cb_submit, cb_join) = cb_dispatcher::start().expect("start dispatcher");
std::mem::forget(cb_join);
let mut state = SchedulerState::new(1, 1, Arc::new(RwLock::new(Vec::new())), cb_submit);
state
.task_id_to_dedupe_mut()
.insert(inner.task_id(), key.clone());
state.offsets_mut().insert(key.clone(), 0);
let entry = RecordEntry::new(inner.clone(), TaskCallbacks::empty());
state
.groups_mut()
.insert(key.clone(), GroupState::new(inner.clone(), entry));
(state, key, total)
}
#[tokio::test]
async fn stray_progress_after_complete_creates_no_orphan_offset() {
let (mut state, key, total) = live_upload_state().await;
assert!(state.groups().contains_key(&key));
handle_worker_event(
WorkerEvent::Completed {
key: key.clone(),
total_size: total,
completion_payload: None,
},
&mut state,
)
.await;
assert!(
!state.groups().contains_key(&key),
"group must be removed after Completed"
);
assert!(
state.offsets().get(&key).is_none(),
"offsets entry must be cleared after Completed"
);
handle_worker_event(
WorkerEvent::Progress {
key: key.clone(),
next_offset: 512,
total_size: total,
},
&mut state,
)
.await;
assert!(
state.offsets().get(&key).is_none(),
"stray Progress after teardown must not resurrect an orphan offsets entry"
);
}
#[tokio::test]
async fn progress_for_live_group_persists_offset() {
let (mut state, key, total) = live_upload_state().await;
handle_worker_event(
WorkerEvent::Progress {
key: key.clone(),
next_offset: 1024,
total_size: total,
},
&mut state,
)
.await;
assert_eq!(
state.offsets().get(&key).copied(),
Some(1024),
"progress for a live group must persist its contiguous offset"
);
}
#[tokio::test]
async fn repro_download_failed_after_progress_zeroes_total() {
let (mut state, key) = live_download_state("repro_failed").await;
let records = attach_capture(&state);
handle_worker_event(
WorkerEvent::Progress {
key: key.clone(),
next_offset: 512,
total_size: 4096,
},
&mut state,
)
.await;
let transmission = wait_for_record(&records, |r| {
matches!(r.status(), TransferStatus::Transmission)
})
.await;
assert_eq!(transmission.total_size(), 4096, "传输中记录应带真实 total");
assert!((transmission.progress() - 0.125).abs() < f32::EPSILON);
handle_worker_event(
WorkerEvent::Failed {
key: key.clone(),
error: MeowError::from_code(InnerErrorCode::Unknown, "boom".to_string()),
total_size: 0,
},
&mut state,
)
.await;
let failed = wait_for_record(&records, |r| {
matches!(r.status(), TransferStatus::Failed(_))
})
.await;
assert_eq!(
failed.total_size(),
4096,
"终态必须保留运行期 total(复现点)"
);
assert!(
(failed.progress() - 0.125).abs() < f32::EPSILON,
"终态必须保留 512/4096 进度(复现点)"
);
}
#[tokio::test]
async fn progress_records_runtime_total_into_scheduler() {
let (mut state, key) = live_download_state("record_total").await;
handle_worker_event(
WorkerEvent::Progress {
key: key.clone(),
next_offset: 512,
total_size: 4096,
},
&mut state,
)
.await;
assert_eq!(state.known_totals().get(&key).copied(), Some(4096));
}
#[tokio::test]
async fn zero_total_progress_does_not_clobber_known_total() {
let (mut state, key) = live_download_state("no_clobber").await;
handle_worker_event(
WorkerEvent::Progress {
key: key.clone(),
next_offset: 512,
total_size: 4096,
},
&mut state,
)
.await;
handle_worker_event(
WorkerEvent::Progress {
key: key.clone(),
next_offset: 600,
total_size: 0,
},
&mut state,
)
.await;
assert_eq!(state.known_totals().get(&key).copied(), Some(4096));
}
#[tokio::test]
async fn terminal_completed_cleans_known_totals() {
let (mut state, key) = live_download_state("complete_cleans").await;
handle_worker_event(
WorkerEvent::Progress {
key: key.clone(),
next_offset: 512,
total_size: 4096,
},
&mut state,
)
.await;
handle_worker_event(
WorkerEvent::Completed {
key: key.clone(),
total_size: 4096,
completion_payload: None,
},
&mut state,
)
.await;
assert!(state.known_totals().get(&key).is_none());
assert!(state.offsets().get(&key).is_none());
}
#[tokio::test]
async fn terminal_failed_cleans_known_totals() {
let (mut state, key) = live_download_state("failed_cleans").await;
handle_worker_event(
WorkerEvent::Progress {
key: key.clone(),
next_offset: 512,
total_size: 4096,
},
&mut state,
)
.await;
handle_worker_event(
WorkerEvent::Failed {
key: key.clone(),
error: MeowError::from_code(InnerErrorCode::Unknown, "boom".to_string()),
total_size: 0,
},
&mut state,
)
.await;
assert!(state.known_totals().get(&key).is_none());
assert!(state.offsets().get(&key).is_none());
}
#[tokio::test]
async fn terminal_canceled_cleans_known_totals() {
let (mut state, key) = live_download_state("canceled_cleans").await;
handle_worker_event(
WorkerEvent::Progress {
key: key.clone(),
next_offset: 512,
total_size: 4096,
},
&mut state,
)
.await;
handle_worker_event(
WorkerEvent::Canceled {
key: key.clone(),
total_size: 0,
},
&mut state,
)
.await;
assert!(state.known_totals().get(&key).is_none());
}
#[tokio::test]
async fn pause_flow_canceled_keeps_known_total() {
let (mut state, key) = live_download_state("pause_keeps").await;
handle_worker_event(
WorkerEvent::Progress {
key: key.clone(),
next_offset: 512,
total_size: 4096,
},
&mut state,
)
.await;
state.paused_set_mut().insert(key.clone());
handle_worker_event(
WorkerEvent::Canceled {
key: key.clone(),
total_size: 0,
},
&mut state,
)
.await;
assert!(state.groups().contains_key(&key), "pause 流不得销毁 group");
assert_eq!(state.known_totals().get(&key).copied(), Some(4096));
}
#[tokio::test]
async fn stray_progress_after_teardown_creates_no_orphan_known_total() {
let (mut state, key) = live_download_state("stray_progress").await;
handle_worker_event(
WorkerEvent::Completed {
key: key.clone(),
total_size: 4096,
completion_payload: None,
},
&mut state,
)
.await;
handle_worker_event(
WorkerEvent::Progress {
key: key.clone(),
next_offset: 512,
total_size: 4096,
},
&mut state,
)
.await;
assert!(state.known_totals().get(&key).is_none());
}
#[tokio::test]
async fn download_failed_prefers_event_total_when_scheduler_has_none() {
let (mut state, key) = live_download_state("failed_event").await;
let records = attach_capture(&state);
handle_worker_event(
WorkerEvent::Failed {
key: key.clone(),
error: MeowError::from_code(InnerErrorCode::Unknown, "boom".to_string()),
total_size: 8192,
},
&mut state,
)
.await;
let failed = wait_for_record(&records, |r| {
matches!(r.status(), TransferStatus::Failed(_))
})
.await;
assert_eq!(failed.total_size(), 8192, "无 known_totals 时终态用事件携带的 total");
}
#[tokio::test]
async fn download_canceled_prefers_event_total_when_scheduler_has_none() {
let (mut state, key) = live_download_state("canceled_event").await;
let records = attach_capture(&state);
handle_worker_event(
WorkerEvent::Canceled {
key: key.clone(),
total_size: 8192,
},
&mut state,
)
.await;
let canceled =
wait_for_record(&records, |r| matches!(r.status(), TransferStatus::Canceled)).await;
assert_eq!(canceled.total_size(), 8192);
assert!(!state.groups().contains_key(&key));
}
#[tokio::test]
async fn download_failed_with_no_total_anywhere_keeps_zero() {
let (mut state, key) = live_download_state("failed_zero").await;
let records = attach_capture(&state);
handle_worker_event(
WorkerEvent::Failed {
key: key.clone(),
error: MeowError::from_code(InnerErrorCode::Unknown, "boom".to_string()),
total_size: 0,
},
&mut state,
)
.await;
let failed = wait_for_record(&records, |r| {
matches!(r.status(), TransferStatus::Failed(_))
})
.await;
assert_eq!(failed.total_size(), 0);
assert_eq!(failed.progress(), 0.0);
}
#[tokio::test]
async fn upload_failed_keeps_preset_total() {
let (mut state, key, total) = live_upload_state().await;
let records = attach_capture(&state);
handle_worker_event(
WorkerEvent::Failed {
key: key.clone(),
error: MeowError::from_code(InnerErrorCode::Unknown, "boom".to_string()),
total_size: 0,
},
&mut state,
)
.await;
let failed = wait_for_record(&records, |r| {
matches!(r.status(), TransferStatus::Failed(_))
})
.await;
assert_eq!(
failed.total_size(),
total,
"上传终态 total 必须保持预设值(文件大小)"
);
}
#[tokio::test]
async fn completed_with_zero_event_total_falls_back_to_known_total() {
let (mut state, key) = live_download_state("complete_fallback").await;
let records = attach_capture(&state);
handle_worker_event(
WorkerEvent::Progress {
key: key.clone(),
next_offset: 4096,
total_size: 4096,
},
&mut state,
)
.await;
handle_worker_event(
WorkerEvent::Completed {
key: key.clone(),
total_size: 0,
completion_payload: None,
},
&mut state,
)
.await;
let complete =
wait_for_record(&records, |r| matches!(r.status(), TransferStatus::Complete)).await;
assert_eq!(complete.total_size(), 4096);
assert!((complete.progress() - 1.0).abs() < f32::EPSILON);
}
#[tokio::test]
async fn zero_total_progress_emit_keeps_known_total() {
let (mut state, key) = live_download_state("emit_no_clobber").await;
let records = attach_capture(&state);
handle_worker_event(
WorkerEvent::Progress {
key: key.clone(),
next_offset: 512,
total_size: 4096,
},
&mut state,
)
.await;
handle_worker_event(
WorkerEvent::Progress {
key: key.clone(),
next_offset: 600,
total_size: 0,
},
&mut state,
)
.await;
let second = wait_for_record(&records, |r| {
matches!(r.status(), TransferStatus::Transmission) && r.progress() > 0.13
})
.await;
assert_eq!(second.total_size(), 4096);
}
#[tokio::test]
async fn stray_failed_after_teardown_creates_no_orphan_known_total() {
let (mut state, key) = live_download_state("stray_failed").await;
handle_worker_event(
WorkerEvent::Completed {
key: key.clone(),
total_size: 4096,
completion_payload: None,
},
&mut state,
)
.await;
handle_worker_event(
WorkerEvent::Failed {
key: key.clone(),
error: MeowError::from_code(InnerErrorCode::Unknown, "late".to_string()),
total_size: 9999,
},
&mut state,
)
.await;
assert!(state.known_totals().get(&key).is_none());
}
}