Documentation
¶
Index ¶
- Constants
- Variables
- func GetDefaultWorkloadContainerNamesToWatch() []string
- func GetUseUUIDForRequestObjName() bool
- func IsPodAdmissionRejected(ps corev1.PodStatus) bool
- func IsPodStuckInitializing(pod *corev1.Pod, k8sTimeConfig *k8sutil.TimeConfig) (bool, types.ICMSInstanceState)
- func NewCommand() *cobra.Command
- func NewK8sComputeBackend(clients *kubeclients.KubeClients, bk8s *BackendK8sCache) (ICMSRequestHelper, K8sArtifactHelper)
- func NewLedgerEventCorrelatorOptions(heartbeatInterval time.Duration) record.CorrelatorOptions
- func SetUseUUIDForRequestObjName(v bool)
- type Agent
- func (a *Agent) IsRequestFromClusterQueue(_ context.Context, req *nvcav2beta1.ICMSRequest) bool
- func (a *Agent) PostICMSInstanceRequestStatusUpdates(ctx context.Context) error
- func (a *Agent) PutICMSRequestAcknowledgement(ctx context.Context) error
- func (a *Agent) PutNVCAStatusUpdate(ctx context.Context) error
- func (a *Agent) RegisterWithICMS(ctx context.Context) (*types.ICMSRegistrationResponse, error)
- func (a *Agent) RenewICMSQueueCreds(ctx context.Context) error
- func (a *Agent) Start(ctx context.Context) error
- func (a *Agent) SyncPeriodicInstanceStatuses(ctx context.Context) error
- type AgentOptions
- type AggregatedInstanceStatus
- type BackendK8sCache
- func (c *BackendK8sCache) ApplyICMSRequestStatusChange(ctx context.Context, req *nvcav2beta1new.ICMSRequest) error
- func (c *BackendK8sCache) CleanupCachedResources(ctx context.Context, allReqs []*nvcav2beta1new.ICMSRequest) error
- func (c *BackendK8sCache) CleanupCreationRequestResources(ctx context.Context, req *nvcav2beta1new.ICMSRequest) error
- func (c *BackendK8sCache) CleanupUnusedResources(ctx context.Context) error
- func (c *BackendK8sCache) CreateICMSCreationMessageRequest(ctx context.Context, msg translate.CreationQueueMessageMetadataGetter, ...) (*nvcav2beta1new.ICMSRequest, error)
- func (c *BackendK8sCache) CreateICMSTerminationMessageRequest(ctx context.Context, tm types.ICMSTerminationMessage, msgReceipt, msgID string) error
- func (c *BackendK8sCache) EmitICMSEvent(req *nvcav2beta1.ICMSRequest, eventType, reason, message string, ...)
- func (c *BackendK8sCache) EmitICMSEventf(req *nvcav2beta1.ICMSRequest, eventType, reason, msgFmt string, ...)
- func (c *BackendK8sCache) FetchICMSRegistrationResponse(ctx context.Context) (*types.ICMSRegistrationResponse, error)
- func (c *BackendK8sCache) ForceSync(ctx context.Context)
- func (c *BackendK8sCache) GetAllBackendGPUs(ctx context.Context) ([]types.BackendGPU, error)
- func (c *BackendK8sCache) GetAllPodsForRequest(ctx context.Context, reqID string) ([]corev1.Pod, error)
- func (c *BackendK8sCache) GetBARTRegistrationResponseSecret(ctx context.Context) (*corev1.Secret, error)
- func (c *BackendK8sCache) GetComponentStatus(ctx context.Context) (hs types.AgentHealth, err error)
- func (c *BackendK8sCache) GetGPUResource(ctx context.Context, gpuName types.GPUName) (types.GPUResource, error)
- func (c *BackendK8sCache) GetNodeFeaturesClient() nodefeatures.Client
- func (c *BackendK8sCache) GetNodeInformer() cache.SharedIndexInformer
- func (c *BackendK8sCache) GetRegisteredBackendGPUs(ctx context.Context, backendGPUs []types.BackendGPU, ...) ([]types.RegistrationGPU, error)
- func (c *BackendK8sCache) ReconcileInstanceStatus(ctx context.Context, is types.ICMSServerInstanceState) (srToReturn *nvcav2beta1new.ICMSRequest, ...)
- func (c *BackendK8sCache) StoreICMSRegistrationResponse(ctx context.Context, res *types.ICMSRegistrationResponse) error
- func (c *BackendK8sCache) StoreUpdatedCredentials(ctx context.Context, qCreds types.QueueCredentials) error
- func (c *BackendK8sCache) SyncAllICMSRequests(ctx context.Context) error
- func (c *BackendK8sCache) SyncICMSRequest(ctx context.Context, nn apitypes.NamespacedName) error
- func (c *BackendK8sCache) SyncICMSRequestByID(ctx context.Context, icmsReqID, messageBatchID string) error
- func (c *BackendK8sCache) SyncNetworkPolicies(ctx context.Context) error
- func (c *BackendK8sCache) UpdateInstanceTypeMetrics(ctx context.Context, gpus []types.RegistrationGPU, ...) error
- func (c *BackendK8sCache) UpdateSchedulerWorkloadMetrics(ctx context.Context)
- type BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) Start(ctx context.Context) (*BackendK8sCache, <-chan *core.Event, error)
- func (b *BackendK8sCacheBuilder) WithCSIVolumeMountOptions(mntOptions []string) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithCachingSupport(enable, encryption bool) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithClients(clients *kubeclients.KubeClients) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithClusterName(clusterName string) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithClusterProvider(cloudProvider string) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithClusterRegion(clusterRegion string) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithComputeBackend(be BackendType) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithConfig(cfg nvcaconfig.Config) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithDynamicNodeFeatureClient(isDynClient bool, opts nodefeatures.DynamicClientOptions) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithEnvOverrides(functionOverrides, taskOverrides map[string]string) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithFNDSClient(client fnds.Client) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithFeatureFlagFetcher(fff featureflag.Fetcher) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithHelmInternalPersistentStorage(enable bool) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithHelmRBACEnforcement(enable bool) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithHelmRepositoryPrefix(hrepo string) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithHelmResourceConstraints(enable bool) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithHelmSharedStorage(enable bool) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithICMSRequestSyncConcurrency(c int) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithImageCredentialHelperImage(image string) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithInfraOverheadGetter(infraOverheadGetter enforce.InfraOverheadGetter) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithLogPosting(enable bool) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithLowLatencyStreaming(enable bool) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithNamespaceLabels(nsl labels.Set) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithNotFoundInstanceStatusReporter(reporter func(ctx context.Context, reqID, instanceID string) error) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithOTelTracer(tracer oteltrace.Tracer) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithPVCRebind(on bool) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithPeriodicInstanceStatusUpdate(enable bool, period time.Duration) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithRequestsNamespace(requestsNamespace string) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithSecretMirrorConfig(namespace, selector string) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithStaticGPUCapacity(gpuCap uint64) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithSystemNamespace(systemNamespace string) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithTimeConfig(cfg *k8sutil.TimeConfig) *BackendK8sCacheBuilder
- func (b *BackendK8sCacheBuilder) WithWorkerDegradationHandler(enable bool) *BackendK8sCacheBuilder
- type BackendType
- type BootstrapTokenString
- type CacheAccessObj
- type ComputeBackend
- type GPUMonitor
- func (m *GPUMonitor) GetComponentStatus(ctx context.Context) (types.AgentHealth, error)
- func (m *GPUMonitor) HasGPUs() bool
- func (m *GPUMonitor) SetHasGPUs(hasGPUs bool)
- func (m *GPUMonitor) SetOnGPUStateChange(cb GPUStateChangeCallback)
- func (m *GPUMonitor) Start(ctx context.Context)
- func (m *GPUMonitor) Stop()
- type GPUMonitorConfig
- type GPUMonitorOption
- type GPUStateChangeCallback
- type ICMSClient
- func (c *ICMSClient) Endpoint() string
- func (c *ICMSClient) GetCreds(ctx context.Context) (*types.ICMSCredentialResponse, error)
- func (c *ICMSClient) GetICMSServerInstanceStatuses(ctx context.Context) (types.ICMSInstanceStatusResponse, error)
- func (c *ICMSClient) PostInstanceStatusUpdate(ctx context.Context, requestID, instanceID string, ...) error
- func (c *ICMSClient) PutHealthStatus(ctx context.Context, hsr *types.HealthStatusRequest) (*types.HealthStatusResponse, error)
- func (c *ICMSClient) PutRequestAcknowledgement(ctx context.Context, icmsReqID string, messageBatchID string, ...) error
- func (c *ICMSClient) Register(ctx context.Context, rreq *types.ICMSRegistrationRequest) (*types.ICMSRegistrationResponse, error)
- type ICMSClientInterface
- type ICMSInstanceReconcileState
- type ICMSRequestHelper
- type InitCacheJobState
- type JWKSUpdater
- type JWKSUpdaterOptions
- type K8sArtifactHelper
- type K8sComputeBackend
- func (c K8sComputeBackend) AggregateInstanceStatuses(ctx context.Context, req *nvcav2beta1.ICMSRequest) AggregatedInstanceStatus
- func (c K8sComputeBackend) AggregatePodInstanceStatus(ctx context.Context, req *nvcav2beta1.ICMSRequest, instanceID string) AggregatedInstanceStatus
- func (c K8sComputeBackend) AllInstancesTerminatedAndReported(ctx context.Context, req *nvcav2beta1.ICMSRequest) bool
- func (c K8sComputeBackend) ApplyCreationMessage(ctx context.Context, req *nvcav2beta1.ICMSRequest) error
- func (c K8sComputeBackend) ApplyTerminationMessage(ctx context.Context, req *nvcav2beta1.ICMSRequest) error
- func (c K8sComputeBackend) CheckInitCacheJobState(ctx context.Context, rwPVCName string, job *batchv1.Job) InitCacheJobState
- func (c K8sComputeBackend) CheckPVCState(ctx context.Context, roPVCName string) (PVCState, error)
- func (c K8sComputeBackend) CleanupModelCachingResources(ctx context.Context, rwPVC *v1.PersistentVolumeClaim, initJobName string) error
- func (c K8sComputeBackend) CleanupModelCachingSetupArtifacts(ctx context.Context, req *nvcav2beta1.ICMSRequest) error
- func (c K8sComputeBackend) ComputeCleanupCacheReferences(ctx context.Context, cacheReferences []string) error
- func (c K8sComputeBackend) CreateConfigMapArtifact(ctx context.Context, a function.LaunchArtifact, mf mutateFunc) error
- func (c K8sComputeBackend) CreatePodArtifact(ctx context.Context, podArt function.LaunchArtifact, mf mutateFunc) error
- func (c K8sComputeBackend) CreatePodArtifactInstances(ctx context.Context, pod *corev1.Pod, req *nvcav2beta1.ICMSRequest, ...) ([]nvcav2beta1.InstanceStatus, error)
- func (c K8sComputeBackend) CreateSecretArtifact(ctx context.Context, a function.LaunchArtifact, mf mutateFunc) error
- func (c K8sComputeBackend) CreateServiceArtifact(ctx context.Context, a function.LaunchArtifact, mf mutateFunc) error
- func (c K8sComputeBackend) GetErroredPodLogs(ctx context.Context, pod *corev1.Pod, prepend string, writeMaxBytes int64) (string, int64, error)
- func (c K8sComputeBackend) GetICMSRequestStatusUpdatesForRequest(ctx context.Context, req *nvcav2beta1.ICMSRequest) ([]types.ICMSRequestUpdateInfo, error)
- func (c K8sComputeBackend) GetICMSRequestUpdatesForCreatePodRequest(ctx context.Context, st nvcav2beta1.InstanceStatus, ...) (types.ICMSRequestUpdateInfo, error)
- func (c K8sComputeBackend) GetICMSRequestUpdatesForCreateRequest(ctx context.Context, req *nvcav2beta1.ICMSRequest) []types.ICMSRequestUpdateInfo
- func (c K8sComputeBackend) GetICMSRequestUpdatesForMiniServiceRequest(ctx context.Context, req *nvcav2beta1.ICMSRequest, ...) (nvcatypes.ICMSRequestUpdateInfo, error)
- func (c K8sComputeBackend) GetICMSRequestUpdatesForTerminationRequest(ctx context.Context, req *nvcav2beta1.ICMSRequest) []types.ICMSRequestUpdateInfo
- func (c K8sComputeBackend) HandleInstanceStatusPreconditionFailure(ctx context.Context, req *nvcav2beta1.ICMSRequest, instID string) error
- func (c K8sComputeBackend) PurgeInstanceID(ctx context.Context, req *nvcav2beta1.ICMSRequest, ...) bool
- func (c K8sComputeBackend) SetupInitCacheJobBlockDevice(ctx context.Context, rwPVCObj *v1.PersistentVolumeClaim, initJob *batchv1.Job, ...) error
- func (c K8sComputeBackend) SetupModelCachingForRequest(ctx context.Context, rwPVC *v1.PersistentVolumeClaim, initJob *batchv1.Job, ...) (ModelCachingState, string)
- func (c K8sComputeBackend) SetupPVCForReaders(ctx context.Context, rwPVC *v1.PersistentVolumeClaim, initJobName string, ...) (ROPVCSetupPhase, error)
- type ModelCachingState
- type PVCState
- type QueueManager
- func (qm *QueueManager) DeleteCreationMessage(ctx context.Context, gpuName, rhdl string) error
- func (qm *QueueManager) DeleteCreationMessageV2(ctx context.Context, rhdl, queueURL string) error
- func (qm *QueueManager) ExtendCreationMessableVisibilityTimeout(ctx context.Context, gpuName, rhdl string) error
- func (qm *QueueManager) ExtendCreationMessableVisibilityTimeoutV2(ctx context.Context, rhdl, queueURL string) error
- func (qm *QueueManager) IsGPUAtCapacity(gpuName types.GPUName) bool
- func (qm *QueueManager) IsPaused() bool
- func (qm *QueueManager) Name() string
- func (qm *QueueManager) Pause()
- func (qm *QueueManager) Resume()
- func (qm *QueueManager) SetGPUAtCapacity(gpuName types.GPUName, atCapacity bool) bool
- func (qm *QueueManager) SetStatusOK(ok bool)
- func (qm *QueueManager) StatusOK() bool
- func (qm *QueueManager) SyncQueues(ctx context.Context) error
- type ROPVCSetupPhase
- type ValidatorSummaryReconciler
Constants ¶
const ( UnusedResourceCleanupDuration = 30 * time.Minute ICMSRequestAckMaxGoroutines = 20 ICMSInstanceRequestStatusUpdatesMaxGoroutines = 20 )
const ( // #nosec G101 EventTickRenewICMSCredentials string = "TICK_RENEW_ICMS_CREDENTIALS" EventTickUpdateHeartbeat string = "TICK_UPDATE_HEARTBEAT" EventTickSyncSQSQueue string = "TICK_SYNC_SQS_QUEUE" EventTickSyncICMSRequestStatus string = "TICK_SYNC_ICMS_REQUEST_STATUS" EventTickSyncICMSRequests string = "TICK_SYNC_ICMS_REQUESTS" EventTickAcknowledgeRequest string = "TICK_ACKNOWLEDGE_REQUEST" EventTickSyncCleanupUnusedResources string = "TICK_SYNC_CLEANUP_UNUSED_RESOURCES" EventTickSyncPeriodicInstanceStatusUpdates string = "TICK_SYNC_PERIODIC_INSTANCE_STATUS_UPDATES" EventTickUpdateICMSRegistration string = "TICK_SYNC_UPDATE_ICMS_REGISTRATION" EventTickSyncNetworkPolicies string = "TICK_SYNC_NETWORK_POLICIES" // metric events EventModelCachingFailed string = "EVENT_MODEL_CACHING_FAILED" EventModelCachingSuccess string = "EVENT_MODEL_CACHING_SUCCESS" EventPVCModelCachingError string = "EVENT_PVC_MODEL_CACHING_ERROR" EventTranslateFunctionError string = "EVENT_TRANSLATE_FUNCTION_ERROR" EventTranslateTaskError string = "EVENT_TRANSLATE_TASK_ERROR" EventWorkloadPaused string = "EVENT_WORKLOAD_PAUSED" )
Controller tick types.
const ( SystemNamespace = "nvca-system" RequestsNamespace = "nvcf-backend" // #nosec G101 RegistrationInfoSecretName = "icms-registration-info" BartServiceAPIKeySecretName = "bart-service-api-key-secret" ResyncInterval = 30 * time.Minute SQSMessageIDKey = "SQSMessageID" K8sNameLabelKey = "kubernetes.io/metadata.name" NVCAFinalizer = "nvca.finalizers.nvidia.io" NVCADeploymentName = "nvca" SecretMirroredFromLabelKey = "nvca.nvcf.nvidia.io/mirrored-from-namespace" )
const ( MaxFailedPodLogLines = int64(20) MaxBytesForPodLogs = int64(1024) )
const ( // GPUMonitorComponentName is the name used in health status reporting. GPUMonitorComponentName = "gpumonitor" // DefaultGPUPollInterval is the default interval between GPU availability checks. // With event-driven monitoring via node informers, this acts as a fallback safety net. DefaultGPUPollInterval = 60 * time.Second // DefaultGPUDebounceTime is the default time a state must be stable before triggering a change. DefaultGPUDebounceTime = 30 * time.Second )
const ( UnexpectedAdmissionErrReason = "UnexpectedAdmissionError" ImagePullIssueReason = "ErrImagePull" ImagePullIssueAlternateReason = "ImagePullBackOff" InferenceContainerName = "inference" InitContainerName = "init" RWPVCSuffix = "rw-pvc" ROPVCSuffix = "ro-pvc" ModelVolumeName = "model-data" )
const ( // NvSnapRestoreFromAnnotation, when present and non-empty, tells // NvSnap's mutating webhook to inject the restore mounts for // the given content-addressed checkpoint hash. NvSnap's webhook // resolves the hash to a manifest ConfigMap, picks a node where // the cache lives (or override target node), and injects the // rootfs volume + sitecustomize plumbing. NvSnapRestoreFromAnnotation = "nvsnap.io/restore-from" // NvSnapCheckpointOnWarmAnnotation marks pods that NVCA's // post-Ready reconciler (Hook B, PR-5) should checkpoint after // the warmup window. NVCA stamps this at pod-create time on // every function-version pod whose NvSnapFunctionState is not // opted out, so the reconciler doesn't need to re-resolve the // FV state at every Ready event. NvSnapCheckpointOnWarmAnnotation = "nvsnap.io/checkpoint-on-warm" // NvSnapFunctionVersionIDAnnotation is the bridge between // ICMSRequest-driven Hook A (which knows the FV id from // req.Spec.FunctionDetails) and the Pod-watching reconciler // (which doesn't see the originating ICMSRequest). Stamped at // the same time as CheckpointOnWarmAnnotation. NvSnapFunctionVersionIDAnnotation = "nvsnap.io/function-version-id" )
Annotation keys NVCA stamps on workload pods to drive the NvSnap integration. NvSnap's mutating webhook reads RestoreFromAnnotation at admission time; Hook B's reconciler (PR-5) watches for pods with CheckpointOnWarmAnnotation to know which to checkpoint after readiness.
const NvSnapCaptureLabel = "nvsnap.io/capture"
NvSnapCaptureLabel marks a pod as a capture candidate.
Unlike the annotations above this has to be a LABEL and it has to be present at ADMISSION: NvSnap's mutating webhook gates the cachedir capture volume (/opt/nvsnap) on it, and a volume cannot be added to a running pod.
nvsnap-server also applies this label once the pod reports warm, which is what drives the agent's rootfs capture watcher. That is correct for the watcher but far too late for the webhook, so cachedir capture was never injected into NVCA-created pods and the agent failed every capture with "cachedir mode: no volume mounted at /opt/nvsnap". Bench manifests never hit this because they carry the label from the start.
Variables ¶
var ( FailedRequestCleanupWindow = 1 * time.Hour CachedRequestCleanupWindow = 1 * time.Hour DefaultTerminationGracePeriodSeconds = 120 // 2 minutes )
var (
ROAccessMode = []v1.PersistentVolumeAccessMode{v1.ReadOnlyMany}
)
skip option only for UT, real cluster detach check is must
var SkippedEventsInSelfDestructMode = map[string]bool{ EventTickRenewICMSCredentials: true, EventTickUpdateHeartbeat: true, EventTickSyncSQSQueue: true, EventTickSyncICMSRequestStatus: true, EventTickAcknowledgeRequest: true, EventTickUpdateICMSRegistration: true, EventTickSyncPeriodicInstanceStatusUpdates: true, }
Functions ¶
func GetDefaultWorkloadContainerNamesToWatch ¶
func GetDefaultWorkloadContainerNamesToWatch() []string
func GetUseUUIDForRequestObjName ¶
func GetUseUUIDForRequestObjName() bool
GetUseUUIDForRequestObjName returns whether request object names use UUID (true) or RequestID (false).
func IsPodAdmissionRejected ¶
func IsPodStuckInitializing ¶
func IsPodStuckInitializing(pod *corev1.Pod, k8sTimeConfig *k8sutil.TimeConfig) (bool, types.ICMSInstanceState)
func NewCommand ¶
func NewK8sComputeBackend ¶
func NewK8sComputeBackend(clients *kubeclients.KubeClients, bk8s *BackendK8sCache) (ICMSRequestHelper, K8sArtifactHelper)
func NewLedgerEventCorrelatorOptions ¶
func NewLedgerEventCorrelatorOptions(heartbeatInterval time.Duration) record.CorrelatorOptions
NewLedgerEventCorrelatorOptions builds client-go Event correlator options so multi-instance ICMSRequest Events do not share one spam budget or collapse into annotation-less aggregates.
For a usable heartbeat interval, aggregation MaxInterval is set just below the periodic status heartbeat (same config source) so each re-report starts a fresh window. When no such window exists (interval <=1s), aggregation is disabled outright via MaxEvents: client-go treats MaxIntervalInSeconds==0 as its 10m default, so a zero here would re-enable aggregation and drop ledger annotations.
func SetUseUUIDForRequestObjName ¶
func SetUseUUIDForRequestObjName(v bool)
SetUseUUIDForRequestObjName sets the naming convention for request objects. Used by tests.
Types ¶
type Agent ¶
type Agent struct {
*AgentOptions
// contains filtered or unexported fields
}
func (*Agent) IsRequestFromClusterQueue ¶
func (a *Agent) IsRequestFromClusterQueue(_ context.Context, req *nvcav2beta1.ICMSRequest) bool
func (*Agent) PostICMSInstanceRequestStatusUpdates ¶
func (*Agent) PutICMSRequestAcknowledgement ¶
func (*Agent) RegisterWithICMS ¶
RegisterWithICMS is re-entrant from the ICMS perspective. BART Agent will Register unconditionally whenever it restarts And store the updated the Credentials in the local Secret
type AgentOptions ¶
type AgentOptions struct {
nvcaconfig.Config
nvcaauth.TokenFetcherOptions
FeatureFlagFetcher featureflag.Fetcher
NCAId string
ClusterName string
ClusterID string
ClusterDescription string
ClusterGroupName string
ClusterGroupID string
ClusterRegion string
ClusterAttributes featureflag.Attributes
// ICMSURL is the ICMS service URL.
ICMSURL string
ICMSHostHeaderOverride string
CloudProvider string
KubeConfigPath string
NVCASvcAddress string
NVCAAdminAddr string
NVCADebugAddr string
SystemNamespace string
RequestsNamespace string
NamespaceLabels labels.Set
ComputeBackend string
CredRenewInterval time.Duration
HeartbeatInterval time.Duration
SyncQueueInterval time.Duration
SyncRequestStatusInterval time.Duration
SyncAcknowledgeRequestInterval time.Duration
PeriodicInstanceStatusInterval time.Duration
// ICMSRequestACKInterval is the interval for ICMS request acknowledgements.
ICMSRequestACKInterval time.Duration
GPUCapacity uint64
K8sVersion string
// ImageCredentialHelperImage is the image tag for "nvcf-image-credential-helper",
// for third party registry cred updates.
ImageCredentialHelperImage string
// K8sTimeConfig configures intervals, timeouts, thresholds for various K8s occurrences.
K8sTimeConfig *k8sutil.TimeConfig
// MinHealthcheckRefreshWait forces the NVCA internal healthchecker
// to wait at least this long between refresh calls,
// in case the healthchecker is too chatty in specific instances.
MinHealthcheckRefreshWait time.Duration
// Feature flags
LogPostingEnabled bool
CachingSupportEnabled bool
NVMeshEncryptionEnabled bool
PeriodicInstanceStatusUpdateEnabled bool
HelmRBACEnforcementEnabled bool
HelmResourceConstraintsEnabled bool
DynamicGPUDiscoveryEnabled bool
MultipleGPUTypesAllowed bool
UniformInstanceLabelsEnabled bool
AutoPurgeDegradedWorkers bool
ClusterTargetingEnabled bool
GXCacheEnabled bool
LowLatencyStreamingEnabled bool
PVCRebindEnabled bool
MultiNodeWorkloadsEnabled bool
// ClientMetricsEnabled turns on the OTel-based semconv metrics for outbound
// dependency clients. When false, the metrics pipeline installs a no-op meter
// provider and instrumented clients emit nothing.
ClientMetricsEnabled bool
// MaintenanceMode indicates the operational mode of NVCA
MaintenanceMode types.MaintenanceMode
// Self-destruct control flags
SkipSelfDestruct bool // Skip self-destruct even if ICMS sends SELF_DESTRUCT
ForceSelfDestruct bool // Force self-destruct mode for testing
// HelmRepository restriction
HelmRepositoryPrefix string
// QueueManager options
EndpointURL string
NATSURL string
NATSHostOverride string
// ReVal service config
HelmReValServiceURL string
HelmReValServiceHostHeaderOverride string
HelmReValStageOAuthTokenURL string
HelmReValStageOAuthPublicKeysetEndpoint string
HelmReValProdOAuthTokenURL string
HelmReValProdOAuthPublicKeysetEndpoint string
// CSIVolumeMountOptions for PVC provisioning
CSIVolumeMountOptions []string
// Function Deployment Stages service config
FunctionDeploymentStagesServiceURL string
FunctionDeploymentStagesStageOAuthTokenURL string
FunctionDeploymentStagesStageOAuthPublicKeysetEndpoint string
FunctionDeploymentStagesProdOAuthTokenURL string
FunctionDeploymentStagesProdOAuthPublicKeysetEndpoint string
// ICMSRequestAckRetryTimeout is the timeout for retrying ICMS request acknowledgements.
ICMSRequestAckRetryTimeout time.Duration
// NVCA Operator version
NVCAOperatorVersion string
// NVCA Agent version
NVCAAgentVersion string
// Secret Mirror
SecretMirrorSourceNamespace string
SecretMirrorLabelSelector string
// Environment variable overrides for workloads
// These are maps of env var key-value pairs applied before translation
FunctionEnvOverrides map[string]string
TaskEnvOverrides map[string]string
// GPU Monitor configuration (used when GracefulNoGPU feature flag is enabled)
// GPUPollInterval is the interval between GPU availability checks
GPUPollInterval time.Duration
// GPUDebounceTime is the time a GPU state must be stable before triggering a change
GPUDebounceTime time.Duration
// MetricsRegisterer allows tests to use a custom prometheus registry
MetricsRegisterer prometheus.Registerer
StartControllerManager func(context.Context, *kubeclients.KubeClients) error
}
func (*AgentOptions) EffectiveICMSRequestACKInterval ¶
func (o *AgentOptions) EffectiveICMSRequestACKInterval() time.Duration
EffectiveICMSRequestACKInterval returns the configured ICMS request acknowledgement interval.
func (*AgentOptions) EffectiveICMSRequestAckRetryTimeout ¶
func (o *AgentOptions) EffectiveICMSRequestAckRetryTimeout() time.Duration
EffectiveICMSRequestAckRetryTimeout returns the configured ICMS request acknowledgement retry timeout.
func (*AgentOptions) EffectiveICMSURL ¶
func (o *AgentOptions) EffectiveICMSURL() string
EffectiveICMSURL returns the ICMS service URL.
func (*AgentOptions) GetOTelAttributes ¶
func (o *AgentOptions) GetOTelAttributes() []otelattr.KeyValue
type AggregatedInstanceStatus ¶
type AggregatedInstanceStatus int
const ( AggregatedInstanceStatusUnknown AggregatedInstanceStatus = iota AggregatedInstanceStatusModelCachingInProgress AggregatedInstanceStatusScheduling AggregatedInstanceStatusPending AggregatedInstanceStatusSucceeded AggregatedInstanceStatusFailed )
func (AggregatedInstanceStatus) String ¶
func (s AggregatedInstanceStatus) String() string
type BackendK8sCache ¶
type BackendK8sCache struct {
// contains filtered or unexported fields
}
BackendK8sCache encapsulates the IO interactions from BART to backend K8s clusters, it contains both query and mutating API calls. BackendK8sCache is immutable after creation, therefore it is thread-safe.
func (*BackendK8sCache) ApplyICMSRequestStatusChange ¶
func (c *BackendK8sCache) ApplyICMSRequestStatusChange(ctx context.Context, req *nvcav2beta1new.ICMSRequest) error
func (*BackendK8sCache) CleanupCachedResources ¶
func (c *BackendK8sCache) CleanupCachedResources(ctx context.Context, allReqs []*nvcav2beta1new.ICMSRequest) error
func (*BackendK8sCache) CleanupCreationRequestResources ¶
func (c *BackendK8sCache) CleanupCreationRequestResources(ctx context.Context, req *nvcav2beta1new.ICMSRequest) error
func (*BackendK8sCache) CleanupUnusedResources ¶
func (c *BackendK8sCache) CleanupUnusedResources(ctx context.Context) error
func (*BackendK8sCache) CreateICMSCreationMessageRequest ¶
func (c *BackendK8sCache) CreateICMSCreationMessageRequest(ctx context.Context, msg translate.CreationQueueMessageMetadataGetter, msgReceipt, msgID, queueURL string, ) (*nvcav2beta1new.ICMSRequest, error)
func (*BackendK8sCache) CreateICMSTerminationMessageRequest ¶
func (c *BackendK8sCache) CreateICMSTerminationMessageRequest(ctx context.Context, tm types.ICMSTerminationMessage, msgReceipt, msgID string) error
func (*BackendK8sCache) EmitICMSEvent ¶
func (c *BackendK8sCache) EmitICMSEvent( req *nvcav2beta1.ICMSRequest, eventType, reason, message string, update *types.ICMSRequestUpdateInfo, )
EmitICMSEvent is EmitICMSEventf without formatting args.
func (*BackendK8sCache) EmitICMSEventf ¶
func (c *BackendK8sCache) EmitICMSEventf( req *nvcav2beta1.ICMSRequest, eventType, reason, msgFmt string, update *types.ICMSRequestUpdateInfo, args ...any, )
EmitICMSEventf emits a Kubernetes Event on the ICMSRequest with FnDs ledger annotations. Pass update=nil for request-level events; pass an update with InstanceID (and status payload fields when applicable) for instance-level events.
func (*BackendK8sCache) FetchICMSRegistrationResponse ¶
func (c *BackendK8sCache) FetchICMSRegistrationResponse(ctx context.Context) (*types.ICMSRegistrationResponse, error)
func (*BackendK8sCache) ForceSync ¶
func (c *BackendK8sCache) ForceSync(ctx context.Context)
func (*BackendK8sCache) GetAllBackendGPUs ¶
func (c *BackendK8sCache) GetAllBackendGPUs(ctx context.Context) ([]types.BackendGPU, error)
func (*BackendK8sCache) GetAllPodsForRequest ¶
func (c *BackendK8sCache) GetAllPodsForRequest(ctx context.Context, reqID string) ([]corev1.Pod, error)
GetAllPodsForRequest returns workload pods for the given logical request ID in the pod instance namespace. It matches the icms-request-id label.
func (*BackendK8sCache) GetBARTRegistrationResponseSecret ¶
func (*BackendK8sCache) GetComponentStatus ¶
func (c *BackendK8sCache) GetComponentStatus(ctx context.Context) (hs types.AgentHealth, err error)
func (*BackendK8sCache) GetGPUResource ¶
func (c *BackendK8sCache) GetGPUResource(ctx context.Context, gpuName types.GPUName) (types.GPUResource, error)
func (*BackendK8sCache) GetNodeFeaturesClient ¶
func (c *BackendK8sCache) GetNodeFeaturesClient() nodefeatures.Client
GetNodeFeaturesClient returns the node features client for GPU monitoring.
func (*BackendK8sCache) GetNodeInformer ¶
func (c *BackendK8sCache) GetNodeInformer() cache.SharedIndexInformer
GetNodeInformer returns the shared node informer for event-driven GPU monitoring.
func (*BackendK8sCache) GetRegisteredBackendGPUs ¶
func (c *BackendK8sCache) GetRegisteredBackendGPUs(ctx context.Context, backendGPUs []types.BackendGPU, multiNodeWorkloadsEnabled bool, ) ([]types.RegistrationGPU, error)
func (*BackendK8sCache) ReconcileInstanceStatus ¶
func (c *BackendK8sCache) ReconcileInstanceStatus(ctx context.Context, is types.ICMSServerInstanceState) ( srToReturn *nvcav2beta1new.ICMSRequest, isToReturn nvcav2beta1new.InstanceStatus, state ICMSInstanceReconcileState, found bool, )
func (*BackendK8sCache) StoreICMSRegistrationResponse ¶
func (c *BackendK8sCache) StoreICMSRegistrationResponse(ctx context.Context, res *types.ICMSRegistrationResponse) error
func (*BackendK8sCache) StoreUpdatedCredentials ¶
func (c *BackendK8sCache) StoreUpdatedCredentials(ctx context.Context, qCreds types.QueueCredentials) error
func (*BackendK8sCache) SyncAllICMSRequests ¶
func (c *BackendK8sCache) SyncAllICMSRequests(ctx context.Context) error
func (*BackendK8sCache) SyncICMSRequest ¶
func (c *BackendK8sCache) SyncICMSRequest(ctx context.Context, nn apitypes.NamespacedName) error
// ICMS request state transition "" (ACK to ICMS) ->
"RequestStatusPending" -> "RequestStatusInProgress" -> "RequestCompleted" / "RequestFailed" -> "RequestCompletionACK / "RequestFailureACK"
func (*BackendK8sCache) SyncICMSRequestByID ¶
func (c *BackendK8sCache) SyncICMSRequestByID(ctx context.Context, icmsReqID, messageBatchID string) error
func (*BackendK8sCache) SyncNetworkPolicies ¶
func (c *BackendK8sCache) SyncNetworkPolicies(ctx context.Context) error
func (*BackendK8sCache) UpdateInstanceTypeMetrics ¶
func (c *BackendK8sCache) UpdateInstanceTypeMetrics( ctx context.Context, gpus []types.RegistrationGPU, gpuNodeClassifications []types.GPUNodeClassification, ) error
func (*BackendK8sCache) UpdateSchedulerWorkloadMetrics ¶
func (c *BackendK8sCache) UpdateSchedulerWorkloadMetrics(ctx context.Context)
UpdateSchedulerWorkloadMetrics recomputes the scheduler workload gauge from live cluster state. It counts active (non-terminal) ICMSRequests grouped by workload kind (function/task) and the actual scheduler observed on their pods. This handles mixed clusters where some workloads were deployed before the KAIScheduler flag was enabled and others after. Because it scans actual pod state, the metric is correct even after NVCA restarts.
type BackendK8sCacheBuilder ¶
type BackendK8sCacheBuilder struct {
*BackendK8sCache
// contains filtered or unexported fields
}
BackendK8sCacheBuilder builds Backendk8sCache and start related edge K8s informers, monitored K8s events are sent to a event channel that is returned by Start()
func NewBackendk8sCacheBuilder ¶
func NewBackendk8sCacheBuilder() *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) Start ¶
func (b *BackendK8sCacheBuilder) Start(ctx context.Context) (*BackendK8sCache, <-chan *core.Event, error)
func (*BackendK8sCacheBuilder) WithCSIVolumeMountOptions ¶
func (b *BackendK8sCacheBuilder) WithCSIVolumeMountOptions(mntOptions []string) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithCachingSupport ¶
func (b *BackendK8sCacheBuilder) WithCachingSupport(enable, encryption bool) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithClients ¶
func (b *BackendK8sCacheBuilder) WithClients(clients *kubeclients.KubeClients) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithClusterName ¶
func (b *BackendK8sCacheBuilder) WithClusterName(clusterName string) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithClusterProvider ¶
func (b *BackendK8sCacheBuilder) WithClusterProvider(cloudProvider string) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithClusterRegion ¶
func (b *BackendK8sCacheBuilder) WithClusterRegion(clusterRegion string) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithComputeBackend ¶
func (b *BackendK8sCacheBuilder) WithComputeBackend(be BackendType) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithConfig ¶
func (b *BackendK8sCacheBuilder) WithConfig(cfg nvcaconfig.Config) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithDynamicNodeFeatureClient ¶
func (b *BackendK8sCacheBuilder) WithDynamicNodeFeatureClient( isDynClient bool, opts nodefeatures.DynamicClientOptions, ) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithEnvOverrides ¶
func (b *BackendK8sCacheBuilder) WithEnvOverrides(functionOverrides, taskOverrides map[string]string) *BackendK8sCacheBuilder
WithEnvOverrides sets the environment variable overrides for function and task workloads
func (*BackendK8sCacheBuilder) WithFNDSClient ¶
func (b *BackendK8sCacheBuilder) WithFNDSClient(client fnds.Client) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithFeatureFlagFetcher ¶
func (b *BackendK8sCacheBuilder) WithFeatureFlagFetcher(fff featureflag.Fetcher) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithHelmInternalPersistentStorage ¶
func (b *BackendK8sCacheBuilder) WithHelmInternalPersistentStorage(enable bool) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithHelmRBACEnforcement ¶
func (b *BackendK8sCacheBuilder) WithHelmRBACEnforcement(enable bool) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithHelmRepositoryPrefix ¶
func (b *BackendK8sCacheBuilder) WithHelmRepositoryPrefix(hrepo string) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithHelmResourceConstraints ¶
func (b *BackendK8sCacheBuilder) WithHelmResourceConstraints(enable bool) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithHelmSharedStorage ¶
func (b *BackendK8sCacheBuilder) WithHelmSharedStorage(enable bool) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithICMSRequestSyncConcurrency ¶
func (b *BackendK8sCacheBuilder) WithICMSRequestSyncConcurrency(c int) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithImageCredentialHelperImage ¶
func (b *BackendK8sCacheBuilder) WithImageCredentialHelperImage(image string) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithInfraOverheadGetter ¶
func (b *BackendK8sCacheBuilder) WithInfraOverheadGetter(infraOverheadGetter enforce.InfraOverheadGetter) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithLogPosting ¶
func (b *BackendK8sCacheBuilder) WithLogPosting(enable bool) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithLowLatencyStreaming ¶
func (b *BackendK8sCacheBuilder) WithLowLatencyStreaming(enable bool) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithNamespaceLabels ¶
func (b *BackendK8sCacheBuilder) WithNamespaceLabels(nsl labels.Set) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithNotFoundInstanceStatusReporter ¶
func (b *BackendK8sCacheBuilder) WithNotFoundInstanceStatusReporter( reporter func(ctx context.Context, reqID, instanceID string) error, ) *BackendK8sCacheBuilder
WithNotFoundInstanceStatusReporter wires the callback used to tell ICMS an instance is terminated immediately when its ICMSRequest CR is deleted locally, instead of waiting for the next periodic instance status sync.
func (*BackendK8sCacheBuilder) WithOTelTracer ¶
func (b *BackendK8sCacheBuilder) WithOTelTracer(tracer oteltrace.Tracer) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithPVCRebind ¶
func (b *BackendK8sCacheBuilder) WithPVCRebind(on bool) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithPeriodicInstanceStatusUpdate ¶
func (b *BackendK8sCacheBuilder) WithPeriodicInstanceStatusUpdate(enable bool, period time.Duration) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithRequestsNamespace ¶
func (b *BackendK8sCacheBuilder) WithRequestsNamespace(requestsNamespace string) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithSecretMirrorConfig ¶
func (b *BackendK8sCacheBuilder) WithSecretMirrorConfig(namespace, selector string) *BackendK8sCacheBuilder
WithSecretMirrorSourceNamespace sets the namespace to source secrets from for mirroring
func (*BackendK8sCacheBuilder) WithStaticGPUCapacity ¶
func (b *BackendK8sCacheBuilder) WithStaticGPUCapacity(gpuCap uint64) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithSystemNamespace ¶
func (b *BackendK8sCacheBuilder) WithSystemNamespace(systemNamespace string) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithTimeConfig ¶
func (b *BackendK8sCacheBuilder) WithTimeConfig(cfg *k8sutil.TimeConfig) *BackendK8sCacheBuilder
func (*BackendK8sCacheBuilder) WithWorkerDegradationHandler ¶
func (b *BackendK8sCacheBuilder) WithWorkerDegradationHandler(enable bool) *BackendK8sCacheBuilder
type BackendType ¶
type BackendType string
const ( BackendTypeAce BackendType = "bcp" BackendTypeK8s BackendType = "k8s" )
type BootstrapTokenString ¶
type CacheAccessObj ¶
type ComputeBackend ¶
type ComputeBackend interface {
// core ICMS request applier
ICMSRequestHelper
// auxiliary K8sArtifactHelper
K8sArtifactHelper
}
TODO: aparthasarat Not All Backend Need the Two, split the interface composition in BackendK8sCache to accommodate that
type GPUMonitor ¶
type GPUMonitor struct {
// contains filtered or unexported fields
}
GPUMonitor continuously monitors GPU availability and notifies on state changes. It implements health.ComponentStatusGetter for readiness checks.
func NewGPUMonitor ¶
func NewGPUMonitor(nfClient nodefeatures.Client, opts ...GPUMonitorOption) *GPUMonitor
NewGPUMonitor creates a new GPU monitor.
func (*GPUMonitor) GetComponentStatus ¶
func (m *GPUMonitor) GetComponentStatus(ctx context.Context) (types.AgentHealth, error)
GetComponentStatus implements health.ComponentStatusGetter. Returns unhealthy status when no GPUs are available.
func (*GPUMonitor) HasGPUs ¶
func (m *GPUMonitor) HasGPUs() bool
HasGPUs returns whether GPUs are currently available.
func (*GPUMonitor) SetHasGPUs ¶
func (m *GPUMonitor) SetHasGPUs(hasGPUs bool)
SetHasGPUs sets the GPU availability state directly. This is primarily used during initial startup to set the state before monitoring begins.
func (*GPUMonitor) SetOnGPUStateChange ¶
func (m *GPUMonitor) SetOnGPUStateChange(cb GPUStateChangeCallback)
SetOnGPUStateChange sets the callback for GPU state changes. This allows setting the callback after construction, which is useful when the callback depends on other components that are created later.
func (*GPUMonitor) Start ¶
func (m *GPUMonitor) Start(ctx context.Context)
Start begins the GPU monitoring loop. It performs an initial check and then polls at the configured interval. If a node informer is configured, node events will also trigger GPU checks.
type GPUMonitorConfig ¶
type GPUMonitorConfig struct {
// PollInterval is the interval between GPU availability checks.
PollInterval time.Duration
// DebounceTime is the time a state must be stable before triggering a change.
DebounceTime time.Duration
}
GPUMonitorConfig holds configuration for the GPU monitor.
type GPUMonitorOption ¶
type GPUMonitorOption func(*GPUMonitor)
GPUMonitorOption is a functional option for configuring the GPU monitor.
func WithGPUDebounceTime ¶
func WithGPUDebounceTime(d time.Duration) GPUMonitorOption
WithGPUDebounceTime sets the debounce time.
func WithGPUPollInterval ¶
func WithGPUPollInterval(d time.Duration) GPUMonitorOption
WithGPUPollInterval sets the polling interval.
func WithGPUStateChangeCallback ¶
func WithGPUStateChangeCallback(cb GPUStateChangeCallback) GPUMonitorOption
WithGPUStateChangeCallback sets the callback for state changes.
func WithNodeInformer ¶
func WithNodeInformer(inf cache.SharedIndexInformer) GPUMonitorOption
WithNodeInformer sets a node informer for event-driven GPU checks. When provided, node add/update/delete events will trigger immediate GPU checks (with coalescing) instead of waiting for the next poll interval.
type GPUStateChangeCallback ¶
GPUStateChangeCallback is called when GPU availability changes. The hasGPUs parameter indicates whether GPUs are now available.
type ICMSClient ¶
type ICMSClient struct {
// contains filtered or unexported fields
}
ICMSClient provides access to ICMS endpoints with auth credential from a tokenFetcher
func NewICMSClient ¶
func NewICMSClient(ctx context.Context, clusterID, endpoint string, tokenFetcher nvcaauth.TokenFetcher, tracer oteltrace.Tracer, httpOpts ...cmnhttp.Option) *ICMSClient
func NewICMSClientWithHostHeaderOverride ¶
func NewICMSClientWithHostHeaderOverride( ctx context.Context, clusterID, endpoint, host string, tokenFetcher nvcaauth.TokenFetcher, tracer oteltrace.Tracer, httpOpts ...cmnhttp.Option, ) *ICMSClient
NewICMSClientWithHostHeaderOverride creates an ICMS client with an optional HTTP Host header override. httpOpts are forwarded to the underlying shared HTTP client (for example to add a metrics transport wrapper).
func (*ICMSClient) Endpoint ¶
func (c *ICMSClient) Endpoint() string
Endpoint returns the ICMS service endpoint URL
func (*ICMSClient) GetCreds ¶
func (c *ICMSClient) GetCreds(ctx context.Context) (*types.ICMSCredentialResponse, error)
func (*ICMSClient) GetICMSServerInstanceStatuses ¶
func (c *ICMSClient) GetICMSServerInstanceStatuses(ctx context.Context) (types.ICMSInstanceStatusResponse, error)
func (*ICMSClient) PostInstanceStatusUpdate ¶
func (c *ICMSClient) PostInstanceStatusUpdate(ctx context.Context, requestID, instanceID string, ireq *types.ICMSInstanceStatusUpdateRequest) error
func (*ICMSClient) PutHealthStatus ¶
func (c *ICMSClient) PutHealthStatus(ctx context.Context, hsr *types.HealthStatusRequest) (*types.HealthStatusResponse, error)
func (*ICMSClient) PutRequestAcknowledgement ¶
func (c *ICMSClient) PutRequestAcknowledgement( ctx context.Context, icmsReqID string, messageBatchID string, instanceCount uint64, srTraceCtxCfg nvcav2beta1.ICMSRequestTraceContextConfig, ) error
func (*ICMSClient) Register ¶
func (c *ICMSClient) Register(ctx context.Context, rreq *types.ICMSRegistrationRequest) (*types.ICMSRegistrationResponse, error)
type ICMSClientInterface ¶
type ICMSClientInterface interface {
PutHealthStatus(ctx context.Context, req *types.HealthStatusRequest) (*types.HealthStatusResponse, error)
Register(ctx context.Context, req *types.ICMSRegistrationRequest) (*types.ICMSRegistrationResponse, error)
PostInstanceStatusUpdate(ctx context.Context, requestID, instanceID string, payload *types.ICMSInstanceStatusUpdateRequest) error
GetICMSServerInstanceStatuses(ctx context.Context) (types.ICMSInstanceStatusResponse, error)
PutRequestAcknowledgement(ctx context.Context, icmsReqID, messageBatchID string, instanceCount uint64, srTraceCtxCfg nvcav2beta1.ICMSRequestTraceContextConfig) error
GetCreds(ctx context.Context) (*types.ICMSCredentialResponse, error)
Endpoint() string
}
ICMSClientInterface defines the interface for ICMS client operations
type ICMSInstanceReconcileState ¶
type ICMSInstanceReconcileState string
const ( ICMSInstanceReconcileNoAction ICMSInstanceReconcileState = "NoAction" ICMSInstanceReconcileUpdateOnly ICMSInstanceReconcileState = "UpdateOnly" ICMSInstanceReconcileTerminateAndUpdate ICMSInstanceReconcileState = "TerminateAndUpdate" )
type ICMSRequestHelper ¶
type ICMSRequestHelper interface {
// ApplyCreationMessage fulfills workload request defined in the ICMS request in the compute backend, also record updated status
ApplyCreationMessage(ctx context.Context, req *nvcav2beta1.ICMSRequest) error
// ApplyTerminationMessage terminates the instances in the request, also record updated status in ICMS request
ApplyTerminationMessage(ctx context.Context, req *nvcav2beta1.ICMSRequest) error
// AggregateInstanceStatuses aggregates the statuses of all Instances in req.
AggregateInstanceStatuses(ctx context.Context, req *nvcav2beta1.ICMSRequest) AggregatedInstanceStatus
// GetICMSRequestStatusUpdatesForRequest returns the consolidate payload of List<ICMSRequestUpdateInfo> to be posted to ICMS for InstanceStatusUpdate
GetICMSRequestStatusUpdatesForRequest(ctx context.Context, req *nvcav2beta1.ICMSRequest) ([]types.ICMSRequestUpdateInfo, error)
// GetICMSRequestUpdatesForTerminationRequest returns payload only for all Termination requests to be posted to ICMS
GetICMSRequestUpdatesForTerminationRequest(ctx context.Context, req *nvcav2beta1.ICMSRequest) []types.ICMSRequestUpdateInfo
// GetICMSRequestUpdatesForCreateRequest returns payload only for all Creation requests to be posted to ICMS
GetICMSRequestUpdatesForCreateRequest(ctx context.Context, req *nvcav2beta1.ICMSRequest) []types.ICMSRequestUpdateInfo
// ComputeCleanupCacheReferences cleanUp underlying cache references
ComputeCleanupCacheReferences(ctx context.Context, references []string) error
// AllInstancesTerminatedAndReported returns if all NVCF worker instances for backend are terminated and reported to ICMS
AllInstancesTerminatedAndReported(ctx context.Context, req *nvcav2beta1.ICMSRequest) bool
// HandleInstanceStatusPreconditionFailure purges the underlying instance and updates the status of ICMSRequest
HandleInstanceStatusPreconditionFailure(ctx context.Context, req *nvcav2beta1.ICMSRequest, instID string) error
// PurgeInstanceID purges a specific instance ID and updates the terminated instances map
PurgeInstanceID(ctx context.Context, req *nvcav2beta1.ICMSRequest, terminatedInstances map[string]nvcav2beta1.InstanceStatus, instanceID string) bool
}
ICMSRequestHelper is an interface for managing ICMS requests in compute backends.
type InitCacheJobState ¶
type InitCacheJobState string
const ( InitCacheJobNotFound InitCacheJobState = "InitCacheJobNotFound" InitCacheJobFailed InitCacheJobState = "InitCacheJobFailed" InitCacheJobInProgress InitCacheJobState = "InitCacheJobInProgress" InitCacheJobCompleted InitCacheJobState = "InitCacheJobCompleted" )
type JWKSUpdater ¶
type JWKSUpdater struct {
// contains filtered or unexported fields
}
JWKSUpdater periodically checks the K8s OIDC JWKS and pushes updates to ICMS.
func NewJWKSUpdater ¶
func NewJWKSUpdater(opts JWKSUpdaterOptions) (*JWKSUpdater, error)
NewJWKSUpdater creates a new JWKSUpdater that pushes JWKS changes to ICMS.
Returns an error when the in-cluster K8s CA cert (k8sCACertPath) is unreadable: every supported deployment target (k3d, kind, EKS, GKE, AKS, kubeadm, OpenShift) projects this file via the standard service-account volume, so a missing CA cert reflects a real misconfiguration that must fail closed rather than silently degrade.
func (*JWKSUpdater) Start ¶
func (u *JWKSUpdater) Start(ctx context.Context) error
Start runs the JWKS updater loop and blocks until ctx is cancelled.
The signature matches controller-runtime's {@code manager.Runnable} interface ({@code Start(ctx context.Context) error}). Today NVCA launches the updater as a bare goroutine (see Agent.Start), matching how other non-controller periodic loops in this package are wired. The Runnable-shaped signature means that if NVCA ever goes multi-replica, this can be registered with the controller-runtime manager and gated on leader election (Runnable.NeedLeaderElection → true) in one call — without changing the updater's own code.
Returns nil on clean shutdown (ctx cancellation); never returns a non-nil error today because the poll loop logs and retries transient push failures itself rather than terminating.
type JWKSUpdaterOptions ¶
type JWKSUpdaterOptions struct {
// ICMSURL is the base URL of the ICMS service to push JWKS to.
ICMSURL string
// ICMSHostHeaderOverride optionally overrides the Host header for ICMS pushes.
ICMSHostHeaderOverride string
// ClusterID identifies this cluster in the ICMS push path.
ClusterID string
// TokenPath is the projected SA token file the agent reads to authenticate
// the ICMS push.
TokenPath string
// TransportWrapper, when set, wraps the ICMS push client's transport as its
// outermost layer. The JWKS push is a fifth ICMS route that does not go
// through the shared retryable client, so it needs the wrapper passed in
// explicitly to appear in the client metrics alongside register, heartbeat,
// credentials and instances.
TransportWrapper func(http.RoundTripper) http.RoundTripper
}
JWKSUpdaterOptions configures a JWKSUpdater.
type K8sArtifactHelper ¶
type K8sArtifactHelper interface {
// AggregatePodInstanceStatus returns an aggregated status reflecting if the podID is Scheduled on a Node
AggregatePodInstanceStatus(ctx context.Context, req *nvcav2beta1.ICMSRequest, podID string) AggregatedInstanceStatus
// GetICMSRequestUpdatesForCreatePodRequest returns payload only for all Pod Creation requests to be posted to ICMS
GetICMSRequestUpdatesForCreatePodRequest(ctx context.Context,
st nvcav2beta1.InstanceStatus, req *nvcav2beta1.ICMSRequest) (types.ICMSRequestUpdateInfo, error)
// CreatePodArtifact creates a PodArtifact as specified by the inputs podArt
CreatePodArtifact(ctx context.Context, podArt function.LaunchArtifact, mf mutateFunc) error
// GetErroredPodLogs returns the pod logs of the failed instance container or the static error message for Utils / Init if one exists
GetErroredPodLogs(ctx context.Context, pod *v1.Pod, prepend string, writeMaxBytes int64) (string, int64, error)
// CreateSecretArtifact creates the k8s Secret as specified by the inputs
CreateSecretArtifact(ctx context.Context, a function.LaunchArtifact, mf mutateFunc) error
// CreateConfigMapArtifact creates the k8s ConfigMap as specified by the inputs
CreateConfigMapArtifact(ctx context.Context, a function.LaunchArtifact, mf mutateFunc) error
// CreatePodArtifactInstances creates the standalone PodInstances
CreatePodArtifactInstances(ctx context.Context, pod *v1.Pod,
req *nvcav2beta1.ICMSRequest, mf mutateFunc) ([]nvcav2beta1.InstanceStatus, error)
}
type K8sComputeBackend ¶
type K8sComputeBackend struct {
// contains filtered or unexported fields
}
func (K8sComputeBackend) AggregateInstanceStatuses ¶
func (c K8sComputeBackend) AggregateInstanceStatuses(ctx context.Context, req *nvcav2beta1.ICMSRequest) AggregatedInstanceStatus
func (K8sComputeBackend) AggregatePodInstanceStatus ¶
func (c K8sComputeBackend) AggregatePodInstanceStatus(ctx context.Context, req *nvcav2beta1.ICMSRequest, instanceID string) AggregatedInstanceStatus
func (K8sComputeBackend) AllInstancesTerminatedAndReported ¶
func (c K8sComputeBackend) AllInstancesTerminatedAndReported(ctx context.Context, req *nvcav2beta1.ICMSRequest) bool
func (K8sComputeBackend) ApplyCreationMessage ¶
func (c K8sComputeBackend) ApplyCreationMessage(ctx context.Context, req *nvcav2beta1.ICMSRequest) error
func (K8sComputeBackend) ApplyTerminationMessage ¶
func (c K8sComputeBackend) ApplyTerminationMessage(ctx context.Context, req *nvcav2beta1.ICMSRequest) error
TODO: Harden logic on instance termination failures
func (K8sComputeBackend) CheckInitCacheJobState ¶
func (c K8sComputeBackend) CheckInitCacheJobState(ctx context.Context, rwPVCName string, job *batchv1.Job) InitCacheJobState
func (K8sComputeBackend) CheckPVCState ¶
func (K8sComputeBackend) CleanupModelCachingResources ¶
func (c K8sComputeBackend) CleanupModelCachingResources(ctx context.Context, rwPVC *v1.PersistentVolumeClaim, initJobName string) error
func (K8sComputeBackend) CleanupModelCachingSetupArtifacts ¶
func (c K8sComputeBackend) CleanupModelCachingSetupArtifacts(ctx context.Context, req *nvcav2beta1.ICMSRequest) error
func (K8sComputeBackend) ComputeCleanupCacheReferences ¶
func (c K8sComputeBackend) ComputeCleanupCacheReferences(ctx context.Context, cacheReferences []string) error
references for K8sComputeBackend are that of PVCNames PVCs are created in the podInstanceNamespace
func (K8sComputeBackend) CreateConfigMapArtifact ¶
func (c K8sComputeBackend) CreateConfigMapArtifact(ctx context.Context, a function.LaunchArtifact, mf mutateFunc) error
func (K8sComputeBackend) CreatePodArtifact ¶
func (c K8sComputeBackend) CreatePodArtifact(ctx context.Context, podArt function.LaunchArtifact, mf mutateFunc) error
func (K8sComputeBackend) CreatePodArtifactInstances ¶
func (c K8sComputeBackend) CreatePodArtifactInstances(ctx context.Context, pod *corev1.Pod, req *nvcav2beta1.ICMSRequest, mf mutateFunc) ([]nvcav2beta1.InstanceStatus, error)
func (K8sComputeBackend) CreateSecretArtifact ¶
func (c K8sComputeBackend) CreateSecretArtifact(ctx context.Context, a function.LaunchArtifact, mf mutateFunc) error
func (K8sComputeBackend) CreateServiceArtifact ¶
func (c K8sComputeBackend) CreateServiceArtifact(ctx context.Context, a function.LaunchArtifact, mf mutateFunc) error
func (K8sComputeBackend) GetErroredPodLogs ¶
func (K8sComputeBackend) GetICMSRequestStatusUpdatesForRequest ¶
func (c K8sComputeBackend) GetICMSRequestStatusUpdatesForRequest(ctx context.Context, req *nvcav2beta1.ICMSRequest) ([]types.ICMSRequestUpdateInfo, error)
func (K8sComputeBackend) GetICMSRequestUpdatesForCreatePodRequest ¶
func (c K8sComputeBackend) GetICMSRequestUpdatesForCreatePodRequest(ctx context.Context, st nvcav2beta1.InstanceStatus, req *nvcav2beta1.ICMSRequest) (types.ICMSRequestUpdateInfo, error)
func (K8sComputeBackend) GetICMSRequestUpdatesForCreateRequest ¶
func (c K8sComputeBackend) GetICMSRequestUpdatesForCreateRequest(ctx context.Context, req *nvcav2beta1.ICMSRequest) []types.ICMSRequestUpdateInfo
func (K8sComputeBackend) GetICMSRequestUpdatesForMiniServiceRequest ¶
func (c K8sComputeBackend) GetICMSRequestUpdatesForMiniServiceRequest(ctx context.Context, req *nvcav2beta1.ICMSRequest, st nvcav2beta1.InstanceStatus, ) (nvcatypes.ICMSRequestUpdateInfo, error)
func (K8sComputeBackend) GetICMSRequestUpdatesForTerminationRequest ¶
func (c K8sComputeBackend) GetICMSRequestUpdatesForTerminationRequest(ctx context.Context, req *nvcav2beta1.ICMSRequest) []types.ICMSRequestUpdateInfo
func (K8sComputeBackend) HandleInstanceStatusPreconditionFailure ¶
func (c K8sComputeBackend) HandleInstanceStatusPreconditionFailure(ctx context.Context, req *nvcav2beta1.ICMSRequest, instID string) error
func (K8sComputeBackend) PurgeInstanceID ¶
func (c K8sComputeBackend) PurgeInstanceID(ctx context.Context, req *nvcav2beta1.ICMSRequest, terminatedInstances map[string]nvcav2beta1.InstanceStatus, instanceID string) bool
PurgeInstanceID implements the ICMSRequestHelper interface by delegating to the private purgeInstanceID method
func (K8sComputeBackend) SetupInitCacheJobBlockDevice ¶
func (c K8sComputeBackend) SetupInitCacheJobBlockDevice(ctx context.Context, rwPVCObj *v1.PersistentVolumeClaim, initJob *batchv1.Job, _ *nvcav2beta1.ICMSRequest) error
func (K8sComputeBackend) SetupModelCachingForRequest ¶
func (c K8sComputeBackend) SetupModelCachingForRequest(ctx context.Context, rwPVC *v1.PersistentVolumeClaim, initJob *batchv1.Job, req *nvcav2beta1.ICMSRequest, mf mutateFunc, ) (ModelCachingState, string)
func (K8sComputeBackend) SetupPVCForReaders ¶
func (c K8sComputeBackend) SetupPVCForReaders(ctx context.Context, rwPVC *v1.PersistentVolumeClaim, initJobName string, req *nvcav2beta1.ICMSRequest, mf mutateFunc) (ROPVCSetupPhase, error)
type ModelCachingState ¶
type ModelCachingState string
const ( ModelCachingInProgress ModelCachingState = "ModelCachingInProgress" ModelCachingFailed ModelCachingState = "ModelCachingFailed" ModelCachingCompleted ModelCachingState = "ModelCachingCompleted" ModelCachingCleanupFailed ModelCachingState = "ModelCachingCleanupFailed" )
type PVCState ¶
type PVCState string
const ( PVCNoState PVCState = "PVCNoState" PVCQueryError PVCState = "PVCStateQueryError" PVCNotFound PVCState = "PVCNotFound" PVCFoundBound PVCState = "PVCFoundBound" PVCFoundUnBound PVCState = "PVCFoundUnBound" PVCUpdateFailed PVCState = "PVCUpdateFailed" PVCFoundBindFailed PVCState = "PVCFoundBindFailed" )
type QueueManager ¶
type QueueManager struct {
Client queue.Client
Backendk8s *BackendK8sCache
// contains filtered or unexported fields
}
func NewQueueManager ¶
func NewQueueManager( bk8s *BackendK8sCache, bsc health.StatusGetter, qc queue.Client, qcreds types.QueueCredentials, fff featureflag.FeatureFlagFetcher, maintenanceMode types.MaintenanceMode, metrics *nvcametrics.Metrics, ) *QueueManager
func (*QueueManager) DeleteCreationMessage ¶
func (qm *QueueManager) DeleteCreationMessage(ctx context.Context, gpuName, rhdl string) error
func (*QueueManager) DeleteCreationMessageV2 ¶
func (qm *QueueManager) DeleteCreationMessageV2(ctx context.Context, rhdl, queueURL string) error
func (*QueueManager) ExtendCreationMessableVisibilityTimeout ¶
func (qm *QueueManager) ExtendCreationMessableVisibilityTimeout(ctx context.Context, gpuName, rhdl string) error
func (*QueueManager) ExtendCreationMessableVisibilityTimeoutV2 ¶
func (qm *QueueManager) ExtendCreationMessableVisibilityTimeoutV2(ctx context.Context, rhdl, queueURL string) error
func (*QueueManager) IsGPUAtCapacity ¶
func (qm *QueueManager) IsGPUAtCapacity(gpuName types.GPUName) bool
IsGPUAtCapacity checks if GPU is at capacity
func (*QueueManager) IsPaused ¶
func (qm *QueueManager) IsPaused() bool
IsPaused returns whether the queue manager is currently paused.
func (*QueueManager) Name ¶
func (qm *QueueManager) Name() string
func (*QueueManager) Pause ¶
func (qm *QueueManager) Pause()
Pause pauses queue processing. Creation messages will not be processed, but termination messages will continue to be processed. Existing in-flight work will continue to completion.
func (*QueueManager) Resume ¶
func (qm *QueueManager) Resume()
Resume resumes queue processing after a pause.
func (*QueueManager) SetGPUAtCapacity ¶
func (qm *QueueManager) SetGPUAtCapacity(gpuName types.GPUName, atCapacity bool) bool
SetGPUAtCapacity returns previous value
func (*QueueManager) SetStatusOK ¶
func (qm *QueueManager) SetStatusOK(ok bool)
func (*QueueManager) StatusOK ¶
func (qm *QueueManager) StatusOK() bool
func (*QueueManager) SyncQueues ¶
func (qm *QueueManager) SyncQueues(ctx context.Context) error
type ROPVCSetupPhase ¶
type ROPVCSetupPhase string
const ( ROPVCSetupQueryFailed ROPVCSetupPhase = "ROPVCSetupQueryFailed" ROPVCSetupInProgress ROPVCSetupPhase = "ROPVCSetupInProgress" ROPVUpdateFailed ROPVCSetupPhase = "ROPVUpdateFailed" ROPVCSetupFailed ROPVCSetupPhase = "ROPVCSetupFailed" ROPVCSetupCompleted ROPVCSetupPhase = "ROPVCSetupCompleted" )
type ValidatorSummaryReconciler ¶
type ValidatorSummaryReconciler struct {
// contains filtered or unexported fields
}
ValidatorSummaryReconciler watches the well-known cluster-validator summary ConfigMap and republishes its content as Prometheus metrics on the agent's long-lived /metrics endpoint. The cluster-validator itself runs as a short-lived init container + CronJob, so it can't serve metrics directly — this reconciler bridges the gap.
func NewValidatorSummaryReconciler ¶
func NewValidatorSummaryReconciler( client kubernetes.Interface, namespace string, m *metrics.Metrics, ) *ValidatorSummaryReconciler
NewValidatorSummaryReconciler constructs the reconciler. The watcher is not started until Start() is called.
func (*ValidatorSummaryReconciler) Start ¶
func (r *ValidatorSummaryReconciler) Start(ctx context.Context) error
Start launches the SharedInformer that watches the summary ConfigMap and publishes its content as metrics. The informer's initial List doubles as a bootstrap — if the ConfigMap already exists at startup (e.g. the agent restarted after a successful validator run), an Add event fires immediately and metrics are populated.
Start blocks only briefly (bounded by cacheSyncTimeout) for the initial cache sync so the agent's /metrics endpoint serves consistent data soon after boot, but never hangs agent startup on a misconfigured watch — metrics are an SLI, not a gate.
Source Files
¶
- agent.go
- agent_http_debug.go
- agent_manager.go
- agent_updates.go
- backendk8scache.go
- backendk8scache_gxcache.go
- bart_setup_prod.go
- cli.go
- computebackend.go
- gpu_registration_manager.go
- gpumonitor.go
- icms_client.go
- jwks_updater.go
- k8scomputebackend.go
- k8scomputebackend_miniservice.go
- k8scomputebackend_modelcache.go
- k8scomputebackend_modelcache_rwx_readonly.go
- k8scomputebackend_task_container.go
- ledger_event_correlator.go
- ledger_events.go
- modelcache_storage_selection.go
- nvsnap_coldstart_gate.go
- nvsnap_coldstart_metrics.go
- nvsnap_controller_start.go
- nvsnap_hook.go
- queue_manager.go
- transport_tls.go
- validator_summary_reconciler.go
- workloadwatcher.go
Directories
¶
| Path | Synopsis |
|---|---|
|
Package nvsnap is the NVCA-side client for the NvSnap checkpoint/restore HTTP API exposed by nvsnap-server.
|
Package nvsnap is the NVCA-side client for the NvSnap checkpoint/restore HTTP API exposed by nvsnap-server. |
|
controller
Package controller wires the post-Ready checkpoint reconciler (in pkg/nvca/nvsnap/reconciler) to a Pod informer + workqueue so the agent can drive it from a goroutine at startup.
|
Package controller wires the post-Ready checkpoint reconciler (in pkg/nvca/nvsnap/reconciler) to a Pod informer + workqueue so the agent can drive it from a goroutine at startup. |
|
reconciler
Package reconciler implements Hook B of the NVCA × NvSnap integration: the post-Ready checkpoint loop.
|
Package reconciler implements Hook B of the NVCA × NvSnap integration: the post-Ready checkpoint loop. |