From 9f3d205cc058dc38006a92451aafcb7d6ad7acc6 Mon Sep 17 00:00:00 2001 From: wangbill Date: Tue, 8 Sep 2026 09:22:10 -0700 Subject: [PATCH 1/9] feat(worker): hydrate streamed orchestration history before replay Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: e0af01a5-0dfa-4e71-a660-c4186e65d7e0 --- CHANGELOG.md | 3 + .../src/worker/task-hub-grpc-worker.ts | 90 ++++ .../test/worker-history-streaming.spec.ts | 439 ++++++++++++++++++ .../worker-history-streaming.spec.ts | 190 ++++++++ 4 files changed, 722 insertions(+) create mode 100644 packages/durabletask-js/test/worker-history-streaming.spec.ts create mode 100644 test/e2e-azuremanaged/worker-history-streaming.spec.ts diff --git a/CHANGELOG.md b/CHANGELOG.md index 03fab17..d26008e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,9 @@ ### New +- Add worker history streaming: hydrate service-selected history before version checks, + tracing, and replay. Cancel history streams on shutdown and abandon incomplete + work items on transport errors instead of failing the orchestration. - Add `ConcurrencyOptions` to configure the orchestration, activity, and entity concurrency hints sent by `TaskHubGrpcWorker` to the backend. - Add an optional `newVersion` parameter to `OrchestrationContext.continueAsNew()` for version migrations. diff --git a/packages/durabletask-js/src/worker/task-hub-grpc-worker.ts b/packages/durabletask-js/src/worker/task-hub-grpc-worker.ts index a84d34d..03adc44 100644 --- a/packages/durabletask-js/src/worker/task-hub-grpc-worker.ts +++ b/packages/durabletask-js/src/worker/task-hub-grpc-worker.ts @@ -114,6 +114,7 @@ export class TaskHubGrpcWorker { private _stub: stubs.TaskHubSidecarServiceClient | null; private _logger: Logger; private _pendingWorkItems: Set>; + private _historyCancellations: Set<() => void>; private _shutdownTimeoutMs: number; private _silentDisconnectTimeoutMs: number; private _silentDisconnectTimer: ReturnType | null; @@ -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)) { @@ -717,6 +719,9 @@ export class TaskHubGrpcWorker { const responseStream = this._responseStream; this._stopWorker = true; this._abortController?.abort(); + for (const cancel of this._historyCancellations) { + cancel(); + } this._clearSilentDisconnectTimer(); const streamClosed = responseStream @@ -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), ); @@ -919,6 +925,59 @@ export class TaskHubGrpcWorker { }); } + private async _streamOrchestrationHistory( + req: pb.OrchestratorRequest, + stub: stubs.TaskHubSidecarServiceClient, + signal?: AbortSignal, + ): Promise { + const request = new pb.StreamInstanceHistoryRequest(); + request.setInstanceid(req.getInstanceid()); + request.setExecutionid(req.getExecutionid()); + request.setForworkitemprocessing(true); + return new Promise((resolve, reject) => { + const events: pb.HistoryEvent[] = []; + let stream: grpc.ClientReadableStream | 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. */ @@ -933,6 +992,37 @@ export class TaskHubGrpcWorker { throw new Error(`Could not execute the orchestrator as the instanceId was not provided (${instanceId})`); } + if (req.getRequireshistorystreaming()) { + const signal = this._abortController?.signal; + try { + const pastEvents = await this._streamOrchestrationHistory(req, stub, signal); + signal?.throwIfAborted(); + req.setPasteventsList(pastEvents); + } catch (error) { + // Incomplete history is a work-item transport failure, not an orchestration failure. + // Do not replay or persist any actions; the backend can redeliver the work item. + if (!signal?.aborted) { + const abandonRequest = new pb.AbandonOrchestrationTaskRequest(); + abandonRequest.setCompletiontoken(completionToken); + try { + await callWithMetadata( + stub.abandonTaskOrchestratorWorkItem.bind(stub), + abandonRequest, + this._metadataGenerator, + signal, + ); + } catch (abandonError) { + WorkerLogs.completionError( + this._logger, + instanceId, + abandonError instanceof Error ? abandonError : new Error(String(abandonError)), + ); + } + } + throw error; + } + } + // Check version compatibility if versioning is enabled const versionCheckResult = this._checkVersionCompatibility(req); if (!versionCheckResult.compatible) { diff --git a/packages/durabletask-js/test/worker-history-streaming.spec.ts b/packages/durabletask-js/test/worker-history-streaming.spec.ts new file mode 100644 index 0000000..ce20aad --- /dev/null +++ b/packages/durabletask-js/test/worker-history-streaming.spec.ts @@ -0,0 +1,439 @@ +// Copyright (c) Microsoft Corporation. All rights reserved. +// Licensed under the MIT License. + +import * as grpc from "@grpc/grpc-js"; +import { Empty } from "google-protobuf/google/protobuf/empty_pb"; +import { Timestamp } from "google-protobuf/google/protobuf/timestamp_pb"; +import * as otel from "@opentelemetry/api"; +import { BasicTracerProvider, InMemorySpanExporter, SimpleSpanProcessor } from "@opentelemetry/sdk-trace-base"; +import * as pb from "../src/proto/orchestrator_service_pb"; +import * as stubs from "../src/proto/orchestrator_service_grpc_pb"; +import * as pbh from "../src/utils/pb-helper.util"; +import { NoOpLogger } from "../src/types/logger.type"; +import { OrchestrationContext } from "../src/task/context/orchestration-context"; +import { OrchestrationExecutor } from "../src/worker/orchestration-executor"; +import { TaskHubGrpcWorker } from "../src/worker/task-hub-grpc-worker"; +import { VersionFailureStrategy, VersionMatchStrategy } from "../src/worker/versioning-options"; +import { DurableTaskAttributes } from "../src/tracing"; + +type HistoryCall = grpc.ServerWritableStream; + +async function waitFor(predicate: () => boolean): Promise { + const deadline = Date.now() + 3000; + while (!predicate()) { + if (Date.now() > deadline) { + throw new Error("Timed out waiting for worker history streaming"); + } + await new Promise((resolve) => setTimeout(resolve, 5)); + } +} + +function chunk(events: pb.HistoryEvent[]): pb.HistoryChunk { + return new pb.HistoryChunk().setEventsList(events); +} + +describe("Worker history streaming over gRPC", () => { + const instanceId = "history-instance"; + const executionId = "history-execution"; + const exporter = new InMemorySpanExporter(); + const provider = new BasicTracerProvider(); + let server: grpc.Server; + let worker: TaskHubGrpcWorker; + let subscription: grpc.ServerWritableStream | undefined; + let historyCalls: HistoryCall[]; + let responses: pb.OrchestratorResponse[]; + let abandonments: pb.AbandonOrchestrationTaskRequest[]; + let onHistory: (call: HistoryCall) => void; + let historySpy: jest.SpyInstance; + + beforeAll(() => { + provider.addSpanProcessor(new SimpleSpanProcessor(exporter)); + provider.register(); + }); + + afterAll(async () => { + await provider.shutdown(); + otel.trace.disable(); + }); + + beforeEach(async () => { + exporter.reset(); + subscription = undefined; + historyCalls = []; + responses = []; + abandonments = []; + onHistory = (call) => call.end(); + historySpy = jest.spyOn(stubs.TaskHubSidecarServiceClient.prototype, "streamInstanceHistory"); + server = new grpc.Server(); + const service = { + hello: (_call, callback) => callback(null, new Empty()), + getWorkItems: (call) => { + subscription = call; + call.on("cancelled", () => call.end()); + }, + streamInstanceHistory: (call) => { + historyCalls.push(call); + onHistory(call); + }, + completeOrchestratorTask: (call, callback) => { + responses.push(call.request); + callback(null, new pb.CompleteTaskResponse()); + }, + abandonTaskOrchestratorWorkItem: (call, callback) => { + abandonments.push(call.request); + callback(null, new pb.AbandonOrchestrationTaskResponse()); + }, + } satisfies Pick< + stubs.ITaskHubSidecarServiceServer, + | "hello" + | "getWorkItems" + | "streamInstanceHistory" + | "completeOrchestratorTask" + | "abandonTaskOrchestratorWorkItem" + >; + server.addService(stubs.TaskHubSidecarServiceService, service); + const port = await new Promise((resolve, reject) => { + server.bindAsync("127.0.0.1:0", grpc.ServerCredentials.createInsecure(), (error, boundPort) => { + if (error) reject(error); + else resolve(boundPort); + }); + }); + worker = new TaskHubGrpcWorker({ + hostAddress: `127.0.0.1:${port}`, + logger: new NoOpLogger(), + metadataGenerator: async () => { + const metadata = new grpc.Metadata(); + metadata.set("taskhub", "history-test"); + metadata.set("authorization", "test-token"); + return metadata; + }, + }); + }); + + afterEach(async () => { + if (worker["_isRunning"]) { + await worker.stop(); + } + server.forceShutdown(); + jest.restoreAllMocks(); + }); + + async function start(): Promise { + await worker.start(); + await waitFor(() => subscription !== undefined); + } + + function request(streaming = true): pb.OrchestratorRequest { + return new pb.OrchestratorRequest() + .setInstanceid(instanceId) + .setExecutionid(pbh.getStringValue(executionId)) + .setNeweventsList([pbh.newOrchestratorStartedEvent()]) + .setRequireshistorystreaming(streaming); + } + + function send(req: pb.OrchestratorRequest, token = "history-token"): void { + subscription!.write(new pb.WorkItem().setOrchestratorrequest(req).setCompletiontoken(token)); + } + + async function settled(): Promise { + await waitFor(() => worker["_pendingWorkItems"].size === 0); + expect(worker["_historyCancellations"].size).toBe(0); + } + + function expectStreamCleanedUp(): void { + const stream = historySpy.mock.results[0].value as grpc.ClientReadableStream; + for (const event of ["data", "end", "error", "close"]) { + expect(stream.listenerCount(event)).toBe(0); + } + } + + it("advertises only HistoryStreaming, retaining concurrency hints", async () => { + await start(); + expect(subscription!.request.getCapabilitiesList()).toEqual([ + pb.WorkerCapability.WORKER_CAPABILITY_HISTORY_STREAMING, + ]); + expect(subscription!.request.getMaxconcurrentorchestrationworkitems()).toBeGreaterThan(0); + }); + + it.each([false, true])("replays complete ordered history (streaming=%s) before new events", async (streaming) => { + const replayStates: boolean[] = []; + worker.addOrchestrator(async function* orderedHistory(ctx: OrchestrationContext, input: number): AsyncGenerator { + const activity = yield ctx.callActivity("echo", input); + replayStates.push(ctx.isReplaying); + const first = yield ctx.waitForExternalEvent("signal"); + const second = yield ctx.waitForExternalEvent("signal"); + replayStates.push(ctx.isReplaying); + return { input, activity, first, second }; + }); + const pastEvents = [ + pbh.newOrchestratorStartedEvent(), + pbh.newExecutionStartedEvent("orderedHistory", instanceId, "7", undefined, executionId, "1.0"), + pbh.newTaskScheduledEvent(1, "echo", "7"), + pbh.newTaskCompletedEvent(1, "14"), + pbh.newEventRaisedEvent("signal", '"past"'), + ]; + const traceId = "1234567890abcdef1234567890abcdef"; + const parentSpanId = "1234567890abcdef"; + const replaySpanId = "abcdef1234567890"; + pastEvents[1] + .getExecutionstarted()! + .setParenttracecontext(new pb.TraceContext().setTraceparent(`00-${traceId}-${parentSpanId}-01`)); + const newEvents = [pbh.newOrchestratorStartedEvent(), pbh.newEventRaisedEvent("signal", '"new"')]; + const req = request(streaming) + .setNeweventsList(newEvents) + .setOrchestrationtracecontext( + new pb.OrchestrationTraceContext() + .setSpanid(pbh.getStringValue(replaySpanId)) + .setSpanstarttime(Timestamp.fromDate(new Date("2026-01-01T00:00:00Z"))), + ); + // Streaming replaces inline history; this deliberately invalid prefix must never be replayed. + req.setPasteventsList(streaming ? [pbh.newTaskCompletedEvent(999, "0")] : pastEvents); + onHistory = (call) => { + call.write(chunk(pastEvents.slice(0, 2))); + call.write(chunk([])); + call.write(chunk(pastEvents.slice(2))); + call.end(); + }; + worker["_versioning"] = { version: "1.0", matchStrategy: VersionMatchStrategy.Strict }; + const execute = jest.spyOn(OrchestrationExecutor.prototype, "execute"); + await start(); + send(req); + await waitFor(() => responses.length > 0 || abandonments.length > 0); + await settled(); + + expect(abandonments).toHaveLength(0); + expect(responses).toHaveLength(1); + const response = responses[0]; + expect(response.getCompletiontoken()).toBe("history-token"); + expect(response.getInstanceid()).toBe(instanceId); + const completed = response.getActionsList()[0].getCompleteorchestration()!; + expect(completed.getOrchestrationstatus()).toBe(pb.OrchestrationStatus.ORCHESTRATION_STATUS_COMPLETED); + expect(JSON.parse(completed.getResult()!.getValue())).toEqual({ + input: 7, + activity: 14, + first: "past", + second: "new", + }); + expect(replayStates).toEqual([true, false]); + expect(execute.mock.calls[0][3]).toBe(executionId); + expect(execute.mock.calls[0][1].map((event) => event.toObject())).toEqual( + pastEvents.map((event) => event.toObject()), + ); + expect(execute.mock.calls[0][2].map((event) => event.toObject())).toEqual( + newEvents.map((event) => event.toObject()), + ); + expect(response.hasOrchestrationtracecontext()).toBe(true); + expect(response.getOrchestrationtracecontext()!.getSpanid()!.getValue()).toBe(replaySpanId); + const span = exporter.getFinishedSpans().find((item) => item.name === "orchestration:orderedHistory@(1.0)")!; + expect(span.spanContext().traceId).toBe(traceId); + expect(span.parentSpanId).toBe(parentSpanId); + expect(span.attributes[DurableTaskAttributes.REPLAY_SPAN_ID]).toBe(replaySpanId); + expect(historyCalls).toHaveLength(streaming ? 1 : 0); + if (streaming) { + expect(historyCalls[0].request.toObject()).toEqual({ + instanceid: instanceId, + executionid: { value: executionId }, + forworkitemprocessing: true, + }); + expect(historyCalls[0].metadata.get("taskhub")).toEqual(["history-test"]); + expect(historyCalls[0].metadata.get("authorization")).toEqual(["test-token"]); + expectStreamCleanedUp(); + } + }); + + it.each([false, true])("accepts an empty history stream (empty chunk=%s)", async (emptyChunk) => { + worker.addOrchestrator(async function emptyHistory() { + return "empty-history"; + }); + onHistory = (call) => { + if (emptyChunk) call.write(chunk([])); + call.end(); + }; + await start(); + send( + request().setNeweventsList([ + pbh.newOrchestratorStartedEvent(), + pbh.newExecutionStartedEvent("emptyHistory", instanceId), + ]), + ); + await waitFor(() => responses.length > 0); + await settled(); + expect(historyCalls).toHaveLength(1); + expect(responses[0].getActionsList()[0].getCompleteorchestration()!.getResult()!.getValue()).toBe( + '"empty-history"', + ); + expectStreamCleanedUp(); + }); + + it.each([grpc.status.UNAVAILABLE, grpc.status.CANCELLED])( + "abandons incomplete history on gRPC status %s", + async (code) => { + const orchestrator = jest.fn(async function shouldNotExecute() { + return "incorrect"; + }); + worker.addNamedOrchestrator("shouldNotExecute", orchestrator); + const events = [pbh.newOrchestratorStartedEvent(), pbh.newExecutionStartedEvent("shouldNotExecute", instanceId)]; + let chunksReceived = 0; + onHistory = (call) => { + const stream = historySpy.mock.results[0].value as grpc.ClientReadableStream; + stream.once("data", () => { + chunksReceived++; + call.emit("error", Object.assign(new Error("history transport failed"), { code })); + }); + call.write(chunk(events)); + }; + await start(); + send(request().setPasteventsList(events)); + await waitFor(() => abandonments.length > 0 || responses.length > 0); + await settled(); + expect(responses).toHaveLength(0); + expect(chunksReceived).toBe(1); + expect(orchestrator).not.toHaveBeenCalled(); + expect(abandonments.map((item) => item.getCompletiontoken())).toEqual(["history-token"]); + expect(exporter.getFinishedSpans()).toHaveLength(0); + expectStreamCleanedUp(); + + onHistory = (call) => { + call.write(chunk(events)); + call.end(); + }; + send(request(), "redelivery-token"); + await waitFor(() => responses.length > 0); + await settled(); + expect(orchestrator).toHaveBeenCalledTimes(1); + expect(responses[0].getCompletiontoken()).toBe("redelivery-token"); + }, + ); + + it.each([VersionFailureStrategy.Reject, VersionFailureStrategy.Fail])( + "checks streamed versions before dispatch (strategy=%s)", + async (failureStrategy) => { + worker["_versioning"] = { version: "1", matchStrategy: VersionMatchStrategy.Strict, failureStrategy }; + onHistory = (call) => { + call.write( + chunk([pbh.newExecutionStartedEvent("unregistered", instanceId, undefined, undefined, executionId, "2")]), + ); + call.end(); + }; + const execute = jest.spyOn(OrchestrationExecutor.prototype, "execute"); + await start(); + send(request()); + await waitFor(() => responses.length > 0 || abandonments.length > 0); + await settled(); + expect(historyCalls).toHaveLength(1); + expect(execute).not.toHaveBeenCalled(); + if (failureStrategy === VersionFailureStrategy.Reject) { + expect(abandonments[0].getCompletiontoken()).toBe("history-token"); + } else { + expect(responses[0].getActionsList()[0].getCompleteorchestration()!.getFailuredetails()!.getErrortype()).toBe( + "VersionMismatch", + ); + } + }, + ); + + it("cancels outstanding history on stop without executing or leaving pending work", async () => { + let cancelled = false; + onHistory = (call) => { + call.on("cancelled", () => { + cancelled = true; + call.end(); + }); + call.write(chunk([])); + }; + const execute = jest.spyOn(OrchestrationExecutor.prototype, "execute"); + await start(); + send(request()); + await waitFor(() => historyCalls.length > 0 || responses.length > 0); + expect(historyCalls).toHaveLength(1); + expect(worker["_pendingWorkItems"].size).toBe(1); + await worker.stop(); + expect(cancelled).toBe(true); + expect(execute).not.toHaveBeenCalled(); + expect(responses).toHaveLength(0); + expect(worker["_pendingWorkItems"].size).toBe(0); + expectStreamCleanedUp(); + }); + + it("does not open a history stream after stop while metadata was pending", async () => { + await start(); + let releaseMetadata!: (metadata: grpc.Metadata) => void; + const metadata = new Promise((resolve) => { + releaseMetadata = resolve; + }); + const getMetadata = jest.fn(() => metadata); + worker["_metadataGenerator"] = getMetadata; + const execute = jest.spyOn(OrchestrationExecutor.prototype, "execute"); + send(request()); + await waitFor(() => getMetadata.mock.calls.length > 0); + worker["_shutdownTimeoutMs"] = 50; + await worker.stop(); + expect(historyCalls).toHaveLength(0); + expect(execute).not.toHaveBeenCalled(); + expect(responses).toHaveLength(0); + expect(worker["_pendingWorkItems"].size).toBe(0); + expect(worker["_historyCancellations"].size).toBe(0); + releaseMetadata(new grpc.Metadata()); + await new Promise((resolve) => setImmediate(resolve)); + expect(historyCalls).toHaveLength(0); + expect(abandonments).toHaveLength(0); + }); + + it("abandons on metadata failure without using inline history", async () => { + await start(); + worker["_metadataGenerator"] = jest + .fn() + .mockRejectedValueOnce(new Error("token refresh failed")) + .mockResolvedValue(new grpc.Metadata()); + const execute = jest.spyOn(OrchestrationExecutor.prototype, "execute"); + send(request()); + await waitFor(() => abandonments.length > 0); + await settled(); + expect(historyCalls).toHaveLength(0); + expect(execute).not.toHaveBeenCalled(); + expect(responses).toHaveLength(0); + expect(abandonments[0].getCompletiontoken()).toBe("history-token"); + }); + + it("uses the work item's captured stub when the worker channel is replaced", async () => { + worker.addOrchestrator(async function capturedStub() { + return "original-channel"; + }); + onHistory = (call) => { + call.write(chunk([pbh.newOrchestratorStartedEvent(), pbh.newExecutionStartedEvent("capturedStub", instanceId)])); + call.end(); + }; + await start(); + const originalStub = worker["_stub"]!; + worker["_deferStubClose"](originalStub); + worker["_stub"] = new stubs.TaskHubSidecarServiceClient("127.0.0.1:1", grpc.credentials.createInsecure()); + send(request()); + await waitFor(() => responses.length > 0); + await settled(); + expect(historyCalls).toHaveLength(1); + expect(responses[0].getActionsList()[0].getCompleteorchestration()!.getResult()!.getValue()).toBe( + '"original-channel"', + ); + }); + + it("tracks concurrent streams without adding per-work-item abort listeners", async () => { + const warnings: Error[] = []; + const onWarning = (warning: Error) => warnings.push(warning); + process.on("warning", onWarning); + onHistory = (call) => call.on("cancelled", () => call.end()); + try { + await start(); + for (let i = 0; i < 12; i++) { + send(request().setInstanceid(`concurrent-${i}`), `token-${i}`); + } + await waitFor(() => historyCalls.length === 12); + expect(worker["_pendingWorkItems"].size).toBe(12); + await worker.stop(); + expect(worker["_pendingWorkItems"].size).toBe(0); + expect(warnings.filter((warning) => warning.name === "MaxListenersExceededWarning")).toHaveLength(0); + expect(responses).toHaveLength(0); + } finally { + process.removeListener("warning", onWarning); + } + }); +}); diff --git a/test/e2e-azuremanaged/worker-history-streaming.spec.ts b/test/e2e-azuremanaged/worker-history-streaming.spec.ts new file mode 100644 index 0000000..777319c --- /dev/null +++ b/test/e2e-azuremanaged/worker-history-streaming.spec.ts @@ -0,0 +1,190 @@ +// Copyright (c) Microsoft Corporation. All rights reserved. +// Licensed under the MIT License. + +/** + * Opt-in, bounded Azure test. Uses a dedicated task hub; never injects work items. + * DTS_HISTORY_STREAMING_E2E=1 and DTS_CONNECTION_STRING are required. + * DTS_HISTORY_STREAMING_ROUNDS defaults to 4 and is capped at 8. + * + * Each activity input/output is below DTS's 1 MiB payload limit. The accumulated + * history is larger, but the service decides whether to stream it. This test + * fails if it does not observe genuine worker history streaming. + */ +import { randomBytes, randomUUID } from "crypto"; +import * as grpc from "@grpc/grpc-js"; +import { ActivityContext, NoOpLogger, OrchestrationContext, OrchestrationStatus } from "@microsoft/durabletask-js"; +import { + DurableTaskAzureManagedClientBuilder, + DurableTaskAzureManagedWorkerBuilder, +} from "@microsoft/durabletask-js-azuremanaged"; +import * as pb from "../../packages/durabletask-js/src/proto/orchestrator_service_pb"; +import * as stubs from "../../packages/durabletask-js/src/proto/orchestrator_service_grpc_pb"; + +const describeAzure = process.env.DTS_HISTORY_STREAMING_E2E === "1" ? describe : describe.skip; + +describeAzure("Azure worker history streaming negotiation", () => { + it("hydrates service-selected history and completes a replayed durable outcome", async () => { + const connectionString = process.env.DTS_CONNECTION_STRING; + if (!connectionString || !/Endpoint=https:\/\//i.test(connectionString)) { + throw new Error("Set DTS_CONNECTION_STRING to a dedicated Azure HTTPS task hub."); + } + const rounds = Number(process.env.DTS_HISTORY_STREAMING_ROUNDS ?? 4); + if (!Number.isInteger(rounds) || rounds < 1 || rounds > 8) { + throw new Error("DTS_HISTORY_STREAMING_ROUNDS must be an integer from 1 to 8."); + } + + const instanceId = `js-history-streaming-${randomUUID()}`; + const workItems: { instanceId: string; executionId?: string; streaming: boolean; past: number; new: number }[] = []; + const histories: { + instanceId: string; + executionId?: string; + forWorkItemProcessing: boolean; + chunks: number; + events: number; + bytes: number; + ended: boolean; + error?: string; + }[] = []; + const getWorkItems = stubs.TaskHubSidecarServiceClient.prototype.getWorkItems; + const streamHistory = stubs.TaskHubSidecarServiceClient.prototype.streamInstanceHistory; + jest.spyOn(stubs.TaskHubSidecarServiceClient.prototype, "getWorkItems").mockImplementation(function ( + this: stubs.TaskHubSidecarServiceClient, + request: pb.GetWorkItemsRequest, + metadata?: grpc.Metadata, + options?: Partial, + ) { + const stream = getWorkItems.call(this, request, metadata, options); + stream.on("data", (item: pb.WorkItem) => { + const req = item.getOrchestratorrequest(); + if (req?.getInstanceid() === instanceId) { + workItems.push({ + instanceId: req.getInstanceid(), + executionId: req.getExecutionid()?.getValue(), + streaming: req.getRequireshistorystreaming(), + past: req.getPasteventsList().length, + new: req.getNeweventsList().length, + }); + } + }); + return stream; + }); + jest.spyOn(stubs.TaskHubSidecarServiceClient.prototype, "streamInstanceHistory").mockImplementation(function ( + this: stubs.TaskHubSidecarServiceClient, + request: pb.StreamInstanceHistoryRequest, + metadata?: grpc.Metadata, + options?: Partial, + ) { + const stream = streamHistory.call(this, request, metadata, options); + const observation: (typeof histories)[number] = { + instanceId: request.getInstanceid(), + executionId: request.getExecutionid()?.getValue(), + forWorkItemProcessing: request.getForworkitemprocessing(), + chunks: 0, + events: 0, + bytes: 0, + ended: false, + }; + histories.push(observation); + const onData = (historyChunk: pb.HistoryChunk) => { + observation.chunks++; + observation.events += historyChunk.getEventsList().length; + for (const event of historyChunk.getEventsList()) { + observation.bytes += event.serializeBinary().length; + } + }; + stream.on("data", onData); + stream.once("end", () => { + observation.ended = true; + stream.removeListener("data", onData); + }); + stream.once("error", (error: grpc.ServiceError) => { + observation.error = `${error.code}: ${error.details}`; + stream.removeListener("data", onData); + }); + return stream; + }); + + const client = new DurableTaskAzureManagedClientBuilder() + .connectionString(connectionString) + .logger(new NoOpLogger()) + .build(); + const worker = new DurableTaskAzureManagedWorkerBuilder() + .connectionString(connectionString) + .logger(new NoOpLogger()) + .build(); + type Echo = { payload: string; index: number }; + const echo = async function historyStreamingEcho(_ctx: ActivityContext, input: Echo): Promise { + return input; + }; + const orchestrator = async function* boundedStreamingHistory( + ctx: OrchestrationContext, + input: { payload: string; rounds: number }, + ): AsyncGenerator { + const indices: number[] = []; + for (let index = 0; index < input.rounds; index++) { + const result: Echo = yield ctx.callActivity(echo, { payload: input.payload, index }); + if (result.payload !== input.payload || result.index !== index) { + throw new Error("Replayed activity payload or ordering was corrupted."); + } + indices.push(result.index); + } + // Ensure the last activity result is persisted into past history for another replay. + yield ctx.createTimer(new Date(ctx.currentUtcDateTime.getTime() + 1000)); + return { indices, payloadBytes: input.payload.length }; + }; + worker.addActivity(echo); + worker.addOrchestrator(orchestrator); + const payload = randomBytes(576 * 1024).toString("base64"); + let started = false; + let scheduled = false; + let completed = false; + let outcome: unknown; + try { + await worker.start(); + started = true; + await client.scheduleNewOrchestration(orchestrator, { payload, rounds }, { instanceId }); + scheduled = true; + const state = await client.waitForOrchestrationCompletion(instanceId, undefined, 90); + completed = state?.runtimeStatus === OrchestrationStatus.COMPLETED; + outcome = { status: state?.runtimeStatus, output: state?.serializedOutput, failure: state?.failureDetails }; + expect(completed).toBe(true); + expect(JSON.parse(state!.serializedOutput!)).toEqual({ + indices: Array.from({ length: rounds }, (_, i) => i), + payloadBytes: payload.length, + }); + const streamedWorkItems = workItems.filter((item) => item.streaming); + if (streamedWorkItems.length === 0) { + throw new Error(`Azure did not request history streaming within the bounded ${rounds}-activity workload.`); + } + expect(histories).toHaveLength(streamedWorkItems.length); + for (const history of histories) { + expect(history.instanceId).toBe(instanceId); + expect(history.executionId).toBeTruthy(); + expect(streamedWorkItems.some((item) => item.executionId === history.executionId)).toBe(true); + expect(history.forWorkItemProcessing).toBe(true); + expect(history.ended).toBe(true); + expect(history.error).toBeUndefined(); + expect(history.events).toBeGreaterThan(0); + } + expect(histories.some((history) => history.chunks > 1)).toBe(true); + } finally { + console.log( + "HISTORY_STREAMING_EVIDENCE", + JSON.stringify({ instanceId, rounds, payloadBytes: payload.length, workItems, histories, outcome }), + ); + try { + if (scheduled) { + if (!completed) { + await client.terminateOrchestration(instanceId, "history-streaming-test-cleanup"); + await client.waitForOrchestrationCompletion(instanceId, undefined, 20); + } + await client.purgeOrchestration(instanceId); + } + } finally { + if (started) await worker.stop(); + await client.stop(); + jest.restoreAllMocks(); + } + } + }, 150000); +}); From c75eef7f4a6fb6269da8c40e2785d5f6a1065b16 Mon Sep 17 00:00:00 2001 From: wangbill Date: Tue, 8 Sep 2026 09:28:46 -0700 Subject: [PATCH 2/9] test(worker): use proven minimal Azure history workload Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: e0af01a5-0dfa-4e71-a660-c4186e65d7e0 --- test/e2e-azuremanaged/worker-history-streaming.spec.ts | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/test/e2e-azuremanaged/worker-history-streaming.spec.ts b/test/e2e-azuremanaged/worker-history-streaming.spec.ts index 777319c..febb4cf 100644 --- a/test/e2e-azuremanaged/worker-history-streaming.spec.ts +++ b/test/e2e-azuremanaged/worker-history-streaming.spec.ts @@ -4,7 +4,7 @@ /** * Opt-in, bounded Azure test. Uses a dedicated task hub; never injects work items. * DTS_HISTORY_STREAMING_E2E=1 and DTS_CONNECTION_STRING are required. - * DTS_HISTORY_STREAMING_ROUNDS defaults to 4 and is capped at 8. + * DTS_HISTORY_STREAMING_ROUNDS defaults to 1 and is capped at 8. * * Each activity input/output is below DTS's 1 MiB payload limit. The accumulated * history is larger, but the service decides whether to stream it. This test @@ -28,7 +28,7 @@ describeAzure("Azure worker history streaming negotiation", () => { if (!connectionString || !/Endpoint=https:\/\//i.test(connectionString)) { throw new Error("Set DTS_CONNECTION_STRING to a dedicated Azure HTTPS task hub."); } - const rounds = Number(process.env.DTS_HISTORY_STREAMING_ROUNDS ?? 4); + const rounds = Number(process.env.DTS_HISTORY_STREAMING_ROUNDS ?? 1); if (!Number.isInteger(rounds) || rounds < 1 || rounds > 8) { throw new Error("DTS_HISTORY_STREAMING_ROUNDS must be an integer from 1 to 8."); } From 2e0cd53deeef507251ec4204bc513a29d269cf62 Mon Sep 17 00:00:00 2001 From: wangbill Date: Tue, 8 Sep 2026 09:42:34 -0700 Subject: [PATCH 3/9] fix(worker): abandon streamed history missing execution start Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: e0af01a5-0dfa-4e71-a660-c4186e65d7e0 --- .../src/worker/task-hub-grpc-worker.ts | 6 ++ .../test/worker-history-streaming.spec.ts | 98 ++++++++++++++----- 2 files changed, 81 insertions(+), 23 deletions(-) diff --git a/packages/durabletask-js/src/worker/task-hub-grpc-worker.ts b/packages/durabletask-js/src/worker/task-hub-grpc-worker.ts index 03adc44..ee70fed 100644 --- a/packages/durabletask-js/src/worker/task-hub-grpc-worker.ts +++ b/packages/durabletask-js/src/worker/task-hub-grpc-worker.ts @@ -997,6 +997,12 @@ export class TaskHubGrpcWorker { try { const pastEvents = await this._streamOrchestrationHistory(req, stub, signal); signal?.throwIfAborted(); + if ( + !pastEvents.some((event) => event.hasExecutionstarted()) && + !req.getNeweventsList().some((event) => event.hasExecutionstarted()) + ) { + throw new Error("The provided orchestration history was incomplete (missing ExecutionStarted)."); + } req.setPasteventsList(pastEvents); } catch (error) { // Incomplete history is a work-item transport failure, not an orchestration failure. diff --git a/packages/durabletask-js/test/worker-history-streaming.spec.ts b/packages/durabletask-js/test/worker-history-streaming.spec.ts index ce20aad..6a544b3 100644 --- a/packages/durabletask-js/test/worker-history-streaming.spec.ts +++ b/packages/durabletask-js/test/worker-history-streaming.spec.ts @@ -265,6 +265,52 @@ describe("Worker history streaming over gRPC", () => { expectStreamCleanedUp(); }); + it.each<[VersionMatchStrategy, boolean]>([ + [VersionMatchStrategy.None, false], + [VersionMatchStrategy.Strict, false], + [VersionMatchStrategy.None, true], + [VersionMatchStrategy.Strict, true], + ])("abandons OK history missing ExecutionStarted (strategy=%s, nonempty=%s)", async (matchStrategy, nonempty) => { + worker["_versioning"] = { version: "1", matchStrategy, failureStrategy: VersionFailureStrategy.Fail }; + worker.addOrchestrator(async function* resumedHistory(ctx: OrchestrationContext, input: number): AsyncGenerator { + return yield ctx.callActivity("echo", input); + }); + const pastEvents = [ + pbh.newOrchestratorStartedEvent(), + pbh.newExecutionStartedEvent("resumedHistory", instanceId, "7", undefined, executionId, "1"), + pbh.newTaskScheduledEvent(1, "echo", "7"), + ]; + const req = request() + .setPasteventsList(pastEvents) + .setNeweventsList([pbh.newOrchestratorStartedEvent(), pbh.newTaskCompletedEvent(1, "14")]); + onHistory = (call) => { + if (nonempty) call.write(chunk([pastEvents[2]])); + call.end(); + }; + const execute = jest.spyOn(OrchestrationExecutor.prototype, "execute"); + await start(); + send(req); + await waitFor(() => responses.length > 0 || abandonments.length > 0); + await settled(); + expect(abandonments.map((item) => item.getCompletiontoken())).toEqual(["history-token"]); + expect(responses).toHaveLength(0); + expect(execute).not.toHaveBeenCalled(); + expect(exporter.getFinishedSpans()).toHaveLength(0); + expectStreamCleanedUp(); + + onHistory = (call) => { + call.write(chunk(pastEvents)); + call.end(); + }; + send(req, "complete-history-redelivery"); + await waitFor(() => responses.length > 0); + await settled(); + expect(responses[0].getCompletiontoken()).toBe("complete-history-redelivery"); + const completed = responses[0].getActionsList()[0].getCompleteorchestration()!; + expect(completed.getOrchestrationstatus()).toBe(pb.OrchestrationStatus.ORCHESTRATION_STATUS_COMPLETED); + expect(completed.getResult()!.getValue()).toBe("14"); + }); + it.each([grpc.status.UNAVAILABLE, grpc.status.CANCELLED])( "abandons incomplete history on gRPC status %s", async (code) => { @@ -355,29 +401,35 @@ describe("Worker history streaming over gRPC", () => { expectStreamCleanedUp(); }); - it("does not open a history stream after stop while metadata was pending", async () => { - await start(); - let releaseMetadata!: (metadata: grpc.Metadata) => void; - const metadata = new Promise((resolve) => { - releaseMetadata = resolve; - }); - const getMetadata = jest.fn(() => metadata); - worker["_metadataGenerator"] = getMetadata; - const execute = jest.spyOn(OrchestrationExecutor.prototype, "execute"); - send(request()); - await waitFor(() => getMetadata.mock.calls.length > 0); - worker["_shutdownTimeoutMs"] = 50; - await worker.stop(); - expect(historyCalls).toHaveLength(0); - expect(execute).not.toHaveBeenCalled(); - expect(responses).toHaveLength(0); - expect(worker["_pendingWorkItems"].size).toBe(0); - expect(worker["_historyCancellations"].size).toBe(0); - releaseMetadata(new grpc.Metadata()); - await new Promise((resolve) => setImmediate(resolve)); - expect(historyCalls).toHaveLength(0); - expect(abandonments).toHaveLength(0); - }); + it.each([false, true])( + "stops pending metadata hydration before late settlement (reject=%s)", + async (rejectMetadata) => { + await start(); + let releaseMetadata!: () => void; + const metadata = new Promise((resolve, reject) => { + releaseMetadata = () => { + if (rejectMetadata) reject(new Error("late metadata failure")); + else resolve(new grpc.Metadata()); + }; + }); + const getMetadata = jest.fn(() => metadata); + worker["_metadataGenerator"] = getMetadata; + const execute = jest.spyOn(OrchestrationExecutor.prototype, "execute"); + send(request()); + await waitFor(() => getMetadata.mock.calls.length > 0); + worker["_shutdownTimeoutMs"] = 50; + await worker.stop(); + expect(historyCalls).toHaveLength(0); + expect(execute).not.toHaveBeenCalled(); + expect(responses).toHaveLength(0); + expect(worker["_pendingWorkItems"].size).toBe(0); + expect(worker["_historyCancellations"].size).toBe(0); + releaseMetadata(); + await new Promise((resolve) => setImmediate(resolve)); + expect(historyCalls).toHaveLength(0); + expect(abandonments).toHaveLength(0); + }, + ); it("abandons on metadata failure without using inline history", async () => { await start(); From a94d574fcc54ed574f2ac0bc16ea67af024f2e16 Mon Sep 17 00:00:00 2001 From: wangbill Date: Tue, 8 Sep 2026 10:03:31 -0700 Subject: [PATCH 4/9] test(worker): include streaming negotiation in existing CI history suite Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: e0af01a5-0dfa-4e71-a660-c4186e65d7e0 --- test/e2e-azuremanaged/history.spec.ts | 169 +++++++++++++++- .../worker-history-streaming.spec.ts | 190 ------------------ 2 files changed, 168 insertions(+), 191 deletions(-) delete mode 100644 test/e2e-azuremanaged/worker-history-streaming.spec.ts diff --git a/test/e2e-azuremanaged/history.spec.ts b/test/e2e-azuremanaged/history.spec.ts index 2a35ad6..bbd551b 100644 --- a/test/e2e-azuremanaged/history.spec.ts +++ b/test/e2e-azuremanaged/history.spec.ts @@ -2,18 +2,22 @@ // Licensed under the MIT License. /** - * E2E tests for getOrchestrationHistory in Durable Task Scheduler (DTS). + * E2E tests for client history queries and worker history streaming in DTS. * * Environment variables (choose one): * - DTS_CONNECTION_STRING: Full connection string (e.g., "Endpoint=https://...;Authentication=DefaultAzure;TaskHub=...") * OR * - ENDPOINT: The endpoint for the DTS emulator (default: http://localhost:8080) * - TASKHUB: The task hub name (default: default) + * - DTS_HISTORY_STREAMING_ROUNDS: Bounded streaming workload (default: 1, maximum: 8) */ +import { randomBytes, randomUUID } from "crypto"; +import * as grpc from "@grpc/grpc-js"; import { TaskHubGrpcClient, TaskHubGrpcWorker, + OrchestrationStatus, getName, whenAll, ActivityContext, @@ -35,6 +39,8 @@ import { DurableTaskAzureManagedClientBuilder, DurableTaskAzureManagedWorkerBuilder, } from "@microsoft/durabletask-js-azuremanaged"; +import * as pb from "../../packages/durabletask-js/src/proto/orchestrator_service_pb"; +import * as stubs from "../../packages/durabletask-js/src/proto/orchestrator_service_grpc_pb"; // Read environment variables const connectionString = process.env.DTS_CONNECTION_STRING; @@ -59,6 +65,167 @@ function createWorker(): TaskHubGrpcWorker { .build(); } +describe("Worker history streaming negotiation", () => { + // Keep the bounded scenario in this CI-selected file so Node 22/24 exercise it automatically. + it("hydrates service-selected history and completes a replayed durable outcome", async () => { + const rounds = Number(process.env.DTS_HISTORY_STREAMING_ROUNDS ?? 1); + if (!Number.isInteger(rounds) || rounds < 1 || rounds > 8) { + throw new Error("DTS_HISTORY_STREAMING_ROUNDS must be an integer from 1 to 8."); + } + + const instanceId = `js-history-streaming-${randomUUID()}`; + const workItems: { instanceId: string; executionId?: string; streaming: boolean; past: number; new: number }[] = []; + const histories: { + instanceId: string; + executionId?: string; + forWorkItemProcessing: boolean; + chunks: number; + events: number; + bytes: number; + ended: boolean; + error?: string; + }[] = []; + const getWorkItems = stubs.TaskHubSidecarServiceClient.prototype.getWorkItems; + const streamHistory = stubs.TaskHubSidecarServiceClient.prototype.streamInstanceHistory; + jest.spyOn(stubs.TaskHubSidecarServiceClient.prototype, "getWorkItems").mockImplementation(function ( + this: stubs.TaskHubSidecarServiceClient, + request: pb.GetWorkItemsRequest, + metadata?: grpc.Metadata, + options?: Partial, + ) { + const stream = getWorkItems.call(this, request, metadata, options); + stream.on("data", (item: pb.WorkItem) => { + const req = item.getOrchestratorrequest(); + if (req?.getInstanceid() === instanceId) { + workItems.push({ + instanceId: req.getInstanceid(), + executionId: req.getExecutionid()?.getValue(), + streaming: req.getRequireshistorystreaming(), + past: req.getPasteventsList().length, + new: req.getNeweventsList().length, + }); + } + }); + return stream; + }); + jest.spyOn(stubs.TaskHubSidecarServiceClient.prototype, "streamInstanceHistory").mockImplementation(function ( + this: stubs.TaskHubSidecarServiceClient, + request: pb.StreamInstanceHistoryRequest, + metadata?: grpc.Metadata, + options?: Partial, + ) { + const stream = streamHistory.call(this, request, metadata, options); + const observation: (typeof histories)[number] = { + instanceId: request.getInstanceid(), + executionId: request.getExecutionid()?.getValue(), + forWorkItemProcessing: request.getForworkitemprocessing(), + chunks: 0, + events: 0, + bytes: 0, + ended: false, + }; + histories.push(observation); + const onData = (historyChunk: pb.HistoryChunk) => { + observation.chunks++; + observation.events += historyChunk.getEventsList().length; + for (const event of historyChunk.getEventsList()) { + observation.bytes += event.serializeBinary().length; + } + }; + stream.on("data", onData); + stream.once("end", () => { + observation.ended = true; + stream.removeListener("data", onData); + }); + stream.once("error", (error: grpc.ServiceError) => { + observation.error = `${error.code}: ${error.details}`; + stream.removeListener("data", onData); + }); + return stream; + }); + + const client = createClient(); + const worker = createWorker(); + type Echo = { payload: string; index: number }; + const echo = async function historyStreamingEcho(_ctx: ActivityContext, input: Echo): Promise { + return input; + }; + const orchestrator = async function* boundedStreamingHistory( + ctx: OrchestrationContext, + input: { payload: string; rounds: number }, + ): AsyncGenerator { + const indices: number[] = []; + for (let index = 0; index < input.rounds; index++) { + const result: Echo = yield ctx.callActivity(echo, { payload: input.payload, index }); + if (result.payload !== input.payload || result.index !== index) { + throw new Error("Replayed activity payload or ordering was corrupted."); + } + indices.push(result.index); + } + // Ensure the last activity result is persisted into past history for another replay. + yield ctx.createTimer(new Date(ctx.currentUtcDateTime.getTime() + 1000)); + return { indices, payloadBytes: input.payload.length }; + }; + worker.addActivity(echo); + worker.addOrchestrator(orchestrator); + // Keep each payload below 1 MiB while accumulating a larger orchestration history. + const payload = randomBytes(576 * 1024).toString("base64"); + let started = false; + let scheduled = false; + let completed = false; + let outcome: unknown; + try { + await worker.start(); + started = true; + await client.scheduleNewOrchestration(orchestrator, { payload, rounds }, { instanceId }); + scheduled = true; + const state = await client.waitForOrchestrationCompletion(instanceId, undefined, 90); + completed = state?.runtimeStatus === OrchestrationStatus.COMPLETED; + outcome = { status: state?.runtimeStatus, output: state?.serializedOutput, failure: state?.failureDetails }; + expect(completed).toBe(true); + expect(JSON.parse(state!.serializedOutput!)).toEqual({ + indices: Array.from({ length: rounds }, (_, i) => i), + payloadBytes: payload.length, + }); + const streamedWorkItems = workItems.filter((item) => item.streaming); + if (streamedWorkItems.length === 0) { + throw new Error( + `The service did not request history streaming within the bounded ${rounds}-activity workload.`, + ); + } + expect(histories).toHaveLength(streamedWorkItems.length); + for (const history of histories) { + expect(history.instanceId).toBe(instanceId); + expect(history.executionId).toBeTruthy(); + expect(streamedWorkItems.some((item) => item.executionId === history.executionId)).toBe(true); + expect(history.forWorkItemProcessing).toBe(true); + expect(history.ended).toBe(true); + expect(history.error).toBeUndefined(); + expect(history.events).toBeGreaterThan(0); + } + expect(histories.some((history) => history.chunks > 1)).toBe(true); + } finally { + console.log( + "HISTORY_STREAMING_EVIDENCE", + JSON.stringify({ instanceId, rounds, payloadBytes: payload.length, workItems, histories, outcome }), + ); + try { + if (scheduled) { + if (!completed) { + await client.terminateOrchestration(instanceId, "history-streaming-test-cleanup"); + await client.waitForOrchestrationCompletion(instanceId, undefined, 20); + } + await client.purgeOrchestration(instanceId); + } + } finally { + if (started) await worker.stop(); + await client.stop(); + jest.restoreAllMocks(); + } + } + }, 150000); +}); + describe("getOrchestrationHistory E2E Tests", () => { let taskHubClient: TaskHubGrpcClient; let taskHubWorker: TaskHubGrpcWorker; diff --git a/test/e2e-azuremanaged/worker-history-streaming.spec.ts b/test/e2e-azuremanaged/worker-history-streaming.spec.ts deleted file mode 100644 index febb4cf..0000000 --- a/test/e2e-azuremanaged/worker-history-streaming.spec.ts +++ /dev/null @@ -1,190 +0,0 @@ -// Copyright (c) Microsoft Corporation. All rights reserved. -// Licensed under the MIT License. - -/** - * Opt-in, bounded Azure test. Uses a dedicated task hub; never injects work items. - * DTS_HISTORY_STREAMING_E2E=1 and DTS_CONNECTION_STRING are required. - * DTS_HISTORY_STREAMING_ROUNDS defaults to 1 and is capped at 8. - * - * Each activity input/output is below DTS's 1 MiB payload limit. The accumulated - * history is larger, but the service decides whether to stream it. This test - * fails if it does not observe genuine worker history streaming. - */ -import { randomBytes, randomUUID } from "crypto"; -import * as grpc from "@grpc/grpc-js"; -import { ActivityContext, NoOpLogger, OrchestrationContext, OrchestrationStatus } from "@microsoft/durabletask-js"; -import { - DurableTaskAzureManagedClientBuilder, - DurableTaskAzureManagedWorkerBuilder, -} from "@microsoft/durabletask-js-azuremanaged"; -import * as pb from "../../packages/durabletask-js/src/proto/orchestrator_service_pb"; -import * as stubs from "../../packages/durabletask-js/src/proto/orchestrator_service_grpc_pb"; - -const describeAzure = process.env.DTS_HISTORY_STREAMING_E2E === "1" ? describe : describe.skip; - -describeAzure("Azure worker history streaming negotiation", () => { - it("hydrates service-selected history and completes a replayed durable outcome", async () => { - const connectionString = process.env.DTS_CONNECTION_STRING; - if (!connectionString || !/Endpoint=https:\/\//i.test(connectionString)) { - throw new Error("Set DTS_CONNECTION_STRING to a dedicated Azure HTTPS task hub."); - } - const rounds = Number(process.env.DTS_HISTORY_STREAMING_ROUNDS ?? 1); - if (!Number.isInteger(rounds) || rounds < 1 || rounds > 8) { - throw new Error("DTS_HISTORY_STREAMING_ROUNDS must be an integer from 1 to 8."); - } - - const instanceId = `js-history-streaming-${randomUUID()}`; - const workItems: { instanceId: string; executionId?: string; streaming: boolean; past: number; new: number }[] = []; - const histories: { - instanceId: string; - executionId?: string; - forWorkItemProcessing: boolean; - chunks: number; - events: number; - bytes: number; - ended: boolean; - error?: string; - }[] = []; - const getWorkItems = stubs.TaskHubSidecarServiceClient.prototype.getWorkItems; - const streamHistory = stubs.TaskHubSidecarServiceClient.prototype.streamInstanceHistory; - jest.spyOn(stubs.TaskHubSidecarServiceClient.prototype, "getWorkItems").mockImplementation(function ( - this: stubs.TaskHubSidecarServiceClient, - request: pb.GetWorkItemsRequest, - metadata?: grpc.Metadata, - options?: Partial, - ) { - const stream = getWorkItems.call(this, request, metadata, options); - stream.on("data", (item: pb.WorkItem) => { - const req = item.getOrchestratorrequest(); - if (req?.getInstanceid() === instanceId) { - workItems.push({ - instanceId: req.getInstanceid(), - executionId: req.getExecutionid()?.getValue(), - streaming: req.getRequireshistorystreaming(), - past: req.getPasteventsList().length, - new: req.getNeweventsList().length, - }); - } - }); - return stream; - }); - jest.spyOn(stubs.TaskHubSidecarServiceClient.prototype, "streamInstanceHistory").mockImplementation(function ( - this: stubs.TaskHubSidecarServiceClient, - request: pb.StreamInstanceHistoryRequest, - metadata?: grpc.Metadata, - options?: Partial, - ) { - const stream = streamHistory.call(this, request, metadata, options); - const observation: (typeof histories)[number] = { - instanceId: request.getInstanceid(), - executionId: request.getExecutionid()?.getValue(), - forWorkItemProcessing: request.getForworkitemprocessing(), - chunks: 0, - events: 0, - bytes: 0, - ended: false, - }; - histories.push(observation); - const onData = (historyChunk: pb.HistoryChunk) => { - observation.chunks++; - observation.events += historyChunk.getEventsList().length; - for (const event of historyChunk.getEventsList()) { - observation.bytes += event.serializeBinary().length; - } - }; - stream.on("data", onData); - stream.once("end", () => { - observation.ended = true; - stream.removeListener("data", onData); - }); - stream.once("error", (error: grpc.ServiceError) => { - observation.error = `${error.code}: ${error.details}`; - stream.removeListener("data", onData); - }); - return stream; - }); - - const client = new DurableTaskAzureManagedClientBuilder() - .connectionString(connectionString) - .logger(new NoOpLogger()) - .build(); - const worker = new DurableTaskAzureManagedWorkerBuilder() - .connectionString(connectionString) - .logger(new NoOpLogger()) - .build(); - type Echo = { payload: string; index: number }; - const echo = async function historyStreamingEcho(_ctx: ActivityContext, input: Echo): Promise { - return input; - }; - const orchestrator = async function* boundedStreamingHistory( - ctx: OrchestrationContext, - input: { payload: string; rounds: number }, - ): AsyncGenerator { - const indices: number[] = []; - for (let index = 0; index < input.rounds; index++) { - const result: Echo = yield ctx.callActivity(echo, { payload: input.payload, index }); - if (result.payload !== input.payload || result.index !== index) { - throw new Error("Replayed activity payload or ordering was corrupted."); - } - indices.push(result.index); - } - // Ensure the last activity result is persisted into past history for another replay. - yield ctx.createTimer(new Date(ctx.currentUtcDateTime.getTime() + 1000)); - return { indices, payloadBytes: input.payload.length }; - }; - worker.addActivity(echo); - worker.addOrchestrator(orchestrator); - const payload = randomBytes(576 * 1024).toString("base64"); - let started = false; - let scheduled = false; - let completed = false; - let outcome: unknown; - try { - await worker.start(); - started = true; - await client.scheduleNewOrchestration(orchestrator, { payload, rounds }, { instanceId }); - scheduled = true; - const state = await client.waitForOrchestrationCompletion(instanceId, undefined, 90); - completed = state?.runtimeStatus === OrchestrationStatus.COMPLETED; - outcome = { status: state?.runtimeStatus, output: state?.serializedOutput, failure: state?.failureDetails }; - expect(completed).toBe(true); - expect(JSON.parse(state!.serializedOutput!)).toEqual({ - indices: Array.from({ length: rounds }, (_, i) => i), - payloadBytes: payload.length, - }); - const streamedWorkItems = workItems.filter((item) => item.streaming); - if (streamedWorkItems.length === 0) { - throw new Error(`Azure did not request history streaming within the bounded ${rounds}-activity workload.`); - } - expect(histories).toHaveLength(streamedWorkItems.length); - for (const history of histories) { - expect(history.instanceId).toBe(instanceId); - expect(history.executionId).toBeTruthy(); - expect(streamedWorkItems.some((item) => item.executionId === history.executionId)).toBe(true); - expect(history.forWorkItemProcessing).toBe(true); - expect(history.ended).toBe(true); - expect(history.error).toBeUndefined(); - expect(history.events).toBeGreaterThan(0); - } - expect(histories.some((history) => history.chunks > 1)).toBe(true); - } finally { - console.log( - "HISTORY_STREAMING_EVIDENCE", - JSON.stringify({ instanceId, rounds, payloadBytes: payload.length, workItems, histories, outcome }), - ); - try { - if (scheduled) { - if (!completed) { - await client.terminateOrchestration(instanceId, "history-streaming-test-cleanup"); - await client.waitForOrchestrationCompletion(instanceId, undefined, 20); - } - await client.purgeOrchestration(instanceId); - } - } finally { - if (started) await worker.stop(); - await client.stop(); - jest.restoreAllMocks(); - } - } - }, 150000); -}); From 0cbd520284a99f4906ddbc83737baf758d97e112 Mon Sep 17 00:00:00 2001 From: wangbill Date: Tue, 8 Sep 2026 10:22:46 -0700 Subject: [PATCH 5/9] refactor(worker): share orchestration abandonment delivery Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: e0af01a5-0dfa-4e71-a660-c4186e65d7e0 Copilot-Session: 7b853e01-1e42-41c2-92a7-e2bbab87d215 --- .../src/worker/task-hub-grpc-worker.ts | 27 +++++------- .../test/worker-history-streaming.spec.ts | 44 ++++++++++++++++++- 2 files changed, 55 insertions(+), 16 deletions(-) diff --git a/packages/durabletask-js/src/worker/task-hub-grpc-worker.ts b/packages/durabletask-js/src/worker/task-hub-grpc-worker.ts index ee70fed..8d6f1b4 100644 --- a/packages/durabletask-js/src/worker/task-hub-grpc-worker.ts +++ b/packages/durabletask-js/src/worker/task-hub-grpc-worker.ts @@ -925,6 +925,16 @@ export class TaskHubGrpcWorker { }); } + private async _abandonOrchestrationWorkItem( + stub: stubs.TaskHubSidecarServiceClient, + completionToken: string, + signal?: AbortSignal, + ): Promise { + const request = new pb.AbandonOrchestrationTaskRequest(); + request.setCompletiontoken(completionToken); + await callWithMetadata(stub.abandonTaskOrchestratorWorkItem.bind(stub), request, this._metadataGenerator, signal); + } + private async _streamOrchestrationHistory( req: pb.OrchestratorRequest, stub: stubs.TaskHubSidecarServiceClient, @@ -1008,15 +1018,8 @@ export class TaskHubGrpcWorker { // Incomplete history is a work-item transport failure, not an orchestration failure. // Do not replay or persist any actions; the backend can redeliver the work item. if (!signal?.aborted) { - const abandonRequest = new pb.AbandonOrchestrationTaskRequest(); - abandonRequest.setCompletiontoken(completionToken); try { - await callWithMetadata( - stub.abandonTaskOrchestratorWorkItem.bind(stub), - abandonRequest, - this._metadataGenerator, - signal, - ); + await this._abandonOrchestrationWorkItem(stub, completionToken, signal); } catch (abandonError) { WorkerLogs.completionError( this._logger, @@ -1077,13 +1080,7 @@ export class TaskHubGrpcWorker { ); try { - const abandonRequest = new pb.AbandonOrchestrationTaskRequest(); - abandonRequest.setCompletiontoken(completionToken); - await callWithMetadata( - stub.abandonTaskOrchestratorWorkItem.bind(stub), - abandonRequest, - this._metadataGenerator, - ); + await this._abandonOrchestrationWorkItem(stub, completionToken); } catch (e: unknown) { const error = e instanceof Error ? e : new Error(String(e)); WorkerLogs.completionError(this._logger, instanceId, error); diff --git a/packages/durabletask-js/test/worker-history-streaming.spec.ts b/packages/durabletask-js/test/worker-history-streaming.spec.ts index 6a544b3..c1bda3b 100644 --- a/packages/durabletask-js/test/worker-history-streaming.spec.ts +++ b/packages/durabletask-js/test/worker-history-streaming.spec.ts @@ -43,7 +43,9 @@ describe("Worker history streaming over gRPC", () => { let historyCalls: HistoryCall[]; let responses: pb.OrchestratorResponse[]; let abandonments: pb.AbandonOrchestrationTaskRequest[]; + let abandonmentMetadata: grpc.Metadata[]; let onHistory: (call: HistoryCall) => void; + let onAbandon: stubs.ITaskHubSidecarServiceServer["abandonTaskOrchestratorWorkItem"]; let historySpy: jest.SpyInstance; beforeAll(() => { @@ -62,7 +64,9 @@ describe("Worker history streaming over gRPC", () => { historyCalls = []; responses = []; abandonments = []; + abandonmentMetadata = []; onHistory = (call) => call.end(); + onAbandon = (_call, callback) => callback(null, new pb.AbandonOrchestrationTaskResponse()); historySpy = jest.spyOn(stubs.TaskHubSidecarServiceClient.prototype, "streamInstanceHistory"); server = new grpc.Server(); const service = { @@ -81,7 +85,8 @@ describe("Worker history streaming over gRPC", () => { }, abandonTaskOrchestratorWorkItem: (call, callback) => { abandonments.push(call.request); - callback(null, new pb.AbandonOrchestrationTaskResponse()); + abandonmentMetadata.push(call.metadata); + onAbandon(call, callback); }, } satisfies Pick< stubs.ITaskHubSidecarServiceServer, @@ -314,6 +319,8 @@ describe("Worker history streaming over gRPC", () => { it.each([grpc.status.UNAVAILABLE, grpc.status.CANCELLED])( "abandons incomplete history on gRPC status %s", async (code) => { + const abandon = jest.fn(worker["_abandonOrchestrationWorkItem"].bind(worker)); + worker["_abandonOrchestrationWorkItem"] = abandon; const orchestrator = jest.fn(async function shouldNotExecute() { return "incorrect"; }); @@ -336,6 +343,10 @@ describe("Worker history streaming over gRPC", () => { expect(chunksReceived).toBe(1); expect(orchestrator).not.toHaveBeenCalled(); expect(abandonments.map((item) => item.getCompletiontoken())).toEqual(["history-token"]); + expect(abandon).toHaveBeenCalledTimes(1); + expect(abandon).toHaveBeenCalledWith(worker["_stub"], "history-token", worker["_abortController"]!.signal); + expect(abandonmentMetadata[0].get("taskhub")).toEqual(["history-test"]); + expect(abandonmentMetadata[0].get("authorization")).toEqual(["test-token"]); expect(exporter.getFinishedSpans()).toHaveLength(0); expectStreamCleanedUp(); @@ -354,6 +365,8 @@ describe("Worker history streaming over gRPC", () => { it.each([VersionFailureStrategy.Reject, VersionFailureStrategy.Fail])( "checks streamed versions before dispatch (strategy=%s)", async (failureStrategy) => { + const abandon = jest.fn(worker["_abandonOrchestrationWorkItem"].bind(worker)); + worker["_abandonOrchestrationWorkItem"] = abandon; worker["_versioning"] = { version: "1", matchStrategy: VersionMatchStrategy.Strict, failureStrategy }; onHistory = (call) => { call.write( @@ -370,7 +383,12 @@ describe("Worker history streaming over gRPC", () => { expect(execute).not.toHaveBeenCalled(); if (failureStrategy === VersionFailureStrategy.Reject) { expect(abandonments[0].getCompletiontoken()).toBe("history-token"); + expect(abandon).toHaveBeenCalledTimes(1); + expect(abandon).toHaveBeenCalledWith(worker["_stub"], "history-token"); + expect(abandonmentMetadata[0].get("taskhub")).toEqual(["history-test"]); + expect(abandonmentMetadata[0].get("authorization")).toEqual(["test-token"]); } else { + expect(abandon).not.toHaveBeenCalled(); expect(responses[0].getActionsList()[0].getCompleteorchestration()!.getFailuredetails()!.getErrortype()).toBe( "VersionMismatch", ); @@ -378,6 +396,30 @@ describe("Worker history streaming over gRPC", () => { }, ); + it("cancels history-failure abandonment on stop without waiting for a late response", async () => { + let completeAbandon!: () => void; + onAbandon = (_call, callback) => { + completeAbandon = () => callback(null, new pb.AbandonOrchestrationTaskResponse()); + }; + const abandon = jest.spyOn(stubs.TaskHubSidecarServiceClient.prototype, "abandonTaskOrchestratorWorkItem"); + const execute = jest.spyOn(OrchestrationExecutor.prototype, "execute"); + await start(); + send(request()); + await waitFor(() => abandonments.length > 0); + const call = abandon.mock.results[0].value as grpc.ClientUnaryCall; + const cancel = jest.spyOn(call, "cancel"); + worker["_shutdownTimeoutMs"] = 50; + await worker.stop(); + expect(cancel).toHaveBeenCalledTimes(1); + expect(worker["_pendingWorkItems"].size).toBe(0); + expect(worker["_historyCancellations"].size).toBe(0); + expect(execute).not.toHaveBeenCalled(); + expect(responses).toHaveLength(0); + completeAbandon(); + await new Promise((resolve) => setImmediate(resolve)); + expect(responses).toHaveLength(0); + }); + it("cancels outstanding history on stop without executing or leaving pending work", async () => { let cancelled = false; onHistory = (call) => { From 1ea22361d34d52c3882a18108083e009175a6ff8 Mon Sep 17 00:00:00 2001 From: wangbill Date: Tue, 8 Sep 2026 10:45:41 -0700 Subject: [PATCH 6/9] fix(worker): cancel abandonment metadata waits on shutdown Reuse the shared gRPC cancellation prerequisite from 7219129 and cover late metadata settlement after a history transport failure. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: e0af01a5-0dfa-4e71-a660-c4186e65d7e0 Copilot-Session: 7b853e01-1e42-41c2-92a7-e2bbab87d215 --- .../src/utils/grpc-helper.util.ts | 117 ++++++++++-------- .../test/worker-history-streaming.spec.ts | 40 ++++++ 2 files changed, 102 insertions(+), 55 deletions(-) diff --git a/packages/durabletask-js/src/utils/grpc-helper.util.ts b/packages/durabletask-js/src/utils/grpc-helper.util.ts index 9af65a7..a8a8b0c 100644 --- a/packages/durabletask-js/src/utils/grpc-helper.util.ts +++ b/packages/durabletask-js/src/utils/grpc-helper.util.ts @@ -1,55 +1,62 @@ -// Copyright (c) Microsoft Corporation. All rights reserved. -// Licensed under the MIT License. - -import * as grpc from "@grpc/grpc-js"; - -/** - * Type for a function that generates gRPC metadata (e.g., for taskhub, auth tokens). - */ -export type MetadataGenerator = () => Promise; - -/** - * Promisifies a gRPC unary call with metadata support. - * - * @param method The gRPC method to call (must be bound to the stub). - * @param req The request object. - * @param metadataGenerator Optional function to generate metadata for the call. - * @param signal Optional signal that cancels the call. - * @returns A promise that resolves with the response or rejects with an error. - */ -export async function callWithMetadata( - method: ( - req: TReq, - metadata: grpc.Metadata, - callback: (error: grpc.ServiceError | null, response: TRes) => void, - ) => grpc.ClientUnaryCall, - req: TReq, - metadataGenerator?: MetadataGenerator, - signal?: AbortSignal, -): Promise { - const metadata = metadataGenerator ? await metadataGenerator() : new grpc.Metadata(); - if (signal?.aborted) { - throw signal.reason; - } - - let onAbort = () => {}; - try { - return await new Promise((resolve, reject) => { - let call: grpc.ClientUnaryCall | undefined = undefined; - onAbort = () => { - reject(signal?.reason); - call?.cancel(); - }; - signal?.addEventListener("abort", onAbort, { once: true }); - call = method(req, metadata, (error, response) => { - if (error) { - reject(error); - } else { - resolve(response); - } - }); - }); - } finally { - signal?.removeEventListener("abort", onAbort); - } -} +// Copyright (c) Microsoft Corporation. All rights reserved. +// Licensed under the MIT License. + +import * as grpc from "@grpc/grpc-js"; + +/** + * Type for a function that generates gRPC metadata (e.g., for taskhub, auth tokens). + */ +export type MetadataGenerator = () => Promise; + +/** + * Promisifies a gRPC unary call with metadata support. + * + * @param method The gRPC method to call (must be bound to the stub). + * @param req The request object. + * @param metadataGenerator Optional function to generate metadata for the call. + * @param signal Optional signal that cancels waiting for metadata and the call. + * @returns A promise that resolves with the response or rejects with an error. + */ +export async function callWithMetadata( + method: ( + req: TReq, + metadata: grpc.Metadata, + callback: (error: grpc.ServiceError | null, response: TRes) => void, + ) => grpc.ClientUnaryCall, + req: TReq, + metadataGenerator?: MetadataGenerator, + signal?: AbortSignal, +): Promise { + if (signal?.aborted) { + throw signal.reason; + } + + let onAbort = () => {}; + try { + return await new Promise((resolve, reject) => { + let call: grpc.ClientUnaryCall | undefined = undefined; + onAbort = () => { + reject(signal?.reason); + call?.cancel(); + }; + signal?.addEventListener("abort", onAbort, { once: true }); + const invoke = async () => { + const metadata = metadataGenerator ? await metadataGenerator() : new grpc.Metadata(); + // Metadata generation may finish after cancellation; never start a late RPC. + if (signal?.aborted) { + throw signal.reason; + } + call = method(req, metadata, (error, response) => { + if (error) { + reject(error); + } else { + resolve(response); + } + }); + }; + invoke().catch(reject); + }); + } finally { + signal?.removeEventListener("abort", onAbort); + } +} diff --git a/packages/durabletask-js/test/worker-history-streaming.spec.ts b/packages/durabletask-js/test/worker-history-streaming.spec.ts index c1bda3b..d0841f7 100644 --- a/packages/durabletask-js/test/worker-history-streaming.spec.ts +++ b/packages/durabletask-js/test/worker-history-streaming.spec.ts @@ -473,6 +473,46 @@ describe("Worker history streaming over gRPC", () => { }, ); + it.each([false, true])( + "stops pending abandonment metadata before late settlement (reject=%s)", + async (rejectMetadata) => { + await start(); + let releaseMetadata!: () => void; + const metadata = new Promise((resolve, reject) => { + releaseMetadata = () => { + if (rejectMetadata) reject(new Error("late abandonment metadata failure")); + else resolve(new grpc.Metadata()); + }; + }); + const getMetadata = jest.fn(() => metadata).mockResolvedValueOnce(new grpc.Metadata()); + worker["_metadataGenerator"] = getMetadata; + onHistory = (call) => { + call.emit("error", Object.assign(new Error("history unavailable"), { code: grpc.status.UNAVAILABLE })); + }; + const abandon = jest.spyOn(stubs.TaskHubSidecarServiceClient.prototype, "abandonTaskOrchestratorWorkItem"); + const execute = jest.spyOn(OrchestrationExecutor.prototype, "execute"); + try { + send(request()); + await waitFor(() => getMetadata.mock.calls.length === 2); + expect(historyCalls).toHaveLength(1); + expect(worker["_pendingWorkItems"].size).toBe(1); + worker["_shutdownTimeoutMs"] = 50; + await worker.stop(); + expect(worker["_pendingWorkItems"].size).toBe(0); + expect(worker["_historyCancellations"].size).toBe(0); + expect(abandon).not.toHaveBeenCalled(); + expect(execute).not.toHaveBeenCalled(); + expect(responses).toHaveLength(0); + } finally { + releaseMetadata(); + await new Promise((resolve) => setImmediate(resolve)); + } + expect(abandon).not.toHaveBeenCalled(); + expect(abandonments).toHaveLength(0); + expect(worker["_pendingWorkItems"].size).toBe(0); + }, + ); + it("abandons on metadata failure without using inline history", async () => { await start(); worker["_metadataGenerator"] = jest From 4458843d6e0f67e682988fcaad50a46de5b7f30c Mon Sep 17 00:00:00 2001 From: wangbill Date: Tue, 8 Sep 2026 10:57:33 -0700 Subject: [PATCH 7/9] refactor(worker): align shared abandonment helper merge anchor Move the unchanged helper before the orchestration execution JSDoc so the independently based response-retry implementation overlaps at the same insertion point. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: e0af01a5-0dfa-4e71-a660-c4186e65d7e0 Copilot-Session: 7b853e01-1e42-41c2-92a7-e2bbab87d215 --- .../src/worker/task-hub-grpc-worker.ts | 20 +++++++++---------- 1 file changed, 10 insertions(+), 10 deletions(-) diff --git a/packages/durabletask-js/src/worker/task-hub-grpc-worker.ts b/packages/durabletask-js/src/worker/task-hub-grpc-worker.ts index 8d6f1b4..325c064 100644 --- a/packages/durabletask-js/src/worker/task-hub-grpc-worker.ts +++ b/packages/durabletask-js/src/worker/task-hub-grpc-worker.ts @@ -911,6 +911,16 @@ export class TaskHubGrpcWorker { this._pendingWorkItems.add(handledPromise); } + private async _abandonOrchestrationWorkItem( + stub: stubs.TaskHubSidecarServiceClient, + completionToken: string, + signal?: AbortSignal, + ): Promise { + const request = new pb.AbandonOrchestrationTaskRequest(); + request.setCompletiontoken(completionToken); + await callWithMetadata(stub.abandonTaskOrchestratorWorkItem.bind(stub), request, this._metadataGenerator, signal); + } + /** * Executes an orchestrator request and tracks it as a pending work item. */ @@ -925,16 +935,6 @@ export class TaskHubGrpcWorker { }); } - private async _abandonOrchestrationWorkItem( - stub: stubs.TaskHubSidecarServiceClient, - completionToken: string, - signal?: AbortSignal, - ): Promise { - const request = new pb.AbandonOrchestrationTaskRequest(); - request.setCompletiontoken(completionToken); - await callWithMetadata(stub.abandonTaskOrchestratorWorkItem.bind(stub), request, this._metadataGenerator, signal); - } - private async _streamOrchestrationHistory( req: pb.OrchestratorRequest, stub: stubs.TaskHubSidecarServiceClient, From 7a41ddbf2a2d55be69b973b3cf41528438246da5 Mon Sep 17 00:00:00 2001 From: wangbill Date: Thu, 10 Sep 2026 14:42:07 -0700 Subject: [PATCH 8/9] Align worker history failures with .NET Complete non-shutdown history failures as Failed without running partial history. Remove history-specific abandonment and its shared helper, restore the upstream version-rejection flow, and keep completion cancellation scoped to streamed work items. Update loopback failure and shutdown coverage without changing the existing service-negotiation E2E scenario or client-history cases. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: e0af01a5-0dfa-4e71-a660-c4186e65d7e0 Copilot-Session: 7b853e01-1e42-41c2-92a7-e2bbab87d215 --- CHANGELOG.md | 6 +- .../src/worker/task-hub-grpc-worker.ts | 65 ++++--- .../test/worker-history-streaming.spec.ts | 182 +++++++----------- 3 files changed, 111 insertions(+), 142 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index cf47466..2bf0d0e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,9 +2,9 @@ ### New -- Add worker history streaming: hydrate service-selected history before version checks, - tracing, and replay. Cancel history streams on shutdown and abandon incomplete - work items on transport errors instead of failing the orchestration. +- 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. diff --git a/packages/durabletask-js/src/worker/task-hub-grpc-worker.ts b/packages/durabletask-js/src/worker/task-hub-grpc-worker.ts index 325c064..7341a81 100644 --- a/packages/durabletask-js/src/worker/task-hub-grpc-worker.ts +++ b/packages/durabletask-js/src/worker/task-hub-grpc-worker.ts @@ -911,16 +911,6 @@ export class TaskHubGrpcWorker { this._pendingWorkItems.add(handledPromise); } - private async _abandonOrchestrationWorkItem( - stub: stubs.TaskHubSidecarServiceClient, - completionToken: string, - signal?: AbortSignal, - ): Promise { - const request = new pb.AbandonOrchestrationTaskRequest(); - request.setCompletiontoken(completionToken); - await callWithMetadata(stub.abandonTaskOrchestratorWorkItem.bind(stub), request, this._metadataGenerator, signal); - } - /** * Executes an orchestrator request and tracks it as a pending work item. */ @@ -1002,33 +992,40 @@ export class TaskHubGrpcWorker { throw new Error(`Could not execute the orchestrator as the instanceId was not provided (${instanceId})`); } + const historySignal = req.getRequireshistorystreaming() ? this._abortController?.signal : undefined; if (req.getRequireshistorystreaming()) { - const signal = this._abortController?.signal; try { - const pastEvents = await this._streamOrchestrationHistory(req, stub, signal); - signal?.throwIfAborted(); + const pastEvents = await this._streamOrchestrationHistory(req, stub, historySignal); + historySignal?.throwIfAborted(); if ( !pastEvents.some((event) => event.hasExecutionstarted()) && !req.getNeweventsList().some((event) => event.hasExecutionstarted()) ) { - throw new Error("The provided orchestration history was incomplete (missing ExecutionStarted)."); + throw new Error("The provided orchestration history was incomplete"); } req.setPasteventsList(pastEvents); - } catch (error) { - // Incomplete history is a work-item transport failure, not an orchestration failure. - // Do not replay or persist any actions; the backend can redeliver the work item. - if (!signal?.aborted) { - try { - await this._abandonOrchestrationWorkItem(stub, completionToken, signal); - } catch (abandonError) { - WorkerLogs.completionError( - this._logger, - instanceId, - abandonError instanceof Error ? abandonError : new Error(String(abandonError)), - ); - } + } catch (e: unknown) { + if (historySignal?.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 callWithMetadata(stub.completeOrchestratorTask.bind(stub), res, this._metadataGenerator, historySignal); + } catch (e: unknown) { + const error = e instanceof Error ? e : new Error(String(e)); + WorkerLogs.completionError(this._logger, instanceId, error); } - throw error; + return; } } @@ -1064,7 +1061,7 @@ export class TaskHubGrpcWorker { res.setActionsList(actions); try { - await callWithMetadata(stub.completeOrchestratorTask.bind(stub), res, this._metadataGenerator); + await callWithMetadata(stub.completeOrchestratorTask.bind(stub), res, this._metadataGenerator, historySignal); } catch (e: unknown) { const error = e instanceof Error ? e : new Error(String(e)); WorkerLogs.completionError(this._logger, instanceId, error); @@ -1080,7 +1077,13 @@ export class TaskHubGrpcWorker { ); try { - await this._abandonOrchestrationWorkItem(stub, completionToken); + const abandonRequest = new pb.AbandonOrchestrationTaskRequest(); + abandonRequest.setCompletiontoken(completionToken); + await callWithMetadata( + stub.abandonTaskOrchestratorWorkItem.bind(stub), + abandonRequest, + this._metadataGenerator, + ); } catch (e: unknown) { const error = e instanceof Error ? e : new Error(String(e)); WorkerLogs.completionError(this._logger, instanceId, error); @@ -1187,7 +1190,7 @@ export class TaskHubGrpcWorker { } try { - await callWithMetadata(stub.completeOrchestratorTask.bind(stub), res, this._metadataGenerator); + await callWithMetadata(stub.completeOrchestratorTask.bind(stub), res, this._metadataGenerator, historySignal); } catch (e: unknown) { const error = e instanceof Error ? e : new Error(String(e)); WorkerLogs.completionError(this._logger, req.getInstanceid(), error); diff --git a/packages/durabletask-js/test/worker-history-streaming.spec.ts b/packages/durabletask-js/test/worker-history-streaming.spec.ts index d0841f7..1f53cf7 100644 --- a/packages/durabletask-js/test/worker-history-streaming.spec.ts +++ b/packages/durabletask-js/test/worker-history-streaming.spec.ts @@ -42,10 +42,10 @@ describe("Worker history streaming over gRPC", () => { let subscription: grpc.ServerWritableStream | undefined; let historyCalls: HistoryCall[]; let responses: pb.OrchestratorResponse[]; + let responseMetadata: grpc.Metadata[]; let abandonments: pb.AbandonOrchestrationTaskRequest[]; let abandonmentMetadata: grpc.Metadata[]; let onHistory: (call: HistoryCall) => void; - let onAbandon: stubs.ITaskHubSidecarServiceServer["abandonTaskOrchestratorWorkItem"]; let historySpy: jest.SpyInstance; beforeAll(() => { @@ -63,10 +63,10 @@ describe("Worker history streaming over gRPC", () => { subscription = undefined; historyCalls = []; responses = []; + responseMetadata = []; abandonments = []; abandonmentMetadata = []; onHistory = (call) => call.end(); - onAbandon = (_call, callback) => callback(null, new pb.AbandonOrchestrationTaskResponse()); historySpy = jest.spyOn(stubs.TaskHubSidecarServiceClient.prototype, "streamInstanceHistory"); server = new grpc.Server(); const service = { @@ -81,12 +81,13 @@ describe("Worker history streaming over gRPC", () => { }, completeOrchestratorTask: (call, callback) => { responses.push(call.request); + responseMetadata.push(call.metadata); callback(null, new pb.CompleteTaskResponse()); }, abandonTaskOrchestratorWorkItem: (call, callback) => { abandonments.push(call.request); abandonmentMetadata.push(call.metadata); - onAbandon(call, callback); + callback(null, new pb.AbandonOrchestrationTaskResponse()); }, } satisfies Pick< stubs.ITaskHubSidecarServiceServer, @@ -152,6 +153,20 @@ describe("Worker history streaming over gRPC", () => { } } + function expectHistoryFailure(message: string): void { + expect(abandonments).toHaveLength(0); + expect(responses).toHaveLength(1); + const response = responses[0]; + expect(response.getInstanceid()).toBe(instanceId); + expect(response.getCompletiontoken()).toBe("history-token"); + expect(response.getActionsList()).toHaveLength(1); + const completed = response.getActionsList()[0].getCompleteorchestration()!; + expect(completed.getOrchestrationstatus()).toBe(pb.OrchestrationStatus.ORCHESTRATION_STATUS_FAILED); + expect(completed.getFailuredetails()!.getErrortype()).toBe("Error"); + expect(completed.getFailuredetails()!.getErrormessage()).toContain(message); + expect(completed.getFailuredetails()!.getStacktrace()!.getValue()).toBeTruthy(); + } + it("advertises only HistoryStreaming, retaining concurrency hints", async () => { await start(); expect(subscription!.request.getCapabilitiesList()).toEqual([ @@ -275,7 +290,7 @@ describe("Worker history streaming over gRPC", () => { [VersionMatchStrategy.Strict, false], [VersionMatchStrategy.None, true], [VersionMatchStrategy.Strict, true], - ])("abandons OK history missing ExecutionStarted (strategy=%s, nonempty=%s)", async (matchStrategy, nonempty) => { + ])("fails OK history missing ExecutionStarted (strategy=%s, nonempty=%s)", async (matchStrategy, nonempty) => { worker["_versioning"] = { version: "1", matchStrategy, failureStrategy: VersionFailureStrategy.Fail }; worker.addOrchestrator(async function* resumedHistory(ctx: OrchestrationContext, input: number): AsyncGenerator { return yield ctx.callActivity("echo", input); @@ -297,30 +312,15 @@ describe("Worker history streaming over gRPC", () => { send(req); await waitFor(() => responses.length > 0 || abandonments.length > 0); await settled(); - expect(abandonments.map((item) => item.getCompletiontoken())).toEqual(["history-token"]); - expect(responses).toHaveLength(0); + expectHistoryFailure("The provided orchestration history was incomplete"); expect(execute).not.toHaveBeenCalled(); expect(exporter.getFinishedSpans()).toHaveLength(0); expectStreamCleanedUp(); - - onHistory = (call) => { - call.write(chunk(pastEvents)); - call.end(); - }; - send(req, "complete-history-redelivery"); - await waitFor(() => responses.length > 0); - await settled(); - expect(responses[0].getCompletiontoken()).toBe("complete-history-redelivery"); - const completed = responses[0].getActionsList()[0].getCompleteorchestration()!; - expect(completed.getOrchestrationstatus()).toBe(pb.OrchestrationStatus.ORCHESTRATION_STATUS_COMPLETED); - expect(completed.getResult()!.getValue()).toBe("14"); }); it.each([grpc.status.UNAVAILABLE, grpc.status.CANCELLED])( - "abandons incomplete history on gRPC status %s", + "fails incomplete history on non-shutdown gRPC status %s", async (code) => { - const abandon = jest.fn(worker["_abandonOrchestrationWorkItem"].bind(worker)); - worker["_abandonOrchestrationWorkItem"] = abandon; const orchestrator = jest.fn(async function shouldNotExecute() { return "incorrect"; }); @@ -335,38 +335,25 @@ describe("Worker history streaming over gRPC", () => { }); call.write(chunk(events)); }; + const execute = jest.spyOn(OrchestrationExecutor.prototype, "execute"); await start(); send(request().setPasteventsList(events)); await waitFor(() => abandonments.length > 0 || responses.length > 0); await settled(); - expect(responses).toHaveLength(0); + expectHistoryFailure("history transport failed"); expect(chunksReceived).toBe(1); expect(orchestrator).not.toHaveBeenCalled(); - expect(abandonments.map((item) => item.getCompletiontoken())).toEqual(["history-token"]); - expect(abandon).toHaveBeenCalledTimes(1); - expect(abandon).toHaveBeenCalledWith(worker["_stub"], "history-token", worker["_abortController"]!.signal); - expect(abandonmentMetadata[0].get("taskhub")).toEqual(["history-test"]); - expect(abandonmentMetadata[0].get("authorization")).toEqual(["test-token"]); + expect(execute).not.toHaveBeenCalled(); + expect(responseMetadata[0].get("taskhub")).toEqual(["history-test"]); + expect(responseMetadata[0].get("authorization")).toEqual(["test-token"]); expect(exporter.getFinishedSpans()).toHaveLength(0); expectStreamCleanedUp(); - - onHistory = (call) => { - call.write(chunk(events)); - call.end(); - }; - send(request(), "redelivery-token"); - await waitFor(() => responses.length > 0); - await settled(); - expect(orchestrator).toHaveBeenCalledTimes(1); - expect(responses[0].getCompletiontoken()).toBe("redelivery-token"); }, ); it.each([VersionFailureStrategy.Reject, VersionFailureStrategy.Fail])( "checks streamed versions before dispatch (strategy=%s)", async (failureStrategy) => { - const abandon = jest.fn(worker["_abandonOrchestrationWorkItem"].bind(worker)); - worker["_abandonOrchestrationWorkItem"] = abandon; worker["_versioning"] = { version: "1", matchStrategy: VersionMatchStrategy.Strict, failureStrategy }; onHistory = (call) => { call.write( @@ -383,12 +370,10 @@ describe("Worker history streaming over gRPC", () => { expect(execute).not.toHaveBeenCalled(); if (failureStrategy === VersionFailureStrategy.Reject) { expect(abandonments[0].getCompletiontoken()).toBe("history-token"); - expect(abandon).toHaveBeenCalledTimes(1); - expect(abandon).toHaveBeenCalledWith(worker["_stub"], "history-token"); expect(abandonmentMetadata[0].get("taskhub")).toEqual(["history-test"]); expect(abandonmentMetadata[0].get("authorization")).toEqual(["test-token"]); } else { - expect(abandon).not.toHaveBeenCalled(); + expect(abandonments).toHaveLength(0); expect(responses[0].getActionsList()[0].getCompleteorchestration()!.getFailuredetails()!.getErrortype()).toBe( "VersionMismatch", ); @@ -396,30 +381,6 @@ describe("Worker history streaming over gRPC", () => { }, ); - it("cancels history-failure abandonment on stop without waiting for a late response", async () => { - let completeAbandon!: () => void; - onAbandon = (_call, callback) => { - completeAbandon = () => callback(null, new pb.AbandonOrchestrationTaskResponse()); - }; - const abandon = jest.spyOn(stubs.TaskHubSidecarServiceClient.prototype, "abandonTaskOrchestratorWorkItem"); - const execute = jest.spyOn(OrchestrationExecutor.prototype, "execute"); - await start(); - send(request()); - await waitFor(() => abandonments.length > 0); - const call = abandon.mock.results[0].value as grpc.ClientUnaryCall; - const cancel = jest.spyOn(call, "cancel"); - worker["_shutdownTimeoutMs"] = 50; - await worker.stop(); - expect(cancel).toHaveBeenCalledTimes(1); - expect(worker["_pendingWorkItems"].size).toBe(0); - expect(worker["_historyCancellations"].size).toBe(0); - expect(execute).not.toHaveBeenCalled(); - expect(responses).toHaveLength(0); - completeAbandon(); - await new Promise((resolve) => setImmediate(resolve)); - expect(responses).toHaveLength(0); - }); - it("cancels outstanding history on stop without executing or leaving pending work", async () => { let cancelled = false; onHistory = (call) => { @@ -470,65 +431,70 @@ describe("Worker history streaming over gRPC", () => { await new Promise((resolve) => setImmediate(resolve)); expect(historyCalls).toHaveLength(0); expect(abandonments).toHaveLength(0); + expect(responses).toHaveLength(0); }, ); - it.each([false, true])( - "stops pending abandonment metadata before late settlement (reject=%s)", - async (rejectMetadata) => { - await start(); - let releaseMetadata!: () => void; - const metadata = new Promise((resolve, reject) => { - releaseMetadata = () => { - if (rejectMetadata) reject(new Error("late abandonment metadata failure")); - else resolve(new grpc.Metadata()); - }; + it("fails on history metadata failure without using inline history", async () => { + await start(); + worker["_metadataGenerator"] = jest + .fn() + .mockRejectedValueOnce(new Error("token refresh failed")) + .mockResolvedValue(new grpc.Metadata()); + const execute = jest.spyOn(OrchestrationExecutor.prototype, "execute"); + send(request()); + await waitFor(() => responses.length > 0 || abandonments.length > 0); + await settled(); + expect(historyCalls).toHaveLength(0); + expect(execute).not.toHaveBeenCalled(); + expectHistoryFailure("token refresh failed"); + }); + + it.each(["history failure", "orchestration completion", "version failure"])( + "does not submit %s after stop during response metadata", + async (outcome) => { + worker.addOrchestrator(async function shutdownHistory() { + return "complete"; }); - const getMetadata = jest.fn(() => metadata).mockResolvedValueOnce(new grpc.Metadata()); - worker["_metadataGenerator"] = getMetadata; onHistory = (call) => { - call.emit("error", Object.assign(new Error("history unavailable"), { code: grpc.status.UNAVAILABLE })); + call.write( + chunk([pbh.newExecutionStartedEvent("shutdownHistory", instanceId, undefined, undefined, executionId, "2")]), + ); + call.end(); }; - const abandon = jest.spyOn(stubs.TaskHubSidecarServiceClient.prototype, "abandonTaskOrchestratorWorkItem"); - const execute = jest.spyOn(OrchestrationExecutor.prototype, "execute"); + if (outcome === "version failure") { + worker["_versioning"] = { + version: "1", + matchStrategy: VersionMatchStrategy.Strict, + failureStrategy: VersionFailureStrategy.Fail, + }; + } + await start(); + let releaseMetadata!: (metadata: grpc.Metadata) => void; + const metadata = new Promise((resolve) => { + releaseMetadata = resolve; + }); + const getMetadata = jest.fn(() => metadata); + if (outcome === "history failure") getMetadata.mockRejectedValueOnce(new Error("token refresh failed")); + else getMetadata.mockResolvedValueOnce(new grpc.Metadata()); + worker["_metadataGenerator"] = getMetadata; + const complete = jest.spyOn(stubs.TaskHubSidecarServiceClient.prototype, "completeOrchestratorTask"); try { send(request()); await waitFor(() => getMetadata.mock.calls.length === 2); - expect(historyCalls).toHaveLength(1); - expect(worker["_pendingWorkItems"].size).toBe(1); worker["_shutdownTimeoutMs"] = 50; await worker.stop(); - expect(worker["_pendingWorkItems"].size).toBe(0); - expect(worker["_historyCancellations"].size).toBe(0); - expect(abandon).not.toHaveBeenCalled(); - expect(execute).not.toHaveBeenCalled(); - expect(responses).toHaveLength(0); } finally { - releaseMetadata(); - await new Promise((resolve) => setImmediate(resolve)); + releaseMetadata(new grpc.Metadata()); } - expect(abandon).not.toHaveBeenCalled(); + await settled(); + expect(complete).not.toHaveBeenCalled(); + expect(responses).toHaveLength(0); expect(abandonments).toHaveLength(0); - expect(worker["_pendingWorkItems"].size).toBe(0); + expect(historyCalls).toHaveLength(outcome === "history failure" ? 0 : 1); }, ); - it("abandons on metadata failure without using inline history", async () => { - await start(); - worker["_metadataGenerator"] = jest - .fn() - .mockRejectedValueOnce(new Error("token refresh failed")) - .mockResolvedValue(new grpc.Metadata()); - const execute = jest.spyOn(OrchestrationExecutor.prototype, "execute"); - send(request()); - await waitFor(() => abandonments.length > 0); - await settled(); - expect(historyCalls).toHaveLength(0); - expect(execute).not.toHaveBeenCalled(); - expect(responses).toHaveLength(0); - expect(abandonments[0].getCompletiontoken()).toBe("history-token"); - }); - it("uses the work item's captured stub when the worker channel is replaced", async () => { worker.addOrchestrator(async function capturedStub() { return "original-channel"; From baf52c26f45762e34ec3036743ad88506767dd9f Mon Sep 17 00:00:00 2001 From: wangbill Date: Fri, 11 Sep 2026 13:02:31 -0700 Subject: [PATCH 9/9] fix(worker): cancel every response send on stop like .NET Use the dispatch-captured run signal for initial response RPCs, retries, and backoff. Remove the first-send exemption and streamed-versus-inline signal policy. Preserve history hydration and retry behavior. Update all response-type cancellation coverage, real gRPC initial-call cases, and shutdown documentation. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: e0af01a5-0dfa-4e71-a660-c4186e65d7e0 Copilot-Session: 7b853e01-1e42-41c2-92a7-e2bbab87d215 --- CHANGELOG.md | 2 + README.md | 9 +- .../src/worker/task-hub-grpc-worker.ts | 26 +- .../worker-response-delivery-grpc.spec.ts | 46 ++++ .../test/worker-response-delivery.spec.ts | 226 ++++++++++-------- 5 files changed, 194 insertions(+), 115 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index a143d21..f0f2144 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -19,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. diff --git a/README.md b/README.md index e7c04f8..7fde1f9 100644 --- a/README.md +++ b/README.md @@ -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 diff --git a/packages/durabletask-js/src/worker/task-hub-grpc-worker.ts b/packages/durabletask-js/src/worker/task-hub-grpc-worker.ts index 853de92..126a706 100644 --- a/packages/durabletask-js/src/worker/task-hub-grpc-worker.ts +++ b/packages/durabletask-js/src/worker/task-hub-grpc-worker.ts @@ -914,8 +914,7 @@ export class TaskHubGrpcWorker { private async _deliverResponse( method: Parameters>[0], request: TReq, - retrySignal?: AbortSignal, - initialSignal?: AbortSignal, + signal?: AbortSignal, ): Promise { const backoff = new ExponentialBackoff({ initialDelayMs: 200, @@ -926,13 +925,7 @@ export class TaskHubGrpcWorker { }); for (;;) { try { - // Initial responses drain during shutdown unless streamed history requires immediate cancellation. - return await callWithMetadata( - method, - request, - this._metadataGenerator, - backoff.attemptCount === 0 ? initialSignal : retrySignal, - ); + return await callWithMetadata(method, request, this._metadataGenerator, signal); } catch (error) { const status = error instanceof Error ? this._getGrpcStatus(error) : undefined; if ( @@ -944,7 +937,7 @@ export class TaskHubGrpcWorker { ) { throw error; } - await backoff.wait(retrySignal); + await backoff.wait(signal); } } } @@ -1032,11 +1025,10 @@ export class TaskHubGrpcWorker { throw new Error(`Could not execute the orchestrator as the instanceId was not provided (${instanceId})`); } - const historySignal = req.getRequireshistorystreaming() ? retrySignal : undefined; if (req.getRequireshistorystreaming()) { try { - const pastEvents = await this._streamOrchestrationHistory(req, stub, historySignal); - historySignal?.throwIfAborted(); + const pastEvents = await this._streamOrchestrationHistory(req, stub, retrySignal); + retrySignal?.throwIfAborted(); if ( !pastEvents.some((event) => event.hasExecutionstarted()) && !req.getNeweventsList().some((event) => event.hasExecutionstarted()) @@ -1045,7 +1037,7 @@ export class TaskHubGrpcWorker { } req.setPasteventsList(pastEvents); } catch (e: unknown) { - if (historySignal?.aborted) return; + 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(); @@ -1060,7 +1052,7 @@ export class TaskHubGrpcWorker { ), ]); try { - await this._deliverResponse(stub.completeOrchestratorTask.bind(stub), res, retrySignal, historySignal); + 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); @@ -1101,7 +1093,7 @@ export class TaskHubGrpcWorker { res.setActionsList(actions); try { - await this._deliverResponse(stub.completeOrchestratorTask.bind(stub), res, retrySignal, historySignal); + 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); @@ -1226,7 +1218,7 @@ export class TaskHubGrpcWorker { } try { - await this._deliverResponse(stub.completeOrchestratorTask.bind(stub), res, retrySignal, historySignal); + 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, req.getInstanceid(), error); diff --git a/packages/durabletask-js/test/worker-response-delivery-grpc.spec.ts b/packages/durabletask-js/test/worker-response-delivery-grpc.spec.ts index c22732c..fe4a644 100644 --- a/packages/durabletask-js/test/worker-response-delivery-grpc.spec.ts +++ b/packages/durabletask-js/test/worker-response-delivery-grpc.spec.ts @@ -117,6 +117,52 @@ describe("Worker response retries over gRPC", () => { expect(logger.error).not.toHaveBeenCalled(); }); + it.each(["metadata", "RPC"] as const)("stop cancels the first response during %s", async (phase) => { + const received = deferred(); + const cancelled = deferred(); + const metadataStarted = deferred(); + const metadata = deferred(); + const requests: pb.ActivityResponse[] = []; + let finishResponse = () => {}; + const stream = await start((call, callback) => { + requests.push(call.request); + received.resolve(); + if (phase === "RPC") { + finishResponse = () => callback(null, new pb.CompleteTaskResponse()); + call.once("cancelled", () => cancelled.resolve()); + } else callback(null, new pb.CompleteTaskResponse()); + }); + if (phase === "metadata") { + worker["_metadataGenerator"] = () => { + metadataStarted.resolve(); + return metadata.promise; + }; + } + const send = jest.spyOn(worker["_stub"]!, "completeActivityTask"); + let cancel: jest.SpyInstance | undefined; + try { + stream.write(workItem()); + await withTimeout(phase === "metadata" ? metadataStarted.promise : received.promise, 5000); + if (phase === "RPC") { + const call = send.mock.results[0].value as grpc.ClientUnaryCall; + cancel = jest.spyOn(call, "cancel"); + } + await worker.stop(); + if (phase === "RPC") expect(cancel).toHaveBeenCalledTimes(1); + } finally { + metadata.resolve(new grpc.Metadata()); + finishResponse(); + } + await withTimeout(Promise.all(worker["_pendingWorkItems"]), 5000); + if (phase === "RPC") { + await withTimeout(cancelled.promise, 5000); + } + expect(send).toHaveBeenCalledTimes(phase === "RPC" ? 1 : 0); + expect(requests).toHaveLength(phase === "RPC" ? 1 : 0); + expect(activity).toHaveBeenCalledTimes(1); + expect(worker["_pendingWorkItems"].size).toBe(0); + }); + it("bounds SDK sends to ten without overriding configured transport retries", async () => { const wait = ExponentialBackoff.prototype.wait; jest.spyOn(ExponentialBackoff.prototype, "wait").mockImplementation(function ( diff --git a/packages/durabletask-js/test/worker-response-delivery.spec.ts b/packages/durabletask-js/test/worker-response-delivery.spec.ts index 57856ff..568377e 100644 --- a/packages/durabletask-js/test/worker-response-delivery.spec.ts +++ b/packages/durabletask-js/test/worker-response-delivery.spec.ts @@ -154,7 +154,7 @@ describe("Worker response retries", () => { }, ); - it("lets running activity work send its first response during graceful shutdown", async () => { + it("does not send a running activity's first response after shutdown starts", async () => { let finish!: () => void; worker.addNamedActivity("answer", () => new Promise((resolve) => (finish = resolve))); const send = jest @@ -177,113 +177,149 @@ describe("Worker response retries", () => { finish(); await jest.advanceTimersByTimeAsync(1000); await stopping; - expect(send).toHaveBeenCalledTimes(1); + expect(send).not.toHaveBeenCalled(); expect(close).toHaveBeenCalledTimes(1); - expect(logger.error).not.toHaveBeenCalled(); + expect(logger.error).toHaveBeenCalledTimes(1); expect(logger.warn).not.toHaveBeenCalled(); }); - it.each(["activity", "orchestrator", "entity-v1", "entity-v2", "version-fail", "version-abandon"] as const)( - "retries the same %s response without rerunning user code and keeps its dispatch-time run signal", - async (kind) => { - const mismatch = kind.startsWith("version"); - const executed = jest.fn(); - if (mismatch) - worker = new TaskHubGrpcWorker({ - logger, - versioning: { - version: "2", - matchStrategy: VersionMatchStrategy.Strict, - failureStrategy: kind === "version-fail" ? VersionFailureStrategy.Fail : VersionFailureStrategy.Reject, - }, - }); - worker["_abortController"] = controller; - worker.addNamedActivity("answer", () => { - executed(); - return 42; + it.each( + (["activity", "orchestrator", "entity-v1", "entity-v2", "version-fail", "version-abandon"] as const).flatMap( + (kind) => + (["retry", "already aborted", "initial metadata", "initial RPC"] as const).map((phase) => ({ kind, phase })), + ), + )("$kind response uses its dispatch-time run signal during $phase", async ({ kind, phase }) => { + const mismatch = kind.startsWith("version"); + const executed = jest.fn(); + if (mismatch) + worker = new TaskHubGrpcWorker({ + logger, + shutdownTimeoutMs: 100, + versioning: { + version: "2", + matchStrategy: VersionMatchStrategy.Strict, + failureStrategy: kind === "version-fail" ? VersionFailureStrategy.Fail : VersionFailureStrategy.Reject, + }, }); - worker.addNamedOrchestrator("answer", async () => { + worker["_abortController"] = controller; + worker.addNamedActivity("answer", () => { + executed(); + return 42; + }); + worker.addNamedOrchestrator("answer", async () => { + executed(); + return 42; + }); + class Counter extends TaskEntity { + increment() { executed(); - return 42; - }); - class Counter extends TaskEntity { - increment() { - executed(); - return ++this.state; - } - protected initializeState() { - return 0; - } + return ++this.state; } - worker.addNamedEntity("counter", () => new Counter()); - const requests: Array< - pb.ActivityResponse | pb.OrchestratorResponse | pb.EntityBatchResult | pb.AbandonOrchestrationTaskRequest - > = []; - function complete(result: T) { - return ( - request: (typeof requests)[number], - _metadata: grpc.Metadata, - optionsOrCallback: Partial | Callback, - callback?: Callback, - ) => { - requests.push(request); - const respond = typeof optionsOrCallback === "function" ? optionsOrCallback : callback!; - respond(requests.length === 1 ? grpcError(grpc.status.INTERNAL) : null, result); - return unaryCall(); - }; + protected initializeState() { + return 0; } - jest.spyOn(stub, "completeActivityTask").mockImplementation(complete(new pb.CompleteTaskResponse())); - jest.spyOn(stub, "completeOrchestratorTask").mockImplementation(complete(new pb.CompleteTaskResponse())); - jest.spyOn(stub, "completeEntityTask").mockImplementation(complete(new pb.CompleteTaskResponse())); - jest - .spyOn(stub, "abandonTaskOrchestratorWorkItem") - .mockImplementation(complete(new pb.AbandonOrchestrationTaskResponse())); - const item = new pb.WorkItem().setCompletiontoken("token"); - if (kind === "activity") item.setActivityrequest(activityRequest()); - else if (kind === "orchestrator" || mismatch) - item.setOrchestratorrequest( - new pb.OrchestratorRequest() - .setInstanceid("instance") - .setNeweventsList([ - new pb.HistoryEvent() - .setTimestamp(Timestamp.fromDate(new Date())) - .setOrchestratorstarted(new pb.OrchestratorStartedEvent()), - new pb.HistoryEvent().setExecutionstarted( - new pb.ExecutionStartedEvent().setName("answer").setVersion(new StringValue().setValue("1")), - ), - ]), - ); - else if (kind === "entity-v1") - item.setEntityrequest( - new pb.EntityBatchRequest() - .setInstanceid("@counter@key") - .setOperationsList([new pb.OperationRequest().setOperation("increment").setRequestid("req")]), - ); - else - item.setEntityrequestv2( - new pb.EntityRequest() - .setInstanceid("@counter@key") - .setOperationrequestsList([ - new pb.HistoryEvent().setEntityoperationsignaled( - new pb.EntityOperationSignaledEvent().setOperation("increment").setRequestid("req"), - ), - ]), - ); - worker["_dispatchWorkItem"](item, stub); - const wait = jest.spyOn(controller.signal, "addEventListener"); - worker["_abortController"] = new AbortController(); - await jest.runAllTimersAsync(); - await Promise.all(worker["_pendingWorkItems"]); + } + worker.addNamedEntity("counter", () => new Counter()); + let finishMetadata!: (metadata: grpc.Metadata) => void; + const metadata = new Promise((resolve) => (finishMetadata = resolve)); + const generateMetadata = jest.fn(async () => new grpc.Metadata()); + if (phase === "initial metadata") generateMetadata.mockImplementation(() => metadata); + worker["_metadataGenerator"] = generateMetadata; + const cancel = jest.fn(); + let finishResponse = () => {}; + const requests: Array< + pb.ActivityResponse | pb.OrchestratorResponse | pb.EntityBatchResult | pb.AbandonOrchestrationTaskRequest + > = []; + function complete(result: T) { + return ( + request: (typeof requests)[number], + _metadata: grpc.Metadata, + optionsOrCallback: Partial | Callback, + callback?: Callback, + ) => { + requests.push(request); + const respond = typeof optionsOrCallback === "function" ? optionsOrCallback : callback!; + if (phase === "initial RPC") { + finishResponse = () => respond(null, result); + } else { + respond(phase === "retry" && requests.length === 1 ? grpcError(grpc.status.INTERNAL) : null, result); + } + return unaryCall(cancel); + }; + } + jest.spyOn(stub, "completeActivityTask").mockImplementation(complete(new pb.CompleteTaskResponse())); + jest.spyOn(stub, "completeOrchestratorTask").mockImplementation(complete(new pb.CompleteTaskResponse())); + jest.spyOn(stub, "completeEntityTask").mockImplementation(complete(new pb.CompleteTaskResponse())); + jest + .spyOn(stub, "abandonTaskOrchestratorWorkItem") + .mockImplementation(complete(new pb.AbandonOrchestrationTaskResponse())); + const item = new pb.WorkItem().setCompletiontoken("token"); + if (kind === "activity") item.setActivityrequest(activityRequest()); + else if (kind === "orchestrator" || mismatch) + item.setOrchestratorrequest( + new pb.OrchestratorRequest() + .setInstanceid("instance") + .setNeweventsList([ + new pb.HistoryEvent() + .setTimestamp(Timestamp.fromDate(new Date())) + .setOrchestratorstarted(new pb.OrchestratorStartedEvent()), + new pb.HistoryEvent().setExecutionstarted( + new pb.ExecutionStartedEvent().setName("answer").setVersion(new StringValue().setValue("1")), + ), + ]), + ); + else if (kind === "entity-v1") + item.setEntityrequest( + new pb.EntityBatchRequest() + .setInstanceid("@counter@key") + .setOperationsList([new pb.OperationRequest().setOperation("increment").setRequestid("req")]), + ); + else + item.setEntityrequestv2( + new pb.EntityRequest() + .setInstanceid("@counter@key") + .setOperationrequestsList([ + new pb.HistoryEvent().setEntityoperationsignaled( + new pb.EntityOperationSignaledEvent().setOperation("increment").setRequestid("req"), + ), + ]), + ); + const wait = jest.spyOn(controller.signal, "addEventListener"); + if (phase === "already aborted") controller.abort(); + worker["_dispatchWorkItem"](item, stub); + if (phase === "initial metadata" || phase === "initial RPC") { + await jest.advanceTimersByTimeAsync(0); + expect(generateMetadata).toHaveBeenCalledTimes(1); + expect(requests).toHaveLength(phase === "initial RPC" ? 1 : 0); + worker["_stub"] = stub; + worker["_isRunning"] = true; + const stopping = worker.stop(); + await jest.advanceTimersByTimeAsync(1100); + await stopping; + } + worker["_abortController"] = new AbortController(); + finishMetadata(new grpc.Metadata()); + finishResponse(); + await jest.runAllTimersAsync(); + await Promise.all(worker["_pendingWorkItems"]); + expect(executed).toHaveBeenCalledTimes(mismatch ? 0 : 1); + if (phase === "retry") { expect(requests).toHaveLength(2); expect(requests[1]).toBe(requests[0]); expect(requests[1].getCompletiontoken()).toBe("token"); - expect(executed).toHaveBeenCalledTimes(mismatch ? 0 : 1); expect(wait).toHaveBeenCalledWith("abort", expect.any(Function), { once: true }); expect(logger.error).not.toHaveBeenCalled(); - }, - ); + } else { + expect(requests).toHaveLength(phase === "initial RPC" ? 1 : 0); + expect(cancel).toHaveBeenCalledTimes(phase === "initial RPC" ? 1 : 0); + expect(controller.signal.aborted).toBe(true); + expect(worker["_abortController"]!.signal.aborted).toBe(false); + } + expect(worker["_pendingWorkItems"].size).toBe(0); + expect(jest.getTimerCount()).toBe(0); + }); - it("does not revive an old run's retries when its activity returns after restart", async () => { + it("does not send an old run's first response when its activity returns after restart", async () => { let finish!: () => void; worker.addNamedActivity("answer", () => new Promise((resolve) => (finish = resolve))); const send = jest @@ -309,7 +345,7 @@ describe("Worker response retries", () => { finish(); await jest.runAllTimersAsync(); await Promise.all(worker["_pendingWorkItems"]); - expect(send).toHaveBeenCalledTimes(1); + expect(send).not.toHaveBeenCalled(); expect(controller.signal.aborted).toBe(true); expect(worker["_abortController"]!.signal.aborted).toBe(false); } finally {