This repository has no description
1package db
2
3import (
4 "database/sql"
5 "fmt"
6
7 "tangled.org/core/eventstream"
8 "tangled.org/core/notifier"
9)
10
11// enough to rebuild the fencing token and workflow identity after a restart
12type MillLease struct {
13 LeaseID string
14 NodeID string
15 Epoch string
16 Engine string
17 Knot string
18 Rkey string
19 Workflow string
20 State string
21}
22
23type ExecutorCursor struct {
24 NodeID string
25 Epoch string
26 AckedSeqno uint64
27}
28
29type OutboxRow struct {
30 Epoch string
31 Seqno uint64
32 Payload []byte
33 ByteSize int64
34 Control bool
35}
36type OutboxDeletion struct {
37 Rows int64
38 Bytes int64
39}
40
41func (d *DB) SaveMillLease(l MillLease) error {
42 _, err := d.Exec(
43 `insert into mill_leases (
44 lease_id, node_id, epoch, engine, knot, rkey, workflow, state
45 ) values (?, ?, ?, ?, ?, ?, ?, ?)
46 on conflict(lease_id) do update set state = excluded.state`,
47 l.LeaseID, l.NodeID, l.Epoch, l.Engine, l.Knot, l.Rkey, l.Workflow, l.State,
48 )
49 return err
50}
51
52func (d *DB) DeleteMillLease(leaseID string) error {
53 _, err := d.Exec(`delete from mill_leases where lease_id = ?`, leaseID)
54 return err
55}
56
57func (d *DB) ListMillLeases() ([]MillLease, error) {
58 rows, err := d.Query(`
59 select lease_id, node_id, epoch, engine, knot, rkey, workflow, state
60 from mill_leases
61 `)
62 if err != nil {
63 return nil, err
64 }
65 defer rows.Close()
66
67 var leases []MillLease
68 for rows.Next() {
69 var l MillLease
70 if err := rows.Scan(
71 &l.LeaseID, &l.NodeID, &l.Epoch, &l.Engine, &l.Knot, &l.Rkey, &l.Workflow, &l.State,
72 ); err != nil {
73 return nil, err
74 }
75 leases = append(leases, l)
76 }
77 return leases, rows.Err()
78}
79
80func (d *DB) ListExecutorCursors() ([]ExecutorCursor, error) {
81 rows, err := d.Query(`select node_id, epoch, acked_seqno from mill_executor_cursors`)
82 if err != nil {
83 return nil, err
84 }
85 defer rows.Close()
86
87 var cursors []ExecutorCursor
88 for rows.Next() {
89 var c ExecutorCursor
90 if err := rows.Scan(&c.NodeID, &c.Epoch, &c.AckedSeqno); err != nil {
91 return nil, err
92 }
93 cursors = append(cursors, c)
94 }
95 return cursors, rows.Err()
96}
97
98func (d *DB) SetOutboxEpoch(epoch string) error {
99 tx, err := d.Begin()
100 if err != nil {
101 return err
102 }
103 defer tx.Rollback()
104
105 if _, err := tx.Exec(`delete from mill_outbox_rows`); err != nil {
106 return err
107 }
108 if _, err := tx.Exec(`delete from mill_outbox_state`); err != nil {
109 return err
110 }
111 if _, err := tx.Exec(`insert into mill_outbox_state (epoch, next_seqno) values (?, 1)`, epoch); err != nil {
112 return err
113 }
114 return tx.Commit()
115}
116
117func (d *DB) AppendOutboxRow(payload []byte, control bool) (uint64, error) {
118 tx, err := d.Begin()
119 if err != nil {
120 return 0, err
121 }
122 defer tx.Rollback()
123
124 var epoch string
125 var nextSeqno uint64
126 err = tx.QueryRow(`select epoch, next_seqno from mill_outbox_state limit 1`).Scan(&epoch, &nextSeqno)
127 if err == sql.ErrNoRows {
128 return 0, fmt.Errorf("no outbox epoch set")
129 } else if err != nil {
130 return 0, err
131 }
132
133 byteSize := int64(len(payload))
134 controlVal := 0
135 if control {
136 controlVal = 1
137 }
138
139 if _, err := tx.Exec(
140 `insert into mill_outbox_rows (epoch, seqno, payload, byte_size, control) values (?, ?, ?, ?, ?)`,
141 epoch, nextSeqno, payload, byteSize, controlVal,
142 ); err != nil {
143 return 0, err
144 }
145
146 if _, err := tx.Exec(
147 `update mill_outbox_state set next_seqno = ? where epoch = ?`,
148 nextSeqno+1, epoch,
149 ); err != nil {
150 return 0, err
151 }
152
153 if err := tx.Commit(); err != nil {
154 return 0, err
155 }
156 return nextSeqno, nil
157}
158
159func (d *DB) DeleteOutboxPrefix(ackedSeqno uint64) (OutboxDeletion, error) {
160 tx, err := d.Begin()
161 if err != nil {
162 return OutboxDeletion{}, err
163 }
164 defer tx.Rollback()
165
166 var epoch string
167 err = tx.QueryRow(`select epoch from mill_outbox_state limit 1`).Scan(&epoch)
168 if err == sql.ErrNoRows {
169 return OutboxDeletion{}, nil
170 }
171 if err != nil {
172 return OutboxDeletion{}, err
173 }
174
175 var deleted OutboxDeletion
176 if err := tx.QueryRow(`
177 select count(*),
178 coalesce(sum(byte_size), 0)
179 from mill_outbox_rows
180 where epoch = ? and seqno <= ?
181 `, epoch, ackedSeqno).Scan(&deleted.Rows, &deleted.Bytes); err != nil {
182 return OutboxDeletion{}, err
183 }
184 if _, err := tx.Exec(
185 `delete from mill_outbox_rows where epoch = ? and seqno <= ?`,
186 epoch, ackedSeqno,
187 ); err != nil {
188 return OutboxDeletion{}, err
189 }
190 if err := tx.Commit(); err != nil {
191 return OutboxDeletion{}, err
192 }
193 return deleted, nil
194}
195
196func (d *DB) ListOutboxRows() ([]OutboxRow, error) {
197 rows, err := d.Query(`select epoch, seqno, payload, byte_size, control from mill_outbox_rows order by seqno`)
198 if err != nil {
199 return nil, err
200 }
201 defer rows.Close()
202
203 var out []OutboxRow
204 for rows.Next() {
205 var r OutboxRow
206 var controlVal int
207 if err := rows.Scan(&r.Epoch, &r.Seqno, &r.Payload, &r.ByteSize, &controlVal); err != nil {
208 return nil, err
209 }
210 r.Control = (controlVal != 0)
211 out = append(out, r)
212 }
213 return out, rows.Err()
214}
215func (d *DB) ListOutboxRowsAfter(seqno uint64, limit int) ([]OutboxRow, error) {
216 rows, err := d.Query(`
217 select epoch, seqno, payload, byte_size, control
218 from mill_outbox_rows
219 where epoch = (select epoch from mill_outbox_state limit 1)
220 and seqno > ?
221 order by seqno
222 limit ?
223 `, seqno, limit)
224 if err != nil {
225 return nil, err
226 }
227 defer rows.Close()
228
229 var out []OutboxRow
230 for rows.Next() {
231 var row OutboxRow
232 var control int
233 if err := rows.Scan(&row.Epoch, &row.Seqno, &row.Payload, &row.ByteSize, &control); err != nil {
234 return nil, err
235 }
236 row.Control = control != 0
237 out = append(out, row)
238 }
239 return out, rows.Err()
240}
241
242func (d *DB) GetOutboxState() (string, uint64, error) {
243 var epoch string
244 var nextSeqno uint64
245 err := d.QueryRow(`select epoch, next_seqno from mill_outbox_state limit 1`).Scan(&epoch, &nextSeqno)
246 if err == sql.ErrNoRows {
247 return "", 0, nil
248 }
249 return epoch, nextSeqno, err
250}
251
252type EventBatchTx struct {
253 tx *sql.Tx
254 db *DB
255}
256
257func (tx *EventBatchTx) InsertStatusEvent(pipelineAtUri, workflow, status string, workflowError *string, exitCode *int64) error {
258 event, err := statusEvent(pipelineAtUri, workflow, status, workflowError, exitCode)
259 if err != nil {
260 return err
261 }
262 return eventstream.Insert(tx.tx, event, nil)
263}
264
265func (tx *EventBatchTx) DeleteLease(leaseID string) error {
266 _, err := tx.tx.Exec(`delete from mill_leases where lease_id = ?`, leaseID)
267 return err
268}
269
270func (tx *EventBatchTx) AdvanceCursor(nodeID, epoch string, seqno uint64) error {
271 _, err := tx.tx.Exec(
272 `insert into mill_executor_cursors (node_id, epoch, acked_seqno) values (?, ?, ?)
273 on conflict(node_id, epoch) do update set acked_seqno = excluded.acked_seqno`,
274 nodeID, epoch, seqno,
275 )
276 return err
277}
278
279func (tx *EventBatchTx) InsertArtifactRef(leaseID, workflow, ref, hash string) error {
280 _, err := tx.tx.Exec(
281 `insert into mill_artifacts (lease_id, workflow, ref, hash)
282 values (?, ?, ?, ?)`,
283 leaseID, workflow, ref, hash,
284 )
285 return err
286}
287
288type PendingArtifact struct {
289 LeaseID string
290 Workflow string
291 Status string
292 Error string
293 ExitCode int64
294 Ref string
295 Hash string
296}
297
298func (d *DB) SavePendingArtifact(leaseID, workflow, status, errStr string, exitCode int64, ref, hash string) error {
299 _, err := d.Exec(
300 `insert into executor_pending_artifacts (lease_id, workflow, status, error, exit_code, ref, hash)
301 values (?, ?, ?, ?, ?, ?, ?)
302 on conflict(lease_id) do update set
303 workflow = excluded.workflow,
304 status = excluded.status,
305 error = excluded.error,
306 exit_code = excluded.exit_code,
307 ref = excluded.ref,
308 hash = excluded.hash`,
309 leaseID, workflow, status, errStr, exitCode, ref, hash,
310 )
311 return err
312}
313
314func (d *DB) RemovePendingArtifact(leaseID string) error {
315 _, err := d.Exec(`delete from executor_pending_artifacts where lease_id = ?`, leaseID)
316 return err
317}
318
319func (d *DB) ListPendingArtifacts() ([]PendingArtifact, error) {
320 rows, err := d.Query(`select lease_id, workflow, status, error, exit_code, ref, hash from executor_pending_artifacts`)
321 if err != nil {
322 return nil, err
323 }
324 defer rows.Close()
325
326 var res []PendingArtifact
327 for rows.Next() {
328 var p PendingArtifact
329 if err := rows.Scan(&p.LeaseID, &p.Workflow, &p.Status, &p.Error, &p.ExitCode, &p.Ref, &p.Hash); err != nil {
330 return nil, err
331 }
332 res = append(res, p)
333 }
334 return res, rows.Err()
335}
336
337func (d *DB) ApplyEventBatch(n *notifier.Notifier, fn func(tx *EventBatchTx) error) error {
338 tx, err := d.Begin()
339 if err != nil {
340 return err
341 }
342 defer tx.Rollback()
343
344 batchTx := &EventBatchTx{
345 tx: tx,
346 db: d,
347 }
348
349 if err := fn(batchTx); err != nil {
350 return err
351 }
352
353 if err := tx.Commit(); err != nil {
354 return err
355 }
356
357 if n != nil {
358 n.NotifyAll()
359 }
360 return nil
361}