首页 / 文章 / 实用笔记:代理架构——第13篇:人机协同模式

实用笔记:代理架构——第13篇:人机协同模式

《实用笔记:代理架构》操作指南——第13篇:人机协同模式:为采用该模式的团队提供的契约、校验机制及可直接插入的代码模块。

4302 词

以下内容为围绕“代理架构——第13篇:人机协同模式”所设计的实用路径。重点在于契约、校验机制以及可直接插入的代码占位符,而非动机性阐述。 在完成概览阶段时,首先明确契约内容:所需输入、成功信号以及部分失败时的处理方式。这样的清单能确保后续的代码修改保持一致性。 将此阶段视为输入与经过验证的输出之间的契约。为相关成果命名,定义成功校验标准,并杜绝无声的半完成状态。

此处包含的内容

“你将发现什么”阶段若被视为可度量的界面,则效果最佳。在扩大范围之前,先记录一份优秀的测试用例、一个失败案例以及回滚说明。在功能结果旁同时记录执行时间以及令牌或查询成本。提前了解成本情况,可避免在从演示环境过渡到共享环境时出现意外费用。要保持图表状态简洁且类型明确,嵌套的数据块会掩盖是哪个节点修改了哪个字段,还会在出现中断后导致流程无法继续。

如何划定界限

将“舞台绘制位置”视为可测量的表面时,其效果最佳。在扩大范围之前,先记录一个成功的案例、一个失败案例以及回滚说明。 将配置置于应用程序代码之外。环境文件、密钥存储和功能标志应集中存放于一个位置,这样操作人员无需查看整个结构即可进行审计。 保持图结构的扁平化与类型化。嵌套的数据块会掩盖哪个节点修改了哪个字段的信息,且在中断后会导致恢复失败。

+--------------------+------------------------+------------------------+
|                    | Low Error Cost         | High Error Cost        |
+--------------------+------------------------+------------------------+
| Reversible         | AUTONOMOUS             | OVERSIGHT SAMPLING     |
|                    | (agent acts freely)    | (act, log, review some)|
+--------------------+------------------------+------------------------+
| Irreversible       | APPROVAL FOR NOVEL     | APPROVAL GATE          |
|                    | (approve first time,   | (always require human  |
|                    |  then autonomous)      |  approval before act)  |
+--------------------+------------------------+------------------------+
# harness/hitl/action_policy.py
from enum import Enum
from dataclasses import dataclass
from typing import Optional, Callable
class ApprovalPolicy(Enum):
    AUTONOMOUS = "autonomous"              # act freely
    OVERSIGHT_SAMPLING = "sampling"        # act, log, review a sample
    APPROVAL_FOR_NOVEL = "approval_novel"  # approve first occurrence, then auto
    APPROVAL_GATE = "approval_gate"        # always require approval
    DUAL_APPROVAL = "dual_approval"        # require two approvers

@dataclass
class ActionPolicy:
    action_name: str
    policy: ApprovalPolicy
    reason: str
    approver_role: Optional[str] = None    # required role to approve
    timeout_seconds: int = 3600
    timeout_action: str = "reject"         # "reject" | "escalate" | "proceed"

# Example policy configuration for a customer support agent
SUPPORT_AGENT_POLICIES = {
    "read_customer_record": ActionPolicy(
        action_name="read_customer_record",
        policy=ApprovalPolicy.AUTONOMOUS,
        reason="Read-only, reversible, low risk",
    ),
    "draft_response": ActionPolicy(
        action_name="draft_response",
        policy=ApprovalPolicy.OVERSIGHT_SAMPLING,
        reason="Reversible but customer-facing quality matters",
    ),
    "send_customer_email": ActionPolicy(
        action_name="send_customer_email",
        policy=ApprovalPolicy.APPROVAL_GATE,
        reason="Irreversible, customer-facing, reputation risk",
        approver_role="support_agent",
        timeout_seconds=1800,
        timeout_action="reject",
    ),
    "issue_refund": ActionPolicy(
        action_name="issue_refund",
        policy=ApprovalPolicy.DUAL_APPROVAL,
        reason="Financial impact, irreversible",
        approver_role="support_manager",
        timeout_seconds=7200,
        timeout_action="escalate",
    ),
    "delete_account": ActionPolicy(
        action_name="delete_account",
        policy=ApprovalPolicy.DUAL_APPROVAL,
        reason="Catastrophic and irreversible",
        approver_role="senior_manager",
        timeout_seconds=86400,
        timeout_action="reject",
    ),
}

class ActionPolicyEngine:
    """
    Determines whether an action requires human approval and how.
    Consulted at the tool execution layer before any action runs.
    """
    def __init__(self, policies: dict):
        self.policies = policies
    def get_policy(self, action_name: str) -> ActionPolicy:
        return self.policies.get(
            action_name,
            # Default to approval gate for unknown actions: fail safe
            ActionPolicy(
                action_name=action_name,
                policy=ApprovalPolicy.APPROVAL_GATE,
                reason="Unknown action, defaulting to approval required",
            )
        )
    def requires_approval(self, action_name: str, is_novel: bool = False) -> bool:
        policy = self.get_policy(action_name)
        if policy.policy == ApprovalPolicy.AUTONOMOUS:
            return False
        if policy.policy == ApprovalPolicy.OVERSIGHT_SAMPLING:
            return False  # acts first, reviewed after
        if policy.policy == ApprovalPolicy.APPROVAL_FOR_NOVEL:
            return is_novel
        return True  # APPROVAL_GATE and DUAL_APPROVAL always require approval

模式1:审批关卡

将模式1的审批关卡视为可度量的对象时,其效果最佳。在扩大范围之前,需记录一份理想状态下的处理过程、一个失败案例以及回滚说明。同时将正常流程与恢复流程都记录下来。重试机制、人工审核环节以及错误处理都属于产品本身的功能,而非后续需要补充的内容。要保持图表状态的简洁性与类型一致性,嵌套的数据结构会掩盖具体是哪个节点修改了哪个字段,还会导致中断后无法继续处理。将模式1的审批关卡视为输入与已验证输出之间的契约,为相关文档命名、明确成功标准,绝不允许出现无声无息的半完成状态。

# harness/hitl/approval_gate.py
import boto3
import time
import uuid
import json
from typing import Optional, Literal
from dataclasses import dataclass
from langgraph.types import interrupt, Command
@dataclass
class ApprovalRequest:
    request_id: str
    run_id: str
    action_name: str
    action_args: dict
    context_summary: str
    requested_at: float
    approver_role: str
    status: str                # PENDING | APPROVED | REJECTED | TIMEOUT
    decided_by: Optional[str] = None
    decided_at: Optional[float] = None
    decision_note: Optional[str] = None

class ApprovalGateManager:
    """
    Manages approval requests using LangGraph interrupts for pause/resume
    and DynamoDB for durable request state.
    """
    def __init__(
        self,
        table_name: str = "agent-approval-requests",
        region: str = "us-east-1",
    ):
        dynamodb = boto3.resource("dynamodb", region_name=region)
        self.table = dynamodb.Table(table_name)
        self.sns = boto3.client("sns", region_name=region)
    def create_approval_request(
        self,
        run_id: str,
        action_name: str,
        action_args: dict,
        context_summary: str,
        approver_role: str,
        notification_topic_arn: Optional[str] = None,
    ) -> ApprovalRequest:
        """
        Creates a pending approval request and notifies approvers.
        """
        request = ApprovalRequest(
            request_id=str(uuid.uuid4()),
            run_id=run_id,
            action_name=action_name,
            action_args=action_args,
            context_summary=context_summary,
            requested_at=time.time(),
            approver_role=approver_role,
            status="PENDING",
        )
        self.table.put_item(Item={
            "request_id": request.request_id,
            "run_id": request.run_id,
            "action_name": request.action_name,
            "action_args": json.dumps(request.action_args),
            "context_summary": request.context_summary,
            "requested_at": int(request.requested_at),
            "approver_role": request.approver_role,
            "status": "PENDING",
        })
        # Notify approvers
        if notification_topic_arn:
            self.sns.publish(
                TopicArn=notification_topic_arn,
                Subject=f"Approval needed: {action_name}",
                Message=json.dumps({
                    "request_id": request.request_id,
                    "action": action_name,
                    "context": context_summary,
                    "approver_role": approver_role,
                }),
            )
        return request
    def record_decision(
        self,
        request_id: str,
        decision: Literal["APPROVED", "REJECTED"],
        decided_by: str,
        decision_note: Optional[str] = None,
    ) -> bool:
        """
        Records a human's approval decision.
        Called by the approval interface when a human responds.
        """
        self.table.update_item(
            Key={"request_id": request_id},
            UpdateExpression=(
                "SET #s = :status, decided_by = :by, "
                "decided_at = :at, decision_note = :note"
            ),
            ExpressionAttributeNames={"#s": "status"},
            ExpressionAttributeValues={
                ":status": decision,
                ":by": decided_by,
                ":at": int(time.time()),
                ":note": decision_note or "",
            },
            # Only allow decision on PENDING requests: prevents double-decision
            ConditionExpression="#s = :pending",
            ExpressionAttributeValues2={":pending": "PENDING"} if False else None,
        )
        return True
    def get_decision(self, request_id: str) -> Optional[str]:
        """Polls the current decision status for a request."""
        response = self.table.get_item(Key={"request_id": request_id})
        item = response.get("Item")
        return item.get("status") if item else None

# The LangGraph node that implements the approval gate
def approval_gate_node(state: dict) -> dict:
    """
    A LangGraph node that pauses execution and waits for human approval.
    Uses interrupt() to suspend the graph. The graph state is persisted
    and can be resumed when the approval decision arrives.
    """
    pending_action = state.get("pending_action")
    if not pending_action:
        return state
    # interrupt() pauses the graph and surfaces this data to the caller.
    # The caller (your application) presents it to a human and resumes
    # the graph with the decision.
    decision = interrupt({
        "type": "approval_required",
        "action": pending_action["name"],
        "args": pending_action["args"],
        "context": state.get("context_summary", ""),
        "run_id": state.get("agent_run_id"),
    })
    # When resumed, decision contains the human's response
    if decision.get("approved"):
        return {
            **state,
            "action_approved": True,
            "approved_by": decision.get("approver"),
        }
    else:
        return {
            **state,
            "action_approved": False,
            "rejection_reason": decision.get("reason", "Rejected by human reviewer"),
        }

模式2:升级处理

在模式2的升级阶段,应在修改代码之前明确输入参数、该步骤的负责人以及结束标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。应在功能结果旁记录执行时间以及令牌或查询成本。提前了解成本情况,可避免在流程从演示环境转向共享环境时出现意外费用。对于会产生支出或更改生产数据的操作,必须经过人工审批。编译时的配置并不等同于业务功能的完整性。

# harness/hitl/escalation.py
import boto3
import time
import json
from typing import Optional
from dataclasses import dataclass
from langchain_aws import ChatBedrock
from langchain_core.messages import SystemMessage, HumanMessage
ESCALATION_TRIGGERS = {
    "low_confidence": "Agent confidence in its solution is below threshold",
    "conflicting_information": "Retrieved information contradicts itself",
    "policy_ambiguity": "The correct action is genuinely ambiguous under policy",
    "high_stakes_uncertainty": "High-stakes decision with insufficient certainty",
    "repeated_failure": "Agent has failed the same task multiple times",
    "explicit_user_request": "User asked to speak with a human",
}

@dataclass
class EscalationEvent:
    escalation_id: str
    run_id: str
    trigger: str
    agent_context: str
    agent_attempted_solution: Optional[str]
    confidence: float
    escalated_at: float
    assigned_to: Optional[str] = None
    resolution: Optional[str] = None

class EscalationManager:
    """
    Handles cases where the agent should hand off to a human.
    Distinct from approval gates: escalation is triggered by the agent
    recognizing its own limitations.
    """
    def __init__(
        self,
        table_name: str = "agent-escalations",
        region: str = "us-east-1",
    ):
        dynamodb = boto3.resource("dynamodb", region_name=region)
        self.table = dynamodb.Table(table_name)
        bedrock = boto3.client("bedrock-runtime", region_name=region)
        self.confidence_model = ChatBedrock(
            client=bedrock,
            model_id="anthropic.claude-haiku-4-5",
            model_kwargs={"temperature": 0, "max_tokens": 256},
        )
    def should_escalate(
        self,
        task: str,
        proposed_solution: str,
        attempts: int,
    ) -> tuple:
        """
        Assesses whether the agent should escalate to a human.
        Returns (should_escalate, trigger, confidence).
        """
        # Repeated failure is a deterministic trigger
        if attempts >= 3:
            return True, "repeated_failure", 0.0
        # Ask the model to self-assess confidence
        response = self.confidence_model.invoke([
            SystemMessage(content="""
Assess your confidence in a proposed solution. Be honest about uncertainty.
Consider: Is the information sufficient? Is the answer unambiguous?
Are there conflicting considerations? Is this high-stakes?
Return JSON:
{
  "confidence": 0.0 to 1.0,
  "should_escalate": true | false,
  "trigger": "low_confidence | conflicting_information | policy_ambiguity | high_stakes_uncertainty | none"
}
"""),
            HumanMessage(content=f"Task: {task}\n\nProposed solution: {proposed_solution}")
        ])
        try:
            assessment = json.loads(response.content)
            return (
                assessment.get("should_escalate", False),
                assessment.get("trigger", "low_confidence"),
                assessment.get("confidence", 0.5),
            )
        except json.JSONDecodeError:
            # If we cannot assess confidence, escalate to be safe
            return True, "low_confidence", 0.0
    def create_escalation(
        self,
        run_id: str,
        trigger: str,
        agent_context: str,
        attempted_solution: Optional[str],
        confidence: float,
        notification_topic_arn: Optional[str] = None,
    ) -> EscalationEvent:
        import uuid
        event = EscalationEvent(
            escalation_id=str(uuid.uuid4()),
            run_id=run_id,
            trigger=trigger,
            agent_context=agent_context,
            agent_attempted_solution=attempted_solution,
            confidence=confidence,
            escalated_at=time.time(),
        )
        self.table.put_item(Item={
            "escalation_id": event.escalation_id,
            "run_id": event.run_id,
            "trigger": event.trigger,
            "agent_context": event.agent_context[:2000],
            "attempted_solution": (attempted_solution or "")[:2000],
            "confidence": str(confidence),
            "escalated_at": int(event.escalated_at),
            "status": "OPEN",
        })
        if notification_topic_arn:
            boto3.client("sns").publish(
                TopicArn=notification_topic_arn,
                Subject=f"Agent escalation: {trigger}",
                Message=event.agent_context[:1000],
            )
        return event

模式3:协作编辑

在模式3的协作编辑阶段,应在修改代码之前明确输入内容、该步骤的负责人以及结束标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。 配置信息应置于应用程序代码之外。环境文件、密钥存储以及功能标志应集中存放于一个位置,以便操作人员无需查看整个系统结构即可进行审计。 对于涉及资金支出或修改生产数据的操作,必须经过人工审批。编译时的连接方式并不等同于业务功能的完整性。

# harness/hitl/collaborative.py
import boto3
import time
import json
from typing import Optional
from dataclasses import dataclass
@dataclass
class CollaborativeEdit:
    edit_id: str
    run_id: str
    original_draft: str
    human_edited: str
    edit_distance: float        # how much was changed
    edit_categories: list       # "tone", "factual", "structure", "detail"
    edited_by: str
    edited_at: float

class CollaborativeEditManager:
    """
    Manages the draft-edit-finalize flow where a human refines agent output.
    Captures edits as learning signals for future improvement.
    """
    def __init__(
        self,
        table_name: str = "agent-collaborative-edits",
        region: str = "us-east-1",
    ):
        dynamodb = boto3.resource("dynamodb", region_name=region)
        self.table = dynamodb.Table(table_name)
    def record_edit(
        self,
        run_id: str,
        original_draft: str,
        human_edited: str,
        edited_by: str,
        edit_categories: list = None,
    ) -> CollaborativeEdit:
        """
        Records a human edit to agent output.
        The edit distance and categories become learning signals.
        """
        import uuid
        edit_distance = self._compute_edit_distance(original_draft, human_edited)
        edit = CollaborativeEdit(
            edit_id=str(uuid.uuid4()),
            run_id=run_id,
            original_draft=original_draft,
            human_edited=human_edited,
            edit_distance=edit_distance,
            edit_categories=edit_categories or [],
            edited_by=edited_by,
            edited_at=time.time(),
        )
        self.table.put_item(Item={
            "edit_id": edit.edit_id,
            "run_id": edit.run_id,
            "original_draft": original_draft[:5000],
            "human_edited": human_edited[:5000],
            "edit_distance": str(edit_distance),
            "edit_categories": edit.edit_categories,
            "edited_by": edited_by,
            "edited_at": int(edit.edited_at),
        })
        return edit
    def analyze_edit_patterns(
        self,
        run_ids: list = None,
        min_edits: int = 20,
    ) -> dict:
        """
        Analyzes accumulated edits to find systematic patterns.
        If humans consistently edit for the same reason, the agent's
        system prompt or procedure should be updated.
        """
        response = self.table.scan(Limit=200)
        edits = response.get("Items", [])
        if len(edits) < min_edits:
            return {"sufficient_data": False, "edit_count": len(edits)}
        # Aggregate edit categories
        category_counts = {}
        total_distance = 0.0
        for edit in edits:
            for cat in edit.get("edit_categories", []):
                category_counts[cat] = category_counts.get(cat, 0) + 1
            total_distance += float(edit.get("edit_distance", 0))
        avg_distance = total_distance / len(edits)
        dominant_category = max(category_counts, key=category_counts.get) if category_counts else None
        return {
            "sufficient_data": True,
            "edit_count": len(edits),
            "average_edit_distance": round(avg_distance, 3),
            "category_distribution": category_counts,
            "dominant_edit_category": dominant_category,
            "recommendation": self._recommendation(dominant_category, avg_distance),
        }
    def _recommendation(self, dominant_category: Optional[str], avg_distance: float) -> str:
        if avg_distance < 0.1:
            return "Edits are minor. Agent output quality is high."
        if dominant_category == "tone":
            return "Humans frequently adjust tone. Update system prompt with tone guidance."
        if dominant_category == "factual":
            return "Humans frequently correct facts. Review agent's information sources."
        if dominant_category == "structure":
            return "Humans frequently restructure. Add output format guidance to system prompt."
        return "Review edit patterns manually for systematic improvements."
    def _compute_edit_distance(self, a: str, b: str) -> float:
        """Normalized Levenshtein distance, 0 (identical) to 1 (completely different)."""
        if not a and not b:
            return 0.0
        # Simple word-level distance for efficiency on long text
        a_words, b_words = a.split(), b.split()
        max_len = max(len(a_words), len(b_words))
        if max_len == 0:
            return 0.0
        # Count matching words in order (simplified)
        common = len(set(a_words) & set(b_words))
        return 1.0 - (common / max_len)

模式4:监督抽样

在模式4的监督抽样阶段,修改代码之前需明确输入参数、该步骤的负责人以及终止标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。 需同时记录正常流程与异常恢复流程。重试机制、人工审核环节以及错误处理都是产品本身的组成部分,而非后续需要补充的内容。 对于涉及资金支出或修改生产数据的操作,必须经过人工审批。编译时的逻辑连接并不等同于业务功能的完整性。 在模式4的监督抽样阶段,修改代码之前需明确输入参数、该步骤的负责人以及终止标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。 应将此阶段视为输入与验证后输出之间的契约。需为相关成果命名、明确成功判定标准,并杜绝无声的半完成状态。

# harness/hitl/oversight_sampling.py
import boto3
import random
import time
import json
from typing import Optional
class OversightSampler:
    """
    Routes a sample of autonomous actions to human review after execution.
    Does not block the action. Catches quality drift and systematic errors.
    Feeds into the Article 8 evaluation pipeline.
    """
    def __init__(
        self,
        base_sample_rate: float = 0.05,
        table_name: str = "agent-oversight-samples",
        region: str = "us-east-1",
    ):
        self.base_sample_rate = base_sample_rate
        dynamodb = boto3.resource("dynamodb", region_name=region)
        self.table = dynamodb.Table(table_name)
    def should_sample(
        self,
        action_name: str,
        agent_confidence: float = 1.0,
        is_novel: bool = False,
    ) -> bool:
        """
        Decides whether this action should be sampled for review.
        Samples more aggressively for low-confidence and novel actions.
        """
        rate = self.base_sample_rate
        # Increase sampling for low confidence
        if agent_confidence < 0.7:
            rate = min(rate * 3, 1.0)
        # Always sample novel actions
        if is_novel:
            return True
        return random.random() < rate
    def record_for_review(
        self,
        run_id: str,
        action_name: str,
        action_args: dict,
        action_result: str,
        agent_confidence: float,
    ):
        """Queues an executed action for human review."""
        import uuid
        self.table.put_item(Item={
            "sample_id": str(uuid.uuid4()),
            "run_id": run_id,
            "action_name": action_name,
            "action_args": json.dumps(action_args)[:2000],
            "action_result": action_result[:2000],
            "agent_confidence": str(agent_confidence),
            "sampled_at": int(time.time()),
            "review_status": "PENDING",
        })
    def record_review_outcome(
        self,
        sample_id: str,
        was_correct: bool,
        reviewer: str,
        notes: Optional[str] = None,
    ):
        """
        Records the human's assessment of a sampled action.
        A pattern of incorrect actions should trigger a policy review.
        """
        self.table.update_item(
            Key={"sample_id": sample_id},
            UpdateExpression=(
                "SET review_status = :s, was_correct = :c, "
                "reviewed_by = :r, review_notes = :n, reviewed_at = :at"
            ),
            ExpressionAttributeValues={
                ":s": "REVIEWED",
                ":c": was_correct,
                ":r": reviewer,
                ":n": notes or "",
                ":at": int(time.time()),
            }
        )
    def get_error_rate(self, action_name: str, days: int = 7) -> dict:
        """
        Computes the error rate for a sampled action over a window.
        If this exceeds threshold, the action's policy should be tightened.
        """
        cutoff = int(time.time()) - (days * 86400)
        response = self.table.scan(
            FilterExpression="action_name = :a AND sampled_at > :c AND review_status = :s",
            ExpressionAttributeValues={
                ":a": action_name,
                ":c": cutoff,
                ":s": "REVIEWED",
            }
        )
        reviews = response.get("Items", [])
        if not reviews:
            return {"sufficient_data": False}
        incorrect = sum(1 for r in reviews if not r.get("was_correct", True))
        error_rate = incorrect / len(reviews)
        return {
            "sufficient_data": len(reviews) >= 10,
            "sample_size": len(reviews),
            "error_rate": round(error_rate, 3),
            "recommendation": (
                "TIGHTEN_POLICY" if error_rate > 0.1 else "MAINTAIN"
            ),
        }

超时与回退处理

在处理超时与回退处理阶段时,首先明确相关规范:所需的输入参数、成功信号,以及部分失败时的处理方式。这样的清单能确保后续的代码修改保持一致性。 在功能结果旁记录执行时间以及令牌或查询成本。提前了解成本情况,可避免在从演示环境切换到共享环境时出现意外费用。 在耗时较高的步骤之后设置检查点。当操作员重新尝试后续节点时,恢复流程不应再次对同一次大型语言模型调用收费。

# harness/hitl/timeout_handler.py
import boto3
import time
from typing import Literal
class ApprovalTimeoutHandler:
    """
    Handles approval requests that exceed their timeout.
    The timeout action is policy-defined per action type.
    Runs as a scheduled Lambda, checking for expired pending requests.
    """
    def __init__(
        self,
        approval_table: str = "agent-approval-requests",
        region: str = "us-east-1",
    ):
        dynamodb = boto3.resource("dynamodb", region_name=region)
        self.table = dynamodb.Table(approval_table)
        self.sns = boto3.client("sns", region_name=region)
    def process_expired_requests(self, policies: dict) -> dict:
        """
        Scans for pending requests past their timeout and applies
        the policy-defined timeout action.
        """
        now = int(time.time())
        response = self.table.scan(
            FilterExpression="#s = :pending",
            ExpressionAttributeNames={"#s": "status"},
            ExpressionAttributeValues={":pending": "PENDING"},
        )
        results = {"rejected": 0, "escalated": 0, "proceeded": 0}
        for item in response.get("Items", []):
            requested_at = int(item.get("requested_at", now))
            action_name = item.get("action_name", "")
            policy = policies.get(action_name)
            if not policy:
                continue
            age = now - requested_at
            if age < policy.timeout_seconds:
                continue  # not expired yet
            timeout_action = policy.timeout_action
            if timeout_action == "reject":
                self._apply_timeout(item["request_id"], "REJECTED",
                                    "Timed out without approval")
                results["rejected"] += 1
            elif timeout_action == "escalate":
                self._escalate_expired(item)
                results["escalated"] += 1
            elif timeout_action == "proceed":
                # Only for low-risk actions where the gate is advisory
                self._apply_timeout(item["request_id"], "APPROVED",
                                    "Auto-approved after timeout per policy")
                results["proceeded"] += 1
        return results
    def _apply_timeout(self, request_id: str, status: str, note: str):
        self.table.update_item(
            Key={"request_id": request_id},
            UpdateExpression="SET #s = :status, decision_note = :note, decided_at = :at",
            ExpressionAttributeNames={"#s": "status"},
            ExpressionAttributeValues={
                ":status": status,
                ":note": note,
                ":at": int(time.time()),
            }
        )
    def _escalate_expired(self, item: dict):
        # Notify a higher tier and extend the deadline
        self.table.update_item(
            Key={"request_id": item["request_id"]},
            UpdateExpression="SET escalated = :t, escalated_at = :at",
            ExpressionAttributeValues={":t": True, ":at": int(time.time())},
        )

审计追踪

在处理审计追踪阶段时,首先写下合同的相关内容:所需输入、成功信号以及部分失败时的处理方式。这样的检查清单能确保后续的代码修改保持透明。 将配置信息与应用程序代码分开存放。环境文件、密钥存储以及功能开关应集中于一个位置,这样操作人员无需查看整个系统结构即可进行审计。 在耗时较高的步骤之后设置检查点。当操作人员重新尝试后续节点时,恢复流程不应再次对同一次大型语言模型调用收费。

# harness/hitl/audit.py
import boto3
import time
import json
import hashlib
from typing import Optional
class HITLAuditLog:
    """
    Immutable audit trail for all human-in-the-loop decisions.
    Uses a hash chain so tampering is detectable.
    Writes to DynamoDB with a separate append-only access pattern.
    """
    def __init__(
        self,
        table_name: str = "agent-hitl-audit",
        region: str = "us-east-1",
    ):
        dynamodb = boto3.resource("dynamodb", region_name=region)
        self.table = dynamodb.Table(table_name)
    def record_decision(
        self,
        run_id: str,
        decision_type: str,       # "approval" | "rejection" | "escalation" | "edit"
        action_name: str,
        decided_by: str,
        decision_details: dict,
        previous_hash: Optional[str] = None,
    ) -> str:
        """
        Records a decision in the audit log with a hash chain.
        Returns the hash of this entry for chaining the next one.
        """
        timestamp = int(time.time())
        entry = {
            "run_id": run_id,
            "decision_type": decision_type,
            "action_name": action_name,
            "decided_by": decided_by,
            "decision_details": json.dumps(decision_details),
            "timestamp": timestamp,
            "previous_hash": previous_hash or "genesis",
        }
        # Compute hash of this entry chained to the previous
        entry_content = json.dumps(entry, sort_keys=True)
        entry_hash = hashlib.sha256(entry_content.encode()).hexdigest()
        self.table.put_item(Item={
            "audit_id": f"{run_id}#{timestamp}#{entry_hash[:8]}",
            **entry,
            "entry_hash": entry_hash,
        })
        return entry_hash
    def verify_chain(self, run_id: str) -> bool:
        """
        Verifies the hash chain for a run's audit entries.
        Returns True if the chain is intact, False if tampering is detected.
        """
        response = self.table.query(
            KeyConditionExpression="run_id = :rid",
            ExpressionAttributeValues={":rid": run_id},
            ScanIndexForward=True,
        )
        entries = sorted(response.get("Items", []), key=lambda x: x["timestamp"])
        previous_hash = "genesis"
        for entry in entries:
            if entry.get("previous_hash") != previous_hash:
                return False
            # Recompute and verify
            check_entry = {
                k: entry[k] for k in
                ["run_id", "decision_type", "action_name", "decided_by",
                 "decision_details", "timestamp", "previous_hash"]
            }
            recomputed = hashlib.sha256(
                json.dumps(check_entry, sort_keys=True).encode()
            ).hexdigest()
            if recomputed != entry.get("entry_hash"):
                return False
            previous_hash = entry["entry_hash"]
        return True

生产环境现实检验

在开展生产环境可行性验证阶段时,首先需明确相关约定:所需的输入参数、成功标志以及部分失败时的处理方式。这份清单能确保后续的代码修改保持一致性。 同时记录正常流程与故障恢复路径。重试机制、人工审核环节以及错误消息处理都是产品功能的一部分,而非后续需要补充的内容。 在成本较高的步骤之后设置检查点。当操作员重新尝试某个节点时,恢复流程不应再次调用相同的大型语言模型接口。 在开展生产环境可行性验证阶段时,首先需明确相关约定:所需的输入参数、成功标志以及部分失败时的处理方式。这份清单能确保后续的代码修改保持一致性。 将此阶段视为输入与验证后输出之间的契约。为相关成果命名,定义成功判定标准,杜绝无声的半完成状态。

参考架构

将参考架构阶段视为可度量的对象来处理,效果最佳。在扩大范围之前,先记录一份最优案例、一个故障场景以及回滚说明。在功能结果旁同时记录执行时间以及令牌或查询成本。提前了解成本情况,可避免在从演示环境过渡到共享环境时出现意外费用。需保持图表状态简洁且具有明确类型;嵌套的数据块会掩盖哪个节点修改了哪个字段的信息,且在中断后会导致流程无法继续。

Agent reaches an action
          |
          v
+---------------------------+
|   ActionPolicyEngine      |
|   Look up action policy   |
+---------------------------+
          |
    +-----+-----+-----------+-----------+
    |     |     |           |           |
    v     v     v           v           v
 AUTONO  SAMP  APPROVAL   APPROVAL   DUAL
 MOUS    LING  FOR NOVEL  GATE       APPROVAL
    |     |     |           |           |
    |     |     v           v           v
    |     |  +--------------------------------+
    |     |  |  ApprovalGateManager           |
    |     |  |  - create request              |
    |     |  |  - LangGraph interrupt()       |
    |     |  |  - notify approvers (SNS)      |
    |     |  |  - persist state (DynamoDB)    |
    |     |  +--------------------------------+
    |     |           |
    |     |     Human decides via interface
    |     |           |
    |     |     +-----+-----+
    |     |     |           |
    |     |  APPROVED   REJECTED / TIMEOUT
    |     |     |           |
    v     v     v           v
+--------------------------------+
|   Execute or Abort Action      |
+--------------------------------+
          |
          v
+--------------------------------+
|   HITLAuditLog (hash chain)    |
|   Immutable decision record    |
+--------------------------------+
          |
          v
+--------------------------------+
|   Feed to Article 8 eval +     |
|   collaborative edit learning  |
+--------------------------------+

参考基础设施栈

将“参考基础设施堆栈”阶段视为可度量的对象来处理时,其效果最佳。在扩大范围之前,先记录一份完美的操作日志、一个故障案例以及回滚说明。 将配置与应用程序代码分开。环境文件、密钥存储和功能标志应集中存放于一个位置,以便操作人员无需查看整个系统结构即可进行审计。 保持系统状态的扁平化与类型化。嵌套的数据结构会掩盖哪个节点修改了哪个字段的信息,且在中断后会导致无法继续执行。

+-----------------------------+---------------------+------------------------------+
| Component                   | Technology          | Role                         |
+-----------------------------+---------------------+------------------------------+
| Action Policy Engine        | Custom              | Per-action approval routing  |
|                             |                     | based on risk profile        |
+-----------------------------+---------------------+------------------------------+
| Pause / Resume              | LangGraph interrupt | Durable suspend while        |
|                             | + checkpointing     | waiting for human            |
+-----------------------------+---------------------+------------------------------+
| Approval Requests           | DynamoDB            | Pending request state,       |
|                             |                     | decision records             |
+-----------------------------+---------------------+------------------------------+
| Notifications               | SNS                 | Alert approvers when a        |
|                             |                     | decision is needed           |
+-----------------------------+---------------------+------------------------------+
| Escalation                  | Custom + Haiku      | Agent self-assessment of     |
|                             | confidence model    | when to hand off to human    |
+-----------------------------+---------------------+------------------------------+
| Collaborative Edits         | DynamoDB            | Draft-edit-finalize flow,    |
|                             |                     | edit pattern learning        |
+-----------------------------+---------------------+------------------------------+
| Oversight Sampling          | Custom + DynamoDB   | Post-hoc review of a sample  |
|                             |                     | of autonomous actions        |
+-----------------------------+---------------------+------------------------------+
| Timeout Handling            | Scheduled Lambda    | Policy-defined action on     |
|                             |                     | expired approval requests    |
+-----------------------------+---------------------+------------------------------+
| Audit Trail                 | DynamoDB hash chain | Immutable, tamper-evident    |
|                             |                     | record of all decisions      |
+-----------------------------+---------------------+------------------------------+
| Auth for Approvers          | CredentialManager   | Role-based approval          |
|                             | (Article 5)         | authorization                |
+-----------------------------+---------------------+------------------------------+
| Learning Loop               | Article 8 eval      | Edits and reviews feed the   |
|                             | pipeline            | improvement pipeline         |
+-----------------------------+---------------------+------------------------------+
| Local Dev Alternative       | Docker Compose +    | Full HITL flow testable      |
|                             | mock approver UI    | without AWS                  |
+-----------------------------+---------------------+------------------------------+

操作检查清单

在“操作检查清单”阶段,应在修改代码之前明确输入内容、该步骤的负责人以及完成标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏的状态信息。

应优先选择小型、可测试的单元,而非庞大的脚本。当某个步骤出错时,故障应指向单一责任点,而非复杂的流程链。

对于涉及资金支出或修改生产数据的环节,必须经过人工审批。编译时的逻辑连接并不等同于业务功能的完整性。

编写简短的操作手册:说明如何轮换密钥、如何清空队列、以及如何回滚上一次的数据导入操作。

将此阶段视为输入数据与经过验证的输出结果之间的契约。为相关产物命名,明确成功标准,绝不允许出现无声的、不完整的处理结果。

对于涉及资金支出或修改生产数据的环节,必须经过人工审批。编译时的逻辑连接并不等同于业务功能的完整性。

在推广该技术栈之前,应先冻结版本,为关键流程记录标准输出日志,并明确回滚步骤。共享环境需要设置速率限制、租户验证机制,以及负责密钥轮换的明确责任人。与其展示花哨的一次性演示,不如注重扎实的可靠性。

c9f5fabd2c2d 的批量处理说明:不要将提供商密钥放入代码仓库,为每个会话设置令牌使用上限,并将日志存储在评估用示例文件旁边,以便后续模型更换时保持数据可比性。