Skip to content

Commit 32d3535

Browse files
authored
fix(core): publish batched deltas before the next block starts (anomalyco#51105)
1 parent 9810d98 commit 32d3535

2 files changed

Lines changed: 58 additions & 1 deletion

File tree

‎packages/core/src/session/runner/publish-llm-event.ts‎

Lines changed: 18 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -193,7 +193,14 @@ export const createLLMEventPublisher = (bus: Pick<Bus.Interface, "publish">, inp
193193
const flush = Effect.fnUntraced(function* () {
194194
for (const id of Array.from(chunks.keys())) yield* end(id)
195195
})
196-
return { start, append, end, flush, has: (id: string) => chunks.has(id) }
196+
/** Publish batched deltas now, keeping every fragment open. */
197+
const publishPending = Effect.fnUntraced(function* () {
198+
for (const [id, current] of Array.from(chunks)) {
199+
if (current.timer) yield* Fiber.interrupt(current.timer)
200+
yield* publishDelta(id)
201+
}
202+
})
203+
return { start, append, end, flush, publishPending, has: (id: string) => chunks.has(id) }
197204
}
198205

199206
const text = fragments(
@@ -255,6 +262,13 @@ export const createLLMEventPublisher = (bus: Pick<Bus.Interface, "publish">, inp
255262
}),
256263
)
257264

265+
// Deltas are batched, but block starts are not. Publishing held deltas first keeps each
266+
// block's content ahead of the next block, so the published order matches the model's.
267+
const publishPendingDeltas = Effect.fnUntraced(function* () {
268+
yield* text.publishPending()
269+
yield* reasoning.publishPending()
270+
})
271+
258272
const flushFragments = Effect.fnUntraced(function* () {
259273
yield* text.flush()
260274
yield* reasoning.flush()
@@ -276,6 +290,7 @@ export const createLLMEventPublisher = (bus: Pick<Bus.Interface, "publish">, inp
276290
}
277291
tools.set(event.id, tool)
278292
yield* toolInput.start(event.id)
293+
yield* publishPendingDeltas()
279294
yield* bus.publish(SessionEvent.Tool.Input.Started, {
280295
sessionID: input.sessionID,
281296
assistantMessageID,
@@ -390,6 +405,7 @@ export const createLLMEventPublisher = (bus: Pick<Bus.Interface, "publish">, inp
390405
case "text-start":
391406
outputStarted = true
392407
const startedTextOrdinal = yield* text.start(event.id, providerState(event.providerMetadata))
408+
yield* publishPendingDeltas()
393409
yield* bus.publish(SessionEvent.Text.Started, {
394410
sessionID: input.sessionID,
395411
assistantMessageID: yield* startAssistant(),
@@ -405,6 +421,7 @@ export const createLLMEventPublisher = (bus: Pick<Bus.Interface, "publish">, inp
405421
case "reasoning-start":
406422
outputStarted = true
407423
const startedReasoningOrdinal = yield* reasoning.start(event.id, providerState(event.providerMetadata))
424+
yield* publishPendingDeltas()
408425
yield* bus.publish(SessionEvent.Reasoning.Started, {
409426
sessionID: input.sessionID,
410427
assistantMessageID: yield* startAssistant(),

‎packages/core/test/session-runner-tool-events.test.ts‎

Lines changed: 40 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -387,6 +387,46 @@ it.effect("batches text deltas and flushes pending text before the terminal even
387387
}),
388388
)
389389

390+
it.effect("publishes batched deltas before the next block starts", () =>
391+
Effect.gen(function* () {
392+
const { published, publisher } = capture()
393+
const types = () =>
394+
published
395+
.map((event) => event.type)
396+
.filter((type) => type !== "session.step.started.1" && type !== "session.step.streamed")
397+
yield* Effect.forEach(
398+
[
399+
LLMEvent.reasoningStart({ id: "reasoning" }),
400+
LLMEvent.reasoningDelta({ id: "reasoning", text: "Plan the edits." }),
401+
LLMEvent.textStart({ id: "text" }),
402+
LLMEvent.textDelta({ id: "text", text: "Now the edits:" }),
403+
LLMEvent.toolInputStart({ id: "call", name: "edit" }),
404+
],
405+
publisher.publish,
406+
{ discard: true },
407+
)
408+
expect(types()).toEqual([
409+
"session.reasoning.started.1",
410+
"session.reasoning.delta",
411+
"session.text.started.1",
412+
"session.text.delta",
413+
"session.tool.input.started.1",
414+
])
415+
416+
// Blocks stay open: later chunks still batch, and nothing is published twice.
417+
yield* publisher.publish(LLMEvent.textDelta({ id: "text", text: " more" }))
418+
yield* TestClock.adjust("1 second")
419+
yield* publisher.publish(LLMEvent.textEnd({ id: "text" }))
420+
expect(published.filter((event) => event.type === "session.text.delta").map((event) => event.data)).toMatchObject([
421+
{ delta: "Now the edits:" },
422+
{ delta: " more" },
423+
])
424+
expect(published.find((event) => event.type === "session.text.ended.1")?.data).toMatchObject({
425+
text: "Now the edits: more",
426+
})
427+
}),
428+
)
429+
390430
it.effect("retains new chunks and orders text-end behind an in-flight timer publication", () =>
391431
Effect.gen(function* () {
392432
const entered = yield* Deferred.make<void>()

0 commit comments

Comments
 (0)