Skip to content

Commit b65d31a

Browse files
committed
fix(core): close event task races
1 parent 8ac51b3 commit b65d31a

4 files changed

Lines changed: 73 additions & 21 deletions

File tree

packages/junior/src/chat/event-tasks/ingest.ts

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -5,9 +5,9 @@ import {
55
type ResourceEvent,
66
} from "@sentry/junior-plugin-api";
77
import { dispatchTask } from "@/chat/agent-dispatch/context";
8-
import { getDb } from "@/chat/db";
8+
import { getDb, getSqlExecutor } from "@/chat/db";
99
import {
10-
createOrGetEventTaskMatch,
10+
claimActiveEventTaskMatch,
1111
findMatchingEventTasks,
1212
markEventTaskMatchDispatched,
1313
} from "@/chat/event-tasks/store";
@@ -77,7 +77,8 @@ export async function ingestEventTasks(
7777
status: "pending",
7878
taskId: task.id,
7979
};
80-
const match = await createOrGetEventTaskMatch(db, pending);
80+
const match = await claimActiveEventTaskMatch(getSqlExecutor(), pending);
81+
if (!match) continue;
8182
if (match.status === "dispatched") continue;
8283
const credentialSubject =
8384
task.credentialMode === "creator"

packages/junior/src/chat/event-tasks/store.ts

Lines changed: 31 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ import {
66
} from "@sentry/junior-plugin-api";
77
import { and, asc, eq, ne, sql } from "drizzle-orm";
88
import { z } from "zod";
9-
import type { JuniorDatabase } from "@/db/db";
9+
import type { JuniorDatabase, JuniorSqlDatabase } from "@/db/db";
1010
import {
1111
juniorEventTaskMatches,
1212
juniorEventTasks,
@@ -211,23 +211,38 @@ async function getEventTaskMatch(
211211
return rows[0] ? parseMatch(rows[0].match) : undefined;
212212
}
213213

214-
/** Create a deduplication record or return the existing match. */
215-
export async function createOrGetEventTaskMatch(
216-
db: JuniorDatabase,
214+
/** Claim a deduplication record while its event task remains active. */
215+
export async function claimActiveEventTaskMatch(
216+
sqlDb: JuniorSqlDatabase,
217217
match: EventTaskMatch,
218-
): Promise<EventTaskMatch> {
218+
): Promise<EventTaskMatch | undefined> {
219219
const parsed = eventTaskMatchSchema.parse(match);
220-
await db
221-
.insert(juniorEventTaskMatches)
222-
.values({
223-
id: parsed.id,
224-
taskId: parsed.taskId,
225-
status: parsed.status,
226-
createdAtMs: parsed.createdAtMs,
227-
match: parsed,
228-
})
229-
.onConflictDoNothing();
230-
return (await getEventTaskMatch(db, parsed.id)) ?? parsed;
220+
return await sqlDb.transaction(async () => {
221+
const db = sqlDb.db();
222+
const active = await db
223+
.select({ id: juniorEventTasks.id })
224+
.from(juniorEventTasks)
225+
.where(
226+
and(
227+
eq(juniorEventTasks.id, parsed.taskId),
228+
eq(juniorEventTasks.status, "active"),
229+
),
230+
)
231+
.for("update")
232+
.limit(1);
233+
if (!active[0]) return undefined;
234+
await db
235+
.insert(juniorEventTaskMatches)
236+
.values({
237+
id: parsed.id,
238+
taskId: parsed.taskId,
239+
status: parsed.status,
240+
createdAtMs: parsed.createdAtMs,
241+
match: parsed,
242+
})
243+
.onConflictDoNothing();
244+
return (await getEventTaskMatch(db, parsed.id)) ?? parsed;
245+
});
231246
}
232247

233248
/** Mark a pending event task match after durable agent dispatch succeeds. */

packages/junior/src/chat/tools/event-tasks.ts

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -312,12 +312,12 @@ export function createUpdateEventTaskTool(context: ToolRuntimeContext) {
312312
if (
313313
input.task === undefined &&
314314
input.trigger === undefined &&
315-
input.credentialMode === undefined
315+
input.credentialMode == null
316316
) {
317317
throw new ToolInputError("Event task update requires a change.");
318318
}
319319
const changesExecution =
320-
input.task !== undefined || input.trigger !== undefined;
320+
input.task !== undefined && input.task !== current.task.text;
321321
const next: EventTask = {
322322
...current,
323323
credentialMode:

packages/junior/tests/integration/event-tasks.test.ts

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ import { getDispatchRecord } from "@/chat/agent-dispatch/store";
55
import { migrateSchema } from "@/chat/conversations/sql/migrations";
66
import { ingestEventTasks } from "@/chat/event-tasks/ingest";
77
import {
8+
claimActiveEventTaskMatch,
89
deleteActiveEventTask,
910
getEventTask,
1011
updateActiveEventTask,
@@ -178,6 +179,31 @@ describe("event tasks", () => {
178179
};
179180
expect(listed.tasks).toHaveLength(2);
180181

182+
await execute(createUpdateEventTaskTool(context("U999")), {
183+
taskId: first.task.id,
184+
trigger: {
185+
provider: "github",
186+
resourceRef: "github:pull_request:getsentry/junior#1174",
187+
resourceType: "pull_request",
188+
label: "Updated GitHub PR label",
189+
events: ["review.changes_requested"],
190+
},
191+
});
192+
await execute(createUpdateEventTaskTool(context("U999")), {
193+
taskId: first.task.id,
194+
task: "Address the requested changes.",
195+
});
196+
expect(await getEventTask(fixture.sql.db(), first.task.id)).toMatchObject({
197+
credentialMode: "creator",
198+
trigger: { label: "Updated GitHub PR label" },
199+
});
200+
await expect(
201+
execute(createUpdateEventTaskTool(context("U999")), {
202+
taskId: first.task.id,
203+
credentialMode: null,
204+
}),
205+
).rejects.toThrow("Event task update requires a change.");
206+
181207
await execute(createUpdateEventTaskTool(context("U999")), {
182208
taskId: first.task.id,
183209
task: "Only summarize the requested changes.",
@@ -248,5 +274,15 @@ describe("event tasks", () => {
248274
status: "deleted",
249275
task: { text: "Post the first summary." },
250276
});
277+
await expect(
278+
claimActiveEventTaskMatch(fixture.sql, {
279+
id: "evmatch_stale_deleted_task",
280+
createdAtMs: deleted.updatedAtMs + 1,
281+
eventKey: "github:delivery-after-delete",
282+
eventType: "review.changes_requested",
283+
status: "pending",
284+
taskId: created.task.id,
285+
}),
286+
).resolves.toBeUndefined();
251287
});
252288
});

0 commit comments

Comments
 (0)