mirror of
https://github.com/saphid/frame-control.git
synced 2026-10-06 07:00:37 +02:00
Add native PC capture adapters behind the shared panel viewer
This commit is contained in:
1 parent
559caa2dd0
commit
c5fc552c60
19 files changed
+1865
-36
No files matched your search
+17
-7
@@ -105,9 +105,9 @@ class MacViewError(Exception):
|
||||
pass
|
||||
|
||||
|
||||
def panel_id(src):
|
||||
def panel_id(src, host="mac"):
|
||||
"""A stable panel id per source, in the range panel-on-frame.sh uses."""
|
||||
return 2_001_000_000 + zlib.crc32(f"mac:{src}".encode()) % 1_000_000
|
||||
return 2_001_000_000 + zlib.crc32(f"{host}:{src}".encode()) % 1_000_000
|
||||
|
||||
|
||||
def fit(w, h, box=PANEL_BOX):
|
||||
@@ -116,6 +116,8 @@ def fit(w, h, box=PANEL_BOX):
|
||||
|
||||
|
||||
class MacView:
|
||||
host = "mac"
|
||||
|
||||
def __init__(self, tunnel_ssh, run, frame, track=None):
|
||||
self.tunnel_ssh = list(tunnel_ssh)
|
||||
self.run = run
|
||||
@@ -166,10 +168,17 @@ class MacView:
|
||||
def _stale(self):
|
||||
try:
|
||||
built = AGENT.stat().st_mtime
|
||||
return any(p.stat().st_mtime > built for p in (SOURCES / "Sources").glob("*.swift"))
|
||||
return any(p.stat().st_mtime > built for p in [*(SOURCES / "Sources").glob("*.swift"),
|
||||
ROOT / "desktop" / "controller.c", ROOT / "desktop" / "controller.h"])
|
||||
except OSError:
|
||||
return False
|
||||
|
||||
def agent_command(self):
|
||||
return [str(AGENT)]
|
||||
|
||||
def agent_environment(self):
|
||||
return {**os.environ, "FRAME_MAC_VIEW_TOKEN": self.token}
|
||||
|
||||
def ensure_agent(self):
|
||||
with self.lock:
|
||||
if self.agent and self.agent.poll() is None:
|
||||
@@ -178,11 +187,11 @@ class MacView:
|
||||
if reason:
|
||||
raise MacViewError(reason)
|
||||
self.build()
|
||||
env = {**os.environ, "FRAME_MAC_VIEW_TOKEN": self.token}
|
||||
env = self.agent_environment()
|
||||
# The same port as before when restarting, so a running tunnel still
|
||||
# fits; otherwise (or if it's gone) whatever the system gives.
|
||||
for port in dict.fromkeys([self.port or 0, 0]):
|
||||
self.agent = subprocess.Popen([str(AGENT), "serve", "--port", str(port), "--page", str(PAGE),
|
||||
self.agent = subprocess.Popen([*self.agent_command(), "serve", "--port", str(port), "--page", str(PAGE),
|
||||
"--exit-on-eof"], env=env, stdin=subprocess.PIPE,
|
||||
stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True)
|
||||
line = self.agent.stdout.readline()
|
||||
@@ -362,7 +371,7 @@ class MacView:
|
||||
# the agent no longer counts them, so only move when none are out there.
|
||||
others = self.shown - {src}
|
||||
self.ensure_tunnel(allow_new_port=not others)
|
||||
appid = panel_id(src)
|
||||
appid = panel_id(src, self.host)
|
||||
# Unique per launch, so a new window is never confused with an old one.
|
||||
tag = "fc" + secrets.token_hex(4)
|
||||
# A single-use ticket for this source, not Frame Control's key: the URL
|
||||
@@ -415,7 +424,8 @@ class MacView:
|
||||
status = self.call("/status")
|
||||
windows = self.call("/windows").get("windows", []) if status.get("screen") else []
|
||||
displays = self.call("/displays").get("displays", [])
|
||||
return {"available": True, "screen": status.get("screen", False),
|
||||
return {"available": True, "host": self.host, "selecting": status.get("selecting", False),
|
||||
"selectionError": status.get("selectionError", ""), "screen": status.get("screen", False),
|
||||
"accessibility": status.get("accessibility", False), "streams": status.get("streams", []),
|
||||
"windows": windows, "displays": displays, "tunnel": self.tunnel_up(),
|
||||
"route": self.route}
|
||||
|
||||
@@ -0,0 +1,507 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Frame Control's Windows/Linux host. Same ticket, H.264/JPEG, timing and
|
||||
input protocol as frame-mac-view; serves the same ui/mac-view.html.
|
||||
|
||||
Only loopback is bound. Frame Control owns the master token, the Frame gets
|
||||
one-source tickets and reconnect keys. Native libraries are bundled.
|
||||
"""
|
||||
import argparse
|
||||
import base64
|
||||
import ctypes as C
|
||||
import hashlib
|
||||
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
|
||||
import json
|
||||
import math
|
||||
import os
|
||||
from pathlib import Path
|
||||
import secrets
|
||||
import socket
|
||||
import struct
|
||||
import sys
|
||||
import threading
|
||||
import time
|
||||
from urllib.parse import parse_qs, urlsplit
|
||||
|
||||
from frame_pc_capture import Native, Controller, Windows, WindowsInput, PortalInput, Encoded, GATE, dimensions, pipeline
|
||||
from frame_stream_stats import Stats
|
||||
|
||||
|
||||
class Grants:
|
||||
"""Caller holds the agent lock, including replacement and Stop."""
|
||||
def __init__(self, token):
|
||||
self.token, self.tickets, self.keys = token, {}, {}
|
||||
|
||||
def master(self, key):
|
||||
return isinstance(key, str) and secrets.compare_digest(key, self.token)
|
||||
|
||||
def ticket(self, src):
|
||||
self.tickets = {k: v for k, v in self.tickets.items() if v[1] > time.monotonic()}
|
||||
if len(self.tickets) >= 128:
|
||||
raise ValueError('Too many pending viewers')
|
||||
ticket = secrets.token_urlsafe(24)
|
||||
self.tickets[ticket] = (src, time.monotonic()+60, None)
|
||||
return ticket
|
||||
|
||||
def redeem(self, src, q):
|
||||
if self.master(q.get('k')):
|
||||
key = secrets.token_urlsafe(24)
|
||||
self.keys[key] = src
|
||||
return key
|
||||
entry = self.tickets.get(q.get('t'))
|
||||
if entry and entry[0] == src and entry[1] > time.monotonic():
|
||||
key = entry[2] or secrets.token_urlsafe(24)
|
||||
self.tickets[q['t']] = (src, entry[1], key)
|
||||
self.keys[key] = src
|
||||
return key
|
||||
key = q.get('r')
|
||||
return key if key and self.keys.get(key) == src else None
|
||||
|
||||
def ack(self, key):
|
||||
self.tickets = {k: v for k, v in self.tickets.items() if v[2] != key}
|
||||
|
||||
def revoke(self, src=None):
|
||||
self.tickets = {k: v for k, v in self.tickets.items() if src is not None and v[0] != src}
|
||||
self.keys = {k: v for k, v in self.keys.items() if src is not None and v != src}
|
||||
|
||||
|
||||
class WebSocket:
|
||||
def __init__(self, handler):
|
||||
self.sock, self.reader = handler.connection, handler.rfile
|
||||
self.lock = threading.Lock()
|
||||
self.closed = False
|
||||
self.sock.settimeout(5)
|
||||
self.sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1)
|
||||
|
||||
def send(self, data, opcode=1):
|
||||
if not isinstance(data, bytes):
|
||||
data = json.dumps(data, separators=(',', ':')).encode()
|
||||
n = len(data)
|
||||
head = bytes([0x80 | opcode, n]) if n < 126 else bytes([0x80 | opcode, 126]) + struct.pack('!H', n) if n <= 65535 else bytes([0x80 | opcode, 127]) + struct.pack('!Q', n)
|
||||
with self.lock:
|
||||
if self.closed:
|
||||
raise ConnectionError('Viewer disconnected')
|
||||
self.sock.sendall(head + data)
|
||||
|
||||
def exact(self, n):
|
||||
data = self.reader.read(n)
|
||||
if len(data) != n:
|
||||
raise ConnectionError('Viewer disconnected')
|
||||
return data
|
||||
|
||||
def receive(self):
|
||||
a, b = self.exact(2)
|
||||
opcode, size = a & 15, b & 127
|
||||
if a & 0x70 or not a & 0x80 or not b & 0x80 or opcode not in (1, 8, 9, 10):
|
||||
raise ValueError('Unsupported WebSocket frame')
|
||||
if size == 126:
|
||||
size = struct.unpack('!H', self.exact(2))[0]
|
||||
elif size == 127:
|
||||
size = struct.unpack('!Q', self.exact(8))[0]
|
||||
if size > 65536 or (opcode >= 8 and size > 125):
|
||||
raise ValueError('WebSocket message too large')
|
||||
mask = self.exact(4)
|
||||
data = bytes(v ^ mask[i % 4] for i, v in enumerate(self.exact(size)))
|
||||
if opcode == 8:
|
||||
raise ConnectionError('Viewer closed')
|
||||
if opcode == 9:
|
||||
self.send(data, 10)
|
||||
if opcode != 1:
|
||||
return {}
|
||||
message = json.loads(data)
|
||||
if not isinstance(message, dict):
|
||||
raise ValueError('Expected an input object')
|
||||
return message
|
||||
|
||||
def close(self):
|
||||
self.closed = True
|
||||
try:
|
||||
self.sock.shutdown(socket.SHUT_RDWR)
|
||||
except OSError:
|
||||
pass
|
||||
|
||||
|
||||
class Session:
|
||||
def __init__(self, agent, ws, key, source, query):
|
||||
self.agent, self.native, self.ws, self.key, self.source = agent, agent.native, ws, key, source
|
||||
self.src = source['src']
|
||||
self.codec = query.get('codec', 'h264')
|
||||
if self.codec not in ('h264', 'jpeg'):
|
||||
raise ValueError('Unsupported codec')
|
||||
self.fps = min(120, max(5, int(query.get('fps', 60))))
|
||||
bpp = float(query.get('bpp', .1))
|
||||
if not math.isfinite(bpp):
|
||||
raise ValueError('Invalid bitrate')
|
||||
self.w, self.h = dimensions(source['w'], source['h'], min(3840, max(320, int(query.get('max', 1920)))))
|
||||
self.bitrate = max(300000, int(self.w*self.h*self.fps*min(.5, max(.02, bpp))))
|
||||
self.controller = Controller(self.native.lib, self.fps, self.bitrate)
|
||||
self.stats = Stats(self.native.now)
|
||||
self.stop_event, self.key_event = threading.Event(), threading.Event()
|
||||
self.lock, self.input_lock = threading.RLock(), threading.RLock()
|
||||
self.pending, self.last_submit = {}, 0
|
||||
self.input = None if self.src == 'test' else WindowsInput(agent.windows, source) if agent.windows else PortalInput(self.native, source)
|
||||
self.writer = None
|
||||
|
||||
def gate(self, stage, pts, capture, arrived):
|
||||
# Exceptions cannot cross a ctypes callback boundary.
|
||||
try:
|
||||
with self.lock:
|
||||
if self.stop_event.is_set():
|
||||
return 0
|
||||
if stage == 1:
|
||||
if pts in self.pending:
|
||||
self.pending[pts]['e0'] = self.native.now()
|
||||
return 1
|
||||
self.stats.captured += 1
|
||||
self.controller.call('capture', arrived)
|
||||
state = self.controller.state()
|
||||
if self.pending or arrived-self.last_submit < 1000000/state['fps'] or not self.controller.call('gate', arrived, 1):
|
||||
self.stats.skipped += 1
|
||||
return 0
|
||||
self.pending[pts] = dict(cap=capture, arr=arrived, e0=arrived, tier=state['tier'], br=state['target'])
|
||||
self.last_submit = arrived
|
||||
return 1
|
||||
except Exception:
|
||||
self.stop_event.set()
|
||||
return 0
|
||||
|
||||
def produce(self):
|
||||
native, capture = self.native.lib, None
|
||||
try:
|
||||
encoder = "x264enc" if self.src == "test" and self.native.has("x264enc") else self.agent.encoder
|
||||
description = pipeline(self.source, sys.platform, encoder, self.w, self.h, self.fps, self.bitrate, self.codec)
|
||||
self.callback = GATE(self.gate)
|
||||
error = C.create_string_buffer(1024)
|
||||
capture = native.fc_capture_open(description.encode(), self.callback, error, len(error))
|
||||
if not capture:
|
||||
raise RuntimeError(error.value.decode(errors='replace'))
|
||||
last_update = last_stats = self.native.now()
|
||||
last_output = last_update
|
||||
while not self.stop_event.is_set():
|
||||
output = Encoded()
|
||||
result = native.fc_capture_pull(capture, C.byref(output))
|
||||
now = self.native.now()
|
||||
if result < 0:
|
||||
raise RuntimeError(native.fc_capture_error(capture).decode(errors='replace'))
|
||||
if result:
|
||||
last_output = now
|
||||
with self.lock:
|
||||
# Exactly one raw frame is in flight and B-frames are disabled.
|
||||
# x264 may offset PTS; associate by that single frame,
|
||||
# keeping the actual pre-encode capture timestamp.
|
||||
raw_pts = next(iter(self.pending), None)
|
||||
record = self.pending.get(raw_pts)
|
||||
if record is None:
|
||||
raise RuntimeError('Encoder changed frame timestamps; timing cannot be matched')
|
||||
data = C.string_at(output.data, output.size)
|
||||
f = self.stats.add(**record, e1=now, snd=now, b=len(data), k=output.key,
|
||||
w=output.width or self.w, h=output.height or self.h)
|
||||
self.controller.call('sent', f['s'], len(data)+17, now)
|
||||
self.ws.send(struct.pack('!BQII', output.key, max(0, f['cap']), f['s'], f['echo'])+data, 2)
|
||||
with self.lock:
|
||||
f['wire'] = self.native.now()
|
||||
self.pending.pop(raw_pts, None)
|
||||
if now-last_output > 10000000:
|
||||
raise RuntimeError('No encoded frames for 10 seconds; check capture permissions and the encoder')
|
||||
if self.key_event.is_set():
|
||||
self.key_event.clear()
|
||||
if self.codec == 'h264':
|
||||
native.fc_capture_key(capture)
|
||||
if now-last_update >= 100000:
|
||||
target = self.controller.update(now)
|
||||
# x264 supports live bitrate changes. Hardware elements
|
||||
# advertise their mutability; the initial bitrate always
|
||||
# applies. Frame gating remains active for every encoder.
|
||||
if target and self.codec == 'h264' and encoder == 'x264enc':
|
||||
native.fc_capture_bitrate(capture, target)
|
||||
last_update = now
|
||||
if now-last_stats >= 1000000:
|
||||
self.ws.send(dict(self.stats.summary(), t='stats', bitrate=self.bitrate,
|
||||
size='%dx%d' % (self.w, self.h), tier=self.controller.state()['tier']))
|
||||
last_stats = now
|
||||
except Exception as e:
|
||||
try:
|
||||
self.ws.send({'t': 'error', 'message': str(e)})
|
||||
except OSError:
|
||||
pass
|
||||
finally:
|
||||
self.stop_event.set()
|
||||
if capture:
|
||||
native.fc_capture_close(capture)
|
||||
self.ws.close()
|
||||
|
||||
def start(self):
|
||||
if self.stop_event.is_set():
|
||||
self.controller.close()
|
||||
return
|
||||
self.ws.send(dict(t='hello', r=self.key))
|
||||
self.ws.send(dict(t='info', src=self.src, title=self.source.get('title', self.source.get('name', 'Test pattern')),
|
||||
app='PC', codec=self.codec, input=self.src == 'test' or self.source.get('devices', 3) == 3,
|
||||
inputMessage='Allow pointer and keyboard control in the host sharing dialog.',
|
||||
aspect=self.w/self.h, warm=0))
|
||||
with self.lock:
|
||||
if self.stop_event.is_set():
|
||||
self.controller.close()
|
||||
return
|
||||
self.writer = threading.Thread(target=self.produce, daemon=True)
|
||||
self.writer.start()
|
||||
try:
|
||||
while not self.stop_event.is_set():
|
||||
m = self.ws.receive()
|
||||
t = m.get('t')
|
||||
if t == 'ping':
|
||||
self.ws.send(dict(t='pong', c=m.get('c', 0), a=self.native.now()))
|
||||
elif t == 'ack':
|
||||
with self.agent.lock:
|
||||
self.agent.grants.ack(self.key)
|
||||
elif t == 'key-frame':
|
||||
self.key_event.set()
|
||||
elif t in ('rx', 'fd', 'clock'):
|
||||
self.stats.report(m)
|
||||
if t == 'rx' and isinstance(m.get('s'), int):
|
||||
self.controller.call('ack', m['s'] & 0xffffffff, self.native.now())
|
||||
elif t in ('m', 'wheel', 'k', 'text', 'release'):
|
||||
with self.input_lock:
|
||||
if self.stop_event.is_set():
|
||||
break
|
||||
if self.input:
|
||||
self.input.handle(m)
|
||||
self.stats.input(m)
|
||||
finally:
|
||||
self.end()
|
||||
self.writer.join(6)
|
||||
if not self.writer.is_alive():
|
||||
self.controller.close()
|
||||
|
||||
def end(self):
|
||||
with self.lock:
|
||||
self.stop_event.set()
|
||||
self.ws.close()
|
||||
with self.input_lock:
|
||||
if self.input:
|
||||
try:
|
||||
self.input.release()
|
||||
except (RuntimeError, OSError):
|
||||
pass
|
||||
|
||||
|
||||
class Agent:
|
||||
def __init__(self, native, token, page):
|
||||
self.native, self.page = native, page
|
||||
self.lock = threading.RLock()
|
||||
self.grants = Grants(token)
|
||||
self.sessions, self.sources = {}, {}
|
||||
self.next_id, self.selecting, self.selection_error = 1, False, ''
|
||||
self.windows = Windows() if sys.platform == 'win32' else None
|
||||
self.encoder = 'mfh264enc' if self.windows else 'vah264enc' if native.has('vah264enc') else 'x264enc'
|
||||
self.shutting_down = False
|
||||
|
||||
def lists(self):
|
||||
if self.windows:
|
||||
windows, displays = self.windows.sources()
|
||||
self.sources = {s['src']: s for s in windows + displays}
|
||||
return windows, displays
|
||||
return [{k: v for k, v in s.items() if k not in ('portal', 'fd', 'node')} for s in self.sources.values()], []
|
||||
|
||||
def source(self, src):
|
||||
if src == 'test':
|
||||
return dict(src='test', title='Test pattern', w=1280, h=720)
|
||||
self.lists()
|
||||
if src not in self.sources:
|
||||
raise ValueError('Choose a window or screen on this computer first')
|
||||
return dict(self.sources[src])
|
||||
|
||||
def select(self):
|
||||
if self.windows or self.selecting or self.shutting_down:
|
||||
return
|
||||
if len(self.sources) >= 8:
|
||||
raise ValueError('Stop a panel before sharing another source')
|
||||
self.selecting, self.selection_error = True, ''
|
||||
def choose():
|
||||
error = C.create_string_buffer(1024)
|
||||
portal = self.native.lib.fc_portal_select(error, len(error))
|
||||
with self.lock:
|
||||
self.selecting = False
|
||||
if not portal:
|
||||
self.selection_error = error.value.decode(errors='replace')
|
||||
elif self.shutting_down:
|
||||
self.native.lib.fc_portal_close(portal)
|
||||
else:
|
||||
fd, node, w, h, devices = [self.native.lib.fc_portal_value(portal, i) for i in range(5)]
|
||||
src = 'window:' + secrets.token_hex(8)
|
||||
self.sources[src] = dict(src=src, title='Shared window or screen', app='Linux portal',
|
||||
w=w, h=h, fd=fd, node=node, portal=portal, devices=devices)
|
||||
threading.Thread(target=choose, daemon=True).start()
|
||||
|
||||
def stop(self, src=None):
|
||||
with self.lock:
|
||||
self.grants.revoke(src)
|
||||
sessions = [s for s in self.sessions.values() if src is None or s.src == src]
|
||||
for session in sessions:
|
||||
try:
|
||||
session.ws.send(dict(t='close'))
|
||||
except OSError:
|
||||
pass
|
||||
session.end()
|
||||
for session in sessions:
|
||||
if session.writer and session.writer is not threading.current_thread():
|
||||
session.writer.join(6)
|
||||
# The capture is stopped before releasing its PipeWire fd/session.
|
||||
with self.lock:
|
||||
if not self.windows:
|
||||
for key, source in list(self.sources.items()):
|
||||
if src is None or key == src:
|
||||
if any(s.src == key and s.writer and s.writer.is_alive() for s in sessions):
|
||||
continue
|
||||
self.native.lib.fc_portal_close(source['portal'])
|
||||
self.sources.pop(key, None)
|
||||
return {'closed': len(sessions)}
|
||||
|
||||
|
||||
class Handler(BaseHTTPRequestHandler):
|
||||
protocol_version = 'HTTP/1.1'
|
||||
|
||||
def log_message(self, *args):
|
||||
pass # URLs contain credentials
|
||||
|
||||
def reply(self, data, status=200, content='application/json'):
|
||||
if not isinstance(data, bytes):
|
||||
data = json.dumps(data).encode()
|
||||
self.send_response(status)
|
||||
self.send_header('Content-Type', content)
|
||||
self.send_header('Content-Length', str(len(data)))
|
||||
self.send_header('Cache-Control', 'no-store')
|
||||
self.send_header('Connection', 'close')
|
||||
self.end_headers()
|
||||
self.wfile.write(data)
|
||||
self.close_connection = True
|
||||
|
||||
def do_GET(self):
|
||||
self.dispatch('GET')
|
||||
|
||||
def do_POST(self):
|
||||
self.dispatch('POST')
|
||||
|
||||
def dispatch(self, method):
|
||||
agent = self.server.agent
|
||||
url = urlsplit(self.path)
|
||||
q = {k: v[-1] for k, v in parse_qs(url.query).items()}
|
||||
try:
|
||||
if method == 'GET' and url.path == '/ping':
|
||||
return self.reply(b'frame-mac-view', content='text/plain')
|
||||
if method == 'GET' and url.path == '/view':
|
||||
return self.reply(agent.page.read_bytes(), content='text/html; charset=utf-8')
|
||||
if method == 'GET' and url.path == '/stream':
|
||||
return self.stream(q)
|
||||
if not agent.grants.master(q.get('k', self.headers.get('X-Token'))):
|
||||
return self.reply({'error': 'forbidden'}, 403)
|
||||
if method == 'POST' and url.path == '/close':
|
||||
return self.reply(agent.stop(q.get('src')))
|
||||
with agent.lock:
|
||||
windows, displays = agent.lists()
|
||||
if method == 'GET' and url.path == '/status':
|
||||
data = dict(version=1, host='windows' if agent.windows else 'linux', screen=True,
|
||||
accessibility=True, selecting=agent.selecting, selectionError=agent.selection_error,
|
||||
encoder=agent.encoder, finished=[], streams=[dict(id=i, src=s.src,
|
||||
title=s.source.get('title', ''), stats=s.stats.summary(), controller=s.controller.state())
|
||||
for i, s in agent.sessions.items() if not s.stop_event.is_set()])
|
||||
elif method == 'GET' and url.path in ('/windows', '/displays'):
|
||||
data = dict(windows=windows, displays=displays, screen=True)
|
||||
elif method == 'POST' and url.path == '/ticket':
|
||||
agent.source(q.get('src'))
|
||||
data = {'ticket': agent.grants.ticket(q['src'])}
|
||||
elif method == 'POST' and url.path == '/permissions':
|
||||
agent.select()
|
||||
data = {'selecting': agent.selecting}
|
||||
elif method == 'GET' and url.path == '/stats':
|
||||
data = dict(now=agent.native.now(), streams=[dict(id=i, src=s.src, controller=s.controller.state(),
|
||||
events=list(s.controller.events), **s.stats.snapshot(max(0, int(q.get('since', 0))), max(0, int(q.get('settle', 1500000)))))
|
||||
for i, s in agent.sessions.items() if q.get('id', str(i)) == str(i) and not s.stop_event.is_set()])
|
||||
elif method == 'POST' and url.path == '/bench':
|
||||
data = {'sent': 0}
|
||||
for s in agent.sessions.values():
|
||||
if s.src == q.get('src'):
|
||||
m = {k: v for k, v in q.items() if k not in ('k', 'src')}
|
||||
for k in ('x', 'y', 'interval'):
|
||||
if k in m:
|
||||
m[k] = float(m[k])
|
||||
s.ws.send(dict(m, t='bench'))
|
||||
data['sent'] += 1
|
||||
else:
|
||||
return self.reply({'error': 'not found'}, 404)
|
||||
self.reply(data)
|
||||
except (ValueError, RuntimeError) as e:
|
||||
self.reply({'error': str(e)}, 400)
|
||||
except (OSError, ConnectionError):
|
||||
self.close_connection = True
|
||||
|
||||
def stream(self, query):
|
||||
agent = self.server.agent
|
||||
with agent.lock:
|
||||
key = agent.grants.redeem(query.get('src'), query)
|
||||
if not key:
|
||||
return self.reply({'error': 'forbidden'}, 403)
|
||||
source = agent.source(query.get('src'))
|
||||
if self.headers.get('Upgrade', '').lower() != 'websocket' or self.headers.get('Sec-WebSocket-Version') != '13':
|
||||
return self.reply({'error': 'expected WebSocket'}, 400)
|
||||
wskey = self.headers.get('Sec-WebSocket-Key', '')
|
||||
if len(base64.b64decode(wskey, validate=True)) != 16:
|
||||
raise ValueError('Bad WebSocket key')
|
||||
if len(agent.sessions) >= 8:
|
||||
raise ValueError('At most eight panels may be open')
|
||||
# Stop can revoke and close only a registered session. Register
|
||||
# under the same lock as redemption, before any capture starts.
|
||||
session = Session(agent, WebSocket(self), key, source, query)
|
||||
for old in list(agent.sessions.values()):
|
||||
if old.key == key:
|
||||
old.end()
|
||||
ident = agent.next_id
|
||||
agent.next_id += 1
|
||||
agent.sessions[ident] = session
|
||||
accept = base64.b64encode(hashlib.sha1((wskey+'258EAFA5-E914-47DA-95CA-C5AB0DC85B11').encode()).digest()).decode()
|
||||
self.send_response(101)
|
||||
self.send_header('Upgrade', 'websocket')
|
||||
self.send_header('Connection', 'Upgrade')
|
||||
self.send_header('Sec-WebSocket-Accept', accept)
|
||||
self.end_headers()
|
||||
try:
|
||||
session.start()
|
||||
except (OSError, ValueError, RuntimeError):
|
||||
session.end()
|
||||
finally:
|
||||
with agent.lock:
|
||||
agent.sessions.pop(ident, None)
|
||||
self.close_connection = True
|
||||
|
||||
|
||||
def main():
|
||||
parser = argparse.ArgumentParser()
|
||||
parser.add_argument('command', choices=['serve'])
|
||||
parser.add_argument('--port', type=int, default=0)
|
||||
parser.add_argument('--page', type=Path, required=True)
|
||||
parser.add_argument('--exit-on-eof', action='store_true')
|
||||
args = parser.parse_args()
|
||||
token = os.environ.get('FRAME_MAC_VIEW_TOKEN')
|
||||
if not token:
|
||||
raise SystemExit('Frame Control must provide a private token')
|
||||
native = Native()
|
||||
server = ThreadingHTTPServer(('127.0.0.1', args.port), Handler)
|
||||
agent = server.agent = Agent(native, token, args.page)
|
||||
if args.exit_on_eof:
|
||||
def eof():
|
||||
sys.stdin.buffer.read()
|
||||
with agent.lock:
|
||||
agent.shutting_down = True
|
||||
agent.stop()
|
||||
server.shutdown()
|
||||
threading.Thread(target=eof, daemon=True).start()
|
||||
print('listening on 127.0.0.1:%d' % server.server_port, flush=True)
|
||||
try:
|
||||
server.serve_forever()
|
||||
finally:
|
||||
agent.shutting_down = True
|
||||
agent.stop()
|
||||
server.server_close()
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
main()
|
||||
@@ -0,0 +1,340 @@
|
||||
"""Native PC capture and input bindings. No desktop-streaming app required.
|
||||
|
||||
The packaged native library contains our adapter and the shared rate controller;
|
||||
GStreamer supplies WGC/Media Foundation and PipeWire/VA-API/x264 libraries.
|
||||
"""
|
||||
import ctypes as C
|
||||
import os
|
||||
from pathlib import Path
|
||||
import re
|
||||
import sys
|
||||
import threading
|
||||
|
||||
ROOT = Path(__file__).resolve().parent.parent
|
||||
NATIVE = Path(os.environ.get('FRAME_PC_NATIVE', ROOT / 'desktop' / 'bundle'))
|
||||
LIBRARY = NATIVE / ('pc-host.dll' if sys.platform == 'win32' else 'pc-host.so')
|
||||
|
||||
|
||||
class Encoded(C.Structure):
|
||||
_fields_ = [('data', C.c_void_p), ('size', C.c_int), ('key', C.c_int),
|
||||
('width', C.c_int), ('height', C.c_int), ('pts', C.c_int64)]
|
||||
|
||||
|
||||
GATE = C.CFUNCTYPE(C.c_int, C.c_int, C.c_int64, C.c_int64, C.c_int64)
|
||||
|
||||
|
||||
class Native:
|
||||
def __init__(self, path=LIBRARY):
|
||||
self.dll_dirs = []
|
||||
if sys.platform == 'win32':
|
||||
for folder in (path.parent, path.parent / 'bin'):
|
||||
self.dll_dirs.append(os.add_dll_directory(str(folder)))
|
||||
self.lib = C.CDLL(str(path))
|
||||
signatures = {
|
||||
'fc_now': (C.c_int64, []), 'fc_gst_init': (None, []),
|
||||
'fc_has_element': (C.c_int, [C.c_char_p]),
|
||||
'fc_new': (C.c_void_p, [C.c_int, C.c_int]), 'fc_free': (None, [C.c_void_p]),
|
||||
'fc_ceiling': (None, [C.c_void_p, C.c_int]),
|
||||
'fc_gate': (C.c_int, [C.c_void_p, C.c_int64, C.c_int]),
|
||||
'fc_capture': (None, [C.c_void_p, C.c_int64]),
|
||||
'fc_sent': (None, [C.c_void_p, C.c_uint32, C.c_int, C.c_int64]),
|
||||
'fc_ack': (C.c_int, [C.c_void_p, C.c_uint32, C.c_int64]),
|
||||
'fc_update': (C.c_int, [C.c_void_p, C.c_int64]),
|
||||
'fc_value': (C.c_int64, [C.c_void_p, C.c_int]),
|
||||
'fc_capture_open': (C.c_void_p, [C.c_char_p, GATE, C.c_char_p, C.c_int]),
|
||||
'fc_capture_pull': (C.c_int, [C.c_void_p, C.POINTER(Encoded)]),
|
||||
'fc_capture_error': (C.c_char_p, [C.c_void_p]),
|
||||
'fc_capture_bitrate': (None, [C.c_void_p, C.c_int]),
|
||||
'fc_capture_key': (None, [C.c_void_p]), 'fc_capture_close': (None, [C.c_void_p]),
|
||||
}
|
||||
if sys.platform.startswith('linux'):
|
||||
signatures.update({
|
||||
'fc_portal_select': (C.c_void_p, [C.c_char_p, C.c_int]),
|
||||
'fc_portal_close': (None, [C.c_void_p]),
|
||||
'fc_portal_value': (C.c_int, [C.c_void_p, C.c_int]),
|
||||
'fc_portal_input': (C.c_int, [C.c_void_p, C.c_int, C.c_double, C.c_double, C.c_int, C.c_int]),
|
||||
})
|
||||
for name, (result, args) in signatures.items():
|
||||
fn = getattr(self.lib, name)
|
||||
fn.restype, fn.argtypes = result, args
|
||||
self.lib.fc_gst_init()
|
||||
|
||||
def has(self, name):
|
||||
return bool(self.lib.fc_has_element(name.encode()))
|
||||
|
||||
def now(self):
|
||||
return self.lib.fc_now()
|
||||
|
||||
|
||||
class Controller:
|
||||
def __init__(self, lib, fps, ceiling):
|
||||
self.lib, self.lock = lib, threading.RLock()
|
||||
self.enabled = os.environ.get('FRAME_MAC_VIEW_ADAPT') != '0'
|
||||
self.ptr = lib.fc_new(fps, self.enabled)
|
||||
if not self.ptr:
|
||||
raise MemoryError('Unable to allocate the streaming controller')
|
||||
lib.fc_ceiling(self.ptr, ceiling)
|
||||
self.events = []
|
||||
|
||||
def state(self):
|
||||
with self.lock:
|
||||
names = ['target', 'ceiling', 'tier', 'fps', 'scale', 'baseRtt', 'inFlight', 'slack']
|
||||
out = {k: self.lib.fc_value(self.ptr, i) for i, k in enumerate(names)}
|
||||
out.update(scale=out['scale'] / 100, baseRtt=out['baseRtt'] / 1000,
|
||||
slack=out['slack'] / 1000, adapt=self.enabled)
|
||||
return out
|
||||
|
||||
def call(self, name, *args):
|
||||
with self.lock:
|
||||
return getattr(self.lib, 'fc_' + name)(self.ptr, *args)
|
||||
|
||||
def update(self, now):
|
||||
with self.lock:
|
||||
old = self.state()
|
||||
result = self.lib.fc_update(self.ptr, now)
|
||||
state = self.state()
|
||||
if state['target'] < old['target'] or state['tier'] != old['tier']:
|
||||
self.events.append({'t': now, 'e': 'target %s bit/s; tier %s' % (state['target'], state['tier'])})
|
||||
self.events = self.events[-200:]
|
||||
return result
|
||||
|
||||
def close(self):
|
||||
with self.lock:
|
||||
if self.ptr:
|
||||
self.lib.fc_free(self.ptr)
|
||||
self.ptr = None
|
||||
|
||||
|
||||
def dimensions(w, h, maximum):
|
||||
scale = min(1, maximum / max(w, h))
|
||||
return max(2, int(w * scale) // 2 * 2), max(2, int(h * scale) // 2 * 2)
|
||||
|
||||
|
||||
def pipeline(source, platform, encoder, w, h, fps, bitrate, codec='h264'):
|
||||
"""Only locally constructed numeric source IDs enter the pipeline parser."""
|
||||
if source['src'] == 'test':
|
||||
capture = 'videotestsrc is-live=true pattern=ball'
|
||||
elif platform == 'win32':
|
||||
kind, ident = source['src'].split(':')
|
||||
if kind not in ('window', 'display') or not re.fullmatch(r'[0-9]+', ident):
|
||||
raise ValueError('Invalid Windows capture source')
|
||||
prop = 'window-handle' if kind == 'window' else 'monitor-handle'
|
||||
capture = 'd3d11screencapturesrc capture-api=wgc show-cursor=true show-border=true %s=%d' % (prop, int(ident))
|
||||
else:
|
||||
capture = 'pipewiresrc fd=%d path=%d do-timestamp=true' % (source['fd'], source['node'])
|
||||
# The source gate runs before conversion/encoding. There is no leaky queue
|
||||
# of H.264 frames; every encoded reference frame reaches the socket.
|
||||
raw = '%s ! video/x-raw,framerate=%d/1 ! identity name=gate ! videoconvert ! videoscale ! video/x-raw,width=%d,height=%d' % (capture, fps, w, h)
|
||||
if codec == 'jpeg':
|
||||
enc = 'jpegenc name=enc quality=80'
|
||||
parse = ''
|
||||
else:
|
||||
choices = {
|
||||
'mfh264enc': 'mfh264enc name=enc low-latency=true bframes=0 gop-size=60',
|
||||
'vah264enc': 'vah264enc name=enc b-frames=0 key-int-max=60',
|
||||
'x264enc': 'x264enc name=enc tune=zerolatency speed-preset=ultrafast bframes=0 key-int-max=60',
|
||||
}
|
||||
enc = choices[encoder] + ' bitrate=%d' % max(1, bitrate // 1000)
|
||||
raw += ',format=NV12' if encoder != 'x264enc' else ',format=I420'
|
||||
parse = ' ! h264parse config-interval=-1 ! video/x-h264,stream-format=byte-stream,alignment=au'
|
||||
return raw + ' ! ' + enc + parse + ' ! appsink name=out sync=false max-buffers=2 drop=false'
|
||||
|
||||
|
||||
# X11 keysyms work on Wayland through the RemoteDesktop portal too. The host's
|
||||
# keyboard layout handles physical keys; Unicode text is sent as Unicode keysyms.
|
||||
KEYSYMS = {'Enter': 0xff0d, 'Tab': 0xff09, 'Backspace': 0xff08, 'Escape': 0xff1b,
|
||||
'Delete': 0xffff, 'ArrowLeft': 0xff51, 'ArrowUp': 0xff52, 'ArrowRight': 0xff53,
|
||||
'ArrowDown': 0xff54, 'Home': 0xff50, 'End': 0xff57, 'PageUp': 0xff55,
|
||||
'PageDown': 0xff56, 'Shift': 0xffe1, 'Control': 0xffe3, 'Alt': 0xffe9, 'Meta': 0xffeb}
|
||||
|
||||
|
||||
class PortalInput:
|
||||
def __init__(self, native, source):
|
||||
self.lib, self.source = native.lib, source
|
||||
self.buttons, self.keys = set(), set()
|
||||
|
||||
def emit(self, kind, x=0, y=0, code=0, down=0):
|
||||
if not self.lib.fc_portal_input(self.source['portal'], kind, x, y, code, down):
|
||||
raise RuntimeError('The desktop portal refused input; check the sharing permission')
|
||||
|
||||
def release(self):
|
||||
for b in list(self.buttons):
|
||||
self.emit(1, code=b, down=0)
|
||||
self.buttons.discard(b)
|
||||
for k in list(self.keys):
|
||||
self.emit(3, code=k, down=0)
|
||||
self.keys.discard(k)
|
||||
|
||||
def handle(self, m):
|
||||
t = m['t']
|
||||
if t == 'release':
|
||||
return self.release()
|
||||
if t in ('m', 'wheel'):
|
||||
x, y = (max(0, min(1, float(m.get(k, 0)))) for k in ('x', 'y'))
|
||||
self.emit(0, x * (self.source['w'] - 1), y * (self.source['h'] - 1))
|
||||
if t == 'm' and m.get('e') in ('up', 'down'):
|
||||
b = {0: 0x110, 1: 0x112, 2: 0x111}.get(m.get('b', 0))
|
||||
if b is not None:
|
||||
down = m['e'] == 'down'
|
||||
self.emit(1, code=b, down=down)
|
||||
(self.buttons.add if down else self.buttons.discard)(b)
|
||||
elif t == 'wheel':
|
||||
self.emit(2, max(-4096, min(4096, float(m.get('dx', 0)))), max(-4096, min(4096, float(m.get('dy', 0)))))
|
||||
elif t == 'k':
|
||||
key = str(m.get('key', ''))
|
||||
k = KEYSYMS.get(key) or (ord(key) if len(key) == 1 and ord(key) < 256 else
|
||||
0x1000000 + ord(key) if len(key) == 1 else 0)
|
||||
if k:
|
||||
down = m.get('e') == 'down'
|
||||
self.emit(3, code=k, down=down)
|
||||
(self.keys.add if down else self.keys.discard)(k)
|
||||
elif t == 'text':
|
||||
for ch in str(m.get('s', ''))[:4096]:
|
||||
k = ord(ch) if ord(ch) < 256 else 0x1000000 + ord(ch)
|
||||
self.emit(3, code=k, down=1)
|
||||
self.emit(3, code=k, down=0)
|
||||
|
||||
|
||||
class Windows:
|
||||
"""WGC source handles and SendInput. No elevation or global input hook."""
|
||||
def __init__(self):
|
||||
from ctypes import wintypes as W
|
||||
self.W = W
|
||||
self.user = C.WinDLL('user32', use_last_error=True)
|
||||
self.dwm = C.WinDLL('dwmapi')
|
||||
self.user.SetProcessDpiAwarenessContext.argtypes = [C.c_void_p]
|
||||
self.user.SetProcessDpiAwarenessContext(C.c_void_p(-4)) # per-monitor v2
|
||||
self.user.IsWindow.argtypes = [W.HWND]
|
||||
self.user.IsWindowVisible.argtypes = [W.HWND]
|
||||
self.user.IsIconic.argtypes = [W.HWND]
|
||||
self.user.GetWindowTextLengthW.argtypes = [W.HWND]
|
||||
self.user.GetWindowTextW.argtypes = [W.HWND, W.LPWSTR, C.c_int]
|
||||
self.user.GetWindowRect.argtypes = [W.HWND, C.POINTER(W.RECT)]
|
||||
self.user.SetForegroundWindow.argtypes = [W.HWND]
|
||||
self.dwm.DwmGetWindowAttribute.argtypes = [W.HWND, W.DWORD, C.c_void_p, W.DWORD]
|
||||
self.callback = C.WINFUNCTYPE(W.BOOL, W.HWND, W.LPARAM)
|
||||
self.monitor_callback = C.WINFUNCTYPE(W.BOOL, W.HMONITOR, W.HDC, C.POINTER(W.RECT), W.LPARAM)
|
||||
self.user.EnumWindows.argtypes = [self.callback, W.LPARAM]
|
||||
self.user.EnumDisplayMonitors.argtypes = [W.HDC, C.c_void_p, self.monitor_callback, W.LPARAM]
|
||||
class Mouse(C.Structure):
|
||||
_fields_ = [('dx', W.LONG), ('dy', W.LONG), ('data', W.DWORD), ('flags', W.DWORD),
|
||||
('time', W.DWORD), ('extra', C.c_size_t)]
|
||||
class Key(C.Structure):
|
||||
_fields_ = [('vk', W.WORD), ('scan', W.WORD), ('flags', W.DWORD), ('time', W.DWORD), ('extra', C.c_size_t)]
|
||||
class Union(C.Union):
|
||||
_fields_ = [('mouse', Mouse), ('key', Key)]
|
||||
class Input(C.Structure):
|
||||
_fields_ = [('type', W.DWORD), ('data', Union)]
|
||||
self.Mouse, self.Key, self.Input = Mouse, Key, Input
|
||||
self.user.SendInput.argtypes = [W.UINT, C.POINTER(Input), C.c_int]
|
||||
self.user.SendInput.restype = W.UINT
|
||||
|
||||
def rect(self, hwnd):
|
||||
r = self.W.RECT()
|
||||
if not self.user.IsWindow(hwnd) or self.user.IsIconic(hwnd):
|
||||
raise RuntimeError('The captured window closed or was minimized')
|
||||
# WGC captures the visible extended frame, excluding invisible resize borders.
|
||||
if self.dwm.DwmGetWindowAttribute(hwnd, 9, C.byref(r), C.sizeof(r)):
|
||||
if not self.user.GetWindowRect(hwnd, C.byref(r)):
|
||||
raise RuntimeError('Cannot locate the captured window')
|
||||
return r.left, r.top, r.right - r.left, r.bottom - r.top
|
||||
|
||||
def sources(self):
|
||||
windows, displays = [], []
|
||||
def window(hwnd, _):
|
||||
if not self.user.IsWindowVisible(hwnd) or self.user.IsIconic(hwnd):
|
||||
return True
|
||||
n = self.user.GetWindowTextLengthW(hwnd)
|
||||
if n:
|
||||
title = C.create_unicode_buffer(n + 1)
|
||||
self.user.GetWindowTextW(hwnd, title, n + 1)
|
||||
try:
|
||||
x, y, w, h = self.rect(hwnd)
|
||||
if w > 0 and h > 0:
|
||||
windows.append(dict(src='window:%d' % hwnd, id=hwnd, title=title.value,
|
||||
app='Windows', w=w, h=h, x=x, y=y))
|
||||
except RuntimeError:
|
||||
pass
|
||||
return True
|
||||
def monitor(handle, dc, rect, data):
|
||||
r = rect.contents
|
||||
displays.append(dict(src='display:%d' % handle, name='Display %d' % (len(displays) + 1),
|
||||
x=r.left, y=r.top, w=r.right-r.left, h=r.bottom-r.top))
|
||||
return True
|
||||
self.user.EnumWindows(self.callback(window), 0)
|
||||
self.user.EnumDisplayMonitors(None, None, self.monitor_callback(monitor), 0)
|
||||
return windows, displays
|
||||
|
||||
def send(self, mouse=None, key=None):
|
||||
inp = self.Input()
|
||||
inp.type = 0 if mouse else 1
|
||||
if mouse:
|
||||
inp.data.mouse = self.Mouse(*mouse, 0, 0)
|
||||
else:
|
||||
inp.data.key = self.Key(*key, 0, 0)
|
||||
if self.user.SendInput(1, C.byref(inp), C.sizeof(inp)) != 1:
|
||||
raise RuntimeError('Windows refused input (elevated apps and the secure desktop cannot be controlled)')
|
||||
|
||||
|
||||
class WindowsInput:
|
||||
def __init__(self, host, source):
|
||||
self.host, self.source = host, source
|
||||
self.buttons, self.keys = set(), set()
|
||||
|
||||
def release(self):
|
||||
for button in list(self.buttons):
|
||||
self.host.send(mouse=(0, 0, 0, {0: 4, 1: 0x40, 2: 0x10}[button]))
|
||||
self.buttons.discard(button)
|
||||
for vk in list(self.keys):
|
||||
self.host.send(key=(vk, 0, 2))
|
||||
self.keys.discard(vk)
|
||||
|
||||
def handle(self, m):
|
||||
t, u = m['t'], self.host.user
|
||||
if t == 'release':
|
||||
return self.release()
|
||||
if t in ('m', 'wheel'):
|
||||
source = self.source
|
||||
if source['src'].startswith('window:'):
|
||||
hwnd = int(source['src'].split(':')[1])
|
||||
x, y, w, h = self.host.rect(hwnd)
|
||||
if m.get('e') == 'down':
|
||||
u.SetForegroundWindow(hwnd)
|
||||
else:
|
||||
x, y, w, h = (source[k] for k in ('x', 'y', 'w', 'h'))
|
||||
x += max(0, min(1, float(m.get('x', 0)))) * (w - 1)
|
||||
y += max(0, min(1, float(m.get('y', 0)))) * (h - 1)
|
||||
left, top, width, height = (u.GetSystemMetrics(i) for i in (76, 77, 78, 79))
|
||||
self.host.send(mouse=(round((x-left)*65535/max(1,width-1)), round((y-top)*65535/max(1,height-1)), 0, 0xc001))
|
||||
if t == 'm' and m.get('e') in ('up', 'down'):
|
||||
b = m.get('b', 0)
|
||||
if b in (0, 1, 2):
|
||||
down = m['e'] == 'down'
|
||||
self.host.send(mouse=(0, 0, 0, {0: 2, 1: 0x20, 2: 8}[b] * (1 if down else 2)))
|
||||
(self.buttons.add if down else self.buttons.discard)(b)
|
||||
elif t == 'wheel':
|
||||
for name, flag, direction in (('dy', 0x800, -1), ('dx', 0x1000, 1)):
|
||||
delta = int(max(-4096, min(4096, float(m.get(name, 0))))) * direction
|
||||
if delta:
|
||||
self.host.send(mouse=(0, 0, delta & 0xffffffff, flag))
|
||||
elif t == 'k':
|
||||
code, down = str(m.get('code', '')), m.get('e') == 'down'
|
||||
special = {'Enter': 13, 'Escape': 27, 'Tab': 9, 'Backspace': 8, 'Space': 32,
|
||||
'ArrowLeft': 37, 'ArrowUp': 38, 'ArrowRight': 39, 'ArrowDown': 40,
|
||||
'Delete': 46, 'Home': 36, 'End': 35, 'PageUp': 33, 'PageDown': 34,
|
||||
'ShiftLeft': 160, 'ShiftRight': 161, 'ControlLeft': 162, 'ControlRight': 163,
|
||||
'AltLeft': 164, 'AltRight': 165, 'MetaLeft': 91, 'MetaRight': 92}
|
||||
vk = special.get(code, 0)
|
||||
if re.fullmatch(r'Key[A-Z]', code) or re.fullmatch(r'Digit[0-9]', code):
|
||||
vk = ord(code[-1])
|
||||
if vk:
|
||||
self.host.send(key=(vk, 0, 0 if down else 2))
|
||||
(self.keys.add if down else self.keys.discard)(vk)
|
||||
elif down and len(str(m.get('key', ''))) == 1:
|
||||
self.handle({'t': 'text', 's': m['key']})
|
||||
elif t == 'text':
|
||||
raw = str(m.get('s', ''))[:4096].encode('utf-16-le')
|
||||
for i in range(0, len(raw), 2):
|
||||
unit = int.from_bytes(raw[i:i+2], 'little')
|
||||
self.host.send(key=(0, unit, 4))
|
||||
self.host.send(key=(0, unit, 6))
|
||||
@@ -0,0 +1,57 @@
|
||||
"""PC host selection behind the existing MacView tunnel/panel controller.
|
||||
|
||||
The public /api/macview name and viewer URL remain compatible with the Mac
|
||||
base branch. Only helper launch, platform availability and source validation
|
||||
are different. No second viewer, SSH supervisor or benchmark launcher.
|
||||
"""
|
||||
import os
|
||||
import sys
|
||||
from frame_macview import MacView, MacViewError, ROOT
|
||||
from frame_pc_capture import LIBRARY, NATIVE
|
||||
|
||||
|
||||
class PCView(MacView):
|
||||
host = 'windows' if sys.platform == 'win32' else 'linux'
|
||||
|
||||
def __init__(self, *args, **kwargs):
|
||||
super().__init__(*args, **kwargs)
|
||||
self.prefer_usb = False # Mac networksetup probe is platform-specific
|
||||
|
||||
def unavailable(self):
|
||||
if sys.platform not in ('win32', 'linux'):
|
||||
return 'PC streaming needs Windows or a Linux desktop.'
|
||||
if not LIBRARY.is_file():
|
||||
return 'The PC streaming libraries are missing from this build. See docs/pc-in-headset.md.'
|
||||
return None
|
||||
|
||||
def state(self):
|
||||
result = super().state()
|
||||
result["host"] = self.host
|
||||
return result
|
||||
|
||||
def build(self):
|
||||
pass # native libraries are built and bundled with the app
|
||||
|
||||
def agent_command(self):
|
||||
return [sys.executable, str(ROOT / 'ui' / 'frame_pc_agent.py')]
|
||||
|
||||
def agent_environment(self):
|
||||
env = super().agent_environment()
|
||||
env['GST_PLUGIN_PATH_1_0'] = str(NATIVE / 'lib' / 'gstreamer-1.0')
|
||||
env['GST_PLUGIN_SYSTEM_PATH_1_0'] = ''
|
||||
env['GST_REGISTRY_1_0'] = str(NATIVE.parent / 'registry.bin') if os.access(NATIVE.parent, os.W_OK) else os.path.join(os.path.expanduser('~'), '.cache', 'frame-control-gst.bin')
|
||||
env['GST_REGISTRY_FORK'] = 'no'
|
||||
if sys.platform == 'win32':
|
||||
env['PATH'] = str(NATIVE / 'bin') + os.pathsep + env.get('PATH', '')
|
||||
else:
|
||||
env['LD_LIBRARY_PATH'] = str(NATIVE / 'lib') + os.pathsep + env.get('LD_LIBRARY_PATH', '')
|
||||
return env
|
||||
|
||||
def _show(self, src, quality, width, height):
|
||||
if src.startswith('separate:'):
|
||||
raise MacViewError('Separate virtual displays are a Mac-only feature. Choose a window or screen.')
|
||||
return super()._show(src, quality, width, height)
|
||||
|
||||
|
||||
def host_view(*args, **kwargs):
|
||||
return (MacView if sys.platform == 'darwin' else PCView)(*args, **kwargs)
|
||||
@@ -0,0 +1,88 @@
|
||||
"""The desktop stream's timing records, in the Mac viewer/bench schema.
|
||||
|
||||
Times are host monotonic microseconds. Zero means unknown, never a fabricated
|
||||
capture/display timestamp. Storage is bounded to the same 4096/512 records as
|
||||
Stats.swift. The existing macview-bench.py grades these records unchanged.
|
||||
"""
|
||||
from collections import OrderedDict
|
||||
import threading
|
||||
|
||||
|
||||
def distribution(values):
|
||||
values = sorted(v for v in values if v is not None)
|
||||
if not values:
|
||||
return {}
|
||||
return {key: round(values[min(len(values)-1, int((len(values)-1)*p + .5))], 1)
|
||||
for key, p in (('p50', .5), ('p95', .95))}
|
||||
|
||||
|
||||
class Stats:
|
||||
def __init__(self, now):
|
||||
self.now, self.lock = now, threading.RLock()
|
||||
self.frames, self.inputs = OrderedDict(), OrderedDict()
|
||||
self.captured = self.skipped = self.dropped = 0
|
||||
self.rtt, self.decoder, self.synced = 0, '', False
|
||||
self.sequence = 0
|
||||
|
||||
def add(self, **values):
|
||||
with self.lock:
|
||||
self.sequence = (self.sequence + 1) & 0xffffffff
|
||||
record = dict.fromkeys(('s', 'k', 'b', 'w', 'h', 'cap', 'arr', 'e0', 'e1', 'snd',
|
||||
'wire', 'rx', 'dec', 'drw', 'vs', 'echo', 'tier', 'br'), 0)
|
||||
record.update(values, s=self.sequence)
|
||||
self.frames[self.sequence] = record
|
||||
while len(self.frames) > 4096:
|
||||
self.frames.popitem(last=False)
|
||||
for event in self.inputs.values():
|
||||
if not event['frame'] and event['inj'] <= record['cap']:
|
||||
event.update(frame=self.sequence, cap=record['cap'])
|
||||
record['echo'] = event['id']
|
||||
return record
|
||||
|
||||
def report(self, message):
|
||||
with self.lock:
|
||||
if message['t'] == 'rx':
|
||||
frame = self.frames.get(message.get('s'))
|
||||
if frame and isinstance(message.get('r'), (int, float)):
|
||||
frame['rx'] = int(message['r'])
|
||||
elif message['t'] == 'fd':
|
||||
for item in message.get('f', [])[:4096]:
|
||||
if isinstance(item, list) and len(item) >= 4:
|
||||
frame = self.frames.get(item[0])
|
||||
if frame:
|
||||
for name, val in zip(('dec', 'drw', 'vs'), item[1:4]):
|
||||
if isinstance(val, (int, float)):
|
||||
frame[name] = int(val)
|
||||
self.dropped += max(0, int(message.get('drop', 0)))
|
||||
elif message['t'] == 'clock':
|
||||
self.rtt = float(message.get('rtt', 0))
|
||||
self.decoder = str(message.get('dec', ''))[:1024]
|
||||
self.synced = True
|
||||
|
||||
def input(self, m):
|
||||
if not isinstance(m.get('i'), int) or not m['i']:
|
||||
return
|
||||
with self.lock:
|
||||
self.inputs[m['i']] = dict(id=m['i'], kind=m['t'], tv=m.get('tv') or 0,
|
||||
inj=self.now(), frame=0, cap=0)
|
||||
while len(self.inputs) > 512:
|
||||
self.inputs.popitem(last=False)
|
||||
|
||||
def summary(self):
|
||||
with self.lock:
|
||||
now = self.now()
|
||||
frames = [f for f in self.frames.values() if f['cap'] >= now-2000000]
|
||||
out = dict(captured=self.captured, skipped=self.skipped, dropped=self.dropped,
|
||||
rtt=self.rtt, decoder=self.decoder, synced=self.synced,
|
||||
fps=sum(f['vs'] > 0 for f in frames)/2, sentFps=len(frames)/2,
|
||||
mbps=round(sum(f['b'] for f in frames)*4/1e6, 2))
|
||||
for key, start, end in (('capture', 'cap', 'arr'), ('queue', 'arr', 'e0'),
|
||||
('encode', 'e0', 'e1'), ('network', 'e1', 'rx'),
|
||||
('decode', 'rx', 'dec'), ('draw', 'dec', 'drw'), ('total', 'cap', 'drw')):
|
||||
out[key] = distribution([(f[end]-f[start])/1000 for f in frames if f[start] and f[end]])
|
||||
return out
|
||||
|
||||
def snapshot(self, since=0, settle=1500000):
|
||||
with self.lock:
|
||||
return dict(frames=[dict(f) for f in self.frames.values() if f['s'] > since and f['cap'] <= self.now()-settle],
|
||||
inputs=[dict(i) for i in self.inputs.values()], captured=self.captured, summary=self.summary())
|
||||
+21
-7
@@ -651,7 +651,7 @@
|
||||
</section>
|
||||
|
||||
<section class="panel" id="macview" hidden>
|
||||
<div class="shelf-head"><h2>Mac in the headset</h2><span class="count" id="mvCount"></span><span class="spacer"></span>
|
||||
<div class="shelf-head"><h2 id="mvTitle">Mac in the headset</h2><span class="count" id="mvCount"></span><span class="spacer"></span>
|
||||
<select id="mvQuality" title="Picture quality and bandwidth">
|
||||
<option value="sharp">Sharp</option><option value="balanced" selected>Balanced</option>
|
||||
<option value="light">Light</option><option value="compatible">Compatible (JPEG)</option>
|
||||
@@ -668,7 +668,7 @@
|
||||
<button class="small danger" data-mv="stopall">Stop all</button>
|
||||
</div>
|
||||
<div class="hint">Each window or screen becomes its own panel in the headset. Place it with the SteamVR dashboard
|
||||
(Float in World, Move, Size). Point and click with the laser, scroll with the thumbstick, and type on the Mac's keyboard.</div>
|
||||
(Float in World, Move, Size). Point and click with the laser, scroll with the thumbstick, and type on your computer’s keyboard.</div>
|
||||
</section>
|
||||
</div>
|
||||
</div>
|
||||
@@ -2218,8 +2218,16 @@ async function loadMacView(restart) {
|
||||
if (seq !== mvSeq) return;
|
||||
const box = $("macview");
|
||||
if (err) { box.hidden = false; return failed($("mvWindows"), err); }
|
||||
if (!s.available) { box.hidden = true; return; }
|
||||
if (!s.available) {
|
||||
box.hidden = !s.host;
|
||||
if (s.host) { $("mvTitle").textContent = "PC in the headset"; $("mvWindows").textContent = s.reason; }
|
||||
return;
|
||||
}
|
||||
box.hidden = false;
|
||||
const pc = s.host && s.host !== "mac";
|
||||
$("mvTitle").textContent = pc ? "PC in the headset" : "Mac in the headset";
|
||||
$("mvSeparate").closest("label").hidden = pc;
|
||||
$("mvSeparate").disabled = pc;
|
||||
const live = new Set((s.streams || []).map(x => x.src));
|
||||
const byStream = src => (s.streams || []).find(x => x.src === src);
|
||||
$("mvCount").textContent = live.size ? `${live.size} showing${s.route === "usb" ? " · over USB-C" : ""}` : "";
|
||||
@@ -2229,6 +2237,12 @@ async function loadMacView(restart) {
|
||||
$("mvPerm").hidden = !perm.length;
|
||||
$("mvPerm").innerHTML = perm.join(" ") + `<div class="row"><button class="small action" data-mv="perm">Allow…</button>
|
||||
<span class="sub">Then press Refresh. macOS may ask you to reopen Frame Control.</span></div>`;
|
||||
if (pc) {
|
||||
$("mvPerm").hidden = s.host !== "linux";
|
||||
$("mvPerm").innerHTML = `<button class="small action" data-mv="perm" ${s.selecting ? "disabled" : ""}>${s.selecting ? "Choose on your computer…" : "Choose a window or screen…"}</button>
|
||||
<span class="sub">Your desktop asks what to share and whether to allow pointer and keyboard control.</span>` +
|
||||
(s.selectionError ? `<div class="sub">${esc(s.selectionError)}</div>` : "");
|
||||
}
|
||||
const liveAs = src => live.has(src) ? src : src.startsWith("window:") && live.has("separate:" + src.slice(7)) ? "separate:" + src.slice(7) : "";
|
||||
$("mvDisplays").innerHTML = (s.displays || []).map(d =>
|
||||
mvRow(d.src, `Whole screen: ${d.name}`, `${d.w}×${d.h}${d.main ? " · main display" : ""}`, d.w, d.h, liveAs(d.src), byStream(liveAs(d.src)))).join("");
|
||||
@@ -2236,11 +2250,11 @@ async function loadMacView(restart) {
|
||||
$("mvWinCount").textContent = wins.length ? `${wins.length}` : "";
|
||||
$("mvWindows").innerHTML = !s.screen ? `<div class="sub">Windows appear here once Screen Recording is allowed.</div>`
|
||||
: wins.length ? wins.map(w => mvRow(w.src, w.title || w.app, `${w.title ? w.app + " · " : ""}${w.w}×${w.h}`, w.w, w.h, liveAs(w.src), byStream(liveAs(w.src)))).join("")
|
||||
: `<div class="sub">No windows open on this Mac's current desktop.</div>`;
|
||||
: `<div class="sub">No shared windows yet.</div>`;
|
||||
// Keep the numbers fresh while something is in the headset and the card is on screen.
|
||||
clearTimeout(mvTimer);
|
||||
const tick = () => { mvTimer = document.hidden || !box.offsetParent ? setTimeout(tick, 3000) : (loadMacView(), 0); };
|
||||
if (live.size) mvTimer = setTimeout(tick, 3000);
|
||||
if (live.size || s.selecting) mvTimer = setTimeout(tick, 3000);
|
||||
}
|
||||
$("mvRefresh").onclick = () => loadMacView(true);
|
||||
$("mvSeparate").checked = localStorage.getItem("mvSeparate") !== "0";
|
||||
@@ -2249,7 +2263,7 @@ $("macview").onclick = async e => {
|
||||
const b = e.target.closest("[data-mv]"); if (!b) return;
|
||||
const src = b.dataset.src;
|
||||
if (b.dataset.mv === "show") {
|
||||
const as = $("mvSeparate").checked && src.startsWith("window:") ? "separate:" + src.slice(7) : src;
|
||||
const as = !$("mvSeparate").disabled && $("mvSeparate").checked && src.startsWith("window:") ? "separate:" + src.slice(7) : src;
|
||||
await act(`Show ${b.dataset.title} in the headset`, async () => {
|
||||
const r = await api("/api/macview", { action: "show", src: as, quality: $("mvQuality").value,
|
||||
w: +b.dataset.w || undefined, h: +b.dataset.h || undefined });
|
||||
@@ -2260,7 +2274,7 @@ $("macview").onclick = async e => {
|
||||
} else if (b.dataset.mv === "stopall") {
|
||||
await act("Stop everything in the headset", () => api("/api/macview", { action: "stop" }), b);
|
||||
} else if (b.dataset.mv === "perm") {
|
||||
await act("Ask macOS for permission", () => api("/api/macview", { action: "permissions" }), b);
|
||||
await act("Choose what to share", () => api("/api/macview", { action: "permissions" }), b);
|
||||
}
|
||||
setTimeout(loadMacView, 800);
|
||||
};
|
||||
|
||||
+5
-5
@@ -22,7 +22,7 @@
|
||||
</head>
|
||||
<body>
|
||||
<canvas id="c" width="16" height="9"></canvas>
|
||||
<div id="note">Connecting to the Mac…</div>
|
||||
<div id="note">Connecting to your computer…</div>
|
||||
<div id="hud" hidden></div>
|
||||
<script>
|
||||
"use strict";
|
||||
@@ -243,14 +243,14 @@ async function connect() {
|
||||
// Refused before ever getting a key: the ticket expired or was revoked,
|
||||
// and retrying can't help.
|
||||
if (!reconnectKey && ++failedStarts >= 3) {
|
||||
show("This view has expired. Press Show in Frame Control on the Mac to open it again.");
|
||||
show("This view has expired. Press Show in Frame Control on your computer to open it again.");
|
||||
return;
|
||||
}
|
||||
hideNoteOnFrame = true;
|
||||
// The Mac may be asleep or the tunnel restarting: keep trying.
|
||||
retry = Math.min(retry + 1, 6);
|
||||
// Keep the Mac's reason (a missing permission, a closed window) on screen.
|
||||
show(lastError ? `${lastError} Retrying…` : "Lost the Mac. Reconnecting…");
|
||||
show(lastError ? `${lastError} Retrying…` : "Lost the computer. Reconnecting…");
|
||||
setTimeout(connect, 500 * 2 ** retry);
|
||||
};
|
||||
}
|
||||
@@ -269,7 +269,7 @@ function control(m) {
|
||||
document.title = tag ? `${name} [${tag}]` : name;
|
||||
if (!m.input) {
|
||||
hideNoteOnFrame = false; // a warning, not a connection notice: let it stay its 8 s
|
||||
show("Clicks and keys need Accessibility permission on the Mac (Frame Control asks for it).", 8000);
|
||||
show(m.inputMessage || "Clicks and keys need Accessibility permission on the Mac (Frame Control asks for it).", 8000);
|
||||
}
|
||||
} else if (m.t === "error") {
|
||||
stats.errors.push(m.message);
|
||||
@@ -391,7 +391,7 @@ addEventListener("keydown", e => onKey(e, true));
|
||||
addEventListener("keyup", e => onKey(e, false));
|
||||
addEventListener("blur", () => send({ t: "release" }));
|
||||
|
||||
show("Connecting to the Mac…");
|
||||
show("Connecting to your computer…");
|
||||
connect();
|
||||
</script>
|
||||
</body>
|
||||
|
||||
+2
-1
@@ -41,6 +41,7 @@ import frame_apk_versions # noqa: E402
|
||||
import frame_catalog # noqa: E402
|
||||
import frame_host # noqa: E402
|
||||
import frame_macview # noqa: E402
|
||||
import frame_pcview # noqa: E402
|
||||
import frame_store # noqa: E402
|
||||
import frame_titles # noqa: E402
|
||||
import frame_webinstall # noqa: E402
|
||||
@@ -1232,7 +1233,7 @@ def _sweep_one(prefix, d):
|
||||
|
||||
# The tunnel gets its own connection: the shared master's options would win
|
||||
# over anything added after them.
|
||||
macview = frame_macview.MacView(["ssh", "-o", "BatchMode=yes", "-o", "ConnectTimeout=8"],
|
||||
macview = frame_pcview.host_view(["ssh", "-o", "BatchMode=yes", "-o", "ConnectTimeout=8"],
|
||||
lambda remote, stdin=None, timeout=30: ssh(remote, stdin=stdin, timeout=timeout),
|
||||
FRAME, track=_live_tunnels.add)
|
||||
|
||||
|
||||
Reference in new issue
Block a user