diff --git a/appview/db/pipeline.go b/appview/db/pipeline.go deleted file mode 100644 --- a/appview/db/pipeline.go +++ /dev/null @@ -1,447 +0,0 @@ -package 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 -} diff --git a/appview/db/pipeline_test.go b/appview/db/pipeline_test.go deleted file mode 100644 --- a/appview/db/pipeline_test.go +++ /dev/null @@ -1,116 +0,0 @@ -package db - -import ( - "context" - "strings" - "testing" - "time" - - "github.com/bluesky-social/indigo/atproto/syntax" - "tangled.org/core/appview/models" - "tangled.org/core/orm" - spindle "tangled.org/core/spindle/models" - "tangled.org/core/workflow" -) - -// seedPipeline inserts a trigger + pipeline row and returns the pipeline. -func seedPipeline(t *testing.T, d *DB, knot, rkey, repoDid string) models.Pipeline { - t.Helper() - sha := strings.Repeat("a", 40) - ref := "refs/heads/main" - newSha := sha - oldSha := strings.Repeat("0", 40) - trigger := models.Trigger{ - Kind: workflow.TriggerKindPush, - PushRef: &ref, - PushNewSha: &newSha, - PushOldSha: &oldSha, - } - tx, err := d.Begin() - if err != nil { - t.Fatalf("Begin: %v", err) - } - triggerID, err := AddTrigger(tx, trigger) - if err != nil { - tx.Rollback() - t.Fatalf("AddTrigger: %v", err) - } - pipeline := models.Pipeline{ - Knot: knot, - Rkey: rkey, - RepoOwner: syntax.DID("did:plc:owner"), - RepoName: "repo", - RepoDid: repoDid, - TriggerId: int(triggerID), - Sha: sha, - } - if err := AddPipeline(tx, pipeline); err != nil { - tx.Rollback() - t.Fatalf("AddPipeline: %v", err) - } - if err := tx.Commit(); err != nil { - t.Fatalf("Commit: %v", err) - } - return pipeline -} - -// seedStatus inserts a pipeline_status row directly. -func seedStatus(t *testing.T, d *DB, spindleInstance, rkey, pipelineKnot, pipelineRkey, workflow string) { - t.Helper() - status := models.PipelineStatus{ - Spindle: spindleInstance, - Rkey: rkey, - PipelineKnot: pipelineKnot, - PipelineRkey: pipelineRkey, - Workflow: workflow, - Status: spindle.StatusKindSuccess, - Created: time.Now(), - } - if err := AddPipelineStatus(context.Background(), d, status); err != nil { - t.Fatalf("AddPipelineStatus: %v", err) - } -} - -// TestGetPipelineStatuses_SpindleValidation verifies that GetPipelineStatuses -// only returns statuses emitted by the spindle registered for the pipeline's -// repo, and silently drops statuses from a rogue spindle. -func TestGetPipelineStatuses_SpindleValidation(t *testing.T) { - d := newTestDB(t) - - const ( - knot = "knot.example.com" - correctSpindle = "spindle.example.com" - rogueSpindle = "evil.example.com" - repoDid = "did:plc:testrepo" - pipelineRkey = "pipeline1" - ) - - // seed repo with the correct spindle - repo := seedRepo(t, d, "did:plc:owner", knot, "repo", "repo", repoDid) - if err := UpdateSpindle(d, repo.RepoDid, &[]string{correctSpindle}[0]); err != nil { - t.Fatalf("UpdateSpindle: %v", err) - } - - // seed the pipeline for this repo - seedPipeline(t, d, knot, pipelineRkey, repoDid) - - // insert one status from the correct spindle, one from a rogue spindle - seedStatus(t, d, correctSpindle, "status-valid", knot, pipelineRkey, "build") - seedStatus(t, d, rogueSpindle, "status-rogue", knot, pipelineRkey, "build") - - pipelines, err := GetPipelineStatuses(d, 10, orm.FilterEq("p.repo_did", repoDid)) - if err != nil { - t.Fatalf("GetPipelineStatuses: %v", err) - } - if len(pipelines) != 1 { - t.Fatalf("expected 1 pipeline, got %d", len(pipelines)) - } - - statuses := pipelines[0].Statuses["build"].Data - if len(statuses) != 1 { - t.Fatalf("expected 1 status (from correct spindle), got %d", len(statuses)) - } - if statuses[0].Spindle != correctSpindle { - t.Errorf("expected spindle %q, got %q", correctSpindle, statuses[0].Spindle) - } -} diff --git a/appview/pages/pages.go b/appview/pages/pages.go --- a/appview/pages/pages.go +++ b/appview/pages/pages.go @@ -890,7 +890,7 @@ Raw bool EmailToDid map[string]string VerifiedCommits commitverify.VerifiedCommits Languages []types.RepoLanguageDetails - Pipelines map[string]models.Pipeline + Pipelines map[string]*tangled.CiDefs_Pipeline NeedsKnotUpgrade bool KnotUnreachable bool types.RepoIndexResponse @@ -959,7 +959,7 @@ TagMap map[string][]string Active string EmailToDid map[string]string VerifiedCommits commitverify.VerifiedCommits - Pipelines map[string]models.Pipeline + Pipelines map[string]*tangled.CiDefs_Pipeline types.RepoLogResponse } @@ -974,7 +974,7 @@ BaseParams RepoInfo repoinfo.RepoInfo Active string EmailToDid map[string]string - Pipeline *models.Pipeline + Pipeline *tangled.CiDefs_Pipeline DiffOpts types.DiffOpts // singular because it's always going to be just one @@ -1374,7 +1374,7 @@ FilterState string FilterQuery string BaseFilterQuery string Stacks []models.Stack - Pipelines map[string]models.Pipeline + Pipelines map[string]tangled.CiDefs_Pipeline LabelDefs map[string]*models.LabelDefinition Page pagination.Page PullCount int @@ -1414,7 +1414,7 @@ Backlinks []models.RichReferenceLink BranchDeleteStatus *models.BranchDeleteStatus MergeCheck types.MergeCheckResponse ResubmitCheck ResubmitResult - Pipelines map[string]models.Pipeline + Pipelines map[string]tangled.CiDefs_Pipeline Diff types.DiffRenderer DiffOpts types.DiffOpts ActiveRound int @@ -1586,7 +1586,7 @@ type PipelinesParams struct { BaseParams RepoInfo repoinfo.RepoInfo - Pipelines []models.Pipeline + Pipelines []*tangled.CiDefs_Pipeline Active string FilterKind string Total int64 @@ -1640,7 +1640,7 @@ type WorkflowParams struct { BaseParams RepoInfo repoinfo.RepoInfo - Pipeline models.Pipeline + Pipeline *tangled.CiDefs_Pipeline Workflow string LogUrl string Active string diff --git a/appview/pipelines/logs.go b/appview/pipelines/logs.go --- a/appview/pipelines/logs.go +++ b/appview/pipelines/logs.go @@ -1,15 +1,15 @@ package pipelines import ( + "errors" "html/template" - "path" "regexp" "strings" + "time" terminal "github.com/buildkite/terminal-to-html/v3" "github.com/gorilla/websocket" "tangled.org/core/appview/pages/markup/sanitizer" - "tangled.org/core/hostutil" ) // matches any ANSI escape sequence: ESC [ m @@ -50,40 +50,32 @@ return template.HTML(sanitized) } -type LogEvent struct { - Msg []byte - Err error -} - -func (ev *LogEvent) IsCloseError() bool { - return websocket.IsCloseError( - ev.Err, - websocket.CloseNormalClosure, - websocket.CloseGoingAway, - websocket.CloseAbnormalClosure, - ) -} - -func ReadLogs(conn *websocket.Conn, ch chan LogEvent) { - defer close(ch) - for { - if conn == nil { - return - } - _, msg, err := conn.ReadMessage() - if err != nil { - ch <- LogEvent{Err: err} - return +// isExpectedClose reports whether err is a clean websocket close (or nil). +func isExpectedClose(err error) bool { + if err == nil { + return true + } + var ce *websocket.CloseError + if errors.As(err, &ce) { + switch ce.Code { + case websocket.CloseNormalClosure, websocket.CloseGoingAway, websocket.CloseAbnormalClosure: + return true } - ch <- LogEvent{Msg: msg} } + return false } -func SpindleURL(spindle, knot, rkey, workflow string) string { - url, err := hostutil.EnsureWsScheme(spindle) - if err != nil { +func derefStr(s *string) string { + if s == nil { return "" } + return *s +} - return url + path.Join("/logs", knot, rkey, workflow) +func parseRFC3339(s string) time.Time { + t, err := time.Parse(time.RFC3339, s) + if err != nil { + return time.Time{} + } + return t } diff --git a/appview/pipelines/notifier.go b/appview/pipelines/notifier.go deleted file mode 100644 --- a/appview/pipelines/notifier.go +++ /dev/null @@ -1,52 +0,0 @@ -package pipelines - -import ( - "sync" - - "github.com/bluesky-social/indigo/atproto/syntax" - "tangled.org/core/notifier" -) - -// StatusNotifier is a keyed broadcast notifier for pipeline status changes, keyed by the pipeline's AT URI. -// -// subscribers are notified whenever a status update arrives for that pipeline -type StatusNotifier struct { - mu sync.Mutex - keys map[syntax.ATURI]*notifier.Notifier -} - -func NewStatusNotifier() *StatusNotifier { - return &StatusNotifier{ - keys: make(map[syntax.ATURI]*notifier.Notifier), - } -} - -func (n *StatusNotifier) Publish(uri syntax.ATURI) { - n.mu.Lock() - p, ok := n.keys[uri] - n.mu.Unlock() - if ok { - p.NotifyAll() - } -} - -func (n *StatusNotifier) Subscribe(uri syntax.ATURI) chan struct{} { - n.mu.Lock() - p, ok := n.keys[uri] - if !ok { - nb := notifier.New() - p = &nb - n.keys[uri] = p - } - n.mu.Unlock() - return p.Subscribe() -} - -func (n *StatusNotifier) Unsubscribe(uri syntax.ATURI, ch chan struct{}) { - n.mu.Lock() - p, ok := n.keys[uri] - n.mu.Unlock() - if ok { - p.Unsubscribe(ch) - } -} diff --git a/appview/pipelines/pipelines.go b/appview/pipelines/pipelines.go --- a/appview/pipelines/pipelines.go +++ b/appview/pipelines/pipelines.go @@ -3,42 +3,39 @@ import ( "bytes" "context" - "encoding/json" - "fmt" "log/slog" "net/http" + "sync" "time" "tangled.org/core/api/tangled" "tangled.org/core/appview/config" "tangled.org/core/appview/db" "tangled.org/core/appview/middleware" - "tangled.org/core/appview/models" "tangled.org/core/appview/oauth" "tangled.org/core/appview/pages" "tangled.org/core/appview/reporesolver" - "tangled.org/core/eventconsumer" "tangled.org/core/hostutil" "tangled.org/core/idresolver" + "tangled.org/core/lexutil" "tangled.org/core/orm" "tangled.org/core/rbac" - spindlemodel "tangled.org/core/spindle/models" + "github.com/bluesky-social/indigo/atproto/syntax" + indigoxrpc "github.com/bluesky-social/indigo/xrpc" "github.com/go-chi/chi/v5" "github.com/gorilla/websocket" ) type Pipelines struct { - repoResolver *reporesolver.RepoResolver - idResolver *idresolver.Resolver - config *config.Config - oauth *oauth.OAuth - pages *pages.Pages - spindlestream *eventconsumer.Consumer - pipelineNotifier *StatusNotifier - db *db.DB - enforcer *rbac.Enforcer - logger *slog.Logger + repoResolver *reporesolver.RepoResolver + idResolver *idresolver.Resolver + config *config.Config + oauth *oauth.OAuth + pages *pages.Pages + db *db.DB + enforcer *rbac.Enforcer + logger *slog.Logger } func (p *Pipelines) Router(mw *middleware.Middleware) http.Handler { @@ -48,7 +45,7 @@ r.Get("/{pipeline}/workflow/{workflow}", p.Workflow) r.Get("/{pipeline}/workflow/{workflow}/logs", p.Logs) r. With(mw.RepoPermissionMiddleware("repo:owner")). - Post("/{pipeline}/workflow/{workflow}/cancel", p.Cancel) + Post("/{pipeline}/workflow/{workflow}/cancel", p.CancelWorkflow) return r } @@ -57,8 +54,6 @@ func New( oauth *oauth.OAuth, repoResolver *reporesolver.RepoResolver, pages *pages.Pages, - spindlestream *eventconsumer.Consumer, - pipelineNotifier *StatusNotifier, idResolver *idresolver.Resolver, db *db.DB, config *config.Config, @@ -66,16 +61,14 @@ enforcer *rbac.Enforcer, logger *slog.Logger, ) *Pipelines { return &Pipelines{ - oauth: oauth, - repoResolver: repoResolver, - pages: pages, - idResolver: idResolver, - config: config, - spindlestream: spindlestream, - pipelineNotifier: pipelineNotifier, - db: db, - enforcer: enforcer, - logger: logger, + oauth: oauth, + repoResolver: repoResolver, + pages: pages, + idResolver: idResolver, + config: config, + db: db, + enforcer: enforcer, + logger: logger, } } @@ -103,28 +96,20 @@ // no filters otherwise, default to "all" filterKind = "all" } - ps, err := db.GetPipelineStatuses( - p.db, - 30, - filters..., - ) + // sh.tangled.ci.queryPipelines(repo, kind, limit=30) + xrpcc := indigoxrpc.Client{Host: f.Spindle} + out, err := tangled.CiQueryPipelines(r.Context(), &xrpcc, nil, "", 1, f.RepoDid) if err != nil { - l.Error("failed to query db", "err", err) - return - } - - total, err := db.GetPipelineCount(p.db, filters...) - if err != nil { - l.Error("failed to query db", "err", err) - return + l.Error("failed to fetch pipelines", "err", err) + panic("unimplemented") // spindle failure, appview should not fail. } p.pages.Pipelines(w, pages.PipelinesParams{ BaseParams: pages.BaseParamsFromContext(r.Context()), RepoInfo: p.repoResolver.GetRepoInfo(r, user), - Pipelines: ps, + Pipelines: out.Pipelines, FilterKind: filterKind, - Total: total, + Total: out.Total, }) } @@ -135,44 +120,57 @@ f, err := p.repoResolver.Resolve(r) if err != nil { l.Error("failed to get repo and knot", "err", err) + p.pages.Error404(w) return } - pipelineId := chi.URLParam(r, "pipeline") - if pipelineId == "" { - l.Error("empty pipeline ID") + pipelineId, err := syntax.ParseTID(chi.URLParam(r, "pipeline")) + if err != nil { + l.Debug("invalid pipeline id", "id", pipelineId) + p.pages.Error404(w) return } - workflow := chi.URLParam(r, "workflow") - if workflow == "" { - l.Error("empty workflow name") + workflowName := chi.URLParam(r, "workflow") + if workflowName == "" { + l.Debug("empty workflow name") + p.pages.Error404(w) return } - ps, err := db.GetPipelineStatuses( - p.db, - 1, - orm.FilterEq("p.repo_did", f.RepoDid), - orm.FilterEq("p.id", pipelineId), - ) + l = l.With("pipeline", pipelineId, "workflow", workflowName) + + // TODO: change url path to: + // /{owner}/{slug}/pipelines/{spindle-did}/{pipeline-id}/workflow/{workflow-id} + + xrpcc := &indigoxrpc.Client{Host: f.Spindle} + out, err := tangled.CiGetPipeline(r.Context(), xrpcc, pipelineId.String()) if err != nil { - l.Error("failed to query db", "err", err) + // TODO(boltless): change behavior based on error + l.Debug("failed to get pipeline", "err", err) + p.pages.Error404(w) return } - if len(ps) != 1 { - l.Error("invalid number of pipelines", "len", len(ps)) + // ensure workflow exists + exist := false + for _, workflow := range out.Workflows { + if workflow.Name == workflowName { + exist = true + break + } + } + if !exist { + l.Debug("workflow doesn't exist in pipeline") + p.pages.Error404(w) return } - singlePipeline := ps[0] - p.pages.Workflow(w, pages.WorkflowParams{ BaseParams: pages.BaseParamsFromContext(r.Context()), RepoInfo: p.repoResolver.GetRepoInfo(r, user), - Pipeline: singlePipeline, - Workflow: workflow, + Pipeline: out, + Workflow: workflowName, }) } @@ -181,6 +179,25 @@ ReadBufferSize: 1024, WriteBufferSize: 1024, } +type webLogScheduler struct { + ch chan *tangled.CiPipelineSubscribeLogs_Event +} + +var _ lexutil.Scheduler[tangled.CiPipelineSubscribeLogs_Event] = (*webLogScheduler)(nil) + +// AddWork implements [lexutil.Scheduler]. +func (w *webLogScheduler) AddWork(ctx context.Context, _ string, val *tangled.CiPipelineSubscribeLogs_Event) error { + select { + case w.ch <- val: + return nil + case <-ctx.Done(): + return ctx.Err() + } +} + +// Shutdown implements [lexutil.Scheduler]. +func (w *webLogScheduler) Shutdown() { close(w.ch) } + func (p *Pipelines) Logs(w http.ResponseWriter, r *http.Request) { l := p.logger.With("handler", "logs") @@ -191,76 +208,91 @@ http.Error(w, "bad repo/knot", http.StatusBadRequest) return } - pipelineId := chi.URLParam(r, "pipeline") - workflow := chi.URLParam(r, "workflow") - if pipelineId == "" || workflow == "" { - http.Error(w, "missing pipeline ID or workflow", http.StatusBadRequest) + if f.Spindle == "" { + http.Error(w, "invalid repo info", http.StatusBadRequest) return } - ps, err := db.GetPipelineStatuses( - p.db, - 1, - orm.FilterEq("p.repo_did", f.RepoDid), - orm.FilterEq("p.id", pipelineId), - ) - if err != nil || len(ps) != 1 { - l.Error("pipeline query failed", "err", err, "count", len(ps)) - http.Error(w, "pipeline not found", http.StatusNotFound) + pipelineId, err := syntax.ParseTID(chi.URLParam(r, "pipeline")) + if err != nil { + l.Debug("invalid pipeline id", "id", pipelineId) + http.Error(w, "invalid pipeline id", http.StatusBadRequest) return } - singlePipeline := ps[0] - spindle := f.Spindle - knot := f.Knot - rkey := singlePipeline.Rkey - - statusCh := p.pipelineNotifier.Subscribe(singlePipeline.AtUri()) - defer p.pipelineNotifier.Unsubscribe(singlePipeline.AtUri(), statusCh) - - if spindle == "" || knot == "" || rkey == "" { - http.Error(w, "invalid repo info", http.StatusBadRequest) + workflowName := chi.URLParam(r, "workflow") + if workflowName == "" { + l.Debug("empty workflow name") + http.Error(w, "invalid workflow name", http.StatusBadRequest) return } - url := SpindleURL(spindle, knot, rkey, workflow) - if url == "" { - http.Error(w, "invalid spindle hostname", http.StatusBadRequest) - return - } - l = l.With("url", url) - clientConn, err := upgrader.Upgrade(w, r, nil) if err != nil { l.Error("websocket upgrade failed", "err", err) return } - defer func() { - _ = clientConn.WriteControl( - websocket.CloseMessage, - websocket.FormatCloseMessage(websocket.CloseNormalClosure, "log stream complete"), - time.Now().Add(time.Second), - ) - clientConn.Close() - }() + defer clientConn.Close() ctx, cancel := context.WithCancel(r.Context()) defer cancel() - l.Info("logs endpoint hit") + evChan := make(chan *tangled.CiPipelineSubscribeLogs_Event, 100) + done := make(chan error, 1) + sched := &webLogScheduler{ch: evChan} + xrpcc := &lexutil.Client{Client: indigoxrpc.Client{Host: f.Spindle}} + go func() { + done <- tangled.CiPipelineSubscribeLogs(ctx, xrpcc, pipelineId.String(), []string{workflowName}, sched) + }() + + var lastWriteLk sync.Mutex + lastWrite := time.Now() + + // Start a goroutine to ping the client periodically to check if it's still + // alive. If the client doesn't respond to a ping within 5 seconds, we'll + // close the connection and teardown the consumer. + go func() { + ticker := time.NewTicker(30 * time.Second) + defer ticker.Stop() + for { + select { + case <-ticker.C: + lastWriteLk.Lock() + lw := lastWrite + lastWriteLk.Unlock() + if time.Since(lw) < 30*time.Second { + continue + } + if err := clientConn.WriteControl(websocket.PingMessage, nil, time.Now().Add(5*time.Second)); err != nil { + l.Warn("failed to ping client", "err", err) + cancel() + return + } + case <-ctx.Done(): + return + } + } + }() - spindleConn, _, err := websocket.DefaultDialer.Dial(url, nil) - if err != nil { - l.Error("websocket dial failed", "err", err) - return - } - defer spindleConn.Close() + clientConn.SetPingHandler(func(message string) error { + err := clientConn.WriteControl(websocket.PongMessage, []byte(message), time.Now().Add(60*time.Second)) + if err == websocket.ErrCloseSent { + return nil + } + return err + }) - // create a channel for incoming messages - evChan := make(chan LogEvent, 100) - // start a goroutine to read from spindle - go ReadLogs(spindleConn, evChan) + // Start a goroutine to read messages from the client and discard them. + go func() { + for { + if _, _, err := clientConn.ReadMessage(); err != nil { + cancel() + return + } + } + }() + // Main loop: sole writer of data frames to the client. stepStartTimes := make(map[int]time.Time) stepAnsi := make(map[int]*ansiState) var fragment bytes.Buffer @@ -272,61 +304,55 @@ return case ev, ok := <-evChan: if !ok { - continue - } - - if ev.Err != nil && ev.IsCloseError() { - l.Debug("graceful shutdown, tail complete", "err", err) - return - } - if ev.Err != nil { - l.Error("error reading from spindle", "err", err) + // Stream ended: Shutdown closed the upstream channel. + if err := <-done; !isExpectedClose(err) { + l.Error("spindle stream error", "err", err) + } return } - var logLine spindlemodel.LogLine - if err = json.Unmarshal(ev.Msg, &logLine); err != nil { - l.Error("failed to parse logline", "err", err) - continue - } + fragment.Reset() - fragment.Reset() + switch { + case ev.Error != nil: + l.Error("spindle error frame", "err", ev.Error.Error, "msg", ev.Error.Message) + return - switch logLine.Kind { - case spindlemodel.LogKindControl: - switch logLine.StepStatus { - case spindlemodel.StepStatusStart: - stepStartTimes[logLine.StepId] = logLine.Time - collapsed := false - if logLine.StepKind == spindlemodel.StepKindSystem { - collapsed = true - } + case ev.Control != nil: + c := ev.Control + step := int(c.Step) + switch derefStr(c.Status) { + case "start": + t := parseRFC3339(c.Time) + stepStartTimes[step] = t + // "system" steps are injected by the CI runner; collapse them. + collapsed := derefStr(c.Kind) == "system" err = p.pages.LogBlock(&fragment, pages.LogBlockParams{ - Id: logLine.StepId, - Name: logLine.Content, - Command: logLine.StepCommand, + Id: step, + Name: c.Content, + Command: derefStr(c.Command), Collapsed: collapsed, - StartTime: logLine.Time, + StartTime: t, }) - case spindlemodel.StepStatusEnd: - startTime := stepStartTimes[logLine.StepId] - endTime := logLine.Time + case "end": err = p.pages.LogBlockEnd(&fragment, pages.LogBlockEndParams{ - Id: logLine.StepId, - StartTime: startTime, - EndTime: endTime, + Id: step, + StartTime: stepStartTimes[step], + EndTime: parseRFC3339(c.Time), }) } - case spindlemodel.LogKindData: - ansi, ok := stepAnsi[logLine.StepId] + case ev.Data != nil: + d := ev.Data + step := int(d.Step) + ansi, ok := stepAnsi[step] if !ok { ansi = NewAnsiState() - stepAnsi[logLine.StepId] = ansi + stepAnsi[step] = ansi } err = p.pages.LogLine(&fragment, pages.LogLineParams{ - Id: logLine.StepId, - Content: ansi.Render(logLine.Content), + Id: step, + Content: ansi.Render(d.Content), }) } if err != nil { @@ -338,95 +364,48 @@ if err = clientConn.WriteMessage(websocket.TextMessage, fragment.Bytes()); err != nil { l.Error("error writing to client", "err", err) return } - - case _, ok := <-statusCh: - if !ok { - continue - } - fresh, err := db.GetPipelineStatuses( - p.db, - 1, - orm.FilterEq("p.repo_did", f.RepoDid), - orm.FilterEq("p.id", pipelineId), - ) - if err != nil || len(fresh) == 0 { - continue - } - for name, ws := range fresh[0].Statuses { - fragment.Reset() - if err = p.pages.WorkflowSymbolOOB(&fragment, pages.WorkflowSymbolOOBParams{ - Name: name, - Statuses: ws, - }); err != nil { - l.Error("failed to render workflow symbol OOB", "err", err) - continue - } - if err = clientConn.WriteMessage(websocket.TextMessage, fragment.Bytes()); err != nil { - l.Error("error writing workflow symbol to client", "err", err) - return - } - } - - case <-time.After(30 * time.Second): - l.Debug("sent keepalive") - if err = clientConn.WriteControl(websocket.PingMessage, []byte{}, time.Now().Add(time.Second)); err != nil { - l.Error("failed to write control", "err", err) - return - } + lastWriteLk.Lock() + lastWrite = time.Now() + lastWriteLk.Unlock() } } } -func (p *Pipelines) Cancel(w http.ResponseWriter, r *http.Request) { - l := p.logger.With("handler", "Cancel") - - var ( - pipelineId = chi.URLParam(r, "pipeline") - workflow = chi.URLParam(r, "workflow") - ) - if pipelineId == "" || workflow == "" { - http.Error(w, "missing pipeline ID or workflow", http.StatusBadRequest) - return - } +func (p *Pipelines) CancelWorkflow(w http.ResponseWriter, r *http.Request) { + l := p.logger.With("handler", "CancelWorkflow") + errorId := "workflow-error" f, err := p.repoResolver.Resolve(r) if err != nil { l.Error("failed to get repo and knot", "err", err) - http.Error(w, "bad repo/knot", http.StatusBadRequest) + p.pages.Notice(w, errorId, "Failed to cancel workflow") return } + l = l.With("repo", f.RepoDid) - pipeline, err := func() (models.Pipeline, error) { - ps, err := db.GetPipelineStatuses( - p.db, - 1, - orm.FilterEq("p.repo_did", f.RepoDid), - orm.FilterEq("p.id", pipelineId), - ) - if err != nil { - return models.Pipeline{}, err - } - if len(ps) != 1 { - return models.Pipeline{}, fmt.Errorf("wrong pipeline count %d", len(ps)) - } - return ps[0], nil - }() + if f.Spindle == "" { + l.Debug("spindle is empty") + p.pages.Notice(w, errorId, "Failed to cancel workflow") + return + } + + pipelineId, err := syntax.ParseTID(chi.URLParam(r, "pipeline")) if err != nil { - l.Error("pipeline query failed", "err", err) - http.Error(w, "pipeline not found", http.StatusNotFound) + l.Debug("invalid pipeline id", "id", pipelineId) + p.pages.Error404(w) + return } - var ( - spindle = f.Spindle - knot = f.Knot - rkey = pipeline.Rkey - ) - if spindle == "" || knot == "" || rkey == "" { - http.Error(w, "invalid repo info", http.StatusBadRequest) + workflowName := chi.URLParam(r, "workflow") + if workflowName == "" { + l.Debug("empty workflow name") + p.pages.Error404(w) return } - hostname, noTLS, err := hostutil.ParseHostname(spindle) + l = l.With("pipeline", pipelineId, "workflow", workflowName) + + hostname, noTLS, err := hostutil.ParseHostname(f.Spindle) if err != nil { http.Error(w, "invalid spindle hostname", http.StatusBadRequest) return @@ -440,20 +419,18 @@ oauth.WithDev(noTLS), oauth.WithTimeout(time.Second*30), // workflow cleanup usually takes time ) - err = tangled.PipelineCancelPipeline( + if err := tangled.PipelineCancelPipeline( r.Context(), spindleClient, &tangled.PipelineCancelPipeline_Input{ Repo: string(f.RepoAt()), - Pipeline: pipeline.AtUri().String(), - Workflow: workflow, + Pipeline: pipelineId.String(), + Workflow: workflowName, }, - ) - errorId := "workflow-error" - if err != nil { + ); err != nil { l.Error("failed to cancel workflow", "err", err) p.pages.Notice(w, errorId, "Failed to cancel workflow") return } - l.Debug("canceled pipeline", "uri", pipeline.AtUri()) + l.Debug("canceled workflow") } diff --git a/appview/pipelines/ssh/cihelpers.go b/appview/pipelines/ssh/cihelpers.go new file mode 100644 --- /dev/null +++ b/appview/pipelines/ssh/cihelpers.go @@ -0,0 +1,43 @@ +package ssh + +import ( + "time" + + "tangled.org/core/api/tangled" +) + +// helper functions against generated code + +func workflowElapsed(wf *tangled.CiDefs_Workflow, now time.Time) time.Duration { + if wf.StartedAt == nil { + return 0 + } + started, err := time.Parse(time.RFC3339, *wf.StartedAt) + if err != nil { + return 0 + } + if wf.FinishedAt == nil { + return now.Sub(started) + } + finished, err := time.Parse(time.RFC3339, *wf.FinishedAt) + if err != nil { + return 0 + } + return finished.Sub(started) +} + +var finishedStatuses = map[string]bool{ + "failed": true, + "timeout": true, + "cancelled": true, + "success": true, +} + +func pipelineFinished(p *tangled.CiDefs_Pipeline) bool { + for _, wf := range p.Workflows { + if !finishedStatuses[wf.Status] { + return false + } + } + return true +} diff --git a/appview/pipelines/ssh/logstream.go b/appview/pipelines/ssh/logstream.go --- a/appview/pipelines/ssh/logstream.go +++ b/appview/pipelines/ssh/logstream.go @@ -1,19 +1,16 @@ package ssh import ( + "context" "time" - tea "github.com/charmbracelet/bubbletea" - "github.com/gorilla/websocket" - "tangled.org/core/appview/pipelines" - spindlemodel "tangled.org/core/spindle/models" + "tangled.org/core/api/tangled" ) type step struct { - id int + id int64 name string command string - kind spindlemodel.StepKind lines []string startTime time.Time endTime time.Time @@ -21,27 +18,30 @@ finished bool } type logDoneMsg struct { - workflow string - err error + err error } type logEventMsg struct { - workflow string - ev pipelines.LogEvent - conn *websocket.Conn - ch chan pipelines.LogEvent + ev *tangled.CiPipelineSubscribeLogs_Event + events chan *tangled.CiPipelineSubscribeLogs_Event + done chan error } -func readNextCmd(workflow string, conn *websocket.Conn, ch chan pipelines.LogEvent) tea.Cmd { - return func() tea.Msg { - return readNextLogEvent(workflow, conn, ch) - } +type eventScheduler struct { + ch chan *tangled.CiPipelineSubscribeLogs_Event +} + +func newEventScheduler() *eventScheduler { + return &eventScheduler{ch: make(chan *tangled.CiPipelineSubscribeLogs_Event, 1024)} } -func readNextLogEvent(workflow string, conn *websocket.Conn, ch chan pipelines.LogEvent) tea.Msg { - ev, ok := <-ch - if !ok { - return logDoneMsg{workflow: workflow} +func (s *eventScheduler) AddWork(ctx context.Context, _ string, v *tangled.CiPipelineSubscribeLogs_Event) error { + select { + case s.ch <- v: + return nil + case <-ctx.Done(): + return ctx.Err() } - return logEventMsg{workflow: workflow, ev: ev, conn: conn, ch: ch} } + +func (s *eventScheduler) Shutdown() { close(s.ch) } diff --git a/appview/pipelines/ssh/server.go b/appview/pipelines/ssh/server.go --- a/appview/pipelines/ssh/server.go +++ b/appview/pipelines/ssh/server.go @@ -10,18 +10,16 @@ "github.com/charmbracelet/wish" tea "github.com/charmbracelet/wish/bubbletea" "tangled.org/core/appview/config" "tangled.org/core/appview/db" - "tangled.org/core/appview/pipelines" ) type Server struct { - db *db.DB - config *config.Config - pipelineNotifier *pipelines.StatusNotifier - logger *slog.Logger + db *db.DB + config *config.Config + logger *slog.Logger } -func New(db *db.DB, cfg *config.Config, pn *pipelines.StatusNotifier, logger *slog.Logger) *Server { - return &Server{db: db, config: cfg, pipelineNotifier: pn, logger: logger} +func New(db *db.DB, cfg *config.Config, logger *slog.Logger) *Server { + return &Server{db: db, config: cfg, logger: logger} } func (s *Server) ListenAndServe(ctx context.Context) error { diff --git a/appview/pipelines/ssh/session.go b/appview/pipelines/ssh/session.go --- a/appview/pipelines/ssh/session.go +++ b/appview/pipelines/ssh/session.go @@ -3,10 +3,15 @@ import ( "fmt" + "github.com/bluesky-social/indigo/atproto/syntax" + indigoxrpc "github.com/bluesky-social/indigo/xrpc" tea "github.com/charmbracelet/bubbletea" "github.com/charmbracelet/ssh" wishtea "github.com/charmbracelet/wish/bubbletea" + "tangled.org/core/api/tangled" "tangled.org/core/appview/db" + "tangled.org/core/hostutil" + extlexutil "tangled.org/core/lexutil" "tangled.org/core/orm" ) @@ -29,18 +34,39 @@ sha := args[1] l = l.With("repoDID", repoDID, "sha", sha) - pipelines, err := db.GetPipelineStatuses(s.db, 1, - orm.FilterEq("p.repo_did", repoDID), - orm.FilterEq("p.sha", sha), - ) - if err != nil || len(pipelines) == 0 { + repo, err := db.GetRepo(s.db, orm.FilterEq("repo_did", repoDID)) + if err != nil { + l.Warn("repo not found", "err", err) + return newErrorModel(renderer, fmt.Sprintf("repo %s not found", repoDID)), wishtea.MakeOptions(sess) + } + if repo.Spindle == "" { + l.Warn("no spindle configured") + return newErrorModel(renderer, "no spindle configured for this repo"), wishtea.MakeOptions(sess) + } + + l = l.With("spindle", repo.Spindle) + + host, err := hostutil.EnsureHttpScheme(repo.Spindle) + if err != nil { + l.Warn("invalid spindlie hostname", "err", err) + return newErrorModel(renderer, fmt.Sprintf("invalid spindle host %q", repo.Spindle)), wishtea.MakeOptions(sess) + } + + xrpcc := extlexutil.Client{Client: indigoxrpc.Client{Host: host}} + out, err := tangled.CiQueryPipelines(sess.Context(), &xrpcc, []string{sha}, "", 1, repoDID) + if err != nil || len(out.Pipelines) == 0 { l.Warn("pipeline not found", "err", err) return newErrorModel(renderer, fmt.Sprintf("pipeline not found for repo %s @ %s", repoDID, sha)), wishtea.MakeOptions(sess) } - pipeline := pipelines[0] - l.Info("serving pipeline", "workflows", len(pipeline.Statuses)) + pipeline := out.Pipelines[0] + if _, err := syntax.ParseTID(pipeline.Id); err != nil { + l.Warn("invalid pipeline id", "id", pipeline.Id, "err", err) + return newErrorModel(renderer, fmt.Sprintf("invalid pipeline id %q", pipeline.Id)), wishtea.MakeOptions(sess) + } + + l.Info("serving pipeline", "pipeline", pipeline.Id, "workflows", len(pipeline.Workflows)) pty, _, _ := sess.Pty() opts := append(wishtea.MakeOptions(sess), tea.WithAltScreen()) - return newPipelineModel(renderer, s, pipeline, pty.Window.Width, pty.Window.Height), opts + return newPipelineModel(renderer, &xrpcc, pipeline, pty.Window.Width, pty.Window.Height), opts } diff --git a/appview/pipelines/ssh/tui.go b/appview/pipelines/ssh/tui.go --- a/appview/pipelines/ssh/tui.go +++ b/appview/pipelines/ssh/tui.go @@ -1,7 +1,8 @@ package ssh import ( - "encoding/json" + "context" + "errors" "fmt" "strings" "time" @@ -11,11 +12,8 @@ "github.com/charmbracelet/bubbles/viewport" tea "github.com/charmbracelet/bubbletea" "github.com/charmbracelet/lipgloss" "github.com/gorilla/websocket" - "tangled.org/core/appview/db" - "tangled.org/core/appview/models" - "tangled.org/core/appview/pipelines" - "tangled.org/core/orm" - spindlemodel "tangled.org/core/spindle/models" + "tangled.org/core/api/tangled" + extlexutil "tangled.org/core/lexutil" ) var ( @@ -27,100 +25,104 @@ type tickMsg time.Time type statusUpdateMsg struct { - pipeline models.Pipeline + pipeline *tangled.CiDefs_Pipeline } type statusUpdateErrMsg struct{ err error } type pipelineModel struct { - renderer *lipgloss.Renderer - server *Server - pipeline models.Pipeline - workflows []string - selected int - logs map[string]*workflowLogs - statusCh chan struct{} - spinner spinner.Model - width int - height int + renderer *lipgloss.Renderer + xrpcc *extlexutil.Client + pipeline *tangled.CiDefs_Pipeline + selected int + logs map[string]*workflowLogs + + // pipeline log stream: cancel tears down the consumer goroutine on quit. + // the event/done channels are threaded through log messages, not stored here. + cancel context.CancelFunc + streamDone bool + streamErr error + + spinner spinner.Model + width int + height int } type workflowLogs struct { steps []step - stepIndex map[int]int + stepIndex map[int64]int // stepId -> index map vp viewport.Model ready bool - done bool - err error } -func newPipelineModel(renderer *lipgloss.Renderer, s *Server, pipeline models.Pipeline, width, height int) *pipelineModel { - workflows := pipeline.Workflows() - logs := make(map[string]*workflowLogs, len(workflows)) - for _, wf := range workflows { - logs[wf] = &workflowLogs{stepIndex: make(map[int]int)} +func newPipelineModel(renderer *lipgloss.Renderer, xrpcc *extlexutil.Client, pipeline *tangled.CiDefs_Pipeline, width, height int) *pipelineModel { + logs := make(map[string]*workflowLogs, len(pipeline.Workflows)) + for _, wf := range pipeline.Workflows { + logs[wf.Name] = &workflowLogs{stepIndex: make(map[int64]int)} } - statusCh := s.pipelineNotifier.Subscribe(pipeline.AtUri()) sp := spinner.New(spinner.WithSpinner(spinner.Line)) return &pipelineModel{ - renderer: renderer, - server: s, - pipeline: pipeline, - workflows: workflows, - logs: logs, - statusCh: statusCh, - spinner: sp, - width: width, - height: height, + renderer: renderer, + xrpcc: xrpcc, + pipeline: pipeline, + logs: logs, + spinner: sp, + width: width, + height: height, } } func (m *pipelineModel) Init() tea.Cmd { - cmds := []tea.Cmd{tick(), m.spinner.Tick, m.waitForStatusUpdate(m.statusCh)} - for _, wf := range m.workflows { - cmds = append(cmds, m.connectCmd(wf)) - } - return tea.Batch(cmds...) + return tea.Batch(tick(), m.spinner.Tick, m.subscribeCmd()) } func tick() tea.Cmd { return tea.Tick(time.Second, func(t time.Time) tea.Msg { return tickMsg(t) }) } -// waitForStatusUpdate blocks on the notifier channel, re-fetches pipeline statuses, and returns the result as a tea.Msg. -func (m *pipelineModel) waitForStatusUpdate(ch chan struct{}) tea.Cmd { - knot := m.pipeline.Knot - rkey := m.pipeline.Rkey +// subscribeCmd opens the ci.pipeline.subscribeLogs stream for current pipeline. +// A consumer goroutine pushes decoded events onto the scheduler channel; the +// returned command yields the first event into the bubbletea loop. +func (m *pipelineModel) subscribeCmd() tea.Cmd { + // cancel existing subscriptions just in case + if m.cancel != nil { + m.cancel() + } + sched := newEventScheduler() + done := make(chan error, 1) + ctx, cancel := context.WithCancel(context.Background()) + m.cancel = cancel + + pipelineId := m.pipeline.Id + go func() { + err := tangled.CiPipelineSubscribeLogs(ctx, m.xrpcc, pipelineId, nil, sched) + done <- err + }() + + return readEventCmd(sched.ch, done) +} + +func readEventCmd(events chan *tangled.CiPipelineSubscribeLogs_Event, done chan error) tea.Cmd { return func() tea.Msg { - if _, ok := <-ch; !ok { - return nil - } - ps, err := db.GetPipelineStatuses(m.server.db, 1, - orm.FilterEq("p.knot", knot), - orm.FilterEq("p.rkey", rkey), - ) - if err != nil || len(ps) == 0 { - return statusUpdateErrMsg{err: fmt.Errorf("refreshing pipeline: %w", err)} + ev, ok := <-events + if !ok { + return logDoneMsg{err: <-done} } - return statusUpdateMsg{pipeline: ps[0]} + return logEventMsg{ev: ev, events: events, done: done} } } -// connectCmd dials the spindle websocket for the given workflow and starts streaming log events. -func (m *pipelineModel) connectCmd(workflow string) tea.Cmd { +// fetchStatusCmd re-fetches the pipeline +func (m *pipelineModel) fetchStatusCmd() tea.Cmd { + pipelineId := m.pipeline.Id return func() tea.Msg { - ws, ok := m.pipeline.Statuses[workflow] - if !ok || len(ws.Data) == 0 { - return logDoneMsg{workflow: workflow} - } - url := pipelines.SpindleURL(ws.Data[0].Spindle, m.pipeline.Knot, m.pipeline.Rkey, workflow) - conn, _, err := websocket.DefaultDialer.Dial(url, nil) + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + out, err := tangled.CiGetPipeline(ctx, m.xrpcc, pipelineId) if err != nil { - return logDoneMsg{workflow: workflow, err: fmt.Errorf("connecting to spindle: %w", err)} + return statusUpdateErrMsg{err: fmt.Errorf("refreshing pipeline: %w", err)} } - ch := make(chan pipelines.LogEvent, 100) - go pipelines.ReadLogs(conn, ch) - return readNextLogEvent(workflow, conn, ch) + return statusUpdateMsg{pipeline: out} } } @@ -153,6 +155,8 @@ m.width, m.height = msg.Width, msg.Height m.resizeViewports() case tickMsg: + // re-render running workflows so elapsed times advance + m.refreshRunning() return m, tick() case spinner.TickMsg: @@ -163,13 +167,15 @@ case tea.KeyMsg: switch msg.String() { case "q", "ctrl+c": - m.server.pipelineNotifier.Unsubscribe(m.pipeline.AtUri(), m.statusCh) + if m.cancel != nil { + m.cancel() + } return m, tea.Quit case "tab", "right", "l": - m.selected = (m.selected + 1) % len(m.workflows) + m.selected = (m.selected + 1) % len(m.pipeline.Workflows) return m, nil case "shift+tab", "left", "h": - m.selected = (m.selected - 1 + len(m.workflows)) % len(m.workflows) + m.selected = (m.selected - 1 + len(m.pipeline.Workflows)) % len(m.pipeline.Workflows) return m, nil } if wl := m.selectedLogs(); wl != nil && wl.ready { @@ -193,48 +199,46 @@ return m, cmd } case logEventMsg: - return m, m.handleLogEvent(msg) + m.applyEvent(msg.ev) + return m, readEventCmd(msg.events, msg.done) case logDoneMsg: - if wl, ok := m.logs[msg.workflow]; ok { - wl.done, wl.err = true, msg.err - m.initViewport(wl) - wl.vp.SetContent(renderLogs(m.renderer, wl, m.width)) - wl.vp.GotoBottom() + m.streamDone = true + if !isExpectedClose(msg.err) { + m.streamErr = msg.err } + m.refreshAll() + // resolve final workflow statuses once now that the stream has ended + return m, m.fetchStatusCmd() case statusUpdateMsg: - // detect any workflows that are new since the last update - known := make(map[string]bool, len(m.workflows)) - for _, wf := range m.workflows { - known[wf] = true + m.pipeline = msg.pipeline + known := make(map[string]bool, len(m.pipeline.Workflows)) + for _, wf := range m.pipeline.Workflows { + known[wf.Name] = true } - m.pipeline = msg.pipeline - var newCmds []tea.Cmd - for _, wf := range msg.pipeline.Workflows() { - if !known[wf] { - m.workflows = append(m.workflows, wf) - m.logs[wf] = &workflowLogs{stepIndex: make(map[int]int)} - newCmds = append(newCmds, m.connectCmd(wf)) + for name := range m.logs { + if !known[name] { + delete(m.logs, name) } } - // re-subscribe for the next update - newCmds = append(newCmds, m.waitForStatusUpdate(m.statusCh)) - return m, tea.Batch(newCmds...) + if m.selected >= len(m.pipeline.Workflows) { + m.selected = max(len(m.pipeline.Workflows)-1, 0) + } + m.refreshAll() case statusUpdateErrMsg: - // re-subscribe even on error so we don't stop listening - return m, m.waitForStatusUpdate(m.statusCh) + // best-effort final status refresh; ignore failures } return m, nil } func (m *pipelineModel) selectedLogs() *workflowLogs { - if len(m.workflows) == 0 { + if len(m.pipeline.Workflows) < 1+m.selected { return nil } - return m.logs[m.workflows[m.selected]] + return m.logs[m.pipeline.Workflows[m.selected].Name] } func (m *pipelineModel) initViewport(wl *workflowLogs) { @@ -245,65 +249,107 @@ wl.vp = viewport.New(m.width, m.vpHeight()) wl.ready = true } -// handleLogEvent processes a single log event, updates the step state, and re-renders the viewport. -func (m *pipelineModel) handleLogEvent(msg logEventMsg) tea.Cmd { - wl, ok := m.logs[msg.workflow] +// ensureWorkflow returns the log state for a workflow, lazily creating its routing entry. +func (m *pipelineModel) ensureWorkflow(name string) *workflowLogs { + wl, ok := m.logs[name] if !ok { - return nil + wl = &workflowLogs{stepIndex: make(map[int64]int)} + m.logs[name] = wl } - if msg.ev.Err != nil { - wl.done = true - if !msg.ev.IsCloseError() { - wl.err = msg.ev.Err - } - m.initViewport(wl) - wl.vp.SetContent(renderLogs(m.renderer, wl, m.width)) - return nil - } - var line spindlemodel.LogLine - if err := json.Unmarshal(msg.ev.Msg, &line); err != nil { - return readNextCmd(msg.workflow, msg.conn, msg.ch) - } - applyLogLine(wl, line) + return wl +} + +// renderWorkflow re-renders a workflow's viewport, preserving bottom-stickiness. +func (m *pipelineModel) renderWorkflow(wl *workflowLogs) { m.initViewport(wl) atBottom := wl.vp.AtBottom() wl.vp.SetContent(renderLogs(m.renderer, wl, m.width)) if atBottom { wl.vp.GotoBottom() } - return readNextCmd(msg.workflow, msg.conn, msg.ch) +} + +// refreshAll re-renders every initialized viewport. +func (m *pipelineModel) refreshAll() { + for _, wl := range m.logs { + m.renderWorkflow(wl) + } } -// applyLogLine mutates wl by appending the log line to the appropriate step. -func applyLogLine(wl *workflowLogs, line spindlemodel.LogLine) { - switch line.Kind { - case spindlemodel.LogKindControl: - switch line.StepStatus { - case spindlemodel.StepStatusStart: - idx := len(wl.steps) - wl.stepIndex[line.StepId] = idx +// refreshRunning re-renders workflows with unfinished steps so elapsed times advance. +func (m *pipelineModel) refreshRunning() { + if m.streamDone { + return + } + for _, wl := range m.logs { + if !wl.ready { + continue + } + for i := range wl.steps { + if !wl.steps[i].finished { + m.renderWorkflow(wl) + break + } + } + } +} + +// applyEvent routes a decoded subscribeLogs event into the matching workflow. +func (m *pipelineModel) applyEvent(ev *tangled.CiPipelineSubscribeLogs_Event) { + switch { + case ev.Error != nil: + if ev.Error.Message != "" { + m.streamErr = fmt.Errorf("%s: %s", ev.Error.Error, ev.Error.Message) + } else { + m.streamErr = fmt.Errorf("%s", ev.Error.Error) + } + + case ev.Control != nil: + c := ev.Control + wl := m.ensureWorkflow(c.Workflow) + switch derefStr(c.Status) { + case "start": + wl.stepIndex[c.Step] = len(wl.steps) wl.steps = append(wl.steps, step{ - id: line.StepId, name: line.Content, command: line.StepCommand, - kind: line.StepKind, startTime: line.Time, + id: c.Step, name: c.Content, command: derefStr(c.Command), startTime: parseRFC3339(c.Time), }) - case spindlemodel.StepStatusEnd: - if idx, ok := wl.stepIndex[line.StepId]; ok { - wl.steps[idx].endTime, wl.steps[idx].finished = line.Time, true + case "end": + if idx, ok := wl.stepIndex[c.Step]; ok { + wl.steps[idx].endTime, wl.steps[idx].finished = parseRFC3339(c.Time), true } } - case spindlemodel.LogKindData: - if idx, ok := wl.stepIndex[line.StepId]; ok { - wl.steps[idx].lines = append(wl.steps[idx].lines, line.Content) + m.renderWorkflow(wl) + + case ev.Data != nil: + d := ev.Data + wl := m.ensureWorkflow(d.Workflow) + if idx, ok := wl.stepIndex[d.Step]; ok { + wl.steps[idx].lines = append(wl.steps[idx].lines, d.Content) } + m.renderWorkflow(wl) } } +// isExpectedClose reports whether err is a clean websocket close (or nil). +func isExpectedClose(err error) bool { + if err == nil { + return true + } + var ce *websocket.CloseError + if errors.As(err, &ce) { + switch ce.Code { + case websocket.CloseNormalClosure, websocket.CloseGoingAway, websocket.CloseAbnormalClosure: + return true + } + } + return false +} + // renderLogs builds the full log content string for a workflow, used as viewport content. func renderLogs(r *lipgloss.Renderer, wl *workflowLogs, width int) string { headerStyle := r.NewStyle().Foreground(colorFg).Bold(true) cmdStyle := r.NewStyle().Foreground(colorBlue).Width(width) dimStyle := r.NewStyle().Faint(true) - now := time.Now() var sb strings.Builder for i := range wl.steps { st := &wl.steps[i] @@ -311,7 +357,7 @@ dur := "" if st.finished { dur = st.endTime.Sub(st.startTime).Round(time.Millisecond).String() } else if !st.startTime.IsZero() { - dur = now.Sub(st.startTime).Round(time.Second).String() + dur = time.Since(st.startTime).Round(time.Second).String() } // build overlay: "── name ──...── dur ──" nameStr := headerStyle.Render(st.name + " ") @@ -330,9 +376,6 @@ sb.WriteString(l + "\n") } sb.WriteString("\n") } - if wl.done && wl.err != nil { - sb.WriteString("error: " + wl.err.Error() + "\n") - } return sb.String() } @@ -340,6 +383,9 @@ func (m *pipelineModel) View() string { body := "" if wl := m.selectedLogs(); wl != nil && wl.ready { body = wl.vp.View() + } + if m.streamErr != nil { + body = lipgloss.JoinVertical(lipgloss.Left, body, m.renderer.NewStyle().Foreground(colorBlue).Render("stream error: "+m.streamErr.Error())) } return lipgloss.JoinVertical(lipgloss.Left, m.topbarView(), "", body) } @@ -352,20 +398,10 @@ now := time.Now() var tabs strings.Builder - for i, wf := range m.workflows { - status := spindlemodel.StatusKindPending - elapsed := "" - if ws, ok := m.pipeline.Statuses[wf]; ok { - latest := ws.Latest() - status = latest.Status - if t := ws.TimeTaken(); t > 0 { - elapsed = t.Round(time.Second).String() - } else { - elapsed = now.Sub(latest.Created).Round(time.Second).String() - } - } - dim := r.NewStyle().Faint(true) - base := " " + statusIcon(status, m.spinner.View()) + " " + wf + for i, wf := range m.pipeline.Workflows { + status := wf.Status + elapsed := workflowElapsed(wf, now).Round(time.Second).String() + base := " " + statusIcon(status, m.spinner.View()) + " " + wf.Name if i == m.selected { tab := base if elapsed != "" { @@ -376,6 +412,7 @@ tabs.WriteString(activeStyle.Render(tab)) } else { tabs.WriteString(base) if elapsed != "" { + dim := r.NewStyle().Faint(true) tabs.WriteString(" " + dim.Render(elapsed)) } tabs.WriteString(" ") @@ -383,7 +420,7 @@ } } tabsStr := tabs.String() - infoStr := triggerLine(r, m.pipeline.Trigger, m.pipeline.Sha) + " · " + helpText(r) + infoStr := triggerLine(r, m.pipeline.Trigger, m.pipeline.Commit) + " · " + helpText(r) gap := max(m.width-lipgloss.Width(tabsStr)-lipgloss.Width(infoStr), 1) @@ -410,40 +447,55 @@ } return sha } -func triggerLine(r *lipgloss.Renderer, t *models.Trigger, sha string) string { +func triggerLine(r *lipgloss.Renderer, t *tangled.CiDefs_Pipeline_Trigger, sha string) string { hash := shortSha(sha) dim := r.NewStyle().Faint(true) if t == nil { return dim.Render(hash) } - if t.IsPush() { - return t.TargetRef() + dim.Render("@"+hash) + dim.Render(" (push)") + if t.CiTrigger_Push != nil { + return t.CiTrigger_Push.Ref + dim.Render("@"+hash) + dim.Render(" (push)") } - if t.IsPullRequest() { + if t.CiTrigger_PullRequest != nil { source := "" - if t.PRSourceBranch != nil { - source = *t.PRSourceBranch + if t.CiTrigger_PullRequest.SourceBranch != nil { + source = *t.CiTrigger_PullRequest.SourceBranch } - return t.TargetRef() + dim.Render(" <- "+source+"@"+hash) + dim.Render(" (pull-request)") + return t.CiTrigger_PullRequest.TargetBranch + dim.Render(" <- "+source+"@"+hash) + dim.Render(" (pull-request)") } return dim.Render(hash) } -func statusIcon(status spindlemodel.StatusKind, spinnerFrame string) string { +func statusIcon(status string, spinnerFrame string) string { switch status { - case spindlemodel.StatusKindSuccess: + case "success": return "✓" - case spindlemodel.StatusKindFailed: + case "failed": return "×" - case spindlemodel.StatusKindRunning: + case "running": return spinnerFrame - case spindlemodel.StatusKindPending: + case "pending": return "·" - case spindlemodel.StatusKindTimeout: + case "timeout": return "⌀" - case spindlemodel.StatusKindCancelled: + case "cancelled": return "-" default: return "?" } } + +func derefStr(s *string) string { + if s == nil { + return "" + } + return *s +} + +func parseRFC3339(s string) time.Time { + t, err := time.Parse(time.RFC3339, s) + if err != nil { + return time.Time{} + } + return t +} diff --git a/appview/pulls/list.go b/appview/pulls/list.go --- a/appview/pulls/list.go +++ b/appview/pulls/list.go @@ -14,6 +14,7 @@ "tangled.org/core/appview/searchquery" "tangled.org/core/orm" "github.com/bluesky-social/indigo/atproto/syntax" + indigoxrpc "github.com/bluesky-social/indigo/xrpc" ) func (s *Pulls) RepoPulls(w http.ResponseWriter, r *http.Request) { @@ -261,20 +262,24 @@ slices.Reverse(stack) stacks = append(stacks, stack) } - ps, err := db.GetPipelineStatuses( - s.db, - len(shas), - orm.FilterEq("p.repo_did", f.RepoDid), - orm.FilterIn("p.sha", shas), - ) - if err != nil { - l.Warn("failed to fetch pipeline statuses", "err", err) - // non-fatal - } - m := make(map[string]models.Pipeline) - for _, p := range ps { - m[p.Sha] = p - } + // commitId -> latest pipeline + pipelines := func(ctx context.Context, shas []string) map[string]tangled.CiDefs_Pipeline { + xrpcc := &indigoxrpc.Client{Host: f.Spindle} + out, err := tangled.CiQueryPipelines(ctx, xrpcc, shas, "", 0, f.RepoDid) + if err != nil { + l.Error("failed to fetch pipelines", "err", err) + } + + m := make(map[string]tangled.CiDefs_Pipeline) + + for _, pipeline := range out.Pipelines { + if pipeline == nil { + continue + } + m[pipeline.Commit] = *pipeline + } + return m + }(r.Context(), shas) labelDefs, err := db.GetLabelDefinitions( s.db, @@ -317,7 +322,7 @@ LabelDefs: defs, FilterState: filterState, FilterQuery: query.String(), Stacks: stacks, - Pipelines: m, + Pipelines: pipelines, Page: page, PullCount: totalPulls, VouchRelationships: vouchRelationships, diff --git a/appview/pulls/single.go b/appview/pulls/single.go --- a/appview/pulls/single.go +++ b/appview/pulls/single.go @@ -1,6 +1,7 @@ package pulls import ( + "context" "fmt" "net/http" "strconv" @@ -150,8 +151,6 @@ // can be nil if this pull is not stacked stack, _ := r.Context().Value("stack").(models.Stack) - m := make(map[string]models.Pipeline) - var shas []string for _, s := range pull.Submissions { shas = append(shas, s.SourceRev) @@ -160,20 +159,24 @@ for _, p := range stack { shas = append(shas, p.LatestSha()) } - ps, err := db.GetPipelineStatuses( - s.db, - len(shas), - orm.FilterEq("p.repo_did", f.RepoDid), - orm.FilterIn("p.sha", shas), - ) - if err != nil { - l.Error("failed to fetch pipeline statuses", "err", err) - // non-fatal - } + // commitId -> latest pipeline + pipelines := func(ctx context.Context) map[string]tangled.CiDefs_Pipeline { + xrpcc := &indigoxrpc.Client{Host: f.Spindle} + out, err := tangled.CiQueryPipelines(ctx, xrpcc, shas, "", 0, f.RepoDid) + if err != nil { + l.Error("failed to fetch pipelines", "err", err) + } - for _, p := range ps { - m[p.Sha] = p - } + m := make(map[string]tangled.CiDefs_Pipeline) + + for _, pipeline := range out.Pipelines { + if pipeline == nil { + continue + } + m[pipeline.Commit] = *pipeline + } + return m + }(r.Context()) entities := []syntax.ATURI{pull.AtUri()} for _, s := range pull.Submissions { @@ -256,7 +259,7 @@ Backlinks: backlinks, BranchDeleteStatus: nil, MergeCheck: types.MergeCheckResponse{}, ResubmitCheck: pages.Unknown, - Pipelines: m, + Pipelines: pipelines, Diff: diff, DiffOpts: diffOpts, ActiveRound: roundIdInt, diff --git a/appview/repo/index.go b/appview/repo/index.go --- a/appview/repo/index.go +++ b/appview/repo/index.go @@ -177,7 +177,7 @@ var shas []string for _, c := range commitsTrunc { shas = append(shas, c.Hash.String()) } - pipelines, err := getPipelineStatuses(rp.db, f, shas) + pipelines, err := getPipelineStatuses(r.Context(), f, shas) if err != nil { l.Error("failed to fetch pipeline statuses", "err", err) // non-fatal diff --git a/appview/repo/log.go b/appview/repo/log.go --- a/appview/repo/log.go +++ b/appview/repo/log.go @@ -11,7 +11,6 @@ "tangled.org/core/api/tangled" "tangled.org/core/appview/commitverify" "tangled.org/core/appview/db" - "tangled.org/core/appview/models" "tangled.org/core/appview/pages" xrpcclient "tangled.org/core/appview/xrpcclient" "tangled.org/core/types" @@ -178,7 +177,7 @@ var shas []string for _, c := range xrpcResp.Commits { shas = append(shas, c.Hash.String()) } - pipelines, err := getPipelineStatuses(rp.db, f, shas) + pipelines, err := getPipelineStatuses(r.Context(), f, shas) if err != nil { l.Error("failed to getPipelineStatuses", "err", err) // non-fatal @@ -250,14 +249,14 @@ l.Error("failed to GetVerifiedCommits", "err", err) } user := rp.oauth.GetMultiAccountUser(r) - pipelines, err := getPipelineStatuses(rp.db, f, []string{result.Diff.Commit.This}) + pipelines, err := getPipelineStatuses(r.Context(), f, []string{result.Diff.Commit.This}) if err != nil { l.Error("failed to getPipelineStatuses", "err", err) // non-fatal } - var pipeline *models.Pipeline + var pipeline *tangled.CiDefs_Pipeline if p, ok := pipelines[result.Diff.Commit.This]; ok { - pipeline = &p + pipeline = p } rp.pages.RepoCommit(w, pages.RepoCommitParams{ diff --git a/appview/repo/repo.go b/appview/repo/repo.go --- a/appview/repo/repo.go +++ b/appview/repo/repo.go @@ -29,7 +29,6 @@ "tangled.org/core/appview/reporesolver" "tangled.org/core/appview/sites" xrpcclient "tangled.org/core/appview/xrpcclient" "tangled.org/core/consts" - "tangled.org/core/eventconsumer" "tangled.org/core/idresolver" "tangled.org/core/ogre" "tangled.org/core/orm" @@ -46,28 +45,26 @@ "github.com/go-chi/chi/v5" ) type Repo struct { - repoResolver *reporesolver.RepoResolver - idResolver *idresolver.Resolver - config *config.Config - oauth *oauth.OAuth - pages *pages.Pages - spindlestream *eventconsumer.Consumer - db *db.DB - enforcer *rbac.Enforcer - acl *knotacl.Service - notifier notify.Notifier - logger *slog.Logger - serviceAuth *serviceauth.ServiceAuth - cfClient *cloudflare.Client - ogreClient *ogre.Client - codesearch *codesearch.CodeSearch + repoResolver *reporesolver.RepoResolver + idResolver *idresolver.Resolver + config *config.Config + oauth *oauth.OAuth + pages *pages.Pages + db *db.DB + enforcer *rbac.Enforcer + acl *knotacl.Service + notifier notify.Notifier + logger *slog.Logger + serviceAuth *serviceauth.ServiceAuth + cfClient *cloudflare.Client + ogreClient *ogre.Client + codesearch *codesearch.CodeSearch } func New( oauth *oauth.OAuth, repoResolver *reporesolver.RepoResolver, pages *pages.Pages, - spindlestream *eventconsumer.Consumer, idResolver *idresolver.Resolver, db *db.DB, config *config.Config, @@ -79,20 +76,19 @@ cfClient *cloudflare.Client, codesearch *codesearch.CodeSearch, ) *Repo { return &Repo{ - oauth: oauth, - repoResolver: repoResolver, - pages: pages, - idResolver: idResolver, - config: config, - spindlestream: spindlestream, - db: db, - notifier: notifier, - enforcer: enforcer, - acl: acl, - logger: logger, - cfClient: cfClient, - ogreClient: ogre.NewClient(config.Ogre.Host), - codesearch: codesearch, + oauth: oauth, + repoResolver: repoResolver, + pages: pages, + idResolver: idResolver, + config: config, + db: db, + notifier: notifier, + enforcer: enforcer, + acl: acl, + logger: logger, + cfClient: cfClient, + ogreClient: ogre.NewClient(config.Ogre.Host), + codesearch: codesearch, } } @@ -171,23 +167,6 @@ if err != nil { fail("Failed to update spindle, unable to save to PDS.", err) return - } - - oldSpindle := f.Spindle - if oldSpindle != "" && oldSpindle != newSpindle { - remaining, qErr := db.GetRepos(rp.db, orm.FilterEq("spindle", oldSpindle)) - if qErr != nil { - l.Warn("failed to count repos using old spindle", "err", qErr) - } else if len(remaining) == 0 { - rp.spindlestream.RemoveSource(eventconsumer.NewSpindleSource(oldSpindle)) - } - } - - if !removingSpindle { - rp.spindlestream.AddSource( - context.Background(), - eventconsumer.NewSpindleSource(newSpindle), - ) } rp.pages.HxRefresh(w) diff --git a/appview/repo/repo_util.go b/appview/repo/repo_util.go --- a/appview/repo/repo_util.go +++ b/appview/repo/repo_util.go @@ -1,14 +1,15 @@ package repo import ( + "context" "maps" "slices" "sort" "strings" - "tangled.org/core/appview/db" + indigoxrpc "github.com/bluesky-social/indigo/xrpc" + "tangled.org/core/api/tangled" "tangled.org/core/appview/models" - "tangled.org/core/orm" "tangled.org/core/types" ) @@ -90,28 +91,24 @@ // grab pipelines from DB and munge that into a hashmap with commit sha as key // // golang is so blessed that it requires 35 lines of imperative code for this func getPipelineStatuses( - d *db.DB, + ctx context.Context, repo *models.Repo, shas []string, -) (map[string]models.Pipeline, error) { - m := make(map[string]models.Pipeline) +) (map[string]*tangled.CiDefs_Pipeline, error) { + m := make(map[string]*tangled.CiDefs_Pipeline) if len(shas) == 0 { return m, nil } - ps, err := db.GetPipelineStatuses( - d, - len(shas), - orm.FilterEq("p.repo_did", repo.RepoDid), - orm.FilterIn("p.sha", shas), - ) + xrpcc := &indigoxrpc.Client{Host: repo.Spindle} + out, err := tangled.CiQueryPipelines(ctx, xrpcc, shas, "", 0, repo.RepoDid) if err != nil { return nil, err } - for _, p := range ps { - m[p.Sha] = p + for _, p := range out.Pipelines { + m[p.Commit] = p } return m, nil diff --git a/appview/state/router.go b/appview/state/router.go --- a/appview/state/router.go +++ b/appview/state/router.go @@ -147,6 +147,8 @@ func (s *State) UserRouter(mw *middleware.Middleware) http.Handler { r := chi.NewRouter() r.Use(mw.InjectBaseParams) + // TODO: workflow status update requests (30s polling) + r.With(mw.ResolveIdent()).Route("/{user}", func(r chi.Router) { r.Get("/", s.Profile) r.Get("/feed.atom", s.AtomFeedPage) @@ -400,7 +402,6 @@ repo := repo.New( s.oauth, s.repoResolver, s.pages, - s.spindlestream, s.idResolver, s.db, s.config, @@ -419,8 +420,6 @@ pipes := pipelines.New( s.oauth, s.repoResolver, s.pages, - s.spindlestream, - s.pipelineNotifier, s.idResolver, s.db, s.config, diff --git a/appview/state/spindlestream.go b/appview/state/spindlestream.go deleted file mode 100644 --- a/appview/state/spindlestream.go +++ /dev/null @@ -1,182 +0,0 @@ -package state - -import ( - "context" - "encoding/json" - "fmt" - "strings" - "time" - - "github.com/bluesky-social/indigo/atproto/syntax" - "tangled.org/core/api/tangled" - "tangled.org/core/appview/config" - "tangled.org/core/appview/db" - "tangled.org/core/appview/models" - "tangled.org/core/appview/pipelines" - ec "tangled.org/core/eventconsumer" - "tangled.org/core/eventstream" - "tangled.org/core/log" - "tangled.org/core/orm" - "tangled.org/core/rbac" - spindle "tangled.org/core/spindle/models" - "tangled.org/core/workflow" -) - -func Spindlestream(ctx context.Context, c *config.Config, d *db.DB, enforcer *rbac.Enforcer, pn *pipelines.StatusNotifier) (*ec.Consumer, error) { - spindles, err := db.GetSpindles(ctx, d, orm.FilterIsNot("verified", "null")) - if err != nil { - return nil, err - } - - hosts := make([]string, len(spindles)) - for i, s := range spindles { - hosts[i] = s.Instance - } - - return bootstrapStream( - ctx, "spindlestream", ec.KindSpindle, hosts, c.Redis.Addr, - c.Spindlestream, - spindleIngester(d, pn), - ), nil -} - -func spindleIngester(d *db.DB, pn *pipelines.StatusNotifier) ec.ProcessFunc { - return func(ctx context.Context, source ec.Source, msg eventstream.Event) error { - switch msg.Nsid { - case tangled.PipelineNSID: - return ingestPipeline(ctx, d, source, msg) - case tangled.PipelineStatusNSID: - return ingestPipelineStatus(ctx, d, pn, source, msg) - } - return nil - } -} - -func ingestPipeline(ctx context.Context, d *db.DB, source ec.Source, msg eventstream.Event) error { - l := log.FromContext(ctx) - - var record tangled.Pipeline - if err := json.Unmarshal(msg.EventJson, &record); err != nil { - return fmt.Errorf("unmarshal pipeline: %w", err) - } - - if record.TriggerMetadata == nil { - return fmt.Errorf("empty trigger metadata: nsid %s, rkey %s", msg.Nsid, msg.Rkey) - } - - if record.TriggerMetadata.Repo == nil { - return fmt.Errorf("empty repo: nsid %s, rkey %s", msg.Nsid, msg.Rkey) - } - - repoName := "" - if record.TriggerMetadata.Repo.Repo != nil { - repoName = *record.TriggerMetadata.Repo.Repo - } - - repo, lookupErr := resolveRepo(d, record.TriggerMetadata.Repo.RepoDid, record.TriggerMetadata.Repo.Did, repoName) - if lookupErr != nil { - return fmt.Errorf("failed to look up repo: %w", lookupErr) - } - if repo.Spindle == "" { - return fmt.Errorf("repo does not have a spindle configured yet: nsid %s, rkey %s", msg.Nsid, msg.Rkey) - } - - // trigger info - var trigger models.Trigger - var sha string - trigger.Kind = workflow.TriggerKind(record.TriggerMetadata.Kind) - switch trigger.Kind { - case workflow.TriggerKindPush: - trigger.PushRef = &record.TriggerMetadata.Push.Ref - trigger.PushNewSha = &record.TriggerMetadata.Push.NewSha - trigger.PushOldSha = &record.TriggerMetadata.Push.OldSha - sha = *trigger.PushNewSha - case workflow.TriggerKindPullRequest: - trigger.PRSourceBranch = &record.TriggerMetadata.PullRequest.SourceBranch - trigger.PRTargetBranch = &record.TriggerMetadata.PullRequest.TargetBranch - trigger.PRSourceSha = &record.TriggerMetadata.PullRequest.SourceSha - trigger.PRAction = &record.TriggerMetadata.PullRequest.Action - sha = *trigger.PRSourceSha - } - - tx, err := d.Begin() - if err != nil { - return fmt.Errorf("failed to start txn: %w", err) - } - - triggerId, err := db.AddTrigger(tx, trigger) - if err != nil { - return fmt.Errorf("failed to add trigger entry: %w", err) - } - - // TODO: we shouldn't even use knot to identify pipelines - knot := record.TriggerMetadata.Repo.Knot - pipeline := models.Pipeline{ - Rkey: msg.Rkey, - Knot: knot, - RepoOwner: syntax.DID(record.TriggerMetadata.Repo.Did), - RepoName: repoName, - RepoDid: repo.RepoDid, - TriggerId: int(triggerId), - Sha: sha, - } - - err = db.AddPipeline(tx, pipeline) - if err != nil { - return fmt.Errorf("failed to add pipeline: %w", err) - } - - err = tx.Commit() - if err != nil { - return fmt.Errorf("failed to commit txn: %w", err) - } - - l.Info("added pipeline", "pipeline", pipeline) - - return nil -} - -func ingestPipelineStatus(ctx context.Context, d *db.DB, pn *pipelines.StatusNotifier, source ec.Source, msg eventstream.Event) error { - var record tangled.PipelineStatus - err := json.Unmarshal(msg.EventJson, &record) - if err != nil { - return err - } - - pipelineUri, err := syntax.ParseATURI(record.Pipeline) - if err != nil { - return err - } - - exitCode := 0 - if record.ExitCode != nil { - exitCode = int(*record.ExitCode) - } - - // pick the record creation time if possible, or use time.Now - created := time.Now() - if t, err := time.Parse(time.RFC3339, record.CreatedAt); err == nil && created.After(t) { - created = t - } - - status := models.PipelineStatus{ - Spindle: source.Host, - Rkey: msg.Rkey, - PipelineKnot: strings.TrimPrefix(pipelineUri.Authority().String(), "did:web:"), - PipelineRkey: pipelineUri.RecordKey().String(), - Created: created, - Workflow: record.Workflow, - Status: spindle.StatusKind(record.Status), - Error: record.Error, - ExitCode: exitCode, - } - - err = db.AddPipelineStatus(ctx, d, status) - if err != nil { - return fmt.Errorf("failed to add pipeline status: %w", err) - } - - pn.Publish(pipelineUri) - - return nil -} diff --git a/appview/state/spindlestream_test.go b/appview/state/spindlestream_test.go deleted file mode 100644 --- a/appview/state/spindlestream_test.go +++ /dev/null @@ -1,144 +0,0 @@ -package state - -import ( - "context" - "io" - "log/slog" - "net/http" - "net/http/httptest" - "path/filepath" - "strings" - "testing" - "time" - - "tangled.org/core/appview/db" - "tangled.org/core/appview/pipelines" - ec "tangled.org/core/eventconsumer" - "tangled.org/core/eventconsumer/cursor" - "tangled.org/core/eventstream" - "tangled.org/core/notifier" - spindledb "tangled.org/core/spindle/db" - spindlemodels "tangled.org/core/spindle/models" -) - -func TestColdStart_SpindleEventsRebuildPipelineStatuses(t *testing.T) { - ctx := t.Context() - - spindleDB, err := spindledb.Make(ctx, filepath.Join(t.TempDir(), "spindle.db")) - if err != nil { - t.Fatalf("spindle Make: %v", err) - } - t.Cleanup(func() { spindleDB.Close() }) - - n := notifier.New() - workflowId := spindlemodels.WorkflowId{ - PipelineId: spindlemodels.PipelineId{Knot: "knot.boltless.example", Rkey: "pipeline-rk1"}, - Name: "build", - } - for _, step := range []func() error{ - func() error { return spindleDB.StatusPending(workflowId, &n) }, - func() error { return spindleDB.StatusRunning(workflowId, &n) }, - func() error { return spindleDB.StatusSuccess(workflowId, &n) }, - } { - if err := step(); err != nil { - t.Fatalf("seed spindle event: %v", err) - } - } - - mux := http.NewServeMux() - mux.HandleFunc("/events", func(w http.ResponseWriter, r *http.Request) { - _ = eventstream.Stream(w, r, eventstream.StreamConfig{ - Backend: spindleDB, - Notifier: &n, - Logger: slog.New(slog.NewTextHandler(io.Discard, nil)), - }) - }) - srv := httptest.NewServer(mux) - t.Cleanup(srv.Close) - source := ec.Source{Kind: "test", Host: strings.TrimPrefix(srv.URL, "http://"), NoTLS: true} - - appviewDB, err := db.Make(ctx, filepath.Join(t.TempDir(), "appview.db")) - if err != nil { - t.Fatalf("appview Make: %v", err) - } - t.Cleanup(func() { appviewDB.Close() }) - - logger := slog.New(slog.NewTextHandler(io.Discard, nil)) - processFunc := spindleIngester(appviewDB, pipelines.NewStatusNotifier()) - - cfg := ec.ConsumerConfig{ - ProcessFunc: processFunc, - WorkerCount: 1, - QueueSize: 16, - ConnectionTimeout: 2 * time.Second, - CursorStore: &cursor.MemoryStore{}, - Logger: logger, - } - c := ec.NewConsumer(cfg) - - consumerCtx, cancel := context.WithCancel(ctx) - defer cancel() - c.Start(consumerCtx) - c.AddSource(consumerCtx, source) - - deadline := time.Now().Add(3 * time.Second) - for time.Now().Before(deadline) { - var n int - if err := appviewDB.QueryRow(`select count(*) from pipeline_statuses`).Scan(&n); err != nil { - t.Fatalf("count: %v", err) - } - if n >= 3 { - break - } - time.Sleep(20 * time.Millisecond) - } - - rows, err := appviewDB.Query(` - select spindle, pipeline_knot, pipeline_rkey, workflow, status - from pipeline_statuses - order by created asc - `) - if err != nil { - t.Fatalf("query: %v", err) - } - defer rows.Close() - - type rec struct { - spindle, knot, rkey, workflow, status string - } - var got []rec - for rows.Next() { - var r rec - if err := rows.Scan(&r.spindle, &r.knot, &r.rkey, &r.workflow, &r.status); err != nil { - t.Fatalf("scan: %v", err) - } - got = append(got, r) - } - - if len(got) != 3 { - t.Fatalf("pipeline_statuses rows = %d, want 3: %+v", len(got), got) - } - - wantStatuses := []string{"pending", "running", "success"} - gotStatuses := map[string]bool{} - for _, r := range got { - gotStatuses[r.status] = true - if r.spindle != source.Host { - t.Errorf("spindle = %q, want %q", r.spindle, source.Host) - } - if r.knot != workflowId.Knot { - t.Errorf("pipeline_knot = %q, want %q", r.knot, workflowId.Knot) - } - if r.rkey != workflowId.Rkey { - t.Errorf("pipeline_rkey = %q, want %q", r.rkey, workflowId.Rkey) - } - if r.workflow != workflowId.Name { - t.Errorf("workflow = %q, want %q", r.workflow, workflowId.Name) - } - } - for _, want := range wantStatuses { - if !gotStatuses[want] { - t.Errorf("missing status %q in projection", want) - } - } -} diff --git a/appview/state/state.go b/appview/state/state.go --- a/appview/state/state.go +++ b/appview/state/state.go @@ -31,7 +31,6 @@ phnotify "tangled.org/core/appview/notify/posthog" whnotify "tangled.org/core/appview/notify/webhook" "tangled.org/core/appview/oauth" "tangled.org/core/appview/pages" - "tangled.org/core/appview/pipelines" pipelinessh "tangled.org/core/appview/pipelines/ssh" "tangled.org/core/appview/reporesolver" "tangled.org/core/appview/repoverify" @@ -71,8 +70,6 @@ config *config.Config repoResolver *reporesolver.RepoResolver aclService *knotacl.Service knotstream *eventconsumer.Consumer - spindlestream *eventconsumer.Consumer - pipelineNotifier *pipelines.StatusNotifier logger *slog.Logger cfClient *cloudflare.Client codesearch *codesearch.CodeSearch @@ -220,14 +217,6 @@ return nil, fmt.Errorf("failed to start knotstream consumer: %w", err) } knotstream.Start(ctx) - pipelineNotifier := pipelines.NewStatusNotifier() - - spindlestream, err := Spindlestream(ctx, config, d, enforcer, pipelineNotifier) - if err != nil { - return nil, fmt.Errorf("failed to start spindlestream consumer: %w", err) - } - spindlestream.Start(ctx) - state := &State{ db: d, notifier: notifier, @@ -244,8 +233,6 @@ config: config, repoResolver: repoResolver, aclService: aclService, knotstream: knotstream, - spindlestream: spindlestream, - pipelineNotifier: pipelineNotifier, logger: logger, cfClient: cfClient, codesearch: &codesearch.CodeSearch{Host: config.CodeSearch.ZoektUrl}, @@ -263,7 +250,7 @@ return s.db.Close() } func (s *State) NewSSHServer() *pipelinessh.Server { - return pipelinessh.New(s.db, s.config, s.pipelineNotifier, log.SubLogger(s.logger, "pipelinessh")) + return pipelinessh.New(s.db, s.config, log.SubLogger(s.logger, "pipelinessh")) } func (s *State) SecurityTxt(w http.ResponseWriter, r *http.Request) {