diff --git a/knot.go b/knot.go index 8ce1fa1..77ece3c 100644 --- a/knot.go +++ b/knot.go @@ -29,7 +29,6 @@ import ( "encoding/json" "fmt" "log/slog" - "time" "tangled.org/core/api/tangled" "tangled.org/core/eventconsumer" @@ -74,11 +73,13 @@ type KnotConsumer interface { type knotConsumer struct { c *eventconsumer.Consumer log *slog.Logger - // br is how we publish synthesized sh.tangled.pipeline.status - // records back out to /events subscribers. Today it's driven by - // the fake-job stand-in in process(); once we hook up Buildkite, - // the webhook handler will be the primary publisher. - br *broker + + // provider dispatches each incoming pipeline trigger to whatever + // backend actually runs it (today: the fake provider; tomorrow: + // Buildkite). The consumer doesn't care which — it just hands + // over the decoded record and lets the provider publish status + // records back through its own broker connection. + provider Provider } // Compile-time interface conformance check. @@ -94,7 +95,7 @@ var _ KnotConsumer = (*knotConsumer)(nil) // restart is harmless. When we start translating triggers into real // Buildkite builds, this should switch to a SQLite-backed cursor store // to avoid duplicate builds. -func startKnotConsumer(ctx context.Context, cfg config, st *store, br *broker) (*knotConsumer, error) { +func startKnotConsumer(ctx context.Context, cfg config, st *store, provider Provider) (*knotConsumer, error) { logger := loggerFrom(ctx).With("component", "knotconsumer") knots, err := st.KnotsForSpindle(ctx, cfg.Hostname) @@ -102,7 +103,7 @@ func startKnotConsumer(ctx context.Context, cfg config, st *store, br *broker) ( return nil, fmt.Errorf("load known knots: %w", err) } - kc := &knotConsumer{log: logger, br: br} + kc := &knotConsumer{log: logger, provider: provider} ccfg := eventconsumer.NewConsumerConfig() ccfg.Logger = logger @@ -200,12 +201,11 @@ func (k *knotConsumer) process(ctx context.Context, src eventconsumer.Source, ms "workflows", len(p.Workflows), ) - // Stand-in for the real Buildkite integration. Spawn one fake - // job per workflow so the /events fan-out has something to - // emit and the appview can show progress end to end. We hand - // each goroutine the worker ctx (app-scoped) so they survive - // process() returning but exit cleanly on shutdown. - k.spawnFakeJobs(ctx, src.Key(), msg.Rkey, p.Workflows) + // Hand the trigger to whichever Provider was configured. + // Spawn is non-blocking — it fans out into provider-owned + // goroutines so this worker can move on to the next event. + // The provider keeps ctx around for shutdown coordination. + k.provider.Spawn(ctx, src.Key(), msg.Rkey, p.Workflows) default: // Knots may publish other record types over the same stream; we @@ -221,109 +221,3 @@ func (k *knotConsumer) process(ctx context.Context, src eventconsumer.Source, ms return nil } - -// fakeJob constants. Pulled out so it's obvious where the timing -// numbers come from, and trivially adjustable when we want to dial the -// fake up or down. -const ( - // fakeJobDuration is the wall-clock length of a fake run. Total - // publishes per workflow = (fakeJobDuration / fakeJobInterval) + 1 - // (one final "success"). - fakeJobDuration = 30 * time.Second - // fakeJobInterval is how often we emit a "running" heartbeat. - fakeJobInterval = 5 * time.Second -) - -// spawnFakeJobs starts a goroutine per workflow. They each emit a -// stream of sh.tangled.pipeline.status records via the broker until -// either the fake duration elapses (success) or ctx is cancelled. -// -// This is a deliberate stand-in: it lets us validate the entire -// jetstream → knot → broker → /events → appview pipeline before the -// real Buildkite plumbing is in place. -func (k *knotConsumer) spawnFakeJobs(ctx context.Context, knot, pipelineRkey string, workflows []*tangled.Pipeline_Workflow) { - if len(workflows) == 0 { - // Nothing to fake — without a workflow name there's no valid - // pipeline.status record to publish. - k.log.Warn("pipeline has no workflows; skipping fake run", - "knot", knot, "rkey", pipelineRkey, - ) - return - } - for _, wf := range workflows { - if wf == nil || wf.Name == "" { - continue - } - go k.runFakeJob(ctx, knot, pipelineRkey, wf.Name) - } -} - -// runFakeJob emits a "running" status every fakeJobInterval for -// fakeJobDuration, then a final "success". It returns early if ctx is -// cancelled (shutdown) — without doing a final publish, since we'd be -// writing to a broker whose store may be closing. -func (k *knotConsumer) runFakeJob(ctx context.Context, knot, pipelineRkey, workflow string) { - // pipelineURI is what the appview parses out of the status record - // to associate it with the originating pipeline. Format mirrors - // what the upstream spindle emits: at://did:web:// - // — the appview strips the did:web: prefix and uses the hostname - // as the knot identifier. - pipelineURI := fmt.Sprintf("at://did:web:%s/%s/%s", - knot, tangled.PipelineNSID, pipelineRkey, - ) - - logger := k.log.With( - "knot", knot, - "pipeline_rkey", pipelineRkey, - "workflow", workflow, - ) - - // Heartbeat phase. seq doubles as a per-workflow disambiguator in - // the synthesized status rkey so multiple fakes don't collide. - deadline := time.Now().Add(fakeJobDuration) - seq := 0 - for time.Now().Before(deadline) { - if err := k.publishStatus(ctx, pipelineURI, workflow, "running", seq); err != nil { - logger.Error("publish fake running status", "err", err, "seq", seq) - return - } - seq++ - select { - case <-ctx.Done(): - logger.Debug("fake job cancelled mid-run", "seq", seq) - return - case <-time.After(fakeJobInterval): - } - } - - // Terminal status. Marked as "success" using the upstream - // StatusKind enum's success label (see tangled.org/core/spindle/models). - if err := k.publishStatus(ctx, pipelineURI, workflow, "success", seq); err != nil { - logger.Error("publish fake success status", "err", err, "seq", seq) - return - } - logger.Info("fake job complete") -} - -// publishStatus assembles a tangled.PipelineStatus, marshals it, and -// hands it to the broker for persistence + fan-out. The rkey we mint -// is purely synthetic — it just needs to be unique across our event -// log; the appview keys its rows on (spindle, rkey). -func (k *knotConsumer) publishStatus(ctx context.Context, pipelineURI, workflow, status string, seq int) error { - rec := tangled.PipelineStatus{ - LexiconTypeID: tangled.PipelineStatusNSID, - Pipeline: pipelineURI, - Workflow: workflow, - Status: status, - CreatedAt: time.Now().UTC().Format(time.RFC3339), - } - body, err := json.Marshal(rec) - if err != nil { - return fmt.Errorf("marshal pipeline.status: %w", err) - } - rkey := fmt.Sprintf("fake-%d-%s-%d", time.Now().UnixNano(), workflow, seq) - if _, err := k.br.Publish(ctx, rkey, tangled.PipelineStatusNSID, body); err != nil { - return fmt.Errorf("publish pipeline.status: %w", err) - } - return nil -} diff --git a/main.go b/main.go index 3b8e42d..7a1b007 100644 --- a/main.go +++ b/main.go @@ -114,11 +114,19 @@ func main() { // them to publish synthetic status events at startup. br := newBroker(st) + // Provider that turns Tangled pipeline triggers into + // pipeline.status events. The fake provider stands in for a real + // CI integration: it emits synthetic running/success heartbeats + // over the broker so the entire jetstream → knot → /events flow + // is exercisable end-to-end. Swap this for a Buildkite-backed + // implementation once that lands. + provider := newFakeProvider(br, logger) + // Start the knot event-stream consumer first so the jetstream // loop has somewhere to register newly-observed knots into. It - // gets the broker so its (currently fake) pipeline runner can - // publish sh.tangled.pipeline.status events back out via /events. - knots, err := startKnotConsumer(ctx, cfg, st, br) + // gets the provider so each incoming pipeline trigger has + // something to dispatch to. + knots, err := startKnotConsumer(ctx, cfg, st, provider) if err != nil { logger.Error("failed to start knot consumer", "err", err) os.Exit(1) diff --git a/provider.go b/provider.go new file mode 100644 index 0000000..e807438 --- /dev/null +++ b/provider.go @@ -0,0 +1,46 @@ +package main + +// Provider is the abstraction over "the thing that turns a Tangled +// pipeline trigger into pipeline.status events". It exists so the rest +// of tack can stay agnostic to whether a given trigger is dispatched to +// Buildkite, run by a stub for testing, or anything else we plug in later. + +import ( + "context" + + "tangled.org/core/api/tangled" +) + +// Provider dispatches a Tangled pipeline trigger to whatever backend +// actually runs the workflows. +// +// Implementations are responsible for publishing +// sh.tangled.pipeline.status records back through whatever channel +// they were constructed with. +type Provider interface { + // Spawn kicks off a pipeline run for every workflow in workflows. + // + // It MUST be non-blocking: the caller is the eventconsumer worker + // that's shared across all knot subscriptions, so per-pipeline + // work has to live on its own goroutine. A typical implementation + // fans out into a goroutine per workflow and returns immediately. + // + // ctx is the consumer's app-scoped context (lives until shutdown, + // not just for the duration of one event). Implementations are + // expected to honour cancellation: in-flight runs should wind + // down without issuing further publishes once ctx is done. + // + // knot is the knot hostname the trigger arrived on; it's the + // authority half of the pipeline ATURI that pipeline.status + // records reference. pipelineRkey is the trigger record's rkey + // on that knot. workflows is the unmodified slice from the + // decoded sh.tangled.pipeline record; implementations should + // tolerate nil entries and zero-length names defensively, since + // the lexicon doesn't enforce either. + Spawn( + ctx context.Context, + knot string, + pipelineRkey string, + workflows []*tangled.Pipeline_Workflow, + ) +} diff --git a/provider_fake.go b/provider_fake.go new file mode 100644 index 0000000..0748a9d --- /dev/null +++ b/provider_fake.go @@ -0,0 +1,157 @@ +package main + +// fakeProvider is a stand-in Provider implementation: it doesn't talk +// to any external CI. For each workflow in a triggered pipeline it +// spawns a goroutine that emits a fixed-cadence stream of +// sh.tangled.pipeline.status records — "running" every five seconds +// for thirty seconds, then a final "success" — through the broker. +// +// The point is to exercise the entire trigger → broker → /events → +// appview path end-to-end before any real CI integration exists. Once +// the Buildkite provider lands, this one stays around as a reference +// implementation and as the test double of choice when a test wants +// "something that publishes plausible status updates" without the +// timing weight of real builds. + +import ( + "context" + "encoding/json" + "fmt" + "log/slog" + "time" + + "tangled.org/core/api/tangled" +) + +// Fake-job timing knobs. Pulled out as constants so it's obvious +// where the numbers come from and they can be tuned independently of +// the rest of the file. Total publishes per workflow = +// (fakeJobDuration / fakeJobInterval) heartbeats + 1 final success. +const ( + fakeJobDuration = 30 * time.Second + fakeJobInterval = 5 * time.Second +) + +// fakeProvider implements Provider against the in-process broker. +type fakeProvider struct { + br *broker + log *slog.Logger +} + +// Compile-time interface check — keeps the fake honest if Provider +// ever gains additional methods. +var _ Provider = (*fakeProvider)(nil) + +// newFakeProvider constructs a fakeProvider bound to br. The provided +// logger is annotated with component=provider so its output stands +// apart from the knot-consumer / jetstream noise. +func newFakeProvider(br *broker, log *slog.Logger) *fakeProvider { + return &fakeProvider{ + br: br, + log: log.With("component", "provider", "kind", "fake"), + } +} + +// Spawn satisfies Provider. It kicks off one runWorkflow goroutine per +// workflow, returning immediately so the eventconsumer worker that +// invoked us isn't blocked. Goroutines inherit ctx (app-scoped) and +// will exit early on cancellation. +func (p *fakeProvider) Spawn( + ctx context.Context, + knot string, + pipelineRkey string, + workflows []*tangled.Pipeline_Workflow, +) { + if len(workflows) == 0 { + // Without a workflow name there's no valid pipeline.status + // record to publish. Log loudly enough that an operator + // staring at the logs can tell the trigger arrived but + // produced no fake activity. + p.log.Warn("pipeline has no workflows; skipping fake run", + "knot", knot, "rkey", pipelineRkey, + ) + return + } + for _, wf := range workflows { + // Defensive: the lexicon allows pointer entries and doesn't + // enforce non-empty names. We can't publish a status for an + // unnamed workflow, so just skip it. + if wf == nil || wf.Name == "" { + continue + } + go p.runWorkflow(ctx, knot, pipelineRkey, wf.Name) + } +} + +// runWorkflow emits a "running" status every fakeJobInterval until +// fakeJobDuration elapses, then a final "success". On ctx +// cancellation it returns without issuing the terminal publish — the +// broker's underlying store may already be closing during shutdown. +func (p *fakeProvider) runWorkflow(ctx context.Context, knot, pipelineRkey, workflow string) { + // pipelineURI is what the appview parses out of the status record + // to associate it back with the originating pipeline. Format + // mirrors the upstream spindle's emission: + // at://did:web://. The appview strips the + // did:web: prefix and treats the remainder as the knot identifier. + pipelineURI := fmt.Sprintf("at://did:web:%s/%s/%s", + knot, tangled.PipelineNSID, pipelineRkey, + ) + + logger := p.log.With( + "knot", knot, + "pipeline_rkey", pipelineRkey, + "workflow", workflow, + ) + + // Heartbeat phase. seq doubles as a per-workflow disambiguator + // in the synthesized status rkey so concurrent fakes (across + // workflows or pipelines) don't collide. + deadline := time.Now().Add(fakeJobDuration) + seq := 0 + for time.Now().Before(deadline) { + if err := p.publishStatus(ctx, pipelineURI, workflow, "running", seq); err != nil { + logger.Error("publish fake running status", "err", err, "seq", seq) + return + } + seq++ + select { + case <-ctx.Done(): + logger.Debug("fake job cancelled mid-run", "seq", seq) + return + case <-time.After(fakeJobInterval): + } + } + + // Terminal publish. "success" matches the upstream StatusKind + // enum (see tangled.org/core/spindle/models) — the appview + // routes status strings through that same enum. + if err := p.publishStatus(ctx, pipelineURI, workflow, "success", seq); err != nil { + logger.Error("publish fake success status", "err", err, "seq", seq) + return + } + logger.Info("fake job complete") +} + +// publishStatus assembles a tangled.PipelineStatus record, marshals +// it, and pushes it through the broker. The synthesized rkey just +// needs to be unique within our event log; the appview keys its rows +// on (spindle, rkey) so we mix in time + workflow + sequence to avoid +// collisions across concurrent workflows on the same pipeline. +func (p *fakeProvider) publishStatus(ctx context.Context, pipelineURI, workflow, status string, seq int) error { + rec := tangled.PipelineStatus{ + LexiconTypeID: tangled.PipelineStatusNSID, + Pipeline: pipelineURI, + Workflow: workflow, + Status: status, + CreatedAt: time.Now().UTC().Format(time.RFC3339), + } + body, err := json.Marshal(rec) + if err != nil { + return fmt.Errorf("marshal pipeline.status: %w", err) + } + rkey := fmt.Sprintf("fake-%d-%s-%d", time.Now().UnixNano(), workflow, seq) + if _, err := p.br.Publish(ctx, rkey, tangled.PipelineStatusNSID, body); err != nil { + return fmt.Errorf("publish pipeline.status: %w", err) + } + return nil +}