perf(native-worker): stop shipping the full item payload on every poll

The 1s snapshot poll re-read and re-parsed the whole worker state file
— which retains every completed item's result (history row, possibly
full synced lyrics; 1.5-5 KB each) for the entire run — then
re-serialized it across the channel and re-iterated it on the Dart UI
isolate. Late in a 1000-track batch that is megabytes parsed several
times per second, on the Android main thread.

Pollers now echo back the state_serial of the last snapshot they fully
processed; when the state has not advanced, the native side serves a
cached compact header (no items, no results, no settings) merged with
the small progress delta — steady-state polls are O(1) regardless of
batch size, and the full payload is delivered exactly once per state
transition. The channel handler also reads and parses off the main
thread now, matching its neighbours.
This commit is contained in:
zarzet
2026-07-10 13:08:57 +07:00
parent 87b46a9694
commit 0d39002af0
4 changed files with 91 additions and 12 deletions
@@ -153,7 +153,15 @@ class DownloadService : Service() {
context.startService(intent)
}
fun getNativeWorkerSnapshot(context: Context): String {
// Header of the last written state snapshot (all fields except the
// per-item payload). Lets steady-state polls skip re-reading and
// 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
fun getNativeWorkerSnapshot(context: Context, sinceStateSerial: Long = 0L): String {
synchronized(NATIVE_WORKER_STATE_FILE_LOCK) {
val stateFile = File(context.filesDir, NATIVE_WORKER_STATE_FILE)
if (!stateFile.exists()) {
@@ -165,12 +173,35 @@ class DownloadService : Service() {
.put("completed", 0)
.put("failed", 0)
.put("skipped", 0)
.put("state_serial", 0L)
.put("items", JSONArray())
.toString()
}
val state = AtomicFile(stateFile).openRead().bufferedReader(Charsets.UTF_8).use {
it.readText()
}.let { JSONObject(it) }
val headerJson = lastStateHeaderJson
val headerSerial = lastStateHeaderSerial
val state: JSONObject
if (sinceStateSerial > 0 &&
headerSerial in 1..sinceStateSerial &&
headerJson != null
) {
// Caller already consumed this items payload; serve the
// cached compact header instead of re-parsing the file.
state = JSONObject(headerJson)
} else {
state = AtomicFile(stateFile).openRead().bufferedReader(Charsets.UTF_8).use {
it.readText()
}.let { JSONObject(it) }
val stateSerial = state.optLong("snapshot_serial", 0L)
state.put("state_serial", stateSerial)
if (sinceStateSerial > 0 && stateSerial in 1..sinceStateSerial) {
state.remove("items")
state.remove("item_ids")
state.remove("last_result")
state.remove("settings_json")
}
}
val progressFile = File(context.filesDir, NATIVE_WORKER_PROGRESS_FILE)
if (progressFile.exists()) {
try {
@@ -1128,8 +1159,13 @@ class DownloadService : Service() {
.put("message", message)
.put("updated_at", System.currentTimeMillis())
.put("snapshot_serial", snapshotSerial)
.put("item_ids", nativeWorkerItemIds())
.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 {
@@ -1159,6 +1195,10 @@ class DownloadService : Service() {
stream = null
if (includeItems) {
latestCommittedStateSnapshotSerial = snapshotSerial
if (headerCandidate != null) {
lastStateHeaderJson = headerCandidate
lastStateHeaderSerial = snapshotSerial
}
} else {
latestCommittedProgressSnapshotSerial = snapshotSerial
}
@@ -3094,7 +3094,19 @@ class MainActivity: FlutterFragmentActivity() {
result.success(null)
}
"getNativeDownloadWorkerSnapshot" -> {
result.success(parseJsonPayload(DownloadService.getNativeWorkerSnapshot(this@MainActivity)))
val sinceStateSerial =
(call.argument<Number>("since_state_serial") ?: 0L).toLong()
// The snapshot can be megabytes late in a large
// batch; read and parse it off the main thread.
val payload = withContext(Dispatchers.IO) {
parseJsonPayload(
DownloadService.getNativeWorkerSnapshot(
this@MainActivity,
sinceStateSerial
)
)
}
result.success(payload)
}
"preWarmTrackCache" -> {
val tracksJson = call.argument<String>("tracks") ?: "[]"
+23 -2
View File
@@ -3809,6 +3809,17 @@ class DownloadQueueNotifier extends Notifier<DownloadQueueState> {
return version == DownloadRequestPayload.nativeWorkerContractVersion;
}
/// Serial of the last fully-processed state snapshot, echoed back to the
/// native side so steady-state polls skip the per-item payload. Falls back
/// to the previous value for snapshots without the field.
int _snapshotStateSerial(Map<String, dynamic> snapshot, int previous) {
final serial = snapshot['state_serial'];
if (serial is num && serial > 0) {
return serial.toInt();
}
return previous;
}
bool _isNativeWorkerSnapshotForRun(
Map<String, dynamic> snapshot,
String runId,
@@ -4019,9 +4030,12 @@ class DownloadQueueNotifier extends Notifier<DownloadQueueState> {
String runId,
) async {
var deadServicePolls = 0;
var lastStateSerial = 0;
try {
while (true) {
final snapshot = await PlatformBridge.getNativeDownloadWorkerSnapshot();
final snapshot = await PlatformBridge.getNativeDownloadWorkerSnapshot(
sinceStateSerial: lastStateSerial,
);
final matchesRun = _isNativeWorkerSnapshotForRun(snapshot, runId);
if (!matchesRun || snapshot['is_running'] == true) {
if (await _isNativeWorkerServiceAlive()) {
@@ -4054,6 +4068,7 @@ class DownloadQueueNotifier extends Notifier<DownloadQueueState> {
}
}
if (!matchesRun) {
lastStateSerial = 0;
await Future<void>.delayed(const Duration(seconds: 1));
continue;
}
@@ -4068,6 +4083,7 @@ class DownloadQueueNotifier extends Notifier<DownloadQueueState> {
reconciledIds,
settings,
);
lastStateSerial = _snapshotStateSerial(snapshot, lastStateSerial);
if (snapshot['is_running'] != true) {
await _clearNativeWorkerRunId(runId);
// Items may have been requeued during reconciliation (e.g. batch
@@ -4160,9 +4176,13 @@ class DownloadQueueNotifier extends Notifier<DownloadQueueState> {
);
final runStartWait = Stopwatch()..start();
var lastStateSerial = 0;
while (true) {
final snapshot = await PlatformBridge.getNativeDownloadWorkerSnapshot();
final snapshot = await PlatformBridge.getNativeDownloadWorkerSnapshot(
sinceStateSerial: lastStateSerial,
);
if (!_isNativeWorkerSnapshotForRun(snapshot, runId)) {
lastStateSerial = 0;
if (runStartWait.elapsed > const Duration(seconds: 30)) {
throw _NativeWorkerStartupTimeout();
}
@@ -4175,6 +4195,7 @@ class DownloadQueueNotifier extends Notifier<DownloadQueueState> {
reconciledIds,
settings,
);
lastStateSerial = _snapshotStateSerial(snapshot, lastStateSerial);
if (snapshot['is_running'] != true) {
await _clearNativeWorkerRunId(runId);
break;
+10 -4
View File
@@ -1096,10 +1096,16 @@ class PlatformBridge {
await _channel.invokeMethod('cancelNativeDownloadWorker');
}
static Future<Map<String, dynamic>> getNativeDownloadWorkerSnapshot() async {
final result = await _channel.invokeMethod(
'getNativeDownloadWorkerSnapshot',
);
/// [sinceStateSerial]: pass the `state_serial` of the last fully-processed
/// snapshot; when the native state hasn't advanced past it, the heavy
/// per-item payload (which grows with every completed item) is omitted and
/// only compact progress fields are returned.
static Future<Map<String, dynamic>> getNativeDownloadWorkerSnapshot({
int sinceStateSerial = 0,
}) async {
final result = await _channel.invokeMethod('getNativeDownloadWorkerSnapshot', {
'since_state_serial': sinceStateSerial,
});
return _decodeMapResult(result);
}