当低代码平台遇上 AI Agent,如何让"自然语言驱动配置"既好用又可靠?本文记录 TOS 数据工厂 Starter 在两周内完成的两阶段演进:先是落地"悬浮球 + Agent 循环内核"的 AI 助手 1.0,随后将后端内核从硬编码的线性循环重构为可观测、可并行、可挂起恢复的 Graph Runtime,并同步优化了 LLM 客户端的 SSE 流式链路。
一、背景:让数据工厂"开口说话"
TOS 数据工厂是一个基于元数据驱动的低代码数据管理 Starter:字典、字典视图、表元数据配置构成了它的核心配置体系。过去这些配置全靠人工在管理界面逐项完成,而 2026 年 8 月的第一周,我们给它加上了一位"配置助手"——用户点击页面右下角的悬浮球,用自然语言即可完成:
- "我需要配置一个字典,字典清单如下……" → AI 创建字典头与字典项
- "基于某个 SQL 字典,帮我生成字典视图" → AI 创建 dict_view 及字段映射
- "我后台建了一张表 xxx,帮我配置元数据管理,识别可配置项并处理好" → AI 读取表结构、生成元数据配置、主动识别枚举列并推荐字典绑定
项目分两个阶段推进:
| 阶段 | 日期 | 交付内容 |
|---|---|---|
| AI 助手 1.0 | 2026-08-01 | 悬浮球前端、线性 Loop Agent 内核、9 个内置工具(SPI 可扩展)、写操作确认令牌、AI 配置管理页 |
| Graph Engine 升级 | 2026-08-07 起 | 内核 Graph 化、ai_graph_run/ai_graph_step 执行可观测、只读工具并行、确认挂起与恢复 |
技术栈:Java 17 / Spring Boot 3.5 / MyBatis-Plus / RestClient(后端),React 19 / TypeScript / Ant Design 6 / react-markdown(前端)。模型接入采用 OpenAI 兼容协议直连(baseUrl/apiKey/model 全部可配),零 SDK 依赖。
二、总体架构
整体链路可以概括为"前端悬浮球 → SSE 流式对话 → Graph Runtime 编排 → LLM + 工具执行 → 三层数据落库":
几条贯穿始终的设计原则:
- 工具复用现有 Service——AI 工具内部直接调用
DictService/MetadataService,天然继承现有安全约束,不另开 SQL 通道; - 对外协议稳定——Graph 升级全程保持
chat/confirmAPI、SSE 事件协议、前端组件零改动; - 渐进式开关——
tos.data.ai.enabled+ai_config.enabled双重控制,关闭后前端悬浮球消失、后端接口 404,业务零影响。
三、核心升级:从 runLoop 到 Graph Runtime
3.1 线性循环的五个痛点
1.0 版本的内核由 AiAgentService.runLoop 驱动:一个 for 循环控制 LLM → tools → LLM,节点概念隐含在代码分支里。功能可用,但工程上存在明显天花板:
- 工具按 LLM 返回顺序串行执行,多个只读查询无法并行;
- 写工具确认是单一挂起点,确认后只能依赖历史消息"重建"继续执行;
- 执行状态散落在内存局部变量中,没有显式 run/step 轨迹;
- 排查问题只能靠日志 + 消息表 + SSE 表现反推执行路径,缺少结构化可观测数据;
- 后续要做复杂工具编排、条件分支、失败定位时,循环体只会越改越乱。
3.2 Graph 化:把隐式循环拆成显式节点
升级引入 com.iauzre.data.starter.ai.graph 子包,把整条执行链路拆为 10 个显式节点,AiAgentService 保留入口职责(锁、限流、归属校验),编排职责全部移交 AiGraphRuntime:
每个节点执行都被记录到 ai_graph_step(node_id、status、耗时、输入/输出摘要),每次 chat/confirm 会话记录到 ai_graph_run(trigger_type、status、current_node、挂起信息)。这就是"可追踪执行流程"的落地——排查一次异常对话,只需查表:
SELECT seq_no, node_id, tool_name, status, duration_ms, error_message
FROM ai_graph_step WHERE run_code = ? ORDER BY seq_no;
两张表都做了敏感信息与体积控制:不存 API Key、不存完整 system prompt、不存完整用户输入和大结果,摘要按固定长度截断(2048 字符)。完整对话仍由 ai_message 负责,工具结果受 tool-result-max-length(默认 4000 字符)约束。
3.3 只读工具安全并行:提速不破坏语义
LLM 单轮返回多个只读工具调用(比如同时 describe_table + list_dicts)是常态。Graph Runtime 将同一批连续只读工具提交到独立线程池并行执行,但遵守一条铁律——并行只改变执行方式,不改变回注顺序:
// AiGraphRuntime#executeReadonlyBatch(节选)
for (ReadonlyCall item : batch) {
AiGraphStep step = recorder.startStep(runCode, sessionCode,
GraphConstants.NODE_READONLY_PARALLEL_TOOLS, "ReadonlyParallelToolsNode",
item.call().getId(), item.tool().name(), item.call().getArguments());
emitJson(state.getSink(), "tool_start", Map.of("toolName", item.tool().name(), "args", item.args()));
Future<ToolExecutionResult> future = toolExecutor.submit(() -> {
ToolExecutionResult result;
try {
AiToolContext ctx = new AiToolContext(state.getUserId(), state.getUsername(),
state.isAdmin(), dataSource);
result = ToolExecutionResult.builder()
.toolCallId(item.call().getId()).toolName(item.tool().name())
.content(truncate(item.tool().execute(item.args(), ctx))) // 统一截断
.success(true).build();
} catch (Exception e) {
result = ToolExecutionResult.builder() // 单工具失败不拖垮同批其他工具
.toolCallId(item.call().getId()).toolName(item.tool().name())
.content(truncate("工具执行失败: " + rootMessage(e)))
.success(false).build();
}
finishReadonlyExecution(state, step, finalized, result); // 独立记录 step
return result;
});
executions.put(item.call().getId(), new ReadonlyExecution(item, future, step, finalized));
}
// 收集阶段:future.get(timeout) 逐个带超时等待,
// 最终 flushReadonlyBatch 仍按 LLM 原始 tool_call 顺序写入 ai_message
稳定性细节远不止"开个线程池":
- 独立超时:每个工具单独应用
tool-timeout,超时future.cancel(true)并生成"工具执行超时,请缩小查询范围"的失败结果回注给 LLM,由模型自行调整方案; - 故障隔离:单个工具失败/超时/不存在只生成失败 tool result,不取消其他只读工具;
- 顺序稳定:SSE 的
tool_start/tool_end可以按实际完成时序推送(前端只展示活动状态,不依赖顺序),但ai_message与 LLM history 的回注顺序必须严格对齐 assistant 的tool_calls数组——否则部分模型服务会直接报协议错误; - 同轮去重缓存:实现层还追加了一个设计文档之外的小增强——同一轮内若 LLM 发出参数完全相同的重复只读调用(参数归一化为排序 JSON 作为 cache key),直接复用首次结果并回放 SSE 事件,避免无谓的重复查询。
3.4 可靠的确认工作流:suspended 状态与断点恢复
写工具的安全模型是本次升级中最值得展开的部分。整个流程分三段:
几道防线层层叠加:
- 权限前置 + 二次校验。
WriteConfirmNode挂起前先检查ROLE_ADMIN,非管理员直接生成拒绝类 tool result 回注(由 AI 向用户解释需要管理员权限);confirm 恢复执行时WriteToolNode再次校验用户、会话归属与权限; - 令牌一次性消费 + 绑定图上下文。
PendingToolCall在升级后扩展了runCode/stepId/nodeId等 graph 恢复字段,确认后 Runtime 能精确定位到挂起的 run 从断点续跑,而不是"靠历史消息猜"; - 挂起即跳过。写工具之后的 tool call 不会抢跑,统一回注"该调用已跳过(等待前序写操作确认)",交由确认后的下一轮 LLM 重新决策;
- 历史自愈。重建 LLM messages 时,
completeToolResponses会为所有未闭合的 tool_call 补上一条"该工具调用未完成(操作未确认或已取消)",避免悬空的 tool_call 破坏 OpenAI 协议约束; - 会话并发锁兜底(详见第五节),防止同一确认令牌被并发消费造成副作用。
3.5 兼容性:前端无感的内核手术
这次重构对外承诺"三个不变":chat/confirm 的 URL 与参数不变、SSE 事件(delta/tool_start/tool_end/confirmation_required/done/error)不变、前端组件不变。验收标准里明确要求:前端无需修改即可完成对话、只读工具调用、写工具确认全链路,且多只读工具并行时总耗时小于严格串行。升级完成后,原有的 AiAgentServiceTest 兼容行为测试全部通过——架构升级的最高境界,就是用户和前端同事都感知不到它发生过。
四、LLM 客户端优化:JdkClientHttpRequestFactory 迁移
OpenAiLlmClient 直连 OpenAI 兼容服务,走 stream: true 的 SSE 流式响应。早期基于默认的 SimpleClientHttpRequestFactory(底层 HttpURLConnection)时,有两类问题直接影响流式链路的稳定性与可排障性:
- 4xx/5xx 时丢失错误响应体——
HttpURLConnection在错误状态下读取 body 的行为受限,模型服务返回的错误详情(比如 "api key invalid"、"insufficient quota")会被吞掉,排障只能靠猜; text/event-stream无法走常规转换器管道——retrieve().body(...)依赖HttpMessageConverter体系,对 SSE 内容类型没有开箱支持。
修复方案分两步。第一步,把请求工厂切换到 JdkClientHttpRequestFactory(底层 java.net.http.HttpClient),错误响应体可以完整读出:
// OpenAiLlmClient(节选)
// 使用 JdkClientHttpRequestFactory:SimpleClientHttpRequestFactory(HttpURLConnection)
// 在 4xx/5xx 时会丢失错误响应体,不利于排障
HttpClient.Builder httpClientBuilder = HttpClient.newBuilder();
if (settings.getConnectTimeout() != null) {
httpClientBuilder.connectTimeout(settings.getConnectTimeout());
}
JdkClientHttpRequestFactory factory = new JdkClientHttpRequestFactory(httpClientBuilder.build());
if (settings.getReadTimeout() != null) {
factory.setReadTimeout(settings.getReadTimeout());
}
RestClient client = RestClient.builder()
.baseUrl(baseUrl)
.requestFactory(factory)
.defaultHeader(HttpHeaders.AUTHORIZATION, "Bearer " + settings.getApiKey())
.build();
第二步,用 exchange(close=false) 直接拿原始 InputStream,绕开消息转换器管道,手工交给 SseStreamParser 逐行解析:
// exchange(close=false) 直接获取原始响应流,绕过 HttpMessageConverter 管道
// (retrieve().body(InputStream.class) 不支持 text/event-stream 内容类型)
try (InputStream is = client.post()
.uri("/chat/completions")
.contentType(MediaType.APPLICATION_JSON)
.accept(MediaType.TEXT_EVENT_STREAM, MediaType.APPLICATION_JSON)
.body(body.toString())
.exchange((request, response) -> {
if (response.getStatusCode().isError()) {
String errBody;
try (InputStream err = response.getBody()) {
errBody = new String(err.readAllBytes(), StandardCharsets.UTF_8);
}
throw new IllegalStateException("LLM调用返回 " + response.getStatusCode() + ": " + errBody);
}
return response.getBody();
}, false)) {
LlmChatResult result = SseStreamParser.parse(
new BufferedReader(new InputStreamReader(is, StandardCharsets.UTF_8)).lines().iterator(),
deltaConsumer, cancellationChecker);
...
}
配套的稳定性处理还包括:cancellationChecker(BooleanSupplier)在解析每个 chunk 时检查 SSE 客户端是否已断开,及时中断读取避免白白消耗模型 token;对含 "closed" 的异常链统一转成 LlmStreamInterruptedException,让 Graph Runtime 能把 run 标记为中断而非失败,与真实异常区分开。SseStreamParser 做成纯函数式的独立类(Iterator<String> 进、LlmChatResult 出),可以脱离网络做单测,finish_reason=length 时追加"[被截断,发'继续'接着说]"提示的逻辑也在单测覆盖内。
另外值得一点:每次调用都按当前 ai_config 临时构建短连接 RestClient,看似"浪费",实则刻意为之——模型参数支持管理界面热更新,复用旧 client 会串用过期配置,正确的权衡是把"每次生效最新配置"放在"连接复用"之前。
五、安全与性能:写操作的纵深防御
AI 助手最容易被质疑的就是"它会不会乱改我的数据"。我们用一组机制把答案钉死为"不会":
5.1 写操作权限与确认双闸门
- 所有
/api/ai/**接口复用现有 Token 鉴权过滤器; - 写工具执行前校验
ROLE_ADMIN(挂起前、恢复后各校验一次),非管理员收到结构化拒绝消息,由 AI 用自然语言解释; - 写工具永远走"预览卡片 → 用户确认 → 令牌校验 → 执行"流程,令牌 TTL 10 分钟、一次性消费、绑定发起人与会话。
5.2 API 密钥掩码处理
ai_config 中的 apiKey 在任何 API 出口都只返回脱敏副本(保留前 3 位 + ****),更新时留空或原样回传脱敏值均视为"不修改":
// AiConfigService(节选)
public AiConfig getMaskedConfig() {
AiConfig config = getConfig();
...
AiConfig masked = copy(config);
masked.setApiKey(mask(config.getApiKey())); // "sk-abc...xyz" -> "sk-****"
return masked;
}
private String mask(String apiKey) {
if (!isNotBlank(apiKey)) return "";
return apiKey.length() <= 3 ? "****" : apiKey.substring(0, 3) + "****";
}
5.3 会话并发锁
同一会话同时只允许一个 chat/confirm 在跑,防止双开 Tab 交叉写入历史。锁基于 CacheService.setIfAbsent(Redis 可用走 Redis,停服自动降级内存),TTL 300 秒——刻意与 SseEmitter 的 300s 超时对齐,保证锁最迟在连接超时的同一时刻自动过期,不会出现"连接已死、锁还赖着";释放时校验持有者 ID(Redis JSON 反序列化可能是 Integer/Long,用 Number 兼容比对),避免误删他人的锁:
// AiAgentService(节选)
private boolean tryLock(String sessionCode, Long userId) {
return cacheService.setIfAbsent(LOCK_PREFIX + sessionCode, userId,
LOCK_TTL_SECONDS, TimeUnit.SECONDS);
}
private void unlock(String sessionCode, Long userId) {
String key = LOCK_PREFIX + sessionCode;
Object holder = cacheService.get(key, Object.class);
// Redis路径经JSON反序列化可能为Integer/Long,用Number兼容比对
if (holder instanceof Number n && n.longValue() == userId) {
cacheService.delete(key);
}
}
5.4 用户级限流与失控防护
限流器按"用户 + 分钟"在 CacheService 上计数(默认每分钟 20 次,可配),超限直接返回 SSE error 事件。整个链路还有三重失控防护:maxAgentRounds 轮数上限(超限结束并明示用户)、历史按完整轮次截断(不拆散 tool_call 与 tool_result 对)、工具结果统一 4000 字符截断(防大结果撞穿上下文窗口):
// AiRateLimiter(节选)
public boolean tryAcquire(Long userId) {
String minute = String.valueOf(Instant.now().truncatedTo(ChronoUnit.MINUTES).getEpochSecond());
String key = KEY_PREFIX + userId + ":" + minute;
Integer count = cacheService.get(key, Integer.class);
int current = count == null ? 0 : count;
if (current >= limitPerMinute) return false;
cacheService.set(key, current + 1, 2, TimeUnit.MINUTES);
return true;
}
此外,只读工具中只有 test_sql_dict 底层真正执行 SQL,安全边界完全复用 DictService.testSqlDict 的既有约束;表结构查询仅走 JDBC DatabaseMetaData,不给 LLM 任何执行任意 SQL 的通道。会话数据则由 AiDataCleanupJob 每日凌晨 3 点按保留天数(默认 90 天)先软删会话、再物理清理消息。
六、前端体验:悬浮球里的完整对话产品
前端是标准的 React + Ant Design 实现,但几个交互决策值得记录:
悬浮球入口。AiAssistant 挂载在 Layout 最外层,不随 Tab 页签切换销毁;支持拖拽(位移超过 4px 判定为拖拽而非点击),位置通过 localStorage 持久化;启动时请求 GET /api/ai/status,功能未启用时整个组件不渲染——前端对"AI 功能不存在"的世界一无所知。
// AiAssistant.tsx(节选):拖拽与点击的判定
const up = () => {
const d = dragRef.current;
if (!d) return;
dragRef.current = null;
if (d.moved) {
// 拖拽:持久化最终位置
localStorage.setItem(POS_KEY, JSON.stringify({ x: d.lastX, y: d.lastY }));
} else {
setOpen(true); // 点击:打开抽屉
}
};
抽屉式对话框。480px 宽的 Drawer 内整合了会话列表(切换/删除)、消息流、确认卡片与输入框。消息流按 ChatItem 联合类型渲染:assistant 消息经 react-markdown 渲染,工具执行渲染为折叠状态条,写操作弹结构化确认卡片(参数预览表格,大列表自动折叠只显示前 N 条);历史消息按 seq_no 向前翻页(滚动到顶部触发加载)。
SSE 实时通信。因为 chat 是 POST 且需携带 Token 头,EventSource 不可用,前端用 fetch + ReadableStream 手写了一个 60 行的 SSE 解析器,按 event:/data: 行协议分发增量渲染;流中断时展示重试按钮而非自动重连,把重试决策留给用户。
确认后的无缝续流。用户点击确认卡片后,confirm 接口返回一条新的 SSE 流,前端在原聊天窗口继续追加渲染——对用户而言,整个"AI 请示 → 我批准 → AI 干完"的过程就是一段连续的对话。
七、实施路径回顾
整个项目按 15 个 Task 推进,每个 Task 独立提交、独立验证(后端 mvn -q test -Dtest=类名,前端 npm run build):
| 阶段 | Task | 关键节点 |
|---|---|---|
| 地基 | 1–2 | 三表 DDL(MySQL/H2 双方言)、实体 Mapper、TosDataAiAutoConfiguration 条件装配骨架 |
| LLM 链路 | 3–4 | 协议模型与 SseStreamParser(纯函数可单测)→ OpenAiLlmClient(此处完成 JdkClientHttpRequestFactory 迁移与 exchange 流式读取) |
| 工具体系 | 5–7 | AiTool SPI + 注册表 → 5 个只读工具 → 4 个写工具 |
| 内核与管控 | 8–11 | 限流器/管理员判定/确认令牌存储 → 会话与配置服务 → Agent 循环内核(TDD,Mock LlmClient 覆盖只读循环/写确认挂起/轮数上限) → 清理任务、Controller、自动装配收口 |
| 前端 | 12–14 | 类型与 SSE 解析 → 悬浮球 + 抽屉对话组件 → AI 配置管理页(apiKey 脱敏回显、连通性测试) |
| 收口 | 15 | 全量验证 + 端到端验收:自然语言建字典 → 生成字典视图 → 配置元数据并识别字典绑定 |
| 升级 | Graph | 内核 Graph 化(ai_graph_run/ai_graph_step、并行工具、断点恢复),12 项专项测试含"并行耗时 < 串行耗时且消息顺序稳定" |
测试策略上有两个值得复制的实践:SseStreamParser 与 Agent 内核全部通过 Mock LlmClient 做确定性单测(不依赖真实模型服务);专项测试覆盖了并发双 Tab 发消息、工具慢查询超时、确认令牌过期重试、Redis 停服降级内存缓存等真实故障场景。
八、总结与展望
回看这次升级,三条主线贯穿始终:
- 稳定性:写操作双闸门(权限 + 确认令牌)、会话并发锁与 Emitter 超时对齐、工具独立超时与故障隔离、结果截断与轮数上限、历史自愈补齐——每一层都有独立的失效模式与兜底;
- 可扩展性:工具开放 SPI(
AiToolBean 自动注册)、模型参数管理界面热更新、Graph 节点显式化后,条件分支节点、多 Agent Graph、可视化拓扑都有了挂载点; - 可观测性:
ai_graph_run/ai_graph_step让每一次 AI 行为都有结构化轨迹可查,AI 不再是黑盒。
后续版本的方向已经在路线图上:Graph 调试 API 与 run 可视化拓扑、条件分支节点与可配置 Graph DSL、服务重启后的幂等恢复、工具副作用防重放,以及按业务场景定义多个 Agent Graph。从"一个能对话的悬浮球"到"一台可审计的执行引擎",这次升级告诉我们:AI Agent 要在生产环境站稳脚跟,模型能力只是一半,另一半是把传统后端工程里的状态管理、并发控制、权限模型和可观测性,原封不动地搬进 Agent 的世界。
从线性循环到图引擎:数据工厂 AI 助手能力升级实录
https://iauzre.com.cn/archives/f0BCkh9r
Comments