fix subblock updates

This commit is contained in:
Vikhyath Mondreti
2025-07-11 20:11:54 -07:00
parent 5f9bfdde06
commit 91a4c6d588
3 changed files with 27 additions and 116 deletions
-8
View File
@@ -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,
+26 -68
View File
@@ -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)
+1 -40
View File
@@ -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<QueuedOperation, 'timestamp' | 'retryCount' | 'status'>) => 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<string, NodeJS.Timeout>()
const operationTimeouts = new Map<string, NodeJS.Timeout>()
// 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<OperationQueueState>((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<OperationQueueState>((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<OperationQueueState>((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<OperationQueueState>((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<OperationQueueState>((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<OperationQueueState>((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<OperationQueueState>((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<OperationQueueState>((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<OperationQueueState>((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<OperationQueueState>((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<OperationQueueState>((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<OperationQueueState>((set, get) =>
},
}))
// Hook wrapper for easier usage
export function useOperationQueue() {
const store = useOperationQueueStore()