From 3e0142d6230fd684d2f5cab4634185c9005f697d Mon Sep 17 00:00:00 2001 From: Jaasthi Santosh Naidu Lakshman Senna Date: Thu, 10 Sep 2026 02:21:23 -0700 Subject: [PATCH 1/4] Fix playlist import: pin YouTubeKit 0.4.9 (ANDROID_VR stream URLs 403) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Every imported track resolved, then every ranged download got HTTP 403 → IngestError.streamURLExpired → re-resolve → 403 again → scheduleRetry kept the row .pending through 5 whole-track attempts (minutes of spinner) → .failed (orange retry badge). Reproduced on the iOS Simulator with a live probe: 0/8 tracks ready, every attempt "prep failed … streamURLExpired". Root cause: Packages/ContinuityKit/Package.resolved pinned YouTubeKit at 7cc8190 (2026-07-05), whose stream URLs come from the ANDROID_VR InnerTube client. Since mid-August 2026 YouTube serves only the first ~1 MB of those URLs and 403s the rest, so the app's 1 MiB ranged downloader dies on chunk ("Update YouTube Changes (August 2026)") switches to the visionOS/web clients plus an embed fallback; its URLs return 206 for every chunk. - Package.swift: depend on YouTubeKit `exact: "0.4.9"` instead of floating on `branch: "main"` (which resolved differently per machine and let the app silently fall behind). - Package.resolved: 0.4.9 / e5b7d03. - Tests/IngestTests/LiveIngestProbeTests.swift: opt-in live-network probe (skipped unless CONTINUITY_LIVE_PROBE=1) that drives the real PreparationQueue and prints the failing stage/error per track — the tool that found this. With 0.4.9: 8/8 and 6/6 tracks ready on the simulator. - AGENTS.md / CLAUDE.md: record the gotcha and the diagnosis recipe. Co-Authored-By: Claude Fable 5.1 --- AGENTS.md | 11 ++ CLAUDE.md | 10 ++ Packages/ContinuityKit/Package.resolved | 6 +- Packages/ContinuityKit/Package.swift | 2 +- .../IngestTests/LiveIngestProbeTests.swift | 116 ++++++++++++++++++ 5 files changed, 141 insertions(+), 4 deletions(-) create mode 100644 Packages/ContinuityKit/Tests/IngestTests/LiveIngestProbeTests.swift diff --git a/AGENTS.md b/AGENTS.md index 11216ee..f196336 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -104,6 +104,17 @@ or commit it. Bundle id `com.sanylax.continuity` (share extension - **Stem separation on the Simulator** is also CPU-only (and slow). Do not "fix" perceived hangs by enabling CoreML on sim — sim CoreML has no ANE/GPU and routes through a ~100× slower serial CPU queue. +- **YouTubeKit is pinned `exact:` on purpose — bump it deliberately, never float on `main`.** + Stream URLs come from whichever InnerTube client the pinned YouTubeKit asks; YouTube retires + clients without notice (Aug 2026: ANDROID_VR URLs started 403-ing after the first ~1 MB, so + every 1 MiB ranged download died on chunk #2 → `streamURLExpired` on every track → minutes of + spinner, then the orange retry badge). The signature of "YouTube changed again" is *every* + imported track failing with `prep failed for …: streamURLExpired` (or `.network`) in Console; + playlist/search scrapes still succeed. First response: check upstream YouTubeKit for a newer + tag, bump `exact:` in `Packages/ContinuityKit/Package.swift`, and re-run the opt-in probe + `CONTINUITY_LIVE_PROBE=1 … -only-testing:IngestTests/LiveIngestProbeTests` (simulator) — + it prints the failing stage and exact error per track. Beware: a scratch SwiftPM harness + that depends on `branch: "main"` silently resolves upstream HEAD, not the app's pin. - **Scrapers are fragile by design.** YouTube/Spotify change their embedded JSON shapes without notice (YouTube moved playlists to `lockupViewModel` mid-project). Parsers handle multiple shapes and are pinned by tests against real fixtures. Resolvers retry transient failures diff --git a/CLAUDE.md b/CLAUDE.md index a1dd47a..7be9f71 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -94,6 +94,16 @@ transfer or a captive-portal page used to be cached as "the model" forever). - `Player.prepare`/`restore` stay **metadata-only** (no engine build, no `notifyUpcoming()`) — see AGENTS.md jetsam gotcha. +### Playlist import "spins forever, then orange retry" (Sept 2026, resolved) +Every imported track resolved fine, then every ranged download got HTTP 403 → mapped to +`streamURLExpired` → re-resolve → 403 again → `scheduleRetry` kept the row `.pending` through +5 whole-track attempts (minutes of spinner) → `.failed`. Root cause: the pinned YouTubeKit +(7cc8190, July) fetched stream URLs via the ANDROID_VR InnerTube client, which YouTube stopped +serving past the first chunk in mid-August 2026. Fix: pin YouTubeKit `exact: "0.4.9"` +(visionOS/web clients + embed fallback). Verified with the opt-in live probe test on the +simulator: 0/8 tracks ready before, 8/8 after (full files, BPM analysed). See the AGENTS.md +gotcha for the diagnosis recipe. + ### Catalog search (PR #108) iTunes Search API (no key) for songs/albums; custom in-app keyboard with `CatalogAutocorrect` (ContinuityCore, Linux-tested) learning vocabulary from results + the diff --git a/Packages/ContinuityKit/Package.resolved b/Packages/ContinuityKit/Package.resolved index 7294f6a..cd33aa7 100644 --- a/Packages/ContinuityKit/Package.resolved +++ b/Packages/ContinuityKit/Package.resolved @@ -1,5 +1,5 @@ { - "originHash" : "02c652374298720cdaa2e24809bc6bf8dfa33c931c8b07272ab5a1c896df433a", + "originHash" : "9510f39439c4eba8d14bf450fcc3c836fb7203d1498224e004303a913e7ccc15", "pins" : [ { "identity" : "onnxruntime-swift-package-manager", @@ -15,8 +15,8 @@ "kind" : "remoteSourceControl", "location" : "https://github.com/alexeichhorn/YouTubeKit", "state" : { - "branch" : "main", - "revision" : "e5b7d0396ce12bf3444f0d209e8436c83373b7af" + "revision" : "e5b7d0396ce12bf3444f0d209e8436c83373b7af", + "version" : "0.4.9" } } ], diff --git a/Packages/ContinuityKit/Package.swift b/Packages/ContinuityKit/Package.swift index acd1741..4662772 100644 --- a/Packages/ContinuityKit/Package.swift +++ b/Packages/ContinuityKit/Package.swift @@ -16,7 +16,7 @@ let package = Package( ], dependencies: [ .package(path: "../ContinuityCore"), - .package(url: "https://github.com/alexeichhorn/YouTubeKit", branch: "main"), + .package(url: "https://github.com/alexeichhorn/YouTubeKit", exact: "0.4.9"), .package(url: "https://github.com/microsoft/onnxruntime-swift-package-manager", exact: "1.20.0"), ], targets: [ diff --git a/Packages/ContinuityKit/Tests/IngestTests/LiveIngestProbeTests.swift b/Packages/ContinuityKit/Tests/IngestTests/LiveIngestProbeTests.swift new file mode 100644 index 0000000..d80ade2 --- /dev/null +++ b/Packages/ContinuityKit/Tests/IngestTests/LiveIngestProbeTests.swift @@ -0,0 +1,116 @@ +import XCTest +import SwiftData +import OSLog +import Domain +@testable import Ingest + +/// Opt-in, live-network diagnostic for the ingest pipeline. Skipped unless +/// `CONTINUITY_LIVE_PROBE=1` is in the environment, so normal test runs stay hermetic. +/// +/// Drives the real `PreparationQueue` (YouTubeKit resolve → ranged download → analysis) for a +/// handful of well-known videos plus two search-query tracks, then prints every +/// `com.continuity.app` log line the run produced. When YouTube changes something and every +/// imported track starts spinning and then failing, this shows the failing stage and the exact +/// error in one run — no device, no debugger. Run it against an iOS Simulator destination: +/// +/// cd Packages/ContinuityKit && CONTINUITY_LIVE_PROBE=1 \ +/// DEVELOPER_DIR=/Applications/Xcode.app/Contents/Developer xcodebuild test \ +/// -scheme ContinuityKit-Package -destination 'platform=iOS Simulator,id=' \ +/// -only-testing:IngestTests/LiveIngestProbeTests +/// +/// (The env var must reach the test process: xcodebuild forwards `TEST_RUNNER_`-prefixed +/// variables, so `TEST_RUNNER_CONTINUITY_LIVE_PROBE=1` also works.) +/// +/// History: on 2026-09-10 this reproduced "every track spins, then shows the orange retry +/// badge" — every ranged download 403'd (`streamURLExpired`) because the pinned YouTubeKit +/// still asked the ANDROID_VR client for stream URLs, which YouTube stopped serving beyond the +/// first chunk in mid-August 2026. YouTubeKit 0.4.9 (visionOS/web clients) fixed it: 8/8 ready. +@MainActor +final class LiveIngestProbeTests: XCTestCase { + + private static let videoIDs = ["dQw4w9WgXcQ", "9bZkp7q19f0", "kJQP7kiw5Fk", "JGwWNGJdvx8"] + private static let queries = ["Blinding Lights The Weeknd", "bad guy Billie Eilish"] + + override func setUp() async throws { + let env = ProcessInfo.processInfo.environment + try XCTSkipUnless(env["CONTINUITY_LIVE_PROBE"] == "1" || env["TEST_RUNNER_CONTINUITY_LIVE_PROBE"] == "1", + "live-network probe; set CONTINUITY_LIVE_PROBE=1 to run") + } + + func testLiveIngest() async throws { + let start = Date() + let budget: TimeInterval = Double(ProcessInfo.processInfo.environment["CONTINUITY_LIVE_PROBE_SECONDS"] ?? "") ?? 300 + print("PROBE start os=\(ProcessInfo.processInfo.operatingSystemVersionString) budget=\(Int(budget))s") + + // Playlist page scrape (metadata only) — the first stage of a YouTube playlist import. + do { + let resolved = try await YouTubePlaylistResolver().resolvePlaylist(playlistID: "PLFgquLnL59alCl_2TQvOiD5Vgm1hCaGSI") + print("PROBE playlist resolve OK: \(resolved.items.count) items, title=\(resolved.title ?? "-")") + } catch { + print("PROBE playlist resolve FAILED: \(error)") + } + + let schema = Schema([Playlist.self, Track.self]) + let container = try ModelContainer(for: schema, configurations: [ModelConfiguration(isStoredInMemoryOnly: true)]) + let context = container.mainContext + let queue = PreparationQueue() + + let playlist = Playlist(title: "Probe", subtitle: "probe", gradientSeed: 1) + context.insert(playlist) + var tracks: [Track] = [] + for (i, id) in Self.videoIDs.enumerated() { + // Force a real download: evict any cached copy from an earlier run. + for container in ["m4a", "webm", "mp4"] { + try? FileManager.default.removeItem(at: AudioCache.fileURL(videoID: id, container: container)) + } + let track = Track(title: "vid \(id)", artist: "probe", durationSeconds: 0, gradientSeed: i, sortIndex: i, + prepState: .pending, youtubeVideoID: id, + sourceURLString: "https://www.youtube.com/watch?v=\(id)") + playlist.tracks.append(track); context.insert(track); tracks.append(track) + } + for (i, query) in Self.queries.enumerated() { + let track = Track(title: query, artist: "probe", durationSeconds: 0, gradientSeed: 100 + i, sortIndex: 100 + i, + prepState: .pending, searchQuery: query) + playlist.tracks.append(track); context.insert(track); tracks.append(track) + } + try context.save() + for track in tracks { queue.enqueue(track, in: context) } + + var lastStates = "" + while Date().timeIntervalSince(start) < budget { + try await Task.sleep(nanoseconds: 3_000_000_000) + let states = tracks.map { "\($0.youtubeVideoID ?? $0.searchQuery ?? "?")=\($0.prepState.rawValue)" }.joined(separator: " ") + if states != lastStates { + print("PROBE t+\(Int(Date().timeIntervalSince(start)))s \(states)") + lastStates = states + } + if tracks.allSatisfy({ $0.prepState == .ready || $0.prepState == .failed }) { break } + } + + print("PROBE ===== LOG DUMP (com.continuity.app) =====") + do { + let store = try OSLogStore(scope: .currentProcessIdentifier) + let entries = try store.getEntries(at: store.position(date: start)) + var printed = 0 + for case let entry as OSLogEntryLog in entries where entry.subsystem == "com.continuity.app" { + printed += 1 + if printed > 400 { print("PROBE ... (truncated)"); break } + let ts = String(format: "%.1f", entry.date.timeIntervalSince(start)) + print("PROBE LOG t+\(ts)s [\(entry.category)] \(entry.level.rawValue): \(entry.composedMessage)") + } + print("PROBE log entries printed: \(printed)") + } catch { + print("PROBE OSLogStore unavailable: \(error)") + } + + print("PROBE ===== FINAL =====") + for track in tracks { + let path = track.localRelativePath.map { AudioCache.url(forRelativePath: $0).path } ?? "-" + let bytes = (try? FileManager.default.attributesOfItem(atPath: path)[.size] as? Int) ?? 0 + print("PROBE FINAL \(track.youtubeVideoID ?? track.searchQuery ?? "?") state=\(track.prepState.rawValue) dur=\(Int(track.durationSeconds))s bytes=\(bytes) bpm=\(track.bpm ?? 0)") + } + let ready = tracks.filter { $0.prepState == .ready }.count + print("PROBE RESULT ready=\(ready)/\(tracks.count) elapsed=\(Int(Date().timeIntervalSince(start)))s") + XCTAssertEqual(ready, tracks.count, "not every probe track became ready — see PROBE LOG lines above") + } +} From a54be668685c51ee1e4686a6ad21f7036b96db44 Mon Sep 17 00:00:00 2001 From: Jaasthi Santosh Naidu Lakshman Senna Date: Thu, 10 Sep 2026 00:38:47 -0700 Subject: [PATCH 2/4] Fix TestFlight archive: declare a shared Continuity scheme Every TestFlight run since PR #127 died in the Archive step after ~30s with xcodebuild exit 65: "The project named "Continuity" does not contain a scheme named "Continuity"". project.yml declared no schemes, and XcodeGen only emits schemes a target asks for; the runs had been relying on xcodebuild auto-creating one, which Xcode 26.x no longer does (same XcodeGen 2.46.0 + Xcode 26.6 locally reproduces the failure, and the declared scheme makes `xcodebuild archive -scheme Continuity` succeed). Records the scheme gotcha in AGENTS.md / CLAUDE.md so it isn't rediscovered. Co-Authored-By: Claude Fable 5.1 --- AGENTS.md | 5 +++++ CLAUDE.md | 4 +++- project.yml | 4 ++++ 3 files changed, 12 insertions(+), 1 deletion(-) diff --git a/AGENTS.md b/AGENTS.md index f196336..8a25933 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -81,6 +81,11 @@ or commit it. Bundle id `com.sanylax.continuity` (share extension - **New files need `xcodegen generate`** before `xcodebuild`, or you get "cannot find X in scope." XcodeGen uses explicit file lists. +- **Schemes come only from `project.yml`.** XcodeGen emits no scheme unless a target declares + `scheme:` (the Continuity target does — keep it). `xcodebuild` on Xcode 26.x does **not** + auto-create schemes the way the Xcode GUI does, so without it `-scheme Continuity` fails in + ~30s with exit 65 "does not contain a scheme named" — which is exactly how every TestFlight + run from PR #127 to #137 died. - **onnxruntime is a static `.framework` (ar archive), not a dylib.** Xcode still embeds a broken ~50 KB stub into `Continuity.app/Frameworks`. Never "fix" that stub by patching `MinimumOSVersion` and re-signing — that cured ITMS upload checks while leaving a poison diff --git a/CLAUDE.md b/CLAUDE.md index 7be9f71..60aa6b0 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -124,7 +124,9 @@ existing `searchQuery` → YouTube ingest path. automatically after processing. - Known failure modes already hit: App-Manager-role key (fixed — Admin key created); empty/placeholder secrets from copy-pasted commands; **ITMS error 90382 "Upload limit - reached"** = Apple's per-app daily cap — wait for the 24h window, nothing to fix. + reached"** = Apple's per-app daily cap — wait for the 24h window, nothing to fix; + **Archive step exits 65 after ~30s** = no `Continuity` scheme in the generated project + (see the AGENTS.md scheme gotcha) — `project.yml` must keep the target-level `scheme:`. - App Store Connect still has a legacy **Xcode Cloud "Archive – iOS"** workflow producing `action_required` checks on PRs; it's ASC-side, unrelated to code, and competes for upload quota — worth disabling in ASC. diff --git a/project.yml b/project.yml index b501311..c4f791f 100644 --- a/project.yml +++ b/project.yml @@ -66,6 +66,10 @@ targets: - CFBundleURLName: com.sanylax.continuity CFBundleURLSchemes: - continuity + # Shared scheme so `xcodebuild -scheme Continuity` works on a fresh checkout / CI without + # relying on Xcode's per-user scheme auto-creation (which xcodebuild does not perform). + scheme: + gatherCoverageData: false settings: base: PRODUCT_BUNDLE_IDENTIFIER: com.sanylax.continuity From a7d8fe4fca69dc6b33e9db51ddaf03b91285e486 Mon Sep 17 00:00:00 2001 From: sanylaxq Date: Fri, 11 Sep 2026 14:29:12 -0700 Subject: [PATCH 3/4] Add download-first priority and a live Downloads screen Let songs and playlists jump the ingest queue, surface byte-level progress, and restore playback before the launch resume pass so a large library doesn't block the first frame. Co-authored-by: Cursor --- App/Continuity/Views/DownloadsView.swift | 102 ++++++++++++++++++ App/Continuity/Views/LibrarySheetView.swift | 35 ++++++ App/Continuity/Views/LibraryView.swift | 28 +++++ App/Continuity/Views/PlaylistDetailView.swift | 17 +++ App/Continuity/Views/RootView.swift | 5 +- App/Continuity/Views/UpNextView.swift | 12 +++ .../Sources/Ingest/AudioDownloader.swift | 16 ++- .../Sources/Ingest/ConcurrencyLimiter.swift | 64 +++++++++-- .../Sources/Ingest/IngestContracts.swift | 45 +++++++- .../Ingest/PreparationQueue+Jobs.swift | 100 +++++++++++++++++ .../Sources/Ingest/PreparationQueue.swift | 65 ++++++++--- .../IngestTests/ConcurrencyLimiterTests.swift | 72 +++++++++++++ 12 files changed, 536 insertions(+), 25 deletions(-) create mode 100644 App/Continuity/Views/DownloadsView.swift create mode 100644 Packages/ContinuityKit/Sources/Ingest/PreparationQueue+Jobs.swift create mode 100644 Packages/ContinuityKit/Tests/IngestTests/ConcurrencyLimiterTests.swift diff --git a/App/Continuity/Views/DownloadsView.swift b/App/Continuity/Views/DownloadsView.swift new file mode 100644 index 0000000..22ff1ad --- /dev/null +++ b/App/Continuity/Views/DownloadsView.swift @@ -0,0 +1,102 @@ +import SwiftUI +import Domain +import Ingest +import SwiftData + +/// Live ingest queue: every track currently downloading, analysing, or waiting for a slot. +struct DownloadsView: View { + @Environment(PreparationQueue.self) private var prepQueue + @Environment(\.modelContext) private var modelContext + @Environment(\.dismiss) private var dismiss + + var body: some View { + NavigationStack { + Group { + if prepQueue.ingestJobs.isEmpty { + ContentUnavailableView( + "Nothing downloading", + systemImage: "arrow.down.circle", + description: Text("Imported songs show up here until their audio is ready.") + ) + } else { + List { + ForEach(prepQueue.ingestJobs) { job in + jobRow(job) + } + } + .listStyle(.insetGrouped) + } + } + .navigationTitle("Downloads") + .navigationBarTitleDisplayMode(.inline) + .toolbar { + ToolbarItem(placement: .confirmationAction) { + Button("Done") { dismiss() } + } + } + } + } + + private func jobRow(_ job: IngestJob) -> some View { + HStack(spacing: 12) { + VStack(alignment: .leading, spacing: 4) { + HStack(spacing: 6) { + Text(job.title).lineLimit(1) + if job.isPrioritized { + Image(systemName: "arrow.up") + .font(.caption2.weight(.bold)) + .foregroundStyle(.tint) + .accessibilityLabel("Prioritized") + } + } + Text(job.artist) + .font(.caption) + .foregroundStyle(.secondary) + .lineLimit(1) + phaseLabel(job) + .font(.caption2) + .foregroundStyle(.tertiary) + if job.phase == .downloading, let fraction = job.fraction { + ProgressView(value: fraction) + .padding(.top, 2) + } else if job.phase != .queued { + ProgressView() + .padding(.top, 2) + } + } + Spacer(minLength: 8) + if !job.isPrioritized { + Button { + prioritize(job) + } label: { + Image(systemName: "arrow.up.to.line") + } + .buttonStyle(.borderless) + .accessibilityLabel("Download first") + } + } + .padding(.vertical, 4) + } + + private func phaseLabel(_ job: IngestJob) -> Text { + switch job.phase { + case .queued: + return Text("Waiting") + case .downloading: + if let fraction = job.fraction { + return Text("Downloading \(Int((fraction * 100).rounded()))%") + } + return Text("Downloading") + case .analyzing: + return Text("Analyzing") + } + } + + private func prioritize(_ job: IngestJob) { + let jobID = job.id + var descriptor = FetchDescriptor(predicate: #Predicate { $0.id == jobID }) + descriptor.fetchLimit = 1 + guard let track = try? modelContext.fetch(descriptor).first else { return } + prepQueue.prioritize(track, in: modelContext) + } +} diff --git a/App/Continuity/Views/LibrarySheetView.swift b/App/Continuity/Views/LibrarySheetView.swift index 54b988c..6ed28b8 100644 --- a/App/Continuity/Views/LibrarySheetView.swift +++ b/App/Continuity/Views/LibrarySheetView.swift @@ -13,6 +13,7 @@ struct LibrarySheetView: View { @State private var showingSearch = false @State private var showingLocalImport = false @State private var showingAppleMusic = false + @State private var showingDownloads = false /// Non-nil while a picked folder/files are being scanned + copied in. @State private var isImportingLocal = false @@ -22,6 +23,9 @@ struct LibrarySheetView: View { .miniPlayerDock() .navigationTitle("Continuity") .toolbar { + ToolbarItem(placement: .topBarLeading) { + DownloadsToolbarButton(showingDownloads: $showingDownloads) + } // Every action is a primaryAction so nothing collapses into a dead "…" // overflow menu (secondaryAction items did, and looked broken). ToolbarItem(placement: .primaryAction) { @@ -69,6 +73,9 @@ struct LibrarySheetView: View { .sheet(isPresented: $showingAdd) { AddMusicView() } + .sheet(isPresented: $showingDownloads) { + DownloadsView() + } .sheet(isPresented: $showingAppleMusic) { AppleMusicImportView() } @@ -111,3 +118,31 @@ private struct AddBadgeIcon: View { .padding(.trailing, 4) // room for the badge inside the tap target } } + +/// Isolated so byte-level download progress only invalidates this control, not the library grid. +private struct DownloadsToolbarButton: View { + @Environment(PreparationQueue.self) private var prepQueue + @Binding var showingDownloads: Bool + + var body: some View { + let count = prepQueue.ingestJobs.count + return Button { + showingDownloads = true + } label: { + Image(systemName: count == 0 ? "arrow.down.circle" : "arrow.down.circle.fill") + } + .accessibilityLabel("Downloads") + .overlay(alignment: .topTrailing) { + if count > 0 { + Text("\(count)") + .font(.system(size: 9, weight: .bold)) + .padding(.horizontal, 4) + .padding(.vertical, 1) + .background(.tint, in: Capsule()) + .foregroundStyle(.white) + .offset(x: 8, y: -8) + .accessibilityHidden(true) + } + } + } +} diff --git a/App/Continuity/Views/LibraryView.swift b/App/Continuity/Views/LibraryView.swift index db27072..500905b 100644 --- a/App/Continuity/Views/LibraryView.swift +++ b/App/Continuity/Views/LibraryView.swift @@ -10,6 +10,7 @@ struct LibraryView: View { @Query(sort: \Playlist.createdAt) private var playlists: [Playlist] @Environment(Player.self) private var player @Environment(\.modelContext) private var modelContext + @Environment(PreparationQueue.self) private var prepQueue @State private var searchText = "" /// Playlist awaiting destructive confirmation — set from the context menu, cleared on dismiss. @State private var playlistPendingDelete: Playlist? @@ -69,6 +70,13 @@ struct LibraryView: View { } .buttonStyle(.plain) .contextMenu { + if playlist.tracks.contains(where: { !$0.isDemo && $0.prepState != .ready }) { + Button { + prepQueue.prioritize(playlist: playlist, in: modelContext) + } label: { + Label("Download First", systemImage: "arrow.up.to.line") + } + } Button(role: .destructive) { playlistPendingDelete = playlist } label: { @@ -101,6 +109,8 @@ private struct SearchResultsView: View { let playlists: [Playlist] let query: String @Environment(Player.self) private var player + @Environment(PreparationQueue.self) private var prepQueue + @Environment(\.modelContext) private var modelContext @Environment(MainPagerState.self) private var pagerState private var matchingPlaylists: [Playlist] { @@ -148,6 +158,15 @@ private struct SearchResultsView: View { } } } + .contextMenu { + if playlist.tracks.contains(where: { !$0.isDemo && $0.prepState != .ready }) { + Button { + prepQueue.prioritize(playlist: playlist, in: modelContext) + } label: { + Label("Download First", systemImage: "arrow.up.to.line") + } + } + } } } } @@ -184,6 +203,8 @@ private struct SearchSongRow: View { let track: Track let play: () -> Void @Environment(Player.self) private var player + @Environment(PreparationQueue.self) private var prepQueue + @Environment(\.modelContext) private var modelContext var body: some View { Button(action: play) { @@ -211,6 +232,13 @@ private struct SearchSongRow: View { } label: { Label("Play Next", systemImage: "text.line.first.and.arrowtriangle.forward") } + if !track.isDemo, track.prepState != .ready { + Button { + prepQueue.prioritize(track, in: modelContext) + } label: { + Label("Download First", systemImage: "arrow.up.to.line") + } + } } } } diff --git a/App/Continuity/Views/PlaylistDetailView.swift b/App/Continuity/Views/PlaylistDetailView.swift index f6b396d..45961e2 100644 --- a/App/Continuity/Views/PlaylistDetailView.swift +++ b/App/Continuity/Views/PlaylistDetailView.swift @@ -44,6 +44,13 @@ struct PlaylistDetailView: View { } label: { Label("Play Next", systemImage: "text.line.first.and.arrowtriangle.forward") } + if !track.isDemo, track.prepState != .ready { + Button { + prepQueue.prioritize(track, in: modelContext) + } label: { + Label("Download First", systemImage: "arrow.up.to.line") + } + } } .swipeActions(edge: .trailing) { Button(role: .destructive) { @@ -90,6 +97,16 @@ struct PlaylistDetailView: View { .buttonStyle(.glassProminent) .padding(.top, 4) + if tracks.contains(where: { !$0.isDemo && $0.prepState != .ready }) { + Button { + prepQueue.prioritize(playlist: playlist, in: modelContext) + } label: { + Label("Download First", systemImage: "arrow.up.to.line") + .frame(maxWidth: 200) + } + .buttonStyle(.bordered) + } + // Source-backed playlists mirror a remote list: manual sync + the auto-sync opt-out. if playlist.isSourceBacked { HStack(spacing: 16) { diff --git a/App/Continuity/Views/RootView.swift b/App/Continuity/Views/RootView.swift index 2073650..458bfe2 100644 --- a/App/Continuity/Views/RootView.swift +++ b/App/Continuity/Views/RootView.swift @@ -101,8 +101,11 @@ struct RootView: View { return ids.compactMap { byID[$0] } } LibraryCleanup.sweepOrphanedFiles(in: modelContext) - prepQueue.resumePreparation(in: modelContext) + // Restore the last song before walking the library for resume/sync — otherwise + // a large unfinished import occupies the main actor until the first frame. restorePlaybackSession() + await Task.yield() + await prepQueue.resumePreparation(in: modelContext) // Launch-time polling pass over source-backed playlists (per-playlist opt-out). prepQueue.autoSyncIfNeeded(in: modelContext) } diff --git a/App/Continuity/Views/UpNextView.swift b/App/Continuity/Views/UpNextView.swift index e8ec49a..0c3429d 100644 --- a/App/Continuity/Views/UpNextView.swift +++ b/App/Continuity/Views/UpNextView.swift @@ -1,6 +1,7 @@ import SwiftUI import Playback import Domain +import Ingest import ContinuityCore /// The queue page (below Now Playing): what plays next, with drag-to-reorder and @@ -8,6 +9,8 @@ import ContinuityCore /// key/tempo-compatible DJ sequence. struct UpNextView: View { @Environment(Player.self) private var player + @Environment(PreparationQueue.self) private var prepQueue + @Environment(\.modelContext) private var modelContext @Environment(MainPagerState.self) private var pagerState // Persisted as a mode label; toggling ON reorders once, toggling OFF is not an undo. @AppStorage("flowMode.v1") private var flowMode = false @@ -81,6 +84,15 @@ struct UpNextView: View { Text(track.artist).font(.caption).foregroundStyle(.secondary).lineLimit(1) } } + .contextMenu { + if !track.isDemo, track.prepState != .ready { + Button { + prepQueue.prioritize(track, in: modelContext) + } label: { + Label("Download First", systemImage: "arrow.up.to.line") + } + } + } } /// Reorders only the upcoming tracks. The current track is passed as the chain's anchor — diff --git a/Packages/ContinuityKit/Sources/Ingest/AudioDownloader.swift b/Packages/ContinuityKit/Sources/Ingest/AudioDownloader.swift index 8369c3c..f9a3d65 100644 --- a/Packages/ContinuityKit/Sources/Ingest/AudioDownloader.swift +++ b/Packages/ContinuityKit/Sources/Ingest/AudioDownloader.swift @@ -21,11 +21,15 @@ final class AudioDownloader: AudioFileDownloading { self.maxRetriesPerChunk = maxRetriesPerChunk } - func downloadAudio(_ resolved: ResolvedAudio) async throws -> URL { + func downloadAudio( + _ resolved: ResolvedAudio, + progress: (@Sendable (Int, Int?) -> Void)? + ) async throws -> URL { let destination = AudioCache.fileURL(videoID: resolved.videoID, container: resolved.container) // Cache hit: already on disk. if FileManager.default.fileExists(atPath: destination.path) { + progress?(1, 1) return destination } @@ -34,7 +38,7 @@ final class AudioDownloader: AudioFileDownloading { .appendingPathComponent("continuity-\(resolved.videoID)-\(UUID().uuidString).\(resolved.container)") defer { try? FileManager.default.removeItem(at: tempURL) } - try await downloadRanged(from: resolved.url, to: tempURL) + try await downloadRanged(from: resolved.url, to: tempURL, progress: progress) // Publish. Move (rename) is atomic on the same volume; if a concurrent download already // won the race and the file now exists, treat that as success rather than corrupting it. @@ -42,6 +46,7 @@ final class AudioDownloader: AudioFileDownloading { try FileManager.default.moveItem(at: tempURL, to: destination) } catch { if FileManager.default.fileExists(atPath: destination.path) { + progress?(1, 1) return destination } throw IngestError.downloadFailed(String(describing: error)) @@ -50,7 +55,11 @@ final class AudioDownloader: AudioFileDownloading { } /// Streams `url` into `fileURL` using sequential `Range` requests until the whole file is fetched. - private func downloadRanged(from url: URL, to fileURL: URL) async throws { + private func downloadRanged( + from url: URL, + to fileURL: URL, + progress: (@Sendable (Int, Int?) -> Void)? + ) async throws { FileManager.default.createFile(atPath: fileURL.path, contents: nil) let handle = try FileHandle(forWritingTo: fileURL) do { @@ -70,6 +79,7 @@ final class AudioDownloader: AudioFileDownloading { } try handle.write(contentsOf: data) offset += data.count + progress?(offset, totalSize) } while totalSize == nil || offset < totalSize! try handle.close() diff --git a/Packages/ContinuityKit/Sources/Ingest/ConcurrencyLimiter.swift b/Packages/ContinuityKit/Sources/Ingest/ConcurrencyLimiter.swift index 1ebcef5..40b2e5e 100644 --- a/Packages/ContinuityKit/Sources/Ingest/ConcurrencyLimiter.swift +++ b/Packages/ContinuityKit/Sources/Ingest/ConcurrencyLimiter.swift @@ -1,15 +1,27 @@ import Foundation -/// A minimal async concurrency gate: at most `limit` holders run between `acquire()` and -/// `release()`; the rest suspend (FIFO) until a slot frees. +/// A minimal async concurrency gate: at most `limit` holders run between `acquire` and +/// `release()`; the rest suspend until a slot frees. +/// +/// Waiters are FIFO among equal `priority` values. A later `bump` reorders still-waiting +/// acquirers so a user-prioritized track jumps the ingest queue without cancelling anyone +/// already downloading. /// /// Used by `PreparationQueue` so importing a large playlist doesn't fire dozens of simultaneous /// resolves/downloads (network throttling) or stem separations (each loads a ~158 MB model and /// pegs the CPU — running many at once would thrash memory). actor ConcurrencyLimiter { + private struct Waiter { + let id: UUID + var priority: Int + let sequence: Int + let continuation: CheckedContinuation + } + private let limit: Int private var active = 0 - private var waiters: [CheckedContinuation] = [] + private var waiters: [Waiter] = [] + private var nextSequence = 0 init(limit: Int) { self.limit = max(1, limit) @@ -17,22 +29,60 @@ actor ConcurrencyLimiter { /// Suspends until a slot is available, then claims it. Pair with exactly one `release()`. func acquire() async { + await acquire(id: UUID(), priority: 0) + } + + /// Identified acquire so a later `bump` can move this waiter ahead of equal/lower priority. + func acquire(id: UUID, priority: Int) async { if active < limit { active += 1 return } - await withCheckedContinuation { waiters.append($0) } + await withCheckedContinuation { continuation in + waiters.append(Waiter( + id: id, + priority: priority, + sequence: nextSequence, + continuation: continuation + )) + nextSequence += 1 + sortWaiters() + } // Resumed by `release()`, which hands over its slot without touching `active`. } - /// Frees a slot, waking the longest-waiting acquirer (if any). + /// Raises still-waiting acquirers in `ids` to at least `priority`. No-op for holders + /// already running — those finish; the next slot goes to the new head of the queue. + func bump(ids: Set, to priority: Int) { + guard !ids.isEmpty else { return } + var changed = false + for i in waiters.indices where ids.contains(waiters[i].id) { + if waiters[i].priority < priority { + waiters[i].priority = priority + changed = true + } + } + if changed { sortWaiters() } + } + + /// Frees a slot, waking the highest-priority waiter (FIFO among ties). func release() { if waiters.isEmpty { active = max(0, active - 1) } else { - // Transfer the slot directly to the next waiter; `active` stays the same. let next = waiters.removeFirst() - next.resume() + next.continuation.resume() + } + } + + /// Waiter count — tests use this instead of `Task.yield` to know an `acquire` has queued. + var pendingCount: Int { waiters.count } + + /// Higher priority first; among equals, the earlier `acquire` stays ahead. + private func sortWaiters() { + waiters.sort { a, b in + if a.priority != b.priority { return a.priority > b.priority } + return a.sequence < b.sequence } } } diff --git a/Packages/ContinuityKit/Sources/Ingest/IngestContracts.swift b/Packages/ContinuityKit/Sources/Ingest/IngestContracts.swift index 0a56b41..1b14216 100644 --- a/Packages/ContinuityKit/Sources/Ingest/IngestContracts.swift +++ b/Packages/ContinuityKit/Sources/Ingest/IngestContracts.swift @@ -128,8 +128,51 @@ protocol VideoMetadataResolving: Sendable { func metadata(videoID: String) async throws -> VideoMetadata } +/// One track currently in the ingest pipeline. Surfaced by `PreparationQueue.ingestJobs` so +/// the Downloads screen can show queued vs in-flight work without walking every SwiftData row. +public struct IngestJob: Identifiable, Equatable, Sendable { + public enum Phase: String, Sendable { + case queued, downloading, analyzing + } + + public let id: UUID + public var title: String + public var artist: String + public var phase: Phase + /// 0...1 while `phase == .downloading` and the server reported a size; nil otherwise. + public var fraction: Double? + public var isPrioritized: Bool + + public init( + id: UUID, + title: String, + artist: String, + phase: Phase, + fraction: Double?, + isPrioritized: Bool + ) { + self.id = id + self.title = title + self.artist = artist + self.phase = phase + self.fraction = fraction + self.isPrioritized = isPrioritized + } +} + /// Downloads a resolved stream to local storage and returns the on-disk file URL. /// Implemented by `AudioDownloader`. protocol AudioFileDownloading: Sendable { - func downloadAudio(_ resolved: ResolvedAudio) async throws -> URL + /// `progress` reports `(bytesWritten, totalBytes?)` as ranged chunks land. `totalBytes` is + /// nil until the first `Content-Range` (or a whole-file 200) arrives. + func downloadAudio( + _ resolved: ResolvedAudio, + progress: (@Sendable (Int, Int?) -> Void)? + ) async throws -> URL +} + +extension AudioFileDownloading { + func downloadAudio(_ resolved: ResolvedAudio) async throws -> URL { + try await downloadAudio(resolved, progress: nil) + } } diff --git a/Packages/ContinuityKit/Sources/Ingest/PreparationQueue+Jobs.swift b/Packages/ContinuityKit/Sources/Ingest/PreparationQueue+Jobs.swift new file mode 100644 index 0000000..d527de9 --- /dev/null +++ b/Packages/ContinuityKit/Sources/Ingest/PreparationQueue+Jobs.swift @@ -0,0 +1,100 @@ +import Domain +import Foundation +import SwiftData + +extension PreparationQueue { + static let songPriority = 100 + static let playlistPriority = 50 + + /// Jump `track` to the head of the ingest queue. Failed rows are re-enqueued. + /// Already-ready tracks are a no-op — there's nothing to download. + public func prioritize(_ track: Track, in context: ModelContext) { + guard !track.isDemo, track.prepState != .ready else { return } + ingestPriority[track.id] = max(ingestPriority[track.id] ?? 0, Self.songPriority) + if let idx = ingestJobs.firstIndex(where: { $0.id == track.id }) { + ingestJobs[idx].isPrioritized = true + sortJobs() + } + let id = track.id + Task { await ingestLimiter.bump(ids: [id], to: Self.songPriority) } + if track.prepState == .failed { + enqueue(track, in: context) + } + } + + /// Jump every not-yet-ready track in a playlist/album ahead of the rest of the library. + public func prioritize(playlist: Playlist, in context: ModelContext) { + let tracks = playlist.tracks.filter { !$0.isDemo && $0.prepState != .ready } + guard !tracks.isEmpty else { return } + let ids = Set(tracks.map(\.id)) + for track in tracks { + ingestPriority[track.id] = max(ingestPriority[track.id] ?? 0, Self.playlistPriority) + } + for i in ingestJobs.indices where ids.contains(ingestJobs[i].id) { + ingestJobs[i].isPrioritized = true + } + sortJobs() + Task { await ingestLimiter.bump(ids: ids, to: Self.playlistPriority) } + var enqueued = false + for track in tracks where track.prepState == .failed { + enqueue(track, in: context, saving: false) + enqueued = true + } + if enqueued { try? context.save() } + } + + func upsertJob(_ track: Track, phase: IngestJob.Phase) { + let prioritized = (ingestPriority[track.id] ?? 0) > 0 + if let idx = ingestJobs.firstIndex(where: { $0.id == track.id }) { + ingestJobs[idx].title = track.title + ingestJobs[idx].artist = track.artist + ingestJobs[idx].phase = phase + ingestJobs[idx].isPrioritized = prioritized + if phase != .downloading { ingestJobs[idx].fraction = nil } + } else { + ingestJobs.append(IngestJob( + id: track.id, + title: track.title, + artist: track.artist, + phase: phase, + fraction: nil, + isPrioritized: prioritized + )) + } + sortJobs() + } + + func updateJobProgress(_ id: UUID, bytes: Int, total: Int?) { + guard let idx = ingestJobs.firstIndex(where: { $0.id == id }) else { return } + ingestJobs[idx].phase = .downloading + guard let total, total > 0 else { return } + let fraction = min(1, Double(bytes) / Double(total)) + // Skip sub-percent redraws — ranged chunks are 1 MiB, so this still updates often. + if let current = ingestJobs[idx].fraction, abs(current - fraction) < 0.01 { return } + ingestJobs[idx].fraction = fraction + } + + func removeJob(_ id: UUID) { + ingestJobs.removeAll { $0.id == id } + } + + /// Active downloads first, then analysis, then the waiting queue. Prioritized rows float up + /// within a phase so "Download first" is visible at the top of the screen. + private func sortJobs() { + ingestJobs.sort { a, b in + if a.isPrioritized != b.isPrioritized { return a.isPrioritized && !b.isPrioritized } + let pa = Self.phaseOrder(a.phase) + let pb = Self.phaseOrder(b.phase) + if pa != pb { return pa < pb } + return a.title.localizedCaseInsensitiveCompare(b.title) == .orderedAscending + } + } + + private static func phaseOrder(_ phase: IngestJob.Phase) -> Int { + switch phase { + case .downloading: return 0 + case .analyzing: return 1 + case .queued: return 2 + } + } +} diff --git a/Packages/ContinuityKit/Sources/Ingest/PreparationQueue.swift b/Packages/ContinuityKit/Sources/Ingest/PreparationQueue.swift index 7b91ba3..e281758 100644 --- a/Packages/ContinuityKit/Sources/Ingest/PreparationQueue.swift +++ b/Packages/ContinuityKit/Sources/Ingest/PreparationQueue.swift @@ -43,10 +43,15 @@ public final class PreparationQueue { let downloader: AudioFileDownloading /// Caps simultaneous resolve+download+analyse work (network-bound). - private let ingestLimiter = ConcurrencyLimiter(limit: 3) + let ingestLimiter = ConcurrencyLimiter(limit: 3) /// Caps simultaneous stem separations to one — each is CPU/RAM-heavy, so they queue. let stemLimiter = ConcurrencyLimiter(limit: 1) + /// In-flight ingest work for the Downloads screen. Not persisted; rebuilt by `enqueue`. + public private(set) var ingestJobs: [IngestJob] = [] + /// User-raised ingest priority per track. 100 = one song, 50 = whole playlist/album. + var ingestPriority: [UUID: Int] = [:] + /// Production wiring — the app constructs the queue with no arguments. The parameterized /// initializer stays internal for dependency-injected tests within the module. public convenience init() { self.init(resolver: YouTubeStreamResolver()) } @@ -90,6 +95,7 @@ public final class PreparationQueue { /// keeps failing keeps climbing its curve instead of looping at `base` forever. func enqueueInternal(_ track: Track, in context: ModelContext, saving: Bool = true) { track.prepState = .pending + upsertJob(track, phase: .queued) if saving { try? context.save() } Task { await process(track, in: context) } } @@ -116,7 +122,10 @@ public final class PreparationQueue { /// interrupted mid-ingest (e.g. the app was killed partway through a large import) or whose /// downloaded audio went missing, and finishes stem separation for tracks that have audio but /// no stems yet. `.failed` tracks are left as-is for an explicit retry. - public func resumePreparation(in context: ModelContext) { + /// + /// Yields every 40 rows so a thousand-track library doesn't occupy the main actor for the + /// whole pass before the first frame can land. + public func resumePreparation(in context: ModelContext) async { guard let tracks = try? context.fetch(FetchDescriptor()) else { return } // One directory listing per cache instead of up to five `fileExists` probes per track: // a thousand-track library meant thousands of stat calls on the main thread, at launch, @@ -124,7 +133,7 @@ public final class PreparationQueue { let cacheIndex = CacheIndex.snapshot() // Demo healing used to `save()` once per track; batched into a single save at the end. var needsSave = false - for track in tracks { + for (i, track) in tracks.enumerated() { // Demo tracks have no source and play synthesized audio — there is nothing to ingest // or resume. Without this guard they'd be re-enqueued (they have no audio file), fail // for lack of a source, and show up as retry-able failures. Heal any that already did. @@ -153,6 +162,7 @@ public final class PreparationQueue { case .failed: break } + if i.isMultiple(of: 40) { await Task.yield() } } if needsSave { try? context.save() } } @@ -188,19 +198,33 @@ public final class PreparationQueue { /// at each stage. Any failure (missing video ID, resolve, or download error) lands the /// track in `.failed`; the UI surfaces that as a retry-able badge rather than a crash. private func process(_ track: Track, in context: ModelContext) async { - track.prepState = .preparing - try? context.save() - retryScheduledTrackIDs.remove(track.id) + let trackID = track.id + retryScheduledTrackIDs.remove(trackID) + upsertJob(track, phase: .queued) // A track needs either a direct video ID (YouTube) or a search query (Spotify-sourced). guard track.youtubeVideoID != nil || track.searchQuery != nil else { track.prepState = .failed + removeJob(trackID) try? context.save() return } - // Gate the network/CPU-heavy stage so a playlist import doesn't run all tracks at once. - await ingestLimiter.acquire() + // Stay `.pending` until a limiter slot is actually ours — otherwise a 200-track import + // looks like 200 simultaneous downloads. `bump` can reorder this waiter meanwhile. + await ingestLimiter.acquire(id: trackID, priority: ingestPriority[trackID] ?? 0) + // Deleted while queued. + guard track.modelContext != nil else { + ingestAttempts[trackID] = nil + removeJob(trackID) + await ingestLimiter.release() + return + } + + track.prepState = .preparing + try? context.save() + upsertJob(track, phase: .downloading) + var prepared = false var failure: Error? do { @@ -218,14 +242,22 @@ public final class PreparationQueue { let resolved = try await resolver.resolveAudio(videoID: id) let fileURL: URL do { - fileURL = try await downloader.downloadAudio(resolved) + fileURL = try await downloader.downloadAudio(resolved, progress: { [weak self] done, total in + Task { @MainActor in + self?.updateJobProgress(trackID, bytes: done, total: total) + } + }) } catch let error as IngestError where error.needsFreshStreamURL { // Signed `googlevideo` URLs are short-lived, and a throttled client gets them // invalidated early — so a queued track's URL can be dead by the time its turn // comes. Retrying the dead URL can never work; re-resolve for a fresh one. Logger.ingest.notice("stream URL expired for \(id, privacy: .public) — re-resolving") let refreshed = try await resolver.resolveAudio(videoID: id) - fileURL = try await downloader.downloadAudio(refreshed) + fileURL = try await downloader.downloadAudio(refreshed, progress: { [weak self] done, total in + Task { @MainActor in + self?.updateJobProgress(trackID, bytes: done, total: total) + } + }) } // The track could have been deleted while we were off the main actor; don't write to // (or resurrect) a dead model. @@ -243,6 +275,7 @@ public final class PreparationQueue { // Analyse tempo + key off the main actor (full-track FFTs). Non-fatal: if analysis // fails the track still plays, just without BPM/key metadata. + upsertJob(track, phase: .analyzing) if let analysis = try? await Task.detached(priority: .utility, operation: { try TrackAnalyzer.analyze(fileURL: fileURL) }).value, track.modelContext != nil { @@ -269,21 +302,27 @@ public final class PreparationQueue { // Don't touch a track that was deleted while we worked. guard track.modelContext != nil else { - ingestAttempts[track.id] = nil + ingestAttempts[trackID] = nil + removeJob(trackID) try? context.save() return } if prepared { - ingestAttempts[track.id] = nil + ingestAttempts[trackID] = nil + ingestPriority[trackID] = nil track.prepState = .ready + removeJob(trackID) } else if let failure, scheduleRetry(track, after: failure, in: context) { // Stays `.pending`, not `.failed`: the row keeps its in-progress badge, and if the // app is killed before the retry fires, `resumePreparation` picks it up at launch. track.prepState = .pending + upsertJob(track, phase: .queued) } else { - ingestAttempts[track.id] = nil + ingestAttempts[trackID] = nil + ingestPriority[trackID] = nil track.prepState = .failed + removeJob(trackID) } try? context.save() diff --git a/Packages/ContinuityKit/Tests/IngestTests/ConcurrencyLimiterTests.swift b/Packages/ContinuityKit/Tests/IngestTests/ConcurrencyLimiterTests.swift new file mode 100644 index 0000000..6b6d5ad --- /dev/null +++ b/Packages/ContinuityKit/Tests/IngestTests/ConcurrencyLimiterTests.swift @@ -0,0 +1,72 @@ +import XCTest +@testable import Ingest + +final class ConcurrencyLimiterTests: XCTestCase { + + /// A later waiter with a higher priority runs before an earlier one after `release`. + func testBumpJumpsTheQueue() async { + let limiter = ConcurrencyLimiter(limit: 1) + let firstID = UUID() + let secondID = UUID() + var order: [UUID] = [] + + await limiter.acquire(id: UUID(), priority: 0) + + let first = Task { + await limiter.acquire(id: firstID, priority: 0) + order.append(firstID) + await limiter.release() + } + await waitUntil(limiter, pending: 1) + + let second = Task { + await limiter.acquire(id: secondID, priority: 0) + order.append(secondID) + await limiter.release() + } + await waitUntil(limiter, pending: 2) + + await limiter.bump(ids: [secondID], to: 100) + await limiter.release() + + _ = await first.result + _ = await second.result + XCTAssertEqual(order, [secondID, firstID]) + } + + func testEqualPriorityStaysFIFO() async { + let limiter = ConcurrencyLimiter(limit: 1) + let firstID = UUID() + let secondID = UUID() + var order: [UUID] = [] + + await limiter.acquire() + + let first = Task { + await limiter.acquire(id: firstID, priority: 0) + order.append(firstID) + await limiter.release() + } + await waitUntil(limiter, pending: 1) + + let second = Task { + await limiter.acquire(id: secondID, priority: 0) + order.append(secondID) + await limiter.release() + } + await waitUntil(limiter, pending: 2) + + await limiter.release() + _ = await first.result + _ = await second.result + XCTAssertEqual(order, [firstID, secondID]) + } + + private func waitUntil(_ limiter: ConcurrencyLimiter, pending: Int) async { + for _ in 0..<10_000 { + if await limiter.pendingCount >= pending { return } + await Task.yield() + } + XCTFail("timed out waiting for \(pending) limiter waiters") + } +} From f8e3af560d4a540df36ca8ad189edf4187709e2c Mon Sep 17 00:00:00 2001 From: sanylaxq Date: Fri, 11 Sep 2026 14:39:00 -0700 Subject: [PATCH 4/4] Fix ingestJobs mutation from PreparationQueue+Jobs private(set) is file-scoped, so the Jobs extension could not update the Downloads list and the device build failed. Co-authored-by: Cursor --- Packages/ContinuityKit/Sources/Ingest/PreparationQueue.swift | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/Packages/ContinuityKit/Sources/Ingest/PreparationQueue.swift b/Packages/ContinuityKit/Sources/Ingest/PreparationQueue.swift index e281758..4f865a1 100644 --- a/Packages/ContinuityKit/Sources/Ingest/PreparationQueue.swift +++ b/Packages/ContinuityKit/Sources/Ingest/PreparationQueue.swift @@ -48,7 +48,8 @@ public final class PreparationQueue { let stemLimiter = ConcurrencyLimiter(limit: 1) /// In-flight ingest work for the Downloads screen. Not persisted; rebuilt by `enqueue`. - public private(set) var ingestJobs: [IngestJob] = [] + /// `internal(set)` so `PreparationQueue+Jobs` can mutate it from another file in Ingest. + public internal(set) var ingestJobs: [IngestJob] = [] /// User-raised ingest priority per track. 100 = one song, 50 = whole playlist/album. var ingestPriority: [UUID: Int] = [:]