package 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 replaying func (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() }