Loop Engineering 案例一
案例业务
这是一个合同质量的自动迭代打磨系统。不是简单的”审一下→通过/驳回”,而是让两个AI Agent(Editor和Reviewer)像律师团队一样循环协作:一个负责改,一个负责审,改完再审,审完再改,直到合同质量达标。
业务逻辑

四大审查维度(定义在 skills/contract-improver/SKILL.md):
| 维度 | 检查内容 | 示例 |
|---|---|---|
| A.完整性 | 必备要素是否齐全? | 双方信息、金额(大小写)、付款方式、交付时间、违约金、保密、知识产权、签署信息 |
| B.明确性 | 有没有模糊表述? | “尽快”→”60个工作日内”、”协商解决”→”提交XX仲裁委员会” |
| C.公平性 | 双方权利义务平衡吗? | 违约金双向对等、付款条件合理、知识产权归属公平 |
| D.合规性 | 符合法律和行业标准吗? | 违约金≤20%、保密期限≥合同期+解约后2年 |
完整代码
"""
Loop Engineering — 合同智能完善系统
Python 代码控制循环 | checkpointer 上下文持久化 | 类化工程封装
═══════════════════════════════════════════════════════════════════════════
业务场景:
合同初稿存在条款缺失、表述模糊、条款不平衡等问题。
通过 Editor-Reviewer 循环协作,自动迭代打磨直到合同达到标准。
Loop Engineering 核心设计:
★ Python for 循环控制重试(代码决定,不是 prompt 决定)
★ checkpointer + thread_id 上下文持久化(每轮消息完全一致,LLM 从历史中知道该做什么)
★ ContractImproveLoop 类工程封装(对标业界 Loop 类设计模式)
★ 精简提示词(约束清单式,符合 2025-2026 最佳实践)
★ Maker-Checker 硬权限隔离(Editor 写,Reviewer 只读)
五大构建块:
① 自动化/调度 → while True 持续监控 pending_contracts/
② 工作树/隔离 → FilesystemBackend(virtual_mode=True)
③ 技能/知识 → skills/contract-improver/SKILL.md
④ 连接器/插件 → read_contract / improve_contract / finalize / escalate
⑤ 子代理/协作 → contract-editor(Maker) + contract-reviewer(Checker)
运行:
python.exe contract_improve_loop_python.py
"""
import os, sys, time, glob, shutil, re
from datetime import datetime
from dataclasses import dataclass, field
from deepagents import create_deep_agent
from deepagents.middleware.subagents import SubAgent
from deepagents.backends import FilesystemBackend
from langchain.tools import tool
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.store.memory import InMemoryStore
from my_llm import deepseek_llm, deepseek_llm_flash
# ═══════════════════════════════════════════════════════════════════════
# 路径配置
# ═══════════════════════════════════════════════════════════════════════
CUR_DIR = os.path.dirname(os.path.abspath(__file__))
PENDING = os.path.join(CUR_DIR, "pending_contracts")
COMPLETED = os.path.join(CUR_DIR, "agent_workspace", "completed")
ESCALATED = os.path.join(CUR_DIR, "agent_workspace", "escalated")
SKILLS = os.path.join(CUR_DIR, "skills")
for d in [COMPLETED, ESCALATED]:
os.makedirs(d, exist_ok=True)
# ═══════════════════════════════════════════════════════════════════════
# 记忆系统
# ═══════════════════════════════════════════════════════════════════════
store = InMemoryStore()
# ═══════════════════════════════════════════════════════════════════════
# 工具定义 —— 按权限严格分层
# ═══════════════════════════════════════════════════════════════════════
@tool
def read_contract(contract_id: str) -> str:
"""读取指定合同的内容。输入合同编号(如 'CT-IMP-001'),返回合同全文。"""
file_path = os.path.join(PENDING, f"{contract_id}.txt")
try:
with open(file_path, "r", encoding="utf-8") as f:
return f.read()
except FileNotFoundError:
return f"合同 {contract_id} 不存在于待处理目录"
@tool
def improve_contract(contract_id: str, improved_content: str) -> str:
"""
写入完善后的合同内容(覆盖原文件)。
只有 contract-editor 子代理有权限调用此工具。
"""
file_path = os.path.join(PENDING, f"{contract_id}.txt")
if not os.path.exists(file_path):
return f"错误:合同 {contract_id} 不在待处理目录中"
with open(file_path, "w", encoding="utf-8") as f:
f.write(improved_content)
# 记录修改操作
store.put(("improvements",), f"{contract_id}_edit_{datetime.now().timestamp():.0f}", {
"contract_id": contract_id,
"action": "improved",
"time": datetime.now().isoformat(),
})
return f"合同 {contract_id} 已更新完善,等待审查员复核。"
@tool
def finalize_contract(contract_id: str, comment: str) -> str:
"""合同完善完成,移至已完成目录。只有编排者能调用。"""
_move_and_log(contract_id, COMPLETED, "completed", comment)
return f"合同 {contract_id} 完善完成,已归档。"
@tool
def escalate_contract(contract_id: str, reason: str) -> str:
"""升级给人工处理,移至升级目录。只有编排者能调用。"""
_move_and_log(contract_id, ESCALATED, "escalated", reason)
return f"合同 {contract_id} 已升级人工处理"
def _move_and_log(cid: str, target_dir: str, status: str, note: str):
src = os.path.join(PENDING, f"{cid}.txt")
dst = os.path.join(target_dir, f"{cid}.txt")
if os.path.exists(src):
shutil.move(src, dst)
store.put(("improvements",), f"{cid}_{status}", {
"contract_id": cid, "status": status, "note": note[:200],
"time": datetime.now().isoformat()
})
# ═══════════════════════════════════════════════════════════════════════
# 子代理定义 —— 精简提示词,约束清单式
# ═══════════════════════════════════════════════════════════════════════
contract_editor = SubAgent(
name="contract-editor",
description="合同修改员:读取合同 → 对照四大维度检查 → 动手修改完善",
system_prompt="""你是合同修改员。用 read_contract 读取合同,用 improve_contract 写入完善版。
对照 SKILL.md 四大维度检查并修复:
- 完整性:双方信息、金额(大小写)、付款方式、交付时间、违约金、保密、知识产权、签署信息
- 明确性:替换"尽快""适时""合理期限""协商解决""约""左右"等模糊词为具体描述
- 公平性:违约金双向对等、付款条件合理、知识产权归属公平
- 合规性:违约金≤20%、保密期限≥合同期+解约后2年
原则:保留原意、补充缺失、明确模糊、平衡权利。
⚠️ 约束:
- 根据你检查的结果或者 Reviewer 的建议,统一全面修改合同,修改完成后简要说明改动,然后立即结束。
- 不要审查自己的修改结果。
- 合同只做一轮修改,不要反复修改同一份合同。""",
tools=[read_contract, improve_contract],
model=deepseek_llm,
)
contract_reviewer = SubAgent(
name="contract-reviewer",
description="合同审查员:独立读取合同 → 逐项检查 → 给出审查结论",
system_prompt="""你是合同审查员。只有 read_contract 权限,不能修改合同。
逐项检查:
- 完整性:双方信息、金额(大小写)、付款方式、交付时间、违约金、保密、知识产权、签署信息
- 明确性:无"尽快""适时""合理期限""协商解决""约""左右"等模糊词
- 公平性:违约金双向对等、付款条件合理、知识产权归属公平
- 合规性:违约金≤20%、保密期限≥合同期+解约后2年
通过回复"审查通过"。不通过回复"审查不通过"并逐条列出问题和修正建议。
⚠️ 约束:
- 给出审查结论后立即结束。你的职责是判断质量,不是控制流程。
- 不要写"修改后再来""请修正后重新提交"等循环指令。
- 你只需要独立审查合同,给出"审查通过"/"审查不通过"以及问题和修正建议。""",
tools=[read_contract],
model=deepseek_llm_flash,
)
# ═══════════════════════════════════════════════════════════════════════
# 编排者 Agent
# 编排者不知道"轮次"概念 —— 它只负责每次收到指令时执行一轮 Editor→Reviewer
# 轮次控制完全由外部 Python 代码负责
# ═══════════════════════════════════════════════════════════════════════
def _build_orchestrator():
"""构建编排者 Agent(含 checkpointer 用于上下文持久化)"""
return create_deep_agent(
model=deepseek_llm,
# 两个工具功能:finalize_contract 归档合同,escalate_contract 升级人工处理
tools=[finalize_contract, escalate_contract],
subagents=[contract_editor, contract_reviewer],
skills=[SKILLS],
backend=FilesystemBackend(root_dir=CUR_DIR, virtual_mode=True),
store=store,
checkpointer=InMemorySaver(),
system_prompt="""你是合同完善编排员。你只负责委派子代理,不亲自修改或审查合同。
每次收到指令时,执行恰好一轮 Editor→Reviewer 流程,然后立即结束:
1. 委派 contract-editor 修改合同(首次处理则全面完善,续轮则针对性修正历史中的问题)
2. 委派 contract-reviewer 独立审查当前合同
3. 根据审查结果决定:
- "审查通过" → 调用 finalize_contract 归档
- "审查不通过" → 汇报问题,等待外部程序发出下次指令
工具权限:
- 委派子代理(contract-editor / contract-reviewer)
- finalize_contract:审查通过后归档
- escalate_contract:遇到无法自动修复的根本性问题时升级
硬性约束(违反将导致系统失控):
- 禁止自行循环!审查不通过时,绝对不要再次委派 Editor。只汇报问题然后结束。
- 每次调用最多委派 Editor 一次、Reviewer 一次。
- 即使 Reviewer 列出多个问题,也不要回到步骤1重新开始。外部程序会处理循环。""",
)
# ═══════════════════════════════════════════════════════════════════════
# 统计数据结构
# ═══════════════════════════════════════════════════════════════════════
@dataclass
class ContractStats:
scanned: int = 0 # 已扫描合同数
completed: int = 0 # 已完成合同数
escalated: int = 0 # 已升级合同数
records: list = field(default_factory=list) # 每个合同的处理记录
# ═══════════════════════════════════════════════════════════════════════
# ContractImproveLoop —— Loop Engineering 核心类
# ═══════════════════════════════════════════════════════════════════════
class ContractImproveLoop:
"""合同智能完善循环系统
Loop Engineering 设计要素:
- Objective: 完善合同至四大维度全部达标
- Trigger: while True 定时扫描 pending_contracts/
- Discovery: scan() 扫描待处理合同
- Workspace: FilesystemBackend(virtual_mode=True) 沙盒隔离
- Context: checkpointer + thread_id 上下文持久化 + SKILL.md 知识
- Delegation: Editor(Maker-写) + Reviewer(Checker-只读)
- Verification: 编排者响应文本判断 + 文件状态兜底
- Budget: max_rounds 修正上限(默认3,可配置)
- Escalation: 循环耗尽 → Python 直接调用 escalate_contract
- Exit: 编排者调用 finalize_contract → 文件移至 completed/
"""
def __init__(self):
self.orchestrator = _build_orchestrator()
self.stats = ContractStats()
# ── 发现 ──────────────────────────────────────────────────
def scan(self) -> list[str]:
"""扫描待处理目录,返回合同 ID 列表"""
files = glob.glob(os.path.join(PENDING, "*.txt"))
return [os.path.basename(f).replace(".txt", "") for f in files]
# ── 单轮流式执行 ──────────────────────────────────────────
def _stream_round(self, contract_id: str, thread_id: str) -> str:
"""执行单轮 Editor→Reviewer,流式展示所有 Agent 的输出。
关键设计:
- 每轮发送完全相同的消息,不区分"第几轮""是否最后一轮"
- checkpointer + thread_id 自动保持上下文,LLM 从历史中知道当前处于什么阶段
- 收集编排者本轮产生的文本,返回给调用方用于判断是否完成
"""
# 消息每轮完全一致 —— 上下文由 checkpointer 自动携带
# 编排者从历史中知道:这是首轮→全面完善,还是续轮→针对性修正
message = f"请处理合同 {contract_id}。"
config = {
"configurable": {"thread_id": thread_id},
}
# 流式展示状态
subagent_names = {}
task_args_buf = ""
cur_tool_name = None
last_source = None
# 收集编排者本轮产生的文本(用于后续判断是否完成)
orchestrator_text = ""
try:
for chunk in self.orchestrator.stream(
input={"messages": [{"role": "user", "content": message}]},
config=config,
stream_mode="messages",
subgraphs=True,
version="v2",
):
if chunk.get("type") != "messages":
continue
token, _metadata = chunk["data"]
namespace = chunk.get("ns", ())
# 识别消息来源(主代理 or 子代理)
ns_sub_id = None
for seg in (namespace or ()):
if isinstance(seg, str) and seg.startswith("tools:"):
ns_sub_id = seg.replace("tools:", "")
break
if ns_sub_id:
if ns_sub_id in subagent_names:
source = f"子代理-{subagent_names[ns_sub_id]}"
else:
m = re.search(r'subagent_type["\s:]+([^",}\s]+)', task_args_buf)
name = m.group(1) if m else ns_sub_id[:8] + "..."
subagent_names[ns_sub_id] = name
source = f"子代理-{name}"
task_args_buf = ""
else:
source = "主代理(编排者)"
msg_type = getattr(token, 'type', 'unknown')
# ── 工具调用 ──
if hasattr(token, 'tool_call_chunks') and token.tool_call_chunks:
if last_source and (source != last_source):
print()
for tc in token.tool_call_chunks:
tc_name = tc.get('name')
if tc_name:
cur_tool_name = tc_name
icons = {
"task": "📋", "read_contract": "📖",
"improve_contract": "✏️", "finalize_contract": "✅",
"escalate_contract": "🆘",
}
print(f"\n [{source}] {icons.get(tc_name, '🔧')} 调用: {tc_name} ",
end="", flush=True)
if tc.get('args'):
args_str = tc['args']
if cur_tool_name == 'task' and source == "主代理(编排者)":
task_args_buf += args_str
print(args_str, end="", flush=True)
# ── 工具结果 ──
if msg_type == "tool":
if last_source:
print()
tool_name = getattr(token, 'name', '?')
result = str(getattr(token, 'content', ''))
if tool_name == 'improve_contract' and len(result) > 200:
result = result[:200] + "...[合同内容已写入]"
print(f" [{source}] 📋 返回:\n {result or '(无返回内容)'}")
last_source = None
continue
# ── AI 文本 ──
content_text = ""
if hasattr(token, 'content'):
c = token.content
if isinstance(c, str):
content_text = c
elif isinstance(c, list):
content_text = ''.join(
item.get('text', str(item)) if isinstance(item, dict)
else str(item) for item in c
)
elif c is not None:
content_text = str(c)
has_tool_calls = hasattr(token, 'tool_call_chunks') and token.tool_call_chunks
if content_text and not has_tool_calls:
if source != last_source:
if last_source:
print()
icons = {"editor": "✏️", "reviewer": "🔍"}
icon = "🎯"
for k, v in icons.items():
if k in source:
icon = v; break
print(f"\n [{source}] {icon} ", end="", flush=True)
print(content_text, end="", flush=True)
# ★ 收集编排者的文本(用于本轮结果判断)
if source == "主代理(编排者)":
orchestrator_text += content_text
last_source = source
print()
except Exception as e:
print(f"\n ❌ [错误] 流式处理异常: {e}")
import traceback; traceback.print_exc()
return orchestrator_text
# ── 结果判断 ──────────────────────────────────────────────
def _check_result(self, orchestrator_text: str, contract_id: str) -> str | None:
"""从编排者的响应文本中判断本轮结果。
设计原则:优先从模型输出判断(直接、可读),文件状态作为兜底验证。
返回值:
'completed' — 审查通过,合同已归档
'escalated' — 编排者主动升级
None — 审查未通过,需要继续下一轮
"""
# 方法1:从编排者的响应文本直接判断
# "审查通过" 或调用了 finalize_contract → 完成
# "escalate" 或 "升级人工" → 升级
if "审查通过" in orchestrator_text or "finalize_contract" in orchestrator_text:
return "completed"
if "escalate_contract" in orchestrator_text or "已升级人工" in orchestrator_text:
return "escalated"
# 方法2:文件系统兜底(编排者调用了工具但文本可能被截断的情况)
file_path = os.path.join(PENDING, f"{contract_id}.txt")
if not os.path.exists(file_path):
for d, s in [(COMPLETED, "completed"), (ESCALATED, "escalated")]:
if os.path.exists(os.path.join(d, f"{contract_id}.txt")):
return s
return None
# ── 处理单个合同 ──────────────────────────────────────────
def process_one(self, contract_id: str, max_rounds: int = 3) -> dict:
"""处理单个合同。
Loop Engineering 核心:Python for 循环控制重试
- 每轮发送完全相同的消息 → checkpointer 提供上下文
- 优先从编排者响应文本判断完成 → 文件状态作为兜底
- 假设max_rounds 设 10 但 1 轮就通过 → 立刻退出,不会空转
- 循环耗尽 → Python 直接调用 escalate_contract(不依赖编排者)
"""
thread_id = f"contract-{contract_id}"
file_path = os.path.join(PENDING, f"{contract_id}.txt")
# ── 显示合同概要 ──
print(f"\n {'='*60}")
print(f" 📄 {contract_id}")
try:
with open(file_path, "r", encoding="utf-8") as f:
preview = f.read()
lines = preview.strip().split("\n")
print(f" 📏 {len(lines)} 行, {len(preview)} 字符")
for line in lines[:3]:
print(f" {line[:80]}{'...' if len(line) > 80 else ''}")
if len(lines) > 3:
print(f" ... (共 {len(lines)} 行)")
except FileNotFoundError:
print(f" ⚠️ 文件不存在")
return {"contract_id": contract_id, "status": "not_found"}
print(f" {'='*60}")
# ═══════════════════════════════════════════════════════
# Python 代码控制循环(真正的 Loop Engineering)
# 每轮消息完全一样,LLM 从 checkpointer 的上下文中知道该做什么
# 从编排者响应文本判断完成
# ═══════════════════════════════════════════════════════
for round_num in range(1, max_rounds + 1):
print(f"\n ╔{'═'*50}╗")
print(f" ║ 🔄 Python 循环控制: 第 {round_num}/{max_rounds} 轮")
print(f" ╚{'═'*50}╝")
# 1.执行单轮 —— 每轮消息完全一致
orchestrator_text = self._stream_round(contract_id, thread_id)
# 2.从编排者响应文本判断本轮结果
result = self._check_result(orchestrator_text, contract_id)
if result == "completed":
print(f"\n 🎉 [Python 判断] 编排者确认审查通过 → 合同完善完成!")
self.stats.completed += 1
return {"contract_id": contract_id, "status": "completed", "rounds": round_num}
if result == "escalated":
print(f"\n 🆘 [Python 判断] 编排者主动升级 → 需人工介入")
self.stats.escalated += 1
return {"contract_id": contract_id, "status": "escalated", "rounds": round_num}
# 3.审查未通过 → Python 决定是否继续
if round_num < max_rounds:
print(f" 🔄 [Python 判断] 审查未通过 → 准备第 {round_num + 1} 轮..."
f"(上下文已持久化,编排者将从历史中获取反馈)")
# 4.Python 循环耗尽 → 代码直接升级,不依赖编排者做这个判断
print(f"\n 🆘 [Python 判断] {max_rounds} 轮修正仍未通过 → Python 代码直接升级人工")
escalate_contract.invoke({
"contract_id": contract_id,
"reason": (
f"经过 {max_rounds} 轮 Editor-Reviewer 修正循环后"
f"合同仍未达标,需人工介入。"
),
})
self.stats.escalated += 1
return {"contract_id": contract_id, "status": "escalated", "rounds": max_rounds}
# ── 单次循环 ──────────────────────────────────────────────
def run_one_cycle(self):
"""执行一次完整的扫描→处理循环"""
print(f"\n{'#'*60}")
print(f" 🔁 扫描循环 — {datetime.now().strftime('%H:%M:%S')}")
print(f"{'#'*60}")
contracts = self.scan()
self.stats.scanned = len(contracts)
print(f" 🔍 [发现] {len(contracts)} 个待完善合同")
if not contracts:
print(f" ⏳ 暂无合同,等待下次扫描...")
return
for cid in contracts:
fpath = os.path.join(PENDING, f"{cid}.txt")
try:
with open(fpath, "r", encoding="utf-8") as f:
first_line = f.readline().strip()
print(f" 📝 {cid}: {first_line[:70]}...")
except Exception:
print(f" 📝 {cid}: (无法读取)")
for cid in contracts:
result = self.process_one(cid,3)
self.stats.records.append(result)
# 输出本轮统计
self._report()
# ── 统计报告 ──────────────────────────────────────────────
def _report(self):
"""输出本轮统计"""
all_records = store.search(("improvements",))
edit_count = sum(1 for r in all_records if r.value.get("action") == "improved")
print(f"\n 📊 [统计] 修改{edit_count}次 | "
f"完成{self.stats.completed} | 升级{self.stats.escalated} | "
f"记忆{len(all_records)}条")
# ── 持续运行 ──────────────────────────────────────────────
def run_loop(self, interval: int = 10):
"""持续运行的监控循环(while True)"""
print("=" * 60)
print(" 🔄 Loop Engineering: 合同智能完善系统")
print(" 🐍 Python 代码控制循环 | checkpointer 上下文持久化")
print("=" * 60)
print(f" 📂 监控目录: {PENDING}")
print(f" ⏱️ 扫描间隔: {interval} 秒 | Ctrl+C 停止")
print(f" ✏️ Editor: contract-editor (读+写)")
print(f" 🔍 Reviewer: contract-reviewer (只读)")
print(f" 🎯 Orchestrator: 编排调度 (finalize + escalate)")
print(f" 💾 持久化: InMemorySaver + thread_id(每轮消息完全一致)")
print(f" 🔄 修正上限: 可配置(默认3轮)")
print("=" * 60)
cycle = 0
try:
while True:
cycle += 1
print(f"\n{'#'*60}")
print(f" 🔁 [监控循环 #{cycle}] {datetime.now().strftime('%H:%M:%S')}")
print(f"{'#'*60}")
self.run_one_cycle()
time.sleep(interval)
except KeyboardInterrupt:
print(f"\n\n 🛑 系统已停止。共运行 {cycle} 个监控循环。")
# ═══════════════════════════════════════════════════════════════════════
# 主入口
# ═══════════════════════════════════════════════════════════════════════
if __name__ == "__main__":
import argparse
parser = argparse.ArgumentParser(description="合同智能完善系统 (Loop Engineering)")
parser.add_argument("--interval", type=int, default=10, help="扫描间隔(秒)")
parser.add_argument("--max-rounds", type=int, default=3, help="最大修正轮数")
args = parser.parse_args()
loop = ContractImproveLoop()
loop.run_loop(interval=args.interval)
五大构建块在代码中的映射
此案例代码中,精确标注每个构建块的位置:
# ═══════════════════════════════════════════════════════════════
# ① 自动化/调度 —— 循环的心跳
# ═══════════════════════════════════════════════════════════════
class ContractImproveLoop:
def run_loop(self, interval: int = 10):
while True: # ← 构建块①
cycle += 1
self.run_one_cycle() # 扫描 → 发现 → 处理
time.sleep(interval) # 休息,继续
# ═══════════════════════════════════════════════════════════════
# ② 工作树/隔离 —— Agent 在哪安全操作?
# ═══════════════════════════════════════════════════════════════
orchestrator = create_deep_agent(
backend=FilesystemBackend(
root_dir=CUR_DIR, # ← 构建块②
virtual_mode=True, # 防止 Agent 跳出工作目录
),
...
)
# ═══════════════════════════════════════════════════════════════
# ③ 技能/知识 —— 处理任务的"操作手册"
# ═══════════════════════════════════════════════════════════════
orchestrator = create_deep_agent(
skills=[SKILLS], # ← 构建块③
# 指向 skills/contract-improver/SKILL.md
# 定义四大审查维度 + 标准合同模板 + 角色职责 + 退出条件
...
)
# ═══════════════════════════════════════════════════════════════
# ④ 连接器/插件 —— Agent 怎么操作外部系统?
# ═══════════════════════════════════════════════════════════════
@tool
def read_contract(contract_id: str) -> str: # ← 构建块④
"""读取合同文件"""
...
@tool
def improve_contract(contract_id: str, # ← 构建块④
improved_content: str) -> str:
"""写入完善后的合同(覆盖原文件)"""
...
@tool
def finalize_contract(contract_id: str, # ← 构建块④
comment: str) -> str:
"""归档已完成合同"""
...
@tool
def escalate_contract(contract_id: str, # ← 构建块④
reason: str) -> str:
"""升级人工处理"""
...
# ═══════════════════════════════════════════════════════════════
# ⑤ 子代理/协作 —— Maker-Checker 权限分离
# ═══════════════════════════════════════════════════════════════
contract_editor = SubAgent( # ← 构建块⑤ Maker
name="contract-editor",
tools=[read_contract, improve_contract], # 有写入权限!
...
)
contract_reviewer = SubAgent( # ← 构建块⑤ Checker
name="contract-reviewer",
tools=[read_contract], # 只有只读权限!
...
)
# ═══════════════════════════════════════════════════════════════
# + 记忆/状态 —— 跨轮/跨会话记住什么?
# ═══════════════════════════════════════════════════════════════
store = InMemoryStore() # ← 记忆层
# store.put() 记录每次 improve / finalize / escalate 操作
checkpointer=InMemorySaver() # ← 上下文持久化
# 同一 thread_id 的多次 stream() 调用自动保持上下文
Loop Engineering 案例二
案例业务
一个电商平台每天收到大量退款申请,涉及不同的订单状态:有的还未发货,有的已在运输途中,有的客户已经签收。不同状态对应不同的处理规则——直接退款、先拦截物流再退款、检查退货状态和7天窗口等。人工处理容易出错(金额算错、规则误用、忘记扣运费),且效率低下。
本案例通过 Loop Engineering 架构自动处理退款单:Maker 分析订单并给出方案(通知文案+退款金额),Checker 独立核验,Orchestrator 在审查通过后执行实际操作。循环修正机制确保每笔退款的金额和通知都准确无误。
业务逻辑

完整代码
"""
Loop Engineering — 智能退款处理循环系统
Python 代码控制循环 | checkpointer 上下文持久化 | 类化工程封装
═══════════════════════════════════════════════════════════════════════════
业务场景:
退款订单实时监控。Maker 负责确认通知文案和计算退款金额(只读),
Checker 负责核验 Maker 的方案是否正确(只读),
Orchestrator 仅在 Checker 审查通过后执行实际的通知发送和退款操作。
架构:
1 个主 Agent (Orchestrator) + 2 个子 Agent (Maker + Checker)
- Maker: 只读工具 → 分析订单、确认通知文案、计算退款金额 → 不执行
- Checker: 只读工具 → 验证 Maker 的方案 → 不执行
- Orchestrator: 写入工具 → 仅在 Checker 审查通过后执行实际操作
Loop Engineering 核心设计:
★ Python for 循环控制重试(代码决定,不是 prompt 决定)
★ checkpointer + thread_id 上下文持久化(每轮消息完全一致)
★ Maker-Checker-Orch 三层隔离(提案 → 审核 → 执行)
★ 类化工程封装(类比 contract_improve_loop 设计模式)
五大构建块:
① 自动化/调度 → while True 持续扫描退款队列
② 工作树/隔离 → FilesystemBackend(virtual_mode=True)
③ 技能/知识 → skills/refund-processor/SKILL.md
④ 连接器/插件 → query_refund_order / execute_refund / send_notification 等
⑤ 子代理/协作 → refund-maker(Maker-提案) + refund-checker(Checker-审核)
运行:
cd D:\CC备课\Loop Engineering\agent-loop-engineering\operations_loop
python operations_loop.py
"""
import os, sys, io, time, uuid, re
from datetime import datetime
from dataclasses import dataclass, field
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
from deepagents import create_deep_agent
from deepagents.middleware.subagents import SubAgent
from deepagents.backends import FilesystemBackend
from langchain.tools import tool
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.store.memory import InMemoryStore
from my_llm import deepseek_llm, deepseek_llm_flash
sys.stdout = io.TextIOWrapper(sys.stdout.buffer, encoding="utf-8", line_buffering=True)
# ═══════════════════════════════════════════════════════════════════════
# 路径配置
# ═══════════════════════════════════════════════════════════════════════
CUR_DIR = os.path.dirname(os.path.abspath(__file__))
SKILLS = os.path.join(CUR_DIR, "skills")
# ═══════════════════════════════════════════════════════════════════════
# 数据层 —— 模拟退款订单数据库
# ═══════════════════════════════════════════════════════════════════════
class Database:
"""模拟退款订单数据库"""
def __init__(self):
self.refund_orders = {
"REF-001": {
"order_id": "ORD-001", "product": "iPhone 15 Pro",
"amount": 8999, "order_status": "待发货",
"customer": "张三", "customer_level": "VIP铂金",
"logistics_no": None, "intercept_status": None,
"sign_date": None, "returned": False,
"refund_amount": None, "refund_executed": False,
"notification_sent": False, "notes": [],
},
# "REF-002": {
# "order_id": "ORD-002", "product": "AirPods Pro",
# "amount": 1899, "order_status": "已发货",
# "customer": "李四", "customer_level": "普通",
# "logistics_no": "SF1234567890", "intercept_status": None,
# "sign_date": None, "returned": False,
# "refund_amount": None, "refund_executed": False,
# "notification_sent": False, "notes": [],
# },
# "REF-003": {
# "order_id": "ORD-003", "product": "MacBook Pro",
# "amount": 12999, "order_status": "已签收",
# "customer": "王五", "customer_level": "黄金",
# "logistics_no": None, "intercept_status": None,
# "sign_date": "2026-06-12", "returned": True,
# "refund_amount": None, "refund_executed": False,
# "notification_sent": False, "notes": [],
# },
# "REF-004": {
# "order_id": "ORD-004", "product": "华为Mate 60 Pro",
# "amount": 6999, "order_status": "已签收",
# "customer": "赵六", "customer_level": "普通",
# "logistics_no": None, "intercept_status": None,
# "sign_date": "2026-06-14", "returned": False,
# "refund_amount": None, "refund_executed": False,
# "notification_sent": False, "notes": [],
# },
# "REF-005": {
# "order_id": "ORD-005", "product": "iPad Air",
# "amount": 4799, "order_status": "已签收",
# "customer": "孙七", "customer_level": "普通",
# "logistics_no": None, "intercept_status": None,
# "sign_date": "2026-06-01", "returned": False,
# "refund_amount": None, "refund_executed": False,
# "notification_sent": False, "notes": [],
# },
}
self.processing_log = []
db = Database()
# ═══════════════════════════════════════════════════════════════════════
# 记忆系统 —— 记录每次操作,用于统计报告
# ═══════════════════════════════════════════════════════════════════════
store = InMemoryStore()
# ═══════════════════════════════════════════════════════════════════════
# 工具定义 —— 按权限严格分层
# ═══════════════════════════════════════════════════════════════════════
# ── 只读工具(Maker + Checker 共用)────────────────────────────────
@tool
def query_refund_order(refund_id: str) -> str:
"""查询退款单详情(只读)。输入退款单号如 'REF-001',返回订单全貌。"""
rf = db.refund_orders.get(refund_id)
if not rf:
return f"退款单 {refund_id} 不存在"
info_parts = [
f"退款单号: {refund_id}",
f"原订单号: {rf['order_id']}",
f"商品: {rf['product']}",
f"订单金额: {rf['amount']}元",
f"订单状态: {rf['order_status']}",
f"客户: {rf['customer']}({rf['customer_level']})",
]
if rf["order_status"] == "已发货":
info_parts.append(f"物流单号: {rf['logistics_no']}")
info_parts.append(f"拦截状态: {rf.get('intercept_status') or '未拦截'}")
if rf["order_status"] == "已签收":
info_parts.append(f"签收日期: {rf['sign_date']}")
info_parts.append(f"商品是否退回: {'是' if rf['returned'] else '否'}")
info_parts.append(f"是否已退款: {'是' if rf['refund_executed'] else '否'}")
info_parts.append(f"是否已发通知: {'是' if rf['notification_sent'] else '否'}")
return "\n".join(info_parts)
@tool
def get_refund_policy(order_status: str) -> str:
"""查询退款规则(只读)。输入订单状态(待发货/已发货/已签收),返回对应规则。"""
policies = {
"待发货": (
"【待发货退款规则】\n"
"退款金额: 订单全额(无手续费)\n"
"通知模板: '您的订单{order_id}({product})退款申请已处理,"
"退款金额¥{amount},预计1-3个工作日到账。'\n"
"动作: 发送通知 + 执行退款 → 标记完成"
),
"已发货": (
"【已发货退款规则】\n"
"步骤1(未拦截): 发送物流拦截通知,模板: '您的订单{order_id}"
"({product})已申请退款,我们已通知物流拦截包裹,拦截成功后立即为您退款。'\n"
"步骤2(拦截成功): 执行退款,金额=订单全额(拦截成功无运费损失,不扣费)\n"
" 通知模板: '您的订单{order_id}({product})物流已拦截成功,"
"退款金额¥{amount},预计3-5个工作日到账。'"
),
"已签收": (
"【已签收退款规则】\n"
"步骤1: check_7day_window 检查是否超7天\n"
"未退回+≤7天: 不退款,发提醒退货通知,模板: '您的订单{order_id}"
"({product})退款申请已收到。请您先将商品寄回(运费由您承担¥15),"
"我们收到退货后将为您退款¥{amount-15}元。'\n"
"已退回+≤7天: 退款金额=订单金额-¥15,模板: '您的订单{order_id}"
"({product})退货已收到,退款金额¥{amount}(已扣除¥15运费),"
"预计5-7个工作日到账。'\n"
">7天(不论退回): 不退款,模板: '很抱歉,您的订单{order_id}签收已超过7天,"
"超出退货时效,无法办理退款。如有疑问请联系人工客服。'"
),
}
return policies.get(order_status, f"未知状态: {order_status},需人工处理")
@tool
def check_7day_window(refund_id: str) -> str:
"""检查7天退货窗口(只读)。已签收订单需检查是否在7天退货期内。"""
rf = db.refund_orders.get(refund_id)
if not rf:
return f"退款单 {refund_id} 不存在"
if rf["order_status"] != "已签收":
return f"订单状态为'{rf['order_status']}',无需检查7天窗口"
if not rf.get("sign_date"):
return "缺少签收日期,无法判断"
days = (datetime.now() - datetime.strptime(rf["sign_date"], "%Y-%m-%d")).days
if days <= 7:
return f"签收日期: {rf['sign_date']},距今{days}天 → 在7天退货期内,可办理退货退款"
else:
return f"签收日期: {rf['sign_date']},距今{days}天 → 已超出7天退货期,不可退款"
# ── 写入工具(只有 Orchestrator 能调用)───────────────────────────
@tool
def execute_refund(refund_id: str, amount: float, reason: str) -> str:
"""执行退款(写入,不可逆!)。只有编排者能调用,必须在Checker审查通过后。"""
rf = db.refund_orders.get(refund_id)
if not rf:
return f"退款失败:退款单 {refund_id} 不存在"
if rf["refund_executed"]:
return f"退款失败:退款单 {refund_id} 已退款,不可重复退款"
rid = f"RFD-{uuid.uuid4().hex[:8].upper()}"
rf["refund_amount"] = amount
rf["refund_executed"] = True
db.processing_log.append(f"[REFUND] {refund_id}: ¥{amount} ({reason}) 单号:{rid}")
store.put(("operations",), f"refund_{refund_id}_{datetime.now().timestamp():.0f}", {
"refund_id": refund_id, "action": "refund", "amount": amount,
"reason": reason, "refund_no": rid, "time": datetime.now().isoformat(),
})
return f"✅ 退款成功!退款单号: {rid},金额: ¥{amount},预计3-7个工作日到账。原因: {reason}"
@tool
def send_notification(refund_id: str, message: str) -> str:
"""发送客户通知(写入)。只有编排者能调用,必须在Checker审查通过后。"""
rf = db.refund_orders.get(refund_id)
if not rf:
return f"通知发送失败:退款单 {refund_id} 不存在"
customer = rf["customer"]
rf["notification_sent"] = True
db.processing_log.append(f"[NOTIFY] {refund_id} → {customer}: {message[:80]}...")
store.put(("operations",), f"notify_{refund_id}_{datetime.now().timestamp():.0f}", {
"refund_id": refund_id, "action": "notify", "customer": customer,
"message": message[:200], "time": datetime.now().isoformat(),
})
return f"✅ 已向{customer}发送通知: {message}"
@tool
def send_logistics_intercept(refund_id: str) -> str:
"""发送物流拦截请求(写入)。只有编排者能调用,必须在Checker审查通过后。"""
rf = db.refund_orders.get(refund_id)
if not rf:
return f"拦截失败:退款单 {refund_id} 不存在"
if rf["order_status"] != "已发货":
return f"拦截失败:订单状态为'{rf['order_status']}',不是'已发货',无法拦截"
rf["intercept_status"] = "已拦截"
db.processing_log.append(f"[INTERCEPT] {refund_id}: 物流{rf['logistics_no']}拦截成功")
store.put(("operations",), f"intercept_{refund_id}_{datetime.now().timestamp():.0f}", {
"refund_id": refund_id, "action": "intercept",
"logistics_no": rf["logistics_no"], "time": datetime.now().isoformat(),
})
return f"✅ 物流拦截成功!物流单号: {rf['logistics_no']},包裹将在途中退回。可执行退款(全额,不扣费)。"
@tool
def finalize_refund(refund_id: str, comment: str) -> str:
"""标记退款处理完成。只有编排者能调用。"""
rf = db.refund_orders.get(refund_id)
if not rf:
return f"归档失败:退款单 {refund_id} 不存在"
rf["notes"].append(f"[{datetime.now().isoformat()}] 完成: {comment}")
db.processing_log.append(f"[DONE] {refund_id}: {comment}")
store.put(("operations",), f"done_{refund_id}_{datetime.now().timestamp():.0f}", {
"refund_id": refund_id, "action": "completed", "comment": comment[:200],
"time": datetime.now().isoformat(),
})
return f"✅ 退款单 {refund_id} 处理完成: {comment}"
@tool
def escalate_to_human(refund_id: str, reason: str, priority: str = "normal") -> str:
"""升级给人工客服(写入)。只有编排者能调用。"""
rf = db.refund_orders.get(refund_id, {})
rf["notes"] = rf.get("notes", []) + [f"[升级] {reason}"]
db.processing_log.append(f"[ESCALATE] {refund_id}: {reason}(优先级:{priority})")
store.put(("operations",), f"escalate_{refund_id}_{datetime.now().timestamp():.0f}", {
"refund_id": refund_id, "action": "escalated", "reason": reason[:200],
"priority": priority, "time": datetime.now().isoformat(),
})
return f"🆘 退款单 {refund_id} 已升级人工处理,优先级: {priority},原因: {reason}"
# ═══════════════════════════════════════════════════════════════════════
# 子代理定义 —— Maker + Checker
# ═══════════════════════════════════════════════════════════════════════
refund_maker = SubAgent(
name="refund-maker",
description="退款处理专员(Maker):分析退款单→确认通知文案→计算退款金额。只读权限,不执行任何写入。",
system_prompt="""你是退款处理专员(Maker)。用 query_refund_order 查退款单,用 get_refund_policy 查规则,用 check_7day_window 检查时效。
═══════════════════════════════════════════════════════════
⚠️ 修改范围规则(必须严格遵守)
═══════════════════════════════════════════════════════════
★ 情况A:如果 task 描述中列出了具体的问题清单(如「1. xxx 2. xxx ...」)
→ 仅修正清单中列出的问题。逐条对照、逐条修正。
→ 问题清单中未提到的内容(通知文案、退款金额),一律保持原样不动。
→ 不要重新分析整笔退款单,只做靶向修正。
★ 情况B:如果 task 描述中没有列出具体问题(首轮全面处理)
→ 按照 SKILL.md 退款处理标准流程全面分析并给出方案。
═══════════════════════════════════════════════════════════
你的核心职责(注意:你只提案,不执行!)
═══════════════════════════════════════════════════════════
1. 用 query_refund_order 读取退款单,确定订单状态
2. 根据订单状态对照 SKILL.md / get_refund_policy 确定处理方式
3. 明确给出以下信息(缺一不可):
a) 通知内容:要发送给客户的完整通知文案
b) 退款金额:如需退款,明确写出金额和计算公式
c) 处理动作:告知编排者需要执行什么操作
("发送通知"/"发送通知+执行退款"/"发送物流拦截"/"仅发送通知不退款"等)
你不需要调用 execute_refund、send_notification、send_logistics_intercept——
这些由编排者在 Checker 审查通过后执行。你只需给出方案。
═══════════════════════════════════════════════════════════
退款金额计算速查表
═══════════════════════════════════════════════════════════
待发货 → 全额退款,不扣费
已发货+未拦截 → 不退款,先拦截(等下一轮拦截成功后再退款)
已发货+已拦截 → 全额退款,不扣费(拦截成功无运费损失)
已签收+未退回+≤7天 → 不退款,发退货提醒
已签收+已退回+≤7天 → 退款 = 订单金额 - 15
已签收+>7天 → 不退款,发超期通知
═══════════════════════════════════════════════════════════
🚫 安全红线(必须遵守——这是收敛的关键)
═══════════════════════════════════════════════════════════
❌ 红线1:修正时不得引入新的计算错误!
- 修正退款金额时不能算错扣费规则(只有已签收+已退回才扣¥15)
- 修正通知文案时金额数字必须与退款金额一致
❌ 红线2:不得破坏已正确的部分!
- 如果通知文案已正确,只修正被指出的问题,不要重写整条通知
- 修改只应减少问题,不应增加问题
❌ 红线3:必须区分"已发货未拦截"和"已发货已拦截"!
- 未拦截 → 先拦截,不能直接退款
- 已拦截 → 执行退款,全额不扣费
处理完成后:
- 清晰列出:(1)通知内容(2)退款金额(3)建议动作
- 然后立即结束,不要审查自己的方案(那是 Checker 的职责)""",
tools=[query_refund_order, get_refund_policy, check_7day_window],
# Maker 没有任何写入工具!它只提案,不执行。
model=deepseek_llm,
)
refund_checker = SubAgent(
name="refund-checker",
description="退款审查员(Checker):独立核验Maker的金额计算和通知文案。只有只读权限!",
system_prompt="""你是退款审查员(Checker)。只有只读权限,不能执行任何写入操作。
═══════════════════════════════════════════════════════════
🛑 最重要:Maker 的方案在 task 描述中,不在文件系统里!
═══════════════════════════════════════════════════════════
编排者委派你时,会把 Maker 的完整方案(通知文案 + 退款金额 + 建议动作)
直接写在 task 描述里发给你。你不需要去文件系统找任何"方案文件"!
❌ 绝对禁止:用 glob/ls/grep/read_file 搜索文件系统寻找 Maker 方案!
❌ 绝对禁止:说"我没有找到 Maker 的方案文件"——方案就在本条消息里!
✅ 正确做法:直接阅读编排者发给你的 task 描述,从中提取 Maker 的方案内容,
然后对照 SKILL.md 规则进行审查。只有当需要核实订单原始数据时,
才使用 query_refund_order / get_refund_policy / check_7day_window。
═══════════════════════════════════════════════════════════
🛑 其次阅读:以下情况一律不算问题,绝对不要报告!
═══════════════════════════════════════════════════════════
【通知话术类 —— 可以容忍】
- 通知文案的具体措辞差异(只要包含:订单号/商品名/金额/预计时间/操作说明)
- "元"和"¥"符号混用(只要金额数字正确)
- 敬语/礼貌用语的有无("亲爱的客户""祝您生活愉快"不是必须的)
- 文案长短详略差异(只要关键信息不遗漏)
⚠️ 核心原则:只要金额算对了、规则用对了、关键信息全了,措辞差异不算问题。
不要因为话术不够优美就报"审查不通过"。
═══════════════════════════════════════════════════════════
★ 真正需要检查的问题(只有这些才报告)
═══════════════════════════════════════════════════════════
对照 SKILL.md 逐项核对:
① 退款金额是否正确?
- 待发货 → 全额(不扣费)
- 已发货+已拦截 → 全额(不扣费)
- 已签收+已退回+≤7天 → 金额-15
- 其他情况 → 不应有退款金额
用 query_refund_order 获取订单金额后重新计算,与 Maker 给出的金额对比
② 通知文案是否包含必要信息?
- 必须包含:订单号、商品名
- 如有退款:必须包含退款金额和预计到账时间
- 如不退款:必须说明不退款原因
③ 处理规则是否适用正确?
- 订单状态判断是否正确
- 已签收是否检查了7天窗口
- 已发货是否正确区分了"未拦截"和"已拦截"
④ 动作建议是否与规则一致?
- 该拦截的是否建议了拦截
- 该退款的是否建议了退款
- 该仅通知的是否建议了仅通知
═══════════════════════════════════════════════════════════
★ 输出格式(必须严格遵守,便于编排者提取并传递给 Maker)
═══════════════════════════════════════════════════════════
如果 Maker 的方案全部正确,仅输出:
审查通过。
如果发现问题:
审查不通过。发现以下问题:
1. 【退款-金额错误】已签收已退回应扣¥15运费,退款金额应为XX元,Maker给出YY元 → 修正建议:退款金额 = 订单金额 - 15 = XX元
2. 【通知-信息缺失】退款通知缺少预计到账时间 → 修正建议:在通知末尾添加"预计X-X个工作日到账"
3. 【规则-状态误判】订单已签收但Maker按已发货规则处理 → 修正建议:应检查退货状态和7天窗口后按已签收规则重新处理
...
每条「修正建议」必须具体可操作,给出正确的数值或文案。
示例(正确):
审查不通过。发现以下问题:
1. 【退款-金额错误】REF-003已签收+已退回+≤7天,应扣¥15。12999-15=12984元,但Maker给出12999元 → 修正建议:退款金额改为12984元,通知中金额同步修正
示例(错误——禁止):
审查不通过。退款金额有问题,请重新计算。
⚠️ 约束:
- 给出审查结论后立即结束。不要写"修改后再来"等循环指令。""",
tools=[query_refund_order, get_refund_policy, check_7day_window],
# Checker 没有任何写入工具!它只审核,不执行。
model=deepseek_llm_flash,
)
# ═══════════════════════════════════════════════════════════════════════
# 编排者 Agent
# ═══════════════════════════════════════════════════════════════════════
def _build_orchestrator():
"""构建编排者 Agent(含 checkpointer 用于上下文持久化)"""
return create_deep_agent(
model=deepseek_llm,
tools=[execute_refund, send_notification, send_logistics_intercept,
finalize_refund, escalate_to_human],
subagents=[refund_maker, refund_checker],
skills=[SKILLS],
backend=FilesystemBackend(root_dir=CUR_DIR, virtual_mode=True),
store=store,
checkpointer=InMemorySaver(),
system_prompt="""你是退款编排员。你负责委派子代理和执行实际操作,不亲自分析订单或审查方案。
═══════════════════════════════════════════════════════════
🛑 启动规则:收到任务后直接委派 Maker,不要做任何准备工作!
═══════════════════════════════════════════════════════════
退款单数据在数据库中,只能通过 query_refund_order 工具访问
(该工具已分配给 Maker 和 Checker)。文件系统里没有退款单文件!
❌ 绝对禁止:用 glob/ls/read_file 搜索文件系统找退款数据!
❌ 绝对禁止:重新读取 SKILL.md(skills 已自动加载)!
❌ 绝对禁止:读取 operations_loop.py 文件!
❌ 绝对禁止:在委派 Maker 之前做任何"数据准备"查询!
✅ 正确做法:收到"请处理退款单 XXX"后,立即委派 Maker,
一步都不要多余。Maker 有自己的工具来查数据。
═══════════════════════════════════════════════════════════
⚠️ 核心规则:你必须确保 Checker 的审查意见被传递给 Maker
═══════════════════════════════════════════════════════════
这是系统中最重要的规则。违反它会导致问题无法收敛。
首轮(无历史时):委派 Maker 对退款单做全面分析并给出方案。
续轮(有历史时):你必须从上一轮的 Checker 输出中提取「所有问题清单」,
原样复制到 Maker 的 task 描述中。格式如下:
"请处理退款单 {refund_id},仅修正以下审查发现的具体问题,其他内容保持原样:
1. [从Checker输出中逐条复制的问题和修正建议]
2. [逐条复制...]
...
请逐一修正上述问题,不要修改未列出的内容,也不要引入新的问题。"
❌ 绝对禁止:只说"根据审查意见修改"而不列出具体问题。
✅ 正确做法:逐条列出 Checker 发现的所有问题,让 Maker 逐条修正。
═══════════════════════════════════════════════════════════
★ 每轮执行流程(恰好一轮 Maker→Checker,然后立即结束)
═══════════════════════════════════════════════════════════
1. 委派 refund-maker 分析退款单:
- 首轮:全面分析,给出通知文案 + 退款金额 + 建议动作
- 续轮:将上一轮 Checker 的具体问题清单原样写入 task 描述
2. 委派 refund-checker 独立审查 Maker 的方案
如果你不把 Maker 方案写入 task 描述,Checker 100% 会说"没收到方案"!
然后你被迫重新委派,浪费整整一轮。一步到位,绝不重来。
子代理之间是完全隔离的——Checker 的对话是全新的,看不到 Maker 刚才输出了什么。
你必须把 Maker 的完整方案直接写进 task 的 description 字段。模板如下(直接复制粘贴):
"请审查 Maker 针对退款单 {refund_id} 的以下方案是否完全正确:
【Maker给出的通知内容】
(逐字复制 Maker 输出的完整通知文案)
【Maker计算的退款金额】
(逐字复制 Maker 输出的金额和计算公式)
【Maker建议的执行动作】
(逐字复制 Maker 输出的建议动作)
请用 query_refund_order / get_refund_policy 核实数据后,对照 SKILL.md 逐项核验,
按标准格式输出审查结论。"
⚠️ 自检:task 的 description 字段如果短于 100 个字符,说明你没附方案,立即修正!
3. 根据审查结果决定下一步:
- "审查通过" → ★ 根据 Maker 的方案执行实际操作:
a) Maker 建议"发送通知" → 调用 send_notification
b) Maker 建议"执行退款" → 调用 execute_refund
c) Maker 建议"发送物流拦截" → 调用 send_logistics_intercept
d) Maker 建议"标记完成" → 调用 finalize_refund
e) Maker 建议"仅发送通知不退款"(如超期拒绝)→ 调用 send_notification + finalize_refund
执行完毕后仅输出"审查通过,已执行:[具体操作]",然后立即停止。
不要输出总结表格、处理摘要、markdown 等任何额外内容!
- "审查不通过" → 将 Checker 的全部问题和修正建议输出在响应中,
格式:"审查不通过。发现以下问题:\n1. xxx\n2. xxx\n...",
然后立即结束。外部程序会将这些交给下一轮的 Maker。
═══════════════════════════════════════════════════════════
工具权限说明
═══════════════════════════════════════════════════════════
- 委派子代理(名称映射——必须使用精确名称!):
· refund-maker = Maker(退款处理专员,只读-提案)
· refund-checker = Checker(退款审查员,只读-审核)
委派时 subagent_type 必须填 "refund-maker" 或 "refund-checker"
- execute_refund: ★ 仅在 Checker 审查通过 + Maker 建议退款时调用
- send_notification: ★ 仅在 Checker 审查通过 + Maker 建议发送通知时调用
- send_logistics_intercept: ★ 仅在 Checker 审查通过 + Maker 建议拦截物流时调用
- finalize_refund: 退款执行完毕或拒绝退款后调用
- escalate_to_human: 遇到无法自动修复的问题时升级
⚠️ 绝对不能跳过 Checker 审查直接执行写入操作!
只有 Checker 明确输出"审查通过"后,才能调用写入工具。
═══════════════════════════════════════════════════════════
硬性约束(违反将导致系统失控)
═══════════════════════════════════════════════════════════
- 禁止自行循环!审查不通过时,绝对不要再次委派 Maker。只汇报问题然后结束。
- 每次调用最多委派 Maker 一次、Checker 一次。
- 禁止未经 Checker 审查通过就调用写入工具!""",
)
# ═══════════════════════════════════════════════════════════════════════
# 统计数据结构
# ═══════════════════════════════════════════════════════════════════════
@dataclass
class RefundStats:
scanned: int = 0
discovered: int = 0
completed: int = 0
escalated: int = 0
records: list = field(default_factory=list)
# ═══════════════════════════════════════════════════════════════════════
# RefundLoop —— Loop Engineering 核心类
# ═══════════════════════════════════════════════════════════════════════
class RefundLoop:
"""智能退款处理循环系统
Loop Engineering 设计要素:
- Objective: 自动处理退款单,Maker提案→Checker审核→Orch执行
- Trigger: while True 定时扫描退款队列
- Discovery: scan() 扫描待处理退款单
- Workspace: FilesystemBackend(virtual_mode=True)
- Context: checkpointer + thread_id + SKILL.md
- Delegation: Maker(只读-提案) + Checker(只读-审核) + Orch(写-执行)
- Verification: Checker 审查输出 + Python 文本判断
- Budget: max_rounds(默认3)
- Escalation: 循环耗尽 → Python 直接 escalate
- Exit: 编排者调用 finalize_refund
"""
def __init__(self, max_rounds: int = 3):
self.orchestrator = _build_orchestrator()
self.stats = RefundStats()
self.max_rounds = max_rounds
# ── 发现 ──────────────────────────────────────────────────
def scan(self) -> list[str]:
"""扫描退款队列,返回待处理的退款单ID列表。"""
pending = []
for rid, rf in db.refund_orders.items():
already_done = any("完成" in n for n in rf.get("notes", []))
if not already_done:
pending.append(rid)
return pending
# ── 单轮流式执行 ──────────────────────────────────────────
def _stream_round(self, refund_id: str, thread_id: str) -> str:
"""执行单轮 Maker→Checker,流式展示所有 Agent 的输出。"""
message = f"请处理退款单 {refund_id}。"
config = {"configurable": {"thread_id": thread_id}}
subagent_names = {}
task_args_buf = ""
cur_tool_name = None
last_source = None
orchestrator_text = ""
try:
for chunk in self.orchestrator.stream(
input={"messages": [{"role": "user", "content": message}]},
config=config,
stream_mode="messages",
subgraphs=True,
version="v2",
):
if chunk.get("type") != "messages":
continue
token, _metadata = chunk["data"]
namespace = chunk.get("ns", ())
ns_sub_id = None
for seg in (namespace or ()):
if isinstance(seg, str) and seg.startswith("tools:"):
ns_sub_id = seg.replace("tools:", "")
break
if ns_sub_id:
if ns_sub_id in subagent_names:
source = f"子代理-{subagent_names[ns_sub_id]}"
else:
m = re.search(r'subagent_type["\s:]+([^",}\s]+)', task_args_buf)
name = m.group(1) if m else ns_sub_id[:8] + "..."
subagent_names[ns_sub_id] = name
source = f"子代理-{name}"
task_args_buf = ""
else:
source = "主代理(编排者)"
msg_type = getattr(token, 'type', 'unknown')
# ── 工具调用 ──
if hasattr(token, 'tool_call_chunks') and token.tool_call_chunks:
if last_source and (source != last_source):
print()
for tc in token.tool_call_chunks:
tc_name = tc.get('name')
if tc_name:
cur_tool_name = tc_name
icons = {
"task": "📋", "query_refund_order": "📖",
"get_refund_policy": "📜", "check_7day_window": "📅",
"execute_refund": "💳", "send_notification": "📨",
"send_logistics_intercept": "🚚",
"finalize_refund": "✅", "escalate_to_human": "🆘",
}
print(f"\n [{source}] {icons.get(tc_name, '🔧')} 调用: {tc_name} ",
end="", flush=True)
if tc.get('args'):
args_str = tc['args']
if cur_tool_name == 'task' and source == "主代理(编排者)":
task_args_buf += args_str
print(args_str, end="", flush=True)
# ── 工具结果 ──
if msg_type == "tool":
if last_source:
print()
tool_name = getattr(token, 'name', '?')
result = str(getattr(token, 'content', ''))
if tool_name in ('execute_refund', 'send_notification') and len(result) > 200:
result = result[:200] + "..."
print(f" [{source}] 📋 返回:\n {result or '(无返回内容)'}")
last_source = None
continue
# ── AI 文本 ──
content_text = ""
if hasattr(token, 'content'):
c = token.content
if isinstance(c, str):
content_text = c
elif isinstance(c, list):
content_text = ''.join(
item.get('text', str(item)) if isinstance(item, dict)
else str(item) for item in c
)
elif c is not None:
content_text = str(c)
has_tool_calls = hasattr(token, 'tool_call_chunks') and token.tool_call_chunks
if content_text and not has_tool_calls:
if source != last_source:
if last_source:
print()
icons = {"maker": "📝", "checker": "🔍"}
icon = "🎯"
for k, v in icons.items():
if k in source:
icon = v; break
print(f"\n [{source}] {icon} ", end="", flush=True)
print(content_text, end="", flush=True)
if source == "主代理(编排者)":
orchestrator_text += content_text
last_source = source
print()
except Exception as e:
print(f"\n ❌ [错误] 流式处理异常: {e}")
import traceback; traceback.print_exc()
return orchestrator_text
# ── 结果判断 ──────────────────────────────────────────────
def _check_result(self, orchestrator_text: str, refund_id: str) -> str | None:
"""从编排者的响应文本中判断本轮结果。"""
if "审查通过" in orchestrator_text or "finalize_refund" in orchestrator_text:
return "completed"
if "escalate_to_human" in orchestrator_text or "已升级人工" in orchestrator_text:
return "escalated"
recent_logs = db.processing_log[-3:]
for log in recent_logs:
if f"[DONE] {refund_id}" in log:
return "completed"
if f"[ESCALATE] {refund_id}" in log:
return "escalated"
return None
# ── 处理单个退款单 ──────────────────────────────────────────
def process_one(self, refund_id: str, max_rounds: int = 3) -> dict:
"""处理单个退款单。Python for 循环控制重试。"""
thread_id = f"refund-{refund_id}"
rf = db.refund_orders.get(refund_id, {})
print(f"\n {'='*60}")
print(f" 📋 {refund_id} | {rf.get('product', 'N/A')} | ¥{rf.get('amount', 'N/A')}")
print(f" 📦 状态: {rf.get('order_status', 'N/A')} | "
f"👤 客户: {rf.get('customer', 'N/A')}({rf.get('customer_level', 'N/A')})")
if rf.get("order_status") == "已签收":
print(f" 📅 签收: {rf.get('sign_date', 'N/A')} | "
f"📦 退回: {'是' if rf.get('returned') else '否'}")
if rf.get("order_status") == "已发货":
print(f" 🚚 拦截: {rf.get('intercept_status') or '未拦截'}")
print(f" {'='*60}")
for round_num in range(1, max_rounds + 1):
print(f"\n ╔{'═'*50}╗")
print(f" ║ 🔄 Python 循环控制: 第 {round_num}/{max_rounds} 轮")
print(f" ╚{'═'*50}╝")
orchestrator_text = self._stream_round(refund_id, thread_id)
result = self._check_result(orchestrator_text, refund_id)
if result == "completed":
print(f"\n 🎉 [Python 判断] 审查通过 + 操作已执行 → 退款单处理完成!")
self.stats.completed += 1
return {"refund_id": refund_id, "status": "completed", "rounds": round_num}
if result == "escalated":
print(f"\n 🆘 [Python 判断] 编排者主动升级 → 需人工介入")
self.stats.escalated += 1
return {"refund_id": refund_id, "status": "escalated", "rounds": round_num}
if round_num < max_rounds:
print(f" 🔄 [Python 判断] 审查未通过 → 准备第 {round_num + 1} 轮..."
f"(上下文已持久化,编排者将从历史中获取反馈)")
print(f"\n 🆘 [Python 判断] {max_rounds} 轮修正仍未通过 → Python 代码直接升级人工")
escalate_to_human.invoke({
"refund_id": refund_id,
"reason": f"经过 {max_rounds} 轮 Maker-Checker 处理循环后仍未解决,需人工介入。",
"priority": "high",
})
self.stats.escalated += 1
return {"refund_id": refund_id, "status": "escalated", "rounds": max_rounds}
# ── 单次循环 ──────────────────────────────────────────────
def run_one_cycle(self):
"""执行一次完整的扫描→处理循环"""
print(f"\n{'#'*60}")
print(f" 🔁 扫描循环 — {datetime.now().strftime('%H:%M:%S')}")
print(f"{'#'*60}")
pending = self.scan()
self.stats.scanned = len(db.refund_orders)
self.stats.discovered = len(pending)
print(f" 🔍 [发现] 扫描 {len(db.refund_orders)} 个退款单 → {len(pending)} 个待处理")
if not pending:
print(f" ⏳ 暂无待处理退款单,等待下次扫描...")
return
for rid in pending:
rf = db.refund_orders[rid]
extra = ""
if rf["order_status"] == "已发货":
extra = f" | 拦截:{rf.get('intercept_status') or '未拦截'}"
elif rf["order_status"] == "已签收":
extra = f" | 退回:{'是' if rf.get('returned') else '否'}"
print(f" 📝 {rid}: {rf['product']} | {rf['order_status']}{extra} | ¥{rf['amount']}")
print(f"\n [开始分派] 共 {len(pending)} 个退款单\n")
for i, rid in enumerate(pending):
print(f" ┌─ [{i+1}/{len(pending)}] ──────────────────────────────")
result = self.process_one(rid, max_rounds=self.max_rounds)
self.stats.records.append(result)
icon = "✅" if result["status"] == "completed" else "🆘"
print(f" └─ [{i+1}/{len(pending)}] {icon} "
f"{result['status']}({result['rounds']}轮)")
self._report()
# ── 统计报告 ──────────────────────────────────────────────
def _report(self):
all_records = store.search(("operations",))
refund_count = sum(1 for r in all_records if r.value.get("action") == "refund")
notify_count = sum(1 for r in all_records if r.value.get("action") == "notify")
intercept_count = sum(1 for r in all_records if r.value.get("action") == "intercept")
print(f"\n 📊 [统计] 退款{refund_count}次 | 拦截{intercept_count}次 | "
f"通知{notify_count}次 | 完成{self.stats.completed} | "
f"升级{self.stats.escalated} | 记忆{len(all_records)}条")
print(f" 操作日志(最近5条):")
for log in db.processing_log[-5:]:
print(f" {log}")
# ── 持续运行 ──────────────────────────────────────────────
def run_loop(self, interval: int = 10):
"""持续运行的监控循环(while True)"""
print("=" * 60)
print(" 🔄 Loop Engineering: 智能退款处理循环系统")
print(" 🐍 Python 代码控制循环 | checkpointer 上下文持久化")
print("=" * 60)
print(f" 📊 监控范围: 退款订单队列")
print(f" ⏱️ 扫描间隔: {interval} 秒 | Ctrl+C 停止")
print(f" 📝 refund-maker: 退款方案提案(只读)")
print(f" 🔍 refund-checker: 独立审查员(只读)")
print(f" 🎯 Orchestrator: 编排调度 + 执行操作(写)")
print(f" 💾 持久化: InMemorySaver + thread_id")
print(f" 🔄 修正上限: 可配置(默认3轮)")
print(f" 📂 Skills: {SKILLS}")
print("=" * 60)
cycle = 0
try:
while True:
cycle += 1
print(f"\n{'#'*60}")
print(f" 🔁 [监控循环 #{cycle}] {datetime.now().strftime('%H:%M:%S')}")
print(f"{'#'*60}")
self.run_one_cycle()
time.sleep(interval)
except KeyboardInterrupt:
print(f"\n\n 🛑 系统已停止。共运行 {cycle} 个监控循环。")
# ═══════════════════════════════════════════════════════════════════════
# 主入口
# ═══════════════════════════════════════════════════════════════════════
if __name__ == "__main__":
import argparse
parser = argparse.ArgumentParser(description="智能退款处理循环系统 (Loop Engineering)")
parser.add_argument("--interval", type=int, default=10, help="扫描间隔(秒)")
parser.add_argument("--max-rounds", type=int, default=3, help="最大修正轮数")
args = parser.parse_args()
loop = RefundLoop(max_rounds=args.max_rounds)
loop.run_loop(interval=args.interval)
五大构建块在代码中的映射
此案例代码中,精确标注每个构建块的位置:
# ═══════════════════════════════════════════════════════════════
# ① 自动化/调度 —— 循环的心跳
# ═══════════════════════════════════════════════════════════════
class RefundLoop:
def run_loop(self, interval: int = 10):
"""持续运行的监控循环(while True)"""
cycle = 0
while True: # ← 持续监控
cycle += 1
self.run_one_cycle() # 扫描 → 发现 → 处理 → 报告
time.sleep(interval)
def scan(self) -> list[str]:
"""扫描退款队列,返回待处理退款单ID"""
pending = []
for rid, rf in db.refund_orders.items():
already_done = any("完成" in n for n in rf.get("notes", []))
if not already_done:
pending.append(rid)
return pending
# ═══════════════════════════════════════════════════════════
# ② 工作树/隔离 —— Agent 在哪里安全操作
# ═══════════════════════════════════════════════════════════
orchestrator = create_deep_agent(
backend=FilesystemBackend(
root_dir=CUR_DIR, # ← 限定工作范围
virtual_mode=True, # ← 虚拟模式,防止误操作真实文件
),
...
)
# ═══════════════════════════════════════════════════════════
# ③ 技能/知识 —— 持久化的项目知识
# ═══════════════════════════════════════════════════════════
SKILLS = os.path.join(CUR_DIR, "skills")
orchestrator = create_deep_agent(
skills=[SKILLS], # ← 加载 skills/refund-processor/SKILL.md
...
)
# ═══════════════════════════════════════════════════════════
# ④ 连接器/插件 —— 按权限严格分层
# ═══════════════════════════════════════════════════════════
# ── 只读工具(Maker + Checker 共用)────────────────────
@tool
def query_refund_order(refund_id: str) -> str: # ← 查询退款单
@tool
def get_refund_policy(order_status: str) -> str: # ← 查询退款规则
@tool
def check_7day_window(refund_id: str) -> str: # ← 检查7天退货窗口
# ── 写入工具(只有 Orchestrator 能调用)───────────────
@tool
def execute_refund(refund_id, amount, reason): # ← 执行退款(不可逆!)
@tool
def send_notification(refund_id, message): # ← 发送客户通知
@tool
def send_logistics_intercept(refund_id): # ← 发送物流拦截
@tool
def finalize_refund(refund_id, comment): # ← 标记完成
@tool
def escalate_to_human(refund_id, reason, priority): # ← 升级人工
# ═══════════════════════════════════════════════════════════
# ⑤ 子代理/协作 —— Maker-Checker-Orch 三层分离
# ═══════════════════════════════════════════════════════════
# Maker:只读,只提案,不执行
refund_maker = SubAgent(
name="refund-maker",
tools=[query_refund_order, get_refund_policy, check_7day_window],
# 没有任何写入工具!
model=deepseek_llm, # ← 强模型做提案
)
# Checker:只读,只审核,不执行
refund_checker = SubAgent(
name="refund-checker",
tools=[query_refund_order, get_refund_policy, check_7day_window],
# 没有任何写入工具!
model=deepseek_llm_flash, # ← 弱模型做审核,降成本
)
# Orchestrator:持有所有写入工具,但只在 Checker 通过后才执行
orchestrator = create_deep_agent(
tools=[execute_refund, send_notification,
send_logistics_intercept, finalize_refund,
escalate_to_human], # ← 所有写入工具
subagents=[refund_maker, refund_checker],
checkpointer=InMemorySaver(),
...
)
# ═══════════════════════════════════════════════════════════
# 记忆系统 — 操作日志记录
# ═══════════════════════════════════════════════════════════
# 操作记忆:每次 execute_refund / send_notification / send_logistics_intercept
# / finalize_refund / escalate_to_human 调用后自动写入 store
store.put(("operations",), f"refund_{refund_id}_{timestamp}", {
"refund_id": refund_id, "action": "refund", "amount": amount, ...
})
# 短期记忆:checkpointer 跨轮次持久化(InMemorySaver)
checkpointer=InMemorySaver()
# 统计报告时从 store 中检索操作记录
all_records = store.search(("operations",))