mirror of
https://github.com/zarzet/SpotiFLAC-Mobile.git
synced 2026-08-02 00:58:37 +02:00
refactor(android): split replaygain journal, snapshots, and network policy out of DownloadService
This commit is contained in:
@@ -67,13 +67,13 @@ class DownloadService : Service() {
|
||||
const val EXTRA_SETTINGS_JSON = "settings_json"
|
||||
const val EXTRA_REQUESTS_PATH = "requests_path"
|
||||
const val EXTRA_SETTINGS_PATH = "settings_path"
|
||||
private const val NATIVE_WORKER_STATE_FILE = "native_download_worker_state.json"
|
||||
private const val NATIVE_WORKER_PROGRESS_FILE = "native_download_worker_progress.json"
|
||||
private const val NATIVE_REPLAYGAIN_JOURNAL_FILE = "native_replaygain_journal.json"
|
||||
private const val NATIVE_WORKER_CONTRACT_VERSION = NativeDownloadFinalizer.NATIVE_WORKER_CONTRACT_VERSION
|
||||
private const val NOTIFICATION_PERCENT_TOTAL = 10_000L
|
||||
private val NATIVE_WORKER_STATE_FILE_LOCK = Any()
|
||||
private val NATIVE_REPLAYGAIN_JOURNAL_FILE_LOCK = Any()
|
||||
internal const val NATIVE_WORKER_STATE_FILE = "native_download_worker_state.json"
|
||||
internal const val NATIVE_WORKER_PROGRESS_FILE = "native_download_worker_progress.json"
|
||||
internal const val NATIVE_REPLAYGAIN_JOURNAL_FILE = "native_replaygain_journal.json"
|
||||
internal const val NATIVE_WORKER_CONTRACT_VERSION = NativeDownloadFinalizer.NATIVE_WORKER_CONTRACT_VERSION
|
||||
internal const val NOTIFICATION_PERCENT_TOTAL = 10_000L
|
||||
internal val NATIVE_WORKER_STATE_FILE_LOCK = Any()
|
||||
internal val NATIVE_REPLAYGAIN_JOURNAL_FILE_LOCK = Any()
|
||||
|
||||
private var isRunning = false
|
||||
|
||||
@@ -165,8 +165,8 @@ class DownloadService : Service() {
|
||||
// re-parsing the full state file — which grows with every completed
|
||||
// item (results embed history rows and lyrics) — once the caller has
|
||||
// already consumed that items payload.
|
||||
@Volatile private var lastStateHeaderJson: String? = null
|
||||
@Volatile private var lastStateHeaderSerial = 0L
|
||||
@Volatile internal var lastStateHeaderJson: String? = null
|
||||
@Volatile internal var lastStateHeaderSerial = 0L
|
||||
|
||||
fun getNativeWorkerSnapshot(context: Context, sinceStateSerial: Long = 0L): String {
|
||||
synchronized(NATIVE_WORKER_STATE_FILE_LOCK) {
|
||||
@@ -253,7 +253,7 @@ class DownloadService : Service() {
|
||||
}
|
||||
}
|
||||
|
||||
private data class NativeDownloadRequest(
|
||||
internal data class NativeDownloadRequest(
|
||||
val itemId: String,
|
||||
val requestJson: String,
|
||||
val trackName: String,
|
||||
@@ -261,7 +261,7 @@ class DownloadService : Service() {
|
||||
val itemJson: String
|
||||
)
|
||||
|
||||
private data class NativeWorkerItem(
|
||||
internal data class NativeWorkerItem(
|
||||
val itemId: String,
|
||||
val trackName: String,
|
||||
val artistName: String,
|
||||
@@ -274,41 +274,41 @@ class DownloadService : Service() {
|
||||
var resultJson: JSONObject? = null
|
||||
)
|
||||
|
||||
private data class NativeWorkerCounts(
|
||||
internal data class NativeWorkerCounts(
|
||||
val total: Int,
|
||||
val completed: Int,
|
||||
val failed: Int,
|
||||
val skipped: Int
|
||||
)
|
||||
|
||||
private val serviceScope = CoroutineScope(SupervisorJob() + Dispatchers.IO)
|
||||
private var nativeWorkerJob: Job? = null
|
||||
internal val serviceScope = CoroutineScope(SupervisorJob() + Dispatchers.IO)
|
||||
internal var nativeWorkerJob: Job? = null
|
||||
private var wakeLock: PowerManager.WakeLock? = null
|
||||
private var currentTrackName = ""
|
||||
private var currentArtistName = ""
|
||||
private var currentStatus = "preparing"
|
||||
internal var currentStatus = "preparing"
|
||||
private var queueCount = 0
|
||||
// Signature of the last home-screen widget push; keeps widget updates
|
||||
// event-driven (track/status/queue changes, 25% steps), never per byte.
|
||||
private var widgetSignature = ""
|
||||
private var lastProgress = 0L
|
||||
private var lastTotal = 0L
|
||||
private var nativeWorkerRunId = ""
|
||||
internal var lastProgress = 0L
|
||||
internal var lastTotal = 0L
|
||||
internal var nativeWorkerRunId = ""
|
||||
@Volatile private var nativeWorkerCurrentItemId = ""
|
||||
private val nativeWorkerItems = mutableListOf<NativeWorkerItem>()
|
||||
private val nativeReplayGainEntries = mutableListOf<JSONObject>()
|
||||
private val nativeReplayGainRequestAlbumKeys = mutableMapOf<String, String>()
|
||||
private val snapshotWriteLock = Any()
|
||||
private val snapshotWriteSerial = AtomicLong(0L)
|
||||
private var latestCommittedStateSnapshotSerial = 0L
|
||||
private var latestCommittedProgressSnapshotSerial = 0L
|
||||
internal val nativeWorkerItems = mutableListOf<NativeWorkerItem>()
|
||||
internal val nativeReplayGainEntries = mutableListOf<JSONObject>()
|
||||
internal val nativeReplayGainRequestAlbumKeys = mutableMapOf<String, String>()
|
||||
internal val snapshotWriteLock = Any()
|
||||
internal val snapshotWriteSerial = AtomicLong(0L)
|
||||
internal var latestCommittedStateSnapshotSerial = 0L
|
||||
internal var latestCommittedProgressSnapshotSerial = 0L
|
||||
@Volatile private var nativeWorkerPaused = false
|
||||
@Volatile private var nativeWorkerNetworkPaused = false
|
||||
@Volatile private var nativeWorkerVerificationPaused = false
|
||||
@Volatile internal var nativeWorkerNetworkPaused = false
|
||||
@Volatile internal var nativeWorkerVerificationPaused = false
|
||||
@Volatile private var nativeWorkerCancelRequested = false
|
||||
private var nativeWorkerDownloadNetworkMode = "any"
|
||||
private var nativeWorkerNetworkCallback: ConnectivityManager.NetworkCallback? = null
|
||||
private val nativeWorkerWifiNetworks = mutableSetOf<Network>()
|
||||
internal var nativeWorkerDownloadNetworkMode = "any"
|
||||
internal var nativeWorkerNetworkCallback: ConnectivityManager.NetworkCallback? = null
|
||||
internal val nativeWorkerWifiNetworks = mutableSetOf<Network>()
|
||||
// Bumped every time a new native queue replaces the current one. A worker
|
||||
// coroutine that observes a different generation than its own must stop
|
||||
// without touching the snapshot or the service lifecycle: cancel() alone
|
||||
@@ -622,18 +622,18 @@ class DownloadService : Service() {
|
||||
}
|
||||
}
|
||||
|
||||
private fun isNativeWorkerPaused(): Boolean =
|
||||
internal fun isNativeWorkerPaused(): Boolean =
|
||||
nativeWorkerPaused ||
|
||||
nativeWorkerNetworkPaused ||
|
||||
nativeWorkerVerificationPaused
|
||||
|
||||
private fun nativeWorkerPauseMessage(): String = when {
|
||||
internal fun nativeWorkerPauseMessage(): String = when {
|
||||
nativeWorkerVerificationPaused -> "Verification required"
|
||||
nativeWorkerNetworkPaused -> "Waiting for Wi-Fi"
|
||||
else -> "Paused"
|
||||
}
|
||||
|
||||
private fun cancelActiveNativeItemForPause() {
|
||||
internal fun cancelActiveNativeItemForPause() {
|
||||
var itemIdToCancel = ""
|
||||
synchronized(nativeWorkerItems) {
|
||||
val activeItem = nativeWorkerItems.firstOrNull {
|
||||
@@ -659,146 +659,6 @@ class DownloadService : Service() {
|
||||
NativeDownloadFinalizer.cancelActiveWork()
|
||||
}
|
||||
|
||||
private fun configureNativeWorkerNetworkPolicy(settingsJson: String) {
|
||||
nativeWorkerDownloadNetworkMode = try {
|
||||
JSONObject(settingsJson).optString("download_network_mode", "any")
|
||||
} catch (_: Exception) {
|
||||
"any"
|
||||
}
|
||||
|
||||
unregisterNativeWorkerNetworkCallback()
|
||||
if (!NativeWorkerPolicy.requiresWifi(nativeWorkerDownloadNetworkMode)) {
|
||||
nativeWorkerNetworkPaused = false
|
||||
return
|
||||
}
|
||||
|
||||
nativeWorkerNetworkPaused = !hasUsableWifiConnection()
|
||||
val connectivityManager =
|
||||
getSystemService(Context.CONNECTIVITY_SERVICE) as ConnectivityManager
|
||||
val callback = object : ConnectivityManager.NetworkCallback() {
|
||||
override fun onAvailable(network: Network) {
|
||||
synchronized(nativeWorkerWifiNetworks) {
|
||||
nativeWorkerWifiNetworks.add(network)
|
||||
}
|
||||
refreshNativeWorkerNetworkPause()
|
||||
}
|
||||
|
||||
override fun onLost(network: Network) {
|
||||
synchronized(nativeWorkerWifiNetworks) {
|
||||
nativeWorkerWifiNetworks.remove(network)
|
||||
}
|
||||
refreshNativeWorkerNetworkPause()
|
||||
}
|
||||
|
||||
override fun onCapabilitiesChanged(
|
||||
network: Network,
|
||||
networkCapabilities: NetworkCapabilities,
|
||||
) {
|
||||
synchronized(nativeWorkerWifiNetworks) {
|
||||
if (networkCapabilities.hasTransport(NetworkCapabilities.TRANSPORT_WIFI) &&
|
||||
networkCapabilities.hasCapability(
|
||||
NetworkCapabilities.NET_CAPABILITY_INTERNET,
|
||||
)
|
||||
) {
|
||||
nativeWorkerWifiNetworks.add(network)
|
||||
} else {
|
||||
nativeWorkerWifiNetworks.remove(network)
|
||||
}
|
||||
}
|
||||
refreshNativeWorkerNetworkPause()
|
||||
}
|
||||
}
|
||||
try {
|
||||
val request = NetworkRequest.Builder()
|
||||
.addTransportType(NetworkCapabilities.TRANSPORT_WIFI)
|
||||
.addCapability(NetworkCapabilities.NET_CAPABILITY_INTERNET)
|
||||
.build()
|
||||
connectivityManager.registerNetworkCallback(request, callback)
|
||||
nativeWorkerNetworkCallback = callback
|
||||
} catch (e: Exception) {
|
||||
android.util.Log.w(
|
||||
"DownloadService",
|
||||
"Failed to monitor Wi-Fi for native worker: ${e.message}",
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
private fun hasUsableWifiConnection(): Boolean {
|
||||
val connectivityManager =
|
||||
getSystemService(Context.CONNECTIVITY_SERVICE) as ConnectivityManager
|
||||
if (synchronized(nativeWorkerWifiNetworks) {
|
||||
nativeWorkerWifiNetworks.isNotEmpty()
|
||||
}
|
||||
) {
|
||||
return true
|
||||
}
|
||||
return try {
|
||||
val activeNetwork = connectivityManager.activeNetwork ?: return false
|
||||
val capabilities =
|
||||
connectivityManager.getNetworkCapabilities(activeNetwork) ?: return false
|
||||
capabilities.hasTransport(NetworkCapabilities.TRANSPORT_WIFI) &&
|
||||
capabilities.hasCapability(NetworkCapabilities.NET_CAPABILITY_INTERNET)
|
||||
} catch (_: Exception) {
|
||||
false
|
||||
}
|
||||
}
|
||||
|
||||
private fun refreshNativeWorkerNetworkPause() {
|
||||
if (!NativeWorkerPolicy.requiresWifi(nativeWorkerDownloadNetworkMode)) return
|
||||
|
||||
val shouldPause = NativeWorkerPolicy.shouldPauseForNetwork(
|
||||
nativeWorkerDownloadNetworkMode,
|
||||
hasUsableWifiConnection(),
|
||||
)
|
||||
if (shouldPause == nativeWorkerNetworkPaused) return
|
||||
|
||||
nativeWorkerNetworkPaused = shouldPause
|
||||
if (nativeWorkerJob?.isActive != true) return
|
||||
|
||||
if (shouldPause) {
|
||||
currentStatus = if (nativeWorkerVerificationPaused) {
|
||||
"verification_required"
|
||||
} else {
|
||||
"waiting_wifi"
|
||||
}
|
||||
cancelActiveNativeItemForPause()
|
||||
updateNotification(0L, 0L)
|
||||
} else {
|
||||
currentStatus = if (nativeWorkerVerificationPaused) {
|
||||
"verification_required"
|
||||
} else {
|
||||
"preparing"
|
||||
}
|
||||
updateNotification(0L, 0L)
|
||||
}
|
||||
writeNativeWorkerSnapshotAsync(
|
||||
isRunning = nativeWorkerJob?.isActive == true,
|
||||
isPaused = isNativeWorkerPaused(),
|
||||
currentItemId = "",
|
||||
message = if (isNativeWorkerPaused()) {
|
||||
nativeWorkerPauseMessage()
|
||||
} else {
|
||||
"Wi-Fi restored"
|
||||
},
|
||||
includeItems = true,
|
||||
)
|
||||
}
|
||||
|
||||
private fun unregisterNativeWorkerNetworkCallback() {
|
||||
val callback = nativeWorkerNetworkCallback
|
||||
nativeWorkerNetworkCallback = null
|
||||
synchronized(nativeWorkerWifiNetworks) {
|
||||
nativeWorkerWifiNetworks.clear()
|
||||
}
|
||||
if (callback == null) return
|
||||
try {
|
||||
val connectivityManager =
|
||||
getSystemService(Context.CONNECTIVITY_SERVICE) as ConnectivityManager
|
||||
connectivityManager.unregisterNetworkCallback(callback)
|
||||
} catch (_: Exception) {
|
||||
}
|
||||
}
|
||||
|
||||
private fun parseNativeDownloadRequests(requestsJson: String): List<NativeDownloadRequest> {
|
||||
val array = JSONArray(requestsJson)
|
||||
val requests = ArrayList<NativeDownloadRequest>(array.length())
|
||||
@@ -1182,479 +1042,6 @@ class DownloadService : Service() {
|
||||
}
|
||||
}
|
||||
|
||||
private fun writeNativeAlbumReplayGainIfComplete(): Boolean {
|
||||
val entries = synchronized(nativeReplayGainEntries) {
|
||||
nativeReplayGainEntries.map { JSONObject(it.toString()) }
|
||||
}
|
||||
if (entries.size <= 1) return true
|
||||
|
||||
val statuses = synchronized(nativeWorkerItems) {
|
||||
nativeWorkerItems.associate { it.itemId to it.status }
|
||||
}
|
||||
val requestKeys = synchronized(nativeReplayGainRequestAlbumKeys) {
|
||||
nativeReplayGainRequestAlbumKeys.toMap()
|
||||
}
|
||||
val eligible = buildEligibleNativeAlbumReplayGain(entries, statuses, requestKeys)
|
||||
if (eligible.length() <= 1) {
|
||||
return !hasPendingNativeAlbumReplayGainWork(statuses)
|
||||
}
|
||||
return writeNativeAlbumReplayGainEntries(eligible)
|
||||
}
|
||||
|
||||
private fun buildEligibleNativeAlbumReplayGain(
|
||||
entries: List<JSONObject>,
|
||||
statuses: Map<String, String>,
|
||||
requestKeys: Map<String, String>
|
||||
): JSONArray {
|
||||
val eligibleIndexes = NativeReplayGainPolicy.eligibleEntryIndexes(
|
||||
entryAlbumKeys = entries.map { it.optString("album_key", "") },
|
||||
statuses = statuses,
|
||||
requestAlbumKeys = requestKeys,
|
||||
)
|
||||
val eligible = JSONArray()
|
||||
for (index in eligibleIndexes) {
|
||||
eligible.put(entries[index])
|
||||
}
|
||||
return eligible
|
||||
}
|
||||
|
||||
private fun writeNativeAlbumReplayGainEntries(eligible: JSONArray): Boolean {
|
||||
if (eligible.length() <= 1) return true
|
||||
try {
|
||||
val result = JSONObject(NativeDownloadFinalizer.writeAlbumReplayGain(this, eligible.toString()))
|
||||
return result.optBoolean("success", false)
|
||||
} catch (e: Exception) {
|
||||
android.util.Log.w("DownloadService", "Native album ReplayGain failed: ${e.message}")
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
private fun hasPendingNativeAlbumReplayGainWork(statuses: Map<String, String>): Boolean {
|
||||
return NativeReplayGainPolicy.hasPendingWork(statuses)
|
||||
}
|
||||
|
||||
private fun writeNativeReplayGainJournal() {
|
||||
val requestKeys = synchronized(nativeReplayGainRequestAlbumKeys) {
|
||||
nativeReplayGainRequestAlbumKeys.toMap()
|
||||
}
|
||||
if (requestKeys.isEmpty()) return
|
||||
|
||||
val entries = synchronized(nativeReplayGainEntries) {
|
||||
nativeReplayGainEntries.map { JSONObject(it.toString()) }
|
||||
}
|
||||
val statuses = synchronized(nativeWorkerItems) {
|
||||
nativeWorkerItems.associate { it.itemId to it.status }
|
||||
}
|
||||
synchronized(NATIVE_REPLAYGAIN_JOURNAL_FILE_LOCK) {
|
||||
val file = AtomicFile(File(filesDir, NATIVE_REPLAYGAIN_JOURNAL_FILE))
|
||||
val existing = readNativeReplayGainJournalLocked(file)
|
||||
val mergedEntries = mergeNativeReplayGainJournalEntries(
|
||||
existing?.optJSONArray("entries"),
|
||||
entries,
|
||||
)
|
||||
val mergedRequestKeys = mergeJsonObjectStringMap(
|
||||
existing?.optJSONObject("request_album_keys"),
|
||||
requestKeys,
|
||||
)
|
||||
val mergedStatuses = mergeJsonObjectStringMap(
|
||||
existing?.optJSONObject("statuses"),
|
||||
statuses,
|
||||
)
|
||||
val root = JSONObject()
|
||||
.put("run_id", nativeWorkerRunId)
|
||||
.put("updated_at", System.currentTimeMillis())
|
||||
.put("entries", mergedEntries)
|
||||
.put("request_album_keys", JSONObject(mergedRequestKeys))
|
||||
.put("statuses", JSONObject(mergedStatuses))
|
||||
|
||||
var stream: java.io.FileOutputStream? = null
|
||||
try {
|
||||
stream = file.startWrite()
|
||||
stream.write(root.toString().toByteArray(Charsets.UTF_8))
|
||||
file.finishWrite(stream)
|
||||
stream = null
|
||||
} catch (e: Exception) {
|
||||
android.util.Log.w("DownloadService", "Failed to write native ReplayGain journal: ${e.message}")
|
||||
} finally {
|
||||
if (stream != null) {
|
||||
file.failWrite(stream)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private fun readNativeReplayGainJournalLocked(file: AtomicFile): JSONObject? {
|
||||
return try {
|
||||
if (!file.baseFile.exists()) return null
|
||||
val text = file.openRead().bufferedReader(Charsets.UTF_8).use {
|
||||
it.readText()
|
||||
}
|
||||
JSONObject(text)
|
||||
} catch (e: Exception) {
|
||||
android.util.Log.w("DownloadService", "Failed to merge native ReplayGain journal: ${e.message}")
|
||||
null
|
||||
}
|
||||
}
|
||||
|
||||
private fun mergeNativeReplayGainJournalEntries(
|
||||
existingEntries: JSONArray?,
|
||||
currentEntries: List<JSONObject>
|
||||
): JSONArray {
|
||||
val byKey = linkedMapOf<String, JSONObject>()
|
||||
|
||||
fun add(entry: JSONObject) {
|
||||
val trackId = entry.optString("track_id", "")
|
||||
val path = entry.optString("file_path", "")
|
||||
val key = if (trackId.isNotBlank()) {
|
||||
"track:$trackId"
|
||||
} else {
|
||||
"path:$path"
|
||||
}
|
||||
if (key != "path:") {
|
||||
byKey[key] = JSONObject(entry.toString())
|
||||
}
|
||||
}
|
||||
|
||||
if (existingEntries != null) {
|
||||
for (index in 0 until existingEntries.length()) {
|
||||
existingEntries.optJSONObject(index)?.let(::add)
|
||||
}
|
||||
}
|
||||
for (entry in currentEntries) add(entry)
|
||||
|
||||
return JSONArray().apply {
|
||||
for (entry in byKey.values) put(entry)
|
||||
}
|
||||
}
|
||||
|
||||
private fun mergeJsonObjectStringMap(
|
||||
existing: JSONObject?,
|
||||
current: Map<String, String>
|
||||
): Map<String, String> {
|
||||
val merged = linkedMapOf<String, String>()
|
||||
if (existing != null) {
|
||||
for (key in existing.keys()) {
|
||||
merged[key] = existing.optString(key, "")
|
||||
}
|
||||
}
|
||||
for ((key, value) in current) {
|
||||
merged[key] = value
|
||||
}
|
||||
return merged
|
||||
}
|
||||
|
||||
private fun clearNativeReplayGainJournal() {
|
||||
synchronized(NATIVE_REPLAYGAIN_JOURNAL_FILE_LOCK) {
|
||||
try {
|
||||
AtomicFile(File(filesDir, NATIVE_REPLAYGAIN_JOURNAL_FILE)).delete()
|
||||
} catch (_: Exception) {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private fun flushNativeAlbumReplayGainJournalIfComplete() {
|
||||
val root = synchronized(NATIVE_REPLAYGAIN_JOURNAL_FILE_LOCK) {
|
||||
try {
|
||||
val file = File(filesDir, NATIVE_REPLAYGAIN_JOURNAL_FILE)
|
||||
if (!file.exists()) return
|
||||
val text = AtomicFile(file).openRead().bufferedReader(Charsets.UTF_8).use {
|
||||
it.readText()
|
||||
}
|
||||
JSONObject(text)
|
||||
} catch (e: Exception) {
|
||||
android.util.Log.w("DownloadService", "Failed to read native ReplayGain journal: ${e.message}")
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
val entriesArray = root.optJSONArray("entries") ?: return
|
||||
val entries = mutableListOf<JSONObject>()
|
||||
for (index in 0 until entriesArray.length()) {
|
||||
entriesArray.optJSONObject(index)?.let { entries.add(JSONObject(it.toString())) }
|
||||
}
|
||||
val statusesJson = root.optJSONObject("statuses") ?: JSONObject()
|
||||
val statuses = mutableMapOf<String, String>()
|
||||
for (key in statusesJson.keys()) {
|
||||
statuses[key] = statusesJson.optString(key, "")
|
||||
}
|
||||
val requestKeysJson = root.optJSONObject("request_album_keys") ?: JSONObject()
|
||||
val requestKeys = mutableMapOf<String, String>()
|
||||
for (key in requestKeysJson.keys()) {
|
||||
requestKeys[key] = requestKeysJson.optString(key, "")
|
||||
}
|
||||
|
||||
val eligible = buildEligibleNativeAlbumReplayGain(entries, statuses, requestKeys)
|
||||
if (eligible.length() <= 1 && hasPendingNativeAlbumReplayGainWork(statuses)) {
|
||||
return
|
||||
}
|
||||
if (writeNativeAlbumReplayGainEntries(eligible)) {
|
||||
clearNativeReplayGainJournal()
|
||||
}
|
||||
}
|
||||
|
||||
private fun writeNativeWorkerSnapshot(
|
||||
isRunning: Boolean,
|
||||
isPaused: Boolean,
|
||||
currentItemId: String,
|
||||
message: String,
|
||||
lastResult: JSONObject? = null,
|
||||
settingsJson: String = "",
|
||||
includeItems: Boolean = false,
|
||||
snapshotSerial: Long = snapshotWriteSerial.incrementAndGet()
|
||||
) {
|
||||
try {
|
||||
synchronized(snapshotWriteLock) {
|
||||
if (includeItems) {
|
||||
if (snapshotSerial < latestCommittedStateSnapshotSerial) return
|
||||
} else {
|
||||
if (snapshotSerial < latestCommittedProgressSnapshotSerial) return
|
||||
}
|
||||
|
||||
val counts = nativeWorkerCounts()
|
||||
val snapshot = JSONObject()
|
||||
.put("contract_version", NATIVE_WORKER_CONTRACT_VERSION)
|
||||
.put("run_id", nativeWorkerRunId.ifBlank { readNativeWorkerRunIdFromSnapshotFile() })
|
||||
.put("is_running", isRunning)
|
||||
.put("is_paused", isPaused)
|
||||
.put("total", counts.total)
|
||||
.put("completed", counts.completed)
|
||||
.put("failed", counts.failed)
|
||||
.put("skipped", counts.skipped)
|
||||
.put("current_item_id", currentItemId)
|
||||
.put("message", message)
|
||||
.put("updated_at", System.currentTimeMillis())
|
||||
.put("snapshot_serial", snapshotSerial)
|
||||
.put("state_serial", if (includeItems) snapshotSerial else latestCommittedStateSnapshotSerial)
|
||||
.put("snapshot_mode", if (includeItems) "compact_items" else "delta")
|
||||
// Snapshot of the header before the per-item payload is
|
||||
// attached; served to pollers that already consumed this
|
||||
// items payload (see getNativeWorkerSnapshot).
|
||||
val headerCandidate = if (includeItems) snapshot.toString() else null
|
||||
snapshot.put("item_ids", nativeWorkerItemIds())
|
||||
if (includeItems) {
|
||||
snapshot.put("items", nativeWorkerItemsSnapshot(includeStatic = false))
|
||||
} else {
|
||||
nativeWorkerItemSnapshot(currentItemId, includeStatic = false)?.let {
|
||||
snapshot.put("item_delta", it)
|
||||
}
|
||||
}
|
||||
if (settingsJson.isNotBlank() && includeItems) {
|
||||
snapshot.put("settings_json", settingsJson)
|
||||
}
|
||||
if (lastResult != null) {
|
||||
snapshot.put("last_result", lastResult)
|
||||
}
|
||||
|
||||
synchronized(NATIVE_WORKER_STATE_FILE_LOCK) {
|
||||
val targetFileName = if (includeItems) {
|
||||
NATIVE_WORKER_STATE_FILE
|
||||
} else {
|
||||
NATIVE_WORKER_PROGRESS_FILE
|
||||
}
|
||||
val file = AtomicFile(File(filesDir, targetFileName))
|
||||
var stream: java.io.FileOutputStream? = null
|
||||
try {
|
||||
stream = file.startWrite()
|
||||
stream.write(snapshot.toString().toByteArray(Charsets.UTF_8))
|
||||
file.finishWrite(stream)
|
||||
stream = null
|
||||
if (includeItems) {
|
||||
latestCommittedStateSnapshotSerial = snapshotSerial
|
||||
if (headerCandidate != null) {
|
||||
lastStateHeaderJson = headerCandidate
|
||||
lastStateHeaderSerial = snapshotSerial
|
||||
}
|
||||
} else {
|
||||
latestCommittedProgressSnapshotSerial = snapshotSerial
|
||||
}
|
||||
} finally {
|
||||
if (stream != null) {
|
||||
file.failWrite(stream)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
} catch (e: Exception) {
|
||||
android.util.Log.w("DownloadService", "Failed to write native worker snapshot: ${e.message}")
|
||||
}
|
||||
}
|
||||
|
||||
private fun writeNativeWorkerSnapshotAsync(
|
||||
isRunning: Boolean,
|
||||
isPaused: Boolean,
|
||||
currentItemId: String,
|
||||
message: String,
|
||||
lastResult: JSONObject? = null,
|
||||
settingsJson: String = "",
|
||||
includeItems: Boolean = false
|
||||
) {
|
||||
val snapshotSerial = snapshotWriteSerial.incrementAndGet()
|
||||
serviceScope.launch {
|
||||
writeNativeWorkerSnapshot(
|
||||
isRunning = isRunning,
|
||||
isPaused = isPaused,
|
||||
currentItemId = currentItemId,
|
||||
message = message,
|
||||
lastResult = lastResult,
|
||||
settingsJson = settingsJson,
|
||||
includeItems = includeItems,
|
||||
snapshotSerial = snapshotSerial
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
private fun readNativeWorkerRunIdFromSnapshotFile(): String {
|
||||
return try {
|
||||
synchronized(NATIVE_WORKER_STATE_FILE_LOCK) {
|
||||
val file = File(filesDir, NATIVE_WORKER_STATE_FILE)
|
||||
if (!file.exists()) {
|
||||
""
|
||||
} else {
|
||||
val text = AtomicFile(file).openRead().bufferedReader(Charsets.UTF_8).use {
|
||||
it.readText()
|
||||
}
|
||||
JSONObject(text).optString("run_id", "")
|
||||
}
|
||||
}
|
||||
} catch (_: Exception) {
|
||||
""
|
||||
}
|
||||
}
|
||||
|
||||
private fun updateNativeWorkerItem(itemId: String, updater: (NativeWorkerItem) -> Unit) {
|
||||
synchronized(nativeWorkerItems) {
|
||||
nativeWorkerItems.firstOrNull { it.itemId == itemId }?.let(updater)
|
||||
}
|
||||
}
|
||||
|
||||
private fun updateNativeWorkerItemProgress(itemId: String) {
|
||||
try {
|
||||
val raw = Gobackend.getAllDownloadProgress()
|
||||
val root = JSONObject(raw)
|
||||
val items = root.optJSONObject("items") ?: return
|
||||
val progress = items.optJSONObject(itemId) ?: return
|
||||
val backendStatus = progress.optString("status", "downloading")
|
||||
val bytesReceived = progress.optLong("bytes_received", 0L)
|
||||
val bytesTotal = progress.optLong("bytes_total", 0L)
|
||||
if (backendStatus == "preparing") {
|
||||
currentStatus = "preparing"
|
||||
updateNativeWorkerItem(itemId) {
|
||||
it.status = "preparing"
|
||||
it.progress = 0.0
|
||||
it.bytesReceived = 0L
|
||||
it.bytesTotal = 0L
|
||||
}
|
||||
lastProgress = 0L
|
||||
lastTotal = 0L
|
||||
updateNotification(0L, 0L)
|
||||
return
|
||||
}
|
||||
val progressValue = if (bytesTotal > 0L) {
|
||||
bytesReceived.toDouble() / bytesTotal.toDouble()
|
||||
} else {
|
||||
progress.optDouble("progress", 0.0)
|
||||
}.coerceIn(0.0, 1.0)
|
||||
currentStatus = if (backendStatus == "finalizing") {
|
||||
"finalizing"
|
||||
} else {
|
||||
"downloading"
|
||||
}
|
||||
updateNativeWorkerItem(itemId) {
|
||||
it.status = currentStatus
|
||||
it.progress = progressValue
|
||||
it.bytesReceived = bytesReceived
|
||||
it.bytesTotal = bytesTotal
|
||||
}
|
||||
if (bytesTotal > 0L) {
|
||||
lastProgress = bytesReceived
|
||||
lastTotal = bytesTotal
|
||||
updateNotification(bytesReceived, bytesTotal)
|
||||
} else if (progressValue > 0.0) {
|
||||
val percentProgress = (progressValue * NOTIFICATION_PERCENT_TOTAL).toLong()
|
||||
.coerceIn(0L, NOTIFICATION_PERCENT_TOTAL)
|
||||
lastProgress = percentProgress
|
||||
lastTotal = NOTIFICATION_PERCENT_TOTAL
|
||||
updateNotification(percentProgress, NOTIFICATION_PERCENT_TOTAL)
|
||||
} else {
|
||||
lastProgress = 0L
|
||||
lastTotal = 0L
|
||||
updateNotification(0L, 0L)
|
||||
}
|
||||
} catch (_: Exception) {
|
||||
}
|
||||
}
|
||||
|
||||
private fun nativeWorkerCounts(): NativeWorkerCounts {
|
||||
var total = 0
|
||||
var completed = 0
|
||||
var failed = 0
|
||||
var skipped = 0
|
||||
synchronized(nativeWorkerItems) {
|
||||
total = nativeWorkerItems.size
|
||||
for (item in nativeWorkerItems) {
|
||||
when (item.status) {
|
||||
"completed" -> completed++
|
||||
"failed" -> failed++
|
||||
"skipped" -> skipped++
|
||||
}
|
||||
}
|
||||
}
|
||||
return NativeWorkerCounts(
|
||||
total = total,
|
||||
completed = completed,
|
||||
failed = failed,
|
||||
skipped = skipped
|
||||
)
|
||||
}
|
||||
|
||||
private fun nativeWorkerItemSnapshot(itemId: String, includeStatic: Boolean): JSONObject? {
|
||||
if (itemId.isBlank()) return null
|
||||
synchronized(nativeWorkerItems) {
|
||||
val item = nativeWorkerItems.firstOrNull { it.itemId == itemId } ?: return null
|
||||
return nativeWorkerItemSnapshotLocked(item, includeStatic)
|
||||
}
|
||||
}
|
||||
|
||||
private fun nativeWorkerItemIds(): JSONArray {
|
||||
val array = JSONArray()
|
||||
synchronized(nativeWorkerItems) {
|
||||
for (item in nativeWorkerItems) {
|
||||
array.put(item.itemId)
|
||||
}
|
||||
}
|
||||
return array
|
||||
}
|
||||
|
||||
private fun nativeWorkerItemsSnapshot(includeStatic: Boolean): JSONArray {
|
||||
val array = JSONArray()
|
||||
synchronized(nativeWorkerItems) {
|
||||
for (item in nativeWorkerItems) {
|
||||
array.put(nativeWorkerItemSnapshotLocked(item, includeStatic))
|
||||
}
|
||||
}
|
||||
return array
|
||||
}
|
||||
|
||||
private fun nativeWorkerItemSnapshotLocked(item: NativeWorkerItem, includeStatic: Boolean): JSONObject {
|
||||
val json = JSONObject()
|
||||
.put("item_id", item.itemId)
|
||||
.put("status", item.status)
|
||||
.put("progress", item.progress)
|
||||
.put("bytes_received", item.bytesReceived)
|
||||
.put("bytes_total", item.bytesTotal)
|
||||
if (includeStatic) {
|
||||
json.put("track_name", item.trackName)
|
||||
.put("artist_name", item.artistName)
|
||||
.put("item_json", item.itemJson)
|
||||
}
|
||||
if (item.error.isNotBlank()) {
|
||||
json.put("error", item.error)
|
||||
}
|
||||
item.resultJson?.let { json.put("result", it) }
|
||||
return json
|
||||
}
|
||||
|
||||
@Synchronized
|
||||
private fun ensureWakeLock() {
|
||||
val existingWakeLock = wakeLock
|
||||
if (existingWakeLock?.isHeld == true) {
|
||||
@@ -1731,7 +1118,7 @@ class DownloadService : Service() {
|
||||
}
|
||||
}
|
||||
|
||||
private fun updateNotification(progress: Long, total: Long) {
|
||||
internal fun updateNotification(progress: Long, total: Long) {
|
||||
if (!isRunning) return
|
||||
ensureWakeLock()
|
||||
|
||||
|
||||
@@ -0,0 +1,174 @@
|
||||
package com.zarz.spotiflac
|
||||
|
||||
import android.app.Notification
|
||||
import android.app.NotificationChannel
|
||||
import android.app.NotificationManager
|
||||
import android.app.PendingIntent
|
||||
import android.app.Service
|
||||
import android.content.Context
|
||||
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
|
||||
import android.util.AtomicFile
|
||||
import androidx.core.app.NotificationCompat
|
||||
import gobackend.Gobackend
|
||||
import kotlinx.coroutines.CancellationException
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.Job
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.cancel
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.launch
|
||||
import org.json.JSONArray
|
||||
import org.json.JSONObject
|
||||
import java.io.File
|
||||
import java.util.concurrent.atomic.AtomicLong
|
||||
|
||||
// Wi-Fi-only download policy: network callback, pause and resume.
|
||||
|
||||
internal fun DownloadService.configureNativeWorkerNetworkPolicy(settingsJson: String) {
|
||||
nativeWorkerDownloadNetworkMode = try {
|
||||
JSONObject(settingsJson).optString("download_network_mode", "any")
|
||||
} catch (_: Exception) {
|
||||
"any"
|
||||
}
|
||||
|
||||
unregisterNativeWorkerNetworkCallback()
|
||||
if (!NativeWorkerPolicy.requiresWifi(nativeWorkerDownloadNetworkMode)) {
|
||||
nativeWorkerNetworkPaused = false
|
||||
return
|
||||
}
|
||||
|
||||
nativeWorkerNetworkPaused = !hasUsableWifiConnection()
|
||||
val connectivityManager =
|
||||
getSystemService(Context.CONNECTIVITY_SERVICE) as ConnectivityManager
|
||||
val callback = object : ConnectivityManager.NetworkCallback() {
|
||||
override fun onAvailable(network: Network) {
|
||||
synchronized(nativeWorkerWifiNetworks) {
|
||||
nativeWorkerWifiNetworks.add(network)
|
||||
}
|
||||
refreshNativeWorkerNetworkPause()
|
||||
}
|
||||
|
||||
override fun onLost(network: Network) {
|
||||
synchronized(nativeWorkerWifiNetworks) {
|
||||
nativeWorkerWifiNetworks.remove(network)
|
||||
}
|
||||
refreshNativeWorkerNetworkPause()
|
||||
}
|
||||
|
||||
override fun onCapabilitiesChanged(
|
||||
network: Network,
|
||||
networkCapabilities: NetworkCapabilities,
|
||||
) {
|
||||
synchronized(nativeWorkerWifiNetworks) {
|
||||
if (networkCapabilities.hasTransport(NetworkCapabilities.TRANSPORT_WIFI) &&
|
||||
networkCapabilities.hasCapability(
|
||||
NetworkCapabilities.NET_CAPABILITY_INTERNET,
|
||||
)
|
||||
) {
|
||||
nativeWorkerWifiNetworks.add(network)
|
||||
} else {
|
||||
nativeWorkerWifiNetworks.remove(network)
|
||||
}
|
||||
}
|
||||
refreshNativeWorkerNetworkPause()
|
||||
}
|
||||
}
|
||||
try {
|
||||
val request = NetworkRequest.Builder()
|
||||
.addTransportType(NetworkCapabilities.TRANSPORT_WIFI)
|
||||
.addCapability(NetworkCapabilities.NET_CAPABILITY_INTERNET)
|
||||
.build()
|
||||
connectivityManager.registerNetworkCallback(request, callback)
|
||||
nativeWorkerNetworkCallback = callback
|
||||
} catch (e: Exception) {
|
||||
android.util.Log.w(
|
||||
"DownloadService",
|
||||
"Failed to monitor Wi-Fi for native worker: ${e.message}",
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
internal fun DownloadService.hasUsableWifiConnection(): Boolean {
|
||||
val connectivityManager =
|
||||
getSystemService(Context.CONNECTIVITY_SERVICE) as ConnectivityManager
|
||||
if (synchronized(nativeWorkerWifiNetworks) {
|
||||
nativeWorkerWifiNetworks.isNotEmpty()
|
||||
}
|
||||
) {
|
||||
return true
|
||||
}
|
||||
return try {
|
||||
val activeNetwork = connectivityManager.activeNetwork ?: return false
|
||||
val capabilities =
|
||||
connectivityManager.getNetworkCapabilities(activeNetwork) ?: return false
|
||||
capabilities.hasTransport(NetworkCapabilities.TRANSPORT_WIFI) &&
|
||||
capabilities.hasCapability(NetworkCapabilities.NET_CAPABILITY_INTERNET)
|
||||
} catch (_: Exception) {
|
||||
false
|
||||
}
|
||||
}
|
||||
|
||||
internal fun DownloadService.refreshNativeWorkerNetworkPause() {
|
||||
if (!NativeWorkerPolicy.requiresWifi(nativeWorkerDownloadNetworkMode)) return
|
||||
|
||||
val shouldPause = NativeWorkerPolicy.shouldPauseForNetwork(
|
||||
nativeWorkerDownloadNetworkMode,
|
||||
hasUsableWifiConnection(),
|
||||
)
|
||||
if (shouldPause == nativeWorkerNetworkPaused) return
|
||||
|
||||
nativeWorkerNetworkPaused = shouldPause
|
||||
if (nativeWorkerJob?.isActive != true) return
|
||||
|
||||
if (shouldPause) {
|
||||
currentStatus = if (nativeWorkerVerificationPaused) {
|
||||
"verification_required"
|
||||
} else {
|
||||
"waiting_wifi"
|
||||
}
|
||||
cancelActiveNativeItemForPause()
|
||||
updateNotification(0L, 0L)
|
||||
} else {
|
||||
currentStatus = if (nativeWorkerVerificationPaused) {
|
||||
"verification_required"
|
||||
} else {
|
||||
"preparing"
|
||||
}
|
||||
updateNotification(0L, 0L)
|
||||
}
|
||||
writeNativeWorkerSnapshotAsync(
|
||||
isRunning = nativeWorkerJob?.isActive == true,
|
||||
isPaused = isNativeWorkerPaused(),
|
||||
currentItemId = "",
|
||||
message = if (isNativeWorkerPaused()) {
|
||||
nativeWorkerPauseMessage()
|
||||
} else {
|
||||
"Wi-Fi restored"
|
||||
},
|
||||
includeItems = true,
|
||||
)
|
||||
}
|
||||
|
||||
internal fun DownloadService.unregisterNativeWorkerNetworkCallback() {
|
||||
val callback = nativeWorkerNetworkCallback
|
||||
nativeWorkerNetworkCallback = null
|
||||
synchronized(nativeWorkerWifiNetworks) {
|
||||
nativeWorkerWifiNetworks.clear()
|
||||
}
|
||||
if (callback == null) return
|
||||
try {
|
||||
val connectivityManager =
|
||||
getSystemService(Context.CONNECTIVITY_SERVICE) as ConnectivityManager
|
||||
connectivityManager.unregisterNetworkCallback(callback)
|
||||
} catch (_: Exception) {
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,244 @@
|
||||
package com.zarz.spotiflac
|
||||
|
||||
import android.app.Notification
|
||||
import android.app.NotificationChannel
|
||||
import android.app.NotificationManager
|
||||
import android.app.PendingIntent
|
||||
import android.app.Service
|
||||
import android.content.Context
|
||||
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
|
||||
import android.util.AtomicFile
|
||||
import androidx.core.app.NotificationCompat
|
||||
import gobackend.Gobackend
|
||||
import kotlinx.coroutines.CancellationException
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.Job
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.cancel
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.launch
|
||||
import org.json.JSONArray
|
||||
import org.json.JSONObject
|
||||
import java.io.File
|
||||
import java.util.concurrent.atomic.AtomicLong
|
||||
|
||||
// Album ReplayGain accumulation and journal persistence for the native worker.
|
||||
|
||||
internal fun DownloadService.writeNativeAlbumReplayGainIfComplete(): Boolean {
|
||||
val entries = synchronized(nativeReplayGainEntries) {
|
||||
nativeReplayGainEntries.map { JSONObject(it.toString()) }
|
||||
}
|
||||
if (entries.size <= 1) return true
|
||||
|
||||
val statuses = synchronized(nativeWorkerItems) {
|
||||
nativeWorkerItems.associate { it.itemId to it.status }
|
||||
}
|
||||
val requestKeys = synchronized(nativeReplayGainRequestAlbumKeys) {
|
||||
nativeReplayGainRequestAlbumKeys.toMap()
|
||||
}
|
||||
val eligible = buildEligibleNativeAlbumReplayGain(entries, statuses, requestKeys)
|
||||
if (eligible.length() <= 1) {
|
||||
return !hasPendingNativeAlbumReplayGainWork(statuses)
|
||||
}
|
||||
return writeNativeAlbumReplayGainEntries(eligible)
|
||||
}
|
||||
|
||||
internal fun DownloadService.buildEligibleNativeAlbumReplayGain(
|
||||
entries: List<JSONObject>,
|
||||
statuses: Map<String, String>,
|
||||
requestKeys: Map<String, String>
|
||||
): JSONArray {
|
||||
val eligibleIndexes = NativeReplayGainPolicy.eligibleEntryIndexes(
|
||||
entryAlbumKeys = entries.map { it.optString("album_key", "") },
|
||||
statuses = statuses,
|
||||
requestAlbumKeys = requestKeys,
|
||||
)
|
||||
val eligible = JSONArray()
|
||||
for (index in eligibleIndexes) {
|
||||
eligible.put(entries[index])
|
||||
}
|
||||
return eligible
|
||||
}
|
||||
|
||||
internal fun DownloadService.writeNativeAlbumReplayGainEntries(eligible: JSONArray): Boolean {
|
||||
if (eligible.length() <= 1) return true
|
||||
try {
|
||||
val result = JSONObject(NativeDownloadFinalizer.writeAlbumReplayGain(this, eligible.toString()))
|
||||
return result.optBoolean("success", false)
|
||||
} catch (e: Exception) {
|
||||
android.util.Log.w("DownloadService", "Native album ReplayGain failed: ${e.message}")
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
internal fun DownloadService.hasPendingNativeAlbumReplayGainWork(statuses: Map<String, String>): Boolean {
|
||||
return NativeReplayGainPolicy.hasPendingWork(statuses)
|
||||
}
|
||||
|
||||
internal fun DownloadService.writeNativeReplayGainJournal() {
|
||||
val requestKeys = synchronized(nativeReplayGainRequestAlbumKeys) {
|
||||
nativeReplayGainRequestAlbumKeys.toMap()
|
||||
}
|
||||
if (requestKeys.isEmpty()) return
|
||||
|
||||
val entries = synchronized(nativeReplayGainEntries) {
|
||||
nativeReplayGainEntries.map { JSONObject(it.toString()) }
|
||||
}
|
||||
val statuses = synchronized(nativeWorkerItems) {
|
||||
nativeWorkerItems.associate { it.itemId to it.status }
|
||||
}
|
||||
synchronized(DownloadService.NATIVE_REPLAYGAIN_JOURNAL_FILE_LOCK) {
|
||||
val file = AtomicFile(File(filesDir, DownloadService.NATIVE_REPLAYGAIN_JOURNAL_FILE))
|
||||
val existing = readNativeReplayGainJournalLocked(file)
|
||||
val mergedEntries = mergeNativeReplayGainJournalEntries(
|
||||
existing?.optJSONArray("entries"),
|
||||
entries,
|
||||
)
|
||||
val mergedRequestKeys = mergeJsonObjectStringMap(
|
||||
existing?.optJSONObject("request_album_keys"),
|
||||
requestKeys,
|
||||
)
|
||||
val mergedStatuses = mergeJsonObjectStringMap(
|
||||
existing?.optJSONObject("statuses"),
|
||||
statuses,
|
||||
)
|
||||
val root = JSONObject()
|
||||
.put("run_id", nativeWorkerRunId)
|
||||
.put("updated_at", System.currentTimeMillis())
|
||||
.put("entries", mergedEntries)
|
||||
.put("request_album_keys", JSONObject(mergedRequestKeys))
|
||||
.put("statuses", JSONObject(mergedStatuses))
|
||||
|
||||
var stream: java.io.FileOutputStream? = null
|
||||
try {
|
||||
stream = file.startWrite()
|
||||
stream.write(root.toString().toByteArray(Charsets.UTF_8))
|
||||
file.finishWrite(stream)
|
||||
stream = null
|
||||
} catch (e: Exception) {
|
||||
android.util.Log.w("DownloadService", "Failed to write native ReplayGain journal: ${e.message}")
|
||||
} finally {
|
||||
if (stream != null) {
|
||||
file.failWrite(stream)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
internal fun DownloadService.readNativeReplayGainJournalLocked(file: AtomicFile): JSONObject? {
|
||||
return try {
|
||||
if (!file.baseFile.exists()) return null
|
||||
val text = file.openRead().bufferedReader(Charsets.UTF_8).use {
|
||||
it.readText()
|
||||
}
|
||||
JSONObject(text)
|
||||
} catch (e: Exception) {
|
||||
android.util.Log.w("DownloadService", "Failed to merge native ReplayGain journal: ${e.message}")
|
||||
null
|
||||
}
|
||||
}
|
||||
|
||||
internal fun DownloadService.mergeNativeReplayGainJournalEntries(
|
||||
existingEntries: JSONArray?,
|
||||
currentEntries: List<JSONObject>
|
||||
): JSONArray {
|
||||
val byKey = linkedMapOf<String, JSONObject>()
|
||||
|
||||
fun add(entry: JSONObject) {
|
||||
val trackId = entry.optString("track_id", "")
|
||||
val path = entry.optString("file_path", "")
|
||||
val key = if (trackId.isNotBlank()) {
|
||||
"track:$trackId"
|
||||
} else {
|
||||
"path:$path"
|
||||
}
|
||||
if (key != "path:") {
|
||||
byKey[key] = JSONObject(entry.toString())
|
||||
}
|
||||
}
|
||||
|
||||
if (existingEntries != null) {
|
||||
for (index in 0 until existingEntries.length()) {
|
||||
existingEntries.optJSONObject(index)?.let(::add)
|
||||
}
|
||||
}
|
||||
for (entry in currentEntries) add(entry)
|
||||
|
||||
return JSONArray().apply {
|
||||
for (entry in byKey.values) put(entry)
|
||||
}
|
||||
}
|
||||
|
||||
internal fun DownloadService.mergeJsonObjectStringMap(
|
||||
existing: JSONObject?,
|
||||
current: Map<String, String>
|
||||
): Map<String, String> {
|
||||
val merged = linkedMapOf<String, String>()
|
||||
if (existing != null) {
|
||||
for (key in existing.keys()) {
|
||||
merged[key] = existing.optString(key, "")
|
||||
}
|
||||
}
|
||||
for ((key, value) in current) {
|
||||
merged[key] = value
|
||||
}
|
||||
return merged
|
||||
}
|
||||
|
||||
internal fun DownloadService.clearNativeReplayGainJournal() {
|
||||
synchronized(DownloadService.NATIVE_REPLAYGAIN_JOURNAL_FILE_LOCK) {
|
||||
try {
|
||||
AtomicFile(File(filesDir, DownloadService.NATIVE_REPLAYGAIN_JOURNAL_FILE)).delete()
|
||||
} catch (_: Exception) {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
internal fun DownloadService.flushNativeAlbumReplayGainJournalIfComplete() {
|
||||
val root = synchronized(DownloadService.NATIVE_REPLAYGAIN_JOURNAL_FILE_LOCK) {
|
||||
try {
|
||||
val file = File(filesDir, DownloadService.NATIVE_REPLAYGAIN_JOURNAL_FILE)
|
||||
if (!file.exists()) return
|
||||
val text = AtomicFile(file).openRead().bufferedReader(Charsets.UTF_8).use {
|
||||
it.readText()
|
||||
}
|
||||
JSONObject(text)
|
||||
} catch (e: Exception) {
|
||||
android.util.Log.w("DownloadService", "Failed to read native ReplayGain journal: ${e.message}")
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
val entriesArray = root.optJSONArray("entries") ?: return
|
||||
val entries = mutableListOf<JSONObject>()
|
||||
for (index in 0 until entriesArray.length()) {
|
||||
entriesArray.optJSONObject(index)?.let { entries.add(JSONObject(it.toString())) }
|
||||
}
|
||||
val statusesJson = root.optJSONObject("statuses") ?: JSONObject()
|
||||
val statuses = mutableMapOf<String, String>()
|
||||
for (key in statusesJson.keys()) {
|
||||
statuses[key] = statusesJson.optString(key, "")
|
||||
}
|
||||
val requestKeysJson = root.optJSONObject("request_album_keys") ?: JSONObject()
|
||||
val requestKeys = mutableMapOf<String, String>()
|
||||
for (key in requestKeysJson.keys()) {
|
||||
requestKeys[key] = requestKeysJson.optString(key, "")
|
||||
}
|
||||
|
||||
val eligible = buildEligibleNativeAlbumReplayGain(entries, statuses, requestKeys)
|
||||
if (eligible.length() <= 1 && hasPendingNativeAlbumReplayGainWork(statuses)) {
|
||||
return
|
||||
}
|
||||
if (writeNativeAlbumReplayGainEntries(eligible)) {
|
||||
clearNativeReplayGainJournal()
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,297 @@
|
||||
package com.zarz.spotiflac
|
||||
|
||||
import android.app.Notification
|
||||
import android.app.NotificationChannel
|
||||
import android.app.NotificationManager
|
||||
import android.app.PendingIntent
|
||||
import android.app.Service
|
||||
import android.content.Context
|
||||
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
|
||||
import android.util.AtomicFile
|
||||
import androidx.core.app.NotificationCompat
|
||||
import gobackend.Gobackend
|
||||
import kotlinx.coroutines.CancellationException
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.Job
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.cancel
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.launch
|
||||
import org.json.JSONArray
|
||||
import org.json.JSONObject
|
||||
import java.io.File
|
||||
import java.util.concurrent.atomic.AtomicLong
|
||||
|
||||
// Native-worker item state snapshots for the Flutter side.
|
||||
|
||||
internal fun DownloadService.writeNativeWorkerSnapshot(
|
||||
isRunning: Boolean,
|
||||
isPaused: Boolean,
|
||||
currentItemId: String,
|
||||
message: String,
|
||||
lastResult: JSONObject? = null,
|
||||
settingsJson: String = "",
|
||||
includeItems: Boolean = false,
|
||||
snapshotSerial: Long = snapshotWriteSerial.incrementAndGet()
|
||||
) {
|
||||
try {
|
||||
synchronized(snapshotWriteLock) {
|
||||
if (includeItems) {
|
||||
if (snapshotSerial < latestCommittedStateSnapshotSerial) return
|
||||
} else {
|
||||
if (snapshotSerial < latestCommittedProgressSnapshotSerial) return
|
||||
}
|
||||
|
||||
val counts = nativeWorkerCounts()
|
||||
val snapshot = JSONObject()
|
||||
.put("contract_version", DownloadService.NATIVE_WORKER_CONTRACT_VERSION)
|
||||
.put("run_id", nativeWorkerRunId.ifBlank { readNativeWorkerRunIdFromSnapshotFile() })
|
||||
.put("is_running", isRunning)
|
||||
.put("is_paused", isPaused)
|
||||
.put("total", counts.total)
|
||||
.put("completed", counts.completed)
|
||||
.put("failed", counts.failed)
|
||||
.put("skipped", counts.skipped)
|
||||
.put("current_item_id", currentItemId)
|
||||
.put("message", message)
|
||||
.put("updated_at", System.currentTimeMillis())
|
||||
.put("snapshot_serial", snapshotSerial)
|
||||
.put("state_serial", if (includeItems) snapshotSerial else latestCommittedStateSnapshotSerial)
|
||||
.put("snapshot_mode", if (includeItems) "compact_items" else "delta")
|
||||
// Snapshot of the header before the per-item payload is
|
||||
// attached; served to pollers that already consumed this
|
||||
// items payload (see getNativeWorkerSnapshot).
|
||||
val headerCandidate = if (includeItems) snapshot.toString() else null
|
||||
snapshot.put("item_ids", nativeWorkerItemIds())
|
||||
if (includeItems) {
|
||||
snapshot.put("items", nativeWorkerItemsSnapshot(includeStatic = false))
|
||||
} else {
|
||||
nativeWorkerItemSnapshot(currentItemId, includeStatic = false)?.let {
|
||||
snapshot.put("item_delta", it)
|
||||
}
|
||||
}
|
||||
if (settingsJson.isNotBlank() && includeItems) {
|
||||
snapshot.put("settings_json", settingsJson)
|
||||
}
|
||||
if (lastResult != null) {
|
||||
snapshot.put("last_result", lastResult)
|
||||
}
|
||||
|
||||
synchronized(DownloadService.NATIVE_WORKER_STATE_FILE_LOCK) {
|
||||
val targetFileName = if (includeItems) {
|
||||
DownloadService.NATIVE_WORKER_STATE_FILE
|
||||
} else {
|
||||
DownloadService.NATIVE_WORKER_PROGRESS_FILE
|
||||
}
|
||||
val file = AtomicFile(File(filesDir, targetFileName))
|
||||
var stream: java.io.FileOutputStream? = null
|
||||
try {
|
||||
stream = file.startWrite()
|
||||
stream.write(snapshot.toString().toByteArray(Charsets.UTF_8))
|
||||
file.finishWrite(stream)
|
||||
stream = null
|
||||
if (includeItems) {
|
||||
latestCommittedStateSnapshotSerial = snapshotSerial
|
||||
if (headerCandidate != null) {
|
||||
DownloadService.lastStateHeaderJson = headerCandidate
|
||||
DownloadService.lastStateHeaderSerial = snapshotSerial
|
||||
}
|
||||
} else {
|
||||
latestCommittedProgressSnapshotSerial = snapshotSerial
|
||||
}
|
||||
} finally {
|
||||
if (stream != null) {
|
||||
file.failWrite(stream)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
} catch (e: Exception) {
|
||||
android.util.Log.w("DownloadService", "Failed to write native worker snapshot: ${e.message}")
|
||||
}
|
||||
}
|
||||
|
||||
internal fun DownloadService.writeNativeWorkerSnapshotAsync(
|
||||
isRunning: Boolean,
|
||||
isPaused: Boolean,
|
||||
currentItemId: String,
|
||||
message: String,
|
||||
lastResult: JSONObject? = null,
|
||||
settingsJson: String = "",
|
||||
includeItems: Boolean = false
|
||||
) {
|
||||
val snapshotSerial = snapshotWriteSerial.incrementAndGet()
|
||||
serviceScope.launch {
|
||||
writeNativeWorkerSnapshot(
|
||||
isRunning = isRunning,
|
||||
isPaused = isPaused,
|
||||
currentItemId = currentItemId,
|
||||
message = message,
|
||||
lastResult = lastResult,
|
||||
settingsJson = settingsJson,
|
||||
includeItems = includeItems,
|
||||
snapshotSerial = snapshotSerial
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
internal fun DownloadService.readNativeWorkerRunIdFromSnapshotFile(): String {
|
||||
return try {
|
||||
synchronized(DownloadService.NATIVE_WORKER_STATE_FILE_LOCK) {
|
||||
val file = File(filesDir, DownloadService.NATIVE_WORKER_STATE_FILE)
|
||||
if (!file.exists()) {
|
||||
""
|
||||
} else {
|
||||
val text = AtomicFile(file).openRead().bufferedReader(Charsets.UTF_8).use {
|
||||
it.readText()
|
||||
}
|
||||
JSONObject(text).optString("run_id", "")
|
||||
}
|
||||
}
|
||||
} catch (_: Exception) {
|
||||
""
|
||||
}
|
||||
}
|
||||
|
||||
internal fun DownloadService.updateNativeWorkerItem(itemId: String, updater: (DownloadService.NativeWorkerItem) -> Unit) {
|
||||
synchronized(nativeWorkerItems) {
|
||||
nativeWorkerItems.firstOrNull { it.itemId == itemId }?.let(updater)
|
||||
}
|
||||
}
|
||||
|
||||
internal fun DownloadService.updateNativeWorkerItemProgress(itemId: String) {
|
||||
try {
|
||||
val raw = Gobackend.getAllDownloadProgress()
|
||||
val root = JSONObject(raw)
|
||||
val items = root.optJSONObject("items") ?: return
|
||||
val progress = items.optJSONObject(itemId) ?: return
|
||||
val backendStatus = progress.optString("status", "downloading")
|
||||
val bytesReceived = progress.optLong("bytes_received", 0L)
|
||||
val bytesTotal = progress.optLong("bytes_total", 0L)
|
||||
if (backendStatus == "preparing") {
|
||||
currentStatus = "preparing"
|
||||
updateNativeWorkerItem(itemId) {
|
||||
it.status = "preparing"
|
||||
it.progress = 0.0
|
||||
it.bytesReceived = 0L
|
||||
it.bytesTotal = 0L
|
||||
}
|
||||
lastProgress = 0L
|
||||
lastTotal = 0L
|
||||
updateNotification(0L, 0L)
|
||||
return
|
||||
}
|
||||
val progressValue = if (bytesTotal > 0L) {
|
||||
bytesReceived.toDouble() / bytesTotal.toDouble()
|
||||
} else {
|
||||
progress.optDouble("progress", 0.0)
|
||||
}.coerceIn(0.0, 1.0)
|
||||
currentStatus = if (backendStatus == "finalizing") {
|
||||
"finalizing"
|
||||
} else {
|
||||
"downloading"
|
||||
}
|
||||
updateNativeWorkerItem(itemId) {
|
||||
it.status = currentStatus
|
||||
it.progress = progressValue
|
||||
it.bytesReceived = bytesReceived
|
||||
it.bytesTotal = bytesTotal
|
||||
}
|
||||
if (bytesTotal > 0L) {
|
||||
lastProgress = bytesReceived
|
||||
lastTotal = bytesTotal
|
||||
updateNotification(bytesReceived, bytesTotal)
|
||||
} else if (progressValue > 0.0) {
|
||||
val percentProgress = (progressValue * DownloadService.NOTIFICATION_PERCENT_TOTAL).toLong()
|
||||
.coerceIn(0L, DownloadService.NOTIFICATION_PERCENT_TOTAL)
|
||||
lastProgress = percentProgress
|
||||
lastTotal = DownloadService.NOTIFICATION_PERCENT_TOTAL
|
||||
updateNotification(percentProgress, DownloadService.NOTIFICATION_PERCENT_TOTAL)
|
||||
} else {
|
||||
lastProgress = 0L
|
||||
lastTotal = 0L
|
||||
updateNotification(0L, 0L)
|
||||
}
|
||||
} catch (_: Exception) {
|
||||
}
|
||||
}
|
||||
|
||||
internal fun DownloadService.nativeWorkerCounts(): DownloadService.NativeWorkerCounts {
|
||||
var total = 0
|
||||
var completed = 0
|
||||
var failed = 0
|
||||
var skipped = 0
|
||||
synchronized(nativeWorkerItems) {
|
||||
total = nativeWorkerItems.size
|
||||
for (item in nativeWorkerItems) {
|
||||
when (item.status) {
|
||||
"completed" -> completed++
|
||||
"failed" -> failed++
|
||||
"skipped" -> skipped++
|
||||
}
|
||||
}
|
||||
}
|
||||
return DownloadService.NativeWorkerCounts(
|
||||
total = total,
|
||||
completed = completed,
|
||||
failed = failed,
|
||||
skipped = skipped
|
||||
)
|
||||
}
|
||||
|
||||
internal fun DownloadService.nativeWorkerItemSnapshot(itemId: String, includeStatic: Boolean): JSONObject? {
|
||||
if (itemId.isBlank()) return null
|
||||
synchronized(nativeWorkerItems) {
|
||||
val item = nativeWorkerItems.firstOrNull { it.itemId == itemId } ?: return null
|
||||
return nativeWorkerItemSnapshotLocked(item, includeStatic)
|
||||
}
|
||||
}
|
||||
|
||||
internal fun DownloadService.nativeWorkerItemIds(): JSONArray {
|
||||
val array = JSONArray()
|
||||
synchronized(nativeWorkerItems) {
|
||||
for (item in nativeWorkerItems) {
|
||||
array.put(item.itemId)
|
||||
}
|
||||
}
|
||||
return array
|
||||
}
|
||||
|
||||
internal fun DownloadService.nativeWorkerItemsSnapshot(includeStatic: Boolean): JSONArray {
|
||||
val array = JSONArray()
|
||||
synchronized(nativeWorkerItems) {
|
||||
for (item in nativeWorkerItems) {
|
||||
array.put(nativeWorkerItemSnapshotLocked(item, includeStatic))
|
||||
}
|
||||
}
|
||||
return array
|
||||
}
|
||||
|
||||
internal fun DownloadService.nativeWorkerItemSnapshotLocked(item: DownloadService.NativeWorkerItem, includeStatic: Boolean): JSONObject {
|
||||
val json = JSONObject()
|
||||
.put("item_id", item.itemId)
|
||||
.put("status", item.status)
|
||||
.put("progress", item.progress)
|
||||
.put("bytes_received", item.bytesReceived)
|
||||
.put("bytes_total", item.bytesTotal)
|
||||
if (includeStatic) {
|
||||
json.put("track_name", item.trackName)
|
||||
.put("artist_name", item.artistName)
|
||||
.put("item_json", item.itemJson)
|
||||
}
|
||||
if (item.error.isNotBlank()) {
|
||||
json.put("error", item.error)
|
||||
}
|
||||
item.resultJson?.let { json.put("result", it) }
|
||||
return json
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user