mirror of
https://github.com/zapstore/zapstore.git
synced 2026-10-05 12:38:24 +00:00
Improve resume and navigation responsiveness by fixing log flush races, incremental catalog resolution, and coalesced Android package scans
This commit is contained in:
+34
-1
@@ -32,6 +32,8 @@ import java.security.MessageDigest
|
||||
import java.security.cert.CertificateFactory
|
||||
import java.util.concurrent.ConcurrentHashMap
|
||||
import java.util.concurrent.CountDownLatch
|
||||
import java.util.concurrent.ExecutorService
|
||||
import java.util.concurrent.Executors
|
||||
import java.util.concurrent.TimeUnit
|
||||
|
||||
private const val TAG = "AndroidPackageManager"
|
||||
@@ -110,6 +112,8 @@ class AndroidPackageManagerPlugin :
|
||||
private lateinit var eventChannel: EventChannel
|
||||
private lateinit var context: Context
|
||||
private val mainHandler = Handler(Looper.getMainLooper())
|
||||
private var packageQueryExecutor: ExecutorService? = null
|
||||
@Volatile private var isEngineAttached = false
|
||||
|
||||
/** Per-app watchdog generation counters (cancels old scheduled callbacks) */
|
||||
private val watchdogGen = mutableMapOf<String, Int>()
|
||||
@@ -337,6 +341,8 @@ class AndroidPackageManagerPlugin :
|
||||
context = binding.applicationContext
|
||||
appContext = context
|
||||
instance = this
|
||||
isEngineAttached = true
|
||||
packageQueryExecutor = Executors.newSingleThreadExecutor()
|
||||
|
||||
methodChannel = MethodChannel(binding.binaryMessenger, "android_package_manager")
|
||||
methodChannel.setMethodCallHandler(this)
|
||||
@@ -355,6 +361,9 @@ class AndroidPackageManagerPlugin :
|
||||
}
|
||||
|
||||
override fun onDetachedFromEngine(binding: FlutterPlugin.FlutterPluginBinding) {
|
||||
isEngineAttached = false
|
||||
packageQueryExecutor?.shutdownNow()
|
||||
packageQueryExecutor = null
|
||||
mainHandler.post { ProcessLifecycleOwner.get().lifecycle.removeObserver(this) }
|
||||
unregisterSessionCallback()
|
||||
methodChannel.setMethodCallHandler(null)
|
||||
@@ -680,7 +689,31 @@ class AndroidPackageManagerPlugin :
|
||||
"requestInstallPermission" -> requestInstallPermission(result)
|
||||
"getInstalledApps" -> {
|
||||
val includeSystem = call.argument<Boolean>("includeSystemApps") ?: false
|
||||
result.success(getInstalledApps(includeSystem))
|
||||
val executor = packageQueryExecutor
|
||||
if (executor == null) {
|
||||
result.error("PLUGIN_DETACHED", "Package manager is unavailable", null)
|
||||
return
|
||||
}
|
||||
executor.execute {
|
||||
try {
|
||||
val installedApps = getInstalledApps(includeSystem)
|
||||
mainHandler.post {
|
||||
if (isEngineAttached) {
|
||||
result.success(installedApps)
|
||||
}
|
||||
}
|
||||
} catch (error: Throwable) {
|
||||
mainHandler.post {
|
||||
if (isEngineAttached) {
|
||||
result.error(
|
||||
"PACKAGE_QUERY_FAILED",
|
||||
error.message ?: "Could not query installed apps",
|
||||
null
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
"uninstall" -> {
|
||||
val packageName = call.argument<String>("packageName")
|
||||
|
||||
+1
-10
@@ -200,20 +200,11 @@ final routerProvider = Provider<GoRouter>((ref) {
|
||||
previousPath = currentPath;
|
||||
|
||||
Future.microtask(() {
|
||||
// Sync installed packages on every navigation to catch sideloads,
|
||||
// external installs/uninstalls, and self-updating apps.
|
||||
// This is a local-only platform channel call (~100-500ms, no network).
|
||||
unawaited(
|
||||
ref.read(packageManagerProvider.notifier).syncInstalledPackages(),
|
||||
);
|
||||
|
||||
// Re-derive catalog from local DB when arriving at updates tab so
|
||||
// data written by other screens or the background service is visible
|
||||
// without waiting for the next poll cycle.
|
||||
if (isUpdatesRoute && !wasUpdatesRoute) {
|
||||
unawaited(
|
||||
ref.read(updatePollerProvider.notifier).refreshFromLocal(),
|
||||
);
|
||||
unawaited(ref.read(updatePollerProvider.notifier).refreshFromLocal());
|
||||
}
|
||||
|
||||
// Clear completed operations when navigating away from updates
|
||||
|
||||
@@ -375,11 +375,9 @@ class LogService {
|
||||
await f;
|
||||
continue;
|
||||
}
|
||||
if (_pending.isEmpty) return;
|
||||
// No flush in flight but entries pending — kick one off and
|
||||
// record it as the active future so concurrent callers join.
|
||||
_flushFuture = _runFlush();
|
||||
await _flushFuture;
|
||||
if (_pending.isEmpty || _diskDisabled) return;
|
||||
// No flush in flight but entries pending — kick one off.
|
||||
await _startFlush();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -449,21 +447,47 @@ class LogService {
|
||||
_flushScheduled = false;
|
||||
// Another caller (e.g. `flush()`) may have started the flush
|
||||
// already; if so, do nothing — entries added after that will
|
||||
// re-schedule in the `_runFlush` finally block.
|
||||
// re-schedule when the in-flight flush completes.
|
||||
if (_flushFuture != null) return;
|
||||
_flushFuture = _runFlush();
|
||||
_startFlush();
|
||||
});
|
||||
}
|
||||
|
||||
/// Sole owner of [_flushFuture]: assigned exactly once here and cleared
|
||||
/// exactly once on completion. `_runFlush` must never touch it — it used
|
||||
/// to null it synchronously on an early return, after which the caller's
|
||||
/// `_flushFuture = _runFlush()` assignment resurrected the already
|
||||
/// completed future. That stale non-null future made `flush()` spin its
|
||||
/// `await`-loop on the microtask queue forever, permanently starving the
|
||||
/// UI event loop (app frozen on resume, kill required).
|
||||
///
|
||||
/// `whenComplete` on an async function's future always fires in a later
|
||||
/// microtask, so the clear below is guaranteed to run after the
|
||||
/// assignment — even when `_runFlush` returns without suspending.
|
||||
Future<void> _startFlush() {
|
||||
late final Future<void> future;
|
||||
future = _runFlush().whenComplete(() {
|
||||
if (identical(_flushFuture, future)) {
|
||||
_flushFuture = null;
|
||||
}
|
||||
// If new entries arrived while we were flushing, schedule again.
|
||||
if (_pending.isNotEmpty && !_diskDisabled) _scheduleFlush();
|
||||
});
|
||||
_flushFuture = future;
|
||||
return future;
|
||||
}
|
||||
|
||||
Future<void> _runFlush() async {
|
||||
if (_pending.isEmpty || _diskDisabled) {
|
||||
_flushFuture = null;
|
||||
if (_diskDisabled) {
|
||||
// Entries can never reach disk this session; drop them so callers
|
||||
// of flush() terminate. The ring buffer still holds them.
|
||||
_pending.clear();
|
||||
return;
|
||||
}
|
||||
if (_pending.isEmpty) return;
|
||||
final file = _activeFile;
|
||||
if (file == null) {
|
||||
_pending.clear();
|
||||
_flushFuture = null;
|
||||
return;
|
||||
}
|
||||
final batch = List<LogEntry>.from(_pending);
|
||||
@@ -475,10 +499,6 @@ class LogService {
|
||||
_onDiskWriteFailure(e);
|
||||
} catch (_) {
|
||||
// Swallow — logging must never crash the app.
|
||||
} finally {
|
||||
_flushFuture = null;
|
||||
// If new entries arrived while we were flushing, schedule again.
|
||||
if (_pending.isNotEmpty) _scheduleFlush();
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
import 'dart:async';
|
||||
import 'dart:io';
|
||||
|
||||
import 'package:flutter/foundation.dart';
|
||||
import 'package:flutter/services.dart';
|
||||
import 'package:models/models.dart';
|
||||
import 'package:zapstore/services/c1_proof_verification.dart';
|
||||
@@ -75,6 +76,7 @@ final class AndroidPackageManager extends PackageManager {
|
||||
static const _eventChannel = EventChannel('android_package_manager/events');
|
||||
bool _supportsSilentInstall = false;
|
||||
int _syncGeneration = 0;
|
||||
Future<void>? _syncInFlight;
|
||||
StreamSubscription<dynamic>? _eventSubscription;
|
||||
|
||||
/// Tracks appIds where we've already attempted to abort orphaned sessions.
|
||||
@@ -713,7 +715,21 @@ final class AndroidPackageManager extends PackageManager {
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
|
||||
@override
|
||||
Future<void> syncInstalledPackages() async {
|
||||
Future<void> syncInstalledPackages() {
|
||||
final inFlight = _syncInFlight;
|
||||
if (inFlight != null) return inFlight;
|
||||
|
||||
late final Future<void> operation;
|
||||
operation = _performInstalledPackagesSync().whenComplete(() {
|
||||
if (identical(_syncInFlight, operation)) {
|
||||
_syncInFlight = null;
|
||||
}
|
||||
});
|
||||
_syncInFlight = operation;
|
||||
return operation;
|
||||
}
|
||||
|
||||
Future<void> _performInstalledPackagesSync() async {
|
||||
final syncGen = ++_syncGeneration;
|
||||
state = state.copyWith(isScanning: true);
|
||||
try {
|
||||
@@ -729,7 +745,7 @@ final class AndroidPackageManager extends PackageManager {
|
||||
) ??
|
||||
[];
|
||||
|
||||
if (syncGen != _syncGeneration) return;
|
||||
if (!mounted || syncGen != _syncGeneration) return;
|
||||
|
||||
final packages = <String, PackageInfo>{};
|
||||
var anyCanInstallSilently = false;
|
||||
@@ -762,7 +778,7 @@ final class AndroidPackageManager extends PackageManager {
|
||||
}
|
||||
}
|
||||
|
||||
if (syncGen != _syncGeneration) return;
|
||||
if (!mounted || syncGen != _syncGeneration) return;
|
||||
|
||||
// Derive general silent install capability from per-app data
|
||||
_supportsSilentInstall = anyCanInstallSilently;
|
||||
@@ -770,8 +786,10 @@ final class AndroidPackageManager extends PackageManager {
|
||||
// Trust the native source as the single source of truth.
|
||||
// Android's SUCCESS broadcast only fires after the package is committed,
|
||||
// so there's no race condition with queryable state.
|
||||
state = state.copyWith(installed: packages);
|
||||
await InstalledPackagesSnapshot.save(state.installed);
|
||||
if (!mapEquals(state.installed, packages)) {
|
||||
state = state.copyWith(installed: packages);
|
||||
await InstalledPackagesSnapshot.save(packages);
|
||||
}
|
||||
|
||||
// Clear operations for apps where the installed version matches the target version
|
||||
// This catches installs that succeeded but we missed the event
|
||||
@@ -827,7 +845,9 @@ final class AndroidPackageManager extends PackageManager {
|
||||
} catch (_) {
|
||||
// Don't clobber state on transient errors
|
||||
} finally {
|
||||
state = state.copyWith(isScanning: false);
|
||||
if (mounted) {
|
||||
state = state.copyWith(isScanning: false);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -152,11 +152,11 @@ class UpdatePollerNotifier extends StateNotifier<UpdatePollerState> {
|
||||
// Fire-and-forget: a remote check blocks for tens of seconds offline,
|
||||
// but must not delay the local-first UI render that [refreshFromLocal]
|
||||
// already produced.
|
||||
unawaited(checkNow());
|
||||
unawaited(checkNow(refreshInstalledPackages: false));
|
||||
}
|
||||
|
||||
/// Trigger an update check. Called by timer and pull-to-refresh.
|
||||
Future<void> checkNow() async {
|
||||
Future<void> checkNow({bool refreshInstalledPackages = true}) async {
|
||||
if (state.isChecking) return;
|
||||
|
||||
if (state.lastCheckTime != null) {
|
||||
@@ -167,7 +167,9 @@ class UpdatePollerNotifier extends StateNotifier<UpdatePollerState> {
|
||||
state = state.copyWith(isChecking: true);
|
||||
|
||||
try {
|
||||
await ref.read(packageManagerProvider.notifier).syncInstalledPackages();
|
||||
if (refreshInstalledPackages) {
|
||||
await ref.read(packageManagerProvider.notifier).syncInstalledPackages();
|
||||
}
|
||||
await _fetchCatalog();
|
||||
state = state.copyWith(
|
||||
isChecking: false,
|
||||
|
||||
@@ -134,7 +134,11 @@ class LatestReleasesNotifier extends StateNotifier<LatestReleasesState> {
|
||||
|
||||
final keepAssets = [...assets, ...filteredOlder];
|
||||
|
||||
final appsByAssetId = _resolveAppsFromLocal(storage, keepAssets);
|
||||
final appsByAssetId = _resolveAppsFromLocal(
|
||||
storage,
|
||||
keepAssets,
|
||||
existing: state.appsByAssetId,
|
||||
);
|
||||
|
||||
state = state.copyWith(
|
||||
firstPage: assets,
|
||||
@@ -224,9 +228,17 @@ class LatestReleasesNotifier extends StateNotifier<LatestReleasesState> {
|
||||
|
||||
state = state.copyWith(
|
||||
olderPages: [...state.olderPages, ...unique],
|
||||
appsByAssetId: _resolveAppsFromLocal(storage, keepAssets),
|
||||
appsByAssetId: _resolveAppsFromLocal(
|
||||
storage,
|
||||
keepAssets,
|
||||
existing: state.appsByAssetId,
|
||||
),
|
||||
isLoadingMore: false,
|
||||
hasMore: assets.length >= _kPageSize,
|
||||
// A page that yields no NEW assets is end-of-list. Mirrors the
|
||||
// shared PagedSubscriptionNotifier contract — without the
|
||||
// `unique.isNotEmpty` half, a page of all-duplicates keeps
|
||||
// `hasMore` true and spins the viewport auto-fill loop.
|
||||
hasMore: unique.isNotEmpty && assets.length >= _kPageSize,
|
||||
);
|
||||
} catch (_) {
|
||||
state = state.copyWith(isLoadingMore: false);
|
||||
@@ -235,13 +247,30 @@ class LatestReleasesNotifier extends StateNotifier<LatestReleasesState> {
|
||||
|
||||
/// Synchronously map asset IDs -> parent App using local storage only.
|
||||
/// Never awaits network; returns an empty entry for ids not yet in cache.
|
||||
///
|
||||
/// INCREMENTAL: reuses already-resolved entries from [existing] and only
|
||||
/// runs a `querySync` for assets not yet resolved. This is critical for
|
||||
/// UI responsiveness — this method runs on the UI isolate on EVERY
|
||||
/// subscription emission, and emissions arrive in bursts on resume
|
||||
/// (relay reconnect re-streams the catalog). Re-resolving the entire
|
||||
/// accumulated asset list on every emission is O(loaded assets) of
|
||||
/// synchronous DB work per emission, which pegs the isolate and freezes
|
||||
/// the app. Resolving only new assets keeps steady-state emissions ~O(0).
|
||||
Map<String, App> _resolveAppsFromLocal(
|
||||
StorageNotifier storage,
|
||||
List<SoftwareAsset> assets,
|
||||
) {
|
||||
final result = <String, App>{};
|
||||
List<SoftwareAsset> assets, {
|
||||
Map<String, App> existing = const {},
|
||||
}) {
|
||||
final assetIds = assets.map((asset) => asset.id).toSet();
|
||||
// Carry over previously resolved entries whose asset is still present;
|
||||
// drop entries for assets no longer in the list.
|
||||
final result = <String, App>{
|
||||
for (final entry in existing.entries)
|
||||
if (assetIds.contains(entry.key)) entry.key: entry.value,
|
||||
};
|
||||
for (final asset in assets) {
|
||||
if (asset.appIdentifier.isEmpty) continue;
|
||||
if (result.containsKey(asset.id)) continue; // already resolved
|
||||
final matches = storage.querySync(
|
||||
RequestFilter<App>(
|
||||
authors: {asset.event.pubkey},
|
||||
@@ -339,6 +368,13 @@ class LatestReleasesContainer extends HookConsumerWidget {
|
||||
|
||||
final combinedApps = [...pinnedApps, ...dedupedApps];
|
||||
|
||||
// Tracks the visible (deduped) app count the last time we auto-filled an
|
||||
// underfilled viewport. Guards against an unbounded loop: when many
|
||||
// SoftwareAssets collapse to few app cards, `maxScrollExtent` stays <= 0
|
||||
// and the post-frame check would otherwise call `loadMore` every frame,
|
||||
// walking the entire local asset history synchronously on the UI isolate.
|
||||
final lastAutoFillLength = useRef<int?>(null);
|
||||
|
||||
useEffect(() {
|
||||
if (state == null) return null;
|
||||
void onScroll() {
|
||||
@@ -356,9 +392,23 @@ class LatestReleasesContainer extends HookConsumerWidget {
|
||||
final position = scrollController.position;
|
||||
final s = ref.read(latestReleasesProvider);
|
||||
if (s.isLoadingMore || !s.hasMore) return;
|
||||
if (position.maxScrollExtent <= 0 ||
|
||||
|
||||
// User scrolled near the bottom of a scrollable list — normal
|
||||
// infinite scroll. Only possible once content overflows the viewport.
|
||||
if (position.maxScrollExtent > 0 &&
|
||||
position.pixels >= position.maxScrollExtent - 300) {
|
||||
ref.read(latestReleasesProvider.notifier).loadMore();
|
||||
return;
|
||||
}
|
||||
|
||||
// Viewport not filled yet — auto-fill so infinite scroll can engage.
|
||||
// Only attempt again once the visible app count has actually grown;
|
||||
// otherwise a page of assets that maps only to already-shown apps
|
||||
// would keep this firing forever.
|
||||
if (position.maxScrollExtent <= 0) {
|
||||
if (lastAutoFillLength.value == combinedApps.length) return;
|
||||
lastAutoFillLength.value = combinedApps.length;
|
||||
ref.read(latestReleasesProvider.notifier).loadMore();
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
import 'dart:io';
|
||||
import 'dart:isolate';
|
||||
|
||||
import 'package:flutter_test/flutter_test.dart';
|
||||
import 'package:path/path.dart' as p;
|
||||
@@ -18,6 +19,18 @@ Future<({LogService log, Directory dir})> _newService(
|
||||
|
||||
File _activeFile(Directory dir) => File(p.join(dir.path, 'zapstore.log'));
|
||||
|
||||
/// Isolate entry for the flushSync-race regression test. Reports 'ok'
|
||||
/// once flush() returns; if flush() spins forever the parent times out.
|
||||
Future<void> _flushAfterFlushSyncRace(SendPort report) async {
|
||||
final dir = await Directory.systemTemp.createTemp('log_service_race_');
|
||||
final svc = LogService.forTesting(directory: dir, isolate: 'race');
|
||||
svc.info('entry'); // schedules the flush microtask
|
||||
svc.flushSync(); // drains _pending before the microtask runs
|
||||
await Future<void>.delayed(Duration.zero); // let the microtask fire
|
||||
await svc.flush(); // must terminate
|
||||
report.send('ok');
|
||||
}
|
||||
|
||||
void main() {
|
||||
group('LogService basics', () {
|
||||
test('records to ring buffer and disk', () async {
|
||||
@@ -288,6 +301,37 @@ void main() {
|
||||
});
|
||||
|
||||
group('Concurrency', () {
|
||||
// Regression: log() schedules a flush microtask; if flushSync() (the
|
||||
// uncaught-error path) drains _pending before that microtask runs,
|
||||
// the microtask used to leave a stale, already-completed
|
||||
// _flushFuture behind. The next flush() — called on every app
|
||||
// pause — then spun `await`ing that completed future forever,
|
||||
// starving the UI event loop with microtasks (app frozen on
|
||||
// resume, kill required).
|
||||
//
|
||||
// A regression cannot be caught by a plain `await flush()` here:
|
||||
// the spin never yields to the event loop, so even a timeout timer
|
||||
// would never fire in this isolate. Run the scenario in a separate
|
||||
// isolate and time out from the outside.
|
||||
test('flush completes after flushSync races a scheduled flush',
|
||||
() async {
|
||||
final done = ReceivePort();
|
||||
final isolate = await Isolate.spawn(
|
||||
_flushAfterFlushSyncRace,
|
||||
done.sendPort,
|
||||
);
|
||||
try {
|
||||
final result = await done.first.timeout(
|
||||
const Duration(seconds: 10),
|
||||
onTimeout: () => 'flush() never completed — microtask spin',
|
||||
);
|
||||
expect(result, 'ok');
|
||||
} finally {
|
||||
isolate.kill(priority: Isolate.immediate);
|
||||
done.close();
|
||||
}
|
||||
});
|
||||
|
||||
test('many overlapping flushes never corrupt a line', () async {
|
||||
final (:log, :dir) = await _newService();
|
||||
|
||||
|
||||
@@ -0,0 +1,54 @@
|
||||
import 'dart:async';
|
||||
|
||||
import 'package:flutter/services.dart';
|
||||
import 'package:flutter_test/flutter_test.dart';
|
||||
import 'package:hooks_riverpod/hooks_riverpod.dart';
|
||||
import 'package:zapstore/services/package_manager/android_package_manager.dart';
|
||||
import 'package:zapstore/services/package_manager/package_manager.dart';
|
||||
|
||||
void main() {
|
||||
TestWidgetsFlutterBinding.ensureInitialized();
|
||||
|
||||
const methodChannel = MethodChannel('android_package_manager');
|
||||
const eventChannel = MethodChannel('android_package_manager/events');
|
||||
|
||||
tearDown(() async {
|
||||
final messenger =
|
||||
TestDefaultBinaryMessengerBinding.instance.defaultBinaryMessenger;
|
||||
messenger.setMockMethodCallHandler(methodChannel, null);
|
||||
messenger.setMockMethodCallHandler(eventChannel, null);
|
||||
});
|
||||
|
||||
test('coalesces overlapping installed-package scans', () async {
|
||||
final scanResult = Completer<List<Object?>>();
|
||||
var scanCalls = 0;
|
||||
final messenger =
|
||||
TestDefaultBinaryMessengerBinding.instance.defaultBinaryMessenger;
|
||||
|
||||
messenger.setMockMethodCallHandler(eventChannel, (_) async => null);
|
||||
messenger.setMockMethodCallHandler(methodChannel, (call) async {
|
||||
expect(call.method, 'getInstalledApps');
|
||||
scanCalls++;
|
||||
return scanResult.future;
|
||||
});
|
||||
|
||||
final container = ProviderContainer(
|
||||
overrides: [
|
||||
packageManagerProvider.overrideWith(AndroidPackageManager.new),
|
||||
],
|
||||
);
|
||||
addTearDown(container.dispose);
|
||||
|
||||
final packageManager = container.read(packageManagerProvider.notifier);
|
||||
final first = packageManager.syncInstalledPackages();
|
||||
final second = packageManager.syncInstalledPackages();
|
||||
|
||||
expect(identical(first, second), isTrue);
|
||||
expect(scanCalls, 1);
|
||||
|
||||
scanResult.complete(const []);
|
||||
await Future.wait([first, second]);
|
||||
|
||||
expect(scanCalls, 1);
|
||||
});
|
||||
}
|
||||
Reference in New Issue
Block a user