Files
LLMGuardX Dev d534865b33 0.11.4: 旗舰版第五轮完善(智能体节点任务下发/个人智能体安全策略)
- 节点任务:管理端向节点池下发(Prompt/HTTP/MCP/Skill/数字员工/自定义),
  指定节点或池路由,认领 SKIP LOCKED + 15 分钟租约,认领令牌防重放上报,
  失败 30s×次数退避重入队,达上限 failed,支持取消/重试,完成与失败站内信。
- 个人智能体安全策略:auto_approve_tools 跳过个人调用审批门;
  rate_limit_multiplier 按 (tool,user) 独立窗口放宽个人限流(全局额度不受影响)。
- 修复存量缺陷:/v1/agent/nodes/ 未挂 publicMux,节点心跳/认领端点在部署
  拓扑下不可达。
- 迁移 000046;任务全链路集成测试连真实库通过,HTTP 端到端验证
  (下发→认领→伪造令牌拒绝→上报→succeeded,列表不泄露认领令牌);
  25 包测试通过,前后端构建通过。
2026-08-13 14:30:37 +08:00

118 lines
4.9 KiB
Go

package agentnode
import (
"context"
"encoding/json"
"net"
"os"
"testing"
"time"
"aigateway.local/core/internal/platform/config"
"aigateway.local/core/internal/platform/database"
)
// TestAgentTaskPostgreSQLLifecycle 模拟一个节点跑通任务全链路:
// 下发 → 心跳上线 → 认领(SKIP LOCKED) → 上报成功 → 状态校验;
// 以及失败重试入队与并发认领互斥。
func TestAgentTaskPostgreSQLLifecycle(t *testing.T) {
databaseURL := os.Getenv("AGENT_NODE_TEST_DATABASE_URL")
if databaseURL == "" {
t.Skip("AGENT_NODE_TEST_DATABASE_URL is not set")
}
ctx := context.Background()
pool, err := database.Open(ctx, config.Database{URL: databaseURL, MaxConns: 8, MinConns: 0})
if err != nil {
t.Fatal(err)
}
defer pool.Close()
_, _ = pool.Exec(ctx, `DELETE FROM gateway.agent_nodes WHERE code='task-node-integration'`)
_, _ = pool.Exec(ctx, `DELETE FROM gateway.agent_tasks WHERE payload->>'tag'='integration'`)
defer pool.Exec(ctx, `DELETE FROM gateway.agent_nodes WHERE code='task-node-integration'`)
defer pool.Exec(ctx, `DELETE FROM gateway.agent_tasks WHERE payload->>'tag'='integration'`)
store := NewStore(pool)
// 1. 登记节点并心跳上线。
node, token, err := store.Create(ctx, CreateInput{Code: "task-node-integration", Name: "Task Node", PoolType: "private", PoolCode: "test", Enabled: true}, "")
if err != nil {
t.Fatal(err)
}
if _, err := store.Heartbeat(ctx, node.Code, token, net.ParseIP("192.0.2.11"), HeartbeatInput{Version: "test", Capabilities: map[string]any{"tool_exec": true}}); err != nil {
t.Fatal(err)
}
// 2. 下发两条任务(同池)。
payload, _ := json.Marshal(map[string]any{"tag": "integration", "seq": 1})
first, err := store.CreateTask(ctx, TaskInput{TaskType: "prompt", Payload: payload, PoolType: "private", PoolCode: "test", MaxAttempts: 2}, "")
if err != nil || first.Status != "queued" {
t.Fatalf("create task=%+v err=%v", first, err)
}
payload2, _ := json.Marshal(map[string]any{"tag": "integration", "seq": 2})
second, err := store.CreateTask(ctx, TaskInput{TaskType: "http", Payload: payload2, PoolType: "private", PoolCode: "test", MaxAttempts: 2}, "")
if err != nil {
t.Fatal(err)
}
// 3. 节点认领:两条任务按创建顺序各被认领一次。
claimed1, err := store.ClaimTask(ctx, node.Code, token)
if err != nil || claimed1.ID != first.ID || claimed1.Status != "claimed" {
t.Fatalf("claim1 task=%+v err=%v", claimed1, err)
}
claimed2, err := store.ClaimTask(ctx, node.Code, token)
if err != nil || claimed2.ID != second.ID {
t.Fatalf("claim2 task=%+v err=%v", claimed2, err)
}
if empty, err := store.ClaimTask(ctx, node.Code, token); err != nil || empty.ID != "" {
t.Fatalf("third claim should be empty, got %+v err=%v", empty, err)
}
// 4. 错误 claim_token 上报必须失败(防串扰)。
badResult, _ := json.Marshal(map[string]string{"ok": "true"})
if _, err := store.CompleteTask(ctx, node.Code, token, first.ID, "forged-token", badResult, ""); err == nil {
t.Fatal("complete with forged claim token should fail")
}
// 5. 成功上报第一条。
done, err := store.CompleteTask(ctx, node.Code, token, first.ID, claimed1.ClaimToken, badResult, "")
if err != nil || done.Status != "succeeded" || done.Attempts != 1 {
t.Fatalf("complete task=%+v err=%v", done, err)
}
// 6. 第二条上报失败 → attempts<max 应回队重试。
if _, err := store.CompleteTask(ctx, node.Code, token, second.ID, claimed2.ClaimToken, nil, "node crashed"); err != nil {
t.Fatal(err)
}
retried, err := store.GetTask(ctx, second.ID)
if err != nil || retried.Status != "queued" || retried.Attempts != 1 || retried.AvailableAt.After(time.Now().Add(2*time.Minute)) {
t.Fatalf("retry task=%+v err=%v", retried, err)
}
// 7. 失败任务按退避延迟可认领:拨回 available_at 后重新认领,再次失败
// 达 max_attempts 后标记 failed。
if _, err := pool.Exec(ctx, `UPDATE gateway.agent_tasks SET available_at=clock_timestamp()-interval '1 minute' WHERE id=$1`, second.ID); err != nil {
t.Fatal(err)
}
claimed3, err := store.ClaimTask(ctx, node.Code, token)
if err != nil || claimed3.ID != second.ID {
t.Fatalf("reclaim task=%+v err=%v", claimed3, err)
}
if _, err := store.CompleteTask(ctx, node.Code, token, second.ID, claimed3.ClaimToken, nil, "node crashed again"); err != nil {
t.Fatal(err)
}
failed, err := store.GetTask(ctx, second.ID)
if err != nil || failed.Status != "failed" || failed.Attempts != 2 {
t.Fatalf("failed task=%+v err=%v", failed, err)
}
// 8. 取消已结束任务应冲突;重试失败任务应重新入队。
if err := store.CancelTask(ctx, first.ID); err == nil {
t.Fatal("cancel succeeded task should conflict")
}
if err := store.RetryTask(ctx, second.ID); err != nil {
t.Fatal(err)
}
if retriedAgain, err := store.GetTask(ctx, second.ID); err != nil || retriedAgain.Status != "queued" {
t.Fatalf("retry again task=%+v err=%v", retriedAgain, err)
}
}