diff --git a/CHANGELOG.md b/CHANGELOG.md index 86db689..57c949d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,9 @@ ### New +- Add an optional `reason` to client `suspendOrchestration()` and `resumeOrchestration()`, + forwarded unchanged to the service and recorded in in-memory test history. Empty strings + are preserved; omitted reasons remain absent. - Expose readonly `ActivityContext.name` and `ActivityContext.version` from the activity request, preserving the requested version even when an unversioned implementation handles it. Existing two-argument context construction remains supported, with empty name/version defaults. diff --git a/packages/azure-functions-durable/CHANGELOG.md b/packages/azure-functions-durable/CHANGELOG.md index b1543fd..09018ef 100644 --- a/packages/azure-functions-durable/CHANGELOG.md +++ b/packages/azure-functions-durable/CHANGELOG.md @@ -14,6 +14,8 @@ ### Fixes +- Forward the optional `reason` from the classic `suspend()` / `resume()` aliases to the + core client instead of ignoring it, preserving literal strings including empty strings. - Preserve nested task `innerFailure` details in `durable-functions/testing` results. - Inherit core version-aware replay dispatch and worker child defaults in the embedded `DurableFunctionsWorker`, including classic-context wrappers. Host `app.*` function registrations diff --git a/packages/azure-functions-durable/src/client.ts b/packages/azure-functions-durable/src/client.ts index 37cdb53..13c9763 100644 --- a/packages/azure-functions-durable/src/client.ts +++ b/packages/azure-functions-durable/src/client.ts @@ -434,12 +434,11 @@ export class DurableFunctionsClient extends TaskHubGrpcClient { * * @deprecated Use {@link suspendOrchestration} instead. * @param instanceId - The orchestration instance to suspend. - * @param _reason - Accepted for classic v3 signature compatibility; ignored (the core engine does - * not record a suspend reason). + * @param reason - Optional reason sent unchanged to the service, including an empty string. */ - async suspend(instanceId: string, _reason?: string): Promise { + async suspend(instanceId: string, reason?: string): Promise { try { - await this.suspendOrchestration(instanceId); + await this.suspendOrchestration(instanceId, reason); } catch (error) { await this._mapControlPlaneError(error, instanceId, "suspend"); } @@ -450,11 +449,11 @@ export class DurableFunctionsClient extends TaskHubGrpcClient { * * @deprecated Use {@link resumeOrchestration} instead. * @param instanceId - The orchestration instance to resume. - * @param _reason - Accepted for classic v3 signature compatibility; ignored. + * @param reason - Optional reason sent unchanged to the service, including an empty string. */ - async resume(instanceId: string, _reason?: string): Promise { + async resume(instanceId: string, reason?: string): Promise { try { - await this.resumeOrchestration(instanceId); + await this.resumeOrchestration(instanceId, reason); } catch (error) { await this._mapControlPlaneError(error, instanceId, "resume"); } diff --git a/packages/azure-functions-durable/test/unit/client.spec.ts b/packages/azure-functions-durable/test/unit/client.spec.ts index 7bebfed..5e2f52b 100644 --- a/packages/azure-functions-durable/test/unit/client.spec.ts +++ b/packages/azure-functions-durable/test/unit/client.spec.ts @@ -87,12 +87,17 @@ describe("DurableFunctionsClient", () => { expect(terminate).toHaveBeenCalledWith("id-1", "cancelled"); const suspend = jest.spyOn(client, "suspendOrchestration").mockResolvedValue(undefined); - await client.suspend("id-1", "ignored-reason"); - expect(suspend).toHaveBeenCalledWith("id-1"); - const resume = jest.spyOn(client, "resumeOrchestration").mockResolvedValue(undefined); - await client.resume("id-1", "ignored-reason"); - expect(resume).toHaveBeenCalledWith("id-1"); + for (const reason of [' "maintenance"\n\u6682\u505c ', "", undefined, null]) { + await Reflect.apply(client.suspend, client, ["id-1", reason]); + expect(suspend).toHaveBeenLastCalledWith("id-1", reason); + await Reflect.apply(client.resume, client, ["id-1", reason]); + expect(resume).toHaveBeenLastCalledWith("id-1", reason); + } + await client.suspend("id-1"); + expect(suspend).toHaveBeenLastCalledWith("id-1", undefined); + await client.resume("id-1"); + expect(resume).toHaveBeenLastCalledWith("id-1", undefined); const rewind = jest.spyOn(client, "rewindInstance").mockResolvedValue(undefined); await client.rewind("id-1", "retrying"); diff --git a/packages/durabletask-js/src/client/client.ts b/packages/durabletask-js/src/client/client.ts index 38567b0..6593619 100644 --- a/packages/durabletask-js/src/client/client.ts +++ b/packages/durabletask-js/src/client/client.ts @@ -572,13 +572,22 @@ export class TaskHubGrpcClient { ); } - async suspendOrchestration(instanceId: string): Promise { + /** + * Suspends an orchestration instance. + * + * @param instanceId - The orchestration instance to suspend. + * @param reason - Optional reason sent unchanged to the service, including an empty string. + */ + async suspendOrchestration(instanceId: string, reason?: string): Promise { if (!instanceId) { throw new Error("instanceId is required"); } const req = new pb.SuspendRequest(); req.setInstanceid(instanceId); + if (reason != null) { + req.setReason(new StringValue().setValue(reason)); + } ClientLogs.suspendingInstance(this._logger, instanceId); @@ -589,13 +598,22 @@ export class TaskHubGrpcClient { ); } - async resumeOrchestration(instanceId: string): Promise { + /** + * Resumes a suspended orchestration instance. + * + * @param instanceId - The orchestration instance to resume. + * @param reason - Optional reason sent unchanged to the service, including an empty string. + */ + async resumeOrchestration(instanceId: string, reason?: string): Promise { if (!instanceId) { throw new Error("instanceId is required"); } const req = new pb.ResumeRequest(); req.setInstanceid(instanceId); + if (reason != null) { + req.setReason(new StringValue().setValue(reason)); + } ClientLogs.resumingInstance(this._logger, instanceId); diff --git a/packages/durabletask-js/src/testing/in-memory-backend.ts b/packages/durabletask-js/src/testing/in-memory-backend.ts index 0859cb2..d833231 100644 --- a/packages/durabletask-js/src/testing/in-memory-backend.ts +++ b/packages/durabletask-js/src/testing/in-memory-backend.ts @@ -314,7 +314,7 @@ export class InMemoryOrchestrationBackend { /** * Suspends an orchestration instance. */ - suspend(instanceId: string): void { + suspend(instanceId: string, reason?: string): void { const instance = this.instances.get(instanceId); if (!instance) { throw new Error(`Orchestration instance '${instanceId}' not found`); @@ -332,7 +332,7 @@ export class InMemoryOrchestrationBackend { // suspend RPC transitions the orchestration to SUSPENDED right away. instance.status = pb.OrchestrationStatus.ORCHESTRATION_STATUS_SUSPENDED; - const event = pbh.newSuspendEvent(); + const event = pbh.newSuspendEvent(reason); instance.pendingEvents.push(event); instance.lastUpdatedAt = new Date(); @@ -346,7 +346,7 @@ export class InMemoryOrchestrationBackend { /** * Resumes a suspended orchestration instance. */ - resume(instanceId: string): void { + resume(instanceId: string, reason?: string): void { const instance = this.instances.get(instanceId); if (!instance) { throw new Error(`Orchestration instance '${instanceId}' not found`); @@ -364,7 +364,7 @@ export class InMemoryOrchestrationBackend { // Transition from SUSPENDED back to RUNNING to match real sidecar behavior. instance.status = pb.OrchestrationStatus.ORCHESTRATION_STATUS_RUNNING; - const event = pbh.newResumeEvent(); + const event = pbh.newResumeEvent(reason); instance.pendingEvents.push(event); instance.lastUpdatedAt = new Date(); diff --git a/packages/durabletask-js/src/testing/test-client.ts b/packages/durabletask-js/src/testing/test-client.ts index 14a9383..712072f 100644 --- a/packages/durabletask-js/src/testing/test-client.ts +++ b/packages/durabletask-js/src/testing/test-client.ts @@ -145,16 +145,22 @@ export class TestOrchestrationClient { /** * Suspends an orchestration. + * + * @param instanceId - The orchestration instance to suspend. + * @param reason - Optional reason recorded unchanged in the suspension history event. */ - async suspendOrchestration(instanceId: string): Promise { - this.backend.suspend(instanceId); + async suspendOrchestration(instanceId: string, reason?: string): Promise { + this.backend.suspend(instanceId, reason); } /** * Resumes a suspended orchestration. + * + * @param instanceId - The orchestration instance to resume. + * @param reason - Optional reason recorded unchanged in the resumption history event. */ - async resumeOrchestration(instanceId: string): Promise { - this.backend.resume(instanceId); + async resumeOrchestration(instanceId: string, reason?: string): Promise { + this.backend.resume(instanceId, reason); } /** diff --git a/packages/durabletask-js/src/utils/pb-helper.util.ts b/packages/durabletask-js/src/utils/pb-helper.util.ts index 34721eb..8993601 100644 --- a/packages/durabletask-js/src/utils/pb-helper.util.ts +++ b/packages/durabletask-js/src/utils/pb-helper.util.ts @@ -316,24 +316,30 @@ export function newEventSentEvent(eventId: number, instanceId: string, name: str return event; } -export function newSuspendEvent(): pb.HistoryEvent { +export function newSuspendEvent(reason?: string): pb.HistoryEvent { + const executionSuspendedEvent = new pb.ExecutionSuspendedEvent(); + executionSuspendedEvent.setInput(getStringValueIfDefined(reason ?? undefined)); + const ts = new Timestamp(); const event = new pb.HistoryEvent(); event.setEventid(-1); event.setTimestamp(ts); - event.setExecutionsuspended(new pb.ExecutionSuspendedEvent()); + event.setExecutionsuspended(executionSuspendedEvent); return event; } -export function newResumeEvent(): pb.HistoryEvent { +export function newResumeEvent(reason?: string): pb.HistoryEvent { + const executionResumedEvent = new pb.ExecutionResumedEvent(); + executionResumedEvent.setInput(getStringValueIfDefined(reason ?? undefined)); + const ts = new Timestamp(); const event = new pb.HistoryEvent(); event.setEventid(-1); event.setTimestamp(ts); - event.setExecutionresumed(new pb.ExecutionResumedEvent()); + event.setExecutionresumed(executionResumedEvent); return event; } diff --git a/packages/durabletask-js/test/client-suspend-resume.spec.ts b/packages/durabletask-js/test/client-suspend-resume.spec.ts new file mode 100644 index 0000000..b02051e --- /dev/null +++ b/packages/durabletask-js/test/client-suspend-resume.spec.ts @@ -0,0 +1,96 @@ +// Copyright (c) Microsoft Corporation. All rights reserved. +// Licensed under the MIT License. + +import * as grpc from "@grpc/grpc-js"; +import { TaskHubGrpcClient } from "../src/client/client"; +import { NoOpLogger } from "../src/types/logger.type"; +import * as pb from "../src/proto/orchestrator_service_pb"; +import { TaskHubSidecarServiceService } from "../src/proto/orchestrator_service_grpc_pb"; + +describe("Client suspend/resume reasons over real gRPC", () => { + const server = new grpc.Server(); + let client: TaskHubGrpcClient; + let requests: (pb.SuspendRequest | pb.ResumeRequest)[]; + let error: grpc.ServerErrorResponse | null; + + beforeAll(async () => { + server.addService(TaskHubSidecarServiceService, { + suspendInstance: ( + call: grpc.ServerUnaryCall, + callback: grpc.sendUnaryData, + ) => { + requests.push(call.request); + callback(error, new pb.SuspendResponse()); + }, + resumeInstance: ( + call: grpc.ServerUnaryCall, + callback: grpc.sendUnaryData, + ) => { + requests.push(call.request); + callback(error, new pb.ResumeResponse()); + }, + }); + 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); + }); + }); + client = new TaskHubGrpcClient({ + hostAddress: `127.0.0.1:${port}`, + logger: new NoOpLogger(), + }); + }); + + beforeEach(() => { + requests = []; + error = null; + }); + + afterAll(async () => { + await client.stop(); + server.forceShutdown(); + }); + + describe.each(["suspendOrchestration", "resumeOrchestration"] as const)("%s", (method) => { + it.each([' "maintenance"\n\u6682\u505c ', "", undefined, null])( + "preserves the value and presence of reason %p", + async (reason) => { + await Reflect.apply(client[method], client, ["instance-1", reason]); + + expect(requests).toHaveLength(1); + expect(requests[0].getInstanceid()).toBe("instance-1"); + expect(requests[0].hasReason()).toBe(reason != null); + expect(requests[0].getReason()?.getValue()).toBe(reason ?? undefined); + }, + ); + + it("keeps one-argument calls valid and omits the reason", async () => { + await client[method]("instance-1"); + + expect(requests).toHaveLength(1); + expect(requests[0].getInstanceid()).toBe("instance-1"); + expect(requests[0].hasReason()).toBe(false); + }); + + it.each(["", undefined, null])("rejects invalid instanceId %p before sending an RPC", async (instanceId) => { + await expect(Reflect.apply(client[method], client, [instanceId, "maintenance"])).rejects.toThrow( + "instanceId is required", + ); + expect(requests).toHaveLength(0); + }); + + it("preserves service failures", async () => { + error = Object.assign(new Error("Invalid instance state"), { + code: grpc.status.FAILED_PRECONDITION, + details: "Invalid instance state", + }); + + await expect(Reflect.apply(client[method], client, ["instance-1", "maintenance"])).rejects.toMatchObject({ + code: grpc.status.FAILED_PRECONDITION, + details: "Invalid instance state", + }); + expect(requests).toHaveLength(1); + }); + }); +}); diff --git a/packages/durabletask-js/test/in-memory-backend.spec.ts b/packages/durabletask-js/test/in-memory-backend.spec.ts index 2683fb3..6d12ae0 100644 --- a/packages/durabletask-js/test/in-memory-backend.spec.ts +++ b/packages/durabletask-js/test/in-memory-backend.spec.ts @@ -855,6 +855,35 @@ describe("In-Memory Backend", () => { }); describe("suspend and resume status", () => { + it.each([' "maintenance"\n\u6682\u505c ', "", undefined, null])( + "preserves suspend/resume reason %p in history", + async (reason) => { + const orchestrator: TOrchestrator = async function* (ctx: OrchestrationContext) { + yield ctx.waitForExternalEvent("proceed"); + return "done"; + }; + worker.addOrchestrator(orchestrator); + await worker.start(); + const id = await client.scheduleNewOrchestration(orchestrator); + await client.waitForOrchestrationStart(id, false, 10); + + await Reflect.apply(client.suspendOrchestration, client, [id, reason]); + await client.raiseOrchestrationEvent(id, "proceed"); + await Reflect.apply(client.resumeOrchestration, client, [id, reason]); + + const state = await client.waitForOrchestrationCompletion(id, true, 10); + expect(state?.runtimeStatus).toBe(OrchestrationStatus.COMPLETED); + const history = backend.getInstance(id)!.history; + const suspended = history.find((event) => event.hasExecutionsuspended())?.getExecutionsuspended(); + const resumed = history.find((event) => event.hasExecutionresumed())?.getExecutionresumed(); + for (const event of [suspended, resumed]) { + expect(event).toBeDefined(); + expect(event!.hasInput()).toBe(reason != null); + expect(event!.getInput()?.getValue()).toBe(reason ?? undefined); + } + }, + ); + it("should update status to SUSPENDED when suspend is called", async () => { const orchestrator: TOrchestrator = async function* (ctx: OrchestrationContext): any { yield ctx.waitForExternalEvent("proceed");