improvement(db): remove deprecated 'state' column from workflow table (#994)

* improvement(db): remove deprecated  column from workflow table

* removed extraneous logs

* update sockets envvar
This commit is contained in:
Waleed Latif
2025-08-16 13:04:49 -07:00
committed by GitHub
parent 133a32e6d3
commit 8748e1d5f9
17 changed files with 6026 additions and 720 deletions
@@ -80,7 +80,6 @@ export async function POST(request: NextRequest, { params }: { params: Promise<{
workspaceId: workspaceId,
name: `${templateData.name} (copy)`,
description: templateData.description,
state: templateData.state,
color: templateData.color,
userId: session.user.id,
createdAt: now,
@@ -158,9 +157,6 @@ export async function POST(request: NextRequest, { params }: { params: Promise<{
}))
}
// Update the workflow with the corrected state
await tx.update(workflow).set({ state: updatedState }).where(eq(workflow.id, newWorkflowId))
// Insert blocks and edges
if (blockEntries.length > 0) {
await tx.insert(workflowBlocks).values(blockEntries)
@@ -7,7 +7,7 @@ import { createLogger } from '@/lib/logs/console/logger'
import { getUserEntityPermissions } from '@/lib/permissions/utils'
import { db } from '@/db'
import { workflow, workflowBlocks, workflowEdges, workflowSubflows } from '@/db/schema'
import type { LoopConfig, ParallelConfig, WorkflowState } from '@/stores/workflows/workflow/types'
import type { LoopConfig, ParallelConfig } from '@/stores/workflows/workflow/types'
const logger = createLogger('WorkflowDuplicateAPI')
@@ -90,7 +90,6 @@ export async function POST(req: NextRequest, { params }: { params: Promise<{ id:
folderId: folderId || source.folderId,
name,
description: description || source.description,
state: source.state, // We'll update this later with new block IDs
color: color || source.color,
lastSynced: now,
createdAt: now,
@@ -112,9 +111,6 @@ export async function POST(req: NextRequest, { params }: { params: Promise<{ id:
// Create a mapping from old block IDs to new block IDs
const blockIdMapping = new Map<string, string>()
// Initialize state for updating with new block IDs
let updatedState: WorkflowState = source.state as WorkflowState
if (sourceBlocks.length > 0) {
// First pass: Create all block ID mappings
sourceBlocks.forEach((block) => {
@@ -265,86 +261,10 @@ export async function POST(req: NextRequest, { params }: { params: Promise<{ id:
)
}
// Update the JSON state to use new block IDs
if (updatedState && typeof updatedState === 'object') {
updatedState = JSON.parse(JSON.stringify(updatedState)) as WorkflowState
// Update blocks object keys
if (updatedState.blocks && typeof updatedState.blocks === 'object') {
const newBlocks = {} as Record<string, (typeof updatedState.blocks)[string]>
for (const [oldId, blockData] of Object.entries(updatedState.blocks)) {
const newId = blockIdMapping.get(oldId) || oldId
newBlocks[newId] = {
...blockData,
id: newId,
// Update data.parentId and extent in the JSON state as well
data: (() => {
const block = blockData as any
if (block.data && typeof block.data === 'object' && block.data.parentId) {
return {
...block.data,
parentId: blockIdMapping.get(block.data.parentId) || block.data.parentId,
extent: 'parent', // Ensure extent is set for child blocks
}
}
return block.data
})(),
}
}
updatedState.blocks = newBlocks
}
// Update edges array
if (updatedState.edges && Array.isArray(updatedState.edges)) {
updatedState.edges = updatedState.edges.map((edge) => ({
...edge,
id: crypto.randomUUID(),
source: blockIdMapping.get(edge.source) || edge.source,
target: blockIdMapping.get(edge.target) || edge.target,
}))
}
// Update loops and parallels if they exist
if (updatedState.loops && typeof updatedState.loops === 'object') {
const newLoops = {} as Record<string, (typeof updatedState.loops)[string]>
for (const [oldId, loopData] of Object.entries(updatedState.loops)) {
const newId = blockIdMapping.get(oldId) || oldId
const loopConfig = loopData as any
newLoops[newId] = {
...loopConfig,
id: newId,
// Update node references in loop config
nodes: loopConfig.nodes
? loopConfig.nodes.map((nodeId: string) => blockIdMapping.get(nodeId) || nodeId)
: [],
}
}
updatedState.loops = newLoops
}
if (updatedState.parallels && typeof updatedState.parallels === 'object') {
const newParallels = {} as Record<string, (typeof updatedState.parallels)[string]>
for (const [oldId, parallelData] of Object.entries(updatedState.parallels)) {
const newId = blockIdMapping.get(oldId) || oldId
const parallelConfig = parallelData as any
newParallels[newId] = {
...parallelConfig,
id: newId,
// Update node references in parallel config
nodes: parallelConfig.nodes
? parallelConfig.nodes.map((nodeId: string) => blockIdMapping.get(nodeId) || nodeId)
: [],
}
}
updatedState.parallels = newParallels
}
}
// Update the workflow state with the new block IDs
// Update the workflow timestamp
await tx
.update(workflow)
.set({
state: updatedState,
updatedAt: now,
})
.where(eq(workflow.id, newWorkflowId))
+24 -4
View File
@@ -89,7 +89,14 @@ describe('Workflow By ID API Route', () => {
userId: 'user-123',
name: 'Test Workflow',
workspaceId: null,
state: { blocks: {}, edges: [] },
}
const mockNormalizedData = {
blocks: {},
edges: [],
loops: {},
parallels: {},
isFromNormalizedTables: true,
}
vi.doMock('@/lib/auth', () => ({
@@ -110,6 +117,10 @@ describe('Workflow By ID API Route', () => {
},
}))
vi.doMock('@/lib/workflows/db-helpers', () => ({
loadWorkflowFromNormalizedTables: vi.fn().mockResolvedValue(mockNormalizedData),
}))
const req = new NextRequest('http://localhost:3000/api/workflows/workflow-123')
const params = Promise.resolve({ id: 'workflow-123' })
@@ -127,7 +138,14 @@ describe('Workflow By ID API Route', () => {
userId: 'other-user',
name: 'Test Workflow',
workspaceId: 'workspace-456',
state: { blocks: {}, edges: [] },
}
const mockNormalizedData = {
blocks: {},
edges: [],
loops: {},
parallels: {},
isFromNormalizedTables: true,
}
vi.doMock('@/lib/auth', () => ({
@@ -148,6 +166,10 @@ describe('Workflow By ID API Route', () => {
},
}))
vi.doMock('@/lib/workflows/db-helpers', () => ({
loadWorkflowFromNormalizedTables: vi.fn().mockResolvedValue(mockNormalizedData),
}))
vi.doMock('@/lib/permissions/utils', () => ({
getUserEntityPermissions: vi.fn().mockResolvedValue('read'),
hasAdminPermission: vi.fn().mockResolvedValue(false),
@@ -170,7 +192,6 @@ describe('Workflow By ID API Route', () => {
userId: 'other-user',
name: 'Test Workflow',
workspaceId: 'workspace-456',
state: { blocks: {}, edges: [] },
}
vi.doMock('@/lib/auth', () => ({
@@ -213,7 +234,6 @@ describe('Workflow By ID API Route', () => {
userId: 'user-123',
name: 'Test Workflow',
workspaceId: null,
state: { blocks: {}, edges: [] },
}
const mockNormalizedData = {
+22 -31
View File
@@ -120,8 +120,6 @@ export async function GET(request: NextRequest, { params }: { params: Promise<{
logger.debug(`[${requestId}] Attempting to load workflow ${workflowId} from normalized tables`)
const normalizedData = await loadWorkflowFromNormalizedTables(workflowId)
const finalWorkflowData = { ...workflowData }
if (normalizedData) {
logger.debug(`[${requestId}] Found normalized data for workflow ${workflowId}:`, {
blocksCount: Object.keys(normalizedData.blocks).length,
@@ -131,38 +129,31 @@ export async function GET(request: NextRequest, { params }: { params: Promise<{
loops: normalizedData.loops,
})
// Use normalized table data - reconstruct complete state object
// First get any existing state properties, then override with normalized data
const existingState =
workflowData.state && typeof workflowData.state === 'object' ? workflowData.state : {}
finalWorkflowData.state = {
// Default values for expected properties
deploymentStatuses: {},
hasActiveWebhook: false,
// Preserve any existing state properties
...existingState,
// Override with normalized data (this takes precedence)
blocks: normalizedData.blocks,
edges: normalizedData.edges,
loops: normalizedData.loops,
parallels: normalizedData.parallels,
lastSaved: Date.now(),
isDeployed: workflowData.isDeployed || false,
deployedAt: workflowData.deployedAt,
// Construct response object with workflow data and state from normalized tables
const finalWorkflowData = {
...workflowData,
state: {
// Default values for expected properties
deploymentStatuses: {},
hasActiveWebhook: false,
// Data from normalized tables
blocks: normalizedData.blocks,
edges: normalizedData.edges,
loops: normalizedData.loops,
parallels: normalizedData.parallels,
lastSaved: Date.now(),
isDeployed: workflowData.isDeployed || false,
deployedAt: workflowData.deployedAt,
},
}
logger.info(`[${requestId}] Loaded workflow ${workflowId} from normalized tables`)
} else {
// Fallback to JSON blob
logger.info(
`[${requestId}] Using JSON blob for workflow ${workflowId} - no normalized data found`
)
const elapsed = Date.now() - startTime
logger.info(`[${requestId}] Successfully fetched workflow ${workflowId} in ${elapsed}ms`)
return NextResponse.json({ data: finalWorkflowData }, { status: 200 })
}
const elapsed = Date.now() - startTime
logger.info(`[${requestId}] Successfully fetched workflow ${workflowId} in ${elapsed}ms`)
return NextResponse.json({ data: finalWorkflowData }, { status: 200 })
return NextResponse.json({ error: 'Workflow has no normalized data' }, { status: 400 })
} catch (error: any) {
const elapsed = Date.now() - startTime
logger.error(`[${requestId}] Error fetching workflow ${workflowId} after ${elapsed}ms`, error)
@@ -220,7 +220,6 @@ export async function PUT(request: NextRequest, { params }: { params: Promise<{
.set({
lastSynced: new Date(),
updatedAt: new Date(),
state: saveResult.jsonBlob, // Also update JSON blob for backward compatibility
})
.where(eq(workflow.id, workflowId))
@@ -18,14 +18,12 @@ import { db } from '@/db'
import { workflowCheckpoints, workflow as workflowTable } from '@/db/schema'
import { generateLoopBlocks, generateParallelBlocks } from '@/stores/workflows/workflow/utils'
// Sim Agent API configuration
const SIM_AGENT_API_URL = env.SIM_AGENT_API_URL || SIM_AGENT_API_URL_DEFAULT
export const dynamic = 'force-dynamic'
const logger = createLogger('WorkflowYamlAPI')
// Request schema for YAML workflow operations
const YamlWorkflowRequestSchema = z.object({
yamlContent: z.string().min(1, 'YAML content is required'),
description: z.string().optional(),
@@ -647,14 +645,13 @@ export async function PUT(request: NextRequest, { params }: { params: Promise<{
.set({
lastSynced: new Date(),
updatedAt: new Date(),
state: saveResult.jsonBlob,
})
.where(eq(workflowTable.id, workflowId))
// Notify socket server for real-time collaboration (for copilot and editor)
if (source === 'copilot' || source === 'editor') {
try {
const socketUrl = process.env.SOCKET_URL || 'http://localhost:3002'
const socketUrl = env.SOCKET_SERVER_URL || 'http://localhost:3002'
await fetch(`${socketUrl}/api/copilot-workflow-edit`, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
-1
View File
@@ -151,7 +151,6 @@ export async function POST(req: NextRequest) {
folderId: folderId || null,
name,
description,
state: initialState,
color,
lastSynced: now,
createdAt: now,
@@ -85,14 +85,10 @@ export async function GET(request: NextRequest) {
edgesCount: normalizedData.edges.length,
})
// Use normalized table data - reconstruct complete state object
const existingState =
workflowData.state && typeof workflowData.state === 'object' ? workflowData.state : {}
// Use normalized table data - construct state from normalized tables
workflowState = {
deploymentStatuses: {},
hasActiveWebhook: false,
...existingState,
blocks: normalizedData.blocks,
edges: normalizedData.edges,
loops: normalizedData.loops,
@@ -116,33 +112,10 @@ export async function GET(request: NextRequest) {
logger.info(`[${requestId}] Loaded workflow ${workflowId} from normalized tables`)
} else {
// Fallback to JSON blob
logger.info(
`[${requestId}] Using JSON blob for workflow ${workflowId} - no normalized data found`
return NextResponse.json(
{ success: false, error: 'Workflow has no normalized data' },
{ status: 400 }
)
if (!workflowData.state || typeof workflowData.state !== 'object') {
return NextResponse.json(
{ success: false, error: 'Workflow has no valid state data' },
{ status: 400 }
)
}
workflowState = workflowData.state as any
// Extract subblock values from JSON blob state
if (workflowState.blocks) {
Object.entries(workflowState.blocks).forEach(([blockId, block]: [string, any]) => {
subBlockValues[blockId] = {}
if (block.subBlocks) {
Object.entries(block.subBlocks).forEach(([subBlockId, subBlock]: [string, any]) => {
if (subBlock && typeof subBlock === 'object' && 'value' in subBlock) {
subBlockValues[blockId][subBlockId] = subBlock.value
}
})
}
})
}
}
// Gather block registry and utilities for sim-agent
-59
View File
@@ -113,64 +113,6 @@ async function createWorkspace(userId: string, name: string) {
// Create initial workflow for the workspace with start block
const starterId = crypto.randomUUID()
const initialState = {
blocks: {
[starterId]: {
id: starterId,
type: 'starter',
name: 'Start',
position: { x: 100, y: 100 },
subBlocks: {
startWorkflow: {
id: 'startWorkflow',
type: 'dropdown',
value: 'manual',
},
webhookPath: {
id: 'webhookPath',
type: 'short-input',
value: '',
},
webhookSecret: {
id: 'webhookSecret',
type: 'short-input',
value: '',
},
scheduleType: {
id: 'scheduleType',
type: 'dropdown',
value: 'daily',
},
minutesInterval: {
id: 'minutesInterval',
type: 'short-input',
value: '',
},
minutesStartingAt: {
id: 'minutesStartingAt',
type: 'short-input',
value: '',
},
},
outputs: {
response: { type: { input: 'any' } },
},
enabled: true,
horizontalHandles: true,
isWide: false,
advancedMode: false,
height: 95,
},
},
edges: [],
subflows: {},
variables: {},
metadata: {
version: '1.0.0',
createdAt: now.toISOString(),
updatedAt: now.toISOString(),
},
}
// Create the workflow
await tx.insert(workflow).values({
@@ -180,7 +122,6 @@ async function createWorkspace(userId: string, name: string) {
folderId: null,
name: 'default-agent',
description: 'Your first workflow - start building here!',
state: initialState,
color: '#3972F6',
lastSynced: now,
createdAt: now,
@@ -0,0 +1 @@
ALTER TABLE "workflow" DROP COLUMN "state";
File diff suppressed because it is too large Load Diff
@@ -519,6 +519,13 @@
"when": 1755304368539,
"tag": "0074_abnormal_dreadnoughts",
"breakpoints": true
},
{
"idx": 75,
"version": "7",
"when": 1755319635487,
"tag": "0075_lush_moonstone",
"breakpoints": true
}
]
}
-3
View File
@@ -121,8 +121,6 @@ export const workflow = pgTable(
folderId: text('folder_id').references(() => workflowFolder.id, { onDelete: 'set null' }),
name: text('name').notNull(),
description: text('description'),
// DEPRECATED: Use normalized tables (workflow_blocks, workflow_edges, workflow_subflows) instead
state: json('state').notNull(),
color: text('color').notNull().default('#3972F6'),
lastSynced: timestamp('last_synced').notNull(),
createdAt: timestamp('created_at').notNull(),
@@ -130,7 +128,6 @@ export const workflow = pgTable(
isDeployed: boolean('is_deployed').notNull().default(false),
deployedState: json('deployed_state'),
deployedAt: timestamp('deployed_at'),
// When set, only this API key is authorized for execution
pinnedApiKey: text('pinned_api_key'),
collaborators: json('collaborators').notNull().default('[]'),
runCount: integer('run_count').notNull().default(0),
@@ -352,24 +352,8 @@ async function getUserWorkflow(workflowId: string): Promise<string> {
}
})
})
} else if (workflowRecord.state) {
// Fallback to JSON blob
const jsonState = workflowRecord.state as any
workflowState = {
blocks: jsonState.blocks || {},
edges: jsonState.edges || [],
loops: jsonState.loops || {},
parallels: jsonState.parallels || {},
}
// For JSON blob, subblock values are embedded in the block state
Object.entries((workflowState.blocks as any) || {}).forEach(([blockId, block]) => {
subBlockValues[blockId] = {}
Object.entries((block as any).subBlocks || {}).forEach(([subBlockId, subBlock]) => {
if ((subBlock as any).value !== undefined) {
subBlockValues[blockId][subBlockId] = (subBlock as any).value
}
})
})
} else {
throw new Error('Workflow has no normalized data')
}
if (!workflowState || !workflowState.blocks) {
-171
View File
@@ -1,171 +0,0 @@
#!/usr/bin/env bun
import { db } from '@/db'
import { user, workflow, workspace } from '@/db/schema'
const testWorkflowState = {
blocks: {
'start-block-123': {
id: 'start-block-123',
type: 'starter',
name: 'Start',
position: {
x: 100,
y: 100,
},
subBlocks: {
startWorkflow: {
id: 'startWorkflow',
type: 'dropdown',
value: 'manual',
},
},
outputs: {
response: {
input: 'any',
},
},
enabled: true,
horizontalHandles: true,
isWide: false,
advancedMode: false,
height: 90,
},
'loop-block-456': {
id: 'loop-block-456',
type: 'loop',
name: 'For Loop',
position: {
x: 400,
y: 100,
},
subBlocks: {},
outputs: {},
enabled: true,
horizontalHandles: true,
isWide: false,
advancedMode: false,
height: 0,
data: {
width: 400,
height: 200,
type: 'loopNode',
},
},
'function-block-789': {
id: 'function-block-789',
type: 'function',
name: 'Return X',
position: {
x: 50,
y: 50,
},
subBlocks: {
code: {
id: 'code',
type: 'code',
value: "return 'X'",
},
},
outputs: {
response: {
result: 'any',
stdout: 'string',
},
},
enabled: true,
horizontalHandles: true,
isWide: false,
advancedMode: false,
height: 144,
data: {
parentId: 'loop-block-456',
extent: 'parent',
},
},
},
edges: [
{
id: 'edge-start-to-loop',
source: 'start-block-123',
target: 'loop-block-456',
sourceHandle: 'source',
targetHandle: 'target',
},
{
id: 'edge-loop-to-function',
source: 'loop-block-456',
target: 'function-block-789',
sourceHandle: 'loop-start-source',
targetHandle: 'target',
},
],
loops: {
'loop-block-456': {
id: 'loop-block-456',
nodes: ['function-block-789'],
iterations: 3,
loopType: 'for',
forEachItems: '',
},
},
parallels: {},
lastSaved: Date.now(),
isDeployed: false,
}
async function insertTestWorkflow() {
try {
console.log('🔍 Finding first workspace and user...')
// Get the first workspace
const workspaces = await db.select().from(workspace).limit(1)
if (workspaces.length === 0) {
throw new Error('No workspaces found. Please create a workspace first.')
}
// Get the first user
const users = await db.select().from(user).limit(1)
if (users.length === 0) {
throw new Error('No users found. Please create a user first.')
}
const workspaceId = workspaces[0].id
const userId = users[0].id
console.log(`✅ Using workspace: ${workspaceId}`)
console.log(`✅ Using user: ${userId}`)
// Insert workflow with old JSON state format
const testWorkflowId = `test-migration-workflow-${Date.now()}`
const now = new Date()
await db.insert(workflow).values({
id: testWorkflowId,
name: 'Test Migration Workflow (Old JSON Format)',
workspaceId: workspaceId,
userId: userId,
state: testWorkflowState, // This is the old JSON format
lastSynced: now,
createdAt: now,
updatedAt: now,
isDeployed: false,
isPublished: false,
})
console.log(`✅ Inserted test workflow with old JSON format: ${testWorkflowId}`)
console.log(`🌐 Access it at: http://localhost:3000/w/${testWorkflowId}`)
console.log('')
console.log('📋 Test steps:')
console.log('1. Open the workflow in your browser')
console.log('2. Verify it renders correctly with all blocks and connections')
console.log('3. Try editing some subblock values')
console.log('4. Run the migration script')
console.log('5. Verify it still works after migration')
} catch (error) {
console.error('❌ Error inserting test workflow:', error)
process.exit(1)
}
}
insertTestWorkflow()
-306
View File
@@ -1,306 +0,0 @@
#!/usr/bin/env bun
import { readFileSync } from 'fs'
import { and, eq, inArray, isNotNull } from 'drizzle-orm'
import { nanoid } from 'nanoid'
import { db } from '@/db'
import { workflow, workflowBlocks, workflowEdges, workflowSubflows } from '@/db/schema'
interface WorkflowState {
blocks: Record<string, any>
edges: any[]
loops?: Record<string, any>
parallels?: Record<string, any>
lastSaved?: number
isDeployed?: boolean
}
async function migrateWorkflowStates(specificWorkflowIds?: string[] | null) {
try {
if (specificWorkflowIds) {
console.log(`🔍 Finding ${specificWorkflowIds.length} specific workflows...`)
} else {
console.log('🔍 Finding workflows with old JSON state format...')
}
// Build the where condition based on whether we have specific IDs
const whereCondition = specificWorkflowIds
? and(
isNotNull(workflow.state), // Has JSON state
inArray(workflow.id, specificWorkflowIds) // Only specific IDs
)
: and(
isNotNull(workflow.state) // Has JSON state
// We'll check for normalized data existence per workflow
)
// Find workflows that have state but no normalized table entries
const workflowsToMigrate = await db
.select({
id: workflow.id,
name: workflow.name,
state: workflow.state,
})
.from(workflow)
.where(whereCondition)
console.log(`📊 Found ${workflowsToMigrate.length} workflows with JSON state`)
if (specificWorkflowIds) {
const foundIds = workflowsToMigrate.map((w) => w.id)
const missingIds = specificWorkflowIds.filter((id) => !foundIds.includes(id))
if (missingIds.length > 0) {
console.log(`⚠️ Warning: ${missingIds.length} specified workflow IDs not found:`)
missingIds.forEach((id) => console.log(` - ${id}`))
}
console.log('')
}
let migratedCount = 0
let skippedCount = 0
let errorCount = 0
for (const wf of workflowsToMigrate) {
try {
// Check if this workflow already has normalized data
const existingBlocks = await db
.select({ id: workflowBlocks.id })
.from(workflowBlocks)
.where(eq(workflowBlocks.workflowId, wf.id))
.limit(1)
if (existingBlocks.length > 0) {
console.log(`⏭️ Skipping ${wf.name} (${wf.id}) - already has normalized data`)
skippedCount++
continue
}
console.log(`🔄 Migrating ${wf.name} (${wf.id})...`)
const state = wf.state as WorkflowState
if (!state || !state.blocks) {
console.log(`⚠️ Skipping ${wf.name} - invalid state format`)
skippedCount++
continue
}
// Clean up invalid blocks (those without an id field) before migration
const originalBlockCount = Object.keys(state.blocks).length
const validBlocks: Record<string, any> = {}
let removedBlockCount = 0
for (const [blockKey, block] of Object.entries(state.blocks)) {
if (block && typeof block === 'object' && block.id) {
// Valid block - has an id field
validBlocks[blockKey] = block
} else {
// Invalid block - missing id field
console.log(` 🗑️ Removing invalid block ${blockKey} (no id field)`)
removedBlockCount++
}
}
if (removedBlockCount > 0) {
console.log(
` 🧹 Cleaned up ${removedBlockCount} invalid blocks (${originalBlockCount} → ${Object.keys(validBlocks).length})`
)
state.blocks = validBlocks
}
await db.transaction(async (tx) => {
// Migrate blocks - generate new IDs and create mapping
const blocks = Object.values(state.blocks)
console.log(` 📦 Migrating ${blocks.length} blocks...`)
// Create mapping from old block IDs to new block IDs
const blockIdMapping: Record<string, string> = {}
for (const block of blocks) {
const newBlockId = nanoid()
blockIdMapping[block.id] = newBlockId
await tx.insert(workflowBlocks).values({
id: newBlockId,
workflowId: wf.id,
type: block.type,
name: block.name,
positionX: String(block.position?.x || 0),
positionY: String(block.position?.y || 0),
enabled: block.enabled ?? true,
horizontalHandles: block.horizontalHandles ?? true,
isWide: block.isWide ?? false,
advancedMode: block.advancedMode ?? false,
triggerMode: block.triggerMode ?? false,
height: String(block.height || 0),
subBlocks: block.subBlocks || {},
outputs: block.outputs || {},
data: block.data || {},
parentId: block.data?.parentId ? blockIdMapping[block.data.parentId] || null : null,
})
}
// Migrate edges - use new block IDs
const edges = state.edges || []
console.log(` 🔗 Migrating ${edges.length} edges...`)
for (const edge of edges) {
const newSourceId = blockIdMapping[edge.source]
const newTargetId = blockIdMapping[edge.target]
// Skip edges that reference blocks that don't exist in our mapping
if (!newSourceId || !newTargetId) {
console.log(` ⚠️ Skipping edge ${edge.id} - references missing blocks`)
continue
}
await tx.insert(workflowEdges).values({
id: nanoid(),
workflowId: wf.id,
sourceBlockId: newSourceId,
targetBlockId: newTargetId,
sourceHandle: edge.sourceHandle || null,
targetHandle: edge.targetHandle || null,
})
}
// Migrate loops - update node IDs to use new block IDs
const loops = state.loops || {}
const loopIds = Object.keys(loops)
console.log(` 🔄 Migrating ${loopIds.length} loops...`)
for (const loopId of loopIds) {
const loop = loops[loopId]
// Map old node IDs to new block IDs
const updatedNodes = (loop.nodes || [])
.map((nodeId: string) => blockIdMapping[nodeId])
.filter(Boolean)
await tx.insert(workflowSubflows).values({
id: nanoid(),
workflowId: wf.id,
type: 'loop',
config: {
id: loop.id,
nodes: updatedNodes,
iterationCount: loop.iterations || 5,
iterationType: loop.loopType || 'for',
collection: loop.forEachItems || '',
},
})
}
// Migrate parallels - update node IDs to use new block IDs
const parallels = state.parallels || {}
const parallelIds = Object.keys(parallels)
console.log(` ⚡ Migrating ${parallelIds.length} parallels...`)
for (const parallelId of parallelIds) {
const parallel = parallels[parallelId]
// Map old node IDs to new block IDs
const updatedNodes = (parallel.nodes || [])
.map((nodeId: string) => blockIdMapping[nodeId])
.filter(Boolean)
await tx.insert(workflowSubflows).values({
id: nanoid(),
workflowId: wf.id,
type: 'parallel',
config: {
id: parallel.id,
nodes: updatedNodes,
parallelCount: 2, // Default parallel count
collection: parallel.distribution || '',
},
})
}
})
console.log(`✅ Successfully migrated ${wf.name}`)
migratedCount++
} catch (error) {
console.error(`❌ Error migrating ${wf.name} (${wf.id}):`, error)
errorCount++
}
}
console.log('')
console.log('📊 Migration Summary:')
console.log(`✅ Migrated: ${migratedCount} workflows`)
console.log(`⏭️ Skipped: ${skippedCount} workflows`)
console.log(`❌ Errors: ${errorCount} workflows`)
console.log('')
if (migratedCount > 0) {
console.log('🎉 Migration completed successfully!')
console.log('')
console.log('📋 Next steps:')
console.log('1. Test the migrated workflows in your browser')
console.log('2. Verify all blocks, edges, and subflows work correctly')
console.log('3. Check that editing and collaboration still work')
console.log('4. Once confirmed, the workflow.state JSON field can be deprecated')
}
} catch (error) {
console.error('❌ Migration failed:', error)
process.exit(1)
}
}
// Add command line argument parsing
const args = process.argv.slice(2)
const dryRun = args.includes('--dry-run')
const showHelp = args.includes('--help') || args.includes('-h')
if (showHelp) {
console.log('🔄 Workflow State Migration Script')
console.log('')
console.log('Usage:')
console.log(' bun run scripts/migrate-workflow-states.ts [options]')
console.log('')
console.log('Options:')
console.log(' --dry-run Show what would be migrated without making changes')
console.log(' --file <path> Migrate only workflow IDs listed in file (comma-separated)')
console.log(' --help, -h Show this help message')
console.log('')
console.log('Examples:')
console.log(' bun run scripts/migrate-workflow-states.ts')
console.log(' bun run scripts/migrate-workflow-states.ts --dry-run')
console.log(' bun run scripts/migrate-workflow-states.ts --file workflow-ids.txt')
console.log(' bun run scripts/migrate-workflow-states.ts --dry-run --file workflow-ids.txt')
console.log('')
console.log('File format (workflow-ids.txt):')
console.log(' abc-123,def-456,ghi-789')
console.log('')
process.exit(0)
}
// Parse --file flag for workflow IDs
let specificWorkflowIds: string[] | null = null
const fileIndex = args.findIndex((arg) => arg === '--file')
if (fileIndex !== -1 && args[fileIndex + 1]) {
const filePath = args[fileIndex + 1]
try {
console.log(`📁 Reading workflow IDs from file: ${filePath}`)
const fileContent = readFileSync(filePath, 'utf-8')
specificWorkflowIds = fileContent
.split(',')
.map((id) => id.trim())
.filter((id) => id.length > 0)
console.log(`📋 Found ${specificWorkflowIds.length} workflow IDs in file`)
console.log('')
} catch (error) {
console.error(`❌ Error reading file ${filePath}:`, error)
process.exit(1)
}
}
if (dryRun) {
console.log('🔍 DRY RUN MODE - No changes will be made')
console.log('')
}
if (specificWorkflowIds) {
console.log('🎯 TARGETED MIGRATION - Only migrating specified workflow IDs')
console.log('')
}
migrateWorkflowStates(specificWorkflowIds)
@@ -125,15 +125,11 @@ export async function getWorkflowState(workflowId: string) {
if (normalizedData) {
// Use normalized data as source of truth
const existingState = workflowData[0].state || {}
const finalState = {
// Default values for expected properties
deploymentStatuses: {},
hasActiveWebhook: false,
// Preserve any existing state properties
...existingState,
// Override with normalized data (this takes precedence)
// Data from normalized tables
blocks: normalizedData.blocks,
edges: normalizedData.edges,
loops: normalizedData.loops,