Browse Source

Merge pull request #598 from alvinunreal/task/597-resource-audit

Fix dashboard shutdown resource cleanup
Alvin 1 month ago
parent
commit
5949e0db07
2 changed files with 179 additions and 10 deletions
  1. 137 1
      src/interview/dashboard.test.ts
  2. 42 9
      src/interview/dashboard.ts

+ 137 - 1
src/interview/dashboard.test.ts

@@ -1,6 +1,6 @@
 import { describe, expect, test } from 'bun:test';
 import * as fs from 'node:fs/promises';
-import { createServer } from 'node:http';
+import { createServer, get } from 'node:http';
 import * as path from 'node:path';
 import { createDashboardServer } from './dashboard';
 
@@ -47,7 +47,143 @@ async function createTempInterviewDir() {
   return tempDir;
 }
 
+async function openSseConnection(
+  baseUrl: string,
+  authToken: string,
+  interviewId: string,
+) {
+  return new Promise<{
+    firstChunk: Promise<string>;
+    closed: Promise<void>;
+  }>((resolve, reject) => {
+    let responded = false;
+    let sawStateEvent = false;
+
+    let firstChunkResolve!: (value: string) => void;
+    let firstChunkReject!: (reason: Error) => void;
+    const firstChunk = new Promise<string>((resolveFirst, rejectFirst) => {
+      firstChunkResolve = resolveFirst;
+      firstChunkReject = rejectFirst;
+    });
+
+    let closedResolve!: () => void;
+    let closedReject!: (reason: Error) => void;
+    const closed = new Promise<void>((resolveClosed, rejectClosed) => {
+      closedResolve = resolveClosed;
+      closedReject = rejectClosed;
+    });
+
+    const rejectStreams = (error: Error) => {
+      firstChunkReject(error);
+      closedReject(error);
+    };
+
+    const onError = (error: Error) => {
+      rejectStreams(error);
+      if (!responded) {
+        responded = true;
+        reject(error);
+      }
+    };
+
+    const request = get(
+      `${baseUrl}/api/interviews/${interviewId}/events?token=${authToken}`,
+      (response) => {
+        responded = true;
+        response.setEncoding('utf8');
+
+        let buffer = '';
+        response.on('data', (chunk: string) => {
+          buffer += chunk;
+          if (!sawStateEvent && buffer.includes('event: state')) {
+            sawStateEvent = true;
+            firstChunkResolve(buffer);
+          }
+        });
+
+        response.once('close', () => {
+          if (!sawStateEvent) {
+            firstChunkReject(
+              new Error('SSE response closed before initial state event'),
+            );
+          }
+          closedResolve();
+        });
+
+        response.once('error', onError);
+        resolve({ firstChunk, closed });
+      },
+    );
+
+    request.once('error', onError);
+  });
+}
+
 describe('dashboard server', () => {
+  describe('server lifecycle', () => {
+    test('close is safe before start and repeated', () => {
+      const dashboard = createDashboardServer({
+        port: 0,
+        outputFolder: 'interview',
+      });
+
+      expect(() => dashboard.close()).not.toThrow();
+      expect(() => dashboard.close()).not.toThrow();
+    });
+
+    test('closes active SSE responses on shutdown', async () => {
+      const { baseUrl, authToken, dashboard, cleanup } = await startDashboard();
+      try {
+        dashboard.pushState({
+          interviewId: 'lifecycle-sse',
+          sessionID: 'session-lifecycle',
+          idea: 'Lifecycle SSE',
+          mode: 'awaiting-user',
+          summary: 'Test',
+          title: 'Lifecycle SSE',
+          questions: [],
+          pendingAnswers: null,
+          lastUpdatedAt: Date.now(),
+          filePath: '',
+          nudgeAction: null,
+        });
+
+        const { firstChunk, closed } = await openSseConnection(
+          baseUrl,
+          authToken,
+          'lifecycle-sse',
+        );
+
+        expect(await firstChunk).toContain('event: state');
+
+        dashboard.close();
+        dashboard.close();
+
+        await closed;
+      } finally {
+        cleanup();
+      }
+    });
+
+    test('rejects SSE firstChunk when response closes before state', async () => {
+      const { baseUrl, authToken, cleanup } = await startDashboard();
+      try {
+        const connection = await openSseConnection(
+          baseUrl,
+          authToken,
+          'missing-lifecycle-sse',
+        );
+
+        await expect(connection.firstChunk).rejects.toThrow(
+          'SSE response closed before initial state event',
+        );
+        await connection.closed;
+      } finally {
+        cleanup();
+      }
+    });
+  });
+
   describe('health endpoint', () => {
     test('returns 200 with status ok and counts', async () => {
       const { baseUrl, cleanup } = await startDashboard();

+ 42 - 9
src/interview/dashboard.ts

@@ -263,15 +263,20 @@ export function createDashboardServer(config: DashboardConfig): {
   ]);
   const CACHE_TTL_MS = 24 * 60 * 60 * 1000;
   const CLEANUP_INTERVAL_MS = 60 * 60 * 1000;
-  const cleanupTimer = setInterval(() => {
-    const cutoff = Date.now() - CACHE_TTL_MS;
-    for (const [id, entry] of stateCache) {
-      if (TERMINAL_MODES.has(entry.mode) && entry.lastUpdatedAt < cutoff) {
-        stateCache.delete(id);
+  function createCleanupTimer(): ReturnType<typeof setInterval> {
+    const timer = setInterval(() => {
+      const cutoff = Date.now() - CACHE_TTL_MS;
+      for (const [id, entry] of stateCache) {
+        if (TERMINAL_MODES.has(entry.mode) && entry.lastUpdatedAt < cutoff) {
+          stateCache.delete(id);
+        }
       }
-    }
-  }, CLEANUP_INTERVAL_MS);
-  cleanupTimer.unref();
+    }, CLEANUP_INTERVAL_MS);
+    timer.unref();
+    return timer;
+  }
+
+  let cleanupTimer: ReturnType<typeof setInterval> | null = null;
 
   // File scan cache (TTL 10s)
   let fileCache: { items: InterviewFileItem[]; at: number } | null = null;
@@ -1288,6 +1293,10 @@ export function createDashboardServer(config: DashboardConfig): {
   function start(): Promise<string> {
     if (baseUrl) return Promise.resolve(baseUrl);
 
+    if (!cleanupTimer) {
+      cleanupTimer = createCleanupTimer();
+    }
+
     return new Promise((resolve, reject) => {
       const server = createServer((request, response) => {
         handleRequest(request, response).catch((error: unknown) => {
@@ -1325,8 +1334,32 @@ export function createDashboardServer(config: DashboardConfig): {
   }
 
   function close(): void {
+    if (cleanupTimer) {
+      clearInterval(cleanupTimer);
+      cleanupTimer = null;
+    }
+
+    for (const clients of sseClients.values()) {
+      for (const response of clients) {
+        try {
+          if (!response.writableEnded && !response.destroyed) {
+            response.end();
+          }
+        } catch {
+          try {
+            response.destroy();
+          } catch {
+            // Ignore cleanup errors
+          }
+        }
+      }
+    }
+
+    sseClients.clear();
+
+    removeAuthFile(config.port);
+
     if (activeServer) {
-      removeAuthFile(config.port);
       activeServer.closeAllConnections();
       activeServer.close();
       activeServer = null;