diff --git a/knotserver/ingester.go b/knotserver/ingester.go --- a/knotserver/ingester.go +++ b/knotserver/ingester.go @@ -7,7 +7,6 @@ "fmt" "io" "net/http" "net/url" - "path/filepath" "strings" comatproto "github.com/bluesky-social/indigo/api/atproto" @@ -17,10 +16,8 @@ "github.com/bluesky-social/jetstream/pkg/models" securejoin "github.com/cyphar/filepath-securejoin" "tangled.org/core/api/tangled" "tangled.org/core/knotserver/db" - "tangled.org/core/knotserver/git" "tangled.org/core/log" "tangled.org/core/rbac" - "tangled.org/core/workflow" ) func (h *Knot) processPublicKey(ctx context.Context, event *models.Event) error { @@ -85,137 +82,6 @@ return nil } -func (h *Knot) processPull(ctx context.Context, event *models.Event) error { - raw := json.RawMessage(event.Commit.Record) - 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) - - if record.Target == nil { - return fmt.Errorf("ignoring pull record: target repo is nil") - } - - l = l.With("target_repo", record.Target.Repo) - l = l.With("target_branch", record.Target.Branch) - - if record.Source == nil { - return fmt.Errorf("ignoring pull record: not a branch-based pull request") - } - - if record.Source.Repo != nil { - return fmt.Errorf("ignoring pull record: fork based pull") - } - - repoAt, err := syntax.ParseATURI(record.Target.Repo) - if err != nil { - return fmt.Errorf("failed to parse ATURI: %w", err) - } - - // resolve this aturi to extract the repo record - ident, err := h.resolver.ResolveIdent(ctx, repoAt.Authority().String()) - if err != nil || ident.Handle.IsInvalidHandle() { - return fmt.Errorf("failed to resolve handle: %w", err) - } - - xrpcc := xrpc.Client{ - Host: ident.PDSEndpoint(), - } - - resp, err := comatproto.RepoGetRecord(ctx, &xrpcc, "", tangled.RepoNSID, repoAt.Authority().String(), repoAt.RecordKey().String()) - if err != nil { - return fmt.Errorf("failed to resolver repo: %w", err) - } - - repo := resp.Value.Val.(*tangled.Repo) - - if repo.Knot != h.c.Server.Hostname { - return fmt.Errorf("rejected pull record: not this knot, %s != %s", repo.Knot, h.c.Server.Hostname) - } - - didSlashRepo, err := securejoin.SecureJoin(ident.DID.String(), repo.Name) - if err != nil { - return fmt.Errorf("failed to construct relative repo path: %w", err) - } - - repoPath, err := securejoin.SecureJoin(h.c.Repo.ScanPath, didSlashRepo) - if err != nil { - return fmt.Errorf("failed to construct absolute repo path: %w", err) - } - - gr, err := git.Open(repoPath, record.Source.Sha) - if err != nil { - return fmt.Errorf("failed to open git repository: %w", err) - } - - workflowDir, err := gr.FileTree(ctx, workflow.WorkflowDir) - if err != nil { - return 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, - }) - } - - trigger := tangled.Pipeline_PullRequestTriggerData{ - Action: "create", - SourceBranch: record.Source.Branch, - SourceSha: record.Source.Sha, - TargetBranch: record.Target.Branch, - } - - compiler := workflow.Compiler{ - Trigger: tangled.Pipeline_TriggerMetadata{ - Kind: string(workflow.TriggerKindPullRequest), - PullRequest: &trigger, - Repo: &tangled.Pipeline_TriggerRepo{ - Did: ident.DID.String(), - Knot: repo.Knot, - Repo: repo.Name, - }, - }, - } - - cp := compiler.Compile(compiler.Parse(pipeline)) - eventJson, err := json.Marshal(cp) - if err != nil { - return fmt.Errorf("failed to marshal pipeline event: %w", err) - } - - // do not run empty pipelines - if cp.Workflows == nil { - return nil - } - - ev := db.Event{ - Rkey: TID(), - Nsid: tangled.PipelineNSID, - EventJson: string(eventJson), - } - - return h.db.InsertEvent(ev, h.n) -} - // duplicated from add collaborator func (h *Knot) processCollaborator(ctx context.Context, event *models.Event) error { raw := json.RawMessage(event.Commit.Record) @@ -338,8 +204,6 @@ case tangled.PublicKeyNSID: err = h.processPublicKey(ctx, event) case tangled.KnotMemberNSID: err = h.processKnotMember(ctx, event) - case tangled.RepoPullNSID: - err = h.processPull(ctx, event) case tangled.RepoCollaboratorNSID: err = h.processCollaborator(ctx, event) } diff --git a/knotserver/internal.go b/knotserver/internal.go --- a/knotserver/internal.go +++ b/knotserver/internal.go @@ -23,7 +23,6 @@ "tangled.org/core/knotserver/git" "tangled.org/core/log" "tangled.org/core/notifier" "tangled.org/core/rbac" - "tangled.org/core/workflow" ) type InternalHandle struct { @@ -188,12 +187,6 @@ if err != nil { l.Error("failed to reply with compare link", "err", err, "line", line, "did", gitUserDid, "repo", gitRelativeDir) // non-fatal } - - err = h.triggerPipeline(&resp.Messages, line, gitUserDid, repoDid, repoName, pushOptions) - if err != nil { - l.Error("failed to trigger pipeline", "err", err, "line", line, "did", gitUserDid, "repo", gitRelativeDir) - // non-fatal - } } writeJSON(w, resp) @@ -242,108 +235,6 @@ EventJson: string(eventJson), } return errors.Join(errs, h.db.InsertEvent(event, h.n)) -} - -func (h *InternalHandle) triggerPipeline( - clientMsgs *[]string, - line git.PostReceiveLine, - gitUserDid string, - repoDid string, - repoName string, - pushOptions PushOptions, -) error { - if pushOptions.skipCi { - return nil - } - - didSlashRepo, err := securejoin.SecureJoin(repoDid, repoName) - if err != nil { - return err - } - - repoPath, err := securejoin.SecureJoin(h.c.Repo.ScanPath, didSlashRepo) - if err != nil { - return err - } - - 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, - }) - } - - trigger := tangled.Pipeline_PushTriggerData{ - Ref: line.Ref, - OldSha: line.OldSha.String(), - NewSha: line.NewSha.String(), - } - - compiler := workflow.Compiler{ - Trigger: tangled.Pipeline_TriggerMetadata{ - Kind: string(workflow.TriggerKindPush), - Push: &trigger, - Repo: &tangled.Pipeline_TriggerRepo{ - Did: repoDid, - Knot: h.c.Server.Hostname, - Repo: repoName, - }, - }, - } - - 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 := db.Event{ - Rkey: TID(), - Nsid: tangled.PipelineNSID, - EventJson: string(eventJson), - } - - return h.db.InsertEvent(event, h.n) } func (h *InternalHandle) emitCompareLink( diff --git a/knotserver/server.go b/knotserver/server.go --- a/knotserver/server.go +++ b/knotserver/server.go @@ -79,7 +79,6 @@ jc, err := jetstream.NewJetstreamClient(c.Server.JetstreamEndpoint, "knotserver", []string{ tangled.PublicKeyNSID, tangled.KnotMemberNSID, - tangled.RepoPullNSID, tangled.RepoCollaboratorNSID, }, nil, log.SubLogger(logger, "jetstream"), db, true, c.Server.LogDids) if err != nil { diff --git a/spindle/server.go b/spindle/server.go --- a/spindle/server.go +++ b/spindle/server.go @@ -327,46 +327,7 @@ 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 { - return nil - tpl := tangled.Pipeline{} - err := json.Unmarshal(msg.EventJson, &tpl) - if err != nil { - fmt.Println("error unmarshalling", 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.Key() != tpl.TriggerMetadata.Repo.Knot { - return fmt.Errorf("repo knot does not match event source: %s != %s", src.Key(), tpl.TriggerMetadata.Repo.Knot) - } - - // filter by repos - _, err = s.db.GetRepoWithName( - syntax.DID(tpl.TriggerMetadata.Repo.Did), - tpl.TriggerMetadata.Repo.Repo, - ) - if err != nil { - return fmt.Errorf("failed to get repo: %w", err) - } - - pipelineId := models.PipelineId{ - Knot: src.Key(), - Rkey: msg.Rkey, - } - - err = s.processPipeline(ctx, 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)