diff --git a/packages/effect-fc-next/src/Query.ts b/packages/effect-fc-next/src/Query.ts index 6e7d589..ca169be 100644 --- a/packages/effect-fc-next/src/Query.ts +++ b/packages/effect-fc-next/src/Query.ts @@ -1,4 +1,4 @@ -import { type Cause, type Context, Duration, Effect, Equal, Exit, Fiber, Option, Pipeable, Predicate, type Scope, Semaphore, Stream, SubscriptionRef } from "effect" +import { type Cause, type Context, Duration, Effect, Equal, type Equivalence, Exit, Fiber, Option, Pipeable, Predicate, type Scope, Semaphore, Stream, SubscriptionRef } from "effect" import { AsyncResult } from "effect/unstable/reactivity" import * as Lens from "./Lens.js" import * as QueryClient from "./QueryClient.js" @@ -8,51 +8,61 @@ import * as Subscribable from "./Subscribable.js" export const QueryTypeId: unique symbol = Symbol.for("@effect-fc/Query/Query") export type QueryTypeId = typeof QueryTypeId -export interface Query +export interface Query extends Pipeable.Pipeable { readonly [QueryTypeId]: QueryTypeId - readonly context: Context.Context - readonly key: Stream.Stream + readonly context: Context.Context readonly f: (key: K) => Effect.Effect + readonly keyEquivalence: Equivalence.Equivalence readonly staleTime: Duration.Duration readonly refreshOnWindowFocus: boolean readonly latestKey: Subscribable.Subscribable> readonly fiber: Subscribable.Subscribable>> - readonly result: Subscribable.Subscribable> - readonly latestFinalResult: Subscribable.Subscribable | AsyncResult.Failure>> + readonly state: Subscribable.Subscribable> + readonly latestFinalState: Subscribable.Subscribable>> readonly run: Effect.Effect - fetch(key: K): Effect.Effect | AsyncResult.Failure, Cause.NoSuchElementError> - fetchSubscribable(key: K): Effect.Effect>, Cause.NoSuchElementError> - readonly refresh: Effect.Effect | AsyncResult.Failure, Cause.NoSuchElementError> - readonly refreshSubscribable: Effect.Effect>, Cause.NoSuchElementError> + fetch(key: K): Effect.Effect, Cause.NoSuchElementError> + fetchSubscribable(key: K): Effect.Effect>, Cause.NoSuchElementError> + readonly refresh: Effect.Effect, Cause.NoSuchElementError> + readonly refreshSubscribable: Effect.Effect>, Cause.NoSuchElementError> readonly invalidateCache: Effect.Effect invalidateCacheEntry(key: K): Effect.Effect } -export const isQuery = (u: unknown): u is Query => Predicate.hasProperty(u, QueryTypeId) +export interface QueryState { + readonly key: K + readonly result: AsyncResult.AsyncResult +} + +export interface FinalQueryState { + readonly key: K + readonly result: AsyncResult.Success | AsyncResult.Failure +} + +export const isQuery = (u: unknown): u is Query => Predicate.hasProperty(u, QueryTypeId) -export class QueryImpl -extends Pipeable.Class implements Query { +export class QueryImpl +extends Pipeable.Class implements Query { readonly [QueryTypeId]: QueryTypeId = QueryTypeId constructor( - readonly context: Context.Context, - readonly key: Stream.Stream, + readonly context: Context.Context, readonly f: (key: K) => Effect.Effect, + readonly keyEquivalence: Equivalence.Equivalence, readonly staleTime: Duration.Duration, readonly refreshOnWindowFocus: boolean, readonly latestKey: Lens.Lens>, readonly fiber: Lens.Lens>>, - readonly result: Lens.Lens>, - readonly latestFinalResult: Lens.Lens | AsyncResult.Failure>>, + readonly state: Lens.Lens>, + readonly latestFinalState: Lens.Lens>>, readonly runSemaphore: Semaphore.Semaphore, ) { @@ -60,20 +70,15 @@ extends Pipeable.Class implements Query { } get run(): Effect.Effect { - return Effect.all([ - Stream.runForEach(this.key, key => this.fetchSubscribable(key)), - - Effect.promise(() => import("@effect/platform-browser")).pipe( - Effect.flatMap(({ BrowserStream }) => this.refreshOnWindowFocus - ? Stream.runForEach( - BrowserStream.fromEventListenerWindow("focus"), - () => this.refreshSubscribable, - ) - : Effect.void - ), - Effect.catchDefect(() => Effect.void), + return Effect.promise(() => import("@effect/platform-browser")).pipe( + Effect.flatMap(({ BrowserStream }) => this.refreshOnWindowFocus + ? Stream.runForEach( + BrowserStream.fromEventListenerWindow("focus"), + () => this.refreshSubscribable, + ) + : Effect.void ), - ], { concurrency: "unbounded" }).pipe( + Effect.catchDefect(() => Effect.void), Effect.ignore, this.runSemaphore.withPermits(1), Effect.provide(this.context), @@ -87,117 +92,156 @@ extends Pipeable.Class implements Query { })) } - fetch(key: K): Effect.Effect< - AsyncResult.Success | AsyncResult.Failure, - Cause.NoSuchElementError - > { + fetch(key: K): Effect.Effect, Cause.NoSuchElementError> { return this.interrupt.pipe( Effect.andThen(Lens.set(this.latestKey, Option.some(key))), - Effect.andThen(this.startCached(key, AsyncResult.initial())), - Effect.flatMap(state => this.watch(key, state)), + Effect.andThen(this.startCached({ + key, + result: AsyncResult.initial(false) + })), + Effect.flatMap(state => this.watch(state)), Effect.provide(this.context), ) } fetchSubscribable(key: K): Effect.Effect< - Subscribable.Subscribable>, + Subscribable.Subscribable>, Cause.NoSuchElementError > { return this.interrupt.pipe( Effect.andThen(Lens.set(this.latestKey, Option.some(key))), - Effect.andThen(this.startCached(key, AsyncResult.initial(false))), - Effect.tap(state => Effect.forkScoped(this.watch(key, state))), + Effect.andThen(this.startCached({ + key, + result: AsyncResult.initial(false) + })), + Effect.tap(state => Effect.forkScoped(this.watch(state))), Effect.provide(this.context), ) } - get refresh(): Effect.Effect | AsyncResult.Failure, Cause.NoSuchElementError> { + get refresh(): Effect.Effect, Cause.NoSuchElementError> { return this.interrupt.pipe( Effect.andThen(Effect.Do), Effect.bind("latestKey", () => Effect.flatMap(Lens.get(this.latestKey), Effect.fromOption)), - Effect.bind("latestFinalResult", () => Lens.get(this.latestFinalResult)), - Effect.bind("subscribable", ({ latestKey, latestFinalResult }) => - this.startCached(latestKey, Option.getOrElse(latestFinalResult, () => AsyncResult.initial(true))) + Effect.bind("latestFinalState", () => Lens.get(this.latestFinalState)), + Effect.bind("subscribable", ({ latestKey, latestFinalState }) => + this.startCached(Option.getOrElse(latestFinalState, () => ({ + key: latestKey, + result: AsyncResult.initial(true), + }))) ), - Effect.flatMap(({ latestKey, subscribable }) => this.watch(latestKey, subscribable)), + Effect.flatMap(({ subscribable }) => this.watch(subscribable)), Effect.provide(this.context), ) } get refreshSubscribable(): Effect.Effect< - Subscribable.Subscribable>, + Subscribable.Subscribable>, Cause.NoSuchElementError > { return this.interrupt.pipe( Effect.andThen(Effect.Do), Effect.bind("latestKey", () => Effect.flatMap(Lens.get(this.latestKey), Effect.fromOption)), - Effect.bind("latestFinalResult", () => Lens.get(this.latestFinalResult)), - Effect.bind("subscribable", ({ latestKey, latestFinalResult }) => - this.startCached(latestKey, Option.getOrElse(latestFinalResult, () => AsyncResult.initial(true))) + Effect.bind("latestFinalState", () => Lens.get(this.latestFinalState)), + Effect.bind("subscribable", ({ latestKey, latestFinalState }) => + this.startCached(Option.getOrElse(latestFinalState, () => ({ + key: latestKey, + result: AsyncResult.initial(true), + }))) ), - Effect.tap(({ latestKey, subscribable }) => Effect.forkScoped(this.watch(latestKey, subscribable))), + Effect.tap(({ subscribable }) => Effect.forkScoped(this.watch(subscribable))), Effect.map(({ subscribable }) => subscribable), Effect.provide(this.context), ) } startCached( - key: K, - previous: AsyncResult.AsyncResult, + previous: QueryState, ): Effect.Effect< - Subscribable.Subscribable>, + Subscribable.Subscribable>, Cause.NoSuchElementError, Scope.Scope | QueryClient.QueryClient | R > { - return Effect.flatMap(this.getCacheEntry(key), Option.match({ + return Effect.flatMap(this.getCacheEntry(previous.key), Option.match({ onSome: entry => Effect.flatMap( QueryClient.isQueryClientCacheEntryStale(entry), isStale => isStale - ? this.start(key, entry.result as AsyncResult.AsyncResult) + ? this.start({ + key: previous.key, + result: entry.result as AsyncResult.AsyncResult, + }) : Effect.succeed(Subscribable.make({ - get: Effect.succeed(entry.result as AsyncResult.AsyncResult), - get changes() { return Stream.make(entry.result as AsyncResult.AsyncResult) }, + get: Effect.succeed({ + key: previous.key, + result: entry.result as AsyncResult.AsyncResult, + }), + get changes() { + return Stream.make({ + key: previous.key, + result: entry.result as AsyncResult.AsyncResult, + }) + }, })), ), - onNone: () => this.start(key, previous), + onNone: () => this.start(previous), })) } start( - key: K, - previous: AsyncResult.AsyncResult, + previous: QueryState, ): Effect.Effect< - Subscribable.Subscribable>, + Subscribable.Subscribable>, never, Scope.Scope | R > { return Effect.gen({ self: this }, function*() { - const state = Lens.fromSubscriptionRef(yield* SubscriptionRef.make(previous)) + const subscribable = Lens.fromSubscriptionRef(yield* SubscriptionRef.make(previous)) const fiber = yield* Effect.forkScoped(Effect.andThen( - Lens.update(state, AsyncResult.match({ - onInitial: () => AsyncResult.initial(true), - onSuccess: v => AsyncResult.success(v.value, { - waiting: true, + Lens.update(subscribable, previous => AsyncResult.match(previous.result, { + onInitial: () => ({ + key: previous.key, + result: AsyncResult.initial(true), + }), + onSuccess: result => ({ + key: previous.key, + result: AsyncResult.success(result.value, { + waiting: true, + }), + }), + onFailure: result => ({ + key: previous.key, + result: AsyncResult.failure(result.cause, { + waiting: true, + previousSuccess: result.previousSuccess, + }), }), - onFailure: v => AsyncResult.failure(v.cause, { - waiting: true, - previousSuccess: v.previousSuccess, - }) })), - Effect.onExit(this.f(key), exit => Lens.update( - state, + Effect.onExit(this.f(previous.key), exit => Lens.update( + subscribable, previous => Exit.match(exit, { - onSuccess: v => AsyncResult.success(v), - onFailure: c => AsyncResult.match(previous, { - onInitial: () => AsyncResult.failure(c), - onSuccess: v => AsyncResult.failure(c, { - previousSuccess: Option.some(v), + onSuccess: v => ({ + key: previous.key, + result: AsyncResult.success(v), + }), + onFailure: c => AsyncResult.match(previous.result, { + onInitial: () => ({ + key: previous.key, + result: AsyncResult.failure(c), + }), + onSuccess: v => ({ + key: previous.key, + result: AsyncResult.failure(c, { + previousSuccess: Option.some(v), + }), + }), + onFailure: v => ({ + key: previous.key, + result: AsyncResult.failure(c, { + previousSuccess: v.previousSuccess, + }), }), - onFailure: v => AsyncResult.failure(c, { - previousSuccess: v.previousSuccess, - }) }), }), ).pipe( @@ -206,7 +250,7 @@ extends Pipeable.Class implements Query { Lens.get(this.fiber), ])), Effect.flatMap(([fiberId, fiber]) => Option.match(fiber, { - onSome: v => Equal.equals(fiberId, v.id) + onSome: v => fiberId === v.id ? Lens.set(this.fiber, Option.none()) : Effect.void, onNone: () => Effect.void, @@ -215,23 +259,22 @@ extends Pipeable.Class implements Query { )) yield* Lens.set(this.fiber, Option.some(fiber)) - return state + return subscribable }) } watch( - key: K, - state: Subscribable.Subscribable> - ): Effect.Effect | AsyncResult.Failure, never, QueryClient.QueryClient> { - return state.get.pipe( + subscribable: Subscribable.Subscribable> + ): Effect.Effect, never, QueryClient.QueryClient> { + return subscribable.get.pipe( Effect.flatMap(initial => Stream.runFoldEffect( - state.changes, + subscribable.changes, () => initial, - (_, result) => Effect.as(Lens.set(this.result, result), result), - ) as Effect.Effect | AsyncResult.Failure>), - Effect.tap(result => Lens.set(this.latestFinalResult, Option.some(result))), - Effect.tap(result => AsyncResult.isSuccess(result) - ? this.setCacheEntry(key, result) + (_, state) => Effect.as(Lens.set(this.state, state), state), + ) as Effect.Effect>), + Effect.tap(state => Lens.set(this.latestFinalState, Option.some(state))), + Effect.tap(state => AsyncResult.isSuccess(state.result) + ? this.setCacheEntry(state.key, state.result) : Effect.void ), ) @@ -285,27 +328,27 @@ extends Pipeable.Class implements Query { } export declare namespace make { - export interface Options { - readonly key: Stream.Stream - readonly f: (key: NoInfer) => Effect.Effect + export interface Options { + readonly f: (key: K) => Effect.Effect + readonly keyEquivalence?: Equivalence.Equivalence, readonly staleTime?: Duration.Input readonly refreshOnWindowFocus?: boolean } } -export const make = Effect.fnUntraced(function* ( - options: make.Options +export const make = Effect.fnUntraced(function* ( + options: make.Options ): Effect.fn.Return< - Query, + Query, Cause.NoSuchElementError, - Scope.Scope | QueryClient.QueryClient | KR | R + Scope.Scope | QueryClient.QueryClient | R > { const client = yield* QueryClient.QueryClient return new QueryImpl( - yield* Effect.context(), - options.key, + yield* Effect.context(), options.f, + options.keyEquivalence ?? Equal.asEquivalence(), options.staleTime ? yield* Effect.fromOption(Duration.fromInput(options.staleTime)) : client.defaultStaleTime, options.refreshOnWindowFocus ?? client.defaultRefreshOnWindowFocus, @@ -319,12 +362,12 @@ export const make = Effect.fnUntraced(function* ( - options: make.Options +export const service = ( + options: make.Options ): Effect.Effect< - Query, + Query, Cause.NoSuchElementError, - Scope.Scope | QueryClient.QueryClient | KR | R + Scope.Scope | QueryClient.QueryClient | R > => Effect.tap( make(options), query => Effect.forkScoped(query.run),