Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
dc47e335aa | ||
|
|
3572d5ba35 | ||
|
|
d9f3986198 |
@@ -33,3 +33,5 @@ Thumbs.db
|
||||
|
||||
# Claude
|
||||
Claude.md
|
||||
coverage/
|
||||
__pycache__/
|
||||
|
||||
@@ -36,6 +36,70 @@ Reference documentation:
|
||||
- [`web_template/aesthetic_diff.md`](https://code.lotusguild.org/LotusGuild/web_template/src/branch/main/aesthetic_diff.md) — cross-app divergence analysis and convergence guide
|
||||
- [`web_template/node/middleware.js`](https://code.lotusguild.org/LotusGuild/web_template/src/branch/main/node/middleware.js) — Express auth, CSRF, CSP nonce middleware
|
||||
|
||||

|
||||
|
||||
## Why this exists
|
||||
|
||||
I run a six-node Proxmox/Ceph cluster plus a pile of LXC containers. Day-to-day work is the same handful of questions asked of many machines: *what is the disk usage everywhere, restart that service on one node, is it safe to proceed?* SSH-ing into each box does not scale, and a full config-management system is more than I need. PULSE is a small orchestration layer: a server with a web UI, and a lightweight **worker** on each machine that executes what the server dispatches over a WebSocket. It adds what shell loops lack: multi-step **workflows** with conditions and human approval gates, scheduled commands, a full execution history, and an API that other tools call (GANDALF uses it to run network diagnostics).
|
||||
|
||||
```mermaid
|
||||
flowchart LR
|
||||
user["Browser<br/>(Authelia SSO)"] -->|HTTPS| srv["PULSE server<br/>Express + EJS + WebSocket"]
|
||||
gd["GANDALF"] -->|internal API| srv
|
||||
srv --> db[("MariaDB<br/>workflows, executions,<br/>schedules")]
|
||||
srv <-->|"WebSocket<br/>commands / results / heartbeat"| w1["worker: node-01"]
|
||||
srv <-->|WebSocket| w2["worker: node-02"]
|
||||
srv <-->|WebSocket| w3["worker: node-03"]
|
||||
```
|
||||
|
||||
Workflows are plain JSON with five step types: `execute`, `prompt` (pause for a human decision), `wait`, `parse` (pull `KEY=VALUE` lines from output into state) and `route` (branch on parsed state).
|
||||
|
||||
## Gallery
|
||||
|
||||
**Dashboard**: recent executions at a glance and live worker status.
|
||||
|
||||

|
||||
|
||||
**Workers** report system info, memory, load, uptime and task slots on every heartbeat:
|
||||
|
||||

|
||||
|
||||
**Workflows** are reusable multi-step jobs. The editor takes the JSON definition directly:
|
||||
|
||||
| Workflow list | Workflow editor |
|
||||
|---|---|
|
||||
|  |  |
|
||||
|
||||
**Executions**: full history with step-by-step logs, per-command output and durations.
|
||||
|
||||

|
||||
|
||||
| Completed run with step logs | A failed command, with its error |
|
||||
|---|---|
|
||||
|  |  |
|
||||
|
||||
**Approval gates:** a `prompt` step pauses the workflow until someone chooses an option (here, a rolling restart waits for a human to approve):
|
||||
|
||||

|
||||
|
||||
**Quick Command** runs a one-off command on one or many workers at once:
|
||||
|
||||
| Compose | Result |
|
||||
|---|---|
|
||||
|  |  |
|
||||
|
||||
**Scheduler** runs commands on an interval or hourly schedule:
|
||||
|
||||

|
||||
|
||||
**Command palette** (`Ctrl+K`) and **light theme**:
|
||||
|
||||
| Command palette | Light theme |
|
||||
|---|---|
|
||||
|  |  |
|
||||
|
||||
<sub>Screenshots are generated with Playwright against a throwaway local instance with three local workers running harmless commands (`uptime`, `df`, `echo`). See `scripts/demo/`.</sub>
|
||||
|
||||
## Web UI
|
||||
|
||||
The UI is server-rendered with EJS. Every route renders `views/pages/<page>.ejs` into the shared
|
||||
|
||||
|
After Width: | Height: | Size: 24 KiB |
|
After Width: | Height: | Size: 128 KiB |
|
After Width: | Height: | Size: 168 KiB |
|
After Width: | Height: | Size: 567 KiB |
|
After Width: | Height: | Size: 100 KiB |
|
After Width: | Height: | Size: 98 KiB |
|
After Width: | Height: | Size: 91 KiB |
|
After Width: | Height: | Size: 80 KiB |
|
After Width: | Height: | Size: 214 KiB |
|
After Width: | Height: | Size: 139 KiB |
|
After Width: | Height: | Size: 115 KiB |
|
After Width: | Height: | Size: 123 KiB |
|
After Width: | Height: | Size: 116 KiB |
|
After Width: | Height: | Size: 59 KiB |
|
After Width: | Height: | Size: 133 KiB |
@@ -11,6 +11,7 @@
|
||||
"license": "ISC",
|
||||
"description": "",
|
||||
"dependencies": {
|
||||
"axios": "^1.20.0",
|
||||
"cron-parser": "^5.5.0",
|
||||
"dotenv": "^17.2.3",
|
||||
"ejs": "3.1.10",
|
||||
@@ -22,6 +23,6 @@
|
||||
},
|
||||
"devDependencies": {
|
||||
"eslint": "^8.57.1",
|
||||
"jest": "^29.7.0"
|
||||
"jest": "^30.5.2"
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,18 @@
|
||||
# Demo environment (for README screenshots)
|
||||
|
||||
Fictional data only. A throwaway MariaDB on loopback, the PULSE server and three local workers that run harmless commands.
|
||||
|
||||
```bash
|
||||
D=$(mktemp -d); rsync -a --exclude .git . $D/pulse && cd $D/pulse && npm install
|
||||
mariadb-install-db --no-defaults --datadir=$D/db --auth-root-authentication-method=normal --skip-test-db
|
||||
mariadbd --no-defaults --datadir=$D/db --socket=$D/sock --port=3399 --bind-address=127.0.0.1 --user=root &
|
||||
mariadb --no-defaults -S $D/sock -e "CREATE DATABASE pulse; CREATE USER 'pulse'@'127.0.0.1' IDENTIFIED BY 'demo-only'; GRANT ALL ON pulse.* TO 'pulse'@'127.0.0.1'"
|
||||
printf 'PORT=8400\nHOST=127.0.0.1\nDB_HOST=127.0.0.1\nDB_PORT=3399\nDB_NAME=pulse\nDB_USER=pulse\nDB_PASSWORD=demo-only\nWORKER_API_KEY=demo-worker-key\n' > .env
|
||||
node server.js &
|
||||
for n in 1 2 3; do (cd worker && WORKER_NAME=node-0$n PULSE_SERVER=http://127.0.0.1:8400 PULSE_WS=ws://127.0.0.1:8400 WORKER_API_KEY=demo-worker-key node worker.js &); done
|
||||
python3 scripts/demo/seed.py # workflows, executions, schedules (via the real API)
|
||||
python3 scripts/demo/capture.py # screenshots + GIF -> docs/img/
|
||||
```
|
||||
|
||||
The server trusts Authelia `Remote-User` headers, so the scripts send them directly. For the demo only, raise `apiLimiter.max`
|
||||
in `server.js` (the screenshot run exceeds 300 requests / 15 min). Requires `playwright` and `pillow`.
|
||||
@@ -0,0 +1,182 @@
|
||||
"""Regenerates docs/img/* from a LOCAL demo instance (fictional data). See scripts/demo/README.md."""
|
||||
|
||||
# flake8: noqa: E501 (demo data and canned output contain long literal lines)
|
||||
|
||||
import asyncio
|
||||
import io
|
||||
from PIL import Image
|
||||
from playwright.async_api import async_playwright
|
||||
import pathlib
|
||||
|
||||
OUT = str(pathlib.Path(__file__).resolve().parents[2] / "docs" / "img")
|
||||
B = "http://127.0.0.1:8400"
|
||||
H = {
|
||||
"Remote-User": "alex.rivera",
|
||||
"Remote-Name": "Alex Rivera",
|
||||
"Remote-Email": "alex@example.com",
|
||||
"Remote-Groups": "admin,employee",
|
||||
}
|
||||
TOP = (
|
||||
"document.activeElement&&document.activeElement.blur();"
|
||||
"document.documentElement.style.scrollBehavior='auto';window.scrollTo(0,0)"
|
||||
)
|
||||
|
||||
|
||||
async def mk(p, theme="dark", w=1440, h=900):
|
||||
b = await p.chromium.launch()
|
||||
c = await b.new_context(viewport={"width": w, "height": h}, extra_http_headers=H)
|
||||
await c.add_init_script(
|
||||
f"try{{localStorage.setItem('lt_theme','{theme}');sessionStorage.setItem('lt_boot','1')}}catch(e){{}}"
|
||||
)
|
||||
return b, c
|
||||
|
||||
|
||||
async def go(pg, path, wait=3500):
|
||||
await pg.goto(B + path)
|
||||
await pg.wait_for_timeout(wait)
|
||||
await pg.evaluate(TOP)
|
||||
await pg.wait_for_timeout(300)
|
||||
|
||||
|
||||
async def shots():
|
||||
async with async_playwright() as p:
|
||||
b, c = await mk(p)
|
||||
pg = await c.new_page()
|
||||
errs = []
|
||||
pg.on("pageerror", lambda e: errs.append(str(e)[:100]))
|
||||
for path, name in (
|
||||
("/", "dashboard"),
|
||||
("/workers", "workers"),
|
||||
("/workflows", "workflows"),
|
||||
("/executions", "executions"),
|
||||
("/quick", "quick-command"),
|
||||
("/scheduler", "scheduler"),
|
||||
):
|
||||
await go(pg, path, 5000 if path == "/" else 3500)
|
||||
await pg.screenshot(path=f"{OUT}/{name}.png")
|
||||
print(name, "ok")
|
||||
print("errors", errs[:3])
|
||||
await b.close()
|
||||
|
||||
|
||||
async def open_exec(pg, text, idx=0, wait=1500):
|
||||
await go(pg, "/executions", 3000)
|
||||
await pg.locator("tr", has_text=text).nth(idx).click()
|
||||
await pg.wait_for_timeout(wait)
|
||||
|
||||
|
||||
async def scroll_modal(pg, y):
|
||||
await pg.evaluate(
|
||||
"(y)=>{const m=[...document.querySelectorAll('*')].filter(e=>e.scrollHeight>e.clientHeight+40&&['auto','scroll'].includes(getComputedStyle(e).overflowY)&&e.clientHeight>200); m.forEach(e=>e.scrollTop=y)}",
|
||||
y,
|
||||
)
|
||||
await pg.wait_for_timeout(400)
|
||||
|
||||
|
||||
async def shots2():
|
||||
async with async_playwright() as p:
|
||||
b, c = await mk(p)
|
||||
pg = await c.new_page()
|
||||
await open_exec(pg, "Rolling restart", 1)
|
||||
await pg.screenshot(path=f"{OUT}/execution-detail.png")
|
||||
await scroll_modal(pg, 700)
|
||||
await pg.screenshot(path=f"{OUT}/execution-detail-steps.png")
|
||||
await open_exec(pg, "Check backup mount", 0)
|
||||
await pg.screenshot(path=f"{OUT}/execution-failed.png")
|
||||
await open_exec(pg, "Rolling restart", 0, 2500)
|
||||
await pg.screenshot(path=f"{OUT}/execution-awaiting-approval.png")
|
||||
# workflow editor
|
||||
await go(pg, "/workflows", 3000)
|
||||
await pg.locator("tr", has_text="Rolling restart").get_by_text("EDIT", exact=False).first.click()
|
||||
await pg.wait_for_timeout(1200)
|
||||
await pg.screenshot(path=f"{OUT}/workflow-editor.png")
|
||||
# quick command with output
|
||||
await go(pg, "/quick", 3000)
|
||||
try:
|
||||
await pg.get_by_text("Multiple Workers").click()
|
||||
await pg.wait_for_timeout(300)
|
||||
await pg.get_by_text("SELECT ALL", exact=False).first.click()
|
||||
await pg.wait_for_timeout(300)
|
||||
except Exception as e:
|
||||
print("multi fail", str(e)[:60])
|
||||
await pg.locator("textarea").first.fill('echo "== $(hostname) =="; uptime; df -h / | tail -1')
|
||||
await pg.wait_for_timeout(300)
|
||||
await pg.screenshot(path=f"{OUT}/quick-command.png")
|
||||
await pg.keyboard.press("Control+Enter")
|
||||
await pg.wait_for_timeout(3500)
|
||||
await pg.screenshot(path=f"{OUT}/quick-command-result.png")
|
||||
# palette
|
||||
await go(pg, "/", 3500)
|
||||
await pg.keyboard.press("Control+k")
|
||||
await pg.wait_for_timeout(500)
|
||||
await pg.keyboard.type("exec", delay=70)
|
||||
await pg.wait_for_timeout(700)
|
||||
await pg.screenshot(path=f"{OUT}/command-palette.png", clip={"x": 320, "y": 60, "width": 800, "height": 460})
|
||||
await b.close()
|
||||
b, c = await mk(p, "light")
|
||||
pg = await c.new_page()
|
||||
await go(pg, "/", 4500)
|
||||
await pg.screenshot(path=f"{OUT}/dashboard-light.png")
|
||||
await b.close()
|
||||
|
||||
|
||||
async def gif():
|
||||
frames = []
|
||||
async with async_playwright() as p:
|
||||
b, c = await mk(p, "dark", 1280, 720)
|
||||
pg = await c.new_page()
|
||||
|
||||
async def snap(n=1, ms=0):
|
||||
if ms:
|
||||
await pg.wait_for_timeout(ms)
|
||||
im = Image.open(io.BytesIO(await pg.screenshot())).convert("RGB").resize((960, 540), Image.LANCZOS)
|
||||
for _ in range(n):
|
||||
frames.append(im)
|
||||
|
||||
await go(pg, "/", 4500)
|
||||
await snap(3)
|
||||
await go(pg, "/workers", 2500)
|
||||
await snap(2)
|
||||
await go(pg, "/workflows", 2500)
|
||||
await snap(2)
|
||||
await go(pg, "/executions", 3500)
|
||||
await pg.locator("tr", has_text="Rolling restart").nth(1).wait_for()
|
||||
await snap(2)
|
||||
await pg.locator("tr", has_text="Rolling restart").nth(1).click()
|
||||
await snap(3, 1200)
|
||||
await pg.evaluate(
|
||||
"(y)=>{[...document.querySelectorAll('*')].filter(e=>e.scrollHeight>e.clientHeight+40&&['auto','scroll'].includes(getComputedStyle(e).overflowY)&&e.clientHeight>200).forEach(e=>e.scrollTop=y)}",
|
||||
700,
|
||||
)
|
||||
await snap(3, 500)
|
||||
await go(pg, "/executions", 3500)
|
||||
await pg.locator("tr", has_text="Rolling restart").nth(1).wait_for()
|
||||
await pg.locator("tr", has_text="Rolling restart").nth(0).click()
|
||||
await snap(3, 1800)
|
||||
await go(pg, "/quick", 2500)
|
||||
await pg.get_by_text("Multiple Workers").click()
|
||||
await pg.get_by_text("SELECT ALL", exact=False).first.click()
|
||||
await pg.locator("textarea").first.fill("uptime; df -h / | tail -1")
|
||||
await snap(2, 300)
|
||||
await pg.keyboard.press("Control+Enter")
|
||||
await snap(3, 2500)
|
||||
await go(pg, "/", 2500)
|
||||
await pg.keyboard.press("Control+k")
|
||||
await snap(2, 500)
|
||||
for ch in "exec":
|
||||
await pg.keyboard.type(ch)
|
||||
await snap(1, 200)
|
||||
await snap(3, 400)
|
||||
await b.close()
|
||||
pal = [f.quantize(colors=96, method=Image.MEDIANCUT, dither=Image.NONE) for f in frames]
|
||||
pal[0].save(OUT + "/demo.gif", save_all=True, append_images=pal[1:], duration=1100, loop=0, optimize=True)
|
||||
print("gif", len(frames))
|
||||
|
||||
|
||||
async def main():
|
||||
await shots()
|
||||
await shots2()
|
||||
await gif()
|
||||
|
||||
|
||||
asyncio.run(main())
|
||||
@@ -0,0 +1,197 @@
|
||||
# flake8: noqa: E501 (demo data and canned output contain long literal lines)
|
||||
import json
|
||||
import time
|
||||
import urllib.request
|
||||
|
||||
B = "http://127.0.0.1:8400"
|
||||
H = {
|
||||
"Remote-User": "alex.rivera",
|
||||
"Remote-Name": "Alex Rivera",
|
||||
"Remote-Email": "alex@example.com",
|
||||
"Remote-Groups": "admin,employee",
|
||||
"Content-Type": "application/json",
|
||||
}
|
||||
|
||||
|
||||
def call(method, path, body=None):
|
||||
req = urllib.request.Request(
|
||||
B + path, method=method, headers=H, data=json.dumps(body).encode() if body is not None else None
|
||||
)
|
||||
try:
|
||||
with urllib.request.urlopen(req, timeout=30) as r:
|
||||
return json.loads(r.read() or b"{}")
|
||||
except urllib.error.HTTPError as e:
|
||||
return {"error": e.code, "body": e.read().decode()[:200]}
|
||||
|
||||
|
||||
workers = {w["name"]: w["id"] for w in call("GET", "/api/workers")}
|
||||
print(workers)
|
||||
wf = {}
|
||||
wf["health"] = call(
|
||||
"POST",
|
||||
"/api/workflows",
|
||||
{
|
||||
"name": "Node health check",
|
||||
"description": "Uptime, root filesystem and memory on every online worker",
|
||||
"definition": {
|
||||
"steps": [
|
||||
{"id": "uptime", "name": "Uptime", "type": "execute", "command": "uptime", "targets": ["all"]},
|
||||
{
|
||||
"id": "disk",
|
||||
"name": "Root filesystem usage",
|
||||
"type": "execute",
|
||||
"command": "df -h / | tail -1",
|
||||
"targets": ["all"],
|
||||
},
|
||||
{"id": "mem", "name": "Memory", "type": "execute", "command": "free -h | head -2", "targets": ["all"]},
|
||||
]
|
||||
},
|
||||
},
|
||||
)["id"]
|
||||
wf["restart"] = call(
|
||||
"POST",
|
||||
"/api/workflows",
|
||||
{
|
||||
"name": "Rolling restart (with approval)",
|
||||
"description": "Checks state, asks a human to approve, restarts one node, verifies",
|
||||
"definition": {
|
||||
"steps": [
|
||||
{
|
||||
"id": "check",
|
||||
"name": "Check cluster state",
|
||||
"type": "execute",
|
||||
"command": 'echo "all 3 nodes healthy"; date',
|
||||
"targets": ["all"],
|
||||
},
|
||||
{
|
||||
"id": "ask",
|
||||
"name": "Approve restart on node-01?",
|
||||
"type": "prompt",
|
||||
"message": "All nodes are healthy. Restart the service on node-01?",
|
||||
"options": ["Proceed", "Abort"],
|
||||
"key": "approval",
|
||||
"routes": {"Abort": "end"},
|
||||
},
|
||||
{
|
||||
"id": "restart",
|
||||
"name": "Restart service on node-01",
|
||||
"type": "execute",
|
||||
"command": 'echo "stopping service"; sleep 1; echo "starting service"; echo "service active"',
|
||||
"targets": ["node-01"],
|
||||
},
|
||||
{"id": "settle", "name": "Let it settle", "type": "wait", "duration": 2},
|
||||
{
|
||||
"id": "verify",
|
||||
"name": "Verify",
|
||||
"type": "execute",
|
||||
"command": 'echo "verification passed on $(hostname)"',
|
||||
"targets": ["node-01"],
|
||||
},
|
||||
]
|
||||
},
|
||||
},
|
||||
)["id"]
|
||||
wf["route"] = call(
|
||||
"POST",
|
||||
"/api/workflows",
|
||||
{
|
||||
"name": "Capacity check with routing",
|
||||
"description": "Parses KEY=VALUE output and routes on the result",
|
||||
"definition": {
|
||||
"steps": [
|
||||
{
|
||||
"id": "probe",
|
||||
"name": "Probe capacity",
|
||||
"type": "execute",
|
||||
"command": 'echo "FREE_PCT=12"; echo "POOL=app-nvme"',
|
||||
"targets": ["node-02"],
|
||||
},
|
||||
{"id": "parse", "name": "Parse results", "type": "parse"},
|
||||
{
|
||||
"id": "decide",
|
||||
"name": "Decide",
|
||||
"type": "route",
|
||||
"conditions": [
|
||||
{"if": "parseInt(state.free_pct) < 20", "goto": "warn", "label": "Low capacity"},
|
||||
{"default": True, "goto": "end"},
|
||||
],
|
||||
},
|
||||
{
|
||||
"id": "warn",
|
||||
"name": "Raise warning",
|
||||
"type": "execute",
|
||||
"command": 'echo "WARNING: pool is below 20% free"; exit 0',
|
||||
"targets": ["node-02"],
|
||||
},
|
||||
]
|
||||
},
|
||||
},
|
||||
)["id"]
|
||||
wf["broken"] = call(
|
||||
"POST",
|
||||
"/api/workflows",
|
||||
{
|
||||
"name": "Check backup mount",
|
||||
"description": "Fails when the backup mount is missing",
|
||||
"definition": {
|
||||
"steps": [
|
||||
{
|
||||
"id": "mnt",
|
||||
"name": "Backup mount present?",
|
||||
"type": "execute",
|
||||
"command": "ls /mnt/does-not-exist-backup",
|
||||
"targets": ["node-03"],
|
||||
}
|
||||
]
|
||||
},
|
||||
},
|
||||
)["id"]
|
||||
print(wf)
|
||||
|
||||
|
||||
def run(k, params=None):
|
||||
r = call("POST", "/api/executions", {"workflow_id": wf[k], "params": params or {}})
|
||||
time.sleep(1.2)
|
||||
return r
|
||||
|
||||
|
||||
# quick commands first (history)
|
||||
for name, cmd in (
|
||||
("node-01", "uname -a"),
|
||||
("node-02", "df -h /"),
|
||||
("node-03", "free -h"),
|
||||
("node-01", "uptime"),
|
||||
("node-02", "lscpu | head -8"),
|
||||
):
|
||||
call("POST", f"/api/workers/{workers[name]}/command", {"command": cmd})
|
||||
time.sleep(1.2)
|
||||
for k in ("health", "health", "route", "broken"):
|
||||
print(k, run(k).get("id"))
|
||||
time.sleep(3)
|
||||
r = run("restart")
|
||||
time.sleep(2.5)
|
||||
print("restart#1", r.get("id"))
|
||||
print(call("POST", f"/api/executions/{r['id']}/respond", {"response": "Proceed"}))
|
||||
time.sleep(6)
|
||||
run("health")
|
||||
time.sleep(3)
|
||||
r2 = run("restart")
|
||||
print("restart waiting", r2.get("id")) # left waiting for approval on purpose
|
||||
for name, cmd, t, v in (
|
||||
("Disk usage report", "df -h /", "interval", "60"),
|
||||
("Memory snapshot", "free -h | head -2", "hourly", "1"),
|
||||
("Uptime log", "uptime", "interval", "120"),
|
||||
):
|
||||
print(
|
||||
call(
|
||||
"POST",
|
||||
"/api/scheduled-commands",
|
||||
{
|
||||
"name": name,
|
||||
"command": cmd,
|
||||
"worker_ids": [workers["node-01"], workers["node-02"]],
|
||||
"schedule_type": t,
|
||||
"schedule_value": v,
|
||||
},
|
||||
)
|
||||
)
|
||||