- 노드 연결 2단계 핸드셰이크 추가 (Register→NodeReady→DispatchReady) - ConnectionGeneration 기반 연결 세대 관리로 stale 연결 차폐 - configured 노드 catalog 기반 snapshot rebuild (Connected 상태 분리) - provider 리스 소유권 일원화: edge가 소유권 승인·반환 전까지 대기 - 모델 대기열 승인/해제/스냅샷 서비스 구현 - 재연결 준비도 통합 테스트, 아카이브된 하위태스크 8건 포함
489 lines
15 KiB
Go
489 lines
15 KiB
Go
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
|
|
}
|
|
}
|