RemoteNetworkService.swift 31 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834
  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. for chunkIndex in 0..<upload.totalChunks {
  208. let offset = UInt64(chunkIndex) * UInt64(upload.chunkSize)
  209. let preparedChunk = try await prepareUploadChunk(
  210. fileURL: fileURL,
  211. uploadID: upload.uploadId,
  212. chunkIndex: chunkIndex,
  213. chunkSize: upload.chunkSize,
  214. offset: offset
  215. )
  216. let chunkStart = Double(offset) / Double(fileSize)
  217. let chunkSpan = Double(preparedChunk.byteCount) / Double(fileSize)
  218. do {
  219. let _: ChunkUploadProgressResponse = try await client.upload(
  220. "sessions/\(sessionID)/assets/chunk",
  221. fileURL: preparedChunk.url,
  222. fileName: "chunk-\(chunkIndex)",
  223. mimeType: "application/octet-stream",
  224. fields: [
  225. "uploadId": upload.uploadId,
  226. "chunkIndex": String(chunkIndex)
  227. ]
  228. ) { fraction in
  229. progress(min(max(chunkStart + chunkSpan * fraction, 0), 1))
  230. }
  231. try? FileManager.default.removeItem(at: preparedChunk.url)
  232. } catch {
  233. try? FileManager.default.removeItem(at: preparedChunk.url)
  234. throw error
  235. }
  236. }
  237. let completeBody = ChunkUploadCompleteBody(
  238. uploadId: upload.uploadId,
  239. totalChunks: upload.totalChunks
  240. )
  241. let payload: AssetUploadPayload = try await client.request(
  242. "sessions/\(sessionID)/assets/complete",
  243. method: .post,
  244. body: try encoder.encode(completeBody),
  245. authenticated: true
  246. )
  247. progress(1)
  248. return AssetUploadResult(asset: payload.asset, sessionRevision: payload.sessionRevision)
  249. }
  250. private func prepareUploadChunk(
  251. fileURL: URL,
  252. uploadID: String,
  253. chunkIndex: Int,
  254. chunkSize: Int,
  255. offset: UInt64
  256. ) async throws -> (url: URL, byteCount: Int) {
  257. let task = Task.detached(priority: .utility) {
  258. try Task.checkCancellation()
  259. let input = try FileHandle(forReadingFrom: fileURL)
  260. defer { try? input.close() }
  261. try input.seek(toOffset: offset)
  262. guard let data = try input.read(upToCount: chunkSize), !data.isEmpty else {
  263. throw APIError.transport("读取上传分片失败")
  264. }
  265. try Task.checkCancellation()
  266. let temporaryURL = FileManager.default.temporaryDirectory
  267. .appendingPathComponent("celestia-\(uploadID)-\(chunkIndex).chunk")
  268. try data.write(to: temporaryURL, options: .atomic)
  269. return (temporaryURL, data.count)
  270. }
  271. return try await withTaskCancellationHandler {
  272. try await task.value
  273. } onCancel: {
  274. task.cancel()
  275. }
  276. }
  277. private func fileSize(of url: URL) throws -> Int64 {
  278. let values = try url.resourceValues(forKeys: [.fileSizeKey])
  279. guard let size = values.fileSize else {
  280. throw APIError.transport("无法读取待上传文件大小")
  281. }
  282. return Int64(size)
  283. }
  284. }
  285. struct RemoteBoundDevice: Decodable {
  286. let id: String
  287. let name: String
  288. let peripheralUUID: String
  289. }
  290. private struct BindDeviceBody: Encodable {
  291. let name: String
  292. let peripheralUUID: String
  293. let hardwareMAC: String?
  294. let firmwareVersion: String?
  295. let batteryLevel: Int?
  296. let freeStorageMB: Int?
  297. let totalStorageMB: Int?
  298. }
  299. private struct UpdateDeviceBody: Encodable {
  300. let name: String?
  301. let batteryLevel: Int?
  302. let isConnected: Bool?
  303. let firmwareVersion: String?
  304. let freeStorageMB: Int?
  305. let totalStorageMB: Int?
  306. }
  307. final class DeviceCloudService {
  308. static let shared = DeviceCloudService()
  309. private let client: APIClient
  310. init(client: APIClient = .shared) {
  311. self.client = client
  312. }
  313. func register(_ device: BoundDevice) async throws -> RemoteBoundDevice {
  314. let body = BindDeviceBody(
  315. name: device.name,
  316. peripheralUUID: device.peripheralUUID,
  317. hardwareMAC: device.hardwareMAC,
  318. firmwareVersion: device.firmwareVersion,
  319. batteryLevel: device.batteryLevel,
  320. freeStorageMB: device.freeStorageMB,
  321. totalStorageMB: device.totalStorageMB
  322. )
  323. return try await client.request("devices", method: .post, body: try JSONEncoder().encode(body), authenticated: true)
  324. }
  325. func update(_ device: BoundDevice) async throws -> RemoteBoundDevice {
  326. guard let cloudID = device.cloudID else { return try await register(device) }
  327. let body = UpdateDeviceBody(
  328. name: device.name,
  329. batteryLevel: device.batteryLevel,
  330. isConnected: device.isConnected,
  331. firmwareVersion: device.firmwareVersion,
  332. freeStorageMB: device.freeStorageMB,
  333. totalStorageMB: device.totalStorageMB
  334. )
  335. return try await client.request("devices/\(cloudID)", method: .put, body: try JSONEncoder().encode(body), authenticated: true)
  336. }
  337. func remove(cloudID: String) async throws {
  338. try await client.requestVoid("devices/\(cloudID)", method: .delete, authenticated: true)
  339. }
  340. }
  341. final class NetworkStatusMonitor: ObservableObject, @unchecked Sendable {
  342. static let shared = NetworkStatusMonitor()
  343. @Published private(set) var isConnected = true
  344. @Published private(set) var isWiFi = false
  345. @Published private(set) var isCellular = false
  346. private let monitor = NWPathMonitor()
  347. private let queue = DispatchQueue(label: "com.celestia.trace.network-path")
  348. private init() {
  349. monitor.pathUpdateHandler = { [weak self] path in
  350. let connected = path.status == .satisfied
  351. let usesWiFi = path.usesInterfaceType(.wifi)
  352. let usesCellular = path.usesInterfaceType(.cellular)
  353. DispatchQueue.main.async {
  354. self?.isConnected = connected
  355. self?.isWiFi = usesWiFi
  356. self?.isCellular = usesCellular
  357. }
  358. }
  359. monitor.start(queue: queue)
  360. }
  361. var connectionName: String {
  362. if isWiFi { return "Wi-Fi" }
  363. if isCellular { return "蜂窝网络" }
  364. return isConnected ? "当前网络" : "网络"
  365. }
  366. }
  367. @MainActor
  368. final class SyncManager: ObservableObject {
  369. static let shared = SyncManager()
  370. @Published private(set) var isSyncing = false
  371. @Published private(set) var lastSyncDate: Date?
  372. @Published private(set) var lastErrorMessage: String?
  373. @Published private(set) var completedCount = 0
  374. @Published private(set) var totalCount = 0
  375. @Published private(set) var activeSessionID: UUID?
  376. @Published private(set) var syncProgress: Double = 0
  377. private let service: NetworkServiceProtocol
  378. private var activeSyncTask: Task<Bool, Never>?
  379. private var automaticSyncTasks: [UUID: Task<Void, Never>] = [:]
  380. private var automaticSyncGenerations: [UUID: UUID] = [:]
  381. init(service: NetworkServiceProtocol = RemoteNetworkService()) {
  382. self.service = service
  383. lastSyncDate = UserDefaults.standard.object(forKey: "com.celestia.trace.last_server_sync") as? Date
  384. }
  385. func estimatedUploadBytes(for session: CelestiaSession) -> Int64 {
  386. var urls: [URL] = []
  387. if let url = AudioPathHelper.resolveURL(for: session.localAudioPath) {
  388. urls.append(url)
  389. }
  390. urls.append(contentsOf: session.events.compactMap { event in
  391. guard event.eventType == "PHOTO" else { return nil }
  392. return AudioPathHelper.resolveURL(for: event.localFilePath)
  393. })
  394. return urls.reduce(into: 0) { total, url in
  395. let values = try? url.resourceValues(forKeys: [.fileSizeKey])
  396. total += Int64(values?.fileSize ?? 0)
  397. }
  398. }
  399. func scheduleAutomaticSync(
  400. for session: CelestiaSession,
  401. modelContext: ModelContext,
  402. delay: Duration = .seconds(1.5)
  403. ) {
  404. guard session.endTime != nil,
  405. session.syncState != .conflict,
  406. AuthManager.shared.currentUser?.id != nil else { return }
  407. let sessionID = session.id
  408. let generation = UUID()
  409. automaticSyncGenerations[sessionID] = generation
  410. automaticSyncTasks[sessionID]?.cancel()
  411. if isSyncing, activeSessionID == sessionID {
  412. activeSyncTask?.cancel()
  413. }
  414. automaticSyncTasks[sessionID] = Task { @MainActor [weak self] in
  415. do {
  416. try await Task.sleep(for: delay)
  417. while let self,
  418. (self.isSyncing || !NetworkStatusMonitor.shared.isConnected) {
  419. try Task.checkCancellation()
  420. try await Task.sleep(for: .seconds(0.75))
  421. }
  422. guard let self,
  423. !Task.isCancelled,
  424. self.automaticSyncGenerations[sessionID] == generation,
  425. session.needsCloudSync,
  426. let userID = AuthManager.shared.currentUser?.id else { return }
  427. _ = await self.sync(
  428. sessions: [session],
  429. modelContext: modelContext,
  430. userID: userID
  431. )
  432. if self.automaticSyncGenerations[sessionID] == generation {
  433. self.automaticSyncTasks[sessionID] = nil
  434. self.automaticSyncGenerations[sessionID] = nil
  435. }
  436. } catch {
  437. guard let self,
  438. self.automaticSyncGenerations[sessionID] == generation else { return }
  439. self.automaticSyncTasks[sessionID] = nil
  440. self.automaticSyncGenerations[sessionID] = nil
  441. }
  442. }
  443. }
  444. func sync(
  445. sessions: [CelestiaSession],
  446. modelContext: ModelContext,
  447. userID: String
  448. ) async -> Bool {
  449. guard !isSyncing else { return false }
  450. isSyncing = true
  451. let task = Task { @MainActor [self] in
  452. await performSync(
  453. sessions: sessions,
  454. modelContext: modelContext,
  455. userID: userID
  456. )
  457. }
  458. activeSyncTask = task
  459. let result = await withTaskCancellationHandler {
  460. await task.value
  461. } onCancel: {
  462. task.cancel()
  463. }
  464. activeSyncTask = nil
  465. isSyncing = false
  466. activeSessionID = nil
  467. return result
  468. }
  469. func pauseSync(sessionID: UUID) {
  470. guard isSyncing, activeSessionID == sessionID else { return }
  471. activeSyncTask?.cancel()
  472. }
  473. private func performSync(
  474. sessions: [CelestiaSession],
  475. modelContext: ModelContext,
  476. userID: String
  477. ) async -> Bool {
  478. lastErrorMessage = nil
  479. completedCount = 0
  480. activeSessionID = nil
  481. syncProgress = 0
  482. let queuedSessions = sessions.filter {
  483. ($0.ownerUserID == nil || $0.ownerUserID == userID)
  484. && $0.needsCloudSync
  485. && $0.syncState != .conflict
  486. }
  487. totalCount = queuedSessions.count
  488. // Give the UI immediate feedback while the remote index is being fetched.
  489. // The persisted state is changed only when this session actually begins syncing.
  490. activeSessionID = queuedSessions.first?.id
  491. do {
  492. let remoteSessions = try await service.fetchSessions()
  493. try Task.checkCancellation()
  494. let mergedSessions = merge(remoteSessions, into: sessions, modelContext: modelContext, userID: userID)
  495. for (remote, local) in mergedSessions {
  496. try await downloadMissingAssets(for: local, remote: remote)
  497. try Task.checkCancellation()
  498. }
  499. try modelContext.save()
  500. let eligible = sessions.filter {
  501. ($0.ownerUserID == nil || $0.ownerUserID == userID)
  502. && $0.needsCloudSync
  503. && $0.syncState != .conflict
  504. }
  505. var failures: [String] = []
  506. let eligibleCount = max(eligible.count, 1)
  507. for (index, session) in eligible.enumerated() {
  508. activeSessionID = session.id
  509. let sessionStartProgress = Double(index) / Double(eligibleCount)
  510. let sessionProgressSpan = 1.0 / Double(eligibleCount)
  511. syncProgress = sessionStartProgress
  512. session.ownerUserID = userID
  513. session.syncState = .syncing
  514. session.lastSyncError = nil
  515. do {
  516. let previousRemote = remoteSessions.first {
  517. $0.id == session.cloudSessionId
  518. || $0.clientId?.caseInsensitiveCompare(session.id.uuidString) == .orderedSame
  519. }
  520. let localEventIDs = Set(session.events.map { $0.id.uuidString.lowercased() })
  521. let deletedEventClientIDs = previousRemote?.events?
  522. .compactMap(\.clientId)
  523. .filter { !localEventIDs.contains($0.lowercased()) } ?? []
  524. let remote = try await service.syncSession(
  525. session,
  526. deletedEventClientIDs: deletedEventClientIDs
  527. )
  528. try Task.checkCancellation()
  529. syncProgress = sessionStartProgress + sessionProgressSpan * 0.12
  530. session.cloudSessionId = remote.id
  531. session.serverRevision = remote.revision
  532. try await uploadLocalAssets(for: session, remote: remote) { [weak self] assetProgress in
  533. Task { @MainActor [weak self] in
  534. self?.syncProgress = sessionStartProgress
  535. + sessionProgressSpan * (0.12 + assetProgress * 0.83)
  536. }
  537. }
  538. try Task.checkCancellation()
  539. session.isSynced = true
  540. session.syncState = .synced
  541. session.lastSyncedAt = Date()
  542. completedCount += 1
  543. syncProgress = sessionStartProgress + sessionProgressSpan
  544. } catch where Task.isCancelled {
  545. session.isSynced = false
  546. session.syncState = .pending
  547. session.lastSyncError = nil
  548. syncProgress = 0
  549. try? modelContext.save()
  550. return false
  551. } catch APIError.server(let code, _) where code == 409 {
  552. session.isSynced = false
  553. session.syncState = .conflict
  554. session.lastSyncError = "本地版本 \(session.serverRevision) 与云端版本不一致"
  555. } catch {
  556. session.isSynced = false
  557. session.syncState = .failed
  558. session.lastSyncError = error.localizedDescription
  559. failures.append("\(session.title):\(error.localizedDescription)")
  560. }
  561. try modelContext.save()
  562. }
  563. let conflicts = sessions.filter { $0.syncState == .conflict }
  564. if !conflicts.isEmpty {
  565. failures.append(contentsOf: conflicts.map {
  566. "\($0.title):本地和云端版本不一致"
  567. })
  568. }
  569. guard failures.isEmpty else {
  570. lastErrorMessage = failures.joined(separator: "\n")
  571. return false
  572. }
  573. try await service.recordSyncCheckpoint()
  574. let now = Date()
  575. lastSyncDate = now
  576. UserDefaults.standard.set(now, forKey: "com.celestia.trace.last_server_sync")
  577. return true
  578. } catch where Task.isCancelled {
  579. syncProgress = 0
  580. return false
  581. } catch {
  582. lastErrorMessage = error.localizedDescription
  583. return false
  584. }
  585. }
  586. private func uploadLocalAssets(
  587. for session: CelestiaSession,
  588. remote: RemoteSession,
  589. progress: @escaping @Sendable (Double) -> Void
  590. ) async throws {
  591. guard let cloudID = session.cloudSessionId else { throw APIError.missingData }
  592. let existingIDs = Set((remote.assets ?? []).map(\.clientId))
  593. var uploads: [(clientID: String, kind: String, url: URL, size: Int64)] = []
  594. if let path = session.localAudioPath {
  595. guard let url = AudioPathHelper.resolveURL(for: path) else {
  596. throw APIError.transport("本地录音文件不存在:\(path)")
  597. }
  598. let sha256 = try await fileSHA256(of: url)
  599. let hasMatchingAudio = (remote.assets ?? []).contains {
  600. $0.kind.uppercased() == "AUDIO"
  601. && $0.sha256.caseInsensitiveCompare(sha256) == .orderedSame
  602. }
  603. if !hasMatchingAudio {
  604. let clientID = "\(session.id.uuidString)-audio-\(sha256.prefix(16))"
  605. uploads.append((clientID, "AUDIO", url, fileSize(of: url)))
  606. }
  607. }
  608. for event in session.events where event.eventType == "PHOTO" {
  609. guard let path = event.localFilePath else { continue }
  610. guard let url = AudioPathHelper.resolveURL(for: path) else {
  611. throw APIError.transport("照片文件不存在:\(path)")
  612. }
  613. let clientID = "\(event.id.uuidString)-photo"
  614. if !existingIDs.contains(clientID) {
  615. uploads.append((clientID, "PHOTO", url, fileSize(of: url)))
  616. }
  617. }
  618. let totalBytes = max(uploads.reduce(Int64(0)) { $0 + $1.size }, 1)
  619. let requiredBytes = uploads.reduce(Int64(0)) { $0 + $1.size }
  620. if requiredBytes > 0 {
  621. let quota = try await service.fetchStorageQuota()
  622. guard requiredBytes <= quota.remainingBytes else {
  623. throw APIError.storageQuotaExceeded(
  624. requiredBytes: requiredBytes,
  625. remainingBytes: quota.remainingBytes
  626. )
  627. }
  628. }
  629. var completedBytes: Int64 = 0
  630. progress(uploads.isEmpty ? 1 : 0)
  631. for upload in uploads {
  632. let bytesBeforeUpload = completedBytes
  633. let result = try await service.uploadAsset(
  634. sessionID: cloudID,
  635. clientID: upload.clientID,
  636. kind: upload.kind,
  637. fileURL: upload.url
  638. ) { fraction in
  639. let sentBytes = Double(bytesBeforeUpload) + Double(upload.size) * fraction
  640. progress(min(max(sentBytes / Double(totalBytes), 0), 1))
  641. }
  642. if let revision = result.sessionRevision {
  643. session.serverRevision = revision
  644. }
  645. completedBytes += upload.size
  646. progress(Double(completedBytes) / Double(totalBytes))
  647. }
  648. }
  649. private func fileSize(of url: URL) -> Int64 {
  650. let values = try? url.resourceValues(forKeys: [.fileSizeKey])
  651. return Int64(values?.fileSize ?? 0)
  652. }
  653. private func fileSHA256(of url: URL) async throws -> String {
  654. let task = Task.detached(priority: .utility) {
  655. let input = try FileHandle(forReadingFrom: url)
  656. defer { try? input.close() }
  657. var hasher = SHA256()
  658. while let data = try input.read(upToCount: 1_024 * 1_024),
  659. !data.isEmpty {
  660. try Task.checkCancellation()
  661. hasher.update(data: data)
  662. }
  663. return hasher.finalize().map { String(format: "%02x", $0) }.joined()
  664. }
  665. return try await withTaskCancellationHandler {
  666. try await task.value
  667. } onCancel: {
  668. task.cancel()
  669. }
  670. }
  671. private func downloadMissingAssets(for session: CelestiaSession, remote: RemoteSession) async throws {
  672. let assets = remote.assets ?? []
  673. if let audio = assets
  674. .filter({ $0.kind.uppercased() == "AUDIO" })
  675. .max(by: {
  676. ($0.createdAt ?? .distantPast) < ($1.createdAt ?? .distantPast)
  677. }) {
  678. let localURL = AudioPathHelper.resolveURL(for: session.localAudioPath)
  679. let localHash: String?
  680. if let localURL {
  681. localHash = try? await fileSHA256(of: localURL)
  682. } else {
  683. localHash = nil
  684. }
  685. let needsDownload = localURL == nil
  686. || (session.isSynced
  687. && localHash?.caseInsensitiveCompare(audio.sha256) != .orderedSame)
  688. if needsDownload {
  689. let destination = downloadDestination(for: audio)
  690. try await service.downloadAsset(sessionID: remote.id, asset: audio, destinationURL: destination)
  691. session.localAudioPath = AudioPathHelper.relativePath(from: destination.path)
  692. }
  693. }
  694. let photoAssets = Dictionary(uniqueKeysWithValues: assets
  695. .filter { $0.kind.uppercased() == "PHOTO" }
  696. .map { ($0.clientId.lowercased(), $0) })
  697. for event in session.events where event.eventType == "PHOTO" {
  698. guard event.localFilePath == nil || AudioPathHelper.resolveURL(for: event.localFilePath) == nil else { continue }
  699. let key = "\(event.id.uuidString)-photo".lowercased()
  700. guard let asset = photoAssets[key] else { continue }
  701. let destination = downloadDestination(for: asset)
  702. try await service.downloadAsset(sessionID: remote.id, asset: asset, destinationURL: destination)
  703. event.localFilePath = AudioPathHelper.relativePath(from: destination.path)
  704. }
  705. }
  706. private func downloadDestination(for asset: RemoteAsset) -> URL {
  707. let documents = FileManager.default.urls(for: .documentDirectory, in: .userDomainMask)[0]
  708. let fileExtension = (asset.fileName as NSString).pathExtension
  709. let suffix = fileExtension.isEmpty ? "" : ".\(fileExtension.lowercased())"
  710. return documents.appendingPathComponent("cloud_\(asset.id)\(suffix)")
  711. }
  712. private func merge(
  713. _ remoteSessions: [RemoteSession],
  714. into localSessions: [CelestiaSession],
  715. modelContext: ModelContext,
  716. userID: String
  717. ) -> [(RemoteSession, CelestiaSession)] {
  718. var byClientID = Dictionary(uniqueKeysWithValues: localSessions.map { ($0.id.uuidString.lowercased(), $0) })
  719. var merged: [(RemoteSession, CelestiaSession)] = []
  720. for remote in remoteSessions {
  721. guard let clientID = remote.clientId?.lowercased() else { continue }
  722. if let local = byClientID[clientID] {
  723. if remote.deletedAt != nil, local.isSynced {
  724. modelContext.delete(local)
  725. continue
  726. }
  727. if !local.isSynced,
  728. remote.revision > local.serverRevision {
  729. local.syncState = .conflict
  730. local.lastSyncError = "本地版本 \(local.serverRevision) 与云端版本 \(remote.revision) 不一致"
  731. continue
  732. }
  733. if local.isSynced, remote.revision > local.serverRevision {
  734. apply(remote, to: local, userID: userID)
  735. }
  736. if remote.deletedAt == nil {
  737. merged.append((remote, local))
  738. }
  739. } else if remote.deletedAt == nil {
  740. let id = UUID(uuidString: clientID) ?? UUID()
  741. let session = CelestiaSession(id: id, title: remote.title, startTime: remote.startTime)
  742. apply(remote, to: session, userID: userID)
  743. modelContext.insert(session)
  744. byClientID[clientID] = session
  745. merged.append((remote, session))
  746. }
  747. }
  748. try? modelContext.save()
  749. return merged
  750. }
  751. private func apply(_ remote: RemoteSession, to local: CelestiaSession, userID: String) {
  752. local.title = remote.title
  753. local.startTime = remote.startTime
  754. local.endTime = remote.endTime
  755. local.cloudSessionId = remote.id
  756. local.ownerUserID = userID
  757. local.serverRevision = remote.revision
  758. local.isSynced = true
  759. local.syncState = .synced
  760. local.lastSyncError = nil
  761. local.events.removeAll()
  762. for remoteEvent in remote.events ?? [] {
  763. let id = remoteEvent.clientId.flatMap(UUID.init(uuidString:)) ?? UUID()
  764. let event = CelestiaTimelineEvent(id: id, relativeTimeMs: remoteEvent.relativeTimeMs, eventType: remoteEvent.eventType)
  765. event.textContent = remoteEvent.textContent
  766. event.voiceStartOffsetMs = remoteEvent.voiceStartOffsetMs
  767. event.voiceEndOffsetMs = remoteEvent.voiceEndOffsetMs
  768. local.events.append(event)
  769. }
  770. }
  771. }