205 lines
7.6 KiB
Python
205 lines
7.6 KiB
Python
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.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_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", [])),
|
||
}
|
||
|
||
|
||
@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:
|
||
return ApiResponse(
|
||
code=0,
|
||
message="当前没有可创建的检测任务",
|
||
data={
|
||
"action": "start",
|
||
"job": None,
|
||
"poll_after_seconds": 2,
|
||
"refresh_status": True,
|
||
},
|
||
)
|
||
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:
|
||
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={
|
||
"action": "start",
|
||
"job": job_summary,
|
||
"poll_after_seconds": 2,
|
||
"refresh_status": True,
|
||
},
|
||
)
|
||
|
||
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:
|
||
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,
|
||
)
|
||
return ApiResponse(
|
||
code=0 if command_ok else 1,
|
||
message=f"{message};{command_message}" if command_ok else command_message,
|
||
data={
|
||
"action": "start",
|
||
"job": job_summary,
|
||
"poll_after_seconds": 2,
|
||
"refresh_status": True,
|
||
},
|
||
)
|
||
|
||
|
||
@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")
|
||
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", "")},
|
||
)
|
||
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,
|
||
)
|
||
return ApiResponse(
|
||
code=0 if ok else 1,
|
||
message=message,
|
||
data={
|
||
"action": "stop",
|
||
"poll_after_seconds": 2,
|
||
"refresh_status": True,
|
||
},
|
||
)
|