6bf8bebf51
CI / Test and Build (push) Failing after 1s
CI / Migrate Dev DB (push) Has been skipped
CI / Migrate DB (push) Has been skipped
CodeQL / Analyze actions (push) Has been cancelled
CodeQL / Analyze javascript-typescript (push) Has been cancelled
CI / Detect Version (push) Has been cancelled
CI / Detect Desktop Changes (push) Has been cancelled
CI / Build AMD64 (blacksmith-2vcpu-ubuntu-2404, ./docker/cron.Dockerfile, ubuntu-latest, ghcr.io/simstudioai/cron) (push) Has been cancelled
CI / Build AMD64 (blacksmith-2vcpu-ubuntu-2404, ./docker/db.Dockerfile, ECR_MIGRATIONS, ubuntu-latest, ghcr.io/simstudioai/migrations) (push) Has been cancelled
CI / Build AMD64 (blacksmith-4vcpu-ubuntu-2404, ./docker/pii.Dockerfile, ECR_PII, ubuntu-latest, ghcr.io/simstudioai/pii) (push) Has been cancelled
CI / Build AMD64 (blacksmith-4vcpu-ubuntu-2404, ./docker/realtime.Dockerfile, ECR_REALTIME, ubuntu-latest, ghcr.io/simstudioai/realtime) (push) Has been cancelled
CI / Build AMD64 (blacksmith-8vcpu-ubuntu-2404, ./docker/app.Dockerfile, ECR_APP, linux-x64-8-core, ghcr.io/simstudioai/simstudio) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/cron.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/cron) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/db.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/migrations) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/pii.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/pii) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/realtime.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/realtime) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-8vcpu-ubuntu-2404-arm, ./docker/app.Dockerfile, linux-arm64-8-core, ghcr.io/simstudioai/simstudio) (push) Has been cancelled
CI / Check Docs Changes (push) Has been cancelled
Publish CLI Package / publish-npm (push) Has been cancelled
Publish Python SDK / publish-pypi (push) Has been cancelled
CI / Deploy Trigger.dev (Dev) (push) Has been cancelled
Helm Chart / Lint, test, and validate chart (push) Has been cancelled
Helm Chart / Chart version bumped (push) Has been cancelled
Publish TypeScript SDK / publish-npm (push) Has been cancelled
CI / Build Dev ECR (blacksmith-8vcpu-ubuntu-2404, ./docker/app.Dockerfile, ECR_APP, linux-x64-8-core) (push) Has been cancelled
CI / Promote Images (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/cron) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/migrations) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/pii) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/realtime) (push) Has been cancelled
CI / Build Dev ECR (blacksmith-2vcpu-ubuntu-2404, ./docker/db.Dockerfile, ECR_MIGRATIONS, ubuntu-latest) (push) Has been cancelled
CI / Build Dev ECR (blacksmith-4vcpu-ubuntu-2404, ./docker/pii.Dockerfile, ECR_PII, ubuntu-latest) (push) Has been cancelled
CI / Build Dev ECR (blacksmith-4vcpu-ubuntu-2404, ./docker/realtime.Dockerfile, ECR_REALTIME, ubuntu-latest) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/simstudio) (push) Has been cancelled
CI / Process Docs (push) Has been cancelled
CI / Create GitHub Release (push) Has been cancelled
CI / Check Desktop Signing Secrets (push) Has been cancelled
CI / Desktop Release (push) Has been cancelled
CI / Create Desktop Prerelease (push) Has been cancelled
CI / Desktop Prerelease Build (push) Has been cancelled
CI / Publish Desktop Prerelease (push) Has been cancelled
CI / Prune Desktop Prereleases (push) Has been cancelled
Helm Chart / Install on kind and run helm test (push) Has been cancelled
193 lines
6.2 KiB
TypeScript
193 lines
6.2 KiB
TypeScript
import type { BlockStateController } from '@/executor/execution/types'
|
|
import type { BlockState, NormalizedBlockOutput } from '@/executor/types'
|
|
import type { ResolvedSecretTraceProvenanceV1 } from '@/executor/utils/resolved-secret-trace-registry'
|
|
import { SubflowNodeIdCodec } from '@/executor/utils/subflow-node-id-codec'
|
|
import {
|
|
buildOuterBranchScopedId,
|
|
extractOuterBranchIndex,
|
|
stripCloneSuffixes,
|
|
} from '@/executor/utils/subflow-utils'
|
|
|
|
function normalizeLookupId(id: string): string {
|
|
return SubflowNodeIdCodec.normalizeLookupId(id)
|
|
}
|
|
|
|
function extractBranchSuffix(id: string): string {
|
|
return SubflowNodeIdCodec.extractBranchSuffix(id)
|
|
}
|
|
|
|
function extractLoopSuffix(id: string): string {
|
|
return SubflowNodeIdCodec.extractLoopSuffix(id)
|
|
}
|
|
export interface LoopScope {
|
|
iteration: number
|
|
currentIterationOutputs: Map<string, NormalizedBlockOutput>
|
|
allIterationOutputs: NormalizedBlockOutput[][]
|
|
maxIterations?: number
|
|
item?: any
|
|
items?: any[]
|
|
condition?: string
|
|
loopType?: 'for' | 'forEach' | 'while' | 'doWhile'
|
|
skipFirstConditionCheck?: boolean
|
|
skippedAtStart?: boolean
|
|
/** Error message if loop validation failed (e.g., exceeded max iterations) */
|
|
validationError?: string
|
|
inputResolvedSecretTraceProvenance?: ResolvedSecretTraceProvenanceV1
|
|
resolvedSecretTraceProvenance?: ResolvedSecretTraceProvenanceV1
|
|
}
|
|
|
|
export interface ParallelScope {
|
|
parallelId: string
|
|
totalBranches: number
|
|
batchSize?: number
|
|
currentBatchStart?: number
|
|
currentBatchSize?: number
|
|
accumulatedOutputs?: Map<number, NormalizedBlockOutput[]>
|
|
branchOutputs: Map<number, NormalizedBlockOutput[]>
|
|
items?: any[]
|
|
/** Error message if parallel validation failed (e.g., exceeded max branches) */
|
|
validationError?: string
|
|
/** Whether the parallel has an empty distribution and should be skipped */
|
|
isEmpty?: boolean
|
|
inputResolvedSecretTraceProvenance?: ResolvedSecretTraceProvenanceV1
|
|
resolvedSecretTraceProvenance?: ResolvedSecretTraceProvenanceV1
|
|
}
|
|
|
|
export class ExecutionState implements BlockStateController {
|
|
private readonly blockStates: Map<string, BlockState>
|
|
private readonly executedBlocks: Set<string>
|
|
|
|
constructor(blockStates?: Map<string, BlockState>, executedBlocks?: Set<string>) {
|
|
this.blockStates = blockStates ?? new Map()
|
|
this.executedBlocks = executedBlocks ?? new Set()
|
|
}
|
|
|
|
getBlockStates(): ReadonlyMap<string, BlockState> {
|
|
return this.blockStates
|
|
}
|
|
|
|
getExecutedBlocks(): ReadonlySet<string> {
|
|
return this.executedBlocks
|
|
}
|
|
|
|
getBlockState(blockId: string, currentNodeId?: string): BlockState | undefined {
|
|
const normalizedId = normalizeLookupId(blockId)
|
|
if (normalizedId !== blockId) {
|
|
return this.blockStates.get(blockId)
|
|
}
|
|
|
|
if (currentNodeId) {
|
|
const scopedState = this.getScopedBlockState(blockId, currentNodeId)
|
|
if (scopedState !== undefined) {
|
|
return scopedState
|
|
}
|
|
|
|
if (extractOuterBranchIndex(currentNodeId) !== undefined) {
|
|
return undefined
|
|
}
|
|
}
|
|
|
|
const direct = this.blockStates.get(blockId)
|
|
if (direct !== undefined) {
|
|
return direct
|
|
}
|
|
|
|
if (currentNodeId && extractBranchSuffix(currentNodeId) === '') {
|
|
const stableBranchZeroState = this.blockStates.get(buildOuterBranchScopedId(blockId, 0))
|
|
if (stableBranchZeroState !== undefined) {
|
|
return stableBranchZeroState
|
|
}
|
|
|
|
const branchZeroState = this.blockStates.get(
|
|
`${blockId}₍0₎${extractLoopSuffix(currentNodeId)}`
|
|
)
|
|
if (branchZeroState !== undefined) {
|
|
return branchZeroState
|
|
}
|
|
}
|
|
|
|
for (const [storedId, state] of this.blockStates.entries()) {
|
|
if (normalizeLookupId(storedId) === blockId) {
|
|
return state
|
|
}
|
|
}
|
|
|
|
return undefined
|
|
}
|
|
|
|
getBlockOutput(blockId: string, currentNodeId?: string): NormalizedBlockOutput | undefined {
|
|
return this.getBlockState(blockId, currentNodeId)?.output
|
|
}
|
|
|
|
private getScopedBlockState(blockId: string, currentNodeId: string): BlockState | undefined {
|
|
const currentBranchSuffix = extractBranchSuffix(currentNodeId)
|
|
const loopSuffix = extractLoopSuffix(currentNodeId)
|
|
|
|
const currentOuterBranchIndex = extractOuterBranchIndex(currentNodeId)
|
|
if (currentOuterBranchIndex !== undefined) {
|
|
for (const [storedId, state] of this.blockStates.entries()) {
|
|
if (stripCloneSuffixes(storedId) !== blockId) continue
|
|
if (extractOuterBranchIndex(storedId) !== currentOuterBranchIndex) continue
|
|
if (extractBranchSuffix(storedId) !== currentBranchSuffix) continue
|
|
if (extractLoopSuffix(storedId) !== loopSuffix) continue
|
|
|
|
return state
|
|
}
|
|
|
|
const siblingBranchState = this.blockStates.get(`${blockId}₍${currentOuterBranchIndex}₎`)
|
|
if (siblingBranchState !== undefined) {
|
|
return siblingBranchState
|
|
}
|
|
} else {
|
|
const withSuffix = `${blockId}${currentBranchSuffix}${loopSuffix}`
|
|
const suffixedState = this.blockStates.get(withSuffix)
|
|
if (suffixedState !== undefined) {
|
|
return suffixedState
|
|
}
|
|
}
|
|
|
|
return undefined
|
|
}
|
|
|
|
setBlockOutput(
|
|
blockId: string,
|
|
output: NormalizedBlockOutput,
|
|
executionTime = 0,
|
|
resolvedSecretTraceProvenance?: BlockState['resolvedSecretTraceProvenance']
|
|
): void {
|
|
const existingState = this.blockStates.get(blockId)
|
|
const effectiveProvenance =
|
|
resolvedSecretTraceProvenance ??
|
|
(existingState?.output === output ? existingState.resolvedSecretTraceProvenance : undefined)
|
|
this.blockStates.set(blockId, {
|
|
output,
|
|
executed: true,
|
|
executionTime,
|
|
...(effectiveProvenance ? { resolvedSecretTraceProvenance: effectiveProvenance } : {}),
|
|
})
|
|
this.executedBlocks.add(blockId)
|
|
}
|
|
|
|
setBlockState(blockId: string, state: BlockState): void {
|
|
this.blockStates.set(blockId, state)
|
|
if (state.executed) {
|
|
this.executedBlocks.add(blockId)
|
|
} else {
|
|
this.executedBlocks.delete(blockId)
|
|
}
|
|
}
|
|
|
|
deleteBlockState(blockId: string): void {
|
|
this.blockStates.delete(blockId)
|
|
this.executedBlocks.delete(blockId)
|
|
}
|
|
|
|
unmarkExecuted(blockId: string): void {
|
|
this.executedBlocks.delete(blockId)
|
|
}
|
|
|
|
hasExecuted(blockId: string): boolean {
|
|
return this.executedBlocks.has(blockId)
|
|
}
|
|
}
|