Skip to content

Commit 2b1dd89

Browse files
carderneTrigger.dev RepoOps
authored andcommitted
fix(webapp): authorize batch parent runs before attaching waitpoints
Mono-RevId: c69d1cf0ad5e5d77b4c00aa3c5c59f6542dea4ce
1 parent fb3b26d commit 2b1dd89

5 files changed

Lines changed: 150 additions & 10 deletions

File tree

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
---
2+
area: webapp
3+
type: fix
4+
---
5+
6+
Reject batch-and-wait requests when the parent run belongs to another environment

‎apps/webapp/app/runEngine/services/batchTrigger.server.ts‎

Lines changed: 17 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,7 @@ import {
77
parsePacket,
88
TaskRunErrorCodes,
99
} from "@trigger.dev/core/v3";
10-
import { BatchId, RunId } from "@trigger.dev/core/v3/isomorphic";
10+
import { BatchId } from "@trigger.dev/core/v3/isomorphic";
1111
import { type BatchTaskRun, Prisma } from "@trigger.dev/database";
1212
import { Evt } from "evt";
1313
import { z } from "zod";
@@ -28,6 +28,7 @@ import type { RunEngine } from "../../v3/runEngine.server";
2828
import { ServiceValidationError, WithRunEngine } from "../../v3/services/baseService.server";
2929
import { TriggerTaskService } from "../../v3/services/triggerTask.server";
3030
import { startActiveSpan } from "../../v3/tracer.server";
31+
import { resolveBatchParentRun } from "./resolveBatchParentRun.server";
3132
import { TriggerFailedTaskService } from "./triggerFailedTask.server";
3233

3334
const PROCESSING_BATCH_SIZE = 50;
@@ -90,6 +91,13 @@ export class RunEngineBatchTriggerService extends WithRunEngine {
9091
"call()",
9192
environment,
9293
async (span) => {
94+
const parentRunInternalId = await resolveBatchParentRun({
95+
runStore: this._engine.runStore,
96+
environmentId: environment.id,
97+
parentRunId: body.parentRunId,
98+
resumeParentOnCompletion: body.resumeParentOnCompletion,
99+
});
100+
93101
const { friendlyId } = await mintBatchFriendlyId({
94102
environment: {
95103
organizationId: environment.organizationId,
@@ -112,7 +120,8 @@ export class RunEngineBatchTriggerService extends WithRunEngine {
112120
payloadPacket,
113121
environment,
114122
body,
115-
options
123+
options,
124+
parentRunInternalId
116125
);
117126

118127
if (!batch) {
@@ -165,7 +174,8 @@ export class RunEngineBatchTriggerService extends WithRunEngine {
165174
payloadPacket: IOPacket,
166175
environment: AuthenticatedEnvironment,
167176
body: BatchTriggerTaskV2RequestBody,
168-
options: BatchTriggerTaskServiceOptions = {}
177+
options: BatchTriggerTaskServiceOptions = {},
178+
parentRunInternalId?: string
169179
) {
170180
// BatchTaskRun.runtimeEnvironmentId no longer has an FK into RuntimeEnvironment;
171181
// validate env existence app-side (covers both create arms below).
@@ -187,9 +197,9 @@ export class RunEngineBatchTriggerService extends WithRunEngine {
187197

188198
this.onBatchTaskRunCreated.post(batch);
189199

190-
if (body.parentRunId && body.resumeParentOnCompletion) {
200+
if (parentRunInternalId) {
191201
await this._engine.blockRunWithCreatedBatch({
192-
runId: RunId.fromFriendlyId(body.parentRunId),
202+
runId: parentRunInternalId,
193203
batchId: batch.id,
194204
environmentId: environment.id,
195205
projectId: environment.projectId,
@@ -278,9 +288,9 @@ export class RunEngineBatchTriggerService extends WithRunEngine {
278288

279289
this.onBatchTaskRunCreated.post(batch);
280290

281-
if (body.parentRunId && body.resumeParentOnCompletion) {
291+
if (parentRunInternalId) {
282292
await this._engine.blockRunWithCreatedBatch({
283-
runId: RunId.fromFriendlyId(body.parentRunId),
293+
runId: parentRunInternalId,
284294
batchId: batch.id,
285295
environmentId: environment.id,
286296
projectId: environment.projectId,

‎apps/webapp/app/runEngine/services/createBatch.server.ts‎

Lines changed: 10 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,5 @@
11
import type { InitializeBatchOptions } from "@internal/run-engine";
22
import { type CreateBatchRequestBody, type CreateBatchResponse } from "@trigger.dev/core/v3";
3-
import { RunId } from "@trigger.dev/core/v3/isomorphic";
43
import { type BatchTaskRun, Prisma } from "@trigger.dev/database";
54
import { Evt } from "evt";
65
import { prisma, type PrismaClientOrTransaction } from "~/db.server";
@@ -14,6 +13,7 @@ import { ServiceValidationError, WithRunEngine } from "../../v3/services/baseSer
1413
import { BatchRateLimitExceededError, getBatchLimits } from "../concerns/batchLimits.server";
1514
import { DefaultQueueManager } from "../concerns/queues.server";
1615
import { DefaultTriggerTaskValidator } from "../validators/triggerTaskValidator";
16+
import { resolveBatchParentRun } from "./resolveBatchParentRun.server";
1717

1818
export type CreateBatchServiceOptions = {
1919
triggerVersion?: string;
@@ -101,6 +101,13 @@ export class CreateBatchService extends WithRunEngine {
101101
// Note: Queue size limits are validated per-queue when batch items are processed,
102102
// since we don't know which queues items will go to until they're streamed.
103103

104+
const parentRunInternalId = await resolveBatchParentRun({
105+
runStore: this._engine.runStore,
106+
environmentId: environment.id,
107+
parentRunId: body.parentRunId,
108+
resumeParentOnCompletion: body.resumeParentOnCompletion,
109+
});
110+
104111
// BatchTaskRun.runtimeEnvironmentId no longer has an FK into RuntimeEnvironment;
105112
// validate env existence app-side (passthrough when split is off).
106113
await controlPlaneResolver.assertEnvExists(environment.id);
@@ -125,14 +132,14 @@ export class CreateBatchService extends WithRunEngine {
125132
await batchStreamGrants.mint(environment.id, friendlyId);
126133

127134
// Block parent run if this is a batchTriggerAndWait
128-
if (body.parentRunId && body.resumeParentOnCompletion) {
135+
if (parentRunInternalId) {
129136
await this._engine.scheduleExpireBatch({
130137
batchId: batch.id,
131138
availableAt: new Date(Date.now() + env.BATCH_SEAL_TIMEOUT_MS),
132139
});
133140

134141
await this._engine.blockRunWithCreatedBatch({
135-
runId: RunId.fromFriendlyId(body.parentRunId),
142+
runId: parentRunInternalId,
136143
batchId: batch.id,
137144
environmentId: environment.id,
138145
projectId: environment.projectId,
Lines changed: 83 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,83 @@
1+
import type { RunStore } from "@internal/run-store";
2+
import { describe, expect, it, vi } from "vitest";
3+
import type { ServiceValidationError } from "../../v3/services/baseService.server";
4+
import { resolveBatchParentRun } from "./resolveBatchParentRun.server";
5+
6+
type BatchParentRunStore = Pick<RunStore, "findRun" | "findRunOnPrimary">;
7+
8+
function createRunStore({ replicaRun, primaryRun }: { replicaRun?: object; primaryRun?: object }) {
9+
return {
10+
findRun: vi.fn().mockResolvedValue(replicaRun ?? null),
11+
findRunOnPrimary: vi.fn().mockResolvedValue(primaryRun ?? null),
12+
} as unknown as BatchParentRunStore;
13+
}
14+
15+
describe("resolveBatchParentRun", () => {
16+
it("does not resolve an unused parent run", async () => {
17+
const runStore = createRunStore({});
18+
19+
await expect(
20+
resolveBatchParentRun({
21+
runStore,
22+
environmentId: "env_1",
23+
parentRunId: "run_parent",
24+
resumeParentOnCompletion: false,
25+
})
26+
).resolves.toBeUndefined();
27+
28+
expect(runStore.findRun).not.toHaveBeenCalled();
29+
});
30+
31+
it("returns a parent run from the calling environment", async () => {
32+
const runStore = createRunStore({ replicaRun: { id: "parent" } });
33+
34+
await expect(
35+
resolveBatchParentRun({
36+
runStore,
37+
environmentId: "env_1",
38+
parentRunId: "run_parent",
39+
resumeParentOnCompletion: true,
40+
})
41+
).resolves.toBe("parent");
42+
43+
expect(runStore.findRun).toHaveBeenCalledWith(
44+
{ id: "parent", runtimeEnvironmentId: "env_1" },
45+
{ select: { id: true } }
46+
);
47+
expect(runStore.findRunOnPrimary).not.toHaveBeenCalled();
48+
});
49+
50+
it("falls back to the owning primary", async () => {
51+
const runStore = createRunStore({ primaryRun: { id: "parent" } });
52+
53+
await expect(
54+
resolveBatchParentRun({
55+
runStore,
56+
environmentId: "env_1",
57+
parentRunId: "run_parent",
58+
resumeParentOnCompletion: true,
59+
})
60+
).resolves.toBe("parent");
61+
62+
expect(runStore.findRunOnPrimary).toHaveBeenCalledWith(
63+
{ id: "parent", runtimeEnvironmentId: "env_1" },
64+
{ select: { id: true } }
65+
);
66+
});
67+
68+
it("rejects a parent outside the calling environment", async () => {
69+
const runStore = createRunStore({});
70+
71+
const result = resolveBatchParentRun({
72+
runStore,
73+
environmentId: "env_1",
74+
parentRunId: "run_foreign",
75+
resumeParentOnCompletion: true,
76+
});
77+
78+
await expect(result).rejects.toMatchObject<ServiceValidationError>({
79+
message: "Parent run not found in the calling environment",
80+
status: 404,
81+
});
82+
});
83+
});
Lines changed: 34 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,34 @@
1+
import type { RunStore } from "@internal/run-store";
2+
import { RunId } from "@trigger.dev/core/v3/isomorphic";
3+
import { ServiceValidationError } from "../../v3/services/baseService.server";
4+
5+
type BatchParentRunStore = Pick<RunStore, "findRun" | "findRunOnPrimary">;
6+
7+
export async function resolveBatchParentRun({
8+
runStore,
9+
environmentId,
10+
parentRunId,
11+
resumeParentOnCompletion,
12+
}: {
13+
runStore: BatchParentRunStore;
14+
environmentId: string;
15+
parentRunId?: string;
16+
resumeParentOnCompletion?: boolean;
17+
}): Promise<string | undefined> {
18+
if (!parentRunId || !resumeParentOnCompletion) {
19+
return;
20+
}
21+
22+
const runId = RunId.fromFriendlyId(parentRunId);
23+
const where = { id: runId, runtimeEnvironmentId: environmentId };
24+
const args = { select: { id: true } } as const;
25+
26+
const parentRun =
27+
(await runStore.findRun(where, args)) ?? (await runStore.findRunOnPrimary(where, args));
28+
29+
if (!parentRun) {
30+
throw new ServiceValidationError("Parent run not found in the calling environment", 404);
31+
}
32+
33+
return parentRun.id;
34+
}

0 commit comments

Comments
 (0)