142 lines
4 KiB
Dart
142 lines
4 KiB
Dart
// 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<String> _queue = [];
|
|
final Set<String> _set = {};
|
|
Timer? _timer;
|
|
bool _flushing = false;
|
|
bool _disposed = false;
|
|
|
|
/// opzionale init se serve (es. caricare stato persistente)
|
|
Future<void> 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<void> _flush() async {
|
|
if (_disposed) return;
|
|
if (_flushing) return;
|
|
if (_queue.isEmpty) return;
|
|
|
|
_flushing = true;
|
|
try {
|
|
// prendi chunk
|
|
final chunk = <String>[];
|
|
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<void> dispose() async {
|
|
_disposed = true;
|
|
_timer?.cancel();
|
|
_timer = null;
|
|
_queue.clear();
|
|
_set.clear();
|
|
}
|
|
}
|