Közös projekt szinkronizációja
This commit is contained in:
@@ -0,0 +1,435 @@
|
||||
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/device_identity_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<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 {
|
||||
// Eszköz-regiszter frissítése (last_seen_at).
|
||||
if (Get.isRegistered<DeviceIdentityService>()) {
|
||||
await DeviceIdentityService.to.registerDevice();
|
||||
}
|
||||
|
||||
await _discoverMemberProjects();
|
||||
await _push();
|
||||
await _pull();
|
||||
|
||||
// Megosztott rétegek (5. lépés) — ha a service be van kötve.
|
||||
if (Get.isRegistered<LayerSyncService>()) {
|
||||
await LayerSyncService.to.pullAll();
|
||||
}
|
||||
|
||||
lastSyncedAt.value = DateTime.now();
|
||||
} catch (e) {
|
||||
lastError.value = e.toString();
|
||||
} 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);
|
||||
|
||||
for (final row in rows) {
|
||||
await _db.upsertProjectFromRemote(Map<String, dynamic>.from(row));
|
||||
}
|
||||
}
|
||||
|
||||
// ═════════════════════════════════════════════════════════════════
|
||||
// PUSH — a 'pending' sorok feltöltése (upsert, idempotens)
|
||||
// ═════════════════════════════════════════════════════════════════
|
||||
|
||||
Future<void> _push() async {
|
||||
await _pushProjects();
|
||||
await _pushMeasuredPoints();
|
||||
await _pushTracks();
|
||||
await _pushTrackPoints();
|
||||
await _pushNoteItems();
|
||||
}
|
||||
|
||||
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'],
|
||||
'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 {
|
||||
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'],
|
||||
'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 {
|
||||
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)) {
|
||||
await _pullProject(p.id!, p.uuid);
|
||||
}
|
||||
}
|
||||
|
||||
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'],
|
||||
'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'],
|
||||
'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();
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user