Add Effect v4 version
Lint / lint (push) Successful in 48s

This commit is contained in:
Julien Valverdé
2026-06-22 02:04:20 +02:00
parent b7ea35006d
commit 091e102b23
41 changed files with 3832 additions and 2 deletions
+283
View File
@@ -0,0 +1,283 @@
import {
type Cause,
type Context,
type Duration,
Effect,
Equal,
Fiber,
Option,
Pipeable,
Predicate,
type Scope,
Semaphore,
Stream,
SubscriptionRef,
} from "effect"
import { Subscribable } from "effect-lens"
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<in out K extends Query.AnyKey, in out A, in out KE = never, in out KR = never, in out E = never, in out R = never, in out P = never>
extends Pipeable.Pipeable {
readonly [QueryTypeId]: QueryTypeId
readonly context: Context.Context<Scope.Scope | QueryClient.QueryClient | KR | R>
readonly key: Stream.Stream<K, KE, KR>
readonly f: (key: K) => Effect.Effect<A, E, R>
readonly initialProgress: P
readonly staleTime: Duration.Input
readonly refreshOnWindowFocus: boolean
readonly latestKey: Subscribable.Subscribable<Option.Option<K>>
readonly fiber: Subscribable.Subscribable<Option.Option<Fiber.Fiber<A, E>>>
readonly result: Subscribable.Subscribable<Result.Result<A, E, P>>
readonly latestFinalResult: Subscribable.Subscribable<Option.Option<Result.Final<A, E, P>>>
readonly run: Effect.Effect<void>
fetch(key: K): Effect.Effect<Result.Final<A, E, P>>
fetchSubscribable(key: K): Effect.Effect<Subscribable.Subscribable<Result.Result<A, E, P>>>
readonly refresh: Effect.Effect<Result.Final<A, E, P>, Cause.NoSuchElementError>
readonly refreshSubscribable: Effect.Effect<Subscribable.Subscribable<Result.Result<A, E, P>>, Cause.NoSuchElementError>
readonly invalidateCache: Effect.Effect<void>
invalidateCacheEntry(key: K): Effect.Effect<void>
}
export declare namespace Query {
export type AnyKey = readonly any[]
}
export class QueryImpl<in out K extends Query.AnyKey, in out A, in out KE = never, in out KR = never, in out E = never, in out R = never, in out P = never>
extends Pipeable.Class implements Query<K, A, KE, KR, E, R, P> {
readonly [QueryTypeId]: QueryTypeId = QueryTypeId
readonly latestKey: Subscribable.Subscribable<Option.Option<K>>
readonly fiber: Subscribable.Subscribable<Option.Option<Fiber.Fiber<A, E>>>
readonly result: Subscribable.Subscribable<Result.Result<A, E, P>>
readonly latestFinalResult: Subscribable.Subscribable<Option.Option<Result.Final<A, E, P>>>
constructor(
readonly context: Context.Context<Scope.Scope | QueryClient.QueryClient | KR | R>,
readonly key: Stream.Stream<K, KE, KR>,
readonly f: (key: K) => Effect.Effect<A, E, R>,
readonly initialProgress: P,
readonly staleTime: Duration.Input,
readonly refreshOnWindowFocus: boolean,
readonly latestKeyRef: SubscriptionRef.SubscriptionRef<Option.Option<K>>,
readonly fiberRef: SubscriptionRef.SubscriptionRef<Option.Option<Fiber.Fiber<A, E>>>,
readonly resultRef: SubscriptionRef.SubscriptionRef<Result.Result<A, E, P>>,
readonly latestFinalResultRef: SubscriptionRef.SubscriptionRef<Option.Option<Result.Final<A, E, P>>>,
readonly runSemaphore: Semaphore.Semaphore,
) {
super()
this.latestKey = fromSubscriptionRef(latestKeyRef)
this.fiber = fromSubscriptionRef(fiberRef)
this.result = fromSubscriptionRef(resultRef)
this.latestFinalResult = fromSubscriptionRef(latestFinalResultRef)
}
get run(): Effect.Effect<void> {
const focus = this.refreshOnWindowFocus && typeof window !== "undefined"
? Stream.runForEach(Stream.fromEventListener<FocusEvent>(window, "focus"), () => this.refreshSubscribable)
: Effect.succeed(undefined)
return Effect.all([
Stream.runForEach(this.key, key => this.fetchSubscribable(key)),
focus,
], { concurrency: "unbounded", discard: true }).pipe(
Effect.ignore,
this.runSemaphore.withPermits(1),
Effect.provide(this.context),
)
}
get interrupt(): Effect.Effect<void> {
return Effect.flatMap(SubscriptionRef.get(this.fiberRef), Option.match({
onSome: Fiber.interrupt,
onNone: () => Effect.succeed(undefined),
}))
}
fetch(key: K): Effect.Effect<Result.Final<A, E, P>> {
const self = this
return Effect.gen(function*() {
yield* self.interrupt
yield* SubscriptionRef.set(self.latestKeyRef, Option.some(key))
const previous = yield* SubscriptionRef.get(self.latestFinalResultRef)
const sub = yield* self.startCached(key, Option.isSome(previous)
? Result.willFetch(previous.value) as Result.Final<A, E, P>
: Result.initial())
return yield* self.watch(key, sub)
}).pipe(Effect.provide(this.context))
}
fetchSubscribable(key: K): Effect.Effect<Subscribable.Subscribable<Result.Result<A, E, P>>> {
const self = this
return Effect.gen(function*() {
yield* self.interrupt
yield* SubscriptionRef.set(self.latestKeyRef, Option.some(key))
const previous = yield* SubscriptionRef.get(self.latestFinalResultRef)
const sub = yield* self.startCached(key, Option.isSome(previous)
? Result.willFetch(previous.value) as Result.Final<A, E, P>
: Result.initial())
yield* Effect.forkScoped(self.watch(key, sub))
return sub
}).pipe(Effect.provide(this.context))
}
get refresh(): Effect.Effect<Result.Final<A, E, P>, Cause.NoSuchElementError> {
const self = this
return Effect.gen(function*() {
yield* self.interrupt
const key = yield* Effect.fromOption(yield* SubscriptionRef.get(self.latestKeyRef))
const previous = yield* SubscriptionRef.get(self.latestFinalResultRef)
const sub = yield* self.startCached(key, Option.isSome(previous)
? Result.willRefresh(previous.value) as Result.Final<A, E, P>
: Result.initial())
return yield* self.watch(key, sub)
}).pipe(Effect.provide(this.context))
}
get refreshSubscribable(): Effect.Effect<Subscribable.Subscribable<Result.Result<A, E, P>>, Cause.NoSuchElementError> {
const self = this
return Effect.gen(function*() {
yield* self.interrupt
const key = yield* Effect.fromOption(yield* SubscriptionRef.get(self.latestKeyRef))
const previous = yield* SubscriptionRef.get(self.latestFinalResultRef)
const sub = yield* self.startCached(key, Option.isSome(previous)
? Result.willRefresh(previous.value) as Result.Final<A, E, P>
: Result.initial())
yield* Effect.forkScoped(self.watch(key, sub))
return sub
}).pipe(Effect.provide(this.context))
}
startCached(
key: K,
initial: Result.Initial | Result.Final<A, E, P>,
): Effect.Effect<Subscribable.Subscribable<Result.Result<A, E, P>>, never, Scope.Scope | QueryClient.QueryClient | R> {
return Effect.flatMap(this.getCacheEntry(key), Option.match({
onSome: entry => Effect.flatMap(
QueryClient.isQueryClientCacheEntryStale(entry),
isStale => isStale
? this.start(key, Result.willRefresh(entry.result) as Result.Final<A, E, P>)
: Effect.succeed(Subscribable.make({
get: Effect.succeed(entry.result as Result.Result<A, E, P>),
changes: Stream.make(entry.result as Result.Result<A, E, P>),
})),
),
onNone: () => this.start(key, initial),
}))
}
start(
key: K,
initial: Result.Initial | Result.Final<A, E, P>,
): Effect.Effect<Subscribable.Subscribable<Result.Result<A, E, P>>, never, Scope.Scope | R> {
const self = this
return Effect.gen(function*() {
const [sub, fiber] = yield* Result.unsafeForkEffect<A, E, R, P>(
Effect.onExit(self.f(key), () => Effect.flatMap(
Effect.all([Effect.fiberId, SubscriptionRef.get(self.fiberRef)]),
([currentFiberId, current]) => Option.match(current, {
onSome: value => Equal.equals(currentFiberId, value.id)
? SubscriptionRef.set(self.fiberRef, Option.none())
: Effect.succeed(undefined),
onNone: () => Effect.succeed(undefined),
}),
)),
{ initial, initialProgress: self.initialProgress },
)
yield* SubscriptionRef.set(self.fiberRef, Option.some(fiber))
return sub
})
}
watch(
key: K,
sub: Subscribable.Subscribable<Result.Result<A, E, P>>,
): Effect.Effect<Result.Final<A, E, P>, never, QueryClient.QueryClient> {
return Effect.flatMap(sub.get, initial => Stream.runFoldEffect(
Stream.takeUntil(sub.changes, result => Result.isFinal(result) && !Result.hasFlag(result)),
() => initial,
(_, result) => Effect.as(SubscriptionRef.set(this.resultRef, result), result),
) as Effect.Effect<Result.Final<A, E, P>>).pipe(
Effect.tap(result => SubscriptionRef.set(this.latestFinalResultRef, Option.some(result))),
Effect.tap(result => Result.isSuccess(result)
? Effect.asVoid(this.setCacheEntry(key, result))
: Effect.succeed(undefined)),
)
}
makeCacheKey(key: K): QueryClient.QueryClientCacheKey {
return new QueryClient.QueryClientCacheKey(key, this.f as (key: Query.AnyKey) => Effect.Effect<unknown, unknown, unknown>)
}
getCacheEntry(key: K): Effect.Effect<Option.Option<QueryClient.QueryClientCacheEntry>, never, QueryClient.QueryClient> {
return Effect.flatMap(QueryClient.QueryClient, client => client.getCacheEntry(this.makeCacheKey(key)))
}
setCacheEntry(key: K, result: Result.Success<A>): Effect.Effect<QueryClient.QueryClientCacheEntry, never, QueryClient.QueryClient> {
return Effect.flatMap(QueryClient.QueryClient, client => client.setCacheEntry(this.makeCacheKey(key), result, this.staleTime))
}
get invalidateCache(): Effect.Effect<void> {
return Effect.flatMap(
QueryClient.QueryClient,
client => client.invalidateCacheEntries(this.f as (key: Query.AnyKey) => Effect.Effect<unknown, unknown, unknown>),
).pipe(Effect.provide(this.context))
}
invalidateCacheEntry(key: K): Effect.Effect<void> {
return Effect.flatMap(
QueryClient.QueryClient,
client => client.invalidateCacheEntry(this.makeCacheKey(key)),
).pipe(Effect.provide(this.context))
}
}
export const isQuery = (u: unknown): u is Query<readonly unknown[], unknown> => Predicate.hasProperty(u, QueryTypeId)
export declare namespace make {
export interface Options<K extends Query.AnyKey, A, KE = never, KR = never, E = never, R = never, P = never> {
readonly key: Stream.Stream<K, KE, KR>
readonly f: (key: NoInfer<K>) => Effect.Effect<A, E, Result.forkEffect.InputContext<R, NoInfer<P>>>
readonly initialProgress?: P
readonly staleTime?: Duration.Input
readonly refreshOnWindowFocus?: boolean
}
}
export const make = Effect.fnUntraced(function* <K extends Query.AnyKey, A, KE = never, KR = never, E = never, R = never, P = never>(
options: make.Options<K, A, KE, KR, E, R, P>,
): Effect.fn.Return<
Query<K, A, KE, KR, E, Result.forkEffect.OutputContext<R, P>, P>,
never,
Scope.Scope | QueryClient.QueryClient | KR | Result.forkEffect.OutputContext<R, P>
> {
const client = yield* QueryClient.QueryClient
return new QueryImpl<K, A, KE, KR, E, Result.forkEffect.OutputContext<R, P>, P>(
yield* Effect.context<Scope.Scope | QueryClient.QueryClient | KR | Result.forkEffect.OutputContext<R, P>>(),
options.key,
options.f as any,
options.initialProgress as P,
options.staleTime ?? client.defaultStaleTime,
options.refreshOnWindowFocus ?? client.defaultRefreshOnWindowFocus,
yield* SubscriptionRef.make(Option.none<K>()),
yield* SubscriptionRef.make(Option.none<Fiber.Fiber<A, E>>()),
yield* SubscriptionRef.make(Result.initial<A, E, P>()),
yield* SubscriptionRef.make(Option.none<Result.Final<A, E, P>>()),
yield* Semaphore.make(1),
)
})
export const service = <K extends Query.AnyKey, A, KE = never, KR = never, E = never, R = never, P = never>(
options: make.Options<K, A, KE, KR, E, R, P>,
): Effect.Effect<
Query<K, A, KE, KR, E, Result.forkEffect.OutputContext<R, P>, P>,
never,
Scope.Scope | QueryClient.QueryClient | KR | Result.forkEffect.OutputContext<R, P>
> => Effect.tap(make(options), query => Effect.asVoid(Effect.forkScoped(query.run)))
const fromSubscriptionRef = <A>(ref: SubscriptionRef.SubscriptionRef<A>): Subscribable.Subscribable<A> => Subscribable.make({
get: SubscriptionRef.get(ref),
changes: SubscriptionRef.changes(ref),
})