Files
ThothII/tools/tht/internal/backup/create_test.go
T

703 lines
27 KiB
Go

package backup
import (
"archive/tar"
"archive/zip"
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"io"
"os"
"path/filepath"
"sort"
"strings"
"testing"
"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"
)
var requiredTestVolumes = []string{"settings", "pi-state", "workspace-registry", "workspace-secrets", "sessions", "qdrant-data", "embedding-models"}
func TestCreateWritesManifestLastWithConfigurationMetadataAndSevenVolumes(t *testing.T) {
fixture := newBackupFixture(t, "local")
output := filepath.Join(t.TempDir(), "custom.zip")
runner := newBackupRunner(fixture.installation, false)
result, err := createWithDependencies(context.Background(), fixture.installation, CreateRequest{Output: output}, testDependencies(t, runner))
if err != nil {
t.Fatal(err)
}
if result.Path != output {
t.Fatalf("Result.Path = %q, want %q", result.Path, output)
}
archive := readFixtureArchive(t, output)
if got := archive.order[len(archive.order)-1]; got != ManifestPath {
t.Fatalf("last archive entry = %q, want %q", got, ManifestPath)
}
if archive.manifest.InstallationID != fixture.installationID || archive.manifest.SourceRevision != testRevision || archive.manifest.ComposeProject != fixture.installation.ProjectName() {
t.Fatalf("manifest identity = %#v", archive.manifest)
}
if len(archive.manifest.Volumes) != len(requiredTestVolumes) {
t.Fatalf("manifest volumes = %d, want %d", len(archive.manifest.Volumes), len(requiredTestVolumes))
}
for _, logical := range requiredTestVolumes {
path := "volumes/" + logical + ".tar"
if _, exists := archive.files[path]; !exists {
t.Errorf("archive is missing %s", path)
}
}
for _, path := range []string{
"configuration/installation/thothii-installation.yaml",
"configuration/environment/operator.env",
"configuration/pi/models.json",
"configuration/pi/settings.json",
"configuration/generated/current-image.yaml",
} {
if _, exists := archive.files[path]; !exists {
t.Errorf("archive is missing %s", path)
}
}
if len(archive.manifest.Images) < 2 {
t.Fatalf("image identities = %#v", archive.manifest.Images)
}
if runner.streamWhileRunning {
t.Fatal("a volume was streamed before the installation was stopped")
}
}
func TestCreateRestartsAndVerifiesAnInstallationThatWasRunning(t *testing.T) {
fixture := newBackupFixture(t, "local")
runner := newBackupRunner(fixture.installation, true)
output := filepath.Join(t.TempDir(), "running.zip")
if _, err := createWithDependencies(context.Background(), fixture.installation, CreateRequest{Output: output}, testDependencies(t, runner)); err != nil {
t.Fatal(err)
}
if !runner.running || !runner.coreRunning || runner.maintenance {
t.Fatalf("running state was not restored: running=%v core=%v maintenance=%v", runner.running, runner.coreRunning, runner.maintenance)
}
if runner.stopCount != 1 || runner.startCount != 1 || runner.healthChecks == 0 {
t.Fatalf("lifecycle counts: stop=%d start=%d health=%d", runner.stopCount, runner.startCount, runner.healthChecks)
}
}
func TestCreateRefusesActiveSessionsWithoutDrainAndRestoresAdmissions(t *testing.T) {
fixture := newBackupFixture(t, "local")
runner := newBackupRunner(fixture.installation, true)
runner.sessionResponses = []string{activeSessionPayload()}
output := filepath.Join(t.TempDir(), "refused.zip")
_, err := createWithDependencies(context.Background(), fixture.installation, CreateRequest{Output: output}, testDependencies(t, runner))
if !errors.Is(err, ErrActiveSessions) {
t.Fatalf("Create() error = %v, want ErrActiveSessions", err)
}
if runner.stopCount != 0 || !runner.running || runner.maintenance {
t.Fatalf("refusal changed lifecycle state: stop=%d running=%v maintenance=%v", runner.stopCount, runner.running, runner.maintenance)
}
if _, statErr := os.Stat(output); !errors.Is(statErr, os.ErrNotExist) {
t.Fatalf("unsafe backup was published: %v", statErr)
}
}
func TestCreateDrainsActiveSessionsBeforeStoppedSnapshot(t *testing.T) {
fixture := newBackupFixture(t, "local")
runner := newBackupRunner(fixture.installation, true)
runner.sessionResponses = []string{activeSessionPayload(), activeSessionPayload(), `[]`}
dependencies := testDependencies(t, runner)
sleeps := 0
dependencies.sleep = func(time.Duration) { sleeps++ }
if _, err := createWithDependencies(context.Background(), fixture.installation, CreateRequest{Output: filepath.Join(t.TempDir(), "drained.zip"), Drain: true}, dependencies); err != nil {
t.Fatal(err)
}
if sleeps == 0 || runner.streamWhileRunning {
t.Fatalf("drain did not wait for a stopped snapshot: sleeps=%d streamedWhileRunning=%v", sleeps, runner.streamWhileRunning)
}
}
func TestCreateDefaultPathUsesHomeInstallationIDUTCAndSourceRevision(t *testing.T) {
fixture := newBackupFixture(t, "local")
runner := newBackupRunner(fixture.installation, false)
dependencies := testDependencies(t, runner)
home := t.TempDir()
dependencies.homeDir = func() (string, error) { return home, nil }
result, err := createWithDependencies(context.Background(), fixture.installation, CreateRequest{}, dependencies)
if err != nil {
t.Fatal(err)
}
wantDirectory := filepath.Join(home, ".thothii", "backups", fixture.installationID)
if filepath.Dir(result.Path) != wantDirectory {
t.Fatalf("default directory = %q, want %q", filepath.Dir(result.Path), wantDirectory)
}
name := filepath.Base(result.Path)
for _, fragment := range []string{"20260816T081112Z", testRevision} {
if !strings.Contains(name, fragment) {
t.Errorf("default archive name %q does not contain %q", name, fragment)
}
}
}
func TestCreateCleansIncompleteArchiveAndRestoresRunningStateAfterStreamFailure(t *testing.T) {
fixture := newBackupFixture(t, "local")
runner := newBackupRunner(fixture.installation, true)
runner.failStreamAt = 3
directory := t.TempDir()
output := filepath.Join(directory, "failed.zip")
if _, err := createWithDependencies(context.Background(), fixture.installation, CreateRequest{Output: output}, testDependencies(t, runner)); err == nil {
t.Fatal("Create() succeeded after a volume stream failure")
}
if !runner.running || !runner.coreRunning || runner.maintenance {
t.Fatalf("running state was not recovered: running=%v core=%v maintenance=%v", runner.running, runner.coreRunning, runner.maintenance)
}
assertNoBackupArtifacts(t, directory)
}
func TestCreateCleansIncompleteArchiveWhenVolumeHelperReportsANonzeroExit(t *testing.T) {
fixture := newBackupFixture(t, "local")
runner := newBackupRunner(fixture.installation, false)
runner.exitOnlyAt = 3
directory := t.TempDir()
if _, err := createWithDependencies(context.Background(), fixture.installation, CreateRequest{Output: filepath.Join(directory, "failed.zip")}, testDependencies(t, runner)); err == nil {
t.Fatal("Create() succeeded after a nonzero volume-helper exit")
}
assertNoBackupArtifacts(t, directory)
}
func TestCreateCleansIncompleteArchiveAfterAtomicPublishFailure(t *testing.T) {
fixture := newBackupFixture(t, "local")
runner := newBackupRunner(fixture.installation, false)
directory := t.TempDir()
dependencies := testDependencies(t, runner)
dependencies.publishReserved = func(*archiveReservation, string) error { return errors.New("publish failed") }
if _, err := createWithDependencies(context.Background(), fixture.installation, CreateRequest{Output: filepath.Join(directory, "failed.zip")}, dependencies); err == nil {
t.Fatal("Create() succeeded after publish failure")
}
assertNoBackupArtifacts(t, directory)
}
func TestCreateDoesNotOverwriteOrDeleteAnOutputCreatedBeforeReservation(t *testing.T) {
fixture := newBackupFixture(t, "local")
runner := newBackupRunner(fixture.installation, false)
directory := t.TempDir()
output := filepath.Join(directory, "race.zip")
dependencies := testDependencies(t, runner)
reserve := dependencies.reserveOutput
dependencies.reserveOutput = func(path string) (*archiveReservation, error) {
if err := os.WriteFile(path, []byte("created-by-another-process"), 0o600); err != nil {
return nil, err
}
return reserve(path)
}
if _, err := createWithDependencies(context.Background(), fixture.installation, CreateRequest{Output: output}, dependencies); err == nil {
t.Fatal("Create() succeeded after another process claimed the output path")
}
contents, err := os.ReadFile(output)
if err != nil {
t.Fatalf("racing output was removed: %v", err)
}
if got := string(contents); got != "created-by-another-process" {
t.Fatalf("racing output = %q, want unchanged content", got)
}
entries, err := os.ReadDir(directory)
if err != nil {
t.Fatal(err)
}
if len(entries) != 1 || entries[0].Name() != "race.zip" {
t.Fatalf("temporary backup artifacts remain after race: %v", entries)
}
}
func TestCreateExcludesExternalSecretPayloadsByDefaultButRecordsDigests(t *testing.T) {
fixture := newBackupFixture(t, "local")
runner := newBackupRunner(fixture.installation, false)
output := filepath.Join(t.TempDir(), "no-external-secrets.zip")
if _, err := createWithDependencies(context.Background(), fixture.installation, CreateRequest{Output: output}, testDependencies(t, runner)); err != nil {
t.Fatal(err)
}
archive := readFixtureArchive(t, output)
all := bytes.Join(mapValues(archive.files), nil)
if bytes.Contains(all, []byte(fixture.secretValue)) {
t.Fatal("default archive contains an external secret value")
}
secretReferences := 0
workspaceSecrets := false
for _, entry := range archive.manifest.Entries {
if entry.Kind == EntrySecretReference {
secretReferences++
if entry.Archived || entry.SourcePath == "" || entry.SHA256 != DigestBytes([]byte(fixture.secretValue)) {
t.Fatalf("secret reference = %#v", entry)
}
}
if entry.Path == "volumes/workspace-secrets.tar" {
workspaceSecrets = entry.Archived
}
}
if secretReferences != 1 || !workspaceSecrets || archive.manifest.IncludesSecrets {
t.Fatalf("secret behavior: references=%d workspace=%v includes=%v", secretReferences, workspaceSecrets, archive.manifest.IncludesSecrets)
}
if strings.Contains(strings.Join(runner.calls, "\n"), fixture.secretValue) {
t.Fatal("external secret value appeared in a process argument")
}
}
func TestCreateRequiresConfirmationToIncludeSecretsAndUsesOwnerOnlyMode(t *testing.T) {
fixture := newBackupFixture(t, "local")
output := filepath.Join(t.TempDir(), "with-secrets.zip")
request := CreateRequest{Output: output, IncludeSecrets: true}
if _, err := createWithDependencies(context.Background(), fixture.installation, request, testDependencies(t, newBackupRunner(fixture.installation, false))); !errors.Is(err, ErrConfirmationRequired) {
t.Fatalf("Create() error = %v, want ErrConfirmationRequired", err)
}
request.Confirm = true
result, err := createWithDependencies(context.Background(), fixture.installation, request, testDependencies(t, newBackupRunner(fixture.installation, false)))
if err != nil {
t.Fatal(err)
}
if result.Warning == "" || strings.Contains(result.Warning, fixture.secretValue) {
t.Fatalf("custody warning = %q", result.Warning)
}
info, err := os.Stat(output)
if err != nil {
t.Fatal(err)
}
if got := info.Mode().Perm(); got != 0o600 {
t.Fatalf("archive mode = %#o, want 0600", got)
}
archive := readFixtureArchive(t, output)
if !archive.manifest.IncludesSecrets || !bytes.Contains(bytes.Join(mapValues(archive.files), nil), []byte(fixture.secretValue)) {
t.Fatal("confirmed archive does not include the external secret payload")
}
}
func TestCreateRejectsInlineSecretValuesBeforeWritingAnArchive(t *testing.T) {
fixture := newBackupFixture(t, "local")
fixture.environment = append(fixture.environment, "THT_LLM_API_KEY=must-never-enter-an-archive")
fixture.writeEnvironment(t)
directory := t.TempDir()
runner := newBackupRunner(fixture.installation, false)
_, err := createWithDependencies(context.Background(), fixture.installation, CreateRequest{Output: filepath.Join(directory, "unsafe.zip")}, testDependencies(t, runner))
if err == nil || !strings.Contains(err.Error(), "inline secret") || strings.Contains(err.Error(), "must-never") {
t.Fatalf("Create() error = %v, want value-free inline-secret refusal", err)
}
if len(runner.calls) != 0 {
t.Fatalf("Docker was called before unsafe environment refusal: %v", runner.calls)
}
assertNoBackupArtifacts(t, directory)
}
func TestCreateRejectsInlineAuthorizationValuesBeforeWritingAnArchive(t *testing.T) {
fixture := newBackupFixture(t, "local")
fixture.environment = append(fixture.environment, "DWH_AUTHORIZATION=must-never-enter-an-archive")
fixture.writeEnvironment(t)
directory := t.TempDir()
runner := newBackupRunner(fixture.installation, false)
_, err := createWithDependencies(context.Background(), fixture.installation, CreateRequest{Output: filepath.Join(directory, "unsafe.zip")}, testDependencies(t, runner))
if err == nil || !strings.Contains(err.Error(), "inline secret") || strings.Contains(err.Error(), "must-never") {
t.Fatalf("Create() error = %v, want value-free inline-secret refusal", err)
}
if len(runner.calls) != 0 {
t.Fatalf("Docker was called before unsafe environment refusal: %v", runner.calls)
}
assertNoBackupArtifacts(t, directory)
}
func TestCreateIncludesServerPreservationRootsWithoutRecursingIntoBackupRoot(t *testing.T) {
fixture := newBackupFixture(t, "server")
for _, item := range []struct{ variable, name, value string }{
{"THT_DATA_ROOT", "data", "session-state"},
{"THT_PI_STATE_ROOT", "pi-state", "pi-state"},
{"THT_WORKSPACE_REGISTRY_ROOT", "registry", "workspace-registry"},
} {
root := filepath.Join(fixture.root, item.name)
if err := os.MkdirAll(root, 0o700); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(filepath.Join(root, "payload"), []byte(item.value), 0o600); err != nil {
t.Fatal(err)
}
fixture.environment = append(fixture.environment, item.variable+"="+root)
}
backupRoot := filepath.Join(fixture.root, "existing-backups")
if err := os.MkdirAll(backupRoot, 0o700); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(filepath.Join(backupRoot, "old-secret-backup"), []byte("must-not-be-recursed"), 0o600); err != nil {
t.Fatal(err)
}
fixture.environment = append(fixture.environment, "THT_BACKUP_ROOT="+backupRoot)
fixture.writeEnvironment(t)
output := filepath.Join(t.TempDir(), "server.zip")
if _, err := createWithDependencies(context.Background(), fixture.installation, CreateRequest{Output: output}, testDependencies(t, newBackupRunner(fixture.installation, false))); err != nil {
t.Fatal(err)
}
archive := readFixtureArchive(t, output)
all := bytes.Join(mapValues(archive.files), nil)
for _, value := range []string{"session-state", "pi-state", "workspace-registry"} {
if !bytes.Contains(all, []byte(value)) {
t.Errorf("server preservation payload %q is missing", value)
}
}
if bytes.Contains(all, []byte("must-not-be-recursed")) {
t.Fatal("backup destination root was recursively included")
}
}
func TestCreateHonorsTheSharedInstallationLifecycleLock(t *testing.T) {
fixture := newBackupFixture(t, "local")
lock, err := lifecycle.Acquire(fixture.installation)
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = lock.Release() })
runner := newBackupRunner(fixture.installation, false)
_, err = createWithDependencies(context.Background(), fixture.installation, CreateRequest{Output: filepath.Join(t.TempDir(), "locked.zip")}, testDependencies(t, runner))
if !errors.Is(err, lifecycle.ErrLocked) {
t.Fatalf("Create() error = %v, want lifecycle.ErrLocked", err)
}
if len(runner.calls) != 0 {
t.Fatalf("Docker runner was called while lock was held: %v", runner.calls)
}
}
func TestCreateRefusesMutableNonRunningServiceStates(t *testing.T) {
for _, state := range []string{"paused", "restarting", "created", "removing"} {
t.Run(state, func(t *testing.T) {
fixture := newBackupFixture(t, "local")
runner := newBackupRunner(fixture.installation, false)
runner.serviceStates = map[string]string{"core": "running", "frontend": state}
dependencies := testDependencies(t, runner)
dependencies.sleep = func(time.Duration) { t.Fatal("unsafe service state reached the drain loop") }
_, err := createWithDependencies(context.Background(), fixture.installation, CreateRequest{Output: filepath.Join(t.TempDir(), "unsafe.zip")}, dependencies)
if err == nil || !strings.Contains(err.Error(), "not safely quiesced") {
t.Fatalf("Create() error = %v, want unsafe service-state refusal", err)
}
if runner.stopCount != 0 || runner.streams != 0 {
t.Fatalf("unsafe state was not refused before snapshot: stops=%d streams=%d", runner.stopCount, runner.streams)
}
})
}
}
type backupFixture struct {
root string
installationID string
installation config.Installation
environment []string
secretValue string
}
func newBackupFixture(t *testing.T, profile string) *backupFixture {
t.Helper()
temporaryRoot, err := filepath.EvalSymlinks(os.TempDir())
if err != nil {
t.Fatal(err)
}
root, err := os.MkdirTemp(temporaryRoot, "tht-backup-")
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = os.RemoveAll(root) })
for _, directory := range []string{filepath.Join(root, "deploy", "pi"), filepath.Join(root, "deploy", "fixture")} {
if err := os.MkdirAll(directory, 0o700); err != nil {
t.Fatal(err)
}
}
write := func(path, contents string) {
if err := os.WriteFile(path, []byte(contents), 0o600); err != nil {
t.Fatal(err)
}
}
write(filepath.Join(root, "compose.yaml"), "services: {}\n")
write(filepath.Join(root, "deploy", "compose."+profile+".yaml"), "services: {}\n")
write(filepath.Join(root, "deploy", "pi", "models.json"), `{"providers":{}}`)
write(filepath.Join(root, "deploy", "pi", "settings.json"), `{"enabledModels":[]}`)
override := filepath.Join(root, "deploy", "fixture", "extra.yaml")
write(override, "services: {}\n")
installationID := profile + "-fixture"
descriptorDirectory := filepath.Join(root, "deploy", installationID)
if err := os.MkdirAll(descriptorDirectory, 0o700); err != nil {
t.Fatal(err)
}
descriptor := filepath.Join(descriptorDirectory, "thothii-installation.yaml")
environment := filepath.Join(descriptorDirectory, "operator.env")
write(descriptor, "profile: "+profile+"\nprojectDirectory: "+root+"\nenvFile: "+environment+"\n")
secretValue := "external-secret-value-for-backup-test"
secretDirectory, err := os.MkdirTemp(temporaryRoot, "tht-backup-secret-")
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = os.RemoveAll(secretDirectory) })
secretPath := filepath.Join(secretDirectory, "external-secret")
write(secretPath, secretValue)
fixture := &backupFixture{
root: root,
installationID: installationID,
secretValue: secretValue,
environment: []string{
"THT_WORKSPACE_INSTALLATION_ID=" + installationID,
"THT_SECRETS_FILE=" + secretPath,
"THT_WORKSPACE_GIT_REMOTE=https://git.example.invalid/workspaces.git",
"THT_WORKSPACE_GIT_BRANCH=main",
},
installation: config.Installation{
Path: descriptor, Profile: profile, ProjectDirectory: root, EnvFile: environment,
Overrides: []string{override},
},
}
fixture.writeEnvironment(t)
currentImage := fixture.installation.CurrentImageOverridePath()
if err := os.MkdirAll(filepath.Dir(currentImage), 0o700); err != nil {
t.Fatal(err)
}
write(currentImage, "services:\n core:\n image: core:current\n")
return fixture
}
func (fixture *backupFixture) writeEnvironment(t *testing.T) {
t.Helper()
if err := os.WriteFile(fixture.installation.EnvFile, []byte(strings.Join(fixture.environment, "\n")+"\n"), 0o600); err != nil {
t.Fatal(err)
}
}
type fakeBackupRunner struct {
installation config.Installation
running bool
coreRunning bool
maintenance bool
sessionResponses []string
calls []string
streams int
failStreamAt int
exitOnlyAt int
streamWhileRunning bool
stopCount int
startCount int
healthChecks int
serviceStates map[string]string
}
func newBackupRunner(installation config.Installation, running bool) *fakeBackupRunner {
return &fakeBackupRunner{installation: installation, running: running, coreRunning: running}
}
func (runner *fakeBackupRunner) Run(_ context.Context, args []string, _ io.Reader) (compose.Result, error) {
command := strings.Join(args, " ")
runner.calls = append(runner.calls, command)
switch {
case strings.Contains(command, " config --format json"):
volumes := map[string]map[string]string{}
for _, logical := range requiredTestVolumes {
volumes[logical] = map[string]string{"name": runner.installation.ProjectName() + "_" + logical}
}
payload := map[string]any{
"volumes": volumes,
"services": map[string]any{
"core": map[string]any{"image": "thothii-core:test"},
"frontend": map[string]any{"image": "thothii-frontend:test"},
},
}
encoded, _ := json.Marshal(payload)
return compose.Result{Stdout: string(encoded)}, nil
case strings.HasPrefix(command, "volume inspect "):
items := make([]map[string]any, 0, len(requiredTestVolumes))
for _, logical := range requiredTestVolumes {
items = append(items, map[string]any{
"Name": runner.installation.ProjectName() + "_" + logical,
"Driver": "local",
"Labels": map[string]string{"com.docker.compose.project": runner.installation.ProjectName(), "com.docker.compose.volume": logical},
})
}
encoded, _ := json.Marshal(items)
return compose.Result{Stdout: string(encoded)}, nil
case strings.Contains(command, " images --format json"):
return compose.Result{Stdout: "{\"Service\":\"core\",\"Repository\":\"thothii-core\",\"Tag\":\"test\",\"ID\":\"sha256:core\"}\n{\"Service\":\"frontend\",\"Repository\":\"thothii-frontend\",\"Tag\":\"test\",\"ID\":\"sha256:frontend\"}\n"}, nil
case strings.Contains(command, " ps --all --format json"):
runner.healthChecks++
states := runner.serviceStates
if states == nil {
if runner.running {
return compose.Result{Stdout: healthyServicesPayload()}, nil
}
return compose.Result{}, nil
}
services := make([]string, 0, len(states))
for service := range states {
services = append(services, service)
}
sort.Strings(services)
lines := make([]string, 0, len(services))
for _, service := range services {
lines = append(lines, fmt.Sprintf(`{"Service":%q,"State":%q}`, service, states[service]))
}
return compose.Result{Stdout: strings.Join(lines, "\n")}, nil
case strings.Contains(command, "/internal/maintenance/status"):
return compose.Result{Stdout: fmt.Sprintf(`{"active":%t,"admissions":0,"recoveryRequired":false}`, runner.maintenance)}, nil
case strings.Contains(command, "/internal/maintenance/activate"):
runner.maintenance = true
return compose.Result{Stdout: `{"active":true,"admissions":0,"recoveryRequired":false}`}, nil
case strings.Contains(command, "/internal/maintenance/deactivate"):
runner.maintenance = false
return compose.Result{Stdout: `{"active":false,"admissions":0,"recoveryRequired":false}`}, nil
case strings.Contains(command, "/sessions?scope="):
if len(runner.sessionResponses) == 0 {
return compose.Result{Stdout: `[]`}, nil
}
response := runner.sessionResponses[0]
if len(runner.sessionResponses) > 1 {
runner.sessionResponses = runner.sessionResponses[1:]
}
return compose.Result{Stdout: response}, nil
case strings.HasSuffix(command, " stop"):
runner.stopCount++
runner.running, runner.coreRunning = false, false
return compose.Result{}, nil
case strings.HasSuffix(command, " start"):
runner.startCount++
runner.running, runner.coreRunning = true, true
return compose.Result{}, nil
case strings.Contains(command, " ps --all --format json"):
runner.healthChecks++
return compose.Result{Stdout: healthyServicesPayload()}, nil
default:
return compose.Result{}, fmt.Errorf("unexpected fake Docker command: %s", command)
}
}
func (runner *fakeBackupRunner) Stream(_ context.Context, args []string, _ io.Reader, stdout io.Writer) (compose.Result, error) {
runner.calls = append(runner.calls, strings.Join(args, " "))
runner.streams++
if runner.running {
runner.streamWhileRunning = true
}
writer := tar.NewWriter(stdout)
payload := []byte(fmt.Sprintf("volume-%d", runner.streams))
if err := writer.WriteHeader(&tar.Header{Name: "payload", Mode: 0o600, Size: int64(len(payload))}); err != nil {
return compose.Result{}, err
}
if _, err := writer.Write(payload); err != nil {
return compose.Result{}, err
}
if err := writer.Close(); err != nil {
return compose.Result{}, err
}
if runner.failStreamAt == runner.streams {
return compose.Result{ExitCode: 1}, errors.New("fixture volume stream failed")
}
if runner.exitOnlyAt == runner.streams {
return compose.Result{ExitCode: 1}, nil
}
return compose.Result{}, nil
}
func (runner *fakeBackupRunner) SessionInventoryScope() string {
if runner.installation.Profile == "local" {
return "mine"
}
return "all"
}
func testDependencies(t *testing.T, runner archiveRunner) dependencies {
t.Helper()
return dependencies{
runner: runner,
now: func() time.Time { return time.Date(2026, 8, 16, 8, 11, 12, 0, time.UTC) },
homeDir: func() (string, error) { return t.TempDir(), nil },
revision: func(context.Context, string) (string, error) { return testRevision, nil },
sleep: func(time.Duration) {},
reserveOutput: reserveArchiveOutput,
publishReserved: publishReservedArchive,
}
}
type fixtureArchive struct {
order []string
files map[string][]byte
manifest Manifest
}
func readFixtureArchive(t *testing.T, path string) fixtureArchive {
t.Helper()
reader, err := zip.OpenReader(path)
if err != nil {
t.Fatal(err)
}
defer reader.Close()
result := fixtureArchive{files: make(map[string][]byte)}
for _, file := range reader.File {
result.order = append(result.order, file.Name)
opened, err := file.Open()
if err != nil {
t.Fatal(err)
}
contents, err := io.ReadAll(opened)
closeErr := opened.Close()
if err != nil || closeErr != nil {
t.Fatalf("read %s: %v / %v", file.Name, err, closeErr)
}
result.files[file.Name] = contents
}
if err := json.Unmarshal(result.files[ManifestPath], &result.manifest); err != nil {
t.Fatalf("decode manifest: %v", err)
}
return result
}
func mapValues(values map[string][]byte) [][]byte {
keys := make([]string, 0, len(values))
for key := range values {
keys = append(keys, key)
}
sort.Strings(keys)
result := make([][]byte, 0, len(keys))
for _, key := range keys {
result = append(result, values[key])
}
return result
}
func activeSessionPayload() string {
return `[{"status":"running","archived":false}]`
}
func healthyServicesPayload() string {
return strings.Join([]string{
`{"Service":"core","State":"running","Health":"healthy"}`,
`{"Service":"frontend","State":"running","Health":"healthy"}`,
`{"Service":"qdrant","State":"running","Health":"healthy"}`,
`{"Service":"embedding","State":"running","Health":"healthy"}`,
`{"Service":"embedding-model-init","State":"exited","ExitCode":0}`,
}, "\n")
}
func assertNoBackupArtifacts(t *testing.T, directory string) {
t.Helper()
entries, err := os.ReadDir(directory)
if err != nil {
t.Fatal(err)
}
if len(entries) != 0 {
names := make([]string, 0, len(entries))
for _, entry := range entries {
names = append(names, entry.Name())
}
t.Fatalf("incomplete backup artifacts remain: %v", names)
}
}