Skip to content

Commit 765acbd

Browse files
committed
H3 切片 1:方案阶段 Proposer×3 并行提议 + 一次归约 + G1 三选(多智能体并行提议首次进生产;agents/worker 330 + backend/api 310 = 640 全绿,真 PG 311)
ModelPlanningNode 有 Supervisor 时按三个固定视角(机理建模 / 数据驱动 / 运筹优化)各派一个 proposer:<view> 子代理(readonly、预算切片 governor.subagent_slice()、spawn/result 双审计落 TOOL_CALLED),节点内 ThreadPoolExecutor 线程 fan-out、pool.map 保证提案按视角顺序留档与完成顺序无关;多样性靠固定视角不靠温度。归约按拍板走一次 model_planning.reduce 调用(而非设计原文的代码去重):≥2 路成功即归约为 plans[](id A/B/C、role primary/baseline/fallback、source_views、fallback 必带 fallback_condition)+ recommended + rationale + dropped[] + progress_note,归约不合法(id 重复 / 推荐悬空 / 五键空)节点失败,归约调用失败降级为按视角顺序直列候选并记警告;quorum:只剩 1 路 → 单案降级记警告且 G1 标题点明「N 路视角提议未成功」,0 路 → 失败交引擎重试。预算硬停不可被 quorum 吞掉:提议 runner 里的 AgentError E31x/E32x 留住异常对象、fan-out 收束后原样抛出,与单次调用路径同走 _BudgetGuardedNode 出「[E310] …」(test_budget_guard 逼出来的,否则三路各报 exhausted 被当软失败、错误码与追加通道指引全丢)。 G1 三选经 review_meta 声明(第四个消费方):options = approve(推荐案,id 保留 → 既有 e2e / 金轨迹 / 前端 CTA 预选零改)+ adopt:<id> 每一备选(label 带角色 blurb)+ reject,impact 带 plans / proposers{succeeded,failed} / dropped;后端零改动——resolve_approval 对任意非 reject / 非 redo 的 option_id 都进 review_decisions 台账,chosen_plan(planning, review_decisions) 台账 adopt:X 优先 → recommended → 首案,实验 / 验证 ×2 / 论文 / 验证材料五处调用全部传台账;前端零改动——审批卡正向选项 >1 时自动摆单选组(修订门同一机制),CDP 无头走查 10/10。契约 plan-proposal 零改动(plans 本就 1~N),投影只取五键,role / source_views / dropped / proposals 三路原样留在节点 outputs。 控制面并发锁是本刀真正的坑:三路提议从三个线程同时调用模型端口,而 EngineLlmPort 的 on_event(append_event + session.commit())/ on_usage(record_usage)/ 预算 check·charge 与 Supervisor 审计(同一 _process_event)都落在同一个 SQLAlchemy Session 与同一份账本上,此前只在 RunnerThread 串行触发。EngineLlmPort(lock=) 一把 RLock 罩住全部记账 / 事件回调(HTTP 调用不上锁仍并行),engine_glue 把同一把锁交给 _process_event(progress / audit 共用);worker runtime 一把锁罩住 engine.record_external(tools recorder + audit)。无监督者(旧装配 / 单节点测试)或 proposer_views=() 走 v3.21 单次调用 model_planning.default 路径逐字节不变;两个新 prompt id 进 API _REQUIRED_PROMPTS / _PROMPT_NODE_IDS 与 worker REQUIRED_PROMPT_IDS(缺一整链回落 SIM)。 测试:skills +10(并行屏障坐实三路同时在飞 + 双审计 / quorum 2-of-3 / 单案降级 / 全失败 / 归约失败降级 / 归约不合法 / 预算硬停透传 / 无人值守 / 无视角走旧路 / 台账选案);evals +3(worker 金样六条提议审计落在方案步骤区间内;全链 G1 摆出每一案、选 adopt:B 后实验任务卡与论文材料都按 B;一路提议失败 quorum 照常到闸门、提议 4 次归约 1 次);worker 用例桩与审计断言随之更新(三路先后随线程调度、按集合 + 相位断言);backend 主 e2e 加 G1 三选 / 投影五键 / 用量 6 条 / run.log 六条提议审计断言 + 新增 adopt:B 全链用例;e320 用例上限 50 → 100(方案阶段 4 次调用第 4 次前累计 90 不越线、论文第 5 次前 120 越线,语义不变)。export_openapi --check OK(无契约模型改动)。如实备案:选 adopt:B 后 CTA 文案仍是「确认 Agent 当前方案并继续」、方案页行 radio 与审批单选是两套控件、推荐项预选进 CTA 与设计原文「不预填」有差异(前端 _preferred_option 既有行为),三项均记入设计文档待下刀。
1 parent d71a342 commit 765acbd

18 files changed

Lines changed: 1370 additions & 54 deletions

agents/evals/src/omm_agent_evals/__init__.py

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -43,11 +43,14 @@
4343
from .scenario import (
4444
CANNED_ANALYSIS,
4545
CANNED_PLANNING,
46+
CANNED_PROPOSALS_BY_VIEW,
47+
CANNED_REDUCE,
4648
EXPERIMENT_CODE,
4749
GOLDEN_EVENT_TYPES,
4850
PROBLEM_STATEMENT,
4951
build_llm,
5052
build_runtime,
53+
canned_proposer,
5154
)
5255

5356
__all__ = [
@@ -59,6 +62,8 @@
5962
"CANNED_PAPER_OUTLINE",
6063
"CANNED_PLANNING",
6164
"CANNED_PREPARATION",
65+
"CANNED_PROPOSALS_BY_VIEW",
66+
"CANNED_REDUCE",
6267
"CANNED_ROBUSTNESS",
6368
"CANNED_VALIDATION",
6469
"CANNED_VALIDATION_CODE",
@@ -81,6 +86,7 @@
8186
"build_llm",
8287
"build_runtime",
8388
"canned_paper_section",
89+
"canned_proposer",
8490
"canned_sandbox_agent",
8591
"robustness_success",
8692
"sandbox_failure",

agents/evals/src/omm_agent_evals/full_chain.py

Lines changed: 26 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,7 @@
2323
from __future__ import annotations
2424

2525
import json
26-
from collections.abc import Callable, Sequence
26+
from collections.abc import Callable, Mapping, Sequence
2727
from dataclasses import dataclass, field
2828
from typing import Any
2929

@@ -56,7 +56,13 @@
5656
)
5757
from omm_agent_tools import failure_detail, summarize
5858

59-
from .scenario import CANNED_ANALYSIS, CANNED_PLANNING, PROBLEM_STATEMENT
59+
from .scenario import (
60+
CANNED_ANALYSIS,
61+
CANNED_PLANNING,
62+
CANNED_REDUCE,
63+
PROBLEM_STATEMENT,
64+
canned_proposer,
65+
)
6066

6167
# -- canned answers for the four stages scenario.py does not stub -------------
6268
# Same problem domain as CANNED_ANALYSIS/CANNED_PLANNING (freight-volume
@@ -455,7 +461,9 @@ def reply(messages: list[dict[str, str]]) -> str:
455461
return reply
456462

457463

458-
def build_full_chain_llm() -> StubLlmPort:
464+
def build_full_chain_llm(
465+
overrides: Mapping[str, str | Callable[[dict[str, Any]], str]] | None = None,
466+
) -> StubLlmPort:
459467
"""Stub responses for every prompt id the six real nodes may consult.
460468
461469
Two channels, mirroring production: template calls (``complete``) for the
@@ -466,18 +474,24 @@ def build_full_chain_llm() -> StubLlmPort:
466474
467475
The paper stage is a multipass pipeline (outline → sections → finalize);
468476
``paper_writing.default`` stays stubbed so the single-call fallback path
469-
remains exercisable from evals.
477+
remains exercisable from evals. ``overrides`` replaces individual template
478+
stubs (e.g. a proposer that fails for one view to drive the quorum path).
470479
"""
471480
return StubLlmPort(
472481
{
473482
"problem_analysis.default": stub_response(CANNED_ANALYSIS, fenced=True),
474483
"data_preparation.default": stub_response(CANNED_PREPARATION),
484+
# 方案阶段(H3):三路 Proposer 并行 + 一次归约;default 只在无监督者
485+
# 的装配里被消费,本会话有监督者,留着是让回落路径仍可从评测触达
475486
"model_planning.default": stub_response(CANNED_PLANNING),
487+
"model_planning.proposer": canned_proposer,
488+
"model_planning.reduce": stub_response(CANNED_REDUCE),
476489
"validating.default": stub_response(CANNED_VALIDATION),
477490
"paper_outline.default": stub_response(CANNED_PAPER_OUTLINE),
478491
"paper_section.default": canned_paper_section,
479492
"paper_finalize.default": stub_response(CANNED_PAPER_FINALIZE),
480493
"paper_writing.default": stub_response(CANNED_PAPER),
494+
**dict(overrides or {}),
481495
},
482496
chat_scripts={
483497
ExperimentExecutionNode.prompt_id: [
@@ -581,12 +595,17 @@ def build_full_chain_session(
581595

582596
#: Template (``complete``) prompt ids in stage order. The experiment stage is
583597
#: absent by design: it is a sandbox agent now, driven through ``chat_text``
584-
#: conversations (see :data:`FULL_CHAIN_CHAT_SEQUENCE`). The paper stage is a
585-
#: multipass pipeline: outline → one call per chapter → finalize.
598+
#: conversations (see :data:`FULL_CHAIN_CHAT_SEQUENCE`). The planning stage is
599+
#: three parallel proposers (same template id, so the recorded order is stable
600+
#: whichever thread lands first) followed by one reduce call. The paper stage
601+
#: is a multipass pipeline: outline → one call per chapter → finalize.
586602
FULL_CHAIN_PROMPT_SEQUENCE = [
587603
"problem_analysis.default",
588604
"data_preparation.default",
589-
"model_planning.default",
605+
"model_planning.proposer",
606+
"model_planning.proposer",
607+
"model_planning.proposer",
608+
"model_planning.reduce",
590609
"validating.default",
591610
"paper_outline.default",
592611
"paper_section.default",

agents/evals/src/omm_agent_evals/scenario.py

Lines changed: 53 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -73,6 +73,48 @@
7373
"rationale": "数据量小且趋势近线性,方案 A 可解释性与评审友好度更高",
7474
}
7575

76+
#: 方案阶段(H3):三视角 Proposer 并行各回一案,归约桩把它们收成 CANNED_PLANNING
77+
#: 的 A/B 两案(多出归约字段 role / source_views;投影只取契约五键)。
78+
CANNED_PROPOSALS_BY_VIEW = {
79+
"机理建模": {
80+
"name": "库存-运量动力学",
81+
"approach": "把运量当作受季节因子驱动的一阶动力学过程",
82+
"steps": ["辨识季节因子", "拟合动力学参数", "外推下季度"],
83+
"risks": ["外部冲击不可解释"],
84+
"fit": "运量趋势近线性,机理项可退化为线性趋势",
85+
},
86+
"数据驱动": {
87+
"name": "时间序列 + 启发式",
88+
"approach": "ARIMA 预测 + 遗传算法搜索配置",
89+
"steps": ["定阶建模", "编码搜索", "对比验证"],
90+
"risks": ["样本过短导致过拟合"],
91+
"fit": "样本较短,需谨慎定阶",
92+
},
93+
"运筹优化": {
94+
"name": "线性回归 + 整数规划",
95+
"approach": "最小二乘拟合运量趋势,再以 MILP 求最优车辆配置",
96+
"steps": ["拟合线性趋势", "构建配置模型", "灵敏度分析"],
97+
"risks": ["非线性冲击场景失效"],
98+
"fit": "车辆数与预算都是硬约束,整数规划直接可解",
99+
},
100+
}
101+
102+
103+
def canned_proposer(variables: dict) -> str:
104+
"""提议人桩:按 view_name 回对应视角的方案(callable 桩拿到的是渲染变量)。"""
105+
return stub_response(CANNED_PROPOSALS_BY_VIEW[variables["view_name"]])
106+
107+
108+
CANNED_REDUCE = {
109+
**CANNED_PLANNING,
110+
"plans": [
111+
{**CANNED_PLANNING["plans"][0], "role": "primary", "source_views": ["operations_research"]},
112+
{**CANNED_PLANNING["plans"][1], "role": "baseline", "source_views": ["data_driven"]},
113+
],
114+
"dropped": ["机理建模:动力学项退化为线性趋势,已并入方案 A"],
115+
"progress_note": "三路提议归约为两案:推荐线性回归 + 整数规划,时间序列作对照基线。",
116+
}
117+
76118
#: Real python executed in the sandbox: closed-form least squares on y≈2x+1.
77119
EXPERIMENT_CODE = """\
78120
import json
@@ -172,7 +214,9 @@ def build_llm() -> StubLlmPort:
172214
return StubLlmPort(
173215
{
174216
"problem_analysis.default": stub_response(CANNED_ANALYSIS, fenced=True),
175-
"model_planning.default": stub_response(CANNED_PLANNING),
217+
# worker 运行时注入了子代理监督者:方案阶段走三路提议 + 归约
218+
"model_planning.proposer": canned_proposer,
219+
"model_planning.reduce": stub_response(CANNED_REDUCE),
176220
}
177221
)
178222

@@ -210,6 +254,14 @@ def build_runtime(root: Path, require_confirmation: bool = True) -> WorkerRuntim
210254
EventType.STEP_SUCCEEDED,
211255
EventType.STATE_CHANGED, # -> MODEL_PLANNING
212256
EventType.STEP_STARTED,
257+
# 三路 Proposer 子代理并行:每路 spawn + result 两条监督者审计(六条 TOOL_CALLED,
258+
# 三路的先后随线程调度而变,事件类型序列不变)
259+
EventType.TOOL_CALLED,
260+
EventType.TOOL_CALLED,
261+
EventType.TOOL_CALLED,
262+
EventType.TOOL_CALLED,
263+
EventType.TOOL_CALLED,
264+
EventType.TOOL_CALLED,
213265
EventType.STEP_SUCCEEDED,
214266
EventType.REVIEW_REQUESTED,
215267
EventType.REVIEW_RESOLVED, # user approves plan A

agents/evals/tests/test_end_to_end.py

Lines changed: 52 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -72,10 +72,13 @@ def test_tool_call_is_recorded_inside_experiment_step(runtime):
7272
run_id = drive(runtime)
7373
events = runtime.events.load(run_id)
7474

75-
tool_events = [e for e in events if e.event_type is EventType.TOOL_CALLED]
75+
tool_events = [
76+
e
77+
for e in events
78+
if e.event_type is EventType.TOOL_CALLED and e.payload["tool"] == "python_run"
79+
]
7680
assert len(tool_events) == 1
7781
tool_event = tool_events[0]
78-
assert tool_event.payload["tool"] == "python_run"
7982
assert tool_event.payload["status"] == "succeeded"
8083

8184
snapshot = runtime.get_snapshot(run_id)
@@ -101,6 +104,53 @@ def test_tool_call_is_recorded_inside_experiment_step(runtime):
101104
assert types_by_seq[tool_event.seq] is EventType.TOOL_CALLED
102105

103106

107+
def test_planning_fanout_audits_sit_inside_the_planning_step(runtime):
108+
"""三路 Proposer 子代理的 spawn / result 审计都落在方案步骤区间内(H3)。"""
109+
run_id = drive(runtime)
110+
events = runtime.events.load(run_id)
111+
snapshot = runtime.get_snapshot(run_id)
112+
113+
proposer_events = [
114+
e
115+
for e in events
116+
if e.event_type is EventType.TOOL_CALLED
117+
and e.payload["tool"].startswith("subagent:proposer:")
118+
]
119+
assert sorted(e.payload["tool"] for e in proposer_events if e.payload["phase"] == "spawn") == [
120+
"subagent:proposer:data_driven",
121+
"subagent:proposer:mechanism",
122+
"subagent:proposer:operations_research",
123+
]
124+
results = [e for e in proposer_events if e.payload["phase"] == "result"]
125+
assert [e.payload["envelope_status"] for e in results] == ["done"] * 3
126+
127+
planning_step = next(
128+
step for step in snapshot.steps if step.state is TaskState.MODEL_PLANNING
129+
)
130+
started_seq = next(
131+
e.seq
132+
for e in events
133+
if e.event_type is EventType.STEP_STARTED
134+
and e.payload["step_id"] == planning_step.step_id
135+
)
136+
succeeded_seq = next(
137+
e.seq
138+
for e in events
139+
if e.event_type is EventType.STEP_SUCCEEDED
140+
and e.payload["step_id"] == planning_step.step_id
141+
)
142+
assert all(started_seq < e.seq < succeeded_seq for e in proposer_events)
143+
# 归约后的方案卡带角色,投影只取契约五键;被并入 A 的机理提议留在 dropped 里
144+
planning = snapshot.outputs[TaskState.MODEL_PLANNING.value]
145+
assert [plan["role"] for plan in planning["plans"]] == ["primary", "baseline"]
146+
assert planning["dropped"] == ["机理建模:动力学项退化为线性趋势,已并入方案 A"]
147+
assert [proposal["view"] for proposal in planning["proposals"]] == [
148+
"mechanism",
149+
"data_driven",
150+
"operations_research",
151+
]
152+
153+
104154
def test_artifacts_exist_on_disk_with_matching_checksums(runtime):
105155
run_id = drive(runtime)
106156
snapshot = runtime.get_snapshot(run_id)

agents/evals/tests/test_full_chain.py

Lines changed: 93 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,12 +26,15 @@
2626
from omm_agent_evals import (
2727
CANNED_EXPERIMENT_CODE,
2828
CANNED_PAPER,
29+
CANNED_PROPOSALS_BY_VIEW,
30+
CANNED_REDUCE,
2931
CANNED_VALIDATION_CODE,
3032
FULL_CHAIN_CHAT_SEQUENCE,
3133
FULL_CHAIN_GOLDEN_EVENT_TYPES,
3234
FULL_CHAIN_METRICS,
3335
FULL_CHAIN_PROMPT_SEQUENCE,
3436
FULL_CHAIN_ROBUSTNESS_CHECKS,
37+
build_full_chain_llm,
3538
build_full_chain_session,
3639
robustness_success,
3740
sandbox_failure,
@@ -42,6 +45,7 @@
4245
G3_ACCEPT_OPTION_ID,
4346
G4_CONFIRM_OPTION_ID,
4447
PYTHON_TOOL_NAME,
48+
stub_response,
4549
)
4650

4751

@@ -375,6 +379,95 @@ def test_review_rejection_then_retry_replans_and_completes():
375379
assert_replay_matches(session)
376380

377381

382+
# -- 4b. G1 三选(H3):归约出的每一案都是选项,选 B 就按 B 往下做 ----------------
383+
384+
385+
def g1_gate_event(session):
386+
gates = [
387+
event for event in session.sink.events
388+
if event.event_type is EventType.REVIEW_REQUESTED
389+
and (event.payload.get("gate") or {}).get("gate") == "G1"
390+
]
391+
return gates[-1]
392+
393+
394+
def test_g1_gate_offers_every_reduced_plan_and_adopting_b_flows_downstream():
395+
session = build_full_chain_session()
396+
engine, snapshot = session.engine, session.snapshot
397+
398+
outcome = engine.run_until_blocked(snapshot)
399+
assert outcome.status == AdvanceOutcome.REVIEW_REQUESTED
400+
401+
gate = g1_gate_event(session).payload["gate"]
402+
assert gate["decision_type"] == "confirm_plan"
403+
assert [option["id"] for option in gate["options"]] == ["approve", "adopt:B", "reject"]
404+
assert [o["id"] for o in gate["options"] if o.get("recommended")] == ["approve"]
405+
assert gate["title"] == (
406+
"请确认建模方案:推荐 A「线性回归 + 整数规划」;备选 B「时间序列 + 启发式」"
407+
)
408+
assert gate["impact"]["proposers"] == {
409+
"succeeded": ["mechanism", "data_driven", "operations_research"],
410+
"failed": [],
411+
}
412+
assert gate["impact"]["dropped"] == CANNED_REDUCE["dropped"]
413+
planning = snapshot.outputs[TaskState.MODEL_PLANNING.value]
414+
assert [plan["role"] for plan in planning["plans"]] == ["primary", "baseline"]
415+
assert planning["quality_warnings"] == []
416+
417+
# 用户改选 B:决策进台账,实验任务卡 / 论文材料都按 B 走
418+
engine.resolve_review(snapshot, approved=True, reason="adopt:B")
419+
outcome = confirm_delivery(session, engine.run_until_blocked(snapshot))
420+
assert outcome.status == AdvanceOutcome.COMPLETED
421+
assert snapshot.review_decisions[TaskState.MODEL_PLANNING.value] == "adopt:B"
422+
experiment_system = next(
423+
call.messages[0]["content"]
424+
for call in session.llm.chat_calls
425+
if call.label == "experiment_code.sandbox"
426+
)
427+
assert '"id": "B"' in experiment_system and "时间序列 + 启发式" in experiment_system
428+
assert '"id": "A"' not in experiment_system
429+
paper_calls = [call for call in session.llm.calls if call.prompt_id == "paper_outline.default"]
430+
assert "时间序列 + 启发式" in paper_calls[0].variables["chosen_plan"]
431+
432+
assert_replay_matches(session)
433+
434+
435+
def test_planning_quorum_one_failed_proposer_still_reaches_the_gate_with_a_warning():
436+
def flaky_proposer(variables):
437+
if variables["view_name"] == "数据驱动":
438+
return "(该视角的模型抽风了,两次都不是 JSON)"
439+
return stub_response(CANNED_PROPOSALS_BY_VIEW[variables["view_name"]])
440+
441+
session = build_full_chain_session(
442+
llm=build_full_chain_llm({"model_planning.proposer": flaky_proposer})
443+
)
444+
engine, snapshot = session.engine, session.snapshot
445+
446+
outcome = engine.run_until_blocked(snapshot)
447+
assert outcome.status == AdvanceOutcome.REVIEW_REQUESTED
448+
449+
# quorum:2/3 成功照常归约;缺席的那一路点名进警告、卡片与归约输入
450+
planning = snapshot.outputs[TaskState.MODEL_PLANNING.value]
451+
[failure] = planning["proposer_failures"]
452+
assert failure.startswith("视角「数据驱动」未成功(failed")
453+
assert planning["quality_warnings"] == [failure]
454+
assert [proposal["view"] for proposal in planning["proposals"]] == [
455+
"mechanism", "operations_research",
456+
]
457+
gate = g1_gate_event(session).payload["gate"]
458+
assert gate["title"].endswith(";1 路视角提议未成功")
459+
assert gate["impact"]["proposers"]["failed"] == [failure]
460+
prompt_ids = [call.prompt_id for call in session.llm.calls]
461+
# 失败那一路用掉一次修复重试:3 + 1 次提议调用,归约仍只一次
462+
assert prompt_ids.count("model_planning.proposer") == 4
463+
assert prompt_ids.count("model_planning.reduce") == 1
464+
465+
engine.resolve_review(snapshot, approved=True, reason="approve")
466+
outcome = confirm_delivery(session, engine.run_until_blocked(snapshot))
467+
assert outcome.status == AdvanceOutcome.COMPLETED
468+
assert_replay_matches(session)
469+
470+
378471
# -- 5. experiment failure + retry recovery ---------------------------------------
379472

380473

0 commit comments

Comments
 (0)