Skip to content

Commit c246229

Browse files
committed
fix(core): harden event task updates
1 parent dd19640 commit c246229

6 files changed

Lines changed: 214 additions & 42 deletions

File tree

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

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -14,9 +14,9 @@ import {
1414
import type { EventTask, EventTaskMatch } from "@/chat/event-tasks/types";
1515
import type { ConversationWorkQueue } from "@/chat/task-execution/queue";
1616

17-
function matchId(taskId: string, eventKey: string): string {
17+
function matchId(taskId: string, provider: string, eventKey: string): string {
1818
return `evmatch_${createHash("sha256")
19-
.update(`${taskId}\0${eventKey}`)
19+
.update(`${taskId}\0${provider}\0${eventKey}`)
2020
.digest("hex")
2121
.slice(0, 32)}`;
2222
}
@@ -68,7 +68,7 @@ export async function ingestEventTasks(
6868
const errors: unknown[] = [];
6969
for (const task of tasks) {
7070
try {
71-
const id = matchId(task.id, event.eventKey);
71+
const id = matchId(task.id, event.provider, event.eventKey);
7272
const pending: EventTaskMatch = {
7373
id,
7474
createdAtMs: nowMs,

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

Lines changed: 40 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,7 @@ import {
44
slackActorSchema,
55
type ResourceEvent,
66
} from "@sentry/junior-plugin-api";
7-
import { and, asc, eq, ne } from "drizzle-orm";
7+
import { and, asc, eq, ne, sql } from "drizzle-orm";
88
import { z } from "zod";
99
import type { JuniorDatabase } from "@/db/db";
1010
import {
@@ -105,21 +105,55 @@ export async function createEventTask(
105105
return (await getEventTask(db, parsed.id)) ?? parsed;
106106
}
107107

108-
/** Replace one existing event task. */
109-
export async function saveEventTask(
108+
/** Replace an active event task only when the caller read its current version. */
109+
export async function updateActiveEventTask(
110110
db: JuniorDatabase,
111111
task: EventTask,
112-
): Promise<void> {
112+
expectedUpdatedAtMs: number,
113+
): Promise<EventTask | undefined> {
113114
const parsed = eventTaskSchema.parse(task);
114-
await db
115+
if (parsed.status !== "active") {
116+
throw new Error("Event task update requires active status");
117+
}
118+
const rows = await db
115119
.update(juniorEventTasks)
116120
.set({
117121
provider: parsed.trigger.provider,
118122
resourceRef: parsed.trigger.resourceRef,
119123
status: parsed.status,
120124
task: parsed,
121125
})
122-
.where(eq(juniorEventTasks.id, parsed.id));
126+
.where(
127+
and(
128+
eq(juniorEventTasks.id, parsed.id),
129+
eq(juniorEventTasks.status, "active"),
130+
sql`${juniorEventTasks.task}->>'updatedAtMs' = ${String(expectedUpdatedAtMs)}`,
131+
),
132+
)
133+
.returning({ task: juniorEventTasks.task });
134+
return rows[0] ? parseTask(rows[0].task) : undefined;
135+
}
136+
137+
/** Delete an active event task without allowing a stale update to revive it. */
138+
export async function deleteActiveEventTask(
139+
db: JuniorDatabase,
140+
task: EventTask,
141+
): Promise<EventTask | undefined> {
142+
const parsed = eventTaskSchema.parse(task);
143+
if (parsed.status !== "deleted") {
144+
throw new Error("Event task deletion requires deleted status");
145+
}
146+
const rows = await db
147+
.update(juniorEventTasks)
148+
.set({ status: parsed.status, task: parsed })
149+
.where(
150+
and(
151+
eq(juniorEventTasks.id, parsed.id),
152+
eq(juniorEventTasks.status, "active"),
153+
),
154+
)
155+
.returning({ task: juniorEventTasks.task });
156+
return rows[0] ? parseTask(rows[0].task) : undefined;
123157
}
124158

125159
/** List active event tasks in one Slack workspace. */

packages/junior/src/chat/prompt.ts

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -326,7 +326,7 @@ const TOOL_POLICY_RULES = [
326326
"- If a sandbox-backed tool reports that sandbox execution is unavailable, treat that as a blocker for local file/shell inspection; do not pretend host files were inspected.",
327327
"- For user-provided URLs, use `webFetch`; for discovery, use `webSearch` then fetch/read promising sources; for current time/date context, use `systemTime`.",
328328
"- When a tool result includes a subscribable resource, use resource-event subscriptions for high-signal provider changes that serve the user's current intent; do not create scheduled polling tasks for events the subscription can deliver. Use the suggested events when they fit and write a concise intent summary.",
329-
"- Use event tasks when the user wants a durable instruction to dispatch whenever registered resource events occur. When the resource and events are known, create a directly requested event task without redundant confirmation.",
329+
"- Use event tasks only when the user explicitly asks for an event task or durable automated work whenever registered resource events occur. Ordinary requests to watch a resource and report relevant updates use resource-event subscriptions. When an event task's resource and events are known, create it without redundant confirmation.",
330330
"- Event tasks use system credentials by default. If the user denies creator credential use, use system mode without asking. Conditional or tentative credential language is not authorization; when the task depends on creator credentials, ask before creating or updating it and enable creator mode only after explicit authorization for future use.",
331331
"- For code changes, debugging or root-cause analysis, broad refactors, and software architecture decisions, use `handoff` before substantive analysis only when it offers a profile that better matches the task. Do not switch merely because the task involves code.",
332332
"- Run `jr-rpc config get|set|unset|list` for provider defaults and `jr-rpc plugins list` for installed plugin introspection as standalone bash commands; do not chain them with `cd`, `&&`, pipes, or provider commands.",

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

Lines changed: 44 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -3,9 +3,10 @@ import { z } from "zod";
33
import { getDb } from "@/chat/db";
44
import {
55
createEventTask,
6+
deleteActiveEventTask,
67
getEventTask,
78
listEventTasksForTeam,
8-
saveEventTask,
9+
updateActiveEventTask,
910
} from "@/chat/event-tasks/store";
1011
import type {
1112
EventTask,
@@ -22,11 +23,11 @@ const MAX_LISTED_TASKS = 50;
2223

2324
const triggerSchema = z
2425
.object({
25-
provider: z.string().min(1),
26-
resourceRef: z.string().min(1),
27-
resourceType: z.string().min(1),
28-
label: z.string().min(1),
29-
events: z.array(z.string().min(1)).min(1),
26+
provider: z.string().trim().min(1),
27+
resourceRef: z.string().trim().min(1),
28+
resourceType: z.string().trim().min(1),
29+
label: z.string().trim().min(1),
30+
events: z.array(z.string().trim().min(1)).min(1),
3031
})
3132
.strict();
3233

@@ -139,6 +140,10 @@ function cleanEvents(events: string[]): string[] {
139140
return clean;
140141
}
141142

143+
function nextUpdatedAtMs(task: EventTask): number {
144+
return Math.max(Date.now(), task.updatedAtMs + 1);
145+
}
146+
142147
function compactTask(task: EventTask) {
143148
return {
144149
id: task.id,
@@ -174,10 +179,10 @@ export function createEventTaskTool(context: ToolRuntimeContext) {
174179
},
175180
executionMode: "sequential",
176181
description:
177-
"Create a durable Junior task when the user asks for work on registered events from a subscribable provider resource. Copy the provider, resource ref, type, label, and supported event names from the provider tool result; do not ask for redundant confirmation when those details are known.",
182+
"Create a durable reactive Junior task only when the user explicitly asks for an event task or automated work whenever registered provider events occur. Ordinary requests to watch a resource and report relevant updates use resource-event subscriptions instead. Copy the provider, resource ref, type, label, and supported event names from the provider tool result; do not ask for redundant confirmation when those details are known.",
178183
inputSchema: z
179184
.object({
180-
task: z.string().min(1).max(4000),
185+
task: z.string().trim().min(1).max(4000),
181186
trigger: triggerSchema,
182187
credentialMode: z
183188
.enum(["system", "creator"])
@@ -221,12 +226,12 @@ export function createEventTaskTool(context: ToolRuntimeContext) {
221226
destination,
222227
originalRequest: context.userText,
223228
status: "active",
224-
task: { text: input.task.trim() },
229+
task: { text: input.task },
225230
trigger: {
226-
provider: input.trigger.provider.trim(),
227-
resourceRef: input.trigger.resourceRef.trim(),
228-
resourceType: input.trigger.resourceType.trim(),
229-
label: input.trigger.label.trim(),
231+
provider: input.trigger.provider,
232+
resourceRef: input.trigger.resourceRef,
233+
resourceType: input.trigger.resourceType,
234+
label: input.trigger.label,
230235
events: cleanEvents(input.trigger.events),
231236
},
232237
updatedAtMs: nowMs,
@@ -283,7 +288,7 @@ export function createUpdateEventTaskTool(context: ToolRuntimeContext) {
283288
inputSchema: z
284289
.object({
285290
taskId: z.string().min(1),
286-
task: z.string().min(1).max(4000).optional(),
291+
task: z.string().trim().min(1).max(4000).optional(),
287292
trigger: triggerSchema.optional(),
288293
credentialMode: z
289294
.enum(["system", "creator"])
@@ -319,20 +324,29 @@ export function createUpdateEventTaskTool(context: ToolRuntimeContext) {
319324
changesExecution && !isCreator
320325
? "system"
321326
: (input.credentialMode ?? current.credentialMode),
322-
task: input.task ? { text: input.task.trim() } : current.task,
327+
task: input.task ? { text: input.task } : current.task,
323328
trigger: input.trigger
324329
? {
325-
provider: input.trigger.provider.trim(),
326-
resourceRef: input.trigger.resourceRef.trim(),
327-
resourceType: input.trigger.resourceType.trim(),
328-
label: input.trigger.label.trim(),
330+
provider: input.trigger.provider,
331+
resourceRef: input.trigger.resourceRef,
332+
resourceType: input.trigger.resourceType,
333+
label: input.trigger.label,
329334
events: cleanEvents(input.trigger.events),
330335
}
331336
: current.trigger,
332-
updatedAtMs: Date.now(),
337+
updatedAtMs: nextUpdatedAtMs(current),
333338
};
334-
await saveEventTask(getDb(), next);
335-
return success(next);
339+
const saved = await updateActiveEventTask(
340+
getDb(),
341+
next,
342+
current.updatedAtMs,
343+
);
344+
if (!saved) {
345+
throw new ToolInputError(
346+
"Event task changed while it was being updated. List event tasks and try again.",
347+
);
348+
}
349+
return success(saved);
336350
},
337351
});
338352
}
@@ -355,10 +369,15 @@ export function createDeleteEventTaskTool(context: ToolRuntimeContext) {
355369
const next: EventTask = {
356370
...current,
357371
status: "deleted",
358-
updatedAtMs: Date.now(),
372+
updatedAtMs: nextUpdatedAtMs(current),
359373
};
360-
await saveEventTask(getDb(), next);
361-
return success(next);
374+
const deleted = await deleteActiveEventTask(getDb(), next);
375+
if (!deleted) {
376+
throw new ToolInputError(
377+
"Event task was not found in the active Slack conversation.",
378+
);
379+
}
380+
return success(deleted);
362381
},
363382
});
364383
}

packages/junior/tests/component/resource-events/resource-events.test.ts

Lines changed: 68 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -3,8 +3,11 @@ import { createMemoryState } from "@chat-adapter/state-memory";
33
import { githubPlugin } from "@sentry/junior-github";
44
import type { StateAdapter } from "chat";
55
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
6+
import { eq } from "drizzle-orm";
67
import { createApp, defineJuniorPlugins } from "@/app";
7-
import { closeDb } from "@/chat/db";
8+
import { getDispatchRecord } from "@/chat/agent-dispatch/store";
9+
import { closeDb, getDb } from "@/chat/db";
10+
import { createEventTask } from "@/chat/event-tasks/store";
811
import {
912
getConfigDefaults,
1013
setConfigDefaults,
@@ -15,6 +18,10 @@ import { setDashboardConversationLinkOptions } from "@/chat/slack/dashboard-link
1518
import { disconnectStateAdapter, getStateAdapter } from "@/chat/state/adapter";
1619
import { JUNIOR_THREAD_STATE_TTL_MS } from "@/chat/state/ttl";
1720
import { getConversationWorkState } from "@/chat/task-execution/store";
21+
import {
22+
juniorEventTaskMatches,
23+
juniorEventTasks,
24+
} from "@/db/schema/event-tasks";
1825
import { ingestResourceEvent } from "@/chat/resource-events/ingest";
1926
import {
2027
cancelResourceEventSubscription,
@@ -153,6 +160,7 @@ describe("resource event subscriptions", () => {
153160
const previousDashboardOptions =
154161
setDashboardConversationLinkOptions(undefined);
155162
setDashboardConversationLinkOptions(previousDashboardOptions);
163+
let eventTaskId: string | undefined;
156164

157165
try {
158166
process.env.GITHUB_WEBHOOK_SECRET = "test-secret";
@@ -166,6 +174,25 @@ describe("resource event subscriptions", () => {
166174
nowMs,
167175
state,
168176
});
177+
eventTaskId = `evt_webhook_bridge_${nowMs}`;
178+
await createEventTask(getDb(), {
179+
id: eventTaskId,
180+
conversationAccess: { audience: "channel", visibility: "public" },
181+
createdAtMs: nowMs,
182+
createdBy: { slackUserId: "U123" },
183+
credentialMode: "system",
184+
destination: SLACK_DESTINATION,
185+
status: "active",
186+
task: { text: "Summarize the reviewer comment." },
187+
trigger: {
188+
events: ["comment.created"],
189+
label: "GitHub PR getsentry/junior#691",
190+
provider: "github",
191+
resourceRef: "github:pull_request:getsentry/junior#691",
192+
resourceType: "pull_request",
193+
},
194+
updatedAtMs: nowMs,
195+
});
169196
const body = JSON.stringify({
170197
action: "created",
171198
repository: { full_name: "getsentry/junior" },
@@ -205,12 +232,39 @@ describe("resource event subscriptions", () => {
205232
);
206233

207234
expect(response.status).toBe(202);
208-
expect(queue.sentRecords()).toEqual([
209-
{
210-
conversationId: CONVERSATION_ID,
211-
idempotencyKey: `resource-event:${subscription.id}:github:delivery-bridge:comment.created`,
235+
expect(queue.sentRecords()).toEqual(
236+
expect.arrayContaining([
237+
{
238+
conversationId: expect.stringMatching(/^agent-dispatch:/),
239+
idempotencyKey: expect.stringMatching(/^agent-dispatch:/),
240+
},
241+
{
242+
conversationId: CONVERSATION_ID,
243+
idempotencyKey: `resource-event:${subscription.id}:github:delivery-bridge:comment.created`,
244+
},
245+
]),
246+
);
247+
const dispatchRecord = queue
248+
.sentRecords()
249+
.find(({ conversationId }) =>
250+
conversationId.startsWith("agent-dispatch:"),
251+
);
252+
expect(dispatchRecord).toBeDefined();
253+
const dispatch = dispatchRecord
254+
? await getDispatchRecord(
255+
dispatchRecord.conversationId.replace(/^agent-dispatch:/, ""),
256+
)
257+
: undefined;
258+
expect(dispatch).toMatchObject({
259+
input: expect.stringContaining("Summarize the reviewer comment."),
260+
metadata: {
261+
eventType: "comment.created",
262+
provider: "github",
263+
resourceRef: "github:pull_request:getsentry/junior#691",
264+
taskId: eventTaskId,
212265
},
213-
]);
266+
plugin: "junior",
267+
});
214268
const work = await getConversationWorkState({
215269
conversationId: CONVERSATION_ID,
216270
state,
@@ -219,6 +273,14 @@ describe("resource event subscriptions", () => {
219273
"please add regression coverage",
220274
);
221275
} finally {
276+
if (eventTaskId) {
277+
await getDb()
278+
.delete(juniorEventTaskMatches)
279+
.where(eq(juniorEventTaskMatches.taskId, eventTaskId));
280+
await getDb()
281+
.delete(juniorEventTasks)
282+
.where(eq(juniorEventTasks.id, eventTaskId));
283+
}
222284
setPlugins(previousPlugins);
223285
pluginCatalogRuntime.setConfig(previousPluginCatalogConfig);
224286
setConfigDefaults(previousConfigDefaults);

0 commit comments

Comments
 (0)