SyncEngine.swift
1064 lignes · 48640 octets
import Foundation struct SyncStatusResult { var pendingRemoteChanges: Int var response: SyncStatusResponse? = nil var error: String? = nil } /// Op refused by the server and left queued. `message` is the verbatim server message /// (« Réservation refusée : période en conflit avec Entretien »), `nil` when the op never /// even reached the network (unresolved local reference); `code` is -1 for a network failure. struct PushFailure: Equatable { var entityType: String var op: String var code: Int? = nil var message: String? = nil } struct SyncResult { var success: Bool var serverTime: String? = nil var error: String? = nil var conflicts: [AgendaConflict] = [] var pushed: Int = 0 var received: Int = 0 var pushFailures: [PushFailure] = [] } /// Pushes the `SyncOpEntity` queue then pulls `/api/sync/pull`; LWW / append / tombstones. /// The calendar↔store bridge (RDV + reservations/unavailabilities bound to device calendars) /// is delegated to `AgendaSyncCoordinator`; this engine orchestrates the step ordering and /// the push (including the reservation 409 special case). final class SyncEngine { private static let watermarkKey = "watermark" private static let watermarkOriginKey = "watermark_origin" private static let defaultWatermark = "1970-01-01T00:00:00Z" private static let entityContactMedia = "contact_media" private let api: AilianceApi private let db: CrmDatabase private let imageStore: ContactImageSaving? /// « serveur|utilisateur » identity the watermark belongs to. When it differs from the /// stored one (server or account change), the pull restarts from the epoch: otherwise a /// watermark inherited from another server hides any older data (e.g. backdated demo /// data → « 0 reçu » forever). `nil` = no check. private let serverIdentity: String? private let agenda: AgendaSyncCoordinator init( api: AilianceApi, db: CrmDatabase, imageStore: ContactImageSaving? = nil, serverIdentity: String? = nil ) { self.api = api self.db = db self.imageStore = imageStore self.serverIdentity = serverIdentity self.agenda = AgendaSyncCoordinator(db: db) } func checkStatus(ressourcesQuery: String? = nil) async throws -> SyncStatusResult { switch await api.syncStatus(sinceIso: try await watermark(), ressourcesQuery: ressourcesQuery) { case .ok(let response): return SyncStatusResult(pendingRemoteChanges: response.total, response: response) case .err(_, let message): return SyncStatusResult(pendingRemoteChanges: 0, error: message) } } /// Empty `bindings`/nil `bridge`: unchanged CRM v1 path (no calendar bridge, `ressources=` omitted). func syncNow( bindings: [CalendarBinding] = [], bridge: CalendarBridge? = nil ) async throws -> SyncResult { let agendaEnabled = bridge != nil && !bindings.isEmpty if agendaEnabled { try await agenda.agendaToRoom(bindings: bindings, bridge: bridge!) } let push = try await pushOps() let ressourcesQuery = agendaEnabled ? agenda.ressourcesQuery(bindings) : nil switch await api.syncPull(sinceIso: try await watermark(), ressourcesQuery: ressourcesQuery) { case .ok(let pull): let received = try await applyPull(pull) if agendaEnabled { try await agenda.roomToAgenda(bindings: bindings, bridge: bridge!) } try await setWatermark(pull.serverTime) return SyncResult( success: true, serverTime: pull.serverTime, conflicts: push.conflicts, pushed: push.succeeded, received: received, pushFailures: push.failures ) case .err(_, let message): return SyncResult( success: false, error: message, conflicts: push.conflicts, pushed: push.succeeded, pushFailures: push.failures ) } } /// Vide la file de push sans effectuer de pull (mode « transcription serveur » : /// l'audio part dès la fin de l'enregistrement). Équivalent de `pousserEnAttente` Android. func pousserEnAttente() async throws -> SyncResult { let push = try await pushOps() return SyncResult( success: true, pushed: push.succeeded, pushFailures: push.failures ) } /// Relance la transcription d'une interaction sur le serveur (appel direct, hors file). /// Retourne `nil` en cas de succès, le message d'erreur sinon. func relancerTranscriptionInteraction(localId: Int64) async throws -> String? { guard let interaction = try await db.interactionDao.getByLocalId(localId), let interactionId = interaction.serverId else { return "Impossible de retrouver l'interaction." } let result = await api.relancerTranscriptionInteraction( contactId: interaction.contactServerId, interactionId: interactionId ) switch result { case .ok: var updated = interaction updated.transcriptionStatut = "en_attente" updated.transcriptionErreur = nil try await db.interactionDao.upsert(updated) return nil case .err(_, let message): return "Relance impossible : \(message)" } } /// Relance la transcription d'une note de projet sur le serveur (appel direct, hors file). /// Retourne `nil` en cas de succès, le message d'erreur sinon. func relancerTranscriptionNoteProjet(localId: Int64) async throws -> String? { guard let note = try await db.noteProjetDao.listAll().first(where: { $0.localId == localId }), let noteId = note.serverId else { return "Impossible de retrouver la note." } let result = await api.relancerTranscriptionNoteProjet( projetId: note.projetServerId, noteId: noteId ) switch result { case .ok: var updated = note updated.transcriptionStatut = "en_attente" updated.transcriptionErreur = nil try await db.noteProjetDao.upsert(updated) return nil case .err(_, let message): return "Relance impossible : \(message)" } } /// Count of local items not yet pushed (« calendrier local en avance » banner). /// With `bindings`/`bridge`: adds card2vcf-tagged calendar events without a local entity /// (orphans not yet cleaned up on the calendar side, see `AgendaSyncCoordinator.orphanEventCount`). func localAheadCount( bindings: [CalendarBinding] = [], bridge: CalendarBridge? = nil ) async throws -> Int { let dirtyRdv = try await db.rdvDao.listDirty().count let dirtyReservations = try await db.reservationDao.listDirty().count let pendingAgendaOps = try await db.syncOpDao.listAll().filter { $0.entityType == AgendaSyncCoordinator.kindRdv || $0.entityType == AgendaSyncCoordinator.entityReservation }.count let orphanEvents: Int if let bridge, !bindings.isEmpty { orphanEvents = try await agenda.orphanEventCount(bindings: bindings, bridge: bridge) } else { orphanEvents = 0 } return dirtyRdv + dirtyReservations + pendingAgendaOps + orphanEvents } /// The user gives up their conflicting reservation: purge entity + SyncOp + calendar event. func abandonLocalReservation(localId: Int64, bridge: CalendarBridge? = nil) async throws { guard let reservation = try await db.reservationDao.getByLocalId(localId) else { return } if let eventId = reservation.calendarEventId, let bridge { try bridge.deleteEvent(eventId: eventId) } for op in try await db.syncOpDao.listAll() where op.entityType == AgendaSyncCoordinator.entityReservation && op.localId == localId { try await db.syncOpDao.deleteById(op.id) } try await db.reservationDao.deleteByLocalId(localId) } // ---- watermark ---- private func watermark() async throws -> String { guard let stored = try await db.syncMetaDao.get(Self.watermarkKey)?.value else { return Self.defaultWatermark } if let serverIdentity { let origin = try await db.syncMetaDao.get(Self.watermarkOriginKey)?.value if origin != serverIdentity { return Self.defaultWatermark } } return stored } private func setWatermark(_ serverTime: String) async throws { try await db.syncMetaDao.upsert(SyncMetaEntity(key: Self.watermarkKey, value: serverTime)) if let serverIdentity { try await db.syncMetaDao.upsert(SyncMetaEntity(key: Self.watermarkOriginKey, value: serverIdentity)) } } // ---- push ---- /// Outcome of a pushed op. `conflict` is non-nil only for a reservation 409 (`ok` stays /// `false`): that case is arbitrated by the dedicated dialog, not a generic failure. /// `code`/`message` carry the server response for every other refusal, so it can be /// surfaced to the user. private struct PushOutcome { var ok: Bool var conflict: AgendaConflict? = nil var code: Int? = nil var message: String? = nil } /// Summary of a push cycle: attempted / dequeued ops, 409 conflicts and generic failures. private struct PushReport { var attempted: Int var succeeded: Int var conflicts: [AgendaConflict] var failures: [PushFailure] } /// Keeps the server `code`/`message` on refusal, so `pushOps` can surface them to the UI. private static func outcome<T>(_ result: AilianceApiClient.ApiResult<T>) -> PushOutcome { switch result { case .ok: return PushOutcome(ok: true) case .err(let code, let message): return PushOutcome(ok: false, code: code, message: message) } } private func pushOps() async throws -> PushReport { var conflicts: [AgendaConflict] = [] var failures: [PushFailure] = [] var attempted = 0 var succeeded = 0 for op in try await db.syncOpDao.listAll() { attempted += 1 let outcome: PushOutcome switch op.entityType { case "contact": outcome = try await pushContactOp(op) case Self.entityContactMedia: outcome = try await pushContactMediaOp(op) case "entreprise": outcome = try await pushEntrepriseOp(op) case "projet": outcome = try await pushProjetOp(op) case "tache": outcome = try await pushTacheOp(op) case "interaction": outcome = try await pushInteractionOp(op) case "interaction_audio": outcome = try await pushInteractionAudioOp(op) case "note_projet": outcome = try await pushNoteProjetOp(op) case "note_projet_audio": outcome = try await pushNoteProjetAudioOp(op) case AgendaSyncCoordinator.kindRdv: outcome = try await pushRdvOp(op) case AgendaSyncCoordinator.entityReservation: outcome = try await pushReservationOp(op) default: outcome = PushOutcome(ok: true) } if outcome.ok { succeeded += 1 try await db.syncOpDao.deleteById(op.id) } else if outcome.conflict == nil { failures.append( PushFailure(entityType: op.entityType, op: op.op, code: outcome.code, message: outcome.message) ) let error = outcome.message ?? outcome.code.map { "HTTP \($0)" } ?? "réseau" try await db.syncOpDao.markFailure(id: op.id, error: error) } if let conflict = outcome.conflict { conflicts.append(conflict) } } return PushReport(attempted: attempted, succeeded: succeeded, conflicts: conflicts, failures: failures) } private func pushContactOp(_ op: SyncOpEntity) async throws -> PushOutcome { switch op.op { case "create": let result = await api.createContact(jsonBody: op.payloadJson) if case .ok(let body) = result { if let newServerId = Self.parseCreatedId(body), let localId = op.localId { if var contact = try await db.crmContactDao.getById(localId) { contact.serverId = newServerId try await db.crmContactDao.update(contact) } if try await !pushContactImages(localId: localId, serverId: newServerId) { try await enqueueContactMediaRetry(localId: localId, serverId: newServerId) } } } return Self.outcome(result) case "update": guard let serverId = op.serverId else { return PushOutcome(ok: false) } let result = await api.updateContact(id: serverId, jsonBody: op.payloadJson) guard result.isOk else { return Self.outcome(result) } let localId: Int64? if let opLocalId = op.localId { localId = opLocalId } else { localId = try await db.crmContactDao.getByServerId(serverId)?.id } if let localId, try await !pushContactImages(localId: localId, serverId: serverId) { try await enqueueContactMediaRetry(localId: localId, serverId: serverId) } return PushOutcome(ok: true) case "delete": guard let serverId = op.serverId else { return PushOutcome(ok: false) } return Self.outcome(await api.deleteContact(id: serverId)) default: return PushOutcome(ok: true) } } private func pushContactMediaOp(_ op: SyncOpEntity) async throws -> PushOutcome { guard let serverId = op.serverId else { return PushOutcome(ok: true) } let localId: Int64? if let opLocalId = op.localId { localId = opLocalId } else { localId = try await db.crmContactDao.getByServerId(serverId)?.id } guard let localId else { return PushOutcome(ok: true) } return PushOutcome(ok: try await pushContactImages(localId: localId, serverId: serverId)) } /// Uploads the local card/photo files to the server. `true` when nothing to send or all OK. /// `false` when at least one present file failed to upload (retry via `contact_media`). private func pushContactImages(localId: Int64, serverId: String) async throws -> Bool { guard let contact = try await db.crmContactDao.getById(localId) else { return true } var ok = true if let cardPath = nonBlank(contact.cardImagePath), let bytes = Self.readFile(cardPath) { let result = await api.uploadContactCarte(id: serverId, bytes: bytes, filename: "card.jpg") if !result.isOk { ok = false } } if let profilePath = nonBlank(contact.profileImagePath), let bytes = Self.readFile(profilePath) { let result = await api.uploadContactPhoto(id: serverId, bytes: bytes, filename: "profile.jpg") if !result.isOk { ok = false } } return ok } private func enqueueContactMediaRetry(localId: Int64, serverId: String) async throws { let already = try await db.syncOpDao.listAll().contains { $0.entityType == Self.entityContactMedia && $0.localId == localId && $0.serverId == serverId } if already { return } var entity = SyncOpEntity(entityType: Self.entityContactMedia, op: "upload") entity.payloadJson = "{}" entity.localId = localId entity.serverId = serverId entity.createdAt = AgendaSyncCoordinator.currentMillis() try await db.syncOpDao.insert(entity) } private func pushEntrepriseOp(_ op: SyncOpEntity) async throws -> PushOutcome { switch op.op { case "create": let result = await api.createEntreprise(jsonBody: op.payloadJson) if case .ok(let body) = result { if let newServerId = Self.parseCreatedId(body), let localId = op.localId { if var entity = try await db.entrepriseDao.getByLocalId(localId) { entity.serverId = newServerId try await db.entrepriseDao.upsert(entity) } } } return Self.outcome(result) case "update": guard let serverId = op.serverId else { return PushOutcome(ok: false) } return Self.outcome(await api.updateEntreprise(id: serverId, jsonBody: op.payloadJson)) case "delete": guard let serverId = op.serverId else { return PushOutcome(ok: false) } return Self.outcome(await api.deleteEntreprise(id: serverId)) default: return PushOutcome(ok: true) } } private func pushProjetOp(_ op: SyncOpEntity) async throws -> PushOutcome { switch op.op { case "create": let result = await api.createProjet(jsonBody: op.payloadJson) if case .ok(let body) = result { if let newServerId = Self.parseCreatedId(body), let localId = op.localId { if var entity = try await db.projetDao.getByLocalId(localId) { entity.serverId = newServerId try await db.projetDao.upsert(entity) } } } return Self.outcome(result) case "update": guard let serverId = op.serverId else { return PushOutcome(ok: false) } return Self.outcome(await api.updateProjet(id: serverId, jsonBody: op.payloadJson)) case "delete": guard let serverId = op.serverId else { return PushOutcome(ok: false) } return Self.outcome(await api.deleteProjet(id: serverId)) default: return PushOutcome(ok: true) } } /// Resolves the `projetServerId` of a task via `op.localId` (create) or `op.serverId` (update/delete/move). private func resolveTacheProjetId(_ op: SyncOpEntity) async throws -> String? { if let localId = op.localId, let tache = try await db.tacheDao.getByLocalId(localId) { return tache.projetServerId } if let serverId = op.serverId, let tache = try await db.tacheDao.getByServerId(serverId) { return tache.projetServerId } return nil } private func pushTacheOp(_ op: SyncOpEntity) async throws -> PushOutcome { guard let projetId = try await resolveTacheProjetId(op) else { return PushOutcome(ok: false) } switch op.op { case "create": let result = await api.createTache(projetId: projetId, jsonBody: op.payloadJson) if case .ok(let body) = result { if let newServerId = Self.parseCreatedId(body), let localId = op.localId { if var entity = try await db.tacheDao.getByLocalId(localId) { entity.serverId = newServerId try await db.tacheDao.upsert(entity) } } } return Self.outcome(result) case "update": guard let serverId = op.serverId else { return PushOutcome(ok: false) } return Self.outcome(await api.updateTache(projetId: projetId, tacheId: serverId, jsonBody: op.payloadJson)) case "delete": guard let serverId = op.serverId else { return PushOutcome(ok: false) } return Self.outcome(await api.deleteTache(projetId: projetId, tacheId: serverId)) case "move": guard let serverId = op.serverId else { return PushOutcome(ok: false) } return Self.outcome(await api.moveTache(projetId: projetId, tacheId: serverId, jsonBody: op.payloadJson)) case "soustache": // `payloadJson` porte l'identifiant de la sous-tâche à basculer. guard let serverId = op.serverId else { return PushOutcome(ok: false) } return Self.outcome( await api.toggleSousTache( projetId: projetId, tacheId: serverId, sousTacheId: op.payloadJson ) ) default: return PushOutcome(ok: true) } } /// `op.serverId` = contactServerId pour create, interactionServerId pour delete. /// `op.localId` = interaction locale (create uniquement). private func pushInteractionOp(_ op: SyncOpEntity) async throws -> PushOutcome { switch op.op { case "create": guard let contactServerId = op.serverId else { return PushOutcome(ok: false) } let result = await api.createInteraction(contactId: contactServerId, jsonBody: op.payloadJson) if case .ok(let body) = result, let localId = op.localId { if let interactionServerId = Self.parseCreatedId(body) { if var entity = try await db.interactionDao.getByLocalId(localId) { entity.serverId = interactionServerId try await db.interactionDao.upsert(entity) // Tente l'upload audio immédiatement ; enfile un retry en cas d'échec. if let audioPath = entity.audioPath, let bytes = Self.readFile(audioPath) { let uploadResult = await api.uploadInteractionAudio( contactId: contactServerId, interactionId: interactionServerId, bytes: bytes, filename: "note.wav" ) if !uploadResult.isOk { try await enqueueInteractionAudioRetry( interactionServerId: interactionServerId ) } } } } } return Self.outcome(result) case "delete": guard let interactionId = op.serverId else { return PushOutcome(ok: false) } return Self.outcome(await api.deleteInteraction(interactionId: interactionId)) default: return PushOutcome(ok: true) } } /// Upload de l'audio d'une interaction (op de retry `interaction_audio`). private func pushInteractionAudioOp(_ op: SyncOpEntity) async throws -> PushOutcome { guard let interactionServerId = op.serverId, let interaction = try await db.interactionDao.getByServerId(interactionServerId), let audioPath = interaction.audioPath, let bytes = Self.readFile(audioPath) else { return PushOutcome(ok: true) } let result = await api.uploadInteractionAudio( contactId: interaction.contactServerId, interactionId: interactionServerId, bytes: bytes, filename: "note.wav" ) return Self.outcome(result) } private func enqueueInteractionAudioRetry(interactionServerId: String) async throws { let already = try await db.syncOpDao.listAll().contains { $0.entityType == "interaction_audio" && $0.serverId == interactionServerId } if already { return } var entity = SyncOpEntity(entityType: "interaction_audio", op: "upload") entity.payloadJson = "{}" entity.serverId = interactionServerId entity.createdAt = AgendaSyncCoordinator.currentMillis() try await db.syncOpDao.insert(entity) } /// Création ou suppression d'une note de projet. /// Create : `op.serverId` = projetServerId (parent), `op.localId` = note locale. /// Delete : `op.serverId` = noteServerId, `op.payloadJson` = `{"projetServerId":"…"}`. private func pushNoteProjetOp(_ op: SyncOpEntity) async throws -> PushOutcome { switch op.op { case "create": guard let projetServerId = op.serverId else { return PushOutcome(ok: false) } let result = await api.createNoteProjet(projetId: projetServerId, jsonBody: op.payloadJson) if case .ok(let body) = result, let localId = op.localId { if let noteServerId = Self.parseCreatedId(body) { // NoteProjetDao n'a pas getByLocalId ; on passe par listAll. if let note = try await db.noteProjetDao.listAll().first(where: { $0.localId == localId }) { var updated = note updated.serverId = noteServerId try await db.noteProjetDao.upsert(updated) // Upload audio si présent. if let audioPath = updated.audioPath, let bytes = Self.readFile(audioPath) { let uploadResult = await api.uploadNoteProjetAudio( projetId: projetServerId, noteId: noteServerId, bytes: bytes, filename: "note.wav" ) if !uploadResult.isOk { try await enqueueNoteProjetAudioRetry( projetServerId: projetServerId, noteServerId: noteServerId ) } } } } } return Self.outcome(result) case "delete": guard let noteId = op.serverId else { return PushOutcome(ok: false) } guard let ref = try? SyncJson.decode(NoteDeleteRef.self, from: op.payloadJson), let projetServerId = ref.projetServerId else { return PushOutcome(ok: false) } return Self.outcome(await api.deleteNoteProjet(projetId: projetServerId, noteId: noteId)) default: return PushOutcome(ok: true) } } private struct NoteDeleteRef: Codable { var projetServerId: String? } /// Upload de l'audio d'une note de projet (op de retry `note_projet_audio`). private func pushNoteProjetAudioOp(_ op: SyncOpEntity) async throws -> PushOutcome { guard let noteServerId = op.serverId, let note = try await db.noteProjetDao.getByServerId(noteServerId), let audioPath = note.audioPath, let bytes = Self.readFile(audioPath) else { return PushOutcome(ok: true) } let result = await api.uploadNoteProjetAudio( projetId: note.projetServerId, noteId: noteServerId, bytes: bytes, filename: "note.wav" ) return Self.outcome(result) } private func enqueueNoteProjetAudioRetry(projetServerId: String, noteServerId: String) async throws { let already = try await db.syncOpDao.listAll().contains { $0.entityType == "note_projet_audio" && $0.serverId == noteServerId } if already { return } var entity = SyncOpEntity(entityType: "note_projet_audio", op: "upload") entity.payloadJson = "{}" entity.serverId = noteServerId entity.createdAt = AgendaSyncCoordinator.currentMillis() try await db.syncOpDao.insert(entity) } private func pushRdvOp(_ op: SyncOpEntity) async throws -> PushOutcome { switch op.op { case "create": let result = await api.createRdv(jsonBody: op.payloadJson) if case .ok(let body) = result { if let newServerId = Self.parseCreatedId(body), let localId = op.localId { if var entity = try await db.rdvDao.getByLocalId(localId) { entity.serverId = newServerId entity.dirtyLocal = false try await db.rdvDao.upsert(entity) } } } return Self.outcome(result) case "update": guard let serverId = op.serverId else { return PushOutcome(ok: false) } let result = await api.updateRdv(id: serverId, jsonBody: op.payloadJson) if result.isOk, let localId = op.localId, var entity = try await db.rdvDao.getByLocalId(localId) { entity.dirtyLocal = false try await db.rdvDao.upsert(entity) } return Self.outcome(result) case "delete": guard let serverId = op.serverId else { return PushOutcome(ok: false) } return Self.outcome(await api.deleteRdv(id: serverId)) default: return PushOutcome(ok: true) } } /// Resolves `(cibleType, cibleId)` via `op.localId` (create) or `op.serverId` (update/delete), /// falling back to `op.payloadJson` (a `delete` detected after the calendar-side removal has /// already deleted the local entity; see `AgendaSyncCoordinator.cibleRefPayload`). private func resolveReservationCible(_ op: SyncOpEntity) async throws -> (String, String)? { if let localId = op.localId, let reservation = try await db.reservationDao.getByLocalId(localId) { return (reservation.cibleType, reservation.cibleId) } if let serverId = op.serverId, let reservation = try await db.reservationDao.getByServerId(serverId) { return (reservation.cibleType, reservation.cibleId) } return Self.parseCibleRef(op.payloadJson) } /// On 409 (server-side overlap): keep the op queued (see `pushOps`), flag `conflictPending` /// on the local entity and surface an `AgendaConflict` for the « Annuler ma réservation ? » dialog. private func pushReservationOp(_ op: SyncOpEntity) async throws -> PushOutcome { guard let (cibleType, cibleId) = try await resolveReservationCible(op) else { return PushOutcome(ok: false) } switch op.op { case "create": let result = await api.createReservation(genre: cibleType, resourceId: cibleId, jsonBody: op.payloadJson) return try await handleReservationPushResult(result, op: op) case "update": guard let serverId = op.serverId else { return PushOutcome(ok: false) } let result = await api.updateReservation( genre: cibleType, resourceId: cibleId, reservationId: serverId, jsonBody: op.payloadJson ) return try await handleReservationPushResult(result, op: op) case "delete": guard let serverId = op.serverId else { return PushOutcome(ok: false) } return Self.outcome(await api.deleteReservation(genre: cibleType, resourceId: cibleId, reservationId: serverId)) default: return PushOutcome(ok: true) } } private func handleReservationPushResult( _ result: AilianceApiClient.ApiResult<String>, op: SyncOpEntity ) async throws -> PushOutcome { switch result { case .ok(let body): if let localId = op.localId, var local = try await db.reservationDao.getByLocalId(localId) { if op.op == "create" { local.serverId = Self.parseCreatedId(body) ?? local.serverId } local.dirtyLocal = false try await db.reservationDao.upsert(local) } return PushOutcome(ok: true) case .err(let code, let message): if code == 409, let localId = op.localId { if var local = try await db.reservationDao.getByLocalId(localId) { local.conflictPending = true try await db.reservationDao.upsert(local) } return PushOutcome(ok: false, conflict: AgendaConflict(localReservationId: localId, message: message)) } return Self.outcome(result) } } // ---- pull ---- /// Applies the pull and returns the number of received items as a user would count /// them: workflows (referential) and tombstones (deletions) are excluded. private func applyPull(_ pull: SyncPullResponse) async throws -> Int { for dto in pull.contacts { try await applyContact(dto) } for dto in pull.entreprises { try await applyEntreprise(dto) } for dto in pull.projets { try await applyProjet(dto) } for dto in pull.taches { try await applyTache(dto) } for dto in pull.interactions { try await applyInteraction(dto) } for dto in pull.notes { try await applyNoteProjet(dto) } for dto in pull.workflows { try await applyWorkflow(dto) } for dto in pull.rdv { try await agenda.applyRdvPull(dto) } for dto in pull.reservations { try await agenda.applyReservationPull(dto) } for dto in pull.indisponibilites { try await agenda.applyIndisponibilitePull(dto) } for dto in pull.tombstones { try await applyTombstone(dto) } return pull.contacts.count + pull.entreprises.count + pull.projets.count + pull.taches.count + pull.interactions.count + pull.notes.count + pull.rdv.count + pull.reservations.count + pull.indisponibilites.count } private func applyContact(_ dto: ContactDto) async throws { let remoteTs = parseIsoToEpochMs(dto.misAJourLe ?? dto.creeLe) let local = try await db.crmContactDao.getByServerId(dto.id) guard let local else { var entity = CrmContactEntity() entity.serverId = dto.id entity.fullName = "\(dto.prenom) \(dto.nom)".trimmingCharacters(in: .whitespaces) entity.firstName = dto.prenom entity.lastName = dto.nom entity.jobTitle = dto.fonction entity.statut = dto.statut entity.etape = dto.etape entity.tags = dto.tags entity.entrepriseServerId = dto.entrepriseId entity.createdAt = parseIsoToEpochMs(dto.creeLe) entity.updatedAt = remoteTs let newId = try await db.crmContactDao.insert(entity) try await pullContactImages(localId: newId, dto: dto, remoteTs: remoteTs, localUpdatedAt: 0) return } var remoteEntity = local remoteEntity.fullName = "\(dto.prenom) \(dto.nom)".trimmingCharacters(in: .whitespaces) remoteEntity.firstName = dto.prenom remoteEntity.lastName = dto.nom remoteEntity.jobTitle = dto.fonction remoteEntity.statut = dto.statut remoteEntity.etape = dto.etape remoteEntity.tags = dto.tags remoteEntity.entrepriseServerId = dto.entrepriseId remoteEntity.updatedAt = remoteTs let merged = LwwMerger.pickLww(local: local, localTs: local.updatedAt, remote: remoteEntity, remoteTs: remoteTs) if merged != local { try await db.crmContactDao.update(merged) } try await pullContactImages(localId: local.id, dto: dto, remoteTs: remoteTs, localUpdatedAt: local.updatedAt) } /// Downloads card/photo when present server-side and (no local file, or server newer), /// unless a contact/media push is still queued. private func pullContactImages( localId: Int64, dto: ContactDto, remoteTs: Int64, localUpdatedAt: Int64 ) async throws { guard let store = imageStore else { return } if try await hasPendingContactPush(localId: localId, serverId: dto.id) { return } guard var contact = try await db.crmContactDao.getById(localId) else { return } let serverNewer = remoteTs > localUpdatedAt if let carteVisite = dto.carteVisite, !carteVisite.isEmpty { let path = contact.cardImagePath let missing = path == nil || path!.isEmpty || !Self.isFile(path!) if missing || serverNewer { if case .ok(let bytes) = await api.downloadContactCarte(id: dto.id) { let saved = try store.saveBytes(contactId: localId, name: "card.jpg", bytes: bytes) contact.cardImagePath = saved try await db.crmContactDao.update(contact) } } } if let photo = dto.photo, !photo.isEmpty { let path = contact.profileImagePath let missing = path == nil || path!.isEmpty || !Self.isFile(path!) if missing || serverNewer { if case .ok(let bytes) = await api.downloadContactPhoto(id: dto.id) { let saved = try store.saveBytes(contactId: localId, name: "profile.jpg", bytes: bytes) contact.profileImagePath = saved try await db.crmContactDao.update(contact) } } } } private func hasPendingContactPush(localId: Int64, serverId: String) async throws -> Bool { try await db.syncOpDao.listAll().contains { op in switch op.entityType { case "contact", Self.entityContactMedia: return (op.localId == localId || op.serverId == serverId) && op.op != "delete" default: return false } } } private func applyEntreprise(_ dto: EntrepriseDto) async throws { let remoteTs = parseIsoToEpochMs(dto.misAJourLe ?? dto.creeLe) let local = try await db.entrepriseDao.getByServerId(dto.id) guard let local else { var entity = EntrepriseEntity() entity.serverId = dto.id entity.nom = dto.nom entity.secteur = dto.secteur entity.creePar = dto.creePar entity.createdAt = parseIsoToEpochMs(dto.creeLe) entity.updatedAt = remoteTs try await db.entrepriseDao.upsert(entity) return } var remoteEntity = local remoteEntity.nom = dto.nom remoteEntity.secteur = dto.secteur remoteEntity.creePar = dto.creePar remoteEntity.updatedAt = remoteTs let merged = LwwMerger.pickLww(local: local, localTs: local.updatedAt, remote: remoteEntity, remoteTs: remoteTs) if merged != local { try await db.entrepriseDao.upsert(merged) } } private func applyProjet(_ dto: ProjetDto) async throws { let remoteTs = parseIsoToEpochMs(dto.misAJourLe ?? dto.creeLe) let local = try await db.projetDao.getByServerId(dto.id) let membresJson = SyncJson.encodeToString(dto.membres) guard let local else { var entity = ProjetEntity() entity.serverId = dto.id entity.nom = dto.nom entity.description = dto.description entity.workflowServerId = dto.workflowId entity.membresJson = membresJson entity.creePar = dto.creePar entity.planification = dto.planification entity.echeance = dto.echeance entity.joursOuvres = dto.joursOuvres entity.createdAt = parseIsoToEpochMs(dto.creeLe) entity.updatedAt = remoteTs try await db.projetDao.upsert(entity) return } var remoteEntity = local remoteEntity.nom = dto.nom remoteEntity.description = dto.description remoteEntity.workflowServerId = dto.workflowId remoteEntity.membresJson = membresJson remoteEntity.creePar = dto.creePar remoteEntity.planification = dto.planification remoteEntity.echeance = dto.echeance remoteEntity.joursOuvres = dto.joursOuvres remoteEntity.updatedAt = remoteTs let merged = LwwMerger.pickLww(local: local, localTs: local.updatedAt, remote: remoteEntity, remoteTs: remoteTs) if merged != local { try await db.projetDao.upsert(merged) } } private func applyTache(_ dto: TacheSyncDto) async throws { let remoteTs = parseIsoToEpochMs(dto.misAJourLe ?? dto.creeLe) let local = try await db.tacheDao.getByServerId(dto.id) let projetLocalId = try await db.projetDao.getByServerId(dto.projetId)?.localId let dependDeJson = SyncJson.encodeToString(dto.dependDe) let sousTachesJson = SyncJson.encodeToString(dto.sousTaches) guard let local else { var entity = TacheEntity() entity.serverId = dto.id entity.projetLocalId = projetLocalId entity.projetServerId = dto.projetId entity.titre = dto.titre entity.description = dto.description entity.colonneId = dto.colonneId entity.ordre = dto.ordre entity.assigneA = dto.assigneA entity.debut = dto.debut entity.dureeJours = dto.dureeJours entity.etiquette = dto.etiquette entity.parentId = dto.parentId entity.dependDeJson = dependDeJson entity.sousTachesJson = sousTachesJson entity.auteur = dto.auteur entity.createdAt = parseIsoToEpochMs(dto.creeLe) entity.updatedAt = remoteTs try await db.tacheDao.upsert(entity) return } var remoteEntity = local remoteEntity.projetLocalId = projetLocalId ?? local.projetLocalId remoteEntity.projetServerId = dto.projetId remoteEntity.titre = dto.titre remoteEntity.description = dto.description remoteEntity.colonneId = dto.colonneId remoteEntity.ordre = dto.ordre remoteEntity.assigneA = dto.assigneA remoteEntity.debut = dto.debut remoteEntity.dureeJours = dto.dureeJours remoteEntity.etiquette = dto.etiquette remoteEntity.parentId = dto.parentId remoteEntity.dependDeJson = dependDeJson remoteEntity.sousTachesJson = sousTachesJson remoteEntity.auteur = dto.auteur remoteEntity.updatedAt = remoteTs let merged = LwwMerger.pickLww(local: local, localTs: local.updatedAt, remote: remoteEntity, remoteTs: remoteTs) if merged != local { try await db.tacheDao.upsert(merged) } } private func applyInteraction(_ dto: InteractionDto) async throws { let remoteTs = parseIsoToEpochMs(dto.misAJourLe ?? dto.creeLe) let existing = try await db.interactionDao.getByServerId(dto.id) guard let existing else { // Insert : nouvelle interaction inconnue localement. var entity = InteractionEntity() entity.serverId = dto.id entity.contactServerId = dto.contactId entity.type = dto.typeInteraction entity.sujet = dto.sujet entity.description = dto.description entity.creePar = dto.creePar entity.createdAt = parseIsoToEpochMs(dto.creeLe) entity.updatedAt = dto.misAJourLe.map { parseIsoToEpochMs($0) } entity.transcriptionStatut = dto.transcription entity.transcriptionErreur = dto.transcriptionErreur // audioPath : nil à l'insert (fichier local pas encore présent). try await db.interactionDao.upsert(entity) return } // Mise à jour LWW : seulement si le distant est plus récent. let localTs = existing.updatedAt ?? existing.createdAt guard remoteTs > localTs else { return } var updated = existing updated.type = dto.typeInteraction updated.sujet = dto.sujet updated.description = dto.description updated.creePar = dto.creePar updated.updatedAt = remoteTs updated.transcriptionStatut = dto.transcription updated.transcriptionErreur = dto.transcriptionErreur // audioPath local préservé : ne pas écraser avec le nom serveur (pieceJointe). try await db.interactionDao.upsert(updated) } private func applyNoteProjet(_ dto: NoteProjetDto) async throws { let remoteTs = parseIsoToEpochMs(dto.majLe ?? dto.creeLe) let existing = try await db.noteProjetDao.getByServerId(dto.id) guard let existing else { var entity = NoteProjetEntity() entity.serverId = dto.id entity.projetServerId = dto.projetId entity.titre = dto.titre entity.texte = dto.contenu entity.auteur = dto.auteur entity.transcriptionStatut = dto.transcription entity.transcriptionErreur = dto.transcriptionErreur entity.createdAt = parseIsoToEpochMs(dto.creeLe) entity.updatedAt = remoteTs try await db.noteProjetDao.upsert(entity) return } // LWW : mise à jour seulement si le distant est plus récent. let localTs = existing.updatedAt ?? existing.createdAt guard remoteTs > localTs else { return } var updated = existing updated.titre = dto.titre updated.texte = dto.contenu updated.auteur = dto.auteur updated.transcriptionStatut = dto.transcription updated.transcriptionErreur = dto.transcriptionErreur // audioPath local préservé : ne pas écraser avec le nom de fichier serveur. updated.updatedAt = remoteTs try await db.noteProjetDao.upsert(updated) } private func applyWorkflow(_ dto: WorkflowDto) async throws { var entity = WorkflowEntity(serverId: dto.id) entity.nom = dto.nom entity.description = dto.description entity.colonnesJson = SyncJson.encodeToString(dto.colonnes) entity.creePar = dto.creePar entity.createdAt = parseIsoToEpochMs(dto.creeLe) try await db.workflowDao.upsert(entity) } private func applyTombstone(_ dto: TombstoneDto) async throws { switch dto.entityType { case "contact": try await db.crmContactDao.deleteByServerId(dto.id) case "entreprise": try await db.entrepriseDao.deleteByServerId(dto.id) case "projet": for serverId in try await db.tacheDao.listByProjetServerId(dto.id).compactMap({ $0.serverId }) { try await db.tacheDao.deleteByServerId(serverId) } try await db.projetDao.deleteByServerId(dto.id) case "tache": try await db.tacheDao.deleteByServerId(dto.id) case "interaction": try await db.interactionDao.deleteByServerId(dto.id) case "note": try await db.noteProjetDao.deleteByServerId(dto.id) default: try await agenda.applyAgendaTombstone(dto) } } // ---- helpers ---- private func nonBlank(_ value: String?) -> String? { guard let value, !value.isEmpty else { return nil } return value } private static func isFile(_ path: String) -> Bool { var isDirectory: ObjCBool = false return FileManager.default.fileExists(atPath: path, isDirectory: &isDirectory) && !isDirectory.boolValue } private static func readFile(_ path: String) -> Data? { guard isFile(path) else { return nil } return try? Data(contentsOf: URL(fileURLWithPath: path)) } private struct CreatedIdDto: Codable { var id: String? = nil init(from decoder: Decoder) throws { let c = try decoder.container(keyedBy: CodingKeys.self) id = try c.decodeIfPresent(String.self, forKey: .id) } } static func parseCreatedId(_ body: String) -> String? { (try? SyncJson.decode(CreatedIdDto.self, from: body))?.id } static func parseCibleRef(_ payloadJson: String) -> (String, String)? { guard let ref = try? SyncJson.decode(CibleRessourceDto.self, from: payloadJson) else { return nil } return (ref.type, ref.id) } }
GitRust