diff --git a/src/backend/automations/engine.ts b/src/backend/automations/engine.ts index 1998ec871..354f90647 100644 --- a/src/backend/automations/engine.ts +++ b/src/backend/automations/engine.ts @@ -11,6 +11,7 @@ import { import { createCurrentAutomationRepository } from "../database/repositories/factory.js"; import { statsLogger } from "../utils/logger.js"; import { resolveHostById } from "../hosts/host-resolver.js"; +import { notifyAutomationInternalEvent } from "../hosts/automation-events.js"; import { executeStep } from "./actions/index.js"; import type { StepExecutionContext, StepResult } from "./actions/types.js"; import { compare } from "./conditions.js"; @@ -267,21 +268,20 @@ export class AutomationEngine { // its own failure, so the event is not emitted for runs that this event // already started. if (request.triggerType !== "internal_event") { - import("../hosts/automation-events.js") - .then(({ notifyAutomationInternalEvent }) => - notifyAutomationInternalEvent( - "automation_failed", - automation.userId, - undefined, - { - automationId: automation.id, - automationName: automation.name, - runId: run.id, - error: error ?? null, - }, - ), - ) - .catch(() => undefined); + // Emitted onto the bus rather than imported from hosts/: the bus is + // fire-and-forget, so this no longer needs the lazy-import dance that + // was working around the repositories cycle. + notifyAutomationInternalEvent( + "automation_failed", + automation.userId, + undefined, + { + automationId: automation.id, + automationName: automation.name, + runId: run.id, + error: error ?? null, + }, + ); } } diff --git a/src/backend/automations/scheduler.ts b/src/backend/automations/scheduler.ts index 2b2cc9458..6a35c770f 100644 --- a/src/backend/automations/scheduler.ts +++ b/src/backend/automations/scheduler.ts @@ -7,6 +7,10 @@ import { hasDwelled, isCoolingDown } from "./conditions.js"; import { pollDockerEvents } from "./docker-watcher.js"; import { AutomationEngine } from "./engine.js"; import { reconcileHeadlessViewers } from "./headless-viewer.js"; +import { + subscribeAutomationTriggers, + unsubscribeAutomationTriggers, +} from "./triggers.js"; /** * The one timer the automations feature owns. @@ -32,6 +36,10 @@ let ticking = false; export function startAutomationScheduler(): void { if (tickTimer) return; + // Event-driven triggers come through the bus; the interval below only drives + // the schedule-based ones. + subscribeAutomationTriggers(); + startupTimer = setTimeout(() => { void tick(); }, STARTUP_DELAY_MS); @@ -48,6 +56,7 @@ export function stopAutomationScheduler(): void { if (startupTimer) clearTimeout(startupTimer); tickTimer = null; startupTimer = null; + unsubscribeAutomationTriggers(); } /** Exposed for tests; the interval calls this. */ diff --git a/src/backend/automations/triggers.ts b/src/backend/automations/triggers.ts index ed88f5502..6df0751a2 100644 --- a/src/backend/automations/triggers.ts +++ b/src/backend/automations/triggers.ts @@ -16,6 +16,7 @@ import { type MetricsSnapshot, } from "./conditions.js"; import { AutomationEngine } from "./engine.js"; +import { pluginEvents, TOPICS } from "../plugins/events.js"; /** * Matches events against automation triggers and decides what fires. @@ -346,6 +347,43 @@ export async function onInternalEvent(event: InternalEvent): Promise { } } +let unsubscribers: Array<() => void> = []; + +/** + * Subscribes the engine to the event bus. + * + * The hosts modules publish onto the bus and know nothing about automations; + * this is the other half of that arrangement, and it is what keeps the edge + * one-directional (automations may import repositories, hosts modules must not + * import automations). + * + * Handlers are async and the bus is fire-and-forget, so a rejection here is + * logged by the bus rather than surfacing to whichever host feature emitted. + */ +export function subscribeAutomationTriggers(): void { + if (unsubscribers.length > 0) return; + + unsubscribers = [ + pluginEvents.on(TOPICS.hostMetrics, (payload) => + onMetrics(payload as MetricEvent), + ), + pluginEvents.on(TOPICS.hostStatus, (payload) => + onStatus(payload as StatusEvent), + ), + pluginEvents.on(TOPICS.hostHealthCheck, (payload) => + onHealthCheck(payload as HealthEvent), + ), + pluginEvents.on(TOPICS.internalEvent, (payload) => + onInternalEvent(payload as InternalEvent), + ), + ]; +} + +export function unsubscribeAutomationTriggers(): void { + for (const unsubscribe of unsubscribers) unsubscribe(); + unsubscribers = []; +} + /** * Hosts that an enabled automation watches, so the poller knows to collect * metrics for them even when nobody is looking. diff --git a/src/backend/hosts/automation-events.ts b/src/backend/hosts/automation-events.ts index 662b5c612..d4935e36b 100644 --- a/src/backend/hosts/automation-events.ts +++ b/src/backend/hosts/automation-events.ts @@ -1,17 +1,24 @@ /** - * One-line hand-off from any hosts feature to the automations engine for a - * generic named event (not tied to metrics polling). + * One-line hand-off from any hosts feature to anything listening for a generic + * named event (not tied to metrics polling). * - * Fire-and-forget and imported lazily: a failure in the automations layer - * must never disturb the caller, and a static import would create a cycle - * (automations reads repositories, which several hosts modules also pull in). - * - * PLUGIN-EVENT: this whole module is the "publish an internal event" half of - * the phase-2 ctx.events bus. Once that exists, callers should emit onto - * ctx.events instead of calling this directly, and the automations engine - * subscribes there rather than being imported ad hoc. + * This now publishes onto the ctx.events bus rather than reaching into the + * automations engine directly. The two properties callers relied on are + * unchanged: it is fire-and-forget, and it creates no static edge to the + * automations layer (which reads repositories that several hosts modules also + * pull in, so a static import would close a cycle). The automations engine + * subscribes to the bus at start-up instead -- see automations/triggers.ts. */ +import { pluginEvents, TOPICS } from "../plugins/events.js"; + +export interface InternalEventPayload { + event: string; + userId: string; + hostId?: number; + details?: Record; +} + export function notifyAutomationInternalEvent( event: string, userId: string, @@ -19,9 +26,10 @@ export function notifyAutomationInternalEvent( details?: Record, ): void { if (!userId) return; - import("../automations/triggers.js") - .then((triggers) => - triggers.onInternalEvent({ event, userId, hostId, details }), - ) - .catch(() => {}); + pluginEvents.emit(TOPICS.internalEvent, { + event, + userId, + hostId, + details, + } satisfies InternalEventPayload); } diff --git a/src/backend/hosts/host-session-status.ts b/src/backend/hosts/host-session-status.ts index 06a90d2a6..b9406c584 100644 --- a/src/backend/hosts/host-session-status.ts +++ b/src/backend/hosts/host-session-status.ts @@ -1,5 +1,12 @@ +import { pluginEvents, TOPICS } from "../plugins/events.js"; + type HostSessionStatusListener = (hostId: number, online: boolean) => void; +export interface HostSessionStatusPayload { + hostId: number; + online: boolean; +} + export class HostSessionStatus { private counts = new Map(); private listeners = new Set(); @@ -32,6 +39,12 @@ export class HostSessionStatus { private emit(hostId: number, online: boolean): void { for (const listener of this.listeners) listener(hostId, online); + // Also published as a topic so plugins (and anything else on the bus) can + // observe session status without holding a reference to this singleton. + pluginEvents.emit(TOPICS.hostSessionStatus, { + hostId, + online, + } satisfies HostSessionStatusPayload); } } diff --git a/src/backend/hosts/metrics/automation-bridge.ts b/src/backend/hosts/metrics/automation-bridge.ts index a3c6fee89..b559c7af1 100644 --- a/src/backend/hosts/metrics/automation-bridge.ts +++ b/src/backend/hosts/metrics/automation-bridge.ts @@ -1,27 +1,50 @@ /** - * One-line hand-off from the metrics poller to the automations engine. + * One-line hand-off from the metrics poller onto the ctx.events bus. * * The poller is already a very large module, so the hooks it calls live here - * instead. Everything is fire-and-forget and imported lazily: a failure in the - * automations layer must never disturb metric collection, and a static import - * would create a cycle (automations reads repositories, which the metrics - * module also pulls in). + * instead. Everything is fire-and-forget: a failure in a subscriber must never + * disturb metric collection, and publishing onto the bus creates no static edge + * to the automations layer (which reads repositories the metrics module also + * pulls in). * * For the generic (non-metrics) internal event notifier, see * `hosts/automation-events.ts` -- that one is shared across features. */ +import { pluginEvents, TOPICS } from "../../plugins/events.js"; import type { MetricsSnapshot } from "../../automations/conditions.js"; +export interface HostMetricsPayload { + hostId: number; + ownerUserId: string; + metrics: MetricsSnapshot; +} + +export interface HostStatusPayload { + hostId: number; + ownerUserId: string; + online: boolean; +} + +export interface HostHealthCheckPayload { + hostId: number; + userId: string; + checkId: string; + ok: boolean; + detail?: string; +} + export function notifyAutomationMetrics( hostId: number, ownerUserId: string, metrics: MetricsSnapshot, ): void { if (!ownerUserId) return; - import("../../automations/triggers.js") - .then((triggers) => triggers.onMetrics({ hostId, ownerUserId, metrics })) - .catch(() => {}); + pluginEvents.emit(TOPICS.hostMetrics, { + hostId, + ownerUserId, + metrics, + } satisfies HostMetricsPayload); } export function notifyAutomationStatus( @@ -30,9 +53,11 @@ export function notifyAutomationStatus( online: boolean, ): void { if (!ownerUserId) return; - import("../../automations/triggers.js") - .then((triggers) => triggers.onStatus({ hostId, ownerUserId, online })) - .catch(() => {}); + pluginEvents.emit(TOPICS.hostStatus, { + hostId, + ownerUserId, + online, + } satisfies HostStatusPayload); } export function notifyAutomationHealthCheck( @@ -43,9 +68,11 @@ export function notifyAutomationHealthCheck( detail?: string, ): void { if (!userId) return; - import("../../automations/triggers.js") - .then((triggers) => - triggers.onHealthCheck({ hostId, userId, checkId, ok, detail }), - ) - .catch(() => {}); + pluginEvents.emit(TOPICS.hostHealthCheck, { + hostId, + userId, + checkId, + ok, + detail, + } satisfies HostHealthCheckPayload); } diff --git a/src/backend/hosts/metrics/index.ts b/src/backend/hosts/metrics/index.ts index 50bfa28ec..23ec9681a 100644 --- a/src/backend/hosts/metrics/index.ts +++ b/src/backend/hosts/metrics/index.ts @@ -61,8 +61,8 @@ import { registerHostMetricsHistoryRoutes } from "./history-routes.js"; import { registerProxmoxStatsRoutes } from "./proxmox-stats-routes.js"; import { registerProxmoxStatsHistoryRoutes } from "./proxmox-stats-history-routes.js"; import { ProxmoxPollingManager } from "./proxmox-stats-polling.js"; -// PLUGIN-EVENT: terminal session online/offline -> phase-2 ctx.events "host.session.status" topic -import { hostSessionStatus } from "../host-session-status.js"; +import type { HostSessionStatusPayload } from "../host-session-status.js"; +import { pluginEvents, TOPICS } from "../../plugins/events.js"; import { AlertEngine } from "./alert-engine.js"; import { notifyAutomationMetrics, @@ -225,10 +225,12 @@ class PollingManager { private unsubscribeHostSessionStatus: () => void; constructor() { - // PLUGIN-EVENT: subscribe to ctx.events.on("host.session.status", ...) once the - // phase-2 event bus exists, instead of the hostSessionStatus singleton directly. - this.unsubscribeHostSessionStatus = hostSessionStatus.subscribe( - (hostId, online) => this.setTerminalSessionOnline(hostId, online), + this.unsubscribeHostSessionStatus = pluginEvents.on( + TOPICS.hostSessionStatus, + (payload) => { + const { hostId, online } = payload as HostSessionStatusPayload; + this.setTerminalSessionOnline(hostId, online); + }, ); this.viewerCleanupInterval = setInterval(() => { this.cleanupInactiveViewers(); diff --git a/src/backend/plugins/events.ts b/src/backend/plugins/events.ts new file mode 100644 index 000000000..281c6c30f --- /dev/null +++ b/src/backend/plugins/events.ts @@ -0,0 +1,99 @@ +/** + * The process-wide event bus behind ctx.events. + * + * This formalises three ad-hoc publish paths that already existed: + * - hosts/automation-events.ts (generic named internal events) + * - hosts/metrics/automation-bridge.ts (metrics-shaped siblings) + * - hosts/host-session-status.ts (the "host.session.status" topic) + * + * Each of those lazily imported the automations engine and swallowed errors, + * for two reasons that still apply and are preserved here: + * + * 1. Fire-and-forget. A failing subscriber must never disturb the caller. A + * metrics poll or a host delete does not fail because an automation threw. + * 2. No static import of the automations layer. It reads repositories, which + * several hosts modules also pull in, so a static edge would close a + * cycle. Subscribers register themselves here instead of being reached + * into. + * + * Topics are dotted strings. The internal ones the server itself publishes are + * listed in TOPICS; a plugin may emit and subscribe to any topic, but only + * topics it is allowed to see are pushed to it. + */ + +import { pluginLogger } from "../utils/logger.js"; + +export const TOPICS = { + /** A terminal session opened or closed against a host. */ + hostSessionStatus: "host.session.status", + /** A generic named internal event (host_deleted, user_login, ...). */ + internalEvent: "internal.event", + /** A metrics snapshot for one host. */ + hostMetrics: "host.metrics", + /** A host went online or offline as seen by the metrics poller. */ + hostStatus: "host.status", + /** A health check result. */ + hostHealthCheck: "host.health_check", +} as const; + +export type EventListener = (payload: unknown) => void; + +class PluginEventBus { + private readonly listeners = new Map>(); + + on(topic: string, listener: EventListener): () => void { + let set = this.listeners.get(topic); + if (!set) { + set = new Set(); + this.listeners.set(topic, set); + } + set.add(listener); + + return () => { + set!.delete(listener); + if (set!.size === 0) this.listeners.delete(topic); + }; + } + + /** + * Fire-and-forget by contract. Never throws, never returns a promise the + * caller is expected to await, and one bad subscriber cannot stop the others. + */ + emit(topic: string, payload: unknown): void { + const set = this.listeners.get(topic); + if (!set || set.size === 0) return; + + for (const listener of [...set]) { + try { + const result = listener(payload) as unknown; + // A listener may be async; its rejection must not become unhandled. + if (result && typeof (result as Promise).catch === "function") { + void (result as Promise).catch((error) => + this.report(topic, error), + ); + } + } catch (error) { + this.report(topic, error); + } + } + } + + listenerCount(topic: string): number { + return this.listeners.get(topic)?.size ?? 0; + } + + /** Test seam. */ + clear(): void { + this.listeners.clear(); + } + + private report(topic: string, error: unknown): void { + pluginLogger.error( + `Event listener for "${topic}" failed`, + error instanceof Error ? error : new Error(String(error)), + { operation: "plugin_events" }, + ); + } +} + +export const pluginEvents = new PluginEventBus(); diff --git a/src/backend/tests/plugins/events.test.ts b/src/backend/tests/plugins/events.test.ts new file mode 100644 index 000000000..e5ebc5375 --- /dev/null +++ b/src/backend/tests/plugins/events.test.ts @@ -0,0 +1,185 @@ +import { afterEach, describe, expect, it, vi } from "vitest"; + +vi.mock("../../utils/logger.js", () => ({ + pluginLogger: { + debug: vi.fn(), + info: vi.fn(), + warn: vi.fn(), + error: vi.fn(), + success: vi.fn(), + }, +})); + +const { pluginEvents, TOPICS } = await import("../../plugins/events.js"); + +describe("plugin event bus", () => { + afterEach(() => { + pluginEvents.clear(); + vi.clearAllMocks(); + }); + + it("delivers a payload to every subscriber of a topic", () => { + const first = vi.fn(); + const second = vi.fn(); + + pluginEvents.on("topic.a", first); + pluginEvents.on("topic.a", second); + pluginEvents.emit("topic.a", { value: 1 }); + + expect(first).toHaveBeenCalledWith({ value: 1 }); + expect(second).toHaveBeenCalledWith({ value: 1 }); + }); + + it("does not deliver across topics", () => { + const listener = vi.fn(); + pluginEvents.on("topic.a", listener); + pluginEvents.emit("topic.b", {}); + + expect(listener).not.toHaveBeenCalled(); + }); + + it("emitting a topic with no subscribers is a no-op", () => { + expect(() => pluginEvents.emit("nobody.listening", {})).not.toThrow(); + }); + + it("unsubscribes and cleans up the empty topic", () => { + const listener = vi.fn(); + const off = pluginEvents.on("topic.a", listener); + + expect(pluginEvents.listenerCount("topic.a")).toBe(1); + off(); + expect(pluginEvents.listenerCount("topic.a")).toBe(0); + + pluginEvents.emit("topic.a", {}); + expect(listener).not.toHaveBeenCalled(); + }); + + it("a throwing subscriber does not stop the others or the caller", async () => { + const after = vi.fn(); + pluginEvents.on("topic.a", () => { + throw new Error("subscriber exploded"); + }); + pluginEvents.on("topic.a", after); + + expect(() => pluginEvents.emit("topic.a", {})).not.toThrow(); + expect(after).toHaveBeenCalled(); + + const { pluginLogger } = await import("../../utils/logger.js"); + expect(pluginLogger.error).toHaveBeenCalled(); + }); + + it("a rejecting async subscriber is reported, not left unhandled", async () => { + pluginEvents.on("topic.a", async () => { + throw new Error("async subscriber exploded"); + }); + + expect(() => pluginEvents.emit("topic.a", {})).not.toThrow(); + await new Promise((resolve) => setImmediate(resolve)); + + const { pluginLogger } = await import("../../utils/logger.js"); + expect(pluginLogger.error).toHaveBeenCalled(); + }); + + it("a subscriber added during an emit does not receive that same emit", () => { + const late = vi.fn(); + pluginEvents.on("topic.a", () => { + pluginEvents.on("topic.a", late); + }); + + pluginEvents.emit("topic.a", {}); + expect(late).not.toHaveBeenCalled(); + + pluginEvents.emit("topic.a", {}); + expect(late).toHaveBeenCalledTimes(1); + }); + + it("exposes the internal topic names the server publishes", () => { + expect(TOPICS.hostSessionStatus).toBe("host.session.status"); + expect(TOPICS.internalEvent).toBe("internal.event"); + expect(TOPICS.hostMetrics).toBe("host.metrics"); + expect(TOPICS.hostStatus).toBe("host.status"); + expect(TOPICS.hostHealthCheck).toBe("host.health_check"); + }); +}); + +describe("hosts publishers route through the bus", () => { + // No vi.resetModules() here: the bus is a module singleton, and resetting + // would hand these imports a different instance than the one subscribed to. + afterEach(() => { + pluginEvents.clear(); + }); + + it("notifyAutomationInternalEvent emits internal.event", async () => { + const { notifyAutomationInternalEvent } = + await import("../../hosts/automation-events.js"); + const listener = vi.fn(); + pluginEvents.on(TOPICS.internalEvent, listener); + + notifyAutomationInternalEvent("host_deleted", "user-1", 42, { a: 1 }); + + expect(listener).toHaveBeenCalledWith({ + event: "host_deleted", + userId: "user-1", + hostId: 42, + details: { a: 1 }, + }); + }); + + it("notifyAutomationInternalEvent ignores a missing userId", async () => { + const { notifyAutomationInternalEvent } = + await import("../../hosts/automation-events.js"); + const listener = vi.fn(); + pluginEvents.on(TOPICS.internalEvent, listener); + + notifyAutomationInternalEvent("host_deleted", ""); + expect(listener).not.toHaveBeenCalled(); + }); + + it("notifyAutomationStatus and notifyAutomationMetrics emit their topics", async () => { + const { notifyAutomationStatus, notifyAutomationMetrics } = + await import("../../hosts/metrics/automation-bridge.js"); + + const status = vi.fn(); + const metrics = vi.fn(); + pluginEvents.on(TOPICS.hostStatus, status); + pluginEvents.on(TOPICS.hostMetrics, metrics); + + notifyAutomationStatus(42, "user-1", true); + notifyAutomationMetrics(42, "user-1", { cpu: 10 } as never); + + expect(status).toHaveBeenCalledWith({ + hostId: 42, + ownerUserId: "user-1", + online: true, + }); + expect(metrics).toHaveBeenCalledWith({ + hostId: 42, + ownerUserId: "user-1", + metrics: { cpu: 10 }, + }); + }); + + it("hostSessionStatus publishes host.session.status on the refcount edges", async () => { + const { HostSessionStatus } = + await import("../../hosts/host-session-status.js"); + const seen: unknown[] = []; + pluginEvents.on(TOPICS.hostSessionStatus, (p) => seen.push(p)); + + const status = new HostSessionStatus(); + const releaseFirst = status.register(7); + const releaseSecond = status.register(7); + + // Only the 0 -> 1 edge publishes. + expect(seen).toEqual([{ hostId: 7, online: true }]); + + releaseFirst(); + expect(seen).toHaveLength(1); + + // ... and only the 1 -> 0 edge. + releaseSecond(); + expect(seen).toEqual([ + { hostId: 7, online: true }, + { hostId: 7, online: false }, + ]); + }); +});