01

一次请求的旅程

从 HTTP 入口到 ReActAgent,再把答案一段段流回浏览器

先认识这套系统

promo-smart-agent 是一个促销营销智能体:运营在 Web / 飞书里问一句话,它就能做活动数据问答、定时报告。

内核是 AgentScopeReActAgent ——它是四条渠道共同的终点,整套系统都围绕它做编排。

🌐
多渠道接入

Web、飞书、定时任务(Cron)、A2A 协议,四个入口同一个大脑

⚙️
配置驱动编排

主 Agent、子 Agent、模型、工具大多靠 YAML 配置拼出来,而非硬编码

📡
流式输出

答案不是一次性返回,而是一个个 token 边生成边推给前端

🧭
本模块只看一条主干

四条渠道殊途同归。这一节我们只跟最直观的 Web 流式请求,把它从入口一路拆到 ReActAgent。其余渠道在模块 4 展开。

你在界面里敲下一句话,然后呢?

你在 Web 界面输入一句话点发送,浏览器里字一个个往外蹦——这一节就把这条流式链路从入口拆到 ReActAgent。

🛄
隐喻:机场行李传送带

你在值机口(HTTP 端点)交出一句话,它沿传送带经过安检分拣(Service / SessionManager / Factory),最后被装上飞机(ReActAgent)。返回时不是一次性全到,而是一件件行李(SSE token)从传送带口陆续吐出来。

看组件互相「喊话」

这条链上有 5 个角色。点「下一条」,看它们如何把一句话传下去、再把 token 逐个传回来。

主干就这一条

把上面的对话压成一行调用链——记住它,后面所有模块都是在给这条链上的某个角色补细节。

1

AgentController

2

AgentChatService

3

SessionManager.getAgent

4

ReActAgent.stream()

5

Flux<Event> → emitter.send

💡
关键认知:流式 = 逐事件推送

「流式」不是魔法。它靠 reactorFlux<Event> 逐个 Event 触发,每个事件到达就 SseEmitter .send 一次。所以前端才看得到「字一个个出来」,而不是干等到最后一次性全出。

入口长什么样

先看值机口。这是 Web 的 SSE 流式端点——它只做三件事,然后把 SseEmitter 直接交还给浏览器。

CODE · AgentController.java:75
@PostMapping(value = "/chat/stream", produces = "text/event-stream")
public SseEmitter chatStream(@RequestBody @Valid AgentChatRequest request) {
    String modelRef = normalizeAndValidateModelRef(request.getModelRef());
    String sessionId = SessionId.of(CallerSourceEnum.WEB, request.getSessionId()).value();
    return agentChatService.chatStreamWithMetrics(
            request.getMessage(), sessionId, request.getPromptTemplateName(), request.getDeviceId(), modelRef);
}
中文逐行

声明这是 /chat/stream 的 POST 端点,响应类型是事件流——告诉浏览器「我会持续推数据」。

方法返回 SseEmitter:返回它之后连接不关,服务端可以慢慢往里送。

先校验前端传来的模型引用,挑出这次该用哪个模型。

把前端的 sessionId 规范成 web:{...} 格式——渠道前缀让不同入口的会话互不串台。

活全交给 Service 去干。Controller 只是值机口,不碰业务。

🔑
为什么 Controller 这么薄?

这是「 sessionId 规范 + 转交」的瘦控制器模式:入口只负责接参数、定渠道,真正的取 Agent、调模型、推流都压到 Service。这样换渠道(飞书 / A2A)时业务逻辑能复用。

真正调 ReActAgent 的地方

请求落到 Service 后,关键就这一段:拿 Agent 产出 Flux<Event>,并在全部产出完时把累计指标写回。

CODE · AgentChatService.java:698
Flux<Event> events = agent.stream(userMsg, STREAM_OPTIONS)
        .contextWrite(ctx -> writeAgentContext(ctx, resolved))
        .doOnNext(event -> { /* 累积输出文本 */ })
        .doOnComplete(() -> {
            AgentChatMetrics metrics = buildMetrics(resolved.getChatUsageRef().get(),
                    resolved.getChatUsageTurns(), System.currentTimeMillis() - start, modelName);
            resolved.getMetricsRef().set(metrics);
        });
return new AgentChatStream(traceId, events);
中文逐行

agent.stream(...) 让 ReActAgent 开始干活,产出一串会陆续到达的事件 Flux<Event>

把本次会话的上下文写进 reactor 的执行链里,下游处理时能读到。

doOnNext:每来一个 event 就顺手累积一次输出文本——这是「边到边处理」的钩子。

doOnComplete:所有 event 都产完后才触发。

这时算出耗时、token 用量等指标,组装成 metrics。

把 metrics 写回引用——埋点/统计就发生在这一刻。

把事件流包成 AgentChatStream 往上返回,真正的推流在上层逐个发。

📍
埋点该插在哪?记住这里

耗时、token 指标都在 doOnComplete → buildMetrics 这一处生成。想加全链路统计,这是天然的落点——后面 quiz 会用到。

让数据自己跑一遍

点「下一步」,看请求包从浏览器进 Controller,一路到 ReActAgent;再看多个 Event 逐个滑回浏览器——注意是一个一个回,不是一次性。

🌐
浏览器
C
AgentController
S
AgentChatService
SM
SessionManager
RA
ReActAgent
点「下一步」开始
🐛
排障直觉

看懂这条链,「字不出来了」就有了排查顺序:是 Controller 没接住?Service 没拿到 Agent?还是 ReActAgent 没产出 Event?逐段往下排即可。

动手想一想

三道题,都是真实场景。先想再选,错了也有讲解。

用户反馈:点发送后页面一直空白,最后才一次性蹦出全部文字。最可能是哪段坏了?

想给每次请求加「耗时 + token 用量」埋点统计,插在哪层最合适?

从「一句话」到「返回第一个 token」,组件被经过的正确顺序是?

➡️
下一站:演员表

这一节出场的 AgentChatService / AgentSessionManager / DefaultMainAgentFactory / ReActAgent,模块 2 会逐个拉出来单独介绍——它们各自是谁、负责什么。

02

演员表:核心组件与分层

模块 1 那串 Controller → Service → SessionManager → Factory → ReActAgent,到底谁负责什么?这一节给每个角色发名牌。

把这套系统想成一个剧组

三个 Maven 模块 各演一个角色,职责泾渭分明。看清谁是谁,就知道一段代码该落在哪。

interfaces · 舞台与售票口

🎫
HTTP / 飞书 / A2A 接入 · 打包入口

core · 所有演员和导演

🎬
Agent / Session / Cron / Service · 全部业务

api · 剧本与道具清单

📜
Dubbo 接口 + 轻量 DTO · 无业务
点任意一层,看它在剧组里演什么 —— 依赖方向是 interfaces → core → api,越往里越稳定。
💡
依赖只能朝里指

外层(interfaces)依赖内层(core、api),但内层永远不知道外层存在。这条单向规则叫“依赖倒置”—— 换掉舞台不影响演员,换售票口不影响剧本。

三张模块名牌

记一句话就够:api 只发清单、core 全是戏、interfaces 只管进出场

📜

promo-smart-agent-api

对外契约:Dubbo 接口 + 跨进程 DTO。供调用方依赖,不放任何业务逻辑

🎬

promo-smart-agent-core

核心模块:Agent 框架、Cron、A2A 桥接、Session 管理。所有业务逻辑都在此 —— 新建任何东西先想 core。

🎫

promo-smart-agent-interfaces

接入层:HTTP 控制器、Dubbo 实现、飞书、A2A 端点。打包入口(fat JAR),别往这里塞逻辑

这就是 pom.xml 里的声明

CODE · pom.xml:21
<modules>
    <module>promo-smart-agent-api</module>       <!-- 对外 DTO/接口,无业务 -->
    <module>promo-smart-agent-core</module>      <!-- 核心:Agent/Cron/Session/业务 -->
    <module>promo-smart-agent-interfaces</module> <!-- HTTP/飞书/A2A 接入,打包入口 -->
</modules>
大白话

一个父项目下挂三个子模块,构建时按顺序编译。

api:对外那份“清单”,纯 DTO 和接口,没有戏。

core:所有戏都在这 —— Agent、定时任务、会话、业务。

interfaces:把上面的戏接到 HTTP/飞书/A2A,并打成一个可运行的包。

运行期的四个角色

模块拆好了,跑起来时谁先谁后?这四张名牌串起一次对话的核心链路。

🎞️
AgentChatService —— 场记(门面)

门面 Service:封装流式 / 同步对话逻辑,让 HTTP 层不用关心 Agent 内部怎么跑。

🪑
AgentSessionManager —— 后台候场区(双层缓存)

把活跃的 Agent 暂存在 Guava LoadingCache 里,下次同一会话直接复用,不用重新化妆。

💄
DefaultMainAgentFactory —— 化妆间(工厂)

Factory:缓存没命中时,它按配置现场组装一个 Agent —— 装工具、挂子 Agent、接 MCP、跑扩展点。

🎭
ReActAgent —— 演员本体

真正“思考—调工具—再思考”的执行者。Factory 装扮好它,SessionManager 让它候场,场记叫它上台。

场记为什么要存在?

门面把“怎么对话”这件事收口,HTTP / Dubbo / 飞书都只调它 —— 改对话逻辑只改一处。

CODE · AgentChatService.java:59
/**
 * Agent 对话 Service
 * <p>封装流式推送和同步调用逻辑,与 HTTP 层解耦。
 * 所有方法均接收 promptTemplateName,透传至 AgentSessionManager 以确保模板切换运行时生效。</p>
 */
@Slf4j @Service @RequiredArgsConstructor
public class AgentChatService {
大白话

这就是“场记”—— 一个统一的对话入口。

它把流式推送和一次性同步两种调用都包好,外层 HTTP 不用碰内部细节(这叫“解耦”)。

每个方法都带上要用哪套提示词模板,再原样传给候场区,所以换模板立刻生效、不用重启。

@Service 让 Spring 启动时自动托管它,依赖自动注入。

候场区:命中就复用,没命中就现造

这是 L1 缓存 的心脏:按“多久没访问就过期”来回收,miss 时才叫化妆间造一个。

CODE · AgentSessionManager.java:68
@PostConstruct
public void init() {
    int ttlMinutes = agentProperties.getMainAgent().getCacheTtlMinutes();
    sessionCache = CacheBuilder.newBuilder()
            .expireAfterAccess(ttlMinutes, TimeUnit.MINUTES)
            .build(new CacheLoader<String, ReActAgent>() {
                @Override public ReActAgent load(String cacheKey) {
                    return createAndRestore(cacheKey);   // ← miss 时建 Agent
                }
            });
}
大白话

服务启动后自动跑一次,先从配置读“缓存能存活多久”。

建一个缓存,规则是:超过这个分钟数没人访问就自动清掉(TTL 到期回收)。

关键一行:给缓存装一个“加载器”。

取某个会话的 Agent 时,缓存里没有(miss),就调 createAndRestore 现场造一个 —— 这一步背后就是化妆间。

🧠
为什么不每次都新建?

造一个 Agent 很贵:要装工具、挂子 Agent、连 MCP、还要恢复历史记忆。同一会话连续对话时复用现成的,省掉这一整套装配 —— 缓存换的就是这个时间。

化妆间:一个 Agent 是怎么被“装扮”出来的

miss 那一刻,Factory 按配置把零件一件件装上 —— 子 Agent、MCP/Skill、编码扩展点,最后 build。

CODE · DefaultMainAgentFactory.java:141
for (String toolAgentRef : props.getToolAgentRefs()) {
    toolAgentRegistry.registerInto(toolAgentRef, session, toolkit);   // 子 Agent → 工具
}
AgentBuildUtils.registerMcpClients(props.getMcpServers(), toolkit, COMPONENT_NAME, mcpRegistry);
for (AgentConfigCustomizer customizer : customizers) {
    customizer.registerTools(toolkit, context);   // 编码扩展点
}
ReActAgent.Builder builder = ReActAgent.builder()
        .name(props.getName()).sysPrompt(sysPrompt).model(model)
        .toolkit(toolkit).memory(memory).maxIters(props.getMaxIters())
        .hooks(List.of(new McpToolLoggingHook(), new ChatUsageCollectorHook()));
return builder.build();
大白话

先按配置里列的子 Agent,把每个都注册成可调用的工具放进 工具箱(Toolkit)

再把外部 MCP 工具也接进同一个工具箱。

然后跑一遍编码扩展点(Customizer),允许开发者用代码再塞工具进来。

最后用 builder 把名字、 系统提示词、模型、工具箱、 记忆、最大迭代轮数、 钩子拼起来,build 出成品 Agent。

1

挂子 Agent

2

接 MCP / Skill

3

跑 Customizer

4

build 成品

想一想:你会把它放哪?

不背定义,只判断 —— 给你一个场景,凭刚学到的分层直觉做决定。

① 你要新增一个业务枚举和一个工具类,应该放进哪个模块?

② 为什么 SessionManager 要把 Agent 缓存起来,而不是每次对话都新建一个?

③ 飞书链路要用到的 SSE 流式协议类,适合放进 api 模块吗?

④ 你想改“新建 Agent 时多挂一个工具”的逻辑,应该从哪个角色下手?

🎯
记住这把尺子

给 AI 下指令时,你就能说清“把这个加到 core 的 XxxService,别写进 Controller”。下一节:Factory 里的 props.getToolAgentRefs() 和四大 Registry 是怎么来的 —— 答案是全靠 YAML 配置驱动。

03

配置驱动一切

加一个子 Agent、换一个模型 —— 改的是 YAML,不是 Java。

还记得那行 props.getToolAgentRefs() 吗?

模块2里,Factory 拿子 Agent 时读了一个列表。那个列表里的子 Agent —— 是从哪来的?

答案就一个:一段 YAML。

🔌
把它想成配电箱

YAML 就是配电箱面板,每条 models[] / tool-agents[] 是一路开关。要加一路用电器(子 Agent),合一个闸(加一段 YAML)就行 —— 不用重新布线(改代码)。

一段 YAML 的命运

跟着一条 tool-agents[] 配置走完它的一生:从文本,到主 Agent 手里一个能用的工具。

📄
YAML 配置
🧩
AgentProperties
📇
ToolAgentRegistry
🏭
Factory
🛠️
主 Agent 的工具
点 "下一步" 看这段 YAML 怎么活过来

配置的根:一个类,五大列表

整个系统"要装什么",都收敛到 @ConfigurationProperties 标注的这一个类里。

CODE
@Data @Component
@ConfigurationProperties(prefix = "agentscope")
public class AgentProperties {
    private List<ModelProperties> models = ...;
    private List<SessionProperties> sessions = ...;
    private List<MemoryProperties> memories = ...;
    private MainAgentProperties mainAgent;
    private List<ToolAgentProperties> toolAgents = ...;
}
中文翻译

这个注解的意思:凡是 agentscope.* 下的配置,都灌进这个类。

于是 YAML 里的每个区块,都对上一个字段 ——

models:你想用哪些大模型。

sessions / memories:会话怎么存、记忆怎么留。

mainAgent:主 Agent 长什么样。

toolAgents:有哪些子 Agent可以被主 Agent 当工具用。

注意:写完这个类,YAML 还只是数据。谁把它变成能跑的实例?

电工登场:Registry 在启动时接线

四大 Registry 各管一摊。看 ModelRegistry 怎么把 models[] 一条条变成真的模型。

CODE
@PostConstruct
public void init() {
    List<ModelProperties> configs = allModelConfigs();
    for (ModelProperties props : configs) {
        Model model = buildModel(props);
        models.put(props.getName(), model);
        if (props.isDefaultModel()) {
            if (defaultModel != null) {
                throw new IllegalStateException("... multiple models marked as default=true");
            }
            defaultModel = model;
            defaultModelName = props.getName();
        }
    }
}
中文翻译

@PostConstruct:启动时自动跑一次,做初始化。

把配置里所有模型条目拿出来,逐条处理。

buildModelapi-type 选实现 —— Anthropic 还是 OpenAI。

造好后用名字存进表里,以后按名字取。

如果这条标了 default-model,记成默认模型 ——

但若已经有一个默认了,直接启动报错。(记住这条,等下考你)

为什么这套设计很值钱

💡
配置驱动 = 把"要什么"和"怎么造"解耦

YAML 只负责说"要什么",Registry 负责"怎么造"。所以换模型只改 model: 一行,加子 Agent 只加一段配置 —— 造的逻辑早就写好了,你碰不到也不用碰。

models[] ModelRegistry → 一个个 ChatModel 实例
sessions[] SessionRegistry → 会话存储实例
tool-agents[] ToolAgentRegistry → 挂上主 Agent toolkit 的 SubAgent

最后一棒:registerInto 把子 Agent 接进工具箱

会话到来时,ToolAgentRegistry.registerInto() 按配置造出 SubAgent,挂进主 Agent 的工具箱。

CODE
public void registerInto(String name, Session session, Toolkit toolkit) {
    ToolAgentEntry entry = getEntry(name);
    Model model = modelRegistry.get(entry.props.getModel());
    SubAgentConfig subAgentConfig = SubAgentConfig.builder()
            .toolName(entry.props.getName())
            .description(enhanceDescription(...))
            .session(session)
            .forwardEvents(entry.props.isForwardEvents())
            .build();
    toolkit.registration()
            .subAgent(() -> { /* 用配置里 sysPrompt/model/skills build 出 ReActAgent */ }, subAgentConfig)
            .apply();
}
中文翻译

按名字取出这条 tool-agent 的配置条目。

配置里写的是模型名字,向 ModelRegistry 换成真模型实例。

用配置里的名字、描述拼出 SubAgent 的"工具说明"。

让它共享当前会话,并按配置决定是否转发事件。

最后注册进 toolkit —— 主 Agent 从此多了一个工具。

全程没碰一行业务逻辑:你给的只是配置

小测:你能判断"改配置还是动代码"吗?

这三题不考定义,考你怎么下手 —— 也是你将来给 AI 下指令的判断力。

① 产品要你加一个"活动数据查询"子 Agent。最小改动是什么?

② 要把默认模型从 A 换成 B,改哪里?

③ 手滑把两个模型都标了 default-model: true,会怎样?

🎯
带走这句话

下次给 AI 下指令,别说"帮我改代码加个子 Agent"。直接说:"在 tool-agents 加一条引用 xxx 的子 Agent" —— 你已经知道这是纯配置。下一站:配置之外,运行期怎么按"谁来的"和"哪个会话"分流。

04

渠道与会话

同一个 Agent 引擎,如何同时服务 Web、飞书、定时任务、A2A——还不串话、不丢记忆

一家酒店,四个入口

Web 上你和它聊、飞书群里 @它、定时任务里它自己跑——背后是同一套引擎

它怎么不把张三的对话串到李四头上?靠一张sessionId(房卡号)。

🏨
隐喻:酒店前台与房卡

不同入口(大堂 / 电话 / APP = 渠道)的客人都来同一家酒店(同一 Agent 引擎),靠房卡号(sessionId)认人。房间记忆分两层存:前台手边的活跃房卡盒(会过期)和后台档案库(长期保留)。

五个入口,五种房卡号

渠道用一个枚举 CallerSourceEnum 钉死成 5 种。每种渠道的 sessionId 长得不一样:

🖥️

Web

web:{UUID}
每个浏览器会话一个随机号。

💬

飞书(私聊)

feishu:{chatId}_{openId}
会话 + 人,两者一起认人。

👥

飞书群

feishu_group:{chatId}
只认群,不带个人 openId——全群共享一段记忆。

定时任务

cron:{UUID}
它自己定时跑,每次一个随机号。

🤝

A2A

a2a:{contextId}
被别的 Agent 调用,沿用对方传来的上下文号。

💡
关键差异:群里不带 openId

飞书群的房卡号故意只用 chatId、不拼 openId——所以群内每个人共用同一段上下文,A 问完 B 追问能接得上。这正是“串话 / 不串话”的设计开关。

渠道清单:一份枚举钉死五种来源

每个渠道还带一个 group 标志,决定“这次对话算谁的、通知发去哪”。

CODE
public enum CallerSourceEnum {
    WEB("web", "Web端", false),
    FEISHU("feishu", "飞书", false),
    FEISHU_GROUP("feishu_group", "飞书群", true),
    CRON("cron", "定时任务", false),
    A2A("a2a", "A2A Agent间调用", false);
    private final String code; private final String desc; private final boolean group;
中文

声明一份“来源渠道”的固定清单,一共 5 项。

Web 端:网页对话,不是群,所以 group=false。

飞书私聊:一对一,也不是群。

飞书群:group=true——这是唯一的群,归属和通知按群处理。

定时任务:它自己触发,单独一类。

A2A:被别的 Agent 调用时用,也算一对一。

每项都记三样:英文码、中文描述、是不是群。

拼房卡号:一个函数,且“拼两次也不出错”

SessionId.of 把渠道前缀和各部分拼成 {source}:{parts}。妙在它是幂等的——已经带前缀的,不会被拼成 web:web:

CODE
public static SessionId of(CallerSourceEnum source, String... parts) {
    String raw = String.join("_", parts);
    for (CallerSourceEnum existing : CallerSourceEnum.values()) {
        if (raw.startsWith(existing.getCode() + DELIMITER)) {   // 已带前缀就复用(幂等)
            return new SessionId(raw, existing);
        }
    }
    return new SessionId(source.getCode() + DELIMITER + raw, source);
}
中文

传入“哪个渠道”和若干片段,造一个房卡号。

先用下划线把片段拼成一串原始文本。

挨个看 5 种渠道前缀……

如果这串文本已经带了某个前缀(比如 web:)……

就直接拿来用,不再加前缀——这就是幂等,防 web:web:。

都没带前缀,才补上本渠道的前缀,拼成 {渠道}:{内容}。

记忆的一来一回(核心)

记忆分两层:L1(内存里的活跃 Agent,快但会过期)和 L2持久化档案)。点下面看一段记忆怎么存档、又怎么被找回。

L1
L1 内存缓存
(快·会过期)
L2
L2 Session
(持久档案库)
👤
用户 / 下一次请求
点 “下一步” 看记忆如何存档、又如何被找回
💡
L1 过期是正常的,不等于丢记忆

L1 设了 TTL(默认 30 分钟)到点就清——这是省内存的设计。真正保命的是 L2:只要 saveTo 存过、loadIfExists 能读回,记忆就不会丢。

存档这一步,代码长这样

每轮对话后,把 L1 里那个活跃 Agent 的记忆,写进 L2 档案库。

CODE
public void saveState(String sessionId, String promptTemplateName, String modelRef) {
    String cacheKey = buildTemplateCacheKey(sessionId, resolved.resolvedName, modelRef);
    ReActAgent agent = sessionCache.getIfPresent(cacheKey);   // L1
    if (agent == null) return;
    Memory memory = agent.getMemory();
    if (memory != null) {
        memory.saveTo(mainSession, sessionId);                // → L2 Session 持久化
    }
}
中文

对话结束后调用,参数告诉它“是哪个会话、用哪套提示、哪个模型”。

按这几样拼出 L1 缓存的钥匙。

用钥匙去 L1 内存里把活跃 Agent 取出来。

L1 里没有(已过期)就直接返回,不用存——它早被前一次存过了。

拿到 Agent 身上的这段对话记忆。

记忆确实存在的话……

就把它写进 L2 档案库(saveTo)——这一步让记忆能跨过期、跨重启活下来。

飞书没有公网 URL,怎么收到消息?

不是飞书来敲我们的门,而是我们主动拨出去——用 WebSocket 建一条长连接,服务一启动就拨号连上飞书服务器,事件顺着这条管子推过来。

CODE
EventDispatcher dispatcher = EventDispatcher.newBuilder("", "")
    .onP2MessageReceiveV1(new ImService.P2MessageReceiveV1Handler() {
        @Override public void handle(P2MessageReceiveV1 event) {
            feishuMessageHandler.handle(event, botCode);   // ← 消息进入主链路
        }
    }).build();
Client wsClient = new Client.Builder(appId, bot.getAppSecret()).eventHandler(dispatcher).build();
wsClient.start();   // 主动 outbound 连接飞书,无需公网 URL
中文

建一个“事件分发器”,约定不同事件分别交给谁处理。

当收到一条用户消息时……

……就执行这个处理方法。

把消息(带上是哪个机器人 botCode)交给主链路接着跑。

用机器人的账号密钥造一个客户端,挂上刚才的分发器。

start() 主动拨出去连飞书——我们这边没有公网地址也能收消息。

排障演练

用户反馈来了。基于刚才看到的 sessionId 拼法和双层缓存,你会先看哪儿?

场景一

飞书群里 A 问了一个问题、得到回答,B 紧接着追问,机器人却像没看到 A 的对话一样答非所问。最可能的根因在哪?

场景二

用户说:“我刷新(或隔了半小时再来)之后,之前的对话记忆好像没了。” 这更像是 L1 还是 L2 的问题?

场景三

新同事问:“我们这服务没申请公网域名,飞书的消息到底是怎么进来的?” 你怎么解释?

场景四

你在排查一个 SessionId 拼接逻辑,怀疑某处把已带前缀的字符串又拼了一次,得到 web:web:abc。看过 SessionId.of 后,你预期它会怎样?

小结

🎫
一引擎 + sessionId 区分

五渠道复用同一 Agent 引擎;CallerSourceEnum 定来源,SessionId.of 幂等拼房卡号、认人不串话。

🗄️
双层缓存记忆

L1 内存快但会过期;L2 Session 持久——saveTo 存、loadIfExists 取,跨过期跨重启不丢。

🔌
飞书无公网

WebSocket 长连接主动拨出,事件顺管子进主链路。

下一站 → 模块 5:扩展点与聪明工程:渠道隔离的 BundleTool、采集 token 的 Hook、Cron 的锁与重试。

05

扩展点与聪明工程

想加能力,先问"该挂在哪个卡口"——再看 Cron 那几处值得偷师的工程手法

主 Agent 是电动工具,扩展点是卡口

Agent 像一把电钻手柄,本身不干活——能力来自换上不同的"头"。选对卡口,能力即插即用。

🔌
开场难题

"想让它会建定时任务,但只在飞书 / 定时渠道开放、网页端不暴露——这种按渠道隔离的能力怎么挂上去?" 答案是一个叫 BundleTool 的卡口。这一节就拆这些卡口。

四个卡口:想加能力,先认这四个口子

别一股脑塞进消息处理器。加能力前先问一句:"这是哪一类能力?"——然后对号入座。

🧰

BundleTool

系统级工具 · 按渠道隔离 · 自动登记。
何时用:要给主 Agent 加一组能力,但只想在某些渠道开放。
例子:cronCreate 建定时任务,只在飞书 / 定时渠道挂载。

🪛

AgentTool

挂给子 Agent 的工具头。
何时用:某个能力只给某个子 Agent 用,不污染主 Agent。
例子:给"选品子 Agent"配一把查商品库的工具。

⚙️

AgentConfigCustomizer

编码扩展头:手写注工具 / 注 Hook。
何时用:配置 YAML 表达不了、需要写 Java 逻辑时。
例子:在 Agent 构建时,用代码塞一个自定义工具或监听器进去。

🎧

Hook

旁路监听头:不改主流程,只在旁边偷听。
何时用:要采集 / 观察,但不想插手主逻辑。
例子:ChatUsageCollectorHook 旁路统计每轮花了多少 token

💡
记住这条决策线

加能力 → 给谁用?给主 Agent 且要分渠道选 BundleTool;给某个子 Agent选 AgentTool;纯代码逻辑选 Customizer;只想旁观不插手选 Hook。

读卡口①:BundleTool 怎么做到"只在飞书开放"

两件事就够了:用 @Component 自动登记、用 channels() 声明挂载渠道。

CODE
@Slf4j @Component
public class CronBundleTool implements BundleTool {
    @Override public String name() { return "bundle-cron"; }
    @Override public List<CallerSourceEnum> channels() {
        return List.of(CallerSourceEnum.FEISHU, CallerSourceEnum.FEISHU_GROUP, CallerSourceEnum.CRON);
    }
    @Tool(description = "Create a scheduled task. taskName...; cronExpr: cron expression ...")
    public CronResult cronCreate(AgentContext ctx, @ToolParam(name = "taskName") String taskName /*...*/) {
        // ...
    }
}
大白话

@Component:贴上这张"自动登记"标签,启动时系统自己发现它,没人需要手动挂。

name():这套工具的代号叫 bundle-cron

channels():核心一行——声明"我只在飞书、飞书群、定时这三个渠道露面"。网页端列表里根本看不到它。

@Tool:每个带这个标签的方法,就是 Agent 能直接调用的一个工具动作,这里是"建一个定时任务"。

🛡️
为什么这很优雅

渠道准入不是写一堆 if 判断挡出来的,而是工具自己声明能去哪。模块 4 里的渠道枚举(CallerSourceEnum)在这里被直接复用。

聪明工程:定时任务的"三道闸门"

定时任务最怕两件事:同一时刻多台机器都收到、重复跑一遍;以及上一轮还没跑完、下一轮又来了。Cron 用三道闸挡住它们。

幂等锁(防"撞车")

Redis 抢一个号牌,抢到的才往下提交,避免同一时刻多实例都触发同一个任务。

互斥锁(防"叠跑")

独占型任务执行前再抢一把 分布式互斥锁,抢不到说明"上一轮还在跑",本轮直接跳过。

重试 3 次(抗抖动)

真正执行时若失败,自动重试最多 3 次,给网络抖动 / 临时故障一个自愈的机会。

看它们"对话":一个任务怎么通过三道闸

点"下一条",看 Scheduler、Redis、Executor、互斥锁、Handler 之间的对话怎么把任务安全送达。

读第一道闸:幂等锁就是一行"抢号牌"

多台机器同时被定时器叫醒,凭什么只有一台真正执行?看这一行 幂等锁

CODE
for (CronTaskPO task : tasks) {
    String idempotencyRedisKey = IDEMPOTENCY_KEY_PREFIX + task.getId() + ":" + idempotencyKey;
    Boolean acquired = stringRedisTemplate.opsForValue()
            .setIfAbsent(idempotencyRedisKey, "1", Duration.ofMinutes(2));
    if (!Boolean.TRUE.equals(acquired)) continue;   // 锁已存在 → 跳过
    cronTaskExecutor.submitTask(task, CronTriggerType.SCHEDULED);
}
大白话

对每个到点的任务,先拼一个唯一的号牌名:任务 ID + 这一拍的时间戳。

setIfAbsent:去 Redis 抢号牌——"如果没人占,就占下,挂 2 分钟后自动过期"。

没抢到(别的机器先到一步)?continue 直接跳过,绝不重复跑。

抢到了,才把任务提交去执行。这就保证了"同一拍只有一台机器真正触发"。

关键顿悟:两把锁,挡的是两件事

💡
幂等锁 ≠ 互斥锁

幂等锁防的是"同一时刻多实例都收到同一个任务"——横向的撞车。互斥锁防的是"上一轮还没跑完,下一轮又来一轮"——纵向的叠跑。两道闸解决的是两个完全不同的问题,缺一不可。

↔️

幂等锁 · 防撞车

维度:横向 / 同一时刻
问题:3 台机器同拍触发同一任务。
动作:Redis 抢号牌,只放一个过。

↕️

互斥锁 · 防叠跑

维度:纵向 / 跨时间
问题:上一轮慢、还没结束,下轮又起。
动作:抢不到独占锁就 SKIPPED_MUTEX 跳过。

另一处聪明:token 统计为什么用 Hook,不在 Service 里循环加

一次对话主 Agent 会喊好几个子 Agent,还是异步流式的。想把所有层的 token 加全,靠一个 Hook + reactor Context

CODE
public <T extends HookEvent> Mono<T> onEvent(T event) {
    if (!(event instanceof PostReasoningEvent reasoningEvent)) return Mono.just(event);
    ChatUsage usage = reasoningEvent.getReasoningMessage() != null
            ? reasoningEvent.getReasoningMessage().getChatUsage() : null;
    if (usage == null) return Mono.just(event);
    return Mono.deferContextual(ctxView -> {
        if (ctxView.hasKey(USAGE_REF_KEY)) {                 // 上层注入的 ref
            AtomicReference<ChatUsage> ref = ctxView.get(USAGE_REF_KEY);
            ref.accumulateAndGet(usage, ChatUsageCollectorHook::accumulate);   // 累加总量
        }
        return Mono.just(event);
    });
}
大白话

只在 PostReasoningEvent 这个事件点出手,其它事件一律放行。

从这一轮的推理结果里掏出 token 用量;掏不到就直接返回,不添乱。

deferContextual:打开那个随订阅链传下来的"隐形手提箱"(reactor Context)。

箱子里若有上层放的累加器引用,就把本轮用量加进去——主 Agent、子 Agent 全部汇到同一个总数。

⚠️
为什么不能在 Service 里 for 循环加

子 Agent 是异步流式触发的,Service 那层根本拿不到每个子调用的实时用量。只有让累加器沿订阅链随 Context 传播,才能不漏一层地累全。

练一练:你会怎么决策 / 排障?

四道题,全是真实场景。先想清楚再看解析。

① 要加一个"查活动效果"的能力,只想让飞书用户用、网页端不暴露。挂哪个扩展点?

② 用户反馈某定时任务"偶发触发了两次"。你先查什么?

③ 为什么 token 统计要用 Hook + reactor Context,而不是在 Service 里 for 循环把各 Agent 加一遍?

④ 一个独占型任务上一轮跑得很慢、还没结束,下一拍又到点了。互斥锁会怎么处理本轮?

小结:选对卡口,看懂锁

A
加能力先认卡口

BundleTool(主 Agent · 分渠道)/ AgentTool(子 Agent)/ Customizer(编码)/ Hook(旁观),别一股脑塞进消息处理器。

B
两把锁挡两件事

幂等锁防同刻撞车,互斥锁防跨时叠跑——任务"没触发 / 触发两次",从这两关查起。

C
Context 让统计跨层不漏

token 用量靠 Hook + reactor Context 沿订阅链累全,跨主 / 子 Agent 一个不落。

➡️
下一站

扩展点多了也会埋坑。模块 6 看真实排障点 + 把前五节拼成一张全局大图。

06

排障直觉与全局大图

坏了怎么查 —— 错误处理规范、一个教科书级的"两路径不一致"坑,最后用一张大图把六节收口。

前五节把系统拆开了,这节练"坏了怎么查"

排障的第一直觉:报错给谁看?给用户的,是能看懂的体检结论;原始 堆栈 这种"原始化验单",留在后台日志。

🩺
给用户:体检结论(可读)

BizException 的 message、SSE 的 [ERROR] 兜底 —— 一句人话说清"出了什么事"。

📋
进日志:原始病历(堆栈)

异常对象、SQL、类名、调用链 —— 排障时翻日志,永远不甩给终端用户。

🧭
本节两条排障主线

错误处理规范:业务错误 / 数据失败 / 流式异常,各有各的"正确报错姿势"。② 两路径不一致:同一股流,两套实现,改一处忘另一处 —— AI 改代码最容易踩的坑。

教科书级的坑:同一股 Event 流,两套 token 采集

运营说"Web 端 token 统计和接口对不上"。根因藏在这里:HTTP 流式走的 doStream 和 Web 后台流式走的 doStreamWeb 用了两套不一样的采集机制

路径 A · HTTP/通用流式 doStream
// 走 Hook 累加的 chatUsageRef + 分轮 turns
AgentChatMetrics metrics = buildMetrics(
    ctx.getChatUsageRef().get(),
    ctx.getChatUsageTurns(),
    System.currentTimeMillis() - start,
    modelNameRef.get());
token 由 Hook 旁路累加进 chatUsageRef
turns 分轮统计
跨主/子 Agent 累全
路径 B · WEB 后台流式 doStreamWeb
Flux<Event> events = agent.stream(userMsg, WEB_STREAM_OPTIONS)
    .contextWrite(ctx -> writeAgentContext(ctx, context))
    .doOnNext(event -> {
        if (event.getMessage() != null && event.getMessage().getChatUsage() != null) {
            lastUsage.set(event.getMessage().getChatUsage());   // ← 旧机制:直接读 event
        }
    })
    .doOnComplete(() -> {
        AgentChatMetrics metrics = buildMetrics(lastUsage.get(),
                System.currentTimeMillis() - start, modelName);   // 无 turns
        context.getMetricsRef().set(metrics);
    });
token 直接从 eventgetChatUsage()
不走 Hook无 turns
只保留 lastUsage,可能漏算多轮
⚠️
改一处忘另一处 → 指标对不上

同一逻辑有两份实现:你改了 doStream 的采集口径,doStreamWeb 还是老样子,于是 Web 指标和 HTTP 接口对不上。这正是"AI 改完一处、另一处行为变了"的典型。

🔍
排障直觉:先问"有没有第二份实现"

看到 doStream / doStreamWeb 这种成对命名,警觉性拉满 —— 一个症状只在 Web 出现而 HTTP 正常,八成是两套实现分叉了。Review/排障时一眼揪出隐藏的分叉,是高阶直觉。

报错的三种"正确姿势"

不是所有"失败"都一样。三类场景,三种工具,别混用。

🚫

业务错误

校验失败、状态不满足、权限不足 → 抛 BizException(ErrorCode)。message 面向用户。

📦

数据 + 可能失败

既要返回数据、又可能失败 → 返回 StatusResult<T>。调用方 isOk() / unwrap() 取值。

🌊

流式异常

SSE 推流途中炸了 → 推 [ERROR] + userFacingErrorMessage 兜底,堆栈进日志。

🚧
三条红线

禁止裸异常 e.getMessage() 直接透传给用户;禁止在 BizException 里拼堆栈/SQL/类名;禁止用 IllegalArgumentException 表达业务错误(统一用 BizException)。

读代码:业务校验失败,抛 BizException

Cron 任务的归属校验 —— 这就是"业务错误"的标准写法。

CODE
// CronTaskServiceImpl.java:261
private void verifyOwnership(String operatorId, CronTaskPO task) {
    if (task == null) throw new BizException(ErrorCode.TASK_NOT_FOUND);
    if (!operatorId.equals(ownerId)) throw new BizException(ErrorCode.TASK_NOT_FOUND);
}
private void verifyActive(CronTaskPO task) {
    if (task.getTaskStatus() != CronTaskStatus.ACTIVE.getCode()) throw new BizException(ErrorCode.TASK_EXPIRED);
}
大白话

任务不存在?抛业务异常 TASK_NOT_FOUND —— 用户会看到一句可读的"任务不存在"。

不是你的任务?同样报"不存在" —— 既校验权限,又不泄露"任务确实存在但不属于你"。

任务已失效?抛 TASK_EXPIRED。每个错误码都自带一句给用户看的话。

关键:这里不返回数据、纯校验,所以用抛异常而不是 StatusResult。

💡
为什么"不存在"和"不属于你"都报同一个码?

安全考量:如果"不属于你"单独报错,攻击者就能靠错误差异探测出哪些任务 ID 真实存在。统一报"不存在"堵死这条信息泄露。

读代码:流式炸了,怎么优雅收尾

SSE 推到一半 Agent 抛异常 —— 不能让前端干等,也不能把堆栈喷给用户。

CODE
// AgentChatService.java:285
} catch (Exception e) {
    log.error("[AgentChatService] Agent 调用异常, sessionId={}", sessionId, e);
    try { emitter.send(SseEmitter.event().data("[ERROR] " + userFacingErrorMessage(e))); }
    catch (IOException ignored) { }
    emitter.completeWithError(e);
}
大白话

出异常了。第一件事:把完整异常(含堆栈)写进日志,带上 sessionId 方便定位。

给前端推一条 [ERROR] + 用户可读消息 —— 注意是 userFacingErrorMessage(e),不是裸的 e.getMessage()

连推送都失败?吞掉这个 IO 异常(前端可能已断开),别让兜底再炸。

最后正常关闭流,前端拿到 [ERROR] 标记就能展示友好提示而不是卡死。

🎯
排障落点

"用户看到一长串 Java 堆栈"= 有人把 e.getMessage() 直接 send 了。正解永远是:堆栈进日志,userFacingErrorMessage 进前端

收口:一张大图,把六节串成一个心智模型

从"谁来的"到"token 怎么回",点任意节点回顾它在全局里的角色。

① 渠道入口(模块4)

🌐
Web / 飞书 / Cron / A2A
↓ 一句话 + sessionId

② 接入层(模块1·2)

🚪
Controller / 飞书长连接

③ 门面 + 会话(模块2·4)

🎬
AgentChatService(门面)
🗂️
SessionManager(L1+L2)
↓ getAgent miss → 现造

④ 工厂 + 配置驱动(模块3·5)

💄
Factory(组装)
🔌
四大 Registry ← YAML
🧩
扩展点 / BundleTool
↓ build

⑤ 执行引擎(模块1·5)

🤖
ReActAgent + Hook(旁路采集 token)
↑ Flux<Event> 逐个回

⑥ 回流给用户(模块1·6)

📤
SSE token / 飞书卡片 / [ERROR] 兜底
点任意节点,看它在全局里干什么 · 标注:配置驱动贯穿④、Hook 旁路采集贯穿⑤→⑥

排障演练 · 把直觉用起来

四道题,全是真实场景。先想,再看解析。

运营反馈:"Web 端 token 统计和接口对不上",HTTP 接口却正常。你先看哪?

一个 Service 方法做权限校验,校验不通过。该抛什么?

用户在前端看到一长串 Java 堆栈。违反了哪条、怎么改?

把一次飞书提问的链路按真实顺序排好:

通关 · 你现在的排障地图

定位串话 / 丢记忆

顺着 sessionId 怎么拼、L1 是否过期、L2 有没有 saveTo/loadIfExists,一路查到根因。

判断改配置还是改代码

加子 Agent / 换模型 = 改 YAML;加渠道隔离能力 = BundleTool。先分清"配电箱合闸"还是"重新布线"。

看出隐藏的第二份实现

doStream / doStreamWeb 这种成对命名一出现,立刻警觉"改一处忘另一处",这是 AI 改代码最爱埋的坑。

🏁
收尾

六节走完,你手里已经有了一张完整的心智地图:一句话从渠道进来、被门面分流、靠 sessionId 认人、按配置组装出 Agent、在 Hook 旁路采集中执行、最后逐个 token 流回。再遇到"串话""丢记忆""指标对不上""堆栈喷给用户",你不再是从头翻代码,而是直接指向那一层、那一处。能判断改配置还是改代码,能一眼看穿隐藏的第二份实现 —— 这就是排障直觉,也是这门课想留给你的东西。