improvement(executor): correctness-by-construction for workflow logs (#4382)

* improvement(executor): correctness-by-construction for workflow logs

Replace the post-hoc reconciliation layer with a deterministic emission
protocol: drain pending callback promises at terminal boundaries, mint a
per-invocation blockExecutionId, and key console entries by that ID.
Eliminates races between block:* and execution:* events without changing
per-block latency.

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>

* test(executor): remove stale fire-and-forget assertion

The deleted test asserted wrappedOnBlockStart returns before the user
callback resolves. After the Stage 1 drain refactor, wrappedOnBlockStart
must await the user callback so the executor's trackCallback set covers
SSE writes — otherwise the drain at terminal-event time can't guarantee
block:* flushes before execution:*.

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>

* test(console-store): cover blockExecutionId keying and idempotency

Locks in the new console store invariants:
- Primary lookup via entryIdByBlockExecutionId fires no legacy warn
- Unknown blockExecutionId falls back to legacy keying and warns
- No-blockExecutionId updates use legacy path silently
- addConsole twice with same blockExecutionId returns the existing entry
- Distinct blockExecutionIds (loop iterations) produce distinct entries

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>

* fix(executor): capture BlockExecutor locally so finally drains its own instance

Previously this.blockExecutor was overwritten on every buildExecutionPipeline
call. Concurrent or re-entrant execute()/executeFromBlock() calls would have
their finally block drain the wrong instance, allowing the first execution's
block events to land after its terminal event. Returning { engine, blockExecutor }
and capturing both locally makes the drain pinned to the same instance the
engine.run() ran against.

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>

* test(terminal): capture store logger by label, not first warn-bearing mock

The previous capture used `.find` over all createLogger results, which
returned whichever module created a logger first — not the store's
logger — causing the legacy-keying warn assertion to see 0 calls.

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 4.7 <noreply@anthropic.com>
This commit is contained in:
Waleed
2026-05-01 15:54:53 -07:00
committed by GitHub
co-authored by Claude Opus 4.7
parent bdaf112db7
commit add55b4ffa
14 changed files with 384 additions and 458 deletions
@@ -919,7 +919,8 @@ async function handleExecutePost(
blockType: string,
executionOrder: number,
iterationContext?: IterationContext,
childWorkflowContext?: ChildWorkflowContext
childWorkflowContext?: ChildWorkflowContext,
blockExecutionId?: string
) => {
reqLogger.info('onBlockStart called', { blockId, blockName, blockType })
sendEvent({
@@ -945,6 +946,7 @@ async function handleExecutePost(
childWorkflowBlockId: childWorkflowContext.parentBlockId,
childWorkflowName: childWorkflowContext.workflowName,
}),
...(blockExecutionId && { blockExecutionId }),
},
})
}
@@ -955,7 +957,8 @@ async function handleExecutePost(
blockType: string,
callbackData: any,
iterationContext?: IterationContext,
childWorkflowContext?: ChildWorkflowContext
childWorkflowContext?: ChildWorkflowContext,
blockExecutionId?: string
) => {
const hasError = callbackData.output?.error
const childWorkflowData = childWorkflowContext
@@ -969,6 +972,11 @@ async function handleExecutePost(
? { childWorkflowInstanceId: callbackData.childWorkflowInstanceId }
: {}
const resolvedBlockExecutionId = blockExecutionId ?? callbackData.blockExecutionId
const blockExecData = resolvedBlockExecutionId
? { blockExecutionId: resolvedBlockExecutionId }
: {}
if (hasError) {
reqLogger.info('onBlockComplete (error) called', {
blockId,
@@ -1002,6 +1010,7 @@ async function handleExecutePost(
}),
...childWorkflowData,
...instanceData,
...blockExecData,
},
})
} else {
@@ -1036,6 +1045,7 @@ async function handleExecutePost(
}),
...childWorkflowData,
...instanceData,
...blockExecData,
},
})
}
@@ -1165,7 +1175,6 @@ async function handleExecutePost(
data: {
error: timeoutErrorMessage,
duration: result.metadata?.duration || 0,
finalBlockLogs: result.logs,
},
})
finalMetaStatus = 'error'
@@ -1179,7 +1188,6 @@ async function handleExecutePost(
workflowId,
data: {
duration: result.metadata?.duration || 0,
finalBlockLogs: result.logs,
},
})
finalMetaStatus = 'cancelled'
@@ -1244,7 +1252,6 @@ async function handleExecutePost(
data: {
error: executionResult?.error || errorMessage,
duration: executionResult?.metadata?.duration || 0,
finalBlockLogs: executionResult?.logs,
},
})
finalMetaStatus = 'error'
@@ -230,31 +230,25 @@ export function useWorkflowExecution() {
durationMs?: number
blockLogs: BlockLog[]
isPreExecutionError?: boolean
finalBlockLogs?: BlockLog[]
}) => {
if (!params.workflowId) return
sharedHandleExecutionErrorConsole(
{ addConsole, updateConsole, cancelRunningEntries },
{ addConsole, updateConsole },
{ ...params, workflowId: params.workflowId }
)
},
[addConsole, cancelRunningEntries, updateConsole]
[addConsole, updateConsole]
)
const handleExecutionCancelledConsole = useCallback(
(params: {
workflowId?: string
executionId?: string
durationMs?: number
finalBlockLogs?: BlockLog[]
}) => {
(params: { workflowId?: string; executionId?: string; durationMs?: number }) => {
if (!params.workflowId) return
sharedHandleExecutionCancelledConsole(
{ addConsole, updateConsole, cancelRunningEntries },
{ addConsole, updateConsole },
{ ...params, workflowId: params.workflowId }
)
},
[addConsole, cancelRunningEntries, updateConsole]
[addConsole, updateConsole]
)
const buildBlockEventHandlers = useCallback(
@@ -1036,8 +1030,6 @@ export function useWorkflowExecution() {
accumulatedBlockLogs,
accumulatedBlockStates,
executedBlockIds,
consoleMode: 'update',
includeStartConsoleEntry: true,
onBlockCompleteCallback: onBlockComplete,
})
@@ -1240,7 +1232,6 @@ export function useWorkflowExecution() {
durationMs: data.duration,
blockLogs: accumulatedBlockLogs,
isPreExecutionError,
finalBlockLogs: data.finalBlockLogs,
})
if (activeWorkflowId && !isExecutingFromChat) {
@@ -1267,7 +1258,6 @@ export function useWorkflowExecution() {
workflowId: activeWorkflowId,
executionId: executionIdRef.current,
durationMs: data?.duration,
finalBlockLogs: data?.finalBlockLogs,
})
if (activeWorkflowId && !isExecutingFromChat) {
@@ -1684,8 +1674,6 @@ export function useWorkflowExecution() {
accumulatedBlockLogs,
accumulatedBlockStates,
executedBlockIds,
consoleMode: 'update',
includeStartConsoleEntry: true,
})
await executionStream.executeFromBlock({
@@ -1755,7 +1743,6 @@ export function useWorkflowExecution() {
error: data.error,
durationMs: data.duration,
blockLogs: accumulatedBlockLogs,
finalBlockLogs: data.finalBlockLogs,
})
setCurrentExecutionId(workflowId, null)
@@ -1768,7 +1755,6 @@ export function useWorkflowExecution() {
workflowId,
executionId: executionIdRef.current,
durationMs: data?.duration,
finalBlockLogs: data?.finalBlockLogs,
})
setCurrentExecutionId(workflowId, null)
@@ -1915,8 +1901,6 @@ export function useWorkflowExecution() {
accumulatedBlockLogs,
accumulatedBlockStates,
executedBlockIds,
consoleMode: 'update',
includeStartConsoleEntry: true,
})
const capturedExecutionId = executionId
@@ -2017,7 +2001,6 @@ export function useWorkflowExecution() {
executionId: capturedExecutionId,
error: data.error,
blockLogs: accumulatedBlockLogs,
finalBlockLogs: data.finalBlockLogs,
})
},
onExecutionCancelled: (data) => {
@@ -2038,7 +2021,6 @@ export function useWorkflowExecution() {
workflowId: reconnectWorkflowId,
executionId: capturedExecutionId,
durationMs: data?.duration,
finalBlockLogs: data?.finalBlockLogs,
})
},
},
@@ -1,14 +1,12 @@
/**
* @vitest-environment node
*/
import { resetTerminalConsoleMock, terminalConsoleMockFns } from '@sim/testing'
import { resetTerminalConsoleMock } from '@sim/testing'
import { beforeEach, describe, expect, it, vi } from 'vitest'
import {
addExecutionErrorConsoleEntry,
handleExecutionErrorConsole,
reconcileFinalBlockLogs,
} from '@/app/workspace/[workspaceId]/w/[workflowId]/utils/workflow-execution-utils'
import type { BlockLog } from '@/executor/types'
describe('workflow-execution-utils', () => {
beforeEach(() => {
@@ -57,52 +55,6 @@ describe('workflow-execution-utils', () => {
expect(addConsole).not.toHaveBeenCalled()
})
it('skips when console store already has a block-level error for this execution (Fix D)', () => {
terminalConsoleMockFns.mockAddConsole({
workflowId: 'wf-1',
blockId: 'fetchAshbyData',
blockName: 'fetchAshbyData',
blockType: 'function',
executionId: 'exec-1',
executionOrder: 1,
success: false,
error: 'Failed to parse response as JSON',
})
const addConsole = vi.fn()
addExecutionErrorConsoleEntry(addConsole, {
workflowId: 'wf-1',
executionId: 'exec-1',
error: 'Run failed',
blockLogs: [],
})
expect(addConsole).not.toHaveBeenCalled()
})
it('still adds when only existing entries are themselves Run Error rows', () => {
terminalConsoleMockFns.mockAddConsole({
workflowId: 'wf-1',
blockId: 'execution-error',
blockName: 'Run Error',
blockType: 'error',
executionId: 'exec-1',
executionOrder: Number.MAX_SAFE_INTEGER,
success: false,
error: 'previous unrelated error',
})
const addConsole = vi.fn()
addExecutionErrorConsoleEntry(addConsole, {
workflowId: 'wf-1',
executionId: 'exec-1',
error: 'New run failed',
blockLogs: [],
})
expect(addConsole).toHaveBeenCalledTimes(1)
})
it('uses Timeout Error label when error indicates a timeout', () => {
const addConsole = vi.fn()
addExecutionErrorConsoleEntry(addConsole, {
@@ -131,121 +83,12 @@ describe('workflow-execution-utils', () => {
})
})
describe('reconcileFinalBlockLogs', () => {
const makeLog = (over: Partial<BlockLog>): BlockLog => ({
blockId: 'b1',
blockName: 'Function',
blockType: 'function',
startedAt: new Date().toISOString(),
endedAt: new Date().toISOString(),
durationMs: 50,
success: true,
executionOrder: 1,
...over,
})
it('flips a still-running entry to the server-reported success state', () => {
terminalConsoleMockFns.mockAddConsole({
workflowId: 'wf-1',
blockId: 'kb-1',
blockName: 'Knowledge 1',
blockType: 'knowledge',
executionId: 'exec-1',
executionOrder: 2,
isRunning: true,
})
const updateConsole = vi.fn()
reconcileFinalBlockLogs(updateConsole, 'wf-1', 'exec-1', [
makeLog({
blockId: 'kb-1',
blockName: 'Knowledge 1',
blockType: 'knowledge',
executionOrder: 2,
success: true,
output: { items: [] },
}),
])
expect(updateConsole).toHaveBeenCalledTimes(1)
const [blockId, update, executionId] = updateConsole.mock.calls[0]
expect(blockId).toBe('kb-1')
expect(executionId).toBe('exec-1')
expect(update).toMatchObject({
success: true,
isRunning: false,
replaceOutput: { items: [] },
})
})
it('flips a still-running entry to the server-reported error state (Bug 1 reconciliation)', () => {
terminalConsoleMockFns.mockAddConsole({
workflowId: 'wf-1',
blockId: 'fn-1',
blockName: 'Function',
blockType: 'function',
executionId: 'exec-1',
executionOrder: 3,
isRunning: true,
})
const updateConsole = vi.fn()
reconcileFinalBlockLogs(updateConsole, 'wf-1', 'exec-1', [
makeLog({
blockId: 'fn-1',
executionOrder: 3,
success: false,
error: 'JSON parse failed',
}),
])
expect(updateConsole).toHaveBeenCalledTimes(1)
expect(updateConsole.mock.calls[0][1]).toMatchObject({
success: false,
error: 'JSON parse failed',
isRunning: false,
})
})
it('skips entries that are not running', () => {
terminalConsoleMockFns.mockAddConsole({
workflowId: 'wf-1',
blockId: 'fn-1',
blockName: 'Function',
blockType: 'function',
executionId: 'exec-1',
executionOrder: 1,
isRunning: false,
success: true,
})
const updateConsole = vi.fn()
reconcileFinalBlockLogs(updateConsole, 'wf-1', 'exec-1', [makeLog({ blockId: 'fn-1' })])
expect(updateConsole).not.toHaveBeenCalled()
})
it('is a no-op when finalBlockLogs is empty or executionId is missing', () => {
const updateConsole = vi.fn()
reconcileFinalBlockLogs(updateConsole, 'wf-1', 'exec-1', [])
reconcileFinalBlockLogs(updateConsole, 'wf-1', undefined, [makeLog({})])
expect(updateConsole).not.toHaveBeenCalled()
})
})
describe('handleExecutionErrorConsole', () => {
it('cancels running entries before adding the synthetic entry', () => {
const calls: string[] = []
const addConsole = vi.fn(() => {
calls.push('add')
return undefined
})
const cancelRunningEntries = vi.fn(() => {
calls.push('cancel')
})
it('adds a synthetic Run Error entry when no block-level error covers it', () => {
const addConsole = vi.fn()
handleExecutionErrorConsole(
{ addConsole, updateConsole: vi.fn(), cancelRunningEntries },
{ addConsole, updateConsole: vi.fn() },
{
workflowId: 'wf-1',
executionId: 'exec-1',
@@ -254,57 +97,36 @@ describe('workflow-execution-utils', () => {
}
)
expect(calls[0]).toBe('cancel')
expect(calls).toContain('add')
expect(addConsole).toHaveBeenCalledTimes(1)
expect(addConsole.mock.calls[0][0].blockName).toBe('Run Error')
})
it('reconciles finalBlockLogs before sweeping running entries (Fix C)', () => {
terminalConsoleMockFns.mockAddConsole({
workflowId: 'wf-1',
blockId: 'kb-1',
blockName: 'Knowledge 1',
blockType: 'knowledge',
executionId: 'exec-1',
executionOrder: 1,
isRunning: true,
})
const calls: string[] = []
const addConsole = vi.fn(() => {
calls.push('add')
return undefined
})
const cancelRunningEntries = vi.fn(() => {
calls.push('cancel')
})
const updateConsole = vi.fn(() => {
calls.push('update')
})
it('skips the synthetic entry when a block-level error already covers the failure', () => {
const addConsole = vi.fn()
handleExecutionErrorConsole(
{ addConsole, updateConsole, cancelRunningEntries },
{ addConsole, updateConsole: vi.fn() },
{
workflowId: 'wf-1',
executionId: 'exec-1',
error: 'boom',
blockLogs: [],
finalBlockLogs: [
blockLogs: [
{
blockId: 'kb-1',
blockName: 'Knowledge 1',
blockType: 'knowledge',
blockId: 'fn-1',
blockName: 'Function',
blockType: 'function',
success: false,
error: 'JSON parse failed',
startedAt: new Date().toISOString(),
endedAt: new Date().toISOString(),
durationMs: 10,
success: true,
executionOrder: 1,
} as any,
],
}
)
expect(updateConsole).toHaveBeenCalledTimes(1)
expect(calls).toEqual(['update', 'cancel', 'add'])
expect(addConsole).not.toHaveBeenCalled()
})
})
})
@@ -5,13 +5,7 @@ import type {
BlockErrorData,
BlockStartedData,
} from '@/lib/workflows/executor/execution-events'
import type {
BlockLog,
BlockState,
ExecutionResult,
NormalizedBlockOutput,
StreamingExecution,
} from '@/executor/types'
import type { BlockLog, BlockState, ExecutionResult, StreamingExecution } from '@/executor/types'
import { stripCloneSuffixes } from '@/executor/utils/subflow-utils'
import { processSSEStream } from '@/hooks/use-execution-stream'
@@ -112,8 +106,6 @@ export interface BlockEventHandlerConfig {
accumulatedBlockLogs: BlockLog[]
accumulatedBlockStates: Map<string, BlockState>
executedBlockIds: Set<string>
consoleMode: 'update' | 'add'
includeStartConsoleEntry: boolean
onBlockCompleteCallback?: (blockId: string, output: unknown) => Promise<void>
}
@@ -142,8 +134,6 @@ export function createBlockEventHandlers(
accumulatedBlockLogs,
accumulatedBlockStates,
executedBlockIds,
consoleMode,
includeStartConsoleEntry,
onBlockCompleteCallback,
} = config
@@ -184,6 +174,7 @@ export function createBlockEventHandlers(
...('childWorkflowInstanceId' in data && {
childWorkflowInstanceId: data.childWorkflowInstanceId,
}),
...(data.blockExecutionId && { blockExecutionId: data.blockExecutionId }),
})
const createBlockLogEntry = (
@@ -203,61 +194,6 @@ export function createBlockEventHandlers(
endedAt: data.endedAt,
})
const addConsoleEntry = (data: BlockCompletedData, output: NormalizedBlockOutput) => {
if (!workflowId) return
addConsole({
input: data.input || {},
output,
success: true,
durationMs: data.durationMs,
startedAt: data.startedAt,
executionOrder: data.executionOrder,
endedAt: data.endedAt,
workflowId,
blockId: data.blockId,
executionId: executionIdRef.current,
blockName: data.blockName || 'Unknown Block',
blockType: data.blockType || 'unknown',
...extractIterationFields(data),
})
}
const addConsoleErrorEntry = (data: BlockErrorData) => {
if (!workflowId) return
const existingRunningEntry = useTerminalConsoleStore
.getState()
.getWorkflowEntries(workflowId)
.some(
(entry) =>
entry.blockId === data.blockId &&
entry.executionId === executionIdRef.current &&
entry.isRunning
)
if (existingRunningEntry) {
updateConsoleErrorEntry(data)
return
}
addConsole({
input: data.input || {},
output: {},
success: false,
error: data.error,
durationMs: data.durationMs,
startedAt: data.startedAt,
executionOrder: data.executionOrder,
endedAt: data.endedAt,
workflowId,
blockId: data.blockId,
executionId: executionIdRef.current,
blockName: data.blockName || 'Unknown Block',
blockType: data.blockType || 'unknown',
...extractIterationFields(data),
})
}
const updateConsoleEntry = (data: BlockCompletedData) => {
updateConsole(
data.blockId,
@@ -299,7 +235,7 @@ export function createBlockEventHandlers(
if (isStaleExecution()) return
updateActiveBlocks(data.blockId, true)
if (!includeStartConsoleEntry || !workflowId) return
if (!workflowId) return
const startedAt = new Date().toISOString()
addConsole({
@@ -344,20 +280,13 @@ export function createBlockEventHandlers(
const output = data.output as Record<string, any> | undefined
const isEmptySubflow = Array.isArray(output?.results) && output.results.length === 0
if (!isEmptySubflow) {
if (includeStartConsoleEntry) {
updateConsoleEntry(data)
}
updateConsoleEntry(data)
return
}
}
accumulatedBlockLogs.push(createBlockLogEntry(data, { success: true, output: data.output }))
if (consoleMode === 'update') {
updateConsoleEntry(data)
} else {
addConsoleEntry(data, data.output as NormalizedBlockOutput)
}
updateConsoleEntry(data)
if (onBlockCompleteCallback) {
onBlockCompleteCallback(data.blockId, data.output).catch((error) => {
@@ -391,11 +320,7 @@ export function createBlockEventHandlers(
createBlockLogEntry(data, { success: false, output: {}, error: data.error })
)
if (consoleMode === 'update') {
updateConsoleErrorEntry(data)
} else {
addConsoleErrorEntry(data)
}
updateConsoleErrorEntry(data)
}
const onBlockChildWorkflowStarted = (data: {
@@ -424,7 +349,6 @@ export function createBlockEventHandlers(
}
type AddConsoleFn = (entry: Omit<ConsoleEntry, 'id' | 'timestamp'>) => ConsoleEntry | undefined
type CancelRunningEntriesFn = (workflowId: string) => void
type UpdateConsoleFn = (
blockId: string,
update: string | ConsoleUpdate,
@@ -433,49 +357,10 @@ type UpdateConsoleFn = (
/**
* Bundle of console-store actions used by the execution-level handlers.
* Mirrors the deps-object pattern established by `createBlockEventHandlers`.
*/
export interface ExecutionConsoleDeps {
addConsole: AddConsoleFn
updateConsole: UpdateConsoleFn
cancelRunningEntries: CancelRunningEntriesFn
}
/**
* Reconciles still-running console entries with the server's authoritative
* `finalBlockLogs` so that any block whose terminal `block:completed`/`block:error`
* SSE event was lost gets the correct success/error state instead of being
* swept to "canceled".
*/
export function reconcileFinalBlockLogs(
updateConsole: UpdateConsoleFn,
workflowId: string,
executionId: string | undefined,
finalBlockLogs: BlockLog[] | undefined
): void {
if (!finalBlockLogs?.length || !executionId) return
const entries = useTerminalConsoleStore.getState().getWorkflowEntries(workflowId)
for (const log of finalBlockLogs) {
const running = entries.find(
(e) => e.blockId === log.blockId && e.executionId === executionId && e.isRunning
)
if (!running) continue
updateConsole(
log.blockId,
{
executionOrder: log.executionOrder,
replaceOutput: (log.output ?? {}) as Record<string, unknown>,
...(log.input ? { input: log.input } : {}),
success: log.success,
...(log.error ? { error: log.error } : {}),
durationMs: log.durationMs,
startedAt: log.startedAt,
endedAt: log.endedAt,
isRunning: false,
},
executionId
)
}
}
export interface ExecutionTimingFields {
@@ -503,8 +388,6 @@ export interface ExecutionErrorConsoleParams {
durationMs?: number
blockLogs: BlockLog[]
isPreExecutionError?: boolean
/** Server's authoritative per-block terminal states, used to reconcile lost SSE events. */
finalBlockLogs?: BlockLog[]
}
/**
@@ -515,19 +398,7 @@ export function addExecutionErrorConsoleEntry(
addConsole: AddConsoleFn,
params: ExecutionErrorConsoleParams
): void {
const hasBlockErrorInLogs = params.blockLogs.some((log) => log.error)
const hasBlockErrorInConsole = useTerminalConsoleStore
.getState()
.getWorkflowEntries(params.workflowId)
.some(
(entry) =>
entry.executionId === params.executionId &&
entry.error != null &&
entry.error !== '' &&
entry.blockType !== 'error' &&
entry.blockType !== 'validation'
)
const hasBlockError = hasBlockErrorInLogs || hasBlockErrorInConsole
const hasBlockError = params.blockLogs.some((log) => log.error)
const isPreExecutionError = params.isPreExecutionError ?? false
if (!isPreExecutionError && hasBlockError) return
@@ -557,21 +428,12 @@ export function addExecutionErrorConsoleEntry(
}
/**
* Reconciles `finalBlockLogs` against still-running entries, sweeps any
* remaining running entries to canceled, and adds an execution-level error
* console entry when no block-level error already covers it.
* Adds an execution-level error console entry when no block-level error already covers it.
*/
export function handleExecutionErrorConsole(
deps: ExecutionConsoleDeps,
params: ExecutionErrorConsoleParams
): void {
reconcileFinalBlockLogs(
deps.updateConsole,
params.workflowId,
params.executionId,
params.finalBlockLogs
)
deps.cancelRunningEntries(params.workflowId)
addExecutionErrorConsoleEntry(deps.addConsole, params)
}
@@ -612,8 +474,6 @@ export interface CancelledConsoleParams {
workflowId: string
executionId?: string
durationMs?: number
/** Server's authoritative per-block terminal states, used to reconcile lost SSE events. */
finalBlockLogs?: BlockLog[]
}
/**
@@ -642,21 +502,12 @@ export function addCancelledConsoleEntry(
}
/**
* Reconciles `finalBlockLogs` against still-running entries, sweeps any
* remaining running entries to canceled, and adds the execution-level
* cancellation console entry.
* Adds the execution-level cancellation console entry.
*/
export function handleExecutionCancelledConsole(
deps: ExecutionConsoleDeps,
params: CancelledConsoleParams
): void {
reconcileFinalBlockLogs(
deps.updateConsole,
params.workflowId,
params.executionId,
params.finalBlockLogs
)
deps.cancelRunningEntries(params.workflowId)
addCancelledConsoleEntry(deps.addConsole, params)
}
@@ -693,7 +544,7 @@ export async function executeWorkflowWithFullLogging(
}
const executionId = options.executionId || generateId()
const { addConsole, updateConsole, cancelRunningEntries } = useTerminalConsoleStore.getState()
const { addConsole, updateConsole } = useTerminalConsoleStore.getState()
const { setActiveBlocks, setBlockRunStatus, setEdgeRunStatus, setCurrentExecutionId } =
useExecutionStore.getState()
const wfId = targetWorkflowId
@@ -714,8 +565,6 @@ export async function executeWorkflowWithFullLogging(
accumulatedBlockLogs,
accumulatedBlockStates: new Map(),
executedBlockIds: new Set(),
consoleMode: 'update',
includeStartConsoleEntry: true,
onBlockCompleteCallback: options.onBlockComplete,
},
{ addConsole, updateConsole, setActiveBlocks, setBlockRunStatus, setEdgeRunStatus }
@@ -825,12 +674,11 @@ export async function executeWorkflowWithFullLogging(
}
handleExecutionCancelledConsole(
{ addConsole, updateConsole, cancelRunningEntries },
{ addConsole, updateConsole },
{
workflowId: wfId,
executionId: executionIdRef.current,
durationMs: data?.duration,
finalBlockLogs: data?.finalBlockLogs,
}
)
},
@@ -847,7 +695,7 @@ export async function executeWorkflowWithFullLogging(
}
handleExecutionErrorConsole(
{ addConsole, updateConsole, cancelRunningEntries },
{ addConsole, updateConsole },
{
workflowId: wfId,
executionId: executionIdRef.current,
@@ -855,7 +703,6 @@ export async function executeWorkflowWithFullLogging(
durationMs: data.duration || 0,
blockLogs: accumulatedBlockLogs,
isPreExecutionError: accumulatedBlockLogs.length === 0,
finalBlockLogs: data.finalBlockLogs,
}
)
},
+42 -8
View File
@@ -1,5 +1,6 @@
import { createLogger, type Logger } from '@sim/logger'
import { toError } from '@sim/utils/errors'
import { generateId } from '@sim/utils/id'
import { redactApiKeys } from '@/lib/core/security/redaction'
import { getBaseUrl } from '@/lib/core/utils/urls'
import {
@@ -52,6 +53,7 @@ const logger = createLogger('BlockExecutor')
export class BlockExecutor {
private execLogger: Logger
private pendingCallbacks: Set<Promise<void>> = new Set()
constructor(
private blockHandlers: BlockHandler[],
@@ -95,7 +97,13 @@ export class BlockExecutor {
if (!isSentinel) {
blockLog = this.createBlockLog(ctx, node.id, block, node, startedAt)
ctx.blockLogs.push(blockLog)
this.fireBlockStartCallback(ctx, node, block, blockLog.executionOrder)
this.fireBlockStartCallback(
ctx,
node,
block,
blockLog.executionOrder,
blockLog.blockExecutionId
)
}
let resolvedInputs: Record<string, any> = {}
@@ -207,7 +215,8 @@ export class BlockExecutor {
blockLog.startedAt,
blockLog.executionOrder,
blockLog.endedAt,
childWorkflowInstanceId
childWorkflowInstanceId,
blockLog.blockExecutionId
)
}
@@ -316,7 +325,8 @@ export class BlockExecutor {
blockLog.startedAt,
blockLog.executionOrder,
blockLog.endedAt,
childWorkflowInstanceId
childWorkflowInstanceId,
blockLog.blockExecutionId
)
}
@@ -390,6 +400,7 @@ export class BlockExecutor {
return {
blockId,
blockExecutionId: generateId(),
blockName,
blockType: block.metadata?.id ?? DEFAULTS.BLOCK_TYPE,
startedAt,
@@ -467,7 +478,8 @@ export class BlockExecutor {
ctx: ExecutionContext,
node: DAGNode,
block: SerializedBlock,
executionOrder: number
executionOrder: number,
blockExecutionId: string | undefined
): void {
if (!this.contextExtensions.onBlockStart) return
@@ -476,14 +488,15 @@ export class BlockExecutor {
const blockType = block.metadata?.id ?? DEFAULTS.BLOCK_TYPE
const iterationContext = getIterationContext(ctx, node?.metadata)
void this.contextExtensions
const promise = this.contextExtensions
.onBlockStart(
blockId,
blockName,
blockType,
executionOrder,
iterationContext,
ctx.childWorkflowContext
ctx.childWorkflowContext,
blockExecutionId
)
.catch((error) => {
this.execLogger.warn('Block start callback failed', {
@@ -492,6 +505,7 @@ export class BlockExecutor {
error: toError(error).message,
})
})
this.trackCallback(promise)
}
/**
@@ -509,7 +523,8 @@ export class BlockExecutor {
startedAt: string,
executionOrder: number,
endedAt: string,
childWorkflowInstanceId?: string
childWorkflowInstanceId?: string,
blockExecutionId?: string
): void {
if (!this.contextExtensions.onBlockComplete) return
@@ -518,7 +533,7 @@ export class BlockExecutor {
const blockType = block.metadata?.id ?? DEFAULTS.BLOCK_TYPE
const iterationContext = getIterationContext(ctx, node?.metadata)
void this.contextExtensions
const promise = this.contextExtensions
.onBlockComplete(
blockId,
blockName,
@@ -531,6 +546,7 @@ export class BlockExecutor {
executionOrder,
endedAt,
childWorkflowInstanceId,
blockExecutionId,
},
iterationContext,
ctx.childWorkflowContext
@@ -542,6 +558,24 @@ export class BlockExecutor {
error: toError(error).message,
})
})
this.trackCallback(promise)
}
private trackCallback(promise: Promise<void>): void {
this.pendingCallbacks.add(promise)
promise.finally(() => {
this.pendingCallbacks.delete(promise)
})
}
/**
* Resolves once every in-flight `onBlockStart` / `onBlockComplete` callback has settled.
* Drained at terminal-event boundaries so block events land before `execution:*` events.
*/
async awaitPendingCallbacks(): Promise<void> {
while (this.pendingCallbacks.size > 0) {
await Promise.allSettled([...this.pendingCallbacks])
}
}
private preparePauseResumeSelfReference(
+14 -5
View File
@@ -79,8 +79,12 @@ export class DAGExecutor {
const { context, state } = this.createExecutionContext(workflowId, triggerBlockId)
context.subflowParentMap = this.buildSubflowParentMap(dag)
const engine = this.buildExecutionPipeline(context, dag, state)
return await engine.run(triggerBlockId)
const { engine, blockExecutor } = this.buildExecutionPipeline(context, dag, state)
try {
return await engine.run(triggerBlockId)
} finally {
await blockExecutor.awaitPendingCallbacks()
}
}
async continueExecution(
@@ -206,8 +210,12 @@ export class DAGExecutor {
})
context.subflowParentMap = this.buildSubflowParentMap(dag)
const engine = this.buildExecutionPipeline(context, dag, state)
return await engine.run()
const { engine, blockExecutor } = this.buildExecutionPipeline(context, dag, state)
try {
return await engine.run()
} finally {
await blockExecutor.awaitPendingCallbacks()
}
}
private buildExecutionPipeline(context: ExecutionContext, dag: DAG, state: ExecutionState) {
@@ -235,7 +243,8 @@ export class DAGExecutor {
loopOrchestrator,
parallelOrchestrator
)
return new ExecutionEngine(context, dag, edgeManager, nodeOrchestrator)
const engine = new ExecutionEngine(context, dag, edgeManager, nodeOrchestrator)
return { engine, blockExecutor }
}
private createExecutionContext(
+8 -3
View File
@@ -118,7 +118,8 @@ export interface ExecutionCallbacks {
blockType: string,
executionOrder: number,
iterationContext?: IterationContext,
childWorkflowContext?: ChildWorkflowContext
childWorkflowContext?: ChildWorkflowContext,
blockExecutionId?: string
) => Promise<void>
onBlockComplete?: (
blockId: string,
@@ -126,7 +127,8 @@ export interface ExecutionCallbacks {
blockType: string,
output: any,
iterationContext?: IterationContext,
childWorkflowContext?: ChildWorkflowContext
childWorkflowContext?: ChildWorkflowContext,
blockExecutionId?: string
) => Promise<void>
/** Fires immediately after instanceId is generated, before child execution begins. */
onChildWorkflowInstanceReady?: (
@@ -172,7 +174,8 @@ export interface ContextExtensions {
blockType: string,
executionOrder: number,
iterationContext?: IterationContext,
childWorkflowContext?: ChildWorkflowContext
childWorkflowContext?: ChildWorkflowContext,
blockExecutionId?: string
) => Promise<void>
onBlockComplete?: (
blockId: string,
@@ -187,6 +190,8 @@ export interface ContextExtensions {
endedAt: string
/** Per-invocation unique ID linking this workflow block execution to its child block events. */
childWorkflowInstanceId?: string
/** Per-invocation unique ID for this block execution (distinct across loop/parallel iterations). */
blockExecutionId?: string
},
iterationContext?: IterationContext,
childWorkflowContext?: ChildWorkflowContext
+6
View File
@@ -208,6 +208,12 @@ export interface NormalizedBlockOutput {
export interface BlockLog {
blockId: string
/**
* Unique per-invocation ID. Same `blockId` can appear multiple times across loop/parallel
* iterations and across runs; `blockExecutionId` is unique for each individual execution
* and survives across `block:started` → `block:completed | block:error`.
*/
blockExecutionId?: string
blockName?: string
blockType?: string
startedAt: string
@@ -239,37 +239,6 @@ describe('executeWorkflowCore terminal finalization sequencing', () => {
expect(findStartBlockMock).toHaveBeenCalledWith(expect.anything(), 'external', false)
})
it('does not await user block start callback after persistence completes', async () => {
let releaseCallback: (() => void) | undefined
const callbackPromise = new Promise<void>((resolve) => {
releaseCallback = resolve
})
executorExecuteMock.mockResolvedValue({
success: true,
status: 'completed',
output: { done: true },
logs: [],
metadata: { duration: 123, startTime: 'start', endTime: 'end' },
})
await executeWorkflowCore({
snapshot: createSnapshot() as any,
callbacks: {
onBlockStart: vi.fn(() => callbackPromise),
},
loggingSession: loggingSession as any,
})
const contextExtensions = executorConstructorMock.mock.calls[0]?.[0]?.contextExtensions
await expect(
contextExtensions.onBlockStart('block-1', 'Fetch', 'api', 1)
).resolves.toBeUndefined()
releaseCallback?.()
})
it('awaits terminal completion before updating run counts and returning', async () => {
const callOrder: string[] = []
@@ -435,6 +435,7 @@ export async function executeWorkflowCore(
executionTime: number
startedAt: string
endedAt: string
blockExecutionId?: string
},
iterationContext?: IterationContext,
childWorkflowContext?: ChildWorkflowContext
@@ -442,13 +443,14 @@ export async function executeWorkflowCore(
try {
await loggingSession.onBlockComplete(blockId, blockName, blockType, output)
if (onBlockComplete) {
void onBlockComplete(
await onBlockComplete(
blockId,
blockName,
blockType,
output,
iterationContext,
childWorkflowContext
childWorkflowContext,
output.blockExecutionId
).catch((error) => {
logger.warn(`[${requestId}] Block completion callback failed`, {
executionId,
@@ -474,18 +476,20 @@ export async function executeWorkflowCore(
blockType: string,
executionOrder: number,
iterationContext?: IterationContext,
childWorkflowContext?: ChildWorkflowContext
childWorkflowContext?: ChildWorkflowContext,
blockExecutionId?: string
) => {
try {
await loggingSession.onBlockStart(blockId, blockName, blockType, new Date().toISOString())
if (onBlockStart) {
void onBlockStart(
await onBlockStart(
blockId,
blockName,
blockType,
executionOrder,
iterationContext,
childWorkflowContext
childWorkflowContext,
blockExecutionId
).catch((error) => {
logger.warn(`[${requestId}] Block start callback failed`, {
executionId,
@@ -3,7 +3,6 @@ import type {
IterationContext,
ParentIteration,
} from '@/executor/execution/types'
import type { BlockLog } from '@/executor/types'
import type { SubflowType } from '@/stores/workflows/workflow/types'
export type ExecutionEventType =
@@ -78,8 +77,6 @@ export interface ExecutionErrorEvent extends BaseExecutionEvent {
data: {
error: string
duration: number
/** Authoritative per-block terminal states from the server's blockLogs. */
finalBlockLogs?: BlockLog[]
}
}
@@ -88,8 +85,6 @@ export interface ExecutionCancelledEvent extends BaseExecutionEvent {
workflowId: string
data: {
duration: number
/** Authoritative per-block terminal states from the server's blockLogs. */
finalBlockLogs?: BlockLog[]
}
}
@@ -111,6 +106,8 @@ export interface BlockStartedEvent extends BaseExecutionEvent {
parentIterations?: ParentIteration[]
childWorkflowBlockId?: string
childWorkflowName?: string
/** Per-invocation unique ID for this block execution (distinct across loop/parallel iterations). */
blockExecutionId?: string
}
}
@@ -139,6 +136,8 @@ export interface BlockCompletedEvent extends BaseExecutionEvent {
childWorkflowName?: string
/** Per-invocation unique ID for correlating child block events with this workflow block. */
childWorkflowInstanceId?: string
/** Per-invocation unique ID for this block execution (distinct across loop/parallel iterations). */
blockExecutionId?: string
}
}
@@ -167,6 +166,8 @@ export interface BlockErrorEvent extends BaseExecutionEvent {
childWorkflowName?: string
/** Per-invocation unique ID for correlating child block events with this workflow block. */
childWorkflowInstanceId?: string
/** Per-invocation unique ID for this block execution (distinct across loop/parallel iterations). */
blockExecutionId?: string
}
}
@@ -283,7 +284,8 @@ export function createExecutionCallbacks(options: {
blockType: string,
executionOrder: number,
iterationContext?: IterationContext,
childWorkflowContext?: ChildWorkflowContext
childWorkflowContext?: ChildWorkflowContext,
blockExecutionId?: string
) => {
await sendBufferedEvent({
type: 'block:started',
@@ -308,6 +310,7 @@ export function createExecutionCallbacks(options: {
childWorkflowBlockId: childWorkflowContext.parentBlockId,
childWorkflowName: childWorkflowContext.workflowName,
}),
...(blockExecutionId && { blockExecutionId }),
},
})
}
@@ -324,9 +327,11 @@ export function createExecutionCallbacks(options: {
executionOrder: number
endedAt: string
childWorkflowInstanceId?: string
blockExecutionId?: string
},
iterationContext?: IterationContext,
childWorkflowContext?: ChildWorkflowContext
childWorkflowContext?: ChildWorkflowContext,
blockExecutionId?: string
) => {
const hasError = callbackData.output?.error
const iterationData = iterationContext
@@ -351,6 +356,11 @@ export function createExecutionCallbacks(options: {
? { childWorkflowInstanceId: callbackData.childWorkflowInstanceId }
: {}
const blockExecData =
blockExecutionId || callbackData.blockExecutionId
? { blockExecutionId: blockExecutionId ?? callbackData.blockExecutionId }
: {}
if (hasError) {
await sendBufferedEvent({
type: 'block:error',
@@ -370,6 +380,7 @@ export function createExecutionCallbacks(options: {
...iterationData,
...childWorkflowData,
...instanceData,
...blockExecData,
},
})
} else {
@@ -391,6 +402,7 @@ export function createExecutionCallbacks(options: {
...iterationData,
...childWorkflowData,
...instanceData,
...blockExecData,
},
})
}
@@ -1,6 +1,7 @@
/**
* @vitest-environment node
*/
import { createLogger } from '@sim/logger'
import { beforeEach, describe, expect, it, vi } from 'vitest'
vi.unmock('@/stores/terminal')
@@ -8,15 +9,25 @@ vi.unmock('@/stores/terminal/console/store')
import { useTerminalConsoleStore } from '@/stores/terminal/console/store'
const storeLoggerCallIdx = vi
.mocked(createLogger)
.mock.calls.findIndex((call) => call[0] === 'TerminalConsoleStore')
const storeLogger =
storeLoggerCallIdx >= 0
? vi.mocked(createLogger).mock.results[storeLoggerCallIdx]?.value
: undefined
describe('terminal console store', () => {
beforeEach(() => {
useTerminalConsoleStore.setState({
workflowEntries: {},
entryIdsByBlockExecution: {},
entryIdByBlockExecutionId: {},
entryLocationById: {},
isOpen: false,
_hasHydrated: true,
})
storeLogger?.warn.mockClear()
})
it('normalizes oversized payloads when adding console entries', () => {
@@ -117,6 +128,148 @@ describe('terminal console store', () => {
expect(after.getWorkflowEntries('wf-1')[0].output).toMatchObject({ status: 'updated' })
})
describe('blockExecutionId keying', () => {
it('updates an entry via the primary index without firing legacy warn', () => {
useTerminalConsoleStore.getState().addConsole({
workflowId: 'wf-1',
blockId: 'block-1',
blockName: 'Function',
blockType: 'function',
executionId: 'exec-1',
blockExecutionId: 'bex-1',
executionOrder: 1,
isRunning: true,
})
useTerminalConsoleStore.getState().updateConsole(
'block-1',
{
executionOrder: 1,
blockExecutionId: 'bex-1',
success: true,
replaceOutput: { status: 'done' },
},
'exec-1'
)
const [entry] = useTerminalConsoleStore.getState().getWorkflowEntries('wf-1')
expect(entry.success).toBe(true)
expect(entry.output).toMatchObject({ status: 'done' })
expect(storeLogger?.warn).not.toHaveBeenCalled()
})
it('falls back to legacy keying and warns when blockExecutionId is unknown', () => {
useTerminalConsoleStore.getState().addConsole({
workflowId: 'wf-1',
blockId: 'block-1',
blockName: 'Function',
blockType: 'function',
executionId: 'exec-1',
executionOrder: 1,
isRunning: true,
})
useTerminalConsoleStore.getState().updateConsole(
'block-1',
{
executionOrder: 1,
blockExecutionId: 'bex-unknown',
success: true,
replaceOutput: { status: 'done' },
},
'exec-1'
)
const [entry] = useTerminalConsoleStore.getState().getWorkflowEntries('wf-1')
expect(entry.success).toBe(true)
expect(storeLogger?.warn).toHaveBeenCalledWith(
'updateConsole used legacy keying (hydrated or cross-deploy entry)',
expect.objectContaining({ blockExecutionId: 'bex-unknown', blockId: 'block-1' })
)
})
it('uses legacy keying without warning when no blockExecutionId is provided', () => {
useTerminalConsoleStore.getState().addConsole({
workflowId: 'wf-1',
blockId: 'block-1',
blockName: 'Function',
blockType: 'function',
executionId: 'exec-1',
executionOrder: 1,
isRunning: true,
})
useTerminalConsoleStore
.getState()
.updateConsole(
'block-1',
{ executionOrder: 1, success: true, replaceOutput: { status: 'done' } },
'exec-1'
)
const [entry] = useTerminalConsoleStore.getState().getWorkflowEntries('wf-1')
expect(entry.success).toBe(true)
expect(storeLogger?.warn).not.toHaveBeenCalled()
})
})
describe('addConsole idempotency', () => {
it('returns the existing entry when called twice with the same blockExecutionId', () => {
const first = useTerminalConsoleStore.getState().addConsole({
workflowId: 'wf-1',
blockId: 'block-1',
blockName: 'Function',
blockType: 'function',
executionId: 'exec-1',
blockExecutionId: 'bex-1',
executionOrder: 1,
isRunning: true,
})
const second = useTerminalConsoleStore.getState().addConsole({
workflowId: 'wf-1',
blockId: 'block-1',
blockName: 'Function',
blockType: 'function',
executionId: 'exec-1',
blockExecutionId: 'bex-1',
executionOrder: 1,
isRunning: true,
})
const entries = useTerminalConsoleStore.getState().getWorkflowEntries('wf-1')
expect(entries).toHaveLength(1)
expect(second?.id).toBe(first?.id)
})
it('creates distinct entries for different blockExecutionIds (loop iterations)', () => {
useTerminalConsoleStore.getState().addConsole({
workflowId: 'wf-1',
blockId: 'block-1',
blockName: 'Function',
blockType: 'function',
executionId: 'exec-1',
blockExecutionId: 'bex-iter-1',
executionOrder: 1,
isRunning: true,
})
useTerminalConsoleStore.getState().addConsole({
workflowId: 'wf-1',
blockId: 'block-1',
blockName: 'Function',
blockType: 'function',
executionId: 'exec-1',
blockExecutionId: 'bex-iter-2',
executionOrder: 2,
isRunning: true,
})
const entries = useTerminalConsoleStore.getState().getWorkflowEntries('wf-1')
expect(entries).toHaveLength(2)
})
})
describe('cancelRunningEntries', () => {
it('flips a plain running entry to canceled', () => {
useTerminalConsoleStore.getState().addConsole({
+84 -12
View File
@@ -123,10 +123,14 @@ function removeWorkflowIndexes(
workflowId: string,
entries: ConsoleEntry[],
entryIdsByBlockExecution: Record<string, string[]>,
entryIdByBlockExecutionId: Record<string, string>,
entryLocationById: Record<string, ConsoleEntryLocation>
): void {
for (const entry of entries) {
delete entryLocationById[entry.id]
if (entry.blockExecutionId && entryIdByBlockExecutionId[entry.blockExecutionId] === entry.id) {
delete entryIdByBlockExecutionId[entry.blockExecutionId]
}
const blockExecutionKey = getBlockExecutionKey(entry.blockId, entry.executionId)
const existingIds = entryIdsByBlockExecution[blockExecutionKey]
if (!existingIds) {
@@ -146,10 +150,14 @@ function indexWorkflowEntries(
workflowId: string,
entries: ConsoleEntry[],
entryIdsByBlockExecution: Record<string, string[]>,
entryIdByBlockExecutionId: Record<string, string>,
entryLocationById: Record<string, ConsoleEntryLocation>
): void {
entries.forEach((entry, index) => {
entryLocationById[entry.id] = { workflowId, index }
if (entry.blockExecutionId) {
entryIdByBlockExecutionId[entry.blockExecutionId] = entry.id
}
const blockExecutionKey = getBlockExecutionKey(entry.blockId, entry.executionId)
const existingIds = entryIdsByBlockExecution[blockExecutionKey]
if (existingIds) {
@@ -162,35 +170,58 @@ function indexWorkflowEntries(
function rebuildWorkflowStateMaps(workflowEntries: Record<string, ConsoleEntry[]>) {
const entryIdsByBlockExecution: Record<string, string[]> = {}
const entryIdByBlockExecutionId: Record<string, string> = {}
const entryLocationById: Record<string, ConsoleEntryLocation> = {}
Object.entries(workflowEntries).forEach(([workflowId, entries]) => {
indexWorkflowEntries(workflowId, entries, entryIdsByBlockExecution, entryLocationById)
indexWorkflowEntries(
workflowId,
entries,
entryIdsByBlockExecution,
entryIdByBlockExecutionId,
entryLocationById
)
})
return { entryIdsByBlockExecution, entryLocationById }
return { entryIdsByBlockExecution, entryIdByBlockExecutionId, entryLocationById }
}
function replaceWorkflowEntries(
state: ConsoleStore,
workflowId: string,
nextEntries: ConsoleEntry[]
): Pick<ConsoleStore, 'workflowEntries' | 'entryIdsByBlockExecution' | 'entryLocationById'> {
): Pick<
ConsoleStore,
'workflowEntries' | 'entryIdsByBlockExecution' | 'entryIdByBlockExecutionId' | 'entryLocationById'
> {
const workflowEntries = cloneWorkflowEntries(state.workflowEntries)
const entryIdsByBlockExecution = { ...state.entryIdsByBlockExecution }
const entryIdByBlockExecutionId = { ...state.entryIdByBlockExecutionId }
const entryLocationById = { ...state.entryLocationById }
const previousEntries = workflowEntries[workflowId] ?? EMPTY_CONSOLE_ENTRIES
removeWorkflowIndexes(workflowId, previousEntries, entryIdsByBlockExecution, entryLocationById)
removeWorkflowIndexes(
workflowId,
previousEntries,
entryIdsByBlockExecution,
entryIdByBlockExecutionId,
entryLocationById
)
if (nextEntries.length === 0) {
delete workflowEntries[workflowId]
} else {
workflowEntries[workflowId] = nextEntries
indexWorkflowEntries(workflowId, nextEntries, entryIdsByBlockExecution, entryLocationById)
indexWorkflowEntries(
workflowId,
nextEntries,
entryIdsByBlockExecution,
entryIdByBlockExecutionId,
entryLocationById
)
}
return { workflowEntries, entryIdsByBlockExecution, entryLocationById }
return { workflowEntries, entryIdsByBlockExecution, entryIdByBlockExecutionId, entryLocationById }
}
function appendWorkflowEntry(
@@ -198,24 +229,38 @@ function appendWorkflowEntry(
workflowId: string,
newEntry: ConsoleEntry,
trimmedEntries: ConsoleEntry[]
): Pick<ConsoleStore, 'workflowEntries' | 'entryIdsByBlockExecution' | 'entryLocationById'> {
): Pick<
ConsoleStore,
'workflowEntries' | 'entryIdsByBlockExecution' | 'entryIdByBlockExecutionId' | 'entryLocationById'
> {
const workflowEntries = cloneWorkflowEntries(state.workflowEntries)
const previousEntries = workflowEntries[workflowId] ?? EMPTY_CONSOLE_ENTRIES
workflowEntries[workflowId] = trimmedEntries
const entryLocationById = { ...state.entryLocationById }
const entryIdsByBlockExecution = { ...state.entryIdsByBlockExecution }
const entryIdByBlockExecutionId = { ...state.entryIdByBlockExecutionId }
const survivingIds = new Set(trimmedEntries.map((e) => e.id))
const droppedEntries = previousEntries.filter((e) => !survivingIds.has(e.id))
if (droppedEntries.length > 0) {
removeWorkflowIndexes(workflowId, droppedEntries, entryIdsByBlockExecution, entryLocationById)
removeWorkflowIndexes(
workflowId,
droppedEntries,
entryIdsByBlockExecution,
entryIdByBlockExecutionId,
entryLocationById
)
}
trimmedEntries.forEach((entry, index) => {
entryLocationById[entry.id] = { workflowId, index }
})
if (newEntry.blockExecutionId) {
entryIdByBlockExecutionId[newEntry.blockExecutionId] = newEntry.id
}
const blockExecutionKey = getBlockExecutionKey(newEntry.blockId, newEntry.executionId)
const existingIds = entryIdsByBlockExecution[blockExecutionKey]
if (existingIds) {
@@ -226,7 +271,7 @@ function appendWorkflowEntry(
entryIdsByBlockExecution[blockExecutionKey] = [newEntry.id]
}
return { workflowEntries, entryIdsByBlockExecution, entryLocationById }
return { workflowEntries, entryIdsByBlockExecution, entryIdByBlockExecutionId, entryLocationById }
}
interface NotifyBlockErrorParams {
@@ -269,6 +314,7 @@ export const useTerminalConsoleStore = create<ConsoleStore>()(
devtools((set, get) => ({
workflowEntries: {},
entryIdsByBlockExecution: {},
entryIdByBlockExecutionId: {},
entryLocationById: {},
isOpen: false,
_hasHydrated: false,
@@ -278,6 +324,16 @@ export const useTerminalConsoleStore = create<ConsoleStore>()(
return get().getWorkflowEntries(entry.workflowId)[0] as ConsoleEntry | undefined
}
if (entry.blockExecutionId) {
const existingId = get().entryIdByBlockExecutionId[entry.blockExecutionId]
if (existingId) {
const location = get().entryLocationById[existingId]
if (location) {
return get().workflowEntries[location.workflowId]?.[location.index]
}
}
}
const redactedEntry = { ...entry }
if (
!isStreamingOutput(entry.output) &&
@@ -441,11 +497,23 @@ export const useTerminalConsoleStore = create<ConsoleStore>()(
updateConsole: (blockId: string, update: string | ConsoleUpdate, executionId?: string) => {
set((state) => {
const candidateIds =
state.entryIdsByBlockExecution[getBlockExecutionKey(blockId, executionId)] ?? []
const blockExecutionId = typeof update === 'object' ? update.blockExecutionId : undefined
const directId = blockExecutionId
? state.entryIdByBlockExecutionId[blockExecutionId]
: undefined
const candidateIds = directId
? [directId]
: (state.entryIdsByBlockExecution[getBlockExecutionKey(blockId, executionId)] ?? [])
if (candidateIds.length === 0) {
return state
}
if (blockExecutionId && !directId) {
logger.warn('updateConsole used legacy keying (hydrated or cross-deploy entry)', {
blockExecutionId,
blockId,
executionId,
})
}
const workflowId = state.entryLocationById[candidateIds[0]]?.workflowId
if (!workflowId) {
@@ -462,7 +530,7 @@ export const useTerminalConsoleStore = create<ConsoleStore>()(
const source = nextEntries ?? currentEntries
const entry = source[location.index]
if (!entry || entry.id !== candidateId) continue
if (!matchesEntryForUpdate(entry, blockId, executionId, update)) continue
if (!directId && !matchesEntryForUpdate(entry, blockId, executionId, update)) continue
if (!nextEntries) {
nextEntries = [...currentEntries]
@@ -570,6 +638,10 @@ export const useTerminalConsoleStore = create<ConsoleStore>()(
updatedEntry.childWorkflowInstanceId = update.childWorkflowInstanceId
}
if (update.blockExecutionId !== undefined) {
updatedEntry.blockExecutionId = update.blockExecutionId
}
nextEntries[location.index] = updatedEntry
}
@@ -32,6 +32,8 @@ export interface ConsoleEntry {
childWorkflowName?: string
/** Per-invocation unique ID linking this workflow block to its child block events */
childWorkflowInstanceId?: string
/** Per-invocation unique ID for this block execution (distinct across loop/parallel iterations) */
blockExecutionId?: string
}
export interface ConsoleUpdate {
@@ -56,6 +58,7 @@ export interface ConsoleUpdate {
childWorkflowBlockId?: string
childWorkflowName?: string
childWorkflowInstanceId?: string
blockExecutionId?: string
}
export interface ConsoleEntryLocation {
@@ -66,6 +69,7 @@ export interface ConsoleEntryLocation {
export interface ConsoleStore {
workflowEntries: Record<string, ConsoleEntry[]>
entryIdsByBlockExecution: Record<string, string[]>
entryIdByBlockExecutionId: Record<string, string>
entryLocationById: Record<string, ConsoleEntryLocation>
isOpen: boolean
addConsole: (entry: Omit<ConsoleEntry, 'id' | 'timestamp'>) => ConsoleEntry | undefined