This repository has no description
0

Configure Feed

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

core / spindle / db / mill_state.go
8.6 kB 361 lines
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}