mirror of
https://github.com/langgenius/dify.git
synced 2026-09-05 00:31:19 +08:00
fix: ensure leader online to accept graph change
This commit is contained in:
parent
4432c86f57
commit
43b9bcd4e0
@ -79,6 +79,29 @@ class WorkflowCollaborationService:
|
|||||||
if not event_type:
|
if not event_type:
|
||||||
return {"msg": "invalid event type"}, 400
|
return {"msg": "invalid event type"}, 400
|
||||||
|
|
||||||
|
if event_type == "sync_request":
|
||||||
|
leader_sid = self._repository.get_current_leader(workflow_id)
|
||||||
|
if leader_sid and self.is_session_active(workflow_id, leader_sid):
|
||||||
|
target_sid = leader_sid
|
||||||
|
else:
|
||||||
|
if leader_sid:
|
||||||
|
self._repository.delete_leader(workflow_id)
|
||||||
|
target_sid = self._select_active_leader(workflow_id, preferred_sid=sid)
|
||||||
|
if target_sid:
|
||||||
|
self._repository.set_leader(workflow_id, target_sid)
|
||||||
|
self.broadcast_leader_change(workflow_id, target_sid)
|
||||||
|
|
||||||
|
if not target_sid:
|
||||||
|
return {"msg": "no_active_leader"}, 200
|
||||||
|
|
||||||
|
self._socketio.emit(
|
||||||
|
"collaboration_update",
|
||||||
|
{"type": event_type, "userId": user_id, "data": event_data, "timestamp": timestamp},
|
||||||
|
room=target_sid,
|
||||||
|
)
|
||||||
|
|
||||||
|
return {"msg": "sync_request_forwarded"}, 200
|
||||||
|
|
||||||
self._socketio.emit(
|
self._socketio.emit(
|
||||||
"collaboration_update",
|
"collaboration_update",
|
||||||
{"type": event_type, "userId": user_id, "data": event_data, "timestamp": timestamp},
|
{"type": event_type, "userId": user_id, "data": event_data, "timestamp": timestamp},
|
||||||
@ -177,6 +200,14 @@ class WorkflowCollaborationService:
|
|||||||
self._repository.set_leader(workflow_id, sid)
|
self._repository.set_leader(workflow_id, sid)
|
||||||
self.broadcast_leader_change(workflow_id, sid)
|
self.broadcast_leader_change(workflow_id, sid)
|
||||||
|
|
||||||
|
def _select_active_leader(self, workflow_id: str, preferred_sid: str | None = None) -> str | None:
|
||||||
|
session_sids = [sid for sid in self._repository.get_session_sids(workflow_id) if self.is_session_active(workflow_id, sid)]
|
||||||
|
if not session_sids:
|
||||||
|
return None
|
||||||
|
if preferred_sid and preferred_sid in session_sids:
|
||||||
|
return preferred_sid
|
||||||
|
return session_sids[0]
|
||||||
|
|
||||||
def is_session_active(self, workflow_id: str, sid: str) -> bool:
|
def is_session_active(self, workflow_id: str, sid: str) -> bool:
|
||||||
if not sid:
|
if not sid:
|
||||||
return False
|
return False
|
||||||
|
|||||||
@ -79,6 +79,76 @@ class TestWorkflowCollaborationService:
|
|||||||
skip_sid="sid-1",
|
skip_sid="sid-1",
|
||||||
)
|
)
|
||||||
|
|
||||||
|
def test_relay_collaboration_event_sync_request_forwards_to_active_leader(
|
||||||
|
self, service: tuple[WorkflowCollaborationService, Mock, Mock]
|
||||||
|
) -> None:
|
||||||
|
collaboration_service, repository, socketio = service
|
||||||
|
repository.get_sid_mapping.return_value = {"workflow_id": "wf-1", "user_id": "u-1"}
|
||||||
|
repository.get_current_leader.return_value = "sid-leader"
|
||||||
|
payload = {"type": "sync_request", "data": {"reason": "join"}, "timestamp": 123}
|
||||||
|
|
||||||
|
with (
|
||||||
|
patch.object(collaboration_service, "refresh_session_state"),
|
||||||
|
patch.object(collaboration_service, "is_session_active", return_value=True),
|
||||||
|
):
|
||||||
|
result = collaboration_service.relay_collaboration_event("sid-1", payload)
|
||||||
|
|
||||||
|
assert result == ({"msg": "sync_request_forwarded"}, 200)
|
||||||
|
socketio.emit.assert_called_once_with(
|
||||||
|
"collaboration_update",
|
||||||
|
{"type": "sync_request", "userId": "u-1", "data": {"reason": "join"}, "timestamp": 123},
|
||||||
|
room="sid-leader",
|
||||||
|
)
|
||||||
|
repository.set_leader.assert_not_called()
|
||||||
|
|
||||||
|
def test_relay_collaboration_event_sync_request_reelects_active_leader(
|
||||||
|
self, service: tuple[WorkflowCollaborationService, Mock, Mock]
|
||||||
|
) -> None:
|
||||||
|
collaboration_service, repository, socketio = service
|
||||||
|
repository.get_sid_mapping.return_value = {"workflow_id": "wf-1", "user_id": "u-1"}
|
||||||
|
repository.get_current_leader.return_value = "sid-old"
|
||||||
|
repository.get_session_sids.return_value = ["sid-2", "sid-3"]
|
||||||
|
payload = {"type": "sync_request", "data": {"reason": "join"}, "timestamp": 123}
|
||||||
|
|
||||||
|
def _is_session_active(_workflow_id: str, session_sid: str) -> bool:
|
||||||
|
return session_sid != "sid-old"
|
||||||
|
|
||||||
|
with (
|
||||||
|
patch.object(collaboration_service, "refresh_session_state"),
|
||||||
|
patch.object(collaboration_service, "broadcast_leader_change") as broadcast_leader_change,
|
||||||
|
patch.object(collaboration_service, "is_session_active", side_effect=_is_session_active),
|
||||||
|
):
|
||||||
|
result = collaboration_service.relay_collaboration_event("sid-2", payload)
|
||||||
|
|
||||||
|
assert result == ({"msg": "sync_request_forwarded"}, 200)
|
||||||
|
repository.delete_leader.assert_called_once_with("wf-1")
|
||||||
|
repository.set_leader.assert_called_once_with("wf-1", "sid-2")
|
||||||
|
broadcast_leader_change.assert_called_once_with("wf-1", "sid-2")
|
||||||
|
socketio.emit.assert_called_once_with(
|
||||||
|
"collaboration_update",
|
||||||
|
{"type": "sync_request", "userId": "u-1", "data": {"reason": "join"}, "timestamp": 123},
|
||||||
|
room="sid-2",
|
||||||
|
)
|
||||||
|
|
||||||
|
def test_relay_collaboration_event_sync_request_returns_when_no_active_leader(
|
||||||
|
self, service: tuple[WorkflowCollaborationService, Mock, Mock]
|
||||||
|
) -> None:
|
||||||
|
collaboration_service, repository, socketio = service
|
||||||
|
repository.get_sid_mapping.return_value = {"workflow_id": "wf-1", "user_id": "u-1"}
|
||||||
|
repository.get_current_leader.return_value = "sid-old"
|
||||||
|
repository.get_session_sids.return_value = []
|
||||||
|
payload = {"type": "sync_request", "data": {"reason": "join"}, "timestamp": 123}
|
||||||
|
|
||||||
|
with (
|
||||||
|
patch.object(collaboration_service, "refresh_session_state"),
|
||||||
|
patch.object(collaboration_service, "is_session_active", return_value=False),
|
||||||
|
):
|
||||||
|
result = collaboration_service.relay_collaboration_event("sid-2", payload)
|
||||||
|
|
||||||
|
assert result == ({"msg": "no_active_leader"}, 200)
|
||||||
|
repository.delete_leader.assert_called_once_with("wf-1")
|
||||||
|
socketio.emit.assert_not_called()
|
||||||
|
|
||||||
def test_relay_graph_event_unauthorized(self, service: tuple[WorkflowCollaborationService, Mock, Mock]) -> None:
|
def test_relay_graph_event_unauthorized(self, service: tuple[WorkflowCollaborationService, Mock, Mock]) -> None:
|
||||||
# Arrange
|
# Arrange
|
||||||
collaboration_service, repository, _socketio = service
|
collaboration_service, repository, _socketio = service
|
||||||
|
|||||||
Loading…
Reference in New Issue
Block a user