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