430 lines
16 KiB
Python
430 lines
16 KiB
Python
from __future__ import annotations
|
||
|
||
from uuid import uuid4
|
||
|
||
from fastapi import APIRouter
|
||
|
||
from app.schemas.common import ApiResponse
|
||
from app.services.detect_job_service import (
|
||
append_detect_job_event,
|
||
create_detect_job_if_needed,
|
||
get_active_detect_job_summary,
|
||
get_detect_job_summary,
|
||
get_detect_queue_health,
|
||
list_detect_jobs,
|
||
)
|
||
from app.services.detect_service import get_detect_status
|
||
from app.services.detect_run_service import create_detect_run_snapshot, finalize_detect_run, mark_detect_run_stopping
|
||
from app.services.ops_job_service import create_ops_job, list_managed_nodes
|
||
from app.services.settings_service import get_settings_payload, resolve_thread_count
|
||
from app.services.worker_control_service import send_worker_command, start_worker
|
||
|
||
router = APIRouter(tags=["detect"])
|
||
|
||
|
||
def _build_detect_action_result(
|
||
*,
|
||
action: str,
|
||
ok: bool,
|
||
message: str,
|
||
poll_after_seconds: int = 2,
|
||
refresh_status: bool = True,
|
||
data: dict | None = None,
|
||
) -> dict:
|
||
result = {
|
||
"action": action,
|
||
"poll_after_seconds": poll_after_seconds,
|
||
"refresh_status": refresh_status,
|
||
**(data or {}),
|
||
}
|
||
lowered = str(message or "").strip().lower()
|
||
if "ui_level" not in result:
|
||
if not ok:
|
||
result["ui_level"] = "warning" if any(keyword in lowered for keyword in ("当前没有", "无需", "未运行", "未启动")) else "error"
|
||
elif any(keyword in str(message or "") for keyword in ("命令已发送", "已发送检测启动请求", "控制指令")):
|
||
result["ui_level"] = "warning"
|
||
else:
|
||
result["ui_level"] = "success"
|
||
if "poll_schedule_seconds" not in result:
|
||
if ok:
|
||
result["poll_schedule_seconds"] = [1, max(2, poll_after_seconds)]
|
||
elif poll_after_seconds > 0:
|
||
result["poll_schedule_seconds"] = [poll_after_seconds]
|
||
else:
|
||
result["poll_schedule_seconds"] = []
|
||
return result
|
||
|
||
|
||
def _build_settings_summary(settings_payload: dict) -> dict:
|
||
thread_count_resolution = resolve_thread_count(settings_payload=settings_payload)
|
||
return {
|
||
"thread_count": int(thread_count_resolution["effective_thread_count"]),
|
||
"thread_count_default": int(thread_count_resolution["default_thread_count"]),
|
||
"thread_count_source": str(thread_count_resolution["source"]),
|
||
"thread_count_override": thread_count_resolution["override_thread_count"],
|
||
"thread_count_node_code": str(thread_count_resolution["node_code"]),
|
||
"proxy_enable": settings_payload["proxy_config"].get("proxy_enable", False),
|
||
"allow_direct": settings_payload["proxy_config"].get("allow_direct", False),
|
||
"proxy_pool_count": len(settings_payload["proxy_config"].get("proxy_urls", [])),
|
||
}
|
||
|
||
|
||
def _mainland_detect_targets() -> dict[str, list[dict]]:
|
||
controllers: list[dict] = []
|
||
workers: list[dict] = []
|
||
for node in list_managed_nodes():
|
||
if not bool(node.get("is_enabled", True)):
|
||
continue
|
||
if str(node.get("region") or "").strip() != "mainland":
|
||
continue
|
||
if not str(node.get("last_seen_at") or "").strip():
|
||
continue
|
||
role = str(node.get("role") or "").strip()
|
||
if role == "control":
|
||
controllers.append(node)
|
||
elif role == "worker":
|
||
workers.append(node)
|
||
return {"controllers": controllers, "workers": workers}
|
||
|
||
|
||
def _queue_remote_detect_job(
|
||
*,
|
||
node_code: str,
|
||
action: str,
|
||
job_summary: dict,
|
||
cycle_token: str,
|
||
requested_by: str = "api",
|
||
payload: dict | None = None,
|
||
) -> dict:
|
||
ok, message, data = create_ops_job(
|
||
{
|
||
"action": action,
|
||
"target_node_code": node_code,
|
||
"execution_mode": "remote-agent",
|
||
"requested_by": requested_by,
|
||
"auto_approve": True,
|
||
"run_now": False,
|
||
"payload": {
|
||
"job_id": int(job_summary.get("job_id") or 0),
|
||
"job_code": str(job_summary.get("job_code") or "").strip(),
|
||
"cycle_token": cycle_token,
|
||
**dict(payload or {}),
|
||
},
|
||
"metadata": {
|
||
"source": "detect.start",
|
||
"job_id": int(job_summary.get("job_id") or 0),
|
||
"job_code": str(job_summary.get("job_code") or "").strip(),
|
||
"cycle_token": cycle_token,
|
||
},
|
||
}
|
||
)
|
||
return {
|
||
"node_code": node_code,
|
||
"action": action,
|
||
"ok": ok,
|
||
"message": message,
|
||
"job": dict((data or {}).get("job") or {}),
|
||
}
|
||
|
||
|
||
def _dispatch_remote_detect_start(*, job_summary: dict, cycle_token: str) -> dict:
|
||
targets = _mainland_detect_targets()
|
||
queued: list[dict] = []
|
||
|
||
for node in targets["controllers"]:
|
||
node_code = str(node.get("node_code") or "").strip()
|
||
if not node_code:
|
||
continue
|
||
queued.append(
|
||
_queue_remote_detect_job(
|
||
node_code=node_code,
|
||
action="runtime.start_sync_agent",
|
||
job_summary=job_summary,
|
||
cycle_token=cycle_token,
|
||
)
|
||
)
|
||
queued.append(
|
||
_queue_remote_detect_job(
|
||
node_code=node_code,
|
||
action="runtime.pull_tasks",
|
||
job_summary=job_summary,
|
||
cycle_token=cycle_token,
|
||
payload={"limit": int(job_summary.get("items_pending") or job_summary.get("items_total") or 0) or 1000},
|
||
)
|
||
)
|
||
queued.append(
|
||
_queue_remote_detect_job(
|
||
node_code=node_code,
|
||
action="runtime.start_worker",
|
||
job_summary=job_summary,
|
||
cycle_token=cycle_token,
|
||
)
|
||
)
|
||
queued.append(
|
||
_queue_remote_detect_job(
|
||
node_code=node_code,
|
||
action="runtime.start_detection",
|
||
job_summary=job_summary,
|
||
cycle_token=cycle_token,
|
||
)
|
||
)
|
||
|
||
for node in targets["workers"]:
|
||
node_code = str(node.get("node_code") or "").strip()
|
||
if not node_code:
|
||
continue
|
||
queued.append(
|
||
_queue_remote_detect_job(
|
||
node_code=node_code,
|
||
action="runtime.start_worker",
|
||
job_summary=job_summary,
|
||
cycle_token=cycle_token,
|
||
)
|
||
)
|
||
queued.append(
|
||
_queue_remote_detect_job(
|
||
node_code=node_code,
|
||
action="runtime.start_detection",
|
||
job_summary=job_summary,
|
||
cycle_token=cycle_token,
|
||
)
|
||
)
|
||
|
||
success_jobs = [item for item in queued if item.get("ok")]
|
||
failed_jobs = [item for item in queued if not item.get("ok")]
|
||
return {
|
||
"target_summary": {
|
||
"controller_nodes": [str(item.get("node_code") or "") for item in targets["controllers"]],
|
||
"worker_nodes": [str(item.get("node_code") or "") for item in targets["workers"]],
|
||
},
|
||
"queued_jobs": queued,
|
||
"queued_total": len(success_jobs),
|
||
"failed_total": len(failed_jobs),
|
||
}
|
||
|
||
|
||
def _dispatch_remote_detect_stop(*, active_job: dict | None, cycle_token: str = "") -> dict:
|
||
targets = _mainland_detect_targets()
|
||
job_summary = active_job or {}
|
||
queued: list[dict] = []
|
||
|
||
for node in [*targets["controllers"], *targets["workers"]]:
|
||
node_code = str(node.get("node_code") or "").strip()
|
||
if not node_code:
|
||
continue
|
||
queued.append(
|
||
_queue_remote_detect_job(
|
||
node_code=node_code,
|
||
action="runtime.stop_detection",
|
||
job_summary=job_summary,
|
||
cycle_token=cycle_token,
|
||
payload={},
|
||
)
|
||
)
|
||
|
||
success_jobs = [item for item in queued if item.get("ok")]
|
||
failed_jobs = [item for item in queued if not item.get("ok")]
|
||
return {
|
||
"target_summary": {
|
||
"controller_nodes": [str(item.get("node_code") or "") for item in targets["controllers"]],
|
||
"worker_nodes": [str(item.get("node_code") or "") for item in targets["workers"]],
|
||
},
|
||
"queued_jobs": queued,
|
||
"queued_total": len(success_jobs),
|
||
"failed_total": len(failed_jobs),
|
||
}
|
||
|
||
|
||
@router.get("/detect/status", response_model=ApiResponse)
|
||
def detect_status() -> ApiResponse:
|
||
return ApiResponse(data=get_detect_status())
|
||
|
||
|
||
@router.get("/detect/job/active", response_model=ApiResponse)
|
||
def detect_active_job() -> ApiResponse:
|
||
return ApiResponse(data=get_active_detect_job_summary(event_limit=50))
|
||
|
||
|
||
@router.get("/detect/jobs", response_model=ApiResponse)
|
||
def detect_jobs(limit: int = 20) -> ApiResponse:
|
||
return ApiResponse(data={"items": list_detect_jobs(limit=limit), "limit": max(1, min(int(limit or 20), 100))})
|
||
|
||
|
||
@router.get("/detect/queue-summary", response_model=ApiResponse)
|
||
def detect_queue_summary(window_minutes: int = 15) -> ApiResponse:
|
||
return ApiResponse(data=get_detect_queue_health(window_minutes=window_minutes))
|
||
|
||
|
||
@router.get("/detect/jobs/{job_id}", response_model=ApiResponse)
|
||
def detect_job_detail(job_id: int) -> ApiResponse:
|
||
data = get_detect_job_summary(job_id, event_limit=100)
|
||
if not data:
|
||
return ApiResponse(code=1, message="检测任务不存在", data=None)
|
||
return ApiResponse(data=data)
|
||
|
||
|
||
@router.post("/detect/start", response_model=ApiResponse)
|
||
def start_detect() -> ApiResponse:
|
||
job_summary = create_detect_job_if_needed(limit=1000, created_by="api")
|
||
if not job_summary:
|
||
result = _build_detect_action_result(
|
||
action="start",
|
||
ok=False,
|
||
message="当前没有可创建的检测任务",
|
||
data={"job": None},
|
||
)
|
||
return ApiResponse(
|
||
code=0,
|
||
message="当前没有可创建的检测任务",
|
||
data=result,
|
||
)
|
||
cycle_token = uuid4().hex[:10]
|
||
append_detect_job_event(
|
||
job_summary["job_id"],
|
||
event_type="job_dispatch_requested",
|
||
message="控制面已发送检测启动请求",
|
||
payload={
|
||
"cycle_token": cycle_token,
|
||
"status": job_summary.get("status"),
|
||
"items_pending": job_summary.get("items_pending", 0),
|
||
"items_claimed": job_summary.get("items_claimed", 0),
|
||
"items_running": job_summary.get("items_running", 0),
|
||
},
|
||
)
|
||
|
||
ok, message = start_worker()
|
||
if not ok:
|
||
result = _build_detect_action_result(
|
||
action="start",
|
||
ok=False,
|
||
message=message,
|
||
data={"job": job_summary},
|
||
)
|
||
append_detect_job_event(
|
||
job_summary["job_id"],
|
||
event_type="job_dispatch_failed",
|
||
level="error",
|
||
message=f"启动 Worker 失败: {message}",
|
||
payload={"cycle_token": cycle_token},
|
||
)
|
||
return ApiResponse(
|
||
code=1,
|
||
message=message,
|
||
data=result,
|
||
)
|
||
|
||
command_ok, command_message = send_worker_command(
|
||
"start_detection",
|
||
payload={
|
||
"cycle_token": cycle_token,
|
||
"job_id": job_summary["job_id"],
|
||
"job_code": job_summary["job_code"],
|
||
},
|
||
)
|
||
append_detect_job_event(
|
||
job_summary["job_id"],
|
||
event_type="job_dispatch_sent" if command_ok else "job_dispatch_rejected",
|
||
level="info" if command_ok else "error",
|
||
message=command_message,
|
||
payload={"cycle_token": cycle_token},
|
||
)
|
||
snapshot = get_detect_status()
|
||
settings_payload = get_settings_payload()
|
||
settings_summary = _build_settings_summary(settings_payload)
|
||
if command_ok:
|
||
remote_dispatch = _dispatch_remote_detect_start(job_summary=job_summary, cycle_token=cycle_token)
|
||
create_detect_run_snapshot(
|
||
message=f"{message};{command_message}",
|
||
runtime={
|
||
"mode": snapshot.get("worker_mode", ""),
|
||
"running": snapshot.get("worker_online", False),
|
||
"process_count": snapshot.get("worker_process_count", 0),
|
||
"latest_start_time": snapshot.get("worker_latest_start_time", ""),
|
||
"message": snapshot.get("worker_runtime_message", ""),
|
||
},
|
||
progress=snapshot.get("progress", {}),
|
||
settings_summary=settings_summary,
|
||
)
|
||
append_detect_job_event(
|
||
job_summary["job_id"],
|
||
event_type="job_dispatch_remote_queued",
|
||
level="info" if int(remote_dispatch.get("failed_total", 0) or 0) == 0 else "warning",
|
||
message=(
|
||
f"已向大陆节点排队 {int(remote_dispatch.get('queued_total', 0) or 0)} 个远端检测动作"
|
||
if int(remote_dispatch.get("queued_total", 0) or 0) > 0
|
||
else "当前没有可排队的大陆远端检测动作"
|
||
),
|
||
payload={
|
||
"cycle_token": cycle_token,
|
||
"remote_dispatch": remote_dispatch,
|
||
},
|
||
)
|
||
else:
|
||
remote_dispatch = {"queued_jobs": [], "queued_total": 0, "failed_total": 0, "target_summary": {"controller_nodes": [], "worker_nodes": []}}
|
||
response_message = f"{message};{command_message}" if command_ok else command_message
|
||
result = _build_detect_action_result(
|
||
action="start",
|
||
ok=command_ok,
|
||
message=response_message,
|
||
data={"job": job_summary, "remote_dispatch": remote_dispatch},
|
||
)
|
||
return ApiResponse(code=0 if command_ok else 1, message=response_message, data=result)
|
||
|
||
|
||
@router.post("/detect/stop", response_model=ApiResponse)
|
||
def stop_detect() -> ApiResponse:
|
||
active_job = get_active_detect_job_summary(event_limit=10)
|
||
ok, message = send_worker_command("stop_detection")
|
||
cycle_token = str((active_job or {}).get("current_cycle_token") or "").strip()
|
||
remote_dispatch = _dispatch_remote_detect_stop(active_job=active_job, cycle_token=cycle_token)
|
||
if active_job:
|
||
append_detect_job_event(
|
||
active_job["job_id"],
|
||
event_type="job_stop_requested" if ok else "job_stop_request_failed",
|
||
level="info" if ok else "error",
|
||
message=message,
|
||
payload={"cycle_token": active_job.get("current_cycle_token", "")},
|
||
)
|
||
append_detect_job_event(
|
||
active_job["job_id"],
|
||
event_type="job_stop_remote_queued",
|
||
level="info" if int(remote_dispatch.get("failed_total", 0) or 0) == 0 else "warning",
|
||
message=(
|
||
f"已向大陆节点排队 {int(remote_dispatch.get('queued_total', 0) or 0)} 个停止检测动作"
|
||
if int(remote_dispatch.get("queued_total", 0) or 0) > 0
|
||
else "当前没有可排队的大陆停止检测动作"
|
||
),
|
||
payload={"cycle_token": cycle_token, "remote_dispatch": remote_dispatch},
|
||
)
|
||
snapshot = get_detect_status()
|
||
settings_payload = get_settings_payload()
|
||
settings_summary = _build_settings_summary(settings_payload)
|
||
mark_detect_run_stopping(
|
||
message=message,
|
||
runtime={
|
||
"mode": snapshot.get("worker_mode", ""),
|
||
"running": snapshot.get("worker_online", False),
|
||
"process_count": snapshot.get("worker_process_count", 0),
|
||
"latest_start_time": snapshot.get("worker_latest_start_time", ""),
|
||
"message": snapshot.get("worker_runtime_message", ""),
|
||
},
|
||
progress=snapshot.get("progress", {}),
|
||
settings_summary=settings_summary,
|
||
)
|
||
if not snapshot.get("worker_online", False):
|
||
finalize_detect_run(
|
||
message=message,
|
||
runtime={
|
||
"mode": snapshot.get("worker_mode", ""),
|
||
"running": snapshot.get("worker_online", False),
|
||
"process_count": snapshot.get("worker_process_count", 0),
|
||
"latest_start_time": snapshot.get("worker_latest_start_time", ""),
|
||
"message": snapshot.get("worker_runtime_message", ""),
|
||
},
|
||
progress=snapshot.get("progress", {}),
|
||
settings_summary=settings_summary,
|
||
active_job=active_job,
|
||
)
|
||
result = _build_detect_action_result(action="stop", ok=ok, message=message, data={"remote_dispatch": remote_dispatch})
|
||
return ApiResponse(code=0 if ok else 1, message=message, data=result)
|