RemoteNetworkService.swift 28 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735
  1. import Foundation
  2. import SwiftData
  3. import UniformTypeIdentifiers
  4. import Network
  5. import CryptoKit
  6. // Remote contract: Docs/RemoteAPI.openapi.yaml
  7. struct RemoteAsset: Decodable {
  8. let id: String
  9. let clientId: String
  10. let kind: String
  11. let fileName: String
  12. let mimeType: String
  13. let sizeBytes: Int64
  14. let sha256: String
  15. let createdAt: Date?
  16. }
  17. struct AssetUploadResult {
  18. let asset: RemoteAsset
  19. let sessionRevision: Int64?
  20. }
  21. struct StorageQuota: Decodable {
  22. let totalBytes: Int64
  23. let usedBytes: Int64
  24. let remainingBytes: Int64
  25. }
  26. struct RemoteEvent: Decodable {
  27. let id: String
  28. let clientId: String?
  29. let relativeTimeMs: Int64
  30. let eventType: String
  31. let textContent: String?
  32. let voiceStartOffsetMs: Int64?
  33. let voiceEndOffsetMs: Int64?
  34. }
  35. struct RemoteSession: Decodable {
  36. let id: String
  37. let clientId: String?
  38. let title: String
  39. let startTime: Date
  40. let endTime: Date?
  41. let durationMs: Int64
  42. let revision: Int64
  43. let deletedAt: Date?
  44. let events: [RemoteEvent]?
  45. let assets: [RemoteAsset]?
  46. }
  47. private struct SessionUpload: Encodable {
  48. let clientId: String
  49. let title: String
  50. let startTime: Date
  51. let endTime: Date?
  52. let durationMs: Int64
  53. let events: [EventUpload]
  54. let baseRevision: Int64?
  55. let deletedEventClientIds: [String]
  56. }
  57. private struct EventUpload: Encodable {
  58. let clientId: String
  59. let relativeTimeMs: Int64
  60. let eventType: String
  61. let textContent: String?
  62. let voiceStartOffsetMs: Int64?
  63. let voiceEndOffsetMs: Int64?
  64. }
  65. private struct SyncCheckpoint: Decodable {
  66. let syncedAt: Date
  67. }
  68. private struct AssetUploadPayload: Decodable {
  69. let asset: RemoteAsset
  70. let sessionRevision: Int64?
  71. private enum CodingKeys: String, CodingKey {
  72. case asset
  73. case sessionRevision
  74. }
  75. init(from decoder: Decoder) throws {
  76. if let container = try? decoder.container(keyedBy: CodingKeys.self),
  77. container.contains(.asset) {
  78. asset = try container.decode(RemoteAsset.self, forKey: .asset)
  79. sessionRevision = try container.decodeIfPresent(Int64.self, forKey: .sessionRevision)
  80. } else {
  81. asset = try RemoteAsset(from: decoder)
  82. sessionRevision = nil
  83. }
  84. }
  85. }
  86. private struct ChunkUploadInitBody: Encodable {
  87. let clientId: String
  88. let kind: String
  89. let fileName: String
  90. let mimeType: String
  91. let fileSize: Int64
  92. let chunkSize: Int
  93. }
  94. private struct ChunkUploadInitResponse: Decodable {
  95. let uploadId: String
  96. let totalChunks: Int
  97. let chunkSize: Int
  98. }
  99. private struct ChunkUploadProgressResponse: Decodable {
  100. let uploadedChunks: Int
  101. let totalChunks: Int
  102. }
  103. private struct ChunkUploadCompleteBody: Encodable {
  104. let uploadId: String
  105. let totalChunks: Int
  106. }
  107. final class RemoteNetworkService: NetworkServiceProtocol {
  108. private static let chunkedUploadThreshold: Int64 = 16 * 1024 * 1024
  109. private static let preferredChunkSize = 8 * 1024 * 1024
  110. private let client: APIClient
  111. init(client: APIClient = .shared) {
  112. self.client = client
  113. }
  114. func fetchSessions() async throws -> [RemoteSession] {
  115. try await client.request("sessions?includeDeleted=true", authenticated: true)
  116. }
  117. func syncSession(
  118. _ session: CelestiaSession,
  119. deletedEventClientIDs: [String]
  120. ) async throws -> RemoteSession {
  121. let payload = SessionUpload(
  122. clientId: session.id.uuidString,
  123. title: session.title,
  124. startTime: session.startTime,
  125. endTime: session.endTime,
  126. durationMs: session.durationMs,
  127. events: session.events.map {
  128. EventUpload(
  129. clientId: $0.id.uuidString,
  130. relativeTimeMs: $0.relativeTimeMs,
  131. eventType: $0.eventType,
  132. textContent: $0.textContent,
  133. voiceStartOffsetMs: $0.voiceStartOffsetMs,
  134. voiceEndOffsetMs: $0.voiceEndOffsetMs
  135. )
  136. },
  137. baseRevision: session.serverRevision > 0 ? session.serverRevision : nil,
  138. deletedEventClientIds: deletedEventClientIDs
  139. )
  140. let encoder = JSONEncoder()
  141. encoder.dateEncodingStrategy = .iso8601
  142. return try await client.request("sessions", method: .post, body: try encoder.encode(payload), authenticated: true)
  143. }
  144. func fetchStorageQuota() async throws -> StorageQuota {
  145. try await client.request("storage/quota", authenticated: true)
  146. }
  147. func uploadAsset(
  148. sessionID: String,
  149. clientID: String,
  150. kind: String,
  151. fileURL: URL,
  152. progress: @escaping @Sendable (Double) -> Void
  153. ) async throws -> AssetUploadResult {
  154. let mimeType = UTType(filenameExtension: fileURL.pathExtension)?.preferredMIMEType ?? "application/octet-stream"
  155. let size = try fileSize(of: fileURL)
  156. if size >= Self.chunkedUploadThreshold {
  157. return try await uploadAssetInChunks(
  158. sessionID: sessionID,
  159. clientID: clientID,
  160. kind: kind,
  161. fileURL: fileURL,
  162. fileSize: size,
  163. mimeType: mimeType,
  164. progress: progress
  165. )
  166. }
  167. let payload: AssetUploadPayload = try await client.upload(
  168. "sessions/\(sessionID)/assets",
  169. fileURL: fileURL,
  170. fileName: fileURL.lastPathComponent,
  171. mimeType: mimeType,
  172. fields: ["clientId": clientID, "kind": kind],
  173. progress: progress
  174. )
  175. return AssetUploadResult(asset: payload.asset, sessionRevision: payload.sessionRevision)
  176. }
  177. func downloadAsset(sessionID: String, asset: RemoteAsset, destinationURL: URL) async throws {
  178. try await client.download("sessions/\(sessionID)/assets/\(asset.id)", to: destinationURL)
  179. }
  180. func recordSyncCheckpoint() async throws {
  181. let _: SyncCheckpoint = try await client.request("sync/trigger", method: .post, authenticated: true)
  182. }
  183. private func uploadAssetInChunks(
  184. sessionID: String,
  185. clientID: String,
  186. kind: String,
  187. fileURL: URL,
  188. fileSize: Int64,
  189. mimeType: String,
  190. progress: @escaping @Sendable (Double) -> Void
  191. ) async throws -> AssetUploadResult {
  192. let body = ChunkUploadInitBody(
  193. clientId: clientID,
  194. kind: kind,
  195. fileName: fileURL.lastPathComponent,
  196. mimeType: mimeType,
  197. fileSize: fileSize,
  198. chunkSize: Self.preferredChunkSize
  199. )
  200. let encoder = JSONEncoder()
  201. let upload: ChunkUploadInitResponse = try await client.request(
  202. "sessions/\(sessionID)/assets/init",
  203. method: .post,
  204. body: try encoder.encode(body),
  205. authenticated: true
  206. )
  207. let input = try FileHandle(forReadingFrom: fileURL)
  208. defer { try? input.close() }
  209. for chunkIndex in 0..<upload.totalChunks {
  210. let offset = UInt64(chunkIndex) * UInt64(upload.chunkSize)
  211. try input.seek(toOffset: offset)
  212. guard let data = try input.read(upToCount: upload.chunkSize), !data.isEmpty else {
  213. throw APIError.transport("读取上传分片失败")
  214. }
  215. let temporaryURL = FileManager.default.temporaryDirectory
  216. .appendingPathComponent("celestia-\(upload.uploadId)-\(chunkIndex).chunk")
  217. try data.write(to: temporaryURL, options: .atomic)
  218. defer { try? FileManager.default.removeItem(at: temporaryURL) }
  219. let chunkStart = Double(offset) / Double(fileSize)
  220. let chunkSpan = Double(data.count) / Double(fileSize)
  221. let _: ChunkUploadProgressResponse = try await client.upload(
  222. "sessions/\(sessionID)/assets/chunk",
  223. fileURL: temporaryURL,
  224. fileName: "chunk-\(chunkIndex)",
  225. mimeType: "application/octet-stream",
  226. fields: [
  227. "uploadId": upload.uploadId,
  228. "chunkIndex": String(chunkIndex)
  229. ]
  230. ) { fraction in
  231. progress(min(max(chunkStart + chunkSpan * fraction, 0), 1))
  232. }
  233. }
  234. let completeBody = ChunkUploadCompleteBody(
  235. uploadId: upload.uploadId,
  236. totalChunks: upload.totalChunks
  237. )
  238. let payload: AssetUploadPayload = try await client.request(
  239. "sessions/\(sessionID)/assets/complete",
  240. method: .post,
  241. body: try encoder.encode(completeBody),
  242. authenticated: true
  243. )
  244. progress(1)
  245. return AssetUploadResult(asset: payload.asset, sessionRevision: payload.sessionRevision)
  246. }
  247. private func fileSize(of url: URL) throws -> Int64 {
  248. let values = try url.resourceValues(forKeys: [.fileSizeKey])
  249. guard let size = values.fileSize else {
  250. throw APIError.transport("无法读取待上传文件大小")
  251. }
  252. return Int64(size)
  253. }
  254. }
  255. struct RemoteBoundDevice: Decodable {
  256. let id: String
  257. let name: String
  258. let peripheralUUID: String
  259. }
  260. private struct BindDeviceBody: Encodable {
  261. let name: String
  262. let peripheralUUID: String
  263. let hardwareMAC: String?
  264. let firmwareVersion: String?
  265. let batteryLevel: Int?
  266. let freeStorageMB: Int?
  267. let totalStorageMB: Int?
  268. }
  269. private struct UpdateDeviceBody: Encodable {
  270. let name: String?
  271. let batteryLevel: Int?
  272. let isConnected: Bool?
  273. let firmwareVersion: String?
  274. let freeStorageMB: Int?
  275. let totalStorageMB: Int?
  276. }
  277. final class DeviceCloudService {
  278. static let shared = DeviceCloudService()
  279. private let client: APIClient
  280. init(client: APIClient = .shared) {
  281. self.client = client
  282. }
  283. func register(_ device: BoundDevice) async throws -> RemoteBoundDevice {
  284. let body = BindDeviceBody(
  285. name: device.name,
  286. peripheralUUID: device.peripheralUUID,
  287. hardwareMAC: device.hardwareMAC,
  288. firmwareVersion: device.firmwareVersion,
  289. batteryLevel: device.batteryLevel,
  290. freeStorageMB: device.freeStorageMB,
  291. totalStorageMB: device.totalStorageMB
  292. )
  293. return try await client.request("devices", method: .post, body: try JSONEncoder().encode(body), authenticated: true)
  294. }
  295. func update(_ device: BoundDevice) async throws -> RemoteBoundDevice {
  296. guard let cloudID = device.cloudID else { return try await register(device) }
  297. let body = UpdateDeviceBody(
  298. name: device.name,
  299. batteryLevel: device.batteryLevel,
  300. isConnected: device.isConnected,
  301. firmwareVersion: device.firmwareVersion,
  302. freeStorageMB: device.freeStorageMB,
  303. totalStorageMB: device.totalStorageMB
  304. )
  305. return try await client.request("devices/\(cloudID)", method: .put, body: try JSONEncoder().encode(body), authenticated: true)
  306. }
  307. func remove(cloudID: String) async throws {
  308. try await client.requestVoid("devices/\(cloudID)", method: .delete, authenticated: true)
  309. }
  310. }
  311. final class NetworkStatusMonitor: ObservableObject, @unchecked Sendable {
  312. static let shared = NetworkStatusMonitor()
  313. @Published private(set) var isConnected = true
  314. @Published private(set) var isWiFi = false
  315. @Published private(set) var isCellular = false
  316. private let monitor = NWPathMonitor()
  317. private let queue = DispatchQueue(label: "com.celestia.trace.network-path")
  318. private init() {
  319. monitor.pathUpdateHandler = { [weak self] path in
  320. let connected = path.status == .satisfied
  321. let usesWiFi = path.usesInterfaceType(.wifi)
  322. let usesCellular = path.usesInterfaceType(.cellular)
  323. DispatchQueue.main.async {
  324. self?.isConnected = connected
  325. self?.isWiFi = usesWiFi
  326. self?.isCellular = usesCellular
  327. }
  328. }
  329. monitor.start(queue: queue)
  330. }
  331. var connectionName: String {
  332. if isWiFi { return "Wi-Fi" }
  333. if isCellular { return "蜂窝网络" }
  334. return isConnected ? "当前网络" : "网络"
  335. }
  336. }
  337. @MainActor
  338. final class SyncManager: ObservableObject {
  339. static let shared = SyncManager()
  340. @Published private(set) var isSyncing = false
  341. @Published private(set) var lastSyncDate: Date?
  342. @Published private(set) var lastErrorMessage: String?
  343. @Published private(set) var completedCount = 0
  344. @Published private(set) var totalCount = 0
  345. @Published private(set) var activeSessionID: UUID?
  346. @Published private(set) var syncProgress: Double = 0
  347. private let service: NetworkServiceProtocol
  348. private var activeSyncTask: Task<Bool, Never>?
  349. init(service: NetworkServiceProtocol = RemoteNetworkService()) {
  350. self.service = service
  351. lastSyncDate = UserDefaults.standard.object(forKey: "com.celestia.trace.last_server_sync") as? Date
  352. }
  353. func estimatedUploadBytes(for session: CelestiaSession) -> Int64 {
  354. var urls: [URL] = []
  355. if let url = AudioPathHelper.resolveURL(for: session.localAudioPath) {
  356. urls.append(url)
  357. }
  358. urls.append(contentsOf: session.events.compactMap { event in
  359. guard event.eventType == "PHOTO" else { return nil }
  360. return AudioPathHelper.resolveURL(for: event.localFilePath)
  361. })
  362. return urls.reduce(into: 0) { total, url in
  363. let values = try? url.resourceValues(forKeys: [.fileSizeKey])
  364. total += Int64(values?.fileSize ?? 0)
  365. }
  366. }
  367. func sync(
  368. sessions: [CelestiaSession],
  369. modelContext: ModelContext,
  370. userID: String
  371. ) async -> Bool {
  372. guard !isSyncing else { return false }
  373. isSyncing = true
  374. let task = Task { @MainActor [self] in
  375. await performSync(
  376. sessions: sessions,
  377. modelContext: modelContext,
  378. userID: userID
  379. )
  380. }
  381. activeSyncTask = task
  382. let result = await withTaskCancellationHandler {
  383. await task.value
  384. } onCancel: {
  385. task.cancel()
  386. }
  387. activeSyncTask = nil
  388. isSyncing = false
  389. activeSessionID = nil
  390. return result
  391. }
  392. func pauseSync(sessionID: UUID) {
  393. guard isSyncing, activeSessionID == sessionID else { return }
  394. activeSyncTask?.cancel()
  395. }
  396. private func performSync(
  397. sessions: [CelestiaSession],
  398. modelContext: ModelContext,
  399. userID: String
  400. ) async -> Bool {
  401. lastErrorMessage = nil
  402. completedCount = 0
  403. activeSessionID = nil
  404. syncProgress = 0
  405. totalCount = sessions.filter {
  406. ($0.ownerUserID == nil || $0.ownerUserID == userID) && (!$0.isSynced || $0.cloudSessionId == nil)
  407. }.count
  408. do {
  409. let remoteSessions = try await service.fetchSessions()
  410. try Task.checkCancellation()
  411. let mergedSessions = merge(remoteSessions, into: sessions, modelContext: modelContext, userID: userID)
  412. for (remote, local) in mergedSessions {
  413. try await downloadMissingAssets(for: local, remote: remote)
  414. try Task.checkCancellation()
  415. }
  416. try modelContext.save()
  417. let eligible = sessions.filter {
  418. ($0.ownerUserID == nil || $0.ownerUserID == userID)
  419. && (!$0.isSynced || $0.cloudSessionId == nil)
  420. && $0.syncState != .conflict
  421. }
  422. var failures: [String] = []
  423. let eligibleCount = max(eligible.count, 1)
  424. for (index, session) in eligible.enumerated() {
  425. activeSessionID = session.id
  426. let sessionStartProgress = Double(index) / Double(eligibleCount)
  427. let sessionProgressSpan = 1.0 / Double(eligibleCount)
  428. syncProgress = sessionStartProgress
  429. session.ownerUserID = userID
  430. session.syncState = .syncing
  431. session.lastSyncError = nil
  432. do {
  433. let previousRemote = remoteSessions.first {
  434. $0.id == session.cloudSessionId
  435. || $0.clientId?.caseInsensitiveCompare(session.id.uuidString) == .orderedSame
  436. }
  437. let localEventIDs = Set(session.events.map { $0.id.uuidString.lowercased() })
  438. let deletedEventClientIDs = previousRemote?.events?
  439. .compactMap(\.clientId)
  440. .filter { !localEventIDs.contains($0.lowercased()) } ?? []
  441. let remote = try await service.syncSession(
  442. session,
  443. deletedEventClientIDs: deletedEventClientIDs
  444. )
  445. try Task.checkCancellation()
  446. syncProgress = sessionStartProgress + sessionProgressSpan * 0.12
  447. session.cloudSessionId = remote.id
  448. session.serverRevision = remote.revision
  449. try await uploadLocalAssets(for: session, remote: remote) { [weak self] assetProgress in
  450. Task { @MainActor [weak self] in
  451. self?.syncProgress = sessionStartProgress
  452. + sessionProgressSpan * (0.12 + assetProgress * 0.83)
  453. }
  454. }
  455. try Task.checkCancellation()
  456. session.isSynced = true
  457. session.syncState = .synced
  458. session.lastSyncedAt = Date()
  459. completedCount += 1
  460. syncProgress = sessionStartProgress + sessionProgressSpan
  461. } catch where Task.isCancelled {
  462. session.isSynced = false
  463. session.syncState = .pending
  464. session.lastSyncError = nil
  465. syncProgress = 0
  466. try? modelContext.save()
  467. return false
  468. } catch APIError.server(let code, _) where code == 409 {
  469. session.isSynced = false
  470. session.syncState = .conflict
  471. session.lastSyncError = "本地版本 \(session.serverRevision) 与云端版本不一致"
  472. } catch {
  473. session.isSynced = false
  474. session.syncState = .failed
  475. session.lastSyncError = error.localizedDescription
  476. failures.append("\(session.title):\(error.localizedDescription)")
  477. }
  478. try modelContext.save()
  479. }
  480. let conflicts = sessions.filter { $0.syncState == .conflict }
  481. if !conflicts.isEmpty {
  482. failures.append(contentsOf: conflicts.map {
  483. "\($0.title):本地和云端版本不一致"
  484. })
  485. }
  486. guard failures.isEmpty else {
  487. lastErrorMessage = failures.joined(separator: "\n")
  488. return false
  489. }
  490. try await service.recordSyncCheckpoint()
  491. let now = Date()
  492. lastSyncDate = now
  493. UserDefaults.standard.set(now, forKey: "com.celestia.trace.last_server_sync")
  494. return true
  495. } catch where Task.isCancelled {
  496. syncProgress = 0
  497. return false
  498. } catch {
  499. lastErrorMessage = error.localizedDescription
  500. return false
  501. }
  502. }
  503. private func uploadLocalAssets(
  504. for session: CelestiaSession,
  505. remote: RemoteSession,
  506. progress: @escaping @Sendable (Double) -> Void
  507. ) async throws {
  508. guard let cloudID = session.cloudSessionId else { throw APIError.missingData }
  509. let existingIDs = Set((remote.assets ?? []).map(\.clientId))
  510. var uploads: [(clientID: String, kind: String, url: URL, size: Int64)] = []
  511. if let path = session.localAudioPath {
  512. guard let url = AudioPathHelper.resolveURL(for: path) else {
  513. throw APIError.transport("本地录音文件不存在:\(path)")
  514. }
  515. let sha256 = try fileSHA256(of: url)
  516. let hasMatchingAudio = (remote.assets ?? []).contains {
  517. $0.kind.uppercased() == "AUDIO"
  518. && $0.sha256.caseInsensitiveCompare(sha256) == .orderedSame
  519. }
  520. if !hasMatchingAudio {
  521. let clientID = "\(session.id.uuidString)-audio-\(sha256.prefix(16))"
  522. uploads.append((clientID, "AUDIO", url, fileSize(of: url)))
  523. }
  524. }
  525. for event in session.events where event.eventType == "PHOTO" {
  526. guard let path = event.localFilePath else { continue }
  527. guard let url = AudioPathHelper.resolveURL(for: path) else {
  528. throw APIError.transport("照片文件不存在:\(path)")
  529. }
  530. let clientID = "\(event.id.uuidString)-photo"
  531. if !existingIDs.contains(clientID) {
  532. uploads.append((clientID, "PHOTO", url, fileSize(of: url)))
  533. }
  534. }
  535. let totalBytes = max(uploads.reduce(Int64(0)) { $0 + $1.size }, 1)
  536. let requiredBytes = uploads.reduce(Int64(0)) { $0 + $1.size }
  537. if requiredBytes > 0 {
  538. let quota = try await service.fetchStorageQuota()
  539. guard requiredBytes <= quota.remainingBytes else {
  540. throw APIError.storageQuotaExceeded(
  541. requiredBytes: requiredBytes,
  542. remainingBytes: quota.remainingBytes
  543. )
  544. }
  545. }
  546. var completedBytes: Int64 = 0
  547. progress(uploads.isEmpty ? 1 : 0)
  548. for upload in uploads {
  549. let bytesBeforeUpload = completedBytes
  550. let result = try await service.uploadAsset(
  551. sessionID: cloudID,
  552. clientID: upload.clientID,
  553. kind: upload.kind,
  554. fileURL: upload.url
  555. ) { fraction in
  556. let sentBytes = Double(bytesBeforeUpload) + Double(upload.size) * fraction
  557. progress(min(max(sentBytes / Double(totalBytes), 0), 1))
  558. }
  559. if let revision = result.sessionRevision {
  560. session.serverRevision = revision
  561. }
  562. completedBytes += upload.size
  563. progress(Double(completedBytes) / Double(totalBytes))
  564. }
  565. }
  566. private func fileSize(of url: URL) -> Int64 {
  567. let values = try? url.resourceValues(forKeys: [.fileSizeKey])
  568. return Int64(values?.fileSize ?? 0)
  569. }
  570. private func fileSHA256(of url: URL) throws -> String {
  571. let input = try FileHandle(forReadingFrom: url)
  572. defer { try? input.close() }
  573. var hasher = SHA256()
  574. while let data = try input.read(upToCount: 1_024 * 1_024),
  575. !data.isEmpty {
  576. hasher.update(data: data)
  577. }
  578. return hasher.finalize().map { String(format: "%02x", $0) }.joined()
  579. }
  580. private func downloadMissingAssets(for session: CelestiaSession, remote: RemoteSession) async throws {
  581. let assets = remote.assets ?? []
  582. if let audio = assets
  583. .filter({ $0.kind.uppercased() == "AUDIO" })
  584. .max(by: {
  585. ($0.createdAt ?? .distantPast) < ($1.createdAt ?? .distantPast)
  586. }) {
  587. let localURL = AudioPathHelper.resolveURL(for: session.localAudioPath)
  588. let localHash = localURL.flatMap { try? fileSHA256(of: $0) }
  589. let needsDownload = localURL == nil
  590. || (session.isSynced
  591. && localHash?.caseInsensitiveCompare(audio.sha256) != .orderedSame)
  592. if needsDownload {
  593. let destination = downloadDestination(for: audio)
  594. try await service.downloadAsset(sessionID: remote.id, asset: audio, destinationURL: destination)
  595. session.localAudioPath = AudioPathHelper.relativePath(from: destination.path)
  596. }
  597. }
  598. let photoAssets = Dictionary(uniqueKeysWithValues: assets
  599. .filter { $0.kind.uppercased() == "PHOTO" }
  600. .map { ($0.clientId.lowercased(), $0) })
  601. for event in session.events where event.eventType == "PHOTO" {
  602. guard event.localFilePath == nil || AudioPathHelper.resolveURL(for: event.localFilePath) == nil else { continue }
  603. let key = "\(event.id.uuidString)-photo".lowercased()
  604. guard let asset = photoAssets[key] else { continue }
  605. let destination = downloadDestination(for: asset)
  606. try await service.downloadAsset(sessionID: remote.id, asset: asset, destinationURL: destination)
  607. event.localFilePath = AudioPathHelper.relativePath(from: destination.path)
  608. }
  609. }
  610. private func downloadDestination(for asset: RemoteAsset) -> URL {
  611. let documents = FileManager.default.urls(for: .documentDirectory, in: .userDomainMask)[0]
  612. let fileExtension = (asset.fileName as NSString).pathExtension
  613. let suffix = fileExtension.isEmpty ? "" : ".\(fileExtension.lowercased())"
  614. return documents.appendingPathComponent("cloud_\(asset.id)\(suffix)")
  615. }
  616. private func merge(
  617. _ remoteSessions: [RemoteSession],
  618. into localSessions: [CelestiaSession],
  619. modelContext: ModelContext,
  620. userID: String
  621. ) -> [(RemoteSession, CelestiaSession)] {
  622. var byClientID = Dictionary(uniqueKeysWithValues: localSessions.map { ($0.id.uuidString.lowercased(), $0) })
  623. var merged: [(RemoteSession, CelestiaSession)] = []
  624. for remote in remoteSessions {
  625. guard let clientID = remote.clientId?.lowercased() else { continue }
  626. if let local = byClientID[clientID] {
  627. if remote.deletedAt != nil, local.isSynced {
  628. modelContext.delete(local)
  629. continue
  630. }
  631. if !local.isSynced,
  632. remote.revision > local.serverRevision {
  633. local.syncState = .conflict
  634. local.lastSyncError = "本地版本 \(local.serverRevision) 与云端版本 \(remote.revision) 不一致"
  635. continue
  636. }
  637. if local.isSynced, remote.revision > local.serverRevision {
  638. apply(remote, to: local, userID: userID)
  639. }
  640. if remote.deletedAt == nil {
  641. merged.append((remote, local))
  642. }
  643. } else if remote.deletedAt == nil {
  644. let id = UUID(uuidString: clientID) ?? UUID()
  645. let session = CelestiaSession(id: id, title: remote.title, startTime: remote.startTime)
  646. apply(remote, to: session, userID: userID)
  647. modelContext.insert(session)
  648. byClientID[clientID] = session
  649. merged.append((remote, session))
  650. }
  651. }
  652. try? modelContext.save()
  653. return merged
  654. }
  655. private func apply(_ remote: RemoteSession, to local: CelestiaSession, userID: String) {
  656. local.title = remote.title
  657. local.startTime = remote.startTime
  658. local.endTime = remote.endTime
  659. local.cloudSessionId = remote.id
  660. local.ownerUserID = userID
  661. local.serverRevision = remote.revision
  662. local.isSynced = true
  663. local.syncState = .synced
  664. local.lastSyncError = nil
  665. local.events.removeAll()
  666. for remoteEvent in remote.events ?? [] {
  667. let id = remoteEvent.clientId.flatMap(UUID.init(uuidString:)) ?? UUID()
  668. let event = CelestiaTimelineEvent(id: id, relativeTimeMs: remoteEvent.relativeTimeMs, eventType: remoteEvent.eventType)
  669. event.textContent = remoteEvent.textContent
  670. event.voiceStartOffsetMs = remoteEvent.voiceStartOffsetMs
  671. event.voiceEndOffsetMs = remoteEvent.voiceEndOffsetMs
  672. local.events.append(event)
  673. }
  674. }
  675. }