495 lines
22 KiB
Swift
495 lines
22 KiB
Swift
// 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<Void, Never>? // replaces Timer — Task.sleep works correctly
|
|
private let pollInterval: UInt64 = 60_000_000_000 // 60 seconds in nanoseconds
|
|
|
|
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 fmt = DateFormatter()
|
|
fmt.dateFormat = "yyyy-MM-dd'T'HH:mm:ss"
|
|
let dates = notifications.compactMap { fmt.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<PendingPhoto>()) 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<LocalInspection>())) ?? []
|
|
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<LocalIssue>())) ?? []
|
|
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<LocalInspection>()) 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<LocalInspection>())
|
|
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<LocalIssue>()) 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<LocalInspection>())) ?? []
|
|
|
|
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<LocalFacility>())
|
|
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<LocalTemplate>())
|
|
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<LocalInspection>()))?
|
|
.filter { $0.syncStatus == "pending" }.count ?? 0
|
|
let issueCount = (try? context.fetch(FetchDescriptor<LocalIssue>()))?
|
|
.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<LocalArea>()))?
|
|
.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<LocalTemplate>()))?
|
|
.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<LocalIssue>())) ?? []
|
|
var serverIdMap: [Int: LocalIssue] = [:]
|
|
for local in allLocal {
|
|
if let sid = local.serverId { serverIdMap[sid] = local }
|
|
}
|
|
|
|
let isoFmt = ISO8601DateFormatter()
|
|
|
|
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 = isoFmt.date(from: ts) ?? {
|
|
let f = DateFormatter()
|
|
f.dateFormat = "yyyy-MM-dd'T'HH:mm:ss"
|
|
return f.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
|
|
}
|
|
}
|
|
}
|