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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
2 changes: 2 additions & 0 deletions packages/azure-functions-durable/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
13 changes: 6 additions & 7 deletions packages/azure-functions-durable/src/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void> {
async suspend(instanceId: string, reason?: string): Promise<void> {
try {
await this.suspendOrchestration(instanceId);
await this.suspendOrchestration(instanceId, reason);
} catch (error) {
await this._mapControlPlaneError(error, instanceId, "suspend");
}
Expand All @@ -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<void> {
async resume(instanceId: string, reason?: string): Promise<void> {
try {
await this.resumeOrchestration(instanceId);
await this.resumeOrchestration(instanceId, reason);
} catch (error) {
await this._mapControlPlaneError(error, instanceId, "resume");
}
Expand Down
15 changes: 10 additions & 5 deletions packages/azure-functions-durable/test/unit/client.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand Down
22 changes: 20 additions & 2 deletions packages/durabletask-js/src/client/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -572,13 +572,22 @@ export class TaskHubGrpcClient {
);
}

async suspendOrchestration(instanceId: string): Promise<void> {
/**
* 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<void> {
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);

Expand All @@ -589,13 +598,22 @@ export class TaskHubGrpcClient {
);
}

async resumeOrchestration(instanceId: string): Promise<void> {
/**
* 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<void> {
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);

Expand Down
8 changes: 4 additions & 4 deletions packages/durabletask-js/src/testing/in-memory-backend.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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`);
Expand All @@ -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();

Expand All @@ -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`);
Expand All @@ -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();

Expand Down
14 changes: 10 additions & 4 deletions packages/durabletask-js/src/testing/test-client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void> {
this.backend.suspend(instanceId);
async suspendOrchestration(instanceId: string, reason?: string): Promise<void> {
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<void> {
this.backend.resume(instanceId);
async resumeOrchestration(instanceId: string, reason?: string): Promise<void> {
this.backend.resume(instanceId, reason);
}

/**
Expand Down
14 changes: 10 additions & 4 deletions packages/durabletask-js/src/utils/pb-helper.util.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
Expand Down
96 changes: 96 additions & 0 deletions packages/durabletask-js/test/client-suspend-resume.spec.ts
Original file line number Diff line number Diff line change
@@ -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<pb.SuspendRequest, pb.SuspendResponse>,
callback: grpc.sendUnaryData<pb.SuspendResponse>,
) => {
requests.push(call.request);
callback(error, new pb.SuspendResponse());
},
resumeInstance: (
call: grpc.ServerUnaryCall<pb.ResumeRequest, pb.ResumeResponse>,
callback: grpc.sendUnaryData<pb.ResumeResponse>,
) => {
requests.push(call.request);
callback(error, new pb.ResumeResponse());
},
});
const port = await new Promise<number>((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);
});
});
});
29 changes: 29 additions & 0 deletions packages/durabletask-js/test/in-memory-backend.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand Down
Loading