diff --git a/.github/workflows/pc-host.yml b/.github/workflows/pc-host.yml new file mode 100644 index 0000000..4a0699d --- /dev/null +++ b/.github/workflows/pc-host.yml @@ -0,0 +1,53 @@ +name: PC host libraries +on: + pull_request: + paths: ["desktop/**", "ui/**", "tests/**", "mac/frame-mac-view/**", "scripts/macview-bench.py", ".github/workflows/pc-host.yml"] + workflow_dispatch: +jobs: + linux: + strategy: + fail-fast: false + matrix: + os: [ubuntu-latest, ubuntu-24.04-arm] + runs-on: ${{ matrix.os }} + steps: + - uses: actions/checkout@v4 + - name: Native libraries + run: | + sudo apt-get update -qq + sudo apt-get install -y libgstreamer1.0-dev libgstreamer-plugins-base1.0-dev libglib2.0-dev gstreamer1.0-plugins-base gstreamer1.0-plugins-good gstreamer1.0-plugins-bad gstreamer1.0-plugins-ugly gstreamer1.0-pipewire + python3 desktop/build.py + - name: Tests (real x264, fake capture consent/input) + env: + FRAME_PC_REQUIRE_NATIVE: "1" + run: python3 -m unittest discover -s tests -p 'test_pc*.py' -v + - uses: actions/upload-artifact@v4 + with: + name: pc-host-${{ matrix.os }} + path: desktop/bundle + windows: + runs-on: windows-latest + steps: + - uses: actions/checkout@v4 + - uses: msys2/setup-msys2@v2 + with: + msystem: UCRT64 + update: true + install: >- + mingw-w64-ucrt-x86_64-gcc + mingw-w64-ucrt-x86_64-pkgconf + mingw-w64-ucrt-x86_64-python + mingw-w64-ucrt-x86_64-gstreamer + mingw-w64-ucrt-x86_64-gst-plugins-base + mingw-w64-ucrt-x86_64-gst-plugins-good + mingw-w64-ucrt-x86_64-gst-plugins-bad + mingw-w64-ucrt-x86_64-gst-plugins-ugly + - name: Native build and tests (x264; WGC and MF need a desktop/GPU) + shell: msys2 {0} + run: | + python desktop/build.py + FRAME_PC_REQUIRE_NATIVE=1 python -m unittest discover -s tests -p 'test_pc*.py' -v + - uses: actions/upload-artifact@v4 + with: + name: pc-host-windows + path: desktop/bundle diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index 04cfeaf..d8b23ed 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -7,7 +7,7 @@ on: push: tags: ["v*"] pull_request: - paths: ["app/**", "ui/**", "scripts/**", "frame/**", "apk-catalog/**", ".github/workflows/release.yml"] + paths: ["app/**", "ui/**", "desktop/**", "mac/**", "scripts/**", "frame/**", "apk-catalog/**", ".github/workflows/release.yml"] workflow_dispatch: permissions: @@ -28,12 +28,41 @@ jobs: - os: ubuntu-latest script: dist:linux files: app/dist/*.AppImage app/dist/*.deb + - os: ubuntu-24.04-arm + script: dist:linux:arm64 + files: app/dist/*.AppImage app/dist/*.deb runs-on: ${{ matrix.os }} steps: - uses: actions/checkout@v4 - uses: actions/setup-node@v4 with: node-version: "24" + - name: Linux streaming libraries + if: runner.os == 'Linux' + shell: bash + run: | + sudo apt-get update -qq + sudo apt-get install -y libgstreamer1.0-dev libgstreamer-plugins-base1.0-dev libglib2.0-dev gstreamer1.0-plugins-base gstreamer1.0-plugins-good gstreamer1.0-plugins-bad gstreamer1.0-plugins-ugly gstreamer1.0-pipewire + arch=x64; if [ "$(uname -m)" = aarch64 ]; then arch=arm64; fi + python3 desktop/build.py --out "app/build/desktop/linux-$arch" + - uses: msys2/setup-msys2@v2 + if: runner.os == 'Windows' + with: + msystem: UCRT64 + update: true + install: >- + mingw-w64-ucrt-x86_64-gcc + mingw-w64-ucrt-x86_64-pkgconf + mingw-w64-ucrt-x86_64-python + mingw-w64-ucrt-x86_64-gstreamer + mingw-w64-ucrt-x86_64-gst-plugins-base + mingw-w64-ucrt-x86_64-gst-plugins-good + mingw-w64-ucrt-x86_64-gst-plugins-bad + mingw-w64-ucrt-x86_64-gst-plugins-ugly + - name: Windows streaming libraries + if: runner.os == 'Windows' + shell: msys2 {0} + run: python desktop/build.py --out app/build/desktop/win-x64 - name: Build working-directory: app shell: bash diff --git a/README.md b/README.md index 98d0413..f51a5ca 100644 --- a/README.md +++ b/README.md @@ -201,6 +201,7 @@ Frame's software fits together, all checked against a real headset and labelled | [Sideloading Linux and Windows games](docs/sideloading.md) | A .zip, folder or .exe as a Steam Devkit Game, runtime detection | | [Install links for websites](docs/web-install.md) | `frame-control://install` links and manifests, the rules, a button to paste | | [Steam games](docs/steam-games.md) · [VR video](docs/vr-video.md) · [WebXR in Chromium](docs/webxr-chromium.md) | Installing and buying, watching VR180/360, the Chromium build | +| [PC in the headset](docs/pc-in-headset.md) | Windows and Linux host implementation and test coverage | | [Mac in the headset](docs/mac-in-headset.md) | Mac windows and screens as panels in the Frame, with laser and keyboard input | | [SSH](docs/ssh.md) · [Streaming](docs/streaming.md) · [Files](docs/file-transfer.md) · [Panels](docs/panels.md) · [Tailscale](docs/tailscale.md) | Topic notes | | [Frame Control for iPhone](docs/iphone.md) | The iPhone and iPad app, how it runs the server on the Frame, pairing | diff --git a/app/package.json b/app/package.json index a06fee8..da8c2d7 100644 --- a/app/package.json +++ b/app/package.json @@ -11,8 +11,9 @@ "icon": "env -u ELECTRON_RUN_AS_NODE electron build/make-icon.js", "dist": "sh ../mac/frame-mac-view/build.sh && node build/fetch-deps.js mac arm64 && electron-builder --mac --arm64 --publish never", "dist:dir": "sh ../mac/frame-mac-view/build.sh && node build/fetch-deps.js mac arm64 && electron-builder --mac --arm64 --dir", - "dist:linux": "node build/fetch-deps.js linux x64 arm64 && electron-builder --linux --x64 --arm64 --publish never", - "dist:win": "node build/fetch-deps.js win x64 && electron-builder --win --x64 --publish never" + "dist:linux": "node build/fetch-deps.js linux x64 && electron-builder --linux --x64 --publish never", + "dist:win": "node build/fetch-deps.js win x64 && electron-builder --win --x64 --publish never", + "dist:linux:arm64": "node build/fetch-deps.js linux arm64 && electron-builder --linux --arm64 --publish never" }, "devDependencies": { "electron": "^44.4.5", @@ -143,7 +144,16 @@ "entry": { "StartupWMClass": "frame-control" } - } + }, + "extraResources": [ + { + "from": "build/desktop/${os}-${arch}", + "to": "desktop/bundle", + "filter": [ + "**/*" + ] + } + ] }, "deb": { "depends": [ @@ -156,7 +166,16 @@ "zip" ], "icon": "build/icon.png", - "artifactName": "Frame-Control-win-${arch}.${ext}" + "artifactName": "Frame-Control-win-${arch}.${ext}", + "extraResources": [ + { + "from": "build/desktop/${os}-${arch}", + "to": "desktop/bundle", + "filter": [ + "**/*" + ] + } + ] }, "nsis": { "oneClick": false, diff --git a/desktop/build.py b/desktop/build.py new file mode 100644 index 0000000..d231c61 --- /dev/null +++ b/desktop/build.py @@ -0,0 +1,91 @@ +#!/usr/bin/env python3 +"""Build our native PC adapter and bundle its ordinary shared libraries. + +Run on the target architecture after installing GStreamer development packages. +No GStreamer executable or desktop streaming application is shipped or invoked. +Linux libc, display/GPU drivers and the portal service remain platform pieces. +""" +import argparse +import json +import os +from pathlib import Path +import re +import shutil +import subprocess +import sys + +ROOT = Path(__file__).resolve().parent + + +def output(*args): + return subprocess.check_output(args, text=True).strip() + + +def build(destination): + win = sys.platform == 'win32' + destination.mkdir(parents=True, exist_ok=True) + modules = ['gstreamer-1.0', 'gstreamer-app-1.0', 'gstreamer-video-1.0'] + if not win: + modules += ['gio-unix-2.0'] + import shlex + flags = shlex.split(output('pkg-config', '--cflags', '--libs', *modules)) + lib = destination / ('pc-host.dll' if win else 'pc-host.so') + cmd = [os.environ.get('CC', 'gcc' if win else 'cc'), '-std=c11', '-O2', '-shared', '-Wall', '-Wextra'] + if not win: + cmd += ['-fPIC', '-Wl,-rpath,$ORIGIN/lib'] + subprocess.run(cmd + [str(ROOT / f) for f in ('controller.c', 'capture.c', 'portal.c')] + flags + ['-o', str(lib)], check=True) + prefix = Path(output('pkg-config', '--variable=prefix', 'gstreamer-1.0')) + plugin_dir = Path(output('pkg-config', '--variable=pluginsdir', 'gstreamer-1.0')) + plugins_out, libs_out = destination / 'lib' / 'gstreamer-1.0', destination / ('bin' if win else 'lib') + plugins_out.mkdir(parents=True, exist_ok=True) + libs_out.mkdir(parents=True, exist_ok=True) + names = ['coreelements', 'videotestsrc', 'videoconvertscale', 'app', 'jpeg', 'videoparsersbad', 'x264'] + names += ['d3d11', 'mediafoundation'] if win else ['pipewire', 'va'] + pending = [lib] + for name in names: + matches = list(plugin_dir.glob('*gst' + name + ('.dll' if win else '.so'))) + if not matches: + raise SystemExit('Missing GStreamer library: ' + name) + for source in matches: + target = plugins_out / source.name + shutil.copy2(source, target) + pending.append(source) + copied = set() + while pending: + binary = pending.pop() + if win: + imported = re.findall(r'DLL Name:\s*(\S+)', output('objdump', '-p', str(binary))) + dependencies = [prefix / 'bin' / name for name in imported] + else: + dependencies = [Path(p) for p in re.findall(r'=>\s+(/\S+)', output('ldd', str(binary)))] + for dep in dependencies: + if not dep.is_file() or dep.name in copied: + continue + if not win and re.match(r'lib(c|m|dl|rt|pthread|resolv)\.so', dep.name): + continue + shutil.copy2(dep, libs_out / dep.name) + copied.add(dep.name) + pending.append(dep) + # License texts and package provenance accompany the dynamically linked + # libraries. Distribution builders keep the upstream package/source URLs. + licenses = destination / 'licenses' + licenses.mkdir(exist_ok=True) + if win: + source = prefix / 'share' / 'licenses' + if source.exists(): + shutil.copytree(source, licenses, dirs_exist_ok=True) + packages = output('pacman', '-Q') if shutil.which('pacman') else 'See MSYS2 build log' + else: + packages = output('dpkg-query', '-W', '-f=${Package} ${Version}\n') if shutil.which('dpkg-query') else '' + # Debian copyright files contain licenses and source homepage details. + for path in Path('/usr/share/doc').glob('*/copyright'): + shutil.copy2(path, licenses / (path.parent.name + '.copyright')) + (destination / 'libraries.json').write_text(json.dumps(dict(gstreamer=output('pkg-config', '--modversion', 'gstreamer-1.0'), + libraries=sorted(copied), packages=packages), indent=2)) + print(lib) + + +if __name__ == '__main__': + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument('--out', type=Path, default=ROOT / 'bundle') + build(parser.parse_args().out.resolve()) diff --git a/desktop/capture.c b/desktop/capture.c new file mode 100644 index 0000000..07c9379 --- /dev/null +++ b/desktop/capture.c @@ -0,0 +1,97 @@ +/* Frame Control's capture/encode adapter. GStreamer is a bundled library; + * capture uses WGC on Windows and the consented PipeWire fd on Linux. */ +#include "controller.h" +#include +#include +#include +#include + +typedef int (*Gate)(int stage, int64_t pts, int64_t capture, int64_t arrived); +typedef struct { + GstElement *pipeline, *encoder, *sink; + Gate gate; + GstSample *sample; + GstMapInfo map; + int mapped; + char error[512]; +} Capture; +typedef struct { const unsigned char *data; int size, key, width, height; int64_t pts; } Encoded; +FC_API int64_t fc_now(void) {return g_get_monotonic_time();} +FC_API void fc_gst_init(void) {gst_init(NULL,NULL);} +FC_API int fc_has_element(const char *name) { + GstElementFactory *f=gst_element_factory_find(name); + if(!f)return 0;gst_object_unref(f);return 1; +} +static GstPadProbeReturn probe(GstPad *pad,GstPadProbeInfo *info,gpointer data) { + Capture *c=data; GstBuffer *b=GST_PAD_PROBE_INFO_BUFFER(info); + if(!b)return GST_PAD_PROBE_OK; + int64_t now=fc_now(), cap=now; + GstClock *clock=gst_element_get_clock(c->pipeline); + if(clock && GST_BUFFER_PTS_IS_VALID(b)) { + GstClockTime current=gst_clock_get_time(clock),origin=gst_element_get_base_time(c->pipeline); + if(current>=origin && current-origin>=GST_BUFFER_PTS(b)) + cap-= (int64_t)((current-origin-GST_BUFFER_PTS(b))/1000); + } + if(clock)gst_object_unref(clock); + int stage=GPOINTER_TO_INT(g_object_get_data(G_OBJECT(pad),"stage")); + return c->gate(stage,(int64_t)GST_BUFFER_PTS(b),cap,now) ? GST_PAD_PROBE_OK : GST_PAD_PROBE_DROP; +} +FC_API Capture *fc_capture_open(const char *pipeline,Gate gate,char *error,int capacity) { + GError *e=NULL;Capture *c=g_new0(Capture,1);c->gate=gate; + c->pipeline=gst_parse_launch(pipeline,&e); + if(e || !c->pipeline) { + g_strlcpy(error,e?e->message:"No pipeline",capacity); + if(e)g_error_free(e);if(c->pipeline)gst_object_unref(c->pipeline);g_free(c);return NULL; + } + c->encoder=gst_bin_get_by_name(GST_BIN(c->pipeline),"enc"); + c->sink=gst_bin_get_by_name(GST_BIN(c->pipeline),"out"); + GstElement *raw=gst_bin_get_by_name(GST_BIN(c->pipeline),"gate"); + if(!c->encoder || !c->sink || !raw) { + g_strlcpy(error,"Pipeline is missing enc, out or gate",capacity); + if(raw)gst_object_unref(raw);if(c->encoder)gst_object_unref(c->encoder); + if(c->sink)gst_object_unref(c->sink);gst_object_unref(c->pipeline);g_free(c);return NULL; + } + GstPad *p=gst_element_get_static_pad(raw,"src"); + g_object_set_data(G_OBJECT(p),"stage",GINT_TO_POINTER(0)); + gst_pad_add_probe(p,GST_PAD_PROBE_TYPE_BUFFER,probe,c,NULL);gst_object_unref(p);gst_object_unref(raw); + p=gst_element_get_static_pad(c->encoder,"sink"); + g_object_set_data(G_OBJECT(p),"stage",GINT_TO_POINTER(1)); + gst_pad_add_probe(p,GST_PAD_PROBE_TYPE_BUFFER,probe,c,NULL);gst_object_unref(p); + gst_element_set_state(c->pipeline,GST_STATE_PLAYING); + return c; +} +FC_API int fc_capture_pull(Capture *c,Encoded *out) { + if(c->mapped) {gst_buffer_unmap(gst_sample_get_buffer(c->sample),&c->map);c->mapped=0;} + if(c->sample){gst_sample_unref(c->sample);c->sample=NULL;} + GstBus *bus=gst_element_get_bus(c->pipeline); + GstMessage *m=gst_bus_pop_filtered(bus,GST_MESSAGE_ERROR|GST_MESSAGE_EOS);gst_object_unref(bus); + if(m) { + if(GST_MESSAGE_TYPE(m)==GST_MESSAGE_ERROR) { + GError *e=NULL;char *debug=NULL;gst_message_parse_error(m,&e,&debug); + g_strlcpy(c->error,e->message,sizeof(c->error));g_error_free(e);g_free(debug); + } else g_strlcpy(c->error,"The capture source closed",sizeof(c->error)); + gst_message_unref(m);return -1; + } + c->sample=gst_app_sink_try_pull_sample(GST_APP_SINK(c->sink),100*GST_MSECOND); + if(!c->sample)return 0; + GstBuffer *b=gst_sample_get_buffer(c->sample); + if(!gst_buffer_map(b,&c->map,GST_MAP_READ))return 0; + c->mapped=1;out->data=c->map.data;out->size=(int)c->map.size; + out->key=!GST_BUFFER_FLAG_IS_SET(b,GST_BUFFER_FLAG_DELTA_UNIT);out->pts=(int64_t)GST_BUFFER_PTS(b); + const GstStructure *s=gst_caps_get_structure(gst_sample_get_caps(c->sample),0); + gst_structure_get_int(s,"width",&out->width);gst_structure_get_int(s,"height",&out->height);return 1; +} +FC_API const char *fc_capture_error(Capture *c) {return c->error;} +FC_API void fc_capture_bitrate(Capture *c,int bps) { + g_object_set(c->encoder,"bitrate",(guint)MAX(1,bps/1000),NULL); +} +FC_API void fc_capture_key(Capture *c) { + GstPad *p=gst_element_get_static_pad(c->encoder,"src"); + gst_pad_send_event(p,gst_video_event_new_upstream_force_key_unit(GST_CLOCK_TIME_NONE,TRUE,0));gst_object_unref(p); +} +FC_API void fc_capture_close(Capture *c) { + gst_element_set_state(c->pipeline,GST_STATE_NULL); + if(c->mapped)gst_buffer_unmap(gst_sample_get_buffer(c->sample),&c->map); + if(c->sample)gst_sample_unref(c->sample); + gst_object_unref(c->sink);gst_object_unref(c->encoder);gst_object_unref(c->pipeline);g_free(c); +} diff --git a/desktop/portal.c b/desktop/portal.c new file mode 100644 index 0000000..0254ba3 --- /dev/null +++ b/desktop/portal.c @@ -0,0 +1,132 @@ +/* xdg-desktop-portal RemoteDesktop + ScreenCast, one consented session per + * panel. PipeWire receives the portal fd, never the unrestricted daemon fd. + * Input uses only the devices granted in Start's response. */ +#ifndef _WIN32 +#include "controller.h" +#include +#include +#include +#include +#define BUS "org.freedesktop.portal.Desktop" +#define PATH "/org/freedesktop/portal/desktop" +#define RD "org.freedesktop.portal.RemoteDesktop" +#define SC "org.freedesktop.portal.ScreenCast" +typedef struct { + GDBusConnection *bus; + GMainContext *context; + char *session; + int fd; + unsigned node, devices; + int width,height; +} Portal; +typedef struct {GVariant *result;int done;unsigned code;} Response; +static void response(GDBusConnection *c,const char *sender,const char *path,const char *iface, + const char *signal,GVariant *params,gpointer data) { + (void)c;(void)sender;(void)path;(void)iface;(void)signal; + Response *r=data;g_variant_get(params,"(u@a{sv})",&r->code,&r->result);r->done=1; +} +static void option(GVariantBuilder *b,const char *key,GVariant *v) {g_variant_builder_add(b,"{sv}",key,v);} +static GVariant *request(Portal *p,const char *iface,const char *method,GVariant *args, + const char *token,char *error,int capacity) { + char *sender=g_strdup(g_dbus_connection_get_unique_name(p->bus)+1); + for(char *c=sender;*c;c++)if(*c=='.')*c='_'; + char *path=g_strdup_printf(PATH "/request/%s/%s",sender,token);g_free(sender); + Response r={0};GError *e=NULL; + guint sub=g_dbus_connection_signal_subscribe(p->bus,BUS,"org.freedesktop.portal.Request","Response", + path,NULL,G_DBUS_SIGNAL_FLAGS_NONE,response,&r,NULL); + GVariant *reply=g_dbus_connection_call_sync(p->bus,BUS,PATH,iface,method,args,G_VARIANT_TYPE("(o)"), + G_DBUS_CALL_FLAGS_NONE,10000,NULL,&e); + if(reply)g_variant_unref(reply); + int64_t deadline=g_get_monotonic_time()+120000000; + while(!e && !r.done && g_get_monotonic_time()context,FALSE)) {} + g_usleep(10000); + } + if(!r.done) { + GVariant *closed=g_dbus_connection_call_sync(p->bus,BUS,path,"org.freedesktop.portal.Request","Close", + NULL,NULL,G_DBUS_CALL_FLAGS_NONE,2000,NULL,NULL); + if(closed)g_variant_unref(closed); + } + g_dbus_connection_signal_unsubscribe(p->bus,sub);g_free(path); + if(e || !r.done || r.code) { + g_strlcpy(error,e?e->message:!r.done?"Screen sharing request timed out":"Screen sharing was cancelled or refused",capacity); + if(e)g_error_free(e);if(r.result)g_variant_unref(r.result);return NULL; + } + return r.result; +} +FC_API void fc_portal_close(Portal *p) { + if(!p)return; + if(p->session && p->bus) { + GVariant *r=g_dbus_connection_call_sync(p->bus,BUS,p->session,"org.freedesktop.portal.Session","Close", + NULL,NULL,G_DBUS_CALL_FLAGS_NONE,2000,NULL,NULL); + if(r)g_variant_unref(r); + } + if(p->fd>=0)close(p->fd); + g_free(p->session);if(p->bus)g_object_unref(p->bus); + if(p->context)g_main_context_unref(p->context);g_free(p); +} +FC_API Portal *fc_portal_select(char *error,int capacity) { + Portal *p=g_new0(Portal,1);p->fd=-1;p->context=g_main_context_new(); + g_main_context_push_thread_default(p->context); + GError *e=NULL;GVariant *r=NULL; + p->bus=g_bus_get_sync(G_BUS_TYPE_SESSION,NULL,&e); + if(!p->bus) {g_strlcpy(error,e->message,capacity);g_error_free(e);goto fail;} + GVariantBuilder b;char token[64],session[64]; + g_snprintf(session,sizeof(session),"fc_%u",g_random_int()); + g_snprintf(token,sizeof(token),"fc_%u",g_random_int()); + g_variant_builder_init(&b,G_VARIANT_TYPE_VARDICT); + option(&b,"handle_token",g_variant_new_string(token));option(&b,"session_handle_token",g_variant_new_string(session)); + r=request(p,RD,"CreateSession",g_variant_new("(a{sv})",&b),token,error,capacity);if(!r)goto fail; + g_variant_lookup(r,"session_handle","s",&p->session);g_variant_unref(r);r=NULL; + if(!p->session){g_strlcpy(error,"Portal did not return a session",capacity);goto fail;} + g_snprintf(token,sizeof(token),"fc_%u",g_random_int());g_variant_builder_init(&b,G_VARIANT_TYPE_VARDICT); + option(&b,"handle_token",g_variant_new_string(token));option(&b,"types",g_variant_new_uint32(3)); + r=request(p,RD,"SelectDevices",g_variant_new("(oa{sv})",p->session,&b),token,error,capacity);if(!r)goto fail;g_variant_unref(r); + g_snprintf(token,sizeof(token),"fc_%u",g_random_int());g_variant_builder_init(&b,G_VARIANT_TYPE_VARDICT); + option(&b,"handle_token",g_variant_new_string(token));option(&b,"types",g_variant_new_uint32(3)); + option(&b,"multiple",g_variant_new_boolean(FALSE));option(&b,"cursor_mode",g_variant_new_uint32(2)); + r=request(p,SC,"SelectSources",g_variant_new("(oa{sv})",p->session,&b),token,error,capacity);if(!r)goto fail;g_variant_unref(r); + g_snprintf(token,sizeof(token),"fc_%u",g_random_int());g_variant_builder_init(&b,G_VARIANT_TYPE_VARDICT); + option(&b,"handle_token",g_variant_new_string(token)); + r=request(p,RD,"Start",g_variant_new("(osa{sv})",p->session,"",&b),token,error,capacity);if(!r)goto fail; + g_variant_lookup(r,"devices","u",&p->devices); + GVariant *streams=g_variant_lookup_value(r,"streams",G_VARIANT_TYPE("a(ua{sv})")); + if(streams && g_variant_n_children(streams)>0) { + GVariant *entry=g_variant_get_child_value(streams,0),*props=NULL; + g_variant_get(entry,"(u@a{sv})",&p->node,&props); + g_variant_lookup(props,"size","(ii)",&p->width,&p->height); + g_variant_unref(props);g_variant_unref(entry); + } + if(streams)g_variant_unref(streams);g_variant_unref(r);r=NULL; + if(!p->node || p->width<=0 || p->height<=0) {g_strlcpy(error,"Portal returned no stream size",capacity);goto fail;} + g_variant_builder_init(&b,G_VARIANT_TYPE_VARDICT);GUnixFDList *fds=NULL; + r=g_dbus_connection_call_with_unix_fd_list_sync(p->bus,BUS,PATH,SC,"OpenPipeWireRemote", + g_variant_new("(oa{sv})",p->session,&b),G_VARIANT_TYPE("(h)"),G_DBUS_CALL_FLAGS_NONE,10000,NULL,&fds,NULL,&e); + if(!r){g_strlcpy(error,e->message,capacity);g_error_free(e);goto fail;} + int handle;g_variant_get(r,"(h)",&handle);g_variant_unref(r); + p->fd=g_unix_fd_list_get(fds,handle,&e);g_object_unref(fds); + if(p->fd<0){g_strlcpy(error,e->message,capacity);g_error_free(e);goto fail;} + g_main_context_pop_thread_default(p->context);return p; +fail: + g_main_context_pop_thread_default(p->context);fc_portal_close(p);return NULL; +} +FC_API int fc_portal_value(Portal *p,int field) { + switch(field){case 0:return p->fd;case 1:return (int)p->node;case 2:return p->width; + case 3:return p->height;case 4:return (int)p->devices;default:return 0;} +} +FC_API int fc_portal_input(Portal *p,int type,double x,double y,int code,int down) { + const char *method=NULL;GVariant *args=NULL;GVariantBuilder b;g_variant_builder_init(&b,G_VARIANT_TYPE_VARDICT); + if(type==0 && (p->devices&2)) { + method="NotifyPointerMotionAbsolute";args=g_variant_new("(oa{sv}udd)",p->session,&b,p->node,x,y); + } else if(type==1 && (p->devices&2)) { + method="NotifyPointerButton";args=g_variant_new("(oa{sv}iu)",p->session,&b,code,(guint)down); + } else if(type==2 && (p->devices&2)) { + method="NotifyPointerAxis";args=g_variant_new("(oa{sv}dd)",p->session,&b,x,y); + } else if(type==3 && (p->devices&1)) { + method="NotifyKeyboardKeysym";args=g_variant_new("(oa{sv}iu)",p->session,&b,code,(guint)down); + } else return 0; + GError *e=NULL;GVariant *r=g_dbus_connection_call_sync(p->bus,BUS,PATH,RD,method,args,NULL, + G_DBUS_CALL_FLAGS_NONE,2000,NULL,&e); + if(e)g_error_free(e);if(r)g_variant_unref(r);return r!=NULL; +} +#endif diff --git a/docs/pc-in-headset.md b/docs/pc-in-headset.md new file mode 100644 index 0000000..59c3784 --- /dev/null +++ b/docs/pc-in-headset.md @@ -0,0 +1,131 @@ +# PC in the headset + +Windows and Linux hosts use **Tools → PC in the headset**. The host shares a +window or screen, and the existing Frame viewer makes it a SteamVR panel. +Move it with the dashboard's Float in World, Move and Size controls. + +**Inferred / not yet verified on a desktop host:** the Windows and Linux +capture and input paths below. This work was developed on a Mac with no +Windows or Linux desktop VM. A native build or test-pattern test in CI does +not establish that desktop capture, a permission dialog, hardware encoding +or laser input works. Keep this feature in the draft/testing stage until +those paths have been tried on real hosts. + +## Own implementation, platform APIs and bundled libraries + +Frame Control owns the host agent, input routing, authentication, streaming +protocol, panel launcher and adaptation. It does not launch or require +Sunshine, OBS or another desktop-streaming app. GStreamer and its codec +plugins are ordinary libraries bundled with the Windows and Linux app; +users do not install a GStreamer application. The shared library build keeps +license texts and package provenance alongside the libraries. + +First-party alternatives considered (**documented**): Valve Remote Play +streams a game/desktop, rather than providing this per-window panel protocol; +Windows Remote Desktop opens a remote session; Linux's desktop portal is the +consent mechanism for sharing the current desktop. The chosen paths are: + +| Host | Capture | Encoding | Input | +|---|---|---|---| +| Windows | Windows.Graphics.Capture, through `d3d11screencapturesrc capture-api=wgc`; HWND or HMONITOR | Hardware Media Foundation (`mfh264enc`), low latency, no B-frames | `SendInput`, with source bounds and per-monitor DPI awareness | +| Linux | RemoteDesktop + ScreenCast portal, then the returned PipeWire fd/node | VA-API (`vah264enc`) where registered; x264 otherwise | RemoteDesktop portal notifications, using only granted pointer/keyboard devices | + +API choices are **documented**, not device verification: +[Windows capture](https://learn.microsoft.com/en-us/windows/uwp/audio-video-camera/screen-capture), +[GStreamer WGC](https://gstreamer.freedesktop.org/documentation/d3d11/d3d11screencapturesrc.html), +[Media Foundation encoder](https://gstreamer.freedesktop.org/documentation/mediafoundation/mfh264enc.html), +[ScreenCast portal](https://flatpak.github.io/xdg-desktop-portal/docs/doc-org.freedesktop.portal.ScreenCast.html), +[RemoteDesktop portal](https://flatpak.github.io/xdg-desktop-portal/docs/doc-org.freedesktop.portal.RemoteDesktop.html). + +Windows needs a WGC-capable Windows 10/11 desktop and an available hardware +Media Foundation H.264 encoder. Elevated windows and the secure desktop +cannot be driven by an ordinary Frame Control process. Minimized/closed +windows may stop producing frames. Protected content is not supported. + +On Linux, press **Choose a window or screen…** and approve the desktop's +sharing dialog. Choose another source to add another panel. Stop releases +that source's portal session; sharing it again asks for consent again. +A desktop must implement both ScreenCast and RemoteDesktop for this path; +a ScreenCast-only compositor cannot provide laser input through this API. +Cancelling or denying a dialog is reported on the card. No portal permission +is bypassed, and Frame Control does not open `/dev/uinput` or the unrestricted +PipeWire daemon on the host. + +The host's own keyboard still works. Input from the viewer uses normalized +picture coordinates, maps through the selected source's bounds, and releases +held buttons/keys on blur, disconnect and Stop. Linux requires the pointer +and keyboard grants. **Untested:** desktop-specific consent, mixed-DPI +Windows input alignment, multi-monitor layouts, hardware encoder behavior, +window resize/minimize, and non-US keyboard layouts. + +## Shared pieces + +- `ui/frame_macview.py` owns the SSH tunnel, reconnect supervision, quality + presets and panel launch for all hosts. `ui/frame_pcview.py` selects the PC + helper; `/api/macview` remains the compatible endpoint. +- `ui/mac-view.html` is the one viewer. The 17-byte big-endian frame header, + Annex-B H.264/JPEG payloads, `hello`/`ack` reconnect handshake, clock sync, + `rx`/`fd` timing reports and input messages are unchanged. +- `desktop/controller.c` is the rate controller shared by the Mac Swift + binding and the PC Python binding. Capture is gated **before** encoding; + encoded reference frames are never discarded. It keeps the Mac's bitrate + demand protection and tier hysteresis. +- PC records use the existing `Stats.swift` JSON schema, with bounded + 4096-frame/512-input storage in `ui/frame_stream_stats.py`. The benchmark's + analysis, targets and network shaping are shared, not reimplemented. + Capture timestamps describe the native pipeline's source time; they do + not prove the time at which the host compositor displayed the pixels. +- Mac virtual-display separation remains Mac-only. Windows WGC and the + Linux portal share the selected window directly. + +Current adaptation limitation: PC gating follows the shared frame-rate tier. +Live bitrate updates are applied to x264. Hardware encoders retain their +initial bitrate, and PC resolution does not yet follow the controller's +scale tier. This is an explicit remaining gap, not a measured performance +claim. + +## Build and measure + +Packaged Windows/Linux builds include `desktop/bundle/pc-host` and its shared +libraries. Source checkouts build them with `python3 desktop/build.py` after +installing GStreamer development packages (see the `PC host libraries` CI +workflow). The feature reports a missing bundle; it does not download or +install a streaming app on first use. + +The existing benchmark now accepts a PC host: + +```sh +python3 scripts/macview-bench.py run --pc --scenario test --label pc-test +# Linux: select a real source in the desktop's sharing dialog +python3 scripts/macview-bench.py run --pc --scenario capture --source choose --label linux-window +# Windows: use the HWND/monitor source ID shown by the host's /windows or /displays +python3 scripts/macview-bench.py run --pc --scenario capture --source window:12345 --label windows-window +``` + +The synthetic PC pattern uses bundled x264 so headless CI can verify the +wire protocol without claiming that a GPU was exercised. The `capture` +scenario measures the selected real source without injecting input or +assuming that it animates at 60 fps. Mac-only Chrome/virtual-display typing +and scrolling automation is not run on PC hosts. Results retain the same +latency stages and record `host_platform`, `pc_host` and `source`. CPU sampling +on PC hosts is explicitly unavailable. `--net` and `--delay` still use the +same bounded shaping relay, without administrator privileges. + +## Evidence + +- **Verified, Mac, 2026-09-28:** 172 existing unit/integration tests passed + after extracting the common controller, including the real Mac helper's + H.264, ticket, timing and input-echo tests. Eight PC adapter tests passed; + native PC tests were skipped locally because their libraries were absent. +- **Verified, real Frame, 2026-09-28, BUILD_ID 20260925.6191901:** the base + helper's synthetic source created panel `valve.steam.desktopgame.2001639889`, + and the shared Chromium viewer decoded H.264. It recorded 286 frames over + the short probe, with a two-second summary of 19.5 fps shown and total + latency p50/p95 83/156.5 ms. This establishes the existing viewer/transport + route, not Windows/Linux capture, input or a latency target. The probe's + helper, tunnel and viewer were stopped afterward. +- **CI, pending:** Ubuntu x64/ARM64 and Windows native library builds and + x264 test-pattern protocol tests. No Windows or Linux desktop VM was used. +- **Untested:** real Windows WGC → Media Foundation → Frame; real Linux + portal → PipeWire → VA-API/x264 → Frame; physical laser input on either. + No benchmark numbers for those desktop paths are claimed. diff --git a/docs/testing.md b/docs/testing.md index b06577d..1f76458 100644 --- a/docs/testing.md +++ b/docs/testing.md @@ -148,3 +148,14 @@ For example, on 2026-09-27 the smoke test found that Steam's `create-shortcut` refuses ids with a hyphen (`missing/invalid arguments`), which the fake had accepted. The fake now refuses them the same way, and Frame Control makes ids Steam accepts. + +## PC panel host tests + +`tests/test_pcview.py` checks source-bound tickets, reconnect/Stop revocation, +input release, source-coordinate mapping and use of the existing panel +launcher. The `PC host libraries` workflow builds the bundled native adapter +on Windows, Ubuntu x64 and Ubuntu ARM64, then runs a real x264 synthetic +stream through the HTTP/WebSocket agent. `FRAME_PC_REQUIRE_NATIVE=1` makes a +missing native bundle fail CI instead of skipping. These tests do not grant a +portal dialog, capture a desktop, exercise a hardware encoder or validate +laser alignment. See [PC host evidence](pc-in-headset.md#evidence). diff --git a/scripts/macview-bench.py b/scripts/macview-bench.py index cffa8b5..4019d9d 100755 --- a/scripts/macview-bench.py +++ b/scripts/macview-bench.py @@ -55,6 +55,7 @@ from urllib.parse import quote, urlencode # %20, not +: the agent's URLComponen ROOT = Path(__file__).resolve().parent.parent sys.path.insert(0, str(ROOT / "ui")) import frame_macview # noqa: E402 +import frame_pcview # noqa: E402 RESULTS = ROOT / "bench" / "results" PAGES = ROOT / "bench" / "pages" @@ -381,7 +382,10 @@ def run_scenario(args, scenario, agent_port, token, frame_ssh): if args.browser_flag is not None: mv.browser_flags = [f for f in args.browser_flag if f] chrome = None - src = "test" + src = args.source or "test" + cpu = None + if args.pc: + mv.host = "windows" if sys.platform == "win32" else "linux" try: if scenario in ("scroll", "type"): chrome = chrome_window(f"{scenario}.html") @@ -405,9 +409,10 @@ def run_scenario(args, scenario, agent_port, token, frame_ssh): since = max([f["s"] for f in first["frames"]] or [0]) captured_before = first["captured"] start_mac_us = mv.call("/stats", id=stream["id"])["now"] - agent_pid = int(subprocess.run(["pgrep", "-f", "Frame Mac View Lab.app/Contents/MacOS/frame-mac-view"], - capture_output=True, text=True).stdout.split()[0]) - cpu = CpuSampler(frame_ssh, agent_pid) + if not args.pc: + agent_pid = int(subprocess.run(["pgrep", "-f", "Frame Mac View Lab.app/Contents/MacOS/frame-mac-view"], + capture_output=True, text=True).stdout.split()[0]) + cpu = CpuSampler(frame_ssh, agent_pid) if relay: relay.begin() expected = 0 # input events the harness asked the viewer for @@ -443,7 +448,8 @@ def run_scenario(args, scenario, agent_port, token, frame_ssh): if not streams: raise SystemExit(f"{scenario}: the stream ended during the run (did the Frame go to sleep?)") data = streams[0] - cpu.stop() + if cpu: + cpu.stop() if args.raw: Path(f"{args.raw}-{scenario}.json").write_text(json.dumps(data)) result = summarize(scenario, data, start_mac_us, end_mac_us, args, relay, schedule, expected) @@ -452,7 +458,7 @@ def run_scenario(args, scenario, agent_port, token, frame_ssh): for e in data.get("events", []) if e["t"] >= start_mac_us] result["captured"] = (captured_end if captured_end is not None else data["captured"]) - captured_before result["source_fps"] = round(result["captured"] / max(result["duration_s"], 1), 1) - result["cpu"] = cpu.summary() + result["cpu"] = cpu.summary() if cpu else {"host": "not sampled"} result["viewer"] = data["summary"].get("decoder", "") result["show_s"] = show_s result["panel"] = shown.get("panel") @@ -460,6 +466,8 @@ def run_scenario(args, scenario, agent_port, token, frame_ssh): result["route"] = mv.route return result finally: + if cpu: + cpu.stop() try: mv.stop(src) except Exception: # noqa: BLE001 - best effort @@ -670,6 +678,8 @@ def fmt(v): def add_run_args(r): r.add_argument("--scenario", default="test,scroll,type") + r.add_argument("--pc", action="store_true", help="run the bundled PC host; use --scenario test or capture") + r.add_argument("--source", default="", help="PC source: window:ID, display:ID, or choose (Linux portal)") r.add_argument("--duration", type=float, default=20) r.add_argument("--warmup", type=float, default=3) r.add_argument("--quality", default="balanced", choices=list(frame_macview.QUALITY)) @@ -717,7 +727,7 @@ def frame_ssh_for(args): def run_suite(args, frame_ssh, quiet=False): - agent_port, token = lab_agent() + agent_port, token = (args.pc_view.port, args.pc_view.token) if args.pc else lab_agent() results = [] for sc in args.scenario.split(","): if not quiet: @@ -748,7 +758,8 @@ def document(args, frame_ssh, results, **extra): "buffer_ms": args.buffer, "host": args.host or args.frame, "usb": args.usb, "ssh_opts": args.ssh_opt, "encoder": os.environ.get("FRAME_MAC_VIEW_ENCODER", ""), "duration_s": args.duration, "browser_flags": args.browser_flag if args.browser_flag is not None else frame_macview.BROWSER_FLAGS}, - "frame_build": build.strip().partition("=")[2], "mac": platform.mac_ver()[0], "headset": standby[-40:], + "frame_build": build.strip().partition("=")[2], "mac": platform.mac_ver()[0], + "host_platform": platform.platform(), "source": args.source, "pc_host": args.pc, "headset": standby[-40:], "scenarios": results, **extra, } @@ -819,19 +830,46 @@ def main(): print("\nworse:\n " + "\n ".join(regs)) sys.exit(1 if regs else 0) - lab_agent() + if args.pc and (args.cmd == "ab" or args.scenario not in ("test", "capture")): + p.error("PC runs use --scenario test or --scenario capture; Mac automation is not portable") + if args.pc and args.scenario == "capture" and not args.source: + p.error("capture needs --source window:ID, display:ID, or choose") + if not args.pc: + lab_agent() frame_ssh = frame_ssh_for(args) # Keep the Mac's screen awake: virtual displays aren't removed while it sleeps. runs = 1 if args.cmd == "run" else args.repeat * len(args.arm) - awake = subprocess.Popen(["caffeinate", "-d", "-u", "-t", str(int(runs * (args.duration * 4 + 60) + 120))]) + awake = subprocess.Popen(["caffeinate", "-d", "-u", "-t", str(int(runs * (args.duration * 4 + 60) + 120))]) if sys.platform == "darwin" else None + args.pc_view = None try: + if args.pc: + args.pc_view = frame_pcview.PCView(frame_ssh[:-1], ssh_runner(frame_ssh), frame_ssh[-1]) + args.pc_view.ensure_agent() + if args.source == "choose": + args.pc_view.call("/permissions", method="POST") + print("Choose a window or screen in your desktop's sharing dialog.", flush=True) + deadline = time.monotonic()+125 + while time.monotonic() < deadline: + state = args.pc_view.call("/status") + if not state.get("selecting"): + sources = args.pc_view.call("/windows")["windows"] + if not sources: + raise SystemExit(state.get("selectionError") or "Nothing was shared") + args.source = sources[-1]["src"] + break + time.sleep(.5) + else: + raise SystemExit("Sharing dialog timed out") if args.cmd == "ab": arm_runs, table = ab(args, frame_ssh) write(document(args, frame_ssh, [], arms=args.arm, runs=arm_runs, table=table), args) return results = run_suite(args, frame_ssh) finally: - awake.terminate() + if args.pc_view: + args.pc_view.shutdown() + if awake: + awake.terminate() doc = document(args, frame_ssh, results) write(doc, args) bad = [f"{r['scenario']} {k}" for r in results for k, v in r["grades"].items() if v in ("bad", "missing")] diff --git a/tests/test_pcview.py b/tests/test_pcview.py new file mode 100644 index 0000000..6af5734 --- /dev/null +++ b/tests/test_pcview.py @@ -0,0 +1,210 @@ +"""PC adapters: grants, input cleanup, launcher reuse and native codec contract. +Native tests require the bundled libraries; CI sets REQUIRE_NATIVE so a missing +bundle is a failure, never a silent skip. No desktop consent or GPU is faked as +successful real-host capture. +""" +import ctypes +import http.client +import json +import os +from pathlib import Path +import shutil +import struct +import subprocess +import sys +import threading +import time +import unittest +from unittest import mock + +ROOT = Path(__file__).resolve().parent.parent +sys.path.insert(0, str(ROOT / 'ui')) +import frame_pc_agent as agent +import frame_pc_capture as capture +import frame_pcview +import frame_macview +from frame_stream_stats import Stats +from test_macview import WS + + +class GrantsTests(unittest.TestCase): + def test_source_bound_retry_ack_and_stop(self): + g = agent.Grants('master') + ticket = g.ticket('window:1') + self.assertIsNone(g.redeem('window:2', {'t': ticket})) + key = g.redeem('window:1', {'t': ticket}) + self.assertEqual(key, g.redeem('window:1', {'t': ticket})) + g.ack(key) + self.assertIsNone(g.redeem('window:1', {'t': ticket})) + self.assertEqual(key, g.redeem('window:1', {'r': key})) + self.assertIsNone(g.redeem('window:2', {'r': key})) + g.revoke('window:1') + self.assertIsNone(g.redeem('window:1', {'r': key})) + + def test_expiry_and_stop_before_redeem(self): + g = agent.Grants('master') + with mock.patch.object(agent.time, 'monotonic', return_value=1): + ticket = g.ticket('test') + with mock.patch.object(agent.time, 'monotonic', return_value=62): + self.assertIsNone(g.redeem('test', {'t': ticket})) + ticket = g.ticket('test') + g.revoke() + self.assertIsNone(g.redeem('test', {'t': ticket})) + + +class AdapterTests(unittest.TestCase): + def test_windows_uses_wgc_and_media_foundation(self): + for kind in ('window', 'display'): + pipeline = capture.pipeline(dict(src=kind+':42'), 'win32', 'mfh264enc', 640, 360, 30, 1000000) + self.assertIn('capture-api=wgc', pipeline) + self.assertIn('window-handle=42' if kind == 'window' else 'monitor-handle=42', pipeline) + self.assertIn('mfh264enc name=enc low-latency=true bframes=0', pipeline) + self.assertNotIn('drop=true', pipeline) + with self.assertRaises(ValueError): + capture.pipeline(dict(src='window:42 ! fakesink'), 'win32', 'mfh264enc', 640, 360, 30, 1000000) + + def test_linux_uses_only_portal_fd_and_node(self): + p = capture.pipeline(dict(src='window:consented', fd=19, node=47), 'linux', 'x264enc', 640, 360, 30, 1000000) + self.assertIn('pipewiresrc fd=19 path=47', p) + self.assertIn('tune=zerolatency', p) + self.assertEqual(capture.dimensions(1920, 1080, 1280), (1280, 720)) + + def test_portal_input_clamps_and_releases_held_state(self): + native = mock.Mock() + native.lib.fc_portal_input.return_value = 1 + inp = capture.PortalInput(native, dict(portal=123, w=100, h=50)) + inp.handle(dict(t='m', e='down', b=0, x=2, y=-1)) + inp.handle(dict(t='k', e='down', key='Control')) + inp.release() + calls = native.lib.fc_portal_input.call_args_list + self.assertEqual(calls[0].args, (123, 0, 99, 0, 0, 0)) + self.assertIn(mock.call(123, 1, 0, 0, 0x110, 0), calls) + self.assertIn(mock.call(123, 3, 0, 0, 0xffe3, 0), calls) + self.assertFalse(inp.buttons or inp.keys) + + def test_windows_release_and_negative_monitor_origin(self): + host = mock.Mock() + host.user.GetSystemMetrics.side_effect = [-1920, 0, 3840, 1080] + inp = capture.WindowsInput(host, dict(src='display:4', x=-1920, y=0, w=1920, h=1080)) + inp.handle(dict(t='m', e='down', b=2, x=0, y=0)) + self.assertEqual(host.send.call_args_list[0], mock.call(mouse=(0, 0, 0, 0xc001))) + inp.handle(dict(t='k', e='down', code='ControlLeft', key='Control')) + inp.release() + self.assertIn(mock.call(mouse=(0, 0, 0, 0x10)), host.send.call_args_list) + self.assertIn(mock.call(key=(162, 0, 2)), host.send.call_args_list) + + def test_pc_reuses_launcher_and_namespaces_panel_ids(self): + view = frame_pcview.PCView(['ssh'], mock.Mock(return_value='panel created'), 'frame') + view.call = mock.Mock(side_effect=lambda path, **kw: {'screen': True, 'ticket': 'one-use', 'streams': []}) + view.ensure_tunnel = mock.Mock() + view.remote_port = 47900 + view.show('window:12') + self.assertEqual(view.run.call_args.kwargs['stdin'], frame_macview.LAUNCH) + self.assertNotEqual(frame_macview.panel_id('window:12'), frame_macview.panel_id('window:12', 'windows')) + self.assertFalse(view.prefer_usb) + with self.assertRaises(frame_macview.MacViewError): + view.show('separate:12') + + def test_stats_keep_capture_time_and_bound_records(self): + stats = Stats(lambda: 10000000) + stats.input(dict(t='m', i=1, tv=9999990)) + f = stats.add(cap=10000001, arr=10000002, e0=10000003, e1=10000004) + self.assertEqual(f['echo'], 1) + stats.report(dict(t='rx', s=f['s'], r=10000010)) + stats.report(dict(t='fd', f=[[f['s'], 10000011, 10000012, 10000013]])) + self.assertEqual(stats.snapshot(settle=0)['inputs'][0]['frame'], f['s']) + for _ in range(4100): + stats.add(cap=10000001) + self.assertEqual(len(stats.frames), 4096) + + +class NativeTests(unittest.TestCase): + @classmethod + def setUpClass(cls): + if not capture.LIBRARY.exists(): + if os.environ.get('FRAME_PC_REQUIRE_NATIVE') == '1': + raise AssertionError('CI required native PC libraries but they were not built') + raise unittest.SkipTest('native PC libraries not built on this host') + cls.token = 'test-' + os.urandom(8).hex() + env = frame_pcview.PCView([], lambda *a: None, 'frame').agent_environment() + env['FRAME_MAC_VIEW_TOKEN'] = cls.token + cls.proc = subprocess.Popen([sys.executable, str(ROOT/'ui/frame_pc_agent.py'), 'serve', '--port', '0', + '--page', str(ROOT/'ui/mac-view.html'), '--exit-on-eof'], stdin=subprocess.PIPE, stdout=subprocess.PIPE, + stderr=subprocess.PIPE, text=True, env=env) + line = cls.proc.stdout.readline() + if 'listening' not in line: + raise AssertionError('Native agent failed: ' + cls.proc.stderr.read()) + cls.port = int(line.rsplit(':', 1)[1]) + + @classmethod + def tearDownClass(cls): + cls.proc.stdin.close() + try: + rc = cls.proc.wait(10) + except subprocess.TimeoutExpired: + cls.proc.kill() + raise AssertionError('Native agent did not stop after EOF') + error = cls.proc.stderr.read() + cls.proc.stdout.close() + cls.proc.stderr.close() + if rc: + raise AssertionError('Native agent exited %s: %s' % (rc, error)) + + def request(self, path, method='GET'): + conn = http.client.HTTPConnection('127.0.0.1', self.port, timeout=10) + try: + conn.request(method, path) + response = conn.getresponse() + return response.status, response.read() + finally: + conn.close() + + def tearDown(self): + self.request('/close?k='+self.token, 'POST') + + def test_auth_and_real_h264_stats(self): + self.assertEqual(self.request('/status')[0], 403) + status, body = self.request('/ticket?src=test&k='+self.token, 'POST') + self.assertEqual(status, 200) + ticket = json.loads(body)['ticket'] + ws = WS(self.port, '/stream?src=test&fps=30&max=640&t='+ticket) + try: + self.assertEqual(ws.status, 101) + key, got, echo = None, 0, False + ws.send_text(json.dumps(dict(t='m', e='down', i=7, tv=1, x=.2, y=.3))) + deadline = time.monotonic()+10 + while time.monotonic() < deadline and (got < 5 or not echo): + op, data = ws.recv() + if op == 1: + m = json.loads(data) + self.assertNotEqual(m.get('t'), 'error', m) + if m.get('t') == 'hello': + key = m['r'] + ws.send_text(json.dumps(dict(t='ack'))) + elif op == 2: + got += 1 + seq, echoed = struct.unpack('!II', data[9:17]) + now = struct.unpack('!Q', data[1:9])[0] + echo |= echoed == 7 + self.assertIn(b'\x00\x00\x00\x01', data[17:]) + ws.send_text(json.dumps(dict(t='rx', s=seq, r=now+100))) + ws.send_text(json.dumps(dict(t='fd', f=[[seq, now+200, now+300, now+400]]))) + self.assertGreaterEqual(got, 5) + self.assertTrue(echo) + time.sleep(.1) + stats = json.loads(self.request('/stats?settle=0&k='+self.token)[1])['streams'][0] + self.assertGreater(stats['frames'][0]['e1'], stats['frames'][0]['e0']) + self.assertIn('controller', stats) + used = WS(self.port, '/stream?src=test&t='+ticket) + self.assertEqual(used.status, 403) + used.close() + self.request('/close?src=test&k='+self.token, 'POST') + revoked = WS(self.port, '/stream?src=test&r='+key) + self.assertEqual(revoked.status, 403) + revoked.close() + finally: + ws.close() + + +if __name__ == '__main__': + unittest.main() diff --git a/ui/frame_macview.py b/ui/frame_macview.py index b89cf88..ee95686 100644 --- a/ui/frame_macview.py +++ b/ui/frame_macview.py @@ -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} diff --git a/ui/frame_pc_agent.py b/ui/frame_pc_agent.py new file mode 100644 index 0000000..803c51e --- /dev/null +++ b/ui/frame_pc_agent.py @@ -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() diff --git a/ui/frame_pc_capture.py b/ui/frame_pc_capture.py new file mode 100644 index 0000000..a4a0bb3 --- /dev/null +++ b/ui/frame_pc_capture.py @@ -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)) diff --git a/ui/frame_pcview.py b/ui/frame_pcview.py new file mode 100644 index 0000000..317b61f --- /dev/null +++ b/ui/frame_pcview.py @@ -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) diff --git a/ui/frame_stream_stats.py b/ui/frame_stream_stats.py new file mode 100644 index 0000000..513f40b --- /dev/null +++ b/ui/frame_stream_stats.py @@ -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()) diff --git a/ui/index.html b/ui/index.html index d9491e4..b7a9b62 100644 --- a/ui/index.html +++ b/ui/index.html @@ -651,7 +651,7 @@