feat(FN-2673): stream automation schedule events over SSE

- Wire AutomationStore into dashboard SSE setup for both default and project-scoped /api/events streams
- Emit schedule:created, schedule:updated, schedule:deleted, and schedule:run events to connected SSE clients
- Add SSE coverage for automation event subscription, relay payloads, cleanup on disconnect, and graceful behavior without automation store
- Add server events integration coverage for automationStore wiring and document automation schedule events in architecture docs
This commit is contained in:
Fusion
2026-04-27 03:21:59 -07:00
committed by gsxdsm
parent f928f0a92e
commit ae988d0ef7
5 changed files with 185 additions and 2 deletions

View File

@@ -430,7 +430,7 @@ Key server capabilities:
### Real-time channels
- **SSE**: `/api/events` (`sse.ts`)
- Emits `task:*`, mission events, AI session updates
- Emits `task:*`, mission events, AI session updates, and automation schedule events (`schedule:created`, `schedule:updated`, `schedule:deleted`, `schedule:run`)
- Project-scoped: resolves project context from query param or engine manager
- **Chat streaming**: `/api/chat/sessions/:id/messages` (`routes.ts` + `chat.ts`)
- Streams assistant responses as SSE events for chat sessions

View File

@@ -109,6 +109,20 @@ describe("server events endpoint integration", () => {
expect(app).toBeDefined();
});
it("wires automationStore into SSE events endpoint", () => {
const store = createMockStore();
const mockAutomationStore = {
on: vi.fn(),
off: vi.fn(),
};
const app = createServer(store, {
automationStore: mockAutomationStore as any,
});
expect(app).toBeDefined();
});
describe("SSE project-scoped event routing", () => {
// Note: Full SSE streaming tests are complex due to connection timeouts.
// These tests verify the endpoint routes are properly configured.

View File

@@ -1,7 +1,7 @@
import { EventEmitter } from "node:events";
import { afterEach, describe, expect, it, vi } from "vitest";
import type { Request, Response } from "express";
import type { TaskStore } from "@fusion/core";
import type { TaskStore, AutomationStore } from "@fusion/core";
import { createSSE, disconnectSSEClient, getActiveSSEConnections, markSSEClientAlive } from "../sse.js";
class MockSocket extends EventEmitter {
@@ -61,10 +61,140 @@ function openSseConnection(clientId: string, projectId?: string) {
return { req, res, socket, store };
}
function createMockAutomationStore(): AutomationStore {
return {
on: vi.fn(),
off: vi.fn(),
} as unknown as AutomationStore;
}
function openSseConnectionWithAutomation(clientId: string, projectId?: string) {
const store = createMockStore();
const automationStore = createMockAutomationStore();
const socket = new MockSocket();
const req = new EventEmitter() as Request & { query: Record<string, string>; socket: MockSocket };
req.query = projectId ? { clientId, projectId } : { clientId };
req.socket = socket;
const res = new MockResponse(socket);
createSSE(
store,
undefined,
undefined,
undefined,
projectId ? { projectId } : undefined,
undefined,
undefined,
undefined,
automationStore,
)(req, res as unknown as Response);
return { req, res, socket, store, automationStore };
}
afterEach(() => {
vi.useRealTimers();
});
describe("automation store SSE events", () => {
it("subscribes to all automation store events", () => {
const connection = openSseConnectionWithAutomation("automation-subscribe");
expect(connection.automationStore.on).toHaveBeenCalledWith("schedule:created", expect.any(Function));
expect(connection.automationStore.on).toHaveBeenCalledWith("schedule:updated", expect.any(Function));
expect(connection.automationStore.on).toHaveBeenCalledWith("schedule:deleted", expect.any(Function));
expect(connection.automationStore.on).toHaveBeenCalledWith("schedule:run", expect.any(Function));
connection.req.emit("close");
});
it("relays schedule:created events", () => {
const connection = openSseConnectionWithAutomation("automation-created");
const handler = vi.mocked(connection.automationStore.on).mock.calls.find(
([eventName]) => eventName === "schedule:created",
)?.[1] as ((schedule: unknown) => void) | undefined;
const schedule = { id: "sch-1", name: "Nightly" };
handler?.(schedule);
expect(connection.res.write).toHaveBeenCalledWith(
`event: schedule:created\ndata: ${JSON.stringify(schedule)}\n\n`,
);
connection.req.emit("close");
});
it("relays schedule:updated events", () => {
const connection = openSseConnectionWithAutomation("automation-updated");
const handler = vi.mocked(connection.automationStore.on).mock.calls.find(
([eventName]) => eventName === "schedule:updated",
)?.[1] as ((schedule: unknown) => void) | undefined;
const schedule = { id: "sch-2", name: "Weekly" };
handler?.(schedule);
expect(connection.res.write).toHaveBeenCalledWith(
`event: schedule:updated\ndata: ${JSON.stringify(schedule)}\n\n`,
);
connection.req.emit("close");
});
it("relays schedule:deleted events", () => {
const connection = openSseConnectionWithAutomation("automation-deleted");
const handler = vi.mocked(connection.automationStore.on).mock.calls.find(
([eventName]) => eventName === "schedule:deleted",
)?.[1] as ((schedule: unknown) => void) | undefined;
const schedule = { id: "sch-3", name: "Cleanup" };
handler?.(schedule);
expect(connection.res.write).toHaveBeenCalledWith(
`event: schedule:deleted\ndata: ${JSON.stringify(schedule)}\n\n`,
);
connection.req.emit("close");
});
it("relays schedule:run events", () => {
const connection = openSseConnectionWithAutomation("automation-run");
const handler = vi.mocked(connection.automationStore.on).mock.calls.find(
([eventName]) => eventName === "schedule:run",
)?.[1] as ((data: unknown) => void) | undefined;
const payload = {
schedule: { id: "sch-4", name: "Run" },
result: { status: "success", startedAt: "2026-04-27T00:00:00.000Z" },
};
handler?.(payload);
expect(connection.res.write).toHaveBeenCalledWith(
`event: schedule:run\ndata: ${JSON.stringify(payload)}\n\n`,
);
connection.req.emit("close");
});
it("cleans up automation store listeners on disconnect", () => {
const connection = openSseConnectionWithAutomation("automation-cleanup");
connection.req.emit("close");
expect(connection.automationStore.off).toHaveBeenCalledWith("schedule:created", expect.any(Function));
expect(connection.automationStore.off).toHaveBeenCalledWith("schedule:updated", expect.any(Function));
expect(connection.automationStore.off).toHaveBeenCalledWith("schedule:deleted", expect.any(Function));
expect(connection.automationStore.off).toHaveBeenCalledWith("schedule:run", expect.any(Function));
});
it("works without automationStore (graceful degradation)", () => {
const connection = openSseConnection("automation-none");
expect(connection.res.write).toHaveBeenCalledWith(": connected\n\n");
connection.req.emit("close");
});
});
describe("createSSE client cleanup", () => {
it("disconnectSSEClient closes and unregisters the matching stream", () => {
const baseline = getActiveSSEConnections();

View File

@@ -617,6 +617,7 @@ export function createServer(store: TaskStore, options?: ServerOptions): ReturnT
defaultAgentStore,
defaultMessageStore,
chatStore,
options?.automationStore,
)(req, res);
return;
}
@@ -628,12 +629,14 @@ export function createServer(store: TaskStore, options?: ServerOptions): ReturnT
let scopedStore: TaskStore;
let agentStore;
let messageStore: MessageStore | undefined;
let automationStore: AutomationStore | undefined;
if (engineManager) {
const engine = engineManager.getEngine(projectId);
scopedStore = engine?.getTaskStore() ?? await getOrCreateProjectStore(projectId);
// Use the engine's stores if available
agentStore = engine?.getAgentStore();
messageStore = engine?.getMessageStore();
automationStore = engine?.getAutomationStore();
} else {
scopedStore = await getOrCreateProjectStore(projectId);
}
@@ -643,6 +646,9 @@ export function createServer(store: TaskStore, options?: ServerOptions): ReturnT
agentStore = new AgentStoreClass({ rootDir: scopedStore.getFusionDir() });
await agentStore.init();
}
if (!automationStore) {
automationStore = options?.automationStore;
}
createSSE(
scopedStore,
scopedStore.getMissionStore(),
@@ -654,6 +660,7 @@ export function createServer(store: TaskStore, options?: ServerOptions): ReturnT
agentStore,
messageStore,
chatStore,
automationStore,
)(req, res);
} catch (err: unknown) {
sendErrorResponse(res, 500, err instanceof Error ? err.message : "Failed to open project event stream");

View File

@@ -10,6 +10,7 @@ import type {
MissionValidatorRun,
FixFeatureCreatedPayload,
ChatStore,
AutomationStore,
} from "@fusion/core";
import type { AiSessionStore } from "./ai-session-store.js";
@@ -285,6 +286,7 @@ export function createSSE(
agentStore?: AgentStore,
messageStore?: MessageStore,
chatStore?: ChatStore,
automationStore?: AutomationStore,
) {
const { projectId } = options ?? {};
@@ -519,6 +521,23 @@ export function createSSE(
send(`event: chat:message:deleted\ndata: ${JSON.stringify({ id: messageId })}\n\n`);
};
// --- Automation store event handlers ---
const onScheduleCreated = (schedule: unknown) => {
send(`event: schedule:created\ndata: ${JSON.stringify(schedule)}\n\n`);
};
const onScheduleUpdated = (schedule: unknown) => {
send(`event: schedule:updated\ndata: ${JSON.stringify(schedule)}\n\n`);
};
const onScheduleDeleted = (schedule: unknown) => {
send(`event: schedule:deleted\ndata: ${JSON.stringify(schedule)}\n\n`);
};
const onScheduleRun = (data: unknown) => {
send(`event: schedule:run\ndata: ${JSON.stringify(data)}\n\n`);
};
// --- Cleanup (all handlers are defined above, safe to reference) ---
let cleaned = false;
@@ -603,6 +622,12 @@ export function createSSE(
chatStore.off("chat:message:added", onChatMessageAdded);
chatStore.off("chat:message:deleted", onChatMessageDeleted);
}
if (automationStore) {
automationStore.off("schedule:created", onScheduleCreated);
automationStore.off("schedule:updated", onScheduleUpdated);
automationStore.off("schedule:deleted", onScheduleDeleted);
automationStore.off("schedule:run", onScheduleRun);
}
}
function closeConnection(reason: SSECloseReason): void {
@@ -694,6 +719,13 @@ export function createSSE(
chatStore.on("chat:message:deleted", onChatMessageDeleted);
}
if (automationStore) {
automationStore.on("schedule:created", onScheduleCreated);
automationStore.on("schedule:updated", onScheduleUpdated);
automationStore.on("schedule:deleted", onScheduleDeleted);
automationStore.on("schedule:run", onScheduleRun);
}
// Heartbeat every 30s to keep connection alive.
// Sent as a named event so the client's EventSource can detect it
// (SSE comments starting with ":" are silently consumed and never