Files
flutter-template/lib/core/sync/sync_service.dart
T
SkyJourney b8b63282ee fix: 清除 platform-flutter 残留代码,重构文档与开发者体验
P0 代码修复:删除 imSdkAppId、清除内网 IP 与原项目域名、修正 USE_MOCK 默认值
P1 CLAUDE.md 重组:零容忍规则前置,新增 dart-defines 字段对照表
P2 文档优化:先跑再配的 setup 流程,新增 architecture.md 架构图
2026-05-14 13:21:06 +08:00

104 lines
3.2 KiB
Dart
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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<bool> execute();
}
/// Offline-First 后台同步调度器。
///
/// - 监听 [Connectivity] 网络状态变化,恢复联网时触发一轮同步
/// - 业务层通过 [enqueue] 注册 SyncTask
/// - 串行执行(避免并发污染服务端)
///
class SyncService {
SyncService({
required Connectivity connectivity,
}) : _connectivity = connectivity;
final Connectivity _connectivity;
final List<SyncTask> _queue = [];
StreamSubscription<List<ConnectivityResult>>? _connSub;
bool _running = false;
/// 启动后台监听。在 main.dart 完成 ProviderScope 初始化后调用。
///
/// **R-ROB-1 / 2026-05-10**:幂等 — 重复调用会先 cancel 旧 subscription
/// 防止 hot reload / Widget 重建时累积订阅造成 _drain 多次触发。
Future<void> start() async {
await _connSub?.cancel();
_connSub = _connectivity.onConnectivityChanged.listen(_onConnectivity);
final initial = await _connectivity.checkConnectivity();
if (_isOnline(initial)) {
unawaited(_drain());
}
}
Future<void> 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<void> _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<ConnectivityResult> results) {
if (_isOnline(results)) {
appTalker.info('[Sync] 检测到联网,触发 drain');
unawaited(_drain());
}
}
bool _isOnline(List<ConnectivityResult> 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;
}