diff --git a/bun.lock b/bun.lock index 9215bed..e741119 100644 --- a/bun.lock +++ b/bun.lock @@ -23,25 +23,30 @@ "prettier": "catalog:", }, "devDependencies": { + "@opentelemetry/api": "catalog:", + "@opentelemetry/sdk-trace-base": "catalog:", "@testcontainers/postgresql": "catalog:", "@types/bun": "catalog:", "expect-type": "catalog:", "zod": "catalog:", }, "peerDependencies": { + "@opentelemetry/api": ">=1.4 <2", "typescript": "catalog:", }, }, }, "catalog": { + "@opentelemetry/api": "1.9.0", + "@opentelemetry/sdk-trace-base": "1.30.1", "@standard-schema/spec": "1.0.0", "@testcontainers/postgresql": "11.8.1", "@types/bun": "latest", "cron-parser": "5.4.0", "expect-type": "1.2.2", - "oxfmt": "latest", - "oxlint": "latest", - "oxlint-tsgolint": "latest", + "oxfmt": "0.15.0", + "oxlint": "1.30.0", + "oxlint-tsgolint": "0.8.3", "postgres": "3.4.7", "prettier": "3.4.2", "typescript": "5.7.2", @@ -58,6 +63,16 @@ "@js-sdsl/ordered-map": ["@js-sdsl/ordered-map@4.4.2", "", {}, "sha512-iUKgm52T8HOE/makSxjqoWhe95ZJA1/G1sYsGev2JDKUSS14KAgg1LHb+Ba+IPow0xflbnSkOsZcO08C7w1gYw=="], + "@opentelemetry/api": ["@opentelemetry/api@1.9.0", "", {}, "sha512-3giAOQvZiH5F9bMlMiv8+GSPMeqg0dbaeo58/0SlA9sxSqZhnUtxzX9/2FzyhS9sWQf5S0GJE0AKBrFqjpeYcg=="], + + "@opentelemetry/core": ["@opentelemetry/core@1.30.1", "", { "dependencies": { "@opentelemetry/semantic-conventions": "1.28.0" }, "peerDependencies": { "@opentelemetry/api": ">=1.0.0 <1.10.0" } }, "sha512-OOCM2C/QIURhJMuKaekP3TRBxBKxG/TWWA0TL2J6nXUtDnuCtccy49LUJF8xPFXMX+0LMcxFpCo8M9cGY1W6rQ=="], + + "@opentelemetry/resources": ["@opentelemetry/resources@1.30.1", "", { "dependencies": { "@opentelemetry/core": "1.30.1", "@opentelemetry/semantic-conventions": "1.28.0" }, "peerDependencies": { "@opentelemetry/api": ">=1.0.0 <1.10.0" } }, "sha512-5UxZqiAgLYGFjS4s9qm5mBVo433u+dSPUFWVWXmLAD4wB65oMCoXaJP1KJa9DIYYMeHu3z4BZcStG3LC593cWA=="], + + "@opentelemetry/sdk-trace-base": ["@opentelemetry/sdk-trace-base@1.30.1", "", { "dependencies": { "@opentelemetry/core": "1.30.1", "@opentelemetry/resources": "1.30.1", "@opentelemetry/semantic-conventions": "1.28.0" }, "peerDependencies": { "@opentelemetry/api": ">=1.0.0 <1.10.0" } }, "sha512-jVPgBbH1gCy2Lb7X0AVQ8XAfgg0pJ4nvl8/IiQA6nxOsPvS+0zMJaFSs2ltXe0J6C8dqjcnpyqINDJmU30+uOg=="], + + "@opentelemetry/semantic-conventions": ["@opentelemetry/semantic-conventions@1.28.0", "", {}, "sha512-lp4qAiMTD4sNWW4DbKLBkfiMZ4jbAboJIGOQr5DvciMRI494OapieI9qiODpOt0XBr1LjIDy1xAGAnVs5supTA=="], + "@oxfmt/darwin-arm64": ["@oxfmt/darwin-arm64@0.15.0", "", { "os": "darwin", "cpu": "arm64" }, "sha512-M5xiXkqtwG/1yVlNZJaXeFJGs1jVTT7Q6gfdIU3nO8wXIj0lsfcirc55XPObWohA8T73oLh3Cp1oSQluxVXQrQ=="], "@oxfmt/darwin-x64": ["@oxfmt/darwin-x64@0.15.0", "", { "os": "darwin", "cpu": "x64" }, "sha512-HXBZBV1oqmZWcmXXQE+zXNpe8nOXyeY9oLeW6KfflR1HCicR75Tmer6Y5/euVOb/mDfegkkjReRoQnYxJ6CmtQ=="], @@ -128,7 +143,7 @@ "@testcontainers/postgresql": ["@testcontainers/postgresql@11.8.1", "", { "dependencies": { "testcontainers": "^11.8.1" } }, "sha512-uPTFk9IY5v9Dm6HcVTOSvZQQ+HrdGiwDk+/LDG+v67yD81kgBlYpH730JJhZoO72d1pJHKwTAJ+5WPOvecl1pw=="], - "@types/bun": ["@types/bun@1.3.4", "", { "dependencies": { "bun-types": "1.3.4" } }, "sha512-EEPTKXHP+zKGPkhRLv+HI0UEX8/o+65hqARxLy8Ov5rIxMBPNTjeZww00CIihrIQGEQBYg+0roO5qOnS/7boGA=="], + "@types/bun": ["@types/bun@1.4.2", "", { "dependencies": { "bun-types": "1.4.2" } }, "sha512-GimotNn7+ZV0uVArItBbriZsR1oNf0+WTzPkdcFrzShI7k2norL0uzEaJT8T33dWr7O/c9ZDuAFQrctKCi72oQ=="], "@types/docker-modem": ["@types/docker-modem@3.0.6", "", { "dependencies": { "@types/node": "*", "@types/ssh2": "*" } }, "sha512-yKpAGEuKRSS8wwx0joknWxsmLha78wNMe9R2S3UNsVOkZded8UqOrV8KoeDXoXsjndxwyF3eIhyClGbO1SEhEg=="], @@ -186,7 +201,7 @@ "buildcheck": ["buildcheck@0.0.6", "", {}, "sha512-8f9ZJCUXyT1M35Jx7MkBgmBMo3oHTTBIPLiY9xyL0pl3T5RwcPEY8cUHr5LBNfu/fk6c2T4DJZuVM/8ZZT2D2A=="], - "bun-types": ["bun-types@1.3.4", "", { "dependencies": { "@types/node": "*" } }, "sha512-5ua817+BZPZOlNaRgGBpZJOSAQ9RQ17pkwPD0yR7CfJg+r8DgIILByFifDTa+IPDDxzf5VNhtNlcKqFzDgJvlQ=="], + "bun-types": ["bun-types@1.4.2", "", { "dependencies": { "@types/node": "*" } }, "sha512-bxV1FgK7yBIzjRe5zBozIM4Bem11ZJcCXSrjWRG3YWLt8yFDePu4cLjpebO8OvPeIE9trbyPF4fuj3Cia4Fj3w=="], "byline": ["byline@5.0.0", "", {}, "sha512-s6webAy+R4SR8XVuJWt2V2rGvhnrhxN+9S15GNuTK3wKPOXFF6RNc+8ug2XhH+2s4f+uudG4kUVYmYOQWL2g0Q=="], diff --git a/migrations/0000000001_setup.sql b/migrations/0000000001_setup.sql index 1ec5f99..4eb0f21 100644 --- a/migrations/0000000001_setup.sql +++ b/migrations/0000000001_setup.sql @@ -92,6 +92,7 @@ create table pgconductor._private_executions ( failed_at timestamptz, completed_at timestamptz, payload jsonb, + trace_context jsonb, run_at timestamptz default pgconductor._private_current_time() not null, locked_at timestamptz, locked_by uuid, @@ -344,7 +345,8 @@ create type pgconductor.execution_spec as ( dedupe_next_slot boolean, cron_expression text, priority integer, - "group" text + "group" text, + trace_context jsonb ); create type pgconductor.task_spec as ( @@ -603,6 +605,7 @@ begin task_key, queue, payload, + trace_context, run_at, dedupe_key, singleton_on, @@ -615,6 +618,7 @@ begin spec.task_key, coalesce(spec.queue, 'default'), spec.payload, + spec.trace_context, coalesce(spec.run_at, v_now), spec.dedupe_key, case @@ -652,7 +656,8 @@ create or replace function pgconductor.invoke( p_dedupe_next_slot boolean default false, p_cron_expression text default null, p_priority integer default null, - p_group text default null + p_group text default null, + p_trace_context jsonb default null ) returns table(id uuid) language plpgsql @@ -708,6 +713,7 @@ begin task_key, queue, payload, + trace_context, run_at, dedupe_key, singleton_on, @@ -719,6 +725,7 @@ begin p_task_key, p_queue, p_payload, + p_trace_context, v_run_at, p_dedupe_key, v_singleton_on, @@ -741,6 +748,7 @@ begin task_key, queue, payload, + trace_context, run_at, dedupe_key, singleton_on, @@ -752,6 +760,7 @@ begin p_task_key, p_queue, p_payload, + p_trace_context, v_next_singleton_on, p_dedupe_key, v_next_singleton_on, @@ -763,6 +772,7 @@ begin where singleton_on is not null and completed_at is null and failed_at is null and cancelled = false do update set payload = excluded.payload, + trace_context = excluded.trace_context, run_at = excluded.run_at, priority = excluded.priority, cron_expression = excluded.cron_expression, @@ -778,6 +788,7 @@ begin task_key, queue, payload, + trace_context, run_at, dedupe_key, cron_expression, @@ -788,6 +799,7 @@ begin p_task_key, p_queue, p_payload, + p_trace_context, v_run_at, p_dedupe_key, p_cron_expression, @@ -796,6 +808,7 @@ begin ) on conflict (task_key, dedupe_key, queue) do update set payload = excluded.payload, + trace_context = excluded.trace_context, run_at = excluded.run_at, priority = excluded.priority, cron_expression = excluded.cron_expression, diff --git a/package.json b/package.json index 7edec01..bdc6c65 100644 --- a/package.json +++ b/package.json @@ -12,9 +12,11 @@ "zod": "4.1.12", "@standard-schema/spec": "1.0.0", "cron-parser": "5.4.0", - "oxfmt": "latest", - "oxlint": "latest", - "oxlint-tsgolint": "latest" + "@opentelemetry/api": "1.9.0", + "@opentelemetry/sdk-trace-base": "1.30.1", + "oxfmt": "0.15.0", + "oxlint": "1.30.0", + "oxlint-tsgolint": "0.8.3" }, "devDependencies": { "oxfmt": "catalog:", diff --git a/packages/pgconductor-js/package.json b/packages/pgconductor-js/package.json index a3adf91..b82d8ec 100644 --- a/packages/pgconductor-js/package.json +++ b/packages/pgconductor-js/package.json @@ -33,7 +33,7 @@ "provenance": true }, "scripts": { - "build": "bun build src/index.ts --outdir dist --target node && bun build cli/index.ts --outdir dist/cli --target node", + "build": "bun build src/index.ts --outdir dist --target node --external @opentelemetry/api && bun build cli/index.ts --outdir dist/cli --target node --external @opentelemetry/api", "test": "bun test", "test:unit": "bun test tests/unit", "test:integration": "bun test tests/integration", @@ -41,12 +41,15 @@ "perf": "bun perf/run.ts" }, "devDependencies": { + "@opentelemetry/api": "catalog:", + "@opentelemetry/sdk-trace-base": "catalog:", "@testcontainers/postgresql": "catalog:", "@types/bun": "catalog:", "expect-type": "catalog:", "zod": "catalog:", }, "peerDependencies": { + "@opentelemetry/api": ">=1.4 <2", "typescript": "catalog:" }, "dependencies": { diff --git a/packages/pgconductor-js/src/conductor.ts b/packages/pgconductor-js/src/conductor.ts index 8026a2a..9412f9b 100644 --- a/packages/pgconductor-js/src/conductor.ts +++ b/packages/pgconductor-js/src/conductor.ts @@ -25,6 +25,17 @@ import { import { Worker, type WorkerConfig } from "./worker"; import { DefaultLogger, type Logger } from "./lib/logger"; import { SchemaManager } from "./schema-manager"; +import { + carrierForContext, + contextForSpan, + endSpan, + messagingAttributes, + setSpanAttribute, + setSpanError, + startSpan, + runWithSpan, +} from "./telemetry"; +import { SpanKind } from "@opentelemetry/api"; import type { EventDefinition, GenericDatabase, @@ -80,6 +91,8 @@ export type ConductorOptions< context: ExtraContext; logger?: Logger; + /** Disable OpenTelemetry instrumentation. The default is enabled and uses the global API provider. */ + telemetry?: false; }; // similar to inngest client @@ -142,6 +155,7 @@ export class Conductor< database?: TDatabaseSchema; context: TExtraContext; logger?: Logger; + telemetry?: false; }, ): Conductor { return new Conductor(options); @@ -234,6 +248,7 @@ export class Conductor< options.config, this.options.context, this.options.events?.definitions ?? [], + this.options.telemetry !== false, ); } @@ -303,6 +318,15 @@ export class Conductor< const queue = task.queue || "default"; if (Array.isArray(payloadOrItems)) { + const producer = + this.options.telemetry === false + ? null + : startSpan( + `send ${queue}`, + SpanKind.PRODUCER, + messagingAttributes(taskName, queue, "send"), + ); + const carrier = producer ? carrierForContext(contextForSpan(producer)) : null; const specs = payloadOrItems.map((item) => ({ task_key: taskName, queue, @@ -314,16 +338,47 @@ export class Conductor< cron_expression: item.cron_expression, priority: item.priority, group: item.group, + trace_context: carrier, })); - return this.db.invokeBatch(specs); + try { + const ids = await runWithSpan(producer, () => this.db.invokeBatch(specs)); + setSpanAttribute(producer, "messaging.batch.message_count", payloadOrItems.length); + return ids; + } catch (error) { + setSpanError(producer, error); + throw error; + } finally { + endSpan(producer); + } } - return this.db.invoke({ - task_key: taskName, - queue, - payload: payloadOrItems, - ...opts, - }); + const producer = + this.options.telemetry === false + ? null + : startSpan( + `send ${queue}`, + SpanKind.PRODUCER, + messagingAttributes(taskName, queue, "send"), + ); + const carrier = producer ? carrierForContext(contextForSpan(producer)) : null; + try { + const id = await runWithSpan(producer, () => + this.db.invoke({ + task_key: taskName, + queue, + payload: payloadOrItems, + ...opts, + trace_context: carrier, + }), + ); + if (id) setSpanAttribute(producer, "messaging.message.id", id); + return id; + } catch (error) { + setSpanError(producer, error); + throw error; + } finally { + endSpan(producer); + } } /** diff --git a/packages/pgconductor-js/src/database-client.ts b/packages/pgconductor-js/src/database-client.ts index df24d05..3c22472 100644 --- a/packages/pgconductor-js/src/database-client.ts +++ b/packages/pgconductor-js/src/database-client.ts @@ -21,6 +21,7 @@ import { type EmitEventArgs, } from "./query-builder"; import { makeChildLogger, type Logger } from "./lib/logger"; +import type { TraceContextCarrier } from "./internal-types"; export type JsonValue = string | number | boolean | null | Payload | JsonValue[]; export type Payload = { [key: string]: JsonValue }; @@ -39,6 +40,7 @@ export interface ExecutionSpec { parent_execution_id?: string | null; parent_step_key?: string | null; parent_timeout_ms?: number | null; + trace_context?: TraceContextCarrier | null; } export interface TaskSpec { @@ -82,6 +84,7 @@ export interface Execution { dead_letter_error?: string | null; dead_letter_attempts?: number | null; dead_letter_failed_at?: Date | null; + trace_context?: TraceContextCarrier | null; } // todo: move all of this to query-builder too or create new types.ts file diff --git a/packages/pgconductor-js/src/generated/sql.ts b/packages/pgconductor-js/src/generated/sql.ts index eeb078c..85b49e7 100644 --- a/packages/pgconductor-js/src/generated/sql.ts +++ b/packages/pgconductor-js/src/generated/sql.ts @@ -108,6 +108,7 @@ create table pgconductor._private_executions ( failed_at timestamptz, completed_at timestamptz, payload jsonb, + trace_context jsonb, run_at timestamptz default pgconductor._private_current_time() not null, locked_at timestamptz, locked_by uuid, @@ -360,7 +361,8 @@ create type pgconductor.execution_spec as ( dedupe_next_slot boolean, cron_expression text, priority integer, - "group" text + "group" text, + trace_context jsonb ); create type pgconductor.task_spec as ( @@ -619,6 +621,7 @@ begin task_key, queue, payload, + trace_context, run_at, dedupe_key, singleton_on, @@ -631,6 +634,7 @@ begin spec.task_key, coalesce(spec.queue, 'default'), spec.payload, + spec.trace_context, coalesce(spec.run_at, v_now), spec.dedupe_key, case @@ -668,7 +672,8 @@ create or replace function pgconductor.invoke( p_dedupe_next_slot boolean default false, p_cron_expression text default null, p_priority integer default null, - p_group text default null + p_group text default null, + p_trace_context jsonb default null ) returns table(id uuid) language plpgsql @@ -724,6 +729,7 @@ begin task_key, queue, payload, + trace_context, run_at, dedupe_key, singleton_on, @@ -735,6 +741,7 @@ begin p_task_key, p_queue, p_payload, + p_trace_context, v_run_at, p_dedupe_key, v_singleton_on, @@ -757,6 +764,7 @@ begin task_key, queue, payload, + trace_context, run_at, dedupe_key, singleton_on, @@ -768,6 +776,7 @@ begin p_task_key, p_queue, p_payload, + p_trace_context, v_next_singleton_on, p_dedupe_key, v_next_singleton_on, @@ -779,6 +788,7 @@ begin where singleton_on is not null and completed_at is null and failed_at is null and cancelled = false do update set payload = excluded.payload, + trace_context = excluded.trace_context, run_at = excluded.run_at, priority = excluded.priority, cron_expression = excluded.cron_expression, @@ -794,6 +804,7 @@ begin task_key, queue, payload, + trace_context, run_at, dedupe_key, cron_expression, @@ -804,6 +815,7 @@ begin p_task_key, p_queue, p_payload, + p_trace_context, v_run_at, p_dedupe_key, p_cron_expression, @@ -812,6 +824,7 @@ begin ) on conflict (task_key, dedupe_key, queue) do update set payload = excluded.payload, + trace_context = excluded.trace_context, run_at = excluded.run_at, priority = excluded.priority, cron_expression = excluded.cron_expression, diff --git a/packages/pgconductor-js/src/internal-types.ts b/packages/pgconductor-js/src/internal-types.ts new file mode 100644 index 0000000..9f5ea3c --- /dev/null +++ b/packages/pgconductor-js/src/internal-types.ts @@ -0,0 +1,6 @@ +/** The serializable W3C carrier persisted with an execution. */ +export type TraceContextCarrier = { + readonly [key: string]: string | undefined; + readonly traceparent: string; + readonly tracestate?: string; +}; diff --git a/packages/pgconductor-js/src/orchestrator.ts b/packages/pgconductor-js/src/orchestrator.ts index 2488dc9..d9b38c4 100644 --- a/packages/pgconductor-js/src/orchestrator.ts +++ b/packages/pgconductor-js/src/orchestrator.ts @@ -44,7 +44,7 @@ export class Orchestrator { private readonly schemaManager: SchemaManager; private readonly logger: Logger; - private heartbeatTimer: Timer | null = null; + private heartbeatTimer: ReturnType | null = null; private _stopDeferred: Deferred | null = null; private _startDeferred: Deferred | null = null; private _abortController: AbortController | null = null; @@ -71,6 +71,7 @@ export class Orchestrator { options.defaultWorker, options.conductor.options.context, options.conductor.options.events?.definitions ?? [], + options.conductor.options.telemetry !== false, ); this.workers.push(worker); } diff --git a/packages/pgconductor-js/src/query-builder.ts b/packages/pgconductor-js/src/query-builder.ts index dc6c2e6..b4f4b6a 100644 --- a/packages/pgconductor-js/src/query-builder.ts +++ b/packages/pgconductor-js/src/query-builder.ts @@ -1,5 +1,6 @@ import type { PendingQuery, Row, RowList, Sql } from "postgres"; import type { GroupedExecutionResults } from "./database-client"; +import { boundedCarrier } from "./telemetry"; import * as assert from "./lib/assert"; import type { Execution, @@ -357,15 +358,14 @@ export class QueryBuilder { where e.id = c.id and e.queue = ${queueName}::text and e.is_available = true returning e.id, e.task_key, e.queue, e.payload, e.waiting_on_execution_id, e.waiting_step_key, e.cancelled, e.last_error, e.dedupe_key, e.cron_expression, - e.locked_by, e."group", e.priority, e.run_at, e.created_at, - e.dead_letter_source_execution_id, e.dead_letter_source_queue, - e.dead_letter_source_task_key, e.dead_letter_error, - e.dead_letter_attempts, e.dead_letter_failed_at + e.locked_by, e."group", e.trace_context, e.dead_letter_source_execution_id, + e.dead_letter_source_queue, e.dead_letter_source_task_key, e.dead_letter_error, + e.dead_letter_attempts, e.dead_letter_failed_at, e.priority, e.run_at, e.created_at ) select id, task_key, queue, payload, waiting_on_execution_id, waiting_step_key, - cancelled, last_error, dedupe_key, cron_expression, locked_by, "group", - dead_letter_source_execution_id, dead_letter_source_queue, - dead_letter_source_task_key, dead_letter_error, dead_letter_attempts, dead_letter_failed_at + cancelled, last_error, dedupe_key, cron_expression, locked_by, "group", trace_context, + dead_letter_source_execution_id, dead_letter_source_queue, dead_letter_source_task_key, + dead_letter_error, dead_letter_attempts, dead_letter_failed_at from claimed order by priority asc, run_at asc, created_at asc, id asc `; @@ -796,6 +796,7 @@ export class QueryBuilder { p_dedupe_next_slot := ${dedupe_next_slot}::boolean, p_cron_expression := ${spec.cron_expression || null}::text, p_priority := ${spec.priority || null}::integer, + p_trace_context := ${boundedCarrier(spec.trace_context) ? this.sql.json(boundedCarrier(spec.trace_context)) : null}::jsonb, p_group := ${spec.group || null}::text ) `; @@ -902,6 +903,7 @@ export class QueryBuilder { dedupe_seconds, dedupe_next_slot, cron_expression: spec.cron_expression || null, + trace_context: boundedCarrier(spec.trace_context), priority: spec.priority, group: spec.group || null, }; @@ -910,7 +912,7 @@ export class QueryBuilder { return this.sql<{ id: string }[]>` select id from pgconductor.invoke_batch( array( - select jsonb_populate_recordset(null::pgconductor.execution_spec, ${this.sql.json(specsArray)}::jsonb) + select jsonb_populate_recordset(null::pgconductor.execution_spec, ${this.sql.json(JSON.parse(JSON.stringify(specsArray)))}::jsonb) ) ) `; diff --git a/packages/pgconductor-js/src/task-context.ts b/packages/pgconductor-js/src/task-context.ts index 035557e..bc2f3f5 100644 --- a/packages/pgconductor-js/src/task-context.ts +++ b/packages/pgconductor-js/src/task-context.ts @@ -18,6 +18,8 @@ import type { Logger } from "./lib/logger"; import { WindowChecker } from "./lib/window-checker"; import { TypedAbortController } from "./lib/typed-abort-controller"; import { parseDuration, type DurationInput } from "./lib/duration"; +import { SpanKind } from "@opentelemetry/api"; +import { endSpan, stepAttributes, runWithSpan, setSpanError, startSpan } from "./telemetry"; import type { EventDefinition, EventName, @@ -97,6 +99,7 @@ export type TaskContextOptions = { execution: Execution; logger: Logger; window?: [string, string]; + telemetry?: false; }; type ScheduleOptions = { @@ -203,22 +206,38 @@ export class TaskContext< return (cached as { result: T }).result; } - // Execute and save - const result = await fn(); - - await this.opts.db.saveStep( - { - executionId: this.opts.execution.id, - queue: this.opts.execution.queue, - orchestratorId: this.opts.execution.locked_by, - key: name, - result: { result: result as JsonValue }, - runAtMs: undefined, - }, - { signal: this.signal }, - ); - - return result; + // Execute and save. A cache hit intentionally has no span. + const span = + this.opts.telemetry === false + ? null + : startSpan( + `step ${this.opts.execution.task_key}`, + SpanKind.INTERNAL, + stepAttributes(this.opts.execution.task_key, this.opts.execution.queue, name), + ); + try { + const result = await runWithSpan(span, fn); + + await runWithSpan(span, () => + this.opts.db.saveStep( + { + executionId: this.opts.execution.id, + queue: this.opts.execution.queue, + orchestratorId: this.opts.execution.locked_by, + key: name, + result: { result: result as JsonValue }, + runAtMs: undefined, + }, + { signal: this.signal }, + ), + ); + endSpan(span); + return result; + } catch (error) { + setSpanError(span, error); + endSpan(span); + throw error; + } } async checkpoint(): Promise { diff --git a/packages/pgconductor-js/src/telemetry.ts b/packages/pgconductor-js/src/telemetry.ts new file mode 100644 index 0000000..c734fa3 --- /dev/null +++ b/packages/pgconductor-js/src/telemetry.ts @@ -0,0 +1,252 @@ +import { + context, + propagation, + ROOT_CONTEXT, + SpanKind, + SpanStatusCode, + trace, + type Context, + type Link, + type Span, + type SpanContext, +} from "@opentelemetry/api"; +import type { TraceContextCarrier } from "./internal-types"; + +export type { TraceContextCarrier } from "./internal-types"; + +const TRACECONTEXT_SIZE_LIMIT = 1024; +const TRACEPARENT = /^00-([0-9a-f]{32})-([0-9a-f]{16})-([0-9a-f]{2})$/; +// W3C simple keys and multi-tenant vendor keys (for example tenant@vendor). +const TRACESTATE_KEY = + /^(?:[a-z][a-z0-9_*/-]{0,255}|[a-z0-9][a-z0-9_*/-]{0,240}@[a-z][a-z0-9_*/-]{0,13})$/; +const TRACESTATE_VALUE = /^[\x20-\x2b\x2d-\x3c\x3e-\x7e]{1,256}$/; +const INSTRUMENTATION_NAME = "pgconductor-js"; +const INSTRUMENTATION_VERSION = "0.1.0"; +const ZERO_TRACE_ID = "00000000000000000000000000000000"; +const ZERO_SPAN_ID = "0000000000000000"; + +type MessagingAttributes = Record; + +type MessagingOperation = "send" | "process" | "settle"; + +export const messagingAttributes = ( + taskKey: string | undefined, + queue: string, + operation: MessagingOperation, + messageId?: string, + batchMessageCount?: number, +): MessagingAttributes => ({ + "messaging.system": "postgres_conductor", + "messaging.destination.name": queue, + "messaging.operation.name": operation, + "messaging.operation.type": operation, + ...(taskKey === undefined ? {} : { "pgconductor.task.name": taskKey }), + ...(messageId === undefined ? {} : { "messaging.message.id": messageId }), + ...(batchMessageCount === undefined + ? {} + : { "messaging.batch.message_count": batchMessageCount }), + "pgconductor.queue": queue, +}); + +export const stepAttributes = ( + taskKey: string, + queue: string, + stepKey: string, +): MessagingAttributes => ({ + "pgconductor.task.name": taskKey, + "pgconductor.step.name": stepKey, + "pgconductor.queue": queue, +}); + +function safe(fn: () => T, fallback: T): T { + try { + return fn(); + } catch { + return fallback; + } +} + +function property(value: object, key: string): unknown { + return safe(() => Object.getOwnPropertyDescriptor(value, key)?.value, undefined); +} + +function isValidTraceparent(value: unknown): value is string { + return safe(() => { + if (typeof value !== "string" || !TRACEPARENT.test(value)) return false; + const match = TRACEPARENT.exec(value); + const traceId = match?.[1]; + const spanId = match?.[2]; + const traceFlags = match?.[3]; + return ( + !!traceId && + !!spanId && + !!traceFlags && + trace.isSpanContextValid({ + traceId, + spanId, + traceFlags: Number.parseInt(traceFlags, 16), + }) + ); + }, false); +} + +function isValidTracestate(value: unknown): value is string { + return safe(() => { + if (typeof value !== "string" || value.length > 512) return false; + const members = value.split(",").map((member) => member.trim()); + if (members.length > 32 || members.some((member) => member.length === 0)) return false; + const keys = new Set(); + for (const member of members) { + const separator = member.indexOf("="); + if (separator <= 0) return false; + const key = member.slice(0, separator); + const memberValue = member.slice(separator + 1); + if (!TRACESTATE_KEY.test(key) || !TRACESTATE_VALUE.test(memberValue) || keys.has(key)) + return false; + keys.add(key); + } + return true; + }, false); +} + +function carrierSize(traceparent: string, tracestate?: string): number { + return Buffer.byteLength( + JSON.stringify({ traceparent, ...(tracestate === undefined ? {} : { tracestate }) }), + "utf8", + ); +} + +/** Strictly validate a carrier without sanitizing malformed tracestate. */ +export function isValidCarrier(value: unknown): value is TraceContextCarrier { + return safe(() => { + if (!value || typeof value !== "object") return false; + const traceparent = property(value, "traceparent"); + const tracestate = property(value, "tracestate"); + return ( + isValidTraceparent(traceparent) && + (tracestate === undefined || isValidTracestate(tracestate)) && + carrierSize(traceparent, typeof tracestate === "string" ? tracestate : undefined) <= + TRACECONTEXT_SIZE_LIMIT + ); + }, false); +} + +/** Sanitize a value read from the database before giving it to OTel or persisting it. */ +export function boundedCarrier(value: unknown): TraceContextCarrier | null { + if (!value || typeof value !== "object") return null; + const traceparent = property(value, "traceparent"); + if (!isValidTraceparent(traceparent)) return null; + + const candidateTracestate = property(value, "tracestate"); + const tracestate = + typeof candidateTracestate === "string" && + isValidTracestate(candidateTracestate) && + carrierSize(traceparent, candidateTracestate) <= TRACECONTEXT_SIZE_LIMIT + ? candidateTracestate + : undefined; + return { traceparent, ...(tracestate === undefined ? {} : { tracestate }) }; +} + +export function extractCarrier(value: unknown): Context | null { + const carrier = boundedCarrier(value); + if (!carrier) return null; + return safe(() => propagation.extract(ROOT_CONTEXT, carrier), null); +} + +export function startSpan( + name: string, + kind: SpanKind, + attributes: MessagingAttributes, + parent?: Context | null, + links: Link[] = [], +): Span | null { + return safe( + () => + trace.getTracer(INSTRUMENTATION_NAME, INSTRUMENTATION_VERSION).startSpan( + name, + { + kind, + attributes, + links, + }, + // null is intentional: an absent carrier is a true root, rather than + // accidentally inheriting a caller's active span. + parent === undefined ? context.active() : parent || ROOT_CONTEXT, + ), + null, + ); +} + +export function linkContext(ctx: Context | null): Link | null { + if (!ctx) return null; + const value = safe(() => trace.getSpanContext(ctx), null); + return value && value.traceId !== ZERO_TRACE_ID && value.spanId !== ZERO_SPAN_ID + ? { context: value } + : null; +} + +export function spanContext(span: Span | null): SpanContext | null { + if (!span) return null; + return safe(() => { + const value = span.spanContext(); + return value.traceId !== ZERO_TRACE_ID && value.spanId !== ZERO_SPAN_ID ? value : null; + }, null); +} + +export function contextForSpan(span: Span | null, parent: Context = context.active()): Context { + return span ? safe(() => trace.setSpan(parent, span), parent) : parent; +} + +export function runWithSpan(span: Span | null, fn: () => T): T { + if (!span) return fn(); + const active = safe(() => context.active(), ROOT_CONTEXT); + const scoped = safe(() => trace.setSpan(active, span), active); + let callbackStarted = false; + try { + return context.with(scoped, () => { + callbackStarted = true; + return fn(); + }); + } catch (error) { + // A broken context manager may throw before invoking the callback. Only then + // retry without a scope; retrying after start would run user work twice. + if (callbackStarted) throw error; + return fn(); + } +} + +export function endSpan(span: Span | null, error?: unknown): void { + if (!span) return; + if (error) setSpanError(span, error); + safe(() => span.end(), undefined); +} + +export function setSpanError(span: Span | null, error: unknown): void { + if (!span) return; + safe(() => { + span.recordException(error instanceof Error ? error : new Error(String(error))); + span.setStatus({ code: SpanStatusCode.ERROR }); + }, undefined); +} + +export function setSpanAttribute( + span: Span | null, + key: string, + value: string | number | boolean, +): void { + if (!span) return; + safe(() => span.setAttribute(key, value), undefined); +} + +export function carrierForContext(ctx: Context = context.active()): TraceContextCarrier | null { + return safe(() => { + const carrier: Record = {}; + propagation.inject(ctx, carrier); + return boundedCarrier(carrier); + }, null); +} + +export const telemetryConstants = { + instrumentationName: INSTRUMENTATION_NAME, + instrumentationVersion: INSTRUMENTATION_VERSION, +}; diff --git a/packages/pgconductor-js/src/worker.ts b/packages/pgconductor-js/src/worker.ts index efd27e8..b20e780 100644 --- a/packages/pgconductor-js/src/worker.ts +++ b/packages/pgconductor-js/src/worker.ts @@ -31,6 +31,18 @@ import { makeChildLogger, type Logger } from "./lib/logger"; import type { EventDefinition } from "./event-definition"; import { coerceError } from "./lib/coerce-error"; import type { TypedAbortController } from "./lib/typed-abort-controller"; +import { ROOT_CONTEXT, SpanKind, type SpanContext } from "@opentelemetry/api"; +import { + endSpan, + extractCarrier, + linkContext, + messagingAttributes, + runWithSpan, + setSpanAttribute, + setSpanError, + spanContext, + startSpan, +} from "./telemetry"; /** * The configuration options for the Worker. @@ -169,6 +181,7 @@ export class Worker< private _drainDidWork = false; private _runningTasks = new Map>(); private eventProcessingGate: Promise = Promise.resolve(); + private readonly processSpanContexts = new Map(); /** Used by Orchestrator to prevent local event fan-out before every worker registers. */ setEventProcessingGate(gate: Promise): void { @@ -183,6 +196,7 @@ export class Worker< config: Partial = {}, private readonly extraContext: object = {}, private readonly eventDefinitions: readonly EventDefinition[] = [], + private readonly telemetry = true, ) { const maintenanceTask = createMaintenanceTask(this.queueName); this.tasks = tasks.reduce( @@ -375,6 +389,7 @@ export class Worker< this._abortController = null; this.orchestratorId = null; this.eventProcessingGate = Promise.resolve(); + this.processSpanContexts.clear(); } private async processEventBatches({ runOnce }: { runOnce: boolean }): Promise { @@ -747,6 +762,16 @@ export class Worker< resolve(taskAbortController.signal.reason); }); }); + const consumer = this.telemetry + ? startSpan( + `process ${exec.queue}`, + SpanKind.CONSUMER, + messagingAttributes(exec.task_key, exec.queue, "process", exec.id), + extractCarrier(exec.trace_context) || ROOT_CONTEXT, + ) + : null; + const consumerContext = spanContext(consumer); + if (consumerContext) this.processSpanContexts.set(exec.id, consumerContext); try { await this.scheduleNextExecution(exec); @@ -779,31 +804,35 @@ export class Worker< ? { ...this.extraContext, db: this.db, tasks: this.tasks } : this.extraContext; - const output = await Promise.race([ - task.execute( - taskEvent, - TaskContext.create( - { - db: this.db, - abortController: taskAbortController, - execution: exec, - logger: makeChildLogger(this.logger, { - execution_id: exec.id, - orchestrator_id: exec.locked_by, - task_key: exec.task_key, - queue: exec.queue, - }), - window: task.window, - }, - extraContext, + const output = await runWithSpan(consumer, () => + Promise.race([ + task.execute( + taskEvent, + TaskContext.create( + { + db: this.db, + abortController: taskAbortController, + execution: exec, + logger: makeChildLogger(this.logger, { + execution_id: exec.id, + orchestrator_id: exec.locked_by, + task_key: exec.task_key, + queue: exec.queue, + }), + window: task.window, + telemetry: this.telemetry ? undefined : false, + }, + extraContext, + ), ), - ), - abortPromise, - ]); + abortPromise, + ]), + ); if (isTaskAbortReason(output)) { switch (output.reason) { case "wait-for-event": + this.processSpanContexts.delete(exec.id); return null; case "child-invocation": return { @@ -853,6 +882,7 @@ export class Worker< result: output, } as const; } catch (err) { + setSpanError(consumer, err); return { execution_id: exec.id, orchestrator_id: exec.locked_by, @@ -862,6 +892,7 @@ export class Worker< error: coerceError(err).message, } as const; } finally { + endSpan(consumer); // Clean up running task tracking this._runningTasks.delete(exec.id); } @@ -879,6 +910,22 @@ export class Worker< taskKey: string, executions: Execution[], ): Promise { + const links = this.telemetry + ? executions.flatMap((exec) => linkContext(extractCarrier(exec.trace_context)) || []) + : []; + const consumer = this.telemetry + ? startSpan( + `process ${this.queueName}`, + SpanKind.CONSUMER, + messagingAttributes(taskKey, this.queueName, "process", undefined, executions.length), + ROOT_CONTEXT, + links, + ) + : null; + const processContext = spanContext(consumer); + if (processContext) + for (const execution of executions) + this.processSpanContexts.set(execution.id, processContext); // Build event array const events = executions.map((exec) => { if (exec.cron_expression) { @@ -921,7 +968,9 @@ export class Worker< // Schedule next executions for cron tasks await Promise.all(executions.map((exec) => this.scheduleNextExecution(exec))); - const result = await Promise.race([task.execute(events, batchContext), abortPromise]); + const result = await runWithSpan(consumer, () => + Promise.race([task.execute(events, batchContext), abortPromise]), + ); // Handle abort reasons if (isTaskAbortReason(result)) { @@ -982,6 +1031,7 @@ export class Worker< result: result[i], })); } catch (err) { + setSpanError(consumer, err); // Handler threw: all fail together const errorMsg = coerceError(err).message; return executions.map((exec) => ({ @@ -992,6 +1042,8 @@ export class Worker< status: "failed" as const, error: errorMsg, })); + } finally { + endSpan(consumer); } } @@ -1032,7 +1084,7 @@ export class Worker< // --- Stage 3: Flush results to database --- private async flushResults(source: AsyncIterable): Promise { let buffer = new BufferState(); - let flushTimer: Timer | null = null; + let flushTimer: ReturnType | null = null; const flushNow = async (isCleanup = false) => { if (buffer.count === 0) return; @@ -1045,14 +1097,47 @@ export class Worker< flushTimer = null; } + const links = this.telemetry + ? Array.from( + new Set( + [...batch.completed, ...batch.failed, ...batch.released, ...batch.invokeChild] + .map((result) => this.processSpanContexts.get(result.execution_id)) + .filter((value): value is SpanContext => value !== undefined), + ), + ).map((context) => ({ context })) + : []; + const settle = + this.telemetry && links.length + ? startSpan( + `settle ${this.queueName}`, + SpanKind.CLIENT, + messagingAttributes(undefined, this.queueName, "settle", undefined, batch.count), + ROOT_CONTEXT, + links, + ) + : null; + let settled = false; try { batch.orchestratorId = this.orchestratorId || batch.orchestratorId; - await this.db.returnExecutions(batch, { signal: this.signal }); + await runWithSpan(settle, () => this.db.returnExecutions(batch, { signal: this.signal })); + settled = true; + setSpanAttribute(settle, "pgconductor.db.commit.status", "success"); } catch (err) { + setSpanError(settle, err); this.logger.error("Error flushing results:", err); if (!isCleanup) { buffer.restore(batch); } + } finally { + endSpan(settle); + if (settled || isCleanup) + for (const result of [ + ...batch.completed, + ...batch.failed, + ...batch.released, + ...batch.invokeChild, + ]) + this.processSpanContexts.delete(result.execution_id); } }; diff --git a/packages/pgconductor-js/tests/integration/telemetry.test.ts b/packages/pgconductor-js/tests/integration/telemetry.test.ts new file mode 100644 index 0000000..5ada2e7 --- /dev/null +++ b/packages/pgconductor-js/tests/integration/telemetry.test.ts @@ -0,0 +1,380 @@ +import { describe, expect, test } from "bun:test"; +import { AsyncLocalStorage } from "node:async_hooks"; +import { + context, + propagation, + ROOT_CONTEXT, + SpanKind, + SpanStatusCode, + trace, + type Context, + type ContextManager, + type Span, +} from "@opentelemetry/api"; +import { + BasicTracerProvider, + InMemorySpanExporter, + SimpleSpanProcessor, +} from "@opentelemetry/sdk-trace-base"; +import type { Sql } from "postgres"; +import { Conductor } from "../../src/conductor"; +import { Orchestrator } from "../../src/orchestrator"; +import { Task } from "../../src/task"; +import { Worker } from "../../src/worker"; +import { DefaultLogger } from "../../src/lib/logger"; +import { InMemoryDatabaseClient } from "../mocks/in-memory-database-client"; +import type { DatabaseClient } from "../../src/database-client"; +import { + boundedCarrier, + carrierForContext, + contextForSpan, + extractCarrier, + messagingAttributes, + startSpan, + endSpan, +} from "../../src/telemetry"; + +const logger = new DefaultLogger(); + +// sdk-trace-base deliberately does not install a context manager. Use the +// platform async context manager so worker/user-child assertions exercise the +// same active-span behavior as a Node application. +class TestContextManager implements ContextManager { + private readonly storage = new AsyncLocalStorage(); + active() { + return this.storage.getStore() || ROOT_CONTEXT; + } + with(ctx: Context, fn: (...args: any[]) => T, thisArg?: any, ...args: any[]) { + return this.storage.run(ctx, () => fn.apply(thisArg, args)); + } + bind(ctx: Context, target: T): T { + return target; + } + enable() { + return this; + } + disable() { + this.storage.disable(); + return this; + } +} + +const fakeSql = Object.assign((async () => [{ id: "fake-execution" }]) as unknown as Sql, { + json: (value: unknown) => JSON.stringify(value), +}); + +function resetGlobals() { + context.disable(); + propagation.disable(); + trace.disable(); +} + +function installProvider() { + resetGlobals(); + const exporter = new InMemorySpanExporter(); + const provider = new BasicTracerProvider(); + provider.addSpanProcessor(new SimpleSpanProcessor(exporter)); + provider.register(); + context.setGlobalContextManager(new TestContextManager()); + return { exporter, provider }; +} + +async function cleanup(provider: BasicTracerProvider) { + await provider.shutdown(); + resetGlobals(); +} + +function spans(exporter: InMemorySpanExporter, name: string) { + return exporter.getFinishedSpans().filter((span) => span.name === name); +} + +function makeTask( + name: string, + execute: (event: any, ctx: any) => Promise, + config: Record = {}, +) { + return Task.create({ name, ...config } as any, { invocable: true } as any, execute as any); +} + +function makeWorker(db: InMemoryDatabaseClient, task: any, telemetry = true) { + return new Worker( + "default", + [task], + db as unknown as DatabaseClient, + logger, + { pollIntervalMs: 1, flushIntervalMs: 1, fetchBatchSize: 10, flushBatchSize: 10 }, + {}, + [], + telemetry, + ); +} + +const producerCarrier = (name: string) => { + const producer = startSpan( + name, + SpanKind.PRODUCER, + messagingAttributes("task", "default", "send"), + ); + const carrier = carrierForContext(contextForSpan(producer)); + endSpan(producer); + return { carrier, producer }; +}; + +describe.serial("OpenTelemetry instrumentation", () => { + test("a Conductor made before provider registration uses the later global provider", async () => { + const conductor = Conductor.create({ sql: fakeSql, context: {} }); + const { exporter, provider } = installProvider(); + await conductor.invoke({ name: "later-provider" }, { value: "not an attribute" } as any); + expect(spans(exporter, "send default")).toHaveLength(1); + await cleanup(provider); + }); + + test("telemetry false emits no spans and persists null carrier", async () => { + const { exporter, provider } = installProvider(); + const db = new InMemoryDatabaseClient(); + const conductor = Conductor.create({ sql: fakeSql, context: {}, telemetry: false }); + await conductor.invoke({ name: "disabled" }, { secret: "payload" } as any); + const id = await db.invoke({ + task_key: "disabled", + queue: "default", + payload: {}, + trace_context: null, + }); + expect(exporter.getFinishedSpans()).toHaveLength(0); + expect(db.getExecution(id!)?.trace_context).toBeNull(); + await cleanup(provider); + }); + + test("the orchestrator passes telemetry opt-out to its default worker", async () => { + const { exporter, provider } = installProvider(); + const db = new InMemoryDatabaseClient(); + const conductor = Conductor.create({ sql: fakeSql, context: {}, telemetry: false }); + (conductor as any).db = db; + const task = makeTask("orchestrator-disabled", async () => undefined); + const id = (await conductor.invoke({ name: "orchestrator-disabled" }, { + secret: "payload", + } as any)) as unknown as string; + const orchestrator = Orchestrator.create({ + conductor, + tasks: [task], + defaultWorker: { pollIntervalMs: 1, flushIntervalMs: 1 }, + } as any); + + await orchestrator.drain(); + expect(db.getExecution(id)?.state).toBe("completed"); + expect(exporter.getFinishedSpans()).toHaveLength(0); + expect(db.getExecution(id)?.trace_context).toBeNull(); + await cleanup(provider); + }); + + test("producer carrier becomes the real worker consumer parent", async () => { + const { exporter, provider } = installProvider(); + const db = new InMemoryDatabaseClient(); + const { carrier, producer } = producerCarrier("send-parent"); + await db.invoke({ + task_key: "parented", + queue: "default", + payload: {}, + trace_context: carrier, + }); + const task = makeTask("parented", async () => undefined); + await makeWorker(db, task).drain("worker"); + const consumer = spans(exporter, "process default")[0]!; + expect(consumer.parentSpanId).toBe(producer!.spanContext().spanId); + expect(consumer.spanContext().traceId).toBe(producer!.spanContext().traceId); + await cleanup(provider); + }); + + test("batch consumer links every producer context", async () => { + const { exporter, provider } = installProvider(); + const db = new InMemoryDatabaseClient(); + const first = producerCarrier("batch-one"); + const second = producerCarrier("batch-two"); + await db.invokeBatch([ + { task_key: "batched", queue: "default", payload: {}, trace_context: first.carrier }, + { task_key: "batched", queue: "default", payload: {}, trace_context: second.carrier }, + ]); + const task = makeTask("batched", async () => [], { batch: { size: 10, timeoutMs: 1 } }); + await makeWorker(db, task).drain("worker"); + const consumer = spans(exporter, "process default")[0]!; + expect(consumer.links.map((link) => link.context.spanId)).toEqual( + expect.arrayContaining([ + first.producer?.spanContext().spanId, + second.producer?.spanContext().spanId, + ]), + ); + await cleanup(provider); + }); + + test("retry attempts have finite, separate process spans", async () => { + const { exporter, provider } = installProvider(); + const db = new InMemoryDatabaseClient(); + let attempts = 0; + const task = makeTask("retry-span", async () => { + if (++attempts === 1) throw new Error("try again"); + }); + (task as any).maxAttempts = 2; + const id = await db.invoke({ task_key: "retry-span", queue: "default", payload: {} }); + const worker = makeWorker(db, task); + await worker.drain("worker"); + db.advanceTime(16000); + await worker.drain("worker"); + const processSpans = spans(exporter, "process default"); + expect(processSpans).toHaveLength(2); + expect(new Set(processSpans.map((span) => span.spanContext().spanId)).size).toBe(2); + expect(db.getExecution(id!)?.state).toBe("completed"); + await cleanup(provider); + }); + + test("a user-created active child is a child of process", async () => { + const { exporter, provider } = installProvider(); + const db = new InMemoryDatabaseClient(); + const task = makeTask("active-child", async () => { + const child = trace.getTracer("user").startSpan("user child"); + child.end(); + }); + await db.invoke({ task_key: "active-child", queue: "default", payload: {} }); + await makeWorker(db, task).drain("worker"); + const parent = spans(exporter, "process default")[0]!; + const child = spans(exporter, "user child")[0]!; + expect(child.parentSpanId).toBe(parent.spanContext().spanId); + await cleanup(provider); + }); + + test("settle has process links and records database errors", async () => { + const { exporter, provider } = installProvider(); + const db = new InMemoryDatabaseClient(); + await db.invoke({ task_key: "settle-error", queue: "default", payload: {} }); + const original = db.returnExecutions.bind(db); + (db as any).returnExecutions = async () => { + throw new Error("database unavailable"); + }; + await makeWorker( + db, + makeTask("settle-error", async () => undefined), + ).drain("worker"); + const settle = spans(exporter, "settle default")[0]!; + expect(settle.links.length).toBe(1); + expect(settle.status.code).toBe(SpanStatusCode.ERROR); + (db as any).returnExecutions = original; + await cleanup(provider); + }); + + test("step callback is one span and cached replay creates none", async () => { + const { exporter, provider } = installProvider(); + const db = new InMemoryDatabaseClient(); + let attempts = 0; + const task = makeTask("step-cache", async (_event, ctx) => { + const value = await ctx.step("once", async () => 42); + if (++attempts === 1) throw new Error("retry"); + return { value }; + }); + await db.registerWorker({ + queueName: "default", + taskSpecs: [{ key: "step-cache", maxAttempts: 2 }], + cronSchedules: [], + eventSubscriptions: [], + }); + await db.invoke({ task_key: "step-cache", queue: "default", payload: {} }); + const worker = makeWorker(db, task); + await worker.drain("worker"); + db.advanceTime(16000); + await worker.drain("worker"); + expect(spans(exporter, "step step-cache")).toHaveLength(1); + await cleanup(provider); + }); + + test("cron and event executions are root process spans", async () => { + const { exporter, provider } = installProvider(); + const db = new InMemoryDatabaseClient(); + await db.invoke({ + task_key: "root-cron", + queue: "default", + payload: {}, + cron_expression: "* * * * *", + dedupe_key: "scheduled::nightly::1", + }); + await db.invoke({ + task_key: "root-event", + queue: "default", + payload: { event: "user.created", payload: { id: 1 } }, + }); + const worker = makeWorker( + db, + makeTask("root-cron", async () => undefined), + ); + // A second worker is unnecessary: use a task map with both definitions. + const eventWorker = new Worker( + "default", + [makeTask("root-cron", async () => undefined), makeTask("root-event", async () => undefined)], + db as unknown as DatabaseClient, + logger, + { pollIntervalMs: 1, flushIntervalMs: 1 }, + {}, + [], + true, + ); + await eventWorker.drain("worker"); + for (const name of ["process default", "process default"]) + expect(spans(exporter, name)[0]!.parentSpanId).toBeUndefined(); + void worker; + await cleanup(provider); + }); + + test("carrier sanitization preserves valid parents and bounds tracestate", async () => { + const { exporter, provider } = installProvider(); + const traceparent = `00-${"1".repeat(32)}-${"2".repeat(16)}-01`; + expect(() => extractCarrier({ traceparent: "not-valid" })).not.toThrow(); + expect(boundedCarrier({ traceparent, tracestate: "tenant@vendor=value" })).toEqual({ + traceparent, + tracestate: "tenant@vendor=value", + }); + expect(boundedCarrier({ traceparent, tracestate: "not a tracestate" })).toEqual({ + traceparent, + }); + expect(boundedCarrier({ traceparent, tracestate: "x".repeat(2000) })).toEqual({ traceparent }); + expect( + Buffer.byteLength( + JSON.stringify(boundedCarrier({ traceparent, tracestate: "x".repeat(2000) })), + "utf8", + ), + ).toBeLessThanOrEqual(1024); + const db = new InMemoryDatabaseClient(); + await db.invoke({ + task_key: "safe", + queue: "default", + payload: { secret: "do not record" }, + trace_context: { traceparent: "bad" } as any, + }); + await makeWorker( + db, + makeTask("safe", async () => undefined), + ).drain("worker"); + for (const span of exporter.getFinishedSpans()) { + expect(span.attributes).not.toHaveProperty("payload"); + expect(JSON.stringify(span.attributes)).not.toContain("do not record"); + } + await cleanup(provider); + }); + + test("hostile propagators fail open", async () => { + const { exporter, provider } = installProvider(); + propagation.setGlobalPropagator({ + inject() { + throw new Error("inject failed"); + }, + extract() { + throw new Error("extract failed"); + }, + fields() { + return []; + }, + } as any); + expect(() => carrierForContext(ROOT_CONTEXT)).not.toThrow(); + expect(() => + extractCarrier({ traceparent: `00-${"1".repeat(32)}-${"2".repeat(16)}-01` }), + ).not.toThrow(); + expect(exporter.getFinishedSpans()).toHaveLength(0); + await cleanup(provider); + }); +}); diff --git a/packages/pgconductor-js/tests/mocks/in-memory-database-client.ts b/packages/pgconductor-js/tests/mocks/in-memory-database-client.ts index 2a07142..67b5a28 100644 --- a/packages/pgconductor-js/tests/mocks/in-memory-database-client.ts +++ b/packages/pgconductor-js/tests/mocks/in-memory-database-client.ts @@ -69,6 +69,7 @@ interface StoredExecution { dead_letter_error: string | null; dead_letter_attempts: number | null; dead_letter_failed_at: Date | null; + trace_context: Execution["trace_context"]; } interface StoredStep { @@ -470,6 +471,7 @@ export class InMemoryDatabaseClient implements IDatabaseClient { dead_letter_error: exec.dead_letter_error, dead_letter_attempts: exec.dead_letter_attempts, dead_letter_failed_at: exec.dead_letter_failed_at, + trace_context: exec.trace_context, locked_by: exec.orchestrator_id || "", }); @@ -794,6 +796,7 @@ export class InMemoryDatabaseClient implements IDatabaseClient { dead_letter_error: null, dead_letter_attempts: null, dead_letter_failed_at: null, + trace_context: spec.trace_context || null, }; this.executions.set(id, execution); @@ -1381,6 +1384,7 @@ export class InMemoryDatabaseClient implements IDatabaseClient { dead_letter_error: error, dead_letter_attempts: exec.attempts, dead_letter_failed_at: now, + trace_context: null, }); }