import { type Cause, type Context, DateTime, type Duration, Effect, Equal, Equivalence, Fiber, HashMap, identity, Option, Pipeable, Predicate, type Scope, Stream, Subscribable, SubscriptionRef } from "effect" import * as QueryClient from "./QueryClient.js" import * as Result from "./Result.js" export const QueryTypeId: unique symbol = Symbol.for("@effect-fc/Query/Query") export type QueryTypeId = typeof QueryTypeId export interface Query extends Pipeable.Pipeable { readonly [QueryTypeId]: QueryTypeId readonly context: Context.Context readonly key: Stream.Stream readonly f: (key: K) => Effect.Effect readonly initialProgress: P readonly staleTime: Duration.DurationInput readonly latestKey: Subscribable.Subscribable> readonly fiber: Subscribable.Subscribable>> readonly result: Subscribable.Subscribable> readonly latestFinalResult: Subscribable.Subscribable>> readonly run: Effect.Effect fetch(key: K): Effect.Effect> fetchSubscribable(key: K): Effect.Effect>> readonly refresh: Effect.Effect, Cause.NoSuchElementException> readonly refreshSubscribable: Effect.Effect>, Cause.NoSuchElementException> readonly invalidateCache: Effect.Effect invalidateCacheEntry(key: K): Effect.Effect } export declare namespace Query { export type AnyKey = readonly any[] } export class QueryImpl extends Pipeable.Class() implements Query { readonly [QueryTypeId]: QueryTypeId = QueryTypeId constructor( readonly context: Context.Context, readonly key: Stream.Stream, readonly f: (key: K) => Effect.Effect, readonly initialProgress: P, readonly staleTime: Duration.DurationInput, readonly latestKey: SubscriptionRef.SubscriptionRef>, readonly fiber: SubscriptionRef.SubscriptionRef>>, readonly result: SubscriptionRef.SubscriptionRef>, readonly latestFinalResult: SubscriptionRef.SubscriptionRef>>, readonly runSemaphore: Effect.Semaphore, ) { super() } get run(): Effect.Effect { return Effect.provide( Stream.runForEach(this.key, key => this.interrupt.pipe( Effect.andThen(SubscriptionRef.set(this.latestKey, Option.some(key))), Effect.andThen(this.latestFinalResult), Effect.andThen(previous => this.startCached(key, Option.isSome(previous) ? Result.willFetch(previous.value) as Result.Final : Result.initial() )), Effect.andThen(sub => Effect.forkScoped(this.watch(key, sub))), this.runSemaphore.withPermits(1), )), this.context, ) } get interrupt(): Effect.Effect { return Effect.andThen(this.fiber, Option.match({ onSome: Fiber.interrupt, onNone: () => Effect.void, })) } fetch(key: K): Effect.Effect> { return this.interrupt.pipe( Effect.andThen(SubscriptionRef.set(this.latestKey, Option.some(key))), Effect.andThen(this.latestFinalResult), Effect.andThen(previous => this.startCached(key, Option.isSome(previous) ? Result.willFetch(previous.value) as Result.Final : Result.initial() )), Effect.andThen(sub => this.watch(key, sub)), Effect.provide(this.context), ) } fetchSubscribable(key: K): Effect.Effect>> { return this.interrupt.pipe( Effect.andThen(SubscriptionRef.set(this.latestKey, Option.some(key))), Effect.andThen(this.latestFinalResult), Effect.andThen(previous => this.startCached(key, Option.isSome(previous) ? Result.willFetch(previous.value) as Result.Final : Result.initial() )), Effect.tap(sub => Effect.forkScoped(this.watch(key, sub))), Effect.provide(this.context), ) } get refresh(): Effect.Effect, Cause.NoSuchElementException> { return this.interrupt.pipe( Effect.andThen(Effect.Do), Effect.bind("latestKey", () => Effect.andThen(this.latestKey, identity)), Effect.bind("latestFinalResult", () => this.latestFinalResult), Effect.bind("subscribable", ({ latestKey, latestFinalResult }) => this.startCached(latestKey, Option.isSome(latestFinalResult) ? Result.willRefresh(latestFinalResult.value) as Result.Final : Result.initial() ) ), Effect.andThen(({ latestKey, subscribable }) => this.watch(latestKey, subscribable)), Effect.provide(this.context), ) } get refreshSubscribable(): Effect.Effect< Subscribable.Subscribable>, Cause.NoSuchElementException > { return this.interrupt.pipe( Effect.andThen(Effect.Do), Effect.bind("latestKey", () => Effect.andThen(this.latestKey, identity)), Effect.bind("latestFinalResult", () => this.latestFinalResult), Effect.bind("subscribable", ({ latestKey, latestFinalResult }) => this.startCached(latestKey, Option.isSome(latestFinalResult) ? Result.willRefresh(latestFinalResult.value) as Result.Final : Result.initial() ) ), Effect.tap(({ latestKey, subscribable }) => Effect.forkScoped(this.watch(latestKey, subscribable))), Effect.map(({ subscribable }) => subscribable), Effect.provide(this.context), ) } startCached( key: K, initial: Result.Initial | Result.Final, ): Effect.Effect< Subscribable.Subscribable>, never, Scope.Scope | QueryClient.QueryClient | R > { return Effect.andThen(this.getCacheEntry(key), Option.match({ onSome: entry => Effect.andThen( QueryClient.isQueryClientCacheEntryStale(entry, this.staleTime), isStale => isStale ? this.start(key, Result.willRefresh(entry.result) as Result.Final) : Effect.succeed(Subscribable.make({ get: Effect.succeed(entry.result as Result.Result), get changes() { return Stream.make(entry.result as Result.Result) }, })), ), onNone: () => this.start(key, initial), })) } start( key: K, initial: Result.Initial | Result.Final, ): Effect.Effect< Subscribable.Subscribable>, never, Scope.Scope | R > { return Result.unsafeForkEffect( Effect.onExit(this.f(key), () => Effect.andThen( Effect.all([Effect.fiberId, this.fiber]), ([currentFiberId, fiber]) => Option.match(fiber, { onSome: v => Equal.equals(currentFiberId, v.id()) ? SubscriptionRef.set(this.fiber, Option.none()) : Effect.void, onNone: () => Effect.void, }), )), { initial, initialProgress: this.initialProgress, } as Result.unsafeForkEffect.Options, ).pipe( Effect.tap(([, fiber]) => SubscriptionRef.set(this.fiber, Option.some(fiber))), Effect.map(([sub]) => sub), ) } watch( key: K, sub: Subscribable.Subscribable> ): Effect.Effect, never, QueryClient.QueryClient> { return sub.get.pipe( Effect.andThen(initial => Stream.runFoldEffect( sub.changes, initial, (_, result) => Effect.as(SubscriptionRef.set(this.result, result), result), ) as Effect.Effect>), Effect.tap(result => SubscriptionRef.set(this.latestFinalResult, Option.some(result))), Effect.tap(result => Result.isSuccess(result) ? this.updateCacheEntry(key, result) : Effect.void ), ) } makeCacheKey(key: K): QueryClient.QueryClientCacheKey { return new QueryClient.QueryClientCacheKey(key, this.f as (key: Query.AnyKey) => Effect.Effect) } getCacheEntry( key: K ): Effect.Effect, never, QueryClient.QueryClient> { return QueryClient.QueryClient.pipe( Effect.andThen(client => client.cache), Effect.map(HashMap.get(this.makeCacheKey(key))), ) } updateCacheEntry( key: K, result: Result.Success, ): Effect.Effect { return Effect.Do.pipe( Effect.bind("client", () => QueryClient.QueryClient), Effect.bind("now", () => DateTime.now), Effect.let("entry", ({ now }) => new QueryClient.QueryClientCacheEntry(result, now)), Effect.tap(({ client, entry }) => SubscriptionRef.update( client.cache, HashMap.set(this.makeCacheKey(key), entry), )), Effect.map(({ entry }) => entry), ) } get invalidateCache(): Effect.Effect { return QueryClient.QueryClient.pipe( Effect.andThen(client => SubscriptionRef.update( client.cache, HashMap.filter((_, key) => !Equivalence.strict()(key.f, this.f)), )), Effect.provide(this.context), ) } invalidateCacheEntry(key: K): Effect.Effect { return QueryClient.QueryClient.pipe( Effect.andThen(client => SubscriptionRef.update( client.cache, HashMap.remove(this.makeCacheKey(key)), )), Effect.provide(this.context), ) } } export const isQuery = (u: unknown): u is Query => Predicate.hasProperty(u, QueryTypeId) export declare namespace make { export interface Options { readonly key: Stream.Stream readonly f: (key: NoInfer) => Effect.Effect>> readonly initialProgress?: P readonly staleTime?: Duration.DurationInput } } export const make = Effect.fnUntraced(function* ( options: make.Options ): Effect.fn.Return< Query, P>, never, Scope.Scope | QueryClient.QueryClient | Result.forkEffect.OutputContext > { const client = yield* QueryClient.QueryClient return new QueryImpl( yield* Effect.context>(), options.key, options.f as any, options.initialProgress as P, options.staleTime ?? client.defaultStaleTime, yield* SubscriptionRef.make(Option.none()), yield* SubscriptionRef.make(Option.none>()), yield* SubscriptionRef.make(Result.initial()), yield* SubscriptionRef.make(Option.none>()), yield* Effect.makeSemaphore(1), ) }) export const service = ( options: make.Options ): Effect.Effect< Query, P>, never, Scope.Scope | QueryClient.QueryClient | Result.forkEffect.OutputContext > => Effect.tap( make(options), query => Effect.forkScoped(query.run), )