| 1 | import Foundation |
| 2 | import AppKit |
| 3 | import CryptoKit |
| 4 | |
| 5 | struct MediaStatus { |
| 6 | var probing = false |
| 7 | var filmstripReady = false |
| 8 | var proxyReady = false |
| 9 | var proxyProgress: Double = 0 // 0..1 while generating |
| 10 | var failed: String? = nil |
| 11 | } |
| 12 | |
| 13 | /// ffprobe/ffmpeg-based derived-media pipeline: probe, filmstrip, ProRes |
| 14 | /// proxy. Cache is content-addressed and LRU-capped so NAS media gets a |
| 15 | /// bounded local working set. |
| 16 | final class MediaPipeline { |
| 17 | static let shared = MediaPipeline() |
| 18 | |
| 19 | let cacheRoot: URL |
| 20 | /// LRU cap in bytes (default 50 GB). Override in Settings, or `defaults |
| 21 | /// write net.paperclover.Sequencer maxCacheGB -int 100`. |
| 22 | var maxCacheBytes: Int64 { |
| 23 | let gb = UserDefaults.standard.integer(forKey: "maxCacheGB") |
| 24 | return Int64(gb > 0 ? gb : 50) * 1_000_000_000 |
| 25 | } |
| 26 | |
| 27 | private let ffmpeg: String? |
| 28 | private let ffprobe: String? |
| 29 | private let workQueue = OperationQueue() |
| 30 | private var statuses: [UUID: MediaStatus] = [:] // main-thread only |
| 31 | private let thumbCache = NSCache<NSString, CGImage>() |
| 32 | private var stripInfoCache: [String: (interval: Double, count: Int)] = [:] |
| 33 | private var lruTouched: [String: Date] = [:] |
| 34 | |
| 35 | init() { |
| 36 | if let custom = UserDefaults.standard.string(forKey: "cacheDir") { |
| 37 | cacheRoot = URL(fileURLWithPath: (custom as NSString).expandingTildeInPath) |
| 38 | } else { |
| 39 | cacheRoot = FileManager.default.urls(for: .cachesDirectory, in: .userDomainMask)[0] |
| 40 | .appendingPathComponent("Sequencer", isDirectory: true) |
| 41 | } |
| 42 | try? FileManager.default.createDirectory(at: cacheRoot, withIntermediateDirectories: true) |
| 43 | ffmpeg = Self.findExecutable("ffmpeg") |
| 44 | ffprobe = Self.findExecutable("ffprobe") |
| 45 | workQueue.maxConcurrentOperationCount = 2 |
| 46 | thumbCache.countLimit = 2000 |
| 47 | // A previous instance killed mid-build (rebuild relaunch, force quit) |
| 48 | // leaves orphaned ffmpeg encoders holding shared VideoToolbox decode |
| 49 | // sessions — enough of them and every AVPlayer here renders black. |
| 50 | Self.reapOrphans(cacheRoot: cacheRoot) |
| 51 | // Measure the cache once at launch; chunk builds wait on `budgetReady`. |
| 52 | DispatchQueue.main.async { self.reconcileLedger() } |
| 53 | } |
| 54 | |
| 55 | static func findExecutable(_ name: String) -> String? { |
| 56 | var candidates = (ProcessInfo.processInfo.environment["PATH"] ?? "") |
| 57 | .split(separator: ":").map(String.init) |
| 58 | candidates += ["/opt/homebrew/bin", "/usr/local/bin", "/run/current-system/sw/bin", |
| 59 | "\(NSHomeDirectory())/.nix-profile/bin", |
| 60 | "/etc/profiles/per-user/\(NSUserName())/bin"] |
| 61 | for dir in candidates { |
| 62 | let p = "\(dir)/\(name)" |
| 63 | if FileManager.default.isExecutableFile(atPath: p) { return p } |
| 64 | } |
| 65 | return nil |
| 66 | } |
| 67 | |
| 68 | func status(for media: MediaItem) -> MediaStatus { |
| 69 | if let s = statuses[media.id] { return s } |
| 70 | var s = MediaStatus() |
| 71 | s.filmstripReady = FileManager.default.fileExists(atPath: stripInfoURL(media.cacheKey).path) |
| 72 | s.proxyReady = FileManager.default.fileExists(atPath: proxyFileURL(media.cacheKey).path) |
| 73 | statuses[media.id] = s |
| 74 | return s |
| 75 | } |
| 76 | |
| 77 | private var offlineCache: [String: (offline: Bool, until: Date)] = [:] // main-thread only |
| 78 | /// Whether the media's ORIGINAL file is currently unreachable (moved, or on |
| 79 | /// an unmounted NAS). Cached briefly so the viewer can call it every frame |
| 80 | /// without a `stat` each time, and so it auto-recovers when the drive returns. |
| 81 | func isOffline(_ media: MediaItem) -> Bool { |
| 82 | let now = Date() |
| 83 | if let c = offlineCache[media.cacheKey], c.until > now { return c.offline } |
| 84 | let off = !FileManager.default.fileExists(atPath: media.path) |
| 85 | offlineCache[media.cacheKey] = (off, now.addingTimeInterval(2)) |
| 86 | return off |
| 87 | } |
| 88 | |
| 89 | // MARK: - Cache paths |
| 90 | |
| 91 | /// The shape `cacheKey(for:)` produces: exactly 16 lowercase hex chars. A |
| 92 | /// key from another tool (or a hand-edited `.sq`) that doesn't match is not |
| 93 | /// trusted as a directory name. |
| 94 | static func isValidCacheKey(_ key: String) -> Bool { |
| 95 | key.count == 16 && key.allSatisfy { $0.isHexDigit && !$0.isUppercase } |
| 96 | } |
| 97 | |
| 98 | /// Content-hash a seed string into a valid 16-hex cache key. |
| 99 | static func hashedKey(_ seed: String) -> String { |
| 100 | let digest = SHA256.hash(data: Data(seed.utf8)) |
| 101 | return digest.map { String(format: "%02x", $0) }.joined().prefix(16).lowercased() |
| 102 | } |
| 103 | |
| 104 | /// Defensive backstop: never let a blank or malformed key resolve to |
| 105 | /// `cacheRoot` itself or escape it via `/` or `..`. A garbage key is folded |
| 106 | /// to a stable hashed stand-in so its derived assets stay contained. |
| 107 | private func sanitizedKey(_ key: String) -> String { |
| 108 | Self.isValidCacheKey(key) ? key : Self.hashedKey("invalid|\(key)") |
| 109 | } |
| 110 | |
| 111 | private func keyDir(_ key: String) -> URL { |
| 112 | cacheRoot.appendingPathComponent(sanitizedKey(key), isDirectory: true) |
| 113 | } |
| 114 | private func proxyFileURL(_ key: String) -> URL { keyDir(key).appendingPathComponent("proxy.mov") } |
| 115 | private func stripDir(_ key: String) -> URL { keyDir(key).appendingPathComponent("strip", isDirectory: true) } |
| 116 | private func stripInfoURL(_ key: String) -> URL { stripDir(key).appendingPathComponent("info.json") } |
| 117 | |
| 118 | static func cacheKey(for url: URL) -> String { |
| 119 | let attrs = try? FileManager.default.attributesOfItem(atPath: url.path) |
| 120 | let size = (attrs?[.size] as? NSNumber)?.int64Value ?? 0 |
| 121 | let mtime = (attrs?[.modificationDate] as? Date)?.timeIntervalSince1970 ?? 0 |
| 122 | return Self.hashedKey("\(url.path)|\(size)|\(Int(mtime))") |
| 123 | } |
| 124 | |
| 125 | /// A trustworthy cache key for a media item loaded from disk: keep a valid |
| 126 | /// one, otherwise recompute from the file (its content hash), or — when the |
| 127 | /// file is missing — fall back to a stable hash of its path so it still |
| 128 | /// can't collide with keyless siblings or escape the cache root. |
| 129 | static func normalizedCacheKey(for media: MediaItem) -> String { |
| 130 | if isValidCacheKey(media.cacheKey) { return media.cacheKey } |
| 131 | if FileManager.default.fileExists(atPath: media.path) { |
| 132 | return cacheKey(for: media.url) |
| 133 | } |
| 134 | return Self.hashedKey("path|\(media.path)") |
| 135 | } |
| 136 | |
| 137 | /// Proxy URL if the proxy exists (touches LRU). |
| 138 | func proxyURL(for media: MediaItem) -> URL? { |
| 139 | let url = proxyFileURL(media.cacheKey) |
| 140 | guard FileManager.default.fileExists(atPath: url.path) else { return nil } |
| 141 | touchLRU(media.cacheKey) |
| 142 | return url |
| 143 | } |
| 144 | |
| 145 | // MARK: - Import |
| 146 | |
| 147 | /// Probe a file and kick off background filmstrip + proxy generation. |
| 148 | func importFile(_ url: URL, completion: @escaping (MediaItem?) -> Void) { |
| 149 | guard let ffprobe else { |
| 150 | DispatchQueue.main.async { |
| 151 | NotificationCenter.default.post(name: .transientStatus, object: nil, |
| 152 | userInfo: ["text": "ffprobe not found — install ffmpeg"]) |
| 153 | completion(nil) |
| 154 | } |
| 155 | return |
| 156 | } |
| 157 | DispatchQueue.global(qos: .userInitiated).async { |
| 158 | let out = Self.run(ffprobe, ["-v", "quiet", "-print_format", "json", |
| 159 | "-show_format", "-show_streams", url.path]).stdout |
| 160 | guard let data = out.data(using: .utf8), |
| 161 | let json = try? JSONSerialization.jsonObject(with: data) as? [String: Any], |
| 162 | let format = json["format"] as? [String: Any], |
| 163 | let streams = json["streams"] as? [[String: Any]], |
| 164 | let duration = Double((format["duration"] as? String) ?? "") |
| 165 | else { |
| 166 | DispatchQueue.main.async { completion(nil) } |
| 167 | return |
| 168 | } |
| 169 | var item = MediaItem(path: url.path) |
| 170 | item.duration = duration |
| 171 | item.cacheKey = Self.cacheKey(for: url) |
| 172 | item.hasAudio = streams.contains { ($0["codec_type"] as? String) == "audio" } |
| 173 | item.isAudio = item.hasAudio && !streams.contains { |
| 174 | ($0["codec_type"] as? String) == "video" |
| 175 | // Album art shows up as a video stream; ignore it. |
| 176 | && ($0["disposition"] as? [String: Any])?["attached_pic"] as? Int != 1 |
| 177 | } |
| 178 | if let v = streams.first(where: { ($0["codec_type"] as? String) == "video" }) { |
| 179 | item.width = v["width"] as? Int ?? 0 |
| 180 | item.height = v["height"] as? Int ?? 0 |
| 181 | if let r = v["r_frame_rate"] as? String { |
| 182 | let parts = r.split(separator: "/").compactMap { Double($0) } |
| 183 | if parts.count == 2, parts[1] > 0 { |
| 184 | item.fps = min(120, max(1, parts[0] / parts[1])) |
| 185 | } |
| 186 | } |
| 187 | } |
| 188 | DispatchQueue.main.async { |
| 189 | completion(item) |
| 190 | self.enqueueDerivedAssets(for: item) |
| 191 | } |
| 192 | } |
| 193 | } |
| 194 | |
| 195 | /// In-memory duration cache for the drop preview (path → seconds). |
| 196 | private var durationCache: [String: Double] = [:] |
| 197 | |
| 198 | /// Cheap duration-only probe for the drag-and-drop landing preview. Runs |
| 199 | /// ffprobe off-main and memoizes by path; the completion fires on main. |
| 200 | /// Directories / sync.json manifests report nil. |
| 201 | func probeDuration(_ url: URL, completion: @escaping (Double?) -> Void) { |
| 202 | if let cached = durationCache[url.path] { completion(cached); return } |
| 203 | guard let ffprobe, |
| 204 | (try? url.resourceValues(forKeys: [.isDirectoryKey]).isDirectory) != true else { |
| 205 | completion(nil); return |
| 206 | } |
| 207 | DispatchQueue.global(qos: .userInitiated).async { |
| 208 | let out = Self.run(ffprobe, ["-v", "quiet", "-show_entries", "format=duration", |
| 209 | "-of", "csv=p=0", url.path]).stdout |
| 210 | let d = Double(out.trimmingCharacters(in: .whitespacesAndNewlines)) |
| 211 | DispatchQueue.main.async { |
| 212 | if let d, d > 0 { self.durationCache[url.path] = d } |
| 213 | completion(d) |
| 214 | } |
| 215 | } |
| 216 | } |
| 217 | |
| 218 | /// Ensure filmstrip/proxy jobs exist for every media in the project |
| 219 | /// (e.g. after opening a project on a machine with a cold cache). |
| 220 | func ensureDerivedAssets(for project: ProjectModel) { |
| 221 | for m in project.media { enqueueDerivedAssets(for: m) } |
| 222 | } |
| 223 | |
| 224 | private var enqueued: Set<String> = [] |
| 225 | |
| 226 | func enqueueDerivedAssets(for media: MediaItem) { |
| 227 | guard ffmpeg != nil, !enqueued.contains(media.cacheKey) else { return } |
| 228 | enqueued.insert(media.cacheKey) |
| 229 | if media.isAudio { |
| 230 | if !FileManager.default.fileExists(atPath: waveformURL(media.cacheKey).path) { |
| 231 | workQueue.addOperation { self.generateWaveform(media) } |
| 232 | } |
| 233 | return |
| 234 | } |
| 235 | let s = status(for: media) |
| 236 | if !s.filmstripReady { workQueue.addOperation { self.generateFilmstrip(media) } } |
| 237 | // Proxies are chunked and demand-driven — see ChunkManager. Whole-file |
| 238 | // proxy.mov caches from earlier versions keep being used when present. |
| 239 | } |
| 240 | |
| 241 | // MARK: - Waveforms (audio clips) |
| 242 | |
| 243 | private func waveformURL(_ key: String) -> URL { |
| 244 | keyDir(key).appendingPathComponent("waveform.png") |
| 245 | } |
| 246 | |
| 247 | private func generateWaveform(_ media: MediaItem) { |
| 248 | guard let ffmpeg else { return } |
| 249 | try? FileManager.default.createDirectory(at: keyDir(media.cacheKey), |
| 250 | withIntermediateDirectories: true) |
| 251 | let res = Self.run(ffmpeg, [ |
| 252 | "-y", "-i", media.path, |
| 253 | "-filter_complex", |
| 254 | "aformat=channel_layouts=mono,showwavespic=s=2048x200:colors=white", |
| 255 | "-frames:v", "1", waveformURL(media.cacheKey).path, |
| 256 | ]) |
| 257 | let bytes = (try? FileManager.default.attributesOfItem( |
| 258 | atPath: waveformURL(media.cacheKey).path)[.size] as? NSNumber)?.int64Value ?? 0 |
| 259 | DispatchQueue.main.async { |
| 260 | self.touchLRU(media.cacheKey) |
| 261 | if res.exitCode == 0 { |
| 262 | self.noteBytesAdded(bytes) |
| 263 | self.imgStateLock.lock() |
| 264 | self.waveformMissing.remove(media.cacheKey) |
| 265 | self.imgStateLock.unlock() |
| 266 | // `scene: true` — a waveform image changes timeline pixels |
| 267 | // (the tile cache flushes on it; plain chunk churn must not). |
| 268 | NotificationCenter.default.post(name: .mediaStatusChanged, object: nil, |
| 269 | userInfo: ["scene": true]) |
| 270 | } |
| 271 | } |
| 272 | } |
| 273 | |
| 274 | private let waveformCache = NSCache<NSString, CGImage>() |
| 275 | |
| 276 | /// Decode an image file straight into the display's raster format (BGRA |
| 277 | /// premultiplied, screen colorspace). A plain NSImage/CGImageSource image |
| 278 | /// keeps the file's own format (RGB JPEG, generic colorspace), and Quartz |
| 279 | /// then converts it on EVERY draw — per thumbnail, per frame. Converting |
| 280 | /// once at load makes the timeline's image blits plain memory copies. |
| 281 | static func displayImage(contentsOf url: URL) -> CGImage? { |
| 282 | guard let src = CGImageSourceCreateWithURL(url as CFURL, nil), |
| 283 | let raw = CGImageSourceCreateImageAtIndex( |
| 284 | src, 0, [kCGImageSourceShouldCache: false] as CFDictionary) |
| 285 | else { return nil } |
| 286 | let space = NSScreen.main?.colorSpace?.cgColorSpace |
| 287 | ?? CGColorSpace(name: CGColorSpace.sRGB)! |
| 288 | guard let ctx = CGContext( |
| 289 | data: nil, width: raw.width, height: raw.height, |
| 290 | bitsPerComponent: 8, bytesPerRow: 0, space: space, |
| 291 | bitmapInfo: CGImageAlphaInfo.premultipliedFirst.rawValue |
| 292 | | CGBitmapInfo.byteOrder32Little.rawValue) else { return raw } |
| 293 | ctx.draw(raw, in: CGRect(x: 0, y: 0, width: raw.width, height: raw.height)) |
| 294 | return ctx.makeImage() ?? raw |
| 295 | } |
| 296 | |
| 297 | /// Same negative/in-flight bookkeeping as thumbnails: a missing waveform |
| 298 | /// must not re-dispatch a load per audio clip per frame. |
| 299 | private var waveformMissing = Set<String>() |
| 300 | private var waveformLoading = Set<String>() |
| 301 | |
| 302 | func waveformImage(for media: MediaItem) -> CGImage? { |
| 303 | let key = media.cacheKey |
| 304 | if let img = waveformCache.object(forKey: key as NSString) { return img } |
| 305 | imgStateLock.lock() |
| 306 | let skip = waveformMissing.contains(key) || waveformLoading.contains(key) |
| 307 | if !skip { waveformLoading.insert(key) } |
| 308 | imgStateLock.unlock() |
| 309 | guard !skip else { return nil } |
| 310 | let url = waveformURL(key) |
| 311 | DispatchQueue.global(qos: .utility).async { |
| 312 | let img = Self.displayImage(contentsOf: url) |
| 313 | DispatchQueue.main.async { |
| 314 | self.imgStateLock.lock() |
| 315 | self.waveformLoading.remove(key) |
| 316 | if img == nil { self.waveformMissing.insert(key) } |
| 317 | self.imgStateLock.unlock() |
| 318 | if let img { |
| 319 | self.waveformCache.setObject(img, forKey: key as NSString) |
| 320 | self.notifyThumbsCoalesced() |
| 321 | } |
| 322 | } |
| 323 | } |
| 324 | return nil |
| 325 | } |
| 326 | |
| 327 | // MARK: - Jobs |
| 328 | |
| 329 | private func generateFilmstrip(_ media: MediaItem) { |
| 330 | guard let ffmpeg else { return } |
| 331 | let dir = stripDir(media.cacheKey) |
| 332 | try? FileManager.default.createDirectory(at: dir, withIntermediateDirectories: true) |
| 333 | // Target ~600 thumbs max so hour-long recordings stay cheap. |
| 334 | let interval = max(0.5, media.duration / 600) |
| 335 | let res = Self.run(ffmpeg, ["-y", "-hwaccel", "videotoolbox", "-i", media.path, |
| 336 | "-vf", "fps=1/\(interval),scale=240:-2", |
| 337 | "-q:v", "7", dir.appendingPathComponent("%06d.jpg").path]) |
| 338 | let count = (try? FileManager.default.contentsOfDirectory(atPath: dir.path))? |
| 339 | .filter { $0.hasSuffix(".jpg") }.count ?? 0 |
| 340 | if res.exitCode == 0, count > 0 { |
| 341 | let info: [String: Any] = ["interval": interval, "count": count] |
| 342 | if let d = try? JSONSerialization.data(withJSONObject: info) { |
| 343 | try? d.write(to: stripInfoURL(media.cacheKey)) |
| 344 | } |
| 345 | } |
| 346 | let bytes = Self.directorySize(dir) |
| 347 | DispatchQueue.main.async { |
| 348 | var s = self.status(for: media) |
| 349 | s.filmstripReady = res.exitCode == 0 && count > 0 |
| 350 | self.statuses[media.id] = s |
| 351 | self.touchLRU(media.cacheKey) |
| 352 | self.noteBytesAdded(bytes) |
| 353 | self.imgStateLock.lock() |
| 354 | self.thumbMissing.removeAll() // a fresh strip supersedes misses |
| 355 | self.imgStateLock.unlock() |
| 356 | NotificationCenter.default.post(name: .mediaStatusChanged, object: nil, |
| 357 | userInfo: ["scene": true]) |
| 358 | } |
| 359 | } |
| 360 | |
| 361 | // MARK: - Filmstrip access |
| 362 | |
| 363 | func filmstripInfo(_ key: String) -> (interval: Double, count: Int)? { |
| 364 | imgStateLock.lock() |
| 365 | let cached = stripInfoCache[key] |
| 366 | imgStateLock.unlock() |
| 367 | if let cached { return cached } |
| 368 | guard let data = try? Data(contentsOf: stripInfoURL(key)), |
| 369 | let json = try? JSONSerialization.jsonObject(with: data) as? [String: Any], |
| 370 | let interval = json["interval"] as? Double, |
| 371 | let count = json["count"] as? Int else { return nil } |
| 372 | imgStateLock.lock() |
| 373 | stripInfoCache[key] = (interval, count) |
| 374 | imgStateLock.unlock() |
| 375 | return (interval, count) |
| 376 | } |
| 377 | |
| 378 | /// Cached thumbnail nearest to `seconds`; loads async and posts a single |
| 379 | /// coalesced .mediaStatusChanged when thumbs land (a post per thumb |
| 380 | /// cascades into an app-wide refresh storm while a filmstrip streams in). |
| 381 | /// Thumbs that failed to load (file absent/evicted) or are mid-load. |
| 382 | /// Without these a dense timeline re-dispatches a load per missing thumb |
| 383 | /// per FRAME — hundreds of no-op queue hops every draw. Misses are |
| 384 | /// forgotten whenever fresh thumbs land (strips may have regenerated). |
| 385 | /// Lock-guarded: the timeline rasterizes lanes on parallel threads. |
| 386 | private var thumbMissing = Set<String>() |
| 387 | private var thumbLoading = Set<String>() |
| 388 | private let imgStateLock = NSLock() |
| 389 | |
| 390 | func filmstripImage(for media: MediaItem, at seconds: Double) -> CGImage? { |
| 391 | guard let info = filmstripInfo(media.cacheKey) else { return nil } |
| 392 | let index = min(info.count, max(1, Int(seconds / info.interval) + 1)) |
| 393 | let cacheId = "\(media.cacheKey)/\(index)" |
| 394 | if let img = thumbCache.object(forKey: cacheId as NSString) { return img } |
| 395 | imgStateLock.lock() |
| 396 | let skip = thumbMissing.contains(cacheId) || thumbLoading.contains(cacheId) |
| 397 | if !skip { thumbLoading.insert(cacheId) } |
| 398 | imgStateLock.unlock() |
| 399 | guard !skip else { return nil } |
| 400 | let url = stripDir(media.cacheKey).appendingPathComponent(String(format: "%06d.jpg", index)) |
| 401 | DispatchQueue.global(qos: .utility).async { |
| 402 | let img = Self.displayImage(contentsOf: url) |
| 403 | DispatchQueue.main.async { |
| 404 | self.imgStateLock.lock() |
| 405 | self.thumbLoading.remove(cacheId) |
| 406 | if img == nil { self.thumbMissing.insert(cacheId) } |
| 407 | self.imgStateLock.unlock() |
| 408 | if let img { |
| 409 | self.thumbCache.setObject(img, forKey: cacheId as NSString) |
| 410 | self.notifyThumbsCoalesced() |
| 411 | } |
| 412 | } |
| 413 | } |
| 414 | return nil |
| 415 | } |
| 416 | |
| 417 | private var thumbNotifyPending = false |
| 418 | private func notifyThumbsCoalesced() { |
| 419 | guard !thumbNotifyPending else { return } |
| 420 | thumbNotifyPending = true |
| 421 | DispatchQueue.main.asyncAfter(deadline: .now() + 0.1) { |
| 422 | self.thumbNotifyPending = false |
| 423 | self.imgStateLock.lock() |
| 424 | self.thumbMissing.removeAll() // strips may have (re)generated |
| 425 | self.imgStateLock.unlock() |
| 426 | NotificationCenter.default.post(name: .mediaStatusChanged, object: nil, |
| 427 | userInfo: ["scene": true]) |
| 428 | } |
| 429 | } |
| 430 | |
| 431 | // MARK: - LRU eviction |
| 432 | |
| 433 | private func touchLRU(_ key: String) { |
| 434 | let now = Date() |
| 435 | if let last = lruTouched[key], now.timeIntervalSince(last) < 60 { return } |
| 436 | lruTouched[key] = now |
| 437 | let url = keyDir(key).appendingPathComponent("lastUsed") |
| 438 | try? "\(now.timeIntervalSince1970)".write(to: url, atomically: true, encoding: .utf8) |
| 439 | } |
| 440 | |
| 441 | private var evicting = false |
| 442 | |
| 443 | // MARK: - Cache budget ledger |
| 444 | |
| 445 | /// Approximate cache size in bytes, maintained incrementally (builds add, |
| 446 | /// evictions subtract) and reconciled against a real walk at launch and on |
| 447 | /// document open. With the ledger, staying under the cap is enforced BEFORE |
| 448 | /// bytes hit the disk (`tryReserve`), and the per-build full-cache walk the |
| 449 | /// old eviction did is gone. Main-thread only. |
| 450 | private(set) var ledgerBytes: Int64 = 0 |
| 451 | /// Bytes promised to in-flight chunk builds; released on completion. |
| 452 | private var reservedBytes: Int64 = 0 |
| 453 | /// False until the first reconcile walk finishes. Admission stays closed |
| 454 | /// while false, so an unmeasured cache can never be built past the cap. |
| 455 | private(set) var budgetReady = false |
| 456 | private var reconciling = false |
| 457 | |
| 458 | /// When eviction runs at all it frees down to here (not just under the |
| 459 | /// cap), so it works in batches instead of one chunk per build. |
| 460 | var lowWatermarkBytes: Int64 { Int64(Double(maxCacheBytes) * 0.92) } |
| 461 | |
| 462 | /// Re-measure the cache and set the ledger to ground truth. Called at |
| 463 | /// launch, on document open, and when the cap changes; incremental updates |
| 464 | /// keep it honest in between. |
| 465 | /// Anything on disk with an mtime before this was left by a previous |
| 466 | /// session — the test for purging session-transient files (rescue slices, |
| 467 | /// crashed builds' .partial temps) without touching live ones. |
| 468 | private static let processStart = Date() |
| 469 | |
| 470 | func reconcileLedger() { |
| 471 | guard !reconciling else { return } |
| 472 | reconciling = true |
| 473 | let root = cacheRoot |
| 474 | DispatchQueue.global(qos: .utility).async { |
| 475 | Self.purgeSessionTransients(root: root) |
| 476 | let total = Self.directorySize(root) |
| 477 | DispatchQueue.main.async { |
| 478 | NSLog("[cache] reconcile: %.2f GB on disk (cap %.0f GB)", |
| 479 | Double(total) / 1e9, Double(self.maxCacheBytes) / 1e9) |
| 480 | self.ledgerBytes = total |
| 481 | self.budgetReady = true |
| 482 | self.reconciling = false |
| 483 | if self.ledgerBytes + self.reservedBytes > self.maxCacheBytes { |
| 484 | self.evictToFit(need: 0) { _ in } |
| 485 | } |
| 486 | for c in DocumentContext.allLive { c.chunks.budgetChanged() } |
| 487 | } |
| 488 | } |
| 489 | } |
| 490 | |
| 491 | /// Reserve room for a build about to start. The reservation counts against |
| 492 | /// the ceiling alongside bytes already on disk, so two concurrent builds |
| 493 | /// can't both squeeze into the same headroom. Window builds reserve up to |
| 494 | /// the cap; background fill only up to the low watermark — the band in |
| 495 | /// between is slack for playhead work, so fill can never trigger (or |
| 496 | /// refight) an eviction. Main thread. |
| 497 | func tryReserve(bytes: Int64, upTo ceiling: Int64) -> Bool { |
| 498 | guard budgetReady, ledgerBytes + reservedBytes + bytes <= ceiling else { return false } |
| 499 | reservedBytes += bytes |
| 500 | return true |
| 501 | } |
| 502 | |
| 503 | /// A reserved build finished: release its reservation and record what |
| 504 | /// actually landed on disk (new file minus any replaced one; 0 on failure). |
| 505 | func commitBuild(reserved: Int64, delta: Int64) { |
| 506 | reservedBytes = max(0, reservedBytes - reserved) |
| 507 | ledgerBytes = max(0, ledgerBytes + delta) |
| 508 | } |
| 509 | |
| 510 | /// Bytes written outside the reservation flow (filmstrips, waveforms — |
| 511 | /// small, but the ledger should still see them between reconciles). |
| 512 | func noteBytesAdded(_ bytes: Int64) { ledgerBytes += bytes } |
| 513 | |
| 514 | /// Cheap budget check for the old trigger sites (build completions, doc |
| 515 | /// open, manual). `reconcile: true` re-walks the disk first — use it when |
| 516 | /// ground truth matters (doc open, cap change, menu action). |
| 517 | func evictIfNeeded(reconcile: Bool = false) { |
| 518 | if reconcile || !budgetReady { reconcileLedger(); return } |
| 519 | if ledgerBytes + reservedBytes > maxCacheBytes { |
| 520 | evictToFit(need: 0) { _ in } |
| 521 | } |
| 522 | } |
| 523 | |
| 524 | /// Free cache space so `need` more bytes fit under `ceiling` (the cap for |
| 525 | /// playhead work, the low watermark for fill), evicting down to the low |
| 526 | /// watermark once it runs at all. Coldness order: whole cache dirs of |
| 527 | /// projects no window has open (LRU) first, then — when `openDocChunks` — |
| 528 | /// individual cold proxy chunks of open documents, never touching any |
| 529 | /// document's working set (see ChunkManager.evictionCandidates). |
| 530 | /// `completion(true)` on main once the space exists. |
| 531 | func evictToFit(need: Int64, openDocChunks: Bool = true, |
| 532 | upTo ceiling: Int64? = nil, |
| 533 | completion: @escaping (Bool) -> Void) { |
| 534 | let ceiling = ceiling ?? maxCacheBytes |
| 535 | guard budgetReady, !evicting else { completion(false); return } |
| 536 | let deficit = (ledgerBytes + reservedBytes + need) - lowWatermarkBytes |
| 537 | guard deficit > 0 else { completion(true); return } |
| 538 | evicting = true |
| 539 | // Whole-dir eviction must not touch any open document's media — those |
| 540 | // dirs also hold filmstrips/waveforms other windows are showing. Their |
| 541 | // cold CHUNKS are reclaimed individually via the candidates instead. |
| 542 | let inUse = Set(DocumentContext.allLive.flatMap { $0.store.project.media.map(\.cacheKey) }) |
| 543 | var chunkCands: [ChunkManager.EvictionCandidate] = [] |
| 544 | if openDocChunks { |
| 545 | for c in DocumentContext.allLive { chunkCands += c.chunks.evictionCandidates() } |
| 546 | chunkCands.sort { $0.coldness > $1.coldness } |
| 547 | } |
| 548 | let root = cacheRoot |
| 549 | DispatchQueue.global(qos: .utility).async { |
| 550 | let fm = FileManager.default |
| 551 | var freed: Int64 = 0 |
| 552 | var evictedDirs: [String] = [] |
| 553 | var evictedChunks: [String: [Int]] = [:] |
| 554 | // 1. Closed projects' whole cache dirs, least recently used first. |
| 555 | if let keys = try? fm.contentsOfDirectory(atPath: root.path) { |
| 556 | var entries: [(key: String, lastUsed: Double)] = [] |
| 557 | for key in keys where !inUse.contains(key) { |
| 558 | let dir = root.appendingPathComponent(key, isDirectory: true) |
| 559 | var isDir: ObjCBool = false |
| 560 | guard fm.fileExists(atPath: dir.path, isDirectory: &isDir), isDir.boolValue else { continue } |
| 561 | let lastUsed = Double((try? String(contentsOf: dir.appendingPathComponent("lastUsed"), encoding: .utf8)) ?? "") ?? 0 |
| 562 | entries.append((key, lastUsed)) |
| 563 | } |
| 564 | for e in entries.sorted(by: { $0.lastUsed < $1.lastUsed }) { |
| 565 | guard freed < deficit else { break } |
| 566 | let dir = root.appendingPathComponent(e.key, isDirectory: true) |
| 567 | let bytes = Self.directorySize(dir) |
| 568 | try? fm.removeItem(at: dir) |
| 569 | evictedDirs.append(e.key) |
| 570 | freed += bytes |
| 571 | } |
| 572 | } |
| 573 | // 2. Cold chunks of open documents, coldest (farthest from any |
| 574 | // playhead / off-timeline) first. |
| 575 | for c in chunkCands { |
| 576 | guard freed < deficit else { break } |
| 577 | let bytes = (try? fm.attributesOfItem(atPath: c.url.path)[.size] as? NSNumber)?.int64Value ?? 0 |
| 578 | guard bytes > 0 else { continue } |
| 579 | try? fm.removeItem(at: c.url) |
| 580 | evictedChunks[c.key, default: []].append(c.index) |
| 581 | freed += bytes |
| 582 | } |
| 583 | DispatchQueue.main.async { |
| 584 | self.ledgerBytes = max(0, self.ledgerBytes - freed) |
| 585 | NSLog("[cache] evict: freed %.2f GB (%d dirs, %d chunks) → ledger %.2f GB", |
| 586 | Double(freed) / 1e9, evictedDirs.count, |
| 587 | evictedChunks.values.map(\.count).reduce(0, +), |
| 588 | Double(self.ledgerBytes) / 1e9) |
| 589 | self.imgStateLock.lock() |
| 590 | for k in evictedDirs { |
| 591 | self.stripInfoCache[k] = nil |
| 592 | } |
| 593 | self.imgStateLock.unlock() |
| 594 | for k in evictedDirs { |
| 595 | self.enqueued.remove(k) |
| 596 | } |
| 597 | for c in DocumentContext.allLive { |
| 598 | c.chunks.forget(keys: evictedDirs) |
| 599 | for (key, idxs) in evictedChunks { |
| 600 | c.chunks.noteEvicted(key: key, indices: idxs) |
| 601 | } |
| 602 | } |
| 603 | self.evicting = false |
| 604 | let ok = self.ledgerBytes + self.reservedBytes + need <= ceiling |
| 605 | if !evictedDirs.isEmpty || !evictedChunks.isEmpty { |
| 606 | NotificationCenter.default.post(name: .mediaStatusChanged, object: nil) |
| 607 | // Only re-pump when something was actually freed — a no-op |
| 608 | // eviction re-pumping would loop pump → evict → pump forever. |
| 609 | for c in DocumentContext.allLive { c.chunks.budgetChanged() } |
| 610 | } |
| 611 | completion(ok) |
| 612 | } |
| 613 | } |
| 614 | } |
| 615 | |
| 616 | /// Delete leftovers no session references anymore: rescue-slice files |
| 617 | /// (`rNNNNNN.mov` — their state is memory-only, so a relaunch can't know |
| 618 | /// which part of the chunk they cover) and `.partial.mov` temps from |
| 619 | /// builds a crash interrupted. Only files from BEFORE this process |
| 620 | /// started — a slice the current session just built stays. |
| 621 | private static func purgeSessionTransients(root: URL) { |
| 622 | let fm = FileManager.default |
| 623 | guard let en = fm.enumerator(at: root, |
| 624 | includingPropertiesForKeys: [.contentModificationDateKey]) |
| 625 | else { return } |
| 626 | for case let f as URL in en { |
| 627 | let name = f.lastPathComponent |
| 628 | let isRescue = name.hasPrefix("r") && name.hasSuffix(".mov") && name.count == 11 |
| 629 | && f.deletingLastPathComponent().lastPathComponent == "chunks" |
| 630 | let isTemp = name.hasSuffix(".partial.mov") |
| 631 | guard isRescue || isTemp else { continue } |
| 632 | let m = (try? f.resourceValues(forKeys: [.contentModificationDateKey]) |
| 633 | .contentModificationDate) ?? .distantPast |
| 634 | if m < processStart { try? fm.removeItem(at: f) } |
| 635 | } |
| 636 | } |
| 637 | |
| 638 | static func directorySize(_ url: URL) -> Int64 { |
| 639 | var total: Int64 = 0 |
| 640 | if let en = FileManager.default.enumerator(at: url, includingPropertiesForKeys: [.fileSizeKey]) { |
| 641 | for case let f as URL in en { |
| 642 | total += Int64((try? f.resourceValues(forKeys: [.fileSizeKey]).fileSize) ?? 0) |
| 643 | } |
| 644 | } |
| 645 | return total |
| 646 | } |
| 647 | |
| 648 | // MARK: - Process helper |
| 649 | |
| 650 | struct RunResult { var exitCode: Int32; var stdout: String } |
| 651 | |
| 652 | /// Every child process currently running, so app teardown can take them |
| 653 | /// down too. A killed Sequencer otherwise leaves its in-flight ffmpeg |
| 654 | /// encodes running as orphans — and each one holds VideoToolbox decode |
| 655 | /// sessions from a SHARED machine-wide pool. Enough accumulated orphans |
| 656 | /// (a few rebuild-relaunch cycles' worth) and every AVPlayer in the next |
| 657 | /// app instance silently renders BLACK: items park, seeks land, |
| 658 | /// isReadyForDisplay says true, no error anywhere. Diagnosed 2026-07-11 |
| 659 | /// after the viewer went black with all state reporting healthy. |
| 660 | private static var liveChildren: [Int32: Process] = [:] |
| 661 | private static let childLock = NSLock() |
| 662 | |
| 663 | /// Terminate every live child (normal quit AND SIGTERM — run.sh pkills |
| 664 | /// the app on every rebuild, which is exactly how orphans were minted). |
| 665 | static func terminateChildren() { |
| 666 | childLock.lock() |
| 667 | let children = Array(liveChildren.values) |
| 668 | childLock.unlock() |
| 669 | for p in children where p.isRunning { p.terminate() } |
| 670 | } |
| 671 | |
| 672 | /// Kill orphaned ffmpeg processes from a PREVIOUS Sequencer instance |
| 673 | /// (crash, force-quit, kill -9 — anything terminateChildren couldn't |
| 674 | /// catch). Identified by their command line referencing our cache root, |
| 675 | /// so nothing else on the machine can match. Runs once at launch. |
| 676 | static func reapOrphans(cacheRoot: URL) { |
| 677 | let marker = cacheRoot.path |
| 678 | DispatchQueue.global(qos: .utility).async { |
| 679 | let r = run("/usr/bin/pgrep", ["-fl", "ffmpeg"]) |
| 680 | var killed = 0 |
| 681 | for line in r.stdout.split(separator: "\n") where line.contains(marker) { |
| 682 | guard let pid = Int32(line.prefix(while: \.isNumber)) else { continue } |
| 683 | kill(pid, SIGKILL) |
| 684 | killed += 1 |
| 685 | } |
| 686 | if killed > 0 { |
| 687 | SeqLog.log("[cache] reaped %d orphaned ffmpeg encoder(s) from a previous instance", killed) |
| 688 | } |
| 689 | } |
| 690 | } |
| 691 | |
| 692 | @discardableResult |
| 693 | static func run(_ path: String, _ args: [String], |
| 694 | duration: Double? = nil, |
| 695 | progress: ((Double) -> Void)? = nil) -> RunResult { |
| 696 | let p = Process() |
| 697 | p.executableURL = URL(fileURLWithPath: path) |
| 698 | p.arguments = args |
| 699 | let outPipe = Pipe() |
| 700 | p.standardOutput = outPipe |
| 701 | p.standardError = Pipe() // discard; keeps ffmpeg from blocking on tty |
| 702 | var collected = Data() |
| 703 | let wantsProgress = progress != nil && duration != nil && duration! > 0 |
| 704 | outPipe.fileHandleForReading.readabilityHandler = { h in |
| 705 | let chunk = h.availableData |
| 706 | if chunk.isEmpty { return } |
| 707 | collected.append(chunk) |
| 708 | if wantsProgress, let text = String(data: chunk, encoding: .utf8) { |
| 709 | for line in text.split(separator: "\n") { |
| 710 | if line.hasPrefix("out_time_us="), let us = Double(line.dropFirst("out_time_us=".count)) { |
| 711 | progress!(min(1, (us / 1_000_000) / duration!)) |
| 712 | } |
| 713 | } |
| 714 | } |
| 715 | } |
| 716 | do { |
| 717 | try p.run() |
| 718 | childLock.lock() |
| 719 | liveChildren[p.processIdentifier] = p |
| 720 | childLock.unlock() |
| 721 | p.waitUntilExit() |
| 722 | childLock.lock() |
| 723 | liveChildren.removeValue(forKey: p.processIdentifier) |
| 724 | childLock.unlock() |
| 725 | } catch { |
| 726 | return RunResult(exitCode: -1, stdout: "") |
| 727 | } |
| 728 | outPipe.fileHandleForReading.readabilityHandler = nil |
| 729 | if let rest = try? outPipe.fileHandleForReading.readToEnd() { collected.append(rest) } |
| 730 | return RunResult(exitCode: p.terminationStatus, |
| 731 | stdout: String(data: collected, encoding: .utf8) ?? "") |
| 732 | } |
| 733 | } |