import * as Chunk from "@effect/data/Chunk" import type * as Context from "@effect/data/Context" import * as Debug from "@effect/data/Debug" import * as Either from "@effect/data/Either" import * as Equal from "@effect/data/Equal" import { pipe } from "@effect/data/Function" import * as Hash from "@effect/data/Hash" import type * as Option from "@effect/data/Option" import * as Cause from "@effect/io/Cause" import * as Effect from "@effect/io/Effect" import type * as CompletedRequestMap from "@effect/query/CompletedRequestMap" import type * as DataSource from "@effect/query/DataSource" import type * as Described from "@effect/query/Described" import * as completedRequestMap from "@effect/query/internal_effect_untraced/completedRequestMap" import * as described from "@effect/query/internal_effect_untraced/described" import type * as Request from "@effect/query/Request" /** @internal */ const DataSourceSymbolKey = "@effect/query/DataSource" /** @internal */ export const DataSourceTypeId: DataSource.DataSourceTypeId = Symbol.for( DataSourceSymbolKey ) as DataSource.DataSourceTypeId const dataSourceVariance = { _R: (_: never) => _, _A: (_: never) => _ } class DataSourceImpl implements DataSource.DataSource { readonly [DataSourceTypeId] = dataSourceVariance constructor( readonly identifier: string, readonly runAll: ( requests: Chunk.Chunk> ) => Effect.Effect ) {} [Hash.symbol](): number { return Hash.string(this.identifier) } [Equal.symbol](that: unknown): boolean { return isDataSource(that) && this.identifier === that.identifier } } /** @internal */ export const isDataSource = (u: unknown): u is DataSource.DataSource => typeof u === "object" && u != null && DataSourceTypeId in u /** @internal */ export const make = Debug.untracedMethod((restore) => ( identifier: string, runAll: (requests: Chunk.Chunk>) => Effect.Effect ): DataSource.DataSource, A> => new DataSourceImpl(identifier, (requests) => Effect.suspend(() => { const map = completedRequestMap.empty() return Effect.as( Effect.provideService(completedRequestMap.CompletedRequestMap, map)(restore(runAll)(requests)), map ) })) ) /** @internal */ export const makeBatched = Debug.untracedMethod((restore) => >( identifier: string, run: (requests: Chunk.Chunk) => Effect.Effect ): DataSource.DataSource, A> => new DataSourceImpl( identifier, Effect.reduce(completedRequestMap.empty(), (outerMap, requests) => { const newRequests = Chunk.filter(requests, (request) => !completedRequestMap.has(outerMap, request)) if (Chunk.isEmpty(newRequests)) { return Effect.succeed(outerMap) } const innerMap = completedRequestMap.empty() return pipe( restore(run)(newRequests), Effect.provideService(completedRequestMap.CompletedRequestMap, innerMap), Effect.map(() => completedRequestMap.combine(outerMap, innerMap)) ) }) ) ) /** @internal */ export const around = Debug.untracedDual< ( before: Described.Described>, after: Described.Described<(a: A2) => Effect.Effect> ) => ( self: DataSource.DataSource ) => DataSource.DataSource, ( self: DataSource.DataSource, before: Described.Described>, after: Described.Described<(a: A2) => Effect.Effect> ) => DataSource.DataSource >(3, (restore) => (self, before, after) => new DataSourceImpl( `${self.identifier}.around(${before.description}, ${after.description})`, (requests) => Effect.acquireUseRelease( before.value, () => restore(self.runAll)(requests), after.value ) )) /** @internal */ export const batchN = Debug.untracedDual< (n: number) => (self: DataSource.DataSource) => DataSource.DataSource, (self: DataSource.DataSource, n: number) => DataSource.DataSource >(2, (restore) => (self, n) => new DataSourceImpl( `${self.identifier}.batchN(${n})`, (requests) => n < 1 ? Effect.die(Cause.IllegalArgumentException("DataSource.batchN: n must be at least 1")) : restore(self.runAll)( Chunk.reduce( requests, Chunk.empty(), (acc, chunk) => Chunk.concat(acc, Chunk.chunksOf(chunk, n)) ) ) )) /** @internal */ export const contramap = Debug.untracedDual< , B extends Request.Request>( f: Described.Described<(_: B) => A> ) => (self: DataSource.DataSource) => DataSource.DataSource, , B extends Request.Request>( self: DataSource.DataSource, f: Described.Described<(_: B) => A> ) => DataSource.DataSource >(2, (restore) => (self, f) => new DataSourceImpl( `${self.identifier}.contramap(${f.description})`, (requests) => restore(self.runAll)(pipe(requests, Chunk.map(Chunk.map(restore(f.value))))) )) /** @internal */ export const contramapContext = Debug.untracedDual< ( f: Described.Described<(context: Context.Context) => Context.Context> ) => >(self: DataSource.DataSource) => DataSource.DataSource, , R0>( self: DataSource.DataSource, f: Described.Described<(context: Context.Context) => Context.Context> ) => DataSource.DataSource >(2, (restore) => , R0>( self: DataSource.DataSource, f: Described.Described<(context: Context.Context) => Context.Context> ) => new DataSourceImpl( `${self.identifier}.contramapContext(${f.description})`, (requests) => Effect.contramapContext( restore(self.runAll)(requests), (context: Context.Context) => restore(f.value)(context) ) )) /** @internal */ export const contramapEffect = Debug.untracedDual< , R2, B extends Request.Request>( f: Described.Described<(_: B) => Effect.Effect> ) => (self: DataSource.DataSource) => DataSource.DataSource, , R2, B extends Request.Request>( self: DataSource.DataSource, f: Described.Described<(_: B) => Effect.Effect> ) => DataSource.DataSource >(2, (restore) => (self, f) => new DataSourceImpl( `${self.identifier}.contramapEffect(${f.description})`, (requests) => Effect.flatMap( Effect.forEach(requests, Effect.forEachPar(restore(f.value))), (requests) => restore(self.runAll)(requests) ) )) /** @internal */ export const eitherWith = Debug.untracedDual< < A extends Request.Request, R2, B extends Request.Request, C extends Request.Request >( that: DataSource.DataSource, f: Described.Described<(_: C) => Either.Either> ) => (self: DataSource.DataSource) => DataSource.DataSource, < R, A extends Request.Request, R2, B extends Request.Request, C extends Request.Request >( self: DataSource.DataSource, that: DataSource.DataSource, f: Described.Described<(_: C) => Either.Either> ) => DataSource.DataSource >( 3, (restore) => < R, A extends Request.Request, R2, B extends Request.Request, C extends Request.Request >( self: DataSource.DataSource, that: DataSource.DataSource, f: Described.Described<(_: C) => Either.Either> ) => new DataSourceImpl( `${self.identifier}.eitherWith(${that.identifier})(${f.description})`, (batch) => pipe( Effect.forEach(batch, (requests) => { const [as, bs] = pipe( requests, Chunk.partitionMap(restore(f.value)) ) return Effect.zipWithPar( restore(self.runAll)(Chunk.of(as)), restore(that.runAll)(Chunk.of(bs)), (self, that) => completedRequestMap.combine(self, that) ) }), Effect.map(Chunk.reduce( completedRequestMap.empty(), (acc, curr) => completedRequestMap.combine(acc, curr) )) ) ) ) /** @internal */ export const fromFunction = Debug.untracedMethod((restore) => >( name: string, f: (request: A) => Request.Request.Success ): DataSource.DataSource => makeBatched(name, (requests) => Effect.map(completedRequestMap.CompletedRequestMap, (map) => pipe( requests, Chunk.forEach((request) => completedRequestMap.set( map, request, // @ts-expect-error Either.right(restore(f)(request)) ) ) ))) ) /** @internal */ export const fromFunctionBatched = Debug.untracedMethod((restore) => >( name: string, f: (chunk: Chunk.Chunk) => Chunk.Chunk> ): DataSource.DataSource => fromFunctionBatchedEffect(name, (as) => Effect.succeed(restore(f)(as))) ) /** @internal */ export const fromFunctionBatchedEffect = Debug.untracedMethod((restore) => >( name: string, f: (chunk: Chunk.Chunk) => Effect.Effect, Chunk.Chunk>> ): DataSource.DataSource => makeBatched(name, (requests) => Effect.flatMap(completedRequestMap.CompletedRequestMap, (map) => pipe( Effect.match( restore(f)(requests), (e): Chunk.Chunk, Request.Request.Success>]> => pipe(requests, Chunk.map((k) => [k, Either.left(e)] as const)), (bs): Chunk.Chunk, Request.Request.Success>]> => pipe(requests, Chunk.zip(pipe(bs, Chunk.map(Either.right)))) ), Effect.map(Chunk.forEach( ([k, v]) => completedRequestMap.set(map, k, v as any) )) ))) ) /** @internal */ export const fromFunctionBatchedOption = Debug.untracedMethod((restore) => >( name: string, f: (chunk: Chunk.Chunk) => Chunk.Chunk>> ): DataSource.DataSource => fromFunctionBatchedOptionEffect(name, (as) => Effect.succeed(restore(f)(as))) ) /** @internal */ export const fromFunctionBatchedOptionEffect = Debug.untracedMethod((restore) => >( name: string, f: ( chunk: Chunk.Chunk ) => Effect.Effect, Chunk.Chunk>>> ): DataSource.DataSource => makeBatched( name, (requests) => Effect.flatMap(completedRequestMap.CompletedRequestMap, (map) => Effect.map( Effect.match( restore(f)(requests), (e): Chunk.Chunk< readonly [ A, Either.Either, Option.Option>> ] > => pipe(requests, Chunk.map((k) => [k, Either.left(e)] as const)), (bs): Chunk.Chunk< readonly [ A, Either.Either, Option.Option>> ] > => pipe(requests, Chunk.zip(pipe(bs, Chunk.map(Either.right)))) ), Chunk.forEach(([k, v]) => completedRequestMap.setOption(map, k, v as any)) )) ) ) /** @internal */ export const fromFunctionBatchedWith = Debug.untracedMethod((restore) => >( name: string, f: (chunk: Chunk.Chunk) => Chunk.Chunk>, g: (value: Request.Request.Success) => Request.Request> ): DataSource.DataSource => fromFunctionBatchedWithEffect( name, (as) => Effect.succeed(restore(f)(as)), restore(g) ) ) /** @internal */ export const fromFunctionBatchedWithEffect = Debug.untracedMethod((restore) => >( name: string, f: (chunk: Chunk.Chunk) => Effect.Effect, Chunk.Chunk>>, g: (b: Request.Request.Success) => Request.Request, Request.Request.Success> ): DataSource.DataSource => makeBatched(name, (requests) => Effect.flatMap(completedRequestMap.CompletedRequestMap, (map) => Effect.map( Effect.match( restore(f)(requests), (e): Chunk.Chunk< readonly [ Request.Request, Request.Request.Success>, Either.Either, Request.Request.Success> ] > => pipe(requests, Chunk.map((k) => [k, Either.left(e)] as const)), (bs): Chunk.Chunk< readonly [ Request.Request, Request.Request.Success>, Either.Either, Request.Request.Success> ] > => pipe(bs, Chunk.map((b) => [restore(g)(b), Either.right(b)] as const)) ), Chunk.forEach(([k, v]) => completedRequestMap.set(map, k, v)) ))) ) /** @internal */ export const fromFunctionEffect = Debug.untracedMethod((restore) => >( name: string, f: (a: A) => Effect.Effect, Request.Request.Success> ): DataSource.DataSource => makeBatched(name, (requests) => Effect.flatMap(completedRequestMap.CompletedRequestMap, (map) => Effect.map( Effect.forEachPar(requests, (a) => Effect.map( Effect.either(restore(f)(a)), (e) => [a, e] as const )), Chunk.forEach(([k, v]) => completedRequestMap.set(map, k, v as any)) ))) ) /** @internal */ export const fromFunctionOption = Debug.untracedMethod((restore) => >( name: string, f: (a: A) => Option.Option> ): DataSource.DataSource => fromFunctionOptionEffect(name, (a) => Effect.succeed(restore(f)(a))) ) /** @internal */ export const fromFunctionOptionEffect = Debug.untracedMethod((restore) => >( name: string, f: (a: A) => Effect.Effect, Option.Option>> ): DataSource.DataSource => makeBatched(name, (requests) => Effect.flatMap(completedRequestMap.CompletedRequestMap, (map) => Effect.map( Effect.forEachPar( requests, (a) => Effect.map(Effect.either(restore(f)(a)), (e) => [a, e] as const) ), Chunk.forEach(([k, v]) => completedRequestMap.setOption(map, k, v as any)) ))) ) /** @internal */ export const never = Debug.untracedMethod(() => (_: void): DataSource.DataSource => make("never", () => Effect.never()) ) /** @internal */ export const provideContext = Debug.untracedDual< ( context: Described.Described> ) => >( self: DataSource.DataSource ) => DataSource.DataSource, >( self: DataSource.DataSource, context: Described.Described> ) => DataSource.DataSource >(2, () => (self, context) => contramapContext( self, described.make(() => context.value, context.description) )) /** @internal */ export const race = Debug.untracedDual< >( that: DataSource.DataSource ) => >( self: DataSource.DataSource ) => DataSource.DataSource, , R2, A2 extends Request.Request>( self: DataSource.DataSource, that: DataSource.DataSource ) => DataSource.DataSource >( 2, (restore) => (self: DataSource.DataSource, that: DataSource.DataSource) => new DataSourceImpl(`${self.identifier}.race(${that.identifier})`, (requests) => Effect.race( restore(self.runAll)(requests as Chunk.Chunk>), restore(that.runAll)(requests as Chunk.Chunk>) )) )