diff --git a/CHANGELOG.md b/CHANGELOG.md index 8156acd..0254862 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,11 @@ # Changelog +## Unreleased + +- Export effect failure/success payload types and `SerializedError` from the + root and core entry points. Check SQL and Cloudflare callback construction + against the same contracts without changing the delivered messages. + ## 0.14.8 - 2026-09-09 - Hold the actor instance and message claim locks through every fenced commit diff --git a/docs/api.md b/docs/api.md index 5c63621..278c3b6 100644 --- a/docs/api.md +++ b/docs/api.md @@ -69,6 +69,10 @@ createdAtMs }` shape returned by `SolidObjectsRuntime.snapshotWithIncarnation`. `PayloadBroadcasts`, and `PayloadBroadcastValue` describe actor-declared transactional work and typed personalized projections. +`EffectFailurePayload`, `EffectSuccessPayload`, +and `SerializedError` describe effect callback messages. They are also exported +from the browser-safe `solid-objects/core` entry point. + `observables()` returns a flat object. Unwrapped values are invalidation-only: their real values participate in change detection, but only their names enter the durable envelope. Use an explicit marker when wire behavior matters: @@ -155,6 +159,60 @@ row per item. It also cannot strand an entry when the runtime coalesces an occurrence. Prefer it for a large queue of interchangeable items. Prefer `key` when one item needs an alarm that you can move on its own. +### Typing your onFailure handler + +An effect callback is an ordinary actor operation. Its payload always includes +the stable `effectId` and the original serialized `arguments`, including `{}` +when the effect was emitted without arguments. Use the exported types when a +watchdog or failure handler needs to correlate work with the current generation: + +```typescript +import { Actor, type EffectFailurePayload, type EffectSuccessPayload } from "solid-objects" + +type RunArguments = { generation: number } + +class ChatRun extends Actor { + static override readonly actorType = "ChatRun" + generation = 0 + status = "idle" + reply = "" + + start(): void { + this.emit("run_model", { + arguments: { generation: ++this.generation }, + onSuccess: "finishTurn", + onFailure: "failTurn", + }) + } + + failTurn({ arguments: original, error }: EffectFailurePayload): void { + if (original.generation !== this.generation) return + this.status = `${error.name}: ${error.message}` + } + + finishTurn({ + arguments: original, + result, + }: EffectSuccessPayload): void { + if (original.generation !== this.generation) return + this.status = "finished" + this.reply = result.reply + } +} +``` + +Failure payloads contain `error: SerializedError`, with string `name` and +`message` fields. They do not include a stack or cause. Success payloads contain +`result`, which can be any JSON value; an undefined effect return becomes +`null`. The default argument type is `JsonObject` and the default success result +type is `JsonValue`. Declare argument shapes with a JSON-compatible type alias. + +These types describe the SQL and Cloudflare callback envelopes. Error messages +for non-Error throws retain each backend's existing serialization behavior. +The generic parameters express your application's contract; they do not add +runtime validation or infer types from `registerEffect()`. Keep registered +effect results and the handler's declared argument/result types in agreement. + ### Runtime managers Every manager below is available as a property on `SolidObjectsRuntime`; the diff --git a/docs/parity.md b/docs/parity.md index c36a4ca..c6e5d2d 100644 --- a/docs/parity.md +++ b/docs/parity.md @@ -54,7 +54,7 @@ is needed for the JavaScript stale-result race fix. | Bounded claim candidate scan | Native | A configurable ordered scan continues to another ready actor when a worker loses the first candidate's lease race. | | Backpressure and payload caps | Partial | Serialization enforces a shared maximum JSON nesting depth, raising `InvalidPayload`, and an optional caller-supplied `maxBytes` limit, raising `PayloadTooLarge`; reminder names are bounded to 255 characters. Distributed per-actor rate limits and global admission control do not exist yet, matching the open Ruby roadmap item. | | Idle activation cache | Native | Long-running workers retain hydrated actors under renewable fenced leases, restore public state after failed turns, and release on timeout, fairness yield, lease loss, or shutdown. | -| Transactional effects and outcome operations | Native | At-least-once handlers receive immutable stable effect, attempt, source-message, and actor identity; success and failure operations also receive the originally staged arguments for correlation. | +| Transactional effects and outcome operations | Native | At-least-once handlers receive immutable stable effect, attempt, source-message, and actor identity; success and failure operations also receive the originally staged arguments for correlation. Typed callback envelopes are exported. | | Actor-to-actor delivery | Native | `sendTo(reference).operation()` stages delivery in the source actor commit. | | One-shot and recurring reminders | Native | Scheduling, replacement events, catch-up policy, stale-claim recovery, pausing, authorized inspection, and idempotent resume are implemented. | | Same-database commit actions | Native | Registered actions receive source-message identity, mailbox sequence, activation generation, and the fenced transaction connection. | @@ -65,6 +65,11 @@ is needed for the JavaScript stale-result race fix. | Result recovery and sync timeout diagnostics | Native | Status, result, and wait reauthorize the stored operation; terminal failure raises structured `MessageFailed`; whole-call adapter deadlines distinguish enqueue, wait, database, activation, and mailbox blockers. | | Result lookup by request ID | Planned | This is also an open Ruby roadmap item and will be implemented in both runtimes when its authorization shape is settled. | +Effect callback envelopes are typed with `EffectFailurePayload`, +`EffectSuccessPayload`, and `SerializedError` in both SQL and Cloudflare. +Ruby RBS contracts are tracked in cardmagic/solid-objects-ruby#64 and preserve +Ruby field names; this does not change runtime delivery semantics. + ## Operations | Capability | Status | TypeScript shape or remaining work | diff --git a/scripts/release-artifact-smoke.mjs b/scripts/release-artifact-smoke.mjs index fb1ef37..4741f4e 100644 --- a/scripts/release-artifact-smoke.mjs +++ b/scripts/release-artifact-smoke.mjs @@ -1,5 +1,5 @@ import assert from "node:assert/strict" -import { mkdtemp, mkdir, readFile, rm } from "node:fs/promises" +import { mkdtemp, mkdir, readFile, rm, writeFile } from "node:fs/promises" import { tmpdir } from "node:os" import { join, resolve } from "node:path" import { spawn } from "node:child_process" @@ -57,6 +57,29 @@ try { ) assert.equal(installedPackage.version, packageDefinition.version) + const consumerPath = join(projectDirectory, "effect-payload-consumer.mts") + await writeFile( + consumerPath, + await readFile(join(repositoryRoot, "test/fixtures/effect-payload-consumer.mts")), + ) + await run( + process.execPath, + [ + join(repositoryRoot, "node_modules/typescript/bin/tsc"), + "--noEmit", + "--strict", + "--noUncheckedIndexedAccess", + "--exactOptionalPropertyTypes", + "--skipLibCheck", + "--module", + "NodeNext", + "--target", + "ES2024", + consumerPath, + ], + { cwd: projectDirectory }, + ) + const resolvedModule = ( await run( process.execPath, diff --git a/src/cloudflare/engine.ts b/src/cloudflare/engine.ts index c833e4f..15d3e27 100644 --- a/src/cloudflare/engine.ts +++ b/src/cloudflare/engine.ts @@ -28,7 +28,13 @@ import { } from "../errors.js" import { deepCopy, jsonObject, normalizeJson, stableJson } from "../serialization.js" import { evaluateActorTurn, readActorObservables, selectActorBroadcast } from "../turn.js" -import type { JsonObject, JsonValue } from "../types.js" +import type { + EffectFailurePayload, + EffectSuccessPayload, + JsonObject, + JsonValue, + SerializedError, +} from "../types.js" import type { CloudflareSettings } from "./configuration.js" import { actorName, callHost, type ActorIdentity, type HostRequest } from "./protocol.js" import type { Instance, Message, Outbox, Reminder, Subscription } from "./records.js" @@ -805,7 +811,7 @@ export class ActorEngine { this.stageEffectCallback({ instance, outbox, - result, + outcome: { result }, operation: outbox.payload.successOperation, }) }) @@ -815,17 +821,18 @@ export class ActorEngine { const exhausted = error instanceof NonRetryableError || outbox.attempt >= this.settings.maxAttempts outbox.status = exhausted ? "dead" : "pending" - outbox.error = { + const errorRecord: SerializedError = { name: errorName(error), message: error instanceof Error ? error.message : "delivery failed", } + outbox.error = errorRecord outbox.availableAt = Date.now() + this.retryDelay(outbox.attempt) this.store.saveOutbox(outbox) if (exhausted && outbox.kind === "effect") this.stageEffectCallback({ instance, outbox, - result: outbox.error, + outcome: { error: errorRecord }, operation: outbox.payload.failureOperation, }) }) @@ -851,7 +858,7 @@ export class ActorEngine { private stageEffectCallback(options: { instance: Instance outbox: Outbox - result: JsonValue + outcome: Pick | Pick operation: JsonValue | undefined }): void { if (typeof options.operation !== "string") return @@ -868,11 +875,9 @@ export class ActorEngine { operation: options.operation, arguments: { effectId: options.outbox.id, - arguments: options.outbox.payload.arguments!, - ...(options.outbox.status === "dead" - ? { error: options.result } - : { result: options.result }), - }, + arguments: options.outbox.payload.arguments as JsonObject, + ...options.outcome, + } satisfies EffectSuccessPayload | EffectFailurePayload, }, }) } diff --git a/src/index.ts b/src/index.ts index 3a28d97..185eb80 100644 --- a/src/index.ts +++ b/src/index.ts @@ -123,6 +123,8 @@ export type { DeepReadonly, DestroyOptions, EffectContext, + EffectFailurePayload, + EffectSuccessPayload, InvocationOptions, JsonObject, JsonPrimitive, @@ -131,6 +133,7 @@ export type { LongRunningComponent, MessageContext, MessageStatus, + SerializedError, SnapshotOptions, } from "./types.js" export type { diff --git a/src/repository.ts b/src/repository.ts index c66a87b..1fafd48 100644 --- a/src/repository.ts +++ b/src/repository.ts @@ -26,7 +26,14 @@ import type { } from "./records.js" import { jsonObject, normalizeJson } from "./serialization.js" import type { RetentionTarget } from "./retention.js" -import type { JsonObject, JsonValue, MessageStatus } from "./types.js" +import type { + EffectFailurePayload, + EffectSuccessPayload, + JsonObject, + JsonValue, + MessageStatus, + SerializedError, +} from "./types.js" import { VERSION } from "./version.js" export interface SyncDiagnosticsRecord { @@ -1486,7 +1493,7 @@ export class Repository { effectId: effect.id, arguments: jsonObject(JSON.parse(effect.arguments)), result, - }), + } satisfies EffectSuccessPayload), idempotencyKey: `effect:${effect.id}:success`, }) void now @@ -1532,7 +1539,7 @@ export class Repository { effectId: effect.id, arguments: jsonObject(JSON.parse(effect.arguments)), error: errorRecord, - }), + } satisfies EffectFailurePayload), idempotencyKey: `effect:${effect.id}:failure`, }) }) @@ -2188,11 +2195,17 @@ function nextReminderRun(options: { return previousRun + (Math.floor((now - previousRun) / interval) + 1) * interval } -function safeError(error: unknown): Record { +function safeError(error: unknown): SerializedError { if (error instanceof Error) { - return jsonObject({ name: error.name, message: error.message }) - } - return jsonObject({ name: "Error", message: String(normalizeJson(error)) }) + return jsonObject({ + name: error.name, + message: error.message, + } satisfies SerializedError) as SerializedError + } + return jsonObject({ + name: "Error", + message: String(normalizeJson(error)), + } satisfies SerializedError) as SerializedError } function retentionPolicy(options: { diff --git a/src/types.ts b/src/types.ts index 8478592..a70880c 100644 --- a/src/types.ts +++ b/src/types.ts @@ -2,6 +2,26 @@ export type JsonPrimitive = null | boolean | number | string export type JsonValue = JsonPrimitive | JsonValue[] | { [key: string]: JsonValue } export type JsonObject = { [key: string]: JsonValue } +export type SerializedError = { + name: string + message: string +} + +export type EffectFailurePayload = { + effectId: string + arguments: Arguments + error: SerializedError +} + +export type EffectSuccessPayload< + Arguments extends JsonObject = JsonObject, + Result extends JsonValue = JsonValue, +> = { + effectId: string + arguments: Arguments + result: Result +} + export type DeepReadonly = Value extends JsonPrimitive ? Value : Value extends readonly (infer Item)[] diff --git a/test/cloudflare/effect-payloads.test.ts b/test/cloudflare/effect-payloads.test.ts new file mode 100644 index 0000000..1ccd733 --- /dev/null +++ b/test/cloudflare/effect-payloads.test.ts @@ -0,0 +1,58 @@ +import { env } from "cloudflare:test" +import { describe, expect, it } from "vitest" +import { createRuntime, durableObjects } from "../../src/cloudflare/index.js" +import { EffectCallbacks, deliveries } from "./worker.js" + +const authorizationContext = "allowed" +const runtime = () => createRuntime({ backend: durableObjects({ namespace: env.ACTORS }) }) + +describe("Cloudflare effect payloads", () => { + it.each([null, false, 42, "reply", ["reply"], { reply: "done" }])( + "delivers the complete success envelope for %j", + async (result) => { + const reference = runtime().ref(EffectCallbacks, `success-${JSON.stringify(result)}`) + const argumentsValue = { generation: 2, nested: { retained: true }, result } + await reference.with({ authorizationContext }).start(argumentsValue) + await expect + .poll(() => + reference.snapshot({ authorizationContext }).then((snapshot) => snapshot.received), + ) + .toEqual([{ effectId: expect.any(String), arguments: argumentsValue, result }]) + const [payload] = (await reference.snapshot({ authorizationContext })).received + expect(deliveries.get(payload!.effectId)).toBe(1) + }, + ) + + it("delivers empty arguments and normalizes an undefined result to null", async () => { + const reference = runtime().ref(EffectCallbacks, "empty") + await reference.with({ authorizationContext }).startEmpty() + await expect + .poll(() => + reference.snapshot({ authorizationContext }).then((snapshot) => snapshot.received), + ) + .toEqual([{ effectId: expect.any(String), arguments: {}, result: null }]) + }) + + it.each([ + { mode: "retry", attempts: 5, name: "Error", message: "exhausted" }, + { mode: "terminal", attempts: 1, name: "NonRetryableError", message: "terminal" }, + { mode: "non-error", attempts: 5, name: "Error", message: "delivery failed" }, + ])("delivers the complete failure envelope for $mode", async (options) => { + const reference = runtime().ref(EffectCallbacks, options.mode) + const argumentsValue = { mode: options.mode, generation: 3 } + await reference.with({ authorizationContext }).start(argumentsValue) + await expect + .poll(() => + reference.snapshot({ authorizationContext }).then((snapshot) => snapshot.received), + ) + .toEqual([ + { + effectId: expect.any(String), + arguments: argumentsValue, + error: { name: options.name, message: options.message }, + }, + ]) + const [payload] = (await reference.snapshot({ authorizationContext })).received + expect(deliveries.get(payload!.effectId)).toBe(options.attempts) + }) +}) diff --git a/test/cloudflare/worker.ts b/test/cloudflare/worker.ts index 709687f..1f4caec 100644 --- a/test/cloudflare/worker.ts +++ b/test/cloudflare/worker.ts @@ -1,4 +1,11 @@ -import { Actor, broadcastValue, NonRetryableError } from "../../src/core.js" +import { + Actor, + broadcastValue, + NonRetryableError, + type EffectFailurePayload, + type EffectSuccessPayload, + type JsonObject, +} from "../../src/core.js" import { PortableCounter } from "../support/portable-actor.js" import { createDurableObjectsHost, @@ -142,8 +149,33 @@ export class VersionedCounter extends Actor { } } +export class EffectCallbacks extends Actor { + static override readonly actorType = "EffectCallbacks" + received: (EffectSuccessPayload | EffectFailurePayload)[] = [] + + start(argumentsValue: JsonObject): void { + this.emit("callbackValue", { + arguments: argumentsValue, + onSuccess: "succeeded", + onFailure: "failed", + }) + } + + startEmpty(): void { + this.emit("callbackEmpty", { onSuccess: "succeeded" }) + } + + succeeded(payload: EffectSuccessPayload): void { + this.received.push(payload) + } + + failed(payload: EffectFailurePayload): void { + this.received.push(payload) + } +} + export class Actors extends createDurableObjectsHost({ - actors: [Counter, PortableCounter, VersionedCounter], + actors: [Counter, PortableCounter, VersionedCounter, EffectCallbacks], configure: (environment) => ({ backend: durableObjects({ namespace: { @@ -172,6 +204,16 @@ export class Actors extends createDurableObjectsHost({ authorizeAdministration: (input) => input.authorizationContext === "allowed", retryDelayMilliseconds: () => 10, effects: { + callbackEmpty: (_arguments, context) => { + deliveries.set(context.id, context.attempt) + }, + callbackValue: (argumentsValue, context) => { + deliveries.set(context.id, context.attempt) + if (argumentsValue.mode === "retry") throw new Error("exhausted") + if (argumentsValue.mode === "terminal") throw new NonRetryableError("terminal") + if (argumentsValue.mode === "non-error") throw "offline" + return argumentsValue.result + }, increment: () => ({ accepted: true }), slow: async (_arguments, context) => { await waitForGate(context.actorId) diff --git a/test/effect-payloads.types.ts b/test/effect-payloads.types.ts new file mode 100644 index 0000000..e0da12f --- /dev/null +++ b/test/effect-payloads.types.ts @@ -0,0 +1,67 @@ +import { expectTypeOf } from "vitest" +import type { + EffectFailurePayload, + EffectSuccessPayload, + JsonObject, + JsonValue, + SerializedError, +} from "../src/index.js" +import type { + EffectFailurePayload as CoreFailurePayload, + EffectSuccessPayload as CoreSuccessPayload, + SerializedError as CoreSerializedError, +} from "../src/core.js" + +type RunArguments = { generation: number } + +export function checkEffectPayloads( + failure: EffectFailurePayload, + success: EffectSuccessPayload, +): void { + expectTypeOf(failure.arguments.generation).toEqualTypeOf() + expectTypeOf(failure.effectId).toEqualTypeOf() + expectTypeOf(failure.error).toEqualTypeOf() + expectTypeOf(failure.error.name).toEqualTypeOf() + expectTypeOf(failure.error.message).toEqualTypeOf() + expectTypeOf(success.result.reply).toEqualTypeOf() + expectTypeOf().toEqualTypeOf() + expectTypeOf().toEqualTypeOf() + expectTypeOf().toEqualTypeOf() + expectTypeOf().toEqualTypeOf() + expectTypeOf().toEqualTypeOf() + + const error: SerializedError = { name: "Error", message: "failed" } + const argumentsValue = { generation: 1 } + const failurePayload: EffectFailurePayload = { + effectId: "effect-1", + arguments: argumentsValue, + error, + } + const results: EffectSuccessPayload[] = [null, false, 1, "reply", [], {}].map( + (result) => ({ effectId: "effect-1", arguments: argumentsValue, result }), + ) + void failurePayload + void results + + // @ts-expect-error Original arguments are required. + const missingArguments: EffectFailurePayload = { effectId: "effect-1", error } + // @ts-expect-error Failure envelopes require the serialized error. + const missingError: EffectFailurePayload = { effectId: "effect-1", arguments: {} } + // @ts-expect-error Success envelopes require a result, including null for no return value. + const missingResult: EffectSuccessPayload = { effectId: "effect-1", arguments: {} } + // @ts-expect-error Error messages are strings. + const wrongError: SerializedError = { name: "Error", message: 42 } + // @ts-expect-error An effect ID is always present. + const missingId: EffectFailurePayload = { arguments: {}, error } + // @ts-expect-error Original arguments must be JSON-compatible. + type InvalidArguments = EffectFailurePayload<{ generation: bigint }> + // @ts-expect-error Results must be JSON-compatible. + type InvalidResult = EffectSuccessPayload + // @ts-expect-error An application's declared generation remains numeric. + failure.arguments.generation = "wrong" + // @ts-expect-error An application's declared result remains typed. + success.result.reply = 42 + void [missingArguments, missingError, missingResult, wrongError, missingId] + expectTypeOf().not.toBeNever() + expectTypeOf().not.toBeNever() +} diff --git a/test/fixtures/effect-payload-consumer.mts b/test/fixtures/effect-payload-consumer.mts new file mode 100644 index 0000000..8847bb6 --- /dev/null +++ b/test/fixtures/effect-payload-consumer.mts @@ -0,0 +1,29 @@ +import type { EffectFailurePayload, EffectSuccessPayload, SerializedError } from "solid-objects" +import type { EffectFailurePayload as CoreFailurePayload } from "solid-objects/core" + +type RunArguments = { generation: number } + +export function failureGeneration(payload: EffectFailurePayload): number { + const corePayload: CoreFailurePayload = payload + const error: SerializedError = corePayload.error + const message: string = error.message + void message + return corePayload.arguments.generation +} + +export function successResult(payload: EffectSuccessPayload): string { + return payload.result +} + +// @ts-expect-error Packaged declarations require the original arguments. +export const invalidFailure: EffectFailurePayload = { + effectId: "id", + error: { name: "Error", message: "failed" }, +} + +export const invalidArguments: EffectSuccessPayload = { + effectId: "id", + // @ts-expect-error Packaged declarations retain application argument types. + arguments: { generation: "wrong" }, + result: null, +} diff --git a/test/outboxes.test.ts b/test/outboxes.test.ts index bc9c724..aa18ecc 100644 --- a/test/outboxes.test.ts +++ b/test/outboxes.test.ts @@ -4,6 +4,7 @@ import type { BroadcastEvent, SolidObjectsConfiguration } from "../src/configura import { NonRetryableError } from "../src/errors.js" import { configure, type SolidObjectsRuntime } from "../src/runtime.js" import { sqlite } from "../src/database/sqlite.js" +import type { EffectFailurePayload, EffectSuccessPayload, JsonObject } from "../src/index.js" class Checkout extends Actor { static override readonly actorType = "Checkout" @@ -24,21 +25,42 @@ class Checkout extends Actor { paymentSucceeded({ arguments: effectArguments, result, - }: { - effectId: string - arguments: { paymentId: string } - result: { receipt: string } - }): void { + }: EffectSuccessPayload<{ paymentId: string }, { receipt: string }>): void { this.status = "paid" this.effectResult = `${effectArguments.paymentId}:${result.receipt}` } - paymentFailed({ arguments: effectArguments }: { arguments: { paymentId: string } }): void { + paymentFailed({ arguments: effectArguments }: EffectFailurePayload<{ paymentId: string }>): void { this.status = "failed" this.failedPaymentId = effectArguments.paymentId } } +class EffectCallbacks extends Actor { + static override readonly actorType = "EffectCallbacks" + received: (EffectSuccessPayload | EffectFailurePayload)[] = [] + + start({ arguments: argumentsValue }: { arguments: JsonObject }): void { + this.emit("callbackValue", { + arguments: argumentsValue, + onSuccess: "succeeded", + onFailure: "failed", + }) + } + + startEmpty(): void { + this.emit("callbackValue", { onSuccess: "succeeded", onFailure: "failed" }) + } + + succeeded(payload: EffectSuccessPayload): void { + this.received.push(payload) + } + + failed(payload: EffectFailurePayload): void { + this.received.push(payload) + } +} + class DatabaseWriter extends Actor { static override readonly actorType = "DatabaseWriter" @@ -121,6 +143,64 @@ afterEach(async () => { }) describe("durable effects", () => { + it.each([ + { result: undefined, expected: null }, + { result: null, expected: null }, + { result: false, expected: false }, + { result: 42, expected: 42 }, + { result: "reply", expected: "reply" }, + { result: ["reply"], expected: ["reply"] }, + { result: { reply: "done" }, expected: { reply: "done" } }, + ])("delivers the complete success envelope for $result", async ({ result, expected }) => { + runtime = configuredRuntime() + let effectId = "" + runtime.registerEffect("callbackValue", (_arguments, context) => { + effectId = context.id + return result + }) + await runtime.install() + const actor = EffectCallbacks.ref("success") + await actor.start({ arguments: { generation: 2, nested: { retained: true } } }) + await runtime.effectWorker().runUntilIdle() + await runtime.worker().runUntilIdle() + + expect(effectId).not.toBe("") + expect(await actor.received).toEqual([ + { effectId, arguments: { generation: 2, nested: { retained: true } }, result: expected }, + ]) + }) + + it.each([ + { error: new Error("exhausted"), attempts: 2, name: "Error", message: "exhausted" }, + { + error: new NonRetryableError("terminal"), + attempts: 1, + name: "NonRetryableError", + message: "terminal", + }, + { error: "offline", attempts: 2, name: "Error", message: "offline" }, + ])("delivers the complete failure envelope for $name/$message", async (options) => { + runtime = configuredRuntime({ maxAttempts: 2, retryDelayMilliseconds: () => 0 }) + let effectId = "" + let attempts = 0 + runtime.registerEffect("callbackValue", (_arguments, context) => { + effectId = context.id + attempts += 1 + throw options.error + }) + await runtime.install() + const actor = EffectCallbacks.ref("failure") + await actor.startEmpty() + await runtime.effectWorker().runUntilIdle() + await runtime.worker().runUntilIdle() + + expect(effectId).not.toBe("") + expect(attempts).toBe(options.attempts) + expect(await actor.received).toEqual([ + { effectId, arguments: {}, error: { name: options.name, message: options.message } }, + ]) + }) + it("correlates concurrent effect callbacks with their staged arguments", async () => { runtime = configuredRuntime() runtime.registerEffect("correlatedEffect", ({ correlationId }) => {