,
+): BuildEventListenerReturn {
+ // Listeners are matched by `event` at dispatch time; the cast erases the
+ // per-event generic so listeners for different events can share one array.
+ return args as unknown as BuildEventListenerReturn;
+}
diff --git a/packages/vitnode/src/api/lib/module.ts b/packages/vitnode/src/api/lib/module.ts
index 3eeb34cbd..1c8bcbf61 100644
--- a/packages/vitnode/src/api/lib/module.ts
+++ b/packages/vitnode/src/api/lib/module.ts
@@ -1,6 +1,7 @@
import { OpenAPIHono } from "@hono/zod-openapi";
import type { BuildCronReturn } from "./cron";
+import type { BuildEventListenerReturn } from "./events";
import type { BuildQueueTaskReturn } from "./queue";
import type { Route } from "./route";
import type { BuildWebSocketReturn } from "./websocket";
@@ -16,6 +17,7 @@ export interface BaseBuildModuleReturn<
Routes extends Route[] = Route
[],
> {
cronJobs: BuildCronReturn[];
+ events: BuildEventListenerReturn[];
hono: OpenAPIHono;
modules?: BaseBuildModuleReturn
[];
name: M;
@@ -45,10 +47,12 @@ export function buildModule<
name,
modules,
cronJobs = [],
+ events = [],
queueTasks = [],
webSockets = [],
}: {
cronJobs?: BuildCronReturn[];
+ events?: BuildEventListenerReturn[];
modules?: Modules;
name: M;
pluginId: P;
@@ -77,6 +81,7 @@ export function buildModule<
name,
modules,
cronJobs,
+ events,
queueTasks,
webSockets,
};
diff --git a/packages/vitnode/src/api/lib/plugin.ts b/packages/vitnode/src/api/lib/plugin.ts
index 33a43274d..a6effde29 100644
--- a/packages/vitnode/src/api/lib/plugin.ts
+++ b/packages/vitnode/src/api/lib/plugin.ts
@@ -2,6 +2,7 @@ import { OpenAPIHono } from "@hono/zod-openapi";
import type { SearchIndexer } from "../models/search";
import type { CronJobConfig } from "./cron";
+import type { EventListenerConfig } from "./events";
import type { BuildModuleReturn } from "./module";
import type { PermissionStaffConfig } from "./permission-staff";
import type { QueueTaskConfig } from "./queue";
@@ -11,6 +12,7 @@ import { checkPluginId } from "./check-plugin-id";
export interface BuildPluginApiReturn {
cronJobs?: Omit[];
+ events?: Omit[];
hono: OpenAPIHono;
permissionStaff?: PermissionStaffConfig;
pluginId: string;
@@ -35,6 +37,7 @@ export function buildApiPlugin({
const hono = new OpenAPIHono();
const cronJobs: BuildPluginApiReturn["cronJobs"] = [];
+ const events: BuildPluginApiReturn["events"] = [];
const queueTasks: BuildPluginApiReturn["queueTasks"] = [];
const webSockets: BuildPluginApiReturn["webSockets"] = [];
modules.forEach(handler => {
@@ -44,6 +47,10 @@ export function buildApiPlugin
({
cronJobs.push({ ...cron, module: handler.name });
});
+ handler.events?.forEach(listener => {
+ events.push({ ...listener, module: handler.name });
+ });
+
handler.queueTasks?.forEach(task => {
queueTasks.push({ ...task, module: handler.name });
});
@@ -57,6 +64,7 @@ export function buildApiPlugin
({
pluginId,
hono,
cronJobs,
+ events,
queueTasks,
searchIndexers,
webSockets,
diff --git a/packages/vitnode/src/api/middlewares/global.middleware.ts b/packages/vitnode/src/api/middlewares/global.middleware.ts
index 26c662e0d..ea68e6df0 100644
--- a/packages/vitnode/src/api/middlewares/global.middleware.ts
+++ b/packages/vitnode/src/api/middlewares/global.middleware.ts
@@ -6,9 +6,11 @@ import { HTTPException } from "hono/http-exception";
import type { VitNodeApiConfig, VitNodeConfig } from "@/vitnode.config";
import type { VitNodeRealtime } from "@/ws/registry";
+import { LocalEventsAdapter } from "@/api/adapters/events/local";
import { PostgresSearchAdapter } from "@/api/adapters/search/postgres";
import { CacheModel } from "@/api/lib/cache";
import { EmailModel } from "@/api/models/email";
+import { EventsModel } from "@/api/models/events";
import { QueueModel } from "@/api/models/queue";
import { SearchModel } from "@/api/models/search";
import { SessionModel } from "@/api/models/session";
@@ -18,9 +20,11 @@ import { CONFIG } from "@/lib/config";
import { realtime } from "@/ws/registry";
import type { BuildCronReturn } from "../lib/cron";
+import type { EventListenerConfig } from "../lib/events";
import type { PermissionStaffCatalogEntry } from "../lib/permission-staff";
import type { BuildQueueTaskReturn } from "../lib/queue";
import type { WebSocketConfig } from "../lib/websocket";
+import type { EventsApiPlugin } from "../models/events";
import type {
SearchIndexerConfig,
SearchProviderApiPlugin,
@@ -74,6 +78,7 @@ export interface EnvVariablesVitNode {
cron: (BuildCronReturn & { module: string; pluginId: string })[];
cronSecret?: string;
email?: VitNodeApiConfig["email"];
+ events: { adapter: EventsApiPlugin; listeners: EventListenerConfig[] };
// Whether a cron adapter is configured (`buildApiConfig({ cron })`), i.e. an
// in-process scheduler is running the registered jobs automatically. Without
// it, jobs only run when the cron endpoint is triggered externally.
@@ -93,6 +98,7 @@ export interface EnvVariablesVitNode {
};
db: Pick["dbProvider"];
email: EmailModel;
+ events: EventsModel;
ipAddress: string;
log: LoggerMiddlewareType;
plugin: {
@@ -123,6 +129,7 @@ export const globalMiddleware = ({
dbProvider,
captcha,
cron,
+ events,
plugins,
pathToMessages,
search,
@@ -135,6 +142,7 @@ export const globalMiddleware = ({
| "cron"
| "dbProvider"
| "email"
+ | "events"
| "pathToMessages"
| "plugins"
| "search"
@@ -147,6 +155,13 @@ export const globalMiddleware = ({
const cronMetadata = collectCronJobs(plugins);
+ const eventsMetadata: EventListenerConfig[] = plugins.flatMap(plugin =>
+ (plugin.events ?? []).map(listener => ({
+ ...listener,
+ pluginId: plugin.pluginId,
+ })),
+ );
+
const queueMetadata = plugins.flatMap(plugin =>
(plugin.queueTasks ?? []).map(task => ({
pluginId: plugin.pluginId,
@@ -219,6 +234,7 @@ export const globalMiddleware = ({
c.set("db", dbProvider);
c.set("cache", new CacheModel(cacheClient, c));
c.set("email", new EmailModel(c));
+ c.set("events", new EventsModel(c));
c.set("queue", new QueueModel(c));
c.set("search", new SearchModel(c));
c.set("storage", new StorageModel(c));
@@ -228,6 +244,10 @@ export const globalMiddleware = ({
pathToMessages,
metadata,
email,
+ events: {
+ adapter: events?.adapter ?? LocalEventsAdapter(),
+ listeners: eventsMetadata,
+ },
search: { adapter: search?.adapter ?? PostgresSearchAdapter() },
searchIndexers: searchIndexersMetadata,
storage,
diff --git a/packages/vitnode/src/api/models/events.test.ts b/packages/vitnode/src/api/models/events.test.ts
new file mode 100644
index 000000000..f9b834a32
--- /dev/null
+++ b/packages/vitnode/src/api/models/events.test.ts
@@ -0,0 +1,263 @@
+// @vitest-environment node
+import type { Context } from "hono";
+
+import { describe, expect, it, vi } from "vitest";
+
+import type { EventListenerConfig } from "../lib/events";
+import type { EventsApiPlugin } from "./events";
+
+import { LocalEventsAdapter } from "../adapters/events/local";
+import { EventsModel } from "./events";
+
+const makeCtx = (
+ overrides: {
+ adapter?: EventsApiPlugin;
+ admin?: { user: { id: number } };
+ listeners?: EventListenerConfig[];
+ plugin?: { id: string };
+ user?: { id: number };
+ } = {},
+): {
+ ctx: Context;
+ logError: ReturnType;
+} => {
+ const logError = vi.fn().mockResolvedValue(undefined);
+ const store: Record = {
+ admin: overrides.admin ?? null,
+ core: {
+ events: {
+ adapter: overrides.adapter ?? LocalEventsAdapter(),
+ listeners: overrides.listeners ?? [],
+ },
+ },
+ log: { error: logError },
+ plugin: overrides.plugin,
+ user: overrides.user ?? null,
+ };
+
+ return {
+ ctx: { get: (k: string) => store[k] } as unknown as Context,
+ logError,
+ };
+};
+
+const makeListener = (
+ overrides: Partial = {},
+): EventListenerConfig => ({
+ event: "user.created",
+ name: "listener",
+ module: "users",
+ pluginId: "@vitnode/core",
+ handler: vi.fn().mockResolvedValue(undefined),
+ ...overrides,
+});
+
+const PAYLOAD = {
+ userId: 1,
+ email: "a@b.com",
+ name: "Test",
+ emailVerified: true,
+};
+
+describe("EventsModel.emit (Local adapter)", () => {
+ it("runs only listeners matching the event, sequentially in registry order", async () => {
+ const order: string[] = [];
+ const first = makeListener({
+ name: "first",
+ handler: async () => {
+ // Resolve on a later tick so a concurrent dispatch would flip the order.
+ await new Promise(resolve => setTimeout(resolve, 5));
+ order.push("first");
+ },
+ });
+ const second = makeListener({
+ name: "second",
+ handler: () => {
+ order.push("second");
+ },
+ });
+ const other = makeListener({
+ name: "other",
+ event: "other.event" as EventListenerConfig["event"],
+ handler: () => {
+ order.push("other");
+ },
+ });
+ const { ctx } = makeCtx({ listeners: [first, other, second] });
+
+ const result = await new EventsModel(ctx).emit("user.created", PAYLOAD);
+
+ expect(order).toEqual(["first", "second"]);
+ expect(result.delivered).toBe(2);
+ expect(result.failures).toEqual([]);
+ });
+
+ it("passes payload and envelope to the handler", async () => {
+ const handler = vi.fn().mockResolvedValue(undefined);
+ const { ctx } = makeCtx({ listeners: [makeListener({ handler })] });
+
+ await new EventsModel(ctx).emit("user.created", PAYLOAD);
+
+ expect(handler).toHaveBeenCalledWith(
+ ctx,
+ PAYLOAD,
+ expect.objectContaining({ name: "user.created", payload: PAYLOAD }),
+ );
+ });
+
+ it("a throwing listener does not stop the remaining listeners", async () => {
+ const ran: string[] = [];
+ const failing = makeListener({
+ name: "failing",
+ module: "posts",
+ pluginId: "@vitnode/blog",
+ handler: () => {
+ throw new Error("boom");
+ },
+ });
+ const after = makeListener({
+ name: "after",
+ handler: () => {
+ ran.push("after");
+ },
+ });
+ const { ctx, logError } = makeCtx({ listeners: [failing, after] });
+
+ const result = await new EventsModel(ctx).emit("user.created", PAYLOAD);
+
+ expect(ran).toEqual(["after"]);
+ expect(result.delivered).toBe(1);
+ expect(result.failures).toEqual([
+ {
+ pluginId: "@vitnode/blog",
+ module: "posts",
+ listener: "failing",
+ error: "boom",
+ },
+ ]);
+ expect(logError).toHaveBeenCalledTimes(1);
+ expect(logError.mock.calls[0][0]).toContain(
+ '"@vitnode/blog:posts:failing"',
+ );
+ });
+
+ it("resolves even when every listener throws", async () => {
+ const listeners = [
+ makeListener({
+ name: "a",
+ handler: () => {
+ throw new Error("a failed");
+ },
+ }),
+ makeListener({
+ name: "b",
+ handler: async () => Promise.reject(new Error("b failed")),
+ }),
+ ];
+ const { ctx } = makeCtx({ listeners });
+
+ const result = await new EventsModel(ctx).emit("user.created", PAYLOAD);
+
+ expect(result.delivered).toBe(0);
+ expect(result.failures).toHaveLength(2);
+ expect(result.status).toBe("delivered");
+ });
+
+ it("returns an empty result when no listeners match", async () => {
+ const { ctx } = makeCtx();
+
+ const result = await new EventsModel(ctx).emit("user.created", PAYLOAD);
+
+ expect(result).toMatchObject({
+ delivered: 0,
+ failures: [],
+ status: "delivered",
+ });
+ expect(result.eventId).toMatch(/^[0-9a-f-]{36}$/);
+ });
+});
+
+describe("EventsModel.emit envelope", () => {
+ const captureEnvelope = () => {
+ const publish = vi.fn().mockResolvedValue({
+ eventId: "x",
+ status: "delivered",
+ delivered: 0,
+ failures: [],
+ });
+
+ return { adapter: { name: "capture", publish }, publish };
+ };
+
+ it("stamps eventId, emittedAt and defaults pluginId to @vitnode/core", async () => {
+ const { adapter, publish } = captureEnvelope();
+ const { ctx } = makeCtx({ adapter });
+
+ await new EventsModel(ctx).emit("user.created", PAYLOAD);
+
+ const envelope = publish.mock.calls[0][1];
+ expect(envelope.eventId).toMatch(/^[0-9a-f-]{36}$/);
+ expect(envelope.emittedAt).toBeInstanceOf(Date);
+ expect(envelope.pluginId).toBe("@vitnode/core");
+ });
+
+ it("uses the emitting plugin id from context", async () => {
+ const { adapter, publish } = captureEnvelope();
+ const { ctx } = makeCtx({ adapter, plugin: { id: "@vitnode/blog" } });
+
+ await new EventsModel(ctx).emit("user.created", PAYLOAD);
+
+ expect(publish.mock.calls[0][1].pluginId).toBe("@vitnode/blog");
+ });
+
+ it("derives the actor: admin wins over user, then user, then system", async () => {
+ const { adapter, publish } = captureEnvelope();
+
+ const { ctx: adminCtx } = makeCtx({
+ adapter,
+ admin: { user: { id: 7 } },
+ user: { id: 3 },
+ });
+ await new EventsModel(adminCtx).emit("user.created", PAYLOAD);
+ expect(publish.mock.calls[0][1].actor).toEqual({ type: "admin", id: 7 });
+
+ const { ctx: userCtx } = makeCtx({ adapter, user: { id: 3 } });
+ await new EventsModel(userCtx).emit("user.created", PAYLOAD);
+ expect(publish.mock.calls[1][1].actor).toEqual({ type: "user", id: 3 });
+
+ const { ctx: systemCtx } = makeCtx({ adapter });
+ await new EventsModel(systemCtx).emit("user.created", PAYLOAD);
+ expect(publish.mock.calls[2][1].actor).toEqual({ type: "system" });
+ });
+});
+
+describe("EventsModel adapter seam", () => {
+ it("delegates to the configured adapter and passes its result through", async () => {
+ const publish = vi.fn().mockResolvedValue({
+ eventId: "broker-id",
+ status: "queued",
+ delivered: 0,
+ failures: [],
+ });
+ const { ctx } = makeCtx({ adapter: { name: "broker", publish } });
+
+ const result = await new EventsModel(ctx).emit("user.created", PAYLOAD);
+
+ expect(publish).toHaveBeenCalledTimes(1);
+ expect(result.status).toBe("queued");
+ });
+
+ it("never throws when the adapter itself fails; logs and reports it", async () => {
+ const publish = vi.fn().mockRejectedValue(new Error("broker down"));
+ const { ctx, logError } = makeCtx({
+ adapter: { name: "broker", publish },
+ });
+
+ const result = await new EventsModel(ctx).emit("user.created", PAYLOAD);
+
+ expect(result.delivered).toBe(0);
+ expect(result.failures).toHaveLength(1);
+ expect(result.failures[0].error).toBe("broker down");
+ expect(logError).toHaveBeenCalledTimes(1);
+ });
+});
diff --git a/packages/vitnode/src/api/models/events.ts b/packages/vitnode/src/api/models/events.ts
new file mode 100644
index 000000000..39eb12873
--- /dev/null
+++ b/packages/vitnode/src/api/models/events.ts
@@ -0,0 +1,176 @@
+import type { Context } from "hono";
+
+import { randomUUID } from "node:crypto";
+
+import type { EventListenerConfig } from "../lib/events";
+
+/**
+ * Global map of domain events emittable via `c.get("events").emit(...)`. Core
+ * events are declared here; plugins extend the map with module augmentation:
+ *
+ * ```ts
+ * declare module "@vitnode/core/api/models/events" {
+ * interface VitNodeEvents {
+ * "blog.post.created": { categoryId: number; postId: number };
+ * }
+ * }
+ * ```
+ *
+ * Payloads must stay JSON-serializable - a broker adapter (Redis Streams,
+ * NATS, ...) serializes the envelope to move it between processes.
+ */
+export interface VitNodeEvents {
+ "role.created": {
+ roleId: number;
+ };
+ /**
+ * Declared for plugins implementing role deletion - core has no role
+ * deletion flow yet and never emits this itself.
+ */
+ "role.deleted": {
+ roleId: number;
+ };
+ "role.updated": {
+ roleId: number;
+ };
+ "user.created": {
+ email: string;
+ emailVerified: boolean;
+ name: string;
+ userId: number;
+ };
+ /**
+ * Declared for plugins implementing account deletion - core has no user
+ * deletion flow yet and never emits this itself.
+ */
+ "user.deleted": {
+ email: string;
+ userId: number;
+ };
+ "user.updated": {
+ email: string;
+ name: string;
+ userId: number;
+ };
+}
+
+export type VitNodeEventName = keyof VitNodeEvents;
+
+export interface EventActor {
+ id?: number;
+ type: "admin" | "system" | "user";
+}
+
+export interface EventEnvelope {
+ actor: EventActor;
+ emittedAt: Date;
+ eventId: string;
+ name: K;
+ payload: VitNodeEvents[K];
+ /** Plugin that emitted the event. */
+ pluginId: string;
+}
+
+export interface EventEmitFailure {
+ error: string;
+ /** Listener `name` as declared in `buildEventListener`. */
+ listener: string;
+ module: string;
+ /** Plugin that owns the failing listener. */
+ pluginId: string;
+}
+
+export interface EventEmitResult {
+ /** Listeners that ran successfully before `emit()` resolved. */
+ delivered: number;
+ eventId: string;
+ failures: EventEmitFailure[];
+ /**
+ * `delivered` - listeners ran in-process before `emit()` resolved (the
+ * bundled Local adapter). `queued` - the envelope was handed to a broker and
+ * delivery happens out-of-band; `delivered`/`failures` say nothing about the
+ * eventual listener runs.
+ */
+ status: "delivered" | "queued";
+}
+
+/**
+ * A pluggable event transport. The bundled Local adapter dispatches directly
+ * to the listeners registered in `c.get("core").events.listeners`; a broker
+ * adapter publishes the envelope and returns `status: "queued"`.
+ */
+export interface EventsApiPlugin {
+ name: string;
+ publish: (c: Context, envelope: EventEnvelope) => Promise;
+}
+
+export class EventsModel {
+ constructor(c: Context) {
+ this.c = c;
+ }
+
+ protected readonly c: Context;
+
+ private adapter(): EventsApiPlugin {
+ return this.c.get("core").events.adapter;
+ }
+
+ /**
+ * Emit a typed domain event. Never throws: listener failures are caught,
+ * logged to `core_logs`, and reported in the returned result. Emit only
+ * AFTER the writes the event describes have committed - after your awaited
+ * inserts/updates, and after any enclosing `db.transaction` callback has
+ * returned.
+ */
+ async emit(
+ name: K,
+ payload: VitNodeEvents[K],
+ ): Promise {
+ const admin = this.c.get("admin");
+ const user = this.c.get("user");
+ const envelope: EventEnvelope = {
+ eventId: randomUUID(),
+ name,
+ payload,
+ emittedAt: new Date(),
+ pluginId: this.c.get("plugin")?.id ?? "@vitnode/core",
+ actor: admin
+ ? { type: "admin", id: admin.user.id }
+ : user
+ ? { type: "user", id: user.id }
+ : { type: "system" },
+ };
+ const adapter = this.adapter();
+
+ try {
+ return await adapter.publish(this.c, envelope);
+ } catch (err) {
+ const error = err instanceof Error ? err.message : String(err);
+ await this.c
+ .get("log")
+ .error(
+ `Events adapter "${adapter.name}" failed to publish "${name}": ${error}`,
+ );
+
+ return {
+ eventId: envelope.eventId,
+ status: "delivered",
+ delivered: 0,
+ failures: [
+ {
+ pluginId: envelope.pluginId,
+ module: "adapter",
+ listener: adapter.name,
+ error,
+ },
+ ],
+ };
+ }
+ }
+
+ name(): string {
+ return this.adapter().name;
+ }
+}
+
+export type { EventListenerConfig };
diff --git a/packages/vitnode/src/api/models/user/sign-up.ts b/packages/vitnode/src/api/models/user/sign-up.ts
index 36f47f146..82f19e6c6 100644
--- a/packages/vitnode/src/api/models/user/sign-up.ts
+++ b/packages/vitnode/src/api/models/user/sign-up.ts
@@ -130,5 +130,16 @@ export const signUp = async (
// eslint-disable-next-line @typescript-eslint/no-unused-vars
const { password: _, ...user } = data;
+ // The insert above runs on `c.get("db")` (auto-commit), so the row is
+ // committed here - even when a caller (e.g. the SSO callback) wraps this in
+ // `db.transaction`. If this insert ever moves onto a `tx` handle, the emit
+ // must move after that transaction returns.
+ await c.get("events").emit("user.created", {
+ userId: data.id,
+ email: data.email,
+ name: data.name,
+ emailVerified: data.emailVerified,
+ });
+
return user;
};
diff --git a/packages/vitnode/src/api/modules/admin/roles/routes/create.route.ts b/packages/vitnode/src/api/modules/admin/roles/routes/create.route.ts
index ca7b4a77c..3a1dcd468 100644
--- a/packages/vitnode/src/api/modules/admin/roles/routes/create.route.ts
+++ b/packages/vitnode/src/api/modules/admin/roles/routes/create.route.ts
@@ -68,6 +68,8 @@ export const createRoleAdminRoute = buildRoute({
values: name,
});
+ await c.get("events").emit("role.created", { roleId: role.id });
+
return c.json({ id: role.id }, 201);
},
});
diff --git a/packages/vitnode/src/api/modules/admin/roles/routes/update.route.ts b/packages/vitnode/src/api/modules/admin/roles/routes/update.route.ts
index f7293838f..456da7bb5 100644
--- a/packages/vitnode/src/api/modules/admin/roles/routes/update.route.ts
+++ b/packages/vitnode/src/api/modules/admin/roles/routes/update.route.ts
@@ -98,6 +98,8 @@ export const updateRoleAdminRoute = buildRoute({
});
}
+ await c.get("events").emit("role.updated", { roleId });
+
return c.json({ id: roleId }, 200);
},
});
diff --git a/packages/vitnode/src/api/modules/admin/users/routes/update.route.ts b/packages/vitnode/src/api/modules/admin/users/routes/update.route.ts
index f6ce210ac..df6ddc8a3 100644
--- a/packages/vitnode/src/api/modules/admin/users/routes/update.route.ts
+++ b/packages/vitnode/src/api/modules/admin/users/routes/update.route.ts
@@ -238,9 +238,16 @@ export const updateUserAdminRoute = buildRoute({
nameCode: core_users.nameCode,
});
+ await c.get("events").emit("user.updated", {
+ userId: updated.id,
+ email: updated.email,
+ name: updated.name,
+ });
+
return c.json(updated, 200);
}
+ // No column change, but roles may still have been reassigned above.
const [current] = await db
.select({
id: core_users.id,
@@ -252,6 +259,12 @@ export const updateUserAdminRoute = buildRoute({
.where(eq(core_users.id, user.id))
.limit(1);
+ await c.get("events").emit("user.updated", {
+ userId: current.id,
+ email: current.email,
+ name: current.name,
+ });
+
return c.json(current, 200);
},
});
diff --git a/packages/vitnode/src/vitnode.config.ts b/packages/vitnode/src/vitnode.config.ts
index 9bbc02a69..195661ffc 100644
--- a/packages/vitnode/src/vitnode.config.ts
+++ b/packages/vitnode/src/vitnode.config.ts
@@ -7,6 +7,7 @@ import type React from "react";
import type { CronAdapter } from "./api/lib/cron";
import type { BuildPluginApiReturn } from "./api/lib/plugin";
import type { EmailApiPlugin } from "./api/models/email";
+import type { EventsApiPlugin } from "./api/models/events";
import type { SearchProviderApiPlugin } from "./api/models/search";
import type { SSOApiPlugin } from "./api/models/sso";
import type { StorageApiPlugin } from "./api/models/storage";
@@ -61,6 +62,16 @@ export interface VitNodeApiConfig {
logo?: DefaultTemplateEmailProps["templateProps"]["logo"];
tailwindConfig?: DefaultTemplateEmailProps["templateProps"]["tailwindConfig"];
};
+ /**
+ * Transport for domain events emitted via `c.get("events").emit(...)`. Ships
+ * a zero-config Local adapter used when `adapter` is omitted: listeners run
+ * sequentially in the emitting request, on the emitting instance only
+ * (single-process delivery). Swap the adapter to publish events to an
+ * external broker (e.g. Redis Streams, NATS) for cross-instance delivery.
+ */
+ events?: {
+ adapter?: EventsApiPlugin;
+ };
metadata: {
shortTitle?: string;
title: string;
diff --git a/plugins/blog/src/api/lib/events.ts b/plugins/blog/src/api/lib/events.ts
new file mode 100644
index 000000000..b9932c9e8
--- /dev/null
+++ b/plugins/blog/src/api/lib/events.ts
@@ -0,0 +1,40 @@
+import { buildEventListener } from "@vitnode/core/api/lib/events";
+
+declare module "@vitnode/core/api/models/events" {
+ interface VitNodeEvents {
+ "blog.category.created": {
+ categoryId: number;
+ };
+ "blog.category.deleted": {
+ categoryId: number;
+ postIds: number[];
+ };
+ "blog.category.updated": {
+ categoryId: number;
+ };
+ "blog.post.created": {
+ categoryId: number;
+ postId: number;
+ };
+ "blog.post.deleted": {
+ categoryId: number;
+ postId: number;
+ };
+ "blog.post.updated": {
+ categoryId: number;
+ postId: number;
+ };
+ }
+}
+
+export const cleanupCategorySearchListener = buildEventListener({
+ event: "blog.category.deleted",
+ name: "cleanup-category-search",
+ description:
+ "Remove search index rows of posts cascade-deleted with a category",
+ handler: async (c, payload) => {
+ for (const postId of payload.postIds) {
+ await c.get("search").delete("blog_post", postId);
+ }
+ },
+});
diff --git a/plugins/blog/src/api/modules/admin/admin.module.ts b/plugins/blog/src/api/modules/admin/admin.module.ts
index 86aa8b2de..4113aca19 100644
--- a/plugins/blog/src/api/modules/admin/admin.module.ts
+++ b/plugins/blog/src/api/modules/admin/admin.module.ts
@@ -1,6 +1,7 @@
import { buildModule } from "@vitnode/core/api/lib/module";
import { CONFIG_PLUGIN } from "../../../const";
+import { cleanupCategorySearchListener } from "../../lib/events";
import { categoriesAdminModule } from "./categories/categories.admin.module";
import { postsAdminModule } from "./posts/posts.admin.module";
@@ -9,4 +10,8 @@ export const adminModule = buildModule({
name: "admin",
modules: [categoriesAdminModule, postsAdminModule],
routes: [],
+ // Event listeners are only collected from top-level modules (like cronJobs
+ // and queueTasks), so they are registered here rather than on the nested
+ // categories module.
+ events: [cleanupCategorySearchListener],
});
diff --git a/plugins/blog/src/api/modules/admin/categories/routes/create.route.ts b/plugins/blog/src/api/modules/admin/categories/routes/create.route.ts
index bc0dba2e6..7f452cf0b 100644
--- a/plugins/blog/src/api/modules/admin/categories/routes/create.route.ts
+++ b/plugins/blog/src/api/modules/admin/categories/routes/create.route.ts
@@ -59,6 +59,10 @@ export const createCategoryRoute = buildRoute({
await saveCategoryTranslations(c, category.id, { title });
+ await c.get("events").emit("blog.category.created", {
+ categoryId: category.id,
+ });
+
return c.json(category, 201);
},
});
diff --git a/plugins/blog/src/api/modules/admin/categories/routes/delete.route.ts b/plugins/blog/src/api/modules/admin/categories/routes/delete.route.ts
index 82d5d3c28..8daa567e5 100644
--- a/plugins/blog/src/api/modules/admin/categories/routes/delete.route.ts
+++ b/plugins/blog/src/api/modules/admin/categories/routes/delete.route.ts
@@ -5,6 +5,7 @@ import { HTTPException } from "hono/http-exception";
import { CONFIG_PLUGIN } from "@/const";
import { blog_categories } from "@/database/categories";
+import { blog_posts } from "@/database/posts";
export const deleteCategoryRoute = buildRoute({
pluginId: CONFIG_PLUGIN.pluginId,
@@ -29,6 +30,14 @@ export const deleteCategoryRoute = buildRoute({
handler: async c => {
const { id } = c.req.valid("param");
+ // Capture the posts the category's `onDelete: "cascade"` is about to
+ // remove, so listeners (e.g. search index cleanup) know what went away.
+ const posts = await c
+ .get("db")
+ .select({ id: blog_posts.id })
+ .from(blog_posts)
+ .where(eq(blog_posts.categoryId, id));
+
const result = await c
.get("db")
.delete(blog_categories)
@@ -39,6 +48,11 @@ export const deleteCategoryRoute = buildRoute({
throw new HTTPException(404);
}
+ await c.get("events").emit("blog.category.deleted", {
+ categoryId: id,
+ postIds: posts.map(post => post.id),
+ });
+
return c.body(null, 204);
},
});
diff --git a/plugins/blog/src/api/modules/admin/categories/routes/edit.route.ts b/plugins/blog/src/api/modules/admin/categories/routes/edit.route.ts
index 47cf8a99e..8b1228269 100644
--- a/plugins/blog/src/api/modules/admin/categories/routes/edit.route.ts
+++ b/plugins/blog/src/api/modules/admin/categories/routes/edit.route.ts
@@ -75,6 +75,10 @@ export const editCategoryRoute = buildRoute({
await saveCategoryTranslations(c, id, { title });
+ await c.get("events").emit("blog.category.updated", {
+ categoryId: id,
+ });
+
return c.json(category);
},
});
diff --git a/plugins/blog/src/api/modules/admin/posts/routes/create.route.ts b/plugins/blog/src/api/modules/admin/posts/routes/create.route.ts
index ef63b2e6e..c6321ca04 100644
--- a/plugins/blog/src/api/modules/admin/posts/routes/create.route.ts
+++ b/plugins/blog/src/api/modules/admin/posts/routes/create.route.ts
@@ -123,6 +123,11 @@ export const createPostRoute = buildRoute({
await savePostTranslations(c, post.id, { title, content, friendlyUrl });
await reindexBlogPost(c, post);
+ await c.get("events").emit("blog.post.created", {
+ postId: post.id,
+ categoryId: post.categoryId,
+ });
+
return c.json(post, 201);
},
});
diff --git a/plugins/blog/src/api/modules/admin/posts/routes/delete.route.ts b/plugins/blog/src/api/modules/admin/posts/routes/delete.route.ts
index 3757c731f..6337e7c73 100644
--- a/plugins/blog/src/api/modules/admin/posts/routes/delete.route.ts
+++ b/plugins/blog/src/api/modules/admin/posts/routes/delete.route.ts
@@ -41,6 +41,11 @@ export const deletePostRoute = buildRoute({
await c.get("search").delete("blog_post", id);
+ await c.get("events").emit("blog.post.deleted", {
+ postId: id,
+ categoryId: result[0].categoryId,
+ });
+
return c.body(null, 204);
},
});
diff --git a/plugins/blog/src/api/modules/admin/posts/routes/edit.route.ts b/plugins/blog/src/api/modules/admin/posts/routes/edit.route.ts
index 67a746a26..f3d80dc6d 100644
--- a/plugins/blog/src/api/modules/admin/posts/routes/edit.route.ts
+++ b/plugins/blog/src/api/modules/admin/posts/routes/edit.route.ts
@@ -129,6 +129,11 @@ export const editPostRoute = buildRoute({
await savePostTranslations(c, id, { title, content, friendlyUrl });
await reindexBlogPost(c, post);
+ await c.get("events").emit("blog.post.updated", {
+ postId: post.id,
+ categoryId: post.categoryId,
+ });
+
return c.json(post);
},
});