Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
41 changes: 41 additions & 0 deletions apps/sim/executor/execution/block-executor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1352,6 +1352,47 @@ describe('BlockExecutor streaming pump', () => {
)
})

it('projects a selected structured string into live sink deltas', async () => {
const handler = createAgentEventsStreamingHandler({
events: [
{ type: 'text_delta', text: '{"answer":"Hello ', turn: 'pending' },
{ type: 'text_delta', text: 'world","score":1}', turn: 'pending' },
{ type: 'turn_end', turn: 'final' },
],
})
const { executor, block, state } = createExecutor(handler)
block.config.params = {
responseFormat: {
schema: {
type: 'object',
properties: { answer: { type: 'string' }, score: { type: 'number' } },
},
},
}
const ctx = createContext(state)
ctx.stream = true
ctx.selectedOutputs = [`${block.id}_answer`]
const sinkText: string[] = []
let forwardedText = ''

ctx.onStream = async (streamingExec) => {
expect(streamingExec.clientStreamTransformed).toBe(true)
expect(streamingExec.clientSinkTransformed).toBe(true)
streamingExec.subscribe?.({
onEvent: (event) => {
if (event.type === 'text_delta') sinkText.push(event.text)
},
})
forwardedText = await new Response(streamingExec.stream).text()
}

await executor.execute(ctx, createNode(block), block)

expect(sinkText).toEqual(['Hello ', 'world'])
expect(forwardedText).toBe('Hello world')
expect(state.getBlockOutput(block.id)?.answer).toBe('Hello world')
})

it('forwards the stable block ID for streams from expanded branch nodes', async () => {
const handler = createAgentEventsStreamingHandler({
events: [{ type: 'text_delta', text: 'branch answer', turn: 'final' }],
Expand Down
14 changes: 12 additions & 2 deletions apps/sim/executor/execution/block-executor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1197,6 +1197,15 @@ export class BlockExecutor {
selectedOutputs,
responseFormat
)
const clientStreamTransformed = processedClientStream !== pump.textStream
const projectedSubscribe = clientStreamTransformed
? streamingResponseFormatProcessor.processEventSubscription(
pump.subscribe,
blockId,
selectedOutputs,
responseFormat
)
: undefined

// Start onStream without awaiting so a sync `subscribe(sink)` can run before
// the first provider pull, then read the projected text stream concurrently
Expand All @@ -1208,10 +1217,11 @@ export class BlockExecutor {
...(executionOrder !== undefined ? { executionOrder } : {}),
stream: processedClientStream,
streamFormat: 'text',
subscribe: pump.subscribe,
subscribe: projectedSubscribe ?? pump.subscribe,
// processStream returns the input stream identity when no
// response-format extraction applies.
clientStreamTransformed: processedClientStream !== pump.textStream,
clientStreamTransformed,
clientSinkTransformed: Boolean(projectedSubscribe),
displayResolvedSecretTraceProvenance:
ctx.resolvedSecretTraceRegistry?.exportCommittedProvenanceForValue(resolvedInputs),
})
Expand Down
3 changes: 3 additions & 0 deletions apps/sim/executor/handlers/workflow/workflow-handler.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2233,6 +2233,9 @@ describe('WorkflowBlockHandler', () => {
const childStream = {
blockId: 'agent-1',
stream: new ReadableStream(),
subscribe: vi.fn(),
clientStreamTransformed: true,
clientSinkTransformed: true,
execution: { success: true, output: {} },
}
await extensions.onStream(childStream)
Expand Down
5 changes: 3 additions & 2 deletions apps/sim/executor/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -666,10 +666,11 @@ export interface StreamingExecution {
/**
* True when {@link stream} is a response-format projection (selected JSON
* fields extracted from structured output) rather than raw answer text. Sink
* `text_delta` events then do NOT match the byte stream, so consumers must
* keep sourcing answer text from {@link stream} instead of the sink.
* `text_delta` events match it only when {@link clientSinkTransformed} is true.
*/
clientStreamTransformed?: boolean
/** True when sink text deltas are projected to match a transformed client stream. */
clientSinkTransformed?: boolean
/** Internal provenance for the exact block input that initiated this live stream. */
displayResolvedSecretTraceProvenance?: ResolvedSecretTraceProvenanceV1
/** Internal source registry retained only for sanitizing failures while the stream drains. */
Expand Down
114 changes: 114 additions & 0 deletions apps/sim/executor/utils.test.ts
Original file line number Diff line number Diff line change
@@ -1,8 +1,10 @@
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
import { formatInternalOutputSelector } from '@/lib/workflows/streaming/output-selector'
import {
StreamingResponseFormatProcessor,
streamingResponseFormatProcessor,
} from '@/executor/utils'
import type { AgentStreamEvent, AgentStreamSink } from '@/providers/stream-events'

describe('StreamingResponseFormatProcessor', () => {
let processor: StreamingResponseFormatProcessor
Expand Down Expand Up @@ -170,6 +172,118 @@ describe('StreamingResponseFormatProcessor', () => {
expect(result).toBe('charlie')
})

it('projects a selected string field from live event deltas', async () => {
let sourceSink: AgentStreamSink | undefined
const projectedEvents: AgentStreamEvent[] = []
const subscribe = processor.processEventSubscription(
(sink) => {
sourceSink = sink
return () => {}
},
'block-1',
['block-1_answer'],
JSON.stringify({
type: 'object',
properties: {
meta: { type: 'object' },
answer: { type: 'string' },
score: { type: 'number' },
},
})
)

expect(subscribe).toBeDefined()
subscribe?.({
onEvent: async (event) => {
projectedEvents.push(event)
},
})

await sourceSink?.onEvent({
type: 'text_delta',
text: '{"meta":{"source":"test"},"answer":"Hello ',
turn: 'pending',
})
expect(projectedEvents).toEqual([{ type: 'text_delta', text: 'Hello ', turn: 'pending' }])

await sourceSink?.onEvent({
type: 'text_delta',
text: 'world","score":1}',
turn: 'pending',
})
await sourceSink?.onEvent({ type: 'thinking_delta', text: 'done' })
await sourceSink?.onEvent({ type: 'turn_end', turn: 'final' })

expect(projectedEvents).toEqual([
{ type: 'text_delta', text: 'Hello ', turn: 'pending' },
{ type: 'text_delta', text: 'world', turn: 'pending' },
{ type: 'thinking_delta', text: 'done' },
{ type: 'turn_end', turn: 'final' },
])
})

it('matches encoded selectors for block IDs containing underscores', async () => {
let sourceSink: AgentStreamSink | undefined
const projectedText: string[] = []
const subscribe = processor.processEventSubscription(
(sink) => {
sourceSink = sink
return () => {}
},
'answer_agent',
[formatInternalOutputSelector('answer_agent', 'answer')],
{ schema: { properties: { answer: { type: 'string' } } } }
)

subscribe?.({
onEvent: (event) => {
if (event.type === 'text_delta') projectedText.push(event.text)
},
})
await sourceSink?.onEvent({
type: 'text_delta',
text: '{"answer":"Streamed"}',
turn: 'final',
})

expect(projectedText).toEqual(['Streamed'])
})

it('holds incomplete JSON escapes until they can be decoded', async () => {
let sourceSink: AgentStreamSink | undefined
const projectedText: string[] = []
const subscribe = processor.processEventSubscription(
(sink) => {
sourceSink = sink
return () => {}
},
'block-1',
['block-1_answer'],
{ schema: { properties: { answer: { type: 'string' } } } }
)

subscribe?.({
onEvent: (event) => {
if (event.type === 'text_delta') projectedText.push(event.text)
},
})
await sourceSink?.onEvent({ type: 'text_delta', text: '{"answer":"line\\', turn: 'final' })
await sourceSink?.onEvent({ type: 'text_delta', text: 'nnext"}', turn: 'final' })

expect(projectedText).toEqual(['line', '\nnext'])
})

it('does not claim live sink projection for non-string fields', () => {
const subscribe = processor.processEventSubscription(
() => () => {},
'block-1',
['block-1_score'],
{ schema: { properties: { score: { type: 'number' } } } }
)

expect(subscribe).toBeUndefined()
})

it.concurrent('should handle missing fields gracefully', async () => {
const mockStream = new ReadableStream({
start(controller) {
Expand Down
Loading
Loading