RemoteNetworkService.swift 39 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005
  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. let locationName: String?
  35. let locationAddress: String?
  36. let latitude: Double?
  37. let longitude: Double?
  38. }
  39. struct RemoteSession: Decodable {
  40. let id: String
  41. let clientId: String?
  42. let title: String
  43. let startTime: Date
  44. let endTime: Date?
  45. let durationMs: Int64
  46. let revision: Int64
  47. let deletedAt: Date?
  48. let events: [RemoteEvent]?
  49. let assets: [RemoteAsset]?
  50. }
  51. private struct SessionUpload: Encodable {
  52. let clientId: String
  53. let title: String
  54. let startTime: Date
  55. let endTime: Date?
  56. let durationMs: Int64
  57. let events: [EventUpload]
  58. let baseRevision: Int64?
  59. let deletedEventClientIds: [String]
  60. }
  61. private struct EventUpload: Encodable {
  62. let clientId: String
  63. let relativeTimeMs: Int64
  64. let eventType: String
  65. let textContent: String?
  66. let voiceStartOffsetMs: Int64?
  67. let voiceEndOffsetMs: Int64?
  68. let locationName: String?
  69. let locationAddress: String?
  70. let latitude: Double?
  71. let longitude: Double?
  72. }
  73. private struct SyncCheckpoint: Decodable {
  74. let syncedAt: Date
  75. }
  76. private struct AssetUploadPayload: Decodable {
  77. let asset: RemoteAsset
  78. let sessionRevision: Int64?
  79. private enum CodingKeys: String, CodingKey {
  80. case asset
  81. case sessionRevision
  82. }
  83. init(from decoder: Decoder) throws {
  84. if let container = try? decoder.container(keyedBy: CodingKeys.self),
  85. container.contains(.asset) {
  86. asset = try container.decode(RemoteAsset.self, forKey: .asset)
  87. sessionRevision = try container.decodeIfPresent(Int64.self, forKey: .sessionRevision)
  88. } else {
  89. asset = try RemoteAsset(from: decoder)
  90. sessionRevision = nil
  91. }
  92. }
  93. }
  94. private struct ChunkUploadInitBody: Encodable {
  95. let clientId: String
  96. let kind: String
  97. let fileName: String
  98. let mimeType: String
  99. let fileSize: Int64
  100. let chunkSize: Int
  101. }
  102. private struct ChunkUploadInitResponse: Decodable {
  103. let uploadId: String
  104. let totalChunks: Int
  105. let chunkSize: Int
  106. }
  107. private struct ChunkUploadProgressResponse: Decodable {
  108. let uploadedChunks: Int
  109. let totalChunks: Int
  110. }
  111. private struct ChunkUploadCompleteBody: Encodable {
  112. let uploadId: String
  113. let totalChunks: Int
  114. }
  115. final class RemoteNetworkService: NetworkServiceProtocol {
  116. private static let chunkedUploadThreshold: Int64 = 16 * 1024 * 1024
  117. private static let preferredChunkSize = 8 * 1024 * 1024
  118. private let client: APIClient
  119. init(client: APIClient = .shared) {
  120. self.client = client
  121. }
  122. func fetchSessions() async throws -> [RemoteSession] {
  123. try await client.request("sessions?includeDeleted=true", authenticated: true)
  124. }
  125. func syncSession(
  126. _ session: CelestiaSession,
  127. deletedEventClientIDs: [String]
  128. ) async throws -> RemoteSession {
  129. let payload = SessionUpload(
  130. clientId: session.id.uuidString,
  131. title: session.title,
  132. startTime: session.startTime,
  133. endTime: session.endTime,
  134. durationMs: session.durationMs,
  135. events: session.events.map {
  136. EventUpload(
  137. clientId: $0.id.uuidString,
  138. relativeTimeMs: $0.relativeTimeMs,
  139. eventType: $0.eventType,
  140. textContent: $0.textContent,
  141. voiceStartOffsetMs: $0.voiceStartOffsetMs,
  142. voiceEndOffsetMs: $0.voiceEndOffsetMs,
  143. locationName: $0.locationName,
  144. locationAddress: $0.locationAddress,
  145. latitude: $0.latitude,
  146. longitude: $0.longitude
  147. )
  148. },
  149. baseRevision: session.serverRevision > 0 ? session.serverRevision : nil,
  150. deletedEventClientIds: deletedEventClientIDs
  151. )
  152. let encoder = JSONEncoder()
  153. encoder.dateEncodingStrategy = .iso8601
  154. return try await client.request("sessions", method: .post, body: try encoder.encode(payload), authenticated: true)
  155. }
  156. func deleteSession(sessionID: String) async throws {
  157. try await client.requestVoid(
  158. "sessions/\(sessionID)",
  159. method: .delete,
  160. authenticated: true
  161. )
  162. }
  163. func fetchStorageQuota() async throws -> StorageQuota {
  164. try await client.request("storage/quota", authenticated: true)
  165. }
  166. func uploadAsset(
  167. sessionID: String,
  168. clientID: String,
  169. kind: String,
  170. fileURL: URL,
  171. progress: @escaping @Sendable (Double) -> Void
  172. ) async throws -> AssetUploadResult {
  173. let mimeType = UTType(filenameExtension: fileURL.pathExtension)?.preferredMIMEType ?? "application/octet-stream"
  174. let size = try fileSize(of: fileURL)
  175. if size >= Self.chunkedUploadThreshold {
  176. return try await uploadAssetInChunks(
  177. sessionID: sessionID,
  178. clientID: clientID,
  179. kind: kind,
  180. fileURL: fileURL,
  181. fileSize: size,
  182. mimeType: mimeType,
  183. progress: progress
  184. )
  185. }
  186. let payload: AssetUploadPayload = try await client.upload(
  187. "sessions/\(sessionID)/assets",
  188. fileURL: fileURL,
  189. fileName: fileURL.lastPathComponent,
  190. mimeType: mimeType,
  191. fields: ["clientId": clientID, "kind": kind],
  192. progress: progress
  193. )
  194. return AssetUploadResult(asset: payload.asset, sessionRevision: payload.sessionRevision)
  195. }
  196. func downloadAsset(sessionID: String, asset: RemoteAsset, destinationURL: URL) async throws {
  197. try await client.download("sessions/\(sessionID)/assets/\(asset.id)", to: destinationURL)
  198. }
  199. func recordSyncCheckpoint() async throws {
  200. let _: SyncCheckpoint = try await client.request("sync/trigger", method: .post, authenticated: true)
  201. }
  202. private func uploadAssetInChunks(
  203. sessionID: String,
  204. clientID: String,
  205. kind: String,
  206. fileURL: URL,
  207. fileSize: Int64,
  208. mimeType: String,
  209. progress: @escaping @Sendable (Double) -> Void
  210. ) async throws -> AssetUploadResult {
  211. let body = ChunkUploadInitBody(
  212. clientId: clientID,
  213. kind: kind,
  214. fileName: fileURL.lastPathComponent,
  215. mimeType: mimeType,
  216. fileSize: fileSize,
  217. chunkSize: Self.preferredChunkSize
  218. )
  219. let encoder = JSONEncoder()
  220. let upload: ChunkUploadInitResponse = try await client.request(
  221. "sessions/\(sessionID)/assets/init",
  222. method: .post,
  223. body: try encoder.encode(body),
  224. authenticated: true
  225. )
  226. for chunkIndex in 0..<upload.totalChunks {
  227. let offset = UInt64(chunkIndex) * UInt64(upload.chunkSize)
  228. let preparedChunk = try await prepareUploadChunk(
  229. fileURL: fileURL,
  230. uploadID: upload.uploadId,
  231. chunkIndex: chunkIndex,
  232. chunkSize: upload.chunkSize,
  233. offset: offset
  234. )
  235. let chunkStart = Double(offset) / Double(fileSize)
  236. let chunkSpan = Double(preparedChunk.byteCount) / Double(fileSize)
  237. do {
  238. let _: ChunkUploadProgressResponse = try await client.upload(
  239. "sessions/\(sessionID)/assets/chunk",
  240. fileURL: preparedChunk.url,
  241. fileName: "chunk-\(chunkIndex)",
  242. mimeType: "application/octet-stream",
  243. fields: [
  244. "uploadId": upload.uploadId,
  245. "chunkIndex": String(chunkIndex)
  246. ]
  247. ) { fraction in
  248. progress(min(max(chunkStart + chunkSpan * fraction, 0), 1))
  249. }
  250. try? FileManager.default.removeItem(at: preparedChunk.url)
  251. } catch {
  252. try? FileManager.default.removeItem(at: preparedChunk.url)
  253. throw error
  254. }
  255. }
  256. let completeBody = ChunkUploadCompleteBody(
  257. uploadId: upload.uploadId,
  258. totalChunks: upload.totalChunks
  259. )
  260. let payload: AssetUploadPayload = try await client.request(
  261. "sessions/\(sessionID)/assets/complete",
  262. method: .post,
  263. body: try encoder.encode(completeBody),
  264. authenticated: true
  265. )
  266. progress(1)
  267. return AssetUploadResult(asset: payload.asset, sessionRevision: payload.sessionRevision)
  268. }
  269. private func prepareUploadChunk(
  270. fileURL: URL,
  271. uploadID: String,
  272. chunkIndex: Int,
  273. chunkSize: Int,
  274. offset: UInt64
  275. ) async throws -> (url: URL, byteCount: Int) {
  276. let task = Task.detached(priority: .utility) {
  277. try Task.checkCancellation()
  278. let input = try FileHandle(forReadingFrom: fileURL)
  279. defer { try? input.close() }
  280. try input.seek(toOffset: offset)
  281. guard let data = try input.read(upToCount: chunkSize), !data.isEmpty else {
  282. throw APIError.transport("读取上传分片失败")
  283. }
  284. try Task.checkCancellation()
  285. let temporaryURL = FileManager.default.temporaryDirectory
  286. .appendingPathComponent("celestia-\(uploadID)-\(chunkIndex).chunk")
  287. try data.write(to: temporaryURL, options: .atomic)
  288. return (temporaryURL, data.count)
  289. }
  290. return try await withTaskCancellationHandler {
  291. try await task.value
  292. } onCancel: {
  293. task.cancel()
  294. }
  295. }
  296. private func fileSize(of url: URL) throws -> Int64 {
  297. let values = try url.resourceValues(forKeys: [.fileSizeKey])
  298. guard let size = values.fileSize else {
  299. throw APIError.transport("无法读取待上传文件大小")
  300. }
  301. return Int64(size)
  302. }
  303. }
  304. struct RemoteBoundDevice: Decodable {
  305. let id: String
  306. let name: String
  307. let peripheralUUID: String
  308. }
  309. private struct BindDeviceBody: Encodable {
  310. let name: String
  311. let peripheralUUID: String
  312. let hardwareMAC: String?
  313. let firmwareVersion: String?
  314. let batteryLevel: Int?
  315. let freeStorageMB: Int?
  316. let totalStorageMB: Int?
  317. }
  318. private struct UpdateDeviceBody: Encodable {
  319. let name: String?
  320. let batteryLevel: Int?
  321. let isConnected: Bool?
  322. let firmwareVersion: String?
  323. let freeStorageMB: Int?
  324. let totalStorageMB: Int?
  325. }
  326. final class DeviceCloudService {
  327. static let shared = DeviceCloudService()
  328. private let client: APIClient
  329. init(client: APIClient = .shared) {
  330. self.client = client
  331. }
  332. func register(_ device: BoundDevice) async throws -> RemoteBoundDevice {
  333. let body = BindDeviceBody(
  334. name: device.name,
  335. peripheralUUID: device.peripheralUUID,
  336. hardwareMAC: device.hardwareMAC,
  337. firmwareVersion: device.firmwareVersion,
  338. batteryLevel: device.batteryLevel,
  339. freeStorageMB: device.freeStorageMB,
  340. totalStorageMB: device.totalStorageMB
  341. )
  342. return try await client.request("devices", method: .post, body: try JSONEncoder().encode(body), authenticated: true)
  343. }
  344. func update(_ device: BoundDevice) async throws -> RemoteBoundDevice {
  345. guard let cloudID = device.cloudID else { return try await register(device) }
  346. let body = UpdateDeviceBody(
  347. name: device.name,
  348. batteryLevel: device.batteryLevel,
  349. isConnected: device.isConnected,
  350. firmwareVersion: device.firmwareVersion,
  351. freeStorageMB: device.freeStorageMB,
  352. totalStorageMB: device.totalStorageMB
  353. )
  354. return try await client.request("devices/\(cloudID)", method: .put, body: try JSONEncoder().encode(body), authenticated: true)
  355. }
  356. func remove(cloudID: String) async throws {
  357. try await client.requestVoid("devices/\(cloudID)", method: .delete, authenticated: true)
  358. }
  359. }
  360. final class NetworkStatusMonitor: ObservableObject, @unchecked Sendable {
  361. static let shared = NetworkStatusMonitor()
  362. @Published private(set) var isConnected = true
  363. @Published private(set) var isWiFi = false
  364. @Published private(set) var isCellular = false
  365. private let monitor = NWPathMonitor()
  366. private let queue = DispatchQueue(label: "com.celestia.trace.network-path")
  367. private init() {
  368. monitor.pathUpdateHandler = { [weak self] path in
  369. let connected = path.status == .satisfied
  370. let usesWiFi = path.usesInterfaceType(.wifi)
  371. let usesCellular = path.usesInterfaceType(.cellular)
  372. DispatchQueue.main.async {
  373. self?.isConnected = connected
  374. self?.isWiFi = usesWiFi
  375. self?.isCellular = usesCellular
  376. }
  377. }
  378. monitor.start(queue: queue)
  379. }
  380. var connectionName: String {
  381. if isWiFi { return "Wi-Fi" }
  382. if isCellular { return "蜂窝网络" }
  383. return isConnected ? "当前网络" : "网络"
  384. }
  385. }
  386. @MainActor
  387. final class SyncManager: ObservableObject {
  388. static let shared = SyncManager()
  389. @Published private(set) var isSyncing = false
  390. @Published private(set) var lastSyncDate: Date?
  391. @Published private(set) var lastErrorMessage: String?
  392. @Published private(set) var completedCount = 0
  393. @Published private(set) var totalCount = 0
  394. @Published private(set) var activeSessionID: UUID?
  395. @Published private(set) var syncProgress: Double = 0
  396. private let service: NetworkServiceProtocol
  397. private var activeSyncTask: Task<Bool, Never>?
  398. private var automaticSyncTasks: [UUID: Task<Void, Never>] = [:]
  399. private var automaticSyncGenerations: [UUID: UUID] = [:]
  400. init(service: NetworkServiceProtocol = RemoteNetworkService()) {
  401. self.service = service
  402. lastSyncDate = UserDefaults.standard.object(forKey: "com.celestia.trace.last_server_sync") as? Date
  403. }
  404. func estimatedUploadBytes(for session: CelestiaSession) -> Int64 {
  405. var urls: [URL] = []
  406. if let url = AudioPathHelper.resolveURL(for: session.localAudioPath) {
  407. urls.append(url)
  408. }
  409. urls.append(contentsOf: session.events.compactMap { event in
  410. guard event.eventType == "PHOTO" else { return nil }
  411. return AudioPathHelper.resolveURL(for: event.localFilePath)
  412. })
  413. return urls.reduce(into: 0) { total, url in
  414. let values = try? url.resourceValues(forKeys: [.fileSizeKey])
  415. total += Int64(values?.fileSize ?? 0)
  416. }
  417. }
  418. func scheduleAutomaticSync(
  419. for session: CelestiaSession,
  420. modelContext: ModelContext,
  421. delay: Duration = .seconds(1.5)
  422. ) {
  423. guard session.endTime != nil,
  424. session.syncState != .conflict,
  425. AuthManager.shared.currentUser?.id != nil else { return }
  426. let sessionID = session.id
  427. let generation = UUID()
  428. automaticSyncGenerations[sessionID] = generation
  429. automaticSyncTasks[sessionID]?.cancel()
  430. if isSyncing, activeSessionID == sessionID {
  431. activeSyncTask?.cancel()
  432. }
  433. automaticSyncTasks[sessionID] = Task { @MainActor [weak self] in
  434. do {
  435. try await Task.sleep(for: delay)
  436. while let self,
  437. (self.isSyncing || !NetworkStatusMonitor.shared.isConnected) {
  438. try Task.checkCancellation()
  439. try await Task.sleep(for: .seconds(0.75))
  440. }
  441. guard let self,
  442. !Task.isCancelled,
  443. self.automaticSyncGenerations[sessionID] == generation,
  444. session.needsCloudSync,
  445. let userID = AuthManager.shared.currentUser?.id else { return }
  446. _ = await self.sync(
  447. sessions: [session],
  448. modelContext: modelContext,
  449. userID: userID
  450. )
  451. if self.automaticSyncGenerations[sessionID] == generation {
  452. self.automaticSyncTasks[sessionID] = nil
  453. self.automaticSyncGenerations[sessionID] = nil
  454. }
  455. } catch {
  456. guard let self,
  457. self.automaticSyncGenerations[sessionID] == generation else { return }
  458. self.automaticSyncTasks[sessionID] = nil
  459. self.automaticSyncGenerations[sessionID] = nil
  460. }
  461. }
  462. }
  463. func sync(
  464. sessions: [CelestiaSession],
  465. modelContext: ModelContext,
  466. userID: String
  467. ) async -> Bool {
  468. guard !isSyncing else {
  469. DeveloperLogStore.log("现场记录同步", "已有同步任务正在运行,忽略重复请求", level: .warning)
  470. return false
  471. }
  472. DeveloperLogStore.log("现场记录同步", "开始同步,本地共 \(sessions.count) 条记录")
  473. isSyncing = true
  474. let task = Task { @MainActor [self] in
  475. await performSync(
  476. sessions: sessions,
  477. modelContext: modelContext,
  478. userID: userID
  479. )
  480. }
  481. activeSyncTask = task
  482. let result = await withTaskCancellationHandler {
  483. await task.value
  484. } onCancel: {
  485. task.cancel()
  486. }
  487. activeSyncTask = nil
  488. isSyncing = false
  489. activeSessionID = nil
  490. DeveloperLogStore.log(
  491. "现场记录同步",
  492. result ? "同步任务完成" : "同步任务未完成",
  493. level: result ? .success : .warning
  494. )
  495. return result
  496. }
  497. func pauseSync(sessionID: UUID) {
  498. guard isSyncing, activeSessionID == sessionID else { return }
  499. activeSyncTask?.cancel()
  500. }
  501. private func performSync(
  502. sessions: [CelestiaSession],
  503. modelContext: ModelContext,
  504. userID: String
  505. ) async -> Bool {
  506. lastErrorMessage = nil
  507. completedCount = 0
  508. activeSessionID = nil
  509. syncProgress = 0
  510. let queuedSessions = deduplicatedSessions(sessions.filter {
  511. ($0.ownerUserID == nil || $0.ownerUserID == userID)
  512. && $0.needsCloudSync
  513. && $0.syncState != .conflict
  514. })
  515. totalCount = queuedSessions.count
  516. DeveloperLogStore.log("现场记录同步", "待上传 \(queuedSessions.count) 条记录,正在获取云端索引")
  517. // Give the UI immediate feedback while the remote index is being fetched.
  518. // The persisted state is changed only when this session actually begins syncing.
  519. activeSessionID = queuedSessions.first?.id
  520. do {
  521. let remoteSessions = try await service.fetchSessions()
  522. DeveloperLogStore.log("现场记录同步", "获取到 \(remoteSessions.count) 条云端记录", level: .success)
  523. try Task.checkCancellation()
  524. // Detail and post-recording syncs may upload only one session, while
  525. // fetchSessions always returns the complete cloud history. Use the
  526. // complete local store as the merge baseline to avoid inserting the
  527. // other cloud sessions again with duplicate client IDs.
  528. let allLocalSessions = try modelContext.fetch(FetchDescriptor<CelestiaSession>())
  529. let mergedSessions = try await merge(
  530. remoteSessions,
  531. into: allLocalSessions,
  532. modelContext: modelContext,
  533. userID: userID
  534. )
  535. for (remote, local) in mergedSessions {
  536. try await downloadMissingAssets(for: local, remote: remote)
  537. try Task.checkCancellation()
  538. await Task.yield()
  539. }
  540. try modelContext.save()
  541. let eligible = deduplicatedSessions(sessions.filter {
  542. ($0.ownerUserID == nil || $0.ownerUserID == userID)
  543. && $0.needsCloudSync
  544. && $0.syncState != .conflict
  545. })
  546. var failures: [String] = []
  547. let eligibleCount = max(eligible.count, 1)
  548. for (index, session) in eligible.enumerated() {
  549. try Task.checkCancellation()
  550. await Task.yield()
  551. activeSessionID = session.id
  552. let sessionStartProgress = Double(index) / Double(eligibleCount)
  553. let sessionProgressSpan = 1.0 / Double(eligibleCount)
  554. syncProgress = sessionStartProgress
  555. session.ownerUserID = userID
  556. session.syncState = .syncing
  557. session.lastSyncError = nil
  558. DeveloperLogStore.log(
  559. "现场记录同步",
  560. "正在同步 \(index + 1)/\(eligible.count):\(session.title)"
  561. )
  562. do {
  563. let previousRemote = remoteSessions.first {
  564. $0.id == session.cloudSessionId
  565. || $0.clientId?.caseInsensitiveCompare(session.id.uuidString) == .orderedSame
  566. }
  567. let localEventIDs = Set(session.events.map { $0.id.uuidString.lowercased() })
  568. let deletedEventClientIDs = previousRemote?.events?
  569. .compactMap(\.clientId)
  570. .filter { !localEventIDs.contains($0.lowercased()) } ?? []
  571. let remote = try await service.syncSession(
  572. session,
  573. deletedEventClientIDs: deletedEventClientIDs
  574. )
  575. try Task.checkCancellation()
  576. syncProgress = sessionStartProgress + sessionProgressSpan * 0.12
  577. session.cloudSessionId = remote.id
  578. session.serverRevision = remote.revision
  579. try await uploadLocalAssets(for: session, remote: remote) { [weak self] assetProgress in
  580. Task { @MainActor [weak self] in
  581. self?.syncProgress = sessionStartProgress
  582. + sessionProgressSpan * (0.12 + assetProgress * 0.83)
  583. }
  584. }
  585. try Task.checkCancellation()
  586. session.isSynced = true
  587. session.syncState = .synced
  588. session.lastSyncedAt = Date()
  589. completedCount += 1
  590. syncProgress = sessionStartProgress + sessionProgressSpan
  591. DeveloperLogStore.log(
  592. "现场记录同步",
  593. "\(session.title) 同步成功,云端版本 v\(session.serverRevision)",
  594. level: .success
  595. )
  596. } catch where Task.isCancelled {
  597. session.isSynced = false
  598. session.syncState = .pending
  599. session.lastSyncError = nil
  600. syncProgress = 0
  601. try? modelContext.save()
  602. return false
  603. } catch APIError.server(let code, _) where code == 409 {
  604. session.isSynced = false
  605. session.syncState = .conflict
  606. session.lastSyncError = "本地版本 \(session.serverRevision) 与云端版本不一致"
  607. DeveloperLogStore.log(
  608. "现场记录同步",
  609. "\(session.title) 发生版本冲突",
  610. level: .error
  611. )
  612. } catch {
  613. session.isSynced = false
  614. session.syncState = .failed
  615. session.lastSyncError = error.localizedDescription
  616. failures.append("\(session.title):\(error.localizedDescription)")
  617. DeveloperLogStore.log(
  618. "现场记录同步",
  619. "\(session.title) 同步失败:\(error.localizedDescription)",
  620. level: .error
  621. )
  622. }
  623. try modelContext.save()
  624. }
  625. let conflicts = sessions.filter { $0.syncState == .conflict }
  626. if !conflicts.isEmpty {
  627. failures.append(contentsOf: conflicts.map {
  628. "\($0.title):本地和云端版本不一致"
  629. })
  630. }
  631. guard failures.isEmpty else {
  632. lastErrorMessage = failures.joined(separator: "\n")
  633. DeveloperLogStore.log("现场记录同步", "同步结束,存在 \(failures.count) 个问题", level: .error)
  634. return false
  635. }
  636. try await service.recordSyncCheckpoint()
  637. let now = Date()
  638. lastSyncDate = now
  639. UserDefaults.standard.set(now, forKey: "com.celestia.trace.last_server_sync")
  640. DeveloperLogStore.log("现场记录同步", "云端检查点已更新", level: .success)
  641. return true
  642. } catch where Task.isCancelled {
  643. syncProgress = 0
  644. return false
  645. } catch {
  646. lastErrorMessage = error.localizedDescription
  647. DeveloperLogStore.log("现场记录同步", "同步中断:\(error.localizedDescription)", level: .error)
  648. return false
  649. }
  650. }
  651. private func uploadLocalAssets(
  652. for session: CelestiaSession,
  653. remote: RemoteSession,
  654. progress: @escaping @Sendable (Double) -> Void
  655. ) async throws {
  656. guard let cloudID = session.cloudSessionId else { throw APIError.missingData }
  657. let existingIDs = Set((remote.assets ?? []).map(\.clientId))
  658. var uploads: [(clientID: String, kind: String, url: URL, size: Int64)] = []
  659. if let path = session.localAudioPath {
  660. guard let url = AudioPathHelper.resolveURL(for: path) else {
  661. throw APIError.transport("本地录音文件不存在:\(path)")
  662. }
  663. let sha256 = try await fileSHA256(of: url)
  664. let hasMatchingAudio = (remote.assets ?? []).contains {
  665. $0.kind.uppercased() == "AUDIO"
  666. && $0.sha256.caseInsensitiveCompare(sha256) == .orderedSame
  667. }
  668. if !hasMatchingAudio {
  669. let clientID = "\(session.id.uuidString)-audio-\(sha256.prefix(16))"
  670. uploads.append((clientID, "AUDIO", url, fileSize(of: url)))
  671. }
  672. }
  673. for event in session.events where event.eventType == "PHOTO" {
  674. guard let path = event.localFilePath else { continue }
  675. guard let url = AudioPathHelper.resolveURL(for: path) else {
  676. throw APIError.transport("照片文件不存在:\(path)")
  677. }
  678. let clientID = "\(event.id.uuidString)-photo"
  679. if !existingIDs.contains(clientID) {
  680. uploads.append((clientID, "PHOTO", url, fileSize(of: url)))
  681. }
  682. }
  683. let totalBytes = max(uploads.reduce(Int64(0)) { $0 + $1.size }, 1)
  684. let requiredBytes = uploads.reduce(Int64(0)) { $0 + $1.size }
  685. if requiredBytes > 0 {
  686. let quota = try await service.fetchStorageQuota()
  687. guard requiredBytes <= quota.remainingBytes else {
  688. throw APIError.storageQuotaExceeded(
  689. requiredBytes: requiredBytes,
  690. remainingBytes: quota.remainingBytes
  691. )
  692. }
  693. }
  694. var completedBytes: Int64 = 0
  695. progress(uploads.isEmpty ? 1 : 0)
  696. for upload in uploads {
  697. let bytesBeforeUpload = completedBytes
  698. let result = try await service.uploadAsset(
  699. sessionID: cloudID,
  700. clientID: upload.clientID,
  701. kind: upload.kind,
  702. fileURL: upload.url
  703. ) { fraction in
  704. let sentBytes = Double(bytesBeforeUpload) + Double(upload.size) * fraction
  705. progress(min(max(sentBytes / Double(totalBytes), 0), 1))
  706. }
  707. if let revision = result.sessionRevision {
  708. session.serverRevision = revision
  709. }
  710. completedBytes += upload.size
  711. progress(Double(completedBytes) / Double(totalBytes))
  712. }
  713. }
  714. private func fileSize(of url: URL) -> Int64 {
  715. let values = try? url.resourceValues(forKeys: [.fileSizeKey])
  716. return Int64(values?.fileSize ?? 0)
  717. }
  718. private func fileSHA256(of url: URL) async throws -> String {
  719. let task = Task.detached(priority: .utility) {
  720. let input = try FileHandle(forReadingFrom: url)
  721. defer { try? input.close() }
  722. var hasher = SHA256()
  723. while let data = try input.read(upToCount: 1_024 * 1_024),
  724. !data.isEmpty {
  725. try Task.checkCancellation()
  726. hasher.update(data: data)
  727. }
  728. return hasher.finalize().map { String(format: "%02x", $0) }.joined()
  729. }
  730. return try await withTaskCancellationHandler {
  731. try await task.value
  732. } onCancel: {
  733. task.cancel()
  734. }
  735. }
  736. private func downloadMissingAssets(for session: CelestiaSession, remote: RemoteSession) async throws {
  737. let assets = remote.assets ?? []
  738. if let audio = assets
  739. .filter({ $0.kind.uppercased() == "AUDIO" })
  740. .max(by: {
  741. ($0.createdAt ?? .distantPast) < ($1.createdAt ?? .distantPast)
  742. }) {
  743. let localURL = AudioPathHelper.resolveURL(for: session.localAudioPath)
  744. let localHash: String?
  745. if let localURL {
  746. localHash = try? await fileSHA256(of: localURL)
  747. } else {
  748. localHash = nil
  749. }
  750. let needsDownload = localURL == nil
  751. || (session.isSynced
  752. && localHash?.caseInsensitiveCompare(audio.sha256) != .orderedSame)
  753. if needsDownload {
  754. let destination = downloadDestination(for: audio)
  755. try await service.downloadAsset(sessionID: remote.id, asset: audio, destinationURL: destination)
  756. session.localAudioPath = AudioPathHelper.relativePath(from: destination.path)
  757. }
  758. }
  759. let photoAssets = Dictionary(
  760. assets
  761. .filter { $0.kind.uppercased() == "PHOTO" }
  762. .map { ($0.clientId.lowercased(), $0) },
  763. uniquingKeysWith: { current, candidate in
  764. (candidate.createdAt ?? .distantPast) > (current.createdAt ?? .distantPast)
  765. ? candidate
  766. : current
  767. }
  768. )
  769. let remotePhotoCount = assets.lazy.filter { $0.kind.uppercased() == "PHOTO" }.count
  770. if photoAssets.count < remotePhotoCount {
  771. DeveloperLogStore.log(
  772. "现场记录同步",
  773. "云端照片资源存在重复 clientId,已选择最新资源继续同步",
  774. level: .warning
  775. )
  776. }
  777. for event in session.events where event.eventType == "PHOTO" {
  778. guard event.localFilePath == nil || AudioPathHelper.resolveURL(for: event.localFilePath) == nil else { continue }
  779. let key = "\(event.id.uuidString)-photo".lowercased()
  780. guard let asset = photoAssets[key] else { continue }
  781. let destination = downloadDestination(for: asset)
  782. try await service.downloadAsset(sessionID: remote.id, asset: asset, destinationURL: destination)
  783. event.localFilePath = AudioPathHelper.relativePath(from: destination.path)
  784. }
  785. }
  786. private func downloadDestination(for asset: RemoteAsset) -> URL {
  787. let documents = FileManager.default.urls(for: .documentDirectory, in: .userDomainMask)[0]
  788. let fileExtension = (asset.fileName as NSString).pathExtension
  789. let suffix = fileExtension.isEmpty ? "" : ".\(fileExtension.lowercased())"
  790. return documents.appendingPathComponent("cloud_\(asset.id)\(suffix)")
  791. }
  792. private func merge(
  793. _ remoteSessions: [RemoteSession],
  794. into localSessions: [CelestiaSession],
  795. modelContext: ModelContext,
  796. userID: String
  797. ) async throws -> [(RemoteSession, CelestiaSession)] {
  798. var byClientID = Dictionary(
  799. localSessions.map { ($0.id.uuidString.lowercased(), $0) },
  800. uniquingKeysWith: preferredLocalSession
  801. )
  802. if byClientID.count < localSessions.count {
  803. DeveloperLogStore.log(
  804. "现场记录同步",
  805. "检测到 \(localSessions.count - byClientID.count) 条重复本地记录,已选择云端版本较新的记录继续同步",
  806. level: .warning
  807. )
  808. }
  809. let activeRemoteSessions = remoteSessions.filter { $0.deletedAt == nil }
  810. let activeRemoteIDs = Set(activeRemoteSessions.map(\.id))
  811. let activeRemoteClientIDs = Set(activeRemoteSessions.compactMap { $0.clientId?.lowercased() })
  812. var deletedLocalObjects: Set<ObjectIdentifier> = []
  813. var deletedFileURLs: Set<URL> = []
  814. var merged: [(RemoteSession, CelestiaSession)] = []
  815. for (index, remote) in remoteSessions.enumerated() {
  816. if index.isMultiple(of: 20) {
  817. await Task.yield()
  818. }
  819. guard let clientID = remote.clientId?.lowercased() else { continue }
  820. if let local = byClientID[clientID] {
  821. if remote.deletedAt != nil, local.hasConfirmedCloudSync {
  822. deletedFileURLs.formUnion(localFileURLs(for: local))
  823. deletedLocalObjects.insert(ObjectIdentifier(local))
  824. modelContext.delete(local)
  825. DeveloperLogStore.log(
  826. "现场记录同步",
  827. "云端记录已删除,同步移除本地记录:\(local.title)"
  828. )
  829. continue
  830. }
  831. if !local.isSynced,
  832. remote.revision > local.serverRevision {
  833. local.syncState = .conflict
  834. local.lastSyncError = "本地版本 \(local.serverRevision) 与云端版本 \(remote.revision) 不一致"
  835. continue
  836. }
  837. if local.isSynced, remote.revision > local.serverRevision {
  838. await apply(remote, to: local, userID: userID)
  839. }
  840. if remote.deletedAt == nil {
  841. merged.append((remote, local))
  842. }
  843. } else if remote.deletedAt == nil {
  844. let id = UUID(uuidString: clientID) ?? UUID()
  845. let session = CelestiaSession(id: id, title: remote.title, startTime: remote.startTime)
  846. await apply(remote, to: session, userID: userID)
  847. modelContext.insert(session)
  848. byClientID[clientID] = session
  849. merged.append((remote, session))
  850. }
  851. }
  852. // Some servers may permanently remove a record instead of returning a
  853. // tombstone. A previously confirmed cloud record that is absent by both
  854. // server ID and client ID should mirror that deletion locally. Records
  855. // with pending/failed/conflicting local content are deliberately kept.
  856. for (index, local) in localSessions.enumerated() {
  857. if index.isMultiple(of: 20) {
  858. await Task.yield()
  859. }
  860. guard !deletedLocalObjects.contains(ObjectIdentifier(local)),
  861. local.ownerUserID == nil || local.ownerUserID == userID,
  862. local.hasConfirmedCloudSync else { continue }
  863. let existsRemotely =
  864. local.cloudSessionId.map(activeRemoteIDs.contains) == true
  865. || activeRemoteClientIDs.contains(local.id.uuidString.lowercased())
  866. guard !existsRemotely else { continue }
  867. deletedFileURLs.formUnion(localFileURLs(for: local))
  868. deletedLocalObjects.insert(ObjectIdentifier(local))
  869. modelContext.delete(local)
  870. DeveloperLogStore.log(
  871. "现场记录同步",
  872. "云端不存在已同步记录,同步移除本地记录:\(local.title)"
  873. )
  874. }
  875. try modelContext.save()
  876. for fileURL in deletedFileURLs {
  877. try? FileManager.default.removeItem(at: fileURL)
  878. }
  879. return merged
  880. }
  881. private func localFileURLs(for session: CelestiaSession) -> Set<URL> {
  882. let storedPaths =
  883. [session.localAudioPath].compactMap { $0 }
  884. + session.audioChunks.map(\.localFilePath)
  885. + session.events.compactMap(\.localFilePath)
  886. return Set(storedPaths.compactMap { AudioPathHelper.resolveURL(for: $0) })
  887. }
  888. private func deduplicatedSessions(_ sessions: [CelestiaSession]) -> [CelestiaSession] {
  889. Array(Dictionary(
  890. sessions.map { ($0.id, $0) },
  891. uniquingKeysWith: preferredLocalSession
  892. ).values)
  893. }
  894. private func preferredLocalSession(
  895. _ current: CelestiaSession,
  896. _ candidate: CelestiaSession
  897. ) -> CelestiaSession {
  898. if (current.cloudSessionId == nil) != (candidate.cloudSessionId == nil) {
  899. return current.cloudSessionId == nil ? candidate : current
  900. }
  901. if current.serverRevision != candidate.serverRevision {
  902. return current.serverRevision > candidate.serverRevision ? current : candidate
  903. }
  904. return (current.lastSyncedAt ?? .distantPast) >= (candidate.lastSyncedAt ?? .distantPast)
  905. ? current
  906. : candidate
  907. }
  908. private func apply(_ remote: RemoteSession, to local: CelestiaSession, userID: String) async {
  909. local.title = remote.title
  910. local.startTime = remote.startTime
  911. local.endTime = remote.endTime
  912. local.cloudSessionId = remote.id
  913. local.ownerUserID = userID
  914. local.serverRevision = remote.revision
  915. local.isSynced = true
  916. local.syncState = .synced
  917. local.lastSyncError = nil
  918. local.events.removeAll()
  919. for (index, remoteEvent) in (remote.events ?? []).enumerated() {
  920. if index.isMultiple(of: 50) {
  921. await Task.yield()
  922. }
  923. let id = remoteEvent.clientId.flatMap(UUID.init(uuidString:)) ?? UUID()
  924. let event = CelestiaTimelineEvent(id: id, relativeTimeMs: remoteEvent.relativeTimeMs, eventType: remoteEvent.eventType)
  925. event.textContent = remoteEvent.textContent
  926. event.voiceStartOffsetMs = remoteEvent.voiceStartOffsetMs
  927. event.voiceEndOffsetMs = remoteEvent.voiceEndOffsetMs
  928. event.locationName = remoteEvent.locationName
  929. event.locationAddress = remoteEvent.locationAddress
  930. event.latitude = remoteEvent.latitude
  931. event.longitude = remoteEvent.longitude
  932. local.events.append(event)
  933. }
  934. }
  935. }