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 { @override Future 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 json = jsonDecode(record.rawState) as Map; return SyncState.fromJson(json); } on FormatException { await store.clear(); return SyncState.empty; } } Future enqueue(ChangeEnvelope change) async { final SyncState current = await _currentState(); final List pending = [ change, ...current.pendingChanges, ]; await _setState(current.copyWith(pendingChanges: pending)); } Future enqueueAll(Iterable changes) async { final List values = changes.toList(growable: false); if (values.isEmpty) { return; } final SyncState current = await _currentState(); final List pending = [ ...values, ...current.pendingChanges, ]; await _setState(current.copyWith(pendingChanges: pending)); } Future applyPushResult(SyncPushResult result, {DateTime? at}) async { final SyncState current = await _currentState(); final DateTime now = at ?? DateTime.now(); final Set rejectedMutationIds = { ...result.rejected.map( (MutationRejection rejection) => rejection.clientMutationId, ), }; final Set completedMutationIds = { ...result.accepted.map((MutationAck ack) => ack.clientMutationId), ...rejectedMutationIds, }; final List rejectedChanges = current.pendingChanges .where( (ChangeEnvelope change) => rejectedMutationIds.contains(change.clientMutationId), ) .toList(growable: false); final List pending = current.pendingChanges .where( (ChangeEnvelope change) => !completedMutationIds.contains(change.clientMutationId), ) .toList(growable: false); final List 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 requeueRejectedChanges() async { final SyncState current = await _currentState(); if (current.lastRejectedChanges.isEmpty) { return; } final List pending = [ ...current.lastRejectedChanges, ...current.pendingChanges, ]; final List dedupedPending = _dedupeByMutationId(pending); await _setState( current.copyWith( pendingChanges: dedupedPending, lastRejected: const [], lastRejectedChanges: const [], lastError: null, ), ); } Future 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 [], lastRejectedChanges: const [], ), ); } Future 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 clearQueue() async { final SyncState current = await _currentState(); await _setState( current.copyWith( pendingChanges: const [], lastRejected: const [], lastRejectedChanges: const [], ), ); } /// Clears the latest rejection payload shown to the user. Future 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 [], lastRejectedChanges: const [], lastError: lastError, ), ); } List _dedupeByMutationId(List changes) { final Set seen = {}; final List deduped = []; for (final ChangeEnvelope change in changes) { if (seen.add(change.clientMutationId)) { deduped.add(change); } } return deduped; } Future _currentState() async { final SyncState? value = state.asData?.value; if (value != null) { return value; } return future; } Future _setState(SyncState next) async { state = AsyncData(next); await ref .read(syncStateStoreProvider) .write(rawState: jsonEncode(next.toJson())); } } final AsyncNotifierProvider syncQueueRepositoryProvider = AsyncNotifierProvider( SyncQueueRepository.new, );