|
@@ -162,6 +162,7 @@ void WebSocketServer::onDisconnect(struct lws* wsi) {
|
|
|
// erased, but the roster can only be republished after - and publishing
|
|
// erased, but the roster can only be republished after - and publishing
|
|
|
// takes the same lock. So: collect, drop the lock, then publish.
|
|
// takes the same lock. So: collect, drop the lock, then publish.
|
|
|
std::vector<std::string> was_watching;
|
|
std::vector<std::string> was_watching;
|
|
|
|
|
+ std::string gone_client_id;
|
|
|
|
|
|
|
|
{
|
|
{
|
|
|
std::unique_lock lock(clients_mutex_);
|
|
std::unique_lock lock(clients_mutex_);
|
|
@@ -169,6 +170,7 @@ void WebSocketServer::onDisconnect(struct lws* wsi) {
|
|
|
auto it = clients_.find(wsi);
|
|
auto it = clients_.find(wsi);
|
|
|
if (it != clients_.end()) {
|
|
if (it != clients_.end()) {
|
|
|
LOG_DEBUG("WebSocket client disconnected: {}", it->second->id);
|
|
LOG_DEBUG("WebSocket client disconnected: {}", it->second->id);
|
|
|
|
|
+ gone_client_id = it->second->id;
|
|
|
for (const auto& sub : it->second->subscriptions) {
|
|
for (const auto& sub : it->second->subscriptions) {
|
|
|
std::string workflow_id = presenceWorkflowId(sub);
|
|
std::string workflow_id = presenceWorkflowId(sub);
|
|
|
if (!workflow_id.empty()) was_watching.push_back(workflow_id);
|
|
if (!workflow_id.empty()) was_watching.push_back(workflow_id);
|
|
@@ -178,9 +180,17 @@ void WebSocketServer::onDisconnect(struct lws* wsi) {
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+ // A lock outlives its holder only if nothing notices they have gone, which
|
|
|
|
|
+ // is exactly how an editor ends up with a node nobody can touch again.
|
|
|
|
|
+ std::vector<std::string> unlocked;
|
|
|
|
|
+ releaseLocksOf(gone_client_id, &unlocked);
|
|
|
|
|
+
|
|
|
for (const auto& workflow_id : was_watching) {
|
|
for (const auto& workflow_id : was_watching) {
|
|
|
publishPresence(workflow_id);
|
|
publishPresence(workflow_id);
|
|
|
}
|
|
}
|
|
|
|
|
+ for (const auto& workflow_id : unlocked) {
|
|
|
|
|
+ publishLocks(workflow_id);
|
|
|
|
|
+ }
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
int WebSocketServer::onReceive(struct lws* wsi, const char* data, size_t len) {
|
|
int WebSocketServer::onReceive(struct lws* wsi, const char* data, size_t len) {
|
|
@@ -247,6 +257,12 @@ void WebSocketServer::processMessage(WebSocketClient& client, const nlohmann::js
|
|
|
handleSubscribe(client, message);
|
|
handleSubscribe(client, message);
|
|
|
} else if (type == "unsubscribe") {
|
|
} else if (type == "unsubscribe") {
|
|
|
handleUnsubscribe(client, message);
|
|
handleUnsubscribe(client, message);
|
|
|
|
|
+ } else if (type == "lock") {
|
|
|
|
|
+ handleLock(client, message);
|
|
|
|
|
+ } else if (type == "unlock") {
|
|
|
|
|
+ handleUnlock(client, message);
|
|
|
|
|
+ } else if (type == "node_moved") {
|
|
|
|
|
+ handleNodeMoved(client, message);
|
|
|
} else if (message_handler_) {
|
|
} else if (message_handler_) {
|
|
|
message_handler_(client.id, message);
|
|
message_handler_(client.id, message);
|
|
|
}
|
|
}
|
|
@@ -270,6 +286,10 @@ void WebSocketServer::handleAuth(WebSocketClient& client, const nlohmann::json&
|
|
|
response["type"] = "auth_success";
|
|
response["type"] = "auth_success";
|
|
|
response["userId"] = client.user_id;
|
|
response["userId"] = client.user_id;
|
|
|
response["username"] = client.username;
|
|
response["username"] = client.username;
|
|
|
|
|
+ // This connection's own id. A lock is per connection, so an editor has
|
|
|
|
|
+ // to be able to tell its own lock from one held by the same person in
|
|
|
|
|
+ // another tab.
|
|
|
|
|
+ response["clientId"] = client.id;
|
|
|
sendToClient(client.id, response);
|
|
sendToClient(client.id, response);
|
|
|
|
|
|
|
|
LOG_DEBUG("WebSocket client authenticated: {}", client.id);
|
|
LOG_DEBUG("WebSocket client authenticated: {}", client.id);
|
|
@@ -298,6 +318,12 @@ void WebSocketServer::handleSubscribe(WebSocketClient& client, const nlohmann::j
|
|
|
// waiting for somebody else to come or go.
|
|
// waiting for somebody else to come or go.
|
|
|
std::string workflow_id = presenceWorkflowId(channel);
|
|
std::string workflow_id = presenceWorkflowId(channel);
|
|
|
if (!workflow_id.empty()) publishPresence(workflow_id);
|
|
if (!workflow_id.empty()) publishPresence(workflow_id);
|
|
|
|
|
+
|
|
|
|
|
+ // The same for locks, which are otherwise only published when they
|
|
|
|
|
+ // change: an editor arriving after somebody started working would see
|
|
|
|
|
+ // an unlocked canvas and walk straight into a node already taken.
|
|
|
|
|
+ std::string locks_workflow = channelWorkflowId(channel, ".locks");
|
|
|
|
|
+ if (!locks_workflow.empty()) publishLocks(locks_workflow);
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
|
|
|
|
|
@@ -314,9 +340,9 @@ std::string WebSocketServer::presenceChannel(const std::string& workflow_id) {
|
|
|
return "workflows." + workflow_id + ".presence";
|
|
return "workflows." + workflow_id + ".presence";
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
-std::string WebSocketServer::presenceWorkflowId(const std::string& channel) {
|
|
|
|
|
|
|
+std::string WebSocketServer::channelWorkflowId(const std::string& channel,
|
|
|
|
|
+ const std::string& suffix) {
|
|
|
static const std::string prefix = "workflows.";
|
|
static const std::string prefix = "workflows.";
|
|
|
- static const std::string suffix = ".presence";
|
|
|
|
|
if (!channel.starts_with(prefix) || !channel.ends_with(suffix)) return "";
|
|
if (!channel.starts_with(prefix) || !channel.ends_with(suffix)) return "";
|
|
|
|
|
|
|
|
// A wildcard subscription is a listener, not somebody with the workflow
|
|
// A wildcard subscription is a listener, not somebody with the workflow
|
|
@@ -327,6 +353,10 @@ std::string WebSocketServer::presenceWorkflowId(const std::string& channel) {
|
|
|
return id;
|
|
return id;
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+std::string WebSocketServer::presenceWorkflowId(const std::string& channel) {
|
|
|
|
|
+ return channelWorkflowId(channel, ".presence");
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
void WebSocketServer::publishPresence(const std::string& workflow_id) {
|
|
void WebSocketServer::publishPresence(const std::string& workflow_id) {
|
|
|
const std::string channel = presenceChannel(workflow_id);
|
|
const std::string channel = presenceChannel(workflow_id);
|
|
|
|
|
|
|
@@ -360,6 +390,143 @@ void WebSocketServer::publishPresence(const std::string& workflow_id) {
|
|
|
broadcast(channel, {{"workflowId", workflow_id}, {"viewers", viewers}});
|
|
broadcast(channel, {{"workflowId", workflow_id}, {"viewers", viewers}});
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+
|
|
|
|
|
+// ---------------------------------------------------------------------------
|
|
|
|
|
+// Node locks and live movement
|
|
|
|
|
+//
|
|
|
|
|
+// While somebody is dragging or configuring a node, everyone else is kept off
|
|
|
|
|
+// that one node - not off the workflow. Two people working on different parts
|
|
|
|
|
+// of the same flow is the normal case and should stay possible; two people
|
|
|
|
|
+// dragging the same node is the one that produces nonsense.
|
|
|
|
|
+// ---------------------------------------------------------------------------
|
|
|
|
|
+
|
|
|
|
|
+void WebSocketServer::handleLock(WebSocketClient& client, const nlohmann::json& message) {
|
|
|
|
|
+ if (!client.authenticated) return;
|
|
|
|
|
+
|
|
|
|
|
+ const std::string workflow_id = message.value("workflowId", "");
|
|
|
|
|
+ const std::string node_id = message.value("nodeId", "");
|
|
|
|
|
+ const std::string kind = message.value("kind", "editing");
|
|
|
|
|
+ if (workflow_id.empty() || node_id.empty()) return;
|
|
|
|
|
+
|
|
|
|
|
+ bool granted = false;
|
|
|
|
|
+ std::string holder_username;
|
|
|
|
|
+ {
|
|
|
|
|
+ std::lock_guard<std::mutex> lock(locks_mutex_);
|
|
|
|
|
+ auto& nodes = locks_[workflow_id];
|
|
|
|
|
+ auto it = nodes.find(node_id);
|
|
|
|
|
+ if (it == nodes.end() || it->second.client_id == client.id) {
|
|
|
|
|
+ // Re-claiming your own lock is how the kind changes from dragging
|
|
|
|
|
+ // to editing without a release in between.
|
|
|
|
|
+ nodes[node_id] = NodeLock{client.id, client.user_id, client.username, kind};
|
|
|
|
|
+ granted = true;
|
|
|
|
|
+ } else {
|
|
|
|
|
+ holder_username = it->second.username;
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ if (!granted) {
|
|
|
|
|
+ // Said out loud rather than ignored: a click that silently does nothing
|
|
|
|
|
+ // reads as the editor being broken.
|
|
|
|
|
+ nlohmann::json response;
|
|
|
|
|
+ response["type"] = "lock_denied";
|
|
|
|
|
+ response["workflowId"] = workflow_id;
|
|
|
|
|
+ response["nodeId"] = node_id;
|
|
|
|
|
+ response["heldBy"] = holder_username;
|
|
|
|
|
+ sendToClient(client.id, response);
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ publishLocks(workflow_id);
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void WebSocketServer::handleUnlock(WebSocketClient& client, const nlohmann::json& message) {
|
|
|
|
|
+ const std::string workflow_id = message.value("workflowId", "");
|
|
|
|
|
+ const std::string node_id = message.value("nodeId", "");
|
|
|
|
|
+ if (workflow_id.empty() || node_id.empty()) return;
|
|
|
|
|
+
|
|
|
|
|
+ bool changed = false;
|
|
|
|
|
+ {
|
|
|
|
|
+ std::lock_guard<std::mutex> lock(locks_mutex_);
|
|
|
|
|
+ auto wf = locks_.find(workflow_id);
|
|
|
|
|
+ if (wf != locks_.end()) {
|
|
|
|
|
+ auto it = wf->second.find(node_id);
|
|
|
|
|
+ // Only the holder can let go - otherwise a stale release from an
|
|
|
|
|
+ // editor that has already moved on would free somebody else's node.
|
|
|
|
|
+ if (it != wf->second.end() && it->second.client_id == client.id) {
|
|
|
|
|
+ wf->second.erase(it);
|
|
|
|
|
+ changed = true;
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ if (changed) publishLocks(workflow_id);
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void WebSocketServer::handleNodeMoved(WebSocketClient& client, const nlohmann::json& message) {
|
|
|
|
|
+ if (!client.authenticated) return;
|
|
|
|
|
+
|
|
|
|
|
+ const std::string workflow_id = message.value("workflowId", "");
|
|
|
|
|
+ const std::string node_id = message.value("nodeId", "");
|
|
|
|
|
+ if (workflow_id.empty() || node_id.empty()) return;
|
|
|
|
|
+
|
|
|
|
|
+ // Only the editor holding the node may say where it is. Without this,
|
|
|
|
|
+ // anyone could move anybody's node by sending the message directly.
|
|
|
|
|
+ {
|
|
|
|
|
+ std::lock_guard<std::mutex> lock(locks_mutex_);
|
|
|
|
|
+ auto wf = locks_.find(workflow_id);
|
|
|
|
|
+ if (wf == locks_.end()) return;
|
|
|
|
|
+ auto it = wf->second.find(node_id);
|
|
|
|
|
+ if (it == wf->second.end() || it->second.client_id != client.id) return;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Not stored: a position in flight is worth nothing once it has been
|
|
|
|
|
+ // delivered, and the authoritative one is whatever gets saved.
|
|
|
|
|
+ broadcast("workflows." + workflow_id + ".moves", {
|
|
|
|
|
+ {"workflowId", workflow_id},
|
|
|
|
|
+ {"nodeId", node_id},
|
|
|
|
|
+ {"position", message.value("position", nlohmann::json::object())},
|
|
|
|
|
+ {"userId", client.user_id},
|
|
|
|
|
+ {"clientId", client.id},
|
|
|
|
|
+ });
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void WebSocketServer::releaseLocksOf(const std::string& client_id,
|
|
|
|
|
+ std::vector<std::string>* touched_workflows) {
|
|
|
|
|
+ std::lock_guard<std::mutex> lock(locks_mutex_);
|
|
|
|
|
+ for (auto& [workflow_id, nodes] : locks_) {
|
|
|
|
|
+ bool changed = false;
|
|
|
|
|
+ for (auto it = nodes.begin(); it != nodes.end();) {
|
|
|
|
|
+ if (it->second.client_id == client_id) {
|
|
|
|
|
+ it = nodes.erase(it);
|
|
|
|
|
+ changed = true;
|
|
|
|
|
+ } else {
|
|
|
|
|
+ ++it;
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ if (changed && touched_workflows) touched_workflows->push_back(workflow_id);
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void WebSocketServer::publishLocks(const std::string& workflow_id) {
|
|
|
|
|
+ nlohmann::json held = nlohmann::json::array();
|
|
|
|
|
+ {
|
|
|
|
|
+ std::lock_guard<std::mutex> lock(locks_mutex_);
|
|
|
|
|
+ auto wf = locks_.find(workflow_id);
|
|
|
|
|
+ if (wf != locks_.end()) {
|
|
|
|
|
+ for (const auto& [node_id, holder] : wf->second) {
|
|
|
|
|
+ held.push_back({
|
|
|
|
|
+ {"nodeId", node_id},
|
|
|
|
|
+ {"userId", holder.user_id},
|
|
|
|
|
+ {"username", holder.username},
|
|
|
|
|
+ {"clientId", holder.client_id},
|
|
|
|
|
+ {"kind", holder.kind},
|
|
|
|
|
+ });
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ broadcast("workflows." + workflow_id + ".locks",
|
|
|
|
|
+ {{"workflowId", workflow_id}, {"locks", held}});
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
void WebSocketServer::broadcast(const std::string& channel, const nlohmann::json& data) {
|
|
void WebSocketServer::broadcast(const std::string& channel, const nlohmann::json& data) {
|
|
|
nlohmann::json message;
|
|
nlohmann::json message;
|
|
|
message["channel"] = channel;
|
|
message["channel"] = channel;
|