Files
coreci-chat/apps/control-plane/ws-server.ts
T
CIAgent 86f7dcc10b feat(P4): Wave D relay agent — Go binary, install script, WebSocket, SSH whitelist hook
REQ-010: modular install script (detect_os/install_binary/write_systemd_unit/register_target)
REQ-011: outbound WebSocket from relay agent to SaaS
REQ-012: register with tenant + target metadata
REQ-013: heartbeat + exponential backoff reconnect (max 5 → alert)
REQ-026 partial: SSH whitelist file + CheckCommand hook (G-004 contract lock,
  G-007 shadow exec.Cmd test, G-008 scope statement)
G-009: unsupported-OS abort with actionable error

---ci---
phase: 4
milestone: v0.1
status: execute
---/ci---
2026-08-25 02:03:04 +00:00

306 lines
11 KiB
TypeScript

/**
* control-plane ws-server — standalone WebSocket server for the Relay Agent
* (REQ-011, REQ-012, REQ-013, Wave D Task 2).
*
* Next.js App Router route handlers are HTTP-only (no native WebSocket
* upgrade), so for M1 the relay WebSocket server runs as a small standalone
* HTTP server alongside Next.js. It binds 0.0.0.0:${CORECI_WS_PORT:-3001} and
* upgrades requests to /api/relay/ws.
*
* Lifecycle per connection:
* 1. On upgrade: read Authorization: Bearer <token>, verify the JWT relay
* token via @coreci/secrets verifyRelayToken (G-005). If invalid → close
* with code 4001 (REQ-012 Edge 12).
* 2. On {type:"register"} message: INSERT into targets
* (tenant_id, hostname, os_name, os_version, ip_address, agent_version,
* last_seen_at=now()) under withTenant + audit append (provision event).
* Respond {type:"registered", targetId}.
* 3. On {type:"ping"} message: UPDATE targets SET last_seen_at=now() WHERE
* id=targetId under withTenant. Respond {type:"pong", ts}.
*
* Connected agents are tracked in-memory for the dashboard (Wave E consumes
* this via `getConnectedAgents()`).
*
* Run standalone: `tsx ws-server.ts` (next to `next dev`/`next start`).
* The systemd unit / docker-compose starts both processes.
*/
import { createServer, type IncomingMessage } from "node:http";
import { WebSocketServer, type WebSocket } from "ws";
import { randomUUID } from "node:crypto";
import { verifyRelayToken, RelayTokenError } from "@coreci/secrets";
import { withTenant, appendAudit, type ScopedClient } from "@coreci/db";
import { getDb } from "./lib/db.js";
/** RELAY_TOKEN_SIGNING_KEY — infra bootstrap signing key (G-010 tier a). */
function signingKey(): string {
const k = process.env.RELAY_TOKEN_SIGNING_KEY;
if (!k) throw new Error("RELAY_TOKEN_SIGNING_KEY not configured");
return k;
}
const WS_PATH = "/api/relay/ws";
const WS_PORT = Number(process.env.CORECI_WS_PORT ?? 3001);
// ─── In-memory connected-agent registry (Wave E consumes this) ───────────────
export interface ConnectedAgent {
tenantId: string;
targetId: string;
hostname: string;
osName: string;
connectedAt: number;
lastSeenAt: number;
}
const connectedAgents = new Map<WebSocket, ConnectedAgent>();
/** Snapshot of connected agents (Wave E dashboard reads this). */
export function getConnectedAgents(): ConnectedAgent[] {
return Array.from(connectedAgents.values());
}
// ─── Wire message types ──────────────────────────────────────────────────────
interface RegisterMessage {
type: "register";
tenantId?: string; // ignored — resolved from the verified JWT
hostname: string;
os: string;
osVersion: string;
ip: string;
agentVersion: string;
}
interface PingMessage {
type: "ping";
ts: number;
}
interface RegisteredResponse {
type: "registered";
targetId: string;
}
interface PongResponse {
type: "pong";
ts: number;
}
interface ErrorResponse {
type: "error";
error: string;
code?: number;
detail?: string;
}
// ─── WS server bootstrap ─────────────────────────────────────────────────────
export async function startWsServer(port = WS_PORT): Promise<{ server: ReturnType<typeof createServer>; wss: WebSocketServer }> {
// Ensure the DB is bootstrapped (loads migrations) before accepting connections.
await getDb();
const server = createServer();
const wss = new WebSocketServer({ noServer: true });
server.on("upgrade", (req, socket, head) => {
const url = new URL(req.url ?? "", `http://${req.headers.host ?? "localhost"}`);
if (url.pathname !== WS_PATH) {
socket.destroy();
return;
}
const token = bearerToken(req);
if (!token) {
socket.destroy();
return;
}
let tenantId: string;
try {
const claims = verifyRelayToken(token, signingKey());
tenantId = claims.tenantId;
} catch (err) {
// REQ-012 Edge 12: registration/token failure → close with 4001.
logWarn(`relay ws auth failed: ${err instanceof Error ? err.message : String(err)}`);
socket.destroy();
return;
}
wss.handleUpgrade(req, socket, head, (ws) => {
// Stash the resolved tenantId on the socket for the message handlers.
(ws as WebSocket & { tenantId: string }).tenantId = tenantId;
wss.emit("connection", ws, req);
});
});
wss.on("connection", (ws) => {
const tenantId = (ws as WebSocket & { tenantId: string }).tenantId;
logInfo(`relay ws connected (tenant=${tenantId})`);
ws.on("message", (data) => {
handleMessage(ws, tenantId, data)
.catch((err) => {
logWarn(`relay ws message error: ${err instanceof Error ? err.message : String(err)}`);
safeSend(ws, { type: "error", error: "internal_error", detail: err instanceof Error ? err.message : String(err) } satisfies ErrorResponse);
});
});
ws.on("close", () => {
const agent = connectedAgents.get(ws);
if (agent) {
logInfo(`relay ws closed (tenant=${agent.tenantId} target=${agent.targetId})`);
connectedAgents.delete(ws);
}
});
ws.on("error", (err) => {
logWarn(`relay ws socket error: ${err.message}`);
});
});
return new Promise((resolve, reject) => {
server.on("error", reject);
server.listen(port, "0.0.0.0", () => {
logInfo(`relay ws server listening on ws://0.0.0.0:${port}${WS_PATH}`);
resolve({ server, wss });
});
});
}
// ─── Message handling ────────────────────────────────────────────────────────
async function handleMessage(ws: WebSocket, tenantId: string, data: unknown): Promise<void> {
let msg: Record<string, unknown>;
try {
if (typeof data !== "string" && !(data instanceof Buffer) && !Array.isArray(data)) {
throw new Error("non-text message");
}
const text = data instanceof Buffer ? data.toString("utf8") : Array.isArray(data) ? Buffer.concat(data as Buffer[]).toString("utf8") : String(data);
msg = JSON.parse(text) as Record<string, unknown>;
} catch (err) {
safeSend(ws, { type: "error", error: "invalid_json", detail: err instanceof Error ? err.message : String(err) } satisfies ErrorResponse);
return;
}
switch (msg.type) {
case "register":
await handleRegister(ws, tenantId, msg as unknown as RegisterMessage);
return;
case "ping":
await handlePing(ws, tenantId, msg as unknown as PingMessage);
return;
default:
safeSend(ws, { type: "error", error: `unknown message type: ${String(msg.type)}` } satisfies ErrorResponse);
}
}
interface TargetRow {
id: string;
}
async function handleRegister(ws: WebSocket, tenantId: string, msg: RegisterMessage): Promise<void> {
if (!msg.hostname || !msg.os || !msg.agentVersion) {
safeSend(ws, { type: "error", error: "register: missing required fields (hostname, os, agentVersion)" } satisfies ErrorResponse);
return;
}
// INSERT the target under withTenant (RLS enforces tenant scoping) and
// append a provision audit event in the same transaction (Edge 7: audit
// failure halts the operation, so the target INSERT rolls back too).
const targetId = await withTenant(tenantId, async (c: ScopedClient) => {
const insertRes = await c.query<TargetRow>(
`INSERT INTO targets (tenant_id, hostname, os_name, os_version, ip_address, agent_version, last_seen_at)
VALUES ($1, $2, $3, $4, $5, $6, now())
RETURNING id`,
[tenantId, msg.hostname, msg.os, msg.osVersion || "", msg.ip || null, msg.agentVersion],
);
const id = insertRes.rows[0]?.id;
if (!id) throw new Error("register: target INSERT returned no id");
await appendAudit(c, {
tenantId,
eventType: "provision",
payload: {
action: "relay_target_registered",
hostname: msg.hostname,
os: msg.os,
osVersion: msg.osVersion,
ip: msg.ip,
agentVersion: msg.agentVersion,
targetId: id,
},
targetId: id,
});
return id;
});
// Track in-memory for the dashboard (Wave E).
connectedAgents.set(ws, {
tenantId,
targetId,
hostname: msg.hostname,
osName: msg.os,
connectedAt: Date.now(),
lastSeenAt: Date.now(),
});
const response: RegisteredResponse = { type: "registered", targetId };
safeSend(ws, response);
logInfo(`relay target registered (tenant=${tenantId} target=${targetId} hostname=${msg.hostname})`);
}
async function handlePing(ws: WebSocket, tenantId: string, msg: PingMessage): Promise<void> {
const agent = connectedAgents.get(ws);
if (!agent) {
safeSend(ws, { type: "error", error: "ping before register" } satisfies ErrorResponse);
return;
}
// UPDATE last_seen under withTenant (RLS enforces scoping).
await withTenant(tenantId, async (c: ScopedClient) => {
await c.query(
`UPDATE targets SET last_seen_at = now() WHERE id = $1`,
[agent.targetId],
);
});
agent.lastSeenAt = Date.now();
const response: PongResponse = { type: "pong", ts: msg.ts ?? Date.now() };
safeSend(ws, response);
}
// ─── Helpers ─────────────────────────────────────────────────────────────────
function bearerToken(req: IncomingMessage): string | null {
const header = req.headers["authorization"];
if (!header || typeof header !== "string") return null;
const m = header.match(/^Bearer\s+(.+)$/i);
return m ? m[1]!.trim() : null;
}
function safeSend(ws: WebSocket, msg: RegisteredResponse | PongResponse | ErrorResponse): void {
if (ws.readyState !== ws.OPEN) return;
try {
ws.send(JSON.stringify(msg));
} catch (err) {
logWarn(`relay ws send failed: ${err instanceof Error ? err.message : String(err)}`);
}
}
function logInfo(msg: string): void {
// eslint-disable-next-line no-console
console.log(`[relay-ws] ${msg}`);
}
function logWarn(msg: string): void {
// eslint-disable-next-line no-console
console.warn(`[relay-ws] ${msg}`);
}
// silence unused import in type-only contexts
void RelayTokenError;
void randomUUID;
// ─── Entrypoint (run standalone: `tsx ws-server.ts`) ──────────────────────────
if (import.meta.url === `file://${process.argv[1]}`) {
startWsServer().catch((err) => {
console.error(`[relay-ws] fatal: ${err instanceof Error ? err.message : String(err)}`);
process.exit(1);
});
}