feat: run annotations (#3241)

Adds an `annotations` JSONB column to task runs that captures where and
how each run was triggered.
This enables filtering and analyzing trigger origins without querying up
the run tree. Also enables making scheduling decisions based on the
trigger source, e.g., use separate affinities for scheduled runs.

Each run records:
- **triggerSource**: who initiated it (sdk, api, dashboard, cli, mcp,
schedule)
- **triggerAction**: what kind of action (trigger, replay, test)
- **rootTriggerSource**: the trigger source of the root ancestor,
propagated through the entire run
 tree
- **rootScheduleId**: schedule id, in case the run tree was triggered
from a schedule

Currently the main motivation for annotations it to determine whether a
run is part of a schedule-originated tree without traversing ancestors.

### A couple of design considerations
- **Decoupled source from method**: triggerSource and triggerAction are
separate fields to avoid
combinatorial explosion (every new source × every new action)
- **Server-side first**: all annotation values are primarily determined
on the server, only a minor SDK change needed
- **Forward-compatible**: annotation fields use
`z.enum([...]).or(anyString)` so new values can be
added without breaking validation; we currently don't need an explicit
version field for annotations.

Note: `metadata` would have been a more fitting name for the db column,
as it is consistent with other tables where we store this type of
information. It is already in use to store user metadata though, so we
go with `annotations` instead.
This commit is contained in:
Saadi Myftija
2026-03-23 16:07:30 +01:00
committed by GitHub
parent 88f755082f
commit d4772b5f60
31 changed files with 366 additions and 28 deletions
+8
View File
@@ -0,0 +1,8 @@
---
"@trigger.dev/redis-worker": patch
"@trigger.dev/sdk": patch
"trigger.dev": patch
"@trigger.dev/core": patch
---
Adapted the CLI API client to propagate the trigger source via http headers.
@@ -5,6 +5,7 @@ import { prisma } from "~/db.server";
import { authenticateApiRequest } from "~/services/apiAuth.server";
import { logger } from "~/services/logger.server";
import { ReplayTaskRunService } from "~/v3/services/replayTaskRun.server";
import { sanitizeTriggerSource } from "~/utils/triggerSource";
const ParamsSchema = z.object({
/* This is the run friendly ID */
@@ -41,8 +42,11 @@ export async function action({ request, params }: ActionFunctionArgs) {
return json({ error: "Run not found" }, { status: 404 });
}
const triggerSource =
sanitizeTriggerSource(request.headers.get("x-trigger-source")) ?? "api";
const service = new ReplayTaskRunService();
const newRun = await service.call(taskRun);
const newRun = await service.call(taskRun, { triggerSource });
if (!newRun) {
return json({ error: "Failed to create new run" }, { status: 400 });
@@ -19,6 +19,7 @@ import {
handleRequestIdempotency,
saveRequestIdempotency,
} from "~/utils/requestIdempotency.server";
import { sanitizeTriggerSource } from "~/utils/triggerSource";
import { ServiceValidationError } from "~/v3/services/baseService.server";
import { OutOfEntitlementError, TriggerTaskService } from "~/v3/services/triggerTask.server";
@@ -36,6 +37,7 @@ export const HeadersSchema = z.object({
"x-trigger-engine-version": RunEngineVersionSchema.nullish(),
"x-trigger-request-idempotency-key": z.string().nullish(),
"x-trigger-realtime-streams-version": z.string().nullish(),
"x-trigger-source": z.string().nullish(),
traceparent: z.string().optional(),
tracestate: z.string().optional(),
});
@@ -67,6 +69,7 @@ const { action, loader } = createActionApiRoute(
"x-trigger-engine-version": engineVersion,
"x-trigger-request-idempotency-key": requestIdempotencyKey,
"x-trigger-realtime-streams-version": realtimeStreamsVersion,
"x-trigger-source": triggerSourceHeader,
} = headers;
const cachedResponse = await handleRequestIdempotency(requestIdempotencyKey, {
@@ -119,6 +122,8 @@ const { action, loader } = createActionApiRoute(
realtimeStreamsVersion: determineRealtimeStreamsVersion(
realtimeStreamsVersion ?? undefined
),
triggerSource: isFromWorker ? "sdk" : sanitizeTriggerSource(triggerSourceHeader) ?? "api",
triggerAction: "trigger",
},
engineVersion ?? undefined
);
@@ -15,6 +15,7 @@ import {
BatchTriggerV3Service,
} from "~/v3/services/batchTriggerV3.server";
import { OutOfEntitlementError } from "~/v3/services/triggerTask.server";
import { sanitizeTriggerSource } from "~/utils/triggerSource";
import { HeadersSchema } from "./api.v1.tasks.$taskId.trigger";
import { determineRealtimeStreamsVersion } from "~/services/realtime/v1StreamsGlobal.server";
import { extractJwtSigningSecretKey } from "~/services/realtime/jwtAuth.server";
@@ -72,6 +73,7 @@ const { action, loader } = createActionApiRoute(
"x-trigger-engine-version": engineVersion,
"batch-processing-strategy": batchProcessingStrategy,
"x-trigger-realtime-streams-version": realtimeStreamsVersion,
"x-trigger-source": triggerSourceHeader,
traceparent,
tracestate,
} = headers;
@@ -113,6 +115,8 @@ const { action, loader } = createActionApiRoute(
realtimeStreamsVersion: determineRealtimeStreamsVersion(
realtimeStreamsVersion ?? undefined
),
triggerSource: isFromWorker ? "sdk" : sanitizeTriggerSource(triggerSourceHeader) ?? "api",
triggerAction: "trigger",
});
const $responseHeaders = await responseHeaders(
@@ -17,6 +17,7 @@ import {
import { ServiceValidationError } from "~/v3/services/baseService.server";
import { BatchProcessingStrategy } from "~/v3/services/batchTriggerV3.server";
import { OutOfEntitlementError } from "~/v3/services/triggerTask.server";
import { sanitizeTriggerSource } from "~/utils/triggerSource";
import { HeadersSchema } from "./api.v1.tasks.$taskId.trigger";
import { determineRealtimeStreamsVersion } from "~/services/realtime/v1StreamsGlobal.server";
import { extractJwtSigningSecretKey } from "~/services/realtime/jwtAuth.server";
@@ -62,6 +63,7 @@ const { action, loader } = createActionApiRoute(
"batch-processing-strategy": batchProcessingStrategy,
"x-trigger-request-idempotency-key": requestIdempotencyKey,
"x-trigger-realtime-streams-version": realtimeStreamsVersion,
"x-trigger-source": triggerSourceHeader,
traceparent,
tracestate,
} = headers;
@@ -127,6 +129,8 @@ const { action, loader } = createActionApiRoute(
realtimeStreamsVersion: determineRealtimeStreamsVersion(
realtimeStreamsVersion ?? undefined
),
triggerSource: isFromWorker ? "sdk" : sanitizeTriggerSource(triggerSourceHeader) ?? "api",
triggerAction: "trigger",
});
const $responseHeaders = await responseHeaders(
+3
View File
@@ -13,6 +13,7 @@ import {
} from "~/utils/requestIdempotency.server";
import { ServiceValidationError } from "~/v3/services/baseService.server";
import { OutOfEntitlementError } from "~/v3/services/triggerTask.server";
import { sanitizeTriggerSource } from "~/utils/triggerSource";
import { HeadersSchema } from "./api.v1.tasks.$taskId.trigger";
import { determineRealtimeStreamsVersion } from "~/services/realtime/v1StreamsGlobal.server";
import { extractJwtSigningSecretKey } from "~/services/realtime/jwtAuth.server";
@@ -65,6 +66,7 @@ const { action, loader } = createActionApiRoute(
"x-trigger-worker": isFromWorker,
"x-trigger-client": triggerClient,
"x-trigger-realtime-streams-version": realtimeStreamsVersion,
"x-trigger-source": triggerSourceHeader,
traceparent,
tracestate,
} = headers;
@@ -132,6 +134,7 @@ const { action, loader } = createActionApiRoute(
realtimeStreamsVersion: determineRealtimeStreamsVersion(
realtimeStreamsVersion ?? undefined
),
triggerSource: isFromWorker ? "sdk" : sanitizeTriggerSource(triggerSourceHeader) ?? "api",
});
const $responseHeaders = await responseHeaders(
@@ -214,6 +214,7 @@ export const action: ActionFunction = async ({ request, params }) => {
ttlSeconds: submission.value.ttlSeconds,
version: submission.value.version,
prioritySeconds: submission.value.prioritySeconds,
triggerSource: "dashboard",
});
if (!newRun) {
@@ -48,6 +48,8 @@ export type BatchTriggerTaskServiceOptions = {
spanParentAsLink?: boolean;
oneTimeUseToken?: string;
realtimeStreamsVersion?: "v1" | "v2";
triggerSource?: string;
triggerAction?: string;
};
/**
@@ -678,6 +680,8 @@ export class RunEngineBatchTriggerService extends WithRunEngine {
batchId: batch.id,
batchIndex: currentIndex,
realtimeStreamsVersion: options?.realtimeStreamsVersion,
triggerSource: options?.triggerSource ?? "api",
triggerAction: options?.triggerAction ?? "trigger",
},
"V2"
);
@@ -17,6 +17,7 @@ export type CreateBatchServiceOptions = {
spanParentAsLink?: boolean;
oneTimeUseToken?: string;
realtimeStreamsVersion?: "v1" | "v2";
triggerSource?: string;
};
/**
@@ -143,6 +144,7 @@ export class CreateBatchService extends WithRunEngine {
idempotencyKey: body.idempotencyKey,
processingConcurrency: config.processingConcurrency,
planType,
triggerSource: options.triggerSource,
};
await this._engine.initializeBatch(initOptions);
@@ -6,6 +6,7 @@ import {
import { Tracer } from "@opentelemetry/api";
import { tryCatch } from "@trigger.dev/core/utils";
import {
RunAnnotations,
TaskRunError,
taskRunErrorEnhancer,
taskRunErrorToString,
@@ -289,6 +290,17 @@ export class RunEngineTriggerTaskService {
const workerQueue = await this.queueConcern.getWorkerQueue(environment, body.options?.region);
// Build annotations for this run
const triggerSource = options.triggerSource ?? "api";
const triggerAction = options.triggerAction ?? "trigger";
const parentAnnotations = RunAnnotations.safeParse(parentRun?.annotations).data;
const annotations = {
triggerSource,
triggerAction,
rootTriggerSource: parentAnnotations?.rootTriggerSource ?? triggerSource,
rootScheduleId: parentAnnotations?.rootScheduleId || options.scheduleId || undefined,
};
try {
return await this.traceEventConcern.traceRun(
triggerRequest,
@@ -369,6 +381,7 @@ export class RunEngineTriggerTaskService {
planType,
realtimeStreamsVersion: options.realtimeStreamsVersion,
debounce: body.options?.debounce,
annotations,
// When debouncing with triggerAndWait, create a span for the debounced trigger
onDebounced:
body.options?.debounce && body.options?.resumeParentOnCompletion
+9
View File
@@ -0,0 +1,9 @@
const ALLOWED_TRIGGER_SOURCES = new Set(["sdk", "cli", "mcp"]);
/** Validates a client-provided trigger source header against the allowlist. */
export function sanitizeTriggerSource(value: string | null | undefined): string | undefined {
if (value && ALLOWED_TRIGGER_SOURCES.has(value)) {
return value;
}
return undefined;
}
@@ -750,6 +750,8 @@ export function setupBatchQueueCallbacks() {
batchIndex: itemIndex,
realtimeStreamsVersion: meta.realtimeStreamsVersion,
planType: meta.planType,
triggerSource: meta.parentRunId ? "sdk" : meta.triggerSource ?? "api",
triggerAction: "trigger",
},
"V2"
);
@@ -106,6 +106,8 @@ function createScheduleEngine() {
scheduleInstanceId,
queueTimestamp: exactScheduleTime,
overrideCreatedAt: exactScheduleTime,
triggerSource: "schedule",
triggerAction: "trigger",
}
);
@@ -57,6 +57,8 @@ export type BatchTriggerTaskServiceOptions = {
spanParentAsLink?: boolean;
oneTimeUseToken?: string;
realtimeStreamsVersion?: "v1" | "v2";
triggerSource?: string;
triggerAction?: string;
};
type RunItemData = {
@@ -853,6 +855,8 @@ export class BatchTriggerV3Service extends BaseService {
skipChecks: true,
runFriendlyId: task.runId,
realtimeStreamsVersion: options?.realtimeStreamsVersion,
triggerSource: options?.triggerSource ?? "api",
triggerAction: options?.triggerAction ?? "trigger",
}
);
@@ -242,6 +242,7 @@ export class BulkActionService extends BaseService {
const [error, result] = await tryCatch(
replayService.call(run, {
bulkActionId: bulkActionId,
triggerSource: "dashboard",
})
);
if (error) {
@@ -27,7 +27,7 @@ export class PerformBulkActionService extends BaseService {
switch (item.group.type) {
case "REPLAY": {
const service = new ReplayTaskRunService(this._prisma);
const result = await service.call(item.sourceRun);
const result = await service.call(item.sourceRun, { triggerSource: "dashboard" });
await this._prisma.bulkActionItem.update({
where: { id: item.id },
@@ -18,6 +18,7 @@ type OverrideOptions = {
payload?: string;
metadata?: unknown;
bulkActionId?: string;
triggerSource?: string;
} & RunOptionsData;
export class ReplayTaskRunService extends BaseService {
@@ -123,6 +124,8 @@ export class ReplayTaskRunService extends BaseService {
realtimeStreamsVersion: determineRealtimeStreamsVersion(
existingTaskRun.realtimeStreamsVersion
),
triggerSource: overrideOptions.triggerSource ?? "api",
triggerAction: "replay",
}
);
+29 -22
View File
@@ -11,28 +11,35 @@ export class TestTaskService extends BaseService {
switch (triggerSource) {
case "STANDARD": {
const result = await triggerTaskService.call(data.taskIdentifier, environment, {
payload: data.payload,
options: {
test: true,
metadata: data.metadata,
delay: data.delaySeconds ? new Date(Date.now() + data.delaySeconds * 1000) : undefined,
ttl: data.ttlSeconds,
idempotencyKey: data.idempotencyKey,
idempotencyKeyTTL: data.idempotencyKeyTTLSeconds
? `${data.idempotencyKeyTTLSeconds}s`
: undefined,
queue: data.queue ? { name: data.queue } : undefined,
concurrencyKey: data.concurrencyKey,
maxAttempts: data.maxAttempts,
maxDuration: data.maxDurationSeconds,
tags: data.tags,
machine: data.machine,
region: data.region,
lockToVersion: data.version === "latest" ? undefined : data.version,
priority: data.prioritySeconds,
const result = await triggerTaskService.call(
data.taskIdentifier,
environment,
{
payload: data.payload,
options: {
test: true,
metadata: data.metadata,
delay: data.delaySeconds
? new Date(Date.now() + data.delaySeconds * 1000)
: undefined,
ttl: data.ttlSeconds,
idempotencyKey: data.idempotencyKey,
idempotencyKeyTTL: data.idempotencyKeyTTLSeconds
? `${data.idempotencyKeyTTLSeconds}s`
: undefined,
queue: data.queue ? { name: data.queue } : undefined,
concurrencyKey: data.concurrencyKey,
maxAttempts: data.maxAttempts,
maxDuration: data.maxDurationSeconds,
tags: data.tags,
machine: data.machine,
region: data.region,
lockToVersion: data.version === "latest" ? undefined : data.version,
priority: data.prioritySeconds,
},
},
});
{ triggerSource: "dashboard", triggerAction: "test" }
);
return result?.run;
}
@@ -72,7 +79,7 @@ export class TestTaskService extends BaseService {
priority: data.prioritySeconds,
},
},
{ customIcon: "scheduled" }
{ customIcon: "scheduled", triggerSource: "dashboard", triggerAction: "test" }
);
return result?.run;
@@ -33,6 +33,8 @@ export type TriggerTaskServiceOptions = {
replayedFromTaskRunFriendlyId?: string;
planType?: string;
realtimeStreamsVersion?: "v1" | "v2";
triggerSource?: string;
triggerAction?: string;
};
export class OutOfEntitlementError extends Error {
@@ -0,0 +1,2 @@
-- AlterTable
ALTER TABLE "public"."TaskRun" ADD COLUMN "annotations" JSONB;
@@ -837,6 +837,9 @@ model TaskRun {
metadataType String @default("application/json")
metadataVersion Int @default(1)
/// Structured annotations: triggerSource, triggerAction, rootTriggerSource, rootScheduleId
annotations Json?
/// Run output
output String?
outputType String @default("application/json")
@@ -296,6 +296,7 @@ export class BatchQueue {
realtimeStreamsVersion: options.realtimeStreamsVersion,
idempotencyKey: options.idempotencyKey,
processingConcurrency: options.processingConcurrency,
triggerSource: options.triggerSource,
};
// Store metadata in completion tracker
@@ -79,6 +79,8 @@ export const BatchMeta = z.object({
processingConcurrency: z.number().optional(),
/** Plan type for billing (e.g., "free", "paid") - used when skipChecks is enabled */
planType: z.string().optional(),
/** Trigger source for run annotations (e.g., "sdk", "cli", "mcp") */
triggerSource: z.string().optional(),
});
export type BatchMeta = z.infer<typeof BatchMeta>;
@@ -168,6 +170,8 @@ export type InitializeBatchOptions = {
processingConcurrency?: number;
/** Plan type for billing (e.g., "free", "paid") - used when skipChecks is enabled */
planType?: string;
/** Trigger source for run annotations (e.g., "sdk", "cli", "mcp") */
triggerSource?: string;
};
/**
@@ -495,6 +495,7 @@ export class RunEngine {
planType,
realtimeStreamsVersion,
debounce,
annotations,
onDebounced,
}: TriggerParams,
tx?: PrismaClientOrTransaction
@@ -668,6 +669,7 @@ export class RunEngine {
createdAt: new Date(),
}
: undefined,
annotations,
executionSnapshots: {
create: {
engine: "V2",
@@ -328,4 +328,203 @@ describe("RunEngine trigger()", () => {
await engine.quit();
}
});
containerTest("Annotations are stored on the run", async ({ prisma, redisOptions }) => {
const authenticatedEnvironment = await setupAuthenticatedEnvironment(prisma, "PRODUCTION");
const engine = new RunEngine({
prisma,
worker: {
redis: redisOptions,
workers: 1,
tasksPerWorker: 10,
pollIntervalMs: 100,
},
queue: {
redis: redisOptions,
masterQueueConsumersDisabled: true,
processWorkerQueueDebounceMs: 50,
},
runLock: { redis: redisOptions },
machines: {
defaultMachine: "small-1x",
machines: {
"small-1x": { name: "small-1x" as const, cpu: 0.5, memory: 0.5, centsPerMs: 0.0001 },
},
baseCostInCents: 0.0001,
},
tracer: trace.getTracer("test", "0.0.0"),
});
try {
const taskIdentifier = "test-task";
await setupBackgroundWorker(engine, authenticatedEnvironment, taskIdentifier);
const run = await engine.trigger(
{
number: 1,
friendlyId: "run_ann1234",
environment: authenticatedEnvironment,
taskIdentifier,
payload: "{}",
payloadType: "application/json",
context: {},
traceContext: {},
traceId: "t12345",
spanId: "s12345",
workerQueue: "main",
queue: `task/${taskIdentifier}`,
isTest: false,
tags: [],
annotations: {
triggerSource: "schedule",
triggerAction: "trigger",
rootTriggerSource: "schedule",
rootScheduleId: "sched_abc123",
},
},
prisma
);
const runFromDb = await prisma.taskRun.findUnique({
where: { id: run.id },
});
expect(runFromDb).toBeDefined();
expect(runFromDb?.annotations).toEqual({
triggerSource: "schedule",
triggerAction: "trigger",
rootTriggerSource: "schedule",
rootScheduleId: "sched_abc123",
});
} finally {
await engine.quit();
}
});
containerTest(
"Annotations propagation pattern (parent → child)",
async ({ prisma, redisOptions }) => {
const authenticatedEnvironment = await setupAuthenticatedEnvironment(prisma, "PRODUCTION");
const engine = new RunEngine({
prisma,
worker: {
redis: redisOptions,
workers: 1,
tasksPerWorker: 10,
pollIntervalMs: 100,
},
queue: {
redis: redisOptions,
masterQueueConsumersDisabled: true,
processWorkerQueueDebounceMs: 50,
},
runLock: { redis: redisOptions },
machines: {
defaultMachine: "small-1x",
machines: {
"small-1x": {
name: "small-1x" as const,
cpu: 0.5,
memory: 0.5,
centsPerMs: 0.0001,
},
},
baseCostInCents: 0.0001,
},
tracer: trace.getTracer("test", "0.0.0"),
});
try {
const parentTask = "parent-task";
const childTask = "child-task";
await setupBackgroundWorker(engine, authenticatedEnvironment, [parentTask, childTask]);
// Trigger parent with schedule annotations
const parentRun = await engine.trigger(
{
number: 1,
friendlyId: "run_p1234ann",
environment: authenticatedEnvironment,
taskIdentifier: parentTask,
payload: "{}",
payloadType: "application/json",
context: {},
traceContext: {},
traceId: "t12345",
spanId: "s12345",
workerQueue: "main",
queue: `task/${parentTask}`,
isTest: false,
tags: [],
annotations: {
triggerSource: "schedule",
triggerAction: "trigger",
rootTriggerSource: "schedule",
rootScheduleId: "sched_abc123",
},
},
prisma
);
// Trigger child — simulating what RunEngineTriggerTaskService builds:
// triggerSource is "sdk" (child triggered from within parent),
// but rootTriggerSource and rootScheduleId are propagated from parent
const childRun = await engine.trigger(
{
number: 2,
friendlyId: "run_c1234ann",
environment: authenticatedEnvironment,
taskIdentifier: childTask,
payload: "{}",
payloadType: "application/json",
context: {},
traceContext: {},
traceId: "t12345",
spanId: "s12346",
workerQueue: "main",
queue: `task/${childTask}`,
isTest: false,
tags: [],
parentTaskRunId: parentRun.id,
resumeParentOnCompletion: true,
annotations: {
triggerSource: "sdk",
triggerAction: "trigger",
rootTriggerSource: "schedule",
rootScheduleId: "sched_abc123",
},
},
prisma
);
const parentFromDb = await prisma.taskRun.findUnique({
where: { id: parentRun.id },
});
const childFromDb = await prisma.taskRun.findUnique({
where: { id: childRun.id },
});
// Parent: schedule-triggered
expect(parentFromDb?.annotations).toEqual({
triggerSource: "schedule",
triggerAction: "trigger",
rootTriggerSource: "schedule",
rootScheduleId: "sched_abc123",
});
// Child: sdk-triggered but root is still schedule
expect(childFromDb?.annotations).toEqual({
triggerSource: "sdk",
triggerAction: "trigger",
rootTriggerSource: "schedule",
rootScheduleId: "sched_abc123",
});
} finally {
await engine.quit();
}
}
);
});
@@ -222,6 +222,12 @@ export type TriggerParams = {
mode?: "leading" | "trailing";
maxDelay?: string;
};
annotations?: {
triggerSource: string;
triggerAction: string;
rootTriggerSource: string;
rootScheduleId?: string;
};
/**
* Called when a run is debounced (existing delayed run found with triggerAndWait).
* Return spanIdToComplete to enable span closing when the run completes.
+5 -1
View File
@@ -60,16 +60,19 @@ import { VERSION } from "./version.js";
export class CliApiClient {
private engineURL: string;
private source: "cli" | "mcp";
constructor(
public readonly apiURL: string,
// TODO: consider making this required
public readonly accessToken?: string,
public readonly branch?: string
public readonly branch?: string,
options?: { source?: "cli" | "mcp" }
) {
this.apiURL = apiURL.replace(/\/$/, "");
this.engineURL = this.apiURL;
this.branch = branch;
this.source = options?.source ?? "cli";
}
async createAuthorizationCode() {
@@ -819,6 +822,7 @@ export class CliApiClient {
const headers: Record<string, string> = {
Authorization: `Bearer ${this.accessToken}`,
"Content-Type": "application/json",
"x-trigger-source": this.source,
};
if (this.branch) {
+3 -1
View File
@@ -109,7 +109,9 @@ export class McpContext {
);
}
return new ApiClient(cliApiClient.apiURL, jwt.data.token);
return new ApiClient(cliApiClient.apiURL, jwt.data.token, undefined, {
additionalHeaders: { "x-trigger-source": "mcp" },
});
}
public async getCwd() {
+4 -1
View File
@@ -40,7 +40,9 @@ export type ZodFetchOptions<TData = any> = {
export type AnyZodFetchOptions = ZodFetchOptions<any>;
export type ApiRequestOptions = Pick<ZodFetchOptions, "retry">;
export type ApiRequestOptions = Pick<ZodFetchOptions, "retry"> & {
additionalHeaders?: Record<string, string>;
};
type KeysEnum<T> = { [P in keyof Required<T>]: true };
@@ -49,6 +51,7 @@ type KeysEnum<T> = { [P in keyof Required<T>]: true };
// compiler such that any missing / extraneous keys will cause an error.
const requestOptionsKeys: KeysEnum<ApiRequestOptions> = {
retry: true,
additionalHeaders: true,
};
export const isRequestOptions = (obj: unknown): obj is ApiRequestOptions => {
+16 -1
View File
@@ -189,6 +189,7 @@ export class ApiClient {
public readonly accessToken: string;
public readonly previewBranch?: string;
public readonly futureFlags: ApiClientFutureFlags;
private readonly additionalHeaders?: Record<string, string>;
private readonly defaultRequestOptions: ZodFetchOptions;
constructor(
@@ -201,7 +202,9 @@ export class ApiClient {
this.accessToken = accessToken;
this.baseUrl = baseUrl.replace(/\/$/, "");
this.previewBranch = previewBranch;
this.defaultRequestOptions = mergeRequestOptions(DEFAULT_ZOD_FETCH_OPTIONS, requestOptions);
const { additionalHeaders, ...restRequestOptions } = requestOptions;
this.additionalHeaders = additionalHeaders;
this.defaultRequestOptions = mergeRequestOptions(DEFAULT_ZOD_FETCH_OPTIONS, restRequestOptions);
this.futureFlags = futureFlags;
}
@@ -1540,6 +1543,18 @@ export class ApiClient {
),
};
if (this.additionalHeaders) {
for (const [key, value] of Object.entries(this.additionalHeaders)) {
if (!(key in headers)) {
headers[key] = value;
}
}
}
if (!headers["x-trigger-source"]) {
headers["x-trigger-source"] = "sdk";
}
if (this.previewBranch) {
headers["x-trigger-branch"] = this.previewBranch;
}
+19
View File
@@ -540,6 +540,25 @@ export const DeploymentTriggeredVia = z
export type DeploymentTriggeredVia = z.infer<typeof DeploymentTriggeredVia>;
export const TriggerSource = z
.enum(["sdk", "api", "dashboard", "cli", "mcp", "schedule"])
.or(anyString);
export type TriggerSource = z.infer<typeof TriggerSource>;
export const TriggerAction = z.enum(["trigger", "replay", "test"]).or(anyString);
export type TriggerAction = z.infer<typeof TriggerAction>;
export const RunAnnotations = z.object({
triggerSource: TriggerSource,
triggerAction: TriggerAction,
rootTriggerSource: TriggerSource,
rootScheduleId: z.string().optional(),
});
export type RunAnnotations = z.infer<typeof RunAnnotations>;
export const UpsertBranchRequestBody = z.object({
git: GitMeta.optional(),
env: z.enum(["preview"]),