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.
This commit is contained in:
edde746
2026-08-21 19:23:44 +02:00
parent 7fc559e64b
commit 7652d060e2
5 changed files with 148 additions and 27 deletions
+12 -4
View File
@@ -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<void> _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) {
@@ -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<void> purgeWriteQueueForService(TrackerService service) async {
+26 -20
View File
@@ -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<void> 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<void> 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<TrackerWriteQueueItem>.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<TrackerWriteQueueItem>.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<void> 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);
@@ -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<void> _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'));
@@ -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 = <TrackerWriteQueueItem>[];
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<bool> 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);
}
}