Associate child runs with the parent span ID (#1352)

* Add safe rootTaskRunId index and a README to @trigger.dev/database

* Associate child runs with the span ID of the span in the parent run that triggered the child run

* Update deprecation notice doc links
This commit is contained in:
Eric Allam
2024-09-24 15:47:44 +01:00
committed by GitHub
parent bff38dceb5
commit d361e24bee
10 changed files with 85 additions and 10 deletions
+5
View File
@@ -0,0 +1,5 @@
---
"@trigger.dev/core": patch
---
Add otel propagation headers "below" the API fetch span, to attribute the child runs with the proper parent span ID
+3 -2
View File
@@ -835,7 +835,8 @@ export class EventRepository {
options: TraceEventOptions & { incomplete?: boolean },
callback: (
e: EventBuilder,
traceContext: Record<string, string | undefined>
traceContext: Record<string, string | undefined>,
traceparent?: { traceId: string; spanId: string }
) => Promise<TResult>
): Promise<TResult> {
const propagatedContext = extractContextFromCarrier(options.context ?? {});
@@ -892,7 +893,7 @@ export class EventRepository {
},
};
const result = await callback(eventBuilder, traceContext);
const result = await callback(eventBuilder, traceContext, propagatedContext?.traceparent);
const duration = process.hrtime.bigint() - start;
@@ -240,7 +240,7 @@ export class TriggerTaskService extends BaseService {
incomplete: true,
immediate: true,
},
async (event, traceContext) => {
async (event, traceContext, traceparent) => {
const run = await autoIncrementCounter.incrementInTransaction(
`v3-run:${environment.id}:${taskId}`,
async (num, tx) => {
@@ -307,6 +307,8 @@ export class TriggerTaskService extends BaseService {
traceContext: traceContext,
traceId: event.traceId,
spanId: event.spanId,
parentSpanId:
options.parentAsLinkType === "replay" ? undefined : traceparent?.spanId,
lockedToVersionId: lockedToBackgroundWorker?.id,
concurrencyKey: body.options?.concurrencyKey,
queue: queueName,
+3 -3
View File
@@ -228,7 +228,7 @@ function validateConfig(config: TriggerConfig, warn = true) {
if (config.additionalFiles && config.additionalFiles.length > 0) {
warn &&
prettyWarning(
`The "additionalFiles" option is deprecated and will be removed. Use the "additionalFiles" build extension instead. See https://trigger.dev/docs/guides/new-build-system-preview#additionalfiles for more information.`
`The "additionalFiles" option is deprecated and will be removed. Use the "additionalFiles" build extension instead. See https://trigger.dev/docs/config/config-file#additionalfiles for more information.`
);
config.build ??= {};
@@ -239,7 +239,7 @@ function validateConfig(config: TriggerConfig, warn = true) {
if (config.additionalPackages && config.additionalPackages.length > 0) {
warn &&
prettyWarning(
`The "additionalPackages" option is deprecated and will be removed. Use the "additionalPackages" build extension instead. See https://trigger.dev/docs/guides/new-build-system-preview#additionalpackages for more information.`
`The "additionalPackages" option is deprecated and will be removed. Use the "additionalPackages" build extension instead. See https://trigger.dev/docs/config/config-file#additionalpackages for more information.`
);
config.build ??= {};
@@ -275,7 +275,7 @@ function validateConfig(config: TriggerConfig, warn = true) {
if ("resolveEnvVars" in config && typeof config.resolveEnvVars === "function") {
warn &&
prettyWarning(
`The "resolveEnvVars" option is deprecated and will be removed. Use the "syncEnvVars" build extension instead. See https://trigger.dev/docs/guides/new-build-system-preview#resolveenvvars for more information.`
`The "resolveEnvVars" option is deprecated and will be removed. Use the "syncEnvVars" build extension instead. See https://trigger.dev/docs/config/config-file#syncenvvars for more information.`
);
const resolveEnvVarsFn = config.resolveEnvVars as ResolveEnvironmentVariablesFunction;
+24 -2
View File
@@ -4,7 +4,7 @@ import { RetryOptions } from "../schemas/index.js";
import { calculateNextRetryDelay } from "../utils/retries.js";
import { ApiConnectionError, ApiError, ApiSchemaValidationError } from "./errors.js";
import { Attributes, Span } from "@opentelemetry/api";
import { Attributes, Span, context, propagation } from "@opentelemetry/api";
import { SemanticInternalAttributes } from "../semanticInternalAttributes.js";
import { TriggerTracer } from "../tracer.js";
import { accessoryAttributes } from "../utils/styleAttributes.js";
@@ -184,9 +184,11 @@ async function _doZodFetch<TResponseBodySchema extends z.ZodTypeAny>(
requestInit?: PromiseOrValue<RequestInit>,
options?: ZodFetchOptions
): Promise<ZodFetchResult<z.output<TResponseBodySchema>>> {
const $requestInit = await requestInit;
let $requestInit = await requestInit;
return traceZodFetch({ url, requestInit: $requestInit, options }, async (span) => {
$requestInit = injectPropagationHeadersIfInWorker($requestInit);
const result = await _doZodFetchWithRetries(schema, url, $requestInit, options);
if (options?.onResponseBody && span) {
@@ -577,3 +579,23 @@ export function isEmptyObj(obj: Object | null | undefined): boolean {
export function hasOwn(obj: Object, key: string): boolean {
return Object.prototype.hasOwnProperty.call(obj, key);
}
// If the requestInit has a header x-trigger-worker = true, then we will do
// propagation.inject(context.active(), headers);
// and return the new requestInit.
function injectPropagationHeadersIfInWorker(requestInit?: RequestInit): RequestInit | undefined {
const headers = new Headers(requestInit?.headers);
if (headers.get("x-trigger-worker") !== "true") {
return requestInit;
}
const headersObject = Object.fromEntries(headers.entries());
propagation.inject(context.active(), headersObject);
return {
...requestInit,
headers: new Headers(headersObject),
};
}
-2
View File
@@ -1,4 +1,3 @@
import { context, propagation } from "@opentelemetry/api";
import { z } from "zod";
import {
AddTagsRequestBody,
@@ -509,7 +508,6 @@ export class ApiClient {
// Only inject the context if we are inside a task
if (taskContext.isInsideTask) {
headers["x-trigger-worker"] = "true";
propagation.inject(context.active(), headers);
if (spanParentAsLink) {
headers["x-trigger-span-parent-as-link"] = "1";
+38
View File
@@ -0,0 +1,38 @@
## @trigger.dev/database
This is the internal database package for the Trigger.dev project. It exports a generated prisma client that can be instantiated with a connection string.
### How to add a new index on a large table
1. Modify the Prisma.schema with a single index change (no other changes, just one index at a time)
2. Create a Prisma migration using `cd packages/database && pnpm run db:migrate:dev --create-only`
3. Modify the SQL file: add IF NOT EXISTS to it and CONCURRENTLY:
```sql
CREATE INDEX CONCURRENTLY IF NOT EXISTS "JobRun_eventId_idx" ON "JobRun" ("eventId");
```
4. Dont apply the Prisma migration locally yet. This is a good opportunity to test the flow.
5. Manually apply the index to your database, by running the index command.
6. Then locally run `pnpm run db:migrate:deploy`
#### Before deploying
Run the index creation before deploying
```sql
CREATE INDEX CONCURRENTLY IF NOT EXISTS "JobRun_eventId_idx" ON "JobRun" ("eventId");
```
These commands are useful:
```sql
-- creates an index safely, this can take a long time (2 mins maybe)
CREATE INDEX CONCURRENTLY IF NOT EXISTS "JobRun_eventId_idx" ON "JobRun" ("eventId");
-- checks the status of an index
SELECT * FROM pg_stat_progress_create_index WHERE relid = '"JobRun"'::regclass;
-- checks if the index is there
SELECT * FROM pg_indexes WHERE tablename = 'JobRun' AND indexname = 'JobRun_eventId_idx';
```
Now, when you deploy and prisma runs the migration, it will skip the index creation because it already exists. If you don't do this, there will be pain.
@@ -0,0 +1,2 @@
-- CreateIndex
CREATE INDEX CONCURRENTLY IF NOT EXISTS "TaskRun_rootTaskRunId_idx" ON "TaskRun"("rootTaskRunId");
@@ -0,0 +1,2 @@
-- AlterTable
ALTER TABLE "TaskRun" ADD COLUMN "parentSpanId" TEXT;
+5
View File
@@ -1748,9 +1748,14 @@ model TaskRun {
/// The depth of this task run in the task run hierarchy
depth Int @default(0)
/// The span ID of the "trigger" span in the parent task run
parentSpanId String?
@@unique([runtimeEnvironmentId, taskIdentifier, idempotencyKey])
// Finding child runs
@@index([parentTaskRunId])
// Finding ancestor runs
@@index([rootTaskRunId])
// Task activity graph
@@index([projectId, createdAt, taskIdentifier])
//Runs list