feat: route internal events through a shared plugin event bus

Formalises three ad-hoc publish paths that already existed:
automation-events.ts, metrics/automation-bridge.ts and the
hostSessionStatus singleton, now the host.session.status topic.

Both invariants the old code relied on are preserved. Publishing stays
fire-and-forget, so a failing subscriber cannot disturb a metrics poll or
a host delete. The automations engine subscribes to the bus at scheduler
start instead of being imported ad hoc, which keeps the edge
one-directional: automations may import repositories, hosts modules must
not import automations.
This commit is contained in:
LukeGus committed 2026-09-21 07:19:01 -05:00
1 parent 0a771c9600
commit 1393fa7581
9 files changed
+433 -52

No files matched your search

+15 -15
View File
@@ -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,
},
);
}
}
+9
View File
@@ -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. */
+38
View File
@@ -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<void> {
}
}
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.
+23 -15
View File
@@ -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<string, unknown>;
}
export function notifyAutomationInternalEvent(
event: string,
userId: string,
@@ -19,9 +26,10 @@ export function notifyAutomationInternalEvent(
details?: Record<string, unknown>,
): 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);
}
+13
View File
@@ -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<number, number>();
private listeners = new Set<HostSessionStatusListener>();
@@ -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);
}
}
+43 -16
View File
@@ -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);
}
+8 -6
View File
@@ -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();
+99
View File
@@ -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<string, Set<EventListener>>();
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<void>).catch === "function") {
void (result as Promise<void>).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();
+185
View File
@@ -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 },
]);
});
});