Something went wrong. Try again.
Monorepo for Tangled
Something went wrong. Try again.
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448package db
import ( "context" "database/sql" "fmt" "slices" "strings" "time"
"github.com/bluesky-social/indigo/atproto/syntax" "tangled.org/core/appview/models" "tangled.org/core/orm")
func GetPipelines(e Execer, filters ...orm.Filter) ([]models.Pipeline, error) { var pipelines []models.Pipeline
var conditions []string var args []any for _, filter := range filters { conditions = append(conditions, filter.Condition()) args = append(args, filter.Arg()...) }
whereClause := "" if conditions != nil { whereClause = " where " + strings.Join(conditions, " and ") }
query := fmt.Sprintf(`select id, rkey, knot, repo_owner, repo_name, sha, created, repo_did from pipelines %s`, whereClause)
rows, err := e.Query(query, args...)
if err != nil { return nil, err } defer rows.Close()
for rows.Next() { var pipeline models.Pipeline var createdAt string var repoDid sql.NullString err = rows.Scan( &pipeline.Id, &pipeline.Rkey, &pipeline.Knot, &pipeline.RepoOwner, &pipeline.RepoName, &pipeline.Sha, &createdAt, &repoDid, ) if err != nil { return nil, err }
if t, err := time.Parse(time.RFC3339, createdAt); err == nil { pipeline.Created = t } if repoDid.Valid { pipeline.RepoDid = repoDid.String }
pipelines = append(pipelines, pipeline) }
if err = rows.Err(); err != nil { return nil, err }
return pipelines, nil}
func AddPipeline(e Execer, pipeline models.Pipeline) error { var repoDid *string if pipeline.RepoDid != "" { repoDid = &pipeline.RepoDid }
args := []any{ pipeline.Rkey, pipeline.Knot, pipeline.RepoOwner, pipeline.RepoName, pipeline.TriggerId, pipeline.Sha, repoDid, }
placeholders := make([]string, len(args)) for i := range placeholders { placeholders[i] = "?" }
query := fmt.Sprintf(` insert or ignore into pipelines ( rkey, knot, repo_owner, repo_name, trigger_id, sha, repo_did ) values (%s) `, strings.Join(placeholders, ","))
_, err := e.Exec(query, args...)
return err}
func AddTrigger(e Execer, trigger models.Trigger) (int64, error) { args := []any{ trigger.Kind, trigger.PushRef, trigger.PushNewSha, trigger.PushOldSha, trigger.PRSourceBranch, trigger.PRTargetBranch, trigger.PRSourceSha, trigger.PRAction, }
placeholders := make([]string, len(args)) for i := range placeholders { placeholders[i] = "?" }
query := fmt.Sprintf(`insert or ignore into triggers ( kind, push_ref, push_new_sha, push_old_sha, pr_source_branch, pr_target_branch, pr_source_sha, pr_action ) values (%s)`, strings.Join(placeholders, ","))
res, err := e.Exec(query, args...) if err != nil { return 0, err }
return res.LastInsertId()}
func AddPipelineStatus(ctx context.Context, e Execer, status models.PipelineStatus) error { args := []any{ status.Spindle, status.Rkey, status.PipelineKnot, status.PipelineRkey, status.Workflow, status.Status, status.Error, status.ExitCode, status.Created.Format(time.RFC3339), }
placeholders := make([]string, len(args)) for i := range placeholders { placeholders[i] = "?" }
query := fmt.Sprintf(` insert or ignore into pipeline_statuses ( spindle, rkey, pipeline_knot, pipeline_rkey, workflow, status, error, exit_code, created ) values (%s) `, strings.Join(placeholders, ","))
_, err := e.ExecContext(ctx, query, args...) return err}
// this is a mega query, but the most useful one:// get N pipelines, for each one get the latest status of its N workflows//// the pipelines table is aliased to `p`// the triggers table is aliased to `t`func GetPipelineStatuses(e Execer, limit int, filters ...orm.Filter) ([]models.Pipeline, error) { var conditions []string var args []any for _, filter := range filters { conditions = append(conditions, filter.Condition()) args = append(args, filter.Arg()...) }
whereClause := "" if conditions != nil { whereClause = " where " + strings.Join(conditions, " and ") }
query := fmt.Sprintf(` select p.id, p.knot, p.rkey, p.repo_owner, p.repo_name, p.sha, p.created, p.repo_did, t.id, t.kind, t.push_ref, t.push_new_sha, t.push_old_sha, t.pr_source_branch, t.pr_target_branch, t.pr_source_sha, t.pr_action from pipelines p join triggers t ON p.trigger_id = t.id %s order by p.created desc limit %d `, whereClause, limit)
rows, err := e.Query(query, args...) if err != nil { return nil, err } defer rows.Close()
pipelines := make(map[syntax.ATURI]models.Pipeline) for rows.Next() { var p models.Pipeline var t models.Trigger var created string var repoDid sql.NullString
err := rows.Scan( &p.Id, &p.Knot, &p.Rkey, &p.RepoOwner, &p.RepoName, &p.Sha, &created, &repoDid, &p.TriggerId, &t.Kind, &t.PushRef, &t.PushNewSha, &t.PushOldSha, &t.PRSourceBranch, &t.PRTargetBranch, &t.PRSourceSha, &t.PRAction, ) if err != nil { return nil, err }
p.Created, err = time.Parse(time.RFC3339, created) if err != nil { return nil, fmt.Errorf("invalid pipeline created timestamp %q: %w", created, err) } if repoDid.Valid { p.RepoDid = repoDid.String }
t.Id = p.TriggerId p.Trigger = &t p.Statuses = make(map[string]models.WorkflowStatus)
pipelines[p.AtUri()] = p }
// get all statuses // the where clause here is of the form: // // and ( // (ps.pipeline_knot = k1 and ps.pipeline_rkey = r1) // or (ps.pipeline_knot = k2 and ps.pipeline_rkey = r2) // ) // // the join on pipelines and repos enforces that the status was emitted // by the spindle that is actually registered for the pipeline's repo. conditions = nil args = nil for _, p := range pipelines { knotFilter := orm.FilterEq("ps.pipeline_knot", p.Knot) rkeyFilter := orm.FilterEq("ps.pipeline_rkey", p.Rkey) conditions = append(conditions, fmt.Sprintf("(%s and %s)", knotFilter.Condition(), rkeyFilter.Condition())) args = append(args, p.Knot) args = append(args, p.Rkey) } whereClause = "" if conditions != nil { whereClause = "and (" + strings.Join(conditions, " or ") + ")" } query = fmt.Sprintf(` select ps.id, ps.spindle, ps.rkey, ps.pipeline_knot, ps.pipeline_rkey, ps.created, ps.workflow, ps.status, ps.error, ps.exit_code from pipeline_statuses ps join pipelines p on p.knot = ps.pipeline_knot and p.rkey = ps.pipeline_rkey join repos r on r.repo_did = p.repo_did where ps.spindle = r.spindle %s `, whereClause)
rows, err = e.Query(query, args...) if err != nil { return nil, err } defer rows.Close()
for rows.Next() { var ps models.PipelineStatus var created string
err := rows.Scan( &ps.ID, &ps.Spindle, &ps.Rkey, &ps.PipelineKnot, &ps.PipelineRkey, &created, &ps.Workflow, &ps.Status, &ps.Error, &ps.ExitCode, ) if err != nil { return nil, err }
ps.Created, err = time.Parse(time.RFC3339, created) if err != nil { return nil, fmt.Errorf("invalid status created timestamp %q: %w", created, err) }
pipelineAt := ps.PipelineAt()
// extract pipeline, ok := pipelines[pipelineAt] if !ok { continue } statuses, _ := pipeline.Statuses[ps.Workflow] if !ok { pipeline.Statuses[ps.Workflow] = models.WorkflowStatus{} }
// append statuses.Data = append(statuses.Data, ps)
// reassign pipeline.Statuses[ps.Workflow] = statuses pipelines[pipelineAt] = pipeline }
var all []models.Pipeline for _, p := range pipelines { for _, s := range p.Statuses { slices.SortFunc(s.Data, func(a, b models.PipelineStatus) int { if a.Created.After(b.Created) { return 1 } if a.Created.Before(b.Created) { return -1 } if a.ID > b.ID { return 1 } if a.ID < b.ID { return -1 } return 0 }) } all = append(all, p) }
// sort pipelines by date slices.SortFunc(all, func(a, b models.Pipeline) int { if a.Created.After(b.Created) { return -1 } return 1 })
return all, nil}
// the pipelines table is aliased to `p`// the triggers table is aliased to `t`func GetPipelineCount(e Execer, filters ...orm.Filter) (int64, error) { var conditions []string var args []any for _, filter := range filters { conditions = append(conditions, filter.Condition()) args = append(args, filter.Arg()...) }
whereClause := "" if conditions != nil { whereClause = " where " + strings.Join(conditions, " and ") }
query := fmt.Sprintf(` select count(1) from pipelines p join triggers t ON p.trigger_id = t.id %s `, whereClause)
rows, err := e.Query(query, args...) if err != nil { return 0, err } defer rows.Close()
for rows.Next() { var count int64 err := rows.Scan(&count) if err != nil { return 0, err }
return count, nil }
// unreachable return 0, nil}