This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / appview / db / entity_state.go
6.8 kB 273 lines
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}