// redis's admin client — redis's own code, living in the module (novox/hq ADR 0044). It speaks // RESP directly over a raw TCP socket (node:net) so the module carries no npm dependency beyond // @novox/mesh-sdk: no redis driver, no redis-cli in the image. Both this module's tools and its // provisioner import it, and nothing outside redis does. // // The client keeps a single connection and runs one command at a time in FIFO order — enough for // an admin surface (PING, INFO, ACL SETUSER/DELUSER, arbitrary commands). Replies come back in the // order requests were sent, which is what the queue below relies on. import { createConnection, type Socket } from "node:net"; import { randomBytes } from "node:crypto"; import { readFileSync } from "node:fs"; /** A parsed RESP value. Errors are surfaced as rejected commands, not as this type. */ export type RespValue = string | number | null | RespValue[]; interface Waiter { resolve: (v: RespValue) => void; reject: (e: Error) => void; } export class RedisClient { private socket: Socket | null = null; private buffer = Buffer.alloc(0); private queue: Waiter[] = []; constructor( readonly host: string, readonly port: number, private readonly password: string, private readonly username = "default", ) {} /** * Build from the module's resolved environment. Reads MESH_REDIS_* first (the documented * names), falling back to the MESH_PROVISION_* keys the manifest already sets on the * provisioner container so the module runs unchanged there. Throws if it cannot find a host and * an admin password — the right failure, because without them nothing it does can work. */ static fromEnv(env: NodeJS.ProcessEnv = process.env): RedisClient { const endpoint = env.MESH_PROVISION_REDIS ?? ""; // "host:port" const host = env.MESH_REDIS_HOST ?? (endpoint ? endpoint.split(":")[0] : undefined); const port = Number(env.MESH_REDIS_PORT ?? (endpoint.includes(":") ? endpoint.split(":")[1] : "") ?? "6379") || 6379; const username = env.MESH_REDIS_USERNAME ?? "default"; const password = env.MESH_REDIS_PASSWORD ?? readSecretFile(env.MESH_PROVISION_PASSWORD_FILE); if (!host || !password) { throw new Error("redis host or admin password is not set — redis's own code cannot reach the server"); } return new RedisClient(host, port, password, username); } /** Open the connection (idempotent) and authenticate. Reconnects if the socket has gone away. */ async connect(): Promise { if (this.socket && !this.socket.destroyed) return; await new Promise((resolve, reject) => { const sock = createConnection({ host: this.host, port: this.port }); this.socket = sock; sock.on("data", (chunk: Buffer | string) => this.onData(typeof chunk === "string" ? Buffer.from(chunk) : chunk)); sock.on("error", (err) => { this.failAll(err); reject(err); }); sock.on("close", () => this.failAll(new Error("redis connection closed"))); sock.once("connect", () => resolve()); }); if (this.password) { const args = this.username && this.username !== "default" ? ["AUTH", this.username, this.password] : ["AUTH", this.password]; await this.send(args); } } /** Run one command and return its parsed reply. A RESP error reply rejects the promise. */ async command(...args: (string | number)[]): Promise { await this.connect(); return this.send(args.map(String)); } async ping(): Promise { return (await this.command("PING")) === "PONG"; } /** Server INFO, returned both raw and parsed into the flat key/value map redis emits. */ async info(section?: string): Promise<{ raw: string; fields: Record }> { const raw = String((await this.command("INFO", ...(section ? [section] : []))) ?? ""); const fields: Record = {}; for (const line of raw.split(/\r?\n/)) { if (!line || line.startsWith("#")) continue; const idx = line.indexOf(":"); if (idx > 0) fields[line.slice(0, idx)] = line.slice(idx + 1); } return { raw, fields }; } /** * Create (or reset to a known state) an ACL user scoped to one keyspace prefix. `reset` first * clears any prior rules so the call is idempotent, then the user is enabled with the given * password, confined to keys matching `:*`, and allowed the ordinary command set. The * consumer connects as this user and can touch nothing outside its prefix. */ async createAclUser(username: string, password: string, keyspacePrefix: string): Promise { await this.command("ACL", "SETUSER", username, "reset", "on", `>${password}`, `~${keyspacePrefix}:*`, "+@all"); } async deleteAclUser(username: string): Promise { await this.command("ACL", "DELUSER", username); } close(): void { if (this.socket) { this.socket.destroy(); this.socket = null; } } // --- connection plumbing --- /** Send an already-connected command, queuing its reply against the FIFO of in-flight requests. */ private send(args: string[]): Promise { const sock = this.socket; if (!sock) return Promise.reject(new Error("redis socket is not connected")); return new Promise((resolve, reject) => { this.queue.push({ resolve, reject }); sock.write(encodeCommand(args)); }); } /** Feed incoming bytes to the parser, resolving as many queued replies as the buffer completes. */ private onData(chunk: Buffer): void { this.buffer = Buffer.concat([this.buffer, chunk]); while (this.queue.length > 0) { let parsed: { value: RespValue | Error; next: number } | null; try { parsed = parseReply(this.buffer, 0); } catch (err) { const waiter = this.queue.shift(); waiter?.reject(err instanceof Error ? err : new Error(String(err))); this.buffer = Buffer.alloc(0); continue; } if (!parsed) break; // one full reply not yet in the buffer this.buffer = this.buffer.subarray(parsed.next); const waiter = this.queue.shift(); if (!waiter) break; if (parsed.value instanceof Error) waiter.reject(parsed.value); else waiter.resolve(parsed.value); } } /** A socket error or close fails every pending command and forces a fresh connect next time. */ private failAll(err: Error): void { const pending = this.queue; this.queue = []; for (const waiter of pending) waiter.reject(err); this.buffer = Buffer.alloc(0); if (this.socket) { this.socket.destroy(); this.socket = null; } } } /** Generate a URL-safe password with no RESP-hostile characters. */ export function generatePassword(): string { return randomBytes(24).toString("base64url"); } function readSecretFile(path: string | undefined): string | undefined { if (!path) return undefined; try { return readFileSync(path, "utf8").trim(); } catch { return undefined; } } // --- RESP wire format --- /** Encode a command as a RESP array of bulk strings. */ function encodeCommand(args: string[]): Buffer { let head = `*${args.length}\r\n`; for (const a of args) head += `$${Buffer.byteLength(a)}\r\n${a}\r\n`; return Buffer.from(head, "utf8"); } /** * Parse one RESP reply starting at `i`. Returns the value and the offset just past it, or null if * the buffer does not yet hold a complete reply (the caller waits for more bytes). A `-` error * reply is returned as an Error value; the client turns that into a rejection. */ function parseReply(buf: Buffer, i: number): { value: RespValue | Error; next: number } | null { if (i >= buf.length) return null; const type = buf[i]; const eol = buf.indexOf("\r\n", i + 1); if (eol === -1) return null; // header line not yet complete const line = buf.toString("utf8", i + 1, eol); const after = eol + 2; switch (type) { case 0x2b: // '+' simple string return { value: line, next: after }; case 0x2d: // '-' error return { value: new Error(line), next: after }; case 0x3a: // ':' integer return { value: Number(line), next: after }; case 0x24: { // '$' bulk string const len = Number(line); if (len === -1) return { value: null, next: after }; const end = after + len; if (buf.length < end + 2) return null; // body not fully arrived return { value: buf.toString("utf8", after, end), next: end + 2 }; } case 0x2a: { // '*' array const count = Number(line); if (count === -1) return { value: null, next: after }; const arr: RespValue[] = []; let cursor = after; for (let k = 0; k < count; k++) { const el = parseReply(buf, cursor); if (!el) return null; // array not fully arrived if (el.value instanceof Error) throw el.value; arr.push(el.value); cursor = el.next; } return { value: arr, next: cursor }; } default: return { value: new Error(`unexpected RESP type byte 0x${type.toString(16)}`), next: after }; } }