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

1098 lines
42 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 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 TestCreateTransactionCapabilityRefusesUnlockedOrForeignInstallation(t *testing.T) {
fixture := newBackupFixture(t, "local")
runner := newBackupRunner(fixture.installation, false)
request := CreateRequest{Output: filepath.Join(t.TempDir(), "transaction.zip")}
if _, err := createWithDependenciesTransaction(context.Background(), nil, fixture.installation, request, testDependencies(t, runner)); !errors.Is(err, lifecycle.ErrTransactionInactive) {
t.Fatalf("nil transaction error = %v, want ErrTransactionInactive", err)
}
if len(runner.calls) != 0 {
t.Fatalf("Docker runner was called without a transaction capability: %v", runner.calls)
}
otherRoot := t.TempDir()
other := config.Installation{ProjectDirectory: otherRoot, Path: filepath.Join(otherRoot, "thothii-installation.yaml")}
transaction, err := lifecycle.AcquireTransaction(other)
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = transaction.Release() })
if _, err := createWithDependenciesTransaction(context.Background(), transaction, fixture.installation, request, testDependencies(t, runner)); !errors.Is(err, lifecycle.ErrTransactionInstallation) {
t.Fatalf("foreign transaction error = %v, want ErrTransactionInstallation", err)
}
if len(runner.calls) != 0 {
t.Fatalf("Docker runner was called with a foreign transaction capability: %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)
}
}