2026-09-18 11:29:07 +08:00

476 lines
20 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.

# -*- coding: utf-8 -*-
"""pbl_agent_runtime 挂载入口M4a 唯一实现)。
历史说明QC 退回 #7 处置):
旧「九关裁决 + 4 表」实现tables.py / tool_registry.py / verdict.py / api.py
已**删除**本文件不再引用M4am4a_*.py8 步裁决 + 5 表)是本模块唯一实现,
`load_pbl_agent_runtime()` 与 `load_pbl_agent_runtime_m4a()` 均委托 M4a 装配。
load_pbl_agent_runtime(env=None, tenant_id=None)
1) ensure_tables —— 建 5 表幂等pbl_agent_trace / pbl_agent_trace_stage 为 append-only
2) seed_agentsdesigner / criticcritic write_allowed=0+ seed_tools13 启用 / 9 禁用)
3) **挂载即 self_check()**数量契约、Critic 零写、8 步裁决链、四类审批、默认裁决 DENY
任一不过 → 抛 RuntimeError宿主应用启动失败fail-closed绝不带病上线
4) 向 ServerEnv 注册 wwwroot/api/*.dspy 所需的 8 个处理函数 + 只读契约
5) 返回 default_verdict= 'DENY'
后端sqlor/ahserver未挂载时自动降级 MemoryStore建表/seed/裁决/审批全链路可用,
数据仅存于进程内重启即失不抛错——便于自测与离线演示US-05 兜底同源策略)。
"""
from .m4a_kernel import (
ADJUDICATION_STEPS,
TRACE_ELEMENTS,
PblError,
TenantContext,
ACTOR_SYSTEM,
get_store,
)
from .m4a_tables import TABLES, APPEND_ONLY_TABLES, ensure_tables
from .m4a_registry import (
ENABLED_TOOL_CODES,
DISABLED_TOOL_CODES,
APPROVAL_REQUIRED_TOOLS,
PERM_PLATFORM_ADMIN,
list_agents,
list_tools,
get_tool,
register_tool,
set_tool_status,
seed_agents,
seed_tools,
)
from .m4a_adjudicate import adjudicate, adjudicate_report, explain_chain
from .m4a_approval import MANDATORY_APPROVAL_TYPES
from .m4a_trace import get_trace, list_traces
from .m4a import M4aApi, load_m4a, make_admin_ctx
from .m4a_init import MODULE_NAME, get_dbname
import functools
DEFAULT_VERDICT = "DENY" # fail-closed任何异常/未知一律拒绝
EXPECTED_ENABLED = 13
EXPECTED_DISABLED = 9
EXPECTED_STEPS = 8
EXPECTED_APPROVAL_TYPES = 4
EXPECTED_TABLES = 5
POLICY_VERSION = "m4a-1.0.0"
AGENTS = ("designer", "critic")
# ---------------------------------------------------------------------------
# 自检(挂载即执行,不过即抛)
# ---------------------------------------------------------------------------
def self_check(tenant_id=None, store=None, run_probes=True):
"""M4a 契约自检。返回 (all_ok, msgs)。
校验项:
C1 工具数量契约 22 = 13 启用 + 9 禁用
C2 自动发布类工具必须禁用blueprint.publish_auto / marketplace.create_listing
C3 Critic 零写权限write_allowed=0 且无任何写工具 allowed_agents 含 critic
C4 裁决链恰为 8 步且顺序为 S1..S8
C5 四类强制人工审批齐备
C6 5 张表齐备且 append-only 表带标记
C7 默认裁决 = DENY
C8 探针:禁用工具/未注册工具/Critic 写 三类调用必须被拒run_probes=True 时)
"""
msgs = []
all_ok = True
st = store or get_store()
tid = tenant_id or "__platform__"
admin = make_admin_ctx(tid)
tools = list_tools(ctx=admin, store=st)
enabled = [t for t in tools if t.get("status") == "enabled"]
disabled = [t for t in tools if t.get("status") == "disabled"]
# C1
if len(tools) != EXPECTED_ENABLED + EXPECTED_DISABLED:
all_ok = False
msgs.append("FAIL C1 工具总数=%d 应为 %d"
% (len(tools), EXPECTED_ENABLED + EXPECTED_DISABLED))
if len(enabled) != EXPECTED_ENABLED:
all_ok = False
msgs.append("FAIL C1 启用工具=%d 应为 %d" % (len(enabled), EXPECTED_ENABLED))
if len(disabled) != EXPECTED_DISABLED:
all_ok = False
msgs.append("FAIL C1 禁用工具=%d 应为 %d" % (len(disabled), EXPECTED_DISABLED))
msgs.append("C1 工具契约:%d = %d 启用 + %d 禁用"
% (len(tools), len(enabled), len(disabled)))
# C2 自动发布/自动改课必须禁用
must_disabled = ("blueprint.publish_auto", "curriculum.modify_auto",
"marketplace.create_listing", "kdb.write", "billing.charge")
codes_disabled = set(DISABLED_TOOL_CODES)
for code in must_disabled:
if code not in codes_disabled:
all_ok = False
msgs.append("FAIL C2 %s 必须处于禁用清单" % code)
msgs.append("C2 高危工具禁用:%s" % ", ".join(must_disabled))
# C3 Critic 零写
agents = list_agents(ctx=admin, store=st)
by_code = {a.get("agent_code"): a for a in agents}
if set(by_code) != set(AGENTS):
all_ok = False
msgs.append("FAIL C3 Agent 定义=%s 应为 %s" % (sorted(by_code), list(AGENTS)))
critic = by_code.get("critic") or {}
if critic and int(critic.get("write_allowed") or 0) != 0:
all_ok = False
msgs.append("FAIL C3 critic.write_allowed=%r 必须为 014.1 零写权限)"
% critic.get("write_allowed"))
write_tools_for_critic = [
t.get("tool_code") for t in enabled
if int(t.get("write_operation") or 0) == 1
and "critic" in (t.get("allowed_agents") or [])
]
if write_tools_for_critic:
all_ok = False
msgs.append("FAIL C3 critic 持有写工具:%s" % write_tools_for_critic)
msgs.append("C3 Critic 零写权限write_allowed=0写工具授权数=0")
# C4 8 步裁决链
steps = [s[0] for s in ADJUDICATION_STEPS]
if len(steps) != EXPECTED_STEPS or steps != ["S%d" % i for i in range(1, 9)]:
all_ok = False
msgs.append("FAIL C4 裁决链=%s 应为 S1..S88 步顺序固定)" % steps)
else:
msgs.append("C4 fail-closed 裁决链 8 步:%s"
% "".join("%s %s" % (s[0], s[1]) for s in ADJUDICATION_STEPS))
# C5 四类审批
if len(MANDATORY_APPROVAL_TYPES) != EXPECTED_APPROVAL_TYPES:
all_ok = False
msgs.append("FAIL C5 强制审批类型=%d 应为 %d"
% (len(MANDATORY_APPROVAL_TYPES), EXPECTED_APPROVAL_TYPES))
appr_tools = [t.get("tool_code") for t in enabled
if int(t.get("require_approval") or 0) == 1]
if set(appr_tools) != set(APPROVAL_REQUIRED_TOOLS):
all_ok = False
msgs.append("FAIL C5 require_approval 工具=%s 应为 %s"
% (appr_tools, list(APPROVAL_REQUIRED_TOOLS)))
msgs.append("C5 四类强制人工审批:%s;需审批工具=%s"
% (", ".join(MANDATORY_APPROVAL_TYPES), ", ".join(appr_tools)))
# C6 5 表 + append-only
tbl_names = [t.get("tblname") for t in TABLES]
if len(tbl_names) != EXPECTED_TABLES:
all_ok = False
msgs.append("FAIL C6 表数=%d 应为 %d%s"
% (len(tbl_names), EXPECTED_TABLES, tbl_names))
for name in ("pbl_agent_trace", "pbl_agent_trace_stage"):
if name not in APPEND_ONLY_TABLES:
all_ok = False
msgs.append("FAIL C6 %s 必须标记 append-only" % name)
msgs.append("C6 数据表 %d 张:%sappend-only%s"
% (len(tbl_names), ", ".join(tbl_names), ", ".join(APPEND_ONLY_TABLES)))
# C7 默认裁决
if DEFAULT_VERDICT != "DENY":
all_ok = False
msgs.append("FAIL C7 DEFAULT_VERDICT=%r 必须为 'DENY'" % DEFAULT_VERDICT)
else:
msgs.append("C7 默认裁决 = DENYfail-closed")
# C8 探针
if run_probes:
probes = [
("designer", "blueprint.publish_auto", "禁用工具(自动发布)"),
("designer", "pbl.publish", "未注册工具(白名单外 default-deny"),
("critic", "blueprint.update", "Critic 写蓝图14.1"),
("critic", "publish.request", "Critic 发布14.1"),
("designer", "compile.trigger", "未审批编译执行S7 审批门)"),
]
for agent_code, tool_code, why in probes:
step = code = None
try:
rep = adjudicate_report(admin, agent_code, tool_code, {}, store=st)
# Verdict.to_dict()allowed=False 即拒绝fail-closed 语义)
denied = rep.get("allowed") is False
step = rep.get("step_reached")
code = rep.get("reason_code")
except PblError as e:
denied, code = True, e.code
except Exception: # noqa: BLE001 —— 任何异常都视为拒绝fail-closed
denied, code = True, "EXCEPTION_AS_DENY"
if not denied:
all_ok = False
msgs.append("FAIL C8 探针未被拒:%s%s%s" % (agent_code, tool_code, why))
else:
msgs.append("C8 探针拒绝:%s%s @%s %s%s"
% (agent_code, tool_code, step, code, why))
msgs.append("轨迹 7 要素:%s" % ", ".join(TRACE_ELEMENTS))
if all_ok:
msgs.append("SELF_CHECK %s: PASS%d 工具 / %d 表 / %d 步 / %d 类审批)"
% (MODULE_NAME, len(tools), len(tbl_names),
len(ADJUDICATION_STEPS), len(MANDATORY_APPROVAL_TYPES)))
return all_ok, msgs
def _wrap(fn):
"""dspy 处理函数统一包装PblError → 结构化拒绝响应fail-closed不抛 500
任何未预期异常同样按拒绝返回default_verdict=DENY 同源策略),
并带 error_code=PBL_E_FORBIDDEN避免异常细节泄漏到前端。
"""
@functools.wraps(fn)
async def _inner(**params_kw):
try:
data = await fn(**params_kw)
except PblError as e:
return dict(e.to_dict(), ok=False, default_verdict=DEFAULT_VERDICT)
except Exception as e: # noqa: BLE001
return {"ok": False, "error_code": "PBL_E_FORBIDDEN",
"error_message": "内部异常,按 fail-closed 拒绝:%s" % type(e).__name__,
"default_verdict": DEFAULT_VERDICT}
if isinstance(data, dict):
data.setdefault("ok", True)
return data
return _inner
# ---------------------------------------------------------------------------
# dspy 处理函数wwwroot/api/*.dspy 直接调用,名称必须与 dspy 内一致)
# ---------------------------------------------------------------------------
def _tid(params_kw):
tid = (params_kw or {}).get("tenant_id")
if not tid:
# fail-closed租户上下文缺失即拒绝S1
raise PblError("PBL_E_TENANT_MISSING", "tenant_id 强制打头,缺失即拒绝")
return tid
def _api(params_kw):
return load_m4a(tenant_id=_tid(params_kw))
@_wrap
async def pbl_agent_designer_run(**params_kw):
"""US-01/03/04Designer 生成 / 修改 / 澄清action 分派,缺省 generate"""
api = _api(params_kw)
action = params_kw.get("action", "generate")
if action == "modify":
return api.designer_modify(params_kw.get("blueprint_id"),
params_kw.get("instruction"),
session_no=params_kw.get("session_no"),
version_no=params_kw.get("version_no"))
if action == "clarify":
return api.designer_clarify(params_kw.get("blueprint_id"),
params_kw.get("missing_fields") or [],
session_no=params_kw.get("session_no"),
round_no=params_kw.get("round_no"))
return api.designer_generate(params_kw.get("intent_text") or "",
owner_teacher_id=params_kw.get("owner_teacher_id"),
class_id=params_kw.get("class_id"),
session_no=params_kw.get("session_no"))
@_wrap
async def pbl_agent_critic_run(**params_kw):
"""US-18Critic 只读评审(零写权限,写操作一律被 S6 拒)。"""
api = _api(params_kw)
if params_kw.get("trace_no"):
return api.get_critic_report(params_kw["trace_no"])
blueprint_id = params_kw.get("blueprint_id")
if not blueprint_id:
# 后端pbl_blueprint 权威服务)未挂载 / Designer 兜底提案尚未落库时,
# Critic 无可评审对象:返回结构化降级说明(不抛 500不伪造评审结论
return {"ok": False, "degraded": True,
"error_code": "PBL_E_BACKEND_UNAVAILABLE",
"error_message": "blueprint_id 缺失或蓝图后端未挂载Critic 无可评审对象;"
"Designer 兜底提案落 pending_backend待后端挂载后重试",
"agent_code": "critic", "write_allowed": 0}
return api.critic_review(blueprint_id,
version_no=params_kw.get("version_no"),
session_no=params_kw.get("session_no"))
@_wrap
async def pbl_tool_adjudicate(**params_kw):
"""fail-closed 8 步裁决预检只裁不执行。tool_code 为空 → 返回裁决链说明。"""
api = _api(params_kw)
tool_code = params_kw.get("tool_code")
if not tool_code:
return {"steps": explain_chain(), "default_verdict": DEFAULT_VERDICT,
"policy_version": POLICY_VERSION}
return api.adjudicate(params_kw.get("agent_code") or "designer", tool_code,
args=params_kw.get("args") or {})
@_wrap
async def pbl_agent_trace_list(**params_kw):
"""F-AG-04轨迹查询可查不可篡改"""
api = _api(params_kw)
if params_kw.get("trace_no"):
return {"trace": api.get_trace(params_kw["trace_no"]),
"completeness": api.trace_completeness(params_kw["trace_no"])}
return api.list_traces(params_kw.get("filters") or {},
page=int(params_kw.get("page") or 1),
size=int(params_kw.get("size") or 20))
def _trace_write(api, action, agent_code, params_kw):
from . import m4a_trace as _t
ctx = api.agent_ctx(agent_code)
if action == "start":
return _t.start_trace(agent_code, session_no=params_kw.get("session_no"),
input_context=params_kw.get("input_context"),
ctx=ctx, store=api.store)
if action == "append":
return _t.append_trace(params_kw.get("trace_no"),
params_kw.get("stage") or params_kw.get("element"),
params_kw.get("payload") or {},
ctx=ctx, store=api.store)
raise PblError("PBL_E_VALIDATION",
"action 非法:%rappend-only 仅支持 start / append" % action)
@_wrap
async def pbl_approval_create(**params_kw):
"""T10Agent 发起人工审批提案Agent 不可自批)。"""
api = _api(params_kw)
return api.request_approval(
action_type=params_kw.get("action_type"),
object_id=params_kw.get("object_id"),
agent_code=params_kw.get("agent_code") or "designer",
tool_code=params_kw.get("tool_code"),
object_type=params_kw.get("object_type"),
action_payload=params_kw.get("action_payload") or {},
approver_id=params_kw.get("approver_id"),
trace_no=params_kw.get("trace_no"),
)
@_wrap
async def pbl_approval_decide(**params_kw):
"""14.2:审批裁决 —— 仅 actor_type=user人类可调Agent 调用被拒。"""
api = _api(params_kw)
return api.decide_approval(params_kw.get("approval_no"),
params_kw.get("status"),
params_kw.get("user_id"),
comment=params_kw.get("comment"))
@_wrap
async def pbl_approval_list(**params_kw):
"""待办审批工作台approver_id 必填)。"""
api = _api(params_kw)
approver_id = params_kw.get("approver_id")
if not approver_id:
raise PblError("PBL_E_VALIDATION", "approver_id 必填")
return api.list_pending_approvals(approver_id)
# 修正 pbl_agent_trace_write 的 start 分支(避免占位返回)
@_wrap
async def pbl_agent_trace_write(**params_kw):
"""轨迹写入append-only。action=start 开轨迹action=append 追加 7 要素之一。"""
api = _api(params_kw)
return _trace_write(api, params_kw.get("action", "append"),
params_kw.get("agent_code") or "designer", params_kw)
# ---------------------------------------------------------------------------
# 对外契约
# ---------------------------------------------------------------------------
def api():
"""模块对外 API 契约(供宿主应用 / 其它模块经 ServerEnv 调用)。"""
return {
"adjudicate": adjudicate,
"adjudicate_report": adjudicate_report,
"adjudication_steps": explain_chain,
"list_tools": list_tools,
"get_tool": get_tool,
"register_tool": register_tool,
"set_tool_status": set_tool_status,
"list_agents": list_agents,
"get_trace": get_trace,
"list_traces": list_traces,
"self_check": self_check,
"load_m4a": load_m4a,
"M4aApi": M4aApi,
"PblError": PblError,
"DEFAULT_VERDICT": DEFAULT_VERDICT,
"POLICY_VERSION": POLICY_VERSION,
"ADJUDICATION_STEPS": ADJUDICATION_STEPS,
"MANDATORY_APPROVAL_TYPES": MANDATORY_APPROVAL_TYPES,
"ENABLED_TOOL_CODES": ENABLED_TOOL_CODES,
"DISABLED_TOOL_CODES": DISABLED_TOOL_CODES,
}
def load_pbl_agent_runtime(env=None, tenant_id=None, force_memory=False,
do_seed=True, run_self_check=True):
"""宿主应用标准挂载入口(委托 M4a 装配)。
返回 default_verdict'DENY')。自检不过抛 RuntimeErrorfail-closed
"""
srv = env
if srv is None:
try:
from ahserver.serverenv import ServerEnv
srv = ServerEnv()
except Exception: # noqa: BLE001 —— 无宿主环境时降级 MemoryStore
srv = None
st = get_store(force_memory=force_memory or srv is None)
ensure_tables(st)
tid = tenant_id or "__platform__"
seed = None
if do_seed:
ctx = TenantContext(tenant_id=tid, actor_type=ACTOR_SYSTEM, actor_id="seed",
permissions={PERM_PLATFORM_ADMIN})
seed = {"agents": seed_agents(ctx=ctx, store=st),
"tools": seed_tools(ctx=ctx, store=st)}
if run_self_check:
ok, msgs = self_check(tenant_id=tid, store=st)
for m in msgs:
_log(srv, m)
if not ok:
raise RuntimeError("%s self_check FAILEDfail-closed拒绝启动%s"
% (MODULE_NAME,
"; ".join([m for m in msgs if m.startswith("FAIL")][:8])))
if srv is not None:
setattr(srv, "pbl_agent_tools", list_tools(ctx=make_admin_ctx(tid), store=st))
setattr(srv, "pbl_agent_enabled_tools", list(ENABLED_TOOL_CODES))
setattr(srv, "pbl_agent_disabled_tools", list(DISABLED_TOOL_CODES))
setattr(srv, "pbl_agent_api", api())
setattr(srv, "pbl_agent_dbname", get_dbname(srv))
setattr(srv, "pbl_agent_seed", seed)
# dspy 处理函数注册wwwroot/api/*.dspy 直接按名调用)
setattr(srv, "pbl_agent_designer_run", pbl_agent_designer_run)
setattr(srv, "pbl_agent_critic_run", pbl_agent_critic_run)
setattr(srv, "pbl_tool_adjudicate", pbl_tool_adjudicate)
setattr(srv, "pbl_agent_trace_list", pbl_agent_trace_list)
setattr(srv, "pbl_agent_trace_write", pbl_agent_trace_write)
setattr(srv, "pbl_approval_create", pbl_approval_create)
setattr(srv, "pbl_approval_decide", pbl_approval_decide)
setattr(srv, "pbl_approval_list", pbl_approval_list)
modules = getattr(srv, "modules", None)
if isinstance(modules, list) and MODULE_NAME not in modules:
modules.append(MODULE_NAME)
return DEFAULT_VERDICT
def _log(srv, msg):
logger = getattr(srv, "logger", None) if srv is not None else None
if logger is not None and hasattr(logger, "info"):
try:
logger.info("[%s] %s" % (MODULE_NAME, msg))
return
except Exception: # noqa: BLE001
pass
print("[%s] %s" % (MODULE_NAME, msg))
if __name__ == "__main__":
_ok, _msgs = self_check(run_probes=True)
for _m in _msgs:
print(_m)
print("RESULT: %s" % ("PASS" if _ok else "FAIL"))
raise SystemExit(0 if _ok else 1)