internal/store/mirrors.go

v1.41.0
gitbay/internal/store/mirrors.go history · blame · raw

146 lines · 4056 bytes

  1package store
  2
  3import (
  4	"errors"
  5	"fmt"
  6)
  7
  8// ErrExists marks unique-constraint refusals callers turn into messages.
  9var ErrExists = errors.New("already exists")
 10
 11// Mirror propagates refs to (push) or from (pull) a foreign remote. The
 12// token is stored server-side — unlike import, mirroring is recurring —
 13// and must never be echoed back in listings.
 14type Mirror struct {
 15	ID        int64
 16	RepoID    int64
 17	Direction string // push | pull
 18	URL       string
 19	Username  string
 20	Token     string
 21	Dirty     bool
 22	LastSync  string
 23	LastError string
 24}
 25
 26// AddMirror stores the mirror, then seals its token under the new row's
 27// id in the same transaction.
 28func (s *Store) AddMirror(repoID int64, direction, url, username, token string) (int64, error) {
 29	tx, err := s.DB.Begin()
 30	if err != nil {
 31		return 0, err
 32	}
 33	defer tx.Rollback()
 34	res, err := tx.Exec(
 35		"INSERT INTO mirrors (repo_id, direction, url, username, token) VALUES (?, ?, ?, ?, '')",
 36		repoID, direction, url, username)
 37	if err != nil {
 38		if isUniqueErr(err) {
 39			return 0, ErrExists
 40		}
 41		return 0, err
 42	}
 43	id, err := res.LastInsertId()
 44	if err != nil {
 45		return 0, err
 46	}
 47	if token != "" {
 48		sealed, err := s.sealValue(mirrorAAD(id), token)
 49		if err != nil {
 50			return 0, err
 51		}
 52		if _, err := tx.Exec("UPDATE mirrors SET token = ? WHERE id = ?", sealed, id); err != nil {
 53			return 0, err
 54		}
 55	}
 56	return id, tx.Commit()
 57}
 58
 59const mirrorSelect = `
 60	SELECT id, repo_id, direction, url, username, token, dirty, last_sync, last_error
 61	FROM mirrors`
 62
 63func (s *Store) scanMirror(row interface{ Scan(...any) error }) (Mirror, error) {
 64	var m Mirror
 65	if err := row.Scan(&m.ID, &m.RepoID, &m.Direction, &m.URL, &m.Username, &m.Token,
 66		&m.Dirty, &m.LastSync, &m.LastError); err != nil {
 67		return m, err
 68	}
 69	var err error
 70	if m.Token, err = s.openValue(mirrorAAD(m.ID), m.Token); err != nil {
 71		return m, fmt.Errorf("mirror %d: %w", m.ID, err)
 72	}
 73	return m, nil
 74}
 75
 76func (s *Store) mirrorQuery(q string, args ...any) ([]Mirror, error) {
 77	rows, err := s.DB.Query(q, args...)
 78	if err != nil {
 79		return nil, err
 80	}
 81	defer rows.Close()
 82	var out []Mirror
 83	for rows.Next() {
 84		m, err := s.scanMirror(rows)
 85		if err != nil {
 86			return nil, err
 87		}
 88		out = append(out, m)
 89	}
 90	return out, rows.Err()
 91}
 92
 93func (s *Store) ListMirrors(repoID int64) ([]Mirror, error) {
 94	return s.mirrorQuery(mirrorSelect+" WHERE repo_id = ? ORDER BY id", repoID)
 95}
 96
 97// DueMirrors returns mirrors needing a sync: anything dirty, plus pull
 98// mirrors whose last sync is older than intervalSeconds.
 99func (s *Store) DueMirrors(intervalSeconds int) ([]Mirror, error) {
100	return s.mirrorQuery(mirrorSelect+`
101		WHERE dirty = 1
102		   OR (direction = 'pull' AND (last_sync = ''
103		       OR strftime('%s','now') - strftime('%s', last_sync) > ?))
104		ORDER BY id`, intervalSeconds)
105}
106
107func (s *Store) RemoveMirror(repoID, id int64) error {
108	res, err := s.DB.Exec("DELETE FROM mirrors WHERE repo_id = ? AND id = ?", repoID, id)
109	if err != nil {
110		return err
111	}
112	if n, _ := res.RowsAffected(); n == 0 {
113		return ErrNotFound
114	}
115	return nil
116}
117
118// MarkMirrorsDirty schedules a sync. An empty direction marks both.
119func (s *Store) MarkMirrorsDirty(repoID int64, direction string) error {
120	q := "UPDATE mirrors SET dirty = 1 WHERE repo_id = ?"
121	args := []any{repoID}
122	if direction != "" {
123		q += " AND direction = ?"
124		args = append(args, direction)
125	}
126	_, err := s.DB.Exec(q, args...)
127	return err
128}
129
130// SetMirrorResult records a sync outcome and clears the dirty flag.
131func (s *Store) SetMirrorResult(id int64, syncErr string) error {
132	_, err := s.DB.Exec(`
133		UPDATE mirrors SET dirty = 0, last_error = ?,
134		       last_sync = strftime('%Y-%m-%dT%H:%M:%fZ','now')
135		WHERE id = ?`, syncErr, id)
136	return err
137}
138
139// PullMirrored reports whether the repo has a pull mirror, which makes it
140// read-only locally: its refs belong to the upstream.
141func (s *Store) PullMirrored(repoID int64) (bool, error) {
142	var n int
143	err := s.DB.QueryRow(
144		"SELECT COUNT(*) FROM mirrors WHERE repo_id = ? AND direction = 'pull'", repoID).Scan(&n)
145	return n > 0, err
146}