using System.Collections.ObjectModel; namespace FluidExplorer.Services.Operations; /// /// 可异步等待的自动/手动复位事件:用于"暂停"语义。 /// 关键点: 必须带 RunContinuationsAsynchronously, /// 否则 会在调用线程(通常是 UI 线程)上同步执行拷贝循环的续体, /// 从而把磁盘 IO 拖回 UI 线程。 /// internal sealed class AsyncManualResetEvent { private TaskCompletionSource _tcs; public AsyncManualResetEvent(bool initialState = true) => _tcs = Create(initialState); public bool IsSet => _tcs.Task.IsCompleted; public Task WaitAsync() => _tcs.Task; public void Set() => _tcs.TrySetResult(); public void Reset() { while (true) { var current = _tcs; if (!current.Task.IsCompleted) return; var fresh = Create(false); if (ReferenceEquals(Interlocked.CompareExchange(ref _tcs, fresh, current), current)) return; } } private static TaskCompletionSource Create(bool completed) { var tcs = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); if (completed) tcs.SetResult(); return tcs; } } /// 引擎内部对一个作业的控制块:取消令牌 + 暂停闸门 + 调度车道标记。 internal sealed class JobContext { public JobContext(FileOperationJob job) => Job = job; public FileOperationJob Job { get; } public CancellationTokenSource Cts { get; } = new(); /// 初始为 Set(未暂停)。Reset 即暂停,Set 即继续。 public AsyncManualResetEvent PauseGate { get; } = new(initialState: true); /// /// true = 走"瞬时车道":同卷 Move / 重命名 / 新建文件夹这类不搬字节的操作, /// 可以和其他作业并行,不必排队等大拷贝。 /// public bool InstantLane { get; set; } public bool Started { get; set; } } /// /// 作业队列:后台 worker 严格按入队顺序调度。 /// 默认串行(避免多任务同时读盘造成磁盘抖动、拖慢整体吞吐), /// 但"瞬时操作"(同卷 Move / 重命名 / 新建文件夹)允许并行,因为它们只做元数据操作。 /// internal sealed class JobQueue : IDisposable { /// 瞬时车道最大并行度,防止一次排入上千个瞬时作业时线程爆炸。 private const int MaxInstantParallelism = 4; private readonly Func _executor; private readonly ObservableCollection _jobs = []; private readonly List _pending = []; private readonly Dictionary _contexts = []; private readonly HashSet _instantTasks = []; private readonly SemaphoreSlim _signal = new(0); private readonly CancellationTokenSource _shutdown = new(); private readonly object _gate = new(); private Task? _dispatcher; private bool _disposed; public JobQueue(Func executor) => _executor = executor; public ObservableCollection Jobs => _jobs; public event EventHandler? JobsChanged; public void Start() { lock (_gate) { _dispatcher ??= Task.Run(DispatcherLoopAsync); } } public void Enqueue(FileOperationJob job, bool instantLane) { lock (_gate) { if (_disposed) return; var ctx = new JobContext(job) { InstantLane = instantLane }; _contexts[job.Id] = ctx; _pending.Add(ctx); _jobs.Add(job); } _signal.Release(); Raise(); } public void Pause(Guid jobId) { JobContext? ctx; lock (_gate) _contexts.TryGetValue(jobId, out ctx); if (ctx is null) return; // 只有已经在跑的作业才能暂停(Queued 状态的作业由模型定义为 CanPause=false)。 if (!ctx.Job.CanPause) return; ctx.PauseGate.Reset(); if (ctx.Job.State == JobState.Running) ctx.Job.State = JobState.Paused; Raise(); } public void Resume(Guid jobId) { JobContext? ctx; lock (_gate) _contexts.TryGetValue(jobId, out ctx); if (ctx is null || ctx.Job.IsFinished) return; ctx.PauseGate.Set(); if (ctx.Job.State == JobState.Paused) ctx.Job.State = JobState.Running; Raise(); } public void Cancel(Guid jobId) { JobContext? ctx; lock (_gate) _contexts.TryGetValue(jobId, out ctx); if (ctx is null || ctx.Job.IsFinished) return; ctx.Cts.Cancel(); // 让处于暂停中的拷贝循环立刻被唤醒,从而在同一粒度内观察到取消。 ctx.PauseGate.Set(); if (ctx.Job.State == JobState.Queued) { // 还没开始跑:直接落终态,调度循环只挑 Queued 的作业,因此会自然跳过。 ctx.Job.State = JobState.Cancelled; ctx.Job.Error = AppendLine(ctx.Job.Error, "已取消(尚未开始)。"); } Raise(); } public void ClearFinished() { lock (_gate) { for (var i = _jobs.Count - 1; i >= 0; i--) { if (!_jobs[i].IsFinished) continue; _contexts.Remove(_jobs[i].Id); _jobs.RemoveAt(i); } _pending.RemoveAll(c => c.Job.IsFinished); } Raise(); } public void Raise() => JobsChanged?.Invoke(this, EventArgs.Empty); internal static string AppendLine(string? existing, string line) => string.IsNullOrEmpty(existing) ? line : existing + Environment.NewLine + line; private async Task DispatcherLoopAsync() { var token = _shutdown.Token; while (!token.IsCancellationRequested) { try { await _signal.WaitAsync(token).ConfigureAwait(false); } catch (OperationCanceledException) { return; } catch (ObjectDisposedException) { return; } while (!token.IsCancellationRequested) { JobContext? ctx; lock (_gate) { ctx = _pending.FirstOrDefault(c => c.Job.State == JobState.Queued); if (ctx is not null && ctx.InstantLane && _instantTasks.Count >= MaxInstantParallelism) ctx = null; // 车道满了:本轮不挑,等下一轮(等价于退化成串行) if (ctx is not null) ctx.Started = true; } if (ctx is null) break; if (ctx.InstantLane) { var task = Task.Run(() => RunSafelyAsync(ctx), CancellationToken.None); lock (_gate) _instantTasks.Add(task); _ = task.ContinueWith( t => { lock (_gate) _instantTasks.Remove(t); _signal.Release(); // 唤醒调度循环,检查是否还有排队的作业 }, CancellationToken.None, TaskContinuationOptions.ExecuteSynchronously, TaskScheduler.Default); } else { await RunSafelyAsync(ctx).ConfigureAwait(false); } lock (_gate) _pending.Remove(ctx); } } } private async Task RunSafelyAsync(JobContext ctx) { try { await _executor(ctx).ConfigureAwait(false); } catch (OperationCanceledException) { if (!ctx.Job.IsFinished) ctx.Job.State = JobState.Cancelled; } catch (Exception ex) { ctx.Job.Error = AppendLine(ctx.Job.Error, ex.Message); if (!ctx.Job.IsFinished) ctx.Job.State = JobState.Failed; } finally { Raise(); } } public void Dispose() { lock (_gate) { if (_disposed) return; _disposed = true; } _shutdown.Cancel(); try { _signal.Release(); } catch (Exception) { /* 忽略 */ } lock (_gate) { foreach (var ctx in _contexts.Values) { try { ctx.Cts.Cancel(); } catch (Exception) { /* 忽略 */ } ctx.Cts.Dispose(); } _contexts.Clear(); } _signal.Dispose(); _shutdown.Dispose(); } }