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