home-assistant: its broker and its Sonarr, Radarr and Lidarr come from the mesh
Home Assistant reached mosquitto and the three Servarr apps at 127.0.0.1 and a port typed into its own storage. It now requires mqtt-topic (asking for every topic: discovery and the devices' topics are its job), sonarr-api, radarr-api and lidarr-api, and a run-once `provisions` step — declared last, restarted when a binding or pair credential changes — makes Home Assistant's config entries say what the mesh bound, through Home Assistant's own config flows and never its .storage: - MQTT: the broker is asked first whether it takes the delivered login; then the integration's reconfigure flow sets broker, port, username and password, every other setting sent back as Home Assistant pre-filled it, and Home Assistant's own connection test must pass. A digest of what was written makes a rerun "already as the mesh says". Refused anywhere, nothing is written and Home Assistant keeps the login it has. - Sonarr/Radarr/Lidarr: the bound key is tried against the app (a minted key is never written; the failure names the `secret accept`); no entry is made through the user flow; a reauth Home Assistant started is finished with the bound key (and URL where the integration asks); an entry already at the bound URL, or at another URL reaching the same running app, is left as it is. These integrations have no reconfigure flow, so a working entry elsewhere is refused loudly — the step never removes an entry. The sidecar's URL now uses the port it was given.
This commit is contained in:
@@ -0,0 +1,479 @@
|
||||
// How home-assistant's provisions step brings Home Assistant's integrations in line with what the
|
||||
// mesh bound: the MQTT integration to `mqtt-topic`, the Sonarr, Radarr and Lidarr integrations to
|
||||
// `sonarr-api`, `radarr-api` and `lidarr-api`. Pure logic over two seams — Home Assistant's config
|
||||
// flows (hass.ts) and the broker/apps — so it is tested against fakes (test/provisions.test.ts).
|
||||
//
|
||||
// The half that reads files and talks HTTP lives beside it (mesh.ts, hass.ts, probe.ts, index.ts).
|
||||
|
||||
import { createHash } from "node:crypto";
|
||||
|
||||
import type { Hass, SchemaField } from "./hass.js";
|
||||
import type { Probe } from "./probe.js";
|
||||
|
||||
/** What the mesh wrote at `binds.<provision>` (the controller's binding document). */
|
||||
export interface Binding {
|
||||
provision?: string;
|
||||
from?: string;
|
||||
at?: string;
|
||||
as?: string;
|
||||
serves?: Record<string, unknown>;
|
||||
}
|
||||
|
||||
/** How one provision came out. Never carries a credential. */
|
||||
export type Outcome =
|
||||
| { what: string; result: "unchanged"; note?: string }
|
||||
| { what: string; result: "written"; fields: string[]; note?: string }
|
||||
| { what: string; result: "equivalent"; note: string }
|
||||
| { what: string; result: "refused"; problem: string };
|
||||
|
||||
/** A port the binding serves, or undefined when it names none usable. */
|
||||
export function portOf(serves: Record<string, unknown> | undefined): number | undefined {
|
||||
const port = Number(serves?.port);
|
||||
return Number.isInteger(port) && port > 0 && port <= 65535 ? port : undefined;
|
||||
}
|
||||
|
||||
/** A host as it goes into a URL: an IPv6 literal bracketed. */
|
||||
export function urlHost(host: string): string {
|
||||
return host.includes(":") && !host.startsWith("[") ? `[${host}]` : host;
|
||||
}
|
||||
|
||||
/**
|
||||
* What this step last wrote, per target, as a digest: the only way to know "already as the mesh
|
||||
* says" for a credential Home Assistant will not show back. A sha256 over the target and the values,
|
||||
* never the values; kept in the module's own placed directory.
|
||||
*/
|
||||
export interface Marks {
|
||||
get(name: string): Promise<string | undefined>;
|
||||
set(name: string, digest: string): Promise<void>;
|
||||
}
|
||||
|
||||
export function digest(...parts: (string | number)[]): string {
|
||||
return createHash("sha256").update(parts.map(String).join("\u0000")).digest("hex");
|
||||
}
|
||||
|
||||
/** An error as text with the credential taken out, raw and URL-encoded. */
|
||||
export function scrub(err: unknown, ...secrets: (string | undefined)[]): string {
|
||||
let text = err instanceof Error ? err.message : String(err);
|
||||
for (const s of secrets) {
|
||||
if (!s) continue;
|
||||
for (const form of new Set([s, encodeURIComponent(s)])) text = text.split(form).join("***");
|
||||
}
|
||||
return text;
|
||||
}
|
||||
|
||||
/**
|
||||
* What a form would submit if a person pressed "submit" without touching it: each field's
|
||||
* suggested value (what Home Assistant pre-fills from the entry), else its default; a section's
|
||||
* fields nested under its name. The step lays only the connection fields over this, so every other
|
||||
* choice the entry carries is sent back exactly as Home Assistant showed it.
|
||||
*/
|
||||
export function formValues(schema: readonly SchemaField[] | null | undefined): Record<string, unknown> {
|
||||
const out: Record<string, unknown> = {};
|
||||
for (const field of schema ?? []) {
|
||||
if (Array.isArray(field.schema)) {
|
||||
out[field.name] = formValues(field.schema);
|
||||
continue;
|
||||
}
|
||||
const suggested = field.description?.suggested_value;
|
||||
if (suggested !== undefined && suggested !== null) out[field.name] = suggested;
|
||||
else if (field.default !== undefined) out[field.name] = field.default;
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
/** Whether a form has a field of this name at its top level. */
|
||||
export function hasField(schema: readonly SchemaField[] | null | undefined, name: string): boolean {
|
||||
return (schema ?? []).some((f) => f.name === name);
|
||||
}
|
||||
|
||||
// ---- MQTT ----
|
||||
|
||||
// Home Assistant's MQTT integration, pointed at the broker the mesh bound — `mqtt-topic`.
|
||||
//
|
||||
// **Why a step.** Home Assistant keeps its broker, login and password in its MQTT config entry
|
||||
// (`.storage/core.config_entries`), not in a file the mesh could fill with `${bound:mqtt-topic:at}`.
|
||||
// So this reads the binding and the pair credential and makes the entry say the same thing, through
|
||||
// the MQTT integration's own reconfigure flow — the flow its "Reconfigure" button runs, which tests
|
||||
// the connection itself and saves nothing it could not connect with.
|
||||
//
|
||||
// **Only the connection, and only when it differs.** Broker, port, username, password. The protocol
|
||||
// version, client id, keepalive, TLS choices and discovery options the entry holds are sent back
|
||||
// exactly as Home Assistant pre-filled them. Whether the password already matches cannot be read
|
||||
// back (Home Assistant never shows a stored password), so the step keeps a digest of what it last
|
||||
// wrote: equal broker/port/username and an equal digest is "already as the mesh says".
|
||||
//
|
||||
// **Nothing loses its connection without someone seeing it.** Before Home Assistant is touched the
|
||||
// broker itself is asked whether it takes the delivered login (the provisioner creates it within
|
||||
// seconds of the grant): if not, nothing is written and the step fails saying why, and Home
|
||||
// Assistant keeps the login it has — the carried `luffy` on ace, which mosquitto keeps. If Home
|
||||
// Assistant's own connection test refuses the new settings, the flow saves nothing, and the step
|
||||
// fails with Home Assistant's reason. A login that may not subscribe to the discovery topics is said
|
||||
// as a warning: discovery would find nothing.
|
||||
|
||||
export const MQTT_PROVISION = "mqtt-topic";
|
||||
/** Home Assistant's discovery prefix, subscribed to whenever discovery is on (the default). */
|
||||
export const DISCOVERY_FILTER = "homeassistant/#";
|
||||
|
||||
export interface MqttWanted {
|
||||
host: string;
|
||||
port: number;
|
||||
username: string;
|
||||
password: string;
|
||||
}
|
||||
|
||||
export type Wanted = { ok: true; want: MqttWanted } | { ok: false; problem: string };
|
||||
|
||||
/** The broker, port and login the mesh says Home Assistant uses. */
|
||||
export function wantedMqtt(binding: Binding | undefined, credential: string | undefined): Wanted {
|
||||
if (!binding) return { ok: false, problem: `no binding for ${MQTT_PROVISION} was delivered — the mesh writes it before this step runs` };
|
||||
const host = typeof binding.at === "string" ? binding.at.trim() : "";
|
||||
if (!host) return { ok: false, problem: `the ${MQTT_PROVISION} binding names no host (at)` };
|
||||
const port = portOf(binding.serves);
|
||||
if (port === undefined) return { ok: false, problem: `the ${MQTT_PROVISION} binding serves no usable port (${String(binding.serves?.port)})` };
|
||||
const scheme = binding.serves?.scheme;
|
||||
if (scheme !== undefined && scheme !== "mqtt") {
|
||||
return { ok: false, problem: `the ${MQTT_PROVISION} binding serves scheme ${String(scheme)}; this step writes plain MQTT` };
|
||||
}
|
||||
const username = typeof binding.as === "string" ? binding.as.trim() : "";
|
||||
if (!username) return { ok: false, problem: `the ${MQTT_PROVISION} binding names no login (as)` };
|
||||
const password = (credential ?? "").replace(/\n$/, "");
|
||||
if (!password) return { ok: false, problem: `the ${MQTT_PROVISION} credential is empty or was not delivered` };
|
||||
return { ok: true, want: { host, port, username, password } };
|
||||
}
|
||||
|
||||
export interface MqttDeps {
|
||||
hass: Hass;
|
||||
probe: Probe;
|
||||
marks: Marks;
|
||||
}
|
||||
|
||||
const markFor = (entryId: string, w: MqttWanted): string => digest("mqtt", entryId, w.host, w.port, w.username, w.password);
|
||||
|
||||
/** Bring Home Assistant's MQTT entry in line with the mesh. Never throws: every failure is an outcome. */
|
||||
export async function reconcileMqtt(deps: MqttDeps, binding: Binding | undefined, credential: string | undefined): Promise<Outcome> {
|
||||
const what = "mqtt";
|
||||
const w = wantedMqtt(binding, credential);
|
||||
if ("problem" in w) return { what, result: "refused", problem: w.problem };
|
||||
const want = w.want;
|
||||
|
||||
// The broker first: a login it does not take is never written into Home Assistant.
|
||||
let note: string | undefined;
|
||||
try {
|
||||
const probe = await deps.probe(want.host, want.port, want.username, want.password, DISCOVERY_FILTER);
|
||||
if (probe.connack === 4 || probe.connack === 5) {
|
||||
return {
|
||||
what,
|
||||
result: "refused",
|
||||
problem:
|
||||
`the broker at ${want.host}:${want.port} does not (yet) take the login ${want.username} with the delivered ` +
|
||||
`password (CONNACK ${probe.connack}); mosquitto's provisioner creates it from the grant — nothing was ` +
|
||||
`written, and Home Assistant keeps the broker login it has`,
|
||||
};
|
||||
}
|
||||
if (probe.connack !== 0) {
|
||||
return { what, result: "refused", problem: `the broker at ${want.host}:${want.port} answered CONNACK ${probe.connack}; nothing was written` };
|
||||
}
|
||||
if (probe.suback === 0x80) {
|
||||
note =
|
||||
`warning: ${want.username} may not subscribe to ${DISCOVERY_FILTER} — MQTT discovery will find nothing; ` +
|
||||
`grant it with the mqtt-topic contribution's \`topics\``;
|
||||
}
|
||||
} catch (err) {
|
||||
return {
|
||||
what,
|
||||
result: "refused",
|
||||
problem: `the broker at ${want.host}:${want.port} could not be asked: ${scrub(err, want.password)}; nothing was written`,
|
||||
};
|
||||
}
|
||||
|
||||
try {
|
||||
const entries = (await deps.hass.entries("mqtt")).filter((e) => e.domain === "mqtt");
|
||||
if (entries.length > 1) {
|
||||
return { what, result: "refused", problem: `Home Assistant has ${entries.length} MQTT entries; which one the mesh owns is not guessed` };
|
||||
}
|
||||
if (entries.length === 0) return await createEntry(deps, want, note);
|
||||
|
||||
const entry = entries[0];
|
||||
const flow = await deps.hass.startFlow("mqtt", entry.entry_id);
|
||||
if (flow.type !== "form" || !flow.flow_id || !flow.data_schema) {
|
||||
if (flow.flow_id) await deps.hass.abortFlow(flow.flow_id);
|
||||
return { what, result: "refused", problem: `Home Assistant's MQTT reconfigure flow answered ${flow.type}${flow.reason ? ` (${flow.reason})` : ""}` };
|
||||
}
|
||||
const current = formValues(flow.data_schema);
|
||||
const fields: string[] = [];
|
||||
if (String(current.broker ?? "") !== want.host) fields.push("broker");
|
||||
if (Number(current.port ?? 0) !== want.port) fields.push("port");
|
||||
if (String(current.username ?? "") !== want.username) fields.push("username");
|
||||
if ((await deps.marks.get("mqtt")) !== markFor(entry.entry_id, want)) fields.push("password");
|
||||
if (fields.length === 0) {
|
||||
await deps.hass.abortFlow(flow.flow_id);
|
||||
return note ? { what, result: "unchanged", note } : { what, result: "unchanged" };
|
||||
}
|
||||
|
||||
const saved = await deps.hass.stepFlow(flow.flow_id, {
|
||||
...current,
|
||||
broker: want.host,
|
||||
port: want.port,
|
||||
username: want.username,
|
||||
password: want.password,
|
||||
});
|
||||
if (saved.type === "abort" && saved.reason === "reconfigure_successful") {
|
||||
await deps.marks.set("mqtt", markFor(entry.entry_id, want));
|
||||
return { what, result: "written", fields, ...(note ? { note } : {}) };
|
||||
}
|
||||
if (saved.flow_id) await deps.hass.abortFlow(saved.flow_id);
|
||||
return {
|
||||
what,
|
||||
result: "refused",
|
||||
problem:
|
||||
`Home Assistant's own connection test refused ${want.username}@${want.host}:${want.port} ` +
|
||||
`(${describe(saved)}); its MQTT entry is unchanged`,
|
||||
};
|
||||
} catch (err) {
|
||||
return { what, result: "refused", problem: scrub(err, want.password) };
|
||||
}
|
||||
}
|
||||
|
||||
/** A fresh Home Assistant has no MQTT entry: made through the integration's user flow. */
|
||||
async function createEntry(deps: MqttDeps, want: MqttWanted, note?: string): Promise<Outcome> {
|
||||
const what = "mqtt";
|
||||
let flow = await deps.hass.startFlow("mqtt");
|
||||
if (flow.type === "form" && flow.step_id !== "broker" && flow.flow_id) {
|
||||
// Anything before the broker form (none outside the Supervisor) is not this step's to answer.
|
||||
await deps.hass.abortFlow(flow.flow_id);
|
||||
return { what, result: "refused", problem: `Home Assistant's MQTT user flow asked ${flow.step_id} before the broker` };
|
||||
}
|
||||
if (flow.type !== "form" || !flow.flow_id) {
|
||||
return { what, result: "refused", problem: `Home Assistant's MQTT user flow answered ${describe(flow)}` };
|
||||
}
|
||||
flow = await deps.hass.stepFlow(flow.flow_id, {
|
||||
...formValues(flow.data_schema),
|
||||
broker: want.host,
|
||||
port: want.port,
|
||||
username: want.username,
|
||||
password: want.password,
|
||||
});
|
||||
if (flow.type === "create_entry") {
|
||||
const id = (flow.result as { entry_id?: string } | undefined)?.entry_id;
|
||||
if (id) await deps.marks.set("mqtt", markFor(id, want));
|
||||
return { what, result: "written", fields: ["entry"], ...(note ? { note } : {}) };
|
||||
}
|
||||
if (flow.flow_id) await deps.hass.abortFlow(flow.flow_id);
|
||||
return { what, result: "refused", problem: `Home Assistant refused a new MQTT entry for ${want.host}:${want.port} (${describe(flow)})` };
|
||||
}
|
||||
|
||||
export function describe(r: { type: string; reason?: string; errors?: Record<string, string> | null }): string {
|
||||
const errors = r.errors ? Object.entries(r.errors).map(([k, v]) => `${k}: ${v}`).join(", ") : "";
|
||||
return [r.type, r.reason, errors].filter(Boolean).join(" — ");
|
||||
}
|
||||
|
||||
// ---- Sonarr, Radarr, Lidarr ----
|
||||
|
||||
// Home Assistant's Sonarr, Radarr and Lidarr integrations, pointed at the apps the mesh bound —
|
||||
// `sonarr-api`, `radarr-api`, `lidarr-api` (their providers: mesh-catalog #156).
|
||||
//
|
||||
// **What Home Assistant lets anyone change, and what it does not.** Each integration keeps a URL and
|
||||
// an API key in its config entry. None of the three has a reconfigure flow: Home Assistant changes
|
||||
// them only through the flow its UI runs —
|
||||
// - a **user flow** makes a new entry (validated against the app);
|
||||
// - a **reauth flow**, which Home Assistant starts by itself when the app refuses the key it holds,
|
||||
// takes a new key (Sonarr) or a new URL and key (Radarr, Lidarr);
|
||||
// - anything else — the URL of a working entry — only by removing the integration and adding it
|
||||
// again, which throws away its entities' names, areas and history links. **This step never
|
||||
// removes an entry.**
|
||||
// So, per app:
|
||||
// 1. The bound key is tried against the bound app first. Refused, nothing is written: until the
|
||||
// operator accepts the app's own key for this pair, the mesh delivers a value it minted, which
|
||||
// no Servarr app takes (novox/hq ADR 0092) — the failure names the `secret accept` that fixes it.
|
||||
// 2. No entry: one is made through the user flow.
|
||||
// 3. A reauth flow Home Assistant started for the entry: finished with the bound key (and URL,
|
||||
// where the integration's reauth asks for one).
|
||||
// 4. An entry whose URL (read from the device the integration registered, `configuration_url`)
|
||||
// is the bound one and which is loaded: already as the mesh says. The key needs no digest here:
|
||||
// a Servarr app has one key, so an entry loaded against the app holds the key the app took.
|
||||
// 5. A working entry at a different URL that reaches **the same app** — the same process, by the
|
||||
// app's own status (start time, data folder, version) — is left as it is and said: ace's entries
|
||||
// say `127.0.0.1:<port>` and the binding says `ace.internal:<port>`, one Sonarr either way.
|
||||
// 6. Anything else is refused, loudly, with what the operator can do; nothing is removed.
|
||||
|
||||
|
||||
export interface ServarrApp {
|
||||
/** The integration's domain, also the app. */
|
||||
domain: "sonarr" | "radarr" | "lidarr";
|
||||
/** The provision it is required as: the `requires`, `binds` and `secrets` key. */
|
||||
provision: string;
|
||||
/** The app's status endpoint: answers 401 to a wrong key, and says which process answered. */
|
||||
statusPath: string;
|
||||
}
|
||||
|
||||
export const APPS: readonly ServarrApp[] = [
|
||||
{ domain: "sonarr", provision: "sonarr-api", statusPath: "/api/v3/system/status" },
|
||||
{ domain: "radarr", provision: "radarr-api", statusPath: "/api/v3/system/status" },
|
||||
{ domain: "lidarr", provision: "lidarr-api", statusPath: "/api/v1/system/status" },
|
||||
];
|
||||
|
||||
/** The HTTP the step needs toward the apps, so a test can stand fakes in. */
|
||||
export interface Http {
|
||||
fetch(url: string, init?: { method?: string; headers?: Record<string, string> }): Promise<{ status: number; text(): Promise<string> }>;
|
||||
}
|
||||
|
||||
export type AppWanted = { ok: true; url: string; key: string; from: string } | { ok: false; problem: string };
|
||||
|
||||
/** The URL and key the mesh says Home Assistant uses for this app. */
|
||||
export function wantedApp(spec: ServarrApp, binding: Binding | undefined, credential: string | undefined): AppWanted {
|
||||
if (!binding) return { ok: false, problem: `no binding for ${spec.provision} was delivered — the mesh writes it before this step runs` };
|
||||
const at = typeof binding.at === "string" ? binding.at.trim() : "";
|
||||
if (!at) return { ok: false, problem: `the ${spec.provision} binding names no host (at)` };
|
||||
const port = portOf(binding.serves);
|
||||
if (port === undefined) return { ok: false, problem: `the ${spec.provision} binding serves no usable port (${String(binding.serves?.port)})` };
|
||||
const scheme = typeof binding.serves?.scheme === "string" && binding.serves.scheme ? binding.serves.scheme : "http";
|
||||
if (scheme !== "http" && scheme !== "https") return { ok: false, problem: `the ${spec.provision} binding serves scheme ${scheme}` };
|
||||
const base = typeof binding.serves?.["url-base"] === "string" ? String(binding.serves["url-base"]).trim().replace(/^\/+|\/+$/g, "") : "";
|
||||
const key = (credential ?? "").trim();
|
||||
if (!key) return { ok: false, problem: `the ${spec.provision} credential is empty or was not delivered` };
|
||||
return {
|
||||
ok: true,
|
||||
url: `${scheme}://${urlHost(at)}:${port}${base ? `/${base}` : ""}`,
|
||||
key,
|
||||
from: typeof binding.from === "string" ? binding.from : "",
|
||||
};
|
||||
}
|
||||
|
||||
/** Two URLs naming the same place: scheme, host, port (explicit or default) and base path. */
|
||||
export function sameUrl(a: string | null | undefined, b: string): boolean {
|
||||
if (!a) return false;
|
||||
try {
|
||||
const x = new URL(a);
|
||||
const y = new URL(b);
|
||||
const port = (u: URL) => u.port || (u.protocol === "https:" ? "443" : "80");
|
||||
const path = (u: URL) => u.pathname.replace(/\/+$/, "");
|
||||
return x.protocol === y.protocol && x.hostname.toLowerCase() === y.hostname.toLowerCase() && port(x) === port(y) && path(x) === path(y);
|
||||
} catch {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
type Status = { taken: true; status: Record<string, unknown> } | { taken: false };
|
||||
|
||||
/** The app's status with this key: `taken: false` when it refuses the key; throws when it cannot be asked. */
|
||||
export async function appStatus(http: Http, spec: ServarrApp, url: string, key: string): Promise<Status> {
|
||||
const res = await http.fetch(`${url.replace(/\/+$/, "")}${spec.statusPath}`, {
|
||||
method: "GET",
|
||||
headers: { "X-Api-Key": key, Accept: "application/json" },
|
||||
});
|
||||
if (res.status === 401 || res.status === 403) return { taken: false };
|
||||
if (res.status < 200 || res.status >= 300) throw new Error(`${spec.domain} answered ${res.status} at ${spec.statusPath}`);
|
||||
return { taken: true, status: JSON.parse(await res.text()) as Record<string, unknown> };
|
||||
}
|
||||
|
||||
/** Whether two status answers came from one running app. */
|
||||
export function sameInstance(a: Record<string, unknown>, b: Record<string, unknown>): boolean {
|
||||
const facts = ["startTime", "appData", "version"];
|
||||
return facts.every((k) => a[k] !== undefined && a[k] !== null && a[k] === b[k]);
|
||||
}
|
||||
|
||||
/** The remedy for a refused key, in the controller's words (ADR 0092). */
|
||||
export function acceptRemedy(spec: ServarrApp, from: string): string {
|
||||
return (
|
||||
`${spec.domain} refuses the ${spec.provision} credential the mesh delivered, so nothing was written into ` +
|
||||
`Home Assistant. A Servarr app has one API key and the mesh cannot make it: accept ${spec.domain}'s own key ` +
|
||||
`for this pair — \`secret accept <this node> home-assistant ${spec.provision} --provider ${from || "<its node>"} ` +
|
||||
`--from <file holding ${spec.domain}'s ApiKey>\``
|
||||
);
|
||||
}
|
||||
|
||||
export interface ServarrDeps {
|
||||
hass: Hass;
|
||||
http: Http;
|
||||
}
|
||||
|
||||
/** The input a Servarr form takes: what it shows, with the URL (where asked) and the key laid over. */
|
||||
function servarrInput(schema: readonly SchemaField[] | null | undefined, url: string, key: string): Record<string, unknown> {
|
||||
const input = formValues(schema);
|
||||
if (hasField(schema, "url")) input.url = url;
|
||||
if (hasField(schema, "api_key")) input.api_key = key;
|
||||
return input;
|
||||
}
|
||||
|
||||
/** Bring Home Assistant's entry for one app in line with the mesh. Never throws. */
|
||||
export async function reconcileApp(deps: ServarrDeps, spec: ServarrApp, binding: Binding | undefined, credential: string | undefined): Promise<Outcome> {
|
||||
const what = spec.domain;
|
||||
const w = wantedApp(spec, binding, credential);
|
||||
if ("problem" in w) return { what, result: "refused", problem: w.problem };
|
||||
|
||||
let bound: Status;
|
||||
try {
|
||||
bound = await appStatus(deps.http, spec, w.url, w.key);
|
||||
} catch (err) {
|
||||
return { what, result: "refused", problem: `${spec.domain} could not be asked at ${w.url}: ${scrub(err, w.key)}` };
|
||||
}
|
||||
if (!bound.taken) return { what, result: "refused", problem: acceptRemedy(spec, w.from) };
|
||||
|
||||
try {
|
||||
const entries = (await deps.hass.entries(spec.domain)).filter((e) => e.domain === spec.domain);
|
||||
if (entries.length > 1) {
|
||||
return { what, result: "refused", problem: `Home Assistant has ${entries.length} ${spec.domain} entries; which one the mesh owns is not guessed` };
|
||||
}
|
||||
|
||||
// No entry: made, through the integration's own user flow, which validates the key itself.
|
||||
if (entries.length === 0) {
|
||||
const flow = await deps.hass.startFlow(spec.domain);
|
||||
if (flow.type !== "form" || !flow.flow_id) return { what, result: "refused", problem: `Home Assistant's ${spec.domain} user flow answered ${describe(flow)}` };
|
||||
const made = await deps.hass.stepFlow(flow.flow_id, servarrInput(flow.data_schema, w.url, w.key));
|
||||
if (made.type === "create_entry") return { what, result: "written", fields: ["entry"] };
|
||||
if (made.flow_id) await deps.hass.abortFlow(made.flow_id);
|
||||
return { what, result: "refused", problem: `Home Assistant refused a new ${spec.domain} entry at ${w.url} (${describe(made)})` };
|
||||
}
|
||||
|
||||
const entry = entries[0];
|
||||
if (entry.disabled_by) return { what, result: "unchanged", note: `the ${spec.domain} entry is disabled (by ${entry.disabled_by}); left alone` };
|
||||
|
||||
// A reauth Home Assistant started because the app refused its key: finished with the bound one.
|
||||
const reauth = (await deps.hass.flowsInProgress()).find(
|
||||
(f) => f.handler === spec.domain && f.context?.source === "reauth" && f.context?.entry_id === entry.entry_id,
|
||||
);
|
||||
if (reauth) {
|
||||
let step = await deps.hass.stepFlow(reauth.flow_id, {}); // reauth_confirm: a confirmation, no fields
|
||||
if (step.type === "form" && step.flow_id && step.step_id !== "reauth_confirm") {
|
||||
const input = servarrInput(step.data_schema, w.url, w.key);
|
||||
const fields = ["api_key", ...(hasField(step.data_schema, "url") ? ["url"] : [])];
|
||||
step = await deps.hass.stepFlow(step.flow_id, input);
|
||||
if (step.type === "abort" && step.reason === "reauth_successful") return { what, result: "written", fields };
|
||||
}
|
||||
return { what, result: "refused", problem: `Home Assistant's ${spec.domain} reauth did not take the bound key and URL (${describe(step)})` };
|
||||
}
|
||||
|
||||
const device = (await deps.hass.devices()).find((d) => d.config_entries?.includes(entry.entry_id) && d.configuration_url);
|
||||
const current = device?.configuration_url ?? undefined;
|
||||
if (entry.state === "loaded" && sameUrl(current, w.url)) return { what, result: "unchanged" };
|
||||
|
||||
if (entry.state === "loaded" && current) {
|
||||
let there: Status | undefined;
|
||||
try {
|
||||
there = await appStatus(deps.http, spec, current, w.key);
|
||||
} catch {
|
||||
there = undefined;
|
||||
}
|
||||
if (there?.taken && sameInstance(there.status, bound.status)) {
|
||||
return {
|
||||
what,
|
||||
result: "equivalent",
|
||||
note:
|
||||
`Home Assistant reaches ${spec.domain} at ${current}, the same running app the mesh bound at ${w.url}; ` +
|
||||
`Home Assistant has no way to change a working ${spec.domain} entry's URL short of removing it, so it is left as it is`,
|
||||
};
|
||||
}
|
||||
}
|
||||
return {
|
||||
what,
|
||||
result: "refused",
|
||||
problem:
|
||||
`Home Assistant's ${spec.domain} entry (${entry.state ?? "unknown state"}) points at ${current ?? "an unknown URL"}, ` +
|
||||
`not the ${spec.domain} the mesh bound at ${w.url}. Home Assistant only lets a working entry's URL change by ` +
|
||||
`removing and re-adding the integration, which this step never does: remove it in Home Assistant ` +
|
||||
`(Settings → Devices & services → ${spec.domain}) and the next run adds it at the bound URL`,
|
||||
};
|
||||
} catch (err) {
|
||||
return { what, result: "refused", problem: scrub(err, w.key) };
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,160 @@
|
||||
// Home Assistant's own configuration API, as the provisions step uses it — the supported way to
|
||||
// change an integration's connection. Home Assistant keeps every integration in
|
||||
// `.storage/core.config_entries`, a file it owns and rewrites; the mesh may not write it, and it is
|
||||
// not a file the mesh could merge into. What Home Assistant offers instead is the same thing its UI
|
||||
// uses: **config flows** over REST (`/api/config/config_entries/flow`) — a user flow creates an
|
||||
// entry, a reconfigure flow changes one, a reauth flow (which Home Assistant starts itself when a
|
||||
// credential stops working) replaces its credential — each validated by the integration's own
|
||||
// connection test before anything is saved. The two things REST does not answer (which flows Home
|
||||
// Assistant has started, which device an entry made) come over its WebSocket API.
|
||||
//
|
||||
// Nothing here reads `.storage`. Authenticated with the module's accepted long-lived access token.
|
||||
|
||||
/** A config entry as `GET /api/config/config_entries/entry` lists it — no data, no credentials. */
|
||||
export interface ConfigEntry {
|
||||
entry_id: string;
|
||||
domain: string;
|
||||
title?: string;
|
||||
source?: string;
|
||||
state?: string;
|
||||
disabled_by?: string | null;
|
||||
}
|
||||
|
||||
/** One field of a flow's form, as Home Assistant serializes a voluptuous schema. */
|
||||
export interface SchemaField {
|
||||
name: string;
|
||||
type?: string;
|
||||
required?: boolean;
|
||||
optional?: boolean;
|
||||
default?: unknown;
|
||||
description?: { suggested_value?: unknown } | null;
|
||||
/** A section (`type: "expandable"`) carries its own fields. */
|
||||
schema?: SchemaField[];
|
||||
}
|
||||
|
||||
/** What a flow answered: another form, an entry made, or the flow ended (abort). */
|
||||
export interface FlowResult {
|
||||
type: string;
|
||||
flow_id?: string;
|
||||
handler?: string;
|
||||
step_id?: string;
|
||||
data_schema?: SchemaField[] | null;
|
||||
errors?: Record<string, string> | null;
|
||||
reason?: string;
|
||||
result?: { entry_id?: string } | unknown;
|
||||
}
|
||||
|
||||
/** A flow in progress that Home Assistant started itself (a reauth, a discovery). */
|
||||
export interface FlowProgress {
|
||||
flow_id: string;
|
||||
handler: string;
|
||||
step_id?: string;
|
||||
context?: { source?: string; entry_id?: string };
|
||||
}
|
||||
|
||||
/** A device from the device registry; an integration names where its app is as configuration_url. */
|
||||
export interface DeviceEntry {
|
||||
id: string;
|
||||
config_entries?: string[];
|
||||
configuration_url?: string | null;
|
||||
}
|
||||
|
||||
export interface Hass {
|
||||
entries(domain: string): Promise<ConfigEntry[]>;
|
||||
/** A user flow for `handler`, or — given an entry — a reconfigure flow for it. */
|
||||
startFlow(handler: string, entryId?: string): Promise<FlowResult>;
|
||||
stepFlow(flowId: string, input: Record<string, unknown>): Promise<FlowResult>;
|
||||
abortFlow(flowId: string): Promise<void>;
|
||||
flowsInProgress(): Promise<FlowProgress[]>;
|
||||
devices(): Promise<DeviceEntry[]>;
|
||||
}
|
||||
|
||||
/** Home Assistant over HTTP: REST for entries and flows, one short WebSocket session per question. */
|
||||
export class HassApi implements Hass {
|
||||
private readonly base: string;
|
||||
|
||||
constructor(url: string, private readonly token: string) {
|
||||
this.base = url.replace(/\/$/, "");
|
||||
}
|
||||
|
||||
private async rest(method: string, path: string, body?: unknown): Promise<unknown> {
|
||||
const res = await fetch(`${this.base}${path}`, {
|
||||
method,
|
||||
headers: {
|
||||
Authorization: `Bearer ${this.token}`,
|
||||
Accept: "application/json",
|
||||
...(body !== undefined ? { "Content-Type": "application/json" } : {}),
|
||||
},
|
||||
body: body !== undefined ? JSON.stringify(body) : undefined,
|
||||
});
|
||||
const text = await res.text();
|
||||
if (!res.ok) {
|
||||
// Home Assistant's error text names fields, never echoes their values.
|
||||
throw new Error(`Home Assistant ${method} ${path} answered ${res.status}${text ? `: ${text.slice(0, 200)}` : ""}`);
|
||||
}
|
||||
return text ? (JSON.parse(text) as unknown) : undefined;
|
||||
}
|
||||
|
||||
async entries(domain: string): Promise<ConfigEntry[]> {
|
||||
return ((await this.rest("GET", `/api/config/config_entries/entry?domain=${encodeURIComponent(domain)}`)) ??
|
||||
[]) as ConfigEntry[];
|
||||
}
|
||||
|
||||
async startFlow(handler: string, entryId?: string): Promise<FlowResult> {
|
||||
return (await this.rest("POST", "/api/config/config_entries/flow", {
|
||||
handler,
|
||||
show_advanced_options: true,
|
||||
...(entryId ? { entry_id: entryId } : {}),
|
||||
})) as FlowResult;
|
||||
}
|
||||
|
||||
async stepFlow(flowId: string, input: Record<string, unknown>): Promise<FlowResult> {
|
||||
return (await this.rest("POST", `/api/config/config_entries/flow/${encodeURIComponent(flowId)}`, input)) as FlowResult;
|
||||
}
|
||||
|
||||
async abortFlow(flowId: string): Promise<void> {
|
||||
await this.rest("DELETE", `/api/config/config_entries/flow/${encodeURIComponent(flowId)}`).catch(() => undefined);
|
||||
}
|
||||
|
||||
async flowsInProgress(): Promise<FlowProgress[]> {
|
||||
return (await this.ws("config_entries/flow/progress")) as FlowProgress[];
|
||||
}
|
||||
|
||||
async devices(): Promise<DeviceEntry[]> {
|
||||
return (await this.ws("config/device_registry/list")) as DeviceEntry[];
|
||||
}
|
||||
|
||||
/** One WebSocket command: connect, authenticate, ask, close. */
|
||||
private ws(type: string): Promise<unknown> {
|
||||
const url = `${this.base.replace(/^http/, "ws")}/api/websocket`;
|
||||
return new Promise((resolve, reject) => {
|
||||
const socket = new WebSocket(url);
|
||||
const timer = setTimeout(() => {
|
||||
socket.close();
|
||||
reject(new Error(`Home Assistant's WebSocket did not answer ${type} within 30s`));
|
||||
}, 30_000);
|
||||
const done = (fn: () => void): void => {
|
||||
clearTimeout(timer);
|
||||
socket.close();
|
||||
fn();
|
||||
};
|
||||
socket.onerror = () => done(() => reject(new Error(`Home Assistant's WebSocket at ${url} failed`)));
|
||||
socket.onmessage = (event: { data: unknown }) => {
|
||||
const msg = JSON.parse(String(event.data)) as {
|
||||
type: string;
|
||||
id?: number;
|
||||
success?: boolean;
|
||||
result?: unknown;
|
||||
error?: { message?: string };
|
||||
};
|
||||
if (msg.type === "auth_required") socket.send(JSON.stringify({ type: "auth", access_token: this.token }));
|
||||
else if (msg.type === "auth_invalid") done(() => reject(new Error("Home Assistant refused the token")));
|
||||
else if (msg.type === "auth_ok") socket.send(JSON.stringify({ id: 1, type }));
|
||||
else if (msg.type === "result" && msg.id === 1) {
|
||||
if (msg.success) done(() => resolve(msg.result));
|
||||
else done(() => reject(new Error(`Home Assistant ${type}: ${msg.error?.message ?? "failed"}`)));
|
||||
}
|
||||
};
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,84 @@
|
||||
// home-assistant's provisions step — run once by the host after Home Assistant starts, and again
|
||||
// whenever a binding or pair credential it reads changes (the container's `restart-on`, novox/hq
|
||||
// ADR 0099). It points Home Assistant's MQTT integration at the `mqtt-topic` broker and its Sonarr,
|
||||
// Radarr and Lidarr integrations at the `sonarr-api`, `radarr-api` and `lidarr-api` apps, through
|
||||
// Home Assistant's own config flows (connections.ts). It connects to no mesh broker.
|
||||
//
|
||||
// Exits non-zero when anything could not be put right, so the node reports the step failed and the
|
||||
// host runs it again on the next apply. Declared last in the manifest, so its failing gates nothing
|
||||
// else of home-assistant's (novox/hq ADR 0136). Never prints a key or password.
|
||||
|
||||
import { join } from "node:path";
|
||||
|
||||
import { APPS, MQTT_PROVISION, reconcileApp, reconcileMqtt, type Outcome } from "./connections.js";
|
||||
import { HassApi } from "./hass.js";
|
||||
import { marksIn, readBinding, readIfThere } from "./mesh.js";
|
||||
import { probeBroker } from "./probe.js";
|
||||
|
||||
const dir = process.env.MESH_PROVISIONS_DIR ?? "/run/provisions";
|
||||
const url = process.env.MESH_HOMEASSISTANT_URL ?? "http://127.0.0.1:8123";
|
||||
const token = (await readIfThere(process.env.MESH_HOMEASSISTANT_TOKEN_FILE))?.trim() ?? "";
|
||||
const marks = marksIn(process.env.MESH_WRITTEN_DIR ?? "/var/lib/home-assistant-provisions");
|
||||
const waitSeconds = Number(process.env.MESH_HOMEASSISTANT_WAIT_SECONDS ?? "300");
|
||||
|
||||
if (!token) {
|
||||
console.error("[hass-provisions] no Home Assistant token — home-assistant's own `token` secret has not been accepted");
|
||||
process.exit(1);
|
||||
}
|
||||
|
||||
/** Home Assistant answers /api/ with 200 once it is up and the token is good. */
|
||||
async function ready(): Promise<boolean> {
|
||||
const until = Date.now() + waitSeconds * 1000;
|
||||
for (;;) {
|
||||
try {
|
||||
const res = await fetch(`${url.replace(/\/$/, "")}/api/`, { headers: { Authorization: `Bearer ${token}` } });
|
||||
if (res.status === 200) return true;
|
||||
if (res.status === 401 || res.status === 403) {
|
||||
console.error("[hass-provisions] Home Assistant refuses the token — accept a long-lived access token it issued");
|
||||
return false;
|
||||
}
|
||||
} catch {
|
||||
// not listening yet
|
||||
}
|
||||
if (Date.now() >= until) return false;
|
||||
await new Promise((r) => setTimeout(r, 3000));
|
||||
}
|
||||
}
|
||||
|
||||
if (!(await ready())) {
|
||||
console.error(`[hass-provisions] Home Assistant did not answer at ${url} within ${waitSeconds}s`);
|
||||
process.exit(1);
|
||||
}
|
||||
|
||||
const hass = new HassApi(url, token);
|
||||
const read = async (p: string) => [await readBinding(join(dir, `${p}.json`)), await readIfThere(join(dir, `${p}.secret`))] as const;
|
||||
|
||||
const outcomes: Outcome[] = [];
|
||||
{
|
||||
const [binding, secret] = await read(MQTT_PROVISION);
|
||||
outcomes.push(await reconcileMqtt({ hass, probe: probeBroker, marks }, binding, secret));
|
||||
}
|
||||
for (const spec of APPS) {
|
||||
const [binding, secret] = await read(spec.provision);
|
||||
outcomes.push(await reconcileApp({ hass, http: { fetch: (u, init) => fetch(u, init) } }, spec, binding, secret));
|
||||
}
|
||||
|
||||
let failed = 0;
|
||||
for (const o of outcomes) {
|
||||
switch (o.result) {
|
||||
case "unchanged":
|
||||
console.log(`[hass-provisions] ${o.what}: already as the mesh says${o.note ? ` — ${o.note}` : ""}`);
|
||||
break;
|
||||
case "written":
|
||||
console.log(`[hass-provisions] ${o.what}: wrote ${o.fields.join(", ")}; Home Assistant's own test passed${o.note ? ` — ${o.note}` : ""}`);
|
||||
break;
|
||||
case "equivalent":
|
||||
console.log(`[hass-provisions] ${o.what}: ${o.note}`);
|
||||
break;
|
||||
case "refused":
|
||||
failed++;
|
||||
console.error(`[hass-provisions] ${o.what}: ${o.problem}`);
|
||||
break;
|
||||
}
|
||||
}
|
||||
process.exitCode = failed > 0 ? 1 : 0;
|
||||
@@ -0,0 +1,42 @@
|
||||
// What the mesh delivered to home-assistant's provisions step, and the step's own small memory.
|
||||
//
|
||||
// Per provision it requires, the mesh writes two files beside each other (the manifest's `binds` and
|
||||
// `secrets`): `<provision>.json`, the binding — where the provider is (`at`), what it serves (`port`,
|
||||
// `scheme`, …) and the login this module presents (`as`) — and `<provision>.secret`, the pair
|
||||
// credential. Nothing here guesses a host, a port or a key.
|
||||
|
||||
import { mkdir, readFile, rename, writeFile } from "node:fs/promises";
|
||||
import { join } from "node:path";
|
||||
|
||||
import type { Binding, Marks } from "./connections.js";
|
||||
|
||||
/** A file the mesh wrote, or undefined when it is not there. */
|
||||
export async function readIfThere(path: string | undefined): Promise<string | undefined> {
|
||||
if (!path) return undefined;
|
||||
return readFile(path, "utf8").catch(() => undefined);
|
||||
}
|
||||
|
||||
/** A binding file parsed, or undefined when absent or not JSON. */
|
||||
export async function readBinding(path: string): Promise<Binding | undefined> {
|
||||
const raw = await readIfThere(path);
|
||||
if (raw === undefined) return undefined;
|
||||
try {
|
||||
return JSON.parse(raw) as Binding;
|
||||
} catch {
|
||||
return undefined;
|
||||
}
|
||||
}
|
||||
|
||||
export function marksIn(dir: string): Marks {
|
||||
return {
|
||||
async get(name) {
|
||||
return (await readIfThere(join(dir, `${name}.digest`)))?.trim() || undefined;
|
||||
},
|
||||
async set(name, value) {
|
||||
await mkdir(dir, { recursive: true, mode: 0o700 });
|
||||
const path = join(dir, `${name}.digest`);
|
||||
await writeFile(`${path}.tmp`, `${value}\n`, { mode: 0o600 });
|
||||
await rename(`${path}.tmp`, path);
|
||||
},
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,117 @@
|
||||
// Ask the broker, before Home Assistant is told anything, whether it takes the login and password
|
||||
// the mesh delivered — and whether that login may subscribe to Home Assistant's discovery topics.
|
||||
//
|
||||
// One MQTT 3.1.1 session: CONNECT (clean, a throwaway client id, so Home Assistant's own session is
|
||||
// never taken over), read the CONNACK, optionally SUBSCRIBE once and read the SUBACK, DISCONNECT.
|
||||
// No dependency: the handful of bytes MQTT needs for this are written here.
|
||||
|
||||
import { randomBytes } from "node:crypto";
|
||||
import { connect } from "node:net";
|
||||
|
||||
export interface ProbeResult {
|
||||
/** 0 accepted; 4 bad username or password; 5 not authorised. */
|
||||
connack: number;
|
||||
/** The SUBACK return code for the filter asked about: 0–2 granted, 0x80 refused. */
|
||||
suback?: number;
|
||||
}
|
||||
|
||||
export type Probe = (host: string, port: number, username: string, password: string, subscribe?: string) => Promise<ProbeResult>;
|
||||
|
||||
function str(v: string): Buffer {
|
||||
const b = Buffer.from(v, "utf8");
|
||||
const len = Buffer.alloc(2);
|
||||
len.writeUInt16BE(b.length);
|
||||
return Buffer.concat([len, b]);
|
||||
}
|
||||
|
||||
function packet(type: number, body: Buffer): Buffer {
|
||||
let remaining = body.length;
|
||||
const lenBytes: number[] = [];
|
||||
do {
|
||||
let byte = remaining % 128;
|
||||
remaining = Math.floor(remaining / 128);
|
||||
if (remaining > 0) byte |= 0x80;
|
||||
lenBytes.push(byte);
|
||||
} while (remaining > 0);
|
||||
return Buffer.concat([Buffer.from([type, ...lenBytes]), body]);
|
||||
}
|
||||
|
||||
/** The first complete packet in `buf`: its type byte, its body, and how many bytes it took. */
|
||||
export function firstPacket(buf: Buffer): { type: number; body: Buffer; used: number } | undefined {
|
||||
if (buf.length < 2) return undefined;
|
||||
let length = 0;
|
||||
let multiplier = 1;
|
||||
let i = 1;
|
||||
for (;;) {
|
||||
if (i >= buf.length) return undefined;
|
||||
const byte = buf[i++];
|
||||
length += (byte & 0x7f) * multiplier;
|
||||
if ((byte & 0x80) === 0) break;
|
||||
multiplier *= 128;
|
||||
if (i > 4) throw new Error("malformed MQTT remaining length");
|
||||
}
|
||||
if (buf.length < i + length) return undefined;
|
||||
return { type: buf[0], body: buf.subarray(i, i + length), used: i + length };
|
||||
}
|
||||
|
||||
export const probeBroker: Probe = (host, port, username, password, subscribe) => {
|
||||
const connectBody = Buffer.concat([
|
||||
str("MQTT"),
|
||||
Buffer.from([4, 0xc2, 0, 10]), // level 4 (3.1.1); username + password + clean session; keepalive 10s
|
||||
str(`mesh-probe-${randomBytes(6).toString("hex")}`),
|
||||
str(username),
|
||||
str(password),
|
||||
]);
|
||||
return new Promise((resolve, reject) => {
|
||||
const socket = connect({ host, port });
|
||||
let buf = Buffer.alloc(0);
|
||||
const result: ProbeResult = { connack: -1 };
|
||||
const timer = setTimeout(() => {
|
||||
socket.destroy();
|
||||
reject(new Error(`no answer from the broker at ${host}:${port} within 10s`));
|
||||
}, 10_000);
|
||||
const finish = (): void => {
|
||||
clearTimeout(timer);
|
||||
if (result.connack === 0) socket.end(Buffer.from([0xe0, 0]));
|
||||
else socket.destroy();
|
||||
resolve(result);
|
||||
};
|
||||
socket.on("connect", () => socket.write(packet(0x10, connectBody)));
|
||||
socket.on("data", (chunk) => {
|
||||
buf = Buffer.concat([buf, chunk]);
|
||||
for (;;) {
|
||||
let p;
|
||||
try {
|
||||
p = firstPacket(buf);
|
||||
} catch (err) {
|
||||
clearTimeout(timer);
|
||||
socket.destroy();
|
||||
reject(err);
|
||||
return;
|
||||
}
|
||||
if (!p) return;
|
||||
buf = buf.subarray(p.used);
|
||||
const kind = p.type >> 4;
|
||||
if (kind === 2) {
|
||||
result.connack = p.body[1] ?? -1;
|
||||
if (result.connack !== 0 || !subscribe) return finish();
|
||||
// SUBSCRIBE, packet id 1, one filter at QoS 0.
|
||||
socket.write(packet(0x82, Buffer.concat([Buffer.from([0, 1]), str(subscribe), Buffer.from([0])])));
|
||||
} else if (kind === 9) {
|
||||
result.suback = p.body[2];
|
||||
return finish();
|
||||
}
|
||||
}
|
||||
});
|
||||
socket.on("error", (err) => {
|
||||
clearTimeout(timer);
|
||||
reject(err);
|
||||
});
|
||||
socket.on("close", () => {
|
||||
if (result.connack === -1) {
|
||||
clearTimeout(timer);
|
||||
reject(new Error(`the broker at ${host}:${port} closed the connection without answering`));
|
||||
}
|
||||
});
|
||||
});
|
||||
};
|
||||
Reference in New Issue
Block a user