using FieldLogger.Models; using FieldLogger.Services.Data; using Microsoft.Extensions.Logging; using SyncCore = FieldLogger.Sync; namespace FieldLogger.Services.Sync; /// /// Bridges the MAUI capture database to the workload-free durable MQTT engine. Network loss never /// deletes a capture: the local point and outbound queue are separate durable records, and an /// application acknowledgement is the only path that marks the local point synced. /// public sealed class MqttSyncService : IMqttSyncService { private readonly SettingsService _settings; private readonly AppDatabase _database; private readonly ILogger _logger; private readonly SyncCore.SqliteOutboundStore _store; private readonly SemaphoreSlim _connectionGate = new(1, 1); private readonly SemaphoreSlim _pumpGate = new(1, 1); private CancellationTokenSource? _lifetime; private Task? _backgroundLoop; private SyncCore.MqttnetTransport? _transport; private SyncCore.MqttSyncEngine? _engine; private bool _started; public bool IsConnected => _transport?.IsConnected == true; public string Status { get; private set; } = "Not started"; public event EventHandler? StatusChanged; public event EventHandler? PointRejected; public MqttSyncService(SettingsService settings, AppDatabase database, ILogger logger) { _settings = settings; _database = database; _logger = logger; _store = new SyncCore.SqliteOutboundStore( Path.Combine(FileSystem.AppDataDirectory, "fieldlogger-sync.db3")); } public async Task StartAsync(CancellationToken cancellationToken = default) { if (_started) return; _started = true; await _store.InitAsync(); _lifetime = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); _backgroundLoop = RunAsync(_lifetime.Token); if (_settings.MqttEnabled) { try { await ConnectAsync(cancellationToken); } catch (Exception ex) { _logger.LogWarning(ex, "Initial MQTT connection failed; durable retry loop remains active"); SetStatus($"Waiting to reconnect: {ex.Message}"); } } else { SetStatus("Sync disabled"); } } public async Task StopAsync() { _lifetime?.Cancel(); if (_backgroundLoop is not null) { try { await _backgroundLoop; } catch (OperationCanceledException) { } } await DisconnectTransportAsync(); SetStatus("Stopped"); } public async Task ConnectAsync(CancellationToken cancellationToken = default) { await ConnectCoreAsync(forceReconnect: true, cancellationToken); await QueueUnsyncedAsync(cancellationToken); await DrainAsync(cancellationToken); } public async Task DisconnectAsync() { await DisconnectTransportAsync(); SetStatus("Disconnected"); } public async Task PublishPointAsync(LoggedPoint point, CancellationToken cancellationToken = default) { await _pumpGate.WaitAsync(cancellationToken); try { await QueuePointAsync(point, cancellationToken); if (IsConnected && _engine is not null) await _engine.DrainOnceAsync(cancellationToken); } finally { _pumpGate.Release(); } } private async Task QueuePointAsync(LoggedPoint point, CancellationToken cancellationToken) { if (string.IsNullOrWhiteSpace(point.SyncPointId)) { point.SyncPointId = Guid.CreateVersion7().ToString(); point.Synced = false; await _database.UpdatePointAsync(point); } if (string.IsNullOrWhiteSpace(_settings.MqttOrgId)) { await RecordFailureAsync(point, "SYNC_NOT_CONFIGURED"); return; } var job = await _database.GetJobAsync(point.JobId); if (job is null) { await RecordFailureAsync(point, "LOCAL_JOB_NOT_FOUND"); return; } try { var record = ToSyncPoint(point); var topic = $"ul/{_settings.MqttOrgId}/app/{_settings.MqttClientId}/log/points"; string? remoteJobId = string.IsNullOrWhiteSpace(job.RemoteId) ? null : job.RemoteId; string? ticket = remoteJobId is null ? TicketFor(job) : null; var outbound = SyncCore.OutboundMessage.FromPoint( record, topic, DateTimeOffset.UtcNow, jobId: remoteJobId, ticket: ticket); await _store.EnqueueAsync(outbound); point.SyncError = null; await _database.UpdatePointAsync(point); } catch (SyncCore.PointNotPublishableException ex) { await RecordFailureAsync(point, ex.ReasonCode); } catch (Exception ex) { _logger.LogError(ex, "Failed to enqueue point {PointId}", point.SyncPointId); await RecordFailureAsync(point, "QUEUE_ERROR"); } cancellationToken.ThrowIfCancellationRequested(); } private async Task QueueUnsyncedAsync(CancellationToken cancellationToken) { foreach (var point in await _database.GetUnsyncedPointsAsync()) { cancellationToken.ThrowIfCancellationRequested(); // Permanent capture/server rejections require operator action, not an automatic loop. if (point.SyncError is not null && point.SyncError is not "SYNC_NOT_CONFIGURED" and not "QUEUE_ERROR") continue; await QueuePointAsync(point, cancellationToken); } } private async Task ConnectCoreAsync(bool forceReconnect, CancellationToken cancellationToken) { if (!_settings.MqttEnabled) throw new InvalidOperationException("MQTT sync is disabled in Settings."); if (string.IsNullOrWhiteSpace(_settings.MqttHost) || string.IsNullOrWhiteSpace(_settings.MqttOrgId)) throw new InvalidOperationException("MQTT host and organization id are required."); string password = await _settings.GetMqttPasswordAsync() ?? throw new InvalidOperationException("MQTT password is required."); await _connectionGate.WaitAsync(cancellationToken); try { if (IsConnected && !forceReconnect) return; await DisconnectTransportAsync(); _transport = new SyncCore.MqttnetTransport(new SyncCore.MqttBrokerConfig { Host = _settings.MqttHost, Port = _settings.MqttPort, WebSocketUri = _settings.MqttPort == 443 ? $"wss://{_settings.MqttHost}/mqtt" : null, ClientId = _settings.MqttClientId, Username = _settings.MqttOrgId, Password = password, UseTls = true, }); _engine = new SyncCore.MqttSyncEngine( _transport, _store, new SyncCore.SyncOptions { OrgId = _settings.MqttOrgId, ClientId = _settings.MqttClientId, }); _engine.PointAccepted += HandleAcceptedAsync; _engine.PointRejected += HandleRejectedAsync; await _engine.ConnectAsync(cancellationToken); SetStatus("Connected"); } finally { _connectionGate.Release(); } } private async Task RunAsync(CancellationToken cancellationToken) { using var timer = new PeriodicTimer(TimeSpan.FromSeconds(5)); while (await timer.WaitForNextTickAsync(cancellationToken)) { if (!_settings.MqttEnabled) { SetStatus("Sync disabled"); continue; } try { if (!IsConnected) { SetStatus("Connecting…"); await ConnectCoreAsync(forceReconnect: false, cancellationToken); await QueueUnsyncedAsync(cancellationToken); } await DrainAsync(cancellationToken); } catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) { throw; } catch (Exception ex) { _logger.LogWarning(ex, "MQTT sync iteration failed; queued data retained"); SetStatus($"Offline — retrying: {ex.Message}"); } } } private async Task DrainAsync(CancellationToken cancellationToken) { if (_engine is null || !IsConnected) return; await _pumpGate.WaitAsync(cancellationToken); try { int published = await _engine.DrainOnceAsync(cancellationToken); if (published > 0) SetStatus($"Connected — {published} point(s) awaiting ack"); } finally { _pumpGate.Release(); } } private async Task HandleAcceptedAsync(string pointId) { await _database.MarkPointSyncedAsync(pointId); SetStatus("Connected — synced"); } private async Task HandleRejectedAsync(SyncCore.OutboundMessage message) { string reason = message.LastReason ?? "REJECTED"; await _database.MarkPointSyncErrorAsync(message.PointId, reason); var point = await _database.GetPointBySyncIdAsync(message.PointId); if (point is not null) PointRejected?.Invoke(this, new PointSyncFailureEventArgs(point, reason)); SetStatus($"Point rejected: {reason}"); } private async Task RecordFailureAsync(LoggedPoint point, string reasonCode) { point.SyncError = reasonCode; point.Synced = false; await _database.UpdatePointAsync(point); PointRejected?.Invoke(this, new PointSyncFailureEventArgs(point, reasonCode)); SetStatus($"Point not queued: {reasonCode}"); } private async Task DisconnectTransportAsync() { if (_transport is null) return; try { await _transport.DisposeAsync(); } finally { _transport = null; _engine = null; } } private void SetStatus(string value) { if (Status == value) return; Status = value; StatusChanged?.Invoke(this, value); } private static string TicketFor(Job job) { string value = $"FL-{job.Id}-{job.Name}".Trim(); return value.Length <= 64 ? value : value[..64]; } private static SyncCore.PointRecord ToSyncPoint(LoggedPoint point) { var timestamp = new DateTimeOffset(DateTime.SpecifyKind(point.TimestampUtc, DateTimeKind.Utc)); DateTimeOffset? positionEpoch = point.GpsUnixTimestamp > 0 ? DateTimeOffset.FromUnixTimeSeconds(point.GpsUnixTimestamp) : null; return new SyncCore.PointRecord { PointId = point.SyncPointId!, Origin = SyncCore.PointOrigin.APP, UploadPath = SyncCore.UploadPath.APP_MQTT, CaptureTrigger = SyncCore.CaptureTrigger.LOCATOR_BUTTON, CreatedAt = timestamp, Position = point.GpsValid ? new SyncCore.PositionGroup { Lat = point.Latitude, Lon = point.Longitude, EllipsoidalHeight = point.Altitude, OrthometricHeight = point.AltitudeCorrected, PositionEpoch = positionEpoch, } : null, Gnss = new SyncCore.GnssGroup { FixType = FixType(point.FixStatusEnum), SatsUsed = point.SatellitesUsed, Hdop = point.Hdop, Hrms = point.Hrms, Vrms = point.Vrms, CorrectionAge = point.CorrectionAgeSeconds, ReceiverSerial = point.GpsSerialNumber, TiltAngle = point.TiltAngle, Source = string.IsNullOrWhiteSpace(point.GpsSerialNumber) ? "PHONE" : "MAGLINK", }, Locate = new SyncCore.LocateGroup { Depth = point.DepthMeters, DepthUnits = point.DepthMeters is null ? null : "m", SignalCurrent = point.CurrentMilliamps, SignalStrength = point.Signal, Frequency = point.Frequency, Gain = point.GainDb, LocateMode = LocateMode(point), PhaseDegrees = point.LdPhase, CompassDegrees = point.CompassAngle, LocatorModel = point.LocatorModel, LocatorSerial = point.LocatorSerialNumber, TelemetryEpoch = timestamp, }, Attributes = new SyncCore.AttributesGroup { UtilityType = UtilityType(point.UtilityEnum) }, Quality = new SyncCore.QualityGroup { QualityFlag = point.GpsValid && point.Hrms > 0 && point.Hrms <= 0.10 ? "IN_SPEC" : "OUT_OF_SPEC", }, }; } private static string FixType(GnssFixStatus status) => status switch { GnssFixStatus.Single => "AUTONOMOUS", GnssFixStatus.Dgps => "DGPS", GnssFixStatus.RtkFloat => "FLOAT", GnssFixStatus.RtkFixed => "FIXED", _ => "NO_FIX", }; private static string UtilityType(UmUtility utility) => utility switch { UmUtility.Gas => "GAS", UmUtility.Power => "ELECTRIC", UmUtility.Communications => "TELECOM", UmUtility.Water => "WATER", UmUtility.Sewer => "SEWER", UmUtility.Fiber => "FIBER", _ => "UNKNOWN", }; private static string LocateMode(LoggedPoint point) => point.FreqType == (int)UmFreqType.Sonde ? "SONDE" : (UmMode)point.Mode switch { UmMode.Null => "NULL", UmMode.Omni or UmMode.TwinOmni => "BROAD_PEAK", _ => "PEAK", }; }