Explorar o código

fix: relay every edit to the other editors, not only drags

Only a drag was sent live, so pressing Auto Layout rearranged the canvas
for the person who pressed it and nobody else, and the same hole silently
covered node configuration changes, adds, deletes, rewiring, renames and
undo - every one of them invisible until somebody saved.

Rather than sending from each of the forty-odd places that change the
graph, this watches the graph and sends what changed once it settles. It
is the choice the undo history already makes, and for the same reason:
the call site that forgets is the one whose edits never leave the
browser, and an auto-layout, a config change, a delete and an undo are
all just the graph looking different afterwards.

Sent as a difference rather than a whole workflow. A flow with fifty
configured nodes is hundreds of kilobytes, and moving one node would
otherwise post all of it to everybody several times a second.

An edit that arrives from somebody else is not echoed back, or two
editors talk in circles - and applying one marks this canvas as having
unsaved changes, because it now does.

Having the workflow open is the check the server makes on an incoming
edit. A lock per edit would buy nothing, since anyone who can see a
workflow can already save anything to it, and it would make an
auto-layout - which moves every node and holds no lock - impossible to
share. Identity is stamped by the server, never taken from the sender.

The loader and the live path now build nodes through one shared function.
Two copies would drift, and the symptom would be a node behaving
differently depending on whether it was loaded or received.

Verified with two browsers: adding nodes, Auto Layout, undo and a node
configuration change in one all appear in the other, and the changed
setting reads back correctly when the node is opened there.
fszontagh hai 1 mes
pai
achega
28fcfc76e6

+ 30 - 0
src/webserver/websocket_server.cpp

@@ -263,6 +263,8 @@ void WebSocketServer::processMessage(WebSocketClient& client, const nlohmann::js
         handleUnlock(client, message);
     } else if (type == "node_moved") {
         handleNodeMoved(client, message);
+    } else if (type == "graph_edit") {
+        handleGraphEdit(client, message);
     } else if (message_handler_) {
         message_handler_(client.id, message);
     }
@@ -489,6 +491,34 @@ void WebSocketServer::handleNodeMoved(WebSocketClient& client, const nlohmann::j
     });
 }
 
+
+void WebSocketServer::handleGraphEdit(WebSocketClient& client, const nlohmann::json& message) {
+    if (!client.authenticated) return;
+
+    const std::string workflow_id = message.value("workflowId", "");
+    if (workflow_id.empty()) return;
+
+    // Having the workflow open is the check that matters here. Anyone who can
+    // see it can already save whatever they like to it, so demanding a lock per
+    // edit would buy nothing and would make an auto-layout - which moves every
+    // node at once and holds no lock at all - impossible to share.
+    const std::string channel = "workflows." + workflow_id + ".edits";
+    {
+        std::shared_lock lock(clients_mutex_);
+        if (!client.subscriptions.contains(channel)) return;
+    }
+
+    nlohmann::json out = message.value("delta", nlohmann::json::object());
+    out["workflowId"] = workflow_id;
+    // Stamped here rather than trusted from the sender, so an edit cannot be
+    // attributed to somebody else.
+    out["userId"] = client.user_id;
+    out["username"] = client.username;
+    out["clientId"] = client.id;
+
+    broadcast(channel, out);
+}
+
 void WebSocketServer::releaseLocksOf(const std::string& client_id,
                                      std::vector<std::string>* touched_workflows) {
     std::lock_guard<std::mutex> lock(locks_mutex_);

+ 1 - 0
src/webserver/websocket_server.hpp

@@ -92,6 +92,7 @@ private:
     void handleLock(WebSocketClient& client, const nlohmann::json& message);
     void handleUnlock(WebSocketClient& client, const nlohmann::json& message);
     void handleNodeMoved(WebSocketClient& client, const nlohmann::json& message);
+    void handleGraphEdit(WebSocketClient& client, const nlohmann::json& message);
     bool matchesChannel(const std::string& subscription, const std::string& channel);
 
     WebSocketServerConfig config_;

+ 35 - 4
webui/src/hooks/useWorkflowLocks.ts

@@ -27,7 +27,24 @@ export interface NodeLock {
 // and sending them costs everybody bandwidth for no benefit.
 const MOVE_INTERVAL_MS = 50
 
-export function useWorkflowLocks(workflowId: string | undefined) {
+export interface RemoteEdit {
+  changed?: any[]
+  removed?: string[]
+  connections?: any[]
+  name?: string
+  username?: string
+  clientId?: string
+}
+
+export function useWorkflowLocks(
+  workflowId: string | undefined,
+  onRemoteEdit?: (edit: RemoteEdit) => void
+) {
+  // Held in a ref so a new callback identity on every render does not tear the
+  // subscription down and build it again.
+  const onRemoteEditRef = useRef(onRemoteEdit)
+  onRemoteEditRef.current = onRemoteEdit
+
   const [locks, setLocks] = useState<Record<string, NodeLock>>({})
   const [livePositions, setLivePositions] = useState<Record<string, { x: number; y: number }>>({})
   const [deniedNode, setDeniedNode] = useState<{ nodeId: string; heldBy: string } | null>(null)
@@ -43,6 +60,7 @@ export function useWorkflowLocks(workflowId: string | undefined) {
 
     const lockChannel = `workflows.${workflowId}.locks`
     const moveChannel = `workflows.${workflowId}.moves`
+    const editChannel = `workflows.${workflowId}.edits`
 
     const onLocks = (data: any) => {
       const next: Record<string, NodeLock> = {}
@@ -59,6 +77,12 @@ export function useWorkflowLocks(workflowId: string | undefined) {
       setLivePositions((current) => ({ ...current, [data.nodeId]: { x: p.x, y: p.y } }))
     }
 
+    const onEdit = (data: any) => {
+      // Your own edit is already on your own canvas.
+      if (data?.clientId && data.clientId === clientIdRef.current) return
+      onRemoteEditRef.current?.(data)
+    }
+
     const onDenied = (data: any) => {
       if (data?.workflowId !== workflowId) return
       setDeniedNode({ nodeId: data.nodeId, heldBy: data.heldBy || 'somebody else' })
@@ -70,9 +94,10 @@ export function useWorkflowLocks(workflowId: string | undefined) {
 
     wsClient.on(lockChannel, onLocks)
     wsClient.on(moveChannel, onMove)
+    wsClient.on(editChannel, onEdit)
     wsClient.on('lock_denied', onDenied)
     wsClient.on('auth_success', onIdentity)
-    wsClient.subscribe([lockChannel, moveChannel])
+    wsClient.subscribe([lockChannel, moveChannel, editChannel])
 
     return () => {
       // Letting go on the way out. The server would release these anyway when
@@ -83,9 +108,10 @@ export function useWorkflowLocks(workflowId: string | undefined) {
         wsClient.send({ type: 'unlock', workflowId, nodeId })
       }
       heldRef.current.clear()
-      wsClient.unsubscribe([lockChannel, moveChannel])
+      wsClient.unsubscribe([lockChannel, moveChannel, editChannel])
       wsClient.off(lockChannel, onLocks)
       wsClient.off(moveChannel, onMove)
+      wsClient.off(editChannel, onEdit)
       wsClient.off('lock_denied', onDenied)
       wsClient.off('auth_success', onIdentity)
       setLocks({})
@@ -112,6 +138,11 @@ export function useWorkflowLocks(workflowId: string | undefined) {
     })
   }, [workflowId])
 
+  const broadcastEdit = useCallback((delta: any) => {
+    if (!workflowId) return
+    wsClient.send({ type: 'graph_edit', workflowId, delta })
+  }, [workflowId])
+
   const reportMove = useCallback((nodeId: string, position: { x: number; y: number }) => {
     if (!workflowId) return
     const now = performance.now()
@@ -132,5 +163,5 @@ export function useWorkflowLocks(workflowId: string | undefined) {
 
   const dismissDenied = useCallback(() => setDeniedNode(null), [])
 
-  return { locks, lockedByOther, livePositions, claim, release, reportMove, deniedNode, dismissDenied }
+  return { locks, lockedByOther, livePositions, claim, release, reportMove, broadcastEdit, deniedNode, dismissDenied }
 }

+ 177 - 59
webui/src/pages/WorkflowEditorPage.tsx

@@ -32,7 +32,7 @@ import { AlertTriangle } from 'lucide-react'
 import { wsClient } from '../api/client'
 import { useWorkflowPresence } from '../hooks/useWorkflowPresence'
 import { useWorkflowLocks } from '../hooks/useWorkflowLocks'
-import { diffWorkflows, mergeWorkflows, summariseChanges, type Change } from '../utils/workflowDiff'
+import { diffWorkflows, graphDelta, mergeWorkflows, summariseChanges, type Change } from '../utils/workflowDiff'
 import type { AvailableField } from '../components/ConditionBuilder'
 import type { FieldInfo } from '../components/AvailableDataPanel'
 import { ExecutionListPanel } from '../components/workflow/ExecutionListPanel'
@@ -591,10 +591,6 @@ function WorkflowEditorInner() {
   const loadedVersionRef = useRef<number | null>(null)
   const { others: otherViewers, myOtherTabs, savedByOther, dismissSavedByOther } =
     useWorkflowPresence(id, loadedVersionRef)
-  const {
-    lockedByOther, livePositions, claim: claimNode, release: releaseNode,
-    reportMove, deniedNode, dismissDenied,
-  } = useWorkflowLocks(id)
   // The save handler needs the latest value without being rebuilt for it.
   const savedByOtherRef = useRef(savedByOther)
   savedByOtherRef.current = savedByOther
@@ -1339,6 +1335,61 @@ function WorkflowEditorInner() {
     return map
   }, [nodeDefs])
 
+  // What every editor watching this workflow should currently have. Edits are
+  // sent as the difference from this, and it moves forward both when this
+  // canvas broadcasts and when somebody else's edit is applied.
+  const sharedGraphRef = useRef<any>(null)
+  const applyingRemoteRef = useRef(false)
+
+  const applyRemoteEdit = useCallback((edit: any) => {
+    applyingRemoteRef.current = true
+
+    if (edit.changed?.length || edit.removed?.length) {
+      const removed = new Set<string>(edit.removed ?? [])
+      const changedById = new Map<string, any>((edit.changed ?? []).map((n: any) => [n.id, n]))
+      setNodes((current) => {
+        const kept = current
+          .filter((n) => !removed.has(n.id))
+          .map((n) => {
+            const incoming = changedById.get(n.id)
+            if (!incoming) return n
+            changedById.delete(n.id)
+            return {
+              ...n,
+              position: incoming.position ?? n.position,
+              data: {
+                ...n.data,
+                label: incoming.name ?? n.data.label,
+                config: incoming.config ?? n.data.config,
+                disabled: incoming.disabled === true,
+              },
+            }
+          })
+        // Whatever is left is new to this canvas. Built the same way as a node
+        // that arrives from a load, so it behaves like any other.
+        const added = [...changedById.values()].map((n: any) =>
+          makeNodeFromStored(n, nodeDefsMap, executeWorkflow, executeTrigger, executionState)
+        )
+        return [...kept, ...added]
+      })
+    }
+
+    if (edit.connections) {
+      setEdges(() => buildEdgesFromStored(edit.connections))
+    }
+
+    if (typeof edit.name === 'string') setWorkflowName(edit.name)
+
+    // Their change is now an unsaved change here too, and saying otherwise
+    // would let it be lost on the way out.
+    setHasChanges(true)
+  }, [setNodes, setEdges, nodeDefsMap, executionState])
+
+  const {
+    lockedByOther, livePositions, claim: claimNode, release: releaseNode,
+    reportMove, broadcastEdit, deniedNode, dismissDenied,
+  } = useWorkflowLocks(id, applyRemoteEdit)
+
   // Compute viewed nodes and edges when in view mode (from workflow snapshot)
   const { viewedNodes, viewedEdges } = useMemo(() => {
     if (!isViewingExecution || !pinnedExecution?.workflowSnapshot) {
@@ -1422,60 +1473,12 @@ function WorkflowEditorInner() {
   // Initialize nodes and edges from workflow
   useEffect(() => {
     if (workflow && nodeDefs.length > 0) {
-      const flowNodes: Node[] = workflow.nodes.map((node: any) => {
-        const nodeDef = nodeDefsMap[node.type]
-        const outputs: NodeOutput[] = nodeDef?.outputs?.length
-          ? nodeDef.outputs
-          : [{ name: 'main', displayName: 'Output', type: 'any' }]
-        return {
-          id: node.id,
-          type: 'workflowNode',
-          position: node.position || { x: 0, y: 0 },
-          data: {
-            label: node.name || node.type,
-            type: node.type,
-            config: node.config,
-            isTrigger: nodeDef?.isTrigger || false,
-            icon: nodeDef?.icon,
-            outputs,
-            dynamicOutputs: nodeDef?.dynamicOutputs,
-            inputs: nodeDef?.inputs,
-            nodeId: node.id,
-            disabled: node.disabled === true,
-            onExecute: () => executeWorkflow(),
-            onExecuteTrigger: executeTrigger,
-            executionState: executionState.nodeStates[node.id],
-          },
-          selected: false,
-        }
-      })
-
-      const flowEdges: Edge[] = workflow.connections.map((conn: any, idx: number) => {
-        let edgeColor = '#cbd5e1' // Lighter default for better visibility
-        if (conn.sourceOutput === 'true' || conn.sourceOutput === 'loop') {
-          edgeColor = '#22c55e'
-        } else if (conn.sourceOutput === 'false') {
-          edgeColor = '#ef4444'
-        } else if (conn.sourceOutput === 'done') {
-          edgeColor = '#3b82f6'
-        }
-
-        return {
-          id: `e${idx}`,
-          source: conn.sourceNodeId,
-          sourceHandle: conn.sourceOutput,
-          target: conn.targetNodeId,
-          targetHandle: conn.targetInput,
-          type: 'smart',
-          style: { strokeWidth: 2, stroke: edgeColor },
-          animated: executionState.status === 'running',
-          label: conn.sourceOutput !== 'main' ? conn.sourceOutput : undefined,
-          labelStyle: { fill: edgeColor, fontWeight: 500, fontSize: 10 },
-          labelBgStyle: { fill: 'var(--color-bg)', fillOpacity: 0.9 },
-          data: { handleType: conn.sourceOutput || 'main' },
-          zIndex: 1,
-        }
-      })
+      const flowNodes: Node[] = workflow.nodes.map((node: any) =>
+        makeNodeFromStored(node, nodeDefsMap, () => executeWorkflow(), executeTrigger, executionState)
+      )
+      const flowEdges: Edge[] = buildEdgesFromStored(
+        workflow.connections, executionState.status === 'running'
+      )
 
       setNodes(flowNodes)
       setEdges(flowEdges)
@@ -1487,9 +1490,53 @@ function WorkflowEditorInner() {
       setHasChanges(loadedIsUnsavedRef.current)
       loadedIsUnsavedRef.current = false
       setGraphEpoch((n) => n + 1)
+      // Everybody who loads this record already agrees on it, so it is the
+      // point future edits are measured from.
+      sharedGraphRef.current = {
+        name: workflow.name,
+        nodes: workflow.nodes || [],
+        connections: workflow.connections || [],
+        settings: workflow.settings || {},
+      }
     }
   }, [workflow, nodeDefs, nodeDefsMap, setNodes, setEdges])
 
+  // Send what changed here to everybody else watching.
+  //
+  // Watching the graph rather than calling this from each place that changes it
+  // is the same choice the undo history makes, and for the same reason: there
+  // are around forty such places, and the one that forgets is the one whose
+  // edits silently never leave the browser. An auto-layout, a config change, a
+  // delete and an undo are all just the graph looking different afterwards.
+  useEffect(() => {
+    if (isViewingExecution || !id) return
+    // A drag has its own stream, several times a second; this would only send
+    // the same thing again more slowly.
+    if (nodes.some((n) => (n as any).dragging)) return
+
+    const timer = setTimeout(() => {
+      const now = currentGraphRef.current()
+
+      // An edit that arrived from somebody else must not be echoed back to
+      // them, or two editors will talk in circles.
+      if (applyingRemoteRef.current) {
+        applyingRemoteRef.current = false
+        sharedGraphRef.current = now
+        return
+      }
+      if (!sharedGraphRef.current) {
+        sharedGraphRef.current = now
+        return
+      }
+      const delta = graphDelta(sharedGraphRef.current, now)
+      if (delta) {
+        broadcastEdit(delta)
+        sharedGraphRef.current = now
+      }
+    }, 200)
+    return () => clearTimeout(timer)
+  }, [nodes, edges, workflowName, id, isViewingExecution, broadcastEdit])
+
   // Put the lock state and anybody else's in-flight drag onto the canvas.
   // Separate from loading, because these arrive while somebody is working and
   // must not disturb anything else about the node.
@@ -3305,6 +3352,77 @@ function WorkflowEditorInner() {
   )
 }
 
+
+/**
+ * A stored node or connection turned into something ReactFlow can draw.
+ *
+ * Pulled out of the loader because a live edit arriving from another editor has
+ * to build a node exactly the same way. Two copies of this would drift, and the
+ * symptom would be a node that behaves subtly differently depending on whether
+ * it was loaded or received.
+ */
+function makeNodeFromStored(
+  node: any,
+  nodeDefsMap: Record<string, NodeDefinition>,
+  onExecute: () => void,
+  onExecuteTrigger: (nodeId: string) => void,
+  executionState: { nodeStates: Record<string, any> }
+): Node {
+  const nodeDef = nodeDefsMap[node.type]
+  const outputs: NodeOutput[] = nodeDef?.outputs?.length
+    ? nodeDef.outputs
+    : [{ name: 'main', displayName: 'Output', type: 'any' }]
+  return {
+    id: node.id,
+    type: 'workflowNode',
+    position: node.position || { x: 0, y: 0 },
+    data: {
+      label: node.name || node.type,
+      type: node.type,
+      config: node.config,
+      isTrigger: nodeDef?.isTrigger || false,
+      icon: nodeDef?.icon,
+      outputs,
+      dynamicOutputs: nodeDef?.dynamicOutputs,
+      inputs: nodeDef?.inputs,
+      nodeId: node.id,
+      disabled: node.disabled === true,
+      onExecute,
+      onExecuteTrigger,
+      executionState: executionState.nodeStates[node.id],
+    },
+    selected: false,
+  }
+}
+
+function buildEdgesFromStored(connections: any[], animated = false): Edge[] {
+  return (connections || []).map((conn: any, idx: number) => {
+    let edgeColor = '#cbd5e1'
+    if (conn.sourceOutput === 'true' || conn.sourceOutput === 'loop') {
+      edgeColor = '#22c55e'
+    } else if (conn.sourceOutput === 'false') {
+      edgeColor = '#ef4444'
+    } else if (conn.sourceOutput === 'done') {
+      edgeColor = '#3b82f6'
+    }
+    return {
+      id: `e${idx}`,
+      source: conn.sourceNodeId,
+      sourceHandle: conn.sourceOutput,
+      target: conn.targetNodeId,
+      targetHandle: conn.targetInput,
+      type: 'smart',
+      style: { strokeWidth: 2, stroke: edgeColor },
+      animated,
+      label: conn.sourceOutput !== 'main' ? conn.sourceOutput : undefined,
+      labelStyle: { fill: edgeColor, fontWeight: 500, fontSize: 10 },
+      labelBgStyle: { fill: 'var(--color-bg)', fillOpacity: 0.9 },
+      data: { handleType: conn.sourceOutput || 'main' },
+      zIndex: 1,
+    }
+  })
+}
+
 export default function WorkflowEditorPage() {
   return (
     <ReactFlowProvider>

+ 43 - 0
webui/src/utils/workflowDiff.ts

@@ -309,3 +309,46 @@ export function mergeWorkflows(
 
   return { graph, conflicts, differsFromStored: diffWorkflows(theirs, graph).length > 0 }
 }
+
+/**
+ * The smallest description of how a graph got from one state to another, for
+ * sending to the other editors.
+ *
+ * Sending the whole workflow on every change would be simpler and is what a
+ * first attempt usually does, but a workflow with fifty configured nodes is
+ * hundreds of kilobytes, and moving one node would post all of it to everybody
+ * several times a second.
+ */
+export interface GraphDelta {
+  changed: any[]
+  removed: string[]
+  /** Sent whole when it differs at all - the wiring is small and fiddly to patch. */
+  connections?: any[]
+  name?: string
+}
+
+export function graphDelta(base: WorkflowGraph, next: WorkflowGraph): GraphDelta | null {
+  const baseById = new Map((base.nodes || []).map((n: any) => [n.id, n]))
+  const nextById = new Map((next.nodes || []).map((n: any) => [n.id, n]))
+
+  const changed: any[] = []
+  for (const [id, node] of nextById) {
+    if (!sameNode(baseById.get(id), node)) changed.push(node)
+  }
+  const removed = [...baseById.keys()].filter((id) => !nextById.has(id))
+
+  const wiringChanged =
+    JSON.stringify((base.connections || []).map(connectionKey).sort()) !==
+    JSON.stringify((next.connections || []).map(connectionKey).sort())
+
+  const renamed = base.name !== next.name ? next.name : undefined
+
+  if (!changed.length && !removed.length && !wiringChanged && renamed === undefined) return null
+
+  return {
+    changed,
+    removed,
+    ...(wiringChanged ? { connections: next.connections || [] } : {}),
+    ...(renamed !== undefined ? { name: renamed } : {}),
+  }
+}