This repository has no description
1package db
2
3import (
4 "database/sql"
5 "errors"
6 "fmt"
7
8 "github.com/bluesky-social/indigo/atproto/syntax"
9 "tangled.org/core/appview/models"
10)
11
12type stateTable struct {
13 name string
14 valueCol string
15}
16
17var (
18 issueStateTable = stateTable{name: "issue_states", valueCol: "state"}
19 pullStateTable = stateTable{name: "pull_states", valueCol: "status"}
20)
21
22func putStateRecord(tx *sql.Tx, t stateTable, rec models.StateRecord) (syntax.ATURI, error) {
23 var priorSubject string
24 err := tx.QueryRow(
25 fmt.Sprintf(`select subject from %s where did = ? and rkey = ?`, t.name),
26 rec.Did, rec.Rkey,
27 ).Scan(&priorSubject)
28 switch {
29 case errors.Is(err, sql.ErrNoRows):
30 priorSubject = ""
31 case err != nil:
32 return "", err
33 }
34
35 if _, err := tx.Exec(fmt.Sprintf(`
36 insert into %s (did, rkey, subject, %s, created_micros)
37 values (?, ?, ?, ?, ?)
38 on conflict(did, rkey) do update set
39 subject = excluded.subject,
40 %s = excluded.%s,
41 created_micros = excluded.created_micros
42 `, t.name, t.valueCol, t.valueCol, t.valueCol),
43 rec.Did, rec.Rkey, rec.Subject, string(rec.Value), rec.SortMicros); err != nil {
44 return "", err
45 }
46
47 if priorSubject != "" && syntax.ATURI(priorSubject) != rec.Subject {
48 return syntax.ATURI(priorSubject), nil
49 }
50 return "", nil
51}
52
53func deleteStateRecord(tx *sql.Tx, t stateTable, did, rkey string) (syntax.ATURI, error) {
54 var subject string
55 err := tx.QueryRow(
56 fmt.Sprintf(`select subject from %s where did = ? and rkey = ?`, t.name),
57 did, rkey,
58 ).Scan(&subject)
59 switch {
60 case errors.Is(err, sql.ErrNoRows):
61 return "", nil
62 case err != nil:
63 return "", err
64 }
65
66 if _, err := tx.Exec(
67 fmt.Sprintf(`delete from %s where did = ? and rkey = ?`, t.name),
68 did, rkey,
69 ); err != nil {
70 return "", err
71 }
72 return syntax.ATURI(subject), nil
73}
74
75func stateWinner(e Execer, t stateTable, subject syntax.ATURI) (models.StateValue, bool, error) {
76 var v string
77 err := e.QueryRow(fmt.Sprintf(`
78 select %s from %s
79 where subject = ?
80 order by created_micros desc, at_uri desc
81 limit 1
82 `, t.valueCol, t.name), subject).Scan(&v)
83 switch {
84 case errors.Is(err, sql.ErrNoRows):
85 return "", false, nil
86 case err != nil:
87 return "", false, err
88 }
89 return models.StateValue(v), true, nil
90}
91
92func PutIssueState(tx *sql.Tx, rec models.StateRecord) (syntax.ATURI, error) {
93 return putStateRecord(tx, issueStateTable, rec)
94}
95
96func PutPullStatus(tx *sql.Tx, rec models.StateRecord) (syntax.ATURI, error) {
97 return putStateRecord(tx, pullStateTable, rec)
98}
99
100type PendingStateRecord struct {
101 Did string
102 Rkey string
103 Nsid string
104 Subject syntax.ATURI
105 Record []byte
106}
107
108func ParkStateRecord(tx *sql.Tx, p PendingStateRecord) error {
109 _, err := tx.Exec(`
110 insert into pending_state_records (did, rkey, nsid, subject, record)
111 values (?, ?, ?, ?, ?)
112 on conflict(did, rkey, nsid) do update set
113 subject = excluded.subject,
114 record = excluded.record
115 `, p.Did, p.Rkey, p.Nsid, p.Subject, p.Record)
116 return err
117}
118
119func UnparkStateRecord(tx *sql.Tx, did, rkey, nsid string) error {
120 _, err := tx.Exec(
121 `delete from pending_state_records where did = ? and rkey = ? and nsid = ?`,
122 did, rkey, nsid,
123 )
124 return err
125}
126
127func distinctPendingSubjects(e Execer, nsid string) ([]syntax.ATURI, error) {
128 query := `select distinct subject from pending_state_records`
129 var args []any
130 if nsid != "" {
131 query += ` where nsid = ?`
132 args = append(args, nsid)
133 }
134 query += ` order by subject asc`
135
136 rows, err := e.Query(query, args...)
137 if err != nil {
138 return nil, err
139 }
140 defer rows.Close()
141
142 var subjects []syntax.ATURI
143 for rows.Next() {
144 var s string
145 if err := rows.Scan(&s); err != nil {
146 return nil, err
147 }
148 subjects = append(subjects, syntax.ATURI(s))
149 }
150 return subjects, rows.Err()
151}
152
153func DistinctPendingStateSubjects(e Execer) ([]syntax.ATURI, error) {
154 return distinctPendingSubjects(e, "")
155}
156
157func PendingStateSubjectsForNsid(e Execer, nsid string) ([]syntax.ATURI, error) {
158 return distinctPendingSubjects(e, nsid)
159}
160
161func EvictStalePendingStateRecords(e Execer, before string) (int64, error) {
162 res, err := e.Exec(
163 `delete from pending_state_records where created < ?`,
164 before,
165 )
166 if err != nil {
167 return 0, err
168 }
169 return res.RowsAffected()
170}
171
172func PendingStateRecordsForSubject(e Execer, subject syntax.ATURI) ([]PendingStateRecord, error) {
173 rows, err := e.Query(
174 `select did, rkey, nsid, subject, record from pending_state_records where subject = ? order by id asc`,
175 subject,
176 )
177 if err != nil {
178 return nil, err
179 }
180 defer rows.Close()
181
182 var pending []PendingStateRecord
183 for rows.Next() {
184 var p PendingStateRecord
185 var subj string
186 if err := rows.Scan(&p.Did, &p.Rkey, &p.Nsid, &subj, &p.Record); err != nil {
187 return nil, err
188 }
189 p.Subject = syntax.ATURI(subj)
190 pending = append(pending, p)
191 }
192 return pending, rows.Err()
193}
194
195func DeleteIssueState(tx *sql.Tx, did, rkey string) (syntax.ATURI, error) {
196 return deleteStateRecord(tx, issueStateTable, did, rkey)
197}
198
199func DeletePullStatus(tx *sql.Tx, did, rkey string) (syntax.ATURI, error) {
200 return deleteStateRecord(tx, pullStateTable, did, rkey)
201}
202
203func setIssueOpen(tx *sql.Tx, subject syntax.ATURI, open bool) error {
204 v := 0
205 if open {
206 v = 1
207 }
208 _, err := tx.Exec(`update issues set open = ? where at_uri = ?`, v, subject)
209 return err
210}
211
212func applyIssueState(tx *sql.Tx, subject syntax.ATURI, resetWhenEmpty bool) error {
213 winner, ok, err := stateWinner(tx, issueStateTable, subject)
214 if err != nil {
215 return err
216 }
217 if !ok {
218 if resetWhenEmpty {
219 return setIssueOpen(tx, subject, true)
220 }
221 return nil
222 }
223 return setIssueOpen(tx, subject, winner == models.StateOpen)
224}
225
226func ResolveIssueState(tx *sql.Tx, subject syntax.ATURI) error {
227 return applyIssueState(tx, subject, false)
228}
229
230func RecomputeIssueState(tx *sql.Tx, subject syntax.ATURI) error {
231 return applyIssueState(tx, subject, true)
232}
233
234func pullStateFromValue(v models.StateValue) models.PullState {
235 switch v {
236 case models.StateMerged:
237 return models.PullMerged
238 case models.StateClosed:
239 return models.PullClosed
240 default:
241 return models.PullOpen
242 }
243}
244
245func setPullState(tx *sql.Tx, subject syntax.ATURI, st models.PullState) error {
246 _, err := tx.Exec(
247 `update pulls set state = ? where at_uri = ? and state <> ?`,
248 st, subject, models.PullAbandoned,
249 )
250 return err
251}
252
253func applyPullStatus(tx *sql.Tx, subject syntax.ATURI, resetWhenEmpty bool) error {
254 winner, ok, err := stateWinner(tx, pullStateTable, subject)
255 if err != nil {
256 return err
257 }
258 if !ok {
259 if resetWhenEmpty {
260 return setPullState(tx, subject, models.PullOpen)
261 }
262 return nil
263 }
264 return setPullState(tx, subject, pullStateFromValue(winner))
265}
266
267func ResolvePullStatus(tx *sql.Tx, subject syntax.ATURI) error {
268 return applyPullStatus(tx, subject, false)
269}
270
271func RecomputePullStatus(tx *sql.Tx, subject syntax.ATURI) error {
272 return applyPullStatus(tx, subject, true)
273}