867 lines
34 KiB
Go
867 lines
34 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 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 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)
|
|
}
|
|
}
|