// 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 } transaction, err := lifecycle.AcquireTransaction(installation) if err != nil { return Result{}, err } defer func() { if releaseErr := transaction.Release(); releaseErr != nil { result = Result{} resultErr = errors.Join(resultErr, fmt.Errorf("release backup lifecycle lock: %w", releaseErr)) } }() return createWithDependenciesTransaction(ctx, transaction, installation, request, dependencies) } // createWithDependenciesTransaction performs backup creation only with an active, opaque // installation-bound lifecycle capability. Restore passes the capability it acquired for the // enclosing transaction, so a recovery checkpoint cannot run without the same lifecycle lock. func createWithDependenciesTransaction(ctx context.Context, transaction *lifecycle.Transaction, installation config.Installation, request CreateRequest, dependencies dependencies) (result Result, resultErr error) { if err := transaction.Verify(installation); err != nil { return Result{}, err } 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 } 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 } if request.IncludeSecrets { if err := safeio.ProtectPrivateRegular(reservation.path); err != nil { _ = reservation.RemoveIfOwned() return Result{}, errors.New("protect secret-bearing backup output") } } 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) }