package node_test import ( "sync" "sync/atomic" "testing" toki "git.toki-labs.com/toki/proto-socket/go" edgenode "iop/apps/edge/internal/node" "iop/packages/go/config" ) func TestRegistry_RegisterAndCount(t *testing.T) { reg := edgenode.NewRegistry() entry := &edgenode.NodeEntry{ NodeID: "node-001", Alias: "local-node", Client: &toki.TcpClient{}, } reg.Register(entry) if reg.Count() != 1 { t.Fatalf("expected 1 node, got %d", reg.Count()) } } func TestRegistryRegisterIfAbsentRejectsConcurrentDuplicate(t *testing.T) { reg := edgenode.NewRegistry() var accepted int32 var wg sync.WaitGroup start := make(chan struct{}) for i := 0; i < 32; i++ { wg.Add(1) go func() { defer wg.Done() <-start if reg.RegisterIfAbsent(&edgenode.NodeEntry{ NodeID: "node-dup", Alias: "dup", Client: &toki.TcpClient{}, }) { atomic.AddInt32(&accepted, 1) } }() } close(start) wg.Wait() if accepted != 1 { t.Fatalf("accepted registrations: got %d want 1", accepted) } if reg.Count() != 1 { t.Fatalf("registry count: got %d want 1", reg.Count()) } } func TestRegistryUnregisterIfClientIgnoresStaleConnection(t *testing.T) { reg := edgenode.NewRegistry() live := &toki.TcpClient{} stale := &toki.TcpClient{} reg.Register(&edgenode.NodeEntry{NodeID: "node-1", Alias: "alias-1", Client: live}) if _, ok := reg.UnregisterIfClient("node-1", stale); ok { t.Fatal("stale client should not unregister live entry") } entry, ok := reg.Get("node-1") if !ok { t.Fatal("live entry was removed by stale client") } if entry.Client != live { t.Fatal("live entry client changed unexpectedly") } gen, ok := reg.UnregisterIfClient("node-1", live) if !ok { t.Fatal("live client should unregister entry") } if gen != entry.ConnectionGeneration { t.Fatalf("unregister returned generation %d, want the live owner's %d", gen, entry.ConnectionGeneration) } if reg.Count() != 0 { t.Fatalf("registry count after live unregister: got %d want 0", reg.Count()) } } // TestRegistryAssignsMonotonicConnectionGeneration pins that each accepted // registration of a node id draws a strictly higher generation, that a reconnect // after unregister keeps climbing, and that a rejected duplicate never advances // the counter or steals the live owner's generation. func TestRegistryAssignsMonotonicConnectionGeneration(t *testing.T) { reg := edgenode.NewRegistry() first := &edgenode.NodeEntry{NodeID: "node-gen", Alias: "gen", Client: &toki.TcpClient{}} if !reg.RegisterIfAbsent(first) { t.Fatal("first registration should be accepted") } if first.ConnectionGeneration == 0 { t.Fatal("accepted registration must be assigned a non-zero generation") } if gen, ok := reg.CurrentGeneration("node-gen"); !ok || gen != first.ConnectionGeneration { t.Fatalf("CurrentGeneration = %d,%v; want %d", gen, ok, first.ConnectionGeneration) } if !reg.IsCurrentOwnerGeneration("node-gen", first.ConnectionGeneration) { t.Fatal("first owner generation should be current") } // A rejected duplicate must not advance the counter or become current. duplicate := &edgenode.NodeEntry{NodeID: "node-gen", Alias: "gen", Client: &toki.TcpClient{}} if reg.RegisterIfAbsent(duplicate) { t.Fatal("duplicate registration should be rejected") } if gen, _ := reg.CurrentGeneration("node-gen"); gen != first.ConnectionGeneration { t.Fatalf("rejected duplicate changed the live generation to %d", gen) } // Reconnect after unregister climbs strictly. if _, ok := reg.UnregisterIfClient("node-gen", first.Client); !ok { t.Fatal("live client should unregister") } if reg.IsCurrentOwnerGeneration("node-gen", first.ConnectionGeneration) { t.Fatal("no owner should be current after unregister") } second := &edgenode.NodeEntry{NodeID: "node-gen", Alias: "gen", Client: &toki.TcpClient{}} if !reg.RegisterIfAbsent(second) { t.Fatal("reconnect registration should be accepted") } if second.ConnectionGeneration <= first.ConnectionGeneration { t.Fatalf("reconnect generation %d must exceed previous %d", second.ConnectionGeneration, first.ConnectionGeneration) } // The old generation is no longer the current owner; the stale connection's // late callback cannot masquerade as the reconnect. if reg.IsCurrentOwnerGeneration("node-gen", first.ConnectionGeneration) { t.Fatal("old generation must not be current after reconnect") } if !reg.IsCurrentOwnerGeneration("node-gen", second.ConnectionGeneration) { t.Fatal("reconnect generation should be current") } } func TestRegistry_Resolve_ByAliasOrID(t *testing.T) { reg := edgenode.NewRegistry() entry := &edgenode.NodeEntry{ NodeID: "node-001", Alias: "local-node", } reg.Register(entry) // Resolve by ID if e, err := reg.Resolve("node-001"); err != nil || e.NodeID != "node-001" { t.Errorf("failed to resolve by ID: %v", err) } // Resolve by Alias if e, err := reg.Resolve("local-node"); err != nil || e.NodeID != "node-001" { t.Errorf("failed to resolve by Alias: %v", err) } // Resolve non-existent if _, err := reg.Resolve("unknown"); err == nil { t.Error("expected error for unknown node") } } func TestRegistry_Resolve_ByDisplayLabel(t *testing.T) { reg := edgenode.NewRegistry() reg.Register(&edgenode.NodeEntry{NodeID: "node-001", Alias: "local-node"}) reg.Register(&edgenode.NodeEntry{NodeID: "node-002", Alias: "remote-node"}) if e, err := reg.Resolve("node0"); err != nil || e.NodeID != "node-001" { t.Fatalf("Resolve(node0) = %+v, %v; want node-001", e, err) } if e, err := reg.Resolve("node1"); err != nil || e.NodeID != "node-002" { t.Fatalf("Resolve(node1) = %+v, %v; want node-002", e, err) } } func TestRegistry_Resolve_ImplicitSingleNodeOnly(t *testing.T) { reg := edgenode.NewRegistry() // Empty if _, err := reg.Resolve(""); err == nil { t.Error("expected error on empty registry") } // Single node reg.Register(&edgenode.NodeEntry{NodeID: "node-1"}) if e, err := reg.Resolve(""); err != nil || e.NodeID != "node-1" { t.Errorf("failed implicit resolve for single node: %v", err) } // Multiple nodes reg.Register(&edgenode.NodeEntry{NodeID: "node-2"}) if _, err := reg.Resolve(""); err == nil { t.Error("expected error for implicit resolve with multiple nodes") } } func TestRegistryRegisterPreservesAgentKind(t *testing.T) { reg := edgenode.NewRegistry() reg.Register(&edgenode.NodeEntry{ NodeID: "node-explicit", Alias: "explicit", AgentKind: config.AgentKindGenericNode, }) entry, ok := reg.Get("node-explicit") if !ok { t.Fatal("expected explicit entry to be registered") } if entry.AgentKind != config.AgentKindGenericNode { t.Fatalf("agent kind: got %q want %q", entry.AgentKind, config.AgentKindGenericNode) } if entry.LifecycleState != edgenode.LifecycleConnected { t.Fatalf("lifecycle: got %q want %q", entry.LifecycleState, edgenode.LifecycleConnected) } } func TestRegistryRegisterDefaultsKindAndLifecycle(t *testing.T) { reg := edgenode.NewRegistry() entry := &edgenode.NodeEntry{NodeID: "node-001", Alias: "local-node"} reg.Register(entry) if entry.AgentKind != config.AgentKindGenericNode { t.Fatalf("agent kind default: got %q want %q", entry.AgentKind, config.AgentKindGenericNode) } if entry.LifecycleState != edgenode.LifecycleConnected { t.Fatalf("lifecycle default: got %q want %q", entry.LifecycleState, edgenode.LifecycleConnected) } } func TestRegistry_UnregisterRemovesAlias(t *testing.T) { reg := edgenode.NewRegistry() reg.Register(&edgenode.NodeEntry{NodeID: "node-1", Alias: "alias-1"}) reg.Unregister("node-1") if reg.Count() != 0 { t.Fatalf("expected 0 nodes, got %d", reg.Count()) } if _, err := reg.Resolve("alias-1"); err == nil { t.Error("alias should be removed after unregister") } } func TestRegistryUpdateLifecycle(t *testing.T) { reg := edgenode.NewRegistry() if ok := reg.UpdateLifecycle("non-existent", edgenode.LifecycleOnline); ok { t.Error("expected false for updating non-existent node") } entry := &edgenode.NodeEntry{NodeID: "node-1", Alias: "alias-1"} reg.Register(entry) if ok := reg.UpdateLifecycle("node-1", edgenode.LifecycleOnline); !ok { t.Error("expected true for updating existing node") } retrieved, ok := reg.Get("node-1") if !ok { t.Fatal("failed to get node") } if retrieved.LifecycleState != edgenode.LifecycleOnline { t.Errorf("lifecycle: got %q, want %q", retrieved.LifecycleState, edgenode.LifecycleOnline) } } // TestRegistryRegisterIfAbsentPendingUntilReady pins the two-phase transport // handshake: an accepted registration claims the id (so duplicates are rejected) // but stays out of every dispatch-ready lookup until MarkDispatchReadyIfClient // flips it, at which point it becomes visible to AllReady/ResolveReady/GetReady. func TestRegistryRegisterIfAbsentPendingUntilReady(t *testing.T) { reg := edgenode.NewRegistry() client := &toki.TcpClient{} entry := &edgenode.NodeEntry{NodeID: "node-ready", Alias: "ready", Client: client} if !reg.RegisterIfAbsent(entry) { t.Fatal("expected RegisterIfAbsent to accept the first registration") } if entry.DispatchReady { t.Fatal("RegisterIfAbsent entry must be pending, not dispatch-ready") } // Pending: present for ownership (Get), absent from every ready lookup. if _, ok := reg.Get("node-ready"); !ok { t.Fatal("pending entry must be present for duplicate-ownership checks") } if _, ok := reg.GetReady("node-ready"); ok { t.Fatal("pending entry must not be dispatch-ready via GetReady") } if len(reg.AllReady()) != 0 { t.Fatalf("pending entry must not appear in AllReady, got %d", len(reg.AllReady())) } if _, err := reg.ResolveReady("node-ready"); err == nil { t.Fatal("pending entry must not resolve via ResolveReady") } if _, err := reg.ResolveReady(""); err == nil { t.Fatal("single pending node must not satisfy implicit ready resolve") } gen, transitioned, ok := reg.MarkDispatchReadyIfClient("node-ready", client) if !ok || !transitioned { t.Fatalf("first ready must transition the current owner: ok=%v transitioned=%v", ok, transitioned) } if gen != entry.ConnectionGeneration { t.Fatalf("ready generation = %d, want owner's %d", gen, entry.ConnectionGeneration) } // Ready: now visible everywhere. if _, ok := reg.GetReady("node-ready"); !ok { t.Fatal("ready entry must be visible via GetReady") } if len(reg.AllReady()) != 1 { t.Fatalf("ready entry must appear in AllReady, got %d", len(reg.AllReady())) } if got, err := reg.ResolveReady("node-ready"); err != nil || got.NodeID != "node-ready" { t.Fatalf("ready entry must resolve via ResolveReady: got=%v err=%v", got, err) } if got, err := reg.ResolveReady(""); err != nil || got.NodeID != "node-ready" { t.Fatalf("single ready node must satisfy implicit ready resolve: got=%v err=%v", got, err) } } // TestRegistryMarkDispatchReadyIdempotentForCurrentOwner pins that a duplicate // ready for an already-ready owner reports ok=true, transitioned=false so the // transport acks success without repeating pump/event, while a stale client is // rejected outright. func TestRegistryMarkDispatchReadyIdempotentForCurrentOwner(t *testing.T) { reg := edgenode.NewRegistry() owner := &toki.TcpClient{} stale := &toki.TcpClient{} reg.RegisterIfAbsent(&edgenode.NodeEntry{NodeID: "node-1", Client: owner}) if _, transitioned, ok := reg.MarkDispatchReadyIfClient("node-1", owner); !ok || !transitioned { t.Fatalf("first ready: ok=%v transitioned=%v, want true/true", ok, transitioned) } if _, transitioned, ok := reg.MarkDispatchReadyIfClient("node-1", owner); !ok || transitioned { t.Fatalf("duplicate ready: ok=%v transitioned=%v, want true/false", ok, transitioned) } if _, _, ok := reg.MarkDispatchReadyIfClient("node-1", stale); ok { t.Fatal("stale client ready must be rejected (ok=false)") } if _, _, ok := reg.MarkDispatchReadyIfClient("missing", owner); ok { t.Fatal("ready for an unregistered id must be rejected (ok=false)") } } // TestRegistryReadyOnlyLookupsExcludePending pins that with a mix of pending and // ready connections, the ready-only lookups return only the ready one and the // implicit resolve is unambiguous. func TestRegistryReadyOnlyLookupsExcludePending(t *testing.T) { reg := edgenode.NewRegistry() readyClient := &toki.TcpClient{} readyEntry := &edgenode.NodeEntry{NodeID: "node-ready", Alias: "ready-alias", Client: readyClient} reg.RegisterIfAbsent(readyEntry) reg.MarkDispatchReadyIfClient("node-ready", readyClient) reg.RegisterIfAbsent(&edgenode.NodeEntry{NodeID: "node-pending", Alias: "pending-alias", Client: &toki.TcpClient{}}) ready := reg.AllReady() if len(ready) != 1 || ready[0].NodeID != "node-ready" { t.Fatalf("AllReady = %v, want only node-ready", ready) } if got, err := reg.ResolveReady(""); err != nil || got.NodeID != "node-ready" { t.Fatalf("implicit ready resolve should pick the only ready node: got=%v err=%v", got, err) } if _, err := reg.ResolveReady("pending-alias"); err == nil { t.Fatal("resolving a pending node by alias via ResolveReady must fail") } // Non-ready lookups still see both. if len(reg.All()) != 2 { t.Fatalf("All must include pending entries, got %d", len(reg.All())) } } func TestRegistryMarkDispatchReadyOwnerAndWithCurrentOwner(t *testing.T) { reg := edgenode.NewRegistry() client1 := &toki.TcpClient{} client2 := &toki.TcpClient{} entry := &edgenode.NodeEntry{NodeID: "node-1", Alias: "alias-1", Client: client1} if !reg.RegisterIfAbsent(entry) { t.Fatal("expected RegisterIfAbsent to succeed") } // 1. Verify MarkDispatchReadyOwner transitions and returns a clone. snap, transitioned, ok := reg.MarkDispatchReadyOwner("node-1", client1) if !ok || !transitioned || snap == nil { t.Fatalf("expected ok=true, transitioned=true: ok=%v transitioned=%v snap=%v", ok, transitioned, snap) } if snap.NodeID != "node-1" || snap.ConnectionGeneration != entry.ConnectionGeneration || !snap.DispatchReady { t.Errorf("incorrect snapshot fields: %+v", snap) } // 2. Verify WithCurrentOwner runs function for matching snapshot. run := false ok = reg.WithCurrentOwner(snap, func() { run = true }) if !ok || !run { t.Errorf("expected WithCurrentOwner to execute callback: ok=%v run=%v", ok, run) } // 3. Verify WithCurrentOwner is no-op if owner unregistered. reg.Unregister("node-1") run = false ok = reg.WithCurrentOwner(snap, func() { run = true }) if ok || run { t.Errorf("expected WithCurrentOwner to skip callback after unregister: ok=%v run=%v", ok, run) } // 4. Verify WithCurrentOwner is no-op for stale generation after reconnect. entry2 := &edgenode.NodeEntry{NodeID: "node-1", Alias: "alias-1", Client: client2} if !reg.RegisterIfAbsent(entry2) { t.Fatal("expected RegisterIfAbsent for reconnect to succeed") } // Note: entry2 has a higher connection generation now. run = false ok = reg.WithCurrentOwner(snap, func() { run = true }) if ok || run { t.Errorf("expected WithCurrentOwner to skip callback for stale generation: ok=%v run=%v", ok, run) } }