sec: authenticate theta-agent enrollment; directory + discovery fixes (v1.29.0)
SECURITY /api/agent/ws authenticated nothing. There was no agent registry, so any client reaching the SSO could register as a node, publish discovery and 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 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. Tokens are minted by POST /api/agent/enroll and shown once. Revoke and rotate drop the live socket immediately. All agent actions are audited. The Ed25519 command-signing key was generated in the AgentManager constructor, so it changed on every restart and the public_key pinned in an agent's agent.yml stopped matching. It now lives in OpenBao at secret/agent/signing-key; if it cannot be loaded the SSO refuses to send high-risk commands rather than signing with a key no agent has seen. DIRECTORY Agents bind to a host resource instead of being matched by hostname, and a bound agent's discovery is written onto that resource -- previously the one source running ON the host contributed nothing to the directory. The resource tree is collapsible, with state persisted per browser. DISCOVERY The Proxmox plugin zipped MACs and IPs from two flat lists by index, attributing addresses to the wrong NIC on multi-NIC guests. NICs are now keyed by MAC. Adds an endpoint resource parenting each node, sourceId/ vmid/node identity, container-interface filtering, node IP/MAC, and offline-node handling. The reconciler could make a resource its own parent, named hosts after their MAC address, had a dead isIp() regex (\\. matches a backslash), merged across kinds, and re-read the whole inventory per resource. Dockerfile.test-runner never copied nodejs/plugins, so every plugin test suite failed in CI as "Cannot find module". Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
@@ -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 };
|
||||
+172
-99
@@ -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;
|
||||
|
||||
Reference in New Issue
Block a user