From f1e7f3f39d1702bb9e050f70630d417743841cc0 Mon Sep 17 00:00:00 2001 From: Seongmin Lee Date: Thu, 25 Dec 2025 14:56:10 +0000 Subject: [PATCH] appview: listen for pipeline events from spindlestream Signed-off-by: Seongmin Lee --- appview/state/knotstream.go | 86 -------------------------------------------------------------------------------------- appview/state/spindlestream.go | 89 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ 2 file(s) changed, 89 insertion(s)(+), 86 deletion(s)(-) diff --git a/appview/state/knotstream.go b/appview/state/knotstream.go --- a/appview/state/knotstream.go +++ b/appview/state/knotstream.go @@ -18,9 +18,7 @@ "tangled.org/core/log" "tangled.org/core/orm" "tangled.org/core/rbac" - "tangled.org/core/workflow" - "github.com/bluesky-social/indigo/atproto/syntax" "github.com/go-git/go-git/v5/plumbing" "github.com/posthog/posthog-go" ) @@ -67,8 +65,6 @@ switch msg.Nsid { case tangled.GitRefUpdateNSID: return ingestRefUpdate(d, enforcer, posthog, dev, source, msg) - case tangled.PipelineNSID: - return ingestPipeline(d, source, msg) } return nil @@ -189,86 +185,4 @@ } return tx.Commit() -} - -func ingestPipeline(d *db.DB, source ec.Source, msg ec.Message) error { - var record tangled.Pipeline - err := json.Unmarshal(msg.EventJson, &record) - if err != nil { - return 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) - } - - // does this repo have a spindle configured? - repos, err := db.GetRepos( - d, - 0, - orm.FilterEq("did", record.TriggerMetadata.Repo.Did), - orm.FilterEq("name", record.TriggerMetadata.Repo.Repo), - ) - if err != nil { - return fmt.Errorf("failed to look for repo in DB: nsid %s, rkey %s, %w", msg.Nsid, msg.Rkey, err) - } - if len(repos) != 1 { - return fmt.Errorf("incorrect number of repos returned: %d (expected 1)", len(repos)) - } - if repos[0].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) - } - - pipeline := models.Pipeline{ - Rkey: msg.Rkey, - Knot: source.Key(), - RepoOwner: syntax.DID(record.TriggerMetadata.Repo.Did), - RepoName: record.TriggerMetadata.Repo.Repo, - 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) - } - - return nil } diff --git a/appview/state/spindlestream.go b/appview/state/spindlestream.go --- a/appview/state/spindlestream.go +++ b/appview/state/spindlestream.go @@ -20,6 +20,7 @@ "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) (*ec.Consumer, error) { @@ -62,12 +63,100 @@ func spindleIngester(ctx context.Context, logger *slog.Logger, d *db.DB) ec.ProcessFunc { return func(ctx context.Context, source ec.Source, msg ec.Message) error { switch msg.Nsid { + case tangled.PipelineNSID: + return ingestPipeline(logger, d, source, msg) case tangled.PipelineStatusNSID: return ingestPipelineStatus(ctx, logger, d, source, msg) } return nil } +} + +func ingestPipeline(l *slog.Logger, d *db.DB, source ec.Source, msg ec.Message) error { + var record tangled.Pipeline + err := json.Unmarshal(msg.EventJson, &record) + if err != nil { + return 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) + } + + // does this repo have a spindle configured? + repos, err := db.GetRepos( + d, + 0, + orm.FilterEq("did", record.TriggerMetadata.Repo.Did), + orm.FilterEq("name", record.TriggerMetadata.Repo.Repo), + ) + if err != nil { + return fmt.Errorf("failed to look for repo in DB: nsid %s, rkey %s, %w", msg.Nsid, msg.Rkey, err) + } + if len(repos) != 1 { + return fmt.Errorf("incorrect number of repos returned: %d (expected 1)", len(repos)) + } + if repos[0].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: record.TriggerMetadata.Repo.Repo, + 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, logger *slog.Logger, d *db.DB, source ec.Source, msg ec.Message) error { -- tangled.sh