iop/apps/edge/internal/service/model_queue_snapshot.go
toki 4dcf6f0cf9 feat(edge): provider 리소스 승인 소유권 정렬 구현
- 노드 연결 2단계 핸드셰이크 추가 (Register→NodeReady→DispatchReady)
- ConnectionGeneration 기반 연결 세대 관리로 stale 연결 차폐
- configured 노드 catalog 기반 snapshot rebuild (Connected 상태 분리)
- provider 리스 소유권 일원화: edge가 소유권 승인·반환 전까지 대기
- 모델 대기열 승인/해제/스냅샷 서비스 구현
- 재연결 준비도 통합 테스트, 아카이브된 하위태스크 8건 포함
2026-07-22 18:10:54 +09:00

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
}
}