// // CameraTransferPipeline.swift // suixinkan // // Created by Codex on 2026/6/29. // import Foundation /// 相机传输编排器,负责监听新照片、下载到本地并按需调用上传 Sink。 @MainActor final class CameraTransferPipeline { var onTasksUpdated: (([CameraTransferTask]) -> Void)? private(set) var tasks: [CameraTransferTask] = [] private let cameraService: any CameraServiceProtocol private var uploadSink: (any CameraAssetUploadSink)? private var uploadEnabled = true private var activeUploadAssetIDs: Set = [] private var pendingDeferredNotifyTask: Task? private var lastDeferredNotifyDate = Date.distantPast private var lastNotifiedProgressByAssetID: [String: Int] = [:] private var inFlightAssetIDs: Set = [] private let maxRetries = 3 private let deferredNotifyInterval: TimeInterval = 0.15 private let maxConcurrentUploads = 3 /// 初始化传输管道,并注入相机服务。 init(cameraService: any CameraServiceProtocol) { self.cameraService = cameraService cameraService.onNewAsset = { [weak self] asset in Task { @MainActor in await self?.handleNewAsset(asset, skipIfExists: false) } } } /// 配置上传 Sink 与是否自动上传。 func configure(uploadSink: (any CameraAssetUploadSink)?, uploadEnabled: Bool) { self.uploadSink = uploadSink self.uploadEnabled = uploadEnabled } /// 启动相机连接。 func connect() async { await cameraService.connect() } /// 断开相机连接。 func disconnect() async { await cameraService.disconnect() } /// 当前相机连接状态。 var connectionState: CameraConnectionState { cameraService.connectionState } /// 设置连接状态变化回调。 var onConnectionStateChange: ((CameraConnectionState) -> Void)? { get { nil } set { cameraService.onConnectionStateChange = newValue } } /// 重试所有失败的上传任务。 func retryFailedUploads() async { guard uploadEnabled, uploadSink != nil else { return } for index in tasks.indices where tasks[index].status == .failed { tasks[index].status = .downloaded tasks[index].errorMessage = nil } notify() processUploadQueue() } /// 手动上传指定 asset 列表(拍完再传模式)。 func uploadAssets(withIDs assetIDs: [String]) async { guard uploadSink != nil else { return } for assetID in assetIDs { if activeUploadAssetIDs.contains(assetID) { continue } guard let index = tasks.firstIndex(where: { $0.assetID == assetID }) else { continue } if tasks[index].status == .uploaded { continue } if tasks[index].localURL == nil { await downloadAsset(assetID: assetID) } guard let refreshedIndex = tasks.firstIndex(where: { $0.assetID == assetID }) else { continue } if tasks[refreshedIndex].status == .downloaded || tasks[refreshedIndex].status == .failed { tasks[refreshedIndex].status = .downloaded tasks[refreshedIndex].errorMessage = nil await uploadTask(assetID: assetID) } } } /// 重传单张照片;若管道中无任务则根据本地相对路径恢复。 func retryUpload(assetID: String, filename: String, localPath: String) async { guard uploadSink != nil else { return } guard !activeUploadAssetIDs.contains(assetID) else { return } if let index = tasks.firstIndex(where: { $0.assetID == assetID }) { if tasks[index].status == .uploaded { return } if tasks[index].localPath?.isEmpty != false, !localPath.isEmpty { tasks[index].localPath = localPath } tasks[index].status = .downloaded tasks[index].errorMessage = nil } else { guard !localPath.isEmpty, let localURL = CameraDownloadStorage.resolveLocalURL(from: localPath), FileManager.default.fileExists(atPath: localURL.path) else { return } tasks.insert( CameraTransferTask( assetID: assetID, filename: filename, localPath: localPath, capturedAt: CameraTransferTask.capturedAtString(forLocalPath: localPath), status: .downloaded ), at: 0 ) } notify() await uploadTask(assetID: assetID) } /// 同步相机内已有照片到任务列表(不自动上传)。 func syncExistingPhotos() async { guard cameraService.connectionState.isConnected else { return } do { let assets = try await cameraService.listAssets() for asset in assets { await handleNewAsset(asset, skipIfExists: true) } } catch { CameraTetheringLogger.log("syncExistingPhotos 失败: \(error.localizedDescription)", logger: CameraTetheringLogger.transfer) } } /// 清空所有传输记录及本地文件。 func clearAllTransfers() { for task in tasks { guard let localURL = task.localURL, FileManager.default.fileExists(atPath: localURL.path) else { continue } try? FileManager.default.removeItem(at: localURL) } tasks.removeAll() notify() } /// 按 assetID 查找任务。 func task(forAssetID assetID: String) -> CameraTransferTask? { tasks.first { $0.assetID == assetID } } private func handleNewAsset(_ asset: CameraAsset, skipIfExists: Bool) async { if skipIfExists, let existing = tasks.first(where: { $0.assetID == asset.id }), existing.status == .uploaded || existing.status == .downloaded { return } if inFlightAssetIDs.contains(asset.id) { return } if tasks.contains(where: { $0.assetID == asset.id && $0.status == .downloading }) { return } inFlightAssetIDs.insert(asset.id) defer { inFlightAssetIDs.remove(asset.id) } let assetID = asset.id if let existingIndex = tasks.firstIndex(where: { $0.assetID == assetID }) { tasks[existingIndex].status = .downloading tasks[existingIndex].errorMessage = nil } else { let task = CameraTransferTask( assetID: asset.id, filename: asset.filename, capturedAt: CameraTransferTask.capturedAtString(from: asset.creationDate), status: .downloading ) tasks.insert(task, at: 0) } notify() await downloadAsset(assetID: assetID, asset: asset) if uploadEnabled, uploadSink != nil, tasks.first(where: { $0.assetID == assetID })?.status == .downloaded { processUploadQueue() } } /// 下载指定 asset 到本地;始终按 assetID 定位任务,避免连拍时 index 失效。 private func downloadAsset(assetID: String, asset: CameraAsset? = nil) async { guard tasks.contains(where: { $0.assetID == assetID }) else { return } do { let resolvedAsset: CameraAsset if let asset { resolvedAsset = asset } else if let cached = try? await cameraService.listAssets().first(where: { $0.id == assetID }) { resolvedAsset = cached } else { throw CameraServiceError.assetNotFound } let localURL = try await cameraService.downloadAsset(resolvedAsset) let destination: URL if CameraDownloadStorage.isInDownloadsDirectory(localURL) { destination = localURL } else { destination = CameraDownloadStorage.uniqueLocalURL(for: resolvedAsset.filename) if FileManager.default.fileExists(atPath: destination.path) { try FileManager.default.removeItem(at: destination) } try FileManager.default.moveItem(at: localURL, to: destination) } guard let index = tasks.firstIndex(where: { $0.assetID == assetID }) else { return } tasks[index].localPath = CameraDownloadStorage.relativePath(for: destination) if tasks[index].capturedAt?.isEmpty != false { tasks[index].capturedAt = CameraTransferTask.capturedAtString(from: resolvedAsset.creationDate) } tasks[index].status = .downloaded tasks[index].progress = 0 notify() } catch { guard let index = tasks.firstIndex(where: { $0.assetID == assetID }) else { return } tasks[index].status = .failed tasks[index].errorMessage = error.localizedDescription notify() } } /// 并发上传已下载任务;边拍边传时限制并发数,避免网络与内存被打满。 private func processUploadQueue() { guard uploadEnabled, uploadSink != nil else { return } while activeUploadAssetIDs.count < maxConcurrentUploads, let task = tasks.first(where: { $0.status == .downloaded && !activeUploadAssetIDs.contains($0.assetID) }) { startConcurrentUpload(assetID: task.assetID) } } private func startConcurrentUpload(assetID: String) { guard !activeUploadAssetIDs.contains(assetID) else { return } activeUploadAssetIDs.insert(assetID) Task { @MainActor in await self.uploadTask(assetID: assetID) self.activeUploadAssetIDs.remove(assetID) self.processUploadQueue() } } private func uploadTask(assetID: String) async { guard let index = tasks.firstIndex(where: { $0.assetID == assetID }), let sink = uploadSink, let localURL = tasks[index].localURL else { return } tasks[index].status = .uploading tasks[index].errorMessage = nil notify() let filename = tasks[index].filename let fileType = filename.lowercased().hasSuffix(".mp4") ? 1 : 2 var attempt = 0 while attempt < maxRetries { do { let remoteURL = try await sink.upload( localURL: localURL, fileName: filename, fileType: fileType, progress: { [weak self] progress in Task { @MainActor in guard let self, let taskIndex = self.tasks.firstIndex(where: { $0.assetID == assetID }) else { return } let lastNotifiedProgress = self.lastNotifiedProgressByAssetID[assetID] ?? -1 guard progress == 100 || progress - lastNotifiedProgress >= 10 else { return } self.lastNotifiedProgressByAssetID[assetID] = progress self.tasks[taskIndex].progress = progress self.notify(progress == 100 ? .immediate : .deferred) } } ) guard let taskIndex = tasks.firstIndex(where: { $0.assetID == assetID }) else { return } tasks[taskIndex].status = .uploaded tasks[taskIndex].remoteURL = remoteURL tasks[taskIndex].progress = 100 lastNotifiedProgressByAssetID.removeValue(forKey: assetID) notify() return } catch { attempt += 1 guard let taskIndex = tasks.firstIndex(where: { $0.assetID == assetID }) else { return } if attempt >= maxRetries { tasks[taskIndex].status = .failed tasks[taskIndex].errorMessage = error.localizedDescription lastNotifiedProgressByAssetID.removeValue(forKey: assetID) notify() } else { let delay = UInt64(pow(2.0, Double(attempt))) try? await Task.sleep(nanoseconds: delay * 1_000_000_000) } } } } private func notify(_ mode: CameraTransferNotifyMode = .immediate) { switch mode { case .immediate: pendingDeferredNotifyTask?.cancel() pendingDeferredNotifyTask = nil onTasksUpdated?(tasks) case .deferred: let elapsed = Date().timeIntervalSince(lastDeferredNotifyDate) if elapsed >= deferredNotifyInterval { lastDeferredNotifyDate = Date() onTasksUpdated?(tasks) return } guard pendingDeferredNotifyTask == nil else { return } let delay = UInt64((deferredNotifyInterval - elapsed) * 1_000_000_000) pendingDeferredNotifyTask = Task { @MainActor in try? await Task.sleep(nanoseconds: delay) guard !Task.isCancelled else { return } self.lastDeferredNotifyDate = Date() self.pendingDeferredNotifyTask = nil self.onTasksUpdated?(self.tasks) } } } } private enum CameraTransferNotifyMode { case immediate case deferred }