import 'dart:async'; import 'package:connectivity_plus/connectivity_plus.dart'; import 'package:sunny_mochi/core/observability/talker_setup.dart'; import 'package:sunny_mochi/core/sync/sync_status.dart'; import 'package:riverpod_annotation/riverpod_annotation.dart'; part 'sync_service.g.dart'; /// 单条 dirty 记录的同步执行器。各业务模块实现此接口注册到 SyncService。 abstract class SyncTask { /// 唯一标识(用于日志和去重) String get key; /// 执行同步:成功返回 true,失败返回 false(保留 pending 等下轮)。 /// 冲突时由实现方写入 conflictPayload 并返回 false。 Future execute(); } /// Offline-First 后台同步调度器。 /// /// - 监听 [Connectivity] 网络状态变化,恢复联网时触发一轮同步 /// - 业务层通过 [enqueue] 注册 SyncTask /// - 串行执行(避免并发污染服务端) /// /// 详见 docs/flutter-architecture-design.md §五.21。 class SyncService { SyncService({ required Connectivity connectivity, }) : _connectivity = connectivity; final Connectivity _connectivity; final List _queue = []; StreamSubscription>? _connSub; bool _running = false; /// 启动后台监听。在 main.dart 完成 ProviderScope 初始化后调用。 /// /// **R-ROB-1 / 2026-05-10**:幂等 — 重复调用会先 cancel 旧 subscription, /// 防止 hot reload / Widget 重建时累积订阅造成 _drain 多次触发。 Future start() async { await _connSub?.cancel(); _connSub = _connectivity.onConnectivityChanged.listen(_onConnectivity); final initial = await _connectivity.checkConnectivity(); if (_isOnline(initial)) { unawaited(_drain()); } } Future dispose() async { await _connSub?.cancel(); _connSub = null; } /// 业务层注册一个 dirty 任务(同步后由任务自身从队列移除)。 void enqueue(SyncTask task) { _queue.add(task); appTalker.verbose('[Sync] enqueue ${task.key} (queue=${_queue.length})'); unawaited(_drain()); } Future _drain() async { if (_running || _queue.isEmpty) return; _running = true; try { while (_queue.isNotEmpty) { final task = _queue.removeAt(0); try { final ok = await task.execute(); appTalker.info( '[Sync] ${task.key} ${ok ? "✓" : "✗ (status=${SyncStatus.conflict.name} 待解决)"}', ); } on Object catch (e, st) { appTalker.warning('[Sync] ${task.key} 异常:$e', e, st); // 失败任务不重新入队(保持 dirty 在 DB 等下轮启动) } } } finally { _running = false; } } void _onConnectivity(List results) { if (_isOnline(results)) { appTalker.info('[Sync] 检测到联网,触发 drain'); unawaited(_drain()); } } bool _isOnline(List results) => results.any((r) => r != ConnectivityResult.none); } @Riverpod(keepAlive: true) Connectivity connectivity(Ref ref) => Connectivity(); @Riverpod(keepAlive: true) SyncService syncService(Ref ref) { final svc = SyncService(connectivity: ref.watch(connectivityProvider)); ref.onDispose(() => unawaited(svc.dispose())); return svc; }