713 lines
26 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

#!/usr/bin/env python3
"""dispatch-agent-task CLI:把冻结任务包派发给 Agent 框架子代理并自动留痕。
流程(全部失败关闭):
加载 spec -> 装配任务包 -> 登记 example_run -> 事件账本 run.started ->
框架适配器执行(事件流逐条入账)-> 证据原子落库(raw + 逐回合 llm_call)
-> 结构化输出校验 -> run.completed + 运行终态 -> 打印回执。
业务决策(补证、重写、下一步)不属于本入口:那是框架里模型的事。
用法:
.venv/bin/python dispatch_agent_task.py --spec task.json \
[--provider P] [--model M] [--thinking low] [--run-id ID]
"""
from __future__ import annotations
import argparse
import json
import os
import re
import sys
from dataclasses import replace
from pathlib import Path
from typing import Any, Callable, Iterable, Mapping
SCRIPT_DIR = Path(__file__).resolve().parent
if str(SCRIPT_DIR) not in sys.path:
sys.path.insert(0, str(SCRIPT_DIR))
EVIDENCE_DIR = (SCRIPT_DIR.parent.parent / "record-run-evidence" / "scripts").resolve()
if str(EVIDENCE_DIR) not in sys.path:
sys.path.insert(0, str(EVIDENCE_DIR))
from agent_task import ( # noqa: E402
AgentTaskSpec,
OutputInvalidError,
TaskSpecError,
build_task_package,
load_spec,
validate_structured_output,
)
from agent_trace import AgentTraceWriter, persist_agent_evidence # noqa: E402
from persist_raw import _check_no_secrets # noqa: E402
import read_tools
from pi_runner import ( # noqa: E402
DEFAULT_PI_BIN,
AgentStreamOutcome,
ExecutionPolicy,
FrameworkError,
PiAgentRunner,
)
from run_registry import finish_run, new_run_id, start_run # noqa: E402
EXIT_OK = 0
EXIT_SPEC_INVALID = 2
EXIT_FRAMEWORK_FAILED = 3
EXIT_OUTPUT_INVALID = 4
EXIT_EVIDENCE_FAILED = 5
DEFAULT_RUN_DIR_ROOT = Path("/tmp/muse-agent-runs")
_RUN_ID_PATTERN = re.compile(r"^[A-Za-z0-9_.-]{1,64}$")
def _usage_int(usage: Mapping[str, Any], *keys: str) -> int:
"""读取第一种存在的 usage 字段;脏值与负值按 0 聚合。"""
for key in keys:
if key not in usage:
continue
try:
return max(0, int(usage.get(key) or 0))
except (TypeError, ValueError, OverflowError):
return 0
return 0
def _usage_totals(outcome: AgentStreamOutcome) -> dict[str, int]:
"""聚合全部模型回合的 token 用量;输入口径包含 cache 读写。"""
totals = {"inputTokens": 0, "outputTokens": 0, "cachedTokens": 0, "reasoningTokens": 0}
for call in outcome.model_calls:
usage = call.usage or {}
cached = _usage_int(usage, "cacheRead", "cache_read_input_tokens", "cached_tokens")
cache_write = _usage_int(usage, "cacheWrite", "cache_creation_input_tokens")
totals["inputTokens"] += _usage_int(usage, "input", "input_tokens", "prompt_tokens") + cached + cache_write
totals["outputTokens"] += _usage_int(usage, "output", "output_tokens", "completion_tokens")
totals["cachedTokens"] += cached
totals["reasoningTokens"] += _usage_int(usage, "reasoning", "reasoning_tokens")
return totals
def _cost_totals(outcome: AgentStreamOutcome) -> tuple[float | None, bool]:
"""已知成本求和;任一回合未知则 total 为 None 并标记不完整。"""
known: list[float] = []
complete = True
for call in outcome.model_calls:
if call.cost_usd is None:
complete = False
else:
known.append(call.cost_usd)
if not complete:
return (sum(known) if known else None), False
return sum(known), True
def _model_call_rows(outcome: AgentStreamOutcome) -> list[dict[str, Any]]:
"""把模型回合映射成 agent_trace.persist_agent_evidence 的输入行。"""
return [
{
"actual_model_id": call.actual_model_id,
"usage": dict(call.usage or {}),
"stop_reason": call.stop_reason,
"cost_usd": call.cost_usd,
"duration_ms": None,
}
for call in outcome.model_calls
]
def _write_private_text(path: Path, text: str) -> None:
"""创建仅当前用户可读写的运行审计文件。"""
fd = os.open(path, os.O_WRONLY | os.O_CREAT | os.O_TRUNC, 0o600)
with os.fdopen(fd, "w", encoding="utf-8") as handle:
handle.write(text)
path.chmod(0o600)
def _write_private_json(path: Path, value: Mapping[str, Any]) -> None:
_write_private_text(path, json.dumps(value, ensure_ascii=False, indent=2))
# 工具 server 扩展路径与内建只读文件工具(探索白名单的两类合法成员)。
READ_TOOL_EXTENSION = Path(__file__).resolve().parent / "muse_read_tools_extension.ts"
BUILTIN_TOOL_NAMES = frozenset({"read", "grep", "find", "ls"})
def run_dispatch(
spec_path: str | Path,
*,
repo_root: str | Path,
policy: ExecutionPolicy,
run_id: str | None = None,
run_dir: str | Path | None = None,
connect_factory: Callable[..., Any] | None = None,
launcher: Callable[..., Iterable] | None = None,
trigger_source: str = "user",
trigger_detail: Mapping[str, Any] | None = None,
session_id: str | None = None,
session_dir: str | Path | None = None,
enable_read_tools: bool = False,
) -> tuple[dict[str, Any], int]:
"""执行一次完整派发;所有可控失败都返回稳定回执与退出码。"""
repo_root_path = Path(repo_root).resolve()
try:
spec = load_spec(spec_path)
package = build_task_package(spec, repo_root_path)
_check_no_secrets(package.system_prompt)
_check_no_secrets(package.user_message)
except TaskSpecError as exc:
return (
{
"status": "failed",
"errorCode": "SPEC_INVALID",
"error": str(exc),
"specPath": str(spec_path),
},
EXIT_SPEC_INVALID,
)
except ValueError as exc:
return (
{
"status": "failed",
"errorCode": "SPEC_INVALID",
"error": f"任务包不可安全执行: {type(exc).__name__}",
"specPath": str(spec_path),
},
EXIT_SPEC_INVALID,
)
disabled_read_tools = sorted(
{
tool
for tool in package.spec.tool_allowlist
if tool in read_tools.TOOL_REGISTRY and not enable_read_tools
}
)
if disabled_read_tools:
return (
{
"status": "failed",
"errorCode": "READ_TOOLS_NOT_ENABLED",
"error": f"任务包含探索工具但未启用工具 server: {disabled_read_tools}",
},
EXIT_SPEC_INVALID,
)
policy_kwargs: dict[str, Any] = {"cwd": policy.cwd or str(repo_root_path)}
if session_id is not None:
policy_kwargs["session_id"] = session_id
policy_kwargs["session_dir"] = str(session_dir) if session_dir is not None else None
if enable_read_tools:
# 工具 server 扩展经环境变量定位只读实现;框架进程继承该环境。
os.environ["MUSE_READ_TOOLS_PYTHON"] = sys.executable
os.environ["MUSE_READ_TOOLS_SCRIPT"] = str(
Path(__file__).resolve().parent / "read_tools.py"
)
policy_kwargs["extension_path"] = str(READ_TOOL_EXTENSION)
effective_policy = replace(policy, **policy_kwargs)
if (
package.role_contract.model_policy == "fixed-opus"
and "opus" not in effective_policy.model.lower()
):
return (
{
"status": "failed",
"errorCode": "ROLE_MODEL_POLICY_MISMATCH",
"error": f"角色 {spec.role} 要求 fixed-opus 模型策略",
"requestedModelId": effective_policy.requested_model_id,
},
EXIT_SPEC_INVALID,
)
run_id = run_id or new_run_id(
f"agent-{spec.role}", work_id=spec.work_id, target_chapter=spec.target_chapter
)
if not isinstance(run_id, str) or _RUN_ID_PATTERN.fullmatch(run_id) is None:
return (
{
"status": "failed",
"errorCode": "RUN_ID_INVALID",
"error": "run_id 只能包含 ASCII 字母、数字、点、下划线和短横线,长度不超过 64",
},
EXIT_SPEC_INVALID,
)
run_dir_path = Path(run_dir) if run_dir is not None else DEFAULT_RUN_DIR_ROOT / run_id
identity = package.as_identity()
def _receipt(**fields: Any) -> dict[str, Any]:
base = {
"runId": run_id,
"role": spec.role,
"framework": effective_policy.framework,
"requestedModelId": effective_policy.requested_model_id,
"outputSchemaId": spec.output_schema_id,
"runDir": str(run_dir_path),
**identity,
}
base.update(fields)
return base
try:
if run_dir is None:
DEFAULT_RUN_DIR_ROOT.mkdir(parents=True, mode=0o700, exist_ok=True)
DEFAULT_RUN_DIR_ROOT.chmod(0o700)
run_dir_path.mkdir(parents=True, mode=0o700, exist_ok=False)
run_dir_path.chmod(0o700)
_write_private_text(
run_dir_path / "task-spec.json", Path(spec_path).read_text(encoding="utf-8")
)
_write_private_text(run_dir_path / "system-prompt.txt", package.system_prompt)
_write_private_text(run_dir_path / "user-message.txt", package.user_message)
except (OSError, UnicodeError) as exc:
return (
_receipt(
status="failed",
errorCode="AUDIT_WRITE_FAILED",
error=f"运行审计目录写入失败: {type(exc).__name__}",
),
EXIT_EVIDENCE_FAILED,
)
try:
detail = dict(trigger_detail or {})
except (TypeError, ValueError) as exc:
receipt = _receipt(
status="failed",
errorCode="TRIGGER_DETAIL_INVALID",
error=f"trigger_detail 不是对象: {type(exc).__name__}",
)
_write_private_json(run_dir_path / "receipt.json", receipt)
return receipt, EXIT_SPEC_INVALID
detail.setdefault(
"dispatch",
{
"framework": effective_policy.framework,
"requestedModelId": effective_policy.requested_model_id,
"specSha256": package.spec_sha256,
},
)
try:
detail_json = json.dumps(
detail, ensure_ascii=False, sort_keys=True, default=str, allow_nan=False
)
_check_no_secrets(detail_json)
except (TypeError, ValueError) as exc:
receipt = _receipt(
status="failed",
errorCode="TRIGGER_DETAIL_INVALID",
error=f"trigger_detail 不可安全记录: {type(exc).__name__}",
)
_write_private_json(run_dir_path / "receipt.json", receipt)
return receipt, EXIT_SPEC_INVALID
try:
run_record = start_run(
connect=connect_factory,
run_id=run_id,
work_id=spec.work_id,
target_chapter=spec.target_chapter,
trigger_source=trigger_source,
trigger_detail=detail,
)
except Exception as exc: # noqa: BLE001 - 注册失败不能启动外部 Agent。
receipt = _receipt(
status="failed",
errorCode="RUN_REGISTRY_START_FAILED",
error=f"运行登记失败: {type(exc).__name__}",
)
_write_private_json(run_dir_path / "receipt.json", receipt)
return receipt, EXIT_EVIDENCE_FAILED
if run_record["status"] != "started":
receipt = _receipt(
status="failed",
errorCode="RUN_ID_EXISTS",
error="run_id 已存在,拒绝覆盖既有运行证据",
)
_write_private_json(run_dir_path / "receipt.json", receipt)
return receipt, EXIT_EVIDENCE_FAILED
writer = AgentTraceWriter(
run_id=run_id,
framework=effective_policy.framework,
agent_role=spec.role,
connect=connect_factory,
)
runner = PiAgentRunner(launcher=launcher)
def _failed(
error_code: str,
error: str,
exit_code: int,
*,
agent_failed: bool = False,
session_id: str | None = None,
final_message: str | None = None,
evidence: Mapping[str, Any] | None = None,
cause_error_code: str | None = None,
) -> tuple[dict[str, Any], int]:
"""尽力闭合失败终态;留痕本身失败时升级为证据错误,保留原始原因码。"""
failures: list[str] = []
if final_message is not None:
try:
_write_private_text(run_dir_path / "final-message.txt", final_message)
except OSError as exc:
failures.append(type(exc).__name__)
if agent_failed:
try:
writer.emit(
"agent.failed",
status="error",
requested_model_id=effective_policy.requested_model_id,
details={"errorCode": error_code},
)
except Exception as exc: # noqa: BLE001 - 继续尝试写 run.failed/终态。
failures.append(type(exc).__name__)
try:
writer.emit(
"run.failed",
status="error",
requested_model_id=effective_policy.requested_model_id,
raw_ref=(evidence or {}).get("transcriptId"),
details={"errorCode": error_code},
)
except Exception as exc: # noqa: BLE001 - 继续尝试闭合 example_run。
failures.append(type(exc).__name__)
try:
finish_run(
run_id,
"failed",
trigger_detail={"errorCode": error_code},
connect=connect_factory,
)
except Exception as exc: # noqa: BLE001 - 回执必须揭示终态未能闭合。
failures.append(type(exc).__name__)
fields: dict[str, Any] = {
"status": "failed",
"errorCode": error_code,
"error": error,
"sessionId": session_id,
}
if evidence is not None:
fields["evidence"] = dict(evidence)
if cause_error_code is not None:
fields["causeErrorCode"] = cause_error_code
if failures:
fields.update(
{
"causeErrorCode": cause_error_code or error_code,
"errorCode": "EVIDENCE_PERSIST_FAILED",
"error": "失败终态留痕未完整写入",
"finalizationErrorTypes": sorted(set(failures)),
}
)
exit_code = EXIT_EVIDENCE_FAILED
receipt = _receipt(**fields)
try:
_write_private_json(run_dir_path / "receipt.json", receipt)
except OSError:
receipt["receiptFileWritten"] = False
exit_code = EXIT_EVIDENCE_FAILED
return receipt, exit_code
try:
writer.emit(
"run.started",
status="ok",
requested_model_id=effective_policy.requested_model_id,
details={"specSha256": package.spec_sha256, "toolAllowlist": list(spec.tool_allowlist)},
)
except Exception as exc: # noqa: BLE001 - 未留起始事件时不得启动框架。
return _failed(
"EVIDENCE_PERSIST_FAILED",
f"运行起始事件写入失败: {type(exc).__name__}",
EXIT_EVIDENCE_FAILED,
)
transcript_path = run_dir_path / "transcript.jsonl"
def _remove_local_raw() -> None:
for path in (transcript_path, run_dir_path / "final-message.txt"):
try:
path.unlink(missing_ok=True)
except OSError:
pass
def _persist_failure_evidence(
failed_outcome: AgentStreamOutcome | None,
) -> tuple[dict[str, Any] | None, str | None]:
"""失败也尽量把已有转录和模型回合写入同一套 raw 证据。"""
if failed_outcome is None:
return None, None
try:
transcript_text = transcript_path.read_text(encoding="utf-8")
except (OSError, UnicodeError):
return None, "EVIDENCE_PERSIST_FAILED"
if not transcript_text.strip():
return None, None
try:
_check_no_secrets(transcript_text)
except ValueError:
_remove_local_raw()
return None, "RAW_SECRET_DETECTED"
try:
evidence_result = persist_agent_evidence(
run_id=run_id,
agent_role=spec.role,
system_prompt=package.system_prompt,
user_message=package.user_message,
final_message=failed_outcome.final_text,
transcript=transcript_text,
model_calls=_model_call_rows(failed_outcome),
requested_model_id=effective_policy.requested_model_id,
connect=connect_factory,
)
return evidence_result, None
except Exception:
return None, "EVIDENCE_PERSIST_FAILED"
try:
transcript_fd = os.open(
transcript_path,
os.O_WRONLY | os.O_CREAT | os.O_EXCL,
0o600,
)
with os.fdopen(transcript_fd, "wb") as transcript_file:
outcome = runner.run(
package,
effective_policy,
writer,
timeout_seconds=spec.max_duration_seconds,
raw_sink=transcript_file.write,
)
transcript_path.chmod(0o600)
except FrameworkError as exc:
failed_outcome = exc.outcome
failure_evidence, evidence_error = _persist_failure_evidence(failed_outcome)
if evidence_error is not None:
return _failed(
evidence_error,
"框架失败证据未能安全落库",
EXIT_EVIDENCE_FAILED,
agent_failed=True,
session_id=failed_outcome.session_id if failed_outcome else None,
cause_error_code=exc.error_code,
)
return _failed(
exc.error_code,
str(exc),
EXIT_FRAMEWORK_FAILED,
agent_failed=True,
session_id=failed_outcome.session_id if failed_outcome else None,
final_message=failed_outcome.final_text if failed_outcome else None,
evidence=failure_evidence,
)
except Exception as exc: # noqa: BLE001 - 事件/本地转录失败属于证据失败。
return _failed(
"EVIDENCE_PERSIST_FAILED",
f"框架执行留痕失败: {type(exc).__name__}",
EXIT_EVIDENCE_FAILED,
agent_failed=True,
)
try:
transcript_text = transcript_path.read_text(encoding="utf-8")
_check_no_secrets(transcript_text)
except ValueError as exc:
_remove_local_raw()
return _failed(
"RAW_SECRET_DETECTED",
"框架转录含疑似凭据,已拒绝留存",
EXIT_EVIDENCE_FAILED,
)
except (OSError, UnicodeError) as exc:
return _failed(
"EVIDENCE_PERSIST_FAILED",
f"框架转录回读失败: {type(exc).__name__}",
EXIT_EVIDENCE_FAILED,
session_id=outcome.session_id,
)
if not transcript_text.strip():
return _failed(
"TRANSCRIPT_EMPTY",
"框架转录为空",
EXIT_FRAMEWORK_FAILED,
session_id=outcome.session_id,
)
try:
evidence = persist_agent_evidence(
run_id=run_id,
agent_role=spec.role,
system_prompt=package.system_prompt,
user_message=package.user_message,
final_message=outcome.final_text,
transcript=transcript_text,
model_calls=_model_call_rows(outcome),
requested_model_id=effective_policy.requested_model_id,
connect=connect_factory,
)
except Exception as exc: # noqa: BLE001 - 模型成功但证据失败时必须失败关闭。
return _failed(
"EVIDENCE_PERSIST_FAILED",
f"框架证据落库失败: {type(exc).__name__}",
EXIT_EVIDENCE_FAILED,
session_id=outcome.session_id,
final_message=outcome.final_text,
)
try:
structured = validate_structured_output(outcome.final_text or "", spec)
except OutputInvalidError as exc:
return _failed(
"OUTPUT_SCHEMA_INVALID",
str(exc),
EXIT_OUTPUT_INVALID,
session_id=outcome.session_id,
final_message=outcome.final_text,
evidence=evidence,
)
try:
_write_private_text(run_dir_path / "final-message.txt", outcome.final_text or "")
_write_private_json(run_dir_path / "output.json", structured)
except OSError as exc:
return _failed(
"AUDIT_WRITE_FAILED",
f"运行结果审计文件写入失败: {type(exc).__name__}",
EXIT_EVIDENCE_FAILED,
session_id=outcome.session_id,
evidence=evidence,
)
# 依赖清单:本次运行实际读取了什么(工具调用序列与参数),链路透视的读侧材料。
dependency_count = len(outcome.tool_calls)
if dependency_count:
dependencies = [
{
"seq": index + 1,
"tool": call.name,
"args": dict(call.args) if call.args else {},
"isError": call.is_error,
}
for index, call in enumerate(outcome.tool_calls)
]
try:
_write_private_json(run_dir_path / "dependencies.json", dependencies)
except OSError as exc:
return _failed(
"AUDIT_WRITE_FAILED",
f"依赖清单写入失败: {type(exc).__name__}",
EXIT_EVIDENCE_FAILED,
session_id=outcome.session_id,
evidence=evidence,
)
usage = _usage_totals(outcome)
total_cost, cost_complete = _cost_totals(outcome)
try:
writer.emit(
"run.completed",
status="ok",
requested_model_id=effective_policy.requested_model_id,
actual_model_id=outcome.model_calls[-1].actual_model_id,
usage={
"input": usage["inputTokens"] - usage["cachedTokens"],
"output": usage["outputTokens"],
"cacheRead": usage["cachedTokens"],
},
cost_usd=total_cost,
raw_ref=evidence.get("transcriptId"),
details={
"sessionId": outcome.session_id,
"turns": outcome.turns,
"toolCalls": len(outcome.tool_calls),
"costComplete": cost_complete,
"leaseId": evidence.get("leaseId"),
"llmCallIds": evidence.get("llmCallIds"),
},
)
finish_run(run_id, "completed", connect=connect_factory)
except Exception as exc: # noqa: BLE001 - 成功终态与终态事件必须一起可见。
return _failed(
"RUN_FINALIZE_FAILED",
f"成功终态写入失败: {type(exc).__name__}",
EXIT_EVIDENCE_FAILED,
session_id=outcome.session_id,
evidence=evidence,
)
receipt = _receipt(
status="completed",
sessionId=outcome.session_id,
durationMs=outcome.duration_ms,
turns=outcome.turns,
toolCallCount=len(outcome.tool_calls),
dependencies=(
{"count": dependency_count, "file": "dependencies.json"}
if dependency_count
else None
),
modelCallCount=len(outcome.model_calls),
actualModelIds=[call.actual_model_id for call in outcome.model_calls],
usage=usage,
totalCostUsd=round(total_cost, 6) if total_cost is not None else None,
costComplete=cost_complete,
finalMessageSha256=_sha256_bare(outcome.final_text or ""),
structuredOutputSha256=_sha256_bare(
json.dumps(structured, ensure_ascii=False, sort_keys=True)
),
evidence=evidence,
)
try:
_write_private_json(run_dir_path / "receipt.json", receipt)
except OSError:
receipt["receiptFileWritten"] = False
return receipt, EXIT_OK
def _sha256_bare(text: str) -> str:
import hashlib
return hashlib.sha256(text.encode("utf-8")).hexdigest()
def main(argv: list[str] | None = None) -> int:
parser = argparse.ArgumentParser(description="把冻结任务包派发给 Agent 框架子代理并自动留痕")
parser.add_argument("--spec", required=True, help="AgentTaskSpec JSON 文件路径")
parser.add_argument("--provider", required=True, help="显式框架 provider")
parser.add_argument("--model", required=True, help="显式框架模型")
parser.add_argument("--thinking", default=None, help="思考等级 off/low/medium/high")
parser.add_argument("--pi-bin", default=DEFAULT_PI_BIN, help="框架二进制(默认 pi)")
parser.add_argument("--repo-root", default=".", help="仓库根(解析 .agent/agents 角色文件)")
parser.add_argument("--run-id", default=None, help="指定 run_id(默认自动生成)")
parser.add_argument("--run-dir", default=None, help="运行目录(默认 /tmp/muse-agent-runs/<run_id>)")
parser.add_argument("--trigger-source", default="user", choices=["user", "replay_eval", "diagnostic"])
parser.add_argument("--session-id", default=None, help="框架会话 ID(复用=继续原会话,缺省=一次性会话)")
parser.add_argument("--session-dir", default=None, help="框架会话目录(与 --session-id 同给)")
parser.add_argument("--enable-read-tools", action="store_true", help="启用只读探索工具 server")
args = parser.parse_args(argv)
policy = ExecutionPolicy(
provider=args.provider, model=args.model, thinking=args.thinking, pi_bin=args.pi_bin
)
receipt, code = run_dispatch(
args.spec,
repo_root=args.repo_root,
policy=policy,
run_id=args.run_id,
run_dir=args.run_dir,
trigger_source=args.trigger_source,
session_id=args.session_id,
session_dir=args.session_dir,
enable_read_tools=args.enable_read_tools,
)
print(json.dumps(receipt, ensure_ascii=False, indent=2))
return code
if __name__ == "__main__":
raise SystemExit(main())