diff --git a/knotserver/ingester.go b/knotserver/ingester.go --- a/knotserver/ingester.go +++ b/knotserver/ingester.go @@ -4,25 +4,14 @@ "context" "encoding/json" "fmt" - "io" - "net/http" - "net/url" - "path/filepath" "strings" - comatproto "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/atproto/syntax" - "github.com/bluesky-social/indigo/xrpc" jmodels "github.com/bluesky-social/jetstream/pkg/models" "tangled.org/core/api/tangled" - "tangled.org/core/appview/models" - "tangled.org/core/eventstream" "tangled.org/core/knotserver/db" - "tangled.org/core/knotserver/git" knotxrpc "tangled.org/core/knotserver/xrpc" "tangled.org/core/log" - "tangled.org/core/tid" - "tangled.org/core/workflow" ) func (h *Knot) processPublicKey(ctx context.Context, event *jmodels.Event) error { @@ -63,276 +52,6 @@ RepoName string RepoDid string DefaultBranch string // default branch -} - -func (h *Knot) validatePullRecord(ctx context.Context, record *tangled.RepoPull) (*targetRepo, error) { - if record.Target == nil { - return nil, fmt.Errorf("ignoring pull record: target repo is nil") - } - - l := log.FromContext(ctx).With("handler", "validatePullRecord") - l = l.With("target_repo", record.Target.Repo) - l = l.With("target_branch", record.Target.Branch) - - if record.Source == nil { - return nil, fmt.Errorf("ignoring pull record: not a branch-based pull request") - } - - if record.Source.Repo != nil { - return nil, fmt.Errorf("ignoring pull record: fork based pull") - } - - var repoPath, ownerDid, repoName, repoDid string - switch { - case strings.HasPrefix(record.Target.Repo, "did:"): - repoDid = record.Target.Repo - var lookupErr error - repoPath, ownerDid, repoName, lookupErr = h.db.ResolveRepoDIDOnDisk(h.c.Repo.ScanPath, repoDid) - if lookupErr != nil { - return nil, fmt.Errorf("unknown target repo DID %s: %w", repoDid, lookupErr) - } - - case strings.Contains(record.Target.Repo, "/"): - // TODO: get rid of this PDS fetch once all repos have DIDs - repoAt, parseErr := syntax.ParseATURI(record.Target.Repo) - if parseErr != nil { - return nil, fmt.Errorf("failed to parse ATURI: %w", parseErr) - } - - ident, resolveErr := h.resolver.ResolveIdent(ctx, repoAt.Authority().String()) - if resolveErr != nil || ident.Handle.IsInvalidHandle() { - return nil, fmt.Errorf("failed to resolve handle: %w", resolveErr) - } - - xrpcc := xrpc.Client{ - Host: ident.PDSEndpoint(), - } - - resp, getErr := comatproto.RepoGetRecord(ctx, &xrpcc, "", tangled.RepoNSID, repoAt.Authority().String(), repoAt.RecordKey().String()) - if getErr != nil { - return nil, fmt.Errorf("failed to resolve repo: %w", getErr) - } - - repo, ok := resp.Value.Val.(*tangled.Repo) - if !ok { - return nil, fmt.Errorf("record at %s is not a tangled.Repo", repoAt) - } - - if repo.Knot != h.c.Server.Hostname { - return nil, fmt.Errorf("rejected pull record: not this knot, %s != %s", repo.Knot, h.c.Server.Hostname) - } - - ownerDid = ident.DID.String() - repoName = repoAt.RecordKey().String() - - repoDid, didErr := h.db.GetRepoDid(ownerDid, repoName) - if didErr != nil { - return nil, fmt.Errorf("failed to resolve repo DID for %s/%s: %w", ownerDid, repoName, didErr) - } - - var lookupErr error - repoPath, _, _, lookupErr = h.db.ResolveRepoDIDOnDisk(h.c.Repo.ScanPath, repoDid) - if lookupErr != nil { - return nil, fmt.Errorf("failed to resolve repo on disk: %w", lookupErr) - } - - default: - return nil, fmt.Errorf("ignoring pull record: target repo has unrecognized format: %s", record.Target.Repo) - } - - gr, err := git.Open(repoPath, record.Source.Branch) - if err != nil { - return nil, fmt.Errorf("failed to open git repository: %w", err) - } - - defaultBranch, _ := gr.FindMainBranch() - - return &targetRepo{ - RepoPath: repoPath, - OwnerDid: ownerDid, - RepoName: repoName, - RepoDid: repoDid, - DefaultBranch: defaultBranch, - }, nil -} - -func (h *Knot) fetchLatestSubmission(ctx context.Context, did, rkey string, record *tangled.RepoPull) (*models.PullSubmission, error) { - // resolve the PR owner's identity to fetch the blob from their PDS - prOwnerIdent, err := h.resolver.ResolveIdent(ctx, did) - if err != nil || prOwnerIdent.Handle.IsInvalidHandle() { - return nil, fmt.Errorf("failed to resolve PR owner handle: %w", err) - } - - if len(record.Rounds) == 0 { - return nil, fmt.Errorf("failed to fetch latest submission, no rounds in record") - } - - roundNumber := len(record.Rounds) - 1 - round := record.Rounds[roundNumber] - - // fetch the blob from the PR owner's PDS - prOwnerPds := prOwnerIdent.PDSEndpoint() - blobUrl, err := url.Parse(fmt.Sprintf("%s/xrpc/com.atproto.sync.getBlob", prOwnerPds)) - if err != nil { - return nil, fmt.Errorf("failed to construct blob URL: %w", err) - } - q := blobUrl.Query() - q.Set("cid", round.PatchBlob.Ref.String()) - q.Set("did", did) - blobUrl.RawQuery = q.Encode() - - req, err := http.NewRequestWithContext(ctx, http.MethodGet, blobUrl.String(), nil) - if err != nil { - return nil, fmt.Errorf("failed to create blob request: %w", err) - } - req.Header.Set("Content-Type", "application/json") - - blobResp, err := http.DefaultClient.Do(req) - if err != nil { - return nil, fmt.Errorf("failed to fetch blob: %w", err) - } - defer blobResp.Body.Close() - - blob := io.ReadCloser(blobResp.Body) - latestSubmission, err := models.PullSubmissionFromRecord(did, rkey, roundNumber, round, &blob) - if err != nil { - return nil, fmt.Errorf("failed to parse submission: %w", err) - } - - return latestSubmission, nil -} - -func (h *Knot) discoverWorkflows(ctx context.Context, repoPath, sha string) (workflow.RawPipeline, error) { - gr, err := git.Open(repoPath, sha) - if err != nil { - return nil, fmt.Errorf("failed to open git repository: %w", err) - } - - workflowDir, err := gr.FileTree(ctx, workflow.WorkflowDir) - if err != nil { - return nil, fmt.Errorf("failed to open workflow directory: %w", err) - } - - var pipeline 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 { - continue - } - - pipeline = append(pipeline, workflow.RawWorkflow{ - Name: e.Name, - Contents: contents, - }) - } - - return pipeline, nil -} - -func (h *Knot) compilePipeline(ctx context.Context, targetRepo *targetRepo, sourceBranch, sourceSha, targetBranch string, rawPipeline workflow.RawPipeline) tangled.Pipeline { - l := log.FromContext(ctx) - - trigger := tangled.Pipeline_PullRequestTriggerData{ - Action: "create", - SourceBranch: sourceBranch, - SourceSha: sourceSha, - TargetBranch: targetBranch, - } - - compiler := workflow.Compiler{ - Trigger: tangled.Pipeline_TriggerMetadata{ - Kind: string(workflow.TriggerKindPullRequest), - PullRequest: &trigger, - Repo: &tangled.Pipeline_TriggerRepo{ - Knot: h.c.Server.Hostname, - RepoDid: &targetRepo.RepoDid, - Did: targetRepo.OwnerDid, - Repo: &targetRepo.RepoName, - DefaultBranch: targetRepo.DefaultBranch, - }, - }, - } - - l.Info("raw", "raw", rawPipeline) - parsed := compiler.Parse(rawPipeline) - l.Info("parsed", "parsed", parsed) - compiled := compiler.Compile(parsed) - - l.Info("compiler diagnostics", "diagnostics", compiler.Diagnostics) - - return compiled -} - -func (h *Knot) processPull(ctx context.Context, event *jmodels.Event) error { - raw := json.RawMessage(event.Commit.Record) - rkey := event.Commit.RKey - did := event.Did - - var record tangled.RepoPull - if err := json.Unmarshal(raw, &record); err != nil { - return fmt.Errorf("failed to unmarshal record: %w", err) - } - - l := log.FromContext(ctx) - l = l.With("handler", "processPull") - l = l.With("did", did) - - l.Info("validating pull record") - targetRepo, err := h.validatePullRecord(ctx, &record) - if err != nil { - l.Warn("pull record did not validate, skipping...") - return err - } - - l = l.With("target_repo", record.Target.Repo) - l = l.With("target_branch", record.Target.Branch) - - l.Info("fetching latest submission") - latestSubmission, err := h.fetchLatestSubmission(ctx, did, rkey, &record) - if err != nil { - return err - } - - sha := latestSubmission.SourceRev - if sha == "" { - return fmt.Errorf("failed to extract source SHA from pull submission") - } - l = l.With("sha", sha) - - l.Info("discovering workflows", "repo_path", targetRepo.RepoPath) - pipeline, err := h.discoverWorkflows(ctx, targetRepo.RepoPath, sha) - if err != nil { - return err - } - - l.Info("compiling pipeline", "workflow_count", len(pipeline)) - cp := h.compilePipeline(ctx, targetRepo, record.Source.Branch, sha, record.Target.Branch, pipeline) - - // do not run empty pipelines - if cp.Workflows == nil { - l.Info("skipping empty pipeline") - return nil - } - - l.Info("marshaling pipeline event") - eventJson, err := json.Marshal(cp) - if err != nil { - return fmt.Errorf("failed to marshal pipeline event: %w", err) - } - - ev := eventstream.Event{ - Rkey: tid.TID(), - Nsid: tangled.PipelineNSID, - EventJson: eventJson, - } - - l.Info("inserting pipeline event") - return h.db.InsertEvent(ev, h.n) } func (h *Knot) processRepo(ctx context.Context, event *jmodels.Event) error { @@ -407,8 +126,6 @@ err = h.processPublicKey(ctx, event) case tangled.RepoNSID: err = h.processRepo(ctx, event) - case tangled.RepoPullNSID: - err = h.processPull(ctx, event) } default: return nil diff --git a/knotserver/internal.go b/knotserver/internal.go --- a/knotserver/internal.go +++ b/knotserver/internal.go @@ -28,7 +28,6 @@ "tangled.org/core/notifier" "tangled.org/core/rbac" "tangled.org/core/tid" - "tangled.org/core/workflow" ) type InternalHandle struct { @@ -188,11 +187,6 @@ fmt.Fprint(w, diskRelative) } -type PushOptions struct { - skipCi bool - verboseCi bool -} - func (h *InternalHandle) PostReceiveHook(w http.ResponseWriter, r *http.Request) { l := h.l.With("handler", "PostReceiveHook") @@ -263,9 +257,13 @@ l.Error("failed to reply with pull request link", "err", err, "line", line, "did", gitUserDid, "repo", gitRelativeDir) } - err = h.triggerPipeline(&resp.Messages, line, gitUserDid, ownerDid, repoName, repoDid, pushOptions) - if err != nil { - l.Error("failed to trigger pipeline", "err", err, "line", line, "did", gitUserDid, "repo", gitRelativeDir) + // emit pipeline logs link + if h.c.LogsAddr != "" { + host, port, err := net.SplitHostPort(h.c.LogsAddr) + if err == nil { + resp.Messages = append(resp.Messages, "→ Browse CI logs in your terminal:") + resp.Messages = append(resp.Messages, fmt.Sprintf(" ssh -t -p %s %s %s %s", port, host, repoDid, line.NewSha)) + } } } @@ -319,133 +317,6 @@ Rkey: tid.TID(), Nsid: tangled.GitRefUpdateNSID, EventJson: eventJson, - } - - return h.db.InsertEvent(event, h.n) -} - -func (h *InternalHandle) triggerPipeline( - clientMsgs *[]string, - line git.PostReceiveLine, - gitUserDid string, - ownerDid string, - repoName string, - repoDid string, - pushOptionsRaw []string, -) error { - var pushOptions PushOptions - for _, option := range pushOptionsRaw { - if option == "skip-ci" || option == "ci-skip" { - pushOptions.skipCi = true - } - if option == "verbose-ci" || option == "ci-verbose" { - pushOptions.verboseCi = true - } - } - if pushOptions.skipCi { - return nil - } - - repoPath, _, _, resolveErr := h.db.ResolveRepoDIDOnDisk(h.c.Repo.ScanPath, repoDid) - if resolveErr != nil { - return fmt.Errorf("failed to resolve repo on disk: %w", resolveErr) - } - - gr, err := git.Open(repoPath, line.Ref) - if err != nil { - return err - } - - workflowDir, err := gr.FileTree(context.Background(), workflow.WorkflowDir) - if err != nil { - return err - } - - var pipeline 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 { - continue - } - - pipeline = append(pipeline, workflow.RawWorkflow{ - Name: e.Name, - Contents: contents, - }) - } - - defaultBranch, _ := gr.FindMainBranch() - - trigger := tangled.Pipeline_PushTriggerData{ - Ref: line.Ref, - OldSha: line.OldSha.String(), - NewSha: line.NewSha.String(), - } - - triggerRepo := &tangled.Pipeline_TriggerRepo{ - Did: ownerDid, - Knot: h.c.Server.Hostname, - Repo: &repoName, - RepoDid: &repoDid, - DefaultBranch: defaultBranch, - } - - changedFiles, err := gr.ChangedFilesBetween(line.OldSha.String(), line.NewSha.String()) - if err != nil { - return fmt.Errorf("getting changed files: %w", err) - } - - compiler := workflow.Compiler{ - Trigger: tangled.Pipeline_TriggerMetadata{ - Kind: string(workflow.TriggerKindPush), - Push: &trigger, - Repo: triggerRepo, - }, - ChangedFiles: changedFiles, - } - - cp := compiler.Compile(compiler.Parse(pipeline)) - eventJson, err := json.Marshal(cp) - if err != nil { - return err - } - - for _, e := range compiler.Diagnostics.Errors { - *clientMsgs = append(*clientMsgs, e.String()) - } - - if pushOptions.verboseCi { - if compiler.Diagnostics.IsEmpty() { - *clientMsgs = append(*clientMsgs, "success: pipeline compiled with no diagnostics") - } - - for _, w := range compiler.Diagnostics.Warnings { - *clientMsgs = append(*clientMsgs, w.String()) - } - } - - // do not run empty pipelines - if cp.Workflows == nil { - return nil - } - - event := eventstream.Event{ - Rkey: tid.TID(), - Nsid: tangled.PipelineNSID, - EventJson: eventJson, - } - - if h.c.LogsAddr != "" { - host, port, err := net.SplitHostPort(h.c.LogsAddr) - if err == nil { - *clientMsgs = append(*clientMsgs, "→ Browse CI logs in your terminal:") - *clientMsgs = append(*clientMsgs, fmt.Sprintf(" ssh -t -p %s %s %s %s", port, host, repoDid, line.NewSha)) - } } return h.db.InsertEvent(event, h.n) diff --git a/knotserver/server.go b/knotserver/server.go --- a/knotserver/server.go +++ b/knotserver/server.go @@ -96,7 +96,6 @@ jc, err := jetstream.NewJetstreamClient(c.Server.JetstreamEndpoint, "knotserver", []string{ tangled.PublicKeyNSID, tangled.RepoNSID, - tangled.RepoPullNSID, }, nil, log.SubLogger(logger, "jetstream"), db, true, c.Server.LogDids) if err != nil { logger.Error("failed to setup jetstream", "error", err) diff --git a/spindle/server.go b/spindle/server.go --- a/spindle/server.go +++ b/spindle/server.go @@ -404,42 +404,7 @@ func (s *Spindle) processKnotStream(ctx context.Context, src eventconsumer.Source, msg eventstream.Event) 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 { - return nil - tpl := tangled.Pipeline{} - err := json.Unmarshal(msg.EventJson, &tpl) - if err != nil { - s.l.Error("failed to unmarshal pipeline event", "err", err) - return err - } - - if tpl.TriggerMetadata == nil { - return fmt.Errorf("no trigger metadata found") - } - - if tpl.TriggerMetadata.Repo == nil { - return fmt.Errorf("no repo data found") - } - - if src.Host != tpl.TriggerMetadata.Repo.Knot { - return fmt.Errorf("repo knot does not match event source: %s != %s", src.Host, tpl.TriggerMetadata.Repo.Knot) - } - - repoDid, err := s.resolvePipelineRepoDid(tpl.TriggerMetadata.Repo) - if err != nil { - return err - } - - pipelineId := models.PipelineId{ - Knot: src.Host, - Rkey: msg.Rkey, - } - - err = s.processPipeline(ctx, repoDid, tpl, pipelineId) - if err != nil { - return err - } - } else if msg.Nsid == tangled.GitRefUpdateNSID { + if msg.Nsid == tangled.GitRefUpdateNSID { event := tangled.GitRefUpdate{} if err := json.Unmarshal(msg.EventJson, &event); err != nil { l.Error("error unmarshalling", "err", err) @@ -452,6 +417,10 @@ repo, err := s.db.GetRepoByDid(repoDid) if err != nil { return fmt.Errorf("unknown repoDid %s: %w", repoDid, err) + } + + if src.Host != repo.Knot { + return fmt.Errorf("repo knot does not match event source: %s != %s", src.Host, repo.Knot) } // NOTE: we are blindly trusting the knot that it will return only repos it own