Files
felhom-controller/controller/internal/sync/sync.go
T
admin 8a0e0a59ad
gates / gates (push) Successful in 12s
v0.235.0: freeze the version, keep the fixes flowing (operator ruling 2026-09-06)
Slice 3. R-447 was BLOCKED because R-438 established that RestartStack's use of
up -d to pick up template changes was CHOSEN and written down in its own comment.
The operator ruled Option 1, and this implements it.

The rule: while the catalog offers the same version you run, its fixes flow to
you; the moment it moves to a newer version you are frozen until you update.

NOTHING was added to any of the thirteen compose up -d call sites. Most of them
are repairs - the boot reconciler, the drive-return gate, the app-stop guard -
and a repair path that refuses to repair leaves a customer's app down, which is
worse than the problem. They are made safe by removing the reason.

app.yaml gains pinned_images: what the app is SUPPOSED to run. It is NOT
installed_images, which is an observation; letting a reading become a deployment
is the R-166 category error one field over. Four writers, each also storing the
exact definition as applied-compose.yml. UpdateStack advances the pin and
re-renders BEFORE the pull, because pull and up -d act on the file on disk, and a
pin set afterwards would pull the frozen version and report success.

The syncer renders instead of copying, through one nil-safe seam. Catalog images
equal the pin -> verbatim, so fixes and self-healing both survive; they differ ->
the WHOLE stored definition, never a substitution of refs into a newer template
(wger 2.6 needs a DB config the older template cannot supply). This is
deliberately not 'skip deployed apps', which was option B and was rejected.

AdoptPins runs once at boot after the backfill, files only, and skips loudly
rather than inventing a pin. syncer.Start() moved to after it: the initial sync
would otherwise run while every app was unpinned and overwrite a deployed app's
version once per boot.

THE BADGE HAD TO CHANGE OR SLICE 2 WOULD HAVE INVERTED SILENTLY. TemplateImages
reads the LIVE compose file, which is now the frozen one, so the comparison would
have answered Naprakesz on exactly the apps that are behind - with every test
green, because the new field has the same type. It now reads CatalogImages.

+16 tests (1729 -> 1745), 28 packages green. Three red-proofs run and reverted.
A test also caught the syncer writing an empty compose file over a live app.
2026-09-06 09:45:34 +02:00

588 lines
21 KiB
Go

package sync
import (
"bytes"
"context"
"crypto/sha256"
"encoding/hex"
"fmt"
"io"
"log"
"os"
"os/exec"
"path/filepath"
"regexp"
"strings"
"sync"
"time"
"gitea.dooplex.hu/admin/felhom-controller/internal/config"
// stacks is imported for its PURE helpers only — RenderPlan and ParseComposeImages. The syncer
// must not reach for the Manager through it: everything it needs about an app arrives through
// renderPlanFn. Importing the package rather than re-writing a compose parser is deliberate;
// a second image parser is exactly the duplication REUSE.md exists to prevent.
"gitea.dooplex.hu/admin/felhom-controller/internal/stacks"
)
// gitCmdTimeout bounds each individual git subprocess (campaign finding F3, 2026-07-06):
// without a deadline, a hung remote parks the sync goroutine inside cmd.Run() — the doSync
// defer never runs, `syncing` stays true, and every manual + periodic sync is refused with
// "Szinkronizálás már folyamatban" until a controller restart. 120s is generous for the
// shallow catalog clone/fetch.
const gitCmdTimeout = 120 * time.Second
// Syncer handles periodic git sync of the app catalog to the local stacks directory.
type Syncer struct {
cfg *config.Config
logger *log.Logger
cacheDir string // local git clone
rescanFn func() error
postSyncHook func(updated []string) // called after sync with names of updated stacks
// renderPlanFn answers, for one app: is it deployed, is it pinned, and where is its stored
// applied definition (v0.235.0). It is the ONLY way this package learns any of that: the syncer
// must not import the stack manager and must NEVER read app.yaml itself — that file is the
// manager's and carries encrypted values.
//
// NIL-SAFE BY DESIGN: a nil seam means "copy verbatim", i.e. byte-for-byte the pre-v0.235.0
// behaviour. Every existing test that constructs a Syncer without one keeps passing, and a
// wiring mistake degrades to the old product rather than to a broken one.
renderPlanFn func(appName string) stacks.RenderPlan
mu sync.Mutex
lastSync time.Time
lastErr error
syncing bool
stopCh chan struct{}
stopOnce sync.Once
}
// SyncStatus holds information about the last sync operation.
type SyncStatus struct {
LastSync time.Time `json:"last_sync"`
LastStatus string `json:"last_status"` // "ok", "error", "disabled", "never"
LastError string `json:"last_error,omitempty"`
Syncing bool `json:"syncing"`
}
// SyncResult holds the result of a single sync operation.
type SyncResult struct {
OK bool `json:"ok"`
NewApps []string `json:"new_apps,omitempty"`
Updated []string `json:"updated,omitempty"`
Message string `json:"message"`
}
// New creates a new Syncer. rescanFn is called after a successful sync to trigger ScanStacks().
// postSyncHook is called with names of updated stacks (may be nil).
func New(cfg *config.Config, logger *log.Logger, rescanFn func() error, postSyncHook func([]string)) *Syncer {
cacheDir := filepath.Join(cfg.Paths.DataDir, "catalog-cache")
return &Syncer{
cfg: cfg,
logger: logger,
cacheDir: cacheDir,
rescanFn: rescanFn,
postSyncHook: postSyncHook,
stopCh: make(chan struct{}),
}
}
// SetRenderPlanFn injects the per-app render plan (v0.235.0). Exported and separate from New for the
// same reason SetSambaRunProbe is: New's signature is called from tests in several packages, and the
// seam has to be optional so a Syncer without one keeps the pre-v0.235.0 behaviour exactly.
func (s *Syncer) SetRenderPlanFn(fn func(appName string) stacks.RenderPlan) { s.renderPlanFn = fn }
// isDebug returns true if the logging level is set to "debug".
func (s *Syncer) isDebug() bool { return s.cfg.Logging.Level == "debug" }
// maskRepoURL masks credentials in a git URL for safe logging.
// e.g., "https://user:token@host/path" → "https://user:***@host/path"
var reURLCreds = regexp.MustCompile(`(https?://)([^:]+):([^@]+)@`)
func maskRepoURL(url string) string {
return reURLCreds.ReplaceAllString(url, "${1}${2}:***@")
}
// Start begins the periodic sync loop. Call Stop() to terminate.
func (s *Syncer) Start() {
if s.cfg.Git.RepoURL == "" {
s.logger.Println("[WARN] [sync] Git repo URL is empty — sync disabled (manual mode)")
return
}
interval, err := time.ParseDuration(s.cfg.Git.SyncInterval)
if err != nil {
s.logger.Printf("[WARN] [sync] Invalid sync_interval %q, defaulting to 15m", s.cfg.Git.SyncInterval)
interval = 15 * time.Minute
}
s.logger.Printf("[INFO] [sync] Starting catalog sync (repo: %s, interval: %s)", s.cfg.Git.RepoURL, interval)
// Initial sync on startup
go func() {
result := s.doSync()
s.logger.Printf("[INFO] [sync] Initial sync: %s", result.Message)
}()
go func() {
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-s.stopCh:
s.logger.Println("[INFO] [sync] Sync loop stopped")
return
case <-ticker.C:
result := s.doSync()
s.logger.Printf("[INFO] [sync] Periodic sync: %s", result.Message)
}
}
}()
}
// Stop terminates the periodic sync loop. Safe to call multiple times.
func (s *Syncer) Stop() {
s.stopOnce.Do(func() {
close(s.stopCh)
})
}
// TriggerSync performs an immediate sync. Returns the result.
// Debounce: refuses if last sync was less than 30 seconds ago.
func (s *Syncer) TriggerSync() SyncResult {
if s.cfg.Git.RepoURL == "" {
return SyncResult{OK: false, Message: "Git sync is disabled (no repo_url configured)"}
}
s.mu.Lock()
if s.syncing {
s.mu.Unlock()
return SyncResult{OK: false, Message: "Szinkronizálás már folyamatban"}
}
if time.Since(s.lastSync) < 30*time.Second {
s.mu.Unlock()
return SyncResult{OK: false, Message: "Túl gyakori szinkronizálás — várj 30 másodpercet"}
}
s.syncing = true
s.mu.Unlock()
return s.doSync()
}
// Status returns the current sync status.
func (s *Syncer) Status() SyncStatus {
s.mu.Lock()
defer s.mu.Unlock()
status := SyncStatus{
LastSync: s.lastSync,
Syncing: s.syncing,
}
if s.cfg.Git.RepoURL == "" {
status.LastStatus = "disabled"
} else if s.lastSync.IsZero() {
status.LastStatus = "never"
} else if s.lastErr != nil {
status.LastStatus = "error"
status.LastError = s.lastErr.Error()
} else {
status.LastStatus = "ok"
}
return status
}
// doSync performs the actual git clone/pull + file copy.
func (s *Syncer) doSync() SyncResult {
s.mu.Lock()
s.syncing = true
s.mu.Unlock()
defer func() {
s.mu.Lock()
s.syncing = false
s.mu.Unlock()
}()
result := SyncResult{OK: true}
s.logger.Printf("[INFO] [sync] Starting catalog sync")
// Step 1: Clone or pull
if err := s.gitCloneOrPull(); err != nil {
s.logger.Printf("[ERROR] [sync] Catalog sync failed: %v", err)
s.mu.Lock()
s.lastErr = err
s.lastSync = time.Now()
s.mu.Unlock()
return SyncResult{OK: false, Message: fmt.Sprintf("Git hiba: %v", err)}
}
// Step 2: Copy templates to stacks dir
newApps, updated, err := s.copyTemplates()
if err != nil {
s.logger.Printf("[ERROR] [sync] Catalog sync failed: %v", err)
s.mu.Lock()
s.lastErr = err
s.lastSync = time.Now()
s.mu.Unlock()
return SyncResult{OK: false, Message: fmt.Sprintf("Másolási hiba: %v", err)}
}
result.NewApps = newApps
result.Updated = updated
// Step 3: Trigger rescan if anything changed
if len(newApps) > 0 || len(updated) > 0 {
if err := s.rescanFn(); err != nil {
s.logger.Printf("[WARN] [sync] Rescan after sync failed: %v", err)
}
}
// Step 4: Inject missing deploy fields for updated stacks
if len(updated) > 0 && s.postSyncHook != nil {
if s.isDebug() {
s.logger.Printf("[DEBUG] [sync] Post-sync hook: triggering missing field injection for %d stack(s): %v", len(updated), updated)
}
s.postSyncHook(updated)
}
// Build message
parts := []string{}
if len(newApps) > 0 {
parts = append(parts, fmt.Sprintf("új: %s", strings.Join(newApps, ", ")))
}
if len(updated) > 0 {
parts = append(parts, fmt.Sprintf("frissítve: %s", strings.Join(updated, ", ")))
}
if len(parts) == 0 {
result.Message = "Sablonok naprakészek — nincs változás"
} else {
result.Message = "Sablonok frissítve — " + strings.Join(parts, "; ")
}
s.logger.Printf("[INFO] [sync] Catalog sync complete")
s.mu.Lock()
s.lastErr = nil
s.lastSync = time.Now()
s.mu.Unlock()
return result
}
// gitCloneOrPull clones the repo if not yet cloned, or pulls latest changes.
func (s *Syncer) gitCloneOrPull() error {
if err := os.MkdirAll(filepath.Dir(s.cacheDir), 0755); err != nil {
return fmt.Errorf("creating cache parent dir: %w", err)
}
gitDir := filepath.Join(s.cacheDir, ".git")
if _, err := os.Stat(gitDir); os.IsNotExist(err) {
// Clone
s.logger.Printf("[INFO] [sync] Cloning %s → %s", s.cfg.Git.RepoURL, s.cacheDir)
args := []string{"clone", "--depth", "1", "--branch", s.cfg.Git.Branch}
repoURL := s.buildRepoURL()
args = append(args, repoURL, s.cacheDir)
if s.isDebug() {
s.logger.Printf("[DEBUG] [sync] git clone URL: %s, branch: %s, cacheDir: %s", maskRepoURL(repoURL), s.cfg.Git.Branch, s.cacheDir)
}
return s.runGit(args...)
}
// Remove stale git lock files left behind by interrupted operations
s.removeGitLockFiles()
// Pull
s.logger.Printf("[INFO] [sync] Pulling latest from %s (branch: %s)", s.cfg.Git.RepoURL, s.cfg.Git.Branch)
if s.isDebug() {
s.logger.Printf("[DEBUG] [sync] git fetch --depth 1 origin %s in %s", s.cfg.Git.Branch, s.cacheDir)
}
if err := s.gitCmd(s.cacheDir, "fetch", "--depth", "1", "origin", s.cfg.Git.Branch); err != nil {
return fmt.Errorf("git fetch: %w", err)
}
if err := s.gitCmd(s.cacheDir, "reset", "--hard", "origin/"+s.cfg.Git.Branch); err != nil {
return fmt.Errorf("git reset: %w", err)
}
return nil
}
// removeGitLockFiles removes stale .git/*.lock files that may have been left
// behind if a previous git operation was interrupted (e.g. container restart).
// These lock files prevent all subsequent git operations from succeeding.
func (s *Syncer) removeGitLockFiles() {
gitDir := filepath.Join(s.cacheDir, ".git")
lockFiles := []string{"index.lock", "shallow.lock", "HEAD.lock"}
for _, name := range lockFiles {
lockPath := filepath.Join(gitDir, name)
if _, err := os.Stat(lockPath); err == nil {
s.logger.Printf("[WARN] [sync] Removing stale lock file: %s", lockPath)
if err := os.Remove(lockPath); err != nil {
s.logger.Printf("[WARN] [sync] Failed to remove lock file %s: %v", lockPath, err)
}
}
}
}
// buildRepoURL constructs the repo URL with optional auth credentials.
func (s *Syncer) buildRepoURL() string {
url := s.cfg.Git.RepoURL
if s.cfg.Git.Username != "" && s.cfg.Git.Token != "" {
// Inject credentials into HTTPS URL: https://user:token@host/path
url = strings.Replace(url, "https://", fmt.Sprintf("https://%s:%s@", s.cfg.Git.Username, s.cfg.Git.Token), 1)
}
return url
}
// copyTemplates copies docker-compose.yml and .felhom.yml from the catalog cache
// to the stacks directory. Never overwrites app.yaml.
func (s *Syncer) copyTemplates() (newApps []string, updated []string, err error) {
templatesDir := filepath.Join(s.cacheDir, "templates")
entries, err := os.ReadDir(templatesDir)
if err != nil {
return nil, nil, fmt.Errorf("reading templates dir: %w", err)
}
for _, entry := range entries {
if !entry.IsDir() {
continue
}
appName := entry.Name()
srcDir := filepath.Join(templatesDir, appName)
dstDir := filepath.Join(s.cfg.Paths.StacksDir, appName)
isNew := false
if _, err := os.Stat(dstDir); os.IsNotExist(err) {
isNew = true
if err := os.MkdirAll(dstDir, 0755); err != nil {
s.logger.Printf("[WARN] [sync] Failed to create stack dir %s: %v", dstDir, err)
continue
}
}
// Files to sync (only template files, never app.yaml)
syncFiles := []string{"docker-compose.yml", ".felhom.yml"}
anyChanged := false
for _, filename := range syncFiles {
src := filepath.Join(srcDir, filename)
dst := filepath.Join(dstDir, filename)
if _, err := os.Stat(src); os.IsNotExist(err) {
if s.isDebug() {
s.logger.Printf("[DEBUG] [sync] %s/%s: source not found, skipping", appName, filename)
}
continue
}
// v0.235.0 — the RENDER. `.felhom.yml` is always copied verbatim (it holds no image);
// only the compose file can be frozen. See renderSource for the whole table.
if filename == "docker-compose.yml" {
frozenSrc, skip := s.renderSource(appName, src)
if skip {
continue
}
src = frozenSrc
}
changed, err := copyIfChanged(src, dst)
if err != nil {
s.logger.Printf("[WARN] [sync] Failed to copy catalog file %s/%s: %v", appName, filename, err)
continue
}
if changed {
anyChanged = true
s.logger.Printf("[INFO] [sync] Updated %s/%s", appName, filename)
if s.isDebug() {
s.logFileHashes(appName, filename, src, dst)
}
} else if s.isDebug() {
s.logger.Printf("[DEBUG] [sync] %s/%s: hash match, skipped", appName, filename)
}
}
if isNew {
newApps = append(newApps, appName)
} else if anyChanged {
updated = append(updated, appName)
}
}
return newApps, updated, nil
}
// renderSource decides WHICH file becomes this app's live docker-compose.yml, and is the whole of
// the v0.235.0 ruling: "freeze the version, keep the fixes flowing" (operator, 2026-09-06).
//
// It returns the path to copy FROM, and whether to skip the app this cycle.
//
// ── THE TABLE, COMPLETE ──────────────────────────────────────────────────────────────────────
//
// not deployed / protected / no seam → the catalog template (today's behaviour)
// deployed, UNPINNED → the catalog template (today's behaviour) + one DEBUG
// deployed, pinned, catalog images == → the catalog template — FIXES FLOW, SELF-HEALING WORKS
// deployed, pinned, catalog images != → the STORED definition — the app is frozen WHOLE
// deployed, pinned, differ, none stored → the catalog template + one WARN
// mid-deploy → skip the compose file this cycle
//
// ── WHY THE FROZEN BRANCH WRITES A WHOLE FILE AND NEVER A SUBSTITUTION ───────────────────────
//
// The obvious-looking alternative — take the new template and put the old image refs back — creates
// a third state nobody chose: `wger 2.6` needs a full DB configuration the older template cannot
// supply, so a new template around an old image is broken in a way neither version is. When the
// catalog has moved, the WHOLE stored definition is used.
//
// ── AND WHY THIS IS NOT SIMPLY "SKIP DEPLOYED APPS" ──────────────────────────────────────────
//
// That was option B and it was rejected: it also stops health-check fixes, memory limits and new
// deploy fields reaching a deployed app, and it destroys the self-healing measured in
// SPIKE-app-update-2026-09-01 §3 — a hand-broken compose file repaired itself within 15 minutes.
// Both halves were worth keeping; only the version change was not.
func (s *Syncer) renderSource(appName, catalogSrc string) (src string, skip bool) {
if s.renderPlanFn == nil {
return catalogSrc, false // pre-v0.235.0 behaviour, byte for byte
}
plan := s.renderPlanFn(appName)
if !plan.Deployed || plan.Protected {
return catalogSrc, false
}
if plan.Deploying {
// Do not race an in-flight deploy for its own compose file.
if s.isDebug() {
s.logger.Printf("[DEBUG] [sync] %s: mid-deploy, compose file left alone this cycle", appName)
}
return "", true
}
if len(plan.Pinned) == 0 {
if s.isDebug() {
s.logger.Printf("[DEBUG] [sync] %s: deployed but UNPINNED — catalog copied verbatim (pre-v0.235.0 behaviour)", appName)
}
return catalogSrc, false
}
catalogImages, err := stacks.ParseComposeImages(catalogSrc)
if err != nil {
// CANNOT TELL whether the catalog has moved. Leave the app's file alone rather than guess in
// either direction — an unreadable catalog template must not be able to unfreeze an app.
s.logger.Printf("[WARN] [sync] %s: cannot read the catalog template's images (%v) — compose file left alone this cycle", appName, err)
return "", true
}
if samePin(plan.Pinned, catalogImages) {
// The catalog still offers what this app runs: everything else in the template is a FIX and
// is delivered, exactly as before v0.235.0. This branch is why the feature is not a freeze.
return catalogSrc, false
}
if plan.AppliedPath == "" {
// We cannot freeze what we do not have, and we must not invent it.
s.logger.Printf("[WARN] [sync] %s: the catalog has moved past this app's pinned version, but no stored definition exists — copying the catalog verbatim (pre-v0.235.0 behaviour). The app will take the new version on its next start.", appName)
return catalogSrc, false
}
// DEFENCE IN DEPTH, and this line was added because a test demanded it: the manager's
// RenderPlanFor already refuses to hand over a path whose file is missing or empty, but the
// syncer is what WRITES, and writing an empty compose file over a live app takes that app down.
// The cost of re-reading a small file once per app per cycle is nothing next to that.
if data, err := os.ReadFile(plan.AppliedPath); err != nil || len(strings.TrimSpace(string(data))) == 0 {
s.logger.Printf("[WARN] [sync] %s: the stored definition is missing or empty (%v) — copying the catalog verbatim rather than writing an empty compose file", appName, err)
return catalogSrc, false
}
if s.isDebug() {
s.logger.Printf("[DEBUG] [sync] %s: catalog has moved past the pin — rendering the stored applied definition", appName)
}
return plan.AppliedPath, false
}
// samePin compares a pin against a template's images. Local to the syncer so this package needs
// nothing from the manager beyond the plan it is handed.
func samePin(pinned, catalog map[string]string) bool {
if len(pinned) != len(catalog) {
return false
}
for svc, ref := range pinned {
if got, ok := catalog[svc]; !ok || got != ref {
return false
}
}
return true
}
// logFileHashes logs the source and destination file hashes for debugging.
func (s *Syncer) logFileHashes(appName, filename, src, dst string) {
srcData, err := os.ReadFile(src)
if err != nil {
return
}
srcHash := sha256.Sum256(srcData)
dstData, err := os.ReadFile(dst)
if err != nil {
s.logger.Printf("[DEBUG] [sync] %s/%s: src=%s, dst=new file", appName, filename, hex.EncodeToString(srcHash[:8]))
return
}
dstHash := sha256.Sum256(dstData)
s.logger.Printf("[DEBUG] [sync] %s/%s: src=%s, dst=%s (changed)", appName, filename, hex.EncodeToString(srcHash[:8]), hex.EncodeToString(dstHash[:8]))
}
// copyIfChanged copies src to dst only if the content differs.
// Returns true if the file was actually written.
func copyIfChanged(src, dst string) (bool, error) {
srcData, err := os.ReadFile(src)
if err != nil {
return false, fmt.Errorf("reading %s: %w", src, err)
}
// Check if dst exists and has same content
dstData, err := os.ReadFile(dst)
if err == nil {
srcHash := sha256.Sum256(srcData)
dstHash := sha256.Sum256(dstData)
if srcHash == dstHash {
return false, nil // No change
}
}
if err := os.WriteFile(dst, srcData, 0644); err != nil {
return false, fmt.Errorf("writing %s: %w", dst, err)
}
return true, nil
}
// runGit executes a git command with the given args under the standard per-command deadline.
func (s *Syncer) runGit(args ...string) error {
return s.gitCmd("", args...)
}
// gitCmd runs one git command in dir under a fresh gitCmdTimeout deadline (per-command,
// not per-sync — a multi-step sync gets a full budget for each subprocess).
func (s *Syncer) gitCmd(dir string, args ...string) error {
ctx, cancel := context.WithTimeout(context.Background(), gitCmdTimeout)
defer cancel()
return s.runGitInDir(ctx, dir, args...)
}
// runGitInDir executes a git command in the specified directory, killed when ctx expires.
func (s *Syncer) runGitInDir(ctx context.Context, dir string, args ...string) error {
cmd := exec.CommandContext(ctx, "git", args...)
if dir != "" {
cmd.Dir = dir
}
var stderr bytes.Buffer
cmd.Stdout = io.Discard
cmd.Stderr = &stderr
s.logger.Printf("[DEBUG] [sync] Running: git %s", maskRepoURL(strings.Join(args, " ")))
if err := cmd.Run(); err != nil {
if ctxErr := ctx.Err(); ctxErr != nil {
// Subprocess was killed by the deadline/cancellation — name it explicitly so the
// failure is diagnosable from SyncResult/logs (never a silent hang).
return fmt.Errorf("git %s: %v (subprocess killed at the %s deadline): %w\nstderr: %s",
maskRepoURL(strings.Join(args, " ")), ctxErr, gitCmdTimeout, err, stderr.String())
}
return fmt.Errorf("git %s: %w\nstderr: %s", maskRepoURL(strings.Join(args, " ")), err, stderr.String())
}
return nil
}