From 290796e5cafba37c0c445b2181c00ef05c3078ad Mon Sep 17 00:00:00 2001 From: Leandro Barbosa Date: Thu, 19 Feb 2026 18:16:23 -0300 Subject: [PATCH] test: add async spawn, announcer, dedup, and process scope tests 27 tests covering Phase 4 components: Spawn (10): accepted, concurrency limit, cascade stop, context timeout, parallel fan-out, announcer deliver/drain, pending, back-pressure, concurrent delivery, cleanup Dedup (10): first call, duplicate, different keys, expired entry, check with result, size, concurrent access (100 goroutines), sweep, spawn key deterministic, announce key format Process scope (7): register/owns, deregister, cross-session isolation, list PIDs filters dead, kill all, cleanup, empty session --- pkg/multiagent/dedup_test.go | 160 ++++++++++++++ pkg/multiagent/spawn_test.go | 378 ++++++++++++++++++++++++++++++++ pkg/tools/process_scope_test.go | 136 ++++++++++++ 3 files changed, 674 insertions(+) create mode 100644 pkg/multiagent/dedup_test.go create mode 100644 pkg/multiagent/spawn_test.go create mode 100644 pkg/tools/process_scope_test.go diff --git a/pkg/multiagent/dedup_test.go b/pkg/multiagent/dedup_test.go new file mode 100644 index 000000000..6969dd992 --- /dev/null +++ b/pkg/multiagent/dedup_test.go @@ -0,0 +1,160 @@ +package multiagent + +import ( + "sync" + "testing" + "time" +) + +func TestDedupCache_FirstCallNotDuplicate(t *testing.T) { + dc := NewDedupCache(5 * time.Minute) + defer dc.Stop() + + if dc.Check("key-1") { + t.Error("first call should not be a duplicate") + } +} + +func TestDedupCache_SecondCallIsDuplicate(t *testing.T) { + dc := NewDedupCache(5 * time.Minute) + defer dc.Stop() + + dc.Check("key-1") + if !dc.Check("key-1") { + t.Error("second call with same key should be a duplicate") + } +} + +func TestDedupCache_DifferentKeysNotDuplicate(t *testing.T) { + dc := NewDedupCache(5 * time.Minute) + defer dc.Stop() + + dc.Check("key-1") + if dc.Check("key-2") { + t.Error("different key should not be a duplicate") + } +} + +func TestDedupCache_ExpiredEntryNotDuplicate(t *testing.T) { + dc := NewDedupCache(50 * time.Millisecond) // very short TTL + defer dc.Stop() + + dc.Check("key-1") + time.Sleep(100 * time.Millisecond) // wait for expiry + + if dc.Check("key-1") { + t.Error("expired entry should not be treated as duplicate") + } +} + +func TestDedupCache_CheckWithResult(t *testing.T) { + dc := NewDedupCache(5 * time.Minute) + defer dc.Stop() + + // First call: not a duplicate + result, isDup := dc.CheckWithResult("key-1") + if isDup || result != "" { + t.Error("first call should not be a duplicate") + } + + // Set result + dc.SetResult("key-1", "cached-result") + + // Second call: duplicate with cached result + result, isDup = dc.CheckWithResult("key-1") + if !isDup { + t.Error("second call should be a duplicate") + } + if result != "cached-result" { + t.Errorf("expected cached-result, got %q", result) + } +} + +func TestDedupCache_Size(t *testing.T) { + dc := NewDedupCache(5 * time.Minute) + defer dc.Stop() + + if dc.Size() != 0 { + t.Error("expected size 0") + } + + dc.Check("key-1") + dc.Check("key-2") + dc.Check("key-3") + + if dc.Size() != 3 { + t.Errorf("expected size 3, got %d", dc.Size()) + } +} + +func TestDedupCache_ConcurrentAccess(t *testing.T) { + dc := NewDedupCache(5 * time.Minute) + defer dc.Stop() + + var wg sync.WaitGroup + duplicates := 0 + var mu sync.Mutex + + // 100 goroutines all trying the same key + for i := 0; i < 100; i++ { + wg.Add(1) + go func() { + defer wg.Done() + if dc.Check("same-key") { + mu.Lock() + duplicates++ + mu.Unlock() + } + }() + } + wg.Wait() + + // Exactly 99 should be duplicates (first one registers) + if duplicates != 99 { + t.Errorf("expected 99 duplicates, got %d", duplicates) + } +} + +func TestDedupCache_Sweep(t *testing.T) { + dc := NewDedupCache(50 * time.Millisecond) + defer dc.Stop() + + dc.Check("key-1") + dc.Check("key-2") + dc.Check("key-3") + + if dc.Size() != 3 { + t.Fatalf("expected 3, got %d", dc.Size()) + } + + // Wait for entries to expire + time.Sleep(100 * time.Millisecond) + + // Manually trigger sweep + dc.sweep() + + if dc.Size() != 0 { + t.Errorf("expected 0 after sweep, got %d", dc.Size()) + } +} + +func TestBuildSpawnKey_Deterministic(t *testing.T) { + k1 := BuildSpawnKey("main", "worker", "do X") + k2 := BuildSpawnKey("main", "worker", "do X") + k3 := BuildSpawnKey("main", "worker", "do Y") + + if k1 != k2 { + t.Error("same inputs should produce same key") + } + if k1 == k3 { + t.Error("different task should produce different key") + } +} + +func TestBuildAnnounceKey_Format(t *testing.T) { + key := BuildAnnounceKey("child-session", "run-123") + expected := "announce:v1:child-session:run-123" + if key != expected { + t.Errorf("expected %q, got %q", expected, key) + } +} diff --git a/pkg/multiagent/spawn_test.go b/pkg/multiagent/spawn_test.go new file mode 100644 index 000000000..7c4007e23 --- /dev/null +++ b/pkg/multiagent/spawn_test.go @@ -0,0 +1,378 @@ +package multiagent + +import ( + "context" + "fmt" + "sync" + "testing" + "time" + + "github.com/sipeed/picoclaw/pkg/providers" + "github.com/sipeed/picoclaw/pkg/tools" +) + +// mockAgentResolver provides test agents. +type mockAgentResolver struct { + agents map[string]*AgentInfo +} + +func (r *mockAgentResolver) GetAgentInfo(id string) *AgentInfo { + return r.agents[id] +} + +func (r *mockAgentResolver) ListAgents() []AgentInfo { + var list []AgentInfo + for _, a := range r.agents { + list = append(list, *a) + } + return list +} + +// mockLLMProvider returns a fixed response after a configurable delay. +type mockLLMProvider struct { + response string + delay time.Duration +} + +func (m *mockLLMProvider) Chat(ctx context.Context, messages []providers.Message, t []providers.ToolDefinition, model string, opts map[string]interface{}) (*providers.LLMResponse, error) { + if m.delay > 0 { + select { + case <-ctx.Done(): + return nil, ctx.Err() + case <-time.After(m.delay): + } + } + return &providers.LLMResponse{Content: m.response}, nil +} + +func (m *mockLLMProvider) GetDefaultModel() string { return "mock" } + +func newTestResolver() *mockAgentResolver { + toolReg := tools.NewToolRegistry() + return &mockAgentResolver{ + agents: map[string]*AgentInfo{ + "worker-a": { + ID: "worker-a", + Name: "Worker A", + Role: "test worker", + Provider: &mockLLMProvider{response: "result from A", delay: 50 * time.Millisecond}, + Tools: toolReg, + MaxIter: 3, + }, + "worker-b": { + ID: "worker-b", + Name: "Worker B", + Role: "test worker", + Provider: &mockLLMProvider{response: "result from B", delay: 50 * time.Millisecond}, + Tools: toolReg, + MaxIter: 3, + }, + }, + } +} + +func TestAsyncSpawn_Accepted(t *testing.T) { + registry := NewRunRegistry() + announcer := NewAnnouncer(10) + sm := NewSpawnManager(registry, announcer, 5, 10*time.Second) + resolver := newTestResolver() + board := NewBlackboard() + + result := sm.AsyncSpawn(context.Background(), resolver, board, SpawnRequest{ + FromAgentID: "main", + ToAgentID: "worker-a", + Task: "do something", + ParentRunKey: "parent-session", + }, "test", "chat1") + + if result.Status != "accepted" { + t.Fatalf("expected accepted, got %s: %s", result.Status, result.Error) + } + if result.RunID == "" { + t.Error("expected non-empty RunID") + } + + // Wait for completion + time.Sleep(200 * time.Millisecond) + + // Check announcement was delivered + anns := announcer.Drain("parent-session") + if len(anns) == 0 { + t.Fatal("expected at least 1 announcement") + } + if anns[0].AgentID != "worker-a" { + t.Errorf("expected agent worker-a, got %s", anns[0].AgentID) + } + if anns[0].Outcome == nil || !anns[0].Outcome.Success { + t.Error("expected successful outcome") + } +} + +func TestAsyncSpawn_ConcurrencyLimit(t *testing.T) { + registry := NewRunRegistry() + announcer := NewAnnouncer(10) + sm := NewSpawnManager(registry, announcer, 2, 10*time.Second) // max 2 concurrent + + slowProvider := &mockLLMProvider{response: "slow result", delay: 500 * time.Millisecond} + toolReg := tools.NewToolRegistry() + resolver := &mockAgentResolver{ + agents: map[string]*AgentInfo{ + "worker": { + ID: "worker", Name: "Worker", Provider: slowProvider, + Tools: toolReg, MaxIter: 3, + }, + }, + } + board := NewBlackboard() + + // Spawn 2 (should succeed — at limit) + r1 := sm.AsyncSpawn(context.Background(), resolver, board, SpawnRequest{ + FromAgentID: "main", ToAgentID: "worker", Task: "task 1", ParentRunKey: "parent", + }, "test", "chat1") + r2 := sm.AsyncSpawn(context.Background(), resolver, board, SpawnRequest{ + FromAgentID: "main", ToAgentID: "worker", Task: "task 2", ParentRunKey: "parent", + }, "test", "chat1") + + if r1.Status != "accepted" || r2.Status != "accepted" { + t.Fatalf("first 2 should be accepted, got r1=%s r2=%s", r1.Status, r2.Status) + } + + // Spawn 3rd (should be rejected — over limit) + r3 := sm.AsyncSpawn(context.Background(), resolver, board, SpawnRequest{ + FromAgentID: "main", ToAgentID: "worker", Task: "task 3", ParentRunKey: "parent", + }, "test", "chat1") + + if r3.Status != "rejected" { + t.Errorf("3rd spawn should be rejected, got %s", r3.Status) + } + + // Wait for first 2 to complete + time.Sleep(700 * time.Millisecond) + + // Now should be able to spawn again + r4 := sm.AsyncSpawn(context.Background(), resolver, board, SpawnRequest{ + FromAgentID: "main", ToAgentID: "worker", Task: "task 4", ParentRunKey: "parent", + }, "test", "chat1") + + if r4.Status != "accepted" { + t.Errorf("4th spawn should be accepted after slots freed, got %s: %s", r4.Status, r4.Error) + } +} + +func TestAsyncSpawn_CascadeStop(t *testing.T) { + registry := NewRunRegistry() + announcer := NewAnnouncer(10) + sm := NewSpawnManager(registry, announcer, 5, 10*time.Second) + + slowProvider := &mockLLMProvider{response: "should not complete", delay: 2 * time.Second} + toolReg := tools.NewToolRegistry() + resolver := &mockAgentResolver{ + agents: map[string]*AgentInfo{ + "worker": { + ID: "worker", Name: "Worker", Provider: slowProvider, + Tools: toolReg, MaxIter: 3, + }, + }, + } + board := NewBlackboard() + + r := sm.AsyncSpawn(context.Background(), resolver, board, SpawnRequest{ + FromAgentID: "main", ToAgentID: "worker", Task: "long task", ParentRunKey: "parent", + }, "test", "chat1") + + if r.Status != "accepted" { + t.Fatalf("expected accepted, got %s", r.Status) + } + + // Give goroutine time to start + time.Sleep(50 * time.Millisecond) + + // Verify it's registered + if registry.ActiveCount() == 0 { + t.Fatal("expected at least 1 active run") + } + + // Cascade stop should cancel it + killed := registry.CascadeStop(r.SessionKey) + if killed == 0 { + t.Error("expected cascade stop to kill at least 1 run") + } + + // Wait for goroutine to clean up + time.Sleep(200 * time.Millisecond) +} + +func TestAsyncSpawn_ContextTimeout(t *testing.T) { + registry := NewRunRegistry() + announcer := NewAnnouncer(10) + sm := NewSpawnManager(registry, announcer, 5, 200*time.Millisecond) // very short timeout + + slowProvider := &mockLLMProvider{response: "too slow", delay: 5 * time.Second} + toolReg := tools.NewToolRegistry() + resolver := &mockAgentResolver{ + agents: map[string]*AgentInfo{ + "worker": { + ID: "worker", Name: "Worker", Provider: slowProvider, + Tools: toolReg, MaxIter: 3, + }, + }, + } + board := NewBlackboard() + + sm.AsyncSpawn(context.Background(), resolver, board, SpawnRequest{ + FromAgentID: "main", ToAgentID: "worker", Task: "slow task", ParentRunKey: "parent", + }, "test", "chat1") + + // Wait for timeout + cleanup + time.Sleep(500 * time.Millisecond) + + // Should have an announcement with failure + anns := announcer.Drain("parent") + if len(anns) == 0 { + t.Fatal("expected announcement after timeout") + } + if anns[0].Outcome.Success { + t.Error("expected failure after timeout") + } +} + +func TestAsyncSpawn_ParallelFanOut(t *testing.T) { + registry := NewRunRegistry() + announcer := NewAnnouncer(20) + sm := NewSpawnManager(registry, announcer, 10, 10*time.Second) + resolver := newTestResolver() + board := NewBlackboard() + + // Fan-out: spawn multiple agents in parallel (Google MapReduce pattern) + var results []*SpawnResult + for i := 0; i < 5; i++ { + target := "worker-a" + if i%2 == 1 { + target = "worker-b" + } + r := sm.AsyncSpawn(context.Background(), resolver, board, SpawnRequest{ + FromAgentID: "main", + ToAgentID: target, + Task: fmt.Sprintf("parallel task %d", i), + ParentRunKey: "parent", + }, "test", "chat1") + results = append(results, r) + } + + // All should be accepted + for i, r := range results { + if r.Status != "accepted" { + t.Errorf("spawn %d should be accepted, got %s", i, r.Status) + } + } + + // Fan-in: wait for all to complete and collect results + time.Sleep(500 * time.Millisecond) + + anns := announcer.Drain("parent") + if len(anns) != 5 { + t.Errorf("expected 5 announcements (fan-in), got %d", len(anns)) + } +} + +// TestAnnouncer tests + +func TestAnnouncer_DeliverAndDrain(t *testing.T) { + a := NewAnnouncer(10) + + a.Deliver("session-1", &Announcement{ + RunID: "run-1", + AgentID: "worker-a", + Content: "result 1", + }) + a.Deliver("session-1", &Announcement{ + RunID: "run-2", + AgentID: "worker-b", + Content: "result 2", + }) + + results := a.Drain("session-1") + if len(results) != 2 { + t.Fatalf("expected 2 announcements, got %d", len(results)) + } + + // Drain again should return empty + results2 := a.Drain("session-1") + if len(results2) != 0 { + t.Errorf("expected 0 after drain, got %d", len(results2)) + } +} + +func TestAnnouncer_Pending(t *testing.T) { + a := NewAnnouncer(10) + + if a.Pending("session-1") != 0 { + t.Error("expected 0 pending for new session") + } + + a.Deliver("session-1", &Announcement{RunID: "r1"}) + a.Deliver("session-1", &Announcement{RunID: "r2"}) + + if a.Pending("session-1") != 2 { + t.Errorf("expected 2 pending, got %d", a.Pending("session-1")) + } +} + +func TestAnnouncer_BackPressure(t *testing.T) { + a := NewAnnouncer(2) // tiny buffer + + // Fill buffer + a.Deliver("session-1", &Announcement{RunID: "r1", Content: "first"}) + a.Deliver("session-1", &Announcement{RunID: "r2", Content: "second"}) + + // Overflow — should drop oldest + a.Deliver("session-1", &Announcement{RunID: "r3", Content: "third"}) + + results := a.Drain("session-1") + if len(results) != 2 { + t.Fatalf("expected 2 after back-pressure, got %d", len(results)) + } + // Most recent should be present + hasThird := false + for _, r := range results { + if r.RunID == "r3" { + hasThird = true + } + } + if !hasThird { + t.Error("expected the newest announcement (r3) to be present after back-pressure") + } +} + +func TestAnnouncer_ConcurrentDelivery(t *testing.T) { + a := NewAnnouncer(100) + + var wg sync.WaitGroup + for i := 0; i < 50; i++ { + wg.Add(1) + go func(n int) { + defer wg.Done() + a.Deliver("session-1", &Announcement{ + RunID: fmt.Sprintf("r-%d", n), + Content: fmt.Sprintf("result %d", n), + }) + }(i) + } + wg.Wait() + + results := a.Drain("session-1") + if len(results) != 50 { + t.Errorf("expected 50 concurrent deliveries, got %d", len(results)) + } +} + +func TestAnnouncer_Cleanup(t *testing.T) { + a := NewAnnouncer(10) + a.Deliver("session-1", &Announcement{RunID: "r1"}) + a.Cleanup("session-1") + + // After cleanup, pending should be 0 (new channel) + if a.Pending("session-1") != 0 { + t.Error("expected 0 pending after cleanup") + } +} diff --git a/pkg/tools/process_scope_test.go b/pkg/tools/process_scope_test.go new file mode 100644 index 000000000..ac74419ba --- /dev/null +++ b/pkg/tools/process_scope_test.go @@ -0,0 +1,136 @@ +package tools + +import ( + "os" + "os/exec" + "testing" +) + +func TestProcessScope_RegisterAndOwns(t *testing.T) { + ps := NewProcessScope() + + ps.Register("session-1", 12345) + ps.Register("session-1", 12346) + ps.Register("session-2", 99999) + + if !ps.Owns("session-1", 12345) { + t.Error("session-1 should own PID 12345") + } + if !ps.Owns("session-1", 12346) { + t.Error("session-1 should own PID 12346") + } + if ps.Owns("session-1", 99999) { + t.Error("session-1 should NOT own PID 99999") + } + if !ps.Owns("session-2", 99999) { + t.Error("session-2 should own PID 99999") + } +} + +func TestProcessScope_Deregister(t *testing.T) { + ps := NewProcessScope() + + ps.Register("session-1", 12345) + ps.Deregister("session-1", 12345) + + if ps.Owns("session-1", 12345) { + t.Error("should not own after deregister") + } +} + +func TestProcessScope_CrossSessionIsolation(t *testing.T) { + ps := NewProcessScope() + + ps.Register("session-a", 100) + ps.Register("session-b", 200) + + if ps.Owns("session-a", 200) { + t.Error("session-a should not see session-b's processes") + } + if ps.Owns("session-b", 100) { + t.Error("session-b should not see session-a's processes") + } +} + +func TestProcessScope_ListPIDs_FiltersDeadProcesses(t *testing.T) { + ps := NewProcessScope() + + // Register current PID (alive) and a fake PID (dead) + ps.Register("session-1", os.Getpid()) + ps.Register("session-1", 999999999) // almost certainly not a real PID + + live := ps.ListPIDs("session-1") + + // Current process should be in the list + found := false + for _, pid := range live { + if pid == os.Getpid() { + found = true + } + if pid == 999999999 { + t.Error("dead PID should have been filtered out") + } + } + if !found { + t.Error("current process PID should be in live list") + } +} + +func TestProcessScope_KillAll(t *testing.T) { + ps := NewProcessScope() + + // Start a real process we can kill + cmd := exec.Command("sleep", "60") + if err := cmd.Start(); err != nil { + t.Skipf("cannot start test process: %v", err) + } + pid := cmd.Process.Pid + + ps.Register("session-1", pid) + + killed := ps.KillAll("session-1") + if killed != 1 { + t.Errorf("expected 1 killed, got %d", killed) + } + + // Reap the child process to prevent zombie (zombie still responds to signal 0). + // cmd.Wait() blocks until the process exits and is reaped by the OS. + err := cmd.Wait() + if err == nil { + t.Error("expected wait to return non-nil error after SIGTERM") + } + + // After reaping, process should no longer be in the process table + if isProcessAlive(pid) { + t.Error("process should have been killed") + } +} + +func TestProcessScope_Cleanup(t *testing.T) { + ps := NewProcessScope() + + ps.Register("session-1", os.Getpid()) + ps.Cleanup("session-1") + + if ps.Owns("session-1", os.Getpid()) { + t.Error("should not own after cleanup") + } +} + +func TestProcessScope_EmptySession(t *testing.T) { + ps := NewProcessScope() + + if ps.Owns("nonexistent", 12345) { + t.Error("nonexistent session should not own anything") + } + + pids := ps.ListPIDs("nonexistent") + if len(pids) != 0 { + t.Error("nonexistent session should have no PIDs") + } + + killed := ps.KillAll("nonexistent") + if killed != 0 { + t.Error("killing nonexistent session should kill 0") + } +}