Дезайн ыыыы
This commit is contained in:
parent
403d7086d4
commit
e7c02c7feb
174 changed files with 15802 additions and 9725 deletions
|
|
@ -1,6 +1,7 @@
|
|||
use std::sync::Arc;
|
||||
use std::sync::atomic::{AtomicU32, AtomicUsize, Ordering};
|
||||
use std::time::{Duration, Instant};
|
||||
use tauri::{AppHandle, Emitter};
|
||||
use tokio::time::sleep;
|
||||
|
||||
use crate::domain::transfer::{TaskState, TransferKind, TransferTask};
|
||||
|
|
@ -11,6 +12,44 @@ use crate::transfer::queue::TransferQueue;
|
|||
|
||||
const POLL_INTERVAL_MS: u64 = 200;
|
||||
|
||||
#[derive(Clone, serde::Serialize)]
|
||||
struct ProgressPayload {
|
||||
id: String,
|
||||
transferred_bytes: u64,
|
||||
speed: u64,
|
||||
eta_secs: Option<u64>,
|
||||
}
|
||||
|
||||
#[derive(Clone, serde::Serialize)]
|
||||
struct StatePayload {
|
||||
id: String,
|
||||
state: TaskState,
|
||||
}
|
||||
|
||||
fn emit_progress(app: &AppHandle, queue: &TransferQueue, id: &str) {
|
||||
if let Some(t) = queue.get(id) {
|
||||
let _ = app.emit(
|
||||
"transfer-progress",
|
||||
ProgressPayload {
|
||||
id: t.id,
|
||||
transferred_bytes: t.transferred_bytes,
|
||||
speed: t.speed.unwrap_or(0),
|
||||
eta_secs: t.eta_secs,
|
||||
},
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
fn emit_state(app: &AppHandle, id: &str, state: TaskState) {
|
||||
let _ = app.emit(
|
||||
"transfer-state-changed",
|
||||
StatePayload {
|
||||
id: id.to_string(),
|
||||
state,
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
/// Уменьшает счётчик активных передач при выходе из скоупа — в том числе при
|
||||
/// панике внутри задачи (Drop отрабатывает и во время размотки стека), иначе
|
||||
/// упавшая задача навсегда съедала бы один слот параллелизма.
|
||||
|
|
@ -30,6 +69,7 @@ pub fn spawn_worker(
|
|||
registry: Arc<RemoteRegistry>,
|
||||
rt_handle: tokio::runtime::Handle,
|
||||
max_concurrent: Arc<AtomicU32>,
|
||||
app: AppHandle,
|
||||
) {
|
||||
let in_flight = Arc::new(AtomicUsize::new(0));
|
||||
rt_handle.spawn(async move {
|
||||
|
|
@ -44,14 +84,16 @@ pub fn spawn_worker(
|
|||
break;
|
||||
};
|
||||
queue.update_state(&task.id, TaskState::Running);
|
||||
emit_state(&app, &task.id, TaskState::Running);
|
||||
in_flight.fetch_add(1, Ordering::Relaxed);
|
||||
|
||||
let queue = queue.clone();
|
||||
let registry = registry.clone();
|
||||
let guard_counter = in_flight.clone();
|
||||
let app = app.clone();
|
||||
tokio::spawn(async move {
|
||||
let _guard = InFlightGuard(guard_counter);
|
||||
run_transfer(task, queue, registry).await;
|
||||
run_transfer(task, queue, registry, app).await;
|
||||
});
|
||||
}
|
||||
sleep(Duration::from_millis(POLL_INTERVAL_MS)).await;
|
||||
|
|
@ -60,11 +102,21 @@ pub fn spawn_worker(
|
|||
}
|
||||
|
||||
/// Выполняет одну задачу передачи целиком и записывает итог в очередь.
|
||||
async fn run_transfer(task: TransferTask, queue: TransferQueue, registry: Arc<RemoteRegistry>) {
|
||||
async fn run_transfer(
|
||||
task: TransferTask,
|
||||
queue: TransferQueue,
|
||||
registry: Arc<RemoteRegistry>,
|
||||
app: AppHandle,
|
||||
) {
|
||||
let fs = match registry.get(&task.connection_id) {
|
||||
Some(fs) => fs,
|
||||
None => {
|
||||
queue.update_state(&task.id, TaskState::Failed("connection not found".into()));
|
||||
emit_state(
|
||||
&app,
|
||||
&task.id,
|
||||
TaskState::Failed("connection not found".into()),
|
||||
);
|
||||
return;
|
||||
}
|
||||
};
|
||||
|
|
@ -73,6 +125,7 @@ async fn run_transfer(task: TransferTask, queue: TransferQueue, registry: Arc<Re
|
|||
// текущую скорость и реагирует на отмену/паузу, выставленные пользователем.
|
||||
let queue_for_progress = queue.clone();
|
||||
let task_id_for_progress = task.id.clone();
|
||||
let app_for_progress = app.clone();
|
||||
let throttle = Arc::new(std::sync::Mutex::new(ProgressThrottle::default()));
|
||||
let last_sample = Arc::new(std::sync::Mutex::new((Instant::now(), 0u64)));
|
||||
let on_progress: Option<Box<dyn Fn(u64) -> ProgressAction + Send>> =
|
||||
|
|
@ -100,6 +153,7 @@ async fn run_transfer(task: TransferTask, queue: TransferQueue, registry: Arc<Re
|
|||
};
|
||||
if throttle.lock().unwrap().should_emit() {
|
||||
queue_for_progress.update_progress(&task_id_for_progress, transferred, speed);
|
||||
emit_progress(&app_for_progress, &queue_for_progress, &task_id_for_progress);
|
||||
}
|
||||
ProgressAction::Continue
|
||||
}));
|
||||
|
|
@ -119,12 +173,15 @@ async fn run_transfer(task: TransferTask, queue: TransferQueue, registry: Arc<Re
|
|||
Ok(()) => {
|
||||
queue.update_state(&task.id, TaskState::Completed);
|
||||
queue.update_progress(&task.id, task.total_bytes, 0);
|
||||
emit_progress(&app, &queue, &task.id);
|
||||
emit_state(&app, &task.id, TaskState::Completed);
|
||||
}
|
||||
Err(e) => {
|
||||
let msg = e.to_string();
|
||||
if msg == "cancelled" {
|
||||
queue.update_state(&task.id, TaskState::Cancelled);
|
||||
queue.remove(&task.id);
|
||||
emit_state(&app, &task.id, TaskState::Cancelled);
|
||||
} else if msg == "paused" {
|
||||
let transferred = queue
|
||||
.get(&task.id)
|
||||
|
|
@ -132,8 +189,10 @@ async fn run_transfer(task: TransferTask, queue: TransferQueue, registry: Arc<Re
|
|||
.unwrap_or(0);
|
||||
queue.update_state(&task.id, TaskState::Paused);
|
||||
queue.update_progress(&task.id, transferred, 0);
|
||||
emit_state(&app, &task.id, TaskState::Paused);
|
||||
} else {
|
||||
queue.update_state(&task.id, TaskState::Failed(msg));
|
||||
queue.update_state(&task.id, TaskState::Failed(msg.clone()));
|
||||
emit_state(&app, &task.id, TaskState::Failed(msg));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue