把 80 屏 Agent 日志变成一张图:Spring Boot SSE + React Flow 的 Run Canvas 实现

前一篇聊到一个我很认同的方向:复杂 Agent 不能只靠聊天窗口展示状态。

这篇不继续讲概念,直接做一个能跑的最小版本。

目标很简单:

Spring Boot后端
→持续产生Agent事件
→SSE推给浏览器
→React Flow实时更新节点
→右侧显示Artifact、预算和待审批事项

不是为了做一个漂亮 Demo,而是验证一个设计:Agent UI 应该是运行状态的投影,而不是把日志重新排版。

最终要看到什么

假设 Agent 正在执行“升级 Java 25”任务。

前端节点:

Assess       DONE
Plan         DONE
Modify       RUNNING
Test         PENDING
Review       PENDING
Deploy       BLOCKED

事件进来以后,节点颜色和标签实时更新。

右侧同时显示:

Changed files: 6
Tests: 128 passed / 3 failed
Model calls: 14
Cost: $2.31
Approval: deploy-canary

用户不需要翻聊天记录。


技术栈

前端使用当前 React Flow 包:

npm install @xyflow/react

后端用 Spring WebFlux 的 ServerSentEvent

为什么这里选 SSE,不选 WebSocket?

因为 Run Canvas 的主要实时方向是:

Server → Browser

用户的审批、取消、重试仍然可以走普通 HTTP。

SSE 足够简单。

如果后面需要高频双向协作,再换 WebSocket。

先定义事件,不要先写 UI

这是最重要的一步。

public enum AgentRunEventType {
    RUN_CREATED,
    PLAN_CREATED,
    STEP_STARTED,
    STEP_COMPLETED,
    STEP_FAILED,
    ARTIFACT_CREATED,
    APPROVAL_REQUESTED,
    APPROVAL_RESOLVED,
    BUDGET_UPDATED,
    RUN_COMPLETED,
    RUN_FAILED
}

事件结构:

public record AgentRunEvent(
        String eventId,
        String runId,
        long sequence,
        AgentRunEventType type,
        String stepId,
        JsonNode payload,
        Instant occurredAt) {
}

sequence 很重要。

不要只用时间戳排序。

多个 Worker 的时钟不一定严格一致。

事件应该先持久化,再推送

错误设计:

Agent执行
→直接SSE推前端

用户刷新页面以后,前面的状态全没了。

正确设计:

Agent执行
→Event Store
→Projection
→SSE

页面第一次打开:

GET /runs/{runId}

拿当前 Projection。

然后:

GET /runs/{runId}/events

订阅增量。

Projection 模型

public record AgentRunView(
        String runId,
        String title,
        RunStatus status,
        long lastSequence,
        List steps,
        List artifacts,
        BudgetView budget,
        List approvals) {
}

Step:

public record StepView(
        String stepId,
        String name,
        StepStatus status,
        List dependencies,
        String summary,
        Instant startedAt,
        Instant completedAt) {
}

Event Store 最小表

create table agent_run_event (
    event_id varchar(128) primary key,
    run_id varchar(128) not null,
    sequence bigint not null,
    event_type varchar(64) not null,
    step_id varchar(128),
    payload jsonb not null,
    occurred_at timestamptz not null,
    unique(run_id, sequence)
);

create index idx_run_event_run_seq
    on agent_run_event(run_id, sequence);

生产系统中,Sequence 应由数据库、Run Actor 或专门序列服务生成,不能用:

System.currentTimeMillis()

凑。

Spring Boot 的 SSE Endpoint

@RestController
@RequestMapping("/api/runs")
public class AgentRunStreamController {

    private final AgentRunEventStream stream;

    public AgentRunStreamController(
            AgentRunEventStream stream) {
        this.stream = stream;
    }

    @GetMapping(
            value = "/{runId}/events",
            produces = MediaType.TEXT_EVENT_STREAM_VALUE)
    public Flux> events(
            @PathVariable String runId,
            @RequestHeader(
                    value = "Last-Event-ID",
                    required = false)
            String lastEventId) {

        return stream.subscribe(runId, lastEventId)
                .map(event ->
                        ServerSentEvent
                                .builder(event)
                                .id(event.eventId())
                                .event(event.type().name())
                                .build());
    }
}

不要只用一个内存 Sinks.Many

Demo 最容易这么写:

Sinks.Many sink =
        Sinks.many()
             .multicast()
             .onBackpressureBuffer();

单机演示没问题。

但生产会遇到:

用户连到Pod A
Agent事件产生在Pod B

所以事件事实最好先进入数据库/Kafka/Redis Stream,再由所有 Web 节点消费。

一个简单的 Stream Service

@Service
public class AgentRunEventStream {

    private final EventRepository repository;
    private final LiveEventBus liveEventBus;

    public Flux subscribe(
            String runId,
            String lastEventId) {

        Mono startSequence =
                repository.resolveSequence(
                        runId,
                        lastEventId)
                        .defaultIfEmpty(0L);

        return startSequence.flatMapMany(seq -> {
            Flux history =
                    repository.findAfter(
                            runId,
                            seq);

            Flux live =
                    liveEventBus.events(runId);

            return Flux.concat(history, live)
                    .distinct(
                            AgentRunEvent::eventId);
        });
    }
}

重点是:

先补历史
再接实时

还要用 Event ID 去重。

SSE 断线恢复

浏览器 EventSource 会自动重连。

服务端要让客户端能够告诉你:

我最后收到哪个Event

SSE 原生有 Last-Event-ID 语义。

所以每个事件都设置:

.id(event.eventId())

重连后只补缺失事件。

还需要 Heartbeat

如果几分钟没有 Agent 事件,代理或网关可能断开长连接。

可以每 15 秒发一个 heartbeat:

Flux> heartbeat =
        Flux.interval(Duration.ofSeconds(15))
                .map(i ->
                        ServerSentEvent
                                .builder()
                                .comment("heartbeat")
                                .build());

与事件流合并。

前端安装 React Flow

npm install @xyflow/react

样式:

import '@xyflow/react/dist/style.css';

如果使用 Tailwind 4 和最新 React Flow UI,样式组织方式可以按当前文档放到全局 CSS。

前端事件类型

export type RunEvent = {
  eventId: string;
  runId: string;
  sequence: number;
  type: string;
  stepId?: string;
  payload: Record;
  occurredAt: string;
};

Projection 不要完全依赖后端每次传整棵图

后端可以传事件:

{
  "type": "STEP_COMPLETED",
  "stepId": "inspect-build",
  "sequence": 18
}

前端 Reducer 应用增量。

function reduceRun(
  state: RunView,
  event: RunEvent
): RunView {
  switch (event.type) {
    case 'STEP_STARTED':
      return updateStep(
        state,
        event.stepId!,
        { status: 'RUNNING' }
      );

    case 'STEP_COMPLETED':
      return updateStep(
        state,
        event.stepId!,
        { status: 'DONE' }
      );

    case 'STEP_FAILED':
      return updateStep(
        state,
        event.stepId!,
        { status: 'FAILED' }
      );

    default:
      return state;
  }
}

React Flow 节点转换

import type { Node, Edge } from '@xyflow/react';

export function toFlow(
  run: RunView
): { nodes: Node[]; edges: Edge[] } {

  const nodes: Node[] = run.steps.map(
    (step, index) => ({
      id: step.stepId,
      position: {
        x: index * 240,
        y: 120,
      },
      data: {
        label: step.name,
        status: step.status,
        summary: step.summary,
      },
      type: 'agentStep',
    })
  );

  const edges: Edge[] = [];

  for (const step of run.steps) {
    for (const dep of step.dependencies) {
      edges.push({
        id: `${dep}-${step.stepId}`,
        source: dep,
        target: step.stepId,
      });
    }
  }

  return { nodes, edges };
}

自定义节点

function AgentStepNode({ data }: NodeProps) {
  return (


        {data.label}



        {data.status}


      {data.summary && (

          {data.summary}

      )}

  );
}

不要只靠颜色表达状态。

同时显示文字:

PENDING
RUNNING
DONE
FAILED
BLOCKED
WAITING_APPROVAL

EventSource Hook

浏览器原生 EventSource 不能方便地设置任意 Header。

如果认证依赖 Cookie,可以直接用:

useEffect(() => {
  const source = new EventSource(
    `/api/runs/${runId}/events`,
    { withCredentials: true }
  );

  source.onmessage = (message) => {
    const event = JSON.parse(message.data);
    dispatch(event);
  };

  source.onerror = () => {
    console.warn('run stream disconnected');
  };

  return () => source.close();
}, [runId]);

如果你使用 Bearer Token,通常需要 fetch streaming、同源 Session 或专门的 SSE Client,而不是把 Token 塞进 URL。

页面第一次加载怎么做

不要等 SSE 从头重放所有事件。

const initial = await fetch(
  `/api/runs/${runId}`
).then(r => r.json());

setRun(initial);

这个接口返回当前 Snapshot。

然后 SSE 只接:

lastSequence之后的事件

Snapshot + Event 是比较实用的组合

只有 Event:

一个运行三天的Agent
打开页面要重放几十万条事件

只有 Snapshot:

不知道历史发生了什么

组合:

Snapshot负责快速打开
Event负责历史与增量

审批不要通过 SSE 回传

SSE 是服务端到客户端。

用户审批走普通接口:

POST /api/runs/{runId}/approvals/{approvalId}

Body:

{
  "decision": "APPROVE",
  "expectedRunVersion": 29
}

服务端检查:

用户权限
Approval状态
Run Version
Token
过期时间

成功后产生:

APPROVAL_RESOLVED

SSE 再把新状态推回来。

右侧 Artifact 面板

不要把 Artifact 正文全部塞进 Node。

Node 只显示:

6 files changed

点击后右侧打开:

pom.xml
Dockerfile
.github/workflows/build.yml

Artifact 后端模型:

public record ArtifactView(
        String artifactId,
        String type,
        String title,
        String contentHash,
        String previewUrl,
        Instant createdAt) {
}

Budget 也应该是一等公民

页面顶部:

Tokens      184K / 300K
Cost        $2.31 / $5.00
Model calls 14 / 25
Tool calls  28 / 50

事件:

{
  "type": "BUDGET_UPDATED",
  "payload": {
    "cost": 2.31,
    "tokens": 184302
  }
}

这样用户会第一次真正意识到:

Agent执行是有预算的

前端不要显示完整 Chain-of-Thought

Run Canvas 需要展示的是:

  • 步骤;
  • 结果;
  • 证据;
  • Tool;
  • Artifact;
  • 错误;
  • 决策。

不是模型内部私有推理。

例如:

Reason:
3个Auth测试失败,均与Refresh Token并发有关。

就够了。

状态图和聊天怎么联动

用户在 Chat 里说:

先不要改Dockerfile。

后端不是只把它追加成消息。

应该产生:

CONSTRAINT_ADDED
PLAN_UPDATED
STEP_CANCELLED

Canvas 立即看到 Plan 变化。

这才叫共享状态。

一个真正需要处理的坑:乱序事件

如果多 Worker 并发发布:

sequence 21
sequence 23
sequence 22

前端不能按到达顺序应用。

最简单的办法是服务端保证 Run 内 Sequence 有序后再发。

否则前端要维护 Buffer:

if (event.sequence === lastSequence + 1) {
  apply(event);
} else {
  buffer(event);
  refetchSnapshot();
}

我更倾向服务端处理。

第二个坑:页面开了几个小时

Run 可能已经归档。

服务端发送:

RUN_COMPLETED

前端收到终态以后:

source.close();

不要让数千个完成任务继续占着 SSE 连接。

第三个坑:Event Store 无限增长

保留策略可以是:

近期Run:完整事件
长期:压缩事件+关键Audit
事故Run:永久保留

每 N 个事件写一个 Snapshot。

推荐的接口集合

POST /api/runs
GET  /api/runs/{runId}
GET  /api/runs/{runId}/events
GET  /api/runs/{runId}/artifacts
POST /api/runs/{runId}/cancel
POST /api/runs/{runId}/approvals/{id}
POST /api/runs/{runId}/retry/{stepId}

前端基本就够用了。

测试别只测页面能动

至少覆盖:

SSE断线重连
Last-Event-ID补发
重复事件去重
乱序
Snapshot与Event一致
多用户权限
审批版本冲突
Run完成关闭连接
Artifact越权
Budget实时更新

我会怎么继续扩展这个 Demo

第一阶段先做:

Run
Step
SSE
React Flow

第二阶段加:

Artifact
Approval
Budget

第三阶段再加:

Trace
Diff
Replay
Multi-Agent

不要第一天就做一个“AI IDE”。

最后

这套实现真正有价值的地方,不是 React Flow。

React Flow 只是把节点画出来。

真正的基础是:

结构化Run State
+不可变Event
+Snapshot
+审批
+Artifact

一旦这些后端语义存在,你可以把前端换成:

  • React Flow;
  • 看板;
  • 时间线;
  • 表格;
  • IDE 插件。

都没关系。

Agent 产品真正该摆脱的是:

所有状态都藏在聊天字符串里

当执行过程变成结构化状态以后,人才能真正看得懂、插得上手,也才敢让 Agent 跑更长的任务。


更多企业级 AI 应用、Agent、RAG 与模型工程化内容,我会继续整理在 智元界

https://www.zyentor.com/