1103 lines
37 KiB
Go
1103 lines
37 KiB
Go
// Package backup creates portable, transactional ThothII installation archives.
|
|
package backup
|
|
|
|
import (
|
|
"archive/zip"
|
|
"bytes"
|
|
"context"
|
|
"crypto/sha256"
|
|
"encoding/hex"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"os"
|
|
"os/exec"
|
|
"path/filepath"
|
|
"sort"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/aritmolab/thothii/tools/tht/internal/compose"
|
|
"github.com/aritmolab/thothii/tools/tht/internal/config"
|
|
"github.com/aritmolab/thothii/tools/tht/internal/lifecycle"
|
|
"github.com/aritmolab/thothii/tools/tht/internal/safeio"
|
|
"github.com/aritmolab/thothii/tools/tht/internal/service"
|
|
"gopkg.in/yaml.v3"
|
|
)
|
|
|
|
var (
|
|
// ErrActiveSessions keeps the refusal machine-readable for the host CLI.
|
|
ErrActiveSessions = errors.New("active sessions require --drain before backup")
|
|
// ErrConfirmationRequired guards backups that would contain external secret payloads.
|
|
ErrConfirmationRequired = errors.New("--include-secrets requires --yes")
|
|
)
|
|
|
|
var requiredVolumes = []string{
|
|
"settings",
|
|
"pi-state",
|
|
"workspace-registry",
|
|
"workspace-secrets",
|
|
"sessions",
|
|
"qdrant-data",
|
|
"embedding-models",
|
|
}
|
|
|
|
var requiredServerVolumes = []string{"qdrant-data", "embedding-models"}
|
|
|
|
func requiredBackupVolumes(installation config.Installation) []string {
|
|
if installation.Profile == "server" {
|
|
return requiredServerVolumes
|
|
}
|
|
return requiredVolumes
|
|
}
|
|
|
|
const (
|
|
helperImage = "busybox:1.36.1"
|
|
drainPollInterval = time.Second
|
|
maxDrainPolls = 300
|
|
// Cleanup must survive a caller timeout or a lost Docker response, but it must not run
|
|
// indefinitely after the command has returned. Archive recovery can require volume work, so
|
|
// keep this deliberately longer than an individual health check.
|
|
cleanupOperationTimeout = 5 * time.Minute
|
|
)
|
|
|
|
// CreateRequest controls one explicit backup request.
|
|
type CreateRequest struct {
|
|
Output string
|
|
IncludeSecrets bool
|
|
Confirm bool
|
|
Drain bool
|
|
}
|
|
|
|
// Result describes a published archive without exposing its contents.
|
|
type Result struct {
|
|
Path string
|
|
Warning string
|
|
}
|
|
|
|
// archiveRunner separates the only streaming Docker call from regular Compose commands. Tests
|
|
// provide a deterministic in-memory implementation; production always passes argument arrays.
|
|
type archiveRunner interface {
|
|
compose.Runner
|
|
Stream(context.Context, []string, io.Reader, io.Writer) (compose.Result, error)
|
|
SessionInventoryScope() string
|
|
}
|
|
|
|
type dependencies struct {
|
|
runner archiveRunner
|
|
now func() time.Time
|
|
homeDir func() (string, error)
|
|
revision func(context.Context, string) (string, error)
|
|
sleep func(time.Duration)
|
|
reserveOutput func(string) (*archiveReservation, error)
|
|
publishReserved func(*archiveReservation, string) error
|
|
}
|
|
|
|
// Create creates an archive with the real Docker command boundary. It performs no shell
|
|
// interpolation and never supplies secret contents in process arguments.
|
|
func Create(ctx context.Context, installation config.Installation, request CreateRequest) (Result, error) {
|
|
return createWithDependencies(ctx, installation, request, productionCreateDependencies(installation))
|
|
}
|
|
|
|
func productionCreateDependencies(installation config.Installation) dependencies {
|
|
runner := hostRunner{runner: compose.NewRunner(""), binary: "docker", profile: installation.Profile}
|
|
return dependencies{
|
|
runner: runner,
|
|
now: time.Now,
|
|
homeDir: os.UserHomeDir,
|
|
revision: func(ctx context.Context, directory string) (string, error) {
|
|
command := exec.CommandContext(ctx, "git", "-C", directory, "rev-parse", "HEAD")
|
|
value, err := command.Output()
|
|
if err != nil {
|
|
return "", errors.New("source revision is unavailable")
|
|
}
|
|
return strings.TrimSpace(string(value)), nil
|
|
},
|
|
sleep: time.Sleep,
|
|
reserveOutput: reserveArchiveOutput,
|
|
publishReserved: publishReservedArchive,
|
|
}
|
|
}
|
|
|
|
func createWithDependencies(ctx context.Context, installation config.Installation, request CreateRequest, dependencies dependencies) (result Result, resultErr error) {
|
|
if dependencies.runner == nil || dependencies.now == nil || dependencies.homeDir == nil || dependencies.revision == nil || dependencies.sleep == nil || dependencies.reserveOutput == nil || dependencies.publishReserved == nil {
|
|
return Result{}, errors.New("backup dependencies are incomplete")
|
|
}
|
|
if request.IncludeSecrets && !request.Confirm {
|
|
return Result{}, ErrConfirmationRequired
|
|
}
|
|
lock, err := lifecycle.Acquire(installation)
|
|
if err != nil {
|
|
return Result{}, err
|
|
}
|
|
defer func() {
|
|
if releaseErr := lock.Release(); releaseErr != nil {
|
|
result = Result{}
|
|
resultErr = errors.Join(resultErr, fmt.Errorf("release backup lifecycle lock: %w", releaseErr))
|
|
}
|
|
}()
|
|
return createWithDependenciesLockHeld(ctx, installation, request, dependencies)
|
|
}
|
|
|
|
// createWithDependenciesLockHeld performs backup creation while the caller owns the installation
|
|
// lifecycle lock. It must never acquire a lifecycle lock itself: Restore uses this primitive to
|
|
// create its recovery checkpoint inside its already-locked transaction.
|
|
func createWithDependenciesLockHeld(ctx context.Context, installation config.Installation, request CreateRequest, dependencies dependencies) (result Result, resultErr error) {
|
|
if err := rejectInlineSecretValues(installation.EnvFile); err != nil {
|
|
return Result{}, err
|
|
}
|
|
installationID, err := backupInstallationID(installation)
|
|
if err != nil {
|
|
return Result{}, err
|
|
}
|
|
secretPaths, err := installation.SecretFiles()
|
|
if err != nil {
|
|
return Result{}, fmt.Errorf("installation external secret references could not be read: %w", err)
|
|
}
|
|
authenticationPaths, err := authenticationConfigFiles(installation)
|
|
if err != nil {
|
|
return Result{}, err
|
|
}
|
|
revision, err := dependencies.revision(ctx, installation.ProjectDirectory)
|
|
if err != nil {
|
|
return Result{}, err
|
|
}
|
|
if !revisionPattern.MatchString(revision) {
|
|
return Result{}, errors.New("source revision is invalid")
|
|
}
|
|
output, err := backupOutputPath(request.Output, installationID, revision, dependencies.now().UTC(), dependencies.homeDir)
|
|
if err != nil {
|
|
return Result{}, err
|
|
}
|
|
if err := ensureNewArchivePath(output); err != nil {
|
|
return Result{}, err
|
|
}
|
|
reservation, err := dependencies.reserveOutput(output)
|
|
if err != nil {
|
|
return Result{}, err
|
|
}
|
|
published := false
|
|
defer func() {
|
|
if !published {
|
|
_ = reservation.RemoveIfOwned()
|
|
}
|
|
}()
|
|
|
|
rendered, err := renderedConfiguration(ctx, installation, dependencies.runner)
|
|
if err != nil {
|
|
return Result{}, err
|
|
}
|
|
volumes, err := inspectRequiredVolumes(ctx, installation, dependencies.runner, rendered)
|
|
if err != nil {
|
|
return Result{}, err
|
|
}
|
|
images, err := imageIdentities(ctx, installation, dependencies.runner, rendered)
|
|
if err != nil {
|
|
return Result{}, err
|
|
}
|
|
|
|
wasRunning, err := installationRunning(ctx, installation, dependencies.runner)
|
|
if err != nil {
|
|
return Result{}, err
|
|
}
|
|
// Record attempted mutations before invoking Docker. Docker can apply a mutation and then
|
|
// lose its response, so a successful return is not evidence that compensation is unnecessary.
|
|
maintenanceAttempted := false
|
|
stopAttempted := false
|
|
defer func() {
|
|
var cleanupErr error
|
|
restartCompleted := true
|
|
if wasRunning && stopAttempted {
|
|
startErr, started := retryBoundedCleanup(func(cleanupContext context.Context) error {
|
|
return composeStartAndVerify(cleanupContext, installation, dependencies.runner)
|
|
})
|
|
if startErr != nil {
|
|
cleanupErr = errors.Join(cleanupErr, fmt.Errorf("backup maintenance cleanup restart: %w", startErr))
|
|
}
|
|
if started {
|
|
stopAttempted = false
|
|
} else {
|
|
restartCompleted = false
|
|
}
|
|
}
|
|
// Do not reopen admissions while the installation is known to be stopped. If a retry could
|
|
// not establish a running core, retain the durable barrier and report every cleanup error.
|
|
if maintenanceAttempted && restartCompleted {
|
|
deactivateErr, deactivated := retryBoundedCleanup(func(cleanupContext context.Context) error {
|
|
return maintenance(cleanupContext, installation, dependencies.runner, false)
|
|
})
|
|
if deactivateErr != nil {
|
|
cleanupErr = errors.Join(cleanupErr, fmt.Errorf("backup maintenance cleanup: %w", deactivateErr))
|
|
}
|
|
if deactivated {
|
|
maintenanceAttempted = false
|
|
}
|
|
}
|
|
if cleanupErr != nil {
|
|
result = Result{}
|
|
resultErr = errors.Join(resultErr, cleanupErr)
|
|
}
|
|
}()
|
|
|
|
if wasRunning {
|
|
maintenanceAttempted = true
|
|
if err := maintenance(ctx, installation, dependencies.runner, true); err != nil {
|
|
return Result{}, err
|
|
}
|
|
if err := waitForNoActiveSessions(ctx, installation, dependencies.runner, request.Drain, dependencies.sleep); err != nil {
|
|
return Result{}, err
|
|
}
|
|
stopAttempted = true
|
|
if err := runCompose(ctx, installation, dependencies.runner, "stop"); err != nil {
|
|
return Result{}, err
|
|
}
|
|
}
|
|
|
|
manifest := Manifest{
|
|
SchemaVersion: CurrentSchemaVersion,
|
|
InstallationID: installationID,
|
|
CreatedAt: dependencies.now().UTC(),
|
|
SourceRevision: revision,
|
|
IncludesSecrets: request.IncludeSecrets && len(secretPaths)+len(authenticationPaths) > 0,
|
|
ComposeProject: installation.ProjectName(),
|
|
Images: images,
|
|
Volumes: volumes,
|
|
}
|
|
if err := writeArchive(ctx, output, reservation, installation, request, secretPaths, authenticationPaths, manifest, volumes, dependencies); err != nil {
|
|
return Result{}, err
|
|
}
|
|
published = true
|
|
if wasRunning {
|
|
if err := composeStartAndVerify(ctx, installation, dependencies.runner); err != nil {
|
|
return Result{}, err
|
|
}
|
|
stopAttempted = false
|
|
if err := maintenance(ctx, installation, dependencies.runner, false); err != nil {
|
|
return Result{}, err
|
|
}
|
|
maintenanceAttempted = false
|
|
}
|
|
result = Result{Path: output}
|
|
if manifest.IncludesSecrets {
|
|
result.Warning = "The archive contains external secret files, including authentication configuration. Protect its custody and access."
|
|
}
|
|
return result, nil
|
|
}
|
|
|
|
type hostRunner struct {
|
|
runner compose.Runner
|
|
binary string
|
|
profile string
|
|
}
|
|
|
|
func (runner hostRunner) Run(ctx context.Context, args []string, stdin io.Reader) (compose.Result, error) {
|
|
return runner.runner.Run(ctx, args, stdin)
|
|
}
|
|
|
|
func (runner hostRunner) Stream(ctx context.Context, args []string, stdin io.Reader, stdout io.Writer) (compose.Result, error) {
|
|
command := exec.CommandContext(ctx, runner.binary, args...)
|
|
command.Stdin = stdin
|
|
command.Stdout = stdout
|
|
var stderr bytes.Buffer
|
|
command.Stderr = &stderr
|
|
err := command.Run()
|
|
result := compose.Result{Stderr: stderr.String()}
|
|
if exitErr := new(exec.ExitError); errors.As(err, &exitErr) {
|
|
result.ExitCode = exitErr.ExitCode()
|
|
}
|
|
return result, err
|
|
}
|
|
|
|
func (runner hostRunner) SessionInventoryScope() string {
|
|
if runner.profile == "local" {
|
|
return "mine"
|
|
}
|
|
return "all"
|
|
}
|
|
|
|
func backupInstallationID(installation config.Installation) (string, error) {
|
|
value := filepath.Base(filepath.Dir(installation.Path))
|
|
if value == "" || value == "." || value == ".." || strings.ContainsAny(value, `/\\`) {
|
|
return "", errors.New("installation ID is invalid")
|
|
}
|
|
return value, nil
|
|
}
|
|
|
|
func backupOutputPath(requested, installationID, revision string, now time.Time, homeDir func() (string, error)) (string, error) {
|
|
if requested != "" {
|
|
if !filepath.IsAbs(requested) {
|
|
absolute, err := filepath.Abs(requested)
|
|
if err != nil {
|
|
return "", errors.New("backup output path is unavailable")
|
|
}
|
|
requested = absolute
|
|
}
|
|
return filepath.Clean(requested), nil
|
|
}
|
|
home, err := homeDir()
|
|
if err != nil || home == "" {
|
|
return "", errors.New("home directory is unavailable for the default backup path")
|
|
}
|
|
name := fmt.Sprintf("thothii-%s-%s-%s.zip", installationID, now.UTC().Format("20060102T150405Z"), revision)
|
|
return filepath.Join(home, ".thothii", "backups", installationID, name), nil
|
|
}
|
|
|
|
func ensureNewArchivePath(output string) error {
|
|
if output == "" || !filepath.IsAbs(output) {
|
|
return errors.New("backup output path must be absolute")
|
|
}
|
|
if err := os.MkdirAll(filepath.Dir(output), 0o700); err != nil {
|
|
return fmt.Errorf("create backup directory: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// archiveReservation claims an output path without replacing an existing archive. Its file
|
|
// identity is retained so cleanup never removes a path another process took over.
|
|
type archiveReservation struct {
|
|
path string
|
|
file *os.File
|
|
info os.FileInfo
|
|
}
|
|
|
|
func reserveArchiveOutput(output string) (*archiveReservation, error) {
|
|
file, err := os.OpenFile(output, os.O_WRONLY|os.O_CREATE|os.O_EXCL, 0o600)
|
|
if errors.Is(err, os.ErrExist) {
|
|
return nil, errors.New("backup output already exists")
|
|
}
|
|
if err != nil {
|
|
return nil, fmt.Errorf("reserve backup output: %w", err)
|
|
}
|
|
info, statErr := file.Stat()
|
|
if statErr != nil {
|
|
_ = file.Close()
|
|
return nil, fmt.Errorf("inspect reserved backup output: %w", statErr)
|
|
}
|
|
return &archiveReservation{path: output, file: file, info: info}, nil
|
|
}
|
|
|
|
func publishReservedArchive(reservation *archiveReservation, temporaryPath string) error {
|
|
if reservation == nil || reservation.file == nil {
|
|
return errors.New("backup output reservation is unavailable")
|
|
}
|
|
temporary, err := os.Open(temporaryPath)
|
|
if err != nil {
|
|
return fmt.Errorf("open temporary backup archive: %w", err)
|
|
}
|
|
defer temporary.Close()
|
|
if err := reservation.file.Truncate(0); err != nil {
|
|
return fmt.Errorf("prepare reserved backup output: %w", err)
|
|
}
|
|
if _, err := reservation.file.Seek(0, io.SeekStart); err != nil {
|
|
return fmt.Errorf("seek reserved backup output: %w", err)
|
|
}
|
|
if _, err := io.Copy(reservation.file, temporary); err != nil {
|
|
return fmt.Errorf("write reserved backup output: %w", err)
|
|
}
|
|
if err := reservation.file.Sync(); err != nil {
|
|
return fmt.Errorf("fsync reserved backup output: %w", err)
|
|
}
|
|
if err := reservation.file.Close(); err != nil {
|
|
return fmt.Errorf("close reserved backup output: %w", err)
|
|
}
|
|
reservation.file = nil
|
|
current, err := os.Stat(reservation.path)
|
|
if err != nil || !os.SameFile(reservation.info, current) {
|
|
return errors.New("backup output ownership changed before publication")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// RemoveIfOwned removes the reservation only when the output path still names the exact file
|
|
// created by this invocation. It is safe to call after another process has claimed the path.
|
|
func (reservation *archiveReservation) RemoveIfOwned() error {
|
|
if reservation == nil {
|
|
return nil
|
|
}
|
|
if reservation.file != nil {
|
|
if err := reservation.file.Close(); err != nil {
|
|
return err
|
|
}
|
|
reservation.file = nil
|
|
}
|
|
current, err := os.Stat(reservation.path)
|
|
if errors.Is(err, os.ErrNotExist) {
|
|
return nil
|
|
}
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if !os.SameFile(reservation.info, current) {
|
|
return nil
|
|
}
|
|
return os.Remove(reservation.path)
|
|
}
|
|
|
|
type renderedCompose struct {
|
|
Volumes map[string]struct {
|
|
Name string `json:"name"`
|
|
} `json:"volumes"`
|
|
Services map[string]struct {
|
|
Image string `json:"image"`
|
|
ContainerName string `json:"container_name"`
|
|
} `json:"services"`
|
|
}
|
|
|
|
func renderedConfiguration(ctx context.Context, installation config.Installation, runner archiveRunner) (renderedCompose, error) {
|
|
result, err := runner.Run(ctx, installation.ComposeArgs("config", "--format", "json"), nil)
|
|
if err != nil {
|
|
return renderedCompose{}, dockerError("render Compose configuration", result, err)
|
|
}
|
|
var rendered renderedCompose
|
|
if json.Unmarshal([]byte(result.Stdout), &rendered) != nil {
|
|
return renderedCompose{}, errors.New("Docker Compose returned invalid backup configuration")
|
|
}
|
|
for _, logical := range requiredBackupVolumes(installation) {
|
|
if rendered.Volumes[logical].Name == "" {
|
|
return renderedCompose{}, fmt.Errorf("required backup volume %q is not configured", logical)
|
|
}
|
|
}
|
|
return rendered, nil
|
|
}
|
|
|
|
func inspectRequiredVolumes(ctx context.Context, installation config.Installation, runner archiveRunner, rendered renderedCompose) ([]VolumeMetadata, error) {
|
|
required := requiredBackupVolumes(installation)
|
|
names := make([]string, 0, len(required))
|
|
for _, logical := range required {
|
|
names = append(names, rendered.Volumes[logical].Name)
|
|
}
|
|
result, err := runner.Run(ctx, append([]string{"volume", "inspect"}, names...), nil)
|
|
if err != nil {
|
|
return nil, dockerError("inspect backup volumes", result, err)
|
|
}
|
|
var inspected []struct {
|
|
Name string `json:"Name"`
|
|
Driver string `json:"Driver"`
|
|
Labels map[string]string `json:"Labels"`
|
|
}
|
|
if json.Unmarshal([]byte(result.Stdout), &inspected) != nil {
|
|
return nil, errors.New("Docker returned invalid volume metadata")
|
|
}
|
|
byName := make(map[string]struct {
|
|
Name string
|
|
Driver string
|
|
Labels map[string]string
|
|
}, len(inspected))
|
|
for _, item := range inspected {
|
|
byName[item.Name] = struct {
|
|
Name string
|
|
Driver string
|
|
Labels map[string]string
|
|
}{item.Name, item.Driver, item.Labels}
|
|
}
|
|
volumes := make([]VolumeMetadata, 0, len(required))
|
|
for _, logical := range required {
|
|
name := rendered.Volumes[logical].Name
|
|
item, ok := byName[name]
|
|
if !ok || item.Driver == "" {
|
|
return nil, fmt.Errorf("required backup volume %q is unavailable", logical)
|
|
}
|
|
volumes = append(volumes, VolumeMetadata{LogicalName: logical, Name: name, Driver: item.Driver, Labels: item.Labels})
|
|
}
|
|
return volumes, nil
|
|
}
|
|
|
|
func imageIdentities(ctx context.Context, installation config.Installation, runner archiveRunner, rendered renderedCompose) ([]ImageIdentity, error) {
|
|
services := make([]string, 0, len(rendered.Services))
|
|
for name, definition := range rendered.Services {
|
|
if definition.Image != "" {
|
|
services = append(services, name)
|
|
}
|
|
}
|
|
sort.Strings(services)
|
|
images := make([]ImageIdentity, 0, len(services))
|
|
for _, name := range services {
|
|
result, err := runner.Run(ctx, installation.ComposeArgs("images", "--format", "json", name), nil)
|
|
if err != nil {
|
|
return nil, dockerError("inspect image identities", result, err)
|
|
}
|
|
inspected, err := decodeComposeImageIdentities(result.Stdout)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
expectedContainer := rendered.Services[name].ContainerName
|
|
if expectedContainer == "" {
|
|
expectedContainer = installation.ProjectName() + "-" + name + "-1"
|
|
}
|
|
id, err := selectComposeImageIdentity(name, expectedContainer, inspected)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
images = append(images, ImageIdentity{Service: name, Reference: rendered.Services[name].Image, ID: id})
|
|
}
|
|
return images, nil
|
|
}
|
|
|
|
type composeImageIdentity struct {
|
|
Service string `json:"Service"`
|
|
ContainerName string `json:"ContainerName"`
|
|
ID string `json:"ID"`
|
|
}
|
|
|
|
func decodeComposeImageIdentities(value string) ([]composeImageIdentity, error) {
|
|
trimmed := strings.TrimSpace(value)
|
|
if trimmed == "" || trimmed == "null" {
|
|
return nil, nil
|
|
}
|
|
var identities []composeImageIdentity
|
|
if strings.HasPrefix(trimmed, "[") {
|
|
if json.Unmarshal([]byte(trimmed), &identities) != nil {
|
|
return nil, errors.New("Docker Compose returned invalid image identities")
|
|
}
|
|
} else {
|
|
decoder := json.NewDecoder(strings.NewReader(trimmed))
|
|
for {
|
|
var item composeImageIdentity
|
|
err := decoder.Decode(&item)
|
|
if errors.Is(err, io.EOF) {
|
|
break
|
|
}
|
|
if err != nil {
|
|
return nil, errors.New("Docker Compose returned invalid image identities")
|
|
}
|
|
identities = append(identities, item)
|
|
}
|
|
}
|
|
for _, item := range identities {
|
|
if item.ID == "" || item.Service == "" && item.ContainerName == "" {
|
|
return nil, errors.New("Docker Compose returned invalid image identities")
|
|
}
|
|
}
|
|
return identities, nil
|
|
}
|
|
|
|
func selectComposeImageIdentity(service string, expectedContainer string, identities []composeImageIdentity) (string, error) {
|
|
if len(identities) == 0 {
|
|
return "", nil
|
|
}
|
|
matching := make([]composeImageIdentity, 0, len(identities))
|
|
for _, item := range identities {
|
|
if item.Service == service || item.Service == "" && item.ContainerName == expectedContainer {
|
|
matching = append(matching, item)
|
|
}
|
|
}
|
|
if len(matching) != 1 {
|
|
return "", errors.New("Docker Compose returned invalid image identities")
|
|
}
|
|
return matching[0].ID, nil
|
|
}
|
|
|
|
func installationRunning(ctx context.Context, installation config.Installation, runner archiveRunner) (bool, error) {
|
|
result, err := runner.Run(ctx, installation.ComposeArgs("ps", "--all", "--format", "json"), nil)
|
|
if err != nil {
|
|
return false, dockerError("inspect service states", result, err)
|
|
}
|
|
states, err := backupServiceStates(result.Stdout)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
running := false
|
|
coreRunning := false
|
|
for _, service := range states {
|
|
switch service.State {
|
|
case "exited", "dead":
|
|
continue
|
|
case "running":
|
|
running = true
|
|
if service.Service == "core" {
|
|
coreRunning = true
|
|
}
|
|
default:
|
|
return false, fmt.Errorf("service %q is %q and is not safely quiesced; stop the installation before backup", service.Service, service.State)
|
|
}
|
|
}
|
|
if running && !coreRunning {
|
|
return false, errors.New("core is not running while other services are active; stop the installation before backup")
|
|
}
|
|
return running, nil
|
|
}
|
|
|
|
type backupServiceState struct {
|
|
Service string `json:"Service"`
|
|
State string `json:"State"`
|
|
}
|
|
|
|
func backupServiceStates(value string) ([]backupServiceState, error) {
|
|
trimmed := strings.TrimSpace(value)
|
|
if trimmed == "" || trimmed == "[]" {
|
|
return nil, nil
|
|
}
|
|
var array []backupServiceState
|
|
if err := json.Unmarshal([]byte(trimmed), &array); err == nil {
|
|
return normalizeBackupServiceStates(array)
|
|
}
|
|
decoder := json.NewDecoder(strings.NewReader(trimmed))
|
|
var states []backupServiceState
|
|
for {
|
|
var state backupServiceState
|
|
err := decoder.Decode(&state)
|
|
if errors.Is(err, io.EOF) {
|
|
break
|
|
}
|
|
if err != nil {
|
|
return nil, errors.New("Docker Compose returned invalid service states")
|
|
}
|
|
states = append(states, state)
|
|
}
|
|
return normalizeBackupServiceStates(states)
|
|
}
|
|
|
|
func normalizeBackupServiceStates(states []backupServiceState) ([]backupServiceState, error) {
|
|
for index := range states {
|
|
states[index].Service = strings.TrimSpace(states[index].Service)
|
|
states[index].State = strings.ToLower(strings.TrimSpace(states[index].State))
|
|
if states[index].Service == "" || states[index].State == "" {
|
|
return nil, errors.New("Docker Compose returned invalid service states")
|
|
}
|
|
}
|
|
return states, nil
|
|
}
|
|
|
|
func maintenance(ctx context.Context, installation config.Installation, runner archiveRunner, activate bool) error {
|
|
action := "maintenance-deactivate"
|
|
if activate {
|
|
action = "maintenance-activate"
|
|
}
|
|
result, err := runner.Run(ctx, installation.ComposeArgs(
|
|
"exec", "-T", "core", "node", "/app/backend/dist/operator-command.js", action,
|
|
), nil)
|
|
if err != nil {
|
|
return dockerError("change maintenance admissions", result, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func waitForNoActiveSessions(ctx context.Context, installation config.Installation, runner archiveRunner, drain bool, sleep func(time.Duration)) error {
|
|
for attempt := 0; attempt < maxDrainPolls; attempt++ {
|
|
result, err := runner.Run(ctx, installation.ComposeArgs(
|
|
"exec", "-T", "core", "node", "/app/backend/dist/operator-command.js",
|
|
"session-inventory", runner.SessionInventoryScope(),
|
|
), nil)
|
|
if err != nil {
|
|
return dockerError("inspect active sessions", result, err)
|
|
}
|
|
active, err := activeSessions(result.Stdout)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if !active {
|
|
return nil
|
|
}
|
|
if !drain {
|
|
return ErrActiveSessions
|
|
}
|
|
select {
|
|
case <-ctx.Done():
|
|
return ctx.Err()
|
|
default:
|
|
sleep(drainPollInterval)
|
|
}
|
|
}
|
|
return errors.New("active sessions did not drain before the backup deadline")
|
|
}
|
|
|
|
func activeSessions(value string) (bool, error) {
|
|
var sessions []struct {
|
|
Status string `json:"status"`
|
|
Archived bool `json:"archived"`
|
|
}
|
|
if json.Unmarshal([]byte(value), &sessions) != nil {
|
|
return false, errors.New("session inventory is invalid")
|
|
}
|
|
for _, session := range sessions {
|
|
if !session.Archived && (session.Status == "" || session.Status == "running" || session.Status == "active") {
|
|
return true, nil
|
|
}
|
|
}
|
|
return false, nil
|
|
}
|
|
|
|
func composeStartAndVerify(ctx context.Context, installation config.Installation, runner archiveRunner) error {
|
|
if err := runCompose(ctx, installation, runner, "start"); err != nil {
|
|
return err
|
|
}
|
|
return service.WaitForHealthy(ctx, installation, runner)
|
|
}
|
|
|
|
func boundedCleanupContext() (context.Context, context.CancelFunc) {
|
|
return context.WithTimeout(context.Background(), cleanupOperationTimeout)
|
|
}
|
|
|
|
// retryBoundedCleanup retries an idempotent compensating mutation once. It preserves a lost
|
|
// response as part of the returned error while reporting whether the retry established the final
|
|
// state needed by the next cleanup action.
|
|
func retryBoundedCleanup(operation func(context.Context) error) (error, bool) {
|
|
run := func() error {
|
|
cleanupContext, cancel := boundedCleanupContext()
|
|
defer cancel()
|
|
return operation(cleanupContext)
|
|
}
|
|
firstErr := run()
|
|
if firstErr == nil {
|
|
return nil, true
|
|
}
|
|
retryErr := run()
|
|
return errors.Join(firstErr, retryErr), retryErr == nil
|
|
}
|
|
|
|
func runCompose(ctx context.Context, installation config.Installation, runner archiveRunner, command ...string) error {
|
|
result, err := runner.Run(ctx, installation.ComposeArgs(command...), nil)
|
|
if err != nil {
|
|
return dockerError("run Docker Compose", result, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func writeArchive(ctx context.Context, output string, reservation *archiveReservation, installation config.Installation, request CreateRequest, secretPaths, authenticationPaths []string, manifest Manifest, volumes []VolumeMetadata, dependencies dependencies) (resultErr error) {
|
|
directory := filepath.Dir(output)
|
|
temporary, err := os.CreateTemp(directory, ".tht-backup-*.tmp")
|
|
if err != nil {
|
|
return fmt.Errorf("create temporary backup archive: %w", err)
|
|
}
|
|
temporaryPath := temporary.Name()
|
|
closedFile := false
|
|
defer func() {
|
|
var closeErr error
|
|
if !closedFile {
|
|
closeErr = temporary.Close()
|
|
}
|
|
if closeErr != nil && resultErr == nil {
|
|
resultErr = closeErr
|
|
}
|
|
_ = os.Remove(temporaryPath)
|
|
}()
|
|
if err := temporary.Chmod(0o600); err != nil {
|
|
return fmt.Errorf("protect temporary backup archive: %w", err)
|
|
}
|
|
writer := zip.NewWriter(temporary)
|
|
closedWriter := false
|
|
defer func() {
|
|
if !closedWriter {
|
|
_ = writer.Close()
|
|
}
|
|
}()
|
|
|
|
add := func(name string, mode os.FileMode, input io.Reader) (Entry, error) {
|
|
header := &zip.FileHeader{Name: filepath.ToSlash(name), Method: zip.Deflate}
|
|
header.SetMode(mode)
|
|
created, err := writer.CreateHeader(header)
|
|
if err != nil {
|
|
return Entry{}, err
|
|
}
|
|
hash := sha256.New()
|
|
size, err := io.Copy(io.MultiWriter(created, hash), input)
|
|
if err != nil {
|
|
return Entry{}, err
|
|
}
|
|
return Entry{Path: filepath.ToSlash(name), SHA256: "sha256:" + hex.EncodeToString(hash.Sum(nil)), Size: size, Archived: true, Mode: uint32(mode.Perm())}, nil
|
|
}
|
|
addFile := func(name, owner string, sensitive bool, source string) error {
|
|
file, err := os.Open(source)
|
|
if err != nil {
|
|
return fmt.Errorf("read backup input: %w", err)
|
|
}
|
|
defer file.Close()
|
|
info, err := file.Stat()
|
|
if err != nil || !info.Mode().IsRegular() {
|
|
return errors.New("backup input is not a regular file")
|
|
}
|
|
entry, err := add(name, info.Mode(), file)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
entry.Kind, entry.Owner, entry.Sensitive = EntryFile, owner, sensitive
|
|
manifest.Entries = append(manifest.Entries, entry)
|
|
return nil
|
|
}
|
|
|
|
for _, input := range configurationInputs(installation) {
|
|
if input.optional {
|
|
if _, err := os.Stat(input.source); errors.Is(err, os.ErrNotExist) {
|
|
continue
|
|
}
|
|
}
|
|
if err := addFile(input.name, input.owner, false, input.source); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
for index, source := range secretPaths {
|
|
if request.IncludeSecrets {
|
|
if err := addFile(fmt.Sprintf("external-secrets/%03d", index), "external-secret", true, source); err != nil {
|
|
return err
|
|
}
|
|
entry := &manifest.Entries[len(manifest.Entries)-1]
|
|
entry.Kind, entry.SourcePath, entry.Sensitive = EntryExternalSecret, source, true
|
|
continue
|
|
}
|
|
file, err := os.Open(source)
|
|
if err != nil {
|
|
return errors.New("external secret file is unavailable")
|
|
}
|
|
info, statErr := file.Stat()
|
|
if statErr != nil || !info.Mode().IsRegular() {
|
|
_ = file.Close()
|
|
return errors.New("external secret file is unavailable")
|
|
}
|
|
hash := sha256.New()
|
|
size, copyErr := io.Copy(hash, file)
|
|
closeErr := file.Close()
|
|
if copyErr != nil || closeErr != nil {
|
|
return errors.New("external secret file is unavailable")
|
|
}
|
|
manifest.Entries = append(manifest.Entries, Entry{
|
|
Path: fmt.Sprintf("external-secrets/%03d", index), Kind: EntrySecretReference, Owner: "external-secret",
|
|
SourcePath: source, SHA256: "sha256:" + hex.EncodeToString(hash.Sum(nil)), Size: size, Sensitive: true,
|
|
})
|
|
}
|
|
for index, source := range authenticationPaths {
|
|
name := fmt.Sprintf("authentication-secrets/%03d-%s", index, filepath.Base(source))
|
|
if request.IncludeSecrets {
|
|
if err := addFile(name, "authentication-configuration", true, source); err != nil {
|
|
return err
|
|
}
|
|
entry := &manifest.Entries[len(manifest.Entries)-1]
|
|
entry.Kind, entry.SourcePath, entry.Sensitive = EntryExternalSecret, source, true
|
|
continue
|
|
}
|
|
if filepath.Base(source) != "auth.yaml" {
|
|
continue
|
|
}
|
|
if err := addAuthenticationReference(&manifest, name, source); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
for _, volume := range volumes {
|
|
headerName := "volumes/" + volume.LogicalName + ".tar"
|
|
header := &zip.FileHeader{Name: headerName, Method: zip.Deflate}
|
|
header.SetMode(0o600)
|
|
archiveWriter, err := writer.CreateHeader(header)
|
|
if err != nil {
|
|
return fmt.Errorf("create volume archive entry %s: %w", volume.LogicalName, err)
|
|
}
|
|
hash := sha256.New()
|
|
sizeWriter := &countingWriter{writer: io.MultiWriter(archiveWriter, hash)}
|
|
streamResult, streamErr := dependencies.runner.Stream(ctx, volumeArchiveCommand(volume.Name), nil, sizeWriter)
|
|
if streamErr != nil || streamResult.ExitCode != 0 {
|
|
if streamErr == nil {
|
|
streamErr = errors.New("Docker volume helper returned a nonzero exit status")
|
|
}
|
|
return fmt.Errorf("archive volume %s: %w", volume.LogicalName, dockerError("stream volume", streamResult, streamErr))
|
|
}
|
|
entry := Entry{
|
|
Path: headerName, Kind: EntryVolume, Owner: "volume:" + volume.LogicalName, LogicalName: volume.LogicalName,
|
|
SHA256: "sha256:" + hex.EncodeToString(hash.Sum(nil)), Size: sizeWriter.size, Mode: 0o600, Archived: true,
|
|
Sensitive: volume.LogicalName == "workspace-secrets",
|
|
}
|
|
manifest.Entries = append(manifest.Entries, entry)
|
|
}
|
|
if installation.Profile == "server" {
|
|
if err := archivePreservationRoots(installation, secretPaths, addFile); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
manifestJSON, err := manifest.JSON()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if _, err := add(ManifestPath, 0o600, bytes.NewReader(manifestJSON)); err != nil {
|
|
return err
|
|
}
|
|
if err := writer.Close(); err != nil {
|
|
return fmt.Errorf("close backup archive: %w", err)
|
|
}
|
|
closedWriter = true
|
|
if err := temporary.Sync(); err != nil {
|
|
return fmt.Errorf("fsync backup archive: %w", err)
|
|
}
|
|
if err := temporary.Close(); err != nil {
|
|
return fmt.Errorf("close backup archive: %w", err)
|
|
}
|
|
closedFile = true
|
|
if err := dependencies.publishReserved(reservation, temporaryPath); err != nil {
|
|
return fmt.Errorf("publish backup archive: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
type configurationInput struct {
|
|
name, owner, source string
|
|
optional bool
|
|
}
|
|
|
|
func configurationInputs(installation config.Installation) []configurationInput {
|
|
inputs := []configurationInput{
|
|
{"configuration/installation/thothii-installation.yaml", "installation", installation.Path, false},
|
|
{"configuration/environment/operator.env", "installation", installation.EnvFile, false},
|
|
{"configuration/pi/models.json", "pi", filepath.Join(installation.ProjectDirectory, "deploy", "pi", "models.json"), false},
|
|
{"configuration/pi/settings.json", "pi", filepath.Join(installation.ProjectDirectory, "deploy", "pi", "settings.json"), false},
|
|
{"configuration/generated/current-image.yaml", "installation", installation.CurrentImageOverridePath(), true},
|
|
}
|
|
for index, source := range installation.Overrides {
|
|
inputs = append(inputs, configurationInput{fmt.Sprintf("configuration/overrides/%02d-%s", index, filepath.Base(source)), "installation", source, false})
|
|
}
|
|
return inputs
|
|
}
|
|
|
|
const maxAuthenticationConfigurationBytes = 1 << 20
|
|
|
|
func authenticationConfigFiles(installation config.Installation) ([]string, error) {
|
|
directory := installation.AuthenticationDirectory()
|
|
if directory == "" {
|
|
return nil, nil
|
|
}
|
|
if err := safeio.ValidatePrivateDirectory(directory); err != nil {
|
|
return nil, errors.New("authentication configuration directory is unavailable or unsafe")
|
|
}
|
|
authPath := filepath.Join(directory, "auth.yaml")
|
|
contents, err := safeio.ReadCanonicalRegular(authPath, maxAuthenticationConfigurationBytes)
|
|
if err != nil {
|
|
return nil, errors.New("authentication configuration is unavailable or unsafe")
|
|
}
|
|
var configuration struct {
|
|
Mode string `yaml:"mode"`
|
|
}
|
|
if err := yaml.Unmarshal(contents, &configuration); err != nil || (configuration.Mode != "local" && configuration.Mode != "oidc") {
|
|
return nil, errors.New("authentication configuration is unavailable or invalid")
|
|
}
|
|
paths := []string{authPath}
|
|
if configuration.Mode == "local" {
|
|
usersPath := filepath.Join(directory, "users.yaml")
|
|
if _, err := safeio.ReadCanonicalRegular(usersPath, maxAuthenticationConfigurationBytes); err != nil {
|
|
return nil, errors.New("authentication user registry is unavailable or unsafe")
|
|
}
|
|
paths = append(paths, usersPath)
|
|
}
|
|
return paths, nil
|
|
}
|
|
|
|
func addAuthenticationReference(manifest *Manifest, name, source string) error {
|
|
contents, err := safeio.ReadCanonicalRegular(source, maxAuthenticationConfigurationBytes)
|
|
if err != nil {
|
|
return errors.New("authentication configuration is unavailable or unsafe")
|
|
}
|
|
digest := sha256.Sum256(contents)
|
|
manifest.Entries = append(manifest.Entries, Entry{
|
|
Path: name, Kind: EntrySecretReference, Owner: "authentication-configuration", SourcePath: source,
|
|
SHA256: "sha256:" + hex.EncodeToString(digest[:]), Size: int64(len(contents)), Sensitive: true,
|
|
})
|
|
return nil
|
|
}
|
|
|
|
func volumeArchiveCommand(volume string) []string {
|
|
return []string{"run", "--rm", "--network", "none", "--mount", "type=volume,src=" + volume + ",dst=/source,readonly", helperImage, "tar", "--numeric-owner", "-C", "/source", "-cf", "-", "."}
|
|
}
|
|
|
|
type countingWriter struct {
|
|
writer io.Writer
|
|
size int64
|
|
}
|
|
|
|
func (writer *countingWriter) Write(value []byte) (int, error) {
|
|
written, err := writer.writer.Write(value)
|
|
writer.size += int64(written)
|
|
return written, err
|
|
}
|
|
|
|
func archivePreservationRoots(installation config.Installation, secretPaths []string, addFile func(string, string, bool, string) error) error {
|
|
backupRoot, err := installation.EnvironmentValue("THT_BACKUP_ROOT")
|
|
if err != nil {
|
|
return errors.New("server preservation roots are unavailable")
|
|
}
|
|
secretSet := make(map[string]bool, len(secretPaths))
|
|
for _, secret := range secretPaths {
|
|
secretSet[secret] = true
|
|
}
|
|
for index, variable := range []string{"THT_DATA_ROOT", "THT_PI_STATE_ROOT", "THT_WORKSPACE_REGISTRY_ROOT"} {
|
|
root, err := installation.EnvironmentValue(variable)
|
|
if err != nil || root == "" {
|
|
return errors.New("server preservation roots are unavailable")
|
|
}
|
|
info, err := os.Stat(root)
|
|
if err != nil || !info.IsDir() {
|
|
return errors.New("server preservation root is unavailable")
|
|
}
|
|
authStateRoot := ""
|
|
if variable == "THT_DATA_ROOT" {
|
|
authStateRoot = filepath.Join(root, "auth")
|
|
}
|
|
err = filepath.WalkDir(root, func(path string, entry os.DirEntry, walkErr error) error {
|
|
if walkErr != nil {
|
|
return walkErr
|
|
}
|
|
if path == backupRoot {
|
|
if entry.IsDir() {
|
|
return filepath.SkipDir
|
|
}
|
|
return nil
|
|
}
|
|
if authStateRoot != "" && path == authStateRoot {
|
|
if entry.IsDir() {
|
|
return filepath.SkipDir
|
|
}
|
|
return nil
|
|
}
|
|
if entry.Type()&os.ModeSymlink != 0 {
|
|
return errors.New("server preservation root contains a symlink")
|
|
}
|
|
if entry.IsDir() || secretSet[path] {
|
|
return nil
|
|
}
|
|
if !entry.Type().IsRegular() {
|
|
return errors.New("server preservation root contains an unsupported entry")
|
|
}
|
|
relative, err := filepath.Rel(root, path)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
name := fmt.Sprintf("preservation/%02d-%s/%s", index, filepath.Base(root), filepath.ToSlash(relative))
|
|
return addFile(name, "preservation-root:"+variable, false, path)
|
|
})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func rejectInlineSecretValues(environmentPath string) error {
|
|
contents, err := os.ReadFile(environmentPath)
|
|
if err != nil {
|
|
return errors.New("installation environment could not be read")
|
|
}
|
|
for _, line := range strings.Split(string(contents), "\n") {
|
|
line = strings.TrimSpace(line)
|
|
if line == "" || strings.HasPrefix(line, "#") {
|
|
continue
|
|
}
|
|
line = strings.TrimPrefix(line, "export ")
|
|
key, value, found := strings.Cut(line, "=")
|
|
if !found || strings.TrimSpace(value) == "" {
|
|
continue
|
|
}
|
|
key = strings.ToUpper(strings.TrimSpace(key))
|
|
if strings.HasSuffix(key, "_FILE") || strings.HasSuffix(key, "_SOURCE") {
|
|
continue
|
|
}
|
|
if strings.Contains(key, "SECRET") || strings.Contains(key, "PASSWORD") || strings.Contains(key, "TOKEN") || strings.Contains(key, "API_KEY") || strings.Contains(key, "CREDENTIAL") || strings.Contains(key, "AUTHORIZATION") || strings.HasSuffix(key, "_AUTH") {
|
|
return errors.New("installation environment contains an inline secret value")
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func dockerError(action string, result compose.Result, err error) error {
|
|
if result.ExitCode != 0 {
|
|
return fmt.Errorf("%s: Docker exited with status %d", action, result.ExitCode)
|
|
}
|
|
if err != nil {
|
|
return fmt.Errorf("%s: %w", action, err)
|
|
}
|
|
return fmt.Errorf("%s failed", action)
|
|
}
|