Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280package db
import ( "context" "database/sql" "encoding/json" "time"
"tangled.org/core/notifier" "tangled.org/core/spindle/models")
// StatusRow is one append-only workflow status transition. Its id is the cursor// the mill executor's observe loop follows.type StatusRow struct { Id int64 Pipeline string // pipeline id Workflow string // workflow name Status string Error *string ExitCode *int64}
func insertStatusTx(tx DBTX, wid models.WorkflowId, kind models.StatusKind, workflowError *string, exitCode *int64) error { _, err := tx.Exec( `insert into workflow_statuses (pipeline_id, workflow, status, error, exit_code, created_at) values (?, ?, ?, ?, ?, ?)`, wid.PipelineId, wid.Name, string(kind), workflowError, exitCode, time.Now().Format(time.RFC3339Nano), ) return err}
func (d *DB) EventHighWater() (int64, error) { var id int64 err := d.QueryRow(`select coalesce(max(id), 0) from workflow_statuses`).Scan(&id) return id, err}
func (d *DB) GetEvents(cursor int64, limit int) ([]StatusRow, error) { rows, err := d.Query( `select id, pipeline_id, workflow, status, error, exit_code from workflow_statuses where id > ? order by id asc limit ?`, cursor, limit, ) if err != nil { return nil, err } defer rows.Close()
var out []StatusRow for rows.Next() { var r StatusRow if err := rows.Scan(&r.Id, &r.Pipeline, &r.Workflow, &r.Status, &r.Error, &r.ExitCode); err != nil { return nil, err } out = append(out, r) } return out, rows.Err()}
func (d *DB) createStatusEvent( workflowId models.WorkflowId, statusKind models.StatusKind, workflowError *string, exitCode *int64, n *notifier.Notifier,) error { if err := insertStatusTx(d, workflowId, statusKind, workflowError, exitCode); err != nil { return err } n.NotifyAll() return nil}
// deleting the lease in the same transaction prevents the terminal event// from replayingfunc (d *DB) CompleteMillLease( leaseID string, workflowId models.WorkflowId, status string, workflowError *string, exitCode *int64, n *notifier.Notifier,) error { return d.ApplyEventBatch(n, func(tx *EventBatchTx) error { if err := tx.InsertStatusEvent(workflowId, status, workflowError, exitCode); err != nil { return err } return tx.DeleteLease(leaseID) })}
func (d *DB) GetStatus(workflowId models.WorkflowId) (models.StatusKind, error) { pipelineId := workflowId.PipelineId
var status string err := d.QueryRow( ` select status from workflow_statuses where pipeline_id = ? and workflow = ? order by id desc limit 1 `, pipelineId, workflowId.Name, ).Scan(&status) if err != nil { if err == sql.ErrNoRows { return "", nil } return "", err }
return models.StatusKind(status), nil}
type statusQueryer interface { QueryRowContext(ctx context.Context, query string, args ...any) *sql.Row}
func workflowStartupDelay(ctx context.Context, q statusQueryer, wid models.WorkflowId) (time.Duration, bool, error) { var pending, running sql.NullString err := q.QueryRowContext(ctx, ` select min(case when status = 'pending' then created_at end), min(case when status = 'running' then created_at end) from workflow_statuses where pipeline_id = ? and workflow = ? `, string(wid.PipelineId), wid.Name).Scan(&pending, &running) if err != nil { return 0, false, err } if !pending.Valid || !running.Valid { return 0, false, nil } pendingAt, err := time.Parse(time.RFC3339Nano, pending.String) if err != nil { return 0, false, err } runningAt, err := time.Parse(time.RFC3339Nano, running.String) if err != nil { return 0, false, err } delay := runningAt.Sub(pendingAt) if delay < 0 { delay = 0 } return delay, true, nil}
func (d *DB) WorkflowStartupDelay(ctx context.Context, wid models.WorkflowId) (time.Duration, bool, error) { return workflowStartupDelay(ctx, d, wid)}
func (tx *EventBatchTx) WorkflowStartupDelay(ctx context.Context, wid models.WorkflowId) (time.Duration, bool, error) { return workflowStartupDelay(ctx, tx.tx, wid)}
func (tx *EventBatchTx) HasWorkflowStatus(ctx context.Context, wid models.WorkflowId, status string) (bool, error) { var present bool err := tx.tx.QueryRowContext(ctx, ` select exists( select 1 from workflow_statuses where pipeline_id = ? and workflow = ? and status = ? ) `, string(wid.PipelineId), wid.Name, status).Scan(&present) return present, err}
func (d *DB) StatusPending(workflowId models.WorkflowId, n *notifier.Notifier) error { return d.createStatusEvent(workflowId, models.StatusKindPending, nil, nil, n)}
func (d *DB) StatusRunning(workflowId models.WorkflowId, n *notifier.Notifier) error { return d.createStatusEvent(workflowId, models.StatusKindRunning, nil, nil, n)}
func (d *DB) StatusFailed(workflowId models.WorkflowId, workflowError string, exitCode int64, n *notifier.Notifier) error { return d.createStatusEvent(workflowId, models.StatusKindFailed, &workflowError, &exitCode, n)}
func (d *DB) StatusCancelled(workflowId models.WorkflowId, workflowError string, exitCode int64, n *notifier.Notifier) error { return d.createStatusEvent(workflowId, models.StatusKindCancelled, &workflowError, &exitCode, n)}
func (d *DB) StatusSuccess(workflowId models.WorkflowId, n *notifier.Notifier) error { return d.createStatusEvent(workflowId, models.StatusKindSuccess, nil, nil, n)}
func (d *DB) StatusTimeout(workflowId models.WorkflowId, n *notifier.Notifier) error { return d.createStatusEvent(workflowId, models.StatusKindTimeout, nil, nil, n)}
type PipelineWorkflow struct { PipelineID models.PipelineId Name string}
func (d *DB) ListPipelineWorkflows(repoDid string) ([]PipelineWorkflow, error) { rows, err := d.Query(`select pipeline_id, payload from pipelines where repo_did = ?`, repoDid) if err != nil { return nil, err } defer rows.Close()
var out []PipelineWorkflow for rows.Next() { var pipelineID, raw string if err := rows.Scan(&pipelineID, &raw); err != nil { return nil, err } var p models.PipelineRecord if err := json.Unmarshal([]byte(raw), &p); err != nil { continue } for _, wf := range p.Workflows { if wf != nil { out = append(out, PipelineWorkflow{PipelineID: models.PipelineId(pipelineID), Name: wf.Name}) } } } return out, rows.Err()}
func (d *DB) DeletePipelinesByRepo(repoDid string, extra []models.PipelineId) error { tx, err := d.Begin() if err != nil { return err } defer tx.Rollback()
rows, err := tx.Query(`select pipeline_id from pipelines where repo_did = ?`, repoDid) if err != nil { return err } seen := make(map[string]bool) var pipelineIDs []string add := func(pipelineID string) { if pipelineID != "" && !seen[pipelineID] { seen[pipelineID] = true pipelineIDs = append(pipelineIDs, pipelineID) } } for rows.Next() { var pipelineID string if err := rows.Scan(&pipelineID); err != nil { rows.Close() return err } add(pipelineID) } if err := rows.Err(); err != nil { rows.Close() return err } rows.Close() for _, id := range extra { add(string(id)) }
for _, pipelineID := range pipelineIDs { if _, err := tx.Exec(`delete from workflow_statuses where pipeline_id = ?`, pipelineID); err != nil { return err } } if _, err := tx.Exec(`delete from pipelines where repo_did = ?`, repoDid); err != nil { return err } return tx.Commit()}