@@ -4,7 +4,10 @@ import { LifecycleScope } from '#/app/scopes';
44import { IEventBus } from '#/app/event/eventBus' ;
55import { AgentActivityUpdated } from '#/agent/activityView/activityView' ;
66import { TurnStarted } from '#/agent/loop/turnEvents' ;
7- import { TurnEnded } from '#/agent/loop/turnOps' ;
7+ import { TurnEnded , turnKey } from '#/agent/loop/turnOps' ;
8+ import { ContextUndone } from '#/agent/undo/undoService' ;
9+ import { IAgentStateService } from '#/agent/state/agentState' ;
10+ import { IEventDispatcher } from '#/state/eventDispatcher' ;
811import {
912 IAgentLifecycleService ,
1013 MAIN_AGENT_ID ,
@@ -18,16 +21,18 @@ export class SessionOutcomeMirror extends Disposable implements ISessionOutcomeM
1821 declare readonly _serviceBrand : undefined ;
1922
2023 private lastPersisted : SessionTurnOutcome | undefined ;
24+ private lastPersistedTurnId : number | undefined ;
2125 private adopted = false ;
2226 private turnStartedHere = false ;
2327 private mainSubscription : DisposableStore | undefined ;
28+ private readonly metadataReady : Promise < void > ;
2429
2530 constructor (
2631 @IAgentLifecycleService private readonly agents : IAgentLifecycleService ,
2732 @ISessionMetadata private readonly metadata : ISessionMetadata ,
2833 ) {
2934 super ( ) ;
30- void this . metadata
35+ this . metadataReady = this . metadata
3136 . read ( )
3237 . then ( ( meta ) => {
3338 if ( ! this . adopted ) this . lastPersisted = meta . lastTurnReason ;
@@ -52,24 +57,33 @@ export class SessionOutcomeMirror extends Disposable implements ISessionOutcomeM
5257
5358 private attachMain ( ) : void {
5459 if ( this . mainSubscription !== undefined ) return ;
55- const bus = this . agents . handleOf ( MAIN_AGENT_ID ) ?. accessor . get ( IEventBus ) as
56- | IEventBus
57- | undefined ;
60+ const handle = this . agents . handleOf ( MAIN_AGENT_ID ) ;
61+ const bus = handle ?. accessor . get ( IEventBus ) as IEventBus | undefined ;
5862 if ( bus === undefined ) return ;
5963 const subscription = new DisposableStore ( ) ;
6064 this . mainSubscription = subscription ;
65+ const dispatcher = handle ?. accessor . get ( IEventDispatcher ) as IEventDispatcher | undefined ;
66+ const agentStates = handle ?. accessor . get ( IAgentStateService ) as IAgentStateService | undefined ;
67+ if ( dispatcher !== undefined && agentStates !== undefined ) {
68+ subscription . add (
69+ dispatcher . hooks . onDidRestore . register ( 'session-outcome-mirror' , async ( _ctx , next ) => {
70+ await next ( ) ;
71+ await this . reconcileAfterRestore ( agentStates ) ;
72+ } ) ,
73+ ) ;
74+ }
6175 subscription . add (
6276 bus . subscribe ( TurnEnded , ( event ) => {
6377 if ( event . reason === 'completed' ) {
64- this . write ( 'completed' ) ;
78+ this . write ( 'completed' , { turnId : event . turnId } ) ;
6579 return ;
6680 }
6781 if ( event . reason === 'failed' || event . reason === 'blocked' ) {
68- this . write ( 'failed' ) ;
82+ this . write ( 'failed' , { turnId : event . turnId } ) ;
6983 return ;
7084 }
7185 if ( event . reason === 'cancelled' && event . interruptReason === 'user_cancelled' ) {
72- this . write ( 'cancelled' ) ;
86+ this . write ( 'cancelled' , { turnId : event . turnId } ) ;
7387 }
7488 } ) ,
7589 ) ;
@@ -79,31 +93,67 @@ export class SessionOutcomeMirror extends Disposable implements ISessionOutcomeM
7993 this . write ( undefined ) ;
8094 } ) ,
8195 ) ;
96+ subscription . add (
97+ bus . subscribe ( ContextUndone , ( event ) => {
98+ if (
99+ event . fromTurnId !== undefined &&
100+ this . lastPersistedTurnId !== undefined &&
101+ this . lastPersistedTurnId < event . fromTurnId
102+ ) {
103+ return ;
104+ }
105+ this . write ( undefined ) ;
106+ } ) ,
107+ ) ;
82108 subscription . add (
83109 bus . subscribe ( AgentActivityUpdated , ( event ) => {
84110 if ( this . turnStartedHere ) return ;
85111 if ( this . lastPersisted !== undefined ) return ;
86112 const reason = event . lastTurn ?. reason ;
87113 if ( reason === 'completed' || reason === 'cancelled' ) {
88- this . write ( reason , { touchUpdatedAt : false } ) ;
114+ this . write ( reason , { touchUpdatedAt : false , turnId : event . lastTurn ?. turnId } ) ;
89115 } else if ( reason === 'failed' || reason === 'blocked' ) {
90- this . write ( 'failed' , { touchUpdatedAt : false } ) ;
116+ this . write ( 'failed' , { touchUpdatedAt : false , turnId : event . lastTurn ?. turnId } ) ;
91117 }
92118 } ) ,
93119 ) ;
94120 }
95121
122+ private async reconcileAfterRestore ( agentStates : IAgentStateService ) : Promise < void > {
123+ await this . metadataReady ;
124+ if ( this . lastPersisted === undefined ) return ;
125+ if ( this . turnStartedHere ) return ;
126+ if ( ! agentStates . has ( turnKey ) ) return ;
127+ const lastEnded = agentStates . get ( turnKey ) . lastEnded ;
128+ if ( lastEnded === undefined ) {
129+ this . write ( undefined , { touchUpdatedAt : false } ) ;
130+ return ;
131+ }
132+ if ( this . lastPersistedTurnId === undefined ) this . lastPersistedTurnId = lastEnded . turnId ;
133+ }
134+
96135 private write (
97136 outcome : SessionTurnOutcome | undefined ,
98- opts ?: { readonly touchUpdatedAt ?: boolean } ,
137+ opts ?: { readonly touchUpdatedAt ?: boolean ; readonly turnId ?: number } ,
99138 ) : void {
100- if ( outcome === this . lastPersisted ) return ;
139+ if ( outcome === this . lastPersisted ) {
140+ if ( opts ?. turnId !== undefined ) this . lastPersistedTurnId = opts . turnId ;
141+ return ;
142+ }
101143 this . adopted = true ;
102144 const previous = this . lastPersisted ;
145+ const previousTurnId = this . lastPersistedTurnId ;
103146 this . lastPersisted = outcome ;
104- void this . metadata . update ( { lastTurnReason : outcome } , opts ) . catch ( ( ) => {
105- if ( this . lastPersisted === outcome ) this . lastPersisted = previous ;
106- } ) ;
147+ this . lastPersistedTurnId =
148+ outcome === undefined ? undefined : ( opts ?. turnId ?? this . lastPersistedTurnId ) ;
149+ void this . metadata
150+ . update ( { lastTurnReason : outcome } , { touchUpdatedAt : opts ?. touchUpdatedAt } )
151+ . catch ( ( ) => {
152+ if ( this . lastPersisted === outcome ) {
153+ this . lastPersisted = previous ;
154+ this . lastPersistedTurnId = previousTurnId ;
155+ }
156+ } ) ;
107157 }
108158}
109159
0 commit comments