feat(supervisor): wide events on dequeue + warm-start trace

This commit is contained in:
nicktrn
2026-05-19 13:50:39 +01:00
parent ec3abbc398
commit 3591aa9d32
+214 -134
View File
@@ -28,6 +28,14 @@ import { FailedPodHandler } from "./services/failedPodHandler.js";
import { getWorkerToken } from "./workerToken.js";
import { OtlpTraceService } from "./services/otlpTraceService.js";
import { extractTraceparent, getRestoreRunnerId } from "./util.js";
import {
fromContext,
recordPhaseSince,
runWideEvent,
setExtra,
setMeta,
type WideEventOptions,
} from "./wideEvents/index.js";
if (env.METRICS_COLLECT_DEFAULTS) {
collectDefaultMetrics({ register });
@@ -50,6 +58,12 @@ class ManagedSupervisor {
private readonly isKubernetes = isKubernetesEnvironment(env.KUBERNETES_FORCE_ENABLED);
private readonly warmStartUrl = env.TRIGGER_WARM_START_URL;
private readonly wideEventOpts: WideEventOptions = {
service: "supervisor",
env: { nodeId: env.TRIGGER_WORKER_INSTANCE_NAME },
enabled: env.TRIGGER_WIDE_EVENTS_ENABLED,
};
constructor() {
const {
TRIGGER_WORKER_TOKEN,
@@ -239,149 +253,202 @@ class ManagedSupervisor {
async ({ time, message, dequeueResponseMs, pollingIntervalMs }) => {
this.logger.verbose(`Received message with timestamp ${time.toLocaleString()}`, message);
if (message.completedWaitpoints.length > 0) {
this.logger.debug("Run has completed waitpoints", {
runId: message.run.id,
completedWaitpoints: message.completedWaitpoints.length,
});
}
const traceparent = extractTraceparent(message.run.traceContext);
if (!message.image) {
this.logger.error("Run has no image", { runId: message.run.id });
return;
}
await runWideEvent(
{
...this.wideEventOpts,
traceparent,
setup: (state) => {
setMeta(state, "run_id", message.run.id);
setMeta(state, "env_id", message.environment.id);
setMeta(state, "org_id", message.organization.id);
setMeta(state, "project_id", message.project.id);
if (message.deployment.friendlyId) {
setMeta(state, "deployment_id", message.deployment.friendlyId);
}
setMeta(state, "machine_preset", message.run.machine.name);
state.extras.iteration = "dequeue";
state.extras.dequeue_response_ms = dequeueResponseMs;
state.extras.polling_interval_ms = pollingIntervalMs;
state.extras.completed_waitpoints = message.completedWaitpoints.length;
},
},
async () => {
if (message.completedWaitpoints.length > 0) {
this.logger.debug("Run has completed waitpoints", {
runId: message.run.id,
completedWaitpoints: message.completedWaitpoints.length,
});
}
const { checkpoint, ...rest } = message;
if (!message.image) {
setExtra(fromContext(), "path_taken", "skipped_no_image");
this.logger.error("Run has no image", { runId: message.run.id });
return;
}
// Register trace context early so snapshot spans work for all paths
// (cold create, restore, warm start). Re-registration on restore is safe
// since dequeue always provides fresh context.
if (this.computeManager?.traceSpansEnabled) {
const traceparent = extractTraceparent(message.run.traceContext);
const { checkpoint, ...rest } = message;
if (traceparent) {
this.workloadServer.registerRunTraceContext(message.run.friendlyId, {
traceparent,
envId: message.environment.id,
orgId: message.organization.id,
projectId: message.project.id,
});
}
}
if (checkpoint) {
this.logger.debug("Restoring run", { runId: message.run.id });
if (this.computeManager) {
try {
const runnerId = getRestoreRunnerId(message.run.friendlyId, checkpoint.id);
const didRestore = await this.computeManager.restore({
snapshotId: checkpoint.location,
runnerId,
runFriendlyId: message.run.friendlyId,
snapshotFriendlyId: message.snapshot.friendlyId,
machine: message.run.machine,
traceContext: message.run.traceContext,
// Register trace context early so snapshot spans work for all paths
// (cold create, restore, warm start). Re-registration on restore is safe
// since dequeue always provides fresh context.
if (this.computeManager?.traceSpansEnabled && traceparent) {
this.workloadServer.registerRunTraceContext(message.run.friendlyId, {
traceparent,
envId: message.environment.id,
orgId: message.organization.id,
projectId: message.project.id,
dequeuedAt: message.dequeuedAt,
});
}
if (didRestore) {
this.logger.debug("Compute restore successful", {
runId: message.run.id,
runnerId,
});
} else {
this.logger.error("Compute restore failed", { runId: message.run.id, runnerId });
if (checkpoint) {
setExtra(fromContext(), "path_taken", "restore");
this.logger.debug("Restoring run", { runId: message.run.id });
if (this.computeManager) {
const restoreStart = performance.now();
try {
const runnerId = getRestoreRunnerId(message.run.friendlyId, checkpoint.id);
const didRestore = await this.computeManager.restore({
snapshotId: checkpoint.location,
runnerId,
runFriendlyId: message.run.friendlyId,
snapshotFriendlyId: message.snapshot.friendlyId,
machine: message.run.machine,
traceContext: message.run.traceContext,
envId: message.environment.id,
orgId: message.organization.id,
projectId: message.project.id,
dequeuedAt: message.dequeuedAt,
});
recordPhaseSince("restore", restoreStart, undefined);
setExtra(fromContext(), "did_restore", didRestore);
if (didRestore) {
this.logger.debug("Compute restore successful", {
runId: message.run.id,
runnerId,
});
} else {
this.logger.error("Compute restore failed", {
runId: message.run.id,
runnerId,
});
}
} catch (error) {
recordPhaseSince(
"restore",
restoreStart,
error instanceof Error ? error : new Error(String(error))
);
this.logger.error("Failed to restore run (compute)", { error });
}
return;
}
if (!this.checkpointClient) {
this.logger.error("No checkpoint client", { runId: message.run.id });
return;
}
const restoreStart = performance.now();
try {
const didRestore = await this.checkpointClient.restoreRun({
runFriendlyId: message.run.friendlyId,
snapshotFriendlyId: message.snapshot.friendlyId,
body: {
...rest,
checkpoint,
},
});
recordPhaseSince("restore", restoreStart, undefined);
setExtra(fromContext(), "did_restore", didRestore);
if (didRestore) {
this.logger.debug("Restore successful", { runId: message.run.id });
} else {
this.logger.error("Restore failed", { runId: message.run.id });
}
} catch (error) {
recordPhaseSince(
"restore",
restoreStart,
error instanceof Error ? error : new Error(String(error))
);
this.logger.error("Failed to restore run", { error });
}
return;
}
this.logger.debug("Scheduling run", { runId: message.run.id });
const warmStartStart = performance.now();
const didWarmStart = await this.tryWarmStart(message, traceparent);
const warmStartCheckMs = Math.round(performance.now() - warmStartStart);
recordPhaseSince("warm_start", warmStartStart, undefined);
setExtra(fromContext(), "did_warm_start", didWarmStart);
if (didWarmStart) {
setExtra(fromContext(), "path_taken", "warm_start");
this.logger.debug("Warm start successful", { runId: message.run.id });
return;
}
setExtra(fromContext(), "path_taken", "cold_create");
const createStart = performance.now();
try {
if (!message.deployment.friendlyId) {
// mostly a type guard, deployments always exists for deployed environments
// a proper fix would be to use a discriminated union schema to differentiate between dequeued runs in dev and in deployed environments.
throw new Error("Deployment is missing");
}
await this.workloadManager.create({
dequeuedAt: message.dequeuedAt,
dequeueResponseMs,
pollingIntervalMs,
warmStartCheckMs,
envId: message.environment.id,
envType: message.environment.type,
image: message.image,
machine: message.run.machine,
orgId: message.organization.id,
projectId: message.project.id,
deploymentFriendlyId: message.deployment.friendlyId,
deploymentVersion: message.backgroundWorker.version,
runId: message.run.id,
runFriendlyId: message.run.friendlyId,
version: message.version,
nextAttemptNumber: message.run.attemptNumber,
snapshotId: message.snapshot.id,
snapshotFriendlyId: message.snapshot.friendlyId,
placementTags: message.placementTags,
traceContext: message.run.traceContext,
annotations: message.run.annotations,
hasPrivateLink: message.organization.hasPrivateLink,
});
recordPhaseSince("workload_create", createStart, undefined);
// Disabled for now
// this.resourceMonitor.blockResources({
// cpu: message.run.machine.cpu,
// memory: message.run.machine.memory,
// });
} catch (error) {
this.logger.error("Failed to restore run (compute)", { error });
recordPhaseSince(
"workload_create",
createStart,
error instanceof Error ? error : new Error(String(error))
);
this.logger.error("Failed to create workload", { error });
}
return;
}
if (!this.checkpointClient) {
this.logger.error("No checkpoint client", { runId: message.run.id });
return;
}
try {
const didRestore = await this.checkpointClient.restoreRun({
runFriendlyId: message.run.friendlyId,
snapshotFriendlyId: message.snapshot.friendlyId,
body: {
...rest,
checkpoint,
},
});
if (didRestore) {
this.logger.debug("Restore successful", { runId: message.run.id });
} else {
this.logger.error("Restore failed", { runId: message.run.id });
}
} catch (error) {
this.logger.error("Failed to restore run", { error });
}
return;
}
this.logger.debug("Scheduling run", { runId: message.run.id });
const warmStartStart = performance.now();
const didWarmStart = await this.tryWarmStart(message);
const warmStartCheckMs = Math.round(performance.now() - warmStartStart);
if (didWarmStart) {
this.logger.debug("Warm start successful", { runId: message.run.id });
return;
}
try {
if (!message.deployment.friendlyId) {
// mostly a type guard, deployments always exists for deployed environments
// a proper fix would be to use a discriminated union schema to differentiate between dequeued runs in dev and in deployed environments.
throw new Error("Deployment is missing");
}
await this.workloadManager.create({
dequeuedAt: message.dequeuedAt,
dequeueResponseMs,
pollingIntervalMs,
warmStartCheckMs,
envId: message.environment.id,
envType: message.environment.type,
image: message.image,
machine: message.run.machine,
orgId: message.organization.id,
projectId: message.project.id,
deploymentFriendlyId: message.deployment.friendlyId,
deploymentVersion: message.backgroundWorker.version,
runId: message.run.id,
runFriendlyId: message.run.friendlyId,
version: message.version,
nextAttemptNumber: message.run.attemptNumber,
snapshotId: message.snapshot.id,
snapshotFriendlyId: message.snapshot.friendlyId,
placementTags: message.placementTags,
traceContext: message.run.traceContext,
annotations: message.run.annotations,
hasPrivateLink: message.organization.hasPrivateLink,
});
// Disabled for now
// this.resourceMonitor.blockResources({
// cpu: message.run.machine.cpu,
// memory: message.run.machine.memory,
// });
} catch (error) {
this.logger.error("Failed to create workload", { error });
}
);
}
);
@@ -404,6 +471,7 @@ class ManagedSupervisor {
checkpointClient: this.checkpointClient,
computeManager: this.computeManager,
tracing: this.tracing,
wideEventOpts: this.wideEventOpts,
});
this.workloadServer.on("runConnected", this.onRunConnected.bind(this));
@@ -420,19 +488,31 @@ class ManagedSupervisor {
this.workerSession.unsubscribeFromRunNotifications([run.friendlyId]);
}
private async tryWarmStart(dequeuedMessage: DequeuedMessage): Promise<boolean> {
private async tryWarmStart(
dequeuedMessage: DequeuedMessage,
traceparent: string | undefined
): Promise<boolean> {
if (!this.warmStartUrl) {
return false;
}
const warmStartUrlWithPath = new URL("/warm-start", this.warmStartUrl);
const headers: Record<string, string> = {
"Content-Type": "application/json",
};
// Propagate the inbound W3C traceparent so the upstream warm-start
// receiver continues the same trace instead of minting a new one. Gated
// by the same kill switch as the wide-event emission so the whole PR is
// a no-op on the wire when disabled.
if (this.wideEventOpts.enabled && traceparent) {
headers.traceparent = traceparent;
}
try {
const res = await fetch(warmStartUrlWithPath.href, {
method: "POST",
headers: {
"Content-Type": "application/json",
},
headers,
body: JSON.stringify({ dequeuedMessage }),
});