From f021a1609863b5d68783516c1965c6b14d632faf Mon Sep 17 00:00:00 2001 From: Hongkuan Zhou <6771308+tedzhouhk@users.noreply.github.com> Date: Tue, 11 Aug 2026 20:18:30 -0700 Subject: [PATCH 1/6] refactor: share BeeCount Cloud provider instance --- lib/pages/cloud/devices_page.dart | 22 ++++-------- lib/providers/sync_providers.dart | 25 +++++++------ .../beecount_cloud_auth_provider_test.dart | 36 +++++++++++++++++++ 3 files changed, 58 insertions(+), 25 deletions(-) create mode 100644 test/providers/beecount_cloud_auth_provider_test.dart diff --git a/lib/pages/cloud/devices_page.dart b/lib/pages/cloud/devices_page.dart index 7d1923949..9db47439a 100644 --- a/lib/pages/cloud/devices_page.dart +++ b/lib/pages/cloud/devices_page.dart @@ -82,17 +82,11 @@ class _DevicesPageState extends ConsumerState { /// 获取 BeeCountCloudProvider 实例(仅 beecountCloud 后端可用) Future _getCloudProvider() async { - final config = await ref.read(activeCloudConfigProvider.future); - if (!config.valid || config.type != CloudBackendType.beecountCloud) { - throw StateError( - AppLocalizations.of(context).cloudCollabUnavailableMessage); - } - final services = await createCloudServices(config); - if (services.provider == null || services.provider is! BeeCountCloudProvider) { - throw StateError( - AppLocalizations.of(context).cloudCollabUnavailableMessage); - } - return services.provider as BeeCountCloudProvider; + final unavailableMessage = + AppLocalizations.of(context).cloudCollabUnavailableMessage; + final provider = await ref.read(beecountCloudProviderInstance.future); + if (provider == null) throw StateError(unavailableMessage); + return provider; } Future _reload({bool keepLoadingState = true}) async { @@ -104,11 +98,9 @@ class _DevicesPageState extends ConsumerState { _scopeDenied = false; }); try { - final auth = await ref.read(authServiceProvider.future); - final user = await auth.currentUser; - final currentDeviceId = user?.metadata?['deviceId']?.toString(); - final provider = await _getCloudProvider(); + final user = await provider.auth.currentUser; + final currentDeviceId = user?.metadata?['deviceId']?.toString(); final devices = await provider.listDevices( view: _showAllSessions ? 'sessions' : 'deduped', activeWithinDays: 30, diff --git a/lib/providers/sync_providers.dart b/lib/providers/sync_providers.dart index 0dd78abe2..255b644de 100644 --- a/lib/providers/sync_providers.dart +++ b/lib/providers/sync_providers.dart @@ -145,19 +145,26 @@ final s3ConfigProvider = FutureProvider((ref) async { }); final authServiceProvider = FutureProvider((ref) async { - final activeAsync = ref.watch(activeCloudConfigProvider); - if (!activeAsync.hasValue) { - return NoopAuthService(); - } - - final config = activeAsync.value!; + final config = await ref.watch(activeCloudConfigProvider.future); if (!config.valid || config.type == CloudBackendType.local) { return NoopAuthService(); } try { + // BeeCount Cloud 必须复用同步引擎持有的唯一 provider/auth 实例。 + // 多个实例虽然共用 SharedPreferences,却各自缓存 session;这会让 2FA + // 登录成功后同步实例仍停留在未登录状态。 + if (config.type == CloudBackendType.beecountCloud) { + final provider = await ref.watch(beecountCloudProviderInstance.future); + return provider?.auth ?? NoopAuthService(); + } + final services = await createCloudServices(config); if (services.auth != null) { + final provider = services.provider; + if (provider != null) { + ref.onDispose(() => unawaited(provider.dispose())); + } return services.auth!; } } catch (e) { @@ -492,10 +499,7 @@ final syncServiceProvider = Provider((ref) { /// 用于 SyncEngine 和其他需要直接访问 BeeCount Cloud API 的场景 final beecountCloudProviderInstance = FutureProvider((ref) async { - final configAsync = ref.watch(activeCloudConfigProvider); - if (!configAsync.hasValue) return null; - - final config = configAsync.value!; + final config = await ref.watch(activeCloudConfigProvider.future); if (!config.valid || config.type != CloudBackendType.beecountCloud) { return null; } @@ -504,6 +508,7 @@ final beecountCloudProviderInstance = final services = await createCloudServices(config); if (services.provider is! BeeCountCloudProvider) return null; final provider = services.provider as BeeCountCloudProvider; + ref.onDispose(() => unawaited(provider.dispose())); final email = config.beecountCloudEmail; final password = config.beecountCloudPassword; diff --git a/test/providers/beecount_cloud_auth_provider_test.dart b/test/providers/beecount_cloud_auth_provider_test.dart new file mode 100644 index 000000000..3a43ea036 --- /dev/null +++ b/test/providers/beecount_cloud_auth_provider_test.dart @@ -0,0 +1,36 @@ +import 'package:flutter_cloud_sync/flutter_cloud_sync.dart'; +import 'package:flutter_riverpod/flutter_riverpod.dart'; +import 'package:flutter_test/flutter_test.dart'; +import 'package:shared_preferences/shared_preferences.dart'; + +import 'package:beecount/providers/sync_providers.dart'; + +void main() { + TestWidgetsFlutterBinding.ensureInitialized(); + + setUp(() { + SharedPreferences.setMockInitialValues({}); + }); + + test('BeeCount Cloud UI auth and SyncEngine share one auth instance', + () async { + const config = CloudServiceConfig( + type: CloudBackendType.beecountCloud, + name: 'BeeCount Cloud', + beecountCloudBaseUrl: 'https://cloud.example.com', + beecountCloudApiPrefix: '/api/v1', + ); + final container = ProviderContainer( + overrides: [ + activeCloudConfigProvider.overrideWith((ref) async => config), + ], + ); + addTearDown(container.dispose); + + final provider = await container.read(beecountCloudProviderInstance.future); + final auth = await container.read(authServiceProvider.future); + + expect(provider, isNotNull); + expect(identical(auth, provider!.auth), isTrue); + }); +} From 40be2fb30533a750a49b7f1069ae320dbe049552 Mon Sep 17 00:00:00 2001 From: Hongkuan Zhou <6771308+tedzhouhk@users.noreply.github.com> Date: Tue, 11 Aug 2026 20:50:00 -0700 Subject: [PATCH 2/6] fix: avoid disposing providers before cached engines --- lib/providers/sync_providers.dart | 5 ----- 1 file changed, 5 deletions(-) diff --git a/lib/providers/sync_providers.dart b/lib/providers/sync_providers.dart index 255b644de..70dae4dbe 100644 --- a/lib/providers/sync_providers.dart +++ b/lib/providers/sync_providers.dart @@ -161,10 +161,6 @@ final authServiceProvider = FutureProvider((ref) async { final services = await createCloudServices(config); if (services.auth != null) { - final provider = services.provider; - if (provider != null) { - ref.onDispose(() => unawaited(provider.dispose())); - } return services.auth!; } } catch (e) { @@ -508,7 +504,6 @@ final beecountCloudProviderInstance = final services = await createCloudServices(config); if (services.provider is! BeeCountCloudProvider) return null; final provider = services.provider as BeeCountCloudProvider; - ref.onDispose(() => unawaited(provider.dispose())); final email = config.beecountCloudEmail; final password = config.beecountCloudPassword; From 3c4f799e73ebe693a4b55e30011d8261eee9aba7 Mon Sep 17 00:00:00 2001 From: Hongkuan Zhou <6771308+tedzhouhk@users.noreply.github.com> Date: Sat, 15 Aug 2026 13:22:56 -0700 Subject: [PATCH 3/6] refactor: narrow shared provider change --- lib/pages/cloud/devices_page.dart | 22 +++++++++++++++------- 1 file changed, 15 insertions(+), 7 deletions(-) diff --git a/lib/pages/cloud/devices_page.dart b/lib/pages/cloud/devices_page.dart index 9db47439a..7d1923949 100644 --- a/lib/pages/cloud/devices_page.dart +++ b/lib/pages/cloud/devices_page.dart @@ -82,11 +82,17 @@ class _DevicesPageState extends ConsumerState { /// 获取 BeeCountCloudProvider 实例(仅 beecountCloud 后端可用) Future _getCloudProvider() async { - final unavailableMessage = - AppLocalizations.of(context).cloudCollabUnavailableMessage; - final provider = await ref.read(beecountCloudProviderInstance.future); - if (provider == null) throw StateError(unavailableMessage); - return provider; + final config = await ref.read(activeCloudConfigProvider.future); + if (!config.valid || config.type != CloudBackendType.beecountCloud) { + throw StateError( + AppLocalizations.of(context).cloudCollabUnavailableMessage); + } + final services = await createCloudServices(config); + if (services.provider == null || services.provider is! BeeCountCloudProvider) { + throw StateError( + AppLocalizations.of(context).cloudCollabUnavailableMessage); + } + return services.provider as BeeCountCloudProvider; } Future _reload({bool keepLoadingState = true}) async { @@ -98,9 +104,11 @@ class _DevicesPageState extends ConsumerState { _scopeDenied = false; }); try { - final provider = await _getCloudProvider(); - final user = await provider.auth.currentUser; + final auth = await ref.read(authServiceProvider.future); + final user = await auth.currentUser; final currentDeviceId = user?.metadata?['deviceId']?.toString(); + + final provider = await _getCloudProvider(); final devices = await provider.listDevices( view: _showAllSessions ? 'sessions' : 'deduped', activeWithinDays: 30, From b1174653d505fb8c23d429d20aa58a44fe0c1ee4 Mon Sep 17 00:00:00 2001 From: Hongkuan Zhou <6771308+tedzhouhk@users.noreply.github.com> Date: Mon, 7 Sep 2026 15:09:20 -0700 Subject: [PATCH 4/6] fix: reuse BeeCount Cloud provider for connection tests --- lib/pages/cloud/cloud_service_page.dart | 11 ++-- .../_fakes/fake_beecount_cloud_provider.dart | 4 ++ test/pages/cloud/cloud_service_page_test.dart | 58 +++++++++++++++++++ .../beecount_cloud_auth_provider_test.dart | 15 ++++- 4 files changed, 81 insertions(+), 7 deletions(-) create mode 100644 test/pages/cloud/cloud_service_page_test.dart diff --git a/lib/pages/cloud/cloud_service_page.dart b/lib/pages/cloud/cloud_service_page.dart index da8134b83..8d3b871e2 100644 --- a/lib/pages/cloud/cloud_service_page.dart +++ b/lib/pages/cloud/cloud_service_page.dart @@ -1883,14 +1883,17 @@ class _CloudServicePageState extends ConsumerState { break; case CloudBackendType.beecountCloud: - // BeeCount Cloud 连接测试 - 调用健康检查接口 + // BeeCount Cloud 连接测试必须复用当前同步引擎的 provider。 + // 单独 createCloudServices 会再创建一套 auth/session,2FA 登录 + // 后可能与主实例分裂,并触发额外的静默登录。 try { - final services = await createCloudServices(config); - if (services.provider == null) { + final provider = + await ref.read(beecountCloudProviderInstance.future); + if (provider == null) { throw Exception('BeeCount Cloud provider 初始化失败'); } // 尝试列出文件验证连接 - await services.provider!.storage.list(path: ''); + await provider.storage.list(path: ''); connectionSuccess = true; } catch (e) { String errorMsg = e.toString(); diff --git a/test/cloud/sync/_fakes/fake_beecount_cloud_provider.dart b/test/cloud/sync/_fakes/fake_beecount_cloud_provider.dart index 2c69ac8f9..e93779489 100644 --- a/test/cloud/sync/_fakes/fake_beecount_cloud_provider.dart +++ b/test/cloud/sync/_fakes/fake_beecount_cloud_provider.dart @@ -60,6 +60,7 @@ class FakeBeeCountCloudAuthService extends BeeCountCloudAuthService { class FakeBeeCountCloudStorageService implements CloudStorageService { final Map _files = {}; final Map?> _metadata = {}; + int listCallCount = 0; /// 测试 helper:模拟 server 端账本列表(`storage.list(path: '')` 返回) final List ledgerSnapshots = []; @@ -88,6 +89,7 @@ class FakeBeeCountCloudStorageService implements CloudStorageService { @override Future> list({required String path}) async { + listCallCount++; // 测试关注的是"远端账本列表",由 [ledgerSnapshots] 控制 return List.unmodifiable(ledgerSnapshots); } @@ -147,6 +149,8 @@ class FakeBeeCountCloudProvider extends BeeCountCloudProvider { /// 控制 storage.list 是否抛错 Exception? storageListError; + int get storageListCallCount => _fakeStorage.listCallCount; + final StreamController _realtimeController = StreamController.broadcast(); diff --git a/test/pages/cloud/cloud_service_page_test.dart b/test/pages/cloud/cloud_service_page_test.dart new file mode 100644 index 000000000..370548d1d --- /dev/null +++ b/test/pages/cloud/cloud_service_page_test.dart @@ -0,0 +1,58 @@ +import 'package:flutter/material.dart'; +import 'package:flutter_cloud_sync/flutter_cloud_sync.dart'; +import 'package:flutter_riverpod/flutter_riverpod.dart'; +import 'package:flutter_test/flutter_test.dart'; +import 'package:shared_preferences/shared_preferences.dart'; + +import 'package:beecount/l10n/app_localizations.dart'; +import 'package:beecount/pages/cloud/cloud_service_page.dart'; +import 'package:beecount/providers/sync_providers.dart'; + +import '../../cloud/sync/_fakes/fake_beecount_cloud_provider.dart'; + +void main() { + TestWidgetsFlutterBinding.ensureInitialized(); + + setUp(() { + SharedPreferences.setMockInitialValues({'multi_device_sync': false}); + }); + + testWidgets('BeeCount Cloud 测试连接复用主 provider', (tester) async { + const config = CloudServiceConfig( + type: CloudBackendType.beecountCloud, + name: 'BeeCount Cloud', + beecountCloudBaseUrl: 'https://cloud.example.com', + beecountCloudApiPrefix: '/api/v1', + ); + final fakeProvider = FakeBeeCountCloudProvider(); + + await tester.pumpWidget( + ProviderScope( + overrides: [ + activeCloudConfigProvider.overrideWith((ref) async => config), + beecountCloudConfigProvider.overrideWith((ref) async => config), + webdavConfigProvider.overrideWith((ref) async => null), + s3ConfigProvider.overrideWith((ref) async => null), + beecountCloudProviderInstance + .overrideWith((ref) async => fakeProvider), + ], + child: MaterialApp( + localizationsDelegates: AppLocalizations.localizationsDelegates, + supportedLocales: AppLocalizations.supportedLocales, + locale: const Locale('zh'), + home: const CloudServicePage(), + ), + ), + ); + await tester.pumpAndSettle(); + + await tester.tap(find.byIcon(Icons.wifi_find)); + await tester.pump(); + await tester.pump(const Duration(milliseconds: 300)); + + expect(fakeProvider.storageListCallCount, 1); + + await tester.tap(find.text('确定')); + await tester.pumpAndSettle(); + }); +} diff --git a/test/providers/beecount_cloud_auth_provider_test.dart b/test/providers/beecount_cloud_auth_provider_test.dart index 3a43ea036..38e22373c 100644 --- a/test/providers/beecount_cloud_auth_provider_test.dart +++ b/test/providers/beecount_cloud_auth_provider_test.dart @@ -27,10 +27,19 @@ void main() { ); addTearDown(container.dispose); - final provider = await container.read(beecountCloudProviderInstance.future); - final auth = await container.read(authServiceProvider.future); + // 真机上 Mine/云服务页会先 watch UI auth,同步引擎随后 + // eager-load provider。主线在这个顺序下会各自创建一套 auth。 + final authFuture = container.read(authServiceProvider.future); + final providerFuture = container.read(beecountCloudProviderInstance.future); + + final auth = await authFuture; + final provider = await providerFuture; expect(provider, isNotNull); - expect(identical(auth, provider!.auth), isTrue); + expect( + identical(auth, provider!.auth), + isTrue, + reason: 'UI 与同步引擎不应持有两套独立 session', + ); }); } From 3f130b668d192bba868f62c403476ad4e9a5262d Mon Sep 17 00:00:00 2001 From: Hongkuan Zhou <6771308+tedzhouhk@users.noreply.github.com> Date: Sat, 12 Sep 2026 21:13:34 -0700 Subject: [PATCH 5/6] refactor: centralize cloud session ownership and lifecycle --- lib/cloud/cloud_session_manager.dart | 143 ++++++++++++++++ lib/cloud/sync/sync_engine.dart | 21 ++- lib/cloud/sync/sync_engine_profile.dart | 3 + lib/cloud/sync/sync_engine_realtime.dart | 3 + lib/cloud/sync/sync_providers.dart | 12 +- lib/cloud/transactions_sync_manager.dart | 6 +- lib/pages/cloud/cloud_service_page.dart | 11 +- lib/pages/cloud/devices_page.dart | 15 +- lib/providers/sync_providers.dart | 99 ++++------- .../lib/src/config/provider_factory.dart | 9 +- .../providers/beecount_cloud_provider.dart | 41 ++++- test/cloud/cloud_session_manager_test.dart | 160 ++++++++++++++++++ test/pages/cloud/cloud_service_page_test.dart | 5 +- test/providers/cloud_session_auth_test.dart | 122 +++++++++++++ test/providers/cloud_session_engine_test.dart | 56 ++++++ 15 files changed, 599 insertions(+), 107 deletions(-) create mode 100644 lib/cloud/cloud_session_manager.dart create mode 100644 test/cloud/cloud_session_manager_test.dart create mode 100644 test/providers/cloud_session_auth_test.dart create mode 100644 test/providers/cloud_session_engine_test.dart diff --git a/lib/cloud/cloud_session_manager.dart b/lib/cloud/cloud_session_manager.dart new file mode 100644 index 000000000..d902630f5 --- /dev/null +++ b/lib/cloud/cloud_session_manager.dart @@ -0,0 +1,143 @@ +import 'dart:async'; +import 'dart:convert'; + +import 'package:flutter_cloud_sync/flutter_cloud_sync.dart'; + +typedef CloudServices = ({CloudProvider? provider, CloudAuthService? auth}); +typedef CloudServicesFactory = Future Function( + CloudServiceConfig); + +/// Owns one active configuration's services. The factory remains uncached. +class CloudSession { + CloudSession(this.config, this.services); + + final CloudServiceConfig config; + final CloudServices services; + final Set _onClose = {}; + Future? _closing; + bool get isClosed => _closing != null; + + /// Consumers stop scheduling work before their network resources are released. + void Function() onClose(void Function() callback) { + if (isClosed) { + callback(); + } else { + _onClose.add(callback); + } + return () => _onClose.remove(callback); + } + + Future close() { + if (_closing != null) return _closing!; + final done = Completer(); + _closing = done.future; + Future(() async { + Object? firstError; + StackTrace? firstStack; + try { + for (final callback in _onClose.toList()) { + try { + callback(); + } catch (error, stack) { + firstError ??= error; + firstStack ??= stack; + } + } + } finally { + _onClose.clear(); + await services.provider?.dispose(); + } + if (firstError != null) { + Error.throwWithStackTrace(firstError!, firstStack!); + } + }).then(done.complete, onError: done.completeError); + return done.future; + } +} + +class CloudSessionManager { + CloudSessionManager( + {CloudServicesFactory factory = createCloudServices, + CloudServicesFactory? temporaryFactory}) + : _factory = factory, + _temporaryFactory = temporaryFactory ?? + ((config) => createCloudServices(config, persistSession: false)); + + final CloudServicesFactory _factory; + final CloudServicesFactory _temporaryFactory; + CloudSession? _active; + Future _tail = Future.value(); + Future? _pending; + String? _requestedKey; + int _generation = 0; + bool _disposed = false; + + Future activate(CloudServiceConfig config) { + if (_disposed) return Future.error(StateError('Session manager closed')); + // Includes credentials: editing credentials must replace the session too. + // Never log this key, since configuration can contain passwords. + final key = jsonEncode(config.toJson()); + if (_requestedKey == key && _pending != null) return _pending!; + _requestedKey = key; + final generation = ++_generation; + final next = _tail.then((_) async { + _checkGeneration(generation); + final previous = _active; + _active = null; + await previous?.close(); + _checkGeneration(generation); + final services = await _factory(config); + final session = CloudSession(config, services); + if (_disposed || generation != _generation) { + await session.close(); + throw StateError('Cloud configuration changed during initialization'); + } + // Backend-specific recovery belongs in session assembly, not UI providers. + final auth = services.auth; + if (auth is BeeCountCloudAuthService) { + auth.setRecoveryCredentials( + email: config.beecountCloudEmail, + password: config.beecountCloudPassword, + ); + } + _active = session; + return session; + }); + _pending = next; + _tail = next.then((_) {}, onError: (Object error, StackTrace stack) { + if (generation == _generation) { + _requestedKey = null; + _pending = null; + } + }); + return next; + } + + void _checkGeneration(int generation) { + if (_disposed || generation != _generation) { + throw StateError('Cloud configuration changed'); + } + } + + /// Draft configurations never replace the active session. Always released. + Future withTemporarySession(CloudServiceConfig config, + Future Function(CloudSession) action) async { + if (_disposed) throw StateError('Session manager closed'); + final session = CloudSession(config, await _temporaryFactory(config)); + try { + if (_disposed) throw StateError('Session manager closed'); + return await action(session); + } finally { + await session.close(); + } + } + + Future dispose() async { + _disposed = true; + ++_generation; + await _tail; + final previous = _active; + _active = null; + await previous?.close(); + } +} diff --git a/lib/cloud/sync/sync_engine.dart b/lib/cloud/sync/sync_engine.dart index f487163dd..1b7b23e3d 100644 --- a/lib/cloud/sync/sync_engine.dart +++ b/lib/cloud/sync/sync_engine.dart @@ -78,6 +78,7 @@ class SyncEngine implements app.SyncService { /// 状态缓存 final Map _statusCache = {}; bool _localChanged = false; + bool _disposed = false; /// WebSocket 实时监听 StreamSubscription? _realtimeSubscription; @@ -361,6 +362,9 @@ class SyncEngine implements app.SyncService { /// 释放资源 void dispose() { + if (_disposed) return; + _disposed = true; + ledgerIdResolver = null; stopListeningRealtime(); _eventsController.close(); } @@ -369,6 +373,7 @@ class SyncEngine implements app.SyncService { /// 执行完整同步(先 push 后 pull) Future sync({required String ledgerId}) async { + if (_disposed) return const SyncResult(error: 'Cloud session closed'); logger.info('SyncEngine', '开始同步 ledger=$ledgerId'); try { final ledgerIdInt = int.tryParse(ledgerId) ?? -1; @@ -532,14 +537,11 @@ class SyncEngine implements app.SyncService { /// 此时设备全局 cursor 可能已经前移、普通 `_pull` 再也拉不回历史。 /// /// 返回新增(非已存在)的账本数,调用方可据此决定要不要 bump 刷新信号。 - /// 并发互斥锁 — **static** 跨 SyncEngine 实例共享。 - /// 关键 bug:join page 拿 syncEngineProvider(family) 的 engine,WS listener - /// 拿 cloudSyncServiceProvider 创建的 engine,两个不同 instance!instance-level - /// 字段互不知道,各跑各的。改 static 后整个进程同一时间只有一个 fetch-then-write - /// 在跑。 - static Completer? _syncLedgersInFlight; + /// 同一会话共用一个 engine;切换会话不能复用旧服务器的 in-flight 结果。 + Completer? _syncLedgersInFlight; Future syncLedgersFromServer() async { + if (_disposed) throw StateError('Cloud session closed'); final existing = _syncLedgersInFlight; if (existing != null) { logger.info('SyncEngine', 'syncLedgersFromServer 已在执行中,等待 in-flight 结果'); @@ -563,6 +565,7 @@ class SyncEngine implements app.SyncService { logger.info('SyncEngine', 'syncLedgersFromServer start'); try { final remote = await provider.readLedgers(); + if (_disposed) return 0; int upserted = 0; int inserted = 0; // 新设备登录场景:Editor 已是 server LedgerMember 但本地 ledgers 表为空。 @@ -572,6 +575,7 @@ class SyncEngine implements app.SyncService { // 也会让单个失败影响其它账本。 final newSharedLedgerSyncIds = []; for (final r in remote) { + if (_disposed) return inserted; final syncId = r.ledgerId; if (syncId.isEmpty) continue; // 用 get() 不用 getSingleOrNull() — 历史可能已经产生过同 syncId 多行 @@ -711,6 +715,7 @@ class SyncEngine implements app.SyncService { /// 否则单飞失效。这俩内部应该只处理 ledger-scope change(transaction / budget / /// ledger / ledger_snapshot)。 Future pushUserGlobalEntities() async { + if (_disposed) throw StateError('Cloud session closed'); final inFlight = _userGlobalPushInFlight; if (inFlight != null) { logger.info('SyncEngine', 'pushUserGlobalEntities 已在执行,复用 in-flight'); @@ -904,6 +909,7 @@ class SyncEngine implements app.SyncService { /// user-global change(account / category / tag)由 [pushUserGlobalEntities] 统一推 /// (在 [_doPush] 开头调用),避免多账本场景下并行 push 重复推送 user-global。 Future push(String ledgerId) async { + if (_disposed) throw StateError('Cloud session closed'); final inFlight = _pushInFlight[ledgerId]; if (inFlight != null) { logger.info('SyncEngine', 'push(ledger=$ledgerId) 已在执行,复用 in-flight'); @@ -1062,6 +1068,7 @@ class SyncEngine implements app.SyncService { /// 轮) /// - replay(sinceOverride 非 null)语义独立,等 in-flight 完成后再自己跑 Future pull(String ledgerId, {int? sinceOverride}) async { + if (_disposed) throw StateError('Cloud session closed'); // 1. in-flight 单飞 final inFlight = _pullInFlight; if (inFlight != null) { @@ -1214,6 +1221,7 @@ class SyncEngine implements app.SyncService { /// - SQLite busy/locked → 单条 retry 2 次 Future<_PullPageOutcome> _applyPullPage( List changes) async { + if (_disposed) throw StateError('Cloud session closed'); int applied = 0; int skipped = 0; BeeCountCloudSyncChange? failingChange; @@ -1221,6 +1229,7 @@ class SyncEngine implements app.SyncService { try { await db.transaction(() async { for (final ch in changes) { + if (_disposed) throw StateError('Cloud session closed'); failingChange = ch; final ok = await _applyOneWithBusyRetry(ch); if (ok) { diff --git a/lib/cloud/sync/sync_engine_profile.dart b/lib/cloud/sync/sync_engine_profile.dart index 2692f9bbb..779b3abf3 100644 --- a/lib/cloud/sync/sync_engine_profile.dart +++ b/lib/cloud/sync/sync_engine_profile.dart @@ -15,12 +15,14 @@ extension SyncEngineProfile on SyncEngine { /// PR 3:不再接 callback,所有字段更新走 [events] stream emit /// `ProfileFieldApplied` 事件,UI 通过 syncEventStreamProvider 订阅处理。 Future syncMyProfile() async { + if (_disposed) return false; final localVersion = await AvatarService.getStoredRemoteVersion(); logger.info('avatar_sync', 'syncMyProfile start, localVersion=$localVersion'); bool anyChanged = false; try { final profile = await provider.getMyProfile(); + if (_disposed) return false; // === theme_primary_color === final theme = profile.themePrimaryColor; @@ -87,6 +89,7 @@ extension SyncEngineProfile on SyncEngine { userId: profile.userId, version: remoteVersion > 0 ? remoteVersion : null, ); + if (_disposed) return false; logger.info('avatar_sync', 'downloaded size=${bytes.length}B'); await AvatarService.saveAvatarFromBytes(bytes); await AvatarService.setStoredRemoteVersion(remoteVersion); diff --git a/lib/cloud/sync/sync_engine_realtime.dart b/lib/cloud/sync/sync_engine_realtime.dart index ba9bc776f..506541e4c 100644 --- a/lib/cloud/sync/sync_engine_realtime.dart +++ b/lib/cloud/sync/sync_engine_realtime.dart @@ -8,6 +8,7 @@ part of 'sync_engine.dart'; extension SyncEngineRealtime on SyncEngine { /// 开始监听 WebSocket 实时事件,收到变更通知时自动触发 pull void startListeningRealtime() { + if (_disposed) return; _realtimeSubscription?.cancel(); // 启动 WebSocket 连接,否则 realtimeEvents 流永远为空 provider.startRealtime().catchError((e) { @@ -66,6 +67,7 @@ extension SyncEngineRealtime on SyncEngine { /// 2 秒防抖:WiFi ↔ 移动网络切换、或 WS reconnect 接着 connectivity 事件 /// 这种"连续上线信号"只触发 1 次 sync。 void _scheduleAutoSync({required String reason}) { + if (_disposed) return; _autoSyncDebounce?.cancel(); _autoSyncDebounce = Timer(const Duration(seconds: 2), () async { if (_autoSyncing) { @@ -617,6 +619,7 @@ extension SyncEngineRealtime on SyncEngine { /// 防抖调度 pull(1 秒内多次触发只执行一次) void _schedulePull(String? ledgerId) { + if (_disposed) return; _pullDebounce?.cancel(); _pullDebounce = Timer(const Duration(seconds: 1), () async { if (_autoPulling) return; diff --git a/lib/cloud/sync/sync_providers.dart b/lib/cloud/sync/sync_providers.dart index 38572a5b8..780ca4a42 100644 --- a/lib/cloud/sync/sync_providers.dart +++ b/lib/cloud/sync/sync_providers.dart @@ -3,6 +3,7 @@ import 'package:flutter_cloud_sync/flutter_cloud_sync.dart' hide SyncStatus; import '../../providers/database_providers.dart'; +import '../../providers/sync_providers.dart' show activeCloudServicesProvider; import 'change_tracker.dart'; import 'sync_engine.dart'; @@ -23,6 +24,11 @@ final changeTrackerProvider = Provider((ref) { /// 装配 callback 后才启动,但 dispose 由 Riverpod GC family entry 时统一触发。 final syncEngineProvider = Provider.family( (ref, provider) { + final session = ref.watch(activeCloudServicesProvider).valueOrNull; + if (session == null || session.isClosed || + !identical(session.services.provider, provider)) { + throw StateError('SyncEngine requires the active cloud session'); + } final db = ref.watch(databaseProvider); final tracker = ref.watch(changeTrackerProvider); final repo = ref.watch(repositoryProvider); @@ -32,7 +38,11 @@ final syncEngineProvider = Provider.family( changeTracker: tracker, repo: repo, ); - ref.onDispose(() => engine.dispose()); + final unregister = session.onClose(engine.dispose); + ref.onDispose(() { + unregister(); + engine.dispose(); + }); return engine; }, ); diff --git a/lib/cloud/transactions_sync_manager.dart b/lib/cloud/transactions_sync_manager.dart index db30e5be5..8bb92b74b 100644 --- a/lib/cloud/transactions_sync_manager.dart +++ b/lib/cloud/transactions_sync_manager.dart @@ -11,6 +11,7 @@ import '../services/data_import_service.dart'; import '../services/system/logger_service.dart'; import 'sync_diff_service.dart'; import 'sync_service.dart'; +import 'cloud_session_manager.dart'; import 'transactions_json.dart'; /// 账本交易的云同步管理器 @@ -20,6 +21,7 @@ class TransactionsSyncManager implements SyncService { final fcs.CloudServiceConfig config; final BeeDatabase db; final BaseRepository repo; + final CloudSession session; fcs.CloudSyncManager? _syncManager; fcs.CloudProvider? _provider; @@ -34,6 +36,7 @@ class TransactionsSyncManager implements SyncService { required this.config, required this.db, required this.repo, + required this.session, }); @override @@ -47,6 +50,7 @@ class TransactionsSyncManager implements SyncService { /// 确保服务已初始化(延迟初始化) Future _ensureInitialized() async { + if (session.isClosed) throw StateError('Cloud session closed'); if (_isInitialized) return; if (_isInitializing) { // 等待初始化完成 @@ -67,7 +71,7 @@ class TransactionsSyncManager implements SyncService { /// 初始化 CloudProvider 和 SyncManager Future _initialize() async { - final services = await fcs.createCloudServices(config); + final services = session.services; _provider = services.provider; if (_provider == null) { diff --git a/lib/pages/cloud/cloud_service_page.dart b/lib/pages/cloud/cloud_service_page.dart index 8d3b871e2..95539d5b1 100644 --- a/lib/pages/cloud/cloud_service_page.dart +++ b/lib/pages/cloud/cloud_service_page.dart @@ -1883,12 +1883,10 @@ class _CloudServicePageState extends ConsumerState { break; case CloudBackendType.beecountCloud: - // BeeCount Cloud 连接测试必须复用当前同步引擎的 provider。 - // 单独 createCloudServices 会再创建一套 auth/session,2FA 登录 - // 后可能与主实例分裂,并触发额外的静默登录。 + // 已激活服务的认证与连接测试共用会话。 try { - final provider = - await ref.read(beecountCloudProviderInstance.future); + final session = await ref.read(activeCloudServicesProvider.future); + final provider = session.services.provider; if (provider == null) { throw Exception('BeeCount Cloud provider 初始化失败'); } @@ -1922,7 +1920,8 @@ class _CloudServicePageState extends ConsumerState { logger.info('CloudServicePage', 'S3 连接测试开始: endpoint=${cleanedConfig.s3Endpoint}, bucket=${cleanedConfig.s3Bucket}'); - final services = await createCloudServices(cleanedConfig); + final session = await ref.read(activeCloudServicesProvider.future); + final services = session.services; logger.info('CloudServicePage', 'S3 provider 创建结果: ${services.provider != null ? "成功" : "失败"}'); diff --git a/lib/pages/cloud/devices_page.dart b/lib/pages/cloud/devices_page.dart index 7d1923949..8a54360bd 100644 --- a/lib/pages/cloud/devices_page.dart +++ b/lib/pages/cloud/devices_page.dart @@ -82,17 +82,12 @@ class _DevicesPageState extends ConsumerState { /// 获取 BeeCountCloudProvider 实例(仅 beecountCloud 后端可用) Future _getCloudProvider() async { - final config = await ref.read(activeCloudConfigProvider.future); - if (!config.valid || config.type != CloudBackendType.beecountCloud) { - throw StateError( - AppLocalizations.of(context).cloudCollabUnavailableMessage); + final unavailable = AppLocalizations.of(context).cloudCollabUnavailableMessage; + final provider = await ref.read(beecountCloudProviderInstance.future); + if (provider == null) { + throw StateError(unavailable); } - final services = await createCloudServices(config); - if (services.provider == null || services.provider is! BeeCountCloudProvider) { - throw StateError( - AppLocalizations.of(context).cloudCollabUnavailableMessage); - } - return services.provider as BeeCountCloudProvider; + return provider; } Future _reload({bool keepLoadingState = true}) async { diff --git a/lib/providers/sync_providers.dart b/lib/providers/sync_providers.dart index 70dae4dbe..7f5696d31 100644 --- a/lib/providers/sync_providers.dart +++ b/lib/providers/sync_providers.dart @@ -9,6 +9,7 @@ import 'package:flutter_riverpod/flutter_riverpod.dart'; import 'package:shared_preferences/shared_preferences.dart'; import 'package:flutter_cloud_sync/flutter_cloud_sync.dart' hide SyncStatus; import '../cloud/sync_service.dart'; +import '../cloud/cloud_session_manager.dart'; import 'shared_ledger_providers.dart'; import '../cloud/sync/sync_coordinator.dart'; import '../cloud/sync/sync_engine.dart'; @@ -144,30 +145,26 @@ final s3ConfigProvider = FutureProvider((ref) async { return store.loadS3(); }); -final authServiceProvider = FutureProvider((ref) async { - final config = await ref.watch(activeCloudConfigProvider.future); - if (!config.valid || config.type == CloudBackendType.local) { - return NoopAuthService(); - } - - try { - // BeeCount Cloud 必须复用同步引擎持有的唯一 provider/auth 实例。 - // 多个实例虽然共用 SharedPreferences,却各自缓存 session;这会让 2FA - // 登录成功后同步实例仍停留在未登录状态。 - if (config.type == CloudBackendType.beecountCloud) { - final provider = await ref.watch(beecountCloudProviderInstance.future); - return provider?.auth ?? NoopAuthService(); - } +/// Owns services across projection rebuilds and configuration switches. +final cloudSessionManagerProvider = Provider((ref) { + final manager = CloudSessionManager(); + ref.onDispose(() => unawaited(manager.dispose().catchError( + (Object e, StackTrace st) => logger.warning('CloudSession', 'Dispose failed: $e', st)))); + return manager; +}); - final services = await createCloudServices(config); - if (services.auth != null) { - return services.auth!; - } - } catch (e) { - // 初始化失败,返回 NoopAuthService - } +final activeCloudServicesProvider = FutureProvider((ref) async { + final manager = ref.watch(cloudSessionManagerProvider); + var cancelled = false; + ref.onDispose(() => cancelled = true); + final config = await ref.watch(activeCloudConfigProvider.future); + if (cancelled) throw StateError('Cloud configuration changed'); + return manager.activate(config); +}); - return NoopAuthService(); +final authServiceProvider = FutureProvider((ref) async { + final session = await ref.watch(activeCloudServicesProvider.future); + return session.services.auth ?? NoopAuthService(); }); // 防重入锁:避免 Provider 重建导致多个自动同步并发执行 @@ -175,7 +172,7 @@ bool _autoSyncInProgress = false; final syncServiceProvider = Provider((ref) { final activeAsync = ref.watch(activeCloudConfigProvider); - if (!activeAsync.hasValue) return LocalOnlySyncService(); + if (activeAsync.isLoading || !activeAsync.hasValue) return LocalOnlySyncService(); final config = activeAsync.value!; if (!config.valid || config.type == CloudBackendType.local) { @@ -185,7 +182,7 @@ final syncServiceProvider = Provider((ref) { // BeeCount Cloud → SyncEngine(增量同步) if (config.type == CloudBackendType.beecountCloud) { final providerAsync = ref.watch(beecountCloudProviderInstance); - if (!providerAsync.hasValue || providerAsync.value == null) { + if (providerAsync.isLoading || !providerAsync.hasValue || providerAsync.value == null) { // Provider 尚未初始化,返回 LocalOnly 等待 return LocalOnlySyncService(); } @@ -488,57 +485,21 @@ final syncServiceProvider = Provider((ref) { // 其他 provider → TransactionsSyncManager(快照同步) final db = ref.watch(databaseProvider); final repo = ref.watch(repositoryProvider); - return TransactionsSyncManager(config: config, db: db, repo: repo); + final servicesAsync = ref.watch(activeCloudServicesProvider); + if (servicesAsync.isLoading) return LocalOnlySyncService(); + final session = servicesAsync.valueOrNull; + if (session == null || session.isClosed) return LocalOnlySyncService(); + return TransactionsSyncManager(config: config, db: db, repo: repo, + session: session); }); /// 已初始化的 BeeCountCloudProvider 实例 /// 用于 SyncEngine 和其他需要直接访问 BeeCount Cloud API 的场景 final beecountCloudProviderInstance = FutureProvider((ref) async { - final config = await ref.watch(activeCloudConfigProvider.future); - if (!config.valid || config.type != CloudBackendType.beecountCloud) { - return null; - } - - try { - final services = await createCloudServices(config); - if (services.provider is! BeeCountCloudProvider) return null; - final provider = services.provider as BeeCountCloudProvider; - - final email = config.beecountCloudEmail; - final password = config.beecountCloudPassword; - - // 把邮密交给 auth service,让它在任何时刻发现 session 失效都能自动重登。 - // 这是解决"token 过期后必须到配置页点一下才能恢复"的关键:auth service - // 内部会在 currentUser / requireAccessToken 触发时尝试恢复,不再等 Provider - // 重建。 - if (services.auth is BeeCountCloudAuthService) { - (services.auth as BeeCountCloudAuthService).setRecoveryCredentials( - email: email, - password: password, - ); - } - - // 双重保险:构造之后也触发一次 currentUser,让 initialize() 没恢复出 - // session 的场景立刻走一次恢复登录(email+password 有时),减少用户第一次 - // 操作时的卡顿感。currentUser 内部已经自带 _tryRecoveryLogin。 - if (services.auth != null) { - try { - final user = await services.auth!.currentUser; - if (user != null) { - logger.info('CloudSync', 'BeeCount Cloud session ready: ${user.email}'); - } else if (email != null && email.isNotEmpty) { - logger.info('CloudSync', 'BeeCount Cloud 未登录,等首次 API 触发恢复'); - } - } catch (e, st) { - logger.warning('CloudSync', 'BeeCount Cloud 初始 currentUser 失败: $e', st); - } - } - return provider; - } catch (e, st) { - logger.error('CloudSync', 'BeeCountCloudProvider 初始化失败', e, st); - } - return null; + final session = await ref.watch(activeCloudServicesProvider.future); + final provider = session.services.provider; + return provider is BeeCountCloudProvider ? provider : null; }); /// BeeCount Cloud 服务端版本号。Mine 页面 / 云同步页都能直接用;失败就 diff --git a/packages/flutter_cloud_sync/lib/src/config/provider_factory.dart b/packages/flutter_cloud_sync/lib/src/config/provider_factory.dart index e0c2a0f1f..0c5388acf 100644 --- a/packages/flutter_cloud_sync/lib/src/config/provider_factory.dart +++ b/packages/flutter_cloud_sync/lib/src/config/provider_factory.dart @@ -19,8 +19,9 @@ import 'cloud_service_config.dart'; /// - Supabase: 使用独立包内的初始化逻辑 /// - WebDAV: 创建新的 WebDAV provider Future<({CloudProvider? provider, CloudAuthService? auth})> createCloudServices( - CloudServiceConfig config, -) async { + CloudServiceConfig config, { + bool persistSession = true, +}) async { if (!config.valid) { return (provider: null, auth: null); } @@ -34,6 +35,8 @@ Future<({CloudProvider? provider, CloudAuthService? auth})> createCloudServices( await provider.initialize({ 'baseUrl': config.beecountCloudBaseUrl!, 'apiPrefix': config.beecountCloudApiPrefix ?? '/api/v1', + 'persistSession': persistSession, + 'email': config.beecountCloudEmail, }); return (provider: provider, auth: provider.auth); @@ -88,7 +91,7 @@ Future<({CloudProvider? provider, CloudAuthService? auth})> createCloudServices( // S3 初始化 - 不捕获异常,让错误向上传递以便调试 final provider = S3Provider(); await provider.initialize({ - 'endpoint': config.s3Endpoint!, + 'endpoint': config.s3Endpoint!.replaceFirst(RegExp(r'^https?://'), ''), 'region': config.s3Region ?? 'us-east-1', 'accessKey': config.s3AccessKey!, 'secretKey': config.s3SecretKey!, diff --git a/packages/flutter_cloud_sync/lib/src/providers/beecount_cloud_provider.dart b/packages/flutter_cloud_sync/lib/src/providers/beecount_cloud_provider.dart index 8856906cf..74c63c82e 100644 --- a/packages/flutter_cloud_sync/lib/src/providers/beecount_cloud_provider.dart +++ b/packages/flutter_cloud_sync/lib/src/providers/beecount_cloud_provider.dart @@ -126,8 +126,9 @@ class BeeCountCloudProvider implements CloudProvider { baseUrl: baseUrl, apiPrefix: apiPrefix, twoFactorHandler: BeeCountCloudProvider.globalTwoFactorHandler, + persistSession: config['persistSession'] != false, ); - await authService.initialize(); + await authService.initialize(expectedEmail: config['email'] as String?); _auth = authService; final storage = BeeCountCloudStorageService( @@ -1119,6 +1120,7 @@ class BeeCountCloudAuthService implements CloudAuthService { required this.apiPrefix, http.Client? httpClient, TwoFactorChallengeHandler? twoFactorHandler, + this.persistSession = true, }) : _httpClient = httpClient ?? http.Client(), _twoFactorHandler = twoFactorHandler; @@ -1126,6 +1128,8 @@ class BeeCountCloudAuthService implements CloudAuthService { final String apiPrefix; final http.Client _httpClient; final TwoFactorChallengeHandler? _twoFactorHandler; + final bool persistSession; + bool _disposed = false; final StreamController _authStateController = StreamController.broadcast(); @@ -1174,7 +1178,8 @@ class BeeCountCloudAuthService implements CloudAuthService { return 'beecount_cloud_local_device_id_$digest'; } - Future initialize() async { + Future initialize({String? expectedEmail}) async { + if (!persistSession) return; final prefs = await SharedPreferences.getInstance(); final raw = prefs.getString(_sessionStorageKey); if (raw == null || raw.isEmpty) { @@ -1183,7 +1188,11 @@ class BeeCountCloudAuthService implements CloudAuthService { try { final json = jsonDecode(raw) as Map; - _session = _BeeCountCloudSession.fromJson(json); + final saved = _BeeCountCloudSession.fromJson(json); + final email = expectedEmail?.trim().toLowerCase(); + if (email != null && email.isNotEmpty && + saved.email?.trim().toLowerCase() != email) return; + _session = saved; if (_isAccessTokenExpired(_session!)) { await _refreshSessionOrClear(); } else { @@ -1199,6 +1208,7 @@ class BeeCountCloudAuthService implements CloudAuthService { @override Future get currentUser async { + if (_disposed) return null; final session = _session; if (session == null) { // 完全没 session(从没登过 / session 被清了):只有带了恢复凭证才尝试 @@ -1218,6 +1228,7 @@ class BeeCountCloudAuthService implements CloudAuthService { } Future requireAccessToken() async { + if (_disposed) throw CloudNotAuthenticatedException('Cloud session closed'); final session = _session; if (session == null) { final recovered = await _tryRecoveryLogin(); @@ -1495,6 +1506,7 @@ class BeeCountCloudAuthService implements CloudAuthService { } Future _resolveOrCreateLocalDeviceId() async { + if (!persistSession) return _generateLocalDeviceId(); final prefs = await SharedPreferences.getInstance(); final existing = _trimOrNull(prefs.getString(_localDeviceIdStorageKey)); if (existing != null) { @@ -1594,6 +1606,8 @@ class BeeCountCloudAuthService implements CloudAuthService { } void dispose() { + if (_disposed) return; + _disposed = true; _authStateController.close(); _httpClient.close(); } @@ -1746,12 +1760,16 @@ class BeeCountCloudAuthService implements CloudAuthService { } Future _saveSession(_BeeCountCloudSession session) async { + if (_disposed) throw CloudNotAuthenticatedException('Cloud session closed'); _session = session; // 任何成功登录路径都清掉静默恢复冷却,避免之前的失败状态拖到现在。 _silentRecoveryCooldownUntil = null; - final prefs = await SharedPreferences.getInstance(); - await prefs.setString(_sessionStorageKey, jsonEncode(session.toJson())); - await prefs.setString(_localDeviceIdStorageKey, session.deviceId); + if (persistSession) { + final prefs = await SharedPreferences.getInstance(); + if (_disposed) throw CloudNotAuthenticatedException('Cloud session closed'); + await prefs.setString(_sessionStorageKey, jsonEncode(session.toJson())); + await prefs.setString(_localDeviceIdStorageKey, session.deviceId); + } final metadata = _deviceMetadataCache; if (metadata != null && metadata.deviceId != session.deviceId) { _deviceMetadataCache = _BeeCountDeviceMetadata( @@ -1767,13 +1785,17 @@ class BeeCountCloudAuthService implements CloudAuthService { } Future _clearSession() async { + if (_disposed) return; _session = null; - final prefs = await SharedPreferences.getInstance(); - await prefs.remove(_sessionStorageKey); - _authStateController.add(null); + if (persistSession) { + final prefs = await SharedPreferences.getInstance(); + await prefs.remove(_sessionStorageKey); + } + if (!_disposed) _authStateController.add(null); } void _emitCurrentUser() { + if (_disposed) return; final session = _session; if (session == null) { _authStateController.add(null); @@ -1805,6 +1827,7 @@ class BeeCountCloudAuthService implements CloudAuthService { Map? body, String? accessToken, }) async { + if (_disposed) throw CloudNotAuthenticatedException('Cloud session closed'); final uri = Uri.parse('$baseUrl$apiPrefix$path'); final request = http.Request(method, uri); request.headers['Content-Type'] = 'application/json'; diff --git a/test/cloud/cloud_session_manager_test.dart b/test/cloud/cloud_session_manager_test.dart new file mode 100644 index 000000000..9c8aba24d --- /dev/null +++ b/test/cloud/cloud_session_manager_test.dart @@ -0,0 +1,160 @@ +import 'dart:async'; + +import 'package:flutter_cloud_sync/flutter_cloud_sync.dart'; +import 'package:flutter_test/flutter_test.dart'; +import 'package:shared_preferences/shared_preferences.dart'; +import 'package:beecount/cloud/cloud_session_manager.dart'; + +import 'sync/_fakes/fake_beecount_cloud_provider.dart'; + +const configA = CloudServiceConfig( + type: CloudBackendType.beecountCloud, + name: 'A', + beecountCloudBaseUrl: 'https://a.example'); +const configB = CloudServiceConfig( + type: CloudBackendType.beecountCloud, + name: 'B', + beecountCloudBaseUrl: 'https://b.example'); + +class TrackedProvider extends FakeBeeCountCloudProvider { + int closed = 0; + @override + Future dispose() async { + closed++; + } +} + +void main() { + TestWidgetsFlutterBinding.ensureInitialized(); + setUp(() => SharedPreferences.setMockInitialValues({})); + + test('concurrent consumers and equal reloaded configurations share services', + () async { + var creates = 0; + final provider = TrackedProvider(); + final manager = CloudSessionManager(factory: (_) async { + creates++; + return (provider: provider, auth: provider.auth); + }); + final sessions = await Future.wait([ + manager.activate(configA), + manager.activate(configA), + manager.activate(CloudServiceConfig.fromJson(configA.toJson())), + ]); + expect(creates, 1); + expect(sessions.every((s) => identical(s, sessions.first)), isTrue); + await manager.dispose(); + await manager.dispose(); + expect(provider.closed, 1); + }); + + test( + 'switch stops consumers and releases old services before creating new ones', + () async { + final events = []; + final first = TrackedProvider(); + final second = TrackedProvider(); + final manager = CloudSessionManager(factory: (cfg) async { + if (cfg == configB) { + expect(first.closed, 1); + expect(events, ['stop engine']); + } + final p = cfg == configA ? first : second; + return (provider: p, auth: p.auth); + }); + final old = await manager.activate(configA); + old.onClose(() => events.add('stop engine')); + await manager.activate(configB); + expect(old.isClosed, isTrue); + await manager.dispose(); + expect(second.closed, 1); + }); + + test( + 'late initialization from superseded config is disposed, never published', + () async { + final started = Completer(); + final release = Completer(); + final first = TrackedProvider(); + final second = TrackedProvider(); + final manager = CloudSessionManager(factory: (cfg) async { + if (cfg == configA) { + started.complete(); + await release.future; + } + final p = cfg == configA ? first : second; + return (provider: p, auth: p.auth); + }); + final obsolete = manager.activate(configA); + final rejected = expectLater(obsolete, throwsStateError); + await started.future; + final current = manager.activate(configB); + release.complete(); + await rejected; + expect((await current).services.provider, second); + expect(first.closed, 1); + await manager.dispose(); + }); + + test('failed initialization can retry the same configuration', () async { + var attempts = 0; + final p = TrackedProvider(); + final manager = CloudSessionManager(factory: (_) async { + if (++attempts == 1) throw StateError('offline'); + return (provider: p, auth: p.auth); + }); + await expectLater(manager.activate(configA), throwsStateError); + expect((await manager.activate(configA)).services.provider, p); + await manager.dispose(); + }); + + test('a failing close listener does not skip remaining cleanup', () async { + final provider = TrackedProvider(); + final session = CloudSession(configA, + (provider: provider, auth: provider.auth)); + var stopped = false; + session.onClose(() => throw StateError('listener failed')); + session.onClose(() => stopped = true); + await expectLater(session.close(), throwsStateError); + expect(stopped, isTrue); + expect(provider.closed, 1); + }); + + test('shutdown during creation disposes late services', () async { + final started = Completer(); + final release = Completer(); + final p = TrackedProvider(); + final manager = CloudSessionManager(factory: (_) async { + started.complete(); + await release.future; + return (provider: p, auth: p.auth); + }); + final pending = expectLater(manager.activate(configA), throwsStateError); + await started.future; + final closing = manager.dispose(); + release.complete(); + await pending; + await closing; + expect(p.closed, 1); + }); + + test('temporary failure releases draft without replacing active services', + () async { + final p = TrackedProvider(); + final draft = TrackedProvider(); + final manager = CloudSessionManager( + factory: (_) async => (provider: p, auth: p.auth), + temporaryFactory: (_) async => (provider: draft, auth: draft.auth), + ); + final active = await manager.activate(configA); + await expectLater( + manager.withTemporarySession(configB, (_) async { + throw StateError('test failed'); + }), + throwsStateError); + expect(draft.closed, 1); + expect(p.closed, 0); + expect(await manager.activate(configA), same(active)); + await manager.dispose(); + }); +} diff --git a/test/pages/cloud/cloud_service_page_test.dart b/test/pages/cloud/cloud_service_page_test.dart index 370548d1d..2cd4a8dd0 100644 --- a/test/pages/cloud/cloud_service_page_test.dart +++ b/test/pages/cloud/cloud_service_page_test.dart @@ -5,6 +5,7 @@ import 'package:flutter_test/flutter_test.dart'; import 'package:shared_preferences/shared_preferences.dart'; import 'package:beecount/l10n/app_localizations.dart'; +import 'package:beecount/cloud/cloud_session_manager.dart'; import 'package:beecount/pages/cloud/cloud_service_page.dart'; import 'package:beecount/providers/sync_providers.dart'; @@ -33,8 +34,8 @@ void main() { beecountCloudConfigProvider.overrideWith((ref) async => config), webdavConfigProvider.overrideWith((ref) async => null), s3ConfigProvider.overrideWith((ref) async => null), - beecountCloudProviderInstance - .overrideWith((ref) async => fakeProvider), + activeCloudServicesProvider.overrideWith((ref) async => CloudSession( + config, (provider: fakeProvider, auth: fakeProvider.auth))), ], child: MaterialApp( localizationsDelegates: AppLocalizations.localizationsDelegates, diff --git a/test/providers/cloud_session_auth_test.dart b/test/providers/cloud_session_auth_test.dart new file mode 100644 index 000000000..fc869e502 --- /dev/null +++ b/test/providers/cloud_session_auth_test.dart @@ -0,0 +1,122 @@ +import 'dart:convert'; + +import 'package:beecount/cloud/cloud_session_manager.dart'; +import 'package:beecount/providers/sync_providers.dart'; +import 'package:flutter_cloud_sync/flutter_cloud_sync.dart'; +import 'package:flutter_riverpod/flutter_riverpod.dart'; +import 'package:flutter_test/flutter_test.dart'; +import 'package:http/http.dart' as http; +import 'package:http/testing.dart'; +import 'package:shared_preferences/shared_preferences.dart'; +import 'package:package_info_plus/package_info_plus.dart'; + +import '../cloud/sync/_fakes/fake_beecount_cloud_provider.dart'; + +class AuthProvider extends FakeBeeCountCloudProvider { + AuthProvider(this.sessionAuth); + final BeeCountCloudAuthService sessionAuth; + @override + CloudAuthService get auth => sessionAuth; + @override + Future dispose() async => sessionAuth.dispose(); +} + +void main() { + TestWidgetsFlutterBinding.ensureInitialized(); + setUp(() { + SharedPreferences.setMockInitialValues({}); + PackageInfo.setMockInitialValues( + appName: 'test', + packageName: 'test', + version: '1', + buildNumber: '1', + buildSignature: 'test'); + }); + + Map tokens(String user) => { + 'user': {'id': user, 'email': '$user@example.com'}, + 'access_token': 'test-access', + 'refresh_token': 'test-refresh', + 'expires_in': 3600, + 'device_id': 'test-device', + }; + + test('UI initialized before 2FA sees login immediately without invalidation', + () async { + final requests = []; + final auth = BeeCountCloudAuthService( + baseUrl: 'https://cloud.example', + apiPrefix: '/api/v1', + httpClient: MockClient((request) async { + requests.add(request.url.path); + return http.Response( + jsonEncode(request.url.path.endsWith('/login') + ? {'requires_2fa': true, 'challenge_token': 'test-challenge'} + : tokens('owner')), + 200); + }), + twoFactorHandler: (challenge) async => + await challenge.verify('totp', '123456') == null, + ); + final provider = AuthProvider(auth); + final manager = CloudSessionManager( + factory: (_) async => (provider: provider, auth: auth)); + final container = ProviderContainer(overrides: [ + cloudSessionManagerProvider.overrideWithValue(manager), + activeCloudConfigProvider.overrideWith((ref) async => + const CloudServiceConfig( + type: CloudBackendType.beecountCloud, + name: 'test', + beecountCloudBaseUrl: 'https://cloud.example')), + ]); + addTearDown(() async { + container.dispose(); + await manager.dispose(); + }); + final uiAuth = await container.read(authServiceProvider.future); + expect(await uiAuth.currentUser, isNull); + final cloud = await container.read(beecountCloudProviderInstance.future); + await cloud!.auth + .signInWithEmail(email: 'owner@example.com', password: 'test'); + expect((await uiAuth.currentUser)?.id, 'owner'); + container.invalidate(authServiceProvider); + container.invalidate(beecountCloudProviderInstance); + final rebuilt = await container.read(authServiceProvider.future); + expect((await rebuilt.currentUser)?.id, 'owner'); + expect(requests, ['/api/v1/auth/login', '/api/v1/auth/2fa/verify']); + }); + + test('temporary BeeCount auth neither reads nor overwrites persisted session', + () async { + BeeCountCloudAuthService makeAuth(bool persist, String user) => + BeeCountCloudAuthService( + baseUrl: 'https://cloud.example', + apiPrefix: '/api/v1', + persistSession: persist, + httpClient: MockClient( + (_) async => http.Response(jsonEncode(tokens(user)), 200)), + ); + final active = makeAuth(true, 'owner'); + final draft = makeAuth(false, 'draft'); + addTearDown(active.dispose); + addTearDown(draft.dispose); + await active.signInWithEmail(email: 'owner@example.com', password: 'test'); + final prefs = await SharedPreferences.getInstance(); + final before = {for (final key in prefs.getKeys()) key: prefs.get(key)}; + await draft.initialize(); + expect(await draft.currentUser, isNull); + await draft.signInWithEmail(email: 'draft@example.com', password: 'test'); + expect((await draft.currentUser)?.id, 'draft'); + expect({for (final key in prefs.getKeys()) key: prefs.get(key)}, before); + expect((await active.currentUser)?.id, 'owner'); + final otherAccount = makeAuth(true, 'draft'); + final sameAccount = makeAuth(true, 'owner'); + addTearDown(otherAccount.dispose); + addTearDown(sameAccount.dispose); + await otherAccount.initialize(expectedEmail: 'draft@example.com'); + expect(await otherAccount.currentUser, isNull); + await sameAccount.initialize(expectedEmail: ' OWNER@example.com '); + expect((await sameAccount.currentUser)?.id, 'owner'); + expect({for (final key in prefs.getKeys()) key: prefs.get(key)}, before); + }); +} diff --git a/test/providers/cloud_session_engine_test.dart b/test/providers/cloud_session_engine_test.dart new file mode 100644 index 000000000..c0278fce7 --- /dev/null +++ b/test/providers/cloud_session_engine_test.dart @@ -0,0 +1,56 @@ +import 'package:drift/native.dart'; +import 'package:flutter_cloud_sync/flutter_cloud_sync.dart'; +import 'package:flutter_riverpod/flutter_riverpod.dart'; +import 'package:flutter_test/flutter_test.dart'; +import 'package:shared_preferences/shared_preferences.dart'; +import 'package:beecount/cloud/cloud_session_manager.dart'; +import 'package:beecount/cloud/sync/sync_providers.dart'; +import 'package:beecount/data/db.dart'; +import 'package:beecount/data/repositories/local/local_repository.dart'; +import 'package:beecount/providers/database_providers.dart'; +import 'package:beecount/providers/sync_providers.dart'; +import '../cloud/sync/_fakes/fake_beecount_cloud_provider.dart'; + +void main() { + TestWidgetsFlutterBinding.ensureInitialized(); + test('configuration switch closes the cached engine and rejects its old key', + () async { + SharedPreferences.setMockInitialValues({}); + final db = BeeDatabase.forTesting(NativeDatabase.memory()); + var config = const CloudServiceConfig( + type: CloudBackendType.beecountCloud, + name: 'A', + beecountCloudBaseUrl: 'https://a.example'); + final manager = CloudSessionManager(factory: (_) async { + final p = FakeBeeCountCloudProvider(); + return (provider: p, auth: p.auth); + }); + final container = ProviderContainer(overrides: [ + databaseProvider.overrideWithValue(db), + repositoryProvider.overrideWithValue(LocalRepository(db)), + cloudSessionManagerProvider.overrideWithValue(manager), + activeCloudConfigProvider.overrideWith((ref) async => config), + ]); + addTearDown(() async { + container.dispose(); + await manager.dispose(); + await db.close(); + }); + final oldProvider = + (await container.read(beecountCloudProviderInstance.future))!; + final oldEngine = container.read(syncEngineProvider(oldProvider)); + config = const CloudServiceConfig( + type: CloudBackendType.beecountCloud, + name: 'B', + beecountCloudBaseUrl: 'https://b.example'); + container.invalidate(activeCloudConfigProvider); + final newProvider = + (await container.read(beecountCloudProviderInstance.future))!; + expect(newProvider, isNot(same(oldProvider))); + expect((await oldEngine.sync(ledgerId: '1')).hasError, isTrue); + expect(() => container.read(syncEngineProvider(oldProvider)), + throwsStateError); + expect(container.read(syncEngineProvider(newProvider)), + isNot(same(oldEngine))); + }); +} From 362a17a34061fc1dfe046983fdff2e1c3c426e29 Mon Sep 17 00:00:00 2001 From: Hongkuan Zhou <6771308+tedzhouhk@users.noreply.github.com> Date: Sat, 12 Sep 2026 21:31:34 -0700 Subject: [PATCH 6/6] fix: serialize ledger sync across sessions and retry initialization --- lib/cloud/sync/sync_engine.dart | 14 ++++- lib/cloud/sync/sync_providers.dart | 4 +- lib/pages/cloud/cloud_service_page.dart | 9 +++ .../sync/sync_engine_session_queue_test.dart | 62 +++++++++++++++++++ test/pages/cloud/cloud_service_page_test.dart | 53 ++++++++++++++++ test/providers/cloud_session_engine_test.dart | 4 ++ 6 files changed, 144 insertions(+), 2 deletions(-) create mode 100644 test/cloud/sync/sync_engine_session_queue_test.dart diff --git a/lib/cloud/sync/sync_engine.dart b/lib/cloud/sync/sync_engine.dart index 1b7b23e3d..9101495a1 100644 --- a/lib/cloud/sync/sync_engine.dart +++ b/lib/cloud/sync/sync_engine.dart @@ -540,6 +540,11 @@ class SyncEngine implements app.SyncService { /// 同一会话共用一个 engine;切换会话不能复用旧服务器的 in-flight 结果。 Completer? _syncLedgersInFlight; + // Different engines using the same database must not interleave ledger + // lookup/insert/GC. Queue work, not results: a new session must fetch its own + // server's ledger list after the previous operation has finished. + static final _ledgerSyncTails = Expando>(); + Future syncLedgersFromServer() async { if (_disposed) throw StateError('Cloud session closed'); final existing = _syncLedgersInFlight; @@ -550,7 +555,14 @@ class SyncEngine implements app.SyncService { final completer = Completer(); _syncLedgersInFlight = completer; try { - final n = await _syncLedgersFromServerLocked(); + final previous = _ledgerSyncTails[db] ?? Future.value(); + final work = previous.then((_) async { + if (_disposed) return 0; + return _syncLedgersFromServerLocked(); + }); + _ledgerSyncTails[db] = work.then((_) {}, + onError: (Object error, StackTrace stack) {}); + final n = await work; completer.complete(n); return n; } catch (e, st) { diff --git a/lib/cloud/sync/sync_providers.dart b/lib/cloud/sync/sync_providers.dart index 780ca4a42..c5c36c656 100644 --- a/lib/cloud/sync/sync_providers.dart +++ b/lib/cloud/sync/sync_providers.dart @@ -24,7 +24,9 @@ final changeTrackerProvider = Provider((ref) { /// 装配 callback 后才启动,但 dispose 由 Riverpod GC family entry 时统一触发。 final syncEngineProvider = Provider.family( (ref, provider) { - final session = ref.watch(activeCloudServicesProvider).valueOrNull; + // Loading with the same previous value is not a new session. + final session = ref.watch(activeCloudServicesProvider.select( + (value) => value.valueOrNull)); if (session == null || session.isClosed || !identical(session.services.provider, provider)) { throw StateError('SyncEngine requires the active cloud session'); diff --git a/lib/pages/cloud/cloud_service_page.dart b/lib/pages/cloud/cloud_service_page.dart index 95539d5b1..8a93b5bfb 100644 --- a/lib/pages/cloud/cloud_service_page.dart +++ b/lib/pages/cloud/cloud_service_page.dart @@ -1809,6 +1809,15 @@ class _CloudServicePageState extends ConsumerState { String? errorDetail; try { + if (config.type == CloudBackendType.beecountCloud || + config.type == CloudBackendType.s3) { + final active = ref.read(activeCloudServicesProvider); + // Retry initialization errors only. Keep healthy sessions and an + // already-running initialization shared with all other consumers. + if (active.hasError && !active.isLoading) { + ref.invalidate(activeCloudServicesProvider); + } + } switch (config.type) { case CloudBackendType.local: break; diff --git a/test/cloud/sync/sync_engine_session_queue_test.dart b/test/cloud/sync/sync_engine_session_queue_test.dart new file mode 100644 index 000000000..90fa16e7f --- /dev/null +++ b/test/cloud/sync/sync_engine_session_queue_test.dart @@ -0,0 +1,62 @@ +import 'dart:async'; + +import 'package:beecount/cloud/sync/change_tracker.dart'; +import 'package:beecount/cloud/sync/sync_engine.dart'; +import 'package:beecount/data/db.dart'; +import 'package:beecount/data/repositories/local/local_repository.dart'; +import 'package:drift/native.dart'; +import 'package:flutter_cloud_sync/flutter_cloud_sync.dart'; +import 'package:flutter_test/flutter_test.dart'; +import 'package:shared_preferences/shared_preferences.dart'; + +import '_fakes/fake_beecount_cloud_provider.dart'; + +class GatedProvider extends FakeBeeCountCloudProvider { + final started = Completer(); + final release = Completer(); + int calls = 0; + + @override + Future> readLedgers() async { + calls++; + started.complete(); + await release.future; + return super.readLedgers(); + } +} + +void main() { + TestWidgetsFlutterBinding.ensureInitialized(); + + test('replacement engine waits but fetches its own server result', () async { + SharedPreferences.setMockInitialValues({}); + final db = BeeDatabase.forTesting(NativeDatabase.memory()); + final first = GatedProvider(); + final second = GatedProvider(); + SyncEngine makeEngine(GatedProvider provider) => SyncEngine( + db: db, + provider: provider, + changeTracker: ChangeTracker(db), + repo: LocalRepository(db), + ); + final oldEngine = makeEngine(first); + final newEngine = makeEngine(second); + addTearDown(() async { + oldEngine.dispose(); + newEngine.dispose(); + await db.close(); + }); + final oldWork = oldEngine.syncLedgersFromServer(); + await first.started.future; + oldEngine.dispose(); + final newWork = newEngine.syncLedgersFromServer(); + await Future.delayed(Duration.zero); + expect(second.calls, 0); + first.release.complete(); + await oldWork; + await second.started.future; + expect(second.calls, 1); + second.release.complete(); + await newWork; + }); +} diff --git a/test/pages/cloud/cloud_service_page_test.dart b/test/pages/cloud/cloud_service_page_test.dart index 2cd4a8dd0..cad9fdc8f 100644 --- a/test/pages/cloud/cloud_service_page_test.dart +++ b/test/pages/cloud/cloud_service_page_test.dart @@ -56,4 +56,57 @@ void main() { await tester.tap(find.text('确定')); await tester.pumpAndSettle(); }); + + testWidgets('S3 connection test retries failed initialization once', + (tester) async { + const config = CloudServiceConfig( + type: CloudBackendType.s3, + name: 'S3', + s3Endpoint: 's3.example.com', + s3AccessKey: 'test', + s3SecretKey: 'test', + s3Bucket: 'test', + ); + final fake = FakeBeeCountCloudProvider(); + var attempts = 0; + final manager = CloudSessionManager(factory: (_) async { + if (++attempts == 1) throw StateError('offline'); + return (provider: fake, auth: fake.auth); + }); + final container = ProviderContainer(overrides: [ + cloudSessionManagerProvider.overrideWithValue(manager), + activeCloudConfigProvider.overrideWith((ref) async => config), + s3ConfigProvider.overrideWith((ref) async => config), + beecountCloudConfigProvider.overrideWith((ref) async => null), + webdavConfigProvider.overrideWith((ref) async => null), + ]); + await expectLater(container.read(activeCloudServicesProvider.future), + throwsStateError); + await tester.pumpWidget(UncontrolledProviderScope( + container: container, + child: MaterialApp( + localizationsDelegates: AppLocalizations.localizationsDelegates, + supportedLocales: AppLocalizations.supportedLocales, + locale: const Locale('zh'), + home: const CloudServicePage(), + ), + )); + await tester.pumpAndSettle(); + for (var i = 0; i < 2; i++) { + await tester.tap(find.byIcon(Icons.wifi_find)); + await tester.pump(); + await tester.pump(const Duration(milliseconds: 300)); + expect(attempts, 2); + expect(fake.storageListCallCount, i + 1); + await tester.tap(find.text('确定')); + await tester.pumpAndSettle(); + } + await tester.pumpWidget(const SizedBox()); + // Flush the logger's delayed write timer before fake-async teardown. + await tester.pump(const Duration(seconds: 3)); + container.dispose(); + final closing = manager.dispose(); + await tester.pumpAndSettle(); + await closing; + }); } diff --git a/test/providers/cloud_session_engine_test.dart b/test/providers/cloud_session_engine_test.dart index c0278fce7..f8e6aa7bf 100644 --- a/test/providers/cloud_session_engine_test.dart +++ b/test/providers/cloud_session_engine_test.dart @@ -39,6 +39,10 @@ void main() { final oldProvider = (await container.read(beecountCloudProviderInstance.future))!; final oldEngine = container.read(syncEngineProvider(oldProvider)); + container.invalidate(activeCloudConfigProvider); + expect(await container.read(beecountCloudProviderInstance.future), + same(oldProvider)); + expect(container.read(syncEngineProvider(oldProvider)), same(oldEngine)); config = const CloudServiceConfig( type: CloudBackendType.beecountCloud, name: 'B',