Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554package db
import ( "context" "database/sql" "fmt" "time"
"tangled.org/core/notifier" "tangled.org/core/spindle/models")
// recovery uses persisted identity, fencing and quota statetype MillLease struct { LeaseID string NodeID string Epoch string Engine string PipelineID string Workflow string State string QuotaReservationID string OwnerDID string RepoDID string MillRecordsTerminalMetrics bool}
type ExecutorCursor struct { NodeID string Epoch string AckedSeqno uint64}
type OutboxRow struct { Epoch string Seqno uint64 Payload []byte ByteSize int64 Control bool}type OutboxDeletion struct { Rows int64 Bytes int64}
func (d *DB) SaveMillLease(l MillLease) error { _, err := d.Exec( `insert into mill_leases ( lease_id, node_id, epoch, engine, pipeline_id, workflow, state, quota_reservation_id, owner_did, repo_did, mill_records_terminal_metrics ) values (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) on conflict(lease_id) do update set state = excluded.state, quota_reservation_id = excluded.quota_reservation_id, owner_did = excluded.owner_did, repo_did = excluded.repo_did, mill_records_terminal_metrics = excluded.mill_records_terminal_metrics`, l.LeaseID, l.NodeID, l.Epoch, l.Engine, l.PipelineID, l.Workflow, l.State, l.QuotaReservationID, l.OwnerDID, l.RepoDID, l.MillRecordsTerminalMetrics, ) return err}
func (d *DB) DeleteMillLease(leaseID string) error { _, err := d.Exec(`delete from mill_leases where lease_id = ?`, leaseID) return err}
func (d *DB) ListMillLeases() ([]MillLease, error) { rows, err := d.Query(` select lease_id, node_id, epoch, engine, pipeline_id, workflow, state, coalesce(quota_reservation_id, ''), coalesce(owner_did, ''), coalesce(repo_did, ''), mill_records_terminal_metrics from mill_leases `) if err != nil { return nil, err } defer rows.Close()
var leases []MillLease for rows.Next() { var l MillLease if err := rows.Scan( &l.LeaseID, &l.NodeID, &l.Epoch, &l.Engine, &l.PipelineID, &l.Workflow, &l.State, &l.QuotaReservationID, &l.OwnerDID, &l.RepoDID, &l.MillRecordsTerminalMetrics, ); err != nil { return nil, err } leases = append(leases, l) } return leases, rows.Err()}
func (d *DB) ListExecutorCursors() ([]ExecutorCursor, error) { rows, err := d.Query(`select node_id, epoch, acked_seqno from mill_executor_cursors`) if err != nil { return nil, err } defer rows.Close()
var cursors []ExecutorCursor for rows.Next() { var c ExecutorCursor if err := rows.Scan(&c.NodeID, &c.Epoch, &c.AckedSeqno); err != nil { return nil, err } cursors = append(cursors, c) } return cursors, rows.Err()}
func (d *DB) SetOutboxEpoch(epoch string) error { return retrySQLite(func() error { tx, err := d.Begin() if err != nil { return err } defer tx.Rollback()
if _, err := tx.Exec(`delete from mill_outbox_rows`); err != nil { return err } if _, err := tx.Exec(`delete from mill_outbox_state`); err != nil { return err } if _, err := tx.Exec(`insert into mill_outbox_state (epoch, next_seqno) values (?, 1)`, epoch); err != nil { return err } return tx.Commit() })}
func (d *DB) AppendOutboxRow(payload []byte, control bool) (uint64, error) { var nextSeqno uint64 err := retrySQLite(func() error { tx, err := d.Begin() if err != nil { return err } defer tx.Rollback()
var epoch string err = tx.QueryRow(` update mill_outbox_state set next_seqno = next_seqno + 1 where epoch = (select epoch from mill_outbox_state limit 1) returning epoch, next_seqno - 1 `).Scan(&epoch, &nextSeqno) if err == sql.ErrNoRows { return fmt.Errorf("no outbox epoch set") } else if err != nil { return err }
byteSize := int64(len(payload)) controlVal := 0 if control { controlVal = 1 }
if _, err := tx.Exec( `insert into mill_outbox_rows (epoch, seqno, payload, byte_size, control) values (?, ?, ?, ?, ?)`, epoch, nextSeqno, payload, byteSize, controlVal, ); err != nil { return err }
return tx.Commit() }) if err != nil { return 0, err } return nextSeqno, nil}
func (d *DB) DeleteOutboxPrefix(ackedSeqno uint64) (OutboxDeletion, error) { var deleted OutboxDeletion err := retrySQLite(func() error { tx, err := d.Begin() if err != nil { return err } defer tx.Rollback()
var epoch string err = tx.QueryRow(`select epoch from mill_outbox_state limit 1`).Scan(&epoch) if err == sql.ErrNoRows { return nil } if err != nil { return err }
deleted = OutboxDeletion{} if err := tx.QueryRow(` select count(*), coalesce(sum(byte_size), 0) from mill_outbox_rows where epoch = ? and seqno <= ? `, epoch, ackedSeqno).Scan(&deleted.Rows, &deleted.Bytes); err != nil { return err } if _, err := tx.Exec( `delete from mill_outbox_rows where epoch = ? and seqno <= ?`, epoch, ackedSeqno, ); err != nil { return err } return tx.Commit() }) if err != nil { return OutboxDeletion{}, err } return deleted, nil}
func (d *DB) ListOutboxRows() ([]OutboxRow, error) { rows, err := d.Query(`select epoch, seqno, payload, byte_size, control from mill_outbox_rows order by seqno`) if err != nil { return nil, err } defer rows.Close()
var out []OutboxRow for rows.Next() { var r OutboxRow var controlVal int if err := rows.Scan(&r.Epoch, &r.Seqno, &r.Payload, &r.ByteSize, &controlVal); err != nil { return nil, err } r.Control = (controlVal != 0) out = append(out, r) } return out, rows.Err()}func (d *DB) ListOutboxRowsAfter(seqno uint64, limit int) ([]OutboxRow, error) { rows, err := d.Query(` select epoch, seqno, payload, byte_size, control from mill_outbox_rows where epoch = (select epoch from mill_outbox_state limit 1) and seqno > ? order by seqno limit ? `, seqno, limit) if err != nil { return nil, err } defer rows.Close()
var out []OutboxRow for rows.Next() { var row OutboxRow var control int if err := rows.Scan(&row.Epoch, &row.Seqno, &row.Payload, &row.ByteSize, &control); err != nil { return nil, err } row.Control = control != 0 out = append(out, row) } return out, rows.Err()}
func (d *DB) GetOutboxState() (string, uint64, error) { var epoch string var nextSeqno uint64 err := d.QueryRow(`select epoch, next_seqno from mill_outbox_state limit 1`).Scan(&epoch, &nextSeqno) if err == sql.ErrNoRows { return "", 0, nil } return epoch, nextSeqno, err}
type EventBatchTx struct { tx *sql.Tx db *DB}
func (tx *EventBatchTx) InsertStatusEvent(wid models.WorkflowId, status string, workflowError *string, exitCode *int64) error { return insertStatusTx(tx.tx, wid, models.StatusKind(status), workflowError, exitCode)}
func (tx *EventBatchTx) DeleteLease(leaseID string) error { _, err := tx.tx.Exec(`delete from mill_leases where lease_id = ?`, leaseID) return err}
func (tx *EventBatchTx) AdvanceCursor(nodeID, epoch string, seqno uint64) error { _, err := tx.tx.Exec( `insert into mill_executor_cursors (node_id, epoch, acked_seqno) values (?, ?, ?) on conflict(node_id, epoch) do update set acked_seqno = excluded.acked_seqno where excluded.acked_seqno > mill_executor_cursors.acked_seqno`, nodeID, epoch, seqno, ) return err}
// start from zero if a row is missingfunc (d *DB) GetExecutorCursor(nodeID, epoch string) (uint64, error) { var seqno uint64 err := d.QueryRow( `select acked_seqno from mill_executor_cursors where node_id = ? and epoch = ?`, nodeID, epoch, ).Scan(&seqno) if err == sql.ErrNoRows { return 0, nil } return seqno, err}
// forgets everything the mill applied from this node, picked up on the// executor's next reconnect since cursors are re-read at session startfunc (d *DB) DeleteExecutorCursors(nodeID string) (int64, error) { res, err := d.Exec(`delete from mill_executor_cursors where node_id = ?`, nodeID) if err != nil { return 0, err } return res.RowsAffected()}
// marks everything at or below seqno as applied, for skipping a backlog// the mill can never receivefunc (d *DB) SetExecutorCursors(nodeID string, seqno uint64) (int64, error) { res, err := d.Exec(`update mill_executor_cursors set acked_seqno = ? where node_id = ?`, seqno, nodeID) if err != nil { return 0, err } return res.RowsAffected()}
// clears artifacts waiting on leases from the old stream, replaying them// under a new epoch would trip the mill's lease-epoch checkfunc (d *DB) ClearPendingArtifacts() error { _, err := d.Exec(`delete from executor_pending_artifacts`) return err}
func (tx *EventBatchTx) InsertArtifactRef(leaseID, repoDid string, wid models.WorkflowId, ref, hash string) error { if err := validateWorkflowIdentity(wid); err != nil { return err } _, err := tx.tx.Exec( `insert into mill_artifacts (lease_id, repo_did, pipeline_id, workflow, ref, hash) values (?, ?, ?, ?, ?, ?)`, leaseID, repoDid, string(wid.PipelineId), wid.Name, ref, hash, ) return err}
type PendingArtifact struct { LeaseID string PipelineID string Workflow string Status string Error string ExitCode int64 Ref string Hash string FailureClass string FailureReason string MillRecordsTerminalMetrics bool}
func (d *DB) SavePendingArtifact( leaseID string, wid models.WorkflowId, status, errStr string, exitCode int64, ref, hash, failureClass, failureReason string, millRecordsTerminalMetrics bool,) error { _, err := d.Exec( `insert into executor_pending_artifacts ( lease_id, pipeline_id, workflow, status, error, exit_code, ref, hash, failure_class, failure_reason, mill_records_terminal_metrics ) values (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) on conflict(lease_id) do update set pipeline_id = excluded.pipeline_id, workflow = excluded.workflow, status = excluded.status, error = excluded.error, exit_code = excluded.exit_code, ref = excluded.ref, hash = excluded.hash, failure_class = excluded.failure_class, failure_reason = excluded.failure_reason, mill_records_terminal_metrics = excluded.mill_records_terminal_metrics`, leaseID, string(wid.PipelineId), wid.Name, status, errStr, exitCode, ref, hash, failureClass, failureReason, millRecordsTerminalMetrics, ) return err}
func (d *DB) RemovePendingArtifact(leaseID string) error { _, err := d.Exec(`delete from executor_pending_artifacts where lease_id = ?`, leaseID) return err}
func (d *DB) ListPendingArtifacts() ([]PendingArtifact, error) { rows, err := d.Query(` select lease_id, pipeline_id, workflow, status, error, exit_code, ref, hash, failure_class, failure_reason, mill_records_terminal_metrics from executor_pending_artifacts `) if err != nil { return nil, err } defer rows.Close()
var res []PendingArtifact for rows.Next() { var p PendingArtifact if err := rows.Scan( &p.LeaseID, &p.PipelineID, &p.Workflow, &p.Status, &p.Error, &p.ExitCode, &p.Ref, &p.Hash, &p.FailureClass, &p.FailureReason, &p.MillRecordsTerminalMetrics, ); err != nil { return nil, err } res = append(res, p) } return res, rows.Err()}
// ApplyEventBatch runs fn in a transaction and retries it while sqlite reports// transient lock errors. fn may run more than once, so non-transactional side// effects must be returned through its result: the result is published only// from the attempt that commits, and a failed attempt's result is discarded.func ApplyEventBatch[T any](d *DB, n *notifier.Notifier, fn func(tx *EventBatchTx) (T, error)) (T, error) { return ApplyEventBatchCtx(context.Background(), d, n, fn)}
func ApplyEventBatchCtx[T any](ctx context.Context, d *DB, n *notifier.Notifier, fn func(tx *EventBatchTx) (T, error)) (T, error) { var result T if err := ctx.Err(); err != nil { var zero T return zero, err } err := retrySQLiteCtx(ctx, 20*time.Second, func() error { if err := ctx.Err(); err != nil { return err } tx, err := d.Begin() if err != nil { return err } defer tx.Rollback()
batchTx := &EventBatchTx{ tx: tx, db: d, } v, err := fn(batchTx) if err != nil { return err } if err := tx.Commit(); err != nil { return err } result = v return nil }) if err != nil { var zero T return zero, err } if n != nil { n.NotifyAll() } return result, nil}
// ApplyEventBatch is the package-level ApplyEventBatch for callers without a// result.func (d *DB) ApplyEventBatch(n *notifier.Notifier, fn func(tx *EventBatchTx) error) error { return d.ApplyEventBatchCtx(context.Background(), n, fn)}
func (d *DB) ApplyEventBatchCtx(ctx context.Context, n *notifier.Notifier, fn func(tx *EventBatchTx) error) error { _, err := ApplyEventBatchCtx(ctx, d, n, func(tx *EventBatchTx) (struct{}, error) { return struct{}{}, fn(tx) }) return err}
func (d *DB) DeleteMillLeasesByRepo(repoDid string) ([]string, error) { tx, err := d.Begin() if err != nil { return nil, err } defer tx.Rollback()
rows, err := tx.Query(`select lease_id from mill_leases where repo_did = ?`, repoDid) if err != nil { return nil, err } var leaseIDs []string for rows.Next() { var id string if err := rows.Scan(&id); err != nil { rows.Close() return nil, err } leaseIDs = append(leaseIDs, id) } if err := rows.Err(); err != nil { rows.Close() return nil, err } rows.Close()
for _, id := range leaseIDs { if _, err := tx.Exec(`delete from executor_pending_artifacts where lease_id = ?`, id); err != nil { return nil, err } } if _, err := tx.Exec(`delete from mill_leases where repo_did = ?`, repoDid); err != nil { return nil, err } // return ids only after the deletion commits if err := tx.Commit(); err != nil { return nil, err } return leaseIDs, nil}
func (d *DB) ListArtifactRefsByRepo(repoDid string) ([]string, error) { rows, err := d.Query(` select ref from mill_artifacts where repo_did = ? union select pa.ref from executor_pending_artifacts pa join mill_leases ml on ml.lease_id = pa.lease_id where ml.repo_did = ?`, repoDid, repoDid, ) if err != nil { return nil, err } defer rows.Close()
var refs []string for rows.Next() { var ref string if err := rows.Scan(&ref); err != nil { return nil, err } refs = append(refs, ref) } return refs, rows.Err()}
func (d *DB) DeleteArtifactRefsByRepo(repoDid string) error { _, err := d.Exec(`delete from mill_artifacts where repo_did = ?`, repoDid) return err}