From 65a42fd5490b318724c811bb0f412804d1da9d1e Mon Sep 17 00:00:00 2001 From: Dax Raad Date: Thu, 16 Jul 2026 13:34:29 -0400 Subject: [PATCH] fix(core): debounce state reloads --- packages/core/src/state.ts | 47 ++++++++++++++++++++++++++---- packages/core/test/agent.test.ts | 5 +++- packages/core/test/catalog.test.ts | 5 +++- packages/core/test/state.test.ts | 31 +++++++++++++++++++- 4 files changed, 80 insertions(+), 8 deletions(-) diff --git a/packages/core/src/state.ts b/packages/core/src/state.ts index 88e9984d4a2..b6d18818014 100644 --- a/packages/core/src/state.ts +++ b/packages/core/src/state.ts @@ -1,6 +1,6 @@ export * as State from "./state" -import { Context, Effect, Scope, Semaphore } from "effect" +import { Clock, Context, Deferred, Effect, Scope, Semaphore } from "effect" /** * A replayable transform applied to a draft during reload. @@ -34,6 +34,7 @@ type Batch = { const CurrentBatch = Context.Reference("@opencode/State/CurrentBatch", { defaultValue: () => undefined, }) +const reloadDebounce = 500 export function batch(effect: Effect.Effect) { return Effect.gen(function* () { @@ -73,6 +74,10 @@ export interface Interface extends Transformable { 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) { @@ -97,7 +102,39 @@ export function create(options: Options): Inte yield* commit(next) }) - const reload = () => semaphore.withPermit(materialize()) + 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, @@ -117,7 +154,7 @@ export function create(options: Options): Inte return Effect.gen(function* () { const batch = yield* CurrentBatch if (batch?.active) { - batch.reloads.add(reload) + batch.reloads.add(materializeReload) return } yield* materialize() @@ -132,8 +169,8 @@ export function create(options: Options): Inte ) yield* Scope.addFinalizer(scope, dispose) const batch = yield* CurrentBatch - if (batch?.active) batch.reloads.add(reload) - else yield* reload() + if (batch?.active) batch.reloads.add(materializeReload) + else yield* materializeReload() return { dispose } }), ) diff --git a/packages/core/test/agent.test.ts b/packages/core/test/agent.test.ts index b60b6e560ca..6a0ca3df7e0 100644 --- a/packages/core/test/agent.test.ts +++ b/packages/core/test/agent.test.ts @@ -1,5 +1,6 @@ import { describe, expect } from "bun:test" import { Effect, Exit, Fiber, Layer, Scope, Stream } from "effect" +import { TestClock } from "effect/testing" import { AgentV2 } from "@opencode-ai/core/agent" import { EventV2 } from "@opencode-ai/core/event" import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder" @@ -76,7 +77,9 @@ describe("AgentV2", () => { ) description = "New description" hidden = false - yield* agent.reload() + const reload = yield* agent.reload().pipe(Effect.forkChild({ startImmediately: true })) + yield* TestClock.adjust("500 millis") + yield* Fiber.join(reload) expect(yield* agent.get(id)).toMatchObject({ description: "New description", hidden: false }) }), diff --git a/packages/core/test/catalog.test.ts b/packages/core/test/catalog.test.ts index ad667f288c2..f14feaee5f4 100644 --- a/packages/core/test/catalog.test.ts +++ b/packages/core/test/catalog.test.ts @@ -1,6 +1,7 @@ import { describe, expect } from "bun:test" import { Money } from "@opencode-ai/schema/money" import { Effect, Fiber, Layer, Stream } from "effect" +import { TestClock } from "effect/testing" import { Catalog } from "@opencode-ai/core/catalog" import { Integration } from "@opencode-ai/core/integration" import { Credential } from "@opencode-ai/core/credential" @@ -261,7 +262,9 @@ describe("CatalogV2", () => { expect((yield* catalog.model.default())?.id).toBe(old) configured = false - yield* catalog.reload() + const reload = yield* catalog.reload().pipe(Effect.forkChild({ startImmediately: true })) + yield* TestClock.adjust("500 millis") + yield* Fiber.join(reload) expect((yield* catalog.model.default())?.id).toBe(newest) }), ) diff --git a/packages/core/test/state.test.ts b/packages/core/test/state.test.ts index 505cd724e7e..b6d45cdc5f3 100644 --- a/packages/core/test/state.test.ts +++ b/packages/core/test/state.test.ts @@ -1,6 +1,7 @@ import { describe, expect } from "bun:test" import { State } from "@opencode-ai/core/state" import { Deferred, Effect, Exit, Fiber, Layer, Scope } from "effect" +import { TestClock } from "effect/testing" import { testEffect } from "./lib/effect" const it = testEffect(Layer.empty) @@ -51,7 +52,9 @@ describe("State", () => { expect(state.get().values).toEqual(["first"]) value = "second" - yield* state.reload() + const reload = yield* state.reload().pipe(Effect.forkChild({ startImmediately: true })) + yield* TestClock.adjust("500 millis") + yield* Fiber.join(reload) expect(state.get().values).toEqual(["second"]) }), ) @@ -112,4 +115,30 @@ describe("State", () => { expect(finalized).toBe(2) }), ) + + it.effect("debounces reload bursts", () => + Effect.gen(function* () { + let finalized = 0 + const state = State.create({ + initial: () => ({ values: [] as string[] }), + draft: (draft) => ({ add: (item: string) => draft.values.push(item) }), + finalize: () => Effect.sync(() => finalized++), + }) + yield* state.transform((draft) => { + draft.add("value") + }) + finalized = 0 + + const first = yield* state.reload().pipe(Effect.forkChild({ startImmediately: true })) + yield* TestClock.adjust("250 millis") + const second = yield* state.reload().pipe(Effect.forkChild({ startImmediately: true })) + yield* TestClock.adjust("499 millis") + expect(finalized).toBe(0) + yield* TestClock.adjust("1 millis") + yield* Fiber.join(first) + yield* Fiber.join(second) + + expect(finalized).toBe(1) + }), + ) })