mirror of
https://github.com/saphid/frame-control.git
synced 2026-10-06 04:04:21 +02:00
Mac in the headset: measure every frame, adapt to the network, separate windows
- Per-frame timing on the Mac's clock (capture, encode, network, decode, draw), viewer clock sync and reports, input echo, /stats and a HUD. - scripts/macview-bench.py: repeatable runs on the real Frame, a shaping relay (no sudo), interleaved A/B between agent settings; results in bench/results/. - Adaptive controller: ack-based send gate with jitter-aware slack, AIMD bitrate that knows when a stream is app-limited, fps then size tiers. On a 50->3->50 Mbit/s step, scroll p95 went from 4.7 s to 72 ms; no cost on a clean link. - Separate mode: real AppKit event loop (HiDPI and NSScreen now work), cropped capture for fixed-size windows, windows kept on their display, graceful quit restores windows; stop/start races fixed. - Encoder timeline clamp (no oversized frame after a pause). - Frame Control shows each live stream's fps, delay, bitrate and tier. Reviewed by GPT-6 Astra xhigh (read-only), 7 rounds; findings fixed. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
1 parent
a3c6e5c984
commit
9e4dcdd147
29 files changed
+21405
-104
No files matched your search
@@ -0,0 +1,28 @@
|
||||
// Private CoreGraphics classes for virtual displays (as used by DeskPad,
|
||||
// BetterDisplay and quest-display). Not in any SDK; see
|
||||
// docs/mac-window-separation.md.
|
||||
#import <Foundation/Foundation.h>
|
||||
#import <CoreGraphics/CoreGraphics.h>
|
||||
@interface CGVirtualDisplayDescriptor : NSObject
|
||||
@property (retain, nonatomic) dispatch_queue_t queue;
|
||||
@property (retain, nonatomic) NSString *name;
|
||||
@property (nonatomic) unsigned int maxPixelsHigh;
|
||||
@property (nonatomic) unsigned int maxPixelsWide;
|
||||
@property (nonatomic) CGSize sizeInMillimeters;
|
||||
@property (nonatomic) unsigned int serialNum;
|
||||
@property (nonatomic) unsigned int productID;
|
||||
@property (nonatomic) unsigned int vendorID;
|
||||
@property (copy, nonatomic) void (^terminationHandler)(id, id);
|
||||
@end
|
||||
@interface CGVirtualDisplayMode : NSObject
|
||||
- (instancetype)initWithWidth:(unsigned int)width height:(unsigned int)height refreshRate:(double)refreshRate;
|
||||
@end
|
||||
@interface CGVirtualDisplaySettings : NSObject
|
||||
@property (retain, nonatomic) NSArray *modes;
|
||||
@property (nonatomic) unsigned int hiDPI;
|
||||
@end
|
||||
@interface CGVirtualDisplay : NSObject
|
||||
@property (readonly, nonatomic) unsigned int displayID;
|
||||
- (instancetype)initWithDescriptor:(CGVirtualDisplayDescriptor *)descriptor;
|
||||
- (BOOL)applySettings:(CGVirtualDisplaySettings *)settings;
|
||||
@end
|
||||
@@ -10,6 +10,7 @@ import ScreenCaptureKit
|
||||
enum Source: Equatable {
|
||||
case window(CGWindowID)
|
||||
case display(CGDirectDisplayID)
|
||||
case separate(CGWindowID) // the window on a display of its own
|
||||
case test
|
||||
|
||||
init?(_ s: String) {
|
||||
@@ -17,6 +18,7 @@ enum Source: Equatable {
|
||||
switch (parts.first, parts.count > 1 ? UInt32(parts[1]) : nil) {
|
||||
case ("window", let id?): self = .window(id)
|
||||
case ("display", let id?): self = .display(id)
|
||||
case ("separate", let id?): self = .separate(id)
|
||||
case ("test", _): self = .test
|
||||
default: return nil
|
||||
}
|
||||
@@ -26,6 +28,7 @@ enum Source: Equatable {
|
||||
switch self {
|
||||
case .window(let id): return "window:\(id)"
|
||||
case .display(let id): return "display:\(id)"
|
||||
case .separate(let id): return "separate:\(id)"
|
||||
case .test: return "test"
|
||||
}
|
||||
}
|
||||
@@ -76,14 +79,20 @@ struct WindowInfo {
|
||||
}
|
||||
|
||||
protocol CaptureSource: AnyObject {
|
||||
var onFrame: ((CVPixelBuffer, CMTime) -> Void)? { get set }
|
||||
/// The picture, its presentation time, and when the Mac composited it (µs, host clock).
|
||||
var onFrame: ((CVPixelBuffer, CMTime, Int64) -> Void)? { get set }
|
||||
var onEnded: ((String) -> Void)? { get set }
|
||||
var onChange: (() -> Void)? { get set } // title or size changed
|
||||
/// True once the window or display is gone for good.
|
||||
var gone: Bool { get }
|
||||
/// Points on the Mac's global display space that the picture covers.
|
||||
var frameRect: CGRect { get }
|
||||
var title: String { get }
|
||||
var app: String { get }
|
||||
var pid: pid_t? { get }
|
||||
func start(maxLong: Int, fps: Int, completion: @escaping (String?) -> Void)
|
||||
/// A smaller or larger picture from now on (the long side, in pixels).
|
||||
func setMaxLong(_ maxLong: Int)
|
||||
func stop()
|
||||
func pointer(x: Double, y: Double, text: String?) // for the test pattern
|
||||
}
|
||||
@@ -95,7 +104,7 @@ extension CaptureSource {
|
||||
/// ScreenCaptureKit, for a window or a display.
|
||||
final class SCKSource: NSObject, CaptureSource, SCStreamOutput, SCStreamDelegate {
|
||||
let source: Source
|
||||
var onFrame: ((CVPixelBuffer, CMTime) -> Void)?
|
||||
var onFrame: ((CVPixelBuffer, CMTime, Int64) -> Void)?
|
||||
var onEnded: ((String) -> Void)?
|
||||
private(set) var frameRect = CGRect.zero
|
||||
private(set) var title = ""
|
||||
@@ -113,9 +122,19 @@ final class SCKSource: NSObject, CaptureSource, SCStreamOutput, SCStreamDelegate
|
||||
/// The window or display no longer exists, so retrying can't help.
|
||||
private(set) var gone = false
|
||||
var onChange: (() -> Void)? // title or size changed
|
||||
/// For a display: capture only this part of it (display points, top-left
|
||||
/// origin). Set before start().
|
||||
var crop: CGRect?
|
||||
|
||||
init(_ source: Source) { self.source = source }
|
||||
|
||||
/// What's captured, in global points: the display, or the cropped part of it.
|
||||
private func displayRect(_ id: CGDirectDisplayID) -> CGRect {
|
||||
let b = CGDisplayBounds(id)
|
||||
guard let c = crop else { return b }
|
||||
return c.offsetBy(dx: b.minX, dy: b.minY).intersection(b)
|
||||
}
|
||||
|
||||
func start(maxLong: Int, fps: Int, completion: @escaping (String?) -> Void) {
|
||||
self.maxLong = maxLong
|
||||
SCShareableContent.getExcludingDesktopWindows(true, onScreenWindowsOnly: false) { [self] content, error in
|
||||
@@ -148,11 +167,17 @@ final class SCKSource: NSObject, CaptureSource, SCStreamOutput, SCStreamDelegate
|
||||
filter = SCContentFilter(display: d, excludingWindows: [])
|
||||
title = displayName(id)
|
||||
app = "Mac"
|
||||
frameRect = CGDisplayBounds(id)
|
||||
case .test:
|
||||
frameRect = displayRect(id)
|
||||
if crop != nil { config.sourceRect = frameRect.offsetBy(dx: -CGDisplayBounds(id).minX, dy: -CGDisplayBounds(id).minY) }
|
||||
case .test, .separate:
|
||||
return completion("not a ScreenCaptureKit source")
|
||||
}
|
||||
scale = Double(filter.pointPixelScale)
|
||||
// A display's own mode knows best: ScreenCaptureKit can still say 1
|
||||
// for a virtual display that has only just switched to HiDPI.
|
||||
if case .display(let id) = source, let mode = CGDisplayCopyDisplayMode(id), mode.width > 0 {
|
||||
scale = max(scale, Double(mode.pixelWidth) / Double(mode.width))
|
||||
}
|
||||
let (w, h) = fitSize(width: frameRect.width * scale, height: frameRect.height * scale, maxLong: maxLong)
|
||||
config.width = w
|
||||
config.height = h
|
||||
@@ -195,8 +220,10 @@ final class SCKSource: NSObject, CaptureSource, SCStreamOutput, SCStreamDelegate
|
||||
}
|
||||
}
|
||||
|
||||
/// Windows move, resize, retitle and close; ScreenCaptureKit doesn't say.
|
||||
/// Windows move, resize, retitle and close, and displays change mode;
|
||||
/// ScreenCaptureKit doesn't say.
|
||||
private func startPolling() {
|
||||
if case .display(let id) = source { return pollDisplay(id) }
|
||||
guard case .window(let id) = source else { return }
|
||||
let t = DispatchSource.makeTimerSource(queue: queue)
|
||||
t.schedule(deadline: .now() + 1, repeating: 1)
|
||||
@@ -224,14 +251,48 @@ final class SCKSource: NSObject, CaptureSource, SCStreamOutput, SCStreamDelegate
|
||||
poll = t
|
||||
}
|
||||
|
||||
func setMaxLong(_ n: Int) {
|
||||
queue.async {
|
||||
self.maxLong = n
|
||||
guard self.stream != nil else { return }
|
||||
let (w, h) = fitSize(width: self.frameRect.width * self.scale, height: self.frameRect.height * self.scale, maxLong: n)
|
||||
guard w != self.config.width || h != self.config.height else { return }
|
||||
self.config.width = w
|
||||
self.config.height = h
|
||||
self.stream?.updateConfiguration(self.config) { _ in }
|
||||
}
|
||||
}
|
||||
|
||||
private func pollDisplay(_ id: CGDirectDisplayID) {
|
||||
let t = DispatchSource.makeTimerSource(queue: queue)
|
||||
t.schedule(deadline: .now() + 0.5, repeating: 1)
|
||||
t.setEventHandler { [weak self] in
|
||||
guard let self, let mode = CGDisplayCopyDisplayMode(id), mode.width > 0 else { return }
|
||||
let bounds = self.displayRect(id), scale = Double(mode.pixelWidth) / Double(mode.width)
|
||||
let (w, h) = fitSize(width: bounds.width * scale, height: bounds.height * scale, maxLong: self.maxLong)
|
||||
self.frameRect = bounds
|
||||
guard w != self.config.width || h != self.config.height else { return }
|
||||
self.scale = scale
|
||||
self.config.width = w
|
||||
self.config.height = h
|
||||
self.stream?.updateConfiguration(self.config) { _ in }
|
||||
self.onChange?()
|
||||
}
|
||||
t.resume()
|
||||
poll = t
|
||||
}
|
||||
|
||||
func stream(_ stream: SCStream, didOutputSampleBuffer sample: CMSampleBuffer, of type: SCStreamOutputType) {
|
||||
guard type == .screen, sample.isValid, let pb = CMSampleBufferGetImageBuffer(sample) else { return }
|
||||
// Only complete frames carry new pixels; idle and blank ones don't.
|
||||
if let atts = CMSampleBufferGetSampleAttachmentsArray(sample, createIfNecessary: false) as? [[SCStreamFrameInfo: Any]],
|
||||
let raw = atts.first?[.status] as? Int, let status = SCFrameStatus(rawValue: raw), status != .complete {
|
||||
let info = (CMSampleBufferGetSampleAttachmentsArray(sample, createIfNecessary: false) as? [[SCStreamFrameInfo: Any]])?.first
|
||||
if let raw = info?[.status] as? Int, let status = SCFrameStatus(rawValue: raw), status != .complete {
|
||||
return
|
||||
}
|
||||
onFrame?(pb, CMSampleBufferGetPresentationTimeStamp(sample))
|
||||
let pts = CMSampleBufferGetPresentationTimeStamp(sample)
|
||||
// When WindowServer composited it: the moment it showed on the Mac.
|
||||
let shown = (info?[.displayTime] as? UInt64).map(hostUs) ?? hostUs(pts)
|
||||
onFrame?(pb, pts, shown)
|
||||
}
|
||||
|
||||
func stream(_ stream: SCStream, didStopWithError error: Error) {
|
||||
@@ -251,8 +312,10 @@ func displayName(_ id: CGDirectDisplayID) -> String {
|
||||
/// A moving test card: bars, a clock and a frame counter, plus a dot where the
|
||||
/// viewer's pointer is and the last key it sent, so input can be checked too.
|
||||
final class TestSource: CaptureSource {
|
||||
var onFrame: ((CVPixelBuffer, CMTime) -> Void)?
|
||||
var onFrame: ((CVPixelBuffer, CMTime, Int64) -> Void)?
|
||||
var onEnded: ((String) -> Void)?
|
||||
var onChange: (() -> Void)?
|
||||
let gone = false
|
||||
let frameRect = CGRect(x: 0, y: 0, width: 1280, height: 720)
|
||||
let title = "Test pattern"
|
||||
let app = "Frame Control"
|
||||
@@ -267,10 +330,7 @@ final class TestSource: CaptureSource {
|
||||
|
||||
func start(maxLong: Int, fps: Int, completion: @escaping (String?) -> Void) {
|
||||
size = fitSize(width: 1280, height: 720, maxLong: maxLong)
|
||||
let attrs: [CFString: Any] = [kCVPixelBufferPixelFormatTypeKey: kCVPixelFormatType_32BGRA,
|
||||
kCVPixelBufferWidthKey: size.0, kCVPixelBufferHeightKey: size.1,
|
||||
kCVPixelBufferIOSurfacePropertiesKey: [:] as CFDictionary]
|
||||
CVPixelBufferPoolCreate(nil, nil, attrs as CFDictionary, &pool)
|
||||
makePool()
|
||||
let t = DispatchSource.makeTimerSource(queue: queue)
|
||||
t.schedule(deadline: .now(), repeating: 1.0 / Double(fps))
|
||||
t.setEventHandler { [weak self] in self?.draw() }
|
||||
@@ -284,6 +344,21 @@ final class TestSource: CaptureSource {
|
||||
timer = nil
|
||||
}
|
||||
|
||||
func setMaxLong(_ n: Int) {
|
||||
queue.async {
|
||||
self.size = fitSize(width: 1280, height: 720, maxLong: n)
|
||||
self.makePool()
|
||||
}
|
||||
}
|
||||
|
||||
private func makePool() {
|
||||
let attrs: [CFString: Any] = [kCVPixelBufferPixelFormatTypeKey: kCVPixelFormatType_32BGRA,
|
||||
kCVPixelBufferWidthKey: size.0, kCVPixelBufferHeightKey: size.1,
|
||||
kCVPixelBufferIOSurfacePropertiesKey: [:] as CFDictionary]
|
||||
pool = nil
|
||||
CVPixelBufferPoolCreate(nil, nil, attrs as CFDictionary, &pool)
|
||||
}
|
||||
|
||||
func pointer(x: Double, y: Double, text: String?) {
|
||||
queue.async {
|
||||
if x >= 0 { self.dot = (x, y) }
|
||||
@@ -334,6 +409,6 @@ final class TestSource: CaptureSource {
|
||||
ctx.fillEllipse(in: CGRect(x: CGFloat(px) * CGFloat(w) - r, y: (1 - CGFloat(py)) * CGFloat(h) - r, width: 2 * r, height: 2 * r))
|
||||
}
|
||||
n += 1
|
||||
onFrame?(pb, CMClockGetTime(CMClockGetHostTimeClock()))
|
||||
onFrame?(pb, CMClockGetTime(CMClockGetHostTimeClock()), nowUs())
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,234 @@
|
||||
// Adapts each stream to its network, latency first. The viewer acknowledges
|
||||
// every frame as it arrives ("rx"); from those acks the controller knows how
|
||||
// long frames take to get through and how fast the link delivers them.
|
||||
//
|
||||
// - The gate: a new frame is sent only while the oldest unacknowledged one is
|
||||
// younger than the path's usual round trip plus a little slack. So frames
|
||||
// never queue up in SSH, TCP or the Wi-Fi driver; while the link is stuck
|
||||
// the newest picture waits and goes out as soon as it moves again.
|
||||
// - The bitrate: when frames start queueing (the round trip grows) or the
|
||||
// gate has to hold frames back, it drops to a bit under what the link
|
||||
// actually delivered; once things are clear it probes up again slowly, never
|
||||
// above the quality setting's bitrate (the ceiling).
|
||||
// - The tier: as the bitrate falls, fewer frames per second (60, 45, 30), then
|
||||
// a smaller picture (75%, then 50% of the panel's pixels).
|
||||
// See docs/mac-in-headset.md ("Adapting to the network"). Thread-safe.
|
||||
import Foundation
|
||||
|
||||
final class RateController {
|
||||
struct Tier: Equatable {
|
||||
let fps: Int
|
||||
let scale: Double // of the picture's long side
|
||||
}
|
||||
|
||||
static let tiers = [Tier(fps: 60, scale: 1), Tier(fps: 45, scale: 1), Tier(fps: 30, scale: 1),
|
||||
Tier(fps: 30, scale: 0.75), Tier(fps: 30, scale: 0.5)]
|
||||
/// A tier is used while the target bitrate is at least this share of the ceiling.
|
||||
static let floors = [0.45, 0.28, 0.16, 0.08, 0]
|
||||
static let enabled = ProcessInfo.processInfo.environment["FRAME_MAC_VIEW_ADAPT"] != "0"
|
||||
|
||||
let maxFps: Int
|
||||
private let lock = NSLock()
|
||||
private var ceiling = 0 // bits/s at full size and frame rate
|
||||
private(set) var target = 0
|
||||
private(set) var tier = 0
|
||||
private var unacked: [(seq: UInt32, sent: Int64, bytes: Int)] = []
|
||||
/// Round trips (send -> ack arrives here), for the baseline: the lowest
|
||||
/// in the last 10 s is the path without any queue.
|
||||
private var rtts: [(t: Int64, v: Int64)] = []
|
||||
private var acked: [(t: Int64, bytes: Int)] = [] // the last second
|
||||
private var sentLog: [(t: Int64, bytes: Int)] = []
|
||||
private var captures: [Int64] = [] // the last second
|
||||
private var frameBytes = 0 // average recent frame, kept while the gate holds everything back
|
||||
private var held = 0 // frames the gate held back since the last update
|
||||
private var lastSignal = false
|
||||
private var sawAck = false
|
||||
private var lastDecrease: Int64 = 0
|
||||
private var lastIncrease: Int64 = 0
|
||||
private var belowSince: Int64 = 0, aboveSince: Int64 = 0
|
||||
/// What changed, for the timeline: (time, event).
|
||||
private(set) var events: [(Int64, String)] = []
|
||||
|
||||
init(maxFps: Int) { self.maxFps = maxFps }
|
||||
|
||||
private func locked<T>(_ f: () -> T) -> T { lock.lock(); defer { lock.unlock() }; return f() }
|
||||
|
||||
/// The quality setting's bitrate at full size; the first call also starts there.
|
||||
func setCeiling(_ bps: Int) {
|
||||
locked {
|
||||
if target == 0 || target > bps { target = bps }
|
||||
ceiling = bps
|
||||
}
|
||||
}
|
||||
|
||||
var fps: Int { locked { fpsLocked } }
|
||||
private var fpsLocked: Int { min(maxFps, RateController.tiers[tier].fps) }
|
||||
var scale: Double { locked { RateController.tiers[tier].scale } }
|
||||
var baseRtt: Int64 { locked { baseline() } }
|
||||
|
||||
private func baseline() -> Int64 { rtts.map(\.v).min() ?? 0 }
|
||||
|
||||
/// How late a frame may be before the gate holds the next one: one frame
|
||||
/// interval, plus room for the jitter this link normally has (1.5 times
|
||||
/// its recent spread), so ordinary Wi-Fi jitter doesn't cost frames but a
|
||||
/// real queue does. Updated in update().
|
||||
private var slack: Int64 = 40_000
|
||||
|
||||
/// Whether a frame may be sent now without queueing behind earlier ones.
|
||||
/// `counts`: a held frame is a sign of congestion (not when merely
|
||||
/// re-checking whether a held frame can go yet).
|
||||
func maySend(now: Int64, counts: Bool = true) -> Bool {
|
||||
locked {
|
||||
guard RateController.enabled, sawAck else { return true } // not heard from the viewer yet
|
||||
// Unacknowledged for 2 s: gone with a reconnection, not in a queue.
|
||||
unacked.removeAll { now - $0.sent > 2_000_000 }
|
||||
guard let oldest = unacked.first else { return true }
|
||||
// Age is what bounds latency. The count only stops a burst, and it
|
||||
// allows a full round trip of frames, so a long but clear path
|
||||
// (100 ms away) still gets every frame.
|
||||
let interval = Int64(1_000_000 / max(1, fpsLocked))
|
||||
let window = max(3, Int((baseline() + slack) / interval) + 1)
|
||||
if unacked.count < window, now - oldest.sent <= baseline() + slack { return true }
|
||||
if counts { held += 1 }
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
/// A picture was captured (sent or not): with the frame sizes, what this
|
||||
/// stream would send if the link allowed.
|
||||
func captured(at t: Int64) {
|
||||
locked {
|
||||
captures.append(t)
|
||||
if captures.count > 256 { captures.removeFirst(captures.count - 256) }
|
||||
}
|
||||
}
|
||||
|
||||
func sent(seq: UInt32, bytes: Int, at t: Int64) {
|
||||
locked {
|
||||
// Bounded even if the viewer never acknowledges (an old viewer, or
|
||||
// the controller is off).
|
||||
unacked.append((seq, t, bytes))
|
||||
if unacked.count > 512 { unacked.removeFirst(unacked.count - 512) }
|
||||
sentLog.append((t, bytes))
|
||||
if sentLog.count > 1024 { sentLog.removeFirst(sentLog.count - 1024) }
|
||||
}
|
||||
}
|
||||
|
||||
/// The viewer has frame `seq`. Returns true if that may let a held frame go.
|
||||
func acked(seq: UInt32, at now: Int64) -> Bool {
|
||||
locked {
|
||||
sawAck = true
|
||||
guard let i = unacked.firstIndex(where: { $0.seq == seq }) else { return false }
|
||||
let f = unacked[i]
|
||||
unacked.removeSubrange(0...i) // TCP delivers in order: earlier ones arrived too
|
||||
rtts.append((now, now - f.sent))
|
||||
acked.append((now, f.bytes))
|
||||
return true
|
||||
}
|
||||
}
|
||||
|
||||
/// Called every 100 ms. Returns the new bitrate target, or nil if the
|
||||
/// controller is off.
|
||||
func update(now: Int64) -> Int? {
|
||||
locked {
|
||||
rtts.removeAll { now - $0.t > 10_000_000 }
|
||||
acked.removeAll { now - $0.t > 500_000 }
|
||||
sentLog.removeAll { now - $0.t > 500_000 }
|
||||
unacked.removeAll { now - $0.sent > 2_000_000 }
|
||||
captures.removeAll { now - $0 > 1_000_000 }
|
||||
guard RateController.enabled, ceiling > 0 else { return nil }
|
||||
let base = baseline()
|
||||
let spread = rtts.filter { now - $0.t < 2_000_000 }.map(\.v).sorted()
|
||||
let jitter = spread.isEmpty ? 0 : spread[spread.count * 9 / 10] - base
|
||||
let interval = Int64(1_000_000 / max(1, fpsLocked))
|
||||
slack = interval + min(max(jitter * 3 / 2, 25_000), 80_000)
|
||||
let recent = rtts.filter { now - $0.t < 300_000 }.map(\.v).sorted()
|
||||
let queueing = recent.isEmpty ? 0 : recent[recent.count / 2] - base
|
||||
let oldestAge = unacked.first.map { now - $0.sent } ?? 0
|
||||
let stuck = oldestAge > base + 100_000
|
||||
let delivered = acked.reduce(0) { $0 + $1.bytes } * 16 // bits/s over the last half second
|
||||
let sending = sentLog.reduce(0) { $0 + $1.bytes } * 16
|
||||
// Demand: captures per second (up to the tier's rate) times the
|
||||
// average frame. A test card or a mostly still window wants far
|
||||
// less than its budget; when the link hiccups, cutting its bitrate
|
||||
// can't help, and it would only look link-limited afterwards.
|
||||
if !sentLog.isEmpty { frameBytes = sentLog.reduce(0) { $0 + $1.bytes } / sentLog.count }
|
||||
let demand = min(captures.count, fpsLocked) * frameBytes * 8
|
||||
// Delay alone isn't our queue: Wi-Fi jitters by itself. It only
|
||||
// counts while this stream uses a good part of its budget (so its
|
||||
// own data could be what's queueing). Frames the gate had to hold,
|
||||
// or one stuck in flight, show demand the link isn't carrying
|
||||
// whatever was sent (the gate itself keeps what's sent low).
|
||||
let busy = sending >= target / 2
|
||||
// Twice in a row (200 ms), so one late ack doesn't count.
|
||||
let signal = (busy && queueing > 40_000) || held >= 3 || stuck
|
||||
let congested = signal && lastSignal
|
||||
lastSignal = signal
|
||||
let heldNow = held
|
||||
held = 0
|
||||
let floorBps = 300_000
|
||||
if congested, now - lastDecrease > 300_000 {
|
||||
// Down to a bit under what got through: at least a fifth off, at
|
||||
// most half (a stall delivers nothing, but the link is still there).
|
||||
let measured = Int(Double(delivered) * 0.9)
|
||||
var next = max(floorBps, min(target * 4 / 5, max(measured, target / 2)))
|
||||
// App-limited (it wants about half its budget or less): never
|
||||
// below twice what it wants, however many cuts in a row. The
|
||||
// extra quarter is hysteresis, so frame sizes wobbling at the
|
||||
// floor don't switch the protection off.
|
||||
if demand > 0, demand * 2 <= target * 5 / 4 { next = max(next, min(target, demand * 2)) }
|
||||
target = next
|
||||
lastDecrease = now
|
||||
events.append((now, "down to \(target / 1000) kbit/s: queue \(queueing / 1000) ms, held \(heldNow), "
|
||||
+ "oldest \(oldestAge / 1000) ms, base \(base / 1000) ms, sent \(sending / 1000) got \(delivered / 1000) "
|
||||
+ "wants \(demand / 1000)"))
|
||||
} else if !congested, now - lastDecrease > 1_000_000, now - lastIncrease > 250_000, target < ceiling,
|
||||
sending > target * 6 / 10 || now - lastDecrease > 3_000_000 {
|
||||
// Clear for a second and using what it has: probe up.
|
||||
target = min(ceiling, Int(Double(target) * 1.1) + 50_000)
|
||||
lastIncrease = now
|
||||
}
|
||||
// Fewer frames or pixels only help a stream that fills its budget;
|
||||
// a small one (a still window, a test card) keeps its tier.
|
||||
retier(now: now, linkLimited: sending >= target * 7 / 10)
|
||||
if events.count > 200 { events.removeFirst(events.count - 200) }
|
||||
return target
|
||||
}
|
||||
}
|
||||
|
||||
/// Steps down quickly, straight to the tier the bitrate supports, and back
|
||||
/// up one tier at a time only when there's clearly room (hysteresis).
|
||||
private func retier(now: Int64, linkLimited: Bool) {
|
||||
let share = Double(target) / Double(max(ceiling, 1))
|
||||
if tier < RateController.tiers.count - 1, share < RateController.floors[tier], linkLimited {
|
||||
if belowSince == 0 { belowSince = now }
|
||||
if now - belowSince > 500_000 {
|
||||
tier = RateController.floors.firstIndex { share >= $0 } ?? RateController.tiers.count - 1
|
||||
belowSince = 0
|
||||
events.append((now, "tier \(tier)"))
|
||||
}
|
||||
} else {
|
||||
belowSince = 0
|
||||
}
|
||||
if tier > 0, share > RateController.floors[tier - 1] * 1.25 {
|
||||
if aboveSince == 0 { aboveSince = now }
|
||||
if now - aboveSince > 2_000_000 {
|
||||
tier -= 1
|
||||
aboveSince = 0
|
||||
events.append((now, "tier \(tier)"))
|
||||
}
|
||||
} else {
|
||||
aboveSince = 0
|
||||
}
|
||||
}
|
||||
|
||||
func state() -> [String: Any] {
|
||||
locked {
|
||||
["target": target, "ceiling": ceiling, "tier": tier, "fps": min(maxFps, RateController.tiers[tier].fps),
|
||||
"scale": RateController.tiers[tier].scale, "baseRtt": Double(baseline()) / 1000,
|
||||
"inFlight": unacked.count, "slack": Double(slack) / 1000, "adapt": RateController.enabled]
|
||||
}
|
||||
}
|
||||
|
||||
func eventList() -> [[String: Any]] { locked { events.map { ["t": $0.0, "e": $0.1] } } }
|
||||
}
|
||||
@@ -18,11 +18,21 @@ final class Encoder {
|
||||
private var session: VTCompressionSession?
|
||||
private var forceKey = true
|
||||
private var lastPts = CMTime.invalid
|
||||
/// Called on VideoToolbox's thread with one access unit (or JPEG) per frame.
|
||||
var onFrame: ((Data, Bool, CMTime) -> Void)?
|
||||
private var lastGiven = CMTime.invalid
|
||||
/// Bits per second to aim for; nil means from the size, frame rate and
|
||||
/// bits per pixel. Changing it takes effect on the next frame.
|
||||
private var target: Int?
|
||||
private(set) var bitrate = 0
|
||||
/// Called on VideoToolbox's thread with one access unit (or JPEG) per
|
||||
/// frame, and the sequence number it was submitted with.
|
||||
var onFrame: ((Data, Bool, CMTime, UInt32) -> Void)?
|
||||
var onError: ((String) -> Void)?
|
||||
/// Called once per frame handed to VideoToolbox, however it went.
|
||||
var onDone: (() -> Void)?
|
||||
var onDone: ((UInt32) -> Void)?
|
||||
/// Experiment switches (FRAME_MAC_VIEW_ENCODER, comma-separated), so a
|
||||
/// benchmark can compare them without a rebuild.
|
||||
static let options = Set((ProcessInfo.processInfo.environment["FRAME_MAC_VIEW_ENCODER"] ?? "")
|
||||
.split(separator: ",").map(String.init))
|
||||
|
||||
init(codec: Codec, fps: Int, bitsPerPixel: Double) {
|
||||
self.codec = codec
|
||||
@@ -30,6 +40,27 @@ final class Encoder {
|
||||
self.bitsPerPixel = bitsPerPixel
|
||||
}
|
||||
|
||||
/// The bitrate the size and quality setting would give.
|
||||
func defaultBitrate(width w: Int, height h: Int) -> Int {
|
||||
Int(min(max(Double(w * h * fps) * bitsPerPixel, 2_000_000), 60_000_000))
|
||||
}
|
||||
|
||||
/// Live, without a new keyframe. Call on the encoding queue.
|
||||
func setBitrate(_ bps: Int?) {
|
||||
target = bps
|
||||
guard codec == .h264, let s = session else { return }
|
||||
let b = bps ?? defaultBitrate(width: width, height: height)
|
||||
guard b != bitrate else { return }
|
||||
bitrate = b
|
||||
VTSessionSetProperty(s, key: kVTCompressionPropertyKey_AverageBitRate, value: b as CFTypeRef)
|
||||
// A hard ceiling too: no more than 200 ms' worth of bits in any 200 ms,
|
||||
// so a keyframe can't hold the link for long.
|
||||
if Encoder.options.contains("cap") {
|
||||
VTSessionSetProperty(s, key: kVTCompressionPropertyKey_DataRateLimits,
|
||||
value: [b / 8 / 5, 0.2] as CFArray)
|
||||
}
|
||||
}
|
||||
|
||||
deinit { invalidate() }
|
||||
|
||||
func requestKeyFrame() { forceKey = true }
|
||||
@@ -46,6 +77,7 @@ final class Encoder {
|
||||
invalidate()
|
||||
var spec: [CFString: Any] = [:]
|
||||
if codec == .h264 { spec[kVTVideoEncoderSpecification_EnableLowLatencyRateControl] = true }
|
||||
if Encoder.options.contains("hw") { spec[kVTVideoEncoderSpecification_RequireHardwareAcceleratedVideoEncoder] = true }
|
||||
var s: VTCompressionSession?
|
||||
let type = codec == .h264 ? kCMVideoCodecType_H264 : kCMVideoCodecType_JPEG
|
||||
func create(_ spec: [CFString: Any]) -> OSStatus {
|
||||
@@ -68,11 +100,11 @@ final class Encoder {
|
||||
set(kVTCompressionPropertyKey_TransferFunction, kCVImageBufferTransferFunction_ITU_R_709_2)
|
||||
set(kVTCompressionPropertyKey_YCbCrMatrix, kCVImageBufferYCbCrMatrix_ITU_R_709_2)
|
||||
if codec == .h264 {
|
||||
let bps = min(max(Double(w * h * fps) * bitsPerPixel, 2_000_000), 60_000_000)
|
||||
set(kVTCompressionPropertyKey_ProfileLevel, kVTProfileLevel_H264_ConstrainedHigh_AutoLevel)
|
||||
set(kVTCompressionPropertyKey_AllowFrameReordering, false)
|
||||
set(kVTCompressionPropertyKey_AverageBitRate, Int(bps))
|
||||
set(kVTCompressionPropertyKey_ExpectedFrameRate, fps)
|
||||
if Encoder.options.contains("nodelay") { set(kVTCompressionPropertyKey_MaxFrameDelayCount, 0) }
|
||||
if Encoder.options.contains("speed") { set(kVTCompressionPropertyKey_PrioritizeEncodingSpeedOverQuality, true) }
|
||||
// A keyframe every 10 s at most, so a viewer that lost one recovers
|
||||
// even if it never asks. Viewers ask for one when they start.
|
||||
set(kVTCompressionPropertyKey_MaxKeyFrameIntervalDuration, 10)
|
||||
@@ -82,24 +114,37 @@ final class Encoder {
|
||||
VTCompressionSessionPrepareToEncodeFrames(s)
|
||||
session = s
|
||||
lastPts = .invalid
|
||||
lastGiven = .invalid
|
||||
width = w
|
||||
height = h
|
||||
forceKey = true
|
||||
bitrate = 0
|
||||
setBitrate(target)
|
||||
return true
|
||||
}
|
||||
|
||||
/// False if the frame never reached VideoToolbox (then onDone won't come).
|
||||
@discardableResult
|
||||
func encode(_ pb: CVPixelBuffer, pts given: CMTime) -> Bool {
|
||||
// VideoToolbox needs strictly increasing timestamps; a resent picture
|
||||
// stamped "now" can be followed by a capture stamped a moment earlier.
|
||||
var pts = given
|
||||
if lastPts.isValid, CMTimeCompare(pts, lastPts) <= 0 { pts = CMTimeAdd(lastPts, CMTime(value: 1, timescale: 1_000_000)) }
|
||||
func encode(_ pb: CVPixelBuffer, pts given: CMTime, seq: UInt32) -> Bool {
|
||||
// The encoder's own timeline: strictly increasing (a resent picture
|
||||
// stamped "now" can be followed by a capture stamped a moment earlier),
|
||||
// and never more than two frames on from the last one. Rate control
|
||||
// budgets bits by elapsed time, so after a pause (the link held frames
|
||||
// back, or nothing changed) one frame would otherwise get a quarter
|
||||
// second's worth of bits: 200 KB that then hold a slow link for half a second.
|
||||
let w = CVPixelBufferGetWidth(pb), h = CVPixelBufferGetHeight(pb)
|
||||
if session == nil || w != width || h != height {
|
||||
guard makeSession(width: w, height: h) else { return false }
|
||||
guard makeSession(width: w, height: h) else { return false } // a new timeline too
|
||||
}
|
||||
guard let s = session else { return false }
|
||||
var pts = given
|
||||
if lastPts.isValid, lastGiven.isValid {
|
||||
let gap = CMTimeSubtract(given, lastGiven)
|
||||
let most = CMTime(value: 2, timescale: CMTimeScale(fps))
|
||||
let step = CMTimeCompare(gap, most) > 0 ? most : gap
|
||||
pts = CMTimeAdd(lastPts, CMTimeMaximum(step, CMTime(value: 1, timescale: 1_000_000)))
|
||||
}
|
||||
lastGiven = given
|
||||
var props: CFDictionary?
|
||||
if forceKey {
|
||||
props = [kVTEncodeFrameOptionKey_ForceKeyFrame: true] as CFDictionary
|
||||
@@ -110,12 +155,12 @@ final class Encoder {
|
||||
let status = VTCompressionSessionEncodeFrame(s, imageBuffer: pb, presentationTimeStamp: pts, duration: .invalid,
|
||||
frameProperties: props, infoFlagsOut: nil) { [weak self] status, _, sample in
|
||||
guard let self else { return }
|
||||
defer { self.onDone?() }
|
||||
defer { self.onDone?(seq) }
|
||||
guard status == noErr, let sample else { return }
|
||||
if codec == .jpeg {
|
||||
if let data = Self.bytes(sample) { self.onFrame?(data, true, pts) }
|
||||
if let data = Self.bytes(sample) { self.onFrame?(data, true, pts, seq) }
|
||||
} else if let (data, key) = Self.annexB(sample) {
|
||||
self.onFrame?(data, key, pts)
|
||||
self.onFrame?(data, key, pts, seq)
|
||||
}
|
||||
}
|
||||
if status != noErr { forceKey = true }
|
||||
|
||||
@@ -0,0 +1,295 @@
|
||||
// "Separate" mode: a streamed window gets a virtual display of its own. The
|
||||
// window is moved onto it and sized to fill it, and the whole display is
|
||||
// captured, so its menus, sheets, popovers and tooltips come along, nothing
|
||||
// can cover it, and clicks land on it without raising anything. On stop the
|
||||
// window goes back where it was. See docs/mac-window-separation.md.
|
||||
import AppKit
|
||||
import ApplicationServices
|
||||
import CoreMedia
|
||||
import Foundation
|
||||
|
||||
@_silgen_name("_AXUIElementGetWindow")
|
||||
private func axWindowID(_ element: AXUIElement, _ id: UnsafeMutablePointer<CGWindowID>) -> AXError
|
||||
|
||||
/// One of our virtual displays, HiDPI (2 pixels per point). Main thread only.
|
||||
final class VirtualDisplay {
|
||||
static let vendor: UInt32 = 0xF0C0
|
||||
/// How often macOS composites our displays (FRAME_MAC_VIEW_VD_HZ to experiment).
|
||||
static let refreshRate = Double(ProcessInfo.processInfo.environment["FRAME_MAC_VIEW_VD_HZ"] ?? "") ?? 60
|
||||
private var display: CGVirtualDisplay?
|
||||
let id: CGDirectDisplayID
|
||||
/// Removal is deferred while the Mac's screen sleeps, and an identity
|
||||
/// can't be reused until it's gone, so each display gets a fresh serial.
|
||||
private static var serial = (UInt32(truncatingIfNeeded: ProcessInfo.processInfo.processIdentifier) & 0xFFFF) << 12
|
||||
|
||||
init?(name: String, width: Int, height: Int) {
|
||||
var made: CGVirtualDisplay?
|
||||
for _ in 0..<8 where made == nil {
|
||||
VirtualDisplay.serial &+= 1
|
||||
let d = CGVirtualDisplayDescriptor()
|
||||
d.queue = DispatchQueue.main
|
||||
d.name = name
|
||||
d.maxPixelsWide = UInt32(2 * width)
|
||||
d.maxPixelsHigh = UInt32(2 * height)
|
||||
d.sizeInMillimeters = CGSize(width: Double(width) * 0.2646, height: Double(height) * 0.2646)
|
||||
d.vendorID = VirtualDisplay.vendor
|
||||
d.productID = 0x5E9A
|
||||
d.serialNum = VirtualDisplay.serial
|
||||
d.terminationHandler = { _, _ in }
|
||||
guard let vd = CGVirtualDisplay(descriptor: d) else { continue }
|
||||
let s = CGVirtualDisplaySettings()
|
||||
s.hiDPI = 1
|
||||
let hz = VirtualDisplay.refreshRate
|
||||
s.modes = [CGVirtualDisplayMode(width: UInt32(2 * width), height: UInt32(2 * height), refreshRate: hz)!,
|
||||
CGVirtualDisplayMode(width: UInt32(width), height: UInt32(height), refreshRate: hz)!]
|
||||
if vd.apply(s), vd.displayID != 0 { made = vd }
|
||||
}
|
||||
guard let made else { return nil }
|
||||
display = made
|
||||
id = made.displayID
|
||||
// Pick the HiDPI mode: `width` points drawn with twice the pixels.
|
||||
let modes = (CGDisplayCopyAllDisplayModes(id, [kCGDisplayShowDuplicateLowResolutionModes: true] as CFDictionary)
|
||||
as? [CGDisplayMode]) ?? []
|
||||
if let hidpi = modes.first(where: {
|
||||
$0.width == width && $0.height == height && $0.pixelWidth == 2 * width && $0.refreshRate == VirtualDisplay.refreshRate
|
||||
}) ?? modes.first(where: { $0.width == width && $0.height == height && $0.pixelWidth == 2 * width }) {
|
||||
var cfg: CGDisplayConfigRef?
|
||||
CGBeginDisplayConfiguration(&cfg)
|
||||
CGConfigureDisplayWithDisplayMode(cfg, id, hidpi, nil)
|
||||
CGCompleteDisplayConfiguration(cfg, .forSession)
|
||||
}
|
||||
}
|
||||
|
||||
var bounds: CGRect { CGDisplayBounds(id) }
|
||||
|
||||
var isHiDPI: Bool { CGDisplayCopyDisplayMode(id).map { $0.pixelWidth > $0.width } ?? false }
|
||||
|
||||
var screen: NSScreen? {
|
||||
NSScreen.screens.first {
|
||||
($0.deviceDescription[NSDeviceDescriptionKey("NSScreenNumber")] as? NSNumber)?.uint32Value == id
|
||||
}
|
||||
}
|
||||
|
||||
/// Where a window may go, in global (top-left origin) points: the display
|
||||
/// minus its menu bar and anything else macOS reserves there.
|
||||
var usableRect: CGRect {
|
||||
let b = bounds
|
||||
guard let s = screen else { return b }
|
||||
let top = s.frame.maxY - s.visibleFrame.maxY, bottom = s.visibleFrame.minY - s.frame.minY
|
||||
let left = s.visibleFrame.minX - s.frame.minX, right = s.frame.maxX - s.visibleFrame.maxX
|
||||
return CGRect(x: b.minX + left, y: b.minY + top, width: b.width - left - right, height: b.height - top - bottom)
|
||||
}
|
||||
|
||||
func release() { display = nil }
|
||||
|
||||
static func isOurs(_ id: CGDirectDisplayID) -> Bool { CGDisplayVendorNumber(id) == vendor }
|
||||
}
|
||||
|
||||
/// A window through the Accessibility API.
|
||||
struct AXWindow {
|
||||
let element: AXUIElement
|
||||
|
||||
init?(windowID: CGWindowID, pid: pid_t) {
|
||||
let app = AXUIElementCreateApplication(pid)
|
||||
var value: CFTypeRef?
|
||||
guard AXUIElementCopyAttributeValue(app, kAXWindowsAttribute as CFString, &value) == .success,
|
||||
let windows = value as? [AXUIElement] else { return nil }
|
||||
guard let match = windows.first(where: { w in
|
||||
var id: CGWindowID = 0
|
||||
return axWindowID(w, &id) == .success && id == windowID
|
||||
}) else { return nil }
|
||||
element = match
|
||||
}
|
||||
|
||||
var frame: CGRect? {
|
||||
var p: CFTypeRef?, s: CFTypeRef?
|
||||
guard AXUIElementCopyAttributeValue(element, kAXPositionAttribute as CFString, &p) == .success,
|
||||
AXUIElementCopyAttributeValue(element, kAXSizeAttribute as CFString, &s) == .success else { return nil }
|
||||
var point = CGPoint.zero, size = CGSize.zero
|
||||
AXValueGetValue(p as! AXValue, .cgPoint, &point)
|
||||
AXValueGetValue(s as! AXValue, .cgSize, &size)
|
||||
return CGRect(origin: point, size: size)
|
||||
}
|
||||
|
||||
var isFullScreen: Bool {
|
||||
var v: CFTypeRef?
|
||||
return AXUIElementCopyAttributeValue(element, "AXFullScreen" as CFString, &v) == .success && (v as? Bool) == true
|
||||
}
|
||||
|
||||
/// Moving to another display: position first (so the size fits there),
|
||||
/// then size, then position again in case the size change moved it.
|
||||
func setFrame(_ r: CGRect) {
|
||||
var origin = r.origin, size = r.size
|
||||
let pos = AXValueCreate(.cgPoint, &origin)!, sz = AXValueCreate(.cgSize, &size)!
|
||||
AXUIElementSetAttributeValue(element, kAXPositionAttribute as CFString, pos)
|
||||
AXUIElementSetAttributeValue(element, kAXSizeAttribute as CFString, sz)
|
||||
AXUIElementSetAttributeValue(element, kAXPositionAttribute as CFString, pos)
|
||||
}
|
||||
}
|
||||
|
||||
/// A window on a display of its own, captured as that display.
|
||||
final class SeparateSource: CaptureSource {
|
||||
let windowID: CGWindowID
|
||||
var onFrame: ((CVPixelBuffer, CMTime, Int64) -> Void)?
|
||||
var onEnded: ((String) -> Void)?
|
||||
var onChange: (() -> Void)?
|
||||
private(set) var title = ""
|
||||
private(set) var app = ""
|
||||
private(set) var pid: pid_t?
|
||||
private(set) var gone = false
|
||||
private var display: VirtualDisplay?
|
||||
private var inner: SCKSource?
|
||||
private var window: AXWindow?
|
||||
private var original: CGRect?
|
||||
private var crop: CGRect? // display points, when only part of the display is captured
|
||||
private var poll: DispatchSourceTimer?
|
||||
private var stopped = false
|
||||
private let queue = DispatchQueue(label: "frame-mac-view.separate")
|
||||
|
||||
init(windowID: CGWindowID) { self.windowID = windowID }
|
||||
|
||||
var frameRect: CGRect { inner?.frameRect ?? .zero }
|
||||
var displayID: CGDirectDisplayID? { display?.id }
|
||||
|
||||
func start(maxLong: Int, fps: Int, completion: @escaping (String?) -> Void) {
|
||||
guard AXIsProcessTrusted() else {
|
||||
return completion("Separate windows need the Accessibility permission (to move the window)")
|
||||
}
|
||||
guard let info = WindowInfo.find(windowID) else {
|
||||
gone = true
|
||||
return completion("that window has closed")
|
||||
}
|
||||
title = info.title
|
||||
app = info.app
|
||||
pid = info.pid
|
||||
DispatchQueue.main.async { [self] in
|
||||
guard !stopped else { return }
|
||||
guard let w = AXWindow(windowID: windowID, pid: info.pid), let frame = w.frame else {
|
||||
return completion("couldn't reach that window through Accessibility")
|
||||
}
|
||||
if w.isFullScreen { return completion("take the window out of full screen first") }
|
||||
// A display the window's size, plus room for the menu bar, in
|
||||
// points; within what a panel can show sharply.
|
||||
let width = min(max(Int(frame.width.rounded()), 640), 2560) / 2 * 2
|
||||
let height = min(max(Int(frame.height.rounded()) + 40, 480), 1600) / 2 * 2
|
||||
guard let vd = VirtualDisplay(name: app.isEmpty ? "Frame Control" : "\(app) (Frame)", width: width, height: height) else {
|
||||
return completion("macOS wouldn't create a display for this window")
|
||||
}
|
||||
display = vd
|
||||
window = w
|
||||
original = frame
|
||||
// The new screen shows up in NSScreen a moment later; its menu bar
|
||||
// decides where the window can go.
|
||||
placeWindow(attempt: 0, maxLong: maxLong, fps: fps, completion: completion)
|
||||
}
|
||||
}
|
||||
|
||||
private func placeWindow(attempt: Int, maxLong: Int, fps: Int, completion: @escaping (String?) -> Void) {
|
||||
guard !stopped, let vd = display, let w = window else { return }
|
||||
// Wait for the screen to appear, and for its HiDPI mode, so the
|
||||
// first frames are captured sharp.
|
||||
if vd.screen == nil || !vd.isHiDPI, attempt < 30 {
|
||||
DispatchQueue.main.asyncAfter(deadline: .now() + 0.1) {
|
||||
self.placeWindow(attempt: attempt + 1, maxLong: maxLong, fps: fps, completion: completion)
|
||||
}
|
||||
return
|
||||
}
|
||||
let target = vd.usableRect
|
||||
w.setFrame(target)
|
||||
guard let now = w.frame, vd.bounds.intersects(now) else {
|
||||
restore()
|
||||
return completion("macOS didn't let the window move to its own display (Stage Manager can prevent this)")
|
||||
}
|
||||
let sck = SCKSource(.display(vd.id))
|
||||
// A window that keeps its own size (Calculator, some dialogs) sits in
|
||||
// the display's top-left corner: send only that corner, menu bar
|
||||
// included, with room around it for menus, not a panel of wallpaper.
|
||||
let b = vd.bounds
|
||||
if now.width < target.width - 8 || now.height < target.height - 8 {
|
||||
let w = min(b.width, max(now.maxX - b.minX, 480)), h = min(b.height, max(now.maxY - b.minY, 360))
|
||||
crop = CGRect(x: 0, y: 0, width: (w / 2).rounded(.up) * 2, height: (h / 2).rounded(.up) * 2)
|
||||
sck.crop = crop
|
||||
}
|
||||
sck.onFrame = { [weak self] pb, pts, shown in self?.onFrame?(pb, pts, shown) }
|
||||
sck.onEnded = { [weak self] reason in self?.onEnded?(reason) }
|
||||
inner = sck
|
||||
sck.start(maxLong: maxLong, fps: fps) { [weak self] error in
|
||||
// Back on main, where stop() runs too, so a viewer that left during
|
||||
// startup can't leave a capture running.
|
||||
DispatchQueue.main.async {
|
||||
guard let self, !self.stopped else { return sck.stop() }
|
||||
if let error {
|
||||
self.restore()
|
||||
return completion(error)
|
||||
}
|
||||
self.startPolling()
|
||||
completion(nil)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// The window can close or be retitled; ScreenCaptureKit won't say.
|
||||
private func startPolling() {
|
||||
let t = DispatchSource.makeTimerSource(queue: queue)
|
||||
t.schedule(deadline: .now() + 1, repeating: 1)
|
||||
t.setEventHandler { [weak self] in
|
||||
guard let self else { return }
|
||||
guard let info = WindowInfo.find(self.windowID) else {
|
||||
self.gone = true
|
||||
self.onEnded?("the window closed")
|
||||
return
|
||||
}
|
||||
if !info.title.isEmpty, info.title != self.title {
|
||||
self.title = info.title
|
||||
self.onChange?()
|
||||
}
|
||||
// Displays added at the same moment can rearrange each other, and
|
||||
// someone may drag the window away: put it back on its display.
|
||||
if !info.bounds.isEmpty {
|
||||
DispatchQueue.main.async { self.keepOnDisplay(info.bounds) }
|
||||
}
|
||||
}
|
||||
t.resume()
|
||||
poll = t
|
||||
}
|
||||
|
||||
private func keepOnDisplay(_ now: CGRect) {
|
||||
guard !stopped, let vd = display, let w = window else { return }
|
||||
let b = vd.bounds
|
||||
if let c = crop {
|
||||
// Only part of the display is sent: the window must stay inside it.
|
||||
guard !c.offsetBy(dx: b.minX, dy: b.minY).insetBy(dx: -2, dy: -2).contains(now) else { return }
|
||||
} else {
|
||||
guard !b.contains(CGPoint(x: now.midX, y: now.minY + 10)) else { return }
|
||||
}
|
||||
let target = vd.usableRect
|
||||
w.setFrame(CGRect(origin: target.origin, size: now.size.width < target.width - 8 ? now.size : target.size))
|
||||
}
|
||||
|
||||
/// Lifecycle state (stopped, inner, poll, window, display) is only touched on main.
|
||||
func stop() {
|
||||
DispatchQueue.main.async { [self] in
|
||||
stopped = true
|
||||
poll?.cancel()
|
||||
poll = nil
|
||||
inner?.stop()
|
||||
restore()
|
||||
}
|
||||
}
|
||||
|
||||
/// Puts the window back where it was, then removes the display. In that
|
||||
/// order: removing a display first would let macOS pick where it goes.
|
||||
private func restore() {
|
||||
if let w = window, let r = original, WindowInfo.find(windowID) != nil { w.setFrame(r) }
|
||||
window = nil
|
||||
original = nil
|
||||
let vd = display
|
||||
display = nil
|
||||
// Give WindowServer a moment to move the window before the display goes.
|
||||
DispatchQueue.main.asyncAfter(deadline: .now() + 0.3) { vd?.release() }
|
||||
}
|
||||
|
||||
func setMaxLong(_ n: Int) { DispatchQueue.main.async { self.inner?.setMaxLong(n) } }
|
||||
|
||||
func pointer(x: Double, y: Double, text: String?) {}
|
||||
}
|
||||
@@ -165,9 +165,10 @@ final class WebSocket {
|
||||
func sendJSON(_ obj: Any) {
|
||||
if let d = try? JSONSerialization.data(withJSONObject: obj), let s = String(data: d, encoding: .utf8) { sendText(s) }
|
||||
}
|
||||
func sendBinary(_ d: Data) { send(opcode: 2, d) }
|
||||
/// `taken` runs once the network stack has taken the whole frame.
|
||||
func sendBinary(_ d: Data, taken: (() -> Void)? = nil) { send(opcode: 2, d, taken: taken) }
|
||||
|
||||
func send(opcode: UInt8, _ payload: Data) {
|
||||
func send(opcode: UInt8, _ payload: Data, taken: (() -> Void)? = nil) {
|
||||
var frame = Data([0x80 | opcode])
|
||||
let n = payload.count
|
||||
if n < 126 {
|
||||
@@ -185,6 +186,7 @@ final class WebSocket {
|
||||
conn.send(content: frame, completion: .contentProcessed { [weak self] _ in
|
||||
guard let self else { return }
|
||||
self.lock.lock(); self._pending -= size; self.lock.unlock()
|
||||
taken?()
|
||||
self.onDrain?()
|
||||
})
|
||||
}
|
||||
|
||||
@@ -0,0 +1,180 @@
|
||||
// Per-frame timing, so every change to the stream can be judged by numbers.
|
||||
// All times are the agent's host clock in microseconds (the clock
|
||||
// ScreenCaptureKit stamps frames with). The viewer syncs to it over the
|
||||
// WebSocket and reports when each frame arrived, decoded and was drawn.
|
||||
import CoreMedia
|
||||
import Foundation
|
||||
|
||||
private let timebase: mach_timebase_info_data_t = {
|
||||
var t = mach_timebase_info_data_t()
|
||||
mach_timebase_info(&t)
|
||||
return t
|
||||
}()
|
||||
|
||||
/// Now on the host clock, in microseconds.
|
||||
func nowUs() -> Int64 { hostUs(mach_absolute_time()) }
|
||||
|
||||
/// A mach_absolute_time value (as in SCStreamFrameInfo.displayTime) in microseconds.
|
||||
func hostUs(_ ticks: UInt64) -> Int64 { Int64(ticks / 1000 * UInt64(timebase.numer) / UInt64(timebase.denom)) }
|
||||
|
||||
/// A host-clock CMTime (ScreenCaptureKit's presentation times) in microseconds.
|
||||
func hostUs(_ t: CMTime) -> Int64 { t.isValid ? Int64(CMTimeGetSeconds(t) * 1_000_000) : nowUs() }
|
||||
|
||||
/// One frame's journey. Zero means "didn't happen" or "not reported yet".
|
||||
struct FrameRecord {
|
||||
var seq: UInt32 = 0
|
||||
var key = false
|
||||
var bytes = 0
|
||||
var width = 0, height = 0
|
||||
var capture: Int64 = 0 // the Mac composited it (display time)
|
||||
var arrived: Int64 = 0 // ScreenCaptureKit handed it to us
|
||||
var encodeStart: Int64 = 0
|
||||
var encodeEnd: Int64 = 0
|
||||
var sent: Int64 = 0 // handed to the WebSocket
|
||||
var wire: Int64 = 0 // the socket took it
|
||||
var received: Int64 = 0 // viewer, agent clock
|
||||
var decoded: Int64 = 0
|
||||
var drawn: Int64 = 0
|
||||
var vsync: Int64 = 0 // the viewer's next animation frame after drawing it
|
||||
var echo: UInt32 = 0 // the first frame after input `echo`
|
||||
var tier = 0
|
||||
var bitrate = 0
|
||||
|
||||
var json: [String: Any] {
|
||||
["s": seq, "k": key ? 1 : 0, "b": bytes, "w": width, "h": height, "cap": capture, "arr": arrived,
|
||||
"e0": encodeStart, "e1": encodeEnd, "snd": sent, "wire": wire, "rx": received, "dec": decoded,
|
||||
"drw": drawn, "vs": vsync, "echo": echo, "tier": tier, "br": bitrate]
|
||||
}
|
||||
}
|
||||
|
||||
/// One input event from the viewer and what came of it.
|
||||
struct InputRecord {
|
||||
var id: UInt32
|
||||
var kind: String
|
||||
var viewer: Int64 // the viewer's event time, agent clock
|
||||
var injected: Int64 = 0 // posted as a CGEvent
|
||||
var frame: UInt32 = 0 // first frame captured after it (0: none yet)
|
||||
var capture: Int64 = 0
|
||||
|
||||
var json: [String: Any] {
|
||||
["id": id, "kind": kind, "tv": viewer, "inj": injected, "frame": frame, "cap": capture]
|
||||
}
|
||||
}
|
||||
|
||||
func percentile(_ sorted: [Double], _ p: Double) -> Double? {
|
||||
guard !sorted.isEmpty else { return nil }
|
||||
let i = min(sorted.count - 1, max(0, Int((p * Double(sorted.count - 1)).rounded())))
|
||||
return sorted[i]
|
||||
}
|
||||
|
||||
/// Frames and inputs of one stream, kept for the last few thousand frames.
|
||||
/// Any thread.
|
||||
final class StreamStats {
|
||||
static let capacity = 4096
|
||||
private let lock = NSLock()
|
||||
private var frames = [FrameRecord?](repeating: nil, count: StreamStats.capacity)
|
||||
private var inputs: [UInt32: InputRecord] = [:]
|
||||
private var inputOrder: [UInt32] = []
|
||||
private(set) var nextSeq: UInt32 = 1
|
||||
var captured = 0 // frames ScreenCaptureKit delivered
|
||||
var skipped = 0 // held back from encoding (link, encoder, pacing or gate busy); only the newest may go later
|
||||
var viewerDropped = 0 // not shown by the viewer (behind, or waiting for a keyframe)
|
||||
var rtt: Double = 0 // ms, the viewer's best recent ping
|
||||
var clockSynced = false
|
||||
var decoder = "" // what the viewer reports about its decoder
|
||||
|
||||
func withLock<T>(_ f: () -> T) -> T { lock.lock(); defer { lock.unlock() }; return f() }
|
||||
|
||||
func newFrame(_ r: FrameRecord) -> UInt32 {
|
||||
withLock {
|
||||
var r = r
|
||||
r.seq = nextSeq
|
||||
nextSeq &+= 1
|
||||
frames[Int(r.seq) % StreamStats.capacity] = r
|
||||
return r.seq
|
||||
}
|
||||
}
|
||||
|
||||
func update(_ seq: UInt32, _ f: (inout FrameRecord) -> Void) {
|
||||
withLock {
|
||||
let i = Int(seq) % StreamStats.capacity
|
||||
guard var r = frames[i], r.seq == seq else { return }
|
||||
f(&r)
|
||||
frames[i] = r
|
||||
}
|
||||
}
|
||||
|
||||
func frame(_ seq: UInt32) -> FrameRecord? {
|
||||
withLock {
|
||||
let r = frames[Int(seq) % StreamStats.capacity]
|
||||
return r?.seq == seq ? r : nil
|
||||
}
|
||||
}
|
||||
|
||||
func addInput(_ r: InputRecord) {
|
||||
withLock {
|
||||
inputs[r.id] = r
|
||||
inputOrder.append(r.id)
|
||||
if inputOrder.count > 512 { inputs[inputOrder.removeFirst()] = nil }
|
||||
}
|
||||
}
|
||||
|
||||
func updateInput(_ id: UInt32, _ f: (inout InputRecord) -> Void) {
|
||||
withLock { if var r = inputs[id] { f(&r); inputs[id] = r } }
|
||||
}
|
||||
|
||||
/// Frames with seq > `since` that were sent at least `settle` µs ago, so
|
||||
/// the viewer has had time to report on them.
|
||||
func settled(since: UInt32, settle: Int64 = 1_500_000) -> [FrameRecord] {
|
||||
let cutoff = nowUs() - settle
|
||||
return withLock {
|
||||
frames.compactMap { $0 }.filter { $0.seq > since && $0.sent > 0 && $0.sent < cutoff }.sorted { $0.seq < $1.seq }
|
||||
}
|
||||
}
|
||||
|
||||
func inputList() -> [InputRecord] { withLock { inputOrder.compactMap { inputs[$0] } } }
|
||||
|
||||
/// Percentiles over the last `window` µs, for /status and the overlay.
|
||||
func summary(window: Int64 = 2_000_000) -> [String: Any] {
|
||||
let now = nowUs(), from = now - window
|
||||
let recent = withLock { frames.compactMap { $0 }.filter { $0.sent > from } }
|
||||
func ms(_ pick: (FrameRecord) -> (Int64, Int64)) -> [String: Double] {
|
||||
let v = recent.compactMap { r -> Double? in
|
||||
let (a, b) = pick(r)
|
||||
return a > 0 && b > 0 ? Double(b - a) / 1000 : nil
|
||||
}.sorted()
|
||||
guard !v.isEmpty else { return [:] }
|
||||
return ["p50": (percentile(v, 0.5)! * 10).rounded() / 10, "p95": (percentile(v, 0.95)! * 10).rounded() / 10]
|
||||
}
|
||||
let secs = Double(window) / 1_000_000
|
||||
let shown = recent.filter { $0.vsync > 0 }.count
|
||||
let bytes = recent.reduce(0) { $0 + $1.bytes }
|
||||
var s: [String: Any] = [
|
||||
"capture": ms { ($0.capture, $0.arrived) },
|
||||
"queue": ms { ($0.arrived, $0.encodeStart) },
|
||||
"encode": ms { ($0.encodeStart, $0.encodeEnd) },
|
||||
"network": ms { ($0.encodeEnd, $0.received) },
|
||||
"decode": ms { ($0.received, $0.decoded) },
|
||||
"draw": ms { ($0.decoded, $0.vsync) },
|
||||
"total": ms { ($0.capture, $0.vsync) },
|
||||
"fps": (Double(shown) / secs * 10).rounded() / 10,
|
||||
"sentFps": (Double(recent.count) / secs * 10).rounded() / 10,
|
||||
"mbps": (Double(bytes * 8) / secs / 1e5).rounded() / 10,
|
||||
"rtt": (rtt * 10).rounded() / 10,
|
||||
"synced": clockSynced,
|
||||
]
|
||||
withLock {
|
||||
s["captured"] = captured
|
||||
s["skipped"] = skipped
|
||||
s["dropped"] = viewerDropped
|
||||
s["decoder"] = decoder
|
||||
let done = inputOrder.suffix(20).compactMap { inputs[$0] }.compactMap { i -> Double? in
|
||||
guard i.frame != 0, let f = frames[Int(i.frame) % StreamStats.capacity], f.seq == i.frame, f.vsync > 0
|
||||
else { return nil }
|
||||
return Double(f.vsync - i.viewer) / 1000
|
||||
}.sorted()
|
||||
if !done.isEmpty { s["input"] = ["p50": percentile(done, 0.5)!, "n": done.count] }
|
||||
}
|
||||
return s
|
||||
}
|
||||
}
|
||||
@@ -16,6 +16,8 @@
|
||||
// POST /ticket?src=... ?k=: a ticket for one viewer of src
|
||||
// POST /close[?src=...] ?k=: end those streams, close their windows
|
||||
// POST /permissions ?k=: show macOS's permission prompts
|
||||
// GET /stats[?id=&since=] ?k=: per-frame timing records (for benchmarks)
|
||||
// POST /bench?src=&action=... ?k=: ask viewers of src to type, click or show their overlay
|
||||
// src is window:<CGWindowID>, display:<CGDirectDisplayID> or test.
|
||||
import AppKit
|
||||
import ApplicationServices
|
||||
@@ -43,7 +45,9 @@ func displaysJSON() -> [[String: Any]] {
|
||||
var n: UInt32 = 0
|
||||
// Online, not active: a display that's asleep is still one you can stream.
|
||||
CGGetOnlineDisplayList(16, &ids, &n)
|
||||
return ids.prefix(Int(n)).filter { CGDisplayMirrorsDisplay($0) == kCGNullDirectDisplay }.map { id in
|
||||
return ids.prefix(Int(n)).filter {
|
||||
CGDisplayMirrorsDisplay($0) == kCGNullDirectDisplay && !VirtualDisplay.isOurs($0) // not a separated window's
|
||||
}.map { id in
|
||||
let b = CGDisplayBounds(id)
|
||||
return ["id": id, "src": "display:\(id)", "name": displayName(id), "w": Int(b.width), "h": Int(b.height),
|
||||
"main": CGDisplayIsMain(id) != 0]
|
||||
@@ -59,6 +63,9 @@ func printJSON(_ obj: Any) {
|
||||
FileHandle.standardOutput.write(d + Data("\n".utf8))
|
||||
}
|
||||
|
||||
/// A number from a JSON message, whatever its JSON type.
|
||||
func num(_ v: Any?) -> Double? { (v as? NSNumber)?.doubleValue }
|
||||
|
||||
/// One viewer watching one source.
|
||||
final class Session {
|
||||
let id: Int
|
||||
@@ -68,11 +75,28 @@ final class Session {
|
||||
let ws: WebSocket
|
||||
let codec: Codec
|
||||
let reconnectKey: String
|
||||
let stats = StreamStats()
|
||||
let controller: RateController
|
||||
private var appliedScale = 1.0
|
||||
private var lastSubmit: Int64 = 0
|
||||
private var paceScheduled = false
|
||||
private let lock = NSLock()
|
||||
private var last: (CVPixelBuffer, CMTime)?
|
||||
/// A captured picture waiting to be encoded (or the last one encoded).
|
||||
private struct Picture {
|
||||
let pb: CVPixelBuffer
|
||||
let pts: CMTime
|
||||
let capture: Int64
|
||||
let arrived: Int64
|
||||
var echo: UInt32 // the input this picture is the first reply to
|
||||
}
|
||||
private var last: Picture?
|
||||
private var skipped = false
|
||||
private var stopped = false
|
||||
private var inFlight = 0
|
||||
/// Input posted to the Mac whose effect hasn't been captured yet: the
|
||||
/// next picture composited after `after` is tagged with it.
|
||||
private var pendingEcho: (id: UInt32, after: Int64)?
|
||||
private var statsTimer: DispatchSourceTimer?
|
||||
/// Frames already queued for the network; beyond this, or with two frames
|
||||
/// already in the encoder, new frames are skipped (before encoding, so no
|
||||
/// reference frame goes missing) and the newest picture is sent once
|
||||
@@ -93,23 +117,63 @@ final class Session {
|
||||
self.codec = codec
|
||||
self.reconnectKey = reconnectKey
|
||||
encoder = Encoder(codec: codec, fps: fps, bitsPerPixel: bitsPerPixel)
|
||||
controller = RateController(maxFps: fps)
|
||||
maxPending = codec == .jpeg ? 3 << 20 : 1 << 20
|
||||
}
|
||||
|
||||
/// Frames come faster than the current tier's frame rate: this one waits,
|
||||
/// and goes when its turn comes (unless a newer one replaces it).
|
||||
/// Call with `lock` held.
|
||||
private func paced(now: Int64) -> Bool {
|
||||
let fps = controller.fps
|
||||
guard fps < controller.maxFps, lastSubmit > 0 else { return false }
|
||||
let interval = Int64(1_000_000 / fps), due = lastSubmit + interval * 9 / 10
|
||||
guard now < due else { return false }
|
||||
if !paceScheduled {
|
||||
paceScheduled = true
|
||||
encodeQueue.asyncAfter(deadline: .now() + .microseconds(Int(due - now))) { [weak self] in
|
||||
guard let self else { return }
|
||||
self.lock.lock(); self.paceScheduled = false; self.lock.unlock()
|
||||
self.resendIfRoom()
|
||||
}
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
/// Encodes `pb` unless the link or the encoder is busy, in which case it
|
||||
/// becomes the picture to send next.
|
||||
private func offer(_ pb: CVPixelBuffer, pts: CMTime) {
|
||||
private func offer(_ pb: CVPixelBuffer, pts: CMTime, capture shown: Int64) {
|
||||
let arrived = nowUs()
|
||||
stats.withLock { stats.captured += 1 }
|
||||
controller.captured(at: arrived)
|
||||
lock.lock()
|
||||
last = (pb, pts)
|
||||
// A reply to input stays with whichever picture ends up being sent.
|
||||
var echo = skipped ? last?.echo ?? 0 : 0
|
||||
if let p = pendingEcho, shown >= p.after {
|
||||
echo = p.id
|
||||
pendingEcho = nil
|
||||
}
|
||||
let picture = Picture(pb: pb, pts: pts, capture: shown, arrived: arrived, echo: echo)
|
||||
last = picture
|
||||
let busy = stopped || ws.pendingBytes > maxPending || inFlight >= Session.maxInFlight
|
||||
if busy { skipped = true } else { inFlight += 1 }
|
||||
|| paced(now: arrived) || !controller.maySend(now: arrived)
|
||||
if busy { skipped = true } else { inFlight += 1; last?.echo = 0; lastSubmit = arrived }
|
||||
lock.unlock()
|
||||
if !busy { submit(pb, pts: pts) }
|
||||
if busy { stats.withLock { stats.skipped += 1 } } else { submit(picture) }
|
||||
}
|
||||
|
||||
private func submit(_ pb: CVPixelBuffer, pts: CMTime) {
|
||||
private func submit(_ p: Picture) {
|
||||
encodeQueue.async {
|
||||
if !self.encoder.encode(pb, pts: pts) { self.finished() }
|
||||
var r = FrameRecord()
|
||||
r.capture = p.capture
|
||||
r.arrived = p.arrived
|
||||
r.echo = p.echo
|
||||
r.bitrate = self.encoder.bitrate
|
||||
r.tier = self.controller.tier
|
||||
r.encodeStart = nowUs()
|
||||
let seq = self.stats.newFrame(r)
|
||||
if p.echo != 0 { self.stats.updateInput(p.echo) { $0.frame = seq; $0.capture = p.capture } }
|
||||
if !self.encoder.encode(p.pb, pts: p.pts, seq: seq) { self.finished() }
|
||||
}
|
||||
}
|
||||
|
||||
@@ -124,28 +188,61 @@ final class Session {
|
||||
|
||||
private func resendIfRoom() {
|
||||
lock.lock()
|
||||
var next: (CVPixelBuffer, CMTime)?
|
||||
if !stopped, skipped, ws.pendingBytes <= maxPending / 2, inFlight < Session.maxInFlight, let l = last {
|
||||
var next: Picture?
|
||||
let now = nowUs()
|
||||
if !stopped, skipped, ws.pendingBytes <= maxPending / 2, inFlight < Session.maxInFlight, let l = last,
|
||||
!paced(now: now), controller.maySend(now: now, counts: false) {
|
||||
next = l
|
||||
skipped = false
|
||||
inFlight += 1
|
||||
last?.echo = 0
|
||||
lastSubmit = now
|
||||
}
|
||||
lock.unlock()
|
||||
if let (pb, _) = next { submit(pb, pts: CMClockGetTime(CMClockGetHostTimeClock())) }
|
||||
// Stamped now: VideoToolbox needs increasing times, and the capture
|
||||
// time stays in the record, so the wait counts as latency.
|
||||
if let p = next {
|
||||
submit(Picture(pb: p.pb, pts: CMClockGetTime(CMClockGetHostTimeClock()), capture: p.capture,
|
||||
arrived: p.arrived, echo: p.echo))
|
||||
}
|
||||
}
|
||||
|
||||
func start(maxLong: Int, fps: Int) {
|
||||
encoder.onFrame = { [weak self] data, key, pts in
|
||||
encoder.onFrame = { [weak self] data, key, pts, seq in
|
||||
guard let self else { return }
|
||||
var msg = Data([key ? 1 : 0])
|
||||
let t = nowUs()
|
||||
// The quality setting's bitrate, at full size, is the most it gets.
|
||||
if self.appliedScale == 1 {
|
||||
self.controller.setCeiling(self.encoder.defaultBitrate(width: self.encoder.width, height: self.encoder.height))
|
||||
}
|
||||
self.controller.sent(seq: seq, bytes: data.count + 17, at: t)
|
||||
var echo: UInt32 = 0
|
||||
let (w, h) = (self.encoder.width, self.encoder.height)
|
||||
self.stats.update(seq) {
|
||||
$0.encodeEnd = t
|
||||
$0.sent = t
|
||||
$0.bytes = data.count
|
||||
$0.key = key
|
||||
$0.width = w
|
||||
$0.height = h
|
||||
echo = $0.echo
|
||||
}
|
||||
// Header: flags (1 = keyframe), pts µs, sequence number, and the
|
||||
// input this frame is the first reply to; all big-endian.
|
||||
var msg = Data(capacity: data.count + 17)
|
||||
msg.append(key ? 1 : 0)
|
||||
var us = UInt64(max(0, CMTimeGetSeconds(pts)) * 1_000_000).bigEndian
|
||||
msg.append(Data(bytes: &us, count: 8))
|
||||
var s = seq.bigEndian, e = echo.bigEndian
|
||||
msg.append(Data(bytes: &s, count: 4))
|
||||
msg.append(Data(bytes: &e, count: 4))
|
||||
msg.append(data)
|
||||
self.ws.sendBinary(msg)
|
||||
let stats = self.stats
|
||||
self.ws.sendBinary(msg) { stats.update(seq) { $0.wire = nowUs() } }
|
||||
}
|
||||
encoder.onDone = { [weak self] in self?.finished() }
|
||||
encoder.onDone = { [weak self] _ in self?.finished() }
|
||||
encoder.onError = { [weak self] message in self?.ws.sendJSON(["t": "error", "message": message]) }
|
||||
capture.onFrame = { [weak self] pb, pts in self?.offer(pb, pts: pts) }
|
||||
capture.onFrame = { [weak self] pb, pts, shown in self?.offer(pb, pts: pts, capture: shown) }
|
||||
capture.onEnded = { [weak self] reason in
|
||||
guard let self else { return }
|
||||
self.ws.sendJSON(["t": "closed", "reason": reason])
|
||||
@@ -158,13 +255,30 @@ final class Session {
|
||||
ws.start()
|
||||
// Reconnect with this (kept in the page's memory, never on a command line).
|
||||
ws.sendJSON(["t": "hello", "r": reconnectKey])
|
||||
if let sck = capture as? SCKSource {
|
||||
sck.onChange = { [weak self] in self?.sendInfo() }
|
||||
capture.onChange = { [weak self] in self?.sendInfo() }
|
||||
// The numbers, for the viewer's overlay.
|
||||
// Adapt ten times a second; the numbers go to the viewer once a second.
|
||||
let t = DispatchSource.makeTimerSource(queue: encodeQueue)
|
||||
t.schedule(deadline: .now() + 0.1, repeating: 0.1)
|
||||
var ticks = 0
|
||||
t.setEventHandler { [weak self] in
|
||||
guard let self, !self.isStopped else { return }
|
||||
self.adapt(maxLong: maxLong)
|
||||
ticks += 1
|
||||
guard ticks % 10 == 0 else { return }
|
||||
var s = self.stats.summary()
|
||||
s.merge(self.controller.state()) { a, _ in a }
|
||||
s["t"] = "stats"
|
||||
s["bitrate"] = self.encoder.bitrate
|
||||
s["size"] = "\(self.encoder.width)×\(self.encoder.height)"
|
||||
self.ws.sendJSON(s)
|
||||
}
|
||||
t.resume()
|
||||
statsTimer = t
|
||||
capture.start(maxLong: maxLong, fps: fps) { [weak self] error in
|
||||
guard let self else { return }
|
||||
if let error {
|
||||
if (self.capture as? SCKSource)?.gone == true {
|
||||
if self.capture.gone {
|
||||
// Final, like a window closing mid-stream: the viewer closes
|
||||
// and its key is revoked, rather than retrying for ever.
|
||||
self.ws.sendJSON(["t": "closed", "reason": error])
|
||||
@@ -179,9 +293,26 @@ final class Session {
|
||||
}
|
||||
}
|
||||
|
||||
/// Applies the controller's bitrate and tier. On the encoding queue.
|
||||
private func adapt(maxLong: Int) {
|
||||
guard let bps = controller.update(now: nowUs()) else { return }
|
||||
encoder.setBitrate(bps)
|
||||
let scale = controller.scale
|
||||
if scale != appliedScale {
|
||||
appliedScale = scale
|
||||
capture.setMaxLong(max(320, Int(Double(maxLong) * scale)))
|
||||
}
|
||||
resendIfRoom() // a lower tier may let a held frame go now
|
||||
}
|
||||
|
||||
/// While someone is typing or clicking, the viewer sends a tiny message this
|
||||
/// often (ms, 0 = off), so the Frame's Wi-Fi doesn't doze between the input
|
||||
/// and the frame that answers it (power saving is on there).
|
||||
static let warmMs = Int(ProcessInfo.processInfo.environment["FRAME_MAC_VIEW_WARM"] ?? "") ?? 0
|
||||
|
||||
func sendInfo() {
|
||||
ws.sendJSON(["t": "info", "src": source.key, "title": capture.title, "app": capture.app, "codec": codec.rawValue,
|
||||
"input": source == .test || Input.allowed,
|
||||
"input": source == .test || Input.allowed, "warm": Session.warmMs,
|
||||
"aspect": Double(capture.frameRect.width / max(capture.frameRect.height, 1))])
|
||||
}
|
||||
|
||||
@@ -194,6 +325,7 @@ final class Session {
|
||||
last = nil
|
||||
lock.unlock()
|
||||
guard !was else { return }
|
||||
statsTimer?.cancel()
|
||||
capture.stop()
|
||||
encodeQueue.async { self.encoder.invalidate() }
|
||||
let owner = id
|
||||
@@ -208,9 +340,55 @@ final class Session {
|
||||
return CGPoint(x: r.minX + min(max(x, 0), 1) * r.width, y: r.minY + min(max(y, 0), 1) * r.height)
|
||||
}
|
||||
|
||||
/// Input with an id (from the viewer) has been posted: the next picture
|
||||
/// the Mac composites is its first chance to show.
|
||||
private func injected(_ m: [String: Any], kind: String) {
|
||||
guard let iid = (m["i"] as? NSNumber)?.uint32Value, iid != 0 else { return }
|
||||
let now = nowUs()
|
||||
stats.addInput(InputRecord(id: iid, kind: kind, viewer: Int64(num(m["tv"]) ?? 0), injected: now))
|
||||
lock.lock(); pendingEcho = (iid, now); lock.unlock()
|
||||
}
|
||||
|
||||
/// Timing reports from the viewer, in the agent's clock.
|
||||
private func report(_ t: String, _ m: [String: Any]) -> Bool {
|
||||
switch t {
|
||||
case "ping":
|
||||
ws.sendJSON(["t": "pong", "c": m["c"] ?? 0, "a": nowUs()])
|
||||
case "rx":
|
||||
guard let s = (m["s"] as? NSNumber)?.uint32Value else { break }
|
||||
if let r = num(m["r"]) { stats.update(s) { $0.received = Int64(r) } }
|
||||
if controller.acked(seq: s, at: nowUs()) { resendIfRoom() }
|
||||
case "fd":
|
||||
for f in m["f"] as? [[NSNumber]] ?? [] where f.count >= 4 {
|
||||
stats.update(f[0].uint32Value) {
|
||||
$0.decoded = f[1].int64Value
|
||||
$0.drawn = f[2].int64Value
|
||||
$0.vsync = f[3].int64Value
|
||||
}
|
||||
}
|
||||
let dropped = Int(num(m["drop"]) ?? 0)
|
||||
stats.withLock { stats.viewerDropped += dropped }
|
||||
case "w":
|
||||
break // keep-warm filler, see sendInfo
|
||||
case "clock":
|
||||
stats.withLock {
|
||||
stats.rtt = num(m["rtt"]) ?? 0
|
||||
stats.clockSynced = true
|
||||
if let d = m["dec"] as? String { stats.decoder = d }
|
||||
}
|
||||
default:
|
||||
return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
/// Asks the viewer to do something for a benchmark (type, click, show its overlay).
|
||||
func bench(_ m: [String: Any]) { ws.sendJSON(m.merging(["t": "bench"]) { a, _ in a }) }
|
||||
|
||||
private func handle(_ text: String) {
|
||||
guard !isStopped, let d = text.data(using: .utf8),
|
||||
let m = try? JSONSerialization.jsonObject(with: d) as? [String: Any], let t = m["t"] as? String else { return }
|
||||
if report(t, m) { return }
|
||||
let x = m["x"] as? Double ?? -1, y = m["y"] as? Double ?? -1
|
||||
if t == "ack" {
|
||||
onAck?()
|
||||
@@ -233,6 +411,7 @@ final class Session {
|
||||
default: desc = nil
|
||||
}
|
||||
capture.pointer(x: x, y: y, text: desc)
|
||||
if desc != nil { injected(m, kind: t) }
|
||||
return
|
||||
}
|
||||
DispatchQueue.main.async { [self] in
|
||||
@@ -245,13 +424,17 @@ final class Session {
|
||||
Input.focus(window: wid, pid: pid)
|
||||
}
|
||||
Input.mouse(kind, button: b, at: p, owner: id)
|
||||
if kind == "down" { injected(m, kind: "click") }
|
||||
case "wheel":
|
||||
Input.scroll(dx: m["dx"] as? Double ?? 0, dy: m["dy"] as? Double ?? 0, at: point(x, y), owner: id)
|
||||
injected(m, kind: "wheel")
|
||||
case "k":
|
||||
let down = (m["e"] as? String) == "down"
|
||||
Input.key(code: m["code"] as? String ?? "", key: m["key"] as? String ?? "",
|
||||
down: (m["e"] as? String) == "down", mods: m["mods"] as? [String] ?? [], owner: id)
|
||||
down: down, mods: m["mods"] as? [String] ?? [], owner: id)
|
||||
if down { injected(m, kind: "key") }
|
||||
case "text":
|
||||
if let s = m["s"] as? String, s.count <= 4096 { Input.text(s) }
|
||||
if let s = m["s"] as? String, s.count <= 4096 { Input.text(s); injected(m, kind: "text") }
|
||||
case "release":
|
||||
Input.releaseAll(owner: id)
|
||||
case "focus":
|
||||
@@ -292,6 +475,21 @@ final class Agent {
|
||||
/// Sources that ended for good (the window closed), for Frame Control.
|
||||
private var finished = Set<String>()
|
||||
|
||||
/// Ends every stream, then exits once separated windows are back on the
|
||||
/// Mac's own screens (they're moved back, then their displays go 0.3 s later).
|
||||
/// Idempotent: SIGTERM and the parent going away can both call it, and a
|
||||
/// session already ended by /close may still be putting its window back.
|
||||
private var quitting = false
|
||||
func quit() {
|
||||
lock.lock()
|
||||
let all = Array(sessions.values), again = quitting
|
||||
quitting = true
|
||||
lock.unlock()
|
||||
guard !again else { return }
|
||||
for s in all { s.ws.sendJSON(["t": "close"]); s.end() }
|
||||
DispatchQueue.main.asyncAfter(deadline: .now() + 0.8) { exit(0) }
|
||||
}
|
||||
|
||||
/// While anyone watches, keep the Mac's display on: a sleeping display
|
||||
/// stops being drawn, so there'd be nothing to capture (and it would lock).
|
||||
private func keepDisplayAwake(_ on: Bool) {
|
||||
@@ -355,7 +553,13 @@ final class Agent {
|
||||
switch (req.method, req.path) {
|
||||
case ("GET", "/status"):
|
||||
lock.lock()
|
||||
let list = sessions.values.map { ["id": $0.id, "src": $0.source.key, "title": $0.capture.title, "app": $0.capture.app] }
|
||||
let list = sessions.values.map { s -> [String: Any] in
|
||||
var e: [String: Any] = ["id": s.id, "src": s.source.key, "title": s.capture.title, "app": s.capture.app,
|
||||
"stats": s.stats.summary(), "bitrate": s.encoder.bitrate,
|
||||
"controller": s.controller.state()]
|
||||
if let d = (s.capture as? SeparateSource)?.displayID { e["display"] = d }
|
||||
return e
|
||||
}
|
||||
lock.unlock()
|
||||
var s = permissionsJSON()
|
||||
s["version"] = version
|
||||
@@ -386,6 +590,33 @@ final class Agent {
|
||||
for s in matching { s.ws.sendJSON(["t": "close"]) }
|
||||
DispatchQueue.global().asyncAfter(deadline: .now() + 0.3) { for s in matching { s.end() } }
|
||||
c.respond(json: ["closed": matching.count])
|
||||
case ("GET", "/stats"):
|
||||
// Frames the viewer has had time to report on, oldest first.
|
||||
let since = UInt32(req.query["since"] ?? "") ?? 0
|
||||
let settle = Int64(req.query["settle"] ?? "") ?? 1_500_000
|
||||
lock.lock()
|
||||
let list = sessions.values.filter { req.query["id"] == nil || "\($0.id)" == req.query["id"] }
|
||||
lock.unlock()
|
||||
c.respond(json: ["now": nowUs(), "streams": list.map { s -> [String: Any] in
|
||||
["id": s.id, "src": s.source.key, "frames": s.stats.settled(since: since, settle: settle).map(\.json),
|
||||
"inputs": s.stats.inputList().map(\.json), "summary": s.stats.summary(),
|
||||
"captured": s.stats.withLock { s.stats.captured }, "controller": s.controller.state(),
|
||||
"events": s.controller.eventList()]
|
||||
}])
|
||||
case ("POST", "/bench"):
|
||||
lock.lock()
|
||||
let list = sessions.values.filter { $0.source.key == req.query["src"] }
|
||||
lock.unlock()
|
||||
var m: [String: Any] = [:]
|
||||
for (k, v) in req.query where k != "k" && k != "src" { m[k] = Double(v) ?? v as Any }
|
||||
for s in list { s.bench(m) }
|
||||
c.respond(json: ["sent": list.count])
|
||||
case ("GET", "/snapshot"):
|
||||
// One JPEG of a display (for checks and thumbnails): ?display=ID
|
||||
guard let id = UInt32(req.query["display"] ?? "") else { return c.respond(400, text: "display=ID") }
|
||||
snapshot(display: id) { data, error in
|
||||
if let data { c.respond(200, body: data, type: "image/jpeg") } else { c.respond(500, text: error ?? "failed") }
|
||||
}
|
||||
case ("POST", "/permissions"):
|
||||
DispatchQueue.main.async { requestPermissions() }
|
||||
c.respond(json: permissionsJSON())
|
||||
@@ -400,7 +631,12 @@ final class Agent {
|
||||
let fps = min(max(Int(req.query["fps"] ?? "") ?? 60, 5), 120)
|
||||
let bpp = min(max(Double(req.query["bpp"] ?? "") ?? 0.1, 0.02), 0.5)
|
||||
guard let ws = c.upgrade(req) else { return }
|
||||
let capture: CaptureSource = src == .test ? TestSource() : SCKSource(src)
|
||||
let capture: CaptureSource
|
||||
switch src {
|
||||
case .test: capture = TestSource()
|
||||
case .separate(let id): capture = SeparateSource(windowID: id)
|
||||
default: capture = SCKSource(src)
|
||||
}
|
||||
lock.lock()
|
||||
// One live viewer per key: a reconnection (or a second use of an
|
||||
// unacknowledged ticket) replaces the one before.
|
||||
@@ -428,6 +664,24 @@ final class Agent {
|
||||
}
|
||||
}
|
||||
|
||||
func snapshot(display id: CGDirectDisplayID, completion: @escaping (Data?, String?) -> Void) {
|
||||
SCShareableContent.getExcludingDesktopWindows(false, onScreenWindowsOnly: true) { content, error in
|
||||
guard let d = content?.displays.first(where: { $0.displayID == id }) else {
|
||||
return completion(nil, error?.localizedDescription ?? "no such display")
|
||||
}
|
||||
let filter = SCContentFilter(display: d, excludingWindows: [])
|
||||
let cfg = SCStreamConfiguration()
|
||||
cfg.width = Int(Double(d.width) * Double(filter.pointPixelScale))
|
||||
cfg.height = Int(Double(d.height) * Double(filter.pointPixelScale))
|
||||
cfg.showsCursor = true
|
||||
SCScreenshotManager.captureImage(contentFilter: filter, configuration: cfg) { image, error in
|
||||
guard let image else { return completion(nil, error?.localizedDescription ?? "no image") }
|
||||
let rep = NSBitmapImageRep(cgImage: image)
|
||||
completion(rep.representation(using: .jpeg, properties: [.compressionFactor: 0.85]), nil)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func requestPermissions() {
|
||||
if !CGPreflightScreenCaptureAccess() { CGRequestScreenCaptureAccess() }
|
||||
if !AXIsProcessTrusted() {
|
||||
@@ -460,13 +714,19 @@ case "serve":
|
||||
}
|
||||
let page = argument("--page", in: args).map { URL(fileURLWithPath: $0) }
|
||||
let agent = Agent(token: token, page: page)
|
||||
// Quit when the parent goes away (it holds our stdin open).
|
||||
// Quit when the parent goes away (it holds our stdin open), or on SIGTERM,
|
||||
// but first end every stream, so separated windows go back where they
|
||||
// were before their displays disappear.
|
||||
if args.contains("--exit-on-eof") {
|
||||
DispatchQueue.global().async {
|
||||
while FileHandle.standardInput.availableData.count > 0 {}
|
||||
exit(0)
|
||||
agent.quit()
|
||||
}
|
||||
}
|
||||
signal(SIGTERM, SIG_IGN)
|
||||
let sigterm = DispatchSource.makeSignalSource(signal: SIGTERM, queue: .main)
|
||||
sigterm.setEventHandler { agent.quit() }
|
||||
sigterm.resume()
|
||||
let server: Server
|
||||
do {
|
||||
server = try Server(port: port) { req, c in agent.handle(req, c) }
|
||||
@@ -482,8 +742,13 @@ case "serve":
|
||||
print("frame-mac-view listening on 127.0.0.1:\(server.port ?? port)")
|
||||
fflush(stdout)
|
||||
}
|
||||
_ = NSApplication.shared // AppKit for NSScreen names and app activation
|
||||
withExtendedLifetime(server) { RunLoop.main.run() }
|
||||
// A real (background, no Dock icon) AppKit event loop, not just a run
|
||||
// loop: without it this process never hears that displays were added or
|
||||
// changed mode, so NSScreen and CGDisplayCopyDisplayMode stay stale for
|
||||
// the displays Separate mode creates.
|
||||
let app = NSApplication.shared
|
||||
app.setActivationPolicy(.prohibited)
|
||||
withExtendedLifetime((server, sigterm)) { app.run() }
|
||||
default:
|
||||
FileHandle.standardError.write(Data("usage: frame-mac-view serve|windows|displays|permissions|request-permissions\n".utf8))
|
||||
exit(2)
|
||||
|
||||
@@ -6,5 +6,6 @@ here=$(cd "$(dirname "$0")" && pwd)
|
||||
out=${1:-$here/../bin/frame-mac-view}
|
||||
mkdir -p "$(dirname "$out")"
|
||||
xcrun swiftc -O -swift-version 5 -target arm64-apple-macos14.0 \
|
||||
-import-objc-header "$here/Sources/CGVirtualDisplay.h" \
|
||||
-o "$out" "$here"/Sources/*.swift
|
||||
echo "built $out"
|
||||
Executable
+50
@@ -0,0 +1,50 @@
|
||||
#!/bin/sh
|
||||
# Wraps the agent in "Frame Mac View Lab.app" for testing from a terminal or
|
||||
# an agent session: macOS then asks for Screen Recording and Accessibility
|
||||
# for this app rather than for the terminal. Signed with an Apple Development
|
||||
# identity if there is one, so the permissions survive rebuilds.
|
||||
# mac/frame-mac-view/lab.sh build the app into mac/bin/
|
||||
# mac/frame-mac-view/lab.sh serve also start it (token, port and log in ~/Library/Caches/frame-mac-view-lab,
|
||||
# readable only by you: the token opens every window and injects input)
|
||||
# mac/frame-mac-view/lab.sh serve-only restart it without rebuilding
|
||||
set -eu
|
||||
here=$(cd "$(dirname "$0")" && pwd)
|
||||
app="$here/../bin/Frame Mac View Lab.app"
|
||||
if [ "${1:-}" != serve-only ]; then
|
||||
mkdir -p "$app/Contents/MacOS"
|
||||
sh "$here/build.sh" "$app/Contents/MacOS/frame-mac-view" >/dev/null
|
||||
cat > "$app/Contents/Info.plist" <<PLIST
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<!DOCTYPE plist PUBLIC "-//Apple//DTD PLIST 1.0//EN" "http://www.apple.com/DTDs/PropertyList-1.0.dtd">
|
||||
<plist version="1.0"><dict>
|
||||
<key>CFBundleIdentifier</key><string>com.saphid.frame-mac-view-lab</string>
|
||||
<key>CFBundleExecutable</key><string>frame-mac-view</string>
|
||||
<key>CFBundleName</key><string>Frame Mac View Lab</string>
|
||||
<key>CFBundlePackageType</key><string>APPL</string>
|
||||
<key>CFBundleShortVersionString</key><string>1</string>
|
||||
<key>LSUIElement</key><true/>
|
||||
<key>NSScreenCaptureUsageDescription</key><string>Tests streaming Mac windows into the Steam Frame.</string>
|
||||
</dict></plist>
|
||||
PLIST
|
||||
id=$(security find-identity -v -p codesigning 2>/dev/null | awk '/Apple Development/ {print $2; exit}')
|
||||
codesign -f -s "${id:--}" "$app" >/dev/null 2>&1
|
||||
echo "built $app (signed ${id:-ad hoc})"
|
||||
fi
|
||||
if [ "${1:-}" = serve ] || [ "${1:-}" = serve-only ]; then
|
||||
pkill -f "Frame Mac View Lab.app/Contents/MacOS/frame-mac-view" 2>/dev/null || true
|
||||
dir="$HOME/Library/Caches/frame-mac-view-lab"
|
||||
mkdir -p "$dir" && chmod 700 "$dir"
|
||||
umask 077
|
||||
token=$(openssl rand -hex 16)
|
||||
rm -f "$dir/token" "$dir/port" "$dir/log"
|
||||
echo "$token" > "$dir/token"
|
||||
: > "$dir/log"
|
||||
# Experiment switches pass through: FRAME_MAC_VIEW_ENCODER (Encoder.swift),
|
||||
# _VD_HZ (Separate.swift), _ADAPT (Controller.swift), _WARM (main.swift).
|
||||
open -n "$app" --env "FRAME_MAC_VIEW_TOKEN=$token" --env "FRAME_MAC_VIEW_ENCODER=${FRAME_MAC_VIEW_ENCODER:-}" --env "FRAME_MAC_VIEW_VD_HZ=${FRAME_MAC_VIEW_VD_HZ:-}" \
|
||||
--env "FRAME_MAC_VIEW_ADAPT=${FRAME_MAC_VIEW_ADAPT:-}" --env "FRAME_MAC_VIEW_WARM=${FRAME_MAC_VIEW_WARM:-}" --stdout "$dir/log" --stderr "$dir/log" \
|
||||
--args serve --port 0 --page "$here/../../ui/mac-view.html"
|
||||
for _ in 1 2 3 4 5 6 7 8 9 10; do grep -q listening "$dir/log" && break; sleep 0.3; done
|
||||
sed -n 's/.*127\.0\.0\.1:\([0-9]*\).*/\1/p' "$dir/log" | head -n 1 > "$dir/port"
|
||||
echo "lab on 127.0.0.1:$(cat "$dir/port")"
|
||||
fi
|
||||
Reference in new issue
Block a user