117 lines
4.1 KiB
TypeScript
117 lines
4.1 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;
|
|
}
|
|
|
|
/** 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 {}; }
|
|
}
|
|
|
|
export class InfluxDBClient {
|
|
readonly baseUrl: string;
|
|
|
|
constructor(
|
|
url: string,
|
|
private readonly token: string,
|
|
private 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"}`;
|
|
const token = cfg.token ?? env.MESH_INFLUXDB_TOKEN;
|
|
if (!token) throw new Error("no InfluxDB token — set MESH_INFLUXDB_TOKEN");
|
|
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;
|
|
}
|
|
|
|
/** 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;
|
|
});
|
|
}
|