rpa_core 运行时链路图
全文用一个厨房类比贯穿:workflow 是菜谱,orchestrator 是厨师长,executor 是各工种厨师。 每个环节列出「做了什么」的具体动作,末尾一行保留设计原因。
⓪ 名词速查(先看这页再往下读)
一份 JSON 写的自动化剧本:先做什么、遇到条件怎么办、循环几次。≈ 菜谱。包含唯一 id、默认输入 inputs 和一棵节点树 root。
Abstract Syntax Tree,抽象语法树:workflow JSON 加载进内存后的树状结构——「sequence 依次做」是父节点,if / forEach 是分支 / 循环节点,action(调一条命令)是叶子。运行时不逐字读 JSON,而是递归遍历这棵树来执行。≈ 把菜谱拆成有嵌套关系的步骤树。
每条命令的「操作规范卡」:叫什么(id 如 browser.click)、哪个执行器接活、输入输出长什么样(JSON Schema)、对外部世界有什么影响(effect)、能不能重试。≈ 每道工序的工艺卡。
启动时把 commands/ 下全部 manifest 加载成的只读总表:命令 id → 规范。附带 sha256 digest(指纹)——内容变一个字节指纹就变,用于检测「恢复时契约是否还是当初那份」。≈ 整本装订好的规范册 + 封条编号。
开跑前的静态检查员:遍历 workflow 树,查引用是否存在、命令是否在 catalog 里、能力是否被授权、危险命令有没有配重试。检查通过才发「执行单」。≈ 开灶前对照规范册逐项核对食材和工具。
检查通过后冻结的执行依据:workflow 树 + catalog 指纹 + 所需能力。之后 orchestrator 只认这份单据,不再回头翻 JSON。≈ 盖章封存的单子,中途不许改。
整个运行的总指挥:拿着执行单从树根开始走,把每个 action 派给对应执行器;管超时、重试、取消、进度检查点和终态判定。它自己不直接操作浏览器/窗口/文件。≈ 厨师长:只派活、看表、记账,不掌勺。
真正动手干活的组件:browser.playwright 开浏览器点页面、desktop.uia / desktop.win32 操作 Windows 窗口控件、python.worker 在子进程里跑用户 Python。输入一条 CommandInvocation(这次调用谁、参数是什么),回一张 CommandResult(结果回执),仅此而已。≈ 各工种厨师。
一张排班表:执行器名字 → 执行器实例(如 "browser.playwright" → PlaywrightExecutor)。orchestrator 从 manifest 里读到「这活归 browser.playwright」,就按名字来这里要人。CLI 启动时组装一次。
厨师长派活的工单(Invocation:命令、参数、第几次尝试)和厨师交回的回执(Result:成功与否、产出值、副作用证据 effects、诊断信息)。规则:厨师只填回执,绝不直接改厨师长的账本。
运行中的备菜台:inputs 放用户原料、steps 放每步产出、loop 放当前循环轮次。workflow 里写 ${steps.step1.outputs.title},resolver 负责整牌替换——按这个取料牌去备菜台取值,找不到立刻报错。
每条命令声明三件事:effect.kind 对外界的影响类型(只读 / 幂等写 / 危险写 / 会话);replay 重做一次是否安全;idempotency 幂等键从哪来。失败回执里若带 unknown 副作用证据 = 「做没做成不知道」,运行时绝不自动重试。≈ 工艺卡上的风险等级。
流水线日志(每动一下追加一行)和最终质检报告(run 的终态、产出、错误)。两个都落盘 run 才算成功。≈ 操作日志 + 质检单。
每个 action 干完后拍的进度快照:哪些步骤确认完成、备菜台上是什么值。进程崩了或想接着跑,就从最近一份快照恢复,已完成的活不再重做。≈ 阶段完工拍照存档。
叫停铃和总时限。取消用 asyncio Event 传播:先通知执行器自行收尾清理,等它退出再终判。deadline 到点整棵执行树判超时。每次重试前都会查「还来不来得及」。
RunHandle 是启动后拿到的遥控器(cancel / wait);RunResult 是终态成绩单(状态、起止时间、outputs、返回值、错误、digest)。
① 名词关系图(谁持有谁、谁调用谁)
② 一次 run 从启动到结束,每个环节做了什么
- validate:只编译不执行,打印
{"valid": true, "catalogDigest": ...} - run:创建 4 个执行器 → 组装
ExecutorRegistry(排班表)→orchestrator.run(plan)→ 打印结果 JSON → 状态非 succeeded 退出码 1 - resume:同上但调
orchestrator.resume(plan, --run-id);checkpoint/恢复拒绝时打印错误并退出码 2;--allow-indeterminate是人工确认开关
--artifacts 工件根目录 输出:控制台 JSON + 退出码- 读
commands/下所有*.json,逐个用 pydantic 校验成CommandManifest(id 格式、版本号、effect policy 三维自洽) - 检查:id 与文件名一致、id 不重复、input/output schema 本身是合法 JSON Schema
- 把全部 manifest 规范化排序 → JSON 序列化 → 算 sha256 作为
catalog.digest - 返回只读映射
{命令id → manifest},运行期间内容不变
CommandCatalog(含 digest)- 把 workflow JSON 校验成
Workflow模型(节点按type判别,形成 AST) - 递归遍历 AST 做静态检查:节点 id 不重复;action 引用的命令必须存在于 catalog;
${...}引用——inputs 键必须存在、steps.*不能向前引用、loop 变量必须在作用域内、try 错误变量必须先声明 replay=unsafe的命令禁止配retry_count > 0- 汇总所有命令需要的 capabilities,与授权集合比对,缺一个就拒绝
- 全部通过 → 产出冻结的
ExecutionPlan(盖章执行单)
ExecutionPlan(workflow + digest + capabilities,不可变)- start:生成 uuid 形式的
run_id→ 建 asyncio task + 取消Event→ 返回RunHandle(遥控器:可cancel()/wait()) - resume:先过 5 道恢复门(见 ③),全过之后同样建 task,但 scopes / 已完成步骤从 checkpoint 恢复
- 准备工作目录、写初始 scopes(备菜台):
inputs = workflow 默认输入 + 用户传参、steps = {}、loop = {} - 写第一条事件:
runStarted(或runResumed+ 已完成步骤清单)
RunHandle,最终 await 到 RunResult- 每次进入节点先做两件事:取消信号置位 → 立即中止;节点路径键 ∈
completedSteps→ 直接跳过(恢复时用) - sequence:按顺序执行 children
- if:resolver 求条件值 → 走 then 或 else 分支
- forEach:解析 items(必须是 list/tuple)→ 每轮设置
loop.item / loop.index→ 执行 children → 结束后恢复原 loop - try:执行 children,出错(取消除外)→ 把错误对象存进错误变量 → 执行 catch → 结束恢复变量
- return:解析返回值 → 以异常方式终止整棵树
#i 迭代段,forEach 才能按轮恢复,不会误跳过未跑的轮次。ADR 0004- 从 catalog 取该命令的 manifest;unsafe 命令配了重试 → 直接拒绝
- 解析
with参数(resolver 按取料牌去备菜台取值)→ jsonschema 校验输入,失败INVALID_INPUT - 写
stepStarted事件 - 重试循环(最多
retry_count+1次):每次检查 workflow deadline 余量 → 计算 attempt 超时 → 把CommandInvocation(工单)交给执行器 → 收CommandResult(回执)→ 成功跳出;失败且可重试 → 指数退避后重试;带 unknown effect 的失败 → 绝不重试,直接进indeterminate - 成功后校验回执:非 pure 命令必须有 kind 匹配、status=committed 的 effect 证据(否则
INVALID_OUTPUT);输出过 jsonschema - 写
scopes["steps"][节点id] = {value, outputs, effects, diagnostics}(后续节点用${steps.节点id.outputs...}引用) - checkpoint 原子落盘(临时文件+replace),成功后才把路径键记入已完成 → 写
stepCompleted事件
CommandResult、状态单向流经 scopes(可测试、可审计);checkpoint 先于事件落盘,崩溃窗口最多只丢一条事件、绝不重放副作用。规则 3 / ADR 0004launch起浏览器并生成 UUID sessionId- 其余操作必须带 sessionId,在对应页面上导航/点击/输入/等待/读文本/查询元素
attachWindow按标题/类名/进程绑窗口 → sessionId- UIA 语义定位控件,执行点击/输入/读文本
- 同上,但走 Win32 消息/句柄路线
- 额外支持
hotkey、menuSelect(菜单/对话框)
- 起子进程跑
workers/python_worker.py,stdin/stdout 传 JSON - 限定 workspace 目录内读写
CommandResult(回执);绝不触碰编排器内部状态events.jsonl:每次状态变化追加一行(seq 递增):runStarted → stepStarted → stepAttemptStarted → stepCompleted/stepFailed/stepRetried → runFinishedresult.json:终态时写一次完整RunResult(状态、起止时间、outputs、返回值、error、digest)- 写事件失败 →
RunPersistenceError立即中止 run(证据比执行优先);写 result.json 失败 → 绝不返回 success,改判PERSISTENCE_FAILED
run_artifacts/<run_id>/ 下的 events.jsonl + result.json + checkpoint.json- 算 workflow deadline = 当前时间 +
workflow.timeout_seconds,用asyncio.timeout_at包住整棵执行树 - 按异常类型映射终态:
WorkflowReturn→succeeded;unknown effect→indeterminate;checkpoint 写失败→recovery_required;超时→failed+TIMEOUT;取消→cancelled;命令错误→failed;意外异常→failed+INTERNAL - 写
runFinished事件 +result.json - 最后补写一次最终 checkpoint(失败忽略——终态证据已落盘,恢复时用上一份边界 checkpoint)
RunResult(正常 run 和恢复 run 走同一条收口路径,证据格式完全一致)③ checkpoint 里有什么,恢复时每道门检查什么
version / workflowId / catalogDigest——格式版本与归属标识completedSteps——已完成 action 的路径键列表(如loop/#0/step)scopes——inputs / steps / loop的值快照(恢复后所有${...}引用从这里取值)returnValue——workflow 已 return 时的返回值,平时是 null
stepCompleted 事件之前,原子写(临时文件 + replace)CheckpointErrorindeterminate 且未传 allow_indeterminate=True → 拒绝,要求人工确认timeout_seconds 计算(monotonic 时间跨进程无意义);重走 AST,已完成路径键直接跳过,其余照常执行