Replace patchelf crack with Freeloader LD_PRELOAD approach
- Multi-stage Dockerfile: discover patterns from PMS binary (capstone), compile .so with zig (musl), layer onto lscr.io/linuxserver/plex - Uses LD_PRELOAD instead of patchelf (which corrupts Plex's musl loader) - Auto-discovery: broad structural patterns with string-anchored fallback (//feature) and relationship-based fallback (BITSET_REF within BS_INIT) - hook.cpp uses __has_include for generated patterns with hardcoded fallbacks - Custom wrapper.sh (no traffic_logger preload) - Vendored Freeloader source (github.com/authrequest/Freeloader, AGPL-3.0) - Removed stale plexmediaserver_crack.so binary - Supports Plex 1.43.3+ (verified against 1.43.2 and 1.43.3)
This commit is contained in:
1 parent
4399a8288d
commit
72f4661bdc
72 files changed
+77927
-17
No files matched your search
@@ -0,0 +1,7 @@
|
||||
__pycache__/
|
||||
*.pyc
|
||||
.pytest_cache/
|
||||
*.egg-info/
|
||||
build/
|
||||
dist/
|
||||
.venv/
|
||||
@@ -0,0 +1,178 @@
|
||||
<!-- SPDX-License-Identifier: AGPL-3.0-or-later -->
|
||||
# plex_relay
|
||||
|
||||
A reverse-engineered, runnable reimplementation of the **Plex Media Server**
|
||||
`RelayController` — the component that makes a server reachable remotely by
|
||||
opening a **reverse SSH tunnel out to a Plex-operated relay host**.
|
||||
|
||||
Reconstructed from the `Plex Media Server` **1.43.2.10687** binary (Linux
|
||||
x86-64). Ships **no Plex code**, embeds no keys, and authenticates to nothing on
|
||||
its own. It is a behavioural model for interoperability research on
|
||||
infrastructure **you operate yourself**.
|
||||
|
||||
---
|
||||
|
||||
## What Plex Relay does (reversed)
|
||||
|
||||
```
|
||||
plex.tv ──"startRelay"(host,port)──▶ ServerEventManager_handle_pubsub_event (0xF5BE3C)
|
||||
│ gate: signed_in && published && relay_enabled
|
||||
▼
|
||||
RelayController_connect (0x12307F2)
|
||||
│ 1 dedup an already-active tunnel for host
|
||||
│ 2 refresh relay host key (≤24h cache)
|
||||
│ 3 pin [host]:443 in relayHostKey.txt
|
||||
│ 4 spawn the ssh reverse tunnel
|
||||
│ 5 arm the 300s inactive-connection reaper
|
||||
▼
|
||||
ssh -p <port> -N -R 0:127.0.0.1:<PMS port>
|
||||
-o UserKnownHostsFile=<datadir>/relayHostKey.txt
|
||||
-o LogLevel=VERBOSE -o PreferredAuthentications=password
|
||||
-o PubkeyAuthentication=no -l <myplex-id> -F /dev/null <relay-host>
|
||||
password = MyPlex token, delivered via the PLEXTOKEN env var + SSH_ASKPASS
|
||||
```
|
||||
|
||||
The relay binds an ephemeral remote port (`-R 0:…`) and forwards inbound
|
||||
remote-client traffic back down the tunnel to the local PMS service. The relay
|
||||
host's SSH key is **pinned**: PMS downloads it at most once per day from
|
||||
`https://downloads.plex.tv/relay/relay_v1.pub` and writes a per-endpoint
|
||||
known_hosts line. `relayHostKey.txt` doubles as PMS's cache and the file handed
|
||||
to `ssh` (the `#` lines are valid known_hosts comments).
|
||||
|
||||
---
|
||||
|
||||
## Architecture
|
||||
|
||||
High cohesion (one reason to change per module) and low coupling (dependencies
|
||||
point inward to abstractions, never outward to I/O):
|
||||
|
||||
```
|
||||
cli composition root / argument parsing
|
||||
└─ controller orchestration; depends ONLY on the two protocols below
|
||||
├─ store HostKeyTrust ── composes ↓↓
|
||||
│ ├─ keys RelayKeyProvider (HTTPS fetch + TTL cache)
|
||||
│ └─ cache HostKeyCache (relayHostKey.txt, atomic, 0o600)
|
||||
└─ tunnel TunnelFactory ── build_ssh_argv (pure) + SSH_ASKPASS + child process
|
||||
models immutable domain values + parsing (no I/O, thread-safe)
|
||||
config immutable, validated configuration
|
||||
errors one rooted exception hierarchy
|
||||
```
|
||||
|
||||
**Dependency inversion.** `RelayController` names what it needs as `Protocol`s —
|
||||
`HostKeyTrust` (store) and `TunnelFactory`/`Tunnel` (tunnel) — and is handed
|
||||
concrete adapters by `RelayController.from_config`, the single composition root.
|
||||
Every external concern (HTTP, filesystem, subprocess, clock) sits behind an
|
||||
injected seam, so the orchestration is unit-tested with fakes and **no network,
|
||||
disk, or process is touched** in the suite.
|
||||
|
||||
| Module | Responsibility | Depends on |
|
||||
| --- | --- | --- |
|
||||
| `errors` | exception taxonomy | — |
|
||||
| `models` | `HostKey`, `RelayKey`, `parse_relay_pub`, endpoint formatting | `errors` |
|
||||
| `config` | frozen, validated `RelayConfig` (token redacted from `repr`) | `errors` |
|
||||
| `keys` | fetch `relay_v1.pub` (HTTPS-only, byte-capped, timed) + TTL cache | `errors`, `models` |
|
||||
| `cache` | load/save `relayHostKey.txt` atomically at `0o600` | `errors`, `models` |
|
||||
| `store` | `HostKeyManager`: compose key + cache, pin endpoints | `keys`, `cache`, `models` |
|
||||
| `tunnel` | `build_ssh_argv` (pure) + askpass + `SubprocessTunnel` | `config`, `errors` |
|
||||
| `controller` | connect / reap / stop lifecycle, thread-safety | the protocols above |
|
||||
| `cli` | wire adapters, parse args | everything |
|
||||
|
||||
### Binary → code map
|
||||
|
||||
| Binary symbol | Address | Code |
|
||||
| --- | --- | --- |
|
||||
| `RelayController` ctor (loads cache) | `0x122FDDA` | `HostKeyCache.load` + `HostKeyManager.__init__` |
|
||||
| `RelayController_connect` | `0x12307F2` | `RelayController.connect` + `store` + `tunnel` |
|
||||
| stopRelay | `0x123068C` | `RelayController.stop` |
|
||||
| inactive-connection reaper (300s) | `0x12320EE` | `RelayController.reap_once` / `_reaper_tick` |
|
||||
| `relayHostKey.txt` path | `0x1231FEE` | `RelayConfig.cache_path` |
|
||||
| `startRelay` PubSub dispatch | `0xF5BE3C` | `RelayController.start_relay` |
|
||||
|
||||
---
|
||||
|
||||
## Error model
|
||||
|
||||
All failures derive from `RelayError`, so callers catch the subsystem broadly or
|
||||
a specific mode. Adapter exceptions (`urllib`, `OSError`, `subprocess`) are
|
||||
caught at the boundary and re-raised as domain errors — they never leak.
|
||||
|
||||
- `ConfigError` — invalid configuration (raised eagerly in `RelayConfig`).
|
||||
- `RelayKeyError` — relay key fetch/parse (non-HTTPS, oversize, transport, bad format).
|
||||
- `HostKeyCacheError` — unreadable/malformed cache (load auto-rebuilds, as PMS does).
|
||||
- `TunnelError` — `ssh` could not be launched.
|
||||
|
||||
**Contract:** `connect()` *raises* on failure (library callers decide).
|
||||
`start_relay()` — the plex.tv event entry — is *resilient*: it logs and returns
|
||||
`False`, mirroring PMS so a bad event can't kill an event loop.
|
||||
|
||||
---
|
||||
|
||||
## Security
|
||||
|
||||
- **Credential never on the command line.** The token is passed to `ssh` only
|
||||
via the `PLEXTOKEN` env var, read by a generated `SSH_ASKPASS` helper. The
|
||||
helper is mode `0o700` and removed on stop *and* via a `weakref.finalize`, so a
|
||||
crash can't leak it. `RelayConfig` excludes the token from `repr`.
|
||||
- **Fetch hardening.** The relay-key URL is HTTPS-only by default (configurable
|
||||
URLs are an SSRF surface), the response is byte-capped, and the request is
|
||||
time-limited.
|
||||
- **Trust-store integrity.** `relayHostKey.txt` is written atomically at `0o600`;
|
||||
it pins the host key `ssh` verifies, so it must not be world-writable.
|
||||
- **No shell.** Processes are spawned from an argv list, never a shell string.
|
||||
|
||||
---
|
||||
|
||||
## Performance & scale
|
||||
|
||||
- **At most one key fetch per TTL**, behind a lock, on a **monotonic** clock
|
||||
(immune to wall-clock jumps).
|
||||
- **Disk writes only on change** — re-pinning an unchanged endpoint is a no-op.
|
||||
- **O(1) liveness** via `Popen.poll()`; the reaper is a single background
|
||||
`threading.Timer` that sweeps O(n) connections every 300s and reschedules
|
||||
itself only while connections remain (idle controllers spawn no timers).
|
||||
- **Immutable domain + value objects** are freely shareable across threads;
|
||||
mutable state lives behind one `RLock`.
|
||||
|
||||
---
|
||||
|
||||
## Install & test
|
||||
|
||||
```bash
|
||||
cd plex_relay
|
||||
python -m pip install -e ".[test]"
|
||||
python -m pytest -q # 51 tests + 1 POSIX-only; no network, no ssh
|
||||
```
|
||||
|
||||
## Use
|
||||
|
||||
```bash
|
||||
export PLEX_RELAY_TOKEN=... # keep the secret off the cmdline
|
||||
plex-relay show --host RELAY_HOST --user MY_ID # dry run: resolve key + print argv
|
||||
plex-relay connect --host RELAY_HOST --user MY_ID --local-port 32400
|
||||
```
|
||||
|
||||
```python
|
||||
from plex_relay import RelayConfig, RelayController
|
||||
|
||||
with RelayController.from_config(RelayConfig(token="…", ssh_user="my-id")) as ctrl:
|
||||
ctrl.start_relay("relay.example.net", 443) # gated + resilient, like PMS
|
||||
```
|
||||
|
||||
For tests or custom transports, bypass the composition root and inject your own
|
||||
collaborators: `RelayController(config, hostkeys=…, tunnels=…)`.
|
||||
|
||||
---
|
||||
|
||||
## Fidelity & limitations
|
||||
|
||||
- **Faithful:** ssh argv (order + flags), the `relayHostKey.txt` format, the 24h
|
||||
key cache, per-endpoint pinning, the dedup / 300s reaper lifecycle, and the
|
||||
`startRelay` gating.
|
||||
- **Adapted, with intent:** the known_hosts endpoint is keyed by the actual ssh
|
||||
port (the binary hardcodes `:443`); the local forward target is configurable
|
||||
(the binary reads the PMS port from its own config); TTLs use a monotonic clock.
|
||||
- **POSIX only:** password delivery uses `SSH_ASKPASS`, which Windows OpenSSH
|
||||
does not honour.
|
||||
- **Not a turnkey relay swap:** this is only the *server→relay* leg. Plex brokers
|
||||
both ends, so pointing it at your own relay also needs a relay `sshd` you
|
||||
control plus client-discovery redirection (see the parent project's notes).
|
||||
@@ -0,0 +1,5 @@
|
||||
# SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
import sys
|
||||
from pathlib import Path
|
||||
|
||||
sys.path.insert(0, str(Path(__file__).parent / "src"))
|
||||
@@ -0,0 +1,39 @@
|
||||
# SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
[build-system]
|
||||
requires = ["setuptools>=61"]
|
||||
build-backend = "setuptools.build_meta"
|
||||
|
||||
[project]
|
||||
name = "plex-relay"
|
||||
version = "0.2.0"
|
||||
description = "Reverse-engineered reimplementation of the Plex Media Server RelayController (educational / RE use)."
|
||||
readme = "README.md"
|
||||
requires-python = ">=3.10"
|
||||
license = { text = "AGPL-3.0-or-later" }
|
||||
authors = [{ name = "Plex_Patch RE notes" }]
|
||||
dependencies = [] # stdlib only
|
||||
|
||||
[project.optional-dependencies]
|
||||
test = ["pytest>=7"]
|
||||
|
||||
[project.scripts]
|
||||
plex-relay = "plex_relay.cli:main"
|
||||
|
||||
[tool.setuptools.packages.find]
|
||||
where = ["src"]
|
||||
|
||||
[tool.setuptools.package-data]
|
||||
plex_relay = ["py.typed"]
|
||||
|
||||
[tool.pytest.ini_options]
|
||||
testpaths = ["tests"]
|
||||
|
||||
[tool.mypy]
|
||||
python_version = "3.10"
|
||||
warn_unused_ignores = true
|
||||
warn_redundant_casts = true
|
||||
disallow_untyped_defs = true
|
||||
|
||||
[tool.ruff]
|
||||
line-length = 100
|
||||
target-version = "py310"
|
||||
@@ -0,0 +1,46 @@
|
||||
# SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
"""plex_relay -- a reverse-engineered reimplementation of the Plex Media Server
|
||||
``RelayController`` (Linux x86-64, build 1.43.2.10687).
|
||||
|
||||
Layering (high cohesion, dependency inversion top-to-bottom):
|
||||
|
||||
cli composition root / argument parsing
|
||||
controller orchestration; depends only on the protocols below
|
||||
store HostKeyTrust: composes the key provider + the disk cache
|
||||
keys relay key acquisition (HTTPS fetch + TTL cache)
|
||||
cache relayHostKey.txt persistence (atomic, 0o600)
|
||||
tunnel ssh argv (pure) + SSH_ASKPASS + child-process tunnel
|
||||
models immutable domain values + parsing (no I/O)
|
||||
config immutable, validated configuration
|
||||
errors one rooted exception hierarchy
|
||||
|
||||
Binary provenance of the key symbols is documented in each module and in
|
||||
``README.md``. Ships no Plex code; authenticates to nothing on its own.
|
||||
"""
|
||||
from .config import RelayConfig
|
||||
from .controller import RelayController
|
||||
from .errors import (
|
||||
ConfigError,
|
||||
HostKeyCacheError,
|
||||
RelayError,
|
||||
RelayKeyError,
|
||||
TunnelError,
|
||||
)
|
||||
from .models import HostKey, RelayKey, known_hosts_endpoint, parse_relay_pub
|
||||
from .tunnel import build_ssh_argv
|
||||
|
||||
__all__ = [
|
||||
"RelayConfig",
|
||||
"RelayController",
|
||||
"HostKey",
|
||||
"RelayKey",
|
||||
"known_hosts_endpoint",
|
||||
"parse_relay_pub",
|
||||
"build_ssh_argv",
|
||||
"RelayError",
|
||||
"ConfigError",
|
||||
"RelayKeyError",
|
||||
"HostKeyCacheError",
|
||||
"TunnelError",
|
||||
]
|
||||
__version__ = "0.2.0"
|
||||
@@ -0,0 +1,97 @@
|
||||
# SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
"""On-disk known_hosts cache (``relayHostKey.txt``) -- file I/O only.
|
||||
|
||||
This module owns the file format and the filesystem; it knows nothing about
|
||||
HTTP, TTLs, or ssh. The file doubles as PMS's cache and the OpenSSH
|
||||
``UserKnownHostsFile`` (``#`` lines are valid known_hosts comments).
|
||||
|
||||
Writes are atomic (write-temp-then-rename) and the file is mode ``0o600`` --
|
||||
it pins the keys ssh will trust, so it must not be world-writable.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import os
|
||||
import stat
|
||||
import tempfile
|
||||
from pathlib import Path
|
||||
from typing import Iterable, Mapping
|
||||
|
||||
from .errors import HostKeyCacheError
|
||||
from .models import HostKey
|
||||
|
||||
log = logging.getLogger("plex_relay.cache")
|
||||
|
||||
|
||||
def parse_known_hosts(lines: Iterable[str]) -> dict[str, HostKey]:
|
||||
"""Parse ``# <endpoint>`` + ``<host> <keytype> <keydata>`` line pairs.
|
||||
|
||||
:raises HostKeyCacheError: on any structural violation (mirrors the binary,
|
||||
which discards and rebuilds a malformed file).
|
||||
"""
|
||||
items = [ln.rstrip("\n") for ln in lines if ln.strip()]
|
||||
entries: dict[str, HostKey] = {}
|
||||
i = 0
|
||||
while i < len(items):
|
||||
comment = items[i]
|
||||
if not comment.startswith("#"):
|
||||
raise HostKeyCacheError(f"expected '# <endpoint>' marker, got {comment!r}")
|
||||
if i + 1 >= len(items):
|
||||
raise HostKeyCacheError("comment marker without a following data line")
|
||||
tokens = items[i + 1].split()
|
||||
if len(tokens) != 3:
|
||||
raise HostKeyCacheError(f"data line part count incorrect: {items[i + 1]!r}")
|
||||
endpoint = comment[1:].strip() or tokens[0]
|
||||
host, keytype, keydata = tokens
|
||||
entries[endpoint] = HostKey(host, keytype, keydata)
|
||||
i += 2
|
||||
return entries
|
||||
|
||||
|
||||
class HostKeyCache:
|
||||
"""Load/store relay host keys from a single known_hosts-format file."""
|
||||
|
||||
def __init__(self, path: Path) -> None:
|
||||
self._path = Path(path)
|
||||
|
||||
@property
|
||||
def path(self) -> Path:
|
||||
return self._path
|
||||
|
||||
def load(self) -> dict[str, HostKey]:
|
||||
"""Return cached entries; a malformed file is dropped and rebuilt empty."""
|
||||
if not self._path.exists():
|
||||
return {}
|
||||
try:
|
||||
text = self._path.read_text("utf-8", "replace")
|
||||
except OSError as exc:
|
||||
raise HostKeyCacheError(f"cannot read {self._path}: {exc}") from exc
|
||||
try:
|
||||
entries = parse_known_hosts(text.splitlines())
|
||||
except HostKeyCacheError as exc:
|
||||
log.warning("host key file malformed, rebuilding: %s", exc)
|
||||
self.save({})
|
||||
return {}
|
||||
log.info("read %d cached host key entries", len(entries))
|
||||
return entries
|
||||
|
||||
def save(self, entries: Mapping[str, HostKey]) -> None:
|
||||
"""Atomically write ``entries`` to the cache file with mode 0o600."""
|
||||
try:
|
||||
self._path.parent.mkdir(parents=True, exist_ok=True)
|
||||
blocks = "".join(entries[k].cache_block() for k in sorted(entries))
|
||||
fd, tmp_name = tempfile.mkstemp(
|
||||
dir=self._path.parent, prefix=self._path.name + ".", suffix=".tmp"
|
||||
)
|
||||
tmp = Path(tmp_name)
|
||||
try:
|
||||
os.write(fd, blocks.encode("utf-8"))
|
||||
finally:
|
||||
os.close(fd)
|
||||
os.chmod(tmp, stat.S_IRUSR | stat.S_IWUSR) # 0o600
|
||||
os.replace(tmp, self._path)
|
||||
except OSError as exc:
|
||||
raise HostKeyCacheError(f"cannot write {self._path}: {exc}") from exc
|
||||
|
||||
|
||||
__all__ = ["HostKeyCache", "parse_known_hosts"]
|
||||
@@ -0,0 +1,126 @@
|
||||
# SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
"""Command-line composition root.
|
||||
|
||||
plex-relay show --host H [--port 443] ... # resolve key + print ssh argv (no spawn)
|
||||
plex-relay connect --host H [--port 443] ... # establish and hold the tunnel
|
||||
|
||||
The token is read from ``$PLEX_RELAY_TOKEN`` by default so it never appears in
|
||||
the process list; ``--token`` overrides. ``show`` is a safe dry run.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import logging
|
||||
import os
|
||||
import signal
|
||||
import sys
|
||||
import threading
|
||||
from pathlib import Path
|
||||
|
||||
from .config import DEFAULT_LOCAL_HOST, DEFAULT_LOCAL_PORT, DEFAULT_RELAY_KEY_URL, RelayConfig
|
||||
from .errors import RelayError
|
||||
from .keys import HttpsRelayKeyFetcher, RelayKeyProvider
|
||||
from .models import known_hosts_endpoint
|
||||
from .tunnel import build_ssh_argv
|
||||
|
||||
log = logging.getLogger("plex_relay.cli")
|
||||
|
||||
|
||||
def _add_common(p: argparse.ArgumentParser) -> None:
|
||||
p.add_argument("--host", required=True, help="relay host to dial")
|
||||
p.add_argument("--port", type=int, default=443, help="relay SSH port (default 443)")
|
||||
p.add_argument("--user", required=True, help="relay SSH login (PMS: MyPlex identity)")
|
||||
p.add_argument("--token", default=os.environ.get("PLEX_RELAY_TOKEN", ""),
|
||||
help="relay password (default: $PLEX_RELAY_TOKEN)")
|
||||
p.add_argument("--local-host", default=DEFAULT_LOCAL_HOST)
|
||||
p.add_argument("--local-port", type=int, default=DEFAULT_LOCAL_PORT)
|
||||
p.add_argument("--data-dir", type=Path, default=Path.home() / ".plex_relay")
|
||||
p.add_argument("--key-url", default=DEFAULT_RELAY_KEY_URL)
|
||||
p.add_argument("--insecure-key-url", action="store_true",
|
||||
help="permit a non-HTTPS relay key URL (e.g. file:// for testing)")
|
||||
p.add_argument("--ssh", default="ssh", help="ssh binary")
|
||||
|
||||
|
||||
def _config(args: argparse.Namespace) -> RelayConfig:
|
||||
return RelayConfig(
|
||||
token=args.token or "dry-run",
|
||||
ssh_user=args.user,
|
||||
local_host=args.local_host,
|
||||
local_port=args.local_port,
|
||||
relay_key_url=args.key_url,
|
||||
data_dir=args.data_dir,
|
||||
ssh_binary=args.ssh,
|
||||
allow_insecure_key_url=args.insecure_key_url,
|
||||
)
|
||||
|
||||
|
||||
def _cmd_show(args: argparse.Namespace) -> int:
|
||||
cfg = _config(args)
|
||||
provider = RelayKeyProvider(
|
||||
cfg.relay_key_url,
|
||||
cfg.key_ttl_seconds,
|
||||
fetcher=HttpsRelayKeyFetcher(
|
||||
timeout=cfg.connect_timeout_seconds, allow_insecure=cfg.allow_insecure_key_url
|
||||
),
|
||||
)
|
||||
endpoint = known_hosts_endpoint(args.host, args.port)
|
||||
try:
|
||||
key = provider.get()
|
||||
print(f"relay key : {key.keytype} {key.keydata[:24]}... (from {cfg.relay_key_url})")
|
||||
print(f"known_hosts : {endpoint} {key.keytype} {key.keydata[:24]}...")
|
||||
except RelayError as exc:
|
||||
print(f"relay key : <unavailable: {exc}>")
|
||||
argv = build_ssh_argv(cfg, args.host, args.port, cfg.cache_path)
|
||||
print("ssh argv :\n " + " ".join(argv))
|
||||
print("env : PLEXTOKEN=*** SSH_ASKPASS=<helper> SSH_ASKPASS_REQUIRE=force")
|
||||
return 0
|
||||
|
||||
|
||||
def _cmd_connect(args: argparse.Namespace) -> int:
|
||||
if not args.token:
|
||||
print("error: no token (set $PLEX_RELAY_TOKEN or --token)", file=sys.stderr)
|
||||
return 2
|
||||
from .controller import RelayController
|
||||
|
||||
cfg = _config(args)
|
||||
stop = threading.Event()
|
||||
with RelayController.from_config(cfg) as ctrl:
|
||||
if not ctrl.start_relay(args.host, args.port):
|
||||
print("error: relay did not start (gating, already active, or failure)", file=sys.stderr)
|
||||
return 1
|
||||
print(f"relay up: {args.host}:{args.port} -> {cfg.local_host}:{cfg.local_port} (Ctrl-C to stop)")
|
||||
signal.signal(signal.SIGINT, lambda *_: stop.set())
|
||||
signal.signal(signal.SIGTERM, lambda *_: stop.set())
|
||||
while not stop.is_set():
|
||||
stop.wait(2.0)
|
||||
if not ctrl.active_hosts:
|
||||
print("relay connection ended", file=sys.stderr)
|
||||
return 1
|
||||
print("relay stopped")
|
||||
return 0
|
||||
|
||||
|
||||
def main(argv: list[str] | None = None) -> int:
|
||||
parser = argparse.ArgumentParser(prog="plex-relay", description=__doc__)
|
||||
parser.add_argument("-v", "--verbose", action="store_true")
|
||||
sub = parser.add_subparsers(dest="cmd", required=True)
|
||||
for name, func, help_ in (("show", _cmd_show, "resolve key + print ssh argv (no spawn)"),
|
||||
("connect", _cmd_connect, "establish and hold the relay tunnel")):
|
||||
sp = sub.add_parser(name, help=help_)
|
||||
_add_common(sp)
|
||||
sp.set_defaults(func=func)
|
||||
|
||||
args = parser.parse_args(argv)
|
||||
logging.basicConfig(
|
||||
level=logging.DEBUG if args.verbose else logging.INFO,
|
||||
format="%(levelname)s %(name)s: %(message)s",
|
||||
)
|
||||
try:
|
||||
return int(args.func(args))
|
||||
except RelayError as exc:
|
||||
print(f"error: {exc}", file=sys.stderr)
|
||||
return 1
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
@@ -0,0 +1,85 @@
|
||||
# SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
"""Immutable, validated runtime configuration.
|
||||
|
||||
Frozen so it can be shared freely across threads and never mutated behind a
|
||||
collaborator's back. The relay credential is excluded from ``repr`` so it cannot
|
||||
leak into logs or tracebacks.
|
||||
|
||||
Defaults mirror constants observed in ``RelayController_connect`` (``0x12307F2``):
|
||||
relay key URL, the 24h key TTL (``86400000000`` us) and the 300s reaper
|
||||
(``300000000`` us). The reverse-forward target (``0:127.0.0.1:%d``) is the PMS
|
||||
service port, which the binary reads from its own config -- hence configurable.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass, field
|
||||
from pathlib import Path
|
||||
|
||||
from .errors import ConfigError
|
||||
|
||||
DEFAULT_RELAY_KEY_URL = "https://downloads.plex.tv/relay/relay_v1.pub"
|
||||
DEFAULT_KEY_TTL_SECONDS = 86_400.0
|
||||
DEFAULT_REAP_INTERVAL_SECONDS = 300.0
|
||||
DEFAULT_LOCAL_HOST = "127.0.0.1"
|
||||
DEFAULT_LOCAL_PORT = 32400
|
||||
CACHE_FILENAME = "relayHostKey.txt"
|
||||
|
||||
|
||||
def _default_data_dir() -> Path:
|
||||
return Path.home() / ".plex_relay"
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class RelayConfig:
|
||||
"""Everything :class:`~plex_relay.controller.RelayController` needs to run."""
|
||||
|
||||
token: str = field(repr=False) # SSH password (PLEXTOKEN); never logged
|
||||
ssh_user: str
|
||||
|
||||
local_host: str = DEFAULT_LOCAL_HOST
|
||||
local_port: int = DEFAULT_LOCAL_PORT
|
||||
|
||||
relay_key_url: str = DEFAULT_RELAY_KEY_URL
|
||||
key_ttl_seconds: float = DEFAULT_KEY_TTL_SECONDS
|
||||
reap_interval_seconds: float = DEFAULT_REAP_INTERVAL_SECONDS
|
||||
|
||||
data_dir: Path = field(default_factory=_default_data_dir)
|
||||
ssh_binary: str = "ssh"
|
||||
ssh_interface: str = "tailscale0"
|
||||
|
||||
connect_timeout_seconds: float = 15.0
|
||||
stop_timeout_seconds: float = 5.0
|
||||
allow_insecure_key_url: bool = False
|
||||
|
||||
# startRelay gating, mirroring ServerEventManager_handle_pubsub_event:
|
||||
# signin_state == 4 && PublishServerOnPlexOnlineKey && RelayEnabled
|
||||
signed_in: bool = True
|
||||
published: bool = True
|
||||
relay_enabled: bool = True
|
||||
|
||||
def __post_init__(self) -> None:
|
||||
# Frozen dataclass: normalise/validate via object.__setattr__.
|
||||
object.__setattr__(self, "data_dir", Path(self.data_dir))
|
||||
if not self.token:
|
||||
raise ConfigError("token is required (relay SSH password)")
|
||||
if not self.ssh_user:
|
||||
raise ConfigError("ssh_user is required (relay SSH login)")
|
||||
if not 0 < self.local_port < 65536:
|
||||
raise ConfigError(f"local_port out of range: {self.local_port}")
|
||||
for name in ("key_ttl_seconds", "reap_interval_seconds",
|
||||
"connect_timeout_seconds", "stop_timeout_seconds"):
|
||||
if getattr(self, name) <= 0:
|
||||
raise ConfigError(f"{name} must be positive")
|
||||
|
||||
@property
|
||||
def cache_path(self) -> Path:
|
||||
"""``<data_dir>/relayHostKey.txt`` (the binary's ``sub_1231FEE``)."""
|
||||
return self.data_dir / CACHE_FILENAME
|
||||
|
||||
def gating_ok(self) -> bool:
|
||||
"""True iff all three startRelay preconditions hold."""
|
||||
return self.signed_in and self.published and self.relay_enabled
|
||||
|
||||
|
||||
__all__ = ["RelayConfig", "DEFAULT_RELAY_KEY_URL", "DEFAULT_LOCAL_HOST",
|
||||
"DEFAULT_LOCAL_PORT", "CACHE_FILENAME"]
|
||||
@@ -0,0 +1,162 @@
|
||||
# SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
"""Relay controller -- the orchestration layer.
|
||||
|
||||
Depends only on abstractions (:class:`HostKeyTrust`, :class:`TunnelFactory`),
|
||||
so the network, disk, and process concerns are all injected and substitutable.
|
||||
:meth:`RelayController.from_config` is the composition root that wires the
|
||||
default adapters together.
|
||||
|
||||
Concurrency: a single ``RLock`` (the binary's ``recursive_mutex``) guards the
|
||||
connection table and the reaper, which runs on a background ``threading.Timer``
|
||||
and reschedules itself only while connections remain.
|
||||
|
||||
Error contract: :meth:`connect` raises :class:`RelayError` on failure (callers
|
||||
decide). :meth:`start_relay` -- the plex.tv event entry point -- is resilient:
|
||||
it logs and returns ``False`` rather than letting an exception escape an event
|
||||
loop, matching PMS's ``ServerEventManager`` behaviour.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import threading
|
||||
|
||||
from .cache import HostKeyCache
|
||||
from .config import RelayConfig
|
||||
from .errors import RelayError
|
||||
from .keys import RelayKeyProvider
|
||||
from .store import HostKeyManager, HostKeyTrust
|
||||
from .tunnel import SubprocessTunnelFactory, Tunnel, TunnelFactory
|
||||
|
||||
log = logging.getLogger("plex_relay.controller")
|
||||
|
||||
|
||||
class RelayController:
|
||||
"""Manages relay tunnels for one server: connect, reap, stop."""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
config: RelayConfig,
|
||||
hostkeys: HostKeyTrust,
|
||||
tunnels: TunnelFactory,
|
||||
) -> None:
|
||||
self._config = config
|
||||
self._hostkeys = hostkeys
|
||||
self._tunnels = tunnels
|
||||
self._lock = threading.RLock()
|
||||
self._connections: dict[str, Tunnel] = {}
|
||||
self._reaper: threading.Timer | None = None
|
||||
self._closed = False
|
||||
|
||||
@classmethod
|
||||
def from_config(cls, config: RelayConfig) -> "RelayController":
|
||||
"""Composition root: wire the default network/disk/process adapters."""
|
||||
provider = RelayKeyProvider(
|
||||
config.relay_key_url,
|
||||
config.key_ttl_seconds,
|
||||
fetcher=_default_fetcher(config),
|
||||
)
|
||||
manager = HostKeyManager(provider, HostKeyCache(config.cache_path))
|
||||
return cls(config, manager, SubprocessTunnelFactory(config))
|
||||
|
||||
# -- entry points ------------------------------------------------------
|
||||
|
||||
def start_relay(self, host: str, port: int = 443) -> bool:
|
||||
"""plex.tv ``startRelay`` handler: gated and resilient."""
|
||||
if not self._config.gating_ok():
|
||||
log.info(
|
||||
"startRelay ignored (signed_in=%s published=%s relay_enabled=%s)",
|
||||
self._config.signed_in, self._config.published, self._config.relay_enabled,
|
||||
)
|
||||
return False
|
||||
try:
|
||||
return self.connect(host, port)
|
||||
except RelayError as exc:
|
||||
log.error("relay to %s failed: %s", host, exc)
|
||||
return False
|
||||
|
||||
def connect(self, host: str, port: int = 443) -> bool:
|
||||
"""Establish a tunnel to ``host``. Returns False if already active.
|
||||
|
||||
:raises RelayError: if the key cannot be obtained or ssh cannot launch.
|
||||
"""
|
||||
with self._lock:
|
||||
if self._closed:
|
||||
raise RelayError("controller is closed")
|
||||
existing = self._connections.get(host)
|
||||
if existing is not None and existing.is_alive():
|
||||
log.info("already have an active relay connection to %s", host)
|
||||
return False
|
||||
self._hostkeys.ensure_trusted(host, port)
|
||||
tunnel = self._tunnels(host, port, self._hostkeys.known_hosts_path)
|
||||
tunnel.start()
|
||||
self._connections[host] = tunnel
|
||||
self._arm_reaper()
|
||||
return True
|
||||
|
||||
def stop(self) -> None:
|
||||
"""Cancel the reaper and tear down every tunnel. Idempotent."""
|
||||
with self._lock:
|
||||
self._closed = True
|
||||
self._cancel_reaper()
|
||||
connections, self._connections = self._connections, {}
|
||||
for tunnel in connections.values():
|
||||
tunnel.stop(self._config.stop_timeout_seconds)
|
||||
|
||||
# -- reaper ------------------------------------------------------------
|
||||
|
||||
def reap_once(self) -> list[str]:
|
||||
"""Drop finished tunnels; return the hosts removed."""
|
||||
with self._lock:
|
||||
dead = [h for h, t in self._connections.items() if not t.is_alive()]
|
||||
stopped = [self._connections.pop(h) for h in dead]
|
||||
for tunnel in stopped:
|
||||
log.info("cleaning up inactive relay connection to %s", tunnel.host)
|
||||
tunnel.stop(self._config.stop_timeout_seconds)
|
||||
return dead
|
||||
|
||||
def _arm_reaper(self) -> None:
|
||||
if self._reaper is None and self._connections and not self._closed:
|
||||
self._schedule_reaper()
|
||||
|
||||
def _schedule_reaper(self) -> None:
|
||||
timer = threading.Timer(self._config.reap_interval_seconds, self._reaper_tick)
|
||||
timer.daemon = True
|
||||
self._reaper = timer
|
||||
timer.start()
|
||||
|
||||
def _reaper_tick(self) -> None:
|
||||
self.reap_once()
|
||||
with self._lock:
|
||||
self._reaper = None
|
||||
if self._connections and not self._closed:
|
||||
self._schedule_reaper()
|
||||
|
||||
def _cancel_reaper(self) -> None:
|
||||
if self._reaper is not None:
|
||||
self._reaper.cancel()
|
||||
self._reaper = None
|
||||
|
||||
# -- introspection -----------------------------------------------------
|
||||
|
||||
@property
|
||||
def active_hosts(self) -> list[str]:
|
||||
with self._lock:
|
||||
return sorted(h for h, t in self._connections.items() if t.is_alive())
|
||||
|
||||
def __enter__(self) -> "RelayController":
|
||||
return self
|
||||
|
||||
def __exit__(self, *exc: object) -> None:
|
||||
self.stop()
|
||||
|
||||
|
||||
def _default_fetcher(config: RelayConfig):
|
||||
from .keys import HttpsRelayKeyFetcher
|
||||
|
||||
return HttpsRelayKeyFetcher(
|
||||
timeout=config.connect_timeout_seconds,
|
||||
allow_insecure=config.allow_insecure_key_url,
|
||||
)
|
||||
|
||||
|
||||
__all__ = ["RelayController"]
|
||||
@@ -0,0 +1,38 @@
|
||||
# SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
"""Exception hierarchy for :mod:`plex_relay`.
|
||||
|
||||
A single rooted hierarchy lets callers catch the whole subsystem
|
||||
(``except RelayError``) or a specific failure mode, and keeps adapter-specific
|
||||
exceptions (``urllib``, ``OSError``, ``subprocess``) from leaking across module
|
||||
boundaries.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
|
||||
class RelayError(Exception):
|
||||
"""Base class for every error raised by this package."""
|
||||
|
||||
|
||||
class ConfigError(RelayError):
|
||||
"""Invalid configuration (bad port, empty credential, ...)."""
|
||||
|
||||
|
||||
class RelayKeyError(RelayError):
|
||||
"""The relay host key could not be fetched or parsed."""
|
||||
|
||||
|
||||
class HostKeyCacheError(RelayError):
|
||||
"""The on-disk known_hosts cache is unreadable or malformed."""
|
||||
|
||||
|
||||
class TunnelError(RelayError):
|
||||
"""The ssh relay tunnel could not be launched."""
|
||||
|
||||
|
||||
__all__ = [
|
||||
"RelayError",
|
||||
"ConfigError",
|
||||
"RelayKeyError",
|
||||
"HostKeyCacheError",
|
||||
"TunnelError",
|
||||
]
|
||||
@@ -0,0 +1,105 @@
|
||||
# SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
"""Relay host-key acquisition: fetch ``relay_v1.pub`` and cache it with a TTL.
|
||||
|
||||
Separated from on-disk caching (:mod:`plex_relay.cache`) so the network policy
|
||||
(HTTPS-only, byte-capped, time-limited) and the freshness policy (TTL) live in
|
||||
one cohesive place and can be swapped wholesale in tests via the injected
|
||||
``fetcher``/``clock`` seams.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import threading
|
||||
import time
|
||||
import urllib.parse
|
||||
import urllib.request
|
||||
from typing import Callable, Protocol
|
||||
|
||||
from .errors import RelayKeyError
|
||||
from .models import RelayKey, parse_relay_pub
|
||||
|
||||
log = logging.getLogger("plex_relay.keys")
|
||||
|
||||
#: Fetches the raw body of a relay-key URL. The seam that tests stub.
|
||||
Fetcher = Callable[[str], str]
|
||||
#: Monotonic time source (seconds). Monotonic so TTLs survive wall-clock jumps.
|
||||
Clock = Callable[[], float]
|
||||
|
||||
_MAX_KEY_BYTES = 64 * 1024 # a host key is a few hundred bytes; cap to bound I/O
|
||||
|
||||
|
||||
class _Opener(Protocol):
|
||||
def __call__(self, url: str, timeout: float): ... # pragma: no cover
|
||||
|
||||
|
||||
class HttpsRelayKeyFetcher:
|
||||
"""Default fetcher: HTTPS-only, time-limited, response-size-capped.
|
||||
|
||||
Hardened against the obvious abuse of a configurable URL: non-HTTPS schemes
|
||||
are refused unless explicitly allowed, the read is bounded, and every
|
||||
transport failure is normalised to :class:`RelayKeyError`.
|
||||
"""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
*,
|
||||
timeout: float = 15.0,
|
||||
max_bytes: int = _MAX_KEY_BYTES,
|
||||
allow_insecure: bool = False,
|
||||
opener: _Opener = urllib.request.urlopen,
|
||||
) -> None:
|
||||
self._timeout = timeout
|
||||
self._max_bytes = max_bytes
|
||||
self._allow_insecure = allow_insecure
|
||||
self._opener = opener
|
||||
|
||||
def __call__(self, url: str) -> str:
|
||||
scheme = urllib.parse.urlsplit(url).scheme.lower()
|
||||
if scheme != "https" and not (self._allow_insecure and scheme in ("http", "file")):
|
||||
raise RelayKeyError(f"refusing non-HTTPS relay key URL: {url!r}")
|
||||
try:
|
||||
with self._opener(url, timeout=self._timeout) as resp:
|
||||
data = resp.read(self._max_bytes + 1)
|
||||
except (OSError, ValueError) as exc:
|
||||
raise RelayKeyError(f"failed to fetch relay key from {url!r}: {exc}") from exc
|
||||
if len(data) > self._max_bytes:
|
||||
raise RelayKeyError(f"relay key response exceeds {self._max_bytes} bytes")
|
||||
return data.decode("utf-8", "replace")
|
||||
|
||||
|
||||
class RelayKeyProvider:
|
||||
"""Provides the relay :class:`RelayKey`, refreshing past a TTL.
|
||||
|
||||
Thread-safe: concurrent ``get()`` calls serialise on a lock so the key is
|
||||
fetched at most once per TTL window even under contention.
|
||||
"""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
url: str,
|
||||
ttl_seconds: float,
|
||||
fetcher: Fetcher | None = None,
|
||||
clock: Clock = time.monotonic,
|
||||
) -> None:
|
||||
self._url = url
|
||||
self._ttl = ttl_seconds
|
||||
self._fetch = fetcher or HttpsRelayKeyFetcher()
|
||||
self._clock = clock
|
||||
self._lock = threading.Lock()
|
||||
self._key: RelayKey | None = None
|
||||
self._fetched_at = 0.0
|
||||
|
||||
def get(self, *, force: bool = False) -> RelayKey:
|
||||
"""Return the relay key, fetching only when stale or ``force``."""
|
||||
with self._lock:
|
||||
now = self._clock()
|
||||
if not force and self._key is not None and (now - self._fetched_at) < self._ttl:
|
||||
log.debug("relay key reused (age %.0fs)", now - self._fetched_at)
|
||||
return self._key
|
||||
key = parse_relay_pub(self._fetch(self._url))
|
||||
self._key, self._fetched_at = key, now
|
||||
log.info("relay key refreshed from %s", self._url)
|
||||
return key
|
||||
|
||||
|
||||
__all__ = ["Fetcher", "Clock", "HttpsRelayKeyFetcher", "RelayKeyProvider"]
|
||||
@@ -0,0 +1,110 @@
|
||||
# SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
"""Pure domain model: value objects and parsing, no I/O.
|
||||
|
||||
Everything here is deterministic and side-effect free, so it is trivially
|
||||
testable and safe to share across threads (all types are immutable).
|
||||
|
||||
Provenance: in ``Plex Media Server`` 1.43.2.10687, ``RelayController_connect``
|
||||
(``0x12307F2``) downloads ``relay_v1.pub``, splits it into exactly three
|
||||
whitespace tokens (rejecting otherwise -- the "part count incorrect" log), and
|
||||
writes a per-endpoint OpenSSH known_hosts line ``[host]:443 <keytype> <keydata>``
|
||||
into ``relayHostKey.txt``.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass
|
||||
|
||||
from .errors import RelayKeyError
|
||||
|
||||
# Key types OpenSSH recognises. Used to disambiguate a public-key line
|
||||
# ("keytype keydata comment") from a known_hosts line ("host keytype keydata").
|
||||
_KEY_TYPES: frozenset[str] = frozenset(
|
||||
{
|
||||
"ssh-rsa",
|
||||
"ssh-dss",
|
||||
"ssh-ed25519",
|
||||
"ecdsa-sha2-nistp256",
|
||||
"ecdsa-sha2-nistp384",
|
||||
"ecdsa-sha2-nistp521",
|
||||
"sk-ssh-ed25519@openssh.com",
|
||||
"sk-ecdsa-sha2-nistp256@openssh.com",
|
||||
"ssh-rsa-cert-v01@openssh.com",
|
||||
"ssh-ed25519-cert-v01@openssh.com",
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
def known_hosts_endpoint(host: str, port: int = 443) -> str:
|
||||
"""Return the OpenSSH known_hosts host field for ``host:port``.
|
||||
|
||||
OpenSSH uses bracket notation for any non-default port. The binary hardcodes
|
||||
``[host]:443``; keying by the actual port keeps the entry valid for relays on
|
||||
other ports while remaining identical at 443.
|
||||
"""
|
||||
host = host.strip()
|
||||
if not host:
|
||||
raise RelayKeyError("empty relay host")
|
||||
return host if port == 22 else f"[{host}]:{port}"
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class RelayKey:
|
||||
"""An SSH host key: an algorithm and its base64-encoded blob."""
|
||||
|
||||
keytype: str
|
||||
keydata: str
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class HostKey:
|
||||
"""A single OpenSSH known_hosts entry for a relay endpoint."""
|
||||
|
||||
endpoint: str
|
||||
keytype: str
|
||||
keydata: str
|
||||
|
||||
@classmethod
|
||||
def for_endpoint(cls, endpoint: str, key: RelayKey) -> "HostKey":
|
||||
return cls(endpoint=endpoint, keytype=key.keytype, keydata=key.keydata)
|
||||
|
||||
@property
|
||||
def key(self) -> RelayKey:
|
||||
return RelayKey(self.keytype, self.keydata)
|
||||
|
||||
def known_hosts_line(self) -> str:
|
||||
"""The line OpenSSH consumes: ``<endpoint> <keytype> <keydata>``."""
|
||||
return f"{self.endpoint} {self.keytype} {self.keydata}"
|
||||
|
||||
def cache_block(self) -> str:
|
||||
"""PMS ``relayHostKey.txt`` representation: ``# <endpoint>`` + data line."""
|
||||
return f"# {self.endpoint}\n{self.known_hosts_line()}\n"
|
||||
|
||||
|
||||
def parse_relay_pub(body: str) -> RelayKey:
|
||||
"""Parse a ``relay_v1.pub`` payload into a :class:`RelayKey`.
|
||||
|
||||
Accepts both the OpenSSH public-key form (``keytype keydata comment``) and
|
||||
the known_hosts form (``host keytype keydata``), disambiguated by which token
|
||||
is a recognised key type.
|
||||
|
||||
:raises RelayKeyError: if no usable key line is present (mirrors the binary's
|
||||
"part count incorrect" rejection).
|
||||
"""
|
||||
for raw in body.splitlines():
|
||||
line = raw.strip()
|
||||
if not line or line.startswith("#"):
|
||||
continue
|
||||
tokens = line.split()
|
||||
if len(tokens) < 3:
|
||||
raise RelayKeyError(
|
||||
f"relay key: part count incorrect ({len(tokens)} tokens)"
|
||||
)
|
||||
if tokens[0] in _KEY_TYPES: # keytype keydata comment
|
||||
return RelayKey(tokens[0], tokens[1])
|
||||
if tokens[1] in _KEY_TYPES: # host keytype keydata
|
||||
return RelayKey(tokens[1], tokens[2])
|
||||
raise RelayKeyError("relay key: no recognised key type in payload")
|
||||
raise RelayKeyError("relay key: payload contained no key line")
|
||||
|
||||
|
||||
__all__ = ["RelayKey", "HostKey", "known_hosts_endpoint", "parse_relay_pub"]
|
||||
Whitespace-only changes.
@@ -0,0 +1,65 @@
|
||||
# SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
"""Host-key trust management -- composes the key provider and the disk cache.
|
||||
|
||||
This is the seam the controller depends on (the :class:`HostKeyTrust` protocol):
|
||||
"make sure ssh will trust this relay endpoint, and tell me which file to hand
|
||||
it." It owns no I/O of its own; it wires together :class:`RelayKeyProvider`
|
||||
(network + TTL) and :class:`HostKeyCache` (disk), keeping each collaborator
|
||||
single-purpose and independently testable.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import threading
|
||||
from pathlib import Path
|
||||
from typing import Protocol
|
||||
|
||||
from .cache import HostKeyCache
|
||||
from .keys import RelayKeyProvider
|
||||
from .models import HostKey, known_hosts_endpoint
|
||||
|
||||
log = logging.getLogger("plex_relay.store")
|
||||
|
||||
|
||||
class HostKeyTrust(Protocol):
|
||||
"""What the controller requires to make a relay endpoint trusted by ssh."""
|
||||
|
||||
def ensure_trusted(self, host: str, port: int = 443, *, force: bool = False) -> HostKey: ...
|
||||
|
||||
@property
|
||||
def known_hosts_path(self) -> Path: ...
|
||||
|
||||
|
||||
class HostKeyManager:
|
||||
"""Default :class:`HostKeyTrust`: refresh key, pin endpoint, persist on change."""
|
||||
|
||||
def __init__(self, provider: RelayKeyProvider, cache: HostKeyCache) -> None:
|
||||
self._provider = provider
|
||||
self._cache = cache
|
||||
self._lock = threading.Lock()
|
||||
self._entries: dict[str, HostKey] = cache.load()
|
||||
|
||||
@property
|
||||
def known_hosts_path(self) -> Path:
|
||||
return self._cache.path
|
||||
|
||||
@property
|
||||
def entries(self) -> dict[str, HostKey]:
|
||||
with self._lock:
|
||||
return dict(self._entries)
|
||||
|
||||
def ensure_trusted(self, host: str, port: int = 443, *, force: bool = False) -> HostKey:
|
||||
"""Ensure ``relayHostKey.txt`` trusts ``host:port``; persist iff changed."""
|
||||
key = self._provider.get(force=force)
|
||||
endpoint = known_hosts_endpoint(host, port)
|
||||
entry = HostKey.for_endpoint(endpoint, key)
|
||||
with self._lock:
|
||||
if self._entries.get(endpoint) == entry:
|
||||
return entry
|
||||
self._entries[endpoint] = entry
|
||||
self._cache.save(self._entries)
|
||||
log.info("pinned relay host key for %s", endpoint)
|
||||
return entry
|
||||
|
||||
|
||||
__all__ = ["HostKeyTrust", "HostKeyManager"]
|
||||
@@ -0,0 +1,207 @@
|
||||
# SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
"""The relay tunnel adapter: build the ssh argv and run it as a child process.
|
||||
|
||||
``build_ssh_argv`` is a pure function (the byte-for-byte reproduction of the
|
||||
argv ``RelayController_connect`` assembles), kept separate from process control
|
||||
so it is testable without spawning anything.
|
||||
|
||||
Security: the relay credential is delivered to ssh via the ``PLEXTOKEN``
|
||||
environment variable read by a generated ``SSH_ASKPASS`` helper -- it never
|
||||
appears on a command line or in the argv list. The helper is mode ``0o700`` and
|
||||
is removed on stop *and* via a finaliser, so a crash cannot leak it.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import os
|
||||
import stat
|
||||
import subprocess
|
||||
import tempfile
|
||||
import weakref
|
||||
from pathlib import Path
|
||||
from typing import Callable, Mapping, Protocol
|
||||
|
||||
from .config import RelayConfig
|
||||
from .errors import TunnelError
|
||||
|
||||
log = logging.getLogger("plex_relay.tunnel")
|
||||
|
||||
_ASKPASS_SCRIPT = "#!/bin/sh\nprintf '%s' \"$PLEXTOKEN\"\n"
|
||||
|
||||
|
||||
class ProcessHandle(Protocol):
|
||||
"""The slice of ``subprocess.Popen`` the tunnel relies on."""
|
||||
|
||||
def poll(self) -> int | None: ...
|
||||
def terminate(self) -> None: ...
|
||||
def kill(self) -> None: ...
|
||||
def wait(self, timeout: float | None = ...) -> int: ...
|
||||
|
||||
|
||||
#: Launches a child process from an argv + environment. The seam tests stub.
|
||||
Spawner = Callable[[list[str], Mapping[str, str]], ProcessHandle]
|
||||
|
||||
|
||||
def build_ssh_argv(
|
||||
config: RelayConfig, host: str, port: int, known_hosts_path: Path
|
||||
) -> list[str]:
|
||||
"""Assemble the ssh reverse-tunnel argv, exactly as the binary does.
|
||||
|
||||
The known_hosts path uses POSIX separators (identical on Linux; OpenSSH
|
||||
accepts forward slashes everywhere).
|
||||
"""
|
||||
return [
|
||||
config.ssh_binary,
|
||||
"-p", str(port),
|
||||
"-N",
|
||||
"-R", f"0:{config.local_host}:{config.local_port}",
|
||||
"-o", f"UserKnownHostsFile={Path(known_hosts_path).as_posix()}",
|
||||
"-o", "LogLevel=VERBOSE",
|
||||
"-o", "PreferredAuthentications=password",
|
||||
"-o", "PubkeyAuthentication=no",
|
||||
"-l", config.ssh_user,
|
||||
"-F", "/dev/null",
|
||||
host,
|
||||
]
|
||||
|
||||
|
||||
class _Askpass:
|
||||
"""A short-lived, self-cleaning SSH_ASKPASS helper script."""
|
||||
|
||||
def __init__(self, token: str) -> None:
|
||||
fd, name = tempfile.mkstemp(prefix="plex_relay_askpass_", suffix=".sh")
|
||||
try:
|
||||
os.write(fd, _ASKPASS_SCRIPT.encode("ascii"))
|
||||
finally:
|
||||
os.close(fd)
|
||||
self.path = Path(name)
|
||||
self.path.chmod(stat.S_IRWXU) # 0o700
|
||||
self._token = token
|
||||
self._finalizer = weakref.finalize(self, _unlink, self.path)
|
||||
|
||||
def env(self, base: Mapping[str, str]) -> dict[str, str]:
|
||||
env = dict(base)
|
||||
env.update(
|
||||
PLEXTOKEN=self._token,
|
||||
SSH_ASKPASS=str(self.path),
|
||||
SSH_ASKPASS_REQUIRE="force",
|
||||
DISPLAY=base.get("DISPLAY", ":0"),
|
||||
)
|
||||
return env
|
||||
|
||||
def cleanup(self) -> None:
|
||||
self._finalizer()
|
||||
|
||||
|
||||
def _unlink(path: Path) -> None:
|
||||
try:
|
||||
path.unlink()
|
||||
except OSError:
|
||||
pass
|
||||
|
||||
|
||||
def _default_spawner(argv: list[str], env: Mapping[str, str]) -> ProcessHandle:
|
||||
# Detach from the controlling tty so ssh uses SSH_ASKPASS for the password.
|
||||
return subprocess.Popen( # noqa: S603 - argv is fully built; no shell
|
||||
argv,
|
||||
env=dict(env),
|
||||
stdin=subprocess.DEVNULL,
|
||||
stdout=subprocess.PIPE,
|
||||
stderr=subprocess.STDOUT,
|
||||
start_new_session=True,
|
||||
)
|
||||
|
||||
|
||||
class Tunnel(Protocol):
|
||||
"""Lifecycle of a single relay tunnel."""
|
||||
|
||||
@property
|
||||
def host(self) -> str: ...
|
||||
def start(self) -> None: ...
|
||||
def is_alive(self) -> bool: ...
|
||||
def stop(self, timeout: float | None = ...) -> None: ...
|
||||
|
||||
|
||||
class TunnelFactory(Protocol):
|
||||
def __call__(self, host: str, port: int, known_hosts_path: Path) -> Tunnel: ...
|
||||
|
||||
|
||||
class SubprocessTunnel:
|
||||
"""A relay tunnel backed by a child ``ssh`` process."""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
config: RelayConfig,
|
||||
host: str,
|
||||
port: int,
|
||||
known_hosts_path: Path,
|
||||
spawner: Spawner,
|
||||
) -> None:
|
||||
self._config = config
|
||||
self._host = host
|
||||
self._port = port
|
||||
self._known_hosts = Path(known_hosts_path)
|
||||
self._spawn = spawner
|
||||
self._proc: ProcessHandle | None = None
|
||||
self._askpass: _Askpass | None = None
|
||||
|
||||
@property
|
||||
def host(self) -> str:
|
||||
return self._host
|
||||
|
||||
@property
|
||||
def argv(self) -> list[str]:
|
||||
return build_ssh_argv(self._config, self._host, self._port, self._known_hosts)
|
||||
|
||||
def start(self) -> None:
|
||||
if self.is_alive():
|
||||
return
|
||||
askpass = _Askpass(self._config.token)
|
||||
try:
|
||||
log.info("starting relay tunnel to %s:%d", self._host, self._port)
|
||||
self._proc = self._spawn(self.argv, askpass.env(os.environ))
|
||||
except OSError as exc:
|
||||
askpass.cleanup()
|
||||
raise TunnelError(f"failed to launch ssh for {self._host}: {exc}") from exc
|
||||
self._askpass = askpass
|
||||
|
||||
def is_alive(self) -> bool:
|
||||
return self._proc is not None and self._proc.poll() is None
|
||||
|
||||
def stop(self, timeout: float | None = 5.0) -> None:
|
||||
proc, self._proc = self._proc, None
|
||||
if proc is not None:
|
||||
log.info("stopping relay tunnel to %s", self._host)
|
||||
try:
|
||||
proc.terminate()
|
||||
proc.wait(timeout=timeout)
|
||||
except Exception: # noqa: BLE001 - escalate to kill on any wait failure
|
||||
try:
|
||||
proc.kill()
|
||||
except OSError:
|
||||
pass
|
||||
if self._askpass is not None:
|
||||
self._askpass.cleanup()
|
||||
self._askpass = None
|
||||
|
||||
|
||||
class SubprocessTunnelFactory:
|
||||
"""Default :class:`TunnelFactory` producing :class:`SubprocessTunnel`."""
|
||||
|
||||
def __init__(self, config: RelayConfig, spawner: Spawner | None = None) -> None:
|
||||
self._config = config
|
||||
self._spawner = spawner or _default_spawner
|
||||
|
||||
def __call__(self, host: str, port: int, known_hosts_path: Path) -> SubprocessTunnel:
|
||||
return SubprocessTunnel(self._config, host, port, known_hosts_path, self._spawner)
|
||||
|
||||
|
||||
__all__ = [
|
||||
"ProcessHandle",
|
||||
"Spawner",
|
||||
"Tunnel",
|
||||
"TunnelFactory",
|
||||
"SubprocessTunnel",
|
||||
"SubprocessTunnelFactory",
|
||||
"build_ssh_argv",
|
||||
]
|
||||
@@ -0,0 +1,59 @@
|
||||
# SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
import os
|
||||
import stat
|
||||
|
||||
import pytest
|
||||
|
||||
from plex_relay.cache import HostKeyCache, parse_known_hosts
|
||||
from plex_relay.errors import HostKeyCacheError
|
||||
from plex_relay.models import HostKey
|
||||
|
||||
|
||||
def _entries():
|
||||
return {
|
||||
"[1.2.3.4]:443": HostKey("[1.2.3.4]:443", "ssh-ed25519", "AAAAfirst"),
|
||||
"[5.6.7.8]:443": HostKey("[5.6.7.8]:443", "ssh-ed25519", "AAAAsecond"),
|
||||
}
|
||||
|
||||
|
||||
def test_save_load_roundtrip(tmp_path):
|
||||
cache = HostKeyCache(tmp_path / "relayHostKey.txt")
|
||||
cache.save(_entries())
|
||||
loaded = HostKeyCache(tmp_path / "relayHostKey.txt").load()
|
||||
assert loaded == _entries()
|
||||
|
||||
|
||||
def test_save_is_sorted_and_blocked(tmp_path):
|
||||
p = tmp_path / "relayHostKey.txt"
|
||||
HostKeyCache(p).save(_entries())
|
||||
text = p.read_text()
|
||||
assert text.index("[1.2.3.4]") < text.index("[5.6.7.8]")
|
||||
assert "# [1.2.3.4]:443\n[1.2.3.4]:443 ssh-ed25519 AAAAfirst\n" in text
|
||||
|
||||
|
||||
@pytest.mark.skipif(os.name != "posix", reason="POSIX file modes only")
|
||||
def test_save_is_0600(tmp_path):
|
||||
p = tmp_path / "relayHostKey.txt"
|
||||
HostKeyCache(p).save(_entries())
|
||||
assert stat.S_IMODE(os.stat(p).st_mode) == 0o600
|
||||
|
||||
|
||||
def test_missing_file_loads_empty(tmp_path):
|
||||
assert HostKeyCache(tmp_path / "nope.txt").load() == {}
|
||||
|
||||
|
||||
def test_malformed_file_is_rebuilt_empty(tmp_path):
|
||||
p = tmp_path / "relayHostKey.txt"
|
||||
p.write_text("not a comment\ngarbage line\n")
|
||||
assert HostKeyCache(p).load() == {}
|
||||
assert p.read_text() == ""
|
||||
|
||||
|
||||
@pytest.mark.parametrize("lines", [
|
||||
["data without marker"],
|
||||
["# marker only"],
|
||||
["# marker", "too many tokens here now"],
|
||||
])
|
||||
def test_parse_rejects_malformed(lines):
|
||||
with pytest.raises(HostKeyCacheError):
|
||||
parse_known_hosts(lines)
|
||||
@@ -0,0 +1,44 @@
|
||||
# SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
import dataclasses
|
||||
|
||||
import pytest
|
||||
|
||||
from plex_relay.config import RelayConfig
|
||||
from plex_relay.errors import ConfigError
|
||||
|
||||
|
||||
def test_valid_config_and_cache_path(tmp_path):
|
||||
cfg = RelayConfig(token="t", ssh_user="u", data_dir=tmp_path)
|
||||
assert cfg.cache_path == tmp_path / "relayHostKey.txt"
|
||||
assert cfg.gating_ok() is True
|
||||
|
||||
|
||||
@pytest.mark.parametrize("kwargs", [
|
||||
{"token": "", "ssh_user": "u"},
|
||||
{"token": "t", "ssh_user": ""},
|
||||
{"token": "t", "ssh_user": "u", "local_port": 0},
|
||||
{"token": "t", "ssh_user": "u", "local_port": 70000},
|
||||
{"token": "t", "ssh_user": "u", "key_ttl_seconds": 0},
|
||||
{"token": "t", "ssh_user": "u", "reap_interval_seconds": -1},
|
||||
])
|
||||
def test_invalid_config_raises(kwargs):
|
||||
with pytest.raises(ConfigError):
|
||||
RelayConfig(**kwargs)
|
||||
|
||||
|
||||
def test_gating_requires_all_three():
|
||||
base = dict(token="t", ssh_user="u")
|
||||
assert RelayConfig(**base, relay_enabled=False).gating_ok() is False
|
||||
assert RelayConfig(**base, published=False).gating_ok() is False
|
||||
assert RelayConfig(**base, signed_in=False).gating_ok() is False
|
||||
|
||||
|
||||
def test_token_is_not_in_repr():
|
||||
cfg = RelayConfig(token="SUPERSECRET", ssh_user="u")
|
||||
assert "SUPERSECRET" not in repr(cfg)
|
||||
|
||||
|
||||
def test_config_is_frozen():
|
||||
cfg = RelayConfig(token="t", ssh_user="u")
|
||||
with pytest.raises(dataclasses.FrozenInstanceError):
|
||||
cfg.local_port = 1 # type: ignore[misc]
|
||||
@@ -0,0 +1,157 @@
|
||||
# SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
|
||||
from plex_relay.config import RelayConfig
|
||||
from plex_relay.controller import RelayController
|
||||
from plex_relay.errors import RelayError, TunnelError
|
||||
from plex_relay.models import HostKey, RelayKey, known_hosts_endpoint
|
||||
|
||||
|
||||
class FakeTunnel:
|
||||
def __init__(self, host, fail=False):
|
||||
self._host = host
|
||||
self._fail = fail
|
||||
self.alive = False
|
||||
self.stopped = False
|
||||
|
||||
@property
|
||||
def host(self):
|
||||
return self._host
|
||||
|
||||
def start(self):
|
||||
if self._fail:
|
||||
raise TunnelError("spawn failed")
|
||||
self.alive = True
|
||||
|
||||
def is_alive(self):
|
||||
return self.alive
|
||||
|
||||
def stop(self, timeout=None):
|
||||
self.stopped = True
|
||||
self.alive = False
|
||||
|
||||
|
||||
class FakeFactory:
|
||||
def __init__(self):
|
||||
self.fail = False
|
||||
self.created = []
|
||||
|
||||
def __call__(self, host, port, known_hosts_path):
|
||||
t = FakeTunnel(host, fail=self.fail)
|
||||
self.created.append(t)
|
||||
return t
|
||||
|
||||
|
||||
class FakeTrust:
|
||||
def __init__(self):
|
||||
self.calls = []
|
||||
|
||||
def ensure_trusted(self, host, port=443, *, force=False):
|
||||
self.calls.append((host, port))
|
||||
return HostKey.for_endpoint(known_hosts_endpoint(host, port), RelayKey("ssh-ed25519", "k"))
|
||||
|
||||
@property
|
||||
def known_hosts_path(self):
|
||||
return Path("/k")
|
||||
|
||||
|
||||
def build(**cfgkw):
|
||||
cfg = RelayConfig(token="t", ssh_user="u", reap_interval_seconds=999.0, **cfgkw)
|
||||
trust, factory = FakeTrust(), FakeFactory()
|
||||
return RelayController(cfg, trust, factory), trust, factory
|
||||
|
||||
|
||||
def test_connect_tracks_and_trusts():
|
||||
ctrl, trust, factory = build()
|
||||
try:
|
||||
assert ctrl.connect("relay.example", 443) is True
|
||||
assert ctrl.active_hosts == ["relay.example"]
|
||||
assert trust.calls == [("relay.example", 443)]
|
||||
assert len(factory.created) == 1
|
||||
finally:
|
||||
ctrl.stop()
|
||||
|
||||
|
||||
def test_connect_dedup_when_alive():
|
||||
ctrl, _, factory = build()
|
||||
try:
|
||||
assert ctrl.connect("h") is True
|
||||
assert ctrl.connect("h") is False
|
||||
assert len(factory.created) == 1
|
||||
finally:
|
||||
ctrl.stop()
|
||||
|
||||
|
||||
def test_reconnect_after_death():
|
||||
ctrl, _, factory = build()
|
||||
try:
|
||||
ctrl.connect("h")
|
||||
factory.created[0].alive = False # tunnel died
|
||||
assert ctrl.connect("h") is True
|
||||
assert len(factory.created) == 2
|
||||
finally:
|
||||
ctrl.stop()
|
||||
|
||||
|
||||
def test_connect_raises_on_tunnel_failure():
|
||||
ctrl, _, factory = build()
|
||||
factory.fail = True
|
||||
with pytest.raises(RelayError):
|
||||
ctrl.connect("h")
|
||||
assert ctrl.active_hosts == []
|
||||
ctrl.stop()
|
||||
|
||||
|
||||
def test_start_relay_blocked_by_gating():
|
||||
for kw in ({"relay_enabled": False}, {"published": False}, {"signed_in": False}):
|
||||
ctrl, trust, _ = build(**kw)
|
||||
assert ctrl.start_relay("h") is False
|
||||
assert trust.calls == []
|
||||
ctrl.stop()
|
||||
|
||||
|
||||
def test_start_relay_is_resilient_to_failure():
|
||||
ctrl, _, factory = build()
|
||||
factory.fail = True
|
||||
assert ctrl.start_relay("h") is False # logs + swallows, does not raise
|
||||
ctrl.stop()
|
||||
|
||||
|
||||
def test_start_relay_ok():
|
||||
ctrl, _, _ = build()
|
||||
try:
|
||||
assert ctrl.start_relay("h", 443) is True
|
||||
assert ctrl.active_hosts == ["h"]
|
||||
finally:
|
||||
ctrl.stop()
|
||||
|
||||
|
||||
def test_reaper_removes_dead():
|
||||
ctrl, _, factory = build()
|
||||
try:
|
||||
ctrl.connect("a")
|
||||
ctrl.connect("b")
|
||||
factory.created[0].alive = False
|
||||
assert ctrl.reap_once() == ["a"]
|
||||
assert ctrl.active_hosts == ["b"]
|
||||
assert factory.created[0].stopped is True
|
||||
finally:
|
||||
ctrl.stop()
|
||||
|
||||
|
||||
def test_stop_terminates_all_and_closes():
|
||||
ctrl, _, factory = build()
|
||||
ctrl.connect("a")
|
||||
ctrl.connect("b")
|
||||
ctrl.stop()
|
||||
assert ctrl.active_hosts == []
|
||||
assert all(t.stopped for t in factory.created)
|
||||
with pytest.raises(RelayError):
|
||||
ctrl.connect("c")
|
||||
|
||||
|
||||
def test_from_config_builds_controller(tmp_path):
|
||||
cfg = RelayConfig(token="t", ssh_user="u", data_dir=tmp_path)
|
||||
assert isinstance(RelayController.from_config(cfg), RelayController)
|
||||
@@ -0,0 +1,70 @@
|
||||
# SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
import pytest
|
||||
|
||||
from plex_relay.errors import RelayKeyError
|
||||
from plex_relay.keys import HttpsRelayKeyFetcher, RelayKeyProvider
|
||||
from plex_relay.models import RelayKey
|
||||
|
||||
PUB = "ssh-ed25519 AAAAkeydata comment"
|
||||
KEY = RelayKey("ssh-ed25519", "AAAAkeydata")
|
||||
|
||||
|
||||
def test_provider_respects_ttl_and_force():
|
||||
calls = []
|
||||
now = [1000.0]
|
||||
provider = RelayKeyProvider(
|
||||
"https://x/relay_v1.pub", 86_400.0,
|
||||
fetcher=lambda url: (calls.append(url), PUB)[1],
|
||||
clock=lambda: now[0],
|
||||
)
|
||||
assert provider.get() == KEY
|
||||
now[0] += 3600 # within TTL -> reuse
|
||||
provider.get()
|
||||
assert len(calls) == 1
|
||||
now[0] += 86_400 # past TTL -> refetch
|
||||
provider.get()
|
||||
assert len(calls) == 2
|
||||
provider.get(force=True) # force -> refetch
|
||||
assert len(calls) == 3
|
||||
|
||||
|
||||
# --- HttpsRelayKeyFetcher ---------------------------------------------------
|
||||
|
||||
class _FakeResp:
|
||||
def __init__(self, data: bytes):
|
||||
self._data = data
|
||||
|
||||
def read(self, n: int = -1) -> bytes:
|
||||
return self._data[:n] if n >= 0 else self._data
|
||||
|
||||
def __enter__(self):
|
||||
return self
|
||||
|
||||
def __exit__(self, *exc):
|
||||
return False
|
||||
|
||||
|
||||
def test_fetcher_rejects_non_https():
|
||||
f = HttpsRelayKeyFetcher(opener=lambda *a, **k: _FakeResp(b""))
|
||||
with pytest.raises(RelayKeyError):
|
||||
f("http://insecure/relay_v1.pub")
|
||||
|
||||
|
||||
def test_fetcher_allows_insecure_when_opted_in():
|
||||
f = HttpsRelayKeyFetcher(allow_insecure=True, opener=lambda *a, **k: _FakeResp(PUB.encode()))
|
||||
assert f("file:///tmp/relay_v1.pub") == PUB
|
||||
|
||||
|
||||
def test_fetcher_caps_response_size():
|
||||
big = b"x" * 100
|
||||
f = HttpsRelayKeyFetcher(max_bytes=10, opener=lambda *a, **k: _FakeResp(big))
|
||||
with pytest.raises(RelayKeyError, match="exceeds"):
|
||||
f("https://x/relay_v1.pub")
|
||||
|
||||
|
||||
def test_fetcher_wraps_transport_errors():
|
||||
def boom(*a, **k):
|
||||
raise OSError("connection refused")
|
||||
f = HttpsRelayKeyFetcher(opener=boom)
|
||||
with pytest.raises(RelayKeyError, match="failed to fetch"):
|
||||
f("https://x/relay_v1.pub")
|
||||
@@ -0,0 +1,45 @@
|
||||
# SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
import pytest
|
||||
|
||||
from plex_relay.errors import RelayKeyError
|
||||
from plex_relay.models import HostKey, RelayKey, known_hosts_endpoint, parse_relay_pub
|
||||
|
||||
PUB = "ssh-ed25519 AAAAC3NzaC1lZDI1NTE5AAAAIabc relay@plex" # keytype keydata comment
|
||||
KH = "* ssh-ed25519 AAAAC3NzaC1lZDI1NTE5AAAAIabc" # host keytype keydata
|
||||
EXPECT = RelayKey("ssh-ed25519", "AAAAC3NzaC1lZDI1NTE5AAAAIabc")
|
||||
|
||||
|
||||
def test_parse_pubkey_form():
|
||||
assert parse_relay_pub(PUB) == EXPECT
|
||||
|
||||
|
||||
def test_parse_known_hosts_form():
|
||||
assert parse_relay_pub(KH) == EXPECT
|
||||
|
||||
|
||||
def test_parse_skips_comments_and_blanks():
|
||||
assert parse_relay_pub(f"# header\n\n{PUB}\n") == EXPECT
|
||||
|
||||
|
||||
@pytest.mark.parametrize("bad", ["ssh-ed25519 onlytwo", "aaa bbb ccc", "", "# only comment\n"])
|
||||
def test_parse_rejects_bad_payloads(bad):
|
||||
with pytest.raises(RelayKeyError):
|
||||
parse_relay_pub(bad)
|
||||
|
||||
|
||||
def test_endpoint_bracket_notation():
|
||||
assert known_hosts_endpoint("1.2.3.4", 443) == "[1.2.3.4]:443"
|
||||
assert known_hosts_endpoint("relay.example", 2222) == "[relay.example]:2222"
|
||||
assert known_hosts_endpoint("relay.example", 22) == "relay.example"
|
||||
|
||||
|
||||
def test_endpoint_rejects_empty_host():
|
||||
with pytest.raises(RelayKeyError):
|
||||
known_hosts_endpoint(" ", 443)
|
||||
|
||||
|
||||
def test_hostkey_serialization():
|
||||
hk = HostKey.for_endpoint("[1.2.3.4]:443", EXPECT)
|
||||
assert hk.known_hosts_line() == "[1.2.3.4]:443 ssh-ed25519 AAAAC3NzaC1lZDI1NTE5AAAAIabc"
|
||||
assert hk.cache_block() == f"# [1.2.3.4]:443\n{hk.known_hosts_line()}\n"
|
||||
assert hk.key == EXPECT
|
||||
@@ -0,0 +1,61 @@
|
||||
# SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
from pathlib import Path
|
||||
|
||||
from plex_relay.keys import RelayKeyProvider
|
||||
from plex_relay.store import HostKeyManager
|
||||
|
||||
|
||||
class FakeCache:
|
||||
def __init__(self):
|
||||
self.entries = {}
|
||||
self.saves = 0
|
||||
|
||||
@property
|
||||
def path(self) -> Path:
|
||||
return Path("/tmp/relayHostKey.txt")
|
||||
|
||||
def load(self):
|
||||
return dict(self.entries)
|
||||
|
||||
def save(self, entries):
|
||||
self.saves += 1
|
||||
self.entries = dict(entries)
|
||||
|
||||
|
||||
def _provider(text_box):
|
||||
return RelayKeyProvider("https://x/relay_v1.pub", 86_400.0,
|
||||
fetcher=lambda u: text_box[0], clock=lambda: 0.0)
|
||||
|
||||
|
||||
def test_ensure_trusted_pins_and_persists():
|
||||
cache = FakeCache()
|
||||
mgr = HostKeyManager(_provider(["ssh-ed25519 AAAAfirst c"]), cache)
|
||||
entry = mgr.ensure_trusted("1.2.3.4", 443)
|
||||
assert entry.known_hosts_line() == "[1.2.3.4]:443 ssh-ed25519 AAAAfirst"
|
||||
assert cache.saves == 1
|
||||
assert "[1.2.3.4]:443" in cache.entries
|
||||
|
||||
|
||||
def test_ensure_trusted_is_idempotent():
|
||||
cache = FakeCache()
|
||||
mgr = HostKeyManager(_provider(["ssh-ed25519 AAAAfirst c"]), cache)
|
||||
mgr.ensure_trusted("1.2.3.4", 443)
|
||||
mgr.ensure_trusted("1.2.3.4", 443) # unchanged -> no extra write
|
||||
assert cache.saves == 1
|
||||
|
||||
|
||||
def test_ensure_trusted_rewrites_on_key_change():
|
||||
cache = FakeCache()
|
||||
box = ["ssh-ed25519 AAAAfirst c"]
|
||||
mgr = HostKeyManager(_provider(box), cache)
|
||||
mgr.ensure_trusted("h", 443)
|
||||
box[0] = "ssh-ed25519 AAAAsecond c"
|
||||
entry = mgr.ensure_trusted("h", 443, force=True)
|
||||
assert entry.keydata == "AAAAsecond"
|
||||
assert cache.saves == 2
|
||||
|
||||
|
||||
def test_known_hosts_path_is_cache_path():
|
||||
cache = FakeCache()
|
||||
mgr = HostKeyManager(_provider(["ssh-ed25519 k c"]), cache)
|
||||
assert mgr.known_hosts_path == cache.path
|
||||
@@ -0,0 +1,91 @@
|
||||
# SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
|
||||
from plex_relay.config import RelayConfig
|
||||
from plex_relay.errors import TunnelError
|
||||
from plex_relay.tunnel import SubprocessTunnel, SubprocessTunnelFactory, build_ssh_argv
|
||||
|
||||
|
||||
def cfg(**kw):
|
||||
base = dict(token="secret", ssh_user="machineid", local_host="127.0.0.1", local_port=32400)
|
||||
base.update(kw)
|
||||
return RelayConfig(**base)
|
||||
|
||||
|
||||
def test_argv_matches_binary_layout():
|
||||
argv = build_ssh_argv(cfg(), "relay.example", 443, Path("/data/relayHostKey.txt"))
|
||||
assert argv == [
|
||||
"ssh", "-p", "443", "-N", "-R", "0:127.0.0.1:32400",
|
||||
"-o", "UserKnownHostsFile=/data/relayHostKey.txt",
|
||||
"-o", "LogLevel=VERBOSE",
|
||||
"-o", "PreferredAuthentications=password",
|
||||
"-o", "PubkeyAuthentication=no",
|
||||
"-l", "machineid", "-F", "/dev/null", "relay.example",
|
||||
]
|
||||
|
||||
|
||||
def test_argv_honours_port_and_target():
|
||||
argv = build_ssh_argv(cfg(local_host="10.0.0.5", local_port=32500), "h", 2222, Path("/k"))
|
||||
assert "0:10.0.0.5:32500" in argv
|
||||
assert argv[argv.index("-p") + 1] == "2222"
|
||||
|
||||
|
||||
class FakeProc:
|
||||
def __init__(self):
|
||||
self.alive = True
|
||||
self.terminated = False
|
||||
|
||||
def poll(self):
|
||||
return None if self.alive else 0
|
||||
|
||||
def terminate(self):
|
||||
self.terminated = True
|
||||
self.alive = False
|
||||
|
||||
def kill(self):
|
||||
self.alive = False
|
||||
|
||||
def wait(self, timeout=None):
|
||||
self.alive = False
|
||||
return 0
|
||||
|
||||
|
||||
def test_tunnel_start_sets_secret_env_and_cleans_askpass():
|
||||
captured = {}
|
||||
|
||||
def spawner(argv, env):
|
||||
captured["argv"] = argv
|
||||
captured["env"] = dict(env)
|
||||
return FakeProc()
|
||||
|
||||
t = SubprocessTunnel(cfg(), "relay.example", 443, Path("/k"), spawner)
|
||||
t.start()
|
||||
assert t.is_alive()
|
||||
assert captured["env"]["PLEXTOKEN"] == "secret"
|
||||
assert "SSH_ASKPASS" in captured["env"]
|
||||
askpass = Path(captured["env"]["SSH_ASKPASS"])
|
||||
assert askpass.exists()
|
||||
assert "secret" not in captured["argv"] # never on the command line
|
||||
t.stop()
|
||||
assert not t.is_alive()
|
||||
assert not askpass.exists() # helper removed on stop
|
||||
|
||||
|
||||
def test_tunnel_start_wraps_spawn_failure():
|
||||
def boom(argv, env):
|
||||
raise OSError("ssh not found")
|
||||
t = SubprocessTunnel(cfg(), "h", 443, Path("/k"), boom)
|
||||
with pytest.raises(TunnelError):
|
||||
t.start()
|
||||
assert not t.is_alive()
|
||||
|
||||
|
||||
def test_factory_builds_tunnel():
|
||||
factory = SubprocessTunnelFactory(cfg(), spawner=lambda a, e: FakeProc())
|
||||
t = factory("h", 443, Path("/k"))
|
||||
assert t.host == "h"
|
||||
t.start()
|
||||
assert t.is_alive()
|
||||
t.stop()
|
||||
Reference in new issue
Block a user