runner

package
v0.15.0 Latest Latest
Warning

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

Go to latest
Published: Sep 9, 2026 License: Apache-2.0 Imports: 51 Imported by: 0

Documentation

Index

Constants

View Source
const (
	DefaulWorkers = 3
)

Variables

This section is empty.

Functions

This section is empty.

Types

type CapabilityBoot added in v0.15.0

type CapabilityBoot[T any] struct {
	Component *boot.Component
	Output    boot.Output[T]
}

CapabilityBoot owns one independently started runner capability.

func NewNodePresenceBoot added in v0.15.0

func NewNodePresenceBoot(host boot.Output[*SandboxHost], storageAgent, sandboxAgent *boot.Component, stopTimeout time.Duration) *CapabilityBoot[*NodePresence]

NewNodePresenceBoot publishes node readiness only after both agents are active.

func NewSandboxAgentBoot added in v0.15.0

func NewSandboxAgentBoot(host boot.Output[*SandboxHost], stopTimeout time.Duration, dependencies ...*boot.Component) *CapabilityBoot[*SandboxAgent]

NewSandboxAgentBoot starts sandbox reconciliation after the restored host and any composition-specific order-only dependencies are ready.

func NewStorageAgentBoot added in v0.15.0

func NewStorageAgentBoot(storage boot.Output[*NodeStorage], host *boot.Component, stopTimeout time.Duration) *CapabilityBoot[*StorageAgent]

NewStorageAgentBoot starts storage reconciliation after storage and the sandbox host are both ready.

type ClusterAccess added in v0.15.0

type ClusterAccess struct {
	RunnerConfig
	Log *slog.Logger
	// contains filtered or unexported fields
}

ClusterAccess owns the RPC state and cluster-backed capabilities used by a runner. It has no container or host-network responsibilities.

func NewClusterAccess added in v0.15.0

func NewClusterAccess(log *slog.Logger, deps RunnerDeps, cfg RunnerConfig) (*ClusterAccess, error)

NewClusterAccess constructs the runner's connection to cluster-owned state and capabilities.

func (*ClusterAccess) Close added in v0.15.0

func (r *ClusterAccess) Close() error

func (*ClusterAccess) Start added in v0.15.0

func (r *ClusterAccess) Start(ctx context.Context) (retErr error)

func (*ClusterAccess) WorkloadIssuer added in v0.15.0

func (r *ClusterAccess) WorkloadIssuer() workloadidentity.TokenIssuer

WorkloadIssuer returns the issuer this runner mints identity tokens through, or nil if none is available.

It is only meaningful after Start, which is where a distributed runner acquires its issuer from the coordinator. Callers that build something needing tokens before then should hold a source they can arm afterwards rather than reading this early and caching a nil.

type NodePresence added in v0.15.0

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

NodePresence owns the session-scoped node registration. The graph starts it only after the node's storage and sandbox agents are running.

func NewNodePresence added in v0.15.0

func NewNodePresence(host *SandboxHost) *NodePresence

NewNodePresence constructs the node registration boundary for a restored host.

func (*NodePresence) Close added in v0.15.0

func (r *NodePresence) Close() error

func (*NodePresence) Drain added in v0.15.0

func (r *NodePresence) Drain(ctx context.Context) error

Drain sets the runner's node status to disabled and stops all running sandboxes

func (*NodePresence) Start added in v0.15.0

func (r *NodePresence) Start(ctx context.Context) error

Start advertises the node only after the graph has started its storage and sandbox agents.

type NodeStorage added in v0.15.0

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

NodeStorage owns local disk recovery and the controllers that maintain disk state for one runner. Start recovers the durable state but deliberately does not begin watching desired state; StorageAgent owns that later transition.

func NewNodeStorage added in v0.15.0

func NewNodeStorage(access *ClusterAccess, deps RunnerDeps, config RunnerConfig) (*NodeStorage, error)

NewNodeStorage constructs durable node storage on top of cluster access.

func (*NodeStorage) Close added in v0.15.0

func (s *NodeStorage) Close() error

func (*NodeStorage) SetRestartMode added in v0.15.0

func (s *NodeStorage) SetRestartMode(v bool)

func (*NodeStorage) Start added in v0.15.0

func (s *NodeStorage) Start(ctx context.Context) error

type Runner

type Runner struct {
	Access       *ClusterAccess
	Storage      *NodeStorage
	Host         *SandboxHost
	StorageAgent *StorageAgent
	SandboxAgent *SandboxAgent
	Presence     *NodePresence
}

Runner reconstitutes the independently startable pieces of the runner role. The standalone runner command keeps this convenience API; the server boot graph starts and stops the pieces as separate owned components.

func NewRunner

func NewRunner(log *slog.Logger, deps RunnerDeps, cfg RunnerConfig) (*Runner, error)

func (*Runner) Close

func (r *Runner) Close() error

func (*Runner) ContainerdContainerForSandbox

func (r *Runner) ContainerdContainerForSandbox(ctx context.Context, id entity.Id) (containerd.Container, error)

func (*Runner) ContainerdNamespace

func (r *Runner) ContainerdNamespace() string

func (*Runner) Drain

func (r *Runner) Drain(ctx context.Context) error

func (*Runner) SetRestartMode added in v0.3.1

func (r *Runner) SetRestartMode(v bool)

func (*Runner) Start

func (r *Runner) Start(ctx context.Context, eg ...*errgroup.Group) error

func (*Runner) WorkloadIssuer added in v0.13.0

func (r *Runner) WorkloadIssuer() workloadidentity.TokenIssuer

type RunnerConfig

type RunnerConfig struct {
	Id            string `json:"id" cbor:"id" yaml:"id"`
	Name          string `json:"name" cbor:"name" yaml:"name"`
	ListenAddress string `json:"listen_address" cbor:"listen_address" yaml:"listen_address"`
	Workers       int    `json:"workers" cbor:"workers" yaml:"workers"`
	DataPath      string `json:"data_path" cbor:"data_path" yaml:"data_path"`

	// Optional RPC configuration for advanced setups
	// If not provided, a default insecure connection will be used
	// to connect to the server address.
	Config *clientconfig.Config `json:"config" cbor:"config" yaml:"config"`

	// Optional cloud authentication configuration for disk replication
	CloudAuth *coordinate.CloudAuthConfig `json:"cloud_auth,omitempty" cbor:"cloud_auth,omitempty" yaml:"cloud_auth,omitempty"`

	// DiskMode configures disk I/O mode ("", "auto", "universal", "accelerator")
	DiskMode string `json:"disk_mode,omitempty" cbor:"disk_mode,omitempty" yaml:"disk_mode,omitempty"`
}

type RunnerDeps added in v0.3.0

type RunnerDeps struct {
	CC        *containerd.Client
	Namespace string
	Bridge    string
	Tempdir   string
	Subnet    *netdb.Subnet

	// Network dependencies
	NetServ *network.ServiceManager

	// Observability dependencies
	LogsMaintainer *observability.LogsMaintainer
	LogWriter      observability.LogWriter
	StatusMon      *observability.StatusMonitor

	// MetricsWriter is also where host-level resource series are published. It
	// is the same writer the sandbox collectors use, taken directly rather than
	// through them because node series are labeled by node rather than by
	// sandbox. Nil disables host metrics, as it does for sandbox metrics.
	MetricsWriter *metrics.VictoriaMetricsWriter

	// Network config
	IPv4Routable    netip.Prefix
	ServicePrefixes []netip.Prefix
	DisableLocalNet bool

	// Resolver
	Resolver netresolve.Resolver

	// Sandbox metrics
	SandboxMetrics *sandbox.Metrics

	// IsCoordinator indicates this runner is the coordinator node.
	// Affects scheduling: stateful sandboxes are routed to the coordinator.
	IsCoordinator bool

	// Flannel network configuration (for distributed runners)
	// If EtcdEndpoints is non-empty, the runner will join the Flannel network
	EtcdEndpoints []string
	EtcdPrefix    string

	// TLS configuration for etcd mTLS (for distributed runners, file paths)
	EtcdTLSCertFile string // Client certificate file path
	EtcdTLSKeyFile  string // Client private key file path
	EtcdTLSCAFile   string // CA certificate file path

	// WorkloadIssuer mints workload identity tokens for sandbox containers. On
	// the coordinator this is the concrete *workloadidentity.Issuer; on a
	// distributed runner it is a remote issuer that proxies minting to the
	// coordinator over RPC.
	WorkloadIssuer workloadidentity.TokenIssuer

	// ApiAddress is where sandboxes on this host reach the cluster API, as a
	// literal IP:port. On the coordinator that is the local bridge router; on a
	// distributed runner it is the coordinator itself. Empty disables
	// in-cluster API access.
	ApiAddress string

	// CACert is the cluster CA in PEM form, mounted into sandboxes so they can
	// verify the API certificate.
	CACert []byte

	// Secrets materializes the secret references a sandbox spec carries, at
	// container creation. On the coordinator this is the local backend
	// registry, where the keyring lives. A distributed runner holds no key
	// material, so it must resolve over RPC instead; until that exists, a
	// sandbox on such a runner fails rather than starting without its secret.
	Secrets secret.Resolver
}

RunnerDeps holds dependencies needed by the Runner to construct controllers.

type SandboxAgent added in v0.15.0

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

SandboxAgent owns desired-state reconciliation for sandboxes and services. SandboxHost has already adopted surviving containers before this starts.

func NewSandboxAgent added in v0.15.0

func NewSandboxAgent(host *SandboxHost) *SandboxAgent

NewSandboxAgent constructs reconciliation over a restored sandbox host.

func (*SandboxAgent) Close added in v0.15.0

func (a *SandboxAgent) Close() error

func (*SandboxAgent) Start added in v0.15.0

func (a *SandboxAgent) Start(ctx context.Context) error

type SandboxHost added in v0.15.0

type SandboxHost struct {
	RunnerConfig

	Log *slog.Logger
	// contains filtered or unexported fields
}

SandboxHost owns the node-local execution substrate.

func NewSandboxHost added in v0.15.0

func NewSandboxHost(access *ClusterAccess, storage *NodeStorage, deps RunnerDeps, cfg RunnerConfig) (*SandboxHost, error)

NewSandboxHost constructs the part of a runner that owns local workload execution. Starting it joins the distributed network when needed, adopts surviving containers, and exposes the node-local exec endpoints. It does not publish the node as schedulable; NodePresence owns that later transition.

func (*SandboxHost) Close added in v0.15.0

func (r *SandboxHost) Close() error

func (*SandboxHost) ContainerdContainerForSandbox added in v0.15.0

func (r *SandboxHost) ContainerdContainerForSandbox(ctx context.Context, id entity.Id) (containerd.Container, error)

func (*SandboxHost) ContainerdNamespace added in v0.15.0

func (r *SandboxHost) ContainerdNamespace() string

func (*SandboxHost) SetupControllers added in v0.15.0

func (r *SandboxHost) SetupControllers(
	ctx context.Context,
	eas *es.EntityAccessClient,
	rs *rpc.Server,
) (
	_ *controller.ControllerManager,
	retErr error,
)

func (*SandboxHost) Start added in v0.15.0

func (r *SandboxHost) Start(ctx context.Context, eg ...*errgroup.Group) error

Start restores the local execution substrate without making the node schedulable. The optional errgroup runs distributed-network background work. The optional errgroup parameter is used for running background tasks like the Flannel network. If eg is nil and the runner needs to join a Flannel network, an internal errgroup will be created.

type StorageAgent added in v0.15.0

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

StorageAgent owns desired-state reconciliation for restored node storage.

func NewStorageAgent added in v0.15.0

func NewStorageAgent(storage *NodeStorage) *StorageAgent

NewStorageAgent constructs reconciliation over restored node storage.

func (*StorageAgent) Close added in v0.15.0

func (a *StorageAgent) Close() error

func (*StorageAgent) Start added in v0.15.0

func (a *StorageAgent) Start(ctx context.Context) error

Jump to

Keyboard shortcuts

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