From 5945953555685d5288b2588b70303490282ea84e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Julien=20Valverd=C3=A9?= Date: Wed, 24 Jun 2026 00:54:10 +0200 Subject: [PATCH] Refactor QueryClient --- packages/effect-fc-next/src/QueryClient.ts | 167 ++++++++++----------- 1 file changed, 79 insertions(+), 88 deletions(-) diff --git a/packages/effect-fc-next/src/QueryClient.ts b/packages/effect-fc-next/src/QueryClient.ts index b811138..61dcf59 100644 --- a/packages/effect-fc-next/src/QueryClient.ts +++ b/packages/effect-fc-next/src/QueryClient.ts @@ -1,24 +1,8 @@ -import { - Context, - DateTime, - Duration, - Effect, - Equal, - Equivalence, - Hash, - HashMap, - Layer, - Option, - Pipeable, - Predicate, - Schedule, - Scope, - Semaphore, - SubscriptionRef, -} from "effect" -import { Subscribable } from "effect-lens" +import { type Cause, Context, DateTime, Duration, Effect, Equal, Equivalence, Hash, HashMap, type Option, Pipeable, Predicate, Schedule, type Scope, Semaphore, SubscriptionRef } from "effect" +import type { AsyncResult } from "effect/unstable/reactivity" +import * as Lens from "./Lens.js" import type * as Query from "./Query.js" -import type * as Result from "./Result.js" +import type * as Subscribable from "./Subscribable.js" export const QueryClientServiceTypeId: unique symbol = Symbol.for("@effect-fc/QueryClient/QueryClientService") @@ -26,89 +10,87 @@ export type QueryClientServiceTypeId = typeof QueryClientServiceTypeId export interface QueryClientService extends Pipeable.Pipeable { readonly [QueryClientServiceTypeId]: QueryClientServiceTypeId + readonly cache: Subscribable.Subscribable> - readonly cacheGcTime: Duration.Input - readonly defaultStaleTime: Duration.Input + readonly cacheGcTime: Duration.Duration + readonly defaultStaleTime: Duration.Duration readonly defaultRefreshOnWindowFocus: boolean - readonly run: Effect.Effect + + readonly run: Effect.Effect getCacheEntry(key: QueryClientCacheKey): Effect.Effect> - setCacheEntry(key: QueryClientCacheKey, result: Result.Success, staleTime: Duration.Input): Effect.Effect + setCacheEntry( + key: QueryClientCacheKey, + result: AsyncResult.Success, + staleTime: Duration.Duration, + ): Effect.Effect invalidateCacheEntries(f: (key: Query.Query.AnyKey) => Effect.Effect): Effect.Effect invalidateCacheEntry(key: QueryClientCacheKey): Effect.Effect } export class QueryClient extends Context.Service()( - "@effect-fc/QueryClient/QueryClient", -) { - static get Default(): Layer.Layer { - return Layer.effect(QueryClient)(service()) - } -} + "@effect-fc/QueryClient/QueryClient" +) {} -export class QueryClientServiceImpl extends Pipeable.Class implements QueryClientService { +export class QueryClientServiceImpl +extends Pipeable.Class +implements QueryClientService { readonly [QueryClientServiceTypeId]: QueryClientServiceTypeId = QueryClientServiceTypeId - readonly cache: Subscribable.Subscribable> constructor( - readonly cacheRef: SubscriptionRef.SubscriptionRef>, - readonly cacheGcTime: Duration.Input, - readonly defaultStaleTime: Duration.Input, + readonly cache: Lens.Lens>, + readonly cacheGcTime: Duration.Duration, + readonly defaultStaleTime: Duration.Duration, readonly defaultRefreshOnWindowFocus: boolean, readonly runSemaphore: Semaphore.Semaphore, ) { super() - this.cache = Subscribable.make({ - get: SubscriptionRef.get(cacheRef), - changes: SubscriptionRef.changes(cacheRef), - }) } - get run(): Effect.Effect { + get run(): Effect.Effect { return this.runSemaphore.withPermits(1)(Effect.repeat( - Effect.flatMap(DateTime.now, now => SubscriptionRef.update(this.cacheRef, HashMap.filter(entry => - Duration.isLessThan( - DateTime.distance(entry.lastAccessedAt, now), - Duration.sum(Duration.fromInputUnsafe(entry.staleTime), Duration.fromInputUnsafe(this.cacheGcTime)), - ) - ))), - Schedule.spaced("30 seconds"), + Effect.flatMap( + DateTime.now, + now => Lens.update(this.cache, HashMap.filter(entry => + Duration.isLessThan( + DateTime.distance(entry.lastAccessedAt, now), + Duration.sum(entry.staleTime, this.cacheGcTime), + ) + )), + ), + Schedule.spaced("30 second"), )) } getCacheEntry(key: QueryClientCacheKey): Effect.Effect> { - const self = this - return Effect.gen(function*() { - const entry = HashMap.get(yield* SubscriptionRef.get(self.cacheRef), key) - if (Option.isNone(entry)) return Option.none() - const now = yield* DateTime.now - const accessed = new QueryClientCacheEntry( - entry.value.result, - entry.value.staleTime, - entry.value.createdAt, - now, - ) - yield* SubscriptionRef.update(self.cacheRef, HashMap.set(key, accessed)) - return Option.some(accessed) - }) + return Effect.all([ + DateTime.now, + Effect.flatMap( + Effect.map(Lens.get(this.cache), HashMap.get(key)), + Effect.fromOption, + ), + ]).pipe( + Effect.map(([now, entry]) => new QueryClientCacheEntry(entry.result, entry.staleTime, entry.createdAt, now)), + Effect.tap(entry => Lens.update(this.cache, HashMap.set(key, entry))), + Effect.option, + ) } setCacheEntry( key: QueryClientCacheKey, - result: Result.Success, - staleTime: Duration.Input, + result: AsyncResult.Success, + staleTime: Duration.Duration, ): Effect.Effect { - return Effect.flatMap(DateTime.now, now => { - const entry = new QueryClientCacheEntry(result, staleTime, now, now) - return Effect.as(SubscriptionRef.update(this.cacheRef, HashMap.set(key, entry)), entry) - }) + return DateTime.now.pipe( + Effect.map(now => new QueryClientCacheEntry(result, staleTime, now, now)), + Effect.tap(entry => Lens.update(this.cache, HashMap.set(key, entry))), + ) } invalidateCacheEntries(f: (key: Query.Query.AnyKey) => Effect.Effect): Effect.Effect { - return SubscriptionRef.update(this.cacheRef, HashMap.filter((_, key) => !Equivalence.strictEqual()(key.f, f))) + return Lens.update(this.cache, HashMap.filter((_, key) => !Equivalence.strictEqual()(key.f, f))) } - invalidateCacheEntry(key: QueryClientCacheKey): Effect.Effect { - return SubscriptionRef.update(this.cacheRef, HashMap.remove(key)) + return Lens.update(this.cache, HashMap.remove(key)) } } @@ -122,11 +104,13 @@ export declare namespace make { } } -export const make = Effect.fnUntraced(function* (options: make.Options = {}): Effect.fn.Return { +export const make = Effect.fnUntraced(function* ( + options: make.Options = {} +): Effect.fn.Return { return new QueryClientServiceImpl( - yield* SubscriptionRef.make(HashMap.empty()), - options.cacheGcTime ?? "5 minutes", - options.defaultStaleTime ?? "0 minutes", + Lens.fromSubscriptionRef(yield* SubscriptionRef.make(HashMap.empty())), + yield* Effect.fromOption(Duration.fromInput(options.cacheGcTime ?? "5 minutes")), + yield* Effect.fromOption(Duration.fromInput(options.defaultStaleTime ?? "0 minutes")), options.defaultRefreshOnWindowFocus ?? true, yield* Semaphore.make(1), ) @@ -136,15 +120,20 @@ export declare namespace service { export interface Options extends make.Options {} } -export const service = (options?: service.Options): Effect.Effect => Effect.tap( +export const service = ( + options?: service.Options +): Effect.Effect => Effect.tap( make(options), - client => Effect.asVoid(Effect.forkScoped(client.run)), + client => Effect.forkScoped(client.run), ) + export const QueryClientCacheKeyTypeId: unique symbol = Symbol.for("@effect-fc/QueryClient/QueryClientCacheKey") export type QueryClientCacheKeyTypeId = typeof QueryClientCacheKeyTypeId -export class QueryClientCacheKey extends Pipeable.Class implements Equal.Equal { +export class QueryClientCacheKey +extends Pipeable.Class +implements Pipeable.Pipeable, Equal.Equal { readonly [QueryClientCacheKeyTypeId]: QueryClientCacheKeyTypeId = QueryClientCacheKeyTypeId constructor( @@ -154,28 +143,28 @@ export class QueryClientCacheKey extends Pipeable.Class implements Equal.Equal { super() } - [Equal.symbol](that: Equal.Equal): boolean { - return isQueryClientCacheKey(that) - && Equivalence.Array(Equal.asEquivalence())(this.key, that.key) - && Equivalence.strictEqual()(this.f, that.f) + [Equal.symbol](that: Equal.Equal) { + return isQueryClientCacheKey(that) && Equivalence.Array(Equal.asEquivalence())(this.key, that.key) && Equivalence.strictEqual()(this.f, that.f) } - - [Hash.symbol](): number { + [Hash.symbol]() { return Hash.combine(Hash.hash(this.f))(Hash.array(this.key)) } } export const isQueryClientCacheKey = (u: unknown): u is QueryClientCacheKey => Predicate.hasProperty(u, QueryClientCacheKeyTypeId) + export const QueryClientCacheEntryTypeId: unique symbol = Symbol.for("@effect-fc/QueryClient/QueryClientCacheEntry") export type QueryClientCacheEntryTypeId = typeof QueryClientCacheEntryTypeId -export class QueryClientCacheEntry extends Pipeable.Class { +export class QueryClientCacheEntry +extends Pipeable.Class +implements Pipeable.Pipeable { readonly [QueryClientCacheEntryTypeId]: QueryClientCacheEntryTypeId = QueryClientCacheEntryTypeId constructor( - readonly result: Result.Success, - readonly staleTime: Duration.Input, + readonly result: AsyncResult.Success, + readonly staleTime: Duration.Duration, readonly createdAt: DateTime.DateTime, readonly lastAccessedAt: DateTime.DateTime, ) { @@ -185,7 +174,9 @@ export class QueryClientCacheEntry extends Pipeable.Class { export const isQueryClientCacheEntry = (u: unknown): u is QueryClientCacheEntry => Predicate.hasProperty(u, QueryClientCacheEntryTypeId) -export const isQueryClientCacheEntryStale = (self: QueryClientCacheEntry): Effect.Effect => Effect.map( +export const isQueryClientCacheEntryStale = ( + self: QueryClientCacheEntry +): Effect.Effect => Effect.map( DateTime.now, - now => Duration.isGreaterThanOrEqualTo(DateTime.distance(self.createdAt, now), Duration.fromInputUnsafe(self.staleTime)), + now => Duration.isGreaterThanOrEqualTo(DateTime.distance(self.createdAt, now), self.staleTime), )