diff --git a/internal/api/authz_test.go b/internal/api/authz_test.go index 6afa0b8..1e62104 100644 --- a/internal/api/authz_test.go +++ b/internal/api/authz_test.go @@ -390,6 +390,56 @@ func TestRestoreBackupRequiresWriteOnTargetDeployment(t *testing.T) { } } +func TestIsolatedRestoreRequiresNewDeploymentName(t *testing.T) { + gin.SetMode(gin.TestMode) + tmpDir := t.TempDir() + createTestDeployment(t, tmpDir, "source-app", &models.ServiceMetadata{Name: "source-app"}) + backupManager, err := backup.NewManager(tmpDir) + if err != nil { + t.Fatal(err) + } + created, err := backupManager.CreateBackup(context.Background(), "source-app", nil) + if err != nil { + t.Fatal(err) + } + server := &Server{backupManager: backupManager} + router := gin.New() + router.POST("/backups/:id/restore", server.restoreBackup) + req := httptest.NewRequest(http.MethodPost, "/backups/"+created.ID+"/restore", bytes.NewBufferString("{\"isolated\":true}")) + req.Header.Set("Content-Type", "application/json") + w := httptest.NewRecorder() + router.ServeHTTP(w, req) + if w.Code != http.StatusBadRequest { + t.Fatalf("expected 400, got %d: %s", w.Code, w.Body.String()) + } +} + +func TestIsolatedRestoreRejectsExistingTargetBeforeStarting(t *testing.T) { + gin.SetMode(gin.TestMode) + tmpDir := t.TempDir() + createTestDeployment(t, tmpDir, "source-app", &models.ServiceMetadata{Name: "source-app"}) + createTestDeployment(t, tmpDir, "target-app", &models.ServiceMetadata{Name: "target-app"}) + backupManager, err := backup.NewManager(tmpDir) + if err != nil { + t.Fatal(err) + } + created, err := backupManager.CreateBackup(context.Background(), "source-app", nil) + if err != nil { + t.Fatal(err) + } + server := &Server{backupManager: backupManager, manager: docker.NewManager(tmpDir)} + router := gin.New() + router.POST("/backups/:id/restore", server.restoreBackup) + body := bytes.NewBufferString(`{"isolated":true,"deployment_name":"target-app"}`) + req := httptest.NewRequest(http.MethodPost, "/backups/"+created.ID+"/restore", body) + req.Header.Set("Content-Type", "application/json") + w := httptest.NewRecorder() + router.ServeHTTP(w, req) + if w.Code != http.StatusConflict { + t.Fatalf("expected 409, got %d: %s", w.Code, w.Body.String()) + } +} + func TestCreateScheduledTaskRequiresWriteDeploymentAccess(t *testing.T) { gin.SetMode(gin.TestMode) diff --git a/internal/api/backup_destination_selection_test.go b/internal/api/backup_destination_selection_test.go new file mode 100644 index 0000000..3a13fb6 --- /dev/null +++ b/internal/api/backup_destination_selection_test.go @@ -0,0 +1,81 @@ +package api + +import ( + "encoding/json" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "strings" + "testing" + + "github.com/flatrun/agent/internal/docker" + "github.com/flatrun/agent/pkg/config" + "github.com/flatrun/agent/pkg/models" + "github.com/gin-gonic/gin" + "gopkg.in/yaml.v3" +) + +func TestDeploymentBackupDestinationsThroughHTTP(t *testing.T) { + 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 app:\n image: nginx\n"), 0644); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(deploymentDir, "service.yml"), []byte("name: app\n"), 0644); err != nil { + t.Fatal(err) + } + disabled := false + server := &Server{ + manager: docker.NewManager(root), + config: &config.Config{Backup: config.BackupConfig{Destinations: []config.BackupDestination{ + {Name: "primary", Kind: "external", CredentialID: "private-credential"}, + {Name: "disabled", Kind: "external", Enabled: &disabled}, + }}}, + } + router := gin.New() + router.GET("/deployments/:name/backup-destinations", server.listDeploymentBackupDestinationOptions) + router.PUT("/deployments/:name/backup-config", server.updateDeploymentBackupConfig) + + options := httptest.NewRecorder() + router.ServeHTTP(options, httptest.NewRequest(http.MethodGet, "/deployments/app/backup-destinations", nil)) + if options.Code != http.StatusOK || strings.Contains(options.Body.String(), "private-credential") || strings.Contains(options.Body.String(), "disabled") { + t.Fatalf("options response = %d %s", options.Code, options.Body.String()) + } + + unknown := httptest.NewRecorder() + router.ServeHTTP(unknown, httptest.NewRequest(http.MethodPut, "/deployments/app/backup-config", strings.NewReader(`{"destinations":["missing"]}`))) + if unknown.Code != http.StatusBadRequest { + t.Fatalf("unknown destination status = %d: %s", unknown.Code, unknown.Body.String()) + } + + saved := httptest.NewRecorder() + router.ServeHTTP(saved, httptest.NewRequest(http.MethodPut, "/deployments/app/backup-config", strings.NewReader(`{"destinations":["primary"]}`))) + if saved.Code != http.StatusOK { + t.Fatalf("save status = %d: %s", saved.Code, saved.Body.String()) + } + data, err := os.ReadFile(filepath.Join(deploymentDir, "service.yml")) + if err != nil { + t.Fatal(err) + } + var metadata models.ServiceMetadata + if err := yaml.Unmarshal(data, &metadata); err != nil { + t.Fatal(err) + } + if metadata.Backup == nil || len(metadata.Backup.Destinations) != 1 || metadata.Backup.Destinations[0] != "primary" { + t.Fatalf("saved backup config = %#v", metadata.Backup) + } + + var response struct { + BackupConfig models.BackupSpec `json:"backup_config"` + } + if err := json.Unmarshal(saved.Body.Bytes(), &response); err != nil { + t.Fatal(err) + } + if len(response.BackupConfig.Destinations) != 1 || response.BackupConfig.Destinations[0] != "primary" { + t.Fatalf("response backup config = %#v", response.BackupConfig) + } +} diff --git a/internal/api/backup_destinations.go b/internal/api/backup_destinations.go index 88da8ea..23c50db 100644 --- a/internal/api/backup_destinations.go +++ b/internal/api/backup_destinations.go @@ -157,6 +157,20 @@ func (s *Server) listBackupDestinations(c *gin.Context) { c.JSON(http.StatusOK, gin.H{"destinations": dests}) } +func (s *Server) listDeploymentBackupDestinationOptions(c *gin.Context) { + type option struct { + Name string `json:"name"` + Kind string `json:"kind"` + } + options := make([]option, 0) + for _, destination := range s.config.Backup.Destinations { + if destination.IsEnabled() { + options = append(options, option{Name: destination.Name, Kind: destination.StoreKind()}) + } + } + c.JSON(http.StatusOK, gin.H{"destinations": options}) +} + func (s *Server) findDestinationByName(name string) (config.BackupDestination, bool) { for _, d := range s.config.Backup.Destinations { if d.Name == name { diff --git a/internal/api/backup_handlers.go b/internal/api/backup_handlers.go index d0eb387..3c5adf2 100644 --- a/internal/api/backup_handlers.go +++ b/internal/api/backup_handlers.go @@ -2,15 +2,107 @@ package api import ( "context" + "fmt" "net/http" "strconv" "github.com/flatrun/agent/internal/auth" "github.com/flatrun/agent/internal/backup" + "github.com/flatrun/agent/internal/scheduler" "github.com/flatrun/agent/pkg/models" "github.com/gin-gonic/gin" ) +type deploymentBackupPolicy struct { + Config *backup.BackupSpec `json:"config"` + Schedules []scheduler.ScheduledTask `json:"schedules"` + BackupCount int `json:"backup_count"` + LocalBytes int64 `json:"local_bytes"` + FailedCount int `json:"failed_count"` + SizeAlert bool `json:"size_alert"` + CleanupPreview *backup.CleanupPreview `json:"cleanup_preview"` +} + +func (s *Server) getDeploymentBackupPolicy(c *gin.Context) { + name := c.Param("name") + deployment, err := s.manager.GetDeployment(name) + if err != nil { + c.JSON(http.StatusNotFound, gin.H{"error": "Deployment not found"}) + return + } + spec := s.effectiveBackupSpec(deployment) + backups, err := s.backupManager.ListBackups(&backup.BackupListFilter{DeploymentName: name}) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + policy := deploymentBackupPolicy{Config: spec, Schedules: []scheduler.ScheduledTask{}, BackupCount: len(backups)} + for _, item := range backups { + if containsLocation(item.Locations, "local") { + policy.LocalBytes += item.Size + } + if item.Status == backup.BackupStatusFailed || item.Status == backup.BackupStatusPartial || item.Status == backup.BackupStatusLocalOnly { + policy.FailedCount++ + } + } + if s.schedulerManager != nil { + tasks, taskErr := s.schedulerManager.GetTasksByDeployment(name) + if taskErr == nil { + for _, task := range tasks { + if task.Type == scheduler.TaskTypeBackup { + policy.Schedules = append(policy.Schedules, task) + } + } + } + } + keep := spec.RetentionCount + if keep < 1 { + keep = 7 + } + policy.CleanupPreview, _ = s.backupManager.PreviewCleanup(name, keep) + policy.SizeAlert = spec.SizeAlertBytes > 0 && policy.LocalBytes >= spec.SizeAlertBytes + c.JSON(http.StatusOK, gin.H{"policy": policy}) +} + +func containsLocation(locations []string, wanted string) bool { + for _, location := range locations { + if location == wanted { + return true + } + } + return false +} + +func (s *Server) previewDeploymentBackupCleanup(c *gin.Context) { + keep, err := strconv.Atoi(c.DefaultQuery("keep", "7")) + if err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": "Invalid retention count"}) + return + } + preview, err := s.backupManager.PreviewCleanup(c.Param("name"), keep) + if err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"preview": preview}) +} + +func (s *Server) cleanupDeploymentBackups(c *gin.Context) { + var req struct { + Keep int `json:"keep" binding:"required,min=1"` + } + if err := c.ShouldBindJSON(&req); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + deleted, err := s.backupManager.CleanupOldBackups(c.Param("name"), req.Keep) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"deleted": deleted}) +} + func (s *Server) retryBackupPublication(c *gin.Context) { if s.backupManager == nil { c.JSON(http.StatusServiceUnavailable, gin.H{"error": "Backup manager not enabled"}) @@ -265,6 +357,14 @@ func (s *Server) updateDeploymentBackupConfig(c *gin.Context) { c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) return } + if err := s.validateBackupDestinations(spec.Destinations); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + if spec.RetentionCount < 0 || spec.SizeAlertBytes < 0 { + c.JSON(http.StatusBadRequest, gin.H{"error": "Retention and size alert values cannot be negative"}) + return + } if deployment.Metadata == nil { deployment.Metadata = &models.ServiceMetadata{} @@ -279,6 +379,26 @@ func (s *Server) updateDeploymentBackupConfig(c *gin.Context) { c.JSON(http.StatusOK, gin.H{"backup_config": spec}) } +func (s *Server) validateBackupDestinations(names []string) error { + enabled := make(map[string]bool) + for _, destination := range s.config.Backup.Destinations { + if destination.IsEnabled() { + enabled[destination.Name] = true + } + } + seen := make(map[string]bool, len(names)) + for _, name := range names { + if name == "" || !enabled[name] { + return fmt.Errorf("backup destination %q is unavailable", name) + } + if seen[name] { + return fmt.Errorf("backup destination %q is selected more than once", name) + } + seen[name] = true + } + return nil +} + func (s *Server) restoreBackup(c *gin.Context) { if s.backupManager == nil { c.JSON(http.StatusServiceUnavailable, gin.H{"error": "Backup manager not enabled"}) @@ -305,11 +425,31 @@ func (s *Server) restoreBackup(c *gin.Context) { if !s.requireDeploymentAccess(c, b.DeploymentName, auth.AccessLevelRead) { return } - if !s.requireDeploymentAccess(c, targetDeployment, auth.AccessLevelWrite) { + if req.Isolated { + if req.DeploymentName == "" || req.DeploymentName == b.DeploymentName { + c.JSON(http.StatusBadRequest, gin.H{"error": "Isolated restore requires a new deployment name"}) + return + } + actor := auth.GetActorFromContext(c) + if actor != nil && !actor.HasPermission(auth.PermDeploymentsWrite) { + c.JSON(http.StatusForbidden, gin.H{"error": "Deployment write permission required"}) + return + } + if _, lookupErr := s.manager.GetDeployment(targetDeployment); lookupErr == nil { + c.JSON(http.StatusConflict, gin.H{"error": "Deployment already exists"}) + return + } + } else if !s.requireDeploymentAccess(c, targetDeployment, auth.AccessLevelWrite) { return } - jobID := s.backupManager.StartRestoreJob(&req) + actor := auth.GetActorFromContext(c) + jobID := s.backupManager.StartRestoreJob(&req, func() error { + if !req.Isolated || s.authManager == nil || actor == nil || actor.User == nil || actor.Role == auth.RoleAdmin { + return nil + } + return s.authManager.AssignDeployment(actor.User.ID, targetDeployment, auth.AccessLevelAdmin, actor.User.ID) + }) c.JSON(http.StatusAccepted, gin.H{"job_id": jobID, "message": "Restore job started"}) } diff --git a/internal/api/database_attach.go b/internal/api/database_attach.go new file mode 100644 index 0000000..5d70109 --- /dev/null +++ b/internal/api/database_attach.go @@ -0,0 +1,103 @@ +package api + +import ( + "net/http" + "os" + "path/filepath" + + "github.com/flatrun/agent/internal/docker" + "github.com/flatrun/agent/pkg/models" + "github.com/gin-gonic/gin" +) + +func (s *Server) attachDeploymentDatabase(c *gin.Context) { + name := c.Param("name") + var req DatabaseConfigRequest + if err := c.ShouldBindJSON(&req); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + if err := req.Validate(); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + deployment, err := s.manager.GetDeployment(name) + if err != nil { + c.JSON(http.StatusNotFound, gin.H{"error": "Deployment not found"}) + return + } + envVars, configs, err := s.createDatabasesForDeployment(name, []DatabaseConfigRequest{req}) + if err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + deployDir := filepath.Join(s.config.DeploymentsPath, name) + merged, err := readDeploymentEnv(deployDir) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + merged = mergeEnvVars(merged, envVars) + if err := s.writeEnvFile(name, merged); err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + composePath := filepath.Join(deployDir, "docker-compose.yml") + content, err := os.ReadFile(composePath) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + updated, err := docker.EnsureServiceEnvFile(string(content), ".env.flatrun") + if err == nil && (req.Mode == "shared" || req.Mode == "existing") { + updated = s.addDatabaseNetwork(updated) + } + if err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + if err := os.WriteFile(composePath, []byte(updated), 0644); err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + if deployment.Metadata == nil { + deployment.Metadata = &models.ServiceMetadata{Name: name} + } + deployment.Metadata.Databases = append(deployment.Metadata.Databases, configs...) + if err := s.manager.SaveMetadata(name, deployment.Metadata); err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"databases": configs, "message": "Database attached. Restart the deployment to apply."}) +} + +func readDeploymentEnv(deployDir string) ([]EnvVar, error) { + var merged []EnvVar + for _, filename := range []string{".env", ".env.flatrun"} { + current, err := os.ReadFile(filepath.Join(deployDir, filename)) + if os.IsNotExist(err) { + continue + } + if err != nil { + return nil, err + } + merged = mergeEnvVars(merged, parseEnvContent(string(current))) + } + return merged, nil +} + +func mergeEnvVars(current, additions []EnvVar) []EnvVar { + index := make(map[string]int, len(current)) + for i, item := range current { + index[item.Key] = i + } + for _, item := range additions { + if i, ok := index[item.Key]; ok { + current[i] = item + } else { + index[item.Key] = len(current) + current = append(current, item) + } + } + return current +} diff --git a/internal/api/database_attach_test.go b/internal/api/database_attach_test.go new file mode 100644 index 0000000..239a200 --- /dev/null +++ b/internal/api/database_attach_test.go @@ -0,0 +1,48 @@ +package api + +import ( + "os" + "path/filepath" + "testing" +) + +func TestMergeDatabaseEnvironmentPreservesUnrelatedValues(t *testing.T) { + got := mergeEnvVars( + []EnvVar{{Key: "APP_ENV", Value: "production"}, {Key: "DB_HOST", Value: "old"}}, + []EnvVar{{Key: "DB_HOST", Value: "database"}, {Key: "DB_NAME", Value: "shop"}}, + ) + want := map[string]string{"APP_ENV": "production", "DB_HOST": "database", "DB_NAME": "shop"} + for _, item := range got { + if want[item.Key] != item.Value { + t.Fatalf("unexpected environment value: %#v", got) + } + delete(want, item.Key) + } + if len(want) != 0 { + t.Fatalf("missing environment values: %#v", want) + } +} + +func TestReadDeploymentEnvPreservesImportedValues(t *testing.T) { + dir := t.TempDir() + if err := os.WriteFile(filepath.Join(dir, ".env"), []byte("APP_ENV=production\nDB_HOST=imported\n"), 0600); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(dir, ".env.flatrun"), []byte("DB_HOST=managed\n"), 0600); err != nil { + t.Fatal(err) + } + got, err := readDeploymentEnv(dir) + if err != nil { + t.Fatal(err) + } + want := map[string]string{"APP_ENV": "production", "DB_HOST": "managed"} + for _, item := range got { + if want[item.Key] != item.Value { + t.Fatalf("unexpected environment value: %#v", got) + } + delete(want, item.Key) + } + if len(want) != 0 { + t.Fatalf("missing environment values: %#v", want) + } +} diff --git a/internal/api/deployment_actions.go b/internal/api/deployment_actions.go index e27a171..96c0390 100644 --- a/internal/api/deployment_actions.go +++ b/internal/api/deployment_actions.go @@ -166,6 +166,9 @@ func (s *Server) mutateDomainUpdate(deployment *models.Deployment, domainID stri if err := s.validateDomainAccess(updated.Access); err != nil { return err } + if updated.Domain == "" { + return apiErrf(http.StatusBadRequest, "Domain is required") + } if deployment.Metadata == nil || len(deployment.Metadata.Domains) == 0 { return apiErrf(http.StatusNotFound, "Domain not found") } @@ -180,6 +183,11 @@ func (s *Server) mutateDomainUpdate(deployment *models.Deployment, domainID stri for i, d := range deployment.Metadata.Domains { if d.ID == domainID { + for _, existing := range deployment.Metadata.Domains { + if existing.ID != domainID && existing.Domain == updated.Domain && existing.PathPrefix == updated.PathPrefix { + return apiErrf(http.StatusConflict, "Domain %s%s already exists", updated.Domain, updated.PathPrefix) + } + } updated.ID = domainID if updated.Service == "" { updated.Service = d.Service diff --git a/internal/api/domains_test.go b/internal/api/domains_test.go index 5ecfa49..174fd15 100644 --- a/internal/api/domains_test.go +++ b/internal/api/domains_test.go @@ -554,3 +554,52 @@ func TestAddDomain(t *testing.T) { } }) } + +func TestDomainListChangesPersistThroughHTTP(t *testing.T) { + server, tmpDir, cleanup := setupDomainsTestServer(t) + defer cleanup() + + createTestDeployment(t, tmpDir, "domain-list", &models.ServiceMetadata{ + Name: "domain-list", + Type: "web", + Domains: []models.DomainConfig{{ + ID: "primary", + Service: "web", + ContainerPort: 80, + Domain: "app.example.com", + }}, + }) + + router := gin.New() + router.POST("/deployments/:name/domains", server.addDomain) + router.PUT("/deployments/:name/domains/:domainId", server.updateDomain) + + addBody := `{"service":"web","container_port":80,"domain":"alias.example.com"}` + add := httptest.NewRecorder() + router.ServeHTTP(add, httptest.NewRequest(http.MethodPost, "/deployments/domain-list/domains", strings.NewReader(addBody))) + if add.Code != http.StatusCreated { + t.Fatalf("add status = %d: %s", add.Code, add.Body.String()) + } + + updateBody := `{"service":"web","container_port":80,"domain":"new.example.com"}` + update := httptest.NewRecorder() + router.ServeHTTP(update, httptest.NewRequest(http.MethodPut, "/deployments/domain-list/domains/primary", strings.NewReader(updateBody))) + if update.Code != http.StatusOK { + t.Fatalf("update status = %d: %s", update.Code, update.Body.String()) + } + + data, err := os.ReadFile(filepath.Join(tmpDir, "domain-list", "service.yml")) + if err != nil { + t.Fatal(err) + } + var metadata models.ServiceMetadata + if err := yaml.Unmarshal(data, &metadata); err != nil { + t.Fatal(err) + } + if len(metadata.Domains) != 2 { + t.Fatalf("domains = %+v, want two saved domains", metadata.Domains) + } + if metadata.Domains[0].Domain != "new.example.com" || metadata.Domains[1].Domain != "alias.example.com" { + t.Fatalf("domains = %+v", metadata.Domains) + } +} diff --git a/internal/api/migration_handlers.go b/internal/api/migration_handlers.go new file mode 100644 index 0000000..abe539d --- /dev/null +++ b/internal/api/migration_handlers.go @@ -0,0 +1,124 @@ +package api + +import ( + "net" + "net/http" + "strings" + "time" + + "github.com/flatrun/agent/pkg/models" + "github.com/gin-gonic/gin" +) + +type migrationStatus struct { + Plan *models.MigrationSpec `json:"plan"` + RetirementReady bool `json:"retirement_ready"` + Blockers []string `json:"blockers"` +} + +func (s *Server) getDeploymentMigration(c *gin.Context) { + deployment, err := s.manager.GetDeployment(c.Param("name")) + if err != nil { + c.JSON(http.StatusNotFound, gin.H{"error": "Deployment not found"}) + return + } + var plan *models.MigrationSpec + if deployment.Metadata != nil { + plan = deployment.Metadata.Migration + } + c.JSON(http.StatusOK, gin.H{"migration": buildMigrationStatus(plan)}) +} + +func (s *Server) updateDeploymentMigration(c *gin.Context) { + name := c.Param("name") + deployment, err := s.manager.GetDeployment(name) + if err != nil { + c.JSON(http.StatusNotFound, gin.H{"error": "Deployment not found"}) + return + } + var plan models.MigrationSpec + if err := c.ShouldBindJSON(&plan); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + seen := map[string]bool{} + for _, site := range plan.Sites { + hostname := strings.TrimSpace(strings.ToLower(site.Hostname)) + if hostname == "" || seen[hostname] { + c.JSON(http.StatusBadRequest, gin.H{"error": "Migration hostnames must be present and unique"}) + return + } + seen[hostname] = true + } + if deployment.Metadata == nil { + deployment.Metadata = &models.ServiceMetadata{Name: name} + } + deployment.Metadata.Migration = &plan + if err := s.manager.SaveMetadata(name, deployment.Metadata); err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"migration": buildMigrationStatus(&plan)}) +} + +func (s *Server) checkDeploymentMigrationDNS(c *gin.Context) { + name := c.Param("name") + deployment, err := s.manager.GetDeployment(name) + if err != nil || deployment.Metadata == nil || deployment.Metadata.Migration == nil { + c.JSON(http.StatusNotFound, gin.H{"error": "Migration plan not found"}) + return + } + plan := deployment.Metadata.Migration + for i := range plan.Sites { + addresses, lookupErr := net.LookupHost(plan.Sites[i].Hostname) + if lookupErr != nil { + addresses = []string{} + } + plan.Sites[i].Resolved = addresses + plan.Sites[i].DNSPropagated = containsString(addresses, plan.ExpectedAddress) + } + if err := s.manager.SaveMetadata(name, deployment.Metadata); err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"migration": buildMigrationStatus(plan)}) +} + +func containsString(values []string, wanted string) bool { + if wanted == "" { + return false + } + for _, value := range values { + if value == wanted { + return true + } + } + return false +} + +func buildMigrationStatus(plan *models.MigrationSpec) migrationStatus { + status := migrationStatus{Plan: plan, Blockers: []string{}} + if plan == nil { + status.Blockers = append(status.Blockers, "migration plan is missing") + return status + } + if !plan.InventoryComplete { + status.Blockers = append(status.Blockers, "source inventory is incomplete") + } + for _, site := range plan.Sites { + if !site.Transferred { + status.Blockers = append(status.Blockers, site.Hostname+" has not completed its initial transfer") + } + if !site.DNSPropagated && plan.CutoverAt != nil { + status.Blockers = append(status.Blockers, site.Hostname+" DNS has not propagated") + } + } + if plan.LastSyncAt == nil || time.Since(*plan.LastSyncAt) > 15*time.Minute { + status.Blockers = append(status.Blockers, "a recent final synchronization is required") + } + if plan.CutoverAt == nil { + status.Blockers = append(status.Blockers, "cutover has not been recorded") + } + status.RetirementReady = len(status.Blockers) == 0 + return status +} diff --git a/internal/api/migration_handlers_test.go b/internal/api/migration_handlers_test.go new file mode 100644 index 0000000..0ce2a6a --- /dev/null +++ b/internal/api/migration_handlers_test.go @@ -0,0 +1,50 @@ +package api + +import ( + "bytes" + "encoding/json" + "net/http" + "net/http/httptest" + "testing" + + "github.com/flatrun/agent/internal/docker" + "github.com/flatrun/agent/pkg/models" + "github.com/gin-gonic/gin" +) + +func TestMigrationWorkflowPersistsThroughHTTP(t *testing.T) { + gin.SetMode(gin.TestMode) + dir := t.TempDir() + createTestDeployment(t, dir, "shop", &models.ServiceMetadata{Name: "shop"}) + server := &Server{manager: docker.NewManager(dir)} + router := gin.New() + router.PUT("/deployments/:name/migration", server.updateDeploymentMigration) + router.GET("/deployments/:name/migration", server.getDeploymentMigration) + + body := bytes.NewBufferString(`{"source":"legacy","inventory_complete":true,"expected_address":"192.0.2.10","sites":[{"hostname":"shop.example.com","source_path":"/srv/shop","bytes":42,"transferred":true}]}`) + req := httptest.NewRequest(http.MethodPut, "/deployments/shop/migration", body) + req.Header.Set("Content-Type", "application/json") + res := httptest.NewRecorder() + router.ServeHTTP(res, req) + if res.Code != http.StatusOK { + t.Fatalf("update returned %d: %s", res.Code, res.Body.String()) + } + + res = httptest.NewRecorder() + router.ServeHTTP(res, httptest.NewRequest(http.MethodGet, "/deployments/shop/migration", nil)) + if res.Code != http.StatusOK { + t.Fatalf("get returned %d: %s", res.Code, res.Body.String()) + } + var response struct { + Migration migrationStatus `json:"migration"` + } + if err := json.Unmarshal(res.Body.Bytes(), &response); err != nil { + t.Fatal(err) + } + if response.Migration.Plan == nil || response.Migration.Plan.Sites[0].Bytes != 42 { + t.Fatalf("migration was not persisted: %s", res.Body.String()) + } + if response.Migration.RetirementReady { + t.Fatal("migration became ready before final sync and cutover") + } +} diff --git a/internal/api/openapi.json b/internal/api/openapi.json index 60d3edc..0bb172c 100644 --- a/internal/api/openapi.json +++ b/internal/api/openapi.json @@ -4289,6 +4289,85 @@ "x-permission": "deployments:write" } }, + "/api/deployments/{name}/backup-cleanup": { + "post": { + "operationId": "post-deployments-by-name-backup-cleanup", + "parameters": [ + { + "in": "path", + "name": "name", + "required": true, + "schema": { + "type": "string" + } + } + ], + "requestBody": { + "content": { + "application/json": { + "schema": { + "properties": { + "keep": { + "type": "integer" + } + }, + "required": [ + "keep" + ], + "type": "object", + "x-columns": [ + "keep" + ], + "x-property-order": [ + "keep" + ] + } + } + }, + "required": true + }, + "responses": { + "200": { + "description": "Success" + } + }, + "tags": [ + "deployments" + ], + "x-permission": "backups:delete" + } + }, + "/api/deployments/{name}/backup-cleanup-preview": { + "get": { + "operationId": "get-deployments-by-name-backup-cleanup-preview", + "parameters": [ + { + "in": "path", + "name": "name", + "required": true, + "schema": { + "type": "string" + } + }, + { + "in": "query", + "name": "keep", + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "Success" + } + }, + "tags": [ + "deployments" + ], + "x-permission": "backups:read" + } + }, "/api/deployments/{name}/backup-config": { "get": { "operationId": "get-deployments-by-name-backup-config", @@ -4345,6 +4424,54 @@ "x-permission": "backups:write" } }, + "/api/deployments/{name}/backup-destinations": { + "get": { + "operationId": "get-deployments-by-name-backup-destinations", + "parameters": [ + { + "in": "path", + "name": "name", + "required": true, + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "Success" + } + }, + "tags": [ + "deployments" + ], + "x-permission": "backups:read" + } + }, + "/api/deployments/{name}/backup-policy": { + "get": { + "operationId": "get-deployments-by-name-backup-policy", + "parameters": [ + { + "in": "path", + "name": "name", + "required": true, + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "Success" + } + }, + "tags": [ + "deployments" + ], + "x-permission": "backups:read" + } + }, "/api/deployments/{name}/backups": { "get": { "operationId": "get-deployments-by-name-backups", @@ -4894,6 +5021,40 @@ "x-permission": "deployments:write" } }, + "/api/deployments/{name}/databases/attach": { + "post": { + "operationId": "post-deployments-by-name-databases-attach", + "parameters": [ + { + "in": "path", + "name": "name", + "required": true, + "schema": { + "type": "string" + } + } + ], + "requestBody": { + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/api.DatabaseConfigRequest" + } + } + }, + "required": true + }, + "responses": { + "200": { + "description": "Success" + } + }, + "tags": [ + "deployments" + ], + "x-permission": "deployments:write" + } + }, "/api/deployments/{name}/deploy": { "post": { "operationId": "post-deployments-by-name-deploy", @@ -5815,6 +5976,86 @@ "x-permission": "deployments:write" } }, + "/api/deployments/{name}/migration": { + "get": { + "operationId": "get-deployments-by-name-migration", + "parameters": [ + { + "in": "path", + "name": "name", + "required": true, + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "Success" + } + }, + "tags": [ + "deployments" + ], + "x-permission": "deployments:read" + }, + "put": { + "operationId": "put-deployments-by-name-migration", + "parameters": [ + { + "in": "path", + "name": "name", + "required": true, + "schema": { + "type": "string" + } + } + ], + "requestBody": { + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/models.MigrationSpec" + } + } + }, + "required": true + }, + "responses": { + "200": { + "description": "Success" + } + }, + "tags": [ + "deployments" + ], + "x-permission": "deployments:write" + } + }, + "/api/deployments/{name}/migration/check-dns": { + "post": { + "operationId": "post-deployments-by-name-migration-check-dns", + "parameters": [ + { + "in": "path", + "name": "name", + "required": true, + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "Success" + } + }, + "tags": [ + "deployments" + ], + "x-permission": "deployments:write" + } + }, "/api/deployments/{name}/mkdir/{path}": { "post": { "operationId": "post-deployments-by-name-mkdir-by-path", @@ -12820,6 +13061,9 @@ "backup.Backup": { "type": "object", "properties": { + "checksum": { + "type": "string" + }, "cleanup_results": { "type": "array", "items": { @@ -12855,6 +13099,12 @@ "$ref": "#/components/schemas/backup.DestinationResult" } }, + "destinations": { + "type": "array", + "items": { + "type": "string" + } + }, "error": { "type": "string" }, @@ -12886,6 +13136,7 @@ "deployment_name", "status", "size", + "checksum", "path", "components", "error", @@ -12895,13 +13146,15 @@ "locations", "component_results", "cleanup_results", - "destination_results" + "destination_results", + "destinations" ], "x-columns": [ "id", "deployment_name", "status", "size", + "checksum", "error", "created_at", "completed_at", @@ -12967,6 +13220,9 @@ "backup.DestinationResult": { "type": "object", "properties": { + "checksum": { + "type": "string" + }, "error": { "type": "string" }, @@ -12975,17 +13231,24 @@ }, "status": { "type": "string" + }, + "verified": { + "type": "boolean" } }, "x-property-order": [ "name", "status", - "error" + "error", + "checksum", + "verified" ], "x-columns": [ "name", "status", - "error" + "error", + "checksum", + "verified" ] }, "backup.RestoreBackupRequest": { @@ -12997,6 +13260,9 @@ "deployment_name": { "type": "string" }, + "isolated": { + "type": "boolean" + }, "restore_data": { "type": "boolean" }, @@ -13010,6 +13276,7 @@ "x-property-order": [ "backup_id", "deployment_name", + "isolated", "restore_data", "restore_db", "stop_first" @@ -13017,12 +13284,10 @@ "x-columns": [ "backup_id", "deployment_name", + "isolated", "restore_data", "restore_db", "stop_first" - ], - "required": [ - "backup_id" ] }, "cluster.AggregatedResponse": { @@ -13834,6 +14099,12 @@ "$ref": "#/components/schemas/models.DatabaseBackupSpec" } }, + "destinations": { + "type": "array", + "items": { + "type": "string" + } + }, "exclude_patterns": { "type": "array", "items": { @@ -13851,6 +14122,12 @@ "items": { "$ref": "#/components/schemas/models.BackupHookSpec" } + }, + "retention_count": { + "type": "integer" + }, + "size_alert_bytes": { + "type": "integer" } }, "x-property-order": [ @@ -13858,7 +14135,14 @@ "databases", "pre_hooks", "post_hooks", - "exclude_patterns" + "exclude_patterns", + "destinations", + "retention_count", + "size_alert_bytes" + ], + "x-columns": [ + "retention_count", + "size_alert_bytes" ] }, "models.Certificate": { @@ -14402,6 +14686,104 @@ "builtin" ] }, + "models.MigrationSite": { + "type": "object", + "properties": { + "bytes": { + "type": "integer" + }, + "dns_propagated": { + "type": "boolean" + }, + "hostname": { + "type": "string" + }, + "last_synced_at": { + "type": "string", + "format": "date-time" + }, + "resolved": { + "type": "array" + }, + "source_path": { + "type": "string" + }, + "transferred": { + "type": "boolean" + } + }, + "x-property-order": [ + "hostname", + "source_path", + "bytes", + "transferred", + "last_synced_at", + "resolved", + "dns_propagated" + ], + "x-columns": [ + "hostname", + "source_path", + "bytes", + "transferred", + "last_synced_at", + "dns_propagated" + ] + }, + "models.MigrationSpec": { + "type": "object", + "properties": { + "cutover_at": { + "type": "string", + "format": "date-time" + }, + "expected_address": { + "type": "string" + }, + "initial_transfer_at": { + "type": "string", + "format": "date-time" + }, + "inventory_complete": { + "type": "boolean" + }, + "last_sync_at": { + "type": "string", + "format": "date-time" + }, + "notes": { + "type": "string" + }, + "sites": { + "type": "array", + "items": { + "$ref": "#/components/schemas/models.MigrationSite" + } + }, + "source": { + "type": "string" + } + }, + "x-property-order": [ + "source", + "sites", + "inventory_complete", + "initial_transfer_at", + "last_sync_at", + "cutover_at", + "expected_address", + "notes" + ], + "x-columns": [ + "source", + "inventory_complete", + "initial_transfer_at", + "last_sync_at", + "cutover_at", + "expected_address", + "notes" + ] + }, "models.NetworkingConfig": { "type": "object", "properties": { @@ -14730,6 +15112,9 @@ "$ref": "#/components/schemas/models.LogSource" } }, + "migration": { + "$ref": "#/components/schemas/models.MigrationSpec" + }, "name": { "type": "string" }, @@ -14783,6 +15168,7 @@ "quick_actions", "security", "backup", + "migration", "scaling", "protected_mode", "require_plan", diff --git a/internal/api/server.go b/internal/api/server.go index aeaa369..bb24260 100644 --- a/internal/api/server.go +++ b/internal/api/server.go @@ -524,6 +524,7 @@ func (s *Server) setupRoutes() { protected.DELETE("/deployments/:name/logs", s.authMiddleware.RequirePermission(auth.PermDeploymentsWrite), s.authMiddleware.RequireDeploymentAccess(auth.AccessLevelWrite), s.deleteDeploymentLogs) protected.GET("/deployments/:name/log-sources", s.authMiddleware.RequirePermission(auth.PermDeploymentsRead), s.authMiddleware.RequireDeploymentAccess(auth.AccessLevelRead), s.getDeploymentLogSources) protected.PUT("/deployments/:name/log-sources", s.authMiddleware.RequirePermission(auth.PermDeploymentsWrite), s.authMiddleware.RequireDeploymentAccess(auth.AccessLevelWrite), s.updateDeploymentLogSources) + protected.POST("/deployments/:name/databases/attach", s.authMiddleware.RequirePermission(auth.PermDeploymentsWrite), s.authMiddleware.RequireDeploymentAccess(auth.AccessLevelWrite), s.attachDeploymentDatabase) protected.GET("/deployments/:name/compose", s.authMiddleware.RequirePermission(auth.PermDeploymentsRead), s.authMiddleware.RequireDeploymentAccess(auth.AccessLevelRead), s.getDeploymentCompose) protected.POST("/deployments/:name/compose/mount", s.authMiddleware.RequirePermission(auth.PermDeploymentsWrite), s.authMiddleware.RequireDeploymentAccess(auth.AccessLevelWrite), s.addDeploymentComposeMount) protected.POST("/deployments/:name/compose/unmount", s.authMiddleware.RequirePermission(auth.PermDeploymentsWrite), s.authMiddleware.RequireDeploymentAccess(auth.AccessLevelWrite), s.removeDeploymentComposeMount) @@ -836,6 +837,13 @@ func (s *Server) setupRoutes() { 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) + protected.GET("/deployments/:name/backup-destinations", s.authMiddleware.RequirePermission(auth.PermBackupsRead), s.authMiddleware.RequireDeploymentAccess(auth.AccessLevelRead), s.listDeploymentBackupDestinationOptions) + protected.GET("/deployments/:name/backup-policy", s.authMiddleware.RequirePermission(auth.PermBackupsRead), s.authMiddleware.RequireDeploymentAccess(auth.AccessLevelRead), s.getDeploymentBackupPolicy) + protected.GET("/deployments/:name/backup-cleanup-preview", s.authMiddleware.RequirePermission(auth.PermBackupsRead), s.authMiddleware.RequireDeploymentAccess(auth.AccessLevelRead), s.previewDeploymentBackupCleanup) + protected.POST("/deployments/:name/backup-cleanup", s.authMiddleware.RequirePermission(auth.PermBackupsDelete), s.authMiddleware.RequireDeploymentAccess(auth.AccessLevelAdmin), s.cleanupDeploymentBackups) + protected.GET("/deployments/:name/migration", s.authMiddleware.RequirePermission(auth.PermDeploymentsRead), s.authMiddleware.RequireDeploymentAccess(auth.AccessLevelRead), s.getDeploymentMigration) + protected.PUT("/deployments/:name/migration", s.authMiddleware.RequirePermission(auth.PermDeploymentsWrite), s.authMiddleware.RequireDeploymentAccess(auth.AccessLevelWrite), s.updateDeploymentMigration) + protected.POST("/deployments/:name/migration/check-dns", s.authMiddleware.RequirePermission(auth.PermDeploymentsWrite), s.authMiddleware.RequireDeploymentAccess(auth.AccessLevelWrite), s.checkDeploymentMigrationDNS) protected.POST("/backups/:id/restore", s.authMiddleware.RequirePermission(auth.PermBackupsWrite), s.restoreBackup) protected.GET("/backups/jobs", s.authMiddleware.RequirePermission(auth.PermBackupsRead), s.listBackupJobs) protected.GET("/backups/jobs/:id", s.authMiddleware.RequirePermission(auth.PermBackupsRead), s.getBackupJob) @@ -6365,12 +6373,15 @@ func (s *Server) updateDomain(c *gin.Context) { return } - result, err := s.proxyOrchestrator.SetupDeployment(deployment) - if err != nil { - _ = s.manager.SaveMetadata(name, originalMetadata) - _ = s.manager.UpdateDeployment(name, originalCompose) - c.JSON(http.StatusConflict, gin.H{"error": "Failed to configure proxy: " + err.Error()}) - return + var result *proxy.SetupResult + if s.proxyOrchestrator != nil { + result, err = s.proxyOrchestrator.SetupDeployment(deployment) + if err != nil { + _ = s.manager.SaveMetadata(name, originalMetadata) + _ = s.manager.UpdateDeployment(name, originalCompose) + c.JSON(http.StatusConflict, gin.H{"error": "Failed to configure proxy: " + err.Error()}) + return + } } c.JSON(http.StatusOK, gin.H{ diff --git a/internal/backup/isolation_test.go b/internal/backup/isolation_test.go new file mode 100644 index 0000000..cb26e80 --- /dev/null +++ b/internal/backup/isolation_test.go @@ -0,0 +1,67 @@ +package backup + +import ( + "os" + "path/filepath" + "strings" + "testing" +) + +func TestIsolateRestoredDeploymentRemovesExternalConnectivity(t *testing.T) { + dir := t.TempDir() + compose := `name: production-project +services: + app: + image: example/app + container_name: production-app + network_mode: host + ports: + - "8080:80" + extra_hosts: + - "host.docker.internal:host-gateway" + volumes: + - /srv/production:/data + - production-data:/cache + db: + image: postgres:16 +volumes: + production-data: + external: true + name: production-data +networks: + production: + external: true +` + if err := os.WriteFile(filepath.Join(dir, "docker-compose.yml"), []byte(compose), 0644); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(dir, ".env"), []byte("COMPOSE_PROJECT_NAME=production\nAPP_ENV=staging\n"), 0600); err != nil { + t.Fatal(err) + } + if err := isolateRestoredDeployment(dir); err != nil { + t.Fatal(err) + } + content, err := os.ReadFile(filepath.Join(dir, "docker-compose.yml")) + if err != nil { + t.Fatal(err) + } + got := string(content) + for _, forbidden := range []string{"container_name:", "network_mode:", "ports:", "extra_hosts:", "external:", "/srv/production", "name: production"} { + if strings.Contains(got, forbidden) { + t.Fatalf("isolated compose still contains %q:\n%s", forbidden, got) + } + } + if !strings.Contains(got, "internal: true") || strings.Count(got, "flatrun_isolated") < 3 { + t.Fatalf("isolated network is missing:\n%s", got) + } + if !strings.Contains(got, "isolated-mounts/app-0") { + t.Fatalf("absolute mount was not isolated:\n%s", got) + } + env, err := os.ReadFile(filepath.Join(dir, ".env")) + if err != nil { + t.Fatal(err) + } + if strings.Contains(string(env), "COMPOSE_PROJECT_NAME") || !strings.Contains(string(env), "APP_ENV=staging") { + t.Fatalf("isolated environment is invalid: %s", env) + } +} diff --git a/internal/backup/jobs.go b/internal/backup/jobs.go index 3013b49..d379b09 100644 --- a/internal/backup/jobs.go +++ b/internal/backup/jobs.go @@ -173,7 +173,7 @@ func (m *Manager) StartBackupJob(deploymentName string, spec *BackupSpec) string return jobID } -func (m *Manager) StartRestoreJob(req *RestoreBackupRequest) string { +func (m *Manager) StartRestoreJob(req *RestoreBackupRequest, afterRestore ...func() error) string { backup, err := m.GetBackup(req.BackupID) if err != nil { jobID := generateJobID("restore", req.BackupID) @@ -198,6 +198,12 @@ func (m *Manager) StartRestoreJob(req *RestoreBackupRequest) string { m.jobs.SetError(jobID, err) return } + for _, callback := range afterRestore { + if err := callback(); err != nil { + m.jobs.SetError(jobID, err) + return + } + } m.jobs.UpdateStatus(jobID, JobStatusCompleted, "Restore completed") }() diff --git a/internal/backup/manager.go b/internal/backup/manager.go index 393fdb2..52c0f92 100644 --- a/internal/backup/manager.go +++ b/internal/backup/manager.go @@ -4,6 +4,7 @@ import ( "archive/tar" "compress/gzip" "context" + "crypto/sha256" "encoding/json" "errors" "fmt" @@ -18,6 +19,7 @@ import ( "time" "github.com/flatrun/agent/pkg/version" + "gopkg.in/yaml.v3" ) type Manager struct { @@ -61,6 +63,9 @@ func (m *Manager) CreateBackup(ctx context.Context, deploymentName string, spec CreatedAt: time.Now(), Components: []string{}, } + if spec != nil { + backup.Destinations = append([]string(nil), spec.Destinations...) + } record := func(kind, name string, required bool, err error) { result := ComponentResult{Name: name, Kind: kind, Required: required, Status: ResultStatusCompleted} if err != nil { @@ -174,6 +179,10 @@ func (m *Manager) CreateBackup(ctx context.Context, deploymentName string, spec if info, statErr := os.Stat(archivePath); statErr == nil { backup.Size = info.Size() } + backup.Checksum, err = checksumFile(archivePath) + if err != nil { + captureErrors = append(captureErrors, fmt.Errorf("failed to checksum backup archive: %w", err)) + } } } @@ -197,7 +206,7 @@ func (m *Manager) CreateBackup(ctx context.Context, deploymentName string, spec return backup, errors.Join(captureErrors...) } - backup.DestinationResults = m.mirrorToRemotes(ctx, deploymentName, backupID, archivePath, backup.Size) + backup.DestinationResults = m.mirrorToRemotes(ctx, deploymentName, backupID, archivePath, backup.Size, backup.Checksum, backup.Destinations) succeeded := 0 for _, result := range backup.DestinationResults { if result.Status == ResultStatusCompleted { @@ -221,6 +230,19 @@ func (m *Manager) CreateBackup(ctx context.Context, deploymentName string, spec return backup, nil } +func checksumFile(path string) (string, error) { + f, err := os.Open(path) + if err != nil { + return "", err + } + defer f.Close() + hash := sha256.New() + if _, err := io.Copy(hash, f); err != nil { + return "", err + } + return fmt.Sprintf("sha256:%x", hash.Sum(nil)), nil +} + func (m *Manager) backupComposeFile(deploymentPath, tempDir string, metadata *BackupMetadata) error { composePath := filepath.Join(deploymentPath, "docker-compose.yml") if _, err := os.Stat(composePath); err != nil { @@ -789,10 +811,29 @@ func (m *Manager) GetBackupPath(backupID string) (string, error) { // applies to local disk only; remote copies are governed by the destination's // own lifecycle policy and are never deleted here. func (m *Manager) CleanupOldBackups(deploymentName string, keepCount int) (int, error) { - backups, err := m.listLocalBackups(&BackupListFilter{DeploymentName: deploymentName}) + preview, err := m.PreviewCleanup(deploymentName, keepCount) if err != nil { return 0, err } + deleted := 0 + for _, backupID := range preview.DeleteIDs { + if err := m.deleteLocalBackup(backupID); err != nil { + log.Printf("Failed to delete old backup %s: %v", backupID, err) + continue + } + deleted++ + } + return deleted, nil +} + +func (m *Manager) PreviewCleanup(deploymentName string, keepCount int) (*CleanupPreview, error) { + if keepCount < 1 { + return nil, errors.New("retention count must be at least 1") + } + backups, err := m.listLocalBackups(&BackupListFilter{DeploymentName: deploymentName}) + if err != nil { + return nil, err + } usable := backups[:0] for _, backup := range backups { @@ -801,19 +842,14 @@ func (m *Manager) CleanupOldBackups(deploymentName string, keepCount int) (int, } } if len(usable) <= keepCount { - return 0, nil + return &CleanupPreview{KeepCount: keepCount, DeleteIDs: []string{}}, nil } - - deleted := 0 + preview := &CleanupPreview{KeepCount: keepCount, DeleteIDs: make([]string, 0, len(usable)-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 - } - deleted++ + preview.DeleteIDs = append(preview.DeleteIDs, backup.ID) + preview.ReclaimedBytes += backup.Size } - - return deleted, nil + return preview, nil } func (m *Manager) RestoreBackup(ctx context.Context, req *RestoreBackupRequest) error { @@ -826,6 +862,16 @@ func (m *Manager) RestoreBackup(ctx context.Context, req *RestoreBackupRequest) if req.DeploymentName != "" { deploymentName = req.DeploymentName } + if req.Isolated { + if req.DeploymentName == "" || deploymentName == backup.DeploymentName { + return errors.New("isolated restore requires a new deployment name") + } + if _, statErr := os.Stat(filepath.Join(m.deploymentsPath, deploymentName)); statErr == nil { + return fmt.Errorf("deployment already exists: %s", deploymentName) + } else if !os.IsNotExist(statErr) { + return fmt.Errorf("failed to inspect restore destination: %w", statErr) + } + } deploymentPath := filepath.Join(m.deploymentsPath, deploymentName) @@ -893,6 +939,12 @@ func (m *Manager) RestoreBackup(ctx context.Context, req *RestoreBackupRequest) } } + if req.Isolated { + if err := isolateRestoredDeployment(deploymentPath); err != nil { + return fmt.Errorf("failed to isolate restored deployment: %w", err) + } + } + if req.RestoreData && len(metadata.Components.MountedData) > 0 { if err := m.restoreMountedData(tempDir, deploymentPath, metadata.Components.MountedData); err != nil { log.Printf("Restore: warning - failed to restore mounted data: %v", err) @@ -933,6 +985,122 @@ func (m *Manager) RestoreBackup(ctx context.Context, req *RestoreBackupRequest) return nil } +func isolateRestoredDeployment(deploymentPath string) error { + composePath := filepath.Join(deploymentPath, "docker-compose.yml") + content, err := os.ReadFile(composePath) + if err != nil { + return err + } + var document map[string]interface{} + if err := yaml.Unmarshal(content, &document); err != nil { + return err + } + delete(document, "name") + services, ok := document["services"].(map[string]interface{}) + if !ok || len(services) == 0 { + return errors.New("compose file has no services") + } + for name, raw := range services { + service, ok := raw.(map[string]interface{}) + if !ok { + return fmt.Errorf("service %s is invalid", name) + } + delete(service, "container_name") + delete(service, "network_mode") + delete(service, "ports") + delete(service, "dns") + delete(service, "dns_search") + delete(service, "extra_hosts") + if volumes, exists := service["volumes"].([]interface{}); exists { + for i, rawVolume := range volumes { + isolated, err := isolateComposeVolume(deploymentPath, name, i, rawVolume) + if err != nil { + return err + } + volumes[i] = isolated + } + service["volumes"] = volumes + } + service["networks"] = []string{"flatrun_isolated"} + services[name] = service + } + if volumes, exists := document["volumes"].(map[string]interface{}); exists { + for name, raw := range volumes { + definition, ok := raw.(map[string]interface{}) + if !ok { + definition = map[string]interface{}{} + } + delete(definition, "external") + delete(definition, "name") + volumes[name] = definition + } + } + document["networks"] = map[string]interface{}{ + "flatrun_isolated": map[string]interface{}{"internal": true}, + } + isolated, err := yaml.Marshal(document) + if err != nil { + return err + } + if err := os.WriteFile(composePath, isolated, 0644); err != nil { + return err + } + for _, filename := range []string{".env", ".env.flatrun"} { + path := filepath.Join(deploymentPath, filename) + content, err := os.ReadFile(path) + if os.IsNotExist(err) { + continue + } + if err != nil { + return err + } + lines := strings.Split(string(content), "\n") + filtered := lines[:0] + for _, line := range lines { + if !strings.HasPrefix(strings.TrimSpace(line), "COMPOSE_PROJECT_NAME=") { + filtered = append(filtered, line) + } + } + if err := os.WriteFile(path, []byte(strings.Join(filtered, "\n")), 0600); err != nil { + return err + } + } + return nil +} + +func isolateComposeVolume(deploymentPath, service string, index int, raw interface{}) (interface{}, error) { + isolatedSource := filepath.ToSlash(filepath.Join(".", "isolated-mounts", fmt.Sprintf("%s-%d", service, index))) + createSource := func() error { + return os.MkdirAll(filepath.Join(deploymentPath, "isolated-mounts", fmt.Sprintf("%s-%d", service, index)), 0755) + } + switch volume := raw.(type) { + case string: + parts := strings.SplitN(volume, ":", 3) + if len(parts) < 2 || !filepath.IsAbs(parts[0]) { + return raw, nil + } + if err := createSource(); err != nil { + return nil, err + } + parts[0] = isolatedSource + return strings.Join(parts, ":"), nil + case map[string]interface{}: + source, _ := volume["source"].(string) + typeName, _ := volume["type"].(string) + if typeName != "bind" && !filepath.IsAbs(source) { + return raw, nil + } + if err := createSource(); err != nil { + return nil, err + } + volume["type"] = "bind" + volume["source"] = isolatedSource + return volume, nil + default: + return nil, fmt.Errorf("service %s has an invalid volume entry", service) + } +} + func (m *Manager) extractArchive(archivePath, destDir string) error { file, err := os.Open(archivePath) if err != nil { diff --git a/internal/backup/manager_test.go b/internal/backup/manager_test.go index f7ebfd4..1c2b412 100644 --- a/internal/backup/manager_test.go +++ b/internal/backup/manager_test.go @@ -362,6 +362,29 @@ func TestCleanupOldBackups(t *testing.T) { } } +func TestPreviewCleanupDoesNotDeleteBackups(t *testing.T) { + m, deploymentPath := setupTestManager(t) + defer os.RemoveAll(deploymentPath) + seedDeployment(t, deploymentPath, "test-deployment") + for i := 0; i < 3; i++ { + if _, err := m.CreateBackup(context.Background(), "test-deployment", nil); err != nil { + t.Fatal(err) + } + time.Sleep(time.Millisecond) + } + preview, err := m.PreviewCleanup("test-deployment", 1) + if err != nil { + t.Fatal(err) + } + if len(preview.DeleteIDs) != 2 || preview.ReclaimedBytes == 0 { + t.Fatalf("unexpected preview: %#v", preview) + } + backups, err := m.ListBackups(&BackupListFilter{DeploymentName: "test-deployment"}) + if err != nil || len(backups) != 3 { + t.Fatalf("preview changed backups under %s: count=%d err=%v", deploymentPath, len(backups), err) + } +} + func TestCleanupOldBackups_DoesNotCountFailedAttempts(t *testing.T) { m, tmpDir := setupTestManager(t) defer os.RemoveAll(tmpDir) diff --git a/internal/backup/remotes.go b/internal/backup/remotes.go index df7946f..54a58a8 100644 --- a/internal/backup/remotes.go +++ b/internal/backup/remotes.go @@ -2,6 +2,7 @@ package backup import ( "context" + "crypto/sha256" "errors" "fmt" "io" @@ -47,14 +48,33 @@ func deploymentFromID(backupID string) string { return parts[0] } -func (m *Manager) mirrorToRemotes(ctx context.Context, deploymentName, backupID, archivePath string, size int64) []DestinationResult { +func (m *Manager) mirrorToRemotes(ctx context.Context, deploymentName, backupID, archivePath string, size int64, checksum string, destinationNames []string) []DestinationResult { remotes := m.getRemotes() - if len(remotes) == 0 { + if len(remotes) == 0 && len(destinationNames) == 0 { return nil } + var results []DestinationResult + if len(destinationNames) > 0 { + byName := make(map[string]Store, len(remotes)) + for _, remote := range remotes { + byName[remote.Name()] = remote + } + selected := make([]Store, 0, len(destinationNames)) + for _, name := range destinationNames { + remote, ok := byName[name] + if !ok { + results = append(results, DestinationResult{Name: name, Status: ResultStatusFailed, Error: "destination unavailable"}) + continue + } + selected = append(selected, remote) + } + remotes = selected + if len(remotes) == 0 { + return results + } + } key := backupKey(deploymentName, backupID) - var results []DestinationResult for _, r := range remotes { result := DestinationResult{Name: r.Name(), Status: ResultStatusCompleted} f, err := os.Open(archivePath) @@ -74,6 +94,24 @@ func (m *Manager) mirrorToRemotes(ctx context.Context, deploymentName, backupID, results = append(results, result) continue } + remote, err := r.Open(ctx, key) + if err != nil { + result.Status = ResultStatusFailed + result.Error = "verification failed" + results = append(results, result) + continue + } + hash := sha256.New() + _, hashErr := io.Copy(hash, remote) + closeErr := remote.Close() + result.Checksum = fmt.Sprintf("sha256:%x", hash.Sum(nil)) + if hashErr != nil || closeErr != nil || result.Checksum != checksum { + result.Status = ResultStatusFailed + result.Error = "verification failed" + results = append(results, result) + continue + } + result.Verified = true results = append(results, result) log.Printf("Backup mirrored: %s -> %s", backupID, r.Name()) } @@ -88,7 +126,13 @@ func (m *Manager) RetryRemotePublication(ctx context.Context, backupID string) ( 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) + if backup.Checksum == "" { + backup.Checksum, err = checksumFile(backup.Path) + if err != nil { + return nil, err + } + } + backup.DestinationResults = m.mirrorToRemotes(ctx, backup.DeploymentName, backup.ID, backup.Path, backup.Size, backup.Checksum, backup.Destinations) backup.Locations = []string{locationLocal} succeeded := 0 for _, result := range backup.DestinationResults { diff --git a/internal/backup/remotes_test.go b/internal/backup/remotes_test.go index a043cad..b906e89 100644 --- a/internal/backup/remotes_test.go +++ b/internal/backup/remotes_test.go @@ -192,6 +192,40 @@ func TestCreateBackup_MirrorsToRemote(t *testing.T) { } } +func TestCreateBackup_MirrorsOnlyToSelectedDestinations(t *testing.T) { + m, tmpDir := setupTestManager(t) + defer os.RemoveAll(tmpDir) + + primary := newFakeStore("primary") + secondary := newFakeStore("secondary") + m.SetRemotes([]Store{primary, secondary}) + seedDeployment(t, tmpDir, "app") + + b, err := m.CreateBackup(context.Background(), "app", &BackupSpec{Destinations: []string{"secondary"}}) + if err != nil { + t.Fatalf("create backup: %v", err) + } + key := backupKey("app", b.ID) + if primary.has(key) || !secondary.has(key) { + t.Fatalf("primary has backup = %v, secondary has backup = %v", primary.has(key), secondary.has(key)) + } + if len(b.DestinationResults) != 1 || b.DestinationResults[0].Name != "secondary" { + t.Fatalf("destination results = %#v", b.DestinationResults) + } + if b.Checksum == "" || !b.DestinationResults[0].Verified || b.DestinationResults[0].Checksum != b.Checksum { + t.Fatalf("backup verification = %#v", b) + } + + secondary.failPut = true + retried, err := m.RetryRemotePublication(context.Background(), b.ID) + if err != nil { + t.Fatalf("retry publication: %v", err) + } + if len(retried.DestinationResults) != 1 || retried.DestinationResults[0].Name != "secondary" { + t.Fatalf("retry destination results = %#v", retried.DestinationResults) + } +} + func TestListAndGetBackup_RemoteOnlyAfterLocalPruned(t *testing.T) { m, tmpDir := setupTestManager(t) defer os.RemoveAll(tmpDir) diff --git a/internal/backup/types.go b/internal/backup/types.go index f0708bd..98af580 100644 --- a/internal/backup/types.go +++ b/internal/backup/types.go @@ -34,9 +34,11 @@ type ComponentResult struct { } type DestinationResult struct { - Name string `json:"name"` - Status ResultStatus `json:"status"` - Error string `json:"error,omitempty"` + Name string `json:"name"` + Status ResultStatus `json:"status"` + Error string `json:"error,omitempty"` + Checksum string `json:"checksum,omitempty"` + Verified bool `json:"verified"` } type BackupSpec = models.BackupSpec @@ -49,6 +51,7 @@ type Backup struct { DeploymentName string `json:"deployment_name"` Status BackupStatus `json:"status"` Size int64 `json:"size"` + Checksum string `json:"checksum,omitempty"` Path string `json:"path" cli:"-"` Components []string `json:"components"` Error string `json:"error,omitempty"` @@ -62,6 +65,7 @@ type Backup struct { ComponentResults []ComponentResult `json:"component_results,omitempty"` CleanupResults []ComponentResult `json:"cleanup_results,omitempty"` DestinationResults []DestinationResult `json:"destination_results,omitempty"` + Destinations []string `json:"destinations,omitempty"` } type BackupMetadata struct { @@ -90,8 +94,9 @@ type CreateBackupRequest struct { } type RestoreBackupRequest struct { - BackupID string `json:"backup_id" binding:"required"` + BackupID string `json:"backup_id,omitempty"` DeploymentName string `json:"deployment_name,omitempty"` + Isolated bool `json:"isolated"` RestoreData bool `json:"restore_data"` RestoreDB bool `json:"restore_db"` StopFirst bool `json:"stop_first"` @@ -103,3 +108,9 @@ type BackupListFilter struct { Limit int Offset int } + +type CleanupPreview struct { + KeepCount int `json:"keep_count"` + DeleteIDs []string `json:"delete_ids"` + ReclaimedBytes int64 `json:"reclaimed_bytes"` +} diff --git a/internal/nginx/manager.go b/internal/nginx/manager.go index 0dfcb6c..4fb163c 100644 --- a/internal/nginx/manager.go +++ b/internal/nginx/manager.go @@ -1104,6 +1104,12 @@ server { ssl_certificate /etc/letsencrypt/live/{{.Domain}}/fullchain.pem; ssl_certificate_key /etc/letsencrypt/live/{{.Domain}}/privkey.pem; + location /.well-known/acme-challenge/ { + root {{.ContainerWebrootPath}}; + access_log off; + log_not_found off; + } + ssl_protocols TLSv1.2 TLSv1.3; ssl_ciphers ECDHE-ECDSA-AES128-GCM-SHA256:ECDHE-RSA-AES128-GCM-SHA256:ECDHE-ECDSA-AES256-GCM-SHA384:ECDHE-RSA-AES256-GCM-SHA384; ssl_prefer_server_ciphers off; @@ -1306,6 +1312,12 @@ server { ssl_certificate /etc/letsencrypt/live/{{.SSLDomain}}/fullchain.pem; ssl_certificate_key /etc/letsencrypt/live/{{.SSLDomain}}/privkey.pem; + location /.well-known/acme-challenge/ { + root {{$.ContainerWebrootPath}}; + access_log off; + log_not_found off; + } + ssl_protocols TLSv1.2 TLSv1.3; ssl_ciphers ECDHE-ECDSA-AES128-GCM-SHA256:ECDHE-RSA-AES128-GCM-SHA256:ECDHE-ECDSA-AES256-GCM-SHA384:ECDHE-RSA-AES256-GCM-SHA384; ssl_prefer_server_ciphers off; diff --git a/internal/nginx/manager_test.go b/internal/nginx/manager_test.go index cb36360..1ffbc16 100644 --- a/internal/nginx/manager_test.go +++ b/internal/nginx/manager_test.go @@ -234,6 +234,13 @@ func TestGenerateConfig_ContainerWebrootPath(t *testing.T) { if strings.Contains(configContent, hostWebrootLine) { t.Errorf("config should not contain host webroot path %q", hostWebrootLine) } + if tt.sslEnabled { + challengeLocations := strings.Count(configContent, "location /.well-known/acme-challenge/") + serverBlocks := strings.Count(configContent, "server {") + if challengeLocations != serverBlocks { + t.Errorf("challenge locations = %d, server blocks = %d\nConfig:\n%s", challengeLocations, serverBlocks, configContent) + } + } }) } } @@ -1923,6 +1930,11 @@ func TestGenerateMultiDomainConfig_SSLToggle(t *testing.T) { if !strings.Contains(config, "ssl_certificate") { t.Error("all-SSL multi-domain config must contain ssl_certificate directive") } + challengeLocations := strings.Count(config, "location /.well-known/acme-challenge/") + serverBlocks := strings.Count(config, "server {") + if challengeLocations != serverBlocks { + t.Errorf("challenge locations = %d, server blocks = %d\nConfig:\n%s", challengeLocations, serverBlocks, config) + } }) t.Run("all domains SSL disabled", func(t *testing.T) { diff --git a/internal/scheduler/executor.go b/internal/scheduler/executor.go index c8cea8e..d7e2813 100644 --- a/internal/scheduler/executor.go +++ b/internal/scheduler/executor.go @@ -57,8 +57,9 @@ func (e *Executor) ExecuteBackup(ctx context.Context, deploymentName string, con return "", err } - if config.RetentionCount > 0 { - deleted, err := e.backupManager.CleanupOldBackups(deploymentName, config.RetentionCount) + retentionCount := backupRetentionCount(config, spec) + if retentionCount > 0 { + deleted, err := e.backupManager.CleanupOldBackups(deploymentName, retentionCount) if err != nil { log.Printf("Scheduler: failed to cleanup old backups: %v", err) } else if deleted > 0 { @@ -69,6 +70,13 @@ func (e *Executor) ExecuteBackup(ctx context.Context, deploymentName string, con return fmt.Sprintf("Backup created: %s (%d bytes)", b.ID, b.Size), nil } +func backupRetentionCount(config *BackupTaskConfig, spec *backup.BackupSpec) int { + if spec != nil && spec.RetentionCount > 0 { + return spec.RetentionCount + } + return config.RetentionCount +} + func (e *Executor) ExecuteCommand(ctx context.Context, deploymentName string, config *CommandTaskConfig) (string, error) { if e.dockerManager == nil { return "", fmt.Errorf("docker manager not available") diff --git a/internal/scheduler/executor_test.go b/internal/scheduler/executor_test.go index 96736e5..e68fda4 100644 --- a/internal/scheduler/executor_test.go +++ b/internal/scheduler/executor_test.go @@ -7,6 +7,7 @@ import ( "strings" "testing" + "github.com/flatrun/agent/internal/backup" "github.com/flatrun/agent/internal/docker" ) @@ -187,3 +188,14 @@ func TestExecuteBackup_NilBackupManager(t *testing.T) { t.Errorf("Unexpected error: %v", err) } } + +func TestBackupRetentionCountPrefersDeploymentPolicy(t *testing.T) { + config := &BackupTaskConfig{RetentionCount: 3} + spec := &backup.BackupSpec{RetentionCount: 10} + if got := backupRetentionCount(config, spec); got != 10 { + t.Fatalf("retention count = %d, want 10", got) + } + if got := backupRetentionCount(config, nil); got != 3 { + t.Fatalf("fallback retention count = %d, want 3", got) + } +} diff --git a/pkg/models/deployment.go b/pkg/models/deployment.go index 1c9c665..7dd04ca 100644 --- a/pkg/models/deployment.go +++ b/pkg/models/deployment.go @@ -40,6 +40,7 @@ type ServiceMetadata struct { QuickActions []QuickAction `yaml:"quick_actions,omitempty" json:"quick_actions,omitempty"` Security *DeploymentSecurityConfig `yaml:"security,omitempty" json:"security,omitempty"` Backup *BackupSpec `yaml:"backup,omitempty" json:"backup,omitempty"` + Migration *MigrationSpec `yaml:"migration,omitempty" json:"migration,omitempty"` Scaling *ScalingConfig `yaml:"scaling,omitempty" json:"scaling,omitempty"` ProtectedMode *ProtectedModeConfig `yaml:"protected_mode,omitempty" json:"protected_mode,omitempty"` RequirePlan bool `yaml:"require_plan,omitempty" json:"require_plan,omitempty"` @@ -49,6 +50,27 @@ type ServiceMetadata struct { Databases []DatabaseConfig `yaml:"databases,omitempty" json:"databases,omitempty"` } +type MigrationSpec struct { + Source string `yaml:"source" json:"source"` + Sites []MigrationSite `yaml:"sites,omitempty" json:"sites,omitempty"` + InventoryComplete bool `yaml:"inventory_complete" json:"inventory_complete"` + InitialTransferAt *time.Time `yaml:"initial_transfer_at,omitempty" json:"initial_transfer_at,omitempty"` + LastSyncAt *time.Time `yaml:"last_sync_at,omitempty" json:"last_sync_at,omitempty"` + CutoverAt *time.Time `yaml:"cutover_at,omitempty" json:"cutover_at,omitempty"` + ExpectedAddress string `yaml:"expected_address,omitempty" json:"expected_address,omitempty"` + Notes string `yaml:"notes,omitempty" json:"notes,omitempty"` +} + +type MigrationSite struct { + Hostname string `yaml:"hostname" json:"hostname"` + SourcePath string `yaml:"source_path,omitempty" json:"source_path,omitempty"` + Bytes int64 `yaml:"bytes,omitempty" json:"bytes,omitempty"` + Transferred bool `yaml:"transferred" json:"transferred"` + LastSyncedAt *time.Time `yaml:"last_synced_at,omitempty" json:"last_synced_at,omitempty"` + Resolved []string `yaml:"resolved,omitempty" json:"resolved,omitempty"` + DNSPropagated bool `yaml:"dns_propagated,omitempty" json:"dns_propagated"` +} + type ScalingConfig struct { Service string `yaml:"service" json:"service"` Stateless bool `yaml:"stateless" json:"stateless"` @@ -265,6 +287,9 @@ type BackupSpec struct { PreHooks []BackupHookSpec `yaml:"pre_hooks,omitempty" json:"pre_hooks,omitempty"` PostHooks []BackupHookSpec `yaml:"post_hooks,omitempty" json:"post_hooks,omitempty"` ExcludePatterns []string `yaml:"exclude_patterns,omitempty" json:"exclude_patterns,omitempty"` + Destinations []string `yaml:"destinations,omitempty" json:"destinations,omitempty"` + RetentionCount int `yaml:"retention_count,omitempty" json:"retention_count,omitempty"` + SizeAlertBytes int64 `yaml:"size_alert_bytes,omitempty" json:"size_alert_bytes,omitempty"` } type ContainerBackupPath struct {