210 lines
8.0 KiB
Python
210 lines
8.0 KiB
Python
#!/usr/bin/env python
|
|
"""Verify the AI Agent Admin runtime through public HTTP APIs."""
|
|
from __future__ import annotations
|
|
|
|
import argparse
|
|
import asyncio
|
|
import json
|
|
import sys
|
|
from typing import Any
|
|
from urllib import error, request
|
|
|
|
|
|
def _preview(value: Any, limit: int = 220) -> str:
|
|
if value is None:
|
|
return ""
|
|
text = value if isinstance(value, str) else json.dumps(value, ensure_ascii=False, default=str)
|
|
return text if len(text) <= limit else f"{text[:limit]}..."
|
|
|
|
|
|
def _pick_token(payload: dict[str, Any]) -> str:
|
|
for key in ("accessToken", "access_token", "token"):
|
|
token = payload.get(key)
|
|
if token:
|
|
return str(token)
|
|
data = payload.get("data") or {}
|
|
for key in ("accessToken", "access_token", "token"):
|
|
token = data.get(key)
|
|
if token:
|
|
return str(token)
|
|
raise RuntimeError("login succeeded but token was not found")
|
|
|
|
|
|
class ApiClient:
|
|
def __init__(self, base_url: str, timeout: float) -> None:
|
|
self.base_url = base_url.rstrip("/")
|
|
self.timeout = timeout
|
|
self.token = ""
|
|
|
|
def _request(self, method: str, path: str, payload: dict[str, Any] | None = None) -> Any:
|
|
body = None
|
|
headers = {"Accept": "application/json"}
|
|
if payload is not None:
|
|
body = json.dumps(payload, ensure_ascii=False).encode("utf-8")
|
|
headers["Content-Type"] = "application/json"
|
|
if self.token:
|
|
headers["Authorization"] = f"Bearer {self.token}"
|
|
|
|
req = request.Request(
|
|
f"{self.base_url}{path}",
|
|
data=body,
|
|
headers=headers,
|
|
method=method,
|
|
)
|
|
try:
|
|
with request.urlopen(req, timeout=self.timeout) as resp:
|
|
content_type = resp.headers.get("Content-Type", "")
|
|
raw = resp.read().decode("utf-8", errors="replace")
|
|
except error.HTTPError as exc:
|
|
detail = exc.read().decode("utf-8", errors="replace")
|
|
raise RuntimeError(f"{method} {path} failed: HTTP {exc.code} {detail}") from exc
|
|
if "application/json" in content_type:
|
|
return json.loads(raw)
|
|
return raw
|
|
|
|
async def get(self, path: str) -> Any:
|
|
return await asyncio.to_thread(self._request, "GET", path, None)
|
|
|
|
async def post(self, path: str, payload: dict[str, Any] | None = None) -> Any:
|
|
return await asyncio.to_thread(self._request, "POST", path, payload or {})
|
|
|
|
|
|
def _unwrap(payload: Any) -> Any:
|
|
if isinstance(payload, dict) and "data" in payload:
|
|
return payload["data"]
|
|
return payload
|
|
|
|
|
|
async def _find_completed_issue_run(
|
|
client: ApiClient,
|
|
workflow_id: str,
|
|
issue: str,
|
|
) -> dict[str, Any]:
|
|
runs_data = _unwrap(await client.get("/api/ai/workflow/runs?page=1&pageSize=50"))
|
|
runs = runs_data.get("items", runs_data) if isinstance(runs_data, dict) else runs_data
|
|
candidates = [
|
|
run
|
|
for run in runs
|
|
if run.get("workflow_id") == workflow_id and run.get("status") == "completed"
|
|
]
|
|
issue_text = issue.lower()
|
|
for run in candidates:
|
|
run_detail = _unwrap(await client.get(f"/api/ai/workflow/runs/{run['id']}"))
|
|
searchable = " ".join(
|
|
_preview(run_detail.get(key), 2000)
|
|
for key in ("inputs", "outputs", "execution_log")
|
|
).lower()
|
|
if issue_text in searchable:
|
|
return run_detail
|
|
raise RuntimeError(f"completed run not found for issue: {issue}")
|
|
|
|
|
|
async def verify(args: argparse.Namespace) -> int:
|
|
base_url = args.base_url.rstrip("/")
|
|
client = ApiClient(base_url=base_url, timeout=args.timeout)
|
|
login_payload = await client.post(
|
|
"/api/core/login",
|
|
{"username": args.username, "password": args.password},
|
|
)
|
|
client.token = _pick_token(login_payload)
|
|
|
|
providers_data = _unwrap(await client.get("/api/ai/provider/list?page=1&pageSize=100"))
|
|
providers = providers_data.get("items", providers_data) if isinstance(providers_data, dict) else providers_data
|
|
provider = next(
|
|
(item for item in providers if item.get("name") == args.provider or item.get("code") == args.provider),
|
|
None,
|
|
)
|
|
if not provider:
|
|
raise RuntimeError(f"provider not found: {args.provider}")
|
|
|
|
provider_detail = _unwrap(await client.get(f"/api/ai/provider/{provider['id']}"))
|
|
provider_test = _unwrap(await client.post(f"/api/ai/provider/{provider['id']}/test"))
|
|
if not provider_test.get("success"):
|
|
raise RuntimeError(f"provider test failed: {provider_test}")
|
|
|
|
agent = _unwrap(await client.get(f"/api/ai/agent/code/{args.agent_code}"))
|
|
if agent.get("status") != "published":
|
|
raise RuntimeError(f"agent is not published: {agent.get('status')}")
|
|
|
|
chat_text = await client.post(
|
|
f"/api/ai/agent/{agent['id']}/chat",
|
|
{
|
|
"message": "用一句话回复:AI Agent Admin 验证脚本连通性测试",
|
|
"session_id": "verify-ai-agent-admin",
|
|
},
|
|
)
|
|
if "llm_chunk" not in chat_text and "content" not in chat_text:
|
|
raise RuntimeError("agent chat stream did not return content")
|
|
|
|
workflow = _unwrap(await client.get(f"/api/ai/workflow/code/{args.workflow_code}"))
|
|
if workflow.get("status") != "published":
|
|
raise RuntimeError(f"workflow is not published: {workflow.get('status')}")
|
|
|
|
run_detail = await _find_completed_issue_run(client, workflow.get("id"), args.issue)
|
|
execution_log = run_detail.get("execution_log") or []
|
|
if len(execution_log) < args.min_steps:
|
|
raise RuntimeError(f"execution log too short: {len(execution_log)} < {args.min_steps}")
|
|
|
|
agent_codes = sorted(
|
|
{
|
|
log.get("metadata", {}).get("agent_code")
|
|
for log in execution_log
|
|
if isinstance(log, dict) and log.get("metadata", {}).get("agent_code")
|
|
}
|
|
)
|
|
if len(agent_codes) < args.min_agents:
|
|
raise RuntimeError(f"not enough collaborating agents: {agent_codes}")
|
|
|
|
report = {
|
|
"base_url": base_url,
|
|
"provider": {
|
|
"id": provider_detail.get("id"),
|
|
"name": provider_detail.get("name"),
|
|
"api_base": provider_detail.get("api_base"),
|
|
"api_key_masked": provider_detail.get("api_key_masked"),
|
|
"model_count": provider_test.get("model_count"),
|
|
},
|
|
"agent": {
|
|
"id": agent.get("id"),
|
|
"code": agent.get("code"),
|
|
"status": agent.get("status"),
|
|
},
|
|
"workflow": {
|
|
"id": workflow.get("id"),
|
|
"code": workflow.get("code"),
|
|
"status": workflow.get("status"),
|
|
"run_id": run_detail.get("id"),
|
|
"run_status": run_detail.get("status"),
|
|
"total_steps": run_detail.get("total_steps"),
|
|
"total_tokens": run_detail.get("total_tokens"),
|
|
"execution_log_count": len(execution_log),
|
|
"agent_codes": agent_codes,
|
|
"final_output_preview": _preview(run_detail.get("outputs"), 320),
|
|
},
|
|
}
|
|
print(json.dumps(report, ensure_ascii=False, indent=2, default=str))
|
|
return 0
|
|
|
|
|
|
def parse_args() -> argparse.Namespace:
|
|
parser = argparse.ArgumentParser(description="Verify AI Agent Admin runtime APIs.")
|
|
parser.add_argument("--base-url", default="http://127.0.0.1:18000")
|
|
parser.add_argument("--username", default="superadmin")
|
|
parser.add_argument("--password", default="123456")
|
|
parser.add_argument("--provider", default="codex")
|
|
parser.add_argument("--agent-code", default="codex")
|
|
parser.add_argument("--workflow-code", default="multica_org_collaboration_flow")
|
|
parser.add_argument("--issue", default="OC-69")
|
|
parser.add_argument("--min-steps", type=int, default=8)
|
|
parser.add_argument("--min-agents", type=int, default=3)
|
|
parser.add_argument("--timeout", type=float, default=60)
|
|
return parser.parse_args()
|
|
|
|
|
|
if __name__ == "__main__":
|
|
try:
|
|
raise SystemExit(asyncio.run(verify(parse_args())))
|
|
except Exception as exc:
|
|
print(f"verify failed: {exc}", file=sys.stderr)
|
|
raise SystemExit(1)
|