Documentation
¶
Overview ¶
Package runtime abstracts where branch instances run (Docker now, K8s in P3).
Index ¶
- Constants
- Variables
- func DockerHostFromCLIContext() string
- func IsInUse(err error) bool
- func IsUnavailable(err error) bool
- func NewKubeClient(kubeconfig string) (kubernetes.Interface, error)
- func ValidateVolumeRoot(dir string) error
- type BranchSpec
- type CSIConfig
- type ContainerInfo
- type DockerDriver
- func (d *DockerDriver) CheckVolumeRoot(ctx context.Context) error
- func (d *DockerDriver) CloneVolume(ctx context.Context, src, dst string, labels map[string]string) error
- func (d *DockerDriver) CreateVolume(ctx context.Context, name string, labels map[string]string) error
- func (d *DockerDriver) EnsureImage(ctx context.Context, ref string) error
- func (d *DockerDriver) Exec(ctx context.Context, id string, cmd []string) error
- func (d *DockerDriver) ExecOutput(ctx context.Context, id string, cmd []string) (string, error)
- func (d *DockerDriver) Inspect(ctx context.Context, id string) (ContainerInfo, error)
- func (d *DockerDriver) ListHelpers(ctx context.Context) ([]ContainerInfo, error)
- func (d *DockerDriver) ListManaged(ctx context.Context) ([]ContainerInfo, error)
- func (d *DockerDriver) ListManagedVolumes(ctx context.Context, instanceID string) ([]VolumeInfo, error)
- func (d *DockerDriver) RemoveVolume(ctx context.Context, name string) error
- func (d *DockerDriver) RunHelper(ctx context.Context, spec HelperSpec) (string, error)
- func (d *DockerDriver) StartBranch(ctx context.Context, spec BranchSpec) (string, error)
- func (d *DockerDriver) StopRemove(ctx context.Context, id string) error
- func (d *DockerDriver) VolumeRoot() string
- type DockerOption
- type Driver
- type HelperSpec
- type KubeDriver
- func (d *KubeDriver) CloneVolume(ctx context.Context, src, dst string, labels map[string]string) error
- func (d *KubeDriver) CreateVolume(ctx context.Context, name string, labels map[string]string) error
- func (d *KubeDriver) DiagnoseContainer(ctx context.Context, id string) (string, bool)
- func (d *KubeDriver) EnsureImage(ctx context.Context, image string) error
- func (d *KubeDriver) Exec(ctx context.Context, id string, cmd []string) error
- func (d *KubeDriver) ExecOutput(ctx context.Context, id string, cmd []string) (string, error)
- func (d *KubeDriver) Inspect(ctx context.Context, id string) (ContainerInfo, error)
- func (d *KubeDriver) ListHelpers(ctx context.Context) ([]ContainerInfo, error)
- func (d *KubeDriver) ListManaged(ctx context.Context) ([]ContainerInfo, error)
- func (d *KubeDriver) ListManagedVolumes(ctx context.Context, instanceID string) ([]VolumeInfo, error)
- func (d *KubeDriver) RemoveVolume(ctx context.Context, name string) error
- func (d *KubeDriver) RunHelper(ctx context.Context, spec HelperSpec) (string, error)
- func (d *KubeDriver) StartBranch(ctx context.Context, spec BranchSpec) (string, error)
- func (d *KubeDriver) StopRemove(ctx context.Context, id string) error
- type KubeOption
- type Mount
- type MountKind
- type VolumeInfo
Constants ¶
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.
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.
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.
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 ¶
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.
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 ¶
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 ¶
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 ¶
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) ExecOutput ¶
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 ¶
CreateVolume provisions an empty volume (hostPath: node dir via helper pod; csi: PVC) carrying the given labels.
func (*KubeDriver) DiagnoseContainer ¶
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 ¶
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 ¶
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 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.