// The InfluxDB API client — influxdb's own code, living in the module (novox/hq ADR 0039). Only // this module's tools import it. Talks to the InfluxDB 2.x HTTP API (/api/v2) with a token. import { readFileSync } from "node:fs"; export interface InfluxHealth { name?: string; status?: string; message?: string; version?: string; } export interface InfluxBucket { id: string; name: string; orgID?: string; retentionSeconds?: number; } /** One permission of an authorization, as InfluxDB represents it: an action on a resource type, * in one org, optionally narrowed to one resource by id (no id = every resource of that type). */ export interface InfluxPermission { action: "read" | "write"; resource: { type: string; orgID?: string; id?: string; name?: string; org?: string }; } /** A v1-compatibility ("legacy") authorization: a username (InfluxDB calls it `token`) and a * password the caller chooses, scoped by permissions. The one credential InfluxDB 2.x lets a * caller set to a value it did not generate — which is what a mesh-minted password needs. */ export interface LegacyAuthorization { id: string; token: string; orgID: string; status?: "active" | "inactive"; description?: string; permissions: InfluxPermission[]; } /** The settings-merged config the mesh delivers (novox/hq ADR 0046): { url, apiKey, token, password, user, ... }. */ function meshConfig(file?: string): Record { if (!file) return {}; try { return JSON.parse(readFileSync(file, "utf8")) as Record; } catch { return {}; } } /** A secret delivered as a file, trimmed; undefined when there is none, so the caller can fall back. */ function tokenFromFile(file?: string): string | undefined { if (!file) return undefined; try { return readFileSync(file, "utf8").trim() || undefined; } catch { return undefined; } } export class InfluxDBClient { readonly baseUrl: string; constructor( url: string, private readonly token: string, readonly org: string, ) { this.baseUrl = url.replace(/\/$/, ""); } /** * Build from the module's resolved environment. The token is the InfluxDB API token (the admin * token the server was initialised with, or a scoped one) — required, since every /api/v2 call * is token-authenticated. The org scopes bucket listing and queries. */ static fromEnv(env: NodeJS.ProcessEnv = process.env): InfluxDBClient { const cfg = meshConfig(env.MESH_INFLUXDB_CONFIG_FILE); const url = cfg.url ?? env.MESH_INFLUXDB_URL ?? `http://127.0.0.1:${env.INFLUXDB_PORT ?? "8086"}`; // The token reaches the process as a file (novox/hq ADR 0086); the environment variable stays // only for a workstation running the tools by hand. const token = cfg.token ?? tokenFromFile(env.MESH_INFLUXDB_TOKEN_FILE) ?? env.MESH_INFLUXDB_TOKEN; if (!token) throw new Error("no InfluxDB token — set MESH_INFLUXDB_TOKEN_FILE"); const org = cfg.org ?? env.MESH_INFLUXDB_ORG ?? "mesh"; return new InfluxDBClient(url, token, org); } private async request(path: string, init?: RequestInit): Promise { const res = await fetch(`${this.baseUrl}${path}`, { ...init, headers: { Authorization: `Token ${this.token}`, ...(init?.headers ?? {}), }, }); if (!res.ok) throw new Error(`InfluxDB API ${path}: ${res.status} ${await res.text()}`); return res; } /** Like request, but the answer is returned whatever its status, for the caller to read. */ private async raw(path: string, init?: RequestInit): Promise { return fetch(`${this.baseUrl}${path}`, { ...init, headers: { Authorization: `Token ${this.token}`, ...(init?.headers ?? {}) }, }); } private async send(path: string, method: string, body?: unknown): Promise { return this.request(path, { method, headers: { "Content-Type": "application/json" }, body: body === undefined ? undefined : JSON.stringify(body), }); } /** The id of the org of this name, or undefined when there is none. */ async orgID(name: string): Promise { const res = await this.raw(`/api/v2/orgs?org=${encodeURIComponent(name)}`); if (res.status === 404) return undefined; if (!res.ok) throw new Error(`InfluxDB API /api/v2/orgs: ${res.status} ${await res.text()}`); const body = (await res.json()) as { orgs?: { id: string; name: string }[] }; return body.orgs?.find((o) => o.name === name)?.id; } /** The bucket of exactly this name in the org, or undefined. */ async findBucket(orgID: string, name: string): Promise { const res = await this.raw(`/api/v2/buckets?orgID=${encodeURIComponent(orgID)}&name=${encodeURIComponent(name)}`); if (res.status === 404) return undefined; if (!res.ok) throw new Error(`InfluxDB API /api/v2/buckets: ${res.status} ${await res.text()}`); const body = (await res.json()) as { buckets?: { id: string; name: string; orgID?: string }[] }; const b = body.buckets?.find((x) => x.name === name); return b ? { id: b.id, name: b.name, orgID: b.orgID } : undefined; } /** Create a bucket that keeps its data for ever — retention is the operator's choice, never the mesh's. */ async createBucket(orgID: string, name: string, description: string): Promise { const b = (await (await this.send("/api/v2/buckets", "POST", { orgID, name, description, retentionRules: [], })).json()) as { id: string; name: string; orgID?: string }; return { id: b.id, name: b.name, orgID: b.orgID }; } /** The v1 authorization whose username is exactly this, or undefined. */ async findLegacy(username: string): Promise { const path = `/private/legacy/authorizations?token=${encodeURIComponent(username)}`; const res = await this.raw(path); // InfluxDB answers a filter matching nothing with 404, not an empty list. if (res.status === 404) return undefined; if (!res.ok) throw new Error(`InfluxDB API ${path}: ${res.status} ${await res.text()}`); const body = (await res.json()) as { authorizations?: LegacyAuthorization[] }; return body.authorizations?.find((a) => a.token === username); } async createLegacy(a: Omit): Promise { return (await (await this.send("/private/legacy/authorizations", "POST", a)).json()) as LegacyAuthorization; } /** Set a v1 authorization's password. InfluxDB keeps only a hash of it, so it can be set, never read. */ async setLegacyPassword(id: string, password: string): Promise { await this.send(`/private/legacy/authorizations/${encodeURIComponent(id)}/password`, "POST", { password }); } async updateLegacy(id: string, patch: { status?: "active" | "inactive"; description?: string }): Promise { await this.send(`/private/legacy/authorizations/${encodeURIComponent(id)}`, "PATCH", patch); } async deleteLegacy(id: string): Promise { await this.send(`/private/legacy/authorizations/${encodeURIComponent(id)}`, "DELETE"); } /** * Whether this username and password sign in on the v1 API — the consumer's own view. Asked with * a statement that reads nothing (`SHOW DATABASES` lists only what the credential may read), sent * with Basic auth so the password is never in a URL. 401 is a wrong password or no such user; * anything else that is not a server error means InfluxDB knew who was asking. */ async legacySignsIn(username: string, password: string): Promise { const res = await fetch(`${this.baseUrl}/query?q=${encodeURIComponent("SHOW DATABASES")}`, { headers: { Authorization: `Basic ${Buffer.from(`${username}:${password}`).toString("base64")}` }, }); await res.arrayBuffer(); if (res.status === 401) return false; if (res.status >= 500) throw new Error(`InfluxDB v1 /query: ${res.status}`); return true; } /** Server health — the one endpoint that needs no token, but we send it anyway. */ async health(): Promise { return (await (await this.request("/health")).json()) as InfluxHealth; } async listBuckets(): Promise { const body = (await (await this.request("/api/v2/buckets")).json()) as { buckets?: unknown[] }; return (body.buckets ?? []).map((b) => { const bucket = b as Record; const rules = (bucket.retentionRules as { everySeconds?: number }[] | undefined) ?? []; return { id: String(bucket.id), name: String(bucket.name), orgID: bucket.orgID ? String(bucket.orgID) : undefined, retentionSeconds: rules[0]?.everySeconds, }; }); } /** * Run a read-only Flux query and return the raw CSV InfluxDB answers with, plus a light parse * into rows. Read-only: Flux has no write verb, and the token's own permissions bound the rest — * this client never calls the write endpoint. */ async query(flux: string): Promise<{ csv: string; rows: Record[] }> { const res = await this.request(`/api/v2/query?org=${encodeURIComponent(this.org)}`, { method: "POST", headers: { "Content-Type": "application/vnd.flux", Accept: "application/csv", }, body: flux, }); const csv = await res.text(); return { csv, rows: parseAnnotatedCsv(csv) }; } } /** Parse InfluxDB's annotated CSV into rows keyed by column header. Annotation lines (starting * with #) and blanks are skipped; the first non-annotation line is the header. */ function parseAnnotatedCsv(csv: string): Record[] { const lines = csv.split("\n").filter((l) => l.trim() && !l.startsWith("#")); if (lines.length < 2) return []; const header = lines[0].split(","); return lines.slice(1).map((line) => { const cells = line.split(","); const row: Record = {}; header.forEach((h, i) => { if (h) row[h] = cells[i] ?? ""; }); return row; }); }