@fyeeme/pi-dynamic-workflows

extension

Deterministic TypeScript workflow orchestration for pi. Fuses the pi-dynamic-workflows design (declarative graph, 10 step primitives, heuristic planner, outcome collectors) with Claude Code's workflow engine coordination mechanisms (deterministic sandbox,

by — · v2.0.1 · published 2w ago

$ pi install npm:@fyeeme/pi-dynamic-workflows
downloads/mo
311
stars
—
last push
—
open issues
—

Signals

license: MITtestspi manifest: missinginstall size: —deps: 0peer deps: 0

Download trend

No downloads in the last 12 weeks.

README

@fyeeme/pi-dynamic-workflows

为 pi 打造的确定性 TypeScript 工作流编排。

把工作流定义成一份类型化的声明式步骤列表,运行后即可获得可恢复、受预算约束、可中止的执行。融合 pi-dynamic-workflows 设计(10 个步骤原语 + 启发式 planner + outcome 收集器)与 Claude Code 工作流引擎的协调机制(确定性沙箱、缓存键恢复、按 agent 中止、动态预算、失控上限)。

语言:English | 中文


为什么需要它

一次运行 = 一份步骤列表(agent / code / log / fan_out / loop_until / loop_until_dry / adversarial / tournament / classify_route / sub_workflow)。引擎保证:

  • 确定性 —— workflow .ts 文件经 AST 守卫,禁止 Date.now() / Math.random() / new Date();run id 是 (timestamp, sequence) 的纯函数。
  • 恢复即不重派 —— 每个 agent 调用以 sha256(workflow + prompt + signature) 为键写入 journal;重跑同一 workflow 会回放缓存的 agent(零子进程派发)。
  • 按 agent 中止 —— 每个在途 agent 持有自己的 AbortController;skipAgent/retryAgent 只针对一个调用,不打扰兄弟调用。
  • 预算 + 失控上限 —— maxAgents / maxTokens 由实时池强制;MAX_BATCH=4096、MAX_LIFETIME_AGENTS=1000 超限抛 BudgetExceededError(绝不静默截断)。
  • 无需 pi 即可测试 —— agent 派发可注入;测试传一个 fake dispatch,无需二进制、无需 provider API、无需 token。

开箱即用:本包的扩展工厂会组合 @fyeeme/pi-subagents (实时代理 UI——编辑器上方 widget、FleetView、/agents——以及 subagent 工具),且版本钉定在自身依赖副本上,工作流运行的可观测性开箱即得。 单独安装 pi-subagents 是可选的,二者可共存(组合幂等)。

安装

这是一个 pi 扩展包(workspace / 本地),尚未发布到 npm。在 pi workspace 中:

npm install --ignore-scripts   # 水合(本包是 workspace 依赖)

这会解析 npm registry 上的 @fyeeme/pi-subagents(无需保持同级仓库目录结构)。

随后从包根模块导入公共 API(TypeScript barrel,包直接以 .ts 源码分发):

import { defineWorkflow, runWorkflow } from "@fyeeme/pi-dynamic-workflows/src/index.ts";

包的 pi.extensions 入口注册 run_workflow 工具,并接入 pi-subagents 的共享子代理 UI(编辑器上方实时 agent widget、下方 FleetView、/agents 转录查看器——每个工作流 agent 以其 step id 出现在其中;需先安装 @fyeeme/pi-subagents,见上方提示)。引擎本身也可经上述导入直接使用。


快速上手

import { defineWorkflow, runWorkflow } from "@fyeeme/pi-dynamic-workflows/src/index.ts";

const wf = defineWorkflow({
	name: "draft-and-refine",
	steps: [
		{ id: "draft", type: "agent", prompt: "起草一段发布说明。" },
		{ id: "refine", type: "agent", prompt: (ctx) => `把下面改写得更精炼:\n\n${ctx.step("draft").results}` },
	],
});

const result = await runWorkflow({ workflow: wf, cwd: process.cwd(), now: Date.now() });
console.log(result.status, result.steps[1].results);

runWorkflow 默认每次 agent 调用派生一个 pi --mode json -p --no-session 子进程(默认 dispatch),因此需要 pi 在 PATH 上并配置好 provider。测试或离线运行时注入一个 fake dispatch 即可(见教程)。


使用教程

1. 定义工作流

defineWorkflow 是一个类型化恒等助手——让你对 steps 判别联合获得完整类型检查。

const wf = defineWorkflow({
	name: "research",
	budget: { maxAgents: 10, maxTokens: 50_000 },
	steps: [
		{ id: "gather", type: "agent", prompt: "列出关于主题 X 的 3 个来源。" },
		{ id: "summarize", type: "agent", prompt: (ctx) => `总结:\n${ctx.step("gather").results}` },
	],
});

ctx.input 是本次运行的初始输入;ctx.step(id) 返回某个已执行步骤的 { results, stats }(若该 id 尚未执行则抛错)。

2. 运行

const result = await runWorkflow({
	workflow: wf,
	cwd: process.cwd(),
	now: 1700000000000,   // 确定性起始时间(ms),同时是 journal/run-id 的种子
	input: "主题 X",
});
// result.status: "completed" | "failed" | "aborted"
// result.steps:  StepResult[](按顺序,每个已执行步骤一条)
// result.stats:  汇总 { tokens, cost, durationMs, agents, failures }
// result.journalFile: 该 workflow 的 JSONL journal 路径

now 必填且确定性——传入本次运行的起始时间;引擎绝不读时钟来生成身份。相同的 (workflow, prompts) 永远生成相同的缓存键。

3. fan_out —— 并行 agent + 合并

const wf = defineWorkflow({
	name: "parallel-research",
	steps: [
		{
			id: "fan",
			type: "fan_out",
			over: () => ["alpha", "beta", "gamma"],
			agent: (topic) => ({ prompt: `研究 ${topic}。` }),
			parallelism: 3,
			merge: (results) => results.join("\n---\n"),
		},
	],
});

fan_out 会先预检整批是否在预算内(MAX_BATCH=4096);每个 item 独立缓存键、独立可中止。

4. loop_until —— 迭代到条件 / 预算

const wf = defineWorkflow({
	name: "refine-loop",
	steps: [
		{
			id: "loop",
			type: "loop_until",
			prompt: (ctx, i) => `第 ${i + 1} 稿。当前:\n${ctx.step("loop")?.results ?? ctx.input}`,
			until: (ctx, i) => i >= 3,
			maxIterations: 5,
		},
	],
});

每次迭代都是独立的缓存键 agent 调用;maxIterations 与预算共同约束循环。

5. 组合模式 —— adversarial / tournament / classify_route

它们构建在同一个 dispatchAgentCall 之上,因此天然享有缓存恢复、预算与中止。

// 生成候选,再由 N 个评判者按 rubric 打分并汇总。
defineWorkflow({
	name: "review",
	steps: [
		{
			id: "adv",
			type: "adversarial",
			produce: { prompt: "写这个函数。" },
			rubric: ["正确性", "处理空输入", "无 off-by-one"],
			judges: 3,             // 默认;minPass 默认为过半数
		},
	],
});
// results: { candidate, passed, passCount, minPass, judges: [{pass, reason}] }

// N 个不同候选,M 个评判者排名,选出多数赢家。
defineWorkflow({
	name: "pick",
	steps: [{ id: "tmt", type: "tournament", candidates: 3, judges: 2, produce: { prompt: "解决 X。" } }],
});
// results: { candidates, winner, judges: [{winner, reason}] }

// 分类输入,再运行匹配路由的子步骤。
defineWorkflow({
	name: "route",
	steps: [
		{
			id: "cr",
			type: "classify_route",
			classifier: { prompt: (ctx) => `分类意图:${ctx.input}` },
			routes: {
				bug: [{ id: "file", type: "agent", prompt: "提一个 bug 报告。" }],
				faq: [{ id: "answer", type: "agent", prompt: "回答这个 FAQ。" }],
			},
			fallback: [{ id: "escalate", type: "agent", prompt: "转给人工。" }],
		},
	],
});
// results: { category, matched, route: StepResult[], routeStatus }

评判/分类的 JSON 采用宽松解析(LLM 常把 "true"/"0" 当字符串返回);路由嵌套有深度上限以防循环。

6. 恢复 —— 重跑零派发

journal 位于 <cwd>/.pi/workflows/<workflow.name>/journal.jsonl(按 workflow 而非按 run,因此同一 workflow 重跑可跨 run 命中缓存;run id 不计入键):

const first  = await runWorkflow({ workflow: wf, cwd, now: T0 }); // 每个 agent 都派发
const second = await runWorkflow({ workflow: wf, cwd, now: T1 }); // 零派发——全部缓存命中

prompt 变了 → 该 agent 的缓存键变 → 重新派发;没变的则回放。

7. 按 agent 中止与跳过

runner 持有一个 AgentSpawnRegistry。拿到它即可针对单个在途调用:

import { createSpawnRegistry, skipAgent } from "@fyeeme/pi-dynamic-workflows/sessions/spawn.ts";

const registry = createSpawnRegistry();
const runP = runWorkflow({ workflow: fanOutWf, cwd, now, registry });
// ...等 `fan#2` 进入在途状态后:
skipAgent(registry, "fan#2");   // 只中止这一个;兄弟调用继续
const result = await runP;       // status "completed"——这批除第 2 项外都跑完了

abortAgent(registry, callId) 中止一个调用;retryAgent(registry, callId) 中止以便 runner 重新派发。调用 id 形如 ${step.id}#${n}(从 1 起)。

8. 预算强制

const wf = defineWorkflow({
	name: "capped",
	budget: { maxAgents: 2 },
	steps: [{ id: "fan", type: "fan_out", over: () => [1, 2, 3], agent: (i) => ({ prompt: `${i}` }) }],
});
const result = await runWorkflow({ workflow: wf, cwd, now });
// result.status === "failed",result.error 匹配 /budget|exhausted/i

maxAgents 在派发整批/单个 agent 前检查;maxTokens 在 agent 落定后由 BudgetPool.isExhausted 强制。二者都抛 BudgetExceededError——绝不静默截断。

9. Outcome 收集器

从 agent 的文本输出里抽取结构化值:

import { collect } from "@fyeeme/pi-dynamic-workflows/src/index.ts";

const urls  = collect<string[]>({ kind: "url" }, result.steps[0].results as string);
const json  = collect({ kind: "json" }, agentText);     // 第一个平衡的 JSON 值
const paths = collect<string[]>({ kind: "file_path" }, agentText);

url / file_path / json 都是文本的纯函数——可对任意 StepResult.results 使用。

10. 启发式 planner

按关键词把目标草拟成单步工作流脚手架(compare → tournament、review → adversarial、classify → classify_route,其余 → agent)。这是一个待你打磨的起点,不是真正的 NL 规划器:

import { heuristicallyPlan } from "@fyeeme/pi-dynamic-workflows/src/index.ts";

const wf = heuristicallyPlan("比较三种排序方案", { judges: 3 });
// wf.steps[0].type === "tournament"

11. 加载 .ts 工作流文件

import { loadWorkflowModule } from "@fyeeme/pi-dynamic-workflows/src/index.ts";

const mod = await loadWorkflowModule<{ workflow: ReturnType<typeof defineWorkflow> }>({
	filePath: "./my-workflow.ts",
});
const wf = mod.workflow;

loader 在 jiti 加载之前跑确定性 AST 守卫——workflow 体内若调用 Date.now() / Math.random() / new Date() 会在加载时被拒(这些会让缓存键失稳)。注意:守卫只扫描入口文件;请让 workflow 单文件,或单独守卫被引入的 helper。

12. 无需 pi 即可测试

注入一个 fake dispatch——无二进制、无 provider、无 token。本包自带的 108 个测试就是这样跑的:

import { runWorkflow, type AgentDispatch } from "@fyeeme/pi-dynamic-workflows/src/index.ts";

const fake: AgentDispatch = async (_registry, opts) => ({
	callId: opts.callId,
	exitCode: 0,
	messages: [{ role: "assistant", content: [{ type: "text", text: `out:${opts.task}` }], /* ...其余字段 */ } as never],
	stderr: "",
	usage: { input: 10, output: 5, cacheRead: 0, cacheWrite: 0, cost: 0, contextTokens: 15, turns: 1 },
	model: "fake",
	stopReason: "stop",
	aborted: false,
});

const result = await runWorkflow({ workflow: wf, cwd: tempDir, now: 1000, dispatch: fake });

步骤类型速查

typepayload 要点结果
agentprompt: string | (ctx)=>string、model?、tools?、systemPrompt?最后一条 assistant 文本
codetransform: (ctx) => unknown(纯函数、不派发、不缓存)transform 的返回值
logmessage: string | (ctx)=>string(叙事行、零派发零 token)该消息(触发 onLog)
fan_outover()、agent(item,i)、parallelism?、merge?合并后的数组(或 merge 的输出)
loop_untilprompt(ctx,i)、until(ctx,i)、maxIterations?每轮输出的数组
adversarialproduce、rubric[]、judges?、minPass?{ candidate, passed, passCount, judges }
tournamentcandidates、judges、produce{ candidates, winner, judges }
classify_routeclassifier、routes: Record<cat, Step[]>、fallback?{ category, matched, route, routeStatus }
sub_workflowworkflow: WorkflowDefinition、input?、inheritBudget?{ steps, status, workflowName, error }
loop_until_dryagent(item, i)、keyOf?、merge?、maxRounds?、dryThreshold?发现项组成的数组

每个步骤都接受 id、retry?: { maxRetries } 与 onBudgetExhaust?: "throw" | "null"——"null" 下预算耗尽时该步骤返回 null (降级)而不是中止运行;运行结果用 degradedSteps 记录降级步骤。默认 "throw" 保持 fail-fast 保证。


API 参考

runWorkflow(opts) → Promise<RunResult>

选项
workflowWorkflowDefinition必填
cwdstring必填(journal 基目录)
nownumber必填——确定性起始 ms
input?unknownctx.input
budget?Budget覆盖 workflow.budget
signal?AbortSignal整运行中止信号
listeners?AgentLifecycleListenersonAgentStart/End/Skip/Retry/CacheHit + onLog(log 步骤)+ onUpdate(流式 delta)
dispatch?AgentDispatch默认 = 真实 spawnAgent
maxPromptBytes?numberA6 尺寸守卫——解析后 prompt 超过该字节数在派发前抛 size-limit(默认 256 KB)
policyGate?(wf) => { allow, reason? }A6 门控——返回 { allow: false } 在任何派发前中止运行
registry?AgentSpawnRegistry供外部 skipAgent/abortAgent
journalDir?string默认 <cwd>/.pi/workflows/<name>
sequence?numberrun-id 消歧

RunResult = { runId, status, steps: StepResult[], stats: StepStats, journalFile?, error?, errorCategory?, degradedSteps? }。

同时导出

defineWorkflow、loadWorkflowModule、collect(含 urlCollector/filePathCollector/jsonCollector/parseFirstJson)、heuristicallyPlan、createSpawnRegistry/abortAgent/skipAgent/retryAgent(来自 sessions/spawn.ts),以及全部步骤/结果/上下文类型。


设计 —— Claude Code 融合

机制模块作用
确定性沙箱src/determinism/ast-guard.tsAST 级禁止 workflow 源码中的非确定性 API
确定性 run idsrc/state/names.tsgenerateRunId({timestamp, sequence}) 为纯函数
缓存键恢复src/cache/{key,journal}.tssha256(workflow+prompt+signature) + 每运行 JSONL journal
按 agent 中止sessions/spawn.tsMap<callId, ChildProcess> + 每调用 AbortController;中止 → 对单个进程 SIGTERM
预算 + 上限src/budget/{pool,caps}.ts实时 BudgetPool + MAX_BATCH/MAX_LIFETIME_AGENTS

三处范式冲突(CC 命令式 ↔ 声明式图)已化解:预算循环变成 fan_out 读取的预检值;进程内 AbortController 变成子进程表;vm 沙箱变成对 jiti 加载源码的加载期 AST 守卫。


测试

node_modules/.bin/tsc -p packages/extensions/pi-dynamic-workflows/tsconfig.json --noEmit   # 类型检查
node_modules/.bin/vitest --run packages/extensions/pi-dynamic-workflows                    # 108 个测试

真实 pi 子进程冒烟(默认 dispatch)位于 examples/smoke-real-pi.ts——在 pi 与 provider 配置妥当后手动运行。

License:MIT。