Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions docs/cli/byo-clients.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`.

Expand Down
1 change: 0 additions & 1 deletion packages/cli/pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,6 @@ dependencies = [
"rich>=13.0",
"pyyaml>=6.0",
"boto3>=1.35",
"psutil>=5.9",
]

[project.optional-dependencies]
Expand Down
92 changes: 23 additions & 69 deletions packages/cli/src/lablink_cli/commands/doctor.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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)
Expand Down Expand Up @@ -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


Expand All @@ -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(
Expand Down
170 changes: 26 additions & 144 deletions packages/cli/src/lablink_cli/commands/register.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand Down Expand Up @@ -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)."
)
Expand All @@ -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(
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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]}) …"
)
Expand All @@ -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)
Loading
Loading