Workflow 基础
前面几章我们讲的智能体,本质都是一个模型 + 一圈工具。
模型决定下一步做什么,工具负责执行,走向哪里全靠模型即兴发挥。这种动态决策在简单任务上很好用,但任务一旦变长、变重,问题就来了。
- 你想让流程确定地走 A→B→C,而不是让模型每次重新决定去哪,做不到,模型可能第 3 轮就跑偏。
- 你想把一个智能体的输出精确地喂给下一个智能体,中间还要插一道校验、一道人工审批,模型自己协调不了。
- 你想并行跑三路分析再汇总,单智能体的工具调用循环是串行的,并行得自己造。
这类需求有一个共同特征,流程的形状是你预先定义好的,模型只在某些节点上负责动脑,这就是 Workflow 要解决的事。
官方文档里有句话很到位,我直接转述:Agent 和 Workflow 都可以包含多个步骤来达成目标,但它们处于不同的抽象层级。
| Agent(智能体) | Workflow(工作流) | |
|---|---|---|
| 谁来决定走向 | 模型。根据上下文和工具,LLM 动态决定下一步 | 你。流程的边、条件都是预先定义好的 |
| 步骤是否固定 | 不固定,模型每轮重新决策 | 固定,执行路径是图结构 |
| 模型的角色 | 全程主导 | 只在某些节点上负责推理,其余节点是普通代码 |
| 擅长 | 开放式、探索性任务(编码、研究) | 业务流程明确、要强可控的任务(审批、ETL、多 agent 编排) |
| 可控性 | 弱(模型可能跑偏) | 强(路径写死,可校验、可检查点) |
Agent 是模型驱动流程,Workflow 是流程驱动模型。
Workflow 把模型降级成图里的一个个节点(executor),而流程的骨架由你用代码(或 YAML)钉死。这换来的是确定性、可观测性、可恢复性,这三件事在长链路、生产级任务里比灵活重要得多。
值得强调的是,Workflow 不是 Agent 的替代品,而是包含关系,Workflow 里的节点可以是普通代码,也可以是完整的 Agent(甚至是一个 Harness)。
Workflow 的核心能力在 Microsoft.Agents.AI.Workflows 包里,在进行后续实验时,请先安装这个 nuget 包。
核心三件套 Executor、Edge、WorkflowBuilder
理解 Workflow 只需要记住 Executor、Edge、WorkflowBuilder 三个东西,后面所有内容都是围绕它们展开的。
- Executor 决定每个节点干什么。
- Edge 决定消息往哪流、流不流。
- WorkflowBuilder 把上面两者钉成一张图,交给运行时去跑。
把这三件套记住,后面所有的高级特性(条件路由、并行、检查点、HITL)都只是在这三者上加料。
Executor 节点
Executor(执行器)是工作流里的处理单元,它接收一条消息,干点活,然后产出一到多条消息或产出工作流输出。你可以把它理解成流程图里的一个方框。
源码里它的定义很朴素:
// 源码位置:src/Microsoft.Agents.AI.Workflows/Executor.cs
public abstract class Executor : IIdentified
{
public string Id { get; }
protected Executor(string id, ExecutorOptions? options = null, ...);
...
}
每个 Executor 有一个唯一 Id,这是它在流程里的门牌号,流程的 Edge (流程连线)就是靠 Id 找到目标的。
Executor 可以是:
- 一段普通代码:例如将内容转大写、调 API、写日志、跑校验。
- 一个 AI Agent:让模型在这个节点上推理。
Edge 连线
Edge(边)定义了 Executor 之间的连接,决定了消息怎么流。 一条边从源 executor指向目标 executor,源产出的消息会沿着边送到目标。你可以把边理解成方框之间的箭头。
builder.AddEdge(executorA, executorB); // 连一条边
边可以是无条件的(源有产出就一定送过去),也可以带条件(只有满足某个谓词时才送)。
除了最基础的直连边,SDK 还内置了 FanOut(一对多分发)和 FanIn(多对一汇聚)两种特殊边,这两种留到后面再讲,本章只看最简单的直连。
WorkflowBuilder 构建流程
WorkflowBuilder 把 Executor 和 Edge 拼成一张有向图,然后 Build() 出一个不可变的 Workflow 实例。
WorkflowBuilder builder = new(startExecutor); // 指定起点
builder.AddEdge(executorA, executorB); // 连一条边
Workflow workflow = builder.Build(); // 编译成可执行的图
Build() 出来的 Workflow 是不可变的,节点、边在构建期就钉死了,运行期不会变。
这点很重要,它意味着流程可以被序列化、做检查点、被多个 session 复用。
定义 Executor 的三种方式
Executor 相当于流程的节点,定义流程节点主要有三种方式,官方推荐用 [MessageHandler] 特性 + partial 类,因为它走编译时源码生成,性能更好、能在编译期做校验、还兼容 Native AOT。
下面把这三种方式都过一遍。
方式一 [MessageHandler]
类要标 partial、继承 Executor,方法上打 [MessageHandler] 特性。源码生成器会在编译期自动注册这些 handler:
using Microsoft.Agents.AI.Workflows;
internal sealed partial class UppercaseExecutor() : Executor("UppercaseExecutor")
{
[MessageHandler]
private ValueTask<string> HandleAsync(string message, IWorkflowContext context)
{
// 返回值同样会自动沿边发送
return ValueTask.FromResult(message.ToUpperInvariant());
}
}
[MessageHandler] 支持的方法签名很灵活(见 MessageHandlerAttribute.cs 的 XML 注释),涵盖同步/异步、有返回值/无返回值:
void Handler(TMessage, IWorkflowContext)
void Handler(TMessage, IWorkflowContext, CancellationToken)
ValueTask Handler(TMessage, IWorkflowContext)
ValueTask Handler(TMessage, IWorkflowContext, CancellationToken)
TResult Handler(TMessage, IWorkflowContext)
TResult Handler(TMessage, IWorkflowContext, CancellationToken)
ValueTask<TResult> Handler(TMessage, IWorkflowContext)
ValueTask<TResult> Handler(TMessage, IWorkflowContext, CancellationToken)
它最大的好处是一个 Executor 能处理多种消息类型,只要写多个 [MessageHandler] 方法,每个方法的第一个参数类型不同就行:
internal sealed partial class SampleExecutor() : Executor("SampleExecutor")
{
[MessageHandler]
private ValueTask<string> HandleStringAsync(string message, IWorkflowContext context)
=> ValueTask.FromResult(message.ToUpperInvariant());
[MessageHandler]
private ValueTask<int> HandleIntAsync(int message, IWorkflowContext context)
=> ValueTask.FromResult(message * 2);
}
这种方式下,SampleExecutor 既能吃 string 也能吃 int,路由是按消息类型自动分发的。
方式二 继承 Executor<TInput, TOutput>
定义 Executor 还可以继承泛型定义,固定输入参数和输出参数。
public abstract class Executor<TInput, TOutput>(string id, ...)
: Executor(id, ...), IMessageHandler<TInput, TOutput>
{
public abstract ValueTask<TOutput> HandleAsync(
TInput message, IWorkflowContext context, CancellationToken cancellationToken = default);
}
方式三 函数式(FunctionExecutor)
有时候你只想把一个 lambda 包成节点,不想专门写个类。SDK 提供了 FunctionExecutor<TInput, TOutput>,它的构造函数直接吃一个委托:
// 直接用 lambda 构造一个 executor
var uppercase = new FunctionExecutor<string, string>(
"UppercaseExecutor",
(message, context, ct) => ValueTask.FromResult(message.ToUpperInvariant()));
定义一个完整的 Workflow
理论讲多了容易晕,先上一个能跑的最小例子。
这个流程很简单,输入一段文本 → 第一个节点转大写 → 第二个节点反转字符串 → 输出,这个案例来源于官方。
继承 Executor<TInput, TOutput> 并重写 HandleAsync:
using Microsoft.Agents.AI.Workflows;
// 节点 1:把输入文本转大写
internal sealed class UppercaseExecutor() : Executor<string, string>("UppercaseExecutor")
{
public override ValueTask<string> HandleAsync(
string message, IWorkflowContext context, CancellationToken cancellationToken = default)
{
// 返回值会被自动当作消息,沿边送给下一个 executor
return ValueTask.FromResult(message.ToUpperInvariant());
}
}
// 节点 2:反转字符串
internal sealed class ReverseTextExecutor() : Executor<string, string>("ReverseTextExecutor")
{
public override ValueTask<string> HandleAsync(
string message, IWorkflowContext context, CancellationToken cancellationToken = default)
=> ValueTask.FromResult(string.Concat(message.Reverse()));
}
定义一个 Workflow 并且把两个节点连接起来,构造一个完整的流程。
using Microsoft.Agents.AI.Workflows;
// 1. 实例化节点
var uppercase = new UppercaseExecutor();
var reverse = new ReverseTextExecutor();
// 2. 拼图:uppercase 为起点,连一条到 reverse 的边
WorkflowBuilder builder = new(uppercase);
builder.AddEdge(uppercase, reverse);
// 3. 声明 reverse 是输出节点(它的产出会作为工作流的最终结果)
builder.WithOutputFrom(reverse);
// 4. 编译成不可变的 Workflow
Workflow workflow = builder.Build();
// 5. 流式执行,边跑边看每个节点的产出
await using StreamingRun run = await InProcessExecution.RunStreamingAsync(workflow, input: "Hello, World!");
await foreach (WorkflowEvent evt in run.WatchStreamAsync())
{
if (evt is ExecutorCompletedEvent completed)
{
Console.WriteLine($"{completed.ExecutorId}: {completed.Data}");
}
else if (evt is WorkflowOutputEvent output)
{
Console.WriteLine($"最终结果: {output.Data}");
}
}
输出:

这个例子虽然小,但已经包含了 Workflow 的全部骨架:两个 Executor、一条 Edge、一个起点、一个输出声明。
这里有个细节,使用 AddEdge(uppercase, reverse) 连线后,如果要抓取返回值或者记录输出,需要使用 .WithOutputFrom(reverse)。
SendMessageAsync 触发节点
正常情况下,一个 Executor 的返回值传递给下游节点,下游节点会触发一次。
但是如果我们希望多次触发下游节点,那么可以不返回数据,而是使用 context.SendMessageAsync() 多次发送消息,次数下游 Executor 也会被多次触发。
internal sealed partial class SplitExecutor() : Executor("SplitExecutor")
{
[MessageHandler]
private async ValueTask HandleAsync(string message, IWorkflowContext context)
{
// 一条输入,主动往图内管道放两条消息(下游会收到两份)
await context.SendMessageAsync(message.ToUpperInvariant());
await context.SendMessageAsync(message.ToLowerInvariant());
}
}
当你要发不止一条消息、或者要在处理中途分批发的时候。返回值只能给一个,使用SendMessageAsync 想发几条发几条。发给某个节点的所有消息,会被逐条送进 handler,收到 N 条 = handler 跑 N 次,框架不会把它们合并成一次调用。
不过这引出一个容易出现,如果下游节点如果有状态,会在多次调用间累积。比如一个带计数器的节点,上游一次发 3 条,计数器会在同一个超级步内涨 3 次。
其实,返回值的本质也是 SendMessageAsync()。
例如,你手动返回一个值:
public override ValueTask<string> HandleAsync(string message, ...) =>
ValueTask.FromResult(message.ToUpperInvariant()); // 等价于 SendMessage(大写后的 message)
实际上 SDK 源码会把返回值手动使用 SendMessageAsync() 触发下游节点。
也就是说,本质底层都是 SendMessageAsync()。
// 源码位置:src/Microsoft.Agents.AI.Workflows/Executor.cs , ExecuteCoreAsync 内部
// result 是你的 handler 执行后的返回值
if (result.Result is not null && this.Options.AutoSendMessageHandlerResultObject)
{
// 框架替你调了一次 SendMessageAsync , 这就是"返回值自动发送"的全部秘密
await context.SendMessageAsync(result.Result, cancellationToken: cancellationToken);
}
监听流程的执行
前面有一段示例代码,我们运行一个流程后,会监听 ExecutorCompletedEvent、WorkflowOutputEvent 两个事件,ExecutorCompletedEvent 是节点结束事件,可以接收节点运行信息和返回值,WorkflowOutputEvent 则是监听 context.YieldOutputAsync(X) 。
await foreach (WorkflowEvent evt in run.WatchStreamAsync())
{
if (evt is ExecutorCompletedEvent completed)
{
Console.WriteLine($"{completed.ExecutorId}: {completed.Data}");
}
else if (evt is WorkflowOutputEvent output)
{
Console.WriteLine($"最终结果: {output.Data}");
}
}
当我们运行工作流时,我们无法知道流程内部的细节,不知道运行到哪里了,除非执行结束或者异常,否则我们不能连接流程运行细节。
context.YieldOutputAsync(X) 的作用就是让我们在流程之外能够收到 Executor 的运行细节。
在 Executor 使用 context.YieldOutputAsync() 输出状态。
// 节点 1:把输入文本转大写
internal sealed class UppercaseExecutor() : Executor<string, string>("UppercaseExecutor")
{
public override ValueTask<string> HandleAsync(
string message, IWorkflowContext context, CancellationToken cancellationToken = default)
{
context.YieldOutputAsync("开始转大写");
return ValueTask.FromResult(message.ToUpperInvariant());
}
}
// 节点 2:反转字符串
internal sealed class ReverseTextExecutor() : Executor<string, string>("ReverseTextExecutor")
{
public override ValueTask<string> HandleAsync(
string message, IWorkflowContext context, CancellationToken cancellationToken = default)
{
context.YieldOutputAsync("开始反转字符串");
return ValueTask.FromResult(string.Concat(message.Reverse()));
}
}
定义流程并接收信息:
var uppercase = new UppercaseExecutor();
var reverse = new ReverseTextExecutor();
WorkflowBuilder builder = new(uppercase);
builder.AddEdge(uppercase, reverse);
// 必须 使用 WithOutputFrom,否则 YieldOutputAsync 不起效
builder.WithOutputFrom(uppercase);
builder.WithOutputFrom(reverse);
Workflow workflow = builder.Build();
await using StreamingRun run = await InProcessExecution.RunStreamingAsync(workflow, input: "Hello, World!");
await foreach (WorkflowEvent evt in run.WatchStreamAsync())
{
if (evt is ExecutorCompletedEvent completed)
{
Console.WriteLine($"当前节点 {completed.ExecutorId} 输出结果 {completed.Data}");
}
else if (evt is WorkflowOutputEvent output)
{
Console.WriteLine($"外部接收信息: {output.Data}");
}
}

实际上,Executor 的返回值会被包装为 ExecutorCompletedEvent,所以我们可以在流程外使用 ExecutorCompletedEvent 读取每个节点的执行状态。
不过使用 context.YieldOutputAsync() 一定要搭配 WithOutputFrom(),否则流程不会触发这个节点的 WorkflowOutputEvent 事件。
在 Workflow 中使用 Agent
Workflow 支持把 Agent 当作 Executor。
本节我们使用简化版的股票情绪识别做一个炒股 Workflow。
由于 Agent 的输出是不确定的,但是 Workflow 的 Executor 要求节点之间的数据流通提前确认结构,因此我们需要为 Agent 节点和下游节点之间定义 Json 结构,避免 Agent 回复纯文本,下游难以使用。
using System.Text.Json.Serialization;
public enum Sentiment { Bullish, Bearish, Neutral } // 利好 / 利空 / 中性
public sealed class SentimentResult
{
[JsonPropertyName("sentiment")]
public Sentiment Sentiment { get; set; }
[JsonPropertyName("score")] // -1.0(极空) ~ 1.0(极多)
public double Score { get; set; }
[JsonPropertyName("reason")]
public string Reason { get; set; } = "";
}
定义一个 Agent。
AIAgent agent = new OpenAIClient(
credential: new ApiKeyCredential("1234"),
options: new OpenAIClientOptions { Endpoint = new Uri("http://127.0.0.1:1234/v1") })
.GetChatClient("qwen/qwen3.5-9b")
.AsIChatClient()
.AsAIAgent(new ChatClientAgentOptions
{
Name = "SentimentAgent",
ChatOptions = new()
{
Instructions = """
你是 A 股舆情分析师。判断给定文本对相关公司的情绪。
- bullish: 看好; bearish: 看空; neutral: 中性
score 范围 -1.0 到 1.0。reason 一句话说明依据。
""",
ResponseFormat = Microsoft.Extensions.AI.ChatResponseFormat.ForJsonSchema<SentimentResult>()
}
});
定义上下游节点:
// 节点 1:适配器节点 —— 把 string 转成 ChatMessage 并发送 TurnToken 触发 agent 执行。
// agent 节点遵循 Chat Protocol:收到消息只是缓存,收到 TurnToken 才真正运行。
[SendsMessage(typeof(ChatMessage))]
[SendsMessage(typeof(TurnToken))]
internal sealed partial class InputNode() : Executor<string>("Input")
{
[MessageHandler]
public override async ValueTask HandleAsync(string text, IWorkflowContext context, CancellationToken cancellationToken = default)
{
// 转成 agent 能理解的对话消息
await context.SendMessageAsync(new ChatMessage(ChatRole.User, text), cancellationToken);
// 发送回合令牌,通知 agent 处理已累积的消息(emitEvents: false 表示非流式,避免增量输出刷屏)
await context.SendMessageAsync(new TurnToken(emitEvents: false), cancellationToken);
}
}
// 节点 2:把 agent 的输出(List<ChatMessage>,内容为 JSON)解析成 SentimentResult 并输出
internal sealed partial class OutputNode() : Executor<List<ChatMessage>, string>("Output")
{
[MessageHandler]
public override ValueTask<string> HandleAsync(List<ChatMessage> messages, IWorkflowContext context, CancellationToken cancellationToken = default)
{
string json = string.Concat(messages.Select(m => m.Text ?? ""));
// 兼容模型把 JSON 包在 ```json ... ``` 里的情况
json = json.Trim();
if (json.StartsWith("```"))
{
int start = json.IndexOf('\n');
int end = json.LastIndexOf("```", StringComparison.Ordinal);
if (start >= 0 && end > start)
{
json = json[(start + 1)..end].Trim();
}
}
SentimentResult? result;
try
{
result = JsonSerializer.Deserialize<SentimentResult>(json,
new JsonSerializerOptions
{
PropertyNameCaseInsensitive = true,
Converters = { new JsonStringEnumConverter() }
});
}
catch (JsonException exception)
{
throw new InvalidOperationException($"模型未返回有效 JSON。原始响应: {json}", exception);
}
return ValueTask.FromResult(
$"[{result!.Sentiment}] 分数 {result.Score:F2}:{result.Reason}");
}
}
把 Agent 当作节点注册,并且构建完整 Workflow。
AIAgent sentimentAgent = new OpenAIClient(
credential: new ApiKeyCredential("1234"),
options: new OpenAIClientOptions { Endpoint = new Uri("http://127.0.0.1:1234/v1") })
.GetChatClient("qwen/qwen3.5-9b")
.AsIChatClient()
.AsAIAgent(new ChatClientAgentOptions
{
Name = "SentimentAgent",
ChatOptions = new()
{
Instructions = """
你是 A 股舆情分析师。判断给定文本对相关公司的情绪。
- bullish: 看好; bearish: 看空; neutral: 中性
score 范围 -1.0 到 1.0。reason 一句话说明依据。
""",
ResponseFormat = Microsoft.Extensions.AI.ChatResponseFormat.ForJsonSchema<SentimentResult>()
}
});
// ── 拼图 ──
var input = new InputNode();
var output = new OutputNode();
ExecutorBinding sentimentNode = sentimentAgent.BindAsExecutor(
new AIAgentHostOptions { ForwardIncomingMessages = false });
Workflow workflow = new WorkflowBuilder(input)
.AddEdge(input, sentimentNode)
.AddEdge(sentimentNode, output)
.WithOutputFrom(output)
.Build();
// ── 跑 ──
await using StreamingRun run = await InProcessExecution.RunStreamingAsync(
workflow, input: "宁德时代新电池能量密度创新高,产能爬坡顺利!");
await foreach (WorkflowEvent evt in run.WatchStreamAsync())
{
if (evt is ExecutorCompletedEvent completed)
{
Console.WriteLine($"当前节点 {completed.ExecutorId} 输出结果 {completed.Data}");
}
else if (evt is WorkflowOutputEvent o)
{
Console.WriteLine($"外部接收信息: {o.Data}");
}
else if (evt is ExecutorFailedEvent failed)
{
Console.Error.WriteLine($"节点 {failed.ExecutorId} 执行失败: {failed.Data}");
}
else if (evt is WorkflowErrorEvent error)
{
Console.Error.WriteLine($"工作流执行失败: {error.Exception}");
}
}
动态构造与跨进程恢复
到目前为止,所有例子都是代码写死的固定图,new WorkflowBuilder(A).AddEdge(A,B).Build(),拓扑在编译期就定死了。但真实业务里,有一类需求很常见,而且恰恰是 workflow 最有价值的场景:
- 拓扑由数据决定:有一批预定义的执行器(抓微博、抓新闻、分析情绪、出报告),但这次该连哪几个、按什么顺序连,要等运行时看输入才知道。比如用户选了"只看微博"就连
微博→情绪→报告;选了"微博+新闻"就用 FanOut 同时连两个数据源。 - 后台长流程:workflow 要在后台跑很久(几分钟、几小时),中间程序可能重启,重启后要能从上次的断点接着跑,而不是从头来过。
这两个能力 SDK 都支持。但要分清楚,它们是两件独立的事,而且有一个关键限制必须先讲明白,否则会用错。
WorkflowBuilder 不是什么"只能在启动时调一次"的特殊对象,它就是个普通 class,所有 AddEdge/AddFanOutEdge/WithOutputFrom 都是 public 方法。你完全可以在运行时,根据输入数据,用代码拼出一张图再 Build:
// 预定义的执行器池(各自独立、可复用)
var fetchWeibo = new FetchWeiboNode();
var fetchNews = new FetchNewsNode();
var sentiment = new SentimentNode();
var report = new ReportNode();
// 根据用户选择,运行时决定拓扑
Workflow BuildWorkflow(UserOptions opts)
{
var builder = new WorkflowBuilder(sentiment); // 情绪分析为起点
if (opts.UseWeibo && opts.UseNews)
{
// 两个数据源都开 → FanOut 同时抓,都汇到情绪分析
builder.AddFanOutEdge(/*某个分发源*/, [fetchWeibo, fetchNews]);
builder.AddFanInBarrierEdge([fetchWeibo, fetchNews], sentiment);
}
else if (opts.UseWeibo)
{
builder.AddEdge(fetchWeibo, sentiment); // 只抓微博
}
else
{
builder.AddEdge(fetchNews, sentiment); // 只抓新闻
}
builder.AddEdge(sentiment, report);
return builder.WithOutputFrom(report).Build();
}
// 每次请求,动态构造一张图
Workflow wf = BuildWorkflow(currentUserOptions);
关键是建立这个认知:WorkflowBuilder 产出的 Workflow 是不可变的,但怎么构造它是完全自由的。你可以把构造逻辑包成一个函数,根据配置/数据库/用户输入动态拼接,拓扑是运行时决定的,不是写死的。
还有一种更精细的玩法,占位符 + 后绑定。AddEdge 接收的 ExecutorBinding 可以由字符串隐式转换而来,所以你能先用字符串 ID 把骨架连好,之后再调 builder.BindExecutor(...) 用真实执行器替换它。Build() 时如果有占位符没被替换,会报错。
var builder = new WorkflowBuilder("start");
builder.AddEdge("start", "analyzer"); // "analyzer" 现在只是个占位 ID
builder.AddEdge("analyzer", "report");
// ... 后期决定 "analyzer" 到底用哪个实现 ...
builder.BindExecutor(realAnalyzerBinding); // 用同 ID 的真实 binding 替换占位符
Workflow wf = builder.Build(); // 全部占位符都已替换,校验通过
这种 "先占位、后绑定" 的模式适合骨架必须先于实现确定"的场景。但说实话,日常更多用的是前面那种(直接根据条件构造整张图),占位符主要在框架内部或可视化编排器里才用得上。了解它存在即可,不用强求。
重启恢复
MAF 的 Workflow 可以做到重启恢复执行,我们将 Workflow 执行到某个位置称为检查点。
然后根据检查点位置进行持久化保存和恢复。
using Microsoft.Agents.AI.Workflows.Checkpointing;
// checkpoint 存到磁盘目录(每个 checkpoint 一个 .json 文件 + 一个 index.jsonl 索引)
DirectoryInfo checkpointDir = Directory.CreateDirectory("./checkpoints");
using FileSystemJsonCheckpointStore store = new(checkpointDir);
CheckpointManager mgr = CheckpointManager.CreateJson(store);
// 关键:把 mgr 传给执行入口,每个超级步结束自动存盘
await using StreamingRun run = await InProcessExecution.RunStreamingAsync(
workflow, input, mgr);
// 监听 SuperStepCompletedEvent,拿到 checkpoint 引用(用来后续恢复)
List<CheckpointInfo> checkpoints = [];
await foreach (WorkflowEvent evt in run.WatchStreamAsync())
{
if (evt is SuperStepCompletedEvent step && step.CompletionInfo?.Checkpoint is { } cp)
checkpoints.Add(cp);
}
CheckpointInfo 只是一个轻量引用(sessionId + checkpointId),真正的状态在 store 的文件里。你可以把它存进数据库、关联到你的业务记录,下次用这个引用去恢复。
ResumeStreamingAsync 接收的是 Workflow 实例 + CheckpointInfo + CheckpointManager,所以可以根据检查点恢复执行。
// ===== 可能是几天后、另一次进程启动 =====
// 1. 重新构造一张拓扑一致的 workflow(用和当初一样的构造逻辑!)
Workflow wf = BuildWorkflow(savedUserOptions);
// 2. 用同一个 store(指向同一个磁盘目录)
using FileSystemJsonCheckpointStore store = new(checkpointDir);
CheckpointManager mgr = CheckpointManager.CreateJson(store);
// 3. 用当初存的 checkpoint 引用,从断点恢复
StreamingRun resumed = await InProcessExecution.ResumeStreamingAsync(
wf, savedCheckpointInfo, mgr);
// 接着消费事件,就像没断过一样
await foreach (WorkflowEvent evt in resumed.WatchStreamAsync())
{
if (evt is WorkflowOutputEvent o)
Console.WriteLine(o.Data);
}