From 825193400773cceab9b92f6bc63247d6dde3d580 Mon Sep 17 00:00:00 2001 From: Dax Raad Date: Sat, 15 Aug 2026 19:26:16 -0400 Subject: [PATCH] fix(core): batch initial streamed delta --- packages/core/src/session/runner/publish-llm-event.ts | 4 ++++ packages/core/test/session-runner-tool-events.test.ts | 11 ++++------- packages/core/test/session-runner.test.ts | 4 ++-- 3 files changed, 10 insertions(+), 9 deletions(-) diff --git a/packages/core/src/session/runner/publish-llm-event.ts b/packages/core/src/session/runner/publish-llm-event.ts index 1527a0f961e..7b17266ac1a 100644 --- a/packages/core/src/session/runner/publish-llm-event.ts +++ b/packages/core/src/session/runner/publish-llm-event.ts @@ -150,6 +150,10 @@ export const createLLMEventPublisher = (bus: Pick, inp if (!current) return yield* Effect.die(new Error(`${name} delta before start: ${id}`)) if (!current.pending) return undefined const now = yield* Clock.currentTimeMillis + if (!force && current.publishedAt === undefined) { + current.publishedAt = now + return undefined + } if (!force && current.publishedAt !== undefined && now - current.publishedAt < deltaBatchInterval) return undefined yield* delta(id, current.pending, current.ordinal) diff --git a/packages/core/test/session-runner-tool-events.test.ts b/packages/core/test/session-runner-tool-events.test.ts index ed6146465e7..0f36123698b 100644 --- a/packages/core/test/session-runner-tool-events.test.ts +++ b/packages/core/test/session-runner-tool-events.test.ts @@ -217,16 +217,13 @@ it.effect("batches text deltas and flushes pending text before the terminal even { discard: true }, ) - expect(published.filter((event) => event.type === "session.text.delta").map((event) => event.data)).toMatchObject([ - { delta: "one" }, - ]) + expect(published.filter((event) => event.type === "session.text.delta")).toHaveLength(0) yield* TestClock.adjust("99 millis") - expect(published.filter((event) => event.type === "session.text.delta")).toHaveLength(1) + expect(published.filter((event) => event.type === "session.text.delta")).toHaveLength(0) yield* TestClock.adjust("1 millis") yield* publisher.publish(LLMEvent.textDelta({ id: "text", text: " four" })) expect(published.filter((event) => event.type === "session.text.delta").map((event) => event.data)).toMatchObject([ - { delta: "one" }, - { delta: " two three four" }, + { delta: "one two three four" }, ]) yield* publisher.publish(LLMEvent.textDelta({ id: "text", text: " five" })) @@ -253,7 +250,7 @@ it.effect("batches reasoning deltas and flushes pending reasoning before the ter expect( published.filter((event) => event.type === "session.reasoning.delta").map((event) => event.data), - ).toMatchObject([{ delta: "one" }, { delta: " two three" }]) + ).toMatchObject([{ delta: "one two three" }]) expect(published.slice(-2).map((event) => event.type)).toEqual([ "session.reasoning.delta", "session.reasoning.ended.1", diff --git a/packages/core/test/session-runner.test.ts b/packages/core/test/session-runner.test.ts index c718b414b96..e55f1d8e278 100644 --- a/packages/core/test/session-runner.test.ts +++ b/packages/core/test/session-runner.test.ts @@ -767,7 +767,7 @@ const verifyEphemeralDeltas = (kind: FragmentKind) => yield* admit(session, prompt) const bus = yield* Bus.Service const live = fixture.delta - ? yield* bus.subscribe(fixture.delta).pipe(Stream.take(2), Stream.runCollect, Effect.forkScoped) + ? yield* bus.subscribe(fixture.delta).pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped) : undefined yield* Effect.yieldNow yield* TestLLM.push(fixture.completeEvents) @@ -785,7 +785,7 @@ const verifyEphemeralDeltas = (kind: FragmentKind) => : [] if (live) { const streamed = Array.from(yield* Fiber.join(live)) - expect(streamed).toHaveLength(2) + expect(streamed).toHaveLength(1) expect( streamed .map((event) => {