diff --git a/internal/api/authz.go b/internal/api/authz.go index 4907969..ced71a2 100644 --- a/internal/api/authz.go +++ b/internal/api/authz.go @@ -1,12 +1,14 @@ package api import ( + "fmt" "net/http" "os/exec" "strings" "github.com/flatrun/agent/internal/auth" "github.com/flatrun/agent/internal/contextkeys" + "github.com/flatrun/agent/pkg/models" "github.com/gin-gonic/gin" ) @@ -108,14 +110,14 @@ func (s *Server) requireContainerAccess(c *gin.Context, containerID, level strin } if actor.Role == auth.RoleAdmin { // Admins can see missing-container errors; non-admins below get a non-enumerating 403. - if _, err := containerDeploymentName(containerID); err != nil { + if _, err := s.containerDeploymentName(containerID); err != nil { c.JSON(http.StatusNotFound, gin.H{"error": "Container not found"}) return false } return true } - deploymentName, err := containerDeploymentName(containerID) + deploymentName, err := s.containerDeploymentName(containerID) if err != nil || deploymentName == "" { c.JSON(http.StatusForbidden, gin.H{"error": "No access to this container"}) return false @@ -135,7 +137,7 @@ func (s *Server) actorCanAccessContainer(c *gin.Context, containerID, level stri return true } - deploymentName, err := containerDeploymentName(containerID) + deploymentName, err := s.containerDeploymentName(containerID) if err != nil || deploymentName == "" { return false } @@ -143,17 +145,60 @@ func (s *Server) actorCanAccessContainer(c *gin.Context, containerID, level stri return actor.CanAccessDeployment(deploymentName, level) } -func containerDeploymentName(containerID string) (string, error) { - cmd := exec.Command("docker", "inspect", "--format", "{{ index .Config.Labels \""+composeProjectLabel+"\" }}", containerID) +func inspectContainerIdentity(containerID string) (string, string, string, error) { + format := "{{.Id}}\n{{.Name}}\n{{ index .Config.Labels \"" + composeProjectLabel + "\" }}" + cmd := exec.Command("docker", "inspect", "--format", format, containerID) output, err := cmd.Output() if err != nil { - return "", err + return "", "", "", err } - - deploymentName := strings.TrimSpace(string(output)) + parts := strings.SplitN(strings.TrimSpace(string(output)), "\n", 3) + if len(parts) != 3 { + return "", "", "", fmt.Errorf("unexpected container inspection result") + } + canonicalID := strings.TrimSpace(parts[0]) + containerName := strings.TrimPrefix(strings.TrimSpace(parts[1]), "/") + deploymentName := strings.TrimSpace(parts[2]) if deploymentName == "" { - return "", nil + deploymentName = "" + } + return canonicalID, containerName, deploymentName, nil +} + +func (s *Server) containerDeploymentName(containerID string) (string, error) { + canonicalID, _, label, inspectErr := inspectContainerIdentity(containerID) + if s.manager == nil { + return label, inspectErr + } + if label != "" { + if deployment, err := s.manager.GetDeployment(label); err == nil && deploymentContainsContainer(deployment, canonicalID) { + return deployment.Name, nil + } + } + deployments, err := s.manager.FindDeployments() + if err != nil { + return "", err + } + for _, candidate := range deployments { + deployment, getErr := s.manager.GetDeployment(candidate.Name) + if getErr == nil && deploymentContainsContainer(deployment, canonicalID) { + return deployment.Name, nil + } + } + if inspectErr != nil { + return "", inspectErr } + return "", fmt.Errorf("container does not belong to a deployment") +} - return deploymentName, nil +func deploymentContainsContainer(deployment *models.Deployment, containerID string) bool { + for _, service := range deployment.Services { + if len(service.ContainerID) < 12 || len(containerID) < 12 { + continue + } + if service.ContainerID == containerID || strings.HasPrefix(service.ContainerID, containerID) || strings.HasPrefix(containerID, service.ContainerID) { + return true + } + } + return false } diff --git a/internal/api/authz_test.go b/internal/api/authz_test.go index c7f6849..6afa0b8 100644 --- a/internal/api/authz_test.go +++ b/internal/api/authz_test.go @@ -39,6 +39,17 @@ func actorMiddleware(actor *auth.ActorContext) gin.HandlerFunc { } } +func TestDeploymentContainsContainerRejectsShortPrefixes(t *testing.T) { + deployment := &models.Deployment{Services: []models.Service{{ContainerID: "abcdef123456"}}} + + if deploymentContainsContainer(deployment, "abcdef") { + t.Fatal("short container ID matched a deployment container") + } + if !deploymentContainsContainer(deployment, "abcdef1234567890") { + t.Fatal("canonical container ID did not match its Docker short ID") + } +} + func TestClusterServiceCredentialsRejectUnscopedSensitiveResources(t *testing.T) { gin.SetMode(gin.TestMode) actor := &auth.ActorContext{ diff --git a/internal/api/backup_handlers.go b/internal/api/backup_handlers.go index b59fbe6..d0eb387 100644 --- a/internal/api/backup_handlers.go +++ b/internal/api/backup_handlers.go @@ -1,6 +1,7 @@ package api import ( + "context" "net/http" "strconv" @@ -10,6 +11,19 @@ import ( "github.com/gin-gonic/gin" ) +func (s *Server) retryBackupPublication(c *gin.Context) { + if s.backupManager == nil { + c.JSON(http.StatusServiceUnavailable, gin.H{"error": "Backup manager not enabled"}) + return + } + result, err := s.backupManager.RetryRemotePublication(context.Background(), c.Param("id")) + if err != nil { + c.JSON(http.StatusBadGateway, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"backup": result}) +} + func (s *Server) listBackups(c *gin.Context) { if s.backupManager == nil { c.JSON(http.StatusServiceUnavailable, gin.H{"error": "Backup manager not enabled"}) diff --git a/internal/api/backup_reliability_test.go b/internal/api/backup_reliability_test.go new file mode 100644 index 0000000..aea4bd4 --- /dev/null +++ b/internal/api/backup_reliability_test.go @@ -0,0 +1,130 @@ +package api + +import ( + "encoding/json" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/flatrun/agent/internal/backup" + "github.com/flatrun/agent/internal/docker" + "github.com/gin-gonic/gin" +) + +func runFailingBackupThroughHTTP(t *testing.T, metadata, dockerScript string, setup func(string)) (backup.Job, string) { + t.Helper() + root := t.TempDir() + deploymentDir := filepath.Join(root, "app") + if err := os.MkdirAll(deploymentDir, 0755); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(deploymentDir, "docker-compose.yml"), []byte("services:\n web:\n image: nginx\n"), 0644); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(deploymentDir, "service.yml"), []byte(metadata), 0644); err != nil { + t.Fatal(err) + } + if setup != nil { + setup(deploymentDir) + } + + binDir := t.TempDir() + logPath := filepath.Join(binDir, "docker.log") + script := "#!/bin/sh\nprintf '%s\\n' \"$*\" >> " + logPath + "\n" + dockerScript + "\n" + if err := os.WriteFile(filepath.Join(binDir, "docker"), []byte(script), 0755); err != nil { + t.Fatal(err) + } + t.Setenv("PATH", binDir+string(os.PathListSeparator)+os.Getenv("PATH")) + + backupManager, err := backup.NewManager(root) + if err != nil { + t.Fatal(err) + } + server := &Server{manager: docker.NewManager(root), backupManager: backupManager} + router := gin.New() + router.POST("/deployments/:name/backups", server.createDeploymentBackup) + router.GET("/deployments/:name/backups/jobs/:id", server.getBackupJob) + + created := httptest.NewRecorder() + router.ServeHTTP(created, httptest.NewRequest(http.MethodPost, "/deployments/app/backups", nil)) + if created.Code != http.StatusAccepted { + t.Fatalf("create status = %d, body = %s", created.Code, created.Body.String()) + } + var createResponse struct { + JobID string `json:"job_id"` + } + if err := json.Unmarshal(created.Body.Bytes(), &createResponse); err != nil { + t.Fatal(err) + } + + var job backup.Job + deadline := time.Now().Add(3 * time.Second) + for time.Now().Before(deadline) { + response := httptest.NewRecorder() + router.ServeHTTP(response, httptest.NewRequest(http.MethodGet, "/deployments/app/backups/jobs/"+createResponse.JobID, nil)) + var body struct { + Job backup.Job `json:"job"` + } + if err := json.Unmarshal(response.Body.Bytes(), &body); err != nil { + t.Fatal(err) + } + job = body.Job + if job.Status == backup.JobStatusFailed { + break + } + time.Sleep(10 * time.Millisecond) + } + if job.Status != backup.JobStatusFailed { + t.Fatalf("job status = %s, want failed", job.Status) + } + logBytes, _ := os.ReadFile(logPath) + return job, string(logBytes) +} + +func TestBackupRequiredDatabaseFailureThroughHTTP(t *testing.T) { + metadata := "name: app\nbackup:\n databases:\n - service: db\n type: unsupported\n post_hooks:\n - service: web\n command: resume\n" + job, dockerLog := runFailingBackupThroughHTTP(t, metadata, "exit 0", nil) + if len(job.ComponentResults) == 0 || job.ComponentResults[len(job.ComponentResults)-1].Status != backup.ResultStatusFailed { + t.Fatalf("component results = %#v", job.ComponentResults) + } + if !strings.Contains(dockerLog, "resume") { + t.Fatalf("cleanup invocation = %q", dockerLog) + } +} + +func TestBackupPreparationFailureStillRunsCleanupThroughHTTP(t *testing.T) { + metadata := "name: app\nbackup:\n databases:\n - service: db\n type: unsupported\n pre_hooks:\n - service: web\n command: prepare\n post_hooks:\n - service: web\n command: resume\n" + _, dockerLog := runFailingBackupThroughHTTP(t, metadata, "case \"$*\" in *prepare*) exit 1;; *) exit 0;; esac", nil) + if !strings.Contains(dockerLog, "prepare") || !strings.Contains(dockerLog, "resume") { + t.Fatalf("hook invocations = %q", dockerLog) + } +} + +func TestBackupRequiredFileFailureThroughHTTP(t *testing.T) { + metadata := "name: app\nbackup:\n databases:\n - service: db\n type: unsupported\n post_hooks:\n - service: web\n command: resume\n" + job, dockerLog := runFailingBackupThroughHTTP(t, metadata, "exit 0", func(deploymentDir string) { + dataDir := filepath.Join(deploymentDir, "data") + if err := os.MkdirAll(dataDir, 0755); err != nil { + t.Fatal(err) + } + if err := os.Symlink(filepath.Join(dataDir, "missing"), filepath.Join(dataDir, "broken")); err != nil { + t.Fatal(err) + } + }) + fileFailed := false + for _, result := range job.ComponentResults { + if result.Kind == "files" && result.Status == backup.ResultStatusFailed { + fileFailed = true + } + } + if !fileFailed { + t.Fatalf("component results = %#v", job.ComponentResults) + } + if !strings.Contains(dockerLog, "resume") { + t.Fatalf("cleanup invocation = %q", dockerLog) + } +} diff --git a/internal/api/compose_validation_integration_test.go b/internal/api/compose_validation_integration_test.go index 6a41122..f029e19 100644 --- a/internal/api/compose_validation_integration_test.go +++ b/internal/api/compose_validation_integration_test.go @@ -109,3 +109,38 @@ networks: t.Errorf("validateComposeContent with relative env_file in deployment dir = %v, want nil", err) } } + +func TestValidateNewComposeContent_UsesSuppliedRequiredVariable(t *testing.T) { + s := &Server{config: &config.Config{Infrastructure: config.InfrastructureConfig{DefaultProxyNetwork: "proxy"}}} + compose := `name: required-env +services: + app: + image: ${APP_IMAGE:?APP_IMAGE is required} +` + if err := s.validateNewComposeContent(compose, "required-env", []EnvVar{{Key: "APP_IMAGE", Value: "nginx:alpine"}}, t.TempDir()); err != nil { + t.Fatalf("validation with supplied required variable: %v", err) + } +} + +func TestValidateComposeContent_PrefersManagedEnvironment(t *testing.T) { + base := t.TempDir() + name := "managed-env" + deploymentDir := filepath.Join(base, name) + if err := os.MkdirAll(deploymentDir, 0755); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(deploymentDir, ".env"), []byte("OTHER=value\n"), 0600); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(deploymentDir, ".env.flatrun"), []byte("APP_IMAGE=nginx:alpine\n"), 0600); err != nil { + t.Fatal(err) + } + s := &Server{ + config: &config.Config{Infrastructure: config.InfrastructureConfig{DefaultProxyNetwork: "proxy"}}, + manager: docker.NewManager(base), + } + compose := "services:\n app:\n image: ${APP_IMAGE:?APP_IMAGE is required}\n" + if err := s.validateComposeContent(compose, name); err != nil { + t.Fatalf("validation with managed environment: %v", err) + } +} diff --git a/internal/api/container_exec.go b/internal/api/container_exec.go index 83323e8..844c225 100644 --- a/internal/api/container_exec.go +++ b/internal/api/container_exec.go @@ -88,7 +88,7 @@ func (s *Server) containerExec(c *gin.Context) { sendError(conn, "No access to this container") return } - if deploymentName, err := containerDeploymentName(containerID); err == nil && deploymentName != "" { + if deploymentName, err := s.containerDeploymentName(containerID); err == nil && deploymentName != "" { if blocked, reason, err := s.protectedDeploymentActionBlocked(deploymentName, protectedActionTerminal); err != nil { sendError(conn, "Failed to check protected mode: "+err.Error()) return @@ -324,7 +324,7 @@ func (s *Server) containerExecHTTP(c *gin.Context) { } commandLine := strings.Join(append([]string{req.Command}, req.Args...), " ") - if deploymentName, err := containerDeploymentName(containerID); err == nil && deploymentName != "" { + if deploymentName, err := s.containerDeploymentName(containerID); err == nil && deploymentName != "" { if blocked, reason, err := s.protectedDeploymentActionBlocked(deploymentName, protectedActionExec); err != nil { c.JSON(http.StatusInternalServerError, gin.H{"error": "Failed to check protected mode: " + err.Error()}) return diff --git a/internal/api/container_identity_integration_test.go b/internal/api/container_identity_integration_test.go new file mode 100644 index 0000000..c875d07 --- /dev/null +++ b/internal/api/container_identity_integration_test.go @@ -0,0 +1,67 @@ +package api + +import ( + "fmt" + "os" + "os/exec" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/flatrun/agent/internal/docker" +) + +func TestContainerDeploymentNameUsesOwningDirectory(t *testing.T) { + if _, err := exec.LookPath("docker"); err != nil { + t.Skip("docker is unavailable") + } + image := os.Getenv("FLATRUN_DOCKER_TEST_IMAGE") + if image == "" { + image = "busybox:latest" + } + if err := exec.Command("docker", "image", "inspect", image).Run(); err != nil { + t.Skipf("docker test image %s is unavailable", image) + } + + suffix := fmt.Sprintf("%d", time.Now().UnixNano()) + deploymentName := "actual-" + suffix + projectName := "label-" + suffix + containerName := "container-" + suffix + root := t.TempDir() + deploymentDir := filepath.Join(root, deploymentName) + if err := os.MkdirAll(deploymentDir, 0755); err != nil { + t.Fatal(err) + } + compose := "name: " + projectName + "\nservices:\n web:\n image: " + image + "\n" + if err := os.WriteFile(filepath.Join(deploymentDir, "docker-compose.yml"), []byte(compose), 0644); err != nil { + t.Fatal(err) + } + + output, err := exec.Command("docker", "create", + "--name", containerName, + "--label", composeProjectLabel+"="+projectName, + "--label", "com.docker.compose.service=web", + image, "sleep", "60").CombinedOutput() + if err != nil { + t.Fatalf("create container: %v: %s", err, output) + } + containerID := strings.TrimSpace(string(output)) + t.Cleanup(func() { _ = exec.Command("docker", "rm", "-f", containerID).Run() }) + if output, err := exec.Command("docker", "start", containerID).CombinedOutput(); err != nil { + t.Fatalf("start container: %v: %s", err, output) + } + + server := &Server{manager: docker.NewManager(root)} + resolved, err := server.containerDeploymentName(containerID) + if err != nil { + t.Fatal(err) + } + if resolved != deploymentName { + t.Fatalf("deployment = %q, want %q", resolved, deploymentName) + } + resolved, err = server.containerDeploymentName(containerName) + if err != nil || resolved != deploymentName { + t.Fatalf("deployment by name = %q, error = %v", resolved, err) + } +} diff --git a/internal/api/openapi.json b/internal/api/openapi.json index d210ad1..a50d502 100644 --- a/internal/api/openapi.json +++ b/internal/api/openapi.json @@ -4543,6 +4543,38 @@ "x-permission": "backups:write" } }, + "/api/deployments/{name}/backups/{id}/retry-publication": { + "post": { + "operationId": "post-deployments-by-name-backups-by-id-retry-publication", + "parameters": [ + { + "in": "path", + "name": "name", + "required": true, + "schema": { + "type": "string" + } + }, + { + "in": "path", + "name": "id", + "required": true, + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "Success" + } + }, + "tags": [ + "deployments" + ], + "x-permission": "backups:write" + } + }, "/api/deployments/{name}/certificates/renew": { "post": { "operationId": "post-deployments-by-name-certificates-renew", @@ -12788,10 +12820,22 @@ "backup.Backup": { "type": "object", "properties": { + "cleanup_results": { + "type": "array", + "items": { + "$ref": "#/components/schemas/backup.ComponentResult" + } + }, "completed_at": { "type": "string", "format": "date-time" }, + "component_results": { + "type": "array", + "items": { + "$ref": "#/components/schemas/backup.ComponentResult" + } + }, "components": { "type": "array", "items": { @@ -12805,6 +12849,12 @@ "deployment_name": { "type": "string" }, + "destination_results": { + "type": "array", + "items": { + "$ref": "#/components/schemas/backup.DestinationResult" + } + }, "error": { "type": "string" }, @@ -12842,7 +12892,10 @@ "created_at", "completed_at", "expires_at", - "locations" + "locations", + "component_results", + "cleanup_results", + "destination_results" ], "x-columns": [ "id", @@ -12855,6 +12908,40 @@ "expires_at" ] }, + "backup.ComponentResult": { + "type": "object", + "properties": { + "error": { + "type": "string" + }, + "kind": { + "type": "string" + }, + "name": { + "type": "string" + }, + "required": { + "type": "boolean" + }, + "status": { + "type": "string" + } + }, + "x-property-order": [ + "name", + "kind", + "required", + "status", + "error" + ], + "x-columns": [ + "name", + "kind", + "required", + "status", + "error" + ] + }, "backup.CreateBackupRequest": { "type": "object", "properties": { @@ -12877,6 +12964,30 @@ "deployment_name" ] }, + "backup.DestinationResult": { + "type": "object", + "properties": { + "error": { + "type": "string" + }, + "name": { + "type": "string" + }, + "status": { + "type": "string" + } + }, + "x-property-order": [ + "name", + "status", + "error" + ], + "x-columns": [ + "name", + "status", + "error" + ] + }, "backup.RestoreBackupRequest": { "type": "object", "properties": { diff --git a/internal/api/protected_mode.go b/internal/api/protected_mode.go index 5d67cd5..0bb87d1 100644 --- a/internal/api/protected_mode.go +++ b/internal/api/protected_mode.go @@ -225,7 +225,7 @@ func protectedCommandRuleMatchesCommand(rule models.ProtectedCommandRule, comman } func (s *Server) protectedContainerCommandBlocked(containerID, command string) (bool, *models.ProtectedCommandRule, error) { - deploymentName, err := containerDeploymentName(containerID) + deploymentName, err := s.containerDeploymentName(containerID) if err != nil || deploymentName == "" { return false, nil, err } diff --git a/internal/api/server.go b/internal/api/server.go index 3c43498..aeaa369 100644 --- a/internal/api/server.go +++ b/internal/api/server.go @@ -25,6 +25,7 @@ import ( "sync" "time" + "github.com/compose-spec/compose-go/v2/dotenv" "github.com/compose-spec/compose-go/v2/loader" composetypes "github.com/compose-spec/compose-go/v2/types" "github.com/flatrun/agent/internal/access" @@ -831,6 +832,7 @@ func (s *Server) setupRoutes() { protected.DELETE("/deployments/:name/backups/:id", s.authMiddleware.RequirePermission(auth.PermBackupsDelete), s.authMiddleware.RequireDeploymentAccess(auth.AccessLevelAdmin), s.requireBackupDeployment, s.deleteBackup) protected.GET("/deployments/:name/backups/:id/download", s.authMiddleware.RequirePermission(auth.PermBackupsRead), s.authMiddleware.RequireDeploymentAccess(auth.AccessLevelRead), s.requireBackupDeployment, s.downloadBackup) protected.POST("/deployments/:name/backups/:id/restore", s.authMiddleware.RequirePermission(auth.PermBackupsWrite), s.authMiddleware.RequireDeploymentAccess(auth.AccessLevelWrite), s.requireBackupDeployment, s.restoreBackup) + protected.POST("/deployments/:name/backups/:id/retry-publication", s.authMiddleware.RequirePermission(auth.PermBackupsWrite), s.authMiddleware.RequireDeploymentAccess(auth.AccessLevelWrite), s.requireBackupDeployment, s.retryBackupPublication) protected.GET("/deployments/:name/backups/jobs/:id", s.authMiddleware.RequirePermission(auth.PermBackupsRead), s.authMiddleware.RequireDeploymentAccess(auth.AccessLevelRead), s.requireBackupJobDeployment, s.getBackupJob) protected.GET("/deployments/:name/backup-config", s.authMiddleware.RequirePermission(auth.PermBackupsRead), s.authMiddleware.RequireDeploymentAccess(auth.AccessLevelRead), s.getDeploymentBackupConfig) protected.PUT("/deployments/:name/backup-config", s.authMiddleware.RequirePermission(auth.PermBackupsWrite), s.authMiddleware.RequireDeploymentAccess(auth.AccessLevelWrite), s.updateDeploymentBackupConfig) @@ -5007,13 +5009,33 @@ func (s *Server) composeValidationDir(name string) string { } func validateComposeWithComposeGo(content, workingDir string) error { + environment := make(map[string]string) + dotEnvPath := filepath.Join(workingDir, ".env.flatrun") + if _, err := os.Stat(dotEnvPath); os.IsNotExist(err) { + dotEnvPath = filepath.Join(workingDir, ".env") + } + if _, err := os.Stat(dotEnvPath); err == nil { + values, loadErr := dotenv.GetEnvFromFile(environment, []string{dotEnvPath}) + if loadErr != nil { + return fmt.Errorf("invalid compose environment: %w", loadErr) + } + for key, value := range values { + environment[key] = value + } + } + for _, entry := range os.Environ() { + key, value, ok := strings.Cut(entry, "=") + if ok { + environment[key] = value + } + } configDetails := composetypes.ConfigDetails{ ConfigFiles: []composetypes.ConfigFile{{ Filename: "compose.yml", Content: []byte(content), }}, WorkingDir: workingDir, - Environment: map[string]string{}, + Environment: environment, } _, err := loader.LoadWithContext(context.Background(), configDetails, func(o *loader.Options) { o.SetProjectName("flatrun-validation", true) diff --git a/internal/backup/jobs.go b/internal/backup/jobs.go index 840d5eb..3013b49 100644 --- a/internal/backup/jobs.go +++ b/internal/backup/jobs.go @@ -16,19 +16,24 @@ const ( JobStatusPending JobStatus = "pending" JobStatusRunning JobStatus = "running" JobStatusCompleted JobStatus = "completed" + JobStatusPartial JobStatus = "partial" + JobStatusLocalOnly JobStatus = "local_only" JobStatusFailed JobStatus = "failed" ) type Job struct { - ID string `json:"id"` - Type JobType `json:"type"` - Status JobStatus `json:"status"` - DeploymentName string `json:"deployment_name"` - BackupID string `json:"backup_id,omitempty"` - Progress string `json:"progress,omitempty"` - Error string `json:"error,omitempty"` - StartedAt time.Time `json:"started_at"` - CompletedAt *time.Time `json:"completed_at,omitempty"` + ID string `json:"id"` + Type JobType `json:"type"` + Status JobStatus `json:"status"` + DeploymentName string `json:"deployment_name"` + BackupID string `json:"backup_id,omitempty"` + Progress string `json:"progress,omitempty"` + Error string `json:"error,omitempty"` + StartedAt time.Time `json:"started_at"` + CompletedAt *time.Time `json:"completed_at,omitempty"` + ComponentResults []ComponentResult `json:"component_results,omitempty"` + CleanupResults []ComponentResult `json:"cleanup_results,omitempty"` + DestinationResults []DestinationResult `json:"destination_results,omitempty"` } type JobTracker struct { @@ -70,13 +75,24 @@ func (t *JobTracker) UpdateStatus(id string, status JobStatus, progress string) if job, ok := t.jobs[id]; ok { job.Status = status job.Progress = progress - if status == JobStatusCompleted || status == JobStatusFailed { + if status == JobStatusCompleted || status == JobStatusPartial || status == JobStatusLocalOnly || status == JobStatusFailed { now := time.Now() job.CompletedAt = &now } } } +func (t *JobTracker) SetBackup(id string, backup *Backup) { + t.mu.Lock() + defer t.mu.Unlock() + if job, ok := t.jobs[id]; ok { + job.BackupID = backup.ID + job.ComponentResults = backup.ComponentResults + job.CleanupResults = backup.CleanupResults + job.DestinationResults = backup.DestinationResults + } +} + func (t *JobTracker) SetError(id string, err error) { t.mu.Lock() defer t.mu.Unlock() @@ -136,13 +152,22 @@ func (m *Manager) StartBackupJob(deploymentName string, spec *BackupSpec) string m.jobs.UpdateStatus(jobID, JobStatusRunning, "Starting backup") backup, err := m.CreateBackup(context.Background(), deploymentName, spec) + if backup != nil { + m.jobs.SetBackup(jobID, backup) + } if err != nil { m.jobs.SetError(jobID, err) return } - m.jobs.SetBackupID(jobID, backup.ID) - m.jobs.UpdateStatus(jobID, JobStatusCompleted, "Backup completed") + switch backup.Status { + case BackupStatusPartial: + m.jobs.UpdateStatus(jobID, JobStatusPartial, "Backup captured with incomplete protection") + case BackupStatusLocalOnly: + m.jobs.UpdateStatus(jobID, JobStatusLocalOnly, "Backup retained locally") + default: + m.jobs.UpdateStatus(jobID, JobStatusCompleted, "Backup completed") + } }() return jobID diff --git a/internal/backup/manager.go b/internal/backup/manager.go index 4d14c29..393fdb2 100644 --- a/internal/backup/manager.go +++ b/internal/backup/manager.go @@ -5,6 +5,7 @@ import ( "compress/gzip" "context" "encoding/json" + "errors" "fmt" "io" "log" @@ -47,7 +48,7 @@ func (m *Manager) CreateBackup(ctx context.Context, deploymentName string, spec return nil, fmt.Errorf("deployment not found: %s", deploymentName) } - backupID := fmt.Sprintf("%s_%s", deploymentName, time.Now().Format("20060102_150405")) + backupID := fmt.Sprintf("%s_%s", deploymentName, time.Now().Format("20060102_150405.000000000")) backupDir := filepath.Join(m.backupsPath, deploymentName) if err := os.MkdirAll(backupDir, 0755); err != nil { return nil, fmt.Errorf("failed to create backup directory: %w", err) @@ -60,6 +61,14 @@ func (m *Manager) CreateBackup(ctx context.Context, deploymentName string, spec CreatedAt: time.Now(), Components: []string{}, } + record := func(kind, name string, required bool, err error) { + result := ComponentResult{Name: name, Kind: kind, Required: required, Status: ResultStatusCompleted} + if err != nil { + result.Status = ResultStatusFailed + result.Error = err.Error() + } + backup.ComponentResults = append(backup.ComponentResults, result) + } tempDir, err := os.MkdirTemp("", "flatrun-backup-*") if err != nil { @@ -76,28 +85,41 @@ func (m *Manager) CreateBackup(ctx context.Context, deploymentName string, spec Components: BackupComponents{}, } - if spec != nil { - if err := m.executeHooks(ctx, deploymentName, spec.PreHooks); err != nil { - log.Printf("Backup: pre-hook warning: %v", err) + var captureErrors []error + if spec != nil && len(spec.PreHooks) > 0 { + err := m.executeHooks(ctx, deploymentName, spec.PreHooks) + record("preparation", "pre_hooks", true, err) + if err != nil { + captureErrors = append(captureErrors, err) } } if err := m.backupComposeFile(deploymentPath, tempDir, &metadata); err != nil { - log.Printf("Backup: compose file warning: %v", err) + record("configuration", "compose", true, err) + captureErrors = append(captureErrors, err) } else { + record("configuration", "compose", true, nil) backup.Components = append(backup.Components, "compose") } if err := m.backupEnvFile(deploymentPath, tempDir, &metadata); err != nil { - log.Printf("Backup: env file warning: %v", err) + record("configuration", "environment", true, err) + captureErrors = append(captureErrors, err) } else { - backup.Components = append(backup.Components, "env") + record("configuration", "environment", true, nil) + if metadata.Components.EnvFile { + backup.Components = append(backup.Components, "env") + } } if err := m.backupMetadataFile(deploymentPath, tempDir, &metadata); err != nil { - log.Printf("Backup: metadata file warning: %v", err) + record("configuration", "metadata", true, err) + captureErrors = append(captureErrors, err) } else { - backup.Components = append(backup.Components, "metadata") + record("configuration", "metadata", true, nil) + if metadata.Components.Metadata { + backup.Components = append(backup.Components, "metadata") + } } var excludes []string @@ -105,15 +127,20 @@ func (m *Manager) CreateBackup(ctx context.Context, deploymentName string, spec excludes = spec.ExcludePatterns } if err := m.backupMountedData(deploymentPath, tempDir, &metadata, excludes); err != nil { - log.Printf("Backup: mounted data warning: %v", err) + record("files", "mounted_data", true, err) + captureErrors = append(captureErrors, err) + } else { + record("files", "mounted_data", true, nil) } if len(metadata.Components.MountedData) > 0 { backup.Components = append(backup.Components, "mounted_data") } if spec != nil && len(spec.ContainerPaths) > 0 { - if err := m.backupContainerData(ctx, deploymentName, spec.ContainerPaths, tempDir, &metadata); err != nil { - log.Printf("Backup: container data warning: %v", err) + results, err := m.backupContainerData(ctx, deploymentName, spec.ContainerPaths, tempDir, &metadata) + backup.ComponentResults = append(backup.ComponentResults, results...) + if err != nil { + captureErrors = append(captureErrors, err) } if len(metadata.Components.ContainerData) > 0 { backup.Components = append(backup.Components, "container_data") @@ -121,46 +148,76 @@ func (m *Manager) CreateBackup(ctx context.Context, deploymentName string, spec } if spec != nil && len(spec.Databases) > 0 { - if err := m.backupDatabases(ctx, deploymentName, spec.Databases, tempDir, &metadata); err != nil { - log.Printf("Backup: database warning: %v", err) + results, err := m.backupDatabases(ctx, deploymentName, spec.Databases, tempDir, &metadata) + backup.ComponentResults = append(backup.ComponentResults, results...) + if err != nil { + captureErrors = append(captureErrors, err) } if len(metadata.Components.Databases) > 0 { backup.Components = append(backup.Components, "databases") } } - metadataJSON, _ := json.MarshalIndent(metadata, "", " ") - if err := os.WriteFile(filepath.Join(tempDir, "backup.json"), metadataJSON, 0644); err != nil { - return nil, fmt.Errorf("failed to write backup metadata: %w", err) - } - archivePath := filepath.Join(backupDir, backupID+".tar.gz") - if err := m.createArchive(tempDir, archivePath); err != nil { - backup.Status = BackupStatusFailed - backup.Error = err.Error() - return backup, fmt.Errorf("failed to create backup archive: %w", err) - } - - if spec != nil { - if err := m.executeHooks(ctx, deploymentName, spec.PostHooks); err != nil { - log.Printf("Backup: post-hook warning: %v", err) + metadata.ComponentResults = backup.ComponentResults + if len(captureErrors) == 0 { + metadataJSON, err := json.MarshalIndent(metadata, "", " ") + if err != nil { + captureErrors = append(captureErrors, err) + } else if err := os.WriteFile(filepath.Join(tempDir, "backup.json"), metadataJSON, 0600); err != nil { + captureErrors = append(captureErrors, fmt.Errorf("failed to write backup metadata: %w", err)) + } else if err := m.createArchive(tempDir, archivePath); err != nil { + captureErrors = append(captureErrors, fmt.Errorf("failed to create backup archive: %w", err)) + } else { + backup.Path = archivePath + backup.Locations = []string{locationLocal} + if info, statErr := os.Stat(archivePath); statErr == nil { + backup.Size = info.Size() + } } } - info, _ := os.Stat(archivePath) - if info != nil { - backup.Size = info.Size() + var cleanupErr error + if spec != nil && len(spec.PostHooks) > 0 { + cleanupErr = m.executeHooks(ctx, deploymentName, spec.PostHooks) + result := ComponentResult{Name: "post_hooks", Kind: "cleanup", Required: true, Status: ResultStatusCompleted} + if cleanupErr != nil { + result.Status = ResultStatusFailed + result.Error = cleanupErr.Error() + } + backup.CleanupResults = append(backup.CleanupResults, result) } - backup.Path = archivePath - backup.Status = BackupStatusCompleted now := time.Now() backup.CompletedAt = &now + if len(captureErrors) > 0 { + backup.Status = BackupStatusFailed + backup.Error = errors.Join(captureErrors...).Error() + _ = m.saveBackupRecord(backup) + return backup, errors.Join(captureErrors...) + } - backup.Locations = []string{locationLocal} - backup.Locations = append(backup.Locations, m.mirrorToRemotes(ctx, deploymentName, backupID, archivePath, backup.Size)...) + backup.DestinationResults = m.mirrorToRemotes(ctx, deploymentName, backupID, archivePath, backup.Size) + succeeded := 0 + for _, result := range backup.DestinationResults { + if result.Status == ResultStatusCompleted { + succeeded++ + backup.Locations = append(backup.Locations, result.Name) + } + } + switch { + case len(backup.DestinationResults) > 0 && succeeded == 0: + backup.Status = BackupStatusLocalOnly + case cleanupErr != nil || succeeded < len(backup.DestinationResults): + backup.Status = BackupStatusPartial + default: + backup.Status = BackupStatusCompleted + } + if err := m.saveBackupRecord(backup); err != nil { + return backup, fmt.Errorf("failed to persist backup result: %w", err) + } - log.Printf("Backup completed: %s (%d bytes)", backupID, backup.Size) + log.Printf("Backup finished with status %s: %s (%d bytes)", backup.Status, backupID, backup.Size) return backup, nil } @@ -194,9 +251,10 @@ func (m *Manager) backupEnvFile(deploymentPath, tempDir string, metadata *Backup envPath := filepath.Join(deploymentPath, envFile) if _, err := os.Stat(envPath); err == nil { destPath := filepath.Join(envDir, envFile) - if err := copyFile(envPath, destPath); err == nil { - found = true + if err := copyFile(envPath, destPath); err != nil { + return fmt.Errorf("failed to copy %s: %w", envFile, err) } + found = true } } @@ -237,8 +295,7 @@ func (m *Manager) backupMountedData(deploymentPath, tempDir string, metadata *Ba if info, err := os.Stat(srcPath); err == nil && info.IsDir() { destPath := filepath.Join(dataDir, dir) if err := copyDir(srcPath, destPath); err != nil { - log.Printf("Backup: failed to copy %s: %v", dir, err) - continue + return fmt.Errorf("failed to copy %s: %w", dir, err) } metadata.Components.MountedData = append(metadata.Components.MountedData, dir) } @@ -261,13 +318,17 @@ func matchesExclude(name string, patterns []string) bool { return false } -func (m *Manager) backupContainerData(ctx context.Context, deploymentName string, paths []ContainerPath, tempDir string, metadata *BackupMetadata) error { +func (m *Manager) backupContainerData(ctx context.Context, deploymentName string, paths []ContainerPath, tempDir string, metadata *BackupMetadata) ([]ComponentResult, error) { containerDir := filepath.Join(tempDir, "container_data") if err := os.MkdirAll(containerDir, 0755); err != nil { - return fmt.Errorf("failed to create container data backup directory: %w", err) + return nil, fmt.Errorf("failed to create container data backup directory: %w", err) } + var results []ComponentResult + var requiredErrors []error for _, path := range paths { + name := fmt.Sprintf("%s:%s", path.Service, path.ContainerPath) + result := ComponentResult{Name: name, Kind: "files", Required: path.Required, Status: ResultStatusCompleted} containerName := fmt.Sprintf("%s-%s", deploymentName, path.Service) if path.Service == deploymentName || path.Service == "" { containerName = deploymentName @@ -275,31 +336,50 @@ func (m *Manager) backupContainerData(ctx context.Context, deploymentName string destPath := filepath.Join(containerDir, path.Service, filepath.Base(path.ContainerPath)) if err := os.MkdirAll(filepath.Dir(destPath), 0755); err != nil { - return fmt.Errorf("failed to create directory for %s: %w", path.ContainerPath, err) + result.Status = ResultStatusFailed + result.Error = fmt.Sprintf("failed to prepare destination: %v", err) + results = append(results, result) + if path.Required { + requiredErrors = append(requiredErrors, fmt.Errorf("failed to create directory for %s: %w", path.ContainerPath, err)) + } + continue } cmd := exec.CommandContext(ctx, "docker", "cp", fmt.Sprintf("%s:%s", containerName, path.ContainerPath), destPath) if err := cmd.Run(); err != nil { + result.Status = ResultStatusFailed + result.Error = "copy failed" + results = append(results, result) if path.Required { - return fmt.Errorf("failed to copy %s from container %s: %w", path.ContainerPath, containerName, err) + requiredErrors = append(requiredErrors, fmt.Errorf("failed to copy %s from container %s: %w", path.ContainerPath, containerName, err)) + } + if !path.Required { + log.Printf("Backup: optional container path %s not available: %v", path.ContainerPath, err) } - log.Printf("Backup: optional container path %s not available: %v", path.ContainerPath, err) continue } - metadata.Components.ContainerData = append(metadata.Components.ContainerData, fmt.Sprintf("%s:%s", path.Service, path.ContainerPath)) + results = append(results, result) + metadata.Components.ContainerData = append(metadata.Components.ContainerData, name) } - return nil + return results, errors.Join(requiredErrors...) } -func (m *Manager) backupDatabases(ctx context.Context, deploymentName string, databases []DatabaseSpec, tempDir string, metadata *BackupMetadata) error { +func (m *Manager) backupDatabases(ctx context.Context, deploymentName string, databases []DatabaseSpec, tempDir string, metadata *BackupMetadata) ([]ComponentResult, error) { dbDir := filepath.Join(tempDir, "databases") if err := os.MkdirAll(dbDir, 0755); err != nil { - return fmt.Errorf("failed to create databases backup directory: %w", err) + return nil, fmt.Errorf("failed to create databases backup directory: %w", err) } + var results []ComponentResult + var databaseErrors []error for _, db := range databases { + name := db.Service + if name == "" { + name = deploymentName + } + result := ComponentResult{Name: name, Kind: "database", Required: true, Status: ResultStatusCompleted} var dumpPath string var err error @@ -309,19 +389,46 @@ func (m *Manager) backupDatabases(ctx context.Context, deploymentName string, da case "postgresql", "postgres": dumpPath, err = m.dumpPostgres(ctx, deploymentName, &db, dbDir) default: - log.Printf("Backup: unsupported database type: %s", db.Type) - continue + err = fmt.Errorf("unsupported database type: %s", db.Type) } if err != nil { - log.Printf("Backup: failed to dump database %s: %v", db.Service, err) + result.Status = ResultStatusFailed + result.Error = "dump failed" + results = append(results, result) + databaseErrors = append(databaseErrors, fmt.Errorf("failed to dump database %s: %w", name, err)) continue } + results = append(results, result) metadata.Components.Databases = append(metadata.Components.Databases, filepath.Base(dumpPath)) } - return nil + return results, errors.Join(databaseErrors...) +} + +func (m *Manager) backupRecordPath(deploymentName, backupID string) string { + return filepath.Join(m.backupsPath, deploymentName, backupID+".json") +} + +func (m *Manager) saveBackupRecord(backup *Backup) error { + data, err := json.MarshalIndent(backup, "", " ") + if err != nil { + return err + } + return os.WriteFile(m.backupRecordPath(backup.DeploymentName, backup.ID), data, 0600) +} + +func (m *Manager) readBackupRecord(deploymentName, backupID string) (*Backup, error) { + data, err := os.ReadFile(m.backupRecordPath(deploymentName, backupID)) + if err != nil { + return nil, err + } + var backup Backup + if err := json.Unmarshal(data, &backup); err != nil { + return nil, err + } + return &backup, nil } func (m *Manager) dumpMySQL(ctx context.Context, deploymentName string, db *DatabaseSpec, dbDir string) (string, error) { @@ -530,6 +637,30 @@ func (m *Manager) listLocalBackups(filter *BackupListFilter) ([]Backup, error) { continue } + seen := make(map[string]bool) + for _, file := range files { + if !strings.HasSuffix(file.Name(), ".json") { + continue + } + backupID := strings.TrimSuffix(file.Name(), ".json") + backup, err := m.readBackupRecord(deploymentDir.Name(), backupID) + if err != nil { + continue + } + if backup.Path != "" { + if info, statErr := os.Stat(backup.Path); statErr == nil { + backup.Size = info.Size() + } else { + backup.Path = "" + backup.Locations = removeLocation(backup.Locations, locationLocal) + } + } + if filter.Status == "" || backup.Status == filter.Status { + backups = append(backups, *backup) + } + seen[backupID] = true + } + for _, file := range files { if !strings.HasSuffix(file.Name(), ".tar.gz") { continue @@ -541,6 +672,9 @@ func (m *Manager) listLocalBackups(filter *BackupListFilter) ([]Backup, error) { } backupID := strings.TrimSuffix(file.Name(), ".tar.gz") + if seen[backupID] { + continue + } backup := Backup{ ID: backupID, DeploymentName: deploymentDir.Name(), @@ -551,7 +685,9 @@ func (m *Manager) listLocalBackups(filter *BackupListFilter) ([]Backup, error) { Locations: []string{locationLocal}, } - backups = append(backups, backup) + if filter.Status == "" || backup.Status == filter.Status { + backups = append(backups, backup) + } } } @@ -574,6 +710,19 @@ func (m *Manager) getLocalBackup(backupID string) (*Backup, error) { deploymentName := parts[0] backupPath := filepath.Join(m.backupsPath, deploymentName, backupID+".tar.gz") + if backup, recordErr := m.readBackupRecord(deploymentName, backupID); recordErr == nil { + if info, statErr := os.Stat(backupPath); statErr == nil { + backup.Path = backupPath + backup.Size = info.Size() + if !containsLocation(backup.Locations, locationLocal) { + backup.Locations = append([]string{locationLocal}, backup.Locations...) + } + return backup, nil + } + backup.Path = "" + backup.Locations = removeLocation(backup.Locations, locationLocal) + return backup, nil + } info, err := os.Stat(backupPath) if err != nil { @@ -597,7 +746,32 @@ func (m *Manager) deleteLocalBackup(backupID string) error { return err } - return os.Remove(backup.Path) + if backup.Path != "" { + if err := os.Remove(backup.Path); err != nil && !os.IsNotExist(err) { + return err + } + } + _ = os.Remove(m.backupRecordPath(backup.DeploymentName, backup.ID)) + return nil +} + +func containsLocation(locations []string, target string) bool { + for _, location := range locations { + if location == target { + return true + } + } + return false +} + +func removeLocation(locations []string, target string) []string { + filtered := locations[:0] + for _, location := range locations { + if location != target { + filtered = append(filtered, location) + } + } + return filtered } func (m *Manager) GetBackupPath(backupID string) (string, error) { @@ -605,6 +779,9 @@ func (m *Manager) GetBackupPath(backupID string) (string, error) { if err != nil { return "", err } + if backup.Path == "" { + return "", fmt.Errorf("backup archive is unavailable: %s", backupID) + } return backup.Path, nil } @@ -617,12 +794,18 @@ func (m *Manager) CleanupOldBackups(deploymentName string, keepCount int) (int, return 0, err } - if len(backups) <= keepCount { + usable := backups[:0] + for _, backup := range backups { + if backup.Path != "" && backup.Status != BackupStatusFailed { + usable = append(usable, backup) + } + } + if len(usable) <= keepCount { return 0, nil } deleted := 0 - for _, backup := range backups[keepCount:] { + for _, backup := range usable[keepCount:] { if err := m.deleteLocalBackup(backup.ID); err != nil { log.Printf("Failed to delete old backup %s: %v", backup.ID, err) continue diff --git a/internal/backup/manager_test.go b/internal/backup/manager_test.go index 06c0277..f7ebfd4 100644 --- a/internal/backup/manager_test.go +++ b/internal/backup/manager_test.go @@ -5,6 +5,7 @@ import ( "os" "path/filepath" "testing" + "time" ) func setupTestManager(t *testing.T) (*Manager, string) { @@ -125,6 +126,38 @@ services: } } +func TestCreateBackup_CleanupFailureIsPartial(t *testing.T) { + m, tmpDir := setupTestManager(t) + defer os.RemoveAll(tmpDir) + seedDeployment(t, tmpDir, "app") + + binDir := t.TempDir() + if err := os.WriteFile(filepath.Join(binDir, "docker"), []byte("#!/bin/sh\nexit 1\n"), 0755); err != nil { + t.Fatal(err) + } + t.Setenv("PATH", binDir+string(os.PathListSeparator)+os.Getenv("PATH")) + + created, err := m.CreateBackup(context.Background(), "app", &BackupSpec{ + PostHooks: []HookSpec{{Service: "web", Command: "resume"}}, + }) + if err != nil { + t.Fatalf("create backup: %v", err) + } + if created.Status != BackupStatusPartial { + t.Fatalf("status = %s, want partial", created.Status) + } + if len(created.CleanupResults) != 1 || created.CleanupResults[0].Status != ResultStatusFailed { + t.Fatalf("cleanup results = %#v", created.CleanupResults) + } + persisted, err := m.GetBackup(created.ID) + if err != nil { + t.Fatal(err) + } + if persisted.Status != BackupStatusPartial { + t.Fatalf("persisted status = %s, want partial", persisted.Status) + } +} + func TestListBackups_Empty(t *testing.T) { m, tmpDir := setupTestManager(t) defer os.RemoveAll(tmpDir) @@ -329,6 +362,75 @@ func TestCleanupOldBackups(t *testing.T) { } } +func TestCleanupOldBackups_DoesNotCountFailedAttempts(t *testing.T) { + m, tmpDir := setupTestManager(t) + defer os.RemoveAll(tmpDir) + seedDeployment(t, tmpDir, "app") + + usable, err := m.CreateBackup(context.Background(), "app", nil) + if err != nil { + t.Fatal(err) + } + failed := &Backup{ + ID: "app_99991231_235959.000000000", + DeploymentName: "app", + Status: BackupStatusFailed, + CreatedAt: time.Now().Add(time.Hour), + } + if err := m.saveBackupRecord(failed); err != nil { + t.Fatal(err) + } + + deleted, err := m.CleanupOldBackups("app", 1) + if err != nil { + t.Fatal(err) + } + if deleted != 0 { + t.Fatalf("deleted = %d, want 0", deleted) + } + if _, err := os.Stat(usable.Path); err != nil { + t.Fatalf("usable recovery point was removed: %v", err) + } +} + +func TestCreateBackup_RecordsEachConfiguredSource(t *testing.T) { + m, tmpDir := setupTestManager(t) + defer os.RemoveAll(tmpDir) + seedDeployment(t, tmpDir, "app") + + binDir := t.TempDir() + if err := os.WriteFile(filepath.Join(binDir, "docker"), []byte("#!/bin/sh\nexit 1\n"), 0755); err != nil { + t.Fatal(err) + } + t.Setenv("PATH", binDir+string(os.PathListSeparator)+os.Getenv("PATH")) + + created, err := m.CreateBackup(context.Background(), "app", &BackupSpec{ + ContainerPaths: []ContainerPath{ + {Service: "web", ContainerPath: "/optional", Required: false}, + {Service: "web", ContainerPath: "/required", Required: true}, + }, + Databases: []DatabaseSpec{ + {Service: "primary", Type: "unsupported"}, + {Service: "analytics", Type: "unsupported"}, + }, + }) + if err == nil { + t.Fatal("required source failures were accepted") + } + if len(created.ComponentResults) != 8 { + t.Fatalf("component results = %#v", created.ComponentResults) + } + if created.ComponentResults[4].Required || created.ComponentResults[4].Status != ResultStatusFailed { + t.Fatalf("optional result = %#v", created.ComponentResults[4]) + } + if !created.ComponentResults[5].Required || created.ComponentResults[5].Status != ResultStatusFailed { + t.Fatalf("required result = %#v", created.ComponentResults[5]) + } + if created.ComponentResults[6].Name != "primary" || created.ComponentResults[7].Name != "analytics" { + t.Fatalf("database results = %#v", created.ComponentResults[6:]) + } +} + func TestStartBackupJob(t *testing.T) { m, tmpDir := setupTestManager(t) defer os.RemoveAll(tmpDir) diff --git a/internal/backup/remotes.go b/internal/backup/remotes.go index 3397518..df7946f 100644 --- a/internal/backup/remotes.go +++ b/internal/backup/remotes.go @@ -47,34 +47,77 @@ func deploymentFromID(backupID string) string { return parts[0] } -// mirrorToRemotes uploads a freshly written local archive to every configured -// remote, returning the names of the destinations that accepted it. A failed -// upload is logged and skipped so a remote outage never fails the backup, whose -// local copy already succeeded. -func (m *Manager) mirrorToRemotes(ctx context.Context, deploymentName, backupID, archivePath string, size int64) []string { +func (m *Manager) mirrorToRemotes(ctx context.Context, deploymentName, backupID, archivePath string, size int64) []DestinationResult { remotes := m.getRemotes() if len(remotes) == 0 { return nil } key := backupKey(deploymentName, backupID) - var locations []string + var results []DestinationResult for _, r := range remotes { + result := DestinationResult{Name: r.Name(), Status: ResultStatusCompleted} f, err := os.Open(archivePath) if err != nil { log.Printf("Backup: mirror to %s failed to open archive: %v", r.Name(), err) + result.Status = ResultStatusFailed + result.Error = err.Error() + results = append(results, result) continue } err = r.Put(ctx, key, f, size) f.Close() if err != nil { log.Printf("Backup: mirror to %s failed: %v", r.Name(), err) + result.Status = ResultStatusFailed + result.Error = "upload failed" + results = append(results, result) continue } - locations = append(locations, r.Name()) + results = append(results, result) log.Printf("Backup mirrored: %s -> %s", backupID, r.Name()) } - return locations + return results +} + +func (m *Manager) RetryRemotePublication(ctx context.Context, backupID string) (*Backup, error) { + backup, err := m.GetBackup(backupID) + if err != nil { + return nil, err + } + if backup.Path == "" { + return nil, fmt.Errorf("local backup archive is unavailable") + } + backup.DestinationResults = m.mirrorToRemotes(ctx, backup.DeploymentName, backup.ID, backup.Path, backup.Size) + backup.Locations = []string{locationLocal} + succeeded := 0 + for _, result := range backup.DestinationResults { + if result.Status == ResultStatusCompleted { + succeeded++ + backup.Locations = append(backup.Locations, result.Name) + } + } + switch { + case len(backup.DestinationResults) > 0 && succeeded == 0: + backup.Status = BackupStatusLocalOnly + case succeeded < len(backup.DestinationResults) || hasFailedResult(backup.CleanupResults): + backup.Status = BackupStatusPartial + default: + backup.Status = BackupStatusCompleted + } + if err := m.saveBackupRecord(backup); err != nil { + return nil, err + } + return backup, nil +} + +func hasFailedResult(results []ComponentResult) bool { + for _, result := range results { + if result.Status == ResultStatusFailed { + return true + } + } + return false } // ListBackups returns backups from the local disk merged with those in every diff --git a/internal/backup/remotes_test.go b/internal/backup/remotes_test.go index d9c3b9f..a043cad 100644 --- a/internal/backup/remotes_test.go +++ b/internal/backup/remotes_test.go @@ -20,6 +20,7 @@ type fakeStore struct { objects map[string][]byte modtimes map[string]time.Time failList bool + failPut bool } func newFakeStore(name string) *fakeStore { @@ -29,6 +30,9 @@ func newFakeStore(name string) *fakeStore { func (f *fakeStore) Name() string { return f.name } func (f *fakeStore) Put(_ context.Context, key string, r io.Reader, _ int64) error { + if f.failPut { + return fmt.Errorf("put failed") + } data, err := io.ReadAll(r) if err != nil { return err @@ -40,6 +44,61 @@ func (f *fakeStore) Put(_ context.Context, key string, r io.Reader, _ int64) err return nil } +func TestCreateBackup_ReportsPartialAndLocalOnlyPublication(t *testing.T) { + m, tmpDir := setupTestManager(t) + defer os.RemoveAll(tmpDir) + seedDeployment(t, tmpDir, "app") + + working := newFakeStore("working") + failing := newFakeStore("failing") + failing.failPut = true + m.SetRemotes([]Store{working, failing}) + + partial, err := m.CreateBackup(context.Background(), "app", nil) + if err != nil { + t.Fatalf("create partial backup: %v", err) + } + if partial.Status != BackupStatusPartial { + t.Fatalf("status = %s, want partial", partial.Status) + } + + working.failPut = true + localOnly, err := m.CreateBackup(context.Background(), "app", nil) + if err != nil { + t.Fatalf("create local backup: %v", err) + } + if localOnly.Status != BackupStatusLocalOnly { + t.Fatalf("status = %s, want local_only", localOnly.Status) + } + if !containsStr(localOnly.Locations, locationLocal) { + t.Fatalf("local archive was not retained: %v", localOnly.Locations) + } +} + +func TestRetryRemotePublication_ReusesLocalArchive(t *testing.T) { + m, tmpDir := setupTestManager(t) + defer os.RemoveAll(tmpDir) + seedDeployment(t, tmpDir, "app") + + remote := newFakeStore("remote") + remote.failPut = true + m.SetRemotes([]Store{remote}) + created, err := m.CreateBackup(context.Background(), "app", nil) + if err != nil { + t.Fatalf("create backup: %v", err) + } + archivePath := created.Path + + remote.failPut = false + retried, err := m.RetryRemotePublication(context.Background(), created.ID) + if err != nil { + t.Fatalf("retry publication: %v", err) + } + if retried.Status != BackupStatusCompleted || retried.Path != archivePath { + t.Fatalf("retry result = %#v", retried) + } +} + func (f *fakeStore) Open(_ context.Context, key string) (io.ReadCloser, error) { f.mu.Lock() defer f.mu.Unlock() diff --git a/internal/backup/types.go b/internal/backup/types.go index 17908da..f0708bd 100644 --- a/internal/backup/types.go +++ b/internal/backup/types.go @@ -12,9 +12,33 @@ const ( BackupStatusPending BackupStatus = "pending" BackupStatusInProgress BackupStatus = "in_progress" BackupStatusCompleted BackupStatus = "completed" + BackupStatusPartial BackupStatus = "partial" + BackupStatusLocalOnly BackupStatus = "local_only" BackupStatusFailed BackupStatus = "failed" ) +type ResultStatus string + +const ( + ResultStatusCompleted ResultStatus = "completed" + ResultStatusSkipped ResultStatus = "skipped" + ResultStatusFailed ResultStatus = "failed" +) + +type ComponentResult struct { + Name string `json:"name"` + Kind string `json:"kind"` + Required bool `json:"required"` + Status ResultStatus `json:"status"` + Error string `json:"error,omitempty"` +} + +type DestinationResult struct { + Name string `json:"name"` + Status ResultStatus `json:"status"` + Error string `json:"error,omitempty"` +} + type BackupSpec = models.BackupSpec type ContainerPath = models.ContainerBackupPath type DatabaseSpec = models.DatabaseBackupSpec @@ -34,17 +58,21 @@ type Backup struct { // Locations lists where this backup exists: "local" and/or remote // destination names. A backup may live remotely only if local retention // has pruned the on-disk copy. - Locations []string `json:"locations,omitempty"` + Locations []string `json:"locations,omitempty"` + ComponentResults []ComponentResult `json:"component_results,omitempty"` + CleanupResults []ComponentResult `json:"cleanup_results,omitempty"` + DestinationResults []DestinationResult `json:"destination_results,omitempty"` } type BackupMetadata struct { - ID string `json:"id"` - DeploymentName string `json:"deployment_name"` - DeploymentPath string `json:"deployment_path"` - CreatedAt time.Time `json:"created_at"` - AgentVersion string `json:"agent_version"` - Components BackupComponents `json:"components"` - ContainerStates map[string]string `json:"container_states,omitempty"` + ID string `json:"id"` + DeploymentName string `json:"deployment_name"` + DeploymentPath string `json:"deployment_path"` + CreatedAt time.Time `json:"created_at"` + AgentVersion string `json:"agent_version"` + Components BackupComponents `json:"components"` + ContainerStates map[string]string `json:"container_states,omitempty"` + ComponentResults []ComponentResult `json:"component_results,omitempty"` } type BackupComponents struct { diff --git a/internal/networks/manager.go b/internal/networks/manager.go index 8e64e3d..a89e8bc 100644 --- a/internal/networks/manager.go +++ b/internal/networks/manager.go @@ -184,19 +184,18 @@ func (m *Manager) DisconnectContainer(networkName, containerName string) error { } func (m *Manager) IsContainerOnNetwork(networkName, containerName string) bool { - cmd := exec.Command("docker", "network", "inspect", networkName, - "--format", "{{range .Containers}}{{.Name}} {{end}}") + cmd := exec.Command("docker", "inspect", containerName, + "--format", "{{json .NetworkSettings.Networks}}") output, err := cmd.Output() if err != nil { return false } - containers := strings.Fields(string(output)) - for _, c := range containers { - if c == containerName { - return true - } + var attached map[string]json.RawMessage + if err := json.Unmarshal(output, &attached); err != nil { + return false } - return false + _, ok := attached[networkName] + return ok } func (m *Manager) EnsureContainerOnNetwork(networkName, containerName string) error { diff --git a/internal/networks/manager_integration_test.go b/internal/networks/manager_integration_test.go new file mode 100644 index 0000000..ea40c22 --- /dev/null +++ b/internal/networks/manager_integration_test.go @@ -0,0 +1,49 @@ +package networks + +import ( + "fmt" + "os" + "os/exec" + "strings" + "testing" + "time" +) + +func TestEnsureContainerOnNetworkAcceptsNameAndID(t *testing.T) { + if _, err := exec.LookPath("docker"); err != nil { + t.Skip("docker is unavailable") + } + image := os.Getenv("FLATRUN_DOCKER_TEST_IMAGE") + if image == "" { + image = "busybox:latest" + } + if err := exec.Command("docker", "image", "inspect", image).Run(); err != nil { + t.Skipf("docker test image %s is unavailable", image) + } + + suffix := fmt.Sprintf("%d", time.Now().UnixNano()) + networkName := "flatrun-network-test-" + suffix + containerName := "flatrun-container-test-" + suffix + if output, err := exec.Command("docker", "network", "create", networkName).CombinedOutput(); err != nil { + t.Fatalf("create network: %v: %s", err, output) + } + t.Cleanup(func() { _ = exec.Command("docker", "network", "rm", networkName).Run() }) + + output, err := exec.Command("docker", "create", "--name", containerName, image, "sleep", "60").CombinedOutput() + if err != nil { + t.Fatalf("create container: %v: %s", err, output) + } + containerID := strings.TrimSpace(string(output)) + t.Cleanup(func() { _ = exec.Command("docker", "rm", "-f", containerID).Run() }) + + manager := NewManager() + if err := manager.EnsureContainerOnNetwork(networkName, containerName); err != nil { + t.Fatalf("attach by name: %v", err) + } + if err := manager.EnsureContainerOnNetwork(networkName, containerID); err != nil { + t.Fatalf("repeat attachment by ID: %v", err) + } + if !manager.IsContainerOnNetwork(networkName, containerName) || !manager.IsContainerOnNetwork(networkName, containerID) { + t.Fatal("container was not attached for both identifiers") + } +}