Documentation
¶
Index ¶
- Constants
- type GPUAffinityScore
- type GPUDiscovery
- func (gd *GPUDiscovery) CheckGPUHealth(ctx context.Context, gpus []*common.GPUDevice) ([]*common.GPUDevice, error)
- func (gd *GPUDiscovery) CheckGPUOversubscription(gpus []*common.GPUDevice) bool
- func (gd *GPUDiscovery) ClearGPUCache(ctx context.Context) error
- func (gd *GPUDiscovery) DiscoverGPUs(ctx context.Context) ([]*common.GPUDevice, error)
- func (gd *GPUDiscovery) DiscoverGPUsFromK8s(ctx context.Context) ([]*common.GPUDevice, error)
- func (gd *GPUDiscovery) FilterGPUsByMemory(gpus []*common.GPUDevice, minGBMemory float64) []*common.GPUDevice
- func (gd *GPUDiscovery) FilterGPUsByType(gpus []*common.GPUDevice, gpuType string) []*common.GPUDevice
- func (gd *GPUDiscovery) FilterGPUsByUtilization(gpus []*common.GPUDevice, maxPercent float64) []*common.GPUDevice
- func (gd *GPUDiscovery) GetGPUInventory(ctx context.Context) (*GPUInventory, error)
- func (gd *GPUDiscovery) IsGPUAvailable(gpus []*common.GPUDevice, index int) bool
- func (gd *GPUDiscovery) MonitorGPUMetrics(ctx context.Context, interval time.Duration)
- func (gd *GPUDiscovery) NodeHasGPUs(ctx context.Context) (bool, error)
- func (gd *GPUDiscovery) SumGPUMemory(gpus []*common.GPUDevice) float64
- func (gd *GPUDiscovery) ValidateGPUIndices(gpus []*common.GPUDevice, indices []int) error
- type GPUInventory
- type GPUTopologyManager
- func (gtm *GPUTopologyManager) ClearTopologyCache(ctx context.Context) error
- func (gtm *GPUTopologyManager) DetectNUMAMapping(ctx context.Context) (map[int]int, error)
- func (gtm *GPUTopologyManager) DetectNVLink(ctx context.Context) ([][]int, error)
- func (gtm *GPUTopologyManager) DetectNVSwitchDomains(gpuCount int, nvlinkPairs [][]int, nvlinkCounts map[string]int) map[int][]int
- func (gtm *GPUTopologyManager) DetectPCIeGeneration(ctx context.Context) (map[int]int, error)
- func (gtm *GPUTopologyManager) DetectTopology(ctx context.Context) (*common.GPUTopology, error)
- func (gtm *GPUTopologyManager) FindNvidiaSMI() string
- func (gtm *GPUTopologyManager) ScoreGPUPlacement(gpu1Index int, gpu2Index int, topology *common.GPUTopology) *GPUAffinityScore
- func (gtm *GPUTopologyManager) ScoreGPUSet(gpuIndices []int, topology *common.GPUTopology) float64
- func (gtm *GPUTopologyManager) SelectBestGPUSet(ctx context.Context, availableGPUs []*common.GPUDevice, requestedCount int, ...) ([]int, *GPUAffinityScore, error)
Constants ¶
const ( CacheKeyGPUDevicesPrefix = "ares:node:gpu:devices" CacheKeyGPUTopology = "ares:node:gpu:topology" CacheKeyGPUHealth = "ares:node:gpu:health" DefaultCacheTTL = 30 * time.Second HealthCheckInterval = 60 * time.Second )
Cache keys
const ( CacheKeyTopologyPrefix = "ares:node:gpu:topology" NVLinkBandwidth = 900.0 PCIeGen4Bandwidth = 32.0 PCIeGen5Bandwidth = 64.0 NVLinkMultiplier = 28.1 SameNUMABonus = 1.5 TopologyCacheTTL = 60 * time.Second )
Cache and constants
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type GPUAffinityScore ¶
type GPUDiscovery ¶
type GPUDiscovery struct {
// contains filtered or unexported fields
}
GPUDiscovery: Discovers and monitors local GPUs Primary: Uses Kubernetes API to query node allocatable resources Fallback: Uses nvidia-smi for detection Caches results in Redis Thread-safe: All operations protected by context cancellation
func NewGPUDiscovery ¶
func NewGPUDiscovery(redisClient *redis.RedisClient) *GPUDiscovery
Original constructor (for backward compatibility) NewGPUDiscovery: Create new GPU discovery service (nvidia-smi only)
func NewGPUDiscoveryWithK8s ¶
func NewGPUDiscoveryWithK8s(redisClient *redis.RedisClient, k8sClient kubernetes.Interface, nodeName string) *GPUDiscovery
✅ NEW: Constructor with K8s client support
func (*GPUDiscovery) CheckGPUHealth ¶
func (gd *GPUDiscovery) CheckGPUHealth(ctx context.Context, gpus []*common.GPUDevice) ([]*common.GPUDevice, error)
CheckGPUHealth: Run health checks on all GPUs Returns: Array of GPUs with health status updated Checks: Accessibility, memory errors, temperature, power
func (*GPUDiscovery) CheckGPUOversubscription ¶
func (gd *GPUDiscovery) CheckGPUOversubscription(gpus []*common.GPUDevice) bool
CheckGPUOversubscription: Detect if GPUs are oversubscribed (bad for latency) Returns: True if any GPU utilization > 90%
func (*GPUDiscovery) ClearGPUCache ¶
func (gd *GPUDiscovery) ClearGPUCache(ctx context.Context) error
ClearGPUCache: Clear GPU cache (for testing or forced refresh)
func (*GPUDiscovery) DiscoverGPUs ¶
✅ NEW: Enhanced DiscoverGPUs with fallback strategy Try K8s API first, fall back to nvidia-smi if needed
func (*GPUDiscovery) DiscoverGPUsFromK8s ¶
NEW: DiscoverGPUsFromK8s - Query Kubernetes API for GPU allocatable resources This is the PREFERRED method when K8s client is available Returns: Array of GPU devices based on node allocatable resources Latency: ~50-100ms (K8s API call)
func (*GPUDiscovery) FilterGPUsByMemory ¶
func (gd *GPUDiscovery) FilterGPUsByMemory(gpus []*common.GPUDevice, minGBMemory float64) []*common.GPUDevice
FilterGPUsByMemory: Get GPUs with minimum memory available Returns: Array of GPUs with at least minGBMemory available
func (*GPUDiscovery) FilterGPUsByType ¶
func (gd *GPUDiscovery) FilterGPUsByType(gpus []*common.GPUDevice, gpuType string) []*common.GPUDevice
FilterGPUsByType: Get GPUs matching specific type Returns: Array of GPUs of given type Example: FilterGPUsByType("A100") -> all A100 GPUs
func (*GPUDiscovery) FilterGPUsByUtilization ¶
func (gd *GPUDiscovery) FilterGPUsByUtilization(gpus []*common.GPUDevice, maxPercent float64) []*common.GPUDevice
FilterGPUsByUtilization: Get least utilized GPUs Returns: Array of healthy GPUs with utilization <= maxPercent
func (*GPUDiscovery) GetGPUInventory ¶
func (gd *GPUDiscovery) GetGPUInventory(ctx context.Context) (*GPUInventory, error)
GetGPUInventory: Get summary of all GPUs on node Returns: Total GPUs, types, and capacity summary Latency: <5ms (cached)
func (*GPUDiscovery) IsGPUAvailable ¶
func (gd *GPUDiscovery) IsGPUAvailable(gpus []*common.GPUDevice, index int) bool
IsGPUAvailable: Check if specific GPU is healthy
func (*GPUDiscovery) MonitorGPUMetrics ¶
func (gd *GPUDiscovery) MonitorGPUMetrics(ctx context.Context, interval time.Duration)
MonitorGPUMetrics: Continuously monitor GPU metrics Runs in background, updates cache periodically Call as: go discovery.MonitorGPUMetrics(ctx)
func (*GPUDiscovery) NodeHasGPUs ¶
func (gd *GPUDiscovery) NodeHasGPUs(ctx context.Context) (bool, error)
NodeHasGPUs: Check if node has any GPUs
func (*GPUDiscovery) SumGPUMemory ¶
func (gd *GPUDiscovery) SumGPUMemory(gpus []*common.GPUDevice) float64
SumGPUMemory: Get total available memory across GPUs
func (*GPUDiscovery) ValidateGPUIndices ¶
func (gd *GPUDiscovery) ValidateGPUIndices(gpus []*common.GPUDevice, indices []int) error
ValidateGPUIndices: Check if requested GPU indices exist and are healthy Used before scheduling job to ensure GPUs available
type GPUInventory ¶
type GPUInventory struct {
TotalGPUs int `json:"total_gpus"`
HealthyGPUs int `json:"healthy_gpus"`
GPUsByType map[string]int `json:"gpus_by_type"`
TotalMemoryGB float64 `json:"total_memory_gb"`
Devices []*common.GPUDevice `json:"devices"`
}
GPUInventory: Summary of node GPU capacity
func (*GPUInventory) AvailableGPUs ¶
func (inv *GPUInventory) AvailableGPUs() int
AvailableGPUs: Count of healthy GPUs
type GPUTopologyManager ¶
type GPUTopologyManager struct {
// contains filtered or unexported fields
}
GPUTopologyManager: Manages GPU topology and affinity
func NewGPUTopologyManager ¶
func NewGPUTopologyManager(redisClient *redis.RedisClient, discovery *GPUDiscovery) *GPUTopologyManager
NewGPUTopologyManager: Create new topology manager
func (*GPUTopologyManager) ClearTopologyCache ¶
func (gtm *GPUTopologyManager) ClearTopologyCache(ctx context.Context) error
func (*GPUTopologyManager) DetectNUMAMapping ¶
func (*GPUTopologyManager) DetectNVLink ¶
func (gtm *GPUTopologyManager) DetectNVLink(ctx context.Context) ([][]int, error)
func (*GPUTopologyManager) DetectNVSwitchDomains ¶
func (gtm *GPUTopologyManager) DetectNVSwitchDomains( gpuCount int, nvlinkPairs [][]int, nvlinkCounts map[string]int, ) map[int][]int
DetectNVSwitchDomains: Identify groups of GPUs that are mutually connected at full NVLink bandwidth (one "domain" = one max-bandwidth island).
The number of domains is HARDWARE-DEPENDENT — this is the common confusion:
- Full-crossbar NVSwitch (DGX A100 / p4d.24xlarge): every pair is NV12, so union-find collapses to a SINGLE domain of all 8 GPUs. Topology-aware placement is effectively a no-op there — any GPU set is equally optimal.
- Hybrid cube-mesh, no NVSwitch (DGX-1 V100): pairs are a mix of NV1/NV2/NV6, so thresholding splits GPUs into multiple domains and placement matters.
Algorithm: GPUs connected at the highest link count form a domain. Uses union-find on pairs with link count >= threshold (max_links * 0.8).
func (*GPUTopologyManager) DetectPCIeGeneration ¶
func (*GPUTopologyManager) DetectTopology ¶
func (gtm *GPUTopologyManager) DetectTopology(ctx context.Context) (*common.GPUTopology, error)
func (*GPUTopologyManager) FindNvidiaSMI ¶
func (gtm *GPUTopologyManager) FindNvidiaSMI() string
func (*GPUTopologyManager) ScoreGPUPlacement ¶
func (gtm *GPUTopologyManager) ScoreGPUPlacement( gpu1Index int, gpu2Index int, topology *common.GPUTopology, ) *GPUAffinityScore
func (*GPUTopologyManager) ScoreGPUSet ¶
func (gtm *GPUTopologyManager) ScoreGPUSet( gpuIndices []int, topology *common.GPUTopology, ) float64