From c3b287b10e3551efbb0a5072bcb635e53d3cd315 Mon Sep 17 00:00:00 2001 From: Anirudh Oppiliappan Date: Thu, 12 Jun 2025 21:28:03 +0300 Subject: [PATCH] spindle/queue: enqueue pipeline jobs Signed-off-by: Anirudh Oppiliappan --- spindle/queue/queue.go | 37 ++++++++++++++++++++++++++++++++ spindle/server.go | 48 ++++++++++++++++++++++++++++++------------ 2 files changed, 71 insertions(+), 14 deletions(-) create mode 100644 spindle/queue/queue.go diff --git a/spindle/queue/queue.go b/spindle/queue/queue.go new file mode 100644 index 00000000..c1234fa4 --- /dev/null +++ b/spindle/queue/queue.go @@ -0,0 +1,37 @@ +package queue + +type Job struct { + Run func() error + OnFail func(error) +} + +type Queue struct { + jobs chan Job +} + +func NewQueue(size int) *Queue { + return &Queue{ + jobs: make(chan Job, size), + } +} + +func (q *Queue) Enqueue(job Job) bool { + select { + case q.jobs <- job: + return true + default: + return false + } +} + +func (q *Queue) StartRunner() { + go func() { + for job := range q.jobs { + if err := job.Run(); err != nil { + if job.OnFail != nil { + job.OnFail(err) + } + } + } + }() +} diff --git a/spindle/server.go b/spindle/server.go index 81158e00..42d10042 100644 --- a/spindle/server.go +++ b/spindle/server.go @@ -1,13 +1,13 @@ package spindle import ( + "context" "encoding/json" "fmt" "log/slog" "net/http" "github.com/go-chi/chi/v5" - "golang.org/x/net/context" "tangled.sh/tangled.sh/core/api/tangled" "tangled.sh/tangled.sh/core/jetstream" "tangled.sh/tangled.sh/core/knotclient" @@ -17,6 +17,7 @@ import ( "tangled.sh/tangled.sh/core/spindle/config" "tangled.sh/tangled.sh/core/spindle/db" "tangled.sh/tangled.sh/core/spindle/engine" + "tangled.sh/tangled.sh/core/spindle/queue" ) type Spindle struct { @@ -26,6 +27,7 @@ type Spindle struct { l *slog.Logger n *notifier.Notifier eng *engine.Engine + jq *queue.Queue } func Run(ctx context.Context) error { @@ -58,6 +60,11 @@ func Run(ctx context.Context) error { return err } + jq := queue.NewQueue(100) + + // starts a job queue runner in the background + jq.StartRunner() + spindle := Spindle{ jc: jc, e: e, @@ -65,8 +72,12 @@ func Run(ctx context.Context) error { l: logger, n: &n, eng: eng, + jq: jq, } + // for each incoming sh.tangled.pipeline, we execute + // spindle.processPipeline, which in turn enqueues the pipeline + // job in the above registered queue. go func() { logger.Info("starting event consumer") knotEventSource := knotclient.NewEventSource("localhost:5555") @@ -74,7 +85,7 @@ func Run(ctx context.Context) error { ccfg := knotclient.NewConsumerConfig() ccfg.Logger = logger ccfg.Dev = cfg.Server.Dev - ccfg.ProcessFunc = spindle.exec + ccfg.ProcessFunc = spindle.processPipeline ccfg.AddEventSource(knotEventSource) ec := knotclient.NewEventConsumer(*ccfg) @@ -96,7 +107,7 @@ func (s *Spindle) Router() http.Handler { return mux } -func (s *Spindle) exec(ctx context.Context, src knotclient.EventSource, msg knotclient.Message) error { +func (s *Spindle) processPipeline(ctx context.Context, src knotclient.EventSource, msg knotclient.Message) error { if msg.Nsid == tangled.PipelineNSID { pipeline := tangled.Pipeline{} err := json.Unmarshal(msg.EventJson, &pipeline) @@ -105,17 +116,26 @@ func (s *Spindle) exec(ctx context.Context, src knotclient.EventSource, msg knot return err } - // this is a "fake" at uri for now - pipelineAtUri := fmt.Sprintf("at://%s/did:web:%s/%s", tangled.PipelineNSID, pipeline.TriggerMetadata.Repo.Knot, msg.Rkey) - - rkey := TID() - err = s.eng.SetupPipeline(ctx, &pipeline, pipelineAtUri, rkey) - if err != nil { - return err - } - err = s.eng.StartWorkflows(ctx, &pipeline, rkey) - if err != nil { - return err + ok := s.jq.Enqueue(queue.Job{ + Run: func() error { + // this is a "fake" at uri for now + pipelineAtUri := fmt.Sprintf("at://%s/did:web:%s/%s", tangled.PipelineNSID, pipeline.TriggerMetadata.Repo.Knot, msg.Rkey) + + rkey := TID() + err = s.eng.SetupPipeline(ctx, &pipeline, pipelineAtUri, rkey) + if err != nil { + return err + } + return s.eng.StartWorkflows(ctx, &pipeline, rkey) + }, + OnFail: func(error) { + s.l.Error("pipeline setup failed", "error", err) + }, + }) + if ok { + s.l.Info("pipeline enqueued successfully", "id", msg.Rkey) + } else { + s.l.Error("failed to enqueue pipeline: queue is full") } } -- 2.51.2