44 changed files with 1030 additions and 5167 deletions
Generated
+674 -119
View File
File diff suppressed because it is too large Load Diff
+4 -27
View File
@@ -34,32 +34,9 @@ ureq = { version = "3", default-features = false, features = ["rustls"] }
toml = "1"
chrono = { version = "0.4", default-features = false, features = ["clock", "serde"] }
eframe = { version = "0.34.2", default-features = false, features = ["glow", "default_fonts", "wayland", "x11"], optional = true }
# Desktop notifications on viewer join/leave. Default features give the
# pure-Rust zbus backend (no system libdbus, no image crate).
notify-rust = { version = "4", optional = true }
# System-tray icon (StatusNotifierItem over D-Bus). Pure-Rust, riding the same
# zbus stack notify-rust already pulls — no GTK, no libappindicator/C libdbus.
ksni = { version = "0.3", optional = true }
# Hand-rolled windowing stack for the GUI (replaces eframe::run_native) so we
# can drop the OS window on "hide to tray" — the only way to truly hide a
# toplevel on Wayland — and recreate it on Show. All of these are already pulled
# in transitively by eframe; making them direct adds no new crates to vet.
# eframe is kept for its egui re-export + icon_data PNG decoder. egui_glow needs
# its (non-default) `winit` feature for the `EguiGlow` integration type; eframe
# pulls egui_glow but without that feature, so we enable it here.
egui_glow = { version = "0.34.2", default-features = false, features = ["winit", "wayland", "x11"], optional = true }
# winit's default set minus `wayland-csd-adwaita`: KWin (and most desktop
# compositors) draw server-side decorations, and eframe never enabled CSD
# either, so dropping it keeps the dependency tree identical to before (no
# sctk-adwaita / tiny-skia / ttf-parser pulled in just for a fallback titlebar).
winit = { version = "0.30", default-features = false, features = ["rwh_06", "x11", "wayland", "wayland-dlopen"], optional = true }
glutin = { version = "0.32", optional = true }
glutin-winit = { version = "0.5", optional = true }
# QR-encode the host ticket so a phone (or a second laptop with a webcam) can
# pick it up without typing 140 chars. default-features = false to skip the
# `image` crate dep tree — we render the modules to an `egui::ColorImage`
# directly.
qrcode = { version = "0.14", default-features = false, optional = true }
tray-icon = { version = "0.24.0", optional = true }
notify-rust = { version = "4.17.0", optional = true }
gtk = { version = "0.18.2", optional = true }
[profile.release]
lto = "thin"
@@ -69,4 +46,4 @@ strip = "symbols"
[features]
# Opt-in graphical front-end (pixelpass --gui). Default-off so the headless
# build never pulls the GUI toolkit tree.
gui = ["dep:eframe", "dep:notify-rust", "dep:ksni", "dep:egui_glow", "dep:winit", "dep:glutin", "dep:glutin-winit", "dep:qrcode"]
gui = ["dep:eframe", "dep:tray-icon", "dep:notify-rust", "dep:gtk"]
-202
View File
@@ -1,202 +0,0 @@
Apache License
Version 2.0, January 2004
http://www.apache.org/licenses/
TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION
1. Definitions.
"License" shall mean the terms and conditions for use, reproduction,
and distribution as defined by Sections 1 through 9 of this document.
"Licensor" shall mean the copyright owner or entity authorized by
the copyright owner that is granting the License.
"Legal Entity" shall mean the union of the acting entity and all
other entities that control, are controlled by, or are under common
control with that entity. For the purposes of this definition,
"control" means (i) the power, direct or indirect, to cause the
direction or management of such entity, whether by contract or
otherwise, or (ii) ownership of fifty percent (50%) or more of the
outstanding shares, or (iii) beneficial ownership of such entity.
"You" (or "Your") shall mean an individual or Legal Entity
exercising permissions granted by this License.
"Source" form shall mean the preferred form for making modifications,
including but not limited to software source code, documentation
source, and configuration files.
"Object" form shall mean any form resulting from mechanical
transformation or translation of a Source form, including but
not limited to compiled object code, generated documentation,
and conversions to other media types.
"Work" shall mean the work of authorship, whether in Source or
Object form, made available under the License, as indicated by a
copyright notice that is included in or attached to the work
(an example is provided in the Appendix below).
"Derivative Works" shall mean any work, whether in Source or Object
form, that is based on (or derived from) the Work and for which the
editorial revisions, annotations, elaborations, or other modifications
represent, as a whole, an original work of authorship. For the purposes
of this License, Derivative Works shall not include works that remain
separable from, or merely link (or bind by name) to the interfaces of,
the Work and Derivative Works thereof.
"Contribution" shall mean any work of authorship, including
the original version of the Work and any modifications or additions
to that Work or Derivative Works thereof, that is intentionally
submitted to Licensor for inclusion in the Work by the copyright owner
or by an individual or Legal Entity authorized to submit on behalf of
the copyright owner. For the purposes of this definition, "submitted"
means any form of electronic, verbal, or written communication sent
to the Licensor or its representatives, including but not limited to
communication on electronic mailing lists, source code control systems,
and issue tracking systems that are managed by, or on behalf of, the
Licensor for the purpose of discussing and improving the Work, but
excluding communication that is conspicuously marked or otherwise
designated in writing by the copyright owner as "Not a Contribution."
"Contributor" shall mean Licensor and any individual or Legal Entity
on behalf of whom a Contribution has been received by Licensor and
subsequently incorporated within the Work.
2. Grant of Copyright License. Subject to the terms and conditions of
this License, each Contributor hereby grants to You a perpetual,
worldwide, non-exclusive, no-charge, royalty-free, irrevocable
copyright license to reproduce, prepare Derivative Works of,
publicly display, publicly perform, sublicense, and distribute the
Work and such Derivative Works in Source or Object form.
3. Grant of Patent License. Subject to the terms and conditions of
this License, each Contributor hereby grants to You a perpetual,
worldwide, non-exclusive, no-charge, royalty-free, irrevocable
(except as stated in this section) patent license to make, have made,
use, offer to sell, sell, import, and otherwise transfer the Work,
where such license applies only to those patent claims licensable
by such Contributor that are necessarily infringed by their
Contribution(s) alone or by combination of their Contribution(s)
with the Work to which such Contribution(s) was submitted. If You
institute patent litigation against any entity (including a
cross-claim or counterclaim in a lawsuit) alleging that the Work
or a Contribution incorporated within the Work constitutes direct
or contributory patent infringement, then any patent licenses
granted to You under this License for that Work shall terminate
as of the date such litigation is filed.
4. Redistribution. You may reproduce and distribute copies of the
Work or Derivative Works thereof in any medium, with or without
modifications, and in Source or Object form, provided that You
meet the following conditions:
(a) You must give any other recipients of the Work or
Derivative Works a copy of this License; and
(b) You must cause any modified files to carry prominent notices
stating that You changed the files; and
(c) You must retain, in the Source form of any Derivative Works
that You distribute, all copyright, patent, trademark, and
attribution notices from the Source form of the Work,
excluding those notices that do not pertain to any part of
the Derivative Works; and
(d) If the Work includes a "NOTICE" text file as part of its
distribution, then any Derivative Works that You distribute must
include a readable copy of the attribution notices contained
within such NOTICE file, excluding those notices that do not
pertain to any part of the Derivative Works, in at least one
of the following places: within a NOTICE text file distributed
as part of the Derivative Works; within the Source form or
documentation, if provided along with the Derivative Works; or,
within a display generated by the Derivative Works, if and
wherever such third-party notices normally appear. The contents
of the NOTICE file are for informational purposes only and
do not modify the License. You may add Your own attribution
notices within Derivative Works that You distribute, alongside
or as an addendum to the NOTICE text from the Work, provided
that such additional attribution notices cannot be construed
as modifying the License.
You may add Your own copyright statement to Your modifications and
may provide additional or different license terms and conditions
for use, reproduction, or distribution of Your modifications, or
for any such Derivative Works as a whole, provided Your use,
reproduction, and distribution of the Work otherwise complies with
the conditions stated in this License.
5. Submission of Contributions. Unless You explicitly state otherwise,
any Contribution intentionally submitted for inclusion in the Work
by You to the Licensor shall be under the terms and conditions of
this License, without any additional terms or conditions.
Notwithstanding the above, nothing herein shall supersede or modify
the terms of any separate license agreement you may have executed
with Licensor regarding such Contributions.
6. Trademarks. This License does not grant permission to use the trade
names, trademarks, service marks, or product names of the Licensor,
except as required for reasonable and customary use in describing the
origin of the Work and reproducing the content of the NOTICE file.
7. Disclaimer of Warranty. Unless required by applicable law or
agreed to in writing, Licensor provides the Work (and each
Contributor provides its Contributions) on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or
implied, including, without limitation, any warranties or conditions
of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A
PARTICULAR PURPOSE. You are solely responsible for determining the
appropriateness of using or redistributing the Work and assume any
risks associated with Your exercise of permissions under this License.
8. Limitation of Liability. In no event and under no legal theory,
whether in tort (including negligence), contract, or otherwise,
unless required by applicable law (such as deliberate and grossly
negligent acts) or agreed to in writing, shall any Contributor be
liable to You for damages, including any direct, indirect, special,
incidental, or consequential damages of any character arising as a
result of this License or out of the use or inability to use the
Work (including but not limited to damages for loss of goodwill,
work stoppage, computer failure or malfunction, or any and all
other commercial damages or losses), even if such Contributor
has been advised of the possibility of such damages.
9. Accepting Warranty or Additional Liability. While redistributing
the Work or Derivative Works thereof, You may choose to offer,
and charge a fee for, acceptance of support, warranty, indemnity,
or other liability obligations and/or rights consistent with this
License. However, in accepting such obligations, You may act only
on Your own behalf and on Your sole responsibility, not on behalf
of any other Contributor, and only if You agree to indemnify,
defend, and hold each Contributor harmless for any liability
incurred by, or claims asserted against, such Contributor by reason
of your accepting any such warranty or additional liability.
END OF TERMS AND CONDITIONS
APPENDIX: How to apply the Apache License to your work.
To apply the Apache License to your work, attach the following
boilerplate notice, with the fields enclosed by brackets "[]"
replaced with your own identifying information. (Don't include
the brackets!) The text should be enclosed in the appropriate
comment syntax for the file format. We also recommend that a
file or class name and description of purpose be included on the
same "printed page" as the copyright notice for easier
identification within third-party archives.
Copyright 2026 mollusk
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at
http://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
-21
View File
@@ -1,21 +0,0 @@
MIT License
Copyright (c) 2026 mollusk
Permission is hereby granted, free of charge, to any person obtaining a copy
of this software and associated documentation files (the "Software"), to deal
in the Software without restriction, including without limitation the rights
to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
copies of the Software, and to permit persons to whom the Software is
furnished to do so, subject to the following conditions:
The above copyright notice and this permission notice shall be included in all
copies or substantial portions of the Software.
THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
SOFTWARE.
+2 -68
View File
@@ -83,15 +83,8 @@ pixelpass --gui
```
Host: pick quality / max-viewers / options, click **Start hosting**, and the
share code appears with a copy button. Connected viewers are listed with a
**Kick** button each, and a desktop notification fires as they join or leave.
View: paste a code, pick mpv or VLC, click **Connect** and the player launches.
A system-tray icon shows current status. **Settings → "Keep running in the
tray when I close the window"** (off by default) makes the close button hide
the window — truly, by dropping it — while any active stream keeps running in
the child; reopen it from the tray. (Plain close still quits when the option is
off, or when no system tray is present.)
share code appears with a copy button alongside a live viewer count. View:
paste a code, pick mpv or VLC, click **Connect** and the player launches.
The window is a thin driver — it runs the same headless `pixelpass` as a
child process and reads its event stream, so the GUI is purely additive and
@@ -223,65 +216,6 @@ measured_at = "2026-05-21T20:41:16Z"
- Skip is sticky — once you skip the test, pixelpass won't ask again
unless you reconfigure.
## Relay
By default pixelpass uses iroh's bundled relay servers to coordinate the
P2P connection (peers still hole-punch a direct UDP path when they can; the
relay is the fallback and the rendezvous point). You can point it at a
different relay — a self-hosted one, or n0's staging/production servers —
with either:
```bash
pixelpass --relay https://relay.example/ # host or viewer
PIXELPASS_RELAY=https://relay.example/ pixelpass … # env-var form
```
The flag applies to both host and viewer and takes precedence over the
environment variable. The env-var form is handy for the `--gui` front-end,
since the GUI's child host/viewer processes inherit it; the `--gui --relay`
flag form is forwarded to them too. Both ends must use the same relay to
find each other.
## Themes
The `--gui` front-end ships three colour themes — **Default Dark**,
**Catppuccin Mocha**, and **Catppuccin Latte** — and you can add your own.
Pick one under **Settings → Appearance**; the choice is remembered.
A theme is a small TOML file of named colours:
```toml
name = "My Theme"
dark = true # base egui defaults to start from (dark or light)
window_bg = "#1b1b1f" # window background
panel_bg = "#242429" # panels / frames
input_bg = "#141417" # text fields, the ticket box
text = "#e6e6ea" # primary text
weak_text = "#a0a0a8" # hints, secondary text
accent = "#5aa0f2" # selection, links, the active control
button_bg = "#33333a" # buttons at rest
button_hovered = "#44444d"
streaming = "#6fdc8c" # "● Streaming"
waiting = "#f2c14e" # "● Waiting for viewers…"
success = "#6fdc8c" # "✓ Copied", valid-code confirmation
warning = "#f0a85a" # non-fatal warnings
error = "#f2756f" # errors
```
Colours are `#rrggbb` hex strings. Any field you leave out falls back to
Default Dark, so partial files are fine.
Two ways to make one:
- **In the app:** Settings → Appearance → *Edit / create a theme* gives you a
colour picker per field with a live preview, and **Save** writes a `.toml`.
- **By hand:** drop a `.toml` into `~/.config/pixelpass/themes/` (the XDG
config dir). It appears in the picker next time you open Settings.
Sharing a theme is just sending someone the file. A user theme whose `name`
matches a built-in overrides that built-in.
## Audio
By default pixelpass captures the default sink's monitor — the viewer
-93
View File
@@ -1,93 +0,0 @@
Copyright 2022 The Noto Project Authors (https://github.com/notofonts/latin-greek-cyrillic)
This Font Software is licensed under the SIL Open Font License, Version 1.1.
This license is copied below, and is also available with a FAQ at:
https://scripts.sil.org/OFL
-----------------------------------------------------------
SIL OPEN FONT LICENSE Version 1.1 - 26 February 2007
-----------------------------------------------------------
PREAMBLE
The goals of the Open Font License (OFL) are to stimulate worldwide
development of collaborative font projects, to support the font creation
efforts of academic and linguistic communities, and to provide a free and
open framework in which fonts may be shared and improved in partnership
with others.
The OFL allows the licensed fonts to be used, studied, modified and
redistributed freely as long as they are not sold by themselves. The
fonts, including any derivative works, can be bundled, embedded,
redistributed and/or sold with any software provided that any reserved
names are not used by derivative works. The fonts and derivatives,
however, cannot be released under any other type of license. The
requirement for fonts to remain under this license does not apply
to any document created using the fonts or their derivatives.
DEFINITIONS
"Font Software" refers to the set of files released by the Copyright
Holder(s) under this license and clearly marked as such. This may
include source files, build scripts and documentation.
"Reserved Font Name" refers to any names specified as such after the
copyright statement(s).
"Original Version" refers to the collection of Font Software components as
distributed by the Copyright Holder(s).
"Modified Version" refers to any derivative made by adding to, deleting,
or substituting -- in part or in whole -- any of the components of the
Original Version, by changing formats or by porting the Font Software to a
new environment.
"Author" refers to any designer, engineer, programmer, technical
writer or other person who contributed to the Font Software.
PERMISSION & CONDITIONS
Permission is hereby granted, free of charge, to any person obtaining
a copy of the Font Software, to use, study, copy, merge, embed, modify,
redistribute, and sell modified and unmodified copies of the Font
Software, subject to the following conditions:
1) Neither the Font Software nor any of its individual components,
in Original or Modified Versions, may be sold by itself.
2) Original or Modified Versions of the Font Software may be bundled,
redistributed and/or sold with any software, provided that each copy
contains the above copyright notice and this license. These can be
included either as stand-alone text files, human-readable headers or
in the appropriate machine-readable metadata fields within text or
binary files as long as those fields can be easily viewed by the user.
3) No Modified Version of the Font Software may use the Reserved Font
Name(s) unless explicit written permission is granted by the corresponding
Copyright Holder. This restriction only applies to the primary font name as
presented to the users.
4) The name(s) of the Copyright Holder(s) or the Author(s) of the Font
Software shall not be used to promote, endorse or advertise any
Modified Version, except to acknowledge the contribution(s) of the
Copyright Holder(s) and the Author(s) or with their explicit written
permission.
5) The Font Software, modified or unmodified, in part or in whole,
must be distributed entirely under this license, and must not be
distributed under any other license. The requirement for fonts to
remain under this license does not apply to any document created
using the Font Software.
TERMINATION
This license becomes null and void if any of the above conditions are
not met.
DISCLAIMER
THE FONT SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND,
EXPRESS OR IMPLIED, INCLUDING BUT NOT LIMITED TO ANY WARRANTIES OF
MERCHANTABILITY, FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT
OF COPYRIGHT, PATENT, TRADEMARK, OR OTHER RIGHT. IN NO EVENT SHALL THE
COPYRIGHT HOLDER BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY,
INCLUDING ANY GENERAL, SPECIAL, INDIRECT, INCIDENTAL, OR CONSEQUENTIAL
DAMAGES, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING
FROM, OUT OF THE USE OR INABILITY TO USE THE FONT SOFTWARE OR FROM
OTHER DEALINGS IN THE FONT SOFTWARE.
Binary file not shown.
Binary file not shown.

Before

Width:  |  Height:  |  Size: 7.9 KiB

-12
View File
@@ -1,12 +0,0 @@
[Desktop Entry]
Type=Application
Name=pixelpass
GenericName=Screen Sharing
Comment=P2P screen sharing over iroh — no port forwarding, no signup
Exec=pixelpass --gui
Icon=pixelpass
StartupWMClass=pixelpass
Terminal=false
Categories=Network;RemoteAccess;
Keywords=screen;share;sharing;remote;p2p;iroh;cast;
StartupNotify=true
-16
View File
@@ -1,16 +0,0 @@
<svg xmlns="http://www.w3.org/2000/svg" width="256" height="256" viewBox="0 0 256 256">
<title>pixelpass</title>
<defs>
<linearGradient id="bg" x1="0" y1="0" x2="256" y2="256" gradientUnits="userSpaceOnUse">
<stop offset="0" stop-color="#4338ca"/>
<stop offset="1" stop-color="#7c3aed"/>
</linearGradient>
</defs>
<rect x="8" y="8" width="240" height="240" rx="56" fill="url(#bg)"/>
<!-- pixel stream: squares fading cyan -> white, "passed" toward the arrow -->
<rect x="50.25" y="181.29" width="16" height="16" rx="3.2" fill="#2dd5ef"/>
<rect x="67" y="124.44" width="20" height="20" rx="4" fill="#62dff2"/>
<rect x="96.25" y="82.09" width="24" height="24" rx="4.8" fill="#98e8f6"/>
<rect x="134.04" y="55.77" width="28" height="28" rx="5.6" fill="#c9f1f9"/>
<path d="M 225.5 57.5 L 183.1 85.8 L 176.9 42.2 Z" fill="#f8fafc"/>
</svg>

Before

Width:  |  Height:  |  Size: 875 B

-3
View File
@@ -1,3 +0,0 @@
.tools/
AppDir/
*.AppImage
-20
View File
@@ -1,20 +0,0 @@
#!/bin/sh
# AppRun for the PixelPass AppImage.
#
# PixelPass is an orchestrator: it shells out to gst-launch-1.0, pactl, and a
# player (mpv/vlc) found on the host PATH. We prepend our own usr/bin so any
# bundled helpers win, but the host's tools remain reachable — that's why this
# app suits AppImage (no sandbox) better than a Flatpak.
HERE="$(dirname "$(readlink -f "$0")")"
export PATH="$HERE/usr/bin:$PATH"
BIN="$HERE/usr/bin/pixelpass"
# With no arguments and no controlling terminal — i.e. launched from a file
# manager or the .desktop entry — open the GUI. From a terminal, or with any
# argument (a ticket, --host, --gui, --repair, …), pass through so the CLI and
# the interactive menu both work.
if [ "$#" -eq 0 ] && [ ! -t 0 ]; then
exec "$BIN" --gui
fi
exec "$BIN" "$@"
-72
View File
@@ -1,72 +0,0 @@
# PixelPass AppImage
A "thin" AppImage: the gui-enabled `pixelpass` binary, a launcher (`AppRun`),
and the desktop entry + icon. Run `./build-appimage.sh` to produce
`pixelpass-<version>-x86_64.AppImage`.
## Why thin
`pixelpass` is an orchestrator — it links almost nothing (only `libpipewire`,
which is excludelisted because it must match the host daemon) and instead
**shells out** to `gst-launch-1.0`, `pactl`, and a player (`mpv`/`vlc`) found on
the host `PATH`. The GUI's graphics libraries (`libGL`, `libwayland-*`,
`libxkbcommon`, X11) are dlopen'd at runtime and are likewise on the AppImage
excludelist — every desktop already has a matching set. So there is nothing
useful to bundle, and bundling the graphics stack would only risk driver
mismatches. The AppImage therefore carries just the binary.
This also explains why PixelPass suits AppImage better than Flatpak: the
no-sandbox model lets the bundled binary freely spawn the host's `gst-launch`,
`pactl`, and player, which a Flatpak sandbox would block.
## Host requirements
The AppImage runs on any reasonably current glibc-based distro that has:
- **GStreamer + plugins** — `gst-launch-1.0`/`gst-inspect-1.0` plus base,
good/bad/ugly, libav, and the PipeWire plugin (the binary tells you the exact
package names for your distro if something is missing).
- **PipeWire** (with the PulseAudio shim, for `pactl`).
- **A player** — `mpv` (preferred) or `vlc` — for the viewer side.
- For X11 single-window capture: `xwininfo`.
These are the same dependencies the Arch package lists; the AppImage just spares
you the Rust toolchain.
## Building for broad compatibility (lower glibc baseline)
An AppImage requires a host glibc **at least as new** as the build host's. Built
straight on a rolling distro (e.g. CachyOS, glibc 2.43) the AppImage only runs
on equally-new systems. Build inside an older base for wider reach. The script
honours `CARGO_TARGET_DIR`, so an isolated toolchain won't clobber your host's
`target/`:
```sh
# One-time: an Ubuntu 24.04 distrobox (docker or podman backend).
distrobox create --yes --image ubuntu:24.04 --name pixelpass-build
distrobox enter pixelpass-build -- sudo apt-get update
distrobox enter pixelpass-build -- sudo apt-get install -y \
build-essential cmake clang libclang-dev pkg-config \
libpipewire-0.3-dev libspa-0.2-dev curl ca-certificates file
# Install rustup inside the box (edition 2024 needs rustc >= 1.85), then:
distrobox enter pixelpass-build -- env \
CARGO_TARGET_DIR=~/.cache/pixelpass-ubuntu/target \
./packaging/appimage/build-appimage.sh
```
**Why Ubuntu 24.04 and not something older:** PixelPass's `pipewire` crate
binds the system's PipeWire headers via bindgen, and anything older than ~PW 1.0
(e.g. Ubuntu 22.04's 0.3.48) fails to compile (missing struct fields / wrong
types). And since PixelPass *is* a PipeWire/portal/Wayland app, it can only run
on distros new enough to have modern PipeWire anyway — so an ancient glibc base
buys nothing. 24.04 (glibc 2.39, PW 1.0.5) is the sweet spot.
The 24.04-built binary's baseline is **glibc 2.39** — and the only 2.39 symbols
are two *weak* `pidfd_*` references from Rust std's process spawning (everything
else is ≤ 2.35). That covers Ubuntu 24.04+, Debian 13+, Fedora 40+, and current
rolling distros.
## Caveats
- **Hardware encode (VAAPI `vah264enc`)** uses the host GPU driver; it can't be
bundled. The software path (`--no-hwencode`, x264) always works.
-59
View File
@@ -1,59 +0,0 @@
#!/usr/bin/env bash
# Build a "thin" PixelPass AppImage: the gui-enabled release binary plus only
# its non-excludelisted shared libraries. The graphics stack (libGL, wayland,
# xkbcommon, X11) is intentionally left to the host — those libs are on the
# AppImage excludelist because they must match the host driver — and the
# runtime tools PixelPass shells out to (gst-launch-1.0, pactl, mpv/vlc) are
# expected on the host PATH, the same contract the Arch package documents.
#
# Usage: packaging/appimage/build-appimage.sh
# Output: packaging/appimage/pixelpass-x86_64.AppImage
set -euo pipefail
here="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
repo="$(cd "$here/../.." && pwd)"
tools="$here/.tools"
appdir="$here/AppDir"
mkdir -p "$tools"
# linuxdeploy is itself an AppImage; run it without FUSE so this works on hosts
# (and CI) that lack libfuse2.
export APPIMAGE_EXTRACT_AND_RUN=1
# Embed the version from Cargo.toml into the AppImage filename metadata.
VERSION="$(grep -m1 '^version' "$repo/Cargo.toml" | sed -E 's/.*"(.*)".*/\1/')"
export VERSION
echo ">> building release binary (--features gui)"
( cd "$repo" && cargo build --release --features gui )
# Honour CARGO_TARGET_DIR so an isolated build (e.g. inside an old-glibc
# distrobox) doesn't have to clobber the host's target/.
bin="${CARGO_TARGET_DIR:-$repo/target}/release/pixelpass"
echo ">> fetching linuxdeploy"
ld="$tools/linuxdeploy-x86_64.AppImage"
if [ ! -x "$ld" ]; then
curl -fL --retry 3 -o "$ld" \
"https://github.com/linuxdeploy/linuxdeploy/releases/download/continuous/linuxdeploy-x86_64.AppImage"
chmod +x "$ld"
fi
echo ">> assembling AppDir"
rm -rf "$appdir"
mkdir -p "$appdir/usr/bin"
install -m755 "$bin" "$appdir/usr/bin/pixelpass"
echo ">> running linuxdeploy (bundles libs, builds the AppImage)"
# -e: analyse this binary for libraries to bundle (only libpipewire et al. that
# aren't excludelisted will be copied; glibc + graphics libs are skipped).
# -d/-i: desktop entry + icon for desktop integration.
# --custom-apprun: our launcher that opens --gui from a file manager.
( cd "$here" && OUTPUT="pixelpass-${VERSION}-x86_64.AppImage" "$ld" \
--appdir "$appdir" \
-e "$bin" \
-d "$repo/assets/pixelpass.desktop" \
-i "$repo/assets/pixelpass-256.png" \
--icon-filename pixelpass \
--custom-apprun "$here/AppRun" \
--output appimage )
echo ">> done: $here/pixelpass-${VERSION}-x86_64.AppImage"
-6
View File
@@ -1,6 +0,0 @@
# makepkg build artifacts
src/
pkg/
/pixelpass/
*.pkg.tar.*
*.log
-65
View File
@@ -1,65 +0,0 @@
# Maintainer: mollusk <jitty+lc1iz0dc@protonmail.com>
#
# Local versioned package, built from the local git repo on `main`.
# For a tagged release, switch the source fragment to `#tag=v0.1.0`.
pkgname=pixelpass
pkgver=0.1.0
pkgrel=1
pkgdesc='P2P screen sharing over iroh — no port forwarding, no signup'
arch=('x86_64')
url='file:///home/mollusk/git/butter/pixelpass'
license=('MIT' 'Apache-2.0' 'OFL-1.1')
depends=(
'gstreamer' # gst-launch-1.0 / gst-inspect-1.0
'gst-plugins-base' # videoscale (quality-preset downscale)
'gst-plugins-good' # ximagesrc (X11 capture) + pulsesrc
'gst-plugins-bad' # h264parse, mpegtsmux, aacparse
'gst-libav' # avenc_aac (audio encode)
'gst-plugin-va' # vah264enc (default hardware H.264 encoder)
'libpulse' # pactl (audio routing / device control)
'hicolor-icon-theme' # owns the scalable icon dir
'libglvnd' # libGL for the egui (glow) GUI
'libxkbcommon' # GUI keyboard handling (winit)
'wayland' # GUI Wayland backend libs
)
optdepends=(
'mpv: recommended stream viewer (the GUI launches mpv)'
'vlc: alternative stream viewer'
'gst-plugins-ugly: software x264 encoding for `pixelpass --no-hwencode`'
'gst-plugin-pipewire: screen capture on Wayland sessions'
'xorg-xwininfo: share a single window on X11 (`pixelpass --window`)'
)
makedepends=('cargo' 'git')
options=('!lto')
_branch='main'
source=("$pkgname::git+file:///home/mollusk/git/butter/pixelpass#branch=$_branch")
sha256sums=('SKIP')
prepare() {
cd "$srcdir/$pkgname"
export RUSTUP_TOOLCHAIN=stable
cargo fetch --locked --target "$(rustc -vV | sed -n 's/host: //p')"
}
build() {
cd "$srcdir/$pkgname"
export RUSTUP_TOOLCHAIN=stable
export CARGO_TARGET_DIR=target
# --features gui so the .desktop launcher (pixelpass --gui) works.
cargo build --frozen --release --features gui
}
package() {
cd "$srcdir/$pkgname"
install -Dm0755 "target/release/$pkgname" "$pkgdir/usr/bin/$pkgname"
install -Dm0644 assets/pixelpass.desktop \
"$pkgdir/usr/share/applications/$pkgname.desktop"
install -Dm0644 assets/pixelpass.svg \
"$pkgdir/usr/share/icons/hicolor/scalable/apps/$pkgname.svg"
install -Dm0644 README.md "$pkgdir/usr/share/doc/$pkgname/README.md"
install -Dm0644 LICENSE-MIT "$pkgdir/usr/share/licenses/$pkgname/LICENSE-MIT"
install -Dm0644 LICENSE-APACHE "$pkgdir/usr/share/licenses/$pkgname/LICENSE-APACHE"
install -Dm0644 assets/NotoSans-OFL.txt \
"$pkgdir/usr/share/licenses/$pkgname/NotoSans-OFL.txt"
}
+1 -17
View File
@@ -68,13 +68,6 @@ pub struct Cli {
pub port: u16,
// ── global ────────────────────────────────────────────────────────
/// Relay server URL to use instead of the bundled defaults, e.g.
/// `https://relay.example/`. Applies to both host and viewer. Falls back
/// to the `PIXELPASS_RELAY` environment variable. Use this to get off the
/// pre-release default relays or to point at a self-hosted relay.
#[arg(long, value_name = "URL")]
pub relay: Option<String>,
/// Launch the graphical front-end (a window with Host/View controls)
/// instead of the terminal menu. Requires a build with `--features gui`.
#[arg(long)]
@@ -147,16 +140,12 @@ pub struct HostOpts {
pub no_hwencode: bool,
pub max_viewers: Option<u32>,
pub interactive: bool,
/// Relay override (resolved from `--relay` / `PIXELPASS_RELAY`); None = defaults.
pub relay: Option<String>,
}
#[derive(Debug, Clone)]
pub struct ViewerOpts {
pub port: u16,
pub interactive: bool,
/// Relay override (resolved from `--relay` / `PIXELPASS_RELAY`); None = defaults.
pub relay: Option<String>,
}
impl Cli {
@@ -174,15 +163,10 @@ impl Cli {
no_hwencode: self.no_hwencode,
max_viewers: self.max_viewers,
interactive,
relay: crate::common::endpoint::relay_override(self.relay.as_deref()),
}
}
pub fn into_viewer_opts(self, interactive: bool) -> ViewerOpts {
ViewerOpts {
port: self.port,
interactive,
relay: crate::common::endpoint::relay_override(self.relay.as_deref()),
}
ViewerOpts { port: self.port, interactive }
}
}
+1 -10
View File
@@ -1,14 +1,5 @@
/// ALPN identifying the pixelpass video wire protocol on the iroh tunnel.
/// ALPN identifying the pixelpass wire protocol on the iroh tunnel.
///
/// Bump the version suffix whenever the wire format changes. Today the wire is
/// "raw MPEG-TS bytes copied bidirectionally," so bumps will be rare.
pub const ALPN: &[u8] = b"pixelpass/0";
/// ALPN for the friends control plane — the always-on presence endpoint that
/// carries friend requests and shared codes between peers' GUIs. Separate from
/// [`ALPN`] so the same machine can run a control endpoint and a video endpoint
/// without their accept loops colliding, and so a control dial never lands on a
/// bare video host (which doesn't speak this protocol). GUI-only, like the rest
/// of the friends stack.
#[cfg(feature = "gui")]
pub const CONTROL_ALPN: &[u8] = b"pixelpass/ctrl/0";
+1 -33
View File
@@ -65,41 +65,9 @@ pub fn measure_upstream_blocking() -> Result<Measurement> {
/// via SAFETY_FACTOR) into a recommended viewer count. Floors to at least 1.
pub fn recommended_max_viewers(safe_mbps: f64, bitrate_kbps: u32) -> u32 {
let per_viewer_mbps = (bitrate_kbps as f64) / 1000.0;
// Guard non-finite / non-positive inputs (only reachable from a corrupted
// config): a NaN safe_mbps would cast to 0 and an infinite one to u32::MAX,
// both of which break the "at least 1" contract.
if !safe_mbps.is_finite() || safe_mbps <= 0.0 || per_viewer_mbps <= 0.0 {
if per_viewer_mbps <= 0.0 {
return 1;
}
let n = (safe_mbps / per_viewer_mbps).floor();
if n < 1.0 { 1 } else { n as u32 }
}
#[cfg(test)]
mod tests {
use super::recommended_max_viewers;
#[test]
fn divides_bandwidth_by_per_viewer_bitrate() {
// 8 Mbps safe / 2 Mbps each = 4 viewers.
assert_eq!(recommended_max_viewers(8.0, 2000), 4);
// Floors the fractional part: 7.9 / 2 = 3.95 -> 3.
assert_eq!(recommended_max_viewers(7.9, 2000), 3);
}
#[test]
fn floors_to_at_least_one() {
// Not even enough for one viewer still allows one (best effort).
assert_eq!(recommended_max_viewers(0.5, 2000), 1);
// Zero / unknown bitrate can't size a budget; floor to one.
assert_eq!(recommended_max_viewers(8.0, 0), 1);
}
#[test]
fn degenerate_inputs_floor_to_one() {
// A corrupted config must not yield 0 (NaN) or u32::MAX (Inf).
assert_eq!(recommended_max_viewers(f64::NAN, 2000), 1);
assert_eq!(recommended_max_viewers(f64::INFINITY, 2000), 1);
assert_eq!(recommended_max_viewers(-5.0, 2000), 1);
}
}
+13 -59
View File
@@ -1,7 +1,8 @@
//! Persistent user-level config at `~/.config/pixelpass/config.toml`.
//!
//! It tracks the bandwidth pre-flight result and the GUI's preferences.
//! Further settings can hang off the same file under their own `[section]`.
//! Right now this only tracks the bandwidth pre-flight result. Future
//! preferences (default player, default bitrate, etc.) can hang off the
//! same file under their own `[section]`.
use anyhow::{Context, Result};
use chrono::{DateTime, Utc};
@@ -15,59 +16,6 @@ use std::path::PathBuf;
pub struct Config {
#[serde(default)]
pub bandwidth: BandwidthEntry,
#[serde(default)]
pub gui: GuiSettings,
}
/// Preferences for the `pixelpass --gui` front-end.
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct GuiSettings {
/// When true, the window's close button hides the app to the system tray
/// (keeping any live stream running) instead of quitting. Defaults to
/// false — closing quits, which is what people expect.
#[serde(default)]
pub close_to_tray: bool,
/// When true, the host screen renders a QR-code panel for the ticket.
/// Defaults to true; the toggle exists for users who prefer the plain
/// text-only host screen.
#[serde(default = "default_true")]
pub show_qr: bool,
/// Name of the active GUI colour theme (a built-in, or a user file in
/// `~/.config/pixelpass/themes/`). Defaults to the built-in Default Dark.
#[serde(default = "default_theme")]
pub theme: String,
/// The display name shown to friends (in requests and shared codes).
/// Seeded from the login name; editable in Settings.
#[serde(default = "default_display_name")]
pub display_name: String,
}
impl Default for GuiSettings {
fn default() -> Self {
Self {
close_to_tray: false,
show_qr: true,
theme: default_theme(),
display_name: default_display_name(),
}
}
}
fn default_true() -> bool {
true
}
fn default_theme() -> String {
"Default Dark".to_string()
}
/// Seed the friends display name from the login name, falling back to a
/// generic label when `$USER` isn't set.
fn default_display_name() -> String {
std::env::var("USER")
.ok()
.filter(|s| !s.trim().is_empty())
.unwrap_or_else(|| "PixelPass user".to_string())
}
/// Result of the first-run upstream measurement.
@@ -89,15 +37,19 @@ pub struct BandwidthEntry {
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
#[derive(Default)]
pub enum BandwidthStatus {
#[default]
Unmeasured,
Measured,
Skipped,
Failed,
}
impl Default for BandwidthStatus {
fn default() -> Self {
Self::Unmeasured
}
}
fn default_status() -> BandwidthStatus {
BandwidthStatus::Unmeasured
}
@@ -129,9 +81,11 @@ pub fn save(cfg: &Config) -> Result<()> {
let parent = path
.parent()
.context("config path has no parent directory")?;
fs::create_dir_all(parent).with_context(|| format!("failed to create {}", parent.display()))?;
fs::create_dir_all(parent)
.with_context(|| format!("failed to create {}", parent.display()))?;
let serialized = toml::to_string_pretty(cfg).context("failed to serialize config to TOML")?;
let serialized =
toml::to_string_pretty(cfg).context("failed to serialize config to TOML")?;
let tmp = parent.join(format!(".config.toml.tmp.{}", std::process::id()));
{
-257
View File
@@ -1,257 +0,0 @@
//! Friends control-plane protocol and service.
//!
//! This is the always-on presence channel that rides the [`CONTROL_ALPN`]
//! endpoint (bound with the persistent identity — see
//! [`super::endpoint::bind_control`]). It's how two peers' GUIs exchange friend
//! requests and pushed share-codes, independent of any video session.
//!
//! Wire shape: **one message per connection.** The sender opens a bi-stream,
//! writes the JSON-encoded [`ControlMsg`], and finishes its send side (EOF
//! delimits the message — no length framing needed). The receiver reads to EOF,
//! parses, hands the message up, then writes a one-byte [`ACK`] back so the
//! sender knows it was delivered *and* parsed. That delivery signal is what
//! lets the host-side code-push queue (a later phase) tell "sent" from "friend
//! was offline." A friend's *reply* (accept/decline) is a separate later
//! connection in the other direction, because acceptance can happen minutes
//! after the request — not a response on the same stream.
use std::time::Duration;
use anyhow::{Context, Result, bail};
use iroh::endpoint::{Incoming, VarInt};
use iroh::{Endpoint, EndpointAddr, EndpointId};
use serde::{Deserialize, Serialize};
use tokio::sync::mpsc;
use super::alpn::CONTROL_ALPN;
/// Upper bound on a single control message. Generous for a display name plus a
/// share-code ticket (~150 chars); rejects a peer trying to make us buffer a
/// huge blob.
const MAX_MSG: usize = 64 * 1024;
/// One-byte application acknowledgement the receiver returns once it has parsed
/// a message. ASCII ACK (0x06).
const ACK: &[u8] = b"\x06";
/// Bound on each phase of the send handshake, so a half-dead peer or relay
/// can't park a sender (or an inbound handler) forever.
const IO_TIMEOUT: Duration = Duration::from_secs(10);
/// A message on the friends control plane.
///
/// `#[serde(tag = "type")]` keeps the JSON self-describing and lets us add
/// variants without breaking older peers (an unknown tag fails to parse and is
/// logged, rather than being silently misread as another variant).
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum ControlMsg {
/// "I'm online; here's my current display name." A presence/name refresh.
Hello { name: String },
/// Ask the recipient to become friends.
FriendRequest { name: String },
/// Accept a request the recipient previously sent us.
FriendAccept { name: String },
/// Decline a pending request, or cancel an outgoing one.
FriendDecline,
/// A host pushing a freshly generated share-code to an accepted friend.
ShareCode { name: String, ticket: String },
}
/// A received control message, paired with the *authenticated* sender id (the
/// connection's verified remote public key — not a value the peer can spoof in
/// the payload, which is why no variant carries a sender id).
#[derive(Debug, Clone)]
pub struct Inbound {
pub from: EndpointId,
pub msg: ControlMsg,
}
fn encode(msg: &ControlMsg) -> Result<Vec<u8>> {
serde_json::to_vec(msg).context("failed to encode control message")
}
fn decode(bytes: &[u8]) -> Result<ControlMsg> {
serde_json::from_slice(bytes).context("failed to decode control message")
}
/// Deliver one message to `peer` over `endpoint`, returning once the recipient
/// has acknowledged it. An error means it was *not* delivered (peer offline,
/// unreachable, or rejected the stream) — the caller can queue and retry.
///
/// `peer` is usually a bare [`EndpointId`] — friends store only the stable id,
/// and n0 DNS discovery resolves it to a live address. The full [`EndpointAddr`]
/// form exists for callers that already hold one (and for hermetic tests).
pub async fn send(
endpoint: &Endpoint,
peer: impl Into<EndpointAddr>,
msg: &ControlMsg,
) -> Result<()> {
let payload = encode(msg)?;
let conn = tokio::time::timeout(IO_TIMEOUT, endpoint.connect(peer, CONTROL_ALPN))
.await
.context("timed out connecting to peer")?
.context("failed to connect to peer")?;
let io = async {
let (mut send, mut recv) = conn
.open_bi()
.await
.context("failed to open control stream")?;
send.write_all(&payload)
.await
.context("failed to write control message")?;
send.finish().context("failed to finish control stream")?;
// Read the peer's ACK. read_to_end returns once the peer finishes its
// send side, so this also serves as "the peer is done with us."
let ack = recv
.read_to_end(ACK.len() + 1)
.await
.context("peer closed the control stream without acknowledging")?;
if ack != ACK {
bail!(
"peer sent an unexpected acknowledgement ({} bytes)",
ack.len()
);
}
Ok(())
};
let result = tokio::time::timeout(IO_TIMEOUT, io)
.await
.context("timed out sending control message")?;
// Clean close so the peer's `closed().await` returns promptly either way.
conn.close(VarInt::from_u32(0), b"done");
result
}
/// Run the control-plane accept loop, forwarding every received message to
/// `tx`. Returns when the endpoint stops accepting (i.e. it was closed).
pub async fn serve(endpoint: Endpoint, tx: mpsc::Sender<Inbound>) {
while let Some(incoming) = endpoint.accept().await {
let tx = tx.clone();
tokio::spawn(async move {
if let Err(e) = handle(incoming, &tx).await {
tracing::warn!("control: inbound connection failed: {e:#}");
}
});
}
tracing::info!("control: endpoint stopped accepting");
}
async fn handle(incoming: Incoming, tx: &mpsc::Sender<Inbound>) -> Result<()> {
let conn = incoming
.await
.context("inbound control connection failed")?;
let from = conn.remote_id();
let msg = async {
let (mut send, mut recv) = conn
.accept_bi()
.await
.context("failed to accept control stream")?;
let bytes = recv
.read_to_end(MAX_MSG)
.await
.context("failed to read control message")?;
let msg = decode(&bytes)?;
// ACK only after a successful parse, so the sender's delivery signal
// means "received and understood."
send.write_all(ACK).await.context("failed to write ack")?;
send.finish().context("failed to finish ack stream")?;
Ok::<_, anyhow::Error>(msg)
};
let msg = tokio::time::timeout(IO_TIMEOUT, msg)
.await
.context("timed out reading control message")??;
// Wait (briefly) for the sender's close so our ACK flushes before the
// connection is dropped at the end of this scope.
let _ = tokio::time::timeout(IO_TIMEOUT, conn.closed()).await;
tx.send(Inbound { from, msg })
.await
.map_err(|_| anyhow::anyhow!("control: receiver dropped"))?;
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn control_msg_round_trips() {
let cases = [
ControlMsg::Hello {
name: "alice".into(),
},
ControlMsg::FriendRequest { name: "bob".into() },
ControlMsg::FriendAccept {
name: "carol".into(),
},
ControlMsg::FriendDecline,
ControlMsg::ShareCode {
name: "dave".into(),
ticket: "endpointaa…".into(),
},
];
for msg in cases {
let bytes = encode(&msg).unwrap();
assert_eq!(decode(&bytes).unwrap(), msg);
}
}
#[test]
fn unknown_tag_is_rejected() {
assert!(decode(br#"{"type":"nonsense"}"#).is_err());
}
/// Bind a control-plane endpoint with a *fresh* random key, so two of them
/// in one test get distinct ids (two real machines each have their own
/// persistent key; `bind_control` would give both the same one here, and
/// iroh refuses "connecting to ourself").
async fn bind_test_control() -> Endpoint {
iroh::Endpoint::builder(iroh::endpoint::presets::N0)
.secret_key(iroh::SecretKey::generate())
.alpns(vec![CONTROL_ALPN.to_vec()])
.bind()
.await
.unwrap()
}
/// End-to-end over two real iroh endpoints on this machine. Ignored by
/// default — it binds endpoints and waits on the relay, so it's slow and
/// network-dependent. Run with `cargo test -- --ignored control`.
#[tokio::test]
#[ignore = "binds real iroh endpoints; run on demand"]
async fn loopback_delivers_and_acks() {
let server = bind_test_control().await;
let client = bind_test_control().await;
// Connect by full addr so the test doesn't depend on DNS discovery.
server.online().await;
client.online().await;
let server_addr = server.addr();
let (tx, mut rx) = mpsc::channel(4);
let server_ep = server.clone();
let serve_task = tokio::spawn(async move { serve(server_ep, tx).await });
let msg = ControlMsg::FriendRequest {
name: "tester".into(),
};
// Full addr (not just the id) so the test doesn't depend on DNS discovery.
send(&client, server_addr.clone(), &msg).await.unwrap();
let got = tokio::time::timeout(Duration::from_secs(15), rx.recv())
.await
.expect("no inbound within 15s")
.expect("channel closed");
assert_eq!(got.msg, msg);
assert_eq!(got.from, client.addr().id);
server.close().await;
client.close().await;
serve_task.abort();
}
}
+14 -40
View File
@@ -85,24 +85,20 @@ fn install_hint_for_bin(bin: &str) -> String {
let distro = detect_distro();
let pkg = match bin {
"gst-launch-1.0" | "gst-inspect-1.0" => match distro.as_deref() {
Some("arch" | "cachyos" | "manjaro" | "endeavouros" | "artix" | "garuda") => {
"gstreamer gst-plugins-base"
}
Some("arch" | "cachyos" | "manjaro" | "endeavouros") => "gstreamer gst-plugins-base",
Some("debian" | "ubuntu" | "pop" | "linuxmint") => "gstreamer1.0-tools",
Some("fedora" | "nobara") => "gstreamer1 gstreamer1-plugins-base-tools",
_ => "gstreamer + tools",
},
"pactl" => match distro.as_deref() {
Some("arch" | "cachyos" | "manjaro" | "endeavouros" | "artix" | "garuda") => "libpulse",
Some("arch" | "cachyos" | "manjaro" | "endeavouros") => "libpulse",
Some("debian" | "ubuntu" | "pop" | "linuxmint") => "pulseaudio-utils",
Some("fedora" | "nobara") => "pulseaudio-utils",
Some("opensuse" | "opensuse-tumbleweed" | "opensuse-leap") => "pulseaudio-utils",
_ => "pulseaudio-utils (provides `pactl`)",
},
"xwininfo" => match distro.as_deref() {
Some("arch" | "cachyos" | "manjaro" | "endeavouros" | "artix" | "garuda") => {
"xorg-xwininfo"
}
Some("arch" | "cachyos" | "manjaro" | "endeavouros") => "xorg-xwininfo",
Some("debian" | "ubuntu" | "pop" | "linuxmint") => "x11-utils",
Some("fedora" | "nobara") => "xorg-x11-utils",
Some("opensuse" | "opensuse-tumbleweed" | "opensuse-leap") => "xwininfo",
@@ -117,74 +113,56 @@ fn install_hint_for_gst_element(name: &str) -> String {
let distro = detect_distro();
let pkg = match name {
"pipewiresrc" => match distro.as_deref() {
Some("arch" | "cachyos" | "manjaro" | "endeavouros" | "artix" | "garuda") => {
"gst-plugin-pipewire"
}
Some("arch" | "cachyos" | "manjaro" | "endeavouros") => "gst-plugin-pipewire",
Some("debian" | "ubuntu" | "pop" | "linuxmint") => "gstreamer1.0-pipewire",
Some("fedora" | "nobara") => "pipewire-gstreamer",
Some("opensuse" | "opensuse-tumbleweed" | "opensuse-leap") => "pipewire-gstreamer",
_ => "the GStreamer PipeWire plugin",
},
"vah264enc" => match distro.as_deref() {
Some("arch" | "cachyos" | "manjaro" | "endeavouros" | "artix" | "garuda") => {
"gst-plugin-va"
}
Some("arch" | "cachyos" | "manjaro" | "endeavouros") => "gst-plugin-va",
Some("debian" | "ubuntu" | "pop" | "linuxmint") => "gstreamer1.0-plugins-bad",
Some("fedora" | "nobara") => "gstreamer1-plugins-bad-free",
Some("opensuse" | "opensuse-tumbleweed" | "opensuse-leap") => "gstreamer-plugins-bad",
_ => {
"the GStreamer VA-API plugin (requires an H.264-capable GPU; almost all modern GPUs)"
}
_ => "the GStreamer VA-API plugin (requires an H.264-capable GPU; almost all modern GPUs)",
},
"x264enc" => match distro.as_deref() {
Some("arch" | "cachyos" | "manjaro" | "endeavouros" | "artix" | "garuda") => {
"gst-plugins-ugly"
}
Some("arch" | "cachyos" | "manjaro" | "endeavouros") => "gst-plugins-ugly",
Some("debian" | "ubuntu" | "pop" | "linuxmint") => "gstreamer1.0-plugins-ugly",
Some("fedora" | "nobara") => "gstreamer1-plugins-ugly",
Some("opensuse" | "opensuse-tumbleweed" | "opensuse-leap") => "gstreamer-plugins-ugly",
_ => "the GStreamer x264 plugin (plugins-ugly)",
},
"ximagesrc" => match distro.as_deref() {
Some("arch" | "cachyos" | "manjaro" | "endeavouros" | "artix" | "garuda") => {
"gst-plugins-good"
}
Some("arch" | "cachyos" | "manjaro" | "endeavouros") => "gst-plugins-good",
Some("debian" | "ubuntu" | "pop" | "linuxmint") => "gstreamer1.0-plugins-good",
Some("fedora" | "nobara") => "gstreamer1-plugins-good",
Some("opensuse" | "opensuse-tumbleweed" | "opensuse-leap") => "gstreamer-plugins-good",
_ => "the GStreamer X11 plugin (plugins-good)",
},
"videoscale" => match distro.as_deref() {
Some("arch" | "cachyos" | "manjaro" | "endeavouros" | "artix" | "garuda") => {
"gst-plugins-base"
}
Some("arch" | "cachyos" | "manjaro" | "endeavouros") => "gst-plugins-base",
Some("debian" | "ubuntu" | "pop" | "linuxmint") => "gstreamer1.0-plugins-base",
Some("fedora" | "nobara") => "gstreamer1-plugins-base",
Some("opensuse" | "opensuse-tumbleweed" | "opensuse-leap") => "gstreamer-plugins-base",
_ => "the GStreamer plugins-base set",
},
"h264parse" | "mpegtsmux" | "aacparse" => match distro.as_deref() {
Some("arch" | "cachyos" | "manjaro" | "endeavouros" | "artix" | "garuda") => {
"gst-plugins-bad"
}
Some("arch" | "cachyos" | "manjaro" | "endeavouros") => "gst-plugins-bad",
Some("debian" | "ubuntu" | "pop" | "linuxmint") => "gstreamer1.0-plugins-bad",
Some("fedora" | "nobara") => "gstreamer1-plugins-bad-free",
Some("opensuse" | "opensuse-tumbleweed" | "opensuse-leap") => "gstreamer-plugins-bad",
_ => "the GStreamer plugins-bad set",
},
"pulsesrc" => match distro.as_deref() {
Some("arch" | "cachyos" | "manjaro" | "endeavouros" | "artix" | "garuda") => {
"gst-plugins-good"
}
Some("arch" | "cachyos" | "manjaro" | "endeavouros") => "gst-plugins-good",
Some("debian" | "ubuntu" | "pop" | "linuxmint") => "gstreamer1.0-pulseaudio",
Some("fedora" | "nobara") => "gstreamer1-plugins-good",
Some("opensuse" | "opensuse-tumbleweed" | "opensuse-leap") => "gstreamer-plugins-good",
_ => "the GStreamer PulseAudio plugin",
},
"avenc_aac" => match distro.as_deref() {
Some("arch" | "cachyos" | "manjaro" | "endeavouros" | "artix" | "garuda") => {
"gst-libav"
}
Some("arch" | "cachyos" | "manjaro" | "endeavouros") => "gst-libav",
Some("debian" | "ubuntu" | "pop" | "linuxmint") => "gstreamer1.0-libav",
Some("fedora" | "nobara") => "gstreamer1-libav",
Some("opensuse" | "opensuse-tumbleweed" | "opensuse-leap") => "gstreamer-libav",
@@ -197,14 +175,10 @@ fn install_hint_for_gst_element(name: &str) -> String {
fn install_command(distro: &Option<String>, pkg: &str) -> String {
let cmd = match distro.as_deref() {
Some("arch" | "cachyos" | "manjaro" | "endeavouros" | "artix" | "garuda") => {
format!("sudo pacman -S {pkg}")
}
Some("arch" | "cachyos" | "manjaro" | "endeavouros") => format!("sudo pacman -S {pkg}"),
Some("debian" | "ubuntu" | "pop" | "linuxmint") => format!("sudo apt install {pkg}"),
Some("fedora" | "nobara") => format!("sudo dnf install {pkg}"),
Some("opensuse" | "opensuse-tumbleweed" | "opensuse-leap") => {
format!("sudo zypper install {pkg}")
}
Some("opensuse" | "opensuse-tumbleweed" | "opensuse-leap") => format!("sudo zypper install {pkg}"),
_ => format!("install the `{pkg}` package via your distro's package manager"),
};
format!("Install hint: {cmd}")
-84
View File
@@ -1,84 +0,0 @@
//! Shared iroh endpoint construction.
//!
//! Two planes, two identities:
//!
//! * The **video** plane (host/viewer sessions) binds with an *ephemeral*
//! keypair — a fresh `EndpointId` per run. Each session is a throwaway tunnel,
//! and keeping its id ephemeral means a screen-share leaks no stable
//! fingerprint.
//! * The **control** plane (the always-on friends presence service) binds with
//! the machine's *persistent* identity (see [`identity`]), so peers can find
//! and recognise each other across launches.
//!
//! They must use different identities because both can be live at once on the
//! same machine (the GUI's control endpoint while a host session runs), and
//! iroh routes by `EndpointId` — two live endpoints sharing one id would make
//! relay delivery ambiguous.
use std::str::FromStr;
use anyhow::{Context, Result};
use iroh::endpoint::presets;
use iroh::{Endpoint, RelayMap, RelayMode, RelayUrl};
use super::alpn::ALPN;
/// Environment variable consulted when `--relay` isn't passed. Lets the GUI's
/// child processes and scripted runs inherit a relay choice without a flag.
pub const RELAY_ENV: &str = "PIXELPASS_RELAY";
/// Resolve the relay override: explicit `--relay` wins, else `PIXELPASS_RELAY`,
/// else `None` (use the bundled defaults).
pub fn relay_override(flag: Option<&str>) -> Option<String> {
flag.map(str::to_owned).or_else(|| {
std::env::var(RELAY_ENV)
.ok()
.filter(|s| !s.trim().is_empty())
})
}
/// Bind a **video-plane** endpoint (host/viewer) with an ephemeral identity.
///
/// With no `relay` override we use [`presets::N0`] — n0 DNS discovery, the
/// library's default relays, and the chosen crypto provider. With an override
/// we keep all of that but swap in a single custom relay via
/// [`RelayMode::Custom`]; this is how a user gets off the rc's bundled
/// (canary-grade) relays or points at a self-hosted one. Discovery is
/// unchanged, so peers still resolve each other by endpoint id.
pub async fn bind(relay: Option<&str>) -> Result<Endpoint> {
// No `secret_key` set → iroh mints a fresh ephemeral keypair for this run.
bind_with(relay, None, ALPN).await
}
/// Bind the **control-plane** endpoint with the machine's persistent identity
/// (see [`super::identity`]) and the friends [`super::alpn::CONTROL_ALPN`]. Its
/// `EndpointId` is the stable id friends know you by.
#[cfg(feature = "gui")]
pub async fn bind_control(relay: Option<&str>) -> Result<Endpoint> {
let secret_key = super::identity::load_or_create()?;
bind_with(relay, Some(secret_key), super::alpn::CONTROL_ALPN).await
}
/// Shared builder: optional persistent key (None → ephemeral) + the plane's ALPN.
async fn bind_with(
relay: Option<&str>,
key: Option<iroh::SecretKey>,
alpn: &[u8],
) -> Result<Endpoint> {
let mut builder = Endpoint::builder(presets::N0).alpns(vec![alpn.to_vec()]);
if let Some(key) = key {
builder = builder.secret_key(key);
}
if let Some(url) = relay {
let url = RelayUrl::from_str(url).with_context(|| {
format!("invalid relay URL {url:?} (expected e.g. https://relay.example/)")
})?;
builder = builder.relay_mode(RelayMode::Custom(RelayMap::from(url)));
}
builder
.bind()
.await
.context("failed to bind the iroh endpoint")
}
-296
View File
@@ -1,296 +0,0 @@
//! Persistent friends store at `~/.config/pixelpass/friends.toml`.
//!
//! Kept in its own file rather than a `[friends]` section of `config.toml` so
//! the headless CLI — which never manages friends and would round-trip the
//! config without this knowledge — can't drop the list on a `--reconfigure`.
//! Same reasoning as the separate `identity.key`.
//!
//! A friend is identified by their stable control-plane [`EndpointId`] (the id
//! from [`super::endpoint::bind_control`]). `EndpointId` serialises as its
//! string form in TOML, so the file is human-readable and hand-editable.
use anyhow::{Context, Result};
use iroh::EndpointId;
use serde::{Deserialize, Serialize};
use std::fs;
use std::io::Write;
use std::path::PathBuf;
/// Where a friendship sits in the mutual-consent handshake.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum FriendState {
/// We've sent them a request and are waiting for them to accept.
PendingOutgoing,
/// They've requested us; waiting for the local user to accept or decline.
PendingIncoming,
/// Both sides have agreed — a real friend.
Accepted,
}
/// One entry in the friends list.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Friend {
pub id: EndpointId,
/// Display name — seeded from the name the peer reported, locally editable.
pub name: String,
pub state: FriendState,
/// Whether the host auto-shares its session code with this friend. Toggled
/// on the host's share picker; persisted here so the choice survives a
/// restart. Defaults to `true` so a newly added friend is included (and an
/// older `friends.toml` without the field loads as share-with-all).
#[serde(default = "default_share")]
pub share: bool,
}
fn default_share() -> bool {
true
}
/// The persisted friends list. Serialises as a TOML array of tables
/// (`[[friends]]`).
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct FriendStore {
#[serde(default)]
pub friends: Vec<Friend>,
}
/// Returns `~/.config/pixelpass/friends.toml`. Shares the config directory with
/// [`super::config`]; the parent is created on save.
pub fn friends_path() -> Result<PathBuf> {
Ok(super::config::config_path()?
.parent()
.context("config path has no parent directory")?
.join("friends.toml"))
}
/// Load the store, or a default (empty) one if the file doesn't exist yet.
/// Parse errors bubble up so a hand-edit being debugged isn't silently
/// overwritten.
pub fn load() -> Result<FriendStore> {
let path = friends_path()?;
match fs::read_to_string(&path) {
Ok(s) => toml::from_str(&s).with_context(|| format!("failed to parse {}", path.display())),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(FriendStore::default()),
Err(e) => Err(e).with_context(|| format!("failed to read {}", path.display())),
}
}
impl FriendStore {
/// Atomic write via tempfile-in-same-dir + rename (mirrors
/// [`super::config::save`]).
pub fn save(&self) -> Result<()> {
let path = friends_path()?;
let parent = path
.parent()
.context("friends path has no parent directory")?;
fs::create_dir_all(parent)
.with_context(|| format!("failed to create {}", parent.display()))?;
let serialized = toml::to_string_pretty(self).context("failed to serialize friends")?;
let tmp = parent.join(format!(".friends.toml.tmp.{}", std::process::id()));
{
let mut f = fs::File::create(&tmp)
.with_context(|| format!("failed to create {}", tmp.display()))?;
f.write_all(serialized.as_bytes())
.with_context(|| format!("failed to write {}", tmp.display()))?;
f.sync_all().ok();
}
fs::rename(&tmp, &path)
.with_context(|| format!("failed to rename {} -> {}", tmp.display(), path.display()))?;
Ok(())
}
pub fn find(&self, id: &EndpointId) -> Option<&Friend> {
self.friends.iter().find(|f| &f.id == id)
}
pub fn find_mut(&mut self, id: &EndpointId) -> Option<&mut Friend> {
self.friends.iter_mut().find(|f| &f.id == id)
}
/// True iff this id is a fully-accepted friend — the gate the code-push
/// (Phase 4) and "is this a known friend?" checks use.
pub fn is_accepted(&self, id: &EndpointId) -> bool {
matches!(
self.find(id),
Some(Friend {
state: FriendState::Accepted,
..
})
)
}
/// Insert a new friend, or update an existing one's `name`/`state` in place.
/// Returns a mutable reference to the stored entry.
pub fn upsert(&mut self, id: EndpointId, name: String, state: FriendState) -> &mut Friend {
if let Some(idx) = self.friends.iter().position(|f| f.id == id) {
let f = &mut self.friends[idx];
f.name = name;
f.state = state;
f
} else {
self.friends.push(Friend {
id,
name,
state,
share: true,
});
self.friends.last_mut().expect("just pushed")
}
}
/// Remove a friend by id. Returns whether an entry was removed.
pub fn remove(&mut self, id: &EndpointId) -> bool {
let before = self.friends.len();
self.friends.retain(|f| &f.id != id);
self.friends.len() != before
}
/// Apply an inbound friend request. Returns `true` if it *completes a mutual
/// match* — we'd already sent them one, so they're now [`Accepted`] and the
/// caller should reply with a `FriendAccept`. Otherwise it's recorded as
/// [`PendingIncoming`] for the user to act on and `false` is returned.
///
/// [`Accepted`]: FriendState::Accepted
/// [`PendingIncoming`]: FriendState::PendingIncoming
pub fn on_friend_request(&mut self, id: EndpointId, name: String) -> bool {
if matches!(
self.find(&id).map(|f| f.state),
Some(FriendState::PendingOutgoing)
) {
self.upsert(id, name, FriendState::Accepted);
true
} else {
self.upsert(id, name, FriendState::PendingIncoming);
false
}
}
/// Apply an inbound acceptance of a request we sent. Returns `true` if it
/// advanced a friendship to [`Accepted`] (i.e. we actually knew this peer);
/// an accept from a stranger is ignored.
///
/// [`Accepted`]: FriendState::Accepted
pub fn on_friend_accept(&mut self, id: EndpointId, name: String) -> bool {
if self.find(&id).is_some() {
self.upsert(id, name, FriendState::Accepted);
true
} else {
false
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn sample_id() -> EndpointId {
iroh::SecretKey::generate().public()
}
#[test]
fn round_trips_through_toml() {
let mut store = FriendStore::default();
store.upsert(sample_id(), "Alice".into(), FriendState::Accepted);
store.upsert(sample_id(), "Bob".into(), FriendState::PendingIncoming);
let toml = toml::to_string_pretty(&store).unwrap();
let back: FriendStore = toml::from_str(&toml).unwrap();
assert_eq!(back.friends, store.friends);
}
#[test]
fn new_friends_default_to_shared_and_survive_round_trip() {
let mut store = FriendStore::default();
let id = sample_id();
store.upsert(id, "Alice".into(), FriendState::Accepted);
assert!(store.find(&id).unwrap().share, "new friends start shared");
// An older friends.toml predating the field loads as share-with-all.
let toml = format!("[[friends]]\nid = \"{id}\"\nname = \"Legacy\"\nstate = \"accepted\"\n");
let back: FriendStore = toml::from_str(&toml).unwrap();
assert!(back.friends[0].share);
}
#[test]
fn upsert_preserves_share_across_refresh() {
let mut store = FriendStore::default();
let id = sample_id();
store.upsert(id, "Alice".into(), FriendState::Accepted);
store.find_mut(&id).unwrap().share = false;
// A later name/presence refresh re-upserts the same peer; the share
// choice must not be reset by it.
store.upsert(id, "Alice (new name)".into(), FriendState::Accepted);
assert!(!store.find(&id).unwrap().share);
}
#[test]
fn upsert_updates_in_place() {
let mut store = FriendStore::default();
let id = sample_id();
store.upsert(id, "Old".into(), FriendState::PendingOutgoing);
store.upsert(id, "New".into(), FriendState::Accepted);
assert_eq!(store.friends.len(), 1);
let f = store.find(&id).unwrap();
assert_eq!(f.name, "New");
assert_eq!(f.state, FriendState::Accepted);
}
#[test]
fn is_accepted_only_for_accepted_state() {
let mut store = FriendStore::default();
let pending = sample_id();
let friend = sample_id();
store.upsert(pending, "P".into(), FriendState::PendingOutgoing);
store.upsert(friend, "F".into(), FriendState::Accepted);
assert!(!store.is_accepted(&pending));
assert!(store.is_accepted(&friend));
assert!(!store.is_accepted(&sample_id()));
}
#[test]
fn remove_reports_whether_present() {
let mut store = FriendStore::default();
let id = sample_id();
store.upsert(id, "X".into(), FriendState::Accepted);
assert!(store.remove(&id));
assert!(!store.remove(&id));
assert!(store.friends.is_empty());
}
#[test]
fn incoming_request_from_stranger_is_pending() {
let mut store = FriendStore::default();
let id = sample_id();
let mutual = store.on_friend_request(id, "Stranger".into());
assert!(!mutual);
assert_eq!(store.find(&id).unwrap().state, FriendState::PendingIncoming);
}
#[test]
fn incoming_request_matching_our_outgoing_is_mutual() {
let mut store = FriendStore::default();
let id = sample_id();
// We asked them first…
store.upsert(id, "Pal".into(), FriendState::PendingOutgoing);
// …then their request arrives — that's a mutual match.
let mutual = store.on_friend_request(id, "Pal".into());
assert!(mutual);
assert_eq!(store.find(&id).unwrap().state, FriendState::Accepted);
}
#[test]
fn accept_advances_known_peer_only() {
let mut store = FriendStore::default();
let known = sample_id();
store.upsert(known, "Known".into(), FriendState::PendingOutgoing);
assert!(store.on_friend_accept(known, "Known".into()));
assert_eq!(store.find(&known).unwrap().state, FriendState::Accepted);
// An accept from someone we never asked is ignored.
let stranger = sample_id();
assert!(!store.on_friend_accept(stranger, "Nope".into()));
assert!(store.find(&stranger).is_none());
}
}
-141
View File
@@ -1,141 +0,0 @@
//! Persistent node identity at `~/.config/pixelpass/identity.key`.
//!
//! Without this, [`super::endpoint::bind`] would let iroh mint a fresh random
//! keypair on every launch, so a peer's `EndpointId` would change each run.
//! The friends system identifies people by that id (it's the public key already
//! embedded in every share code), so it must stay stable across launches — and
//! across roles: the same machine gets the same id whether it's hosting,
//! viewing, or just sitting in the GUI.
//!
//! The key is the ed25519 secret (32 bytes) stored as hex on its own line, in a
//! `0600` file separate from `config.toml` — it's a secret, not a preference,
//! and keeping it out of the TOML means a hand-edit or a config reset can't
//! clobber your identity.
use anyhow::{Context, Result, bail};
use iroh::SecretKey;
use std::fs;
use std::io::Write;
use std::path::PathBuf;
/// Returns `~/.config/pixelpass/identity.key` (or the XDG equivalent). Shares
/// the config directory with [`super::config`]; the parent is created on save.
pub fn identity_path() -> Result<PathBuf> {
Ok(super::config::config_path()?
.parent()
.context("config path has no parent directory")?
.join("identity.key"))
}
/// Load the persisted secret key, or generate-and-save one on first run.
///
/// A malformed file is a hard error rather than a silent regenerate: silently
/// minting a new identity would orphan every friend who has the old id, so we'd
/// rather fail loud and let the user notice (and decide) than lose it quietly.
pub fn load_or_create() -> Result<SecretKey> {
let path = identity_path()?;
match fs::read_to_string(&path) {
Ok(s) => parse_key(s.trim())
.with_context(|| format!("failed to parse the identity key at {}", path.display())),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
let key = SecretKey::generate();
save(&key)?;
tracing::info!(id = %key.public(), "generated a new persistent identity");
Ok(key)
}
Err(e) => Err(e).with_context(|| format!("failed to read {}", path.display())),
}
}
fn parse_key(hex: &str) -> Result<SecretKey> {
let bytes = decode_hex(hex)?;
let arr: [u8; 32] = bytes
.try_into()
.map_err(|_| anyhow::anyhow!("identity key must be 32 bytes (64 hex chars)"))?;
Ok(SecretKey::from_bytes(&arr))
}
/// Atomic, `0600` write: tempfile-in-same-dir, chmod, then rename. Same
/// approach as [`super::config::save`], but with restrictive perms applied
/// before the rename so the secret is never briefly world-readable.
pub fn save(key: &SecretKey) -> Result<()> {
let path = identity_path()?;
let parent = path
.parent()
.context("identity path has no parent directory")?;
fs::create_dir_all(parent).with_context(|| format!("failed to create {}", parent.display()))?;
let tmp = parent.join(format!(".identity.key.tmp.{}", std::process::id()));
{
let mut f = fs::File::create(&tmp)
.with_context(|| format!("failed to create {}", tmp.display()))?;
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
f.set_permissions(fs::Permissions::from_mode(0o600))
.with_context(|| format!("failed to chmod {}", tmp.display()))?;
}
f.write_all(encode_hex(&key.to_bytes()).as_bytes())
.with_context(|| format!("failed to write {}", tmp.display()))?;
f.write_all(b"\n").ok();
f.sync_all().ok();
}
fs::rename(&tmp, &path)
.with_context(|| format!("failed to rename {} -> {}", tmp.display(), path.display()))?;
Ok(())
}
fn encode_hex(bytes: &[u8]) -> String {
let mut s = String::with_capacity(bytes.len() * 2);
for b in bytes {
s.push_str(&format!("{b:02x}"));
}
s
}
fn decode_hex(s: &str) -> Result<Vec<u8>> {
if !s.len().is_multiple_of(2) {
bail!("hex string has an odd length");
}
(0..s.len())
.step_by(2)
.map(|i| {
u8::from_str_radix(&s[i..i + 2], 16)
.with_context(|| format!("invalid hex byte at offset {i}"))
})
.collect()
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn hex_round_trips() {
let bytes: Vec<u8> = (0u8..=255).collect();
let encoded = encode_hex(&bytes);
assert_eq!(encoded.len(), bytes.len() * 2);
assert_eq!(decode_hex(&encoded).unwrap(), bytes);
}
#[test]
fn key_round_trips_through_hex() {
let key = SecretKey::generate();
let hex = encode_hex(&key.to_bytes());
let parsed = parse_key(&hex).unwrap();
assert_eq!(parsed.to_bytes(), key.to_bytes());
assert_eq!(parsed.public(), key.public());
}
#[test]
fn rejects_wrong_length() {
assert!(parse_key("dead").is_err());
assert!(parse_key("").is_err());
}
#[test]
fn rejects_odd_and_nonhex() {
assert!(decode_hex("abc").is_err());
assert!(decode_hex("zz").is_err());
}
}
-10
View File
@@ -1,18 +1,8 @@
pub mod alpn;
pub mod bandwidth;
pub mod config;
// The friends stack (persistent identity + control plane) is GUI-only — a
// headless CLI host runs no presence service — so it's gated with the feature
// that pulls the rest of the GUI, keeping the headless build lean.
#[cfg(feature = "gui")]
pub mod control;
pub mod deps;
pub mod display;
pub mod endpoint;
#[cfg(feature = "gui")]
pub mod friends;
#[cfg(feature = "gui")]
pub mod identity;
pub mod output;
pub mod process;
pub mod signal;
+3 -9
View File
@@ -22,11 +22,7 @@ pub fn set_json(enabled: bool) {
JSON_ENABLED.store(enabled, Ordering::Relaxed);
}
/// Whether the JSON event stream is on — i.e. we're being driven by a
/// machine front-end (the `--gui` shell-out) rather than a human terminal.
/// Gates features that only make sense under that front-end, like the
/// stdin command channel the host reads `kick` requests from.
pub fn json_enabled() -> bool {
fn json_enabled() -> bool {
JSON_ENABLED.load(Ordering::Relaxed)
}
@@ -46,11 +42,9 @@ pub enum Event<'a> {
max_viewers: u32,
max_viewers_source: &'a str,
},
/// A viewer joined. `id` is the viewer's endpoint id; `active` is the new
/// total after the join.
/// A new viewer joined.
ViewerJoined { id: &'a str, active: u32, max: u32 },
/// A viewer left — disconnected on their own or kicked by the host. `id`
/// is the viewer's endpoint id; `active` is the new total after.
/// A viewer disconnected.
ViewerLeft { id: &'a str, active: u32, max: u32 },
/// Capture pipeline lifecycle (spawned on first viewer, torn down on last).
Capture { state: CaptureState },
+3 -10
View File
@@ -7,17 +7,10 @@ pub fn install_ctrl_c() -> CancellationToken {
let token = CancellationToken::new();
let trigger = token.clone();
tokio::spawn(async move {
if let Err(e) = tokio::signal::ctrl_c().await {
// Installing the handler failed — ctrl-c won't trigger a graceful
// shutdown. Say so instead of failing silently; the user can still
// kill the process, and the second-ctrl-c arm below would only fail
// the same way, so bail out of the task.
tracing::warn!("could not install ctrl-c handler: {e}; ctrl-c won't shut down cleanly");
return;
if tokio::signal::ctrl_c().await.is_ok() {
tracing::info!("ctrl-c received, shutting down");
trigger.cancel();
}
tracing::info!("ctrl-c received, shutting down");
trigger.cancel();
if tokio::signal::ctrl_c().await.is_ok() {
tracing::warn!("second ctrl-c — exiting now");
std::process::exit(130);
+46 -79
View File
@@ -6,18 +6,17 @@
//! egui app drains each frame. stderr is captured into a small ring so a
//! failed launch (e.g. a missing gst plugin) can be surfaced in the window.
use std::io::{BufRead, BufReader, Write};
use std::process::{Child, ChildStdin, Command, Stdio};
use std::io::{BufRead, BufReader};
use std::process::{Child, Command, Stdio};
use std::sync::mpsc::Receiver;
use std::sync::{Arc, Mutex};
use std::time::Duration;
use eframe::egui;
use nix::sys::signal::{Signal, kill};
use nix::unistd::Pid;
use serde::Deserialize;
use super::Waker;
/// One parsed event from the child's stdout. Owned mirror of
/// [`crate::common::output::Event`] (which borrows for emit); kept separate so
/// the wire format and the parser can evolve independently.
@@ -67,36 +66,26 @@ pub enum CaptureState {
const STDERR_TAIL_MAX: usize = 60;
pub struct ChildProc {
/// `Some` while the child is owned here; `Drop` takes it to hand off to a
/// detached reaper thread (see the `Drop` impl).
child: Option<Child>,
child: Child,
pub rx: Receiver<ChildEvent>,
stderr_tail: Arc<Mutex<Vec<String>>>,
/// Write end of the child's stdin, for the line-based command channel
/// (see [`ChildProc::send_command`]). `None` once it's been closed.
stdin: Option<ChildStdin>,
stdin: std::process::ChildStdin,
}
impl ChildProc {
/// Spawn `pixelpass <args>` as a child, wiring up the event reader. The
/// `waker` is pinged whenever an event arrives so the UI thread wakes to
/// drain it — this wakes the winit event loop directly (via an
/// `EventLoopProxy`), so it works even when the window is hidden to the tray
/// and no frames are running (egui's own repaint callback would not fire
/// repeatedly in that idle state — see [`super::Waker`]).
pub fn spawn(args: &[String], waker: Waker) -> std::io::Result<Self> {
/// Spawn `pixelpass <args>` as a child, wiring up the event reader. `ctx`
/// is repainted whenever an event arrives so the UI updates live.
pub fn spawn(args: &[String], ctx: egui::Context) -> std::io::Result<Self> {
let exe = std::env::current_exe()?;
let mut child = Command::new(exe)
.args(args)
// Piped so we can send line commands (e.g. `kick <id>`); the host
// only reads it when driven this way (`--output json`).
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()?;
let stdin = child.stdin.take();
let (tx, rx) = std::sync::mpsc::channel();
let stdin = child.stdin.take().expect("stdin piped");
let stdout = child.stdout.take().expect("stdout piped");
std::thread::spawn(move || {
let reader = BufReader::new(stdout);
@@ -109,14 +98,9 @@ impl ChildProc {
if tx.send(ev).is_err() {
break; // app gone
}
waker.wake();
ctx.request_repaint();
}
}
// stdout closed → the child has exited (player closed, connection
// ended, or a failed launch). Wake once more so the UI reaps it and
// clears the "running" view, even if no final event was emitted and
// the window is hidden to the tray.
waker.wake();
});
let stderr_tail = Arc::new(Mutex::new(Vec::<String>::new()));
@@ -135,64 +119,55 @@ impl ChildProc {
});
Ok(Self {
child: Some(child),
child,
rx,
stderr_tail,
stdin,
})
}
/// Send one newline-terminated command to the child over its stdin (the
/// host parses these as `kick <endpoint-id>`). Best-effort: a closed pipe
/// (child already gone) just drops the command.
pub fn send_command(&mut self, cmd: &str) {
let Some(stdin) = self.stdin.as_mut() else {
return;
};
if let Err(e) = writeln!(stdin, "{cmd}") {
tracing::warn!("failed to send command to host child: {e}");
self.stdin = None; // pipe is dead; stop trying
}
}
/// Whether the child is still running.
pub fn is_alive(&mut self) -> bool {
matches!(self.child.as_mut().map(Child::try_wait), Some(Ok(None)))
matches!(self.child.try_wait(), Ok(None))
}
/// The last captured stderr lines, joined — for error display.
pub fn stderr_tail(&self) -> String {
self.stderr_tail.lock().unwrap().join("\n")
}
/// Send a newline-terminated command to the child.
pub fn send_command(&mut self, cmd: &str) {
use std::io::Write;
if let Err(e) = writeln!(self.stdin, "{cmd}") {
tracing::warn!("failed to send command to child: {e}");
}
}
/// Gracefully stop the child: SIGINT (so the host runs its ctrl-c teardown
/// — tears down capture, closes the endpoint), with a ~2 s grace period
/// before a hard kill. Idempotent.
pub fn stop(&mut self) {
if matches!(self.child.try_wait(), Ok(Some(_))) {
return; // already exited
}
let _ = kill(Pid::from_raw(self.child.id() as i32), Signal::SIGINT);
for _ in 0..40 {
if matches!(self.child.try_wait(), Ok(Some(_))) {
return;
}
std::thread::sleep(Duration::from_millis(50));
}
let _ = self.child.kill();
let _ = self.child.wait();
}
}
impl Drop for ChildProc {
fn drop(&mut self) {
// Leaving a host/viewer screen, or closing the window, must not orphan
// a live child — but it must also not *block*. eframe runs this drop
// synchronously while it destroys the window, so a grace-period wait
// here freezes the window mid-close: the first click looks like it did
// nothing (the stream just drops) and the window only goes away on a
// second click. So SIGINT now — synchronously, so the host always gets
// its ctrl-c teardown (capture down, endpoint closed) even if we exit
// right after — then reap on a detached thread instead of waiting.
let Some(mut child) = self.child.take() else {
return;
};
if matches!(child.try_wait(), Ok(Some(_))) {
return; // already exited; nothing to signal or reap
}
let _ = kill(Pid::from_raw(child.id() as i32), Signal::SIGINT);
std::thread::spawn(move || {
for _ in 0..40 {
if matches!(child.try_wait(), Ok(Some(_))) {
return;
}
std::thread::sleep(Duration::from_millis(50));
}
let _ = child.kill();
let _ = child.wait();
});
// Closing the window (dropping the app, hence the session) must not
// orphan a live host child streaming to viewers.
self.stop();
}
}
@@ -247,7 +222,7 @@ mod tests {
}
#[test]
fn viewer_join_leave_round_trip() {
fn viewer_events_round_trips() {
assert!(matches!(
parse(Event::ViewerJoined { id: "nodeXYZ", active: 2, max: 4 }),
ChildEvent::ViewerJoined { id, active: 2, max: 4 } if id == "nodeXYZ"
@@ -261,20 +236,12 @@ mod tests {
#[test]
fn capture_state_round_trips() {
assert!(matches!(
parse(Event::Capture {
state: EmitState::Started
}),
ChildEvent::Capture {
state: CaptureState::Started
}
parse(Event::Capture { state: EmitState::Started }),
ChildEvent::Capture { state: CaptureState::Started }
));
assert!(matches!(
parse(Event::Capture {
state: EmitState::Stopped
}),
ChildEvent::Capture {
state: CaptureState::Stopped
}
parse(Event::Capture { state: EmitState::Stopped }),
ChildEvent::Capture { state: CaptureState::Stopped }
));
}
-88
View File
@@ -1,88 +0,0 @@
//! Share-code wrapping: carrying the host's stable friend id alongside the
//! one-shot video ticket.
//!
//! A bare video ticket identifies only the host's *ephemeral* video endpoint,
//! so two people who meet over one can't learn each other's stable friend id —
//! the thing the friends system needs. The GUI host therefore wraps its ticket
//! with its control-plane [`EndpointId`]; the viewer unwraps it, dials the
//! video ticket as before, and now also knows who to befriend (and announces
//! itself back over the control plane so the host learns the viewer in turn).
//!
//! Format: `pixelpassF1:<host-control-id>.<bare-ticket>`. Both the id and the
//! ticket are base32 text with no `.`, so a single `.` separator is
//! unambiguous. [`unwrap`] is lenient: anything without the prefix is treated
//! as a bare ticket, so a plain CLI ticket pasted into the GUI still works (it
//! just offers no friend option). The host name isn't carried here — the
//! viewer's announcement triggers a name exchange over the control plane.
use std::str::FromStr;
use iroh::EndpointId;
/// Prefix marking a wrapped friend code. The `F1` is the wrap-format version,
/// bumped if the layout ever changes.
const MAGIC: &str = "pixelpassF1:";
/// Wrap a bare ticket with the host's control id, for display/copy/QR.
pub fn wrap(host_id: EndpointId, ticket: &str) -> String {
format!("{MAGIC}{host_id}.{ticket}")
}
/// Split an input into `(host control id if it was a wrapped code, bare
/// ticket)`. A bare or unrecognised input yields `(None, trimmed input)` so the
/// viewer path stays identical to before for plain tickets.
pub fn unwrap(code: &str) -> (Option<EndpointId>, String) {
let code = code.trim();
if let Some(rest) = code.strip_prefix(MAGIC)
&& let Some((id_str, ticket)) = rest.split_once('.')
&& let Ok(id) = EndpointId::from_str(id_str)
&& !ticket.is_empty()
{
return (Some(id), ticket.to_string());
}
(None, code.to_string())
}
#[cfg(test)]
mod tests {
use super::*;
fn sample_id() -> EndpointId {
iroh::SecretKey::generate().public()
}
#[test]
fn wrap_unwrap_round_trips() {
let id = sample_id();
let ticket = "endpointaabwxjexzensznfvuudiapn5tyzws3angd2merarm";
let code = wrap(id, ticket);
let (got_id, got_ticket) = unwrap(&code);
assert_eq!(got_id, Some(id));
assert_eq!(got_ticket, ticket);
}
#[test]
fn bare_ticket_passes_through() {
let ticket = "endpointaabwxjexzensznfvuudiapn5tyzws3angd2merarm";
let (id, got) = unwrap(ticket);
assert_eq!(id, None);
assert_eq!(got, ticket);
}
#[test]
fn trims_surrounding_whitespace() {
let ticket = "endpointaabwxjex";
let (id, got) = unwrap(&format!(" {} ", wrap(sample_id(), ticket)));
assert!(id.is_some());
assert_eq!(got, ticket);
}
#[test]
fn malformed_wrapped_code_falls_back_to_bare() {
// Prefix present but the id isn't a valid EndpointId → treat the whole
// thing as a (doomed) bare ticket rather than panicking.
let (id, got) = unwrap("pixelpassF1:not-an-id.endpointaa");
assert_eq!(id, None);
assert_eq!(got, "pixelpassF1:not-an-id.endpointaa");
}
}
+130 -1984
View File
File diff suppressed because it is too large Load Diff
-299
View File
@@ -1,299 +0,0 @@
//! The always-on friends presence service.
//!
//! A control-plane iroh endpoint ([`endpoint::bind_control`]) that lives for the
//! whole GUI session on its own thread with a current-thread tokio runtime — the
//! GUI is a synchronous winit/egui loop, so iroh's async work can't run on it
//! (the same reason [`super::tray`] has its own thread + runtime).
//!
//! Inbound control messages are forwarded over a std mpsc channel the UI drains
//! each [`super::PixelPassApp::tick`]; the [`Waker`] is pinged on arrival so a
//! message wakes the loop even while the window is hidden to the tray — the same
//! trick the headless-child reader uses.
use std::sync::Arc;
use std::sync::mpsc::{self, Receiver};
use std::thread;
use iroh::{Endpoint, EndpointId};
use tokio::sync::mpsc as tmpsc;
use super::Waker;
use crate::common::{
control::{self, ControlMsg, Inbound},
endpoint, identity,
};
/// A command the UI hands the presence service over [`PresenceHandle`].
enum Command {
/// Deliver one message, once, fire-and-forget (friend request/accept/decline
/// and the presence `Hello`). A failure is logged, not retried.
Send { peer: EndpointId, msg: ControlMsg },
/// Begin — or replace — a share campaign: push `msg` (a
/// [`ControlMsg::ShareCode`]) to every peer in `peers`, retrying the ones
/// that are offline until they're reached or the campaign is stopped. Each
/// success emits a [`PresenceEvent::ShareDelivered`]. Replaces any campaign
/// already running (a fresh host session supersedes the previous code).
StartShare {
msg: ControlMsg,
peers: Vec<EndpointId>,
},
/// Stop the active share campaign — the host stopped or left the screen, so
/// the perishable code is no longer valid and offline friends shouldn't keep
/// being chased.
StopShare,
}
/// Something the service surfaces to the UI, drained each tick.
pub enum PresenceEvent {
/// A control message arrived from a peer.
Message(Inbound),
/// A share-campaign code reached `peer` (its ACK came back). Lets the host
/// screen flip that friend's row from "retrying" to "delivered."
ShareDelivered { peer: EndpointId },
}
/// How long to wait before re-attempting delivery to friends who were offline
/// on the previous round of a share campaign.
const SHARE_RETRY: std::time::Duration = std::time::Duration::from_secs(5);
/// Handle the GUI holds for the presence service. Dropping it doesn't stop the
/// service (the thread is detached; the endpoint closes when the process exits)
/// — it just stops the UI from draining inbound messages.
pub struct PresenceHandle {
/// Our stable control-plane id — what friends know us by, and what we embed
/// in a wrapped share code so a viewer can find us.
id: EndpointId,
/// Service events (inbound messages + share receipts), drained by
/// [`PresenceHandle::drain`] each tick.
rx: Receiver<PresenceEvent>,
/// Commands handed to the service thread. Unbounded tokio sender so the sync
/// UI can enqueue without blocking or being inside the runtime.
out_tx: tmpsc::UnboundedSender<Command>,
}
impl PresenceHandle {
/// Our stable control-plane id.
pub fn id(&self) -> EndpointId {
self.id
}
/// Pull every service event received since the last call. Collected by the
/// caller so it can take `&mut self` while handling them.
pub fn drain(&self) -> Vec<PresenceEvent> {
std::iter::from_fn(|| self.rx.try_recv().ok()).collect()
}
/// Enqueue a one-shot message for delivery to `peer`. Fire-and-forget from
/// the UI's view; the service connects, delivers, and logs a failure. A send
/// error here only means the service thread is gone.
pub fn send(&self, peer: EndpointId, msg: ControlMsg) {
self.command(Command::Send { peer, msg });
}
/// Begin (or replace) a share campaign pushing `msg` to `peers`, retrying
/// offline friends until [`PresenceHandle::stop_share`] or the next call.
pub fn start_share(&self, msg: ControlMsg, peers: Vec<EndpointId>) {
self.command(Command::StartShare { msg, peers });
}
/// Stop the active share campaign (host stopped — the code is now stale).
pub fn stop_share(&self) {
self.command(Command::StopShare);
}
fn command(&self, cmd: Command) {
if self.out_tx.send(cmd).is_err() {
tracing::warn!("presence: service thread gone; dropping command");
}
}
}
/// Start the presence service. Returns `None` if the persistent identity can't
/// be loaded — the GUI then simply runs without friends features rather than
/// refusing to start. The endpoint binds asynchronously on the spawned thread;
/// our id is known immediately because it derives from the saved key, so we can
/// fail-fast and log it without waiting on the relay handshake.
pub fn start(waker: Waker, relay: Option<String>) -> Option<PresenceHandle> {
let id: EndpointId = match identity::load_or_create() {
Ok(key) => key.public(),
Err(e) => {
tracing::warn!("presence: no identity, friends features disabled: {e:#}");
return None;
}
};
tracing::info!(%id, "presence: starting control service");
let (tx, rx) = mpsc::channel::<PresenceEvent>();
let (out_tx, out_rx) = tmpsc::unbounded_channel::<Command>();
thread::Builder::new()
.name("pixelpass-presence".into())
.spawn(move || run(relay, id, tx, out_rx, waker))
.map_err(|e| tracing::warn!("presence: could not spawn service thread: {e}"))
.ok()?;
Some(PresenceHandle { id, rx, out_tx })
}
/// Thread body: a current-thread tokio runtime that binds the control endpoint,
/// runs the accept loop, bridges inbound messages to the UI channel, and
/// delivers outbound messages the UI enqueues.
fn run(
relay: Option<String>,
id: EndpointId,
tx: mpsc::Sender<PresenceEvent>,
mut out_rx: tmpsc::UnboundedReceiver<Command>,
waker: Waker,
) {
let rt = match tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
{
Ok(rt) => rt,
Err(e) => {
tracing::error!("presence: failed to build runtime: {e}");
return;
}
};
rt.block_on(async move {
let ep = match endpoint::bind_control(relay.as_deref()).await {
Ok(ep) => ep,
Err(e) => {
tracing::error!("presence: failed to bind control endpoint: {e:#}");
return;
}
};
tracing::info!(%id, "presence: control endpoint online");
// One async→sync bridge for *everything* the UI sees: every producer
// (the accept loop and the share campaign) pushes a `PresenceEvent` into
// `ui_tx`; this task drains it onto the std channel and wakes the loop so
// the event lands even while the window is hidden to the tray.
let (ui_tx, mut ui_rx) = tmpsc::channel::<PresenceEvent>(64);
let forward = tokio::spawn(async move {
while let Some(event) = ui_rx.recv().await {
if tx.send(event).is_err() {
break; // UI gone
}
waker.wake();
}
});
// Wrap inbound control messages as events and feed the bridge.
let (itx, mut irx) = tmpsc::channel::<Inbound>(32);
let inbound_ui = ui_tx.clone();
let inbound = tokio::spawn(async move {
while let Some(msg) = irx.recv().await {
if inbound_ui.send(PresenceEvent::Message(msg)).await.is_err() {
break;
}
}
});
// Handle UI commands: one-shot sends each on their own task, and a single
// abortable share campaign (StartShare replaces it, StopShare cancels it).
let cmd_ep = ep.clone();
let commands = tokio::spawn(async move {
let mut share: Option<tokio::task::JoinHandle<()>> = None;
while let Some(cmd) = out_rx.recv().await {
match cmd {
Command::Send { peer, msg } => {
let ep = cmd_ep.clone();
tokio::spawn(async move {
if let Err(e) = control::send(&ep, peer, &msg).await {
tracing::warn!(%peer, "presence: outbound send failed: {e:#}");
}
});
}
Command::StartShare { msg, peers } => {
if let Some(t) = share.take() {
t.abort();
}
let ep = cmd_ep.clone();
let ui = ui_tx.clone();
share = Some(tokio::spawn(run_share(ep, msg, peers, ui)));
}
Command::StopShare => {
if let Some(t) = share.take() {
t.abort();
}
}
}
}
});
control::serve(ep, itx).await;
forward.abort();
inbound.abort();
commands.abort();
});
}
/// Push `msg` to every peer in `peers`, retrying the ones that are offline every
/// [`SHARE_RETRY`] until all are delivered (or the task is aborted by a
/// StartShare/StopShare). Emits one [`PresenceEvent::ShareDelivered`] per peer
/// the moment its ACK comes back — that ACK *is* the delivery signal.
///
/// Each round fires all still-pending peers **concurrently**, so a single
/// offline friend's ~10s connect timeout doesn't serialise the whole round
/// (which it did when peers were tried one at a time).
async fn run_share(
ep: Endpoint,
msg: ControlMsg,
mut pending: Vec<EndpointId>,
ui: tmpsc::Sender<PresenceEvent>,
) {
// The code is immutable for the campaign's life; share it across the
// per-peer tasks via an `Arc` rather than re-cloning the payload each round.
let msg = Arc::new(msg);
while !pending.is_empty() {
let mut round = tokio::task::JoinSet::new();
for peer in pending {
let ep = ep.clone();
let msg = Arc::clone(&msg);
round.spawn(async move {
match control::send(&ep, peer, &msg).await {
Ok(()) => (peer, true),
Err(e) => {
tracing::debug!(%peer, "presence: share not yet delivered: {e:#}");
(peer, false)
}
}
});
}
let mut still = Vec::new();
while let Some(joined) = round.join_next().await {
let (peer, delivered) = match joined {
Ok(outcome) => outcome,
// A send task panicking is unexpected; log and drop that peer
// from the campaign rather than abort the whole round. (A
// campaign-level abort drops this future entirely — we never
// observe that as a JoinError here.)
Err(e) => {
tracing::warn!("presence: share task failed: {e}");
continue;
}
};
if delivered {
tracing::info!(%peer, "presence: shared code delivered");
if ui
.send(PresenceEvent::ShareDelivered { peer })
.await
.is_err()
{
return; // UI gone — nothing left to report to
}
} else {
still.push(peer);
}
}
if still.is_empty() {
break;
}
pending = still;
tokio::time::sleep(SHARE_RETRY).await;
}
tracing::info!("presence: share campaign complete");
}
-456
View File
@@ -1,456 +0,0 @@
//! User-customisable colour themes for the GUI.
//!
//! A theme is a small, curated *semantic* palette — backgrounds, text, an
//! accent, and the handful of status colours the app uses (streaming, waiting,
//! success, warning, error). That's deliberately a fixed set rather than a
//! passthrough of every [`egui::Visuals`] field: it's easy to author by hand,
//! covers the whole look of the app, and stays stable across egui upgrades.
//!
//! Themes serialise to TOML with colours as `#rrggbb` hex strings. Three
//! themes ship built in; users drop their own `*.toml` files in
//! `~/.config/pixelpass/themes/` (or save one from the in-app editor) and they
//! show up alongside the built-ins. A user file whose `name` matches a built-in
//! overrides it.
use std::path::PathBuf;
use anyhow::{Context, Result};
use directories::ProjectDirs;
use eframe::egui::{self, Color32};
use serde::{Deserialize, Serialize};
/// One colour theme: a curated semantic palette.
///
/// `#[serde(default)]` on the container means any field missing from a TOML
/// file falls back to the corresponding field of [`Theme::default`] (the
/// built-in Default Dark), so a partial or hand-trimmed file still loads.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(default)]
pub struct Theme {
/// Display name, shown in the picker and used as the file stem on save.
pub name: String,
/// Base egui defaults to start from before applying the palette overrides.
pub dark: bool,
// ── chrome ────────────────────────────────────────────────────────
/// Window background.
#[serde(with = "hex")]
pub window_bg: Color32,
/// Panel / frame background.
#[serde(with = "hex")]
pub panel_bg: Color32,
/// Text-input and read-only field background (the ticket box, etc.).
#[serde(with = "hex")]
pub input_bg: Color32,
/// Primary text.
#[serde(with = "hex")]
pub text: Color32,
/// Secondary / de-emphasised text (hints, the version line).
#[serde(with = "hex")]
pub weak_text: Color32,
/// Accent: selection, hyperlinks, and the active/pressed widget fill.
#[serde(with = "hex")]
pub accent: Color32,
/// Button (and other interactive widget) resting background.
#[serde(with = "hex")]
pub button_bg: Color32,
/// Button background on hover.
#[serde(with = "hex")]
pub button_hovered: Color32,
// ── semantic status colours ───────────────────────────────────────
/// "● Streaming" indicator.
#[serde(with = "hex")]
pub streaming: Color32,
/// "● Waiting for viewers…" indicator.
#[serde(with = "hex")]
pub waiting: Color32,
/// Success notes, e.g. "✓ Copied to clipboard".
#[serde(with = "hex")]
pub success: Color32,
/// Non-fatal warnings, e.g. a host-full refusal.
#[serde(with = "hex")]
pub warning: Color32,
/// Errors.
#[serde(with = "hex")]
pub error: Color32,
}
impl Default for Theme {
fn default() -> Self {
default_dark()
}
}
impl Theme {
/// Build the egui [`Visuals`](egui::Visuals) this theme describes. Starts
/// from egui's dark or light defaults (so anything the palette doesn't name
/// stays sensible) and overrides the curated fields.
pub fn visuals(&self) -> egui::Visuals {
use egui::{Stroke, Visuals};
let mut v = if self.dark {
Visuals::dark()
} else {
Visuals::light()
};
v.dark_mode = self.dark;
v.window_fill = self.window_bg;
v.panel_fill = self.panel_bg;
v.faint_bg_color = self.panel_bg;
v.extreme_bg_color = self.input_bg;
v.override_text_color = Some(self.text);
// `.weak()` text resolves via `weak_text_color()`, which derives from
// `text` unless this is set — so without it the weak-text field is dead.
v.weak_text_color = Some(self.weak_text);
v.hyperlink_color = self.accent;
v.error_fg_color = self.error;
v.warn_fg_color = self.warning;
// A translucent accent reads well as a selection highlight on either a
// light or dark base.
v.selection.bg_fill =
Color32::from_rgba_unmultiplied(self.accent.r(), self.accent.g(), self.accent.b(), 96);
v.selection.stroke = Stroke::new(1.0, self.accent);
let text_stroke = Stroke::new(1.0, self.text);
let weak_stroke = Stroke::new(1.0, self.weak_text);
v.widgets.noninteractive.bg_fill = self.panel_bg;
v.widgets.noninteractive.weak_bg_fill = self.panel_bg;
v.widgets.noninteractive.fg_stroke = weak_stroke;
v.widgets.inactive.bg_fill = self.button_bg;
v.widgets.inactive.weak_bg_fill = self.button_bg;
v.widgets.inactive.fg_stroke = text_stroke;
v.widgets.hovered.bg_fill = self.button_hovered;
v.widgets.hovered.weak_bg_fill = self.button_hovered;
v.widgets.hovered.fg_stroke = text_stroke;
v.widgets.active.bg_fill = self.accent;
v.widgets.active.weak_bg_fill = self.accent;
v.widgets.active.fg_stroke = text_stroke;
v
}
}
// ── built-in themes ───────────────────────────────────────────────────────
/// Names of the built-in themes, in picker order.
pub const BUILTIN_NAMES: [&str; 3] = ["Default Dark", "Catppuccin Mocha", "Catppuccin Latte"];
/// Parse a built-in's hex literal, panicking on a typo (these are compile-time
/// constants we control, so a bad value is a bug, not user input).
fn c(hex: &str) -> Color32 {
parse_hex(hex).expect("built-in theme hex is valid")
}
/// The default theme — a neutral dark palette. Also [`Theme::default`].
pub fn default_dark() -> Theme {
Theme {
name: "Default Dark".to_string(),
dark: true,
window_bg: c("#1b1b1f"),
panel_bg: c("#242429"),
input_bg: c("#141417"),
text: c("#e6e6ea"),
weak_text: c("#a0a0a8"),
accent: c("#5aa0f2"),
button_bg: c("#33333a"),
button_hovered: c("#44444d"),
streaming: c("#6fdc8c"),
waiting: c("#f2c14e"),
success: c("#6fdc8c"),
warning: c("#f0a85a"),
error: c("#f2756f"),
}
}
/// Catppuccin Mocha (dark). <https://github.com/catppuccin/catppuccin>
fn catppuccin_mocha() -> Theme {
Theme {
name: "Catppuccin Mocha".to_string(),
dark: true,
window_bg: c("#1e1e2e"),
panel_bg: c("#181825"),
input_bg: c("#11111b"),
text: c("#cdd6f4"),
weak_text: c("#a6adc8"),
accent: c("#cba6f7"),
button_bg: c("#313244"),
button_hovered: c("#45475a"),
streaming: c("#a6e3a1"),
waiting: c("#f9e2af"),
success: c("#a6e3a1"),
warning: c("#fab387"),
error: c("#f38ba8"),
}
}
/// Catppuccin Latte (light). <https://github.com/catppuccin/catppuccin>
fn catppuccin_latte() -> Theme {
Theme {
name: "Catppuccin Latte".to_string(),
dark: false,
window_bg: c("#eff1f5"),
panel_bg: c("#e6e9ef"),
input_bg: c("#dce0e8"),
text: c("#4c4f69"),
weak_text: c("#6c6f85"),
accent: c("#8839ef"),
button_bg: c("#ccd0da"),
button_hovered: c("#bcc0cc"),
streaming: c("#40a02b"),
waiting: c("#df8e1d"),
success: c("#40a02b"),
warning: c("#fe640b"),
error: c("#d20f39"),
}
}
/// The built-in themes, in [`BUILTIN_NAMES`] order.
pub fn builtins() -> Vec<Theme> {
vec![default_dark(), catppuccin_mocha(), catppuccin_latte()]
}
/// Whether `name` is one of the built-ins (which are read-only — the editor
/// nudges you to save under a new name).
pub fn is_builtin(name: &str) -> bool {
BUILTIN_NAMES.contains(&name)
}
// ── on-disk themes ──────────────────────────────────────────────────────────
/// `~/.config/pixelpass/themes/` (or the XDG equivalent). Not created until a
/// theme is saved.
pub fn themes_dir() -> Result<PathBuf> {
let dirs = ProjectDirs::from("", "", "pixelpass")
.context("could not locate a config directory for pixelpass")?;
Ok(dirs.config_dir().join("themes"))
}
/// Parse every `*.toml` in the themes dir into a [`Theme`]. A file that fails
/// to parse is logged and skipped rather than aborting the whole list, so one
/// bad file can't hide the rest. Returns themes sorted by name.
pub fn list_user_themes() -> Vec<Theme> {
let Ok(dir) = themes_dir() else {
return Vec::new();
};
let Ok(entries) = std::fs::read_dir(&dir) else {
return Vec::new(); // dir doesn't exist yet → no user themes
};
let mut out = Vec::new();
for entry in entries.flatten() {
let path = entry.path();
if path.extension().and_then(|e| e.to_str()) != Some("toml") {
continue;
}
match std::fs::read_to_string(&path) {
Ok(s) => match toml::from_str::<Theme>(&s) {
Ok(mut t) => {
// Fall back to the file stem if the file omits a name.
if t.name.trim().is_empty() {
t.name = path
.file_stem()
.and_then(|s| s.to_str())
.unwrap_or("Unnamed")
.to_string();
}
out.push(t);
}
Err(e) => tracing::warn!("skipping theme {}: {e}", path.display()),
},
Err(e) => tracing::warn!("could not read theme {}: {e}", path.display()),
}
}
out.sort_by_key(|t| t.name.to_lowercase());
out
}
/// Built-ins plus user themes, in picker order: built-ins first (a user file
/// with a matching `name` overrides the built-in's colours in place), then any
/// remaining user themes alphabetically.
pub fn all_themes() -> Vec<Theme> {
let users = list_user_themes();
let mut out: Vec<Theme> = builtins()
.into_iter()
.map(|b| {
users
.iter()
.find(|u| u.name == b.name)
.cloned()
.unwrap_or(b)
})
.collect();
for u in users {
if !is_builtin(&u.name) {
out.push(u);
}
}
out
}
/// The theme with this `name`, or Default Dark if it can't be found (e.g. the
/// config names a theme whose file was deleted).
pub fn load_named(name: &str) -> Theme {
all_themes()
.into_iter()
.find(|t| t.name == name)
.unwrap_or_else(default_dark)
}
/// Write `theme` to `<themes_dir>/<slug>.toml` and return the path. Overwrites
/// an existing file with the same slug (i.e. saving a tweaked theme under the
/// same name updates it in place).
pub fn save_theme(theme: &Theme) -> Result<PathBuf> {
let dir = themes_dir()?;
std::fs::create_dir_all(&dir).with_context(|| format!("failed to create {}", dir.display()))?;
let slug = slugify(&theme.name);
let path = dir.join(format!("{slug}.toml"));
let body = toml::to_string_pretty(theme).context("failed to serialise theme to TOML")?;
let contents = format!(
"# PixelPass theme. Colours are #rrggbb hex strings.\n\
# Edit and re-pick it in Settings, or drop more .toml files in this folder.\n\n\
{body}"
);
std::fs::write(&path, contents)
.with_context(|| format!("failed to write {}", path.display()))?;
Ok(path)
}
/// Lowercase, replace runs of non-alphanumerics with a single hyphen, trim
/// hyphens. Empty input becomes `theme`.
fn slugify(name: &str) -> String {
let mut slug = String::new();
let mut prev_hyphen = false;
for ch in name.trim().chars() {
if ch.is_ascii_alphanumeric() {
slug.push(ch.to_ascii_lowercase());
prev_hyphen = false;
} else if !prev_hyphen {
slug.push('-');
prev_hyphen = true;
}
}
let slug = slug.trim_matches('-').to_string();
if slug.is_empty() {
"theme".to_string()
} else {
slug
}
}
// ── hex colour parsing ────────────────────────────────────────────────────
/// Parse `#rrggbb` into an opaque [`Color32`] (the leading `#` is optional).
/// An 8-digit `#rrggbbaa` is accepted leniently but its alpha is ignored —
/// theme colours are opaque, and `Color32`'s premultiplied storage can't
/// round-trip a straight alpha losslessly anyway. Returns `None` on malformed
/// input.
pub fn parse_hex(s: &str) -> Option<Color32> {
let s = s.trim();
let s = s.strip_prefix('#').unwrap_or(s);
if !matches!(s.len(), 6 | 8) || !s.bytes().all(|b| b.is_ascii_hexdigit()) {
return None;
}
let byte = |i: usize| u8::from_str_radix(&s[i..i + 2], 16).ok();
Some(Color32::from_rgb(byte(0)?, byte(2)?, byte(4)?))
}
/// Format a [`Color32`] as opaque `#rrggbb`.
pub fn to_hex(c: Color32) -> String {
let [r, g, b, _] = c.to_srgba_unmultiplied();
format!("#{r:02x}{g:02x}{b:02x}")
}
/// serde adaptor so `Color32` fields round-trip as hex strings in TOML.
mod hex {
use super::{parse_hex, to_hex};
use eframe::egui::Color32;
use serde::{Deserialize, Deserializer, Serializer, de::Error};
pub fn serialize<S: Serializer>(c: &Color32, s: S) -> Result<S::Ok, S::Error> {
s.serialize_str(&to_hex(*c))
}
pub fn deserialize<'de, D: Deserializer<'de>>(d: D) -> Result<Color32, D::Error> {
let s = String::deserialize(d)?;
parse_hex(&s).ok_or_else(|| {
D::Error::custom(format!(
"invalid hex colour {s:?} (expected #rrggbb or #rrggbbaa)"
))
})
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn hex_round_trips() {
for (input, expect) in [
("#1e1e2e", Color32::from_rgb(0x1e, 0x1e, 0x2e)),
("aabbcc", Color32::from_rgb(0xaa, 0xbb, 0xcc)),
// 8-digit is accepted but the alpha is dropped (opaque rgb).
("#11223344", Color32::from_rgb(0x11, 0x22, 0x33)),
] {
assert_eq!(parse_hex(input).expect("parses"), expect);
}
assert_eq!(to_hex(Color32::from_rgb(0x1e, 0x1e, 0x2e)), "#1e1e2e");
// Opaque colours round-trip exactly.
let c = Color32::from_rgb(0xab, 0xcd, 0xef);
assert_eq!(parse_hex(&to_hex(c)), Some(c));
}
#[test]
fn hex_rejects_garbage() {
for bad in ["", "#fff", "#12345", "nothex", "#gggggg", "#1234567"] {
assert!(parse_hex(bad).is_none(), "{bad:?} should not parse");
}
}
#[test]
fn theme_toml_round_trips() {
let original = catppuccin_mocha();
let toml = toml::to_string_pretty(&original).unwrap();
let parsed: Theme = toml::from_str(&toml).unwrap();
assert_eq!(original, parsed);
// Colours serialise as hex strings, not RGBA tables.
assert!(toml.contains("window_bg = \"#1e1e2e\""), "{toml}");
}
#[test]
fn partial_toml_fills_from_default() {
// Only a name and one colour; everything else must fall back to Default Dark.
let parsed: Theme = toml::from_str("name = \"Partial\"\naccent = \"#ff0000\"").unwrap();
let base = default_dark();
assert_eq!(parsed.name, "Partial");
assert_eq!(parsed.accent, Color32::from_rgb(0xff, 0, 0));
assert_eq!(parsed.window_bg, base.window_bg); // filled from default
assert_eq!(parsed.text, base.text);
}
#[test]
fn slugify_is_filesystem_safe() {
assert_eq!(slugify("Catppuccin Mocha"), "catppuccin-mocha");
assert_eq!(slugify(" My Theme!! "), "my-theme");
assert_eq!(slugify("***"), "theme");
assert_eq!(slugify("Solarized/Dark"), "solarized-dark");
}
#[test]
fn builtins_match_names() {
let names: Vec<String> = builtins().iter().map(|t| t.name.clone()).collect();
let expected: Vec<String> = BUILTIN_NAMES.iter().map(|s| s.to_string()).collect();
assert_eq!(names, expected);
for t in builtins() {
assert!(is_builtin(&t.name));
}
}
}
-257
View File
@@ -1,257 +0,0 @@
//! System-tray (StatusNotifierItem) integration for the GUI.
//!
//! The tray runs on its **own dedicated thread** with its own current-thread
//! tokio runtime, fully decoupled from the winit event loop (which owns the
//! main thread) and from the process-wide `#[tokio::main]` runtime. It talks to
//! the egui app purely over winit's event channel and a status channel:
//!
//! * tray → app: a [`super::UserEvent::Tray`] carrying a [`TrayAction`]
//! (Show / Quit), pushed through the [`winit::event_loop::EventLoopProxy`].
//! Using the proxy (not egui's repaint) is essential: a tray click must
//! wake the winit loop even when the window has been **dropped** (hidden to
//! tray), so the loop can recreate it.
//! * app → tray: [`TrayStatus`] (idle / hosting / viewing), pushed on change.
//!
//! Why a separate thread instead of `Handle::current().spawn`: updating the
//! tray from the egui thread would need `block_on`, which panics when called
//! from inside the running runtime. Keeping ksni's async wholly on its own
//! runtime sidesteps that and keeps the frame loop non-blocking.
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use ksni::TrayMethods;
use winit::event_loop::EventLoopProxy;
use super::UserEvent;
/// What the user picked from the tray icon or its menu (tray thread → app),
/// delivered as a [`UserEvent::Tray`].
pub enum TrayAction {
/// Left-click, or the "Show window" item: bring the window back.
Show,
/// The "Quit" item: really exit (the close button only hides to tray).
Quit,
}
/// What the tray icon's tooltip/menu reflect (app → tray thread).
#[derive(Clone, Copy, PartialEq, Eq)]
pub enum TrayStatus {
Idle,
Hosting { active: u32, max: u32 },
Viewing,
}
fn status_text(status: TrayStatus) -> String {
match status {
TrayStatus::Idle => "Idle".to_string(),
TrayStatus::Hosting { active, max } => {
format!("Hosting — {active} of {max} viewer(s) connected")
}
TrayStatus::Viewing => "Viewing a stream".to_string(),
}
}
/// Handle held by the egui app for the lifetime of the window. Dropping it
/// closes the app→tray channel, which ends the tray thread and removes the icon.
pub struct TrayHandle {
status_tx: tokio::sync::mpsc::UnboundedSender<TrayStatus>,
/// Set true once the tray actually registered with a StatusNotifier host.
/// The app must not divert the window's close to a tray that never appeared.
registered: Arc<AtomicBool>,
/// Last status pushed, so we don't spam D-Bus with no-op updates.
last_sent: Option<TrayStatus>,
}
impl TrayHandle {
/// Whether a system tray is actually showing our icon. Until this is true,
/// hiding the window would strand it with no way back.
pub fn registered(&self) -> bool {
self.registered.load(Ordering::Acquire)
}
/// Push a status change to the tray, deduped against the last one sent.
pub fn set_status(&mut self, status: TrayStatus) {
if self.last_sent != Some(status) {
let _ = self.status_tx.send(status);
self.last_sent = Some(status);
}
}
}
struct PixelPassTray {
status: TrayStatus,
/// ARGB pixmap, so the icon shows even where the themed "pixelpass" name
/// can't be resolved (e.g. running the dev binary before `make install`).
icon: Vec<ksni::Icon>,
/// Wakes the winit loop and delivers the action — works even when the
/// window has been dropped to the tray (no egui frame is running then).
proxy: EventLoopProxy<UserEvent>,
/// Shared with [`TrayHandle`]; kept in sync with the watcher's presence via
/// the `watcher_online`/`watcher_offline` callbacks so the app never diverts
/// a close to a tray that has since disappeared.
registered: Arc<AtomicBool>,
}
impl PixelPassTray {
fn notify(&self, action: TrayAction) {
let _ = self.proxy.send_event(UserEvent::Tray(action));
}
}
impl ksni::Tray for PixelPassTray {
fn id(&self) -> String {
"pixelpass".to_string()
}
fn title(&self) -> String {
"PixelPass".to_string()
}
// Themed icon (matches the installed hicolor/scalable/apps/pixelpass.svg);
// icon_pixmap below is the always-works fallback.
fn icon_name(&self) -> String {
"pixelpass".to_string()
}
fn icon_pixmap(&self) -> Vec<ksni::Icon> {
self.icon.clone()
}
fn status(&self) -> ksni::Status {
ksni::Status::Active
}
fn tool_tip(&self) -> ksni::ToolTip {
ksni::ToolTip {
title: "PixelPass".to_string(),
description: status_text(self.status),
icon_name: "pixelpass".to_string(),
icon_pixmap: Vec::new(),
}
}
fn activate(&mut self, _x: i32, _y: i32) {
self.notify(TrayAction::Show);
}
/// The StatusNotifierWatcher came back (e.g. the panel restarted). Mark the
/// tray live again so close-to-tray can resume hiding the window.
fn watcher_online(&self) {
self.registered.store(true, Ordering::Release);
}
/// The watcher went away (panel restart, tray plugin disabled, …). Clear the
/// flag so a subsequent close quits normally instead of destroying the window
/// into a tray that no longer exists, and force the window back now in case
/// it was already hidden (otherwise it'd be stranded with no way to restore).
/// Returning `true` keeps the service alive so it re-registers if the watcher
/// returns.
fn watcher_offline(&self, reason: ksni::OfflineReason) -> bool {
tracing::warn!("tray: StatusNotifierWatcher offline ({reason:?}); restoring window");
self.registered.store(false, Ordering::Release);
self.notify(TrayAction::Show);
true
}
fn menu(&self) -> Vec<ksni::MenuItem<Self>> {
use ksni::menu::{MenuItem, StandardItem};
vec![
// Non-clickable status line.
StandardItem {
label: status_text(self.status),
enabled: false,
..Default::default()
}
.into(),
MenuItem::Separator,
StandardItem {
label: "Show window".to_string(),
activate: Box::new(|t: &mut Self| t.notify(TrayAction::Show)),
..Default::default()
}
.into(),
StandardItem {
label: "Quit PixelPass".to_string(),
icon_name: "application-exit".to_string(),
activate: Box::new(|t: &mut Self| t.notify(TrayAction::Quit)),
..Default::default()
}
.into(),
]
}
}
/// Decode the embedded PNG (RGBA) and convert to the ARGB pixmap ksni wants.
/// Reuses eframe's PNG decoder so we don't take a direct `image` dependency.
fn load_icon() -> Option<Vec<ksni::Icon>> {
let icon =
eframe::icon_data::from_png_bytes(include_bytes!("../../assets/pixelpass-256.png")).ok()?;
let mut data = icon.rgba; // RGBA8, row-major
for px in data.chunks_exact_mut(4) {
px.rotate_right(1); // [r,g,b,a] -> [a,r,g,b], network byte order
}
Some(vec![ksni::Icon {
width: icon.width as i32,
height: icon.height as i32,
data,
}])
}
/// Start the tray on its own thread. Returns a handle for the app to drive it,
/// or `None` if the icon couldn't be decoded or the thread couldn't spawn (in
/// which case the GUI simply runs without a tray — close behaves as before).
pub fn start(proxy: EventLoopProxy<UserEvent>) -> Option<TrayHandle> {
let icon = load_icon()?;
let (status_tx, mut status_rx) = tokio::sync::mpsc::unbounded_channel::<TrayStatus>();
let registered = Arc::new(AtomicBool::new(false));
let registered_thread = registered.clone();
std::thread::Builder::new()
.name("pixelpass-tray".to_string())
.spawn(move || {
let rt = match tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
{
Ok(rt) => rt,
Err(e) => {
tracing::warn!("tray: could not build runtime: {e}");
return;
}
};
rt.block_on(async move {
let tray = PixelPassTray {
status: TrayStatus::Idle,
icon,
proxy,
registered: registered_thread.clone(),
};
let handle = match tray.spawn().await {
Ok(handle) => handle,
Err(e) => {
// No StatusNotifier host (no system tray) — degrade
// gracefully: the window keeps its normal close.
tracing::warn!("tray: not available, running without it: {e}");
return;
}
};
registered_thread.store(true, Ordering::Release);
// Apply status changes until the app drops its sender (on quit),
// which ends this loop, the runtime, the thread, and the icon.
while let Some(status) = status_rx.recv().await {
let _ = handle
.update(move |t: &mut PixelPassTray| t.status = status)
.await;
}
});
})
.ok()?;
Some(TrayHandle {
status_tx,
registered,
last_sent: None,
})
}
+32 -34
View File
@@ -52,8 +52,9 @@ impl Routing {
let pid = std::process::id();
let sink_name = format!("pixelpass_capture_{pid}");
let sink_module = load_module(&["module-null-sink", &format!("sink_name={sink_name}")])
.context("failed to load module-null-sink")?;
let sink_module =
load_module(&["module-null-sink", &format!("sink_name={sink_name}")])
.context("failed to load module-null-sink")?;
// 20ms loopback latency keeps the mirrored audio tight; pactl's
// default of 200ms is enough to be perceptible.
@@ -140,10 +141,7 @@ impl Routing {
/// Stop the stream router (if any), then unload loopback (if still
/// loaded), then unload the null-sink. Order matters: PipeWire can
/// leave zombie links if you destroy a sink with active inputs.
///
/// Every step is a `take()`, so this is idempotent — `Drop` calls it again
/// as a backstop and the second run is a no-op.
fn cleanup(&mut self) {
pub fn shutdown(mut self) {
if let Some(router) = self.stream_router.take() {
router.shutdown();
}
@@ -157,17 +155,22 @@ impl Routing {
unload_module(id);
}
}
/// Consume the routing and tear it all down now. `Drop` is the backstop;
/// the real work lives in [`cleanup`](Self::cleanup).
pub fn shutdown(mut self) {
self.cleanup();
}
}
impl Drop for Routing {
fn drop(&mut self) {
self.cleanup();
if let Some(router) = self.stream_router.take() {
router.shutdown();
}
if let Some(task) = self.event_task.take() {
task.abort();
}
if let Some(id) = self.loopback_module.lock().unwrap().take() {
unload_module(id);
}
if let Some(id) = self.sink_module.take() {
unload_module(id);
}
}
}
@@ -204,9 +207,7 @@ fn parse_sink_inputs(stdout: &[u8]) -> Result<Vec<App>> {
serde_json::from_slice(stdout).context("pactl returned unparseable JSON")?;
let mut counts: BTreeMap<String, u32> = BTreeMap::new();
for entry in entries {
let Some(name) = entry.properties.application_name else {
continue;
};
let Some(name) = entry.properties.application_name else { continue };
let trimmed = name.trim();
if trimmed.is_empty() {
continue;
@@ -360,14 +361,16 @@ fn run_router(
) -> Result<()> {
use pipewire::{self as pw, types::ObjectType};
let main_loop =
pw::main_loop::MainLoopRc::new(None).context("pw main loop construction failed")?;
let context =
pw::context::ContextRc::new(&main_loop, None).context("pw context construction failed")?;
let main_loop = pw::main_loop::MainLoopRc::new(None)
.context("pw main loop construction failed")?;
let context = pw::context::ContextRc::new(&main_loop, None)
.context("pw context construction failed")?;
let core = context
.connect_rc(None)
.context("pw core connect failed (is the daemon running?)")?;
let registry = core.get_registry_rc().context("pw get_registry failed")?;
let registry = core
.get_registry_rc()
.context("pw get_registry failed")?;
let state = Rc::new(RefCell::new(RouterState {
sink_serial: None,
@@ -408,21 +411,20 @@ fn run_router(
let _reg_listener = registry
.add_listener_local()
.global(move |obj| {
let Some(reg) = registry_weak.upgrade() else {
return;
};
let Some(reg) = registry_weak.upgrade() else { return };
match obj.type_ {
ObjectType::Node => {
let Some(props) = obj.props.as_ref() else {
return;
};
let Some(props) = obj.props.as_ref() else { return };
if props.get("node.name") == Some(sink_name_owned.as_str()) {
if let Some(serial) = props
.get("object.serial")
.and_then(|s| s.parse::<u32>().ok())
{
state_for_reg.borrow_mut().sink_serial = Some(serial);
tracing::info!(serial, "audio routing: pixelpass sink registered");
tracing::info!(
serial,
"audio routing: pixelpass sink registered"
);
try_flush(&state_for_reg, &event_tx_for_reg);
}
return;
@@ -430,9 +432,7 @@ fn run_router(
if props.get("media.class") != Some("Stream/Output/Audio") {
return;
}
let Some(app) = props.get("application.name") else {
return;
};
let Some(app) = props.get("application.name") else { return };
if !app.eq_ignore_ascii_case(&filter_lower) {
return;
}
@@ -445,9 +445,7 @@ fn run_router(
try_flush(&state_for_reg, &event_tx_for_reg);
}
ObjectType::Metadata => {
let Some(props) = obj.props.as_ref() else {
return;
};
let Some(props) = obj.props.as_ref() else { return };
if props.get("metadata.name") != Some("default") {
return;
}
+43 -111
View File
@@ -7,33 +7,29 @@ mod wayland;
mod x11;
use anyhow::{Result, bail};
use iroh::endpoint::Connection;
use iroh::endpoint::{Connection, presets};
use iroh::{Endpoint, EndpointAddr};
use iroh_tickets::endpoint::EndpointTicket;
use std::collections::HashMap;
use std::time::Duration;
use tokio::io::{AsyncBufReadExt, BufReader};
use tokio::sync::{mpsc, oneshot};
use tokio_util::sync::CancellationToken;
use crate::cli::HostOpts;
use crate::common::{
bandwidth, config, config::BandwidthStatus, deps, display::DisplayServer, endpoint, output,
alpn::ALPN, bandwidth, config, config::BandwidthStatus, deps, display::DisplayServer, output,
signal, tunnel,
};
use self::pipeline::CaptureHandle;
use self::quality::EffectiveQuality;
/// Messages from per-viewer tasks (and the GUI command channel) to the
/// capture supervisor.
// The shared `Viewer` suffix is the point — these are all viewer lifecycle
// messages — so keep the descriptive names.
#[allow(clippy::enum_variant_names)]
/// Messages from per-viewer tasks to the capture supervisor.
enum SupervisorMsg {
/// A new viewer wants in. Supervisor replies with the local capture HTTP
/// port to connect to, or an error string if the host is full or capture
/// spawn failed. `cancel` is the viewer's own token — the supervisor keeps
/// it so a later `KickViewer` can tear this viewer's stream down.
/// A new viewer wants in. Supervisor replies with the local capture
/// HTTP port to connect to, or an error string if the host is full or
/// capture spawn failed.
AddViewer {
id: String,
cancel: CancellationToken,
@@ -42,9 +38,7 @@ enum SupervisorMsg {
/// A viewer's session ended. Supervisor decrements the count and tears
/// down capture if it just hit zero.
RemoveViewer { id: String },
/// Host asked (via the GUI command channel) to disconnect a viewer by
/// endpoint id. Cancels that viewer's token; the normal teardown path then
/// emits the `ViewerLeft`.
/// Request to kick a specific viewer.
KickViewer { id: String },
}
@@ -74,7 +68,10 @@ pub async fn run(opts: HostOpts) -> Result<()> {
let cancel = signal::install_ctrl_c();
let endpoint = endpoint::bind(opts.relay.as_deref()).await?;
let endpoint = Endpoint::builder(presets::N0)
.alpns(vec![ALPN.to_vec()])
.bind()
.await?;
// Relay-only ticket: wait for the home relay to connect, then keep only
// the endpoint id + relay URL and drop the direct IP candidates. The relay
@@ -124,14 +121,16 @@ pub async fn run(opts: HostOpts) -> Result<()> {
sup_rx,
));
// Command channel for the GUI front-end: read `kick <endpoint-id>` lines
// off stdin. Only when machine-driven (`--output json`) — a human host has
// nothing to type here, and we don't want to swallow terminal input. Runs
// on a plain OS thread (not a tokio task) so a read parked on stdin can't
// hold up runtime shutdown on Ctrl+C; the thread dies with the process.
if output::json_enabled() {
spawn_kick_listener(sup_tx.clone());
}
// Stdin listener for "kick <id>"
let stdin_sup_tx = sup_tx.clone();
tokio::spawn(async move {
let mut lines = BufReader::new(tokio::io::stdin()).lines();
while let Ok(Some(line)) = lines.next_line().await {
if let Some(id) = line.strip_prefix("kick ") {
let _ = stdin_sup_tx.send(SupervisorMsg::KickViewer { id: id.trim().to_string() }).await;
}
}
});
accept_loop(&endpoint, sup_tx.clone(), cancel.clone()).await;
@@ -179,18 +178,11 @@ async fn handle_peer(
cancel: CancellationToken,
) {
let remote = conn.remote_id();
let id = remote.to_string();
// This viewer's own kill switch: the supervisor holds a clone so a `kick`
// can cancel it, and the stream select! below watches it.
let id_str = remote.to_string();
let peer_cancel = CancellationToken::new();
let (reply_tx, reply_rx) = oneshot::channel();
let add = SupervisorMsg::AddViewer {
id: id.clone(),
cancel: peer_cancel.clone(),
reply: reply_tx,
};
if sup_tx.send(add).await.is_err() {
if sup_tx.send(SupervisorMsg::AddViewer { id: id_str.clone(), cancel: peer_cancel.clone(), reply: reply_tx }).await.is_err() {
tracing::warn!(%remote, "supervisor channel closed; dropping peer");
return;
}
@@ -211,7 +203,7 @@ async fn handle_peer(
Ok(s) => s,
Err(e) => {
tracing::warn!(%remote, "accept_bi failed: {e:#}");
let _ = sup_tx.send(SupervisorMsg::RemoveViewer { id }).await;
let _ = sup_tx.send(SupervisorMsg::RemoveViewer { id: id_str }).await;
return;
}
};
@@ -222,7 +214,7 @@ async fn handle_peer(
Ok(t) => t,
Err(e) => {
tracing::warn!(%remote, "connect_to_capture failed: {e:#}");
let _ = sup_tx.send(SupervisorMsg::RemoveViewer { id }).await;
let _ = sup_tx.send(SupervisorMsg::RemoveViewer { id: id_str }).await;
return;
}
};
@@ -242,30 +234,7 @@ async fn handle_peer(
}
eprintln!("[pixelpass] viewer disconnected: {remote}");
let _ = sup_tx.send(SupervisorMsg::RemoveViewer { id }).await;
}
/// Read `kick <endpoint-id>` lines off stdin and forward them to the
/// supervisor. Runs on a detached OS thread (see the call site for why). Ends
/// when stdin hits EOF (the GUI closed the pipe) or the supervisor is gone.
fn spawn_kick_listener(sup_tx: mpsc::Sender<SupervisorMsg>) {
use std::io::BufRead;
std::thread::spawn(move || {
let stdin = std::io::stdin();
for line in stdin.lock().lines().map_while(Result::ok) {
let Some(id) = line.trim().strip_prefix("kick ") else {
continue;
};
let msg = SupervisorMsg::KickViewer {
id: id.trim().to_string(),
};
// blocking_send is valid here: this is a plain thread, not inside
// the tokio runtime. An Err means the supervisor closed — stop.
if sup_tx.blocking_send(msg).is_err() {
break;
}
}
});
let _ = sup_tx.send(SupervisorMsg::RemoveViewer { id: id_str }).await;
}
/// Owns the single shared CaptureHandle and the active viewer count. Spawns
@@ -280,15 +249,12 @@ async fn supervise(
mut rx: mpsc::Receiver<SupervisorMsg>,
) {
let mut handle: Option<CaptureHandle> = None;
// Active viewers, keyed by endpoint id, holding each one's kill switch.
// The count is just `viewers.len()`. (A given endpoint connecting twice is
// a non-case here: each viewer process uses a fresh ephemeral identity.)
let mut count: u32 = 0;
let mut viewers: HashMap<String, CancellationToken> = HashMap::new();
while let Some(msg) = rx.recv().await {
match msg {
SupervisorMsg::AddViewer { id, cancel, reply } => {
let count = viewers.len() as u32;
if count >= max_viewers {
let reason =
format!("host is full ({count} of {max_viewers} viewers connected)");
@@ -314,30 +280,18 @@ async fn supervise(
}
let port = handle.as_ref().expect("handle was just set").local_port();
count += 1;
viewers.insert(id.clone(), cancel);
let active = viewers.len() as u32;
let _ = reply.send(Ok(port));
output::emit(output::Event::ViewerJoined {
id: &id,
active,
max: max_viewers,
});
tracing::info!(active, cap = max_viewers, "viewer joined");
output::emit(output::Event::ViewerJoined { id: &id, active: count, max: max_viewers });
tracing::info!(active = count, cap = max_viewers, "viewer joined");
}
SupervisorMsg::RemoveViewer { id } => {
// A given viewer task only ever sends RemoveViewer once, but the
// map remove is the source of truth either way.
if viewers.remove(&id).is_none() {
continue;
}
let active = viewers.len() as u32;
output::emit(output::Event::ViewerLeft {
id: &id,
active,
max: max_viewers,
});
tracing::info!(active, cap = max_viewers, "viewer left");
if active == 0
viewers.remove(&id);
count = count.saturating_sub(1);
output::emit(output::Event::ViewerLeft { id: &id, active: count, max: max_viewers });
tracing::info!(active = count, cap = max_viewers, "viewer left");
if count == 0
&& let Some(h) = handle.take()
{
tracing::info!("last viewer left — tearing down capture");
@@ -348,14 +302,9 @@ async fn supervise(
}
}
SupervisorMsg::KickViewer { id } => {
match viewers.get(&id) {
// Cancel the viewer's token; its handle_peer select! wakes,
// sends RemoveViewer, and the leave is emitted there.
Some(cancel) => {
tracing::info!(%id, "kicking viewer");
cancel.cancel();
}
None => tracing::debug!(%id, "kick for unknown/already-gone viewer"),
if let Some(cancel) = viewers.get(&id) {
tracing::info!(%id, "kicking viewer");
cancel.cancel();
}
}
}
@@ -382,25 +331,10 @@ fn print_host_banner(
eprintln!("┌─ PixelPass · host ─────────────────────────────────────────");
eprintln!("│ display server : {display:?}");
eprintln!("│ capture : {}", capture_summary(opts));
eprintln!(
"│ quality : {}{}",
quality.label,
quality.dimensions_summary()
);
eprintln!("│ quality : {}{}", quality.label, quality.dimensions_summary());
eprintln!("│ ({})", quality.note);
eprintln!(
"│ hw encode : {}",
if opts.no_hwencode {
"off (software x264)"
} else {
"on (VAAPI H.264)"
}
);
eprintln!(
"│ max viewers : {} ({})",
resolution.value,
resolution.source.label()
);
eprintln!("│ hw encode : {}", if opts.no_hwencode { "off (software x264)" } else { "on (VAAPI H.264)" });
eprintln!("│ max viewers : {} ({})", resolution.value, resolution.source.label());
eprintln!("");
if clipboard_ok {
eprintln!("│ Your share code has been copied to your clipboard.");
@@ -464,9 +398,7 @@ fn resolve_max_viewers(opts: &HostOpts, effective_bitrate: u32) -> MaxViewersRes
let n = bandwidth::recommended_max_viewers(upstream, effective_bitrate);
return MaxViewersResolution {
value: n,
source: MaxViewersSource::BandwidthMeasurement {
safe_mbps: upstream,
},
source: MaxViewersSource::BandwidthMeasurement { safe_mbps: upstream },
};
}
MaxViewersResolution {
+1 -3
View File
@@ -308,9 +308,7 @@ async fn default_audio_monitor() -> Result<String> {
.arg("get-default-sink")
.output()
.await
.context(
"failed to run `pactl get-default-sink` (install pulseaudio-utils or pipewire-pulse)",
)?;
.context("failed to run `pactl get-default-sink` (install pulseaudio-utils or pipewire-pulse)")?;
if !output.status.success() {
bail!(
"pactl get-default-sink failed: {}",
+7 -32
View File
@@ -27,26 +27,10 @@ impl Quality {
/// values and resolves to one of the others at runtime (see [`resolve_auto`]).
fn preset(self) -> Option<Preset> {
let p = match self {
Quality::Source => Preset {
max_height: None,
bitrate: 6000,
framerate: 30,
},
Quality::High => Preset {
max_height: Some(1080),
bitrate: 4000,
framerate: 30,
},
Quality::Medium => Preset {
max_height: Some(720),
bitrate: 2500,
framerate: 30,
},
Quality::Low => Preset {
max_height: Some(480),
bitrate: 1000,
framerate: 30,
},
Quality::Source => Preset { max_height: None, bitrate: 6000, framerate: 30 },
Quality::High => Preset { max_height: Some(1080), bitrate: 4000, framerate: 30 },
Quality::Medium => Preset { max_height: Some(720), bitrate: 2500, framerate: 30 },
Quality::Low => Preset { max_height: Some(480), bitrate: 1000, framerate: 30 },
Quality::Auto => return None,
};
Some(p)
@@ -65,12 +49,7 @@ impl Quality {
/// Fixed presets in descending quality order — Auto walks this to find the
/// best one whose per-viewer bitrate fits the measured upstream budget.
const AUTO_LADDER: [Quality; 4] = [
Quality::Source,
Quality::High,
Quality::Medium,
Quality::Low,
];
const AUTO_LADDER: [Quality; 4] = [Quality::Source, Quality::High, Quality::Medium, Quality::Low];
/// Auto's fallback when there is no usable bandwidth measurement.
const AUTO_FALLBACK: Quality = Quality::Medium;
@@ -167,9 +146,7 @@ fn resolve_auto(safe_mbps: Option<f64>, sizing_viewers: u32) -> (Preset, String,
(
preset,
format!("Auto → {}", chosen.name()),
format!(
"auto: {safe_mbps:.1} Mbps safe ÷ {n} viewer(s) = {budget_mbps:.1} Mbps each"
),
format!("auto: {safe_mbps:.1} Mbps safe ÷ {n} viewer(s) = {budget_mbps:.1} Mbps each"),
)
}
None => {
@@ -177,8 +154,7 @@ fn resolve_auto(safe_mbps: Option<f64>, sizing_viewers: u32) -> (Preset, String,
(
preset,
format!("Auto → {}", AUTO_FALLBACK.name()),
"auto fallback — no bandwidth measurement (run `pixelpass --reconfigure`)"
.to_string(),
"auto fallback — no bandwidth measurement (run `pixelpass --reconfigure`)".to_string(),
)
}
}
@@ -213,7 +189,6 @@ mod tests {
no_hwencode: false,
max_viewers,
interactive: false,
relay: None,
}
}
+2 -5
View File
@@ -130,11 +130,8 @@ async fn run_accept_loop(listener: TcpListener, tx: broadcast::Sender<Arc<Vec<u8
let sock = match listener.accept().await {
Ok((s, _)) => s,
Err(e) => {
// Most accept errors are transient (EMFILE from a brief FD spike,
// EINTR, etc.). Bailing on the first one would kill the entire
// viewer fanout for the rest of the session.
tracing::warn!("capture HTTP accept failed (continuing): {e}");
continue;
tracing::warn!("capture HTTP accept failed: {e}");
return;
}
};
let rx = tx.subscribe();
+5 -15
View File
@@ -26,11 +26,7 @@ pub async fn start(opts: &HostOpts, quality: &EffectiveQuality) -> Result<Captur
.context("could not reach the xdg-desktop-portal ScreenCast interface")?;
let session = proxy.create_session().await?;
let source = if opts.window {
SourceType::Window
} else {
SourceType::Monitor
};
let source = if opts.window { SourceType::Window } else { SourceType::Monitor };
proxy
.select_sources(
&session,
@@ -74,16 +70,10 @@ pub async fn start(opts: &HostOpts, quality: &EffectiveQuality) -> Result<Captur
"do-timestamp=true".to_string(),
];
pipeline::spawn(
opts,
quality,
Some((w as u32, h as u32)),
source_args,
move || {
// Parent no longer needs the pipewire fd — gst inherited its own copy.
let _ = close(raw_fd);
},
)
pipeline::spawn(opts, quality, Some((w as u32, h as u32)), source_args, move || {
// Parent no longer needs the pipewire fd — gst inherited its own copy.
let _ = close(raw_fd);
})
.await
}
+10 -11
View File
@@ -13,7 +13,7 @@ pub async fn run(cli: Cli) -> Result<()> {
let theme = ColorfulTheme::default();
let choice = Select::with_theme(&theme)
.with_prompt("What do you want to do?")
.items([
.items(&[
"Host (share my screen)",
"View (watch someone else's screen)",
])
@@ -112,7 +112,7 @@ fn pick_quality(theme: &ColorfulTheme) -> Result<Quality> {
let choice = Select::with_theme(theme)
.with_prompt("What quality should the viewer(s) get?")
.items(items)
.items(&items)
.default(0)
.interact()?;
@@ -138,7 +138,7 @@ pub async fn run_reconfigure() -> Result<()> {
async fn preflight_if_needed(theme: &ColorfulTheme) {
let mut cfg = config::load().unwrap_or_default();
match cfg.bandwidth.status {
config::BandwidthStatus::Measured | config::BandwidthStatus::Skipped => (),
config::BandwidthStatus::Measured | config::BandwidthStatus::Skipped => return,
config::BandwidthStatus::Unmeasured => {
eprintln!();
eprintln!("First-time setup");
@@ -154,7 +154,7 @@ async fn preflight_if_needed(theme: &ColorfulTheme) {
let Ok(choice) = Select::with_theme(theme)
.with_prompt("What would you like to do?")
.items([
.items(&[
"Run the bandwidth test (recommended)",
"Skip — use the conservative default",
])
@@ -180,7 +180,10 @@ async fn preflight_if_needed(theme: &ColorfulTheme) {
eprintln!();
let Ok(choice) = Select::with_theme(theme)
.with_prompt("Last bandwidth test failed. Try again?")
.items(["Yes — retry now", "No — use the conservative default"])
.items(&[
"Yes — retry now",
"No — use the conservative default",
])
.default(0)
.interact()
else {
@@ -338,12 +341,8 @@ pub fn prompt_player() -> Result<Player> {
let theme = ColorfulTheme::default();
let choice = Select::with_theme(&theme)
.with_prompt("Connected. Pick a player to launch")
.items(["mpv", "VLC"])
.items(&["mpv", "VLC"])
.default(0)
.interact()?;
Ok(if choice == 0 {
Player::Mpv
} else {
Player::Vlc
})
Ok(if choice == 0 { Player::Mpv } else { Player::Vlc })
}
+2 -6
View File
@@ -25,7 +25,7 @@ async fn main() -> Result<()> {
if cli.gui {
#[cfg(feature = "gui")]
{
return gui::run(cli.relay);
return gui::run();
}
#[cfg(not(feature = "gui"))]
{
@@ -73,11 +73,7 @@ async fn main() -> Result<()> {
}
fn init_tracing(verbose: bool) {
let default = if verbose {
"pixelpass=trace,iroh=info"
} else {
"pixelpass=info,iroh=warn"
};
let default = if verbose { "pixelpass=trace,iroh=info" } else { "pixelpass=info,iroh=warn" };
let filter = EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new(default));
// Tracing MUST write to stderr. `tracing_subscriber::fmt()` defaults its
// writer to stdout, but with `--output json` stdout carries the JSON event
+4 -4
View File
@@ -110,7 +110,9 @@ pub async fn run() -> Result<()> {
}
if live_skipped > 0 {
println!("[pixelpass] --repair: left {live_skipped} live pixelpass host(s) alone.");
println!(
"[pixelpass] --repair: left {live_skipped} live pixelpass host(s) alone."
);
}
if failed > 0 {
@@ -152,9 +154,7 @@ fn list_modules() -> Result<Vec<Module>> {
for line in text.lines() {
let mut parts = line.splitn(4, '\t');
let Some(id_str) = parts.next() else { continue };
let Ok(id) = id_str.parse::<u32>() else {
continue;
};
let Ok(id) = id_str.parse::<u32>() else { continue };
let Some(name) = parts.next() else { continue };
let args = parts.next().unwrap_or("").to_string();
modules.push(Module {
+32 -34
View File
@@ -1,10 +1,12 @@
use anyhow::{Context, Result, bail};
use iroh::Endpoint;
use iroh::endpoint::presets;
use iroh_tickets::endpoint::EndpointTicket;
use std::time::Duration;
use tokio::net::TcpListener;
use crate::cli::ViewerOpts;
use crate::common::{alpn::ALPN, endpoint, output, signal};
use crate::common::{alpn::ALPN, output, signal};
/// Cap on the initial QUIC connect. `endpoint.connect()` has no built-in
/// deadline, so an offline host / stale code / unreachable relay otherwise
@@ -15,7 +17,10 @@ const CONNECT_TIMEOUT: Duration = Duration::from_secs(15);
pub async fn run(ticket: EndpointTicket, opts: ViewerOpts) -> Result<()> {
let cancel = signal::install_ctrl_c();
let endpoint = endpoint::bind(opts.relay.as_deref()).await?;
let endpoint = Endpoint::builder(presets::N0)
.alpns(vec![ALPN.to_vec()])
.bind()
.await?;
let addr = ticket.endpoint_addr().clone();
tracing::info!(remote = %addr.id, "connecting to host");
@@ -45,41 +50,34 @@ pub async fn run(ticket: EndpointTicket, opts: ViewerOpts) -> Result<()> {
}
},
};
// Everything past the established connection runs in one block so any error
// (open_bi, bind, local_addr, accept) is captured rather than `?`-propagated
// straight out of the function — that would skip the close below and leak the
// endpoint. The connect-phase arms above close explicitly for the same reason.
let result = async {
let (quic_send, quic_recv) = conn.open_bi().await?;
let (quic_send, quic_recv) = conn.open_bi().await?;
let listener = TcpListener::bind(("127.0.0.1", opts.port)).await?;
let port = listener.local_addr()?.port();
let url = format!("http://127.0.0.1:{port}");
output::emit(output::Event::Connected { url: &url });
let listener = TcpListener::bind(("127.0.0.1", opts.port)).await?;
let port = listener.local_addr()?.port();
let url = format!("http://127.0.0.1:{port}");
output::emit(output::Event::Connected { url: &url });
if opts.interactive {
let player = crate::interactive::prompt_player()?;
player
.spawn(&url)
.with_context(|| "failed to launch player")?;
print_viewer_banner_interactive();
} else {
print_viewer_banner(&url);
}
tokio::select! {
accepted = listener.accept() => {
let (tcp, peer) = accepted?;
tracing::info!(%peer, "local viewer connected");
crate::common::tunnel::bridge(quic_send, quic_recv, tcp).await
}
_ = cancel.cancelled() => {
tracing::info!("ctrl-c received before local viewer connected");
Ok(())
}
}
if opts.interactive {
let player = crate::interactive::prompt_player()?;
player
.spawn(&url)
.with_context(|| "failed to launch player")?;
print_viewer_banner_interactive();
} else {
print_viewer_banner(&url);
}
.await;
let result = tokio::select! {
accepted = listener.accept() => {
let (tcp, peer) = accepted?;
tracing::info!(%peer, "local viewer connected");
crate::common::tunnel::bridge(quic_send, quic_recv, tcp).await
}
_ = cancel.cancelled() => {
tracing::info!("ctrl-c received before local viewer connected");
Ok(())
}
};
endpoint.close().await;
result