Compare commits

...

14 Commits

Author SHA1 Message Date
Alexandre
2afef8cf05 fix(claude_desktop): preserve backend logging 2026-08-02 17:58:22 +02:00
Alexandre
b59beb2aad fix(claude_desktop): keep Headroom logs in addon output 2026-08-02 17:57:47 +02:00
Alexandre
4f5d0fa6b6 fix(claude_desktop): fail open on Headroom API changes 2026-08-02 17:57:33 +02:00
github-actions[bot]
2497729f14 docs(claude_desktop): document memory reduction 2026-08-02 15:55:56 +00:00
Alexandre
57aac487c5 ci: prepare Claude Desktop changelog entry 2026-08-02 17:55:41 +02:00
Alexandre
0664e00ad1 fix(claude_desktop): expire stale keepalive sessions 2026-08-02 17:53:40 +02:00
Alexandre
307e2687f2 fix(claude_desktop): resolve persistent bashrc explicitly 2026-08-02 17:52:02 +02:00
Alexandre
c7583ac4ce fix(claude_desktop): avoid recursive cache chown 2026-08-02 17:51:49 +02:00
Alexandre
519fc02b60 fix(claude_desktop): remove RAM-backed tmpfs 2026-08-02 17:50:47 +02:00
Alexandre
4084a57b6f fix(claude_desktop): start Headroom backend on demand 2026-08-02 17:50:28 +02:00
Alexandre
504d86fa50 fix(claude_desktop): move application cache off tmpfs 2026-08-02 17:50:16 +02:00
Alexandre
0ce79bf3ec fix(claude_desktop): route Headroom MCP through adapter 2026-08-02 17:50:04 +02:00
Alexandre
e464608155 fix(claude_desktop): delegate MCP compression 2026-08-02 17:49:54 +02:00
Alexandre
800e4a9d89 fix(claude_desktop): lazy-load Headroom proxy 2026-08-02 17:49:37 +02:00
7 changed files with 686 additions and 23 deletions

View File

@@ -1,3 +1,9 @@
## 1.38 (02-08-2026)
- Reduce Claude Desktop add-on RAM use in three places without removing Headroom, RTK, TokenSave, Cowork, Dispatch, or the streamed desktop. Headroom's heavy HTTP proxy is no longer started at container boot: a standard-library TCP gate stays on port 8787, starts the real proxy on the first request, and stops it after 15 minutes without traffic (`HEADROOM_IDLE_TIMEOUT_SECONDS` remains overridable through `env_vars`). This releases the proxy's Python, ONNX Runtime, tokenizer, and Kompress model allocations while the add-on is idle; a later request starts a clean backend transparently.
- Prevent Headroom from loading a second Kompress model inside every Claude Desktop/Claude Code MCP process. The add-on now intercepts `headroom mcp serve` with a lightweight adapter that retains upstream MCP retrieval/statistics behavior but delegates `headroom_compress` to the shared proxy's loopback `/v1/compress` endpoint. The MCP process explicitly disables local Kompress and never imports `headroom.compress`, so only the proxy backend can own the ML runtime.
- Remove the add-on-wide `tmpfs: true` mount and undo the shared Selkies script's `/tmp/cache` redirection for this add-on. Electron/Chromium, Mesa, and application caches now live under the selected persistent home (`$HOME/.cache`) as reclaimable filesystem cache instead of RAM-backed cgroup shmem. Runtime sockets and XDG runtime state remain under `/run`.
## ubunturesolute-version-3a10bef7 (2026-08-01)
- Update to latest version from linuxserver/docker-baseimage-selkies (changelog : https://github.com/linuxserver/docker-baseimage-selkies/releases)

View File

@@ -119,8 +119,7 @@ schema:
tokensave_project_paths:
- str
slug: claude_desktop
tmpfs: true
udev: true
url: https://github.com/alexbelgium/hassio-addons
version: "ubunturesolute-version-3a10bef7"
version: "1.38"
video: true

View File

@@ -0,0 +1,58 @@
#!/usr/bin/with-contenv bashio
# shellcheck shell=bash
set -e
# 20-folders.sh is shared with the Webtop add-ons and historically redirects
# XDG_CACHE_HOME to /tmp/cache. This Claude-specific follow-up runs before the
# graphical longruns and moves general application caches back under the
# persistent home. Removing config.yaml's `tmpfs: true` then ensures Chromium,
# Electron, Mesa and other cache pages are reclaimable filesystem cache instead
# of permanently charged cgroup shmem.
LOCATION="$(getent passwd abc 2> /dev/null | cut -d: -f6 || true)"
if [ -z "$LOCATION" ] || [ "$LOCATION" = "/" ]; then
bashio::log.warning "Unable to resolve abc home; leaving XDG cache configuration unchanged"
exit 0
fi
CACHE_DIR="$LOCATION/.cache"
if [ -L "$CACHE_DIR" ]; then
rm -f "$CACHE_DIR"
fi
mkdir -p "$CACHE_DIR"
chown "$(id -u abc):$(id -g abc)" "$CACHE_DIR"
chmod 700 "$CACHE_DIR"
CACHE_DIR="$CACHE_DIR" LOCATION="$LOCATION" python3 - <<'PY'
import os
import re
from pathlib import Path
cache = os.environ["CACHE_DIR"]
quoted = cache.replace("\\", "\\\\").replace('"', '\\"')
replacement = f'export XDG_CACHE_HOME="{quoted}"'
for path in Path("/etc/s6-overlay/s6-rc.d").glob("*/run"):
try:
text = path.read_text()
except (OSError, UnicodeDecodeError):
continue
updated = re.sub(r"^export XDG_CACHE_HOME=.*$", replacement, text, flags=re.MULTILINE)
if updated != text:
path.write_text(updated)
bashrc = Path(os.environ["LOCATION"]) / ".bashrc"
if bashrc.exists():
text = bashrc.read_text()
updated = re.sub(r"^export XDG_CACHE_HOME=.*$", replacement, text, flags=re.MULTILINE)
if updated != text:
bashrc.write_text(updated)
PY
S6_ENVDIR="/run/s6/container_environment"
mkdir -p "$S6_ENVDIR"
printf '%s' "$CACHE_DIR" > "$S6_ENVDIR/XDG_CACHE_HOME"
# Safe here: no graphical longrun has started yet, and the former directory was
# only a boot-created target for the now-removed persistent-home symlink.
rm -rf /tmp/cache
bashio::log.info "Application cache moved from RAM-backed /tmp to $CACHE_DIR"

View File

@@ -0,0 +1,59 @@
#!/usr/bin/with-contenv bashio
# shellcheck shell=bash
set -e
# Keep every existing Headroom CLI command unchanged, but intercept the MCP
# server entrypoint so it uses the add-on's single-runtime adapter. The real
# binary is intentionally outside /usr/local/bin; writing the wrapper there
# makes it the command 82-claude_tools.sh registers in Claude Desktop/Code.
REAL_HEADROOM=""
for candidate in /usr/bin/headroom /lsiopy/bin/headroom; do
if [ -x "$candidate" ]; then
REAL_HEADROOM="$candidate"
break
fi
done
if [ -z "$REAL_HEADROOM" ]; then
bashio::log.warning "Headroom executable was not found; MCP adapter wrapper was not installed"
exit 0
fi
# The adapter intentionally uses a small, stable subset of the upstream MCP
# server. Since Headroom is installed unpinned at image build time, verify that
# subset before replacing the command. A future incompatible Headroom release
# therefore keeps its native MCP server instead of breaking Claude startup.
if ! /lsiopy/bin/python3 - <<'PY'
from headroom.ccr.mcp_server import HeadroomMCPServer
for name in ("run_stdio", "cleanup", "_compress_content"):
if not callable(getattr(HeadroomMCPServer, name, None)):
raise SystemExit(f"HeadroomMCPServer.{name} is unavailable")
PY
then
bashio::log.warning "Installed Headroom is incompatible with the single-runtime MCP adapter; preserving the native MCP server"
exit 0
fi
wrapper="$(mktemp /usr/local/bin/.headroom-wrapper.XXXXXX)"
cleanup() {
rm -f "$wrapper"
}
trap cleanup EXIT
cat > "$wrapper" <<EOF
#!/bin/sh
REAL_HEADROOM="$REAL_HEADROOM"
if [ "\${1:-}" = "mcp" ] && [ "\${2:-}" = "serve" ]; then
shift 2
unset HF_HOME
export HEADROOM_DISABLE_KOMPRESS=1
exec /lsiopy/bin/python3 /usr/local/bin/headroom-mcp-proxy.py "\$@"
fi
exec "\$REAL_HEADROOM" "\$@"
EOF
chmod 0755 "$wrapper"
mv -f "$wrapper" /usr/local/bin/headroom
trap - EXIT
bashio::log.info "Headroom MCP compression is delegated to the shared lazy proxy runtime"

View File

@@ -1,32 +1,43 @@
#!/usr/bin/with-contenv bashio
# Headroom optimization proxy — local backend for Claude Desktop MCP and Claude Code.
# Headroom optimization proxy — lazy backend for Claude Desktop MCP and Claude Code.
declare port=8787
declare backend_port=8789
declare host=127.0.0.1
# The dashboard is unauthenticated. Keep it container-local by default and bind all
# interfaces only when the user explicitly opts in and maps port 8787.
# The dashboard is unauthenticated. Keep the gate container-local by default and
# bind all interfaces only when the user explicitly opts in and maps port 8787.
if bashio::config.true 'expose_headroom_dashboard'; then
host=0.0.0.0
fi
if bashio::config.true 'install_headroom' && command -v headroom > /dev/null 2>&1; then
# Kompress (the ONNX compression engine) needs its model in the local HF cache: the
# proxy's startup preload is deliberately cache-only, and the default HF cache lands
# under ~/.cache, which the add-on points at tmpfs (/tmp/cache) — wiped on every
# restart. Without a warm persistent cache the proxy ran forever in "deferred" mode
# and recorded zero compression savings. Point the cache at persistent storage;
# nothing else is needed here — the proxy's own request path already downloads a
# missing model in the background on first use (ensure_background_load) and passes
# requests through uncompressed until it lands, so this self-heals within a couple of
# requests on the first boot and loads instantly (eager preload) on every boot after.
# A synchronous pre-warm was tried here and removed: it blocked the port bind for up
# to the download's duration, which left the settings-managed ANTHROPIC_BASE_URL
# (see 82-claude_tools.sh) pointing at a proxy that wasn't listening yet.
export HF_HOME="${HOME}/.headroom/hf"
mkdir -p "$HF_HOME"
chown abc:abc "$HF_HOME" 2> /dev/null || true
bashio::log.info "svc-headroom: starting local Headroom proxy on ${host}:${port}"
exec s6-setuidgid abc headroom proxy --host "${host}" --port "${port}" --code-aware
if bashio::config.true 'install_headroom'; then
real_headroom=""
for candidate in /usr/bin/headroom /lsiopy/bin/headroom; do
if [ -x "$candidate" ]; then
real_headroom="$candidate"
break
fi
done
if [ -n "$real_headroom" ]; then
# Keep model artifacts persistent, but do not import Headroom or load the
# model in this longrun. The standard-library gate starts the real proxy
# on the first request and terminates it after the idle timeout, releasing
# Python/ONNX/model allocations. HEADROOM_IDLE_TIMEOUT_SECONDS is
# overridable through env_vars; 900 seconds is the default.
export HF_HOME="${HOME}/.headroom/hf"
mkdir -p "$HF_HOME" "${HOME}/.headroom"
chown abc:abc "${HOME}/.headroom" "$HF_HOME" 2> /dev/null || true
bashio::log.info "svc-headroom: starting lazy gate on ${host}:${port} (backend ${backend_port}, idle timeout ${HEADROOM_IDLE_TIMEOUT_SECONDS:-900}s)"
exec s6-setuidgid abc env \
HEADROOM_REAL_BIN="$real_headroom" \
HEADROOM_GATE_HOST="$host" \
HEADROOM_GATE_PORT="$port" \
HEADROOM_BACKEND_HOST=127.0.0.1 \
HEADROOM_BACKEND_PORT="$backend_port" \
HEADROOM_IDLE_TIMEOUT_SECONDS="${HEADROOM_IDLE_TIMEOUT_SECONDS:-900}" \
HF_HOME="$HF_HOME" \
/lsiopy/bin/python3 /usr/local/bin/headroom-proxy-gate.py
fi
fi
bashio::log.info "svc-headroom: install_headroom disabled or headroom not found; idling"

View File

@@ -0,0 +1,142 @@
#!/lsiopy/bin/python3
"""Run Headroom MCP with compression delegated to the single proxy backend.
Upstream Headroom MCP normally performs ``headroom_compress`` in its own
process. That imports the compression pipeline and can load a second copy of the
Kompress ONNX model in addition to the HTTP proxy. This adapter preserves the
upstream MCP protocol and retrieve/stats implementations, but replaces only its
local compression method with a call to the proxy's loopback-only
``/v1/compress`` endpoint.
The proxy gate starts the heavy backend on this first request and later unloads
it after the configured idle timeout. The MCP process therefore remains a
lightweight protocol bridge and never imports ``headroom.compress``.
"""
from __future__ import annotations
import argparse
import asyncio
import json
import os
import sys
from types import MethodType
from typing import Any
def parse_args() -> argparse.Namespace:
parser = argparse.ArgumentParser(description="Headroom MCP single-runtime adapter")
parser.add_argument(
"--proxy-url",
default=os.environ.get("HEADROOM_PROXY_URL", "http://127.0.0.1:8787"),
)
parser.add_argument("--transport", default="stdio")
parser.add_argument("--host", default="127.0.0.1")
parser.add_argument("--port", type=int, default=8788)
parser.add_argument("--path", default="/mcp")
parser.add_argument("--debug", action="store_true")
parser.add_argument("--direct", action="store_true")
args, unknown = parser.parse_known_args()
if unknown:
print(f"headroom-mcp-proxy: ignoring unsupported arguments: {unknown}", file=sys.stderr)
if args.transport.lower() != "stdio":
parser.error("only stdio transport is supported by the add-on adapter")
return args
def _number(data: dict[str, Any], key: str) -> int:
value = data.get(key, 0)
try:
return int(value)
except (TypeError, ValueError):
return 0
def make_proxy_compressor(proxy_url: str):
endpoint = f"{proxy_url.rstrip('/')}/v1/compress"
model = os.environ.get("HEADROOM_MCP_MODEL", "claude-sonnet-4-5-20250929")
def compress_via_proxy(_self, content: str) -> dict[str, Any]:
# Imported here, not at module import, so MCP startup stays small. httpx is
# already a dependency of Headroom's MCP server for retrieve/stats.
import httpx
response = httpx.post(
endpoint,
json={
"messages": [{"role": "tool", "content": content}],
"model": model,
},
timeout=httpx.Timeout(180.0, connect=75.0),
)
response.raise_for_status()
data = response.json()
if not isinstance(data, dict):
raise RuntimeError("Headroom proxy returned a non-object compression response")
messages = data.get("messages")
compressed: Any = content
if isinstance(messages, list) and messages:
last = messages[-1]
if isinstance(last, dict) and "content" in last:
compressed = last["content"]
if not isinstance(compressed, str):
compressed = json.dumps(compressed, ensure_ascii=False)
hashes = data.get("ccr_hashes")
hash_key = next((item for item in hashes if isinstance(item, str)), None) if isinstance(hashes, list) else None
before = _number(data, "tokens_before")
after = _number(data, "tokens_after")
saved = _number(data, "tokens_saved")
if saved <= 0:
saved = max(0, before - after)
savings_percent = round(saved / before * 100, 1) if before > 0 else 0.0
transforms = data.get("transforms_applied")
if not isinstance(transforms, list):
transforms = []
note = "Compression was executed by the shared Headroom proxy runtime."
if hash_key:
note += (
f" Original stored with hash={hash_key}. "
"Use mcp__headroom__headroom_retrieve to recover it."
)
return {
"compressed": compressed,
"hash": hash_key,
"original_tokens": before,
"compressed_tokens": after,
"tokens_saved": saved,
"savings_percent": savings_percent,
"transforms": transforms,
"note": note,
}
return compress_via_proxy
async def run() -> None:
args = parse_args()
# Defense in depth. The adapter never calls the local compressor, but keep
# the MCP process explicitly unable to initialize Kompress if upstream code
# changes or an unrelated import probes the compression pipeline.
os.environ.pop("HF_HOME", None)
os.environ["HEADROOM_DISABLE_KOMPRESS"] = "1"
os.environ["HEADROOM_PROXY_URL"] = args.proxy_url
from headroom.ccr.mcp_server import HeadroomMCPServer
server = HeadroomMCPServer(proxy_url=args.proxy_url, check_proxy=True)
server._compress_content = MethodType(make_proxy_compressor(args.proxy_url), server)
try:
await server.run_stdio()
finally:
await server.cleanup()
if __name__ == "__main__":
try:
asyncio.run(run())
except KeyboardInterrupt:
pass

View File

@@ -0,0 +1,388 @@
#!/lsiopy/bin/python3
"""Lazy TCP gate for the Headroom HTTP proxy.
The gate remains resident on the public Headroom port while the heavy Headroom
proxy (and its optional ONNX Kompress model) is started only for real traffic.
After an idle period the backend process is terminated, releasing its Python,
ONNX and model allocations. The next connection starts a fresh backend.
Only the Python standard library is imported here deliberately: the idle path
must not import Headroom, ONNX Runtime, transformers or the MCP SDK.
"""
from __future__ import annotations
import asyncio
import contextlib
import json
import os
import signal
import socket
import sys
import time
from pathlib import Path
def _int_env(name: str, default: int, minimum: int = 1) -> int:
try:
value = int(os.environ.get(name, str(default)))
except (TypeError, ValueError):
value = default
return max(minimum, value)
GATE_HOST = os.environ.get("HEADROOM_GATE_HOST", "127.0.0.1")
GATE_PORT = _int_env("HEADROOM_GATE_PORT", 8787)
BACKEND_HOST = os.environ.get("HEADROOM_BACKEND_HOST", "127.0.0.1")
BACKEND_PORT = _int_env("HEADROOM_BACKEND_PORT", 8789)
HEADROOM_BIN = os.environ.get("HEADROOM_REAL_BIN", "/usr/bin/headroom")
IDLE_TIMEOUT = _int_env("HEADROOM_IDLE_TIMEOUT_SECONDS", 900)
START_TIMEOUT = _int_env("HEADROOM_START_TIMEOUT_SECONDS", 60)
STOP_TIMEOUT = _int_env("HEADROOM_STOP_TIMEOUT_SECONDS", 10)
HF_HOME = os.environ.get("HF_HOME", str(Path.home() / ".headroom/hf"))
MAX_HEADER_BYTES = 128 * 1024
COPY_CHUNK = 64 * 1024
def log(message: str) -> None:
stamp = time.strftime("%Y-%m-%d %H:%M:%S")
print(f"[{stamp}] headroom-gate: {message}", file=sys.stderr, flush=True)
class BackendUnavailable(RuntimeError):
pass
class ProxyGate:
def __init__(self) -> None:
self._lock = asyncio.Lock()
self._process: asyncio.subprocess.Process | None = None
self._last_activity = time.monotonic()
self._active_connections = 0
self._stopping = False
self._watch_task: asyncio.Task[None] | None = None
@property
def state(self) -> str:
process = self._process
if process is None:
return "dormant"
if process.returncode is None:
return "running"
return "stopped"
def touch(self) -> None:
self._last_activity = time.monotonic()
async def _backend_healthy(self) -> bool:
for path in ("/livez", "/health"):
try:
reader, writer = await asyncio.wait_for(
asyncio.open_connection(BACKEND_HOST, BACKEND_PORT), timeout=1.0
)
request = (
f"GET {path} HTTP/1.1\r\n"
f"Host: {BACKEND_HOST}:{BACKEND_PORT}\r\n"
"Connection: close\r\n\r\n"
).encode()
writer.write(request)
await writer.drain()
line = await asyncio.wait_for(reader.readline(), timeout=1.0)
writer.close()
with contextlib.suppress(Exception):
await writer.wait_closed()
if line.startswith(b"HTTP/") and b" 2" in line[:16]:
return True
except (OSError, asyncio.TimeoutError):
continue
return False
async def ensure_backend(self) -> None:
if await self._backend_healthy():
self.touch()
return
async with self._lock:
if await self._backend_healthy():
self.touch()
return
process = self._process
if process is not None and process.returncode is None:
await self._wait_until_healthy()
self.touch()
return
if not os.path.isfile(HEADROOM_BIN) or not os.access(HEADROOM_BIN, os.X_OK):
raise BackendUnavailable(f"Headroom executable is unavailable: {HEADROOM_BIN}")
Path(HF_HOME).mkdir(parents=True, exist_ok=True)
environment = os.environ.copy()
environment.update(
{
"HF_HOME": HF_HOME,
"HEADROOM_PROXY_GATE": "1",
}
)
command = [
HEADROOM_BIN,
"proxy",
"--host",
BACKEND_HOST,
"--port",
str(BACKEND_PORT),
"--code-aware",
]
log(f"starting backend on {BACKEND_HOST}:{BACKEND_PORT}")
try:
self._process = await asyncio.create_subprocess_exec(
*command,
env=environment,
start_new_session=True,
)
except Exception as exc:
self._process = None
raise BackendUnavailable(f"unable to start Headroom: {exc}") from exc
self._watch_task = asyncio.create_task(self._watch_backend(self._process))
try:
await self._wait_until_healthy()
except Exception:
await self._terminate_backend_locked("startup failure")
raise
self.touch()
async def _wait_until_healthy(self) -> None:
deadline = time.monotonic() + START_TIMEOUT
while time.monotonic() < deadline:
process = self._process
if process is not None and process.returncode is not None:
raise BackendUnavailable(
f"Headroom exited during startup with status {process.returncode}; "
"see the add-on log"
)
if await self._backend_healthy():
log("backend is ready")
return
await asyncio.sleep(0.25)
raise BackendUnavailable(
f"Headroom did not become ready within {START_TIMEOUT}s; see the add-on log"
)
async def _watch_backend(self, process: asyncio.subprocess.Process) -> None:
returncode = await process.wait()
async with self._lock:
was_current = self._process is process
if was_current:
self._process = None
if not self._stopping and was_current:
log(f"backend exited with status {returncode}")
async def stop_backend(self, reason: str) -> None:
async with self._lock:
await self._terminate_backend_locked(reason)
async def _terminate_backend_locked(self, reason: str) -> None:
process = self._process
if process is None:
return
if process.returncode is not None:
self._process = None
return
log(f"stopping backend ({reason})")
try:
os.killpg(process.pid, signal.SIGTERM)
except ProcessLookupError:
pass
try:
await asyncio.wait_for(process.wait(), timeout=STOP_TIMEOUT)
except asyncio.TimeoutError:
log("backend did not stop after SIGTERM; sending SIGKILL")
with contextlib.suppress(ProcessLookupError):
os.killpg(process.pid, signal.SIGKILL)
with contextlib.suppress(Exception):
await process.wait()
self._process = None
async def idle_monitor(self) -> None:
interval = max(5, min(30, IDLE_TIMEOUT // 4))
while not self._stopping:
await asyncio.sleep(interval)
process = self._process
idle_for = time.monotonic() - self._last_activity
if (
process is not None
and process.returncode is None
and idle_for >= IDLE_TIMEOUT
):
await self.stop_backend(f"idle for {int(idle_for)}s")
def status_payload(self) -> bytes:
process = self._process
payload = {
"status": self.state,
"backend_pid": process.pid if process is not None and process.returncode is None else None,
"active_connections": self._active_connections,
"idle_seconds": round(time.monotonic() - self._last_activity, 1),
"idle_timeout_seconds": IDLE_TIMEOUT,
"backend": f"{BACKEND_HOST}:{BACKEND_PORT}",
}
body = json.dumps(payload, separators=(",", ":")).encode()
return (
b"HTTP/1.1 200 OK\r\n"
b"Content-Type: application/json\r\n"
+ f"Content-Length: {len(body)}\r\n".encode()
+ b"Connection: close\r\n\r\n"
+ body
)
async def handle_client(
self,
client_reader: asyncio.StreamReader,
client_writer: asyncio.StreamWriter,
) -> None:
self._active_connections += 1
real_traffic = False
peer = client_writer.get_extra_info("peername")
try:
try:
initial = await asyncio.wait_for(
client_reader.readuntil(b"\r\n\r\n"), timeout=15.0
)
except (
asyncio.IncompleteReadError,
asyncio.LimitOverrunError,
asyncio.TimeoutError,
):
return
if len(initial) > MAX_HEADER_BYTES:
await self._send_error(client_writer, 431, "request headers too large")
return
first_line = initial.split(b"\r\n", 1)[0]
parts = first_line.split(b" ")
target = parts[1].split(b"?", 1)[0] if len(parts) >= 2 else b""
if target == b"/gate/status":
client_writer.write(self.status_payload())
await client_writer.drain()
return
real_traffic = True
self.touch()
try:
await self.ensure_backend()
backend_reader, backend_writer = await asyncio.wait_for(
asyncio.open_connection(BACKEND_HOST, BACKEND_PORT), timeout=5.0
)
except (BackendUnavailable, OSError, asyncio.TimeoutError) as exc:
log(f"backend unavailable for {peer}: {exc}")
await self._send_error(client_writer, 503, str(exc))
return
backend_writer.write(initial)
await backend_writer.drain()
self.touch()
upstream = asyncio.create_task(self._pipe(client_reader, backend_writer))
downstream = asyncio.create_task(self._pipe(backend_reader, client_writer))
done, pending = await asyncio.wait(
{upstream, downstream}, return_when=asyncio.FIRST_COMPLETED
)
for task in pending:
task.cancel()
for task in done | pending:
with contextlib.suppress(
asyncio.CancelledError,
ConnectionError,
OSError,
):
await task
backend_writer.close()
with contextlib.suppress(Exception):
await backend_writer.wait_closed()
finally:
self._active_connections = max(0, self._active_connections - 1)
if real_traffic:
self.touch()
client_writer.close()
with contextlib.suppress(Exception):
await client_writer.wait_closed()
async def _pipe(
self,
reader: asyncio.StreamReader,
writer: asyncio.StreamWriter,
) -> None:
while True:
data = await reader.read(COPY_CHUNK)
if not data:
with contextlib.suppress(Exception):
writer.write_eof()
return
writer.write(data)
await writer.drain()
self.touch()
@staticmethod
async def _send_error(
writer: asyncio.StreamWriter,
status: int,
detail: str,
) -> None:
reason = (
"Service Unavailable"
if status == 503
else "Request Header Fields Too Large"
)
body = json.dumps({"error": detail}).encode()
response = (
f"HTTP/1.1 {status} {reason}\r\n"
"Content-Type: application/json\r\n"
f"Content-Length: {len(body)}\r\n"
"Connection: close\r\n\r\n"
).encode() + body
writer.write(response)
with contextlib.suppress(Exception):
await writer.drain()
async def shutdown(self) -> None:
self._stopping = True
await self.stop_backend("gate shutdown")
async def async_main() -> None:
gate = ProxyGate()
loop = asyncio.get_running_loop()
stop_event = asyncio.Event()
for sig in (signal.SIGTERM, signal.SIGINT):
with contextlib.suppress(NotImplementedError):
loop.add_signal_handler(sig, stop_event.set)
server = await asyncio.start_server(
gate.handle_client,
GATE_HOST,
GATE_PORT,
limit=MAX_HEADER_BYTES + 1,
family=socket.AF_INET,
)
addresses = ", ".join(str(sock.getsockname()) for sock in server.sockets or [])
log(
f"listening on {addresses}; backend is lazy and stops after {IDLE_TIMEOUT}s idle; "
"status endpoint: /gate/status"
)
monitor = asyncio.create_task(gate.idle_monitor())
async with server:
await stop_event.wait()
monitor.cancel()
with contextlib.suppress(asyncio.CancelledError):
await monitor
await gate.shutdown()
if __name__ == "__main__":
try:
asyncio.run(async_main())
except KeyboardInterrupt:
pass