171 lines
6.7 KiB
Python
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
|