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'; import 'package:terepi_seged/services/app_logger.dart'; import 'package:terepi_seged/services/contact_service.dart'; 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'; 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(); 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 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 syncNow() async { if (_client.auth.currentUser == null) return; if (isSyncing.value) { _syncRequestedWhileRunning = true; return; } isSyncing.value = true; lastError.value = ''; try { lastError.value = ''; // Eszköz-regiszter frissítése (last_seen_at). if (Get.isRegistered()) { await _isolate( 'eszköz regisztációja', DeviceIdentityService.to.registerDevice); } await _isolate('tagság felderítése', _discoverMemberProjects); await _isolate('feltöltés', _push); await _isolate('letöltés', _pull); // Megosztott rétegek (5. lépés) — ha a service be van kötve. if (Get.isRegistered()) { await _isolate('rétegek', LayerSyncService.to.pullAll); } if (Get.isRegistered()) { await _isolate('kitűzés', StakeoutSyncService.to.sync); } if (Get.isRegistered()) { await _isolate('kapcsolatok', () => ContactService.to.flush()); } lastSyncedAt.value = DateTime.now(); } catch (e, s) { lastError.value = e.toString(); AppLogger.e('TsSyncService - SyncNow', lastError.value, error: e, stack: s); } 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 _discoverMemberProjects() async { final rows = await _client .from('terepi_seged_shared_projects') .select() .eq('is_member', true); final remoteUuids = {}; for (final row in rows) { final map = Map.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()) { await ProjectService.to.reloadProjects(); } } // ═════════════════════════════════════════════════════════════════ // PUSH — a 'pending' sorok feltöltése (upsert, idempotens) // ═════════════════════════════════════════════════════════════════ Future _push() async { 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); } Future _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 _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'], '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 pushTracks() async { 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'], '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 pushTrackPoints() async { 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 _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 _pull() async { // Minden szinkronizált (nem lokális) projekt. final projects = await AppDatabase.instance.listProjects(); for (final p in projects.where((p) => !p.isLocalOnly)) { await _isolate('letöltés (${p.name})', () => _pullProject(p.id!, p.uuid)); } } Future _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'], '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'], '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 = {}; final toInsert = >[]; 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(); } /// 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 _isolate(String label, Future 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'; } } }