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
5 changes: 5 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,9 @@

### New

- Add .NET-aligned worker history streaming: hydrate service-selected history before
version checks and replay. History errors produce a Failed completion; shutdown
cancels without submitting completion.
- Add an optional per-call `AbortSignal` to client start and completion waits.
- Add `ConcurrencyOptions` to configure the orchestration, activity, and entity concurrency
hints sent by `TaskHubGrpcWorker` to the backend.
Expand All @@ -16,6 +19,8 @@

### Fixes

- Align worker response cancellation with .NET: `stop()` cancels initial sends as well
as retries and backoff for all work items. Work finishing after stop no longer sends a response.
- Retry worker completion and version-rejection responses on transient gRPC failures, reusing
the computed response without rerunning user code. Bound SDK sends to ten with shutdown-aware backoff.
- Cancel pending client wait RPCs on timeout or cancellation without terminating the orchestration.
Expand Down
9 changes: 6 additions & 3 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -138,9 +138,12 @@ cap before adding 0-20% jitter. Permanent errors and exhausted attempts use the
error logs. Configured gRPC transport retries remain enabled, so ten SDK sends can involve
more than ten network attempts.

`stop()` cancels retry backoff and in-flight retry RPCs. Already-running work can still send
its first response during the existing bounded shutdown wait; user code and metadata
generation are not canceled. Channel retirement and backend lock durations are unchanged.
`stop()` cancels all response RPCs, including the initial send, and retry backoff using
the worker run's signal captured when the work item was dispatched. This applies equally
to inline and streamed orchestrations, activities, entities, and version-failure/rejection
responses. Work finishing after stop cannot send its first response, even after a restart.
User code and metadata generation are not canceled; if metadata finishes after stop,
the response RPC is not started. Channel retirement and backend lock durations are unchanged.
Retries do not guarantee connection recovery, acceptance of expired tokens, or exactly-once execution.

### Reusing orchestration instance IDs
Expand Down
107 changes: 98 additions & 9 deletions packages/durabletask-js/src/worker/task-hub-grpc-worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -114,6 +114,7 @@ export class TaskHubGrpcWorker {
private _stub: stubs.TaskHubSidecarServiceClient | null;
private _logger: Logger;
private _pendingWorkItems: Set<Promise<void>>;
private _historyCancellations: Set<() => void>;
private _shutdownTimeoutMs: number;
private _silentDisconnectTimeoutMs: number;
private _silentDisconnectTimer: ReturnType<typeof setTimeout> | null;
Expand Down Expand Up @@ -214,6 +215,7 @@ export class TaskHubGrpcWorker {
this._stub = null;
this._logger = resolvedLogger ?? new ConsoleLogger();
this._pendingWorkItems = new Set();
this._historyCancellations = new Set();
this._shutdownTimeoutMs = resolvedShutdownTimeoutMs ?? DEFAULT_SHUTDOWN_TIMEOUT_MS;
const silentDisconnectTimeoutMs = resolvedSilentDisconnectTimeoutMs ?? DEFAULT_SILENT_DISCONNECT_TIMEOUT_MS;
if (!Number.isFinite(silentDisconnectTimeoutMs)) {
Expand Down Expand Up @@ -717,6 +719,9 @@ export class TaskHubGrpcWorker {
const responseStream = this._responseStream;
this._stopWorker = true;
this._abortController?.abort();
for (const cancel of this._historyCancellations) {
Comment thread
YunchuWang marked this conversation as resolved.
cancel();
}
this._clearSilentDisconnectTimer();

const streamClosed = responseStream
Expand Down Expand Up @@ -789,6 +794,7 @@ export class TaskHubGrpcWorker {
*/
private _buildGetWorkItemsRequest(): pb.GetWorkItemsRequest {
const request = new pb.GetWorkItemsRequest();
request.setCapabilitiesList([pb.WorkerCapability.WORKER_CAPABILITY_HISTORY_STREAMING]);
request.setMaxconcurrentactivityworkitems(
Math.min(this._concurrency.maximumConcurrentActivityWorkItems, MAX_PROTOCOL_CONCURRENCY),
);
Expand Down Expand Up @@ -908,7 +914,7 @@ export class TaskHubGrpcWorker {
private async _deliverResponse<TReq, TRes>(
method: Parameters<typeof callWithMetadata<TReq, TRes>>[0],
request: TReq,
retrySignal?: AbortSignal,
signal?: AbortSignal,
): Promise<TRes> {
const backoff = new ExponentialBackoff({
initialDelayMs: 200,
Expand All @@ -919,13 +925,7 @@ export class TaskHubGrpcWorker {
});
for (;;) {
try {
// Allow the initial response to drain during shutdown; only retries use the captured run signal.
return await callWithMetadata(
method,
request,
this._metadataGenerator,
backoff.attemptCount === 0 ? undefined : retrySignal,
);
return await callWithMetadata(method, request, this._metadataGenerator, signal);
} catch (error) {
const status = error instanceof Error ? this._getGrpcStatus(error) : undefined;
if (
Expand All @@ -937,7 +937,7 @@ export class TaskHubGrpcWorker {
) {
throw error;
}
await backoff.wait(retrySignal);
await backoff.wait(signal);
}
}
}
Expand All @@ -957,6 +957,59 @@ export class TaskHubGrpcWorker {
});
}

private async _streamOrchestrationHistory(
req: pb.OrchestratorRequest,
stub: stubs.TaskHubSidecarServiceClient,
signal?: AbortSignal,
): Promise<pb.HistoryEvent[]> {
const request = new pb.StreamInstanceHistoryRequest();
request.setInstanceid(req.getInstanceid());
request.setExecutionid(req.getExecutionid());
request.setForworkitemprocessing(true);
return new Promise<pb.HistoryEvent[]>((resolve, reject) => {
const events: pb.HistoryEvent[] = [];
let stream: grpc.ClientReadableStream<pb.HistoryChunk> | undefined;
let settled = false;
const finish = (error?: Error) => {
if (settled) return;
settled = true;
stream?.removeListener("data", onData);
stream?.removeListener("end", onEnd);
stream?.removeListener("error", onError);
stream?.removeListener("close", onClose);
this._historyCancellations.delete(cancel);
stream?.destroy();
if (error) reject(error);
else resolve(events);
};
const onData = (chunk: pb.HistoryChunk) => {
for (const event of chunk.getEventsList()) {
events.push(event);
}
};
const onEnd = () => finish();
const onError = (error: unknown) => finish(error instanceof Error ? error : new Error(String(error)));
const onClose = () => finish(new Error("Orchestration history stream closed before all history was received."));
const cancel = () => {
if (stream) stream.cancel();
else finish(new Error("Orchestration history hydration was cancelled."));
};
// Track metadata acquisition too: a stalled token refresh must not outlive shutdown.
this._historyCancellations.add(cancel);
this._getMetadata()
.then((metadata) => {
if (settled) return;
signal?.throwIfAborted();
stream = stub.streamInstanceHistory(request, metadata);
stream.on("data", onData);
stream.once("end", onEnd);
stream.once("error", onError);
stream.once("close", onClose);
})
.catch(onError);
});
}

/**
* Internal implementation of orchestrator execution.
*/
Expand All @@ -972,6 +1025,42 @@ export class TaskHubGrpcWorker {
throw new Error(`Could not execute the orchestrator as the instanceId was not provided (${instanceId})`);
}

if (req.getRequireshistorystreaming()) {
try {
const pastEvents = await this._streamOrchestrationHistory(req, stub, retrySignal);
retrySignal?.throwIfAborted();
if (
!pastEvents.some((event) => event.hasExecutionstarted()) &&
!req.getNeweventsList().some((event) => event.hasExecutionstarted())
) {
throw new Error("The provided orchestration history was incomplete");
}
Comment thread
YunchuWang marked this conversation as resolved.
req.setPasteventsList(pastEvents);
} catch (e: unknown) {
if (retrySignal?.aborted) return;
const error = e instanceof Error ? e : new Error(String(e));
WorkerLogs.executionError(this._logger, instanceId, error);
const res = new pb.OrchestratorResponse();
res.setInstanceid(instanceId);
res.setCompletiontoken(completionToken);
res.setActionsList([
pbh.newCompleteOrchestrationAction(
-1,
pb.OrchestrationStatus.ORCHESTRATION_STATUS_FAILED,
undefined,
pbh.newFailureDetails(error),
),
]);
try {
await this._deliverResponse(stub.completeOrchestratorTask.bind(stub), res, retrySignal);
} catch (e: unknown) {
const error = e instanceof Error ? e : new Error(String(e));
WorkerLogs.completionError(this._logger, instanceId, error);
}
return;
}
}

// Check version compatibility if versioning is enabled
const versionCheckResult = this._checkVersionCompatibility(req);
if (!versionCheckResult.compatible) {
Expand Down
Loading
Loading