fix(computer-use): maintain long-lived window capture streams

This commit is contained in:
程序员阿江(Relakkes)
2026-08-31 11:25:38 +08:00
parent 6ace7f865b
commit 42d7cb1c93
10 changed files with 1584 additions and 69 deletions
+2
View File
@@ -49,6 +49,8 @@ let package = Package(
linkerSettings: [
.linkedFramework("AppKit"),
.linkedFramework("CoreGraphics"),
.linkedFramework("CoreMedia"),
.linkedFramework("CoreVideo"),
.linkedFramework("QuartzCore"),
.linkedFramework("ApplicationServices"),
.linkedFramework("ScreenCaptureKit"),
@@ -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)
}
@@ -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
@@ -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()
}
@@ -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. \
@@ -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<Outcome, Never>?
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<Outcome, Never>?
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
)
}
}
@@ -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
}
}
@@ -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..<source.endIndex
).map { end in String(source[start.lowerBound..<end.lowerBound]) }
}
)
XCTAssertTrue(body.contains("router.invalidateWindowCaptureStream()"))
}
}
@@ -73,6 +73,47 @@ final class TargetVisibilityPolicyTests: XCTestCase {
return try String(contentsOf: root.appendingPathComponent(name), encoding: .utf8)
}
func testCoveredWindowWithStreamProviderNeverUsesOneShotFallback() {
XCTAssertFalse(TargetVisibilityPolicy.permitsOneShotFallback(
windowIsCovered: true,
streamProviderInstalled: true
))
}
func testCaptureTargetMustRemainTheSameKeyWindow() {
XCTAssertTrue(TargetVisibilityPolicy.captureTargetStillMatches(
snapshotWindowID: 100,
currentWindowID: 100
))
XCTAssertTrue(TargetVisibilityPolicy.captureTargetStillMatches(
snapshotWindowID: nil,
currentWindowID: nil
))
XCTAssertFalse(TargetVisibilityPolicy.captureTargetStillMatches(
snapshotWindowID: 100,
currentWindowID: 101
))
XCTAssertFalse(TargetVisibilityPolicy.captureTargetStillMatches(
snapshotWindowID: 100,
currentWindowID: nil
))
XCTAssertFalse(TargetVisibilityPolicy.captureTargetStillMatches(
snapshotWindowID: nil,
currentWindowID: 100
))
}
func testVisibleOrLegacyWindowCanUseOneShotFallback() {
XCTAssertTrue(TargetVisibilityPolicy.permitsOneShotFallback(
windowIsCovered: false,
streamProviderInstalled: true
))
XCTAssertTrue(TargetVisibilityPolicy.permitsOneShotFallback(
windowIsCovered: true,
streamProviderInstalled: false
))
}
/// Regression for session ba87fe5e-a2dc-4d93-932f-4c868df4c05e.
///
/// `withForegroundLease` used to reject every mutation when another app
@@ -114,15 +155,30 @@ final class TargetVisibilityPolicyTests: XCTestCase {
XCTAssertTrue(body.contains("coveredCaptureNotice"))
}
func testFirstCoveredCaptureWarnsWithoutBlockingInput() {
let notice = TargetVisibilityPolicy.coveredCaptureNotice
XCTAssertTrue(notice.contains("may be stale"))
func testCoveredLiveStreamExplainsFreshFramesWithoutBlockingInput() {
let notice = TargetVisibilityPolicy.coveredCaptureNotice(
liveStreamActive: true
)
XCTAssertTrue(notice.contains("long-lived window stream"))
XCTAssertTrue(notice.contains("latest complete frame"))
XCTAssertTrue(notice.contains("does not block"))
XCTAssertTrue(notice.contains("continue the task"))
XCTAssertFalse(notice.contains("may be stale"))
XCTAssertFalse(notice.contains("paused its renderer"))
XCTAssertFalse(notice.contains("refused"))
XCTAssertFalse(notice.contains("Uncover the window"))
}
func testCoveredStreamFailureOmitsPotentiallyStaleOneShotPixels() {
let notice = TargetVisibilityPolicy.coveredCaptureNotice(
liveStreamActive: false
)
XCTAssertTrue(notice.contains("live window-stream frame was unavailable"))
XCTAssertTrue(notice.contains("one-shot screenshot is intentionally not used"))
XCTAssertTrue(notice.contains("compositor-cached pixels"))
XCTAssertTrue(notice.contains("accessibility state"))
XCTAssertTrue(notice.contains("does not block"))
}
/// An explicit AX Raise remains available, but it must not activate the app.
func testUncoveringNeverActivatesTheApplication() throws {
let source = try source("AXAction.swift")
@@ -144,12 +200,17 @@ final class TargetVisibilityPolicyTests: XCTestCase {
func testTheIdenticalCaptureNoticeNamesTheToggleTrap() {
for covered in [true, false] {
let notice = TargetVisibilityPolicy.identicalCaptureNotice(windowIsCovered: covered)
XCTAssertTrue(notice.contains("byte-for-byte identical"))
XCTAssertTrue(notice.contains("toggle"))
XCTAssertTrue(notice.contains("undo the first"))
XCTAssertFalse(notice.contains("refused"))
XCTAssertFalse(notice.contains("Uncover the window"))
for liveStreamActive in [true, false] {
let notice = TargetVisibilityPolicy.identicalCaptureNotice(
windowIsCovered: covered,
liveStreamActive: liveStreamActive
)
XCTAssertTrue(notice.contains("byte-for-byte identical"))
XCTAssertTrue(notice.contains("toggle"))
XCTAssertTrue(notice.contains("undo the first"))
XCTAssertFalse(notice.contains("refused"))
XCTAssertFalse(notice.contains("Uncover the window"))
}
}
}
@@ -159,15 +220,23 @@ final class TargetVisibilityPolicyTests: XCTestCase {
// time, and the model spent four minutes chasing the possibility we had
// handed it. Coverage is something we compute, so it must not be
// presented to the model as an open question.
let visible = TargetVisibilityPolicy.identicalCaptureNotice(windowIsCovered: false)
let visible = TargetVisibilityPolicy.identicalCaptureNotice(
windowIsCovered: false,
liveStreamActive: true
)
XCTAssertTrue(visible.contains("no visible pixel change"))
XCTAssertTrue(visible.contains("not fully covered"))
XCTAssertFalse(visible.contains("repainting"))
XCTAssertFalse(visible.contains("paused its renderer"))
let covered = TargetVisibilityPolicy.identicalCaptureNotice(windowIsCovered: true)
XCTAssertTrue(covered.contains("paused its renderer"))
let covered = TargetVisibilityPolicy.identicalCaptureNotice(
windowIsCovered: true,
liveStreamActive: true
)
XCTAssertTrue(covered.contains("long-lived window stream"))
XCTAssertTrue(covered.contains("newest complete frame"))
XCTAssertTrue(covered.contains("does not block"))
XCTAssertFalse(covered.contains("no visible pixel change"))
XCTAssertTrue(covered.contains("no visible pixel change"))
XCTAssertFalse(covered.contains("paused its renderer"))
}
}
@@ -0,0 +1,457 @@
import CoreMedia
import ScreenCaptureKit
import XCTest
@testable import cc_haha_computer_use
@MainActor
final class WindowCaptureStreamTests: XCTestCase {
func testSameTargetReusesStreamAndReturnsNewestFrameAcrossCoveredAction() async {
let factory = FakeWindowCaptureStreamFactory { source, _ in
source.startFrame = makeFrame(for: source.targetKey, sequence: 1, uptime: 10, byte: 1)
}
let manager = makeManager(factory: factory)
let target = makeTarget(windowID: 10)
let first = await manager.frame(for: target, newerThanUptime: nil)
XCTAssertEqual(first?.sequence, 1)
XCTAssertEqual(factory.sources.count, 1)
factory.sources[0].publish(
makeFrame(for: target.key, sequence: 3, uptime: 12, byte: 3)
)
let afterCoveredAction = await manager.frame(
for: target,
newerThanUptime: 11
)
XCTAssertEqual(afterCoveredAction?.sequence, 3)
XCTAssertEqual(afterCoveredAction?.bytes.first, 3)
XCTAssertEqual(factory.sources.count, 1)
XCTAssertEqual(factory.sources[0].startCount, 1)
XCTAssertEqual(factory.sources[0].retireCount, 0)
}
func testIdenticalPixelsWithANewerSequenceSatisfyFreshness() async {
let factory = FakeWindowCaptureStreamFactory { source, _ in
source.startFrame = makeFrame(for: source.targetKey, sequence: 1, uptime: 10, byte: 7)
}
let manager = makeManager(factory: factory)
let target = makeTarget(windowID: 11)
let first = await manager.frame(for: target, newerThanUptime: nil)
XCTAssertEqual(first?.sequence, 1)
let source = factory.sources[0]
source.onLatestRead = { source, readCount in
guard readCount >= 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
)
}