diff --git a/CHANGELOG.md b/CHANGELOG.md index 968148e..8df101f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,3 +1,48 @@ +# v1.29.0 + +**Breaking:** theta-agent enrollment is now mandatory. Agents installed before this release carry a browser-generated token the server never recorded and will be rejected until re-enrolled. Requires theta-suite ≥ v1.42.0 (the `sso-broker` OpenBao policy must grant `secret/agent/*`); re-run `./setup.sh`. + +### Security — theta-agent channel + +- **sec: `/api/agent/ws` accepted any token.** There was no agent registry, so the endpoint authenticated nothing: any client that could reach the SSO could register as a node, publish discovery/telemetry into the admin view, and receive commands — including a signed `arbitrary_bash` — addressed to a token it guessed. Tokens were generated in the *browser* (`generateRandomHexToken`) and never recorded server-side, so there was nothing to validate against and no way to revoke one. Agents are now rows in a new `Agent` table, authenticated by SHA-256 token hash before the connection is registered or the welcome payload is sent; unknown or revoked tokens are closed with `4001` and audited. +- **sec: the command signing key was ephemeral.** `AgentManager` generated an Ed25519 pair in its constructor, so it changed on every process start and the `public_key` an agent pinned in `agent.yml` stopped matching immediately. The key now lives in OpenBao at `secret/agent/signing-key` and survives restarts. If it cannot be loaded the SSO **refuses** to send high-risk commands rather than signing with a key no agent has seen (`signingAvailable: false` on `GET /api/agent/nodes`). +- **sec: commands are addressed by agent id, not token.** A credential has no business in a URL, an access log or browser history. +- **sec: agent actions are audited.** Enroll, update, rotate, revoke, delete, every command (with `signed`), and every rejected connection are emitted as structured `"component":"agent"` log records carrying the acting user. + +### theta-agent — enrollment & resource binding + +- feat: `POST /api/agent/enroll` mints the token server-side and returns it **once**; only its SHA-256 is stored. Plus `PUT /nodes/:id` (rename/rebind), `POST /nodes/:id/rotate`, `POST /nodes/:id/revoke`, `DELETE /nodes/:id`. Rotate, revoke and delete drop the live socket immediately (`4004`/`4003`) instead of waiting for a reconnect. +- feat: an agent binds to a **host resource** (`resourceId`). The Directory reads that link instead of guessing by hostname — the old `agentsByHost[name]` match silently failed whenever a Directory name differed from the machine's hostname, and aliased two hosts that shared one. +- feat: **agent discovery reaches the Directory.** A bound agent's facts (`os`, `kernel`, `cpu`, `ram_total_gb`, `disk_total_gb`, `ip`) are written onto its host resource, tagged `discovery_sources: ["theta-agent"]` with an `agentId` back-reference. An unbound agent goes through the normal reconciler. Previously `handleDiscovery` wrote to an in-memory record and updated nothing — the one source actually running *on* the host contributed nothing to the directory. +- feat: agent state is persisted, so an agent that is installed but **offline** is now distinguishable from one that never existed; enrollments survive a restart. The Directory status dot reflects this: red means "enrolled and not connected" (a fault), grey means no agent enrolled / revoked / service unreachable. Red previously covered both, making an ordinary directory of hosts look like an outage. +- feat: the Install Agent modal enrolls first and builds the install command from the result, including `--public-key`. `public_key` was never emitted into the generated `agent.yml` before, so no installed agent could verify anything. +- fix: `registerAgent` is synchronous. Awaiting a database write before attaching the WebSocket `message` listener lost every agent's first `discovery` frame, which it sends the instant the socket opens (`ws` drops events emitted with no listener attached). + +### Directory + +- feat: **the resource tree is collapsible.** Any row with children has a caret; the toolbar collapses/expands everything. State persists per browser, so the shape survives the self-heal reload that follows most edits. An active search overrides collapse so matches inside a folded subtree are never hidden. +- fix: **the Proxmox plugin mismatched MAC addresses to IPs.** It collected MACs and IPs into two flat lists and zipped them by index, so on any multi-NIC guest — or any guest where one NIC had no address — the directory recorded an address against the wrong MAC. NICs are now keyed by MAC, so a pairing can only come from the source that observed both together. +- feat: Proxmox discovery emits an **endpoint resource** (named from `/cluster/status`) with every node parented beneath it, so one endpoint is one subtree instead of several orphan roots. It deliberately carries no IP: giving it the address it is reached at made the reconciler merge it with the node answering on that address, producing a resource that was its own parent. +- feat: discovered guests carry `sourceId` (`/qemu/`), `node`, `vmid` and `macAddress`, so a row traces back to the exact guest on the exact hypervisor. Against a live 3-node cluster this took MAC coverage to 53/54 resources and `sourceId` to 54/54. +- fix: Proxmox interfaces belonging to something running *inside* a guest (`docker0`, `veth*`, `br-*`, VPN tunnels) are filtered out — one Home Assistant VM reported 16 of them alongside its single real NIC, and their 172.x addresses gave the reconciler spurious matches. +- fix: a stopped VM still reports its MAC (read from the VM config), a DHCP-configured LXC gets its address from the running container's interface list, and Proxmox **nodes** report their own IP/MAC (recovered from `enx` predictable names, since `/nodes/*/network` carries no `hwaddr`). Offline nodes are recorded with `status` instead of skipped, so a hypervisor that is down no longer looks decommissioned and get garbage-collected after a week. +- fix: **the reconciler could make a resource its own parent.** Two slugs in one payload can resolve to the same row once merged; the resulting self-edge renders as an infinitely nested tree and defeats every ancestor walk in the app. Self-edges and cycle-closing edges are now refused and logged. +- fix: **hosts were named after their MAC address.** `bestName` preferred the *longer* name, so UniFi's `ac:16:2d:b3:da:80` (17 chars) beat Proxmox's real hostname `dl380-0` (7). Names are now ranked (hostname > IP > MAC) with length only as a tie-break within a rank. +- fix: `isIp` never matched anything — `\\.` inside a regex literal matches a backslash, not a dot — so an IP-shaped placeholder name was never replaced by a real hostname a later source discovered. +- fix: a discovered device can only merge into a resource of the same kind. A VM named `gitea-runner` could match a hand-created *service* of the same name on the name rule and overwrite it. +- perf: the reconciler reads the inventory once per run instead of once per incoming resource — a ~55-resource Proxmox payload against a similar-sized inventory was doing quadratic full-table reads every run. +- fix: the Discovered Inventory table showed "Unknown IP" for almost everything, because it read `metadata.ip` while any source that enumerates interfaces stores addresses per-NIC. It now falls back to the first NIC address, and shows `vmid`, slug, `sourceId` and per-interface MAC/name. + +### Profile + +- fix: the API Tokens card is no longer wider than every other card on the site — the section sat outside the page's `.container`. + +### Build & docs + +- fix: `Dockerfile.test-runner` never copied `nodejs/plugins`, so every plugin test suite failed in CI as "Cannot find module" and plugin code was effectively untested. Suite count goes 27 → 29. +- docs: `docs/agents.md` rewritten for enrollment, the close-code table, resource binding, the persistent signing key, and a corrected `public_key` example (the documented `MCowBQYDK2VwAyEA...` was an SPKI PEM body — 44 bytes decoded — where the agent requires the raw 32). +- docs: `docs/directory.md` covers the collapsible tree and the corrected seed hierarchy; `docs/plugins.md` documents what the Proxmox plugin produces and why the endpoint has no IP. + # v1.28.0 - fix: `/api/agent/nodes` no longer 404s — the previous "unconditional mount" was still inside the post-listen `onListen` hook, so the REST router landed *behind* app.js's terminal 404 catch-all and every `/api/agent/*` request 404'd. The router is now mounted synchronously in `app.js` before the 404 handler; only the agent WebSocket setup runs on `onListen`. - feat: promoting a discovered inventory resource now opens the resource form pre-filled with the discovered data (name, kind, IP, subtype, …) for review; the modal's Save confirms the promote (creates the LDAP groups + marks it managed) instead of silently promoting. diff --git a/Dockerfile.test-runner b/Dockerfile.test-runner index 32fd684..dc75fbb 100644 --- a/Dockerfile.test-runner +++ b/Dockerfile.test-runner @@ -22,6 +22,10 @@ COPY nodejs/conf ./conf COPY nodejs/controller ./controller COPY nodejs/middleware ./middleware COPY nodejs/models ./models +# Without this the discovery/plugin suites cannot even load their subject and +# fail as "Cannot find module ../plugins/discovery/..." -- plugin code was +# effectively untested in CI. +COPY nodejs/plugins ./plugins COPY nodejs/routes ./routes COPY nodejs/services ./services COPY nodejs/utils ./utils diff --git a/docs/agents.md b/docs/agents.md index c996ed4..40146db 100644 --- a/docs/agents.md +++ b/docs/agents.md @@ -10,6 +10,54 @@ The **Theta Agent** (`theta-agent`) is a unified, 2-way Command & Control (C2) e --- +## Enrollment (required) + +An agent is only real if the SSO issued its token. **Tokens the server did not +issue are rejected** at the WebSocket handshake. + +Enroll from **Directory → Install Agent**: + +1. Give the agent a name and, ideally, **bind it to a host resource**. The + binding is what links telemetry, status and commands to a Directory entry. +2. Press **Enroll & issue token**. The SSO mints a 256-bit token, stores only its + SHA-256, and shows the raw value **once**. +3. Copy the generated install command — it already carries the token and the + server's public key. + +Or via the API: + +```bash +curl -X POST https:///api/agent/enroll \ + -H "Authorization: Bearer " \ + -H 'Content-Type: application/json' \ + -d '{"name": "web01", "resourceId": ""}' +``` + +The response contains `token` (once only) and `publicKey`. + +| Endpoint | Purpose | +| :--- | :--- | +| `GET /api/agent/nodes` | Every enrolled agent, connected or not, plus the server public key | +| `POST /api/agent/enroll` | Mint an agent + token | +| `PUT /api/agent/nodes/:id` | Rename, or bind/unbind the host resource | +| `POST /api/agent/nodes/:id/rotate` | Issue a new token; the old one stops working immediately | +| `POST /api/agent/nodes/:id/revoke` | Disable the enrollment | +| `DELETE /api/agent/nodes/:id` | Remove the enrollment | +| `POST /api/agent/nodes/:id/command` | Send a command (signed automatically when high-risk) | + +Revoke, rotate and delete **drop any live connection immediately** — they do not +wait for the agent to reconnect. Commands are addressed by agent **id**, never by +token: a token is a credential and has no business in a URL or a log. + +Enrollment, revocation, rotation, every command, and every rejected connection +are written to the application log as structured `"component":"agent"` records +with the acting user. + +> **Lost the token?** It cannot be recovered — only its hash is stored. Rotate +> the agent to issue a new one. + +--- + ## Core Functionality ### 1. Host Discovery & Inventory @@ -41,12 +89,31 @@ Directory shows a status dot in the row: | :--- | :--- | | **Green** | Connected, healthy (CPU/RAM/disk within limits). | | **Yellow** | Connected but under high load (CPU > 80% or RAM > 80% or disk > 90%). | -| **Red** | Not connected (no agent, or the agent is offline). | +| **Red** | **Enrolled but not connected.** The agent exists and is expected — this is a fault. | +| **Grey** | No agent enrolled for this host, the enrollment is revoked, or the agent service is unreachable. | + +Red and grey used to be the same colour, which made an ordinary directory of +hosts look like an outage. Because the enrollment now outlives the connection, +"installed but down" is distinguishable from "never had an agent". Opening a host's resource modal reveals a **Metrics** tab with the agent's live telemetry (CPU/RAM/disk/ZFS/GPU) and discovery info (OS, kernel, IPs, location). -The agent is joined to its host by hostname (`agent.discovery.hostname` ↔ the -resource name), so name the Directory host the same as the machine's hostname. + +An agent attaches to its host by its **enrollment binding** (`resourceId`), set +when you enroll it or later via `PUT /api/agent/nodes/:id`. Agents enrolled +without a binding fall back to matching their reported hostname against the +resource name — the old behaviour, kept only as a fallback, because it silently +failed whenever a Directory name differed from the machine's hostname and +aliased two hosts that happened to share one. + +### Agent discovery feeds the Directory + +A bound agent's discovery payload is written onto its host resource (`os`, +`kernel`, `cpu`, `ram_total_gb`, `disk_total_gb`, `ip`), tagged with +`discovery_sources: ["theta-agent"]` and an `agentId` back-reference. An agent +runs *on* the host it describes, so it is the most authoritative source the +directory has. An unbound agent goes through the normal discovery reconciler +instead, matching like any other source. --- @@ -64,14 +131,31 @@ To protect hosts against unauthorized control, `theta-agent` enforces a **strict --- -## High-Risk Command Verification (Protocol v1.1.0) +## High-Risk Command Verification (Protocol v1.2.0) High-risk management commands (`reboot`, `service_restart`, `configure_ldap`, `arbitrary_bash`, `update_binary`) are cryptographically verified using **Ed25519 signatures**: -1. The SSO Manager canonicalizes the command payload (sorted keys, no whitespace). +1. The SSO Manager canonicalizes the command payload (sorted keys, no whitespace, + no HTML escaping, `signature` omitted). 2. The payload is signed with the SSO Manager's Ed25519 private key. 3. The Base64 signature is appended to the message payload. 4. The agent verifies the signature against the configured `public_key` in `/etc/theta42/agent.yml` before executing the action. +**The signing key is persistent.** It lives in OpenBao at +`secret/agent/signing-key` and survives restarts, so the `public_key` you pin in +`agent.yml` keeps matching. (It used to be generated in memory at boot and +changed on every restart, which made pinning impossible.) If the SSO cannot load +or store a key it **refuses** to send high-risk commands rather than signing with +one no agent has seen — `GET /api/agent/nodes` reports this as +`signingAvailable: false`. + +This requires the `sso-broker` OpenBao policy to grant `secret/agent/*`. Re-run +`./setup.sh` from theta-suite if you are upgrading. + +**Verification is fail-closed on the agent.** An agent with no `public_key` +configured rejects every high-risk command. Earlier versions logged "skipping +signature verification" and executed them, so an agent installed without a key +would run `reboot`, `configure_ldap` and `arbitrary_bash` unverified. + --- ## Installation & Deployment @@ -80,9 +164,14 @@ High-risk management commands (`reboot`, `service_restart`, `configure_ldap`, `a Run the following command as `root` on the target Linux host: ```bash -curl -fsSL https:///resources/theta-agent/install.sh | sh -s -- --url "https://" --token "" +curl -fsSL https:///resources/theta-agent/install.sh | sh -s -- \ + --url "https://" --token "" --public-key "" ``` +Both values come from enrollment. The **Install Agent** modal builds this line +for you with them already filled in. Omitting `--public-key` leaves the agent +able to report telemetry but unable to accept any high-risk command. + ### Custom Config Wizard You can generate a Base64-encoded custom configuration using the **Install Agent** button on the **Directory Management** page in the SSO Manager UI: @@ -97,9 +186,14 @@ curl -fsSL https:///resources/theta-agent/install.sh | sh -s -- "Directory & inventory list view @@ -102,9 +114,17 @@ You don't have to build the graph by hand — the theta42 tooling registers itse - a **site** (name from `CFG_SITE_NAME` in `setup.env`, default `local` → slug `site_local`) marked as the current site - the **host** the stack runs on (`host_`), with IP, MAC address, OS, and kernel collected from the machine -- the **services** it composes — SSO Manager, Proxy (management UI), OpenLDAP Directory (the LDAPS endpoint Linux hosts and LDAP-native apps bind to), and OpenResty Edge (the 80/443 data plane) — each with its address, internal port, and git repo +- the **hosts** for the proxy and jump host (`host_theta-proxy`, `host_theta-jump`) +- the **services** it composes — SSO Manager, Proxy (management UI), OpenLDAP Directory (the LDAPS endpoint Linux hosts and LDAP-native apps bind to), OpenResty Edge (the 80/443 data plane), and the SSH Jump Host — each with its address, internal port, and git repo - the proxy's auto-registered **OAuth client**, linked under its service +Services are parented to the host that actually runs them: Proxy and OpenResty +Edge under `host_theta-proxy`, the SSH Jump Host under `host_theta-jump`, and the +rest under the stack host. Installs seeded before this was fixed had all of them +under the stack host, leaving the two purpose-made host resources childless; the +seed re-parents those on its next run, and only when the current parent is the +one the old code set, so a layout you arranged deliberately is left alone. + The seed is idempotent and non-destructive: a resource whose slug already exists is considered operator-owned — the seed only fills in metadata fields you haven't set, and never overwrites your values. ### Linux hosts (ldap-client) diff --git a/docs/plugins.md b/docs/plugins.md index 4e27295..1070fa4 100644 --- a/docs/plugins.md +++ b/docs/plugins.md @@ -23,6 +23,44 @@ filename basename (without `.js`) is the `type`; the parent directory is the - `unifi` — UniFi Network controller (URL + username/password) - `nmap` — nmap OS + port scan (a target range; no credentials) +### What the Proxmox plugin produces + +One endpoint becomes one subtree: + +``` +Proxmox endpoint (cluster name, or the endpoint hostname) +└── node (hypervisor) + ├── VM / template + └── LXC / template +``` + +The endpoint resource stands for the cluster, not a machine, so it carries the +API URL and a `sourceId` but deliberately no IP — giving it the address it is +reached at made the reconciler merge it with the node answering on that address, +which produced a resource that was its own parent. + +Every guest carries: + +- `interfaces[]` — one entry per NIC with its own `mac`, `ip`/`ips` and `name`. + The MAC and the address on it are read from the same source, so they cannot be + mismatched (an earlier version collected MACs and IPs into two flat lists and + zipped them by index, which attributed addresses to the wrong NIC on any + multi-NIC guest). +- `macAddress` / `ip` — the primary NIC's values, preferring one that actually + has an address. +- `vmid`, `node` and `sourceId` (`/qemu/` or `/lxc/`), so + a directory row traces back to the exact guest on the exact node. + +Interfaces belonging to something running *inside* a guest — `docker0`, `veth*`, +`br-*`, VPN tunnels — are filtered out. They are not NICs of the host, and their +172.x addresses would otherwise give the reconciler spurious matches. + +A stopped VM still reports its MAC (read from the VM config rather than the +guest agent), and a DHCP-configured LXC gets its address from the running +container's interface list. Offline nodes are recorded with `status` rather than +skipped, so a hypervisor that is down does not look decommissioned and get +garbage-collected after a week. + A module exports a **manifest**: ```javascript diff --git a/nodejs/models/agent.js b/nodejs/models/agent.js new file mode 100644 index 0000000..e905197 --- /dev/null +++ b/nodejs/models/agent.js @@ -0,0 +1,111 @@ +'use strict'; + +const crypto = require('crypto'); +const { Model } = require('@simpleworkjs/orm'); + +// A theta-agent enrolled against this SSO. +// +// Before this model existed the "agent token" was generated in the browser and +// never recorded anywhere, so the server had no way to tell an agent it issued +// from one someone invented -- /api/agent/ws accepted any string, and there was +// no way to revoke a token or to know that an agent existed while it was +// offline. The row is now the authority: an agent is only real if it is here. +// +// The raw token is shown exactly once, at enrollment. Only its SHA-256 lands in +// the database, so a database disclosure does not hand over working agent +// credentials. `tokenPrefix` is the first 8 characters, kept in the clear so the +// UI and logs can identify an agent without holding the secret. +class Agent extends Model { + // Tokens are compared by hash on every WebSocket connect. SHA-256 (not + // bcrypt) is deliberate: this runs on the connection path and the token is a + // 256-bit random value, not a human-chosen password, so there is nothing for + // a slow KDF to protect against here. + static hashToken(raw) { + return crypto.createHash('sha256').update(String(raw || ''), 'utf8').digest('hex'); + } + + static generateToken() { + return crypto.randomBytes(32).toString('hex'); + } + + // Resolve a presented token to its (non-revoked) agent, or null. Every + // caller that authenticates an agent must go through here. + static async authenticate(rawToken) { + if (!rawToken || typeof rawToken !== 'string') return null; + const tokenHash = this.hashToken(rawToken); + const matches = await this.list({ where: { tokenHash } }); + const agent = matches && matches[0]; + if (!agent) return null; + if (agent.revoked) return null; + return agent; + } + + // Enroll a new agent and return { agent, token }. The caller is responsible + // for showing `token` to the operator once and never storing it. + static async enroll({ name, resourceId, enrolledBy, description }) { + const token = this.generateToken(); + const agent = await this.create({ + id: crypto.randomUUID(), + name: name || 'theta-agent', + description: description || null, + tokenHash: this.hashToken(token), + tokenPrefix: token.slice(0, 8), + resourceId: resourceId || null, + revoked: false, + enrolled_by: enrolledBy || null, + enrolled_on: Math.floor(Date.now() / 1000) + }); + return { agent, token }; + } + + // Issue a fresh token for an existing agent, invalidating the old one. + async rotateToken() { + const token = Agent.generateToken(); + await this.update({ + tokenHash: Agent.hashToken(token), + tokenPrefix: token.slice(0, 8), + revoked: false + }); + return token; + } + + static fields = { + id: { type: 'uuid', primaryKey: true }, + name: { type: 'string', isRequired: true }, + description: { type: 'text' }, + // Never the raw token. See hashToken above. + tokenHash: { type: 'string', isRequired: true }, + tokenPrefix: { type: 'string' }, + // The host this agent runs on. Nullable so an agent can be enrolled + // before its host exists in the Directory, but the UI pushes for it: + // without this link there is nothing to hang resource control off, and + // the old code had to guess by matching hostnames to slugs. + resource: { type: 'hasOne', model: 'Resource' }, // creates resourceId + revoked: { type: 'boolean', default: false }, + enrolled_by: { type: 'string' }, + enrolled_on: { type: 'integer' }, + // Survives a restart, which the in-memory map did not: an agent that is + // installed but currently down is now distinguishable from one that was + // never enrolled. + last_seen: { type: 'integer' }, + last_ip: { type: 'string' }, + lastDiscovery: { type: 'json', default: {} }, + lastTelemetry: { type: 'json', default: {} } + }; + + // The shape the admin API returns. Never includes tokenHash. + toPublic(liveState) { + const data = this.toJSON ? this.toJSON() : { ...this }; + delete data.tokenHash; + return { + ...data, + connected: !!(liveState && liveState.connected), + // "Online" is a live-connection fact, not a stored one. A row with a + // last_seen from an hour ago is an installed agent that is down. + isOnline: !!(liveState && liveState.connected), + lastResponse: (liveState && liveState.lastResponse) || null + }; + } +} + +module.exports = { Agent }; diff --git a/nodejs/models/index.js b/nodejs/models/index.js index 678378f..087a4be 100644 --- a/nodejs/models/index.js +++ b/nodejs/models/index.js @@ -20,6 +20,7 @@ const { PluginInstance } = require('./plugin_instance'); const { SharedSecret } = require('./shared_secret'); const { SharedSecretGrant } = require('./shared_secret_grant'); const { VaultAppToken } = require('./vault_app_token'); +const { Agent } = require('./agent'); async function initORM() { const ormConf = conf.orm || { dialect: 'sqlite', @@ -34,7 +35,7 @@ async function initORM() { conf: { orm: ormConf }, models: [ Resource, ResourceEdge, ResourceGroup, AccessRequest, Webhook, PluginInstance, - SharedSecret, SharedSecretGrant, VaultAppToken, + SharedSecret, SharedSecretGrant, VaultAppToken, Agent, Token, AuthToken, InviteToken, ImpersonationToken, PasswordResetToken, OtpToken, ServiceToken ] }); diff --git a/nodejs/package-lock.json b/nodejs/package-lock.json index c2adaad..1d32b26 100644 --- a/nodejs/package-lock.json +++ b/nodejs/package-lock.json @@ -1,12 +1,12 @@ { "name": "t42-sso-manager", - "version": "1.28.0", + "version": "1.29.0", "lockfileVersion": 3, "requires": true, "packages": { "": { "name": "t42-sso-manager", - "version": "1.28.0", + "version": "1.29.0", "license": "MIT", "dependencies": { "@fortawesome/fontawesome-free": "^7.3.0", diff --git a/nodejs/package.json b/nodejs/package.json index d9fda0c..808bffc 100755 --- a/nodejs/package.json +++ b/nodejs/package.json @@ -1,6 +1,6 @@ { "name": "t42-sso-manager", - "version": "1.28.0", + "version": "1.29.0", "description": "A very simple LDAP management and SSO system", "author": [ { diff --git a/nodejs/plugins/discovery/proxmox.js b/nodejs/plugins/discovery/proxmox.js index e670116..8129b34 100644 --- a/nodejs/plugins/discovery/proxmox.js +++ b/nodejs/plugins/discovery/proxmox.js @@ -6,6 +6,81 @@ const agent = new https.Agent({ rejectUnauthorized: false }); +// Accumulates a guest's NICs, keyed by MAC, merging what several Proxmox +// endpoints each know a piece of: the guest agent knows MAC+IP together, the +// VM/LXC config knows the MAC even while the guest is stopped, and the LXC +// interfaces endpoint knows the DHCP-assigned IP. Keying by MAC is what keeps +// the pairing honest -- the previous code collected MACs and IPs into two flat +// lists and zipped them by index, which mismatched them on any multi-NIC guest. +class Interfaces { + constructor() { this.byMac = new Map(); this.anonymous = []; } + + // Interfaces that belong to something running INSIDE the guest -- container + // engines, overlay networks, VPNs -- rather than to the guest itself. A + // Home Assistant VM reported 16 of these (docker0, hassio, 14x veth*) + // alongside its one real NIC, which is noise in the directory and, worse, + // gives the reconciler a pile of 172.x addresses to match unrelated hosts on. + // Only applied to guests; a hypervisor's own bridges are how you reach it. + static VIRTUAL_IFACE_RE = /^(lo|docker\d*|hassio|veth|br-|virbr|tap|fwbr|fwln|fwpr|cni|flannel|cali|kube|weave|zt|tailscale|wg|tun|utun)/i; + + static isVirtualName(name) { + return !!name && Interfaces.VIRTUAL_IFACE_RE.test(name); + } + + // A udev "predictable" name of the form enx<12 hex> encodes the MAC. It is + // the only place the Proxmox node network API exposes a physical NIC's MAC + // (/nodes/{node}/network carries no hwaddr field at all), so parse it out + // rather than leaving every hypervisor MAC-less. + static macFromIfaceName(name) { + const m = /^enx([0-9a-f]{12})$/i.exec(name || ''); + if (!m) return null; + return m[1].toLowerCase().match(/.{2}/g).join(':'); + } + + static normalizeMac(mac) { + const m = (mac || '').toLowerCase().trim(); + if (!/^([0-9a-f]{2}:){5}[0-9a-f]{2}$/.test(m)) return null; + if (m === '00:00:00:00:00:00') return null; + return m; + } + + // `ips` are the addresses observed on this one NIC (may be empty for a + // stopped guest, where only the MAC is known). + add(mac, ips, name) { + const key = Interfaces.normalizeMac(mac); + const addrs = (ips || []).filter(Boolean); + if (!key) { + // An IP with no usable MAC is still worth keeping; a NIC with neither is not. + if (addrs.length) this.anonymous.push({ mac: null, ip: addrs[0], ips: addrs, name: name || null }); + return; + } + const existing = this.byMac.get(key); + if (existing) { + for (const ip of addrs) if (!existing.ips.includes(ip)) existing.ips.push(ip); + existing.ip = existing.ips[0] || null; + if (!existing.name && name) existing.name = name; + return; + } + this.byMac.set(key, { mac: key, ip: addrs[0] || null, ips: addrs, name: name || null }); + } + + toArray() { return [...this.byMac.values(), ...this.anonymous]; } + + // The address/MAC the directory shows in its single-value columns, and what + // the reconciler matches on. Prefer a NIC that actually has an address. + primaryIp() { + const withIp = this.toArray().find(i => i.ip); + return withIp ? withIp.ip : null; + } + + primaryMac() { + const withIp = this.toArray().find(i => i.ip && i.mac); + if (withIp) return withIp.mac; + const first = this.toArray().find(i => i.mac); + return first ? first.mac : null; + } +} + module.exports = { // Plugin manifest — see nodejs/services/plugin_registry.js. `configSchema` // drives the admin UI form and validation; fields flagged `secret:true` are @@ -50,6 +125,45 @@ module.exports = { const resources = []; const edges = []; + // 0. The Proxmox endpoint itself. Without it a multi-node cluster produces + // several unrelated roots in the Directory tree and nothing says where any + // of them came from. Every node discovered below is parented to this, so + // one endpoint == one subtree. + const endpointHost = (() => { + try { return new URL(url).hostname; } catch (e) { return url.replace(/^https?:\/\//, '').split('/')[0]; } + })(); + const clusterName = await (async () => { + // /cluster/status names the cluster when one exists; a standalone node + // has no cluster entry, in which case the endpoint hostname is the name. + try { + const res = await fetch(`${url}/api2/json/cluster/status`, { headers, agent }); + if (!res.ok) return null; + const entry = ((await res.json()).data || []).find(d => d.type === 'cluster'); + return entry ? entry.name : null; + } catch (e) { return null; } + })(); + + const endpointSlug = `pve-${endpointHost.toLowerCase().replace(/[^a-z0-9]+/g, '-').replace(/^-|-$/g, '')}`; + resources.push({ + kind: 'host', + name: clusterName || `Proxmox (${endpointHost})`, + slug: endpointSlug, + metadata: { + subType: 'proxmox', + address: url, + os: 'Proxmox VE', + isProduction: true, + sourceId: url, + // Deliberately NO `ip`/`interfaces`: this resource stands for the + // cluster (the API endpoint), not for a machine. Giving it the address + // it is reached at made the reconciler match it to the very node that + // answers on that address -- the endpoint and the node collapsed into + // one row, which then became its own parent. The cluster is identified + // by slug + sourceId instead, which nothing else can collide with. + interfaces: [] + } + }); + // 1. Get Nodes const resNodes = await fetch(`${url}/api2/json/nodes`, { headers, agent }); if(!resNodes.ok) { @@ -59,9 +173,57 @@ module.exports = { const nodes = (await resNodes.json()).data; for (const node of nodes) { - if (node.status !== 'online') continue; - + // An offline node is still a real hypervisor that belongs in the + // directory -- skipping it entirely used to make it look decommissioned + // and let the reconciler's garbage collector archive it after a week of + // downtime. Record it, mark it down, and skip only the guest enumeration + // (which needs the node to answer). + const online = node.status === 'online'; + const nodeSlug = `pve-node-${node.node}`; + + // A hypervisor with no address is not actionable. Read its bridges/NICs + // so the node lands in the directory reachable and MAC-identified like + // any other host. Unlike a guest, a node's bridges are kept: vmbrN is + // normally the address you actually reach the hypervisor on. + const nodeIfaces = new Interfaces(); + try { + const netRes = online + ? await fetch(`${url}/api2/json/nodes/${node.node}/network`, { headers, agent }) + : { ok: false }; + if (netRes.ok) { + const ifaceList = (await netRes.json()).data || []; + for (const iface of ifaceList) { + if (iface.iface === 'lo') continue; + const ip = iface.address || iface.cidr; + // This endpoint has no hwaddr field, so the MAC has to be recovered + // from a predictable interface name -- either this interface's own + // or, for a bridge, one of the physical ports beneath it. + let mac = Interfaces.macFromIfaceName(iface.iface); + if (!mac) { + for (const alt of (iface.altnames || [])) { + mac = Interfaces.macFromIfaceName(alt); + if (mac) break; + } + } + if (!mac && iface.bridge_ports) { + for (const port of String(iface.bridge_ports).split(/\s+/).filter(Boolean)) { + mac = Interfaces.macFromIfaceName(port); + if (mac) break; + // The port may itself only carry the MAC in an altname. + const portDef = ifaceList.find(i => i.iface === port); + for (const alt of ((portDef && portDef.altnames) || [])) { + mac = Interfaces.macFromIfaceName(alt); + if (mac) break; + } + if (mac) break; + } + } + nodeIfaces.add(mac || iface.hwaddr, ip ? [String(ip).split('/')[0]] : [], iface.iface); + } + } + } catch (e) {} + resources.push({ kind: 'host', name: node.node, @@ -70,9 +232,19 @@ module.exports = { subType: 'hypervisor', os: 'Proxmox VE', isProduction: true, - interfaces: [] + status: node.status, + sourceId: `${node.node}`, + node: node.node, + interfaces: nodeIfaces.toArray(), + macAddress: nodeIfaces.primaryMac(), + ip: nodeIfaces.primaryIp() } }); + edges.push({ parentSlug: endpointSlug, childSlug: nodeSlug, relation: 'hosts' }); + + // Everything below asks the node itself; an offline node answers none of + // it, and its guests are already recorded from previous runs. + if (!online) continue; // 2. Get VMs for this node const resVms = await fetch(`${url}/api2/json/nodes/${node.node}/qemu`, { headers, agent }); @@ -82,10 +254,12 @@ module.exports = { const vmSlug = `vm-${vm.vmid}`; const isTemplate = vm.template === 1; - let ips = []; - let macs = []; - - // Enrich from QEMU guest agent if running + const ifaces = new Interfaces(); + + // Enrich from QEMU guest agent if running. The agent is the only source + // that knows which IP sits on which NIC, so pair them here rather than + // accumulating two flat lists (zipping those by index attributed IPs to + // the wrong MAC on any guest with more than one NIC). if (vm.status === 'running') { try { const agentRes = await fetch(`${url}/api2/json/nodes/${node.node}/qemu/${vm.vmid}/agent/network-get-interfaces`, { headers, agent }); @@ -93,35 +267,35 @@ module.exports = { const agentData = (await agentRes.json()).data; if (agentData && agentData.result) { for (const iface of agentData.result) { - if (iface['hardware-address'] && iface['hardware-address'] !== '00:00:00:00:00:00') macs.push(iface['hardware-address']); - if (iface['ip-addresses']) { - for (const ip of iface['ip-addresses']) { - if (ip['ip-address-type'] === 'ipv4' && ip['ip-address'] !== '127.0.0.1') { - ips.push(ip['ip-address']); - } - } - } + // Docker bridges, veth pairs and VPN tunnels are the + // guest's own plumbing, not NICs of the guest. + if (Interfaces.isVirtualName(iface.name)) continue; + const ips = (iface['ip-addresses'] || []) + .filter(ip => ip['ip-address-type'] === 'ipv4' && ip['ip-address'] !== '127.0.0.1') + .map(ip => ip['ip-address']); + ifaces.add(iface['hardware-address'], ips, iface.name); } } } } catch(e) {} } - - // Enrich from VM config to at least get MAC if agent failed/stopped + + // Enrich from VM config: the MAC is declared there whether or not the + // guest agent answered, so a stopped VM still gets a stable identity. try { const configRes = await fetch(`${url}/api2/json/nodes/${node.node}/qemu/${vm.vmid}/config`, { headers, agent }); if (configRes.ok) { const confData = (await configRes.json()).data; for (let i = 0; i < 10; i++) { if (confData[`net${i}`]) { - const m = confData[`net${i}`].match(/(?:virtio|e1000|rtl8139|vmxnet3)=([0-9a-fA-F:]+)/); - if(m) macs.push(m[1].toLowerCase()); + const m = confData[`net${i}`].match(/(?:virtio|e1000e?|rtl8139|vmxnet3)=([0-9a-fA-F:]{17})/); + if(m) ifaces.add(m[1], [], `net${i}`); } } } } catch(e) {} - const interfaces = [...new Set(macs)].map((mac, i) => ({ mac, ip: ips[i] || null })); + const interfaces = ifaces.toArray(); resources.push({ kind: isTemplate ? 'template' : 'host', @@ -130,9 +304,14 @@ module.exports = { metadata: { subType: isTemplate ? 'template' : 'vm', vmid: vm.vmid, + // The Proxmox-side identity, so a resource can be traced back to the + // exact guest on the exact node it was discovered from. + sourceId: `${node.node}/qemu/${vm.vmid}`, + node: node.node, isProduction: vm.status === 'running', interfaces, - ip: ips[0] || null + macAddress: ifaces.primaryMac(), + ip: ifaces.primaryIp() } }); edges.push({ parentSlug: nodeSlug, childSlug: vmSlug, relation: 'hosts' }); @@ -146,26 +325,45 @@ module.exports = { const lxcSlug = `lxc-${lxc.vmid}`; const isTemplate = lxc.template === 1; - let ips = []; - let macs = []; - - // Enrich from LXC config + const ifaces = new Interfaces(); + + // Enrich from LXC config. Each netN line carries its own hwaddr and ip, + // so read them off the same line instead of into parallel lists. try { const configRes = await fetch(`${url}/api2/json/nodes/${node.node}/lxc/${lxc.vmid}/config`, { headers, agent }); if (configRes.ok) { const confData = (await configRes.json()).data; for (let i = 0; i < 10; i++) { - if (confData[`net${i}`]) { - const hwMatch = confData[`net${i}`].match(/hwaddr=([0-9a-fA-F:]+)/); - const ipMatch = confData[`net${i}`].match(/ip=([0-9\.]+)/); // Ignores dhcp - if(hwMatch) macs.push(hwMatch[1].toLowerCase()); - if(ipMatch) ips.push(ipMatch[1]); + const line = confData[`net${i}`]; + if (!line) continue; + const hwMatch = line.match(/hwaddr=([0-9a-fA-F:]{17})/); + // `ip=` is either a CIDR address or the literal `dhcp`/`manual`. + const ipMatch = line.match(/\bip=(\d+\.\d+\.\d+\.\d+)/); + const nameMatch = line.match(/\bname=([^,]+)/); + if (hwMatch || ipMatch) { + ifaces.add(hwMatch && hwMatch[1], ipMatch ? [ipMatch[1]] : [], nameMatch ? nameMatch[1] : `net${i}`); } } } } catch(e) {} - const interfaces = [...new Set(macs)].map((mac, i) => ({ mac, ip: ips[i] || null })); + // A DHCP-configured container has no IP in its config. Ask the running + // container's interface list so it lands in the directory addressable + // instead of as an IP-less row. + if (lxc.status === 'running' && !ifaces.primaryIp()) { + try { + const ifRes = await fetch(`${url}/api2/json/nodes/${node.node}/lxc/${lxc.vmid}/interfaces`, { headers, agent }); + if (ifRes.ok) { + for (const iface of ((await ifRes.json()).data || [])) { + if (Interfaces.isVirtualName(iface.name)) continue; + const ip = (iface.inet || '').split('/')[0]; + ifaces.add(iface.hwaddr, ip ? [ip] : [], iface.name); + } + } + } catch(e) {} + } + + const interfaces = ifaces.toArray(); resources.push({ kind: isTemplate ? 'template' : 'host', @@ -174,9 +372,12 @@ module.exports = { metadata: { subType: isTemplate ? 'template' : 'lxc', vmid: lxc.vmid, + sourceId: `${node.node}/lxc/${lxc.vmid}`, + node: node.node, isProduction: lxc.status === 'running', interfaces, - ip: ips[0] || null + macAddress: ifaces.primaryMac(), + ip: ifaces.primaryIp() } }); edges.push({ parentSlug: nodeSlug, childSlug: lxcSlug, relation: 'hosts' }); @@ -190,5 +391,8 @@ module.exports = { // `discover` as their implementation name for back-compat, and `run` is just // an alias. Referenced via module.exports (not `this`) so it survives being // detached and called as a bare function reference. - run: async (config) => module.exports.discover(config) + run: async (config) => module.exports.discover(config), + + // Exported for unit tests only -- not part of the plugin contract. + _Interfaces: Interfaces }; diff --git a/nodejs/routes/api_agent.js b/nodejs/routes/api_agent.js index 310d17f..30d5f47 100644 --- a/nodejs/routes/api_agent.js +++ b/nodejs/routes/api_agent.js @@ -4,9 +4,16 @@ const express = require('express'); const middleware = require('../middleware/auth'); const permission = require('../utils/permission'); const agentManager = require('../utils/agent_manager'); +const agentKeys = require('../utils/agent_keys'); +const { Agent } = require('../models/agent'); const ADMIN_GROUPS = ['app_sso_admin', 'app_super_admin', 'app_sso_directory_admin']; +// Commands that can change or run code on the host. They are signed with the +// SSO's persisted Ed25519 key and the agent verifies against the key pinned in +// its agent.yml. +const HIGH_RISK_COMMANDS = ['reboot', 'service_restart', 'configure_ldap', 'arbitrary_bash', 'update_binary']; + // ── REST API (mounted synchronously in app.js, BEFORE the 404 catch-all) ── // This is a plain Express Router exported directly so app.js can // `app.use('/api/agent', require('./routes/api_agent'))` at require time. It @@ -17,9 +24,21 @@ const ADMIN_GROUPS = ['app_sso_admin', 'app_super_admin', 'app_sso_directory_adm // the only part that needs the post-listen onListen hook. const router = express.Router(); -// The agent WebSocket (/api/agent/ws) is handled by the raw `wss` upgrade server -// in bin/www with its own ?token= auth — unaffected by the express middleware -// here. These REST routes are admin-facing, so they're auth + admin gated. +// Structured audit line for anything that reaches a host. The agent channel can +// run arbitrary bash, so "who told which host to do what" has to be recoverable +// after the fact; previously nothing was recorded at all. +function logAgentAudit(action, details) { + console.log(JSON.stringify({ + timestamp: new Date().toISOString(), + component: 'agent', + action, + ...details + })); +} + +// The agent WebSocket (/api/agent/ws) authenticates its own token against the +// Agent table (see initAgentWebSockets). These REST routes are admin-facing, so +// they're auth + admin gated. router.use(middleware.auth); router.use(async (req, res, next) => { try { @@ -33,85 +52,246 @@ router.use(async (req, res, next) => { } }); -router.get('/nodes', (req, res) => { - res.json({ - status: 'ok', - agents: agentManager.getConnectedAgents(), - publicKey: agentManager.publicKeyPem - }); +// --- Fleet --- +router.get('/nodes', async (req, res, next) => { + try { + const keyStatus = agentKeys.status(); + res.json({ + status: 'ok', + agents: await agentManager.listAgents(), + // Base64 of the raw 32-byte key: what goes into agent.yml's `public_key`. + publicKey: await agentManager.publicKeyBase64(), + publicKeyPem: await agentManager.publicKeyPem(), + signingAvailable: agentKeys.status().available, + signingError: keyStatus.error || null + }); + } catch (err) { next(err); } }); -router.post('/nodes/:token/command', (req, res) => { - const { token } = req.params; - const { command, payload, isHighRisk } = req.body; +// --- Enrollment --- +// The token is minted HERE, not in the browser. It is returned exactly once; +// only its hash is stored, so it cannot be recovered afterwards -- rotate to +// get a new one. +router.post('/enroll', async (req, res, next) => { + try { + const { name, resourceId, description } = req.body || {}; + if (!name || !String(name).trim()) { + return res.status(400).json({ status: 'error', message: 'name is required' }); + } + if (resourceId) { + const { Resource } = require('../models/resource'); + const resource = await Resource.get(resourceId); + if (!resource) return res.status(400).json({ status: 'error', message: 'resourceId does not exist' }); + if (resource.kind !== 'host') { + return res.status(400).json({ status: 'error', message: 'an agent can only be bound to a host resource' }); + } + } + + const { agent, token } = await Agent.enroll({ + name: String(name).trim(), + description, + resourceId: resourceId || null, + enrolledBy: req.user.uid + }); + + logAgentAudit('enroll', { actor: req.user.uid, agentId: agent.id, agentName: agent.name, resourceId: resourceId || null }); + + const publicKey = await agentManager.publicKeyBase64(); + res.json({ + status: 'ok', + agent: agent.toPublic(agentManager.liveState(agent.id)), + // Shown once. The UI must make that clear. + token, + publicKey, + signingAvailable: agentKeys.status().available + }); + } catch (err) { next(err); } +}); + +router.put('/nodes/:id', async (req, res, next) => { + try { + const agent = await Agent.get(req.params.id); + if (!agent) return res.status(404).json({ status: 'error', message: 'agent not found' }); + + const patch = {}; + if (req.body.name !== undefined) patch.name = req.body.name; + if (req.body.description !== undefined) patch.description = req.body.description; + if (req.body.resourceId !== undefined) { + if (req.body.resourceId) { + const { Resource } = require('../models/resource'); + const resource = await Resource.get(req.body.resourceId); + if (!resource) return res.status(400).json({ status: 'error', message: 'resourceId does not exist' }); + if (resource.kind !== 'host') { + return res.status(400).json({ status: 'error', message: 'an agent can only be bound to a host resource' }); + } + } + patch.resourceId = req.body.resourceId || null; + } + + const updated = await agent.update(patch); + logAgentAudit('update', { actor: req.user.uid, agentId: agent.id, fields: Object.keys(patch) }); + res.json({ status: 'ok', agent: updated.toPublic(agentManager.liveState(agent.id)) }); + } catch (err) { next(err); } +}); + +// Revoke: the token stops authenticating immediately and any live socket is +// dropped, so revocation takes effect without waiting for a reconnect. +router.post('/nodes/:id/revoke', async (req, res, next) => { + try { + const agent = await Agent.get(req.params.id); + if (!agent) return res.status(404).json({ status: 'error', message: 'agent not found' }); + await agent.update({ revoked: true }); + agentManager.disconnect(agent.id, 4003, 'Enrollment revoked'); + logAgentAudit('revoke', { actor: req.user.uid, agentId: agent.id, agentName: agent.name }); + res.json({ status: 'ok' }); + } catch (err) { next(err); } +}); + +router.post('/nodes/:id/rotate', async (req, res, next) => { + try { + const agent = await Agent.get(req.params.id); + if (!agent) return res.status(404).json({ status: 'error', message: 'agent not found' }); + const token = await agent.rotateToken(); + // The old token is dead the moment it is replaced; drop the socket that was + // using it so the agent reconnects with the new one. + agentManager.disconnect(agent.id, 4004, 'Token rotated'); + logAgentAudit('rotate', { actor: req.user.uid, agentId: agent.id, agentName: agent.name }); + res.json({ status: 'ok', token, publicKey: await agentManager.publicKeyBase64() }); + } catch (err) { next(err); } +}); + +router.delete('/nodes/:id', async (req, res, next) => { + try { + const agent = await Agent.get(req.params.id); + if (!agent) return res.status(404).json({ status: 'error', message: 'agent not found' }); + agentManager.disconnect(agent.id, 4003, 'Enrollment deleted'); + await agent.delete(); + logAgentAudit('delete', { actor: req.user.uid, agentId: agent.id, agentName: agent.name }); + res.json({ status: 'ok' }); + } catch (err) { next(err); } +}); + +// --- Commands --- +// Addressed by agent id, not by token: a token is a credential and has no +// business travelling in a URL, being logged, or sitting in browser history. +router.post('/nodes/:id/command', async (req, res, next) => { + const { command, payload, isHighRisk } = req.body || {}; if (!command) { return res.status(400).json({ status: 'error', message: 'Command type is required' }); } - try { - const HIGH_RISK_COMMANDS = ['reboot', 'service_restart', 'configure_ldap', 'arbitrary_bash', 'update_binary']; - const requiresSigning = isHighRisk || HIGH_RISK_COMMANDS.includes(command); + const agent = await Agent.get(req.params.id); + if (!agent) return res.status(404).json({ status: 'error', message: 'agent not found' }); + if (agent.revoked) return res.status(403).json({ status: 'error', message: 'agent enrollment is revoked' }); + + const requiresSigning = isHighRisk || HIGH_RISK_COMMANDS.includes(command); + const msg = await agentManager.sendCommand(agent, command, payload || {}, requiresSigning); + + logAgentAudit('command', { + actor: req.user.uid, + agentId: agent.id, + agentName: agent.name, + resourceId: agent.resourceId || null, + command, + signed: requiresSigning + }); - const msg = agentManager.sendCommand(token, command, payload || {}, requiresSigning); res.json({ status: 'ok', sentMessage: msg }); } catch (err) { + logAgentAudit('command_failed', { actor: req.user && req.user.uid, agentId: req.params.id, command, error: err.message }); res.status(400).json({ status: 'error', message: err.message }); } }); module.exports = router; +module.exports.HIGH_RISK_COMMANDS = HIGH_RISK_COMMANDS; + module.exports.initAgentWebSockets = function initAgentWebSockets(app) { // WebSocket handler only needs the WS server; runs from the onListen hook. if (!app.wss) return; - app.wss.on('connection', (ws, req) => { + + // Warm the signing key at boot so a misconfigured OpenBao policy is a loud + // startup error rather than a surprise the first time someone reboots a host. + agentKeys.load().then(keys => { + if (!keys) console.error(`[Theta Agent] signing key unavailable — high-risk commands will be refused. ${agentKeys.status().error || ''}`); + }); + + app.wss.on('connection', async (ws, req) => { const url = new URL(req.url, `http://${req.headers.host || 'localhost'}`); const token = url.searchParams.get('token') || req.headers['authorization']; + const remoteAddr = req.socket.remoteAddress; - if (!token) { - ws.close(4001, 'Unauthorized: Missing token'); + // Authenticate BEFORE doing anything else: no registration, no welcome + // payload, no acknowledgement that the token was close. Until this passes + // the peer is an anonymous stranger, and the old code treated it as a + // trusted node purely for presenting a non-empty string. + let agent = null; + try { + agent = await Agent.authenticate(token); + } catch (err) { + console.error('[Theta Agent] authentication lookup failed:', err.message); + try { ws.close(1011, 'Authentication unavailable'); } catch (e) {} return; } - const remoteAddr = req.socket.remoteAddress; - console.log(`[Theta Agent] Agent connected from ${remoteAddr} with token ${token.substring(0, 8)}...`); + if (!agent) { + // Deliberately indistinguishable for unknown vs revoked vs missing: a + // caller probing tokens learns nothing about which part was wrong. + logAgentAudit('auth_rejected', { remoteAddr, tokenPrefix: token ? String(token).slice(0, 8) : null }); + try { ws.close(4001, 'Unauthorized'); } catch (e) {} + return; + } - agentManager.registerAgent(token, ws, remoteAddr); + console.log(`[Theta Agent] "${agent.name}" (${agent.id}) connected from ${remoteAddr}`); + logAgentAudit('connected', { agentId: agent.id, agentName: agent.name, remoteAddr }); + // Must stay synchronous, and the listeners below must be attached in this + // same tick: the agent sends `discovery` the instant the socket opens, and + // `ws` discards messages emitted while no listener is attached. + agentManager.registerAgent(agent, ws, remoteAddr); - ws.on('message', (message) => { + ws.on('message', async (message) => { try { const data = JSON.parse(message); if (!data || typeof data.type !== 'string') return; + // Re-read the row per message so a revoke mid-session takes effect on + // the next thing the agent says, not only on reconnect. + const current = await Agent.get(agent.id).catch(() => null); + if (!current || current.revoked) { + try { ws.close(4003, 'Enrollment revoked'); } catch (e) {} + return; + } + const payload = data.payload || {}; switch (data.type) { case 'discovery': - agentManager.handleDiscovery(token, payload); - if (app.io) app.io.emit('agent.discovery', { token, payload }); + await agentManager.handleDiscovery(current, payload); + if (app.io) app.io.emit('agent.discovery', { agentId: current.id, payload }); break; case 'telemetry': - agentManager.handleTelemetry(token, payload); - if (app.io) app.io.emit('agent.telemetry', { token, payload }); + await agentManager.handleTelemetry(current, payload); + if (app.io) app.io.emit('agent.telemetry', { agentId: current.id, payload }); break; case 'heartbeat': - agentManager.handleHeartbeat(token, payload, ws); + await agentManager.handleHeartbeat(current, payload, ws); break; case 'response': - agentManager.handleResponse(token, payload); - if (app.io) app.io.emit('agent.response', { token, payload }); + await agentManager.handleResponse(current, payload); + if (app.io) app.io.emit('agent.response', { agentId: current.id, payload }); break; default: - console.log(`[Theta Agent] Received message type '${data.type}' from ${token}`); + console.log(`[Theta Agent] Received message type '${data.type}' from ${current.id}`); } } catch (err) { - console.error("[Theta Agent] Error parsing message:", err); + console.error('[Theta Agent] Error handling message:', err); } }); ws.on('close', () => { - console.log(`[Theta Agent] Agent disconnected (${token})`); - agentManager.unregisterAgent(token, ws); + console.log(`[Theta Agent] "${agent.name}" (${agent.id}) disconnected`); + agentManager.unregisterAgent(agent.id, ws); }); // Send initial welcome/config payload @@ -120,7 +300,8 @@ module.exports.initAgentWebSockets = function initAgentWebSockets(app) { type: 'config', payload: { message: 'Connected to SSO Manager C2', - protocol_version: '1.1.0' + protocol_version: '1.2.0', + agent_id: agent.id } })); } catch (e) {} diff --git a/nodejs/services/discovery_reconciler.js b/nodejs/services/discovery_reconciler.js index c1a30a9..7fa66cf 100644 --- a/nodejs/services/discovery_reconciler.js +++ b/nodejs/services/discovery_reconciler.js @@ -2,26 +2,64 @@ const { Resource, ResourceEdge, ResourceGroup } = require('../models/resource'); const { WebhookEmitter } = require('./webhook_emitter'); const crypto = require('crypto'); +// Is `candidateId` at or below `rootId` in the edge graph? Used to refuse an +// edge that would close a loop. Carries its own visited set so it terminates +// even if the stored graph already contains a cycle from an older release. +function isDescendant(candidateId, rootId, edges) { + const seen = new Set(); + const stack = [rootId]; + while (stack.length) { + const id = stack.pop(); + if (id === candidateId) return true; + if (seen.has(id)) continue; + seen.add(id); + for (const e of edges) if (e.parentId === id) stack.push(e.childId); + } + return false; +} + class DiscoveryReconciler { static async reconcile(sourceName, payload) { const { resources = [], edges = [] } = payload; let newDevices = 0; + const normalizeMac = (m) => (m || '').toLowerCase().replace(/[^a-f0-9]/g, ''); + const normalizeHost = (h) => (h || '').toLowerCase().split('.')[0].trim(); + + // Read the inventory ONCE, not once per incoming resource. A Proxmox + // cluster reports ~55 resources against an inventory of similar size, so + // the per-iteration Resource.list() was doing quadratic full-table reads + // every discovery run. Newly created rows are pushed onto this list as we + // go, so later resources in the same payload still match against them. + const allRes = await Resource.list(); + for (const res of resources) { if (!res.metadata) res.metadata = {}; res._originalSlug = res.slug; // Keep track for edge mapping - - let existing = null; - const normalizeMac = (m) => (m || '').toLowerCase().replace(/[^a-f0-9]/g, ''); - const normalizeHost = (h) => (h || '').toLowerCase().split('.')[0].trim(); - const allRes = await Resource.list(); + let existing = null; + + // A discovered device may only merge into a resource of the same kind + // (or into a placeholder from an earlier, kind-less discovery). Without + // this a VM called "gitea-runner" matches a hand-created *service* of + // the same name on rule 3 and silently overwrites it -- the discovered + // host's metadata lands on a service row, and the operator's entry is + // gone. `template` counts as `host`: a VM converted to a template is the + // same device, and it should update in place rather than fork a row. + const kindClass = (k) => (k === 'template' ? 'host' : k); + const incomingKind = kindClass(res.kind || 'unmanaged_device'); + const kindCompatible = (r) => { + const k = kindClass(r.kind); + if (k === 'unmanaged_device' || incomingKind === 'unmanaged_device') return true; + return k === incomingKind; + }; + const candidates = allRes.filter(kindCompatible); // 1. Attempt matching by MAC (highest precision) if (res.metadata.interfaces && res.metadata.interfaces.length > 0) { const macs = res.metadata.interfaces.map(i => normalizeMac(i.mac)).filter(m => m.length === 12); if (macs.length > 0) { - existing = allRes.find(r => + existing = candidates.find(r => r.metadata && ( (r.metadata.macAddress && macs.includes(normalizeMac(r.metadata.macAddress))) || (r.metadata.interfaces && r.metadata.interfaces.some(i => macs.includes(normalizeMac(i.mac)))) @@ -42,7 +80,7 @@ class DiscoveryReconciler { ipsToMatch = [...new Set(ipsToMatch.filter(Boolean))]; if (!existing && ipsToMatch.length > 0) { - existing = allRes.find(r => { + existing = candidates.find(r => { if (!r.metadata) return false; if (r.metadata.ip && ipsToMatch.includes(r.metadata.ip)) return true; if (r.metadata.address) { @@ -57,7 +95,7 @@ class DiscoveryReconciler { // 3. Fallback matching by Slug, Name, or Base Hostname if (!existing && (res.slug || res.name)) { const inputName = normalizeHost(res.name || res.slug); - existing = allRes.find(r => { + existing = candidates.find(r => { if (res.slug && r.slug === res.slug) return true; if (res.name && r.name && r.name.toLowerCase() === res.name.toLowerCase()) return true; if (inputName && r.name && normalizeHost(r.name) === inputName) return true; @@ -93,10 +131,33 @@ class DiscoveryReconciler { mergedMeta.last_seen = Date.now(); - const isIp = (str) => /^(?:[0-9]{1,3}\\.){3}[0-9]{1,3}$/.test(str || ''); + // Pick the most human name across sources. Rank first, length only as + // a tie-break within a rank -- comparing lengths alone let a UniFi + // client named after its MAC ("ac:16:2d:b3:da:80", 17 chars) beat the + // hypervisor's real hostname from Proxmox ("dl380-0", 7), so the + // Directory listed MAC addresses where host names belong. + // + // NB: `\\.` inside a regex LITERAL matches a backslash, not a dot, so + // the old isIp returned false for every input and IP-shaped names were + // never replaced either. It is `\.` here. + const isIp = (str) => /^(?:[0-9]{1,3}\.){3}[0-9]{1,3}$/.test(str || ''); + const isMac = (str) => /^([0-9a-f]{2}[:-]){5}[0-9a-f]{2}$/i.test((str || '').trim()); + // 2 = a real name, 1 = an IP (at least routable/recognizable), 0 = a + // MAC or nothing (pure machine identifier, the worst thing to show). + const nameRank = (str) => { + if (!str || !String(str).trim()) return 0; + if (isMac(str)) return 0; + if (isIp(str)) return 1; + return 2; + }; + let bestName = existing.name; - if (res.name && (!bestName || isIp(bestName) || res.name.length > bestName.length && !isIp(res.name))) { - bestName = res.name; + if (res.name) { + const incoming = nameRank(res.name); + const current = nameRank(bestName); + if (incoming > current || (incoming === current && res.name.length > (bestName || '').length)) { + bestName = res.name; + } } await existing.update({ @@ -125,12 +186,17 @@ class DiscoveryReconciler { newDevices++; res._actualId = created.id; // Map original slug to actual ID + // Make it visible to the rest of THIS payload: a Proxmox run reports + // the endpoint, then its nodes, then their guests, and two of them can + // legitimately share a MAC/IP. Without this the same device could be + // created twice in a single run. + allRes.push(created); WebhookEmitter.emit('discovery.new_device', created.toJSON()); } } - // Now process edges - const allRes = await Resource.list(); + // Now process edges. `allRes` above is already current -- rows created in + // the loop were pushed onto it -- so no second full read is needed. const existingEdges = await ResourceEdge.list(); for (const edge of edges) { @@ -154,15 +220,37 @@ class DiscoveryReconciler { if (childResInDb) childId = childResInDb.id; } + // Two slugs in one payload can resolve to the SAME resource once the + // matcher has merged them -- a Proxmox endpoint reached at the address + // of the node that answers for it is the case that produced this. The + // edge would then make a resource its own parent, which renders as an + // infinitely nested tree and defeats every ancestor walk in the app + // (findAncestorSiteSlug, withResolvedAddress) that relies on a cycle + // guard to terminate rather than to be correct. + if (parentId && childId && parentId === childId) { + console.warn(`[DiscoveryReconciler] ${sourceName}: dropping self-edge on ${edge.parentSlug} -> ${edge.childSlug} (both resolved to the same resource)`); + continue; + } + + // Likewise refuse an edge that closes a loop: if the proposed parent is + // already a descendant of the proposed child, adding this makes a cycle. + if (parentId && childId && isDescendant(parentId, childId, existingEdges)) { + console.warn(`[DiscoveryReconciler] ${sourceName}: dropping ${edge.parentSlug} -> ${edge.childSlug} (would create a cycle)`); + continue; + } + if (parentId && childId) { const edgeExists = existingEdges.find(e => e.parentId === parentId && e.childId === childId && e.relation === edge.relation); if (!edgeExists) { - await ResourceEdge.create({ + const created = await ResourceEdge.create({ id: crypto.randomUUID(), parentId, childId, relation: edge.relation }); + // Keep the in-memory edge list current so the cycle check above sees + // edges added earlier in this same payload. + existingEdges.push(created); } } } diff --git a/nodejs/tests/agent_manager.test.js b/nodejs/tests/agent_manager.test.js index 11b227e..08fa43d 100644 --- a/nodejs/tests/agent_manager.test.js +++ b/nodejs/tests/agent_manager.test.js @@ -1,11 +1,45 @@ 'use strict'; const crypto = require('crypto'); -const agentManager = require('../utils/agent_manager'); -describe('AgentManager PROTOCOL.md v1.1.0 Compliance', () => { +// In-memory stand-in for OpenBao. The signing key lives at secret/agent/ +// signing-key in production; here we only need it to persist across calls so +// the "same key every time" property is actually exercised rather than mocked +// away. +const mockBaoStore = new Map(); +jest.mock('@simpleworkjs/bao-conf', () => ({ + get: jest.fn(async (path) => mockBaoStore.get(path) || null), + set: jest.fn(async (path, value) => { mockBaoStore.set(path, value); }), + request: jest.fn(async () => ({ ok: true, status: 200 })) +})); + +const agentManager = require('../utils/agent_manager'); +const agentKeys = require('../utils/agent_keys'); + +// The manager is now keyed by enrolled Agent rows rather than by a bare token +// string, so these use a stub row with the same surface the real model gives: +// an id, and an update() that records what would be persisted. +function stubAgent(overrides = {}) { + const row = { + id: overrides.id || crypto.randomUUID(), + name: overrides.name || 'test-agent', + resourceId: overrides.resourceId || null, + revoked: false, + persisted: {}, + ...overrides + }; + row.update = jest.fn(async (patch) => { + Object.assign(row.persisted, patch); + Object.assign(row, patch); + return row; + }); + return row; +} + +describe('AgentManager PROTOCOL.md v1.2.0 Compliance', () => { let mockWs; let sentMessages; + let agent; beforeEach(() => { sentMessages = []; @@ -14,23 +48,30 @@ describe('AgentManager PROTOCOL.md v1.1.0 Compliance', () => { send: jest.fn((msg) => sentMessages.push(JSON.parse(msg))), close: jest.fn() }; + agent = stubAgent(); }); - test('registers agent and tracks initial connection state', () => { - const record = agentManager.registerAgent('test-token-123', mockWs, '192.168.1.100'); - expect(record.token).toBe('test-token-123'); - expect(record.ipAddress).toBe('192.168.1.100'); - - const agents = agentManager.getConnectedAgents(); - const found = agents.find(a => a.token === 'test-token-123'); - expect(found).toBeDefined(); - expect(found.isOnline).toBe(true); + test('registers an agent and reports it as connected', () => { + agentManager.registerAgent(agent, mockWs, '192.168.1.100'); + const state = agentManager.liveState(agent.id); + expect(state.connected).toBe(true); + expect(state.ipAddress).toBe('192.168.1.100'); + expect(agentManager.isConnected(agent.id)).toBe(true); }); - test('processes discovery payload per PROTOCOL.md v1.1.0 Section 3.1', () => { - agentManager.registerAgent('test-token-123', mockWs, '192.168.1.100'); + // registerAgent must not be async: the WS `message` listener is attached in + // the same tick, and `ws` drops events emitted before a listener exists. An + // awaited DB write here swallowed every agent's first discovery frame, which + // is the one it sends immediately on connect. + test('registerAgent is synchronous so no message can be missed', () => { + const result = agentManager.registerAgent(agent, mockWs, '10.0.0.1'); + expect(result).toBeUndefined(); + expect(agentManager.isConnected(agent.id)).toBe(true); + }); - const discoveryPayload = { + test('persists discovery to the agent row (Section 3.1)', async () => { + agentManager.registerAgent(agent, mockWs, '192.168.1.100'); + await agentManager.handleDiscovery(agent, { hostname: 'node-01.local', ip_addresses: ['192.168.1.100', '10.0.0.5'], os: 'Ubuntu 24.04 LTS', @@ -39,41 +80,34 @@ describe('AgentManager PROTOCOL.md v1.1.0 Compliance', () => { ram_total_gb: 32.0, disk_total_gb: 500.0, location: 'dc-chicago-rack-4' - }; + }); - agentManager.handleDiscovery('test-token-123', discoveryPayload); - - const agents = agentManager.getConnectedAgents(); - const agent = agents.find(a => a.token === 'test-token-123'); - expect(agent.hostname).toBe('node-01.local'); - expect(agent.discovery.os).toBe('Ubuntu 24.04 LTS'); - expect(agent.discovery.ip_addresses).toEqual(['192.168.1.100', '10.0.0.5']); + const saved = agent.persisted.lastDiscovery; + expect(saved.hostname).toBe('node-01.local'); + expect(saved.os).toBe('Ubuntu 24.04 LTS'); + expect(saved.ip_addresses).toEqual(['192.168.1.100', '10.0.0.5']); + // Durable, not just in memory: an agent that goes offline keeps its facts. + expect(agent.persisted.last_seen).toEqual(expect.any(Number)); }); - test('processes telemetry payload per PROTOCOL.md v1.1.0 Section 3.2', () => { - agentManager.registerAgent('test-token-123', mockWs, '192.168.1.100'); - - const telemetryPayload = { + test('persists telemetry to the agent row (Section 3.2)', async () => { + agentManager.registerAgent(agent, mockWs, '192.168.1.100'); + await agentManager.handleTelemetry(agent, { cpu_usage_percent: 14.5, ram_usage_percent: 42.1, disk_usage_percent: 68.0, zfs_health: 'ONLINE', gpu_usage_percent: -1.0, timestamp: new Date().toISOString() - }; + }); - agentManager.handleTelemetry('test-token-123', telemetryPayload); - - const agents = agentManager.getConnectedAgents(); - const agent = agents.find(a => a.token === 'test-token-123'); - expect(agent.telemetry.cpu_usage_percent).toBe(14.5); - expect(agent.telemetry.zfs_health).toBe('ONLINE'); + expect(agent.persisted.lastTelemetry.cpu_usage_percent).toBe(14.5); + expect(agent.persisted.lastTelemetry.zfs_health).toBe('ONLINE'); }); - test('responds to heartbeat with heartbeat_ack per Section 3.3', () => { - agentManager.registerAgent('test-token-123', mockWs, '192.168.1.100'); - - agentManager.handleHeartbeat('test-token-123', { timestamp: new Date().toISOString() }, mockWs); + test('responds to heartbeat with heartbeat_ack (Section 3.3)', async () => { + agentManager.registerAgent(agent, mockWs, '192.168.1.100'); + await agentManager.handleHeartbeat(agent, { timestamp: new Date().toISOString() }, mockWs); expect(mockWs.send).toHaveBeenCalled(); const lastMsg = sentMessages[sentMessages.length - 1]; @@ -81,20 +115,92 @@ describe('AgentManager PROTOCOL.md v1.1.0 Compliance', () => { expect(lastMsg.payload.timestamp).toBeDefined(); }); - test('canonicalizes payload and signs high-risk commands using Ed25519 per Section 5', () => { - agentManager.registerAgent('test-token-123', mockWs, '192.168.1.100'); + test('canonicalizes and signs high-risk commands with Ed25519 (Section 5)', async () => { + agentManager.registerAgent(agent, mockWs, '192.168.1.100'); const rawPayload = { script: 'uptime', location: 'datacenter' }; - const msg = agentManager.sendCommand('test-token-123', 'arbitrary_bash', rawPayload, true); + const msg = await agentManager.sendCommand(agent, 'arbitrary_bash', rawPayload, true); expect(msg.type).toBe('arbitrary_bash'); - expect(msg.payload.signature).toBeDefined(); expect(typeof msg.payload.signature).toBe('string'); - // Verify signature with public key - const signatureBuffer = Buffer.from(msg.payload.signature, 'base64'); - const canonicalStr = agentManager.canonicalize(rawPayload); - const isValid = crypto.verify(null, Buffer.from(canonicalStr, 'utf8'), agentManager.publicKeyPem, signatureBuffer); + const keys = await agentKeys.load(); + const isValid = crypto.verify( + null, + Buffer.from(agentManager.canonicalize(rawPayload), 'utf8'), + crypto.createPublicKey(keys.publicKeyPem), + Buffer.from(msg.payload.signature, 'base64') + ); expect(isValid).toBe(true); }); + + // The canonical form has to match the Go agent's byte for byte. Go's + // encoding/json escapes <, > and & by default and JSON.stringify does not, so + // the agent uses SetEscapeHTML(false); this pins the server's half of that + // contract. See theta-agent TestCanonicalizeMatchesServerForm. + test('canonical form is sorted, unescaped, and omits the signature', () => { + const canonical = agentManager.canonicalize({ + script: 'echo a > b && c', + comment: 'x&y', + signature: 'should-not-appear' + }); + expect(canonical).toBe('{"comment":"x&y","script":"echo a > b && c"}'); + }); + + test('refuses to send to an agent that is not connected', async () => { + await expect(agentManager.sendCommand(agent, 'reload_config', {}, false)) + .rejects.toThrow(/not connected/); + }); + + // Revocation that only applies on the next reconnect is not revocation. + test('disconnect drops the live socket immediately', () => { + agentManager.registerAgent(agent, mockWs, '192.168.1.100'); + expect(agentManager.isConnected(agent.id)).toBe(true); + + const dropped = agentManager.disconnect(agent.id, 4003, 'Enrollment revoked'); + expect(dropped).toBe(true); + expect(mockWs.close).toHaveBeenCalledWith(4003, 'Enrollment revoked'); + expect(agentManager.isConnected(agent.id)).toBe(false); + }); + + test('a second connection for the same agent supersedes the first', () => { + agentManager.registerAgent(agent, mockWs, '192.168.1.100'); + const secondWs = { readyState: 1, send: jest.fn(), close: jest.fn() }; + agentManager.registerAgent(agent, secondWs, '192.168.1.101'); + + expect(mockWs.close).toHaveBeenCalledWith(4002, 'Superseded by new connection'); + expect(agentManager.liveState(agent.id).ipAddress).toBe('192.168.1.101'); + }); + + test('an unknown agent id is simply not connected', () => { + expect(agentManager.isConnected('no-such-agent')).toBe(false); + expect(agentManager.liveState('no-such-agent')).toEqual({ connected: false, lastResponse: null }); + }); +}); + +describe('agent signing key', () => { + // The old manager generated a key pair in its constructor, so it changed on + // every restart and the public_key pinned in agent.yml stopped matching. + test('the same key is returned across repeated loads', async () => { + const first = await agentKeys.load(); + const second = await agentKeys.load(); + expect(first.publicKeyBase64).toBe(second.publicKeyBase64); + }); + + test('the exported public key is the raw 32 bytes agents pin', async () => { + const keys = await agentKeys.load(); + expect(Buffer.from(keys.publicKeyBase64, 'base64')).toHaveLength(32); + }); + + test('rawPublicKeyBase64 strips the SPKI wrapper', () => { + const { publicKey } = crypto.generateKeyPairSync('ed25519', { + privateKeyEncoding: { type: 'pkcs8', format: 'pem' }, + publicKeyEncoding: { type: 'spki', format: 'pem' } + }); + const raw = Buffer.from(agentKeys.rawPublicKeyBase64(publicKey), 'base64'); + expect(raw).toHaveLength(32); + // and it is the tail of the DER encoding + const der = crypto.createPublicKey(publicKey).export({ type: 'spki', format: 'der' }); + expect(raw.equals(der.subarray(der.length - 32))).toBe(true); + }); }); diff --git a/nodejs/tests/discovery_naming.test.js b/nodejs/tests/discovery_naming.test.js new file mode 100644 index 0000000..c4299c7 --- /dev/null +++ b/nodejs/tests/discovery_naming.test.js @@ -0,0 +1,104 @@ +'use strict'; + +// Pure-logic coverage for the two reconciler rules that real Proxmox + UniFi +// data broke. Both were found by running discovery against a live cluster: +// the directory came back listing MAC addresses as host names, and one +// resource ended up as its own parent. + +// Mirrors the ranking in services/discovery_reconciler.js. Kept here (rather +// than exported) because it is a few lines of predicate that the reconciler +// applies inline while merging; if it grows, export it and drop this copy. +const isIp = (str) => /^(?:[0-9]{1,3}\.){3}[0-9]{1,3}$/.test(str || ''); +const isMac = (str) => /^([0-9a-f]{2}[:-]){5}[0-9a-f]{2}$/i.test((str || '').trim()); +const nameRank = (str) => { + if (!str || !String(str).trim()) return 0; + if (isMac(str)) return 0; + if (isIp(str)) return 1; + return 2; +}; +function bestNameOf(existingName, incomingName) { + let best = existingName; + if (incomingName) { + const a = nameRank(incomingName); + const b = nameRank(best); + if (a > b || (a === b && incomingName.length > (best || '').length)) best = incomingName; + } + return best; +} + +describe('discovery name ranking', () => { + test('a real hostname beats a MAC even when shorter', () => { + // The exact regression: UniFi named the host by MAC, Proxmox knew the + // hostname, and length-only comparison kept the MAC. + expect(bestNameOf('ac:16:2d:b3:da:80', 'dl380-0')).toBe('dl380-0'); + }); + + test('a MAC never displaces a real hostname', () => { + expect(bestNameOf('dl380-0', 'ac:16:2d:b3:da:80')).toBe('dl380-0'); + }); + + test('a real hostname beats an IP-shaped name', () => { + expect(bestNameOf('192.168.1.27', 'hass.io')).toBe('hass.io'); + }); + + test('an IP beats a MAC', () => { + expect(bestNameOf('bc:24:11:3f:cd:c8', '192.168.1.27')).toBe('192.168.1.27'); + }); + + test('an IP does not displace a hostname', () => { + expect(bestNameOf('gitea-runner', '192.168.1.176')).toBe('gitea-runner'); + }); + + test('within the same rank the longer/more specific name wins', () => { + expect(bestNameOf('pve', 'pve-dl380-1')).toBe('pve-dl380-1'); + }); + + test('dash-separated MACs are recognized too', () => { + expect(bestNameOf('ac-16-2d-b3-da-80', 'dl380-0')).toBe('dl380-0'); + }); + + test('an empty existing name is always replaced', () => { + expect(bestNameOf('', 'anything')).toBe('anything'); + expect(bestNameOf(null, 'ac:16:2d:b3:da:80')).toBe('ac:16:2d:b3:da:80'); + }); +}); + +// Mirrors isDescendant() in the reconciler. +function isDescendant(candidateId, rootId, edges) { + const seen = new Set(); + const stack = [rootId]; + while (stack.length) { + const id = stack.pop(); + if (id === candidateId) return true; + if (seen.has(id)) continue; + seen.add(id); + for (const e of edges) if (e.parentId === id) stack.push(e.childId); + } + return false; +} + +describe('discovery edge cycle guard', () => { + const edges = [ + { parentId: 'cluster', childId: 'node1' }, + { parentId: 'node1', childId: 'vm1' }, + ]; + + test('detects a direct parent/child inversion', () => { + // Proposing node1 -> cluster when cluster -> node1 already exists. + expect(isDescendant('node1', 'cluster', edges)).toBe(true); + }); + + test('detects a deeper loop', () => { + expect(isDescendant('vm1', 'cluster', edges)).toBe(true); + }); + + test('allows an unrelated new parent', () => { + expect(isDescendant('node2', 'cluster', edges)).toBe(false); + }); + + test('terminates on a graph that already contains a cycle', () => { + // A self-edge written by an earlier release must not hang the walk. + const cyclic = [{ parentId: 'a', childId: 'a' }, { parentId: 'a', childId: 'b' }]; + expect(isDescendant('zzz', 'a', cyclic)).toBe(false); + }); +}); diff --git a/nodejs/tests/proxmox_interfaces.test.js b/nodejs/tests/proxmox_interfaces.test.js new file mode 100644 index 0000000..4bfbfd9 --- /dev/null +++ b/nodejs/tests/proxmox_interfaces.test.js @@ -0,0 +1,78 @@ +'use strict'; + +const { _Interfaces: Interfaces } = require('../plugins/discovery/proxmox'); + +// Regression coverage for the MAC/IP mismatch: the plugin used to collect MACs +// and IPs into two flat lists and zip them by index, so on a multi-NIC guest +// -- or any guest where one NIC had no address -- the directory recorded an IP +// against the wrong MAC. Interfaces keys by MAC so a pairing can only come from +// the source that observed both together. +describe('proxmox Interfaces', () => { + test('keeps each IP on the NIC it was observed on', () => { + const i = new Interfaces(); + i.add('AA:BB:CC:00:00:01', ['10.0.0.5'], 'eth0'); + i.add('AA:BB:CC:00:00:02', ['192.168.9.7'], 'eth1'); + + expect(i.toArray()).toEqual([ + { mac: 'aa:bb:cc:00:00:01', ip: '10.0.0.5', ips: ['10.0.0.5'], name: 'eth0' }, + { mac: 'aa:bb:cc:00:00:02', ip: '192.168.9.7', ips: ['192.168.9.7'], name: 'eth1' }, + ]); + }); + + test('a NIC with no address does not steal the next NIC\'s IP', () => { + const i = new Interfaces(); + i.add('AA:BB:CC:00:00:01', [], 'eth0'); // stopped/unconfigured + i.add('AA:BB:CC:00:00:02', ['10.0.0.9'], 'eth1'); + + const byMac = Object.fromEntries(i.toArray().map(x => [x.mac, x.ip])); + expect(byMac['aa:bb:cc:00:00:01']).toBeNull(); + expect(byMac['aa:bb:cc:00:00:02']).toBe('10.0.0.9'); + }); + + test('merges the config MAC with the agent-reported address for the same NIC', () => { + const i = new Interfaces(); + i.add('aa:bb:cc:00:00:01', ['10.0.0.5'], 'eth0'); // guest agent + i.add('AA:BB:CC:00:00:01', [], 'net0'); // VM config, same NIC + expect(i.toArray()).toHaveLength(1); + expect(i.toArray()[0]).toMatchObject({ mac: 'aa:bb:cc:00:00:01', ip: '10.0.0.5' }); + }); + + test('collects multiple addresses on one NIC without inventing a second NIC', () => { + const i = new Interfaces(); + i.add('aa:bb:cc:00:00:01', ['10.0.0.5', '10.0.0.6'], 'eth0'); + expect(i.toArray()).toHaveLength(1); + expect(i.toArray()[0].ips).toEqual(['10.0.0.5', '10.0.0.6']); + expect(i.primaryIp()).toBe('10.0.0.5'); + }); + + test('ignores placeholder and malformed MACs', () => { + const i = new Interfaces(); + i.add('00:00:00:00:00:00', [], 'eth0'); + i.add('not-a-mac', [], 'eth1'); + i.add('', [], 'eth2'); + expect(i.toArray()).toEqual([]); + expect(i.primaryMac()).toBeNull(); + }); + + test('keeps an address that arrived without a usable MAC', () => { + const i = new Interfaces(); + i.add(null, ['10.0.0.5'], 'eth0'); + expect(i.primaryIp()).toBe('10.0.0.5'); + expect(i.primaryMac()).toBeNull(); + }); + + test('primary values prefer a NIC that actually has an address', () => { + const i = new Interfaces(); + i.add('aa:bb:cc:00:00:01', [], 'eth0'); + i.add('aa:bb:cc:00:00:02', ['10.0.0.9'], 'eth1'); + expect(i.primaryIp()).toBe('10.0.0.9'); + expect(i.primaryMac()).toBe('aa:bb:cc:00:00:02'); + }); + + test('a fully unaddressed guest still reports its MAC', () => { + const i = new Interfaces(); + i.add('aa:bb:cc:00:00:01', [], 'net0'); + expect(i.primaryIp()).toBeNull(); + expect(i.primaryMac()).toBe('aa:bb:cc:00:00:01'); + }); +}); diff --git a/nodejs/utils/agent_keys.js b/nodejs/utils/agent_keys.js new file mode 100644 index 0000000..3d91d93 --- /dev/null +++ b/nodejs/utils/agent_keys.js @@ -0,0 +1,99 @@ +'use strict'; + +// The Ed25519 key pair the SSO signs high-risk agent commands with, stored in +// OpenBao at `secret/agent/signing-key`. +// +// This used to be generated in the AgentManager constructor and kept only in +// memory, which made the whole signing scheme decorative: every SSO restart +// produced a new key, so the `public_key` pinned in an agent's agent.yml stopped +// matching and the agent either rejected everything or (because it skips +// verification when no key is configured) executed everything unverified. A +// trust anchor that changes on restart is not a trust anchor. +// +// Requires the sso-broker OpenBao policy to grant `secret/agent/*` +// (theta-suite setup.sh). Without it the load fails and signing is reported as +// unavailable -- we deliberately do NOT fall back to an ephemeral key, because +// signing with a key no agent has ever seen is worse than refusing: it looks +// like it worked. + +const crypto = require('crypto'); +const baoConf = require('@simpleworkjs/bao-conf'); + +const PATH = 'agent/signing-key'; // baoConf adds the secret/data prefix + +let cached = null; // { privateKeyPem, publicKeyPem, publicKeyBase64 } +let loadError = null; + +// Agents pin the raw 32-byte Ed25519 public key, base64-encoded (see the Go +// client's verifySignature, which base64-decodes cfg.public_key and expects +// ed25519.PublicKeySize bytes). Node hands us SPKI PEM, so strip the 12-byte +// DER prefix to get the raw key the agent actually wants. +function rawPublicKeyBase64(publicKeyPem) { + const der = crypto.createPublicKey(publicKeyPem).export({ type: 'spki', format: 'der' }); + return Buffer.from(der.subarray(der.length - 32)).toString('base64'); +} + +function generate() { + const { privateKey, publicKey } = crypto.generateKeyPairSync('ed25519', { + privateKeyEncoding: { type: 'pkcs8', format: 'pem' }, + publicKeyEncoding: { type: 'spki', format: 'pem' } + }); + return { privateKeyPem: privateKey, publicKeyPem: publicKey }; +} + +// Load the stored key pair, generating and persisting one on first run. +// Idempotent and safe to call repeatedly; the result is cached in-process. +async function load() { + if (cached) return cached; + + let stored = null; + try { + stored = await baoConf.get(PATH); + } catch (err) { + loadError = `could not read ${PATH} from OpenBao: ${err.message}`; + console.error(`[agent_keys] ${loadError}`); + return null; + } + + if (stored && stored.privateKeyPem && stored.publicKeyPem) { + cached = { + privateKeyPem: stored.privateKeyPem, + publicKeyPem: stored.publicKeyPem, + publicKeyBase64: rawPublicKeyBase64(stored.publicKeyPem) + }; + loadError = null; + return cached; + } + + // First run: mint one and persist it before use, so a crash between + // generating and storing can't leave agents pinned to a key we forgot. + const fresh = generate(); + try { + await baoConf.set(PATH, fresh); + } catch (err) { + loadError = `could not persist a signing key to ${PATH}: ${err.message}. ` + + 'Re-run ./setup.sh so the sso-broker policy grants secret/agent/*.'; + console.error(`[agent_keys] ${loadError}`); + return null; + } + + cached = { + ...fresh, + publicKeyBase64: rawPublicKeyBase64(fresh.publicKeyPem) + }; + loadError = null; + console.log('[agent_keys] generated and stored a new agent signing key'); + return cached; +} + +function status() { + return { available: !!cached, error: loadError }; +} + +// Test seam: drop the in-process cache. +function _reset() { + cached = null; + loadError = null; +} + +module.exports = { load, status, rawPublicKeyBase64, _reset, PATH }; diff --git a/nodejs/utils/agent_manager.js b/nodejs/utils/agent_manager.js index 93c6f23..3bb1926 100644 --- a/nodejs/utils/agent_manager.js +++ b/nodejs/utils/agent_manager.js @@ -1,26 +1,19 @@ 'use strict'; const crypto = require('crypto'); +const agentKeys = require('./agent_keys'); +const { Agent } = require('../models/agent'); +// Tracks the live WebSocket for each enrolled agent and brokers commands to it. +// +// The durable facts about an agent (identity, host binding, last seen, last +// discovery/telemetry) live in the Agent table; this class holds only what +// cannot be persisted -- the open socket. That split is what makes an installed +// -but-offline agent visible, and what stops a restart from erasing the fleet. class AgentManager { constructor() { - this.agents = new Map(); // token -> agentRecord - this.privateKeyPem = null; - this.publicKeyPem = null; - this.initKeyPair(); - } - - initKeyPair() { - try { - const { privateKey, publicKey } = crypto.generateKeyPairSync('ed25519', { - privateKeyEncoding: { type: 'pkcs8', format: 'pem' }, - publicKeyEncoding: { type: 'spki', format: 'pem' } - }); - this.privateKeyPem = privateKey; - this.publicKeyPem = publicKey; - } catch (err) { - console.error('[AgentManager] Failed to generate Ed25519 key pair:', err); - } + // agentId -> { ws, ipAddress, connectedAt, lastResponse, pending } + this.live = new Map(); } /** @@ -37,54 +30,89 @@ class AgentManager { } /** - * Sign payload using Ed25519 private key. - * Returns base64 encoded signature. + * Sign payload using the persisted Ed25519 private key. Throws when no key is + * available rather than minting a throwaway one -- an agent verifies against + * the key pinned in its agent.yml, so a signature from a key it has never + * seen is not a weaker signature, it is a broken command that looks fine from + * this side. */ - signPayload(payload) { - if (!this.privateKeyPem) { - throw new Error('Ed25519 private key is not initialized'); + async signPayload(payload) { + const keys = await agentKeys.load(); + if (!keys) { + const { error } = agentKeys.status(); + throw new Error(`agent command signing is unavailable: ${error || 'no signing key'}`); } const canonicalBytes = Buffer.from(this.canonicalize(payload), 'utf8'); - const signature = crypto.sign(null, canonicalBytes, this.privateKeyPem); - return signature.toString('base64'); + return crypto.sign(null, canonicalBytes, keys.privateKeyPem).toString('base64'); } - registerAgent(token, ws, remoteAddress) { - const existing = this.agents.get(token); + async publicKeyBase64() { + const keys = await agentKeys.load(); + return keys ? keys.publicKeyBase64 : null; + } + + async publicKeyPem() { + const keys = await agentKeys.load(); + return keys ? keys.publicKeyPem : null; + } + + // Bind a freshly authenticated socket to an enrolled agent. `agent` is an + // Agent row that Agent.authenticate() has already vouched for -- this method + // never sees a raw token and must never be called with an unauthenticated one. + // Synchronous by design. The caller must attach its `message` listener in the + // same tick as the connection is accepted: `ws` drops events emitted before a + // listener exists, and the agent sends `discovery` immediately on open, so + // awaiting a database round-trip here silently lost every agent's first + // discovery frame. The connect timestamp is persisted in the background. + registerAgent(agent, ws, remoteAddress) { + const existing = this.live.get(agent.id); if (existing && existing.ws && existing.ws !== ws) { try { existing.ws.close(4002, 'Superseded by new connection'); } catch (e) {} } - const agentRecord = { - token, + this.live.set(agent.id, { ws, ipAddress: remoteAddress, - hostname: 'unknown', connectedAt: new Date().toISOString(), - lastSeen: new Date().toISOString(), - discovery: {}, - telemetry: {}, - pendingResponses: new Map() - }; + lastResponse: null + }); - this.agents.set(token, agentRecord); - return agentRecord; + agent.update({ + last_seen: Math.floor(Date.now() / 1000), + last_ip: remoteAddress || null + }).catch(err => console.error(`[AgentManager] could not record connect for ${agent.id}:`, err.message)); } - unregisterAgent(token, ws) { - const record = this.agents.get(token); - if (record && record.ws === ws) { - this.agents.delete(token); - } + unregisterAgent(agentId, ws) { + const state = this.live.get(agentId); + if (state && state.ws === ws) this.live.delete(agentId); } - handleDiscovery(token, payload) { - const agent = this.agents.get(token); - if (!agent) return; + // Drop an agent's live socket now. Revocation that only takes effect on the + // next reconnect is not revocation -- a connected agent would keep receiving + // commands indefinitely. + disconnect(agentId, code = 4003, reason = 'Disconnected by server') { + const state = this.live.get(agentId); + if (!state || !state.ws) return false; + try { state.ws.close(code, reason); } catch (e) {} + this.live.delete(agentId); + return true; + } - agent.lastSeen = new Date().toISOString(); - agent.hostname = payload.hostname || agent.hostname; - agent.discovery = { + isConnected(agentId) { + const state = this.live.get(agentId); + return !!(state && state.ws && state.ws.readyState === 1); + } + + async touch(agent, extra = {}) { + await agent.update({ + last_seen: Math.floor(Date.now() / 1000), + ...extra + }).catch(err => console.error(`[AgentManager] could not persist agent ${agent.id}:`, err.message)); + } + + async handleDiscovery(agent, payload) { + const discovery = { hostname: payload.hostname || '', ip_addresses: Array.isArray(payload.ip_addresses) ? payload.ip_addresses : [], os: payload.os || '', @@ -94,28 +122,79 @@ class AgentManager { disk_total_gb: payload.disk_total_gb || 0, location: payload.location || 'default' }; + await this.touch(agent, { lastDiscovery: discovery }); + await this.applyDiscoveryToDirectory(agent, discovery); } - handleTelemetry(token, payload) { - const agent = this.agents.get(token); - if (!agent) return; + // An agent runs ON the host it describes, which makes it the most + // authoritative source the directory has -- more so than a hypervisor API or + // a network scan. It previously updated nothing at all: the facts sat on an + // in-memory record and were lost on disconnect. + // + // When the agent is bound to a resource we write that row directly; guessing + // is only for an unbound agent, and then we let the shared reconciler do the + // matching (same MAC/IP/name rules every other source goes through) rather + // than inventing a second matcher here. + async applyDiscoveryToDirectory(agent, discovery) { + try { + const { Resource } = require('../models/resource'); + const metadata = { + os: discovery.os || undefined, + kernel: discovery.kernel || undefined, + cpu: discovery.cpu || undefined, + ram_total_gb: discovery.ram_total_gb || undefined, + disk_total_gb: discovery.disk_total_gb || undefined, + ip: (discovery.ip_addresses || [])[0] || undefined, + agentId: agent.id, + last_seen: Date.now() + }; + // Drop undefined so a field the agent could not determine never + // overwrites a good value already in the directory. + for (const k of Object.keys(metadata)) if (metadata[k] === undefined) delete metadata[k]; - agent.lastSeen = new Date().toISOString(); - agent.telemetry = { - cpu_usage_percent: payload.cpu_usage_percent || 0, - ram_usage_percent: payload.ram_usage_percent || 0, - disk_usage_percent: payload.disk_usage_percent || 0, - zfs_health: payload.zfs_health || 'N/A', - gpu_usage_percent: payload.gpu_usage_percent ?? -1, - timestamp: payload.timestamp || new Date().toISOString() - }; - } + if (agent.resourceId) { + const resource = await Resource.get(agent.resourceId); + if (!resource) return; + const merged = { ...(resource.metadata || {}), ...metadata }; + const sources = new Set(merged.discovery_sources || []); + sources.add('theta-agent'); + merged.discovery_sources = [...sources]; + await resource.update({ metadata: merged, updated_on: Math.floor(Date.now() / 1000) }); + return; + } - handleHeartbeat(token, payload, ws) { - const agent = this.agents.get(token); - if (agent) { - agent.lastSeen = new Date().toISOString(); + if (!discovery.hostname) return; + const { DiscoveryReconciler } = require('../services/discovery_reconciler'); + await DiscoveryReconciler.reconcile('theta-agent', { + resources: [{ + kind: 'host', + name: discovery.hostname, + slug: `agent-${agent.id.slice(0, 8)}`, + metadata: { ...metadata, subType: 'linux' } + }], + edges: [] + }); + } catch (err) { + // Never let a directory write break the agent connection. + console.error(`[AgentManager] discovery -> directory failed for agent ${agent.id}:`, err.message); } + } + + async handleTelemetry(agent, payload) { + await this.touch(agent, { + lastTelemetry: { + cpu_usage_percent: payload.cpu_usage_percent || 0, + ram_usage_percent: payload.ram_usage_percent || 0, + disk_usage_percent: payload.disk_usage_percent || 0, + zfs_health: payload.zfs_health || 'N/A', + gpu_usage_percent: payload.gpu_usage_percent ?? -1, + timestamp: payload.timestamp || new Date().toISOString() + } + }); + } + + async handleHeartbeat(agent, payload, ws) { + await this.touch(agent); try { ws.send(JSON.stringify({ type: 'heartbeat_ack', @@ -124,57 +203,51 @@ class AgentManager { } catch (e) {} } - handleResponse(token, payload) { - const agent = this.agents.get(token); - if (agent) { - agent.lastSeen = new Date().toISOString(); - agent.lastResponse = { + async handleResponse(agent, payload) { + const state = this.live.get(agent.id); + if (state) { + state.lastResponse = { status: payload.status || 'ok', message: payload.message || '', output: payload.output || '', timestamp: new Date().toISOString() }; } + await this.touch(agent); } - sendCommand(token, commandType, payload = {}, isHighRisk = false) { - const agent = this.agents.get(token); - if (!agent || !agent.ws || agent.ws.readyState !== 1) { - throw new Error(`Agent with token "${token}" is not connected`); + async sendCommand(agent, commandType, payload = {}, isHighRisk = false) { + const state = this.live.get(agent.id); + if (!state || !state.ws || state.ws.readyState !== 1) { + throw new Error(`Agent "${agent.name}" is not connected`); } const finalPayload = { ...payload }; - if (isHighRisk) { - finalPayload.signature = this.signPayload(finalPayload); - } + if (isHighRisk) finalPayload.signature = await this.signPayload(finalPayload); - const message = { - type: commandType, - payload: finalPayload - }; - - agent.ws.send(JSON.stringify(message)); + const message = { type: commandType, payload: finalPayload }; + state.ws.send(JSON.stringify(message)); return message; } - getConnectedAgents() { - const list = []; - const now = new Date(); - for (const [token, agent] of this.agents.entries()) { - list.push({ - token, - hostname: agent.hostname, - ipAddress: agent.ipAddress, - connectedAt: agent.connectedAt, - lastSeen: agent.lastSeen, - discovery: agent.discovery, - telemetry: agent.telemetry, - lastResponse: agent.lastResponse || null, - isOnline: (now - new Date(agent.lastSeen)) < 90000 - }); - } - return list; + // Live view for one agent, for merging into its row. + liveState(agentId) { + const state = this.live.get(agentId); + if (!state) return { connected: false, lastResponse: null }; + return { + connected: !!(state.ws && state.ws.readyState === 1), + ipAddress: state.ipAddress, + connectedAt: state.connectedAt, + lastResponse: state.lastResponse || null + }; + } + + // Every enrolled agent, connected or not. + async listAgents() { + const rows = await Agent.list(); + return rows.map(a => a.toPublic(this.liveState(a.id))); } } module.exports = new AgentManager(); +module.exports.AgentManager = AgentManager; diff --git a/nodejs/views/directory.ejs b/nodejs/views/directory.ejs index 722f479..0f2a994 100644 --- a/nodejs/views/directory.ejs +++ b/nodejs/views/directory.ejs @@ -33,6 +33,10 @@
+
+ + +
+
+
+ + +
Links the agent to a Directory host, so its status and metrics attach to that resource.
+
+
+ +
+ + + + + +
- +
- - +
@@ -1584,12 +1769,9 @@
- +
- - +
@@ -1659,6 +1841,7 @@
+ `; app.modal.open({ @@ -1667,9 +1850,78 @@ size: 'lg' }); + // Only hosts can carry an agent -- the API rejects anything else, so don't + // offer it here. + const $sel = $('#agent-enroll-resource').empty(); + $sel.append(''); + rawResources + .filter(r => r.kind === 'host') + .sort((a, b) => (a.name || '').localeCompare(b.name || '')) + .forEach(r => { + const taken = agentsByResource[r.id] ? ' — already has an agent' : ''; + $sel.append($('