Browse Source

Merge pull request #418 from alvinunreal/fix/multiplexer-session-race

Stabilize multiplexer pane lifecycle
Alvin 3 months ago
parent
commit
c269e672c8
2 changed files with 317 additions and 66 deletions
  1. 183 26
      src/multiplexer/session-manager.test.ts
  2. 134 40
      src/multiplexer/session-manager.ts

+ 183 - 26
src/multiplexer/session-manager.test.ts

@@ -255,7 +255,7 @@ describe('MultiplexerSessionManager', () => {
       expect(mockMultiplexer.closePane).not.toHaveBeenCalled();
     });
 
-    test('respawns pane on busy for known prior session', async () => {
+    test('respawns pane on later busy after idle close for resumable session', async () => {
       const ctx = createMockContext();
       const manager = new MultiplexerSessionManager(
         ctx,
@@ -308,73 +308,230 @@ describe('MultiplexerSessionManager', () => {
       expect(mockMultiplexer.closePane).toHaveBeenCalledTimes(1);
     });
 
-    test('does nothing on busy for unknown session', async () => {
+    test('respawns after in-flight idle close when busy resumes same session', async () => {
       const ctx = createMockContext();
       const manager = new MultiplexerSessionManager(
         ctx,
         defaultMultiplexerConfig,
       );
+      const closeDeferred = createDeferred<boolean>();
 
-      await manager.onSessionStatus({
+      mockMultiplexer.spawnPane
+        .mockResolvedValueOnce({
+          success: true,
+          paneId: 'p-close-race',
+        })
+        .mockResolvedValueOnce({
+          success: true,
+          paneId: 'p-close-race-resumed',
+        });
+      mockMultiplexer.closePane.mockImplementationOnce(
+        () => closeDeferred.promise,
+      );
+
+      await manager.onSessionCreated({
+        type: 'session.created',
+        properties: {
+          info: {
+            id: 'child-close-race',
+            parentID: 'parent-close-race',
+            title: 'Worker',
+          },
+        },
+      });
+
+      const idlePromise = manager.onSessionStatus({
         type: 'session.status',
         properties: {
-          sessionID: 'unknown-session',
+          sessionID: 'child-close-race',
+          status: { type: 'idle' },
+        },
+      });
+
+      await Promise.resolve();
+
+      const busyPromise = manager.onSessionStatus({
+        type: 'session.status',
+        properties: {
+          sessionID: 'child-close-race',
           status: { type: 'busy' },
         },
       });
 
-      expect(mockMultiplexer.spawnPane).not.toHaveBeenCalled();
+      expect(mockMultiplexer.spawnPane).toHaveBeenCalledTimes(1);
+
+      closeDeferred.resolve(true);
+      await Promise.all([idlePromise, busyPromise]);
+
+      expect(mockMultiplexer.closePane).toHaveBeenCalledTimes(1);
+      expect(mockMultiplexer.spawnPane).toHaveBeenCalledTimes(2);
+      expect(mockMultiplexer.spawnPane).toHaveBeenLastCalledWith(
+        'child-close-race',
+        'Worker',
+        `http://localhost:${process.env.OPENCODE_PORT ?? '4096'}/`,
+        '/test/directory',
+      );
     });
 
-    test('re-checks tracked sessions after async respawn guard', async () => {
+    test('does not respawn after in-flight close if session is deleted', async () => {
       const ctx = createMockContext();
       const manager = new MultiplexerSessionManager(
         ctx,
         defaultMultiplexerConfig,
       );
+      const closeDeferred = createDeferred<boolean>();
 
       mockMultiplexer.spawnPane
-        .mockResolvedValueOnce({ success: true, paneId: 'p-1' })
         .mockResolvedValueOnce({
           success: true,
-          paneId: 'p-should-not-happen',
+          paneId: 'p-delete-race',
+        })
+        .mockResolvedValueOnce({
+          success: true,
+          paneId: 'p-should-not-respawn',
         });
+      mockMultiplexer.closePane.mockImplementationOnce(
+        () => closeDeferred.promise,
+      );
 
       await manager.onSessionCreated({
         type: 'session.created',
         properties: {
           info: {
-            id: 'child-999',
-            parentID: 'parent-999',
+            id: 'child-delete-race',
+            parentID: 'parent-delete-race',
             title: 'Worker',
-            directory: '/task/dir',
           },
         },
       });
 
-      ctx.client.session.status.mockResolvedValue({
-        data: { 'child-999': { type: 'idle' } },
+      const idlePromise = manager.onSessionStatus({
+        type: 'session.status',
+        properties: {
+          sessionID: 'child-delete-race',
+          status: { type: 'idle' },
+        },
       });
-      await (manager as any).pollSessions();
 
-      const respawnPromise = (manager as any).respawnIfKnown('child-999');
+      await Promise.resolve();
 
-      (manager as any).sessions.set('child-999', {
-        sessionId: 'child-999',
-        paneId: 'p-existing',
-        parentId: 'parent-999',
-        title: 'Worker',
-        directory: '/task/dir',
-        createdAt: Date.now(),
-        lastSeenAt: Date.now(),
+      const busyPromise = manager.onSessionStatus({
+        type: 'session.status',
+        properties: {
+          sessionID: 'child-delete-race',
+          status: { type: 'busy' },
+        },
+      });
+
+      const deletedPromise = manager.onSessionDeleted({
+        type: 'session.deleted',
+        properties: {
+          sessionID: 'child-delete-race',
+        },
       });
 
-      await respawnPromise;
+      closeDeferred.resolve(true);
+      await Promise.all([idlePromise, busyPromise, deletedPromise]);
 
+      expect(mockMultiplexer.closePane).toHaveBeenCalledTimes(1);
       expect(mockMultiplexer.spawnPane).toHaveBeenCalledTimes(1);
-      expect((manager as any).sessions.get('child-999')?.paneId).toBe(
-        'p-existing',
+    });
+
+    test('closes pane on session.deleted using info.id', async () => {
+      const ctx = createMockContext();
+      const manager = new MultiplexerSessionManager(
+        ctx,
+        defaultMultiplexerConfig,
       );
+
+      mockMultiplexer.spawnPane.mockResolvedValueOnce({
+        success: true,
+        paneId: 'p-info-id',
+      });
+
+      await manager.onSessionCreated({
+        type: 'session.created',
+        properties: {
+          info: {
+            id: 'child-info-id',
+            parentID: 'parent-info-id',
+          },
+        },
+      });
+
+      await manager.onSessionDeleted({
+        type: 'session.deleted',
+        properties: {
+          info: { id: 'child-info-id' },
+        },
+      });
+
+      expect(mockMultiplexer.closePane).toHaveBeenCalledWith('p-info-id');
+
+      await manager.onSessionStatus({
+        type: 'session.status',
+        properties: {
+          sessionID: 'child-info-id',
+          status: { type: 'busy' },
+        },
+      });
+
+      expect(mockMultiplexer.spawnPane).toHaveBeenCalledTimes(1);
+    });
+
+    test('closes pane returned by a stale spawn after session deleted', async () => {
+      const ctx = createMockContext();
+      const manager = new MultiplexerSessionManager(
+        ctx,
+        defaultMultiplexerConfig,
+      );
+      const spawnDeferred = createDeferred<{ success: true; paneId: string }>();
+
+      mockMultiplexer.spawnPane.mockImplementationOnce(
+        () => spawnDeferred.promise,
+      );
+
+      const createPromise = manager.onSessionCreated({
+        type: 'session.created',
+        properties: {
+          info: {
+            id: 'child-stale-spawn',
+            parentID: 'parent-stale-spawn',
+          },
+        },
+      });
+
+      await Promise.resolve();
+
+      await manager.onSessionDeleted({
+        type: 'session.deleted',
+        properties: {
+          info: { id: 'child-stale-spawn' },
+        },
+      });
+
+      spawnDeferred.resolve({ success: true, paneId: 'p-stale-spawn' });
+      await createPromise;
+
+      expect(mockMultiplexer.closePane).toHaveBeenCalledWith('p-stale-spawn');
+    });
+
+    test('does nothing on busy for unknown session', async () => {
+      const ctx = createMockContext();
+      const manager = new MultiplexerSessionManager(
+        ctx,
+        defaultMultiplexerConfig,
+      );
+
+      await manager.onSessionStatus({
+        type: 'session.status',
+        properties: {
+          sessionID: 'unknown-session',
+          status: { type: 'busy' },
+        },
+      });
+
+      expect(mockMultiplexer.spawnPane).not.toHaveBeenCalled();
     });
 
     test('does not respawn while initial pane spawn is still in progress', async () => {

+ 134 - 40
src/multiplexer/session-manager.ts

@@ -41,6 +41,8 @@ interface SessionEvent {
   };
 }
 
+type CloseReason = 'idle' | 'deleted' | 'missing' | 'timeout';
+
 const SESSION_TIMEOUT_MS = 10 * 60 * 1000;
 const SESSION_MISSING_GRACE_MS = POLL_INTERVAL_BACKGROUND_MS * 3;
 
@@ -58,6 +60,7 @@ export class MultiplexerSessionManager {
   private sessions = new Map<string, TrackedSession>();
   private knownSessions = new Map<string, KnownSession>();
   private spawningSessions = new Set<string>();
+  private closingSessions = new Map<string, Promise<void>>();
   private pollInterval?: ReturnType<typeof setInterval>;
   private enabled = false;
 
@@ -95,12 +98,6 @@ export class MultiplexerSessionManager {
     const title = info.title ?? 'Subagent';
     const directory = info.directory ?? this.directory;
 
-    this.knownSessions.set(sessionId, {
-      parentId,
-      title,
-      directory,
-    });
-
     if (this.isTrackedOrSpawning(sessionId)) {
       log('[multiplexer-session-manager] session already tracked or spawning', {
         sessionId,
@@ -108,6 +105,17 @@ export class MultiplexerSessionManager {
       return;
     }
 
+    const closing = this.closingSessions.get(sessionId);
+    if (closing) await closing;
+
+    if (this.isTrackedOrSpawning(sessionId)) return;
+
+    this.knownSessions.set(sessionId, {
+      parentId,
+      title,
+      directory,
+    });
+
     this.spawningSessions.add(sessionId);
 
     try {
@@ -119,7 +127,7 @@ export class MultiplexerSessionManager {
         return;
       }
 
-      if (this.sessions.has(sessionId)) {
+      if (this.closingSessions.has(sessionId) || this.sessions.has(sessionId)) {
         return;
       }
 
@@ -141,25 +149,42 @@ export class MultiplexerSessionManager {
           return { success: false, paneId: undefined };
         });
 
-      if (paneResult.success && paneResult.paneId) {
-        const now = Date.now();
-        this.sessions.set(sessionId, {
-          sessionId,
-          paneId: paneResult.paneId,
-          parentId,
-          title,
-          directory,
-          createdAt: now,
-          lastSeenAt: now,
-        });
-
-        log('[multiplexer-session-manager] pane spawned', {
-          sessionId,
-          paneId: paneResult.paneId,
-        });
+      if (!paneResult.success || !paneResult.paneId) return;
 
-        this.startPolling();
+      if (
+        !this.knownSessions.has(sessionId) ||
+        this.closingSessions.has(sessionId)
+      ) {
+        await this.multiplexer.closePane(paneResult.paneId).catch((err) =>
+          log(
+            '[multiplexer-session-manager] closing stale spawned pane failed',
+            {
+              sessionId,
+              paneId: paneResult.paneId,
+              error: String(err),
+            },
+          ),
+        );
+        return;
       }
+
+      const now = Date.now();
+      this.sessions.set(sessionId, {
+        sessionId,
+        paneId: paneResult.paneId,
+        parentId,
+        title,
+        directory,
+        createdAt: now,
+        lastSeenAt: now,
+      });
+
+      log('[multiplexer-session-manager] pane spawned', {
+        sessionId,
+        paneId: paneResult.paneId,
+      });
+
+      this.startPolling();
     } finally {
       this.spawningSessions.delete(sessionId);
     }
@@ -173,7 +198,7 @@ export class MultiplexerSessionManager {
     if (!sessionId) return;
 
     if (event.properties?.status?.type === 'idle') {
-      await this.closeSession(sessionId);
+      await this.closeSession(sessionId, 'idle');
       return;
     }
 
@@ -186,15 +211,14 @@ export class MultiplexerSessionManager {
     if (!this.enabled) return;
     if (event.type !== 'session.deleted') return;
 
-    const sessionId = event.properties?.sessionID;
+    const sessionId = this.getSessionId(event);
     if (!sessionId) return;
 
     log('[multiplexer-session-manager] session deleted, closing pane', {
       sessionId,
     });
 
-    await this.closeSession(sessionId);
-    this.knownSessions.delete(sessionId);
+    await this.closeSession(sessionId, 'deleted');
   }
 
   private startPolling(): void {
@@ -229,7 +253,8 @@ export class MultiplexerSessionManager {
       >;
 
       const now = Date.now();
-      const sessionsToClose: string[] = [];
+      const sessionsToClose: Array<{ sessionId: string; reason: CloseReason }> =
+        [];
 
       for (const [sessionId, tracked] of this.sessions.entries()) {
         const status = allStatuses[sessionId];
@@ -248,38 +273,71 @@ export class MultiplexerSessionManager {
         const isTimedOut = now - tracked.createdAt > SESSION_TIMEOUT_MS;
 
         if (isIdle || missingTooLong || isTimedOut) {
-          sessionsToClose.push(sessionId);
+          sessionsToClose.push({
+            sessionId,
+            reason: isIdle ? 'idle' : isTimedOut ? 'timeout' : 'missing',
+          });
         }
       }
 
-      for (const sessionId of sessionsToClose) {
-        await this.closeSession(sessionId);
+      for (const { sessionId, reason } of sessionsToClose) {
+        await this.closeSession(sessionId, reason);
       }
     } catch (err) {
       log('[multiplexer-session-manager] poll error', { error: String(err) });
     }
   }
 
-  private async closeSession(sessionId: string): Promise<void> {
+  private async closeSession(
+    sessionId: string,
+    reason: CloseReason,
+  ): Promise<void> {
+    if (reason === 'deleted') {
+      this.knownSessions.delete(sessionId);
+    }
+
+    const existingClose = this.closingSessions.get(sessionId);
+    if (existingClose) return existingClose;
+
     const tracked = this.sessions.get(sessionId);
     if (!tracked || !this.multiplexer) return;
 
+    this.sessions.delete(sessionId);
+
     log('[multiplexer-session-manager] closing session pane', {
       sessionId,
       paneId: tracked.paneId,
+      reason,
     });
 
-    await this.multiplexer.closePane(tracked.paneId);
-    this.sessions.delete(sessionId);
+    const closePromise: Promise<void> = this.multiplexer
+      .closePane(tracked.paneId)
+      .then(() => undefined)
+      .catch((err) =>
+        log('[multiplexer-session-manager] failed to close session pane', {
+          sessionId,
+          paneId: tracked.paneId,
+          reason,
+          error: String(err),
+        }),
+      )
+      .finally(() => {
+        this.closingSessions.delete(sessionId);
+        this.updatePolling();
+      });
 
-    if (this.sessions.size === 0) {
-      this.stopPolling();
-    }
+    this.closingSessions.set(sessionId, closePromise);
+    await closePromise;
   }
 
   private async respawnIfKnown(sessionId: string): Promise<void> {
     if (!this.enabled || !this.multiplexer) return;
-    if (this.isTrackedOrSpawning(sessionId)) return;
+    const closing = this.closingSessions.get(sessionId);
+    if (closing) await closing;
+
+    if (this.isTrackedOrSpawning(sessionId)) {
+      return;
+    }
 
     const known = this.knownSessions.get(sessionId);
     if (!known) return;
@@ -299,7 +357,9 @@ export class MultiplexerSessionManager {
         return;
       }
 
-      if (this.sessions.has(sessionId)) return;
+      if (this.sessions.has(sessionId) || this.closingSessions.has(sessionId)) {
+        return;
+      }
 
       log(
         '[multiplexer-session-manager] child session busy again, respawning pane',
@@ -321,6 +381,23 @@ export class MultiplexerSessionManager {
 
       if (!paneResult.success || !paneResult.paneId) return;
 
+      if (
+        !this.knownSessions.has(sessionId) ||
+        this.closingSessions.has(sessionId)
+      ) {
+        await this.multiplexer.closePane(paneResult.paneId).catch((err) =>
+          log(
+            '[multiplexer-session-manager] closing stale respawned pane failed',
+            {
+              sessionId,
+              paneId: paneResult.paneId,
+              error: String(err),
+            },
+          ),
+        );
+        return;
+      }
+
       const now = Date.now();
       this.sessions.set(sessionId, {
         sessionId,
@@ -347,9 +424,25 @@ export class MultiplexerSessionManager {
     return this.sessions.has(sessionId) || this.spawningSessions.has(sessionId);
   }
 
+  private updatePolling(): void {
+    if (this.sessions.size > 0 || this.closingSessions.size > 0) {
+      this.startPolling();
+    } else {
+      this.stopPolling();
+    }
+  }
+
+  private getSessionId(event: SessionEvent): string | undefined {
+    return event.properties?.info?.id ?? event.properties?.sessionID;
+  }
+
   async cleanup(): Promise<void> {
     this.stopPolling();
 
+    if (this.closingSessions.size > 0) {
+      await Promise.all(this.closingSessions.values());
+    }
+
     if (this.sessions.size > 0 && this.multiplexer) {
       log('[multiplexer-session-manager] closing all panes', {
         count: this.sessions.size,
@@ -369,6 +462,7 @@ export class MultiplexerSessionManager {
 
     this.knownSessions.clear();
     this.spawningSessions.clear();
+    this.closingSessions.clear();
 
     log('[multiplexer-session-manager] cleanup complete');
   }