Files
MobilApp/lib/services/ts_sync_service.dart
T

486 lines
18 KiB
Dart
Raw Normal View History

2026-07-07 02:21:08 +02:00
import 'dart:async';
import 'dart:convert';
import 'package:connectivity_plus/connectivity_plus.dart';
import 'package:get/get.dart';
import 'package:supabase_flutter/supabase_flutter.dart';
2026-07-12 02:14:49 +02:00
import 'package:terepi_seged/services/app_logger.dart';
import 'package:terepi_seged/services/contact_service.dart';
2026-07-07 02:21:08 +02:00
import 'package:terepi_seged/services/device_identity_service.dart';
import 'package:terepi_seged/services/project_service.dart';
import 'package:terepi_seged/services/stakeout_sync_service.dart';
2026-07-07 02:21:08 +02:00
import 'app_database.dart';
import 'layer_sync_service.dart';
/// Outbox-alapú, kétirányú szinkron a ts_* Supabase-táblákkal.
///
/// Elvek:
/// * OFFLINE-FIRST: a UI mindig a lokális DB-ből dolgozik, a szinkron
/// háttérfolyamat. Net nélkül minden működik, a változások sorban állnak.
/// * A "csak lokális" (is_local_only = 1) projektek és adataik SOHA nem
/// kerülnek feltöltésre — a kizárást az AppDatabase pending-lekérdezései
/// végzik SQL-szinten, így nem múlhat hívási fegyelmen.
/// * PUSH → PULL sorrend: előbb a saját 'pending' sorok mennek fel
/// (uuid-ra upsert, ezért idempotens — ismételt feltöltés nem duplikál),
/// utána jön le a többiek munkája (updated_at > kurzor).
/// * LWW (last-write-wins) az updated_at alapján, amit a SZERVER állít.
///
/// Triggerek: hálózat visszatérése (connectivity_plus), periodikus időzítő,
/// és kézi syncNow() (pl. pull-to-refresh, csatlakozás után).
class TsSyncService extends GetxService {
static TsSyncService get to => Get.find();
final isSyncing = false.obs;
final lastSyncedAt = Rxn<DateTime>();
final lastError = ''.obs;
final pendingCount = 0.obs;
SupabaseClient get _client => Supabase.instance.client;
AppDatabase get _db => AppDatabase.instance;
StreamSubscription? _connSub;
Timer? _timer;
bool _syncRequestedWhileRunning = false;
@override
void onInit() {
super.onInit();
// Hálózat visszatérésekor azonnali szinkron.
_connSub = Connectivity().onConnectivityChanged.listen((results) {
final online = results.any((r) => r != ConnectivityResult.none);
if (online) syncNow();
});
// Biztonsági háló: 3 percenként akkor is, ha nem volt esemény.
_timer = Timer.periodic(const Duration(minutes: 3), (_) => syncNow());
refreshPendingCount();
}
@override
void onClose() {
_connSub?.cancel();
_timer?.cancel();
super.onClose();
}
Future<void> refreshPendingCount() async {
pendingCount.value = await _db.pendingSyncCount();
}
/// Teljes szinkron-ciklus. Újrahívás futás közben nem indít párhuzamos
/// ciklust, csak megjegyzi, hogy a végén még egyszer le kell futnia.
Future<void> syncNow() async {
if (_client.auth.currentUser == null) return;
if (isSyncing.value) {
_syncRequestedWhileRunning = true;
return;
}
isSyncing.value = true;
lastError.value = '';
try {
2026-07-12 02:14:49 +02:00
lastError.value = '';
2026-07-07 02:21:08 +02:00
// Eszköz-regiszter frissítése (last_seen_at).
if (Get.isRegistered<DeviceIdentityService>()) {
2026-07-12 02:14:49 +02:00
await _isolate(
'eszköz regisztációja', DeviceIdentityService.to.registerDevice);
2026-07-07 02:21:08 +02:00
}
2026-07-12 02:14:49 +02:00
await _isolate('tagság felderítése', _discoverMemberProjects);
await _isolate('feltöltés', _push);
await _isolate('letöltés', _pull);
2026-07-07 02:21:08 +02:00
// Megosztott rétegek (5. lépés) — ha a service be van kötve.
if (Get.isRegistered<LayerSyncService>()) {
2026-07-12 02:14:49 +02:00
await _isolate('rétegek', LayerSyncService.to.pullAll);
2026-07-07 02:21:08 +02:00
}
if (Get.isRegistered<StakeoutSyncService>()) {
2026-07-12 02:14:49 +02:00
await _isolate('kitűzés', StakeoutSyncService.to.sync);
}
if (Get.isRegistered<ContactService>()) {
2026-07-12 02:14:49 +02:00
await _isolate('kapcsolatok', () => ContactService.to.flush());
}
2026-07-07 02:21:08 +02:00
lastSyncedAt.value = DateTime.now();
2026-07-12 02:14:49 +02:00
} catch (e, s) {
2026-07-07 02:21:08 +02:00
lastError.value = e.toString();
2026-07-12 02:14:49 +02:00
AppLogger.e('TsSyncService - SyncNow', lastError.value,
error: e, stack: s);
2026-07-07 02:21:08 +02:00
} finally {
await refreshPendingCount();
isSyncing.value = false;
if (_syncRequestedWhileRunning) {
_syncRequestedWhileRunning = false;
unawaited(syncNow());
}
}
}
// ═════════════════════════════════════════════════════════════════
// Tagság-felderítés
// ═════════════════════════════════════════════════════════════════
/// Ha a felhasználó egy MÁSIK eszközén csatlakozott egy közös
/// projekthez, ez hozza létre a lokális projekt-sort itt is — így a
/// több-eszközös használat magától konzisztens marad.
Future<void> _discoverMemberProjects() async {
final rows = await _client
.from('terepi_seged_shared_projects')
.select()
.eq('is_member', true);
final remoteUuids = <String>{};
2026-07-07 02:21:08 +02:00
for (final row in rows) {
final map = Map<String, dynamic>.from(row);
remoteUuids.add(map['id'] as String);
await _db.upsertProjectFromRemote(map);
}
// Ami tartósan hiányzik erről a listáról (törölve vagy kikerültünk a
// tagságból), azt a helyi gyerek-adatokkal együtt eltávolítjuk.
await _db.reconcileMissingProjects(remoteUuids);
// A ProjectService saját, memóriában tartott listája nem tud
// magától a helyi törlésről — enélkül a projekt-választóban addig
// ottmaradna, amíg valaki újra nem indítja az appot.
if (Get.isRegistered<ProjectService>()) {
await ProjectService.to.reloadProjects();
2026-07-07 02:21:08 +02:00
}
}
// ═════════════════════════════════════════════════════════════════
// PUSH — a 'pending' sorok feltöltése (upsert, idempotens)
// ═════════════════════════════════════════════════════════════════
Future<void> _push() async {
2026-07-12 02:14:49 +02:00
await _isolate('projektek push', _pushProjects);
await _isolate('mérési pontok', _pushMeasuredPoints);
await _isolate('track-ek push', pushTracks);
await _isolate('track-pontok push', pushTrackPoints);
await _isolate('jegyzetek push', _pushNoteItems);
2026-07-07 02:21:08 +02:00
}
Future<void> _pushProjects() async {
final rows = await _db.pendingProjects();
if (rows.isEmpty) return;
final payload = rows
.map((r) => {
'id': r['uuid'],
'name': r['name'],
'client': r['client'],
'description': r['description'],
'crs': r['crs'],
'color': r['color'],
'status': r['status'],
'deleted_at': _toUtc(r['deleted_at'] as String?),
})
.toList();
await _client.from('terepi_seged_projects').upsert(payload);
await _db.markSyncedByUuid(
'projects', rows.map((r) => r['uuid'] as String).toList());
}
Future<void> _pushMeasuredPoints() async {
final rows = await _db.pendingMeasuredPoints();
if (rows.isEmpty) return;
final payload = rows
.map((r) => {
'id': r['uuid'],
'project_id': r['project_uuid'],
'name': r['name'],
'eov_y': r['eov_y'],
'eov_x': r['eov_x'],
'eov_z': r['eov_z'],
'latitude': r['latitude'],
'longitude': r['longitude'],
'altitude': r['altitude'],
'accuracy': r['accuracy'],
'fix_quality': r['fix_quality'],
'measured_at': _toUtc(r['timestamp'] as String?),
'note': r['note'],
'device_id': r['device_id'],
'app_instance_id': r['app_instance_id'],
2026-07-07 02:21:08 +02:00
'deleted_at': _toUtc(r['deleted_at'] as String?),
})
.toList();
await _client.from('terepi_seged_measured_points').upsert(payload);
await _db.markSyncedByUuid(
'measured_points', rows.map((r) => r['uuid'] as String).toList());
}
Future<void> pushTracks() async {
2026-07-07 02:21:08 +02:00
final rows = await _db.pendingTracks();
if (rows.isEmpty) return;
final payload = rows
.map((r) => {
'id': r['uuid'],
'project_id': r['project_uuid'],
'name': r['name'],
'start_time': _toUtc(r['start_time'] as String?),
'end_time': _toUtc(r['end_time'] as String?),
'status': r['status'],
'source': r['source'],
'distance_m': r['distance_m'],
'point_count': r['point_count'],
'device_id': r['device_id'],
'app_instance_id': r['app_instance_id'],
2026-07-07 02:21:08 +02:00
'deleted_at': _toUtc(r['deleted_at'] as String?),
})
.toList();
await _client.from('terepi_seged_tracks').upsert(payload);
await _db.markSyncedByUuid(
'tracks', rows.map((r) => r['uuid'] as String).toList());
}
/// Track-pontok: nagy mennyiség lehet, ezért 500-as batch-ekben, ciklusban,
/// amíg el nem fogy. A szerveren nincs UPDATE policy (append-only), ezért
/// ignoreDuplicates — az ismételt feltöltés csendben kimarad.
Future<void> pushTrackPoints() async {
2026-07-07 02:21:08 +02:00
while (true) {
final rows = await _db.pendingTrackPoints(limit: 500);
if (rows.isEmpty) break;
final payload = rows
.map((r) => {
'id': r['uuid'],
'track_id': r['track_uuid'],
'project_id': r['project_uuid'],
'latitude': r['latitude'],
'longitude': r['longitude'],
'altitude': r['altitude'],
'accuracy': r['accuracy'],
'speed': r['speed'],
'heading': r['heading'],
'recorded_at': _toUtc(r['timestamp'] as String?),
})
.toList();
await _client
.from('terepi_seged_track_points')
.upsert(payload, ignoreDuplicates: true);
await _db.markSyncedByUuid(
'track_points', rows.map((r) => r['uuid'] as String).toList());
}
}
Future<void> _pushNoteItems() async {
final rows = await _db.pendingNoteItems();
if (rows.isEmpty) return;
final payload = rows
.map((r) => {
'id': r['uuid'],
'project_id': r['project_uuid'],
'type': r['type'],
'points_json': jsonDecode(r['points_json'] as String),
'color': r['color'],
'opacity': r['opacity'],
'stroke_width': r['stroke_width'],
'stroke_color': r['stroke_color'],
'label': r['label'],
'deleted_at': _toUtc(r['deleted_at'] as String?),
})
.toList();
await _client.from('terepi_seged_note_items').upsert(payload);
await _db.markSyncedByUuid(
'note_items', rows.map((r) => r['uuid'] as String).toList());
}
// ═════════════════════════════════════════════════════════════════
// PULL — inkrementális letöltés projektenként (updated_at > kurzor)
// ═════════════════════════════════════════════════════════════════
Future<void> _pull() async {
// Minden szinkronizált (nem lokális) projekt.
final projects = await AppDatabase.instance.listProjects();
for (final p in projects.where((p) => !p.isLocalOnly)) {
2026-07-12 02:14:49 +02:00
await _isolate('letöltés (${p.name})', () => _pullProject(p.id!, p.uuid));
2026-07-07 02:21:08 +02:00
}
}
Future<void> _pullProject(int localId, String projectUuid) async {
final cursor =
await _db.getProjectCursor(localId) ?? '1970-01-01T00:00:00Z';
String maxSeen = cursor;
String track(String? ts) {
if (ts != null && DateTime.parse(ts).isAfter(DateTime.parse(maxSeen))) {
maxSeen = ts;
}
return ts ?? '';
}
// ── Bemért pontok ───────────────────────────────────────────────
final mps = await _client
.from('terepi_seged_measured_points')
.select()
.eq('project_id', projectUuid)
.gt('updated_at', cursor)
.order('updated_at', ascending: true)
.limit(1000);
for (final r in mps) {
track(r['updated_at'] as String?);
await _db.applyRemoteRow('measured_points', {
'uuid': r['id'],
'project_id': localId,
'name': r['name'],
'eov_y': r['eov_y'],
'eov_x': r['eov_x'],
'eov_z': r['eov_z'],
'latitude': r['latitude'],
'longitude': r['longitude'],
'altitude': r['altitude'],
'accuracy': r['accuracy'],
'fix_quality': r['fix_quality'],
'timestamp': r['measured_at'],
'note': r['note'] ?? '',
'created_by': r['created_by'],
'device_id': r['device_id'],
'app_instance_id': r['app_instance_id'],
2026-07-07 02:21:08 +02:00
'updated_at': r['updated_at'],
'deleted_at': r['deleted_at'],
'sync_status': 'synced',
});
}
// ── Trackek ─────────────────────────────────────────────────────
final tracks = await _client
.from('terepi_seged_tracks')
.select()
.eq('project_id', projectUuid)
.gt('updated_at', cursor)
.order('updated_at', ascending: true)
.limit(1000);
for (final r in tracks) {
track(r['updated_at'] as String?);
await _db.applyRemoteRow('tracks', {
'uuid': r['id'],
'project_id': localId,
'name': r['name'],
'start_time': r['start_time'],
'end_time': r['end_time'],
'status': r['status'],
'source': r['source'] ?? 'Telefon GPS',
'distance_m': r['distance_m'] ?? 0,
'point_count': r['point_count'] ?? 0,
'is_local_only': 0,
'created_by': r['created_by'],
'device_id': r['device_id'],
'app_instance_id': r['app_instance_id'],
2026-07-07 02:21:08 +02:00
'updated_at': r['updated_at'],
'deleted_at': r['deleted_at'],
'sync_status': 'synced',
});
}
// ── Track-pontok (append-only; created_at a szűrő) ─────────────
final tps = await _client
.from('terepi_seged_track_points')
.select()
.eq('project_id', projectUuid)
.gt('created_at', cursor)
.order('created_at', ascending: true)
.limit(2000);
// Távoli track-uuid → lokális track-id feloldás.
final trackIds = <String, int?>{};
final toInsert = <Map<String, dynamic>>[];
for (final r in tps) {
track(r['created_at'] as String?);
final tUuid = r['track_id'] as String;
trackIds[tUuid] ??= await _db.getTrackIdByUuid(tUuid);
final localTrackId = trackIds[tUuid];
if (localTrackId == null) continue; // a track a következő körben jön
toInsert.add({
'uuid': r['id'],
'track_id': localTrackId,
'latitude': r['latitude'],
'longitude': r['longitude'],
'altitude': r['altitude'],
'accuracy': r['accuracy'],
'speed': r['speed'],
'heading': r['heading'],
'timestamp': r['recorded_at'],
'sync_status': 'synced',
});
}
await _db.insertRemoteTrackPoints(toInsert);
// ── Terepbejárás-elemek ────────────────────────────────────────
final notes = await _client
.from('terepi_seged_note_items')
.select()
.eq('project_id', projectUuid)
.gt('updated_at', cursor)
.order('updated_at', ascending: true)
.limit(1000);
for (final r in notes) {
track(r['updated_at'] as String?);
await _db.applyRemoteRow('note_items', {
'uuid': r['id'],
'project_id': localId,
'type': r['type'],
'points_json': jsonEncode(r['points_json']),
'color': r['color'],
'opacity': r['opacity'],
'stroke_width': r['stroke_width'],
'stroke_color': r['stroke_color'],
'label': r['label'] ?? '',
'created_at': r['created_at'],
'created_by': r['created_by'],
'updated_at': r['updated_at'],
'deleted_at': r['deleted_at'],
'sync_status': 'synced',
});
}
if (maxSeen != cursor) {
await _db.setProjectCursor(localId, maxSeen);
}
}
// ═════════════════════════════════════════════════════════════════
/// A lokális időbélyegek zóna nélküliek (DateTime.now().toIso8601String()),
/// a Postgres timestamptz viszont zóna nélküli értéket UTC-nek értelmezne
/// → küldés előtt explicit UTC-re konvertálunk, különben 1-2 órát csúszna
/// minden időpont.
String? _toUtc(String? localIso) {
if (localIso == null || localIso.isEmpty) return null;
return DateTime.parse(localIso).toUtc().toIso8601String();
}
2026-07-12 02:14:49 +02:00
/// Egy lépés elszigetelt futtatása: hiba vagy időtúllépés esetén NEM
/// dobja tovább — csak feljegyzi és a szinkron a KÖVETKEZŐ lépéssel
/// folytatódik. Enélkül egyetlen hibás sor (típushiba, FK-ütközés stb.)
/// vagy egy beragadt hálózati hívás CSENDBEN leállítaná az összes
/// további lépést, minden ciklusban, örökre.
Future<void> _isolate(String label, Future<void> Function() fn,
{Duration timeout = const Duration(seconds: 25)}) async {
try {
await fn().timeout(timeout);
} catch (e) {
lastError.value = lastError.value.isEmpty
? '$label: $e'
: '${lastError.value} · $label: $e';
}
}
2026-07-07 02:21:08 +02:00
}