Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 11 additions & 2 deletions packages/cli/src/nonInteractive/control/ControlDispatcher.ts
Original file line number Diff line number Diff line change
Expand Up @@ -80,6 +80,8 @@ export class ControlDispatcher implements IPendingRequestRegistry {
private pendingOutgoingRequests: Map<string, PendingOutgoingRequest> =
new Map();

private abortHandler: (() => void) | null = null;

constructor(context: IControlContext) {
this.context = context;

Expand All @@ -102,9 +104,10 @@ export class ControlDispatcher implements IPendingRequestRegistry {
// this.hookController = new HookController(context, this, 'HookController');

// Listen for main abort signal
this.context.abortSignal.addEventListener('abort', () => {
this.abortHandler = () => {
this.shutdown();
});
};
this.context.abortSignal.addEventListener('abort', this.abortHandler);
}

/**
Expand Down Expand Up @@ -240,6 +243,12 @@ export class ControlDispatcher implements IPendingRequestRegistry {
shutdown(): void {
debugLogger.debug('[ControlDispatcher] Shutting down');

// Remove abort listener to prevent memory leak
if (this.abortHandler) {
this.context.abortSignal.removeEventListener('abort', this.abortHandler);
this.abortHandler = null;
}

// Cancel all incoming requests
for (const [
_requestId,
Expand Down
3 changes: 2 additions & 1 deletion packages/cli/src/nonInteractive/session.ts
Original file line number Diff line number Diff line change
Expand Up @@ -408,7 +408,8 @@ class Session {
private handleInterrupt(): void {
debugLogger.info('[Session] Interrupt requested');
this.abortController.abort();
this.abortController = new AbortController();
// Do not create a new AbortController to prevent listener leaks.
// Subsequent queries will check signal.aborted and fail immediately.
}

private setupSignalHandlers(): void {
Expand Down
15 changes: 13 additions & 2 deletions packages/sdk-typescript/src/query/Query.ts
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,7 @@ export class Query implements AsyncIterable<SDKMessage> {
private firstResultReceivedResolve?: () => void;

private readonly isSingleTurn: boolean;
private abortHandler: (() => void) | null = null;

constructor(
transport: Transport,
Expand Down Expand Up @@ -125,12 +126,13 @@ export class Query implements AsyncIterable<SDKMessage> {
logger.error('Error during abort cleanup:', err);
});
} else {
this.abortController.signal.addEventListener('abort', () => {
this.abortHandler = () => {
this.inputStream.error(new AbortError('Query aborted by user'));
this.close().catch((err) => {
logger.error('Error during abort cleanup:', err);
});
});
};
this.abortController.signal.addEventListener('abort', this.abortHandler);
}

this.initialized = this.initialize();
Expand Down Expand Up @@ -719,6 +721,15 @@ export class Query implements AsyncIterable<SDKMessage> {

this.closed = true;

// Remove abort listener to prevent memory leak
if (this.abortHandler) {
this.abortController.signal.removeEventListener(
'abort',
this.abortHandler,
);
this.abortHandler = null;
}

for (const pending of this.pendingControlRequests.values()) {
pending.abortController.abort();
clearTimeout(pending.timeout);
Expand Down