authorclover caruso <git@paperclover.net> 2026-10-05 19:57:26-07:00
committerclover caruso <git@paperclover.net> 2026-10-06 22:08:40-07:00
log9e9eefa9e34d0067f716831ef397d23ecae5c895
tree111af50fdb73d300290e8f01108df5e5cde0a029
parent442a2007a84dd5e45dfe33380a5d483bc0c1dc29
signature Signed by SSH key SHA256:52mNGHRsVFBDED9IAX5pe+LRWUefqTbxEReunq21QvU (~clover)

Enable separate Windows surfaces and reliable native cursor reconnects

Add bounded Windows Graphics Capture surfaces and dirty rectangles, synchronized browser composition, cached caption dragging, and observed stream metrics. Keep the experiment opt-in on Windows 10 build 19041+ and Windows 11, with desktop fallback. Use WebRTC congestion feedback and retransmission, coalesce video work, preserve reliable held-button input, and refresh native cursors on connect. Retire guest channels safely and release CUDA contexts when viewers close. Correct libvirt channel paths for long VM names. Validated 25 native GStreamer/WebRTC tests, 14 VM boundary tests, Windows x86/x64/ARM64 builds and cursor sanitizer checks, frontend check/build, production Windows 10/11 composition and reconnects, GPU release, and unchanged inference processes/driver/Nomad. Assisted-by: gpt-6.1-sol

14 files changed, 1259 insertions(+), 79 deletions(-)

dashboard/web/components/VMConsole.tsx+150-5
...@@ -5,16 +5,30 @@ import Keyboard from "lucide-solid/icons/keyboard";...@@ -5,16 +5,30 @@ import Keyboard from "lucide-solid/icons/keyboard";
5import Clipboard from "lucide-solid/icons/clipboard";5import Clipboard from "lucide-solid/icons/clipboard";
6import Maximize from "lucide-solid/icons/maximize";6import Maximize from "lucide-solid/icons/maximize";
7import RefreshCw from "lucide-solid/icons/refresh-cw";7import RefreshCw from "lucide-solid/icons/refresh-cw";
8import Check from "lucide-solid/icons/check";
8import RemoteKeyboard from "@novnc/keyboard";9import RemoteKeyboard from "@novnc/keyboard";
910
11type WindowSurface = {
12 id: number; x: number; y: number; width: number; height: number; atlas_x: number; atlas_y: number;
13 caption_left: number; caption_top: number; caption_right: number; caption_bottom: number;
14};
15type WindowFrame = { windows: WindowSurface[]; rtpTimestamp: number };
16
10export function VMConsole(props: { name: string; enabled: boolean; active: boolean; toolbar: JSX.Element; content: (toolbar: HTMLDivElement) => JSX.Element }) {17export function VMConsole(props: { name: string; enabled: boolean; active: boolean; toolbar: JSX.Element; content: (toolbar: HTMLDivElement) => JSX.Element }) {
11 const enabled = createMemo(() => props.enabled);18 const enabled = createMemo(() => props.enabled);
12 let attempt = 0;19 let attempt = 0;
13 let video!: HTMLVideoElement;20 let video!: HTMLVideoElement;
21 let canvas!: HTMLCanvasElement;
14 let cursorImage!: HTMLImageElement;22 let cursorImage!: HTMLImageElement;
15 let cursorInvert!: HTMLImageElement;23 let cursorInvert!: HTMLImageElement;
16 let panel!: HTMLDivElement;24 let panel!: HTMLDivElement;
17 let input: RTCDataChannel | undefined;25 let input: RTCDataChannel | undefined;
26 let movement: RTCDataChannel | undefined;
27 let pointerSequence = 0;
28 let windowFrame: WindowFrame | undefined;
29 let pendingFrames: WindowFrame[] = [];
30 let frameTimestamp: number | undefined;
31 let drag: { id: number; startX: number; startY: number; x: number; y: number; dx: number; dy: number; released: number } | undefined;
18 let disconnect = () => {};32 let disconnect = () => {};
19 let width = 0;33 let width = 0;
20 let height = 0;34 let height = 0;
...@@ -55,9 +69,38 @@ export function VMConsole(props: { name: string; enabled: boolean; active: boole...@@ -55,9 +69,38 @@ export function VMConsole(props: { name: string; enabled: boolean; active: boole
55 const [problem, setProblem] = createSignal("");69 const [problem, setProblem] = createSignal("");
56 const [clipboard, setClipboard] = createSignal("");70 const [clipboard, setClipboard] = createSignal("");
57 const [showClipboard, setShowClipboard] = createSignal(false);71 const [showClipboard, setShowClipboard] = createSignal(false);
72 const [windowsAvailable, setWindowsAvailable] = createSignal(false);
73 const [separateWindows, setSeparateWindows] = createSignal(false);
74 const [frameSync, setFrameSync] = createSignal(false);
75 const [windowCapture, setWindowCapture] = createSignal(false);
76 const [showStats, setShowStats] = createSignal(false);
77 const [stats, setStats] = createSignal<{ fps: number; bitrate: number; lost: number; dropped: number; delay: number; rtt: number }>();
78 const [capture, setCapture] = createSignal<{ fps: number; capture: number; scan: number; write: number }>();
58 const [tools, setTools] = createSignal<HTMLDivElement>();79 const [tools, setTools] = createSignal<HTMLDivElement>();
59 const send = (message: object) => { if (input?.readyState === "open") input.send(JSON.stringify(message)); };80 const send = (message: object) => { if (input?.readyState === "open") input.send(JSON.stringify(message)); };
60 const release = () => { buttons = 0; send({ type: "release" }); };81 const compose = () => {
82 const frame = windowFrame;
83 const composed = !!frame?.windows.length;
84 video.classList.toggle("vm-video-composed", composed);
85 canvas.hidden = !composed;
86 if (!frame || !composed || !width || !height || video.readyState < 2) return;
87 if (canvas.width !== width || canvas.height !== height) { canvas.width = width; canvas.height = height; }
88 const context = canvas.getContext("2d", { alpha: false });
89 if (!context) return;
90 context.fillStyle = "#101014";
91 context.fillRect(0, 0, width, height);
92 for (const window of [...frame.windows].reverse()) {
93 let x = window.x, y = window.y;
94 if (drag?.id === window.id) {
95 const acknowledged = Math.abs(window.x - drag.x - drag.dx) <= 3 && Math.abs(window.y - drag.y - drag.dy) <= 3;
96 if (drag.released && (acknowledged || performance.now() - drag.released > 500)) drag = undefined;
97 else { x = drag.x + drag.dx; y = drag.y + drag.dy; }
98 }
99 context.drawImage(video, window.atlas_x, window.atlas_y, window.width, window.height, x, y, window.width, window.height);
100 }
101 canvas.dataset.windows = String(frame.windows.length);
102 };
103 const release = () => { buttons = 0; drag = undefined; compose(); send({ type: "release", sequence: ++pointerSequence }); };
61 const connect = async () => {104 const connect = async () => {
62 const current = ++attempt;105 const current = ++attempt;
63 disconnect();106 disconnect();
...@@ -81,6 +124,45 @@ export function VMConsole(props: { name: string; enabled: boolean; active: boole...@@ -81,6 +124,45 @@ export function VMConsole(props: { name: string; enabled: boolean; active: boole
81 video.addEventListener("focus", focus);124 video.addEventListener("focus", focus);
82 video.addEventListener("blur", blur);125 video.addEventListener("blur", blur);
83 let peer: RTCPeerConnection | undefined;126 let peer: RTCPeerConnection | undefined;
127 let frameCallback: number | undefined;
128 let statsTimer: number | undefined;
129 let previousStats: { time: number; bytes: number; frames: number } | undefined;
130 const frameReady: VideoFrameRequestCallback = (_, metadata) => {
131 if (current !== attempt || disposed) return;
132 frameTimestamp = (metadata as VideoFrameCallbackMetadata & { rtpTimestamp?: number }).rtpTimestamp;
133 setFrameSync(Number.isFinite(frameTimestamp));
134 if (!frameSync()) {
135 frameTimestamp = undefined; windowFrame = undefined; pendingFrames = [];
136 if (separateWindows()) { setSeparateWindows(false); send({ type: "experiment", enabled: false }); }
137 }
138 while (pendingFrames[0] && frameTimestamp !== undefined && ((frameTimestamp - pendingFrames[0].rtpTimestamp) | 0) >= 0) {
139 windowFrame = pendingFrames.shift();
140 }
141 compose();
142 frameCallback = video.requestVideoFrameCallback(frameReady);
143 };
144 const measure = async () => {
145 if (!peer || current !== attempt || disposed) return;
146 try {
147 const report = await peer.getStats();
148 if (current !== attempt || disposed || peer.connectionState === "closed") return;
149 let rtt = 0;
150 report.forEach((value) => { if (value.type === "candidate-pair" && value.state === "succeeded" && value.nominated) rtt = (value.currentRoundTripTime ?? 0) * 1000; });
151 report.forEach((value) => {
152 if (value.type !== "inbound-rtp" || value.kind !== "video") return;
153 const elapsed = previousStats ? (value.timestamp - previousStats.time) / 1000 : 0;
154 setStats({ fps: elapsed ? (value.framesDecoded - previousStats!.frames) / elapsed : 0,
155 bitrate: elapsed ? (value.bytesReceived - previousStats!.bytes) * 8 / elapsed / 1000000 : 0,
156 lost: value.packetsLost ?? 0, dropped: value.framesDropped ?? 0,
157 delay: value.jitterBufferEmittedCount ? value.jitterBufferDelay / value.jitterBufferEmittedCount * 1000 : 0, rtt });
158 previousStats = { time: value.timestamp, bytes: value.bytesReceived, frames: value.framesDecoded };
159 });
160 } catch {
161 if (peer.connectionState !== "closed") return;
162 } finally {
163 if (current === attempt && !disposed && peer.connectionState !== "closed") statsTimer = window.setTimeout(measure, 1000);
164 }
165 };
84 let pending = "";166 let pending = "";
85 let messages = Promise.resolve();167 let messages = Promise.resolve();
86 const decoder = new TextDecoder();168 const decoder = new TextDecoder();
...@@ -101,6 +183,12 @@ export function VMConsole(props: { name: string; enabled: boolean; active: boole...@@ -101,6 +183,12 @@ export function VMConsole(props: { name: string; enabled: boolean; active: boole
101 socket.onclose = null;183 socket.onclose = null;
102 socket.close();184 socket.close();
103 peer?.close();185 peer?.close();
186 clearTimeout(statsTimer);
187 if (frameCallback !== undefined) video.cancelVideoFrameCallback(frameCallback);
188 windowFrame = undefined; pendingFrames = []; frameTimestamp = undefined; drag = undefined;
189 movement = undefined; pointerSequence = 0;
190 setFrameSync(false); setWindowsAvailable(false); setWindowCapture(false); setStats(undefined); setCapture(undefined);
191 compose();
104 video.onplaying = null;192 video.onplaying = null;
105 video.srcObject = null;193 video.srcObject = null;
106 clearTimeout(cursorTimer);194 clearTimeout(cursorTimer);
...@@ -125,11 +213,15 @@ export function VMConsole(props: { name: string; enabled: boolean; active: boole...@@ -125,11 +213,15 @@ export function VMConsole(props: { name: string; enabled: boolean; active: boole
125 if (message.type === "ready") {213 if (message.type === "ready") {
126 peer = new RTCPeerConnection({ iceServers: message.iceServers });214 peer = new RTCPeerConnection({ iceServers: message.iceServers });
127 const transceiver = peer.addTransceiver("video", { direction: "recvonly" });215 const transceiver = peer.addTransceiver("video", { direction: "recvonly" });
128 const codecs = RTCRtpReceiver.getCapabilities("video")?.codecs.filter((codec) => codec.mimeType === "video/H264");216 const codecs = RTCRtpReceiver.getCapabilities("video")?.codecs.filter((codec) => codec.mimeType === "video/H264" || codec.mimeType === "video/rtx");
129 if (codecs?.length) transceiver.setCodecPreferences(codecs);217 if (codecs?.length) transceiver.setCodecPreferences(codecs);
130 input = peer.createDataChannel("input");218 input = peer.createDataChannel("input");
219 movement = peer.createDataChannel("pointer", { ordered: false, maxRetransmits: 0 });
131 keyboard.onkeyevent = (key, _code, down) => send({ type: "key", key, down });220 keyboard.onkeyevent = (key, _code, down) => send({ type: "key", key, down });
132 input.onopen = () => { if (document.activeElement === video) focus(); };221 input.onopen = () => {
222 if (document.activeElement === video) focus();
223 if (windowsAvailable()) send({ type: "experiment", enabled: separateWindows() });
224 };
133 peer.onicecandidate = ({ candidate }) => { if (candidate) signal({ type: "candidate", ...candidate.toJSON() }); };225 peer.onicecandidate = ({ candidate }) => { if (candidate) signal({ type: "candidate", ...candidate.toJSON() }); };
134 peer.onconnectionstatechange = () => {226 peer.onconnectionstatechange = () => {
135 if (peer?.connectionState === "failed" || peer?.connectionState === "disconnected") {227 if (peer?.connectionState === "failed" || peer?.connectionState === "disconnected") {
...@@ -137,7 +229,11 @@ export function VMConsole(props: { name: string; enabled: boolean; active: boole...@@ -137,7 +229,11 @@ export function VMConsole(props: { name: string; enabled: boolean; active: boole
137 }229 }
138 };230 };
139 peer.ontrack = ({ track }) => {231 peer.ontrack = ({ track }) => {
232 const receiver = transceiver.receiver as RTCRtpReceiver & { jitterBufferTarget?: number };
233 if ("jitterBufferTarget" in receiver) receiver.jitterBufferTarget = 0;
140 video.srcObject = new MediaStream([track]);234 video.srcObject = new MediaStream([track]);
235 if (video.requestVideoFrameCallback) frameCallback = video.requestVideoFrameCallback(frameReady);
236 void measure();
141 void video.play().catch(() => {237 void video.play().catch(() => {
142 if (current === attempt && video.srcObject) fail("The browser couldn't play the screen. Reconnect to try again.");238 if (current === attempt && video.srcObject) fail("The browser couldn't play the screen. Reconnect to try again.");
143 });239 });
...@@ -155,6 +251,27 @@ export function VMConsole(props: { name: string; enabled: boolean; active: boole...@@ -155,6 +251,27 @@ export function VMConsole(props: { name: string; enabled: boolean; active: boole
155 height = message.height;251 height = message.height;
156 video.dataset.encoder = message.encoder;252 video.dataset.encoder = message.encoder;
157 video.dataset.source = message.source ?? "qemu";253 video.dataset.source = message.source ?? "qemu";
254 } else if (message.type === "capabilities") {
255 setWindowsAvailable(message.windows);
256 if (message.windows) send({ type: "experiment", enabled: separateWindows() });
257 } else if (message.type === "window-capture") {
258 setWindowCapture(message.available);
259 } else if (message.type === "windows") {
260 pendingFrames.push(message);
261 if (pendingFrames.length > 16) pendingFrames.shift();
262 if (frameTimestamp !== undefined && ((frameTimestamp - message.rtpTimestamp) | 0) >= 0) {
263 windowFrame = message; pendingFrames = []; compose();
264 }
265 } else if (message.type === "window-positions") {
266 if (windowFrame?.windows.length && windowFrame.windows.length === message.windows.length &&
267 windowFrame.windows.every((window) => message.windows.some((next: WindowSurface) => window.id === next.id && window.width === next.width && window.height === next.height))) {
268 windowFrame = { ...windowFrame, windows: message.windows }; compose();
269 }
270 } else if (message.type === "capture-stats") {
271 setCapture({ fps: message.frames * 1000 / message.elapsed_ms,
272 capture: message.frames ? message.capture_us / message.frames / 1000 : 0,
273 scan: message.frames ? message.scan_us / message.frames / 1000 : 0,
274 write: message.frames ? message.write_us / message.frames / 1000 : 0 });
158 } else if (message.type === "cursor") {275 } else if (message.type === "cursor") {
159 if (message.frame === 0) {276 if (message.frame === 0) {
160 clearTimeout(cursorTimer);277 clearTimeout(cursorTimer);
...@@ -193,6 +310,7 @@ export function VMConsole(props: { name: string; enabled: boolean; active: boole...@@ -193,6 +310,7 @@ export function VMConsole(props: { name: string; enabled: boolean; active: boole
193 if (current !== attempt || disposed) return;310 if (current !== attempt || disposed) return;
194 clearTimeout(timeout);311 clearTimeout(timeout);
195 setState("Connected");312 setState("Connected");
313 compose();
196 if (props.active) video.focus();314 if (props.active) video.focus();
197 };315 };
198 } catch (failure) {316 } catch (failure) {
...@@ -225,7 +343,18 @@ export function VMConsole(props: { name: string; enabled: boolean; active: boole...@@ -225,7 +343,18 @@ export function VMConsole(props: { name: string; enabled: boolean; active: boole
225 cursorInside = x >= 0 && y >= 0 && x < width && y < height;343 cursorInside = x >= 0 && y >= 0 && x < width && y < height;
226 drawCursor();344 drawCursor();
227 if ((event.type === "pointerdown" || event.type === "wheel") && (x < 0 || y < 0 || x >= width || y >= height)) return false;345 if ((event.type === "pointerdown" || event.type === "wheel") && (x < 0 || y < 0 || x >= width || y >= height)) return false;
228 send({ type: "pointer", x, y, buttons: mask });346 if (event.type === "pointerdown" && mask === 1) {
347 const window = windowFrame?.windows.find((window) => x >= window.x && y >= window.y && x < window.x + window.width && y < window.y + window.height);
348 if (window && x >= window.caption_left && x < window.caption_right && y >= window.caption_top && y < window.caption_bottom) {
349 drag = { id: window.id, startX: x, startY: y, x: window.x, y: window.y, dx: 0, dy: 0, released: 0 };
350 }
351 }
352 if (drag && event.type === "pointermove" && buttons === 1) { drag.dx = x - drag.startX; drag.dy = y - drag.startY; compose(); }
353 if (drag && event.type === "pointerup" && !mask) { drag.released = performance.now(); window.setTimeout(compose, 510); }
354 const message = JSON.stringify({ type: "pointer", x, y, buttons: mask, sequence: ++pointerSequence });
355 if (event.type === "pointermove" && !mask && movement?.readyState === "open") {
356 if (movement.bufferedAmount < 16384) movement.send(message);
357 } else if (input?.readyState === "open" && (event.type !== "pointermove" || input.bufferedAmount < 16384)) input.send(message);
229 return true;358 return true;
230 };359 };
231 const ctrlAltDel = () => {360 const ctrlAltDel = () => {
...@@ -241,6 +370,12 @@ export function VMConsole(props: { name: string; enabled: boolean; active: boole...@@ -241,6 +370,12 @@ export function VMConsole(props: { name: string; enabled: boolean; active: boole
241 <button disabled={!connected()} onClick={() => setShowClipboard(!showClipboard())}><Clipboard size={14} aria-hidden="true" />Clipboard</button>370 <button disabled={!connected()} onClick={() => setShowClipboard(!showClipboard())}><Clipboard size={14} aria-hidden="true" />Clipboard</button>
242 </VMMenu>371 </VMMenu>
243 <VMMenu label="View">372 <VMMenu label="View">
373 <button disabled={!connected() || !windowsAvailable() || !frameSync()} aria-pressed={separateWindows()}
374 onClick={() => { const enabled = !separateWindows(); setSeparateWindows(enabled); send({ type: "experiment", enabled }); }}>
375 <Check size={14} aria-hidden="true" style={{ visibility: separateWindows() ? "visible" : "hidden" }} />
376 Separate windows (experimental)
377 </button>
378 <button aria-pressed={showStats()} onClick={() => setShowStats(!showStats())}><Check size={14} aria-hidden="true" style={{ visibility: showStats() ? "visible" : "hidden" }} />Stream statistics</button>
244 <button onClick={() => panel.requestFullscreen().catch((failure) => setProblem(reason(failure)))}><Maximize size={14} aria-hidden="true" />Fullscreen</button>379 <button onClick={() => panel.requestFullscreen().catch((failure) => setProblem(reason(failure)))}><Maximize size={14} aria-hidden="true" />Fullscreen</button>
245 <button disabled={!props.enabled || state() === "Connecting…"} onClick={connect}><RefreshCw size={14} aria-hidden="true" />Reconnect</button>380 <button disabled={!props.enabled || state() === "Connecting…"} onClick={connect}><RefreshCw size={14} aria-hidden="true" />Reconnect</button>
246 </VMMenu>381 </VMMenu>
...@@ -248,6 +383,13 @@ export function VMConsole(props: { name: string; enabled: boolean; active: boole...@@ -248,6 +383,13 @@ export function VMConsole(props: { name: string; enabled: boolean; active: boole
248 <span class="vm-console-status" role="status">{props.enabled ? state() : "Off"}</span>383 <span class="vm-console-status" role="status">{props.enabled ? state() : "Off"}</span>
249 </Show>384 </Show>
250 </div>385 </div>
386 <Show when={props.active && showStats()}>
387 <div class="vm-stream-stats">
388 <Show when={stats()}>{(value) => <><span>{value().fps.toFixed(1)} FPS</span><span>{value().bitrate.toFixed(2)} Mbps</span><span>{value().lost} packets lost</span><span>{value().dropped} frames dropped</span><span>{value().rtt.toFixed(0)} ms round trip</span><span>{value().delay.toFixed(0)} ms buffer</span></>}</Show>
389 <Show when={capture()}>{(value) => <><span>Guest {value().fps.toFixed(1)} FPS</span><span>Desktop capture {value().capture.toFixed(1)} ms</span><span>Compare {value().scan.toFixed(1)} ms</span><span>Stream work {value().write.toFixed(1)} ms</span></>}</Show>
390 <Show when={separateWindows()}><span>{windowCapture() ? "Separate windows" : "Desktop fallback"}</span></Show>
391 </div>
392 </Show>
251 <Show when={props.active && problem()}><p class="error vm-console-error" role="alert">{problem()}</p></Show>393 <Show when={props.active && problem()}><p class="error vm-console-error" role="alert">{problem()}</p></Show>
252 <Show when={props.active && showClipboard()}>394 <Show when={props.active && showClipboard()}>
253 <div class="vm-clipboard">395 <div class="vm-clipboard">
...@@ -260,6 +402,7 @@ export function VMConsole(props: { name: string; enabled: boolean; active: boole...@@ -260,6 +402,7 @@ export function VMConsole(props: { name: string; enabled: boolean; active: boole
260 </div>402 </div>
261 </Show>403 </Show>
262 <div class="vm-console-screen" hidden={!props.active || !props.enabled}>404 <div class="vm-console-screen" hidden={!props.active || !props.enabled}>
405 <canvas ref={canvas} class="vm-window-canvas" hidden aria-hidden="true" />
263 <video ref={video} autoplay muted playsinline tabindex="0" aria-label="VM screen"406 <video ref={video} autoplay muted playsinline tabindex="0" aria-label="VM screen"
264 onContextMenu={(event) => event.preventDefault()}407 onContextMenu={(event) => event.preventDefault()}
265 onPointerMove={pointer}408 onPointerMove={pointer}
...@@ -267,11 +410,13 @@ export function VMConsole(props: { name: string; enabled: boolean; active: boole...@@ -267,11 +410,13 @@ export function VMConsole(props: { name: string; enabled: boolean; active: boole
267 onPointerDown={(event) => {410 onPointerDown={(event) => {
268 event.preventDefault();411 event.preventDefault();
269 const mask = buttons | (event.button === 0 ? 1 : event.button === 1 ? 2 : event.button === 2 ? 4 : 0);412 const mask = buttons | (event.button === 0 ? 1 : event.button === 1 ? 2 : event.button === 2 ? 4 : 0);
270 if (!pointer(event, mask)) return;413 if (!pointer(event)) return;
414 pointer(event, mask);
271 buttons = mask;415 buttons = mask;
272 video.focus(); video.setPointerCapture(event.pointerId);416 video.focus(); video.setPointerCapture(event.pointerId);
273 }}417 }}
274 onPointerUp={(event) => {418 onPointerUp={(event) => {
419 if (buttons) pointer(event);
275 buttons &= ~(event.button === 0 ? 1 : event.button === 1 ? 2 : event.button === 2 ? 4 : 0);420 buttons &= ~(event.button === 0 ? 1 : event.button === 1 ? 2 : event.button === 2 ? 4 : 0);
276 pointer(event);421 pointer(event);
277 if (!event.buttons && video.hasPointerCapture(event.pointerId)) video.releasePointerCapture(event.pointerId);422 if (!event.buttons && video.hasPointerCapture(event.pointerId)) video.releasePointerCapture(event.pointerId);
dashboard/web/pages/VMs.css+4-1
...@@ -93,8 +93,11 @@...@@ -93,8 +93,11 @@
93.vm-console-status { color: var(--text-2); font-size: 12px; }93.vm-console-status { color: var(--text-2); font-size: 12px; }
94.vm-password { display: grid; gap: 6px; padding: 8px 10px; color: var(--muted); font-size: 11px; }94.vm-password { display: grid; gap: 6px; padding: 8px 10px; color: var(--muted); font-size: 11px; }
95.vm-password input { width: 26ch; min-width: 0; padding: 5px 6px; color: var(--text); font: 12px monospace; }95.vm-password input { width: 26ch; min-width: 0; padding: 5px 6px; color: var(--text); font: 12px monospace; }
96.vm-console-screen { flex: 1; background: #101014; min-height: 0; }96.vm-console-screen { flex: 1; background: #101014; min-height: 0; position: relative; }
97.vm-console-screen video { display: block; width: 100%; height: 100%; object-fit: contain; outline: none; touch-action: none; }97.vm-console-screen video { display: block; width: 100%; height: 100%; object-fit: contain; outline: none; touch-action: none; }
98.vm-console-screen video.vm-video-composed { opacity: 0; }
99.vm-window-canvas { position: absolute; inset: 0; width: 100%; height: 100%; object-fit: contain; pointer-events: none; }
100.vm-stream-stats { display: flex; flex-wrap: wrap; gap: 4px 16px; padding: 6px 12px; font-size: 12px; font-variant-numeric: tabular-nums; }
98.vm-cursor { display: none; position: fixed; left: 0; top: 0; pointer-events: none; }101.vm-cursor { display: none; position: fixed; left: 0; top: 0; pointer-events: none; }
99.vm-cursor-invert { mix-blend-mode: difference; }102.vm-cursor-invert { mix-blend-mode: difference; }
100.vm-console-error { padding: 8px 12px; margin: 0; font-size: 12px; }103.vm-console-error { padding: 8px 12px; margin: 0; font-size: 12px; }
dashboard/web/styles.css+1-1
...@@ -173,7 +173,7 @@ button { font: inherit; color: inherit; }...@@ -173,7 +173,7 @@ button { font: inherit; color: inherit; }
173.nav-item:hover { background: var(--hover); color: var(--text); }173.nav-item:hover { background: var(--hover); color: var(--text); }
174.nav-item.active { background: var(--accent-wash); color: var(--accent); font-weight: 600; }174.nav-item.active { background: var(--accent-wash); color: var(--accent); font-weight: 600; }
175.nav-item .icon { width: 16px; height: 16px; flex: none; opacity: 0.9; }175.nav-item .icon { width: 16px; height: 16px; flex: none; opacity: 0.9; }
176.nav-item .label { flex: 1; min-width: 0; overflow: hidden; text-overflow: ellipsis; white-space: nowrap; }176.nav-item .label, .nav-sub .label { flex: 1; min-width: 0; overflow: hidden; text-overflow: ellipsis; white-space: nowrap; }
177.nav-item .badge { display: flex; gap: 8px; font-size: 11px; color: var(--muted); font-variant-numeric: tabular-nums; }177.nav-item .badge { display: flex; gap: 8px; font-size: 11px; color: var(--muted); font-variant-numeric: tabular-nums; }
178.nav-item .tally .status-label { gap: 4px; color: var(--text-2); }178.nav-item .tally .status-label { gap: 4px; color: var(--text-2); }
179.nav-item .tally .status { width: 8px; height: 8px; }179.nav-item .tally .status { width: 8px; height: 8px; }
guest/DESIGN.md+20-7
...@@ -1,10 +1,10 @@...@@ -1,10 +1,10 @@
1# Guest cursor stream1# Guest display stream
22
3The executable is freestanding C99. `cursor.c` contains pixel conversion and dirty-rectangle detection; `windows.c` owns Win32 capture and transport. Zig supplies the cross compiler and Windows headers, with no Zig or C runtime in the executable. x86 and x64 target Windows 7; ARM64 targets Windows 10. PE subsystem versions and imported DLLs are checked by `python3 guest/test.py`. Runtime validation currently uses Windows 10, not a Windows 7 fixture.3The executable is freestanding C99. `cursor.c` contains pixel conversion and dirty-rectangle detection; `windows.c` owns desktop capture and transport; `capture.c` implements optional window surfaces. Zig supplies the cross compiler and Windows headers, with no Zig or C runtime in the executable. x86 and x64 target Windows 7; ARM64 targets Windows 10. PE subsystem versions and imported DLLs are checked by `python3 guest/test.py`. Runtime validation uses Windows 10 and 11; Windows 7 runtime behavior is unverified.
44
5Build with `python3 guest/build.py --arch x86_64` (also `x86` and `aarch64`). In an elevated PowerShell session inside the logged-in administrator's desktop, run `./install.ps1 -Binary ./build/snowglobe-guest-x86_64.exe`. The installer uses Task Scheduler COM APIs available to PowerShell 2, installs into protected Program Files, and starts an elevated task at that user's logon. `-Remove` removes that task and executable. Elevation is required by the VirtIO serial driver's device ACL; capture must run in the interactive session, not session 0. The guest needs a compatible VirtIO serial driver and the channel named in `protocol.json`. A standard-user launcher requires a privileged handle broker; it is not part of this first backend.5Build with `python3 guest/build.py --arch x86_64` (also `x86` and `aarch64`). In an elevated PowerShell session inside the logged-in administrator's desktop, run `./install.ps1 -Binary ./build/snowglobe-guest-x86_64.exe`. The installer uses Task Scheduler COM APIs available to PowerShell 2, installs into protected Program Files, and starts an elevated task at that user's logon. `-Remove` removes that task and executable. Elevation is required by the VirtIO serial driver's device ACL; capture must run in the interactive session, not session 0. The guest needs a compatible VirtIO serial driver and the channel named in `protocol.json`. A standard-user launcher requires a privileged handle broker; it is not part of this first backend.
66
7`BitBlt` copies the primary desktop into a CPU DIB without drawing the cursor. RGB comparisons ignore the unused fourth byte. Unchanged pixels produce no screen packets; changed pixels produce one bounding rectangle. Cursor polling runs every 8 ms; screen capture runs at most 30 Hz. The host's existing encoder supplies WebRTC video. Secure desktops, RDP sessions that differ from the physical console, and capture failures return the viewer to QEMU capture.7`BitBlt` copies the primary desktop into a CPU DIB without drawing the cursor. RGB comparisons ignore the unused fourth byte. Unchanged pixels produce no screen packets; changed pixels produce one bounding rectangle. Cursor polling runs every millisecond; desktop capture targets 60 Hz. The host's existing encoder supplies WebRTC video. Secure desktops, RDP sessions that differ from the physical console, and capture failures return the viewer to QEMU capture.
88
9`GetCursorInfo` supplies the shape, visibility, and observed position. `GetIconInfo` and black/white `DrawIconEx` renders recover RGBA and binary XOR planes, including colored XOR. Nonbinary color XOR and unsupported cursor sizes explicitly fall back to the browser's default cursor, preserving video. Animation timing uses the optional `GetCursorFrameInfo` user32 export, resolved by name. Missing or unusable animation metadata falls back to a static frame. Animation starts a browser-local clock; Windows does not expose the current animation phase through this interface. Ordinary shapes use the browser cursor; XOR shapes use local layers with difference blending. Neither follows streamed positions.9`GetCursorInfo` supplies the shape, visibility, and observed position. `GetIconInfo` and black/white `DrawIconEx` renders recover RGBA and binary XOR planes, including colored XOR. Nonbinary color XOR and unsupported cursor sizes explicitly fall back to the browser's default cursor, preserving video. Animation timing uses the optional `GetCursorFrameInfo` user32 export, resolved by name. Missing or unusable animation metadata falls back to a static frame. Animation starts a browser-local clock; Windows does not expose the current animation phase through this interface. Ordinary shapes use the browser cursor; XOR shapes use local layers with difference blending. Neither follows streamed positions.
1010
...@@ -16,10 +16,23 @@ Build with `python3 guest/build.py --arch x86_64` (also `x86` and `aarch64`). In...@@ -16,10 +16,23 @@ Build with `python3 guest/build.py --arch x86_64` (also `x86` and `aarch64`). In
16- Cursor: each frame is a uint32 duration in milliseconds, width × height RGBA pixels, then one XOR bit mask byte per pixel. The schema assigns the channel bits; each selected channel inverts the underlying browser pixel. Zero dimensions and zero frames mean unsupported shape/default cursor.16- Cursor: each frame is a uint32 duration in milliseconds, width × height RGBA pixels, then one XOR bit mask byte per pixel. The schema assigns the channel bits; each selected channel inverts the underlying browser pixel. Zero dimensions and zero frames mean unsupported shape/default cursor.
17- Pointer: coordinates use signed two's-complement values in the uint32 fields. These are observations, not a cursor hit map.17- Pointer: coordinates use signed two's-complement values in the uint32 fields. These are observations, not a cursor hit map.
18- Suspend and heartbeat have empty payloads. A heartbeat arrives at least once per second; four seconds without a packet falls back to QEMU.18- Suspend and heartbeat have empty payloads. A heartbeat arrives at least once per second; four seconds without a packet falls back to QEMU.
19- Refresh is the only host-to-guest message and has an empty payload. The agent validates the entire request. The host requests it on every connection, including reconnects that reuse an open guest port.19- Refresh has an empty payload. The host requests it on every connection and resends cached cursor state when WebRTC connects.
20- Stream pauses or resumes capture; heartbeats continue while paused.
21- Stop has an empty payload. The last viewer closing or service shutdown drains guest writes until the guest acknowledges stop, then closes the channel. The guest sends no further packets after its acknowledgement and closes its handle.
22- Capabilities advertises window capture and stop-handshake support. Older agents without stop support retain the original transport-close behavior; update the agent for safe channel retirement. Experiment requests enable or disable window capture.
23- Windows lists full surface dimensions, desktop position, stacking order, and native caption bounds. Surface contains tightly packed BGRx dirty rectangles; the first update is complete.
24- Capture statistics reports observed frame counts and time spent capturing, comparing, and writing.
2025
21The host broker validates a running VM's libvirt-owned channel path. One guest reader fans out to that VM's viewers. Guest bytes are untrusted and bounded before decoding; viewer buffers are bounded independently. No guest network listener or GPU is required.26Transport writes are limited to 64 KiB so each request fits the VirtIO serial DMA descriptor budget.
2227
23## Later backends and surfaces28The host broker validates a running VM's libvirt-owned channel path. One guest reader fans out to that VM's viewers. Guest bytes are untrusted and bounded before decoding; viewer buffers are bounded independently. No guest network listener or host GPU capture process is required.
2429
25Keep the core independent of an OS and add real platform implementations when their capture and cursor APIs can be tested. X11 cursor notifications, Wayland's compositor/portal permissions, and macOS capture permissions need separate backend decisions, not empty adapters. The next Windows experiment can enumerate HWNDs and attach full surfaces, geometry, stacking, and dirty regions to this transport; occlusion, layered windows, protected content, DPI, and secure desktops need explicit behavior. A window rectangle does not imply a uniform cursor: application hit testing can change the shape anywhere inside it. Region prediction should be based on validated application information or observations with invalidation, not fabricated bounding boxes.30## Separate windows experiment
31
32Windows 10 build 19041 and later can enable separate windows from the viewer’s View menu. Windows Graphics Capture supplies full unoccluded surfaces with cursor capture disabled. Its dependencies are loaded dynamically, preserving the Windows 7 import boundary. Desktop and taskbar surfaces use `PrintWindow` because Windows Graphics Capture rejects those shell HWNDs. Changed rectangles update a single host atlas and one existing video encoder. Metadata describes the atlas and stacking; the browser composes it and moves cached surfaces immediately during native caption drags. Atlas geometry is matched to the decoded RTP frame.
33
34The experiment is off by default. Unsupported, layered, or oversized window sets fall back to desktop video. Secure desktops and session changes suspend guest capture. `PrintWindow` is synchronous; an unresponsive shell can interrupt the guest stream and trigger the host timeout. Protected content and custom caption behavior remain subject to the application’s capture support.
35
36## Later backends
37
38Keep the core independent of an OS and add real platform implementations when their capture and cursor APIs can be tested. X11 cursor notifications, Wayland's compositor/portal permissions, and macOS capture permissions need separate backend decisions, not empty adapters. A window rectangle does not imply a uniform cursor: application hit testing can change the shape anywhere inside it. Region prediction should be based on validated application information or observations with invalidation, not fabricated bounding boxes.
guest/build.py+3-1
...@@ -19,6 +19,8 @@ header += [f"#define SG_XOR_{key.upper()} {value}u" for key, value in schema["xo...@@ -19,6 +19,8 @@ header += [f"#define SG_XOR_{key.upper()} {value}u" for key, value in schema["xo
19for name, packet in schema["packets"].items():19for name, packet in schema["packets"].items():
20 header += [f"#define SG_{name.upper()} {packet['id']}u", f"#define SG_{name.upper()}_WORDS {len(packet['fields'])}u"]20 header += [f"#define SG_{name.upper()} {packet['id']}u", f"#define SG_{name.upper()}_WORDS {len(packet['fields'])}u"]
21 header += [f"#define SG_{name.upper()}_{field.upper()} {index}u" for index, field in enumerate(packet["fields"])]21 header += [f"#define SG_{name.upper()}_{field.upper()} {index}u" for index, field in enumerate(packet["fields"])]
22header += [f"#define SG_WINDOW_WORDS {len(schema['window_fields'])}u"]
23header += [f"#define SG_WINDOW_{field.upper()} {index}u" for index, field in enumerate(schema["window_fields"])]
22header += [f'#define SG_MAGIC "{schema["magic"]}"', f'#define SG_CHANNEL L"{schema["channel"]}"', "#endif"]24header += [f'#define SG_MAGIC "{schema["magic"]}"', f'#define SG_CHANNEL L"{schema["channel"]}"', "#endif"]
23(args.output / "protocol.h").write_text("\n".join(header) + "\n")25(args.output / "protocol.h").write_text("\n".join(header) + "\n")
24subsystem = "10.0" if args.arch == "aarch64" else "6.1"26subsystem = "10.0" if args.arch == "aarch64" else "6.1"
...@@ -29,7 +31,7 @@ subprocess.run([...@@ -29,7 +31,7 @@ subprocess.run([
29 "zig", "cc", "-target", args.arch + ("-windows.win10-gnu" if args.arch == "aarch64" else "-windows.win7-gnu"), "-std=c99", "-Os", "-Wall", "-Wextra", "-Werror",31 "zig", "cc", "-target", args.arch + ("-windows.win10-gnu" if args.arch == "aarch64" else "-windows.win7-gnu"), "-std=c99", "-Os", "-Wall", "-Wextra", "-Werror",
30 "-ffreestanding", "-fno-stack-protector", "-nostdlib", "-D_WIN32_WINNT=0x0601", "-DWINVER=0x0601",32 "-ffreestanding", "-fno-stack-protector", "-nostdlib", "-D_WIN32_WINNT=0x0601", "-DWINVER=0x0601",
31 "-isystem", str(library / "libc/include/any-windows-any"),33 "-isystem", str(library / "libc/include/any-windows-any"),
32 "-I", str(args.output), str(source / "cursor.c"), str(source / "windows.c"),34 "-I", str(args.output), str(source / "cursor.c"), str(source / "windows.c"), str(source / "capture.c"),
33 "-lkernel32", "-luser32", "-lgdi32", "-Wl,--entry,mainCRTStartup", "-Wl,--subsystem,windows",35 "-lkernel32", "-luser32", "-lgdi32", "-Wl,--entry,mainCRTStartup", "-Wl,--subsystem,windows",
34 "-Wl,--major-subsystem-version," + subsystem.split(".")[0], "-Wl,--minor-subsystem-version," + subsystem.split(".")[1],36 "-Wl,--major-subsystem-version," + subsystem.split(".")[0], "-Wl,--minor-subsystem-version," + subsystem.split(".")[1],
35 "-o", str(args.output / ("snowglobe-guest-" + args.arch + ".exe")),37 "-o", str(args.output / ("snowglobe-guest-" + args.arch + ".exe")),
guest/capture.c created+350
...@@ -0,0 +1,350 @@
1#define WIN32_LEAN_AND_MEAN
2#define COBJMACROS
3#define UNICODE
4#define _UNICODE
5#include <windows.h>
6#include <d3d11.h>
7#include "capture.h"
8
9struct size { INT32 width, height; };
10struct object { void **vtable; };
11/* WinRT ABI calls stay here so the Win7 executable imports no newer runtime. */
12#define METHOD(o, slot, signature) ((signature)((struct object *)(o))->vtable[slot])
13typedef HRESULT (WINAPI *query_fn)(void *, const GUID *, void **);
14typedef ULONG (WINAPI *release_fn)(void *);
15typedef HRESULT (WINAPI *close_fn)(void *);
16typedef HRESULT (WINAPI *object_fn)(void *, void **);
17typedef HRESULT (WINAPI *size_fn)(void *, struct size *);
18typedef HRESULT (WINAPI *window_fn)(void *, HWND, const GUID *, void **);
19typedef HRESULT (WINAPI *pool_fn)(void *, void *, int, int, struct size, void **);
20typedef HRESULT (WINAPI *recreate_fn)(void *, void *, int, int, struct size);
21typedef HRESULT (WINAPI *session_fn)(void *, void *, void **);
22typedef HRESULT (WINAPI *bool_fn)(void *, BOOLEAN);
23typedef HRESULT (WINAPI *initialize_fn)(int);
24typedef void (WINAPI *uninitialize_fn)(void);
25typedef HRESULT (WINAPI *string_fn)(PCWSTR, UINT32, void **);
26typedef HRESULT (WINAPI *delete_string_fn)(void *);
27typedef HRESULT (WINAPI *factory_fn)(void *, const GUID *, void **);
28typedef HRESULT (WINAPI *device_fn)(IDXGIDevice *, void **);
29typedef HRESULT (WINAPI *attribute_fn)(HWND, DWORD, PVOID, DWORD);
30
31static const GUID item_iid = {0x79c3f95b,0x31f7,0x4ec2,{0xa4,0x64,0x63,0x2e,0xf5,0xd3,0x07,0x60}};
32static const GUID interop_iid = {0x3628e81b,0x3cac,0x4c60,{0xb7,0xf4,0x23,0xce,0x0e,0x0c,0x33,0x56}};
33static const GUID pool_iid = {0x589b103f,0x6bbc,0x5df5,{0xa9,0x91,0x02,0xe2,0x8b,0x3b,0x66,0xd5}};
34static const GUID cursor_iid = {0x2c39ae40,0x7d2e,0x5044,{0x80,0x4e,0x8b,0x67,0x99,0xd4,0xcf,0x9e}};
35static const GUID close_iid = {0x30d5a829,0x7fa4,0x4026,{0x83,0xbb,0xd7,0x5b,0xae,0x4e,0xa9,0x9e}};
36static const GUID dxgi_iid = {0x54ec77fa,0x1377,0x44e6,{0x8c,0x32,0x88,0xfd,0x5f,0x44,0xc8,0x4c}};
37static const GUID access_iid = {0xa9b3d012,0x3df2,0x4ee3,{0xb8,0xd1,0x86,0x95,0xf4,0x57,0xd3,0xc1}};
38static const GUID texture_iid = {0x6f15aaf2,0xd208,0x4e89,{0x9a,0xb4,0x48,0x95,0x35,0xd3,0x4f,0x9c}};
39
40struct window {
41 HWND hwnd;
42 uint32_t id, width, height;
43 struct size pool_size;
44 void *pool, *session;
45 ID3D11Texture2D *staging;
46 uint8_t *before, *bits;
47 HDC dc;
48 HBITMAP dib, selected;
49 int initial, seen, shell;
50};
51static struct window windows[SG_MAX_WINDOWS];
52static uint32_t layout[SG_WINDOWS_WORDS + SG_MAX_WINDOWS * SG_WINDOW_WORDS];
53static uint32_t prior[SG_WINDOWS_WORDS + SG_MAX_WINDOWS * SG_WINDOW_WORDS];
54static uint32_t count, next_id;
55static size_t allocated;
56static int invalid;
57static HANDLE heap;
58/* Capture workers outlive session disposal; retain their DLLs until process exit. */
59static HMODULE runtime, d3d, dwm, composition;
60static ID3D11Device *device;
61static ID3D11DeviceContext *context;
62static void *capture_device, *interop, *pool_factory;
63static string_fn create_string;
64static delete_string_fn delete_string;
65static factory_fn get_factory;
66static uninitialize_fn uninitialize;
67static attribute_fn attribute;
68static DWORD retry;
69
70static void release(void *object) {
71 if (object) METHOD(object, 2, release_fn)(object);
72}
73static void close_object(void *object) {
74 void *closable = NULL;
75 if (object && SUCCEEDED(METHOD(object, 0, query_fn)(object, &close_iid, &closable))) {
76 METHOD(closable, 6, close_fn)(closable);
77 release(closable);
78 }
79 release(object);
80}
81static void drop(struct window *window) {
82 close_object(window->session); close_object(window->pool);
83 if (window->staging) ID3D11Texture2D_Release(window->staging);
84 if (window->selected) SelectObject(window->dc, window->selected);
85 if (window->dib) DeleteObject(window->dib);
86 if (window->dc) DeleteDC(window->dc);
87 window->dc = NULL; window->dib = window->selected = NULL; window->bits = NULL;
88 if (window->before) { allocated -= (size_t)window->width * window->height; HeapFree(heap, 0, window->before); }
89 window->hwnd = NULL; window->session = window->pool = NULL;
90 window->staging = NULL; window->before = NULL;
91 window->width = window->height = 0;
92}
93static int factory(PCWSTR name, const GUID *iid, void **result) {
94 void *string = NULL;
95 UINT32 length = 0;
96 while (name[length]) ++length;
97 if (FAILED(create_string(name, length, &string))) return 0;
98 HRESULT status = get_factory(string, iid, result);
99 delete_string(string);
100 HMODULE module;
101 return SUCCEEDED(status) && GetModuleHandleExW(GET_MODULE_HANDLE_EX_FLAG_FROM_ADDRESS | GET_MODULE_HANDLE_EX_FLAG_PIN,
102 (LPCWSTR)((struct object *)*result)->vtable[0], &module);
103}
104
105void sg_windows_stop(void) {
106 uint32_t i;
107 for (i = 0; i < SG_MAX_WINDOWS; ++i) drop(windows + i);
108 release(pool_factory); release(interop); release(capture_device);
109 if (context) ID3D11DeviceContext_Release(context);
110 if (device) ID3D11Device_Release(device);
111 pool_factory = interop = capture_device = NULL; context = NULL; device = NULL;
112 if (uninitialize) uninitialize();
113 uninitialize = NULL;
114 count = 0;
115 retry = 0;
116 for (i = 0; i < sizeof(prior) / sizeof(prior[0]); ++i) prior[i] = 0;
117}
118
119int sg_windows_start(void) {
120 initialize_fn initialize;
121 typedef HRESULT (WINAPI *create_fn)(IDXGIAdapter *, D3D_DRIVER_TYPE, HMODULE, UINT, const D3D_FEATURE_LEVEL *, UINT, UINT, ID3D11Device **, D3D_FEATURE_LEVEL *, ID3D11DeviceContext **);
122 create_fn create;
123 device_fn wrap;
124 IDXGIDevice *dxgi = NULL;
125 sg_windows_stop();
126 heap = GetProcessHeap();
127 if (!runtime) runtime = LoadLibraryW(L"combase.dll");
128 if (!d3d) d3d = LoadLibraryW(L"d3d11.dll");
129 if (!dwm) dwm = LoadLibraryW(L"dwmapi.dll");
130 if (!composition) composition = LoadLibraryW(L"dcomp.dll");
131 if (!runtime || !d3d || !dwm || !composition) goto failed;
132 initialize = (initialize_fn)GetProcAddress(runtime, "RoInitialize");
133 create_string = (string_fn)GetProcAddress(runtime, "WindowsCreateString");
134 delete_string = (delete_string_fn)GetProcAddress(runtime, "WindowsDeleteString");
135 get_factory = (factory_fn)GetProcAddress(runtime, "RoGetActivationFactory");
136 create = (create_fn)GetProcAddress(d3d, "D3D11CreateDevice");
137 wrap = (device_fn)GetProcAddress(d3d, "CreateDirect3D11DeviceFromDXGIDevice");
138 attribute = (attribute_fn)GetProcAddress(dwm, "DwmGetWindowAttribute");
139 if (!initialize || !create_string || !delete_string || !get_factory || !create || !wrap || !attribute || FAILED(initialize(1))) goto failed;
140 uninitialize = (uninitialize_fn)GetProcAddress(runtime, "RoUninitialize");
141 if (!factory(L"Windows.Graphics.Capture.GraphicsCaptureItem", &interop_iid, &interop) ||
142 !factory(L"Windows.Graphics.Capture.Direct3D11CaptureFramePool", &pool_iid, &pool_factory)) goto failed;
143 if (FAILED(create(NULL, D3D_DRIVER_TYPE_HARDWARE, NULL, D3D11_CREATE_DEVICE_BGRA_SUPPORT, NULL, 0, D3D11_SDK_VERSION, &device, NULL, &context)) &&
144 FAILED(create(NULL, D3D_DRIVER_TYPE_WARP, NULL, D3D11_CREATE_DEVICE_BGRA_SUPPORT, NULL, 0, D3D11_SDK_VERSION, &device, NULL, &context))) goto failed;
145 if (FAILED(ID3D11Device_QueryInterface(device, &dxgi_iid, (void **)&dxgi))) goto failed;
146 HRESULT status = wrap(dxgi, &capture_device);
147 IDXGIDevice_Release(dxgi);
148 if (FAILED(status)) goto failed;
149 return 1;
150failed:
151 sg_windows_stop();
152 return 0;
153}
154
155static BOOL CALLBACK enumerate(HWND hwnd, LPARAM unused) {
156 RECT rect;
157 DWORD cloaked = 0;
158 uint32_t i;
159 (void)unused;
160 if (!IsWindowVisible(hwnd) || IsIconic(hwnd)) return TRUE;
161 attribute(hwnd, 14, &cloaked, sizeof(cloaked));
162 if (cloaked || !GetWindowRect(hwnd, &rect) || rect.right <= rect.left || rect.bottom <= rect.top) return TRUE;
163 if (rect.left >= (LONG)layout[SG_WINDOWS_WIDTH] || rect.top >= (LONG)layout[SG_WINDOWS_HEIGHT] || rect.right <= 0 || rect.bottom <= 0) return TRUE;
164 if (count == SG_MAX_WINDOWS || (GetWindowLongW(hwnd, GWL_EXSTYLE) & WS_EX_LAYERED)) { invalid = 1; return FALSE; }
165 for (i = 0; i < SG_MAX_WINDOWS; ++i) if (windows[i].hwnd == hwnd) break;
166 if (i == SG_MAX_WINDOWS) {
167 for (i = 0; i < SG_MAX_WINDOWS; ++i) if (!windows[i].hwnd) break;
168 if (i == SG_MAX_WINDOWS) { invalid = 1; return FALSE; }
169 windows[i].hwnd = hwnd; windows[i].id = ++next_id; windows[i].initial = 1;
170 }
171 windows[i].seen = 1;
172 WCHAR class_name[32] = {0};
173 GetClassNameW(hwnd, class_name, 32);
174 windows[i].shell = hwnd == GetShellWindow() ||
175 !lstrcmpW(class_name, L"Shell_TrayWnd") || !lstrcmpW(class_name, L"Shell_SecondaryTrayWnd") || !lstrcmpW(class_name, L"WorkerW");
176 if (FAILED(attribute(hwnd, 9, &rect, sizeof(rect)))) GetWindowRect(hwnd, &rect);
177 uint32_t *record = layout + SG_WINDOWS_WORDS + count++ * SG_WINDOW_WORDS;
178 record[SG_WINDOW_ID] = windows[i].id;
179 record[SG_WINDOW_X] = (uint32_t)rect.left; record[SG_WINDOW_Y] = (uint32_t)rect.top;
180 record[SG_WINDOW_WIDTH] = (uint32_t)(rect.right - rect.left); record[SG_WINDOW_HEIGHT] = (uint32_t)(rect.bottom - rect.top);
181 for (i = SG_WINDOW_CAPTION_LEFT; i <= SG_WINDOW_CAPTION_BOTTOM; ++i) record[i] = 0;
182 if ((GetWindowLongW(hwnd, GWL_STYLE) & WS_CAPTION) == WS_CAPTION) {
183 TITLEBARINFO title = {0};
184 title.cbSize = sizeof(title);
185 if (GetTitleBarInfo(hwnd, &title) && title.rcTitleBar.right > title.rcTitleBar.left && !(title.rgstate[0] & STATE_SYSTEM_INVISIBLE)) {
186 record[SG_WINDOW_CAPTION_LEFT] = (uint32_t)(title.rcTitleBar.left + GetSystemMetrics(SM_CXSIZE));
187 record[SG_WINDOW_CAPTION_TOP] = (uint32_t)title.rcTitleBar.top;
188 record[SG_WINDOW_CAPTION_RIGHT] = (uint32_t)(title.rcTitleBar.right - GetSystemMetrics(SM_CXSIZE) * 3);
189 record[SG_WINDOW_CAPTION_BOTTOM] = (uint32_t)title.rcTitleBar.bottom;
190 }
191 }
192 return TRUE;
193}
194
195static int start_window(struct window *window) {
196 void *item = NULL, *cursor = NULL;
197 struct size size;
198 int success = 0;
199 if (FAILED(METHOD(interop, 3, window_fn)(interop, window->hwnd, &item_iid, &item))) goto done;
200 if (FAILED(METHOD(item, 7, size_fn)(item, &size)) || size.width <= 0 || size.height <= 0 || size.width > (INT32)SG_MAX_DIMENSION || size.height > (INT32)SG_MAX_DIMENSION) goto done;
201 if (FAILED(METHOD(pool_factory, 6, pool_fn)(pool_factory, capture_device, DXGI_FORMAT_B8G8R8A8_UNORM, 2, size, &window->pool))) goto done;
202 if (FAILED(METHOD(window->pool, 10, session_fn)(window->pool, item, &window->session))) goto done;
203 if (FAILED(METHOD(window->session, 0, query_fn)(window->session, &cursor_iid, &cursor))) goto done;
204 if (FAILED(METHOD(cursor, 7, bool_fn)(cursor, FALSE)) || FAILED(METHOD(window->session, 6, close_fn)(window->session))) goto done;
205 window->pool_size = size;
206 success = 1;
207done:
208 release(cursor); release(item);
209 return success;
210}
211
212int sg_windows_poll(sg_emit emit, int refresh) {
213 uint32_t i, n, words[SG_SURFACE_WORDS] = {0};
214 if (retry && (DWORD)(GetTickCount() - retry) < 1000 && !refresh) return 1;
215 retry = 0;
216 count = 0; invalid = 0;
217 layout[SG_WINDOWS_WIDTH] = (uint32_t)GetSystemMetrics(SM_CXSCREEN);
218 layout[SG_WINDOWS_HEIGHT] = (uint32_t)GetSystemMetrics(SM_CYSCREEN);
219 for (i = 0; i < SG_MAX_WINDOWS; ++i) windows[i].seen = 0;
220 EnumWindows(enumerate, 0);
221 for (i = 0; i < SG_MAX_WINDOWS; ++i) if (windows[i].hwnd && !windows[i].seen) drop(windows + i);
222 layout[SG_WINDOWS_COUNT] = invalid ? 0 : count;
223 n = SG_WINDOWS_WORDS + layout[SG_WINDOWS_COUNT] * SG_WINDOW_WORDS;
224 int changed = refresh;
225 for (i = 0; i < n; ++i) if (layout[i] != prior[i]) changed = 1;
226 if (changed && !emit(SG_WINDOWS, layout, n, NULL, 0)) return 0;
227 for (i = 0; i < n; ++i) prior[i] = layout[i];
228 if (invalid) return 1;
229 for (i = 0; i < SG_MAX_WINDOWS; ++i) {
230 struct window *window = windows + i;
231 void *frame = NULL, *surface = NULL, *access = NULL;
232 ID3D11Texture2D *texture = NULL;
233 D3D11_TEXTURE2D_DESC desc;
234 D3D11_MAPPED_SUBRESOURCE mapped;
235 struct size size;
236 struct sg_rect rect;
237 uint32_t y;
238 int success = 1;
239 if (!window->hwnd) continue;
240 uint32_t *record = NULL;
241 for (n = 0; n < count; ++n) {
242 uint32_t *candidate = layout + SG_WINDOWS_WORDS + n * SG_WINDOW_WORDS;
243 if (candidate[SG_WINDOW_ID] == window->id) { record = candidate; break; }
244 }
245 if (!record) continue;
246 if (window->shell) {
247 RECT bounds;
248 if (!GetWindowRect(window->hwnd, &bounds)) goto frame_failed;
249 size.width = bounds.right - bounds.left; size.height = bounds.bottom - bounds.top;
250 mapped.pData = NULL;
251 } else {
252 if (!window->pool && !start_window(window)) goto frame_failed;
253 struct size target = {(INT32)record[SG_WINDOW_WIDTH], (INT32)record[SG_WINDOW_HEIGHT]};
254 if (window->pool_size.width != target.width || window->pool_size.height != target.height) {
255 if (FAILED(METHOD(window->pool, 6, recreate_fn)(window->pool, capture_device, DXGI_FORMAT_B8G8R8A8_UNORM, 2, target))) goto frame_failed;
256 window->pool_size = target;
257 window->initial = 1;
258 }
259 if (FAILED(METHOD(window->pool, 7, object_fn)(window->pool, &frame))) goto frame_failed;
260 if (!frame) continue;
261 if (FAILED(METHOD(frame, 8, size_fn)(frame, &size))) goto frame_failed;
262 if (size.width != target.width || size.height != target.height) goto frame_done;
263 if (FAILED(METHOD(frame, 6, object_fn)(frame, &surface)) ||
264 FAILED(METHOD(surface, 0, query_fn)(surface, &access_iid, &access)) ||
265 FAILED(METHOD(access, 3, query_fn)(access, &texture_iid, (void **)&texture))) goto frame_failed;
266 }
267 if (size.width <= 0 || size.height <= 0 || size.width > (INT32)SG_MAX_DIMENSION || size.height > (INT32)SG_MAX_DIMENSION) goto frame_failed;
268 if (texture) {
269 ID3D11Texture2D_GetDesc(texture, &desc);
270 if (desc.Width < (UINT)size.width || desc.Height < (UINT)size.height) goto frame_failed;
271 }
272 if (window->width != (uint32_t)size.width || window->height != (uint32_t)size.height) {
273 if (window->staging) ID3D11Texture2D_Release(window->staging);
274 window->staging = NULL;
275 if (window->selected) SelectObject(window->dc, window->selected);
276 if (window->dib) DeleteObject(window->dib);
277 window->dib = window->selected = NULL; window->bits = NULL;
278 if (window->before) { allocated -= (size_t)window->width * window->height; HeapFree(heap, 0, window->before); window->before = NULL; }
279 if (allocated + (size_t)size.width * size.height > SG_MAX_PIXELS) goto frame_failed;
280 window->width = (uint32_t)size.width; window->height = (uint32_t)size.height;
281 window->before = HeapAlloc(heap, 0, (size_t)window->width * window->height * 4);
282 if (!window->before) goto frame_failed;
283 allocated += (size_t)window->width * window->height;
284 if (window->shell) {
285 BITMAPINFO info = {0};
286 if (!window->dc) window->dc = CreateCompatibleDC(NULL);
287 info.bmiHeader.biSize = sizeof(BITMAPINFOHEADER); info.bmiHeader.biWidth = size.width;
288 info.bmiHeader.biHeight = -size.height; info.bmiHeader.biPlanes = 1; info.bmiHeader.biBitCount = 32;
289 window->dib = CreateDIBSection(window->dc, &info, DIB_RGB_COLORS, (void **)&window->bits, NULL, 0);
290 if (!window->dc || !window->dib) goto frame_failed;
291 window->selected = SelectObject(window->dc, window->dib);
292 } else {
293 desc.Width = window->width; desc.Height = window->height;
294 desc.BindFlags = 0; desc.MiscFlags = 0; desc.Usage = D3D11_USAGE_STAGING; desc.CPUAccessFlags = D3D11_CPU_ACCESS_READ;
295 if (FAILED(ID3D11Device_CreateTexture2D(device, &desc, NULL, &window->staging))) goto frame_failed;
296 }
297 window->initial = 1;
298 }
299 if (window->shell) {
300 DWORD_PTR result;
301 if (!SendMessageTimeoutW(window->hwnd, WM_NULL, 0, 0, SMTO_ABORTIFHUNG | SMTO_BLOCK, 20, &result) || !PrintWindow(window->hwnd, window->dc, 2)) goto frame_failed;
302 GdiFlush();
303 mapped.pData = window->bits; mapped.RowPitch = window->width * 4;
304 } else {
305 D3D11_BOX box = {0, 0, 0, window->width, window->height, 1};
306 ID3D11DeviceContext_CopySubresourceRegion(context, (ID3D11Resource *)window->staging, 0, 0, 0, 0, (ID3D11Resource *)texture, 0, &box);
307 if (FAILED(ID3D11DeviceContext_Map(context, (ID3D11Resource *)window->staging, 0, D3D11_MAP_READ, 0, &mapped))) goto frame_failed;
308 }
309 rect.x = window->width; rect.y = window->height; rect.width = rect.height = 0;
310 if (refresh || window->initial) { rect.x = rect.y = 0; rect.width = window->width; rect.height = window->height; }
311 else for (y = 0; y < window->height; ++y) {
312 struct sg_rect row;
313 if (sg_difference(window->before + (size_t)y * window->width * 4, (uint8_t *)mapped.pData + (size_t)y * mapped.RowPitch, window->width, 1, &row)) {
314 if (row.x < rect.x) rect.x = row.x;
315 if (row.x + row.width > rect.width) rect.width = row.x + row.width;
316 if (y < rect.y) rect.y = y;
317 rect.height = y + 1;
318 }
319 }
320 if (rect.width && rect.height) {
321 if (!refresh && !window->initial) { rect.width -= rect.x; rect.height -= rect.y; }
322 words[SG_SURFACE_ID] = window->id; words[SG_SURFACE_WIDTH] = window->width; words[SG_SURFACE_HEIGHT] = window->height;
323 words[SG_SURFACE_X] = rect.x; words[SG_SURFACE_Y] = rect.y; words[SG_SURFACE_RECT_WIDTH] = rect.width; words[SG_SURFACE_RECT_HEIGHT] = rect.height;
324 success = emit(SG_SURFACE, words, SG_SURFACE_WORDS, (uint8_t *)mapped.pData + (size_t)rect.y * mapped.RowPitch + rect.x * 4, mapped.RowPitch);
325 for (y = rect.y; y < rect.y + rect.height; ++y) {
326 uint32_t x;
327 uint8_t *src = (uint8_t *)mapped.pData + (size_t)y * mapped.RowPitch + rect.x * 4;
328 uint8_t *dst = window->before + ((size_t)y * window->width + rect.x) * 4;
329 for (x = 0; x < rect.width * 4; ++x) dst[x] = src[x];
330 }
331 window->initial = 0;
332 }
333 if (texture) ID3D11DeviceContext_Unmap(context, (ID3D11Resource *)window->staging, 0);
334 goto frame_done;
335frame_failed:
336 invalid = 1;
337frame_done:
338 if (texture) ID3D11Texture2D_Release(texture);
339 release(access); release(surface); close_object(frame);
340 if (!success) return 0;
341 if (invalid) { drop(window); break; }
342 }
343 if (invalid) {
344 retry = GetTickCount();
345 uint32_t disabled[SG_WINDOWS_WORDS] = {layout[SG_WINDOWS_WIDTH], layout[SG_WINDOWS_HEIGHT], 0};
346 if (!emit(SG_WINDOWS, disabled, SG_WINDOWS_WORDS, NULL, 0)) return 0;
347 prior[SG_WINDOWS_COUNT] = 0;
348 }
349 return 1;
350}
guest/capture.h created+9
...@@ -0,0 +1,9 @@
1#ifndef SG_CAPTURE_H
2#define SG_CAPTURE_H
3#include "cursor.h"
4
5typedef int (*sg_emit)(uint32_t type, const uint32_t *words, uint32_t count, const uint8_t *pixels, uint32_t stride);
6int sg_windows_start(void);
7int sg_windows_poll(sg_emit emit, int refresh);
8void sg_windows_stop(void);
9#endif
guest/protocol.json+10-2
...@@ -2,13 +2,21 @@...@@ -2,13 +2,21 @@
2 "magic": "SGV1",2 "magic": "SGV1",
3 "channel": "net.paperclover.snowglobe.0",3 "channel": "net.paperclover.snowglobe.0",
4 "xor_bits": {"red": 1, "green": 2, "blue": 4},4 "xor_bits": {"red": 1, "green": 2, "blue": 4},
5 "limits": {"dimension": 4096, "pixels": 8294400, "cursor_size": 128, "cursor_frames": 64, "packet_bytes": 33554456},5 "limits": {"dimension": 4096, "pixels": 8294400, "cursor_size": 128, "cursor_frames": 64, "packet_bytes": 33554460, "windows": 32},
6 "window_fields": ["id", "x", "y", "width", "height", "caption_left", "caption_top", "caption_right", "caption_bottom"],
6 "packets": {7 "packets": {
7 "screen": {"id": 1, "fields": ["width", "height", "x", "y", "rect_width", "rect_height"]},8 "screen": {"id": 1, "fields": ["width", "height", "x", "y", "rect_width", "rect_height"]},
8 "cursor": {"id": 2, "fields": ["id", "width", "height", "hot_x", "hot_y", "frames"]},9 "cursor": {"id": 2, "fields": ["id", "width", "height", "hot_x", "hot_y", "frames"]},
9 "pointer": {"id": 3, "fields": ["id", "x", "y", "visible"]},10 "pointer": {"id": 3, "fields": ["id", "x", "y", "visible"]},
10 "suspend": {"id": 4, "fields": []},11 "suspend": {"id": 4, "fields": []},
11 "refresh": {"id": 5, "fields": []},12 "refresh": {"id": 5, "fields": []},
12 "heartbeat": {"id": 6, "fields": []}13 "heartbeat": {"id": 6, "fields": []},
14 "experiment": {"id": 7, "fields": ["enabled"]},
15 "windows": {"id": 8, "fields": ["width", "height", "count"]},
16 "surface": {"id": 9, "fields": ["id", "width", "height", "x", "y", "rect_width", "rect_height"]},
17 "capture_stats": {"id": 10, "fields": ["elapsed_ms", "frames", "capture_us", "scan_us", "write_us"]},
18 "capabilities": {"id": 11, "fields": ["windows", "stop"]},
19 "stream": {"id": 12, "fields": ["enabled"]},
20 "stop": {"id": 13, "fields": []}
13 }21 }
14}22}
guest/windows.c+128-24
...@@ -3,19 +3,29 @@...@@ -3,19 +3,29 @@
3#define _UNICODE3#define _UNICODE
4#include <windows.h>4#include <windows.h>
5#include "cursor.h"5#include "cursor.h"
6#include "capture.h"
67
7typedef HCURSOR (WINAPI *cursor_frame_fn)(HCURSOR, DWORD, DWORD, DWORD *, DWORD *);8typedef HCURSOR (WINAPI *cursor_frame_fn)(HCURSOR, DWORD, DWORD, DWORD *, DWORD *);
8static cursor_frame_fn cursor_frame;9static cursor_frame_fn cursor_frame;
9static HANDLE output;10static HANDLE output;
10static HANDLE heap;11static HANDLE heap;
11static OVERLAPPED sending;12static OVERLAPPED sending;
13static int modern_windows;
14
15static uint32_t microseconds(LONGLONG ticks, LONGLONG frequency) {
16 while (frequency > 0x7fffffff) { frequency >>= 1; ticks >>= 1; }
17 if (ticks > 0x7fffffff) return 0;
18 int value = MulDiv((int)ticks, 1000000, (int)frequency);
19 return value < 0 ? 0 : (uint32_t)value;
20}
1221
13static int write_bytes(const void *data, DWORD length) {22static int write_bytes(const void *data, DWORD length) {
14 const uint8_t *bytes = data;23 const uint8_t *bytes = data;
15 DWORD written;24 DWORD written;
16 while (length) {25 while (length) {
17 ResetEvent(sending.hEvent);26 ResetEvent(sending.hEvent);
18 if (!WriteFile(output, bytes, length, &written, &sending)) {27 DWORD amount = length > 65536 ? 65536 : length;
28 if (!WriteFile(output, bytes, amount, &written, &sending)) {
19 if (GetLastError() != ERROR_IO_PENDING) return 0;29 if (GetLastError() != ERROR_IO_PENDING) return 0;
20 if (WaitForSingleObject(sending.hEvent, 3000) != WAIT_OBJECT_0) {30 if (WaitForSingleObject(sending.hEvent, 3000) != WAIT_OBJECT_0) {
21 CancelIo(output); WaitForSingleObject(sending.hEvent, INFINITE); return 0;31 CancelIo(output); WaitForSingleObject(sending.hEvent, INFINITE); return 0;
...@@ -29,16 +39,34 @@ static int write_bytes(const void *data, DWORD length) {...@@ -29,16 +39,34 @@ static int write_bytes(const void *data, DWORD length) {
29}39}
3040
31static int packet(uint32_t type, uint32_t length, const uint32_t *words, uint32_t count) {41static int packet(uint32_t type, uint32_t length, const uint32_t *words, uint32_t count) {
32 uint8_t header[12], number[4];42 uint8_t header[12 + (SG_WINDOWS_WORDS + SG_MAX_WINDOWS * SG_WINDOW_WORDS) * 4];
33 uint32_t i;43 uint32_t i;
44 if (count > SG_WINDOWS_WORDS + SG_MAX_WINDOWS * SG_WINDOW_WORDS) return 0;
34 header[0] = SG_MAGIC[0]; header[1] = SG_MAGIC[1]; header[2] = SG_MAGIC[2]; header[3] = SG_MAGIC[3];45 header[0] = SG_MAGIC[0]; header[1] = SG_MAGIC[1]; header[2] = SG_MAGIC[2]; header[3] = SG_MAGIC[3];
35 sg_word(header + 4, type); sg_word(header + 8, length);46 sg_word(header + 4, type); sg_word(header + 8, length);
36 if (!write_bytes(header, sizeof(header))) return 0;47 for (i = 0; i < count; ++i) sg_word(header + 12 + i * 4, words[i]);
37 for (i = 0; i < count; ++i) {48 return write_bytes(header, 12 + count * 4);
38 sg_word(number, words[i]);49}
39 if (!write_bytes(number, sizeof(number))) return 0;50
51static int emit(uint32_t type, const uint32_t *words, uint32_t count, const uint8_t *pixels, uint32_t stride) {
52 uint32_t w = 0, h = 0;
53 if (pixels) {
54 w = words[type == SG_SCREEN ? SG_SCREEN_RECT_WIDTH : SG_SURFACE_RECT_WIDTH];
55 h = words[type == SG_SCREEN ? SG_SCREEN_RECT_HEIGHT : SG_SURFACE_RECT_HEIGHT];
40 }56 }
41 return 1;57 if (!packet(type, count * 4 + w * h * 4, words, count)) return 0;
58 if (!pixels) return 1;
59 if (stride == w * 4) return write_bytes(pixels, w * h * 4);
60 uint8_t *packed = HeapAlloc(heap, 0, (size_t)w * h * 4);
61 if (!packed) return 0;
62 for (uint32_t y = 0; y < h; ++y) {
63 const uint8_t *src = pixels + (size_t)y * stride;
64 uint8_t *dst = packed + (size_t)y * w * 4;
65 for (uint32_t x = 0; x < w * 4; ++x) dst[x] = src[x];
66 }
67 int success = write_bytes(packed, w * h * 4);
68 HeapFree(heap, 0, packed);
69 return success;
42}70}
4371
44static int capture_cursor(HCURSOR cursor, uint32_t id) {72static int capture_cursor(HCURSOR cursor, uint32_t id) {
...@@ -108,14 +136,27 @@ static void stream(DWORD session) {...@@ -108,14 +136,27 @@ static void stream(DWORD session) {
108 uint32_t width = 0, height = 0, id = 0;136 uint32_t width = 0, height = 0, id = 0;
109 HCURSOR last_cursor = NULL;137 HCURSOR last_cursor = NULL;
110 POINT last_point = {-1, -1};138 POINT last_point = {-1, -1};
111 DWORD last_flags = ~0u, captured = GetTickCount() - 34;139 DWORD last_flags = ~0u, cursor_retry = 0;
112 DWORD heartbeat = GetTickCount(), received = 0;140 DWORD heartbeat = GetTickCount(), received = 0;
141 LARGE_INTEGER frequency, now, next_capture, started, scanned, written;
142 uint32_t stats[SG_CAPTURE_STATS_WORDS] = {0};
143 HMODULE winmm = LoadLibraryW(L"winmm.dll");
144 typedef UINT (WINAPI *period_fn)(UINT);
145 period_fn begin_period = winmm ? (period_fn)GetProcAddress(winmm, "timeBeginPeriod") : NULL;
146 period_fn end_period = winmm ? (period_fn)GetProcAddress(winmm, "timeEndPeriod") : NULL;
113 OVERLAPPED receiving = {0};147 OVERLAPPED receiving = {0};
114 uint8_t request[12];148 uint8_t request[16];
149 DWORD request_size = 12;
115 int pending = 0;150 int pending = 0;
116 int suspended = 0, initial = 1;151 int suspended = 0, initial = 1, capturing = 1;
152 int windows_enabled = 0, windows_active = 0, cursor_supported = 1;
153 QueryPerformanceFrequency(&frequency); QueryPerformanceCounter(&next_capture);
154 if (begin_period && end_period) begin_period(1);
117 receiving.hEvent = CreateEventW(NULL, TRUE, FALSE, NULL);155 receiving.hEvent = CreateEventW(NULL, TRUE, FALSE, NULL);
118 if (!screen || !dc || !receiving.hEvent) goto done;156 if (!screen || !dc || !receiving.hEvent) goto done;
157 if (frequency.QuadPart <= 0 || frequency.QuadPart > 0xffffffffu) goto done;
158 uint32_t capabilities[SG_CAPABILITIES_WORDS] = {modern_windows, 1};
159 if (!emit(SG_CAPABILITIES, capabilities, SG_CAPABILITIES_WORDS, NULL, 0)) goto done;
119 for (;;) {160 for (;;) {
120 if (session != WTSGetActiveConsoleSessionId()) goto done;161 if (session != WTSGetActiveConsoleSessionId()) goto done;
121 CURSORINFO cursor = {sizeof(CURSORINFO), 0, NULL, {0, 0}};162 CURSORINFO cursor = {sizeof(CURSORINFO), 0, NULL, {0, 0}};
...@@ -125,7 +166,7 @@ static void stream(DWORD session) {...@@ -125,7 +166,7 @@ static void stream(DWORD session) {
125 DWORD amount = 0;166 DWORD amount = 0;
126 if (!pending) {167 if (!pending) {
127 ResetEvent(receiving.hEvent);168 ResetEvent(receiving.hEvent);
128 if (!ReadFile(output, request + received, (DWORD)sizeof(request) - received, &amount, &receiving)) {169 if (!ReadFile(output, request + received, request_size - received, &amount, &receiving)) {
129 if (GetLastError() != ERROR_IO_PENDING) goto done;170 if (GetLastError() != ERROR_IO_PENDING) goto done;
130 pending = 1;171 pending = 1;
131 } else if (!GetOverlappedResult(output, &receiving, &amount, FALSE)) goto done;172 } else if (!GetOverlappedResult(output, &receiving, &amount, FALSE)) goto done;
...@@ -134,27 +175,62 @@ static void stream(DWORD session) {...@@ -134,27 +175,62 @@ static void stream(DWORD session) {
134 if (!pending) {175 if (!pending) {
135 if (!amount) goto done;176 if (!amount) goto done;
136 received += amount;177 received += amount;
137 if (received == sizeof(request)) {178 if (received == 12 && request_size == 12) {
138 uint8_t expected[12] = {SG_MAGIC[0], SG_MAGIC[1], SG_MAGIC[2], SG_MAGIC[3], 0,0,0,0, 0,0,0,0};
139 uint32_t i;179 uint32_t i;
140 sg_word(expected + 4, SG_REFRESH);180 for (i = 0; i < 4; ++i) if (request[i] != SG_MAGIC[i]) goto done;
141 for (i = 0; i < sizeof(request); ++i) if (request[i] != expected[i]) goto done;181 if (request[4] || request[5] || request[6] || request[8] || request[9] || request[10]) goto done;
142 received = 0; initial = 1; last_cursor = NULL; last_flags = ~0u;182 if (request[7] == SG_REFRESH && !request[11]) {
183 received = 0; initial = 1; last_cursor = NULL; last_flags = ~0u;
184 if (!emit(SG_CAPABILITIES, capabilities, SG_CAPABILITIES_WORDS, NULL, 0)) goto done;
185 } else if (request[7] == SG_STOP && !request[11]) {
186 packet(SG_STOP, 0, NULL, 0);
187 goto done;
188 } else if ((request[7] == SG_EXPERIMENT || request[7] == SG_STREAM) && request[11] == 4) request_size = 16;
189 else goto done;
190 }
191 if (received == 16) {
192 if (request[12] || request[13] || request[14] || request[15] > 1) goto done;
193 if (request[7] == SG_STREAM) {
194 if (capturing != request[15]) {
195 capturing = request[15];
196 initial = 1; last_cursor = NULL; last_flags = ~0u;
197 if (!capturing) { sg_windows_stop(); windows_active = 0; }
198 }
199 } else {
200 int enabled = request[15] && modern_windows;
201 if (enabled != windows_enabled) { sg_windows_stop(); windows_active = 0; }
202 windows_enabled = enabled;
203 if (!windows_enabled) {
204 uint32_t disabled[SG_WINDOWS_WORDS] = {0};
205 if (!emit(SG_WINDOWS, disabled, SG_WINDOWS_WORDS, NULL, 0)) goto done;
206 }
207 }
208 received = 0; request_size = 12;
143 }209 }
144 }210 }
145 if ((DWORD)(GetTickCount() - heartbeat) >= 1000) {211 if ((DWORD)(GetTickCount() - heartbeat) >= 1000) {
212 stats[SG_CAPTURE_STATS_ELAPSED_MS] = GetTickCount() - heartbeat;
213 if (!emit(SG_CAPTURE_STATS, stats, SG_CAPTURE_STATS_WORDS, NULL, 0)) goto done;
214 for (uint32_t i = 0; i < SG_CAPTURE_STATS_WORDS; ++i) stats[i] = 0;
146 if (!packet(SG_HEARTBEAT, 0, NULL, 0)) goto done;215 if (!packet(SG_HEARTBEAT, 0, NULL, 0)) goto done;
147 heartbeat = GetTickCount();216 heartbeat = GetTickCount();
148 }217 }
218 if (!capturing) { Sleep(10); continue; }
149 desktop = OpenInputDesktop(0, FALSE, DESKTOP_READOBJECTS);219 desktop = OpenInputDesktop(0, FALSE, DESKTOP_READOBJECTS);
150 int active = desktop && GetUserObjectInformationW(desktop, UOI_NAME, name, sizeof(name), &length) && name[0] == L'D' && name[1] == L'e' && name[2] == L'f' && name[3] == L'a' && name[4] == L'u' && name[5] == L'l' && name[6] == L't' && !name[7];220 int active = desktop && GetUserObjectInformationW(desktop, UOI_NAME, name, sizeof(name), &length) && name[0] == L'D' && name[1] == L'e' && name[2] == L'f' && name[3] == L'a' && name[4] == L'u' && name[5] == L'l' && name[6] == L't' && !name[7];
151 if (desktop) CloseDesktop(desktop);221 if (desktop) CloseDesktop(desktop);
152 if (!active || !GetCursorInfo(&cursor)) {222 if (!active || !GetCursorInfo(&cursor)) {
153 if (!suspended && !packet(SG_SUSPEND, 0, NULL, 0)) goto done;223 if (!suspended && !packet(SG_SUSPEND, 0, NULL, 0)) goto done;
154 suspended = 1; initial = 1; last_cursor = NULL; last_flags = ~0u;224 suspended = 1; initial = 1; last_cursor = NULL; last_flags = ~0u;
225 if (windows_active) sg_windows_stop();
226 windows_active = 0;
155 Sleep(100); continue;227 Sleep(100); continue;
156 }228 }
157 int changed = cursor.hCursor && cursor.hCursor != last_cursor;229 int changed = cursor.hCursor && cursor.hCursor != last_cursor;
230 if (suspended) {
231 if (!emit(SG_CAPABILITIES, capabilities, SG_CAPABILITIES_WORDS, NULL, 0)) goto done;
232 suspended = 0;
233 }
158 if (changed) ++id;234 if (changed) ++id;
159 if (changed || cursor.flags != last_flags || cursor.ptScreenPos.x != last_point.x || cursor.ptScreenPos.y != last_point.y) {235 if (changed || cursor.flags != last_flags || cursor.ptScreenPos.x != last_point.x || cursor.ptScreenPos.y != last_point.y) {
160 uint32_t words[SG_POINTER_WORDS] = {0};236 uint32_t words[SG_POINTER_WORDS] = {0};
...@@ -163,7 +239,7 @@ static void stream(DWORD session) {...@@ -163,7 +239,7 @@ static void stream(DWORD session) {
163 if (!packet(SG_POINTER, sizeof(words), words, SG_POINTER_WORDS)) goto done;239 if (!packet(SG_POINTER, sizeof(words), words, SG_POINTER_WORDS)) goto done;
164 last_flags = cursor.flags; last_point = cursor.ptScreenPos;240 last_flags = cursor.flags; last_point = cursor.ptScreenPos;
165 }241 }
166 if (changed) {242 if (changed || (!cursor_supported && (DWORD)(GetTickCount() - cursor_retry) >= 1000)) {
167 int result = capture_cursor(cursor.hCursor, id);243 int result = capture_cursor(cursor.hCursor, id);
168 if (!result) goto done;244 if (!result) goto done;
169 if (result < 0) {245 if (result < 0) {
...@@ -172,10 +248,14 @@ static void stream(DWORD session) {...@@ -172,10 +248,14 @@ static void stream(DWORD session) {
172 if (!packet(SG_CURSOR, sizeof(words), words, SG_CURSOR_WORDS)) goto done;248 if (!packet(SG_CURSOR, sizeof(words), words, SG_CURSOR_WORDS)) goto done;
173 }249 }
174 last_cursor = cursor.hCursor;250 last_cursor = cursor.hCursor;
251 cursor_supported = result > 0; cursor_retry = GetTickCount();
175 }252 }
176 if ((DWORD)(GetTickCount() - captured) >= 33) {253 QueryPerformanceCounter(&now);
254 if (now.QuadPart >= next_capture.QuadPart) {
255 next_capture.QuadPart = now.QuadPart + (uint32_t)frequency.QuadPart / 60;
256 started = now;
177 uint32_t w = (uint32_t)GetSystemMetrics(SM_CXSCREEN), h = (uint32_t)GetSystemMetrics(SM_CYSCREEN), y;257 uint32_t w = (uint32_t)GetSystemMetrics(SM_CXSCREEN), h = (uint32_t)GetSystemMetrics(SM_CYSCREEN), y;
178 struct sg_rect rect;258 struct sg_rect rect = {0};
179 uint32_t words[SG_SCREEN_WORDS] = {0};259 uint32_t words[SG_SCREEN_WORDS] = {0};
180 if (!w || !h || w > SG_MAX_DIMENSION || h > SG_MAX_DIMENSION || w * h > SG_MAX_PIXELS) goto done;260 if (!w || !h || w > SG_MAX_DIMENSION || h > SG_MAX_DIMENSION || w * h > SG_MAX_PIXELS) goto done;
181 if (w != width || h != height) {261 if (w != width || h != height) {
...@@ -195,22 +275,42 @@ static void stream(DWORD session) {...@@ -195,22 +275,42 @@ static void stream(DWORD session) {
195 suspended = 1; initial = 1; Sleep(100); continue;275 suspended = 1; initial = 1; Sleep(100); continue;
196 }276 }
197 GdiFlush();277 GdiFlush();
278 QueryPerformanceCounter(&scanned);
279 ++stats[SG_CAPTURE_STATS_FRAMES];
280 stats[SG_CAPTURE_STATS_CAPTURE_US] += microseconds(scanned.QuadPart - started.QuadPart, frequency.QuadPart);
198 if (initial) { rect.x = 0; rect.y = 0; rect.width = width; rect.height = height; }281 if (initial) { rect.x = 0; rect.y = 0; rect.width = width; rect.height = height; }
199 else if (!sg_difference(before, pixels, width, height, &rect)) { captured = GetTickCount(); Sleep(8); continue; }282 else if (!sg_difference(before, pixels, width, height, &rect)) rect.width = rect.height = 0;
283 QueryPerformanceCounter(&written);
284 stats[SG_CAPTURE_STATS_SCAN_US] += microseconds(written.QuadPart - scanned.QuadPart, frequency.QuadPart);
200 words[SG_SCREEN_WIDTH] = width; words[SG_SCREEN_HEIGHT] = height;285 words[SG_SCREEN_WIDTH] = width; words[SG_SCREEN_HEIGHT] = height;
201 words[SG_SCREEN_X] = rect.x; words[SG_SCREEN_Y] = rect.y;286 words[SG_SCREEN_X] = rect.x; words[SG_SCREEN_Y] = rect.y;
202 words[SG_SCREEN_RECT_WIDTH] = rect.width; words[SG_SCREEN_RECT_HEIGHT] = rect.height;287 words[SG_SCREEN_RECT_WIDTH] = rect.width; words[SG_SCREEN_RECT_HEIGHT] = rect.height;
203 if (!packet(SG_SCREEN, sizeof(words) + rect.width * rect.height * 4, words, SG_SCREEN_WORDS)) goto done;288 if (rect.width && !emit(SG_SCREEN, words, SG_SCREEN_WORDS, pixels + ((size_t)rect.y * width + rect.x) * 4, width * 4)) goto done;
204 for (y = rect.y; y < rect.y + rect.height; ++y) {289 for (y = rect.y; y < rect.y + rect.height; ++y) {
205 size_t offset = ((size_t)y * width + rect.x) * 4, i;290 size_t offset = ((size_t)y * width + rect.x) * 4, i;
206 if (!write_bytes(pixels + offset, rect.width * 4)) goto done;
207 for (i = 0; i < rect.width * 4; ++i) before[offset + i] = pixels[offset + i];291 for (i = 0; i < rect.width * 4; ++i) before[offset + i] = pixels[offset + i];
208 }292 }
209 initial = 0; suspended = 0; captured = GetTickCount();293 if (windows_enabled) {
294 if (!windows_active) {
295 windows_active = sg_windows_start();
296 if (!windows_active) {
297 uint32_t disabled[SG_WINDOWS_WORDS] = {0};
298 if (!emit(SG_WINDOWS, disabled, SG_WINDOWS_WORDS, NULL, 0)) goto done;
299 windows_enabled = 0;
300 }
301 }
302 if (windows_active && !sg_windows_poll(emit, initial)) goto done;
303 }
304 QueryPerformanceCounter(&now);
305 stats[SG_CAPTURE_STATS_WRITE_US] += microseconds(now.QuadPart - written.QuadPart, frequency.QuadPart);
306 initial = 0; suspended = 0;
210 }307 }
211 Sleep(8);308 Sleep(1);
212 }309 }
213done:310done:
311 sg_windows_stop();
312 if (begin_period && end_period) end_period(1);
313 if (winmm) FreeLibrary(winmm);
214 CancelIo(output);314 CancelIo(output);
215 if (pending) WaitForSingleObject(receiving.hEvent, INFINITE);315 if (pending) WaitForSingleObject(receiving.hEvent, INFINITE);
216 if (receiving.hEvent) CloseHandle(receiving.hEvent);316 if (receiving.hEvent) CloseHandle(receiving.hEvent);
...@@ -229,6 +329,10 @@ void mainCRTStartup(void) {...@@ -229,6 +329,10 @@ void mainCRTStartup(void) {
229 sending.hEvent = CreateEventW(NULL, TRUE, FALSE, NULL);329 sending.hEvent = CreateEventW(NULL, TRUE, FALSE, NULL);
230 if (!sending.hEvent) ExitProcess(1);330 if (!sending.hEvent) ExitProcess(1);
231 SetProcessDPIAware();331 SetProcessDPIAware();
332 typedef LONG (WINAPI *version_fn)(OSVERSIONINFOW *);
333 version_fn version = (version_fn)GetProcAddress(GetModuleHandleW(L"ntdll.dll"), "RtlGetVersion");
334 OSVERSIONINFOW os = {0}; os.dwOSVersionInfoSize = sizeof(os);
335 modern_windows = version && !version(&os) && os.dwMajorVersion >= 10 && os.dwBuildNumber >= 19041;
232 symbol = GetProcAddress(GetModuleHandleW(L"user32.dll"), "GetCursorFrameInfo");336 symbol = GetProcAddress(GetModuleHandleW(L"user32.dll"), "GetCursorFrameInfo");
233 /* This optional user32 export supplies ANI step timing; static capture uses only documented APIs. */337 /* This optional user32 export supplies ANI step timing; static capture uses only documented APIs. */
234 cursor_frame = (cursor_frame_fn)symbol;338 cursor_frame = (cursor_frame_fn)symbol;
nixos/configuration.nix+1-1
...@@ -10,7 +10,7 @@ let...@@ -10,7 +10,7 @@ let
10 };10 };
11 nativePkl = pkgs.callPackage ./pkl.nix { };11 nativePkl = pkgs.callPackage ./pkl.nix { };
12 screenPython = pkgs.python3.withPackages (packages: [ packages.pygobject3 packages.gst-python ]);12 screenPython = pkgs.python3.withPackages (packages: [ packages.pygobject3 packages.gst-python ]);
13 screenPlugins = (with pkgs.gst_all_1; [ (lib.getLib gstreamer) gst-plugins-base gst-plugins-good gst-plugins-bad gst-plugins-ugly ]) ++ [ (lib.getLib pkgs.libnice) ];13 screenPlugins = (with pkgs.gst_all_1; [ (lib.getLib gstreamer) gst-plugins-base gst-plugins-good gst-plugins-bad gst-plugins-ugly gst-plugins-rs ]) ++ [ (lib.getLib pkgs.libnice) ];
14 secureOvmf = pkgs.OVMF.override {14 secureOvmf = pkgs.OVMF.override {
15 secureBoot = true;15 secureBoot = true;
16 msVarsTemplate = true;16 msVarsTemplate = true;
tools/dashboard-vms-test.py+12
...@@ -22,6 +22,18 @@ spec.loader.exec_module(host)...@@ -22,6 +22,18 @@ spec.loader.exec_module(host)
2222
2323
24class VMTests(unittest.TestCase):24class VMTests(unittest.TestCase):
25 def test_guest_channel_uses_libvirts_truncated_domain_name(self):
26 name = 'window-stream-win11-test'
27 path = f'/run/libvirt/qemu/channel/175-{name[:20]}/{vms.GUEST_CHANNEL}'
28 xml = f'<domain id="175"><devices><channel type="unix"><source mode="bind" path="{path}"/><target name="{vms.GUEST_CHANNEL}"/></channel></devices></domain>'
29 with patch.object(sys, 'argv', ['vms.py', 'guest', json.dumps({'name': name})]), \
30 patch.object(vms, 'virsh', side_effect=lambda action, *_: xml if action == 'dumpxml' else 'running'):
31 self.assertEqual(vms.main(), {'path': path})
32 for invalid in ['/tmp/display.sock', path.replace('175-', '176-'), path.replace(name[:20], name)]:
33 with self.subTest(path=invalid), patch.object(vms, 'virsh', side_effect=lambda action, *_: xml.replace(path, invalid) if action == 'dumpxml' else 'running'):
34 with self.assertRaises(ValueError):
35 vms.main()
36
25 def test_nested_installers_are_separate_from_prepared_guests(self):37 def test_nested_installers_are_separate_from_prepared_guests(self):
26 with tempfile.TemporaryDirectory() as temporary:38 with tempfile.TemporaryDirectory() as temporary:
27 root = Path(temporary).resolve()39 root = Path(temporary).resolve()
tools/vm-screen-test.py+219-3
...@@ -3,14 +3,17 @@ import asyncio...@@ -3,14 +3,17 @@ import asyncio
3import base643import base64
4import importlib.util4import importlib.util
5import json5import json
6import os
6from pathlib import Path7from pathlib import Path
7import random8import random
8import struct9import struct
10import sys
11import tempfile
9import threading12import threading
10import unittest13import unittest
11import zlib14import zlib
12from types import SimpleNamespace15from types import SimpleNamespace
13from unittest.mock import Mock16from unittest.mock import AsyncMock, Mock, patch
1417
1518
16spec = importlib.util.spec_from_file_location("vm_screen", Path(__file__).with_name("vm-screen.py"))19spec = importlib.util.spec_from_file_location("vm_screen", Path(__file__).with_name("vm-screen.py"))
...@@ -40,6 +43,34 @@ class Writer:...@@ -40,6 +43,34 @@ class Writer:
40 self.closed = True43 self.closed = True
4144
4245
46class LifecycleTests(unittest.IsolatedAsyncioTestCase):
47 async def test_sigterm_closes_connected_clients(self):
48 with tempfile.TemporaryDirectory() as directory:
49 socket = str(Path(directory) / "screen.sock")
50 process = await asyncio.create_subprocess_exec(sys.executable, str(Path(__file__).with_name("vm-screen.py")),
51 env={**os.environ, "STUDIO_VM_SCREEN_SOCKET": socket})
52 writers = []
53 try:
54 async with asyncio.timeout(5):
55 while not Path(socket).exists():
56 await asyncio.sleep(0.01)
57 for _ in range(4):
58 _, writer = await asyncio.open_unix_connection(socket)
59 writers.append(writer)
60 reader, writer = await asyncio.open_unix_connection(socket)
61 writers.append(writer)
62 self.assertTrue(json.loads(await reader.readline())["fatal"])
63 process.terminate()
64 self.assertEqual(await asyncio.wait_for(process.wait(), 5), 0)
65 finally:
66 for writer in writers:
67 writer.close()
68 await writer.wait_closed()
69 if process.returncode is None:
70 process.kill()
71 await process.wait()
72
73
43class DesktopTests(unittest.IsolatedAsyncioTestCase):74class DesktopTests(unittest.IsolatedAsyncioTestCase):
44 def setUp(self):75 def setUp(self):
45 self.reader, self.writer = asyncio.StreamReader(), Writer()76 self.reader, self.writer = asyncio.StreamReader(), Writer()
...@@ -105,6 +136,47 @@ class DesktopTests(unittest.IsolatedAsyncioTestCase):...@@ -105,6 +136,47 @@ class DesktopTests(unittest.IsolatedAsyncioTestCase):
105 self.desktop.input(message)136 self.desktop.input(message)
106137
107138
139class WindowTests(unittest.TestCase):
140 def window(self, id, x=0, y=0, width=2, height=2):
141 return dict(id=id, x=x, y=y, width=width, height=height,
142 caption_left=x, caption_top=y, caption_right=x + width, caption_bottom=y + 1)
143
144 def surface(self, atlas, id, data, x=0, y=0, width=2, height=2):
145 atlas.update(dict(id=id, width=2, height=2, x=x, y=y, rect_width=width, rect_height=height), data)
146
147 def test_partial_surfaces_and_geometry_reuse(self):
148 atlas = screen.WindowDesktop(4, 2, [self.window(1), self.window(2)])
149 self.assertFalse(atlas.complete)
150 with self.assertRaisesRegex(ValueError, "complete first"):
151 self.surface(atlas, 1, bytes(4), width=1, height=1)
152 self.surface(atlas, 1, b"a" * 16)
153 self.surface(atlas, 2, b"b" * 16)
154 self.assertTrue(atlas.complete)
155 self.surface(atlas, 1, b"c" * 4, x=1, y=1, width=1, height=1)
156 self.assertEqual(atlas.surfaces[1][2], b"a" * 12 + b"c" * 4)
157 moved = screen.WindowDesktop(4, 2, [self.window(2), self.window(1, x=1)], atlas)
158 self.assertFalse(moved.repacked)
159 self.assertIs(moved.pixels, atlas.pixels)
160 self.assertEqual(moved.pixels, b"a" * 8 + b"b" * 8 + b"a" * 4 + b"c" * 4 + b"b" * 8)
161 removed = screen.WindowDesktop(4, 2, [self.window(1)], moved)
162 self.assertTrue(removed.repacked)
163 self.assertTrue(removed.complete)
164 self.assertEqual(removed.pixels[16:24], b"a" * 4 + b"c" * 4)
165
166 def test_untrusted_window_bounds_and_atlas_budget(self):
167 for windows in [[self.window(0)], [self.window(1), self.window(1)],
168 [self.window(1, x=4097)], [self.window(1, width=0)]]:
169 with self.subTest(windows=windows), self.assertRaises(ValueError):
170 screen.WindowDesktop(4, 2, windows)
171 with self.assertRaises(OverflowError):
172 screen.WindowDesktop(4096, 1, [self.window(1, width=4096, height=2048), self.window(2, width=4096, height=2048)])
173 atlas = screen.WindowDesktop(4, 2, [self.window(1)])
174 for values, data in [(dict(id=99, width=2, height=2, x=0, y=0, rect_width=2, rect_height=2), bytes(16)),
175 (dict(id=1, width=2, height=2, x=1, y=1, rect_width=2, rect_height=2), bytes(16))]:
176 with self.assertRaises(ValueError):
177 atlas.update(values, data)
178
179
108class GuestTests(unittest.IsolatedAsyncioTestCase):180class GuestTests(unittest.IsolatedAsyncioTestCase):
109 def setUp(self):181 def setUp(self):
110 self.guest = screen.GuestDesktop("fixture")182 self.guest = screen.GuestDesktop("fixture")
...@@ -112,6 +184,7 @@ class GuestTests(unittest.IsolatedAsyncioTestCase):...@@ -112,6 +184,7 @@ class GuestTests(unittest.IsolatedAsyncioTestCase):
112 self.presented = self.resized = 0184 self.presented = self.resized = 0
113 self.viewer = Mock(desktop=SimpleNamespace(cursor={"type": "cursor", "source": "qemu"}),185 self.viewer = Mock(desktop=SimpleNamespace(cursor={"type": "cursor", "source": "qemu"}),
114 writer=Writer(), send=self.messages.append, resize=self.resize, present=self.present)186 writer=Writer(), send=self.messages.append, resize=self.resize, present=self.present)
187 self.viewer.source = self.guest
115 self.guest.screens = {self.viewer}188 self.guest.screens = {self.viewer}
116189
117 def present(self):190 def present(self):
...@@ -120,6 +193,62 @@ class GuestTests(unittest.IsolatedAsyncioTestCase):...@@ -120,6 +193,62 @@ class GuestTests(unittest.IsolatedAsyncioTestCase):
120 def resize(self):193 def resize(self):
121 self.resized += 1194 self.resized += 1
122195
196 async def test_stop_drains_packet_before_closing_channel(self):
197 reader, writer = asyncio.StreamReader(), Writer()
198 writer.get_extra_info = Mock(return_value=Mock(getsockopt=Mock(return_value=struct.pack("3i", 1, 0, 0))))
199 writer.wait_closed = AsyncMock()
200 reader.feed_data(struct.pack("!I", 2) + b"{}")
201 reader.feed_data(struct.pack("!4sII2I", b"SGV1", screen.GUEST_PROTOCOL["packets"]["capabilities"]["id"], 8, 1, 1))
202 reader.feed_data(struct.pack("!4sII6I", b"SGV1", screen.GUEST_PROTOCOL["packets"]["screen"]["id"], 40, 2, 2, 0, 0, 2, 2) + bytes(8))
203 screen.GUESTS[self.guest.name] = self.guest
204 with patch.object(asyncio, "open_unix_connection", AsyncMock(return_value=(reader, writer))):
205 task = asyncio.create_task(self.guest.run())
206 self.guest.task = task
207 async with asyncio.timeout(1):
208 while self.guest.writer is None or self.guest.windows_available is None:
209 await asyncio.sleep(0)
210 self.guest.stop()
211 self.guest.request(refresh=True)
212 stop = struct.pack("!4sII", b"SGV1", screen.GUEST_PROTOCOL["packets"]["stop"]["id"], 0)
213 await asyncio.sleep(0)
214 self.assertFalse(writer.closed)
215 self.assertFalse(task.done())
216 reader.feed_data(bytes(8) + stop)
217 await asyncio.wait_for(task, 1)
218 self.assertTrue(writer.closed)
219 self.assertEqual(self.presented, 2)
220 self.assertNotIn(self.guest.name, screen.GUESTS)
221
222 async def test_stop_during_broker_negotiation(self):
223 reader, writer = asyncio.StreamReader(), Writer()
224 writer.get_extra_info = Mock(return_value=Mock(getsockopt=Mock(return_value=struct.pack("3i", 1, 0, 0))))
225 writer.wait_closed = AsyncMock()
226 stop = struct.pack("!4sII", b"SGV1", screen.GUEST_PROTOCOL["packets"]["stop"]["id"], 0)
227 with patch.object(asyncio, "open_unix_connection", AsyncMock(return_value=(reader, writer))) as connect:
228 task = asyncio.create_task(self.guest.run())
229 self.guest.task = task
230 async with asyncio.timeout(1):
231 while not writer.data:
232 await asyncio.sleep(0)
233 self.guest.stop()
234 self.assertIsNone(self.guest.writer)
235 self.assertFalse(writer.closed)
236 reader.feed_data(struct.pack("!I", 2) + b"{}" + struct.pack("!4sII2I", b"SGV1", screen.GUEST_PROTOCOL["packets"]["capabilities"]["id"], 8, 1, 1) + stop)
237 await asyncio.wait_for(task, 1)
238 connect.assert_awaited_once()
239 self.assertTrue(writer.closed)
240 self.assertTrue(writer.data.endswith(stop))
241
242 async def test_stop_before_reader_starts(self):
243 self.guest.stop()
244 with patch.object(asyncio, "open_unix_connection", AsyncMock()) as connect:
245 await self.guest.run()
246 connect.assert_not_awaited()
247
248 async def test_unsolicited_stop_is_rejected(self):
249 with self.assertRaisesRegex(ValueError, "Unexpected"):
250 await self.packets(("stop", [], b""))
251
123 async def packets(self, *packets):252 async def packets(self, *packets):
124 reader = asyncio.StreamReader()253 reader = asyncio.StreamReader()
125 for name, words, data in packets:254 for name, words, data in packets:
...@@ -138,7 +267,7 @@ class GuestTests(unittest.IsolatedAsyncioTestCase):...@@ -138,7 +267,7 @@ class GuestTests(unittest.IsolatedAsyncioTestCase):
138 self.assertEqual(self.guest.pixels[24:28], b"\x01\x02\x03\0")267 self.assertEqual(self.guest.pixels[24:28], b"\x01\x02\x03\0")
139 self.assertEqual((self.presented, self.resized), (2, 1))268 self.assertEqual((self.presented, self.resized), (2, 1))
140 self.assertEqual([message["duration"] for message in self.messages if message["type"] == "cursor"], [17, 33])269 self.assertEqual([message["duration"] for message in self.messages if message["type"] == "cursor"], [17, 33])
141 self.assertIn("invert", self.messages[0])270 self.assertIn("invert", next(message for message in self.messages if message["type"] == "cursor"))
142 self.assertEqual(self.guest.pointer["x"], -1)271 self.assertEqual(self.guest.pointer["x"], -1)
143272
144 async def test_suspend_restores_fallback_and_requires_full_refresh(self):273 async def test_suspend_restores_fallback_and_requires_full_refresh(self):
...@@ -171,6 +300,17 @@ class GuestTests(unittest.IsolatedAsyncioTestCase):...@@ -171,6 +300,17 @@ class GuestTests(unittest.IsolatedAsyncioTestCase):
171 self.assertFalse(self.guest.screens)300 self.assertFalse(self.guest.screens)
172 self.assertEqual(len(self.guest.cursor), 2)301 self.assertEqual(len(self.guest.cursor), 2)
173302
303 async def test_capture_pauses_without_closing_guest_channel(self):
304 self.guest.writer = Writer()
305 self.guest.windows_available = False
306 self.guest.request()
307 self.assertEqual(bytes(self.guest.writer.data), struct.pack("!4sIII", b"SGV1", 12, 4, 1))
308 self.guest.writer.data.clear()
309 self.guest.screens.clear()
310 self.guest.request()
311 self.assertEqual(bytes(self.guest.writer.data), struct.pack("!4sIII", b"SGV1", 12, 4, 0))
312 self.assertFalse(self.guest.writer.closed)
313
174 async def test_missing_heartbeat_times_out_but_does_not_encode(self):314 async def test_missing_heartbeat_times_out_but_does_not_encode(self):
175 await self.packets(("screen", [1, 1, 0, 0, 1, 1], bytes(4)), ("heartbeat", [], b""))315 await self.packets(("screen", [1, 1, 0, 0, 1, 1], bytes(4)), ("heartbeat", [], b""))
176 self.assertEqual(self.presented, 1)316 self.assertEqual(self.presented, 1)
...@@ -178,6 +318,33 @@ class GuestTests(unittest.IsolatedAsyncioTestCase):...@@ -178,6 +318,33 @@ class GuestTests(unittest.IsolatedAsyncioTestCase):
178 with self.assertRaises(TimeoutError):318 with self.assertRaises(TimeoutError):
179 await self.guest.read(reader)319 await self.guest.read(reader)
180320
321 async def test_capabilities_and_cursor_resend_on_connected_viewer(self):
322 self.guest.writer = Writer()
323 self.viewer.experiment = True
324 await self.packets(("capabilities", [1], b""),
325 ("screen", [1, 1, 0, 0, 1, 1], bytes(4)),
326 ("cursor", [3, 1, 1, 0, 0, 1], struct.pack("!I", 100) + bytes(5)),
327 ("pointer", [3, 4, 5, 1], b""))
328 self.messages.clear()
329 self.guest.sync(self.viewer)
330 self.assertEqual([m["type"] for m in self.messages], ["capabilities", "cursor-position", "cursor"])
331 self.assertEqual(self.guest.writer.data[-4:], struct.pack("!I", 1))
332
333 async def test_window_layout_waits_for_full_surfaces_and_falls_back(self):
334 self.viewer.experiment = True
335 records = [1, 0, 0, 2, 2, 0, 0, 2, 1, 2, 2, 0, 2, 2, 2, 0, 4, 1]
336 await self.packets(("screen", [4, 2, 0, 0, 4, 2], bytes(32)),
337 ("windows", [4, 2, 2], struct.pack("!18I", *records)),
338 ("surface", [1, 2, 2, 0, 0, 2, 2], b"a" * 16))
339 self.assertFalse(self.guest.layers.complete)
340 self.presented = 0
341 await self.packets(("surface", [2, 2, 2, 0, 0, 2, 2], b"b" * 16))
342 self.assertTrue(self.guest.layers.complete)
343 self.assertEqual(self.presented, 1)
344 await self.packets(("windows", [4, 2, 0], b""))
345 self.assertIsNone(self.guest.layers)
346 self.assertEqual(self.messages[-1], {"type": "window-capture", "available": False})
347
181348
182class StreamTests(unittest.IsolatedAsyncioTestCase):349class StreamTests(unittest.IsolatedAsyncioTestCase):
183 async def asyncSetUp(self):350 async def asyncSetUp(self):
...@@ -187,8 +354,15 @@ class StreamTests(unittest.IsolatedAsyncioTestCase):...@@ -187,8 +354,15 @@ class StreamTests(unittest.IsolatedAsyncioTestCase):
187 self.stream = screen.Screen(self.writer, desktop)354 self.stream = screen.Screen(self.writer, desktop)
188355
189 async def asyncTearDown(self):356 async def asyncTearDown(self):
190 self.writer.close()357 pipeline = self.stream.pipeline.weak_ref()
358 encoder = self.stream.encoder.weak_ref()
191 self.stream.close()359 self.stream.close()
360 self.assertIsNone(pipeline())
361 self.assertIsNone(encoder())
362 self.stream.present()
363 self.stream.resize()
364 self.stream.remote_description(None)
365 self.stream.answer(None)
192366
193 async def test_browser_payload_and_sending_direction(self):367 async def test_browser_payload_and_sending_direction(self):
194 sdp = "\r\n".join([368 sdp = "\r\n".join([
...@@ -208,6 +382,47 @@ class StreamTests(unittest.IsolatedAsyncioTestCase):...@@ -208,6 +382,47 @@ class StreamTests(unittest.IsolatedAsyncioTestCase):
208 self.assertIn("a=rtpmap:103 H264/90000", answer["sdp"])382 self.assertIn("a=rtpmap:103 H264/90000", answer["sdp"])
209 self.assertEqual(self.stream.pipeline.get_by_name("pay").get_property("pt"), 103)383 self.assertEqual(self.stream.pipeline.get_by_name("pay").get_property("pt"), 103)
210384
385 async def test_twcc_and_retransmission_are_negotiated(self):
386 offer = "\r\n".join([
387 "v=0", "o=- 1 1 IN IP4 127.0.0.1", "s=-", "t=0 0", "a=group:BUNDLE 0",
388 "m=video 9 UDP/TLS/RTP/SAVPF 103 104", "c=IN IP4 0.0.0.0", "a=mid:0", "a=recvonly",
389 "a=rtcp-mux", "a=ice-ufrag:abcd", "a=ice-pwd:abcdefghijklmnopqrstuvwxyz",
390 "a=setup:actpass", "a=fingerprint:sha-256 " + ":".join(["00"] * 32),
391 "a=rtpmap:103 H264/90000", "a=rtpmap:104 rtx/90000", "a=fmtp:104 apt=103",
392 "a=fmtp:103 packetization-mode=1;profile-level-id=42e01f;level-asymmetry-allowed=1",
393 "a=rtcp-fb:103 nack", "a=rtcp-fb:103 nack pli", "a=rtcp-fb:103 transport-cc",
394 "a=extmap:3 " + screen.TWCC, "",
395 ])
396 self.stream.signal({"type": "offer", "sdp": offer})
397 async with asyncio.timeout(5):
398 while b'"type":"answer"' not in self.writer.data:
399 await asyncio.sleep(0.01)
400 answer = next(json.loads(line)["sdp"] for line in self.writer.data.splitlines() if json.loads(line)["type"] == "answer")
401 self.assertIn("a=rtpmap:104 rtx/90000", answer)
402 self.assertIn("a=fmtp:104 apt=103", answer)
403 primary = self.stream.pipeline.get_by_name("pay").get_property("ssrc")
404 self.assertIn(f"a=ssrc-group:FID {primary} ", answer)
405 self.assertIn(f"a=ssrc:{primary} msid:", answer)
406 self.assertIn("a=rtcp-fb:103 transport-cc", answer)
407 self.assertIsNotNone(self.stream.bwe)
408 self.stream.bwe.set_property("estimated-bitrate", 1000000)
409 self.assertEqual(self.stream.encoder.get_property("bitrate"), 1000)
410
411 async def test_unreliable_movement_cannot_restore_released_buttons(self):
412 def pointer(sequence, buttons, channel="input"):
413 self.stream.input(json.dumps(dict(type="pointer", x=1, y=1, sequence=sequence, buttons=buttons)), channel)
414 pointer(1, 1)
415 pointer(3, 0, "pointer")
416 self.assertEqual(self.stream.desktop.buttons, 1)
417 pointer(2, 0)
418 self.assertEqual(self.stream.desktop.buttons, 0)
419 before = bytes(self.stream.desktop.writer.data)
420 pointer(2, 1, "pointer")
421 self.assertEqual(self.stream.desktop.writer.data, before)
422 self.stream.input(json.dumps(dict(type="release", sequence=4)))
423 pointer(3, 1, "pointer")
424 self.assertEqual(self.stream.desktop.buttons, 0)
425
211 async def test_stalled_guest_closes_and_ignores_queued_input(self):426 async def test_stalled_guest_closes_and_ignores_queued_input(self):
212 message = json.dumps({"type": "clipboard", "text": "a" * 65536})427 message = json.dumps({"type": "clipboard", "text": "a" * 65536})
213 for _ in range(64):428 for _ in range(64):
...@@ -216,6 +431,7 @@ class StreamTests(unittest.IsolatedAsyncioTestCase):...@@ -216,6 +431,7 @@ class StreamTests(unittest.IsolatedAsyncioTestCase):
216 self.assertLess(len(self.stream.desktop.writer.data), 327680)431 self.assertLess(len(self.stream.desktop.writer.data), 327680)
217432
218 async def test_rtp_packets_fit_tunneled_network(self):433 async def test_rtp_packets_fit_tunneled_network(self):
434 await self.test_twcc_and_retransmission_are_negotiated()
219 Gst = screen.Gst435 Gst = screen.Gst
220 pipeline = self.stream.pipeline436 pipeline = self.stream.pipeline
221 pipeline.set_state(Gst.State.NULL)437 pipeline.set_state(Gst.State.NULL)
tools/vm-screen.py+351-33
...@@ -2,13 +2,16 @@...@@ -2,13 +2,16 @@
2import asyncio2import asyncio
3import base643import base64
4import contextlib4import contextlib
5from collections import deque
5import json6import json
6import logging7import logging
7import os8import os
8from pathlib import Path9from pathlib import Path
10import signal
9import socket11import socket
10import struct12import struct
11import threading13import threading
14import time
12import zlib15import zlib
1316
14import gi17import gi
...@@ -16,7 +19,8 @@ import gi...@@ -16,7 +19,8 @@ import gi
16gi.require_version("Gst", "1.0")19gi.require_version("Gst", "1.0")
17gi.require_version("GstSdp", "1.0")20gi.require_version("GstSdp", "1.0")
18gi.require_version("GstWebRTC", "1.0")21gi.require_version("GstWebRTC", "1.0")
19from gi.repository import GLib, Gst, GstSdp, GstWebRTC22gi.require_version("GstRtp", "1.0")
23from gi.repository import GLib, Gst, GstRtp, GstSdp, GstWebRTC
2024
2125
22MAX_PIXELS = 3840 * 216026MAX_PIXELS = 3840 * 2160
...@@ -24,6 +28,69 @@ STUN = os.environ.get("STUDIO_VM_STUN_SERVER", "stun://stun.cloudflare.com:3478"...@@ -24,6 +28,69 @@ STUN = os.environ.get("STUDIO_VM_STUN_SERVER", "stun://stun.cloudflare.com:3478"
24SLOTS = asyncio.Semaphore(4)28SLOTS = asyncio.Semaphore(4)
25GUEST_PROTOCOL = json.loads(Path(os.environ.get("STUDIO_VM_GUEST_PROTOCOL", Path(__file__).resolve().parent.parent / "guest/protocol.json")).read_text())29GUEST_PROTOCOL = json.loads(Path(os.environ.get("STUDIO_VM_GUEST_PROTOCOL", Path(__file__).resolve().parent.parent / "guest/protocol.json")).read_text())
26GUESTS = {}30GUESTS = {}
31TWCC = "http://www.ietf.org/id/draft-holmer-rmcat-transport-wide-cc-extensions-01"
32
33
34class WindowDesktop:
35 def __init__(self, width, height, windows, previous=None):
36 limits = GUEST_PROTOCOL["limits"]
37 if not 0 < width <= limits["dimension"] or not 0 < height <= limits["dimension"] or width * height > limits["pixels"] or not 0 < len(windows) <= limits["windows"]:
38 raise ValueError("Invalid guest window desktop")
39 self.windows = windows
40 self.surfaces = {}
41 self.width = max(width, *(window["width"] for window in windows))
42 x = y = row_height = 0
43 ids = set()
44 for window in sorted(windows, key=lambda window: window["id"]):
45 w, h = window["width"], window["height"]
46 if not window["id"] or window["id"] in ids or not 0 < w <= limits["dimension"] or not 0 < h <= limits["dimension"] or abs(window["x"]) > limits["dimension"] or abs(window["y"]) > limits["dimension"]:
47 raise ValueError("Invalid guest window bounds")
48 ids.add(window["id"])
49 if x + w > self.width:
50 x, y, row_height = 0, y + row_height, 0
51 window["atlas_x"], window["atlas_y"] = x, y
52 x, row_height = x + w, max(row_height, h)
53 surface = previous.surfaces.get(window["id"]) if previous else None
54 self.surfaces[window["id"]] = surface if surface and (surface[0], surface[1]) == (w, h) else (w, h, None)
55 self.height = y + row_height
56 if self.height > limits["dimension"] or self.width * self.height > limits["pixels"]:
57 raise OverflowError("Guest windows exceed the video atlas")
58 self.repacked = not previous or (self.width, self.height, [(w["id"], w["width"], w["height"]) for w in sorted(windows, key=lambda w: w["id"])]) != (
59 previous.width, previous.height, [(w["id"], w["width"], w["height"]) for w in sorted(previous.windows, key=lambda w: w["id"])])
60 self.pixels = bytearray(self.width * self.height * 4) if self.repacked else previous.pixels
61 for window in windows:
62 surface = self.surfaces[window["id"]]
63 if self.repacked and surface[2] is not None:
64 self.copy(window, 0, 0, window["width"], window["height"], surface[2])
65
66 @property
67 def complete(self):
68 return all(surface[2] is not None for surface in self.surfaces.values())
69
70 def copy(self, window, x, y, width, height, pixels):
71 for row in range(height):
72 offset = ((window["atlas_y"] + y + row) * self.width + window["atlas_x"] + x) * 4
73 self.pixels[offset:offset + width * 4] = pixels[row * width * 4:(row + 1) * width * 4]
74
75 def update(self, values, data):
76 window = next((window for window in self.windows if window["id"] == values["id"]), None)
77 if not window:
78 raise ValueError("Unknown guest window")
79 width, height = window["width"], window["height"]
80 x, y, w, h = values["x"], values["y"], values["rect_width"], values["rect_height"]
81 if (values["width"], values["height"]) != (width, height) or not w or not h or x + w > width or y + h > height or len(data) != w * h * 4:
82 raise ValueError("Invalid guest window rectangle")
83 pixels = self.surfaces[window["id"]][2]
84 if pixels is None:
85 if (x, y, w, h) != (0, 0, width, height):
86 raise ValueError("Guest window needs a complete first frame")
87 pixels = bytearray(data)
88 self.surfaces[window["id"]] = (width, height, pixels)
89 else:
90 for row in range(h):
91 offset = ((y + row) * width + x) * 4
92 pixels[offset:offset + w * 4] = data[row * w * 4:(row + 1) * w * 4]
93 self.copy(window, x, y, w, h, data)
2794
2895
29def png(width, height, pixels):96def png(width, height, pixels):
...@@ -41,14 +108,48 @@ class GuestDesktop:...@@ -41,14 +108,48 @@ class GuestDesktop:
41 self.pixels = bytearray()108 self.pixels = bytearray()
42 self.cursor = []109 self.cursor = []
43 self.pointer = None110 self.pointer = None
111 self.layers = None
112 self.writer = None
113 self.windows_available = None
114 self.stop_supported = None
115 self.stopping = False
116
117 def request(self, refresh=False):
118 if not self.stopping and self.writer and not self.writer.is_closing():
119 if refresh:
120 self.writer.write(struct.pack("!4sII", GUEST_PROTOCOL["magic"].encode(), GUEST_PROTOCOL["packets"]["refresh"]["id"], 0))
121 if self.windows_available is not None:
122 self.writer.write(struct.pack("!4sIII", GUEST_PROTOCOL["magic"].encode(), GUEST_PROTOCOL["packets"]["stream"]["id"], 4, bool(self.screens)))
123 if self.windows_available:
124 enabled = any(screen.experiment for screen in self.screens)
125 self.writer.write(struct.pack("!4sIII", GUEST_PROTOCOL["magic"].encode(), GUEST_PROTOCOL["packets"]["experiment"]["id"], 4, enabled))
126
127 def stop(self):
128 self.stopping = True
129 if self.writer and not self.writer.is_closing():
130 if self.stop_supported:
131 self.writer.write(struct.pack("!4sII", GUEST_PROTOCOL["magic"].encode(), GUEST_PROTOCOL["packets"]["stop"]["id"], 0))
132 elif self.stop_supported is False:
133 self.writer.close()
134
135 def sync(self, screen):
136 screen.send({"type": "capabilities", "windows": bool(self.windows_available)})
137 if self.pixels:
138 if self.pointer:
139 screen.send(self.pointer)
140 for cursor in self.cursor:
141 screen.send(cursor)
44142
45 def suspend(self):143 def suspend(self):
46 active = bool(self.pixels)144 active = bool(self.pixels)
47 self.pixels = bytearray()145 self.pixels = bytearray()
48 self.cursor.clear()146 self.cursor.clear()
49 self.pointer = None147 self.pointer = None
148 self.layers = None
149 for screen in tuple(self.screens):
150 screen.send({"type": "capabilities", "windows": False})
50 if active:151 if active:
51 for screen in self.screens:152 for screen in tuple(self.screens):
52 screen.resize()153 screen.resize()
53 screen.send(screen.desktop.cursor or {"type": "cursor", "source": "qemu", "frame": 0, "fallback": True, "x": 0, "y": 0})154 screen.send(screen.desktop.cursor or {"type": "cursor", "source": "qemu", "frame": 0, "fallback": True, "x": 0, "y": 0})
54 screen.present()155 screen.present()
...@@ -62,11 +163,17 @@ class GuestDesktop:...@@ -62,11 +163,17 @@ class GuestDesktop:
62 if magic != GUEST_PROTOCOL["magic"].encode() or kind not in packets or length > limits["packet_bytes"]:163 if magic != GUEST_PROTOCOL["magic"].encode() or kind not in packets or length > limits["packet_bytes"]:
63 raise ValueError("Invalid guest display packet")164 raise ValueError("Invalid guest display packet")
64 name, fields = packets[kind]165 name, fields = packets[kind]
166 if name == "capabilities" and length == 4:
167 fields = fields[:1]
65 if length < len(fields) * 4:168 if length < len(fields) * 4:
66 raise ValueError("Incomplete guest display packet")169 raise ValueError("Incomplete guest display packet")
67 payload = await reader.readexactly(length)170 payload = await reader.readexactly(length)
68 values = dict(zip(fields, struct.unpack("!" + "I" * len(fields), payload[:len(fields) * 4])))171 values = dict(zip(fields, struct.unpack("!" + "I" * len(fields), payload[:len(fields) * 4])))
69 data = payload[len(fields) * 4:]172 data = payload[len(fields) * 4:]
173 if self.stop_supported is None and name != "capabilities":
174 self.stop_supported = False
175 if self.stopping:
176 return
70 if name == "screen":177 if name == "screen":
71 width, height = values["width"], values["height"]178 width, height = values["width"], values["height"]
72 x, y, w, h = values["x"], values["y"], values["rect_width"], values["rect_height"]179 x, y, w, h = values["x"], values["y"], values["rect_width"], values["rect_height"]
...@@ -84,18 +191,16 @@ class GuestDesktop:...@@ -84,18 +191,16 @@ class GuestDesktop:
84 for screen in tuple(self.screens):191 for screen in tuple(self.screens):
85 if resized:192 if resized:
86 screen.resize()193 screen.resize()
87 if self.pointer:194 self.sync(screen)
88 screen.send(self.pointer)195 if screen.source is self:
89 for cursor in self.cursor:196 screen.present()
90 screen.send(cursor)
91 screen.present()
92 await asyncio.sleep(0)197 await asyncio.sleep(0)
93 elif name == "cursor":198 elif name == "cursor":
94 width, height, count = values["width"], values["height"], values["frames"]199 width, height, count = values["width"], values["height"], values["frames"]
95 if not width and not height and not count and not data and not values["hot_x"] and not values["hot_y"]:200 if not width and not height and not count and not data and not values["hot_x"] and not values["hot_y"]:
96 self.cursor = [{"type": "cursor", "source": "guest", "id": values["id"], "width": 0, "height": 0, "frame": 0, "frames": 0, "fallback": True, "x": 0, "y": 0}]201 self.cursor = [{"type": "cursor", "source": "guest", "id": values["id"], "width": 0, "height": 0, "frame": 0, "frames": 0, "fallback": True, "x": 0, "y": 0}]
97 if self.pixels:202 if self.pixels:
98 for screen in self.screens:203 for screen in tuple(self.screens):
99 screen.send(self.cursor[0])204 screen.send(self.cursor[0])
100 continue205 continue
101 pixels = width * height206 pixels = width * height
...@@ -131,7 +236,7 @@ class GuestDesktop:...@@ -131,7 +236,7 @@ class GuestDesktop:
131 self.pointer = {"type": "cursor-position", "id": values["id"], "visible": bool(values["visible"]),236 self.pointer = {"type": "cursor-position", "id": values["id"], "visible": bool(values["visible"]),
132 "x": (values["x"] ^ 0x80000000) - 0x80000000, "y": (values["y"] ^ 0x80000000) - 0x80000000}237 "x": (values["x"] ^ 0x80000000) - 0x80000000, "y": (values["y"] ^ 0x80000000) - 0x80000000}
133 if self.pixels:238 if self.pixels:
134 for screen in self.screens:239 for screen in tuple(self.screens):
135 screen.send(self.pointer)240 screen.send(self.pointer)
136 elif name == "suspend":241 elif name == "suspend":
137 if data:242 if data:
...@@ -140,11 +245,70 @@ class GuestDesktop:...@@ -140,11 +245,70 @@ class GuestDesktop:
140 elif name == "heartbeat":245 elif name == "heartbeat":
141 if data:246 if data:
142 raise ValueError("Invalid guest display heartbeat")247 raise ValueError("Invalid guest display heartbeat")
248 elif name == "stop":
249 if data or not self.stopping:
250 raise ValueError("Unexpected guest display stop")
251 return
252 elif name == "capabilities":
253 if data or values["windows"] > 1 or values.get("stop", 0) > 1:
254 raise ValueError("Invalid guest display capabilities")
255 self.windows_available = bool(values["windows"])
256 self.stop_supported = bool(values.get("stop", 0))
257 for screen in tuple(self.screens):
258 screen.send({"type": "capabilities", "windows": bool(self.windows_available)})
259 if self.stopping:
260 self.stop()
261 else:
262 self.request()
263 elif name == "capture_stats":
264 if data or not 0 < values["elapsed_ms"] <= 60000 or values["frames"] > 10000 or any(values[key] > 60000000 for key in ("capture_us", "scan_us", "write_us")):
265 raise ValueError("Invalid guest capture statistics")
266 for screen in tuple(self.screens):
267 screen.send({"type": "capture-stats", **values})
268 elif name == "windows":
269 count = values["count"]
270 fields = GUEST_PROTOCOL["window_fields"]
271 if count > limits["windows"] or len(data) != count * len(fields) * 4:
272 raise ValueError("Invalid guest window list")
273 previous = self.layers
274 if count:
275 records = []
276 for record in struct.iter_unpack("!" + "I" * len(fields), data):
277 window = dict(zip(fields, record))
278 for key in ("x", "y", "caption_left", "caption_top", "caption_right", "caption_bottom"):
279 window[key] = (window[key] ^ 0x80000000) - 0x80000000
280 records.append(window)
281 try:
282 self.layers = WindowDesktop(values["width"], values["height"], records, previous)
283 except OverflowError:
284 self.layers = None
285 else:
286 self.layers = None
287 for screen in tuple(self.screens):
288 if screen.experiment:
289 if self.layers and not self.layers.repacked and self.layers.complete and screen.source is self.layers:
290 screen.send({"type": "window-positions", "windows": self.layers.windows})
291 else:
292 screen.resize()
293 screen.present()
294 screen.send({"type": "window-capture", "available": bool(self.layers)})
295 elif name == "surface":
296 if self.layers:
297 complete = self.layers.complete
298 self.layers.update(values, data)
299 for screen in tuple(self.screens):
300 if screen.experiment and self.layers.complete:
301 if not complete:
302 screen.resize()
303 screen.present()
143 else:304 else:
144 raise ValueError("Unexpected guest display request")305 raise ValueError("Unexpected guest display request")
145306
146 async def run(self):307 async def run(self):
147 while True:308 while True:
309 if self.stopping:
310 GUESTS.pop(self.name, None)
311 return
148 writer = None312 writer = None
149 try:313 try:
150 reader, writer = await asyncio.wait_for(asyncio.open_unix_connection(314 reader, writer = await asyncio.wait_for(asyncio.open_unix_connection(
...@@ -158,16 +322,23 @@ class GuestDesktop:...@@ -158,16 +322,23 @@ class GuestDesktop:
158 raise ValueError("Invalid guest display broker response")322 raise ValueError("Invalid guest display broker response")
159 response = json.loads(await reader.readexactly(size))323 response = json.loads(await reader.readexactly(size))
160 if "error" not in response:324 if "error" not in response:
325 self.writer = writer
161 writer.write(struct.pack("!4sII", GUEST_PROTOCOL["magic"].encode(), GUEST_PROTOCOL["packets"]["refresh"]["id"], 0))326 writer.write(struct.pack("!4sII", GUEST_PROTOCOL["magic"].encode(), GUEST_PROTOCOL["packets"]["refresh"]["id"], 0))
162 await self.read(reader)327 await self.read(reader)
163 except (OSError, ValueError, EOFError, asyncio.IncompleteReadError, TimeoutError) as error:328 except (OSError, ValueError, EOFError, asyncio.IncompleteReadError, TimeoutError) as error:
164 logging.info("Guest display %s: %s", self.name, error)329 logging.info("Guest display %s: %s", self.name, error)
165 finally:330 finally:
331 self.writer = None
332 self.windows_available = None
333 self.stop_supported = None
166 self.suspend()334 self.suspend()
167 if writer:335 if writer:
168 writer.close()336 writer.close()
169 with contextlib.suppress(OSError, TimeoutError):337 with contextlib.suppress(OSError, TimeoutError):
170 await asyncio.wait_for(writer.wait_closed(), 2)338 await asyncio.wait_for(writer.wait_closed(), 2)
339 if self.stopping or not self.screens:
340 GUESTS.pop(self.name, None)
341 return
171 await asyncio.sleep(2)342 await asyncio.sleep(2)
172343
173344
...@@ -176,6 +347,7 @@ class Desktop:...@@ -176,6 +347,7 @@ class Desktop:
176 self.reader, self.writer, self.notify = reader, writer, notify347 self.reader, self.writer, self.notify = reader, writer, notify
177 self.keys = set()348 self.keys = set()
178 self.pointer = (0, 0)349 self.pointer = (0, 0)
350 self.buttons = 0
179 self.pixels = bytearray()351 self.pixels = bytearray()
180 self.cursor = None352 self.cursor = None
181353
...@@ -226,11 +398,13 @@ class Desktop:...@@ -226,11 +398,13 @@ class Desktop:
226 if any(type(value) is not int for value in (x, y, buttons)) or not 0 <= buttons <= 255:398 if any(type(value) is not int for value in (x, y, buttons)) or not 0 <= buttons <= 255:
227 raise ValueError("Unsupported pointer event.")399 raise ValueError("Unsupported pointer event.")
228 self.pointer = (max(0, min(x, self.width - 1)), max(0, min(y, self.height - 1)))400 self.pointer = (max(0, min(x, self.width - 1)), max(0, min(y, self.height - 1)))
401 self.buttons = buttons
229 self.writer.write(struct.pack("!BBHH", 5, buttons, *self.pointer))402 self.writer.write(struct.pack("!BBHH", 5, buttons, *self.pointer))
230 elif kind == "release":403 elif kind == "release":
231 for key in self.keys:404 for key in self.keys:
232 self.writer.write(struct.pack("!BB2xI", 4, 0, key))405 self.writer.write(struct.pack("!BB2xI", 4, 0, key))
233 self.keys.clear()406 self.keys.clear()
407 self.buttons = 0
234 self.writer.write(struct.pack("!BBHH", 5, 0, *self.pointer))408 self.writer.write(struct.pack("!BBHH", 5, 0, *self.pointer))
235 elif kind == "clipboard":409 elif kind == "clipboard":
236 text = message.get("text")410 text = message.get("text")
...@@ -310,17 +484,29 @@ class Screen:...@@ -310,17 +484,29 @@ class Screen:
310 self.writer, self.desktop = writer, desktop484 self.writer, self.desktop = writer, desktop
311 self.guest = None485 self.guest = None
312 self.channel = None486 self.channel = None
487 self.pointer_channel = None
488 self.pointer_sequence = -1
489 self.experiment = False
490 self.frame_timer = None
491 self.last_frame = 0
492 self.layouts = {}
493 self.encoding_layouts = deque(maxlen=64)
494 self.last_rtp_pts = None
495 self.sent_layout = None
496 self.bwe = None
313 hardware = Gst.ElementFactory.make("nvh264enc")497 hardware = Gst.ElementFactory.make("nvh264enc")
314 encoder = "nvh264enc name=encoder preset=p4 tune=ultra-low-latency zerolatency=true bframes=0 bitrate=6000 gop-size=30" if hardware else (498 encoder = "nvh264enc name=encoder preset=p4 tune=ultra-low-latency zerolatency=true bframes=0 rc-mode=cbr bitrate=4000 gop-size=60" if hardware else (
315 "x264enc name=encoder tune=zerolatency speed-preset=ultrafast bitrate=6000 key-int-max=30 bframes=0")499 "x264enc name=encoder tune=zerolatency speed-preset=ultrafast bitrate=4000 key-int-max=60 bframes=0")
316 width, height = self.desktop.width, self.desktop.height500 width, height = self.desktop.width, self.desktop.height
501 ssrc = struct.unpack("!I", os.urandom(4))[0] | 1
502 # TWCC and RTX headers are added after payload packetization.
317 self.pipeline = Gst.parse_launch(503 self.pipeline = Gst.parse_launch(
318 f"appsrc name=frames is-live=true format=time do-timestamp=true max-buffers=1 leaky-type=downstream "504 f"appsrc name=frames is-live=true format=time max-buffers=1 leaky-type=downstream "
319 f"caps=video/x-raw,format=BGRx,width={width},height={height},framerate=30/1 "505 f"caps=video/x-raw,format=BGRx,width={width},height={height},framerate=60/1 "
320 f"! videoconvert ! videoscale ! capsfilter name=size caps=video/x-raw,format=NV12,width={width + width % 2},height={height + height % 2} "506 f"! videoconvert ! videoscale ! capsfilter name=size caps=video/x-raw,format=NV12,width={width + width % 2},height={height + height % 2} "
321 f"! {encoder} ! video/x-h264,profile=constrained-baseline ! h264parse "507 f"! {encoder} ! video/x-h264,profile=constrained-baseline ! h264parse "
322 "! rtph264pay name=pay pt=96 mtu=1200 config-interval=-1 aggregate-mode=zero-latency "508 f"! rtph264pay name=pay pt=96 ssrc={ssrc} mtu=1190 config-interval=-1 aggregate-mode=zero-latency "
323 "! capsfilter name=codec caps=application/x-rtp,media=video,encoding-name=H264,clock-rate=90000 "509 f"! capsfilter name=codec caps=application/x-rtp,media=video,encoding-name=H264,clock-rate=90000,ssrc=(uint){ssrc} "
324 "! webrtcbin name=rtc bundle-policy=max-bundle")510 "! webrtcbin name=rtc bundle-policy=max-bundle")
325 self.frames = self.pipeline.get_by_name("frames")511 self.frames = self.pipeline.get_by_name("frames")
326 self.rtc = self.pipeline.get_by_name("rtc")512 self.rtc = self.pipeline.get_by_name("rtc")
...@@ -335,7 +521,11 @@ class Screen:...@@ -335,7 +521,11 @@ class Screen:
335 self.send, {"type": "candidate", "sdpMLineIndex": line, "candidate": candidate}))),521 self.send, {"type": "candidate", "sdpMLineIndex": line, "candidate": candidate}))),
336 (self.rtc, self.rtc.connect("on-data-channel", self.data_channel)),522 (self.rtc, self.rtc.connect("on-data-channel", self.data_channel)),
337 (self.rtc, self.rtc.connect("notify::connection-state", self.connection_state)),523 (self.rtc, self.rtc.connect("notify::connection-state", self.connection_state)),
524 (self.rtc, self.rtc.connect("request-aux-sender", lambda *_: self.bwe)),
338 ]525 ]
526 self.encoder = self.pipeline.get_by_name("encoder")
527 self.encoder_probe = self.encoder.get_static_pad("sink").add_probe(Gst.PadProbeType.BUFFER, self.encoding)
528 self.rtp_probe = self.pipeline.get_by_name("pay").get_static_pad("src").add_probe(Gst.PadProbeType.BUFFER, self.rtp_layout)
339 self.probe = self.frames.get_static_pad("src").add_probe(Gst.PadProbeType.EVENT_UPSTREAM, self.refresh)529 self.probe = self.frames.get_static_pad("src").add_probe(Gst.PadProbeType.EVENT_UPSTREAM, self.refresh)
340 self.bus = self.pipeline.get_bus()530 self.bus = self.pipeline.get_bus()
341 self.bus.add_signal_watch()531 self.bus.add_signal_watch()
...@@ -344,18 +534,25 @@ class Screen:...@@ -344,18 +534,25 @@ class Screen:
344 self.send({"type": "display", "width": width, "height": height, "encoder": "nvenc" if hardware else "x264"})534 self.send({"type": "display", "width": width, "height": height, "encoder": "nvenc" if hardware else "x264"})
345535
346 def close(self):536 def close(self):
537 self.writer.close()
538 if self.frame_timer:
539 self.frame_timer.cancel()
347 self.desktop.notify = None540 self.desktop.notify = None
348 if self.guest:541 if self.guest:
349 self.guest.screens.discard(self)542 self.guest.screens.discard(self)
350 if not self.guest.screens:543 if self.guest.screens:
351 GUESTS.pop(self.guest.name, None)544 self.guest.request()
352 self.guest.task.cancel()545 else:
546 self.guest.stop()
353 for element, handler in self.handlers:547 for element, handler in self.handlers:
354 element.disconnect(handler)548 element.disconnect(handler)
355 self.handlers.clear()549 self.handlers.clear()
356 self.frames.get_static_pad("src").remove_probe(self.probe)550 self.frames.get_static_pad("src").remove_probe(self.probe)
551 self.encoder.get_static_pad("sink").remove_probe(self.encoder_probe)
552 self.pipeline.get_by_name("pay").get_static_pad("src").remove_probe(self.rtp_probe)
357 self.bus.remove_signal_watch()553 self.bus.remove_signal_watch()
358 self.pipeline.set_state(Gst.State.NULL)554 self.pipeline.set_state(Gst.State.NULL)
555 self.pipeline = self.encoder = None
359556
360 def send(self, message):557 def send(self, message):
361 if not self.writer.is_closing():558 if not self.writer.is_closing():
...@@ -373,13 +570,16 @@ class Screen:...@@ -373,13 +570,16 @@ class Screen:
373 self.loop.call_soon_threadsafe(self.fail, "The screen stream stopped. Reconnect to try again.")570 self.loop.call_soon_threadsafe(self.fail, "The screen stream stopped. Reconnect to try again.")
374571
375 def resize(self):572 def resize(self):
573 if self.writer.is_closing():
574 return
376 width, height = self.source.width, self.source.height575 width, height = self.source.width, self.source.height
377 self.frames.set_property("caps", Gst.Caps.from_string(576 self.frames.set_property("caps", Gst.Caps.from_string(
378 f"video/x-raw,format=BGRx,width={width},height={height},framerate=30/1"))577 f"video/x-raw,format=BGRx,width={width},height={height},framerate=60/1"))
379 self.pipeline.get_by_name("size").set_property("caps", Gst.Caps.from_string(578 self.pipeline.get_by_name("size").set_property("caps", Gst.Caps.from_string(
380 f"video/x-raw,format=NV12,width={width + width % 2},height={height + height % 2}"))579 f"video/x-raw,format=NV12,width={width + width % 2},height={height + height % 2}"))
381 self.send({"type": "display", "width": width, "height": height,580 self.send({"type": "display", "width": self.guest.width if self.guest and self.guest.pixels else width,
382 "source": "guest" if self.source is self.guest else "qemu",581 "height": self.guest.height if self.guest and self.guest.pixels else height,
582 "source": "guest" if self.guest and self.guest.pixels else "qemu",
383 "encoder": "nvenc" if self.pipeline.get_by_name("encoder").get_factory().get_name() == "nvh264enc" else "x264"})583 "encoder": "nvenc" if self.pipeline.get_by_name("encoder").get_factory().get_name() == "nvh264enc" else "x264"})
384584
385 def refresh(self, _, info):585 def refresh(self, _, info):
...@@ -393,17 +593,24 @@ class Screen:...@@ -393,17 +593,24 @@ class Screen:
393 if state == GstWebRTC.WebRTCPeerConnectionState.CONNECTED:593 if state == GstWebRTC.WebRTCPeerConnectionState.CONNECTED:
394 self.loop.call_soon_threadsafe(self.present)594 self.loop.call_soon_threadsafe(self.present)
395 self.loop.call_soon_threadsafe(self.desktop.request, False)595 self.loop.call_soon_threadsafe(self.desktop.request, False)
596 if self.guest:
597 self.loop.call_soon_threadsafe(self.guest.sync, self)
598 self.loop.call_soon_threadsafe(self.guest.request, True)
396 elif state in (GstWebRTC.WebRTCPeerConnectionState.FAILED, GstWebRTC.WebRTCPeerConnectionState.CLOSED):599 elif state in (GstWebRTC.WebRTCPeerConnectionState.FAILED, GstWebRTC.WebRTCPeerConnectionState.CLOSED):
397 self.loop.call_soon_threadsafe(self.fail, "The screen disconnected. Reconnect to try again.")600 self.loop.call_soon_threadsafe(self.fail, "The screen disconnected. Reconnect to try again.")
398601
399 def data_channel(self, _, channel):602 def data_channel(self, _, channel):
400 if channel.get_property("label") != "input" or self.channel:603 label = channel.get_property("label")
604 if label not in ("input", "pointer") or (self.channel if label == "input" else self.pointer_channel):
401 channel.emit("close")605 channel.emit("close")
402 return606 return
403 self.channel = channel607 if label == "input":
404 self.handlers.append((channel, channel.connect("on-message-string", lambda _, message: self.loop.call_soon_threadsafe(self.input, message))))608 self.channel = channel
609 else:
610 self.pointer_channel = channel
611 self.handlers.append((channel, channel.connect("on-message-string", lambda _, message: self.loop.call_soon_threadsafe(self.input, message, label))))
405612
406 def input(self, raw):613 def input(self, raw, channel="input"):
407 if self.writer.is_closing():614 if self.writer.is_closing():
408 return615 return
409 try:616 try:
...@@ -412,6 +619,25 @@ class Screen:...@@ -412,6 +619,25 @@ class Screen:
412 message = json.loads(raw)619 message = json.loads(raw)
413 if not isinstance(message, dict):620 if not isinstance(message, dict):
414 raise ValueError("Unsupported input event.")621 raise ValueError("Unsupported input event.")
622 if channel == "pointer" and message.get("type") != "pointer":
623 raise ValueError("Unsupported pointer event.")
624 if message.get("type") == "experiment":
625 if type(message.get("enabled")) is not bool or not self.guest:
626 raise ValueError("Unsupported window experiment.")
627 self.experiment = message["enabled"] and bool(self.guest.windows_available)
628 self.guest.request()
629 self.resize()
630 self.present()
631 return
632 if message.get("type") in ("pointer", "release") and "sequence" in message:
633 sequence = message["sequence"]
634 if type(sequence) is not int or not 0 <= sequence <= 0x1fffffffffffff:
635 raise ValueError("Unsupported pointer event.")
636 if channel == "pointer" and sequence <= self.pointer_sequence:
637 return
638 self.pointer_sequence = max(sequence, self.pointer_sequence)
639 if channel == "pointer":
640 message["buttons"] = self.desktop.buttons
415 self.desktop.input(message)641 self.desktop.input(message)
416 if self.desktop.writer.transport.get_write_buffer_size() > 262144:642 if self.desktop.writer.transport.get_write_buffer_size() > 262144:
417 self.fail("The VM stopped accepting input. Reconnect to try again.")643 self.fail("The VM stopped accepting input. Reconnect to try again.")
...@@ -420,12 +646,47 @@ class Screen:...@@ -420,12 +646,47 @@ class Screen:
420646
421 @property647 @property
422 def source(self):648 def source(self):
649 if self.experiment and self.guest and self.guest.pixels and self.guest.layers and self.guest.layers.complete:
650 return self.guest.layers
423 return self.guest if self.guest and self.guest.pixels else self.desktop651 return self.guest if self.guest and self.guest.pixels else self.desktop
424652
425 def present(self):653 def present(self):
426 if self.rtc.get_property("connection-state") == GstWebRTC.WebRTCPeerConnectionState.CONNECTED:654 if self.writer.is_closing():
427 buffer = Gst.Buffer.new_wrapped(bytes(self.source.pixels))655 return
428 self.frames.emit("push-buffer", buffer)656 if self.frame_timer is None and self.rtc.get_property("connection-state") == GstWebRTC.WebRTCPeerConnectionState.CONNECTED:
657 self.frame_timer = self.loop.call_later(max(0, self.last_frame + 1 / 60 - time.monotonic()), self.push_frame)
658
659 def push_frame(self):
660 self.frame_timer = None
661 if self.writer.is_closing():
662 return
663 self.last_frame = time.monotonic()
664 source = self.source
665 buffer = Gst.Buffer.new_wrapped(bytes(source.pixels))
666 buffer.pts = self.pipeline.get_clock().get_time() - self.pipeline.get_base_time()
667 buffer.duration = Gst.SECOND // 60
668 layout = {"type": "windows", "windows": source.windows if isinstance(source, WindowDesktop) else [], "width": source.width, "height": source.height}
669 self.layouts[buffer.pts] = layout
670 while len(self.layouts) > 64:
671 del self.layouts[next(iter(self.layouts))]
672 self.frames.emit("push-buffer", buffer)
673
674 def encoding(self, _, info):
675 layout = self.layouts.pop(info.get_buffer().pts, None)
676 if layout is not None:
677 self.encoding_layouts.append(layout)
678 return Gst.PadProbeReturn.OK
679
680 def rtp_layout(self, _, info):
681 buffer = info.get_buffer()
682 if buffer.pts != self.last_rtp_pts and self.encoding_layouts:
683 self.last_rtp_pts = buffer.pts
684 layout = self.encoding_layouts.popleft()
685 if layout != self.sent_layout:
686 self.sent_layout = layout
687 timestamp = struct.unpack("!I", buffer.extract_dup(4, 4))[0]
688 self.loop.call_soon_threadsafe(self.send, {**layout, "rtpTimestamp": timestamp})
689 return Gst.PadProbeReturn.OK
429690
430 def signal(self, message):691 def signal(self, message):
431 if message.get("type") == "offer" and isinstance(message.get("sdp"), str):692 if message.get("type") == "offer" and isinstance(message.get("sdp"), str):
...@@ -455,13 +716,40 @@ class Screen:...@@ -455,13 +716,40 @@ class Screen:
455 if not caps:716 if not caps:
456 raise ValueError("This browser doesn't support the VM's H.264 stream. Open it in a current browser.")717 raise ValueError("This browser doesn't support the VM's H.264 stream. Open it in a current browser.")
457 payload = caps.get_structure(0).get_value("payload")718 payload = caps.get_structure(0).get_value("payload")
719 extension = None
720 for index in range(media.attributes_len()):
721 attribute = media.get_attribute(index)
722 if attribute.key == "extmap" and attribute.value.split()[1:2] == [TWCC]:
723 extension = int(attribute.value.split()[0].split("/")[0])
724 if not 1 <= extension <= 255:
725 raise ValueError("Unsupported screen connection extension.")
726 if extension and Gst.ElementFactory.find("rtpgccbwe"):
727 self.bwe = Gst.ElementFactory.make("rtpgccbwe")
728 self.bwe.set_property("min-bitrate", 500000)
729 self.bwe.set_property("max-bitrate", 12000000)
730 self.bwe.set_property("estimated-bitrate", 4000000)
731 self.handlers.append((self.bwe, self.bwe.connect("notify::estimated-bitrate", lambda bwe, _: self.encoder.set_property("bitrate", bwe.get_property("estimated-bitrate") // 1000))))
732 header = GstRtp.RTPHeaderExtension.create_from_uri(TWCC)
733 header.set_id(extension)
734 self.pipeline.get_by_name("pay").emit("add-extension", header)
458 caps = Gst.Caps.from_string(735 caps = Gst.Caps.from_string(
459 f"application/x-rtp,media=video,encoding-name=H264,clock-rate=90000,payload={payload},"736 f"application/x-rtp,media=video,encoding-name=H264,clock-rate=90000,payload={payload},"
737 f"ssrc=(uint){self.pipeline.get_by_name('pay').get_property('ssrc')},"
460 "packetization-mode=(string)1,rtcp-fb-nack=(boolean)true,rtcp-fb-nack-pli=(boolean)true")738 "packetization-mode=(string)1,rtcp-fb-nack=(boolean)true,rtcp-fb-nack-pli=(boolean)true")
739 if self.bwe:
740 caps.set_value("rtcp-fb-transport-cc", True)
741 caps.set_value(f"extmap-{extension}", TWCC)
461 self.pipeline.get_by_name("pay").set_property("pt", payload)742 self.pipeline.get_by_name("pay").set_property("pt", payload)
462 self.pipeline.get_by_name("codec").set_property("caps", caps)743 self.pipeline.get_by_name("codec").set_property("caps", caps)
463 transceiver = self.rtc.get_static_pad("sink_0").get_property("transceiver")744 transceiver = self.rtc.get_static_pad("sink_0").get_property("transceiver")
464 transceiver.set_property("codec-preferences", caps)745 preferences = caps.copy()
746 for index in range(media.formats_len()):
747 offered = media.get_caps_from_media(int(media.get_format(index)))
748 if offered:
749 codec = offered.get_structure(0)
750 if codec.get_string("encoding-name") == "RTX" and codec.get_string("apt") == str(payload):
751 preferences.append(offered)
752 transceiver.set_property("codec-preferences", preferences)
465 transceiver.set_property("direction", GstWebRTC.WebRTCRTPTransceiverDirection.SENDONLY)753 transceiver.set_property("direction", GstWebRTC.WebRTCRTPTransceiverDirection.SENDONLY)
466 transceiver.set_property("do-nack", True)754 transceiver.set_property("do-nack", True)
467 description = GstWebRTC.WebRTCSessionDescription.new(GstWebRTC.WebRTCSDPType.OFFER, sdp)755 description = GstWebRTC.WebRTCSessionDescription.new(GstWebRTC.WebRTCSDPType.OFFER, sdp)
...@@ -476,6 +764,8 @@ class Screen:...@@ -476,6 +764,8 @@ class Screen:
476 raise ValueError("Unsupported screen connection message.")764 raise ValueError("Unsupported screen connection message.")
477765
478 def remote_description(self, promise, *_):766 def remote_description(self, promise, *_):
767 if self.writer.is_closing():
768 return
479 reply = promise.get_reply()769 reply = promise.get_reply()
480 if reply and reply.has_field("error"):770 if reply and reply.has_field("error"):
481 logging.error("VM screen offer: %s", reply.to_string())771 logging.error("VM screen offer: %s", reply.to_string())
...@@ -484,6 +774,8 @@ class Screen:...@@ -484,6 +774,8 @@ class Screen:
484 self.rtc.emit("create-answer", None, Gst.Promise.new_with_change_func(self.answer, None, None))774 self.rtc.emit("create-answer", None, Gst.Promise.new_with_change_func(self.answer, None, None))
485775
486 def answer(self, promise, *_):776 def answer(self, promise, *_):
777 if self.writer.is_closing():
778 return
487 reply = promise.get_reply()779 reply = promise.get_reply()
488 if not reply or not reply.has_field("answer"):780 if not reply or not reply.has_field("answer"):
489 logging.error("VM screen answer: %s", reply.to_string() if reply else "missing reply")781 logging.error("VM screen answer: %s", reply.to_string() if reply else "missing reply")
...@@ -523,6 +815,9 @@ async def serve(reader, writer):...@@ -523,6 +815,9 @@ async def serve(reader, writer):
523 raise ValueError(result["error"])815 raise ValueError(result["error"])
524 await asyncio.wait_for(desktop.start(), 10)816 await asyncio.wait_for(desktop.start(), 10)
525 screen = Screen(writer, desktop)817 screen = Screen(writer, desktop)
818 guest = GUESTS.get(name)
819 if guest and guest.stopping:
820 await asyncio.shield(guest.task)
526 if name not in GUESTS:821 if name not in GUESTS:
527 GUESTS[name] = GuestDesktop(name)822 GUESTS[name] = GuestDesktop(name)
528 GUESTS[name].task = asyncio.create_task(GUESTS[name].run())823 GUESTS[name].task = asyncio.create_task(GUESTS[name].run())
...@@ -587,13 +882,36 @@ async def serve(reader, writer):...@@ -587,13 +882,36 @@ async def serve(reader, writer):
587882
588async def main():883async def main():
589 Gst.init(None)884 Gst.init(None)
590 threading.Thread(target=GLib.MainLoop().run, daemon=True).start()885 glib = GLib.MainLoop()
886 threading.Thread(target=glib.run, daemon=True).start()
887 stopped = asyncio.Event()
888 loop = asyncio.get_running_loop()
889 for signum in (signal.SIGTERM, signal.SIGINT):
890 loop.add_signal_handler(signum, stopped.set)
891 clients = set()
892 def connected(reader, writer):
893 if stopped.is_set():
894 writer.close()
895 return
896 task = asyncio.create_task(serve(reader, writer))
897 clients.add(task)
898 task.add_done_callback(clients.discard)
591 path = Path(os.environ.get("STUDIO_VM_SCREEN_SOCKET", "/run/studio-vm-screen/screen.sock"))899 path = Path(os.environ.get("STUDIO_VM_SCREEN_SOCKET", "/run/studio-vm-screen/screen.sock"))
592 path.unlink(missing_ok=True)900 path.unlink(missing_ok=True)
593 server = await asyncio.start_unix_server(serve, path=str(path), limit=131072)901 server = await asyncio.start_unix_server(connected, path=str(path), limit=131072)
594 path.chmod(0o600)902 path.chmod(0o600)
595 async with server:903 await stopped.wait()
596 await server.serve_forever()904 server.close()
905 guests = tuple(GUESTS.values())
906 for guest in guests:
907 guest.stop()
908 tasks = tuple(clients)
909 for task in tasks:
910 task.cancel()
911 await asyncio.gather(*tasks, return_exceptions=True)
912 await server.wait_closed()
913 await asyncio.gather(*(guest.task for guest in guests), return_exceptions=True)
914 glib.quit()
597915
598916
599if __name__ == "__main__":917if __name__ == "__main__":
tools/vms.py+1-1
...@@ -830,7 +830,7 @@ def main():...@@ -830,7 +830,7 @@ def main():
830 if action == "guest":830 if action == "guest":
831 root = ET.fromstring(virsh("dumpxml", payload["name"]))831 root = ET.fromstring(virsh("dumpxml", payload["name"]))
832 source = root.find(f"./devices/channel[@type='unix']/target[@name='{GUEST_CHANNEL}']/../source")832 source = root.find(f"./devices/channel[@type='unix']/target[@name='{GUEST_CHANNEL}']/../source")
833 expected = f"/run/libvirt/qemu/channel/{root.get('id')}-{payload['name']}/{GUEST_CHANNEL}"833 expected = f"/run/libvirt/qemu/channel/{root.get('id')}-{payload['name'][:20]}/{GUEST_CHANNEL}"
834 if virsh("domstate", payload["name"]).strip() != "running" or source is None or source.get("mode") != "bind" or source.get("path") != expected:834 if virsh("domstate", payload["name"]).strip() != "running" or source is None or source.get("mode") != "bind" or source.get("path") != expected:
835 raise ValueError("This VM has no guest display channel.")835 raise ValueError("This VM has no guest display channel.")
836 return {"path": expected}836 return {"path": expected}