feat: complete ops ai workspace bridge

This commit is contained in:
DCCONSTRUCTIONS
2026-06-13 08:46:55 +03:00
parent 8b1e356494
commit 3806848b07
7 changed files with 2040 additions and 45 deletions
@@ -5,10 +5,14 @@
from django.urls import path
from plane.app.views import (
AIWorkspaceExecutorCheckEndpoint,
AIWorkspaceExecutorDetailEndpoint,
AIWorkspaceExecutorEventsEndpoint,
AIWorkspaceExecutorListEndpoint,
AIWorkspaceExecutorSelectEndpoint,
AIWorkspaceExecutorWindowsAgentEndpoint,
AIWorkspaceSettingsEndpoint,
AIWorkspaceThreadDispatchEndpoint,
AIWorkspaceThreadDetailEndpoint,
AIWorkspaceThreadListEndpoint,
AIWorkspaceThreadMessagesEndpoint,
@@ -16,6 +20,11 @@ from plane.app.views import (
urlpatterns = [
path(
"workspaces/<str:slug>/ai-workspace/settings/",
AIWorkspaceSettingsEndpoint.as_view(),
name="ai-workspace-settings",
),
path(
"workspaces/<str:slug>/ai-workspace/executors/",
AIWorkspaceExecutorListEndpoint.as_view(),
@@ -31,6 +40,16 @@ urlpatterns = [
AIWorkspaceExecutorSelectEndpoint.as_view(),
name="ai-workspace-executor-select",
),
path(
"workspaces/<str:slug>/ai-workspace/executors/<uuid:executor_id>/check/",
AIWorkspaceExecutorCheckEndpoint.as_view(),
name="ai-workspace-executor-check",
),
path(
"workspaces/<str:slug>/ai-workspace/executors/<uuid:executor_id>/events/",
AIWorkspaceExecutorEventsEndpoint.as_view(),
name="ai-workspace-executor-events",
),
path(
"workspaces/<str:slug>/ai-workspace/executors/<uuid:executor_id>/agent/windows.ps1",
AIWorkspaceExecutorWindowsAgentEndpoint.as_view(),
@@ -51,4 +70,9 @@ urlpatterns = [
AIWorkspaceThreadMessagesEndpoint.as_view(),
name="ai-workspace-thread-messages",
),
path(
"workspaces/<str:slug>/ai-workspace/threads/<uuid:thread_id>/dispatch/",
AIWorkspaceThreadDispatchEndpoint.as_view(),
name="ai-workspace-thread-dispatch",
),
]
@@ -175,10 +175,14 @@ from .module.archive import ModuleArchiveUnarchiveEndpoint
from .api import ApiTokenEndpoint
from .ai_workspace import (
AIWorkspaceExecutorCheckEndpoint,
AIWorkspaceExecutorDetailEndpoint,
AIWorkspaceExecutorEventsEndpoint,
AIWorkspaceExecutorListEndpoint,
AIWorkspaceExecutorSelectEndpoint,
AIWorkspaceExecutorWindowsAgentEndpoint,
AIWorkspaceSettingsEndpoint,
AIWorkspaceThreadDispatchEndpoint,
AIWorkspaceThreadDetailEndpoint,
AIWorkspaceThreadListEndpoint,
AIWorkspaceThreadMessagesEndpoint,
@@ -4,12 +4,31 @@ import os
# Third party imports
import requests
from django.http import HttpResponse
from django.utils import timezone
from rest_framework import status
from rest_framework.response import Response
# Module imports
from plane.app.permissions import ROLE, allow_permission
from plane.app.views.base import BaseAPIView
from plane.app.views.codex_agents import agent_path, get_gateway_config, owner_path, require_workspace, validate_project_in_workspace
AI_WORKSPACE_OPS_AGENT_DISPLAY_NAME = "AI Workspace Ops"
AI_WORKSPACE_OPS_MCP_SERVER_NAME = "nodedc_tasker"
AI_WORKSPACE_OPS_SCOPES = [
"workspace:read",
"project:read",
"project:member:add_existing",
"issue:read",
"issue:create",
"issue:update",
"issue:move",
"issue:comment",
"issue:label",
"issue:assign",
"issue:structured_blocks:write",
]
def get_ai_workspace_config():
@@ -49,12 +68,321 @@ def user_headers(user):
}
def ai_workspace_request(request, method, path, payload=None, query_params=None):
def ai_workspace_gateway_json(method, path, payload=None):
base_url, token, timeout = get_gateway_config()
if not base_url or not token:
return None, Response(
{
"ok": False,
"error": "codex_agent_gateway_not_configured",
"message": "NODE.DC Codex Agent Gateway URL/token is not configured.",
},
status=status.HTTP_503_SERVICE_UNAVAILABLE,
)
try:
response = requests.request(
method,
f"{base_url}{path}",
headers={
"Authorization": f"Bearer {token}",
"Accept": "application/json",
},
json=payload,
timeout=timeout,
)
except requests.RequestException:
return None, Response(
{
"ok": False,
"error": "codex_agent_gateway_unavailable",
"message": "NODE.DC Codex Agent Gateway is unavailable.",
},
status=status.HTTP_502_BAD_GATEWAY,
)
try:
data = response.json()
except ValueError:
data = {
"ok": False,
"error": "codex_agent_gateway_invalid_response",
"message": "NODE.DC Codex Agent Gateway returned a non-JSON response.",
}
if response.status_code >= 400:
return None, Response(data, status=response.status_code)
return data, None
def get_ops_agent(request, create=True):
owner = owner_path(request.user)
agents_payload, gateway_error = ai_workspace_gateway_json("GET", f"/api/internal/v1/owners/{owner}/agents")
if gateway_error is not None:
return None, gateway_error
agent = None
for candidate in agents_payload.get("agents") or []:
if (
candidate.get("display_name") == AI_WORKSPACE_OPS_AGENT_DISPLAY_NAME
and candidate.get("status", "active") == "active"
):
agent = candidate
break
if agent is None and create:
agent_payload, gateway_error = ai_workspace_gateway_json(
"POST",
f"/api/internal/v1/owners/{owner}/agents",
{
"display_name": AI_WORKSPACE_OPS_AGENT_DISPLAY_NAME,
"owner_email": request.user.email or None,
"avatar_url": None,
},
)
if gateway_error is not None:
return None, gateway_error
agent = agent_payload.get("agent") or {}
return agent, None
def replace_ops_agent_project_grants(request, agent_id, grants):
owner = owner_path(request.user)
_, gateway_error = ai_workspace_gateway_json(
"POST",
f"/api/internal/v1/owners/{owner}/agents/{agent_path(agent_id)}/grants/replace-projects",
{
"grants": grants,
"scopes": AI_WORKSPACE_OPS_SCOPES,
"mode": "voluntary",
},
)
if gateway_error is not None:
return None, gateway_error
return {"agent_id": agent_id, "grants": grants}, None
def ensure_ops_agent_project_grant(request, route_slug, ops_workspace_slug=None, ops_project_id=None):
ops_workspace_slug = (ops_workspace_slug or request.query_params.get("ops_workspace_slug") or route_slug or "").strip()
ops_project_id = (ops_project_id or request.query_params.get("ops_project_id") or "").strip()
if not ops_project_id:
return {}, None
workspace, workspace_error = require_workspace(ops_workspace_slug)
if workspace_error is not None:
return None, workspace_error
project_error = validate_project_in_workspace(workspace, ops_project_id, request.user)
if project_error is not None:
return None, project_error
agent, gateway_error = get_ops_agent(request, create=True)
if gateway_error is not None:
return None, gateway_error
agent_id = agent.get("id")
if not agent_id:
return None, Response(
{"ok": False, "error": "ai_workspace_ops_agent_missing"},
status=status.HTTP_502_BAD_GATEWAY,
)
grant_state, gateway_error = replace_ops_agent_project_grants(
request,
agent_id,
[{"workspace_slug": ops_workspace_slug, "project_id": ops_project_id}],
)
if gateway_error is not None:
return None, gateway_error
return {**grant_state, "ops_workspace_slug": ops_workspace_slug, "ops_project_id": ops_project_id}, None
def clear_ops_agent_project_grants(request):
agent, gateway_error = get_ops_agent(request, create=False)
if gateway_error is not None:
return None, gateway_error
if not agent:
return {}, None
agent_id = agent.get("id")
if not agent_id:
return {}, None
return replace_ops_agent_project_grants(request, agent_id, [])
def create_ops_mcp_token_for_request(request, agent_id):
owner = owner_path(request.user)
token_payload, gateway_error = ai_workspace_gateway_json(
"POST",
f"/api/internal/v1/owners/{owner}/agents/{agent_path(agent_id)}/tokens",
{"name": "AI Workspace Bridge"},
)
if gateway_error is not None:
return None, gateway_error
mcp_server = ((token_payload.get("setup") or {}).get("mcp_server") or {})
mcp_url = mcp_server.get("url")
mcp_token = token_payload.get("token")
if not mcp_url or not mcp_token:
return None, Response(
{"ok": False, "error": "ai_workspace_ops_mcp_setup_missing"},
status=status.HTTP_502_BAD_GATEWAY,
)
return {
"opsMcpUrl": mcp_url,
"opsMcpToken": mcp_token,
"opsMcpServerName": mcp_server.get("name") or AI_WORKSPACE_OPS_MCP_SERVER_NAME,
}, None
def get_ai_workspace_settings_data(request):
response = ai_workspace_request(request, "GET", "/settings")
if response.status_code >= 400:
return None, response
data = getattr(response, "data", None)
if not isinstance(data, dict):
return None, Response(
{"ok": False, "error": "ai_workspace_assistant_invalid_response"},
status=status.HTTP_502_BAD_GATEWAY,
)
return data, None
def settings_has_ops_mcp_server(settings_data):
if not isinstance(settings_data, dict):
return False
settings = settings_data.get("settings") or {}
if not isinstance(settings, dict):
return False
metadata = settings.get("metadata") or {}
if not isinstance(metadata, dict):
return False
def is_ops_server(server):
if not isinstance(server, dict):
return False
server_name = (server.get("serverName") or server.get("server_name") or server.get("name") or "").strip()
app_id = (server.get("appId") or server.get("app_id") or "").strip().lower()
return bool(server.get("url")) and (app_id == "ops" or server_name == AI_WORKSPACE_OPS_MCP_SERVER_NAME)
def has_ops_server(value):
if isinstance(value, list):
return any(is_ops_server(item) for item in value)
if isinstance(value, dict):
return any(is_ops_server(item) for item in value.values())
return False
app_grants = metadata.get("appGrants") or {}
ops_grant = app_grants.get("ops") if isinstance(app_grants, dict) else {}
if isinstance(ops_grant, dict) and has_ops_server(ops_grant.get("mcpServers")):
return True
return has_ops_server(metadata.get("mcpServers"))
def build_ops_app_grant(active_context, ops_mcp_params=None):
context = active_context if isinstance(active_context, dict) else {}
grant = {
"appId": "ops",
"appTitle": "NODE.DC Ops",
"surface": "ops",
"updatedAt": timezone.now().isoformat(),
"context": {
"opsWorkspaceSlug": context.get("opsWorkspaceSlug") or context.get("workspaceSlug") or "",
"opsRouteWorkspaceSlug": context.get("opsRouteWorkspaceSlug") or "",
"opsProjectId": context.get("opsProjectId") or "",
"opsProjectIdentifier": context.get("opsProjectIdentifier") or "",
"opsProjectSlug": context.get("opsProjectSlug") or "",
"opsProjectTitle": context.get("opsProjectTitle") or "",
},
}
if not ops_mcp_params:
return grant
mcp_url = ops_mcp_params.get("opsMcpUrl")
mcp_token = ops_mcp_params.get("opsMcpToken")
if not mcp_url or not mcp_token:
return grant
grant["mcpServers"] = [
{
"serverName": ops_mcp_params.get("opsMcpServerName") or AI_WORKSPACE_OPS_MCP_SERVER_NAME,
"url": mcp_url,
"enabled": True,
"required": False,
"startupTimeoutSec": 20,
"toolTimeoutSec": 60,
"httpHeaders": {
"Authorization": f"Bearer {mcp_token}",
"Accept": "application/json",
"MCP-Protocol-Version": "2025-06-18",
},
}
]
return grant
def attach_ops_grant_for_active_context(request, slug, payload, settings_data=None):
active_context = payload.get("activeContext") or {}
if not isinstance(active_context, dict):
active_context = {}
ops_project_id = (active_context.get("opsProjectId") or "").strip()
if settings_data is None:
settings_data, settings_error = get_ai_workspace_settings_data(request)
if settings_error is not None:
return None, settings_error
if not ops_project_id:
if settings_has_ops_mcp_server(settings_data):
_, clear_error = clear_ops_agent_project_grants(request)
if clear_error is not None:
return None, clear_error
return attach_ops_app_grant(payload, build_ops_app_grant(active_context)), None
grant_state, grant_error = ensure_ops_agent_project_grant(
request,
slug,
ops_workspace_slug=active_context.get("opsWorkspaceSlug") or active_context.get("workspaceSlug") or slug,
ops_project_id=ops_project_id,
)
if grant_error is not None:
return None, grant_error
ops_mcp_params = None
agent_id = (grant_state or {}).get("agent_id")
if agent_id and not settings_has_ops_mcp_server(settings_data):
ops_mcp_params, ops_mcp_error = create_ops_mcp_token_for_request(request, agent_id)
if ops_mcp_error is not None:
return None, ops_mcp_error
return attach_ops_app_grant(payload, build_ops_app_grant(active_context, ops_mcp_params)), None
def attach_ops_app_grant(payload, app_grant):
if not app_grant:
return payload
metadata = payload.get("metadata") or {}
if not isinstance(metadata, dict):
metadata = {}
app_grants = metadata.get("appGrants") or {}
if not isinstance(app_grants, dict):
app_grants = {}
metadata["appGrants"] = {
**app_grants,
"ops": app_grant,
}
payload["metadata"] = metadata
return payload
def ai_workspace_request(request, method, path, payload=None, query_params=None, timeout_override=None):
config, error_response = require_ai_workspace_config()
if error_response is not None:
return error_response
base_url, token, timeout = config
request_timeout = timeout_override or timeout
try:
response = requests.request(
method,
@@ -66,7 +394,7 @@ def ai_workspace_request(request, method, path, payload=None, query_params=None)
},
params=query_params,
json=payload,
timeout=timeout,
timeout=request_timeout,
)
except requests.RequestException:
return Response(
@@ -95,10 +423,30 @@ def ops_thread_payload(raw_payload, slug):
active_context = payload.get("activeContext") or payload.get("active_context") or {}
if not isinstance(active_context, dict):
active_context = {}
ops_workspace_slug = active_context.get("opsWorkspaceSlug") or active_context.get("workspaceSlug") or slug
active_context = {
**active_context,
"surface": "ops",
"workspaceSlug": slug,
"workspaceSlug": ops_workspace_slug,
"opsWorkspaceSlug": ops_workspace_slug,
"opsRouteWorkspaceSlug": slug,
}
contexts = active_context.get("contexts")
if not isinstance(contexts, dict):
contexts = {}
ops_context = contexts.get("ops")
if not isinstance(ops_context, dict):
ops_context = {}
active_context["contexts"] = {
**contexts,
"ops": {
**ops_context,
**active_context,
"surface": "ops",
"workspaceSlug": ops_workspace_slug,
"opsWorkspaceSlug": ops_workspace_slug,
"opsRouteWorkspaceSlug": slug,
},
}
tool_packs = payload.get("enabledToolPacks") or payload.get("enabled_tool_packs") or []
@@ -113,6 +461,64 @@ def ops_thread_payload(raw_payload, slug):
return payload
def ops_settings_payload(raw_payload, slug):
payload = dict(raw_payload or {})
active_context = payload.get("activeContext") or payload.get("active_context") or {}
if not isinstance(active_context, dict):
active_context = {}
ops_workspace_slug = active_context.get("opsWorkspaceSlug") or active_context.get("workspaceSlug") or slug
active_context = {
**active_context,
"surface": "ops",
"workspaceSlug": ops_workspace_slug,
"opsWorkspaceSlug": ops_workspace_slug,
"opsRouteWorkspaceSlug": slug,
}
contexts = active_context.get("contexts")
if not isinstance(contexts, dict):
contexts = {}
ops_context = contexts.get("ops")
if not isinstance(ops_context, dict):
ops_context = {}
active_context["contexts"] = {
**contexts,
"ops": {
**ops_context,
**active_context,
"surface": "ops",
"workspaceSlug": ops_workspace_slug,
"opsWorkspaceSlug": ops_workspace_slug,
"opsRouteWorkspaceSlug": slug,
},
}
tool_packs = payload.get("enabledToolPacks") or payload.get("enabled_tool_packs") or []
if not isinstance(tool_packs, list):
tool_packs = []
for tool_pack in ("ops", "engine"):
if tool_pack not in tool_packs:
tool_packs.append(tool_pack)
payload["activeContext"] = active_context
payload["enabledToolPacks"] = tool_packs
return payload
class AIWorkspaceSettingsEndpoint(BaseAPIView):
@allow_permission(allowed_roles=[ROLE.ADMIN, ROLE.MEMBER], level="WORKSPACE")
def get(self, request, slug):
return ai_workspace_request(request, "GET", "/settings")
@allow_permission(allowed_roles=[ROLE.ADMIN, ROLE.MEMBER], level="WORKSPACE")
def patch(self, request, slug):
payload = ops_settings_payload(request.data, slug)
payload, ops_grant_error = attach_ops_grant_for_active_context(request, slug, payload)
if ops_grant_error is not None:
return ops_grant_error
return ai_workspace_request(request, "PATCH", "/settings", payload)
class AIWorkspaceExecutorListEndpoint(BaseAPIView):
@allow_permission(allowed_roles=[ROLE.ADMIN, ROLE.MEMBER], level="WORKSPACE")
def get(self, request, slug):
@@ -139,6 +545,34 @@ class AIWorkspaceExecutorSelectEndpoint(BaseAPIView):
return ai_workspace_request(request, "POST", f"/executors/{executor_id}/select")
class AIWorkspaceExecutorCheckEndpoint(BaseAPIView):
@allow_permission(allowed_roles=[ROLE.ADMIN, ROLE.MEMBER], level="WORKSPACE")
def post(self, request, slug, executor_id):
return ai_workspace_request(
request,
"POST",
f"/executors/{executor_id}/check",
request.data,
timeout_override=15,
)
class AIWorkspaceExecutorEventsEndpoint(BaseAPIView):
@allow_permission(allowed_roles=[ROLE.ADMIN, ROLE.MEMBER], level="WORKSPACE")
def get(self, request, slug, executor_id):
query_params = {
"since": request.query_params.get("since") or "0",
"limit": request.query_params.get("limit") or "100",
}
return ai_workspace_request(
request,
"GET",
f"/executors/{executor_id}/events",
query_params=query_params,
timeout_override=15,
)
class AIWorkspaceExecutorWindowsAgentEndpoint(BaseAPIView):
@allow_permission(allowed_roles=[ROLE.ADMIN, ROLE.MEMBER], level="WORKSPACE")
def get(self, request, slug, executor_id):
@@ -150,6 +584,34 @@ class AIWorkspaceExecutorWindowsAgentEndpoint(BaseAPIView):
params = {}
if request.query_params.get("port"):
params["port"] = request.query_params.get("port")
ops_project_id = (request.query_params.get("ops_project_id") or "").strip()
if ops_project_id:
ops_workspace_slug = (request.query_params.get("ops_workspace_slug") or slug or "").strip()
active_context = {
"surface": "ops",
"workspaceSlug": ops_workspace_slug,
"opsWorkspaceSlug": ops_workspace_slug,
"opsRouteWorkspaceSlug": slug,
"opsProjectId": ops_project_id,
}
sync_payload = ops_settings_payload(
{
"activeContext": active_context,
"enabledToolPacks": ["ops", "engine"],
"metadata": {
"source": "ops-ai-workspace-installer",
"updatedAt": timezone.now().isoformat(),
},
},
slug,
)
sync_payload, ops_grant_error = attach_ops_grant_for_active_context(request, slug, sync_payload)
if ops_grant_error is not None:
return ops_grant_error
sync_response = ai_workspace_request(request, "PATCH", "/settings", sync_payload)
if sync_response.status_code >= 400:
return sync_response
try:
response = requests.get(
f"{base_url}/api/ai-workspace/assistant/v1/executors/{executor_id}/agent/windows.ps1",
@@ -185,9 +647,11 @@ class AIWorkspaceThreadListEndpoint(BaseAPIView):
@allow_permission(allowed_roles=[ROLE.ADMIN, ROLE.MEMBER], level="WORKSPACE")
def get(self, request, slug):
query_params = {
"surface": request.query_params.get("surface") or "ops",
"kind": request.query_params.get("kind") or "shared",
"limit": request.query_params.get("limit") or "100",
}
if request.query_params.get("surface"):
query_params["surface"] = request.query_params.get("surface")
return ai_workspace_request(request, "GET", "/threads", query_params=query_params)
@allow_permission(allowed_roles=[ROLE.ADMIN, ROLE.MEMBER], level="WORKSPACE")
@@ -218,3 +682,27 @@ class AIWorkspaceThreadMessagesEndpoint(BaseAPIView):
@allow_permission(allowed_roles=[ROLE.ADMIN, ROLE.MEMBER], level="WORKSPACE")
def post(self, request, slug, thread_id):
return ai_workspace_request(request, "POST", f"/threads/{thread_id}/messages", request.data)
class AIWorkspaceThreadDispatchEndpoint(BaseAPIView):
@allow_permission(allowed_roles=[ROLE.ADMIN, ROLE.MEMBER], level="WORKSPACE")
def post(self, request, slug, thread_id):
payload = dict(request.data or {})
context = payload.get("context") or {}
if not isinstance(context, dict):
context = {}
ops_workspace_slug = context.get("opsWorkspaceSlug") or context.get("workspaceSlug") or slug
payload["context"] = {
**context,
"surface": "ops",
"workspaceSlug": ops_workspace_slug,
"opsWorkspaceSlug": ops_workspace_slug,
"opsRouteWorkspaceSlug": slug,
}
return ai_workspace_request(
request,
"POST",
f"/threads/{thread_id}/dispatch",
payload,
timeout_override=20,
)