import 'dart:async'; import 'dart:collection'; import 'package:drift/drift.dart'; import 'package:flutter/foundation.dart'; import '../data/database.dart'; import '../sources/registry.dart'; import '../sources/source.dart'; import 'cbz_store.dart'; import 'ntfy_notifier.dart'; class DownloadQueue extends ChangeNotifier { DownloadQueue({ required AppDatabase db, required SourceRegistry sources, required this._store, NtfyNotifier? ntfy, }) : _db = db, _sources = sources, _ntfy = ntfy ?? NtfyNotifier(store: _store); final AppDatabase _db; final SourceRegistry _sources; final CbzStore _store; final NtfyNotifier _ntfy; final Queue _chapterIds = Queue(); final Set _cancelled = {}; /// chapterId → jobId for open jobs created at enqueue. final Map _jobByChapter = {}; bool _running = false; int? _activeChapterId; double _progress = 0; String _message = ''; int _sessionOk = 0; int _sessionFailed = 0; bool _sessionHadWork = false; DateTime? _sessionStartedAt; final List> _pendingNtfy = []; int? get activeChapterId => _activeChapterId; double get progress => _progress; String get message => _message; bool get isRunning => _running; int get queuedCount => _chapterIds.length; String _label(String seriesTitle, String chapterLabel) => '$seriesTitle — $chapterLabel'; Future enqueueChapters(Iterable chapterIds) async { for (final id in chapterIds) { if (_chapterIds.contains(id) || id == _activeChapterId) continue; _cancelled.remove(id); final chapter = await _db.getChapter(id); if (chapter == null) continue; final title = await _db.getTitle(chapter.titleId); if (title == null) continue; _chapterIds.add(id); await (_db.update(_db.chapters)..where((c) => c.id.equals(id))).write( const ChaptersCompanion( status: Value('pending'), error: Value(null), ), ); final existing = await _db.findOpenJobForChapter(id); late final int jobId; if (existing != null) { jobId = existing.id; await (_db.update(_db.jobs)..where((j) => j.id.equals(jobId))).write( JobsCompanion( titleId: Value(title.id), chapterId: Value(id), status: const Value('queued'), progress: const Value(0), message: Value(_label(title.title, chapter.label)), finishedAt: const Value(null), ), ); } else { jobId = await _db.into(_db.jobs).insert( JobsCompanion.insert( type: 'download', titleId: Value(title.id), chapterId: Value(id), status: const Value('queued'), message: Value(_label(title.title, chapter.label)), ), ); } _jobByChapter[id] = jobId; } notifyListeners(); unawaited(_pump()); } /// Cancel / dismiss a queue row for [chapterId]. /// /// Always closes open jobs. Only resets chapter status when it was still /// pending/downloading/failed — never wipes an already-downloaded chapter. Future cancelChapter(int chapterId) async { _cancelled.add(chapterId); _chapterIds.removeWhere((id) => id == chapterId); _jobByChapter.remove(chapterId); await (_db.update(_db.jobs) ..where( (j) => j.chapterId.equals(chapterId) & j.status.isIn(['queued', 'running', 'failed']), )) .write( JobsCompanion( status: const Value('done'), message: const Value('Cancelled'), finishedAt: Value(DateTime.now()), ), ); final chapter = await _db.getChapter(chapterId); if (chapter != null && chapter.status != 'downloaded') { await (_db.update(_db.chapters)..where((c) => c.id.equals(chapterId))) .write( const ChaptersCompanion(status: Value('pending'), error: Value(null)), ); } notifyListeners(); } Future retryChapter(int chapterId) async { await enqueueChapters([chapterId]); } /// Remove offline download: delete CBZ + cache, reset chapter to pending. Future deleteDownload(int chapterId) async { final chapter = await _db.getChapter(chapterId); if (chapter == null) return; await _store.deleteDownload( cbzPath: chapter.cbzPath, cacheDir: chapter.cacheDir, ); await _closeOpenJobs(chapterId, message: 'Deleted'); await (_db.update(_db.chapters)..where((c) => c.id.equals(chapterId))).write( const ChaptersCompanion( status: Value('pending'), cbzPath: Value(null), cacheDir: Value(null), downloadedAt: Value(null), error: Value(null), ), ); } /// Mark all open jobs for [chapterId] as done (dismisses Queue rows). Future _closeOpenJobs( int chapterId, { required String message, double? progress, }) async { await (_db.update(_db.jobs) ..where( (j) => j.chapterId.equals(chapterId) & j.status.isIn(['queued', 'running', 'failed']), )) .write( JobsCompanion( status: const Value('done'), progress: progress != null ? Value(progress) : const Value.absent(), message: Value(message), finishedAt: Value(DateTime.now()), ), ); } /// Drop stale queue rows for chapters that are already downloaded. Future reconcileStaleJobs() async { final open = await (_db.select(_db.jobs) ..where((j) => j.status.isIn(['queued', 'running', 'failed']))) .get(); for (final job in open) { final chapterId = job.chapterId; if (chapterId == null) continue; final chapter = await _db.getChapter(chapterId); if (chapter?.status == 'downloaded') { await _closeOpenJobs(chapterId, message: chapter!.label, progress: 1); } } } Future _pump() async { if (_running) return; _running = true; _sessionOk = 0; _sessionFailed = 0; _sessionHadWork = false; _sessionStartedAt = DateTime.now(); _pendingNtfy.clear(); notifyListeners(); try { while (_chapterIds.isNotEmpty) { final chapterId = _chapterIds.removeFirst(); if (_cancelled.contains(chapterId)) { await _closeOpenJobs(chapterId, message: 'Cancelled'); final chapter = await _db.getChapter(chapterId); if (chapter != null && chapter.status != 'downloaded') { await (_db.update(_db.chapters) ..where((c) => c.id.equals(chapterId))) .write( const ChaptersCompanion( status: Value('pending'), error: Value(null), ), ); } _jobByChapter.remove(chapterId); continue; } _sessionHadWork = true; await _downloadOne(chapterId); } } finally { final hadWork = _sessionHadWork; final ok = _sessionOk; final failed = _sessionFailed; final started = _sessionStartedAt; final pending = List>.from(_pendingNtfy); _pendingNtfy.clear(); _running = false; _activeChapterId = null; _progress = 0; _message = ''; _sessionOk = 0; _sessionFailed = 0; _sessionHadWork = false; _sessionStartedAt = null; notifyListeners(); if (hadWork) { // Per-chapter completed/failed webhooks first, then queue-empty. await Future.wait(pending); final duration = started == null ? null : DateTime.now().difference(started); await _ntfy.queueDrained(ok: ok, failed: failed, duration: duration); } } } Future _downloadOne(int chapterId) async { _activeChapterId = chapterId; _progress = 0; final chapter = await _db.getChapter(chapterId); if (chapter == null) { await _closeOpenJobs(chapterId, message: 'Chapter missing'); return; } final title = await _db.getTitle(chapter.titleId); if (title == null) { await _closeOpenJobs(chapterId, message: 'Series missing'); return; } final label = _label(title.title, chapter.label); _message = label; notifyListeners(); var jobId = _jobByChapter[chapterId]; if (jobId == null) { final existing = await _db.findOpenJobForChapter(chapterId); if (existing != null) { jobId = existing.id; } else { jobId = await _db.into(_db.jobs).insert( JobsCompanion.insert( type: 'download', titleId: Value(title.id), chapterId: Value(chapterId), status: const Value('running'), message: Value(label), ), ); } _jobByChapter[chapterId] = jobId; } await (_db.update(_db.jobs)..where((j) => j.id.equals(jobId!))).write( JobsCompanion( titleId: Value(title.id), status: const Value('running'), progress: const Value(0), message: Value(label), finishedAt: const Value(null), ), ); await (_db.update(_db.chapters)..where((c) => c.id.equals(chapterId))).write( const ChaptersCompanion(status: Value('downloading'), error: Value(null)), ); try { final source = _sources.get(SourceId.values.byName(title.source)); final link = ChapterLink( url: chapter.url, title: chapter.label, chapterKey: chapter.chapterKey, chapterNumber: chapter.chapterNumber, ); final result = await _store.downloadChapter( source: source, comicTitle: title.title, chapter: link, onProgress: (p, _) { _progress = p; _message = label; notifyListeners(); unawaited( (_db.update(_db.jobs)..where((j) => j.id.equals(jobId!))).write( JobsCompanion(progress: Value(p), message: Value(label)), ), ); }, isCancelled: () => _cancelled.contains(chapterId), ); await (_db.update(_db.chapters)..where((c) => c.id.equals(chapterId))).write( ChaptersCompanion( status: const Value('downloaded'), cbzPath: Value(result.cbzPath), cacheDir: Value(result.cacheDir), downloadedAt: Value(DateTime.now()), error: const Value(null), ), ); // Close every open job for this chapter (covers stale duplicates). await _closeOpenJobs(chapterId, message: label, progress: 1); _sessionOk++; _pendingNtfy.add( _ntfy.chapterCompleted( ChapterCompletedPayload( seriesTitle: title.title, chapterLabel: chapter.label, imageCount: result.imageCount, rawBytes: result.rawBytes, cbzBytes: result.cbzBytes, cbzPath: result.cbzPath, cacheDir: result.cacheDir, coverUrl: title.coverUrl, ), ), ); } catch (e) { final cancelled = _cancelled.contains(chapterId); await (_db.update(_db.chapters)..where((c) => c.id.equals(chapterId))).write( ChaptersCompanion( status: Value(cancelled ? 'pending' : 'failed'), error: Value(cancelled ? null : e.toString()), ), ); await (_db.update(_db.jobs)..where((j) => j.id.equals(jobId!))).write( JobsCompanion( status: Value(cancelled ? 'done' : 'failed'), message: Value(cancelled ? 'Cancelled' : label), finishedAt: Value(DateTime.now()), ), ); if (!cancelled) { _sessionFailed++; _pendingNtfy.add( _ntfy.jobFailed( seriesTitle: title.title, chapterLabel: chapter.label, error: e.toString(), ), ); } } finally { _cancelled.remove(chapterId); _jobByChapter.remove(chapterId); } } /// Upsert series into library from a remote [SeriesDetail]. Future upsertSeries(SeriesDetail series) async { final existing = await _db.findTitleByUrl(series.listingUrl); late final int titleId; if (existing == null) { titleId = await _db.into(_db.titles).insert( TitlesCompanion.insert( source: series.source.id, url: series.listingUrl, slug: series.slug, title: series.title, coverUrl: Value(series.coverUrl), description: Value(series.description), lastCheckedAt: Value(DateTime.now()), ), ); } else { titleId = existing.id; await (_db.update(_db.titles)..where((t) => t.id.equals(titleId))).write( TitlesCompanion( title: Value(series.title), coverUrl: Value(series.coverUrl), description: Value(series.description), lastCheckedAt: Value(DateTime.now()), ), ); } for (final ch in series.chapters) { final rows = await (_db.select(_db.chapters) ..where( (c) => c.titleId.equals(titleId) & c.chapterKey.equals(ch.chapterKey), )) .get(); if (rows.isEmpty) { await _db.into(_db.chapters).insert( ChaptersCompanion.insert( titleId: titleId, chapterKey: ch.chapterKey, chapterNumber: Value(ch.chapterNumber), label: ch.title, url: ch.url, ), ); } else { await (_db.update(_db.chapters)..where((c) => c.id.equals(rows.first.id))) .write( ChaptersCompanion( label: Value(ch.title), url: Value(ch.url), chapterNumber: Value(ch.chapterNumber), ), ); } } return titleId; } }