refactor(download): simplify finalization and remove the retired worker

This commit is contained in:
zarzet committed 2026-09-30 21:44:31 +07:00
1 parent 3a63918e1f
commit b8e3a4c166
5 files changed
+32 -370

No files matched your search

@@ -10,8 +10,6 @@ import android.content.Intent
import android.content.pm.ServiceInfo
import android.net.ConnectivityManager
import android.net.Network
import android.net.NetworkCapabilities
import android.net.NetworkRequest
import android.os.Build
import android.os.IBinder
import android.os.PowerManager
@@ -1432,295 +1430,6 @@ class DownloadService : Service() {
}
}
private suspend fun runNativeWorker(
requests: List<NativeDownloadRequest>,
settingsJson: String,
generation: Long
) {
val rateLimitAttempts = mutableMapOf<String, Int>()
val progressCoordinatorJob = startNativeWorkerProgressCoordinator(generation)
try {
var requestIndex = 0
while (requestIndex < requests.size) {
val request = requests[requestIndex]
while (isNativeWorkerPaused() &&
!nativeWorkerCancelRequested &&
generation == nativeWorkerGeneration
) {
writeNativeWorkerSnapshot(
isRunning = true,
isPaused = true,
currentItemId = request.itemId,
message = nativeWorkerPauseMessage(),
settingsJson = settingsJson,
includeItems = true
)
delay(500)
}
if (nativeWorkerCancelRequested || generation != nativeWorkerGeneration) {
break
}
var retryCurrentRequest = false
nativeWorkerCurrentItemId = request.itemId
currentTrackName = request.trackName
currentArtistName = request.artistName
currentStatus = "preparing"
lastProgress = 0L
lastTotal = 0L
updateNotification(0, 0)
updateNativeWorkerItem(request.itemId) {
it.status = "preparing"
it.progress = 0.0
it.bytesReceived = 0L
it.bytesTotal = 0L
it.error = ""
it.resultJson = null
}
writeNativeWorkerSnapshot(
isRunning = true,
isPaused = false,
currentItemId = request.itemId,
message = "Preparing",
settingsJson = settingsJson,
includeItems = true
)
var directoryScope: AutoCloseable? = null
try {
directoryScope = coreBackend.openDownloadDirectoryForRequest(request.requestJson)
coreBackend.initItemProgress(request.itemId)
val response = SafDownloadHandler.handle(this, request.requestJson, coreBackend)
if (generation != nativeWorkerGeneration) {
// Superseded while blocked in the download call; the
// new run owns the shared state now.
break
}
var result = JSONObject(response)
if (result.optBoolean("success", false)) {
currentStatus = "finalizing"
updateNativeWorkerItem(request.itemId) {
it.status = "finalizing"
it.progress = 0.95
it.error = ""
}
writeNativeWorkerSnapshot(
isRunning = true,
isPaused = false,
currentItemId = request.itemId,
message = "Finalizing",
settingsJson = settingsJson
)
result = NativeDownloadFinalizer.finalize(
this,
request.itemId,
request.requestJson,
request.itemJson,
result,
settingsJson
) {
nativeWorkerCancelRequested ||
isNativeWorkerPaused() ||
generation != nativeWorkerGeneration
}
}
if (result.optBoolean("success", false)) {
result.optJSONObject("replaygain")?.let { replayGain ->
synchronized(nativeReplayGainEntries) {
nativeReplayGainEntries.add(JSONObject(replayGain.toString()))
}
}
updateNativeWorkerItem(request.itemId) {
it.status = "completed"
it.progress = 1.0
it.error = ""
it.resultJson = result
}
writeNativeReplayGainJournal()
writeNativeAlbumReplayGainIfComplete()
} else {
val errorType = result.optString("error_type")
val errorMessage = result.optString("error")
if (errorType == "cancelled" &&
isNativeWorkerPaused() &&
!nativeWorkerCancelRequested
) {
updateNativeWorkerItem(request.itemId) {
it.status = "queued"
it.progress = 0.0
it.bytesReceived = 0L
it.bytesTotal = 0L
it.error = ""
it.resultJson = null
}
writeNativeWorkerSnapshot(
isRunning = true,
isPaused = true,
currentItemId = request.itemId,
message = "Paused",
settingsJson = settingsJson,
includeItems = true
)
retryCurrentRequest = true
} else if (NativeWorkerPolicy.shouldRetryRateLimit(
errorType = errorType,
errorMessage = errorMessage,
attempts = rateLimitAttempts[request.itemId] ?: 0,
)
) {
rateLimitAttempts[request.itemId] =
(rateLimitAttempts[request.itemId] ?: 0) + 1
val delaySeconds = NativeWorkerPolicy.rateLimitDelaySeconds(
retryAfterSeconds = result
.optInt("retry_after_seconds", 0)
.takeIf { it > 0 },
errorMessage = errorMessage,
)
currentStatus = "rate_limited"
updateNativeWorkerItem(request.itemId) {
it.status = "queued"
it.progress = 0.0
it.bytesReceived = 0L
it.bytesTotal = 0L
it.error = "Rate limited, retrying in ${delaySeconds}s"
it.resultJson = null
}
writeNativeWorkerSnapshot(
isRunning = true,
isPaused = isNativeWorkerPaused(),
currentItemId = request.itemId,
message = "Rate limited, retrying in ${delaySeconds}s",
settingsJson = settingsJson,
includeItems = true,
)
updateNotification(0L, 0L)
delay(delaySeconds * 1000L)
retryCurrentRequest = true
} else if (NativeWorkerPolicy.isVerificationRequired(
errorType = errorType,
errorMessage = errorMessage,
)
) {
nativeWorkerVerificationPaused = true
currentStatus = "verification_required"
updateNativeWorkerItem(request.itemId) {
it.status = "failed"
it.error = errorMessage
it.resultJson = result
}
writeNativeReplayGainJournal()
writeNativeWorkerSnapshot(
isRunning = true,
isPaused = true,
currentItemId = request.itemId,
message = "Verification required",
lastResult = result,
settingsJson = settingsJson,
includeItems = true,
)
// Publish immediately. If Flutter is alive it will
// replace this same notification ID while owning
// the interactive challenge; if Flutter is
// suspended, the native alert remains visible.
showNativeVerificationRequired(request, result)
updateNotification(0L, 0L)
retryCurrentRequest = true
} else {
updateNativeWorkerItem(request.itemId) {
it.status = if (errorType == "cancelled") {
"skipped"
} else {
"failed"
}
it.error = errorMessage
it.resultJson = result
}
writeNativeReplayGainJournal()
}
}
if (!retryCurrentRequest) {
writeNativeWorkerSnapshot(
isRunning = true,
isPaused = false,
currentItemId = request.itemId,
message = if (result.optBoolean("success", false)) "Completed" else "Failed",
lastResult = result,
settingsJson = settingsJson,
includeItems = true
)
}
} catch (e: CancellationException) {
if (nativeWorkerCancelRequested &&
generation == nativeWorkerGeneration
) {
updateNativeWorkerItem(request.itemId) {
it.status = "skipped"
it.error = "Cancelled"
}
}
throw e
} catch (e: Exception) {
updateNativeWorkerItem(request.itemId) {
it.status = "failed"
it.error = e.message ?: "Native download failed"
}
writeNativeReplayGainJournal()
writeNativeWorkerSnapshot(
isRunning = true,
isPaused = false,
currentItemId = request.itemId,
message = e.message ?: "Native download failed",
settingsJson = settingsJson,
includeItems = true
)
} finally {
directoryScope?.close()
finishNativePauseCancellation(request.itemId, generation)
updateNativeWorkerItemProgress(request.itemId)
try {
coreBackend.clearItemProgress(request.itemId)
} catch (_: Exception) {
}
}
if (!retryCurrentRequest) {
if (nativeWorkerCurrentItemId == request.itemId) {
nativeWorkerCurrentItemId = ""
}
requestIndex++
}
}
} finally {
stopNativeWorkerProgressCoordinator(progressCoordinatorJob)
if (generation == nativeWorkerGeneration) {
cancelScheduledNativeWorkerItemsSnapshot()
if (!nativeWorkerCancelRequested) {
flushNativeAlbumReplayGainJournalIfComplete()
}
val counts = nativeWorkerCounts()
val shouldNotifyCompletion =
NativeWorkerPolicy.shouldNotifyQueueComplete(
cancelRequested = nativeWorkerCancelRequested,
completed = counts.completed,
failed = counts.failed,
)
currentStatus = "finalizing"
releaseIdleDownloadMemory()
writeNativeWorkerSnapshot(
isRunning = false,
isPaused = false,
currentItemId = "",
message = if (nativeWorkerCancelRequested) "Cancelled" else "Finished",
settingsJson = settingsJson,
includeItems = true
)
stopForegroundService(cancelNativeWorker = false)
if (shouldNotifyCompletion) {
showNativeQueueComplete(counts)
}
}
}
}
private fun releaseIdleDownloadMemory() {
try {
// All workers and album tagging have finished. Release idle backend
@@ -1150,7 +1150,10 @@ extension _DownloadQueueFinalization on DownloadQueueNotifier {
return null;
}
var preserve = false;
String? producedFileName;
final rawFileName =
(result['file_name'] as String?) ?? context.safFileName ?? 'track';
final newFileName =
'${rawFileName.replaceFirst(RegExp(r'\.[^.]+$'), '')}.flac';
final newUri = await _replaceSafFileVia(
uri: filePath,
treeUri: treeUri,
@@ -1176,13 +1179,6 @@ extension _DownloadQueueFinalization on DownloadQueueNotifier {
'as FLAC and embedding metadata.',
);
await embedFlacMetadata(tempPath);
final rawFileName =
(result['file_name'] as String?) ??
context.safFileName ??
'track';
final baseName = rawFileName.replaceFirst(RegExp(r'\.[^.]+$'), '');
final newFileName = '$baseName.flac';
producedFileName = newFileName;
return (tempPath, newFileName);
}
final flacPath = await FFmpegService.convertM4aToFlac(tempPath);
@@ -1191,13 +1187,6 @@ extension _DownloadQueueFinalization on DownloadQueueNotifier {
}
addCleanup(flacPath);
await embedFlacMetadata(flacPath);
final rawFileName =
(result['file_name'] as String?) ??
context.safFileName ??
'track';
final baseName = rawFileName.replaceFirst(RegExp(r'\.[^.]+$'), '');
final newFileName = '$baseName.flac';
producedFileName = newFileName;
return (flacPath, newFileName);
},
);
@@ -1207,7 +1196,7 @@ extension _DownloadQueueFinalization on DownloadQueueNotifier {
if (newUri == null) {
return null;
}
result['file_name'] = producedFileName;
result['file_name'] = newFileName;
markFinalOutputAsFlac();
return newUri;
}
@@ -5,13 +5,6 @@ part of 'download_queue_provider.dart';
/// the native download call, decrypt/convert/embed finalization, and the
/// history hand-off. Queue scheduling stays in the main file.
extension _SingleItemDownload on DownloadQueueNotifier {
bool _isStorageWriteFailure(Map<String, dynamic> result) {
return isStorageWriteFailure(
errorType: result['error_type']?.toString(),
errorMessage: (result['error'] ?? result['message'])?.toString(),
);
}
Future<String?> _runPostProcessingHooks(
String filePath,
Track track,
@@ -520,7 +513,11 @@ class _DownloadRun {
outputDir: effectiveOutputDir,
);
if (result['success'] != true && n._isStorageWriteFailure(result)) {
if (result['success'] != true &&
isStorageWriteFailure(
errorType: result['error_type']?.toString(),
errorMessage: (result['error'] ?? result['message'])?.toString(),
)) {
if (n._isLocallyCancelled(item.id)) {
_log.i('Download was cancelled before storage fallback, skipping');
return false;
@@ -931,7 +928,7 @@ class _DownloadRun {
if (shouldForceDashSafM4aHandling) {
_log.w(
'SAF file is labeled FLAC but backend returned DASH/M4A stream; converting it back to FLAC.',
'SAF file is labeled FLAC but backend returned an M4A stream; converting it back to FLAC.',
);
}
@@ -1268,10 +1265,8 @@ class _DownloadRun {
: renamedPath;
await file.rename(finalRenamedPath);
targetPath = finalRenamedPath;
filePath = finalRenamedPath;
} else {
filePath = targetPath;
}
filePath = targetPath;
if (metadataEmbeddingEnabled) {
n.updateItemStatus(
@@ -1288,9 +1283,7 @@ class _DownloadRun {
}
Future<void> _convertLocalM4aToFlac(String currentFilePath) async {
_log.d(
'M4A file detected (Hi-Res DASH stream), attempting conversion to FLAC...',
);
_log.d('M4A file detected, attempting conversion to FLAC...');
try {
final file = File(currentFilePath);
@@ -1348,23 +1341,7 @@ class _DownloadRun {
_log.d('Converted to FLAC: $flacPath');
_log.d('Embedding metadata and cover to converted FLAC...');
try {
final backendGenre = result['genre'] as String?;
final backendLabel = result['label'] as String?;
final backendCopyright = result['copyright'] as String?;
if (backendGenre != null ||
backendLabel != null ||
backendCopyright != null) {
_log.d(
'Extended metadata from backend - Genre: $backendGenre, Label: $backendLabel, Copyright: $backendCopyright',
);
}
await _embedFinalMetadata(flacPath, format: 'flac');
_log.d('Metadata and cover embedded successfully');
} catch (e) {
_log.w('Warning: Failed to embed metadata/cover: $e');
}
await _embedFinalMetadata(flacPath, format: 'flac');
} else {
_log.w('FFmpeg conversion returned null, keeping M4A file');
}
@@ -1384,16 +1361,17 @@ class _DownloadRun {
resultOutputExt == '.ogg';
final isMp3File =
currentFilePath.endsWith('.mp3') || resultOutputExt == '.mp3';
final ext = isOpusFile
? (resultOutputExt == '.ogg' ? '.ogg' : '.opus')
final (ext, formatName) = isOpusFile
? (resultOutputExt == '.ogg' ? '.ogg' : '.opus', 'Opus')
: isMp3File
? '.mp3'
: '.flac';
final formatName = isOpusFile
? 'Opus'
: isMp3File
? 'MP3'
: 'FLAC';
? ('.mp3', 'MP3')
: ('.flac', 'FLAC');
// MP3 takes precedence for the embed format when both checks match.
final embedFormat = isMp3File
? 'mp3'
: isOpusFile
? 'opus'
: 'flac';
_log.d(
'SAF $formatName detected, embedding metadata and cover via temp file...',
);
@@ -1416,11 +1394,6 @@ class _DownloadRun {
// so a sidecar .lrc written next to it would be orphaned;
// the SAF .lrc is written by _saveExternalLrc after publish,
// reusing the LRC fetched here via result['lyrics_lrc'].
final embedFormat = isMp3File
? 'mp3'
: isOpusFile
? 'opus'
: 'flac';
final fetchedLrc = await _embedFinalMetadata(
tempPath,
format: embedFormat,
+1 -1
View File
@@ -1,7 +1,7 @@
import 'package:spotiflac_android/models/settings.dart';
import 'package:spotiflac_android/services/library_database.dart';
/// Field group keys understood by the Go re-enrich backend.
/// Field group keys understood by the native re-enrich backend.
class ReEnrichFields {
static const String cover = 'cover';
static const String lyrics = 'lyrics';
+7 -16
View File
@@ -505,16 +505,6 @@ class FFmpegService {
}.contains(normalized);
}
/// Probes the source audio bit depth (bits_per_raw_sample, falling back to
/// bits_per_sample). Returns null when unknown.
static Future<int?> probeBitDepth(String filePath) async {
return (await _probePrimaryAudioProperties(filePath)).bitDepth;
}
static Future<int?> probeSampleRate(String filePath) async {
return (await _probePrimaryAudioProperties(filePath)).sampleRate;
}
/// Returns `true` when [filePath] starts with the native FLAC magic bytes
/// (`fLaC`). Useful to distinguish a real FLAC file from a FLAC-in-MP4
/// container that carries a `.flac` extension or claims codec=flac.
@@ -2102,8 +2092,9 @@ class FFmpegService {
}
/// Convert to uncompressed PCM (WAV or AIFF), preserving bit depth when known.
/// Tags and cover are written natively into an embedded ID3 chunk by the Go
/// backend (RIFF "id3 " for WAV, "ID3 " for AIFF) for full-fidelity tagging.
/// Tags and cover are written natively into an embedded ID3 chunk by the
/// native backend (RIFF "id3 " for WAV, "ID3 " for AIFF) for full-fidelity
/// tagging.
static Future<String?> _convertToPcm({
required String inputPath,
required Map<String, String> metadata,
@@ -2125,7 +2116,7 @@ class FFmpegService {
final outputPath = outputPlan.workingPath;
var depth = targetBitDepth ?? sourceBitDepth;
if (depth == null || depth <= 0) {
depth = await probeBitDepth(inputPath);
depth = (await _probePrimaryAudioProperties(inputPath)).bitDepth;
}
final use24 = depth != null && depth >= 24;
final codec = isAiff
@@ -2198,9 +2189,9 @@ class FFmpegService {
);
}
/// Writes tags + cover into a WAV/AIFF file via the Go native ID3-chunk
/// Writes tags + cover into a WAV/AIFF file via the native backend ID3-chunk
/// writer (PlatformBridge.editFileMetadata). Maps Vorbis-style metadata keys
/// to the lowercase field names the Go editor expects.
/// to the lowercase field names the native editor expects.
static Future<bool> _embedChunkTagsNative(
String path,
Map<String, String> vorbisMetadata,
@@ -2284,7 +2275,7 @@ class FFmpegService {
/// Each track is extracted with `-c copy` (no re-encoding) and metadata is embedded.
/// [audioPath] is the source audio file (FLAC, WAV, etc.)
/// [outputDir] is where individual track files will be saved
/// [tracks] is the list of track split info from the Go CUE parser
/// [tracks] is the list of track split info from the native CUE parser
/// [albumMetadata] contains album-level metadata (artist, album, genre, date)
static Future<List<String>?> splitCueToTracks({
required String audioPath,