Documentation
¶
Index ¶
- Variables
- func GetPodAddr() (string, error)
- func GetProcCurrentCPUMillicores(cpuTime float64, prevCPUTime float64, systemCPUTime float64, ...) float64
- func GetSystemCPU() (float64, error)
- func IsCRIURestoreError(err error) bool
- func IsCheckpointHostIncompatible(err error) bool
- func IsRunscCheckpointVersionMismatch(err error) bool
- func NewBackendRepositoryClient(ctx context.Context, config types.AppConfig, token string) (pb.BackendRepositoryServiceClient, error)
- func NewContainerRepositoryClient(ctx context.Context, config types.AppConfig, token string) (pb.ContainerRepositoryServiceClient, error)
- func NewThunderServiceClient(ctx context.Context, config types.AppConfig, token string) (pb.ThunderServiceClient, error)
- func NewWorkerRepositoryClient(ctx context.Context, config types.AppConfig, token string) (pb.WorkerRepositoryServiceClient, error)
- func RunCacheServer() error
- type AssignedGpuDevices
- type CRIUManager
- type ContainerInstance
- type ContainerLogMessage
- type ContainerLogger
- type ContainerMountManager
- type ContainerNetwork
- type ContainerNetworkManager
- func (m *ContainerNetworkManager) Close() error
- func (m *ContainerNetworkManager) ExposePort(containerId string, hostPort, containerPort int) error
- func (m *ContainerNetworkManager) ExposePorts(containerId string, bindings []PortBinding) error
- func (m *ContainerNetworkManager) HostPortConflict(port int) (bool, error)
- func (m *ContainerNetworkManager) ReleasePortReservations(containerID string)
- func (m *ContainerNetworkManager) ReservePorts(containerID string, count int) ([]int, error)
- func (m *ContainerNetworkManager) Setup(containerId string, spec *specs.Spec, request *types.ContainerRequest) error
- func (m *ContainerNetworkManager) TearDown(containerId string) error
- func (m *ContainerNetworkManager) UpdateNetworkPermissions(containerId string, request *types.ContainerRequest) error
- type ContainerNvidiaManager
- func (c *ContainerNvidiaManager) AssignGPUDevices(ctx context.Context, request *types.ContainerRequest) ([]int, error)
- func (c *ContainerNvidiaManager) CDIDevices(assignedDevices []int) []string
- func (c *ContainerNvidiaManager) GetContainerGPUDevices(containerId string) []int
- func (c *ContainerNvidiaManager) InjectAssignedEnvVars(env []string, assignedDevices []int) []string
- func (c *ContainerNvidiaManager) InjectEnvVars(env []string) []string
- func (c *ContainerNvidiaManager) InjectMounts(mounts []specs.Mount) []specs.Mount
- func (c *ContainerNvidiaManager) UnassignGPUDevices(ctx context.Context, containerId string)
- type ContainerOptions
- type ContainerResources
- type ContainerRuntimeServer
- func (s *ContainerRuntimeServer) ContainerArchive(req *pb.ContainerArchiveRequest, ...) error
- func (s *ContainerRuntimeServer) ContainerCheckpoint(ctx context.Context, in *pb.ContainerCheckpointRequest) (*pb.ContainerCheckpointResponse, error)
- func (s *ContainerRuntimeServer) ContainerExec(ctx context.Context, in *pb.ContainerExecRequest) (*pb.ContainerExecResponse, error)
- func (s *ContainerRuntimeServer) ContainerKill(ctx context.Context, in *pb.ContainerKillRequest) (*pb.ContainerKillResponse, error)
- func (s *ContainerRuntimeServer) ContainerSandboxCreateDirectory(ctx context.Context, in *pb.ContainerSandboxCreateDirectoryRequest) (*pb.ContainerSandboxCreateDirectoryResponse, error)
- func (s *ContainerRuntimeServer) ContainerSandboxDeleteDirectory(ctx context.Context, in *pb.ContainerSandboxDeleteDirectoryRequest) (*pb.ContainerSandboxDeleteDirectoryResponse, error)
- func (s *ContainerRuntimeServer) ContainerSandboxDeleteFile(ctx context.Context, in *pb.ContainerSandboxDeleteFileRequest) (*pb.ContainerSandboxDeleteFileResponse, error)
- func (s *ContainerRuntimeServer) ContainerSandboxDownloadFile(ctx context.Context, in *pb.ContainerSandboxDownloadFileRequest) (*pb.ContainerSandboxDownloadFileResponse, error)
- func (s *ContainerRuntimeServer) ContainerSandboxExec(ctx context.Context, in *pb.ContainerSandboxExecRequest) (*pb.ContainerSandboxExecResponse, error)
- func (s *ContainerRuntimeServer) ContainerSandboxExposePort(ctx context.Context, in *pb.ContainerSandboxExposePortRequest) (*pb.ContainerSandboxExposePortResponse, error)
- func (s *ContainerRuntimeServer) ContainerSandboxFindInFiles(ctx context.Context, in *pb.ContainerSandboxFindInFilesRequest) (*pb.ContainerSandboxFindInFilesResponse, error)
- func (s *ContainerRuntimeServer) ContainerSandboxKill(ctx context.Context, in *pb.ContainerSandboxKillRequest) (*pb.ContainerSandboxKillResponse, error)
- func (s *ContainerRuntimeServer) ContainerSandboxListExposedPorts(ctx context.Context, in *pb.ContainerSandboxListExposedPortsRequest) (*pb.ContainerSandboxListExposedPortsResponse, error)
- func (s *ContainerRuntimeServer) ContainerSandboxListFiles(ctx context.Context, in *pb.ContainerSandboxListFilesRequest) (*pb.ContainerSandboxListFilesResponse, error)
- func (s *ContainerRuntimeServer) ContainerSandboxListProcesses(ctx context.Context, in *pb.ContainerSandboxListProcessesRequest) (*pb.ContainerSandboxListProcessesResponse, error)
- func (s *ContainerRuntimeServer) ContainerSandboxReplaceInFiles(ctx context.Context, in *pb.ContainerSandboxReplaceInFilesRequest) (*pb.ContainerSandboxReplaceInFilesResponse, error)
- func (s *ContainerRuntimeServer) ContainerSandboxStatFile(ctx context.Context, in *pb.ContainerSandboxStatFileRequest) (*pb.ContainerSandboxStatFileResponse, error)
- func (s *ContainerRuntimeServer) ContainerSandboxStatus(ctx context.Context, in *pb.ContainerSandboxStatusRequest) (*pb.ContainerSandboxStatusResponse, error)
- func (s *ContainerRuntimeServer) ContainerSandboxStderr(ctx context.Context, in *pb.ContainerSandboxStderrRequest) (*pb.ContainerSandboxStderrResponse, error)
- func (s *ContainerRuntimeServer) ContainerSandboxStdout(ctx context.Context, in *pb.ContainerSandboxStdoutRequest) (*pb.ContainerSandboxStdoutResponse, error)
- func (s *ContainerRuntimeServer) ContainerSandboxUpdateNetworkPermissions(ctx context.Context, in *pb.ContainerSandboxUpdateNetworkPermissionsRequest) (*pb.ContainerSandboxUpdateNetworkPermissionsResponse, error)
- func (s *ContainerRuntimeServer) ContainerSandboxUploadFile(ctx context.Context, in *pb.ContainerSandboxUploadFileRequest) (*pb.ContainerSandboxUploadFileResponse, error)
- func (s *ContainerRuntimeServer) ContainerSnapshotDisks(ctx context.Context, in *pb.ContainerSnapshotDisksRequest) (*pb.ContainerSnapshotDisksResponse, error)
- func (s *ContainerRuntimeServer) ContainerStatus(ctx context.Context, in *pb.ContainerStatusRequest) (*pb.ContainerStatusResponse, error)
- func (s *ContainerRuntimeServer) ContainerStreamLogs(req *pb.ContainerStreamLogsRequest, ...) error
- func (s *ContainerRuntimeServer) ContainerSyncWorkspace(ctx context.Context, in *pb.SyncContainerWorkspaceRequest) (*pb.SyncContainerWorkspaceResponse, error)
- func (s *ContainerRuntimeServer) Start() error
- func (s *ContainerRuntimeServer) Stop() error
- type ContainerRuntimeServerOpts
- type ContainerThunderManager
- func (c *ContainerThunderManager) AssignGPUDevices(ctx context.Context, request *types.ContainerRequest) ([]int, error)
- func (c *ContainerThunderManager) CDIDevices(assignedDevices []int) []string
- func (c *ContainerThunderManager) GetContainerGPUDevices(containerId string) []int
- func (c *ContainerThunderManager) InjectAssignedEnvVars(env []string, assignedDevices []int) []string
- func (c *ContainerThunderManager) InjectEnvVars(env []string) []string
- func (c *ContainerThunderManager) InjectMounts(mounts []specs.Mount) []specs.Mount
- func (c *ContainerThunderManager) UnassignGPUDevices(ctx context.Context, containerId string)
- type ContainerUsageRecorder
- type CreateCheckpointOpts
- type ErrCRIURestoreFailed
- type ErrCheckpointHostIncompatible
- type ErrCheckpointRuntimeIncompatible
- type ErrRunscCheckpointVersionMismatch
- type FileCacheManager
- type FileLock
- type GPUInfoClient
- type GPUInfoStat
- type GPUManager
- type GPUMemoryUsageStats
- type GvisorResources
- type ImageClient
- func (c *ImageClient) Archive(ctx context.Context, bundlePath *PathInfo, imageId string, ...) error
- func (c *ImageClient) ArchiveLayer(ctx context.Context, request *types.ContainerRequest, upperDir, imageId string, ...) error
- func (c *ImageClient) BuildAndArchiveImage(ctx context.Context, outputLogger *slog.Logger, ...) error
- func (c *ImageClient) Cleanup() error
- func (c *ImageClient) GetCLIPImageMetadata(imageId string) (*clipCommon.ImageMetadata, bool)
- func (c *ImageClient) GetSourceImageRef(imageId string) (string, bool)
- func (c *ImageClient) PullAndArchiveImage(ctx context.Context, outputLogger *slog.Logger, ...) error
- func (c *ImageClient) PullLazy(ctx context.Context, request *types.ContainerRequest) (time.Duration, error)
- type NvidiaCRIUManager
- type NvidiaInfoClient
- type PathInfo
- type PortBinding
- type ProcUtil
- type ProcessMonitor
- type ProcessStats
- type RestoreOpts
- type RuncResources
- type ServiceProxy
- type StagedFile
- type StandardResources
- type Worker
- type WorkerCacheManager
- type WorkerUsageMetrics
- type WorkspaceStorageManager
- func (sm *WorkspaceStorageManager) Cleanup() error
- func (sm *WorkspaceStorageManager) Create(workspaceName string, storage storage.Storage)
- func (sm *WorkspaceStorageManager) Mount(workspaceName string, workspaceStorage *types.WorkspaceStorage) (storage.Storage, error)
- func (sm *WorkspaceStorageManager) Unmount(workspaceName string) error
Constants ¶
This section is empty.
Variables ¶
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 ¶
GetPodAddr gets the IP from the POD_IP env var. Returns an error if it fails to retrieve an IP.
func GetSystemCPU ¶
func IsCRIURestoreError ¶
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
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 ¶
func (s *ContainerRuntimeServer) ContainerArchive(req *pb.ContainerArchiveRequest, stream pb.ContainerService_ContainerArchiveServer) error
ContainerArchive archives a container's filesystem
func (*ContainerRuntimeServer) ContainerCheckpoint ¶
func (s *ContainerRuntimeServer) ContainerCheckpoint(ctx context.Context, in *pb.ContainerCheckpointRequest) (*pb.ContainerCheckpointResponse, error)
ContainerCheckpoint creates a checkpoint of a running container
func (*ContainerRuntimeServer) ContainerExec ¶
func (s *ContainerRuntimeServer) ContainerExec(ctx context.Context, in *pb.ContainerExecRequest) (*pb.ContainerExecResponse, error)
ContainerExec executes a command inside a running container
func (*ContainerRuntimeServer) ContainerKill ¶
func (s *ContainerRuntimeServer) ContainerKill(ctx context.Context, in *pb.ContainerKillRequest) (*pb.ContainerKillResponse, error)
ContainerKill kills and removes a container using the worker's configured runtime
func (*ContainerRuntimeServer) ContainerSandboxCreateDirectory ¶
func (s *ContainerRuntimeServer) ContainerSandboxCreateDirectory(ctx context.Context, in *pb.ContainerSandboxCreateDirectoryRequest) (*pb.ContainerSandboxCreateDirectoryResponse, error)
func (*ContainerRuntimeServer) ContainerSandboxDeleteDirectory ¶
func (s *ContainerRuntimeServer) ContainerSandboxDeleteDirectory(ctx context.Context, in *pb.ContainerSandboxDeleteDirectoryRequest) (*pb.ContainerSandboxDeleteDirectoryResponse, error)
func (*ContainerRuntimeServer) ContainerSandboxDeleteFile ¶
func (s *ContainerRuntimeServer) ContainerSandboxDeleteFile(ctx context.Context, in *pb.ContainerSandboxDeleteFileRequest) (*pb.ContainerSandboxDeleteFileResponse, error)
func (*ContainerRuntimeServer) ContainerSandboxDownloadFile ¶
func (s *ContainerRuntimeServer) ContainerSandboxDownloadFile(ctx context.Context, in *pb.ContainerSandboxDownloadFileRequest) (*pb.ContainerSandboxDownloadFileResponse, error)
func (*ContainerRuntimeServer) ContainerSandboxExec ¶
func (s *ContainerRuntimeServer) ContainerSandboxExec(ctx context.Context, in *pb.ContainerSandboxExecRequest) (*pb.ContainerSandboxExecResponse, error)
func (*ContainerRuntimeServer) ContainerSandboxExposePort ¶
func (s *ContainerRuntimeServer) ContainerSandboxExposePort(ctx context.Context, in *pb.ContainerSandboxExposePortRequest) (*pb.ContainerSandboxExposePortResponse, error)
func (*ContainerRuntimeServer) ContainerSandboxFindInFiles ¶
func (s *ContainerRuntimeServer) ContainerSandboxFindInFiles(ctx context.Context, in *pb.ContainerSandboxFindInFilesRequest) (*pb.ContainerSandboxFindInFilesResponse, error)
func (*ContainerRuntimeServer) ContainerSandboxKill ¶
func (s *ContainerRuntimeServer) ContainerSandboxKill(ctx context.Context, in *pb.ContainerSandboxKillRequest) (*pb.ContainerSandboxKillResponse, error)
func (*ContainerRuntimeServer) ContainerSandboxListExposedPorts ¶
func (s *ContainerRuntimeServer) ContainerSandboxListExposedPorts(ctx context.Context, in *pb.ContainerSandboxListExposedPortsRequest) (*pb.ContainerSandboxListExposedPortsResponse, error)
func (*ContainerRuntimeServer) ContainerSandboxListFiles ¶
func (s *ContainerRuntimeServer) ContainerSandboxListFiles(ctx context.Context, in *pb.ContainerSandboxListFilesRequest) (*pb.ContainerSandboxListFilesResponse, error)
func (*ContainerRuntimeServer) ContainerSandboxListProcesses ¶
func (s *ContainerRuntimeServer) ContainerSandboxListProcesses(ctx context.Context, in *pb.ContainerSandboxListProcessesRequest) (*pb.ContainerSandboxListProcessesResponse, error)
func (*ContainerRuntimeServer) ContainerSandboxReplaceInFiles ¶
func (s *ContainerRuntimeServer) ContainerSandboxReplaceInFiles(ctx context.Context, in *pb.ContainerSandboxReplaceInFilesRequest) (*pb.ContainerSandboxReplaceInFilesResponse, error)
func (*ContainerRuntimeServer) ContainerSandboxStatFile ¶
func (s *ContainerRuntimeServer) ContainerSandboxStatFile(ctx context.Context, in *pb.ContainerSandboxStatFileRequest) (*pb.ContainerSandboxStatFileResponse, error)
func (*ContainerRuntimeServer) ContainerSandboxStatus ¶
func (s *ContainerRuntimeServer) ContainerSandboxStatus(ctx context.Context, in *pb.ContainerSandboxStatusRequest) (*pb.ContainerSandboxStatusResponse, error)
func (*ContainerRuntimeServer) ContainerSandboxStderr ¶
func (s *ContainerRuntimeServer) ContainerSandboxStderr(ctx context.Context, in *pb.ContainerSandboxStderrRequest) (*pb.ContainerSandboxStderrResponse, error)
func (*ContainerRuntimeServer) ContainerSandboxStdout ¶
func (s *ContainerRuntimeServer) ContainerSandboxStdout(ctx context.Context, in *pb.ContainerSandboxStdoutRequest) (*pb.ContainerSandboxStdoutResponse, error)
func (*ContainerRuntimeServer) ContainerSandboxUpdateNetworkPermissions ¶
func (s *ContainerRuntimeServer) ContainerSandboxUpdateNetworkPermissions(ctx context.Context, in *pb.ContainerSandboxUpdateNetworkPermissionsRequest) (*pb.ContainerSandboxUpdateNetworkPermissionsResponse, error)
func (*ContainerRuntimeServer) ContainerSandboxUploadFile ¶
func (s *ContainerRuntimeServer) ContainerSandboxUploadFile(ctx context.Context, in *pb.ContainerSandboxUploadFileRequest) (*pb.ContainerSandboxUploadFileResponse, error)
func (*ContainerRuntimeServer) ContainerSnapshotDisks ¶
func (s *ContainerRuntimeServer) ContainerSnapshotDisks(ctx context.Context, in *pb.ContainerSnapshotDisksRequest) (*pb.ContainerSnapshotDisksResponse, error)
ContainerSnapshotDisks snapshots a running container's durable disks.
func (*ContainerRuntimeServer) ContainerStatus ¶
func (s *ContainerRuntimeServer) ContainerStatus(ctx context.Context, in *pb.ContainerStatusRequest) (*pb.ContainerStatusResponse, error)
ContainerStatus returns the status of a container
func (*ContainerRuntimeServer) ContainerStreamLogs ¶
func (s *ContainerRuntimeServer) ContainerStreamLogs(req *pb.ContainerStreamLogsRequest, stream pb.ContainerService_ContainerStreamLogsServer) error
ContainerStreamLogs streams container logs
func (*ContainerRuntimeServer) ContainerSyncWorkspace ¶
func (s *ContainerRuntimeServer) ContainerSyncWorkspace(ctx context.Context, in *pb.SyncContainerWorkspaceRequest) (*pb.SyncContainerWorkspaceResponse, error)
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 ¶
func (e *ErrCheckpointHostIncompatible) Error() string
type ErrCheckpointRuntimeIncompatible ¶
func (*ErrCheckpointRuntimeIncompatible) Error ¶
func (e *ErrCheckpointRuntimeIncompatible) Error() string
func (*ErrCheckpointRuntimeIncompatible) Unwrap ¶
func (e *ErrCheckpointRuntimeIncompatible) Unwrap() error
type ErrRunscCheckpointVersionMismatch ¶
type ErrRunscCheckpointVersionMismatch struct {
Stderr string
}
func (*ErrRunscCheckpointVersionMismatch) Error ¶
func (e *ErrRunscCheckpointVersionMismatch) Error() string
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 (*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 ¶
type GPUInfoClient ¶
type GPUInfoStat ¶
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 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 (*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 ¶
type PortBinding ¶
type ProcUtil ¶
func NewProcUtil ¶
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 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 (*Worker) IsCRIUAvailable ¶
func (*Worker) RunContainer ¶
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
type WorkerUsageMetrics ¶
type WorkerUsageMetrics struct {
// contains filtered or unexported fields
}
func NewWorkerUsageMetrics ¶
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
Source Files
¶
- cache_coordinator.go
- cache_manager.go
- cache_metadata_store.go
- cache_reconciler.go
- cache_server.go
- container_config.go
- container_environment.go
- container_events.go
- container_network.go
- container_port_proxy.go
- container_resources.go
- container_routes.go
- container_server.go
- container_server_guestfs.go
- criu.go
- criu_nvidia.go
- dockerfile_plan.go
- durable_disk.go
- durable_disk_lifecycle.go
- durable_disk_qcow.go
- durable_disk_snapshot.go
- events.go
- evict.go
- file_cache.go
- git_build.go
- git_dockerfile.go
- gpu_info.go
- guest_layer.go
- image.go
- image_base_cache.go
- image_build_apt.go
- image_build_cache.go
- image_build_layered.go
- image_cache.go
- image_credentials.go
- image_events.go
- image_layers.go
- image_snapshot.go
- lifecycle.go
- logger.go
- mount.go
- network.go
- network_util.go
- nvidia.go
- oom.go
- procutil.go
- repository.go
- resources.go
- runtime_credentials.go
- sandbox.go
- service_proxy.go
- storage_manager.go
- sysfs_linux.go
- thunder.go
- usage.go
- util.go
- utils.go
- worker.go
- worker_events.go
- worker_mode.go