diff --git a/.changeset/actor-error-reporting.md b/.changeset/actor-error-reporting.md new file mode 100644 index 0000000..2a21b49 --- /dev/null +++ b/.changeset/actor-error-reporting.md @@ -0,0 +1,5 @@ +--- +"effect-machine": minor +--- + +Report contained actor failures to the Effect `ErrorReporter`s that were registered where the actor was spawned. Each runtime generation reports its complete defect cause once when it closes, including transition, spawn, task, background, and cleanup failures. Supervised restarts report each generation. A restart step that fails before a new generation exists reports with phase `restart`. This includes a restart recovery and a restart schedule that dies; an exhausted schedule does not report. A cold-start recovery that fails before the first generation runs reports with phase `recovery`, also when a stop interrupts the start. A recovery that throws instead of returning an Effect reports the same way. A schedule defect that settles together with schedule exhaustion still reports. A final output defect reports when the actor completes. Child actors report through the reporters of the handler that spawned them. Normal stops, final states, and pure interruption do not report. Each report annotates the reporting fiber with the actor ID, generation, and defect phase. The actor keeps its spawn-time reporters, so a later `start` or `stop` caller cannot replace them. Handler code inside the actor also runs with the spawn-time reporters, so `Effect.withErrorReporting` in a handler reaches them and no longer reaches reporters that were provided only around `start`. A failure that also reaches a caller boundary, such as `Effect.withErrorReporting`, an RPC server or an HTTP handler, can report again there. A reporter built with `ErrorReporter.make` skips the second report only when the defect is an object; a primitive defect or a raw reporter reports twice. A parent handler that re-raises a child failure, such as a failed child stop or a failed `self.spawn` start, reports that cause again as the parent's defect, with the same deduplication rules. A throwing reporter does not change the actor exit. Effect calls the reporters in order, so a throwing reporter can stop the reporters after it. Cluster entity runtimes report generation defects through the reporters in their allocation context. The generation of an entity report names one allocation, so it is `0` after each reactivation. Without registered reporters, behavior does not change. diff --git a/AGENTS.md b/AGENTS.md index 34393a3..014cd56 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -203,6 +203,14 @@ const otherExit = yield * other.awaitExit; // ActorExit - **Classifier** — `shouldRestart` optionally skips restart for specific defect types - Entity-machine: cluster-supervised via `defectRetryPolicy`, NOT local supervision +## Error Reporting + +- Lifecycle owners report to the Effect `ErrorReporter`s captured at spawn: generation closure reports its combined defect, terminal completion reports a final output defect, and the supervisor reports a restart step that fails before a new generation exists. +- `createActor` pins `CurrentErrorReporters` in the actor service context. Restarted generations allocate in the supervisor fiber, which also carries the start caller's context; do not let that caller replace the spawn-time reporters. +- Report before publishing the exit, and isolate reporter defects. A throwing reporter must not block exit settlement. +- Inside the actor runtime, do not report a cause again where it is only re-raised: stop failures belong to the generation, and implicit-system close failures belong to the child actors. A parent handler that re-raises a child failure is a new parent defect, so the parent generation reports it too. That overlap, like the overlap with caller boundaries, relies on `ErrorReporter.make` identity dedup and is documented as such. +- Wrap user callbacks that return an Effect (such as `lifecycle.recovery.resolve`) in `Effect.suspend` before attaching report hooks, so a synchronous throw becomes a reported defect. + ## Child Actors Spawn children from `.spawn()`/`.background()` handlers via `self.spawn(id, childMachine)`: diff --git a/SKILL.md b/SKILL.md index 37603f6..5cabff4 100644 --- a/SKILL.md +++ b/SKILL.md @@ -313,6 +313,7 @@ An exit animation can keep a screen mounted after a final transition. Keep requi 8. **call vs send**: `send` = fire-and-forget, `call` = request-reply, `ask` = typed reply 9. **Non-Effect code**: Use `actor.client`. React and Solid should use Actor Atoms. 10. **ActorStoppedError**: Pending `call`/`ask` Deferreds settled on stop +11. **Error reporters at spawn**: provide `ErrorReporter.layer` where the actor is spawned. Lifecycle owners report each generation defect, restart-step failure, and final output defect to those reporters. Pure interruption does not report. ## Cluster / Entity Machines diff --git a/docs/persistence-and-supervision.md b/docs/persistence-and-supervision.md index 721ea24..a863355 100644 --- a/docs/persistence-and-supervision.md +++ b/docs/persistence-and-supervision.md @@ -51,6 +51,31 @@ Machine.spawn(machine, { See [`supervision.ts`](../examples/core/src/supervision.ts). +## Error reporting + +The actor lifecycle reports contained failures to Effect `ErrorReporter`s. Register the reporters in the context that spawns the actor: + +```ts +Machine.spawn(machine).pipe(Effect.provide(ErrorReporter.layer([reporter]))); +``` + +- Each generation reports its complete defect once, when it closes. The cause includes transition, spawn, task, background, and cleanup failures of that generation. +- Each supervised restart reports its own generation. A restart step that fails before a new generation exists reports with phase `restart`. This includes a restart recovery that dies and a restart schedule that dies. An exhausted schedule does not report. +- A cold-start `lifecycle.recovery.resolve` that fails before the first generation runs reports with phase `recovery`. It also reports when a stop interrupts the start. +- A final output defect reports once when the actor completes. +- Child actors report their own failures. They capture the reporters of the parent handler that spawned them. +- Normal stops, final states, and pure interruption do not report. +- The actor keeps the reporters that it captured at spawn. A later `start` or `stop` caller cannot add or replace them. +- Handler code inside the actor also runs with the spawn-time reporters. `Effect.withErrorReporting` in a handler reaches them. It does not reach reporters that were provided only around `start`. +- Each report annotates the reporting fiber with `effect_machine.actor.id`, `effect_machine.actor.generation`, and `effect_machine.defect.phase`. Reporters read them from `References.CurrentLogAnnotations`. +- A failure that also reaches a caller can report again at the caller's boundary, for example `Effect.withErrorReporting`, an RPC server or an HTTP handler. Use `ErrorReporter.make`. It skips a cause or an error object that it already reported. A primitive defect, such as `Effect.die("boom")`, is not an object, so it reports twice. A raw reporter that is not built with `ErrorReporter.make` also reports twice. +- A parent handler that re-raises a child failure fails the parent too, so the parent reports that cause again as its own defect. Examples are a scoped child whose stop fails and a `self.spawn` whose start fails. `ErrorReporter.make` skips the second report when the defect is an object. A primitive defect or a raw reporter reports twice. +- A cold-start recovery that throws instead of returning an Effect reports like a recovery that dies. +- A reporter that throws does not change the actor exit. Effect calls the reporters of one set in order, so a throwing reporter can stop the reporters after it. +- A cluster entity report names one allocation. Its generation is `0` after each reactivation. + +Inspection `@machine.error` events stay diagnostics. They do not replace reporting. + ## Local and cluster durability Local actors use lifecycle hooks. Entity machines use the cluster persistence adapter. diff --git a/src/actor.ts b/src/actor.ts index b8a74f7..027ee77 100644 --- a/src/actor.ts +++ b/src/actor.ts @@ -10,14 +10,17 @@ import { Deferred, Cause, Effect, + ErrorReporter, Exit, Fiber, Layer, MutableHashMap, Option, PubSub, + Pull, Queue, Ref, + Result, Schedule, Scope, Semaphore, @@ -44,6 +47,7 @@ import { DuplicateActorError, ActorStoppedError } from "./errors.js"; import { createRuntime, notifyStateListeners, + reportLifecycleFailure, type RuntimeLifecycleHooks, type RuntimeQueuedEvent, type RuntimeHandle, @@ -693,6 +697,7 @@ function activateGeneration( /** * Run the supervision loop for a supervised actor. * Observes exit deferred, applies restart policy, resets cell resources on restart. + * Returns the terminal generation exit. The actor owner completes it. * @internal */ const runSupervisionLoop = < @@ -709,9 +714,10 @@ const runSupervisionLoop = < ) => Effect.Effect>; lifecycle?: Lifecycle; onRestart?: (generation: number, exit: ActorExit) => Effect.Effect; - onTerminal: (exit: RuntimeExit) => Effect.Effect; + /** The reporters captured at spawn. */ + errorReporters: ReadonlySet; }, -) => +): Effect.Effect | undefined> => Effect.gen(function* () { const step = yield* Schedule.toStepWithSleep(options.supervision.schedule); @@ -723,24 +729,30 @@ const runSupervisionLoop = < const generationExit = yield* Deferred.await(currentRuntime.exitDeferred); yield* currentRuntime.awaitClosed; - if (generationExit._tag !== "Defect") { - yield* options.onTerminal(generationExit); - return; - } + if (generationExit._tag !== "Defect") return generationExit; if ( options.supervision.shouldRestart !== undefined && !options.supervision.shouldRestart(generationExit) ) { - yield* options.onTerminal(generationExit); - return; + return generationExit; } const pull = step(generationExit); const scheduleExit = yield* pull.pipe(Effect.exit); if (scheduleExit._tag === "Failure") { - yield* options.onTerminal(generationExit); - return; + // An exhausted schedule halts its pull. A defect can settle with that halt, for example in a + // finalizer, so report whatever remains once the halt is removed. + const scheduleFailure = Pull.filterDone(scheduleExit.cause); + if (Result.isFailure(scheduleFailure)) { + yield* reportLifecycleFailure(scheduleFailure.failure, { + reporters: options.errorReporters, + actorId: cell.id, + generation: cell.generation.current, + phase: "restart", + }); + } + return generationExit; } // Bump generation before restart — recovery.resolve sees the new generation @@ -817,6 +829,7 @@ export const createActor = Effect.fn("effect-machine.actor.spawn")(function* < ) { const lifecycle: Lifecycle | undefined = options.lifecycle; const capturedContext = yield* Effect.context(); + const errorReporters = yield* ErrorReporter.CurrentErrorReporters; // Spawn is cold. The caller has already resolved machine input and hydration. // Recovery runs during start, not allocate. @@ -832,7 +845,12 @@ export const createActor = Effect.fn("effect-machine.actor.spawn")(function* < const localInspector = options.inspect ?? ambientInspector; const systemInspectors = systemInspectorsBySystem.get(system) ?? new Set(); const inspectorValue = makeInspectionDispatcher(localInspector, systemInspectors); - const serviceContext = Context.add(capturedContext, ActorInspection, inspectorValue); + // Pin the spawn-time reporters. Restarted generations allocate in the supervisor + // fiber, which also carries the start caller's context. + const serviceContext = capturedContext.pipe( + Context.add(ActorInspection, inspectorValue), + Context.add(ErrorReporter.CurrentErrorReporters, errorReporters), + ); // Actor-specific state const childrenMap = new Map>(); @@ -929,6 +947,16 @@ export const createActor = Effect.fn("effect-machine.actor.spawn")(function* < terminalCause = cause; terminal = { _tag: "Defect", cause, phase: "cleanup" }; } + // Generations report their own failures, and children report theirs. Only the + // output conversion belongs to the actor owner. + if (Exit.isFailure(converted)) { + yield* reportLifecycleFailure(converted.cause, { + reporters: errorReporters, + actorId: id, + generation: generation.current, + phase: "cleanup", + }); + } yield* SubscriptionRef.set(lifecycleRef, terminal); yield* Deferred.succeed(terminalExitDeferred, terminal); if (terminalCause !== undefined) { @@ -1141,11 +1169,26 @@ export const createActor = Effect.fn("effect-machine.actor.spawn")(function* < }); // Run recovery if lifecycle.recovery exists AND not hydrated (hydrate takes precedence) if (lifecycle?.recovery !== undefined && !isHydrated) { - const resolved = yield* lifecycle.recovery.resolve({ - actorId: id, - generation: generation.current, - machineInitial: options.machineInitial, - }); + const recovery = lifecycle.recovery; + // `suspend` turns a resolver that throws into a defect that `onError` can see. + const resolved = yield* Effect.suspend(() => + recovery.resolve({ + actorId: id, + generation: generation.current, + machineInitial: options.machineInitial, + }), + ).pipe( + // No generation has run yet, so no generation closure can report this failure. + // `onError` reports even when a stop is interrupting the start. + Effect.onError((cause) => + reportLifecycleFailure(cause, { + reporters: errorReporters, + actorId: id, + generation: generation.current, + phase: "recovery", + }), + ), + ); if (stopRequested) return yield* Effect.interrupt; if (Option.isSome(resolved)) { // Update cell stateRef @@ -1176,8 +1219,23 @@ export const createActor = Effect.fn("effect-machine.actor.spawn")(function* < spawnGeneration, lifecycle, onRestart: options.onRestart, - onTerminal: completeTerminal, - }), + errorReporters, + }).pipe( + // A restart step can fail before a new generation exists to report it. `onError` + // reports even when a stop is interrupting the supervisor. + Effect.onError((cause) => + reportLifecycleFailure(cause, { + reporters: errorReporters, + actorId: id, + generation: generation.current, + phase: "restart", + }), + ), + Effect.flatMap((terminalExit) => { + if (terminalExit === undefined) return Effect.void; + return completeTerminal(terminalExit); + }), + ), ); } else { // No supervision — wire terminal exit from the current generation diff --git a/src/internal/runtime.ts b/src/internal/runtime.ts index 34cfced..69492e5 100644 --- a/src/internal/runtime.ts +++ b/src/internal/runtime.ts @@ -21,6 +21,7 @@ import { Cause, Deferred, Effect, + ErrorReporter, Exit, Fiber, Queue, @@ -126,6 +127,40 @@ const RuntimeExit = { }), }; +/** Spawn-time reporters and lifecycle owner context for one failure report. */ +interface LifecycleFailureReport { + readonly reporters: ReadonlySet; + readonly actorId: string; + readonly generation: number; + /** + * `restart` marks a supervision step that failed before a new generation existed. + * `recovery` marks a cold-start recovery that failed before the first generation ran. + */ + readonly phase: DefectPhase | "restart" | "recovery"; +} + +/** + * Report a failure that an actor lifecycle owner settles. The reporters come + * from the spawn context, not from whichever fiber closes the actor. + * @internal + */ +export const reportLifecycleFailure: { + (report: LifecycleFailureReport): (cause: Cause.Cause) => Effect.Effect; + (cause: Cause.Cause, report: LifecycleFailureReport): Effect.Effect; +} = dual(2, (cause: Cause.Cause, report: LifecycleFailureReport): Effect.Effect => { + if (report.reporters.size === 0 || Cause.hasInterruptsOnly(cause)) return Effect.void; + return ErrorReporter.report(cause).pipe( + Effect.annotateLogs({ + "effect_machine.actor.id": report.actorId, + "effect_machine.actor.generation": report.generation, + "effect_machine.defect.phase": report.phase, + }), + Effect.provideService(ErrorReporter.CurrentErrorReporters, report.reporters), + // Reporters are host callbacks. A throwing reporter must not stop lifecycle settlement. + Effect.ignoreCause, + ); +}); + /** @internal */ export interface RuntimeHandle { /** Enqueue a fire-and-forget event */ @@ -247,6 +282,7 @@ export const createRuntime = Effect.fn("effect-machine.runtime.create")(function // Capture services at allocation so delayed start and stop retain them. const services = yield* Effect.context(); + const errorReporters = yield* ErrorReporter.CurrentErrorReporters; const fork = Effect.runForkWith(services); const { stateRef, latestTransitionRef, stoppedRef, eventQueue, listeners } = config.cellResources; @@ -424,15 +460,24 @@ export const createRuntime = Effect.fn("effect-machine.runtime.create")(function } else { closed = Exit.failCause(cleanupCause); } + let finalExit = selectedReason; if (Exit.isFailure(closed)) { let cause = closed.cause; if (selectedReason._tag === "Defect") { cause = Cause.combine(closed.cause)(selectedReason.cause); } - yield* setExit(RuntimeExit.Defect(cause, "cleanup")); - } else { - yield* setExit(selectedReason); + finalExit = RuntimeExit.Defect(cause, "cleanup"); + } + // Report the complete generation failure before any waiter can observe the exit. + if (finalExit._tag === "Defect") { + yield* reportLifecycleFailure(finalExit.cause, { + reporters: errorReporters, + actorId, + generation, + phase: finalExit.phase, + }); } + yield* setExit(finalExit); yield* Deferred.succeed(closedDeferred, closed); if (Exit.isFailure(closed)) { return yield* Effect.failCause(closed.cause).pipe(Effect.orDie); diff --git a/test/error-reporting.test.ts b/test/error-reporting.test.ts new file mode 100644 index 0000000..8912bc8 --- /dev/null +++ b/test/error-reporting.test.ts @@ -0,0 +1,831 @@ +// @effect-diagnostics strictEffectProvide:off - tests are entry points +// @effect-diagnostics anyUnknownInErrorContext:off +import { + Cause, + Data, + Effect, + ErrorReporter, + Layer, + References, + Schedule, + type Scope, +} from "effect"; +import { Entity, ShardingConfig } from "effect/cluster"; +import { describe, expect, it } from "effect-bun-test"; + +import { EntityMachine, toEntity } from "../src/cluster/index.js"; + +import { + ActorSystemDefault, + Event, + Machine, + State, + Supervision, + type ActorExit, + type DefectPhase, +} from "../src/index.js"; + +// ============================================================================ +// Fixtures +// ============================================================================ + +/** A defect with object identity. Equal messages must still report separately. */ +class LifecycleDefect extends Data.TaggedError( + "effect-machine/test/error-reporting.test/LifecycleDefect", +)<{ readonly message: string }> {} + +const S = State({ Idle: {}, Active: {}, Done: {} }); +const E = Event({ Start: {}, Crash: {}, Finish: {} }); + +interface Report { + readonly cause: Cause.Cause; + readonly annotations: Readonly>; +} + +/** A raw reporter. It records every call, so it can prove the absence of duplicate reports. */ +const makeRecorder = () => { + const reports: Report[] = []; + const reporter: ErrorReporter.ErrorReporter = { + [ErrorReporter.TypeId]: ErrorReporter.TypeId, + report: ({ cause, fiber }) => { + reports.push({ cause, annotations: fiber.getRef(References.CurrentLogAnnotations) }); + }, + }; + return { reporter, reports, layer: ErrorReporter.layer([reporter]) }; +}; + +const onlyReport = (reports: ReadonlyArray): Effect.Effect => { + expect(reports).toHaveLength(1); + const [report] = reports; + if (report === undefined) return Effect.die("expected one report"); + return Effect.succeed(report); +}; + +const defectsOf = (cause: Cause.Cause): ReadonlyArray => + cause.reasons.filter(Cause.isDieReason).map((reason) => reason.defect); + +/** Poll a condition that another fiber settles. */ +const eventually = (predicate: () => boolean) => + Effect.sleep("1 millis").pipe( + Effect.repeat({ until: predicate }), + Effect.timeout("1 second"), + Effect.orDie, + ); + +const annotationsFor = (actorId: string, generation: number, phase: string) => ({ + "effect_machine.actor.id": actorId, + "effect_machine.actor.generation": generation, + "effect_machine.defect.phase": phase, +}); + +interface PhaseCase { + readonly name: string; + readonly actorId: string; + readonly phase: DefectPhase; + /** Spawns an actor, makes it fail with `failure`, and returns its terminal exit. */ + readonly run: ( + failure: LifecycleDefect, + ) => Effect.Effect, never, Scope.Scope>; +} + +const phaseCases: ReadonlyArray = [ + { + name: "a transition defect", + actorId: "transition", + phase: "transition", + run: (failure) => + Effect.gen(function* () { + const machine = Machine.make({ state: S, event: E, initial: S.Idle }).on( + S.Idle, + E.Crash, + () => Effect.die(failure), + ); + const actor = yield* Machine.spawn(machine, { id: "transition" }); + yield* actor.start; + yield* actor.send(E.Crash); + return yield* actor.awaitExit; + }), + }, + { + name: "an initial spawn defect", + actorId: "initial-spawn", + phase: "initial-spawn", + run: (failure) => + Effect.gen(function* () { + const machine = Machine.make({ state: S, event: E, initial: S.Idle }).spawn(S.Idle, () => + Effect.die(failure), + ); + const actor = yield* Machine.spawn(machine, { id: "initial-spawn" }); + yield* Effect.exit(actor.start); + return yield* actor.awaitExit; + }), + }, + { + name: "a spawn effect defect", + actorId: "spawn", + phase: "spawn", + run: (failure) => + Effect.gen(function* () { + const machine = Machine.make({ state: S, event: E, initial: S.Idle }) + .on(S.Idle, E.Start, () => S.Active) + .spawn(S.Active, () => Effect.yieldNow.pipe(Effect.andThen(Effect.die(failure)))); + const actor = yield* Machine.spawn(machine, { id: "spawn" }); + yield* actor.start; + yield* actor.send(E.Start); + return yield* actor.awaitExit; + }), + }, + { + name: "a task defect", + actorId: "task", + phase: "spawn", + run: (failure) => + Effect.gen(function* () { + const machine = Machine.make({ state: S, event: E, initial: S.Idle }) + .on(S.Idle, E.Start, () => S.Active) + .on(S.Active, E.Finish, () => S.Done) + .task(S.Active, () => Effect.die(failure), { onSuccess: () => E.Finish }) + .final(S.Done); + const actor = yield* Machine.spawn(machine, { id: "task" }); + yield* actor.start; + yield* actor.send(E.Start); + return yield* actor.awaitExit; + }), + }, + { + name: "a background defect", + actorId: "background", + phase: "background", + run: (failure) => + Effect.gen(function* () { + const machine = Machine.make({ state: S, event: E, initial: S.Idle }).background(() => + Effect.yieldNow.pipe(Effect.andThen(Effect.die(failure))), + ); + const actor = yield* Machine.spawn(machine, { id: "background" }); + yield* actor.start; + return yield* actor.awaitExit; + }), + }, + { + name: "a cleanup defect during stop", + actorId: "cleanup", + phase: "cleanup", + run: (failure) => + Effect.gen(function* () { + const machine = Machine.make({ state: S, event: E, initial: S.Idle }).spawn(S.Idle, () => + Effect.addFinalizer(() => Effect.die(failure)), + ); + const actor = yield* Machine.spawn(machine, { id: "cleanup" }); + yield* actor.start; + yield* Effect.exit(actor.stop); + return yield* actor.awaitExit; + }), + }, +]; + +// ============================================================================ +// Generation closure +// ============================================================================ + +describe("error reporting: generation closure", () => { + for (const phaseCase of phaseCases) { + it.scopedLive(`reports ${phaseCase.name} once with the settled cause`, () => + Effect.gen(function* () { + const recorder = makeRecorder(); + const failure = new LifecycleDefect({ message: phaseCase.name }); + const exit = yield* phaseCase.run(failure).pipe(Effect.provide(recorder.layer)); + + expect(exit._tag).toBe("Defect"); + if (exit._tag !== "Defect") return; + expect(exit.phase).toBe(phaseCase.phase); + const report = yield* onlyReport(recorder.reports); + expect(report.cause).toBe(exit.cause); + expect(defectsOf(report.cause)).toContain(failure); + expect(report.annotations).toMatchObject( + annotationsFor(phaseCase.actorId, 0, phaseCase.phase), + ); + }), + ); + } + + it.scopedLive("reports a transition defect and a cleanup defect in one aggregate", () => + Effect.gen(function* () { + const recorder = makeRecorder(); + const transitionFailure = new LifecycleDefect({ message: "same message" }); + const cleanupFailure = new LifecycleDefect({ message: "same message" }); + const machine = Machine.make({ state: S, event: E, initial: S.Idle }) + .spawn(S.Idle, () => Effect.addFinalizer(() => Effect.die(cleanupFailure))) + .on(S.Idle, E.Crash, () => Effect.die(transitionFailure)); + const exit = yield* Effect.gen(function* () { + const actor = yield* Machine.spawn(machine, { id: "aggregate" }); + yield* actor.start; + yield* actor.send(E.Crash); + return yield* actor.awaitExit; + }).pipe(Effect.provide(recorder.layer)); + + expect(exit._tag).toBe("Defect"); + if (exit._tag !== "Defect") return; + expect(exit.phase).toBe("cleanup"); + const report = yield* onlyReport(recorder.reports); + expect(report.cause).toBe(exit.cause); + const defects = defectsOf(report.cause); + expect(defects).toHaveLength(2); + expect(defects[0]).toBe(transitionFailure); + expect(defects[1]).toBe(cleanupFailure); + expect(report.annotations).toMatchObject(annotationsFor("aggregate", 0, "cleanup")); + }), + ); + + it.scopedLive("keeps equal-message failures distinct in a native reporter", () => + Effect.gen(function* () { + const messages: string[] = []; + const reporter = ErrorReporter.make(({ error }) => { + messages.push(error.message); + }); + const machine = Machine.make({ state: S, event: E, initial: S.Idle }) + .spawn(S.Idle, () => + Effect.addFinalizer(() => Effect.die(new LifecycleDefect({ message: "same message" }))), + ) + .on(S.Idle, E.Crash, () => Effect.die(new LifecycleDefect({ message: "same message" }))); + yield* Effect.gen(function* () { + const actor = yield* Machine.spawn(machine); + yield* actor.start; + yield* actor.send(E.Crash); + yield* actor.awaitExit; + }).pipe(Effect.provide(ErrorReporter.layer([reporter]))); + + expect(messages).toEqual(["same message", "same message"]); + }), + ); + + it.scopedLive("does not report a failure twice when the start caller also reports it", () => + Effect.gen(function* () { + const messages: string[] = []; + const reporter = ErrorReporter.make(({ error }) => { + messages.push(error.message); + }); + const failure = new LifecycleDefect({ message: "initial spawn failure" }); + const machine = Machine.make({ state: S, event: E, initial: S.Idle }).spawn(S.Idle, () => + Effect.die(failure), + ); + const started = yield* Effect.gen(function* () { + const actor = yield* Machine.spawn(machine); + return yield* actor.start.pipe(Effect.withErrorReporting, Effect.exit); + }).pipe(Effect.provide(ErrorReporter.layer([reporter]))); + + expect(started._tag).toBe("Failure"); + expect(messages).toEqual(["initial spawn failure"]); + }), + ); + + it.scopedLive("does not report normal stops or final states", () => + Effect.gen(function* () { + const recorder = makeRecorder(); + const machine = Machine.make({ state: S, event: E, initial: S.Idle }) + .on(S.Idle, E.Finish, () => S.Done) + .spawn(S.Idle, () => Effect.never) + .background(() => Effect.never) + .final(S.Done); + const exits = yield* Effect.gen(function* () { + const stopped = yield* Machine.spawn(machine); + yield* stopped.start; + yield* stopped.stop; + const finished = yield* Machine.spawn(machine); + yield* finished.start; + yield* finished.send(E.Finish); + return [yield* stopped.awaitExit, yield* finished.awaitExit]; + }).pipe(Effect.provide(recorder.layer)); + + expect(exits.map((exit) => exit._tag)).toEqual(["Stopped", "Final"]); + expect(recorder.reports).toHaveLength(0); + }), + ); + + it.scopedLive("does not report a handler that interrupts itself", () => + Effect.gen(function* () { + const recorder = makeRecorder(); + const machine = Machine.make({ state: S, event: E, initial: S.Idle }).on( + S.Idle, + E.Crash, + () => Effect.interrupt, + ); + const exit = yield* Effect.gen(function* () { + const actor = yield* Machine.spawn(machine, { id: "self-interrupt" }); + yield* actor.start; + yield* actor.send(E.Crash); + return yield* actor.awaitExit.pipe(Effect.timeout("1 second")); + }).pipe(Effect.provide(recorder.layer)); + + // The generation still ends as a defect. Only its report is withheld. + expect(exit._tag).toBe("Defect"); + if (exit._tag !== "Defect") return; + expect(exit.cause.reasons.map((reason) => reason._tag)).toEqual(["Interrupt"]); + expect(recorder.reports).toHaveLength(0); + }), + ); + + it.scopedLive("does not report a restart recovery that interrupts itself", () => + Effect.gen(function* () { + const recorder = makeRecorder(); + const failure = new LifecycleDefect({ message: "transition" }); + const machine = Machine.make({ state: S, event: E, initial: S.Idle }).on( + S.Idle, + E.Crash, + () => Effect.die(failure), + ); + yield* Effect.gen(function* () { + const actor = yield* Machine.spawn(machine, { + id: "interrupted-recovery", + supervision: Supervision.restart({ maxRestarts: 1 }), + lifecycle: { + recovery: { + resolve: ({ generation }) => { + if (generation === 0) return Effect.succeedNone; + return Effect.interrupt; + }, + }, + }, + }); + yield* actor.start; + yield* actor.send(E.Crash); + yield* eventually(() => recorder.reports.length >= 1); + // The restart step runs after the generation report. Give it time to report wrongly. + yield* Effect.sleep("50 millis"); + yield* Effect.exit(actor.stop); + }).pipe(Effect.provide(recorder.layer)); + + expect(recorder.reports.map((report) => defectsOf(report.cause))).toEqual([[failure]]); + }), + ); +}); + +// ============================================================================ +// Cold-start recovery +// ============================================================================ + +const coldStartCases = [ + { label: "an unsupervised", supervisionOptions: {} }, + { + label: "a supervised", + supervisionOptions: { supervision: Supervision.restart({ maxRestarts: 1 }) }, + }, +]; + +describe("error reporting: cold-start recovery", () => { + for (const { label, supervisionOptions } of coldStartCases) { + it.scopedLive(`reports a recovery defect that fails the start of ${label} actor`, () => + Effect.gen(function* () { + const recorder = makeRecorder(); + const failure = new LifecycleDefect({ message: "cold recovery" }); + const machine = Machine.make({ state: S, event: E, initial: S.Idle }); + const startExit = yield* Effect.gen(function* () { + const actor = yield* Machine.spawn(machine, { + id: "cold-start", + ...supervisionOptions, + lifecycle: { recovery: { resolve: () => Effect.die(failure) } }, + }); + return yield* Effect.exit(actor.start); + }).pipe(Effect.provide(recorder.layer)); + + expect(startExit._tag).toBe("Failure"); + const report = yield* onlyReport(recorder.reports); + expect(defectsOf(report.cause)).toEqual([failure]); + expect(report.annotations).toMatchObject(annotationsFor("cold-start", 0, "recovery")); + }), + ); + + it.scopedLive( + `reports a recovery that throws while it builds the start of ${label} actor`, + () => + Effect.gen(function* () { + const recorder = makeRecorder(); + const failure = new LifecycleDefect({ message: "synchronous recovery" }); + const machine = Machine.make({ state: S, event: E, initial: S.Idle }); + const startExit = yield* Effect.gen(function* () { + const actor = yield* Machine.spawn(machine, { + id: "throwing-recovery", + ...supervisionOptions, + lifecycle: { + recovery: { + resolve: () => { + // oxlint-disable-next-line effect/noThrowStatement -- a resolver that throws instead of returning an Effect is the defect under test. + throw failure; + }, + }, + }, + }); + return yield* Effect.exit(actor.start); + }).pipe(Effect.provide(recorder.layer)); + + expect(startExit._tag).toBe("Failure"); + const report = yield* onlyReport(recorder.reports); + expect(defectsOf(report.cause)).toEqual([failure]); + expect(report.annotations).toMatchObject( + annotationsFor("throwing-recovery", 0, "recovery"), + ); + }), + ); + + it.scopedLive(`reports a recovery defect that ends while stop cancels ${label} start`, () => + Effect.gen(function* () { + const recorder = makeRecorder(); + const failure = new LifecycleDefect({ message: "recovery during stop" }); + const machine = Machine.make({ state: S, event: E, initial: S.Idle }); + const exit = yield* Effect.gen(function* () { + const actor = yield* Machine.spawn(machine, { + id: "stopped-recovery", + ...supervisionOptions, + lifecycle: { + recovery: { + // Stop cannot cancel this step, so its defect settles while stop waits. + resolve: () => + Effect.sleep("30 millis").pipe( + Effect.andThen(Effect.die(failure)), + Effect.uninterruptible, + ), + }, + }, + }); + yield* Effect.forkChild(Effect.exit(actor.start)); + yield* Effect.sleep("5 millis"); + yield* Effect.exit(actor.stop); + return yield* actor.awaitExit.pipe(Effect.timeout("1 second")); + }).pipe(Effect.provide(recorder.layer)); + + expect(exit._tag).toBe("Defect"); + if (exit._tag !== "Defect") return; + expect(exit.phase).toBe("cleanup"); + expect(defectsOf(exit.cause)).toEqual([failure]); + const report = yield* onlyReport(recorder.reports); + expect(defectsOf(report.cause)).toEqual([failure]); + expect(report.annotations).toMatchObject(annotationsFor("stopped-recovery", 0, "recovery")); + }), + ); + } +}); + +// ============================================================================ +// Captured spawn context and restarts +// ============================================================================ + +describe("error reporting: captured spawn context", () => { + it.scopedLive("reports every restarted generation to the reporters captured at spawn", () => + Effect.gen(function* () { + const spawnRecorder = makeRecorder(); + const callerRecorder = makeRecorder(); + const failures: LifecycleDefect[] = []; + const machine = Machine.make({ state: S, event: E, initial: S.Idle }).spawn(S.Idle, () => + Effect.suspend(() => { + const failure = new LifecycleDefect({ message: `generation ${failures.length}` }); + failures.push(failure); + return Effect.die(failure); + }), + ); + const actor = yield* Machine.spawn(machine, { + id: "restarting", + supervision: Supervision.restart({ maxRestarts: 2 }), + }).pipe(Effect.provide(spawnRecorder.layer)); + const exit = yield* Effect.gen(function* () { + yield* actor.start; + return yield* actor.awaitExit; + }).pipe(Effect.provide(callerRecorder.layer)); + + expect(exit._tag).toBe("Defect"); + expect(failures).toHaveLength(3); + expect(callerRecorder.reports).toHaveLength(0); + expect(spawnRecorder.reports).toHaveLength(3); + spawnRecorder.reports.forEach((report, generation) => { + const defects = defectsOf(report.cause); + expect(defects).toHaveLength(1); + expect(defects[0]).toBe(failures[generation]); + expect(report.annotations).toMatchObject( + annotationsFor("restarting", generation, "initial-spawn"), + ); + }); + }), + ); + + it.scopedLive("does not borrow a caller's reporters when none were captured at spawn", () => + Effect.gen(function* () { + const callerRecorder = makeRecorder(); + const machine = Machine.make({ state: S, event: E, initial: S.Idle }).spawn(S.Idle, () => + Effect.die(new LifecycleDefect({ message: "initial spawn" })), + ); + const actor = yield* Machine.spawn(machine, { + supervision: Supervision.restart({ maxRestarts: 1 }), + }); + const exit = yield* Effect.gen(function* () { + yield* actor.start; + return yield* actor.awaitExit; + }).pipe(Effect.provide(callerRecorder.layer)); + + expect(exit._tag).toBe("Defect"); + expect(callerRecorder.reports).toHaveLength(0); + }), + ); + + it.scopedLive("reports a child failure to the reporters captured by its parent", () => + Effect.gen(function* () { + const recorder = makeRecorder(); + const failure = new LifecycleDefect({ message: "child transition" }); + const child = Machine.make({ state: S, event: E, initial: S.Idle }).on(S.Idle, E.Crash, () => + Effect.die(failure), + ); + const parent = Machine.make({ state: S, event: E, initial: S.Idle }).background(({ self }) => + self.spawn("child", child).pipe( + Effect.orDie, + Effect.flatMap((ref) => ref.send(E.Crash).pipe(Effect.andThen(ref.awaitExit))), + Effect.andThen(Effect.never), + ), + ); + const exit = yield* Effect.gen(function* () { + const actor = yield* Machine.spawn(parent, { id: "parent" }); + yield* actor.start; + yield* eventually(() => actor.children.has("child")); + const childRef = actor.children.get("child"); + if (childRef === undefined) return yield* Effect.die("child was not spawned"); + const childExit = yield* childRef.awaitExit; + yield* actor.stop; + return childExit; + }).pipe(Effect.provide(Layer.merge(recorder.layer, ActorSystemDefault))); + + expect(exit._tag).toBe("Defect"); + const report = yield* onlyReport(recorder.reports); + expect(defectsOf(report.cause)).toEqual([failure]); + expect(report.annotations).toMatchObject(annotationsFor("child", 0, "transition")); + }), + ); +}); + +// ============================================================================ +// Actor-owned failures +// ============================================================================ + +describe("error reporting: actor-owned failures", () => { + it.scopedLive("reports a final output defect once at terminal completion", () => + Effect.gen(function* () { + const recorder = makeRecorder(); + const failure = new LifecycleDefect({ message: "output" }); + const machine = Machine.make({ state: S, event: E, initial: S.Idle }) + .on(S.Idle, E.Finish, () => S.Done) + .final(S.Done, () => { + // oxlint-disable-next-line effect/noThrowStatement -- final output callbacks are synchronous; this throw is the defect under test. + throw failure; + }); + const exit = yield* Effect.gen(function* () { + const actor = yield* Machine.spawn(machine, { id: "output" }); + yield* actor.start; + yield* actor.send(E.Finish); + return yield* actor.awaitExit; + }).pipe(Effect.provide(recorder.layer)); + + expect(exit._tag).toBe("Defect"); + if (exit._tag !== "Defect") return; + const report = yield* onlyReport(recorder.reports); + expect(report.cause).toBe(exit.cause); + expect(defectsOf(exit.cause)).toEqual([failure]); + expect(report.annotations).toMatchObject(annotationsFor("output", 0, "cleanup")); + }), + ); + + it.scopedLive("reports a restart failure that leaves no generation to close", () => + Effect.gen(function* () { + const recorder = makeRecorder(); + const transitionFailure = new LifecycleDefect({ message: "transition" }); + const recoveryFailure = new LifecycleDefect({ message: "recovery" }); + const machine = Machine.make({ state: S, event: E, initial: S.Idle }).on( + S.Idle, + E.Crash, + () => Effect.die(transitionFailure), + ); + yield* Effect.gen(function* () { + const actor = yield* Machine.spawn(machine, { + id: "recovery", + supervision: Supervision.restart({ maxRestarts: 1 }), + lifecycle: { + recovery: { + resolve: ({ generation }) => { + if (generation === 0) return Effect.succeedNone; + return Effect.die(recoveryFailure); + }, + }, + }, + }); + yield* actor.start; + yield* actor.send(E.Crash); + yield* eventually(() => recorder.reports.length >= 2); + yield* Effect.exit(actor.stop); + }).pipe(Effect.provide(recorder.layer)); + + expect(recorder.reports.map((report) => defectsOf(report.cause))).toEqual([ + [transitionFailure], + [recoveryFailure], + ]); + expect(recorder.reports.at(1)?.annotations).toMatchObject( + annotationsFor("recovery", 1, "restart"), + ); + }), + ); + + it.scopedLive("reports a restart recovery defect that ends while stop cancels the restart", () => + Effect.gen(function* () { + const recorder = makeRecorder(); + const transitionFailure = new LifecycleDefect({ message: "transition" }); + const recoveryFailure = new LifecycleDefect({ message: "restart recovery during stop" }); + const machine = Machine.make({ state: S, event: E, initial: S.Idle }).on( + S.Idle, + E.Crash, + () => Effect.die(transitionFailure), + ); + yield* Effect.gen(function* () { + const actor = yield* Machine.spawn(machine, { + id: "stopped-restart", + supervision: Supervision.restart({ maxRestarts: 1 }), + lifecycle: { + recovery: { + // Stop cannot cancel this step, so its defect settles while stop waits. + resolve: ({ generation }) => { + if (generation === 0) return Effect.succeedNone; + return Effect.sleep("30 millis").pipe( + Effect.andThen(Effect.die(recoveryFailure)), + Effect.uninterruptible, + ); + }, + }, + }, + }); + yield* actor.start; + yield* actor.send(E.Crash); + yield* eventually(() => recorder.reports.length >= 1); + yield* Effect.sleep("5 millis"); + yield* Effect.exit(actor.stop); + }).pipe(Effect.provide(recorder.layer)); + + expect(recorder.reports.map((report) => defectsOf(report.cause))).toEqual([ + [transitionFailure], + [recoveryFailure], + ]); + expect(recorder.reports.at(1)?.annotations).toMatchObject( + annotationsFor("stopped-restart", 1, "restart"), + ); + }), + ); + + it.scopedLive("reports a restart schedule defect apart from schedule exhaustion", () => + Effect.gen(function* () { + const recorder = makeRecorder(); + const transitionFailure = new LifecycleDefect({ message: "transition" }); + const scheduleFailure = new LifecycleDefect({ message: "schedule" }); + const machine = Machine.make({ state: S, event: E, initial: S.Idle }).on( + S.Idle, + E.Crash, + () => Effect.die(transitionFailure), + ); + const dyingSchedule = Schedule.fromStep( + Effect.succeed((_now: number, _input: unknown) => Effect.die(scheduleFailure)), + ); + const exit = yield* Effect.gen(function* () { + const actor = yield* Machine.spawn(machine, { + id: "schedule", + supervision: { schedule: dyingSchedule }, + }); + yield* actor.start; + yield* actor.send(E.Crash); + return yield* actor.awaitExit.pipe(Effect.timeout("1 second")); + }).pipe(Effect.provide(recorder.layer)); + + expect(exit._tag).toBe("Defect"); + expect(recorder.reports.map((report) => defectsOf(report.cause))).toEqual([ + [transitionFailure], + [scheduleFailure], + ]); + expect(recorder.reports.at(1)?.annotations).toMatchObject( + annotationsFor("schedule", 0, "restart"), + ); + }), + ); + + it.scopedLive("reports a restart schedule defect that settles with exhaustion", () => + Effect.gen(function* () { + const recorder = makeRecorder(); + const transitionFailure = new LifecycleDefect({ message: "transition" }); + const scheduleFailure = new LifecycleDefect({ message: "schedule finalizer" }); + const machine = Machine.make({ state: S, event: E, initial: S.Idle }).on( + S.Idle, + E.Crash, + () => Effect.die(transitionFailure), + ); + const exhaustingSchedule = Schedule.fromStep( + Effect.succeed((_now: number, _input: unknown) => + Cause.done("exhausted").pipe(Effect.ensuring(Effect.die(scheduleFailure))), + ), + ); + const exit = yield* Effect.gen(function* () { + const actor = yield* Machine.spawn(machine, { + id: "exhausted-schedule-defect", + supervision: { schedule: exhaustingSchedule }, + }); + yield* actor.start; + yield* actor.send(E.Crash); + return yield* actor.awaitExit.pipe(Effect.timeout("1 second")); + }).pipe(Effect.provide(recorder.layer)); + + expect(exit._tag).toBe("Defect"); + expect(recorder.reports.map((report) => defectsOf(report.cause))).toEqual([ + [transitionFailure], + [scheduleFailure], + ]); + expect(recorder.reports.at(1)?.annotations).toMatchObject( + annotationsFor("exhausted-schedule-defect", 0, "restart"), + ); + }), + ); + + it.scopedLive("reports nothing more when the restart schedule is exhausted", () => + Effect.gen(function* () { + const recorder = makeRecorder(); + const failure = new LifecycleDefect({ message: "transition" }); + const machine = Machine.make({ state: S, event: E, initial: S.Idle }).on( + S.Idle, + E.Crash, + () => Effect.die(failure), + ); + const exit = yield* Effect.gen(function* () { + const actor = yield* Machine.spawn(machine, { + id: "exhausted", + supervision: Supervision.restart({ maxRestarts: 0 }), + }); + yield* actor.start; + yield* actor.send(E.Crash); + return yield* actor.awaitExit.pipe(Effect.timeout("1 second")); + }).pipe(Effect.provide(recorder.layer)); + + expect(exit._tag).toBe("Defect"); + expect(recorder.reports.map((report) => defectsOf(report.cause))).toEqual([[failure]]); + }), + ); + + it.scopedLive("settles the actor when a reporter throws", () => + Effect.gen(function* () { + const throwing: ErrorReporter.ErrorReporter = { + [ErrorReporter.TypeId]: ErrorReporter.TypeId, + report: () => { + // oxlint-disable-next-line effect/noThrowStatement -- a host reporter can throw; this throw is the fault under test. + throw new LifecycleDefect({ message: "reporter bug" }); + }, + }; + const machine = Machine.make({ state: S, event: E, initial: S.Idle }).on( + S.Idle, + E.Crash, + () => Effect.die(new LifecycleDefect({ message: "transition" })), + ); + const exit = yield* Effect.gen(function* () { + const actor = yield* Machine.spawn(machine); + yield* actor.start; + yield* actor.send(E.Crash); + return yield* actor.awaitExit.pipe(Effect.timeout("1 second")); + }).pipe(Effect.provide(ErrorReporter.layer([throwing]))); + + expect(exit._tag).toBe("Defect"); + }), + ); +}); + +// ============================================================================ +// Cluster entities +// ============================================================================ + +describe("error reporting: cluster entities", () => { + it.scopedLive("reports an entity generation defect to the reporters of its allocation", () => + Effect.gen(function* () { + const recorder = makeRecorder(); + const failure = new LifecycleDefect({ message: "entity background" }); + const machine = Machine.make({ state: S, event: E, initial: S.Idle }).background(() => + Effect.yieldNow.pipe(Effect.andThen(Effect.die(failure))), + ); + const entity = toEntity(machine, { type: "ReportingEntity" }); + const makeClient = yield* Entity.makeTestClient( + entity, + EntityMachine.layer(entity, machine, {}).pipe( + Layer.provide(Layer.merge(ActorSystemDefault, recorder.layer)), + ), + ); + yield* makeClient("entity-1"); + yield* eventually(() => recorder.reports.length > 0); + + const report = yield* onlyReport(recorder.reports); + expect(defectsOf(report.cause)).toEqual([failure]); + expect(report.annotations).toMatchObject(annotationsFor("entity-1", 0, "background")); + }).pipe( + Effect.provide( + ShardingConfig.layer({ + shardsPerGroup: 300, + entityMailboxCapacity: 10, + entityTerminationTimeout: 0, + entityMessagePollInterval: 5000, + sendRetryInterval: 100, + }), + ), + ), + ); +});