// lib/remote/remote_added_queue.dart // // Queue per gestire eventi "added" dal WS: // - dedupe (dinamico) per evitare dipendenze da file mancanti // - batching byIds (batchSize) // - debounce flush (debounce) // - fallback a progressiveSync se burst troppo grande (tooManyThreshold) // - upsert massivo delegato a RemoteSyncEngine.upsertRemoteEntries // - notifica UI tramite onUiRefresh (debounced lato UI) import 'dart:async'; import 'package:flutter/foundation.dart'; import 'remote_http_api.dart'; import 'remote_sync_engine.dart'; // per SyncManager.instance se serve class RemoteAddedQueue { RemoteAddedQueue({ required this.api, required this.engine, required this.deduper, required this.onUiRefresh, this.batchSize = 200, this.tooManyThreshold = 1200, this.debounce = const Duration(milliseconds: 250), }); final RemoteHttpApi api; final RemoteSyncEngine engine; final dynamic deduper; // dynamic per compatibilità con diverse implementazioni final VoidCallback onUiRefresh; final int batchSize; final int tooManyThreshold; final Duration debounce; final List _queue = []; final Set _set = {}; Timer? _timer; bool _flushing = false; bool _disposed = false; /// opzionale init se serve (es. caricare stato persistente) Future init() async { // placeholder: se vuoi ripristinare queue persistente, fallo qui return; } void enqueue(String? id) { if (_disposed) return; if (id == null) return; final sid = id.toString(); if (_set.contains(sid)) return; // se deduper ha has, evita duplicati già processati try { if (deduper != null && (deduper.has as dynamic) != null) { final already = (deduper.has as dynamic)(sid); if (already == true) return; } } catch (_) {} _set.add(sid); _queue.add(sid); // se burst troppo grande -> fallback a progressiveSync if (_queue.length >= tooManyThreshold) { debugPrint('[remote][addedQueue] TOO_MANY_THRESHOLD reached (${_queue.length}) -> fallback progressiveSync'); _timer?.cancel(); _timer = null; _queue.clear(); _set.clear(); unawaited(engine.progressiveSync()); return; } _timer ??= Timer(debounce, () { _timer = null; unawaited(_flush()); }); } Future _flush() async { if (_disposed) return; if (_flushing) return; if (_queue.isEmpty) return; _flushing = true; try { // prendi chunk final chunk = []; final take = _queue.length < batchSize ? _queue.length : batchSize; for (var i = 0; i < take; i++) { chunk.add(_queue.removeAt(0)); } chunk.forEach(_set.remove); if (chunk.isEmpty) return; debugPrint('[remote][addedQueue] flushing chunk size=${chunk.length}'); // fetch by ids final items = await api.getPhotosByIds(chunk); if (items.isNotEmpty) { // delega all'engine: centralizza upsert + creazione folders + append // engine.upsertRemoteEntries accetta sia Map che RemotePhotoItem (dynamic) await engine.upsertRemoteEntries(items); } // Non tentiamo di leggere event_id come Map: il payload può essere RemotePhotoItem. // Se vuoi markare event_id, fallo dove il server fornisce esplicitamente event_id. // refresh UI una sola volta per batch (debounced lato UI) try { onUiRefresh(); } catch (e, st) { debugPrint('[remote][addedQueue] onUiRefresh error: $e\n$st'); } // se ci sono ancora elementi, schedule immediato per continuare if (_queue.isNotEmpty) { Future.microtask(_flush); } } catch (e, st) { debugPrint('[remote][addedQueue] flush error: $e\n$st'); // svuota e fallback a progressiveSync _queue.clear(); _set.clear(); await engine.progressiveSync(); } finally { _flushing = false; } } Future dispose() async { _disposed = true; _timer?.cancel(); _timer = null; _queue.clear(); _set.clear(); } }