internal/mirror/mirror.go

61564b7c32807deb349e28f2b5c6909cb4143870
gitbay/internal/mirror/mirror.go history · blame · raw

116 lines · 2989 bytes

  1// Package mirror synchronizes repositories with foreign remotes: push
  2// mirrors propagate local refs outward after each receive, pull mirrors
  3// keep a local copy fresh from an upstream. Sync runs in a background
  4// worker, never in the push path; outcomes are recorded per mirror so
  5// `repo mirror list` shows failure states like webhook deliveries do.
  6package mirror
  7
  8import (
  9	"context"
 10	"fmt"
 11	"log/slog"
 12	"os"
 13	"os/exec"
 14	"path/filepath"
 15	"time"
 16
 17	"gitbay.org/gitbay/internal/config"
 18	"gitbay.org/gitbay/internal/control"
 19	"gitbay.org/gitbay/internal/store"
 20	"gitbay.org/gitbay/internal/toolpath"
 21)
 22
 23const askpassScript = `#!/bin/sh
 24case "$1" in
 25  Username*) echo "${GITBAY_MIRROR_USER}" ;;
 26  *)         echo "${GITBAY_MIRROR_TOKEN}" ;;
 27esac
 28`
 29
 30type Worker struct {
 31	St   *store.Store
 32	Cfg  config.Config
 33	Tick time.Duration
 34}
 35
 36func New(st *store.Store, cfg config.Config) *Worker {
 37	tick := 10 * time.Second
 38	if v := os.Getenv("GITBAY_MIRROR_TICK"); v != "" {
 39		if d, err := time.ParseDuration(v); err == nil {
 40			tick = d
 41		}
 42	}
 43	return &Worker{St: st, Cfg: cfg, Tick: tick}
 44}
 45
 46func (w *Worker) Run(ctx context.Context) {
 47	t := time.NewTicker(w.Tick)
 48	defer t.Stop()
 49	for {
 50		select {
 51		case <-ctx.Done():
 52			return
 53		case <-t.C:
 54			w.sweep()
 55		}
 56	}
 57}
 58
 59func (w *Worker) sweep() {
 60	interval := w.Cfg.Mirrors.PullIntervalMinutes * 60
 61	due, err := w.St.DueMirrors(interval)
 62	if err != nil {
 63		slog.Error("mirror: listing due", "err", err)
 64		return
 65	}
 66	for _, m := range due {
 67		if err := w.sync(m); err != nil {
 68			slog.Warn("mirror sync failed", "mirror", m.ID, "url", m.URL, "err", err)
 69			w.St.SetMirrorResult(m.ID, err.Error())
 70		} else {
 71			w.St.SetMirrorResult(m.ID, "")
 72		}
 73	}
 74}
 75
 76func (w *Worker) sync(m store.Mirror) error {
 77	repo, err := w.St.RepoByID(m.RepoID)
 78	if err != nil {
 79		return err
 80	}
 81	dir := control.RepoDir(w.Cfg.Server.Root, repo.OwnerName, repo.Name)
 82
 83	env := []string{"GIT_TERMINAL_PROMPT=0", "HOME=" + w.Cfg.Server.Root}
 84	if m.Token != "" {
 85		askpass := filepath.Join(w.Cfg.Server.Root, "mirror-askpass.sh")
 86		if err := os.WriteFile(askpass, []byte(askpassScript), 0o700); err != nil {
 87			return err
 88		}
 89		user := m.Username
 90		if user == "" {
 91			user = "x-access-token"
 92		}
 93		env = append(env,
 94			"GIT_ASKPASS="+askpass,
 95			"GITBAY_MIRROR_USER="+user,
 96			"GITBAY_MIRROR_TOKEN="+m.Token)
 97	}
 98
 99	ctx, cancel := context.WithTimeout(context.Background(), 10*time.Minute)
100	defer cancel()
101	var args []string
102	if m.Direction == "push" {
103		// Branches and tags only: internal refs (merge-requests) stay home.
104		args = []string{"-C", dir, "push", "--prune", m.URL,
105			"+refs/heads/*:refs/heads/*", "+refs/tags/*:refs/tags/*"}
106	} else {
107		args = []string{"-C", dir, "fetch", "--prune", m.URL,
108			"+refs/heads/*:refs/heads/*", "+refs/tags/*:refs/tags/*"}
109	}
110	cmd := exec.CommandContext(ctx, toolpath.Look("git"), args...)
111	cmd.Env = env
112	if out, err := cmd.CombinedOutput(); err != nil {
113		return fmt.Errorf("git %s: %v: %.300s", m.Direction, err, out)
114	}
115	return nil
116}