From 8f55350f8b37b30d0a66a350e11981fc560a0d0d Mon Sep 17 00:00:00 2001 From: Bram Schuur Date: Thu, 16 Jul 2026 10:59:57 +0200 Subject: [PATCH 1/9] STAC-24630: Add stackgraph-v2 backup/restore --- cmd/root.go | 5 + cmd/settings/restore.go | 2 +- cmd/stackgraph/restore.go | 2 +- cmd/stackgraphv2/abort.go | 67 +++ cmd/stackgraphv2/backfill.go | 74 ++++ cmd/stackgraphv2/check_and_finalize.go | 58 +++ cmd/stackgraphv2/list.go | 102 +++++ cmd/stackgraphv2/restore.go | 415 ++++++++++++++++++ cmd/stackgraphv2/restore_test.go | 125 ++++++ cmd/stackgraphv2/simplejob.go | 134 ++++++ cmd/stackgraphv2/stackgraphv2.go | 21 + cmd/victoriametrics/restore.go | 2 +- internal/clients/s3/filter.go | 17 + internal/orchestration/restore/job.go | 3 + .../scripts/restore-settings-backup.sh | 2 +- .../restore-stackgraph-backup-v2-abort.sh | 13 + .../restore-stackgraph-backup-v2-backfill.sh | 13 + .../restore-stackgraph-backup-v2-env.sh | 30 ++ .../restore-stackgraph-backup-v2-live.sh | 38 ++ .../scripts/restore-stackgraph-backup.sh | 6 +- .../restore-victoria-metrics-backup.sh | 2 +- 21 files changed, 1123 insertions(+), 8 deletions(-) create mode 100644 cmd/stackgraphv2/abort.go create mode 100644 cmd/stackgraphv2/backfill.go create mode 100644 cmd/stackgraphv2/check_and_finalize.go create mode 100644 cmd/stackgraphv2/list.go create mode 100644 cmd/stackgraphv2/restore.go create mode 100644 cmd/stackgraphv2/restore_test.go create mode 100644 cmd/stackgraphv2/simplejob.go create mode 100644 cmd/stackgraphv2/stackgraphv2.go create mode 100644 internal/scripts/scripts/restore-stackgraph-backup-v2-abort.sh create mode 100644 internal/scripts/scripts/restore-stackgraph-backup-v2-backfill.sh create mode 100644 internal/scripts/scripts/restore-stackgraph-backup-v2-env.sh create mode 100644 internal/scripts/scripts/restore-stackgraph-backup-v2-live.sh diff --git a/cmd/root.go b/cmd/root.go index f2d7686..cba0f17 100644 --- a/cmd/root.go +++ b/cmd/root.go @@ -8,6 +8,7 @@ import ( "github.com/stackvista/stackstate-backup-cli/cmd/elasticsearch" "github.com/stackvista/stackstate-backup-cli/cmd/settings" "github.com/stackvista/stackstate-backup-cli/cmd/stackgraph" + "github.com/stackvista/stackstate-backup-cli/cmd/stackgraphv2" "github.com/stackvista/stackstate-backup-cli/cmd/version" "github.com/stackvista/stackstate-backup-cli/cmd/victoriametrics" "github.com/stackvista/stackstate-backup-cli/internal/foundation/config" @@ -42,6 +43,10 @@ func init() { addBackupConfigFlags(stackgraphCmd) rootCmd.AddCommand(stackgraphCmd) + stackgraphv2Cmd := stackgraphv2.Cmd(flags) + addBackupConfigFlags(stackgraphv2Cmd) + rootCmd.AddCommand(stackgraphv2Cmd) + settingsCmd := settings.Cmd(flags) addBackupConfigFlags(settingsCmd) rootCmd.AddCommand(settingsCmd) diff --git a/cmd/settings/restore.go b/cmd/settings/restore.go index 05e09c8..86742b5 100644 --- a/cmd/settings/restore.go +++ b/cmd/settings/restore.go @@ -197,7 +197,7 @@ func buildEnvVar(extraEnvVar []corev1.EnvVar, config *config.Config) []corev1.En commonVar := []corev1.EnvVar{ {Name: "BACKUP_CONFIGURATION_BUCKET_NAME", Value: config.Settings.Bucket}, {Name: "BACKUP_CONFIGURATION_S3_PREFIX", Value: config.Settings.S3Prefix}, - {Name: "MINIO_ENDPOINT", Value: fmt.Sprintf("%s:%d", storageService.Name, storageService.Port)}, + {Name: "S3_ENDPOINT", Value: fmt.Sprintf("%s:%d", storageService.Name, storageService.Port)}, {Name: "STACKSTATE_BASE_URL", Value: config.GetBaseURL()}, {Name: "RECEIVER_BASE_URL", Value: config.GetReceiverBaseURL()}, {Name: "PLATFORM_VERSION", Value: config.GetPlatformVersion()}, diff --git a/cmd/stackgraph/restore.go b/cmd/stackgraph/restore.go index 8391ceb..4c02c06 100644 --- a/cmd/stackgraph/restore.go +++ b/cmd/stackgraph/restore.go @@ -276,7 +276,7 @@ func buildRestoreEnvVars(backupFile string, config *config.Config) []corev1.EnvV {Name: "FORCE_DELETE", Value: purgeStackgraphDataFlag}, {Name: "BACKUP_STACKGRAPH_BUCKET_NAME", Value: config.Stackgraph.Bucket}, {Name: "BACKUP_STACKGRAPH_S3_PREFIX", Value: config.Stackgraph.S3Prefix}, - {Name: "MINIO_ENDPOINT", Value: fmt.Sprintf("%s:%d", storageService.Name, storageService.Port)}, + {Name: "S3_ENDPOINT", Value: fmt.Sprintf("%s:%d", storageService.Name, storageService.Port)}, {Name: "STACKSTATE_BASE_URL", Value: config.GetBaseURL()}, {Name: "RECEIVER_BASE_URL", Value: config.GetReceiverBaseURL()}, {Name: "PLATFORM_VERSION", Value: config.GetPlatformVersion()}, diff --git a/cmd/stackgraphv2/abort.go b/cmd/stackgraphv2/abort.go new file mode 100644 index 0000000..24fa11f --- /dev/null +++ b/cmd/stackgraphv2/abort.go @@ -0,0 +1,67 @@ +package stackgraphv2 + +import ( + "fmt" + "time" + + "github.com/spf13/cobra" + "github.com/stackvista/stackstate-backup-cli/cmd/cmdutils" + "github.com/stackvista/stackstate-backup-cli/internal/app" + "github.com/stackvista/stackstate-backup-cli/internal/clients/k8s" + "github.com/stackvista/stackstate-backup-cli/internal/foundation/config" + "github.com/stackvista/stackstate-backup-cli/internal/foundation/logger" + "github.com/stackvista/stackstate-backup-cli/internal/orchestration/restore" +) + +const ( + abortNameTemplate = "stackgraph-abort-v2" + abortScript = "/backup-restore-scripts/restore-stackgraph-backup-v2-abort.sh" +) + +func abortCmd(globalFlags *config.CLIGlobalFlags) *cobra.Command { + cmd := &cobra.Command{ + Use: "abort", + Short: "Abort a partially backfilled restore.", + Long: "Abort a partially backfilled restore. This will not load (backfill) any data but keep the data as is." + + "Be aware this leaves the restore incomplete, historical data will not be recovered, but leave the instance in a usable state." + + "This should be used when backfilling is failing and the database restore is accepted as-is.", + Run: func(_ *cobra.Command, _ []string) { + cmdutils.Run(globalFlags, runAbort, cmdutils.StorageIsRequired) + }, + } + + return cmd +} + +func runAbort(appCtx *app.Context) error { + // Setup Kubernetes resources for restore job + appCtx.Logger.Println() + if err := restore.EnsureResources(appCtx.K8sClient, appCtx.Namespace, appCtx.Config, appCtx.Logger); err != nil { + return err + } + + // Create restore job + appCtx.Logger.Println() + appCtx.Logger.Infof("Creating abort job.") + + jobName := fmt.Sprintf("%s-%s", abortNameTemplate, time.Now().Format("20060102t150405")) + + if err := createAbortJob(appCtx.K8sClient, appCtx.Namespace, jobName, appCtx.Config); err != nil { + return fmt.Errorf("failed to create abort job: %w", err) + } + + appCtx.Logger.Successf("Abort job created: %s", jobName) + + return waitAndCleanupAbortJob(appCtx.K8sClient, appCtx.Namespace, jobName, appCtx.Logger) +} + +// waitAndCleanupAbortJob waits for job completion and cleans up resources +func waitAndCleanupAbortJob(k8sClient *k8s.Client, namespace, jobName string, log *logger.Logger) error { + restore.PrintWaitingMessage(log, "stackgraph-v2", jobName, namespace) + return restore.WaitAndCleanup(k8sClient, namespace, jobName, log, false) +} + +// createAbortJob creates a Kubernetes Job for aborting a partially backfilled restore +func createAbortJob(k8sClient *k8s.Client, namespace, jobName string, config *config.Config) error { + return createSimpleJob(k8sClient, namespace, jobName, "abort", abortScript, config) +} diff --git a/cmd/stackgraphv2/backfill.go b/cmd/stackgraphv2/backfill.go new file mode 100644 index 0000000..f0a8ff9 --- /dev/null +++ b/cmd/stackgraphv2/backfill.go @@ -0,0 +1,74 @@ +package stackgraphv2 + +import ( + "fmt" + "time" + + "github.com/spf13/cobra" + "github.com/stackvista/stackstate-backup-cli/cmd/cmdutils" + "github.com/stackvista/stackstate-backup-cli/internal/app" + "github.com/stackvista/stackstate-backup-cli/internal/clients/k8s" + "github.com/stackvista/stackstate-backup-cli/internal/foundation/config" + "github.com/stackvista/stackstate-backup-cli/internal/foundation/logger" + "github.com/stackvista/stackstate-backup-cli/internal/orchestration/restore" +) + +const ( + backfillNameTemplate = "stackgraph-backfill-v2" + backfillScript = "/backup-restore-scripts/restore-stackgraph-backup-v2-backfill.sh" +) + +func backfillCmd(globalFlags *config.CLIGlobalFlags) *cobra.Command { + cmd := &cobra.Command{ + Use: "backfill", + Short: "Complete a restore by backfilling old data", + Long: "Complete a restore by backfilling old data. This process can run while the system is already up and running " + + "After the backfill is done, the restore process is complete. ", + Run: func(_ *cobra.Command, _ []string) { + cmdutils.Run(globalFlags, runBackfill, cmdutils.StorageIsRequired) + }, + } + + return cmd +} + +func runBackfill(appCtx *app.Context) error { + // Setup Kubernetes resources for restore job + appCtx.Logger.Println() + if err := restore.EnsureResources(appCtx.K8sClient, appCtx.Namespace, appCtx.Config, appCtx.Logger); err != nil { + return err + } + + // Create restore job + appCtx.Logger.Println() + appCtx.Logger.Infof("Creating backfill job.") + + jobName := fmt.Sprintf("%s-%s", backfillNameTemplate, time.Now().Format("20060102t150405")) + + if err := createBackfillJob(appCtx.K8sClient, appCtx.Namespace, jobName, appCtx.Config); err != nil { + return fmt.Errorf("failed to create backfill job: %w", err) + } + + appCtx.Logger.Successf("Backfill job created: %s", jobName) + + err := waitAndCleanupBackfillJob(appCtx.K8sClient, appCtx.Namespace, jobName, appCtx.Logger) + + if err != nil { + appCtx.Logger.Println() + appCtx.Logger.Infof("Backfill failed. It is possible to restart the backfill without starting a complete restore.") + return err + } + + return nil +} + +// waitAndCleanupBackfillJob waits for job completion and cleans up resources +func waitAndCleanupBackfillJob(k8sClient *k8s.Client, namespace, jobName string, log *logger.Logger) error { + restore.PrintWaitingMessage(log, "stackgraph-v2", jobName, namespace) + return restore.WaitAndCleanup(k8sClient, namespace, jobName, log, false) +} + +// createBackfillJob creates a Kubernetes Job for backfilling old data during a restore +func createBackfillJob(k8sClient *k8s.Client, namespace, jobName string, config *config.Config) error { + return createSimpleJob(k8sClient, namespace, jobName, "backfill", backfillScript, config) +} diff --git a/cmd/stackgraphv2/check_and_finalize.go b/cmd/stackgraphv2/check_and_finalize.go new file mode 100644 index 0000000..39cd794 --- /dev/null +++ b/cmd/stackgraphv2/check_and_finalize.go @@ -0,0 +1,58 @@ +package stackgraphv2 + +import ( + "github.com/spf13/cobra" + "github.com/stackvista/stackstate-backup-cli/cmd/cmdutils" + "github.com/stackvista/stackstate-backup-cli/internal/app" + "github.com/stackvista/stackstate-backup-cli/internal/foundation/config" + "github.com/stackvista/stackstate-backup-cli/internal/orchestration/restore" + "github.com/stackvista/stackstate-backup-cli/internal/orchestration/scale" +) + +// Check and finalize command flags +var ( + checkJobName string + waitForJob bool +) + +func checkAndFinalizeCmd(globalFlags *config.CLIGlobalFlags) *cobra.Command { + cmd := &cobra.Command{ + Use: "check-and-finalize", + Short: "Check and finalize a Stackgraph restore (v2) job", + Long: `Check the status of a background Stackgraph restore job and clean up resources. + +This command is useful when a restore job was started with --background flag or was interrupted (Ctrl+C). +It will check the job status, print logs if it failed, and clean up the job and PVC resources. + +Examples: + # Check job status without waiting + sts-backup stackgraph-v2 check-and-finalize --job stackgraph-restore-20250128t143000 -n my-namespace + + # Wait for job completion and cleanup + sts-backup stackgraph-v2 check-and-finalize --job stackgraph-restore-20250128t143000 --wait -n my-namespace`, + Run: func(_ *cobra.Command, _ []string) { + cmdutils.Run(globalFlags, runCheckAndFinalize, cmdutils.StorageIsRequired) + }, + } + + cmd.Flags().StringVarP(&checkJobName, "job", "j", "", "Stackgraph restore job name (required)") + cmd.Flags().BoolVarP(&waitForJob, "wait", "w", false, "Wait for job to complete before cleanup") + _ = cmd.MarkFlagRequired("job") + + return cmd +} + +func runCheckAndFinalize(appCtx *app.Context) error { + return restore.CheckAndFinalize(restore.CheckAndFinalizeParams{ + K8sClient: appCtx.K8sClient, + Namespace: appCtx.Namespace, + JobName: checkJobName, + ServiceName: "stackgraph", + ScaleUpFn: scale.ScaleUpAndReleaseLock, + ScaleDownFn: scale.ScaleDown, + ScaleSelector: appCtx.Config.Stackgraph.Restore.ScaleDownLabelSelector, + CleanupPVC: true, + WaitForJob: waitForJob, + Log: appCtx.Logger, + }) +} diff --git a/cmd/stackgraphv2/list.go b/cmd/stackgraphv2/list.go new file mode 100644 index 0000000..7700379 --- /dev/null +++ b/cmd/stackgraphv2/list.go @@ -0,0 +1,102 @@ +package stackgraphv2 + +import ( + "context" + "fmt" + "sort" + "strings" + + "github.com/aws/aws-sdk-go-v2/aws" + "github.com/aws/aws-sdk-go-v2/service/s3" + "github.com/spf13/cobra" + "github.com/stackvista/stackstate-backup-cli/cmd/cmdutils" + "github.com/stackvista/stackstate-backup-cli/internal/app" + s3client "github.com/stackvista/stackstate-backup-cli/internal/clients/s3" + "github.com/stackvista/stackstate-backup-cli/internal/foundation/config" + "github.com/stackvista/stackstate-backup-cli/internal/foundation/output" + "github.com/stackvista/stackstate-backup-cli/internal/orchestration/portforward" +) + +const ( + backupFileNameRegex = `^sts-backup-.*\.graph.v2$` +) + +func listCmd(globalFlags *config.CLIGlobalFlags) *cobra.Command { + return &cobra.Command{ + Use: "list", + Short: "List available Stackgraph backups (v2) from S3", + Run: func(_ *cobra.Command, _ []string) { + cmdutils.Run(globalFlags, runList, cmdutils.StorageIsRequired) + }, + } +} + +func runList(appCtx *app.Context) error { + // Setup port-forward to S3-compatible storage + storageService := appCtx.Config.GetStorageService() + serviceName := storageService.Name + remotePort := storageService.Port + + pf, err := portforward.SetupPortForward(appCtx.K8sClient, appCtx.Namespace, serviceName, remotePort, appCtx.Logger) + if err != nil { + return err + } + defer close(pf.StopChan) + + // Create S3 client with actual port + s3Client, err := appCtx.NewS3Client(pf.LocalPort) + if err != nil { + return fmt.Errorf("failed to create S3 client: %w", err) + } + + // List objects in bucket + bucket := appCtx.Config.Stackgraph.Bucket + prefix := appCtx.Config.Stackgraph.S3Prefix + "v2/" + + appCtx.Logger.Infof("Listing Stackgraph backups in bucket '%s' with prefix '%s'...", bucket, prefix) + + input := &s3.ListObjectsV2Input{ + Bucket: aws.String(bucket), + Prefix: aws.String(prefix), + Delimiter: aws.String("/"), + } + + result, err := s3Client.ListObjectsV2(context.Background(), input) + if err != nil { + return fmt.Errorf("failed to list S3 objects: %w", err) + } + + filteredObjects := s3client.ConvertBackupObjects(result.Contents) + + // Filter to only include direct children of the prefix that match the backup filename pattern, + // and strip the prefix from the key + filteredObjects, err = s3client.FilterByPrefixAndRegex(filteredObjects, prefix, backupFileNameRegex) + if err != nil { + return fmt.Errorf("failed to filter objects: %w", err) + } + + // Sort by LastModified time (most recent first) + sort.Slice(filteredObjects, func(i, j int) bool { + return filteredObjects[i].LastModified.After(filteredObjects[j].LastModified) + }) + + if len(filteredObjects) == 0 { + appCtx.Formatter.PrintMessage("No backups found") + return nil + } + + table := output.Table{ + Headers: []string{"NAME", "LAST MODIFIED"}, + Rows: make([][]string, 0, len(filteredObjects)), + } + + for _, obj := range filteredObjects { + row := []string{ + strings.TrimPrefix(obj.Key, prefix), + obj.LastModified.Format("2006-01-02 15:04:05 MST"), + } + table.Rows = append(table.Rows, row) + } + + return appCtx.Formatter.PrintTable(table) +} diff --git a/cmd/stackgraphv2/restore.go b/cmd/stackgraphv2/restore.go new file mode 100644 index 0000000..c42ec8d --- /dev/null +++ b/cmd/stackgraphv2/restore.go @@ -0,0 +1,415 @@ +package stackgraphv2 + +import ( + "context" + "fmt" + "sort" + "strconv" + "strings" + "time" + + "github.com/aws/aws-sdk-go-v2/aws" + "github.com/aws/aws-sdk-go-v2/service/s3" + "github.com/spf13/cobra" + "github.com/stackvista/stackstate-backup-cli/cmd/cmdutils" + "github.com/stackvista/stackstate-backup-cli/internal/app" + "github.com/stackvista/stackstate-backup-cli/internal/clients/k8s" + s3client "github.com/stackvista/stackstate-backup-cli/internal/clients/s3" + "github.com/stackvista/stackstate-backup-cli/internal/foundation/config" + "github.com/stackvista/stackstate-backup-cli/internal/foundation/logger" + "github.com/stackvista/stackstate-backup-cli/internal/orchestration/portforward" + "github.com/stackvista/stackstate-backup-cli/internal/orchestration/restore" + "github.com/stackvista/stackstate-backup-cli/internal/orchestration/scale" + corev1 "k8s.io/api/core/v1" +) + +const ( + jobNameTemplate = "stackgraph-restore-v2" + configMapDefaultFileMode = 0755 + purgeStackgraphDataFlag = "-force" +) + +// Restore command flags +var ( + archiveName string + useLatest bool + skipConfirmation bool + skipStackpacks bool +) + +func restoreCmd(globalFlags *config.CLIGlobalFlags) *cobra.Command { + cmd := &cobra.Command{ + Use: "restore", + Short: "Restore Stackgraph from a backup archive", + Long: "Restore Stackgraph data from a backup archive stored in S3. Automatically also restores " + + "Stackpacks backup that was made at the same time, it can be skipped with --skip-stackpacks. " + + "Can use --latest or --archive to specify which backup to restore." + + "The process is split into two phases: restoring live data, which takes the system down. And a backfill stage which loads historic data.", + Run: func(_ *cobra.Command, _ []string) { + cmdutils.Run(globalFlags, runRestore, cmdutils.StorageIsRequired) + }, + } + + cmd.Flags().StringVar(&archiveName, "archive", "", "Specific archive name to restore (e.g., sts-backup-20210216-0300.graph)") + cmd.Flags().BoolVar(&useLatest, "latest", false, "Restore from the most recent backup") + cmd.Flags().BoolVarP(&skipConfirmation, "yes", "y", false, "Skip confirmation prompt") + cmd.Flags().BoolVar(&skipStackpacks, "skip-stackpacks", false, "Skip restoring stackpacks backup") + cmd.MarkFlagsMutuallyExclusive("archive", "latest") + cmd.MarkFlagsOneRequired("archive", "latest") + + return cmd +} + +func runRestore(appCtx *app.Context) error { + err := liveRestore(appCtx) + if err != nil { + return err + } + + appCtx.Logger.Println() + appCtx.Logger.Infof("The system is back up and accessible. Continuing with backfilling historical data, this will happen while the system is running.") + appCtx.Logger.Infof("During this process not all historical data might be available.") + + return runBackfill(appCtx) +} + +func liveRestore(appCtx *app.Context) error { + // Determine which archive to restore + backupFile := archiveName + if useLatest { + appCtx.Logger.Infof("Finding latest backup...") + latest, err := getLatestBackup(appCtx.K8sClient, appCtx.Namespace, appCtx.Config, appCtx.Logger) + if err != nil { + return err + } + backupFile = latest + appCtx.Logger.Infof("Using latest backup: %s", backupFile) + } + + // Warn user and ask for confirmation + if !skipConfirmation { + appCtx.Logger.Println() + appCtx.Logger.Warningf("WARNING: Restoring from backup will PURGE all existing Stackgraph data!") + appCtx.Logger.Warningf("This operation cannot be undone.") + appCtx.Logger.Println() + appCtx.Logger.Infof("Backup to restore: %s", backupFile) + appCtx.Logger.Infof("Namespace: %s", appCtx.Namespace) + appCtx.Logger.Println() + + if !restore.PromptForConfirmation() { + return fmt.Errorf("restore operation cancelled by user") + } + } + + // Scale down deployments before restore (with lock protection) + appCtx.Logger.Println() + scaleDownLabelSelector := appCtx.Config.Stackgraph.Restore.ScaleDownLabelSelector + scaledDeployments, err := scale.ScaleDownWithLock(scale.ScaleDownWithLockParams{ + K8sClient: appCtx.K8sClient, + Namespace: appCtx.Namespace, + LabelSelector: scaleDownLabelSelector, + Datastore: config.DatastoreStackgraph, + AllSelectors: appCtx.Config.GetAllScaleDownSelectors(), + Log: appCtx.Logger, + }) + if err != nil { + return err + } + + // Ensure deployments are scaled back up and lock released on exit (even if restore fails) + defer func() { + if len(scaledDeployments) > 0 { + appCtx.Logger.Println() + if err := scale.ScaleUpAndReleaseLock(appCtx.K8sClient, appCtx.Namespace, scaleDownLabelSelector, appCtx.Logger); err != nil { + appCtx.Logger.Warningf("Failed to scale up deployments: %v", err) + } + } + }() + + // Setup Kubernetes resources for restore job + appCtx.Logger.Println() + if err := restore.EnsureResources(appCtx.K8sClient, appCtx.Namespace, appCtx.Config, appCtx.Logger); err != nil { + return err + } + + // Create restore job + appCtx.Logger.Println() + appCtx.Logger.Infof("Creating restore of live data job for backup: %s", backupFile) + + jobName := fmt.Sprintf("%s-%s", jobNameTemplate, time.Now().Format("20060102t150405")) + + if err = createRestoreJob(appCtx.K8sClient, appCtx.Namespace, jobName, backupFile, appCtx.Config); err != nil { + return fmt.Errorf("failed to create restore job: %w", err) + } + + appCtx.Logger.Successf("Restore job created: %s", jobName) + + err = waitAndCleanupRestoreJob(appCtx.K8sClient, appCtx.Namespace, jobName, appCtx.Logger) + if err != nil { + return err + } + + appCtx.Logger.Println() + appCtx.Logger.Infof("Successfully restored live data for: %s", backupFile) + return nil +} + +// waitAndCleanupRestoreJob waits for job completion and cleans up resources +func waitAndCleanupRestoreJob(k8sClient *k8s.Client, namespace, jobName string, log *logger.Logger) error { + restore.PrintWaitingMessage(log, "stackgraph-v2", jobName, namespace) + return restore.WaitAndCleanup(k8sClient, namespace, jobName, log, true) +} + +// getLatestBackup retrieves the most recent backup from S3 +func getLatestBackup(k8sClient *k8s.Client, namespace string, config *config.Config, log *logger.Logger) (string, error) { + // Setup port-forward to S3-compatible storage + storageService := config.GetStorageService() + serviceName := storageService.Name + remotePort := storageService.Port + + pf, err := portforward.SetupPortForward(k8sClient, namespace, serviceName, remotePort, log) + if err != nil { + return "", err + } + defer close(pf.StopChan) + + // Create S3 client + endpoint := fmt.Sprintf("http://localhost:%d", pf.LocalPort) + s3Client, err := s3client.NewClient(endpoint, config.GetStorageAccessKey(), config.GetStorageSecretKey()) + if err != nil { + return "", err + } + + // List objects in bucket + bucket := config.Stackgraph.Bucket + prefix := config.Stackgraph.S3Prefix + "v2/" + + input := &s3.ListObjectsV2Input{ + Bucket: aws.String(bucket), + Prefix: aws.String(prefix), + Delimiter: aws.String("/"), + } + + result, err := s3Client.ListObjectsV2(context.Background(), input) + if err != nil { + return "", fmt.Errorf("failed to list S3 objects: %w", err) + } + + filteredObjects := s3client.ConvertBackupObjects(result.Contents) + + // Filter to only include direct children of the prefix that match the backup filename pattern, + // and strip the prefix from the key + filteredObjects, err = s3client.FilterByPrefixAndRegex(filteredObjects, prefix, backupFileNameRegex) + if err != nil { + return "", fmt.Errorf("failed to filter objects: %w", err) + } + + if len(filteredObjects) == 0 { + return "", fmt.Errorf("no backups found in bucket %s", bucket) + } + + // Sort by LastModified time (most recent first) + sort.Slice(filteredObjects, func(i, j int) bool { + return filteredObjects[i].LastModified.After(filteredObjects[j].LastModified) + }) + return filteredObjects[0].Key, nil +} + +// buildPVCSpec builds a PVCSpec from configuration +func buildPVCSpec(name string, config *config.Config, labels map[string]string) k8s.PVCSpec { + pvcConfig := config.Stackgraph.Restore.PVC + + // Convert string access modes to k8s types + accessModes := []corev1.PersistentVolumeAccessMode{corev1.ReadWriteOnce} // default + if len(pvcConfig.AccessModes) > 0 { + accessModes = make([]corev1.PersistentVolumeAccessMode, 0, len(pvcConfig.AccessModes)) + for _, mode := range pvcConfig.AccessModes { + accessModes = append(accessModes, corev1.PersistentVolumeAccessMode(mode)) + } + } + + // Handle storage class (nil if not set) + var storageClass *string + if pvcConfig.StorageClassName != "" { + storageClass = &pvcConfig.StorageClassName + } + + return k8s.PVCSpec{ + Name: name, + Labels: labels, + // This is not derived from the restore pvc size, the tmp pvc for restore only needs a bounded buffer for stackpacks. + StorageSize: "5Gi", + AccessModes: accessModes, + StorageClass: storageClass, + } +} + +// createRestoreJob creates a Kubernetes Job and PVC for restoring from backup +func createRestoreJob(k8sClient *k8s.Client, namespace, jobName, backupFile string, config *config.Config) error { + defaultMode := int32(configMapDefaultFileMode) + + // Merge common labels with resource-specific labels + pvcLabels := k8s.MergeLabels(config.Kubernetes.CommonLabels, map[string]string{}) + jobLabels := k8s.MergeLabels(config.Kubernetes.CommonLabels, config.Stackgraph.Restore.Job.Labels) + + // Create PVC first + pvcSpec := buildPVCSpec(jobName, config, pvcLabels) + pvc, err := k8sClient.CreatePVC(namespace, pvcSpec) + if err != nil { + return fmt.Errorf("failed to create PVC: %w", err) + } + + // Build job spec using configuration + spec := k8s.JobSpec{ + Name: jobName, + Labels: jobLabels, + ImagePullSecrets: k8s.ConvertImagePullSecrets(config.Stackgraph.Restore.Job.ImagePullSecrets), + SecurityContext: k8s.ConvertPodSecurityContext(&config.Stackgraph.Restore.Job.SecurityContext), + NodeSelector: config.Stackgraph.Restore.Job.NodeSelector, + Tolerations: k8s.ConvertTolerations(config.Stackgraph.Restore.Job.Tolerations), + Affinity: k8s.ConvertAffinity(config.Stackgraph.Restore.Job.Affinity), + Containers: buildRestoreContainers(backupFile, config), + InitContainers: buildRestoreInitContainers(config), + Volumes: buildRestoreVolumes(jobName, config, defaultMode), + } + + // Create job + _, err = k8sClient.CreateJob(namespace, spec) + if err != nil { + // Cleanup PVC if job creation fails + _ = k8sClient.DeletePVC(namespace, pvc.Name) + return fmt.Errorf("failed to create job: %w", err) + } + + return nil +} + +// buildRestoreEnvVars constructs environment variables for the restore job +func buildRestoreEnvVars(backupFile string, config *config.Config) []corev1.EnvVar { + storageService := config.GetStorageService() + env := []corev1.EnvVar{ + {Name: "BACKUP_FILE", Value: backupFile}, + {Name: "FORCE_DELETE", Value: purgeStackgraphDataFlag}, + {Name: "BACKUP_STACKGRAPH_BUCKET_NAME", Value: config.Stackgraph.Bucket}, + {Name: "BACKUP_STACKGRAPH_S3_PREFIX", Value: config.Stackgraph.S3Prefix}, + {Name: "S3_ENDPOINT", Value: fmt.Sprintf("%s:%d", storageService.Name, storageService.Port)}, + {Name: "STACKSTATE_BASE_URL", Value: config.GetBaseURL()}, + {Name: "RECEIVER_BASE_URL", Value: config.GetReceiverBaseURL()}, + {Name: "PLATFORM_VERSION", Value: config.GetPlatformVersion()}, + {Name: "ZOOKEEPER_QUORUM", Value: config.Stackgraph.Restore.ZookeeperQuorum}, + {Name: "SKIP_STACKPACKS", Value: strconv.FormatBool(skipStackpacks || config.Stackpacks == nil)}, + } + if config.Stackpacks != nil { + env = append(env, corev1.EnvVar{Name: "CONFIG_FORCE_stackstate_stackPacks_localStackPacksUri", Value: config.Stackpacks.LocalStackPacksURI}) + env = append(env, corev1.EnvVar{Name: "BACKUP_STACKGRAPH_STACKPACKS_DIR", Value: config.Stackpacks.BackupDirectory}) + } + return env +} + +// buildRestoreVolumeMounts constructs volume mounts for the restore job container +func buildRestoreVolumeMounts(config *config.Config) []corev1.VolumeMount { + volumeMounts := []corev1.VolumeMount{ + {Name: "backup-log", MountPath: "/opt/docker/etc_log"}, + {Name: "backup-restore-scripts", MountPath: "/backup-restore-scripts"}, + {Name: "minio-keys", MountPath: "/aws-keys"}, + {Name: "tmp-data", MountPath: "/tmp-data"}, + } + + if config.Stackpacks != nil && config.Stackpacks.PVC != "" && strings.HasPrefix(config.Stackpacks.LocalStackPacksURI, "file://") { + volumeMounts = append(volumeMounts, corev1.VolumeMount{ + Name: "stackpacks-local", + MountPath: strings.TrimPrefix(config.Stackpacks.LocalStackPacksURI, "file://"), + }) + } + + return volumeMounts +} + +// buildRestoreInitContainers constructs init containers for the restore job +func buildRestoreInitContainers(config *config.Config) []corev1.Container { + storageService := config.GetStorageService() + return []corev1.Container{ + { + Name: "wait", + Image: config.Stackgraph.Restore.Job.WaitImage, + ImagePullPolicy: corev1.PullIfNotPresent, + Command: []string{ + "sh", + "-c", + fmt.Sprintf("/entrypoint -c %s:%d -t 300", storageService.Name, storageService.Port), + }, + SecurityContext: k8s.ConvertSecurityContext(config.Stackgraph.Restore.Job.ContainerSecurityContext), + }, + } +} + +// buildRestoreVolumes constructs volumes for the restore job pod +func buildRestoreVolumes(jobName string, config *config.Config, defaultMode int32) []corev1.Volume { + volumes := []corev1.Volume{ + { + Name: "backup-log", + VolumeSource: corev1.VolumeSource{ + ConfigMap: &corev1.ConfigMapVolumeSource{ + LocalObjectReference: corev1.LocalObjectReference{ + Name: config.Stackgraph.Restore.LoggingConfigConfigMapName, + }, + }, + }, + }, + { + Name: "backup-restore-scripts", + VolumeSource: corev1.VolumeSource{ + ConfigMap: &corev1.ConfigMapVolumeSource{ + LocalObjectReference: corev1.LocalObjectReference{ + Name: restore.RestoreScriptsConfigMap, + }, + DefaultMode: &defaultMode, + }, + }, + }, + { + Name: "minio-keys", + VolumeSource: corev1.VolumeSource{ + Secret: &corev1.SecretVolumeSource{ + SecretName: restore.MinioKeysSecretName, + }, + }, + }, + { + Name: "tmp-data", + VolumeSource: corev1.VolumeSource{ + PersistentVolumeClaim: &corev1.PersistentVolumeClaimVolumeSource{ + ClaimName: jobName, + }, + }, + }, + } + if config.Stackpacks != nil && config.Stackpacks.PVC != "" { + volumes = append(volumes, corev1.Volume{ + Name: "stackpacks-local", + VolumeSource: corev1.VolumeSource{ + PersistentVolumeClaim: &corev1.PersistentVolumeClaimVolumeSource{ + ClaimName: config.Stackpacks.PVC, + }, + }, + }) + } + + return volumes +} + +// buildRestoreContainers constructs containers for the restore job +func buildRestoreContainers(backupFile string, config *config.Config) []corev1.Container { + return []corev1.Container{ + { + Name: "restore", + Image: config.Stackgraph.Restore.Job.Image, + ImagePullPolicy: corev1.PullIfNotPresent, + SecurityContext: k8s.ConvertSecurityContext(config.Stackgraph.Restore.Job.ContainerSecurityContext), + Command: []string{"/backup-restore-scripts/restore-stackgraph-backup-v2-live.sh"}, + Env: buildRestoreEnvVars(backupFile, config), + Resources: k8s.ConvertResources(config.Stackgraph.Restore.Job.Resources), + VolumeMounts: buildRestoreVolumeMounts(config), + }, + } +} diff --git a/cmd/stackgraphv2/restore_test.go b/cmd/stackgraphv2/restore_test.go new file mode 100644 index 0000000..42a1939 --- /dev/null +++ b/cmd/stackgraphv2/restore_test.go @@ -0,0 +1,125 @@ +package stackgraphv2 + +import ( + "testing" + + "github.com/stackvista/stackstate-backup-cli/internal/foundation/config" + "github.com/stretchr/testify/assert" +) + +func TestBuildRestoreEnvVars_SkipStackpacksWhenMissing(t *testing.T) { + tests := []struct { + name string + stackpacks *config.StackpacksConfig + skipStackpacksFlag bool + expectedValue string + }{ + { + name: "stackpacks nil and flag false", + stackpacks: nil, + skipStackpacksFlag: false, + expectedValue: "true", + }, + { + name: "stackpacks nil and flag true", + stackpacks: nil, + skipStackpacksFlag: true, + expectedValue: "true", + }, + { + name: "stackpacks present and flag false", + stackpacks: &config.StackpacksConfig{ + LocalStackPacksURI: "/var/stackpacks_local", + BackupDirectory: "stackpacks/", + }, + skipStackpacksFlag: false, + expectedValue: "false", + }, + { + name: "stackpacks present and flag true", + stackpacks: &config.StackpacksConfig{ + LocalStackPacksURI: "/var/stackpacks_local", + BackupDirectory: "stackpacks/", + }, + skipStackpacksFlag: true, + expectedValue: "true", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + // Set the package-level flag + skipStackpacks = tt.skipStackpacksFlag + + cfg := &config.Config{ + Stackpacks: tt.stackpacks, + Storage: config.StorageConfig{ + Service: config.ServiceConfig{Name: "storage", Port: 9000}, + }, + } + envVars := buildRestoreEnvVars("backup.graph", cfg) + + var skipValue string + for _, env := range envVars { + if env.Name == "SKIP_STACKPACKS" { + skipValue = env.Value + break + } + } + assert.Equal(t, tt.expectedValue, skipValue) + }) + } +} + +func TestBuildRestoreVolumeMounts_StackpacksLocalFileURI(t *testing.T) { + tests := []struct { + name string + stackpacks *config.StackpacksConfig + expectStackpacks bool + expectedMountPath string + }{ + { + name: "no stackpacks config", + stackpacks: nil, + expectStackpacks: false, + }, + { + name: "stackpacks with no PVC", + stackpacks: &config.StackpacksConfig{LocalStackPacksURI: "file:///var/stackpacks_local"}, + expectStackpacks: false, + }, + { + name: "stackpacks with file:// URI and PVC", + stackpacks: &config.StackpacksConfig{LocalStackPacksURI: "file:///var/stackpacks_local", PVC: "stackpacks-pvc"}, + expectStackpacks: true, + expectedMountPath: "/var/stackpacks_local", + }, + { + name: "stackpacks with non-file URI and PVC", + stackpacks: &config.StackpacksConfig{LocalStackPacksURI: "s3://my-bucket/stackpacks", PVC: "stackpacks-pvc"}, + expectStackpacks: false, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + cfg := &config.Config{Stackpacks: tt.stackpacks} + mounts := buildRestoreVolumeMounts(cfg) + + var stackpacksMount *struct{ Name, MountPath string } + for _, m := range mounts { + if m.Name == "stackpacks-local" { + stackpacksMount = &struct{ Name, MountPath string }{m.Name, m.MountPath} + break + } + } + + if tt.expectStackpacks { + assert.NotNil(t, stackpacksMount, "expected stackpacks-local volume mount to be present") + assert.Equal(t, tt.expectedMountPath, stackpacksMount.MountPath) + } else { + assert.Nil(t, stackpacksMount, "expected stackpacks-local volume mount to be absent") + } + }) + } +} diff --git a/cmd/stackgraphv2/simplejob.go b/cmd/stackgraphv2/simplejob.go new file mode 100644 index 0000000..e6ca0e2 --- /dev/null +++ b/cmd/stackgraphv2/simplejob.go @@ -0,0 +1,134 @@ +package stackgraphv2 + +import ( + "fmt" + + "github.com/stackvista/stackstate-backup-cli/internal/clients/k8s" + "github.com/stackvista/stackstate-backup-cli/internal/foundation/config" + "github.com/stackvista/stackstate-backup-cli/internal/orchestration/restore" + corev1 "k8s.io/api/core/v1" +) + +// createSimpleJob creates a Kubernetes Job for a "simple" stackgraph restore step +// (abort or backfill). These jobs share the same spec and only differ in their +// container name and command. +func createSimpleJob(k8sClient *k8s.Client, namespace, jobName, containerName, command string, config *config.Config) error { + defaultMode := int32(configMapDefaultFileMode) + + // Merge common labels with resource-specific labels + jobLabels := k8s.MergeLabels(config.Kubernetes.CommonLabels, config.Stackgraph.Restore.Job.Labels) + + // Build job spec using configuration + spec := k8s.JobSpec{ + Name: jobName, + Labels: jobLabels, + ImagePullSecrets: k8s.ConvertImagePullSecrets(config.Stackgraph.Restore.Job.ImagePullSecrets), + SecurityContext: k8s.ConvertPodSecurityContext(&config.Stackgraph.Restore.Job.SecurityContext), + NodeSelector: config.Stackgraph.Restore.Job.NodeSelector, + Tolerations: k8s.ConvertTolerations(config.Stackgraph.Restore.Job.Tolerations), + Affinity: k8s.ConvertAffinity(config.Stackgraph.Restore.Job.Affinity), + Containers: buildSimpleJobContainers(config, containerName, command), + InitContainers: buildSimpleJobInitContainers(config), + Volumes: buildSimpleJobVolumes(config, defaultMode), + } + + // Create job + _, err := k8sClient.CreateJob(namespace, spec) + if err != nil { + return fmt.Errorf("failed to create job: %w", err) + } + + return nil +} + +// buildSimpleJobEnvVars constructs environment variables shared by the abort and backfill jobs +func buildSimpleJobEnvVars(config *config.Config) []corev1.EnvVar { + storageService := config.GetStorageService() + return []corev1.EnvVar{ + {Name: "BACKUP_STACKGRAPH_BUCKET_NAME", Value: config.Stackgraph.Bucket}, + {Name: "BACKUP_STACKGRAPH_S3_PREFIX", Value: config.Stackgraph.S3Prefix}, + {Name: "S3_ENDPOINT", Value: fmt.Sprintf("%s:%d", storageService.Name, storageService.Port)}, + {Name: "STACKSTATE_BASE_URL", Value: config.GetBaseURL()}, + {Name: "RECEIVER_BASE_URL", Value: config.GetReceiverBaseURL()}, + {Name: "PLATFORM_VERSION", Value: config.GetPlatformVersion()}, + {Name: "ZOOKEEPER_QUORUM", Value: config.Stackgraph.Restore.ZookeeperQuorum}, + } +} + +// buildSimpleJobVolumeMounts constructs volume mounts shared by the abort and backfill job containers +func buildSimpleJobVolumeMounts() []corev1.VolumeMount { + return []corev1.VolumeMount{ + {Name: "backup-log", MountPath: "/opt/docker/etc_log"}, + {Name: "backup-restore-scripts", MountPath: "/backup-restore-scripts"}, + {Name: "minio-keys", MountPath: "/aws-keys"}, + } +} + +// buildSimpleJobInitContainers constructs init containers shared by the abort and backfill jobs +func buildSimpleJobInitContainers(config *config.Config) []corev1.Container { + storageService := config.GetStorageService() + return []corev1.Container{ + { + Name: "wait", + Image: config.Stackgraph.Restore.Job.WaitImage, + ImagePullPolicy: corev1.PullIfNotPresent, + Command: []string{ + "sh", + "-c", + fmt.Sprintf("/entrypoint -c %s:%d -t 300", storageService.Name, storageService.Port), + }, + SecurityContext: k8s.ConvertSecurityContext(config.Stackgraph.Restore.Job.ContainerSecurityContext), + }, + } +} + +// buildSimpleJobVolumes constructs volumes shared by the abort and backfill job pods +func buildSimpleJobVolumes(config *config.Config, defaultMode int32) []corev1.Volume { + return []corev1.Volume{ + { + Name: "backup-log", + VolumeSource: corev1.VolumeSource{ + ConfigMap: &corev1.ConfigMapVolumeSource{ + LocalObjectReference: corev1.LocalObjectReference{ + Name: config.Stackgraph.Restore.LoggingConfigConfigMapName, + }, + }, + }, + }, + { + Name: "backup-restore-scripts", + VolumeSource: corev1.VolumeSource{ + ConfigMap: &corev1.ConfigMapVolumeSource{ + LocalObjectReference: corev1.LocalObjectReference{ + Name: restore.RestoreScriptsConfigMap, + }, + DefaultMode: &defaultMode, + }, + }, + }, + { + Name: "minio-keys", + VolumeSource: corev1.VolumeSource{ + Secret: &corev1.SecretVolumeSource{ + SecretName: restore.MinioKeysSecretName, + }, + }, + }, + } +} + +// buildSimpleJobContainers constructs the main container for a simple (abort/backfill) job +func buildSimpleJobContainers(config *config.Config, name, command string) []corev1.Container { + return []corev1.Container{ + { + Name: name, + Image: config.Stackgraph.Restore.Job.Image, + ImagePullPolicy: corev1.PullIfNotPresent, + SecurityContext: k8s.ConvertSecurityContext(config.Stackgraph.Restore.Job.ContainerSecurityContext), + Command: []string{command}, + Env: buildSimpleJobEnvVars(config), + Resources: k8s.ConvertResources(config.Stackgraph.Restore.Job.Resources), + VolumeMounts: buildSimpleJobVolumeMounts(), + }, + } +} diff --git a/cmd/stackgraphv2/stackgraphv2.go b/cmd/stackgraphv2/stackgraphv2.go new file mode 100644 index 0000000..1779a10 --- /dev/null +++ b/cmd/stackgraphv2/stackgraphv2.go @@ -0,0 +1,21 @@ +package stackgraphv2 + +import ( + "github.com/spf13/cobra" + "github.com/stackvista/stackstate-backup-cli/internal/foundation/config" +) + +func Cmd(globalFlags *config.CLIGlobalFlags) *cobra.Command { + cmd := &cobra.Command{ + Use: "stackgraph-v2", + Short: "Stackgraph backup and restore (v2) operations", + } + + cmd.AddCommand(listCmd(globalFlags)) + cmd.AddCommand(restoreCmd(globalFlags)) + cmd.AddCommand(backfillCmd(globalFlags)) + cmd.AddCommand(abortCmd(globalFlags)) + cmd.AddCommand(checkAndFinalizeCmd(globalFlags)) + + return cmd +} diff --git a/cmd/victoriametrics/restore.go b/cmd/victoriametrics/restore.go index 6c27f2d..42e2471 100644 --- a/cmd/victoriametrics/restore.go +++ b/cmd/victoriametrics/restore.go @@ -230,7 +230,7 @@ func createRestoreJob(k8sClient *k8s.Client, namespace, jobName, backupFile stri func buildRestoreEnvVars(config *config.Config) []corev1.EnvVar { storageService := config.GetStorageService() return []corev1.EnvVar{ - {Name: "MINIO_ENDPOINT", Value: fmt.Sprintf("%s:%d", storageService.Name, storageService.Port)}, + {Name: "S3_ENDPOINT", Value: fmt.Sprintf("%s:%d", storageService.Name, storageService.Port)}, } } diff --git a/internal/clients/s3/filter.go b/internal/clients/s3/filter.go index 34e9e1e..839628b 100644 --- a/internal/clients/s3/filter.go +++ b/internal/clients/s3/filter.go @@ -38,6 +38,23 @@ func FilterBackupObjects(objects []s3types.Object) []Object { return filteredObjects } +// ConvertBackupObjects filters out backup part files ending with .digits. +func ConvertBackupObjects(objects []s3types.Object) []Object { + var filteredObjects []Object + + for _, obj := range objects { + key := aws.ToString(obj.Key) + + filteredObjects = append(filteredObjects, Object{ + Key: key, + LastModified: aws.ToTime(obj.LastModified), + Size: aws.ToInt64(obj.Size), + }) + } + + return filteredObjects +} + func hasNumericFileSuffix(key string) bool { if !strings.Contains(key, ".") { return false diff --git a/internal/orchestration/restore/job.go b/internal/orchestration/restore/job.go index dab9971..d5925cf 100644 --- a/internal/orchestration/restore/job.go +++ b/internal/orchestration/restore/job.go @@ -84,6 +84,9 @@ func PrintWaitingMessage(log *logger.Logger, serviceName, jobName, namespace str log.Println() log.Infof("Waiting for restore job to complete (this may take significant amount of time depending on the archive size)...") log.Println() + log.Infof("Monitoring commands:") + log.Infof(" kubectl logs --follow job/%s -n %s", jobName, namespace) + log.Println() log.Infof("You can safely interrupt this command with Ctrl+C.") log.Infof("To check status, scale up the required deployments and cleanup later, run:") log.Infof(" sts-backup %s check-and-finalize --job %s --wait -n %s", serviceName, jobName, namespace) diff --git a/internal/scripts/scripts/restore-settings-backup.sh b/internal/scripts/scripts/restore-settings-backup.sh index c33d4bc..05482b3 100644 --- a/internal/scripts/scripts/restore-settings-backup.sh +++ b/internal/scripts/scripts/restore-settings-backup.sh @@ -17,7 +17,7 @@ download_from_s3() { local dest="$3" local backup_file="$4" echo "=== Downloading Settings backup \"${backup_file}\" from bucket \"${bucket}\"..." - sts-toolbox aws s3 --endpoint "http://${MINIO_ENDPOINT}" --region minio cp "s3://${bucket}/${prefix}${backup_file}" "${dest}/${backup_file}" + sts-toolbox aws s3 --endpoint "http://${S3_ENDPOINT}" --region us-east-1 cp "s3://${bucket}/${prefix}${backup_file}" "${dest}/${backup_file}" } RESTORE_FILE="" diff --git a/internal/scripts/scripts/restore-stackgraph-backup-v2-abort.sh b/internal/scripts/scripts/restore-stackgraph-backup-v2-abort.sh new file mode 100644 index 0000000..6664a55 --- /dev/null +++ b/internal/scripts/scripts/restore-stackgraph-backup-v2-abort.sh @@ -0,0 +1,13 @@ +#!/usr/bin/env bash +set -Eeuo pipefail + +SCRIPT_DIR="$( cd "$( dirname "${BASH_SOURCE[0]}" )" &> /dev/null && pwd )" + +source $SCRIPT_DIR/restore-stackgraph-backup-v2-env.sh + +echo "=== Finalizing historic StackGraph backup data (v2)" + +/opt/docker/bin/stackstate-server -Dlogback.configurationFile=/opt/docker/etc_log/logback.xml -import-v2-abort "s3a://${BACKUP_V2_LOCATION}" +echo "=== StackGraph restore finalized" + + diff --git a/internal/scripts/scripts/restore-stackgraph-backup-v2-backfill.sh b/internal/scripts/scripts/restore-stackgraph-backup-v2-backfill.sh new file mode 100644 index 0000000..5332f9f --- /dev/null +++ b/internal/scripts/scripts/restore-stackgraph-backup-v2-backfill.sh @@ -0,0 +1,13 @@ +#!/usr/bin/env bash +set -Eeuo pipefail + +SCRIPT_DIR="$( cd "$( dirname "${BASH_SOURCE[0]}" )" &> /dev/null && pwd )" + +source $SCRIPT_DIR/restore-stackgraph-backup-v2-env.sh + +echo "=== Backfilling historic StackGraph backup data (v2)" + +/opt/docker/bin/stackstate-server -Dlogback.configurationFile=/opt/docker/etc_log/logback.xml -import-v2-backfill "s3a://${BACKUP_V2_LOCATION}" +echo "=== StackGraph restore complete" + + diff --git a/internal/scripts/scripts/restore-stackgraph-backup-v2-env.sh b/internal/scripts/scripts/restore-stackgraph-backup-v2-env.sh new file mode 100644 index 0000000..e4c0130 --- /dev/null +++ b/internal/scripts/scripts/restore-stackgraph-backup-v2-env.sh @@ -0,0 +1,30 @@ +export AWS_ACCESS_KEY_ID +AWS_ACCESS_KEY_ID="$(cat /aws-keys/accesskey)" +export AWS_SECRET_ACCESS_KEY +AWS_SECRET_ACCESS_KEY="$(cat /aws-keys/secretkey)" + +export BACKUP_V2_LOCATION="${BACKUP_STACKGRAPH_BUCKET_NAME}/${BACKUP_STACKGRAPH_S3_PREFIX}v2/" + +TYPESAFE_ESCAPED_BUCKET=$(echo "${BACKUP_STACKGRAPH_BUCKET_NAME}" | sed 's/_/___/g; s/-/__/g; s/\./_/g') +# HACK: We configure hbase here through typesafe config. However, typesafe does not support a key being both object and +# string (as in `endpoint = "string"` and `endpoint.region = "us-east-1"` +# We build a little hack there, which allows postfixing endpoint as endpoint__, which gets stripped when transforming typesafe to hbase conf +AWS_BUCKET_ENDPOINT_VAR="CONFIG_FORCE_fs_s3a_bucket_${TYPESAFE_ESCAPED_BUCKET}_endpoint___" +export "${AWS_BUCKET_ENDPOINT_VAR}=http://${S3_ENDPOINT}" +AWS_BUCKET_REGION_VAR="CONFIG_FORCE_fs_s3a_bucket_${TYPESAFE_ESCAPED_BUCKET}_endpoint_region" +export "${AWS_BUCKET_REGION_VAR}=us-east-1" +AWS_BUCKET_ACCESS_KEY_VAR="CONFIG_FORCE_fs_s3a_bucket_${TYPESAFE_ESCAPED_BUCKET}_access_key" +export "${AWS_BUCKET_ACCESS_KEY_VAR}=$(cat /aws-keys/accesskey)" +AWS_BUCKET_SECRET_KEY_VAR="CONFIG_FORCE_fs_s3a_bucket_${TYPESAFE_ESCAPED_BUCKET}_secret_key" +export "${AWS_BUCKET_SECRET_KEY_VAR}=$(cat /aws-keys/secretkey)" + +AWS_BUCKET_PATH_STYLE_VAR="CONFIG_FORCE_fs_s3a_bucket_${TYPESAFE_ESCAPED_BUCKET}_path_style_access" +export "${AWS_BUCKET_PATH_STYLE_VAR}=true" +AWS_BUCKET_CONNECTION_SSL_ENABLED_VAR="CONFIG_FORCE_fs_s3a_bucket_${TYPESAFE_ESCAPED_BUCKET}_connection_ssl_enabled" +export "${AWS_BUCKET_CONNECTION_SSL_ENABLED_VAR}=false" +AWS_BUCKET_AWS_CREDENTIALS_PROVIDER_VAR="CONFIG_FORCE_fs_s3a_bucket_${TYPESAFE_ESCAPED_BUCKET}_aws_credentials_provider" +export "${AWS_BUCKET_AWS_CREDENTIALS_PROVIDER_VAR}=org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider" + +# Increase the hbase client timeout for bulk operations. +export CONFIG_FORCE_hbase_rpc_timeout="120000" + diff --git a/internal/scripts/scripts/restore-stackgraph-backup-v2-live.sh b/internal/scripts/scripts/restore-stackgraph-backup-v2-live.sh new file mode 100644 index 0000000..8dd736a --- /dev/null +++ b/internal/scripts/scripts/restore-stackgraph-backup-v2-live.sh @@ -0,0 +1,38 @@ +#!/usr/bin/env bash +set -Eeuo pipefail + +SCRIPT_DIR="$( cd "$( dirname "${BASH_SOURCE[0]}" )" &> /dev/null && pwd )" + +TMP_DIR=/tmp-data + +source $SCRIPT_DIR/restore-stackgraph-backup-v2-env.sh + +echo "=== Importing StackGraph data (v2) from \"${BACKUP_FILE}\"..." + +/opt/docker/bin/stackstate-server -Dlogback.configurationFile=/opt/docker/etc_log/logback.xml -import-v2-live "s3a://${BACKUP_V2_LOCATION}" -backup-name "${BACKUP_FILE}" "${FORCE_DELETE}" + +echo "=== StackGraph live data loaded, please continue with backfill" + +# === StackPacks Restore === +if [ "${SKIP_STACKPACKS:-false}" == "true" ]; then + echo "=== Skipping StackPacks restore (--skip-stackpacks flag set)" +else + # Construct stackpacks backup filename from the original backup file + STACKPACKS_FILE="${BACKUP_FILE}.stackpacks.zip" + + echo "=== Checking for StackPacks backup (v2) \"${STACKPACKS_FILE}\" in bucket \"${BACKUP_STACKGRAPH_BUCKET_NAME}\"..." + + # Check if stackpacks backup exists in S3 + if ! sts-toolbox aws s3 ls --endpoint "http://${S3_ENDPOINT}" --region us-east-1 --bucket "${BACKUP_STACKGRAPH_BUCKET_NAME}" --prefix "${BACKUP_STACKGRAPH_S3_PREFIX}v2/${BACKUP_STACKGRAPH_STACKPACKS_DIR}${STACKPACKS_FILE}" 2>/dev/null | grep -q "${STACKPACKS_FILE}"; then + echo "=== ERROR: StackPacks backup \"${STACKPACKS_FILE}\" not found in S3" + exit 1 + fi + + echo "=== Downloading StackPacks backup..." + sts-toolbox aws s3 cp --endpoint "http://${S3_ENDPOINT}" --region us-east-1 "s3://${BACKUP_STACKGRAPH_BUCKET_NAME}/${BACKUP_STACKGRAPH_S3_PREFIX}v2/${BACKUP_STACKGRAPH_STACKPACKS_DIR}${STACKPACKS_FILE}" "${TMP_DIR}/${STACKPACKS_FILE}" + + echo "=== Restoring StackPacks from \"${STACKPACKS_FILE}\"..." + /opt/docker/bin/stack-packs-backup -Dlogback.configurationFile=/opt/docker/etc_log/logback.xml -restore "${TMP_DIR}/${STACKPACKS_FILE}" + echo "=== StackPacks restore complete" +fi +echo "===" diff --git a/internal/scripts/scripts/restore-stackgraph-backup.sh b/internal/scripts/scripts/restore-stackgraph-backup.sh index 13b1894..7d5b2ff 100644 --- a/internal/scripts/scripts/restore-stackgraph-backup.sh +++ b/internal/scripts/scripts/restore-stackgraph-backup.sh @@ -9,7 +9,7 @@ export AWS_SECRET_ACCESS_KEY AWS_SECRET_ACCESS_KEY="$(cat /aws-keys/secretkey)" echo "=== Downloading StackGraph backup \"${BACKUP_FILE}\" from bucket \"${BACKUP_STACKGRAPH_BUCKET_NAME}\"..." -sts-toolbox aws s3 cp --endpoint "http://${MINIO_ENDPOINT}" --region minio "s3://${BACKUP_STACKGRAPH_BUCKET_NAME}/${BACKUP_STACKGRAPH_S3_PREFIX}${BACKUP_FILE}" "${TMP_DIR}/${BACKUP_FILE}" +sts-toolbox aws s3 cp --endpoint "http://${S3_ENDPOINT}" --region us-east-1 "s3://${BACKUP_STACKGRAPH_BUCKET_NAME}/${BACKUP_STACKGRAPH_S3_PREFIX}${BACKUP_FILE}" "${TMP_DIR}/${BACKUP_FILE}" echo "=== Importing StackGraph data from \"${BACKUP_FILE}\"..." /opt/docker/bin/stackstate-server -Dlogback.configurationFile=/opt/docker/etc_log/logback.xml -import "${TMP_DIR}/${BACKUP_FILE}" "${FORCE_DELETE}" @@ -25,13 +25,13 @@ else echo "=== Checking for StackPacks backup \"${STACKPACKS_FILE}\" in bucket \"${BACKUP_STACKGRAPH_BUCKET_NAME}\"..." # Check if stackpacks backup exists in S3 - if ! sts-toolbox aws s3 ls --endpoint "http://${MINIO_ENDPOINT}" --region minio --bucket "${BACKUP_STACKGRAPH_BUCKET_NAME}" --prefix "${BACKUP_STACKGRAPH_S3_PREFIX}${BACKUP_STACKGRAPH_STACKPACKS_DIR}${STACKPACKS_FILE}" 2>/dev/null | grep -q "${STACKPACKS_FILE}"; then + if ! sts-toolbox aws s3 ls --endpoint "http://${S3_ENDPOINT}" --region us-east-1 --bucket "${BACKUP_STACKGRAPH_BUCKET_NAME}" --prefix "${BACKUP_STACKGRAPH_S3_PREFIX}${BACKUP_STACKGRAPH_STACKPACKS_DIR}${STACKPACKS_FILE}" 2>/dev/null | grep -q "${STACKPACKS_FILE}"; then echo "=== ERROR: StackPacks backup \"${STACKPACKS_FILE}\" not found in S3" exit 1 fi echo "=== Downloading StackPacks backup..." - sts-toolbox aws s3 cp --endpoint "http://${MINIO_ENDPOINT}" --region minio "s3://${BACKUP_STACKGRAPH_BUCKET_NAME}/${BACKUP_STACKGRAPH_S3_PREFIX}${BACKUP_STACKGRAPH_STACKPACKS_DIR}${STACKPACKS_FILE}" "${TMP_DIR}/${STACKPACKS_FILE}" + sts-toolbox aws s3 cp --endpoint "http://${S3_ENDPOINT}" --region us-east-1 "s3://${BACKUP_STACKGRAPH_BUCKET_NAME}/${BACKUP_STACKGRAPH_S3_PREFIX}${BACKUP_STACKGRAPH_STACKPACKS_DIR}${STACKPACKS_FILE}" "${TMP_DIR}/${STACKPACKS_FILE}" echo "=== Restoring StackPacks from \"${STACKPACKS_FILE}\"..." /opt/docker/bin/stack-packs-backup -Dlogback.configurationFile=/opt/docker/etc_log/logback.xml -restore "${TMP_DIR}/${STACKPACKS_FILE}" diff --git a/internal/scripts/scripts/restore-victoria-metrics-backup.sh b/internal/scripts/scripts/restore-victoria-metrics-backup.sh index ee2e4f1..ceb3d34 100644 --- a/internal/scripts/scripts/restore-victoria-metrics-backup.sh +++ b/internal/scripts/scripts/restore-victoria-metrics-backup.sh @@ -14,4 +14,4 @@ AWS_ACCESS_KEY_ID="$(cat /aws-keys/accesskey)" export AWS_SECRET_ACCESS_KEY AWS_SECRET_ACCESS_KEY="$(cat /aws-keys/secretkey)" -/vmrestore-prod -storageDataPath=/storage -src="s3://$S3_LOCATION" -customS3Endpoint="http://$MINIO_ENDPOINT" -httpListenAddr "$METRICS_ADDR" +/vmrestore-prod -storageDataPath=/storage -src="s3://$S3_LOCATION" -customS3Endpoint="http://$S3_ENDPOINT" -httpListenAddr "$METRICS_ADDR" From 1a76699f12b2d7de236e00078f13be8b9264af11 Mon Sep 17 00:00:00 2001 From: Bram Schuur Date: Mon, 20 Jul 2026 09:44:34 +0200 Subject: [PATCH 2/9] STAC-24630: Work out review comments --- cmd/stackgraphv2/backfill.go | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/cmd/stackgraphv2/backfill.go b/cmd/stackgraphv2/backfill.go index f0a8ff9..6ec8566 100644 --- a/cmd/stackgraphv2/backfill.go +++ b/cmd/stackgraphv2/backfill.go @@ -21,9 +21,11 @@ const ( func backfillCmd(globalFlags *config.CLIGlobalFlags) *cobra.Command { cmd := &cobra.Command{ Use: "backfill", - Short: "Complete a restore by backfilling old data", - Long: "Complete a restore by backfilling old data. This process can run while the system is already up and running " + - "After the backfill is done, the restore process is complete. ", + Short: "Complete an interrupted restore by backfilling old data", + Long: "Complete an interrupted restore by backfilling old data. This process can run while the system is already up and running." + + "After the backfill is done, the restore process is complete. " + + "During a normal restore the restore command will take care of backfilling the data, but if that gets interrupted (Ctrl-C) " + + "or somehow crashes due to cluster instability, this command allows retrying the backfill portion of the restore.", Run: func(_ *cobra.Command, _ []string) { cmdutils.Run(globalFlags, runBackfill, cmdutils.StorageIsRequired) }, From 24b1933ef3ebda30b8892cef551c156a434ea391 Mon Sep 17 00:00:00 2001 From: Bram Schuur Date: Tue, 21 Jul 2026 16:31:32 +0200 Subject: [PATCH 3/9] STAC-24630: Add additional message when jobs fail/get interrupted --- cmd/stackgraphv2/abort.go | 11 ++++-- cmd/stackgraphv2/backfill.go | 5 ++- cmd/stackgraphv2/check_and_finalize.go | 12 ++++--- cmd/stackgraphv2/logging.go | 42 ++++++++++++++++++++++ cmd/stackgraphv2/restore.go | 5 +-- internal/orchestration/restore/finalize.go | 1 + 6 files changed, 64 insertions(+), 12 deletions(-) create mode 100644 cmd/stackgraphv2/logging.go diff --git a/cmd/stackgraphv2/abort.go b/cmd/stackgraphv2/abort.go index 24fa11f..492a210 100644 --- a/cmd/stackgraphv2/abort.go +++ b/cmd/stackgraphv2/abort.go @@ -14,7 +14,7 @@ import ( ) const ( - abortNameTemplate = "stackgraph-abort-v2" + abortNameTemplate = "stackgraph-v2-abort" abortScript = "/backup-restore-scripts/restore-stackgraph-backup-v2-abort.sh" ) @@ -52,7 +52,14 @@ func runAbort(appCtx *app.Context) error { appCtx.Logger.Successf("Abort job created: %s", jobName) - return waitAndCleanupAbortJob(appCtx.K8sClient, appCtx.Namespace, jobName, appCtx.Logger) + err := waitAndCleanupAbortJob(appCtx.K8sClient, appCtx.Namespace, jobName, appCtx.Logger) + + if err != nil { + logAfterJobResult(appCtx.Logger, checkJobName, false) + return err + } + + return nil } // waitAndCleanupAbortJob waits for job completion and cleans up resources diff --git a/cmd/stackgraphv2/backfill.go b/cmd/stackgraphv2/backfill.go index 6ec8566..c5eaa8d 100644 --- a/cmd/stackgraphv2/backfill.go +++ b/cmd/stackgraphv2/backfill.go @@ -14,7 +14,7 @@ import ( ) const ( - backfillNameTemplate = "stackgraph-backfill-v2" + backfillNameTemplate = "stackgraph-v2-backfill" backfillScript = "/backup-restore-scripts/restore-stackgraph-backup-v2-backfill.sh" ) @@ -56,8 +56,7 @@ func runBackfill(appCtx *app.Context) error { err := waitAndCleanupBackfillJob(appCtx.K8sClient, appCtx.Namespace, jobName, appCtx.Logger) if err != nil { - appCtx.Logger.Println() - appCtx.Logger.Infof("Backfill failed. It is possible to restart the backfill without starting a complete restore.") + logAfterJobResult(appCtx.Logger, checkJobName, false) return err } diff --git a/cmd/stackgraphv2/check_and_finalize.go b/cmd/stackgraphv2/check_and_finalize.go index 39cd794..9f266fe 100644 --- a/cmd/stackgraphv2/check_and_finalize.go +++ b/cmd/stackgraphv2/check_and_finalize.go @@ -21,15 +21,15 @@ func checkAndFinalizeCmd(globalFlags *config.CLIGlobalFlags) *cobra.Command { Short: "Check and finalize a Stackgraph restore (v2) job", Long: `Check the status of a background Stackgraph restore job and clean up resources. -This command is useful when a restore job was started with --background flag or was interrupted (Ctrl+C). +This command is useful when a restore/backfill/abort job was interrupted (Ctrl+C). It will check the job status, print logs if it failed, and clean up the job and PVC resources. Examples: # Check job status without waiting - sts-backup stackgraph-v2 check-and-finalize --job stackgraph-restore-20250128t143000 -n my-namespace + sts-backup stackgraph-v2 check-and-finalize --job stackgraph-v2-restore-20250128t143000 -n my-namespace # Wait for job completion and cleanup - sts-backup stackgraph-v2 check-and-finalize --job stackgraph-restore-20250128t143000 --wait -n my-namespace`, + sts-backup stackgraph-v2 check-and-finalize --job stackgraph-v2-restore-20250128t143000 --wait -n my-namespace`, Run: func(_ *cobra.Command, _ []string) { cmdutils.Run(globalFlags, runCheckAndFinalize, cmdutils.StorageIsRequired) }, @@ -43,11 +43,11 @@ Examples: } func runCheckAndFinalize(appCtx *app.Context) error { - return restore.CheckAndFinalize(restore.CheckAndFinalizeParams{ + err := restore.CheckAndFinalize(restore.CheckAndFinalizeParams{ K8sClient: appCtx.K8sClient, Namespace: appCtx.Namespace, JobName: checkJobName, - ServiceName: "stackgraph", + ServiceName: "stackgraph-v2", ScaleUpFn: scale.ScaleUpAndReleaseLock, ScaleDownFn: scale.ScaleDown, ScaleSelector: appCtx.Config.Stackgraph.Restore.ScaleDownLabelSelector, @@ -55,4 +55,6 @@ func runCheckAndFinalize(appCtx *app.Context) error { WaitForJob: waitForJob, Log: appCtx.Logger, }) + logAfterJobResult(appCtx.Logger, checkJobName, err == nil) + return err } diff --git a/cmd/stackgraphv2/logging.go b/cmd/stackgraphv2/logging.go new file mode 100644 index 0000000..dfbca44 --- /dev/null +++ b/cmd/stackgraphv2/logging.go @@ -0,0 +1,42 @@ +package stackgraphv2 + +import ( + "strings" + + "github.com/stackvista/stackstate-backup-cli/internal/foundation/logger" +) + +func logAfterJobResult(log *logger.Logger, jobName string, success bool) { + switch { + case strings.HasPrefix(jobName, restoreNameTemplate): + if success { + log.Println() + log.Infof("Job '%s' has successfully restored the live data. The system is running again.", jobName) + log.Infof("Now run `sts-backup stackgraph-v2 backfill` to load the historic data.") + } else { + log.Println() + log.Infof("Job '%s' has failed to restore the live data. The system was brought back up but the data was not restored.", jobName) + log.Infof("Rerun your `sts-backup stackgraph-v2 restore ...` command toretry restoring data.") + } + case strings.HasPrefix(jobName, backfillNameTemplate): + if success { + log.Println() + log.Infof("Job '%s' has successfully backfilled the historic data. The restore is complete.", jobName) + } else { + log.Println() + log.Infof("Job '%s' has failed to backfill all historic data.", jobName) + log.Infof("Run `sts-backup stackgraph-v2 backfill` to retry restoring old data, or") + log.Infof("run `sts-backup stackgraph-v2 abort` to leave the restored data as-is and continue with") + log.Infof("the data currently in the system.") + } + case strings.HasPrefix(jobName, abortNameTemplate): + if success { + log.Println() + log.Infof("Job '%s' has successfully aborted. Historic data is incomplete but the system is functional.", jobName) + } else { + log.Println() + log.Infof("Job '%s' has failed to abort the restore.", jobName) + log.Infof("Rerun `sts-backup stackgraph-v2 abort` again to try again.") + } + } +} diff --git a/cmd/stackgraphv2/restore.go b/cmd/stackgraphv2/restore.go index c42ec8d..e03a2d3 100644 --- a/cmd/stackgraphv2/restore.go +++ b/cmd/stackgraphv2/restore.go @@ -24,7 +24,7 @@ import ( ) const ( - jobNameTemplate = "stackgraph-restore-v2" + restoreNameTemplate = "stackgraph-v2-restore" configMapDefaultFileMode = 0755 purgeStackgraphDataFlag = "-force" ) @@ -136,7 +136,7 @@ func liveRestore(appCtx *app.Context) error { appCtx.Logger.Println() appCtx.Logger.Infof("Creating restore of live data job for backup: %s", backupFile) - jobName := fmt.Sprintf("%s-%s", jobNameTemplate, time.Now().Format("20060102t150405")) + jobName := fmt.Sprintf("%s-%s", restoreNameTemplate, time.Now().Format("20060102t150405")) if err = createRestoreJob(appCtx.K8sClient, appCtx.Namespace, jobName, backupFile, appCtx.Config); err != nil { return fmt.Errorf("failed to create restore job: %w", err) @@ -146,6 +146,7 @@ func liveRestore(appCtx *app.Context) error { err = waitAndCleanupRestoreJob(appCtx.K8sClient, appCtx.Namespace, jobName, appCtx.Logger) if err != nil { + logAfterJobResult(appCtx.Logger, checkJobName, false) return err } diff --git a/internal/orchestration/restore/finalize.go b/internal/orchestration/restore/finalize.go index ae3ba1f..a737b0e 100644 --- a/internal/orchestration/restore/finalize.go +++ b/internal/orchestration/restore/finalize.go @@ -117,6 +117,7 @@ type CheckAndFinalizeParams struct { // CheckAndFinalize checks the status of a background restore job and cleans up resources // This is useful when a restore job was started with --background flag or was interrupted (Ctrl+C) +// Returns whether the job succeeded func CheckAndFinalize(params CheckAndFinalizeParams) error { // Get job params.Log.Infof("Checking status of job: %s", params.JobName) From 5c9a70e40303891dc6fa3c54dcdbc128add78337 Mon Sep 17 00:00:00 2001 From: Bram Schuur Date: Wed, 22 Jul 2026 11:22:28 +0200 Subject: [PATCH 4/9] STAC-24630: Make sure to scale up after a failed restore --- internal/orchestration/restore/finalize.go | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/internal/orchestration/restore/finalize.go b/internal/orchestration/restore/finalize.go index a737b0e..7c7cf16 100644 --- a/internal/orchestration/restore/finalize.go +++ b/internal/orchestration/restore/finalize.go @@ -85,6 +85,11 @@ func WaitAndFinalize(params WaitAndFinalizeParams) error { // Still cleanup even if failed params.Log.Println() _ = CleanupResources(params.K8sClient, params.Namespace, params.JobName, "", params.Log, params.CleanupPVC) + + // Scale up deployments that were scaled down before restore + if errScaleUp := params.ScaleUpFn(params.K8sClient, params.Namespace, params.ScaleSelector, params.Log); err != nil { + params.Log.Warningf("Failed to scale up workload: %v", errScaleUp) + } return err } From 14eb630ff96c9b24e58bc9f680b39b9775644580 Mon Sep 17 00:00:00 2001 From: Bram Schuur Date: Wed, 22 Jul 2026 13:29:19 +0200 Subject: [PATCH 5/9] STAC-24630: Some small fixes --- cmd/stackgraphv2/abort.go | 2 +- cmd/stackgraphv2/backfill.go | 2 +- cmd/stackgraphv2/restore.go | 2 +- 3 files changed, 3 insertions(+), 3 deletions(-) diff --git a/cmd/stackgraphv2/abort.go b/cmd/stackgraphv2/abort.go index 492a210..5470e63 100644 --- a/cmd/stackgraphv2/abort.go +++ b/cmd/stackgraphv2/abort.go @@ -55,7 +55,7 @@ func runAbort(appCtx *app.Context) error { err := waitAndCleanupAbortJob(appCtx.K8sClient, appCtx.Namespace, jobName, appCtx.Logger) if err != nil { - logAfterJobResult(appCtx.Logger, checkJobName, false) + logAfterJobResult(appCtx.Logger, jobName, false) return err } diff --git a/cmd/stackgraphv2/backfill.go b/cmd/stackgraphv2/backfill.go index c5eaa8d..eccb517 100644 --- a/cmd/stackgraphv2/backfill.go +++ b/cmd/stackgraphv2/backfill.go @@ -56,7 +56,7 @@ func runBackfill(appCtx *app.Context) error { err := waitAndCleanupBackfillJob(appCtx.K8sClient, appCtx.Namespace, jobName, appCtx.Logger) if err != nil { - logAfterJobResult(appCtx.Logger, checkJobName, false) + logAfterJobResult(appCtx.Logger, jobName, false) return err } diff --git a/cmd/stackgraphv2/restore.go b/cmd/stackgraphv2/restore.go index e03a2d3..3a869b5 100644 --- a/cmd/stackgraphv2/restore.go +++ b/cmd/stackgraphv2/restore.go @@ -146,7 +146,7 @@ func liveRestore(appCtx *app.Context) error { err = waitAndCleanupRestoreJob(appCtx.K8sClient, appCtx.Namespace, jobName, appCtx.Logger) if err != nil { - logAfterJobResult(appCtx.Logger, checkJobName, false) + logAfterJobResult(appCtx.Logger, jobName, false) return err } From eec5a85828c71c29a1ced580f99b089f4eb7c6e2 Mon Sep 17 00:00:00 2001 From: Vladimir Iliakov Date: Mon, 24 Aug 2026 16:53:34 +0200 Subject: [PATCH 6/9] STAC-24630: order stackgraph backups by name, not S3 LastModified Syncing a bucket restamps LastModified, so synced archives tie and `restore --latest` picked an arbitrary one. --- cmd/stackgraph/list.go | 6 +---- cmd/stackgraph/restore.go | 6 +---- cmd/stackgraphv2/list.go | 6 +---- cmd/stackgraphv2/restore.go | 6 +---- internal/clients/s3/filter.go | 12 ++++++++++ internal/clients/s3/filter_test.go | 37 ++++++++++++++++++++++++++++++ 6 files changed, 53 insertions(+), 20 deletions(-) diff --git a/cmd/stackgraph/list.go b/cmd/stackgraph/list.go index 50c0ed7..2368265 100644 --- a/cmd/stackgraph/list.go +++ b/cmd/stackgraph/list.go @@ -3,7 +3,6 @@ package stackgraph import ( "context" "fmt" - "sort" "strings" "github.com/aws/aws-sdk-go-v2/aws" @@ -74,10 +73,7 @@ func runList(appCtx *app.Context) error { return fmt.Errorf("failed to filter objects: %w", err) } - // Sort by LastModified time (most recent first) - sort.Slice(filteredObjects, func(i, j int) bool { - return filteredObjects[i].LastModified.After(filteredObjects[j].LastModified) - }) + s3client.SortByKeyDescending(filteredObjects) if len(filteredObjects) == 0 { appCtx.Formatter.PrintMessage("No backups found") diff --git a/cmd/stackgraph/restore.go b/cmd/stackgraph/restore.go index 4c02c06..fa82e49 100644 --- a/cmd/stackgraph/restore.go +++ b/cmd/stackgraph/restore.go @@ -3,7 +3,6 @@ package stackgraph import ( "context" "fmt" - "sort" "strconv" "strings" "time" @@ -193,10 +192,7 @@ func getLatestBackup(k8sClient *k8s.Client, namespace string, config *config.Con return "", fmt.Errorf("no backups found in bucket %s", bucket) } - // Sort by LastModified time (most recent first) - sort.Slice(filteredObjects, func(i, j int) bool { - return filteredObjects[i].LastModified.After(filteredObjects[j].LastModified) - }) + s3client.SortByKeyDescending(filteredObjects) return filteredObjects[0].Key, nil } diff --git a/cmd/stackgraphv2/list.go b/cmd/stackgraphv2/list.go index 7700379..f3149d2 100644 --- a/cmd/stackgraphv2/list.go +++ b/cmd/stackgraphv2/list.go @@ -3,7 +3,6 @@ package stackgraphv2 import ( "context" "fmt" - "sort" "strings" "github.com/aws/aws-sdk-go-v2/aws" @@ -75,10 +74,7 @@ func runList(appCtx *app.Context) error { return fmt.Errorf("failed to filter objects: %w", err) } - // Sort by LastModified time (most recent first) - sort.Slice(filteredObjects, func(i, j int) bool { - return filteredObjects[i].LastModified.After(filteredObjects[j].LastModified) - }) + s3client.SortByKeyDescending(filteredObjects) if len(filteredObjects) == 0 { appCtx.Formatter.PrintMessage("No backups found") diff --git a/cmd/stackgraphv2/restore.go b/cmd/stackgraphv2/restore.go index 3a869b5..be2f1ab 100644 --- a/cmd/stackgraphv2/restore.go +++ b/cmd/stackgraphv2/restore.go @@ -3,7 +3,6 @@ package stackgraphv2 import ( "context" "fmt" - "sort" "strconv" "strings" "time" @@ -209,10 +208,7 @@ func getLatestBackup(k8sClient *k8s.Client, namespace string, config *config.Con return "", fmt.Errorf("no backups found in bucket %s", bucket) } - // Sort by LastModified time (most recent first) - sort.Slice(filteredObjects, func(i, j int) bool { - return filteredObjects[i].LastModified.After(filteredObjects[j].LastModified) - }) + s3client.SortByKeyDescending(filteredObjects) return filteredObjects[0].Key, nil } diff --git a/internal/clients/s3/filter.go b/internal/clients/s3/filter.go index 839628b..73534cb 100644 --- a/internal/clients/s3/filter.go +++ b/internal/clients/s3/filter.go @@ -3,6 +3,7 @@ package s3 import ( "fmt" "regexp" + "sort" "strings" "time" @@ -17,6 +18,17 @@ type Object struct { Size int64 } +// SortByKeyDescending orders objects newest-first for backup names that embed a +// lexicographically sortable timestamp (sts-backup-YYYYMMDD-HHMM...). +// LastModified cannot be used as the ordering key: copying a bucket restamps it, +// so `aws s3 sync` leaves every synced archive with the same value and the real +// backup order is lost. +func SortByKeyDescending(objects []Object) { + sort.Slice(objects, func(i, j int) bool { + return objects[i].Key > objects[j].Key + }) +} + // FilterBackupObjects filters out backup part files ending with .digits. func FilterBackupObjects(objects []s3types.Object) []Object { var filteredObjects []Object diff --git a/internal/clients/s3/filter_test.go b/internal/clients/s3/filter_test.go index 03eb3b6..c9b72d0 100644 --- a/internal/clients/s3/filter_test.go +++ b/internal/clients/s3/filter_test.go @@ -375,3 +375,40 @@ func TestFilterByPrefixAndRegex_PreservesMetadata(t *testing.T) { assert.Equal(t, int64(1234567890), result[0].Size) assert.Equal(t, now.Unix(), result[0].LastModified.Unix()) } + +// TestSortByKeyDescending_IdenticalLastModified covers the case that broke `restore --latest`: +// syncing a bucket gives every copied archive the same LastModified, so the ordering has to +// come from the timestamp in the key. +func TestSortByKeyDescending_IdenticalLastModified(t *testing.T) { + synced := time.Now() + + objects := []Object{ + {Key: "sts-backup-20260724-0642.graph.v2", LastModified: synced}, + {Key: "sts-backup-20260724-0701.graph.v2", LastModified: synced}, + {Key: "sts-backup-20260723-0300.graph.v2", LastModified: synced}, + } + + SortByKeyDescending(objects) + + assert.Equal(t, []string{ + "sts-backup-20260724-0701.graph.v2", + "sts-backup-20260724-0642.graph.v2", + "sts-backup-20260723-0300.graph.v2", + }, []string{objects[0].Key, objects[1].Key, objects[2].Key}) +} + +// TestSortByKeyDescending_IgnoresLastModified tests that a restamped LastModified which +// contradicts the backup name does not affect the ordering. +func TestSortByKeyDescending_IgnoresLastModified(t *testing.T) { + now := time.Now() + + objects := []Object{ + {Key: "sts-backup-20260101-0300.graph", LastModified: now}, + {Key: "sts-backup-20260201-0300.graph", LastModified: now.Add(-48 * time.Hour)}, + } + + SortByKeyDescending(objects) + + assert.Equal(t, "sts-backup-20260201-0300.graph", objects[0].Key) + assert.Equal(t, "sts-backup-20260101-0300.graph", objects[1].Key) +} From b692f9e3c30bbcf242442feb2c16704eee20203a Mon Sep 17 00:00:00 2001 From: Vladimir Iliakov Date: Fri, 28 Aug 2026 11:11:31 +0200 Subject: [PATCH 7/9] STAC-25639: retry transient port-forward setup failures (#35) --- .../orchestration/portforward/portforward.go | 66 +++++++++-- .../portforward/portforward_test.go | 105 ++++++++---------- 2 files changed, 100 insertions(+), 71 deletions(-) diff --git a/internal/orchestration/portforward/portforward.go b/internal/orchestration/portforward/portforward.go index 3225603..8c03138 100644 --- a/internal/orchestration/portforward/portforward.go +++ b/internal/orchestration/portforward/portforward.go @@ -2,11 +2,20 @@ package portforward import ( "fmt" + "time" - "github.com/stackvista/stackstate-backup-cli/internal/clients/k8s" "github.com/stackvista/stackstate-backup-cli/internal/foundation/logger" ) +const ( + portForwardMaxAttempts = 3 + portForwardRetryDelay = 2 * time.Second +) + +type portForwardClient interface { + PortForwardService(namespace, serviceName string, remotePort int) (chan struct{}, int, error) +} + // Conn contains the channels needed to manage a port-forward connection type Conn struct { StopChan chan struct{} @@ -18,23 +27,58 @@ type Conn struct { // It returns a Conn containing the stop channel and the actual local port. // The caller is responsible for closing the StopChan when done. func SetupPortForward( - k8sClient *k8s.Client, + k8sClient portForwardClient, namespace string, serviceName string, remotePort int, log *logger.Logger, +) (*Conn, error) { + return setupPortForward( + k8sClient, + namespace, + serviceName, + remotePort, + log, + portForwardMaxAttempts, + portForwardRetryDelay, + ) +} + +func setupPortForward( + k8sClient portForwardClient, + namespace string, + serviceName string, + remotePort int, + log *logger.Logger, + maxAttempts int, + retryDelay time.Duration, ) (*Conn, error) { log.Infof("Setting up port-forward to %s:%d in namespace %s...", serviceName, remotePort, namespace) - stopChan, actualLocalPort, err := k8sClient.PortForwardService(namespace, serviceName, remotePort) - if err != nil { - return nil, fmt.Errorf("failed to setup port-forward: %w", err) - } + for attempt := 1; attempt <= maxAttempts; attempt++ { + stopChan, actualLocalPort, err := k8sClient.PortForwardService(namespace, serviceName, remotePort) + if err == nil { + log.Successf("Port-forward established on localhost:%d", actualLocalPort) - log.Successf("Port-forward established on localhost:%d", actualLocalPort) + return &Conn{ + StopChan: stopChan, + LocalPort: actualLocalPort, + }, nil + } + + if attempt == maxAttempts { + return nil, fmt.Errorf("failed to setup port-forward after %d attempts: %w", attempt, err) + } + + log.Warningf( + "Port-forward attempt %d/%d failed: %v; retrying in %s", + attempt, + maxAttempts, + err, + retryDelay, + ) + time.Sleep(retryDelay) + } - return &Conn{ - StopChan: stopChan, - LocalPort: actualLocalPort, - }, nil + return nil, fmt.Errorf("failed to setup port-forward") } diff --git a/internal/orchestration/portforward/portforward_test.go b/internal/orchestration/portforward/portforward_test.go index af3cffe..f153b8c 100644 --- a/internal/orchestration/portforward/portforward_test.go +++ b/internal/orchestration/portforward/portforward_test.go @@ -1,83 +1,68 @@ package portforward import ( + "errors" "testing" - corev1 "k8s.io/api/core/v1" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/client-go/kubernetes/fake" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" - "github.com/stackvista/stackstate-backup-cli/internal/clients/k8s" "github.com/stackvista/stackstate-backup-cli/internal/foundation/logger" ) -func TestSetupPortForward_ServiceNotFound(t *testing.T) { - fakeClientset := fake.NewSimpleClientset() - client := k8s.NewTestClient(fakeClientset) - log := logger.New(true, false) +type portForwardResult struct { + stopChan chan struct{} + localPort int + err error +} - _, err := SetupPortForward(client, "default", "nonexistent-service", 9200, log) - if err == nil { - t.Fatal("expected error for nonexistent service, got nil") - } +type fakeClient struct { + results []portForwardResult + calls int } -func TestSetupPortForward_NoPodsFound(t *testing.T) { - fakeClientset := fake.NewSimpleClientset( - &corev1.Service{ - ObjectMeta: metav1.ObjectMeta{ - Name: "test-service", - Namespace: "default", - }, - Spec: corev1.ServiceSpec{ - Selector: map[string]string{ - "app": "test", - }, - }, +func (f *fakeClient) PortForwardService(_, _ string, _ int) (chan struct{}, int, error) { + result := f.results[f.calls] + f.calls++ + return result.stopChan, result.localPort, result.err +} + +func TestSetupPortForward_RetriesTransientFailure(t *testing.T) { + stopChan := make(chan struct{}) + client := &fakeClient{ + results: []portForwardResult{ + {err: errors.New("connection reset by peer")}, + {err: errors.New("error upgrading connection")}, + {stopChan: stopChan, localPort: 43210}, }, - ) - client := k8s.NewTestClient(fakeClientset) + } log := logger.New(true, false) - _, err := SetupPortForward(client, "default", "test-service", 9200, log) - if err == nil { - t.Fatal("expected error for service with no pods, got nil") - } + result, err := setupPortForward(client, "default", "test-service", 9200, log, 3, 0) + + require.NoError(t, err) + assert.Equal(t, 3, client.calls) + assert.Equal(t, stopChan, result.StopChan) + assert.Equal(t, 43210, result.LocalPort) } -func TestSetupPortForward_NoRunningPods(t *testing.T) { - fakeClientset := fake.NewSimpleClientset( - &corev1.Service{ - ObjectMeta: metav1.ObjectMeta{ - Name: "test-service", - Namespace: "default", - }, - Spec: corev1.ServiceSpec{ - Selector: map[string]string{ - "app": "test", - }, - }, - }, - &corev1.Pod{ - ObjectMeta: metav1.ObjectMeta{ - Name: "test-pod", - Namespace: "default", - Labels: map[string]string{ - "app": "test", - }, - }, - Status: corev1.PodStatus{ - Phase: corev1.PodPending, - }, +func TestSetupPortForward_ReturnsLastErrorAfterRetries(t *testing.T) { + client := &fakeClient{ + results: []portForwardResult{ + {err: errors.New("first failure")}, + {err: errors.New("second failure")}, + {err: errors.New("last failure")}, }, - ) - client := k8s.NewTestClient(fakeClientset) + } log := logger.New(true, false) - _, err := SetupPortForward(client, "default", "test-service", 9200, log) - if err == nil { - t.Fatal("expected error for service with no running pods, got nil") - } + result, err := setupPortForward(client, "default", "test-service", 9200, log, 3, 0) + + require.Error(t, err) + assert.Nil(t, result) + assert.Equal(t, 3, client.calls) + assert.ErrorContains(t, err, "failed to setup port-forward after 3 attempts") + assert.ErrorContains(t, err, "last failure") } func TestConn_Structure(t *testing.T) { From 0b157e1b602fea85bbe0c606edad2bf1f917a47c Mon Sep 17 00:00:00 2001 From: Vladimir Iliakov Date: Fri, 28 Aug 2026 12:42:08 +0200 Subject: [PATCH 8/9] STAC-25639: elasticsearch restore wait until all shards and indices are healthy --- cmd/elasticsearch/check_and_finalize.go | 287 +++++++++++++--- cmd/elasticsearch/check_and_finalize_test.go | 312 ++++++++++++++++++ cmd/elasticsearch/restore.go | 14 +- cmd/elasticsearch/restore_test.go | 4 +- internal/clients/elasticsearch/client.go | 59 ++-- internal/clients/elasticsearch/client_test.go | 90 +++++ internal/clients/elasticsearch/interface.go | 2 +- internal/orchestration/restore/apirestore.go | 5 +- 8 files changed, 686 insertions(+), 87 deletions(-) create mode 100644 cmd/elasticsearch/check_and_finalize_test.go diff --git a/cmd/elasticsearch/check_and_finalize.go b/cmd/elasticsearch/check_and_finalize.go index f01ec44..cd5cf5c 100644 --- a/cmd/elasticsearch/check_and_finalize.go +++ b/cmd/elasticsearch/check_and_finalize.go @@ -2,12 +2,14 @@ package elasticsearch import ( "fmt" + "time" "github.com/spf13/cobra" "github.com/stackvista/stackstate-backup-cli/cmd/cmdutils" "github.com/stackvista/stackstate-backup-cli/internal/app" es "github.com/stackvista/stackstate-backup-cli/internal/clients/elasticsearch" "github.com/stackvista/stackstate-backup-cli/internal/foundation/config" + "github.com/stackvista/stackstate-backup-cli/internal/foundation/logger" "github.com/stackvista/stackstate-backup-cli/internal/orchestration/portforward" "github.com/stackvista/stackstate-backup-cli/internal/orchestration/restore" "github.com/stackvista/stackstate-backup-cli/internal/orchestration/scale" @@ -15,8 +17,10 @@ import ( // Check-and-finalize command flags var ( - checkOperationID string - checkWait bool + checkOperationID string + checkWait bool + checkNoProgressIn time.Duration + checkFinalizeOnly bool ) func checkAndFinalizeCmd(globalFlags *config.CLIGlobalFlags) *cobra.Command { @@ -25,19 +29,38 @@ func checkAndFinalizeCmd(globalFlags *config.CLIGlobalFlags) *cobra.Command { Short: "Check restore status and finalize if complete", Long: `Check the status of a restore operation and perform finalization (scale up deployments) if complete. If the restore is still running and --wait is specified, wait for completion before finalizing.`, + PreRunE: func(_ *cobra.Command, _ []string) error { + if !checkFinalizeOnly && checkOperationID == "" { + return fmt.Errorf("--operation-id is required unless --finalize-only is set") + } + + return nil + }, Run: func(_ *cobra.Command, _ []string) { cmdutils.Run(globalFlags, runCheckAndFinalize, cmdutils.StorageIsRequired) }, } - cmd.Flags().StringVar(&checkOperationID, "operation-id", "", "Operation ID of the restore operation (required)") + cmd.Flags().StringVar(&checkOperationID, "operation-id", "", + "Snapshot name of the restore operation (required unless --finalize-only)") cmd.Flags().BoolVar(&checkWait, "wait", false, "Wait for restore to complete if still running") - _ = cmd.MarkFlagRequired("operation-id") + cmd.Flags().DurationVar(&checkNoProgressIn, "no-progress-timeout", defaultNoProgressTimeout, noProgressTimeoutUsage) + cmd.Flags().BoolVar(&checkFinalizeOnly, "finalize-only", false, + "Scale the deployments back up and release the restore lock without checking restore status. "+ + "Use when the snapshot can no longer be read and the restore is known to be finished") + cmd.MarkFlagsMutuallyExclusive("finalize-only", "wait") return cmd } func runCheckAndFinalize(appCtx *app.Context) error { + // Finalizing needs Kubernetes only, so it stays available when Elasticsearch or the snapshot + // cannot be reached at all - otherwise nothing in the CLI can release the restore lock. + if checkFinalizeOnly { + appCtx.Logger.Warningf("Skipping the restore status check on request; finalizing") + return finalizeRestore(appCtx) + } + // Setup port-forward to Elasticsearch serviceName := appCtx.Config.Elasticsearch.Service.Name remotePort := appCtx.Config.Elasticsearch.Service.Port @@ -56,13 +79,223 @@ func runCheckAndFinalize(appCtx *app.Context) error { repository := appCtx.Config.Elasticsearch.Restore.Repository - return checkAndFinalize(esClient, appCtx, repository, checkOperationID, checkWait) + return checkAndFinalize(esClient, appCtx, repository, checkOperationID, checkWait, checkNoProgressIn) +} + +const ( + // defaultNoProgressTimeout bounds the wait so a stuck restore fails instead of polling until + // the caller's own timeout with the workloads scaled down and the restore lock held. Generous + // on purpose: a spurious failure aborts a working restore, while the only cost of waiting too + // long is a later failure. Progress is counted per primary shard, so tripping this means no + // shard at all completed in the window. + defaultNoProgressTimeout = 2 * time.Hour + + noProgressTimeoutUsage = "Fail if no primary shard finishes restoring within this duration. " + + "Bounds inactivity, not total restore time; 0 waits indefinitely" + + // restoreStatusMaxErrors tolerates a port-forward dropping mid-restore. Polling now spans the + // whole restore, so a pod restart must not end a restore that is still running server-side. + restoreStatusMaxErrors = 5 +) + +// reconnectingHealthClient rebuilds the port-forward and Elasticsearch client after a failed +// health check. It starts from the caller's client so the happy path opens no extra port-forward. +type reconnectingHealthClient struct { + appCtx *app.Context + client indicesHealthGetter + pf *portforward.Conn +} + +func (r *reconnectingHealthClient) GetIndicesHealth() (map[string]es.IndexHealth, error) { + if r.client == nil { + if err := r.connect(); err != nil { + return nil, err + } + } + + health, err := r.client.GetIndicesHealth() + if err != nil { + // The tunnel is bound to a fixed local port, so a broken one never recovers. Drop it and + // let the next poll dial a fresh pod. + r.disconnect() + return nil, err + } + + return health, nil +} + +func (r *reconnectingHealthClient) connect() error { + pf, err := portforward.SetupPortForward( + r.appCtx.K8sClient, + r.appCtx.Namespace, + r.appCtx.Config.Elasticsearch.Service.Name, + r.appCtx.Config.Elasticsearch.Service.Port, + r.appCtx.Logger, + ) + if err != nil { + return err + } + + client, err := r.appCtx.NewESClient(pf.LocalPort) + if err != nil { + close(pf.StopChan) + return fmt.Errorf("failed to create Elasticsearch client: %w", err) + } + + r.pf, r.client = pf, client + + return nil +} + +// disconnect drops the client and closes only a port-forward this type opened itself; the one the +// caller passed in is closed by the caller. +func (r *reconnectingHealthClient) disconnect() { + if r.pf != nil { + close(r.pf.StopChan) + r.pf = nil + } + + r.client = nil +} + +// expectedRestoredIndices returns the snapshot's indices that this restore recreates. The snapshot +// is re-read rather than passed in so that check-and-finalize works from just a snapshot name; +// snapshots are immutable, so the list cannot drift. +func expectedRestoredIndices(esClient es.Interface, appCtx *app.Context, repository, snapshotName string) ([]string, error) { + snapshot, err := esClient.GetSnapshot(repository, snapshotName) + if err != nil { + return nil, fmt.Errorf("failed to get snapshot details: %w", err) + } + + expected := filterSTSIndices( + snapshot.Indices, + appCtx.Config.Elasticsearch.Restore.IndexPrefix, + appCtx.Config.Elasticsearch.Restore.DatastreamIndexPrefix, + ) + if len(expected) == 0 { + return nil, fmt.Errorf("snapshot %s contains no indices matching the configured STS prefixes", snapshotName) + } + + return expected, nil +} + +// restoreProgress summarises how far a restore has got. Replicas are deliberately ignored: a +// cluster with fewer nodes than replicas keeps them unassigned forever, so requiring green would +// never complete there. +type restoreProgress struct { + // indicesRestored counts expected indices with every primary shard active, and decides completion. + indicesRestored int + // primariesActive counts active primary shards, and drives stall detection. Whole indices are + // too coarse for that: a large multi-shard index can restore for a long time without finishing. + primariesActive int } -func checkAndFinalize(esClient es.Interface, appCtx *app.Context, repository, snapshotName string, waitForComplete bool) error { +func measureRestore(expected []string, health map[string]es.IndexHealth) restoreProgress { + var progress restoreProgress + + for _, index := range expected { + indexHealth, exists := health[index] + if !exists { + continue + } + + progress.primariesActive += indexHealth.ActivePrimaryShards + if indexHealth.NumberOfShards > 0 && indexHealth.ActivePrimaryShards == indexHealth.NumberOfShards { + progress.indicesRestored++ + } + } + + return progress +} + +type indicesHealthGetter interface { + GetIndicesHealth() (map[string]es.IndexHealth, error) +} + +// newRestoreStatusFn builds the status callback used for both the single check and the wait loop. +// Completion is derived from the restored indices themselves: a restore that has been accepted but +// not yet applied is indistinguishable from a finished one when judged by recovery activity alone. +func newRestoreStatusFn( + esClient indicesHealthGetter, + log *logger.Logger, + expected []string, + noProgressTimeout time.Duration, + maxErrors int, +) func() (string, bool, error) { + lastPrimariesActive := -1 + lastProgressAt := time.Now() + errCount := 0 + + return func() (string, bool, error) { + health, err := esClient.GetIndicesHealth() + if err != nil { + errCount++ + if errCount >= maxErrors { + return "", false, err + } + + log.Warningf("Restore status check failed (%d/%d), retrying: %v", errCount, maxErrors, err) + + return es.StatusInProgress, false, nil + } + errCount = 0 + + progress := measureRestore(expected, health) + if progress.indicesRestored == len(expected) { + return es.StatusSuccess, true, nil + } + + if progress.primariesActive != lastPrimariesActive { + lastPrimariesActive = progress.primariesActive + lastProgressAt = time.Now() + } else if noProgressTimeout > 0 && time.Since(lastProgressAt) > noProgressTimeout { + // Deliberately not reported as a failed restore: all this establishes is that no progress + // was observed from here. Restarting a restore deletes every STS index first, so a caller + // that reads this as "failed" and retries would destroy a restore that is merely slow. + return "", false, fmt.Errorf( + "elasticsearch restore stalled: no primary shard finished restoring in %s "+ + "(%d of %d indices complete, %d primaries active); "+ + "it may still be running server-side, so check before restarting it", + noProgressTimeout, progress.indicesRestored, len(expected), progress.primariesActive, + ) + } + + log.Debugf( + "Restored %d of %d indices (%d primaries active)", + progress.indicesRestored, len(expected), progress.primariesActive, + ) + + return es.StatusInProgress, false, nil + } +} + +func checkAndFinalize( + esClient es.Interface, + appCtx *app.Context, + repository, snapshotName string, + waitForComplete bool, + noProgressTimeout time.Duration, +) error { + expected, err := expectedRestoredIndices(esClient, appCtx, repository, snapshotName) + if err != nil { + return err + } + + healthClient := &reconnectingHealthClient{appCtx: appCtx, client: esClient} + defer healthClient.disconnect() + + // Retrying only pays off while polling. A one-shot check that swallowed the error would report + // a dead tunnel as a running restore and exit 0. + maxErrors := restoreStatusMaxErrors + if !waitForComplete { + maxErrors = 1 + } + + statusFn := newRestoreStatusFn(healthClient, appCtx.Logger, expected, noProgressTimeout, maxErrors) + // Get restore status - appCtx.Logger.Infof("Checking restore status for snapshot: %s", snapshotName) - status, isComplete, err := esClient.GetRestoreStatus(repository, snapshotName) + appCtx.Logger.Infof("Checking restore status for snapshot: %s (%d indices)", snapshotName, len(expected)) + status, isComplete, err := statusFn() if err != nil { return fmt.Errorf("failed to get restore status: %w", err) } @@ -72,17 +305,9 @@ func checkAndFinalize(esClient es.Interface, appCtx *app.Context, repository, sn // Handle different scenarios if isComplete { switch status { - case "SUCCESS": + case es.StatusSuccess: appCtx.Logger.Successf("Restore completed successfully") return finalizeRestore(appCtx) - case "NOT_FOUND": - appCtx.Logger.Infof("No restore operation found for snapshot: %s", snapshotName) - appCtx.Logger.Infof("The restore may have already been finalized") - appCtx.Logger.Println() - appCtx.Logger.Infof("Checking if deployments need to be scaled up...") - return attemptScaleUp(appCtx) - case "FAILED": - return fmt.Errorf("restore failed with status: %s", status) default: return fmt.Errorf("restore completed with unexpected status: %s", status) } @@ -93,7 +318,7 @@ func checkAndFinalize(esClient es.Interface, appCtx *app.Context, repository, sn if waitForComplete { appCtx.Logger.Println() - return waitAndFinalize(esClient, appCtx, repository, snapshotName) + return waitAndFinalize(statusFn, appCtx, snapshotName) } // Not waiting - print status and exit @@ -103,15 +328,10 @@ func checkAndFinalize(esClient es.Interface, appCtx *app.Context, repository, sn } // waitAndFinalize waits for restore to complete and finalizes (scale up) -func waitAndFinalize(esClient es.Interface, appCtx *app.Context, repository, snapshotName string) error { +func waitAndFinalize(statusFn func() (string, bool, error), appCtx *app.Context, snapshotName string) error { restore.PrintAPIWaitingMessage("elasticsearch", snapshotName, appCtx.Namespace, appCtx.Logger) - // Wait for restore to complete - checkStatusFn := func() (string, bool, error) { - return esClient.GetRestoreStatus(repository, snapshotName) - } - - if err := restore.WaitForAPIRestore(checkStatusFn, 0, appCtx.Logger); err != nil { + if err := restore.WaitForAPIRestore(statusFn, 0, appCtx.Logger); err != nil { return err } @@ -129,20 +349,3 @@ func finalizeRestore(appCtx *app.Context) error { return restore.FinalizeRestore(scaleUpFn, appCtx.Logger) } - -// attemptScaleUp tries to scale up deployments and release lock (used when restore is not found/already complete) -func attemptScaleUp(appCtx *app.Context) error { - labelSelector := appCtx.Config.Elasticsearch.Restore.ScaleDownLabelSelector - scaleUpFn := func() error { - return scale.ScaleUpAndReleaseLock(appCtx.K8sClient, appCtx.Namespace, labelSelector, appCtx.Logger) - } - - if err := scaleUpFn(); err != nil { - // Don't fail if no deployments found to scale up - appCtx.Logger.Infof("No deployments found to scale up (this is normal if already finalized)") - return nil - } - - appCtx.Logger.Successf("Finalization completed successfully") - return nil -} diff --git a/cmd/elasticsearch/check_and_finalize_test.go b/cmd/elasticsearch/check_and_finalize_test.go new file mode 100644 index 0000000..b55ae03 --- /dev/null +++ b/cmd/elasticsearch/check_and_finalize_test.go @@ -0,0 +1,312 @@ +package elasticsearch + +import ( + "errors" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + es "github.com/stackvista/stackstate-backup-cli/internal/clients/elasticsearch" + "github.com/stackvista/stackstate-backup-cli/internal/foundation/logger" +) + +type healthResult struct { + health map[string]es.IndexHealth + err error +} + +// fakeHealthClient replays results in order, repeating the last one once exhausted. +type fakeHealthClient struct { + results []healthResult + calls int +} + +func (f *fakeHealthClient) GetIndicesHealth() (map[string]es.IndexHealth, error) { + result := f.results[min(f.calls, len(f.results)-1)] + f.calls++ + + return result.health, result.err +} + +func ok(health map[string]es.IndexHealth) healthResult { + return healthResult{health: health} +} + +func restored(shards int) es.IndexHealth { + return es.IndexHealth{Status: "green", NumberOfShards: shards, ActivePrimaryShards: shards} +} + +func TestMeasureRestore(t *testing.T) { + tests := []struct { + name string + expected []string + health map[string]es.IndexHealth + wantIndices int + wantPrimaries int + }{ + { + name: "index not recreated yet", + expected: []string{"sts_topology", "sts_events"}, + health: map[string]es.IndexHealth{"sts_topology": restored(1)}, + wantIndices: 1, + wantPrimaries: 1, + }, + { + name: "primaries still recovering are not counted", + expected: []string{"sts_topology"}, + health: map[string]es.IndexHealth{ + "sts_topology": {Status: "red", NumberOfShards: 3, ActivePrimaryShards: 2, InitializingShards: 1}, + }, + wantIndices: 0, + wantPrimaries: 2, + }, + { + name: "index present with no active primaries is not counted", + expected: []string{"sts_topology"}, + health: map[string]es.IndexHealth{ + "sts_topology": {Status: "red", NumberOfShards: 3, ActivePrimaryShards: 0, UnassignedShards: 3}, + }, + wantIndices: 0, + wantPrimaries: 0, + }, + { + name: "unassigned replicas do not block completion", + expected: []string{"sts_topology"}, + health: map[string]es.IndexHealth{ + "sts_topology": { + Status: "yellow", NumberOfShards: 1, NumberOfReplicas: 1, + ActivePrimaryShards: 1, ActiveShards: 1, UnassignedShards: 1, + }, + }, + wantIndices: 1, + wantPrimaries: 1, + }, + { + name: "all primaries active across several indices", + expected: []string{"sts_topology", ".ds-sts_k8s_logs-000001"}, + health: map[string]es.IndexHealth{ + "sts_topology": restored(3), + ".ds-sts_k8s_logs-000001": restored(2), + "unrelated": restored(1), + }, + wantIndices: 2, + wantPrimaries: 5, + }, + { + name: "index reporting zero shards is not counted", + expected: []string{"sts_topology"}, + health: map[string]es.IndexHealth{"sts_topology": {Status: "green"}}, + wantIndices: 0, + wantPrimaries: 0, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + progress := measureRestore(tt.expected, tt.health) + assert.Equal(t, tt.wantIndices, progress.indicesRestored, "indicesRestored") + assert.Equal(t, tt.wantPrimaries, progress.primariesActive, "primariesActive") + }) + } +} + +var statusFnExpected = []string{"sts_topology", "sts_events"} + +func TestNewRestoreStatusFn_Progress(t *testing.T) { + expected := statusFnExpected + log := logger.New(true, false) + + t.Run("accepted but not yet applied reports in progress", func(t *testing.T) { + client := &fakeHealthClient{results: []healthResult{ok(map[string]es.IndexHealth{})}} + + status, isComplete, err := newRestoreStatusFn(client, log, expected, time.Minute, 5)() + + require.NoError(t, err) + assert.Equal(t, es.StatusInProgress, status) + assert.False(t, isComplete) + }) + + t.Run("complete once every index has all primaries active", func(t *testing.T) { + client := &fakeHealthClient{results: []healthResult{ + ok(map[string]es.IndexHealth{"sts_topology": restored(1)}), + ok(map[string]es.IndexHealth{"sts_topology": restored(1), "sts_events": restored(1)}), + }} + statusFn := newRestoreStatusFn(client, log, expected, time.Minute, 5) + + status, isComplete, err := statusFn() + require.NoError(t, err) + assert.Equal(t, es.StatusInProgress, status) + assert.False(t, isComplete) + + status, isComplete, err = statusFn() + require.NoError(t, err) + assert.Equal(t, es.StatusSuccess, status) + assert.True(t, isComplete) + }) +} + +func TestNewRestoreStatusFn_Deadline(t *testing.T) { + expected := statusFnExpected + log := logger.New(true, false) + + t.Run("stalled restore errors instead of polling forever", func(t *testing.T) { + client := &fakeHealthClient{results: []healthResult{ + ok(map[string]es.IndexHealth{"sts_topology": restored(1)}), + }} + statusFn := newRestoreStatusFn(client, log, expected, time.Nanosecond, 5) + + // First call records the progress it observed; the second sees no change past the deadline. + _, isComplete, err := statusFn() + require.NoError(t, err) + assert.False(t, isComplete) + + _, isComplete, err = statusFn() + require.Error(t, err) + assert.False(t, isComplete) + // Must not read as a failed restore: the caller's retry deletes every STS index first. + assert.ErrorContains(t, err, "stalled") + assert.ErrorContains(t, err, "may still be running") + assert.NotContains(t, err.Error(), "restore failed") + }) + + t.Run("shard progress inside one index is not a stall", func(t *testing.T) { + // A large multi-shard index can restore for a long time without completing. Counting whole + // indices would read this as stalled; counting primaries does not. + client := &fakeHealthClient{results: []healthResult{ + ok(map[string]es.IndexHealth{"sts_topology": {NumberOfShards: 5, ActivePrimaryShards: 1}}), + ok(map[string]es.IndexHealth{"sts_topology": {NumberOfShards: 5, ActivePrimaryShards: 2}}), + ok(map[string]es.IndexHealth{"sts_topology": {NumberOfShards: 5, ActivePrimaryShards: 3}}), + }} + statusFn := newRestoreStatusFn(client, log, expected, time.Nanosecond, 5) + + for range 3 { + status, isComplete, err := statusFn() + require.NoError(t, err) + assert.Equal(t, es.StatusInProgress, status) + assert.False(t, isComplete) + } + }) + + t.Run("progress resets the deadline", func(t *testing.T) { + client := &fakeHealthClient{results: []healthResult{ + ok(map[string]es.IndexHealth{}), + ok(map[string]es.IndexHealth{"sts_topology": restored(1)}), + }} + statusFn := newRestoreStatusFn(client, log, expected, time.Nanosecond, 5) + + _, _, err := statusFn() + require.NoError(t, err) + + status, isComplete, err := statusFn() + require.NoError(t, err) + assert.Equal(t, es.StatusInProgress, status) + assert.False(t, isComplete) + }) + + t.Run("zero waits indefinitely", func(t *testing.T) { + client := &fakeHealthClient{results: []healthResult{ + ok(map[string]es.IndexHealth{"sts_topology": restored(1)}), + }} + statusFn := newRestoreStatusFn(client, log, expected, 0, 5) + + for range 3 { + status, isComplete, err := statusFn() + require.NoError(t, err) + assert.Equal(t, es.StatusInProgress, status) + assert.False(t, isComplete) + } + }) +} + +func TestNewRestoreStatusFn_TransientErrors(t *testing.T) { + expected := statusFnExpected + log := logger.New(true, false) + + t.Run("a dropped port-forward does not end the restore", func(t *testing.T) { + client := &fakeHealthClient{results: []healthResult{ + {err: errors.New("connection refused")}, + {err: errors.New("connection refused")}, + ok(map[string]es.IndexHealth{"sts_topology": restored(1), "sts_events": restored(1)}), + }} + statusFn := newRestoreStatusFn(client, log, expected, time.Minute, 5) + + for range 2 { + status, isComplete, err := statusFn() + require.NoError(t, err) + assert.Equal(t, es.StatusInProgress, status) + assert.False(t, isComplete) + } + + status, isComplete, err := statusFn() + require.NoError(t, err) + assert.Equal(t, es.StatusSuccess, status) + assert.True(t, isComplete) + }) + + t.Run("errors surface once they stop being transient", func(t *testing.T) { + client := &fakeHealthClient{results: []healthResult{{err: errors.New("connection refused")}}} + statusFn := newRestoreStatusFn(client, log, expected, time.Minute, 3) + + for range 2 { + _, _, err := statusFn() + require.NoError(t, err) + } + + _, _, err := statusFn() + require.Error(t, err) + assert.ErrorContains(t, err, "connection refused") + }) + + t.Run("a successful check clears earlier errors", func(t *testing.T) { + client := &fakeHealthClient{results: []healthResult{ + {err: errors.New("blip")}, + ok(map[string]es.IndexHealth{"sts_topology": restored(1)}), + {err: errors.New("blip")}, + {err: errors.New("blip")}, + }} + statusFn := newRestoreStatusFn(client, log, expected, time.Minute, 3) + + for range 4 { + _, _, err := statusFn() + require.NoError(t, err) + } + }) +} + +func TestReconnectingHealthClient(t *testing.T) { + t.Run("passes through the caller's client while it works", func(t *testing.T) { + health := map[string]es.IndexHealth{"sts_topology": restored(1)} + seeded := &fakeHealthClient{results: []healthResult{ok(health)}} + client := &reconnectingHealthClient{client: seeded} + + got, err := client.GetIndicesHealth() + + require.NoError(t, err) + assert.Equal(t, health, got) + assert.Equal(t, 1, seeded.calls) + }) + + t.Run("drops a broken client but never closes a port-forward it did not open", func(t *testing.T) { + // The caller closes the port-forward it passed in, so adopting it here would double-close a + // channel and panic. pf must stay nil until connect() opens one. + seeded := &fakeHealthClient{results: []healthResult{{err: errors.New("connection refused")}}} + client := &reconnectingHealthClient{client: seeded} + + _, err := client.GetIndicesHealth() + + require.Error(t, err) + assert.Nil(t, client.client, "a broken client must be dropped so the next call reconnects") + assert.Nil(t, client.pf, "must not adopt a port-forward it did not open") + }) + + t.Run("disconnect is safe to call repeatedly", func(t *testing.T) { + client := &reconnectingHealthClient{client: &fakeHealthClient{}} + + assert.NotPanics(t, func() { + client.disconnect() + client.disconnect() + }) + }) +} diff --git a/cmd/elasticsearch/restore.go b/cmd/elasticsearch/restore.go index 675f980..77868d4 100644 --- a/cmd/elasticsearch/restore.go +++ b/cmd/elasticsearch/restore.go @@ -26,11 +26,12 @@ const ( // Restore command flags var ( - snapshotName string - useLatest bool - runBackground bool - skipConfirmation bool - allowPartial bool + snapshotName string + useLatest bool + runBackground bool + skipConfirmation bool + allowPartial bool + restoreNoProgressIn time.Duration ) func restoreCmd(globalFlags *config.CLIGlobalFlags) *cobra.Command { @@ -47,6 +48,7 @@ func restoreCmd(globalFlags *config.CLIGlobalFlags) *cobra.Command { cmd.Flags().BoolVar(&runBackground, "background", false, "Run restore in background without waiting for completion") cmd.Flags().BoolVarP(&skipConfirmation, "yes", "y", false, "Skip confirmation prompt") cmd.Flags().BoolVar(&allowPartial, "allow-partial", false, "Allow restoring from a PARTIAL snapshot without extra confirmation") + cmd.Flags().DurationVar(&restoreNoProgressIn, "no-progress-timeout", defaultNoProgressTimeout, noProgressTimeoutUsage) cmd.MarkFlagsMutuallyExclusive("snapshot", "latest") cmd.MarkFlagsOneRequired("snapshot", "latest") return cmd @@ -145,7 +147,7 @@ func runRestore(appCtx *app.Context) error { return nil } - return checkAndFinalize(esClient, appCtx, repository, selectedSnapshot, !runBackground) + return checkAndFinalize(esClient, appCtx, repository, selectedSnapshot, !runBackground, restoreNoProgressIn) } // getLatestSnapshot retrieves the most recent snapshot from the repository diff --git a/cmd/elasticsearch/restore_test.go b/cmd/elasticsearch/restore_test.go index e723b43..e611100 100644 --- a/cmd/elasticsearch/restore_test.go +++ b/cmd/elasticsearch/restore_test.go @@ -94,8 +94,8 @@ func (m *mockESClientForRestore) ConfigureSLMPolicy(_, _, _, _, _, _ string, _, return fmt.Errorf("not implemented") } -func (m *mockESClientForRestore) GetRestoreStatus(_, _ string) (string, bool, error) { - return "SUCCESS", true, nil +func (m *mockESClientForRestore) GetIndicesHealth() (map[string]elasticsearch.IndexHealth, error) { + return nil, fmt.Errorf("not implemented") } // TestRestoreCmd_Unit tests the command structure diff --git a/internal/clients/elasticsearch/client.go b/internal/clients/elasticsearch/client.go index 1619ab8..03db72d 100644 --- a/internal/clients/elasticsearch/client.go +++ b/internal/clients/elasticsearch/client.go @@ -390,52 +390,43 @@ func (c *Client) RestoreSnapshot(repository, snapshotName, indicesPattern string return nil } -// RecoveryInfo represents the recovery status of a shard from _cat/recovery API -type RecoveryInfo struct { - Index string `json:"index"` - Shard string `json:"shard"` - Type string `json:"type"` - Stage string `json:"stage"` - Repository string `json:"repository"` - Snapshot string `json:"snapshot"` +// IndexHealth is the per-index section of the _cluster/health response +type IndexHealth struct { + Status string `json:"status"` + NumberOfShards int `json:"number_of_shards"` + NumberOfReplicas int `json:"number_of_replicas"` + ActivePrimaryShards int `json:"active_primary_shards"` + ActiveShards int `json:"active_shards"` + InitializingShards int `json:"initializing_shards"` + UnassignedShards int `json:"unassigned_shards"` } -// GetRestoreStatus checks the status of a restore operation by examining active shard recoveries. -// When a snapshot is being restored, shards are recovered with type "snapshot". -// Returns: (statusMessage, isComplete, error) -// Status can be: "IN_PROGRESS", "SUCCESS" -func (c *Client) GetRestoreStatus(repository, snapshotName string) (string, bool, error) { - // Use _cat/recovery API to check for active snapshot recoveries. - // This shows shards that are currently being recovered from a snapshot. - res, err := c.es.Cat.Recovery( - c.es.Cat.Recovery.WithContext(context.Background()), - c.es.Cat.Recovery.WithFormat("json"), - c.es.Cat.Recovery.WithActiveOnly(true), - c.es.Cat.Recovery.WithH("index,shard,type,stage,repository,snapshot"), +// GetIndicesHealth returns per-index cluster health keyed by index name. +// Indices that do not exist are simply absent from the result. +func (c *Client) GetIndicesHealth() (map[string]IndexHealth, error) { + res, err := c.es.Cluster.Health( + c.es.Cluster.Health.WithContext(context.Background()), + c.es.Cluster.Health.WithLevel("indices"), ) if err != nil { - return "", false, fmt.Errorf("failed to get recovery status: %w", err) + return nil, fmt.Errorf("failed to get cluster health: %w", err) } defer res.Body.Close() if res.IsError() { - return "", false, fmt.Errorf("elasticsearch returned error: %s", res.String()) + return nil, fmt.Errorf("elasticsearch returned error: %s", res.String()) } - var recoveries []RecoveryInfo - if err := json.NewDecoder(res.Body).Decode(&recoveries); err != nil { - return "", false, fmt.Errorf("failed to decode response: %w", err) + var health struct { + Indices map[string]IndexHealth `json:"indices"` + } + if err := json.NewDecoder(res.Body).Decode(&health); err != nil { + return nil, fmt.Errorf("failed to decode response: %w", err) } - // Check if any active recovery is from the specified snapshot - for _, recovery := range recoveries { - if recovery.Type == "snapshot" && - recovery.Repository == repository && - recovery.Snapshot == snapshotName { - return StatusInProgress, false, nil - } + if health.Indices == nil { + health.Indices = map[string]IndexHealth{} } - // No active recoveries from this snapshot - restore is complete - return StatusSuccess, true, nil + return health.Indices, nil } diff --git a/internal/clients/elasticsearch/client_test.go b/internal/clients/elasticsearch/client_test.go index f8e0e8e..74671a5 100644 --- a/internal/clients/elasticsearch/client_test.go +++ b/internal/clients/elasticsearch/client_test.go @@ -412,6 +412,96 @@ func TestClient_RestoreSnapshot(t *testing.T) { } } +func TestClient_GetIndicesHealth(t *testing.T) { + tests := []struct { + name string + responseBody string + expectedCount int + assertContents func(t *testing.T, health map[string]IndexHealth) + }{ + { + name: "indices are keyed by name", + responseBody: `{ + "cluster_name": "test", + "status": "green", + "indices": { + "sts_topology": { + "status": "green", + "number_of_shards": 3, + "number_of_replicas": 1, + "active_primary_shards": 3, + "active_shards": 6, + "initializing_shards": 0, + "unassigned_shards": 0 + }, + "sts_events": { + "status": "yellow", + "number_of_shards": 1, + "number_of_replicas": 1, + "active_primary_shards": 1, + "active_shards": 1, + "initializing_shards": 0, + "unassigned_shards": 1 + } + } + }`, + expectedCount: 2, + assertContents: func(t *testing.T, health map[string]IndexHealth) { + assert.Equal(t, 3, health["sts_topology"].NumberOfShards) + assert.Equal(t, 3, health["sts_topology"].ActivePrimaryShards) + assert.Equal(t, "yellow", health["sts_events"].Status) + assert.Equal(t, 1, health["sts_events"].UnassignedShards) + }, + }, + { + name: "no indices section yields an empty map", + responseBody: `{"cluster_name": "test", "status": "green"}`, + expectedCount: 0, + assertContents: func(t *testing.T, health map[string]IndexHealth) { + assert.NotNil(t, health) + }, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + server := mockESServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + assert.Equal(t, "/_cluster/health", r.URL.Path) + assert.Equal(t, "indices", r.URL.Query().Get("level")) + + w.WriteHeader(http.StatusOK) + _, _ = w.Write([]byte(tt.responseBody)) + })) + defer server.Close() + + client, err := NewClient(server.URL) + require.NoError(t, err) + + health, err := client.GetIndicesHealth() + + require.NoError(t, err) + assert.Len(t, health, tt.expectedCount) + tt.assertContents(t, health) + }) + } +} + +func TestClient_GetIndicesHealth_ElasticsearchError(t *testing.T) { + server := mockESServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusInternalServerError) + _, _ = w.Write([]byte(`{"error": "boom"}`)) + })) + defer server.Close() + + client, err := NewClient(server.URL) + require.NoError(t, err) + + health, err := client.GetIndicesHealth() + + require.Error(t, err) + assert.Nil(t, health) +} + func TestNewClient(t *testing.T) { client, err := NewClient("http://localhost:9200") require.NoError(t, err) diff --git a/internal/clients/elasticsearch/interface.go b/internal/clients/elasticsearch/interface.go index 5e1afd8..f8e0089 100644 --- a/internal/clients/elasticsearch/interface.go +++ b/internal/clients/elasticsearch/interface.go @@ -7,9 +7,9 @@ type Interface interface { ListSnapshots(repository string) ([]Snapshot, error) GetSnapshot(repository, snapshotName string) (*Snapshot, error) RestoreSnapshot(repository, snapshotName, indicesPattern string, partial bool) error - GetRestoreStatus(repository, snapshotName string) (string, bool, error) // Index operations + GetIndicesHealth() (map[string]IndexHealth, error) ListIndices(pattern string) ([]string, error) ListIndicesDetailed() ([]IndexInfo, error) DeleteIndex(index string) error diff --git a/internal/orchestration/restore/apirestore.go b/internal/orchestration/restore/apirestore.go index 1cd2d7c..5ffe37a 100644 --- a/internal/orchestration/restore/apirestore.go +++ b/internal/orchestration/restore/apirestore.go @@ -30,7 +30,8 @@ func WaitForAPIRestore( <-ticker.C statusMsg, isComplete, err := checkStatusFn() if err != nil { - return fmt.Errorf("failed to check restore status: %w", err) + // Not necessarily a failed check: a status function may also stop the wait deliberately. + return fmt.Errorf("stopped waiting for restore: %w", err) } log.Debugf("Restore status: %s (complete: %v)", statusMsg, isComplete) @@ -53,7 +54,7 @@ func PrintAPIWaitingMessage(serviceName, identifier, namespace string, log *logg log.Println() log.Infof("You can safely interrupt this command with Ctrl+C.") log.Infof("To check status and finalize later, run:") - log.Infof(" sts-backup %s check-and-finalize --operation-id %s -n %s", serviceName, identifier, namespace) + log.Infof(" sts-backup %s check-and-finalize --operation-id %s --wait -n %s", serviceName, identifier, namespace) } // PrintAPIRunningRestoreStatus prints status and instructions for a running restore From 239b4155f38c1d50b3df937a301979a8b2871145 Mon Sep 17 00:00:00 2001 From: Bram Schuur Date: Tue, 1 Sep 2026 11:00:17 +0200 Subject: [PATCH 9/9] STAC-25517: Fix success message on check-and-finalize --- internal/orchestration/restore/finalize.go | 25 +++++++++++++--------- 1 file changed, 15 insertions(+), 10 deletions(-) diff --git a/internal/orchestration/restore/finalize.go b/internal/orchestration/restore/finalize.go index 7c7cf16..802bc3e 100644 --- a/internal/orchestration/restore/finalize.go +++ b/internal/orchestration/restore/finalize.go @@ -39,6 +39,12 @@ type HandleCompletedJobParams struct { // HandleCompletedJob handles a job that's already complete // This includes printing status, fetching logs on failure, scaling up, and cleanup func HandleCompletedJob(params HandleCompletedJobParams) error { + defer func() { + // Cleanup resources + params.Log.Println() + _ = CleanupResources(params.K8sClient, params.Namespace, params.JobName, "", params.Log, params.CleanupPVC) + }() + params.Log.Println() if params.JobSucceeded { params.Log.Successf("Job completed successfully: %s", params.JobName) @@ -48,19 +54,18 @@ func HandleCompletedJob(params HandleCompletedJobParams) error { if err := params.ScaleUpFn(params.K8sClient, params.Namespace, params.ScaleSelector, params.Log); err != nil { params.Log.Warningf("Failed to scale up workload: %v", err) } - } else { - params.Log.Errorf("Job failed: %s", params.JobName) - params.Log.Println() - params.Log.Infof("Fetching logs...") - params.Log.Println() - if err := PrintJobLogs(params.K8sClient, params.Namespace, params.JobName, params.Log); err != nil { - params.Log.Warningf("Failed to fetch logs: %v", err) - } + + return nil } - // Cleanup resources + params.Log.Errorf("Job failed: %s", params.JobName) params.Log.Println() - return CleanupResources(params.K8sClient, params.Namespace, params.JobName, "", params.Log, params.CleanupPVC) + params.Log.Infof("Fetching logs...") + params.Log.Println() + if err := PrintJobLogs(params.K8sClient, params.Namespace, params.JobName, params.Log); err != nil { + params.Log.Warningf("Failed to fetch logs: %v", err) + } + return fmt.Errorf("job failed: %s", params.JobName) } // WaitAndFinalizeParams contains parameters for WaitAndFinalize