Explorar o código

feat: show who else has a workflow open, and refuse a save that would overwrite them

Two people editing the same workflow used to be entirely invisible to each
other, and whoever pressed save second silently threw the other's work
away. This does not make editing collaborative - it makes it honest.

Presence: the editor subscribes to workflows.<id>.presence and the server
publishes the roster of who is on it, derived from the subscriptions
rather than kept as a separate list, so a closed tab or a dropped
connection removes somebody without anything having to remember to say
so. One person with two tabs is still one person. The header shows the
others; a single-user install never sees it.

Conflicts: the editor sends the record version it loaded, and the server
refuses the save with a 409 if the stored workflow has moved on, handing
back the current record. The work stays on the canvas and the user
chooses what happens to it. A save without a version still goes through
unchanged, so nothing else that writes workflows had to change. The check
is a read then a write rather than one atomic compare-and-set, because
this endpoint writes a patch and the atomic form is replace-only - noted
in the code.

Three things had to be fixed for any of it to work, each of which was
silently wrong on its own:

- The WebSocket client dropped any subscribe issued before the socket
  finished connecting, which on a cold page load is all of them. It now
  remembers what it was asked for and replays it once authenticated -
  which also silently fixes execution updates on a freshly loaded editor.
- The login and /me responses carried no user id at all, though the type
  declared one, so "is this me?" always answered no. toJson() is what
  gets written to the database, where the id lives in _id, so the id is
  added to the response rather than to the record.
- transformWorkflow dropped _version, so the editor had no version to
  send back.

Verified against two real users: the roster shows each of them once,
loses somebody when their connection closes, and survives a cold load;
the live notice names the person who saved; a stale save is refused with
the changes still on the canvas and Save still available; an ordinary
save is untouched.
fszontagh hai 1 mes
pai
achega
239b85445d

+ 8 - 1
src/webserver/api/auth_controller.cpp

@@ -67,6 +67,11 @@ void AuthController::login(const httplib::Request& req, httplib::Response& res,
         response["refreshToken"] = login_resp.refresh_token;
         response["expiresAt"] = login_resp.expires_at;
         response["user"] = login_resp.user.toJson();
+        // toJson() is what gets written to the database, where the id lives in
+        // the separate _id key - so it leaves the id out. Over the API the id
+        // is exactly what a client needs to recognise itself, and its absence
+        // meant "is this me?" quietly answered no every time.
+        response["user"]["id"] = login_resp.user.id;
 
         sendJson(res, response);
     } catch (const std::exception& e) {
@@ -127,7 +132,9 @@ void AuthController::me(const httplib::Request& req, httplib::Response& res,
         return;
     }
 
-    sendJson(res, result.value().toJson());
+    auto user_json = result.value().toJson();
+    user_json["id"] = result.value().id;
+    sendJson(res, user_json);
 }
 
 void AuthController::sendJson(httplib::Response& res, const nlohmann::json& data, int status) {

+ 47 - 0
src/webserver/api/workflow_controller.cpp

@@ -262,6 +262,16 @@ void WorkflowController::updateWorkflow(const httplib::Request& req, httplib::Re
         std::string id = req.matches[1];
         auto body = nlohmann::json::parse(req.body);
 
+        // What the editor had loaded when the user started changing it. Two
+        // people with the same workflow open used to mean whoever pressed save
+        // second silently threw the other's work away, with nothing anywhere to
+        // say it had happened.
+        int64_t expected_version = -1;
+        if (body.contains("expectedVersion") && body["expectedVersion"].is_number()) {
+            expected_version = body["expectedVersion"].get<int64_t>();
+        }
+        body.erase("expectedVersion");
+
         // Prevent updating metadata fields (managed by database)
         body.erase("id");
         body.erase("_id");
@@ -283,6 +293,32 @@ void WorkflowController::updateWorkflow(const httplib::Request& req, httplib::Re
             was_active = before.value().value("active", false);
         }
 
+        // Refuse rather than overwrite. The client is told the version it is
+        // behind by and gets the current record back, so it can show what
+        // changed instead of just losing.
+        //
+        // This is a read then a write, not one atomic compare-and-set: the
+        // database can do that, but only for a whole-record replace, and this
+        // endpoint writes a patch - a replace here would drop every field the
+        // client did not send, ownerId among them. The window is the few
+        // milliseconds between the two calls, and what it protects against is
+        // two people editing for minutes, so it catches what it is for.
+        if (expected_version >= 0 && before.ok()) {
+            const int64_t current_version = before.value().value("_version", int64_t{0});
+            if (current_version != expected_version) {
+                nlohmann::json conflict = {
+                    {"error", "This workflow was saved by someone else while you were editing it"},
+                    {"code", "version_conflict"},
+                    {"expectedVersion", expected_version},
+                    {"currentVersion", current_version},
+                    {"current", before.value()},
+                };
+                res.status = 409;
+                res.set_content(conflict.dump(), "application/json");
+                return;
+            }
+        }
+
         // This is a replace, not a merge, so a body that does not mention
         // "active" would drop it - and a workflow that was running quietly
         // stops, with nothing recorded to say why. Saving a change to a node
@@ -328,6 +364,17 @@ void WorkflowController::updateWorkflow(const httplib::Request& req, httplib::Re
             }
 
             ws_server_.broadcast("workflows.updated", workflow.value());
+
+            // Named, and on the workflow's own channel, so another editor with
+            // it open is told who saved and can leave the saver's own window
+            // alone.
+            ws_server_.broadcast("workflows." + id + ".saved", {
+                {"workflowId", id},
+                {"userId", ctx.user_id},
+                {"username", ctx.username},
+                {"version", workflow.value().value("_version", int64_t{0})},
+            });
+
             sendJson(res, workflow.value());
         } else {
             sendJson(res, {{"id", id}, {"version", result.value()}});

+ 80 - 6
src/webserver/websocket_server.cpp

@@ -158,13 +158,28 @@ int WebSocketServer::onConnect(struct lws* wsi) {
 }
 
 void WebSocketServer::onDisconnect(struct lws* wsi) {
-    std::unique_lock lock(clients_mutex_);
+    // Which workflows this client was watching has to be read before it is
+    // erased, but the roster can only be republished after - and publishing
+    // takes the same lock. So: collect, drop the lock, then publish.
+    std::vector<std::string> was_watching;
 
-    auto it = clients_.find(wsi);
-    if (it != clients_.end()) {
-        LOG_DEBUG("WebSocket client disconnected: {}", it->second->id);
-        client_id_map_.erase(it->second->id);
-        clients_.erase(it);
+    {
+        std::unique_lock lock(clients_mutex_);
+
+        auto it = clients_.find(wsi);
+        if (it != clients_.end()) {
+            LOG_DEBUG("WebSocket client disconnected: {}", it->second->id);
+            for (const auto& sub : it->second->subscriptions) {
+                std::string workflow_id = presenceWorkflowId(sub);
+                if (!workflow_id.empty()) was_watching.push_back(workflow_id);
+            }
+            client_id_map_.erase(it->second->id);
+            clients_.erase(it);
+        }
+    }
+
+    for (const auto& workflow_id : was_watching) {
+        publishPresence(workflow_id);
     }
 }
 
@@ -249,10 +264,12 @@ void WebSocketServer::handleAuth(WebSocketClient& client, const nlohmann::json&
     if (result.ok()) {
         client.authenticated = true;
         client.user_id = result.value().user_id;
+        client.username = result.value().username;
 
         nlohmann::json response;
         response["type"] = "auth_success";
         response["userId"] = client.user_id;
+        response["username"] = client.username;
         sendToClient(client.id, response);
 
         LOG_DEBUG("WebSocket client authenticated: {}", client.id);
@@ -276,6 +293,11 @@ void WebSocketServer::handleSubscribe(WebSocketClient& client, const nlohmann::j
     auto channels = message.value("channels", std::vector<std::string>{});
     for (const auto& channel : channels) {
         subscribe(client.id, channel);
+        // Publishing after subscribing means the arriving client is itself on
+        // the channel, so it gets the roster as its first message rather than
+        // waiting for somebody else to come or go.
+        std::string workflow_id = presenceWorkflowId(channel);
+        if (!workflow_id.empty()) publishPresence(workflow_id);
     }
 }
 
@@ -283,9 +305,61 @@ void WebSocketServer::handleUnsubscribe(WebSocketClient& client, const nlohmann:
     auto channels = message.value("channels", std::vector<std::string>{});
     for (const auto& channel : channels) {
         unsubscribe(client.id, channel);
+        std::string workflow_id = presenceWorkflowId(channel);
+        if (!workflow_id.empty()) publishPresence(workflow_id);
     }
 }
 
+std::string WebSocketServer::presenceChannel(const std::string& workflow_id) {
+    return "workflows." + workflow_id + ".presence";
+}
+
+std::string WebSocketServer::presenceWorkflowId(const std::string& channel) {
+    static const std::string prefix = "workflows.";
+    static const std::string suffix = ".presence";
+    if (!channel.starts_with(prefix) || !channel.ends_with(suffix)) return "";
+
+    // A wildcard subscription is a listener, not somebody with the workflow
+    // open, and counting it would put phantom names in the roster.
+    std::string id = channel.substr(prefix.size(),
+                                    channel.size() - prefix.size() - suffix.size());
+    if (id.empty() || id.find('*') != std::string::npos) return "";
+    return id;
+}
+
+void WebSocketServer::publishPresence(const std::string& workflow_id) {
+    const std::string channel = presenceChannel(workflow_id);
+
+    nlohmann::json viewers = nlohmann::json::array();
+    {
+        std::shared_lock lock(clients_mutex_);
+
+        // One person with the editor open in two tabs is one person, so the
+        // roster is by user rather than by connection.
+        std::unordered_map<std::string, size_t> seen;
+        for (const auto& [wsi, client] : clients_) {
+            (void)wsi;
+            if (!client->authenticated) continue;
+            if (!client->subscriptions.contains(channel)) continue;
+
+            auto it = seen.find(client->user_id);
+            if (it != seen.end()) {
+                viewers[it->second]["connections"] =
+                    viewers[it->second]["connections"].get<int>() + 1;
+                continue;
+            }
+            seen.emplace(client->user_id, viewers.size());
+            viewers.push_back({
+                {"userId", client->user_id},
+                {"username", client->username},
+                {"connections", 1},
+            });
+        }
+    }
+
+    broadcast(channel, {{"workflowId", workflow_id}, {"viewers", viewers}});
+}
+
 void WebSocketServer::broadcast(const std::string& channel, const nlohmann::json& data) {
     nlohmann::json message;
     message["channel"] = channel;

+ 8 - 0
src/webserver/websocket_server.hpp

@@ -19,6 +19,7 @@ struct WebSocketClient {
     struct lws* wsi = nullptr;
     std::string id;
     std::string user_id;
+    std::string username;
     bool authenticated = false;
     std::unordered_set<std::string> subscriptions;
     std::vector<std::string> send_queue;
@@ -54,6 +55,13 @@ public:
     void subscribe(const std::string& client_id, const std::string& channel);
     void unsubscribe(const std::string& client_id, const std::string& channel);
 
+    // Who currently has a workflow open. Derived from the subscriptions rather
+    // than stored separately, so a client that drops off never leaves a ghost
+    // behind - there is no second list that can disagree with the connections.
+    void publishPresence(const std::string& workflow_id);
+    static std::string presenceWorkflowId(const std::string& channel);
+    static std::string presenceChannel(const std::string& workflow_id);
+
     // Message handler
     void setMessageHandler(WebSocketMessageHandler handler);
 

+ 14 - 0
webui/src/api/client.ts

@@ -75,6 +75,11 @@ export class WebSocketClient {
   private ws: WebSocket | null = null
   private reconnectTimeout: number | null = null
   private messageHandlers: Map<string, Set<(data: any) => void>> = new Map()
+  // What this client wants to be subscribed to, as opposed to what it has
+  // managed to say so far. send() drops anything queued before the socket is
+  // open, so a page that subscribes while connecting - which is every page,
+  // on a cold load - was silently never subscribed to anything.
+  private desiredChannels = new Set<string>()
   private reconnectAttempts: number = 0
   private maxReconnectAttempts: number = 5
   private baseReconnectDelay: number = 1000
@@ -139,6 +144,13 @@ export class WebSocketClient {
     this.ws.onmessage = (event) => {
       try {
         const message = JSON.parse(event.data)
+
+        // Subscriptions are replayed here rather than in onopen because the
+        // server refuses a subscribe until it has accepted the token.
+        if (message.type === 'auth_success' && this.desiredChannels.size > 0) {
+          this.send({ type: 'subscribe', channels: [...this.desiredChannels] })
+        }
+
         const channel = message.channel || message.type
         const data = message.data || message
 
@@ -213,10 +225,12 @@ export class WebSocketClient {
   }
 
   subscribe(channels: string[]) {
+    channels.forEach((c) => this.desiredChannels.add(c))
     this.send({ type: 'subscribe', channels })
   }
 
   unsubscribe(channels: string[]) {
+    channels.forEach((c) => this.desiredChannels.delete(c))
     this.send({ type: 'unsubscribe', channels })
   }
 

+ 5 - 0
webui/src/api/workflows.ts

@@ -31,6 +31,10 @@ export interface Workflow {
   ownerId?: string
   createdBy?: string
   updatedBy?: string
+  // The database's record version, bumped on every write. The editor sends the
+  // one it loaded back on save so the server can refuse to overwrite somebody
+  // else's work rather than silently discarding it.
+  version?: number
 }
 
 export interface NodeOutput {
@@ -86,6 +90,7 @@ function transformWorkflow(data: any): Workflow {
     // person. Kept here for when the backend starts populating it.
     createdBy: data._created_by || undefined,
     updatedBy: data._updated_by || undefined,
+    version: typeof data._version === 'number' ? data._version : undefined,
   }
 }
 

+ 25 - 0
webui/src/components/workflow/EditorHeader.tsx

@@ -28,6 +28,7 @@ interface EditorHeaderProps {
   onExecuteTrigger: (triggerNodeId: string) => void
   onAddNode: () => void
   onAutoLayout: () => void
+  otherViewers?: { userId: string; username: string }[]
   onUndo: () => void
   onRedo: () => void
   canUndo: boolean
@@ -61,6 +62,7 @@ export function EditorHeader({
   onExecuteTrigger,
   onAddNode,
   onAutoLayout,
+  otherViewers = [],
   onUndo,
   onRedo,
   canUndo,
@@ -184,6 +186,29 @@ export function EditorHeader({
           <Plus className="w-4 h-4" />
           Add Node
         </button>
+        {/* Who else has this open. Only shown when there is somebody, so a
+            single-user install never sees it. */}
+        {otherViewers.length > 0 && (
+          <div
+            className="flex items-center -space-x-2 mr-1"
+            title={`Also editing: ${otherViewers.map((v) => v.username).join(', ')}`}
+          >
+            {otherViewers.slice(0, 3).map((v) => (
+              <span
+                key={v.userId}
+                className="w-6 h-6 rounded-full bg-indigo-500 text-white text-xs flex items-center justify-center ring-2 ring-white dark:ring-slate-800 uppercase"
+              >
+                {v.username.slice(0, 1)}
+              </span>
+            ))}
+            {otherViewers.length > 3 && (
+              <span className="w-6 h-6 rounded-full bg-gray-400 text-white text-xs flex items-center justify-center ring-2 ring-white dark:ring-slate-800">
+                +{otherViewers.length - 3}
+              </span>
+            )}
+          </div>
+        )}
+
         {/* Shown rather than left to the shortcut: an editor that can undo and
             does not say so is one nobody tries it in. */}
         <button

+ 82 - 0
webui/src/hooks/useWorkflowPresence.ts

@@ -0,0 +1,82 @@
+import { useCallback, useEffect, useRef, useState } from 'react'
+import { wsClient } from '../api/client'
+import { useAuthStore } from '../stores/authStore'
+
+/**
+ * Who else has this workflow open, and whether they have saved it.
+ *
+ * Two people editing the same workflow used to be entirely invisible: no sign
+ * anybody else was there, and whoever pressed save second silently threw the
+ * other's work away. This does not make editing collaborative - it makes it
+ * honest. You can see who is here, and you are told when the copy in front of
+ * you has been overtaken.
+ *
+ * The roster comes from the server, derived from who is subscribed to the
+ * workflow's presence channel, so a closed tab or a dropped connection removes
+ * somebody without anything having to remember to say so.
+ */
+
+export interface Viewer {
+  userId: string
+  username: string
+  connections: number
+}
+
+export interface SavedByOther {
+  userId: string
+  username: string
+  version: number
+}
+
+export function useWorkflowPresence(workflowId: string | undefined) {
+  const currentUserId = useAuthStore((s) => s.user?.id)
+
+  const [viewers, setViewers] = useState<Viewer[]>([])
+  const [savedByOther, setSavedByOther] = useState<SavedByOther | null>(null)
+
+  // Read inside the handler rather than captured, so resubscribing is not
+  // needed every time the signed-in user is resolved.
+  const currentUserIdRef = useRef(currentUserId)
+  currentUserIdRef.current = currentUserId
+
+  useEffect(() => {
+    if (!workflowId) return
+
+    const presence = `workflows.${workflowId}.presence`
+    const saved = `workflows.${workflowId}.saved`
+
+    const onPresence = (data: any) => {
+      setViewers(Array.isArray(data?.viewers) ? data.viewers : [])
+    }
+
+    const onSaved = (data: any) => {
+      // Your own save is not news, and would otherwise warn you about yourself
+      // every time you pressed the button.
+      if (data?.userId && data.userId === currentUserIdRef.current) return
+      setSavedByOther({
+        userId: data?.userId ?? '',
+        username: data?.username || 'Someone',
+        version: Number(data?.version) || 0,
+      })
+    }
+
+    wsClient.on(presence, onPresence)
+    wsClient.on(saved, onSaved)
+    // Subscribing before the socket is open is fine: the client remembers what
+    // it was asked for and says so once it has connected and authenticated.
+    wsClient.subscribe([presence, saved])
+
+    return () => {
+      wsClient.unsubscribe([presence, saved])
+      wsClient.off(presence, onPresence)
+      wsClient.off(saved, onSaved)
+      setViewers([])
+      setSavedByOther(null)
+    }
+  }, [workflowId])
+
+  const others = viewers.filter((v) => v.userId !== currentUserId)
+  const dismissSavedByOther = useCallback(() => setSavedByOther(null), [])
+
+  return { viewers, others, savedByOther, dismissSavedByOther }
+}

+ 69 - 2
webui/src/pages/WorkflowEditorPage.tsx

@@ -28,7 +28,9 @@ import ReactFlow, {
 import { BackgroundVariant } from '@reactflow/background'
 import 'reactflow/dist/style.css'
 import { useCallback, useEffect, useState, useRef, useMemo } from 'react'
+import { AlertTriangle } from 'lucide-react'
 import { wsClient } from '../api/client'
+import { useWorkflowPresence } from '../hooks/useWorkflowPresence'
 import type { AvailableField } from '../components/ConditionBuilder'
 import type { FieldInfo } from '../components/AvailableDataPanel'
 import { ExecutionListPanel } from '../components/workflow/ExecutionListPanel'
@@ -559,6 +561,19 @@ function WorkflowEditorInner() {
     enabled: !!id,
   })
 
+  // The version currently on the canvas. Kept in a ref because saving reads it
+  // and must not be re-created every time it changes.
+  const { others: otherViewers, savedByOther, dismissSavedByOther } = useWorkflowPresence(id)
+  // The save handler needs the latest value without being rebuilt for it.
+  const savedByOtherRef = useRef(savedByOther)
+  savedByOtherRef.current = savedByOther
+
+  const loadedVersionRef = useRef<number | null>(null)
+  const [saveConflict, setSaveConflict] = useState<{ username: string; currentVersion: number } | null>(null)
+  useEffect(() => {
+    if (typeof workflow?.version === 'number') loadedVersionRef.current = workflow.version
+  }, [workflow])
+
   const { data: nodeDefinitions } = useQuery({
     queryKey: ['nodes'],
     queryFn: () => nodesApi.list(),
@@ -722,12 +737,27 @@ function WorkflowEditorInner() {
 
   const saveMutation = useMutation({
     mutationFn: (data: Partial<Workflow>) => workflowsApi.update(id!, data),
-    onSuccess: () => {
+    onSuccess: (saved: any) => {
       queryClient.invalidateQueries({ queryKey: ['workflow', id] })
       setHasChanges(false)
+      setSaveConflict(null)
+      dismissSavedByOther()
+      if (typeof saved?.version === 'number') loadedVersionRef.current = saved.version
       showToast('success', 'Workflow saved successfully')
     },
     onError: (error: any) => {
+      // A refused save is not a failure to report and forget: the work is still
+      // on the canvas, and the user has to be given the choice of what happens
+      // to it rather than being told "failed" and left to guess.
+      const conflict = error?.response?.status === 409 ? error.response.data : null
+      if (conflict) {
+        setSaveConflict({
+          username: conflict.current?.updatedByUsername || savedByOtherRef.current?.username || 'Someone else',
+          currentVersion: Number(conflict.currentVersion) || 0,
+        })
+        showToast('error', 'Someone else saved this workflow while you were editing it')
+        return
+      }
       showToast('error', `Failed to save: ${error.message}`)
     },
   })
@@ -1463,7 +1493,11 @@ function WorkflowEditorInner() {
       nodes: workflowNodes,
       connections: workflowConnections,
       settings: workflowSettings,
-    })
+      // The version this editor loaded. The server refuses the save if the
+      // stored workflow has moved on since, rather than overwriting whoever
+      // got there first.
+      expectedVersion: loadedVersionRef.current ?? undefined,
+    } as any)
   }, [nodes, edges, saveMutation, workflowSettings, workflowName])
 
   // Handle workflow rename
@@ -2670,6 +2704,38 @@ function WorkflowEditorInner() {
         isViewingExecution && pinnedExecution ? 'pt-12' : ''
       }`}
     >
+      {/* Someone else has saved this workflow, or refused this editor's save.
+          Both mean the same thing to the person sitting here: what is on the
+          canvas is no longer what is stored. */}
+      {(saveConflict || savedByOther) && !isViewingExecution && (
+        <div className="flex items-center gap-3 px-4 py-2 bg-amber-50 dark:bg-amber-900/30 border-b border-amber-300 dark:border-amber-700 text-sm text-amber-900 dark:text-amber-200">
+          <AlertTriangle className="w-4 h-4 shrink-0" />
+          <span className="flex-1">
+            {saveConflict
+              ? `${saveConflict.username} saved this workflow while you were editing it, so your save was refused. Your changes are still here - reload to take their version, or copy what you need out first.`
+              : `${savedByOther!.username} just saved this workflow. Your copy is out of date${hasChanges ? ', and you have unsaved changes' : ''}.`}
+          </span>
+          <button
+            onClick={() => {
+              // A reload throws away whatever is on the canvas, so it is a
+              // button somebody presses, never something that happens to them
+              // mid-edit.
+              queryClient.invalidateQueries({ queryKey: ['workflow', id] })
+              window.location.reload()
+            }}
+            className="px-2 py-1 rounded bg-amber-600 text-white hover:bg-amber-700 shrink-0"
+          >
+            Reload
+          </button>
+          <button
+            onClick={() => { setSaveConflict(null); dismissSavedByOther() }}
+            className="px-2 py-1 rounded hover:bg-amber-100 dark:hover:bg-amber-800 shrink-0"
+          >
+            Dismiss
+          </button>
+        </div>
+      )}
+
       {/* Execution Viewer Banner */}
       {isViewingExecution && pinnedExecution && (
         <ExecutionViewerBanner
@@ -2706,6 +2772,7 @@ function WorkflowEditorInner() {
         onExecuteTrigger={executeTrigger}
         onAddNode={() => setShowNodePicker(true)}
         onAutoLayout={autoLayout}
+        otherViewers={otherViewers}
         onUndo={() => { if (history.undo()) setHasChanges(true) }}
         onRedo={() => { if (history.redo()) setHasChanges(true) }}
         canUndo={history.canUndo}