@@ -527,115 +527,6 @@ export class CopilotSession {
527527 cancel : async ( runId ) => this . rpc . workflow . cancel ( { runId } ) ,
528528 } ;
529529
530- /**
531- * Resolve a start/resume envelope into the terminal envelope callers expect.
532- *
533- * The CLI may answer `session.factory.run` and `session.factory.resume`
534- * before the run settles, so a non-terminal envelope is followed by a wait
535- * on the run's terminal state.
536- */
537- private settleFactoryRun ( envelope : WireFactoryRunResult ) : Promise < FactoryRunResult > {
538- if ( isFactoryRunTerminal ( envelope . status ) ) {
539- return Promise . resolve ( envelope ) ;
540- }
541- return this . waitForFactoryRun ( envelope . runId ) ;
542- }
543-
544- /**
545- * Resolve when a factory run reaches a terminal status.
546- *
547- * The subscription is installed *before* the first read so a transition
548- * landing between the two cannot be missed, and re-reads are serialized so
549- * overlapping invalidation events cannot interleave — the run's revision
550- * advances once per operation, so a burst of events is common and must
551- * collapse into a single in-flight read. A bounded periodic re-read keeps a
552- * dropped invalidation from leaving the wait pending forever.
553- */
554- private waitForFactoryRun ( runId : string , signal ?: AbortSignal ) : Promise < FactoryRunResult > {
555- const abortError = ( ) : unknown =>
556- signal ?. reason ?? new DOMException ( "Factory run wait was aborted" , "AbortError" ) ;
557- if ( signal ?. aborted === true ) {
558- return Promise . reject ( abortError ( ) ) ;
559- }
560-
561- return new Promise < FactoryRunResult > ( ( resolve , reject ) => {
562- let settled = false ;
563- let reading = false ;
564- let rereadRequested = false ;
565- let pollHandle : ReturnType < typeof setInterval > | undefined ;
566- let unsubscribe : ( ( ) => void ) | undefined ;
567- let onAbort : ( ( ) => void ) | undefined ;
568-
569- const finish = ( complete : ( ) => void ) : void => {
570- if ( settled ) {
571- return ;
572- }
573- settled = true ;
574- if ( pollHandle !== undefined ) {
575- clearInterval ( pollHandle ) ;
576- }
577- unsubscribe ?.( ) ;
578- if ( onAbort !== undefined ) {
579- signal ?. removeEventListener ( "abort" , onAbort ) ;
580- }
581- complete ( ) ;
582- } ;
583-
584- const read = async ( ) : Promise < void > => {
585- if ( settled ) {
586- return ;
587- }
588- if ( reading ) {
589- rereadRequested = true ;
590- return ;
591- }
592- reading = true ;
593- try {
594- do {
595- rereadRequested = false ;
596- const envelope = await this . rpc . factory . getRun ( { runId } ) ;
597- if ( isFactoryRunTerminal ( envelope . status ) ) {
598- finish ( ( ) => resolve ( envelope ) ) ;
599- return ;
600- }
601- } while ( rereadRequested && ! settled ) ;
602- } catch ( error ) {
603- finish ( ( ) => reject ( error ) ) ;
604- } finally {
605- reading = false ;
606- }
607- } ;
608-
609- if ( signal !== undefined ) {
610- onAbort = ( ) : void => finish ( ( ) => reject ( abortError ( ) ) ) ;
611- signal . addEventListener ( "abort" , onAbort , { once : true } ) ;
612- }
613-
614- const subscribe = this . on . bind ( this ) as unknown as (
615- eventType : string ,
616- handler : ( event : SessionEvent ) => void
617- ) => ( ) => void ;
618- const handleRunUpdated = ( event : SessionEvent ) : void => {
619- const data = event . data as { runId ?: unknown } ;
620- if ( data . runId === runId ) {
621- void read ( ) ;
622- }
623- } ;
624- const unsubscribeFactory = subscribe ( "factory.run_updated" , handleRunUpdated ) ;
625- const unsubscribeWorkflow = subscribe ( "workflow.run_updated" , handleRunUpdated ) ;
626- unsubscribe = ( ) : void => {
627- unsubscribeFactory ( ) ;
628- unsubscribeWorkflow ( ) ;
629- } ;
630-
631- pollHandle = setInterval ( ( ) => void read ( ) , 5_000 ) ;
632- // The re-read is a safety net, not work the process owes anyone: an
633- // outstanding wait must never keep Node alive on its own.
634- pollHandle . unref ?.( ) ;
635- void read ( ) ;
636- } ) ;
637- }
638-
639530 private settleWorkflowRun ( envelope : WireWorkflowRunResult ) : Promise < WorkflowRunResult > {
640531 if ( isWorkflowRunTerminal ( envelope . status ) ) {
641532 return Promise . resolve ( envelope ) ;
0 commit comments