From 7652d060e2d65634225930c649d0a769d5a55c4e Mon Sep 17 00:00:00 2001 From: edde746 <86283021+edde746@users.noreply.github.com> Date: Fri, 21 Aug 2026 13:39:49 +0200 Subject: [PATCH] fix(trackers): purge queued writes on session invalidation and lock-order the fallback sweep Queued tracker writes still replayed through the wrong account when the session died by token expiry instead of explicit disconnect: the auth-failure teardown cleared the store and rebound null but never purged the service's retry queue, so rows queued under account A replayed through whichever account connected next. The invalidation callback now purges like the disconnect path, after the rebind so an in-flight failure is dropped by the account-binding check instead of re-queueing behind the purge. Narrower second hole in the same invariant: removeService swept the in-memory fallback outside the queue lock, so an enqueue whose persist failed while a disconnect raced it re-buffered the row after the sweep and the next flush resurrected it. The fallback add and the fallback sweep both run inside the queue lock now, so lock-slot order covers the buffered store the same way it already covered the persisted one. --- lib/providers/trackers_provider.dart | 16 ++++-- .../trackers/tracker_coordinator.dart | 7 +-- .../trackers/tracker_write_queue.dart | 46 +++++++++------- test/providers/trackers_provider_test.dart | 52 ++++++++++++++++++ .../trackers/tracker_write_queue_test.dart | 54 +++++++++++++++++++ 5 files changed, 148 insertions(+), 27 deletions(-) diff --git a/lib/providers/trackers_provider.dart b/lib/providers/trackers_provider.dart index d3a0752cb..13ce5a4e8 100644 --- a/lib/providers/trackers_provider.dart +++ b/lib/providers/trackers_provider.dart @@ -342,10 +342,11 @@ class TrackersProvider extends ChangeNotifier with DisposableChangeNotifierMixin } /// Explicit-disconnect teardown. Every caller is one of the `disconnectX` - /// methods — profile rebinds go through [onActiveProfileChanged] instead — - /// so this is also where the service's queued writes are purged: they were - /// created under the account being dropped and must never replay through - /// whichever account connects to this service next. + /// methods — profile rebinds go through [onActiveProfileChanged] instead. + /// The service's queued writes are purged here: they were created under the + /// account being dropped and must never replay through whichever account + /// connects to this service next. Session invalidation tears down through + /// [_rebind]'s `onInvalidated`, which purges for the same reason. Future _clearAndRebind(_TrackerSlot slot) async { _invalidateConnect(slot.service); final userUuid = _activeUserUuid; @@ -429,6 +430,13 @@ class TrackersProvider extends ChangeNotifier with DisposableChangeNotifierMixin slot.store.clear(boundUuid); slot.session = null; _rebind(slot); + // Same contract as the explicit-disconnect purge in [_clearAndRebind]: + // rows queued under the session that just died must never replay + // through whichever account connects to this service next. The rebind + // above moved the account binding first, so an in-flight write that + // fails after this point is dropped instead of re-queued behind the + // purge. + unawaited(TrackerCoordinator.instance.purgeWriteQueueForService(slot.service)); safeNotifyListeners(); }, onUpdated: (next) { diff --git a/lib/services/trackers/tracker_coordinator.dart b/lib/services/trackers/tracker_coordinator.dart index 7680bd932..e4843efa5 100644 --- a/lib/services/trackers/tracker_coordinator.dart +++ b/lib/services/trackers/tracker_coordinator.dart @@ -153,9 +153,10 @@ class TrackerCoordinator { /// Drop [service]'s queued writes for the active profile. /// - /// Called on explicit disconnect, before another account can bind: a queued - /// row created under the departing account would otherwise be replayed - /// through its successor, silently editing the wrong account's history. + /// Called on explicit disconnect and on session invalidation, before another + /// account can bind: a queued row created under the departing account would + /// otherwise be replayed through its successor, silently editing the wrong + /// account's history. /// Best-effort like [flushWriteQueue]: a purge that cannot be persisted is /// logged and the rows simply keep waiting. Future purgeWriteQueueForService(TrackerService service) async { diff --git a/lib/services/trackers/tracker_write_queue.dart b/lib/services/trackers/tracker_write_queue.dart index 51e4ee60a..5183e1961 100644 --- a/lib/services/trackers/tracker_write_queue.dart +++ b/lib/services/trackers/tracker_write_queue.dart @@ -213,9 +213,13 @@ class TrackerWriteQueue { /// Persist [item] as the surviving intent for its coalesce key; fall back to a /// bounded in-memory buffer when the disk write throws. Buffered items are /// retried at the start of the next [flush]. - Future enqueue(String userUuid, TrackerWriteQueueItem item) async { - try { - await _locked(() async { + /// + /// The fallback add runs inside the queue lock so it is ordered against + /// [removeService]: a purge whose slot is claimed after this enqueue's is + /// guaranteed to also sweep a row that could only be buffered, not persisted. + Future enqueue(String userUuid, TrackerWriteQueueItem item) { + return _locked(() async { + try { final items = await load(userUuid); final claim = item.progressClaim; if (claim != null && @@ -227,20 +231,20 @@ class TrackerWriteQueue { items.removeWhere((queued) => queued.coalesceKey == item.coalesceKey); items.add(item); await _save(userUuid, items); - }); - } catch (e, st) { - appLogger.e( - 'Tracker write queue: persist failed for ${item.service.name} ${item.ctx.ratingKey}, buffering in memory', - error: e, - stackTrace: st, - ); - final fallback = _inMemoryFallbackByUser.putIfAbsent(userUuid, Queue.new); - if (fallback.length >= _maxInMemoryFallback) { - final dropped = fallback.removeFirst(); - appLogger.w('Tracker write queue: in-memory buffer full, dropping ${dropped.service.name}'); + } catch (e, st) { + appLogger.e( + 'Tracker write queue: persist failed for ${item.service.name} ${item.ctx.ratingKey}, buffering in memory', + error: e, + stackTrace: st, + ); + final fallback = _inMemoryFallbackByUser.putIfAbsent(userUuid, Queue.new); + if (fallback.length >= _maxInMemoryFallback) { + final dropped = fallback.removeFirst(); + appLogger.w('Tracker write queue: in-memory buffer full, dropping ${dropped.service.name}'); + } + fallback.addLast(item); } - fallback.addLast(item); - } + }); } /// Drop queued writes a completed direct write has superseded. @@ -267,12 +271,14 @@ class TrackerWriteQueue { /// Drop every queued row for [service] under [userUuid] — persisted and /// in-memory fallback alike. /// - /// Called on explicit disconnect: items carry no tracker-account identity, - /// so a row queued under the departing account would otherwise replay - /// through whichever account connects to this service next. + /// Called on explicit disconnect and on session invalidation: items carry no + /// tracker-account identity, so a row queued under the departing account + /// would otherwise replay through whichever account connects to this service + /// next. The fallback sweep runs inside the lock so it is ordered after any + /// racing [enqueue] whose failed persist buffered its row. Future removeService(String userUuid, TrackerService service) async { - _inMemoryFallbackByUser[userUuid]?.removeWhere((item) => item.service == service); await _locked(() async { + _inMemoryFallbackByUser[userUuid]?.removeWhere((item) => item.service == service); final items = await load(userUuid); final before = items.length; items.removeWhere((item) => item.service == service); diff --git a/test/providers/trackers_provider_test.dart b/test/providers/trackers_provider_test.dart index 58b7cbd2d..9e5a74100 100644 --- a/test/providers/trackers_provider_test.dart +++ b/test/providers/trackers_provider_test.dart @@ -1,6 +1,7 @@ import 'dart:async'; import 'package:flutter_test/flutter_test.dart'; +import 'package:plezy/models/trackers/tracker_context.dart'; import 'package:plezy/providers/trackers_provider.dart'; import 'package:plezy/services/base_shared_preferences_service.dart'; import 'package:plezy/services/trackers/anilist/anilist_tracker.dart'; @@ -8,10 +9,12 @@ import 'package:plezy/services/trackers/tracker_account_store.dart'; import 'package:plezy/services/trackers/tracker_constants.dart'; import 'package:plezy/services/trackers/tracker_coordinator.dart'; import 'package:plezy/services/trackers/tracker_session.dart'; +import 'package:plezy/services/trackers/tracker_write_queue.dart'; import 'package:plezy/services/trackers/mal/mal_tracker.dart'; import 'package:plezy/services/trackers/mdblist/mdblist_tracker.dart'; import 'package:plezy/services/trackers/simkl/simkl_tracker.dart'; import 'package:plezy/services/trackers/trakt/trakt_tracker.dart'; +import 'package:plezy/utils/external_ids.dart'; import '../test_helpers/io_fakes.dart'; import '../test_helpers/prefs.dart'; @@ -56,6 +59,24 @@ Future _bindProfile(TrackersProvider provider, String? userUuid) async { await TrackerCoordinator.instance.flushWriteQueue(); } +/// A queued watched write for [service], as a failed live write would leave it. +TrackerWriteQueueItem _queuedItem(TrackerService service) { + final ctx = TrackerContext.movie( + external: const ExternalIds(tmdb: 456), + anime: null, + ratingKey: 'movie-1', + libraryGlobalKey: 'server-1:8', + ); + return TrackerWriteQueueItem( + service: service, + watched: true, + ctx: ctx, + coalesceKey: trackerItemCoalesceKey(service, ctx, trackerExternalRowIdentity(ctx.external))!, + progressClaim: null, + watchedAtIso: '2026-05-12T00:00:00.000Z', + ); +} + void main() { setUp(() { resetSharedPreferencesForTest(); @@ -228,6 +249,37 @@ void main() { p.dispose(); }); + test('session invalidation purges the service queued writes so the next account cannot replay them', () async { + const uuid = 'profile-invalidated'; + await _simklStore.save(uuid, _simkl(username: 'carol')); + BaseSharedPreferencesService.resetForTesting(); + + // Rows queued under the Simkl account about to be invalidated, plus a + // Trakt control row that must survive the purge. + final queue = TrackerWriteQueue(); + await queue.enqueue(uuid, _queuedItem(TrackerService.simkl)); + await queue.enqueue(uuid, _queuedItem(TrackerService.trakt)); + + final p = TrackersProvider(); + await _bindProfile(p, uuid); + expect(p.isSimklConnected, isTrue); + + // Fire the auth-failure teardown exactly as the client does on a 401. + SimklTracker.instance.client!.onSessionInvalidated(); + expect(p.isSimklConnected, isFalse); + await pumpEventQueue(); + + // Account B connecting starts with a flush; nothing of A's may be waiting. + await TrackerCoordinator.instance.flushWriteQueue(); + final survivors = await queue.load(uuid); + expect(survivors.map((item) => item.service), [ + TrackerService.trakt, + ], reason: 'the invalidated service rows must not wait for the next account'); + expect(await _simklStore.load(uuid), isNull); + + p.dispose(); + }); + test('disconnect during an in-flight profile load keeps the other trackers loaded', () async { const uuid = 'profile-race'; await _malStore.save(uuid, _mal(username: 'alice')); diff --git a/test/services/trackers/tracker_write_queue_test.dart b/test/services/trackers/tracker_write_queue_test.dart index c9a78df57..b45f09040 100644 --- a/test/services/trackers/tracker_write_queue_test.dart +++ b/test/services/trackers/tracker_write_queue_test.dart @@ -7,6 +7,9 @@ import 'package:plezy/services/base_shared_preferences_service.dart'; import 'package:plezy/services/trackers/tracker_constants.dart'; import 'package:plezy/services/trackers/tracker_write_queue.dart'; import 'package:plezy/utils/external_ids.dart'; +import 'package:shared_preferences_platform_interface/in_memory_shared_preferences_async.dart'; +import 'package:shared_preferences_platform_interface/shared_preferences_async_platform_interface.dart'; +import 'package:shared_preferences_platform_interface/types.dart'; import '../../test_helpers/prefs.dart'; @@ -273,6 +276,39 @@ void main() { expect(sent.map((item) => item.service), [TrackerService.mal], reason: 'the purged service must not dispatch'); }); + test('a persist-failure enqueue racing the disconnect purge cannot resurrect the row', () async { + final platform = _FailingQueuePreferences(); + SharedPreferencesAsyncPlatform.instance = platform; + BaseSharedPreferencesService.resetForTesting(); + + final queue = TrackerWriteQueue(); + final ctx = _episode(); + final key = trackerItemCoalesceKey(TrackerService.trakt, ctx, trackerExternalRowIdentity(ctx.external))!; + + // The enqueue claims its lock slot synchronously, then its persist fails + // and the row can only be buffered in memory. The purge claims the next + // slot in the same synchronous segment — the interleaving of a write + // failing while the user disconnects the service. + platform.failQueueWrites = true; + final enqueue = queue.enqueue('user-a', _item(ctx: ctx, coalesceKey: key)); + final purge = queue.removeService('user-a', TrackerService.trakt); + await enqueue; + await purge; + platform.failQueueWrites = false; + + final sent = []; + await queue.flush( + 'user-a', + send: (item) async { + sent.add(item); + return TrackerWriteDisposition.done; + }, + ); + + expect(sent, isEmpty, reason: 'the buffered row was created under the disconnected account'); + expect(await queue.load('user-a'), isEmpty); + }); + test('legacy Trakt rows migrate once with their intent and episode metadata intact', () async { final prefs = await BaseSharedPreferencesService.sharedCache(); const user = 'legacy-user'; @@ -378,3 +414,21 @@ void main() { expect(prefs.getString(archiveKey), corruptPayload); }); } + +/// Fails writes to the tracker queue key while armed, simulating the disk +/// full/revoked-storage case that feeds the in-memory fallback. +/// `SharedPreferencesWithCache` sits above the platform and is not +/// subclassable, so the failure is injected here. +final class _FailingQueuePreferences extends InMemorySharedPreferencesAsync { + _FailingQueuePreferences() : super.empty(); + + bool failQueueWrites = false; + + @override + Future setString(String key, String value, SharedPreferencesOptions options) async { + if (failQueueWrites && key.contains('tracker_write_queue')) { + throw StateError('simulated persist failure'); + } + return super.setString(key, value, options); + } +}