Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
72 changes: 69 additions & 3 deletions openless-all/app/src-tauri/src/asr/local/download.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,20 @@ use tokio::io::{AsyncSeekExt, AsyncWriteExt};

use super::models::{model_dir, ModelId, READY_SENTINEL};

/// 进度事件最小发射间隔(毫秒)。HTTP 每 chunk 回调一次 on_progress,若全量
/// 转发,前端每秒收到上百个 IPC 事件、进度条高频刷新会「抽搐」(issue 见
/// LocalAsr 下载浮层)。按 ≥150ms 节流后肉眼平滑(约 6-7 次/秒),首条进度
/// 与 phase 事件(started/finished/cancelled/failed)不受此限。
pub(crate) const PROGRESS_EMIT_MIN_INTERVAL_MS: u64 = 150;

/// 当前 Unix 毫秒时间戳(进度节流用)。
pub(crate) fn now_millis() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0)
}

/// 下载源镜像。
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "kebab-case")]
Expand Down Expand Up @@ -245,8 +259,20 @@ pub(crate) fn build_client() -> Result<reqwest::Client> {
builder.build().context("build reqwest client failed")
}

/// 判定一个「已存在」的目标文件是否完整可信,纯函数便于单测(#686)。
/// - 大小一致 → 完整;
/// 用户主动取消下载后,清理断点续传产物(`<file>.partial` sparse 文件 +
/// `<file>.partial.idx` 块索引)。`.partial` 按 `set_len` 预分配了目标全长
/// —— 1.7B 模型即使只下了 1% 也占 1.7GB 逻辑大小,不删会让用户以为
/// 「取消失效」且磁盘占用虚高。仅用户取消(非 worker 自 abort)时调用;
/// worker 失败触发的中止保留续传点,重试可直接续传。
pub(crate) fn remove_partial_artifacts(dir: &Path, dest_paths: &[String]) {
for path in dest_paths {
let dest = dir.join(path);
let _ = std::fs::remove_file(dest.with_extension("partial"));
let _ = std::fs::remove_file(dest.with_extension("partial.idx"));
}
}

/// 判定一个「已存在」的目标文件是否完整可信,纯函数便于单测(#686)。/// - 大小一致 → 完整;
/// - 大小不符(截断 / 损坏 / 超大)→ 不完整,应删除重下;
/// - `expected_size == 0`(HF 未给出大小)→ 退回旧行为「存在即信任」,避免对未知大小
/// 的文件反复重下。
Expand Down Expand Up @@ -393,8 +419,18 @@ async fn run_download(
let model_id_emit = model_id_str.clone();
let file_path_emit = file_path.clone();
let in_flight_for_cb = Arc::clone(&in_flight_bytes);
let last_emit = Arc::new(AtomicU64::new(0));
let on_progress: Arc<dyn Fn(u64) + Send + Sync> = Arc::new(move |bytes_in_file| {
in_flight_for_cb[idx].store(bytes_in_file, Ordering::Relaxed);
// 节流:距上次 emit < 150ms 的中间进度直接丢弃(高频事件会让
// 前端进度条抽搐),in_flight 仍照常累计,下次 emit 带的是最新值。
let now = now_millis();
if now - last_emit.load(Ordering::Relaxed)
< PROGRESS_EMIT_MIN_INTERVAL_MS
{
return;
}
last_emit.store(now, Ordering::Relaxed);
let total_in_flight: u64 = in_flight_for_cb
.iter()
.map(|a| a.load(Ordering::Relaxed))
Expand Down Expand Up @@ -461,6 +497,10 @@ async fn run_download(

// 用户主动 cancel(不是我们因为错误自己 set 的)→ Cancelled
if cancel.load(Ordering::SeqCst) && !self_aborted {
// 取消 = 放弃该模型:清掉 .partial/.partial.idx,避免残留稀疏大文件
// 占满磁盘(用户取消意图明确,不留续传点)。
let dest_paths: Vec<String> = info.files.iter().map(|f| f.path.clone()).collect();
remove_partial_artifacts(&dir, &dest_paths);
emit_cancelled(app, model_id, "", 0, file_count, total_bytes);
return Ok(());
}
Expand Down Expand Up @@ -1037,7 +1077,7 @@ fn emit_cancelled(

#[cfg(test)]
mod tests {
use super::existing_file_is_complete;
use super::{existing_file_is_complete, remove_partial_artifacts};

#[test]
fn complete_when_size_matches() {
Expand All @@ -1060,4 +1100,30 @@ mod tests {
assert!(existing_file_is_complete(0, 0));
assert!(existing_file_is_complete(999, 0));
}

#[test]
fn remove_partial_artifacts_deletes_partials_keeps_complete() {
// 用户取消后:`<file>.partial` 与 `<file>.partial.idx` 应被清掉,
// 已完成/完整的目标文件不受影响。
let uniq = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos())
.unwrap_or(0);
let dir = std::env::temp_dir().join(format!("ol-asr-dl-test-{uniq}"));
std::fs::create_dir_all(&dir).unwrap();
let dest = dir.join("model.safetensors");
let partial = dest.with_extension("partial");
let idx = partial.with_extension("partial.idx");
let keep = dir.join("config.json");
for p in [&dest, &partial, &idx, &keep] {
std::fs::write(p, b"x").unwrap();
}
let dest_paths: Vec<String> = vec!["model.safetensors".into()];
remove_partial_artifacts(&dir, &dest_paths);
assert!(!partial.exists(), ".partial 应被删除");
assert!(!idx.exists(), ".partial.idx 应被删除");
assert!(dest.exists(), "完整目标文件不应被删除");
assert!(keep.exists(), "未在清单里的文件不应被删除");
let _ = std::fs::remove_dir_all(&dir);
}
}
20 changes: 19 additions & 1 deletion openless-all/app/src-tauri/src/asr/local/foundry_runtime.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@
mod imp {
use std::path::{Path, PathBuf};
use std::sync::{
atomic::{AtomicBool, Ordering},
atomic::{AtomicBool, AtomicU64, Ordering},
Arc,
};

Expand Down Expand Up @@ -97,6 +97,24 @@ mod imp {
let _lifecycle = self.lifecycle.lock().await;
self.cancel_prepare.store(false, Ordering::SeqCst);
let progress: FoundryPrepareProgressCallback = Arc::new(progress);
// 节流:SDK 的 percent 回调频率不可控(可能远高于前端可感知的
// 刷新率),percent 类事件 ≥150ms 才转发,避免进度浮层抽搐;
// phase 事件(percent=None,如 runtime/model/load 的阶段切换与
// finished/failed)不受限,保证阶段提示不丢。
let raw = Arc::clone(&progress);
let last_emit = Arc::new(AtomicU64::new(0));
let progress: FoundryPrepareProgressCallback = Arc::new(move |payload| {
if payload.percent.is_some() {
let now = crate::asr::local::download::now_millis();
if now - last_emit.load(Ordering::Relaxed)
< crate::asr::local::download::PROGRESS_EMIT_MIN_INTERVAL_MS
{
return;
}
last_emit.store(now, Ordering::Relaxed);
}
raw(payload);
});
let runtime_source = foundry_native::normalize_runtime_source(runtime_source);
Ok(self
.ensure_loaded_locked(alias, runtime_source, progress)
Expand Down
27 changes: 26 additions & 1 deletion openless-all/app/src-tauri/src/asr/local/sherpa_download.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,8 @@ use sha2::{Digest, Sha256};
use tauri::{AppHandle, Emitter};

use super::download::{
build_client, download_one, partial_actual_size, DownloadPhase, DownloadProgress, Mirror,
build_client, download_one, now_millis, partial_actual_size, DownloadPhase,
DownloadProgress, Mirror, PROGRESS_EMIT_MIN_INTERVAL_MS,
};
use super::sherpa;

Expand Down Expand Up @@ -417,8 +418,18 @@ async fn run_download(
}
let app_emit = app.clone();
let in_flight_for_cb = Arc::clone(&in_flight_bytes);
let last_emit = Arc::new(AtomicU64::new(0));
let on_progress: Arc<dyn Fn(u64) + Send + Sync> = Arc::new(move |bytes_in_file| {
in_flight_for_cb[idx].store(bytes_in_file, Ordering::Relaxed);
// 节流(同 download.rs):每 HTTP chunk 回调一次,全量转发会
// 高频刷前端进度条;in_flight 照常累计,只按 ≥150ms 转发最新值。
let now = now_millis();
if now - last_emit.load(Ordering::Relaxed)
< PROGRESS_EMIT_MIN_INTERVAL_MS
{
return;
}
last_emit.store(now, Ordering::Relaxed);
let total_in_flight: u64 = in_flight_for_cb
.iter()
.map(|bytes| bytes.load(Ordering::Relaxed))
Expand Down Expand Up @@ -479,6 +490,10 @@ async fn run_download(
}

if cancel.load(Ordering::SeqCst) && !self_aborted {
// 用户主动取消 = 放弃该模型:清掉 .partial/.partial.idx(同 qwen3 路径,
// 避免稀疏大文件占满磁盘),不留续传点。
let dest_paths: Vec<String> = info.files.iter().map(|f| f.local_path.clone()).collect();
super::download::remove_partial_artifacts(&dir, &dest_paths);
emit_cancelled(app, model_alias, file_count, total_bytes);
return Ok(());
}
Expand Down Expand Up @@ -553,7 +568,14 @@ async fn run_release_archive_download(
let app_emit = app.clone();
let model_alias_emit = model_alias.to_string();
let file_name_emit = archive.file_name.to_string();
let last_emit = Arc::new(AtomicU64::new(0));
let on_progress: Arc<dyn Fn(u64) + Send + Sync> = Arc::new(move |bytes_downloaded| {
// 节流(同 download.rs):release 包下载同样按 ≥150ms 转发进度。
let now = now_millis();
if now - last_emit.load(Ordering::Relaxed) < PROGRESS_EMIT_MIN_INTERVAL_MS {
return;
}
last_emit.store(now, Ordering::Relaxed);
let _ = app_emit.emit(
"sherpa-onnx-asr-download-progress",
DownloadProgress {
Expand Down Expand Up @@ -589,6 +611,9 @@ async fn run_release_archive_download(
.await
};
if cancel.load(Ordering::SeqCst) {
// 用户取消:release 包同样清理 .partial/.partial.idx(与多文件路径一致)。
let _ = std::fs::remove_file(archive_path.with_extension("partial"));
let _ = std::fs::remove_file(archive_path.with_extension("partial.idx"));
emit_cancelled(app, model_alias, file_count, total_bytes);
return Ok(());
}
Expand Down
3 changes: 3 additions & 0 deletions openless-all/app/src/App.tsx
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import { lazy, Suspense, useEffect, useState } from 'react';
import { Capsule } from './components/Capsule';
import { GlobalDownloadProgress } from './components/GlobalDownloadProgress';
import { detectOS, type OS } from './components/WindowChrome';
import {
checkAccessibilityPermission,
Expand Down Expand Up @@ -298,6 +299,8 @@ export function App({ isCapsule, isQa, isSelectionPolishPreview, isLessComputer,
return (
<Suspense fallback={null}>
<HotkeySettingsProvider>
{/* 全局下载进度浮层:主窗口所有页面常驻(自身监听事件,与页面解耦)。 */}
<GlobalDownloadProgress />
{platformCaps?.platform === 'android' && (
<div style={{ display: mobileQaOpen ? 'block' : 'none', height: '100%' }}>
<QaPanel
Expand Down
Loading
Loading