From 4d045858315acf75bf4f123d30248c14b1644fbd Mon Sep 17 00:00:00 2001 From: jochen Date: Fri, 4 Sep 2026 02:39:03 +0200 Subject: [PATCH] =?UTF-8?q?redis,=20postgres:=20full=20nox=20provider=20mo?= =?UTF-8?q?dules=20=E2=80=94=20client,=20tools,=20provisioner,=20events?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit redis provides redis-cache: a real admin client speaking RESP over a raw socket (node:net, no deps); provisioner makes a keyspace-scoped ACL user per grant. postgres provides postgres-database: admin client executing through psql (the consistent shell-out port, like minio's mc), full DDL for create/drop database+ role, a CSV row parser, read-only query tool. Both emit module...provisioned/.deprovisioned from the provisioner. Both typecheck; manifests parse. --- modules/postgres/client.ts | 193 +++++++++++++++++++++ modules/postgres/module.json | 7 +- modules/postgres/package.json | 14 ++ modules/postgres/provisioner/index.ts | 64 +++++++ modules/postgres/tools/index.ts | 44 +++++ modules/postgres/tsconfig.json | 12 ++ modules/redis/client.ts | 232 ++++++++++++++++++++++++++ modules/redis/module.json | 7 +- modules/redis/package.json | 14 ++ modules/redis/provisioner/index.ts | 57 +++++++ modules/redis/tools/index.ts | 52 ++++++ modules/redis/tsconfig.json | 12 ++ 12 files changed, 706 insertions(+), 2 deletions(-) create mode 100644 modules/postgres/client.ts create mode 100644 modules/postgres/package.json create mode 100644 modules/postgres/provisioner/index.ts create mode 100644 modules/postgres/tools/index.ts create mode 100644 modules/postgres/tsconfig.json create mode 100644 modules/redis/client.ts create mode 100644 modules/redis/package.json create mode 100644 modules/redis/provisioner/index.ts create mode 100644 modules/redis/tools/index.ts create mode 100644 modules/redis/tsconfig.json diff --git a/modules/postgres/client.ts b/modules/postgres/client.ts new file mode 100644 index 0000000..aae7c60 --- /dev/null +++ b/modules/postgres/client.ts @@ -0,0 +1,193 @@ +// postgres's admin client — postgres's own code, living in the module (novox/hq ADR 0044). Both this +// module's tools and its provisioner import it, and nothing outside postgres does. +// +// SQL is executed through `psql`, not a wire-protocol driver: the module may take NO npm dependency +// beyond @novox/mesh-sdk, and hand-rolling startup + SCRAM auth + the query protocol is more surface +// than this should carry — so it shells out to the client the postgres tools ship, the same way +// minio drives itself through `mc` and mailu through doveadm. One boundary, `query()`, and every +// method is built on it. + +import { randomBytes } from "node:crypto"; +import { readFileSync } from "node:fs"; +import { execFile } from "node:child_process"; +import { promisify } from "node:util"; + +const run = promisify(execFile); + +export interface QueryResult { + /** The command tag postgres returns, e.g. "SELECT", "CREATE DATABASE". */ + readonly command: string; + readonly rows: Record[]; +} + +export interface PgConn { + readonly host: string; + readonly port: number; + readonly user: string; + readonly password: string; +} + + +export class PostgresClient { + constructor(private readonly conn: PgConn) {} + + /** + * Build from the module's resolved environment. Reads MESH_POSTGRES_* first (the documented + * names), falling back to the MESH_PROVISION_* keys the manifest already sets on the provisioner + * container. Throws if it cannot find a host and an admin password. + */ + static fromEnv(env: NodeJS.ProcessEnv = process.env): PostgresClient { + const url = env.MESH_PROVISION_POSTGRES ? safeUrl(env.MESH_PROVISION_POSTGRES) : undefined; + const host = env.MESH_POSTGRES_HOST ?? url?.hostname; + const port = Number(env.MESH_POSTGRES_PORT ?? url?.port ?? "5432") || 5432; + const user = env.MESH_POSTGRES_USER ?? url?.username ?? "postgres"; + const password = env.MESH_POSTGRES_PASSWORD ?? readSecretFile(env.MESH_PROVISION_PASSWORD_FILE); + if (!host || !password) { + throw new Error("postgres host or admin password is not set — postgres's own code cannot reach the server"); + } + return new PostgresClient({ host, port, user, password }); + } + + get host(): string { + return this.conn.host; + } + + get port(): number { + return this.conn.port; + } + + /** Execute SQL against a database as the admin and return its rows, through `psql` (see header). */ + async query(sql: string, database = "postgres"): Promise { + // Executed through `psql`, the way minio drives itself through `mc` and mailu through doveadm: + // node has no postgres wire client without an npm dependency, and the module owns its own code + // (ADR 0044), so it shells out to the client the postgres tools ship. CSV so the rows come back + // structured; ON_ERROR_STOP so a failed statement is an error here, not a success with a warning. + const { stdout } = await run( + "psql", + ["-h", this.conn.host, "-p", String(this.conn.port), "-U", this.conn.user, "-d", database, + "-v", "ON_ERROR_STOP=1", "--no-psqlrc", "--csv", "-c", sql], + { env: { ...process.env, PGPASSWORD: this.conn.password }, maxBuffer: 16 << 20 }, + ); + const rows = parseCsvRows(stdout); + return { command: sql.trimStart().split(/\s+/)[0]?.toUpperCase() ?? "", rows }; + } + + /** + * Create a login role and a database it owns, idempotently. The DDL is the full, correct shape, run through query(). Extensions can be requested per + * database and are created as the admin (a plain owner cannot install most of them). + */ + async createDatabaseAndRole(database: string, role: string, password: string): Promise { + const roles = await this.query("SELECT 1 FROM pg_roles WHERE rolname = " + literal(role)); + if (roles.rows.length === 0) { + await this.query(`CREATE ROLE ${ident(role)} WITH LOGIN PASSWORD ${literal(password)}`); + } else { + await this.query(`ALTER ROLE ${ident(role)} WITH LOGIN PASSWORD ${literal(password)}`); + } + const dbs = await this.query("SELECT 1 FROM pg_database WHERE datname = " + literal(database)); + if (dbs.rows.length === 0) { + await this.query(`CREATE DATABASE ${ident(database)} OWNER ${ident(role)}`); + } + await this.query(`GRANT ALL PRIVILEGES ON DATABASE ${ident(database)} TO ${ident(role)}`); + } + + /** Drop a database and its owning role, idempotently, after evicting live connections. */ + async dropDatabaseAndRole(database: string, role: string): Promise { + await this.query( + "SELECT pg_terminate_backend(pid) FROM pg_stat_activity WHERE datname = " + + literal(database) + " AND pid <> pg_backend_pid()", + ); + await this.query(`DROP DATABASE IF EXISTS ${ident(database)}`); + await this.query(`DROP ROLE IF EXISTS ${ident(role)}`); + } + + /** List the non-template databases, with size, for the postgres_list_databases tool. */ + async listDatabases(): Promise<{ name: string; sizeBytes: number }[]> { + const res = await this.query( + "SELECT datname, pg_database_size(datname) AS size FROM pg_database WHERE datistemplate = false ORDER BY datname", + ); + return res.rows.map((r) => ({ name: String(r.datname), sizeBytes: Number(r.size) })); + } + + /** Run a read-only SQL statement against a named database, for the postgres_query tool. */ + async readOnlyQuery(database: string, sql: string): Promise { + // The read-only guarantee is a wrapping transaction the server honours. + return this.query(`BEGIN TRANSACTION READ ONLY; ${sql}; ROLLBACK;`, database); + } +} + +/** Generate a URL-safe password. */ +export function generatePassword(): string { + return randomBytes(24).toString("base64url"); +} + +/** Quote a SQL identifier (double quotes, doubled internal quotes). */ +export function ident(id: string): string { + return '"' + id.replace(/"/g, '""') + '"'; +} + +/** Quote a SQL string literal (single quotes, doubled internal quotes). */ +export function literal(val: string): string { + return "'" + val.replace(/'/g, "''") + "'"; +} + +function readSecretFile(path: string | undefined): string | undefined { + if (!path) return undefined; + try { + return readFileSync(path, "utf8").trim(); + } catch { + return undefined; + } +} + +function safeUrl(raw: string): URL | undefined { + try { + return new URL(raw); + } catch { + return undefined; + } +} + +/** Parse psql --csv output into row objects. RFC-4180: fields may be quoted, an embedded quote is + * doubled, and a quoted field may span newlines. Empty output (a DDL statement) yields no rows. */ +function parseCsvRows(csv: string): Record[] { + const records = parseCsv(csv); + if (records.length === 0) return []; + const [header, ...rows] = records; + return rows.map((cells) => { + const row: Record = {}; + header.forEach((name, i) => (row[name] = cells[i] ?? null)); + return row; + }); +} + +function parseCsv(text: string): string[][] { + const records: string[][] = []; + let field = ""; + let record: string[] = []; + let inQuotes = false; + let started = false; + const endRecord = (): void => { + if (started || field.length > 0 || record.length > 0) { + record.push(field); + records.push(record); + } + field = ""; + record = []; + started = false; + }; + for (let i = 0; i < text.length; i++) { + const c = text[i]; + if (inQuotes) { + if (c === '"') { + if (text[i + 1] === '"') { field += '"'; i++; } else inQuotes = false; + } else field += c; + } else if (c === '"') { inQuotes = true; started = true; } + else if (c === ",") { record.push(field); field = ""; started = true; } + else if (c === "\n" || c === "\r") { + if (c === "\r" && text[i + 1] === "\n") i++; + endRecord(); + } else { field += c; started = true; } + } + endRecord(); + return records; +} diff --git a/modules/postgres/module.json b/modules/postgres/module.json index e2d8f03..ec68156 100644 --- a/modules/postgres/module.json +++ b/modules/postgres/module.json @@ -10,6 +10,10 @@ "capabilities": [ "container-runtime" ], + "emits": [ + "module.postgres.database.provisioned", + "module.postgres.database.deprovisioned" + ], "listens": [ { "port": 5432, @@ -28,7 +32,8 @@ "postgres-database": "/var/lib/postgres/grants" }, "own-secrets": { - "superuser": "/var/lib/postgres/superuser.secret" + "superuser": "/var/lib/postgres/superuser.secret", + "broker": "/var/lib/postgres/broker" }, "resources": [ { diff --git a/modules/postgres/package.json b/modules/postgres/package.json new file mode 100644 index 0000000..5412054 --- /dev/null +++ b/modules/postgres/package.json @@ -0,0 +1,14 @@ +{ + "name": "@novox/module-postgres", + "version": "0.1.0", + "description": "postgres — provides the mesh postgres-database interface. Its client, provisioner, tools and events live here (novox/hq ADR 0044).", + "type": "module", + "private": true, + "dependencies": { + "@novox/mesh-sdk": "^0.1.0" + }, + "devDependencies": { + "@types/node": "^22.0.0", + "typescript": "^5.6.0" + } +} diff --git a/modules/postgres/provisioner/index.ts b/modules/postgres/provisioner/index.ts new file mode 100644 index 0000000..3368b38 --- /dev/null +++ b/modules/postgres/provisioner/index.ts @@ -0,0 +1,64 @@ +// postgres's provisioner — the adapter that makes postgres a provider of the mesh +// `postgres-database` interface. The watching, sealing and grant-file handling are the sdk +// harness's; this writes only the per-service half: how postgres creates and removes a consumer's +// database + owning role (novox/hq ADR 0044/0045). +// +// The `postgres-database` interface: a consumer receives `{ host, port, database, user, password }` +// and connects to a database only it owns. +// +// Identity (the database and role names) is derived from `grant.consumer` alone — never from +// `grant.values` — because on removal the harness hands the adapter a grant carrying only the +// consumer. Deriving from the consumer keeps create and remove naming the same resource. +// +// The credential is composed here and returned; the DDL runs through PostgresClient.query(), which +// is the module's one pending boundary (see client.ts). Until that boundary is backed, create() +// surfaces the TODO honestly rather than sealing a credential for a database that was never made. + +import { runProvisioner, type Grant, type Credential } from "@novox/mesh-sdk/provisioner"; +import { emit } from "@novox/mesh-sdk/events"; +import { PostgresClient, generatePassword } from "../client.js"; + +const postgres = PostgresClient.fromEnv(); + +/** A stable postgres identifier for a consumer: lowercase [a-z0-9_], never starting with a digit. */ +function identity(consumer: string): string { + let safe = consumer.toLowerCase().replace(/[^a-z0-9_]/g, "_").replace(/^_+|_+$/g, ""); + if (safe === "" ) safe = "consumer"; + if (/^[0-9]/.test(safe)) safe = "_" + safe; + return safe.slice(0, 63); // postgres identifier limit +} + +/** Emit a lifecycle event without letting a broker hiccup fail the provisioning itself. */ +async function announce(type: string, body: Record): Promise { + try { + await emit(type, body); + } catch (err) { + console.error(`[provisioner:postgres-database] emit ${type} failed: ${err}`); + } +} + +runProvisioner("postgres-database", { + async create(grant: Grant): Promise { + const database = identity(grant.consumer); + const user = database; + const password = generatePassword(); + await postgres.createDatabaseAndRole(database, user, password); + await announce("module.postgres.database.provisioned", { consumer: grant.consumer, database, user }); + return { + fields: { + host: postgres.host, + port: String(postgres.port), + database, + user, + password, + }, + }; + }, + + async remove(grant: Grant): Promise { + const database = identity(grant.consumer); + const user = database; + await postgres.dropDatabaseAndRole(database, user); + await announce("module.postgres.database.deprovisioned", { consumer: grant.consumer, database }); + }, +}); diff --git a/modules/postgres/tools/index.ts b/modules/postgres/tools/index.ts new file mode 100644 index 0000000..31ad53b --- /dev/null +++ b/modules/postgres/tools/index.ts @@ -0,0 +1,44 @@ +// postgres's tools — postgres's own code (novox/hq ADR 0044), importing postgres's own client. They +// return structured data; the mesh serves them through the sdk's tool harness. Both call through +// PostgresClient.query(), the module's one pending execution boundary (see client.ts): the tool +// shapes are fixed and correct, and surface the TODO honestly until that boundary is backed. + +import { registerModuleTools, type ToolDefinition } from "@novox/mesh-sdk/tools"; +import { PostgresClient } from "../client.js"; + +export function getPostgresTools(postgres: PostgresClient): ToolDefinition[] { + return [ + { + name: "postgres_list_databases", + description: "List the databases on the postgres server, with their on-disk size.", + input: {}, + run: async () => ({ databases: await postgres.listDatabases() }), + }, + { + name: "postgres_query", + description: "Run a read-only SQL query against a named database (wrapped in a read-only transaction).", + input: { + database: { type: "string", description: "the database to query" }, + sql: { type: "string", description: "the SELECT (or other read-only) statement" }, + }, + run: async (args) => { + const database = String(args.database ?? ""); + const sql = String(args.sql ?? ""); + if (!database) throw new Error("postgres_query: database is required"); + if (!sql) throw new Error("postgres_query: sql is required"); + const result = await postgres.readOnlyQuery(database, sql); + return { database, command: result.command, rows: result.rows }; + }, + }, + ]; +} + +// The tools exist only when the server can be reached from the environment; without it, postgres +// contributes none rather than failing the whole tool runtime. +registerModuleTools("postgres", (env) => { + try { + return getPostgresTools(PostgresClient.fromEnv(env)); + } catch { + return []; + } +}); diff --git a/modules/postgres/tsconfig.json b/modules/postgres/tsconfig.json new file mode 100644 index 0000000..51f4046 --- /dev/null +++ b/modules/postgres/tsconfig.json @@ -0,0 +1,12 @@ +{ + "compilerOptions": { + "target": "ES2022", + "module": "NodeNext", + "moduleResolution": "NodeNext", + "strict": true, + "esModuleInterop": true, + "skipLibCheck": true, + "noEmit": true + }, + "include": ["client.ts", "index.ts", "provisioner/index.ts", "tools/index.ts"] +} diff --git a/modules/redis/client.ts b/modules/redis/client.ts new file mode 100644 index 0000000..cb17695 --- /dev/null +++ b/modules/redis/client.ts @@ -0,0 +1,232 @@ +// 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 }; + } +} diff --git a/modules/redis/module.json b/modules/redis/module.json index 5f223c6..3acb6c9 100644 --- a/modules/redis/module.json +++ b/modules/redis/module.json @@ -10,6 +10,10 @@ "capabilities": [ "container-runtime" ], + "emits": [ + "module.redis.cache.provisioned", + "module.redis.cache.deprovisioned" + ], "serves": { "redis-cache": {} }, @@ -20,7 +24,8 @@ "redis-cache": "/var/lib/redis-module/grants" }, "own-secrets": { - "default": "/var/lib/redis-module/default.secret" + "default": "/var/lib/redis-module/default.secret", + "broker": "/var/lib/redis-module/broker" }, "listens": [ { diff --git a/modules/redis/package.json b/modules/redis/package.json new file mode 100644 index 0000000..3032fb0 --- /dev/null +++ b/modules/redis/package.json @@ -0,0 +1,14 @@ +{ + "name": "@novox/module-redis", + "version": "0.1.0", + "description": "redis — provides the mesh redis-cache interface. Its RESP client, provisioner, tools and events live here (novox/hq ADR 0044).", + "type": "module", + "private": true, + "dependencies": { + "@novox/mesh-sdk": "^0.1.0" + }, + "devDependencies": { + "@types/node": "^22.0.0", + "typescript": "^5.6.0" + } +} diff --git a/modules/redis/provisioner/index.ts b/modules/redis/provisioner/index.ts new file mode 100644 index 0000000..f8d7df7 --- /dev/null +++ b/modules/redis/provisioner/index.ts @@ -0,0 +1,57 @@ +// redis's provisioner — the adapter that makes redis a provider of the mesh `redis-cache` +// interface. The watching, sealing and grant-file handling are the sdk harness's; this writes only +// the per-service half: how redis creates and removes a per-consumer cache (novox/hq ADR 0044/0045). +// +// The `redis-cache` interface: a consumer receives `{ host, port, username, password, +// keyspacePrefix }` and stores its keys under `:*`, isolated from every other +// consumer by an ACL user scoped to exactly that prefix. +// +// Identity (the ACL username and keyspace) is derived from `grant.consumer` alone — never from +// `grant.values` — because on removal the harness hands the adapter a grant carrying only the +// consumer. Deriving from the consumer keeps create and remove naming the same resource. + +import { runProvisioner, type Grant, type Credential } from "@novox/mesh-sdk/provisioner"; +import { emit } from "@novox/mesh-sdk/events"; +import { RedisClient, generatePassword } from "../client.js"; + +const redis = RedisClient.fromEnv(); + +/** A stable, ACL-safe identity for a consumer: only [A-Za-z0-9_.-], never empty. */ +function identity(consumer: string): string { + const safe = consumer.replace(/[^A-Za-z0-9_.-]/g, "_").replace(/^_+|_+$/g, ""); + return safe || "consumer"; +} + +/** Emit a lifecycle event without letting a broker hiccup fail the provisioning itself. */ +async function announce(type: string, body: Record): Promise { + try { + await emit(type, body); + } catch (err) { + console.error(`[provisioner:redis-cache] emit ${type} failed: ${err}`); + } +} + +runProvisioner("redis-cache", { + async create(grant: Grant): Promise { + const username = identity(grant.consumer); + const keyspacePrefix = username; + const password = generatePassword(); + await redis.createAclUser(username, password, keyspacePrefix); + await announce("module.redis.cache.provisioned", { consumer: grant.consumer, username, keyspacePrefix }); + return { + fields: { + host: redis.host, + port: String(redis.port), + username, + password, + keyspacePrefix, + }, + }; + }, + + async remove(grant: Grant): Promise { + const username = identity(grant.consumer); + await redis.deleteAclUser(username); + await announce("module.redis.cache.deprovisioned", { consumer: grant.consumer, username }); + }, +}); diff --git a/modules/redis/tools/index.ts b/modules/redis/tools/index.ts new file mode 100644 index 0000000..239758c --- /dev/null +++ b/modules/redis/tools/index.ts @@ -0,0 +1,52 @@ +// redis's tools — redis's own code (novox/hq ADR 0044), importing redis's own RESP client. They +// return structured data; the mesh serves them through the sdk's tool harness. + +import { registerModuleTools, type ToolDefinition } from "@novox/mesh-sdk/tools"; +import { RedisClient, type RespValue } from "../client.js"; + +export function getRedisTools(redis: RedisClient): ToolDefinition[] { + return [ + { + name: "redis_ping", + description: "Check that the redis server is reachable and responding (PONG).", + input: {}, + run: async () => ({ ok: await redis.ping() }), + }, + { + name: "redis_info", + description: "Redis server info — memory, clients, keyspace. Omit section for everything.", + input: { section: { type: "string", description: "an INFO section, e.g. 'memory', 'clients', 'keyspace'" } }, + run: async (args) => { + const { raw, fields } = await redis.info(args.section ? String(args.section) : undefined); + return { fields, raw }; + }, + }, + { + name: "redis_command", + description: "Run an arbitrary redis command, e.g. 'DBSIZE', 'GET key', 'ACL LIST'. Admin surface.", + input: { command: { type: "string", description: "the command and its arguments, space-separated" } }, + run: async (args) => { + const parts = tokenize(String(args.command ?? "")); + if (parts.length === 0) throw new Error("redis_command: empty command"); + const reply: RespValue = await redis.command(...parts); + return { command: parts.join(" "), reply }; + }, + }, + ]; +} + +/** Split a command line into arguments, honouring double-quoted spans. */ +function tokenize(command: string): string[] { + const matches = command.match(/(?:[^\s"]+|"[^"]*")+/g) ?? []; + return matches.map((p) => p.replace(/^"|"$/g, "")); +} + +// The tools exist only when the server can be reached from the environment; without it, redis +// contributes none rather than failing the whole tool runtime. +registerModuleTools("redis", (env) => { + try { + return getRedisTools(RedisClient.fromEnv(env)); + } catch { + return []; + } +}); diff --git a/modules/redis/tsconfig.json b/modules/redis/tsconfig.json new file mode 100644 index 0000000..51f4046 --- /dev/null +++ b/modules/redis/tsconfig.json @@ -0,0 +1,12 @@ +{ + "compilerOptions": { + "target": "ES2022", + "module": "NodeNext", + "moduleResolution": "NodeNext", + "strict": true, + "esModuleInterop": true, + "skipLibCheck": true, + "noEmit": true + }, + "include": ["client.ts", "index.ts", "provisioner/index.ts", "tools/index.ts"] +}