Skip to content
Open
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
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,9 @@

### Fixes

- Bound each worker sidecar hello attempt to 30 seconds, retry failed connections, and cancel
pending hello calls and reconnect delays when the worker stops.

## v0.4.0 (2026-07-31)

### Changes
Expand Down
28 changes: 21 additions & 7 deletions packages/durabletask-js/src/utils/backoff.util.ts
Original file line number Diff line number Diff line change
Expand Up @@ -130,18 +130,32 @@ export class ExponentialBackoff {
* Waits for the current backoff delay, then increments the attempt count
* and calculates the next delay.
*
* @returns Promise that resolves after the delay.
* @param signal Optional signal that cancels the pending delay.
* @returns Promise that resolves after the delay or rejects when aborted.
*/
async wait(): Promise<void> {
async wait(signal?: AbortSignal): Promise<void> {
const delay = this._calculateDelayWithJitter();

await new Promise<void>((resolve) => setTimeout(resolve, delay));
if (signal?.aborted) {
throw signal.reason;
}

await new Promise<void>((resolve, reject) => {
const onAbort = () => {
clearTimeout(timeoutId);
signal?.removeEventListener("abort", onAbort);
reject(signal?.reason);
};
const timeoutId = setTimeout(() => {
signal?.removeEventListener("abort", onAbort);
resolve();
}, delay);

signal?.addEventListener("abort", onAbort, { once: true });
});

this._attemptCount++;
this._currentDelayMs = Math.min(
this._currentDelayMs * this._multiplier,
this._maxDelayMs,
);
this._currentDelayMs = Math.min(this._currentDelayMs * this._multiplier, this._maxDelayMs);
}

/**
Expand Down
125 changes: 101 additions & 24 deletions packages/durabletask-js/src/worker/task-hub-grpc-worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@ import {

/** Default timeout in milliseconds for graceful shutdown. */
const DEFAULT_SHUTDOWN_TIMEOUT_MS = 30000;
const HELLO_TIMEOUT_MS = 30000;

/**
* Options for creating a TaskHubGrpcWorker.
Expand Down Expand Up @@ -88,6 +89,8 @@ export class TaskHubGrpcWorker {
private _backoff: ExponentialBackoff;
private _versioning?: VersioningOptions;
private _workItemFilters?: WorkItemFilters | "auto";
private _abortController: AbortController | null;
private _helloCall: grpc.ClientUnaryCall | null;

/**
* Creates a new TaskHubGrpcWorker instance.
Expand Down Expand Up @@ -179,6 +182,8 @@ export class TaskHubGrpcWorker {
});
this._versioning = resolvedVersioning;
this._workItemFilters = resolvedWorkItemFilters;
this._abortController = null;
this._helloCall = null;
}

/**
Expand All @@ -195,13 +200,28 @@ export class TaskHubGrpcWorker {
* Creates a new gRPC client and retries the worker.
* Properly closes the old client to prevent connection leaks.
*/
private async _createNewClientAndRetry(): Promise<void> {
private async _createNewClientAndRetry(signal?: AbortSignal): Promise<void> {
if (signal?.aborted) {
return;
}

// Close the old stub to prevent connection leaks
if (this._stub) {
this._stub.close();
}

await this._backoff.wait();
try {
await this._backoff.wait(signal);
} catch (err) {
if (signal?.aborted) {
return;
}
throw err;
}

if (signal?.aborted) {
return;
}

const newClient = new GrpcClient(
this._hostAddress,
Expand All @@ -212,8 +232,8 @@ export class TaskHubGrpcWorker {
this._stub = newClient.stub;

// Do not await - run in background
this.internalRunWorker(newClient, true).catch((err) => {
if (!this._stopWorker) {
this.internalRunWorker(newClient, true, signal).catch((err) => {
if (!signal?.aborted && !this._stopWorker) {
WorkerLogs.workerError(this._logger, err);
}
});
Expand Down Expand Up @@ -388,30 +408,76 @@ export class TaskHubGrpcWorker {
throw new Error("The worker is already running.");
}

this._stopWorker = false;
this._backoff.reset();
const abortController = new AbortController();
this._abortController = abortController;
const client = new GrpcClient(this._hostAddress, this._grpcChannelOptions, this._tls, this._grpcChannelCredentials);
this._stub = client.stub;

// Run in background but catch any unhandled errors to prevent unhandled rejections
this.internalRunWorker(client).catch((err) => {
this.internalRunWorker(client, false, abortController.signal).catch((err) => {
// Only log if the worker wasn't stopped intentionally
if (!this._stopWorker) {
if (!abortController.signal.aborted) {
WorkerLogs.workerError(this._logger, err);
}
});

this._isRunning = true;
}

async internalRunWorker(client: GrpcClient, isRetry: boolean = false): Promise<void> {
async internalRunWorker(client: GrpcClient, _isRetry: boolean = false, signal?: AbortSignal): Promise<void> {
try {
// send a "Hello" message to the sidecar to ensure that it's listening
await callWithMetadata(client.stub.hello.bind(client.stub), new Empty(), this._metadataGenerator);
const helloMetadata = await this._getMetadata();
if (signal?.aborted) {
return;
}
await new Promise<void>((resolve, reject) => {
let helloCall: grpc.ClientUnaryCall | undefined = undefined;
let settled = false;
const finish = (err?: Error) => {
if (settled) return;
settled = true;
signal?.removeEventListener("abort", onAbort);
if (this._helloCall === helloCall) {
this._helloCall = null;
}
if (err) {
reject(err);
} else {
resolve();
}
};
const onAbort = () => {
helloCall?.cancel();
finish(signal?.reason);
};
signal?.addEventListener("abort", onAbort, { once: true });
try {
helloCall = client.stub.hello(
new Empty(),
helloMetadata,
{ deadline: new Date(Date.now() + HELLO_TIMEOUT_MS) },
(err) => finish(err ?? undefined),
);
if (!settled) {
this._helloCall = helloCall;
}
} catch (err) {
const normalizedError = err instanceof Error ? err : new Error(String(err));
finish(normalizedError);
}
});

// Reset backoff on successful connection
this._backoff.reset();

// Stream work items from the sidecar (pass metadata for insecure connections)
const metadata = await this._getMetadata();
if (signal?.aborted) {
return;
}
const request = this._buildGetWorkItemsRequest();

const stream = client.stub.getWorkItems(request, metadata);
Expand Down Expand Up @@ -449,7 +515,7 @@ export class TaskHubGrpcWorker {

// Wait for the stream to end or error
stream.on("end", () => {
if (this._stopWorker) {
if (signal?.aborted || this._stopWorker) {
WorkerLogs.streamEnded(this._logger);
stream.removeAllListeners();
stream.destroy();
Expand All @@ -460,16 +526,16 @@ export class TaskHubGrpcWorker {
stream.on("error", () => {}); // Prevent unhandled "error" after cleanup
stream.destroy();
WorkerLogs.streamRetry(this._logger, this._backoff.peekNextDelay());
this._createNewClientAndRetry().catch((retryErr) => {
if (!this._stopWorker) {
this._createNewClientAndRetry(signal).catch((retryErr) => {
if (!signal?.aborted && !this._stopWorker) {
WorkerLogs.workerError(this._logger, retryErr instanceof Error ? retryErr : new Error(String(retryErr)));
}
});
});

stream.on("error", (err: Error) => {
// Ignore cancellation errors when the worker is being stopped intentionally
if (this._stopWorker) {
if (signal?.aborted || this._stopWorker) {
return;
}
WorkerLogs.streamErrorInfo(this._logger, err);
Expand All @@ -482,24 +548,21 @@ export class TaskHubGrpcWorker {
stream.on("error", () => {}); // Prevent unhandled "error" after cleanup
stream.destroy();
WorkerLogs.streamRetry(this._logger, this._backoff.peekNextDelay());
this._createNewClientAndRetry().catch((retryErr) => {
if (!this._stopWorker) {
this._createNewClientAndRetry(signal).catch((retryErr) => {
if (!signal?.aborted && !this._stopWorker) {
WorkerLogs.workerError(this._logger, retryErr instanceof Error ? retryErr : new Error(String(retryErr)));
}
});
});
} catch (err) {
if (this._stopWorker) {
if (signal?.aborted || this._stopWorker) {
// ignoring the error because the worker has been stopped
return;
}
const error = err instanceof Error ? err : new Error(String(err));
WorkerLogs.streamError(this._logger, error);
if (!isRetry) {
throw error;
}
WorkerLogs.connectionRetry(this._logger, this._backoff.peekNextDelay());
await this._createNewClientAndRetry();
await this._createNewClientAndRetry(signal);
return;
}
}
Expand All @@ -513,19 +576,27 @@ export class TaskHubGrpcWorker {
throw new Error("The worker is not running.");
}

const abortController = this._abortController;
const helloCall = this._helloCall;
const responseStream = this._responseStream;
this._stopWorker = true;
abortController?.abort();
if (this._helloCall === helloCall) {
helloCall?.cancel();
this._helloCall = null;
}

// Cancel stream first while error handlers are still attached
// This allows the error handler to suppress CANCELLED errors
this._responseStream?.cancel();
responseStream?.cancel();

// Wait for the stream to react to cancellation using events rather than a fixed delay.
// This avoids race conditions caused by relying on timing alone.
if (this._responseStream) {
if (responseStream) {
try {
await withTimeout(
new Promise<void>((resolve) => {
const stream = this._responseStream!;
const stream = responseStream;
// Any of these events indicates the stream has processed cancellation / is closing.
stream.once("end", resolve);
stream.once("close", resolve);
Expand All @@ -540,8 +611,11 @@ export class TaskHubGrpcWorker {
}

// Now safe to remove listeners and destroy
this._responseStream?.removeAllListeners();
this._responseStream?.destroy();
responseStream?.removeAllListeners();
responseStream?.destroy();
if (this._responseStream === responseStream) {
this._responseStream = null;
}

// Wait for pending work items to complete with timeout
if (this._pendingWorkItems.size > 0) {
Expand All @@ -564,6 +638,9 @@ export class TaskHubGrpcWorker {
this._stub.close();
}
this._isRunning = false;
if (this._abortController === abortController) {
this._abortController = null;
}

// Brief pause to allow gRPC cleanup
// https://github.com/grpc/grpc-node/issues/1563#issuecomment-829483711
Expand Down
23 changes: 23 additions & 0 deletions packages/durabletask-js/test/backoff.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,29 @@ describe("ExponentialBackoff", () => {
await backoff.wait();
expect(backoff.currentDelayMs).toBe(100); // 50 * 10 = 500, capped at 100
});

it("should abort without advancing the backoff state", async () => {
jest.useFakeTimers();
const backoff = new ExponentialBackoff({
initialDelayMs: 1000,
multiplier: 2,
jitterFactor: 0,
});
const controller = new AbortController();
const abortError = new Error("stopped");

try {
const waitPromise = backoff.wait(controller.signal);
controller.abort(abortError);

await expect(waitPromise).rejects.toBe(abortError);
expect(backoff.attemptCount).toBe(0);
expect(backoff.currentDelayMs).toBe(1000);
expect(jest.getTimerCount()).toBe(0);
} finally {
jest.useRealTimers();
}
});
});

describe("reset", () => {
Expand Down
Loading
Loading