Files
suixinkan_ios_new/suixinkan/Core/CameraTransfer/CameraTransferPipeline.swift
汉秋 19d681192c 修复 OTG 列表删除照片后边拍边传仍会带回已删项的问题。
删除时同步清理管道任务并维护 excludedAssetIDs,避免 notify 将无本地文件的旧任务重新展示为仅文件名的待上传项。

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

442 lines
17 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
/// asset ID
private(set) var excludedAssetIDs: Set<String> = []
///
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 }
}
/// asset ID sync/notify
func setExcludedAssetIDs(_ ids: Set<String>) {
excludedAssetIDs = ids
}
/// asset
func removeAsset(assetID: String) {
guard !assetID.isEmpty else { return }
excludedAssetIDs.insert(assetID)
tasks.removeAll { $0.assetID == assetID }
autoUploadEligibleAssetIDs.remove(assetID)
activeUploadAssetIDs.remove(assetID)
inFlightAssetIDs.remove(assetID)
lastNotifiedProgressByAssetID.removeValue(forKey: assetID)
}
private func handleNewAsset(_ asset: CameraAsset, skipIfExists: Bool) async {
if excludedAssetIDs.contains(asset.id) { return }
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
}