package service import ( edgenode "iop/apps/edge/internal/node" "iop/packages/go/config" iop "iop/proto/gen/iop" ) // getSnapshotForNodeLocked builds provider snapshots for nodeID without // acquiring m.mu. The caller must hold m.mu before invoking this variant so // that runtime config and queue state are observed atomically under a single // critical section. Used by the status path that linearizes the queue → service // reader-writer boundary. // connected reports the node's current connectivity from the registry. When // false, the catalog entry is preserved (id/adapter/type/category/models) but // every effective value is dropped to zero/offline so the snapshot reflects // disconnect without removing the configured provider from the list. disabled // providers remain disabled regardless of connectivity; only the connected // flag flips enabled providers from available to unavailable/offline. func (m *modelQueueManager) getSnapshotForNodeLocked(nodeID string, rec *edgenode.NodeRecord, connected bool) []*iop.ProviderSnapshot { var snaps []*iop.ProviderSnapshot // Catalog-first: nodes with a providers[] catalog emit only catalog snapshots. // Adapter snapshots (sections 1-4) are omitted to prevent duplicate metrics // when a node has both adapter instances and provider-pool catalog entries. if len(rec.Providers) > 0 { // Precompute current-candidate pressure once so the catalog loop reads // live candidate universe without rescanning the full queue per provider. pressure := m.providerQueuePressureLocked() for _, prov := range rec.Providers { if prov.ID == "" { continue } servedModels := make([]string, len(prov.Models)) copy(servedModels, prov.Models) lifecycleCaps := make([]string, len(prov.LifecycleCapabilities)) copy(lifecycleCaps, prov.LifecycleCapabilities) // Disabled providers appear in the snapshot with status=disabled and // effective capacity 0 so operators can see the switch state. if !config.ProviderEnabled(prov) { snaps = append(snaps, m.buildDisabledProviderSnapshot(&prov)) continue } capVal := prov.Capacity inflight, queued, longInflight, longQueued := m.providerSnapshotStatsLocked(nodeID, prov.ID, pressure) snaps = append(snaps, &iop.ProviderSnapshot{ Adapter: prov.Adapter, Status: effectiveStatus(connected), Health: effectiveHealth(connected, prov.Health), Capacity: int32(effectiveCount(connected, capVal)), InFlight: int32(effectiveCount(connected, inflight)), Queued: int32(effectiveCount(connected, queued)), Id: prov.ID, Type: prov.Type, Category: string(prov.Category), ServedModels: servedModels, LoadRatio: func() float32 { if !connected || capVal <= 0 { return 0 } return float32(inflight) / float32(capVal) }(), LifecycleCapabilities: lifecycleCaps, LongContextCapacity: int32(effectiveCount(connected, prov.LongContextCapacity)), LongInFlight: int32(effectiveCount(connected, longInflight)), LongQueued: int32(effectiveCount(connected, longQueued)), }) } return snaps } // Legacy adapter snapshots for nodes with no providers catalog. concurrencyFallback := 1 if rec.Runtime.Concurrency > 0 { concurrencyFallback = rec.Runtime.Concurrency } // 1. CLI if rec.Adapters.CLI.Enabled { capVal := concurrencyFallback inflight, queued := m.getStatsForAdapterLocked(nodeID, rec, "cli") snaps = append(snaps, &iop.ProviderSnapshot{ Adapter: "cli", Status: "available", Capacity: int32(capVal), InFlight: int32(inflight), Queued: int32(queued), }) } // 2. Ollama for _, inst := range rec.Adapters.OllamaInstances { if !inst.Enabled { continue } name := inst.Name if name == "" { name = "ollama" } capVal := inst.Capacity if capVal <= 0 { capVal = concurrencyFallback } inflight, queued := m.getStatsForAdapterLocked(nodeID, rec, name) snaps = append(snaps, &iop.ProviderSnapshot{ Adapter: name, Status: "available", Capacity: int32(capVal), InFlight: int32(inflight), Queued: int32(queued), }) } // 3. vLLM for _, inst := range rec.Adapters.VllmInstances { if !inst.Enabled { continue } name := inst.Name if name == "" { name = "vllm" } capVal := inst.Capacity if capVal <= 0 { capVal = concurrencyFallback } inflight, queued := m.getStatsForAdapterLocked(nodeID, rec, name) snaps = append(snaps, &iop.ProviderSnapshot{ Adapter: name, Status: "available", Capacity: int32(capVal), InFlight: int32(inflight), Queued: int32(queued), }) } // 4. OpenAI Compat for _, inst := range rec.Adapters.OpenAICompatInstances { if !inst.Enabled { continue } name := inst.Name if name == "" { name = "openai_compat" } capVal := inst.Capacity if capVal <= 0 { capVal = concurrencyFallback } inflight, queued := m.getStatsForAdapterLocked(nodeID, rec, name) snaps = append(snaps, &iop.ProviderSnapshot{ Adapter: name, Status: "available", Capacity: int32(capVal), InFlight: int32(inflight), Queued: int32(queued), }) } return snaps } func (m *modelQueueManager) getSnapshotForNode(nodeID string, rec *edgenode.NodeRecord, connected bool) []*iop.ProviderSnapshot { m.mu.Lock() defer m.mu.Unlock() return m.getSnapshotForNodeLocked(nodeID, rec, connected) } // buildDisabledProviderSnapshot returns a catalog-identity disabled provider // snapshot. Disabled providers keep status/health=disabled and all counters=0 // regardless of connectivity; the catalog entry is always preserved. func (m *modelQueueManager) buildDisabledProviderSnapshot(prov *config.NodeProviderConf) *iop.ProviderSnapshot { servedModels := make([]string, len(prov.Models)) copy(servedModels, prov.Models) lifecycleCaps := make([]string, len(prov.LifecycleCapabilities)) copy(lifecycleCaps, prov.LifecycleCapabilities) return &iop.ProviderSnapshot{ Adapter: prov.Adapter, Status: "disabled", Capacity: 0, InFlight: 0, Queued: 0, Id: prov.ID, Type: prov.Type, Category: string(prov.Category), ServedModels: servedModels, Health: "disabled", LoadRatio: 0, LifecycleCapabilities: lifecycleCaps, LongContextCapacity: 0, LongInFlight: 0, LongQueued: 0, } } // effectiveStatus returns "unavailable" when the node is disconnected, or // "available" when connected. Disabled providers are handled by the caller and // never pass through this helper. func effectiveStatus(connected bool) string { if connected { return "available" } return "unavailable" } // effectiveHealth returns "offline" when the node is disconnected, or the // provider's own health when connected. Disabled providers keep their own // "disabled" health and never flow through this helper. func effectiveHealth(connected bool, health string) string { if connected { return health } return "offline" } // effectiveCount returns zero when disconnected, or val when connected. func effectiveCount(connected bool, val int) int { if connected { return val } return 0 } // providerPressure accumulates queued and long-queued candidate counts for // a single (nodeID, providerID) pair across all model groups. Built once per // snapshot so the catalog loop can read both dimensions without rescanning // the full queue. type providerPressure struct { queued int longQueued int } // providerPressureMap is a nested map keyed by (nodeID, providerID) that // holds queued/longQueued candidate counts for every provider referenced by // pending items. Built once per snapshot via providerQueuePressureLocked. type providerPressureMap map[string]map[string]*providerPressure // providerQueuePressureLocked aggregates candidate pressure across all model // groups using the live candidate resolver for each pending item. One item's // (nodeID, providerID) pair is counted at most once regardless of how many // times it appears in that item's candidate slice. Resolver error or an empty // result excludes the item from the aggregation (no stale snapshot fallback). // Must be called with m.mu held. func (m *modelQueueManager) providerQueuePressureLocked() providerPressureMap { pm := make(providerPressureMap) for _, group := range m.groups { for _, item := range group.queue { candidates, outcome, _ := m.resolveQueuedCandidatesLocked(item) if outcome != resolveOk { continue } seen := make(map[nodeProvKey]bool, len(candidates)) for i := range candidates { c := &candidates[i] if c.entry == nil { continue } key := nodeProvKey{nodeID: c.entry.NodeID, providerID: c.providerID} if seen[key] { continue } seen[key] = true if pm[key.nodeID] == nil { pm[key.nodeID] = make(map[string]*providerPressure) } pp := pm[key.nodeID][key.providerID] if pp == nil { pp = &providerPressure{} pm[key.nodeID][key.providerID] = pp } pp.queued++ if item.long { pp.longQueued++ } } } } return pm } // providerSnapshotStatsLocked returns lease-backed inflight counts and the // pre-aggregated candidate pressure for (nodeID, providerID). The pressure // map is built once per snapshot by providerQueuePressureLocked. Must be // called with m.mu held. func (m *modelQueueManager) providerSnapshotStatsLocked(nodeID, providerID string, pressure providerPressureMap) (inFlight, queued, longInFlight, longQueued int) { key := providerResourceKey{nodeID: nodeID, providerID: providerID} if res, ok := m.resources[key]; ok { inFlight = res.inFlight longInFlight = res.longInFlight } if nodeMap, ok := pressure[nodeID]; ok { if pp, ok := nodeMap[providerID]; ok { queued = pp.queued longQueued = pp.longQueued } } return } // nodeProvKey is a local dedup key inside providerQueuePressureLocked. type nodeProvKey struct { nodeID string providerID string } // getStatsForProviderLocked returns lease-backed in-flight and candidate // pressure for a provider-pool provider identified by (nodeID, providerID). // queued is not a provider-owned queue depth: one Edge pending request is // counted once for every provider candidate it includes, so queued values // across provider snapshots are not additive. Must be called with m.mu held. // // Deprecated: prefer providerQueuePressureLocked + providerSnapshotStatsLocked // for snapshot path. Kept for direct test usage that does not build a pressure // map. func (m *modelQueueManager) getStatsForProviderLocked(nodeID, providerID string) (inFlight, queued int) { key := providerResourceKey{nodeID: nodeID, providerID: providerID} if res, ok := m.resources[key]; ok { inFlight = res.inFlight } for _, group := range m.groups { for _, item := range group.queue { for _, c := range item.candidates { if c.entry.NodeID == nodeID && c.providerID == providerID { queued++ break } } } } return } // getLongStatsForProviderLocked returns lease-backed long-context in-flight and // long-request candidate pressure for a provider-pool provider identified by // (nodeID, providerID). As with queued, one multi-candidate request may contribute // to more than one provider snapshot. Must be called with m.mu held. func (m *modelQueueManager) getLongStatsForProviderLocked(nodeID, providerID string) (longInFlight, longQueued int) { key := providerResourceKey{nodeID: nodeID, providerID: providerID} if res, ok := m.resources[key]; ok { longInFlight = res.longInFlight } for _, group := range m.groups { for _, item := range group.queue { if item.long { for _, c := range item.candidates { if c.entry.NodeID == nodeID && c.providerID == providerID { longQueued++ break } } } } } return } // longInFlightForProvider returns the long-context in-flight count for a provider, // acquiring the manager lock. Exposed for status/snapshot reporting and tests. func (m *modelQueueManager) longInFlightForProvider(nodeID, providerID string) int { m.mu.Lock() defer m.mu.Unlock() longInFlight, _ := m.getLongStatsForProviderLocked(nodeID, providerID) return longInFlight } func (m *modelQueueManager) getStatsForAdapterLocked(nodeID string, rec *edgenode.NodeRecord, adapterName string) (inFlight, queued int) { for _, group := range m.groups { canonical, ok := resolveSnapshotAdapterName(rec, group.adapter, group.target) if ok && canonical == adapterName { if val, ok := group.inflight[nodeID]; ok { inFlight += val } for _, item := range group.queue { for _, c := range item.candidates { if c.entry.NodeID == nodeID { queued++ break } } } } } return } // resolveSnapshotAdapterName returns the adapter/instance key used to match // queue state for a provider snapshot. // When adapterType is empty it falls back to the provider id (target) so that // queue stats are still resolved instead of being dropped on an empty key. func resolveSnapshotAdapterName(rec *edgenode.NodeRecord, adapterType, target string) (string, bool) { // Fallback to provider id when adapter is not set. if adapterType == "" { adapterType = target } if rec == nil { return adapterType, true } // 1. Exact instance Name match (highest priority). for _, inst := range rec.Adapters.OllamaInstances { if inst.Name == adapterType { return adapterType, true } } for _, inst := range rec.Adapters.VllmInstances { if inst.Name == adapterType { return adapterType, true } } for _, inst := range rec.Adapters.OpenAICompatInstances { if inst.Name == adapterType { return adapterType, true } } // 2. Type-name route. switch adapterType { case "ollama": var enabled []string for _, inst := range rec.Adapters.OllamaInstances { if inst.Enabled { name := inst.Name if name == "" { name = "ollama" } enabled = append(enabled, name) } } switch len(enabled) { case 0: return "ollama", true case 1: return enabled[0], true default: return "", false } case "vllm": var enabled []string for _, inst := range rec.Adapters.VllmInstances { if inst.Enabled { name := inst.Name if name == "" { name = "vllm" } enabled = append(enabled, name) } } switch len(enabled) { case 0: return "vllm", true case 1: return enabled[0], true default: return "", false } case "openai_compat": var enabled []string for _, inst := range rec.Adapters.OpenAICompatInstances { if inst.Enabled { name := inst.Name if name == "" { name = "openai_compat" } enabled = append(enabled, name) } } switch len(enabled) { case 0: return "openai_compat", true case 1: return enabled[0], true default: return "", false } case "cli": return "cli", true default: return adapterType, true } }