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 * as View from "./View.js" export const MutationTypeId: unique symbol = Symbol.for("@effect-fc/Mutation/Mutation") export type MutationTypeId = typeof MutationTypeId export interface Mutation extends Pipeable.Pipeable { readonly [MutationTypeId]: MutationTypeId readonly context: Context.Context readonly f: (key: K) => Effect.Effect readonly latestKey: View.View> readonly fiber: View.View>> readonly state: View.View> readonly latestFinalResult: View.View | AsyncResult.Failure>> mutate(key: K): Effect.Effect | AsyncResult.Failure> mutateView(key: K): Effect.Effect>> } export const isMutation = (u: unknown): u is Mutation => Predicate.hasProperty(u, MutationTypeId) export class MutationImpl extends Pipeable.Class implements Mutation { readonly [MutationTypeId]: MutationTypeId = MutationTypeId constructor( readonly context: Context.Context, readonly f: (key: K) => Effect.Effect, readonly latestKey: Lens.Lens>, readonly fiber: Lens.Lens>>, readonly state: Lens.Lens>, readonly latestFinalResult: Lens.Lens | AsyncResult.Failure>>, ) { super() } mutate(key: K): Effect.Effect | AsyncResult.Failure> { return Lens.set(this.latestKey, Option.some(key)).pipe( Effect.andThen(this.start(key)), Effect.flatMap(state => this.watch(state)), Effect.provide(this.context), ) } mutateView(key: K): Effect.Effect>> { return Lens.set(this.latestKey, Option.some(key)).pipe( Effect.andThen(this.start(key)), Effect.tap(state => Effect.forkScoped(this.watch(state))), Effect.provide(this.context), ) } start(key: K): Effect.Effect< View.View>, never, Scope.Scope | R > { 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)) )) 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( state: View.View> ): Effect.Effect | AsyncResult.Failure> { return View.get(state).pipe( Effect.andThen(initial => Stream.runFoldEffect( View.changes(state), () => initial, (_, result) => Effect.as(Lens.set(this.state, result), result), ) as Effect.Effect | AsyncResult.Failure>), Effect.tap(result => Lens.set(this.latestFinalResult, Option.some(result))), ) } } export declare namespace make { export interface Options { readonly f: (key: K) => Effect.Effect } } export const make = Effect.fnUntraced(function* ( options: make.Options ): Effect.fn.Return< Mutation, never, Scope.Scope | R > { return new MutationImpl( yield* Effect.context(), options.f, 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>())), ) })