Files
suixinkan_ios_new/suixinkan/Core/CameraTransfer/CameraTransferPipeline.swift
汉秋 90d135b5ac 完善 OTG 相机有线传图:支持 Sony/Canon 自动识别,页面 detach 保持会话以便再次进入复连。
引入进程级 CameraTetheringSession 与共享 USB Browser;离开传图页仅取消 UI 回调,App 终止时再完整 shutdown 释放 PTP 与 Browser 资源。

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-07-01 15:45:35 +08:00

423 lines
16 KiB
Swift
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

//
// 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 downloadScope = CameraDownloadStorage.Scope(accountKey: "unscoped", albumID: 0)
private var activeUploadAssetIDs: Set<String> = []
private var autoUploadEligibleAssetIDs: Set<String> = []
private var pendingDeferredNotifyTask: Task<Void, Never>?
private var lastDeferredNotifyDate = Date.distantPast
private var lastNotifiedProgressByAssetID: [String: Int] = [:]
private var inFlightAssetIDs: Set<String> = []
private let maxRetries = 3
private let deferredNotifyInterval: TimeInterval = 0.15
private let maxConcurrentUploads = 3
private var isUIAttached = false
///
init(cameraService: any CameraServiceProtocol) {
self.cameraService = cameraService
}
/// UI OTG attach
func attachToUI(
onConnectionStateChange: @escaping (CameraConnectionState) -> Void,
onTasksUpdated: @escaping ([CameraTransferTask]) -> Void
) {
isUIAttached = true
self.onTasksUpdated = onTasksUpdated
cameraService.onConnectionStateChange = onConnectionStateChange
cameraService.onNewAsset = { [weak self] asset in
Task { @MainActor in
await self?.handleNewAsset(asset, skipIfExists: false)
}
}
}
/// UI
func detachFromUI() {
guard isUIAttached else { return }
isUIAttached = false
pendingDeferredNotifyTask?.cancel()
pendingDeferredNotifyTask = nil
activeUploadAssetIDs.removeAll()
inFlightAssetIDs.removeAll()
onTasksUpdated = nil
cameraService.detachFromUI()
}
/// Sink
func configure(
uploadSink: (any CameraAssetUploadSink)?,
uploadEnabled: Bool,
downloadScope: CameraDownloadStorage.Scope? = nil
) {
self.uploadSink = uploadSink
self.uploadEnabled = uploadEnabled
if let downloadScope {
self.downloadScope = downloadScope
}
}
///
func connect() async {
await cameraService.connect()
if cameraService.connectionState.isConnected {
await syncExistingPhotos()
}
}
/// App teardown
func shutdown() async {
detachFromUI()
await cameraService.shutdown()
}
/// shutdown
func disconnect() async {
await shutdown()
}
/// USB
func restartCameraDiscovery() async {
if let autoDetect = cameraService as? AutoDetectCameraService {
await autoDetect.restartDiscovery()
return
}
await cameraService.shutdown()
await cameraService.connect()
}
///
var connectionState: CameraConnectionState {
cameraService.connectionState
}
///
func retryFailedUploads() async {
guard uploadEnabled, uploadSink != nil else { return }
let failedAssetIDs = tasks.filter { $0.status == .failed }.map(\.assetID)
for index in tasks.indices where tasks[index].status == .failed {
tasks[index].status = .downloaded
tasks[index].errorMessage = nil
}
notify()
await uploadAssets(withIDs: failedAssetIDs)
}
/// asset
func uploadAssets(withIDs assetIDs: [String]) async {
guard uploadSink != nil else { return }
var seenAssetIDs: Set<String> = []
let uniqueAssetIDs = assetIDs.filter { seenAssetIDs.insert($0).inserted }
var currentIndex = 0
while currentIndex < uniqueAssetIDs.count {
let upperBound = min(currentIndex + maxConcurrentUploads, uniqueAssetIDs.count)
let batch = uniqueAssetIDs[currentIndex ..< upperBound]
let uploadTasks = batch.map { assetID in
Task { @MainActor in
await self.uploadAssetIfNeeded(assetID: assetID)
}
}
for task in uploadTasks {
await task.value
}
currentIndex = upperBound
}
}
private func uploadAssetIfNeeded(assetID: String) async {
if activeUploadAssetIDs.contains(assetID) { return }
guard let index = tasks.firstIndex(where: { $0.assetID == assetID }) else { return }
if tasks[index].status == .uploaded { return }
activeUploadAssetIDs.insert(assetID)
defer { activeUploadAssetIDs.remove(assetID) }
if tasks[index].localURL == nil {
await downloadAsset(assetID: assetID)
}
guard let refreshedIndex = tasks.firstIndex(where: { $0.assetID == assetID }) else { return }
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()
autoUploadEligibleAssetIDs.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 }
let isKnownAsset = tasks.contains { $0.assetID == asset.id }
let shouldAutoUpload = !skipIfExists
&& uploadEnabled
&& uploadSink != nil
&& (!isKnownAsset || autoUploadEligibleAssetIDs.contains(asset.id))
inFlightAssetIDs.insert(asset.id)
defer { inFlightAssetIDs.remove(asset.id) }
let assetID = asset.id
if shouldAutoUpload {
autoUploadEligibleAssetIDs.insert(assetID)
}
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 shouldAutoUpload,
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.isInOriginalsDirectory(localURL, scope: downloadScope) {
destination = localURL
} else {
destination = CameraDownloadStorage.uniqueLocalURL(for: resolvedAsset.filename, scope: downloadScope)
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 uploadSink != nil else { return }
while activeUploadAssetIDs.count < maxConcurrentUploads,
let task = tasks.first(where: {
$0.status == .downloaded
&& !activeUploadAssetIDs.contains($0.assetID)
&& autoUploadEligibleAssetIDs.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
autoUploadEligibleAssetIDs.remove(assetID)
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
}