Files
memind/docs/agent-job-sse-integration.md
john 1798c07d42 docs: Update documentation and release rules
- 架构和规划文档更新
- 开发、工程、生产发布规则更新
- 服务隔离和升级指南
- README 更新
2026-06-27 08:25:06 +08:00

4.3 KiB

Agent Job SSE 流集成指南(路④)

后端已实现:长任务可异步化,前端通过 Server-Sent Events (SSE) 订阅进度。

后端端点

GET /mindspace/v1/agent/jobs/:jobId/stream

流式返回进度事件,直到任务终态或客户端断开。

事件格式

event: progress
data: {"id":"...", "status":"running", "progress":{"stage":"analyzing"},...}

event: done
data: {"id":"...", "status":"completed", "resultPageId":"..."}

事件类型

event 何时 payload
progress 状态/进度变化时推送 完整 job 对象
done 终态(completed/failed/cancelled/timed_out) 完整 job 对象(含 errorCode)

HTTP headers

  • Content-Type: text/event-stream
  • Cache-Control: no-cache, no-transform
  • X-Accel-Buffering: no (告诉代理别缓冲)

前端集成示例(React)

import { useEffect, useState } from 'react';

function AgentJobProgress({ jobId }) {
  const [job, setJob] = useState(null);
  const [error, setError] = useState(null);

  useEffect(() => {
    if (!jobId) return;

    const eventSource = new EventSource(
      `/mindspace/v1/agent/jobs/${jobId}/stream`
    );

    eventSource.addEventListener('progress', (event) => {
      const data = JSON.parse(event.data);
      setJob(data);
      console.log(`[progress] ${data.status}:`, data.progress);
    });

    eventSource.addEventListener('done', (event) => {
      const data = JSON.parse(event.data);
      setJob(data);
      console.log(`[done] ${data.status}`);
      eventSource.close();
    });

    eventSource.addEventListener('error', (event) => {
      if (eventSource.readyState === EventSource.CLOSED) {
        console.log('[stream closed by server]');
      } else {
        setError(`Stream error: ${event.message || 'unknown'}`);
      }
      eventSource.close();
    });

    return () => eventSource.close();
  }, [jobId]);

  if (error) return <div className="error">{error}</div>;
  if (!job) return <div>等待中...</div>;

  const isTerminal = ['completed', 'failed', 'cancelled', 'timed_out'].includes(
    job.status
  );

  return (
    <div className="job-progress">
      <div className="status">{job.status}</div>
      {job.progress && (
        <div className="stage">{job.progress.stage || job.status}</div>
      )}
      {job.errorMessage && <div className="error">{job.errorMessage}</div>}
      {job.resultPageId && (
        <div className="result">
           结果页面已生成:{job.resultPageId}
        </div>
      )}
      {!isTerminal && <div className="spinner">处理中...</div>}
    </div>
  );
}

export default AgentJobProgress;

使用流程

1. 长任务投递(返回 202 Accepted)

const jobRes = await fetch('/mindspace/v1/agent/jobs', {
  method: 'POST',
  body: JSON.stringify({
    jobType: 'generate_page',
    instruction: '生成一个海报',
    outputCategoryId: '...',
  }),
});
const { id: jobId } = await jobRes.json();

2. 立即返回,前端订阅进度

// 路由切换到进度页面,挂载 <AgentJobProgress jobId={jobId} />
// SSE 连接建立,收到 progress/done 事件

3. 完成后展示结果

// done 事件里有 resultPageId / resultAssetId
// 导航到 /mindspace/pages/:pageId 查看结果

微信场景(异步 ACK)

现有机制(wechat-mp.mjs 的 ACK + asyncTaskId):

  1. 收微信消息 → 立即 ACK
  2. 投递 job → 返回 asyncTaskId
  3. Job 完成 → Webhook 通知用户(via 客服接口)

这里 SSE 是Web/H5 专用 —— 微信端继续用现有异步 ACK + 客服消息回推(已实现)。

错误处理

eventSource.addEventListener('error', (event) => {
  if (eventSource.readyState === EventSource.CLOSED) {
    // 服务器正常关闭(终态)
    return;
  }
  // 真的错误
  console.error('SSE error:', event);
  retryOrShowError();
});

性能提示

  • SSE 连接建立快于 WebSocket,更轻
  • 不需要 polling(没有网络浪费)
  • 浏览器自动重连(如服务器 5s 内无数据)
  • 移动浏览器背景标签页会冻结连接(正常行为)

后端已支持的环境变量

# 轮询间隔(毫秒,默认 1000)
MINDSPACE_AGENT_SSE_POLL_MS=500

测试

# 手工测试(curl)
curl -N http://localhost:3000/mindspace/v1/agent/jobs/JOB_ID/stream

# 应该立即看到 progress 事件,然后在终态时看到 done 事件