grafana's data source and Node-RED's influxdb nodes reached ace's InfluxDB by a LAN IP or a public name nobody routes, with a credential somebody made by hand. Now a consumer requires influxdb-api and is told where it is, which org and default bucket it serves, and signs in with the password the mesh minted for the pair. The credential is a v1-compatibility authorization, made per grant by the new provisioner: InfluxDB 2.x generates API tokens itself and ignores one the caller sends, so a v2 token could only be accepted by hand per pair; a v1 authorization takes a caller-chosen password (8-72 characters, the mesh mints 40) and reads/writes every bucket as a database of its name over InfluxQL and line protocol. A consumer contributes `access` (read, write, read-write) and, for writing, the buckets; a missing bucket is made and never deleted. Only authorizations named mesh_* and marked [mesh] are ever changed or removed; anything else of that name is refused and left alone. The org and default bucket are served facts the assignment's settings set, reaching both the consumers and the provisioner's config.json.
232 lines
9.9 KiB
TypeScript
232 lines
9.9 KiB
TypeScript
// 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<string, string> {
|
|
if (!file) return {};
|
|
try { return JSON.parse(readFileSync(file, "utf8")) as Record<string, string>; }
|
|
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<Response> {
|
|
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<Response> {
|
|
return fetch(`${this.baseUrl}${path}`, {
|
|
...init,
|
|
headers: { Authorization: `Token ${this.token}`, ...(init?.headers ?? {}) },
|
|
});
|
|
}
|
|
|
|
private async send(path: string, method: string, body?: unknown): Promise<Response> {
|
|
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<string | undefined> {
|
|
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<InfluxBucket | undefined> {
|
|
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<InfluxBucket> {
|
|
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<LegacyAuthorization | undefined> {
|
|
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<LegacyAuthorization, "id">): Promise<LegacyAuthorization> {
|
|
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<void> {
|
|
await this.send(`/private/legacy/authorizations/${encodeURIComponent(id)}/password`, "POST", { password });
|
|
}
|
|
|
|
async updateLegacy(id: string, patch: { status?: "active" | "inactive"; description?: string }): Promise<void> {
|
|
await this.send(`/private/legacy/authorizations/${encodeURIComponent(id)}`, "PATCH", patch);
|
|
}
|
|
|
|
async deleteLegacy(id: string): Promise<void> {
|
|
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<boolean> {
|
|
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<InfluxHealth> {
|
|
return (await (await this.request("/health")).json()) as InfluxHealth;
|
|
}
|
|
|
|
async listBuckets(): Promise<InfluxBucket[]> {
|
|
const body = (await (await this.request("/api/v2/buckets")).json()) as { buckets?: unknown[] };
|
|
return (body.buckets ?? []).map((b) => {
|
|
const bucket = b as Record<string, unknown>;
|
|
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<string, string>[] }> {
|
|
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<string, string>[] {
|
|
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<string, string> = {};
|
|
header.forEach((h, i) => {
|
|
if (h) row[h] = cells[i] ?? "";
|
|
});
|
|
return row;
|
|
});
|
|
}
|