import Foundation import SwiftData import UniformTypeIdentifiers import Network import CryptoKit // Remote contract: Docs/RemoteAPI.openapi.yaml struct RemoteAsset: Decodable { let id: String let clientId: String let kind: String let fileName: String let mimeType: String let sizeBytes: Int64 let sha256: String let createdAt: Date? } struct AssetUploadResult { let asset: RemoteAsset let sessionRevision: Int64? } struct StorageQuota: Decodable { let totalBytes: Int64 let usedBytes: Int64 let remainingBytes: Int64 } struct RemoteEvent: Decodable { let id: String let clientId: String? let relativeTimeMs: Int64 let eventType: String let textContent: String? let voiceStartOffsetMs: Int64? let voiceEndOffsetMs: Int64? } struct RemoteSession: Decodable { let id: String let clientId: String? let title: String let startTime: Date let endTime: Date? let durationMs: Int64 let revision: Int64 let deletedAt: Date? let events: [RemoteEvent]? let assets: [RemoteAsset]? } private struct SessionUpload: Encodable { let clientId: String let title: String let startTime: Date let endTime: Date? let durationMs: Int64 let events: [EventUpload] let baseRevision: Int64? let deletedEventClientIds: [String] } private struct EventUpload: Encodable { let clientId: String let relativeTimeMs: Int64 let eventType: String let textContent: String? let voiceStartOffsetMs: Int64? let voiceEndOffsetMs: Int64? } private struct SyncCheckpoint: Decodable { let syncedAt: Date } private struct AssetUploadPayload: Decodable { let asset: RemoteAsset let sessionRevision: Int64? private enum CodingKeys: String, CodingKey { case asset case sessionRevision } init(from decoder: Decoder) throws { if let container = try? decoder.container(keyedBy: CodingKeys.self), container.contains(.asset) { asset = try container.decode(RemoteAsset.self, forKey: .asset) sessionRevision = try container.decodeIfPresent(Int64.self, forKey: .sessionRevision) } else { asset = try RemoteAsset(from: decoder) sessionRevision = nil } } } private struct ChunkUploadInitBody: Encodable { let clientId: String let kind: String let fileName: String let mimeType: String let fileSize: Int64 let chunkSize: Int } private struct ChunkUploadInitResponse: Decodable { let uploadId: String let totalChunks: Int let chunkSize: Int } private struct ChunkUploadProgressResponse: Decodable { let uploadedChunks: Int let totalChunks: Int } private struct ChunkUploadCompleteBody: Encodable { let uploadId: String let totalChunks: Int } final class RemoteNetworkService: NetworkServiceProtocol { private static let chunkedUploadThreshold: Int64 = 16 * 1024 * 1024 private static let preferredChunkSize = 8 * 1024 * 1024 private let client: APIClient init(client: APIClient = .shared) { self.client = client } func fetchSessions() async throws -> [RemoteSession] { try await client.request("sessions?includeDeleted=true", authenticated: true) } func syncSession( _ session: CelestiaSession, deletedEventClientIDs: [String] ) async throws -> RemoteSession { let payload = SessionUpload( clientId: session.id.uuidString, title: session.title, startTime: session.startTime, endTime: session.endTime, durationMs: session.durationMs, events: session.events.map { EventUpload( clientId: $0.id.uuidString, relativeTimeMs: $0.relativeTimeMs, eventType: $0.eventType, textContent: $0.textContent, voiceStartOffsetMs: $0.voiceStartOffsetMs, voiceEndOffsetMs: $0.voiceEndOffsetMs ) }, baseRevision: session.serverRevision > 0 ? session.serverRevision : nil, deletedEventClientIds: deletedEventClientIDs ) let encoder = JSONEncoder() encoder.dateEncodingStrategy = .iso8601 return try await client.request("sessions", method: .post, body: try encoder.encode(payload), authenticated: true) } func fetchStorageQuota() async throws -> StorageQuota { try await client.request("storage/quota", authenticated: true) } func uploadAsset( sessionID: String, clientID: String, kind: String, fileURL: URL, progress: @escaping @Sendable (Double) -> Void ) async throws -> AssetUploadResult { let mimeType = UTType(filenameExtension: fileURL.pathExtension)?.preferredMIMEType ?? "application/octet-stream" let size = try fileSize(of: fileURL) if size >= Self.chunkedUploadThreshold { return try await uploadAssetInChunks( sessionID: sessionID, clientID: clientID, kind: kind, fileURL: fileURL, fileSize: size, mimeType: mimeType, progress: progress ) } let payload: AssetUploadPayload = try await client.upload( "sessions/\(sessionID)/assets", fileURL: fileURL, fileName: fileURL.lastPathComponent, mimeType: mimeType, fields: ["clientId": clientID, "kind": kind], progress: progress ) return AssetUploadResult(asset: payload.asset, sessionRevision: payload.sessionRevision) } func downloadAsset(sessionID: String, asset: RemoteAsset, destinationURL: URL) async throws { try await client.download("sessions/\(sessionID)/assets/\(asset.id)", to: destinationURL) } func recordSyncCheckpoint() async throws { let _: SyncCheckpoint = try await client.request("sync/trigger", method: .post, authenticated: true) } private func uploadAssetInChunks( sessionID: String, clientID: String, kind: String, fileURL: URL, fileSize: Int64, mimeType: String, progress: @escaping @Sendable (Double) -> Void ) async throws -> AssetUploadResult { let body = ChunkUploadInitBody( clientId: clientID, kind: kind, fileName: fileURL.lastPathComponent, mimeType: mimeType, fileSize: fileSize, chunkSize: Self.preferredChunkSize ) let encoder = JSONEncoder() let upload: ChunkUploadInitResponse = try await client.request( "sessions/\(sessionID)/assets/init", method: .post, body: try encoder.encode(body), authenticated: true ) for chunkIndex in 0.. (url: URL, byteCount: Int) { let task = Task.detached(priority: .utility) { try Task.checkCancellation() let input = try FileHandle(forReadingFrom: fileURL) defer { try? input.close() } try input.seek(toOffset: offset) guard let data = try input.read(upToCount: chunkSize), !data.isEmpty else { throw APIError.transport("读取上传分片失败") } try Task.checkCancellation() let temporaryURL = FileManager.default.temporaryDirectory .appendingPathComponent("celestia-\(uploadID)-\(chunkIndex).chunk") try data.write(to: temporaryURL, options: .atomic) return (temporaryURL, data.count) } return try await withTaskCancellationHandler { try await task.value } onCancel: { task.cancel() } } private func fileSize(of url: URL) throws -> Int64 { let values = try url.resourceValues(forKeys: [.fileSizeKey]) guard let size = values.fileSize else { throw APIError.transport("无法读取待上传文件大小") } return Int64(size) } } struct RemoteBoundDevice: Decodable { let id: String let name: String let peripheralUUID: String } private struct BindDeviceBody: Encodable { let name: String let peripheralUUID: String let hardwareMAC: String? let firmwareVersion: String? let batteryLevel: Int? let freeStorageMB: Int? let totalStorageMB: Int? } private struct UpdateDeviceBody: Encodable { let name: String? let batteryLevel: Int? let isConnected: Bool? let firmwareVersion: String? let freeStorageMB: Int? let totalStorageMB: Int? } final class DeviceCloudService { static let shared = DeviceCloudService() private let client: APIClient init(client: APIClient = .shared) { self.client = client } func register(_ device: BoundDevice) async throws -> RemoteBoundDevice { let body = BindDeviceBody( name: device.name, peripheralUUID: device.peripheralUUID, hardwareMAC: device.hardwareMAC, firmwareVersion: device.firmwareVersion, batteryLevel: device.batteryLevel, freeStorageMB: device.freeStorageMB, totalStorageMB: device.totalStorageMB ) return try await client.request("devices", method: .post, body: try JSONEncoder().encode(body), authenticated: true) } func update(_ device: BoundDevice) async throws -> RemoteBoundDevice { guard let cloudID = device.cloudID else { return try await register(device) } let body = UpdateDeviceBody( name: device.name, batteryLevel: device.batteryLevel, isConnected: device.isConnected, firmwareVersion: device.firmwareVersion, freeStorageMB: device.freeStorageMB, totalStorageMB: device.totalStorageMB ) return try await client.request("devices/\(cloudID)", method: .put, body: try JSONEncoder().encode(body), authenticated: true) } func remove(cloudID: String) async throws { try await client.requestVoid("devices/\(cloudID)", method: .delete, authenticated: true) } } final class NetworkStatusMonitor: ObservableObject, @unchecked Sendable { static let shared = NetworkStatusMonitor() @Published private(set) var isConnected = true @Published private(set) var isWiFi = false @Published private(set) var isCellular = false private let monitor = NWPathMonitor() private let queue = DispatchQueue(label: "com.celestia.trace.network-path") private init() { monitor.pathUpdateHandler = { [weak self] path in let connected = path.status == .satisfied let usesWiFi = path.usesInterfaceType(.wifi) let usesCellular = path.usesInterfaceType(.cellular) DispatchQueue.main.async { self?.isConnected = connected self?.isWiFi = usesWiFi self?.isCellular = usesCellular } } monitor.start(queue: queue) } var connectionName: String { if isWiFi { return "Wi-Fi" } if isCellular { return "蜂窝网络" } return isConnected ? "当前网络" : "网络" } } @MainActor final class SyncManager: ObservableObject { static let shared = SyncManager() @Published private(set) var isSyncing = false @Published private(set) var lastSyncDate: Date? @Published private(set) var lastErrorMessage: String? @Published private(set) var completedCount = 0 @Published private(set) var totalCount = 0 @Published private(set) var activeSessionID: UUID? @Published private(set) var syncProgress: Double = 0 private let service: NetworkServiceProtocol private var activeSyncTask: Task? private var automaticSyncTasks: [UUID: Task] = [:] private var automaticSyncGenerations: [UUID: UUID] = [:] init(service: NetworkServiceProtocol = RemoteNetworkService()) { self.service = service lastSyncDate = UserDefaults.standard.object(forKey: "com.celestia.trace.last_server_sync") as? Date } func estimatedUploadBytes(for session: CelestiaSession) -> Int64 { var urls: [URL] = [] if let url = AudioPathHelper.resolveURL(for: session.localAudioPath) { urls.append(url) } urls.append(contentsOf: session.events.compactMap { event in guard event.eventType == "PHOTO" else { return nil } return AudioPathHelper.resolveURL(for: event.localFilePath) }) return urls.reduce(into: 0) { total, url in let values = try? url.resourceValues(forKeys: [.fileSizeKey]) total += Int64(values?.fileSize ?? 0) } } func scheduleAutomaticSync( for session: CelestiaSession, modelContext: ModelContext, delay: Duration = .seconds(1.5) ) { guard session.endTime != nil, session.syncState != .conflict, AuthManager.shared.currentUser?.id != nil else { return } let sessionID = session.id let generation = UUID() automaticSyncGenerations[sessionID] = generation automaticSyncTasks[sessionID]?.cancel() if isSyncing, activeSessionID == sessionID { activeSyncTask?.cancel() } automaticSyncTasks[sessionID] = Task { @MainActor [weak self] in do { try await Task.sleep(for: delay) while let self, (self.isSyncing || !NetworkStatusMonitor.shared.isConnected) { try Task.checkCancellation() try await Task.sleep(for: .seconds(0.75)) } guard let self, !Task.isCancelled, self.automaticSyncGenerations[sessionID] == generation, session.needsCloudSync, let userID = AuthManager.shared.currentUser?.id else { return } _ = await self.sync( sessions: [session], modelContext: modelContext, userID: userID ) if self.automaticSyncGenerations[sessionID] == generation { self.automaticSyncTasks[sessionID] = nil self.automaticSyncGenerations[sessionID] = nil } } catch { guard let self, self.automaticSyncGenerations[sessionID] == generation else { return } self.automaticSyncTasks[sessionID] = nil self.automaticSyncGenerations[sessionID] = nil } } } func sync( sessions: [CelestiaSession], modelContext: ModelContext, userID: String ) async -> Bool { guard !isSyncing else { return false } isSyncing = true let task = Task { @MainActor [self] in await performSync( sessions: sessions, modelContext: modelContext, userID: userID ) } activeSyncTask = task let result = await withTaskCancellationHandler { await task.value } onCancel: { task.cancel() } activeSyncTask = nil isSyncing = false activeSessionID = nil return result } func pauseSync(sessionID: UUID) { guard isSyncing, activeSessionID == sessionID else { return } activeSyncTask?.cancel() } private func performSync( sessions: [CelestiaSession], modelContext: ModelContext, userID: String ) async -> Bool { lastErrorMessage = nil completedCount = 0 activeSessionID = nil syncProgress = 0 let queuedSessions = sessions.filter { ($0.ownerUserID == nil || $0.ownerUserID == userID) && $0.needsCloudSync && $0.syncState != .conflict } totalCount = queuedSessions.count // Give the UI immediate feedback while the remote index is being fetched. // The persisted state is changed only when this session actually begins syncing. activeSessionID = queuedSessions.first?.id do { let remoteSessions = try await service.fetchSessions() try Task.checkCancellation() let mergedSessions = merge(remoteSessions, into: sessions, modelContext: modelContext, userID: userID) for (remote, local) in mergedSessions { try await downloadMissingAssets(for: local, remote: remote) try Task.checkCancellation() } try modelContext.save() let eligible = sessions.filter { ($0.ownerUserID == nil || $0.ownerUserID == userID) && $0.needsCloudSync && $0.syncState != .conflict } var failures: [String] = [] let eligibleCount = max(eligible.count, 1) for (index, session) in eligible.enumerated() { activeSessionID = session.id let sessionStartProgress = Double(index) / Double(eligibleCount) let sessionProgressSpan = 1.0 / Double(eligibleCount) syncProgress = sessionStartProgress session.ownerUserID = userID session.syncState = .syncing session.lastSyncError = nil do { let previousRemote = remoteSessions.first { $0.id == session.cloudSessionId || $0.clientId?.caseInsensitiveCompare(session.id.uuidString) == .orderedSame } let localEventIDs = Set(session.events.map { $0.id.uuidString.lowercased() }) let deletedEventClientIDs = previousRemote?.events? .compactMap(\.clientId) .filter { !localEventIDs.contains($0.lowercased()) } ?? [] let remote = try await service.syncSession( session, deletedEventClientIDs: deletedEventClientIDs ) try Task.checkCancellation() syncProgress = sessionStartProgress + sessionProgressSpan * 0.12 session.cloudSessionId = remote.id session.serverRevision = remote.revision try await uploadLocalAssets(for: session, remote: remote) { [weak self] assetProgress in Task { @MainActor [weak self] in self?.syncProgress = sessionStartProgress + sessionProgressSpan * (0.12 + assetProgress * 0.83) } } try Task.checkCancellation() session.isSynced = true session.syncState = .synced session.lastSyncedAt = Date() completedCount += 1 syncProgress = sessionStartProgress + sessionProgressSpan } catch where Task.isCancelled { session.isSynced = false session.syncState = .pending session.lastSyncError = nil syncProgress = 0 try? modelContext.save() return false } catch APIError.server(let code, _) where code == 409 { session.isSynced = false session.syncState = .conflict session.lastSyncError = "本地版本 \(session.serverRevision) 与云端版本不一致" } catch { session.isSynced = false session.syncState = .failed session.lastSyncError = error.localizedDescription failures.append("\(session.title):\(error.localizedDescription)") } try modelContext.save() } let conflicts = sessions.filter { $0.syncState == .conflict } if !conflicts.isEmpty { failures.append(contentsOf: conflicts.map { "\($0.title):本地和云端版本不一致" }) } guard failures.isEmpty else { lastErrorMessage = failures.joined(separator: "\n") return false } try await service.recordSyncCheckpoint() let now = Date() lastSyncDate = now UserDefaults.standard.set(now, forKey: "com.celestia.trace.last_server_sync") return true } catch where Task.isCancelled { syncProgress = 0 return false } catch { lastErrorMessage = error.localizedDescription return false } } private func uploadLocalAssets( for session: CelestiaSession, remote: RemoteSession, progress: @escaping @Sendable (Double) -> Void ) async throws { guard let cloudID = session.cloudSessionId else { throw APIError.missingData } let existingIDs = Set((remote.assets ?? []).map(\.clientId)) var uploads: [(clientID: String, kind: String, url: URL, size: Int64)] = [] if let path = session.localAudioPath { guard let url = AudioPathHelper.resolveURL(for: path) else { throw APIError.transport("本地录音文件不存在:\(path)") } let sha256 = try await fileSHA256(of: url) let hasMatchingAudio = (remote.assets ?? []).contains { $0.kind.uppercased() == "AUDIO" && $0.sha256.caseInsensitiveCompare(sha256) == .orderedSame } if !hasMatchingAudio { let clientID = "\(session.id.uuidString)-audio-\(sha256.prefix(16))" uploads.append((clientID, "AUDIO", url, fileSize(of: url))) } } for event in session.events where event.eventType == "PHOTO" { guard let path = event.localFilePath else { continue } guard let url = AudioPathHelper.resolveURL(for: path) else { throw APIError.transport("照片文件不存在:\(path)") } let clientID = "\(event.id.uuidString)-photo" if !existingIDs.contains(clientID) { uploads.append((clientID, "PHOTO", url, fileSize(of: url))) } } let totalBytes = max(uploads.reduce(Int64(0)) { $0 + $1.size }, 1) let requiredBytes = uploads.reduce(Int64(0)) { $0 + $1.size } if requiredBytes > 0 { let quota = try await service.fetchStorageQuota() guard requiredBytes <= quota.remainingBytes else { throw APIError.storageQuotaExceeded( requiredBytes: requiredBytes, remainingBytes: quota.remainingBytes ) } } var completedBytes: Int64 = 0 progress(uploads.isEmpty ? 1 : 0) for upload in uploads { let bytesBeforeUpload = completedBytes let result = try await service.uploadAsset( sessionID: cloudID, clientID: upload.clientID, kind: upload.kind, fileURL: upload.url ) { fraction in let sentBytes = Double(bytesBeforeUpload) + Double(upload.size) * fraction progress(min(max(sentBytes / Double(totalBytes), 0), 1)) } if let revision = result.sessionRevision { session.serverRevision = revision } completedBytes += upload.size progress(Double(completedBytes) / Double(totalBytes)) } } private func fileSize(of url: URL) -> Int64 { let values = try? url.resourceValues(forKeys: [.fileSizeKey]) return Int64(values?.fileSize ?? 0) } private func fileSHA256(of url: URL) async throws -> String { let task = Task.detached(priority: .utility) { let input = try FileHandle(forReadingFrom: url) defer { try? input.close() } var hasher = SHA256() while let data = try input.read(upToCount: 1_024 * 1_024), !data.isEmpty { try Task.checkCancellation() hasher.update(data: data) } return hasher.finalize().map { String(format: "%02x", $0) }.joined() } return try await withTaskCancellationHandler { try await task.value } onCancel: { task.cancel() } } private func downloadMissingAssets(for session: CelestiaSession, remote: RemoteSession) async throws { let assets = remote.assets ?? [] if let audio = assets .filter({ $0.kind.uppercased() == "AUDIO" }) .max(by: { ($0.createdAt ?? .distantPast) < ($1.createdAt ?? .distantPast) }) { let localURL = AudioPathHelper.resolveURL(for: session.localAudioPath) let localHash: String? if let localURL { localHash = try? await fileSHA256(of: localURL) } else { localHash = nil } let needsDownload = localURL == nil || (session.isSynced && localHash?.caseInsensitiveCompare(audio.sha256) != .orderedSame) if needsDownload { let destination = downloadDestination(for: audio) try await service.downloadAsset(sessionID: remote.id, asset: audio, destinationURL: destination) session.localAudioPath = AudioPathHelper.relativePath(from: destination.path) } } let photoAssets = Dictionary(uniqueKeysWithValues: assets .filter { $0.kind.uppercased() == "PHOTO" } .map { ($0.clientId.lowercased(), $0) }) for event in session.events where event.eventType == "PHOTO" { guard event.localFilePath == nil || AudioPathHelper.resolveURL(for: event.localFilePath) == nil else { continue } let key = "\(event.id.uuidString)-photo".lowercased() guard let asset = photoAssets[key] else { continue } let destination = downloadDestination(for: asset) try await service.downloadAsset(sessionID: remote.id, asset: asset, destinationURL: destination) event.localFilePath = AudioPathHelper.relativePath(from: destination.path) } } private func downloadDestination(for asset: RemoteAsset) -> URL { let documents = FileManager.default.urls(for: .documentDirectory, in: .userDomainMask)[0] let fileExtension = (asset.fileName as NSString).pathExtension let suffix = fileExtension.isEmpty ? "" : ".\(fileExtension.lowercased())" return documents.appendingPathComponent("cloud_\(asset.id)\(suffix)") } private func merge( _ remoteSessions: [RemoteSession], into localSessions: [CelestiaSession], modelContext: ModelContext, userID: String ) -> [(RemoteSession, CelestiaSession)] { var byClientID = Dictionary(uniqueKeysWithValues: localSessions.map { ($0.id.uuidString.lowercased(), $0) }) var merged: [(RemoteSession, CelestiaSession)] = [] for remote in remoteSessions { guard let clientID = remote.clientId?.lowercased() else { continue } if let local = byClientID[clientID] { if remote.deletedAt != nil, local.isSynced { modelContext.delete(local) continue } if !local.isSynced, remote.revision > local.serverRevision { local.syncState = .conflict local.lastSyncError = "本地版本 \(local.serverRevision) 与云端版本 \(remote.revision) 不一致" continue } if local.isSynced, remote.revision > local.serverRevision { apply(remote, to: local, userID: userID) } if remote.deletedAt == nil { merged.append((remote, local)) } } else if remote.deletedAt == nil { let id = UUID(uuidString: clientID) ?? UUID() let session = CelestiaSession(id: id, title: remote.title, startTime: remote.startTime) apply(remote, to: session, userID: userID) modelContext.insert(session) byClientID[clientID] = session merged.append((remote, session)) } } try? modelContext.save() return merged } private func apply(_ remote: RemoteSession, to local: CelestiaSession, userID: String) { local.title = remote.title local.startTime = remote.startTime local.endTime = remote.endTime local.cloudSessionId = remote.id local.ownerUserID = userID local.serverRevision = remote.revision local.isSynced = true local.syncState = .synced local.lastSyncError = nil local.events.removeAll() for remoteEvent in remote.events ?? [] { let id = remoteEvent.clientId.flatMap(UUID.init(uuidString:)) ?? UUID() let event = CelestiaTimelineEvent(id: id, relativeTimeMs: remoteEvent.relativeTimeMs, eventType: remoteEvent.eventType) event.textContent = remoteEvent.textContent event.voiceStartOffsetMs = remoteEvent.voiceStartOffsetMs event.voiceEndOffsetMs = remoteEvent.voiceEndOffsetMs local.events.append(event) } } }