runtime

package
v1.0.0 Latest Latest
Warning

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

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

Documentation

Overview

Package runtime abstracts where branch instances run (Docker now, K8s in P3).

Index

Constants

View Source
const LabelBranchID = "pgoverlay.branch.id"

LabelBranchID tags a branch container/pod (and its rw volume) with the id of the registry row it belongs to. Reconcile uses it to recognise a container a live saga has started but not yet recorded on the row.

View Source
const LabelInstance = "pgoverlay.instance"

LabelInstance tags every managed resource with the owning registry's instance id (registry.InstanceID()). Reconcile reclaims only resources carrying ITS id, so concurrent pgoverlay instances sharing one Docker daemon (and the parallel IT suite) never GC each other's live containers/volumes.

View Source
const LabelVolumeRoot = "pgoverlay.volume-root"

LabelVolumeRoot records on a volume created under a volume root which root that was. RemoveVolume deletes the directory under the root the volume was created with, so changing or dropping --volume-root later strands nothing.

View Source
const UtilityImage = "alpine:3.24@sha256:294b683cb724975bec92580e1e685676bd4b50bda910ddb8c51d4cabeaec77e6"

UtilityImage is the small shell image every file-level helper runs in: preparing a seed volume, installing a branch entrypoint, measuring disk use, copying and removing hostPath volume directories, the zfs userland.

It is pinned by digest, not by tag. Some of these helpers run as root with the storage node's whole data root mounted, so a re-pointed tag (or a poisoned pull-through cache) would run someone else's code against every branch's data. The digest is the one the Dockerfiles' runtime stage builds on, and TestUtilityImageMatchesDockerfiles keeps the two from drifting: a base-image bump that forgets this constant fails the unit suite.

The kube driver runs helpers that ask for UtilityImage on --kube-helper-image instead when it is set (a mirror in a private or air-gapped registry).

Variables

View Source
var ErrNotFound = errors.New("not found")

ErrNotFound is wrapped into the error Inspect returns when the container or pod does not exist, so callers can tell "gone" from "could not ask" with errors.Is.

View Source
var ErrVolumeExists = errors.New("volume already exists")

ErrVolumeExists is wrapped into the error CreateVolume returns when a volume of that name already exists. Creating never adopts an existing volume: a stale volume left behind under a reused name would otherwise become a new branch's writable layer, with the old branch's writes in it.

Functions

func DockerHostFromCLIContext

func DockerHostFromCLIContext() string

DockerHostFromCLIContext resolves the docker endpoint from the CLI's current context ($DOCKER_CONTEXT, else ~/.docker/config.json, honouring $DOCKER_CONFIG). The Go SDK's FromEnv only honors DOCKER_HOST, so without this, setups like Colima (where /var/run/docker.sock is absent or stale) fail. Returns "" when unresolvable. The endpoint is returned as stored, so it may be an ssh:// URL that the SDK cannot dial (NewDockerDriver rejects those). Exported so integration helpers can prime DOCKER_HOST for libraries (e.g. testcontainers) that don't read docker CLI contexts.

func IsInUse

func IsInUse(err error) bool

IsInUse reports whether err, from a Driver call or a helper it ran, is the runtime refusing to remove something that is still in use: a docker volume a container still mounts (409 Conflict, "volume is in use"), or a busy zfs dataset or mount. Removing it again once that user is gone succeeds.

func IsUnavailable

func IsUnavailable(err error) bool

IsUnavailable reports whether err means the runtime could not be reached or did not answer in time: the Docker daemon is down or its socket is gone, the Kubernetes API server is unavailable, overloaded or timed out, a connection was refused, or the call ran out of time.

func NewKubeClient

func NewKubeClient(kubeconfig string) (kubernetes.Interface, error)

NewKubeClient builds a kubernetes.Interface from the same in-cluster / kubeconfig loading path the kube driver uses (kubeconfig=="" → in-cluster, then KUBECONFIG / ~/.kube/config). branchd reuses it for the leader-election Lease so HA shares the driver's cluster credentials.

func ValidateVolumeRoot

func ValidateVolumeRoot(dir string) error

ValidateVolumeRoot rejects a volume root that cannot be used as one: it must be an absolute, clean path on the (Linux) Docker host, other than /.

Types

type BranchSpec

type BranchSpec struct {
	Name       string // container name, e.g. pgoverlay-br-pr-1
	Image      string
	Env        []string
	Mounts     []Mount
	Entrypoint []string // overrides image entrypoint
	Labels     map[string]string
	Network    string
}

BranchSpec is a long-running branch Postgres container.

type CSIConfig

type CSIConfig struct {
	// StorageClass provisions every pgoverlay PVC; it must support PVC
	// dataSource cloning (or VolumeSnapshots when SnapshotClass is set).
	StorageClass string
	// SnapshotClass switches branch cloning from PVC dataSource clones to
	// VolumeSnapshot + restore ("" = direct PVC clones).
	SnapshotClass string
	// VolumeSize is the storage request of every pgoverlay PVC ("" = 10Gi).
	VolumeSize string
}

CSIConfig configures the csi storage strategy.

type ContainerInfo

type ContainerInfo struct {
	ID      string
	Running bool
	Stopped bool
	Status  string // the runtime's own description of the state, for messages
	Host    string // address the instance is reachable on (127.0.0.1 for docker, pod IP for k8s); "" unless Running
	Port    int    // port on Host (docker: host port mapped to 5432, 0 if none)
	Created time.Time
	Labels  map[string]string
}

ContainerInfo describes one branch or helper instance as the runtime sees it.

Running and Stopped are not complements. Running means the instance is up and Host/Port are current. Stopped means it is down and the runtime will not bring it back on its own: a docker container that exited, died or was created but never started, or a pod in phase Failed (for example evicted) or Succeeded. Neither being set means the runtime is still working on it (docker restarting, a Pending or terminating pod): look again later.

type DockerDriver

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

func NewDockerDriver

func NewDockerDriver(opts ...DockerOption) (*DockerDriver, error)

NewDockerDriver builds a client for the Docker endpoint the docker CLI would use: DOCKER_HOST (with DOCKER_CERT_PATH / DOCKER_TLS_VERIFY), else the current CLI context, including its TLS material. Endpoints the SDK cannot dial (ssh://) are rejected with an explanation instead of failing later with a misleading HTTP error.

func (*DockerDriver) CheckVolumeRoot

func (d *DockerDriver) CheckVolumeRoot(ctx context.Context) error

CheckVolumeRoot verifies that the volume root exists on the Docker host and is a writable directory. branchd runs it at startup so a missing disk fails there, not in the first branch create. No-op without a volume root.

func (*DockerDriver) CloneVolume

func (d *DockerDriver) CloneVolume(ctx context.Context, src, dst string, labels map[string]string) error

CloneVolume provisions dst as a copy of src. Docker named volumes have no copy-on-write clone primitive, so this is a full `cp -a` through a helper container — a generic fallback satisfying the Driver contract; no engine flow uses it on docker today (the overlay and zfs backends have cheaper mechanisms).

func (*DockerDriver) CreateVolume

func (d *DockerDriver) CreateVolume(ctx context.Context, name string, labels map[string]string) error

CreateVolume creates an empty named volume. Docker's VolumeCreate is idempotent on the name — it hands back an existing volume unchanged, labels and data included — so the name is checked first and an existing volume is an error wrapping ErrVolumeExists. Volume names are derived from branch and source names, so adopting silently would give a recreated branch the writes (or the frozen layer) of whatever last used the name.

With a volume root (WithVolumeRoot) the volume is a local-driver bind volume over a fresh directory <root>/<name>; see createRootVolume.

func (*DockerDriver) EnsureImage

func (d *DockerDriver) EnsureImage(ctx context.Context, ref string) error

func (*DockerDriver) Exec

func (d *DockerDriver) Exec(ctx context.Context, id string, cmd []string) error

func (*DockerDriver) ExecOutput

func (d *DockerDriver) ExecOutput(ctx context.Context, id string, cmd []string) (string, error)

ExecOutput runs cmd in the container as execUser and returns its stdout. The exec API multiplexes stdout/stderr over one attached stream (stdcopy); stderr is kept separate so captured output (e.g. a pg_dump) stays clean, and is embedded in the error on non-zero exit.

Success needs both a complete stream and a finished command with exit code 0. A stream that breaks early (a dropped connection to the daemon) is an error rather than truncated output, and the exit code is read only once the daemon reports the exec is no longer running: ExecInspect's ExitCode is 0 until then, which would pass a command that has not finished (or that will fail) as a success.

func (*DockerDriver) Inspect

func (d *DockerDriver) Inspect(ctx context.Context, id string) (ContainerInfo, error)

func (*DockerDriver) ListHelpers

func (d *DockerDriver) ListHelpers(ctx context.Context) ([]ContainerInfo, error)

ListHelpers lists every pgoverlay helper container on the daemon, whichever instance ran it (helpers carry no instance label).

func (*DockerDriver) ListManaged

func (d *DockerDriver) ListManaged(ctx context.Context) ([]ContainerInfo, error)

func (*DockerDriver) ListManagedVolumes

func (d *DockerDriver) ListManagedVolumes(ctx context.Context, instanceID string) ([]VolumeInfo, error)

func (*DockerDriver) RemoveVolume

func (d *DockerDriver) RemoveVolume(ctx context.Context, name string) error

RemoveVolume removes the volume; a missing one is success. A bind volume created under a volume root loses its directory too, once the volume itself is gone (removing a volume a container still uses fails first, and leaves the data alone).

func (*DockerDriver) RunHelper

func (d *DockerDriver) RunHelper(ctx context.Context, spec HelperSpec) (string, error)

func (*DockerDriver) StartBranch

func (d *DockerDriver) StartBranch(ctx context.Context, spec BranchSpec) (string, error)

func (*DockerDriver) StopRemove

func (d *DockerDriver) StopRemove(ctx context.Context, id string) error

StopRemove stops and removes the container and waits until it is actually gone, matching KubeDriver.StopRemove. Idempotent: NotFound is success.

The wait is load-bearing, not politeness. ContainerRemove returns once the daemon has ACCEPTED the removal, while teardown continues in the background and the container keeps its reference on the rw volume until it finishes. Every caller in internal/engine follows StopRemove with RemoveVolume on that volume (saga.go, csi.go, freeze.go, reconcile.go), so returning early makes the next call fail with "volume is in use - [<container id>]". That is a timing-dependent failure: it needs the volume removal to land inside the teardown window, so it passes most runs and fails perhaps one in three. KubeDriver already waits and documents why; this is the same contract, which the interface has always implied because of how its callers are written.

func (*DockerDriver) VolumeRoot

func (d *DockerDriver) VolumeRoot() string

VolumeRoot is the configured volume root ("" = docker-managed volumes).

type DockerOption

type DockerOption func(*DockerDriver)

DockerOption configures a DockerDriver.

func WithVolumeRoot

func WithVolumeRoot(dir string) DockerOption

WithVolumeRoot makes every volume the driver creates a bind volume over a directory under dir on the Docker host (see the comment at the top of this file). "" keeps docker-managed volumes. Volumes that already exist keep working either way: they are mounted by name, and removed according to how they were created.

type Driver

type Driver interface {
	EnsureImage(ctx context.Context, image string) error
	// CreateVolume provisions an empty volume carrying labels. It does not
	// adopt an existing volume of the same name: the docker driver returns an
	// error wrapping ErrVolumeExists, and a csi PVC create fails with
	// AlreadyExists.
	CreateVolume(ctx context.Context, name string, labels map[string]string) error
	RemoveVolume(ctx context.Context, name string) error
	// CloneVolume provisions dst as a copy of src (dst must not exist; labels
	// land on dst). Copy-on-write where the storage supports it (kube csi:
	// PVC dataSource clone / snapshot restore); a full copy elsewhere.
	CloneVolume(ctx context.Context, src, dst string, labels map[string]string) error
	RunHelper(ctx context.Context, spec HelperSpec) (output string, err error)
	StartBranch(ctx context.Context, spec BranchSpec) (id string, err error)
	Exec(ctx context.Context, containerID string, cmd []string) error // error on non-zero exit
	// ExecOutput is Exec with the command's stdout captured and returned
	// (stderr goes into the error on failure). Used where the engine needs
	// in-container command output (pg_dump, row-count probes) against a
	// RUNNING instance — RunHelper can also capture output but spins a new
	// container, which cannot reach an instance's local socket.
	ExecOutput(ctx context.Context, containerID string, cmd []string) (string, error)
	// Inspect reports one instance. A missing container/pod is an error
	// wrapping ErrNotFound.
	Inspect(ctx context.Context, containerID string) (ContainerInfo, error)
	StopRemove(ctx context.Context, containerID string) error
	ListManaged(ctx context.Context) ([]ContainerInfo, error) // labels pgoverlay.managed=true, pgoverlay.role=branch
	// ListHelpers returns every pgoverlay helper container/pod (labels
	// pgoverlay.managed=true, pgoverlay.role=helper), running or not. Helpers
	// are removed by the process that ran them; reconcile uses this to find
	// the ones a crashed process left behind.
	ListHelpers(ctx context.Context) ([]ContainerInfo, error)
	// ListManagedVolumes returns every volume carrying both the
	// pgoverlay.managed=true label AND pgoverlay.instance=<instanceID> (docker
	// named volumes / kube PVCs). Reconcile uses it to find orphaned rw and
	// source-layer volumes owned by THIS registry; scoping by instance id keeps
	// one instance from GC'ing another's volumes on a shared daemon. The zfs
	// backend manages datasets, not driver volumes, so its driver may return an
	// empty list (zfs orphans are GC'd via the per-branch/source paths instead).
	ListManagedVolumes(ctx context.Context, instanceID string) ([]VolumeInfo, error)
}

type HelperSpec

type HelperSpec struct {
	Image   string
	Cmd     []string
	Env     []string
	Mounts  []Mount
	Network string
	User    string // e.g. "postgres" for pg_basebackup so file ownership is uid 999
	// Privileged runs the helper with full privileges, with HostDevices
	// mapped in (zfs backend: /dev/zfs). The docker driver maps the devices
	// explicitly; on kube a privileged container sees host devices anyway.
	Privileged  bool
	HostDevices []string
	// SysAdmin runs the helper with the privileges an overlay branch
	// container has, and no more: CAP_SYS_ADMIN (to mount an overlay) with
	// the AppArmor profile, and on kube the seccomp profile, unconfined. The
	// copy-up probe uses it. Ignored when Privileged is set.
	SysAdmin bool
}

HelperSpec is a one-shot container performing a data operation (seeding, file fixes, measurements). Run blocks until exit and returns the captured combined output; non-zero exit = error including that output.

type KubeDriver

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

KubeDriver runs branches as pods. Where the data lives is a pluggable storage strategy:

  • hostPath (default): "volumes" are subdirectories of dataRoot on one designated storage node; every pod is pinned there with nodeName and branch pods get SYS_ADMIN for their in-container overlay mount (decision 1: single-node dev/test scope).
  • csi: "volumes" are PersistentVolumeClaims and branches are PVC clones; pods schedule anywhere, need no extra capabilities, and run postgres directly on their claim (multi-node scope, decision 4 / Phase 5 D).

Container IDs are pod names.

func NewKubeDriver

func NewKubeDriver(kubeconfig, namespace, nodeName, dataRoot string, opts ...KubeOption) (*KubeDriver, error)

NewKubeDriver connects to the cluster with the hostPath storage strategy (all data under dataRoot on the named storage node).

func NewKubeDriverCSI

func NewKubeDriverCSI(kubeconfig, namespace string, csi CSIConfig, opts ...KubeOption) (*KubeDriver, error)

NewKubeDriverCSI connects to the cluster with the csi storage strategy: volumes are PVCs, branches are PVC clones, pods schedule on any node.

func (*KubeDriver) CloneVolume

func (d *KubeDriver) CloneVolume(ctx context.Context, src, dst string, labels map[string]string) error

CloneVolume provisions dst as a copy of src: a full `cp -a` through a helper pod (hostPath) or a copy-on-write PVC clone / snapshot restore (csi).

func (*KubeDriver) CreateVolume

func (d *KubeDriver) CreateVolume(ctx context.Context, name string, labels map[string]string) error

CreateVolume provisions an empty volume (hostPath: node dir via helper pod; csi: PVC) carrying the given labels.

func (*KubeDriver) DiagnoseContainer

func (d *KubeDriver) DiagnoseContainer(ctx context.Context, id string) (string, bool)

DiagnoseContainer explains why a branch pod is not serving yet: its scheduling condition, the container's waiting reason, or its last termination, plus the tail of its logs when it crashed. fatal reports a state that will not resolve by waiting (image pull back-off, invalid image, bad container config, crash loop). The engine's readiness wait consults it (an optional driver capability) so a hopeless start fails fast with its cause, and the cause lands in the error before the compensation deletes the pod.

func (*KubeDriver) EnsureImage

func (d *KubeDriver) EnsureImage(ctx context.Context, image string) error

EnsureImage is a no-op: the kubelet pulls images on pod start.

func (*KubeDriver) Exec

func (d *KubeDriver) Exec(ctx context.Context, id string, cmd []string) error

Exec runs cmd in the pod's first container and fails on non-zero exit with captured output, matching the docker driver's contract.

func (*KubeDriver) ExecOutput

func (d *KubeDriver) ExecOutput(ctx context.Context, id string, cmd []string) (string, error)

ExecOutput runs cmd in the pod's first container over the SPDY exec subresource and returns the captured stdout; stderr is kept separate and embedded in the error on failure (non-zero exit surfaces as a stream error).

func (*KubeDriver) Inspect

func (d *KubeDriver) Inspect(ctx context.Context, id string) (ContainerInfo, error)

func (*KubeDriver) ListHelpers

func (d *KubeDriver) ListHelpers(ctx context.Context) ([]ContainerInfo, error)

ListHelpers lists the helper pods in the namespace, running or finished.

func (*KubeDriver) ListManaged

func (d *KubeDriver) ListManaged(ctx context.Context) ([]ContainerInfo, error)

func (*KubeDriver) ListManagedVolumes

func (d *KubeDriver) ListManagedVolumes(ctx context.Context, instanceID string) ([]VolumeInfo, error)

ListManagedVolumes returns every pgoverlay-managed volume name (delegated to the storage strategy: hostPath dirs or labelled PVCs).

func (*KubeDriver) RemoveVolume

func (d *KubeDriver) RemoveVolume(ctx context.Context, name string) error

RemoveVolume deletes the volume. Idempotent (removing a missing volume succeeds).

func (*KubeDriver) RunHelper

func (d *KubeDriver) RunHelper(ctx context.Context, spec HelperSpec) (string, error)

RunHelper runs spec to completion in a one-shot helper pod. The helper's environment travels in a short-lived Secret, never in the pod spec.

func (*KubeDriver) StartBranch

func (d *KubeDriver) StartBranch(ctx context.Context, spec BranchSpec) (string, error)

func (*KubeDriver) StopRemove

func (d *KubeDriver) StopRemove(ctx context.Context, id string) error

StopRemove deletes the pod (30s grace, background propagation) and waits until it is gone so a same-name recreate (branch reset) cannot collide. Removing a helper pod (an orphan a crashed branchd left behind) also removes its env Secret. Idempotent: NotFound is success.

type KubeOption

type KubeOption func(*KubeDriver)

KubeOption configures a KubeDriver beyond its storage strategy.

func WithHelperImage

func WithHelperImage(image string) KubeOption

WithHelperImage runs every helper that asks for UtilityImage on image instead (branchd --kube-helper-image): a mirror in a private or air-gapped registry, or a newer pinned digest. "" keeps UtilityImage.

func WithInstanceID

func WithInstanceID(id string) KubeOption

WithInstanceID stamps LabelInstance=id on helper pods and their Secrets, the same instance label branch pods and volumes carry, so a helper left behind by a crashed branchd can be attributed to the registry that ran it.

func WithOwnerPod

func WithOwnerPod(namespace, name, uid string) KubeOption

WithOwnerPod makes branchd's own pod the owner of every helper pod and helper Secret, so Kubernetes garbage collection removes them once that pod is gone: a branchd killed mid-seed (a crash, a rollout past the shutdown budget) would otherwise leave its helper running, and its Secret stored, with nothing to clean them up. namespace, name and uid come from the downward API. It is ignored unless all three are set and namespace is the driver's: an owner in another namespace counts as missing, which would get every helper collected the moment it is created.

type Mount

type Mount struct {
	Kind     MountKind
	Volume   string // volume name (MountVolume) or absolute host path (MountHostPath)
	Target   string
	ReadOnly bool
}

type MountKind

type MountKind string

MountKind selects how Mount.Volume is interpreted.

const (
	// MountVolume (the zero value) names a managed volume: a docker named
	// volume, or a dataRoot subdirectory on the kube storage node.
	MountVolume MountKind = ""
	// MountHostPath bind-mounts an absolute host path. Used by the zfs
	// backend to mount dataset mountpoints; the path must already exist.
	MountHostPath MountKind = "hostpath"
)

type VolumeInfo

type VolumeInfo struct {
	Name string
	// Created is when the volume was created; zero when the runtime cannot
	// tell. Reconcile never garbage-collects a volume younger than its grace
	// period, because a saga in another process may have just created it.
	Created time.Time
	// Leftover marks storage the runtime no longer has a volume for: a
	// directory under the docker driver's volume root (WithVolumeRoot) whose
	// volume is gone, because a removal stopped between deleting the volume
	// and deleting its directory, or someone removed the volume by hand.
	// Nothing can mount it (a container would get a fresh, empty volume of
	// that name), so it counts as missing everywhere except garbage
	// collection, where RemoveVolume deletes the directory.
	Leftover bool
}

VolumeInfo is one managed volume as ListManagedVolumes reports it.

Jump to

Keyboard shortcuts

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