Files
felhom-controller/controller/internal/sync/sync.go
T
admin 0b349e2415 R-607 + R-615: the catalog sync follows a changed repo_url and rescans when the catalog moved
R-615: a changed git.repo_url was inert (the clone's stored origin was fetched for ever). The sync
now compares the configured URL with the clone's origin: another repository -> drop the cache and
clone it; the same repository with new credentials -> set-url.
R-607: a catalog move that changed no stack directory (a deployed, pinned app keeps its stored
definition) skipped the rescan, so CatalogImages and the update badge stayed stale, and the sync
said 'nincs változás'. The sync now notes the cache commit before and after, rescans when it moved,
says so, and returns catalog_moved.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_0159rPz1ZhFKsS53msqPYxtS
2026-10-05 21:55:37 +02:00

717 lines
27 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"`
// CatalogMoved (R-607): the catalog cache's commit changed in this sync, whether or not any stack
// directory did.
CatalogMoved bool `json:"catalog_moved,omitempty"`
}
// 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. R-607: note the cache's commit before and after — the stack-dir copies
// below are NOT the whole story: a deployed, pinned app keeps its stored definition (renderSource),
// so a catalog move can change nothing in the stack dirs and still change what the update badge
// must compare against (`CatalogImages`, read from the cache by the rescan).
headBefore := s.cacheHead()
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
headAfter := s.cacheHead()
catalogMoved := headAfter != "" && headAfter != headBefore
result.CatalogMoved = catalogMoved
if catalogMoved {
s.logger.Printf("[INFO] [sync] catalog moved %s -> %s (%d new, %d stack dir(s) changed)", shortHead(headBefore), shortHead(headAfter), len(newApps), len(updated))
}
// Step 3: Trigger rescan if anything changed — the catalog itself included (R-607: before, a move
// that changed no stack dir left `CatalogImages` stale until the next timed scan, and the badge
// could read „Naprakész" on an app that was behind).
if len(newApps) > 0 || len(updated) > 0 || catalogMoved {
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 && catalogMoved {
// R-607: the catalog DID change; only the installed apps' own files did not (each keeps the
// version it runs until the household updates it). „nincs változás" was false here.
result.Message = "Sablonok frissítve — a katalógus új változata betöltve; a telepített alkalmazások fájljai nem változtak"
} else 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()
// R-615: the clone remembers the repository it was made from. A changed `git.repo_url` used to be
// INERT — every later fetch went to the stored origin and reported success. Compare and follow:
// a different repository (credentials aside) → drop the cache and clone the new one; the same
// repository with different credentials (a rotated token) → point origin at the new URL.
if cur, err := s.gitOutput(s.cacheDir, "config", "--get", "remote.origin.url"); err != nil {
s.logger.Printf("[WARN] [sync] cannot read the catalog cache's origin (%v) — fetching from it as before", err)
} else if want := s.buildRepoURL(); cur != want {
if stripURLCreds(cur) != stripURLCreds(want) {
s.logger.Printf("[WARN] [sync] git.repo_url changed (cache was cloned from %s, config says %s) — re-cloning the catalog cache from the configured repository (R-615)", maskRepoURL(cur), maskRepoURL(want))
if err := os.RemoveAll(s.cacheDir); err != nil {
return fmt.Errorf("removing the catalog cache for a re-clone: %w", err)
}
return s.gitCloneOrPull()
}
s.logger.Printf("[INFO] [sync] catalog repository credentials changed — updating the cache's origin (R-615)")
if err := s.gitCmd(s.cacheDir, "remote", "set-url", "origin", want); err != nil {
return fmt.Errorf("git remote set-url: %w", err)
}
}
// 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
}
var changed bool
var err error
var rendered []byte // v0.269.0: the catalog compose with its TESTED digests (`09` §6.4 part 6)
if filename == "docker-compose.yml" && filepath.Dir(src) == srcDir {
raw, rerr := os.ReadFile(src)
if rerr != nil {
s.logger.Printf("[WARN] [sync] Failed to read catalog file %s/%s: %v", appName, filename, rerr)
continue
}
if s.renderPlanFn != nil && s.renderPlanFn(appName).Deployed {
// A DEPLOYED app keeps the digest it runs; only a guarded update moves it (night
// 2026-09-24 Part B, live finding). A fresh install takes the ladder's tested digest.
current, _ := os.ReadFile(dst)
rendered = stacks.CarryDigests(raw, current)
} else {
rendered = stacks.RenderWithLadderDigests(srcDir, raw)
}
changed, err = writeIfChanged(rendered, dst)
} else {
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 rendered != nil {
if serr := stacks.StoreAppliedDefinition(refreshAppliedIn, rendered); serr != nil {
s.logger.Printf("[WARN] [sync] %s: could not refresh the stored definition: %v", appName, serr)
}
} else 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]))
}
// writeIfChanged writes data to dst only if the content differs. Returns true if written.
func writeIfChanged(data []byte, dst string) (bool, error) {
if cur, err := os.ReadFile(dst); err == nil && sha256.Sum256(cur) == sha256.Sum256(data) {
return false, nil
}
if err := os.WriteFile(dst, data, 0644); err != nil {
return false, fmt.Errorf("writing %s: %w", dst, err)
}
return true, nil
}
// 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
}
// gitOutput runs one git command in dir and returns its trimmed stdout (R-615/R-607 reads).
func (s *Syncer) gitOutput(dir string, args ...string) (string, error) {
ctx, cancel := context.WithTimeout(context.Background(), gitCmdTimeout)
defer cancel()
cmd := exec.CommandContext(ctx, "git", args...)
cmd.Dir = dir
out, err := cmd.Output()
if err != nil {
return "", fmt.Errorf("git %s: %w", maskRepoURL(strings.Join(args, " ")), err)
}
return strings.TrimSpace(string(out)), nil
}
// cacheHead is the catalog cache's current commit, "" when there is no clone yet or it cannot be read.
func (s *Syncer) cacheHead() string {
if _, err := os.Stat(filepath.Join(s.cacheDir, ".git")); err != nil {
return ""
}
h, err := s.gitOutput(s.cacheDir, "rev-parse", "HEAD")
if err != nil {
return ""
}
return h
}
func shortHead(h string) string {
if h == "" {
return "(none)"
}
if len(h) > 12 {
return h[:12]
}
return h
}
// stripURLCreds removes a `user:token@` part, so two URLs naming the same repository compare equal.
func stripURLCreds(u string) string { return reURLCreds.ReplaceAllString(u, "${1}") }
// 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
}