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 DistinctPendingStateSubjects(e Execer) ([]syntax.ATURI, error) {
128 rows, err := e.Query(`select distinct subject from pending_state_records order by subject asc`)
129 if err != nil {
130 return nil, err
131 }
132 defer rows.Close()
133
134 var subjects []syntax.ATURI
135 for rows.Next() {
136 var s string
137 if err := rows.Scan(&s); err != nil {
138 return nil, err
139 }
140 subjects = append(subjects, syntax.ATURI(s))
141 }
142 return subjects, rows.Err()
143}
144
145func EvictStalePendingStateRecords(e Execer, before string) (int64, error) {
146 res, err := e.Exec(`
147 delete from pending_state_records
148 where created < ?
149 and not exists (select 1 from issues where at_uri = pending_state_records.subject)
150 and not exists (select 1 from pulls where at_uri = pending_state_records.subject)
151 `, before)
152 if err != nil {
153 return 0, err
154 }
155 return res.RowsAffected()
156}
157
158func PendingStateRecordsForSubject(e Execer, subject syntax.ATURI) ([]PendingStateRecord, error) {
159 rows, err := e.Query(
160 `select did, rkey, nsid, subject, record from pending_state_records where subject = ? order by id asc`,
161 subject,
162 )
163 if err != nil {
164 return nil, err
165 }
166 defer rows.Close()
167
168 var pending []PendingStateRecord
169 for rows.Next() {
170 var p PendingStateRecord
171 var subj string
172 if err := rows.Scan(&p.Did, &p.Rkey, &p.Nsid, &subj, &p.Record); err != nil {
173 return nil, err
174 }
175 p.Subject = syntax.ATURI(subj)
176 pending = append(pending, p)
177 }
178 return pending, rows.Err()
179}
180
181func DeleteIssueState(tx *sql.Tx, did, rkey string) (syntax.ATURI, error) {
182 return deleteStateRecord(tx, issueStateTable, did, rkey)
183}
184
185func DeletePullStatus(tx *sql.Tx, did, rkey string) (syntax.ATURI, error) {
186 return deleteStateRecord(tx, pullStateTable, did, rkey)
187}
188
189func setIssueOpen(tx *sql.Tx, subject syntax.ATURI, open bool) error {
190 v := 0
191 if open {
192 v = 1
193 }
194 _, err := tx.Exec(`update issues set open = ? where at_uri = ?`, v, subject)
195 return err
196}
197
198func applyIssueState(tx *sql.Tx, subject syntax.ATURI, resetWhenEmpty bool) error {
199 winner, ok, err := stateWinner(tx, issueStateTable, subject)
200 if err != nil {
201 return err
202 }
203 if !ok {
204 if resetWhenEmpty {
205 return setIssueOpen(tx, subject, true)
206 }
207 return nil
208 }
209 return setIssueOpen(tx, subject, winner == models.StateOpen)
210}
211
212func ResolveIssueState(tx *sql.Tx, subject syntax.ATURI) error {
213 return applyIssueState(tx, subject, false)
214}
215
216func RecomputeIssueState(tx *sql.Tx, subject syntax.ATURI) error {
217 return applyIssueState(tx, subject, true)
218}
219
220func pullStateFromValue(v models.StateValue) models.PullState {
221 switch v {
222 case models.StateMerged:
223 return models.PullMerged
224 case models.StateClosed:
225 return models.PullClosed
226 default:
227 return models.PullOpen
228 }
229}
230
231func setPullState(tx *sql.Tx, subject syntax.ATURI, st models.PullState) error {
232 _, err := tx.Exec(
233 `update pulls set state = ? where at_uri = ? and state <> ?`,
234 st, subject, models.PullAbandoned,
235 )
236 return err
237}
238
239func applyPullStatus(tx *sql.Tx, subject syntax.ATURI, resetWhenEmpty bool) error {
240 winner, ok, err := stateWinner(tx, pullStateTable, subject)
241 if err != nil {
242 return err
243 }
244 if !ok {
245 if resetWhenEmpty {
246 return setPullState(tx, subject, models.PullOpen)
247 }
248 return nil
249 }
250 return setPullState(tx, subject, pullStateFromValue(winner))
251}
252
253func ResolvePullStatus(tx *sql.Tx, subject syntax.ATURI) error {
254 return applyPullStatus(tx, subject, false)
255}
256
257func RecomputePullStatus(tx *sql.Tx, subject syntax.ATURI) error {
258 return applyPullStatus(tx, subject, true)
259}