feat: add local-first private AI digest workflow
Migrate app code into canonical feature slices, add phone-only AI digest scheduling and review, wire local notification/background task support, and cover the flow with tests.
This commit is contained in:
@@ -0,0 +1,10 @@
|
||||
# Sync Slice
|
||||
|
||||
This slice owns queued mutation persistence and backend sync orchestration.
|
||||
|
||||
- `application/`: background runner, coordinator, and trigger logic.
|
||||
- `data/`: sync queue repository and persistence adapters.
|
||||
- `presentation/sync_view.dart`: manual sync UI and diagnostics.
|
||||
|
||||
This slice is where future backend sessions should focus once the client data
|
||||
model is stable enough to sync.
|
||||
@@ -0,0 +1,4 @@
|
||||
# Sync Application
|
||||
|
||||
Sync coordination, background triggers, and orchestration logic live here. Use
|
||||
this folder when you need to change how local changes are pushed or retried.
|
||||
@@ -0,0 +1,84 @@
|
||||
import 'dart:math' as math;
|
||||
|
||||
import 'package:flutter/foundation.dart';
|
||||
import 'package:relationship_saver/features/sync/sync_coordinator.dart';
|
||||
|
||||
typedef SyncNowRunner = Future<SyncRunResult> Function();
|
||||
typedef TriggerGate = Future<bool> Function();
|
||||
|
||||
/// Coordinates auto-triggered sync calls with cooldown and concurrency guards.
|
||||
class SyncAutoTriggerController {
|
||||
SyncAutoTriggerController({
|
||||
required SyncNowRunner syncNow,
|
||||
required DateTime Function() now,
|
||||
TriggerGate? canTrigger,
|
||||
this.minimumInterval = const Duration(minutes: 3),
|
||||
}) : _syncNow = syncNow,
|
||||
_now = now,
|
||||
_canTrigger = canTrigger ?? _alwaysAllowed;
|
||||
|
||||
final SyncNowRunner _syncNow;
|
||||
final DateTime Function() _now;
|
||||
final TriggerGate _canTrigger;
|
||||
final Duration minimumInterval;
|
||||
|
||||
bool _running = false;
|
||||
DateTime? _lastTriggeredAt;
|
||||
int _consecutiveFailures = 0;
|
||||
|
||||
@visibleForTesting
|
||||
bool get isRunning => _running;
|
||||
|
||||
@visibleForTesting
|
||||
DateTime? get lastTriggeredAt => _lastTriggeredAt;
|
||||
|
||||
@visibleForTesting
|
||||
int get consecutiveFailures => _consecutiveFailures;
|
||||
|
||||
/// Triggers `syncNow` when not already running and outside cooldown window.
|
||||
Future<bool> trigger({bool ignoreCooldown = false}) async {
|
||||
if (_running) {
|
||||
return false;
|
||||
}
|
||||
|
||||
final DateTime now = _now();
|
||||
final DateTime? last = _lastTriggeredAt;
|
||||
if (!ignoreCooldown &&
|
||||
last != null &&
|
||||
now.difference(last) < _cooldownForCurrentState()) {
|
||||
return false;
|
||||
}
|
||||
|
||||
if (!await _canTrigger()) {
|
||||
return false;
|
||||
}
|
||||
|
||||
_running = true;
|
||||
_lastTriggeredAt = now;
|
||||
try {
|
||||
final SyncRunResult result = await _syncNow();
|
||||
if (result.success) {
|
||||
_consecutiveFailures = 0;
|
||||
} else {
|
||||
_consecutiveFailures += 1;
|
||||
}
|
||||
return true;
|
||||
} catch (_) {
|
||||
_consecutiveFailures += 1;
|
||||
return true;
|
||||
} finally {
|
||||
_running = false;
|
||||
}
|
||||
}
|
||||
|
||||
Duration _cooldownForCurrentState() {
|
||||
if (_consecutiveFailures <= 0) {
|
||||
return minimumInterval;
|
||||
}
|
||||
final int cappedFailures = math.min(_consecutiveFailures, 4);
|
||||
final int multiplier = 1 << cappedFailures;
|
||||
return minimumInterval * multiplier;
|
||||
}
|
||||
|
||||
static Future<bool> _alwaysAllowed() async => true;
|
||||
}
|
||||
@@ -0,0 +1,107 @@
|
||||
import 'dart:async';
|
||||
|
||||
import 'package:flutter/widgets.dart';
|
||||
import 'package:flutter_riverpod/flutter_riverpod.dart';
|
||||
import 'package:relationship_saver/core/config/app_config.dart';
|
||||
import 'package:relationship_saver/core/network/reachability/network_reachability.dart';
|
||||
import 'package:relationship_saver/core/network/reachability/network_reachability_provider.dart';
|
||||
import 'package:relationship_saver/features/sync/sync_auto_trigger_controller.dart';
|
||||
import 'package:relationship_saver/features/sync/sync_coordinator.dart';
|
||||
|
||||
/// Runs lightweight background sync triggers while authenticated.
|
||||
class SyncBackgroundRunner extends ConsumerStatefulWidget {
|
||||
const SyncBackgroundRunner({required this.child, super.key});
|
||||
|
||||
final Widget child;
|
||||
|
||||
@override
|
||||
ConsumerState<SyncBackgroundRunner> createState() =>
|
||||
_SyncBackgroundRunnerState();
|
||||
}
|
||||
|
||||
class _SyncBackgroundRunnerState extends ConsumerState<SyncBackgroundRunner>
|
||||
with WidgetsBindingObserver {
|
||||
late final SyncAutoTriggerController _controller;
|
||||
Timer? _periodicTimer;
|
||||
StreamSubscription<bool>? _reachabilitySubscription;
|
||||
bool? _lastReachable;
|
||||
|
||||
@override
|
||||
void initState() {
|
||||
super.initState();
|
||||
final reachability = ref.read(networkReachabilityProvider);
|
||||
_controller = SyncAutoTriggerController(
|
||||
syncNow: () => ref.read(syncCoordinatorProvider).syncNow(),
|
||||
now: DateTime.now,
|
||||
minimumInterval: Duration(
|
||||
seconds: AppConfig.backgroundSyncIntervalSeconds,
|
||||
),
|
||||
canTrigger: () async {
|
||||
if (AppConfig.useFakeBackend) {
|
||||
return true;
|
||||
}
|
||||
return reachability.isReachable();
|
||||
},
|
||||
);
|
||||
WidgetsBinding.instance.addObserver(this);
|
||||
_configurePeriodicRunner();
|
||||
_configureReachabilityRunner(reachability);
|
||||
if (AppConfig.enableBackgroundSync) {
|
||||
unawaited(_controller.trigger(ignoreCooldown: true));
|
||||
}
|
||||
}
|
||||
|
||||
@override
|
||||
void didChangeAppLifecycleState(AppLifecycleState state) {
|
||||
if (!AppConfig.enableBackgroundSync) {
|
||||
return;
|
||||
}
|
||||
if (state == AppLifecycleState.resumed) {
|
||||
unawaited(_controller.trigger(ignoreCooldown: true));
|
||||
}
|
||||
}
|
||||
|
||||
@override
|
||||
void dispose() {
|
||||
WidgetsBinding.instance.removeObserver(this);
|
||||
_periodicTimer?.cancel();
|
||||
_reachabilitySubscription?.cancel();
|
||||
super.dispose();
|
||||
}
|
||||
|
||||
@override
|
||||
Widget build(BuildContext context) => widget.child;
|
||||
|
||||
void _configurePeriodicRunner() {
|
||||
if (!AppConfig.enableBackgroundSync) {
|
||||
return;
|
||||
}
|
||||
|
||||
final int seconds = AppConfig.backgroundSyncIntervalSeconds;
|
||||
if (seconds <= 0) {
|
||||
return;
|
||||
}
|
||||
|
||||
_periodicTimer = Timer.periodic(Duration(seconds: seconds), (_) {
|
||||
unawaited(_controller.trigger());
|
||||
});
|
||||
}
|
||||
|
||||
void _configureReachabilityRunner(NetworkReachability reachability) {
|
||||
if (!AppConfig.enableBackgroundSync) {
|
||||
return;
|
||||
}
|
||||
_reachabilitySubscription = reachability.watch().listen((bool reachable) {
|
||||
final bool? previous = _lastReachable;
|
||||
_lastReachable = reachable;
|
||||
if (reachable && previous == false) {
|
||||
unawaited(_controller.trigger(ignoreCooldown: true));
|
||||
}
|
||||
});
|
||||
unawaited(_primeReachabilityState(reachability));
|
||||
}
|
||||
|
||||
Future<void> _primeReachabilityState(NetworkReachability reachability) async {
|
||||
_lastReachable = await reachability.isReachable();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,179 @@
|
||||
import 'package:flutter/foundation.dart';
|
||||
import 'package:flutter_riverpod/flutter_riverpod.dart';
|
||||
import 'package:relationship_saver/features/local/local_repository.dart';
|
||||
import 'package:relationship_saver/features/sync/sync_queue_repository.dart';
|
||||
import 'package:relationship_saver/features/sync/sync_state.dart';
|
||||
import 'package:relationship_saver/integrations/backend/backend_gateway_provider.dart';
|
||||
|
||||
/// Result details for one sync action.
|
||||
@immutable
|
||||
class SyncRunResult {
|
||||
const SyncRunResult({
|
||||
required this.success,
|
||||
required this.operation,
|
||||
required this.message,
|
||||
required this.pendingAfter,
|
||||
this.cursor,
|
||||
this.accepted = 0,
|
||||
this.rejected = 0,
|
||||
this.pulled = 0,
|
||||
this.error,
|
||||
});
|
||||
|
||||
final bool success;
|
||||
final String operation;
|
||||
final String message;
|
||||
final int pendingAfter;
|
||||
final String? cursor;
|
||||
final int accepted;
|
||||
final int rejected;
|
||||
final int pulled;
|
||||
final Object? error;
|
||||
}
|
||||
|
||||
/// Coordinates push/pull sync operations across queue, gateway, and local store.
|
||||
class SyncCoordinator {
|
||||
const SyncCoordinator(this._ref);
|
||||
|
||||
final Ref _ref;
|
||||
|
||||
Future<SyncRunResult> pushPending() async {
|
||||
final SyncState queueState = await _ref.read(
|
||||
syncQueueRepositoryProvider.future,
|
||||
);
|
||||
if (queueState.pendingChanges.isEmpty) {
|
||||
return SyncRunResult(
|
||||
success: true,
|
||||
operation: 'push',
|
||||
message: 'No pending local changes to push.',
|
||||
pendingAfter: 0,
|
||||
cursor: queueState.cursor,
|
||||
);
|
||||
}
|
||||
|
||||
final DateTime now = DateTime.now();
|
||||
final syncQueue = _ref.read(syncQueueRepositoryProvider.notifier);
|
||||
try {
|
||||
final result = await _ref
|
||||
.read(backendGatewayProvider)
|
||||
.pushChanges(
|
||||
changes: queueState.pendingChanges,
|
||||
cursor: queueState.cursor,
|
||||
);
|
||||
await syncQueue.applyPushResult(result, at: now);
|
||||
final SyncState after = await _ref.read(
|
||||
syncQueueRepositoryProvider.future,
|
||||
);
|
||||
return SyncRunResult(
|
||||
success: true,
|
||||
operation: 'push',
|
||||
message:
|
||||
'Push complete. accepted=${result.accepted.length}, rejected=${result.rejected.length}.',
|
||||
accepted: result.accepted.length,
|
||||
rejected: result.rejected.length,
|
||||
pendingAfter: after.pendingChanges.length,
|
||||
cursor: after.cursor,
|
||||
);
|
||||
} catch (error) {
|
||||
await syncQueue.markFailure(error, at: now);
|
||||
final SyncState after = await _ref.read(
|
||||
syncQueueRepositoryProvider.future,
|
||||
);
|
||||
return SyncRunResult(
|
||||
success: false,
|
||||
operation: 'push',
|
||||
message: 'Push failed. Local changes stay queued.',
|
||||
pendingAfter: after.pendingChanges.length,
|
||||
cursor: after.cursor,
|
||||
error: error,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Future<SyncRunResult> pullRemote() async {
|
||||
final SyncState queueState = await _ref.read(
|
||||
syncQueueRepositoryProvider.future,
|
||||
);
|
||||
final DateTime now = DateTime.now();
|
||||
final syncQueue = _ref.read(syncQueueRepositoryProvider.notifier);
|
||||
|
||||
try {
|
||||
final result = await _ref
|
||||
.read(backendGatewayProvider)
|
||||
.pullChanges(cursor: queueState.cursor);
|
||||
await _ref
|
||||
.read(localRepositoryProvider.notifier)
|
||||
.applyRemoteChanges(result.changes);
|
||||
await syncQueue.applyPullResult(result, at: now);
|
||||
final SyncState after = await _ref.read(
|
||||
syncQueueRepositoryProvider.future,
|
||||
);
|
||||
return SyncRunResult(
|
||||
success: true,
|
||||
operation: 'pull',
|
||||
message: 'Pulled ${result.changes.length} change(s) from backend.',
|
||||
pulled: result.changes.length,
|
||||
pendingAfter: after.pendingChanges.length,
|
||||
cursor: after.cursor,
|
||||
);
|
||||
} catch (error) {
|
||||
await syncQueue.markFailure(error, at: now);
|
||||
final SyncState after = await _ref.read(
|
||||
syncQueueRepositoryProvider.future,
|
||||
);
|
||||
return SyncRunResult(
|
||||
success: false,
|
||||
operation: 'pull',
|
||||
message: 'Pull failed. Staying in local-first mode.',
|
||||
pendingAfter: after.pendingChanges.length,
|
||||
cursor: after.cursor,
|
||||
error: error,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Future<SyncRunResult> syncNow() async {
|
||||
final SyncRunResult pushResult = await pushPending();
|
||||
if (!pushResult.success) {
|
||||
return SyncRunResult(
|
||||
success: false,
|
||||
operation: 'sync',
|
||||
message: 'Sync paused at push step. ${pushResult.message}',
|
||||
pendingAfter: pushResult.pendingAfter,
|
||||
cursor: pushResult.cursor,
|
||||
accepted: pushResult.accepted,
|
||||
rejected: pushResult.rejected,
|
||||
error: pushResult.error,
|
||||
);
|
||||
}
|
||||
|
||||
final SyncRunResult pullResult = await pullRemote();
|
||||
if (!pullResult.success) {
|
||||
return SyncRunResult(
|
||||
success: false,
|
||||
operation: 'sync',
|
||||
message: 'Sync paused at pull step. ${pullResult.message}',
|
||||
pendingAfter: pullResult.pendingAfter,
|
||||
cursor: pullResult.cursor,
|
||||
accepted: pushResult.accepted,
|
||||
rejected: pushResult.rejected,
|
||||
error: pullResult.error,
|
||||
);
|
||||
}
|
||||
|
||||
return SyncRunResult(
|
||||
success: true,
|
||||
operation: 'sync',
|
||||
message:
|
||||
'Sync complete. accepted=${pushResult.accepted}, rejected=${pushResult.rejected}, pulled=${pullResult.pulled}.',
|
||||
pendingAfter: pullResult.pendingAfter,
|
||||
cursor: pullResult.cursor,
|
||||
accepted: pushResult.accepted,
|
||||
rejected: pushResult.rejected,
|
||||
pulled: pullResult.pulled,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
final Provider<SyncCoordinator> syncCoordinatorProvider =
|
||||
Provider<SyncCoordinator>(SyncCoordinator.new);
|
||||
@@ -0,0 +1,4 @@
|
||||
# Sync Data
|
||||
|
||||
Sync repositories and storage adapters live here. This is the persistence layer
|
||||
for sync state, separate from the app's main local relationship store.
|
||||
@@ -0,0 +1,4 @@
|
||||
# Sync Data Storage
|
||||
|
||||
This folder contains the concrete persistence adapters for sync-specific state.
|
||||
Keep it storage-focused and free of higher-level retry policy logic.
|
||||
@@ -0,0 +1,18 @@
|
||||
/// Raw persisted payload for sync queue state.
|
||||
class SyncStateRecord {
|
||||
const SyncStateRecord({required this.rawState});
|
||||
|
||||
final String rawState;
|
||||
}
|
||||
|
||||
/// Persistence boundary for sync queue metadata and pending changes.
|
||||
abstract interface class SyncStateStore {
|
||||
/// Reads stored sync state payload.
|
||||
Future<SyncStateRecord?> read();
|
||||
|
||||
/// Writes sync state payload.
|
||||
Future<void> write({required String rawState});
|
||||
|
||||
/// Clears stored sync state payload.
|
||||
Future<void> clear();
|
||||
}
|
||||
@@ -0,0 +1,46 @@
|
||||
import 'package:hive_flutter/hive_flutter.dart';
|
||||
import 'package:relationship_saver/features/sync/data/storage/sync_state_store.dart';
|
||||
|
||||
/// Hive-backed sync state store.
|
||||
class HiveSyncStateStore implements SyncStateStore {
|
||||
static const String _boxName = 'relationship_saver_sync';
|
||||
static const String _stateKey = 'sync_state_json';
|
||||
|
||||
static bool _initialized = false;
|
||||
|
||||
@override
|
||||
Future<void> clear() async {
|
||||
final Box<dynamic> box = await _openBox();
|
||||
await box.delete(_stateKey);
|
||||
}
|
||||
|
||||
@override
|
||||
Future<SyncStateRecord?> read() async {
|
||||
final Box<dynamic> box = await _openBox();
|
||||
final String? rawState = box.get(_stateKey) as String?;
|
||||
if (rawState == null || rawState.isEmpty) {
|
||||
return null;
|
||||
}
|
||||
|
||||
return SyncStateRecord(rawState: rawState);
|
||||
}
|
||||
|
||||
@override
|
||||
Future<void> write({required String rawState}) async {
|
||||
final Box<dynamic> box = await _openBox();
|
||||
await box.put(_stateKey, rawState);
|
||||
}
|
||||
|
||||
Future<Box<dynamic>> _openBox() async {
|
||||
if (!_initialized) {
|
||||
await Hive.initFlutter();
|
||||
_initialized = true;
|
||||
}
|
||||
|
||||
if (!Hive.isBoxOpen(_boxName)) {
|
||||
return Hive.openBox<dynamic>(_boxName);
|
||||
}
|
||||
|
||||
return Hive.box<dynamic>(_boxName);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,19 @@
|
||||
import 'package:relationship_saver/features/sync/data/storage/sync_state_store.dart';
|
||||
|
||||
/// In-memory sync state store for tests.
|
||||
class InMemorySyncStateStore implements SyncStateStore {
|
||||
SyncStateRecord? _record;
|
||||
|
||||
@override
|
||||
Future<void> clear() async {
|
||||
_record = null;
|
||||
}
|
||||
|
||||
@override
|
||||
Future<SyncStateRecord?> read() async => _record;
|
||||
|
||||
@override
|
||||
Future<void> write({required String rawState}) async {
|
||||
_record = SyncStateRecord(rawState: rawState);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,14 @@
|
||||
import 'package:flutter_riverpod/flutter_riverpod.dart';
|
||||
import 'package:relationship_saver/core/config/app_config.dart';
|
||||
import 'package:relationship_saver/features/sync/data/storage/sync_state_store.dart';
|
||||
import 'package:relationship_saver/features/sync/data/storage/sync_state_store_hive.dart';
|
||||
import 'package:relationship_saver/features/sync/data/storage/sync_state_store_shared_prefs.dart';
|
||||
|
||||
/// Chooses sync queue persistence backend.
|
||||
final Provider<SyncStateStore> syncStateStoreProvider =
|
||||
Provider<SyncStateStore>((Ref ref) {
|
||||
if (AppConfig.useHiveLocalDb) {
|
||||
return HiveSyncStateStore();
|
||||
}
|
||||
return SharedPrefsSyncStateStore();
|
||||
});
|
||||
@@ -0,0 +1,32 @@
|
||||
import 'package:relationship_saver/features/sync/data/storage/sync_state_store.dart';
|
||||
import 'package:shared_preferences/shared_preferences.dart';
|
||||
|
||||
/// shared_preferences-backed sync state store.
|
||||
class SharedPrefsSyncStateStore implements SyncStateStore {
|
||||
SharedPrefsSyncStateStore({this.stateKey = 'sync_queue_state_v1'});
|
||||
|
||||
final String stateKey;
|
||||
|
||||
@override
|
||||
Future<void> clear() async {
|
||||
final SharedPreferences prefs = await SharedPreferences.getInstance();
|
||||
await prefs.remove(stateKey);
|
||||
}
|
||||
|
||||
@override
|
||||
Future<SyncStateRecord?> read() async {
|
||||
final SharedPreferences prefs = await SharedPreferences.getInstance();
|
||||
final String? rawState = prefs.getString(stateKey);
|
||||
if (rawState == null || rawState.isEmpty) {
|
||||
return null;
|
||||
}
|
||||
|
||||
return SyncStateRecord(rawState: rawState);
|
||||
}
|
||||
|
||||
@override
|
||||
Future<void> write({required String rawState}) async {
|
||||
final SharedPreferences prefs = await SharedPreferences.getInstance();
|
||||
await prefs.setString(stateKey, rawState);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,220 @@
|
||||
import 'dart:convert';
|
||||
|
||||
import 'package:flutter_riverpod/flutter_riverpod.dart';
|
||||
import 'package:relationship_saver/core/config/app_config.dart';
|
||||
import 'package:relationship_saver/features/sync/data/storage/sync_state_store.dart';
|
||||
import 'package:relationship_saver/features/sync/data/storage/sync_state_store_hive.dart';
|
||||
import 'package:relationship_saver/features/sync/data/storage/sync_state_store_provider.dart';
|
||||
import 'package:relationship_saver/features/sync/data/storage/sync_state_store_shared_prefs.dart';
|
||||
import 'package:relationship_saver/features/sync/sync_state.dart';
|
||||
import 'package:relationship_saver/integrations/backend/models/backend_models.dart';
|
||||
|
||||
/// Persists queued local mutations and sync cursor metadata.
|
||||
class SyncQueueRepository extends AsyncNotifier<SyncState> {
|
||||
@override
|
||||
Future<SyncState> build() async {
|
||||
final SyncStateStore store = ref.read(syncStateStoreProvider);
|
||||
SyncStateRecord? record = await store.read();
|
||||
|
||||
if ((record == null || record.rawState.isEmpty) &&
|
||||
AppConfig.useHiveLocalDb &&
|
||||
store is HiveSyncStateStore) {
|
||||
final SharedPrefsSyncStateStore legacyStore = SharedPrefsSyncStateStore();
|
||||
final SyncStateRecord? legacy = await legacyStore.read();
|
||||
if (legacy != null && legacy.rawState.isNotEmpty) {
|
||||
await store.write(rawState: legacy.rawState);
|
||||
await legacyStore.clear();
|
||||
record = await store.read();
|
||||
}
|
||||
}
|
||||
|
||||
if (record == null || record.rawState.isEmpty) {
|
||||
return SyncState.empty;
|
||||
}
|
||||
|
||||
try {
|
||||
final Map<String, dynamic> json =
|
||||
jsonDecode(record.rawState) as Map<String, dynamic>;
|
||||
return SyncState.fromJson(json);
|
||||
} on FormatException {
|
||||
await store.clear();
|
||||
return SyncState.empty;
|
||||
}
|
||||
}
|
||||
|
||||
Future<void> enqueue(ChangeEnvelope change) async {
|
||||
final SyncState current = await _currentState();
|
||||
final List<ChangeEnvelope> pending = <ChangeEnvelope>[
|
||||
change,
|
||||
...current.pendingChanges,
|
||||
];
|
||||
await _setState(current.copyWith(pendingChanges: pending));
|
||||
}
|
||||
|
||||
Future<void> enqueueAll(Iterable<ChangeEnvelope> changes) async {
|
||||
final List<ChangeEnvelope> values = changes.toList(growable: false);
|
||||
if (values.isEmpty) {
|
||||
return;
|
||||
}
|
||||
|
||||
final SyncState current = await _currentState();
|
||||
final List<ChangeEnvelope> pending = <ChangeEnvelope>[
|
||||
...values,
|
||||
...current.pendingChanges,
|
||||
];
|
||||
await _setState(current.copyWith(pendingChanges: pending));
|
||||
}
|
||||
|
||||
Future<void> applyPushResult(SyncPushResult result, {DateTime? at}) async {
|
||||
final SyncState current = await _currentState();
|
||||
final DateTime now = at ?? DateTime.now();
|
||||
final Set<String> rejectedMutationIds = <String>{
|
||||
...result.rejected.map(
|
||||
(MutationRejection rejection) => rejection.clientMutationId,
|
||||
),
|
||||
};
|
||||
final Set<String> completedMutationIds = <String>{
|
||||
...result.accepted.map((MutationAck ack) => ack.clientMutationId),
|
||||
...rejectedMutationIds,
|
||||
};
|
||||
final List<ChangeEnvelope> rejectedChanges = current.pendingChanges
|
||||
.where(
|
||||
(ChangeEnvelope change) =>
|
||||
rejectedMutationIds.contains(change.clientMutationId),
|
||||
)
|
||||
.toList(growable: false);
|
||||
final List<ChangeEnvelope> pending = current.pendingChanges
|
||||
.where(
|
||||
(ChangeEnvelope change) =>
|
||||
!completedMutationIds.contains(change.clientMutationId),
|
||||
)
|
||||
.toList(growable: false);
|
||||
final List<ChangeEnvelope> dedupedRejectedChanges = rejectedChanges
|
||||
.where(
|
||||
(ChangeEnvelope change) => !pending.any(
|
||||
(ChangeEnvelope pendingChange) =>
|
||||
pendingChange.clientMutationId == change.clientMutationId,
|
||||
),
|
||||
)
|
||||
.toList(growable: false);
|
||||
|
||||
final String? rejectionMessage = result.rejected.isEmpty
|
||||
? null
|
||||
: 'Push rejected ${result.rejected.length} change(s). Fix locally and retry.';
|
||||
|
||||
await _setState(
|
||||
current.copyWith(
|
||||
cursor: result.cursor,
|
||||
pendingChanges: pending,
|
||||
lastAttemptAt: now,
|
||||
lastSyncAt: now,
|
||||
lastError: rejectionMessage,
|
||||
lastRejected: result.rejected,
|
||||
lastRejectedChanges: dedupedRejectedChanges,
|
||||
),
|
||||
);
|
||||
}
|
||||
|
||||
/// Requeues rejected changes to pending sync list for manual retry.
|
||||
Future<void> requeueRejectedChanges() async {
|
||||
final SyncState current = await _currentState();
|
||||
if (current.lastRejectedChanges.isEmpty) {
|
||||
return;
|
||||
}
|
||||
|
||||
final List<ChangeEnvelope> pending = <ChangeEnvelope>[
|
||||
...current.lastRejectedChanges,
|
||||
...current.pendingChanges,
|
||||
];
|
||||
final List<ChangeEnvelope> dedupedPending = _dedupeByMutationId(pending);
|
||||
await _setState(
|
||||
current.copyWith(
|
||||
pendingChanges: dedupedPending,
|
||||
lastRejected: const <MutationRejection>[],
|
||||
lastRejectedChanges: const <ChangeEnvelope>[],
|
||||
lastError: null,
|
||||
),
|
||||
);
|
||||
}
|
||||
|
||||
Future<void> applyPullResult(SyncPullResult result, {DateTime? at}) async {
|
||||
final SyncState current = await _currentState();
|
||||
final DateTime now = at ?? DateTime.now();
|
||||
await _setState(
|
||||
current.copyWith(
|
||||
cursor: result.cursor,
|
||||
lastAttemptAt: now,
|
||||
lastSyncAt: now,
|
||||
lastError: null,
|
||||
lastRejected: const <MutationRejection>[],
|
||||
lastRejectedChanges: const <ChangeEnvelope>[],
|
||||
),
|
||||
);
|
||||
}
|
||||
|
||||
Future<void> markFailure(Object error, {DateTime? at}) async {
|
||||
final SyncState current = await _currentState();
|
||||
final DateTime now = at ?? DateTime.now();
|
||||
await _setState(current.copyWith(lastAttemptAt: now, lastError: '$error'));
|
||||
}
|
||||
|
||||
Future<void> clearQueue() async {
|
||||
final SyncState current = await _currentState();
|
||||
await _setState(
|
||||
current.copyWith(
|
||||
pendingChanges: const <ChangeEnvelope>[],
|
||||
lastRejected: const <MutationRejection>[],
|
||||
lastRejectedChanges: const <ChangeEnvelope>[],
|
||||
),
|
||||
);
|
||||
}
|
||||
|
||||
/// Clears the latest rejection payload shown to the user.
|
||||
Future<void> clearRejections() async {
|
||||
final SyncState current = await _currentState();
|
||||
final String? lastError =
|
||||
current.lastError != null &&
|
||||
current.lastError!.startsWith('Push rejected')
|
||||
? null
|
||||
: current.lastError;
|
||||
await _setState(
|
||||
current.copyWith(
|
||||
lastRejected: const <MutationRejection>[],
|
||||
lastRejectedChanges: const <ChangeEnvelope>[],
|
||||
lastError: lastError,
|
||||
),
|
||||
);
|
||||
}
|
||||
|
||||
List<ChangeEnvelope> _dedupeByMutationId(List<ChangeEnvelope> changes) {
|
||||
final Set<String> seen = <String>{};
|
||||
final List<ChangeEnvelope> deduped = <ChangeEnvelope>[];
|
||||
for (final ChangeEnvelope change in changes) {
|
||||
if (seen.add(change.clientMutationId)) {
|
||||
deduped.add(change);
|
||||
}
|
||||
}
|
||||
return deduped;
|
||||
}
|
||||
|
||||
Future<SyncState> _currentState() async {
|
||||
final SyncState? value = state.asData?.value;
|
||||
if (value != null) {
|
||||
return value;
|
||||
}
|
||||
return future;
|
||||
}
|
||||
|
||||
Future<void> _setState(SyncState next) async {
|
||||
state = AsyncData<SyncState>(next);
|
||||
await ref
|
||||
.read(syncStateStoreProvider)
|
||||
.write(rawState: jsonEncode(next.toJson()));
|
||||
}
|
||||
}
|
||||
|
||||
final AsyncNotifierProvider<SyncQueueRepository, SyncState>
|
||||
syncQueueRepositoryProvider =
|
||||
AsyncNotifierProvider<SyncQueueRepository, SyncState>(
|
||||
SyncQueueRepository.new,
|
||||
);
|
||||
@@ -0,0 +1,4 @@
|
||||
# Sync Presentation
|
||||
|
||||
Sync status screens and user-facing repair flows belong here. Use this folder for
|
||||
visibility and control, not for queue or transport internals.
|
||||
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,4 @@
|
||||
# Legacy Sync Storage Exports
|
||||
|
||||
These files are compatibility exports. The canonical sync storage code now lives
|
||||
under `lib/features/sync/data/storage/`, so prefer editing that folder instead.
|
||||
@@ -1,18 +1,2 @@
|
||||
/// Raw persisted payload for sync queue state.
|
||||
class SyncStateRecord {
|
||||
const SyncStateRecord({required this.rawState});
|
||||
|
||||
final String rawState;
|
||||
}
|
||||
|
||||
/// Persistence boundary for sync queue metadata and pending changes.
|
||||
abstract interface class SyncStateStore {
|
||||
/// Reads stored sync state payload.
|
||||
Future<SyncStateRecord?> read();
|
||||
|
||||
/// Writes sync state payload.
|
||||
Future<void> write({required String rawState});
|
||||
|
||||
/// Clears stored sync state payload.
|
||||
Future<void> clear();
|
||||
}
|
||||
// Legacy compatibility export for the sync state store contract.
|
||||
export 'package:relationship_saver/features/sync/data/storage/sync_state_store.dart';
|
||||
|
||||
@@ -1,46 +1,2 @@
|
||||
import 'package:hive_flutter/hive_flutter.dart';
|
||||
import 'package:relationship_saver/features/sync/storage/sync_state_store.dart';
|
||||
|
||||
/// Hive-backed sync state store.
|
||||
class HiveSyncStateStore implements SyncStateStore {
|
||||
static const String _boxName = 'relationship_saver_sync';
|
||||
static const String _stateKey = 'sync_state_json';
|
||||
|
||||
static bool _initialized = false;
|
||||
|
||||
@override
|
||||
Future<void> clear() async {
|
||||
final Box<dynamic> box = await _openBox();
|
||||
await box.delete(_stateKey);
|
||||
}
|
||||
|
||||
@override
|
||||
Future<SyncStateRecord?> read() async {
|
||||
final Box<dynamic> box = await _openBox();
|
||||
final String? rawState = box.get(_stateKey) as String?;
|
||||
if (rawState == null || rawState.isEmpty) {
|
||||
return null;
|
||||
}
|
||||
|
||||
return SyncStateRecord(rawState: rawState);
|
||||
}
|
||||
|
||||
@override
|
||||
Future<void> write({required String rawState}) async {
|
||||
final Box<dynamic> box = await _openBox();
|
||||
await box.put(_stateKey, rawState);
|
||||
}
|
||||
|
||||
Future<Box<dynamic>> _openBox() async {
|
||||
if (!_initialized) {
|
||||
await Hive.initFlutter();
|
||||
_initialized = true;
|
||||
}
|
||||
|
||||
if (!Hive.isBoxOpen(_boxName)) {
|
||||
return Hive.openBox<dynamic>(_boxName);
|
||||
}
|
||||
|
||||
return Hive.box<dynamic>(_boxName);
|
||||
}
|
||||
}
|
||||
// Legacy compatibility export for the Hive sync state store.
|
||||
export 'package:relationship_saver/features/sync/data/storage/sync_state_store_hive.dart';
|
||||
|
||||
@@ -1,19 +1,2 @@
|
||||
import 'package:relationship_saver/features/sync/storage/sync_state_store.dart';
|
||||
|
||||
/// In-memory sync state store for tests.
|
||||
class InMemorySyncStateStore implements SyncStateStore {
|
||||
SyncStateRecord? _record;
|
||||
|
||||
@override
|
||||
Future<void> clear() async {
|
||||
_record = null;
|
||||
}
|
||||
|
||||
@override
|
||||
Future<SyncStateRecord?> read() async => _record;
|
||||
|
||||
@override
|
||||
Future<void> write({required String rawState}) async {
|
||||
_record = SyncStateRecord(rawState: rawState);
|
||||
}
|
||||
}
|
||||
// Legacy compatibility export for the in-memory sync state store.
|
||||
export 'package:relationship_saver/features/sync/data/storage/sync_state_store_in_memory.dart';
|
||||
|
||||
@@ -1,14 +1,2 @@
|
||||
import 'package:flutter_riverpod/flutter_riverpod.dart';
|
||||
import 'package:relationship_saver/core/config/app_config.dart';
|
||||
import 'package:relationship_saver/features/sync/storage/sync_state_store.dart';
|
||||
import 'package:relationship_saver/features/sync/storage/sync_state_store_hive.dart';
|
||||
import 'package:relationship_saver/features/sync/storage/sync_state_store_shared_prefs.dart';
|
||||
|
||||
/// Chooses sync queue persistence backend.
|
||||
final Provider<SyncStateStore> syncStateStoreProvider =
|
||||
Provider<SyncStateStore>((Ref ref) {
|
||||
if (AppConfig.useHiveLocalDb) {
|
||||
return HiveSyncStateStore();
|
||||
}
|
||||
return SharedPrefsSyncStateStore();
|
||||
});
|
||||
// Legacy compatibility export for the sync state store provider.
|
||||
export 'package:relationship_saver/features/sync/data/storage/sync_state_store_provider.dart';
|
||||
|
||||
@@ -1,32 +1,2 @@
|
||||
import 'package:relationship_saver/features/sync/storage/sync_state_store.dart';
|
||||
import 'package:shared_preferences/shared_preferences.dart';
|
||||
|
||||
/// shared_preferences-backed sync state store.
|
||||
class SharedPrefsSyncStateStore implements SyncStateStore {
|
||||
SharedPrefsSyncStateStore({this.stateKey = 'sync_queue_state_v1'});
|
||||
|
||||
final String stateKey;
|
||||
|
||||
@override
|
||||
Future<void> clear() async {
|
||||
final SharedPreferences prefs = await SharedPreferences.getInstance();
|
||||
await prefs.remove(stateKey);
|
||||
}
|
||||
|
||||
@override
|
||||
Future<SyncStateRecord?> read() async {
|
||||
final SharedPreferences prefs = await SharedPreferences.getInstance();
|
||||
final String? rawState = prefs.getString(stateKey);
|
||||
if (rawState == null || rawState.isEmpty) {
|
||||
return null;
|
||||
}
|
||||
|
||||
return SyncStateRecord(rawState: rawState);
|
||||
}
|
||||
|
||||
@override
|
||||
Future<void> write({required String rawState}) async {
|
||||
final SharedPreferences prefs = await SharedPreferences.getInstance();
|
||||
await prefs.setString(stateKey, rawState);
|
||||
}
|
||||
}
|
||||
// Legacy compatibility export for the shared preferences sync state store.
|
||||
export 'package:relationship_saver/features/sync/data/storage/sync_state_store_shared_prefs.dart';
|
||||
|
||||
@@ -1,84 +1,2 @@
|
||||
import 'dart:math' as math;
|
||||
|
||||
import 'package:flutter/foundation.dart';
|
||||
import 'package:relationship_saver/features/sync/sync_coordinator.dart';
|
||||
|
||||
typedef SyncNowRunner = Future<SyncRunResult> Function();
|
||||
typedef TriggerGate = Future<bool> Function();
|
||||
|
||||
/// Coordinates auto-triggered sync calls with cooldown and concurrency guards.
|
||||
class SyncAutoTriggerController {
|
||||
SyncAutoTriggerController({
|
||||
required SyncNowRunner syncNow,
|
||||
required DateTime Function() now,
|
||||
TriggerGate? canTrigger,
|
||||
this.minimumInterval = const Duration(minutes: 3),
|
||||
}) : _syncNow = syncNow,
|
||||
_now = now,
|
||||
_canTrigger = canTrigger ?? _alwaysAllowed;
|
||||
|
||||
final SyncNowRunner _syncNow;
|
||||
final DateTime Function() _now;
|
||||
final TriggerGate _canTrigger;
|
||||
final Duration minimumInterval;
|
||||
|
||||
bool _running = false;
|
||||
DateTime? _lastTriggeredAt;
|
||||
int _consecutiveFailures = 0;
|
||||
|
||||
@visibleForTesting
|
||||
bool get isRunning => _running;
|
||||
|
||||
@visibleForTesting
|
||||
DateTime? get lastTriggeredAt => _lastTriggeredAt;
|
||||
|
||||
@visibleForTesting
|
||||
int get consecutiveFailures => _consecutiveFailures;
|
||||
|
||||
/// Triggers `syncNow` when not already running and outside cooldown window.
|
||||
Future<bool> trigger({bool ignoreCooldown = false}) async {
|
||||
if (_running) {
|
||||
return false;
|
||||
}
|
||||
|
||||
final DateTime now = _now();
|
||||
final DateTime? last = _lastTriggeredAt;
|
||||
if (!ignoreCooldown &&
|
||||
last != null &&
|
||||
now.difference(last) < _cooldownForCurrentState()) {
|
||||
return false;
|
||||
}
|
||||
|
||||
if (!await _canTrigger()) {
|
||||
return false;
|
||||
}
|
||||
|
||||
_running = true;
|
||||
_lastTriggeredAt = now;
|
||||
try {
|
||||
final SyncRunResult result = await _syncNow();
|
||||
if (result.success) {
|
||||
_consecutiveFailures = 0;
|
||||
} else {
|
||||
_consecutiveFailures += 1;
|
||||
}
|
||||
return true;
|
||||
} catch (_) {
|
||||
_consecutiveFailures += 1;
|
||||
return true;
|
||||
} finally {
|
||||
_running = false;
|
||||
}
|
||||
}
|
||||
|
||||
Duration _cooldownForCurrentState() {
|
||||
if (_consecutiveFailures <= 0) {
|
||||
return minimumInterval;
|
||||
}
|
||||
final int cappedFailures = math.min(_consecutiveFailures, 4);
|
||||
final int multiplier = 1 << cappedFailures;
|
||||
return minimumInterval * multiplier;
|
||||
}
|
||||
|
||||
static Future<bool> _alwaysAllowed() async => true;
|
||||
}
|
||||
// Legacy compatibility export for the sync auto-trigger controller.
|
||||
export 'package:relationship_saver/features/sync/application/sync_auto_trigger_controller.dart';
|
||||
|
||||
@@ -1,107 +1,2 @@
|
||||
import 'dart:async';
|
||||
|
||||
import 'package:flutter/widgets.dart';
|
||||
import 'package:flutter_riverpod/flutter_riverpod.dart';
|
||||
import 'package:relationship_saver/core/config/app_config.dart';
|
||||
import 'package:relationship_saver/core/network/reachability/network_reachability.dart';
|
||||
import 'package:relationship_saver/core/network/reachability/network_reachability_provider.dart';
|
||||
import 'package:relationship_saver/features/sync/sync_auto_trigger_controller.dart';
|
||||
import 'package:relationship_saver/features/sync/sync_coordinator.dart';
|
||||
|
||||
/// Runs lightweight background sync triggers while authenticated.
|
||||
class SyncBackgroundRunner extends ConsumerStatefulWidget {
|
||||
const SyncBackgroundRunner({required this.child, super.key});
|
||||
|
||||
final Widget child;
|
||||
|
||||
@override
|
||||
ConsumerState<SyncBackgroundRunner> createState() =>
|
||||
_SyncBackgroundRunnerState();
|
||||
}
|
||||
|
||||
class _SyncBackgroundRunnerState extends ConsumerState<SyncBackgroundRunner>
|
||||
with WidgetsBindingObserver {
|
||||
late final SyncAutoTriggerController _controller;
|
||||
Timer? _periodicTimer;
|
||||
StreamSubscription<bool>? _reachabilitySubscription;
|
||||
bool? _lastReachable;
|
||||
|
||||
@override
|
||||
void initState() {
|
||||
super.initState();
|
||||
final reachability = ref.read(networkReachabilityProvider);
|
||||
_controller = SyncAutoTriggerController(
|
||||
syncNow: () => ref.read(syncCoordinatorProvider).syncNow(),
|
||||
now: DateTime.now,
|
||||
minimumInterval: Duration(
|
||||
seconds: AppConfig.backgroundSyncIntervalSeconds,
|
||||
),
|
||||
canTrigger: () async {
|
||||
if (AppConfig.useFakeBackend) {
|
||||
return true;
|
||||
}
|
||||
return reachability.isReachable();
|
||||
},
|
||||
);
|
||||
WidgetsBinding.instance.addObserver(this);
|
||||
_configurePeriodicRunner();
|
||||
_configureReachabilityRunner(reachability);
|
||||
if (AppConfig.enableBackgroundSync) {
|
||||
unawaited(_controller.trigger(ignoreCooldown: true));
|
||||
}
|
||||
}
|
||||
|
||||
@override
|
||||
void didChangeAppLifecycleState(AppLifecycleState state) {
|
||||
if (!AppConfig.enableBackgroundSync) {
|
||||
return;
|
||||
}
|
||||
if (state == AppLifecycleState.resumed) {
|
||||
unawaited(_controller.trigger(ignoreCooldown: true));
|
||||
}
|
||||
}
|
||||
|
||||
@override
|
||||
void dispose() {
|
||||
WidgetsBinding.instance.removeObserver(this);
|
||||
_periodicTimer?.cancel();
|
||||
_reachabilitySubscription?.cancel();
|
||||
super.dispose();
|
||||
}
|
||||
|
||||
@override
|
||||
Widget build(BuildContext context) => widget.child;
|
||||
|
||||
void _configurePeriodicRunner() {
|
||||
if (!AppConfig.enableBackgroundSync) {
|
||||
return;
|
||||
}
|
||||
|
||||
final int seconds = AppConfig.backgroundSyncIntervalSeconds;
|
||||
if (seconds <= 0) {
|
||||
return;
|
||||
}
|
||||
|
||||
_periodicTimer = Timer.periodic(Duration(seconds: seconds), (_) {
|
||||
unawaited(_controller.trigger());
|
||||
});
|
||||
}
|
||||
|
||||
void _configureReachabilityRunner(NetworkReachability reachability) {
|
||||
if (!AppConfig.enableBackgroundSync) {
|
||||
return;
|
||||
}
|
||||
_reachabilitySubscription = reachability.watch().listen((bool reachable) {
|
||||
final bool? previous = _lastReachable;
|
||||
_lastReachable = reachable;
|
||||
if (reachable && previous == false) {
|
||||
unawaited(_controller.trigger(ignoreCooldown: true));
|
||||
}
|
||||
});
|
||||
unawaited(_primeReachabilityState(reachability));
|
||||
}
|
||||
|
||||
Future<void> _primeReachabilityState(NetworkReachability reachability) async {
|
||||
_lastReachable = await reachability.isReachable();
|
||||
}
|
||||
}
|
||||
// Legacy compatibility export for the sync background runner.
|
||||
export 'package:relationship_saver/features/sync/application/sync_background_runner.dart';
|
||||
|
||||
@@ -1,179 +1,2 @@
|
||||
import 'package:flutter/foundation.dart';
|
||||
import 'package:flutter_riverpod/flutter_riverpod.dart';
|
||||
import 'package:relationship_saver/features/local/local_repository.dart';
|
||||
import 'package:relationship_saver/features/sync/sync_queue_repository.dart';
|
||||
import 'package:relationship_saver/features/sync/sync_state.dart';
|
||||
import 'package:relationship_saver/integrations/backend/backend_gateway_provider.dart';
|
||||
|
||||
/// Result details for one sync action.
|
||||
@immutable
|
||||
class SyncRunResult {
|
||||
const SyncRunResult({
|
||||
required this.success,
|
||||
required this.operation,
|
||||
required this.message,
|
||||
required this.pendingAfter,
|
||||
this.cursor,
|
||||
this.accepted = 0,
|
||||
this.rejected = 0,
|
||||
this.pulled = 0,
|
||||
this.error,
|
||||
});
|
||||
|
||||
final bool success;
|
||||
final String operation;
|
||||
final String message;
|
||||
final int pendingAfter;
|
||||
final String? cursor;
|
||||
final int accepted;
|
||||
final int rejected;
|
||||
final int pulled;
|
||||
final Object? error;
|
||||
}
|
||||
|
||||
/// Coordinates push/pull sync operations across queue, gateway, and local store.
|
||||
class SyncCoordinator {
|
||||
const SyncCoordinator(this._ref);
|
||||
|
||||
final Ref _ref;
|
||||
|
||||
Future<SyncRunResult> pushPending() async {
|
||||
final SyncState queueState = await _ref.read(
|
||||
syncQueueRepositoryProvider.future,
|
||||
);
|
||||
if (queueState.pendingChanges.isEmpty) {
|
||||
return SyncRunResult(
|
||||
success: true,
|
||||
operation: 'push',
|
||||
message: 'No pending local changes to push.',
|
||||
pendingAfter: 0,
|
||||
cursor: queueState.cursor,
|
||||
);
|
||||
}
|
||||
|
||||
final DateTime now = DateTime.now();
|
||||
final syncQueue = _ref.read(syncQueueRepositoryProvider.notifier);
|
||||
try {
|
||||
final result = await _ref
|
||||
.read(backendGatewayProvider)
|
||||
.pushChanges(
|
||||
changes: queueState.pendingChanges,
|
||||
cursor: queueState.cursor,
|
||||
);
|
||||
await syncQueue.applyPushResult(result, at: now);
|
||||
final SyncState after = await _ref.read(
|
||||
syncQueueRepositoryProvider.future,
|
||||
);
|
||||
return SyncRunResult(
|
||||
success: true,
|
||||
operation: 'push',
|
||||
message:
|
||||
'Push complete. accepted=${result.accepted.length}, rejected=${result.rejected.length}.',
|
||||
accepted: result.accepted.length,
|
||||
rejected: result.rejected.length,
|
||||
pendingAfter: after.pendingChanges.length,
|
||||
cursor: after.cursor,
|
||||
);
|
||||
} catch (error) {
|
||||
await syncQueue.markFailure(error, at: now);
|
||||
final SyncState after = await _ref.read(
|
||||
syncQueueRepositoryProvider.future,
|
||||
);
|
||||
return SyncRunResult(
|
||||
success: false,
|
||||
operation: 'push',
|
||||
message: 'Push failed. Local changes stay queued.',
|
||||
pendingAfter: after.pendingChanges.length,
|
||||
cursor: after.cursor,
|
||||
error: error,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Future<SyncRunResult> pullRemote() async {
|
||||
final SyncState queueState = await _ref.read(
|
||||
syncQueueRepositoryProvider.future,
|
||||
);
|
||||
final DateTime now = DateTime.now();
|
||||
final syncQueue = _ref.read(syncQueueRepositoryProvider.notifier);
|
||||
|
||||
try {
|
||||
final result = await _ref
|
||||
.read(backendGatewayProvider)
|
||||
.pullChanges(cursor: queueState.cursor);
|
||||
await _ref
|
||||
.read(localRepositoryProvider.notifier)
|
||||
.applyRemoteChanges(result.changes);
|
||||
await syncQueue.applyPullResult(result, at: now);
|
||||
final SyncState after = await _ref.read(
|
||||
syncQueueRepositoryProvider.future,
|
||||
);
|
||||
return SyncRunResult(
|
||||
success: true,
|
||||
operation: 'pull',
|
||||
message: 'Pulled ${result.changes.length} change(s) from backend.',
|
||||
pulled: result.changes.length,
|
||||
pendingAfter: after.pendingChanges.length,
|
||||
cursor: after.cursor,
|
||||
);
|
||||
} catch (error) {
|
||||
await syncQueue.markFailure(error, at: now);
|
||||
final SyncState after = await _ref.read(
|
||||
syncQueueRepositoryProvider.future,
|
||||
);
|
||||
return SyncRunResult(
|
||||
success: false,
|
||||
operation: 'pull',
|
||||
message: 'Pull failed. Staying in local-first mode.',
|
||||
pendingAfter: after.pendingChanges.length,
|
||||
cursor: after.cursor,
|
||||
error: error,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Future<SyncRunResult> syncNow() async {
|
||||
final SyncRunResult pushResult = await pushPending();
|
||||
if (!pushResult.success) {
|
||||
return SyncRunResult(
|
||||
success: false,
|
||||
operation: 'sync',
|
||||
message: 'Sync paused at push step. ${pushResult.message}',
|
||||
pendingAfter: pushResult.pendingAfter,
|
||||
cursor: pushResult.cursor,
|
||||
accepted: pushResult.accepted,
|
||||
rejected: pushResult.rejected,
|
||||
error: pushResult.error,
|
||||
);
|
||||
}
|
||||
|
||||
final SyncRunResult pullResult = await pullRemote();
|
||||
if (!pullResult.success) {
|
||||
return SyncRunResult(
|
||||
success: false,
|
||||
operation: 'sync',
|
||||
message: 'Sync paused at pull step. ${pullResult.message}',
|
||||
pendingAfter: pullResult.pendingAfter,
|
||||
cursor: pullResult.cursor,
|
||||
accepted: pushResult.accepted,
|
||||
rejected: pushResult.rejected,
|
||||
error: pullResult.error,
|
||||
);
|
||||
}
|
||||
|
||||
return SyncRunResult(
|
||||
success: true,
|
||||
operation: 'sync',
|
||||
message:
|
||||
'Sync complete. accepted=${pushResult.accepted}, rejected=${pushResult.rejected}, pulled=${pullResult.pulled}.',
|
||||
pendingAfter: pullResult.pendingAfter,
|
||||
cursor: pullResult.cursor,
|
||||
accepted: pushResult.accepted,
|
||||
rejected: pushResult.rejected,
|
||||
pulled: pullResult.pulled,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
final Provider<SyncCoordinator> syncCoordinatorProvider =
|
||||
Provider<SyncCoordinator>(SyncCoordinator.new);
|
||||
// Legacy compatibility export for the sync coordinator.
|
||||
export 'package:relationship_saver/features/sync/application/sync_coordinator.dart';
|
||||
|
||||
@@ -1,220 +1,2 @@
|
||||
import 'dart:convert';
|
||||
|
||||
import 'package:flutter_riverpod/flutter_riverpod.dart';
|
||||
import 'package:relationship_saver/core/config/app_config.dart';
|
||||
import 'package:relationship_saver/features/sync/storage/sync_state_store.dart';
|
||||
import 'package:relationship_saver/features/sync/storage/sync_state_store_hive.dart';
|
||||
import 'package:relationship_saver/features/sync/storage/sync_state_store_provider.dart';
|
||||
import 'package:relationship_saver/features/sync/storage/sync_state_store_shared_prefs.dart';
|
||||
import 'package:relationship_saver/features/sync/sync_state.dart';
|
||||
import 'package:relationship_saver/integrations/backend/models/backend_models.dart';
|
||||
|
||||
/// Persists queued local mutations and sync cursor metadata.
|
||||
class SyncQueueRepository extends AsyncNotifier<SyncState> {
|
||||
@override
|
||||
Future<SyncState> build() async {
|
||||
final SyncStateStore store = ref.read(syncStateStoreProvider);
|
||||
SyncStateRecord? record = await store.read();
|
||||
|
||||
if ((record == null || record.rawState.isEmpty) &&
|
||||
AppConfig.useHiveLocalDb &&
|
||||
store is HiveSyncStateStore) {
|
||||
final SharedPrefsSyncStateStore legacyStore = SharedPrefsSyncStateStore();
|
||||
final SyncStateRecord? legacy = await legacyStore.read();
|
||||
if (legacy != null && legacy.rawState.isNotEmpty) {
|
||||
await store.write(rawState: legacy.rawState);
|
||||
await legacyStore.clear();
|
||||
record = await store.read();
|
||||
}
|
||||
}
|
||||
|
||||
if (record == null || record.rawState.isEmpty) {
|
||||
return SyncState.empty;
|
||||
}
|
||||
|
||||
try {
|
||||
final Map<String, dynamic> json =
|
||||
jsonDecode(record.rawState) as Map<String, dynamic>;
|
||||
return SyncState.fromJson(json);
|
||||
} on FormatException {
|
||||
await store.clear();
|
||||
return SyncState.empty;
|
||||
}
|
||||
}
|
||||
|
||||
Future<void> enqueue(ChangeEnvelope change) async {
|
||||
final SyncState current = await _currentState();
|
||||
final List<ChangeEnvelope> pending = <ChangeEnvelope>[
|
||||
change,
|
||||
...current.pendingChanges,
|
||||
];
|
||||
await _setState(current.copyWith(pendingChanges: pending));
|
||||
}
|
||||
|
||||
Future<void> enqueueAll(Iterable<ChangeEnvelope> changes) async {
|
||||
final List<ChangeEnvelope> values = changes.toList(growable: false);
|
||||
if (values.isEmpty) {
|
||||
return;
|
||||
}
|
||||
|
||||
final SyncState current = await _currentState();
|
||||
final List<ChangeEnvelope> pending = <ChangeEnvelope>[
|
||||
...values,
|
||||
...current.pendingChanges,
|
||||
];
|
||||
await _setState(current.copyWith(pendingChanges: pending));
|
||||
}
|
||||
|
||||
Future<void> applyPushResult(SyncPushResult result, {DateTime? at}) async {
|
||||
final SyncState current = await _currentState();
|
||||
final DateTime now = at ?? DateTime.now();
|
||||
final Set<String> rejectedMutationIds = <String>{
|
||||
...result.rejected.map(
|
||||
(MutationRejection rejection) => rejection.clientMutationId,
|
||||
),
|
||||
};
|
||||
final Set<String> completedMutationIds = <String>{
|
||||
...result.accepted.map((MutationAck ack) => ack.clientMutationId),
|
||||
...rejectedMutationIds,
|
||||
};
|
||||
final List<ChangeEnvelope> rejectedChanges = current.pendingChanges
|
||||
.where(
|
||||
(ChangeEnvelope change) =>
|
||||
rejectedMutationIds.contains(change.clientMutationId),
|
||||
)
|
||||
.toList(growable: false);
|
||||
final List<ChangeEnvelope> pending = current.pendingChanges
|
||||
.where(
|
||||
(ChangeEnvelope change) =>
|
||||
!completedMutationIds.contains(change.clientMutationId),
|
||||
)
|
||||
.toList(growable: false);
|
||||
final List<ChangeEnvelope> dedupedRejectedChanges = rejectedChanges
|
||||
.where(
|
||||
(ChangeEnvelope change) => !pending.any(
|
||||
(ChangeEnvelope pendingChange) =>
|
||||
pendingChange.clientMutationId == change.clientMutationId,
|
||||
),
|
||||
)
|
||||
.toList(growable: false);
|
||||
|
||||
final String? rejectionMessage = result.rejected.isEmpty
|
||||
? null
|
||||
: 'Push rejected ${result.rejected.length} change(s). Fix locally and retry.';
|
||||
|
||||
await _setState(
|
||||
current.copyWith(
|
||||
cursor: result.cursor,
|
||||
pendingChanges: pending,
|
||||
lastAttemptAt: now,
|
||||
lastSyncAt: now,
|
||||
lastError: rejectionMessage,
|
||||
lastRejected: result.rejected,
|
||||
lastRejectedChanges: dedupedRejectedChanges,
|
||||
),
|
||||
);
|
||||
}
|
||||
|
||||
/// Requeues rejected changes to pending sync list for manual retry.
|
||||
Future<void> requeueRejectedChanges() async {
|
||||
final SyncState current = await _currentState();
|
||||
if (current.lastRejectedChanges.isEmpty) {
|
||||
return;
|
||||
}
|
||||
|
||||
final List<ChangeEnvelope> pending = <ChangeEnvelope>[
|
||||
...current.lastRejectedChanges,
|
||||
...current.pendingChanges,
|
||||
];
|
||||
final List<ChangeEnvelope> dedupedPending = _dedupeByMutationId(pending);
|
||||
await _setState(
|
||||
current.copyWith(
|
||||
pendingChanges: dedupedPending,
|
||||
lastRejected: const <MutationRejection>[],
|
||||
lastRejectedChanges: const <ChangeEnvelope>[],
|
||||
lastError: null,
|
||||
),
|
||||
);
|
||||
}
|
||||
|
||||
Future<void> applyPullResult(SyncPullResult result, {DateTime? at}) async {
|
||||
final SyncState current = await _currentState();
|
||||
final DateTime now = at ?? DateTime.now();
|
||||
await _setState(
|
||||
current.copyWith(
|
||||
cursor: result.cursor,
|
||||
lastAttemptAt: now,
|
||||
lastSyncAt: now,
|
||||
lastError: null,
|
||||
lastRejected: const <MutationRejection>[],
|
||||
lastRejectedChanges: const <ChangeEnvelope>[],
|
||||
),
|
||||
);
|
||||
}
|
||||
|
||||
Future<void> markFailure(Object error, {DateTime? at}) async {
|
||||
final SyncState current = await _currentState();
|
||||
final DateTime now = at ?? DateTime.now();
|
||||
await _setState(current.copyWith(lastAttemptAt: now, lastError: '$error'));
|
||||
}
|
||||
|
||||
Future<void> clearQueue() async {
|
||||
final SyncState current = await _currentState();
|
||||
await _setState(
|
||||
current.copyWith(
|
||||
pendingChanges: const <ChangeEnvelope>[],
|
||||
lastRejected: const <MutationRejection>[],
|
||||
lastRejectedChanges: const <ChangeEnvelope>[],
|
||||
),
|
||||
);
|
||||
}
|
||||
|
||||
/// Clears the latest rejection payload shown to the user.
|
||||
Future<void> clearRejections() async {
|
||||
final SyncState current = await _currentState();
|
||||
final String? lastError =
|
||||
current.lastError != null &&
|
||||
current.lastError!.startsWith('Push rejected')
|
||||
? null
|
||||
: current.lastError;
|
||||
await _setState(
|
||||
current.copyWith(
|
||||
lastRejected: const <MutationRejection>[],
|
||||
lastRejectedChanges: const <ChangeEnvelope>[],
|
||||
lastError: lastError,
|
||||
),
|
||||
);
|
||||
}
|
||||
|
||||
List<ChangeEnvelope> _dedupeByMutationId(List<ChangeEnvelope> changes) {
|
||||
final Set<String> seen = <String>{};
|
||||
final List<ChangeEnvelope> deduped = <ChangeEnvelope>[];
|
||||
for (final ChangeEnvelope change in changes) {
|
||||
if (seen.add(change.clientMutationId)) {
|
||||
deduped.add(change);
|
||||
}
|
||||
}
|
||||
return deduped;
|
||||
}
|
||||
|
||||
Future<SyncState> _currentState() async {
|
||||
final SyncState? value = state.asData?.value;
|
||||
if (value != null) {
|
||||
return value;
|
||||
}
|
||||
return future;
|
||||
}
|
||||
|
||||
Future<void> _setState(SyncState next) async {
|
||||
state = AsyncData<SyncState>(next);
|
||||
await ref
|
||||
.read(syncStateStoreProvider)
|
||||
.write(rawState: jsonEncode(next.toJson()));
|
||||
}
|
||||
}
|
||||
|
||||
final AsyncNotifierProvider<SyncQueueRepository, SyncState>
|
||||
syncQueueRepositoryProvider =
|
||||
AsyncNotifierProvider<SyncQueueRepository, SyncState>(
|
||||
SyncQueueRepository.new,
|
||||
);
|
||||
// Legacy compatibility export for the sync queue data store.
|
||||
export 'package:relationship_saver/features/sync/data/sync_queue_repository.dart';
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user