From cd4b8ead902d3969fde24e9dad687720afd9056a Mon Sep 17 00:00:00 2001 From: zarzet <42882290+zarzet@users.noreply.github.com> Date: Thu, 17 Sep 2026 02:59:22 +0700 Subject: [PATCH] perf(analysis): serialize and cancel owned FFmpeg sessions --- lib/services/audio_analysis_jobs.dart | 64 ++++++++++++ lib/widgets/audio_analysis_spectrogram.dart | 8 +- lib/widgets/audio_analysis_widget.dart | 108 +++++++++++++++----- test/audio_analysis_jobs_test.dart | 74 ++++++++++++++ 4 files changed, 224 insertions(+), 30 deletions(-) create mode 100644 lib/services/audio_analysis_jobs.dart create mode 100644 test/audio_analysis_jobs_test.dart diff --git a/lib/services/audio_analysis_jobs.dart b/lib/services/audio_analysis_jobs.dart new file mode 100644 index 00000000..dc580814 --- /dev/null +++ b/lib/services/audio_analysis_jobs.dart @@ -0,0 +1,64 @@ +import 'dart:async'; + +class AudioAnalysisCancelled implements Exception { + const AudioAnalysisCancelled(); +} + +/// Serializes native analysis sessions and cancels only the owned session. +/// A replacement waits for native completion before creating its output files. +class AudioAnalysisJobs { + final Future<({int id, Future completed})> Function(List) start; + final Future Function(int) cancel; + Future _tail = Future.value(); + int _generation = 0; + int? _activeId; + bool _disposed = false; + + AudioAnalysisJobs({required this.start, required this.cancel}); + + void invalidate() { + _generation++; + final id = _activeId; + if (id != null) unawaited(_cancel(id)); + } + + void dispose() { + _disposed = true; + invalidate(); + } + + Future _cancel(int id) async { + try { + await cancel(id); + } catch (_) { + // Await the native completion even if the cancellation request fails. + } + } + + Future run(List arguments) { + final generation = _generation; + void checkCurrent() { + if (_disposed || generation != _generation) { + throw const AudioAnalysisCancelled(); + } + } + + final result = _tail.then((_) async { + checkCurrent(); + final session = await start(arguments); + _activeId = session.id; + try { + if (_disposed || generation != _generation) { + await _cancel(session.id); + } + final value = await session.completed; + checkCurrent(); + return value; + } finally { + _activeId = null; + } + }); + _tail = result.then((_) {}, onError: (Object _, StackTrace _) {}); + return result; + } +} diff --git a/lib/widgets/audio_analysis_spectrogram.dart b/lib/widgets/audio_analysis_spectrogram.dart index 3991a56b..f251b695 100644 --- a/lib/widgets/audio_analysis_spectrogram.dart +++ b/lib/widgets/audio_analysis_spectrogram.dart @@ -44,7 +44,9 @@ class _SpectrogramView extends StatelessWidget { ChoiceChip( label: Text(context.l10n.audioAnalysisChannels), selected: selectedChannel < 0, - onSelected: (_) => onChannelChanged(-1), + onSelected: channelLoading + ? null + : (_) => onChannelChanged(-1), visualDensity: VisualDensity.compact, ), for (var channel = 0; channel < channels; channel++) ...[ @@ -52,7 +54,9 @@ class _SpectrogramView extends StatelessWidget { ChoiceChip( label: Text('Ch ${channel + 1}'), selected: selectedChannel == channel, - onSelected: (_) => onChannelChanged(channel), + onSelected: channelLoading + ? null + : (_) => onChannelChanged(channel), visualDensity: VisualDensity.compact, ), ], diff --git a/lib/widgets/audio_analysis_widget.dart b/lib/widgets/audio_analysis_widget.dart index b664c88a..4d86ce84 100644 --- a/lib/widgets/audio_analysis_widget.dart +++ b/lib/widgets/audio_analysis_widget.dart @@ -5,10 +5,12 @@ import 'dart:math' as math; import 'dart:typed_data'; import 'dart:ui' as ui; import 'package:ffmpeg_kit_flutter_new_full/ffmpeg_kit.dart'; +import 'package:ffmpeg_kit_flutter_new_full/ffmpeg_session.dart'; import 'package:ffmpeg_kit_flutter_new_full/ffprobe_kit.dart'; import 'package:ffmpeg_kit_flutter_new_full/return_code.dart'; import 'package:flutter/foundation.dart'; import 'package:flutter/material.dart'; +import 'package:spotiflac_android/services/audio_analysis_jobs.dart'; import 'package:spotiflac_android/widgets/settings_group.dart'; import 'package:path_provider/path_provider.dart'; import 'package:spotiflac_android/l10n/l10n.dart'; @@ -606,6 +608,29 @@ class _AudioAnalysisCardState extends State { bool _spectrogramChannelLoading = false; int _spectrogramRequestId = 0; String? _unsupportedCodec; + final _analysisJobs = AudioAnalysisJobs( + start: (arguments) async { + final completed = Completer(); + final session = await FFmpegKit.executeWithArgumentsAsync( + arguments, + completed.complete, + ); + final id = session.getSessionId(); + if (id == null) { + // Never pass null to cancel: FFmpegKit interprets it as cancel-all. + await completed.future; + throw StateError('Analysis session has no ID'); + } + return (id: id, completed: completed.future); + }, + cancel: (id) => FFmpegKit.cancel(id), + ); + + void _checkAnalysisRequest(int requestId) { + if (!mounted || requestId != _spectrogramRequestId) { + throw const AudioAnalysisCancelled(); + } + } static const _supportedExtensions = { '.flac', @@ -654,6 +679,7 @@ class _AudioAnalysisCardState extends State { } _spectrogramRequestId++; + _analysisJobs.invalidate(); _spectrogramImage?.dispose(); _spectrogramImage = null; _spectrogramChannel = -1; @@ -672,6 +698,7 @@ class _AudioAnalysisCardState extends State { @override void dispose() { _spectrogramRequestId++; + _analysisJobs.dispose(); _spectrogramImage?.dispose(); super.dispose(); } @@ -709,6 +736,7 @@ class _AudioAnalysisCardState extends State { image ??= await _generateAndCacheSpectrogram( filePath: expectedPath, analysisData: cached, + requestId: requestId, ); if (isCurrentRequest()) { setState(() { @@ -729,28 +757,30 @@ class _AudioAnalysisCardState extends State { Future _generateAndCacheSpectrogram({ String? filePath, AudioAnalysisData? analysisData, + required int requestId, }) async { + _checkAnalysisRequest(requestId); final sourcePath = filePath ?? widget.filePath; final data = analysisData ?? _data; + final channel = _spectrogramChannel; final artifact = await _generateSpectrogramForFile( sourcePath, - channel: _spectrogramChannel, + requestId: requestId, + channel: channel, durationSeconds: data?.duration, sampleRate: data?.sampleRate, channels: data?.channels, ); - await _saveSpectrogramToCache( - sourcePath, - artifact.image, - channel: _spectrogramChannel, - ); + await _saveSpectrogramToCache(sourcePath, artifact.image, channel: channel); return artifact.image; } Future _analyze({bool forceRefresh = false}) async { if (_analyzing) return; + final sourcePath = widget.filePath; + final requestId = ++_spectrogramRequestId; + _analysisJobs.invalidate(); setState(() { - _spectrogramRequestId++; _analyzing = true; _spectrogramChannelLoading = false; _error = null; @@ -764,12 +794,11 @@ class _AudioAnalysisCardState extends State { try { if (forceRefresh) { - await _clearCache(widget.filePath); + await _clearCache(sourcePath); } - final cached = forceRefresh - ? null - : await _loadFromCache(widget.filePath); + final cached = forceRefresh ? null : await _loadFromCache(sourcePath); + _checkAnalysisRequest(requestId); AudioAnalysisData data; ui.Image? image; @@ -780,24 +809,28 @@ class _AudioAnalysisCardState extends State { } data = cached; image = await _loadSpectrogramFromCache( - widget.filePath, + sourcePath, channel: _spectrogramChannel, ); } else { - final result = await _runAnalysis(widget.filePath); + final result = await _runAnalysis(sourcePath, requestId); data = result.data; image = result.spectrogramImage; - await _saveToCache(widget.filePath, data); + await _saveToCache(sourcePath, data); await _saveSpectrogramToCache( - widget.filePath, + sourcePath, image, channel: _spectrogramChannel, ); } - image ??= await _generateAndCacheSpectrogram(analysisData: data); + image ??= await _generateAndCacheSpectrogram( + filePath: sourcePath, + analysisData: data, + requestId: requestId, + ); - if (mounted) { + if (mounted && requestId == _spectrogramRequestId) { setState(() { _data = data; _spectrogramImage?.dispose(); @@ -807,8 +840,10 @@ class _AudioAnalysisCardState extends State { } else { image.dispose(); } + } on AudioAnalysisCancelled { + // A new file/request owns the state; cancelled native work is cleaned up. } on _UnsupportedAudioAnalysisCodecException catch (e) { - if (mounted) { + if (mounted && requestId == _spectrogramRequestId) { setState(() { _unsupportedCodec = e.codecLabel; _error = null; @@ -816,7 +851,7 @@ class _AudioAnalysisCardState extends State { }); } } catch (e) { - if (mounted) { + if (mounted && requestId == _spectrogramRequestId) { setState(() { _error = context.friendlyError(e); _analyzing = false; @@ -947,7 +982,11 @@ class _AudioAnalysisCardState extends State { } } - Future<_AudioAnalysisRunResult> _runAnalysis(String filePath) async { + Future<_AudioAnalysisRunResult> _runAnalysis( + String filePath, + int requestId, + ) async { + _checkAnalysisRequest(requestId); String workingPath = filePath; String? tempCopy; if (filePath.startsWith('content://')) { @@ -959,7 +998,9 @@ class _AudioAnalysisCardState extends State { } try { + _checkAnalysisRequest(requestId); final info = await _getMediaInfo(workingPath); + _checkAnalysisRequest(requestId); final unsupported = unsupportedAudioAnalysisCodecLabel(info.codecName); if (unsupported != null) { throw _UnsupportedAudioAnalysisCodecException(unsupported); @@ -971,12 +1012,14 @@ class _AudioAnalysisCardState extends State { : info.duration; spectrogram = await _generateSpectrogram( workingPath, + requestId: requestId, channel: -1, includeCutoffPlane: true, durationSeconds: effectiveDuration, sampleRate: info.sampleRate, channels: info.channels, ); + _checkAnalysisRequest(requestId); final cutoffIntensity = spectrogram.cutoffIntensity; if (cutoffIntensity == null) { throw Exception('FFmpeg spectral cutoff plane was not generated'); @@ -992,6 +1035,7 @@ class _AudioAnalysisCardState extends State { ); final levelMetrics = await _runFullStreamLevelAnalysis( workingPath, + requestId: requestId, durationSeconds: effectiveDuration, ); if (levelMetrics == null) { @@ -1044,11 +1088,13 @@ class _AudioAnalysisCardState extends State { Future<_GeneratedSpectrogram> _generateSpectrogramForFile( String filePath, { + required int requestId, required int channel, double? durationSeconds, int? sampleRate, int? channels, }) async { + _checkAnalysisRequest(requestId); String workingPath = filePath; String? tempCopy; if (filePath.startsWith('content://')) { @@ -1062,6 +1108,7 @@ class _AudioAnalysisCardState extends State { try { return await _generateSpectrogram( workingPath, + requestId: requestId, channel: channel, durationSeconds: durationSeconds, sampleRate: sampleRate, @@ -1078,6 +1125,7 @@ class _AudioAnalysisCardState extends State { Future<_GeneratedSpectrogram> _generateSpectrogram( String inputPath, { + required int requestId, required int channel, bool includeCutoffPlane = false, double? durationSeconds, @@ -1091,7 +1139,8 @@ class _AudioAnalysisCardState extends State { final cutoffPath = includeCutoffPlane ? '$rawPath.cutoff.gray' : null; try { - final session = await FFmpegKit.executeWithArguments( + _checkAnalysisRequest(requestId); + final session = await _analysisJobs.run( buildAudioSpectrogramArguments( inputPath: inputPath, outputPath: rawPath, @@ -1163,6 +1212,8 @@ class _AudioAnalysisCardState extends State { Future _changeSpectrogramChannel(int channel) async { final data = _data; if (data == null || + _spectrogramChannelLoading || + _analyzing || channel == _spectrogramChannel || channel < -1 || channel >= data.channels) { @@ -1170,6 +1221,7 @@ class _AudioAnalysisCardState extends State { } final previousChannel = _spectrogramChannel; + final sourcePath = widget.filePath; final requestId = ++_spectrogramRequestId; setState(() { _spectrogramChannel = channel; @@ -1178,20 +1230,18 @@ class _AudioAnalysisCardState extends State { ui.Image? image; try { - image = await _loadSpectrogramFromCache( - widget.filePath, - channel: channel, - ); + image = await _loadSpectrogramFromCache(sourcePath, channel: channel); if (image == null) { final artifact = await _generateSpectrogramForFile( - widget.filePath, + sourcePath, + requestId: requestId, channel: channel, durationSeconds: data.duration, sampleRate: data.sampleRate, channels: data.channels, ); image = artifact.image; - await _saveSpectrogramToCache(widget.filePath, image, channel: channel); + await _saveSpectrogramToCache(sourcePath, image, channel: channel); } if (!mounted || requestId != _spectrogramRequestId) { @@ -1378,6 +1428,7 @@ class _AudioAnalysisCardState extends State { Future<_LevelMetrics?> _runFullStreamLevelAnalysis( String inputPath, { + required int requestId, required double durationSeconds, }) async { final tempDir = await getTemporaryDirectory(); @@ -1386,7 +1437,8 @@ class _AudioAnalysisCardState extends State { '${DateTime.now().microsecondsSinceEpoch}.txt', ); try { - final session = await FFmpegKit.executeWithArguments( + _checkAnalysisRequest(requestId); + final session = await _analysisJobs.run( buildAudioMetricsArguments( inputPath: inputPath, metadataPath: metadataFile.path, diff --git a/test/audio_analysis_jobs_test.dart b/test/audio_analysis_jobs_test.dart new file mode 100644 index 00000000..1c4303a1 --- /dev/null +++ b/test/audio_analysis_jobs_test.dart @@ -0,0 +1,74 @@ +import 'dart:async'; + +import 'package:flutter_test/flutter_test.dart'; +import 'package:spotiflac_android/services/audio_analysis_jobs.dart'; + +void main() { + test( + 'replacement waits for cancelled session and skips stale queued work', + () async { + final completions = >[]; + final started = []; + final cancelled = []; + final jobs = AudioAnalysisJobs( + start: (args) async { + started.add(args.single); + final done = Completer(); + completions.add(done); + return (id: completions.length, completed: done.future); + }, + cancel: (id) async => cancelled.add(id), + ); + final first = jobs.run(['old']); + final firstCheck = expectLater( + first, + throwsA(isA()), + ); + await Future.delayed(Duration.zero); + final stale = jobs.run(['stale']); + final staleCheck = expectLater( + stale, + throwsA(isA()), + ); + jobs.invalidate(); + final latest = jobs.run(['latest']); + await Future.delayed(Duration.zero); + expect(cancelled, [1]); + expect(started, ['old']); + completions.first.complete('discarded'); + await firstCheck; + await staleCheck; + await Future.delayed(Duration.zero); + expect(started, ['old', 'latest']); + completions.last.complete('result'); + expect(await latest, 'result'); + jobs.dispose(); + }, + ); + + test( + 'dispose during session creation cancels its eventual session ID', + () async { + final creation = Completer<({int id, Future completed})>(); + final completion = Completer(); + final cancelled = []; + final jobs = AudioAnalysisJobs( + start: (_) => creation.future, + cancel: (id) async => cancelled.add(id), + ); + final result = jobs.run(['analysis']); + final check = expectLater(result, throwsA(isA())); + await Future.delayed(Duration.zero); + jobs.dispose(); + creation.complete((id: 42, completed: completion.future)); + await Future.delayed(Duration.zero); + expect(cancelled, [42]); + completion.complete('discarded'); + await check; + await expectLater( + jobs.run(['later']), + throwsA(isA()), + ); + }, + ); +}