Files
JobRadar/agent_runtime/services.py
T

171 lines
6.7 KiB
Python

"""Agent Run 的事务、权限和状态转换服务。"""
from django.db import transaction
from django.db.models import Max
from django.utils import timezone
from common.exceptions import InvalidStateTransition, PermissionDenied
from common.logging import sanitize_summary
from .models import (
AgentRun,
AgentRunEvent,
ApprovalStatus,
HumanApproval,
RunStatus,
ToolCall,
ToolCallStatus,
)
TRANSITIONS = {
RunStatus.PENDING: {RunStatus.RUNNING, RunStatus.CANCELLED},
RunStatus.RUNNING: {RunStatus.WAITING_APPROVAL, RunStatus.SUCCEEDED, RunStatus.FAILED, RunStatus.CANCELLED},
RunStatus.WAITING_APPROVAL: {RunStatus.RUNNING, RunStatus.FAILED, RunStatus.CANCELLED},
}
def runs_for_user(user):
"""普通站点始终限制为当前用户数据,管理员跨用户查询使用 Admin。"""
return AgentRun.objects.filter(owner=user)
def _append_event(run, event_type: str, summary: str, payload=None):
"""在持有 Run 写锁的事务中分配下一个事件序号。"""
last = run.events.aggregate(value=Max("sequence"))["value"] or 0
return AgentRunEvent.objects.create(
run=run,
sequence=last + 1,
event_type=event_type,
summary=summary,
payload_summary=sanitize_summary(payload or {}),
occurred_at=timezone.now(),
)
@transaction.atomic
def create_run(owner, title: str, input_summary=None) -> AgentRun:
"""创建等待运行的记录,并原子追加创建事件。"""
run = AgentRun.objects.create(
owner=owner, title=title.strip(), input_summary=sanitize_summary(input_summary or {})
)
_append_event(run, "run_created", "运行记录已创建")
return run
@transaction.atomic
def transition_run(run_id, target_status: str, *, error_code="", error_summary="", output=None):
"""校验并执行一次状态转换,主记录与审计事件同时提交。"""
run = AgentRun.objects.select_for_update().get(pk=run_id)
if target_status not in TRANSITIONS.get(run.status, set()):
raise InvalidStateTransition(f"不允许从 {run.status} 转换到 {target_status}。")
now = timezone.now()
run.status = target_status
run.lock_version += 1
if target_status == RunStatus.RUNNING and run.started_at is None:
run.started_at = now
if target_status in {RunStatus.SUCCEEDED, RunStatus.FAILED, RunStatus.CANCELLED}:
run.finished_at = now
if run.started_at:
run.duration_ms = max(0, int((now - run.started_at).total_seconds() * 1000))
if target_status == RunStatus.SUCCEEDED:
run.output_summary = sanitize_summary(output or {})
if target_status == RunStatus.FAILED:
run.error_code = error_code[:80]
run.error_summary = error_summary[:2000]
run.save()
_append_event(run, f"run_{target_status}", f"运行状态变更为 {run.get_status_display()}")
return run
@transaction.atomic
def request_approval(run_id, request_key: str, approval_type: str, summary=None):
"""暂停运行并创建唯一的待确认请求。"""
run = AgentRun.objects.select_for_update().get(pk=run_id)
if run.status != RunStatus.RUNNING:
raise InvalidStateTransition("只有运行中的任务可以请求人工确认。")
approval = HumanApproval.objects.create(
run=run,
request_key=request_key,
approval_type=approval_type,
request_summary=sanitize_summary(summary or {}),
requested_at=timezone.now(),
)
run.status = RunStatus.WAITING_APPROVAL
run.lock_version += 1
run.save(update_fields=("status", "lock_version", "updated_at"))
_append_event(run, "approval_requested", "运行等待人工确认", {"request_key": request_key})
return approval
@transaction.atomic
def start_tool_call(run_id, call_id: str, tool_name: str, arguments=None, idempotency_key=""):
"""在外部调用前登记开始状态;唯一约束负责阻止重复调用标识和幂等键。"""
run = AgentRun.objects.select_for_update().get(pk=run_id)
if run.status != RunStatus.RUNNING:
raise InvalidStateTransition("只有运行中的任务可以开始工具调用。")
call = ToolCall.objects.create(
run=run,
call_id=call_id,
tool_name=tool_name,
idempotency_key=idempotency_key,
arguments_summary=sanitize_summary(arguments or {}),
started_at=timezone.now(),
)
_append_event(run, "tool_started", f"工具 {tool_name} 开始执行", {"call_id": call_id})
return call
@transaction.atomic
def finish_tool_call(call_id: int, *, result=None, error_code="", error_summary=""):
"""在外部调用结束后的独立短事务中登记成功或失败结果。"""
call = ToolCall.objects.select_for_update().select_related("run").get(pk=call_id)
if call.status != ToolCallStatus.STARTED:
raise InvalidStateTransition("工具调用已经结束。")
now = timezone.now()
call.finished_at = now
call.duration_ms = max(0, int((now - call.started_at).total_seconds() * 1000))
if error_code:
call.status = ToolCallStatus.FAILED
call.error_code = error_code[:80]
call.error_summary = error_summary[:2000]
event_type, summary = "tool_failed", f"工具 {call.tool_name} 执行失败"
else:
call.status = ToolCallStatus.SUCCEEDED
call.result_summary = sanitize_summary(result or {})
event_type, summary = "tool_completed", f"工具 {call.tool_name} 执行完成"
call.save()
run = AgentRun.objects.select_for_update().get(pk=call.run_id)
_append_event(run, event_type, summary, {"call_id": call.call_id})
return call
@transaction.atomic
def resolve_approval(actor, approval_id, approved: bool, summary=None):
"""只允许所属用户或管理员处理一次待确认请求。"""
approval = HumanApproval.objects.select_for_update().select_related("run").get(pk=approval_id)
if actor != approval.run.owner and not actor.is_staff:
raise PermissionDenied("无权处理该确认请求。")
if approval.status != ApprovalStatus.PENDING:
raise InvalidStateTransition("该确认请求已经处理。")
approval.status = ApprovalStatus.APPROVED if approved else ApprovalStatus.REJECTED
approval.resolved_by = actor
approval.resolved_at = timezone.now()
approval.decision_summary = sanitize_summary(summary or {})
approval.save()
run = AgentRun.objects.select_for_update().get(pk=approval.run_id)
run.status = RunStatus.RUNNING if approved else RunStatus.CANCELLED
run.lock_version += 1
if not approved:
run.finished_at = timezone.now()
run.save()
_append_event(run, "approval_resolved", "人工确认已处理", {"approved": approved})
return approval