From 0283a75b5e467a7d24b204ae5694509810e6b76e Mon Sep 17 00:00:00 2001 From: "sentry-junior[bot]" <264270552+sentry-junior[bot]@users.noreply.github.com> Date: Thu, 10 Sep 2026 18:29:09 +0000 Subject: [PATCH 01/10] fix(events): debounce watch event wakes Co-Authored-By: David Cramer --- packages/junior/src/chat/events/notification.ts | 3 +++ packages/junior/src/chat/task-execution/store.ts | 3 +++ packages/junior/tests/component/events/events.test.ts | 9 ++++++++- 3 files changed, 14 insertions(+), 1 deletion(-) diff --git a/packages/junior/src/chat/events/notification.ts b/packages/junior/src/chat/events/notification.ts index dd470b6682..8e537db3fe 100644 --- a/packages/junior/src/chat/events/notification.ts +++ b/packages/junior/src/chat/events/notification.ts @@ -14,6 +14,8 @@ import { eventGuidance } from "@/chat/events/catalog"; import { getEventCatalog } from "@/chat/events/runtime-catalog"; import { EVENT_AUTHOR_ID } from "@/chat/events/actor"; +const EVENT_WAKE_DELAY_MS = 2_000; + export interface EventNotification { eventKey: string; eventType: string; @@ -182,5 +184,6 @@ export async function enqueueEventNotification(args: { }), queue: args.queue, state: args.state, + wakeDelayMs: EVENT_WAKE_DELAY_MS, }); } diff --git a/packages/junior/src/chat/task-execution/store.ts b/packages/junior/src/chat/task-execution/store.ts index b2dc648e11..0da61ec3d2 100644 --- a/packages/junior/src/chat/task-execution/store.ts +++ b/packages/junior/src/chat/task-execution/store.ts @@ -246,6 +246,7 @@ async function enqueueAfterAppend(args: { nowMs?: number; queue: ConversationWorkQueue; state?: StateAdapter; + wakeDelayMs?: number; }): Promise { const nowMs = args.nowMs ?? now(); const appendResult = args.appendResult; @@ -284,6 +285,7 @@ async function enqueueAfterAppend(args: { conversationId: args.message.conversationId, conversationStore: args.conversationStore, idempotencyKey, + delayMs: args.wakeDelayMs, nowMs, queue: args.queue, state: args.state, @@ -310,6 +312,7 @@ export async function appendAndEnqueueInboundMessage(args: { nowMs?: number; queue: ConversationWorkQueue; state?: StateAdapter; + wakeDelayMs?: number; }): Promise { const nowMs = args.nowMs ?? now(); const result = await enqueueAfterAppend({ diff --git a/packages/junior/tests/component/events/events.test.ts b/packages/junior/tests/component/events/events.test.ts index ef0dfcfc53..f3ee7b9f14 100644 --- a/packages/junior/tests/component/events/events.test.ts +++ b/packages/junior/tests/component/events/events.test.ts @@ -134,6 +134,7 @@ describe("event delivery", () => { expect(queue.sentRecords()).toEqual([ { conversationId: CONVERSATION_ID, + delayMs: 2_000, idempotencyKey: `event:${subscription.id}:delivery-1:check-suite-1`, }, ]); @@ -244,10 +245,12 @@ describe("event delivery", () => { expect.arrayContaining([ { conversationId: CONVERSATION_ID, + delayMs: 2_000, idempotencyKey: `event:${threadWatch.id}:delivery-multi:check-suite-1`, }, { conversationId: "agent:deadbeefcafebabe", + delayMs: 2_000, idempotencyKey: `event:${opaqueWatch.id}:delivery-multi:check-suite-1`, }, ]), @@ -360,6 +363,7 @@ describe("event delivery", () => { }, { conversationId: CONVERSATION_ID, + delayMs: 2_000, idempotencyKey: `event:${subscription.id}:github:delivery-bridge:pull_request.comment.created`, }, ]), @@ -496,7 +500,9 @@ describe("event delivery", () => { ).resolves.toEqual({ enqueued: 1 }); } - expect(queue.sentRecords()).toHaveLength(1); + expect(queue.sentRecords()).toEqual([ + expect.objectContaining({ delayMs: 2_000 }), + ]); await expect( getConversationWorkState({ conversationId: CONVERSATION_ID }), ).resolves.toMatchObject({ messages: [{}, {}] }); @@ -540,6 +546,7 @@ describe("event delivery", () => { expect(queue.sentRecords()).toEqual([ { conversationId: CONVERSATION_ID, + delayMs: 2_000, idempotencyKey: `event:${subscription.id}:delivery-3:check-suite-1`, }, ]); From 7954acbc119dc436dbf299a2ba8f6fb4c28c2a6c Mon Sep 17 00:00:00 2001 From: "sentry-junior[bot]" <264270552+sentry-junior[bot]@users.noreply.github.com> Date: Thu, 10 Sep 2026 18:39:34 +0000 Subject: [PATCH 02/10] Revert "fix(events): debounce watch event wakes" This reverts commit 0283a75b5e467a7d24b204ae5694509810e6b76e. Co-Authored-By: David Cramer --- packages/junior/src/chat/events/notification.ts | 3 --- packages/junior/src/chat/task-execution/store.ts | 3 --- packages/junior/tests/component/events/events.test.ts | 9 +-------- 3 files changed, 1 insertion(+), 14 deletions(-) diff --git a/packages/junior/src/chat/events/notification.ts b/packages/junior/src/chat/events/notification.ts index 8e537db3fe..dd470b6682 100644 --- a/packages/junior/src/chat/events/notification.ts +++ b/packages/junior/src/chat/events/notification.ts @@ -14,8 +14,6 @@ import { eventGuidance } from "@/chat/events/catalog"; import { getEventCatalog } from "@/chat/events/runtime-catalog"; import { EVENT_AUTHOR_ID } from "@/chat/events/actor"; -const EVENT_WAKE_DELAY_MS = 2_000; - export interface EventNotification { eventKey: string; eventType: string; @@ -184,6 +182,5 @@ export async function enqueueEventNotification(args: { }), queue: args.queue, state: args.state, - wakeDelayMs: EVENT_WAKE_DELAY_MS, }); } diff --git a/packages/junior/src/chat/task-execution/store.ts b/packages/junior/src/chat/task-execution/store.ts index 0da61ec3d2..b2dc648e11 100644 --- a/packages/junior/src/chat/task-execution/store.ts +++ b/packages/junior/src/chat/task-execution/store.ts @@ -246,7 +246,6 @@ async function enqueueAfterAppend(args: { nowMs?: number; queue: ConversationWorkQueue; state?: StateAdapter; - wakeDelayMs?: number; }): Promise { const nowMs = args.nowMs ?? now(); const appendResult = args.appendResult; @@ -285,7 +284,6 @@ async function enqueueAfterAppend(args: { conversationId: args.message.conversationId, conversationStore: args.conversationStore, idempotencyKey, - delayMs: args.wakeDelayMs, nowMs, queue: args.queue, state: args.state, @@ -312,7 +310,6 @@ export async function appendAndEnqueueInboundMessage(args: { nowMs?: number; queue: ConversationWorkQueue; state?: StateAdapter; - wakeDelayMs?: number; }): Promise { const nowMs = args.nowMs ?? now(); const result = await enqueueAfterAppend({ diff --git a/packages/junior/tests/component/events/events.test.ts b/packages/junior/tests/component/events/events.test.ts index f3ee7b9f14..ef0dfcfc53 100644 --- a/packages/junior/tests/component/events/events.test.ts +++ b/packages/junior/tests/component/events/events.test.ts @@ -134,7 +134,6 @@ describe("event delivery", () => { expect(queue.sentRecords()).toEqual([ { conversationId: CONVERSATION_ID, - delayMs: 2_000, idempotencyKey: `event:${subscription.id}:delivery-1:check-suite-1`, }, ]); @@ -245,12 +244,10 @@ describe("event delivery", () => { expect.arrayContaining([ { conversationId: CONVERSATION_ID, - delayMs: 2_000, idempotencyKey: `event:${threadWatch.id}:delivery-multi:check-suite-1`, }, { conversationId: "agent:deadbeefcafebabe", - delayMs: 2_000, idempotencyKey: `event:${opaqueWatch.id}:delivery-multi:check-suite-1`, }, ]), @@ -363,7 +360,6 @@ describe("event delivery", () => { }, { conversationId: CONVERSATION_ID, - delayMs: 2_000, idempotencyKey: `event:${subscription.id}:github:delivery-bridge:pull_request.comment.created`, }, ]), @@ -500,9 +496,7 @@ describe("event delivery", () => { ).resolves.toEqual({ enqueued: 1 }); } - expect(queue.sentRecords()).toEqual([ - expect.objectContaining({ delayMs: 2_000 }), - ]); + expect(queue.sentRecords()).toHaveLength(1); await expect( getConversationWorkState({ conversationId: CONVERSATION_ID }), ).resolves.toMatchObject({ messages: [{}, {}] }); @@ -546,7 +540,6 @@ describe("event delivery", () => { expect(queue.sentRecords()).toEqual([ { conversationId: CONVERSATION_ID, - delayMs: 2_000, idempotencyKey: `event:${subscription.id}:delivery-3:check-suite-1`, }, ]); From c9e33c59f03cb26e368f9e4b278cdbea959a5ddd Mon Sep 17 00:00:00 2001 From: "sentry-junior[bot]" <264270552+sentry-junior[bot]@users.noreply.github.com> Date: Thu, 10 Sep 2026 18:45:27 +0000 Subject: [PATCH 03/10] fix(events): renewable, capped debounce for watch event wakes Replace the flat pre-delay with a debounce owned by the event Turn worker: the first event waits 30s, each additional pending event adds 5s, and the wait is capped at 60s measured from the first event. This lets a burst of watch events collect into one Turn without needing a persisted debounce deadline or a second scheduling mechanism. Co-Authored-By: David Cramer --- .../junior/src/chat/events/wake-debounce.ts | 29 ++++ .../chat/task-execution/conversation-turn.ts | 15 ++ .../junior/src/chat/task-execution/worker.ts | 3 + .../conversation-turn-work.test.ts | 134 ++++++++++++++++++ .../tests/unit/events/wake-debounce.test.ts | 68 +++++++++ 5 files changed, 249 insertions(+) create mode 100644 packages/junior/src/chat/events/wake-debounce.ts create mode 100644 packages/junior/tests/unit/events/wake-debounce.test.ts diff --git a/packages/junior/src/chat/events/wake-debounce.ts b/packages/junior/src/chat/events/wake-debounce.ts new file mode 100644 index 0000000000..33d645dda4 --- /dev/null +++ b/packages/junior/src/chat/events/wake-debounce.ts @@ -0,0 +1,29 @@ +/** + * Renewable wait before an event-sourced Turn runs, so a burst of watch + * events can collect into one Turn instead of one Turn per event. + * + * Each additional event in the batch nudges the wait later by a small step. + * The wait is bounded by a cap measured from the first event, so a busy + * watch cannot delay a Turn indefinitely. + */ +const EVENT_WAKE_BASE_DELAY_MS = 30_000; +const EVENT_WAKE_STEP_DELAY_MS = 5_000; +const EVENT_WAKE_MAX_DELAY_MS = 60_000; + +/** + * Return the remaining wait before an event-sourced Turn should run, or + * `undefined` once the batch has waited long enough to run now. + */ +export function remainingEventWakeDelayMs(args: { + eventCount: number; + firstReceivedAtMs: number; + nowMs: number; +}): number | undefined { + const waitMs = Math.min( + EVENT_WAKE_MAX_DELAY_MS, + EVENT_WAKE_BASE_DELAY_MS + + EVENT_WAKE_STEP_DELAY_MS * (args.eventCount - 1), + ); + const remainingMs = args.firstReceivedAtMs + waitMs - args.nowMs; + return remainingMs > 0 ? remainingMs : undefined; +} diff --git a/packages/junior/src/chat/task-execution/conversation-turn.ts b/packages/junior/src/chat/task-execution/conversation-turn.ts index 568885e559..2c2b536c23 100644 --- a/packages/junior/src/chat/task-execution/conversation-turn.ts +++ b/packages/junior/src/chat/task-execution/conversation-turn.ts @@ -82,6 +82,7 @@ import { isEventMailboxMetadata, type EventMailboxMetadata, } from "@/chat/events/notification"; +import { remainingEventWakeDelayMs } from "@/chat/events/wake-debounce"; import { isEventConversationMessage } from "@/chat/events/actor"; function stableHex(...parts: string[]): string { @@ -174,6 +175,20 @@ export function createConversationTurnWorker( context: ConversationWorkerContext, resolved: MailboxTurnWork, ): Promise => { + if (resolved.kind === "mailbox") { + const first = resolved.batch[0]!; + if (isEventMailboxMetadata(first.message.input.metadata)) { + const delayMs = remainingEventWakeDelayMs({ + eventCount: resolved.batch.length, + firstReceivedAtMs: first.message.receivedAtMs, + nowMs: Date.now(), + }); + if (delayMs !== undefined) { + return { status: "deferred", delayMs }; + } + } + } + const lifecycle = new ConversationTurnLifecycleService( getConversationEventStore(), ); diff --git a/packages/junior/src/chat/task-execution/worker.ts b/packages/junior/src/chat/task-execution/worker.ts index 03937b019f..20f78abf54 100644 --- a/packages/junior/src/chat/task-execution/worker.ts +++ b/packages/junior/src/chat/task-execution/worker.ts @@ -64,6 +64,8 @@ export interface InboxAttempt { export interface ConversationWorkerResult { /** `paused` waits for an external wake but must resume if a stop raced it. */ status: "completed" | "deferred" | "lost_lease" | "paused" | "yielded"; + /** Wait before the next wake attempt. Only meaningful when `status` is `deferred`. */ + delayMs?: number; } export interface ConversationWorkProcessResult { @@ -734,6 +736,7 @@ async function processConversationWorkInContext( const wake = await ensureConversationWake({ conversationId, conversationStore: options.conversationStore, + delayMs: result.delayMs, idempotencyKey: nudgeIdempotencyKey( "deferred", conversationId, diff --git a/packages/junior/tests/integration/conversation-turn-work.test.ts b/packages/junior/tests/integration/conversation-turn-work.test.ts index 49df978c08..0a263098a3 100644 --- a/packages/junior/tests/integration/conversation-turn-work.test.ts +++ b/packages/junior/tests/integration/conversation-turn-work.test.ts @@ -658,6 +658,140 @@ describe("Conversation mailbox Turn work", () => { ]); }, 10_000); + it("collects a burst of events into one Turn behind a renewable debounce", async () => { + const { actor, conversationStore, queue, state } = + await createConversationFixture(); + const conversationId = createConversationId({ + actorEmail: actor.email, + idempotencyKey: "event-burst-1", + }); + await recordWebConversationActivity({ + actor, + conversationId, + conversationStore, + nowMs: 1, + }); + await conversationStore.recordActivity({ + activityAtMs: 1, + conversationId, + nowMs: 1, + title: "Events", + }); + const baseMs = 1_700_000_000_000; + const nowSpy = vi.spyOn(Date, "now").mockReturnValue(baseMs); + try { + const subscription = { + conversationId, + id: "resource-subscription-burst", + }; + const firstMessage = createEventInboundMessage({ + event: { + eventKey: "check-1", + eventType: "check_suite.completed", + identifier: "getsentry/junior#1563", + namespace: "github", + occurredAtMs: baseMs, + trustedSummary: "Check one failed", + }, + receivedAtMs: baseMs, + subscription, + text: "Check one failed", + }); + await appendAndEnqueueInboundMessage({ + conversationStore, + message: firstMessage, + queue, + state, + }); + + const agentRuns: AgentRun[] = []; + const worker = createConversationTurnWorker( + createModelAgentRunnerForRun((run) => { + agentRuns.push(run); + return createModelStream([ + { type: "text", text: "Handled both events." }, + ]); + }), + ); + const run = requireConversationTurn(worker); + + // The lone event has not waited the 30s base debounce yet. + await expect( + processConversationQueueMessage(queue.takeMessage(), { + conversationStore, + queue, + run, + state, + }), + ).resolves.toEqual({ status: "pending_requeued" }); + expect(agentRuns).toHaveLength(0); + expect(queue.sentRecords().at(-1)).toMatchObject({ delayMs: 30_000 }); + + // A follow-up event during the wait nudges the wait later by one step + // instead of restarting it, and does not schedule a second wake. + nowSpy.mockReturnValue(baseMs + 200); + const secondMessage = createEventInboundMessage({ + event: { + eventKey: "check-2", + eventType: "check_suite.completed", + identifier: "getsentry/junior#1563", + namespace: "github", + occurredAtMs: baseMs + 200, + trustedSummary: "Check two failed", + }, + receivedAtMs: baseMs + 200, + subscription, + text: "Check two failed", + }); + await expect( + appendAndEnqueueInboundMessage({ + conversationStore, + message: secondMessage, + queue, + state, + }), + ).resolves.toMatchObject({ status: "appended" }); + expect(queue.sentRecords()).toHaveLength(2); + + // The original wake fires at baseMs + 30_000. Two pending events extend + // the wait by one 5s step, so it defers again instead of running early. + nowSpy.mockReturnValue(baseMs + 30_000); + await expect( + processConversationQueueMessage(queue.takeMessage(), { + conversationStore, + queue, + run, + state, + }), + ).resolves.toEqual({ status: "pending_requeued" }); + expect(agentRuns).toHaveLength(0); + expect(queue.sentRecords().at(-1)).toMatchObject({ delayMs: 5_000 }); + + // Once the extended wait elapses, both events run together in one Turn. + nowSpy.mockReturnValue(baseMs + 35_000); + await expect( + processConversationQueueMessage(queue.takeMessage(), { + conversationStore, + queue, + run, + state, + }), + ).resolves.toEqual({ status: "completed" }); + expect(agentRuns).toHaveLength(1); + + const transcript = ( + await getConversationEventStore().loadHistory(conversationId) + ).flatMap((event) => + event.data.type === "message" && event.data.role === "user" + ? [event.data.text] + : [], + ); + expect(transcript).toEqual(["Check one failed\n\nCheck two failed"]); + } finally { + nowSpy.mockRestore(); + } + }, 10_000); + it("delivers a resumed event from the Conversation Location", async () => { const { conversationStore, queue, state } = await createConversationFixture(); diff --git a/packages/junior/tests/unit/events/wake-debounce.test.ts b/packages/junior/tests/unit/events/wake-debounce.test.ts new file mode 100644 index 0000000000..32174da004 --- /dev/null +++ b/packages/junior/tests/unit/events/wake-debounce.test.ts @@ -0,0 +1,68 @@ +import { describe, expect, it } from "vitest"; +import { remainingEventWakeDelayMs } from "@/chat/events/wake-debounce"; + +describe("remainingEventWakeDelayMs", () => { + it("waits the base delay for a single event", () => { + expect( + remainingEventWakeDelayMs({ + eventCount: 1, + firstReceivedAtMs: 1_000, + nowMs: 1_000, + }), + ).toBe(30_000); + }); + + it("adds a step per follow-up event", () => { + expect( + remainingEventWakeDelayMs({ + eventCount: 2, + firstReceivedAtMs: 1_000, + nowMs: 1_000, + }), + ).toBe(35_000); + expect( + remainingEventWakeDelayMs({ + eventCount: 3, + firstReceivedAtMs: 1_000, + nowMs: 1_000, + }), + ).toBe(40_000); + }); + + it("caps the wait measured from the first event", () => { + expect( + remainingEventWakeDelayMs({ + eventCount: 20, + firstReceivedAtMs: 1_000, + nowMs: 1_000, + }), + ).toBe(60_000); + }); + + it("returns undefined once the wait has elapsed", () => { + expect( + remainingEventWakeDelayMs({ + eventCount: 1, + firstReceivedAtMs: 1_000, + nowMs: 31_000, + }), + ).toBeUndefined(); + expect( + remainingEventWakeDelayMs({ + eventCount: 1, + firstReceivedAtMs: 1_000, + nowMs: 30_999, + }), + ).toBe(1); + }); + + it("still runs once the cap elapses even with more follow-up events", () => { + expect( + remainingEventWakeDelayMs({ + eventCount: 100, + firstReceivedAtMs: 1_000, + nowMs: 61_000, + }), + ).toBeUndefined(); + }); +}); From 13ea36656c399dad9d74b7250472624ac440646e Mon Sep 17 00:00:00 2001 From: "sentry-junior[bot]" <264270552+sentry-junior[bot]@users.noreply.github.com> Date: Thu, 10 Sep 2026 18:52:51 +0000 Subject: [PATCH 04/10] test(events): keep debounce coverage focused --- .../conversation-turn-work.test.ts | 134 --------------- .../integration/event-wake-debounce.test.ts | 160 ++++++++++++++++++ 2 files changed, 160 insertions(+), 134 deletions(-) create mode 100644 packages/junior/tests/integration/event-wake-debounce.test.ts diff --git a/packages/junior/tests/integration/conversation-turn-work.test.ts b/packages/junior/tests/integration/conversation-turn-work.test.ts index 0a263098a3..49df978c08 100644 --- a/packages/junior/tests/integration/conversation-turn-work.test.ts +++ b/packages/junior/tests/integration/conversation-turn-work.test.ts @@ -658,140 +658,6 @@ describe("Conversation mailbox Turn work", () => { ]); }, 10_000); - it("collects a burst of events into one Turn behind a renewable debounce", async () => { - const { actor, conversationStore, queue, state } = - await createConversationFixture(); - const conversationId = createConversationId({ - actorEmail: actor.email, - idempotencyKey: "event-burst-1", - }); - await recordWebConversationActivity({ - actor, - conversationId, - conversationStore, - nowMs: 1, - }); - await conversationStore.recordActivity({ - activityAtMs: 1, - conversationId, - nowMs: 1, - title: "Events", - }); - const baseMs = 1_700_000_000_000; - const nowSpy = vi.spyOn(Date, "now").mockReturnValue(baseMs); - try { - const subscription = { - conversationId, - id: "resource-subscription-burst", - }; - const firstMessage = createEventInboundMessage({ - event: { - eventKey: "check-1", - eventType: "check_suite.completed", - identifier: "getsentry/junior#1563", - namespace: "github", - occurredAtMs: baseMs, - trustedSummary: "Check one failed", - }, - receivedAtMs: baseMs, - subscription, - text: "Check one failed", - }); - await appendAndEnqueueInboundMessage({ - conversationStore, - message: firstMessage, - queue, - state, - }); - - const agentRuns: AgentRun[] = []; - const worker = createConversationTurnWorker( - createModelAgentRunnerForRun((run) => { - agentRuns.push(run); - return createModelStream([ - { type: "text", text: "Handled both events." }, - ]); - }), - ); - const run = requireConversationTurn(worker); - - // The lone event has not waited the 30s base debounce yet. - await expect( - processConversationQueueMessage(queue.takeMessage(), { - conversationStore, - queue, - run, - state, - }), - ).resolves.toEqual({ status: "pending_requeued" }); - expect(agentRuns).toHaveLength(0); - expect(queue.sentRecords().at(-1)).toMatchObject({ delayMs: 30_000 }); - - // A follow-up event during the wait nudges the wait later by one step - // instead of restarting it, and does not schedule a second wake. - nowSpy.mockReturnValue(baseMs + 200); - const secondMessage = createEventInboundMessage({ - event: { - eventKey: "check-2", - eventType: "check_suite.completed", - identifier: "getsentry/junior#1563", - namespace: "github", - occurredAtMs: baseMs + 200, - trustedSummary: "Check two failed", - }, - receivedAtMs: baseMs + 200, - subscription, - text: "Check two failed", - }); - await expect( - appendAndEnqueueInboundMessage({ - conversationStore, - message: secondMessage, - queue, - state, - }), - ).resolves.toMatchObject({ status: "appended" }); - expect(queue.sentRecords()).toHaveLength(2); - - // The original wake fires at baseMs + 30_000. Two pending events extend - // the wait by one 5s step, so it defers again instead of running early. - nowSpy.mockReturnValue(baseMs + 30_000); - await expect( - processConversationQueueMessage(queue.takeMessage(), { - conversationStore, - queue, - run, - state, - }), - ).resolves.toEqual({ status: "pending_requeued" }); - expect(agentRuns).toHaveLength(0); - expect(queue.sentRecords().at(-1)).toMatchObject({ delayMs: 5_000 }); - - // Once the extended wait elapses, both events run together in one Turn. - nowSpy.mockReturnValue(baseMs + 35_000); - await expect( - processConversationQueueMessage(queue.takeMessage(), { - conversationStore, - queue, - run, - state, - }), - ).resolves.toEqual({ status: "completed" }); - expect(agentRuns).toHaveLength(1); - - const transcript = ( - await getConversationEventStore().loadHistory(conversationId) - ).flatMap((event) => - event.data.type === "message" && event.data.role === "user" - ? [event.data.text] - : [], - ); - expect(transcript).toEqual(["Check one failed\n\nCheck two failed"]); - } finally { - nowSpy.mockRestore(); - } - }, 10_000); - it("delivers a resumed event from the Conversation Location", async () => { const { conversationStore, queue, state } = await createConversationFixture(); diff --git a/packages/junior/tests/integration/event-wake-debounce.test.ts b/packages/junior/tests/integration/event-wake-debounce.test.ts new file mode 100644 index 0000000000..1562e74b29 --- /dev/null +++ b/packages/junior/tests/integration/event-wake-debounce.test.ts @@ -0,0 +1,160 @@ +import { afterEach, describe, expect, it, vi } from "vitest"; +import { createEventInboundMessage } from "@/chat/events/notification"; +import type { AgentRun } from "@/chat/agent/types"; +import { getConversationEventStore } from "@/chat/db"; +import { + createConversationId, + recordWebConversationActivity, +} from "@/chat/conversations/web-input"; +import { + resolveMailboxTurnWork, + type MailboxTurnWork, +} from "@/chat/task-execution/mailbox-turn"; +import { createConversationTurnWorker } from "@/chat/task-execution/conversation-turn"; +import { appendAndEnqueueInboundMessage } from "@/chat/task-execution/store"; +import { processConversationQueueMessage } from "@/chat/task-execution/vercel-callback"; +import type { + ConversationWorkerContext, + ConversationWorkerResult, +} from "@/chat/task-execution/worker"; +import { + closeConversationFixture, + createConversationFixture, +} from "../fixtures/conversation"; +import { createModelAgentRunnerForRun } from "../fixtures/agent-runner"; +import { createModelStream } from "../fixtures/model-stream"; + +function requireConversationTurn( + worker: ( + context: ConversationWorkerContext, + resolved: MailboxTurnWork, + ) => Promise, +) { + return async ( + context: ConversationWorkerContext, + ): Promise => { + const resolved = await resolveMailboxTurnWork(context); + if (!resolved) throw new Error("Expected Conversation mailbox work"); + return await worker(context, resolved); + }; +} + +describe("event wake debounce", () => { + afterEach(async () => { + await closeConversationFixture(); + vi.restoreAllMocks(); + }); + + it("collects a burst of events into one Turn behind a renewable debounce", async () => { + const { actor, conversationStore, queue, state } = + await createConversationFixture(); + const conversationId = createConversationId({ + actorEmail: actor.email, + idempotencyKey: "event-burst-1", + }); + await recordWebConversationActivity({ + actor, + conversationId, + conversationStore, + nowMs: 1, + }); + await conversationStore.recordActivity({ + activityAtMs: 1, + conversationId, + nowMs: 1, + title: "Events", + }); + + const baseMs = 1_700_000_000_000; + const nowSpy = vi.spyOn(Date, "now").mockReturnValue(baseMs); + const subscription = { + conversationId, + id: "resource-subscription-burst", + }; + const eventMessage = (sequence: number, receivedAtMs: number) => + createEventInboundMessage({ + event: { + eventKey: `check-${sequence}`, + eventType: "check_suite.completed", + identifier: "getsentry/junior#1563", + namespace: "github", + occurredAtMs: receivedAtMs, + trustedSummary: `Check ${sequence} failed`, + }, + receivedAtMs, + subscription, + text: `Check ${sequence} failed`, + }); + + await appendAndEnqueueInboundMessage({ + conversationStore, + message: eventMessage(1, baseMs), + queue, + state, + }); + + const agentRuns: AgentRun[] = []; + const run = requireConversationTurn( + createConversationTurnWorker( + createModelAgentRunnerForRun((agentRun) => { + agentRuns.push(agentRun); + return createModelStream([ + { type: "text", text: "Handled both events." }, + ]); + }), + ), + ); + + await expect( + processConversationQueueMessage(queue.takeMessage(), { + conversationStore, + queue, + run, + state, + }), + ).resolves.toEqual({ status: "pending_requeued" }); + expect(agentRuns).toHaveLength(0); + expect(queue.sentRecords().at(-1)).toMatchObject({ delayMs: 30_000 }); + + nowSpy.mockReturnValue(baseMs + 200); + await appendAndEnqueueInboundMessage({ + conversationStore, + message: eventMessage(2, baseMs + 200), + queue, + state, + }); + expect(queue.sentRecords()).toHaveLength(2); + + nowSpy.mockReturnValue(baseMs + 30_000); + await expect( + processConversationQueueMessage(queue.takeMessage(), { + conversationStore, + queue, + run, + state, + }), + ).resolves.toEqual({ status: "pending_requeued" }); + expect(agentRuns).toHaveLength(0); + expect(queue.sentRecords().at(-1)).toMatchObject({ delayMs: 5_000 }); + + nowSpy.mockReturnValue(baseMs + 35_000); + await expect( + processConversationQueueMessage(queue.takeMessage(), { + conversationStore, + queue, + run, + state, + }), + ).resolves.toEqual({ status: "completed" }); + expect(agentRuns).toHaveLength(1); + + const userMessages = ( + await getConversationEventStore().loadHistory(conversationId) + ).flatMap((event) => + event.data.type === "message" && event.data.role === "user" + ? [event.data.text] + : [], + ); + expect(userMessages).toEqual(["Check 1 failed\n\nCheck 2 failed"]); + }); +}); From 45eee1130174817df00bda7b8b2c3a6b8ed2b100 Mon Sep 17 00:00:00 2001 From: "sentry-junior[bot]" <264270552+sentry-junior[bot]@users.noreply.github.com> Date: Thu, 10 Sep 2026 18:55:01 +0000 Subject: [PATCH 05/10] refactor(events): simplify watch event delay Co-Authored-By: David Cramer --- .../junior/src/chat/events/wake-debounce.ts | 29 -------- .../chat/task-execution/conversation-turn.ts | 20 +++--- ...ounce.test.ts => event-wake-delay.test.ts} | 4 +- .../tests/unit/events/wake-debounce.test.ts | 68 ------------------- 4 files changed, 13 insertions(+), 108 deletions(-) delete mode 100644 packages/junior/src/chat/events/wake-debounce.ts rename packages/junior/tests/integration/{event-wake-debounce.test.ts => event-wake-delay.test.ts} (97%) delete mode 100644 packages/junior/tests/unit/events/wake-debounce.test.ts diff --git a/packages/junior/src/chat/events/wake-debounce.ts b/packages/junior/src/chat/events/wake-debounce.ts deleted file mode 100644 index 33d645dda4..0000000000 --- a/packages/junior/src/chat/events/wake-debounce.ts +++ /dev/null @@ -1,29 +0,0 @@ -/** - * Renewable wait before an event-sourced Turn runs, so a burst of watch - * events can collect into one Turn instead of one Turn per event. - * - * Each additional event in the batch nudges the wait later by a small step. - * The wait is bounded by a cap measured from the first event, so a busy - * watch cannot delay a Turn indefinitely. - */ -const EVENT_WAKE_BASE_DELAY_MS = 30_000; -const EVENT_WAKE_STEP_DELAY_MS = 5_000; -const EVENT_WAKE_MAX_DELAY_MS = 60_000; - -/** - * Return the remaining wait before an event-sourced Turn should run, or - * `undefined` once the batch has waited long enough to run now. - */ -export function remainingEventWakeDelayMs(args: { - eventCount: number; - firstReceivedAtMs: number; - nowMs: number; -}): number | undefined { - const waitMs = Math.min( - EVENT_WAKE_MAX_DELAY_MS, - EVENT_WAKE_BASE_DELAY_MS + - EVENT_WAKE_STEP_DELAY_MS * (args.eventCount - 1), - ); - const remainingMs = args.firstReceivedAtMs + waitMs - args.nowMs; - return remainingMs > 0 ? remainingMs : undefined; -} diff --git a/packages/junior/src/chat/task-execution/conversation-turn.ts b/packages/junior/src/chat/task-execution/conversation-turn.ts index 2c2b536c23..6dc4ea37dd 100644 --- a/packages/junior/src/chat/task-execution/conversation-turn.ts +++ b/packages/junior/src/chat/task-execution/conversation-turn.ts @@ -82,9 +82,12 @@ import { isEventMailboxMetadata, type EventMailboxMetadata, } from "@/chat/events/notification"; -import { remainingEventWakeDelayMs } from "@/chat/events/wake-debounce"; import { isEventConversationMessage } from "@/chat/events/actor"; +const EVENT_WAIT_MS = 30_000; +const EVENT_WAIT_PER_EXTRA_MESSAGE_MS = 5_000; +const EVENT_MAX_WAIT_MS = 60_000; + function stableHex(...parts: string[]): string { return createHash("sha256") .update(parts.join("\u0000")) @@ -178,14 +181,13 @@ export function createConversationTurnWorker( if (resolved.kind === "mailbox") { const first = resolved.batch[0]!; if (isEventMailboxMetadata(first.message.input.metadata)) { - const delayMs = remainingEventWakeDelayMs({ - eventCount: resolved.batch.length, - firstReceivedAtMs: first.message.receivedAtMs, - nowMs: Date.now(), - }); - if (delayMs !== undefined) { - return { status: "deferred", delayMs }; - } + const waitMs = Math.min( + EVENT_MAX_WAIT_MS, + EVENT_WAIT_MS + + EVENT_WAIT_PER_EXTRA_MESSAGE_MS * (resolved.batch.length - 1), + ); + const delayMs = first.message.receivedAtMs + waitMs - Date.now(); + if (delayMs > 0) return { status: "deferred", delayMs }; } } diff --git a/packages/junior/tests/integration/event-wake-debounce.test.ts b/packages/junior/tests/integration/event-wake-delay.test.ts similarity index 97% rename from packages/junior/tests/integration/event-wake-debounce.test.ts rename to packages/junior/tests/integration/event-wake-delay.test.ts index 1562e74b29..546266024e 100644 --- a/packages/junior/tests/integration/event-wake-debounce.test.ts +++ b/packages/junior/tests/integration/event-wake-delay.test.ts @@ -39,13 +39,13 @@ function requireConversationTurn( }; } -describe("event wake debounce", () => { +describe("event wake delay", () => { afterEach(async () => { await closeConversationFixture(); vi.restoreAllMocks(); }); - it("collects a burst of events into one Turn behind a renewable debounce", async () => { + it("waits for a burst of events before running one Turn", async () => { const { actor, conversationStore, queue, state } = await createConversationFixture(); const conversationId = createConversationId({ diff --git a/packages/junior/tests/unit/events/wake-debounce.test.ts b/packages/junior/tests/unit/events/wake-debounce.test.ts deleted file mode 100644 index 32174da004..0000000000 --- a/packages/junior/tests/unit/events/wake-debounce.test.ts +++ /dev/null @@ -1,68 +0,0 @@ -import { describe, expect, it } from "vitest"; -import { remainingEventWakeDelayMs } from "@/chat/events/wake-debounce"; - -describe("remainingEventWakeDelayMs", () => { - it("waits the base delay for a single event", () => { - expect( - remainingEventWakeDelayMs({ - eventCount: 1, - firstReceivedAtMs: 1_000, - nowMs: 1_000, - }), - ).toBe(30_000); - }); - - it("adds a step per follow-up event", () => { - expect( - remainingEventWakeDelayMs({ - eventCount: 2, - firstReceivedAtMs: 1_000, - nowMs: 1_000, - }), - ).toBe(35_000); - expect( - remainingEventWakeDelayMs({ - eventCount: 3, - firstReceivedAtMs: 1_000, - nowMs: 1_000, - }), - ).toBe(40_000); - }); - - it("caps the wait measured from the first event", () => { - expect( - remainingEventWakeDelayMs({ - eventCount: 20, - firstReceivedAtMs: 1_000, - nowMs: 1_000, - }), - ).toBe(60_000); - }); - - it("returns undefined once the wait has elapsed", () => { - expect( - remainingEventWakeDelayMs({ - eventCount: 1, - firstReceivedAtMs: 1_000, - nowMs: 31_000, - }), - ).toBeUndefined(); - expect( - remainingEventWakeDelayMs({ - eventCount: 1, - firstReceivedAtMs: 1_000, - nowMs: 30_999, - }), - ).toBe(1); - }); - - it("still runs once the cap elapses even with more follow-up events", () => { - expect( - remainingEventWakeDelayMs({ - eventCount: 100, - firstReceivedAtMs: 1_000, - nowMs: 61_000, - }), - ).toBeUndefined(); - }); -}); From b200ebfbc04fcebc2b449d1df062e5ee3a6d2c38 Mon Sep 17 00:00:00 2001 From: "sentry-junior[bot]" <264270552+sentry-junior[bot]@users.noreply.github.com> Date: Thu, 10 Sep 2026 19:00:52 +0000 Subject: [PATCH 06/10] fix(events): delay the queue wake Co-Authored-By: David Cramer --- .../junior/src/chat/events/notification.ts | 3 +++ .../chat/task-execution/conversation-turn.ts | 2 +- .../junior/src/chat/task-execution/store.ts | 3 +++ .../tests/component/events/events.test.ts | 5 +++++ .../integration/event-wake-delay.test.ts | 22 ++++++++----------- 5 files changed, 21 insertions(+), 14 deletions(-) diff --git a/packages/junior/src/chat/events/notification.ts b/packages/junior/src/chat/events/notification.ts index dd470b6682..6da19e42d4 100644 --- a/packages/junior/src/chat/events/notification.ts +++ b/packages/junior/src/chat/events/notification.ts @@ -14,6 +14,8 @@ import { eventGuidance } from "@/chat/events/catalog"; import { getEventCatalog } from "@/chat/events/runtime-catalog"; import { EVENT_AUTHOR_ID } from "@/chat/events/actor"; +export const EVENT_WAIT_MS = 30_000; + export interface EventNotification { eventKey: string; eventType: string; @@ -181,6 +183,7 @@ export async function enqueueEventNotification(args: { text, }), queue: args.queue, + queueDelayMs: EVENT_WAIT_MS, state: args.state, }); } diff --git a/packages/junior/src/chat/task-execution/conversation-turn.ts b/packages/junior/src/chat/task-execution/conversation-turn.ts index 6dc4ea37dd..31e2bc08d9 100644 --- a/packages/junior/src/chat/task-execution/conversation-turn.ts +++ b/packages/junior/src/chat/task-execution/conversation-turn.ts @@ -79,12 +79,12 @@ import { import { joinMailboxText } from "@/chat/task-execution/mailbox-input"; import { resolveConversationDestination } from "@/chat/conversations/destination"; import { + EVENT_WAIT_MS, isEventMailboxMetadata, type EventMailboxMetadata, } from "@/chat/events/notification"; import { isEventConversationMessage } from "@/chat/events/actor"; -const EVENT_WAIT_MS = 30_000; const EVENT_WAIT_PER_EXTRA_MESSAGE_MS = 5_000; const EVENT_MAX_WAIT_MS = 60_000; diff --git a/packages/junior/src/chat/task-execution/store.ts b/packages/junior/src/chat/task-execution/store.ts index b2dc648e11..b923cf5367 100644 --- a/packages/junior/src/chat/task-execution/store.ts +++ b/packages/junior/src/chat/task-execution/store.ts @@ -245,6 +245,7 @@ async function enqueueAfterAppend(args: { conversationStore?: ConversationStore; nowMs?: number; queue: ConversationWorkQueue; + queueDelayMs?: number; state?: StateAdapter; }): Promise { const nowMs = args.nowMs ?? now(); @@ -283,6 +284,7 @@ async function enqueueAfterAppend(args: { const wake = await ensureConversationWake({ conversationId: args.message.conversationId, conversationStore: args.conversationStore, + delayMs: args.queueDelayMs, idempotencyKey, nowMs, queue: args.queue, @@ -309,6 +311,7 @@ export async function appendAndEnqueueInboundMessage(args: { conversationStore?: ConversationStore; nowMs?: number; queue: ConversationWorkQueue; + queueDelayMs?: number; state?: StateAdapter; }): Promise { const nowMs = args.nowMs ?? now(); diff --git a/packages/junior/tests/component/events/events.test.ts b/packages/junior/tests/component/events/events.test.ts index ef0dfcfc53..18612e6afa 100644 --- a/packages/junior/tests/component/events/events.test.ts +++ b/packages/junior/tests/component/events/events.test.ts @@ -134,6 +134,7 @@ describe("event delivery", () => { expect(queue.sentRecords()).toEqual([ { conversationId: CONVERSATION_ID, + delayMs: 30_000, idempotencyKey: `event:${subscription.id}:delivery-1:check-suite-1`, }, ]); @@ -244,10 +245,12 @@ describe("event delivery", () => { expect.arrayContaining([ { conversationId: CONVERSATION_ID, + delayMs: 30_000, idempotencyKey: `event:${threadWatch.id}:delivery-multi:check-suite-1`, }, { conversationId: "agent:deadbeefcafebabe", + delayMs: 30_000, idempotencyKey: `event:${opaqueWatch.id}:delivery-multi:check-suite-1`, }, ]), @@ -360,6 +363,7 @@ describe("event delivery", () => { }, { conversationId: CONVERSATION_ID, + delayMs: 30_000, idempotencyKey: `event:${subscription.id}:github:delivery-bridge:pull_request.comment.created`, }, ]), @@ -540,6 +544,7 @@ describe("event delivery", () => { expect(queue.sentRecords()).toEqual([ { conversationId: CONVERSATION_ID, + delayMs: 30_000, idempotencyKey: `event:${subscription.id}:delivery-3:check-suite-1`, }, ]); diff --git a/packages/junior/tests/integration/event-wake-delay.test.ts b/packages/junior/tests/integration/event-wake-delay.test.ts index 546266024e..1af0ef3ecf 100644 --- a/packages/junior/tests/integration/event-wake-delay.test.ts +++ b/packages/junior/tests/integration/event-wake-delay.test.ts @@ -1,5 +1,8 @@ import { afterEach, describe, expect, it, vi } from "vitest"; -import { createEventInboundMessage } from "@/chat/events/notification"; +import { + createEventInboundMessage, + EVENT_WAIT_MS, +} from "@/chat/events/notification"; import type { AgentRun } from "@/chat/agent/types"; import { getConversationEventStore } from "@/chat/db"; import { @@ -90,8 +93,12 @@ describe("event wake delay", () => { conversationStore, message: eventMessage(1, baseMs), queue, + queueDelayMs: EVENT_WAIT_MS, state, }); + expect(queue.sentRecords()).toEqual([ + expect.objectContaining({ delayMs: 30_000 }), + ]); const agentRuns: AgentRun[] = []; const run = requireConversationTurn( @@ -105,17 +112,6 @@ describe("event wake delay", () => { ), ); - await expect( - processConversationQueueMessage(queue.takeMessage(), { - conversationStore, - queue, - run, - state, - }), - ).resolves.toEqual({ status: "pending_requeued" }); - expect(agentRuns).toHaveLength(0); - expect(queue.sentRecords().at(-1)).toMatchObject({ delayMs: 30_000 }); - nowSpy.mockReturnValue(baseMs + 200); await appendAndEnqueueInboundMessage({ conversationStore, @@ -123,7 +119,7 @@ describe("event wake delay", () => { queue, state, }); - expect(queue.sentRecords()).toHaveLength(2); + expect(queue.sentRecords()).toHaveLength(1); nowSpy.mockReturnValue(baseMs + 30_000); await expect( From f6b1729ddf674f16b5b477e420920522e455e998 Mon Sep 17 00:00:00 2001 From: "sentry-junior[bot]" <264270552+sentry-junior[bot]@users.noreply.github.com> Date: Thu, 10 Sep 2026 19:04:27 +0000 Subject: [PATCH 07/10] fix(events): let user messages bypass event delay --- .../src/chat/conversations/web-input.ts | 1 + .../junior/src/chat/task-execution/store.ts | 4 +++ .../junior/src/chat/task-execution/worker.ts | 10 ++++-- .../integration/event-wake-delay.test.ts | 33 ++++++++++++++++--- 4 files changed, 41 insertions(+), 7 deletions(-) diff --git a/packages/junior/src/chat/conversations/web-input.ts b/packages/junior/src/chat/conversations/web-input.ts index e3e739ca13..84bf1eda77 100644 --- a/packages/junior/src/chat/conversations/web-input.ts +++ b/packages/junior/src/chat/conversations/web-input.ts @@ -288,6 +288,7 @@ export async function appendAndEnqueueWebMessage( conversationStore: options.conversationStore, nowMs, queue: options.queue, + replaceExistingWake: true, state: options.state, }); const status = result.status === "appended" ? "accepted" : result.status; diff --git a/packages/junior/src/chat/task-execution/store.ts b/packages/junior/src/chat/task-execution/store.ts index b923cf5367..58c2d98b50 100644 --- a/packages/junior/src/chat/task-execution/store.ts +++ b/packages/junior/src/chat/task-execution/store.ts @@ -246,6 +246,7 @@ async function enqueueAfterAppend(args: { nowMs?: number; queue: ConversationWorkQueue; queueDelayMs?: number; + replaceExistingWake?: true; state?: StateAdapter; }): Promise { const nowMs = args.nowMs ?? now(); @@ -288,6 +289,7 @@ async function enqueueAfterAppend(args: { idempotencyKey, nowMs, queue: args.queue, + replaceExistingWake: args.replaceExistingWake, state: args.state, }); if (wake.status !== "enqueued") { @@ -312,6 +314,7 @@ export async function appendAndEnqueueInboundMessage(args: { nowMs?: number; queue: ConversationWorkQueue; queueDelayMs?: number; + replaceExistingWake?: true; state?: StateAdapter; }): Promise { const nowMs = args.nowMs ?? now(); @@ -336,6 +339,7 @@ export async function appendAndEnqueueExclusiveInboundMessage(args: { conversationStore?: ConversationStore; nowMs?: number; queue: ConversationWorkQueue; + replaceExistingWake?: true; state?: StateAdapter; }): Promise { const nowMs = args.nowMs ?? now(); diff --git a/packages/junior/src/chat/task-execution/worker.ts b/packages/junior/src/chat/task-execution/worker.ts index 20f78abf54..47f76f8bf3 100644 --- a/packages/junior/src/chat/task-execution/worker.ts +++ b/packages/junior/src/chat/task-execution/worker.ts @@ -120,9 +120,13 @@ function selectAttemptMessages(work: ConversationWorkState): InboundMessage[] { if (interrupts.length > 0) { return selectContiguousTurnBatch(interrupts); } - return work.execution.status === "paused" - ? [] - : selectContiguousTurnBatch(messages); + if (work.execution.status === "paused") return []; + const nonEventMessages = messages.filter( + (message) => message.source !== "event", + ); + return selectContiguousTurnBatch( + nonEventMessages.length > 0 ? nonEventMessages : messages, + ); } function nudgeIdempotencyKey( diff --git a/packages/junior/tests/integration/event-wake-delay.test.ts b/packages/junior/tests/integration/event-wake-delay.test.ts index 1af0ef3ecf..819ca8e983 100644 --- a/packages/junior/tests/integration/event-wake-delay.test.ts +++ b/packages/junior/tests/integration/event-wake-delay.test.ts @@ -6,6 +6,7 @@ import { import type { AgentRun } from "@/chat/agent/types"; import { getConversationEventStore } from "@/chat/db"; import { + appendAndEnqueueWebMessage, createConversationId, recordWebConversationActivity, } from "@/chat/conversations/web-input"; @@ -112,6 +113,27 @@ describe("event wake delay", () => { ), ); + nowSpy.mockReturnValue(baseMs + 100); + await appendAndEnqueueWebMessage( + { + actor, + conversationId, + idempotencyKey: "human-message", + message: "Please check now", + }, + { conversationStore, queue, state }, + ); + expect(queue.sentRecords()).toHaveLength(2); + await expect( + processConversationQueueMessage(queue.takeMessage(), { + conversationStore, + queue, + run, + state, + }), + ).resolves.toEqual({ status: "pending_requeued" }); + expect(agentRuns).toHaveLength(1); + nowSpy.mockReturnValue(baseMs + 200); await appendAndEnqueueInboundMessage({ conversationStore, @@ -119,7 +141,7 @@ describe("event wake delay", () => { queue, state, }); - expect(queue.sentRecords()).toHaveLength(1); + expect(queue.sentRecords()).toHaveLength(3); nowSpy.mockReturnValue(baseMs + 30_000); await expect( @@ -130,7 +152,7 @@ describe("event wake delay", () => { state, }), ).resolves.toEqual({ status: "pending_requeued" }); - expect(agentRuns).toHaveLength(0); + expect(agentRuns).toHaveLength(1); expect(queue.sentRecords().at(-1)).toMatchObject({ delayMs: 5_000 }); nowSpy.mockReturnValue(baseMs + 35_000); @@ -142,7 +164,7 @@ describe("event wake delay", () => { state, }), ).resolves.toEqual({ status: "completed" }); - expect(agentRuns).toHaveLength(1); + expect(agentRuns).toHaveLength(2); const userMessages = ( await getConversationEventStore().loadHistory(conversationId) @@ -151,6 +173,9 @@ describe("event wake delay", () => { ? [event.data.text] : [], ); - expect(userMessages).toEqual(["Check 1 failed\n\nCheck 2 failed"]); + expect(userMessages).toEqual([ + "Please check now", + "Check 1 failed\n\nCheck 2 failed", + ]); }); }); From 8af9a5393b523d06edb9248ce95c7f33db0bfadd Mon Sep 17 00:00:00 2001 From: "sentry-junior[bot]" <264270552+sentry-junior[bot]@users.noreply.github.com> Date: Thu, 10 Sep 2026 19:11:16 +0000 Subject: [PATCH 08/10] fix(events): let all human input bypass event wake delay Cursor Bugbot flagged that Slack (and any other non-web) inbound messages went through appendAndEnqueueInboundMessage without replaceExistingWake, so a pending event enqueue marker could swallow their wake and force them to wait out the event delay. Move the default to the shared enqueueAfterAppend choke point: any message whose source is not "event" now defaults to replaceExistingWake, covering Slack, dispatch, and invocation callers without needing every call site to opt in. --- packages/junior/src/chat/task-execution/store.ts | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/packages/junior/src/chat/task-execution/store.ts b/packages/junior/src/chat/task-execution/store.ts index 58c2d98b50..8a1dfd65b7 100644 --- a/packages/junior/src/chat/task-execution/store.ts +++ b/packages/junior/src/chat/task-execution/store.ts @@ -289,7 +289,11 @@ async function enqueueAfterAppend(args: { idempotencyKey, nowMs, queue: args.queue, - replaceExistingWake: args.replaceExistingWake, + // Human input must not wait out an event's queue delay: only event + // messages default to coalescing with an already-pending wake. + replaceExistingWake: + args.replaceExistingWake ?? + (args.message.source === "event" ? undefined : true), state: args.state, }); if (wake.status !== "enqueued") { From b622d60b3cf178471d970ec5544f309e7214b566 Mon Sep 17 00:00:00 2001 From: "sentry-junior[bot]" <264270552+sentry-junior[bot]@users.noreply.github.com> Date: Thu, 10 Sep 2026 19:20:38 +0000 Subject: [PATCH 09/10] fix(events): preserve active conversation leases Non-event input must bypass an event's delayed enqueue marker, but it must not bypass a worker that already owns the conversation lease. Split those controls so Slack steering stays in the active worker without redundant queue nudges. --- .../junior/src/chat/conversations/web-input.ts | 1 - packages/junior/src/chat/task-execution/store.ts | 16 +++++++++------- 2 files changed, 9 insertions(+), 8 deletions(-) diff --git a/packages/junior/src/chat/conversations/web-input.ts b/packages/junior/src/chat/conversations/web-input.ts index 84bf1eda77..e3e739ca13 100644 --- a/packages/junior/src/chat/conversations/web-input.ts +++ b/packages/junior/src/chat/conversations/web-input.ts @@ -288,7 +288,6 @@ export async function appendAndEnqueueWebMessage( conversationStore: options.conversationStore, nowMs, queue: options.queue, - replaceExistingWake: true, state: options.state, }); const status = result.status === "appended" ? "accepted" : result.status; diff --git a/packages/junior/src/chat/task-execution/store.ts b/packages/junior/src/chat/task-execution/store.ts index 8a1dfd65b7..a8240a7d87 100644 --- a/packages/junior/src/chat/task-execution/store.ts +++ b/packages/junior/src/chat/task-execution/store.ts @@ -169,14 +169,16 @@ export function hasRunnableConversationWork( /** * Ensure runnable conversation work has one accepted queue wake-up nudge. * - * Ordinary wakes coalesce on a recent accepted marker. Replacement is only for - * consumed or known-stale deliveries where another queue nudge must exist. + * Ordinary wakes coalesce on a recent accepted marker. A caller can ignore the + * marker without replacing an active lease. Replacement is only for consumed + * or known-stale deliveries where another queue nudge must exist. */ export async function ensureConversationWake(args: { conversationId: string; conversationStore?: ConversationStore; delayMs?: number; idempotencyKey: string; + ignoreEnqueueMarker?: true; nowMs?: number; queue: ConversationWorkQueue; replaceExistingWake?: true; @@ -198,6 +200,7 @@ export async function ensureConversationWake(args: { return { status: "lease_active" }; } if ( + args.ignoreEnqueueMarker !== true && args.replaceExistingWake !== true && hasRecentEnqueueMarker(conversation, nowMs) ) { @@ -287,13 +290,12 @@ async function enqueueAfterAppend(args: { conversationStore: args.conversationStore, delayMs: args.queueDelayMs, idempotencyKey, + // Human input must not wait out an event's queue delay. An active worker + // still owns the mailbox, so do not replace its wake. + ignoreEnqueueMarker: args.message.source === "event" ? undefined : true, nowMs, queue: args.queue, - // Human input must not wait out an event's queue delay: only event - // messages default to coalescing with an already-pending wake. - replaceExistingWake: - args.replaceExistingWake ?? - (args.message.source === "event" ? undefined : true), + replaceExistingWake: args.replaceExistingWake, state: args.state, }); if (wake.status !== "enqueued") { From f8797a43d55b8f6838d8702353c829cd907a961e Mon Sep 17 00:00:00 2001 From: "sentry-junior[bot]" <264270552+sentry-junior[bot]@users.noreply.github.com> Date: Thu, 10 Sep 2026 19:40:26 +0000 Subject: [PATCH 10/10] docs(events): explain the event debounce deferral Add comments explaining why the mailbox worker defers a batch that starts with an event, how the debounce window grows with batch size, and why deferral re-enqueues instead of sleeping. --- .../src/chat/task-execution/conversation-turn.ts | 11 +++++++++++ 1 file changed, 11 insertions(+) diff --git a/packages/junior/src/chat/task-execution/conversation-turn.ts b/packages/junior/src/chat/task-execution/conversation-turn.ts index 31e2bc08d9..95b099a6fd 100644 --- a/packages/junior/src/chat/task-execution/conversation-turn.ts +++ b/packages/junior/src/chat/task-execution/conversation-turn.ts @@ -85,7 +85,9 @@ import { } from "@/chat/events/notification"; import { isEventConversationMessage } from "@/chat/events/actor"; +/** Extra debounce time added per additional event already batched. */ const EVENT_WAIT_PER_EXTRA_MESSAGE_MS = 5_000; +/** Upper bound on the debounce window regardless of batch size. */ const EVENT_MAX_WAIT_MS = 60_000; function stableHex(...parts: string[]): string { @@ -178,6 +180,15 @@ export function createConversationTurnWorker( context: ConversationWorkerContext, resolved: MailboxTurnWork, ): Promise => { + // A resource-event burst (e.g. several check runs on one PR) should + // produce one Turn instead of one per event. When the batch starts with + // an event, wait past its debounce window before running the Turn so + // later events in the same burst still land in this batch. The window + // grows with batch size (more events waiting means a longer burst) up to + // EVENT_MAX_WAIT_MS. If the window has not elapsed, defer without + // running: the worker re-enqueues this wake with the remaining delay + // (see `ensureConversationWake`'s `delayMs`), so this never sleeps in + // process. if (resolved.kind === "mailbox") { const first = resolved.batch[0]!; if (isEventMailboxMetadata(first.message.input.metadata)) {