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.
This commit is contained in:
@@ -24,7 +24,9 @@ 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 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 waitSeconds = Number(process.env.MESH_OMBI_WAIT_SECONDS ?? "180");
|
||||||
|
|
||||||
const http: Http = { fetch: (u, init) => fetch(u, init) };
|
// 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) {
|
if (!apiKey) {
|
||||||
console.error("[ombi-connections] no ombi API key — ombi's own `api-key` secret has not been accepted");
|
console.error("[ombi-connections] no ombi API key — ombi's own `api-key` secret has not been accepted");
|
||||||
|
|||||||
@@ -9,12 +9,20 @@
|
|||||||
// **Which entry is plex's.** ombi may list several Plex servers. The one this provision names is
|
// **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
|
// 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
|
// 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 is this server's,
|
// brought in line; an entry for any other server is never touched.
|
||||||
// one is 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. Whether
|
// When no entry carries that identifier, an entry may still be this server reached another way:
|
||||||
// Plex is enabled in ombi, watchlist import, the selected libraries, the batch size and everything
|
// ace's ombi holds one loaded from an older server and later retyped to plex's public name, so its
|
||||||
// else an operator chose are left exactly as they are.
|
// 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
|
// **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
|
// pair, the mesh delivers a value it minted itself, which plex answers with 401 (or 400 on its own
|
||||||
@@ -116,8 +124,10 @@ export async function plexTakes(http: Http, want: PlexConnection): Promise<{ tak
|
|||||||
}
|
}
|
||||||
|
|
||||||
/** The server's own machineIdentifier, which plex answers without a token. */
|
/** The server's own machineIdentifier, which plex answers without a token. */
|
||||||
export async function plexIdentity(http: Http, want: PlexConnection): Promise<string> {
|
export async function plexIdentity(http: Http, want: Pick<PlexConnection, "ip" | "port" | "ssl" | "subDir">): Promise<string> {
|
||||||
const res = await plexGet(http, want, "/identity", false);
|
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`);
|
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 body = JSON.parse(await res.text()) as { MediaContainer?: { machineIdentifier?: unknown } };
|
||||||
const id = body.MediaContainer?.machineIdentifier;
|
const id = body.MediaContainer?.machineIdentifier;
|
||||||
@@ -138,25 +148,57 @@ export function plexAcceptRemedy(from: string): string {
|
|||||||
/** The batch size ombi's settings screen fills in for a server it adds ("150 by default"). */
|
/** The batch size ombi's settings screen fills in for a server it adds ("150 by default"). */
|
||||||
const EPISODE_BATCH_SIZE = 150;
|
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
|
* ombi's Plex settings with this server's connection laid over them: every entry naming the
|
||||||
* server's machineIdentifier gets the connection, every other entry is left as it was, and when
|
* server's machineIdentifier — or adopted, its address answering as this server — gets the
|
||||||
* none names it one is added. Returns the document to save and the fields that changed.
|
* 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(
|
export function withPlexServer(
|
||||||
document: Record<string, unknown> | undefined,
|
document: Record<string, unknown> | undefined,
|
||||||
machineIdentifier: string,
|
machineIdentifier: string,
|
||||||
want: PlexConnection,
|
want: PlexConnection,
|
||||||
name: string,
|
name: string,
|
||||||
|
adopted: ReadonlySet<number> = new Set(),
|
||||||
): { next: Record<string, unknown>; fields: string[]; added: boolean; entry: Record<string, unknown> } {
|
): { next: Record<string, unknown>; fields: string[]; added: boolean; entry: Record<string, unknown> } {
|
||||||
const doc = document ?? {};
|
const doc = document ?? {};
|
||||||
const servers = Array.isArray(doc.servers) ? (doc.servers as Record<string, unknown>[]) : [];
|
const servers = serversOf(doc);
|
||||||
const fields = new Set<string>();
|
const fields = new Set<string>();
|
||||||
let entry: Record<string, unknown> | undefined;
|
let entry: Record<string, unknown> | undefined;
|
||||||
const next = servers.map((s) => {
|
const next = servers.map((s, i) => {
|
||||||
if (s?.machineIdentifier !== machineIdentifier) return s;
|
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);
|
for (const f of differingPlex(s, want)) fields.add(f);
|
||||||
const laid = { ...s, ip: want.ip, port: want.port, ssl: want.ssl, subDir: want.subDir, plexAuthToken: want.plexAuthToken };
|
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;
|
entry ??= laid;
|
||||||
return laid;
|
return laid;
|
||||||
});
|
});
|
||||||
@@ -201,7 +243,8 @@ export async function reconcilePlex(http: Http, ombi: Ombi, binding: Binding | u
|
|||||||
|
|
||||||
try {
|
try {
|
||||||
const document = (await ombiCall(http, ombi, "GET", "/Settings/Plex")) as Record<string, unknown> | undefined;
|
const document = (await ombiCall(http, ombi, "GET", "/Settings/Plex")) as Record<string, unknown> | undefined;
|
||||||
const laid = withPlexServer(document, machineIdentifier, want, name);
|
const adopted = await answeringAs(http, serversOf(document), machineIdentifier);
|
||||||
|
const laid = withPlexServer(document, machineIdentifier, want, name, adopted);
|
||||||
if (laid.fields.length > 0) {
|
if (laid.fields.length > 0) {
|
||||||
const saved = await ombiCall(http, ombi, "POST", "/Settings/Plex", laid.next);
|
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" };
|
if (saved === false) return { app, result: "refused", problem: "ombi declined to save its Plex settings" };
|
||||||
|
|||||||
@@ -45,6 +45,13 @@ function fakes(plexSettings: Record<string, unknown>, opts: { reachable?: boolea
|
|||||||
text: async () => (value === undefined ? "" : JSON.stringify(value)),
|
text: async () => (value === undefined ? "" : JSON.stringify(value)),
|
||||||
});
|
});
|
||||||
const u = new URL(url);
|
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 (u.port === "32400" || u.hostname === "ace.internal") {
|
||||||
if (opts.reachable === false) throw new Error("connect ECONNREFUSED");
|
if (opts.reachable === false) throw new Error("connect ECONNREFUSED");
|
||||||
if (u.pathname === "/identity") return reply(200, { MediaContainer: { machineIdentifier: MACHINE } });
|
if (u.pathname === "/identity") return reply(200, { MediaContainer: { machineIdentifier: MACHINE } });
|
||||||
@@ -171,3 +178,46 @@ test("every entry naming the server is laid over, not only the first", () => {
|
|||||||
assert.equal(laid.added, false);
|
assert.equal(laid.added, false);
|
||||||
assert.deepEqual((laid.next.servers as { ip: string }[]).map((s) => s.ip), ["ace.internal", "ace.internal"]);
|
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");
|
||||||
|
});
|
||||||
|
|||||||
Reference in New Issue
Block a user