"""把冻结任务包装配后派发给 Agent 执行器,并闭合 Muse 侧证据。 流程(全部失败关闭): 加载 spec -> 装配任务包 -> 登记 example_run -> 事件账本 run.started -> runtime 执行(事件流逐条入账)-> 证据原子落库(raw + 逐回合 llm_call) -> 结构化输出校验 -> run.completed + 运行终态。 业务决策(补证、重写、下一步)不属于本入口。 """ from __future__ import annotations import json import os import sys from dataclasses import replace from pathlib import Path from typing import Any, Callable, Iterable, Mapping PROJECT_ROOT = next( parent for parent in (Path(__file__).resolve().parent, *Path(__file__).resolve().parents) if (parent / "AGENTS.md").is_file() and (parent / ".git").exists() ) if str(PROJECT_ROOT) not in sys.path: sys.path.insert(0, str(PROJECT_ROOT)) DISPATCH_SCRIPTS = ( PROJECT_ROOT / "muse" / "lifecycle" / "dispatch" / "skills" / "dispatch-agent-task" / "scripts" ) EVIDENCE_DIR = PROJECT_ROOT / "muse" / "authority" / "evidence" / "skills" / "record-run-evidence" / "scripts" READ_TOOLS_DIR = PROJECT_ROOT / "muse" / "authority" / "tools" / "read" for _path in (DISPATCH_SCRIPTS, EVIDENCE_DIR, READ_TOOLS_DIR): if str(_path) not in sys.path: sys.path.insert(0, str(_path)) from role_task import ( # noqa: E402 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 framework.adapters.pi.normalization import normalize_pi_transcript # noqa: E402 from framework.adapters.pi.runner import ( # noqa: E402 AgentStreamOutcome, ExecutionPolicy, FrameworkError, ) from role_policy import validate_role_execution_policy # noqa: E402 from run_registry import finish_run, new_run_id, start_run # noqa: E402 from runtime.agent_executor import execute_agent # noqa: E402 from runtime.runs import ( # noqa: E402 DEFAULT_RUN_DIR_ROOT, open_transcript, prepare_run_dir, read_transcript_text, remove_local_raw, seed_run_inputs, validate_run_id, write_private_json, write_private_text, ) EXIT_OK = 0 EXIT_SPEC_INVALID = 2 EXIT_FRAMEWORK_FAILED = 3 EXIT_OUTPUT_INVALID = 4 EXIT_EVIDENCE_FAILED = 5 READ_TOOL_EXTENSION = PROJECT_ROOT / "framework" / "adapters" / "pi" / "mcp_bridge.ts" BUILTIN_TOOL_NAMES = frozenset({"read", "grep", "find", "ls"}) 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 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( READ_TOOLS_DIR / "read_tools.py" ) policy_kwargs["extension_path"] = str(READ_TOOL_EXTENSION) effective_policy = replace(policy, **policy_kwargs) try: validate_role_execution_policy(package, effective_policy) except ValueError as exc: return ( { "status": "failed", "errorCode": "ROLE_MODEL_POLICY_MISMATCH", "error": str(exc), "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 validate_run_id(run_id): 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: run_dir_path = prepare_run_dir(run_id, run_dir) seed_run_inputs( run_dir_path, spec_text=Path(spec_path).read_text(encoding="utf-8"), system_prompt=package.system_prompt, user_message=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, ) 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 _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 = read_transcript_text(run_dir_path) 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(run_dir_path) 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: with open_transcript(run_dir_path) as transcript_file: outcome = execute_agent( package.as_framework_request( session_mode="continue" if session_id is not None else "fresh", request_id=run_id, ), effective_policy, writer, timeout_seconds=spec.max_duration_seconds, raw_sink=transcript_file.write, launcher=launcher, ) 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 = read_transcript_text(run_dir_path) _check_no_secrets(transcript_text) except ValueError as exc: remove_local_raw(run_dir_path) 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: framework_artifact = normalize_pi_transcript( transcript_path, run_dir_path / "framework-events.jsonl", run_id=run_id, framework_version="pi-json-v1", ) except (OSError, UnicodeError, ValueError) as exc: return _failed( "FRAMEWORK_ARTIFACT_INVALID", f"通用框架事件工件生成失败: {type(exc).__name__}", EXIT_EVIDENCE_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), "unknownFrameworkEvents": list(outcome.unknown_event_types), "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), unknownFrameworkEvents=list(outcome.unknown_event_types), 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, frameworkArtifact=framework_artifact, ) 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()