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 568885e559..95b099a6fd 100644 --- a/packages/junior/src/chat/task-execution/conversation-turn.ts +++ b/packages/junior/src/chat/task-execution/conversation-turn.ts @@ -79,11 +79,17 @@ 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"; +/** 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 { return createHash("sha256") .update(parts.join("\u0000")) @@ -174,6 +180,28 @@ 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)) { + 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 }; + } + } + const lifecycle = new ConversationTurnLifecycleService( getConversationEventStore(), ); diff --git a/packages/junior/src/chat/task-execution/store.ts b/packages/junior/src/chat/task-execution/store.ts index b2dc648e11..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) ) { @@ -245,6 +248,8 @@ async function enqueueAfterAppend(args: { conversationStore?: ConversationStore; nowMs?: number; queue: ConversationWorkQueue; + queueDelayMs?: number; + replaceExistingWake?: true; state?: StateAdapter; }): Promise { const nowMs = args.nowMs ?? now(); @@ -283,9 +288,14 @@ async function enqueueAfterAppend(args: { const wake = await ensureConversationWake({ conversationId: args.message.conversationId, 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, + replaceExistingWake: args.replaceExistingWake, state: args.state, }); if (wake.status !== "enqueued") { @@ -309,6 +319,8 @@ export async function appendAndEnqueueInboundMessage(args: { conversationStore?: ConversationStore; nowMs?: number; queue: ConversationWorkQueue; + queueDelayMs?: number; + replaceExistingWake?: true; state?: StateAdapter; }): Promise { const nowMs = args.nowMs ?? now(); @@ -333,6 +345,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 03937b019f..47f76f8bf3 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 { @@ -118,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( @@ -734,6 +740,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/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 new file mode 100644 index 0000000000..819ca8e983 --- /dev/null +++ b/packages/junior/tests/integration/event-wake-delay.test.ts @@ -0,0 +1,181 @@ +import { afterEach, describe, expect, it, vi } from "vitest"; +import { + createEventInboundMessage, + EVENT_WAIT_MS, +} from "@/chat/events/notification"; +import type { AgentRun } from "@/chat/agent/types"; +import { getConversationEventStore } from "@/chat/db"; +import { + appendAndEnqueueWebMessage, + 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 delay", () => { + afterEach(async () => { + await closeConversationFixture(); + vi.restoreAllMocks(); + }); + + it("waits for a burst of events before running one Turn", 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, + queueDelayMs: EVENT_WAIT_MS, + state, + }); + expect(queue.sentRecords()).toEqual([ + expect.objectContaining({ delayMs: 30_000 }), + ]); + + const agentRuns: AgentRun[] = []; + const run = requireConversationTurn( + createConversationTurnWorker( + createModelAgentRunnerForRun((agentRun) => { + agentRuns.push(agentRun); + return createModelStream([ + { type: "text", text: "Handled both events." }, + ]); + }), + ), + ); + + 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, + message: eventMessage(2, baseMs + 200), + queue, + state, + }); + expect(queue.sentRecords()).toHaveLength(3); + + nowSpy.mockReturnValue(baseMs + 30_000); + await expect( + processConversationQueueMessage(queue.takeMessage(), { + conversationStore, + queue, + run, + state, + }), + ).resolves.toEqual({ status: "pending_requeued" }); + expect(agentRuns).toHaveLength(1); + 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(2); + + const userMessages = ( + await getConversationEventStore().loadHistory(conversationId) + ).flatMap((event) => + event.data.type === "message" && event.data.role === "user" + ? [event.data.text] + : [], + ); + expect(userMessages).toEqual([ + "Please check now", + "Check 1 failed\n\nCheck 2 failed", + ]); + }); +});