home-assistant, influxdb: full nox modules (ADR 0044/0046)
home-assistant: states/call-service/config tools, emits state.changed bounded to actuator/contact domains (not attribute ticks), overridable via a watch allowlist. influxdb: health/buckets/flux-query tools — tools-only, since a time-series DB has no lifecycle event to emit here. Typecheck; manifests parse.
This commit is contained in:
@@ -0,0 +1,106 @@
|
||||
// The InfluxDB API client — influxdb's own code, living in the module (novox/hq ADR 0044). Only
|
||||
// this module's tools import it. Talks to the InfluxDB 2.x HTTP API (/api/v2) with a token.
|
||||
|
||||
export interface InfluxHealth {
|
||||
name?: string;
|
||||
status?: string;
|
||||
message?: string;
|
||||
version?: string;
|
||||
}
|
||||
|
||||
export interface InfluxBucket {
|
||||
id: string;
|
||||
name: string;
|
||||
orgID?: string;
|
||||
retentionSeconds?: number;
|
||||
}
|
||||
|
||||
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 url = env.MESH_INFLUXDB_URL ?? `http://127.0.0.1:${env.INFLUXDB_PORT ?? "8086"}`;
|
||||
const token = env.MESH_INFLUXDB_TOKEN;
|
||||
if (!token) throw new Error("no InfluxDB token — set MESH_INFLUXDB_TOKEN");
|
||||
const 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;
|
||||
});
|
||||
}
|
||||
@@ -0,0 +1,14 @@
|
||||
{
|
||||
"name": "@novox/module-influxdb",
|
||||
"version": "0.1.0",
|
||||
"description": "influxdb — time-series database. Its API client and tools 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,46 @@
|
||||
// influxdb's tools — its own code (novox/hq ADR 0044), importing its own client. They return
|
||||
// structured data; the mesh serves them through the sdk's tool harness. Read-only: health, bucket
|
||||
// listing, and Flux queries — no write path is exposed.
|
||||
|
||||
import { registerModuleTools, type ToolDefinition } from "@novox/mesh-sdk/tools";
|
||||
import { InfluxDBClient } from "../client.js";
|
||||
|
||||
export function getInfluxDBTools(influx: InfluxDBClient): ToolDefinition[] {
|
||||
return [
|
||||
{
|
||||
name: "influxdb_health",
|
||||
description: "InfluxDB server health and version.",
|
||||
input: {},
|
||||
run: async () => influx.health(),
|
||||
},
|
||||
{
|
||||
name: "influxdb_list_buckets",
|
||||
description: "List InfluxDB buckets in the org, with their retention.",
|
||||
input: {},
|
||||
run: async () => {
|
||||
const buckets = await influx.listBuckets();
|
||||
return { count: buckets.length, buckets };
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "influxdb_query",
|
||||
description:
|
||||
"Run a read-only Flux query against InfluxDB and return the parsed rows (plus raw CSV). The query is Flux, e.g. from(bucket:\"default\") |> range(start:-1h).",
|
||||
input: { flux: { type: "string", description: "the Flux query to run" } },
|
||||
run: async (args) => {
|
||||
const { csv, rows } = await influx.query(String(args.flux));
|
||||
return { rowCount: rows.length, rows, csv };
|
||||
},
|
||||
},
|
||||
];
|
||||
}
|
||||
|
||||
// The tools exist only when a token is configured; without one, influxdb contributes none rather
|
||||
// than failing the whole runtime.
|
||||
registerModuleTools("influxdb", (env) => {
|
||||
try {
|
||||
return getInfluxDBTools(InfluxDBClient.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", "tools/index.ts"]
|
||||
}
|
||||
Reference in New Issue
Block a user