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 TestDecodeComposeImageIdentitiesAcceptsNonEmptyArrayAndStreamingJSON(t *testing.T) { for name, input := range map[string]string{ "array": `[{"ContainerName":"project-core-1","ID":"sha256:core"},{"ContainerName":"project-frontend-1","ID":"sha256:frontend"}]`, "streaming": "{\"Service\":\"core\",\"ID\":\"sha256:core\"}\n{\"Service\":\"frontend\",\"ID\":\"sha256:frontend\"}\n", } { t.Run(name, func(t *testing.T) { identities, err := decodeComposeImageIdentities(input) if err != nil { t.Fatal(err) } if len(identities) != 2 || identities[0].ID != "sha256:core" || identities[1].ID != "sha256:frontend" { t.Fatalf("image identities = %#v", identities) } }) } } func TestDecodeComposeImageIdentitiesAcceptsNoContainerForProfiledService(t *testing.T) { for name, input := range map[string]string{ "empty output": "", "empty array": "[]", "null": "null", } { t.Run(name, func(t *testing.T) { identities, err := decodeComposeImageIdentities(input) if err != nil || len(identities) != 0 { t.Fatalf("decodeComposeImageIdentities(%q) = %#v, %v", input, identities, err) } }) } } func TestDecodeComposeImageIdentitiesRejectsMalformedOrIncompleteOutput(t *testing.T) { for name, input := range map[string]string{ "malformed JSON": "[", "scalar JSON": "true", "empty object": "{}", "null array element": "[null]", "missing image ID": `[{"Service":"core"}]`, "trailing document": "null\n{}", "trailing malformed bytes": `null garbage`, } { t.Run(name, func(t *testing.T) { if identities, err := decodeComposeImageIdentities(input); err == nil { t.Fatalf("decodeComposeImageIdentities(%q) = %#v, nil; want error", input, identities) } }) } } func TestSelectComposeImageIdentityUsesExactComposeContainerWhenServiceIsAbsent(t *testing.T) { identities := []composeImageIdentity{ {ContainerName: "project-core-1", ID: "sha256:core"}, {ContainerName: "project-llm", ID: "sha256:shared-image"}, } id, err := selectComposeImageIdentity("core", "project-core-1", identities) if err != nil { t.Fatal(err) } if id != "sha256:core" { t.Fatalf("selected image ID = %q, want sha256:core", id) } } func TestSelectComposeImageIdentityDoesNotFailOpen(t *testing.T) { for name, identities := range map[string][]composeImageIdentity{ "unrelated container only": {{ContainerName: "project-llm", ID: "sha256:shared-image"}}, "duplicate service": { {Service: "core", ID: "sha256:first"}, {Service: "core", ID: "sha256:second"}, }, } { t.Run(name, func(t *testing.T) { if id, err := selectComposeImageIdentity("core", "project-core-1", identities); err == nil { t.Fatalf("selectComposeImageIdentity() = %q, nil; want error", id) } }) } } func TestSelectComposeImageIdentityAcceptsSemanticallyEmptyOutput(t *testing.T) { id, err := selectComposeImageIdentity("workspace-maintenance", "project-workspace-maintenance-1", nil) if err != nil || id != "" { t.Fatalf("selectComposeImageIdentity() = %q, %v; want empty identity", id, err) } } 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 TestCreateReferencesAuthFilesByDefaultAndArchivesThemOnlyWithSecretCustody(t *testing.T) { fixture := newBackupFixture(t, "local") authDirectory := filepath.Join(filepath.Dir(fixture.installation.Path), "auth") if err := os.Mkdir(authDirectory, 0o700); err != nil { t.Fatal(err) } authPath := filepath.Join(authDirectory, "auth.yaml") usersPath := filepath.Join(authDirectory, "users.yaml") if err := os.WriteFile(authPath, []byte("version: 1\nmode: local\n"), 0o600); err != nil { t.Fatal(err) } if err := os.WriteFile(usersPath, []byte("users:\n - passwordHash: must-not-be-archived-by-default\n"), 0o600); err != nil { t.Fatal(err) } fixture.installation.Authentication.ConfigDirectory = authDirectory defaultOutput := filepath.Join(t.TempDir(), "default.zip") defaultResult, err := createWithDependencies(context.Background(), fixture.installation, CreateRequest{Output: defaultOutput}, testDependencies(t, newBackupRunner(fixture.installation, false))) if err != nil { t.Fatal(err) } if defaultResult.Warning != "" { t.Fatalf("default backup warning = %q, want no custody warning", defaultResult.Warning) } defaultArchive := readFixtureArchive(t, defaultOutput) defaultBytes := bytes.Join(mapValues(defaultArchive.files), nil) for _, value := range []string{"mode: local", "must-not-be-archived-by-default"} { if bytes.Contains(defaultBytes, []byte(value)) { t.Fatalf("default backup contains authentication content %q", value) } } if !manifestHasReference(defaultArchive.manifest, authPath) || manifestHasReference(defaultArchive.manifest, usersPath) { t.Fatalf("default backup did not record only the auth.yaml configuration path: %#v", defaultArchive.manifest.Entries) } secretOutput := filepath.Join(t.TempDir(), "with-auth-secrets.zip") secretResult, err := createWithDependencies(context.Background(), fixture.installation, CreateRequest{Output: secretOutput, IncludeSecrets: true, Confirm: true}, testDependencies(t, newBackupRunner(fixture.installation, false))) if err != nil { t.Fatal(err) } if !strings.Contains(secretResult.Warning, "custody") { t.Fatalf("secret backup warning = %q, want custody guidance", secretResult.Warning) } secretArchive := readFixtureArchive(t, secretOutput) for _, path := range []string{authPath, usersPath} { if !manifestHasArchivedSecret(secretArchive.manifest, path) { t.Fatalf("secret backup did not archive authentication file %q", path) } } } func manifestHasReference(manifest Manifest, sourcePath string) bool { for _, entry := range manifest.Entries { if entry.Kind == EntrySecretReference && entry.SourcePath == sourcePath && !entry.Archived { return true } } return false } func manifestHasArchivedSecret(manifest Manifest, sourcePath string) bool { for _, entry := range manifest.Entries { if entry.Kind == EntryExternalSecret && entry.SourcePath == sourcePath && entry.Archived && entry.Sensitive { return true } } return false } 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) serverVolumes := make([]string, 0, len(archive.manifest.Volumes)) for _, volume := range archive.manifest.Volumes { serverVolumes = append(serverVolumes, volume.LogicalName) } if strings.Join(serverVolumes, ",") != "embedding-models,qdrant-data" { t.Fatalf("server backup volumes = %v, want only infrastructure named volumes", serverVolumes) } 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) } }) } } func TestCreateCleansMutationsWhenDockerLosesTheResponse(t *testing.T) { activationResponseLost := errors.New("activation response lost") stopResponseLost := errors.New("stop response lost") for _, scenario := range []struct { name string failure *commandFailure wantErr error cancelCaller bool wantStopCount int wantStartCount int wantCleanupCalls int }{ { name: "maintenance activation", wantErr: activationResponseLost, failure: &commandFailure{ match: func(command string) bool { return strings.Contains(command, " maintenance-activate") }, err: activationResponseLost, remaining: 1, }, wantCleanupCalls: 1, }, { name: "stop response loss", wantErr: stopResponseLost, wantStopCount: 1, wantStartCount: 1, failure: &commandFailure{ match: func(command string) bool { return strings.HasSuffix(command, " stop") }, err: stopResponseLost, effect: func() { // A mutating command may have completed before its response was lost. }, remaining: 1, }, wantCleanupCalls: 2, }, { name: "caller cancellation after stop", wantErr: context.Canceled, cancelCaller: true, wantStopCount: 1, wantStartCount: 1, failure: &commandFailure{ match: func(command string) bool { return strings.HasSuffix(command, " stop") }, err: context.Canceled, remaining: 1, }, wantCleanupCalls: 2, }, } { scenario := scenario t.Run(scenario.name, func(t *testing.T) { fixture := newBackupFixture(t, "local") backing := newBackupRunner(fixture.installation, true) caller, cancel := context.WithCancel(context.Background()) t.Cleanup(cancel) failure := *scenario.failure originalEffect := failure.effect failure.effect = func() { if strings.Contains(scenario.name, "activation") { backing.maintenance = true } else { backing.stopCount++ backing.running, backing.coreRunning = false, false } if originalEffect != nil { originalEffect() } if scenario.cancelCaller { cancel() } } runner := &commandFailureRunner{fakeBackupRunner: backing, failures: []*commandFailure{&failure}} _, err := createWithDependencies(caller, fixture.installation, CreateRequest{Output: filepath.Join(t.TempDir(), "backup.zip")}, testDependencies(t, runner)) if !errors.Is(err, scenario.wantErr) { t.Fatalf("Create() error = %v, want %v", err, scenario.wantErr) } if backing.running != true || backing.coreRunning != true || backing.maintenance != false { t.Fatalf("cleanup state = running:%t core:%t maintenance:%t, want running and admitted", backing.running, backing.coreRunning, backing.maintenance) } if backing.stopCount != scenario.wantStopCount || backing.startCount != scenario.wantStartCount { t.Fatalf("lifecycle commands = stop:%d start:%d, want stop:%d start:%d", backing.stopCount, backing.startCount, scenario.wantStopCount, scenario.wantStartCount) } if len(runner.cleanupCommandContexts) != scenario.wantCleanupCalls { t.Fatalf("cleanup command contexts = %d, want %d", len(runner.cleanupCommandContexts), scenario.wantCleanupCalls) } assertIndependentBoundedCleanupContexts(t, runner.cleanupCommandContexts) }) } } func TestCreateJoinsPrimaryAndCleanupFailuresAfterPartialRestart(t *testing.T) { fixture := newBackupFixture(t, "local") backing := newBackupRunner(fixture.installation, true) startResponseLost := errors.New("start response lost") deactivationResponseLost := errors.New("deactivation response lost") runner := &commandFailureRunner{ fakeBackupRunner: backing, failures: []*commandFailure{ { match: func(command string) bool { return strings.HasSuffix(command, " start") }, err: startResponseLost, effect: func() { backing.startCount++ backing.running, backing.coreRunning = true, true }, remaining: 1, }, { match: func(command string) bool { return strings.Contains(command, " maintenance-deactivate") }, err: deactivationResponseLost, effect: func() { backing.maintenance = false }, remaining: 1, }, }, } _, err := createWithDependencies(context.Background(), fixture.installation, CreateRequest{Output: filepath.Join(t.TempDir(), "backup.zip")}, testDependencies(t, runner)) if !errors.Is(err, startResponseLost) || !errors.Is(err, deactivationResponseLost) { t.Fatalf("Create() error = %v, want joined start and deactivation failures", err) } if backing.startCount != 2 || !backing.running || backing.maintenance { t.Fatalf("partial-success cleanup state = starts:%d running:%t maintenance:%t", backing.startCount, backing.running, backing.maintenance) } } 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 } type commandFailure struct { match func(string) bool err error effect func() skip int remaining int } type commandFailureRunner struct { *fakeBackupRunner failures []*commandFailure cleanupCommandContexts []cleanupContextObservation streamFailure func(context.Context) error } type cleanupContextObservation struct { err error deadline time.Time hasDeadline bool } func observeCleanupContext(ctx context.Context) cleanupContextObservation { deadline, hasDeadline := ctx.Deadline() return cleanupContextObservation{err: ctx.Err(), deadline: deadline, hasDeadline: hasDeadline} } func (r *commandFailureRunner) Run(ctx context.Context, args []string, stdin io.Reader) (compose.Result, error) { command := strings.Join(args, " ") if strings.HasSuffix(command, " start") || strings.Contains(command, " maintenance-deactivate") { r.cleanupCommandContexts = append(r.cleanupCommandContexts, observeCleanupContext(ctx)) } for _, failure := range r.failures { if failure.remaining != 0 && failure.match(command) { if failure.skip > 0 { failure.skip-- continue } if failure.remaining > 0 { failure.remaining-- } if failure.effect != nil { failure.effect() } return compose.Result{}, failure.err } } return r.fakeBackupRunner.Run(ctx, args, stdin) } func (r *commandFailureRunner) Stream(ctx context.Context, args []string, stdin io.Reader, stdout io.Writer) (compose.Result, error) { if r.streamFailure != nil { return compose.Result{}, r.streamFailure(ctx) } return r.fakeBackupRunner.Stream(ctx, args, stdin, stdout) } func assertIndependentBoundedCleanupContexts(t *testing.T, contexts []cleanupContextObservation) { t.Helper() if len(contexts) == 0 { t.Fatal("expected cleanup commands") } for _, cleanupContext := range contexts { if cleanupContext.err != nil { t.Fatalf("cleanup used a cancelled context: %v", cleanupContext.err) } if !cleanupContext.hasDeadline { t.Fatal("cleanup context has no deadline") } if remaining := time.Until(cleanupContext.deadline); remaining <= 0 || remaining > 10*time.Minute { t.Fatalf("unexpected cleanup deadline remaining: %s", remaining) } } } 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 requiredBackupVolumes(runner.installation) { 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 "): required := requiredBackupVolumes(runner.installation) items := make([]map[string]any, 0, len(required)) for _, logical := range required { 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, "operator-command.js maintenance-status"): return compose.Result{Stdout: fmt.Sprintf(`{"active":%t,"admissions":0,"recoveryRequired":false}`, runner.maintenance)}, nil case strings.Contains(command, "operator-command.js maintenance-activate"): runner.maintenance = true return compose.Result{Stdout: `{"active":true,"admissions":0,"recoveryRequired":false}`}, nil case strings.Contains(command, "operator-command.js maintenance-deactivate"): runner.maintenance = false return compose.Result{Stdout: `{"active":false,"admissions":0,"recoveryRequired":false}`}, nil case strings.Contains(command, "operator-command.js session-inventory"): 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) } }