diff --git a/android/app/src/main/kotlin/dev/zapstore/app/plugins/AndroidPackageManagerPlugin.kt b/android/app/src/main/kotlin/dev/zapstore/app/plugins/AndroidPackageManagerPlugin.kt index e52ba7f..f6ff0c6 100644 --- a/android/app/src/main/kotlin/dev/zapstore/app/plugins/AndroidPackageManagerPlugin.kt +++ b/android/app/src/main/kotlin/dev/zapstore/app/plugins/AndroidPackageManagerPlugin.kt @@ -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() @@ -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("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("packageName") diff --git a/lib/router.dart b/lib/router.dart index 2965a80..e4e9c01 100644 --- a/lib/router.dart +++ b/lib/router.dart @@ -200,20 +200,11 @@ final routerProvider = Provider((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 diff --git a/lib/services/log_service.dart b/lib/services/log_service.dart index 6defe3d..690740c 100644 --- a/lib/services/log_service.dart +++ b/lib/services/log_service.dart @@ -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 _startFlush() { + late final Future 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 _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.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(); } } diff --git a/lib/services/package_manager/android_package_manager.dart b/lib/services/package_manager/android_package_manager.dart index 8087e86..c51de7f 100644 --- a/lib/services/package_manager/android_package_manager.dart +++ b/lib/services/package_manager/android_package_manager.dart @@ -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? _syncInFlight; StreamSubscription? _eventSubscription; /// Tracks appIds where we've already attempted to abort orphaned sessions. @@ -713,7 +715,21 @@ final class AndroidPackageManager extends PackageManager { // ═══════════════════════════════════════════════════════════════════════════ @override - Future syncInstalledPackages() async { + Future syncInstalledPackages() { + final inFlight = _syncInFlight; + if (inFlight != null) return inFlight; + + late final Future operation; + operation = _performInstalledPackagesSync().whenComplete(() { + if (identical(_syncInFlight, operation)) { + _syncInFlight = null; + } + }); + _syncInFlight = operation; + return operation; + } + + Future _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 = {}; 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); + } } } } diff --git a/lib/services/updates_service.dart b/lib/services/updates_service.dart index f08e41e..c320d38 100644 --- a/lib/services/updates_service.dart +++ b/lib/services/updates_service.dart @@ -152,11 +152,11 @@ class UpdatePollerNotifier extends StateNotifier { // 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 checkNow() async { + Future checkNow({bool refreshInstalledPackages = true}) async { if (state.isChecking) return; if (state.lastCheckTime != null) { @@ -167,7 +167,9 @@ class UpdatePollerNotifier extends StateNotifier { 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, diff --git a/lib/widgets/latest_releases_container.dart b/lib/widgets/latest_releases_container.dart index 2403993..25bac47 100644 --- a/lib/widgets/latest_releases_container.dart +++ b/lib/widgets/latest_releases_container.dart @@ -134,7 +134,11 @@ class LatestReleasesNotifier extends StateNotifier { 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 { 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 { /// 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 _resolveAppsFromLocal( StorageNotifier storage, - List assets, - ) { - final result = {}; + List assets, { + Map 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 = { + 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( 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(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(); } } diff --git a/test/services/log_service_test.dart b/test/services/log_service_test.dart index 484423e..b481b63 100644 --- a/test/services/log_service_test.dart +++ b/test/services/log_service_test.dart @@ -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 _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.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(); diff --git a/test/services/package_manager/android_package_manager_test.dart b/test/services/package_manager/android_package_manager_test.dart new file mode 100644 index 0000000..a0ace2a --- /dev/null +++ b/test/services/package_manager/android_package_manager_test.dart @@ -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>(); + 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); + }); +}