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}");
    }
}


输出:

image-20260818090414362


这个例子虽然小,但已经包含了 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}");
    }
}


image-20260818092839697


实际上,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);
}