From 078dfc7712acc93631164ff8c2ca636ff5faab4e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Julien=20Valverd=C3=A9?= Date: Fri, 17 Jul 2026 00:15:45 +0200 Subject: [PATCH] 0.2.1 + 2.0.0-beta.1 (#7) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: Julien Valverdé Reviewed-on: https://git.valverde.cloud/Thilawyn/effect-lens/pulls/7 --- bun.lock | 16 +- packages/effect-lens-next/README.md | 58 ++-- packages/effect-lens-next/package.json | 12 +- packages/effect-lens-next/src/Lens.test.ts | 16 +- packages/effect-lens-next/src/Lens.ts | 48 ++- packages/effect-lens-next/src/Subscribable.ts | 246 --------------- .../{Subscribable.test.ts => View.test.ts} | 61 ++-- packages/effect-lens-next/src/View.ts | 292 ++++++++++++++++++ packages/effect-lens-next/src/index.ts | 2 +- packages/effect-lens/README.md | 34 +- packages/effect-lens/package.json | 2 +- packages/effect-lens/src/Lens.ts | 12 +- packages/effect-lens/src/Subscribable.test.ts | 23 ++ packages/effect-lens/src/Subscribable.ts | 35 +++ 14 files changed, 490 insertions(+), 367 deletions(-) delete mode 100644 packages/effect-lens-next/src/Subscribable.ts rename packages/effect-lens-next/src/{Subscribable.test.ts => View.test.ts} (67%) create mode 100644 packages/effect-lens-next/src/View.ts diff --git a/bun.lock b/bun.lock index d9d4893..ccd73bd 100644 --- a/bun.lock +++ b/bun.lock @@ -16,7 +16,7 @@ }, "packages/effect-lens": { "name": "effect-lens", - "version": "0.2.0", + "version": "0.2.1", "devDependencies": { "effect": "^3.21.0", }, @@ -26,12 +26,12 @@ }, "packages/effect-lens-next": { "name": "effect-lens-next", - "version": "0.1.0-beta.0", + "version": "2.0.0-beta.1", "devDependencies": { - "effect": "4.0.0-beta.85", + "effect": "4.0.0-beta.98", }, "peerDependencies": { - "effect": "4.0.0-beta.85", + "effect": "4.0.0-beta.98", }, }, "packages/example": { @@ -185,7 +185,7 @@ "pure-rand": ["pure-rand@6.1.0", "", {}, "sha512-bVWawvoZoBYpp6yIoQtQXHZjmz35RSVHnUOTefl8Vcjr8snTPY1wnpSPMWekcFwbxI6gtmT7rSYPFvz71ldiOA=="], - "toml": ["toml@4.1.1", "", {}, "sha512-EBJnVBr3dTXdA89WVFoAIPUqkBjxPMwRqsfuo1r240tKFHXv3zgca4+NJib/h6TyvGF7vOawz0jGuryJCdNHrw=="], + "toml": ["toml@4.3.0", "", {}, "sha512-lVb8X9BsPVuH0M4BKeS91tXAmJvCjQ5UIyAbQFaxkKGyUFK2RPkhwaFSQH8vbpl1d23eu/IBH+dwVMHWaq9A5A=="], "turbo": ["turbo@2.9.18", "", { "optionalDependencies": { "@turbo/darwin-64": "2.9.18", "@turbo/darwin-arm64": "2.9.18", "@turbo/linux-64": "2.9.18", "@turbo/linux-arm64": "2.9.18", "@turbo/windows-64": "2.9.18", "@turbo/windows-arm64": "2.9.18" }, "bin": { "turbo": "bin/turbo" } }, "sha512-bwabv6PupzeavybzEoArBAkwq5fnzwf8OFnRtpHwnviFWuwJPFxtyH+aVp36TmIqK3aYYgtTJ3J0m2ysxxSzQg=="], @@ -203,12 +203,14 @@ "@effect/sql/uuid": ["uuid@11.1.1", "", { "bin": { "uuid": "dist/esm/bin/uuid" } }, "sha512-vIYxrBCC/N/K+Js3qSN88go7kIfNPssr/hHCesKCQNAjmgvYS2oqr69kIufEG+O4+PfezOH4EbIeHCfFov8ZgQ=="], - "effect-lens-next/effect": ["effect@4.0.0-beta.85", "", { "dependencies": { "@standard-schema/spec": "^1.1.0", "fast-check": "^4.8.0", "find-my-way-ts": "^0.1.6", "ini": "^7.0.0", "kubernetes-types": "^1.30.0", "msgpackr": "^2.0.1", "multipasta": "^0.2.7", "toml": "^4.1.1", "uuid": "^14.0.0", "yaml": "^2.9.0" } }, "sha512-Cjv9YQyv4CiIccmRAQIWAoeESCpCpiuHYY8zb5vqiYs3Ac2yE5RQAnKq7z4Ir/3VJYZb8kxh6rS8czDDnpTdkQ=="], + "effect-lens-next/effect": ["effect@4.0.0-beta.98", "", { "dependencies": { "@standard-schema/spec": "^1.1.0", "fast-check": "^4.9.0", "find-my-way-ts": "^0.1.6", "ini": "^7.0.0", "kubernetes-types": "^1.30.0", "msgpackr": "^2.0.4", "multipasta": "^0.2.8", "toml": "^4.1.2", "uuid": "^14.0.1", "yaml": "^2.9.0" } }, "sha512-oz+bsG5h+6RNrw4t5GMfQrk/xBS8ROoqkYsuvRhBr5O7mCOrpvH/hbw+QrDzvKIpX4HJClwm86F94c87W0sJxg=="], - "effect-lens-next/effect/fast-check": ["fast-check@4.8.0", "", { "dependencies": { "pure-rand": "^8.0.0" } }, "sha512-GOJ158CUMnN6cSahsv4+ExARvIDuzzinFjkp0E9WtiBa5zcVeLozVkWaE4IzFcc+Y48Wp1EDlUZsXRyAztQcSg=="], + "effect-lens-next/effect/fast-check": ["fast-check@4.9.0", "", { "dependencies": { "pure-rand": "^8.0.0" } }, "sha512-7ms6T7SybUev/PQITciI0yLM2pOSFy5zpG8Ty7tQofcVaQUvrMXp6CBwqF6fThLCLOrfBtuHAtwq6Yu4XPCllg=="], "effect-lens-next/effect/msgpackr": ["msgpackr@2.0.4", "", { "optionalDependencies": { "msgpackr-extract": "^3.0.4" } }, "sha512-o1C5KRmuRt+apqMr1HuGSqWStZoRBUpEsCsl15uM9VdAF1qHLtvMOU2En747EnTyEl6c4pzPewRMFF31s1CNbA=="], + "effect-lens-next/effect/multipasta": ["multipasta@0.2.8", "", {}, "sha512-ZPWuMKyv0cSO29f7hozp+k6+crZbQijV8ipMvxNxRf2SwtYGTX1ZX89Kd20VV4H9Znonx+EQn+iy1wGQsJ+b+Q=="], + "effect-lens-next/effect/fast-check/pure-rand": ["pure-rand@8.4.0", "", {}, "sha512-IoM8YF/jY0hiugFo/wOWqfmarlE6J0wc6fDK1PhftMk7MGhVZl88sZimmqBBFomLOCSmcCCpsfj7wXASCpvK9A=="], } } diff --git a/packages/effect-lens-next/README.md b/packages/effect-lens-next/README.md index 677bea4..0ce9fc1 100644 --- a/packages/effect-lens-next/README.md +++ b/packages/effect-lens-next/README.md @@ -4,13 +4,13 @@ A Lens type for [Effect](https://effect.website/) to easily manage nested state. ## Install ``` -npm install effect-lens@beta effect@4.0.0-beta.85 -yarn add effect-lens@beta effect@4.0.0-beta.85 -bun add effect-lens@beta effect@4.0.0-beta.85 +npm install effect-lens@beta effect@4.0.0-beta.98 +yarn add effect-lens@beta effect@4.0.0-beta.98 +bun add effect-lens@beta effect@4.0.0-beta.98 ``` ## Peer dependencies -- `effect` 4.0.0-beta.85 +- `effect` 4.0.0-beta.98 ## Quickstart @@ -48,7 +48,7 @@ const lens = Lens.fromSubscriptionRef(ref) // ^ Lens.Lens const value = yield* Lens.get(lens) -yield* Effect.forkScoped(Stream.runForEach(lens.changes, Console.log)) +yield* Effect.forkScoped(Stream.runForEach(Lens.changes(lens), Console.log)) yield* Lens.updateEffect(lens, values => Effect.fromOption(Array.replace(values, 1, 1664)) ) @@ -64,7 +64,7 @@ You can also create Lenses manually using `make` by providing: - `get`: an effect that reads the current value, - `changes`: a stream of value changes, - `commit`: an effectful write primitive, -- `lock`: an effect that produces the lock used to serialize writes. +- `lock`: an effect that produces the lock used to serialize writes and preserve atomicity. You can get pretty creative! Here's an example of a Lens that points to a specific key of the browser `LocalStorage`: ```typescript @@ -101,9 +101,6 @@ const lens = Effect.all([ ) ``` -Note: while Lens supports asynchronous effects for the proxy logic, we would recommend keeping them synchronous to preserve atomicity. - - ### Focusing Lenses can focus on a nested part of the data type they point to. @@ -240,41 +237,44 @@ const nameLens = lens.pipe( ``` -### Subscribable +### View -Effect 4 no longer provides `Subscribable` and `Readable`. This library provides its own `Subscribable` interface, which you can use as a constraint to allow some parts of your app to only read and subscribe to the Lenses you provide them: +A `View` is a read-only, reactive view of a value: it lets you read the current value and observe subsequent changes. Every `Lens` is also a `View`, which you can use as a constraint to allow some parts of your app to only read and subscribe to the Lenses you provide them. It replaces Effect v3's `Subscribable` module: ```typescript const ref = yield* SubscriptionRef.make<{ readonly users: readonly User[] -}>({ users: [...] }) +}>({ users: [] }) -const someFunctionThatShouldOnlyHaveReadonlyAccessToTheState = ( - usersSub: Subscribable.Subscribable -) => Effect.gen(function*() { - // Do whatever - const usersCountSub = Subscribable.map(usersSub, a => a.length) - const users = yield* usersSub.get - yield* Effect.forkScoped(Stream.runForEach(usersSub.changes, ...)) -}) +const logUserCount = (users: View.View) => + Effect.gen(function*() { + const userCount = View.focusArrayLength(users) + yield* Console.log(`There are ${yield* View.get(userCount)} users`) + yield* Stream.runForEach( + View.changes(userCount), + count => Console.log(`There are now ${count} users`), + ) + }) -const lens = ref.pipe( +const usersLens = ref.pipe( Lens.fromSubscriptionRef, Lens.focusObjectOn("users"), ) -yield* someFunctionThatShouldOnlyHaveReadonlyAccessToTheState(lens) +yield* Effect.forkScoped(logUserCount(usersLens)) ``` +Use `Lens.asView` when you want to make that read-only boundary explicit. + #### Focusing -This library provides a `Subscribable` module with value and error transforms, recovery (`catch`, `catchCause`, `orElse`, `orElseSucceed`, and `retry`), and transforms to narrow the focus of `Subscribable`'s, same as Lenses: +Views can be focused in the same way as Lenses. The `View` module also provides value and error transforms, plus recovery (`catch`, `catchCause`, `orElse`, `orElseSucceed`, and `retry`): ```typescript -import { Subscribable } from "effect-lens" +import { View } from "effect-lens" -declare const sub: Subscribable.Subscribable +declare const view: View.View -// \/ Subscribable.Subscribable -const nameSub = sub.pipe( - Subscribable.focusArrayAt(1), - Subscribable.focusObjectOn("name"), +// \/ View.View +const nameView = view.pipe( + View.focusArrayAt(1), + View.focusObjectOn("name"), ) ``` diff --git a/packages/effect-lens-next/package.json b/packages/effect-lens-next/package.json index 85c0a18..894bbcb 100644 --- a/packages/effect-lens-next/package.json +++ b/packages/effect-lens-next/package.json @@ -1,7 +1,7 @@ { "name": "effect-lens-next", "description": "An effectful Lens type to easily manage nested state", - "version": "2.0.0-beta.0", + "version": "2.0.0-beta.1", "type": "module", "files": [ "./README.md", @@ -21,9 +21,9 @@ "types": "./dist/Lens.d.ts", "default": "./dist/Lens.js" }, - "./Subscribable": { - "types": "./dist/Subscribable.d.ts", - "default": "./dist/Subscribable.js" + "./View": { + "types": "./dist/View.d.ts", + "default": "./dist/View.js" } }, "scripts": { @@ -36,9 +36,9 @@ "clean:modules": "rm -rf node_modules" }, "peerDependencies": { - "effect": "4.0.0-beta.85" + "effect": "4.0.0-beta.98" }, "devDependencies": { - "effect": "4.0.0-beta.85" + "effect": "4.0.0-beta.98" } } diff --git a/packages/effect-lens-next/src/Lens.test.ts b/packages/effect-lens-next/src/Lens.test.ts index 8bc85c8..db5602a 100644 --- a/packages/effect-lens-next/src/Lens.test.ts +++ b/packages/effect-lens-next/src/Lens.test.ts @@ -187,7 +187,7 @@ describe("Lens", () => { expect(result[1]).toBe(25) }) - test("Ref and SynchronizedRef adapters read and update their sources", async () => { + test("Ref, SynchronizedRef, and SubscriptionRef adapters read and update their sources", async () => { const result = await Effect.runPromise(Effect.gen(function*() { const ref = yield* Ref.make(1) const refLens = yield* Lens.fromRef(ref) @@ -195,12 +195,20 @@ describe("Lens", () => { const synchronizedRef = yield* SynchronizedRef.make(10) const synchronizedLens = Lens.fromSynchronizedRef(synchronizedRef) - yield* Lens.updateEffect(synchronizedLens, n => Effect.succeed(n + 5)) + yield* Lens.update(synchronizedLens, n => n + 5) - return [yield* Ref.get(ref), yield* SynchronizedRef.get(synchronizedRef)] as const + const subscriptionRef = yield* SubscriptionRef.make(100) + const subscriptionLens = Lens.fromSubscriptionRef(subscriptionRef) + yield* Lens.update(subscriptionLens, n => n + 2) + + return [ + yield* Ref.get(ref), + yield* SynchronizedRef.get(synchronizedRef), + yield* SubscriptionRef.get(subscriptionRef), + ] as const })) - expect(result).toEqual([2, 15]) + expect(result).toEqual([2, 15, 102]) }) test("modifyEffect updates are atomic under concurrency", async () => { diff --git a/packages/effect-lens-next/src/Lens.ts b/packages/effect-lens-next/src/Lens.ts index c7ee35a..ee9ebb4 100644 --- a/packages/effect-lens-next/src/Lens.ts +++ b/packages/effect-lens-next/src/Lens.ts @@ -1,25 +1,8 @@ -import { - Array, - type Cause, - Chunk, - type Context, - Effect, - Function, - identity, - Option, - Pipeable, - Predicate, - PubSub, - Ref, - Semaphore, - Stream, - SubscriptionRef, - SynchronizedRef, -} from "effect" -import * as Subscribable from "./Subscribable.js" +import { Array, type Cause, Chunk, type Context, Effect, Function, identity, Option, Pipeable, Predicate, PubSub, Ref, Semaphore, Stream, SubscriptionRef, SynchronizedRef } from "effect" +import * as View from "./View.js" -export const LensTypeId: unique symbol = Symbol.for("@effect-fc/Lens/v4/Lens") +export const LensTypeId: unique symbol = Symbol.for("@effect-lens/Lens/Lens") export type LensTypeId = typeof LensTypeId /** @@ -29,8 +12,8 @@ export type LensTypeId = typeof LensTypeId * 2. a `changes` stream that emits every subsequent update to `A`, and * 3. a `modify` effect that can transform the current value. */ -export interface Lens -extends Subscribable.Subscribable { +export interface Lens +extends View.View { readonly [LensTypeId]: LensTypeId readonly modifyEffect: ( @@ -52,7 +35,7 @@ export const LensImplTypeId: unique symbol = Symbol.for("@effect-fc/Lens/v4/Lens export type LensImplTypeId = typeof LensImplTypeId export declare namespace LensImpl { - export interface Resolved { + export interface Resolved { readonly value: A readonly commit: ( next: Effect.Effect @@ -64,9 +47,9 @@ export declare namespace LensImpl { } } -export abstract class LensImpl +export abstract class LensImpl extends Pipeable.Class implements Lens { - readonly [Subscribable.SubscribableTypeId]: Subscribable.SubscribableTypeId = Subscribable.SubscribableTypeId + readonly [View.ViewTypeId]: View.ViewTypeId = View.ViewTypeId readonly [LensTypeId]: LensTypeId = LensTypeId readonly [LensImplTypeId]: LensImplTypeId = LensImplTypeId @@ -120,13 +103,13 @@ export const asLensImpl = ( return lens as LensImpl } -export const asSubscribable = ( +export const asView = ( lens: Lens -): Subscribable.Subscribable => lens +): View.View => lens export declare namespace LensLazyImpl { - export interface Source { + export interface Source { readonly get: Effect.Effect readonly changes: Stream.Stream readonly commit: (a: A) => Effect.Effect @@ -134,7 +117,7 @@ export declare namespace LensLazyImpl { } } -export class LensLazyImpl +export class LensLazyImpl extends LensImpl { constructor( readonly source: LensLazyImpl.Source, @@ -163,7 +146,7 @@ export const make = ( ): Lens => new LensLazyImpl(source) -export class UnwrappedLensImpl +export class UnwrappedLensImpl extends LensImpl { constructor( readonly effect: Effect.Effect, E1, R1> @@ -912,6 +895,11 @@ export const focusOption: { */ export const get = (self: Lens): Effect.Effect => self.get +/** + * Returns the stream of changes from a `Lens`. + */ +export const changes = (self: Lens): Stream.Stream => self.changes + /** * Atomically modifies the value of a `Lens` and returns a computed result. */ diff --git a/packages/effect-lens-next/src/Subscribable.ts b/packages/effect-lens-next/src/Subscribable.ts deleted file mode 100644 index bb9b53f..0000000 --- a/packages/effect-lens-next/src/Subscribable.ts +++ /dev/null @@ -1,246 +0,0 @@ -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") -export type SubscribableTypeId = typeof SubscribableTypeId - -export interface Subscribable extends Pipeable.Pipeable { - readonly [SubscribableTypeId]: SubscribableTypeId - readonly get: Effect.Effect - readonly changes: Stream.Stream -} - -export const isSubscribable = (u: unknown): u is Subscribable => Predicate.hasProperty(u, SubscribableTypeId) - - -export const SubscribableImplTypeId: unique symbol = Symbol.for("@effect-fc/Lens/v4/SubscribableImpl") -export type SubscribableImplTypeId = typeof SubscribableImplTypeId - -export declare namespace SubscribableImpl { - export interface Source { - readonly get: Effect.Effect - readonly changes: Stream.Stream - } -} - -export class SubscribableImpl -extends Pipeable.Class implements Subscribable { - readonly [SubscribableTypeId]: SubscribableTypeId = SubscribableTypeId - readonly [SubscribableImplTypeId]: SubscribableImplTypeId = SubscribableImplTypeId - - constructor( - readonly source: SubscribableImpl.Source, - ) { - super() - } - - get get() { return this.source.get } - get changes() { return this.source.changes } -} - -export const isSubscribableImpl = (u: unknown): u is SubscribableImpl => Predicate.hasProperty(u, SubscribableImplTypeId) - -export const asSubscribableImpl = ( - subscribable: Subscribable -): SubscribableImpl => { - if (!isSubscribableImpl(subscribable)) - throw new Error("Not a 'SubscribableImpl'") - return subscribable as SubscribableImpl -} - -export const make = ( - source: SubscribableImpl.Source -): Subscribable => new SubscribableImpl(source) - -export const unwrap = ( - effect: Effect.Effect, E1, R1>, -): Subscribable => make({ - get: Effect.flatMap(effect, self => self.get), - changes: Stream.unwrap(Effect.map(effect, self => self.changes)), -}) - - -export const map: { - (f: (a: NoInfer) => B): (self: Subscribable) => Subscribable - (self: Subscribable, f: (a: NoInfer) => B): Subscribable -} = Function.dual(2, (self: Subscribable, f: (a: NoInfer) => B) => make({ - get get() { return Effect.map(self.get, f) }, - get changes() { return Stream.map(self.changes, f) }, -})) - -export const mapEffect: { - (f: (a: NoInfer) => Effect.Effect): (self: Subscribable) => Subscribable - (self: Subscribable, f: (a: NoInfer) => Effect.Effect): Subscribable -} = Function.dual(2, ( - self: Subscribable, - f: (a: NoInfer) => Effect.Effect, -) => make({ - get get() { return Effect.flatMap(self.get, f) }, - get changes() { return Stream.mapEffect(self.changes, f) }, -})) - -/** Maps over an `Option` value in the `Subscribable`. */ -export const mapOption: { - (f: (a: A) => B): (self: Subscribable, E, R>) => Subscribable, E, R> - (self: Subscribable, E, R>, f: (a: A) => B): Subscribable, E, R> -} = Function.dual(2, (self: Subscribable, E, R>, f: (a: A) => B) => - map(self, Option.map(f)), -) - -/** Maps over an `Option` value in the `Subscribable` with an Effect. */ -export const mapOptionEffect: { - (f: (a: A) => Effect.Effect): (self: Subscribable, E, R>) => Subscribable, E | E2, R | R2> - (self: Subscribable, E, R>, f: (a: A) => Effect.Effect): Subscribable, E | E2, R | R2> -} = Function.dual(2, ( - self: Subscribable, E, R>, - f: (a: A) => Effect.Effect, -) => mapEffect(self, Option.match({ - onSome: a => Effect.map(f(a), Option.some), - onNone: () => Effect.succeed(Option.none()), -}))) - - -/** Maps errors from both the current value and the stream of changes. */ -export const mapError: { - (f: (error: NoInfer) => E2): (self: Subscribable) => Subscribable - (self: Subscribable, f: (error: NoInfer) => E2): Subscribable -} = Function.dual(2, ( - self: Subscribable, - f: (error: NoInfer) => E2, -) => make({ - get get() { return Effect.mapError(self.get, f) }, - get changes() { return Stream.mapError(self.changes, f) }, -})) - -/** Runs an effect when either the current value or the stream of changes fails. */ -export const tapError: { - (f: (error: NoInfer) => Effect.Effect): (self: Subscribable) => Subscribable - (self: Subscribable, f: (error: NoInfer) => Effect.Effect): Subscribable -} = Function.dual(2, ( - self: Subscribable, - f: (error: NoInfer) => Effect.Effect, -) => make({ - get get() { return Effect.tapError(self.get, f) }, - 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: { - (key: K): (self: Subscribable) => Subscribable - (self: Subscribable, key: K): Subscribable -} = Function.dual(2, (self: Subscribable, key: K) => - map(self, a => a[key]), -) - -/** Narrows the focus to an indexed element of an array. */ -export const focusArrayAt: { - (index: number): (self: Subscribable) => Subscribable - (self: Subscribable, index: number): Subscribable -} = Function.dual(2, (self: Subscribable, index: number) => - mapEffect(self, a => Effect.fromOption(Array.get(a, index))), -) - -export const focusArrayLength = ( - self: Subscribable, -): Subscribable => map(self, Array.length) - -/** Narrows the focus to an indexed element of a readonly tuple. */ -export const focusTupleAt: { - (index: I): (self: Subscribable) => Subscribable - (self: Subscribable, index: I): Subscribable -} = Function.dual(2, (self: Subscribable, index: I) => - map(self, Array.getUnsafe(index)), -) - -/** Narrows the focus to an indexed element of `Chunk`. */ -export const focusChunkAt: { - (index: number): (self: Subscribable, E, R>) => Subscribable - (self: Subscribable, E, R>, index: number): Subscribable -} = Function.dual(2, (self: Subscribable, E, R>, index: number) => - mapEffect(self, chunk => Effect.fromOption(Chunk.get(chunk, index))), -) - -export const focusChunkSize = ( - self: Subscribable, E, R>, -): Subscribable => map(self, Chunk.size) - -export const focusIterableSize = , E, R>( - self: Subscribable, -): Subscribable => map(self, Iterable.size) diff --git a/packages/effect-lens-next/src/Subscribable.test.ts b/packages/effect-lens-next/src/View.test.ts similarity index 67% rename from packages/effect-lens-next/src/Subscribable.test.ts rename to packages/effect-lens-next/src/View.test.ts index a79d2b1..9c94b6e 100644 --- a/packages/effect-lens-next/src/Subscribable.test.ts +++ b/packages/effect-lens-next/src/View.test.ts @@ -1,16 +1,16 @@ import { describe, expect, test } from "bun:test" import { Chunk, Effect, Stream, SubscriptionRef } from "effect" import * as Lens from "./Lens.js" -import * as Subscribable from "./Subscribable.js" +import * as View from "./View.js" -describe("Subscribable", () => { +describe("View", () => { test("mapError transforms errors from get and changes", async () => { - const source = Subscribable.make({ + const source = View.make({ get: Effect.fail("get"), changes: Stream.fail("changes"), }) - const mapped = source.pipe(Subscribable.mapError((error: string) => `mapped:${error}`)) + const mapped = source.pipe(View.mapError((error: string) => `mapped:${error}`)) const result = await Effect.runPromise(Effect.gen(function*() { const getError = yield* Effect.flip(mapped.get) @@ -23,11 +23,11 @@ describe("Subscribable", () => { test("tapError observes errors from get and changes", async () => { const observed: Array = [] - const source = Subscribable.make({ + const source = View.make({ get: Effect.fail("get"), changes: Stream.fail("changes"), }) - const tapped = Subscribable.tapError(source, error => Effect.sync(() => observed.push(error))) + const tapped = View.tapError(source, error => Effect.sync(() => observed.push(error))) await Effect.runPromise(Effect.gen(function*() { yield* Effect.flip(tapped.get) @@ -38,11 +38,11 @@ describe("Subscribable", () => { }) test("catch recovers get and changes with the corresponding fallback channel", async () => { - const source = Subscribable.make({ + const source = View.make({ get: Effect.fail("get"), changes: Stream.fail("changes"), }) - const recovered = source.pipe(Subscribable.catch(() => Subscribable.make({ + const recovered = source.pipe(View.catch(() => View.make({ get: Effect.succeed("fallback-get"), changes: Stream.succeed("fallback-changes"), }))) @@ -57,11 +57,11 @@ describe("Subscribable", () => { }) test("orElseSucceed recovers errors from get and changes", async () => { - const source = Subscribable.make({ + const source = View.make({ get: Effect.fail("get"), changes: Stream.fail("changes"), }) - const recovered = Subscribable.orElseSucceed(source, () => "fallback") + const recovered = View.orElseSucceed(source, () => "fallback") const result = await Effect.runPromise(Effect.gen(function*() { const current = yield* recovered.get @@ -72,17 +72,40 @@ describe("Subscribable", () => { expect(result).toEqual(["fallback", ["fallback"]]) }) + test("zipLatestAll combines current values and change streams", async () => { + const zipped = View.zipLatestAll( + View.make({ + get: Effect.succeed(1), + changes: Stream.succeed(2), + }), + View.make({ + get: Effect.succeed("one"), + changes: Stream.succeed("two"), + }), + ) + + const result = await Effect.runPromise(Effect.all([ + zipped.get, + Stream.runCollect(zipped.changes), + ])) + + expect([result[0], Array.from(result[1])]).toEqual([ + [1, "one"], + [[2, "two"]], + ]) + }) + test("focusArrayLength reads the current array length and reflects updates", async () => { const result = await Effect.runPromise( Effect.flatMap( SubscriptionRef.make([1, 2, 3]), parent => { - const sizeSub = Subscribable.focusArrayLength(Lens.fromSubscriptionRef(parent)) + const sizeView = View.focusArrayLength(Lens.fromSubscriptionRef(parent)) return Effect.flatMap( - sizeSub.get, + sizeView.get, initial => Effect.flatMap( SubscriptionRef.set(parent, [1, 2, 3, 4, 5]), - () => Effect.map(sizeSub.get, next => [initial, next] as const), + () => Effect.map(sizeView.get, next => [initial, next] as const), ), ) }, @@ -97,12 +120,12 @@ describe("Subscribable", () => { Effect.flatMap( SubscriptionRef.make(Chunk.make(1, 2) as Chunk.Chunk), parent => { - const sizeSub = Subscribable.focusChunkSize(Lens.fromSubscriptionRef(parent)) + const sizeView = View.focusChunkSize(Lens.fromSubscriptionRef(parent)) return Effect.flatMap( - sizeSub.get, + sizeView.get, initial => Effect.flatMap( SubscriptionRef.set(parent, Chunk.make(1, 2, 3, 4)), - () => Effect.map(sizeSub.get, next => [initial, next] as const), + () => Effect.map(sizeView.get, next => [initial, next] as const), ), ) }, @@ -117,12 +140,12 @@ describe("Subscribable", () => { Effect.flatMap( SubscriptionRef.make([1, 2, 3]), parent => { - const sizeSub = Subscribable.focusIterableSize(Lens.fromSubscriptionRef(parent)) + const sizeView = View.focusIterableSize(Lens.fromSubscriptionRef(parent)) return Effect.flatMap( - sizeSub.get, + sizeView.get, initial => Effect.flatMap( SubscriptionRef.set(parent, [1, 2, 3, 4, 5]), - () => Effect.map(sizeSub.get, next => [initial, next] as const), + () => Effect.map(sizeView.get, next => [initial, next] as const), ), ) }, diff --git a/packages/effect-lens-next/src/View.ts b/packages/effect-lens-next/src/View.ts new file mode 100644 index 0000000..63987ca --- /dev/null +++ b/packages/effect-lens-next/src/View.ts @@ -0,0 +1,292 @@ +import { Array, type Cause, Chunk, Effect, Function, Iterable, Option, Pipeable, Predicate, type Result, type Schedule, Stream } from "effect" + + +export const ViewTypeId: unique symbol = Symbol.for("@effect-lens/View/View") +export type ViewTypeId = typeof ViewTypeId + +export interface View extends Pipeable.Pipeable { + readonly [ViewTypeId]: ViewTypeId + readonly get: Effect.Effect + readonly changes: Stream.Stream +} + +export const isView = (u: unknown): u is View => Predicate.hasProperty(u, ViewTypeId) + + +export const ViewImplTypeId: unique symbol = Symbol.for("@effect-fc/Lens/v4/ViewImpl") +export type ViewImplTypeId = typeof ViewImplTypeId + +export declare namespace ViewImpl { + export interface Source { + readonly get: Effect.Effect + readonly changes: Stream.Stream + } +} + +export class ViewImpl +extends Pipeable.Class implements View { + readonly [ViewTypeId]: ViewTypeId = ViewTypeId + readonly [ViewImplTypeId]: ViewImplTypeId = ViewImplTypeId + + constructor( + readonly source: ViewImpl.Source, + ) { + super() + } + + get get() { return this.source.get } + get changes() { return this.source.changes } +} + +export const isViewImpl = (u: unknown): u is ViewImpl => Predicate.hasProperty(u, ViewImplTypeId) + +export const asViewImpl = ( + view: View +): ViewImpl => { + if (!isViewImpl(view)) + throw new Error("Not a 'ViewImpl'") + return view as ViewImpl +} + +export const make = ( + source: ViewImpl.Source +): View => new ViewImpl(source) + +export const unwrap = ( + effect: Effect.Effect, E1, R1>, +): View => make({ + get: Effect.flatMap(effect, self => self.get), + changes: Stream.unwrap(Effect.map(effect, self => self.changes)), +}) + + +export const map: { + (f: (a: NoInfer) => B): (self: View) => View + (self: View, f: (a: NoInfer) => B): View +} = Function.dual(2, (self: View, f: (a: NoInfer) => B) => make({ + get get() { return Effect.map(self.get, f) }, + get changes() { return Stream.map(self.changes, f) }, +})) + +export const mapEffect: { + (f: (a: NoInfer) => Effect.Effect): (self: View) => View + (self: View, f: (a: NoInfer) => Effect.Effect): View +} = Function.dual(2, ( + self: View, + f: (a: NoInfer) => Effect.Effect, +) => make({ + get get() { return Effect.flatMap(self.get, f) }, + get changes() { return Stream.mapEffect(self.changes, f) }, +})) + +/** Maps over an `Option` value in the `View`. */ +export const mapOption: { + (f: (a: A) => B): (self: View, E, R>) => View, E, R> + (self: View, E, R>, f: (a: A) => B): View, E, R> +} = Function.dual(2, (self: View, E, R>, f: (a: A) => B) => + map(self, Option.map(f)), +) + +/** Maps over an `Option` value in the `View` with an Effect. */ +export const mapOptionEffect: { + (f: (a: A) => Effect.Effect): (self: View, E, R>) => View, E | E2, R | R2> + (self: View, E, R>, f: (a: A) => Effect.Effect): View, E | E2, R | R2> +} = Function.dual(2, ( + self: View, E, R>, + f: (a: A) => Effect.Effect, +) => mapEffect(self, Option.match({ + onSome: a => Effect.map(f(a), Option.some), + onNone: () => Effect.succeed(Option.none()), +}))) + +/** + * Allows transforming only the `changes` stream of a `View`. + */ +export const mapStream: { + ( + f: (changes: Stream.Stream, NoInfer, NoInfer>) => Stream.Stream, NoInfer, NoInfer>, + ): (self: View) => View + ( + self: View, + f: (changes: Stream.Stream, NoInfer, NoInfer>) => Stream.Stream, NoInfer, NoInfer>, + ): View +} = Function.dual(2, ( + self: View, + f: (changes: Stream.Stream, NoInfer, NoInfer>) => Stream.Stream, NoInfer, NoInfer>, +): View => make({ + get get() { return self.get }, + get changes() { return f(self.changes) }, +})) + +/** Converts typed failures from both channels into `Result` values. */ +export const result = ( + self: View, +): View, never, R> => make({ + get get() { return Effect.result(self.get) }, + get changes() { return Stream.result(self.changes) }, +}) + +/** + * Combines the current values and streams of changes from multiple `View` values. + */ +export const zipLatestAll = []>( + ...elements: T +): View< + [T[number]] extends [never] + ? never + : { [K in keyof T]: T[K] extends View ? A : never }, + [T[number]] extends [never] ? never : T[number] extends View ? E : never, + [T[number]] extends [never] ? never : T[number] extends View ? R : never +> => make({ + get: Effect.all(elements.map(view => view.get)), + changes: Stream.zipLatestAll(...elements.map(view => view.changes)), +}) as any + + +/** Maps errors from both the current value and the stream of changes. */ +export const mapError: { + (f: (error: NoInfer) => E2): (self: View) => View + (self: View, f: (error: NoInfer) => E2): View +} = Function.dual(2, ( + self: View, + f: (error: NoInfer) => E2, +) => make({ + get get() { return Effect.mapError(self.get, f) }, + get changes() { return Stream.mapError(self.changes, f) }, +})) + +/** Runs an effect when either the current value or the stream of changes fails. */ +export const tapError: { + (f: (error: NoInfer) => Effect.Effect): (self: View) => View + (self: View, f: (error: NoInfer) => Effect.Effect): View +} = Function.dual(2, ( + self: View, + f: (error: NoInfer) => Effect.Effect, +) => make({ + get get() { return Effect.tapError(self.get, f) }, + get changes() { return Stream.tapError(self.changes, f) }, +})) + +const catch_: { + (f: (error: NoInfer) => View): (self: View) => View + (self: View, f: (error: NoInfer) => View): View +} = Function.dual(2, ( + self: View, + f: (error: NoInfer) => View, +) => 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 `View`. */ +export { catch_ as catch } + +/** Recovers from all failure causes with another `View`. */ +export const catchCause: { + (f: (cause: Cause.Cause>) => View): (self: View) => View + (self: View, f: (cause: Cause.Cause>) => View): View +} = Function.dual(2, ( + self: View, + f: (cause: Cause.Cause>) => View, +) => 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: View) => View + (self: View, f: (cause: Cause.Cause>) => Effect.Effect): View +} = Function.dual(2, ( + self: View, + 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 `View` when either channel fails. */ +export const orElse: { + (that: () => View): (self: View) => View + (self: View, that: () => View): View +} = Function.dual(2, ( + self: View, + that: () => View, +) => catch_(self, that)) + +/** Replaces typed errors from either channel with a lazily evaluated value. */ +export const orElseSucceed: { + (value: () => B): (self: View) => View + (self: View, value: () => B): View +} = Function.dual(2, (self: View, 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: View) => View + (self: View, policy: Schedule.Schedule, E2, R2>): View +} = Function.dual(2, ( + self: View, + policy: Schedule.Schedule, E2, R2>, +) => make({ + get get() { return Effect.retry(self.get, policy) }, + get changes() { return Stream.retry(self.changes, policy) }, +})) + + +/** Narrows the focus to a field of an object. */ +export const focusObjectOn: { + (key: K): (self: View) => View + (self: View, key: K): View +} = Function.dual(2, (self: View, key: K) => + map(self, a => a[key]), +) + +/** Narrows the focus to an indexed element of an array. */ +export const focusArrayAt: { + (index: number): (self: View) => View + (self: View, index: number): View +} = Function.dual(2, (self: View, index: number) => + mapEffect(self, a => Effect.fromOption(Array.get(a, index))), +) + +export const focusArrayLength = ( + self: View, +): View => map(self, Array.length) + +/** Narrows the focus to an indexed element of a readonly tuple. */ +export const focusTupleAt: { + (index: I): (self: View) => View + (self: View, index: I): View +} = Function.dual(2, (self: View, index: I) => + map(self, Array.getUnsafe(index)), +) + +/** Narrows the focus to an indexed element of `Chunk`. */ +export const focusChunkAt: { + (index: number): (self: View, E, R>) => View + (self: View, E, R>, index: number): View +} = Function.dual(2, (self: View, E, R>, index: number) => + mapEffect(self, chunk => Effect.fromOption(Chunk.get(chunk, index))), +) + +export const focusChunkSize = ( + self: View, E, R>, +): View => map(self, Chunk.size) + +export const focusIterableSize = , E, R>( + self: View, +): View => map(self, Iterable.size) + + +/** + * Reads the current value from a `View`. + */ +export const get = (self: View): Effect.Effect => self.get + +/** + * Returns the stream of changes from a `View`. + */ +export const changes = (self: View): Stream.Stream => self.changes diff --git a/packages/effect-lens-next/src/index.ts b/packages/effect-lens-next/src/index.ts index 24e4bb8..d5df55e 100644 --- a/packages/effect-lens-next/src/index.ts +++ b/packages/effect-lens-next/src/index.ts @@ -1,2 +1,2 @@ export * as Lens from "./Lens.js" -export * as Subscribable from "./Subscribable.js" +export * as View from "./View.js" diff --git a/packages/effect-lens/README.md b/packages/effect-lens/README.md index a34dad7..c6bce5b 100644 --- a/packages/effect-lens/README.md +++ b/packages/effect-lens/README.md @@ -2,7 +2,7 @@ A Lens type for [Effect](https://effect.website/) to easily manage nested state. -This version is for Effect v3. For Effect v4, use version 2.X.X-beta.X. +This version is for Effect v3. For Effect v4, use the [2.0.0 beta](https://www.npmjs.com/package/effect-lens/v/2.0.0-beta.1). ## Install ``` @@ -64,7 +64,7 @@ You can also create Lenses manually using `make` by providing: - `get`: an effect that reads the current value, - `changes`: a stream of value changes, - `commit`: an effectful write primitive, -- `lock`: an effect that produces the lock used to serialize writes. +- `lock`: an effect that produces the lock used to serialize writes and preserve atomicity. You can get pretty creative! Here's an example of a Lens that points to a specific key of the browser `LocalStorage`: ```typescript @@ -101,9 +101,6 @@ const lens = Effect.all([ ) ``` -Note: while Lens supports asynchronous effects for the proxy logic, we would recommend keeping them synchronous to preserve atomicity. - - ### Focusing Lenses can focus on a nested part of the data type they point to. @@ -242,30 +239,31 @@ const nameLens = lens.pipe( ### Subscribable -Lens implements both Effect's `Subscribable` and `Readable`, which you can use as a constraint to allow some parts of your app to only read and subscribe to the Lenses you provide them: +Effect's `Subscribable` is a read-only, reactive view of a value: it lets you read the current value and observe subsequent changes. Every `Lens` implements both `Subscribable` and `Readable`, which you can use as constraints to allow some parts of your app to only read and subscribe to the Lenses you provide them: ```typescript const ref = yield* SubscriptionRef.make<{ readonly users: readonly User[] -}>({ users: [...] }) +}>({ users: [] }) -const someFunctionThatShouldOnlyHaveReadonlyAccessToTheState = ( - usersSub: Subscribable.Subscribable -) => Effect.gen(function*() { - // Do whatever - const usersCountSub = Subscribable.map(usersSub, a => a.length) - const users = yield* usersSub.get - yield* Effect.forkScoped(Stream.runForEach(usersSub.changes, ...)) -}) +const logUserCount = (users: Subscribable.Subscribable) => + Effect.gen(function*() { + const userCount = Subscribable.focusArrayLength(users) + yield* Console.log(`There are ${yield* userCount.get} users`) + yield* Stream.runForEach( + userCount.changes, + count => Console.log(`There are now ${count} users`), + ) + }) -const lens = ref.pipe( +const usersLens = ref.pipe( Lens.fromSubscriptionRef, Lens.focusObjectOn("users"), ) -yield* someFunctionThatShouldOnlyHaveReadonlyAccessToTheState(lens) +yield* Effect.forkScoped(logUserCount(usersLens)) ``` #### Focusing -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: +Subscribables can be focused in the same way as Lenses. This library re-exports Effect's `Subscribable` module and adds value and error transforms, plus recovery (`catchAll`, `catchAllCause`, `orElse`, `orElseSucceed`, and `retry`): ```typescript import { Subscribable } from "effect-lens" diff --git a/packages/effect-lens/package.json b/packages/effect-lens/package.json index 38b4137..c1eebf0 100644 --- a/packages/effect-lens/package.json +++ b/packages/effect-lens/package.json @@ -1,7 +1,7 @@ { "name": "effect-lens", "description": "An effectful Lens type to easily manage nested state", - "version": "0.2.1", + "version": "0.2.2", "type": "module", "files": [ "./README.md", diff --git a/packages/effect-lens/src/Lens.ts b/packages/effect-lens/src/Lens.ts index 2b2d06a..4089126 100644 --- a/packages/effect-lens/src/Lens.ts +++ b/packages/effect-lens/src/Lens.ts @@ -13,7 +13,7 @@ export type LensTypeId = typeof LensTypeId * 2. a `changes` stream that emits every subsequent update to `A`, and * 3. a `modify` effect that can transform the current value. */ -export interface Lens +export interface Lens extends Subscribable.Subscribable { readonly [LensTypeId]: LensTypeId @@ -32,7 +32,7 @@ export const LensImplTypeId: unique symbol = Symbol.for("@effect-fc/Lens/LensImp export type LensImplTypeId = typeof LensImplTypeId export declare namespace LensImpl { - export interface Resolved { + export interface Resolved { readonly value: A readonly commit: ( next: Effect.Effect @@ -44,7 +44,7 @@ export declare namespace LensImpl { } } -export abstract class LensImpl +export abstract class LensImpl extends Pipeable.Class() implements Lens { readonly [Readable.TypeId]: Readable.TypeId = Readable.TypeId readonly [Subscribable.TypeId]: Subscribable.TypeId = Subscribable.TypeId @@ -89,7 +89,7 @@ export const asSubscribable = ( export declare namespace LensLazyImpl { - export interface Source { + export interface Source { readonly get: Effect.Effect readonly changes: Stream.Stream readonly commit: (a: A) => Effect.Effect @@ -97,7 +97,7 @@ export declare namespace LensLazyImpl { } } -export class LensLazyImpl +export class LensLazyImpl extends LensImpl { constructor( readonly source: LensLazyImpl.Source, @@ -126,7 +126,7 @@ export const make = ( ): Lens => new LensLazyImpl(source) -export class UnwrappedLensImpl +export class UnwrappedLensImpl extends LensImpl { constructor( readonly effect: Effect.Effect, E1, R1> diff --git a/packages/effect-lens/src/Subscribable.test.ts b/packages/effect-lens/src/Subscribable.test.ts index 447f14f..ca70106 100644 --- a/packages/effect-lens/src/Subscribable.test.ts +++ b/packages/effect-lens/src/Subscribable.test.ts @@ -71,6 +71,29 @@ describe("Subscribable", () => { expect(result).toEqual(["fallback", ["fallback"]]) }) + test("zipLatestAll combines current values and change streams", async () => { + const zipped = Subscribable.zipLatestAll( + Subscribable.make({ + get: Effect.succeed(1), + changes: Stream.succeed(2), + }), + Subscribable.make({ + get: Effect.succeed("one"), + changes: Stream.succeed("two"), + }), + ) + + const result = await Effect.runPromise(Effect.all([ + zipped.get, + Stream.runCollect(zipped.changes), + ])) + + expect([result[0], Array.from(result[1])]).toEqual([ + [1, "one"], + [[2, "two"]], + ]) + }) + 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 c2ac183..53fc83c 100644 --- a/packages/effect-lens/src/Subscribable.ts +++ b/packages/effect-lens/src/Subscribable.ts @@ -40,6 +40,41 @@ export const mapOptionEffect: { onNone: () => Effect.succeed(Option.none()), }))) +/** + * Allows transforming only the `changes` stream of a `Subscribable`. + */ +export const mapStream: { + ( + f: (changes: Stream.Stream, NoInfer, NoInfer>) => Stream.Stream, NoInfer, NoInfer>, + ): (self: Subscribable.Subscribable) => Subscribable.Subscribable + ( + self: Subscribable.Subscribable, + f: (changes: Stream.Stream, NoInfer, NoInfer>) => Stream.Stream, NoInfer, NoInfer>, + ): Subscribable.Subscribable +} = Function.dual(2, ( + self: Subscribable.Subscribable, + f: (changes: Stream.Stream, NoInfer, NoInfer>) => Stream.Stream, NoInfer, NoInfer>, +): Subscribable.Subscribable => Subscribable.make({ + get get() { return self.get }, + get changes() { return f(self.changes) }, +})) + +/** + * Combines the current values and streams of changes from multiple `Subscribable` values. + */ +export const zipLatestAll = []>( + ...elements: T +): Subscribable.Subscribable< + [T[number]] extends [never] + ? never + : { [K in keyof T]: T[K] extends Subscribable.Subscribable ? A : never }, + [T[number]] extends [never] ? never : T[number] extends Subscribable.Subscribable ? E : never, + [T[number]] extends [never] ? never : T[number] extends Subscribable.Subscribable ? R : never +> => Subscribable.make({ + get: Effect.all(elements.map(view => view.get)), + changes: Stream.zipLatestAll(...elements.map(view => view.changes)), +}) as any + /** * Maps errors from both the current value and the stream of changes.