From 91a4c6d588c59daa3668eb456631d72ae815d3d0 Mon Sep 17 00:00:00 2001 From: Vikhyath Mondreti Date: Fri, 11 Jul 2025 20:11:54 -0700 Subject: [PATCH] fix subblock updates --- apps/sim/contexts/socket-context.tsx | 8 -- apps/sim/hooks/use-collaborative-workflow.ts | 94 ++++++-------------- apps/sim/stores/operation-queue/store.ts | 41 +-------- 3 files changed, 27 insertions(+), 116 deletions(-) diff --git a/apps/sim/contexts/socket-context.tsx b/apps/sim/contexts/socket-context.tsx index faed9dd064..11987bb0cb 100644 --- a/apps/sim/contexts/socket-context.tsx +++ b/apps/sim/contexts/socket-context.tsx @@ -473,14 +473,6 @@ export function SocketProvider({ children, user }: SocketProviderProps) { // Emit workflow operations (blocks, edges, subflows) const emitWorkflowOperation = useCallback( (operation: string, target: string, payload: any, operationId?: string) => { - console.log('🚀 Attempting to emit operation', { - hasSocket: !!socket, - currentWorkflowId, - operationId, - operation, - target, - }) - if (!socket || !currentWorkflowId) { console.log('❌ Cannot emit - missing requirements', { hasSocket: !!socket, diff --git a/apps/sim/hooks/use-collaborative-workflow.ts b/apps/sim/hooks/use-collaborative-workflow.ts index 5363d33bfa..79ce111df0 100644 --- a/apps/sim/hooks/use-collaborative-workflow.ts +++ b/apps/sim/hooks/use-collaborative-workflow.ts @@ -98,12 +98,10 @@ export function useCollaborativeWorkflow() { handleSocketReconnection, ]) - // Handle incoming workflow operations from other users useEffect(() => { const handleWorkflowOperation = (data: any) => { const { operation, target, payload, userId } = data - // Don't apply our own operations if (isApplyingRemoteChange.current) return logger.info(`Received ${operation} on ${target} from user ${userId}`) @@ -115,8 +113,6 @@ export function useCollaborativeWorkflow() { if (target === 'block') { switch (operation) { case 'add': - // Use normal addBlock - the collaborative system now sends complete data - // and the validation schema preserves outputs and subBlocks workflowStore.addBlock( payload.id, payload.type, @@ -126,17 +122,13 @@ export function useCollaborativeWorkflow() { payload.parentId, payload.extent ) - // Handle auto-connect edge if present if (payload.autoConnectEdge) { workflowStore.addEdge(payload.autoConnectEdge) } break case 'update-position': { - // Apply position update only if it's newer than the last applied timestamp - // This prevents jagged movement from out-of-order position updates const blockId = payload.id - // Server should always provide timestamp - if missing, skip ordering check if (!data.timestamp) { logger.warn('Position update missing timestamp, applying without ordering check', { blockId, @@ -180,12 +172,9 @@ export function useCollaborativeWorkflow() { workflowStore.setBlockWide(payload.id, payload.isWide) break case 'update-advanced-mode': - // Note: toggleBlockAdvancedMode doesn't take a parameter, it just toggles - // For now, we'll use the existing toggle method workflowStore.toggleBlockAdvancedMode(payload.id) break case 'toggle-handles': { - // Apply the handles toggle - we need to set the specific value to ensure consistency const currentBlock = workflowStore.blocks[payload.id] if (currentBlock && currentBlock.horizontalHandles !== payload.horizontalHandles) { workflowStore.toggleBlockHandles(payload.id) @@ -193,7 +182,6 @@ export function useCollaborativeWorkflow() { break } case 'duplicate': - // Apply the duplicate operation by adding the new block workflowStore.addBlock( payload.id, payload.type, @@ -374,7 +362,6 @@ export function useCollaborativeWorkflow() { const { operationId, error, retryable } = data logger.warn('Operation failed', { operationId, error, retryable }) - // Create a retry function that re-emits the operation using the correct channel const retryFunction = (operation: any) => { const { operation: op, target, payload } = operation.operation @@ -421,7 +408,6 @@ export function useCollaborativeWorkflow() { queue, ]) - // Helper function to execute queued operations const executeQueuedOperation = useCallback( (operation: string, target: string, payload: any, localAction: () => void) => { console.log('🎯 executeQueuedOperation called', { @@ -430,16 +416,12 @@ export function useCollaborativeWorkflow() { isApplyingRemoteChange: isApplyingRemoteChange.current, }) - // Skip if applying remote changes if (isApplyingRemoteChange.current) { - console.log('❌ Skipping - applying remote change') return } - // Generate operation ID const operationId = crypto.randomUUID() - // Add to queue addToQueue({ id: operationId, operation: { @@ -451,34 +433,24 @@ export function useCollaborativeWorkflow() { userId: session?.user?.id || 'unknown', }) - // Apply locally localAction() - // Emit to server with operation ID emitWorkflowOperation(operation, target, payload, operationId) }, [addToQueue, emitWorkflowOperation, session?.user?.id] ) - // Special helper for debounced operations (position updates) - // These are high-frequency, low-importance operations that don't need queue tracking const executeQueuedDebouncedOperation = useCallback( (operation: string, target: string, payload: any, localAction: () => void) => { - // Skip if applying remote changes if (isApplyingRemoteChange.current) return - // Apply locally first (immediate UI feedback) localAction() - // For debounced operations, don't use queue tracking - // The debouncing in socket context handles reliability - // No operation ID needed since we're not tracking these emitWorkflowOperation(operation, target, payload) }, [emitWorkflowOperation] ) - // Collaborative workflow operations const collaborativeAddBlock = useCallback( ( id: string, @@ -490,7 +462,6 @@ export function useCollaborativeWorkflow() { extent?: 'parent', autoConnectEdge?: Edge ) => { - // Create complete block data upfront using the same logic as the store const blockConfig = getBlock(type) // Handle loop/parallel blocks that don't use BlockConfig @@ -513,13 +484,11 @@ export function useCollaborativeWorkflow() { autoConnectEdge, // Include edge data for atomic operation } - // Apply locally first workflowStore.addBlock(id, type, name, position, data, parentId, extent) if (autoConnectEdge) { workflowStore.addEdge(autoConnectEdge) } - // Then broadcast to other clients with complete block data if (!isApplyingRemoteChange.current) { emitWorkflowOperation('add', 'block', completeBlockData) } @@ -545,7 +514,6 @@ export function useCollaborativeWorkflow() { }) } - // Generate outputs using the same logic as the store const outputs = resolveOutputType(blockConfig.outputs) const completeBlockData = { @@ -669,11 +637,9 @@ export function useCollaborativeWorkflow() { const collaborativeToggleBlockAdvancedMode = useCallback( (id: string) => { - // Get the current state before toggling const currentBlock = workflowStore.blocks[id] if (!currentBlock) return - // Calculate the new advancedMode value const newAdvancedMode = !currentBlock.advancedMode executeQueuedOperation( @@ -688,11 +654,9 @@ export function useCollaborativeWorkflow() { const collaborativeToggleBlockHandles = useCallback( (id: string) => { - // Get the current state before toggling const currentBlock = workflowStore.blocks[id] if (!currentBlock) return - // Calculate the new horizontalHandles value const newHorizontalHandles = !currentBlock.horizontalHandles executeQueuedOperation( @@ -717,7 +681,6 @@ export function useCollaborativeWorkflow() { y: sourceBlock.position.y + 20, } - // Generate new name with numbering const match = sourceBlock.name.match(/(.*?)(\d+)?$/) const newName = match?.[2] ? `${match[1]}${Number.parseInt(match[2]) + 1}` @@ -741,7 +704,6 @@ export function useCollaborativeWorkflow() { height: sourceBlock.height || 0, } - // Apply locally first using addBlock to ensure consistent IDs workflowStore.addBlock( newId, sourceBlock.type, @@ -752,7 +714,6 @@ export function useCollaborativeWorkflow() { sourceBlock.data?.extent ) - // Copy subblock values to the new block const activeWorkflowId = useWorkflowRegistry.getState().activeWorkflowId if (activeWorkflowId) { const subBlockValues = @@ -769,7 +730,6 @@ export function useCollaborativeWorkflow() { } executeQueuedOperation('duplicate', 'block', duplicatedBlockData, () => { - // Apply locally - add the duplicated block workflowStore.addBlock( newId, sourceBlock.type, @@ -778,10 +738,8 @@ export function useCollaborativeWorkflow() { sourceBlock.data ? JSON.parse(JSON.stringify(sourceBlock.data)) : {} ) - // Copy subblock values to the new block const subBlockValues = subBlockStore.workflowValues[activeWorkflowId || '']?.[sourceId] if (subBlockValues && activeWorkflowId) { - // Copy each subblock value individually Object.entries(subBlockValues).forEach(([subblockId, value]) => { subBlockStore.setValue(newId, subblockId, value) }) @@ -809,10 +767,8 @@ export function useCollaborativeWorkflow() { const collaborativeSetSubblockValue = useCallback( (blockId: string, subblockId: string, value: any) => { - // Skip if applying remote changes if (isApplyingRemoteChange.current) return - // Check workflow state if (!currentWorkflowId || activeWorkflowId !== currentWorkflowId) { logger.debug('Skipping subblock update - not in active workflow', { currentWorkflowId, @@ -823,17 +779,37 @@ export function useCollaborativeWorkflow() { return } - // Apply locally first + // Generate operation ID for queue tracking + const operationId = crypto.randomUUID() + + // Add to queue for retry mechanism + addToQueue({ + id: operationId, + operation: { + operation: 'subblock-update', + target: 'subblock', + payload: { blockId, subblockId, value }, + }, + workflowId: activeWorkflowId || '', + userId: session?.user?.id || 'unknown', + }) + + // Apply locally first (immediate UI feedback) subBlockStore.setValue(blockId, subblockId, value) - // Emit to server (subblock updates have their own handler with built-in retry) - // No need for operation queue - the subblock handler already has confirmation/failure logic - emitSubblockUpdate(blockId, subblockId, value) + // Emit to server with operation ID for tracking + emitSubblockUpdate(blockId, subblockId, value, operationId) }, - [subBlockStore, emitSubblockUpdate, currentWorkflowId, activeWorkflowId] + [ + subBlockStore, + emitSubblockUpdate, + currentWorkflowId, + activeWorkflowId, + addToQueue, + session?.user?.id, + ] ) - // Collaborative loop/parallel configuration updates const collaborativeUpdateLoopCount = useCallback( (loopId: string, count: number) => { // Get current state BEFORE making changes @@ -866,16 +842,13 @@ export function useCollaborativeWorkflow() { const collaborativeUpdateLoopType = useCallback( (loopId: string, loopType: 'for' | 'forEach') => { - // Get current state BEFORE making changes const currentBlock = workflowStore.blocks[loopId] if (!currentBlock || currentBlock.type !== 'loop') return - // Find child nodes before state changes const childNodes = Object.values(workflowStore.blocks) .filter((b) => b.data?.parentId === loopId) .map((b) => b.id) - // Get current values to preserve them const currentIterations = currentBlock.data?.count || 5 const currentCollection = currentBlock.data?.collection || '' @@ -896,16 +869,13 @@ export function useCollaborativeWorkflow() { const collaborativeUpdateLoopCollection = useCallback( (loopId: string, collection: string) => { - // Get current state BEFORE making changes const currentBlock = workflowStore.blocks[loopId] if (!currentBlock || currentBlock.type !== 'loop') return - // Find child nodes before state changes const childNodes = Object.values(workflowStore.blocks) .filter((b) => b.data?.parentId === loopId) .map((b) => b.id) - // Get current values to preserve them const currentIterations = currentBlock.data?.count || 5 const currentLoopType = currentBlock.data?.loopType || 'for' @@ -926,16 +896,13 @@ export function useCollaborativeWorkflow() { const collaborativeUpdateParallelCount = useCallback( (parallelId: string, count: number) => { - // Get current state BEFORE making changes const currentBlock = workflowStore.blocks[parallelId] if (!currentBlock || currentBlock.type !== 'parallel') return - // Find child nodes before state changes const childNodes = Object.values(workflowStore.blocks) .filter((b) => b.data?.parentId === parallelId) .map((b) => b.id) - // Get current values to preserve them const currentDistribution = currentBlock.data?.collection || '' const currentParallelType = currentBlock.data?.parallelType || 'collection' @@ -959,16 +926,13 @@ export function useCollaborativeWorkflow() { const collaborativeUpdateParallelCollection = useCallback( (parallelId: string, collection: string) => { - // Get current state BEFORE making changes const currentBlock = workflowStore.blocks[parallelId] if (!currentBlock || currentBlock.type !== 'parallel') return - // Find child nodes before state changes const childNodes = Object.values(workflowStore.blocks) .filter((b) => b.data?.parentId === parallelId) .map((b) => b.id) - // Get current values to preserve them const currentCount = currentBlock.data?.count || 5 const currentParallelType = currentBlock.data?.parallelType || 'collection' @@ -992,23 +956,18 @@ export function useCollaborativeWorkflow() { const collaborativeUpdateParallelType = useCallback( (parallelId: string, parallelType: 'count' | 'collection') => { - // Get current state BEFORE making changes const currentBlock = workflowStore.blocks[parallelId] if (!currentBlock || currentBlock.type !== 'parallel') return - // Find child nodes before state changes const childNodes = Object.values(workflowStore.blocks) .filter((b) => b.data?.parentId === parallelId) .map((b) => b.id) - // Calculate new values based on type change let newCount = currentBlock.data?.count || 5 let newDistribution = currentBlock.data?.collection || '' - // Reset values based on type (same logic as the UI) if (parallelType === 'count') { newDistribution = '' - // Keep existing count } else { newCount = 1 newDistribution = newDistribution || '' @@ -1027,7 +986,6 @@ export function useCollaborativeWorkflow() { 'subflow', { id: parallelId, type: 'parallel', config }, () => { - // Apply all changes locally workflowStore.updateParallelType(parallelId, parallelType) workflowStore.updateParallelCount(parallelId, newCount) workflowStore.updateParallelCollection(parallelId, newDistribution) diff --git a/apps/sim/stores/operation-queue/store.ts b/apps/sim/stores/operation-queue/store.ts index a6014ac1a5..a49d7951c4 100644 --- a/apps/sim/stores/operation-queue/store.ts +++ b/apps/sim/stores/operation-queue/store.ts @@ -3,7 +3,6 @@ import { createLogger } from '@/lib/logs/console-logger' const logger = createLogger('OperationQueue') -// Operation queue types export interface QueuedOperation { id: string operation: { @@ -11,7 +10,7 @@ export interface QueuedOperation { target: string payload: any } - workflowId: string // Track which workflow this operation belongs to + workflowId: string timestamp: number retryCount: number status: 'pending' | 'confirmed' | 'failed' @@ -23,7 +22,6 @@ interface OperationQueueState { isProcessing: boolean hasOperationError: boolean - // Actions addToQueue: (operation: Omit) => void confirmOperation: (operationId: string) => void failOperation: (operationId: string, emitFunction: (operation: QueuedOperation) => void) => void @@ -33,11 +31,9 @@ interface OperationQueueState { clearError: () => void } -// Global timeout maps (outside of Zustand store to avoid serialization issues) const retryTimeouts = new Map() const operationTimeouts = new Map() -// Global registry for emit functions and current workflow (set by collaborative workflow hook) let emitWorkflowOperation: | ((operation: string, target: string, payload: any, operationId?: string) => void) | null = null @@ -64,7 +60,6 @@ export const useOperationQueueStore = create((set, get) => addToQueue: (operation) => { const state = get() - // Check if operation already exists in queue const existingOp = state.operations.find((op) => op.id === operation.id) if (existingOp) { console.log('⚠️ Operation already in queue, skipping duplicate', { operationId: operation.id }) @@ -83,15 +78,12 @@ export const useOperationQueueStore = create((set, get) => operation: queuedOp.operation, }) - // Start 5-second timeout to detect unresponsive server - // This will trigger retry mechanism if server doesn't respond at all const timeoutId = setTimeout(() => { logger.warn('Operation timeout - no server response after 5 seconds', { operationId: queuedOp.id, }) operationTimeouts.delete(queuedOp.id) - // Handle timeout directly in store instead of emitting events get().handleOperationTimeout(queuedOp.id) }, 5000) @@ -106,14 +98,12 @@ export const useOperationQueueStore = create((set, get) => const state = get() const newOperations = state.operations.filter((op) => op.id !== operationId) - // Clear any retry timeout for this operation const retryTimeout = retryTimeouts.get(operationId) if (retryTimeout) { clearTimeout(retryTimeout) retryTimeouts.delete(operationId) } - // Clear any operation timeout for this operation const operationTimeout = operationTimeouts.get(operationId) if (operationTimeout) { clearTimeout(operationTimeout) @@ -136,7 +126,6 @@ export const useOperationQueueStore = create((set, get) => return } - // Clear any existing operation timeout since we're handling the failure const operationTimeout = operationTimeouts.get(operationId) if (operationTimeout) { clearTimeout(operationTimeout) @@ -144,7 +133,6 @@ export const useOperationQueueStore = create((set, get) => } if (operation.retryCount < 3) { - // Retry the operation with exponential backoff const newRetryCount = operation.retryCount + 1 const delay = 2 ** newRetryCount * 1000 // 2s, 4s, 8s @@ -154,7 +142,6 @@ export const useOperationQueueStore = create((set, get) => }) const timeout = setTimeout(() => { - // Check if we're still in the same workflow before retrying if (operation.workflowId !== currentWorkflowId) { logger.warn('Cancelling retry - workflow changed', { operationId, @@ -162,41 +149,24 @@ export const useOperationQueueStore = create((set, get) => currentWorkflow: currentWorkflowId, }) retryTimeouts.delete(operationId) - // Remove operation from queue since it's no longer relevant set((state) => ({ operations: state.operations.filter((op) => op.id !== operationId), })) return } - // Re-emit the operation emitFunction(operation) retryTimeouts.delete(operationId) - - // Start a new operation timeout for the retry - const newTimeoutId = setTimeout(() => { - logger.warn('Retry operation timeout - no server response after 5 seconds', { - operationId, - }) - operationTimeouts.delete(operationId) - - // Trigger another retry attempt - get().handleOperationTimeout(operationId) - }, 5000) - - operationTimeouts.set(operationId, newTimeoutId) }, delay) retryTimeouts.set(operationId, timeout) - // Update retry count set((state) => ({ operations: state.operations.map((op) => op.id === operationId ? { ...op, retryCount: newRetryCount } : op ), })) } else { - // Max retries exceeded - trigger offline mode logger.error('Operation failed after max retries, triggering offline mode', { operationId }) get().triggerOfflineMode() } @@ -214,24 +184,20 @@ export const useOperationQueueStore = create((set, get) => operationId, }) - // Create a retry function that re-emits the operation using the correct channel const retryFunction = (operation: any) => { const { operation: op, target, payload } = operation.operation if (op === 'subblock-update' && target === 'subblock') { - // Use subblock-update channel for subblock operations if (emitSubblockUpdate) { emitSubblockUpdate(payload.blockId, payload.subblockId, payload.value, operation.id) } } else { - // Use workflow-operation channel for block/edge/subflow operations if (emitWorkflowOperation) { emitWorkflowOperation(op, target, payload, operation.id) } } } - // Treat timeout as a failure to trigger retry mechanism get().failOperation(operationId, retryFunction) }, @@ -244,7 +210,6 @@ export const useOperationQueueStore = create((set, get) => operationTimeouts.forEach((timeout) => clearTimeout(timeout)) operationTimeouts.clear() - // Keep operations in queue but reset their retry counts and start fresh timeouts const state = get() const resetOperations = state.operations.map((op) => ({ ...op, @@ -258,7 +223,6 @@ export const useOperationQueueStore = create((set, get) => hasOperationError: false, }) - // Start new timeouts for all operations (they'll retry when socket is ready) resetOperations.forEach((operation) => { const timeoutId = setTimeout(() => { logger.warn('Operation timeout after reconnection - no server response after 5 seconds', { @@ -275,13 +239,11 @@ export const useOperationQueueStore = create((set, get) => triggerOfflineMode: () => { logger.error('Operation failed after retries - triggering offline mode') - // Clear all timeouts and queue retryTimeouts.forEach((timeout) => clearTimeout(timeout)) retryTimeouts.clear() operationTimeouts.forEach((timeout) => clearTimeout(timeout)) operationTimeouts.clear() - // Clear queue and trigger error state set({ operations: [], isProcessing: false, @@ -294,7 +256,6 @@ export const useOperationQueueStore = create((set, get) => }, })) -// Hook wrapper for easier usage export function useOperationQueue() { const store = useOperationQueueStore()