Spring Boot实现Webhook Agent任务路由
前面的事件驱动 Agent 讲的是架构。
真正落地时,第一版最容易写成这样:
@PostMapping("/webhook")
public void handle(
@RequestBody String payload) {
agentService.run(payload);
}
这个实现 Demo 能跑,生产几乎一定会出问题。
因为它没有处理:
签名
去重
乱序
Debounce
重放
权限
副作用
下面直接用 Spring Boot 搭一个最小 Webhook Agent Router,把外部事件先变成可靠任务,再交给 Agent。
先定义统一Event Envelope
不同 Provider 的 Webhook 格式差异很大。
不要让 Agent Router 直接理解 Gmail、Slack、GitHub 原始 JSON。
先统一:
public record ExternalEvent(
String provider,
String deliveryId,
String eventType,
String resourceId,
String actorId,
Instant eventTime,
JsonNode metadata) {
}
例如 GitHub:
{
"provider": "github",
"deliveryId": "d-9182",
"eventType": "pull_request.synchronize",
"resourceId": "repo-a/pr-182",
"actorId": "user-17",
"eventTime": "2026-08-26T01:10:00Z"
}
Webhook入口第一件事不是解析业务,而是验签
public interface WebhookSignatureVerifier {
boolean supports(
String provider);
void verify(
HttpHeaders headers,
byte[] rawBody);
}
为什么必须使用 rawBody?
因为很多签名是对:
原始字节
计算。
如果先反序列化再序列化:
空格
字段顺序
Unicode
都可能变化。
Controller先保留原始Payload
@RestController
@RequestMapping("/webhooks")
public class WebhookController {
private final WebhookService service;
@PostMapping("/{provider}")
public ResponseEntity receive(
@PathVariable String provider,
@RequestHeader HttpHeaders headers,
@RequestBody byte[] rawBody) {
service.accept(
provider,
headers,
rawBody);
return ResponseEntity.accepted()
.build();
}
}
第二步:Inbox去重
Entity:
@Entity
@Table(
name = "webhook_inbox",
uniqueConstraints = {
@UniqueConstraint(
columnNames = {
"provider",
"deliveryId"
}
)
}
)
public class WebhookInboxEntity {
@Id
private String id;
private String provider;
private String deliveryId;
private String eventType;
private String resourceId;
private String payloadHash;
@Enumerated(EnumType.STRING)
private InboxStatus status;
private Instant receivedAt;
}
状态:
public enum InboxStatus {
RECEIVED,
NORMALIZED,
DISPATCHED,
COMPLETED,
REJECTED,
FAILED
}
插入失败就是重复
@Transactional
public boolean saveIfNew(
ExternalEvent event,
String payloadHash) {
try {
repository.save(
mapper.toEntity(
event,
payloadHash));
repository.flush();
return true;
} catch (
DataIntegrityViolationException ex) {
return false;
}
}
重复事件:
直接ACK
不再启动Agent
第三步:Normalizer
public interface WebhookNormalizer {
boolean supports(
String provider);
ExternalEvent normalize(
HttpHeaders headers,
byte[] rawBody);
}
GitHub、Slack、Gmail 分别实现。
Router 不关心原始格式。
第四步:Trigger Policy
不是所有事件都值得启动 Agent。
public record TriggerDecision(
boolean accepted,
String taskType,
String reason) {
}
public interface TriggerPolicy {
TriggerDecision evaluate(
ExternalEvent event);
}
例如 PR:
if (!Set.of(
"pull_request.opened",
"pull_request.synchronize")
.contains(event.eventType())) {
return reject("EVENT_NOT_SUPPORTED");
}
Actor也要检查
如果事件来自自己的自动化:
agent-bot
默认忽略。
防止:
Agent评论
→触发Webhook
→再评论
→无限循环
第五步:Debounce
Slack 类事件不适合一条消息一个 Run。
定义:
public record DebounceKey(
String provider,
String resourceId,
String taskType) {
}
数据库:
create table event_debounce (
debounce_key varchar(256) primary key,
first_event_at timestamptz not null,
last_event_at timestamptz not null,
event_count int not null,
dispatch_after timestamptz not null
);
收到新事件:
更新last_event_at
event_count +1
dispatch_after = now + 60s
Scheduler 只发:
到期聚合任务
GitHub PR则适合Latest Wins
PR 连续 Push:
sha1
sha2
sha3
只保留最新:
head_sha
Task 创建前再次查询当前 PR。
如果旧 Run 已经开始:
mark STALE
第六步:Task Dispatcher
不要 Controller 直接调用 Agent。
先生成 Task:
public record AgentTask(
String taskId,
String taskType,
String provider,
String resourceId,
String subjectId,
String triggerEventId,
Map context) {
}
然后进入:
Queue
例如:
public interface AgentTaskQueue {
void publish(
AgentTask task);
}
这样 Webhook Endpoint 可以快速 ACK。
为什么不能在HTTP线程里跑Agent
Agent 可能:
30秒
5分钟
20分钟
Webhook Provider 通常不会等这么久。
如果 Endpoint 超时:
Provider重试
你又得到重复任务。
所以:
Receive
→持久化
→ACK
→异步执行
是基本边界。
第七步:Task创建也要幂等
虽然 Inbox 已经去重,仍然可能出现:
Worker重放
数据库恢复
人工Replay
所以 Task Key:
taskType
+
resourceId
+
resourceVersion
例如:
code-review:repo-a/pr-182:sha-991
唯一约束。
Task表
create table agent_task (
task_id varchar(128) primary key,
idempotency_key varchar(256) unique not null,
task_type varchar(64) not null,
resource_id varchar(256) not null,
status varchar(32) not null,
created_at timestamptz not null,
started_at timestamptz,
completed_at timestamptz
);
第八步:Action Policy不要和Trigger Policy混
Webhook 能启动:
code-review
不代表 Agent 能:
merge
Agent Run 只拿:
repo.read
checks.read
comment.write
如果要 Merge:
新的Approval
+
新的Execution Grant
一个Task Permission Profile
task_type: code-review
capabilities:
- repo.read
- checks.read
- review.comment.write
forbidden:
- repo.merge
- deploy
Task Dispatcher 把这个 Profile 绑定到 Run。
第九步:Side Effect Ledger
如果 Agent 最后写了 Review Comment:
comment.create
要记录。
public record SideEffectEntry(
String idempotencyKey,
String runId,
String capability,
String actionHash,
SideEffectStatus status,
String providerReceipt) {
}
如果 Worker Crash 后 Replay:
先查Ledger
避免重复评论。
第十步:Dead Letter和Replay
任务失败不能丢。
public enum TaskStatus {
PENDING,
RUNNING,
COMPLETED,
FAILED,
DEAD_LETTER,
STALE
}
达到:
max_attempts
进入 DLQ。
Replay:
POST /api/tasks/{taskId}/replay
但 Replay 不是简单把状态改回 PENDING。
必须:
检查资源当前状态
检查Side Effect
检查权限
Replay Policy
public ReplayDecision canReplay(
AgentTask task) {
if (task.isStale()) {
return DENY;
}
if (sideEffects.hasUnknown(task.id())) {
return REQUIRE_RECONCILIATION;
}
return ALLOW;
}
第十一步:Event Audit
每个 Task 要能回到原始 Event:
Task
→ Trigger Event
→ Provider Delivery
保存:
delivery_id
payload_hash
不一定长期保存完整敏感 Payload。
原始正文可以放加密对象存储,数据库只放 Ref。
第十二步:Budget
public record AutomationBudget(
String automationId,
int maxRunsPerHour,
int maxRunsPerDay,
BigDecimal maxCostPerDay) {
}
Dispatch 前检查。
超过:
RATE_LIMITED
而不是无限入队。
一个Budget Gate
if (usage.todayRuns()
>= budget.maxRunsPerDay()) {
return TriggerDecision.reject(
"DAILY_RUN_BUDGET_EXCEEDED");
}
第十三步:同资源串行,不同资源并行
Partition Key:
provider + ":" + resourceId
PR-182 的事件串行。
PR-183 可以并行。
实现可以用:
Kafka Partition
Redis Lock
DB Advisory Lock
不要全局单线程。
第十四步:Metrics
webhook_received_total
webhook_duplicate_total
webhook_signature_failed_total
webhook_rejected_total
agent_task_dispatched_total
agent_task_stale_total
agent_task_replay_total
automation_loop_blocked_total
再看:
Event → Run latency
这是事件驱动自动化最关键的体验指标。
Trace
webhook.receive
├─signature.verify
├─inbox.insert
├─normalize
├─policy.evaluate
└─task.dispatch
└─agent.run
整个链可以从 Delivery ID 查。
最少测试这12个Case
1. 正常事件
2. 同Delivery重复
3. 签名错误
4. 不支持事件类型
5. Automation Actor触发
6. Slack 50条Debounce
7. PR连续3次Push
8. Resource已关闭
9. Task重复创建
10. Agent副作用后Worker Crash
11. UNKNOWN副作用Replay
12. Daily Budget超限
目录结构
webhook/
├── controller
├── signature
├── normalizer
├── inbox
├── policy
├── debounce
├── dispatcher
├── task
├── replay
└── metrics
不要把全部逻辑塞进一个 WebhookController。
Webhook Agent Router 真正做的不是“收个 HTTP 请求”。
它负责把:
不可靠外部事件
转换成:
可去重
可授权
可重放
可预算
可审计
的 Agent Task。
只要这个层次建起来,Gmail、Slack、GitHub 甚至未来更多事件源,都可以复用同一套运行模型。
而且 Agent 仍然只负责它擅长的部分:
理解
推理
生成
决策建议
重复、乱序、幂等、预算和权限,继续交给确定性系统。
更多企业级 AI 应用、Agent、RAG 与模型工程化内容,我会继续整理在 智元界:
https://www.zyentor.com/