From 9a274e22ce14d4ddaa68146e61013f9c1756f479 Mon Sep 17 00:00:00 2001 From: Ben Vargas Date: Thu, 11 Jun 2026 09:09:15 -0600 Subject: [PATCH] fix: abort SSE event subscription on stream close so process can exit The SDK's SSE generator only cancels its fetch reader from an abort-signal handler; closing the iterator alone runs releaseLock() without cancel(), leaving the socket open. Every completed stream therefore held an ESTABLISHED connection to the OpenCode server that kept the Node event loop alive indefinitely after streamText() finished. - Create an AbortController per stream in doStream and pass its signal to client.event.subscribe - Abort it in closeIterator (before iterator.return(), which also unblocks a pending next()) and in the stream's outer finally - Track active subscription controllers in OpencodeClientManager so dispose() tears down SSE connections to non-managed servers --- src/opencode-client-manager.ts | 19 +++++++++++++++++++ src/opencode-language-model.test.ts | 1 + src/opencode-language-model.ts | 14 ++++++++++++++ src/opencode-provider.test.ts | 2 ++ 4 files changed, 36 insertions(+) diff --git a/src/opencode-client-manager.ts b/src/opencode-client-manager.ts index 140fa67..e4e3472 100644 --- a/src/opencode-client-manager.ts +++ b/src/opencode-client-manager.ts @@ -40,6 +40,7 @@ export class OpencodeClientManager { private logger: Logger; private initPromise: Promise | null = null; private isDisposed = false; + private activeEventSubscriptions = new Set(); private cleanupHandlersRegistered = false; private cleanupHandlers: { exit?: () => void; @@ -321,6 +322,17 @@ export class OpencodeClientManager { return this.server !== null; } + /** + * Track an event-stream AbortController so dispose() can tear down SSE + * connections that are still open. Returns an unregister function. + */ + registerEventSubscription(controller: AbortController): () => void { + this.activeEventSubscriptions.add(controller); + return () => { + this.activeEventSubscriptions.delete(controller); + }; + } + /** * Dispose of the client manager, stopping the server if managed. */ @@ -332,6 +344,13 @@ export class OpencodeClientManager { this.isDisposed = true; this.unregisterCleanupHandlers(); + // Abort any SSE event streams still open; stopping a managed server does + // not close them, and for reused servers nothing else would. + for (const controller of this.activeEventSubscriptions) { + controller.abort(); + } + this.activeEventSubscriptions.clear(); + // Release the singleton slot so a later getInstance() builds a fresh // manager instead of returning this disposed one. if (OpencodeClientManager.instance === this) { diff --git a/src/opencode-language-model.test.ts b/src/opencode-language-model.test.ts index dca954e..3337965 100644 --- a/src/opencode-language-model.test.ts +++ b/src/opencode-language-model.test.ts @@ -76,6 +76,7 @@ const mockClientManager = { getClient: vi.fn().mockResolvedValue(mockClient), dispose: vi.fn().mockResolvedValue(undefined), getServerUrl: vi.fn().mockReturnValue("http://127.0.0.1:4096"), + registerEventSubscription: vi.fn().mockReturnValue(() => {}), }; describe("opencode-language-model", () => { diff --git a/src/opencode-language-model.ts b/src/opencode-language-model.ts index 53e6bb0..3884aa3 100644 --- a/src/opencode-language-model.ts +++ b/src/opencode-language-model.ts @@ -403,9 +403,17 @@ export class OpencodeLanguageModel implements LanguageModelV3 { const streamWarnings = [...warnings]; let streamStartEmitted = false; + // The SDK's SSE generator only cancels its fetch reader from an + // abort-signal handler; closing the iterator alone leaves the socket + // open and keeps the Node event loop alive. + const eventAbortController = new AbortController(); + const unregisterEventSubscription = + this.clientManager.registerEventSubscription(eventAbortController); + try { const eventsResult = await client.event.subscribe( directory ? { directory } : undefined, + { signal: eventAbortController.signal }, ); const eventStream = eventsResult.stream; @@ -484,6 +492,10 @@ export class OpencodeLanguageModel implements LanguageModelV3 { return; } iteratorClosed = true; + // Abort before iterator.return(): the abort handler cancels the + // SSE reader, which also unblocks a pending iterator.next() the + // return() call would otherwise wait on. + eventAbortController.abort(); if (typeof iterator.return === "function") { try { await iterator.return(undefined); @@ -581,6 +593,8 @@ export class OpencodeLanguageModel implements LanguageModelV3 { controller.enqueue({ type: "error", error: wrappedError }); } } finally { + eventAbortController.abort(); + unregisterEventSubscription(); controller.close(); } }, diff --git a/src/opencode-provider.test.ts b/src/opencode-provider.test.ts index e2455c7..77e280c 100644 --- a/src/opencode-provider.test.ts +++ b/src/opencode-provider.test.ts @@ -26,6 +26,7 @@ vi.mock("./opencode-client-manager.js", async (importOriginal) => { }), dispose: vi.fn().mockResolvedValue(undefined), getServerUrl: vi.fn().mockReturnValue("http://127.0.0.1:4096"), + registerEventSubscription: vi.fn().mockReturnValue(() => {}), }), resetInstance: vi.fn(), }, @@ -33,6 +34,7 @@ vi.mock("./opencode-client-manager.js", async (importOriginal) => { getClient: vi.fn().mockResolvedValue({}), dispose: vi.fn().mockResolvedValue(undefined), getServerUrl: vi.fn().mockReturnValue("http://127.0.0.1:4096"), + registerEventSubscription: vi.fn().mockReturnValue(() => {}), }), }; });