From e15802976ddb59472a29e66c06d268767bf7e169 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Julien=20Valverd=C3=A9?= Date: Sun, 21 Jun 2026 23:32:29 +0200 Subject: [PATCH] Add additional error handling primitives for Subscribable --- packages/effect-lens-next/README.md | 2 +- .../effect-lens-next/src/Subscribable.test.ts | 35 +++++++ packages/effect-lens-next/src/Subscribable.ts | 78 ++++++++++++++- packages/effect-lens/README.md | 2 +- packages/effect-lens/src/Subscribable.test.ts | 35 +++++++ packages/effect-lens/src/Subscribable.ts | 99 ++++++++++++++++++- 6 files changed, 243 insertions(+), 8 deletions(-) diff --git a/packages/effect-lens-next/README.md b/packages/effect-lens-next/README.md index 0312827..e236800 100644 --- a/packages/effect-lens-next/README.md +++ b/packages/effect-lens-next/README.md @@ -134,4 +134,4 @@ const count = Subscribable.focusArrayLength(users) const current = yield* count.get ``` -The module includes `make`, `map`, `mapEffect`, `mapError`, `tapError`, `unwrap`, and the same readonly focus helpers used by Lens. A `SubscriptionRef` can first be converted with `Lens.fromSubscriptionRef`, since every Lens is already a Subscribable. +The module includes value transforms, error recovery (`catch`, `catchCause`, `orElse`, `orElseSucceed`, and `retry`), error inspection, `unwrap`, and the same readonly focus helpers used by Lens. A `SubscriptionRef` can first be converted with `Lens.fromSubscriptionRef`, since every Lens is already a Subscribable. diff --git a/packages/effect-lens-next/src/Subscribable.test.ts b/packages/effect-lens-next/src/Subscribable.test.ts index 94fab97..a79d2b1 100644 --- a/packages/effect-lens-next/src/Subscribable.test.ts +++ b/packages/effect-lens-next/src/Subscribable.test.ts @@ -37,6 +37,41 @@ describe("Subscribable", () => { expect(observed).toEqual(["get", "changes"]) }) + test("catch recovers get and changes with the corresponding fallback channel", async () => { + const source = Subscribable.make({ + get: Effect.fail("get"), + changes: Stream.fail("changes"), + }) + const recovered = source.pipe(Subscribable.catch(() => Subscribable.make({ + get: Effect.succeed("fallback-get"), + changes: Stream.succeed("fallback-changes"), + }))) + + const result = await Effect.runPromise(Effect.gen(function*() { + const current = yield* recovered.get + const changes = yield* Stream.runCollect(recovered.changes) + return [current, Array.from(changes)] + })) + + expect(result).toEqual(["fallback-get", ["fallback-changes"]]) + }) + + test("orElseSucceed recovers errors from get and changes", async () => { + const source = Subscribable.make({ + get: Effect.fail("get"), + changes: Stream.fail("changes"), + }) + const recovered = Subscribable.orElseSucceed(source, () => "fallback") + + const result = await Effect.runPromise(Effect.gen(function*() { + const current = yield* recovered.get + const changes = yield* Stream.runCollect(recovered.changes) + return [current, Array.from(changes)] + })) + + expect(result).toEqual(["fallback", ["fallback"]]) + }) + test("focusArrayLength reads the current array length and reflects updates", async () => { const result = await Effect.runPromise( Effect.flatMap( diff --git a/packages/effect-lens-next/src/Subscribable.ts b/packages/effect-lens-next/src/Subscribable.ts index 844d633..bb9b53f 100644 --- a/packages/effect-lens-next/src/Subscribable.ts +++ b/packages/effect-lens-next/src/Subscribable.ts @@ -1,4 +1,4 @@ -import { Array, type Cause, Chunk, Effect, Function, Iterable, Option, Pipeable, Predicate, Stream } from "effect" +import { Array, type Cause, Chunk, Effect, Function, Iterable, Option, Pipeable, Predicate, type Result, type Schedule, Stream } from "effect" export const SubscribableTypeId: unique symbol = Symbol.for("@effect-fc/Lens/v4/Subscribable") @@ -124,6 +124,82 @@ export const tapError: { get changes() { return Stream.tapError(self.changes, f) }, })) +const catch_: { + (f: (error: NoInfer) => Subscribable): (self: Subscribable) => Subscribable + (self: Subscribable, f: (error: NoInfer) => Subscribable): Subscribable +} = Function.dual(2, ( + self: Subscribable, + f: (error: NoInfer) => Subscribable, +) => make({ + get get() { return Effect.catch(self.get, error => f(error).get) }, + get changes() { return Stream.catch(self.changes, error => f(error).changes) }, +})) + +/** Recovers from typed errors with another `Subscribable`. */ +export { catch_ as catch } + +/** Recovers from all failure causes with another `Subscribable`. */ +export const catchCause: { + (f: (cause: Cause.Cause>) => Subscribable): (self: Subscribable) => Subscribable + (self: Subscribable, f: (cause: Cause.Cause>) => Subscribable): Subscribable +} = Function.dual(2, ( + self: Subscribable, + f: (cause: Cause.Cause>) => Subscribable, +) => make({ + get get() { return Effect.catchCause(self.get, cause => f(cause).get) }, + get changes() { return Stream.catchCause(self.changes, cause => f(cause).changes) }, +})) + +/** Runs an effect when either channel fails, exposing the complete failure cause. */ +export const tapCause: { + (f: (cause: Cause.Cause>) => Effect.Effect): (self: Subscribable) => Subscribable + (self: Subscribable, f: (cause: Cause.Cause>) => Effect.Effect): Subscribable +} = Function.dual(2, ( + self: Subscribable, + f: (cause: Cause.Cause>) => Effect.Effect, +) => make({ + get get() { return Effect.tapCause(self.get, f) }, + get changes() { return Stream.tapCause(self.changes, f) }, +})) + +/** Falls back to another `Subscribable` when either channel fails. */ +export const orElse: { + (that: () => Subscribable): (self: Subscribable) => Subscribable + (self: Subscribable, that: () => Subscribable): Subscribable +} = Function.dual(2, ( + self: Subscribable, + that: () => Subscribable, +) => catch_(self, that)) + +/** Replaces typed errors from either channel with a lazily evaluated value. */ +export const orElseSucceed: { + (value: () => B): (self: Subscribable) => Subscribable + (self: Subscribable, value: () => B): Subscribable +} = Function.dual(2, (self: Subscribable, value: () => B) => make({ + get get() { return Effect.orElseSucceed(self.get, value) }, + get changes() { return Stream.orElseSucceed(self.changes, value) }, +})) + +/** Retries failures from both channels according to the supplied schedule. */ +export const retry: { + (policy: Schedule.Schedule, E2, R2>): (self: Subscribable) => Subscribable + (self: Subscribable, policy: Schedule.Schedule, E2, R2>): Subscribable +} = Function.dual(2, ( + self: Subscribable, + policy: Schedule.Schedule, E2, R2>, +) => make({ + get get() { return Effect.retry(self.get, policy) }, + get changes() { return Stream.retry(self.changes, policy) }, +})) + +/** Converts typed failures from both channels into `Result` values. */ +export const result = ( + self: Subscribable, +): Subscribable, never, R> => make({ + get get() { return Effect.result(self.get) }, + get changes() { return Stream.result(self.changes) }, +}) + /** Narrows the focus to a field of an object. */ export const focusObjectOn: { diff --git a/packages/effect-lens/README.md b/packages/effect-lens/README.md index cbe6b98..82ba2e5 100644 --- a/packages/effect-lens/README.md +++ b/packages/effect-lens/README.md @@ -263,7 +263,7 @@ yield* someFunctionThatShouldOnlyHaveReadonlyAccessToTheState(lens) ``` #### Focusing -This library re-exports Effect's `Subscribable` module and adds `mapError`, `tapError`, and a few transforms to narrow the focus of `Subscribable`'s, same as Lenses: +This library re-exports Effect's `Subscribable` module and adds error transforms and recovery (`catchAll`, `catchAllCause`, `orElse`, `orElseSucceed`, and `retry`), plus transforms to narrow the focus of `Subscribable`'s, same as Lenses: ```typescript import { Subscribable } from "effect-lens" diff --git a/packages/effect-lens/src/Subscribable.test.ts b/packages/effect-lens/src/Subscribable.test.ts index 1bd36d5..447f14f 100644 --- a/packages/effect-lens/src/Subscribable.test.ts +++ b/packages/effect-lens/src/Subscribable.test.ts @@ -36,6 +36,41 @@ describe("Subscribable", () => { expect(observed).toEqual(["get", "changes"]) }) + test("catchAll recovers get and changes with the corresponding fallback channel", async () => { + const source = Subscribable.make({ + get: Effect.fail("get"), + changes: Stream.fail("changes"), + }) + const recovered = source.pipe(Subscribable.catchAll(() => Subscribable.make({ + get: Effect.succeed("fallback-get"), + changes: Stream.succeed("fallback-changes"), + }))) + + const result = await Effect.runPromise(Effect.gen(function*() { + const current = yield* recovered.get + const changes = yield* Stream.runCollect(recovered.changes) + return [current, Array.from(changes)] + })) + + expect(result).toEqual(["fallback-get", ["fallback-changes"]]) + }) + + test("orElseSucceed recovers errors from get and changes", async () => { + const source = Subscribable.make({ + get: Effect.fail("get"), + changes: Stream.fail("changes"), + }) + const recovered = Subscribable.orElseSucceed(source, () => "fallback") + + const result = await Effect.runPromise(Effect.gen(function*() { + const current = yield* recovered.get + const changes = yield* Stream.runCollect(recovered.changes) + return [current, Array.from(changes)] + })) + + expect(result).toEqual(["fallback", ["fallback"]]) + }) + test("focusArrayLength reads the current array length and reflects updates", async () => { const result = await Effect.runPromise( Effect.flatMap( diff --git a/packages/effect-lens/src/Subscribable.ts b/packages/effect-lens/src/Subscribable.ts index 9aea4e5..c2ac183 100644 --- a/packages/effect-lens/src/Subscribable.ts +++ b/packages/effect-lens/src/Subscribable.ts @@ -1,4 +1,4 @@ -import { Array, Chunk, Effect, Function, Iterable, Option, Stream, Subscribable } from "effect" +import { Array, type Cause, Chunk, Effect, type Either, Function, Iterable, Option, type Schedule, Stream, Subscribable } from "effect" import type { NoSuchElementException } from "effect/Cause" @@ -56,8 +56,20 @@ export const mapError: { self: Subscribable.Subscribable, f: (error: NoInfer) => E2, ): Subscribable.Subscribable => Subscribable.make({ - get: Effect.mapError(self.get, f), - changes: Stream.mapError(self.changes, f), + get get() { return Effect.mapError(self.get, f) }, + get changes() { return Stream.mapError(self.changes, f) }, +})) + +/** Maps complete failure causes from both channels. */ +export const mapErrorCause: { + (f: (cause: Cause.Cause>) => Cause.Cause): (self: Subscribable.Subscribable) => Subscribable.Subscribable + (self: Subscribable.Subscribable, f: (cause: Cause.Cause>) => Cause.Cause): Subscribable.Subscribable +} = Function.dual(2, ( + self: Subscribable.Subscribable, + f: (cause: Cause.Cause>) => Cause.Cause, +): Subscribable.Subscribable => Subscribable.make({ + get get() { return Effect.mapErrorCause(self.get, f) }, + get changes() { return Stream.mapErrorCause(self.changes, f) }, })) /** @@ -75,10 +87,87 @@ export const tapError: { self: Subscribable.Subscribable, f: (error: NoInfer) => Effect.Effect, ): Subscribable.Subscribable => Subscribable.make({ - get: Effect.tapError(self.get, f), - changes: Stream.tapError(self.changes, f), + get get() { return Effect.tapError(self.get, f) }, + get changes() { return Stream.tapError(self.changes, f) }, })) +/** Recovers from typed errors with another `Subscribable`. */ +export const catchAll: { + (f: (error: NoInfer) => Subscribable.Subscribable): (self: Subscribable.Subscribable) => Subscribable.Subscribable + (self: Subscribable.Subscribable, f: (error: NoInfer) => Subscribable.Subscribable): Subscribable.Subscribable +} = Function.dual(2, ( + self: Subscribable.Subscribable, + f: (error: NoInfer) => Subscribable.Subscribable, +): Subscribable.Subscribable => Subscribable.make({ + get get() { return Effect.catchAll(self.get, error => f(error).get) }, + get changes() { return Stream.catchAll(self.changes, error => f(error).changes) }, +})) + +/** Recovers from all failure causes with another `Subscribable`. */ +export const catchAllCause: { + (f: (cause: Cause.Cause>) => Subscribable.Subscribable): (self: Subscribable.Subscribable) => Subscribable.Subscribable + (self: Subscribable.Subscribable, f: (cause: Cause.Cause>) => Subscribable.Subscribable): Subscribable.Subscribable +} = Function.dual(2, ( + self: Subscribable.Subscribable, + f: (cause: Cause.Cause>) => Subscribable.Subscribable, +): Subscribable.Subscribable => Subscribable.make({ + get get() { return Effect.catchAllCause(self.get, cause => f(cause).get) }, + get changes() { return Stream.catchAllCause(self.changes, cause => f(cause).changes) }, +})) + +/** Runs an effect when either channel fails, exposing the complete failure cause. */ +export const tapErrorCause: { + (f: (cause: Cause.Cause>) => Effect.Effect): (self: Subscribable.Subscribable) => Subscribable.Subscribable + (self: Subscribable.Subscribable, f: (cause: Cause.Cause>) => Effect.Effect): Subscribable.Subscribable +} = Function.dual(2, ( + self: Subscribable.Subscribable, + f: (cause: Cause.Cause>) => Effect.Effect, +): Subscribable.Subscribable => Subscribable.make({ + get get() { return Effect.tapErrorCause(self.get, f) }, + get changes() { return Stream.tapErrorCause(self.changes, f) }, +})) + +/** Falls back to another `Subscribable` when either channel fails. */ +export const orElse: { + (that: () => Subscribable.Subscribable): (self: Subscribable.Subscribable) => Subscribable.Subscribable + (self: Subscribable.Subscribable, that: () => Subscribable.Subscribable): Subscribable.Subscribable +} = Function.dual(2, ( + self: Subscribable.Subscribable, + that: () => Subscribable.Subscribable, +): Subscribable.Subscribable => catchAll(self, that)) + +/** Replaces typed errors from either channel with a lazily evaluated value. */ +export const orElseSucceed: { + (value: () => B): (self: Subscribable.Subscribable) => Subscribable.Subscribable + (self: Subscribable.Subscribable, value: () => B): Subscribable.Subscribable +} = Function.dual(2, ( + self: Subscribable.Subscribable, + value: () => B, +): Subscribable.Subscribable => Subscribable.make({ + get get() { return Effect.orElseSucceed(self.get, value) }, + get changes() { return Stream.orElseSucceed(self.changes, value) }, +})) + +/** Retries failures from both channels according to the supplied schedule. */ +export const retry: { + (policy: Schedule.Schedule, R2>): (self: Subscribable.Subscribable) => Subscribable.Subscribable + (self: Subscribable.Subscribable, policy: Schedule.Schedule, R2>): Subscribable.Subscribable +} = Function.dual(2, ( + self: Subscribable.Subscribable, + policy: Schedule.Schedule, R2>, +): Subscribable.Subscribable => Subscribable.make({ + get get() { return Effect.retry(self.get, policy) }, + get changes() { return Stream.retry(self.changes, policy) }, +})) + +/** Converts typed failures from both channels into `Either` values. */ +export const either = ( + self: Subscribable.Subscribable, +): Subscribable.Subscribable, never, R> => Subscribable.make({ + get get() { return Effect.either(self.get) }, + get changes() { return Stream.either(self.changes) }, +}) + /** * Narrows the focus to a field of an object.