diff --git a/knotserver/internal.go b/knotserver/internal.go --- a/knotserver/internal.go +++ b/knotserver/internal.go @@ -176,6 +176,7 @@ } for _, line := range lines { + // TODO: pass pushOptions to refUpdate err := h.insertRefUpdate(line, gitUserDid, repoDid, repoName) if err != nil { l.Error("failed to insert op", "err", err, "line", line, "did", gitUserDid, "repo", gitRelativeDir) diff --git a/spindle/server.go b/spindle/server.go --- a/spindle/server.go +++ b/spindle/server.go @@ -4,6 +4,7 @@ "context" _ "embed" "encoding/json" + "errors" "fmt" "log/slog" "maps" @@ -13,11 +14,13 @@ "github.com/bluesky-social/indigo/atproto/syntax" "github.com/go-chi/chi/v5" + "github.com/go-git/go-git/v5/plumbing/object" "github.com/hashicorp/go-version" "tangled.org/core/api/tangled" "tangled.org/core/eventconsumer" "tangled.org/core/eventconsumer/cursor" "tangled.org/core/idresolver" + kgit "tangled.org/core/knotserver/git" "tangled.org/core/log" "tangled.org/core/notifier" "tangled.org/core/rbac2" @@ -31,6 +34,8 @@ "tangled.org/core/spindle/secrets" "tangled.org/core/spindle/xrpc" "tangled.org/core/tap" + "tangled.org/core/tid" + "tangled.org/core/workflow" "tangled.org/core/xrpc/serviceauth" ) @@ -130,13 +135,13 @@ return nil, fmt.Errorf("failed to setup sqlite3 cursor store: %w", err) } - // for each incoming sh.tangled.pipeline, we execute - // spindle.processPipeline, which in turn enqueues the pipeline - // job in the above registered queue. + // spindle listen to knot stream for sh.tangled.git.refUpdate + // which will sync the local workflow files in spindle and enqueues the + // pipeline job for on-push workflows ccfg := eventconsumer.NewConsumerConfig() ccfg.Logger = log.SubLogger(logger, "eventconsumer") ccfg.Dev = cfg.Server.Dev - ccfg.ProcessFunc = spindle.processPipeline + ccfg.ProcessFunc = spindle.processKnotStream ccfg.CursorStore = cursorStore knownKnots, err := d.Knots() if err != nil { @@ -281,7 +286,7 @@ return x.Router() } -func (s *Spindle) processPipeline(ctx context.Context, src eventconsumer.Source, msg eventconsumer.Message) error { +func (s *Spindle) processKnotStream(ctx context.Context, src eventconsumer.Source, msg eventconsumer.Message) error { l := log.FromContext(ctx).With("handler", "processKnotStream") l = l.With("src", src.Key(), "msg.Nsid", msg.Nsid, "msg.Rkey", msg.Rkey) if msg.Nsid == tangled.PipelineNSID { @@ -319,72 +324,9 @@ Rkey: msg.Rkey, } - workflows := make(map[models.Engine][]models.Workflow) - - // Build pipeline environment variables once for all workflows - pipelineEnv := models.PipelineEnvVars(tpl.TriggerMetadata, pipelineId, s.cfg.Server.Dev) - - for _, w := range tpl.Workflows { - if w != nil { - if _, ok := s.engs[w.Engine]; !ok { - err = s.db.StatusFailed(models.WorkflowId{ - PipelineId: pipelineId, - Name: w.Name, - }, fmt.Sprintf("unknown engine %#v", w.Engine), -1, s.n) - if err != nil { - return fmt.Errorf("db.StatusFailed: %w", err) - } - - continue - } - - eng := s.engs[w.Engine] - - if _, ok := workflows[eng]; !ok { - workflows[eng] = []models.Workflow{} - } - - ewf, err := s.engs[w.Engine].InitWorkflow(*w, tpl) - if err != nil { - return fmt.Errorf("init workflow: %w", err) - } - - // inject TANGLED_* env vars after InitWorkflow - // This prevents user-defined env vars from overriding them - if ewf.Environment == nil { - ewf.Environment = make(map[string]string) - } - maps.Copy(ewf.Environment, pipelineEnv) - - workflows[eng] = append(workflows[eng], *ewf) - - err = s.db.StatusPending(models.WorkflowId{ - PipelineId: pipelineId, - Name: w.Name, - }, s.n) - if err != nil { - return fmt.Errorf("db.StatusPending: %w", err) - } - } - } - - ok := s.jq.Enqueue(queue.Job{ - Run: func() error { - engine.StartWorkflows(log.SubLogger(s.l, "engine"), s.vault, s.cfg, s.db, s.n, ctx, &models.Pipeline{ - RepoOwner: tpl.TriggerMetadata.Repo.Did, - RepoName: tpl.TriggerMetadata.Repo.Repo, - Workflows: workflows, - }, pipelineId) - return nil - }, - OnFail: func(jobError error) { - s.l.Error("pipeline run failed", "error", jobError) - }, - }) - if ok { - s.l.Info("pipeline enqueued successfully", "id", msg.Rkey) - } else { - s.l.Error("failed to enqueue pipeline: queue is full") + err = s.processPipeline(ctx, tpl, pipelineId) + if err != nil { + return err } } else if msg.Nsid == tangled.GitRefUpdateNSID { event := tangled.GitRefUpdate{} @@ -409,9 +351,169 @@ } l.Info("synced git repo") - // TODO: plan the pipeline + compiler := workflow.Compiler{ + Trigger: tangled.Pipeline_TriggerMetadata{ + Kind: string(workflow.TriggerKindPush), + Push: &tangled.Pipeline_PushTriggerData{ + Ref: event.Ref, + OldSha: event.OldSha, + NewSha: event.NewSha, + }, + Repo: &tangled.Pipeline_TriggerRepo{ + Did: repo.Did.String(), + Knot: repo.Knot, + Repo: repo.Name, + }, + }, + } + + // load workflow definitions from rev (without spindle context) + rawPipeline, err := s.loadPipeline(ctx, repoCloneUri, repoPath, event.NewSha) + if err != nil { + return fmt.Errorf("loading pipeline: %w", err) + } + if len(rawPipeline) == 0 { + l.Info("no workflow definition find for the repo. skipping the event") + return nil + } + tpl := compiler.Compile(compiler.Parse(rawPipeline)) + // TODO: pass compile error to workflow log + for _, w := range compiler.Diagnostics.Errors { + l.Error(w.String()) + } + for _, w := range compiler.Diagnostics.Warnings { + l.Warn(w.String()) + } + + pipelineId := models.PipelineId{ + Knot: tpl.TriggerMetadata.Repo.Knot, + Rkey: tid.TID(), + } + if err := s.db.CreatePipelineEvent(pipelineId.Rkey, tpl, s.n); err != nil { + l.Error("failed to create pipeline event", "err", err) + return nil + } + err = s.processPipeline(ctx, tpl, pipelineId) + if err != nil { + return err + } } + return nil +} + +func (s *Spindle) loadPipeline(ctx context.Context, repoUri, repoPath, rev string) (workflow.RawPipeline, error) { + if err := git.SparseSyncGitRepo(ctx, repoUri, repoPath, rev); err != nil { + return nil, fmt.Errorf("syncing git repo: %w", err) + } + gr, err := kgit.Open(repoPath, rev) + if err != nil { + return nil, fmt.Errorf("opening git repo: %w", err) + } + + workflowDir, err := gr.FileTree(ctx, workflow.WorkflowDir) + if errors.Is(err, object.ErrDirectoryNotFound) { + // return empty RawPipeline when directory doesn't exist + return nil, nil + } else if err != nil { + return nil, fmt.Errorf("loading file tree: %w", err) + } + + var rawPipeline workflow.RawPipeline + for _, e := range workflowDir { + if !e.IsFile() { + continue + } + + fpath := filepath.Join(workflow.WorkflowDir, e.Name) + contents, err := gr.RawContent(fpath) + if err != nil { + return nil, fmt.Errorf("reading raw content of '%s': %w", fpath, err) + } + + rawPipeline = append(rawPipeline, workflow.RawWorkflow{ + Name: e.Name, + Contents: contents, + }) + } + + return rawPipeline, nil +} + +func (s *Spindle) processPipeline(ctx context.Context, tpl tangled.Pipeline, pipelineId models.PipelineId) error { + // Build pipeline environment variables once for all workflows + pipelineEnv := models.PipelineEnvVars(tpl.TriggerMetadata, pipelineId, s.cfg.Server.Dev) + + // filter & init workflows + workflows := make(map[models.Engine][]models.Workflow) + for _, w := range tpl.Workflows { + if w == nil { + continue + } + if _, ok := s.engs[w.Engine]; !ok { + err := s.db.StatusFailed(models.WorkflowId{ + PipelineId: pipelineId, + Name: w.Name, + }, fmt.Sprintf("unknown engine %#v", w.Engine), -1, s.n) + if err != nil { + return fmt.Errorf("db.StatusFailed: %w", err) + } + + continue + } + + eng := s.engs[w.Engine] + + if _, ok := workflows[eng]; !ok { + workflows[eng] = []models.Workflow{} + } + + ewf, err := s.engs[w.Engine].InitWorkflow(*w, tpl) + if err != nil { + return fmt.Errorf("init workflow: %w", err) + } + + // inject TANGLED_* env vars after InitWorkflow + // This prevents user-defined env vars from overriding them + if ewf.Environment == nil { + ewf.Environment = make(map[string]string) + } + maps.Copy(ewf.Environment, pipelineEnv) + + workflows[eng] = append(workflows[eng], *ewf) + } + + // enqueue pipeline + ok := s.jq.Enqueue(queue.Job{ + Run: func() error { + engine.StartWorkflows(log.SubLogger(s.l, "engine"), s.vault, s.cfg, s.db, s.n, ctx, &models.Pipeline{ + RepoOwner: tpl.TriggerMetadata.Repo.Did, + RepoName: tpl.TriggerMetadata.Repo.Repo, + Workflows: workflows, + }, pipelineId) + return nil + }, + OnFail: func(jobError error) { + s.l.Error("pipeline run failed", "error", jobError) + }, + }) + if !ok { + return fmt.Errorf("failed to enqueue pipeline: queue is full") + } + s.l.Info("pipeline enqueued successfully", "id", pipelineId) + + // emit StatusPending for all workflows here (after successful enqueue) + for _, ewfs := range workflows { + for _, ewf := range ewfs { + err := s.db.StatusPending(models.WorkflowId{ + PipelineId: pipelineId, + Name: ewf.Name, + }, s.n) + if err != nil { + return fmt.Errorf("db.StatusPending: %w", err) + } + } + } return nil } diff --git a/spindle/tap.go b/spindle/tap.go --- a/spindle/tap.go +++ b/spindle/tap.go @@ -11,7 +11,10 @@ "tangled.org/core/eventconsumer" "tangled.org/core/spindle/db" "tangled.org/core/spindle/git" + "tangled.org/core/spindle/models" "tangled.org/core/tap" + "tangled.org/core/tid" + "tangled.org/core/workflow" ) func (s *Spindle) processEvent(ctx context.Context, evt tap.Event) error { @@ -281,11 +284,95 @@ l.Info("processing pull record") + // only listen to live events + if !evt.Record.Live { + l.Info("skipping backfill event", "event", evt.Record.AtUri()) + return nil + } + switch evt.Record.Action { case tap.RecordCreateAction, tap.RecordUpdateAction: - // TODO + record := tangled.RepoPull{} + if err := json.Unmarshal(evt.Record.Record, &record); err != nil { + l.Error("invalid record", "err", err) + return fmt.Errorf("parsing record: %w", err) + } + + // ignore legacy records + if record.Target == nil { + l.Info("ignoring pull record: target repo is nil") + return nil + } + + // ignore patch-based and fork-based PRs + if record.Source == nil || record.Source.Repo != nil { + l.Info("ignoring pull record: not a branch-based pull request") + return nil + } + + // skip if target repo is unknown + repo, err := s.db.GetRepo(syntax.ATURI(record.Target.Repo)) + if err != nil { + l.Warn("target repo is not ingested yet", "repo", record.Target.Repo, "err", err) + return fmt.Errorf("target repo is unknown") + } + + compiler := workflow.Compiler{ + Trigger: tangled.Pipeline_TriggerMetadata{ + Kind: string(workflow.TriggerKindPullRequest), + PullRequest: &tangled.Pipeline_PullRequestTriggerData{ + Action: "create", + SourceBranch: record.Source.Branch, + SourceSha: record.Source.Sha, + TargetBranch: record.Target.Branch, + }, + Repo: &tangled.Pipeline_TriggerRepo{ + Did: repo.Did.String(), + Knot: repo.Knot, + Repo: repo.Name, + }, + }, + } + + repoUri := s.newRepoCloneUrl(repo.Knot, repo.Did.String(), repo.Name) + repoPath := s.newRepoPath(repo.Did, repo.Rkey) + + // load workflow definitions from rev (without spindle context) + rawPipeline, err := s.loadPipeline(ctx, repoUri, repoPath, record.Source.Sha) + if err != nil { + // don't retry + l.Error("failed loading pipeline", "err", err) + return nil + } + if len(rawPipeline) == 0 { + l.Info("no workflow definition find for the repo. skipping the event") + return nil + } + tpl := compiler.Compile(compiler.Parse(rawPipeline)) + // TODO: pass compile error to workflow log + for _, w := range compiler.Diagnostics.Errors { + l.Error(w.String()) + } + for _, w := range compiler.Diagnostics.Warnings { + l.Warn(w.String()) + } + + pipelineId := models.PipelineId{ + Knot: tpl.TriggerMetadata.Repo.Knot, + Rkey: tid.TID(), + } + if err := s.db.CreatePipelineEvent(pipelineId.Rkey, tpl, s.n); err != nil { + l.Error("failed to create pipeline event", "err", err) + return nil + } + err = s.processPipeline(ctx, tpl, pipelineId) + if err != nil { + // don't retry + l.Error("failed processing pipeline", "err", err) + return nil + } case tap.RecordDeleteAction: - // TODO + // no-op } return nil } diff --git a/nix/modules/spindle.nix b/nix/modules/spindle.nix --- a/nix/modules/spindle.nix +++ b/nix/modules/spindle.nix @@ -136,7 +136,10 @@ "sh.tangled.repo" "sh.tangled.repo.collaborator" "sh.tangled.spindle.member" + "sh.tangled.repo.pull" ]}" + # temporary hack to listen for repo.pull from non-tangled users + "TAP_SIGNAL_COLLECTION=sh.tangled.repo.pull" ]; ExecStart = "${getExe cfg.tap-package} run"; }; diff --git a/spindle/db/events.go b/spindle/db/events.go --- a/spindle/db/events.go +++ b/spindle/db/events.go @@ -70,6 +70,20 @@ return evts, nil } +func (d *DB) CreatePipelineEvent(rkey string, pipeline tangled.Pipeline, n *notifier.Notifier) error { + eventJson, err := json.Marshal(pipeline) + if err != nil { + return err + } + event := Event{ + Rkey: rkey, + Nsid: tangled.PipelineNSID, + Created: time.Now().UnixNano(), + EventJson: string(eventJson), + } + return d.insertEvent(event, n) +} + func (d *DB) createStatusEvent( workflowId models.WorkflowId, statusKind models.StatusKind,