use crate::file_transfer_record::FileTransferRecord;
use crate::ids::TaskId;
use crate::inner::group_state::RecordEntry;
use crate::inner::inner_task::InnerTask;
use crate::inner::scheduler_state::SchedulerState;
use crate::inner::task_callbacks::{CompleteCb, ProgressCb};
use crate::inner::UniqueId;
use crate::transfer_status::TransferStatus;
pub(crate) fn invoke_progress_cb(state: &SchedulerState, cb: &ProgressCb, dto: FileTransferRecord) {
let Some(submit) = state.cb_submit() else {
crate::meow_trace_log!(
"callback",
"progress callback skipped: dispatcher already taken (closing)"
);
return;
};
submit.submit_progress(cb.clone(), dto);
}
pub(crate) fn invoke_complete_cb(
state: &SchedulerState,
cb: &CompleteCb,
task_id: TaskId,
payload: Option<String>,
) {
let Some(submit) = state.cb_submit() else {
crate::meow_warn_log!(
"callback",
"complete callback skipped: dispatcher already taken (closing)"
);
return;
};
submit.submit_complete(cb.clone(), task_id, payload);
}
pub(crate) fn emit_global_progress(state: &SchedulerState, dto: FileTransferRecord) {
let listeners: Vec<ProgressCb> = match state.global_progress_listener().read() {
Ok(g) => g.iter().map(|(_, cb)| cb.clone()).collect(),
Err(_) => {
crate::meow_warn_log!(
"emit_global_progress",
"global listener lock poisoned; skip progress broadcast"
);
return;
}
};
crate::meow_trace_log!(
"emit_global_progress",
"broadcast start: listener_count={} task_id={:?}",
listeners.len(),
dto.task_id()
);
for cb in listeners {
invoke_progress_cb(state, &cb, dto.clone());
}
}
pub(crate) fn effective_total(state: &SchedulerState, key: &UniqueId, inner: &InnerTask) -> u64 {
match state.known_totals().get(key).copied() {
Some(v) if v > 0 => v,
_ => inner.total_size(),
}
}
pub(crate) fn emit_status(
state: &SchedulerState,
entry: &RecordEntry,
status: TransferStatus,
transferred: u64,
total: u64,
) {
crate::meow_trace_log!(
"emit_status",
"status emit start: task_id={:?} status={:?} transferred={} total={}",
entry.inner().task_id(),
status,
transferred,
total
);
let inner = entry.inner();
let dto = FileTransferRecord::new(
inner.task_id(),
inner.file_sign_arc(),
inner.file_name_arc(),
total,
if total == 0 {
0.0
} else {
transferred as f32 / total as f32
},
status,
inner.direction(),
);
if let Some(cb) = &entry.callbacks().progress_cb() {
invoke_progress_cb(state, cb, dto.clone());
}
emit_global_progress(state, dto);
}
#[cfg(test)]
mod tests {
use crate::inner::test_support::{live_download_state, live_download_state_with_preset};
fn inner_of(
state: &crate::inner::scheduler_state::SchedulerState,
key: &crate::inner::UniqueId,
) -> crate::inner::inner_task::InnerTask {
state
.groups()
.get(key)
.expect("group")
.entry()
.inner()
.clone()
}
#[tokio::test]
async fn effective_total_prefers_runtime_over_unset_preset() {
let (mut state, key) = live_download_state("eff_runtime").await;
state.known_totals_mut().insert(key.clone(), 4096);
let inner = inner_of(&state, &key);
assert_eq!(super::effective_total(&state, &key, &inner), 4096);
}
#[tokio::test]
async fn effective_total_returns_zero_when_nothing_known() {
let (state, key) = live_download_state("eff_zero").await;
let inner = inner_of(&state, &key);
assert_eq!(super::effective_total(&state, &key, &inner), 0);
}
#[tokio::test]
async fn effective_total_falls_back_to_preset_when_no_runtime() {
let (state, key) = live_download_state_with_preset("eff_fallback", 8192).await;
let inner = inner_of(&state, &key);
assert_eq!(super::effective_total(&state, &key, &inner), 8192);
}
#[tokio::test]
async fn effective_total_runtime_wins_over_larger_preset() {
let (mut state, key) = live_download_state_with_preset("eff_shrink", 10000).await;
state.known_totals_mut().insert(key.clone(), 4096);
let inner = inner_of(&state, &key);
assert_eq!(
super::effective_total(&state, &key, &inner),
4096,
"运行期探测值必须压过更大的预设值"
);
}
}