worker

package
v0.0.0-...-36533f3 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 30, 2026 License: AGPL-3.0 Imports: 105 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

View Source
var ErrEvictionIncomplete = errors.New("evicted containers did not release their resources in time")

ErrEvictionIncomplete is returned when victims chosen for a request were still holding their resources after the drain and kill windows passed.

Functions

func GetPodAddr

func GetPodAddr() (string, error)

GetPodAddr gets the IP from the POD_IP env var. Returns an error if it fails to retrieve an IP.

func GetProcCurrentCPUMillicores

func GetProcCurrentCPUMillicores(cpuTime float64, prevCPUTime float64, systemCPUTime float64, prevSystemCPUTime float64) float64

func GetSystemCPU

func GetSystemCPU() (float64, error)

func IsCRIURestoreError

func IsCRIURestoreError(err error) bool

func IsCheckpointHostIncompatible

func IsCheckpointHostIncompatible(err error) bool

func IsRunscCheckpointVersionMismatch

func IsRunscCheckpointVersionMismatch(err error) bool

func NewBackendRepositoryClient

func NewBackendRepositoryClient(ctx context.Context, config types.AppConfig, token string) (pb.BackendRepositoryServiceClient, error)

NewBackendRepositoryClient creates a new backend repository client

func NewContainerRepositoryClient

func NewContainerRepositoryClient(ctx context.Context, config types.AppConfig, token string) (pb.ContainerRepositoryServiceClient, error)

NewContainerRepositoryClient creates a new container repository client

func NewThunderServiceClient

func NewThunderServiceClient(ctx context.Context, config types.AppConfig, token string) (pb.ThunderServiceClient, error)

NewThunderServiceClient creates a new Thunder service client.

func NewWorkerRepositoryClient

func NewWorkerRepositoryClient(ctx context.Context, config types.AppConfig, token string) (pb.WorkerRepositoryServiceClient, error)

NewWorkerRepositoryClient creates a new worker repository client

func RunCacheServer

func RunCacheServer() error

RunCacheServer runs the worker cache manager without accepting user containers. It uses the same worker image and cache registration path as a normal worker, but never starts the container runtime.

Types

type AssignedGpuDevices

type AssignedGpuDevices struct {
}

type CRIUManager

type CRIUManager interface {
	Available() bool
	CreateCheckpoint(ctx context.Context, runtime runtime.Runtime, checkpointId string, request *types.ContainerRequest, terminateAfterCheckpoint bool) (string, error)
	RestoreCheckpoint(ctx context.Context, runtime runtime.Runtime, opts *RestoreOpts) (int, error)
}

func InitializeCRIUManager

func InitializeCRIUManager(ctx context.Context, config types.CRIUConfig, checkpointRoot string) (CRIUManager, error)

InitializeCRIUManager initializes a new CRIU manager that can be used to checkpoint and restore containers.

func InitializeNvidiaCRIU

func InitializeNvidiaCRIU(ctx context.Context, config types.CRIUConfig, checkpointRoot string) (CRIUManager, error)

type ContainerInstance

type ContainerInstance struct {
	Id                         string
	StubId                     string
	BundlePath                 string
	Overlay                    *common.ContainerOverlay
	Spec                       *specs.Spec
	Err                        error
	ExitCode                   int
	Port                       int
	OutputWriter               *common.OutputWriter
	LogBuffer                  *common.LogBuffer
	Request                    *types.ContainerRequest
	StopReason                 types.StopContainerReason
	RuntimeStarted             bool
	RuntimePid                 int
	RuntimeStartedAt           int64
	SandboxProcessManager      *goproc.GoProcClient
	SandboxProcessManagerReady bool

	DeferredCPUQuota           *specs.LinuxCPU
	CPUSet                     string
	RestoreCPUAffinityDeferred bool
	ProcessManagerReadyOnce    sync.Once
	ProcessManagerReadyChan    chan struct{}
	ContainerIp                string

	ContainerAddressMap map[int32]string
	Runtime             runtime.Runtime

	StopEscalationStarted     atomic.Bool
	StuckMountRecoveryStarted atomic.Bool
	// contains filtered or unexported fields
}

type ContainerLogMessage

type ContainerLogMessage struct {
	// Only SDK envelopes carry beta9_log; arbitrary JSON belongs to the user.
	Internal bool `json:"beta9_log"`
	// SDK stdout and stderr envelopes share the same transport.
	Stream      string                      `json:"stream"`
	Level       string                      `json:"level"`
	Message     string                      `json:"message"`
	TaskID      *string                     `json:"task_id"`
	RunnerEvent *types.ContainerRunnerEvent `json:"beta9_event"`
}

type ContainerLogger

type ContainerLogger struct {
	// contains filtered or unexported fields
}

func (*ContainerLogger) CaptureLogs

func (r *ContainerLogger) CaptureLogs(request *types.ContainerRequest, logChan chan common.LogRecord) error

func (*ContainerLogger) Log

func (r *ContainerLogger) Log(containerId, stubId string, format string, args ...any) error

func (*ContainerLogger) Read

func (r *ContainerLogger) Read(containerId string, buffer []byte) (int64, error)

type ContainerMountManager

type ContainerMountManager struct {
	// contains filtered or unexported fields
}

func NewContainerMountManager

func NewContainerMountManager(config types.AppConfig, poolConfig ...types.WorkerPoolConfig) *ContainerMountManager

func (*ContainerMountManager) RemoveContainerMounts

func (c *ContainerMountManager) RemoveContainerMounts(containerId string)

RemoveContainerMounts removes all mounts for a container

func (*ContainerMountManager) RequiresWorkspaceStorageMount

func (c *ContainerMountManager) RequiresWorkspaceStorageMount(request *types.ContainerRequest) bool

func (*ContainerMountManager) SetupContainerMounts

func (c *ContainerMountManager) SetupContainerMounts(ctx context.Context, request *types.ContainerRequest, outputLogger *slog.Logger) error

SetupContainerMounts initializes any external storage for a container

type ContainerNetwork

type ContainerNetwork interface {
	Setup(containerId string, spec *specs.Spec, request *types.ContainerRequest) error
	TearDown(containerId string) error
	ExposePort(containerId string, hostPort, containerPort int) error
	ExposePorts(containerId string, bindings []PortBinding) error
	ReservePorts(containerId string, count int) ([]int, error)
	ReleasePortReservations(containerId string)
	UpdateNetworkPermissions(containerId string, request *types.ContainerRequest) error
	ContainerPortAddress(containerId string, binding PortBinding) (string, error)
	ContainerPortAddressMap(containerId string, bindings []PortBinding) (map[int32]string, error)
	Close() error
}

type ContainerNetworkManager

type ContainerNetworkManager struct {
	// contains filtered or unexported fields
}

func NewContainerNetworkManager

func NewContainerNetworkManager(ctx context.Context, workerId, poolName string, workerRepoClient pb.WorkerRepositoryServiceClient, containerRepoClient pb.ContainerRepositoryServiceClient, eventRepo repo.EventRepository, config types.AppConfig, containerInstances *common.SafeMap[*ContainerInstance], poolConfig types.WorkerPoolConfig, containerStartLimit int, slotPreparer func(netnsPath string) error) (*ContainerNetworkManager, error)

func (*ContainerNetworkManager) Close

func (m *ContainerNetworkManager) Close() error

func (*ContainerNetworkManager) ExposePort

func (m *ContainerNetworkManager) ExposePort(containerId string, hostPort, containerPort int) error

func (*ContainerNetworkManager) ExposePorts

func (m *ContainerNetworkManager) ExposePorts(containerId string, bindings []PortBinding) error

func (*ContainerNetworkManager) HostPortConflict

func (m *ContainerNetworkManager) HostPortConflict(port int) (bool, error)

HostPortConflict reports whether PREROUTING already owns port. DNAT runs before socket lookup, so net.Listen alone cannot detect this collision.

func (*ContainerNetworkManager) ReleasePortReservations

func (m *ContainerNetworkManager) ReleasePortReservations(containerID string)

func (*ContainerNetworkManager) ReservePorts

func (m *ContainerNetworkManager) ReservePorts(containerID string, count int) ([]int, error)

func (*ContainerNetworkManager) Setup

func (m *ContainerNetworkManager) Setup(containerId string, spec *specs.Spec, request *types.ContainerRequest) error

func (*ContainerNetworkManager) TearDown

func (m *ContainerNetworkManager) TearDown(containerId string) error

func (*ContainerNetworkManager) UpdateNetworkPermissions

func (m *ContainerNetworkManager) UpdateNetworkPermissions(containerId string, request *types.ContainerRequest) error

type ContainerNvidiaManager

type ContainerNvidiaManager struct {
	// contains filtered or unexported fields
}

func (*ContainerNvidiaManager) AssignGPUDevices

func (c *ContainerNvidiaManager) AssignGPUDevices(ctx context.Context, request *types.ContainerRequest) ([]int, error)

func (*ContainerNvidiaManager) CDIDevices

func (c *ContainerNvidiaManager) CDIDevices(assignedDevices []int) []string

CDIDevices returns index-based CDI device names. The generated CDI spec (nvidia-ctk cdi generate) names devices by index, so UUID selectors would be unresolvable; UUIDs are only used for env var pinning.

func (*ContainerNvidiaManager) GetContainerGPUDevices

func (c *ContainerNvidiaManager) GetContainerGPUDevices(containerId string) []int

func (*ContainerNvidiaManager) InjectAssignedEnvVars

func (c *ContainerNvidiaManager) InjectAssignedEnvVars(env []string, assignedDevices []int) []string

func (*ContainerNvidiaManager) InjectEnvVars

func (c *ContainerNvidiaManager) InjectEnvVars(env []string) []string

func (*ContainerNvidiaManager) InjectMounts

func (c *ContainerNvidiaManager) InjectMounts(mounts []specs.Mount) []specs.Mount

func (*ContainerNvidiaManager) UnassignGPUDevices

func (c *ContainerNvidiaManager) UnassignGPUDevices(ctx context.Context, containerId string)

type ContainerOptions

type ContainerOptions struct {
	BundlePath                  string
	HostBindPort                int
	BindPorts                   []int
	StartupPortBindings         []PortBinding
	InitialSpec                 *specs.Spec
	StartupStartedAt            time.Time
	CheckpointFilesystemRestore *checkpointFilesystemRestore
	// AddressesRegistered is the address-map publication started as soon as
	// host ports were reserved; RUNNING waits on it, and so does teardown.
	AddressesRegistered *addressRegistration
}

type ContainerResources

type ContainerResources interface {
	GetCPU(request *types.ContainerRequest) *specs.LinuxCPU
	GetMemory(request *types.ContainerRequest) *specs.LinuxMemory
}

type ContainerRuntimeServer

type ContainerRuntimeServer struct {
	pb.UnimplementedContainerServiceServer
	// contains filtered or unexported fields
}

func NewContainerRuntimeServer

func NewContainerRuntimeServer(opts *ContainerRuntimeServerOpts) (*ContainerRuntimeServer, error)

NewContainerRuntimeServer creates a new runtime-agnostic container server

func (*ContainerRuntimeServer) ContainerArchive

ContainerArchive archives a container's filesystem

func (*ContainerRuntimeServer) ContainerCheckpoint

ContainerCheckpoint creates a checkpoint of a running container

func (*ContainerRuntimeServer) ContainerExec

ContainerExec executes a command inside a running container

func (*ContainerRuntimeServer) ContainerKill

ContainerKill kills and removes a container using the worker's configured runtime

func (*ContainerRuntimeServer) ContainerSandboxExec

func (*ContainerRuntimeServer) ContainerSandboxKill

func (*ContainerRuntimeServer) ContainerSandboxListFiles

func (*ContainerRuntimeServer) ContainerSandboxStatFile

func (*ContainerRuntimeServer) ContainerSandboxStatus

func (*ContainerRuntimeServer) ContainerSandboxStderr

func (*ContainerRuntimeServer) ContainerSandboxStdout

func (*ContainerRuntimeServer) ContainerSnapshotDisks

ContainerSnapshotDisks snapshots a running container's durable disks.

func (*ContainerRuntimeServer) ContainerStatus

ContainerStatus returns the status of a container

func (*ContainerRuntimeServer) ContainerStreamLogs

ContainerStreamLogs streams container logs

func (*ContainerRuntimeServer) ContainerSyncWorkspace

ContainerSyncWorkspace syncs workspace files

func (*ContainerRuntimeServer) Start

func (s *ContainerRuntimeServer) Start() error

func (*ContainerRuntimeServer) Stop

func (s *ContainerRuntimeServer) Stop() error

type ContainerRuntimeServerOpts

type ContainerRuntimeServerOpts struct {
	PodAddr                 string
	Runtime                 runtime.Runtime // The runtime configured for this worker pool
	ContainerInstances      *common.SafeMap[*ContainerInstance]
	ImageClient             *ImageClient
	ContainerRepoClient     pb.ContainerRepositoryServiceClient
	ContainerNetworkManager ContainerNetwork
	EventRepo               repository.EventRepository
	WorkerID                string
	BackendRoute            backendRouteFunc
	CreateCheckpoint        createCheckpointFunc
	SnapshotDisks           snapshotDisksFunc
}

type ContainerThunderManager

type ContainerThunderManager struct {
	// contains filtered or unexported fields
}

func NewContainerThunderManager

func NewContainerThunderManager(client pb.ThunderServiceClient) *ContainerThunderManager

func (*ContainerThunderManager) AssignGPUDevices

func (c *ContainerThunderManager) AssignGPUDevices(ctx context.Context, request *types.ContainerRequest) ([]int, error)

func (*ContainerThunderManager) CDIDevices

func (c *ContainerThunderManager) CDIDevices(assignedDevices []int) []string

func (*ContainerThunderManager) GetContainerGPUDevices

func (c *ContainerThunderManager) GetContainerGPUDevices(containerId string) []int

func (*ContainerThunderManager) InjectAssignedEnvVars

func (c *ContainerThunderManager) InjectAssignedEnvVars(env []string, assignedDevices []int) []string

func (*ContainerThunderManager) InjectEnvVars

func (c *ContainerThunderManager) InjectEnvVars(env []string) []string

func (*ContainerThunderManager) InjectMounts

func (c *ContainerThunderManager) InjectMounts(mounts []specs.Mount) []specs.Mount

func (*ContainerThunderManager) UnassignGPUDevices

func (c *ContainerThunderManager) UnassignGPUDevices(ctx context.Context, containerId string)

type ContainerUsageRecorder

type ContainerUsageRecorder interface {
	RecordContainerUsage(ctx context.Context, request *types.ContainerRequest, start, end time.Time, costCents *float64) error
}

ContainerUsageRecorder reports billable container usage to an external system. Implementations carry their own attribution; the worker supplies the frozen container interval and an optional quoted cost.

type CreateCheckpointOpts

type CreateCheckpointOpts struct {
	Request                  *types.ContainerRequest
	CheckpointId             string
	ContainerIp              string
	OutputLogger             *slog.Logger
	CheckpointPIDChan        chan int
	WaitForSignal            bool
	TerminateAfterCheckpoint bool
	CheckpointRuntime        string
	// DiskSnapshots reports the durable disk snapshots captured alongside the
	// memory image, so callers can pair the checkpoint with exact generations.
	DiskSnapshots []*types.DiskSnapshot
}

type ErrCRIURestoreFailed

type ErrCRIURestoreFailed struct {
	Stderr string
}

func (*ErrCRIURestoreFailed) Error

func (e *ErrCRIURestoreFailed) Error() string

type ErrCheckpointHostIncompatible

type ErrCheckpointHostIncompatible struct {
	Stderr string
}

func (*ErrCheckpointHostIncompatible) Error

type ErrCheckpointRuntimeIncompatible

type ErrCheckpointRuntimeIncompatible struct {
	RuntimeName string
	PayloadName string
	Err         error
}

func (*ErrCheckpointRuntimeIncompatible) Error

func (*ErrCheckpointRuntimeIncompatible) Unwrap

type ErrRunscCheckpointVersionMismatch

type ErrRunscCheckpointVersionMismatch struct {
	Stderr string
}

func (*ErrRunscCheckpointVersionMismatch) Error

type FileCacheManager

type FileCacheManager struct {
	// contains filtered or unexported fields
}

func NewFileCacheManager

func NewFileCacheManager(config types.AppConfig, client *cache.Client) *FileCacheManager

func (*FileCacheManager) CacheAvailable

func (cm *FileCacheManager) CacheAvailable() bool

CacheAvailable checks if the file cache is available

func (*FileCacheManager) CacheFilesInPath

func (cm *FileCacheManager) CacheFilesInPath(sourcePath string)

CacheFilesInPath caches files from a specified source path

func (*FileCacheManager) EnableVolumeCaching

func (cm *FileCacheManager) EnableVolumeCaching(workspaceName string, volumeCacheMap map[string]string, spec *specs.Spec) error

func (*FileCacheManager) GetClient

func (cm *FileCacheManager) GetClient() *cache.Client

GetClient returns the cache client instance.

type FileLock

type FileLock struct {
	// contains filtered or unexported fields
}

func NewFileLock

func NewFileLock(path string) *FileLock

func (*FileLock) Acquire

func (fl *FileLock) Acquire() error

func (*FileLock) Release

func (fl *FileLock) Release() error

type GPUInfoClient

type GPUInfoClient interface {
	AvailableGPUDevices() ([]int, error)
	DeviceUUIDs() (map[int]string, error)
	GetGPUMemoryUsage(deviceIndex int) (GPUMemoryUsageStats, error)
}

type GPUInfoStat

type GPUInfoStat struct {
	MemoryUsed  uint64
	MemoryTotal uint64
}

type GPUManager

type GPUManager interface {
	AssignGPUDevices(ctx context.Context, request *types.ContainerRequest) ([]int, error)
	GetContainerGPUDevices(containerId string) []int
	UnassignGPUDevices(ctx context.Context, containerId string)
	CDIDevices(assignedDevices []int) []string
	InjectEnvVars(env []string) []string
	InjectAssignedEnvVars(env []string, assignedDevices []int) []string
	InjectMounts(mounts []specs.Mount) []specs.Mount
}

func NewContainerNvidiaManager

func NewContainerNvidiaManager(gpuCount uint32, runtimeName string) GPUManager

type GPUMemoryUsageStats

type GPUMemoryUsageStats struct {
	UsedCapacity  int64
	TotalCapacity int64
}

type GvisorResources

type GvisorResources struct {
	*StandardResources
}

func NewGvisorResources

func NewGvisorResources() *GvisorResources

func (*GvisorResources) GetCPU

func (g *GvisorResources) GetCPU(request *types.ContainerRequest) *specs.LinuxCPU

func (*GvisorResources) GetMemory

func (g *GvisorResources) GetMemory(request *types.ContainerRequest) *specs.LinuxMemory

type ImageClient

type ImageClient struct {
	// contains filtered or unexported fields
}

func NewImageClient

func NewImageClient(config types.AppConfig, workerId, workerPoolName string, workerRepoClient pb.WorkerRepositoryServiceClient, fileCacheManager *FileCacheManager) (*ImageClient, error)

func (*ImageClient) Archive

func (c *ImageClient) Archive(ctx context.Context, bundlePath *PathInfo, imageId string, progressChan chan int) error

Generate and upload archived version of the image for distribution

func (*ImageClient) ArchiveLayer

func (c *ImageClient) ArchiveLayer(ctx context.Context, request *types.ContainerRequest, upperDir, imageId string, progress chan<- int) error

ArchiveLayer publishes a running sandbox's filesystem as a new image: the overlay upper directory becomes one OCI layer appended to the image the sandbox started from. The base layers are already in a registry, so the push uploads only the delta, the indexer skips the layers it has indexed before, and the new layer's content is seeded into the content cache while it is indexed. Archiving the merged root shipped the whole image again for every snapshot.

func (*ImageClient) BuildAndArchiveImage

func (c *ImageClient) BuildAndArchiveImage(ctx context.Context, outputLogger *slog.Logger, request *types.ContainerRequest) error

func (*ImageClient) Cleanup

func (c *ImageClient) Cleanup() error

func (*ImageClient) GetCLIPImageMetadata

func (c *ImageClient) GetCLIPImageMetadata(imageId string) (*clipCommon.ImageMetadata, bool)

GetCLIPImageMetadata extracts CLIP image metadata from the archive

func (*ImageClient) GetSourceImageRef

func (c *ImageClient) GetSourceImageRef(imageId string) (string, bool)

GetSourceImageRef retrieves the cached source image reference for a v2 image

func (*ImageClient) PullAndArchiveImage

func (c *ImageClient) PullAndArchiveImage(ctx context.Context, outputLogger *slog.Logger, request *types.ContainerRequest) error

func (*ImageClient) PullLazy

func (c *ImageClient) PullLazy(ctx context.Context, request *types.ContainerRequest) (time.Duration, error)

type NvidiaCRIUManager

type NvidiaCRIUManager struct {
	// contains filtered or unexported fields
}

func (*NvidiaCRIUManager) Available

func (c *NvidiaCRIUManager) Available() bool

func (*NvidiaCRIUManager) CreateCheckpoint

func (c *NvidiaCRIUManager) CreateCheckpoint(ctx context.Context, rt runtime.Runtime, checkpointId string, request *types.ContainerRequest, terminateAfterCheckpoint bool) (string, error)

func (*NvidiaCRIUManager) RestoreCheckpoint

func (c *NvidiaCRIUManager) RestoreCheckpoint(ctx context.Context, rt runtime.Runtime, opts *RestoreOpts) (int, error)

type NvidiaInfoClient

type NvidiaInfoClient struct {
	// contains filtered or unexported fields
}

func (*NvidiaInfoClient) AvailableGPUDevices

func (c *NvidiaInfoClient) AvailableGPUDevices() ([]int, error)

func (*NvidiaInfoClient) DeviceUUIDs

func (c *NvidiaInfoClient) DeviceUUIDs() (map[int]string, error)

DeviceUUIDs returns the index-to-UUID mapping for all devices from a single nvidia-smi query, so callers can resolve many devices without re-querying.

func (*NvidiaInfoClient) GetGPUMemoryUsage

func (c *NvidiaInfoClient) GetGPUMemoryUsage(deviceIndex int) (GPUMemoryUsageStats, error)

GetGpuMemoryUsage retrieves the memory usage of a specific NVIDIA GPU. It returns the total and used memory in bytes.

type PathInfo

type PathInfo struct {
	Path string
	// contains filtered or unexported fields
}

func NewPathInfo

func NewPathInfo(path string) *PathInfo

func (*PathInfo) GetSize

func (p *PathInfo) GetSize() float64

type PortBinding

type PortBinding struct {
	HostPort      int
	ContainerPort int
}

type ProcUtil

type ProcUtil struct {
	procfs.Proc
}

func NewProcUtil

func NewProcUtil(pid int) (*ProcUtil, error)

type ProcessMonitor

type ProcessMonitor struct {
	// contains filtered or unexported fields
}

func NewProcessMonitor

func NewProcessMonitor(pid int, devices []specs.LinuxDeviceCgroup, gpuDeviceIds []int) *ProcessMonitor

func (*ProcessMonitor) GetStatistics

func (m *ProcessMonitor) GetStatistics() (*ProcessStats, error)

func (*ProcessMonitor) Prime

func (m *ProcessMonitor) Prime()

type ProcessStats

type ProcessStats struct {
	CPU    uint64 // in millicores
	Memory process.MemoryInfoStat
	IO     process.IOCountersStat
	NetIO  gopsutilnet.IOCountersStat
	GPU    GPUInfoStat
}

type RestoreOpts

type RestoreOpts struct {
	// contains filtered or unexported fields
}

type RuncResources

type RuncResources struct {
	*StandardResources
}

func NewRuncResources

func NewRuncResources() *RuncResources

type ServiceProxy

type ServiceProxy struct {
	// contains filtered or unexported fields
}

ServiceProxy lets containers reach TCP services (<name>.<tcp externalHost>: databases and TCP pods, never HTTP apps) where that host does not resolve for them, by pinning it in /etc/hosts to a worker-side listener that forwards to the gateway's TCP listener. Bytes are untouched; the gateway terminates TLS and routes by SNI. Pinning is best effort: if the target is unreachable from this worker, containers keep resolving the host over DNS.

func NewServiceProxy

func NewServiceProxy(ctx context.Context, config types.AppConfig) *ServiceProxy

func (*ServiceProxy) Attach

func (p *ServiceProxy) Attach(request *types.ContainerRequest, spec *specs.Spec) error

Attach pins sibling hostnames from the env to the bridge address via /etc/hosts.

func (*ServiceProxy) Stop

func (p *ServiceProxy) Stop()

type StagedFile

type StagedFile struct {
	Path    string
	Content string
}

type StandardResources

type StandardResources struct {
	// contains filtered or unexported fields
}

func NewStandardResources

func NewStandardResources() *StandardResources

func (*StandardResources) GetCPU

func (r *StandardResources) GetCPU(request *types.ContainerRequest) *specs.LinuxCPU

func (*StandardResources) GetMemory

func (r *StandardResources) GetMemory(request *types.ContainerRequest) *specs.LinuxMemory

type Worker

type Worker struct {
	// contains filtered or unexported fields
}

func NewWorker

func NewWorker() (_ *Worker, err error)

func (*Worker) IsCRIUAvailable

func (s *Worker) IsCRIUAvailable(gpuCount uint32) bool

func (*Worker) Run

func (s *Worker) Run() error

func (*Worker) RunContainer

func (s *Worker) RunContainer(ctx context.Context, request *types.ContainerRequest) error

Spawn a single container and stream output to stdout/stderr

type WorkerCacheManager

type WorkerCacheManager struct {
	// contains filtered or unexported fields
}

func NewWorkerCacheManager

func NewWorkerCacheManager(ctx context.Context, config types.AppConfig, poolConfig types.WorkerPoolConfig, workerRepo pb.WorkerRepositoryServiceClient, eventRepo repo.EventRepository, containerInstances *common.SafeMap[*ContainerInstance], workerID, poolName, podAddr string) *WorkerCacheManager

func (*WorkerCacheManager) CheckpointRoot

func (m *WorkerCacheManager) CheckpointRoot() string

func (*WorkerCacheManager) Close

func (m *WorkerCacheManager) Close() error

func (*WorkerCacheManager) ContentReporter

func (m *WorkerCacheManager) ContentReporter() *cacheContentReporter

ContentReporter returns the required-content reporter, or nil when reconciliation is disabled or unavailable. Callers must tolerate nil.

func (*WorkerCacheManager) Drain

func (m *WorkerCacheManager) Drain() error

func (*WorkerCacheManager) Start

func (m *WorkerCacheManager) Start() (*cache.Client, error)

type WorkerUsageMetrics

type WorkerUsageMetrics struct {
	// contains filtered or unexported fields
}

func NewWorkerUsageMetrics

func NewWorkerUsageMetrics(
	ctx context.Context,
	workerId string,
	config types.AppConfig,
	gpuType string,
	poolMode types.PoolMode,
	usageRecorder ContainerUsageRecorder,
) (*WorkerUsageMetrics, error)

func (*WorkerUsageMetrics) EmitContainerUsage

func (wm *WorkerUsageMetrics) EmitContainerUsage(ctx context.Context, request *types.ContainerRequest)

EmitContainerUsage binds a quote to each interval start. Once an end is captured, a bounded quote refresh cannot change the accounted duration. Any effective-date boundary inside [start,end) becomes a separate segment.

type WorkspaceStorageManager

type WorkspaceStorageManager struct {
	// contains filtered or unexported fields
}

func NewWorkspaceStorageManager

func NewWorkspaceStorageManager(ctx context.Context, config types.StorageConfig, poolConfig types.WorkerPoolConfig, containerInstances *common.SafeMap[*ContainerInstance], cacheClient *cache.Client) (*WorkspaceStorageManager, error)

func (*WorkspaceStorageManager) Cleanup

func (sm *WorkspaceStorageManager) Cleanup() error

func (*WorkspaceStorageManager) Create

func (sm *WorkspaceStorageManager) Create(workspaceName string, storage storage.Storage)

func (*WorkspaceStorageManager) Mount

func (sm *WorkspaceStorageManager) Mount(workspaceName string, workspaceStorage *types.WorkspaceStorage) (storage.Storage, error)

func (*WorkspaceStorageManager) Unmount

func (sm *WorkspaceStorageManager) Unmount(workspaceName string) error

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL