package node import ( "fmt" "sort" "strconv" "strings" "sync" toki "git.toki-labs.com/toki/proto-socket/go" "iop/packages/go/config" ) const ( LifecycleAccepted = "accepted" LifecycleOnline = "online" LifecycleFailed = "failed" LifecycleConnected = "connected" ) // NodeEntry represents one connected node. ConnectionGeneration is the internal // monotonic identity the registry assigns to each accepted connection of a node. // It is not a wire/config field; provider leases and dispatch handoffs bind to it // so a stale (superseded or rejected) connection can be fenced by identity, not // only by client pointer. // // DispatchReady separates an accepted connection that has only claimed ownership // and received its config (pending) from one the node has confirmed is ready to // receive dispatch (config applied, handler installed). A pending entry occupies // the id so duplicate registrations are rejected, but it is excluded from run // dispatch, config-refresh push, and connected snapshots/events until the node's // NodeReadyRequest flips this to true. The flag is only ever mutated under the // registry lock; read it through the registry's ready-aware helpers, never off a // shared *NodeEntry, so the transition stays race-free. type NodeEntry struct { NodeID string Alias string AgentKind string LifecycleState string Client *toki.TcpClient Index int HasIndex bool ConnectionGeneration uint64 DispatchReady bool } // Registry manages all nodes connected to edge. type Registry struct { mu sync.RWMutex byID map[string]*NodeEntry byAlias map[string]*NodeEntry byIndex map[int]*NodeEntry indexByID map[string]int nextIndex int // genByID holds the per-node connection generation counter. It is never // reset — a reconnecting node always draws a strictly higher generation, so // an old connection's late callback can never be mistaken for the new owner. genByID map[string]uint64 } func NewRegistry() *Registry { return &Registry{ byID: make(map[string]*NodeEntry), byAlias: make(map[string]*NodeEntry), byIndex: make(map[int]*NodeEntry), indexByID: make(map[string]int), genByID: make(map[string]uint64), } } // Register inserts entry unconditionally and marks it dispatch-ready. It models a // fully-connected, dispatchable node in one step; the transport's two-phase // accepted→ready handshake uses RegisterIfAbsent (pending) plus // MarkDispatchReadyIfClient instead. func (r *Registry) Register(entry *NodeEntry) { r.mu.Lock() defer r.mu.Unlock() entry.DispatchReady = true r.registerLocked(entry) } // RegisterIfAbsent registers entry only when the node id is not already // connected. The check and insert happen under one lock so concurrent duplicate // registration attempts cannot both be accepted by the transport server. The // entry is left pending (DispatchReady=false): it claims the id so duplicates are // rejected, but it is excluded from dispatch/refresh/connected snapshots until // MarkDispatchReadyIfClient flips it ready on the node's NodeReadyRequest. func (r *Registry) RegisterIfAbsent(entry *NodeEntry) bool { r.mu.Lock() defer r.mu.Unlock() if _, exists := r.byID[entry.NodeID]; exists { return false } r.registerLocked(entry) return true } func (r *Registry) registerLocked(entry *NodeEntry) { if entry.AgentKind == "" { entry.AgentKind = config.AgentKindGenericNode } if entry.LifecycleState == "" { entry.LifecycleState = LifecycleConnected } // A connection generation is minted only for an accepted registration, so a // rejected duplicate (RegisterIfAbsent returns before reaching here) never // bumps the live node's generation. r.genByID[entry.NodeID]++ entry.ConnectionGeneration = r.genByID[entry.NodeID] if idx, ok := r.indexByID[entry.NodeID]; ok { entry.Index = idx entry.HasIndex = true } else { if !entry.HasIndex { entry.Index = r.nextIndex } entry.HasIndex = true r.indexByID[entry.NodeID] = entry.Index if entry.Index >= r.nextIndex { r.nextIndex = entry.Index + 1 } } r.byID[entry.NodeID] = entry r.byIndex[entry.Index] = entry if entry.Alias != "" { r.byAlias[entry.Alias] = entry } } func (r *Registry) UpdateLifecycle(nodeID string, state string) bool { r.mu.Lock() defer r.mu.Unlock() entry, ok := r.byID[nodeID] if !ok { return false } entry.LifecycleState = state return true } func (e *NodeEntry) DisplayLabel() string { if e == nil { return "unknown" } if e.HasIndex { return fmt.Sprintf("node%d", e.Index) } if e.Alias != "" { return e.Alias } if e.NodeID != "" { return e.NodeID } return "unknown" } func (r *Registry) Unregister(nodeID string) { r.mu.Lock() defer r.mu.Unlock() if entry, ok := r.byID[nodeID]; ok { r.unregisterLocked(nodeID, entry) } } // UnregisterIfClient removes nodeID only when the currently registered entry // belongs to client, returning the removed owner's connection generation. Late // disconnect callbacks from rejected or superseded connections must not clear the // live registry entry, and the returned generation lets the transport fence the // exact connection that closed rather than whatever owns the id now. func (r *Registry) UnregisterIfClient(nodeID string, client *toki.TcpClient) (uint64, bool) { r.mu.Lock() defer r.mu.Unlock() entry, ok := r.byID[nodeID] if !ok || entry.Client != client { return 0, false } generation := entry.ConnectionGeneration r.unregisterLocked(nodeID, entry) return generation, true } // CurrentGeneration returns the connection generation of the node id's current // registry owner, or false when no entry is registered. func (r *Registry) CurrentGeneration(nodeID string) (uint64, bool) { r.mu.RLock() defer r.mu.RUnlock() entry, ok := r.byID[nodeID] if !ok { return 0, false } return entry.ConnectionGeneration, true } // IsCurrentOwnerGeneration reports whether generation still matches the node id's // current registry owner. A dispatch path calls it just before sending so a lease // minted for a connection that has since disconnected or been superseded by a // reconnect is fenced instead of dispatched to a dead client. func (r *Registry) IsCurrentOwnerGeneration(nodeID string, generation uint64) bool { r.mu.RLock() defer r.mu.RUnlock() entry, ok := r.byID[nodeID] return ok && entry.ConnectionGeneration == generation } // Clone returns a shallow copy of NodeEntry, preserving pointers but copying fields // for safe, lock-free snapshot reads after registry lock release. func (e *NodeEntry) Clone() *NodeEntry { if e == nil { return nil } return &NodeEntry{ NodeID: e.NodeID, Alias: e.Alias, AgentKind: e.AgentKind, LifecycleState: e.LifecycleState, Client: e.Client, Index: e.Index, HasIndex: e.HasIndex, ConnectionGeneration: e.ConnectionGeneration, DispatchReady: e.DispatchReady, } } // MarkDispatchReadyIfClient transitions the node id's current owner to // dispatch-ready, but only when the still-registered entry belongs to client. It // returns the owner's connection generation, whether this call performed the // pending→ready transition, and whether client is the current owner at all. // // - ok=false: client no longer owns the entry (superseded by a reconnect, or the // entry was already removed by a disconnect). The ready signal is stale and // must be rejected — the transport tells the node to reconnect. // - ok=true, transitioned=true: the first ready for this connection. The caller // opens dispatch eligibility, pumps stranded waiters, and emits the connected // event exactly once. // - ok=true, transitioned=false: a duplicate ready for an already-ready owner. // The caller acks success but runs no additional pump or event. func (r *Registry) MarkDispatchReadyIfClient(nodeID string, client *toki.TcpClient) (generation uint64, transitioned bool, ok bool) { r.mu.Lock() defer r.mu.Unlock() entry, exists := r.byID[nodeID] if !exists || entry.Client != client { return 0, false, false } if entry.DispatchReady { return entry.ConnectionGeneration, false, true } entry.DispatchReady = true return entry.ConnectionGeneration, true, true } // MarkDispatchReadyOwner transitions the node id's current owner to dispatch-ready // under the registry lock, returning a cloned snapshot of the NodeEntry, transitioned flag, and ok. // Use this snapshot to drive downstream lifecycle callbacks and ready notifications safely. func (r *Registry) MarkDispatchReadyOwner(nodeID string, client *toki.TcpClient) (entry *NodeEntry, transitioned bool, ok bool) { r.mu.Lock() defer r.mu.Unlock() entry, exists := r.byID[nodeID] if !exists || entry.Client != client { return nil, false, false } if entry.DispatchReady { return entry.Clone(), false, true } entry.DispatchReady = true return entry.Clone(), true, true } // WithCurrentOwner executes fn under the registry lock only when the currently // registered owner for entry.NodeID exactly matches the provided entry's Client // and ConnectionGeneration. This guarantees the callback runs against the current // active owner, ensuring no stale ready events are emitted. func (r *Registry) WithCurrentOwner(entry *NodeEntry, fn func()) bool { if entry == nil { return false } r.mu.Lock() defer r.mu.Unlock() current, exists := r.byID[entry.NodeID] if !exists || current.Client != entry.Client || current.ConnectionGeneration != entry.ConnectionGeneration { return false } fn() return true } // WithCurrentDispatchOwner executes fn under the registry lock only when the currently // registered owner for nodeID matches client and generation. It prevents check-then-act // races between connection generation checks and dispatch enqueuing/handoff. func (r *Registry) WithCurrentDispatchOwner(nodeID string, client *toki.TcpClient, generation uint64, fn func() error) error { r.mu.Lock() defer r.mu.Unlock() current, exists := r.byID[nodeID] if !exists || current.Client != client || current.ConnectionGeneration != generation { return fmt.Errorf("provider node %q connection changed before dispatch (fenced generation %d)", nodeID, generation) } return fn() } func (r *Registry) unregisterLocked(nodeID string, entry *NodeEntry) { if entry.Alias != "" { delete(r.byAlias, entry.Alias) } delete(r.byIndex, entry.Index) delete(r.byID, nodeID) } func (r *Registry) Get(nodeID string) (*NodeEntry, bool) { r.mu.RLock() defer r.mu.RUnlock() e, ok := r.byID[nodeID] return e, ok } // GetReady returns the entry for nodeID only when it exists and is // dispatch-ready. Connected snapshots and config-refresh push use it so a pending // (accepted-but-not-yet-ready) connection is reported offline and never pushed to. func (r *Registry) GetReady(nodeID string) (*NodeEntry, bool) { r.mu.RLock() defer r.mu.RUnlock() e, ok := r.byID[nodeID] if !ok || !e.DispatchReady { return nil, false } return e, true } func (r *Registry) Resolve(ref string) (*NodeEntry, error) { r.mu.RLock() defer r.mu.RUnlock() return r.resolveLocked(ref, false) } // ResolveReady resolves ref exactly like Resolve but only ever returns a // dispatch-ready entry. Direct-dispatch paths (run/cancel/tunnel/node-command) // use it so a request is never sent to a pending connection whose node has not // installed its handler yet; observational lookups keep using Resolve. func (r *Registry) ResolveReady(ref string) (*NodeEntry, error) { r.mu.RLock() defer r.mu.RUnlock() return r.resolveLocked(ref, true) } func (r *Registry) resolveLocked(ref string, readyOnly bool) (*NodeEntry, error) { eligible := func(entry *NodeEntry) bool { return entry != nil && (!readyOnly || entry.DispatchReady) } if ref != "" { if entry, ok := r.byID[ref]; ok && eligible(entry) { return entry, nil } if entry, ok := r.byAlias[ref]; ok && eligible(entry) { return entry, nil } if entry, ok := r.resolveDisplayLabelLocked(ref); ok && eligible(entry) { return entry, nil } return nil, fmt.Errorf("node %q not found", ref) } var only *NodeEntry count := 0 for _, entry := range r.byID { if !eligible(entry) { continue } count++ only = entry } switch count { case 1: return only, nil case 0: return nil, fmt.Errorf("no nodes connected") default: return nil, fmt.Errorf("multiple nodes connected; select one with /node ") } } func (r *Registry) resolveDisplayLabelLocked(ref string) (*NodeEntry, bool) { idx, ok := parseDisplayNodeIndex(ref) if !ok { return nil, false } entry, ok := r.byIndex[idx] return entry, ok } func parseDisplayNodeIndex(ref string) (int, bool) { if !strings.HasPrefix(ref, "node") { return 0, false } raw := strings.TrimPrefix(ref, "node") if raw == "" { return 0, false } idx, err := strconv.Atoi(raw) if err != nil || idx < 0 { return 0, false } return idx, true } func (r *Registry) All() []*NodeEntry { r.mu.RLock() defer r.mu.RUnlock() return r.collectLocked(false) } // AllReady returns only the dispatch-ready entries, sorted like All. Candidate // resolution and capability listing use it so a pending accepted connection is // never offered as a run/provider-pool dispatch target before its node signals // readiness. func (r *Registry) AllReady() []*NodeEntry { r.mu.RLock() defer r.mu.RUnlock() return r.collectLocked(true) } func (r *Registry) collectLocked(readyOnly bool) []*NodeEntry { out := make([]*NodeEntry, 0, len(r.byID)) for _, e := range r.byID { if readyOnly && !e.DispatchReady { continue } out = append(out, e) } sort.Slice(out, func(i, j int) bool { if out[i].HasIndex != out[j].HasIndex { return out[i].HasIndex } if out[i].HasIndex && out[i].Index != out[j].Index { return out[i].Index < out[j].Index } return out[i].NodeID < out[j].NodeID }) return out } func (r *Registry) Count() int { r.mu.RLock() defer r.mu.RUnlock() return len(r.byID) } // Pick is deprecated: use Resolve("") instead for single-node fallback, or Resolve(ref) for explicit selection. func (r *Registry) Pick() (*NodeEntry, error) { return r.Resolve("") }