From 304fe4ee52f7a7f929beb66456ebc979dc867776 Mon Sep 17 00:00:00 2001 From: 7174Andy Date: Wed, 26 Aug 2026 15:46:34 -0700 Subject: [PATCH 1/2] feat: BYO clients ship their own container logs; drop the CLI's detached shipper BYO log shipping was host-side only: run-locally boxes relied on a detached CLI process that could die silently on operator laptops, and hand-off clients (register --no-run-locally, e.g. Run:AI workloads) had no shipper at all -- the container is its own PID 1 with no docker daemon to tail, so their allocator logs pages stayed empty. Now the container ships its own stream. start.sh routes fd 5 (every service's tagged output) through a new ship_logs worker in the client package when SHIP_LOGS=1, which register writes for every BYO shape: * Passthrough first: each line is written to container stdout before anything else touches it, so `docker logs` output is byte-identical; shipping happens on a separate thread fed by a bounded drop-oldest queue and can never block the read loop (the lablink#304 tail-stall lesson: no file intermediary, read the source stream). * Fail open: a supervisor loop respawns a crashed worker and execs plain `cat` after 3 failures -- worst case is exactly the old plumbing, never a frozen container. Missing env degrades to pure passthrough. The CLI's log_shipper.py (441 lines + 673 test lines), its PID/state files, register's spawn/revive machinery, Docker.follow_logs, and the psutil dependency are all deleted. doctor's shipper check becomes a pgrep inside the container. Release ordering: the client image must ship with the ship_logs entry point before a new CLI reaches operators, or BYO clients log locally but nothing reaches the allocator. Old images ignore SHIP_LOGS. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01AzEx3SZp9EAmTVKe6JEWRa --- docs/cli/byo-clients.md | 5 + packages/cli/pyproject.toml | 1 - .../cli/src/lablink_cli/commands/doctor.py | 92 +-- .../cli/src/lablink_cli/commands/register.py | 170 +---- .../src/lablink_cli/commands/reset_overlay.py | 2 +- packages/cli/src/lablink_cli/docker.py | 26 - packages/cli/src/lablink_cli/log_shipper.py | 441 ------------ packages/cli/tests/conftest.py | 8 +- packages/cli/tests/test_docker.py | 62 +- packages/cli/tests/test_doctor.py | 65 +- packages/cli/tests/test_log_shipper.py | 673 ------------------ packages/cli/tests/test_register.py | 386 +--------- packages/client/pyproject.toml | 1 + .../src/lablink_client_service/ship_logs.py | 233 ++++++ packages/client/start.sh | 33 +- packages/client/tests/test_ship_logs.py | 152 ++++ 16 files changed, 533 insertions(+), 1817 deletions(-) delete mode 100644 packages/cli/src/lablink_cli/log_shipper.py delete mode 100644 packages/cli/tests/test_log_shipper.py create mode 100644 packages/client/src/lablink_client_service/ship_logs.py create mode 100644 packages/client/tests/test_ship_logs.py diff --git a/docs/cli/byo-clients.md b/docs/cli/byo-clients.md index 7073c8ad0..e332ca499 100644 --- a/docs/cli/byo-clients.md +++ b/docs/cli/byo-clients.md @@ -179,6 +179,11 @@ easiest thing to send to someone else who owns a box. the workload. In that case you must also pass `--hostname` and `--machine-identity`, since there's nothing local to auto-detect from. + The printed env includes `SHIP_LOGS=1`: the client container ships its own + output stream to the allocator's per-VM logs page (there is no host-side + shipper in any BYO shape). Keep that line in your paste, or the logs page + stays empty for the client. + If docker is missing on the box, `register` keeps the env file so you can install docker and re-run with `--force`. diff --git a/packages/cli/pyproject.toml b/packages/cli/pyproject.toml index 65db4b736..ccd0b79d0 100644 --- a/packages/cli/pyproject.toml +++ b/packages/cli/pyproject.toml @@ -19,7 +19,6 @@ dependencies = [ "rich>=13.0", "pyyaml>=6.0", "boto3>=1.35", - "psutil>=5.9", ] [project.optional-dependencies] diff --git a/packages/cli/src/lablink_cli/commands/doctor.py b/packages/cli/src/lablink_cli/commands/doctor.py index 7653fe23a..f55eae274 100644 --- a/packages/cli/src/lablink_cli/commands/doctor.py +++ b/packages/cli/src/lablink_cli/commands/doctor.py @@ -366,27 +366,6 @@ def _check_manual_prereqs(*, docker: Docker | None = None) -> None: # client on this machine actually working?" # -------------------------------------------------------------------- -# A shipper that is alive but hasn't shipped in this long is reporting a -# problem no liveness check can see — the process is up and the container is -# healthy, but nothing is reaching the allocator. That combination went -# unnoticed for a week (lablink#428), which is the reason this command exists. -SHIPPER_STALE_AFTER_S = 15 * 60 - - -def _format_age(seconds: float) -> str: - """Coarse human age ("6d", "3h", "20m") for the staleness message. - - A raw minute count reads as noise once it passes a few hours - ("8687 min ago"), and this is the one line an operator scans to decide - whether logs are flowing. - """ - if seconds >= 86400: - return f"{int(seconds // 86400)}d" - if seconds >= 3600: - return f"{int(seconds // 3600)}h" - return f"{int(seconds // 60)}m" - - def _check_client_registered() -> dict: """Check that `lablink client register` has run on this box.""" from lablink_cli.commands.register import DEFAULT_ENV_FILE @@ -409,7 +388,7 @@ def _check_client_container(docker: Docker) -> dict: "daemon_error" when the daemon is unreachable, so a separate probe would only duplicate the same `docker inspect` call. """ - from lablink_cli.log_shipper import CONTAINER_NAME + from lablink_cli.commands.register import CONTAINER_NAME result = {"check": "Client container", "status": "fail"} status = docker.container_status(CONTAINER_NAME) @@ -440,59 +419,34 @@ def _check_client_container(docker: Docker) -> dict: return result -def _check_log_shipper(now: float | None = None) -> dict: - """Check the log shipper is alive AND actually shipping. +def _check_log_shipper(docker: Docker) -> dict: + """Check the in-container ship_logs worker is running. - Liveness alone is not enough. A shipper can sit blocked on a quiet - container with a full buffer, process up, nothing delivered — so this - also reports how long ago a batch last landed. + The shipper lives inside the client container — start.sh feeds every + service's output through the client package's ship_logs worker when + SHIP_LOGS=1 — so the probe is a pgrep inside the container. The old + host-side staleness check (lablink#428's alive-but-not-shipping + hazard) moved in-container with it: a worker that can't reach the + allocator says so in `docker logs lablink-client` + ("ship_logs: dropped N lines after retries"). """ - import time - from datetime import datetime, timezone - - from lablink_cli.commands.register import _shipper_alive - from lablink_cli.log_shipper import STATE_FILE, read_last_shipped_ts + from lablink_cli.commands.register import CONTAINER_NAME result = {"check": "Log shipper", "status": "fail"} - if not _shipper_alive(): - result["detail"] = ( - "Not running — client logs are not reaching the allocator. " - "Run `lablink client register` to restart it." - ) - return result - - last = read_last_shipped_ts(STATE_FILE) - if last is None: - result["status"] = "warn" - result["detail"] = ( - "Running, but has never shipped a batch. Normal for the first " - "minute after registering; otherwise check the allocator URL " - "and client secret." - ) - return result - - try: - shipped_at = datetime.strptime(last, "%Y-%m-%dT%H:%M:%SZ").replace( - tzinfo=timezone.utc - ) - except ValueError: - result["status"] = "warn" - result["detail"] = f"Running; unparseable last-shipped value {last!r}" - return result - - current = now if now is not None else time.time() - age_s = current - shipped_at.timestamp() - if age_s > SHIPPER_STALE_AFTER_S: - result["status"] = "warn" - result["detail"] = ( - f"Running, but last shipped {_format_age(age_s)} ago ({last}). " - "The process is up but nothing is reaching the allocator." - ) + probe = docker.exec_in(CONTAINER_NAME, ["pgrep", "-f", "ship_logs"]) + if probe.ok: + result["status"] = "pass" + result["detail"] = "ship_logs worker running inside the container" return result - result["status"] = "pass" - result["detail"] = f"Running; last shipped {last}" + result["detail"] = ( + "No ship_logs worker inside the container — client logs are not " + "reaching the allocator. The container is either down (see the " + "check above), running an image that predates in-container " + "shipping, or missing SHIP_LOGS=1 in its env. Re-run " + "`lablink client register --force` to recreate it." + ) return result @@ -512,7 +466,7 @@ def run_client_doctor(*, docker: Docker | None = None) -> None: checks = [ _check_client_registered(), _check_client_container(docker), - _check_log_shipper(), + _check_log_shipper(docker), ] _render_checks( diff --git a/packages/cli/src/lablink_cli/commands/register.py b/packages/cli/src/lablink_cli/commands/register.py index d6d730ea7..d9cd86c56 100644 --- a/packages/cli/src/lablink_cli/commands/register.py +++ b/packages/cli/src/lablink_cli/commands/register.py @@ -4,15 +4,10 @@ import base64 import json -import os -import subprocess -import sys from datetime import datetime, timezone from pathlib import Path from urllib.parse import urlparse -import psutil - from rich.console import Console from lablink_cli import byo_detect @@ -36,7 +31,9 @@ # same box doesn't have the client-received copy clobber the operator's # local override. DEFAULT_STARTUP_SCRIPT = Path.home() / ".lablink" / "client-custom-startup.sh" -PID_FILE = Path.home() / ".lablink" / "log_shipper.pid" +# One box runs one client container, under a fixed name. Other commands +# (doctor, reset-overlay) import this to inspect the same container. +CONTAINER_NAME = "lablink-client" # tailscaled's node identity (its state dir). Persisted in a named volume so # recreating the container reuses the SAME tailnet node instead of minting a # new one — a new node cannot claim a MagicDNS name that an existing (even @@ -358,23 +355,22 @@ def run_register( "[dim]This container opens a tunnel to the allocator on " "start — confirm with `docker logs lablink-client`.[/dim]" ) - _start_log_shipper(env_file, console) def _resume(env_file: Path, console: Console, docker: Docker) -> None: """Re-run mode for an already-registered host. Does NOT mint a new client_secret. Restarts the container if stopped, - revives the shipper if dead, otherwise prints a no-op message. + otherwise prints a no-op message. The log shipper runs inside the + container (start.sh's ship_logs worker), so it needs no reviving here. Note: container image is NOT re-pulled — that's `--force` territory. """ - status = docker.container_status("lablink-client") - container_action: str | None = None + status = docker.container_status(CONTAINER_NAME) if status == "missing": console.print( - "[yellow]Already registered, but lablink-client container is " + f"[yellow]Already registered, but {CONTAINER_NAME} container is " "missing.[/yellow] Re-run with [bold]--force[/bold] to recreate " "it (this mints a new client_secret)." ) @@ -385,42 +381,25 @@ def _resume(env_file: Path, console: Console, docker: Docker) -> None: ) raise SystemExit(1) if status in ("exited", "restarting"): - start = docker.start_container("lablink-client") + start = docker.start_container(CONTAINER_NAME) if not start.ok: detail = ( start.stderr.strip() or f"docker start exited {start.returncode}" ) console.print( - f"[red]docker start lablink-client failed:[/red] {detail}" + f"[red]docker start {CONTAINER_NAME} failed:[/red] {detail}" ) raise SystemExit(1) - container_action = "restarted" - - shipper_action: str | None = None - if _shipper_alive(): - if container_action is None: - console.print( - "[green]Already registered. Container and log shipper " - "are running.[/green]" - ) - return - else: - _start_log_shipper(env_file, console) - shipper_action = "restarted" - - if container_action and shipper_action: - console.print( - "[green]Restarted container and log shipper.[/green] " - "To pull a newer client image, re-run with --force." - ) - elif container_action: console.print( "[green]Restarted container.[/green] " "To pull a newer client image, re-run with --force." ) - elif shipper_action: - console.print("[green]Restarted log shipper.[/green]") + return + + console.print( + "[green]Already registered. Container is running.[/green]" + ) def _write_env_file( @@ -465,6 +444,15 @@ def _write_env_file( f"CONNECTIVITY={resp['connectivity']}", f"CLIENT_IMAGE={resp['client_image']}", f"REGISTER_RESPONSE={register_response_json}", + # BYO clients ship their own container logs: start.sh routes every + # service's output through the client package's ship_logs worker, + # which passes lines through to docker logs unchanged and forwards + # copies to /api/vm-logs. There is no host-side shipper in any BYO + # shape — hand-off clients (--no-run-locally) have no docker daemon + # anywhere, and the CLI's old detached shipper for run-locally + # boxes was deleted in favor of this. Requires a client image with + # the ship_logs entry point; older images ignore the variable. + "SHIP_LOGS=1", ] # cfg.machine.repository / cfg.machine.software, shipped by the # allocator's register response. These reach an AWS client through @@ -711,7 +699,7 @@ def _exec_docker(cmd: list[str], console: Console, docker: Docker) -> None: raise SystemExit(1) # Remove any existing container with the target name. Quiet on # success; we don't care if it didn't exist. - docker.remove_container("lablink-client", force=True) + docker.remove_container(CONTAINER_NAME, force=True) console.print( f"Starting client container (image: {cmd[-1]}) …" ) @@ -728,112 +716,6 @@ def _exec_docker(cmd: list[str], console: Console, docker: Docker) -> None: ) raise SystemExit(result.returncode) console.print( - "[green]Container running as lablink-client.[/green] " - "View logs with: docker logs -f lablink-client" - ) - - -def _stop_existing_shipper(console: Console) -> None: - """Terminate any running shipper recorded in the PID file. - - Called before spawning a new shipper so ``--force`` re-register doesn't - leave the old shipper briefly tailing the replaced container and - POSTing duplicates against the same hostname. The cmdline guard - matches ``_shipper_alive`` so an unrelated PID-reused process is left - alone. - """ - if not PID_FILE.exists(): - return - try: - pid = int(PID_FILE.read_text().strip()) - except (OSError, ValueError): - PID_FILE.unlink(missing_ok=True) - return - try: - proc = psutil.Process(pid) - cmdline = proc.cmdline() - except (psutil.NoSuchProcess, psutil.AccessDenied): - PID_FILE.unlink(missing_ok=True) - return - if not any("lablink_cli.log_shipper" in arg for arg in cmdline): - PID_FILE.unlink(missing_ok=True) - return - - console.print(f"[dim]Stopping existing log shipper (PID {pid})...[/dim]") - try: - proc.terminate() - proc.wait(timeout=5) - except psutil.TimeoutExpired: - try: - proc.kill() - except psutil.NoSuchProcess: - pass - except (psutil.NoSuchProcess, psutil.AccessDenied): - pass - # The shipper's SIGTERM handler removes the PID file; if we escalated - # to SIGKILL the handler never ran, so clean up here. - PID_FILE.unlink(missing_ok=True) - - -def _start_log_shipper(env_file: Path, console: Console) -> None: - """Spawn the log shipper as a detached background process. - - The shipper survives this `register` invocation and runs until either - the user does ``docker stop lablink-client`` (shipper's docker-logs - subprocess exits and inspect reports missing) or the host reboots. - """ - _stop_existing_shipper(console) - - log_dir = Path.home() / ".lablink" - log_dir.mkdir(parents=True, exist_ok=True) - shipper_log = log_dir / "log_shipper.log" - # Append-mode handle for the detached child's stdout+stderr. The - # shipper itself writes structured lines to this file via self_log(); - # the open handle here is just a safety net for any stray print. - log_fd = open(shipper_log, "a", buffering=1) - - cmd = [sys.executable, "-m", "lablink_cli.log_shipper", str(env_file)] - - popen_kwargs: dict = { - "stdin": subprocess.DEVNULL, - "stdout": log_fd, - "stderr": log_fd, - "close_fds": True, - } - if os.name == "nt": - # Windows: detach so the child survives the parent's exit. - popen_kwargs["creationflags"] = ( - subprocess.DETACHED_PROCESS - | subprocess.CREATE_NEW_PROCESS_GROUP - ) - else: - # POSIX: new session detaches from the controlling TTY and parent - # process group, matching nohup semantics. - popen_kwargs["start_new_session"] = True - - proc = subprocess.Popen(cmd, **popen_kwargs) - console.print( - f"[green]Log shipping started (PID {proc.pid}).[/green] " - f"Logs: {shipper_log}" + f"[green]Container running as {CONTAINER_NAME}.[/green] " + f"View logs with: docker logs -f {CONTAINER_NAME}" ) - - -def _shipper_alive() -> bool: - """True iff a live log-shipper process matching our PID file exists. - - Two-stage check: PID present in PID file AND that PID belongs to a - process whose cmdline mentions ``lablink_cli.log_shipper``. The - cmdline guard prevents false positives from PID reuse after reboot. - """ - if not PID_FILE.exists(): - return False - try: - pid = int(PID_FILE.read_text().strip()) - except (OSError, ValueError): - return False - try: - proc = psutil.Process(pid) - cmdline = proc.cmdline() - except (psutil.NoSuchProcess, psutil.AccessDenied): - return False - return any("lablink_cli.log_shipper" in arg for arg in cmdline) diff --git a/packages/cli/src/lablink_cli/commands/reset_overlay.py b/packages/cli/src/lablink_cli/commands/reset_overlay.py index 0c0e3c595..77def8544 100644 --- a/packages/cli/src/lablink_cli/commands/reset_overlay.py +++ b/packages/cli/src/lablink_cli/commands/reset_overlay.py @@ -26,7 +26,7 @@ from lablink_cli.commands.register import TAILSCALE_STATE_VOLUME from lablink_cli.docker import Docker, DockerUnavailable, default_docker -from lablink_cli.log_shipper import CONTAINER_NAME +from lablink_cli.commands.register import CONTAINER_NAME TAILNET_ADMIN_URL = "https://login.tailscale.com/admin/machines" diff --git a/packages/cli/src/lablink_cli/docker.py b/packages/cli/src/lablink_cli/docker.py index ccd37b3d7..339df17fa 100644 --- a/packages/cli/src/lablink_cli/docker.py +++ b/packages/cli/src/lablink_cli/docker.py @@ -244,27 +244,6 @@ def logs( stderr=result.stderr or "", ) - def follow_logs( - self, name: str, *, since: str | None = None - ) -> subprocess.Popen: - """Spawn ``docker logs --follow --timestamps [--since ] ``. - - Returns the Popen handle: the log shipper needs ``.terminate()`` and - ``.poll()`` as well as incremental reads from ``.stdout``. - """ - self.require() - argv = ["docker", "logs", "--follow", "--timestamps"] - if since: - argv += ["--since", since] - argv.append(name) - return subprocess.Popen( - argv, - stdout=subprocess.PIPE, - stderr=subprocess.STDOUT, - text=True, - bufsize=1, # line-buffered - ) - # -- escape hatches ------------------------------------------------ def compose( @@ -389,11 +368,6 @@ def logs( ) -> Result: return Result(returncode=1, stderr=self._NOT_FOUND) - def follow_logs( - self, name: str, *, since: str | None = None - ) -> subprocess.Popen: - raise DockerUnavailable(self._NOT_FOUND) - def _run( self, argv: list[str], diff --git a/packages/cli/src/lablink_cli/log_shipper.py b/packages/cli/src/lablink_cli/log_shipper.py deleted file mode 100644 index 8126aba53..000000000 --- a/packages/cli/src/lablink_cli/log_shipper.py +++ /dev/null @@ -1,441 +0,0 @@ -"""Ships docker container logs from a BYO client to the allocator. - -Invoked as: ``python -m lablink_cli.log_shipper `` -by ``lablink register``. Reads CLIENT_SECRET / ALLOCATOR_URL / VM_NAME from -the env file written by register, batches ``docker logs --follow`` output, -and POSTs to ``/api/vm-logs/``. -""" - -from __future__ import annotations - -import json -import os -import queue -import re -import signal -import sys -import threading -import time -from datetime import datetime, timedelta, timezone -from pathlib import Path -from typing import Callable, Literal -from urllib.error import HTTPError, URLError -from urllib.request import Request -from urllib.request import urlopen as _stdlib_urlopen - -from lablink_cli.api import USER_AGENT -from lablink_cli.docker import Docker, default_docker - - -def load_env(env_file: Path) -> dict[str, str]: - """Parse the BYO client.env file (KEY=VALUE per line, # comments).""" - text = Path(env_file).read_text() - env: dict[str, str] = {} - for raw in text.splitlines(): - line = raw.strip() - if not line or line.startswith("#"): - continue - if "=" not in line: - continue - key, _, value = line.partition("=") - env[key.strip()] = value - return env - - -def read_last_shipped_ts(state_file: Path) -> str | None: - """Return the timestamp of the last successfully shipped line, or None. - - Treats missing or corrupt state as None (first-attach behavior). - """ - try: - data = json.loads(Path(state_file).read_text()) - except (FileNotFoundError, json.JSONDecodeError): - return None - ts = data.get("last_shipped_ts") - return ts if isinstance(ts, str) else None - - -def write_last_shipped_ts(state_file: Path, ts: str) -> None: - """Persist last_shipped_ts atomically (write-and-rename).""" - path = Path(state_file) - path.parent.mkdir(parents=True, exist_ok=True) - tmp = path.with_suffix(path.suffix + ".tmp") - tmp.write_text(json.dumps({"last_shipped_ts": ts})) - tmp.replace(path) - - -MAX_RETRIES = 3 -RETRY_BACKOFF_S = (1, 2, 4) # sleep before retry attempt 1, 2, 3 -# The allocator routes logs to the docker_logs column when log_group ends -# with "-docker"; cloud_init otherwise (main.py:851). Manual/BYO clients -# only ship docker container output, so use a name that satisfies the -# suffix check. -LOG_GROUP = "manual-docker" - -PostResult = Literal["ok", "drop", "fatal"] - - -def post_batch( - *, - allocator_url: str, - vm_name: str, - client_secret: str, - messages: list[str], - log_group: str = LOG_GROUP, - urlopen: Callable = _stdlib_urlopen, - sleep: Callable[[float], None] = time.sleep, -) -> PostResult: - """POST a batch of log lines to /api/vm-logs/. - - Returns ``"ok"`` on 2xx, ``"fatal"`` on 4xx (no retry — shipper should - exit), and ``"drop"`` after MAX_RETRIES of 5xx or network failures. - """ - url = f"{allocator_url.rstrip('/')}/api/vm-logs/{vm_name}" - body = json.dumps({"log_group": log_group, "messages": messages}).encode() - headers = { - "Content-Type": "application/json", - "Authorization": f"Bearer {client_secret}", - # urllib's default "Python-urllib/x.y" is blocked with HTTP 403 by - # Cloudflare-proxied allocators (see api.py's USER_AGENT) — post_batch - # treats any 4xx as fatal, so without this every batch kills the - # shipper on its first POST. - "User-Agent": USER_AGENT, - } - - for attempt in range(MAX_RETRIES): - if attempt > 0: - sleep(RETRY_BACKOFF_S[attempt - 1]) - try: - req = Request(url, data=body, headers=headers, method="POST") - with urlopen(req, timeout=10) as resp: - if 200 <= resp.status < 300: - return "ok" - continue - except HTTPError as e: - if 400 <= e.code < 500: - return "fatal" # bad secret / unknown hostname — exit - continue - except URLError: - continue - return "drop" - - -BATCH_SIZE = 50 -FLUSH_INTERVAL_S = 15 - - -def should_flush(*, buffer_len: int, elapsed_s: float) -> bool: - """Return True if the buffer should be flushed now.""" - if buffer_len == 0: - return False - return buffer_len >= BATCH_SIZE or elapsed_s >= FLUSH_INTERVAL_S - - -CONTAINER_NAME = "lablink-client" - -# Strip RFC3339Nano fractional seconds (e.g. -# 2026-05-28T14:23:01.123456789Z → 2026-05-28T14:23:01Z) so admin views -# aren't cluttered with nanosecond noise. Matches log_shipper.sh:101. -_TS_RE = re.compile( - r"^(\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2})(?:\.\d+)?Z (.*)$" -) - - -def parse_docker_line(line: str) -> tuple[str | None, str]: - """Split a ``docker logs --timestamps`` line into ``(ts, message)``. - - Returns ``(None, line)`` if no timestamp prefix is present. - """ - m = _TS_RE.match(line) - if not m: - return None, line - return f"{m.group(1)}Z", m.group(2) - - -LOG_SHIPPER_DIR = Path.home() / ".lablink" -PID_FILE = LOG_SHIPPER_DIR / "log_shipper.pid" -STATE_FILE = LOG_SHIPPER_DIR / "log_shipper.state" -SELF_LOG_FILE = LOG_SHIPPER_DIR / "log_shipper.log" -FIRST_ATTACH_LOOKBACK_S = 60 -CONTAINER_RESTART_WAIT_S = 5 -INSPECT_RETRY_INTERVAL_S = 30 -INSPECT_MAX_RETRIES = 5 -# After this many consecutive "exited" inspections, give up — the container -# has stopped and docker is not restarting it, which means the user invoked -# `docker stop`. With `--restart unless-stopped`, a crashed container goes to -# "restarting" within ms, so consecutive "exited" reliably indicates user -# intent rather than a transient state. -MAX_EXITED_CONSECUTIVE = 3 - -SELF_LOG_MAX_BYTES = 1_000_000 - - -def self_log(log_file: Path, message: str) -> None: - """Append a timestamped line to the shipper's own diagnostic log. - - Rotates to .1 (single rotation) when the file exceeds 1MB. This - is the shipper's only error channel; it runs detached so stdout/stderr - are discarded. - """ - path = Path(log_file) - path.parent.mkdir(parents=True, exist_ok=True) - if path.exists() and path.stat().st_size >= SELF_LOG_MAX_BYTES: - rotated = path.with_suffix(path.suffix + ".1") - path.replace(rotated) - ts = datetime.now(timezone.utc).isoformat(timespec="seconds") - with path.open("a") as f: - f.write(f"{ts} {message}\n") - - -def _initial_since() -> str: - """RFC3339 timestamp ~1 minute ago, for first-ever shipper attach.""" - ts = datetime.now(timezone.utc) - timedelta( - seconds=FIRST_ATTACH_LOOKBACK_S - ) - return ts.strftime("%Y-%m-%dT%H:%M:%SZ") - - -# Yielded when the flush window elapses with no new log line, so the read -# loop wakes up and re-evaluates should_flush(). -TICK = object() - - -def _read_lines_from_popen(proc, *, timeout: float = FLUSH_INTERVAL_S): - """Yield a Popen's stdout line by line, plus TICK on each idle timeout. - - A bare ``for line in proc.stdout`` only wakes when output arrives, so a - container that logs a burst at startup and then goes quiet holds its - buffer forever and the allocator never sees a single line. The shell - shipper avoids this with ``read -t "$FLUSH_INTERVAL"`` - (log_shipper.sh:115); pipes aren't selectable on Windows — where BYO - clients actually run — so a reader thread plus a queue is the portable - equivalent. - """ - if proc.stdout is None: - return - q: queue.Queue = queue.Queue() - eof = object() - - def pump(): - try: - for line in proc.stdout: - q.put(line.rstrip("\n")) - finally: - q.put(eof) - - threading.Thread(target=pump, daemon=True).start() - while True: - try: - item = q.get(timeout=timeout) - except queue.Empty: - yield TICK - continue - if item is eof: - return - yield item - - -def run_shipper( - env_file: Path, - *, - _line_iter: Callable | None = None, - _sleep: Callable[[float], None] = time.sleep, - docker: Docker | None = None, -) -> None: - """Main shipper loop. Returns when shipping should stop.""" - docker = docker or default_docker() - env = load_env(env_file) - allocator_url = env["ALLOCATOR_URL"] - vm_name = env["VM_NAME"] - client_secret = env["CLIENT_SECRET"] - - self_log(SELF_LOG_FILE, f"shipper starting for vm_name={vm_name}") - - since = read_last_shipped_ts(STATE_FILE) or _initial_since() - self_log(SELF_LOG_FILE, f"attaching to docker logs --since {since}") - - # ---- Attach loop: re-runs if container restarts ---- - inspect_failures = 0 - exited_consecutive = 0 - while True: - if _line_iter is not None: - line_source = _line_iter() - proc = None - else: - proc = docker.follow_logs(CONTAINER_NAME, since=since) - line_source = _read_lines_from_popen(proc) - - buffer: list[str] = [] - buffer_first_ts: float | None = None - last_ts_in_batch: str | None = None - - # ---- Inner read loop ---- - try: - for line in line_source: - # TICK means the flush window elapsed with no new line — - # skip straight to the flush check below. should_flush() - # already no-ops on an empty buffer. - if line is not TICK: - ts, msg = parse_docker_line(line) - # Buffer the original tagged line, preserving the - # timestamp prefix so admin views show it (matches - # log_shipper.sh's docker --timestamps + sed pipeline). - if ts is not None: - buffer.append(f"{ts} {msg}") - last_ts_in_batch = ts - else: - buffer.append(msg) - if buffer_first_ts is None: - buffer_first_ts = time.monotonic() - - elapsed = ( - time.monotonic() - buffer_first_ts - if buffer_first_ts is not None - else 0 - ) - if should_flush( - buffer_len=len(buffer), elapsed_s=elapsed - ): - result = post_batch( - allocator_url=allocator_url, - vm_name=vm_name, - client_secret=client_secret, - messages=buffer, - ) - if result == "ok" and last_ts_in_batch: - write_last_shipped_ts( - STATE_FILE, last_ts_in_batch - ) - elif result == "fatal": - self_log( - SELF_LOG_FILE, - "POST returned fatal (4xx); exiting", - ) - return - elif result == "drop": - self_log( - SELF_LOG_FILE, - f"dropped batch of {len(buffer)} after retries", - ) - buffer = [] - buffer_first_ts = None - last_ts_in_batch = None - finally: - if proc is not None: - try: - proc.terminate() - proc.wait(timeout=5) - except Exception: - pass - - # Flush any tail buffer before deciding whether to reconnect. - if buffer: - result = post_batch( - allocator_url=allocator_url, - vm_name=vm_name, - client_secret=client_secret, - messages=buffer, - ) - if result == "ok" and last_ts_in_batch: - write_last_shipped_ts(STATE_FILE, last_ts_in_batch) - elif result == "fatal": - self_log( - SELF_LOG_FILE, - "POST returned fatal during tail flush; exiting", - ) - return - - # docker logs --follow exited. Inspect to decide what to do. - status = docker.container_status(CONTAINER_NAME) - self_log( - SELF_LOG_FILE, f"docker logs ended; container status={status}" - ) - if status == "missing": - self_log(SELF_LOG_FILE, "container missing; exiting") - return - if status == "daemon_error": - inspect_failures += 1 - if inspect_failures >= INSPECT_MAX_RETRIES: - self_log( - SELF_LOG_FILE, - "daemon unreachable after max retries; exiting", - ) - return - _sleep(INSPECT_RETRY_INTERVAL_S) - continue - inspect_failures = 0 - if status == "exited": - exited_consecutive += 1 - if exited_consecutive >= MAX_EXITED_CONSECUTIVE: - self_log( - SELF_LOG_FILE, - f"container stayed exited for {exited_consecutive} " - "consecutive checks; treating as user-initiated stop; " - "exiting", - ) - return - _sleep(CONTAINER_RESTART_WAIT_S) - elif status == "restarting": - # docker is bringing it back — don't count toward the give-up - # threshold, just wait and reconnect. - exited_consecutive = 0 - _sleep(CONTAINER_RESTART_WAIT_S) - else: - # status == "running" → reconnect immediately - exited_consecutive = 0 - # update since for the reconnect so we don't re-ship - new_since = read_last_shipped_ts(STATE_FILE) - if new_since: - since = new_since - # If _line_iter is set (test), exit the outer loop after one pass to - # keep tests deterministic. - if _line_iter is not None: - return - - -def _handle_shutdown(signum, _frame) -> None: - """SIGTERM/SIGINT handler: unlink PID file and exit cleanly.""" - try: - PID_FILE.unlink(missing_ok=True) - except OSError: - pass - self_log(SELF_LOG_FILE, f"received signal {signum}; exiting") - sys.exit(0) - - -def main(argv: list[str] | None = None) -> int: - """Entry point for ``python -m lablink_cli.log_shipper ``.""" - args = argv if argv is not None else sys.argv[1:] - if len(args) != 1: - print( - "usage: python -m lablink_cli.log_shipper ", - file=sys.stderr, - ) - return 2 - - env_file = Path(args[0]) - - LOG_SHIPPER_DIR.mkdir(parents=True, exist_ok=True) - PID_FILE.write_text(str(os.getpid())) - - # Best-effort signal handlers. Windows lacks SIGTERM in the standard - # sense; signal.signal(SIGTERM, ...) works on POSIX but is a no-op or - # raises on Windows for some signals — guard with try/except. - for sig in (signal.SIGTERM, signal.SIGINT): - try: - signal.signal(sig, _handle_shutdown) - except (ValueError, AttributeError): - pass - - try: - run_shipper(env_file) - finally: - try: - PID_FILE.unlink(missing_ok=True) - except OSError: - pass - return 0 - - -if __name__ == "__main__": - raise SystemExit(main()) diff --git a/packages/cli/tests/conftest.py b/packages/cli/tests/conftest.py index 1fcbb94c3..90a139eae 100644 --- a/packages/cli/tests/conftest.py +++ b/packages/cli/tests/conftest.py @@ -29,11 +29,9 @@ def _no_real_docker(request, monkeypatch): a real daemon. This only keeps tests off the docker *daemon*, not off the host in - general: `register._start_log_shipper` spawns - `python -m lablink_cli.log_shipper` via a bare ``subprocess.Popen`` - whose ``argv[0]`` is ``sys.executable``, not ``"docker"`` — this guard - does not (and should not) catch that. A test exercising that path that - forgets to mock `_start_log_shipper` spawns a real detached process. + general: a bare ``subprocess.Popen`` whose ``argv[0]`` isn't ``docker`` + (e.g. spawning a Python module) is not caught — and should not be, this + guard is about missed adapter call sites, not process spawning. """ if request.node.get_closest_marker("integration"): return diff --git a/packages/cli/tests/test_docker.py b/packages/cli/tests/test_docker.py index 56d109a94..927efa6e6 100644 --- a/packages/cli/tests/test_docker.py +++ b/packages/cli/tests/test_docker.py @@ -3,7 +3,7 @@ from __future__ import annotations import subprocess -from unittest.mock import MagicMock, patch +from unittest.mock import patch import pytest @@ -136,59 +136,6 @@ def test_logs_swallows_missing_binary_as_a_failed_result(): assert not result.ok -def test_follow_logs_builds_base_argv(): - # Ported from log_shipper.TestOpenDockerLogs.test_builds_command_with_since - # (Task 5) — the container-name-only shape of that same argv. - with patch("lablink_cli.docker.subprocess.Popen") as mock_popen: - mock_popen.return_value = MagicMock() - Docker().follow_logs("lablink-client") - argv = mock_popen.call_args.args[0] - assert argv == [ - "docker", - "logs", - "--follow", - "--timestamps", - "lablink-client", - ] - - -def test_follow_logs_includes_since_before_name(): - # Ported from - # log_shipper.TestOpenDockerLogs.test_builds_command_with_since. - with patch("lablink_cli.docker.subprocess.Popen") as mock_popen: - mock_popen.return_value = MagicMock() - Docker().follow_logs( - "lablink-client", since="2026-08-12T00:00:00Z" - ) - argv = mock_popen.call_args.args[0] - since_idx = argv.index("--since") - assert argv[since_idx + 1] == "2026-08-12T00:00:00Z" - assert argv.index("lablink-client") > since_idx + 1 - - -def test_follow_logs_omits_since_when_none(): - # Ported from log_shipper.TestOpenDockerLogs.test_omits_since_when_none. - with patch("lablink_cli.docker.subprocess.Popen") as mock_popen: - mock_popen.return_value = MagicMock() - Docker().follow_logs("lablink-client", since=None) - argv = mock_popen.call_args.args[0] - assert "--since" not in argv - - -def test_follow_logs_popen_kwargs(): - """Line-buffering and merged streams are load-bearing for the shipper's - incremental reads (_read_lines_from_popen) — not covered by the argv - tests above, so assert on them directly.""" - with patch("lablink_cli.docker.subprocess.Popen") as mock_popen: - mock_popen.return_value = MagicMock() - Docker().follow_logs("lablink-client") - kwargs = mock_popen.call_args.kwargs - assert kwargs["stdout"] is subprocess.PIPE - assert kwargs["stderr"] is subprocess.STDOUT - assert kwargs["text"] is True - assert kwargs["bufsize"] == 1 - - def test_compose_passes_workdir_as_cwd(): with patch("lablink_cli.docker.subprocess.run") as run: run.return_value = _completed(0) @@ -242,13 +189,6 @@ def test_null_docker_logs_does_not_shell_out(): assert "No such container" in result.stderr -def test_log_shipper_no_longer_defines_inspect_container(): - """The verb lives in the adapter now; log_shipper must not re-export it.""" - import lablink_cli.log_shipper as ls - - assert not hasattr(ls, "inspect_container") - - def test_container_status_running_moved_from_log_shipper(): # Ported from log_shipper.TestInspectContainer.test_running (Task 2) — # same scenario as test_container_status_maps_running above, kept as a diff --git a/packages/cli/tests/test_doctor.py b/packages/cli/tests/test_doctor.py index 81ddcb950..424c82bf7 100644 --- a/packages/cli/tests/test_doctor.py +++ b/packages/cli/tests/test_doctor.py @@ -304,61 +304,30 @@ def test_daemon_error_reported_as_daemon_problem(self): class TestCheckLogShipper: - # 2026-08-05T12:00:00Z - SHIPPED = "2026-08-05T12:00:00Z" - SHIPPED_EPOCH = 1785931200.0 + """The shipper runs inside the container (start.sh's ship_logs + worker), so the check is a pgrep via `docker exec`.""" - def _run(self, *, alive, last, now=None): + def _run(self, *, ok): from lablink_cli.commands.doctor import _check_log_shipper + from lablink_cli.docker import Result - with ( - patch( - "lablink_cli.commands.register._shipper_alive", - return_value=alive, - ), - patch( - "lablink_cli.log_shipper.read_last_shipped_ts", - return_value=last, - ), - ): - return _check_log_shipper(now=now) - - def test_dead_shipper_fails(self): - result = self._run(alive=False, last=self.SHIPPED) - assert result["status"] == "fail" - assert "not reaching the allocator" in result["detail"].lower() + mock_docker = MagicMock() + mock_docker.exec_in.return_value = Result(0 if ok else 1) + result = _check_log_shipper(mock_docker) + return result, mock_docker - def test_alive_and_recent_passes(self): - result = self._run( - alive=True, last=self.SHIPPED, now=self.SHIPPED_EPOCH + 60 - ) + def test_worker_present_passes(self): + result, docker = self._run(ok=True) assert result["status"] == "pass" - - def test_alive_but_stale_warns(self): - """The failure a liveness-only check cannot see: process up, - container healthy, nothing reaching the allocator.""" - result = self._run( - alive=True, last=self.SHIPPED, now=self.SHIPPED_EPOCH + 6 * 86400 + docker.exec_in.assert_called_once_with( + "lablink-client", ["pgrep", "-f", "ship_logs"] ) - assert result["status"] == "warn" - assert "6d ago" in result["detail"] - assert self.SHIPPED in result["detail"] - - def test_age_formatting_stays_readable(self): - from lablink_cli.commands.doctor import _format_age - assert _format_age(20 * 60) == "20m" - assert _format_age(3 * 3600) == "3h" - assert _format_age(6 * 86400) == "6d" - - def test_alive_but_never_shipped_warns(self): - result = self._run(alive=True, last=None) - assert result["status"] == "warn" - assert "never shipped" in result["detail"] - - def test_unparseable_timestamp_warns(self): - result = self._run(alive=True, last="not-a-timestamp") - assert result["status"] == "warn" + def test_worker_absent_fails_with_remedy(self): + result, _ = self._run(ok=False) + assert result["status"] == "fail" + assert "not reaching the allocator" in result["detail"] + assert "--force" in result["detail"] class TestRunClientDoctor: diff --git a/packages/cli/tests/test_log_shipper.py b/packages/cli/tests/test_log_shipper.py deleted file mode 100644 index e4f68105b..000000000 --- a/packages/cli/tests/test_log_shipper.py +++ /dev/null @@ -1,673 +0,0 @@ -"""Tests for lablink_cli.log_shipper.""" - -from __future__ import annotations - -import json - -import pytest - - -class TestLoadEnv: - def test_parses_key_value_lines(self, tmp_path): - from lablink_cli.log_shipper import load_env - - env_file = tmp_path / "client.env" - env_file.write_text( - "# Comment line\n" - "CLIENT_ID=42\n" - "VM_NAME=42\n" - "CLIENT_SECRET=s3cr3t\n" - "ALLOCATOR_URL=https://lablink.example.com\n" - "\n" - ) - - env = load_env(env_file) - - assert env["CLIENT_ID"] == "42" - assert env["VM_NAME"] == "42" - assert env["CLIENT_SECRET"] == "s3cr3t" - assert env["ALLOCATOR_URL"] == "https://lablink.example.com" - assert "# Comment line" not in env - - def test_missing_file_raises(self, tmp_path): - from lablink_cli.log_shipper import load_env - - with pytest.raises(FileNotFoundError): - load_env(tmp_path / "nope.env") - - def test_value_with_equals_sign_preserved(self, tmp_path): - from lablink_cli.log_shipper import load_env - - env_file = tmp_path / "client.env" - env_file.write_text("WEIRD=a=b=c\n") - - assert load_env(env_file)["WEIRD"] == "a=b=c" - - -class TestStateFile: - def test_read_missing_returns_none(self, tmp_path): - from lablink_cli.log_shipper import read_last_shipped_ts - - assert read_last_shipped_ts(tmp_path / "missing.state") is None - - def test_round_trip(self, tmp_path): - from lablink_cli.log_shipper import ( - read_last_shipped_ts, - write_last_shipped_ts, - ) - - state = tmp_path / "log_shipper.state" - write_last_shipped_ts(state, "2026-05-28T14:23:01Z") - - assert read_last_shipped_ts(state) == "2026-05-28T14:23:01Z" - - def test_corrupt_state_returns_none(self, tmp_path): - from lablink_cli.log_shipper import read_last_shipped_ts - - state = tmp_path / "bad.state" - state.write_text("{not json") - - assert read_last_shipped_ts(state) is None - - def test_state_dir_created(self, tmp_path): - from lablink_cli.log_shipper import write_last_shipped_ts - - nested = tmp_path / "subdir" / "log_shipper.state" - write_last_shipped_ts(nested, "2026-05-28T14:23:01Z") - - assert nested.exists() - - -class TestPostBatch: - def _args(self, **overrides): - base = dict( - allocator_url="https://lablink.example.com", - vm_name="42", - client_secret="s3cr3t", - messages=["[start] booting", "[agent] ready"], - ) - base.update(overrides) - return base - - def test_success_returns_ok(self): - from unittest.mock import MagicMock - from lablink_cli.log_shipper import post_batch - - resp = MagicMock() - resp.__enter__.return_value = resp - resp.__exit__.return_value = False - resp.status = 200 - resp.read.return_value = b'{"ok": true}' - urlopen = MagicMock(return_value=resp) - sleep = MagicMock() - - result = post_batch(**self._args(), urlopen=urlopen, sleep=sleep) - - assert result == "ok" - assert urlopen.call_count == 1 - sleep.assert_not_called() - # verify request shape - req = urlopen.call_args.args[0] - assert req.full_url == "https://lablink.example.com/api/vm-logs/42" - assert req.get_method() == "POST" - assert req.get_header("Authorization") == "Bearer s3cr3t" - assert req.get_header("Content-type") == "application/json" - body = json.loads(req.data.decode()) - assert body == { - "log_group": "manual-docker", - "messages": ["[start] booting", "[agent] ready"], - } - - def test_sends_product_user_agent(self): - """Cloudflare-proxied allocators 403 urllib's default UA. - - Regression: post_batch built its Request without a User-Agent, - so every batch was blocked before the client secret was even - checked, and post_batch's own 4xx handling treated that as - fatal and killed the shipper on its first POST. - """ - from unittest.mock import MagicMock - from lablink_cli.api import USER_AGENT - from lablink_cli.log_shipper import post_batch - - resp = MagicMock() - resp.__enter__.return_value = resp - resp.__exit__.return_value = False - resp.status = 200 - resp.read.return_value = b'{"ok": true}' - urlopen = MagicMock(return_value=resp) - - post_batch(**self._args(), urlopen=urlopen, sleep=MagicMock()) - - req = urlopen.call_args.args[0] - assert req.get_header("User-agent") == USER_AGENT - - def test_retries_on_5xx_then_succeeds(self): - from io import BytesIO - from unittest.mock import MagicMock - from urllib.error import HTTPError - from lablink_cli.log_shipper import post_batch - - ok_resp = MagicMock() - ok_resp.__enter__.return_value = ok_resp - ok_resp.__exit__.return_value = False - ok_resp.status = 200 - - urlopen = MagicMock( - side_effect=[ - HTTPError("u", 503, "boom", {}, BytesIO(b"")), - HTTPError("u", 503, "boom", {}, BytesIO(b"")), - ok_resp, - ] - ) - sleep = MagicMock() - - result = post_batch(**self._args(), urlopen=urlopen, sleep=sleep) - - assert result == "ok" - assert urlopen.call_count == 3 - # backoffs: 1s before retry 1, 2s before retry 2 - sleep.assert_any_call(1) - sleep.assert_any_call(2) - - def test_drops_on_3_consecutive_5xx(self): - from io import BytesIO - from unittest.mock import MagicMock - from urllib.error import HTTPError - from lablink_cli.log_shipper import post_batch - - urlopen = MagicMock( - side_effect=HTTPError("u", 503, "boom", {}, BytesIO(b"")) - ) - sleep = MagicMock() - - result = post_batch(**self._args(), urlopen=urlopen, sleep=sleep) - - assert result == "drop" - assert urlopen.call_count == 3 # initial + 2 retries - - def test_drops_on_network_errors(self): - from unittest.mock import MagicMock - from urllib.error import URLError - from lablink_cli.log_shipper import post_batch - - urlopen = MagicMock(side_effect=URLError("connection refused")) - sleep = MagicMock() - - result = post_batch(**self._args(), urlopen=urlopen, sleep=sleep) - - assert result == "drop" - assert urlopen.call_count == 3 - - def test_fatal_on_401(self): - from io import BytesIO - from unittest.mock import MagicMock - from urllib.error import HTTPError - from lablink_cli.log_shipper import post_batch - - urlopen = MagicMock( - side_effect=HTTPError("u", 401, "unauthorized", {}, BytesIO(b"")) - ) - sleep = MagicMock() - - result = post_batch(**self._args(), urlopen=urlopen, sleep=sleep) - - assert result == "fatal" - assert urlopen.call_count == 1 # no retry on 4xx - sleep.assert_not_called() - - def test_fatal_on_404(self): - from io import BytesIO - from unittest.mock import MagicMock - from urllib.error import HTTPError - from lablink_cli.log_shipper import post_batch - - urlopen = MagicMock( - side_effect=HTTPError("u", 404, "not found", {}, BytesIO(b"")) - ) - sleep = MagicMock() - - result = post_batch(**self._args(), urlopen=urlopen, sleep=sleep) - - assert result == "fatal" - assert urlopen.call_count == 1 - - -class TestShouldFlush: - def test_empty_buffer_never_flushes(self): - from lablink_cli.log_shipper import should_flush - - assert should_flush(buffer_len=0, elapsed_s=999) is False - - def test_flushes_at_batch_size(self): - from lablink_cli.log_shipper import should_flush - - assert should_flush(buffer_len=50, elapsed_s=0) is True - assert should_flush(buffer_len=49, elapsed_s=0) is False - - def test_flushes_at_time_threshold(self): - from lablink_cli.log_shipper import should_flush - - assert should_flush(buffer_len=1, elapsed_s=15) is True - assert should_flush(buffer_len=1, elapsed_s=14.9) is False - - -class TestParseDockerLine: - def test_strips_nanoseconds(self): - from lablink_cli.log_shipper import parse_docker_line - - ts, msg = parse_docker_line( - "2026-05-28T14:23:01.123456789Z [agent] hello world" - ) - assert ts == "2026-05-28T14:23:01Z" - assert msg == "[agent] hello world" - - def test_whole_seconds_passthrough(self): - from lablink_cli.log_shipper import parse_docker_line - - ts, msg = parse_docker_line("2026-05-28T14:23:01Z [start] boot") - assert ts == "2026-05-28T14:23:01Z" - assert msg == "[start] boot" - - def test_no_timestamp_returns_none_ts(self): - from lablink_cli.log_shipper import parse_docker_line - - ts, msg = parse_docker_line("no timestamp here") - assert ts is None - assert msg == "no timestamp here" - - -class TestReadLinesFromPopen: - def test_yields_tick_while_source_is_idle(self): - """A source that emits a burst then goes quiet must still tick. - - Without the idle tick the read loop blocks forever and a buffered - startup burst is never flushed to the allocator. - """ - import subprocess - import sys - from lablink_cli.log_shipper import TICK, _read_lines_from_popen - - proc = subprocess.Popen( - [ - sys.executable, - "-c", - "import time; print('burst', flush=True); time.sleep(30)", - ], - stdout=subprocess.PIPE, - stderr=subprocess.STDOUT, - text=True, - bufsize=1, - ) - try: - seen = [] - for item in _read_lines_from_popen(proc, timeout=0.1): - seen.append(item) - # Interpreter startup may outrun the first tick, so wait for - # both the burst and the ticks rather than assuming an order. - if "burst" in seen and seen.count(TICK) >= 2: - break - finally: - proc.terminate() - proc.wait(timeout=5) - - assert "burst" in seen - assert seen.count(TICK) >= 2 - - -class TestSelfLog: - def test_appends_line(self, tmp_path): - from lablink_cli.log_shipper import self_log - - log = tmp_path / "log_shipper.log" - self_log(log, "first") - self_log(log, "second") - - content = log.read_text() - assert "first" in content - assert "second" in content - # one line per call - assert content.count("\n") == 2 - - def test_rotates_at_size_cap(self, tmp_path): - from lablink_cli.log_shipper import self_log - - log = tmp_path / "log_shipper.log" - # Pre-fill above cap - log.write_text("x" * 1_100_000) - - self_log(log, "new line") - - rotated = tmp_path / "log_shipper.log.1" - assert rotated.exists() - assert log.read_text().rstrip().endswith("new line") - # the rotated file holds the old content - assert len(rotated.read_text()) >= 1_000_000 - - -class TestRunShipper: - def _env(self, tmp_path): - env_file = tmp_path / "client.env" - env_file.write_text( - "CLIENT_ID=42\n" - "VM_NAME=42\n" - "CLIENT_SECRET=s3cr3t\n" - "ALLOCATOR_URL=https://lablink.example.com\n" - ) - return env_file - - def test_flushes_full_batch_then_exits_on_missing_container( - self, tmp_path, monkeypatch - ): - from unittest.mock import MagicMock - from lablink_cli.log_shipper import run_shipper - - env_file = self._env(tmp_path) - - # 50 lines → triggers batch flush - lines = [ - f"2026-05-28T14:23:{i:02d}Z [start] line-{i}" for i in range(50) - ] - post_calls = [] - - def fake_post_batch(**kw): - post_calls.append(kw) - return "ok" - - def fake_iter(): - yield from lines - - monkeypatch.setattr( - "lablink_cli.log_shipper.post_batch", fake_post_batch - ) - mock_docker = MagicMock() - mock_docker.container_status.return_value = "missing" - monkeypatch.setattr( - "lablink_cli.log_shipper.default_docker", - lambda: mock_docker, - ) - # state_dir override so tests don't touch ~/.lablink - monkeypatch.setattr( - "lablink_cli.log_shipper.STATE_FILE", - tmp_path / "log_shipper.state", - ) - monkeypatch.setattr( - "lablink_cli.log_shipper.SELF_LOG_FILE", - tmp_path / "log_shipper.log", - ) - - run_shipper(env_file, _line_iter=fake_iter, _sleep=lambda s: None) - - assert len(post_calls) == 1 - assert len(post_calls[0]["messages"]) == 50 - # state updated with the timestamp of the last line in the batch - state = (tmp_path / "log_shipper.state").read_text() - assert "2026-05-28T14:23:49Z" in state - - def test_idle_tick_flushes_a_partial_buffer(self, tmp_path, monkeypatch): - """A short startup burst followed by silence must still be shipped. - - Regression: the flush check used to run only when a new line - arrived, so a container that logged at boot and then went quiet - held its lines forever and the allocator recorded nothing. - """ - from unittest.mock import MagicMock - from lablink_cli.log_shipper import TICK, run_shipper - - env_file = self._env(tmp_path) - post_calls = [] - - def fake_post_batch(**kw): - post_calls.append(kw) - return "ok" - - def fake_iter(): - yield "2026-05-28T14:23:01Z boot line 1" - yield "2026-05-28T14:23:02Z boot line 2" - # ...then the container goes quiet. Ticks are all we get. - yield TICK - yield TICK - - # Clock advances past FLUSH_INTERVAL_S while only ticks arrive. - clock = iter([0, 0, 0, 100, 200]) - monkeypatch.setattr( - "lablink_cli.log_shipper.time.monotonic", lambda: next(clock) - ) - monkeypatch.setattr( - "lablink_cli.log_shipper.post_batch", fake_post_batch - ) - mock_docker = MagicMock() - mock_docker.container_status.return_value = "missing" - monkeypatch.setattr( - "lablink_cli.log_shipper.default_docker", - lambda: mock_docker, - ) - monkeypatch.setattr( - "lablink_cli.log_shipper.STATE_FILE", - tmp_path / "log_shipper.state", - ) - monkeypatch.setattr( - "lablink_cli.log_shipper.SELF_LOG_FILE", - tmp_path / "log_shipper.log", - ) - - run_shipper(env_file, _line_iter=fake_iter, _sleep=lambda s: None) - - assert len(post_calls) == 1 - assert post_calls[0]["messages"] == [ - "2026-05-28T14:23:01Z boot line 1", - "2026-05-28T14:23:02Z boot line 2", - ] - - def test_exits_on_fatal_post(self, tmp_path, monkeypatch): - from unittest.mock import MagicMock - from lablink_cli.log_shipper import run_shipper - - env_file = self._env(tmp_path) - - def fake_iter(): - for i in range(50): - yield f"2026-05-28T14:23:{i:02d}Z [agent] x" - - monkeypatch.setattr( - "lablink_cli.log_shipper.post_batch", - lambda **kw: "fatal", - ) - mock_docker = MagicMock() - mock_docker.container_status.return_value = "running" - monkeypatch.setattr( - "lablink_cli.log_shipper.default_docker", - lambda: mock_docker, - ) - monkeypatch.setattr( - "lablink_cli.log_shipper.STATE_FILE", - tmp_path / "log_shipper.state", - ) - monkeypatch.setattr( - "lablink_cli.log_shipper.SELF_LOG_FILE", - tmp_path / "log_shipper.log", - ) - - run_shipper(env_file, _line_iter=fake_iter, _sleep=lambda s: None) - - log = (tmp_path / "log_shipper.log").read_text() - assert "fatal" in log.lower() or "exiting" in log.lower() - - def test_resumes_from_last_shipped_ts( - self, tmp_path, monkeypatch - ): - """When a state file exists, the docker logs --since arg matches it.""" - from unittest.mock import MagicMock - from lablink_cli.log_shipper import run_shipper - - env_file = self._env(tmp_path) - state_file = tmp_path / "log_shipper.state" - state_file.write_text('{"last_shipped_ts": "2026-05-28T14:00:00Z"}') - - monkeypatch.setattr( - "lablink_cli.log_shipper.STATE_FILE", state_file - ) - monkeypatch.setattr( - "lablink_cli.log_shipper.SELF_LOG_FILE", - tmp_path / "log_shipper.log", - ) - captured_since: list[str | None] = [] - - def fake_follow_logs(name, *, since=None): - captured_since.append(since) - # Yield no lines, end immediately - mock = MagicMock() - mock.stdout = iter([]) - mock.terminate = MagicMock() - mock.wait = MagicMock() - return mock - - mock_docker = MagicMock() - mock_docker.follow_logs.side_effect = fake_follow_logs - mock_docker.container_status.return_value = "missing" - - run_shipper(env_file, _sleep=lambda s: None, docker=mock_docker) - - assert captured_since == ["2026-05-28T14:00:00Z"] - - def test_exits_after_consecutive_exited_inspections( - self, tmp_path, monkeypatch - ): - """`docker stop` leaves the container in 'exited' state; shipper must - exit rather than busy-looping. With --restart unless-stopped, a - crashed container reaches 'restarting' within ms, so consecutive - 'exited' reliably means the user invoked stop.""" - from unittest.mock import MagicMock - from lablink_cli.log_shipper import run_shipper, MAX_EXITED_CONSECUTIVE - - env_file = self._env(tmp_path) - - monkeypatch.setattr( - "lablink_cli.log_shipper.STATE_FILE", - tmp_path / "log_shipper.state", - ) - monkeypatch.setattr( - "lablink_cli.log_shipper.SELF_LOG_FILE", - tmp_path / "log_shipper.log", - ) - - # Each docker-logs attach yields no lines (container is stopped). - def fake_follow_logs(name, *, since=None): - mock = MagicMock() - mock.stdout = iter([]) - return mock - - mock_docker = MagicMock() - mock_docker.follow_logs.side_effect = fake_follow_logs - mock_docker.container_status.return_value = "exited" - - run_shipper(env_file, _sleep=lambda s: None, docker=mock_docker) - - # Should give up after MAX_EXITED_CONSECUTIVE inspections — not loop - # forever. Pre-fix this test would hang indefinitely. - assert mock_docker.container_status.call_count == MAX_EXITED_CONSECUTIVE - - def test_restarting_resets_exited_counter( - self, tmp_path, monkeypatch - ): - """If docker restarts the container (crash + --restart unless-stopped), - the 'restarting' status should reset the exited counter so the shipper - doesn't accumulate exits across reconnects and prematurely give up.""" - from unittest.mock import MagicMock - from lablink_cli.log_shipper import run_shipper, MAX_EXITED_CONSECUTIVE - - env_file = self._env(tmp_path) - - monkeypatch.setattr( - "lablink_cli.log_shipper.STATE_FILE", - tmp_path / "log_shipper.state", - ) - monkeypatch.setattr( - "lablink_cli.log_shipper.SELF_LOG_FILE", - tmp_path / "log_shipper.log", - ) - - def fake_follow_logs(name, *, since=None): - mock = MagicMock() - mock.stdout = iter([]) - return mock - - # Alternate: exited, restarting, exited, restarting, ..., then a run - # of exited that should trigger the give-up. The intervening - # "restarting" results should keep resetting the counter. - sequence = ( - ["exited", "restarting"] * (MAX_EXITED_CONSECUTIVE + 2) - + ["exited"] * MAX_EXITED_CONSECUTIVE - ) - mock_docker = MagicMock() - mock_docker.follow_logs.side_effect = fake_follow_logs - mock_docker.container_status.side_effect = sequence - - run_shipper(env_file, _sleep=lambda s: None, docker=mock_docker) - - # Total inspections = alternating prefix + N trailing exited. - expected = 2 * (MAX_EXITED_CONSECUTIVE + 2) + MAX_EXITED_CONSECUTIVE - assert mock_docker.container_status.call_count == expected - - -class TestMainEntry: - def test_writes_pid_file_on_start(self, tmp_path, monkeypatch): - from unittest.mock import MagicMock - from lablink_cli import log_shipper - - env_file = tmp_path / "client.env" - env_file.write_text( - "CLIENT_ID=1\nVM_NAME=1\nCLIENT_SECRET=s\n" - "ALLOCATOR_URL=https://x\n" - ) - pid_file = tmp_path / "log_shipper.pid" - monkeypatch.setattr(log_shipper, "PID_FILE", pid_file) - monkeypatch.setattr( - log_shipper, "STATE_FILE", tmp_path / "log_shipper.state" - ) - monkeypatch.setattr( - log_shipper, "SELF_LOG_FILE", tmp_path / "log_shipper.log" - ) - monkeypatch.setattr( - log_shipper, "run_shipper", MagicMock() - ) - - log_shipper.main([str(env_file)]) - - # main() should have written the PID and then removed it on exit. - assert log_shipper.run_shipper.call_count == 1 - # PID file cleaned up after run_shipper returns - assert not pid_file.exists() - - def test_signal_handler_unlinks_pid_file( - self, tmp_path, monkeypatch - ): - import signal - from lablink_cli import log_shipper - - pid_file = tmp_path / "log_shipper.pid" - pid_file.write_text("12345") - monkeypatch.setattr(log_shipper, "PID_FILE", pid_file) - monkeypatch.setattr( - log_shipper, "SELF_LOG_FILE", tmp_path / "log_shipper.log" - ) - - with pytest.raises(SystemExit) as exc: - log_shipper._handle_shutdown(signal.SIGTERM, None) - - assert exc.value.code == 0 - assert not pid_file.exists() - - -def test_run_shipper_takes_a_docker_adapter(): - import inspect - - from lablink_cli.log_shipper import run_shipper - - assert "docker" in inspect.signature(run_shipper).parameters - - -def test_open_docker_logs_is_gone(): - import lablink_cli.log_shipper as ls - - assert not hasattr(ls, "open_docker_logs") diff --git a/packages/cli/tests/test_register.py b/packages/cli/tests/test_register.py index 65dc07fd3..fe0bb77ee 100644 --- a/packages/cli/tests/test_register.py +++ b/packages/cli/tests/test_register.py @@ -129,54 +129,28 @@ def _kwargs(env_file, **overrides): class TestResumePath: - @patch("lablink_cli.commands.register._start_log_shipper") - @patch("lablink_cli.commands.register._shipper_alive") - def test_everything_running_is_noop( - self, mock_alive, mock_spawn, tmp_env_file - ): + def test_everything_running_is_noop(self, tmp_env_file, capsys): from lablink_cli.commands.register import run_register tmp_env_file.write_text("CLIENT_ID=42\nCLIENT_SECRET=s\n") docker = RegisterDocker(status="running") - mock_alive.return_value = True run_register(**_kwargs(tmp_env_file), docker=docker) - # Should NOT start a new container, NOT spawn a new shipper. - mock_spawn.assert_not_called() - # And should NOT do a fresh `docker run`. + # Should NOT do a fresh `docker run`. assert docker.detached_argv is None - - @patch("lablink_cli.commands.register._start_log_shipper") - @patch("lablink_cli.commands.register._shipper_alive") - def test_dead_shipper_revived_no_re_register( - self, mock_alive, mock_spawn, tmp_env_file - ): - from lablink_cli.commands.register import run_register - tmp_env_file.write_text("CLIENT_ID=42\nCLIENT_SECRET=s\n") - docker = RegisterDocker(status="running") - mock_alive.return_value = False - - run_register(**_kwargs(tmp_env_file), docker=docker) - - mock_spawn.assert_called_once() + assert "Already registered" in capsys.readouterr().out # The secret didn't change — env file content untouched. assert "CLIENT_SECRET=s" in tmp_env_file.read_text() - @patch("lablink_cli.commands.register._start_log_shipper") - @patch("lablink_cli.commands.register._shipper_alive") - def test_exited_container_restarted( - self, mock_alive, mock_spawn, tmp_env_file, capsys - ): + def test_exited_container_restarted(self, tmp_env_file, capsys): from lablink_cli.commands.register import run_register tmp_env_file.write_text("CLIENT_ID=42\nCLIENT_SECRET=s\n") docker = RegisterDocker(status="exited") - mock_alive.return_value = False run_register(**_kwargs(tmp_env_file), docker=docker) # docker start lablink-client invoked and succeeded. assert "Restarted container" in capsys.readouterr().out - mock_spawn.assert_called_once() def test_exited_container_start_failure_is_descriptive_when_stderr_empty( self, tmp_path, capsys @@ -210,9 +184,7 @@ def test_force_still_re_registers( "lablink_cli.commands.register.RegistrationClient" ) as mock_client_cls, patch( "lablink_cli.commands.register.byo_detect" - ) as mock_detect, patch( - "lablink_cli.commands.register.subprocess.Popen" - ): + ) as mock_detect: mock_detect.detect_hostname.return_value = "byo-01" mock_detect.detect_lan_ip.return_value = "192.168.1.42" mock_detect.resolve_machine_identity.return_value = "mid" @@ -333,17 +305,19 @@ def test_success_skips_docker_and_prints_env( assert "OVERLAY_HOSTNAME=classroom-gpu-3" in out assert "TAILSCALE_AUTHKEY=tskey-abc" in out assert "CLIENT_SECRET=s" in out + # Hand-off clients have no host-side log shipper, so the pasted + # env must tell start.sh to ship the container's own stream. + assert "SHIP_LOGS=1" in out # We cannot mount anything on this path, so the operator must be # told to persist tailscaled's state themselves — otherwise every # workload restart mints a new tailnet node (lablink#404). assert "/var/lib/tailscale" in out - @patch("lablink_cli.commands.register.subprocess.Popen") @patch("lablink_cli.commands.register.RegistrationClient") @patch("lablink_cli.commands.register.byo_detect") def test_run_locally_default_autodetects_and_execs_docker( - self, mock_detect, mock_client_cls, mock_popen, + self, mock_detect, mock_client_cls, tmp_env_file, successful_response, ): """run_locally defaults to True: with --overlay-hostname alone @@ -438,20 +412,10 @@ def _register(tmp_env_file, resp, mock_detect, mock_client_cls): run_register(**_kwargs(tmp_env_file), docker=RegisterDocker()) return tmp_env_file.read_text() - # subprocess.Popen stays mocked (even though docker no longer shells - # out through it) — `_start_log_shipper` still spawns a REAL detached - # `python -m lablink_cli.log_shipper` process at the end of a - # successful run_register. That process constructs its own default - # Docker() adapter (this fake is only wired into THIS process), so - # leaving Popen unmocked here made every one of these tests spawn a - # real subprocess that queried the developer's actual docker daemon - # and appended to the real ~/.lablink/log_shipper.log — confirmed by - # running these tests before adding this patch back. - @patch("lablink_cli.commands.register.subprocess.Popen") @patch("lablink_cli.commands.register.RegistrationClient") @patch("lablink_cli.commands.register.byo_detect") def test_env_file_carries_repository_and_software( - self, mock_detect, mock_client_cls, mock_popen, + self, mock_detect, mock_client_cls, tmp_env_file, successful_response, ): resp = dict( @@ -466,11 +430,10 @@ def test_env_file_carries_repository_and_software( ) in content assert "SUBJECT_SOFTWARE=sleap" in content - @patch("lablink_cli.commands.register.subprocess.Popen") @patch("lablink_cli.commands.register.RegistrationClient") @patch("lablink_cli.commands.register.byo_detect") def test_env_file_omits_empty_repository_and_software( - self, mock_detect, mock_client_cls, mock_popen, + self, mock_detect, mock_client_cls, tmp_env_file, successful_response, ): """An unset cfg.machine.repository must leave the var out entirely @@ -482,11 +445,10 @@ def test_env_file_omits_empty_repository_and_software( assert "TUTORIAL_REPO_TO_CLONE" not in content assert "SUBJECT_SOFTWARE" not in content - @patch("lablink_cli.commands.register.subprocess.Popen") @patch("lablink_cli.commands.register.RegistrationClient") @patch("lablink_cli.commands.register.byo_detect") def test_env_file_omits_vars_when_allocator_predates_fix( - self, mock_detect, mock_client_cls, mock_popen, + self, mock_detect, mock_client_cls, tmp_env_file, successful_response, ): """A newer CLI against an older allocator gets a response with @@ -500,11 +462,10 @@ def test_env_file_omits_vars_when_allocator_predates_fix( class TestSuccessFlow: - @patch("lablink_cli.commands.register.subprocess.Popen") @patch("lablink_cli.commands.register.RegistrationClient") @patch("lablink_cli.commands.register.byo_detect") def test_full_success_writes_env_file_and_execs_docker( - self, mock_detect, mock_client_cls, mock_popen, + self, mock_detect, mock_client_cls, tmp_env_file, successful_response, ): from lablink_cli.commands.register import run_register @@ -563,17 +524,14 @@ def test_full_success_writes_env_file_and_execs_docker( assert cmd[pull_idx + 1] == "always" assert "ghcr.io/talmolab/lablink-client:0.4.0" in cmd - # Log shipper spawned exactly once with the env file - assert mock_popen.call_count == 1 - shipper_cmd = mock_popen.call_args.args[0] - assert "lablink_cli.log_shipper" in shipper_cmd - assert str(tmp_env_file) in shipper_cmd + # The container ships its own logs (start.sh's ship_logs worker); + # the env file must carry the switch that turns that on. + assert "SHIP_LOGS=1" in content - @patch("lablink_cli.commands.register.subprocess.Popen") @patch("lablink_cli.commands.register.RegistrationClient") @patch("lablink_cli.commands.register.byo_detect") def test_user_overrides_beat_detection( - self, mock_detect, mock_client_cls, mock_popen, + self, mock_detect, mock_client_cls, tmp_env_file, successful_response, ): from lablink_cli.commands.register import run_register @@ -620,11 +578,10 @@ def test_missing_lan_ip_aborts(self, mock_detect, tmp_env_file): with pytest.raises(SystemExit): run_register(**_kwargs(tmp_env_file)) - @patch("lablink_cli.commands.register.subprocess.Popen") @patch("lablink_cli.commands.register.RegistrationClient") @patch("lablink_cli.commands.register.byo_detect") def test_gpu_present_override_keeps_detected_model( - self, mock_detect, mock_client_cls, mock_popen, + self, mock_detect, mock_client_cls, tmp_env_file, successful_response, ): """User passes --gpu-present (no --gpu-model); detection still provides @@ -651,11 +608,10 @@ def test_gpu_present_override_keeps_detected_model( gpu_model="NVIDIA A100", # detection-supplied fallback ) - @patch("lablink_cli.commands.register.subprocess.Popen") @patch("lablink_cli.commands.register.RegistrationClient") @patch("lablink_cli.commands.register.byo_detect") def test_gpu_model_override_wins_over_detection( - self, mock_detect, mock_client_cls, mock_popen, + self, mock_detect, mock_client_cls, tmp_env_file, successful_response, ): from lablink_cli.commands.register import run_register @@ -681,11 +637,10 @@ def test_gpu_model_override_wins_over_detection( gpu_model="USER_PROVIDED", ) - @patch("lablink_cli.commands.register.subprocess.Popen") @patch("lablink_cli.commands.register.RegistrationClient") @patch("lablink_cli.commands.register.byo_detect") def test_publishes_agent_and_kasmvnc_ports_not_network_host( - self, mock_detect, mock_client_cls, mock_popen, + self, mock_detect, mock_client_cls, tmp_env_file, successful_response, ): """The allocator reaches the BYO client over the LAN at @@ -779,11 +734,10 @@ def test_aborts_when_cgroup_driver_systemd_and_gpu_present( # docker run must not fire when the cgroup driver check fails. assert docker.detached_argv is None - @patch("lablink_cli.commands.register.subprocess.Popen") @patch("lablink_cli.commands.register.RegistrationClient") @patch("lablink_cli.commands.register.byo_detect") def test_skips_cgroup_check_when_gpu_absent( - self, mock_detect, mock_client_cls, mock_popen, + self, mock_detect, mock_client_cls, tmp_env_file, successful_response, ): """CPU-only BYO clients don't need GPU runtime — never query @@ -913,11 +867,10 @@ def test_reports_nonzero_exit_when_stderr_empty(self, capsys): class TestForceFlag: - @patch("lablink_cli.commands.register.subprocess.Popen") @patch("lablink_cli.commands.register.RegistrationClient") @patch("lablink_cli.commands.register.byo_detect") def test_force_overwrites_env_file( - self, mock_detect, mock_client_cls, mock_popen, + self, mock_detect, mock_client_cls, tmp_env_file, successful_response, ): from lablink_cli.commands.register import run_register @@ -937,11 +890,10 @@ def test_force_overwrites_env_file( assert "CLIENT_ID=42" in tmp_env_file.read_text() assert "CLIENT_ID=99" not in tmp_env_file.read_text() - @patch("lablink_cli.commands.register.subprocess.Popen") @patch("lablink_cli.commands.register.RegistrationClient") @patch("lablink_cli.commands.register.byo_detect") def test_force_removes_existing_container_before_run( - self, mock_detect, mock_client_cls, mock_popen, + self, mock_detect, mock_client_cls, tmp_env_file, successful_response, ): """With --force, the orchestrator must `docker rm -f` the old container @@ -1007,268 +959,6 @@ def test_auth_error_exits_nonzero( assert not tmp_env_file.exists() -class TestStartLogShipper: - @patch("lablink_cli.commands.register.subprocess.Popen") - def test_spawns_detached_python_module( - self, mock_popen, tmp_path, monkeypatch - ): - from lablink_cli.commands.register import _start_log_shipper - from rich.console import Console - - # Point PID_FILE at tmp_path so _stop_existing_shipper doesn't - # touch the real ~/.lablink/log_shipper.pid (would otherwise risk - # terminating a live shipper on the developer's machine). - monkeypatch.setattr( - "lablink_cli.commands.register.PID_FILE", - tmp_path / "log_shipper.pid", - ) - - env_file = tmp_path / "client.env" - env_file.write_text("CLIENT_ID=1\n") - mock_popen.return_value = MagicMock(pid=99999) - - _start_log_shipper(env_file, Console()) - - assert mock_popen.call_count == 1 - cmd = mock_popen.call_args.args[0] - # invoked as: python -m lablink_cli.log_shipper - assert cmd[0].endswith("python") or "python" in cmd[0] - assert "-m" in cmd - assert "lablink_cli.log_shipper" in cmd - assert str(env_file) in cmd - # detached: start_new_session on POSIX OR Windows creationflags - kwargs = mock_popen.call_args.kwargs - assert kwargs.get("start_new_session") is True or ( - kwargs.get("creationflags", 0) != 0 - ) - # stdin closed; stdout/stderr to log file - assert kwargs.get("stdin") is not None # DEVNULL - - @patch("lablink_cli.commands.register._stop_existing_shipper") - @patch("lablink_cli.commands.register.subprocess.Popen") - def test_terminates_existing_shipper_before_spawn( - self, mock_popen, mock_stop, tmp_path - ): - """Guarantees there is no overlap between old and new shippers - (which would POST duplicates under --force re-register).""" - from unittest.mock import Mock - from lablink_cli.commands.register import _start_log_shipper - from rich.console import Console - - env_file = tmp_path / "client.env" - env_file.write_text("CLIENT_ID=1\n") - mock_popen.return_value = MagicMock(pid=99999) - - # Attach both mocks to a parent so we get a single ordered call - # log. Asserting the call names appear in source order catches a - # future refactor that moves Popen above _stop_existing_shipper. - parent = Mock() - parent.attach_mock(mock_stop, "stop") - parent.attach_mock(mock_popen, "popen") - - _start_log_shipper(env_file, Console()) - - call_names = [c[0] for c in parent.mock_calls] - assert call_names == ["stop", "popen"], ( - f"expected stop -> popen, got {call_names}" - ) - - -class TestStopExistingShipper: - """Covers `_stop_existing_shipper` — the kill-old-shipper step that - prevents double-shipper duplicate POSTs during --force re-register.""" - - def _fake_psutil(self, monkeypatch, process_factory): - """Install a fake psutil module whose `Process()` returns whatever - process_factory() yields. Exceptions are re-exported as classes so - ``except psutil.NoSuchProcess`` works in the SUT.""" - from unittest.mock import MagicMock - from lablink_cli.commands import register - - fake = MagicMock() - fake.NoSuchProcess = type("NoSuchProcess", (Exception,), {}) - fake.AccessDenied = type("AccessDenied", (Exception,), {}) - fake.TimeoutExpired = type("TimeoutExpired", (Exception,), {}) - fake.Process.side_effect = process_factory - monkeypatch.setattr(register, "psutil", fake) - return fake - - def test_no_pid_file_is_noop(self, tmp_path, monkeypatch): - from lablink_cli.commands import register - from rich.console import Console - - pid_file = tmp_path / "log_shipper.pid" - monkeypatch.setattr(register, "PID_FILE", pid_file) - - # Should not raise; should not touch psutil at all. - register._stop_existing_shipper(Console()) - - def test_terminates_matching_shipper(self, tmp_path, monkeypatch): - from unittest.mock import MagicMock - from lablink_cli.commands import register - from rich.console import Console - - pid_file = tmp_path / "log_shipper.pid" - pid_file.write_text("12345") - monkeypatch.setattr(register, "PID_FILE", pid_file) - - fake_proc = MagicMock() - fake_proc.cmdline.return_value = [ - "/usr/bin/python", "-m", - "lablink_cli.log_shipper", "/x/client.env", - ] - self._fake_psutil(monkeypatch, lambda _pid: fake_proc) - - register._stop_existing_shipper(Console()) - - fake_proc.terminate.assert_called_once() - fake_proc.wait.assert_called_once() - # PID file cleared so a stale entry never confuses the next run. - assert not pid_file.exists() - - def test_escalates_to_kill_on_timeout(self, tmp_path, monkeypatch): - from unittest.mock import MagicMock - from lablink_cli.commands import register - from rich.console import Console - - pid_file = tmp_path / "log_shipper.pid" - pid_file.write_text("12345") - monkeypatch.setattr(register, "PID_FILE", pid_file) - - fake_proc = MagicMock() - fake_proc.cmdline.return_value = [ - "python", "-m", "lablink_cli.log_shipper", "/x" - ] - fake = self._fake_psutil(monkeypatch, lambda _pid: fake_proc) - # SIGTERM didn't take — wait() raises TimeoutExpired. - fake_proc.wait.side_effect = fake.TimeoutExpired() - - register._stop_existing_shipper(Console()) - - fake_proc.terminate.assert_called_once() - fake_proc.kill.assert_called_once() - # PID file still cleaned up after SIGKILL (handler never ran). - assert not pid_file.exists() - - def test_does_not_kill_unrelated_pid(self, tmp_path, monkeypatch): - """PID file points at a real but unrelated process (e.g. PID reuse - after reboot). Cmdline guard must protect it.""" - from unittest.mock import MagicMock - from lablink_cli.commands import register - from rich.console import Console - - pid_file = tmp_path / "log_shipper.pid" - pid_file.write_text("12345") - monkeypatch.setattr(register, "PID_FILE", pid_file) - - fake_proc = MagicMock() - fake_proc.cmdline.return_value = ["/usr/bin/vim", "notes.txt"] - self._fake_psutil(monkeypatch, lambda _pid: fake_proc) - - register._stop_existing_shipper(Console()) - - fake_proc.terminate.assert_not_called() - fake_proc.kill.assert_not_called() - # Stale PID file dropped so we don't keep skipping forever. - assert not pid_file.exists() - - def test_stale_pid_removes_file_silently(self, tmp_path, monkeypatch): - """PID file references a PID that no longer exists.""" - from lablink_cli.commands import register - from rich.console import Console - - pid_file = tmp_path / "log_shipper.pid" - pid_file.write_text("99999") - monkeypatch.setattr(register, "PID_FILE", pid_file) - - fake = self._fake_psutil(monkeypatch, None) - fake.Process.side_effect = fake.NoSuchProcess() - - register._stop_existing_shipper(Console()) - - assert not pid_file.exists() - - def test_corrupt_pid_file_removed(self, tmp_path, monkeypatch): - from lablink_cli.commands import register - from rich.console import Console - - pid_file = tmp_path / "log_shipper.pid" - pid_file.write_text("not-a-number") - monkeypatch.setattr(register, "PID_FILE", pid_file) - - register._stop_existing_shipper(Console()) - - assert not pid_file.exists() - - -class TestShipperAlive: - def test_no_pid_file_returns_false(self, tmp_path, monkeypatch): - from lablink_cli.commands import register - - pid_file = tmp_path / "log_shipper.pid" - monkeypatch.setattr(register, "PID_FILE", pid_file) - - assert register._shipper_alive() is False - - def test_pid_with_matching_cmdline_returns_true( - self, tmp_path, monkeypatch - ): - from unittest.mock import MagicMock - from lablink_cli.commands import register - - pid_file = tmp_path / "log_shipper.pid" - pid_file.write_text("12345") - monkeypatch.setattr(register, "PID_FILE", pid_file) - - fake_proc = MagicMock() - fake_proc.cmdline.return_value = [ - "/usr/bin/python", "-m", - "lablink_cli.log_shipper", "/home/u/.lablink/client.env", - ] - fake_psutil = MagicMock() - fake_psutil.Process.return_value = fake_proc - fake_psutil.NoSuchProcess = type("NoSuchProcess", (Exception,), {}) - fake_psutil.AccessDenied = type("AccessDenied", (Exception,), {}) - monkeypatch.setattr(register, "psutil", fake_psutil) - - assert register._shipper_alive() is True - - def test_pid_with_wrong_cmdline_returns_false( - self, tmp_path, monkeypatch - ): - from unittest.mock import MagicMock - from lablink_cli.commands import register - - pid_file = tmp_path / "log_shipper.pid" - pid_file.write_text("12345") - monkeypatch.setattr(register, "PID_FILE", pid_file) - - fake_proc = MagicMock() - fake_proc.cmdline.return_value = ["/usr/bin/vim"] - fake_psutil = MagicMock() - fake_psutil.Process.return_value = fake_proc - fake_psutil.NoSuchProcess = type("NoSuchProcess", (Exception,), {}) - fake_psutil.AccessDenied = type("AccessDenied", (Exception,), {}) - monkeypatch.setattr(register, "psutil", fake_psutil) - - assert register._shipper_alive() is False - - def test_dead_pid_returns_false(self, tmp_path, monkeypatch): - from unittest.mock import MagicMock - from lablink_cli.commands import register - - pid_file = tmp_path / "log_shipper.pid" - pid_file.write_text("12345") - monkeypatch.setattr(register, "PID_FILE", pid_file) - - fake_psutil = MagicMock() - fake_psutil.NoSuchProcess = type("NoSuchProcess", (Exception,), {}) - fake_psutil.AccessDenied = type("AccessDenied", (Exception,), {}) - fake_psutil.Process.side_effect = fake_psutil.NoSuchProcess() - monkeypatch.setattr(register, "psutil", fake_psutil) - - assert register._shipper_alive() is False - class TestWriteStartupScript: """Covers `_write_startup_script` — decodes the allocator-shipped @@ -1346,11 +1036,10 @@ class TestDockerRunMountsStartupScript: This is the actual delivery path the bug fix is closing. """ - @patch("lablink_cli.commands.register.subprocess.Popen") @patch("lablink_cli.commands.register.RegistrationClient") @patch("lablink_cli.commands.register.byo_detect") def test_run_register_mounts_script_when_allocator_ships_one( - self, mock_detect, mock_client_cls, mock_popen, + self, mock_detect, mock_client_cls, tmp_env_file, successful_response, ): import base64 @@ -1418,11 +1107,10 @@ def test_run_register_mounts_script_when_allocator_ships_one( f"STARTUP_SUCCESS_CHECK_B64 not forwarded; -e args={env_args}" ) - @patch("lablink_cli.commands.register.subprocess.Popen") @patch("lablink_cli.commands.register.RegistrationClient") @patch("lablink_cli.commands.register.byo_detect") def test_run_register_skips_mount_when_no_script( - self, mock_detect, mock_client_cls, mock_popen, + self, mock_detect, mock_client_cls, tmp_env_file, successful_response, ): """Allocator returned empty payload (script disabled) → no @@ -1470,16 +1158,6 @@ def test_run_register_skips_mount_when_no_script( f"script; got {env_args}" ) - def test_corrupt_pid_file_returns_false(self, tmp_path, monkeypatch): - from lablink_cli.commands import register - - pid_file = tmp_path / "log_shipper.pid" - pid_file.write_text("not-a-number") - monkeypatch.setattr(register, "PID_FILE", pid_file) - - assert register._shipper_alive() is False - - class TestWriteEnvFile: """`_write_env_file` must propagate the full register response so the client container's start.sh can extract the Tier 1 monitoring block @@ -1577,6 +1255,20 @@ def test_overlay_fields_absent_when_not_given(self, tmp_env_file): assert "OVERLAY_HOSTNAME" not in content assert "TAILSCALE_AUTHKEY" not in content + def test_ship_logs_always_written(self, tmp_env_file): + """SHIP_LOGS=1 turns on the client container's own log shipper + (start.sh's ship_logs worker). Every BYO shape needs it — there + is no host-side shipper in any of them — so it is written + unconditionally.""" + from lablink_cli.commands.register import _write_env_file + + _write_env_file( + tmp_env_file, + self._resp_with_monitoring(), + allocator_url="https://lablink.example.com", + ) + assert "SHIP_LOGS=1" in tmp_env_file.read_text() + def test_prefers_caller_url_over_downgraded_response_url(self, tmp_env_file): """Regression (P1 review finding): a mesh-overlay/Funnel allocator's register response derives allocator_url from Flask's diff --git a/packages/client/pyproject.toml b/packages/client/pyproject.toml index 0b1a2e1b0..7365f5c9b 100644 --- a/packages/client/pyproject.toml +++ b/packages/client/pyproject.toml @@ -42,6 +42,7 @@ agent = "lablink_client_service.agent.api:main" check_gpu = "lablink_client_service.check_gpu:main" heartbeat = "lablink_client_service.heartbeat:main" lablink-monitoring = "lablink_client_service.monitoring.__main__:main" +ship_logs = "lablink_client_service.ship_logs:main" update_inuse_status = "lablink_client_service.update_inuse_status:main" [tool.setuptools.package-data] diff --git a/packages/client/src/lablink_client_service/ship_logs.py b/packages/client/src/lablink_client_service/ship_logs.py new file mode 100644 index 000000000..1feee47bb --- /dev/null +++ b/packages/client/src/lablink_client_service/ship_logs.py @@ -0,0 +1,233 @@ +"""Ships this container's own log stream to the allocator. + +Log shipping is host-side in the AWS topology (user_data's +log_shipper.sh tails ``docker logs``), but BYO clients have no +LabLink-controlled host process to rely on: run-locally boxes used a +detached CLI shipper that could die silently on operator laptops, and +hand-off clients (``register --no-run-locally``, e.g. a Run:AI +workload) never had any shipper at all — the container is its own +PID 1 with no docker daemon in sight. So BYO containers ship their own +stream: start.sh points fd 5 — the [tag]-prefixed output of every +service — at this worker's stdin when SHIP_LOGS=1. + +That puts this worker IN the logging path, which imposes two hard +rules (lablink#304's silent tail stall is the cautionary tale): + +* Passthrough first. Every stdin line is written to stdout (container + PID-1 stdout, i.e. ``docker logs``) before anything else touches it, + and the read loop never blocks on the network — shipping happens on + a separate thread fed by a bounded drop-oldest queue. +* Fail open. Missing env degrades to pure passthrough, and start.sh's + supervisor ``exec cat``s after repeated crashes — the worst case is + logs visible in ``docker logs`` but absent from the allocator, never + a frozen container. +""" + +import os +import signal +import sys +import threading +import time +from collections import deque +from datetime import datetime, timezone +from typing import Callable + +import requests + +from lablink_client_service.http_utils import get_auth_headers, sanitize_url + +BATCH_SIZE = 50 +FLUSH_INTERVAL_S = 15 +POLL_INTERVAL_S = 1.0 +POST_TIMEOUT_S = 10 +MAX_RETRIES = 3 +RETRY_BACKOFF_S = (1, 2, 4) +# Bounds shipping memory during output bursts (a pip install can emit +# hundreds of lines/second). Overflow drops the OLDEST lines: when the +# allocator is unreachable for a while, the most recent output is what +# an operator debugging the client needs. +QUEUE_MAX_LINES = 2000 +# The allocator routes a batch to the docker_logs column when log_group +# ends with "-docker" (vm_telemetry.py's receive_vm_logs); the +# "container" prefix names the source for anyone reading the DB. +LOG_GROUP = "container-docker" + + +def read_loop(queue: deque, stdin=None, stdout=None) -> None: + """Forward stdin to stdout line by line, queueing stamped copies. + + The passthrough write happens FIRST and the loop touches nothing + that can block indefinitely besides stdin itself, so ``docker + logs`` output is byte-identical to running without this worker. + Returns on EOF (container shutdown). + """ + stdin = stdin if stdin is not None else sys.stdin + stdout = stdout if stdout is not None else sys.stdout + # iter(readline, ""), not `for line in stdin`: file iteration may + # read ahead into an internal buffer and sit on complete lines, + # which in this seat delays every service's logs. + for line in iter(stdin.readline, ""): + stdout.write(line) + stdout.flush() + # Stamped at read time, not emission time — equivalent here + # (the passthrough seat sees lines the moment services emit + # them) and avoids parsing arbitrary service output. + ts = datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") + queue.append(f"{ts} {line.rstrip()}") + + +def drain(queue: deque, max_items: int = BATCH_SIZE) -> list: + """Pop up to ``max_items`` lines from the left of ``queue``.""" + batch: list = [] + while len(batch) < max_items: + try: + batch.append(queue.popleft()) + except IndexError: + break + return batch + + +def post_batch( + *, + base_url: str, + vm_name: str, + headers: dict, + messages: list, + retries: int = MAX_RETRIES, + _sleep: Callable[[float], None] = time.sleep, +) -> bool: + """POST one batch to /api/vm-logs/. True on 2xx. + + Retries with backoff, then reports failure so the caller drops the + batch. Unlike the deleted CLI shipper it never treats a 4xx as + fatal: this worker is the only log channel BYO clients have and + nothing respawns it for shipping-only failures, so a transient 401 + while the allocator's DB warms up must degrade to a dropped batch, + not a permanently dead shipper. + """ + url = f"{base_url}/api/vm-logs/{vm_name}" + payload = {"log_group": LOG_GROUP, "messages": messages} + for attempt in range(retries): + if attempt: + _sleep(RETRY_BACKOFF_S[attempt - 1]) + try: + resp = requests.post( + url, json=payload, headers=headers, timeout=POST_TIMEOUT_S + ) + if 200 <= resp.status_code < 300: + return True + except requests.exceptions.RequestException: + pass + return False + + +def shipper_loop( + queue: deque, + stop_event: threading.Event, + post_fn: Callable[[list, int], bool], + *, + _monotonic: Callable[[], float] = time.monotonic, + poll_s: float = POLL_INTERVAL_S, +) -> None: + """Flush the queue on the 50-line / 15-second rule until stopped. + + ``post_fn(messages, retries)`` does the actual POST. The stop path + performs a final single-attempt flush (no backoff) so a graceful + ``docker stop`` ships the tail inside docker's grace period. + """ + last_flush = _monotonic() + while True: + stopped = stop_event.wait(poll_s) + now = _monotonic() + if not queue: + last_flush = now + elif ( + stopped + or len(queue) >= BATCH_SIZE + or now - last_flush >= FLUSH_INTERVAL_S + ): + retries = 1 if stopped else MAX_RETRIES + while True: + batch = drain(queue) + if not batch: + break + if not post_fn(batch, retries): + print( + f"ship_logs: dropped {len(batch)} lines after " + "retries", + file=sys.stderr, + flush=True, + ) + last_flush = now + if stopped: + return + + +def main() -> None: + """Entry point for the ``ship_logs`` console script.""" + allocator_url = os.environ.get("ALLOCATOR_URL") + client_secret = os.environ.get("CLIENT_SECRET") + vm_name = os.environ.get("VM_NAME") + if not (allocator_url and client_secret and vm_name): + # Fail open: stay a pure passthrough so logging never depends + # on registration plumbing being complete. + print( + "ship_logs: ALLOCATOR_URL/CLIENT_SECRET/VM_NAME not all set; " + "passing lines through without shipping", + file=sys.stderr, + flush=True, + ) + read_loop(deque(maxlen=1)) + return + + base_url = sanitize_url(allocator_url) + headers = {"Content-Type": "application/json"} + headers.update(get_auth_headers(client_secret)) + + def post_fn(messages: list, retries: int) -> bool: + return post_batch( + base_url=base_url, + vm_name=vm_name, + headers=headers, + messages=messages, + retries=retries, + ) + + queue: deque = deque(maxlen=QUEUE_MAX_LINES) + stop_event = threading.Event() + signaled = threading.Event() + + def shipper() -> None: + shipper_loop(queue, stop_event, post_fn) + if signaled.is_set(): + # The main thread is blocked in readline and only docker's + # SIGKILL would end it; exit here so `docker stop` is fast + # and the final flush above still happened. + os._exit(0) + + thread = threading.Thread(target=shipper, daemon=True) + thread.start() + + def handle_stop(signum, _frame) -> None: + signaled.set() + stop_event.set() + + for sig in (signal.SIGTERM, signal.SIGINT): + try: + signal.signal(sig, handle_stop) + except (ValueError, AttributeError): + pass + + print( + f"ship_logs: forwarding to {base_url}/api/vm-logs/{vm_name}", + file=sys.stderr, + flush=True, + ) + read_loop(queue) + # EOF: every fd-5 writer is gone. Final flush, then exit. + stop_event.set() + thread.join(timeout=POST_TIMEOUT_S + 5) + + +if __name__ == "__main__": + main() diff --git a/packages/client/start.sh b/packages/client/start.sh index 5fdda14eb..d93feb4a9 100644 --- a/packages/client/start.sh +++ b/packages/client/start.sh @@ -11,7 +11,38 @@ CONTAINER_START_TIME=$(date +%s) # services are launched with their own `... | sed ... >&5 &` pipeline, so # the inner sed writes directly to fd 5 and bypasses the [start] tagger # (otherwise lines would be double-tagged as "[start] [agent] ..."). -exec 5>&1 +# When SHIP_LOGS=1 (written into the env file by `lablink client register` +# for BYO clients), fd 5 feeds the in-container ship_logs worker instead of +# going straight to stdout. The worker passes every line through to this +# container's stdout unchanged (`docker logs` output is identical either +# way) and forwards a copy to the allocator's /api/vm-logs. BYO clients +# need this because they have no host-side shipper: hand-off clients +# (--no-run-locally, e.g. Run:AI) are their own PID 1 with no docker daemon +# anywhere, and the CLI's old detached shipper for run-locally boxes died +# silently on operator laptops. +# +# The worker sits IN the logging path, so its failure must degrade, never +# block (lablink#304's silent tail stall is the cautionary tale): the +# supervisor loop respawns a crashed worker — the kernel's pipe buffer +# absorbs the gap — and after 3 crashes it execs plain `cat`, which is +# logging exactly as if SHIP_LOGS were unset. Absolute venv path because +# this runs before the venv activation below, and the process-substitution +# subshell never sees it anyway. +if [ "${SHIP_LOGS:-0}" = "1" ]; then + exec 5> >( + fails=0 + while :; do + /home/client/.venv/bin/ship_logs && break + fails=$((fails+1)) + if [ "$fails" -ge 3 ]; then + echo "ship_logs crashed $fails times; falling back to passthrough" >&2 + exec cat + fi + done + ) +else + exec 5>&1 +fi exec > >(sed -u 's/^/[start] /' >&5) 2>&1 # ----------------------------------------------------------------------- diff --git a/packages/client/tests/test_ship_logs.py b/packages/client/tests/test_ship_logs.py new file mode 100644 index 000000000..7329f5dfd --- /dev/null +++ b/packages/client/tests/test_ship_logs.py @@ -0,0 +1,152 @@ +"""Tests for the in-container log shipper (BYO clients).""" + +import io +import re +import threading +from collections import deque +from unittest.mock import MagicMock, patch + +import requests + +from lablink_client_service import ship_logs +from lablink_client_service.ship_logs import ( + BATCH_SIZE, + LOG_GROUP, + drain, + post_batch, + read_loop, + shipper_loop, +) + + +class TestReadLoop: + def test_passthrough_is_byte_identical_and_ordered(self): + src = "[start] one\n[agent] two\n[kasmvnc] three\n" + stdout = io.StringIO() + q: deque = deque() + read_loop(q, stdin=io.StringIO(src), stdout=stdout) + assert stdout.getvalue() == src + + def test_queued_copies_are_timestamped(self): + q: deque = deque() + read_loop(q, stdin=io.StringIO("hello\n"), stdout=io.StringIO()) + assert len(q) == 1 + assert re.match( + r"^\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}Z hello$", q[0] + ) + + def test_overflow_drops_oldest_but_passthrough_survives(self): + src = "".join(f"line {i}\n" for i in range(10)) + stdout = io.StringIO() + q: deque = deque(maxlen=3) + read_loop(q, stdin=io.StringIO(src), stdout=stdout) + assert stdout.getvalue() == src # passthrough never drops + assert [line.split(" ", 1)[1] for line in q] == [ + "line 7", "line 8", "line 9", + ] + + +class TestDrain: + def test_drains_in_batches_leaving_remainder(self): + q = deque(range(BATCH_SIZE + 2)) + first = drain(q) + assert len(first) == BATCH_SIZE + assert drain(q) == [BATCH_SIZE, BATCH_SIZE + 1] + assert drain(q) == [] + + +class TestPostBatch: + def test_posts_to_vm_logs_with_docker_suffix_group(self): + resp = MagicMock(status_code=200) + with patch.object(ship_logs.requests, "post", return_value=resp) as post: + ok = post_batch( + base_url="http://alloc", + vm_name="runai-client-2", + headers={"Authorization": "Bearer sekrit"}, + messages=["a"], + ) + assert ok + assert post.call_args.args == ("http://alloc/api/vm-logs/runai-client-2",) + payload = post.call_args.kwargs["json"] + # The allocator routes to the docker_logs column only when the + # group ends with "-docker" (receive_vm_logs). + assert payload["log_group"] == LOG_GROUP + assert LOG_GROUP.endswith("-docker") + assert payload["messages"] == ["a"] + + def test_retries_then_reports_failure_without_raising(self): + err = requests.exceptions.ConnectionError("nope") + with patch.object(ship_logs.requests, "post", side_effect=err) as post: + ok = post_batch( + base_url="http://alloc", + vm_name="vm", + headers={}, + messages=["a"], + _sleep=lambda s: None, + ) + assert not ok + assert post.call_count == ship_logs.MAX_RETRIES + + def test_4xx_is_dropped_not_fatal(self): + # A 4xx must not kill this worker — nothing respawns it for + # shipping-only failures. post_batch just reports failure. + resp = MagicMock(status_code=401) + with patch.object(ship_logs.requests, "post", return_value=resp): + ok = post_batch( + base_url="http://alloc", + vm_name="vm", + headers={}, + messages=["a"], + _sleep=lambda s: None, + ) + assert not ok + + def test_single_retry_mode_for_final_flush(self): + err = requests.exceptions.ConnectionError("nope") + with patch.object(ship_logs.requests, "post", side_effect=err) as post: + post_batch( + base_url="http://alloc", + vm_name="vm", + headers={}, + messages=["a"], + retries=1, + _sleep=lambda s: None, + ) + assert post.call_count == 1 + + +class TestShipperLoop: + def test_stop_flushes_everything_with_single_retry(self): + q = deque(f"l{i}" for i in range(BATCH_SIZE + 2)) + stop = threading.Event() + stop.set() + calls = [] + shipper_loop( + q, stop, lambda msgs, retries: calls.append((msgs, retries)) or True, + poll_s=0, + ) + assert [len(m) for m, _ in calls] == [BATCH_SIZE, 2] + assert all(retries == 1 for _, retries in calls) + assert not q + + def test_interval_elapse_flushes_small_buffer(self): + q = deque(["only-line"]) + stop = threading.Event() + calls = [] + clock = iter([0.0, 20.0, 20.0, 40.0]) # jumps past FLUSH_INTERVAL_S + + def post(msgs, retries): + calls.append(msgs) + stop.set() # end the loop after the first flush + return True + + shipper_loop(q, stop, post, _monotonic=lambda: next(clock), poll_s=0) + assert calls and calls[0] == ["only-line"] + + def test_failed_post_drops_batch_and_continues(self, capsys): + q = deque(["a", "b"]) + stop = threading.Event() + stop.set() + shipper_loop(q, stop, lambda msgs, retries: False, poll_s=0) + assert not q # dropped, not retained + assert "dropped 2 lines" in capsys.readouterr().err From 496562b859ca7236a0868fcaa81beb9c00681ff4 Mon Sep 17 00:00:00 2001 From: 7174Andy Date: Wed, 26 Aug 2026 16:15:27 -0700 Subject: [PATCH 2/2] test(client): cover ship_logs' console-script entry point MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit CI's client coverage gate failed at 89% (fail-under=90): every worker primitive was tested but main() itself — env wiring, the shipper thread, signal-handler registration, the EOF final flush, and the missing-env passthrough degrade — was 34 uncovered statements (module at 66%). Two end-to-end tests through main() bring the module to 96%. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01AzEx3SZp9EAmTVKe6JEWRa --- packages/client/tests/test_ship_logs.py | 67 +++++++++++++++++++++++++ 1 file changed, 67 insertions(+) diff --git a/packages/client/tests/test_ship_logs.py b/packages/client/tests/test_ship_logs.py index 7329f5dfd..9b7c4370a 100644 --- a/packages/client/tests/test_ship_logs.py +++ b/packages/client/tests/test_ship_logs.py @@ -2,6 +2,7 @@ import io import re +import signal import threading from collections import deque from unittest.mock import MagicMock, patch @@ -150,3 +151,69 @@ def test_failed_post_drops_batch_and_continues(self, capsys): shipper_loop(q, stop, lambda msgs, retries: False, poll_s=0) assert not q # dropped, not retained assert "dropped 2 lines" in capsys.readouterr().err + + +class TestMain: + """Covers the console-script entry point's wiring.""" + + def _clear_env(self, monkeypatch): + for var in ("ALLOCATOR_URL", "CLIENT_SECRET", "VM_NAME"): + monkeypatch.delenv(var, raising=False) + + def test_missing_env_degrades_to_pure_passthrough( + self, monkeypatch, capsys + ): + """Fail open: logging must never depend on registration plumbing + being complete — lines still reach stdout, nothing is shipped.""" + self._clear_env(monkeypatch) + monkeypatch.setattr("sys.stdin", io.StringIO("a\nb\n")) + with patch.object(ship_logs.requests, "post") as post: + ship_logs.main() + captured = capsys.readouterr() + assert captured.out == "a\nb\n" + assert "without shipping" in captured.err + post.assert_not_called() + + def test_full_run_ships_stream_and_exits_on_eof( + self, monkeypatch, capsys + ): + """End-to-end through main(): env wiring, the shipper thread, the + signal-handler registration, and the EOF final flush.""" + monkeypatch.setenv("ALLOCATOR_URL", "http://alloc/") + monkeypatch.setenv("CLIENT_SECRET", "sekrit") + monkeypatch.setenv("VM_NAME", "runai-client-2") + monkeypatch.setattr("sys.stdin", io.StringIO("one\ntwo\n")) + # Recorder instead of the real signal.signal: pytest owns the + # process's SIGINT handling, and clobbering it would outlive + # this test. + registered = {} + monkeypatch.setattr( + ship_logs.signal, + "signal", + lambda sig, handler: registered.__setitem__(sig, handler), + ) + + resp = MagicMock(status_code=200) + with patch.object(ship_logs.requests, "post", return_value=resp) as post: + ship_logs.main() + + # Passthrough reached stdout; the startup line went to stderr. + captured = capsys.readouterr() + assert captured.out == "one\ntwo\n" + assert "/api/vm-logs/runai-client-2" in captured.err + + # EOF triggered the final flush: both lines, stamped, one batch, + # bearer-authenticated, sanitized URL (no double slash). + assert post.call_count == 1 + assert post.call_args.args == ( + "http://alloc/api/vm-logs/runai-client-2", + ) + messages = post.call_args.kwargs["json"]["messages"] + assert [m.split(" ", 1)[1] for m in messages] == ["one", "two"] + headers = post.call_args.kwargs["headers"] + assert headers["Authorization"] == "Bearer sekrit" + + # Both stop signals were wired to the same handler, and invoking + # it is safe after shutdown (events set on an already-done loop). + assert set(registered) == {signal.SIGTERM, signal.SIGINT} + registered[signal.SIGTERM](signal.SIGTERM, None)