export * as State from "./state" import { Clock, Context, Deferred, Effect, Scope, Semaphore } from "effect" /** * A replayable transform applied to a draft during reload. * * Domain drafts expose readable and writable state while preserving concise * plugin/config code. Transforms synchronously rebuild derived state. */ type TransformCallback = (draft: DraftApi) => void export type MakeDraft = (state: State) => DraftApi export interface Registration { readonly dispose: Effect.Effect } export type Transform = ( transform: TransformCallback, ) => Effect.Effect export type Reload = () => Effect.Effect export interface Transformable { readonly transform: Transform readonly reload: Reload } type Batch = { active: boolean readonly reloads: Set } const CurrentBatch = Context.Reference("@opencode/State/CurrentBatch", { defaultValue: () => undefined, }) const reloadDebounce = 500 export function batch(effect: Effect.Effect) { return Effect.gen(function* () { const current = yield* CurrentBatch if (current?.active) return yield* effect const batch: Batch = { active: true, reloads: new Set() } const exit = yield* effect.pipe(Effect.provideService(CurrentBatch, batch), Effect.exit) batch.active = false yield* Effect.forEach(batch.reloads, (reload) => reload(), { discard: true }) return yield* exit }) } export const inherit = Effect.fnUntraced(function* () { const batch = yield* CurrentBatch return (effect: Effect.Effect) => Effect.provideService(effect, CurrentBatch, batch) }) export interface Options { readonly name?: string /** Creates the base value for initial state and every scoped-transform reload. */ readonly initial: () => State /** Wraps mutable state in a domain-specific draft API. */ readonly draft: MakeDraft /** * Runs after the rebuilt state becomes visible. Update events published here * act as read barriers: subscribers refetching on the event observe the * committed state. */ readonly finalize?: (draft: DraftApi) => Effect.Effect } export interface Interface extends Transformable { readonly get: () => State /** * Registers and applies a scoped transform. Closing the owning Scope removes * the transform and reloads the materialized state. */ } export function create(options: Options): Interface { let state = options.initial() let transforms: { run: TransformCallback }[] = [] let generation = 0 let requestedAt = 0 let running = false let waiters: { generation: number; done: Deferred.Deferred }[] = [] const semaphore = Semaphore.makeUnsafe(1) const commit = Effect.fn("State.commit")(function* (next: State) { state = next if (options.finalize) yield* options.finalize(options.draft(next)) }) const apply = (transform: TransformCallback, draft: DraftApi) => Effect.sync(() => { transform(draft) }) const materialize = Effect.fnUntraced(function* () { const next = options.initial() const api = options.draft(next) for (const transform of transforms) yield* apply(transform.run, api).pipe( Effect.withSpan("State.reload.update", { attributes: { state: options.name ?? "anonymous" } }), ) yield* commit(next) }) const materializeReload = () => semaphore.withPermit(materialize()) const rebuild = (): Effect.Effect => Effect.gen(function* () { const clock = yield* Clock.Clock const remaining = requestedAt + reloadDebounce - clock.currentTimeMillisUnsafe() if (remaining > 0) yield* Effect.sleep(remaining) if (clock.currentTimeMillisUnsafe() < requestedAt + reloadDebounce) return yield* rebuild() const target = generation const exit = yield* materializeReload().pipe(Effect.exit) const completed = waiters.filter((waiter) => waiter.generation <= target) waiters = waiters.filter((waiter) => waiter.generation > target) yield* Effect.forEach(completed, (waiter) => Deferred.done(waiter.done, exit), { concurrency: "unbounded", discard: true, }) if (generation > target) return yield* rebuild() running = false }) const reload = Effect.fnUntraced(function* () { const done = Deferred.makeUnsafe() const clock = yield* Clock.Clock generation++ requestedAt = clock.currentTimeMillisUnsafe() waiters.push({ generation, done }) if (!running) { running = true yield* rebuild().pipe(Effect.forkDetach) } return yield* Deferred.await(done) }) const result: Interface = { get: () => state, transform: Effect.fn("State.transform")(function* (update) { yield* Effect.annotateCurrentSpan("state", options.name ?? "anonymous") const scope = yield* Scope.Scope return yield* Effect.uninterruptible( Effect.gen(function* () { const transform = { run: update } let active = true const dispose = Effect.uninterruptible( semaphore.withPermit( Effect.suspend(() => { if (!active) return Effect.void active = false transforms = transforms.filter((item) => item !== transform) return Effect.gen(function* () { const batch = yield* CurrentBatch if (batch?.active) { batch.reloads.add(materializeReload) return } yield* materialize() }) }), ), ) yield* semaphore.withPermit( Effect.sync(() => { transforms = [...transforms, transform] }), ) yield* Scope.addFinalizer(scope, dispose) const batch = yield* CurrentBatch if (batch?.active) batch.reloads.add(materializeReload) else yield* materializeReload() return { dispose } }), ) }), reload, } return result }