一次请求的旅程
从 HTTP 入口到 ReActAgent,再把答案一段段流回浏览器
先认识这套系统
promo-smart-agent 是一个促销营销智能体:运营在 Web / 飞书里问一句话,它就能做活动数据问答、定时报告。
内核是 AgentScope 的 ReActAgent ——它是四条渠道共同的终点,整套系统都围绕它做编排。
Web、飞书、定时任务(Cron)、A2A 协议,四个入口同一个大脑
主 Agent、子 Agent、模型、工具大多靠 YAML 配置拼出来,而非硬编码
答案不是一次性返回,而是一个个 token 边生成边推给前端
四条渠道殊途同归。这一节我们只跟最直观的 Web 流式请求,把它从入口一路拆到 ReActAgent。其余渠道在模块 4 展开。
你在界面里敲下一句话,然后呢?
你在 Web 界面输入一句话点发送,浏览器里字一个个往外蹦——这一节就把这条流式链路从入口拆到 ReActAgent。
你在值机口(HTTP 端点)交出一句话,它沿传送带经过安检分拣(Service / SessionManager / Factory),最后被装上飞机(ReActAgent)。返回时不是一次性全到,而是一件件行李(SSE token)从传送带口陆续吐出来。
看组件互相「喊话」
这条链上有 5 个角色。点「下一条」,看它们如何把一句话传下去、再把 token 逐个传回来。
主干就这一条
把上面的对话压成一行调用链——记住它,后面所有模块都是在给这条链上的某个角色补细节。
AgentController
AgentChatService
SessionManager.getAgent
ReActAgent.stream()
Flux<Event> → emitter.send
「流式」不是魔法。它靠 reactor 的 Flux<Event> 逐个 Event 触发,每个事件到达就 SseEmitter .send 一次。所以前端才看得到「字一个个出来」,而不是干等到最后一次性全出。
入口长什么样
先看值机口。这是 Web 的 SSE 流式端点——它只做三件事,然后把 SseEmitter 直接交还给浏览器。
@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 只是值机口,不碰业务。
这是「 sessionId 规范 + 转交」的瘦控制器模式:入口只负责接参数、定渠道,真正的取 Agent、调模型、推流都压到 Service。这样换渠道(飞书 / A2A)时业务逻辑能复用。
真正调 ReActAgent 的地方
请求落到 Service 后,关键就这一段:拿 Agent 产出 Flux<Event>,并在全部产出完时把累计指标写回。
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 逐个滑回浏览器——注意是一个一个回,不是一次性。
看懂这条链,「字不出来了」就有了排查顺序:是 Controller 没接住?Service 没拿到 Agent?还是 ReActAgent 没产出 Event?逐段往下排即可。
动手想一想
三道题,都是真实场景。先想再选,错了也有讲解。
用户反馈:点发送后页面一直空白,最后才一次性蹦出全部文字。最可能是哪段坏了?
想给每次请求加「耗时 + token 用量」埋点统计,插在哪层最合适?
从「一句话」到「返回第一个 token」,组件被经过的正确顺序是?
这一节出场的 AgentChatService / AgentSessionManager / DefaultMainAgentFactory / ReActAgent,模块 2 会逐个拉出来单独介绍——它们各自是谁、负责什么。
演员表:核心组件与分层
模块 1 那串 Controller → Service → SessionManager → Factory → ReActAgent,到底谁负责什么?这一节给每个角色发名牌。
把这套系统想成一个剧组
三个 Maven 模块 各演一个角色,职责泾渭分明。看清谁是谁,就知道一段代码该落在哪。
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 里的声明
<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,并打成一个可运行的包。
运行期的四个角色
模块拆好了,跑起来时谁先谁后?这四张名牌串起一次对话的核心链路。
门面 Service:封装流式 / 同步对话逻辑,让 HTTP 层不用关心 Agent 内部怎么跑。
把活跃的 Agent 暂存在 Guava LoadingCache 里,下次同一会话直接复用,不用重新化妆。
Factory:缓存没命中时,它按配置现场组装一个 Agent —— 装工具、挂子 Agent、接 MCP、跑扩展点。
真正“思考—调工具—再思考”的执行者。Factory 装扮好它,SessionManager 让它候场,场记叫它上台。
场记为什么要存在?
门面把“怎么对话”这件事收口,HTTP / Dubbo / 飞书都只调它 —— 改对话逻辑只改一处。
/**
* Agent 对话 Service
* <p>封装流式推送和同步调用逻辑,与 HTTP 层解耦。
* 所有方法均接收 promptTemplateName,透传至 AgentSessionManager 以确保模板切换运行时生效。</p>
*/
@Slf4j @Service @RequiredArgsConstructor
public class AgentChatService {
这就是“场记”—— 一个统一的对话入口。
它把流式推送和一次性同步两种调用都包好,外层 HTTP 不用碰内部细节(这叫“解耦”)。
每个方法都带上要用哪套提示词模板,再原样传给候场区,所以换模板立刻生效、不用重启。
@Service 让 Spring 启动时自动托管它,依赖自动注入。
候场区:命中就复用,没命中就现造
这是 L1 缓存 的心脏:按“多久没访问就过期”来回收,miss 时才叫化妆间造一个。
@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。
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。
挂子 Agent
接 MCP / Skill
跑 Customizer
build 成品
想一想:你会把它放哪?
不背定义,只判断 —— 给你一个场景,凭刚学到的分层直觉做决定。
① 你要新增一个业务枚举和一个工具类,应该放进哪个模块?
② 为什么 SessionManager 要把 Agent 缓存起来,而不是每次对话都新建一个?
③ 飞书链路要用到的 SSE 流式协议类,适合放进 api 模块吗?
④ 你想改“新建 Agent 时多挂一个工具”的逻辑,应该从哪个角色下手?
给 AI 下指令时,你就能说清“把这个加到 core 的 XxxService,别写进 Controller”。下一节:Factory 里的 props.getToolAgentRefs() 和四大 Registry 是怎么来的 —— 答案是全靠 YAML 配置驱动。
配置驱动一切
加一个子 Agent、换一个模型 —— 改的是 YAML,不是 Java。
还记得那行 props.getToolAgentRefs() 吗?
模块2里,Factory 拿子 Agent 时读了一个列表。那个列表里的子 Agent —— 是从哪来的?
答案就一个:一段 YAML。
YAML 就是配电箱面板,每条 models[] / tool-agents[] 是一路开关。要加一路用电器(子 Agent),合一个闸(加一段 YAML)就行 —— 不用重新布线(改代码)。
一段 YAML 的命运
跟着一条 tool-agents[] 配置走完它的一生:从文本,到主 Agent 手里一个能用的工具。
配置的根:一个类,五大列表
整个系统"要装什么",都收敛到 @ConfigurationProperties 标注的这一个类里。
@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[] 一条条变成真的模型。
@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:启动时自动跑一次,做初始化。
把配置里所有模型条目拿出来,逐条处理。
buildModel 按 api-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 的工具箱。
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" —— 你已经知道这是纯配置。下一站:配置之外,运行期怎么按"谁来的"和"哪个会话"分流。
渠道与会话
同一个 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 调用,沿用对方传来的上下文号。
飞书群的房卡号故意只用 chatId、不拼 openId——所以群内每个人共用同一段上下文,A 问完 B 追问能接得上。这正是“串话 / 不串话”的设计开关。
渠道清单:一份枚举钉死五种来源
每个渠道还带一个 group 标志,决定“这次对话算谁的、通知发去哪”。
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:。
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 设了 TTL(默认 30 分钟)到点就清——这是省内存的设计。真正保命的是 L2:只要 saveTo 存过、loadIfExists 能读回,记忆就不会丢。
存档这一步,代码长这样
每轮对话后,把 L1 里那个活跃 Agent 的记忆,写进 L2 档案库。
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 建一条长连接,服务一启动就拨号连上飞书服务器,事件顺着这条管子推过来。
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 后,你预期它会怎样?
小结
五渠道复用同一 Agent 引擎;CallerSourceEnum 定来源,SessionId.of 幂等拼房卡号、认人不串话。
L1 内存快但会过期;L2 Session 持久——saveTo 存、loadIfExists 取,跨过期跨重启不丢。
WebSocket 长连接主动拨出,事件顺管子进主链路。
下一站 → 模块 5:扩展点与聪明工程:渠道隔离的 BundleTool、采集 token 的 Hook、Cron 的锁与重试。
扩展点与聪明工程
想加能力,先问"该挂在哪个卡口"——再看 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() 声明挂载渠道。
@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 次,给网络抖动 / 临时故障一个自愈的机会。
看它们"对话":一个任务怎么通过三道闸
点"下一条",看 Scheduler、Redis、Executor、互斥锁、Handler 之间的对话怎么把任务安全送达。
读第一道闸:幂等锁就是一行"抢号牌"
多台机器同时被定时器叫醒,凭什么只有一台真正执行?看这一行 幂等锁。
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。
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 全部汇到同一个总数。
子 Agent 是异步流式触发的,Service 那层根本拿不到每个子调用的实时用量。只有让累加器沿订阅链随 Context 传播,才能不漏一层地累全。
练一练:你会怎么决策 / 排障?
四道题,全是真实场景。先想清楚再看解析。
① 要加一个"查活动效果"的能力,只想让飞书用户用、网页端不暴露。挂哪个扩展点?
② 用户反馈某定时任务"偶发触发了两次"。你先查什么?
③ 为什么 token 统计要用 Hook + reactor Context,而不是在 Service 里 for 循环把各 Agent 加一遍?
④ 一个独占型任务上一轮跑得很慢、还没结束,下一拍又到点了。互斥锁会怎么处理本轮?
小结:选对卡口,看懂锁
BundleTool(主 Agent · 分渠道)/ AgentTool(子 Agent)/ Customizer(编码)/ Hook(旁观),别一股脑塞进消息处理器。
幂等锁防同刻撞车,互斥锁防跨时叠跑——任务"没触发 / 触发两次",从这两关查起。
token 用量靠 Hook + reactor Context 沿订阅链累全,跨主 / 子 Agent 一个不落。
扩展点多了也会埋坑。模块 6 看真实排障点 + 把前五节拼成一张全局大图。
排障直觉与全局大图
坏了怎么查 —— 错误处理规范、一个教科书级的"两路径不一致"坑,最后用一张大图把六节收口。
前五节把系统拆开了,这节练"坏了怎么查"
排障的第一直觉:报错给谁看?给用户的,是能看懂的体检结论;原始 堆栈 这种"原始化验单",留在后台日志。
BizException 的 message、SSE 的 [ERROR] 兜底 —— 一句人话说清"出了什么事"。
异常对象、SQL、类名、调用链 —— 排障时翻日志,永远不甩给终端用户。
① 错误处理规范:业务错误 / 数据失败 / 流式异常,各有各的"正确报错姿势"。② 两路径不一致:同一股流,两套实现,改一处忘另一处 —— AI 改代码最容易踩的坑。
教科书级的坑:同一股 Event 流,两套 token 采集
运营说"Web 端 token 统计和接口对不上"。根因藏在这里:HTTP 流式走的 doStream 和 Web 后台流式走的 doStreamWeb 用了两套不一样的采集机制。
doStream
// 走 Hook 累加的 chatUsageRef + 分轮 turns
AgentChatMetrics metrics = buildMetrics(
ctx.getChatUsageRef().get(),
ctx.getChatUsageTurns(),
System.currentTimeMillis() - start,
modelNameRef.get());
chatUsageRefdoStreamWeb
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);
});
getChatUsage()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 任务的归属校验 —— 这就是"业务错误"的标准写法。
// 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 抛异常 —— 不能让前端干等,也不能把堆栈喷给用户。
// 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)
② 接入层(模块1·2)
③ 门面 + 会话(模块2·4)
④ 工厂 + 配置驱动(模块3·5)
⑤ 执行引擎(模块1·5)
⑥ 回流给用户(模块1·6)
排障演练 · 把直觉用起来
四道题,全是真实场景。先想,再看解析。
运营反馈:"Web 端 token 统计和接口对不上",HTTP 接口却正常。你先看哪?
一个 Service 方法做权限校验,校验不通过。该抛什么?
用户在前端看到一长串 Java 堆栈。违反了哪条、怎么改?
把一次飞书提问的链路按真实顺序排好:
通关 · 你现在的排障地图
顺着 sessionId 怎么拼、L1 是否过期、L2 有没有 saveTo/loadIfExists,一路查到根因。
加子 Agent / 换模型 = 改 YAML;加渠道隔离能力 = BundleTool。先分清"配电箱合闸"还是"重新布线"。
doStream / doStreamWeb 这种成对命名一出现,立刻警觉"改一处忘另一处",这是 AI 改代码最爱埋的坑。
六节走完,你手里已经有了一张完整的心智地图:一句话从渠道进来、被门面分流、靠 sessionId 认人、按配置组装出 Agent、在 Hook 旁路采集中执行、最后逐个 token 流回。再遇到"串话""丢记忆""指标对不上""堆栈喷给用户",你不再是从头翻代码,而是直接指向那一层、那一处。能判断改配置还是改代码,能一眼看穿隐藏的第二份实现 —— 这就是排障直觉,也是这门课想留给你的东西。