Compare commits
8
Commits
main
...
cockpit-v1
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
6029bfed64 | ||
|
|
cb4d250eb2 | ||
|
|
86055b2e69 | ||
|
|
eb217ea3e1 | ||
|
|
5ad1e02168 | ||
|
|
8af2e04866 | ||
|
|
ff11e0189f | ||
|
|
152db1098e |
@@ -0,0 +1,3 @@
|
||||
__pycache__/
|
||||
*.pyc
|
||||
.venv/
|
||||
@@ -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.
|
||||
@@ -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
|
||||
```
|
||||
|
||||
@@ -0,0 +1,2 @@
|
||||
"""Registry-driven single-GPU pipeline console."""
|
||||
__version__ = "0.1.0"
|
||||
@@ -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.
@@ -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=>({'&':'&','<':'<','>':'>','"':'"',"'":'''}[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 & 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 |
@@ -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
|
||||
@@ -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", []))
|
||||
@@ -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
|
||||
@@ -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
|
||||
@@ -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
|
||||
@@ -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" }
|
||||
Executable
+4
@@ -0,0 +1,4 @@
|
||||
#!/usr/bin/env bash
|
||||
set -euo pipefail
|
||||
root="$(cd "$(dirname "$0")/.." && pwd)"
|
||||
exec python3 "$root/tests/selftest.py" "$root"
|
||||
@@ -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
@@ -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())
|
||||
Reference in New Issue
Block a user