diff --git a/packages/effect-fc-next/src/Mutation.ts b/packages/effect-fc-next/src/Mutation.ts index 0d873eb..5ae14d4 100644 --- a/packages/effect-fc-next/src/Mutation.ts +++ b/packages/effect-fc-next/src/Mutation.ts @@ -1,146 +1,157 @@ -import { type Context, Effect, Equal, type Fiber, Option, Pipeable, Predicate, type Scope, Stream, SubscriptionRef } from "effect" -import { Subscribable } from "effect-lens" -import * as Result from "./Result.js" +import { type Context, Effect, Equal, Exit, type Fiber, Option, Pipeable, Predicate, type Scope, Stream, SubscriptionRef } from "effect" +import { AsyncResult } from "effect/unstable/reactivity" +import * as Lens from "./Lens.js" +import type * as Subscribable from "./Subscribable.js" export const MutationTypeId: unique symbol = Symbol.for("@effect-fc/Mutation/Mutation") export type MutationTypeId = typeof MutationTypeId -export interface Mutation +export interface Mutation extends Pipeable.Pipeable { readonly [MutationTypeId]: MutationTypeId readonly context: Context.Context readonly f: (key: K) => Effect.Effect - readonly initialProgress: P readonly latestKey: Subscribable.Subscribable> readonly fiber: Subscribable.Subscribable>> - readonly result: Subscribable.Subscribable> - readonly latestFinalResult: Subscribable.Subscribable>> + readonly result: Subscribable.Subscribable> + readonly latestFinalResult: Subscribable.Subscribable | AsyncResult.Failure>> - mutate(key: K): Effect.Effect> - mutateSubscribable(key: K): Effect.Effect>> + mutate(key: K): Effect.Effect | AsyncResult.Failure> + mutateSubscribable(key: K): Effect.Effect>> } export declare namespace Mutation { export type AnyKey = readonly any[] } -export class MutationImpl -extends Pipeable.Class implements Mutation { +export class MutationImpl +extends Pipeable.Class implements Mutation { readonly [MutationTypeId]: MutationTypeId = MutationTypeId - readonly latestKey: Subscribable.Subscribable> - readonly fiber: Subscribable.Subscribable>> - readonly result: Subscribable.Subscribable> - readonly latestFinalResult: Subscribable.Subscribable>> constructor( readonly context: Context.Context, readonly f: (key: K) => Effect.Effect, - readonly initialProgress: P, - readonly latestKeyRef: SubscriptionRef.SubscriptionRef>, - readonly fiberRef: SubscriptionRef.SubscriptionRef>>, - readonly resultRef: SubscriptionRef.SubscriptionRef>, - readonly latestFinalResultRef: SubscriptionRef.SubscriptionRef>>, + readonly latestKey: Lens.Lens>, + readonly fiber: Lens.Lens>>, + readonly result: Lens.Lens>, + readonly latestFinalResult: Lens.Lens | AsyncResult.Failure>>, ) { super() - this.latestKey = fromSubscriptionRef(latestKeyRef) - this.fiber = fromSubscriptionRef(fiberRef) - this.result = fromSubscriptionRef(resultRef) - this.latestFinalResult = fromSubscriptionRef(latestFinalResultRef) } - mutate(key: K): Effect.Effect> { - return SubscriptionRef.set(this.latestKeyRef, Option.some(key)).pipe( + mutate(key: K): Effect.Effect | AsyncResult.Failure> { + return Lens.set(this.latestKey, Option.some(key)).pipe( Effect.andThen(this.start(key)), - Effect.andThen(sub => this.watch(sub)), + Effect.flatMap(state => this.watch(state)), Effect.provide(this.context), ) } - mutateSubscribable(key: K): Effect.Effect>> { - return SubscriptionRef.set(this.latestKeyRef, Option.some(key)).pipe( + mutateSubscribable(key: K): Effect.Effect>> { + return Lens.set(this.latestKey, Option.some(key)).pipe( Effect.andThen(this.start(key)), - Effect.tap(sub => Effect.forkScoped(this.watch(sub))), + Effect.tap(state => Effect.forkScoped(this.watch(state))), Effect.provide(this.context), ) } start(key: K): Effect.Effect< - Subscribable.Subscribable>, + Subscribable.Subscribable>, never, Scope.Scope | R > { - const self = this - return Effect.gen(function*() { - const initial = yield* SubscriptionRef.get(self.latestFinalResultRef) - const [sub, fiber] = yield* Result.unsafeForkEffect( - Effect.onExit(self.f(key), () => Effect.andThen( - Effect.all([Effect.fiberId, SubscriptionRef.get(self.fiberRef)]), - ([currentFiberId, fiber]) => Option.match(fiber, { - onSome: v => Equal.equals(currentFiberId, v.id) - ? SubscriptionRef.set(self.fiberRef, Option.none()) - : Effect.succeed(undefined), - onNone: () => Effect.succeed(undefined), - }), - )), + return Effect.gen({ self: this }, function*() { + const previous = yield* Lens.get(this.latestFinalResult) + const state = Lens.fromSubscriptionRef(yield* SubscriptionRef.make>( + Option.getOrElse(previous, () => AsyncResult.initial(false)) + )) - { - initial: Option.isSome(initial) ? Result.willFetch(initial.value) : Result.initial(), - initialProgress: self.initialProgress, - } as Result.unsafeForkEffect.Options, - ) - yield* SubscriptionRef.set(self.fiberRef, Option.some(fiber)) - return sub + const fiber = yield* Effect.forkScoped(Effect.andThen( + Lens.update(state, AsyncResult.match({ + onInitial: () => AsyncResult.initial(true), + onSuccess: v => AsyncResult.success(v.value, { + waiting: true, + }), + onFailure: v => AsyncResult.failure(v.cause, { + waiting: true, + previousSuccess: v.previousSuccess, + }) + })), + + Effect.onExit(this.f(key), exit => Lens.update( + state, + 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), + }), + onFailure: v => AsyncResult.failure(c, { + previousSuccess: v.previousSuccess, + }) + }), + }), + ).pipe( + Effect.andThen(Effect.all([ + Effect.fiberId, + Lens.get(this.fiber), + ])), + Effect.flatMap(([fiberId, fiber]) => Option.match(fiber, { + onSome: v => Equal.equals(fiberId, v.id) + ? Lens.set(this.fiber, Option.none()) + : Effect.void, + onNone: () => Effect.void, + })), + )), + )) + + yield* Lens.set(this.fiber, Option.some(fiber)) + return state }) } watch( - sub: Subscribable.Subscribable> - ): Effect.Effect> { - return sub.get.pipe( + state: Subscribable.Subscribable> + ): Effect.Effect | AsyncResult.Failure> { + return state.get.pipe( Effect.andThen(initial => Stream.runFoldEffect( - Stream.takeUntil(sub.changes, result => Result.isFinal(result) && !Result.hasFlag(result)), + state.changes, () => initial, - (_, result) => Effect.as(SubscriptionRef.set(this.resultRef, result), result), - ) as Effect.Effect>), - Effect.tap(result => SubscriptionRef.set(this.latestFinalResultRef, Option.some(result))), + (_, result) => Effect.as(Lens.set(this.result, result), result), + ) as Effect.Effect | AsyncResult.Failure>), + Effect.tap(result => Lens.set(this.latestFinalResult, Option.some(result))), ) } } -export const isMutation = (u: unknown): u is Mutation => Predicate.hasProperty(u, MutationTypeId) +export const isMutation = (u: unknown): u is Mutation => Predicate.hasProperty(u, MutationTypeId) export declare namespace make { - export interface Options { - readonly f: (key: K) => Effect.Effect>> - readonly initialProgress?: P + export interface Options { + readonly f: (key: K) => Effect.Effect } } -export const make = Effect.fnUntraced(function* ( - options: make.Options +export const make = Effect.fnUntraced(function* ( + options: make.Options ): Effect.fn.Return< - Mutation, P>, + Mutation, never, - Scope.Scope | Result.forkEffect.OutputContext + Scope.Scope | R > { return new MutationImpl( - yield* Effect.context>(), - options.f as any, - options.initialProgress as P, + yield* Effect.context(), + options.f, - yield* SubscriptionRef.make(Option.none()), - yield* SubscriptionRef.make(Option.none>()), - yield* SubscriptionRef.make(Result.initial()), - yield* SubscriptionRef.make(Option.none>()), + Lens.fromSubscriptionRef(yield* SubscriptionRef.make(Option.none())), + Lens.fromSubscriptionRef(yield* SubscriptionRef.make(Option.none>())), + Lens.fromSubscriptionRef(yield* SubscriptionRef.make>(AsyncResult.initial())), + Lens.fromSubscriptionRef(yield* SubscriptionRef.make(Option.none | AsyncResult.Failure>())), ) }) - -const fromSubscriptionRef = (ref: SubscriptionRef.SubscriptionRef): Subscribable.Subscribable => Subscribable.make({ - get: SubscriptionRef.get(ref), - changes: SubscriptionRef.changes(ref), -}) diff --git a/packages/effect-fc-next/src/Query.ts b/packages/effect-fc-next/src/Query.ts index 670f15d..27a6531 100644 --- a/packages/effect-fc-next/src/Query.ts +++ b/packages/effect-fc-next/src/Query.ts @@ -286,7 +286,7 @@ extends Pipeable.Class implements Query { } -export const isQuery = (u: unknown): u is Query => Predicate.hasProperty(u, QueryTypeId) +export const isQuery = (u: unknown): u is Query => Predicate.hasProperty(u, QueryTypeId) export declare namespace make {