238 lines
9.3 KiB
TypeScript
238 lines
9.3 KiB
TypeScript
// redis's admin client — redis's own code, living in the module (novox/hq ADR 0039). 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<void> {
|
|
if (this.socket && !this.socket.destroyed) return;
|
|
await new Promise<void>((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<RespValue> {
|
|
await this.connect();
|
|
return this.send(args.map(String));
|
|
}
|
|
|
|
async ping(): Promise<boolean> {
|
|
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<string, string> }> {
|
|
const raw = String((await this.command("INFO", ...(section ? [section] : []))) ?? "");
|
|
const fields: Record<string, string> = {};
|
|
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 `<prefix>:*`, and allowed the ordinary command set. The
|
|
* consumer connects as this user and can touch nothing outside its prefix.
|
|
*
|
|
* Minus the dangerous category: a key pattern confines commands that name keys, and FLUSHALL,
|
|
* FLUSHDB, CONFIG, SHUTDOWN and the rest of `@dangerous` name none — with `+@all` alone a
|
|
* consumer scoped to its own keys could still wipe the server (novox/hq issue 080). KEYS goes
|
|
* with them; SCAN stays, and is what a consumer should use anyway.
|
|
*/
|
|
async createAclUser(username: string, password: string, keyspacePrefix: string): Promise<void> {
|
|
await this.command("ACL", "SETUSER", username, "reset", "on", `>${password}`, `~${keyspacePrefix}:*`, "+@all", "-@dangerous");
|
|
}
|
|
|
|
async deleteAclUser(username: string): Promise<void> {
|
|
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<RespValue> {
|
|
const sock = this.socket;
|
|
if (!sock) return Promise.reject(new Error("redis socket is not connected"));
|
|
return new Promise<RespValue>((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 };
|
|
}
|
|
}
|