8 Commits
Author SHA1 Message Date
nyra-lab 6029bfed64 open a token-gated front door on the WireGuard address, and only there
The console was loopback-only, which is right for the desktop GUI and the TUI but leaves
the phone with nothing to talk to. --host now takes a comma-separated list and
--token-file names the secret:

- loopback needs no token (it is the machine itself)
- any non-loopback listener requires one, and the server REFUSES TO START if the token
  file is missing, so a front door cannot be opened by accident
- the extra listener retries until its interface exists - wg0 can come up after this
  service does, and a boot race must not look like an auth failure from the phone
- /api/state reports "listeners", so which doors are actually open is visible
- new /api/ping for clients that just want to test the connection

Nine new cases: the token decision table (loopback, ipv6, missing/wrong/right bearer,
query token), token-file reading, and the refusal to open an unauthenticated door.
2026-09-27 16:17:19 +00:00
nyra-lab cb4d250eb2 audit the HTTP reclaim too - it frees VRAM like any command does
The reclaim recipe's HTTP branch called urllib directly and left no row in
audit.jsonl, so the one action that actually frees the card was invisible in the
trail. Runner.note() records side effects that have no subprocess, and the failure
message now goes through friendly_error like every other status line.
2026-09-27 15:27:12 +00:00
nyra-lab 86055b2e69 GUI: stop letting the poll eat button feedback, and shrink the logo
Every button reported its outcome ("job ... accepted", or the error) by writing the
header host line - which render() rewrites on every poll, so the message vanished
inside a second and the buttons looked dead. Notices now have their own element and
linger for 12s.

Also: the logo ships at 160x160 (20 KB) instead of the raw 1024x1024 (523 KB), and a
.gitignore keeps __pycache__/.venv out of the checkout.
2026-09-27 15:26:13 +00:00
nyra-lab eb217ea3e1 make the TUI actually run, and stop leaking developer text into status lines
The TUI defined refresh(), which collides with Textual Widget.refresh(*, repaint).
App.__init__ died on construction, so the terminal front end had never once run -
the lab had no textual installed, so nothing exercised that path. Renamed it to
reload_state, added action_reload for the r binding, and confirmed live frames on
Utumno against textual 8.2.8.

Also in this pass:
- health failures read "connection refused" / "timed out after 3s" instead of
  <urlopen error [Errno 111] Connection refused> (probes.friendly_error)
- the GUI header crops the logo into a rounded badge and gains a favicon
- the process table shows the memory not attributable to any process, so the rows
  reconcile with the card total
2026-09-27 15:24:47 +00:00
nyra-lab 5ad1e02168 switching must not be blocked by a process we did not start
ComfyUI is normally started by the Pictures page, so the cockpit correctly sees it
as external and refuses to kill it - which also blocked every switch. Implicit
stops (the card switch) are now best-effort: report it, leave it running, and let
the reclaim + VRAM wait decide. Explicit stops stay strict and fail loudly.
Three new cases cover it, plus the protected-stop wall and a portable fixture.
2026-09-27 15:18:48 +00:00
nyra-lab 8af2e04866 tests: do not assume the dev VM has no GPU - accept the VRAM verdict (unavailable/wait/below), found by running the suite on the real host 2026-09-27 15:16:33 +00:00
nyra-lab ff11e0189f gpu-cockpit: registry-driven GPU console (server, GUI, TUI, tests)
Built by Codex (nyra-lab, gpt-6-luna, effort high) from the brief in
runs/gpu-cockpit-v1/prompt.md; nyra reviewed the whole tree and fixed three
things before this commit:

- unit probe now asks systemd for LoadState, so a typo'd unit reads
  unknown instead of down
- a protected pipeline can no longer be stopped by a mis-click: explicit
  stop/restart of a live protected pipeline returns 409 and needs force
  (the GUI asks for a deliberate "Force stop"), plus two tests for it
- read-only probes (systemctl show, journalctl, nvidia-smi) still run under
  --dry-run, so dry-run shows real state while changing nothing

Tests: 24 cases pass in the dev VM (no GPU, stub systemd) and on NixOS.
2026-09-27 15:16:20 +00:00
nyra-lab 152db1098e AGENTS.md: repo conventions, commands, boundaries, definition of done 2026-09-27 14:50:08 +00:00
22 changed files with 1269 additions and 1 deletions
+3
View File
@@ -0,0 +1,3 @@
__pycache__/
*.pyc
.venv/
+54
View File
@@ -0,0 +1,54 @@
# gpu-cockpit — agent notes
A registry-driven console for a **single shared inference GPU**. One server owns state and command
execution; the web GUI and the terminal TUI are thin clients of its HTTP API. Nothing in the code
knows about any specific service — that all lives in the TOML registry.
## Layout
- `cockpit/` — the package: registry, runner, GPU probe, unit/proc control, health, state,
arbitration, HTTP server, and the inline GUI under `cockpit/gui/`. **stdlib only.**
- `tui/cockpit.py` — terminal client. The only place `textual` may be imported.
- `tests/`, `scripts/selftest.sh` — fixture-driven; must pass with no GPU and no real systemd.
- `registry.example.toml` — documentation by example. Real registries are deployment data.
## Commands
```bash
python3 -m cockpit --registry registry.example.toml --check # validate, no server
python3 -m cockpit --registry registry.example.toml --port 8770
python3 -m cockpit --registry registry.example.toml --dry-run # executes nothing, reports intent
bash scripts/selftest.sh # the acceptance suite
python3 -m compileall -q cockpit tui
python3 tui/cockpit.py --print # no TTY, no textual needed
```
Run `bash scripts/selftest.sh` before claiming anything works. It is the contract, not a smoke test.
## Boundaries
**Always do**
- Route every executed command through the single runner, so the audit log stays complete.
- Treat a health probe as the only evidence a pipeline is up; a `0` exit from a start command is not.
- Bound every wait, and report a timeout as a failure with the real observed state.
- Bind `127.0.0.1` unless `--host` is passed explicitly.
- Keep registry problems as validation errors/warnings shown in the UI, never a crash at boot.
**Ask first**
- Adding any dependency, including to `tui/`.
- Changing the registry schema or an existing field's meaning.
- Renaming an API route or a key inside `/api/state` (both front ends depend on them).
**Never do**
- `shell=True`, or building a command from an HTTP request's contents. The API accepts pipeline
ids, action names and option ids, and validates each against the registry.
- Stop a `protected = true` pipeline implicitly, or infer it from a `0` exit code.
- Write to the registry file from the API. Selection state is the most the server may write.
- Edit files outside this repository, or assume a deployment host's paths exist.
## Definition of done
1. `bash scripts/selftest.sh` exits 0 and prints one `PASS:` line per case.
2. `python3 -m compileall -q cockpit tui` is silent.
3. `python3 -m cockpit --registry registry.example.toml --check` reports `0 errors, 0 warnings`.
4. The README documents how to add a new pipeline without touching code — and it is true.
+112 -1
View File
@@ -1,3 +1,114 @@
# gpu-cockpit
GPU + pipeline cockpit: registry-driven VRAM/state console with TUI and web GUI
`gpu-cockpit` is a registry-driven console for one shared inference GPU. A single Python server
owns host probes, pipeline state, process/unit control, action jobs and the audit trail. The browser
GUI and terminal TUI are HTTP clients; service-specific details stay in TOML.
## Run
Requires Python 3.11+ (stdlib only for the server and GUI). Copy
[`registry.example.toml`](registry.example.toml) to `~/.config/gpu-cockpit/registry.toml`, then:
```sh
python3 -m cockpit --registry ~/.config/gpu-cockpit/registry.toml --check
python3 -m cockpit --registry ~/.config/gpu-cockpit/registry.toml --port 8770
```
The server binds to `127.0.0.1` by default. Set `--host` to bind elsewhere deliberately. State,
logs, selections and audit records go in `~/.local/state/gpu-cockpit` by default. `COCKPIT_REGISTRY`
can provide the default registry path. `--dry-run` serves the API and performs read-only health
checks while recording intended commands without starting or stopping processes.
Open `http://127.0.0.1:8770/` for the responsive web GUI. The TUI uses the same API:
```sh
python3 tui/cockpit.py --url http://127.0.0.1:8770
python3 tui/cockpit.py --print --url http://127.0.0.1:8770
```
Install `textual` to use the interactive TUI (`pip install textual`). `--print` needs no optional
dependency or terminal.
## Add a pipeline
Add another `[[pipeline]]` block in the registry. No Python change is required. Each pipeline needs
one or more components. Components start in listed order and stop in reverse order. A proc component
is spawned and supervised by the server; a unit component is operated through `systemctl`.
`health` can use `http`, `tcp`, `command` or `none`.
For example, a CPU-only local service can run alongside a card owner:
```toml
[[pipeline]]
id = "cpu-llm"
label = "CPU language model"
description = "A local inference endpoint that does not claim the GPU"
requires_card = false
vram_mib = 0
serving = "Small instruct model"
[[pipeline.component]]
kind = "proc"
id = "server"
label = "Inference server"
argv = ["/opt/local-llm/bin/server", "--port", "8090", "{model}"]
cwd = "/opt/local-llm"
env = { OMP_NUM_THREADS = "8" }
ready_timeout_s = 30
stop_timeout_s = 10
health = { kind = "tcp", host = "127.0.0.1", port = 8090, timeout_s = 2 }
[[pipeline.option]]
id = "model"
label = "Model"
kind = "argv"
default = "small"
[[pipeline.option.choice]]
id = "small"
label = "Small model"
vars = { model = "/models/small.gguf" }
```
Placeholders are substituted within argv elements using values from `argv` options. An empty value
drops the entire argv element. Every placeholder must be declared by an option, and every choice
must provide each variable used by that option. Options can also be `note` kind for display-only
choices. Run `--check` before restarting the server; validation problems remain visible in the API
and GUI.
## Safety and operation
- A pipeline with `requires_card = true` owns the exclusive card group. Starting another such
pipeline stops active peers, in reverse component order, before reclaiming and starting it.
- `protected = true` blocks implicit stops. A user must make a deliberate forced switch.
- Starts count as successful only after their configured health probe passes. Timeouts fail the job
and include the recent pipeline log tail. GPU waits are bounded and report unavailable probes or
timeout conditions honestly.
- Only one action job runs globally at a time. Jobs expose progress in `/api/jobs/<id>` and state.
- Proc pidfiles include Linux process start identity; stale files and reused PIDs do not imply an
owned process. An externally started healthy service is shown as external and cannot be killed.
- Executed commands and dry-run intent are recorded under the state directory. The API never edits
the registry; option selection is persisted in `selected.toml`.
## API
| Route | Purpose |
|---|---|
| `GET /` | Web GUI |
| `GET /static/<file>` | Allowlisted GUI assets |
| `GET /api/state` | Full card, component, pipeline, job and validation snapshot |
| `GET /api/events` | Server-sent state events and heartbeat (`?once=1` for one event) |
| `POST /api/pipelines/<id>/action` | Start, stop or restart; returns a job id |
| `POST /api/card/free` | Run the registry reclaim recipe as a job |
| `GET /api/jobs/<id>` | Job state, progress steps and failure detail |
| `GET /api/logs/<id>?n=200` | Tail proc log or unit journal |
| `GET /api/audit?n=100` | Recent command audit records |
| `GET /api/registry` | Parsed registry and validation report |
| `POST /api/pipelines/<id>/options` | Validate and persist option selection |
## Development checks
```sh
python3 -m compileall -q cockpit tui
python3 -m cockpit --registry registry.example.toml --check
bash scripts/selftest.sh
```
+2
View File
@@ -0,0 +1,2 @@
"""Registry-driven single-GPU pipeline console."""
__version__ = "0.1.0"
+4
View File
@@ -0,0 +1,4 @@
from .server import main
if __name__ == "__main__":
raise SystemExit(main())
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
+28
View File
@@ -0,0 +1,28 @@
<!doctype html><html lang="en"><head><meta charset="utf-8"><meta name="viewport" content="width=device-width,initial-scale=1"><title>GPU Cockpit</title><link rel="icon" href="/static/logo.png">
<style>
:root{color-scheme:dark;--bg:#0c1118;--panel:#141d27;--line:#293747;--text:#e7edf4;--muted:#91a2b5;--green:#47d7a0;--amber:#f3b95f;--red:#ff6c72;--blue:#71b7ff}*{box-sizing:border-box}body{margin:0;background:radial-gradient(ellipse at 40% -10%,#1a2a3c,var(--bg) 55%);color:var(--text);font:15px/1.45 system-ui,sans-serif}header{padding:18px max(20px,calc((100vw - 1440px)/2));display:flex;gap:14px;align-items:center;border-bottom:1px solid var(--line);position:sticky;top:0;background:#0c1118eF;z-index:4}header .brand{width:38px;height:38px;flex:none;border-radius:11px;overflow:hidden;border:1px solid #2b3a4d;background:#14161f;box-shadow:0 0 12px #71b7ff22}header .brand img{width:180%;height:180%;margin:-40% 0 0 -40%;display:block}.unattributed td{color:var(--muted);font-style:italic}h1{font-size:19px;margin:0}header small,.muted{color:var(--muted)}.right{margin-left:auto;display:flex;align-items:center;gap:14px}.wrap{max-width:1440px;margin:auto;padding:22px}.top{display:grid;grid-template-columns:minmax(300px,1fr) minmax(0,2fr);gap:18px}.panel,.card{background:linear-gradient(145deg,#17222e,#121a23);border:1px solid var(--line);border-radius:15px;padding:18px;box-shadow:0 12px 32px #0002}.panel h2{margin:0 0 14px;font-size:16px}.gauge{display:flex;gap:18px;align-items:center}.ring{width:120px;height:120px;flex:none}.ring text{fill:var(--text);font-size:17px;font-weight:700}.metrics{display:grid;grid-template-columns:repeat(3,1fr);gap:10px;flex:1}.metric{border-left:1px solid var(--line);padding-left:12px}.metric b{display:block;font-size:20px}.metric span{font-size:12px;color:var(--muted)}.spark{width:100%;height:40px;margin-top:10px}.process{width:100%;border-collapse:collapse;font-size:13px}.process td,.process th{text-align:left;padding:7px;border-bottom:1px solid #29374788}.external{color:var(--amber)}.section{display:flex;align-items:end;justify-content:space-between;margin:25px 0 12px}.section h2{margin:0}.grid{display:grid;grid-template-columns:repeat(3,minmax(0,1fr));gap:14px}.card h3{margin:0;font-size:17px}.title{display:flex;align-items:center;gap:9px}.dot{width:10px;height:10px;border-radius:50%;background:#778493}.dot.up{background:var(--green);box-shadow:0 0 10px #47d7a088}.dot.partial,.dot.starting,.dot.stopping{background:var(--amber)}.dot.error{background:var(--red)}.desc{min-height:42px;color:var(--muted);font-size:13px;margin:7px 0 12px}.tags{display:flex;gap:7px;flex-wrap:wrap;color:var(--muted);font-size:12px}.serving{margin:10px 0 4px;font-weight:600}.note{color:var(--amber);font-size:12px;min-height:18px}.components{border-top:1px solid var(--line);margin-top:12px;padding-top:9px}.comp{display:flex;justify-content:space-between;gap:8px;font-size:12px;padding:3px 0}.selects{margin-top:10px;display:grid;gap:8px}.choice{display:grid;grid-template-columns:auto 1fr;align-items:center;gap:8px;font-size:12px}select,button{font:inherit;color:var(--text);background:#202e3d;border:1px solid #3a4e63;border-radius:8px;padding:8px}select{min-width:0;width:100%}button{cursor:pointer}button:hover{border-color:var(--blue)}button:focus-visible,a:focus-visible,select:focus-visible{outline:2px solid var(--blue);outline-offset:2px}.primary{background:#216b58;border-color:#399b7e;font-weight:650}.danger{background:#45272c;border-color:#704047}.buttons{display:flex;gap:7px;flex-wrap:wrap;margin-top:13px}.buttons button{font-size:13px}.job{color:var(--blue);font-size:12px;min-height:18px;margin-top:8px}.log{margin-top:8px;color:var(--muted);font-size:12px}details summary{cursor:pointer;color:var(--muted)}.banner{display:none;padding:14px 17px;margin-top:16px;border:1px solid #a77e38;background:#382d1c;border-radius:12px}.banner.show{display:flex;align-items:center;gap:12px;flex-wrap:wrap}.banner span{flex:1}.activity{margin-top:20px}.activity table{width:100%;font-size:12px}#activityContent{max-height:320px;overflow:auto}.status{display:inline-flex;align-items:center;gap:7px}.live{color:var(--green)}.offline{color:var(--red)}.notice{color:var(--amber);font-size:13px;max-width:34ch;overflow:hidden;text-overflow:ellipsis;white-space:nowrap}.notice.bad{color:var(--red)}a{color:var(--blue)}@media(max-width:1000px){.grid{grid-template-columns:repeat(2,minmax(0,1fr))}.top{grid-template-columns:1fr}}@media(max-width:600px){header{padding:12px 14px}.wrap{padding:14px}.grid{grid-template-columns:1fr}.right{gap:7px}header small{display:none}.gauge{align-items:flex-start}.ring{width:95px;height:95px}.metrics{grid-template-columns:repeat(2,1fr)}.metric b{font-size:17px}}
</style></head><body><header><span class="brand"><img src="/static/logo.png" onerror="this.parentNode.hidden=true" alt=""></span><div><h1>GPU Cockpit</h1><small id="host">Connecting…</small></div><div class="right"><span id="notice" class="notice"></span><span id="connection" class="status offline">● offline</span><span id="stream" class="muted">polling</span><button class="primary" onclick="freeCard()">Free the card</button></div></header>
<main class="wrap"><section class="top"><div class="panel"><h2>Card memory</h2><div class="gauge"><svg class="ring" viewBox="0 0 120 120" aria-label="VRAM usage"><circle cx="60" cy="60" r="48" fill="none" stroke="#293747" stroke-width="10"/><circle id="arc" cx="60" cy="60" r="48" fill="none" stroke="#47d7a0" stroke-width="10" stroke-linecap="round" transform="rotate(-90 60 60)" stroke-dasharray="301.6" stroke-dashoffset="301.6"/><text x="60" y="65" text-anchor="middle" id="pct">—</text></svg><div class="metrics"><div class="metric"><b id="used">—</b><span>used / total MiB</span></div><div class="metric"><b id="util">—</b><span>GPU utilization</span></div><div class="metric"><b id="temp">—</b><span>temperature</span></div></div></div><small class="muted" id="floor"></small><svg class="spark" viewBox="0 0 300 40" preserveAspectRatio="none"><polyline id="spark" fill="none" stroke="#71b7ff" stroke-width="2" points=""/></svg></div><div class="panel"><h2>GPU processes</h2><div id="processes" class="muted">No GPU probe</div></div></section>
<div id="banner" class="banner"><span id="bannerText"></span><button id="switchBtn" class="primary">Switch to this</button><button id="forceBtn" class="danger" hidden>Force switch</button></div><div class="section"><h2>Pipelines</h2><span class="muted" id="regline"></span></div><section id="pipes" class="grid"></section>
<details class="panel activity"><summary>Activity · command audit and registry report</summary><div id="activityContent"><h3>Registry validation</h3><div id="validation"></div><h3>Command audit</h3><div id="audit"></div></div></details></main>
<script>
const $=id=>document.getElementById(id);let state=null,busy=false,evt=null;
async function api(path,opts){let r=await fetch(path,opts);let j=await r.json().catch(()=>({}));if(!r.ok)throw Error(j.error||r.statusText);return j}
function esc(x){return String(x??'').replace(/[&<>"']/g,c=>({'&':'&amp;','<':'&lt;','>':'&gt;','"':'&quot;',"'":'&#39;'}[c]))}
function render(s){state=s;$('host').textContent=`${s.host} · ${s.now}`;$('connection').className='status live';$('connection').textContent='● live';let c=s.card;if(!c.available){$('used').textContent='No GPU probe';$('util').textContent='—';$('temp').textContent='—';$('pct').textContent='—';$('arc').style.strokeDashoffset=301.6;$('floor').textContent=c.detail||'nvidia-smi unavailable';$('processes').textContent='No GPU probe';}else{$('used').textContent=`${c.used_mib} / ${c.total_mib}`;$('util').textContent=`${c.util_pct}%`;$('temp').textContent=`${c.temp_c}°C`;$('pct').textContent=`${c.pct_of_usable}%`;$('arc').style.strokeDashoffset=301.6*(1-c.pct_of_usable/100);$('floor').textContent=`Idle floor ${c.idle_floor_mib} MiB · gauge uses ${c.usable_mib} MiB usable`;$('processes').innerHTML=c.processes.length?`<table class="process"><thead><tr><th>PID</th><th>Process</th><th>MiB</th><th>Owner</th></tr></thead><tbody>${c.processes.map(p=>`<tr class="${p.pipeline?'':'external'}"><td>${p.pid}</td><td>${esc(p.name)}</td><td>${p.used_mib}</td><td>${p.pipeline?esc(p.pipeline):'External holder'}</td></tr>`).join('')+((c.used_mib-c.processes.reduce((a,p)=>a+p.used_mib,0))>150?`<tr class="unattributed"><td>—</td><td>desktop &amp; graphics surfaces</td><td>${c.used_mib-c.processes.reduce((a,p)=>a+p.used_mib,0)}</td><td>—</td></tr>`:'')}</tbody></table>`:'No GPU processes reported';let vals=c.history.map(x=>x.used_mib),max=Math.max(1,...vals);$('spark').setAttribute('points',vals.map((v,i)=>`${i*300/Math.max(1,vals.length-1)},${38-v/max*34}`).join(' '));}
$('regline').textContent=`${s.registry.pipeline_count} pipelines · ${s.registry.errors.length} errors · ${s.registry.warnings.length} warnings`;
$('pipes').innerHTML=s.pipelines.map(p=>`<article class="card"><div class="title"><span class="dot ${p.state}"></span><h3>${esc(p.label)}</h3><span class="right muted">${esc(p.state)}</span></div><div class="desc">${esc(p.description)}</div><div class="tags"><span>${p.vram_mib??'—'} MiB expected</span><span>${p.requires_card?'shared card':'CPU / independent'}</span>${p.protected?'<span>protected</span>':''}</div><div class="serving">${esc(p.serving||'')}</div><div class="note">${esc(p.note||'')}</div><div class="components">${p.components.map(c=>`<div class="comp"><span>${esc(c.label)} · ${esc(c.state)}${c.port?` · :${c.port}`:''}</span><span>${esc(c.health?.detail||c.detail||'')}</span></div>`).join('')}</div><div class="selects">${p.options.map(o=>`<label class="choice">${esc(o.label)}<select aria-label="${esc(o.label)}" onchange="choose('${p.id}','${o.id}',this.value)">${o.choices.map(ch=>`<option value="${esc(ch.id)}" ${ch.id===o.selected?'selected':''}>${esc(ch.label)}${ch.verified?' · verified':''}</option>`).join('')}</select>${o.choices.find(x=>x.id===o.selected)?.note?`<small class="muted">${esc(o.choices.find(x=>x.id===o.selected).note)}</small>`:''}</label>`).join('')}</div><div class="buttons"><button class="primary" onclick="action('${p.id}','start')">Start</button><button onclick="stopPipeline('${p.id}')">Stop</button><button onclick="action('${p.id}','restart')">Restart</button><button onclick="switchTo('${p.id}')">Switch to this</button>${p.open_url?`<a href="${esc(p.open_url)}" target="_blank" rel="noopener">Open ↗</a>`:''}</div><div class="job">${p.job?`${esc(p.job.id)} · ${esc(p.job.step)}`:''}</div><details class="log"><summary>Log tail</summary><pre id="log-${p.id}">Open to load</pre></details></article>`).join('');
let active=s.pipelines.find(p=>p.job);busy=!!s.job;$('pipes').querySelectorAll('button,select').forEach(x=>x.disabled=busy);if(active){$('banner').classList.add('show');$('bannerText').textContent=`${active.label} is changing: ${active.job.step}`;$('switchBtn').hidden=true;$('forceBtn').hidden=true;}else{$('banner').classList.remove('show');}
$('validation').innerHTML=[...s.registry.errors.map(x=>`<div class="external">Error: ${esc(x)}</div>`),...s.registry.warnings.map(x=>`<div class="muted">Warning: ${esc(x)}</div>`)].join('')||'<span class="live">No registry problems</span>';api('/api/audit?n=30').then(a=>$('audit').innerHTML=`<table class="process"><tbody>${a.slice().reverse().map(x=>`<tr><td>${esc(x.ts)}</td><td>${esc(x.pipeline)}</td><td>${esc(x.argv.join(' '))}</td><td>${x.rc}</td><td>${x.duration_ms} ms</td></tr>`).join('')}</tbody></table>`).catch(()=>{});
}
async function refresh(){try{render(await api('/api/state'))}catch(e){$('connection').className='status offline';$('connection').textContent='● offline';$('host').textContent='Waiting for server';}}
let noticeTimer=null;
function notice(msg,bad){let n=$('notice');n.textContent=msg;n.className=bad?'notice bad':'notice';clearTimeout(noticeTimer);noticeTimer=setTimeout(()=>{n.textContent=''},12000)}
async function action(id,a,force=false){if(busy)return;let options={};state.pipelines.find(p=>p.id===id)?.options.forEach(o=>options[o.id]=o.selected);try{let j=await api(`/api/pipelines/${id}/action`,{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify({action:a,options,force})});notice(`job ${j.job_id} accepted`);busy=true;refresh()}catch(e){notice(e.message,true);}}
async function choose(pid,oid,value){let p=state.pipelines.find(p=>p.id===pid),options={};p.options.forEach(o=>options[o.id]=o.selected);options[oid]=value;try{await api(`/api/pipelines/${pid}/options`,{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify({options})});refresh()}catch(e){notice(e.message,true)}}
async function switchTo(id){let p=state.pipelines.find(x=>x.id===id),owners=state.pipelines.filter(x=>x.id!==id&&x.requires_card&&x.state!=='down');let blocked=owners.find(x=>x.protected);if(blocked){$('banner').classList.add('show');$('bannerText').textContent=`${blocked.label} is protected and would keep ${p.label} from using the card. Force stops it and loses its current work.`;$('switchBtn').hidden=true;$('forceBtn').hidden=false;$('forceBtn').onclick=()=>action(id,'start',true);return;}if(owners.length){$('banner').classList.add('show');$('bannerText').textContent=`Starting ${p.label} will stop ${owners.map(x=>x.label).join(', ')} and release the card.`;$('switchBtn').hidden=false;$('switchBtn').textContent='Stop them and switch';$('switchBtn').onclick=()=>action(id,'start');}else action(id,'start');}
async function stopPipeline(id){let p=state.pipelines.find(x=>x.id===id);if(p&&p.protected&&p.state!=='down'){$('banner').classList.add('show');$('bannerText').textContent=`${p.label} is protected \u2014 stopping it now would lose work in progress.`;$('switchBtn').hidden=true;$('forceBtn').hidden=false;$('forceBtn').textContent='Force stop';$('forceBtn').onclick=()=>action(id,'stop',true);return;}action(id,'stop');}
async function freeCard(){try{await api('/api/card/free',{method:'POST'});busy=true;refresh()}catch(e){notice(e.message,true)}}
document.addEventListener('click',e=>{let d=e.target.closest('details.log');if(d&&d.open){let id=d.querySelector('pre').id.slice(4);api(`/api/logs/${id}?n=40`).then(x=>$(d.querySelector('pre').id).textContent=x.join('\n')).catch(()=>{})}});
function connect(){if(!('EventSource'in window))return;$('stream').textContent='live stream';evt=new EventSource('/api/events');evt.addEventListener('state',e=>{try{render(JSON.parse(e.data))}catch{}});evt.onerror=()=>{$('stream').textContent='polling fallback';evt.close();evt=null;setTimeout(connect,5000)}}refresh();connect();setInterval(()=>{if(!evt||evt.readyState!==1)refresh()},3000);
</script></body></html>
Binary file not shown.

After

Width:  |  Height:  |  Size: 19 KiB

+63
View File
@@ -0,0 +1,63 @@
"""Read-only host, process and health probes."""
import json, os, re, socket, subprocess, time, urllib.request, urllib.error
def friendly_error(exc,timeout=None):
"""A status line has no room for '<urlopen error [Errno 111] Connection refused>'.
Map the usual socket failures to something a person reads at a glance."""
reason=getattr(exc,"reason",None); errno=getattr(reason,"errno",None) or getattr(exc,"errno",None)
if errno==111: return "connection refused"
if errno==113: return "no route to host"
if errno==98: return "address already in use"
if errno==110 or isinstance(exc,(TimeoutError,socket.timeout)) or "timed out" in str(exc).lower():
return f"timed out after {timeout}s" if timeout else "timed out"
if isinstance(exc,urllib.error.URLError): exc=reason or exc
return str(exc).strip(" <>")[:120]
def health(spec, runner=None, pipeline=None):
if not spec or spec.get("kind")=="none": return {"ok":True,"detail":"no health check","latency_ms":0}
start=time.monotonic(); kind=spec.get("kind"); timeout=min(30,max(.1,float(spec.get("timeout_s",3))))
try:
if kind=="http":
req=urllib.request.Request(spec["url"],method="GET")
with urllib.request.urlopen(req,timeout=timeout) as r: ok=200<=r.status<400; detail=f"HTTP {r.status}"
elif kind=="tcp":
with socket.create_connection((spec["host"],int(spec["port"])),timeout=timeout): pass
ok=True; detail="TCP connected"
elif kind=="command":
if runner:
rc,out,_=runner.run(spec["argv"],pipeline,"health",timeout=timeout,readonly=True)
else:
return {"ok":False,"detail":"command health probe needs runner","latency_ms":0}
ok=rc==0; detail=f"exit {rc}"
else: return {"ok":False,"detail":f"unknown health kind: {kind}","latency_ms":0}
return {"ok":ok,"detail":detail,"latency_ms":round((time.monotonic()-start)*1000)}
except Exception as exc: return {"ok":False,"detail":friendly_error(exc,timeout),"latency_ms":round((time.monotonic()-start)*1000)}
def gpu(runner=None):
try:
if runner:
rc,out,_=runner.run(["nvidia-smi","--query-gpu=name,memory.used,memory.total,utilization.gpu,temperature.gpu","--format=csv,noheader,nounits"],None,"gpu-probe",timeout=3,readonly=True)
else:
p=subprocess.run(["nvidia-smi","--query-gpu=name,memory.used,memory.total,utilization.gpu,temperature.gpu","--format=csv,noheader,nounits"],timeout=3,text=True,capture_output=True); rc,out=p.returncode,p.stdout if p.returncode==0 else p.stderr
if rc: raise RuntimeError(out.strip() or f"exit {rc}")
row=out.strip().splitlines()[0].split(",")
name,used,total,util,temp=[x.strip() for x in row]
procs=[]
if runner: prc,pout,_=runner.run(["nvidia-smi","--query-compute-apps=pid,process_name,used_memory","--format=csv,noheader,nounits"],None,"gpu-process-probe",timeout=3,readonly=True)
else:
pp=subprocess.run(["nvidia-smi","--query-compute-apps=pid,process_name,used_memory","--format=csv,noheader,nounits"],timeout=3,text=True,capture_output=True); prc,pout=pp.returncode,pp.stdout
if prc==0:
for line in pout.splitlines():
x=[z.strip() for z in line.split(",")]
if len(x)>=3:
try: procs.append({"pid":int(x[0]),"name":os.path.basename(x[1]),"used_mib":int(float(x[2])),"pipeline":None})
except ValueError: pass
return {"available":True,"name":name,"used_mib":int(float(used)),"total_mib":int(float(total)),"util_pct":int(float(util)),"temp_c":int(float(temp)),"processes":procs}
except Exception as exc: return {"available":False,"name":None,"used_mib":None,"total_mib":None,"util_pct":None,"temp_c":None,"processes":[],"detail":friendly_error(exc)}
def proc_start(pid):
try:
raw=open(f"/proc/{pid}/stat").read(); fields=raw[raw.rfind(")")+2:].split()
if fields[0] in ("Z","X"): return None
return fields[19]
except Exception: return None
+149
View File
@@ -0,0 +1,149 @@
"""Registry loading and validation. Invalid pipelines are retained for reporting."""
import re
import tomllib
from pathlib import Path
ID = re.compile(r"^[a-z0-9-]+$")
PLACEHOLDER = re.compile(r"\{([A-Za-z_][A-Za-z0-9_]*)\}")
def load(path):
errors, warnings = [], []
try:
with open(path, "rb") as f:
data = tomllib.load(f)
except Exception as exc:
return {"card": {}, "pipeline": []}, [f"registry: {exc}"], warnings
card = data.setdefault("card", {})
if not isinstance(card, dict):
errors.append("card must be a table")
card = data["card"] = {}
if "reclaim" in card and not isinstance(card["reclaim"], dict):
errors.append("card.reclaim must be a table")
card.pop("reclaim", None)
for key, default, minimum in (("total_mib", None, 1), ("idle_floor_mib", 0, 0), ("poll_seconds", 2, .1)):
if key not in card:
if default is not None: card[key] = default
continue
try:
value=float(card[key])
if value < minimum: raise ValueError()
card[key]=int(value) if key != "poll_seconds" else value
except (TypeError, ValueError):
errors.append(f"card.{key} must be a number >= {minimum}")
if default is not None: card[key]=default
else: card.pop(key,None)
rec=card.get("reclaim")
if rec and rec.get("kind") not in ("http", "command"):
errors.append("card.reclaim.kind must be http or command")
if rec:
try:
if int(rec.get("wait_timeout_s",45)) < 1: raise ValueError()
except (TypeError,ValueError): errors.append("card.reclaim.wait_timeout_s must be a positive integer")
if rec.get("kind")=="http" and not isinstance(rec.get("url"),str): errors.append("card.reclaim.url is required for http reclaim")
if rec.get("kind")=="command" and (not isinstance(rec.get("argv"),list) or not rec.get("argv")):
errors.append("card.reclaim.argv is required for command reclaim")
pipes = data.setdefault("pipeline", [])
if not isinstance(pipes, list):
errors.append("pipeline must be an array of tables")
pipes = []
data["pipeline"] = pipes
ids = set()
for i, p in enumerate(pipes):
prefix = f"pipeline[{i}]"
if not isinstance(p, dict):
errors.append(f"{prefix}: pipeline entry must be a table")
continue
pid = p.get("id")
if not isinstance(pid, str) or not ID.fullmatch(pid):
errors.append(f"{prefix}: invalid id")
elif pid in ids:
errors.append(f"{prefix}: duplicate id {pid}")
if isinstance(pid, str):
ids.add(pid)
comps = p.get("component", [])
if not isinstance(comps, list):
errors.append(f"{prefix} {pid}: component must be an array of tables")
comps = []
p["component"] = comps
if not comps:
errors.append(f"{prefix} {pid}: at least one component is required")
opts = p.get("option", [])
if not isinstance(opts, list):
errors.append(f"{prefix} {pid}: option must be an array of tables")
opts = []
p["option"] = opts
clean_opts = [o for o in opts if isinstance(o, dict) and isinstance(o.get("id"),str)]
if len(clean_opts) != len(opts):
errors.append(f"{prefix} {pid}: option entries must be tables with string ids")
p["option"] = clean_opts
for o in clean_opts:
if not isinstance(o.get("id"),str): errors.append(f"{prefix} {pid}: option id must be a string")
if o.get("kind") not in ("argv", "note"):
errors.append(f"{prefix} {pid}: option {o.get('id')} kind must be argv or note")
choices = o.get("choice", [])
if not isinstance(choices, list):
errors.append(f"{prefix} {pid}: option {o.get('id')} choices must be an array")
choices = []
choices = [ch for ch in choices if isinstance(ch, dict) and isinstance(ch.get("id"),str)]
o["choice"] = choices
if any(not isinstance(ch.get("id"),str) for ch in choices): errors.append(f"{prefix} {pid}: choice ids must be strings")
if o.get("kind") == "argv":
names = set()
for ch in choices:
vars_ = ch.get("vars", {})
if not isinstance(vars_, dict):
errors.append(f"{prefix} {pid}: choice {ch.get('id')} vars must be a table")
vars_ = ch["vars"] = {}
names.update(vars_)
for ch in choices:
missing = names - set(ch.get("vars", {}))
if missing:
errors.append(f"{prefix} {pid}: choice {ch.get('id')} missing vars: {', '.join(sorted(missing))}")
if o.get("default") not in [ch.get("id") for ch in choices]:
errors.append(f"{prefix} {pid}: option {o.get('id')} has invalid default")
compids = set()
for c in comps:
if not isinstance(c, dict):
errors.append(f"{prefix} {pid}: component entries must be tables")
continue
if c.get("kind") not in ("proc", "unit"):
errors.append(f"{prefix} {pid}: component kind must be proc or unit")
cid = c.get("id")
if isinstance(cid,str) and cid in compids:
errors.append(f"{prefix} {pid}: duplicate component id {cid}")
if isinstance(cid, str):
compids.add(cid)
if c.get("kind")=="unit" and not isinstance(c.get("unit"),str):
errors.append(f"{prefix} {pid}/{cid}: unit name is required")
if c.get("kind") == "proc":
argv = c.get("argv", [])
if not isinstance(argv, list) or not argv or any(not isinstance(x, str) for x in argv):
errors.append(f"{prefix} {pid}/{cid}: proc argv must be a non-empty string array")
argv = []
elif argv[0].startswith("/") and not Path(argv[0]).exists():
warnings.append(f"{pid}/{cid}: executable path does not exist: {argv[0]}")
providers = {name for o in clean_opts if o.get("kind") == "argv" for ch in o.get("choice", []) for name in ch.get("vars", {})}
for arg in argv:
for name in PLACEHOLDER.findall(arg):
if name not in providers:
errors.append(f"{prefix} {pid}: placeholder {{{name}}} has no option value")
if "env" in c and (not isinstance(c["env"], dict) or any(not isinstance(k, str) or not isinstance(v, str) for k,v in c["env"].items())):
errors.append(f"{prefix} {pid}/{cid}: env must map strings to strings")
if "health" in c and not isinstance(c["health"], dict):
errors.append(f"{prefix} {pid}/{cid}: health must be a table")
c["health"] = {"kind": "none"}
health=c.get("health",{"kind":"none"})
if isinstance(health,dict):
if health.get("kind") not in ("http","tcp","command","none"):
errors.append(f"{prefix} {pid}/{cid}: invalid health kind")
if health.get("kind")=="command" and (not isinstance(health.get("argv"),list) or not health.get("argv")):
errors.append(f"{prefix} {pid}/{cid}: command health requires argv")
cwd = c.get("cwd")
if cwd and not Path(cwd).exists():
warnings.append(f"{pid}/{cid}: configured cwd does not exist: {cwd}")
return data, errors, warnings
def problem_count(reg):
return len(reg.get("pipeline", []))
+66
View File
@@ -0,0 +1,66 @@
"""The sole subprocess boundary and command audit trail."""
import json, os, subprocess, threading, time
import signal
from datetime import datetime, timezone
def stamp(): return datetime.now(timezone.utc).isoformat(timespec="seconds")
class Runner:
def __init__(self, state_dir, dry_run=False, systemctl="systemctl", sudo="sudo"):
self.state_dir, self.dry_run, self.systemctl, self.sudo = state_dir, dry_run, systemctl, sudo
self.lock = threading.Lock(); self.audit_path = os.path.join(state_dir, "audit.jsonl")
os.makedirs(state_dir, exist_ok=True)
def run(self, argv, pipeline=None, action="probe", timeout=15, env=None, cwd=None, readonly=False):
start=time.monotonic(); rc=0; out=""
try:
if self.dry_run and not readonly:
out="dry-run: command not executed"
else:
p=subprocess.run(list(map(str,argv)), cwd=cwd, env=env, text=True, stdout=subprocess.PIPE,
stderr=subprocess.STDOUT, timeout=timeout, check=False)
rc=p.returncode; out=p.stdout or ""
except Exception as exc:
rc=127; out=f"{type(exc).__name__}: {exc}"
row={"ts":stamp(),"pipeline":pipeline,"action":action,"argv":list(argv),"rc":rc,
"duration_ms":round((time.monotonic()-start)*1000),"dry_run":self.dry_run}
with self.lock:
with open(self.audit_path,"a",encoding="utf-8") as f: f.write(json.dumps(row)+"\n")
return rc,out,row
def note(self,argv,action,pipeline=None,rc=0,duration_ms=0):
"""Audit a side effect that has no subprocess of its own - the HTTP reclaim
frees VRAM on the card like any command does, so the trail must carry it."""
row={"ts":stamp(),"pipeline":pipeline,"action":action,"argv":list(map(str,argv)),"rc":rc,
"duration_ms":round(duration_ms),"dry_run":self.dry_run}
with self.lock:
with open(self.audit_path,"a",encoding="utf-8") as f: f.write(json.dumps(row)+"\n")
return row
def audit(self,n=100):
try:
with open(self.audit_path,encoding="utf-8") as f: rows=[json.loads(x) for x in f if x.strip()]
return rows[-max(1,min(1000,n)):]
except OSError: return []
def spawn(self, argv, *, cwd=None, env=None, stdout=None, stderr=None, pipeline=None, action="start"):
"""Start an owned long-lived child; launch intent is audited by run()."""
start=time.monotonic()
try:
child=subprocess.Popen(list(map(str,argv)), cwd=cwd, env=env, stdin=subprocess.DEVNULL,
stdout=stdout, stderr=stderr, start_new_session=True)
except Exception as exc:
self.record_intent(argv,pipeline,action,cwd,rc=127,duration_ms=round((time.monotonic()-start)*1000),error=str(exc))
raise
self.record_intent(argv,pipeline,action,cwd,rc=0,duration_ms=round((time.monotonic()-start)*1000))
return child
def record_intent(self, argv, pipeline=None, action="start", cwd=None, *, rc=0, duration_ms=0, error=None):
row={"ts":stamp(),"pipeline":pipeline,"action":action,"argv":list(argv),"rc":rc,
"duration_ms":duration_ms,"dry_run":False,"cwd":cwd}
if error: row["error"]=error
with self.lock:
with open(self.audit_path,"a",encoding="utf-8") as f: f.write(json.dumps(row)+"\n")
def signal(self, pid, sig, pipeline=None, action="stop"):
start=time.monotonic(); rc=0; out=""
try: os.kill(pid,sig)
except ProcessLookupError: rc=1; out="process no longer exists"
row={"ts":stamp(),"pipeline":pipeline,"action":action,"argv":["signal",signal.Signals(sig).name,str(pid)],"rc":rc,"duration_ms":round((time.monotonic()-start)*1000),"dry_run":False}
with self.lock:
with open(self.audit_path,"a",encoding="utf-8") as f: f.write(json.dumps(row)+"\n")
return rc,out,row
+423
View File
@@ -0,0 +1,423 @@
"""HTTP API, asynchronous action jobs and snapshot assembly."""
import argparse, hmac, json, os, re, secrets, signal, socket, sys, threading, time, urllib.request
from collections import deque
from datetime import datetime, timezone
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from pathlib import Path
from urllib.parse import urlparse, parse_qs, unquote
from . import __version__
from .registry import load, PLACEHOLDER
from .runner import Runner, stamp
from .probes import health, gpu, proc_start, friendly_error
def iso(): return datetime.now(timezone.utc).isoformat(timespec="seconds")
def is_loopback(addr):
"""Loopback is the machine itself: the desktop GUI and the TUI. Everything else has
crossed a network boundary and has to prove it may talk to us."""
a=str(addr or "")
return a=="::1" or a.startswith("127.") or a.startswith("::ffff:127.")
def bearer(headers,query):
got=(headers.get("Authorization") or "").strip()
if got[:7].lower()=="bearer ": return got[7:].strip()
return (query or {}).get("token",[""])[0]
def token_matches(expected,headers,query,client_addr):
"""No token configured == loopback-only: main() refuses to open a front-door listener
without one, so this can never be the thing that leaves a port wide open."""
if not expected: return True
if is_loopback(client_addr): return True
return hmac.compare_digest(bearer(headers,query),expected)
def load_token(path):
if not path: return None
try: return (Path(path).read_text(encoding="utf-8").strip() or None)
except FileNotFoundError: return None
except OSError as exc:
print(f"cannot read token file {path}: {exc}",file=sys.stderr); raise SystemExit(1)
class App:
def __init__(self,args):
self.args=args; self.registry,self.errors,self.warnings=load(args.registry)
self.pipelines={p.get("id"):p for p in self.registry["pipeline"] if isinstance(p,dict) and isinstance(p.get("id"),str) and re.fullmatch(r"[a-z0-9-]+",p.get("id"))}
self.runner=Runner(args.state_dir,args.dry_run,args.systemctl,args.sudo)
self.lock=threading.Lock(); self.action_lock=threading.Lock(); self.job_guard=threading.Lock(); self.jobs={}; self.active_job=None
self.history=deque(maxlen=90); self.last_gpu_sample=0; self.last_change={}; self.selected=self.read_selected()
def read_selected(self):
try:
import tomllib
with open(os.path.join(self.args.state_dir,"selected.toml"),"rb") as f: return tomllib.load(f)
except Exception: return {}
def save_selected(self):
path=os.path.join(self.args.state_dir,"selected.toml"); lines=[]
for pid,opts in self.selected.items():
if not isinstance(opts,dict): continue
lines.append(f'[{json.dumps(pid)}]')
lines.extend(f'{k} = {json.dumps(v)}' for k,v in opts.items())
Path(path).write_text("\n".join(lines)+("\n" if lines else ""),encoding="utf-8")
def component(self,p,c):
result={"id":c.get("id"),"kind":c.get("kind"),"label":c.get("label",c.get("id")),"state":"unknown","detail":"state unavailable","health":None,"port":None,"unit":None}
spec=c.get("health",{}); result["port"]=spec.get("port") if spec.get("kind")=="tcp" else None
if c.get("kind")=="unit":
unit=c.get("unit",""); result["unit"]=unit
rc,out,_=self.runner.run([self.args.systemctl,"show",unit,"--property=LoadState,ActiveState,SubState,ExecMainPID","--no-pager"],p.get("id"),"probe",timeout=4,readonly=True)
vals=dict(x.split("=",1) for x in out.splitlines() if "=" in x)
if rc or not vals: result["state"]="unknown"; result["detail"]=out.strip() or "systemd unit not found"; return result
if vals.get("LoadState")=="not-found": result["state"]="unknown"; result["detail"]=f"systemd unit not found: {unit}"; return result
if vals.get("ActiveState")=="active":
h=health(spec,self.runner,p.get("id")); result.update(state="up" if h["ok"] else "error",detail=f"pid {vals.get('ExecMainPID','?')}",health=h)
else: result.update(state="down",detail=vals.get("ActiveState","inactive"))
return result
pidfile=Path(self.args.state_dir)/"pids"/f"{p.get('id')}-{c.get('id')}.json"
recorded=None
try: recorded=json.loads(pidfile.read_text())
except Exception: pass
if recorded:
pid=recorded.get("pid"); start=proc_start(pid) if isinstance(pid,int) else None
if start is not None and str(start)==str(recorded.get("start")):
h=health(spec,self.runner,p.get("id")); result.update(state="up" if h["ok"] else "error",detail=f"pid {pid}",health=h); return result
h=health(spec,self.runner,p.get("id"))
if spec.get("kind")=="none": h={"ok":False,"detail":"no health check","latency_ms":0}
if h["ok"]: result.update(state="up",detail="external",health=h)
elif recorded: result.update(state="down",detail="stale pidfile; process identity changed",health=h)
else: result.update(state="down",detail="not running",health=h)
return result
def snapshot(self):
g=gpu(self.runner); cardcfg=self.registry.get("card",{}); total=cardcfg.get("total_mib") or g.get("total_mib"); floor=int(cardcfg.get("idle_floor_mib",0)); usable=max(0,(total or 0)-floor)
if g.get("available") and time.monotonic()-self.last_gpu_sample>=max(.1,float(cardcfg.get("poll_seconds",2))):
self.history.append({"t":iso(),"used_mib":g["used_mib"]})
self.last_gpu_sample=time.monotonic()
g.update(total_mib=total,free_mib=max(0,total-g["used_mib"]) if g.get("available") and total else None,idle_floor_mib=floor,usable_mib=usable,pct_of_usable=min(100,round(100*g["used_mib"]/usable)) if g.get("available") and usable else None,history=list(self.history))
pipes=[]; activejob=self.jobs.get(self.active_job) if self.active_job else None
for pid,p in self.pipelines.items():
comps=[self.component(p,c) for c in p.get("component",[])]; ups=sum(c["state"]=="up" for c in comps)
errors=sum(c["state"]=="error" for c in comps)
state="up" if comps and ups==len(comps) else "error" if errors and ups==0 else "partial" if ups else "down"
job=next((j for j in self.jobs.values() if j.get("pipeline")==pid and j.get("state")=="running"),None)
if job: state="stopping" if job.get("action")=="stop" else "starting"
options=[]; selected=self.selected.get(pid,{})
for o in p.get("option",[]):
choice=selected.get(o.get("id"),o.get("default"))
options.append({"id":o.get("id"),"label":o.get("label"),"kind":o.get("kind"),"selected":choice,"choices":[{"id":c.get("id"),"label":c.get("label"),"note":c.get("note"),"verified":c.get("verified"),"is_default":c.get("id")==o.get("default")} for c in o.get("choice",[])]})
pipes.append({"id":pid,"label":p.get("label",pid),"description":p.get("description",""),"state":state,"health":"ok" if state=="up" else "unknown" if state=="down" else "fail","requires_card":p.get("requires_card",False),"vram_mib":p.get("vram_mib"),"protected":p.get("protected",False),"serving":p.get("serving"),"open_url":p.get("open_url"),"note":p.get("note"),"last_change":self.last_change.get(pid),"components":comps,"selected":selected,"options":options,"job":({"id":job["id"],"state":job["state"],"step":job["steps"][-1]["text"] if job["steps"] else "working"} if job else None)})
if g.get("available"):
for pipe in pipes:
for comp in pipe["components"]:
detail=comp.get("detail","")
if comp.get("state")=="up" and detail.startswith("pid "):
try: pidnum=int(detail.split()[1])
except ValueError: continue
for pr in g["processes"]:
if pr["pid"]==pidnum: pr["pipeline"]=pipe["id"]
return {"host":socket.gethostname(),"now":iso(),"version":__version__,"card":g,"pipelines":pipes,"job":({"id":activejob["id"],"pipeline":activejob.get("pipeline"),"state":activejob["state"],"step":activejob["steps"][-1]["text"] if activejob["steps"] else "working","started":activejob["started"]} if activejob else None),"registry":{"path":self.args.registry,"pipeline_count":len(self.pipelines),"errors":self.errors,"warnings":self.warnings},"listeners":list(getattr(self,"listeners",[]))}
def start_job(self,pid,action,options=None,force=False):
with self.job_guard:
if self.active_job:
raise ValueError(f"another action is running ({self.active_job})")
jid=f"j-{int(time.time()*1000)}-{os.getpid()}"; job={"id":jid,"pipeline":pid,"action":action,"state":"running","started":iso(),"steps":[],"error":None}
self.jobs[jid]=job; self.active_job=jid
threading.Thread(target=self.work,args=(job,options or {},force),daemon=True).start(); return jid
def preflight(self,p,action,options,force):
if p.get("protected") and action in ("stop","restart") and not force:
live=[self.component(p,c)["state"] for c in p.get("component",[])]
if any(x in ("up","error","partial") for x in live):
raise PermissionError(f"protected pipeline {p['id']} is active; force is required to stop it")
if action in ("start","restart") and p.get("requires_card") and not force:
for other in self.pipelines.values():
if other["id"] != p["id"] and other.get("requires_card") and other.get("protected"):
if any(self.component(other,c)["state"] in ("up","error","partial") for c in other.get("component",[])):
raise PermissionError(f"protected pipeline {other['id']} is active; force is required")
if action in ("start","restart"):
for opt in p.get("option",[]):
oid=opt.get("id"); value=options.get(oid,self.selected.get(p["id"],{}).get(oid,opt.get("default")))
choice=next((x for x in opt.get("choice",[]) if x.get("id")==value),None)
if choice is None: raise LookupError(f"unknown option value {oid}={value}")
required=set().union(*(set(x.get("vars",{})) for x in opt.get("choice",[]))) if opt.get("kind")=="argv" else set()
if required-set(choice.get("vars",{})): raise LookupError(f"choice {value} is missing required vars")
def step(self,j,text,ok=True): j["steps"].append({"ts":iso(),"text":text,"ok":ok})
def work(self,j,options,force):
try:
if j["action"]=="free":
self.step(j,"asking the configured reclaim recipe to free the card")
self.reclaim(j)
else:
p=self.pipelines[j["pipeline"]]
if j["action"] in ("stop","restart"): self.stop_pipeline(j,p,force)
if j["action"] in ("start","restart"):
self.start_pipeline(j,p,options,force)
j["state"]="done"; self.last_change[j.get("pipeline")]=iso() if j.get("pipeline") else None
except Exception as exc:
j["error"]=str(exc); self.step(j,f"failed: {exc}",False); j["state"]="failed"
finally:
if self.active_job==j["id"]: self.active_job=None
def validate_options(self,p,requested):
out={}; psel=self.selected.setdefault(p["id"],{})
for o in p.get("option",[]):
oid=o.get("id"); choice=requested.get(oid,psel.get(oid,o.get("default")))
found=next((c for c in o.get("choice",[]) if c.get("id")==choice),None)
if not found: raise ValueError(f"unknown option value {oid}={choice}")
allvars=set().union(*(set(c.get("vars",{})) for c in o.get("choice",[]))) if o.get("kind")=="argv" else set()
if allvars-set(found.get("vars",{})): raise ValueError(f"choice {choice} is missing required vars")
psel[oid]=choice
if o.get("kind")=="argv": out.update(found.get("vars",{}))
self.save_selected(); return out
def start_pipeline(self,j,p,requested,force):
vars=self.validate_options(p,requested)
current=[self.component(p,c)["state"] for c in p.get("component",[])]
if current and all(x=="up" for x in current):
self.step(j,f"{p.get('label',p['id'])} is already healthy; no command was run")
return
if p.get("requires_card"):
for other in self.pipelines.values():
if other["id"]==p["id"] or not other.get("requires_card"): continue
state=[self.component(other,c)["state"] for c in other.get("component",[])]
if any(x in ("up","error","partial") for x in state):
if other.get("protected") and not force: raise RuntimeError(f"protected pipeline {other['id']} is active; force is required")
self.step(j,f"stopping {other.get('label',other['id'])} (releasing the card)…")
self.stop_pipeline(j,other,force,implicit=True)
reclaim=self.registry.get("card",{}).get("reclaim")
if reclaim: self.reclaim(j)
threshold=p.get("start_below_mib")
if threshold is not None:
g=gpu(self.runner); used=g.get("used_mib") if g.get("available") else None
deadline=time.monotonic()+min(120,int(reclaim.get("wait_timeout_s",45)) if reclaim else 45)
if used is None: self.step(j,"GPU probe unavailable; cannot verify VRAM before start",False)
else:
while used>=int(threshold) and time.monotonic()<deadline:
self.step(j,f"waiting for used VRAM < {threshold} MiB (now {used})…")
time.sleep(min(max(.1,float(self.registry.get("card",{}).get("poll_seconds",2))),max(0,deadline-time.monotonic())))
g=gpu(self.runner); used=g.get("used_mib") if g.get("available") else None
if used is None: break
if used is None or used>=int(threshold): self.step(j,f"VRAM wait timed out/unavailable (observed {used}); continuing with start",False)
else: self.step(j,f"used VRAM is {used} MiB; below {threshold} MiB")
for c in p.get("component",[]):
self.step(j,f"starting {c.get('label',c.get('id'))}…")
self.control(p,c,"start",vars)
if self.args.dry_run:
self.step(j,f"dry-run: would start {c.get('label',c.get('id'))}; health was not asserted")
continue
timeout=max(1,int(c.get("ready_timeout_s",20))); deadline=time.monotonic()+min(timeout,300); h=health(c.get("health",{}),self.runner,p["id"]); elapsed=0
while not h["ok"] and time.monotonic()<deadline:
time.sleep(min(.5,max(0,deadline-time.monotonic()))); elapsed+=.5; h=health(c.get("health",{}),self.runner,p["id"])
if elapsed and int(elapsed)%5==0: self.step(j,f"waiting for health ({int(elapsed)}s)…")
if not h["ok"]:
tail=self.logs(p["id"],20); raise RuntimeError(f"{c.get('label',c.get('id'))} health failed after {timeout}s: {h['detail']}; last log lines: {' | '.join(tail[-8:]) or 'log file does not exist'}")
self.step(j,f"{c.get('label',c.get('id'))} healthy ({h['detail']})")
def stop_pipeline(self,j,p,force,implicit=False):
if p.get("protected") and implicit and not force: raise RuntimeError(f"protected pipeline {p['id']} is active; force is required")
for c in reversed(p.get("component",[])):
state=self.component(p,c)
if state["state"]=="down": continue
if c.get("kind")=="proc" and state["detail"]=="external":
msg=f"cannot stop {c.get('label')}: not started by the cockpit (no owned pid)"
if implicit:
self.step(j,msg+" - leaving it running; the card must be freed another way",False); continue
raise RuntimeError(msg)
self.step(j,f"stopping {c.get('label',c.get('id'))}…")
self.control(p,c,"stop",{})
states=[self.component(p,c)["state"] for c in p.get("component",[])]
self.step(j,f"stop verification: {', '.join(states)}; GPU probe reports {gpu(self.runner).get('used_mib')} MiB used")
def control(self,p,c,action,vars):
pid=p["id"]; cid=c.get("id"); kind=c.get("kind")
if kind=="unit":
argv=[self.args.systemctl,action,c["unit"]]
if self.args.sudo: argv=[self.args.sudo,"-n",*argv]
rc,out,_=self.runner.run(argv,pid,action,timeout=30)
if rc: raise RuntimeError(f"systemctl {action} {c['unit']} failed: {out.strip()}")
return
pf=Path(self.args.state_dir)/"pids"/f"{pid}-{cid}.json"; pf.parent.mkdir(parents=True,exist_ok=True)
old=None
try: old=json.loads(pf.read_text())
except Exception: pass
if action=="stop":
if self.args.dry_run: self.runner.run(["kill","-TERM",str((old or {}).get("pid","unknown"))],pid,action); return
if not old or not isinstance(old.get("pid"),int) or proc_start(old["pid"]) is None or str(proc_start(old["pid"]))!=str(old.get("start")):
if health(c.get("health",{}),self.runner,pid)["ok"]: raise RuntimeError(f"refusing to stop {c.get('label')}: external process has no verified pid")
return
self.runner.signal(old["pid"],signal.SIGTERM,pid,action)
deadline=time.monotonic()+min(60,int(c.get("stop_timeout_s",20)))
while proc_start(old["pid"]) is not None and time.monotonic()<deadline: time.sleep(.2)
if proc_start(old["pid"]) is not None:
self.runner.signal(old["pid"],signal.SIGKILL,pid,action)
try: pf.unlink()
except OSError: pass
return
argv=[]
for arg in c.get("argv",[]):
v=PLACEHOLDER.sub(lambda m:str(vars.get(m.group(1),"")),str(arg))
if v: argv.append(v)
if self.args.dry_run:
self.runner.run(argv,pid,action,cwd=c.get("cwd")); return
# Start processes through the runner-owned subprocess API, then detach into their logfile.
log=Path(self.args.state_dir)/"logs"/f"{pid}.log"; log.parent.mkdir(parents=True,exist_ok=True)
env=os.environ.copy(); env.update({str(k):str(v) for k,v in c.get("env",{}).items()})
proc=self.runner.spawn(argv,cwd=c.get("cwd"),env=env,stdout=open(log,"ab"),stderr=__import__("subprocess").STDOUT,pipeline=pid,action=action)
time.sleep(.03); start=proc_start(proc.pid)
if start is None: raise RuntimeError(f"process exited during launch; see {log}")
pf.write_text(json.dumps({"pid":proc.pid,"start":start}),encoding="utf-8")
def reclaim(self,j):
rec=self.registry.get("card",{}).get("reclaim")
if not rec: self.step(j,"no card.reclaim recipe is configured",False); return
if rec.get("kind")=="http":
method,url=rec.get("method","POST"),rec.get("url","")
if self.args.dry_run: self.runner.run(["HTTP",method,url],None,"reclaim"); return
started=time.monotonic()
try:
req=urllib.request.Request(url,data=rec.get("body","").encode() if rec.get("body") else None,method=method,headers={"Content-Type":"application/json"})
with urllib.request.urlopen(req,timeout=10) as res:
self.runner.note(["HTTP",method,url],"reclaim",rc=res.status,duration_ms=(time.monotonic()-started)*1000)
self.step(j,f"reclaim endpoint returned HTTP {res.status}")
except Exception as exc:
self.runner.note(["HTTP",method,url],"reclaim",rc=127,duration_ms=(time.monotonic()-started)*1000)
self.step(j,f"reclaim endpoint failed: {friendly_error(exc,10)}",False)
elif rec.get("kind")=="command":
rc,out,_=self.runner.run(rec.get("argv",[]),None,"reclaim",timeout=30)
self.step(j,f"reclaim command exit {rc}: {out.strip()}",rc==0)
else: self.step(j,f"unknown reclaim kind {rec.get('kind')}",False)
if rec.get("wait_below_mib") is not None and not self.args.dry_run:
deadline=time.monotonic()+min(120,int(rec.get("wait_timeout_s",45))); observed=None
while time.monotonic()<deadline:
g=gpu(self.runner); observed=g.get("used_mib") if g.get("available") else None
if observed is not None and observed<int(rec["wait_below_mib"]):
self.step(j,f"used VRAM is {observed} MiB; below reclaim threshold {rec['wait_below_mib']} MiB"); return
if observed is None: break
time.sleep(min(max(.1,float(self.registry.get("card",{}).get("poll_seconds",2))),max(0,deadline-time.monotonic())))
self.step(j,f"reclaim VRAM wait timed out/unavailable (observed {observed} MiB)",False)
def logs(self,pid,n):
p=self.pipelines.get(pid,{})
units=[c.get("unit") for c in p.get("component",[]) if c.get("kind")=="unit" and c.get("unit")]
if units:
result=[]
for unit in units:
rc,out,_=self.runner.run(["journalctl","-u",unit,"-n",str(max(1,min(1000,n))),"--no-pager"],pid,"logs",timeout=8,readonly=True)
if rc: result.append(f"{unit}: journal unavailable: {out.strip()}")
else: result.extend(out.splitlines())
return result[-n:]
path=Path(self.args.state_dir)/"logs"/f"{pid}.log"
try: return path.read_text(errors="replace").splitlines()[-n:]
except OSError: return [f"log file does not exist: {path}"]
class Handler(BaseHTTPRequestHandler):
server_version="gpu-cockpit"
def log_message(self,*a): pass
@property
def app(self): return self.server.app
def sendj(self,obj,status=200):
b=json.dumps(obj).encode(); self.send_response(status); self.send_header("Content-Type","application/json"); self.send_header("Content-Length",str(len(b))); self.end_headers(); self.wfile.write(b)
def do_GET(self):
u=urlparse(self.path); path=u.path
if not token_matches(self.app.auth_token,self.headers,parse_qs(u.query),self.client_address[0]): return self.sendj({"error":"unauthorized"},401)
if path=="/": return self.file(Path(__file__).parent/"gui"/"index.html","text/html; charset=utf-8")
if path.startswith("/static/"):
name=unquote(path[len("/static/"):]); ext=Path(name).suffix.lower()
if "/" in name or ".." in name or ext not in (".png",".svg",".css",".js",".ico"): return self.send_error(404)
return self.file(Path(__file__).parent/"gui"/"static"/name, {".png":"image/png",".svg":"image/svg+xml",".css":"text/css",".js":"text/javascript",".ico":"image/x-icon"}[ext])
if path=="/api/ping": return self.sendj({"ok":True,"host":socket.gethostname(),"version":__version__,"now":iso()})
if path=="/api/state": return self.sendj(self.app.snapshot())
if path=="/api/registry": return self.sendj({"registry":self.app.registry,"errors":self.app.errors,"warnings":self.app.warnings})
if path=="/api/audit": return self.sendj(self.app.runner.audit(int(parse_qs(u.query).get("n",[100])[0])))
if path.startswith("/api/logs/"):
pid=path.rsplit("/",1)[-1]
if pid not in self.app.pipelines: return self.sendj({"error":"unknown pipeline"},404)
return self.sendj(self.app.logs(pid,int(parse_qs(u.query).get("n",[200])[0])))
if path.startswith("/api/jobs/"):
jid=path.rsplit("/",1)[-1]; job=self.app.jobs.get(jid)
return self.sendj(job or {"error":"job not found"},200 if job else 404)
if path=="/api/events":
q=parse_qs(u.query); once=q.get("once")==["1"]; data=json.dumps(self.app.snapshot())
self.send_response(200); self.send_header("Content-Type","text/event-stream"); self.send_header("Cache-Control","no-cache"); self.end_headers()
self.wfile.write(f"event: state\ndata: {data}\n\n".encode()); self.wfile.flush()
if not once:
try:
last=data; heartbeat=time.monotonic()
while True:
time.sleep(1); current=json.dumps(self.app.snapshot())
old_state=json.loads(last); new_state=json.loads(current)
old_state.pop("now",None); new_state.pop("now",None)
if old_state != new_state:
self.wfile.write(f"event: state\ndata: {current}\n\n".encode()); self.wfile.flush(); last=current
elif time.monotonic()-heartbeat>=15:
self.wfile.write(b": heartbeat\n\n"); self.wfile.flush(); heartbeat=time.monotonic()
except (BrokenPipeError,ConnectionResetError): pass
return
self.send_error(404)
def file(self,path,mime):
try: b=path.read_bytes()
except OSError: return self.send_error(404)
self.send_response(200); self.send_header("Content-Type",mime); self.send_header("Content-Length",str(len(b))); self.end_headers(); self.wfile.write(b)
def do_POST(self):
u=urlparse(self.path); path=u.path
if not token_matches(self.app.auth_token,self.headers,parse_qs(u.query),self.client_address[0]): return self.sendj({"error":"unauthorized"},401)
try: body=json.loads(self.rfile.read(int(self.headers.get("Content-Length",0))) or b"{}")
except Exception: return self.sendj({"error":"invalid JSON"},400)
try:
if path=="/api/card/free": jid=self.app.start_job(None,"free")
else:
sm=re.fullmatch(r"/api/pipelines/([a-z0-9-]+)/options",path)
if sm:
pid=sm.group(1); p=self.app.pipelines.get(pid)
if not p: return self.sendj({"error":"unknown pipeline"},404)
opts=body.get("options",{})
if not isinstance(opts,dict): return self.sendj({"error":"options must be an object"},400)
self.app.preflight(p,"start",opts,True)
self.app.validate_options(p,opts)
return self.sendj({"selected":self.app.selected.get(pid,{})})
m=re.fullmatch(r"/api/pipelines/([a-z0-9-]+)/action",path)
if not m: return self.send_error(404)
pid=m.group(1); p=self.app.pipelines.get(pid)
if not p: return self.sendj({"error":"unknown pipeline"},404)
action=body.get("action")
if action not in ("start","stop","restart"): return self.sendj({"error":"invalid action"},400)
options=body.get("options",{})
if not isinstance(options,dict): return self.sendj({"error":"options must be an object"},400)
for opt in p.get("option",[]):
if opt.get("id") in options and options[opt["id"]] not in [x.get("id") for x in opt.get("choice",[])]: return self.sendj({"error":f"unknown option value {opt['id']}={options[opt['id']] }"},400)
self.app.preflight(p,action,options,bool(body.get("force",False)))
jid=self.app.start_job(pid,action,options,bool(body.get("force",False)))
return self.sendj({"job_id":jid},202)
except PermissionError as exc: return self.sendj({"error":str(exc)},409)
except LookupError as exc: return self.sendj({"error":str(exc)},400)
except ValueError as exc: return self.sendj({"error":str(exc)},409)
except Exception as exc: return self.sendj({"error":str(exc)},500)
def serve_extras(app,hosts,port):
"""Bind the front-door listeners, retrying until they exist. A WireGuard address only
appears once wg0 is up, which can happen after this service starts - without the retry a
boot race would close the door permanently and look like an auth problem from the phone."""
pending=list(hosts)
while pending:
for h in list(pending):
try: srv=ThreadingHTTPServer((h,port),Handler)
except OSError as exc:
print(f"front door {h}:{port} not available yet ({exc}) - retrying in 5s",flush=True); continue
srv.app=app; app.listeners.append(f"{h}:{port}")
print(f"listening on http://{h}:{port} (front door)",flush=True)
threading.Thread(target=srv.serve_forever,daemon=True).start()
pending.remove(h)
if pending: time.sleep(5)
def main(argv=None):
home=str(Path.home()); envreg=os.environ.get("COCKPIT_REGISTRY")
ap=argparse.ArgumentParser(); ap.add_argument("--registry",default=envreg or os.path.join(home,".config/gpu-cockpit/registry.toml")); ap.add_argument("--host",default="127.0.0.1",help="comma-separated; anything non-loopback needs a token"); ap.add_argument("--token-file",default=os.path.join(home,".config/gpu-cockpit/token")); ap.add_argument("--port",type=int,default=8770); ap.add_argument("--state-dir",default=os.path.join(home,".local/state/gpu-cockpit")); ap.add_argument("--systemctl",default="systemctl"); ap.add_argument("--sudo",default=None); ap.add_argument("--dry-run",action="store_true"); ap.add_argument("--check",action="store_true"); args=ap.parse_args(argv)
reg,errors,warnings=load(args.registry)
if args.check:
print(f"registry OK: {len(reg.get('pipeline',[]))} pipelines, {len(errors)} errors, {len(warnings)} warnings")
for x in errors: print(f"ERROR: {x}")
for x in warnings: print(f"WARNING: {x}")
return 1 if errors else 0
app=App(args)
app.auth_token=load_token(args.token_file); app.listeners=[]
hosts=[h.strip() for h in str(args.host).split(",") if h.strip()] or ["127.0.0.1"]
extras=[h for h in hosts if not is_loopback(h)]
if extras and not app.auth_token:
print(f"refusing to listen on {', '.join(extras)} with no token ({args.token_file}) - that is an open front door",file=sys.stderr); return 1
try: server=ThreadingHTTPServer((hosts[0],args.port),Handler)
except OSError as exc: print(f"cannot listen on http://{hosts[0]}:{args.port}: {exc}",file=sys.stderr); return 1
server.app=app; app.listeners.append(f"{hosts[0]}:{args.port}")
print(f"listening on http://{hosts[0]}:{args.port}",flush=True)
if extras: threading.Thread(target=serve_extras,args=(app,extras,args.port),daemon=True).start()
try: server.serve_forever()
except KeyboardInterrupt: pass
finally: server.server_close()
return 0
+12
View File
@@ -0,0 +1,12 @@
[Unit]
Description=GPU Cockpit single-card pipeline console
After=network.target
[Service]
Type=simple
ExecStart=/usr/bin/python3 -m cockpit --registry %h/.config/gpu-cockpit/registry.toml --state-dir %h/.local/state/gpu-cockpit --host 127.0.0.1 --port 8770
Restart=on-failure
RestartSec=3
[Install]
WantedBy=default.target
+119
View File
@@ -0,0 +1,119 @@
[card]
label = "RTX 4060"
total_mib = 8188
idle_floor_mib = 1700
poll_seconds = 2
[card.reclaim]
label = "ask the image server to unload models"
kind = "http"
url = "http://127.0.0.1:8188/free"
method = "POST"
body = '{"unload_models": true, "free_memory": true}'
wait_below_mib = 2000
wait_timeout_s = 45
[[pipeline]]
id = "image"
label = "Image generation"
description = "An image generation service on the shared card"
requires_card = true
vram_mib = 7019
protected = false
start_below_mib = 2600
serving = "Z-Image-Turbo Q4_K_M + LoRA ladder"
open_url = "http://127.0.0.1:8188"
note = "Peaks near 7 GB; leave the card alone while it renders."
[[pipeline.component]]
kind = "proc"
id = "image-server"
label = "Image server"
argv = ["/usr/bin/python3", "-m", "http.server", "8188", "--bind", "127.0.0.1", "{extra}"]
cwd = "/tmp"
env = { PYTHONUNBUFFERED = "1" }
ready_timeout_s = 15
stop_timeout_s = 5
health = { kind = "http", url = "http://127.0.0.1:8188/", timeout_s = 2 }
[[pipeline.option]]
id = "model"
label = "Model"
kind = "argv"
default = "turbo"
[[pipeline.option.choice]]
id = "turbo"
label = "Z-Image Turbo"
note = "Fast image generation"
verified = "2026-09-26"
vars = { model = "/models/z-image.gguf", alias = "image-turbo", extra = "" }
[[pipeline.option.choice]]
id = "quality"
label = "Image quality model"
vars = { model = "/models/image-quality.gguf", alias = "image-quality", extra = "--directory ." }
[[pipeline.option]]
id = "profile"
label = "Profile note"
kind = "note"
default = "desktop"
[[pipeline.option.choice]]
id = "desktop"
label = "Desktop preset"
note = "Leaves memory available for the desktop."
[[pipeline.option.choice]]
id = "render"
label = "Render preset"
note = "Use for dedicated render sessions."
[[pipeline]]
id = "voice"
label = "Voice stack"
description = "A systemd managed voice pipeline"
requires_card = true
vram_mib = 3200
protected = false
serving = "Speech and verbalization"
[[pipeline.component]]
kind = "unit"
id = "voice-server"
label = "Voice server"
unit = "gpu-cockpit-example-voice.service"
ready_timeout_s = 20
health = { kind = "command", argv = ["/usr/bin/test", "-e", "/dev/null"], timeout_s = 2 }
[[pipeline]]
id = "trainer"
label = "Trainer"
description = "Explicitly protected card workload"
requires_card = true
protected = true
vram_mib = 7000
serving = "Training job"
[[pipeline.component]]
kind = "proc"
id = "trainer-proc"
label = "Training process"
argv = ["/bin/sleep", "300"]
ready_timeout_s = 2
stop_timeout_s = 2
health = { kind = "none" }
[[pipeline]]
id = "cpu-llm"
label = "CPU model"
description = "Runs without claiming the GPU"
requires_card = false
vram_mib = 0
serving = "CPU inference"
[[pipeline.component]]
kind = "proc"
id = "cpu"
label = "CPU service"
argv = ["/bin/sleep", "300"]
ready_timeout_s = 2
stop_timeout_s = 2
health = { kind = "none" }
+4
View File
@@ -0,0 +1,4 @@
#!/usr/bin/env bash
set -euo pipefail
root="$(cd "$(dirname "$0")/.." && pwd)"
exec python3 "$root/tests/selftest.py" "$root"
+119
View File
@@ -0,0 +1,119 @@
import json, os, pathlib, socket, subprocess, sys, tempfile, threading, time, urllib.error, urllib.request
from http.server import BaseHTTPRequestHandler, HTTPServer
root=pathlib.Path(sys.argv[1]); tmp=pathlib.Path(tempfile.mkdtemp(prefix='cockpit-test-')); cases=0; proc=None
def check(name, test):
global cases
if not test: print(f'FAIL: {name}: assertion failed',flush=True); raise SystemExit(1)
cases+=1; print(f'PASS: {name}',flush=True)
class Reclaim(BaseHTTPRequestHandler):
hits=0
def do_POST(self): Reclaim.hits+=1; self.send_response(200); self.end_headers(); self.wfile.write(b'ok')
def log_message(self,*a): pass
reclaim=HTTPServer(('127.0.0.1',0),Reclaim); threading.Thread(target=reclaim.serve_forever,daemon=True).start()
stub=tmp/'systemctl'; stub.write_text('''#!/usr/bin/env sh
s="$COCKPIT_STUB_STATE"; c="$1"; u="$2"
case "$c" in
show) if grep -qx "$u" "$s" 2>/dev/null; then echo ActiveState=active; echo SubState=running; echo ExecMainPID=1; else echo ActiveState=inactive; echo SubState=dead; echo ExecMainPID=0; fi;;
start|restart) grep -vx "$u" "$s" 2>/dev/null > "$s.tmp" || true; echo "$u" >> "$s.tmp"; mv "$s.tmp" "$s";;
stop) grep -vx "$u" "$s" 2>/dev/null > "$s.tmp" || true; mv "$s.tmp" "$s";;
is-active) grep -qx "$u" "$s" && exit 0 || exit 3;; *) exit 0;; esac
'''); stub.chmod(0o755)
port=socket.socket(); port.bind(('127.0.0.1',0)); pnum=port.getsockname()[1]; port.close()
reg=tmp/'registry.toml'
reg.write_text(f'''[card]\ntotal_mib=8188\nidle_floor_mib=1700\npoll_seconds=0.05\n[card.reclaim]\nkind="http"\nurl="http://127.0.0.1:{reclaim.server_port}/free"\nmethod="POST"\nwait_timeout_s=1\n[[pipeline]]\nid="a"\nlabel="A"\nrequires_card=true\nstart_below_mib=2600\n[[pipeline.component]]\nid="p"\nkind="proc"\nlabel="A process"\nargv=["{sys.executable}","-c","import time;time.sleep(60)"]\nready_timeout_s=2\nstop_timeout_s=1\nhealth={{kind="none"}}\n[[pipeline]]\nid="b"\nlabel="B"\nrequires_card=true\n[[pipeline.component]]\nid="p"\nkind="proc"\nlabel="B process"\nargv=["{sys.executable}","-c","import time;time.sleep(60)"]\nready_timeout_s=2\nstop_timeout_s=1\nhealth={{kind="none"}}\n[[pipeline]]\nid="protected"\nlabel="Protected"\nrequires_card=true\nprotected=true\n[[pipeline.component]]\nid="p"\nkind="proc"\nlabel="Protected process"\nargv=["{sys.executable}","-c","import time;time.sleep(60)"]\nready_timeout_s=2\nstop_timeout_s=1\nhealth={{kind="none"}}\n[[pipeline]]\nid="bad-health"\nlabel="Unreachable"\n[[pipeline.component]]\nid="p"\nkind="proc"\nlabel="Bad process"\nargv=["{sys.executable}","-c","import time;time.sleep(60)"]\nready_timeout_s=1\nstop_timeout_s=1\nhealth={{kind="tcp",host="127.0.0.1",port=1,timeout_s=0.1}}\n[[pipeline]]\nid="options"\nlabel="Options"\n[[pipeline.component]]\nid="p"\nkind="proc"\nlabel="Option process"\nargv=["{sys.executable}","-c","import time;time.sleep(60)","{{extra}}"]\nready_timeout_s=2\nstop_timeout_s=1\nhealth={{kind="none"}}\n[[pipeline.option]]\nid="model"\nlabel="Model"\nkind="argv"\ndefault="one"\n[[pipeline.option.choice]]\nid="one"\nlabel="One"\nvars={{extra=""}}\n[[pipeline.option.choice]]\nid="two"\nlabel="Two"\nvars={{extra="-v"}}\n[[pipeline.option]]\nid="required"\nlabel="Required var"\nkind="argv"\ndefault="present"\n[[pipeline.option.choice]]\nid="present"\nlabel="Present"\nvars={{required="yes"}}\n[[pipeline.option.choice]]\nid="absent"\nlabel="Missing required value"\nvars={{}}\n[[pipeline]]\nid="broken"\nlabel="Broken"\n[[pipeline.component]]\nid="p"\nkind="proc"\nlabel="Broken"\nargv=["{sys.executable}","{{missing}}"]\nhealth={{kind="none"}}\n[[pipeline]]\nid="external"\nlabel="External"\nrequires_card=true\n[[pipeline.component]]\nid="p"\nkind="proc"\nlabel="External process"\nargv=["{sys.executable}","-c","import time;time.sleep(60)"]\nready_timeout_s=2\nhealth={{kind="tcp",host="127.0.0.1",port={reclaim.server_port},timeout_s=1}}\n''')
state=tmp/'state'; env=os.environ.copy(); env['COCKPIT_STUB_STATE']=str(tmp/'units'); log=open(tmp/'server.log','w+')
proc=subprocess.Popen([sys.executable,'-m','cockpit','--registry',str(reg),'--port',str(pnum),'--state-dir',str(state),'--systemctl',str(stub)],cwd=root,stdout=log,stderr=log,env=env)
base=f'http://127.0.0.1:{pnum}'
def get(path):
with urllib.request.urlopen(base+path,timeout=5) as r:return r.status,r.read().decode(),r.headers
def post(path,obj):
req=urllib.request.Request(base+path,data=json.dumps(obj).encode(),headers={'Content-Type':'application/json'},method='POST')
try:
with urllib.request.urlopen(req,timeout=5) as r:return r.status,json.loads(r.read())
except urllib.error.HTTPError as e:return e.code,json.loads(e.read())
def snapshot():return json.loads(get('/api/state')[1])
def job(jid):
end=time.time()+5
while time.time()<end:
j=json.loads(get('/api/jobs/'+jid)[1])
if j['state']!='running':return j
time.sleep(.05)
raise AssertionError('job timed out')
def action(pid,act='start',**kw):
code,r=post(f'/api/pipelines/{pid}/action',{'action':act,**kw})
return job(r['job_id']) if code==202 else (code,r)
try:
deadline=time.time()+5
while time.time()<deadline:
try: status,raw,_=get('/api/state'); break
except Exception: time.sleep(.05)
check('server starts, listens, and state shape',status==200 and all(k in json.loads(raw) for k in ('host','card','pipelines','registry')))
check('fixture pipeline count',len(json.loads(raw)['pipelines'])==7)
check('down state for not-started proc',next(x for x in snapshot()['pipelines'] if x['id']=='a')['state']=='down')
check('unreachable health state is down',next(x for x in snapshot()['pipelines'] if x['id']=='bad-health')['state']=='down')
j=action('a'); steps=' '.join(x['text'] for x in j['steps'])
vram_verdict=('GPU probe unavailable' in steps) or ('VRAM wait' in steps) or ('below 2600 MiB' in steps)
check('start verifies health/PID with an honest VRAM verdict',j['state']=='done' and vram_verdict and 'healthy' in steps and next(x for x in snapshot()['pipelines'] if x['id']=='a')['components'][0]['detail'].startswith('pid '))
j=action('a','stop'); check('stop verifies pipeline down',j['state']=='done' and next(x for x in snapshot()['pipelines'] if x['id']=='a')['state']=='down')
action('a'); j=action('b'); check('arbitration stops A before B',j['state']=='done' and any('stopping A' in x['text'] for x in j['steps']) and next(x for x in snapshot()['pipelines'] if x['id']=='a')['state']=='down')
action('protected'); status,body=post('/api/pipelines/b/action',{'action':'start'}); check('protected switch returns 409 without change',status==409 and 'protected' in body['error'] and next(x for x in snapshot()['pipelines'] if x['id']=='protected')['state']=='up')
status,body=post('/api/pipelines/protected/action',{'action':'stop'}); check('protected stop requires force',status==409 and 'protected' in body['error'] and next(x for x in snapshot()['pipelines'] if x['id']=='protected')['state']=='up')
j=action('protected','stop',force=True); check('force stop of a protected pipeline works',j['state']=='done' and next(x for x in snapshot()['pipelines'] if x['id']=='protected')['state']=='down')
j=action('b','start',force=True); check('force allows protected switch',j['state']=='done')
check('reclaim HTTP recipe called',Reclaim.hits>=2)
action('b')
ext=next(x for x in snapshot()['pipelines'] if x['id']=='external')
check('unowned but answering component reads up/external',ext['state']=='up' and ext['components'][0]['detail']=='external')
j=action('external','stop'); check('explicit stop of an unowned process fails loudly',j['state']=='failed' and 'not started by the cockpit' in j['error'])
j=action('a','start'); check('a switch leaves an unowned holder running and still proceeds',j['state']=='done' and any('leaving it running' in x['text'] for x in j['steps']))
status,_=post('/api/pipelines/options/action',{'action':'start','options':{'model':'invalid'}}); check('unknown option value returns 400',status==400)
status,_=post('/api/pipelines/options/action',{'action':'start','options':{'required':'absent'}}); check('choice missing a declared var returns 400',status==400)
status,body=post('/api/pipelines/options/options',{'options':{'model':'two'}}); check('option selection persists',status==200 and next(x for x in snapshot()['pipelines'] if x['id']=='options')['selected'].get('model')=='two')
j=action('bad-health'); check('ready timeout fails with log tail',j['state']=='failed' and 'last log lines' in j['error'])
pf=state/'pids'/'a-p.json'; pf.parent.mkdir(parents=True,exist_ok=True); pf.write_text('not json'); check('garbage pidfile reads down',next(x for x in snapshot()['pipelines'] if x['id']=='a')['state']=='down')
pf.write_text(json.dumps({'pid':os.getpid(),'start':'wrong'})); check('reused pid mismatch reads down',next(x for x in snapshot()['pipelines'] if x['id']=='a')['state']=='down')
status,raw,_=get('/api/events?once=1'); check('SSE once emits state event',status==200 and 'event: state' in raw)
status,html,_=get('/'); check('GUI markup served',status==200 and 'Free the card' in html)
status,raw,_=get('/api/registry'); check('broken registry validation reported',status==200 and bool(json.loads(raw)['errors']))
# Verify dry-run records intent without creating proc children.
dryport=socket.socket(); dryport.bind(('127.0.0.1',0)); dp=dryport.getsockname()[1]; dryport.close()
dry=subprocess.Popen([sys.executable,'-m','cockpit','--registry',str(reg),'--port',str(dp),'--state-dir',str(tmp/'dry'),'--dry-run'],cwd=root,stdout=subprocess.PIPE,stderr=subprocess.STDOUT,text=True,env=env)
try:
line=dry.stdout.readline(); check('dry-run server starts',line.startswith('listening on'))
req=urllib.request.Request(f'http://127.0.0.1:{dp}/api/pipelines/a/action',data=b'{"action":"start"}',headers={'Content-Type':'application/json'},method='POST')
jid=json.loads(urllib.request.urlopen(req).read())['job_id']; time.sleep(.2)
req=urllib.request.Request(f'http://127.0.0.1:{dp}/api/pipelines/options/action',data=b'{"action":"start","options":{"model":"two"}}',headers={'Content-Type':'application/json'},method='POST')
urllib.request.urlopen(req).read(); time.sleep(.2); rows=[json.loads(x) for x in (tmp/'dry'/'audit.jsonl').read_text().splitlines()]
check('dry-run audits and launches no process',any(x['dry_run'] for x in rows) and not (tmp/'dry'/'pids'/'a-p.json').exists())
check('selected option changes recorded argv',any(x['pipeline']=='options' and '-v' in x['argv'] for x in rows))
finally: dry.terminate(); dry.wait(timeout=3)
# --- the front door: the tunnel listener is token-gated, and cannot be opened without one ---
if str(root) not in sys.path: sys.path.insert(0,str(root))
from cockpit import server as srvmod
check('loopback needs no token (the desktop GUI and the TUI live there)',srvmod.token_matches('sekret',{},{},'127.0.0.1'))
check('ipv6 loopback counts as loopback',srvmod.token_matches('sekret',{},{},'::1'))
check('a tunnel caller with no token is refused',not srvmod.token_matches('sekret',{},{},'192.168.2.1'))
check('a tunnel caller with the right bearer token gets in',srvmod.token_matches('sekret',{'Authorization':'Bearer sekret'},{},'192.168.2.1'))
check('a wrong bearer token is refused',not srvmod.token_matches('sekret',{'Authorization':'Bearer nope'},{},'192.168.2.1'))
check('a query token serves clients that cannot set headers',srvmod.token_matches('sekret',{}, {'token':['sekret']},'192.168.2.1'))
(tmp/'tok').write_text('sekret\n')
check('a token file is read back trimmed',srvmod.load_token(str(tmp/'tok'))=='sekret')
check('a missing token file means no token, not an empty one',srvmod.load_token(str(tmp/'absent')) is None)
spare=socket.socket(); spare.bind(('127.0.0.1',0)); sport=spare.getsockname()[1]; spare.close()
guard=subprocess.run([sys.executable,'-m','cockpit','--registry',str(reg),'--port',str(sport),'--state-dir',str(tmp/'guard'),
'--host','127.0.0.1,10.10.10.1','--token-file',str(tmp/'absent'),'--dry-run'],
cwd=root,stdout=subprocess.PIPE,stderr=subprocess.STDOUT,text=True,env=env,timeout=30)
check('a front door with no token is refused outright, not silently left open',guard.returncode==1 and 'open front door' in guard.stdout)
print(f'ALL TESTS PASSED ({cases} cases)',flush=True)
finally:
if proc:
proc.terminate()
try:proc.wait(timeout=3)
except subprocess.TimeoutExpired:proc.kill();proc.wait()
for f in (state/'pids').glob('*.json') if (state/'pids').exists() else []:
try:
pid=json.loads(f.read_text())['pid']
if pid != os.getpid(): os.kill(pid,15)
except Exception:pass
reclaim.shutdown(); log.close()
+111
View File
@@ -0,0 +1,111 @@
#!/usr/bin/env python3
"""Textual HTTP client for gpu-cockpit; --print needs only the standard library."""
import argparse, json, sys, time, urllib.error, urllib.request
def get(url):
with urllib.request.urlopen(url,timeout=3) as r: return json.loads(r.read())
def plain(url):
try:
s=get(url.rstrip('/')+'/api/state')
except Exception as exc:
print(f"Cannot reach gpu-cockpit at {url}: {exc}. Start it with: python3 -m cockpit --registry ~/.config/gpu-cockpit/registry.toml")
return 2
c=s['card'];
if c['available']: print(f"{s['host']} · GPU {c['used_mib']}/{c['usable_mib']} MiB usable · {c['util_pct']}% · {c['temp_c']}°C")
else: print(f"{s['host']} · no GPU probe ({c.get('detail','unavailable')})")
for p in s['pipelines']: print(f"{p['state']:9} {p['id']:<16} {p['label']} · {p.get('vram_mib') or 0} MiB")
return 0
def main():
ap=argparse.ArgumentParser(); ap.add_argument('--url',default='http://127.0.0.1:8770'); ap.add_argument('--refresh',type=float,default=2); ap.add_argument('--print',dest='plain',action='store_true'); a=ap.parse_args()
if a.plain: return plain(a.url)
try:
from textual.app import App, ComposeResult
from textual.containers import Horizontal, Vertical
from textual.widgets import Button, DataTable, Footer, Header, Select, Static
except ImportError:
print('Textual is required for the interactive TUI. Install it with: pip install textual',file=sys.stderr); return 2
class Cockpit(App):
TITLE='GPU Cockpit'; BINDINGS=[('q','quit','Quit'),('r','reload','Refresh'),('s','start','Start'),('x','stop','Stop'),('R','restart','Restart'),('f','free','Free card'),('l','logs','Logs'),('m','choose','Choose option')]
CSS='Screen{layout:vertical} #content{height:1fr} #pipes{width:40%} #detail{width:60%;border:round $primary;padding:1} #step{height:3;color:$accent}'
def compose(self)->ComposeResult:
yield Header(); yield Static('Connecting…',id='step')
with Horizontal(id='content'):
yield DataTable(id='pipes'); yield Static('Select a pipeline',id='detail')
yield Footer()
def on_mount(self):
self.query_one('#pipes',DataTable).add_columns('State','ID','Pipeline','VRAM MiB'); self.set_interval(a.refresh,self.reload_state)
def reload_state(self):
try: self.snap=get(a.url.rstrip('/')+'/api/state'); table=self.query_one('#pipes',DataTable); table.clear()
except Exception as exc:
self.snap=None; self.query_one('#step',Static).update(f"Server unavailable at {a.url}. Start: python3 -m cockpit --registry ~/.config/gpu-cockpit/registry.toml ({exc})"); return
c=self.snap['card']; head=f"{self.snap['host']} · "+(f"GPU {c['used_mib']}/{c['usable_mib']} MiB · {c['util_pct']}% · {c['temp_c']}°C" if c['available'] else 'no GPU probe')
if self.snap.get('job'): head+=f" · {self.snap['job']['step']}"
self.query_one('#step',Static).update(head)
for p in self.snap['pipelines']: table.add_row({'up':'●','down':'○','partial':'◐','starting':'◒','stopping':'◓'}.get(p['state'],'?'),p['id'],p['label'],str(p.get('vram_mib') or 0),key=p['id'])
def on_data_table_row_highlighted(self,event):
if not self.snap:return
p=next((x for x in self.snap['pipelines'] if x['id']==str(event.row_key.value)),None)
if p:
try: logs='\n'.join(get(a.url.rstrip('/')+f"/api/logs/{p['id']}?n=20"))
except Exception as exc: logs=f"Logs unavailable: {exc}"
opts='\n'.join(f"{o['label']}: {o['selected']}" for o in p['options']) or 'No options'
comps='\n'.join(f"{c['label']}: {c['state']} · {c.get('health',{}).get('detail') if c.get('health') else c['detail']}" for c in p['components'])
detail=f"{p['label']} · {p['state']}\n{p['description']}\nServing: {p.get('serving') or '—'}\nVRAM: {p.get('vram_mib') or 0} MiB\n\nComponents\n{comps}\n\nOptions\n{opts}\n\nLast 20 log lines\n{logs}"
if p.get('job'): detail+=f"\n\nJob {p['job']['id']}: {p['job']['step']}"
self.query_one('#detail',Static).update(detail)
def selected(self):
if not self.snap:return None
key=self.query_one('#pipes',DataTable).cursor_row
try: pid=self.query_one('#pipes',DataTable).get_row_at(key)[1]
except Exception:return None
return next((x for x in self.snap['pipelines'] if x['id']==pid),None)
def post(self,path,data):
req=urllib.request.Request(a.url.rstrip('/')+path,data=json.dumps(data).encode(),headers={'Content-Type':'application/json'},method='POST')
try: urllib.request.urlopen(req,timeout=3).read()
except Exception as exc:self.query_one('#step',Static).update(str(exc))
self.set_timer(.2,self.reload_state)
def action(self,name):
p=self.selected()
if p:self.post(f"/api/pipelines/{p['id']}/action",{'action':name,'options':p['selected']})
def action_reload(self):self.reload_state()
def action_start(self):self.action('start')
def action_stop(self):self.action('stop')
def action_restart(self):self.action('restart')
def action_free(self):self.post('/api/card/free',{})
def action_logs(self):
p=self.selected()
if p:
try:self.push_screen(LogsScreen('\n'.join(get(a.url.rstrip('/')+f"/api/logs/{p['id']}?n=20"))))
except Exception as exc:self.query_one('#step',Static).update(str(exc))
def action_choose(self):
p=self.selected()
if p:self.push_screen(OptionScreen(p))
from textual.screen import ModalScreen
class LogsScreen(ModalScreen):
BINDINGS=[('escape','dismiss','Close')]
def __init__(self,text):super().__init__();self.text=text
def compose(self):yield Static(self.text)
class OptionScreen(ModalScreen):
BINDINGS=[('escape','dismiss','Close')]
CSS='OptionScreen{align: center middle} #box{width:70%;height:auto;max-height:80%;border:round $primary;background:$surface;padding:1 2} Select{margin:1 0} Button{margin-top:1}'
def __init__(self,pipeline):super().__init__();self.pipeline=pipeline;self.values=dict(pipeline['selected'])
def compose(self):
with Vertical(id='box'):
yield Static(f"Choose options · {self.pipeline['label']}")
for option in self.pipeline['options']:
values=[(c['label'],c['id']) for c in option['choices']]
if values: yield Select(values,value=option['selected'],prompt=option['label'],id='option-'+option['id'])
yield Button('Save selection',variant='primary',id='save')
def on_select_changed(self,event):
if event.select.id and event.select.id.startswith('option-') and event.value is not Select.BLANK:
self.values[event.select.id[7:]]=event.value
def on_button_pressed(self,event):
if event.button.id=='save':
path=f"/api/pipelines/{self.pipeline['id']}/options"
req=urllib.request.Request(a.url.rstrip('/')+path,data=json.dumps({'options':self.values}).encode(),headers={'Content-Type':'application/json'},method='POST')
try: urllib.request.urlopen(req,timeout=3).read(); self.dismiss(self.values)
except Exception as exc: self.query_one(Static).update(str(exc))
Cockpit().run(); return 0
if __name__=='__main__': raise SystemExit(main())