From fec93a064b48e96b0831d842e6f0a619038c8a83 Mon Sep 17 00:00:00 2001 From: nicktrn <55853254+nicktrn@users.noreply.github.com> Date: Thu, 23 Jul 2026 16:39:08 +0100 Subject: [PATCH 1/4] feat(supervisor): add prometheus metric for outbound http requests --- .../supervisor-outbound-request-metrics.md | 6 ++++ apps/supervisor/src/index.ts | 35 ++++++++++++++++++- 2 files changed, 40 insertions(+), 1 deletion(-) create mode 100644 .server-changes/supervisor-outbound-request-metrics.md diff --git a/.server-changes/supervisor-outbound-request-metrics.md b/.server-changes/supervisor-outbound-request-metrics.md new file mode 100644 index 0000000000..8a510973c4 --- /dev/null +++ b/.server-changes/supervisor-outbound-request-metrics.md @@ -0,0 +1,6 @@ +--- +area: supervisor +type: improvement +--- + +The supervisor now reports a Prometheus metric for its outbound HTTP requests, so failed calls to upstream services are visible for monitoring. diff --git a/apps/supervisor/src/index.ts b/apps/supervisor/src/index.ts index cf73be90ee..221c7e2892 100644 --- a/apps/supervisor/src/index.ts +++ b/apps/supervisor/src/index.ts @@ -22,7 +22,7 @@ import { isKubernetesEnvironment, } from "@trigger.dev/core/v3/serverOnly"; import { createK8sApi, createApiserverMetricsFetcher } from "./clients/kubernetes.js"; -import { collectDefaultMetrics, Gauge, Histogram } from "prom-client"; +import { collectDefaultMetrics, Counter, Gauge, Histogram } from "prom-client"; import { register } from "./metrics.js"; import { PodCleaner } from "./services/podCleaner.js"; import { FailedPodHandler } from "./services/failedPodHandler.js"; @@ -60,6 +60,13 @@ const workloadCreateDuration = new Histogram({ registers: [register], }); +const outboundRequestsTotal = new Counter({ + name: "supervisor_outbound_request_total", + help: "Count of outbound HTTP requests from the supervisor, by target name, method, response status, and outcome (ok, http_error, invalid_response, network_error).", + labelNames: ["name", "method", "status", "outcome"], + registers: [register], +}); + class ManagedSupervisor { private readonly workerSession: SupervisorSession; private readonly metricsServer?: HttpServer; @@ -700,8 +707,15 @@ class ManagedSupervisor { }); if (!res.ok) { + outboundRequestsTotal.inc({ + name: "warm_start", + method: "POST", + status: String(res.status), + outcome: "http_error", + }); this.logger.error("Warm start failed", { runId: dequeuedMessage.run.id, + statusCode: res.status, }); return false; } @@ -710,6 +724,12 @@ class ManagedSupervisor { const parsedData = z.object({ didWarmStart: z.boolean() }).safeParse(data); if (!parsedData.success) { + outboundRequestsTotal.inc({ + name: "warm_start", + method: "POST", + status: String(res.status), + outcome: "invalid_response", + }); this.logger.error("Warm start response invalid", { runId: dequeuedMessage.run.id, data, @@ -717,8 +737,21 @@ class ManagedSupervisor { return false; } + outboundRequestsTotal.inc({ + name: "warm_start", + method: "POST", + status: String(res.status), + outcome: "ok", + }); + return parsedData.data.didWarmStart; } catch (error) { + outboundRequestsTotal.inc({ + name: "warm_start", + method: "POST", + status: "none", + outcome: "network_error", + }); this.logger.error("Warm start error", { runId: dequeuedMessage.run.id, error, From 5218d3d88f478fcb668fdf77f49cb3af833a42ac Mon Sep 17 00:00:00 2001 From: nicktrn <55853254+nicktrn@users.noreply.github.com> Date: Thu, 23 Jul 2026 17:12:27 +0100 Subject: [PATCH 2/4] feat(supervisor): instrument outbound worker api client requests --- .../supervisor-http-client-request-metrics.md | 5 ++ apps/supervisor/src/index.ts | 1 + .../src/v3/runEngineWorker/supervisor/http.ts | 71 +++++++++++++++---- .../v3/runEngineWorker/supervisor/types.ts | 8 +++ 4 files changed, 71 insertions(+), 14 deletions(-) create mode 100644 .changeset/supervisor-http-client-request-metrics.md diff --git a/.changeset/supervisor-http-client-request-metrics.md b/.changeset/supervisor-http-client-request-metrics.md new file mode 100644 index 0000000000..e0dba6e1cc --- /dev/null +++ b/.changeset/supervisor-http-client-request-metrics.md @@ -0,0 +1,5 @@ +--- +"@trigger.dev/core": patch +--- + +Add an optional `onHttpRequestComplete` callback to the supervisor worker API client so hosts can record metrics for its outbound requests (method, response status, and outcome). diff --git a/apps/supervisor/src/index.ts b/apps/supervisor/src/index.ts index 221c7e2892..3379d4e3d3 100644 --- a/apps/supervisor/src/index.ts +++ b/apps/supervisor/src/index.ts @@ -329,6 +329,7 @@ class ManagedSupervisor { runNotificationsEnabled: env.TRIGGER_WORKLOAD_API_ENABLED, heartbeatIntervalSeconds: env.TRIGGER_WORKER_HEARTBEAT_INTERVAL_SECONDS, sendRunDebugLogs: env.SEND_RUN_DEBUG_LOGS, + onHttpRequestComplete: (metric) => outboundRequestsTotal.inc(metric), preDequeue: async () => { // Synchronous, hot-path-safe cached read; false when no monitors are active. const skipForBackpressure = this.backpressureMonitors.some((m) => m.shouldSkipDequeue()); diff --git a/packages/core/src/v3/runEngineWorker/supervisor/http.ts b/packages/core/src/v3/runEngineWorker/supervisor/http.ts index d6d51b1c2b..e0eeb372ae 100644 --- a/packages/core/src/v3/runEngineWorker/supervisor/http.ts +++ b/packages/core/src/v3/runEngineWorker/supervisor/http.ts @@ -21,9 +21,9 @@ import { WorkerApiSuspendRunResponseBody, WorkerApiRunSnapshotsSinceResponseBody, } from "./schemas.js"; -import type { SupervisorClientCommonOptions } from "./types.js"; +import type { SupervisorClientCommonOptions, SupervisorHttpRequestMetric } from "./types.js"; import { getDefaultWorkerHeaders } from "./util.js"; -import { wrapZodFetch } from "../../zodfetch.js"; +import { wrapZodFetch, type ApiResult, type ZodFetchOptions } from "../../zodfetch.js"; import { createHeaders } from "../util.js"; import { WORKER_HEADERS } from "../consts.js"; import { SimpleStructuredLogger } from "../../utils/structuredLogger.js"; @@ -36,6 +36,7 @@ export class SupervisorHttpClient { private readonly instanceName: string; private readonly defaultHeaders: Record; private readonly sendRunDebugLogs: boolean; + private readonly onHttpRequestComplete?: (metric: SupervisorHttpRequestMetric) => void; private readonly logger = new SimpleStructuredLogger("supervisor-http-client"); @@ -45,6 +46,7 @@ export class SupervisorHttpClient { this.instanceName = opts.instanceName; this.defaultHeaders = getDefaultWorkerHeaders(opts); this.sendRunDebugLogs = opts.sendRunDebugLogs ?? false; + this.onHttpRequestComplete = opts.onHttpRequestComplete; if (!this.apiUrl) { throw new Error("apiURL is required and needs to be a non-empty string"); @@ -59,8 +61,38 @@ export class SupervisorHttpClient { } } + private async request( + name: string, + schema: T, + url: string, + requestInit?: RequestInit, + options?: ZodFetchOptions> + ): Promise>> { + const result = await wrapZodFetch(schema, url, requestInit, options); + + if (this.onHttpRequestComplete) { + const method = requestInit?.method ?? "GET"; + + if (result.success) { + this.onHttpRequestComplete({ name, method, status: "2xx", outcome: "ok" }); + } else if (typeof result.statusCode === "number") { + this.onHttpRequestComplete({ + name, + method, + status: String(result.statusCode), + outcome: "http_error", + }); + } else { + this.onHttpRequestComplete({ name, method, status: "none", outcome: "network_error" }); + } + } + + return result; + } + async connect(body: WorkerApiConnectRequestBody) { - return wrapZodFetch( + return this.request( + "connect", WorkerApiConnectResponseBody, `${this.apiUrl}/engine/v1/worker-actions/connect`, { @@ -75,7 +107,8 @@ export class SupervisorHttpClient { } async dequeue(body: WorkerApiDequeueRequestBody) { - return wrapZodFetch( + return this.request( + "dequeue", WorkerApiDequeueResponseBody, `${this.apiUrl}/engine/v1/worker-actions/dequeue`, { @@ -91,7 +124,8 @@ export class SupervisorHttpClient { /** @deprecated Not currently used */ async dequeueFromVersion(deploymentId: string, maxRunCount = 1, runnerId?: string) { - return wrapZodFetch( + return this.request( + "dequeue_from_version", WorkerApiDequeueResponseBody, `${this.apiUrl}/engine/v1/worker-actions/deployments/${deploymentId}/dequeue?maxRunCount=${maxRunCount}`, { @@ -104,7 +138,8 @@ export class SupervisorHttpClient { } async heartbeatWorker(body: WorkerApiHeartbeatRequestBody) { - return wrapZodFetch( + return this.request( + "heartbeat_worker", WorkerApiHeartbeatResponseBody, `${this.apiUrl}/engine/v1/worker-actions/heartbeat`, { @@ -124,7 +159,8 @@ export class SupervisorHttpClient { body: WorkerApiRunHeartbeatRequestBody, runnerId?: string ) { - return wrapZodFetch( + return this.request( + "heartbeat_run", WorkerApiRunHeartbeatResponseBody, `${this.apiUrl}/engine/v1/worker-actions/runs/${runId}/snapshots/${snapshotId}/heartbeat`, { @@ -146,7 +182,8 @@ export class SupervisorHttpClient { runnerId?: string, environmentId?: string ) { - return wrapZodFetch( + return this.request( + "start_run_attempt", WorkerApiRunAttemptStartResponseBody, `${this.apiUrl}/engine/v1/worker-actions/runs/${runId}/snapshots/${snapshotId}/attempts/start`, { @@ -168,7 +205,8 @@ export class SupervisorHttpClient { runnerId?: string, environmentId?: string ) { - return wrapZodFetch( + return this.request( + "complete_run_attempt", WorkerApiRunAttemptCompleteResponseBody, `${this.apiUrl}/engine/v1/worker-actions/runs/${runId}/snapshots/${snapshotId}/attempts/complete`, { @@ -184,7 +222,8 @@ export class SupervisorHttpClient { } async getLatestSnapshot(runId: string, runnerId?: string, environmentId?: string) { - return wrapZodFetch( + return this.request( + "get_latest_snapshot", WorkerApiRunLatestSnapshotResponseBody, `${this.apiUrl}/engine/v1/worker-actions/runs/${runId}/snapshots/latest`, { @@ -204,7 +243,8 @@ export class SupervisorHttpClient { runnerId?: string, environmentId?: string ) { - return wrapZodFetch( + return this.request( + "get_snapshots_since", WorkerApiRunSnapshotsSinceResponseBody, `${this.apiUrl}/engine/v1/worker-actions/runs/${runId}/snapshots/since/${snapshotId}`, { @@ -224,7 +264,8 @@ export class SupervisorHttpClient { } try { - const res = await wrapZodFetch( + const res = await this.request( + "send_debug_log", z.unknown(), `${this.apiUrl}/engine/v1/worker-actions/runs/${runId}/logs/debug`, { @@ -252,7 +293,8 @@ export class SupervisorHttpClient { runnerId?: string, environmentId?: string ) { - return wrapZodFetch( + return this.request( + "continue_run_execution", WorkerApiContinueRunExecutionRequestBody, `${this.apiUrl}/engine/v1/worker-actions/runs/${runId}/snapshots/${snapshotId}/continue`, { @@ -292,7 +334,8 @@ export class SupervisorHttpClient { runnerId?: string; body: WorkerApiSuspendRunRequestBody; }) { - return wrapZodFetch( + return this.request( + "submit_suspend_completion", WorkerApiSuspendRunResponseBody, `${this.apiUrl}/engine/v1/worker-actions/runs/${runId}/snapshots/${snapshotId}/suspend`, { diff --git a/packages/core/src/v3/runEngineWorker/supervisor/types.ts b/packages/core/src/v3/runEngineWorker/supervisor/types.ts index 4d0479cbca..2e9c01c3e5 100644 --- a/packages/core/src/v3/runEngineWorker/supervisor/types.ts +++ b/packages/core/src/v3/runEngineWorker/supervisor/types.ts @@ -1,5 +1,12 @@ import type { MachineResources } from "../../schemas/runEngine.js"; +export type SupervisorHttpRequestMetric = { + name: string; + method: string; + status: string; + outcome: "ok" | "http_error" | "network_error"; +}; + export type SupervisorClientCommonOptions = { apiUrl: string; workerToken: string; @@ -7,6 +14,7 @@ export type SupervisorClientCommonOptions = { deploymentId?: string; managedWorkerSecret?: string; sendRunDebugLogs?: boolean; + onHttpRequestComplete?: (metric: SupervisorHttpRequestMetric) => void; }; export type PreDequeueFn = () => Promise<{ From afd1b1e14eaeb108a56851ad5db1e3ab93e4cfae Mon Sep 17 00:00:00 2001 From: nicktrn <55853254+nicktrn@users.noreply.github.com> Date: Thu, 23 Jul 2026 17:28:19 +0100 Subject: [PATCH 3/4] fix(supervisor): distinguish invalid response outcome and drop unneeded core changeset --- .changeset/supervisor-http-client-request-metrics.md | 5 ----- .server-changes/supervisor-outbound-request-metrics.md | 2 +- packages/core/src/v3/runEngineWorker/supervisor/http.ts | 2 ++ packages/core/src/v3/runEngineWorker/supervisor/types.ts | 2 +- 4 files changed, 4 insertions(+), 7 deletions(-) delete mode 100644 .changeset/supervisor-http-client-request-metrics.md diff --git a/.changeset/supervisor-http-client-request-metrics.md b/.changeset/supervisor-http-client-request-metrics.md deleted file mode 100644 index e0dba6e1cc..0000000000 --- a/.changeset/supervisor-http-client-request-metrics.md +++ /dev/null @@ -1,5 +0,0 @@ ---- -"@trigger.dev/core": patch ---- - -Add an optional `onHttpRequestComplete` callback to the supervisor worker API client so hosts can record metrics for its outbound requests (method, response status, and outcome). diff --git a/.server-changes/supervisor-outbound-request-metrics.md b/.server-changes/supervisor-outbound-request-metrics.md index 8a510973c4..c66a499b81 100644 --- a/.server-changes/supervisor-outbound-request-metrics.md +++ b/.server-changes/supervisor-outbound-request-metrics.md @@ -3,4 +3,4 @@ area: supervisor type: improvement --- -The supervisor now reports a Prometheus metric for its outbound HTTP requests, so failed calls to upstream services are visible for monitoring. +Improved supervisor observability: it now reports metrics for its outbound requests, making failed calls to upstream services easier to monitor. diff --git a/packages/core/src/v3/runEngineWorker/supervisor/http.ts b/packages/core/src/v3/runEngineWorker/supervisor/http.ts index e0eeb372ae..2dd39186d9 100644 --- a/packages/core/src/v3/runEngineWorker/supervisor/http.ts +++ b/packages/core/src/v3/runEngineWorker/supervisor/http.ts @@ -75,6 +75,8 @@ export class SupervisorHttpClient { if (result.success) { this.onHttpRequestComplete({ name, method, status: "2xx", outcome: "ok" }); + } else if (result.statusCode === 200) { + this.onHttpRequestComplete({ name, method, status: "200", outcome: "invalid_response" }); } else if (typeof result.statusCode === "number") { this.onHttpRequestComplete({ name, diff --git a/packages/core/src/v3/runEngineWorker/supervisor/types.ts b/packages/core/src/v3/runEngineWorker/supervisor/types.ts index 2e9c01c3e5..cad4567830 100644 --- a/packages/core/src/v3/runEngineWorker/supervisor/types.ts +++ b/packages/core/src/v3/runEngineWorker/supervisor/types.ts @@ -4,7 +4,7 @@ export type SupervisorHttpRequestMetric = { name: string; method: string; status: string; - outcome: "ok" | "http_error" | "network_error"; + outcome: "ok" | "http_error" | "invalid_response" | "network_error"; }; export type SupervisorClientCommonOptions = { From 1beda66538bd2cea29af528973c98196265512dc Mon Sep 17 00:00:00 2001 From: nicktrn <55853254+nicktrn@users.noreply.github.com> Date: Thu, 23 Jul 2026 17:43:54 +0100 Subject: [PATCH 4/4] feat(supervisor): add duration histogram for outbound requests --- apps/supervisor/src/index.ts | 53 ++++++++++--------- .../src/v3/runEngineWorker/supervisor/http.ts | 21 ++++++-- .../v3/runEngineWorker/supervisor/types.ts | 1 + 3 files changed, 47 insertions(+), 28 deletions(-) diff --git a/apps/supervisor/src/index.ts b/apps/supervisor/src/index.ts index 3379d4e3d3..f30203df54 100644 --- a/apps/supervisor/src/index.ts +++ b/apps/supervisor/src/index.ts @@ -67,6 +67,14 @@ const outboundRequestsTotal = new Counter({ registers: [register], }); +const outboundRequestDuration = new Histogram({ + name: "supervisor_outbound_request_duration_seconds", + help: "Duration of outbound HTTP requests from the supervisor, by target name and outcome. Includes the HTTP client's internal retries and backoff.", + labelNames: ["name", "outcome"], + buckets: [0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1, 2, 5, 10, 11, 12.5, 15, 20, 30, 60], + registers: [register], +}); + class ManagedSupervisor { private readonly workerSession: SupervisorSession; private readonly metricsServer?: HttpServer; @@ -329,7 +337,10 @@ class ManagedSupervisor { runNotificationsEnabled: env.TRIGGER_WORKLOAD_API_ENABLED, heartbeatIntervalSeconds: env.TRIGGER_WORKER_HEARTBEAT_INTERVAL_SECONDS, sendRunDebugLogs: env.SEND_RUN_DEBUG_LOGS, - onHttpRequestComplete: (metric) => outboundRequestsTotal.inc(metric), + onHttpRequestComplete: ({ name, method, status, outcome, durationMs }) => { + outboundRequestsTotal.inc({ name, method, status, outcome }); + outboundRequestDuration.observe({ name, outcome }, durationMs / 1000); + }, preDequeue: async () => { // Synchronous, hot-path-safe cached read; false when no monitors are active. const skipForBackpressure = this.backpressureMonitors.some((m) => m.shouldSkipDequeue()); @@ -700,6 +711,18 @@ class ManagedSupervisor { headers.traceparent = traceparent; } + const requestStart = performance.now(); + const record = ( + status: string, + outcome: "ok" | "http_error" | "invalid_response" | "network_error" + ) => { + outboundRequestsTotal.inc({ name: "warm_start", method: "POST", status, outcome }); + outboundRequestDuration.observe( + { name: "warm_start", outcome }, + (performance.now() - requestStart) / 1000 + ); + }; + try { const res = await fetch(warmStartUrlWithPath.href, { method: "POST", @@ -708,12 +731,7 @@ class ManagedSupervisor { }); if (!res.ok) { - outboundRequestsTotal.inc({ - name: "warm_start", - method: "POST", - status: String(res.status), - outcome: "http_error", - }); + record(String(res.status), "http_error"); this.logger.error("Warm start failed", { runId: dequeuedMessage.run.id, statusCode: res.status, @@ -725,12 +743,7 @@ class ManagedSupervisor { const parsedData = z.object({ didWarmStart: z.boolean() }).safeParse(data); if (!parsedData.success) { - outboundRequestsTotal.inc({ - name: "warm_start", - method: "POST", - status: String(res.status), - outcome: "invalid_response", - }); + record(String(res.status), "invalid_response"); this.logger.error("Warm start response invalid", { runId: dequeuedMessage.run.id, data, @@ -738,21 +751,11 @@ class ManagedSupervisor { return false; } - outboundRequestsTotal.inc({ - name: "warm_start", - method: "POST", - status: String(res.status), - outcome: "ok", - }); + record(String(res.status), "ok"); return parsedData.data.didWarmStart; } catch (error) { - outboundRequestsTotal.inc({ - name: "warm_start", - method: "POST", - status: "none", - outcome: "network_error", - }); + record("none", "network_error"); this.logger.error("Warm start error", { runId: dequeuedMessage.run.id, error, diff --git a/packages/core/src/v3/runEngineWorker/supervisor/http.ts b/packages/core/src/v3/runEngineWorker/supervisor/http.ts index 2dd39186d9..450be43257 100644 --- a/packages/core/src/v3/runEngineWorker/supervisor/http.ts +++ b/packages/core/src/v3/runEngineWorker/supervisor/http.ts @@ -68,24 +68,39 @@ export class SupervisorHttpClient { requestInit?: RequestInit, options?: ZodFetchOptions> ): Promise>> { + const start = performance.now(); const result = await wrapZodFetch(schema, url, requestInit, options); if (this.onHttpRequestComplete) { + const durationMs = performance.now() - start; const method = requestInit?.method ?? "GET"; if (result.success) { - this.onHttpRequestComplete({ name, method, status: "2xx", outcome: "ok" }); + this.onHttpRequestComplete({ name, method, status: "2xx", outcome: "ok", durationMs }); } else if (result.statusCode === 200) { - this.onHttpRequestComplete({ name, method, status: "200", outcome: "invalid_response" }); + this.onHttpRequestComplete({ + name, + method, + status: "200", + outcome: "invalid_response", + durationMs, + }); } else if (typeof result.statusCode === "number") { this.onHttpRequestComplete({ name, method, status: String(result.statusCode), outcome: "http_error", + durationMs, }); } else { - this.onHttpRequestComplete({ name, method, status: "none", outcome: "network_error" }); + this.onHttpRequestComplete({ + name, + method, + status: "none", + outcome: "network_error", + durationMs, + }); } } diff --git a/packages/core/src/v3/runEngineWorker/supervisor/types.ts b/packages/core/src/v3/runEngineWorker/supervisor/types.ts index cad4567830..aff6ffa7d2 100644 --- a/packages/core/src/v3/runEngineWorker/supervisor/types.ts +++ b/packages/core/src/v3/runEngineWorker/supervisor/types.ts @@ -5,6 +5,7 @@ export type SupervisorHttpRequestMetric = { method: string; status: string; outcome: "ok" | "http_error" | "invalid_response" | "network_error"; + durationMs: number; }; export type SupervisorClientCommonOptions = {