@@ -27,11 +27,16 @@ import type { DAGNode } from '@/executor/dag/builder'
2727import { BlockExecutor } from '@/executor/execution/block-executor'
2828import { ExecutionState } from '@/executor/execution/state'
2929import { AgentBlockHandler } from '@/executor/handlers/agent/agent-handler'
30+ import * as agentMemory from '@/executor/handlers/agent/memory'
3031import type { AgentInputs , Message } from '@/executor/handlers/agent/types'
3132import type { ExecutionContext , StreamingExecution } from '@/executor/types'
3233import { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry'
3334import { VariableResolver } from '@/executor/variables/resolver'
3435import { executeProviderRequest } from '@/providers'
36+ import {
37+ getEncryptedConversationMessage ,
38+ setEncryptedConversationMessage ,
39+ } from '@/providers/conversation-metadata'
3540import { installStreamingCostPolicy } from '@/providers/cost-policy'
3641import { getModelCapabilities , SIM_AUTO_MODEL_ID } from '@/providers/models'
3742import {
@@ -49,10 +54,16 @@ const {
4954 mockDiscoverMcpServerToolsAsExecutor,
5055 mockImportWorkspaceFileSecretProvenanceForModelView,
5156 mockValidateModelProvider,
57+ mockOpenAgentTurnSession,
5258} = vi . hoisted ( ( ) => ( {
5359 mockDiscoverMcpServerToolsAsExecutor : vi . fn ( ) . mockResolvedValue ( [ ] ) ,
5460 mockImportWorkspaceFileSecretProvenanceForModelView : vi . fn ( ) . mockResolvedValue ( true ) ,
5561 mockValidateModelProvider : vi . fn ( ) . mockResolvedValue ( undefined ) ,
62+ mockOpenAgentTurnSession : vi . fn ( ) . mockResolvedValue ( undefined ) ,
63+ } ) )
64+
65+ vi . mock ( '@/lib/memory/agent-turn-session' , ( ) => ( {
66+ openAgentTurnSession : mockOpenAgentTurnSession ,
5667} ) )
5768
5869vi . mock ( '@/lib/internal/mcp/discover-tools' , ( ) => ( {
@@ -77,6 +88,7 @@ vi.mock('@/providers/utils', () => ({
7788 'function' in toolCall &&
7889 ( toolCall as { function ?: unknown } ) . function != null ,
7990 getProviderFromModel : vi . fn ( ) . mockReturnValue ( 'mock-provider' ) ,
91+ isDeepResearchModel : ( model : string ) => model . includes ( 'deep-research' ) ,
8092 transformBlockTool : vi . fn ( ) ,
8193 getBaseModelProviders : vi . fn ( ) . mockReturnValue ( { openai : { } , anthropic : { } } ) ,
8294 getApiKey : vi . fn ( ) . mockReturnValue ( 'mock-api-key' ) ,
@@ -188,6 +200,7 @@ describe('AgentBlockHandler', () => {
188200 beforeEach ( ( ) => {
189201 handler = new AgentBlockHandler ( )
190202 vi . clearAllMocks ( )
203+ mockOpenAgentTurnSession . mockReset ( ) . mockResolvedValue ( undefined )
191204 mockValidateModelProvider . mockReset ( ) . mockResolvedValue ( undefined )
192205 mockDiscoverMcpServerToolsAsExecutor . mockImplementation (
193206 async ( { serverId } : { serverId : string } ) =>
@@ -339,6 +352,208 @@ describe('AgentBlockHandler', () => {
339352 } )
340353 } )
341354
355+ describe ( 'durable conversation lifecycle' , ( ) => {
356+ const inputs : AgentInputs = {
357+ model : 'gpt-4o' ,
358+ memoryType : 'conversation' ,
359+ conversationId : 'conversation-1' ,
360+ messages : [ { role : 'user' , content : 'Continue the work.' } ] ,
361+ userPrompt : 'Keep the answer brief.' ,
362+ }
363+
364+ afterEach ( ( ) => vi . restoreAllMocks ( ) )
365+
366+ it ( 'shares one turn across fallback and preserves private history metadata' , async ( ) => {
367+ const session = {
368+ turnId : 'turn-1' ,
369+ memoryId : 'memory-1' ,
370+ finalize : vi . fn ( ) ,
371+ getFinalResponse : vi . fn ( ) ,
372+ getFinalAssistantContent : vi . fn ( ) ,
373+ }
374+ mockOpenAgentTurnSession . mockResolvedValue ( session )
375+ mockGetProviderFromModel . mockImplementation ( ( model : string ) =>
376+ model . startsWith ( 'claude' ) ? 'anthropic' : 'openai'
377+ )
378+ const assistant : Message = {
379+ role : 'assistant' ,
380+ content : null ,
381+ tool_calls : [
382+ { id : 'call-1' , type : 'function' , function : { name : 'search' , arguments : '{}' } } ,
383+ ] ,
384+ }
385+ setEncryptedConversationMessage ( assistant , 'encrypted-private-native-history' )
386+ const fetch = vi
387+ . spyOn ( agentMemory . memoryService , 'fetchMemoryMessages' )
388+ . mockResolvedValue ( [
389+ { role : 'user' , content : 'Search first.' } ,
390+ assistant ,
391+ { role : 'tool' , content : '{"answer":42}' , name : 'search' , tool_call_id : 'call-1' } ,
392+ ] )
393+ const append = vi . spyOn ( agentMemory . memoryService , 'appendToMemory' ) . mockResolvedValue ( )
394+ mockExecuteProviderRequest . mockRejectedValueOnce ( new Error ( 'overloaded' ) )
395+
396+ await handler . execute (
397+ { ...mockContext , executionId : 'execution-1' } ,
398+ mockBlock ,
399+ { ...inputs , fallbackModels : [ { model : 'claude-sonnet-5' } ] } ,
400+ { nodeId : 'agent-node' , executionOrder : 3 }
401+ )
402+
403+ expect ( mockOpenAgentTurnSession ) . toHaveBeenCalledWith (
404+ expect . objectContaining ( { nodeId : 'agent-node' , executionOrder : 3 } )
405+ )
406+ expect ( fetch ) . toHaveBeenCalledWith (
407+ expect . anything ( ) ,
408+ expect . anything ( ) ,
409+ expect . any ( WeakMap ) ,
410+ { richHistory : true , excludeTurnId : 'turn-1' , memoryId : 'memory-1' }
411+ )
412+ expect ( mockExecuteProviderRequest ) . toHaveBeenCalledTimes ( 2 )
413+ for ( const [ , request , runtime ] of mockExecuteProviderRequest . mock . calls ) {
414+ expect ( runtime . agentConversation ) . toBe ( session )
415+ expect ( request . messages [ 1 ] ) . toEqual ( assistant )
416+ expect ( getEncryptedConversationMessage ( request . messages [ 1 ] ) ) . toBe (
417+ 'encrypted-private-native-history'
418+ )
419+ }
420+ expect ( append . mock . calls . map ( ( call ) => call [ 3 ] ?. appendKey ) ) . toEqual ( [ 'input' , 'user-prompt' ] )
421+ for ( const call of append . mock . calls ) {
422+ expect ( call [ 3 ] ) . toMatchObject ( { memoryId : 'memory-1' , turnId : 'turn-1' } )
423+ }
424+ expect ( session . finalize ) . toHaveBeenCalledWith ( 'Mocked response content' , 'claude-sonnet-5' )
425+ } )
426+
427+ it ( 'deduplicates retry inputs by invocation while another loop iteration gets a new turn' , async ( ) => {
428+ const stored : Message [ ] = [ ]
429+ const turns = new WeakMap < object , string > ( )
430+ const keys = new WeakMap < object , string > ( )
431+ vi . spyOn ( agentMemory , 'getMemoryMessageTurnId' ) . mockImplementation ( ( message ) =>
432+ turns . get ( message )
433+ )
434+ vi . spyOn ( agentMemory , 'getMemoryMessageAppendKey' ) . mockImplementation ( ( message ) =>
435+ keys . get ( message )
436+ )
437+ vi . spyOn ( agentMemory . memoryService , 'fetchMemoryMessages' ) . mockImplementation ( async ( ) => [
438+ ...stored ,
439+ ] )
440+ const append = vi
441+ . spyOn ( agentMemory . memoryService , 'appendToMemory' )
442+ . mockImplementation ( async ( _ctx , _inputs , message , options ) => {
443+ stored . push ( message )
444+ if ( options ) {
445+ turns . set ( message , options . turnId )
446+ keys . set ( message , options . appendKey )
447+ }
448+ } )
449+ const first = {
450+ turnId : 'turn-1' ,
451+ memoryId : 'memory-1' ,
452+ finalize : vi . fn ( ) ,
453+ getFinalResponse : vi . fn ( ) ,
454+ getFinalAssistantContent : vi . fn ( ) ,
455+ }
456+ const second = {
457+ turnId : 'turn-2' ,
458+ memoryId : 'memory-1' ,
459+ finalize : vi . fn ( ) ,
460+ getFinalResponse : vi . fn ( ) ,
461+ getFinalAssistantContent : vi . fn ( ) ,
462+ }
463+ mockOpenAgentTurnSession
464+ . mockResolvedValueOnce ( first )
465+ . mockResolvedValueOnce ( first )
466+ . mockResolvedValueOnce ( second )
467+ mockExecuteProviderRequest . mockRejectedValueOnce ( new Error ( 'retry this block' ) )
468+ const ctx = { ...mockContext , executionId : 'execution-1' }
469+
470+ await expect (
471+ handler . execute ( ctx , mockBlock , inputs , { nodeId : 'agent-node' , executionOrder : 3 } )
472+ ) . rejects . toThrow ( 'retry this block' )
473+ await handler . execute ( ctx , mockBlock , inputs , { nodeId : 'agent-node' , executionOrder : 3 } )
474+ await handler . execute ( ctx , mockBlock , inputs , { nodeId : 'agent-node' , executionOrder : 4 } )
475+
476+ expect ( append . mock . calls . map ( ( call ) => [ call [ 3 ] ?. turnId , call [ 3 ] ?. appendKey ] ) ) . toEqual ( [
477+ [ 'turn-1' , 'seed:0' ] ,
478+ [ 'turn-1' , 'user-prompt' ] ,
479+ [ 'turn-2' , 'input' ] ,
480+ [ 'turn-2' , 'user-prompt' ] ,
481+ ] )
482+ expect ( mockExecuteProviderRequest . mock . calls [ 1 ] [ 1 ] . messages ) . toHaveLength ( 2 )
483+ expect ( mockExecuteProviderRequest . mock . calls [ 2 ] [ 1 ] . messages ) . toHaveLength ( 4 )
484+ } )
485+
486+ it ( 'persists the complete captured answer when structured output removes its content field' , async ( ) => {
487+ const content = '{"answer":42}'
488+ const session = {
489+ turnId : 'turn-1' ,
490+ memoryId : 'memory-1' ,
491+ finalize : vi . fn ( ) ,
492+ getFinalResponse : vi . fn ( ) ,
493+ getFinalAssistantContent : ( ) => content ,
494+ }
495+ mockOpenAgentTurnSession . mockResolvedValue ( session )
496+ vi . spyOn ( agentMemory . memoryService , 'fetchMemoryMessages' ) . mockResolvedValue ( [ ] )
497+ const append = vi . spyOn ( agentMemory . memoryService , 'appendToMemory' ) . mockResolvedValue ( )
498+ mockExecuteProviderRequest . mockResolvedValueOnce ( { content, model : 'gpt-4o' } )
499+ const result = await handler . execute (
500+ { ...mockContext , executionId : 'execution-1' } ,
501+ mockBlock ,
502+ {
503+ ...inputs ,
504+ responseFormat : { type : 'object' , properties : { answer : { type : 'number' } } } ,
505+ } ,
506+ { nodeId : 'agent-node' , executionOrder : 3 }
507+ )
508+ expect ( result ) . toMatchObject ( { answer : 42 } )
509+ expect ( result ) . not . toHaveProperty ( 'content' )
510+ expect ( append . mock . calls . every ( ( call ) => call [ 2 ] . role === 'user' ) ) . toBe ( true )
511+ expect ( session . finalize ) . toHaveBeenCalledWith ( content , 'gpt-4o' )
512+ } )
513+
514+ it ( 'finalizes a tool-only answer without creating an empty public assistant message' , async ( ) => {
515+ const session = {
516+ turnId : 'turn-1' ,
517+ memoryId : 'memory-1' ,
518+ finalize : vi . fn ( ) ,
519+ getFinalResponse : vi . fn ( ) ,
520+ getFinalAssistantContent : vi . fn ( ) ,
521+ }
522+ mockOpenAgentTurnSession . mockResolvedValue ( session )
523+ vi . spyOn ( agentMemory . memoryService , 'fetchMemoryMessages' ) . mockResolvedValue ( [ ] )
524+ const append = vi . spyOn ( agentMemory . memoryService , 'appendToMemory' ) . mockResolvedValue ( )
525+ mockExecuteProviderRequest . mockResolvedValueOnce ( { content : '' , model : 'gpt-4o' } )
526+ await handler . execute ( { ...mockContext , executionId : 'execution-1' } , mockBlock , inputs , {
527+ nodeId : 'agent-node' ,
528+ executionOrder : 3 ,
529+ } )
530+ expect ( append . mock . calls . every ( ( call ) => call [ 2 ] . role === 'user' ) ) . toBe ( true )
531+ expect ( session . finalize ) . toHaveBeenCalledWith ( '' , 'gpt-4o' )
532+ } )
533+
534+ it . each ( [ 'none' , 'deep-research-pro-preview-12-2025' , 'follow-up' ] ) (
535+ 'keeps %s on the existing provider lifecycle' ,
536+ async ( mode ) => {
537+ vi . spyOn ( agentMemory . memoryService , 'fetchMemoryMessages' ) . mockResolvedValue ( [ ] )
538+ vi . spyOn ( agentMemory . memoryService , 'seedMemory' ) . mockResolvedValue ( )
539+ vi . spyOn ( agentMemory . memoryService , 'appendToMemory' ) . mockResolvedValue ( )
540+ await handler . execute (
541+ { ...mockContext , executionId : 'execution-1' } ,
542+ mockBlock ,
543+ {
544+ ...inputs ,
545+ ...( mode === 'none' ? { memoryType : 'none' } : { } ) ,
546+ ...( mode . startsWith ( 'deep-research' ) ? { model : mode } : { } ) ,
547+ ...( mode === 'follow-up' ? { previousInteractionId : 'interaction-1' } : { } ) ,
548+ } ,
549+ { nodeId : 'agent-node' , executionOrder : 3 }
550+ )
551+ expect ( mockOpenAgentTurnSession ) . not . toHaveBeenCalled ( )
552+ expect ( mockExecuteProviderRequest . mock . calls [ 0 ] [ 2 ] . agentConversation ) . toBeUndefined ( )
553+ }
554+ )
555+ } )
556+
342557 describe ( 'conversation attachment replay' , ( ) => {
343558 beforeEach ( ( ) => {
344559 dbChainMockFns . returning . mockResolvedValue ( [ { id : 'memory-1' } ] )
@@ -5841,13 +6056,15 @@ describe('AgentBlockHandler', () => {
58416056 } )
58426057
58436058 describe ( 'wrapStreamForMemoryPersistence envelope' , ( ) => {
5844- it ( 'preserves streamFormat and subscribe via object spread' , ( ) => {
6059+ it ( 'preserves streamFormat, subscribe, and the existing completion callback' , async ( ) => {
58456060 const handler = new AgentBlockHandler ( )
58466061 const subscribe = vi . fn ( )
6062+ const onFullContent = vi . fn ( )
58476063 const streamingExec : StreamingExecution = {
58486064 stream : new ReadableStream ( ) ,
58496065 streamFormat : 'agent-events-v1' ,
58506066 subscribe,
6067+ onFullContent,
58516068 execution : {
58526069 success : true ,
58536070 output : { content : '' } ,
@@ -5871,6 +6088,8 @@ describe('AgentBlockHandler', () => {
58716088 expect ( wrapped . stream ) . toBe ( streamingExec . stream )
58726089 expect ( wrapped . execution ) . toBe ( streamingExec . execution )
58736090 expect ( typeof wrapped . onFullContent ) . toBe ( 'function' )
6091+ await wrapped . onFullContent ?.( '' )
6092+ expect ( onFullContent ) . toHaveBeenCalledWith ( '' )
58746093 } )
58756094 } )
58766095} )
0 commit comments