Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions packages/junior/src/chat/events/notification.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -181,6 +183,7 @@ export async function enqueueEventNotification(args: {
text,
}),
queue: args.queue,
queueDelayMs: EVENT_WAIT_MS,
state: args.state,
});
}
28 changes: 28 additions & 0 deletions packages/junior/src/chat/task-execution/conversation-turn.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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"))
Expand Down Expand Up @@ -174,6 +180,28 @@ export function createConversationTurnWorker(
context: ConversationWorkerContext,
resolved: MailboxTurnWork,
): Promise<ConversationWorkerResult> => {
// 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") {
Comment thread
sentry-junior[bot] marked this conversation as resolved.
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 };
}
Comment thread
sentry-junior[bot] marked this conversation as resolved.
}

const lifecycle = new ConversationTurnLifecycleService(
getConversationEventStore(),
);
Expand Down
17 changes: 15 additions & 2 deletions packages/junior/src/chat/task-execution/store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -198,6 +200,7 @@ export async function ensureConversationWake(args: {
return { status: "lease_active" };
}
if (
args.ignoreEnqueueMarker !== true &&
args.replaceExistingWake !== true &&
hasRecentEnqueueMarker(conversation, nowMs)
) {
Expand Down Expand Up @@ -245,6 +248,8 @@ async function enqueueAfterAppend(args: {
conversationStore?: ConversationStore;
nowMs?: number;
queue: ConversationWorkQueue;
queueDelayMs?: number;
replaceExistingWake?: true;
state?: StateAdapter;
}): Promise<AppendAndEnqueueExclusiveInboundMessageResult> {
const nowMs = args.nowMs ?? now();
Expand Down Expand Up @@ -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") {
Expand All @@ -309,6 +319,8 @@ export async function appendAndEnqueueInboundMessage(args: {
conversationStore?: ConversationStore;
nowMs?: number;
queue: ConversationWorkQueue;
queueDelayMs?: number;
replaceExistingWake?: true;
state?: StateAdapter;
}): Promise<AppendAndEnqueueInboundMessageResult> {
const nowMs = args.nowMs ?? now();
Expand All @@ -333,6 +345,7 @@ export async function appendAndEnqueueExclusiveInboundMessage(args: {
conversationStore?: ConversationStore;
nowMs?: number;
queue: ConversationWorkQueue;
replaceExistingWake?: true;
state?: StateAdapter;
}): Promise<AppendAndEnqueueExclusiveInboundMessageResult> {
const nowMs = args.nowMs ?? now();
Expand Down
13 changes: 10 additions & 3 deletions packages/junior/src/chat/task-execution/worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -734,6 +740,7 @@ async function processConversationWorkInContext(
const wake = await ensureConversationWake({
conversationId,
conversationStore: options.conversationStore,
delayMs: result.delayMs,
idempotencyKey: nudgeIdempotencyKey(
"deferred",
conversationId,
Expand Down
5 changes: 5 additions & 0 deletions packages/junior/tests/component/events/events.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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`,
},
]);
Expand Down Expand Up @@ -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`,
},
]),
Expand Down Expand Up @@ -360,6 +363,7 @@ describe("event delivery", () => {
},
{
conversationId: CONVERSATION_ID,
delayMs: 30_000,
idempotencyKey: `event:${subscription.id}:github:delivery-bridge:pull_request.comment.created`,
},
]),
Expand Down Expand Up @@ -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`,
},
]);
Expand Down
181 changes: 181 additions & 0 deletions packages/junior/tests/integration/event-wake-delay.test.ts
Original file line number Diff line number Diff line change
@@ -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<ConversationWorkerResult>,
) {
return async (
context: ConversationWorkerContext,
): Promise<ConversationWorkerResult> => {
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",
]);
});
});
Loading