// Sync/SyncManager.swift // ---------------------- // Manages connectivity monitoring, reference data sync (Phase A), // the outbox queue for offline inspection/issue submission (Phase B), // and notification polling (Phase C). import Foundation import Network import SwiftData import SwiftUI import Combine import UserNotifications @MainActor class SyncManager: ObservableObject { // ── Published State ─────────────────────────────────────────────────── @Published var isOnline = false @Published var isSyncing = false @Published var lastSyncAt: Date? @Published var syncError: String? @Published var pendingCount = 0 // ── Dependencies ────────────────────────────────────────────────────── private let monitor = NWPathMonitor() private let monitorQueue = DispatchQueue(label: "com.jqc.networkmonitor") var modelContext: ModelContext? // ── Notification polling state ──────────────────────────────────────── // Tracks the timestamp of the most recently fetched notification so each // poll only retrieves newer records. Nil on first launch → server returns // last 50 unread. Reset to nil on logout. private var lastNotificationFetch: Date? private var pollTask: Task? // replaces Timer — Task.sleep works correctly private let pollInterval: UInt64 = 60_000_000_000 // 60 seconds in nanoseconds // ── Shared date formatters ──────────────────────────────────────────── // DateFormatter init is expensive — allocating one per poll call or per // issue would add measurable overhead at sync time. These are created // once and reused across all calls. Both are nonisolated statics so they // can be read from any context without actor-hopping. // // isoFormatter — parses/formats ISO 8601 strings from the server API // e.g. "2026-05-01T14:30:00" // notifFormatter — same format, used to advance the notification poll cursor nonisolated static let isoFormatter: DateFormatter = { let f = DateFormatter() f.locale = Locale(identifier: "en_US_POSIX") f.dateFormat = "yyyy-MM-dd'T'HH:mm:ss" return f }() static let shared = SyncManager() private init() {} // ── Start Monitoring ────────────────────────────────────────────────── func startMonitoring() { monitor.pathUpdateHandler = { [weak self] path in Task { @MainActor [weak self] in guard let self else { return } let wasOffline = !self.isOnline self.isOnline = path.status == .satisfied if self.isOnline { if wasOffline { await self.triggerSync() } self.startPollTask() } else { self.stopPollTask() } } } monitor.start(queue: monitorQueue) } // ── Notification poll task ──────────────────────────────────────────── // Timer.scheduledTimer requires RunLoop.main to be ticking. When called // from inside a Swift Concurrency Task { @MainActor } the current RunLoop // is NOT RunLoop.main — the timer is added to a runloop that never runs, // so it silently fires never. Task + Task.sleep has no such dependency. private func startPollTask() { guard pollTask == nil else { return } // already running pollTask = Task { [weak self] in while !Task.isCancelled { try? await Task.sleep(nanoseconds: 60_000_000_000) guard !Task.isCancelled else { break } await MainActor.run { [weak self] in guard let self, self.isOnline, AuthManager.shared.isAuthenticated else { return } Task { await self.pollNotifications() } } } } } private func stopPollTask() { pollTask?.cancel() pollTask = nil } /// Called on logout so the next login starts a clean fetch. func resetNotificationPoller() { lastNotificationFetch = nil stopPollTask() } // ── Notification polling ────────────────────────────────────────────── func pollNotifications() async { guard isOnline, AuthManager.shared.isAuthenticated else { return } do { let notifications = try await APIClient.shared.fetchNotifications(since: lastNotificationFetch) guard !notifications.isEmpty else { return } // Deliver a local notification for each new item for n in notifications { deliverLocalNotification(n) } // Update the cursor to the newest notification's timestamp let dates = notifications.compactMap { Self.isoFormatter.date(from: $0.createdAt) } if let newest = dates.max() { lastNotificationFetch = newest } // Mark all fetched notifications as read on the server let ids = notifications.map(\.id) try await APIClient.shared.markNotificationsRead(ids: ids) } catch APIError.notAuthenticated { // Token expired and refresh failed — let AuthManager handle it } catch { // Network errors are silent; next poll will retry } } // ── Local notification delivery ─────────────────────────────────────── private func deliverLocalNotification(_ n: APINotification) { let content = UNMutableNotificationContent() content.title = n.title content.body = n.body content.sound = .default // Use the server notification ID as the identifier so duplicate // deliveries (if the same record is fetched twice) replace rather // than stack. let identifier = "jqc-notif-\(n.id)" let request = UNNotificationRequest( identifier: identifier, content: content, trigger: nil // nil = deliver immediately ) UNUserNotificationCenter.current().add(request) { error in if let error { print("[JQC] Local notification delivery failed: \(error)") } } } // ── Full Sync ───────────────────────────────────────────────────────── func triggerSync() async { // Do not sync unless authenticated — avoids 401 loops before // restoreSession() completes on first launch. guard isOnline, let context = modelContext, AuthManager.shared.isAuthenticated else { return } isSyncing = true syncError = nil defer { isSyncing = false } await processPhotoQueue(context: context) await processInspectionQueue(context: context) await processIssueQueue(context: context) await pullReferenceData() await pullAssignedIssues(context: context) // Poll notifications immediately on every sync rather than waiting // for the 60-second timer — ensures the inspector sees assignments // and follow-up requests as soon as the app goes online. await pollNotifications() updatePendingCount(context: context) lastSyncAt = Date() } // ── Outbox: Photos ──────────────────────────────────────────────────── private func processPhotoQueue(context: ModelContext) async { // Fetch all then filter in Swift — #Predicate cannot reference // string literals against PendingPhoto.uploadStatus reliably // when the predicate type is inferred across model boundaries. guard let allPhotos = try? context.fetch(FetchDescriptor()) else { return } let pending = allPhotos .filter { $0.uploadStatus == "pending" } .sorted { $0.createdAt < $1.createdAt } for photo in pending { do { let serverPath = try await APIClient.shared.uploadPhoto( localPath: photo.localFilePath, entityType: photo.entityType ) photo.serverPath = serverPath photo.uploadStatus = "uploaded" // Update parent inspection form field value if photo.entityType == "inspection", let fieldId = photo.fieldId { let entityId = photo.entityLocalId // Fetch-all + filter in Swift — #Predicate with a captured String // variable causes "LocalInspection is ambiguous" under Xcode 26 // SWIFT_DEFAULT_ACTOR_ISOLATION = MainActor (CLAUDE.md rule 3). let allInspections = (try? context.fetch(FetchDescriptor())) ?? [] allInspections.first(where: { $0.localId == entityId })?.setValue(serverPath, forFieldId: fieldId) } // Update parent issue photo paths array if photo.entityType == "issue" { let entityId = photo.entityLocalId let allIssues = (try? context.fetch(FetchDescriptor())) ?? [] if let issue = allIssues.first(where: { $0.localId == entityId }) { var paths = issue.photoServerPaths if !paths.contains(serverPath) { paths.append(serverPath) } issue.photoServerPaths = paths } } try? context.save() } catch { photo.uploadStatus = "failed" try? context.save() } } } // ── Outbox: Inspections ─────────────────────────────────────────────── private func processInspectionQueue(context: ModelContext) async { guard let all = try? context.fetch(FetchDescriptor()) else { return } let pending = all .filter { $0.status == "completed" && $0.syncStatus == "pending" } .sorted { $0.createdAt < $1.createdAt } for inspection in pending { let photosReady = inspection.pendingPhotos.allSatisfy { $0.uploadStatus == "uploaded" || $0.uploadStatus == "failed" } guard photosReady else { continue } do { let inspectionId = try await APIClient.shared.submitInspection(inspection) inspection.serverId = inspectionId inspection.syncStatus = "synced" inspection.status = "synced" // Clear follow-up flag on parent. if let parentLocalId = inspection.parentLocalId { let allInspections = try? context.fetch(FetchDescriptor()) if let parent = allInspections?.first(where: { $0.localId == parentLocalId }) { parent.followUpRequired = false parent.followUpNote = nil } } try? context.save() } catch { inspection.syncRetryCount += 1 inspection.syncErrorMessage = error.localizedDescription if inspection.syncRetryCount >= 5 { inspection.syncStatus = "failed" } syncError = "Failed to sync inspection: \(error.localizedDescription)" try? context.save() } } } // ── Outbox: Issues ──────────────────────────────────────────────────── private func processIssueQueue(context: ModelContext) async { guard let all = try? context.fetch(FetchDescriptor()) else { return } let pending = all .filter { $0.syncStatus == "pending" } .sorted { $0.createdAt < $1.createdAt } // Pre-fetch all inspections to check parent sync status. let allInspections = (try? context.fetch(FetchDescriptor())) ?? [] for issue in pending { // Guard: if the parent inspection permanently failed to sync, // submitting this issue without an inspection_id would create an // orphaned server record. Mark it failed immediately instead. let parentId = issue.inspectionLocalId let parent = allInspections.first(where: { $0.localId == parentId }) if parent?.syncStatus == "failed" { issue.syncStatus = "failed" issue.syncErrorMessage = "Parent inspection failed to sync — issue cannot be submitted." try? context.save() syncError = "Issue \(issue.localId.prefix(8))\u{2026} blocked: parent inspection did not sync." continue } do { let issueId = try await APIClient.shared.submitIssue(issue) issue.serverId = issueId issue.syncStatus = "synced" try? context.save() } catch { issue.syncRetryCount += 1 issue.syncErrorMessage = error.localizedDescription if issue.syncRetryCount >= 5 { issue.syncStatus = "failed" } try? context.save() } } } // ── Reference Data ──────────────────────────────────────────────────── func pullReferenceData() async { guard isOnline, let context = modelContext else { return } do { let facilitiesData: FacilitiesResponseData = try await APIClient.shared.request("/api/v1/facilities") let templatesData: TemplatesResponseData = try await APIClient.shared.request("/api/v1/templates") let existingFacilities = try context.fetch(FetchDescriptor()) let facilityMap = Dictionary( existingFacilities.map { ($0.serverId, $0) }, uniquingKeysWith: { a, _ in a } ) for apiFacility in facilitiesData.facilities { if let existing = facilityMap[apiFacility.id] { existing.update(from: apiFacility) } else { context.insert(LocalFacility(from: apiFacility)) } try await upsertAreas(for: apiFacility.id, context: context) } let existingTemplates = try context.fetch(FetchDescriptor()) let templateMap = Dictionary( existingTemplates.map { ($0.serverId, $0) }, uniquingKeysWith: { a, _ in a } ) for apiSummary in templatesData.templates { if let existing = templateMap[apiSummary.id] { existing.updateSummary(from: apiSummary) } else { context.insert(LocalTemplate(from: apiSummary)) } try await upsertTemplateSchema( id: apiSummary.id, context: context, templateMap: templateMap ) } try context.save() } catch APIError.notAuthenticated { syncError = "Session expired. Please log in again." } catch { syncError = "Sync failed: \(error.localizedDescription)" } } // ── Pending Count ───────────────────────────────────────────────────── func updatePendingCount(context: ModelContext) { let inspCount = (try? context.fetch(FetchDescriptor()))? .filter { $0.syncStatus == "pending" }.count ?? 0 let issueCount = (try? context.fetch(FetchDescriptor()))? .filter { $0.syncStatus == "pending" }.count ?? 0 pendingCount = inspCount + issueCount } // ── Private Helpers ─────────────────────────────────────────────────── private func upsertAreas(for facilityId: Int, context: ModelContext) async throws { let areasData: AreasResponseData = try await APIClient.shared.request( "/api/v1/facilities/\(facilityId)/areas" ) let existing = (try? context.fetch(FetchDescriptor()))? .filter { $0.facilityServerId == facilityId } ?? [] let areaMap = Dictionary(existing.map { ($0.serverId, $0) }, uniquingKeysWith: { a, _ in a }) for apiArea in areasData.areas { if let ex = areaMap[apiArea.id] { ex.update(from: apiArea) } else { context.insert(LocalArea(from: apiArea)) } } } private func upsertTemplateSchema( id: Int, context: ModelContext, templateMap: [Int: LocalTemplate] ) async throws { let detailData: TemplateDetailResponseData = try await APIClient.shared.request( "/api/v1/templates/\(id)" ) if let existing = templateMap[id] { existing.updateSchema(from: detailData.template) } else { (try? context.fetch(FetchDescriptor()))? .first { $0.serverId == id }? .updateSchema(from: detailData.template) } } // ── Pull server-assigned issues ─────────────────────────────────────── // Fetches issues assigned to the current user on the server and upserts // them into SwiftData so IssuesListView shows them alongside device-created issues. // Keyed by serverId — existing records are updated in-place, new ones inserted. // These records carry syncStatus = "synced" and a generated localId so they // are never re-submitted to the server by processIssueQueue. func pullAssignedIssues(context: ModelContext) async { guard isOnline, AuthManager.shared.isAuthenticated else { return } do { let apiIssues = try await APIClient.shared.fetchAssignedIssues() // Do not return early on empty — deletion still needs to run // to remove issues that were unassigned from this inspector. // Build a map of existing LocalIssues by serverId for upsert let allLocal = (try? context.fetch(FetchDescriptor())) ?? [] var serverIdMap: [Int: LocalIssue] = [:] for local in allLocal { if let sid = local.serverId { serverIdMap[sid] = local } } for api in apiIssues { if let existing = serverIdMap[api.id] { // Update mutable fields on existing record existing.issueStatus = api.status existing.severity = api.severity existing.issueDescription = api.description if let fid = api.facilityId { existing.facilityServerId = fid } // Refresh photos in case they were added after first pull var serverPaths: [String] = [] if let p = api.photoPath, !p.isEmpty { serverPaths.append(p) } serverPaths.append(contentsOf: api.resultPhotos) existing.photoServerPaths = serverPaths } else { // Insert new server-pulled issue let local = LocalIssue( inspectionLocalId: "", facilityServerId: api.facilityId ?? 0, severity: api.severity, description: api.description ) local.serverId = api.id local.issueStatus = api.status local.syncStatus = "synced" // never re-submit // Store server photos so IssueDetailView can show them var serverPaths: [String] = [] if let p = api.photoPath, !p.isEmpty { serverPaths.append(p) } serverPaths.append(contentsOf: api.resultPhotos) local.photoServerPaths = serverPaths if let ts = api.reportedAt, let date = Self.isoFormatter.date(from: ts) { local.createdAt = date } context.insert(local) } } // Remove server-pulled records that are no longer in the response. // This happens when an issue is reassigned to a different inspector — // the server stops returning it for this user, so the local copy must // be deleted. Only remove records that were pulled from the server // (syncStatus == "synced" AND serverId != nil AND inspectionLocalId == ""). // Device-created issues (inspectionLocalId != "") are never touched. let returnedServerIds = Set(apiIssues.map { $0.id }) for local in allLocal { guard let sid = local.serverId, local.syncStatus == "synced", local.inspectionLocalId == "" else { continue } if !returnedServerIds.contains(sid) { context.delete(local) } } try? context.save() } catch APIError.notAuthenticated { // Let AuthManager handle session expiry } catch { // Non-fatal — IssuesListView still shows device-created issues } } }