Add additional error handling primitives for Subscribable
Lint / lint (push) Successful in 16s

This commit is contained in:
Julien Valverdé
2026-06-21 23:32:29 +02:00
parent 778a041a13
commit e15802976d
6 changed files with 243 additions and 8 deletions
+1 -1
View File
@@ -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.
@@ -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(
+77 -1
View File
@@ -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_: {
<E, B, E2, R2>(f: (error: NoInfer<E>) => Subscribable<B, E2, R2>): <A, R>(self: Subscribable<A, E, R>) => Subscribable<A | B, E2, R | R2>
<A, E, R, B, E2, R2>(self: Subscribable<A, E, R>, f: (error: NoInfer<E>) => Subscribable<B, E2, R2>): Subscribable<A | B, E2, R | R2>
} = Function.dual(2, <A, E, R, B, E2, R2>(
self: Subscribable<A, E, R>,
f: (error: NoInfer<E>) => Subscribable<B, E2, R2>,
) => 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: {
<E, B, E2, R2>(f: (cause: Cause.Cause<NoInfer<E>>) => Subscribable<B, E2, R2>): <A, R>(self: Subscribable<A, E, R>) => Subscribable<A | B, E2, R | R2>
<A, E, R, B, E2, R2>(self: Subscribable<A, E, R>, f: (cause: Cause.Cause<NoInfer<E>>) => Subscribable<B, E2, R2>): Subscribable<A | B, E2, R | R2>
} = Function.dual(2, <A, E, R, B, E2, R2>(
self: Subscribable<A, E, R>,
f: (cause: Cause.Cause<NoInfer<E>>) => Subscribable<B, E2, R2>,
) => 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: {
<E, B, E2, R2>(f: (cause: Cause.Cause<NoInfer<E>>) => Effect.Effect<B, E2, R2>): <A, R>(self: Subscribable<A, E, R>) => Subscribable<A, E | E2, R | R2>
<A, E, R, B, E2, R2>(self: Subscribable<A, E, R>, f: (cause: Cause.Cause<NoInfer<E>>) => Effect.Effect<B, E2, R2>): Subscribable<A, E | E2, R | R2>
} = Function.dual(2, <A, E, R, B, E2, R2>(
self: Subscribable<A, E, R>,
f: (cause: Cause.Cause<NoInfer<E>>) => Effect.Effect<B, E2, R2>,
) => 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: {
<B, E2, R2>(that: () => Subscribable<B, E2, R2>): <A, E, R>(self: Subscribable<A, E, R>) => Subscribable<A | B, E2, R | R2>
<A, E, R, B, E2, R2>(self: Subscribable<A, E, R>, that: () => Subscribable<B, E2, R2>): Subscribable<A | B, E2, R | R2>
} = Function.dual(2, <A, E, R, B, E2, R2>(
self: Subscribable<A, E, R>,
that: () => Subscribable<B, E2, R2>,
) => catch_(self, that))
/** Replaces typed errors from either channel with a lazily evaluated value. */
export const orElseSucceed: {
<B>(value: () => B): <A, E, R>(self: Subscribable<A, E, R>) => Subscribable<A | B, never, R>
<A, E, R, B>(self: Subscribable<A, E, R>, value: () => B): Subscribable<A | B, never, R>
} = Function.dual(2, <A, E, R, B>(self: Subscribable<A, E, R>, 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: {
<E, X, E2, R2>(policy: Schedule.Schedule<X, NoInfer<E>, E2, R2>): <A, R>(self: Subscribable<A, E, R>) => Subscribable<A, E | E2, R | R2>
<A, E, R, X, E2, R2>(self: Subscribable<A, E, R>, policy: Schedule.Schedule<X, NoInfer<E>, E2, R2>): Subscribable<A, E | E2, R | R2>
} = Function.dual(2, <A, E, R, X, E2, R2>(
self: Subscribable<A, E, R>,
policy: Schedule.Schedule<X, NoInfer<E>, 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 = <A, E, R>(
self: Subscribable<A, E, R>,
): Subscribable<Result.Result<A, E>, 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: {
+1 -1
View File
@@ -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"
@@ -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(
+94 -5
View File
@@ -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<A, E, R>,
f: (error: NoInfer<E>) => E2,
): Subscribable.Subscribable<A, E2, R> => 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: {
<E, E2>(f: (cause: Cause.Cause<NoInfer<E>>) => Cause.Cause<E2>): <A, R>(self: Subscribable.Subscribable<A, E, R>) => Subscribable.Subscribable<A, E2, R>
<A, E, R, E2>(self: Subscribable.Subscribable<A, E, R>, f: (cause: Cause.Cause<NoInfer<E>>) => Cause.Cause<E2>): Subscribable.Subscribable<A, E2, R>
} = Function.dual(2, <A, E, R, E2>(
self: Subscribable.Subscribable<A, E, R>,
f: (cause: Cause.Cause<NoInfer<E>>) => Cause.Cause<E2>,
): Subscribable.Subscribable<A, E2, R> => 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<A, E, R>,
f: (error: NoInfer<E>) => Effect.Effect<B, E2, R2>,
): Subscribable.Subscribable<A, E | E2, R | R2> => 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: {
<E, B, E2, R2>(f: (error: NoInfer<E>) => Subscribable.Subscribable<B, E2, R2>): <A, R>(self: Subscribable.Subscribable<A, E, R>) => Subscribable.Subscribable<A | B, E2, R | R2>
<A, E, R, B, E2, R2>(self: Subscribable.Subscribable<A, E, R>, f: (error: NoInfer<E>) => Subscribable.Subscribable<B, E2, R2>): Subscribable.Subscribable<A | B, E2, R | R2>
} = Function.dual(2, <A, E, R, B, E2, R2>(
self: Subscribable.Subscribable<A, E, R>,
f: (error: NoInfer<E>) => Subscribable.Subscribable<B, E2, R2>,
): Subscribable.Subscribable<A | B, E2, R | R2> => 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: {
<E, B, E2, R2>(f: (cause: Cause.Cause<NoInfer<E>>) => Subscribable.Subscribable<B, E2, R2>): <A, R>(self: Subscribable.Subscribable<A, E, R>) => Subscribable.Subscribable<A | B, E2, R | R2>
<A, E, R, B, E2, R2>(self: Subscribable.Subscribable<A, E, R>, f: (cause: Cause.Cause<NoInfer<E>>) => Subscribable.Subscribable<B, E2, R2>): Subscribable.Subscribable<A | B, E2, R | R2>
} = Function.dual(2, <A, E, R, B, E2, R2>(
self: Subscribable.Subscribable<A, E, R>,
f: (cause: Cause.Cause<NoInfer<E>>) => Subscribable.Subscribable<B, E2, R2>,
): Subscribable.Subscribable<A | B, E2, R | R2> => 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: {
<E, B, E2, R2>(f: (cause: Cause.Cause<NoInfer<E>>) => Effect.Effect<B, E2, R2>): <A, R>(self: Subscribable.Subscribable<A, E, R>) => Subscribable.Subscribable<A, E | E2, R | R2>
<A, E, R, B, E2, R2>(self: Subscribable.Subscribable<A, E, R>, f: (cause: Cause.Cause<NoInfer<E>>) => Effect.Effect<B, E2, R2>): Subscribable.Subscribable<A, E | E2, R | R2>
} = Function.dual(2, <A, E, R, B, E2, R2>(
self: Subscribable.Subscribable<A, E, R>,
f: (cause: Cause.Cause<NoInfer<E>>) => Effect.Effect<B, E2, R2>,
): Subscribable.Subscribable<A, E | E2, R | R2> => 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: {
<B, E2, R2>(that: () => Subscribable.Subscribable<B, E2, R2>): <A, E, R>(self: Subscribable.Subscribable<A, E, R>) => Subscribable.Subscribable<A | B, E2, R | R2>
<A, E, R, B, E2, R2>(self: Subscribable.Subscribable<A, E, R>, that: () => Subscribable.Subscribable<B, E2, R2>): Subscribable.Subscribable<A | B, E2, R | R2>
} = Function.dual(2, <A, E, R, B, E2, R2>(
self: Subscribable.Subscribable<A, E, R>,
that: () => Subscribable.Subscribable<B, E2, R2>,
): Subscribable.Subscribable<A | B, E2, R | R2> => catchAll(self, that))
/** Replaces typed errors from either channel with a lazily evaluated value. */
export const orElseSucceed: {
<B>(value: () => B): <A, E, R>(self: Subscribable.Subscribable<A, E, R>) => Subscribable.Subscribable<A | B, never, R>
<A, E, R, B>(self: Subscribable.Subscribable<A, E, R>, value: () => B): Subscribable.Subscribable<A | B, never, R>
} = Function.dual(2, <A, E, R, B>(
self: Subscribable.Subscribable<A, E, R>,
value: () => B,
): Subscribable.Subscribable<A | B, never, R> => 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: {
<E, X, R2>(policy: Schedule.Schedule<X, NoInfer<E>, R2>): <A, R>(self: Subscribable.Subscribable<A, E, R>) => Subscribable.Subscribable<A, E, R | R2>
<A, E, R, X, R2>(self: Subscribable.Subscribable<A, E, R>, policy: Schedule.Schedule<X, NoInfer<E>, R2>): Subscribable.Subscribable<A, E, R | R2>
} = Function.dual(2, <A, E, R, X, R2>(
self: Subscribable.Subscribable<A, E, R>,
policy: Schedule.Schedule<X, NoInfer<E>, R2>,
): Subscribable.Subscribable<A, E, R | R2> => 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 = <A, E, R>(
self: Subscribable.Subscribable<A, E, R>,
): Subscribable.Subscribable<Either.Either<A, E>, 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.