RemoteNetworkService.swift 32 KB

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