improvement(func-exec): normalize inputs to match schema (#4473)

This commit is contained in:
Vikhyath Mondreti
2026-05-06 10:35:38 -07:00
committed by GitHub
parent 7b6aa728ae
commit 1989a12ff3
18 changed files with 325 additions and 50 deletions
@@ -1,6 +1,7 @@
import { createLogger, type Logger } from '@sim/logger'
import { toError } from '@sim/utils/errors'
import { redactApiKeys } from '@/lib/core/security/redaction'
import { normalizeStringArray } from '@/lib/core/utils/arrays'
import { getBaseUrl } from '@/lib/core/utils/urls'
import {
containsUserFileWithMetadata,
@@ -164,7 +165,7 @@ export class BlockExecutor {
block,
streamingExec,
resolvedInputs,
ctx.selectedOutputs ?? []
normalizeStringArray(ctx.selectedOutputs)
)
}
+5 -3
View File
@@ -1,4 +1,6 @@
import { createLogger, type Logger } from '@sim/logger'
import { normalizeStringArray } from '@/lib/core/utils/arrays'
import { normalizeStringRecord, normalizeWorkflowVariables } from '@/lib/core/utils/records'
import { StartBlockPath } from '@/lib/workflows/triggers/triggers'
import type { DAG } from '@/executor/dag/builder'
import { DAGBuilder } from '@/executor/dag/builder'
@@ -56,9 +58,9 @@ export class DAGExecutor {
constructor(options: DAGExecutorOptions) {
this.workflow = options.workflow
this.environmentVariables = options.envVarValues ?? {}
this.environmentVariables = normalizeStringRecord(options.envVarValues)
this.workflowInput = options.workflowInput ?? {}
this.workflowVariables = options.workflowVariables ?? {}
this.workflowVariables = normalizeWorkflowVariables(options.workflowVariables)
this.contextExtensions = options.contextExtensions ?? {}
this.dagBuilder = new DAGBuilder()
this.execLogger = logger.withMetadata({
@@ -325,7 +327,7 @@ export class DAGExecutor {
: new Set(),
workflow: this.workflow,
stream: this.contextExtensions.stream ?? false,
selectedOutputs: this.contextExtensions.selectedOutputs ?? [],
selectedOutputs: normalizeStringArray(this.contextExtensions.selectedOutputs),
edges: this.contextExtensions.edges ?? [],
onStream: this.contextExtensions.onStream,
onBlockStart: this.contextExtensions.onBlockStart,
@@ -119,8 +119,8 @@ export function serializePauseSnapshot(
executionMetadata,
context.workflow,
{},
context.workflowVariables ?? {},
context.selectedOutputs ?? [],
context.workflowVariables,
context.selectedOutputs,
state
)
@@ -0,0 +1,30 @@
import { describe, expect, it } from 'vitest'
import { ExecutionSnapshot } from '@/executor/execution/snapshot'
import type { ExecutionMetadata } from '@/executor/execution/types'
const metadata: ExecutionMetadata = {
requestId: 'request-1',
executionId: 'execution-1',
workflowId: 'workflow-1',
workspaceId: 'workspace-1',
userId: 'user-1',
triggerType: 'manual',
startTime: '2026-05-06T00:00:00.000Z',
}
describe('ExecutionSnapshot', () => {
it('normalizes untyped persisted execution state at construction', () => {
const variable = { id: 'var-1', name: 'brand', type: 'plain', value: 'myfitness' }
const snapshot = new ExecutionSnapshot(
metadata,
{ blocks: [] },
{},
[variable],
['agent.content', 123, 'function.result']
)
expect(snapshot.workflowVariables).toEqual({ 'var-1': variable })
expect(snapshot.selectedOutputs).toEqual(['agent.content', 'function.result'])
})
})
+23 -7
View File
@@ -1,14 +1,30 @@
import { normalizeStringArray } from '@/lib/core/utils/arrays'
import { normalizeWorkflowVariables } from '@/lib/core/utils/records'
import type { ExecutionMetadata, SerializableExecutionState } from '@/executor/execution/types'
export class ExecutionSnapshot {
public readonly metadata: ExecutionMetadata
public readonly workflow: any
public readonly input: any
public readonly workflowVariables: Record<string, any>
public readonly selectedOutputs: string[]
public readonly state?: SerializableExecutionState
constructor(
public readonly metadata: ExecutionMetadata,
public readonly workflow: any,
public readonly input: any,
public readonly workflowVariables: Record<string, any>,
public readonly selectedOutputs: string[] = [],
public readonly state?: SerializableExecutionState
) {}
metadata: ExecutionMetadata,
workflow: any,
input: any,
workflowVariables: unknown,
selectedOutputs: unknown = [],
state?: SerializableExecutionState
) {
this.metadata = metadata
this.workflow = workflow
this.input = input
this.workflowVariables = normalizeWorkflowVariables(workflowVariables)
this.selectedOutputs = normalizeStringArray(selectedOutputs)
this.state = state
}
toJSON(): string {
return JSON.stringify({
@@ -4,6 +4,7 @@ import { createLogger } from '@sim/logger'
import { toError } from '@sim/utils/errors'
import { sleep } from '@sim/utils/helpers'
import { and, eq, inArray, isNull } from 'drizzle-orm'
import { normalizeStringRecord, normalizeWorkflowVariables } from '@/lib/core/utils/records'
import { createMcpToolId } from '@/lib/mcp/utils'
import { getCustomToolById } from '@/lib/workflows/custom-tools/operations'
import { getAllBlocks } from '@/blocks'
@@ -815,8 +816,8 @@ export class AgentBlockHandler implements BlockHandler {
userId: ctx.userId,
stream: streaming,
messages: messages?.map(({ executionId, ...msg }) => msg),
environmentVariables: ctx.environmentVariables || {},
workflowVariables: ctx.workflowVariables || {},
environmentVariables: normalizeStringRecord(ctx.environmentVariables),
workflowVariables: normalizeWorkflowVariables(ctx.workflowVariables),
blockData,
blockNameMapping,
reasoningEffort: inputs.reasoningEffort,
@@ -885,8 +886,8 @@ export class AgentBlockHandler implements BlockHandler {
userId: ctx.userId,
stream: providerRequest.stream,
messages: 'messages' in providerRequest ? providerRequest.messages : undefined,
environmentVariables: ctx.environmentVariables || {},
workflowVariables: ctx.workflowVariables || {},
environmentVariables: normalizeStringRecord(ctx.environmentVariables),
workflowVariables: normalizeWorkflowVariables(ctx.workflowVariables),
blockData,
blockNameMapping,
isDeployedContext: ctx.isDeployedContext,
@@ -1,4 +1,5 @@
import { createLogger } from '@sim/logger'
import { normalizeStringRecord, normalizeWorkflowVariables } from '@/lib/core/utils/records'
import type { BlockOutput } from '@/blocks/types'
import { BlockType, CONDITION, DEFAULTS, EDGE } from '@/executor/constants'
import type { BlockHandler, ExecutionContext } from '@/executor/types'
@@ -40,8 +41,8 @@ export async function evaluateConditionExpression(
{
code,
timeout: CONDITION_TIMEOUT_MS,
envVars: ctx.environmentVariables || {},
workflowVariables: ctx.workflowVariables || {},
envVars: normalizeStringRecord(ctx.environmentVariables),
workflowVariables: normalizeWorkflowVariables(ctx.workflowVariables),
blockData,
blockNameMapping,
blockOutputSchemas,
@@ -196,6 +196,28 @@ describe('FunctionBlockHandler', () => {
)
})
it('should normalize malformed execution context records before calling function_execute', async () => {
const legacyVariable = { id: 'var-1', name: 'brand', type: 'plain', value: 'myfitness' }
mockContext.workflowVariables = [legacyVariable] as unknown as Record<string, any>
mockContext.environmentVariables = ['invalid-env'] as unknown as Record<string, string>
await handler.execute(mockContext, mockBlock, {
code: 'return "myfitness"',
[FUNCTION_BLOCK_CONTEXT_VARS_KEY]: ['invalid-context'],
})
expect(mockExecuteTool).toHaveBeenCalledWith(
'function_execute',
expect.objectContaining({
envVars: {},
workflowVariables: { 'var-1': legacyVariable },
contextVariables: {},
}),
false,
mockContext
)
})
it('should handle tool error with no specific message', async () => {
const inputs = { code: 'some code' }
const errorResult = { success: false }
@@ -1,3 +1,8 @@
import {
normalizeRecord,
normalizeStringRecord,
normalizeWorkflowVariables,
} from '@/lib/core/utils/records'
import { DEFAULT_EXECUTION_TIMEOUT_MS } from '@/lib/execution/constants'
import { DEFAULT_CODE_LANGUAGE } from '@/lib/execution/languages'
import { BlockType } from '@/executor/constants'
@@ -26,8 +31,7 @@ export class FunctionBlockHandler implements BlockHandler {
const { blockData, blockNameMapping, blockOutputSchemas } = collectBlockData(ctx)
const contextVariables =
(inputs[FUNCTION_BLOCK_CONTEXT_VARS_KEY] as Record<string, unknown> | undefined) ?? {}
const contextVariables = normalizeRecord(inputs[FUNCTION_BLOCK_CONTEXT_VARS_KEY])
const result = await executeTool(
'function_execute',
@@ -35,8 +39,8 @@ export class FunctionBlockHandler implements BlockHandler {
code: codeContent,
language: inputs.language || DEFAULT_CODE_LANGUAGE,
timeout: inputs.timeout || DEFAULT_EXECUTION_TIMEOUT_MS,
envVars: ctx.environmentVariables || {},
workflowVariables: ctx.workflowVariables || {},
envVars: normalizeStringRecord(ctx.environmentVariables),
workflowVariables: normalizeWorkflowVariables(ctx.workflowVariables),
blockData,
blockNameMapping,
blockOutputSchemas,
+13
View File
@@ -0,0 +1,13 @@
import { describe, expect, it } from 'vitest'
import { normalizeStringArray } from '@/lib/core/utils/arrays'
describe('array normalization utilities', () => {
it('normalizes string arrays loaded from untyped state', () => {
expect(normalizeStringArray(['output-1', 2, 'output-2', null])).toEqual([
'output-1',
'output-2',
])
expect(normalizeStringArray('output-1')).toEqual([])
expect(normalizeStringArray(undefined)).toEqual([])
})
})
+10
View File
@@ -0,0 +1,10 @@
/**
* Normalizes optional string-list values loaded from untyped persisted state.
*/
export function normalizeStringArray(value: unknown): string[] {
if (!Array.isArray(value)) {
return []
}
return value.filter((item): item is string => typeof item === 'string')
}
+64
View File
@@ -0,0 +1,64 @@
import { describe, expect, it } from 'vitest'
import {
isPlainRecord,
normalizeRecord,
normalizeRecordMap,
normalizeStringRecord,
normalizeWorkflowVariables,
} from '@/lib/core/utils/records'
describe('record normalization utilities', () => {
it('identifies plain records without accepting arrays or null', () => {
expect(isPlainRecord({})).toBe(true)
expect(isPlainRecord(Object.create(null))).toBe(true)
expect(isPlainRecord([])).toBe(false)
expect(isPlainRecord(null)).toBe(false)
})
it('normalizes unknown values to object records', () => {
expect(normalizeRecord({ value: 1 })).toEqual({ value: 1 })
expect(normalizeRecord([])).toEqual({})
expect(normalizeRecord('not-a-record')).toEqual({})
})
it('normalizes string records for environment-like values', () => {
expect(
normalizeStringRecord({
TOKEN: 'secret',
RETRIES: 3,
ENABLED: true,
EMPTY: null,
})
).toEqual({
TOKEN: 'secret',
RETRIES: '3',
ENABLED: 'true',
})
expect(normalizeStringRecord([])).toEqual({})
})
it('normalizes record maps by dropping malformed entries', () => {
expect(
normalizeRecordMap({
valid: { type: 'string' },
invalid: [],
})
).toEqual({
valid: { type: 'string' },
})
})
it('normalizes legacy workflow variable arrays into records', () => {
const variableWithId = { id: 'var-1', name: 'brand', type: 'plain', value: 'myfitness' }
const variableWithName = { name: 'channel', type: 'plain', value: 'whatsapp' }
expect(normalizeWorkflowVariables([variableWithId, variableWithName, []])).toEqual({
'var-1': variableWithId,
channel: variableWithName,
})
expect(normalizeWorkflowVariables({ existing: variableWithId })).toEqual({
existing: variableWithId,
})
expect(normalizeWorkflowVariables('not-a-record')).toEqual({})
})
})
+90
View File
@@ -0,0 +1,90 @@
export type UnknownRecord = Record<string, unknown>
export type StringRecord = Record<string, string>
/**
* Returns true only for object-map values, excluding arrays and null.
*/
export function isPlainRecord(value: unknown): value is UnknownRecord {
if (typeof value !== 'object' || value === null || Array.isArray(value)) {
return false
}
const prototype = Object.getPrototypeOf(value)
return prototype === Object.prototype || prototype === null
}
/**
* Normalizes optional execution context maps to the record shape expected by
* internal API contracts.
*/
export function normalizeRecord(value: unknown): UnknownRecord {
return isPlainRecord(value) ? value : {}
}
/**
* Normalizes environment-like maps to string values, matching process/env
* semantics at execution boundaries.
*/
export function normalizeStringRecord(value: unknown): StringRecord {
if (!isPlainRecord(value)) {
return {}
}
const normalized: StringRecord = {}
for (const [key, entryValue] of Object.entries(value)) {
if (entryValue === undefined || entryValue === null) {
continue
}
normalized[key] = typeof entryValue === 'string' ? entryValue : String(entryValue)
}
return normalized
}
/**
* Normalizes record-of-record maps such as block output schema maps.
*/
export function normalizeRecordMap(value: unknown): Record<string, UnknownRecord> {
if (!isPlainRecord(value)) {
return {}
}
const normalized: Record<string, UnknownRecord> = {}
for (const [key, entryValue] of Object.entries(value)) {
if (isPlainRecord(entryValue)) {
normalized[key] = entryValue
}
}
return normalized
}
/**
* Workflow variables are stored as a record in current state, while some
* legacy and imported snapshots can carry an array of variable objects.
*/
export function normalizeWorkflowVariables(value: unknown): UnknownRecord {
if (isPlainRecord(value)) {
return value
}
if (!Array.isArray(value)) {
return {}
}
const normalized: UnknownRecord = {}
for (const variable of value) {
if (!isPlainRecord(variable)) {
continue
}
const id = typeof variable.id === 'string' && variable.id.trim() ? variable.id : undefined
const name =
typeof variable.name === 'string' && variable.name.trim() ? variable.name : undefined
const key = id ?? name
if (key) {
normalized[key] = variable
}
}
return normalized
}
@@ -7,6 +7,7 @@ import { createLogger } from '@sim/logger'
import { mergeSubblockStateWithValues } from '@sim/workflow-persistence/subblocks'
import type { Edge } from 'reactflow'
import { z } from 'zod'
import { isPlainRecord } from '@/lib/core/utils/records'
import { getPersonalAndWorkspaceEnv } from '@/lib/environment/utils'
import { clearExecutionCancellation } from '@/lib/execution/cancellation'
import type { LoggingSession } from '@/lib/logs/execution/logging-session'
@@ -581,6 +582,16 @@ export async function executeWorkflowCore(
callChain: metadata.callChain,
}
for (const variable of Object.values(workflowVariables)) {
if (
isPlainRecord(variable) &&
variable.value !== undefined &&
typeof variable.type === 'string'
) {
variable.value = parseVariableValueByType(variable.value, variable.type)
}
}
const executorInstance = new Executor({
workflow: serializedWorkflow,
envVarValues: decryptedEnvVars,
@@ -589,16 +600,6 @@ export async function executeWorkflowCore(
contextExtensions,
})
// Convert initial workflow variables to their native types
if (workflowVariables) {
for (const [varId, variable] of Object.entries(workflowVariables)) {
const v = variable as { value?: unknown; type?: string }
if (v.value !== undefined && v.type) {
v.value = parseVariableValueByType(v.value, v.type)
}
}
}
const result = runFromBlock
? ((await executorInstance.executeFromBlock(
workflowId,
@@ -138,7 +138,7 @@ export async function executeQueuedWorkflowJob(
payload.workflow,
payload.input,
payload.variables,
payload.selectedOutputs ?? []
payload.selectedOutputs
)
let callbacks = {}
+15 -4
View File
@@ -5,6 +5,11 @@ import type { CompletionUsage } from 'openai/resources/completions'
import { dollarsToCredits } from '@/lib/billing/credits/conversion'
import { env } from '@/lib/core/config/env'
import { getBlacklistedProvidersFromEnv, isHosted } from '@/lib/core/config/feature-flags'
import {
normalizeRecord,
normalizeStringRecord,
normalizeWorkflowVariables,
} from '@/lib/core/utils/records'
import {
buildCanonicalIndex,
type CanonicalGroup,
@@ -1166,10 +1171,16 @@ export function prepareToolExecution(
},
}
: {}),
...(request.environmentVariables ? { envVars: request.environmentVariables } : {}),
...(request.workflowVariables ? { workflowVariables: request.workflowVariables } : {}),
...(request.blockData ? { blockData: request.blockData } : {}),
...(request.blockNameMapping ? { blockNameMapping: request.blockNameMapping } : {}),
...(request.environmentVariables
? { envVars: normalizeStringRecord(request.environmentVariables) }
: {}),
...(request.workflowVariables
? { workflowVariables: normalizeWorkflowVariables(request.workflowVariables) }
: {}),
...(request.blockData ? { blockData: normalizeRecord(request.blockData) } : {}),
...(request.blockNameMapping
? { blockNameMapping: normalizeStringRecord(request.blockNameMapping) }
: {}),
...(tool.parameters ? { _toolSchema: tool.parameters } : {}),
}
+12 -6
View File
@@ -1,3 +1,9 @@
import {
normalizeRecord,
normalizeRecordMap,
normalizeStringRecord,
normalizeWorkflowVariables,
} from '@/lib/core/utils/records'
import { DEFAULT_EXECUTION_TIMEOUT_MS } from '@/lib/execution/constants'
import { DEFAULT_CODE_LANGUAGE } from '@/lib/execution/languages'
import type { CodeExecutionInput, CodeExecutionOutput } from '@/tools/function/types'
@@ -123,12 +129,12 @@ export const functionExecuteTool: ToolConfig<CodeExecutionInput, CodeExecutionOu
outputTable: params.outputTable,
outputSandboxPath: params.outputSandboxPath,
outputMimeType: params.outputMimeType,
envVars: params.envVars || {},
workflowVariables: params.workflowVariables || {},
blockData: params.blockData || {},
blockNameMapping: params.blockNameMapping || {},
blockOutputSchemas: params.blockOutputSchemas || {},
contextVariables: params.contextVariables || {},
envVars: normalizeStringRecord(params.envVars),
workflowVariables: normalizeWorkflowVariables(params.workflowVariables),
blockData: normalizeRecord(params.blockData),
blockNameMapping: normalizeStringRecord(params.blockNameMapping),
blockOutputSchemas: normalizeRecordMap(params.blockOutputSchemas),
contextVariables: normalizeRecord(params.contextVariables),
workflowId: params._context?.workflowId,
userId: params._context?.userId,
workspaceId: params._context?.workspaceId,
+9 -6
View File
@@ -1,5 +1,10 @@
import { createLogger } from '@sim/logger'
import { getMaxExecutionTimeout } from '@/lib/core/execution-limits'
import {
normalizeRecord,
normalizeStringRecord,
normalizeWorkflowVariables,
} from '@/lib/core/utils/records'
import type { EnvironmentVariable } from '@/lib/environment/api'
import { getQueryClient } from '@/app/_shell/providers/get-query-client'
import type { CustomToolDefinition } from '@/hooks/queries/custom-tools'
@@ -254,14 +259,12 @@ export function createCustomToolRequestBody(customTool: any, isClient = true, wo
// 1. envVars parameter (passed from provider/agent context)
// 2. Client-side store (if running in browser)
// 3. Empty object (fallback)
const envVars = params.envVars || (isClient ? getClientEnvVars() : {})
const envVars = normalizeStringRecord(params.envVars || (isClient ? getClientEnvVars() : {}))
// Get workflow variables from params (passed from execution context)
const workflowVariables = params.workflowVariables || {}
const workflowVariables = normalizeWorkflowVariables(params.workflowVariables)
// Get block data and mapping from params (passed from execution context)
const blockData = params.blockData || {}
const blockNameMapping = params.blockNameMapping || {}
const blockData = normalizeRecord(params.blockData)
const blockNameMapping = normalizeStringRecord(params.blockNameMapping)
// Include everything needed for execution
return {