redis, postgres: full nox provider modules — client, tools, provisioner, events
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.<x>.<thing>.provisioned/.deprovisioned from the provisioner. Both typecheck; manifests parse.
This commit is contained in:
@@ -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<string, unknown>[];
|
||||||
|
}
|
||||||
|
|
||||||
|
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<QueryResult> {
|
||||||
|
// 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<void> {
|
||||||
|
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<void> {
|
||||||
|
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<QueryResult> {
|
||||||
|
// 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<string, unknown>[] {
|
||||||
|
const records = parseCsv(csv);
|
||||||
|
if (records.length === 0) return [];
|
||||||
|
const [header, ...rows] = records;
|
||||||
|
return rows.map((cells) => {
|
||||||
|
const row: Record<string, unknown> = {};
|
||||||
|
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;
|
||||||
|
}
|
||||||
@@ -10,6 +10,10 @@
|
|||||||
"capabilities": [
|
"capabilities": [
|
||||||
"container-runtime"
|
"container-runtime"
|
||||||
],
|
],
|
||||||
|
"emits": [
|
||||||
|
"module.postgres.database.provisioned",
|
||||||
|
"module.postgres.database.deprovisioned"
|
||||||
|
],
|
||||||
"listens": [
|
"listens": [
|
||||||
{
|
{
|
||||||
"port": 5432,
|
"port": 5432,
|
||||||
@@ -28,7 +32,8 @@
|
|||||||
"postgres-database": "/var/lib/postgres/grants"
|
"postgres-database": "/var/lib/postgres/grants"
|
||||||
},
|
},
|
||||||
"own-secrets": {
|
"own-secrets": {
|
||||||
"superuser": "/var/lib/postgres/superuser.secret"
|
"superuser": "/var/lib/postgres/superuser.secret",
|
||||||
|
"broker": "/var/lib/postgres/broker"
|
||||||
},
|
},
|
||||||
"resources": [
|
"resources": [
|
||||||
{
|
{
|
||||||
|
|||||||
@@ -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"
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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<string, string>): Promise<void> {
|
||||||
|
try {
|
||||||
|
await emit(type, body);
|
||||||
|
} catch (err) {
|
||||||
|
console.error(`[provisioner:postgres-database] emit ${type} failed: ${err}`);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
runProvisioner("postgres-database", {
|
||||||
|
async create(grant: Grant): Promise<Credential> {
|
||||||
|
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<void> {
|
||||||
|
const database = identity(grant.consumer);
|
||||||
|
const user = database;
|
||||||
|
await postgres.dropDatabaseAndRole(database, user);
|
||||||
|
await announce("module.postgres.database.deprovisioned", { consumer: grant.consumer, database });
|
||||||
|
},
|
||||||
|
});
|
||||||
@@ -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 [];
|
||||||
|
}
|
||||||
|
});
|
||||||
@@ -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"]
|
||||||
|
}
|
||||||
@@ -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<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.
|
||||||
|
*/
|
||||||
|
async createAclUser(username: string, password: string, keyspacePrefix: string): Promise<void> {
|
||||||
|
await this.command("ACL", "SETUSER", username, "reset", "on", `>${password}`, `~${keyspacePrefix}:*`, "+@all");
|
||||||
|
}
|
||||||
|
|
||||||
|
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 };
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -10,6 +10,10 @@
|
|||||||
"capabilities": [
|
"capabilities": [
|
||||||
"container-runtime"
|
"container-runtime"
|
||||||
],
|
],
|
||||||
|
"emits": [
|
||||||
|
"module.redis.cache.provisioned",
|
||||||
|
"module.redis.cache.deprovisioned"
|
||||||
|
],
|
||||||
"serves": {
|
"serves": {
|
||||||
"redis-cache": {}
|
"redis-cache": {}
|
||||||
},
|
},
|
||||||
@@ -20,7 +24,8 @@
|
|||||||
"redis-cache": "/var/lib/redis-module/grants"
|
"redis-cache": "/var/lib/redis-module/grants"
|
||||||
},
|
},
|
||||||
"own-secrets": {
|
"own-secrets": {
|
||||||
"default": "/var/lib/redis-module/default.secret"
|
"default": "/var/lib/redis-module/default.secret",
|
||||||
|
"broker": "/var/lib/redis-module/broker"
|
||||||
},
|
},
|
||||||
"listens": [
|
"listens": [
|
||||||
{
|
{
|
||||||
|
|||||||
@@ -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"
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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 `<keyspacePrefix>:*`, 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<string, string>): Promise<void> {
|
||||||
|
try {
|
||||||
|
await emit(type, body);
|
||||||
|
} catch (err) {
|
||||||
|
console.error(`[provisioner:redis-cache] emit ${type} failed: ${err}`);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
runProvisioner("redis-cache", {
|
||||||
|
async create(grant: Grant): Promise<Credential> {
|
||||||
|
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<void> {
|
||||||
|
const username = identity(grant.consumer);
|
||||||
|
await redis.deleteAclUser(username);
|
||||||
|
await announce("module.redis.cache.deprovisioned", { consumer: grant.consumer, username });
|
||||||
|
},
|
||||||
|
});
|
||||||
@@ -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 [];
|
||||||
|
}
|
||||||
|
});
|
||||||
@@ -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"]
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user