@@ -3,7 +3,11 @@ import { SimpleStructuredLogger } from "@trigger.dev/core/v3/utils/structuredLog
33import { formatLogLine , startTelnetLogServer } from "@trigger.dev/core/v3/telnetLogServer" ;
44import { env } from "./env.js" ;
55import { WorkloadServer } from "./workloadServer/index.js" ;
6- import type { WorkloadManagerOptions , WorkloadManager } from "./workloadManager/types.js" ;
6+ import type {
7+ WorkloadManagerCreateOptions ,
8+ WorkloadManagerOptions ,
9+ WorkloadManager ,
10+ } from "./workloadManager/types.js" ;
711import Docker from "dockerode" ;
812import { z } from "zod" ;
913import { type DequeuedMessage } from "@trigger.dev/core/v3" ;
@@ -84,6 +88,7 @@ class ManagedSupervisor {
8488 private readonly workloadManager : WorkloadManager ;
8589 private readonly workloadManagerBackend : "compute" | "kubernetes" | "run-crd" | "docker" ;
8690 private readonly computeManager ?: ComputeWorkloadManager ;
91+ private readonly runCrdManager ?: RunCrdWorkloadManager ;
8792 private readonly logger = new SimpleStructuredLogger ( "managed-supervisor" ) ;
8893 private readonly resourceMonitor : ResourceMonitor ;
8994 private readonly checkpointClient ?: CheckpointClient ;
@@ -192,11 +197,18 @@ class ManagedSupervisor {
192197 this . workloadManager = computeManager ;
193198 this . workloadManagerBackend = "compute" ;
194199 } else if ( this . isKubernetes && env . KUBERNETES_RUN_CRD_ENABLED ) {
195- this . workloadManager = new RunCrdWorkloadManager ( {
200+ const runCrdManager = new RunCrdWorkloadManager ( {
196201 ...workloadManagerOptions ,
197202 namespace : env . KUBERNETES_NAMESPACE ,
198203 runtime : env . KUBERNETES_RUNNER_RUNTIME ,
204+ snapshots : {
205+ enabled : env . COMPUTE_SNAPSHOTS_ENABLED ,
206+ delayMs : env . COMPUTE_SNAPSHOT_DELAY_MS ,
207+ dispatchLimit : env . COMPUTE_SNAPSHOT_DISPATCH_LIMIT ,
208+ } ,
199209 } ) ;
210+ this . runCrdManager = runCrdManager ;
211+ this . workloadManager = runCrdManager ;
200212 this . workloadManagerBackend = "run-crd" ;
201213 } else if ( this . isKubernetes ) {
202214 this . workloadManager = new KubernetesWorkloadManager ( workloadManagerOptions ) ;
@@ -521,6 +533,11 @@ class ManagedSupervisor {
521533 return ;
522534 }
523535
536+ if ( this . runCrdManager ?. restores ( checkpoint ) ) {
537+ await this . restoreRunner ( this . runCrdManager , message , checkpoint ) ;
538+ return ;
539+ }
540+
524541 if ( ! this . checkpointClient ) {
525542 this . logger . error ( "No checkpoint client" , { runId : message . run . id } ) ;
526543 return ;
@@ -608,6 +625,7 @@ class ManagedSupervisor {
608625 workerClient : this . workerSession . httpClient ,
609626 checkpointClient : this . checkpointClient ,
610627 computeManager : this . computeManager ,
628+ runnerSnapshotter : this . runCrdManager ,
611629 tracing : this . tracing ,
612630 snapshotCallbackSecret : workerToken ,
613631 wideEventOpts : this . wideEventOpts ,
@@ -630,57 +648,86 @@ class ManagedSupervisor {
630648 this . workerSession . unsubscribeFromRunNotifications ( [ run . friendlyId ] ) ;
631649 }
632650
633- private async createWorkload ( message : DequeuedMessage , timings : WarmStartTimings ) {
634- const createStart = performance . now ( ) ;
635- try {
636- if ( ! message . deployment . friendlyId ) {
637- // mostly a type guard, deployments always exists for deployed environments
638- // a proper fix would be to use a discriminated union schema to differentiate between dequeued runs in dev and in deployed environments.
639- throw new Error ( "Deployment is missing" ) ;
640- }
651+ private async createOptionsFor (
652+ message : DequeuedMessage ,
653+ timings ?: WarmStartTimings
654+ ) : Promise < WorkloadManagerCreateOptions > {
655+ if ( ! message . deployment . friendlyId ) {
656+ // mostly a type guard, deployments always exists for deployed environments
657+ // a proper fix would be to use a discriminated union schema to differentiate between dequeued runs in dev and in deployed environments.
658+ throw new Error ( "Deployment is missing" ) ;
659+ }
641660
642- if ( ! message . image ) {
643- // same type-guard situation as deployment above
644- throw new Error ( "Image is missing" ) ;
645- }
661+ if ( ! message . image ) {
662+ // same type-guard situation as deployment above
663+ throw new Error ( "Image is missing" ) ;
664+ }
646665
647- const deploymentToken = await mintDeploymentToken ( {
648- deployment : message . deployment . friendlyId ,
649- deployment_version : message . backgroundWorker . version ,
650- environment_id : message . environment . id ,
651- environment_type : message . environment . type ,
652- org_id : message . organization . id ,
653- project_id : message . project . id ,
654- } ) ;
666+ const deploymentToken = await mintDeploymentToken ( {
667+ deployment : message . deployment . friendlyId ,
668+ deployment_version : message . backgroundWorker . version ,
669+ environment_id : message . environment . id ,
670+ environment_type : message . environment . type ,
671+ org_id : message . organization . id ,
672+ project_id : message . project . id ,
673+ } ) ;
655674
656- await this . workloadManager . create ( {
657- dequeuedAt : message . dequeuedAt ,
658- dequeueResponseMs : timings . dequeueResponseMs ,
659- pollingIntervalMs : timings . pollingIntervalMs ,
660- warmStartCheckMs : timings . warmStartCheckMs ,
661- envId : message . environment . id ,
662- envType : message . environment . type ,
663- image : message . image ,
664- machine : message . run . machine ,
665- orgId : message . organization . id ,
666- projectId : message . project . id ,
667- deploymentFriendlyId : message . deployment . friendlyId ,
668- deploymentVersion : message . backgroundWorker . version ,
669- runtime : message . backgroundWorker . runtime ,
670- deploymentToken,
671- runId : message . run . id ,
672- runFriendlyId : message . run . friendlyId ,
673- version : message . version ,
674- nextAttemptNumber : message . run . attemptNumber ,
675- snapshotId : message . snapshot . id ,
676- snapshotFriendlyId : message . snapshot . friendlyId ,
677- // Carry the run's storage route to the cold-start pod so its start request echoes it back.
678- snapshotRoute : message . snapshotRoute ,
679- placementTags : message . placementTags ,
680- traceContext : message . run . traceContext ,
681- annotations : message . run . annotations ,
682- hasPrivateLink : message . organization . hasPrivateLink ,
683- } ) ;
675+ return {
676+ dequeuedAt : message . dequeuedAt ,
677+ dequeueResponseMs : timings ?. dequeueResponseMs ,
678+ pollingIntervalMs : timings ?. pollingIntervalMs ,
679+ warmStartCheckMs : timings ?. warmStartCheckMs ,
680+ envId : message . environment . id ,
681+ envType : message . environment . type ,
682+ image : message . image ,
683+ machine : message . run . machine ,
684+ orgId : message . organization . id ,
685+ projectId : message . project . id ,
686+ deploymentFriendlyId : message . deployment . friendlyId ,
687+ deploymentVersion : message . backgroundWorker . version ,
688+ runtime : message . backgroundWorker . runtime ,
689+ deploymentToken,
690+ runId : message . run . id ,
691+ runFriendlyId : message . run . friendlyId ,
692+ version : message . version ,
693+ nextAttemptNumber : message . run . attemptNumber ,
694+ snapshotId : message . snapshot . id ,
695+ snapshotFriendlyId : message . snapshot . friendlyId ,
696+ // Carry the run's storage route to the runner pod so its start request echoes it back.
697+ snapshotRoute : message . snapshotRoute ,
698+ placementTags : message . placementTags ,
699+ traceContext : message . run . traceContext ,
700+ annotations : message . run . annotations ,
701+ hasPrivateLink : message . organization . hasPrivateLink ,
702+ } ;
703+ }
704+
705+ private async restoreRunner (
706+ manager : RunCrdWorkloadManager ,
707+ message : DequeuedMessage ,
708+ checkpoint : { id : string ; location : string }
709+ ) {
710+ const restoreStart = performance . now ( ) ;
711+ try {
712+ await manager . restore ( await this . createOptionsFor ( message ) , checkpoint ) ;
713+ recordPhaseSince ( "restore" , restoreStart , undefined ) ;
714+ setExtra ( fromContext ( ) , "did_restore" , true ) ;
715+ this . logger . debug ( "Runner restore created" , { runId : message . run . id } ) ;
716+ } catch ( error ) {
717+ recordPhaseSince (
718+ "restore" ,
719+ restoreStart ,
720+ error instanceof Error ? error : new Error ( String ( error ) )
721+ ) ;
722+ setExtra ( fromContext ( ) , "did_restore" , false ) ;
723+ this . logger . error ( "Failed to restore run (run-crd)" , { runId : message . run . id , error } ) ;
724+ }
725+ }
726+
727+ private async createWorkload ( message : DequeuedMessage , timings : WarmStartTimings ) {
728+ const createStart = performance . now ( ) ;
729+ try {
730+ await this . workloadManager . create ( await this . createOptionsFor ( message , timings ) ) ;
684731 recordPhaseSince ( "workload_create" , createStart , undefined ) ;
685732 workloadCreateDuration . observe (
686733 { backend : this . workloadManagerBackend , outcome : "success" } ,
@@ -722,10 +769,8 @@ class ManagedSupervisor {
722769 const headers : Record < string , string > = {
723770 "Content-Type" : "application/json" ,
724771 } ;
725- // Propagate the inbound W3C traceparent so the upstream warm-start
726- // receiver continues the same trace instead of minting a new one. Gated
727- // by the same kill switch as the wide-event emission so the whole PR is
728- // a no-op on the wire when disabled.
772+ // Lets the warm-start receiver continue this trace instead of minting a new one.
773+ // Gated by the wide-events kill switch, so the wire is unchanged when disabled.
729774 if ( this . wideEventOpts . enabled && traceparent ) {
730775 headers . traceparent = traceparent ;
731776 }
0 commit comments