Compare commits

..
Author SHA1 Message Date
jschoubben 0010b6bf21 ombi: adopt the Plex entry that plex itself answers for
ace's ombi holds one Plex entry, loaded from an older server and later
retyped to plex's public name: its stored machineIdentifier is not plex's,
while its address answers as plex. Matching by identifier alone would leave
it and add a second entry for the same server.

When no entry carries the server's identifier, each entry's own address is
asked for /identity, and an entry plex answers for is adopted: the bound
connection laid over, and the identifier corrected (ombi builds its "view in
Plex" links from it). Its name, libraries and every other choice stay. An
entry that cannot be asked, or answers as another server, is left alone;
nothing is guessed. Every call of the step is now bounded, since an entry may
name a host that no longer answers.
2026-09-30 13:14:12 +02:00
jschoubben b068a9d399 ombi: reach plex through the mesh, in the same step as the Servarr apps
ombi reached plex at its public name, typed into its settings screen, so it
depended on plex's public route and on nobody moving plex. ombi now requires
plex-api, and the run-once step that writes its Servarr connections writes
its Plex one too - renamed from `servarr` to `connections`, since it is no
longer only that.

ombi keeps several Plex servers. The entry this provision names is found by
the server's own machineIdentifier (plex answers it at /identity, and ombi
stored it when the server was loaded), and only its host, port, TLS, base
path and token are written, only when they differ. Another server's entry,
the selected libraries, whether Plex is enabled and every other choice are
left alone. An ombi with no entry for the server gets one.

The token is tried against plex first. Until the operator accepts the
server's X-Plex-Token for this pair the mesh delivers a value it minted,
which plex refuses (401, or 400 on a network it trusts); refused, nothing is
written and the step fails naming the secret accept, so a working token in
ombi is never replaced by a dead one.

Tests import the compiled step, as keycloak's do: the step imports its
sibling with the .js specifier the build needs, which type stripping does
not resolve. `npm test` builds first.
2026-09-30 12:59:12 +02:00
19 changed files with 603 additions and 1374 deletions
+1 -4
View File
@@ -13,7 +13,7 @@ ARG RUNTIME_BASE
FROM ${BUILD_BASE} AS build FROM ${BUILD_BASE} AS build
WORKDIR /app/modules/home-assistant WORKDIR /app/modules/home-assistant
COPY . . COPY . .
RUN node /app/node_modules/typescript/bin/tsc client.ts index.ts tools/index.ts provisions/hass.ts provisions/probe.ts provisions/connections.ts provisions/mesh.ts provisions/index.ts \ RUN node /app/node_modules/typescript/bin/tsc client.ts index.ts tools/index.ts \
--module NodeNext --moduleResolution NodeNext --target ES2022 --outDir dist --module NodeNext --moduleResolution NodeNext --target ES2022 --outDir dist
FROM ${RUNTIME_BASE} FROM ${RUNTIME_BASE}
@@ -22,6 +22,3 @@ COPY --from=build /app/modules/home-assistant/dist /app/modules/home-assistant/d
# provider's provisioner runs its reconcile loop in the same process, with the broker connected — # provider's provisioner runs its reconcile loop in the same process, with the broker connected —
# the convention novox/hq issues 060/061 settled. # the convention novox/hq issues 060/061 settled.
ENV MESH_TOOL_MODULES=/app/modules/home-assistant/dist/index.js,/app/modules/home-assistant/dist/tools/index.js ENV MESH_TOOL_MODULES=/app/modules/home-assistant/dist/index.js,/app/modules/home-assistant/dist/tools/index.js
# NOT dist/provisions/index.js: that is a step the host runs to completion, named by the
# `provisions` container's args as `mesh-tools run …` (novox/hq ADR 0052). Listed here it would run
# inside the serving sidecar too, and exit it.
+13 -93
View File
@@ -18,21 +18,7 @@
"port": 8123, "port": 8123,
"protocol": "tcp", "protocol": "tcp",
"from": "mesh", "from": "mesh",
"why": "the dashboard, the API and the companion apps" "why": "the dashboard and the API"
},
{
"name": "sonos-events",
"port": 1400,
"protocol": "tcp",
"from": "mesh",
"why": "the Sonos integration's event callback: speakers push their state changes here"
},
{
"name": "webrtc",
"port": 18555,
"protocol": "tcp",
"from": "mesh",
"why": "the bundled go2rtc's WebRTC port, which camera streams to a browser use"
} }
], ],
"resources": [ "resources": [
@@ -42,33 +28,24 @@
"path": "/var/lib/mesh/home-assistant", "path": "/var/lib/mesh/home-assistant",
"mode": "0700" "mode": "0700"
}, },
{
"id": "state",
"type": "directory",
"mode": "0700",
"place": "."
},
{ {
"id": "config", "id": "config",
"type": "directory", "type": "directory",
"mode": "0700" "path": "/services/home-assistant/config",
}, "mode": "0700",
{ "owner": "1000:1000"
"id": "written",
"type": "directory",
"mode": "0700"
}, },
{ {
"id": "server", "id": "server",
"type": "container", "type": "container",
"name": "home-assistant", "name": "home-assistant",
"image": "ghcr.io/home-assistant/home-assistant@sha256:d8922685169707fd91e8b9729902d975f06157d005e422874d201e0261dda196", "image": "ghcr.io/home-assistant/home-assistant@sha256:14931c6b13756317849f46da1d01b45937a1150db66c081cfe529d48215943fe",
"network": "host", "network": "host",
"env": { "env": {
"TZ": "Etc/UTC" "TZ": "Etc/UTC"
}, },
"volumes": [ "volumes": [
"${dir:config}:/config" "/services/home-assistant/config:/config"
] ]
}, },
{ {
@@ -87,90 +64,33 @@
"volumes": [ "volumes": [
"/var/lib/mesh/home-assistant/broker:/run/secrets/broker:ro", "/var/lib/mesh/home-assistant/broker:/run/secrets/broker:ro",
"/var/lib/mesh/home-assistant/token:/run/secrets/token:ro", "/var/lib/mesh/home-assistant/token:/run/secrets/token:ro",
"/var/lib/mesh/home-assistant/config.json:/run/config/config.json:ro" "/var/lib/mesh/home-assistant/config.json:/run/config/config.json:ro",
"/services/home-assistant/config:/var/lib/home-assistant/config:ro"
], ],
"env": { "env": {
"MESH_BROKER_FILE": "/run/secrets/broker", "MESH_BROKER_FILE": "/run/secrets/broker",
"MESH_HOMEASSISTANT_URL": "http://127.0.0.1:${port:8123}", "MESH_HOMEASSISTANT_URL": "http://127.0.0.1:8123",
"MESH_HOMEASSISTANT_TOKEN_FILE": "/run/secrets/token", "MESH_HOMEASSISTANT_TOKEN_FILE": "/run/secrets/token",
"MESH_HOMEASSISTANT_CONFIG_FILE": "/run/config/config.json" "MESH_HOMEASSISTANT_CONFIG_FILE": "/run/config/config.json",
"MESH_HOMEASSISTANT_CONFIG_DIR": "/var/lib/home-assistant/config"
}, },
"restart-on": [ "restart-on": [
"runtime-config" "runtime-config"
], ],
"artifact": "runtime" "artifact": "runtime"
},
{
"id": "provisions",
"type": "container",
"name": "mesh-home-assistant-provisions",
"network": "host",
"run-once": true,
"volumes": [
"/var/lib/mesh/home-assistant/token:/run/secrets/token:ro",
"${dir:written}:/var/lib/home-assistant-provisions",
"${dir:state}/mqtt-topic.json:/run/provisions/mqtt-topic.json:ro",
"${dir:state}/mqtt-topic.secret:/run/provisions/mqtt-topic.secret:ro",
"${dir:state}/sonarr-api.json:/run/provisions/sonarr-api.json:ro",
"${dir:state}/sonarr-api.secret:/run/provisions/sonarr-api.secret:ro",
"${dir:state}/radarr-api.json:/run/provisions/radarr-api.json:ro",
"${dir:state}/radarr-api.secret:/run/provisions/radarr-api.secret:ro",
"${dir:state}/lidarr-api.json:/run/provisions/lidarr-api.json:ro",
"${dir:state}/lidarr-api.secret:/run/provisions/lidarr-api.secret:ro"
],
"env": {
"MESH_HOMEASSISTANT_URL": "http://127.0.0.1:${port:8123}",
"MESH_HOMEASSISTANT_TOKEN_FILE": "/run/secrets/token",
"MESH_PROVISIONS_DIR": "/run/provisions",
"MESH_WRITTEN_DIR": "/var/lib/home-assistant-provisions"
},
"args": [
"run",
"/app/modules/home-assistant/dist/provisions/index.js"
],
"restart-on": [
"bound-mqtt-topic",
"secret-mqtt-topic",
"bound-sonarr-api",
"secret-sonarr-api",
"bound-radarr-api",
"secret-radarr-api",
"bound-lidarr-api",
"secret-lidarr-api"
],
"artifact": "runtime"
} }
], ],
"requires": [ "requires": [
"lidarr-api", "route"
"mqtt-topic",
"radarr-api",
"route",
"sonarr-api"
], ],
"contributes": { "contributes": {
"mqtt-topic": {
"topics": [
"#"
]
},
"route": { "route": {
"label": "home-assistant", "label": "home-assistant",
"endpoint": "web" "endpoint": "web"
} }
}, },
"binds": { "binds": {
"route": "${dir:state}/route.json", "route": "/var/lib/mesh/home-assistant/route.json"
"mqtt-topic": "${dir:state}/mqtt-topic.json",
"sonarr-api": "${dir:state}/sonarr-api.json",
"radarr-api": "${dir:state}/radarr-api.json",
"lidarr-api": "${dir:state}/lidarr-api.json"
},
"secrets": {
"mqtt-topic": "${dir:state}/mqtt-topic.secret",
"sonarr-api": "${dir:state}/sonarr-api.secret",
"radarr-api": "${dir:state}/radarr-api.secret",
"lidarr-api": "${dir:state}/lidarr-api.secret"
}, },
"build": { "build": {
"on": [ "on": [
+1 -6
View File
@@ -1,14 +1,9 @@
{ {
"name": "@novox/module-home-assistant", "name": "@novox/module-home-assistant",
"version": "0.1.0", "version": "0.1.0",
"description": "home-assistant \u2014 home automation platform. Its API client, tools and events live here (novox/hq ADR 0039).", "description": "home-assistant — home automation platform. Its API client, tools and events live here (novox/hq ADR 0039).",
"type": "module", "type": "module",
"private": true, "private": true,
"scripts": {
"build": "tsc client.ts index.ts tools/index.ts provisions/hass.ts provisions/probe.ts provisions/connections.ts provisions/mesh.ts provisions/index.ts --module NodeNext --moduleResolution NodeNext --target ES2022 --outDir dist",
"typecheck": "tsc -p tsconfig.json",
"test": "node --test --experimental-strip-types 'test/*.test.ts'"
},
"dependencies": { "dependencies": {
"@novox/mesh-sdk": "^0.1.0" "@novox/mesh-sdk": "^0.1.0"
}, },
@@ -1,486 +0,0 @@
// 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)}` };
}
const shown = formValues(flow.data_schema);
// A new entry's form has no value for its two certificate choices (a reconfigure pre-fills them
// from the entry): plain MQTT, so neither a CA nor a client certificate.
const other = (shown.other_settings ?? {}) as Record<string, unknown>;
if (flow.data_schema?.some((f) => f.name === "other_settings")) {
shown.other_settings = { set_ca_cert: "off", set_client_cert: false, ...other };
}
flow = await deps.hass.stepFlow(flow.flow_id, {
...shown,
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) };
}
}
-160
View File
@@ -1,160 +0,0 @@
// 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"}`)));
}
};
});
}
}
@@ -1,84 +0,0 @@
// 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;
-42
View File
@@ -1,42 +0,0 @@
// 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);
},
};
}
-117
View File
@@ -1,117 +0,0 @@
// 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`));
}
});
});
};
@@ -1,293 +0,0 @@
// What holds home-assistant's provisions step (provisions/*.ts): Home Assistant's MQTT entry is
// made to use the broker, port and login the mesh bound — only after the broker takes that login,
// through the reconfigure flow, keeping every other setting as Home Assistant pre-filled it, and not
// again once it already says so; its Sonarr/Radarr/Lidarr entries are made, finished (reauth), left
// alone when they already reach the bound app, and never removed; a key the app refuses (the mesh's
// minted value before the operator accepts the app's) is never written.
//
// Home Assistant and the apps are fakes answering as the real ones do (flow shapes checked against
// ghcr.io/home-assistant/home-assistant 2026.9.3, the build ace runs).
import { test } from "node:test";
import assert from "node:assert/strict";
import type { ConfigEntry, DeviceEntry, FlowProgress, FlowResult, Hass, SchemaField } from "../provisions/hass.ts";
import type { Binding, Marks } from "../provisions/connections.ts";
import { APPS, formValues, reconcileApp, reconcileMqtt, sameUrl, type Http, type ServarrApp } from "../provisions/connections.ts";
import type { Probe } from "../provisions/probe.ts";
const PWD_NOT_CHANGED = "__**password_not_changed**__";
const MINTED = "mesh-minted-password";
function mqttBinding(): Binding {
return { provision: "mqtt-topic", from: "ace", at: "ace.internal", as: "mesh_ace_hass", serves: { scheme: "mqtt", port: 1883 } };
}
/** The MQTT reconfigure form as Home Assistant serializes it, pre-filled from an entry. */
function brokerForm(data: Record<string, unknown>): SchemaField[] {
return [
{ name: "broker", type: "string", required: true, description: { suggested_value: data.broker } },
{ name: "port", type: "integer", required: true, default: 1883, description: { suggested_value: data.port } },
{ name: "protocol", type: "select", required: true, default: "3.1.1", description: { suggested_value: data.protocol } },
{ name: "username", type: "string", optional: true, description: { suggested_value: data.username } },
{ name: "password", type: "string", optional: true, description: { suggested_value: data.password ? PWD_NOT_CHANGED : undefined } },
{
name: "other_settings",
type: "expandable",
required: true,
schema: [
{ name: "keepalive", type: "integer", optional: true, description: { suggested_value: 60 } },
{ name: "transport", type: "select", required: true, default: "tcp", description: { suggested_value: "tcp" } },
{ name: "set_ca_cert", type: "select", required: true, description: { suggested_value: "off" } },
{ name: "set_client_cert", type: "boolean", required: true, description: { suggested_value: false } },
],
},
];
}
interface FakeOpts {
entries?: Record<string, (ConfigEntry & { data: Record<string, unknown> })[]>;
/** What Home Assistant's own connection test accepts. */
accepts?: (data: Record<string, unknown>) => boolean;
reauth?: FlowProgress[];
devices?: DeviceEntry[];
}
function fakeHass(opts: FakeOpts = {}) {
const entries = opts.entries ?? {};
const calls: string[] = [];
const submitted: Record<string, unknown>[] = [];
const flows = new Map<string, { handler: string; entryId?: string; step: string; reauth?: boolean }>();
let n = 0;
const accepts = opts.accepts ?? (() => true);
const form = (id: string, step: string, schema: SchemaField[], errors?: Record<string, string>): FlowResult => ({
type: "form", flow_id: id, step_id: step, data_schema: schema, errors: errors ?? null,
});
const servarrUser: SchemaField[] = [
{ name: "url", type: "string", required: true },
{ name: "api_key", type: "string", required: true },
{ name: "more_options", type: "expandable", required: true, schema: [{ name: "verify_ssl", type: "boolean", optional: true, default: false }] },
];
const hass: Hass = {
async entries(domain) {
calls.push(`entries ${domain}`);
return (entries[domain] ?? []).map(({ data: _d, ...e }) => e);
},
async startFlow(handler, entryId) {
calls.push(`start ${handler}${entryId ? ` ${entryId}` : ""}`);
const id = `f${++n}`;
if (handler === "mqtt") {
const entry = entryId ? entries.mqtt.find((e) => e.entry_id === entryId) : undefined;
if (entryId && !entry) return { type: "abort", reason: "not_found" };
flows.set(id, { handler, entryId, step: "broker" });
return form(id, "broker", brokerForm(entry?.data ?? {}));
}
if (entryId) return { type: "abort", reason: "not_implemented" }; // no reconfigure for Servarr
flows.set(id, { handler, step: "user" });
return form(id, "user", servarrUser);
},
async stepFlow(flowId, input) {
calls.push(`step ${flowId}`);
const flow = flows.get(flowId);
if (!flow) throw new Error(`Home Assistant POST flow/${flowId} answered 404`);
if (flow.step === "reauth_confirm") {
flow.step = "user";
return form(flowId, "user", [
{ name: "url", type: "string", required: true, default: "http://old:1" },
{ name: "api_key", type: "string", optional: true },
{ name: "verify_ssl", type: "boolean", optional: true, default: false },
]);
}
submitted.push(input);
if (flow.handler === "mqtt") {
const entry = entries.mqtt?.find((e) => e.entry_id === flow.entryId);
const data = { ...input, ...(input.password === PWD_NOT_CHANGED ? { password: entry?.data.password } : {}) };
if (!accepts(data)) return form(flowId, "broker", brokerForm(data), { base: "cannot_connect" });
flows.delete(flowId);
if (entry) {
entry.data = data;
return { type: "abort", reason: "reconfigure_successful" };
}
(entries.mqtt ??= []).push({ entry_id: "new-mqtt", domain: "mqtt", state: "loaded", data });
return { type: "create_entry", result: { entry_id: "new-mqtt" } };
}
if (!accepts(input)) return form(flowId, "user", servarrUser, { base: "invalid_auth" });
flows.delete(flowId);
if (flow.reauth) return { type: "abort", reason: "reauth_successful" };
(entries[flow.handler] ??= []).push({ entry_id: `new-${flow.handler}`, domain: flow.handler, state: "loaded", data: input });
return { type: "create_entry", result: { entry_id: `new-${flow.handler}` } };
},
async abortFlow(flowId) {
calls.push(`abort ${flowId}`);
flows.delete(flowId);
},
async flowsInProgress() {
for (const f of opts.reauth ?? []) flows.set(f.flow_id, { handler: f.handler, entryId: f.context?.entry_id, step: "reauth_confirm", reauth: true });
return opts.reauth ?? [];
},
async devices() {
return opts.devices ?? [];
},
};
return { hass, calls, submitted, entries };
}
function memoryMarks(): Marks & { store: Map<string, string> } {
const store = new Map<string, string>();
return { store, get: async (k) => store.get(k), set: async (k, v) => void store.set(k, v) };
}
const takes = (suback = 0): Probe => async (_h, _p, user, pass) => ({ connack: user === "mesh_ace_hass" && pass === MINTED ? 0 : 5, suback });
const aceMqttEntry = () => ({
entry_id: "7d1e", domain: "mqtt", state: "loaded",
data: { broker: "127.0.0.1", port: 1883, protocol: "5", username: "luffy", password: "luffys-password" },
});
test("mqtt: the broker is asked first; a login it does not take is never written", async () => {
const f = fakeHass({ entries: { mqtt: [aceMqttEntry()] } });
const out = await reconcileMqtt({ hass: f.hass, probe: async () => ({ connack: 5 }), marks: memoryMarks() }, mqttBinding(), MINTED);
assert.equal(out.result, "refused");
assert.match((out as { problem: string }).problem, /does not \(yet\) take the login mesh_ace_hass/);
assert.deepEqual(f.calls, []); // Home Assistant not even asked
assert.equal(f.entries.mqtt[0].data.username, "luffy");
});
test("mqtt: ace's entry (127.0.0.1, luffy) is moved to the bound broker and login, every other setting kept", async () => {
const f = fakeHass({ entries: { mqtt: [aceMqttEntry()] } });
const marks = memoryMarks();
const out = await reconcileMqtt({ hass: f.hass, probe: takes(), marks }, mqttBinding(), `${MINTED}\n`);
assert.deepEqual(out, { what: "mqtt", result: "written", fields: ["broker", "username", "password"] });
assert.deepEqual(f.entries.mqtt[0].data, {
broker: "ace.internal", port: 1883, protocol: "5", username: "mesh_ace_hass", password: MINTED,
other_settings: { keepalive: 60, transport: "tcp", set_ca_cert: "off", set_client_cert: false },
});
assert.ok(marks.store.get("mqtt"));
assert.ok(![...marks.store.values()].some((v) => v.includes(MINTED)));
// Run again: nothing differs, the flow is opened to read and closed without submitting.
const before = f.submitted.length;
const again = await reconcileMqtt({ hass: f.hass, probe: takes(), marks }, mqttBinding(), MINTED);
assert.deepEqual(again, { what: "mqtt", result: "unchanged" });
assert.equal(f.submitted.length, before);
assert.match(f.calls.at(-1) ?? "", /^abort /);
});
test("mqtt: a new password alone is written (the digest tells)", async () => {
const f = fakeHass({ entries: { mqtt: [aceMqttEntry()] } });
const marks = memoryMarks();
await reconcileMqtt({ hass: f.hass, probe: takes(), marks }, mqttBinding(), MINTED);
const rotated: Probe = async () => ({ connack: 0, suback: 0 });
const out = await reconcileMqtt({ hass: f.hass, probe: rotated, marks }, mqttBinding(), "rotated");
assert.deepEqual(out, { what: "mqtt", result: "written", fields: ["password"] });
assert.equal(f.entries.mqtt[0].data.password, "rotated");
});
test("mqtt: Home Assistant's own connection test refusing saves nothing and fails loudly", async () => {
const f = fakeHass({ entries: { mqtt: [aceMqttEntry()] }, accepts: () => false });
const marks = memoryMarks();
const out = await reconcileMqtt({ hass: f.hass, probe: takes(), marks }, mqttBinding(), MINTED);
assert.equal(out.result, "refused");
assert.match((out as { problem: string }).problem, /cannot_connect.*unchanged/);
assert.equal(f.entries.mqtt[0].data.username, "luffy");
assert.equal(marks.store.size, 0);
});
test("mqtt: a fresh Home Assistant gets an entry; a grant without the discovery topics is warned about", async () => {
const f = fakeHass();
const out = await reconcileMqtt({ hass: f.hass, probe: takes(0x80), marks: memoryMarks() }, mqttBinding(), MINTED);
assert.equal(out.result, "written");
assert.match((out as { note?: string }).note ?? "", /may not subscribe to homeassistant\/#/);
assert.equal(f.entries.mqtt[0].data.broker, "ace.internal");
});
test("mqtt: two entries, or a binding without a port, are refused rather than guessed", async () => {
const f = fakeHass({ entries: { mqtt: [aceMqttEntry(), { ...aceMqttEntry(), entry_id: "other" }] } });
assert.equal((await reconcileMqtt({ hass: f.hass, probe: takes(), marks: memoryMarks() }, mqttBinding(), MINTED)).result, "refused");
const noPort = { ...mqttBinding(), serves: {} };
assert.match(((await reconcileMqtt({ hass: f.hass, probe: takes(), marks: memoryMarks() }, noPort, MINTED)) as { problem: string }).problem, /no usable port/);
});
// ---- Servarr ----
const SONARR = APPS.find((a) => a.domain === "sonarr") as ServarrApp;
const RADARR = APPS.find((a) => a.domain === "radarr") as ServarrApp;
const KEY = "the-apps-own-key";
const servarrBinding = (port: number, at = "ace.internal"): Binding => ({ provision: "sonarr-api", from: "ace", at, as: "mesh_ace_hass", serves: { scheme: "http", port, "url-base": "" } });
/** One running Sonarr, answering on several addresses (127.0.0.1 and ace.internal are one host). */
function apps(instances: Record<string, { startTime: string }>): Http & { asked: string[] } {
const asked: string[] = [];
return {
asked,
async fetch(url, init) {
asked.push(url);
const u = new URL(url);
const inst = instances[`${u.hostname}:${u.port}`];
if (!inst) throw new Error("connect ECONNREFUSED");
if (init?.headers?.["X-Api-Key"] !== KEY) return { status: 401, text: async () => "" };
return { status: 200, text: async () => JSON.stringify({ version: "4.0.15", appData: "/config", startTime: inst.startTime }) };
},
};
}
const oneSonarr = () => apps({ "ace.internal:8989": { startTime: "t1" }, "127.0.0.1:8989": { startTime: "t1" } });
const sonarrEntry = (state = "loaded") => ({ entry_id: "5a1d", domain: "sonarr", state, data: { url: "http://127.0.0.1:8989", api_key: KEY } });
test("servarr: the mesh's minted key is never written; the remedy names the accept", async () => {
const f = fakeHass({ entries: { sonarr: [sonarrEntry()] } });
const out = await reconcileApp({ hass: f.hass, http: oneSonarr() }, SONARR, servarrBinding(8989), "minted-by-the-mesh");
assert.equal(out.result, "refused");
assert.match((out as { problem: string }).problem, /secret accept <this node> home-assistant sonarr-api --provider ace/);
assert.deepEqual(f.calls, []);
});
test("servarr: ace's entry at 127.0.0.1 reaches the same Sonarr the mesh bound at ace.internal — left, and said", async () => {
const f = fakeHass({ entries: { sonarr: [sonarrEntry()] }, devices: [{ id: "d", config_entries: ["5a1d"], configuration_url: "http://127.0.0.1:8989" }] });
const out = await reconcileApp({ hass: f.hass, http: oneSonarr() }, SONARR, servarrBinding(8989), KEY);
assert.equal(out.result, "equivalent");
assert.equal(f.submitted.length, 0);
});
test("servarr: an entry already at the bound URL is unchanged", async () => {
const f = fakeHass({ entries: { sonarr: [sonarrEntry()] }, devices: [{ id: "d", config_entries: ["5a1d"], configuration_url: "http://ace.internal:8989" }] });
assert.deepEqual(await reconcileApp({ hass: f.hass, http: oneSonarr() }, SONARR, servarrBinding(8989), KEY), { what: "sonarr", result: "unchanged" });
});
test("servarr: a working entry that reaches a different app is refused, and nothing is removed", async () => {
const f = fakeHass({ entries: { sonarr: [sonarrEntry()] }, devices: [{ id: "d", config_entries: ["5a1d"], configuration_url: "http://127.0.0.1:8989" }] });
const two = apps({ "ace.internal:8989": { startTime: "t1" }, "127.0.0.1:8989": { startTime: "another" } });
const out = await reconcileApp({ hass: f.hass, http: two }, SONARR, servarrBinding(8989), KEY);
assert.equal(out.result, "refused");
assert.match((out as { problem: string }).problem, /never does/);
assert.equal(f.entries.sonarr.length, 1);
});
test("servarr: no entry — one is made at the bound URL through the user flow", async () => {
const f = fakeHass({ entries: {} });
const out = await reconcileApp({ hass: f.hass, http: oneSonarr() }, SONARR, servarrBinding(8989), KEY);
assert.deepEqual(out, { what: "sonarr", result: "written", fields: ["entry"] });
assert.deepEqual(f.submitted[0], { url: "http://ace.internal:8989", api_key: KEY, more_options: { verify_ssl: false } });
});
test("servarr: a reauth Home Assistant started is finished with the bound URL and key", async () => {
const entry = { ...sonarrEntry("setup_error"), domain: "radarr", entry_id: "1955" };
const f = fakeHass({
entries: { radarr: [entry] },
reauth: [{ flow_id: "r1", handler: "radarr", step_id: "reauth_confirm", context: { source: "reauth", entry_id: "1955" } }],
});
const radarr = apps({ "ace.internal:7878": { startTime: "t" } });
const out = await reconcileApp({ hass: f.hass, http: radarr }, RADARR, { ...servarrBinding(7878), provision: "radarr-api" }, KEY);
assert.deepEqual(out, { what: "radarr", result: "written", fields: ["api_key", "url"] });
assert.deepEqual(f.submitted[0], { url: "http://ace.internal:7878", api_key: KEY, verify_ssl: false });
});
test("form values: suggested first, then default, sections nested", () => {
assert.deepEqual(formValues(brokerForm({ broker: "b", port: 1, protocol: "5", username: "u", password: "p" })), {
broker: "b", port: 1, protocol: "5", username: "u", password: PWD_NOT_CHANGED,
other_settings: { keepalive: 60, transport: "tcp", set_ca_cert: "off", set_client_cert: false },
});
assert.ok(sameUrl("http://ace.internal:8989/", "http://ace.internal:8989"));
assert.ok(sameUrl("http://ACE.internal", "http://ace.internal:80"));
assert.ok(!sameUrl("http://127.0.0.1:8989", "http://ace.internal:8989"));
});
+1 -10
View File
@@ -8,14 +8,5 @@
"skipLibCheck": true, "skipLibCheck": true,
"noEmit": true "noEmit": true
}, },
"include": [ "include": ["client.ts", "index.ts", "tools/index.ts"]
"client.ts",
"index.ts",
"tools/index.ts",
"provisions/hass.ts",
"provisions/probe.ts",
"provisions/connections.ts",
"provisions/mesh.ts",
"provisions/index.ts"
]
} }
+2 -2
View File
@@ -13,7 +13,7 @@ ARG RUNTIME_BASE
FROM ${BUILD_BASE} AS build FROM ${BUILD_BASE} AS build
WORKDIR /app/modules/ombi WORKDIR /app/modules/ombi
COPY . . COPY . .
RUN node /app/node_modules/typescript/bin/tsc client.ts index.ts tools/index.ts servarr/settings.ts servarr/index.ts \ RUN node /app/node_modules/typescript/bin/tsc client.ts index.ts tools/index.ts servarr/settings.ts plex/settings.ts connections/index.ts \
--module NodeNext --moduleResolution NodeNext --target ES2022 --outDir dist --module NodeNext --moduleResolution NodeNext --target ES2022 --outDir dist
FROM ${RUNTIME_BASE} FROM ${RUNTIME_BASE}
@@ -22,6 +22,6 @@ COPY --from=build /app/modules/ombi/dist /app/modules/ombi/dist
# provider's provisioner runs its reconcile loop in the same process, with the broker connected — # provider's provisioner runs its reconcile loop in the same process, with the broker connected —
# the convention novox/hq issues 060/061 settled. # the convention novox/hq issues 060/061 settled.
ENV MESH_TOOL_MODULES=/app/modules/ombi/dist/index.js,/app/modules/ombi/dist/tools/index.js ENV MESH_TOOL_MODULES=/app/modules/ombi/dist/index.js,/app/modules/ombi/dist/tools/index.js
# NOT dist/servarr/index.js: that is a step the host runs to completion, named by the `servarr` # NOT dist/connections/index.js: that is a step the host runs to completion, named by the `connections`
# container's args as `mesh-tools run …` (novox/hq ADR 0052). Listed here it would run inside the # container's args as `mesh-tools run …` (novox/hq ADR 0052). Listed here it would run inside the
# serving sidecar too, and exit it. # serving sidecar too, and exit it.
+66
View File
@@ -0,0 +1,66 @@
// ombi's connections step — run once by the host after ombi's server starts, and run again whenever
// a binding or pair credential it reads changes (the container's `restart-on`, novox/hq ADR 0099).
// It brings ombi's connections to Sonarr, Radarr, Lidarr (servarr/settings.ts) and Plex
// (plex/settings.ts) in line with what the mesh bound.
//
// **A step, not a loop**, for the reason route-adapter gives: everything it does is a function of
// files the mesh writes, and the host already knows when they change. It connects to no broker.
//
// Exits non-zero when any app could not be put right — a refused credential, an unreachable app, an
// ombi that cannot reach it — so the node reports the step failed and the host runs it again on the
// next apply. It is declared last in the manifest, so its failing gates nothing else of ombi's
// (novox/hq ADR 0136). One app failing does not stop the others being put right.
//
// Reads, per provision, `<dir>/<provision>.json` (the binding) and `<dir>/<provision>.secret` (the
// pair credential), where <dir> is MESH_CONNECTIONS_DIR. Never prints a key or a token.
import { join } from "node:path";
import { PLEX_PROVISION, reconcilePlex } from "../plex/settings.js";
import { APPS, ombiReady, readBinding, readIfThere, reconcileApp, type Http, type Outcome } from "../servarr/settings.js";
const dir = process.env.MESH_CONNECTIONS_DIR ?? "/run/connections";
const url = process.env.MESH_OMBI_URL ?? "http://127.0.0.1:3579";
const apiKey = (await readIfThere(process.env.MESH_OMBI_API_KEY_FILE))?.trim() ?? process.env.MESH_OMBI_API_KEY ?? "";
const waitSeconds = Number(process.env.MESH_OMBI_WAIT_SECONDS ?? "180");
// Every call bounded: an entry ombi keeps may name a host that no longer answers, and a step that
// hangs on it holds the apply.
const http: Http = { fetch: (u, init) => fetch(u, { ...init, signal: AbortSignal.timeout(20_000) }) };
if (!apiKey) {
console.error("[ombi-connections] no ombi API key — ombi's own `api-key` secret has not been accepted");
process.exit(1);
}
const ombi = { url, apiKey };
if (!(await ombiReady(http, ombi, waitSeconds * 1000))) {
console.error(`[ombi-connections] ombi did not answer at ${url} within ${waitSeconds}s`);
process.exit(1);
}
const inputs = async (provision: string) =>
[await readBinding(join(dir, `${provision}.json`)), await readIfThere(join(dir, `${provision}.secret`))] as const;
const outcomes: Outcome[] = [];
for (const spec of APPS) {
outcomes.push(await reconcileApp(http, ombi, spec, ...(await inputs(spec.provision))));
}
outcomes.push(await reconcilePlex(http, ombi, ...(await inputs(PLEX_PROVISION))));
let failed = 0;
for (const outcome of outcomes) {
switch (outcome.result) {
case "unchanged":
console.log(`[ombi-connections] ${outcome.app}: already as the mesh says; connection tested`);
break;
case "written":
console.log(`[ombi-connections] ${outcome.app}: wrote ${outcome.fields.join(", ")}; connection tested`);
break;
case "refused":
failed++;
console.error(`[ombi-connections] ${outcome.app}: ${outcome.problem}`);
break;
}
}
process.exitCode = failed > 0 ? 1 : 0;
+20 -13
View File
@@ -87,28 +87,30 @@
"artifact": "runtime" "artifact": "runtime"
}, },
{ {
"id": "servarr", "id": "connections",
"type": "container", "type": "container",
"name": "mesh-ombi-servarr", "name": "mesh-ombi-connections",
"network": "host", "network": "host",
"run-once": true, "run-once": true,
"volumes": [ "volumes": [
"/var/lib/mesh/ombi/api-key:/run/secrets/api-key:ro", "/var/lib/mesh/ombi/api-key:/run/secrets/api-key:ro",
"${dir:state}/sonarr-api.json:/run/servarr/sonarr-api.json:ro", "${dir:state}/sonarr-api.json:/run/connections/sonarr-api.json:ro",
"${dir:state}/sonarr-api.secret:/run/servarr/sonarr-api.secret:ro", "${dir:state}/sonarr-api.secret:/run/connections/sonarr-api.secret:ro",
"${dir:state}/radarr-api.json:/run/servarr/radarr-api.json:ro", "${dir:state}/radarr-api.json:/run/connections/radarr-api.json:ro",
"${dir:state}/radarr-api.secret:/run/servarr/radarr-api.secret:ro", "${dir:state}/radarr-api.secret:/run/connections/radarr-api.secret:ro",
"${dir:state}/lidarr-api.json:/run/servarr/lidarr-api.json:ro", "${dir:state}/lidarr-api.json:/run/connections/lidarr-api.json:ro",
"${dir:state}/lidarr-api.secret:/run/servarr/lidarr-api.secret:ro" "${dir:state}/lidarr-api.secret:/run/connections/lidarr-api.secret:ro",
"${dir:state}/plex-api.json:/run/connections/plex-api.json:ro",
"${dir:state}/plex-api.secret:/run/connections/plex-api.secret:ro"
], ],
"env": { "env": {
"MESH_OMBI_URL": "http://127.0.0.1:${port:3579}", "MESH_OMBI_URL": "http://127.0.0.1:${port:3579}",
"MESH_OMBI_API_KEY_FILE": "/run/secrets/api-key", "MESH_OMBI_API_KEY_FILE": "/run/secrets/api-key",
"MESH_SERVARR_DIR": "/run/servarr" "MESH_CONNECTIONS_DIR": "/run/connections"
}, },
"args": [ "args": [
"run", "run",
"/app/modules/ombi/dist/servarr/index.js" "/app/modules/ombi/dist/connections/index.js"
], ],
"restart-on": [ "restart-on": [
"bound-sonarr-api", "bound-sonarr-api",
@@ -116,13 +118,16 @@
"bound-radarr-api", "bound-radarr-api",
"secret-radarr-api", "secret-radarr-api",
"bound-lidarr-api", "bound-lidarr-api",
"secret-lidarr-api" "secret-lidarr-api",
"bound-plex-api",
"secret-plex-api"
], ],
"artifact": "runtime" "artifact": "runtime"
} }
], ],
"requires": [ "requires": [
"lidarr-api", "lidarr-api",
"plex-api",
"radarr-api", "radarr-api",
"route", "route",
"sonarr-api" "sonarr-api"
@@ -137,12 +142,14 @@
"route": "${dir:state}/route.json", "route": "${dir:state}/route.json",
"sonarr-api": "${dir:state}/sonarr-api.json", "sonarr-api": "${dir:state}/sonarr-api.json",
"radarr-api": "${dir:state}/radarr-api.json", "radarr-api": "${dir:state}/radarr-api.json",
"lidarr-api": "${dir:state}/lidarr-api.json" "lidarr-api": "${dir:state}/lidarr-api.json",
"plex-api": "${dir:state}/plex-api.json"
}, },
"secrets": { "secrets": {
"sonarr-api": "${dir:state}/sonarr-api.secret", "sonarr-api": "${dir:state}/sonarr-api.secret",
"radarr-api": "${dir:state}/radarr-api.secret", "radarr-api": "${dir:state}/radarr-api.secret",
"lidarr-api": "${dir:state}/lidarr-api.secret" "lidarr-api": "${dir:state}/lidarr-api.secret",
"plex-api": "${dir:state}/plex-api.secret"
}, },
"build": { "build": {
"on": [ "on": [
+2 -2
View File
@@ -5,9 +5,9 @@
"type": "module", "type": "module",
"private": true, "private": true,
"scripts": { "scripts": {
"build": "tsc client.ts index.ts tools/index.ts servarr/settings.ts servarr/index.ts --module NodeNext --moduleResolution NodeNext --target ES2022 --outDir dist", "build": "tsc client.ts index.ts tools/index.ts servarr/settings.ts plex/settings.ts connections/index.ts --module NodeNext --moduleResolution NodeNext --target ES2022 --outDir dist",
"typecheck": "tsc -p tsconfig.json", "typecheck": "tsc -p tsconfig.json",
"test": "node --test --experimental-strip-types 'test/*.test.ts'" "test": "npm run build && node --test --experimental-strip-types 'test/*.test.ts'"
}, },
"dependencies": { "dependencies": {
"@novox/mesh-sdk": "^0.1.0" "@novox/mesh-sdk": "^0.1.0"
+271
View File
@@ -0,0 +1,271 @@
// Where ombi reaches Plex — decided by the mesh, written into ombi by ombi's own API.
//
// **Why this exists.** ombi keeps its Plex servers in its own database (OmbiSettings.db), so the
// mesh has no file to write `${bound:plex-api:at}` into. ombi requires `plex-api`; the mesh delivers
// a binding (where plex is: `at`, and what it serves: `port`, `scheme`) and a pair credential (the
// server owner's X-Plex-Token, accepted by the operator — plex.tv issues it and the mesh cannot
// mint it). This step makes ombi's Plex settings say the same thing, beside its Servarr ones.
//
// **Which entry is plex's.** ombi may list several Plex servers. The one this provision names is
// found by the server's own machineIdentifier, which plex answers at /identity — the same value
// ombi stored when an operator loaded the server in its settings screen. That entry's connection is
// brought in line; an entry for any other server is never touched.
//
// When no entry carries that identifier, an entry may still be this server reached another way:
// ace's ombi holds one loaded from an older server and later retyped to plex's public name, so its
// stored identifier is stale while its address answers as this plex. Each entry's OWN address is
// asked for /identity, and an entry plex itself answers for is this server's — adopted: its
// connection laid over and its identifier corrected (ombi builds its "view in Plex" links from it).
// Nothing is guessed: an entry whose address is unreachable, or answers as another server, is left
// as it was. Only when no entry is this server's either way is one added, named as plex names
// itself — it is the mesh's, so later runs keep it true.
//
// **Only the connection, and only when it differs.** Host, port, TLS, base path and token — and the
// identifier of an adopted entry. Whether Plex is enabled in ombi, watchlist import, the selected
// libraries, the batch size and everything else an operator chose are left exactly as they are.
//
// **A token plex refuses is never written.** Until the operator accepts the server's token for this
// pair, the mesh delivers a value it minted itself, which plex answers with 401 (or 400 on its own
// network). Writing it would replace a working token in ombi with a dead one, so the token is tried
// against plex first; refused, nothing is written and the step fails naming the `secret accept`.
import { isLoopback, ombiCall, type Binding, type Http, type Ombi, type Outcome } from "../servarr/settings.js";
/** The provision ombi requires for Plex — the manifest's `requires`, `binds` and `secrets` key. */
export const PLEX_PROVISION = "plex-api";
/** The connection fields ombi keeps for a Plex server — the only ones this step ever writes. */
export interface PlexConnection {
ip: string;
port: number;
ssl: boolean;
subDir: string | null;
plexAuthToken: string;
}
export type PlexWanted = { ok: true; connection: PlexConnection; from: string } | { ok: false; problem: string };
/** The connection the mesh says ombi should use, from the binding and the pair credential. */
export function wantedPlex(binding: Binding | undefined, credential: string | undefined): PlexWanted {
if (!binding) {
return { ok: false, problem: `no binding for ${PLEX_PROVISION} was delivered — the mesh writes it before this step runs` };
}
const at = typeof binding.at === "string" ? binding.at.trim() : "";
const serves = binding.serves ?? {};
const port = Number(serves.port);
if (!at) return { ok: false, problem: `the ${PLEX_PROVISION} binding names no host (at)` };
if (isLoopback(at)) {
return {
ok: false,
problem:
`the ${PLEX_PROVISION} binding says plex is at ${at}, which from ombi's own container is ombi itself. ` +
`The mesh hands loopback to a machine that is not on the private network; put it on the private ` +
`network so plex has an address ombi can dial`,
};
}
if (!Number.isInteger(port) || port <= 0 || port > 65535) {
return { ok: false, problem: `the ${PLEX_PROVISION} binding serves no usable port (${String(serves.port)})` };
}
const scheme = typeof serves.scheme === "string" && serves.scheme ? serves.scheme : "http";
if (scheme !== "http" && scheme !== "https") {
return { ok: false, problem: `the ${PLEX_PROVISION} binding serves scheme ${scheme}, which ombi cannot dial` };
}
const token = (credential ?? "").trim();
if (!token) return { ok: false, problem: `the ${PLEX_PROVISION} credential is empty or was not delivered` };
return {
ok: true,
from: typeof binding.from === "string" ? binding.from : "",
connection: { ip: at, port, ssl: scheme === "https", subDir: null, plexAuthToken: token },
};
}
/** plex's base URL as the step dials it — the same host and port ombi will be given. */
export function plexUrl(want: PlexConnection): string {
const host = want.ip.includes(":") && !want.ip.startsWith("[") ? `[${want.ip}]` : want.ip;
return `${want.ssl ? "https" : "http"}://${host}:${want.port}`;
}
/** Which connection fields differ between an entry ombi holds and what the mesh says. Names only. */
export function differingPlex(current: Record<string, unknown> | undefined, want: PlexConnection): (keyof PlexConnection)[] {
const now = current ?? {};
const out: (keyof PlexConnection)[] = [];
if (String(now.ip ?? "") !== want.ip) out.push("ip");
if (Number(now.port ?? 0) !== want.port) out.push("port");
if (Boolean(now.ssl) !== want.ssl) out.push("ssl");
const sub = typeof now.subDir === "string" && now.subDir.trim() !== "" ? now.subDir : null;
if (sub !== want.subDir) out.push("subDir");
if (String(now.plexAuthToken ?? "") !== want.plexAuthToken) out.push("plexAuthToken");
return out;
}
async function plexGet(http: Http, want: PlexConnection, path: string, withToken: boolean) {
const headers: Record<string, string> = { Accept: "application/json" };
if (withToken) headers["X-Plex-Token"] = want.plexAuthToken;
return http.fetch(`${plexUrl(want)}${path}`, { method: "GET", headers });
}
/**
* Does plex take this token? `true` it does; `false` it refused it — 401 or 403, or 400, which is
* what plex answers a token it never issued on a network it trusts. A thrown error when plex could
* not be asked.
*/
export async function plexTakes(http: Http, want: PlexConnection): Promise<{ takes: boolean; friendlyName?: string }> {
const res = await plexGet(http, want, "/", true);
if (res.status === 400 || res.status === 401 || res.status === 403) return { takes: false };
if (res.status < 200 || res.status >= 300) throw new Error(`plex answered ${res.status} at /`);
let friendlyName: string | undefined;
try {
const body = JSON.parse(await res.text()) as { MediaContainer?: { friendlyName?: unknown } };
if (typeof body.MediaContainer?.friendlyName === "string") friendlyName = body.MediaContainer.friendlyName;
} catch {
// a name is a nicety for a new entry, not a condition
}
return { takes: true, friendlyName };
}
/** The server's own machineIdentifier, which plex answers without a token. */
export async function plexIdentity(http: Http, want: Pick<PlexConnection, "ip" | "port" | "ssl" | "subDir">): Promise<string> {
const host = want.ip.includes(":") && !want.ip.startsWith("[") ? `[${want.ip}]` : want.ip;
const base = `${want.ssl ? "https" : "http"}://${host}:${want.port}${want.subDir ? `/${want.subDir.replace(/^\/+|\/+$/g, "")}` : ""}`;
const res = await http.fetch(`${base}/identity`, { method: "GET", headers: { Accept: "application/json" } });
if (res.status !== 200) throw new Error(`plex answered ${res.status} at /identity`);
const body = JSON.parse(await res.text()) as { MediaContainer?: { machineIdentifier?: unknown } };
const id = body.MediaContainer?.machineIdentifier;
if (typeof id !== "string" || id === "") throw new Error("plex's /identity names no machineIdentifier");
return id;
}
/** The remedy for a refused token, in the controller's own words (ADR 0092). */
export function plexAcceptRemedy(from: string): string {
return (
`plex refuses the ${PLEX_PROVISION} credential the mesh delivered, so it was not written into ombi. ` +
`plex's token is issued by plex.tv and the mesh cannot make it: accept the server's own token for ` +
`this pair — \`secret accept <this node> ombi ${PLEX_PROVISION} --provider ${from || "<its node>"} ` +
`--from <file holding the server's X-Plex-Token>\``
);
}
/** The batch size ombi's settings screen fills in for a server it adds ("150 by default"). */
const EPISODE_BATCH_SIZE = 150;
/** ombi's server entries, as its settings document holds them (null on a fresh ombi). */
export function serversOf(document: Record<string, unknown> | undefined): Record<string, unknown>[] {
const servers = document?.servers;
return Array.isArray(servers) ? (servers as Record<string, unknown>[]) : [];
}
/**
* Which entries, holding another identifier, plex answers for at their own address — this server,
* reached another way. Asked only when no entry carries the identifier. An entry that cannot be
* asked is not this server's: nothing is guessed.
*/
export async function answeringAs(http: Http, servers: Record<string, unknown>[], machineIdentifier: string): Promise<Set<number>> {
const out = new Set<number>();
if (servers.some((s) => s?.machineIdentifier === machineIdentifier)) return out;
for (const [i, s] of servers.entries()) {
const ip = typeof s?.ip === "string" ? s.ip.trim() : "";
const port = Number(s?.port);
if (!ip || !Number.isInteger(port) || port <= 0 || port > 65535) continue;
const subDir = typeof s.subDir === "string" && s.subDir.trim() !== "" ? s.subDir : null;
try {
if ((await plexIdentity(http, { ip, port, ssl: Boolean(s.ssl), subDir })) === machineIdentifier) out.add(i);
} catch {
// unreachable, or not a plex: not this server's
}
}
return out;
}
/**
* ombi's Plex settings with this server's connection laid over them: every entry naming the
* server's machineIdentifier — or adopted, its address answering as this server — gets the
* connection (an adopted one also the identifier), every other entry is left as it was, and when
* none is this server's one is added. Returns the document to save and the fields that changed.
*/
export function withPlexServer(
document: Record<string, unknown> | undefined,
machineIdentifier: string,
want: PlexConnection,
name: string,
adopted: ReadonlySet<number> = new Set(),
): { next: Record<string, unknown>; fields: string[]; added: boolean; entry: Record<string, unknown> } {
const doc = document ?? {};
const servers = serversOf(doc);
const fields = new Set<string>();
let entry: Record<string, unknown> | undefined;
const next = servers.map((s, i) => {
const adopt = adopted.has(i) && s?.machineIdentifier !== machineIdentifier;
if (s?.machineIdentifier !== machineIdentifier && !adopt) return s;
for (const f of differingPlex(s, want)) fields.add(f);
if (adopt) fields.add("machineIdentifier");
const laid = { ...s, machineIdentifier, ip: want.ip, port: want.port, ssl: want.ssl, subDir: want.subDir, plexAuthToken: want.plexAuthToken };
entry ??= laid;
return laid;
});
if (entry) return { next: { ...doc, servers: next }, fields: [...fields], added: false, entry };
const added: Record<string, unknown> = {
name,
machineIdentifier,
ip: want.ip,
port: want.port,
ssl: want.ssl,
subDir: want.subDir,
plexAuthToken: want.plexAuthToken,
episodeBatchSize: EPISODE_BATCH_SIZE,
plexSelectedLibraries: [],
};
return { next: { ...doc, servers: [...next, added] }, fields: ["server"], added: true, entry: added };
}
/**
* Bring ombi's connection to plex in line with the mesh: check the token against plex, find the
* server's entry by its machineIdentifier, write only the connection fields when they differ (or add
* the entry), then have ombi test the connection from its own container. Never throws.
*/
export async function reconcilePlex(http: Http, ombi: Ombi, binding: Binding | undefined, credential: string | undefined): Promise<Outcome> {
const app = "plex";
const w = wantedPlex(binding, credential);
// `in`, not `!w.ok`: the Dockerfile compiles without strict, where a boolean discriminant does not
// narrow.
if ("problem" in w) return { app, result: "refused", problem: w.problem };
const want = w.connection;
let name: string;
let machineIdentifier: string;
try {
const taken = await plexTakes(http, want);
if (!taken.takes) return { app, result: "refused", problem: plexAcceptRemedy(w.from) };
machineIdentifier = await plexIdentity(http, want);
name = taken.friendlyName || "Plex";
} catch (err) {
return { app, result: "refused", problem: `plex could not be asked whether it takes the token at ${want.ip}:${want.port}: ${message(err)}` };
}
try {
const document = (await ombiCall(http, ombi, "GET", "/Settings/Plex")) as Record<string, unknown> | undefined;
const adopted = await answeringAs(http, serversOf(document), machineIdentifier);
const laid = withPlexServer(document, machineIdentifier, want, name, adopted);
if (laid.fields.length > 0) {
const saved = await ombiCall(http, ombi, "POST", "/Settings/Plex", laid.next);
if (saved === false) return { app, result: "refused", problem: "ombi declined to save its Plex settings" };
}
// ombi's own test, from ombi's own container — the path the step's check above did not take.
const tested = await ombiCall(http, ombi, "POST", "/Tester/plex", laid.entry);
if (tested !== true) {
return {
app,
result: "refused",
problem:
`ombi cannot reach plex at ${want.ip}:${want.port} from its own container` +
(laid.fields.length > 0 ? `; its settings were written (${laid.fields.join(", ")})` : ""),
};
}
return laid.fields.length > 0 ? { app, result: "written", fields: laid.fields } : { app, result: "unchanged" };
} catch (err) {
return { app, result: "refused", problem: message(err) };
}
}
function message(err: unknown): string {
return err instanceof Error ? err.message : String(err);
}
-59
View File
@@ -1,59 +0,0 @@
// ombi's Servarr step — run once by the host after ombi's server starts, and run again whenever a
// binding or pair credential it reads changes (the container's `restart-on`, novox/hq ADR 0099).
//
// **A step, not a loop**, for the reason route-adapter gives: everything it does is a function of
// files the mesh writes, and the host already knows when they change. It connects to no broker.
//
// Exits non-zero when any app could not be put right — a refused credential, an unreachable app, an
// ombi that cannot reach it — so the node reports the step failed and the host runs it again on the
// next apply. It is declared last in the manifest, so its failing gates nothing else of ombi's
// (novox/hq ADR 0136).
//
// Reads, per app, `<dir>/<provision>.json` (the binding) and `<dir>/<provision>.secret` (the pair
// credential), where <dir> is MESH_SERVARR_DIR. Never prints a key.
import { join } from "node:path";
import { APPS, ombiReady, readBinding, readIfThere, reconcileApp, type Http } from "./settings.js";
const dir = process.env.MESH_SERVARR_DIR ?? "/run/servarr";
const url = process.env.MESH_OMBI_URL ?? "http://127.0.0.1:3579";
const apiKey = (await readIfThere(process.env.MESH_OMBI_API_KEY_FILE))?.trim() ?? process.env.MESH_OMBI_API_KEY ?? "";
const waitSeconds = Number(process.env.MESH_OMBI_WAIT_SECONDS ?? "180");
const http: Http = { fetch: (u, init) => fetch(u, init) };
if (!apiKey) {
console.error("[ombi-servarr] no ombi API key — ombi's own `api-key` secret has not been accepted");
process.exit(1);
}
const ombi = { url, apiKey };
if (!(await ombiReady(http, ombi, waitSeconds * 1000))) {
console.error(`[ombi-servarr] ombi did not answer at ${url} within ${waitSeconds}s`);
process.exit(1);
}
let failed = 0;
for (const spec of APPS) {
const outcome = await reconcileApp(
http,
ombi,
spec,
await readBinding(join(dir, `${spec.provision}.json`)),
await readIfThere(join(dir, `${spec.provision}.secret`)),
);
switch (outcome.result) {
case "unchanged":
console.log(`[ombi-servarr] ${outcome.app}: already as the mesh says; connection tested`);
break;
case "written":
console.log(`[ombi-servarr] ${outcome.app}: wrote ${outcome.fields.join(", ")}; connection tested`);
break;
case "refused":
failed++;
console.error(`[ombi-servarr] ${outcome.app}: ${outcome.problem}`);
break;
}
}
process.exitCode = failed > 0 ? 1 : 0;
+2 -2
View File
@@ -115,7 +115,7 @@ export function subDirOf(urlBase: unknown): string | null {
return trimmed === "" ? null : trimmed; return trimmed === "" ? null : trimmed;
} }
function isLoopback(host: string): boolean { export function isLoopback(host: string): boolean {
const h = host.toLowerCase(); const h = host.toLowerCase();
return h === "localhost" || h === "::1" || h === "[::1]" || /^127\./.test(h); return h === "localhost" || h === "::1" || h === "[::1]" || /^127\./.test(h);
} }
@@ -163,7 +163,7 @@ export interface Ombi {
apiKey: string; apiKey: string;
} }
async function ombiCall(http: Http, ombi: Ombi, method: string, path: string, body?: unknown): Promise<unknown> { export async function ombiCall(http: Http, ombi: Ombi, method: string, path: string, body?: unknown): Promise<unknown> {
const res = await http.fetch(`${ombi.url.replace(/\/$/, "")}/api/v1${path}`, { const res = await http.fetch(`${ombi.url.replace(/\/$/, "")}/api/v1${path}`, {
method, method,
headers: { headers: {
+223
View File
@@ -0,0 +1,223 @@
// What holds ombi's Plex step (plex/settings.ts): the entry for the server plex says it is — found
// by machineIdentifier — is made to say what the mesh bound (host, port, TLS, token) and nothing
// else it keeps is touched; an entry for another server is left alone; an ombi with no entry for it
// gets one; nothing is written when nothing differs; and a token plex refuses (the mesh's own minted
// value, before the operator accepts the server's token) is never written, with the `secret accept`
// that fixes it named.
//
// ombi and plex are fakes answering the routes the step touches as the real ones do (checked against
// lscr.io/linuxserver/ombi 4.53.10 and plexinc/pms-docker 1.43.4: plex answers 401 to an unknown
// token from another network and 400 on one it trusts; ombi's /Tester/plex answers a bare boolean).
//
// Imports the compiled step, as keycloak's tests do: plex/settings.ts imports its sibling with the
// `.js` specifier the build needs, which Node's type stripping does not resolve to a `.ts` file.
import { test } from "node:test";
import assert from "node:assert/strict";
import { differingPlex, reconcilePlex, wantedPlex, withPlexServer } from "../dist/plex/settings.js";
import type { Binding, Http } from "../servarr/settings.ts";
const TOKEN = "the-servers-own-token";
const MACHINE = "5c47d9a165d10b622995d55b3ae1f168242f33bd";
const OMBI = { url: "http://127.0.0.1:3579", apiKey: "ombi-key" };
function binding(port = 32400, at = "ace.internal", scheme = "http"): Binding {
return { binding: 1, provision: "plex-api", from: "ace", at, as: "mesh_ace_ombi", serves: { scheme, port } } as Binding;
}
interface Call {
method: string;
url: string;
body?: unknown;
}
function fakes(plexSettings: Record<string, unknown>, opts: { reachable?: boolean; trusted?: boolean } = {}) {
const calls: Call[] = [];
const store = { plex: plexSettings };
const http: Http = {
async fetch(url, init) {
const method = init?.method ?? "GET";
const body = init?.body ? (JSON.parse(init.body) as unknown) : undefined;
calls.push({ method, url, body });
const reply = (status: number, value?: unknown) => ({
status,
text: async () => (value === undefined ? "" : JSON.stringify(value)),
});
const u = new URL(url);
// Other servers an entry may name: a friend's, and plex's own public name (the same server).
if (u.hostname === "10.0.0.9") return reply(200, { MediaContainer: { machineIdentifier: "another-server" } });
if (u.hostname === "gone.example") throw new Error("getaddrinfo ENOTFOUND");
if (u.hostname === "plex.zurag.be") {
if (u.pathname === "/identity") return reply(200, { MediaContainer: { machineIdentifier: MACHINE } });
return reply(401);
}
if (u.port === "32400" || u.hostname === "ace.internal") {
if (opts.reachable === false) throw new Error("connect ECONNREFUSED");
if (u.pathname === "/identity") return reply(200, { MediaContainer: { machineIdentifier: MACHINE } });
const token = init?.headers?.["X-Plex-Token"];
if (token !== TOKEN) return reply(opts.trusted ? 400 : 401);
return reply(200, { MediaContainer: { friendlyName: "ace", machineIdentifier: MACHINE } });
}
if (init?.headers?.ApiKey !== "ombi-key") return reply(401);
if (u.pathname === "/api/v1/Settings/Plex" && method === "GET") return reply(200, store.plex);
if (u.pathname === "/api/v1/Settings/Plex" && method === "POST") {
store.plex = body as Record<string, unknown>;
return reply(200, true);
}
if (u.pathname === "/api/v1/Tester/plex") {
const tried = body as { plexAuthToken?: string; ip?: string };
return reply(200, tried.plexAuthToken === TOKEN && tried.ip === "ace.internal");
}
return reply(404);
},
};
return { http, calls, store };
}
// What an operator's ombi holds: plex loaded through its public name, plus a friend's server.
const operatorPlex = () => ({
enable: true,
enableWatchlistImport: true,
monitorAll: false,
installId: "b358a2a2-2ab0-4025-a3f3-450313c3c418",
servers: [
{
name: "ace", plexAuthToken: TOKEN, machineIdentifier: MACHINE, episodeBatchSize: 150,
serverHostname: "https://app.plex.tv", plexSelectedLibraries: [{ key: "1", title: "Films", enabled: true }],
ssl: true, subDir: null, ip: "plex.zurag.be", port: 443, id: 1,
},
{
name: "a friend", plexAuthToken: "their-token", machineIdentifier: "another-server", episodeBatchSize: 150,
plexSelectedLibraries: [], ssl: false, subDir: null, ip: "10.0.0.9", port: 32400, id: 2,
},
],
id: 4,
});
test("the server's own entry gets the bound connection; its libraries and every other setting stay", async () => {
const f = fakes(operatorPlex());
const out = await reconcilePlex(f.http, OMBI, binding(), `${TOKEN}\n`);
assert.deepEqual(out, { app: "plex", result: "written", fields: ["ip", "port", "ssl"] });
const want = operatorPlex();
Object.assign(want.servers[0], { ip: "ace.internal", port: 32400, ssl: false });
assert.deepEqual(f.store.plex, want);
});
test("another server's entry is never touched", async () => {
const f = fakes(operatorPlex());
await reconcilePlex(f.http, OMBI, binding(), TOKEN);
const servers = f.store.plex.servers as Record<string, unknown>[];
assert.deepEqual(servers[1], operatorPlex().servers[1]);
});
test("nothing is written when ombi already says what the mesh says", async () => {
const doc = operatorPlex();
Object.assign(doc.servers[0], { ip: "ace.internal", port: 32400, ssl: false });
const f = fakes(doc);
const out = await reconcilePlex(f.http, OMBI, binding(), TOKEN);
assert.deepEqual(out, { app: "plex", result: "unchanged" });
assert.equal(f.calls.filter((c) => c.method === "POST" && c.url.includes("/Settings/")).length, 0);
});
test("an ombi with no entry for this server gets one, named as plex names itself", async () => {
const fresh = { enable: false, enableWatchlistImport: false, monitorAll: false, installId: "x", servers: null, id: 0 };
const f = fakes(fresh);
const out = await reconcilePlex(f.http, OMBI, binding(), TOKEN);
assert.deepEqual(out, { app: "plex", result: "written", fields: ["server"] });
assert.deepEqual(f.store.plex, {
...fresh,
servers: [{
name: "ace", machineIdentifier: MACHINE, ip: "ace.internal", port: 32400, ssl: false, subDir: null,
plexAuthToken: TOKEN, episodeBatchSize: 150, plexSelectedLibraries: [],
}],
});
assert.equal(f.store.plex.enable, false, "whether plex is enabled in ombi is the operator's choice");
});
test("a token plex refuses is never written, and the accept that fixes it is named", async () => {
for (const trusted of [false, true]) {
const f = fakes(operatorPlex(), { trusted });
const out = await reconcilePlex(f.http, OMBI, binding(), "a-value-the-mesh-minted");
assert.equal(out.result, "refused");
const problem = (out as { problem: string }).problem;
assert.match(problem, /secret accept <this node> ombi plex-api --provider ace/);
assert.doesNotMatch(problem, /a-value-the-mesh-minted/);
assert.deepEqual(f.store.plex, operatorPlex(), "ombi's working settings were left alone");
assert.equal(f.calls.some((c) => c.url.includes("/api/v1/")), false, "ombi was not even asked");
}
});
test("a plex it cannot reach is reported, and ombi is left alone", async () => {
const f = fakes(operatorPlex(), { reachable: false });
const out = await reconcilePlex(f.http, OMBI, binding(), TOKEN);
assert.equal(out.result, "refused");
assert.match((out as { problem: string }).problem, /could not be asked.*ECONNREFUSED/);
assert.deepEqual(f.store.plex, operatorPlex());
});
test("a loopback binding is refused: from ombi's container it is ombi itself", () => {
const w = wantedPlex(binding(32400, "127.0.0.1"), TOKEN);
assert.equal(w.ok, false);
assert.match((w as { problem: string }).problem, /private network/);
});
test("an https binding sets ombi's ssl flag; an empty subDir is none", () => {
const w = wantedPlex(binding(32400, "ace.internal", "https"), TOKEN);
assert.equal(w.ok && w.connection.ssl, true);
assert.deepEqual(
differingPlex({ ip: "h", port: 1, ssl: false, subDir: "", plexAuthToken: "k" }, { ip: "h", port: 1, ssl: false, subDir: null, plexAuthToken: "k" }),
[],
);
});
test("every entry naming the server is laid over, not only the first", () => {
const doc = { servers: [{ machineIdentifier: MACHINE, ip: "a" }, { machineIdentifier: MACHINE, ip: "b" }] };
const want = { ip: "ace.internal", port: 32400, ssl: false, subDir: null, plexAuthToken: TOKEN };
const laid = withPlexServer(doc, MACHINE, want, "ace");
assert.equal(laid.added, false);
assert.deepEqual((laid.next.servers as { ip: string }[]).map((s) => s.ip), ["ace.internal", "ace.internal"]);
});
// ace's own ombi: its one entry was loaded from an older server (a stale identifier) and retyped to
// plex's public name, so it IS this server, reached another way (read from ace, 2026-09-30).
const acesOmbi = () => ({
enable: true,
enableWatchlistImport: true,
servers: [{
name: "Nami", plexAuthToken: TOKEN, machineIdentifier: "76562198623e708eef85b46aedb72c8f2fe671aa", episodeBatchSize: 0,
plexSelectedLibraries: [1, 2, 3, 4, 5, 6].map((k) => ({ key: String(k), enabled: true })), ssl: true, subDir: null,
ip: "plex.zurag.be", port: 443, id: 1,
}],
id: 4,
});
test("an entry whose own address answers as this server is adopted: connection and identifier, nothing else", async () => {
const f = fakes(acesOmbi());
const out = await reconcilePlex(f.http, OMBI, binding(), TOKEN);
assert.deepEqual(out, { app: "plex", result: "written", fields: ["ip", "port", "ssl", "machineIdentifier"] });
const want = acesOmbi();
Object.assign(want.servers[0], { ip: "ace.internal", port: 32400, ssl: false, machineIdentifier: MACHINE });
assert.deepEqual(f.store.plex, want, "one entry, still named Nami, its six libraries kept; none added");
});
test("an entry answering as another server, or not at all, is not adopted; this server gets its own", async () => {
const doc = {
servers: [
{ name: "a friend", machineIdentifier: "stale-1", ip: "10.0.0.9", port: 32400, ssl: false, plexAuthToken: "theirs" },
{ name: "gone", machineIdentifier: "stale-2", ip: "gone.example", port: 32400, ssl: false, plexAuthToken: "old" },
],
};
const f = fakes(structuredClone(doc));
const out = await reconcilePlex(f.http, OMBI, binding(), TOKEN);
assert.deepEqual(out, { app: "plex", result: "written", fields: ["server"] });
const servers = f.store.plex.servers as Record<string, unknown>[];
assert.deepEqual(servers.slice(0, 2), doc.servers, "both left exactly as they were");
assert.equal(servers[2].machineIdentifier, MACHINE);
});
test("no entry is probed once one carries the server's identifier", async () => {
const f = fakes(operatorPlex());
await reconcilePlex(f.http, OMBI, binding(), TOKEN);
assert.equal(f.calls.some((c) => c.url.startsWith("http://10.0.0.9")), false, "the friend's server was not asked");
});
+1 -1
View File
@@ -8,5 +8,5 @@
"skipLibCheck": true, "skipLibCheck": true,
"noEmit": true "noEmit": true
}, },
"include": ["client.ts", "index.ts", "tools/index.ts", "servarr/settings.ts", "servarr/index.ts"] "include": ["client.ts", "index.ts", "tools/index.ts", "servarr/settings.ts", "plex/settings.ts", "connections/index.ts"]
} }