import 'package:flutter/foundation.dart'; import 'package:aves/model/source/collection_source.dart'; import 'remote_http_api.dart'; import 'remote_models.dart'; import 'remote_repository.dart'; import 'remote_state_store.dart'; import 'package:aves/remote/collection_source_remote_ext.dart'; import 'package:aves/remote/collection_source_remote_ws_ext.dart'; import 'package:aves/remote/remote_controller.dart'; class RemoteSyncEngine { RemoteSyncEngine({ required this.api, required this.repo, required this.source, required this.state, }); final RemoteHttpApi api; final RemoteRepository repo; final CollectionSource source; final RemoteStateStore state; static const int retentionDays = 30; static const Duration overlap = Duration(seconds: 3); bool _syncInFlight = false; bool _pendingSync = false; bool _isTooOld(String iso) { final t = DateTime.tryParse(iso); if (t == null) return true; return DateTime.now().toUtc().difference(t.toUtc()).inDays > retentionDays; } String _applyOverlap(String iso) { final t = DateTime.tryParse(iso); if (t == null) return iso; return t.toUtc().subtract(overlap).toIso8601String(); } // ------------------------------------------------------------ // Public API // ------------------------------------------------------------ Future progressiveSync() async { return _runExclusive(() async { final last = await state.getLastSyncIso(); if (last == null || last.isEmpty || _isTooOld(last)) { await fullSyncImpl(); return; } await _deltaSyncImpl(_applyOverlap(last)); }); } Future progressiveSyncFrom(String sinceIso) async { return _runExclusive(() async { if (sinceIso.isEmpty || _isTooOld(sinceIso)) { await fullSyncImpl(); return; } await _deltaSyncImpl(_applyOverlap(sinceIso)); }); } // ------------------------------------------------------------ // FULL SYNC (senza _runExclusive) // ------------------------------------------------------------ Future fullSyncImpl() async { debugPrint('[remote] fullSync start'); final list = await api.getAllPhotos(); final items = list.map(RemotePhotoItem.fromJson).toList(); // decide se siamo in bootstrap (solo in quel caso deleteAllRemotes) final isBootstrap = !(await RemoteController.instance.bootstrapDone()); final serverIds = items.map((e) => e.id).where((v) => v.isNotEmpty).toSet(); if (isBootstrap) { debugPrint('[remote] fullSync bootstrap -> deleteAllRemotes'); await repo.deleteAllRemotes(); } else { debugPrint('[remote] fullSync incremental -> no deleteAllRemotes'); } // upsert chunked con yield per non bloccare il main thread const int chunk = 40; // prova 40, abbassa a 20 se serve final t0 = DateTime.now(); for (var i = 0; i < items.length; i += chunk) { final end = (i + chunk < items.length) ? i + chunk : items.length; final sub = items.sublist(i, end); final chunkT0 = DateTime.now(); await repo.upsertAll(sub, chunkSize: chunk); debugPrint('[remote] fullSync upsert chunk $i..${end - 1} size=${sub.length} took=${DateTime.now().difference(chunkT0)}'); // yield per permettere redraw UI await Future.delayed(const Duration(milliseconds: 20)); } debugPrint('[remote] fullSync upsertAll total=${items.length} totalTime=${DateTime.now().difference(t0)}'); // se non era bootstrap, rimuovi i remoti mancanti (prune) if (!isBootstrap) { final pruned = await repo.pruneMissingRemotes(serverIds); debugPrint('[remote] fullSync pruneMissingRemotes deleted=$pruned'); } // append remoti in collection (collection batching dovrebbe gestire il carico) await source.appendRemoteEntriesFromDb(); final nextIso = _maxIsoFromItems(items) ?? DateTime.now().toUtc().toIso8601String(); await state.setLastSyncIso(nextIso); final nextMs = DateTime.parse(nextIso).millisecondsSinceEpoch; final sessionId = await state.getSessionId(); final deviceId = await state.getDeviceId(); final ws = source.remoteWsClient; ws?.notifyRecoveryCompleted(); debugPrint('[remote] sending sync_done session=$sessionId device=$deviceId ms=$nextMs'); await api.sendSyncDone( sessionId: sessionId, deviceId: deviceId, lastSyncMs: nextMs, ); debugPrint('[remote] fullSync done items=${items.length} lastSync=$nextIso'); } // ------------------------------------------------------------ // DELTA SYNC (senza _runExclusive) // ------------------------------------------------------------ Future _deltaSyncImpl(String sinceIso, {bool appendAfter = true}) async { debugPrint('[remote] deltaSyncFrom since=$sinceIso'); final List> changed = await api.getChanges(sinceIso); final List> deleted = await api.getDeletedHard(sinceIso); var didChange = false; List changedItems = const []; if (changed.isNotEmpty) { changedItems = changed.map(RemotePhotoItem.fromJson).toList(); const int deltaChunk = 40; const int maxChunkAttempts = 2; // numero di tentativi per chunk for (var i = 0; i < changedItems.length; i += deltaChunk) { final end = (i + deltaChunk < changedItems.length) ? i + deltaChunk : changedItems.length; final sub = changedItems.sublist(i, end); var attempt = 0; var success = false; final chunkT0 = DateTime.now(); while (!success && attempt < maxChunkAttempts) { attempt++; try { await repo.upsertAll(sub, chunkSize: deltaChunk); success = true; debugPrint('[remote] delta upsert chunk $i..${end - 1} size=${sub.length} took=${DateTime.now().difference(chunkT0)} attempt=$attempt'); } catch (e, st) { debugPrint('[remote] delta upsert chunk $i..${end - 1} FAILED attempt=$attempt: $e\n$st'); // piccolo backoff prima del retry if (attempt < maxChunkAttempts) await Future.delayed(const Duration(milliseconds: 50)); } } if (!success) { // segnala ma continua con i chunk successivi debugPrint('[remote] delta upsert chunk $i..${end - 1} permanently failed after $maxChunkAttempts attempts'); } // yield per non bloccare il renderer await Future.delayed(const Duration(milliseconds: 10)); } debugPrint('[remote] delta upsert total=${changedItems.length}'); didChange = true; } if (deleted.isNotEmpty) { final ids = deleted .map((e) => e['id']?.toString()) .whereType() .where((s) => s.isNotEmpty) .toSet(); if (ids.isNotEmpty) { await repo.deleteRemotesByRemoteIds(ids); debugPrint('[remote] delta hardDeleted=${ids.length}'); didChange = true; } } if (appendAfter && didChange) { await source.appendRemoteEntriesFromDb(); } final next = _computeNextSyncIso( changedItems: changedItems, hardDeletedRaw: deleted, ) ?? DateTime.now().toUtc().toIso8601String(); await state.setLastSyncIso(next); debugPrint('[remote] delta lastSync -> $next'); } /// Upsert massivo di items ricevuti dal WS (items: List) /// - chunked per non bloccare il main thread /// - accetta sia Map che RemotePhotoItem /// - assicura append in-memory via source.appendRemoteEntriesFromDb() Future upsertRemoteEntries(List items) async { if (items.isEmpty) return; // Normalizza in RemotePhotoItem final List photos = []; for (final it in items) { try { if (it is RemotePhotoItem) { photos.add(it); } else if (it is Map) { photos.add(RemotePhotoItem.fromJson(it)); } else { // ignora tipi non riconosciuti } } catch (e, st) { debugPrint('[remote] upsertRemoteEntries: item parse error: $e\n$st'); } } if (photos.isEmpty) return; // chunk per non bloccare il renderer const int chunk = 40; for (var i = 0; i < photos.length; i += chunk) { final end = (i + chunk < photos.length) ? i + chunk : photos.length; final sub = photos.sublist(i, end); final chunkT0 = DateTime.now(); try { // usa il repository per l'upsert (repo.upsertAll già chunked) await repo.upsertAll(sub, chunkSize: chunk); debugPrint('[remote] upsertRemoteEntries chunk $i..${end - 1} size=${sub.length} took=${DateTime.now().difference(chunkT0)}'); } catch (e, st) { debugPrint('[remote] upsertRemoteEntries chunk FAILED $i..${end - 1}: $e\n$st'); // continua con i chunk successivi } // piccolo yield per permettere redraw UI await Future.delayed(const Duration(milliseconds: 10)); } // Nota: non facciamo riferimento a repo.ensureFoldersFromEntries o a metodi non esistenti. // Ci affidiamo a source.appendRemoteEntriesFromDb() per aggiornare la cache in-memory e // alla logica esistente che ricostruisce folders se necessario. try { await source.appendRemoteEntriesFromDb(); } catch (e, st) { debugPrint('[remote] appendRemoteEntriesFromDb error: $e\n$st'); } } // ------------------------------------------------------------ // Exclusive wrapper (FIXATO) // ------------------------------------------------------------ Future _runExclusive(Future Function() fn) async { if (_syncInFlight) { _pendingSync = true; debugPrint('[remote] sync already in flight -> pending=true'); return; } _syncInFlight = true; try { await fn(); } finally { _syncInFlight = false; if (_pendingSync) { _pendingSync = false; debugPrint('[remote] run pending sync'); // ❗ NON ricorsivo Future.microtask(() => _runExclusive(fn)); } } } // ------------------------------------------------------------ // Helpers // ------------------------------------------------------------ String? _maxIsoFromItems(List items) { DateTime? maxDt; for (final it in items) { final dt = it.updatedAtUtc ?? it.createdAtUtc; if (dt != null && (maxDt == null || dt.isAfter(maxDt))) { maxDt = dt; } } return maxDt?.toUtc().toIso8601String(); } String? _computeNextSyncIso({ required List changedItems, required List> hardDeletedRaw, }) { DateTime? maxDt; void consider(DateTime? dt) { if (dt == null) return; if (maxDt == null || dt.isAfter(maxDt!)) maxDt = dt; } for (final it in changedItems) { consider(it.updatedAtUtc ?? it.createdAtUtc); } for (final d in hardDeletedRaw) { final s = d['deleted_at']?.toString(); if (s != null && s.isNotEmpty) { consider(DateTime.tryParse(s)?.toUtc()); } } return maxDt?.toUtc().toIso8601String(); } }