2a56f557d0
gates / gates (push) Successful in 14s
Found by the LIVE validation on demo-hp, not by review. Scenario A passed - a non-image catalog change reached the pinned app on the real 15-minute cycle - and that is exactly what exposed the gap: the stored applied-compose.yml is written when the PIN is written, so the fix landed in the live compose file and not in the store. The first time the catalog then moved a version, the freeze would have rendered the pre-fix definition and reverted every fix delivered since - silently undoing the half of the operator's ruling that says fixes keep flowing. The equal-images branch now refreshes the store as it delivers. The images cannot move in that branch by construction, so no version moves and no intent is rewritten. RenderPlan gains StackDir so the syncer can write it. TestFixRefreshesTheStoredDefinition asserts both halves: the fix reaches the store, and it survives the freeze that follows.
605 lines
22 KiB
Go
605 lines
22 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.
|
|
refreshAppliedIn := ""
|
|
if filename == "docker-compose.yml" {
|
|
renderedSrc, skip, refresh := s.renderSource(appName, src)
|
|
if skip {
|
|
continue
|
|
}
|
|
src, refreshAppliedIn = renderedSrc, refresh
|
|
}
|
|
|
|
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)
|
|
// The fix we just delivered becomes part of what this app is pinned TO. Without
|
|
// this the stored definition stays as it was when the pin was written, and the
|
|
// first time the catalog moves the freeze reverts every fix delivered since —
|
|
// silently undoing the half of the ruling that says fixes keep flowing.
|
|
// The IMAGES are unchanged here by construction (this branch only runs when the
|
|
// catalog's images equal the pin), so no version moves and no intent is rewritten.
|
|
if refreshAppliedIn != "" {
|
|
if data, rerr := os.ReadFile(src); rerr != nil {
|
|
s.logger.Printf("[WARN] [sync] %s: could not re-read the template to refresh its stored definition: %v", appName, rerr)
|
|
} else if serr := stacks.StoreAppliedDefinition(refreshAppliedIn, data); serr != nil {
|
|
s.logger.Printf("[WARN] [sync] %s: could not refresh the stored definition: %v", appName, serr)
|
|
} else if s.isDebug() {
|
|
s.logger.Printf("[DEBUG] [sync] %s: stored definition refreshed with the delivered fix", appName)
|
|
}
|
|
}
|
|
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, refreshAppliedIn string) {
|
|
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.
|
|
// The delivered fix also REFRESHES the stored definition — see the call site.
|
|
return catalogSrc, false, plan.StackDir
|
|
}
|
|
|
|
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
|
|
}
|