From 0b746e77022440952269af850d8b2a245c4e7ed7 Mon Sep 17 00:00:00 2001 From: vthwang Date: Wed, 23 Sep 2026 12:21:00 -0700 Subject: [PATCH] feat: support per-component TOML config updates Allow stack owners and platform admins to update one runtime config without restarting other components. - Add syntax checks, scoped writes, and rollback on failed readiness. - Detect repeated startup crashes and cover config jobs with tests. Signed-off-by: vthwang --- go.mod | 2 +- .../handler/admin_platform_stack_admins.go | 41 +- internal/handler/setup_config.go | 521 ++++++++++++++++++ internal/handler/setup_config_test.go | 87 +++ internal/handler/setup_fullstack.go | 11 + internal/handler/setup_vtc.go | 11 + internal/k8s/component_jobs.go | 29 +- internal/k8s/component_jobs_test.go | 30 + internal/k8s/component_readiness_test.go | 40 ++ internal/k8s/component_resources.go | 42 ++ internal/k8s/fullstack_names.go | 14 +- internal/k8s/vta_resources.go | 2 + internal/model/vta_acl.go | 2 +- internal/router/router.go | 6 + 14 files changed, 817 insertions(+), 21 deletions(-) create mode 100644 internal/handler/setup_config.go create mode 100644 internal/handler/setup_config_test.go diff --git a/go.mod b/go.mod index 8599acb..3278ed1 100644 --- a/go.mod +++ b/go.mod @@ -11,6 +11,7 @@ require ( github.com/golang-migrate/migrate/v4 v4.19.1 github.com/google/uuid v1.6.0 github.com/joho/godotenv v1.5.1 + github.com/pelletier/go-toml/v2 v2.3.1 golang.org/x/net v0.55.0 golang.org/x/sync v0.20.0 gorm.io/driver/postgres v1.6.0 @@ -71,7 +72,6 @@ require ( github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect github.com/modern-go/reflect2 v1.0.3-0.20250322232337-35a7c28c31ee // indirect github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect - github.com/pelletier/go-toml/v2 v2.3.1 // indirect github.com/philhofer/fwd v1.2.0 // indirect github.com/quic-go/qpack v0.6.0 // indirect github.com/quic-go/quic-go v0.59.1 // indirect diff --git a/internal/handler/admin_platform_stack_admins.go b/internal/handler/admin_platform_stack_admins.go index 28e354c..e5f3d50 100644 --- a/internal/handler/admin_platform_stack_admins.go +++ b/internal/handler/admin_platform_stack_admins.go @@ -90,7 +90,30 @@ var aclJobLocks = aclJobLockSet{held: make(map[uint]struct{})} // errAclJobBusy is a sentinel so the handler answers 409 (come back in a // minute) rather than 502 (something broke). Nothing is wrong when this fires. var errAclJobBusy = errors.New( - "another ACL operation is running for this VTA — that takes about a minute; try again after it finishes") + "another maintenance operation is running for this stack — try again after it finishes") + +// acquireSessionMaintenance serialises every operation that stops one or more +// components in a stack. The in-process lock is fast; the database timestamp +// closes the same race across API replicas. Callers must invoke the returned +// release function. +func (h *SetupHandler) acquireSessionMaintenance(sessionID uint) (func(), error) { + if !aclJobLocks.TryLock(sessionID) { + return nil, errAclJobBusy + } + startedAt, ok, err := h.acquireVtaAclMaintenance(sessionID) + if err != nil { + aclJobLocks.Unlock(sessionID) + return nil, fmt.Errorf("failed to lock stack maintenance: %w", err) + } + if !ok { + aclJobLocks.Unlock(sessionID) + return nil, errAclJobBusy + } + return func() { + h.releaseVtaAclMaintenance(sessionID, startedAt) + aclJobLocks.Unlock(sessionID) + }, nil +} // platformSession loads the platform stack's session, writing the response and // returning nil when there isn't one. Same two-step lookup GetPlatformStack @@ -183,21 +206,11 @@ func (h *SetupHandler) runVtaAclJob( // rather than Lock: a caller who waits would sit through the other window // and then start their own, so the honest answer is to refuse now and let // them retry once — a queue of these is a queue of outages. - if !aclJobLocks.TryLock(session.ID) { - return "", nil, errAclJobBusy - } - defer aclJobLocks.Unlock(session.ID) - - // The in-process lock above protects callers handled by this replica. This - // row closes the same race across replicas and doubles as snapshot metadata. - lockStartedAt, ok, lockErr := h.acquireVtaAclMaintenance(session.ID) + release, lockErr := h.acquireSessionMaintenance(session.ID) if lockErr != nil { - return "", nil, fmt.Errorf("failed to lock ACL maintenance: %w", lockErr) - } - if !ok { - return "", nil, errAclJobBusy + return "", nil, lockErr } - defer h.releaseVtaAclMaintenance(session.ID, lockStartedAt) + defer release() ns := h.k8s.UserNamespace(fmt.Sprintf("%d", session.UserID)) target := vtaAclTargetFor(session) diff --git a/internal/handler/setup_config.go b/internal/handler/setup_config.go new file mode 100644 index 0000000..45189c7 --- /dev/null +++ b/internal/handler/setup_config.go @@ -0,0 +1,521 @@ +package handler + +import ( + "context" + "errors" + "fmt" + "log" + "net/http" + "strings" + "sync" + "time" + + "github.com/gin-gonic/gin" + "github.com/pelletier/go-toml/v2" + "golang.org/x/sync/errgroup" + + "github.com/ic3software/vtafarm-api/internal/k8s" + "github.com/ic3software/vtafarm-api/internal/model" +) + +const ( + maxConfigBytes = 192 << 10 + configOperationLimit = 12 * time.Minute + configCrashRestartLimit = 3 +) + +type stackConfigTarget struct { + name string + deployment string + selector string + pvc string + mountPath string + configPath string + inputKey string +} + +type stackConfigRequest struct { + Component string `json:"component"` + Content string `json:"content"` +} + +func stackConfigTargets(session *model.SetupSession) []stackConfigTarget { + sessionID := session.ID + if !session.IsFullStack() { + return []stackConfigTarget{{"vta", k8s.VtaDeploymentName(sessionID), fmt.Sprintf("app=vta,session-id=%d", sessionID), k8s.VtaPVCName(sessionID), "/work/vta", "/work/vta/config.toml", "vta.toml"}} + } + return []stackConfigTarget{ + {"vta", k8s.FSVtaName(sessionID), fmt.Sprintf("app=fs-vta,session-id=%d", sessionID), k8s.FSVtaName(sessionID), "/work/vta", "/work/vta/config.toml", "vta.toml"}, + {"mediator", k8s.FSMediatorName(sessionID), fmt.Sprintf("app=fs-mediator,session-id=%d", sessionID), k8s.FSMediatorName(sessionID), "/work/mediator", "/work/mediator/conf/mediator.toml", "mediator.toml"}, + {"dids", k8s.FSDidsName(sessionID), fmt.Sprintf("app=fs-dids,session-id=%d", sessionID), k8s.FSDidsName(sessionID), "/work/dids", "/work/dids/config.toml", "dids.toml"}, + {"vtc", k8s.FSVtcName(sessionID), fmt.Sprintf("app=fs-vtc,session-id=%d", sessionID), k8s.FSVtcName(sessionID), "/work/vtc", "/work/vtc/config.toml", "vtc.toml"}, + } +} + +// Owner-facing routes use the session's own components. A vta_only agent does +// not own the platform components it uses. +func (h *SetupHandler) GetStackConfigs(c *gin.Context) { + if session := h.userSession(c); session != nil { + h.getStackConfigs(c, session) + } +} + +func (h *SetupHandler) ValidateStackConfigs(c *gin.Context) { + if session := h.userSession(c); session != nil { + h.validateStackConfigs(c, session) + } +} + +func (h *SetupHandler) ApplyStackConfigs(c *gin.Context) { + if session := h.userSession(c); session != nil { + h.applyStackConfigs(c, session) + } +} + +// Admin variants deliberately resolve only the platform stack. General admins +// do not gain a TOML editor for customer stacks through this feature. +func (h *SetupHandler) AdminGetPlatformStackConfigs(c *gin.Context) { + if session := h.platformSession(c); session != nil { + h.getStackConfigs(c, session) + } +} + +func (h *SetupHandler) AdminValidatePlatformStackConfigs(c *gin.Context) { + if session := h.platformSession(c); session != nil { + h.validateStackConfigs(c, session) + } +} + +func (h *SetupHandler) AdminApplyPlatformStackConfigs(c *gin.Context) { + if session := h.platformSession(c); session != nil { + h.applyStackConfigs(c, session) + } +} + +func configSessionReady(c *gin.Context, session *model.SetupSession) bool { + if session.Status != "running" { + c.JSON(http.StatusConflict, gin.H{"error": "the session must be running before its configuration can be changed"}) + return false + } + return true +} + +func (h *SetupHandler) getStackConfigs(c *gin.Context, session *model.SetupSession) { + if !configSessionReady(c, session) { + return + } + target, ok := stackConfigTargetFor(session, c.Query("component")) + if !ok { + c.JSON(http.StatusBadRequest, gin.H{"error": "unknown configuration component"}) + return + } + content, err := h.readStackConfig(c.Request.Context(), session, target) + if err != nil { + c.JSON(http.StatusBadGateway, gin.H{"error": "failed to read configuration: " + err.Error()}) + return + } + c.Header("Cache-Control", "no-store") + c.JSON(http.StatusOK, gin.H{"content": content}) +} + +func (h *SetupHandler) validateStackConfigs(c *gin.Context, session *model.SetupSession) { + if !configSessionReady(c, session) { + return + } + req, ok := bindStackConfig(c) + if !ok { + return + } + if validationErrors := validateTOMLConfig(req, stackConfigTargets(session)); len(validationErrors) > 0 { + c.JSON(http.StatusUnprocessableEntity, gin.H{"error": "TOML validation failed", "validation_errors": validationErrors}) + return + } + c.JSON(http.StatusOK, gin.H{"valid": true}) +} + +func bindStackConfig(c *gin.Context) (stackConfigRequest, bool) { + requestLimit := maxConfigBytes + (64 << 10) + c.Request.Body = http.MaxBytesReader(c.Writer, c.Request.Body, int64(requestLimit)) + var req stackConfigRequest + if err := c.ShouldBindJSON(&req); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": "invalid request body: " + err.Error()}) + return req, false + } + return req, true +} + +func stackConfigTargetFor(session *model.SetupSession, name string) (stackConfigTarget, bool) { + for _, target := range stackConfigTargets(session) { + if target.name == name { + return target, true + } + } + return stackConfigTarget{}, false +} + +func validateTOMLConfig(req stackConfigRequest, targets []stackConfigTarget) map[string]string { + validationErrors := make(map[string]string) + allowed := false + for _, target := range targets { + if target.name == req.Component { + allowed = true + break + } + } + if !allowed { + validationErrors["component"] = "unknown component for this session" + return validationErrors + } + if strings.TrimSpace(req.Content) == "" { + validationErrors[req.Component] = "configuration cannot be empty" + } else if len(req.Content) > maxConfigBytes { + validationErrors[req.Component] = fmt.Sprintf("configuration exceeds %d KiB", maxConfigBytes>>10) + } else { + var parsed map[string]any + if err := toml.Unmarshal([]byte(req.Content), &parsed); err != nil { + validationErrors[req.Component] = err.Error() + } + } + return validationErrors +} + +func (h *SetupHandler) applyStackConfigs(c *gin.Context, session *model.SetupSession) { + if !configSessionReady(c, session) { + return + } + req, ok := bindStackConfig(c) + if !ok { + return + } + if validationErrors := validateTOMLConfig(req, stackConfigTargets(session)); len(validationErrors) > 0 { + c.JSON(http.StatusUnprocessableEntity, gin.H{"error": "TOML validation failed", "validation_errors": validationErrors}) + return + } + target, _ := stackConfigTargetFor(session, req.Component) + targets := []stackConfigTarget{target} + if h.k8s == nil { + c.JSON(http.StatusServiceUnavailable, gin.H{"error": "k8s not configured"}) + return + } + + release, err := h.acquireSessionMaintenance(session.ID) + if err != nil { + status := http.StatusBadGateway + if errors.Is(err, errAclJobBusy) { + status = http.StatusConflict + } + c.JSON(status, gin.H{"error": err.Error()}) + return + } + defer release() + var upgradesInFlight int64 + if err := h.db.Model(&model.UpgradeTask{}). + Where("session_id = ? AND status IN ?", session.ID, + []string{model.UpgradeTaskPending, model.UpgradeTaskRunning}). + Count(&upgradesInFlight).Error; err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": "failed to check for an in-progress upgrade"}) + return + } + if upgradesInFlight > 0 { + c.JSON(http.StatusConflict, gin.H{"error": "an upgrade for this session is already in progress"}) + return + } + + currentContent, err := h.readStackConfig(c.Request.Context(), session, target) + if err != nil { + c.JSON(http.StatusBadGateway, gin.H{"error": "failed to read current component configuration: " + err.Error()}) + return + } + if currentContent == req.Content { + c.JSON(http.StatusOK, gin.H{"status": "unchanged"}) + return + } + current := map[string]string{target.name: currentContent} + + ctx, cancel := context.WithTimeout(context.Background(), configOperationLimit) + defer cancel() + ns := h.k8s.UserNamespace(fmt.Sprintf("%d", session.UserID)) + if err := h.stopStack(ctx, ns, targets); err != nil { + recoveryCtx, recoveryCancel := context.WithTimeout(context.Background(), configOperationLimit) + defer recoveryCancel() + restartErr := h.startStack(recoveryCtx, ns, targets) + c.JSON(http.StatusBadGateway, gin.H{"error": joinOperationErrors("failed to stop the component", err, restartErr)}) + return + } + + if err := h.runConfigWriteJob(ctx, ns, session, targets, map[string]string{target.name: req.Content}); err != nil { + rollbackErr := h.recoverStackConfigs(ns, session, targets, current) + body := gin.H{ + "error": "failed to write configuration; the previous configuration was restored", + "write_error": err.Error(), + "rolled_back": rollbackErr == nil, + } + if rollbackErr != nil { + body["error"] = "failed to write configuration and automatic rollback was incomplete" + body["rollback_error"] = rollbackErr.Error() + } + c.JSON(http.StatusBadGateway, body) + return + } + + if err := h.startStack(ctx, ns, targets); err == nil { + h.cleanupConfigBackups(ctx, ns, targets) + c.JSON(http.StatusOK, gin.H{"status": "applied", "validated": true}) + return + } else { + startupErr := err + rollbackErr := h.recoverStackConfigs(ns, session, targets, current) + body := gin.H{ + "error": "the new configuration did not pass startup validation; the previous configuration was restored", + "startup_error": startupErr.Error(), + "rolled_back": rollbackErr == nil, + } + if rollbackErr != nil { + body["error"] = "the new configuration did not pass startup validation and automatic rollback was incomplete" + body["rollback_error"] = rollbackErr.Error() + } + c.JSON(http.StatusUnprocessableEntity, body) + } +} + +func (h *SetupHandler) recoverStackConfigs( + ns string, + session *model.SetupSession, + targets []stackConfigTarget, + previous map[string]string, +) error { + // The apply deadline may already have expired. Recovery needs its own window. + ctx, cancel := context.WithTimeout(context.Background(), configOperationLimit) + defer cancel() + return h.rollbackStackConfigs(ctx, ns, session, targets, previous) +} + +func (h *SetupHandler) readStackConfig(ctx context.Context, session *model.SetupSession, target stackConfigTarget) (string, error) { + if h.k8s == nil { + return "", fmt.Errorf("k8s not configured") + } + ns := h.k8s.UserNamespace(fmt.Sprintf("%d", session.UserID)) + pod, container, err := h.k8s.RunningPod(ctx, ns, target.selector) + if err != nil { + return "", fmt.Errorf("%s: %w", target.name, err) + } + content, err := h.k8s.ExecCapture(ctx, ns, pod, container, + []string{"sh", "-c", "cat " + shellQuote(target.configPath)}) + if err != nil { + return "", fmt.Errorf("%s: %w", target.name, err) + } + return content, nil +} + +func (h *SetupHandler) stopStack(ctx context.Context, ns string, targets []stackConfigTarget) error { + parentCtx := ctx + var errs []error + var mu sync.Mutex + g, ctx := errgroup.WithContext(parentCtx) + for _, target := range targets { + target := target + g.Go(func() error { + if err := h.k8s.ScaleComponentDeployment(ctx, ns, target.deployment, 0); err != nil { + mu.Lock() + errs = append(errs, fmt.Errorf("%s: %w", target.name, err)) + mu.Unlock() + } + return nil + }) + } + _ = g.Wait() + if len(errs) > 0 { + return errors.Join(errs...) + } + + g, ctx = errgroup.WithContext(parentCtx) + for _, target := range targets { + target := target + g.Go(func() error { + if err := h.k8s.WaitForComponentPodsGone(ctx, ns, target.selector, 2*time.Minute); err != nil { + return fmt.Errorf("%s: %w", target.name, err) + } + return nil + }) + } + return g.Wait() +} + +func (h *SetupHandler) startStack(ctx context.Context, ns string, targets []stackConfigTarget) error { + parentCtx := ctx + var errs []error + var mu sync.Mutex + g, ctx := errgroup.WithContext(parentCtx) + for _, target := range targets { + target := target + g.Go(func() error { + if err := h.k8s.ScaleComponentDeployment(ctx, ns, target.deployment, 1); err != nil { + mu.Lock() + errs = append(errs, fmt.Errorf("%s scale: %w", target.name, err)) + mu.Unlock() + } + return nil + }) + } + _ = g.Wait() + + g, ctx = errgroup.WithContext(parentCtx) + for _, target := range targets { + target := target + g.Go(func() error { + if err := h.k8s.WaitForComponentDeploymentReadyOrRestarts(ctx, ns, target.deployment, target.selector, 2*time.Minute, configCrashRestartLimit); err != nil { + mu.Lock() + errs = append(errs, fmt.Errorf("%s readiness: %w", target.name, err)) + mu.Unlock() + } + return nil + }) + } + _ = g.Wait() + return errors.Join(errs...) +} + +func (h *SetupHandler) runConfigWriteJob( + ctx context.Context, + ns string, + session *model.SetupSession, + targets []stackConfigTarget, + configs map[string]string, +) error { + jobName := k8s.FSJobConfigUpdate(session.ID) + h.k8s.DeleteComponentJob(ctx, ns, jobName) + secretData := make(map[string][]byte, len(targets)) + var commands []string + for _, target := range targets { + secretData[target.inputKey] = []byte(configs[target.name]) + path := shellQuote(target.configPath) + next := shellQuote(target.configPath + ".vtafarm-next") + backup := shellQuote(target.configPath + ".vtafarm-backup") + input := shellQuote("/config/" + target.inputKey) + commands = append(commands, + fmt.Sprintf("cp -p %s %s", path, backup), + fmt.Sprintf("cp -p %s %s", path, next), + fmt.Sprintf("cat %s > %s", input, next), + fmt.Sprintf("mv %s %s", next, path), + ) + } + if err := h.k8s.CreateComponentJob(ctx, ns, k8s.ComponentJobSpec{ + Name: jobName, + Image: session.VtaImage, + Command: []string{"sh", "-c", "set -eu; umask 077; " + strings.Join(commands, "; ")}, + WorkingDir: targets[0].mountPath, + ServiceAccount: k8s.VtaServiceAccount, + PVCMounts: configPVCMounts(targets), + SecretData: secretData, + Env: noColorEnv(), + }); err != nil { + h.k8s.DeleteComponentJob(context.Background(), ns, jobName) + return err + } + defer h.k8s.DeleteComponentJob(context.Background(), ns, jobName) + return h.waitConfigJob(ctx, ns, jobName) +} + +func (h *SetupHandler) rollbackStackConfigs( + ctx context.Context, + ns string, + session *model.SetupSession, + targets []stackConfigTarget, + previous map[string]string, +) error { + var errs []error + if err := h.stopStack(ctx, ns, targets); err != nil { + errs = append(errs, fmt.Errorf("stop before rollback: %w", err)) + if restartErr := h.startStack(ctx, ns, targets); restartErr != nil { + errs = append(errs, fmt.Errorf("restart after failed rollback stop: %w", restartErr)) + } + return errors.Join(errs...) + } + jobName := k8s.FSJobConfigRollback(session.ID) + h.k8s.DeleteComponentJob(ctx, ns, jobName) + secretData := make(map[string][]byte, len(targets)) + var commands []string + for _, target := range targets { + secretData[target.inputKey] = []byte(previous[target.name]) + path := shellQuote(target.configPath) + next := shellQuote(target.configPath + ".vtafarm-next") + backup := shellQuote(target.configPath + ".vtafarm-backup") + input := shellQuote("/config/" + target.inputKey) + commands = append(commands, + fmt.Sprintf("rm -f %s", next), + fmt.Sprintf("cp -p %s %s", path, next), + fmt.Sprintf("cat %s > %s", input, next), + fmt.Sprintf("mv %s %s", next, path), + fmt.Sprintf("rm -f %s", backup), + ) + } + if err := h.k8s.CreateComponentJob(ctx, ns, k8s.ComponentJobSpec{ + Name: jobName, + Image: session.VtaImage, + Command: []string{"sh", "-c", "set -eu; " + strings.Join(commands, "; ")}, + WorkingDir: targets[0].mountPath, + ServiceAccount: k8s.VtaServiceAccount, + PVCMounts: configPVCMounts(targets), + SecretData: secretData, + Env: noColorEnv(), + }); err != nil { + h.k8s.DeleteComponentJob(context.Background(), ns, jobName) + errs = append(errs, fmt.Errorf("create rollback job: %w", err)) + } else { + if err := h.waitConfigJob(ctx, ns, jobName); err != nil { + errs = append(errs, fmt.Errorf("rollback job: %w", err)) + } + h.k8s.DeleteComponentJob(context.Background(), ns, jobName) + } + if err := h.startStack(ctx, ns, targets); err != nil { + errs = append(errs, fmt.Errorf("restart after rollback: %w", err)) + } + return errors.Join(errs...) +} + +func configPVCMounts(targets []stackConfigTarget) []k8s.PVCMount { + mounts := make([]k8s.PVCMount, 0, len(targets)) + for _, target := range targets { + mounts = append(mounts, k8s.PVCMount{ + Name: target.name + "-data", ClaimName: target.pvc, MountPath: target.mountPath, + }) + } + return mounts +} + +func (h *SetupHandler) waitConfigJob(ctx context.Context, ns, jobName string) error { + succeeded, failMsg, err := h.k8s.WaitForJob(ctx, ns, jobName) + if err != nil { + return err + } + if !succeeded { + if logs, logErr := h.k8s.JobLogs(ctx, ns, jobName); logErr == nil && strings.TrimSpace(logs) != "" { + failMsg = strings.TrimSpace(logs) + } + return fmt.Errorf("job failed: %s", failMsg) + } + return nil +} + +func (h *SetupHandler) cleanupConfigBackups(ctx context.Context, ns string, targets []stackConfigTarget) { + for _, target := range targets { + pod, container, err := h.k8s.RunningPod(ctx, ns, target.selector) + if err != nil { + log.Printf("[config] warn: find %s pod for backup cleanup: %v", target.name, err) + continue + } + _, err = h.k8s.ExecCapture(ctx, ns, pod, container, + []string{"sh", "-c", "rm -f " + shellQuote(target.configPath+".vtafarm-backup")}) + if err != nil { + log.Printf("[config] warn: clean %s config backup: %v", target.name, err) + } + } +} + +func joinOperationErrors(prefix string, operationErr, restartErr error) string { + message := prefix + ": " + operationErr.Error() + if restartErr != nil { + message += "; WARNING: one or more components did not restart: " + restartErr.Error() + } + return message +} diff --git a/internal/handler/setup_config_test.go b/internal/handler/setup_config_test.go new file mode 100644 index 0000000..9b552a3 --- /dev/null +++ b/internal/handler/setup_config_test.go @@ -0,0 +1,87 @@ +package handler + +import ( + "strings" + "testing" + + "github.com/ic3software/vtafarm-api/internal/model" +) + +func TestValidateTOMLConfig(t *testing.T) { + targets := stackConfigTargets(&model.SetupSession{ID: 7, Mode: model.ModeFullStack}) + for _, name := range []string{"vta", "mediator", "dids", "vtc"} { + if got := validateTOMLConfig(stackConfigRequest{Component: name, Content: "[server]\nport = 8100\n"}, targets); len(got) != 0 { + t.Errorf("valid %s config returned errors: %v", name, got) + } + } + + cases := []struct { + req stackConfigRequest + key string + }{ + {stackConfigRequest{Component: "vta", Content: "[server\n"}, "vta"}, + {stackConfigRequest{Component: "mediator", Content: ""}, "mediator"}, + {stackConfigRequest{Component: "extra", Content: "ok = true"}, "component"}, + {stackConfigRequest{Content: "ok = true"}, "component"}, + } + for _, test := range cases { + if got := validateTOMLConfig(test.req, targets); got[test.key] == "" { + t.Errorf("missing %q validation error for %+v: %v", test.key, test.req, got) + } + } +} + +func TestValidateTOMLConfigRejectsOversize(t *testing.T) { + targets := stackConfigTargets(&model.SetupSession{ID: 7, Mode: model.ModeFullStack}) + req := stackConfigRequest{Component: "vta", Content: "value = \"" + strings.Repeat("x", maxConfigBytes) + "\""} + if got := validateTOMLConfig(req, targets)["vta"]; !strings.Contains(got, "exceeds") { + t.Fatalf("oversize error = %q, want size limit", got) + } +} + +func TestStackConfigTargets(t *testing.T) { + targets := stackConfigTargets(&model.SetupSession{ID: 7, Mode: model.ModeFullStack}) + if len(targets) != 4 { + t.Fatalf("targets = %d, want 4", len(targets)) + } + want := map[string]string{ + "vta": "/work/vta/config.toml", "mediator": "/work/mediator/conf/mediator.toml", + "dids": "/work/dids/config.toml", "vtc": "/work/vtc/config.toml", + } + for _, target := range targets { + if target.configPath != want[target.name] { + t.Errorf("%s path = %q, want %q", target.name, target.configPath, want[target.name]) + } + } +} + +func TestVtaOnlyConfigTargetsAndValidation(t *testing.T) { + targets := stackConfigTargets(&model.SetupSession{ID: 7, Mode: model.ModeVtaOnly}) + if len(targets) != 1 { + t.Fatalf("targets = %d, want 1", len(targets)) + } + target := targets[0] + if target.name != "vta" || target.deployment != "vta-7" || target.pvc != "vta-data-7" || target.selector != "app=vta,session-id=7" || target.configPath != "/work/vta/config.toml" { + t.Fatalf("unexpected VTA-only target: %+v", target) + } + if got := validateTOMLConfig(stackConfigRequest{Component: "vta", Content: "[server]\nport = 8100\n"}, targets); len(got) != 0 { + t.Fatalf("valid VTA-only config returned errors: %v", got) + } + if got := validateTOMLConfig(stackConfigRequest{Component: "mediator", Content: "ok = true"}, targets); got["component"] == "" { + t.Fatalf("shared component was accepted: %v", got) + } + if got := validateTOMLConfig(stackConfigRequest{Component: "vta", Content: "[server\n"}, targets); got["vta"] == "" { + t.Fatalf("invalid VTA-only TOML was accepted: %v", got) + } +} + +func TestStackConfigTargetForSelectsOnlyRequestedComponent(t *testing.T) { + session := &model.SetupSession{ID: 7, Mode: model.ModeFullStack} + target, ok := stackConfigTargetFor(session, "vtc") + if !ok || target.name != "vtc" || target.deployment != "fs-7-vtc" { + t.Fatalf("VTC target = %+v, %v", target, ok) + } + if _, ok := stackConfigTargetFor(session, "unknown"); ok { + t.Fatal("unknown component was accepted") + } +} diff --git a/internal/handler/setup_fullstack.go b/internal/handler/setup_fullstack.go index ccbb472..0ad127c 100644 --- a/internal/handler/setup_fullstack.go +++ b/internal/handler/setup_fullstack.go @@ -2,6 +2,7 @@ package handler import ( "context" + "errors" "fmt" "log" "net/http" @@ -536,6 +537,16 @@ func (h *SetupHandler) reissueDidsEnroll(c *gin.Context, session *model.SetupSes c.JSON(http.StatusServiceUnavailable, gin.H{"error": "k8s not configured"}) return } + release, err := h.acquireSessionMaintenance(session.ID) + if err != nil { + status := http.StatusBadGateway + if errors.Is(err, errAclJobBusy) { + status = http.StatusConflict + } + c.JSON(status, gin.H{"error": err.Error()}) + return + } + defer release() ctx := c.Request.Context() ns := h.k8s.UserNamespace(fmt.Sprintf("%d", session.UserID)) diff --git a/internal/handler/setup_vtc.go b/internal/handler/setup_vtc.go index 091a201..f2f470a 100644 --- a/internal/handler/setup_vtc.go +++ b/internal/handler/setup_vtc.go @@ -2,6 +2,7 @@ package handler import ( "context" + "errors" "fmt" "log" "net/http" @@ -53,6 +54,16 @@ func (h *SetupHandler) reissueVtcInstall(c *gin.Context, session *model.SetupSes c.JSON(http.StatusServiceUnavailable, gin.H{"error": "k8s not configured"}) return } + release, err := h.acquireSessionMaintenance(session.ID) + if err != nil { + status := http.StatusBadGateway + if errors.Is(err, errAclJobBusy) { + status = http.StatusConflict + } + c.JSON(status, gin.H{"error": err.Error()}) + return + } + defer release() ctx := c.Request.Context() ns := h.k8s.UserNamespace(fmt.Sprintf("%d", session.UserID)) diff --git a/internal/k8s/component_jobs.go b/internal/k8s/component_jobs.go index 12fbd0c..f784a32 100644 --- a/internal/k8s/component_jobs.go +++ b/internal/k8s/component_jobs.go @@ -39,6 +39,11 @@ type ComponentJobSpec struct { ConfigMapKey string // e.g. "vta-setup.toml", mounted at /config/ ConfigMapData string + // SecretData is mounted read-only at /config, like the ConfigMap above, but + // keeps credential-bearing inputs (such as a rendered config.toml) out of a + // ConfigMap. It is mutually exclusive with ConfigMapName. + SecretData map[string][]byte + Env []corev1.EnvVar // ActiveDeadlineSeconds defaults to 600 (10 min) when zero. @@ -79,6 +84,9 @@ func (c *Client) CreateComponentPVC(ctx context.Context, ns, name, storageSize s // described by spec. WaitForJob/JobLogs/StreamJobLogs (setup_jobs.go) are // already generic over job name and are reused as-is to drive it. func (c *Client) CreateComponentJob(ctx context.Context, ns string, spec ComponentJobSpec) error { + if spec.ConfigMapName != "" && len(spec.SecretData) > 0 { + return fmt.Errorf("component job %s cannot use both ConfigMapData and SecretData", spec.Name) + } if spec.ConfigMapName != "" { _, err := c.kube.CoreV1().ConfigMaps(ns).Create(ctx, &corev1.ConfigMap{ ObjectMeta: metav1.ObjectMeta{Name: spec.ConfigMapName, Namespace: ns}, @@ -88,6 +96,15 @@ func (c *Client) CreateComponentJob(ctx context.Context, ns string, spec Compone return fmt.Errorf("create configmap %s: %w", spec.ConfigMapName, err) } } + if len(spec.SecretData) > 0 { + _, err := c.kube.CoreV1().Secrets(ns).Create(ctx, &corev1.Secret{ + ObjectMeta: metav1.ObjectMeta{Name: spec.Name, Namespace: ns}, + Data: spec.SecretData, + }, metav1.CreateOptions{}) + if err != nil && !k8serrors.IsAlreadyExists(err) { + return fmt.Errorf("create secret %s: %w", spec.Name, err) + } + } volumes := make([]corev1.Volume, 0, len(spec.PVCMounts)+1) mounts := make([]corev1.VolumeMount, 0, len(spec.PVCMounts)+1) @@ -111,6 +128,15 @@ func (c *Client) CreateComponentJob(ctx context.Context, ns string, spec Compone }) mounts = append(mounts, corev1.VolumeMount{Name: "config", MountPath: "/config"}) } + if len(spec.SecretData) > 0 { + volumes = append(volumes, corev1.Volume{ + Name: "config", + VolumeSource: corev1.VolumeSource{ + Secret: &corev1.SecretVolumeSource{SecretName: spec.Name}, + }, + }) + mounts = append(mounts, corev1.VolumeMount{Name: "config", MountPath: "/config", ReadOnly: true}) + } backoff := int32(0) ttl := int32(3600) @@ -155,9 +181,10 @@ func (c *Client) DeleteComponentJob(ctx context.Context, ns, name string) { opts := metav1.DeleteOptions{PropagationPolicy: &propagation} _ = c.kube.BatchV1().Jobs(ns).Delete(ctx, name, opts) _ = c.kube.CoreV1().ConfigMaps(ns).Delete(ctx, name, metav1.DeleteOptions{}) + _ = c.kube.CoreV1().Secrets(ns).Delete(ctx, name, metav1.DeleteOptions{}) } -// DeleteAllComponentJobs removes every full_stack setup Job (+ ConfigMap) +// DeleteAllComponentJobs removes every full_stack setup Job (+ ConfigMap or Secret) // for a session. Best-effort. func (c *Client) DeleteAllComponentJobs(ctx context.Context, ns string, sessionID uint) { for _, name := range allFSJobNames(sessionID) { diff --git a/internal/k8s/component_jobs_test.go b/internal/k8s/component_jobs_test.go index 06b3cc7..97438cd 100644 --- a/internal/k8s/component_jobs_test.go +++ b/internal/k8s/component_jobs_test.go @@ -52,6 +52,36 @@ func TestCreateComponentPVCRejectsInvalidStorageSize(t *testing.T) { } } +func TestCreateComponentJobMountsSecretConfig(t *testing.T) { + client := &Client{kube: fake.NewSimpleClientset()} + if err := client.CreateComponentJob(context.Background(), "test", ComponentJobSpec{ + Name: "config-update", + Image: "example/vta:test", + SecretData: map[string][]byte{"vta.toml": []byte("ok = true\n")}, + }); err != nil { + t.Fatalf("CreateComponentJob() error = %v", err) + } + secret, err := client.kube.CoreV1().Secrets("test").Get( + context.Background(), "config-update", metav1.GetOptions{}, + ) + if err != nil { + t.Fatalf("get Secret: %v", err) + } + if got := string(secret.Data["vta.toml"]); got != "ok = true\n" { + t.Fatalf("secret data = %q", got) + } + job, err := client.kube.BatchV1().Jobs("test").Get( + context.Background(), "config-update", metav1.GetOptions{}, + ) + if err != nil { + t.Fatalf("get Job: %v", err) + } + mounts := job.Spec.Template.Spec.Containers[0].VolumeMounts + if len(mounts) != 1 || mounts[0].MountPath != "/config" || !mounts[0].ReadOnly { + t.Fatalf("secret mount = %#v", mounts) + } +} + func TestCreateComponentDeploymentUsesConfiguredHealthReadinessProbe(t *testing.T) { client := &Client{kube: fake.NewSimpleClientset()} diff --git a/internal/k8s/component_readiness_test.go b/internal/k8s/component_readiness_test.go index b35955c..ac7d7b8 100644 --- a/internal/k8s/component_readiness_test.go +++ b/internal/k8s/component_readiness_test.go @@ -2,9 +2,12 @@ package k8s import ( "context" + "strings" "testing" + "time" appsv1 "k8s.io/api/apps/v1" + corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/client-go/kubernetes/fake" ) @@ -28,3 +31,40 @@ func TestComponentDeploymentReady(t *testing.T) { t.Fatalf("missing deployment = ready %v, err %v", ready, err) } } + +func TestWaitForComponentDeploymentReadyOrRestarts(t *testing.T) { + for _, test := range []struct { + name string + restarts int32 + podReady bool + readyReplicas int32 + want string + }{ + {name: "ready", podReady: true, readyReplicas: 1}, + {name: "three crashes", restarts: 3, want: "restarted 3 times"}, + {name: "two crashes still waits", restarts: 2, want: "timeout"}, + {name: "stale deployment readiness", restarts: 2, readyReplicas: 1, want: "timeout"}, + } { + t.Run(test.name, func(t *testing.T) { + client := &Client{kube: fake.NewSimpleClientset( + &appsv1.Deployment{ + ObjectMeta: metav1.ObjectMeta{Name: "vta-42", Namespace: "fpp-user-1", Generation: 2}, + Status: appsv1.DeploymentStatus{ObservedGeneration: 2, ReadyReplicas: test.readyReplicas}, + }, + &corev1.Pod{ + ObjectMeta: metav1.ObjectMeta{Name: "vta-42-pod", Namespace: "fpp-user-1", Labels: map[string]string{"app": "vta", "session-id": "42"}}, + Spec: corev1.PodSpec{Containers: []corev1.Container{{Name: "vta"}}}, + Status: corev1.PodStatus{ContainerStatuses: []corev1.ContainerStatus{{Name: "vta", RestartCount: test.restarts, Ready: test.podReady}}}, + }, + )} + err := client.WaitForComponentDeploymentReadyOrRestarts(context.Background(), "fpp-user-1", "vta-42", "app=vta,session-id=42", 20*time.Millisecond, 3) + if test.want == "" { + if err != nil { + t.Fatalf("ready deployment returned error: %v", err) + } + } else if err == nil || !strings.Contains(err.Error(), test.want) { + t.Fatalf("error = %v, want %q", err, test.want) + } + }) + } +} diff --git a/internal/k8s/component_resources.go b/internal/k8s/component_resources.go index c310fe5..07192ea 100644 --- a/internal/k8s/component_resources.go +++ b/internal/k8s/component_resources.go @@ -198,6 +198,48 @@ func (c *Client) WaitForComponentDeploymentReady(ctx context.Context, ns, name s } } +// WaitForComponentDeploymentReadyOrRestarts is for configuration updates: a +// repeatedly crashing new pod should trigger rollback before the full readiness +// timeout, while a pod that simply never becomes ready still has that timeout. +func (c *Client) WaitForComponentDeploymentReadyOrRestarts(ctx context.Context, ns, name, selector string, timeout time.Duration, restartLimit int32) error { + deadline := time.NewTimer(timeout) + defer deadline.Stop() + ticker := time.NewTicker(2 * time.Second) + defer ticker.Stop() + for { + deploy, deployErr := c.kube.AppsV1().Deployments(ns).Get(ctx, name, metav1.GetOptions{}) + pods, podsErr := c.kube.CoreV1().Pods(ns).List(ctx, metav1.ListOptions{LabelSelector: selector}) + podReady := false + if podsErr == nil { + for _, pod := range pods.Items { + if pod.DeletionTimestamp != nil || len(pod.Spec.Containers) == 0 { + continue + } + mainContainer := pod.Spec.Containers[0].Name + for _, status := range pod.Status.ContainerStatuses { + if status.Name != mainContainer { + continue + } + if status.RestartCount >= restartLimit { + return fmt.Errorf("pod %s restarted %d times during configuration startup check", pod.Name, status.RestartCount) + } + podReady = podReady || status.Ready + } + } + } + if deployErr == nil && podsErr == nil && deploy.Status.ObservedGeneration >= deploy.Generation && deploy.Status.ReadyReplicas > 0 && podReady { + return nil + } + select { + case <-ctx.Done(): + return ctx.Err() + case <-deadline.C: + return fmt.Errorf("timeout waiting for deployment %s to become ready", name) + case <-ticker.C: + } + } +} + // ComponentDeploymentReady reports the current readiness state without // waiting. A ready replica means the workload's configured readiness probe has // succeeded; callers use this for explicit admin health checks. diff --git a/internal/k8s/fullstack_names.go b/internal/k8s/fullstack_names.go index 9132c9c..12fb9a7 100644 --- a/internal/k8s/fullstack_names.go +++ b/internal/k8s/fullstack_names.go @@ -73,10 +73,14 @@ func FSJobVtaACL(sessionID uint) string { return fmt.Sprintf("fs-%d-vta-acl", se // full_stack-only Jobs (design §8). FSJobVtcInvite is the reissue // endpoint's `vtc admin invite` Job (POST /setup/:id/vtc/reissue-install), // not a pipeline step — mirrors how FSJobDidsInvite doubles for reissue. -func FSJobVtcSetupKey(sessionID uint) string { return fmt.Sprintf("fs-%d-vtc-setup-key", sessionID) } -func FSJobVtcAclGrant(sessionID uint) string { return fmt.Sprintf("fs-%d-vtc-acl-grant", sessionID) } -func FSJobVtcSetup(sessionID uint) string { return fmt.Sprintf("fs-%d-vtc-setup", sessionID) } -func FSJobVtcInvite(sessionID uint) string { return fmt.Sprintf("fs-%d-vtc-invite", sessionID) } +func FSJobVtcSetupKey(sessionID uint) string { return fmt.Sprintf("fs-%d-vtc-setup-key", sessionID) } +func FSJobVtcAclGrant(sessionID uint) string { return fmt.Sprintf("fs-%d-vtc-acl-grant", sessionID) } +func FSJobVtcSetup(sessionID uint) string { return fmt.Sprintf("fs-%d-vtc-setup", sessionID) } +func FSJobVtcInvite(sessionID uint) string { return fmt.Sprintf("fs-%d-vtc-invite", sessionID) } +func FSJobConfigUpdate(sessionID uint) string { return fmt.Sprintf("fs-%d-config-update", sessionID) } +func FSJobConfigRollback(sessionID uint) string { + return fmt.Sprintf("fs-%d-config-rollback", sessionID) +} // allFSJobNames lists every setup Job name for a session — used by teardown // to best-effort delete each one (and its ConfigMap, where one exists). @@ -99,5 +103,7 @@ func allFSJobNames(sessionID uint) []string { FSJobVtcAclGrant(sessionID), FSJobVtcSetup(sessionID), FSJobVtcInvite(sessionID), + FSJobConfigUpdate(sessionID), + FSJobConfigRollback(sessionID), } } diff --git a/internal/k8s/vta_resources.go b/internal/k8s/vta_resources.go index 948520a..d14bbaf 100644 --- a/internal/k8s/vta_resources.go +++ b/internal/k8s/vta_resources.go @@ -246,5 +246,7 @@ func (c *Client) DeleteVtaResources(ctx context.Context, ns string, sessionID ui _ = c.kube.CoreV1().Services(ns).Delete(ctx, vtaServiceName(sessionID), metav1.DeleteOptions{}) _ = c.kube.BatchV1().Jobs(ns).Delete(ctx, ProvisionJobName(sessionID), opts) _ = c.kube.BatchV1().Jobs(ns).Delete(ctx, VtaACLJobName(sessionID), opts) + c.DeleteComponentJob(ctx, ns, FSJobConfigUpdate(sessionID)) + c.DeleteComponentJob(ctx, ns, FSJobConfigRollback(sessionID)) _ = c.kube.CoreV1().PersistentVolumeClaims(ns).Delete(ctx, VtaPVCName(sessionID), metav1.DeleteOptions{}) } diff --git a/internal/model/vta_acl.go b/internal/model/vta_acl.go index a98e950..ba3c988 100644 --- a/internal/model/vta_acl.go +++ b/internal/model/vta_acl.go @@ -4,7 +4,7 @@ import "time" // VtaAclSnapshot records when the farm last read the complete ACL directly // from a stopped VTA. MaintenanceStartedAt is also the cross-replica lock for -// every offline ACL operation on that session. +// every operation that stops one or more components in that session. type VtaAclSnapshot struct { SessionID uint `json:"-" gorm:"primaryKey;column:session_id"` SyncedAt *time.Time `json:"synced_at"` diff --git a/internal/router/router.go b/internal/router/router.go index c847c2e..bd6809b 100644 --- a/internal/router/router.go +++ b/internal/router/router.go @@ -205,6 +205,9 @@ func Setup( // a domains row for our own zone. adminAuth.POST("/admin/platform-stack", sh.CreatePlatformStack) adminAuth.GET("/admin/platform-stack", sh.GetPlatformStack) + adminAuth.GET("/admin/platform-stack/config", sh.AdminGetPlatformStackConfigs) + adminAuth.POST("/admin/platform-stack/config/validate", sh.AdminValidatePlatformStackConfigs) + adminAuth.PUT("/admin/platform-stack/config", sh.AdminApplyPlatformStackConfigs) // Co-admins on that stack's VTA — self-service, so a second admin can // add the did:key their own `pnm setup` minted instead of asking // whoever holds the credential to run `pnm acl create` for them. @@ -287,6 +290,9 @@ func Setup( // binary wrote to its PVC, and the pods' own logs — no Job logs. userAuth.GET("/setup/:id/export/configs", sh.ExportConfigs) userAuth.GET("/setup/:id/export/logs", sh.ExportLogs) + userAuth.GET("/setup/:id/config", sh.GetStackConfigs) + userAuth.POST("/setup/:id/config/validate", sh.ValidateStackConfigs) + userAuth.PUT("/setup/:id/config", sh.ApplyStackConfigs) // Self-service image upgrade/downgrade — a user can only ever change // their own session (looked up by unique_id AND user_id). userAuth.POST("/setup/:id/upgrade", uph.CreateForSession)