diff --git a/native/cu-helper/Package.swift b/native/cu-helper/Package.swift index 5ac9e322..b33e373a 100644 --- a/native/cu-helper/Package.swift +++ b/native/cu-helper/Package.swift @@ -49,6 +49,8 @@ let package = Package( linkerSettings: [ .linkedFramework("AppKit"), .linkedFramework("CoreGraphics"), + .linkedFramework("CoreMedia"), + .linkedFramework("CoreVideo"), .linkedFramework("QuartzCore"), .linkedFramework("ApplicationServices"), .linkedFramework("ScreenCaptureKit"), diff --git a/native/cu-helper/Sources/cu-helper/Capture.swift b/native/cu-helper/Sources/cu-helper/Capture.swift index 746cc8f1..a309418b 100644 --- a/native/cu-helper/Sources/cu-helper/Capture.swift +++ b/native/cu-helper/Sources/cu-helper/Capture.swift @@ -38,6 +38,26 @@ import ImageIO import ScreenCaptureKit import UniformTypeIdentifiers +enum WindowShotCaptureSource: Equatable, Sendable { + case stream + case screenshotManager + case screenCaptureCLI + + var isLiveStream: Bool { self == .stream } +} + +struct WindowShot: Sendable { + let base64: String + let width: Int + let height: Int + let originX: Double + let originY: Double + let pointWidth: Double + let pointHeight: Double + let windowID: CGWindowID + let source: WindowShotCaptureSource +} + @available(macOS 14.0, *) @MainActor public enum Capture { @@ -431,14 +451,11 @@ public enum Capture { /// (pixels-per-point) to invert image-pixel coordinates back into the /// global-point space that clicks/cursor/glow all live in. `nil` on any /// failure. Never throws; never prompts. - public static func windowShot( + static func windowShot( pid: pid_t, preferredWindowID: CGWindowID? = nil, scale: Double = 0.5 - ) async -> (base64: String, width: Int, height: Int, - originX: Double, originY: Double, - pointWidth: Double, pointHeight: Double, - windowID: CGWindowID)? { + ) async -> WindowShot? { // Passive permission gate — no prompt on the hot path. A denied grant // means SCK would hand us a black frame, so bail to `nil` early and let // the caller fall back to AX-text-only. @@ -461,9 +478,17 @@ public enum Capture { scale: outputScale ) { if let encoded = pngBase64WithSize(image) { - return (encoded.base64, encoded.width, encoded.height, - Double(f.origin.x), Double(f.origin.y), - Double(f.width), Double(f.height), target.windowID) + return WindowShot( + base64: encoded.base64, + width: encoded.width, + height: encoded.height, + originX: Double(f.origin.x), + originY: Double(f.origin.y), + pointWidth: Double(f.width), + pointHeight: Double(f.height), + windowID: target.windowID, + source: .screenshotManager + ) } // SCK produced pixels but PNG/base64 failed — degrade to `nil` // rather than re-capturing; the caller still gets AX text. @@ -474,9 +499,17 @@ public enum Capture { if let raw = screencaptureWindow(windowID: target.windowID) { let scaled = (try? scaleImage(raw, scale: outputScale)) ?? raw if let encoded = pngBase64WithSize(scaled) { - return (encoded.base64, encoded.width, encoded.height, - Double(f.origin.x), Double(f.origin.y), - Double(f.width), Double(f.height), target.windowID) + return WindowShot( + base64: encoded.base64, + width: encoded.width, + height: encoded.height, + originX: Double(f.origin.x), + originY: Double(f.origin.y), + pointWidth: Double(f.width), + pointHeight: Double(f.height), + windowID: target.windowID, + source: .screenCaptureCLI + ) } } @@ -487,7 +520,7 @@ public enum Capture { /// A capture candidate: the CoreGraphics window id plus its global Quartz /// frame (top-left origin) used to size the SCK output buffer. - private struct WindowCandidate { + struct WindowCandidate { let windowID: CGWindowID let frame: CGRect } @@ -505,7 +538,7 @@ public enum Capture { /// helper discards the `windowID` (which SCK needs to match a specific /// `SCWindow`) and filters to layer-0 only. For a window-locked capture we /// want the owner's frontmost window regardless of layer, keyed by id. - private static func bestWindow( + static func bestWindow( forPid pid: pid_t, preferredWindowID: CGWindowID? ) -> WindowCandidate? { @@ -739,7 +772,7 @@ public enum Capture { /// `get_app_state`'s MCP envelope uses `mimeType: "image/png"`, so unlike the /// JPEG `screenshot` path this is lossless PNG. Returns `nil` on any encode /// failure (never throws) so `windowShot` degrades to AX-text-only. - private static func pngBase64WithSize( + static func pngBase64WithSize( _ image: CGImage ) -> (base64: String, width: Int, height: Int)? { let data = NSMutableData() @@ -759,7 +792,7 @@ public enum Capture { /// Backing scale factor (Retina = 2.0) for the screen hosting `frame`. Used /// to size the SCK output buffer in native pixels before the requested /// downscale. Falls back to the main screen, then 2.0. - private static func backingScaleFactor(forWindowFrame frame: CGRect) -> Double { + static func backingScaleFactor(forWindowFrame frame: CGRect) -> Double { if let screen = NSScreen.screens.first(where: { $0.frame.intersects(frame) }) { return Double(screen.backingScaleFactor) } diff --git a/native/cu-helper/Sources/cu-helper/CommandRouter.swift b/native/cu-helper/Sources/cu-helper/CommandRouter.swift index b5258d83..de9bca94 100644 --- a/native/cu-helper/Sources/cu-helper/CommandRouter.swift +++ b/native/cu-helper/Sources/cu-helper/CommandRouter.swift @@ -71,6 +71,7 @@ public final class CommandRouter { private let capabilities: Capabilities private let inputMonitor: PhysicalInputEpochMonitor private let foregroundRuntime: ForegroundLeaseRuntime + private let windowCaptureProvider: (any WindowCaptureProviding)? /// After a left `mouse_down` (decomposed drag) the held point is parked /// here so a following `mouse_up` releases at the same logical location. @@ -82,12 +83,14 @@ public final class CommandRouter { init( cursor: VirtualCursor, capabilities: Capabilities, - inputMonitor: PhysicalInputEpochMonitor + inputMonitor: PhysicalInputEpochMonitor, + windowCaptureProvider: (any WindowCaptureProviding)? = nil ) { self.cursor = cursor self.capabilities = capabilities self.inputMonitor = inputMonitor self.foregroundRuntime = .live(monitor: inputMonitor) + self.windowCaptureProvider = windowCaptureProvider } func resetHeldSessionState() { @@ -100,11 +103,16 @@ public final class CommandRouter { AXTree.resetSessionSnapshots() Self.lastShotTransform.removeAll() Self.lastCaptureDigest.removeAll() + windowCaptureProvider?.invalidate() // Apps we told they were focused must be told they are not, or the // belief outlives the session that needed it. SyntheticWindowFocus.relinquishAll() } + func invalidateWindowCaptureStream() { + windowCaptureProvider?.invalidate() + } + /// Dispatch one command. The payload is the raw decoded JSON object the /// bridge sent (or `.object([:])` for the empty-payload commands). Returns /// the `result` ``JSONValue``; throws ``CUError`` on any failure. @@ -554,13 +562,15 @@ public final class CommandRouter { private static func identicalCaptureNotice( pid: pid_t, base64: String, - windowID: CGWindowID + windowID: CGWindowID, + liveStreamActive: Bool ) -> String? { let digest = base64.hashValue defer { lastCaptureDigest[pid] = digest } guard let previous = lastCaptureDigest[pid], previous == digest else { return nil } return TargetVisibilityPolicy.identicalCaptureNotice( - windowIsCovered: WindowGeometry.isFullyCovered(windowID: windowID) + windowIsCovered: WindowGeometry.isFullyCovered(windowID: windowID), + liveStreamActive: liveStreamActive ) } @@ -572,8 +582,16 @@ public final class CommandRouter { object["axText"] = .string(existing + "\n\n" + notice) } - private func handleGetAppState(_ payload: JSONValue) async throws -> JSONValue { - let disableDiff = try optionalBoolean(payload, key: "disableDiff") ?? false + private func handleGetAppState( + _ payload: JSONValue, + windowChangeRetriesRemaining: Int = 1, + forceFullSnapshot: Bool = false + ) async throws -> JSONValue { + let requestedDisableDiff = try optionalBoolean(payload, key: "disableDiff") ?? false + let disableDiff = Self.effectiveDisableDiff( + requested: requestedDisableDiff, + forceFullSnapshot: forceFullSnapshot + ) let selector = try AppTargetResolver.requiredSelector(payload: payload) try requireAXTrusted() // Resolve a RUNNING match first; if an app was named but isn't running, @@ -631,12 +649,6 @@ public final class CommandRouter { ) if let windowID = snapshotEvidence.keyWindowID { object["windowID"] = .int(Int(windowID)) - if WindowGeometry.isFullyCovered(windowID: windowID) { - Self.appendAXNotice( - TargetVisibilityPolicy.coveredCaptureNotice, - to: &object - ) - } } // A failed or mismatched fresh capture must not leave coordinates from @@ -653,11 +665,96 @@ public final class CommandRouter { appIsBusy: result.axText.contains("progress indicator") ) - if let shot = await Capture.windowShot( + let streamedShot = await windowCaptureProvider?.windowShot( pid: pid, + processIdentity: snapshotEvidence.processIdentity, preferredWindowID: snapshotEvidence.keyWindowID, - scale: 0.5 + scale: 0.5, + newerThanUptime: MutationClock.lastMutation() + ) + guard TargetVisibilityPolicy.captureTargetStillMatches( + snapshotWindowID: snapshotEvidence.keyWindowID, + currentWindowID: AXTree.currentKeyWindowID(pid: pid) + ) else { + windowCaptureProvider?.invalidate() + guard windowChangeRetriesRemaining > 0 else { + throw CUError( + "stale_snapshot", + "The target key window changed while get_app_state was capturing it" + ) + } + return try await handleGetAppState( + payload, + windowChangeRetriesRemaining: windowChangeRetriesRemaining - 1, + forceFullSnapshot: true + ) + } + + var windowIsCovered = snapshotEvidence.keyWindowID.map { + WindowGeometry.isFullyCovered(windowID: $0) + } ?? false + var shot: WindowShot? + if let streamedShot { + shot = streamedShot + } else if TargetVisibilityPolicy.permitsOneShotFallback( + windowIsCovered: windowIsCovered, + streamProviderInstalled: windowCaptureProvider != nil ) { + shot = await Capture.windowShot( + pid: pid, + preferredWindowID: snapshotEvidence.keyWindowID, + scale: 0.5 + ) + } else { + // A one-shot capture can repeat compositor-cached pixels for a + // covered Chromium/CEF window. Returning no image is safer than + // presenting that stale fallback as post-action evidence; AX + // state remains available and the live stream stays installed + // for the next read. + shot = nil + } + + // A fallback capture can itself wait for SCK/CLI. Revalidate both + // identity and coverage after that await so a newly opened sheet or + // newly covering foreground window cannot turn a safe decision into + // a stale screenshot attachment. + guard TargetVisibilityPolicy.captureTargetStillMatches( + snapshotWindowID: snapshotEvidence.keyWindowID, + currentWindowID: AXTree.currentKeyWindowID(pid: pid) + ) else { + windowCaptureProvider?.invalidate() + guard windowChangeRetriesRemaining > 0 else { + throw CUError( + "stale_snapshot", + "The target key window changed while get_app_state was capturing it" + ) + } + return try await handleGetAppState( + payload, + windowChangeRetriesRemaining: windowChangeRetriesRemaining - 1, + forceFullSnapshot: true + ) + } + windowIsCovered = snapshotEvidence.keyWindowID.map { + WindowGeometry.isFullyCovered(windowID: $0) + } ?? false + if windowIsCovered, + windowCaptureProvider != nil, + shot?.source.isLiveStream != true { + shot = nil + } + + if windowIsCovered { + Self.appendAXNotice( + TargetVisibilityPolicy.coveredCaptureNotice( + liveStreamActive: shot?.source.isLiveStream == true + ), + to: &object + ) + } + + if let shot, + AXTree.currentProcessIdentity(pid: pid) == snapshotEvidence.processIdentity { // An identical capture is the other half of the same problem: // the pixels cannot say whether the action missed or the window // is not painting, and the model reading them cannot tell @@ -666,7 +763,8 @@ public final class CommandRouter { if let notice = Self.identicalCaptureNotice( pid: pid, base64: shot.base64, - windowID: shot.windowID + windowID: shot.windowID, + liveStreamActive: shot.source.isLiveStream ) { Self.appendAXNotice(notice, to: &object) } @@ -702,6 +800,15 @@ public final class CommandRouter { return .object(object) } + /// A snapshot discarded by an internal key-window retry was never delivered + /// to the model, so it cannot become the baseline for a returned diff. + static func effectiveDisableDiff( + requested: Bool, + forceFullSnapshot: Bool + ) -> Bool { + requested || forceFullSnapshot + } + /// The inverse of `windowShot`'s image-pixel space: a window's global Quartz /// top-left origin (points) plus image-pixels-per-window-point. Used to map a /// model-supplied image-pixel coordinate back to the global-point space that diff --git a/native/cu-helper/Sources/cu-helper/Daemon.swift b/native/cu-helper/Sources/cu-helper/Daemon.swift index d0e83fe1..89377848 100644 --- a/native/cu-helper/Sources/cu-helper/Daemon.swift +++ b/native/cu-helper/Sources/cu-helper/Daemon.swift @@ -103,11 +103,18 @@ public final class Daemon { self.socketPath = socketPath self.inputMonitor = inputMonitor let cursor = VirtualCursor(headless: false) + let windowCaptureProvider: (any WindowCaptureProviding)? + if #available(macOS 14.0, *) { + windowCaptureProvider = WindowCaptureStreamManager() + } else { + windowCaptureProvider = nil + } self.cursor = cursor self.router = CommandRouter( cursor: cursor, capabilities: Capabilities(headless: false), - inputMonitor: inputMonitor + inputMonitor: inputMonitor, + windowCaptureProvider: windowCaptureProvider ) } @@ -582,6 +589,7 @@ public final class Daemon { explicitOverlayTarget = nil Injection.clearResolvedTarget() cursor.hide() + router.invalidateWindowCaptureStream() } /// The app the cursor should follow: ONLY the app the last injection / @@ -667,6 +675,7 @@ public final class Daemon { ) { [weak self] _ in MainActor.assumeIsolated { guard let self else { return } + self.router.invalidateWindowCaptureStream() // Re-warm cursor windows for the new screen set. self.cursor.preload() } diff --git a/native/cu-helper/Sources/cu-helper/TargetVisibilityPolicy.swift b/native/cu-helper/Sources/cu-helper/TargetVisibilityPolicy.swift index 78929185..419509fe 100644 --- a/native/cu-helper/Sources/cu-helper/TargetVisibilityPolicy.swift +++ b/native/cu-helper/Sources/cu-helper/TargetVisibilityPolicy.swift @@ -1,35 +1,85 @@ +import CoreGraphics import Foundation /// Explains a repeated window capture without turning ordinary occlusion into /// an input failure. /// -/// ScreenCaptureKit captures a target window independently of the desktop's -/// stacking order, and all fallback input is addressed to the target process -/// and window. Another app covering the target is therefore not a reason to -/// raise it, activate it, or reject a mutation. A byte-identical capture is -/// still useful evidence, but only about pixels: it cannot prove that an action -/// failed, and it must never stop background automation. +/// The daemon keeps a desktop-independent SCStream subscribed to the target +/// window throughout a Computer Use turn. That gives Chromium/CEF a continuous +/// WindowServer consumer while another app fully covers it, while all input is +/// still addressed to the target process/window. Occlusion is therefore never +/// a reason to raise, activate, or reject a mutation. enum TargetVisibilityPolicy { - static let coveredCaptureNotice = """ - NOTE: Another application fully covers the target window. The screenshot \ - is captured from that window rather than from the visible desktop, but it \ - may be stale if the target paused its renderer while covered. Coverage does \ - not block Accessibility actions or app- and window-targeted input; continue \ - the task without activating or raising the target just to expose it. - """ + static func captureTargetStillMatches( + snapshotWindowID: CGWindowID?, + currentWindowID: CGWindowID? + ) -> Bool { + snapshotWindowID == currentWindowID + } - static func identicalCaptureNotice(windowIsCovered: Bool) -> String { - let cause = windowIsCovered - ? """ - The target window is covered, so this image may be older than it \ - looks if that app paused its renderer. Coverage does not block \ - app- and window-targeted input. - """ - : """ - The target window is not fully covered, so coverage does not \ - explain the identical pixels, and no visible pixel change was \ - observed. - """ + /// Once the daemon has a stream-capable provider, a covered-window read + /// must never fall back to a one-shot compositor capture: those pixels can + /// predate the action even though the Accessibility mutation succeeded. + /// Older macOS versions without the provider keep their compatibility path. + static func permitsOneShotFallback( + windowIsCovered: Bool, + streamProviderInstalled: Bool + ) -> Bool { + !windowIsCovered || !streamProviderInstalled + } + + static func coveredCaptureNotice(liveStreamActive: Bool) -> String { + if liveStreamActive { + return """ + NOTE: Another application fully covers the target window. A \ + long-lived window stream remains subscribed while it is covered, \ + and this screenshot comes from its latest complete frame rather \ + than the visible desktop. Coverage does not block Accessibility \ + actions or app- and window-targeted input; continue without \ + activating or raising the target. + """ + } + return """ + NOTE: Another application fully covers the target window, and a fresh \ + live window-stream frame was unavailable for this read. A one-shot \ + screenshot is intentionally not used on stream-capable systems because \ + it may contain compositor-cached pixels. Coverage still does not block \ + Accessibility actions or app- and window-targeted input; continue \ + without activating or raising the target, and rely on the accessibility \ + state until the stream produces a fresh frame. + """ + } + + static func identicalCaptureNotice( + windowIsCovered: Bool, + liveStreamActive: Bool + ) -> String { + let cause: String + if liveStreamActive { + cause = windowIsCovered + ? """ + The target is covered, but its long-lived window stream is \ + still active; the newest complete frame contains no visible \ + pixel change. Coverage does not block app- and window-targeted \ + input. + """ + : """ + The target window is not fully covered, and the newest \ + complete stream frame contains no visible pixel change. + """ + } else { + cause = windowIsCovered + ? """ + A live window-stream frame was unavailable while the target \ + was covered, so this fallback image may be stale. Coverage \ + does not block app- and window-targeted input. + """ + : """ + The target window is not fully covered, so coverage does not \ + explain the identical pixels, and no visible pixel change was \ + observed. + """ + } return """ NOTE: This screenshot is byte-for-byte identical to the previous one. \ diff --git a/native/cu-helper/Sources/cu-helper/WindowCaptureStream.swift b/native/cu-helper/Sources/cu-helper/WindowCaptureStream.swift new file mode 100644 index 00000000..92291d1f --- /dev/null +++ b/native/cu-helper/Sources/cu-helper/WindowCaptureStream.swift @@ -0,0 +1,720 @@ +import CoreGraphics +import CoreMedia +import CoreVideo +import Foundation +@preconcurrency import ScreenCaptureKit + +/// The process/window/config identity of one long-lived window stream. +/// +/// A PID or Window Server id can be reused. The proven process identity keeps +/// a replacement process from inheriting frames captured for an earlier +/// lifetime, while the output dimensions force a fresh stream after a resize +/// or backing-scale change. +struct WindowCaptureStreamKey: Equatable, Sendable { + let pid: pid_t + let processIdentity: AXTreeProcessIdentity + let windowID: CGWindowID + let pixelWidth: Int + let pixelHeight: Int +} + +/// Current geometry plus the stable identity/config key used by the stream. +/// Origin is deliberately outside the key: moving an unchanged window does +/// not require a new ScreenCaptureKit consumer, but every delivered shot uses +/// freshly-read geometry so screenshot coordinates still invert correctly. +struct WindowCaptureStreamTarget: Equatable, Sendable { + let key: WindowCaptureStreamKey + let originX: Double + let originY: Double + let pointWidth: Double + let pointHeight: Double +} + +/// An immutable copy of the newest complete BGRA frame. The ScreenCaptureKit +/// pixel buffer is reused after its callback returns, so bytes must be copied +/// before crossing out of the sample queue. +struct WindowCaptureStreamFrame: Equatable, Sendable { + let bytes: Data + let width: Int + let height: Int + let bytesPerRow: Int + let sequence: UInt64 + let receivedUptime: TimeInterval +} + +enum WindowCaptureFrameStatusPolicy { + static func accepts(_ status: SCFrameStatus) -> Bool { + status == .complete || status == .started + } + + static func marksFailure(_ status: SCFrameStatus) -> Bool { + status == .stopped + } +} + +@MainActor +protocol WindowCaptureProviding: AnyObject { + func windowShot( + pid: pid_t, + processIdentity: AXTreeProcessIdentity, + preferredWindowID: CGWindowID?, + scale: Double, + newerThanUptime: TimeInterval? + ) async -> WindowShot? + + /// Synchronously makes every current/late frame unusable, then retires the + /// physical stream asynchronously so turn cleanup never waits on SCK. + func invalidate() +} + +@MainActor +protocol WindowCaptureStreamSource: AnyObject { + var targetKey: WindowCaptureStreamKey { get } + var hasFailed: Bool { get } + func start() async throws + func latestFrame() -> WindowCaptureStreamFrame? + func retire() +} + +@MainActor +protocol WindowCaptureStreamSourceFactory: AnyObject { + func makeSource(for target: WindowCaptureStreamTarget) -> any WindowCaptureStreamSource +} + +/// One daemon-owned capture session. A source starts on the first +/// `get_app_state`, remains alive across mutations, and is reused until the +/// process lifetime, window identity, or output configuration changes. +@MainActor +final class WindowCaptureStreamManager: WindowCaptureProviding { + private struct Entry { + let generation: UInt64 + let source: any WindowCaptureStreamSource + } + + private let factory: any WindowCaptureStreamSourceFactory + private let frameWaitAttempts: Int + private let frameWaitNanoseconds: UInt64 + private var generation: UInt64 = 0 + private var starting: Entry? + private var active: Entry? + + convenience init() { + self.init(factory: ScreenCaptureKitWindowStreamSourceFactory()) + } + + init( + factory: any WindowCaptureStreamSourceFactory, + frameWaitAttempts: Int = 12, + frameWaitNanoseconds: UInt64 = 50_000_000 + ) { + self.factory = factory + self.frameWaitAttempts = max(0, frameWaitAttempts) + self.frameWaitNanoseconds = frameWaitNanoseconds + } + + func windowShot( + pid: pid_t, + processIdentity: AXTreeProcessIdentity, + preferredWindowID: CGWindowID?, + scale: Double, + newerThanUptime: TimeInterval? + ) async -> WindowShot? { + guard Capture.hasScreenRecordingPermission(), + processIdentity.isProven, + AXTree.currentProcessIdentity(pid: pid) == processIdentity else { + invalidate() + return nil + } + + // Geometry may change while the first frame is in flight. Re-resolve + // the exact AX window once and reconfigure instead of binding new + // pixels to an old coordinate transform. + for _ in 0..<2 { + guard let target = Capture.windowCaptureStreamTarget( + pid: pid, + processIdentity: processIdentity, + preferredWindowID: preferredWindowID, + scale: scale + ) else { + invalidate() + return nil + } + guard let frame = await frame( + for: target, + newerThanUptime: newerThanUptime + ) else { + return nil + } + if let preferredWindowID, + !TargetVisibilityPolicy.captureTargetStillMatches( + snapshotWindowID: preferredWindowID, + currentWindowID: AXTree.currentKeyWindowID(pid: pid) + ) { + // The AX snapshot and pixel target must describe the same key + // window. A sheet/dialog can replace the key window while SCK + // is waiting; retire this source so the router can rebuild the + // entire snapshot against the replacement window. + invalidate() + return nil + } + guard AXTree.currentProcessIdentity(pid: pid) == processIdentity, + let current = Capture.windowCaptureStreamTarget( + pid: pid, + processIdentity: processIdentity, + preferredWindowID: target.key.windowID, + scale: scale + ) else { + invalidate() + return nil + } + guard current.key == target.key else { + continue + } + return Capture.windowShot(from: frame, target: current) + } + return nil + } + + /// Internal transition seam used by deterministic tests. No TCC, AX, or + /// real ScreenCaptureKit object is needed to prove lifecycle behavior. + func frame( + for target: WindowCaptureStreamTarget, + newerThanUptime: TimeInterval? + ) async -> WindowCaptureStreamFrame? { + // One bounded restart is allowed after a start/delegate failure. A live + // source that simply has not delivered a fresh frame remains installed + // so the next get_app_state can reuse it. + for restartAttempt in 0..<2 { + guard let source = await source(for: target) else { continue } + if source.hasFailed { + retire(source: source) + continue + } + if let frame = await waitForFrame( + from: source, + target: target, + newerThanUptime: newerThanUptime + ) { + return frame + } + if source.hasFailed { + retire(source: source) + continue + } + if newerThanUptime != nil, restartAttempt == 0 { + // A live stream can become silently suspended without a + // delegate error. One bounded rebuild forces a new `.started` + // frame and prevents every later post-mutation read from + // timing out forever on the same starved source. + retire(source: source) + continue + } + return nil + } + return nil + } + + func invalidate() { + generation &+= 1 + let currentStarting = starting + let currentActive = active + starting = nil + active = nil + + // `retire` first invalidates the callback mailbox synchronously. Old + // queued callbacks therefore cannot populate a later generation even + // if the asynchronous physical stop finishes after a new stream starts. + currentStarting?.source.retire() + if currentActive?.source !== currentStarting?.source { + currentActive?.source.retire() + } + } + + var activeGenerationForTesting: UInt64? { active?.generation } + var activeKeyForTesting: WindowCaptureStreamKey? { active?.source.targetKey } + + private func source( + for target: WindowCaptureStreamTarget + ) async -> (any WindowCaptureStreamSource)? { + if let active, + active.source.targetKey == target.key, + !active.source.hasFailed { + return active.source + } + + if let active { + retire(source: active.source) + } + + generation &+= 1 + let operationGeneration = generation + let source = factory.makeSource(for: target) + starting = Entry(generation: operationGeneration, source: source) + + do { + try await source.start() + } catch { + if starting?.generation == operationGeneration { + starting = nil + } + source.retire() + return nil + } + + guard generation == operationGeneration, + starting?.generation == operationGeneration else { + source.retire() + return nil + } + starting = nil + let installed = Entry(generation: operationGeneration, source: source) + active = installed + return installed.source + } + + private func waitForFrame( + from source: any WindowCaptureStreamSource, + target: WindowCaptureStreamTarget, + newerThanUptime: TimeInterval? + ) async -> WindowCaptureStreamFrame? { + for attempt in 0...frameWaitAttempts { + if source.hasFailed { return nil } + if let frame = source.latestFrame(), + frame.width == target.key.pixelWidth, + frame.height == target.key.pixelHeight, + newerThanUptime.map({ frame.receivedUptime > $0 }) ?? true { + return frame + } + guard attempt < frameWaitAttempts else { break } + if frameWaitNanoseconds == 0 { + await Task.yield() + } else { + try? await Task.sleep(nanoseconds: frameWaitNanoseconds) + } + } + return nil + } + + private func retire(source: any WindowCaptureStreamSource) { + generation &+= 1 + if active?.source === source { active = nil } + if starting?.source === source { starting = nil } + source.retire() + } +} + +@MainActor +private final class ScreenCaptureKitWindowStreamSourceFactory: WindowCaptureStreamSourceFactory { + func makeSource( + for target: WindowCaptureStreamTarget + ) -> any WindowCaptureStreamSource { + ScreenCaptureKitWindowStreamSource(target: target) + } +} + +/// The real SCK stream. Lifecycle stays on MainActor; sample delivery is copied +/// on one serial queue into the lock-backed mailbox below. +@MainActor +final class ScreenCaptureKitWindowStreamSource: WindowCaptureStreamSource { + private static let sampleQueue = DispatchQueue( + label: "dev.cchaha.cu-helper.window-stream.frames", + qos: .userInitiated + ) + private static let startTimeout: TimeInterval = 2.5 + + let targetKey: WindowCaptureStreamKey + private let target: WindowCaptureStreamTarget + private let output = WindowCaptureStreamMailbox() + private var stream: SCStream? + private var retired = false + + init(target: WindowCaptureStreamTarget) { + self.target = target + self.targetKey = target.key + } + + var hasFailed: Bool { output.hasFailed } + + func latestFrame() -> WindowCaptureStreamFrame? { + output.latestFrame() + } + + func start() async throws { + guard !retired else { + throw CUError("capture_failed", "Window stream was retired before it started") + } + + // Bound the whole startup, including SCShareableContent enumeration. + // That API can wedge before an SCStream exists, so timing only + // `startCapture()` would still let the daemon request hang indefinitely. + let gate = SCStreamStartGate() + let startTask = Task { @MainActor in + try await self.startCaptureWithoutTimeout() + } + Task { @MainActor in + do { + try await startTask.value + gate.resolve(.success) + } catch { + gate.resolve(.failure(error.localizedDescription)) + } + } + DispatchQueue.global(qos: .userInitiated).asyncAfter( + deadline: .now() + Self.startTimeout + ) { + gate.resolve(.failure("Timed out starting the window stream")) + } + + switch await gate.wait() { + case .success: + guard !retired else { + throw CUError("capture_failed", "Window stream was retired while starting") + } + case .failure(let message): + startTask.cancel() + retire() + throw CUError("capture_failed", message) + } + } + + private func startCaptureWithoutTimeout() async throws { + let content: SCShareableContent + do { + content = try await SCShareableContent.excludingDesktopWindows( + false, + onScreenWindowsOnly: false + ) + } catch { + throw CUError( + "capture_failed", + "SCShareableContent failed for window stream: \(error.localizedDescription)" + ) + } + + guard !retired, + let window = content.windows.first(where: { + $0.windowID == targetKey.windowID + && $0.owningApplication?.processID == targetKey.pid + }) else { + throw CUError( + "capture_failed", + "The exact target window disappeared before streaming started" + ) + } + + let filter = SCContentFilter(desktopIndependentWindow: window) + let configuration = Self.makeConfiguration(for: target) + let stream = SCStream( + filter: filter, + configuration: configuration, + delegate: output + ) + do { + try stream.addStreamOutput( + output, + type: .screen, + sampleHandlerQueue: Self.sampleQueue + ) + } catch { + throw CUError( + "capture_failed", + "Adding the window stream output failed: \(error.localizedDescription)" + ) + } + self.stream = stream + try await stream.startCapture() + guard !retired else { + throw CUError("capture_failed", "Window stream was retired while starting") + } + } + + func retire() { + guard !retired else { return } + retired = true + output.invalidate() + guard let stream else { return } + self.stream = nil + + // Do not await SCK shutdown on overlay_hide/disconnect. The mailbox is + // already inert, and retaining the stream in this completion closure + // lets ScreenCaptureKit finish cleanup without delaying the turn. + Task { @MainActor in + try? await stream.stopCapture() + } + } + + static func makeConfiguration( + for target: WindowCaptureStreamTarget + ) -> SCStreamConfiguration { + let configuration = SCStreamConfiguration() + configuration.width = max(1, target.key.pixelWidth) + configuration.height = max(1, target.key.pixelHeight) + // Match Codex's long-lived window stream cadence and buffering. This + // keeps a continuous WindowServer consumer for an occluded renderer; + // it is not a polling screenshot throttle. + configuration.minimumFrameInterval = CMTime(value: 1, timescale: 60) + configuration.queueDepth = 5 + configuration.showsCursor = false + configuration.capturesAudio = false + configuration.scalesToFit = true + configuration.preservesAspectRatio = true + configuration.pixelFormat = kCVPixelFormatType_32BGRA + configuration.colorSpaceName = CGColorSpace.sRGB + configuration.captureResolution = .best + configuration.ignoreShadowsSingleWindow = true + configuration.ignoreShadowsDisplay = true + return configuration + } +} + +/// Only pixel-bearing ScreenCaptureKit frames enter this one-frame mailbox. +/// `.started` is the first generated frame after start and `.complete` is a +/// later generated frame. Idle, blank, suspended, and stopped notifications +/// never advance sequence or satisfy a post-mutation freshness watermark. +private final class WindowCaptureStreamMailbox: NSObject, SCStreamOutput, SCStreamDelegate, @unchecked Sendable { + private let lock = NSLock() + private var accepting = true + private var failed = false + private var sequence: UInt64 = 0 + private var frame: WindowCaptureStreamFrame? + + var hasFailed: Bool { + lock.lock() + defer { lock.unlock() } + return failed + } + + func latestFrame() -> WindowCaptureStreamFrame? { + lock.lock() + defer { lock.unlock() } + return accepting ? frame : nil + } + + func invalidate() { + lock.lock() + accepting = false + frame = nil + lock.unlock() + } + + func stream( + _ stream: SCStream, + didOutputSampleBuffer sampleBuffer: CMSampleBuffer, + of outputType: SCStreamOutputType + ) { + guard outputType == .screen, + sampleBuffer.isValid, + sampleBuffer.dataReadiness == .ready, + let status = Self.frameStatus(sampleBuffer) else { + return + } + if WindowCaptureFrameStatusPolicy.marksFailure(status) { + markFailed() + return + } + guard WindowCaptureFrameStatusPolicy.accepts(status), + let pixelBuffer = sampleBuffer.imageBuffer, + CVPixelBufferGetPixelFormatType(pixelBuffer) == kCVPixelFormatType_32BGRA else { + return + } + // Timestamp callback arrival before copying. A mutation can land on the + // main actor while a large pixel buffer is being copied; stamping after + // that copy would make pre-mutation pixels look post-mutation fresh. + let receivedUptime = ProcessInfo.processInfo.systemUptime + + CVPixelBufferLockBaseAddress(pixelBuffer, .readOnly) + defer { CVPixelBufferUnlockBaseAddress(pixelBuffer, .readOnly) } + guard let baseAddress = CVPixelBufferGetBaseAddress(pixelBuffer) else { return } + + let width = CVPixelBufferGetWidth(pixelBuffer) + let height = CVPixelBufferGetHeight(pixelBuffer) + let bytesPerRow = CVPixelBufferGetBytesPerRow(pixelBuffer) + let (minimumBytesPerRow, rowOverflow) = width.multipliedReportingOverflow(by: 4) + let (byteCount, overflow) = bytesPerRow.multipliedReportingOverflow(by: height) + guard width > 0, + height > 0, + !rowOverflow, + bytesPerRow >= minimumBytesPerRow, + !overflow else { + return + } + + let copied = Data(bytes: baseAddress, count: byteCount) + + lock.lock() + guard accepting else { + lock.unlock() + return + } + sequence &+= 1 + frame = WindowCaptureStreamFrame( + bytes: copied, + width: width, + height: height, + bytesPerRow: bytesPerRow, + sequence: sequence, + receivedUptime: receivedUptime + ) + lock.unlock() + } + + func stream(_ stream: SCStream, didStopWithError error: Error) { + markFailed() + } + + private func markFailed() { + lock.lock() + if accepting { + failed = true + frame = nil + } + lock.unlock() + } + + private static func frameStatus(_ sampleBuffer: CMSampleBuffer) -> SCFrameStatus? { + guard + let attachments = CMSampleBufferGetSampleAttachmentsArray( + sampleBuffer, + createIfNecessary: false + ) as? [[SCStreamFrameInfo: Any]], + let rawStatus = attachments.first?[.status] as? Int, + let status = SCFrameStatus(rawValue: rawStatus) + else { + return nil + } + return status + } +} + +private final class SCStreamStartGate: @unchecked Sendable { + enum Outcome: Sendable { + case success + case failure(String) + } + + private let lock = NSLock() + private var outcome: Outcome? + private var continuation: CheckedContinuation? + + func wait() async -> Outcome { + await withCheckedContinuation { continuation in + lock.lock() + if let outcome { + lock.unlock() + continuation.resume(returning: outcome) + } else { + self.continuation = continuation + lock.unlock() + } + } + } + + func resolve(_ outcome: Outcome) { + let continuation: CheckedContinuation? + lock.lock() + guard self.outcome == nil else { + lock.unlock() + return + } + self.outcome = outcome + continuation = self.continuation + self.continuation = nil + lock.unlock() + continuation?.resume(returning: outcome) + } +} + +@available(macOS 14.0, *) +@MainActor +extension Capture { + static func windowCaptureStreamTarget( + pid: pid_t, + processIdentity: AXTreeProcessIdentity, + preferredWindowID: CGWindowID?, + scale: Double + ) -> WindowCaptureStreamTarget? { + guard processIdentity.isProven, + let candidate = bestWindow( + forPid: pid, + preferredWindowID: preferredWindowID + ) else { + return nil + } + let frame = candidate.frame + guard frame.width > 1, frame.height > 1 else { return nil } + + let outputScale = scale > 0 ? scale : 0.5 + let backingScale = backingScaleFactor(forWindowFrame: frame) + let width = max(1, Int(ceil(frame.width * backingScale * outputScale))) + let height = max(1, Int(ceil(frame.height * backingScale * outputScale))) + return WindowCaptureStreamTarget( + key: WindowCaptureStreamKey( + pid: pid, + processIdentity: processIdentity, + windowID: candidate.windowID, + pixelWidth: width, + pixelHeight: height + ), + originX: Double(frame.origin.x), + originY: Double(frame.origin.y), + pointWidth: Double(frame.width), + pointHeight: Double(frame.height) + ) + } + + static func windowShot( + from frame: WindowCaptureStreamFrame, + target: WindowCaptureStreamTarget + ) -> WindowShot? { + guard frame.width == target.key.pixelWidth, + frame.height == target.key.pixelHeight, + let image = image(from: frame), + let encoded = pngBase64WithSize(image) else { + return nil + } + return WindowShot( + base64: encoded.base64, + width: encoded.width, + height: encoded.height, + originX: target.originX, + originY: target.originY, + pointWidth: target.pointWidth, + pointHeight: target.pointHeight, + windowID: target.key.windowID, + source: .stream + ) + } + + private static func image(from frame: WindowCaptureStreamFrame) -> CGImage? { + let (minimumBytesPerRow, rowOverflow) = frame.width.multipliedReportingOverflow(by: 4) + let (minimumByteCount, countOverflow) = frame.bytesPerRow.multipliedReportingOverflow( + by: frame.height + ) + guard frame.width > 0, + frame.height > 0, + !rowOverflow, + !countOverflow, + frame.bytesPerRow >= minimumBytesPerRow, + frame.bytes.count >= minimumByteCount, + let provider = CGDataProvider(data: frame.bytes as CFData), + let colorSpace = CGColorSpace(name: CGColorSpace.sRGB) else { + return nil + } + let bitmapInfo = CGBitmapInfo.byteOrder32Little.union( + CGBitmapInfo(rawValue: CGImageAlphaInfo.premultipliedFirst.rawValue) + ) + return CGImage( + width: frame.width, + height: frame.height, + bitsPerComponent: 8, + bitsPerPixel: 32, + bytesPerRow: frame.bytesPerRow, + space: colorSpace, + bitmapInfo: bitmapInfo, + provider: provider, + decode: nil, + shouldInterpolate: true, + intent: .defaultIntent + ) + } +} diff --git a/native/cu-helper/Tests/CuHelperTests/CommandRouterSafetyTests.swift b/native/cu-helper/Tests/CuHelperTests/CommandRouterSafetyTests.swift index 4194abba..c12bf595 100644 --- a/native/cu-helper/Tests/CuHelperTests/CommandRouterSafetyTests.swift +++ b/native/cu-helper/Tests/CuHelperTests/CommandRouterSafetyTests.swift @@ -126,6 +126,36 @@ final class CommandRouterSafetyTests: XCTestCase { } } + func testSessionResetAlsoInvalidatesTheWindowCaptureProvider() { + let monitor = PhysicalInputEpochMonitor(counterReader: { _ in 0 }) + let provider = WindowCaptureProviderSpy() + let router = CommandRouter( + cursor: VirtualCursor(headless: true), + capabilities: Capabilities(headless: true), + inputMonitor: monitor, + windowCaptureProvider: provider + ) + + router.resetSessionState() + + XCTAssertEqual(provider.invalidateCount, 1) + } + + func testDiscardedWindowSnapshotForcesTheRetryToReturnAFullTree() { + XCTAssertFalse(CommandRouter.effectiveDisableDiff( + requested: false, + forceFullSnapshot: false + )) + XCTAssertTrue(CommandRouter.effectiveDisableDiff( + requested: false, + forceFullSnapshot: true + )) + XCTAssertTrue(CommandRouter.effectiveDisableDiff( + requested: true, + forceFullSnapshot: false + )) + } + func testCoordinateTransformNeverFallsBackAcrossApps() { CommandRouter.clearShotTransformsForTesting() CommandRouter.recordShotTransform( @@ -376,3 +406,22 @@ final class CommandRouterSafetyTests: XCTestCase { } } } + +@MainActor +private final class WindowCaptureProviderSpy: WindowCaptureProviding { + private(set) var invalidateCount = 0 + + func windowShot( + pid: pid_t, + processIdentity: AXTreeProcessIdentity, + preferredWindowID: CGWindowID?, + scale: Double, + newerThanUptime: TimeInterval? + ) async -> WindowShot? { + nil + } + + func invalidate() { + invalidateCount += 1 + } +} diff --git a/native/cu-helper/Tests/CuHelperTests/DaemonOverlayTargetTests.swift b/native/cu-helper/Tests/CuHelperTests/DaemonOverlayTargetTests.swift index daec068f..0c57dc7a 100644 --- a/native/cu-helper/Tests/CuHelperTests/DaemonOverlayTargetTests.swift +++ b/native/cu-helper/Tests/CuHelperTests/DaemonOverlayTargetTests.swift @@ -51,4 +51,23 @@ final class DaemonOverlayTargetTests: XCTestCase { currentIdentity: { _ in nil } )) } + + func testEveryOverlayStopAlsoInvalidatesTheLongLivedWindowStream() throws { + let sourceURL = URL(fileURLWithPath: #filePath) + .deletingLastPathComponent() + .deletingLastPathComponent() + .deletingLastPathComponent() + .appendingPathComponent("Sources/cu-helper/Daemon.swift") + let source = try String(contentsOf: sourceURL, encoding: .utf8) + let body = try XCTUnwrap( + source.range(of: "private func stopOverlaySession()").flatMap { start in + source.range( + of: "private func resolvedInjectionOverlayTarget", + range: start.upperBound..= 2 else { return } + source.publish( + makeFrame(for: target.key, sequence: 2, uptime: 20, byte: 7) + ) + } + let fresh = await manager.frame(for: target, newerThanUptime: 15) + + XCTAssertEqual(fresh?.sequence, 2) + XCTAssertEqual(fresh?.bytes, first?.bytes) + XCTAssertEqual(factory.sources.count, 1) + } + + func testPostMutationTimeoutNeverReturnsThePreMutationFrame() async { + let factory = FakeWindowCaptureStreamFactory { source, _ in + source.startFrame = makeFrame(for: source.targetKey, sequence: 1, uptime: 10, byte: 4) + } + let manager = makeManager(factory: factory) + let target = makeTarget(windowID: 12) + + let initial = await manager.frame(for: target, newerThanUptime: nil) + XCTAssertNotNil(initial) + let stale = await manager.frame(for: target, newerThanUptime: 11) + + XCTAssertNil(stale) + XCTAssertEqual(factory.sources.count, 2) + XCTAssertEqual(factory.sources[0].retireCount, 1) + XCTAssertEqual(factory.sources[1].retireCount, 0) + XCTAssertEqual(manager.activeKeyForTesting, target.key) + } + + func testPostMutationStarvationRebuildsOnceAndReturnsAStartedFrame() async { + let factory = FakeWindowCaptureStreamFactory { source, index in + source.startFrame = makeFrame( + for: source.targetKey, + sequence: 1, + uptime: index == 0 ? 10 : 20, + byte: UInt8(index + 1) + ) + } + let manager = makeManager(factory: factory) + let target = makeTarget(windowID: 13) + + _ = await manager.frame(for: target, newerThanUptime: nil) + let recovered = await manager.frame(for: target, newerThanUptime: 15) + + XCTAssertEqual(recovered?.receivedUptime, 20) + XCTAssertEqual(recovered?.bytes.first, 2) + XCTAssertEqual(factory.sources.count, 2) + XCTAssertEqual(factory.sources[0].retireCount, 1) + XCTAssertEqual(factory.sources[1].retireCount, 0) + } + + func testWindowSwitchRetiresOldStreamAndLateFramesCannotLeak() async { + let factory = FakeWindowCaptureStreamFactory { source, index in + source.startFrame = makeFrame( + for: source.targetKey, + sequence: 1, + uptime: 10 + Double(index), + byte: UInt8(index + 1) + ) + } + let manager = makeManager(factory: factory) + let targetA = makeTarget(windowID: 21) + let targetB = makeTarget(windowID: 22) + + let frameA = await manager.frame(for: targetA, newerThanUptime: nil) + let generationA = manager.activeGenerationForTesting + let frameB = await manager.frame(for: targetB, newerThanUptime: nil) + let generationB = manager.activeGenerationForTesting + + XCTAssertEqual(frameA?.bytes.first, 1) + XCTAssertEqual(frameB?.bytes.first, 2) + XCTAssertNotEqual(generationA, generationB) + XCTAssertEqual(factory.sources.count, 2) + XCTAssertEqual(factory.sources[0].retireCount, 1) + + factory.sources[0].publish( + makeFrame(for: targetA.key, sequence: 99, uptime: 99, byte: 99) + ) + let stillB = await manager.frame(for: targetB, newerThanUptime: nil) + XCTAssertEqual(stillB?.bytes.first, 2) + XCTAssertEqual(manager.activeKeyForTesting, targetB.key) + } + + func testPIDReuseReplacesStreamEvenWhenPIDAndWindowIDMatch() async { + let factory = FakeWindowCaptureStreamFactory { source, index in + source.startFrame = makeFrame( + for: source.targetKey, + sequence: 1, + uptime: 10, + byte: UInt8(index + 1) + ) + } + let manager = makeManager(factory: factory) + let firstLifetime = makeTarget(windowID: 31, launchTime: 100) + let reusedPID = makeTarget(windowID: 31, launchTime: 200) + + _ = await manager.frame(for: firstLifetime, newerThanUptime: nil) + let replacement = await manager.frame(for: reusedPID, newerThanUptime: nil) + + XCTAssertEqual(replacement?.bytes.first, 2) + XCTAssertEqual(factory.sources.count, 2) + XCTAssertEqual(factory.sources[0].retireCount, 1) + XCTAssertEqual(manager.activeKeyForTesting?.processIdentity.launchTime, 200) + } + + func testResizeCannotReuseAFrameFromTheOldConfiguration() async { + let factory = FakeWindowCaptureStreamFactory { source, index in + source.startFrame = makeFrame( + for: source.targetKey, + sequence: 1, + uptime: 10, + byte: UInt8(index + 1) + ) + } + let manager = makeManager(factory: factory) + let before = makeTarget(windowID: 41, pixelWidth: 2, pixelHeight: 2) + let after = makeTarget(windowID: 41, pixelWidth: 4, pixelHeight: 3) + + _ = await manager.frame(for: before, newerThanUptime: nil) + let resized = await manager.frame(for: after, newerThanUptime: nil) + + XCTAssertEqual(resized?.width, 4) + XCTAssertEqual(resized?.height, 3) + XCTAssertEqual(resized?.bytes.first, 2) + XCTAssertEqual(factory.sources.count, 2) + XCTAssertEqual(factory.sources[0].retireCount, 1) + } + + func testMovingTheSameWindowReusesItsStream() async { + let factory = FakeWindowCaptureStreamFactory { source, _ in + source.startFrame = makeFrame( + for: source.targetKey, + sequence: 1, + uptime: 10, + byte: 5 + ) + } + let manager = makeManager(factory: factory) + let before = makeTarget(windowID: 42, originX: 100, originY: 200) + let after = makeTarget(windowID: 42, originX: 500, originY: 600) + + _ = await manager.frame(for: before, newerThanUptime: nil) + let moved = await manager.frame(for: after, newerThanUptime: nil) + + XCTAssertEqual(moved?.bytes.first, 5) + XCTAssertEqual(factory.sources.count, 1) + XCTAssertEqual(factory.sources[0].startCount, 1) + XCTAssertEqual(factory.sources[0].retireCount, 0) + XCTAssertEqual(manager.activeKeyForTesting, after.key) + } + + func testDelegateFailureRestartsOnceAndDropsTheFailedFrame() async { + let factory = FakeWindowCaptureStreamFactory { source, index in + source.startFrame = makeFrame( + for: source.targetKey, + sequence: 1, + uptime: 10, + byte: UInt8(index + 1) + ) + } + let manager = makeManager(factory: factory) + let target = makeTarget(windowID: 51) + + let first = await manager.frame(for: target, newerThanUptime: nil) + XCTAssertEqual(first?.bytes.first, 1) + factory.sources[0].failed = true + + let recovered = await manager.frame(for: target, newerThanUptime: nil) + + XCTAssertEqual(recovered?.bytes.first, 2) + XCTAssertEqual(factory.sources.count, 2) + XCTAssertEqual(factory.sources[0].retireCount, 1) + XCTAssertEqual(factory.sources[1].startCount, 1) + } + + func testOneBoundedRetryRecoversAStartFailure() async { + let factory = FakeWindowCaptureStreamFactory { source, index in + if index == 0 { + source.startError = CUError("capture_failed", "fixture start failure") + } else { + source.startFrame = makeFrame( + for: source.targetKey, + sequence: 1, + uptime: 10, + byte: 8 + ) + } + } + let manager = makeManager(factory: factory) + let target = makeTarget(windowID: 52) + + let recovered = await manager.frame(for: target, newerThanUptime: nil) + + XCTAssertEqual(recovered?.bytes.first, 8) + XCTAssertEqual(factory.sources.count, 2) + XCTAssertEqual(factory.sources[0].retireCount, 1) + XCTAssertEqual(factory.sources[1].startCount, 1) + } + + func testSessionInvalidationStopsTheStreamAndMakesLateFramesInert() async { + let factory = FakeWindowCaptureStreamFactory { source, index in + source.startFrame = makeFrame( + for: source.targetKey, + sequence: 1, + uptime: 10, + byte: UInt8(index + 1) + ) + } + let manager = makeManager(factory: factory) + let target = makeTarget(windowID: 61) + + _ = await manager.frame(for: target, newerThanUptime: nil) + let oldSource = factory.sources[0] + manager.invalidate() + oldSource.publish( + makeFrame(for: target.key, sequence: 100, uptime: 100, byte: 100) + ) + + XCTAssertNil(manager.activeKeyForTesting) + XCTAssertEqual(oldSource.retireCount, 1) + + let nextTurn = await manager.frame(for: target, newerThanUptime: nil) + XCTAssertEqual(nextTurn?.bytes.first, 2) + XCTAssertEqual(factory.sources.count, 2) + } + + func testOnlyPixelBearingFrameStatusesAreAccepted() { + XCTAssertTrue(WindowCaptureFrameStatusPolicy.accepts(.started)) + XCTAssertTrue(WindowCaptureFrameStatusPolicy.accepts(.complete)) + XCTAssertFalse(WindowCaptureFrameStatusPolicy.accepts(.idle)) + XCTAssertFalse(WindowCaptureFrameStatusPolicy.accepts(.blank)) + XCTAssertFalse(WindowCaptureFrameStatusPolicy.accepts(.suspended)) + XCTAssertFalse(WindowCaptureFrameStatusPolicy.accepts(.stopped)) + XCTAssertFalse(WindowCaptureFrameStatusPolicy.marksFailure(.started)) + XCTAssertFalse(WindowCaptureFrameStatusPolicy.marksFailure(.complete)) + XCTAssertFalse(WindowCaptureFrameStatusPolicy.marksFailure(.idle)) + XCTAssertFalse(WindowCaptureFrameStatusPolicy.marksFailure(.blank)) + XCTAssertFalse(WindowCaptureFrameStatusPolicy.marksFailure(.suspended)) + XCTAssertTrue(WindowCaptureFrameStatusPolicy.marksFailure(.stopped)) + } + + func testStreamConfigurationUsesCodexCadenceBufferingAndPixelFormat() { + let target = makeTarget(windowID: 71, pixelWidth: 640, pixelHeight: 480) + let configuration = ScreenCaptureKitWindowStreamSource.makeConfiguration( + for: target + ) + + XCTAssertEqual(configuration.width, 640) + XCTAssertEqual(configuration.height, 480) + XCTAssertEqual( + configuration.minimumFrameInterval, + CMTime(value: 1, timescale: 60) + ) + XCTAssertEqual(configuration.queueDepth, 5) + XCTAssertFalse(configuration.showsCursor) + XCTAssertFalse(configuration.capturesAudio) + XCTAssertTrue(configuration.scalesToFit) + XCTAssertTrue(configuration.preservesAspectRatio) + XCTAssertEqual(configuration.pixelFormat, kCVPixelFormatType_32BGRA) + XCTAssertEqual(configuration.colorSpaceName, CGColorSpace.sRGB) + XCTAssertTrue(configuration.ignoreShadowsSingleWindow) + XCTAssertTrue(configuration.ignoreShadowsDisplay) + } + + func testCopiedBGRAFrameEncodesAsAStreamPNGWithTheSameGeometry() throws { + let target = makeTarget( + windowID: 72, + pixelWidth: 1, + pixelHeight: 1, + originX: 300, + originY: 400 + ) + let frame = WindowCaptureStreamFrame( + bytes: Data([0x00, 0x00, 0xff, 0xff]), + width: 1, + height: 1, + bytesPerRow: 4, + sequence: 1, + receivedUptime: 10 + ) + + let shot = try XCTUnwrap(Capture.windowShot(from: frame, target: target)) + let png = try XCTUnwrap(Data(base64Encoded: shot.base64)) + + XCTAssertEqual(Array(png.prefix(8)), [137, 80, 78, 71, 13, 10, 26, 10]) + XCTAssertEqual(shot.width, 1) + XCTAssertEqual(shot.height, 1) + XCTAssertEqual(shot.originX, 300) + XCTAssertEqual(shot.originY, 400) + XCTAssertEqual(shot.windowID, 72) + XCTAssertEqual(shot.source, .stream) + } + + private func makeManager( + factory: FakeWindowCaptureStreamFactory + ) -> WindowCaptureStreamManager { + WindowCaptureStreamManager( + factory: factory, + frameWaitAttempts: 0, + frameWaitNanoseconds: 0 + ) + } + + private func makeTarget( + windowID: CGWindowID, + launchTime: TimeInterval = 100, + pixelWidth: Int = 2, + pixelHeight: Int = 2, + originX: Double = 100, + originY: Double = 200 + ) -> WindowCaptureStreamTarget { + WindowCaptureStreamTarget( + key: WindowCaptureStreamKey( + pid: 4242, + processIdentity: AXTreeProcessIdentity( + bundleID: "com.example.fixture", + executablePath: "/Applications/Fixture.app/Contents/MacOS/Fixture", + launchTime: launchTime + ), + windowID: windowID, + pixelWidth: pixelWidth, + pixelHeight: pixelHeight + ), + originX: originX, + originY: originY, + pointWidth: Double(pixelWidth), + pointHeight: Double(pixelHeight) + ) + } +} + +@MainActor +private final class FakeWindowCaptureStreamFactory: WindowCaptureStreamSourceFactory { + typealias Configure = (FakeWindowCaptureStreamSource, Int) -> Void + + private let configure: Configure + private(set) var sources: [FakeWindowCaptureStreamSource] = [] + + init(configure: @escaping Configure) { + self.configure = configure + } + + func makeSource( + for target: WindowCaptureStreamTarget + ) -> any WindowCaptureStreamSource { + let source = FakeWindowCaptureStreamSource(targetKey: target.key) + configure(source, sources.count) + sources.append(source) + return source + } +} + +@MainActor +private final class FakeWindowCaptureStreamSource: WindowCaptureStreamSource { + let targetKey: WindowCaptureStreamKey + var failed = false + var startError: CUError? + var startFrame: WindowCaptureStreamFrame? + var onLatestRead: ((FakeWindowCaptureStreamSource, Int) -> Void)? + private(set) var startCount = 0 + private(set) var retireCount = 0 + private(set) var latestReadCount = 0 + private var latest: WindowCaptureStreamFrame? + + init(targetKey: WindowCaptureStreamKey) { + self.targetKey = targetKey + } + + var hasFailed: Bool { failed } + + func start() async throws { + startCount += 1 + if let startError { throw startError } + latest = startFrame + } + + func latestFrame() -> WindowCaptureStreamFrame? { + latestReadCount += 1 + onLatestRead?(self, latestReadCount) + return latest + } + + func retire() { + retireCount += 1 + } + + func publish(_ frame: WindowCaptureStreamFrame) { + latest = frame + } +} + +private func makeFrame( + for key: WindowCaptureStreamKey, + sequence: UInt64, + uptime: TimeInterval, + byte: UInt8 +) -> WindowCaptureStreamFrame { + let bytesPerRow = key.pixelWidth * 4 + return WindowCaptureStreamFrame( + bytes: Data( + repeating: byte, + count: bytesPerRow * key.pixelHeight + ), + width: key.pixelWidth, + height: key.pixelHeight, + bytesPerRow: bytesPerRow, + sequence: sequence, + receivedUptime: uptime + ) +}