Skip to content

Commit cdc1e76

Browse files
authored
fix: abort SSE event subscription on stream close so process can exit (#30)
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
1 parent 0e807d1 commit cdc1e76

4 files changed

Lines changed: 36 additions & 0 deletions

File tree

src/opencode-client-manager.ts

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,7 @@ export class OpencodeClientManager {
4040
private logger: Logger;
4141
private initPromise: Promise<OpencodeClient> | null = null;
4242
private isDisposed = false;
43+
private activeEventSubscriptions = new Set<AbortController>();
4344
private cleanupHandlersRegistered = false;
4445
private cleanupHandlers: {
4546
exit?: () => void;
@@ -321,6 +322,17 @@ export class OpencodeClientManager {
321322
return this.server !== null;
322323
}
323324

325+
/**
326+
* Track an event-stream AbortController so dispose() can tear down SSE
327+
* connections that are still open. Returns an unregister function.
328+
*/
329+
registerEventSubscription(controller: AbortController): () => void {
330+
this.activeEventSubscriptions.add(controller);
331+
return () => {
332+
this.activeEventSubscriptions.delete(controller);
333+
};
334+
}
335+
324336
/**
325337
* Dispose of the client manager, stopping the server if managed.
326338
*/
@@ -332,6 +344,13 @@ export class OpencodeClientManager {
332344
this.isDisposed = true;
333345
this.unregisterCleanupHandlers();
334346

347+
// Abort any SSE event streams still open; stopping a managed server does
348+
// not close them, and for reused servers nothing else would.
349+
for (const controller of this.activeEventSubscriptions) {
350+
controller.abort();
351+
}
352+
this.activeEventSubscriptions.clear();
353+
335354
// Release the singleton slot so a later getInstance() builds a fresh
336355
// manager instead of returning this disposed one.
337356
if (OpencodeClientManager.instance === this) {

src/opencode-language-model.test.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -76,6 +76,7 @@ const mockClientManager = {
7676
getClient: vi.fn().mockResolvedValue(mockClient),
7777
dispose: vi.fn().mockResolvedValue(undefined),
7878
getServerUrl: vi.fn().mockReturnValue("http://127.0.0.1:4096"),
79+
registerEventSubscription: vi.fn().mockReturnValue(() => {}),
7980
};
8081

8182
describe("opencode-language-model", () => {

src/opencode-language-model.ts

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -403,9 +403,17 @@ export class OpencodeLanguageModel implements LanguageModelV3 {
403403
const streamWarnings = [...warnings];
404404
let streamStartEmitted = false;
405405

406+
// The SDK's SSE generator only cancels its fetch reader from an
407+
// abort-signal handler; closing the iterator alone leaves the socket
408+
// open and keeps the Node event loop alive.
409+
const eventAbortController = new AbortController();
410+
const unregisterEventSubscription =
411+
this.clientManager.registerEventSubscription(eventAbortController);
412+
406413
try {
407414
const eventsResult = await client.event.subscribe(
408415
directory ? { directory } : undefined,
416+
{ signal: eventAbortController.signal },
409417
);
410418

411419
const eventStream = eventsResult.stream;
@@ -484,6 +492,10 @@ export class OpencodeLanguageModel implements LanguageModelV3 {
484492
return;
485493
}
486494
iteratorClosed = true;
495+
// Abort before iterator.return(): the abort handler cancels the
496+
// SSE reader, which also unblocks a pending iterator.next() the
497+
// return() call would otherwise wait on.
498+
eventAbortController.abort();
487499
if (typeof iterator.return === "function") {
488500
try {
489501
await iterator.return(undefined);
@@ -581,6 +593,8 @@ export class OpencodeLanguageModel implements LanguageModelV3 {
581593
controller.enqueue({ type: "error", error: wrappedError });
582594
}
583595
} finally {
596+
eventAbortController.abort();
597+
unregisterEventSubscription();
584598
controller.close();
585599
}
586600
},

src/opencode-provider.test.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,13 +26,15 @@ vi.mock("./opencode-client-manager.js", async (importOriginal) => {
2626
}),
2727
dispose: vi.fn().mockResolvedValue(undefined),
2828
getServerUrl: vi.fn().mockReturnValue("http://127.0.0.1:4096"),
29+
registerEventSubscription: vi.fn().mockReturnValue(() => {}),
2930
}),
3031
resetInstance: vi.fn(),
3132
},
3233
createClientManagerFromSettings: vi.fn().mockReturnValue({
3334
getClient: vi.fn().mockResolvedValue({}),
3435
dispose: vi.fn().mockResolvedValue(undefined),
3536
getServerUrl: vi.fn().mockReturnValue("http://127.0.0.1:4096"),
37+
registerEventSubscription: vi.fn().mockReturnValue(() => {}),
3638
}),
3739
};
3840
});

0 commit comments

Comments
 (0)