Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 9 additions & 2 deletions App/Continuity/Views/RootView.swift
Original file line number Diff line number Diff line change
Expand Up @@ -79,7 +79,11 @@ struct RootView: View {
// Natural queue exhaustion loops playback into the listening history — resolve
// the persisted IDs to live tracks in order; deleted tracks simply drop out.
player.onQueueExhausted = { ids in
let tracks = (try? modelContext.fetch(FetchDescriptor<Track>())) ?? []
// ID-map resolution only: don't eagerly hydrate every track's beatTimes
// (hundreds of doubles each) — untouched properties fault in on demand.
var descriptor = FetchDescriptor<Track>()
descriptor.propertiesToFetch = [\.id]
let tracks = (try? modelContext.fetch(descriptor)) ?? []
let byID = Dictionary(uniqueKeysWithValues: tracks.map { ($0.id, $0) })
return ids.compactMap { byID[$0] }
}
Expand Down Expand Up @@ -227,7 +231,10 @@ struct RootView: View {
/// fresh install, stages COMË N GO paused at the start of its playlist.
private func restorePlaybackSession() {
guard player.currentTrack == nil else { return } // already playing (e.g. state restore re-entry)
let tracks = (try? modelContext.fetch(FetchDescriptor<Track>())) ?? []
// ID-map resolution: skip eager beatTimes hydration (see onQueueExhausted).
var descriptor = FetchDescriptor<Track>()
descriptor.propertiesToFetch = [\.id]
let tracks = (try? modelContext.fetch(descriptor)) ?? []

if let state = PlaybackStateStore.load() {
let byID = Dictionary(uniqueKeysWithValues: tracks.map { ($0.id, $0) })
Expand Down
21 changes: 15 additions & 6 deletions Packages/ContinuityCore/Sources/ContinuityCore/OverlapAdd.swift
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,10 @@ public final class StreamingOverlapAdd {

/// First frame index still held in internal state (everything before it has been drained).
private var base = 0
/// Physical index of frame `base` within `acc`/`weight`. Drains advance this cursor and
/// compact only once a segment of dead frames accumulates — removeFirst on every drain
/// memmoved the whole remaining buffer per window.
private var head = 0
/// Per-channel weighted accumulation, index i ↔ absolute frame base + i.
private var acc: [[Float]]
/// Sum of window weights per frame, shared across channels.
Expand All @@ -43,12 +47,12 @@ public final class StreamingOverlapAdd {
precondition(length >= 0 && length <= segment)
let needed = (start - base) + length
if needed > count {
let grow = needed - count
let grow = (head + needed) - weight.count
for c in 0..<channels { acc[c].append(contentsOf: repeatElement(0, count: grow)) }
weight.append(contentsOf: repeatElement(0, count: grow))
count = needed
}
let offset = start - base
let offset = head + (start - base)
for i in 0..<length {
let w = window[i]
for c in 0..<channels { acc[c][offset + i] += sample(c, i) * w }
Expand All @@ -65,14 +69,19 @@ public final class StreamingOverlapAdd {
let n = upTo - base
var out = Array(repeating: [Float](repeating: 0, count: n), count: channels)
for i in 0..<n {
let w = weight[i]
let w = weight[head + i]
guard w > 0 else { continue }
for c in 0..<channels { out[c][i] = acc[c][i] / w }
for c in 0..<channels { out[c][i] = acc[c][head + i] / w }
}
for c in 0..<channels { acc[c].removeFirst(n) }
weight.removeFirst(n)
head += n
base = upTo
count -= n
// Amortized compaction: shed dead frames once a segment's worth piled up.
if head >= segment {
for c in 0..<channels { acc[c].removeFirst(head) }
weight.removeFirst(head)
head = 0
}
return out
}

Expand Down
12 changes: 10 additions & 2 deletions Packages/ContinuityKit/Sources/Ingest/LibraryCleanup.swift
Original file line number Diff line number Diff line change
Expand Up @@ -11,9 +11,17 @@ import SwiftData
/// post-deletion library.
public enum LibraryCleanup {

/// Fetch limited to the fields `stemKey` reads — a cleanup pass has no business eagerly
/// hydrating every track's beat grid.
private static func stemKeyDescriptor() -> FetchDescriptor<Track> {
var descriptor = FetchDescriptor<Track>()
descriptor.propertiesToFetch = [\.youtubeVideoID, \.id]
return descriptor
}

@MainActor
public static func removeOrphanedFiles(keys: [String], in context: ModelContext) {
let tracks = (try? context.fetch(FetchDescriptor<Track>())) ?? []
let tracks = (try? context.fetch(Self.stemKeyDescriptor())) ?? []
let referenced = Set(tracks.map(\.stemKey))
for key in Set(keys) where !referenced.contains(key) {
removeCachedFiles(key: key)
Expand All @@ -25,7 +33,7 @@ public enum LibraryCleanup {
/// on disk after the delete-time cleanup had already run.
@MainActor
public static func sweepOrphanedFiles(in context: ModelContext) {
guard let tracks = try? context.fetch(FetchDescriptor<Track>()) else { return }
guard let tracks = try? context.fetch(Self.stemKeyDescriptor()) else { return }
let referenced = Set(tracks.map(\.stemKey))

if let files = try? FileManager.default.contentsOfDirectory(
Expand Down
30 changes: 27 additions & 3 deletions Packages/ContinuityKit/Sources/Ingest/PreparationQueue+Stems.swift
Original file line number Diff line number Diff line change
Expand Up @@ -33,9 +33,30 @@ extension PreparationQueue {
}
// Budget pass on every queue move (not just post-separation): catches overage that
// accrued outside this code path — e.g. gigabytes of legacy float32 stems.
scheduleBudgetPass()
}

/// Runs the cache budget pass single-flight: while one full directory scan is in
/// progress, further requests coalesce into (at most) one follow-up pass, which reads
/// the protected set fresh when it starts — so the newest play-queue neighborhood is
/// what gets protected, and rapid skips can't stack N concurrent scans.
private func scheduleBudgetPass() {
guard !budgetPassRunning else {
budgetPassQueued = true
return
}
budgetPassRunning = true
let protected = protectedStemKeys
Task.detached(priority: .utility) {
Task.detached(priority: .utility) { [weak self] in
StemCache.enforceBudget(protecting: protected)
await MainActor.run {
guard let self else { return }
self.budgetPassRunning = false
if self.budgetPassQueued {
self.budgetPassQueued = false
self.scheduleBudgetPass()
}
}
}
}

Expand Down Expand Up @@ -98,7 +119,10 @@ extension PreparationQueue {
let accompanimentOut = StemCache.accompanimentURL(key: key)
let limiter = stemLimiter

Task.detached(priority: .utility) { [weak self] in
// weak track: the separation runs for CPU-minutes; holding the @Model (and with it
// its context graph) strongly for that long keeps deleted models alive. All accesses
// already guard on modelContext, so weak just lets a deleted track deallocate.
Task.detached(priority: .utility) { [weak self, weak track] in
await limiter.acquire()
// Re-check after the (possibly long) wait for the slot: another queued separation
// of the same video, or a re-added track, may have written the stems meanwhile —
Expand Down Expand Up @@ -136,7 +160,7 @@ extension PreparationQueue {
let separator = OnnxStemSeparator(modelURL: modelURL)
_ = try separator.separate(inputURL: inputURL, vocalsOut: vocalsOut, accompanimentOut: accompanimentOut)
await MainActor.run {
guard track.modelContext != nil else { return }
guard let track, track.modelContext != nil else { return }
track.vocalsRelativePath = StemCache.relativePath(for: vocalsOut)
track.accompanimentRelativePath = StemCache.relativePath(for: accompanimentOut)
try? context.save()
Expand Down
28 changes: 23 additions & 5 deletions Packages/ContinuityKit/Sources/Ingest/PreparationQueue.swift
Original file line number Diff line number Diff line change
Expand Up @@ -174,8 +174,11 @@ public final class PreparationQueue {
track.localRelativePath = AudioCache.relativePath(for: fileURL)

// Real duration for the row (bare-ID adds start at 0, which renders as "0:00").
if track.durationSeconds <= 0, let file = try? AVAudioFile(forReading: fileURL) {
track.durationSeconds = Double(file.length) / file.processingFormat.sampleRate
// Header parse runs off-main: AVAudioFile(forReading:) opens + parses the
// container synchronously, and this class is MainActor-bound.
if track.durationSeconds <= 0 {
let seconds = await Self.fileDurationSeconds(url: fileURL)
if let seconds, track.modelContext != nil { track.durationSeconds = seconds }
}

// NOTE: the real title/channel (oEmbed) is deliberately NOT fetched here — it
Expand Down Expand Up @@ -223,6 +226,14 @@ public final class PreparationQueue {
}
}

/// Opens the audio file off the main actor and returns its duration, or nil if unreadable.
nonisolated static func fileDurationSeconds(url: URL) async -> Double? {
await Task.detached(priority: .utility) {
guard let file = try? AVAudioFile(forReading: url) else { return nil }
return Double(file.length) / file.processingFormat.sampleRate
}.value
}

/// Whether the track still shows the "YouTube Video (abc123)" placeholder from a bare-ID add.
private static func hasPlaceholderMetadata(_ track: Track) -> Bool {
track.youtubeVideoID != nil && track.title.hasPrefix("YouTube Video (")
Expand Down Expand Up @@ -253,9 +264,11 @@ public final class PreparationQueue {
if needsDuration || needsSilenceScan, let relativePath = track.localRelativePath {
let url = AudioCache.url(forRelativePath: relativePath)
await ingestLimiter.acquire()
if needsDuration, track.modelContext != nil,
let file = try? AVAudioFile(forReading: url) {
track.durationSeconds = Double(file.length) / file.processingFormat.sampleRate
if needsDuration {
// Off-main: a launch backfill over a big library would otherwise do one
// synchronous main-thread file open per track — UI hitches during playback.
let seconds = await Self.fileDurationSeconds(url: url)
if let seconds, track.modelContext != nil { track.durationSeconds = seconds }
}
if needsSilenceScan, track.modelContext != nil {
MemoryFootprint.breadcrumb("silence scan begin")
Expand Down Expand Up @@ -320,6 +333,11 @@ public final class PreparationQueue {
var separationAllowedAt: Date?
/// At most one pending post-hold retry (ensureStems also re-fires on every track change).
var stemsRetryScheduled = false
/// Single-flight state for the stem-cache budget pass: rapid skips fired one full
/// directory scan per queue move, racing each other. One runs at a time; at most one
/// more is queued (re-reading the latest protected set when it starts).
var budgetPassRunning = false
var budgetPassQueued = false


}
24 changes: 17 additions & 7 deletions Packages/ContinuityKit/Sources/Ingest/StemSeparator.swift
Original file line number Diff line number Diff line change
Expand Up @@ -198,9 +198,14 @@ final class OnnxStemSeparator: StemSeparating {
let ola = StreamingOverlapAdd(channels: channels, segment: segment, overlap: segment / 4)
let stride = ola.stride

// Un-flushed mix frames [flushed, flushed + mix count) — needed both as model input and
// to derive accompaniment = mix − vocals after normalization.
// Un-flushed mix frames — needed both as model input and to derive
// accompaniment = mix − vocals after normalization. `mixBase` is the absolute frame
// index of mixL[0]: flushes advance `flushed` and only compact the arrays once a
// whole segment of dead frames has built up, instead of memmoving the remainder on
// every window (removeFirst per flush was ~1.4 MB of copy traffic per window, running
// concurrently with playback).
var mixL = [Float](), mixR = [Float]()
var mixBase = 0
var flushed = 0
var decodedEnd = 0
var atEOF = false
Expand Down Expand Up @@ -229,7 +234,7 @@ final class OnnxStemSeparator: StemSeparating {

// Build the input tensor: shape (1, 2, segment), planar [ch0…, ch1…], zero-padded.
let length = min(segment, decodedEnd - start)
let offset = start - flushed
let offset = start - mixBase
for i in 0..<length {
input[i] = mixL[offset + i]
input[segment + i] = mixR[offset + i]
Expand Down Expand Up @@ -257,15 +262,20 @@ final class OnnxStemSeparator: StemSeparating {
let n = finalUpTo - flushed
var accL = [Float](repeating: 0, count: n)
var accR = [Float](repeating: 0, count: n)
let head = flushed - mixBase
for i in 0..<n {
accL[i] = mixL[i] - voc[0][i]
accR[i] = mixR[i] - voc[1][i]
accL[i] = mixL[head + i] - voc[0][i]
accR[i] = mixR[head + i] - voc[1][i]
}
try append(to: vocalsFile, left: voc[0], right: voc[1])
try append(to: accFile, left: accL, right: accR)
mixL.removeFirst(n)
mixR.removeFirst(n)
flushed = finalUpTo
// Amortized compaction: drop dead frames only once a segment's worth piled up.
if flushed - mixBase >= segment {
mixL.removeFirst(flushed - mixBase)
mixR.removeFirst(flushed - mixBase)
mixBase = flushed
}
}
if atEOF && start >= decodedEnd { break }
}
Expand Down
Loading