@@ -103,6 +103,22 @@ export function agentRebuildFailure(err: unknown): Error {
103103 : new Error ( String ( err ) ) ;
104104}
105105
106+ /**
107+ * Hard-stop interrupt: bump delivery generation, then enqueue the agent rebuild.
108+ * Overlay stays open across interrupt; the bump must happen before enqueue so a
109+ * later accept/decline cannot late-bind into the rebuilt agent.
110+ */
111+ export function startInterruptRebuild ( args : {
112+ deliveryGeneration : { bump : ( ) => void } ;
113+ markSendAborted : ( ) => void ;
114+ enqueue : ( op : ( ) => Promise < void > ) => unknown ;
115+ rebuild : ( ) => Promise < void > ;
116+ } ) : void {
117+ args . deliveryGeneration . bump ( ) ;
118+ args . markSendAborted ( ) ;
119+ void args . enqueue ( args . rebuild ) ;
120+ }
121+
106122export function clearsActiveRun ( kind : SnapshotKind ) : boolean {
107123 return kind === "run-end" ;
108124}
@@ -384,36 +400,39 @@ export async function createRunLifecycle(
384400 // Close it, drain the old stream, and rebuild a fresh agent so the next send
385401 // works.
386402 const interrupt = ( ) : void => {
387- // Overlay stays open across interrupt; bump so a later accept/decline
388- // cannot late-bind into the rebuilt agent.
389- services . deliveryGeneration . bump ( ) ;
390- state . sendAborted = true ;
391- void enqueueOp ( async ( ) => {
392- try {
393- // close() tears down stream consumers before the aborted cycle's
394- // inference.error is delivered, so the recorder never sees a terminal
395- // event for the dead cycle — dispose closes it against stray deltas
396- // and salvages the buffer before that teardown, so it is never lost
397- // or misattributed to the rebuilt agent's next cycle.
398- await services . cycleRecorder . dispose ( "interrupted" ) ;
399- const closedCleanly = await closeAgentForRebuild ( liveAgent ( state ) , "interrupt" ) ;
400- await state . streamPromise ?. catch ( ( err : unknown ) => {
401- tuiLogger . debug ( "stream drain during interrupt teardown failed: {error}" , {
402- error : err instanceof Error ? err . message : String ( err ) ,
403+ startInterruptRebuild ( {
404+ deliveryGeneration : services . deliveryGeneration ,
405+ markSendAborted : ( ) => {
406+ state . sendAborted = true ;
407+ } ,
408+ enqueue : enqueueOp ,
409+ rebuild : async ( ) => {
410+ try {
411+ // close() tears down stream consumers before the aborted cycle's
412+ // inference.error is delivered, so the recorder never sees a terminal
413+ // event for the dead cycle — dispose closes it against stray deltas
414+ // and salvages the buffer before that teardown, so it is never lost
415+ // or misattributed to the rebuilt agent's next cycle.
416+ await services . cycleRecorder . dispose ( "interrupted" ) ;
417+ const closedCleanly = await closeAgentForRebuild ( liveAgent ( state ) , "interrupt" ) ;
418+ await state . streamPromise ?. catch ( ( err : unknown ) => {
419+ tuiLogger . debug ( "stream drain during interrupt teardown failed: {error}" , {
420+ error : err instanceof Error ? err . message : String ( err ) ,
421+ } ) ;
403422 } ) ;
404- } ) ;
405- if ( ! closedCleanly ) {
406- throw new AgentContextLockError ( state . workdir ) ;
423+ if ( ! closedCleanly ) {
424+ throw new AgentContextLockError ( state . workdir ) ;
425+ }
426+ state . currentAgent = await services . buildAgent ( ) ;
427+ services . cycleRecorder . reset ( ) ;
428+ state . streamPromise = consumeStream ( liveAgent ( state ) . stream ( ) , streamSink ) ;
429+ services . workflowController . reattach ( ) ;
430+ state . fatalBuildError = null ;
431+ } catch ( err ) {
432+ recordRunError ( state , err ) ;
433+ state . fatalBuildError = agentRebuildFailure ( err ) ;
407434 }
408- state . currentAgent = await services . buildAgent ( ) ;
409- services . cycleRecorder . reset ( ) ;
410- state . streamPromise = consumeStream ( liveAgent ( state ) . stream ( ) , streamSink ) ;
411- services . workflowController . reattach ( ) ;
412- state . fatalBuildError = null ;
413- } catch ( err ) {
414- recordRunError ( state , err ) ;
415- state . fatalBuildError = agentRebuildFailure ( err ) ;
416- }
435+ } ,
417436 } ) ;
418437 } ;
419438 state . interrupt = interrupt ;
0 commit comments