From 4053d0e305a3bcd1ae65663b02aa703e28dd94d3 Mon Sep 17 00:00:00 2001 From: Lewis Date: Thu, 21 May 2026 10:09:34 +0300 Subject: [PATCH] appview/state: switch streams to eventstream, backfill legacy cursors Lewis: May this revision serve well! --- appview/knots/knots.go | 12 +++ appview/repo/repo.go | 11 ++- appview/state/knotstream.go | 58 ++++------- appview/state/spindlestream.go | 54 +++-------- appview/state/spindlestream_test.go | 145 ++++++++++++++++++++++++++++ appview/state/streams.go | 56 +++++++++++ appview/state/streams_test.go | 56 +++++++++++ 7 files changed, 313 insertions(+), 79 deletions(-) create mode 100644 appview/state/spindlestream_test.go create mode 100644 appview/state/streams.go create mode 100644 appview/state/streams_test.go diff --git a/appview/knots/knots.go b/appview/knots/knots.go index c421db00..4cc3de49 100644 --- a/appview/knots/knots.go +++ b/appview/knots/knots.go @@ -337,6 +337,18 @@ func (k *Knots) delete(w http.ResponseWriter, r *http.Request) { return } + if registration.Registered != nil { + remaining, rErr := db.GetRegistrations(k.Db, + orm.FilterEq("domain", domain), + orm.FilterIsNot("registered", "null"), + ) + if rErr != nil { + l.Warn("failed to check remaining registrations after delete", "err", rErr) + } else if len(remaining) == 0 { + go k.Knotstream.RemoveSource(eventconsumer.NewKnotSource(domain)) + } + } + shouldRedirect := r.Header.Get("shouldRedirect") if shouldRedirect == "true" { k.Pages.HxRedirect(w, "/knots") diff --git a/appview/repo/repo.go b/appview/repo/repo.go index 6009f288..85f74d10 100644 --- a/appview/repo/repo.go +++ b/appview/repo/repo.go @@ -168,8 +168,17 @@ func (rp *Repo) EditSpindle(w http.ResponseWriter, r *http.Request) { return } + oldSpindle := f.Spindle + if oldSpindle != "" && oldSpindle != newSpindle { + remaining, qErr := db.GetRepos(rp.db, orm.FilterEq("spindle", oldSpindle)) + if qErr != nil { + l.Warn("failed to count repos using old spindle", "err", qErr) + } else if len(remaining) == 0 { + rp.spindlestream.RemoveSource(eventconsumer.NewSpindleSource(oldSpindle)) + } + } + if !removingSpindle { - // add this spindle to spindle stream rp.spindlestream.AddSource( context.Background(), eventconsumer.NewSpindleSource(newSpindle), diff --git a/appview/state/knotstream.go b/appview/state/knotstream.go index 3df5985b..aaa994b5 100644 --- a/appview/state/knotstream.go +++ b/appview/state/knotstream.go @@ -14,13 +14,12 @@ import ( "tangled.org/core/appview/notify" "tangled.org/core/api/tangled" - "tangled.org/core/appview/cache" "tangled.org/core/appview/config" "tangled.org/core/appview/db" "tangled.org/core/appview/models" "tangled.org/core/appview/sites" ec "tangled.org/core/eventconsumer" - "tangled.org/core/eventconsumer/cursor" + "tangled.org/core/eventstream" knotdb "tangled.org/core/knotserver/db" "tangled.org/core/log" "tangled.org/core/orm" @@ -33,40 +32,21 @@ import ( ) func Knotstream(ctx context.Context, c *config.Config, d *db.DB, enforcer *rbac.Enforcer, posthog posthog.Client, notifier notify.Notifier, cfClient *cloudflare.Client) (*ec.Consumer, error) { - logger := log.FromContext(ctx) - logger = log.SubLogger(logger, "knotstream") - - knots, err := db.GetRegistrations( - d, - orm.FilterIsNot("registered", "null"), - ) + knots, err := db.GetRegistrations(d, orm.FilterIsNot("registered", "null")) if err != nil { return nil, err } - srcs := make(map[ec.Source]struct{}) - for _, k := range knots { - s := ec.NewKnotSource(k.Domain) - srcs[s] = struct{}{} - } - - cache := cache.New(c.Redis.Addr) - cursorStore := cursor.NewRedisCursorStore(cache) - - cfg := ec.ConsumerConfig{ - Sources: srcs, - ProcessFunc: knotIngester(d, enforcer, posthog, notifier, c.Core.Dev, c, cfClient), - RetryInterval: c.Knotstream.RetryInterval, - MaxRetryInterval: c.Knotstream.MaxRetryInterval, - ConnectionTimeout: c.Knotstream.ConnectionTimeout, - WorkerCount: c.Knotstream.WorkerCount, - QueueSize: c.Knotstream.QueueSize, - Logger: logger, - Dev: c.Core.Dev, - CursorStore: &cursorStore, + hosts := make([]string, len(knots)) + for i, k := range knots { + hosts[i] = k.Domain } - return ec.NewConsumer(cfg), nil + return bootstrapStream( + ctx, "knotstream", ec.KindKnot, hosts, c.Redis.Addr, + c.Knotstream, c.Core.Dev, + knotIngester(d, enforcer, posthog, notifier, c.Core.Dev, c, cfClient), + ), nil } func resolveRepo(d *db.DB, repoDid *string, ownerDid, repoName string) (*models.Repo, error) { @@ -84,7 +64,7 @@ func resolveRepo(d *db.DB, repoDid *string, ownerDid, repoName string) (*models. } func knotIngester(d *db.DB, enforcer *rbac.Enforcer, posthog posthog.Client, notifier notify.Notifier, dev bool, c *config.Config, cfClient *cloudflare.Client) ec.ProcessFunc { - return func(ctx context.Context, source ec.Source, msg ec.Message) error { + return func(ctx context.Context, source ec.Source, msg eventstream.Event) error { switch msg.Nsid { case tangled.GitRefUpdateNSID: return ingestRefUpdate(ctx, d, enforcer, posthog, notifier, dev, c, cfClient, source, msg) @@ -99,7 +79,7 @@ func knotIngester(d *db.DB, enforcer *rbac.Enforcer, posthog posthog.Client, not } // TODO(boltless): remove this. knotmirror should do all sort of indexing -func ingestRefUpdate(ctx context.Context, d *db.DB, enforcer *rbac.Enforcer, pc posthog.Client, notifier notify.Notifier, dev bool, c *config.Config, cfClient *cloudflare.Client, source ec.Source, msg ec.Message) error { +func ingestRefUpdate(ctx context.Context, d *db.DB, enforcer *rbac.Enforcer, pc posthog.Client, notifier notify.Notifier, dev bool, c *config.Config, cfClient *cloudflare.Client, source ec.Source, msg eventstream.Event) error { logger := log.FromContext(ctx) var record tangled.GitRefUpdate @@ -112,12 +92,12 @@ func ingestRefUpdate(ctx context.Context, d *db.DB, enforcer *rbac.Enforcer, pc if err != nil { return err } - if !slices.Contains(knownKnots, source.Key()) { - return fmt.Errorf("%s does not belong to %s, something is fishy", record.CommitterDid, source.Key()) + if !slices.Contains(knownKnots, source.Host) { + return fmt.Errorf("%s does not belong to %s, something is fishy", record.CommitterDid, source.Host) } if record.Repo == "" { - return fmt.Errorf("gitRefUpdate from %s missing repo", source.Key()) + return fmt.Errorf("gitRefUpdate from %s missing repo", source.Host) } repo, lookupErr := db.GetRepoByDid(d, record.Repo) @@ -285,7 +265,7 @@ func updateRepoLanguages(d *db.DB, record tangled.GitRefUpdate) error { return tx.Commit() } -func ingestPipeline(d *db.DB, source ec.Source, msg ec.Message) error { +func ingestPipeline(d *db.DB, source ec.Source, msg eventstream.Event) error { var record tangled.Pipeline err := json.Unmarshal(msg.EventJson, &record) if err != nil { @@ -343,7 +323,7 @@ func ingestPipeline(d *db.DB, source ec.Source, msg ec.Message) error { pipeline := models.Pipeline{ Rkey: msg.Rkey, - Knot: source.Key(), + Knot: source.Host, RepoOwner: syntax.DID(record.TriggerMetadata.Repo.Did), RepoName: repoName, RepoDid: repo.RepoDid, @@ -364,7 +344,7 @@ func ingestPipeline(d *db.DB, source ec.Source, msg ec.Message) error { return nil } -func ingestDIDAssign(d *db.DB, enforcer *rbac.Enforcer, source ec.Source, msg ec.Message, ctx context.Context) error { +func ingestDIDAssign(d *db.DB, enforcer *rbac.Enforcer, source ec.Source, msg eventstream.Event, ctx context.Context) error { logger := log.FromContext(ctx) var record knotdb.RepoDIDAssign @@ -393,7 +373,7 @@ func ingestDIDAssign(d *db.DB, enforcer *rbac.Enforcer, source ec.Source, msg ec return nil } repo := repos[0] - knot := source.Key() + knot := source.Host if repo.Knot != knot { return fmt.Errorf("didAssign from %s for repo hosted on %s, rejecting", knot, repo.Knot) diff --git a/appview/state/spindlestream.go b/appview/state/spindlestream.go index bdd098e1..ebe2643a 100644 --- a/appview/state/spindlestream.go +++ b/appview/state/spindlestream.go @@ -4,75 +4,51 @@ import ( "context" "encoding/json" "fmt" - "log/slog" "strings" "time" "github.com/bluesky-social/indigo/atproto/syntax" "tangled.org/core/api/tangled" - "tangled.org/core/appview/cache" "tangled.org/core/appview/config" "tangled.org/core/appview/db" "tangled.org/core/appview/models" "tangled.org/core/appview/pipelines" ec "tangled.org/core/eventconsumer" - "tangled.org/core/eventconsumer/cursor" - "tangled.org/core/log" + "tangled.org/core/eventstream" "tangled.org/core/orm" "tangled.org/core/rbac" spindle "tangled.org/core/spindle/models" ) func Spindlestream(ctx context.Context, c *config.Config, d *db.DB, enforcer *rbac.Enforcer, pn *pipelines.StatusNotifier) (*ec.Consumer, error) { - logger := log.FromContext(ctx) - logger = log.SubLogger(logger, "spindlestream") - - spindles, err := db.GetSpindles( - ctx, - d, - orm.FilterIsNot("verified", "null"), - ) + spindles, err := db.GetSpindles(ctx, d, orm.FilterIsNot("verified", "null")) if err != nil { return nil, err } - srcs := make(map[ec.Source]struct{}) - for _, s := range spindles { - src := ec.NewSpindleSource(s.Instance) - srcs[src] = struct{}{} + hosts := make([]string, len(spindles)) + for i, s := range spindles { + hosts[i] = s.Instance } - cache := cache.New(c.Redis.Addr) - cursorStore := cursor.NewRedisCursorStore(cache) - - cfg := ec.ConsumerConfig{ - Sources: srcs, - ProcessFunc: spindleIngester(ctx, logger, d, pn), - RetryInterval: c.Spindlestream.RetryInterval, - MaxRetryInterval: c.Spindlestream.MaxRetryInterval, - ConnectionTimeout: c.Spindlestream.ConnectionTimeout, - WorkerCount: c.Spindlestream.WorkerCount, - QueueSize: c.Spindlestream.QueueSize, - Logger: logger, - Dev: c.Core.Dev, - CursorStore: &cursorStore, - } - - return ec.NewConsumer(cfg), nil + return bootstrapStream( + ctx, "spindlestream", ec.KindSpindle, hosts, c.Redis.Addr, + c.Spindlestream, c.Core.Dev, + spindleIngester(d, pn), + ), nil } -func spindleIngester(ctx context.Context, logger *slog.Logger, d *db.DB, pn *pipelines.StatusNotifier) ec.ProcessFunc { - return func(ctx context.Context, source ec.Source, msg ec.Message) error { +func spindleIngester(d *db.DB, pn *pipelines.StatusNotifier) ec.ProcessFunc { + return func(ctx context.Context, source ec.Source, msg eventstream.Event) error { switch msg.Nsid { case tangled.PipelineStatusNSID: - return ingestPipelineStatus(ctx, logger, d, pn, source, msg) + return ingestPipelineStatus(ctx, d, pn, source, msg) } - return nil } } -func ingestPipelineStatus(ctx context.Context, logger *slog.Logger, d *db.DB, pn *pipelines.StatusNotifier, source ec.Source, msg ec.Message) error { +func ingestPipelineStatus(ctx context.Context, d *db.DB, pn *pipelines.StatusNotifier, source ec.Source, msg eventstream.Event) error { var record tangled.PipelineStatus err := json.Unmarshal(msg.EventJson, &record) if err != nil { @@ -96,7 +72,7 @@ func ingestPipelineStatus(ctx context.Context, logger *slog.Logger, d *db.DB, pn } status := models.PipelineStatus{ - Spindle: source.Key(), + Spindle: source.Host, Rkey: msg.Rkey, PipelineKnot: strings.TrimPrefix(pipelineUri.Authority().String(), "did:web:"), PipelineRkey: pipelineUri.RecordKey().String(), diff --git a/appview/state/spindlestream_test.go b/appview/state/spindlestream_test.go new file mode 100644 index 00000000..05cccf57 --- /dev/null +++ b/appview/state/spindlestream_test.go @@ -0,0 +1,145 @@ +package state + +import ( + "context" + "io" + "log/slog" + "net/http" + "net/http/httptest" + "path/filepath" + "strings" + "testing" + "time" + + "tangled.org/core/appview/db" + "tangled.org/core/appview/pipelines" + ec "tangled.org/core/eventconsumer" + "tangled.org/core/eventconsumer/cursor" + "tangled.org/core/eventstream" + "tangled.org/core/notifier" + spindledb "tangled.org/core/spindle/db" + spindlemodels "tangled.org/core/spindle/models" +) + +func TestColdStart_SpindleEventsRebuildPipelineStatuses(t *testing.T) { + ctx := t.Context() + + spindleDB, err := spindledb.Make(ctx, filepath.Join(t.TempDir(), "spindle.db")) + if err != nil { + t.Fatalf("spindle Make: %v", err) + } + t.Cleanup(func() { spindleDB.Close() }) + + n := notifier.New() + workflowId := spindlemodels.WorkflowId{ + PipelineId: spindlemodels.PipelineId{Knot: "knot.boltless.example", Rkey: "pipeline-rk1"}, + Name: "build", + } + for _, step := range []func() error{ + func() error { return spindleDB.StatusPending(workflowId, &n) }, + func() error { return spindleDB.StatusRunning(workflowId, &n) }, + func() error { return spindleDB.StatusSuccess(workflowId, &n) }, + } { + if err := step(); err != nil { + t.Fatalf("seed spindle event: %v", err) + } + } + + mux := http.NewServeMux() + mux.HandleFunc("/events", func(w http.ResponseWriter, r *http.Request) { + _ = eventstream.Stream(w, r, eventstream.StreamConfig{ + Backend: spindleDB, + Notifier: &n, + Logger: slog.New(slog.NewTextHandler(io.Discard, nil)), + }) + }) + srv := httptest.NewServer(mux) + t.Cleanup(srv.Close) + source := ec.Source{Kind: "test", Host: strings.TrimPrefix(srv.URL, "http://")} + + appviewDB, err := db.Make(ctx, filepath.Join(t.TempDir(), "appview.db")) + if err != nil { + t.Fatalf("appview Make: %v", err) + } + t.Cleanup(func() { appviewDB.Close() }) + + logger := slog.New(slog.NewTextHandler(io.Discard, nil)) + processFunc := spindleIngester(appviewDB, pipelines.NewStatusNotifier()) + + cfg := ec.ConsumerConfig{ + ProcessFunc: processFunc, + WorkerCount: 1, + QueueSize: 16, + ConnectionTimeout: 2 * time.Second, + CursorStore: &cursor.MemoryStore{}, + URLFunc: ec.DefaultURL(true), + Logger: logger, + } + c := ec.NewConsumer(cfg) + + consumerCtx, cancel := context.WithCancel(ctx) + defer cancel() + c.Start(consumerCtx) + c.AddSource(consumerCtx, source) + + deadline := time.Now().Add(3 * time.Second) + for time.Now().Before(deadline) { + var n int + if err := appviewDB.QueryRow(`select count(*) from pipeline_statuses`).Scan(&n); err != nil { + t.Fatalf("count: %v", err) + } + if n >= 3 { + break + } + time.Sleep(20 * time.Millisecond) + } + + rows, err := appviewDB.Query(` + select spindle, pipeline_knot, pipeline_rkey, workflow, status + from pipeline_statuses + order by created asc + `) + if err != nil { + t.Fatalf("query: %v", err) + } + defer rows.Close() + + type rec struct { + spindle, knot, rkey, workflow, status string + } + var got []rec + for rows.Next() { + var r rec + if err := rows.Scan(&r.spindle, &r.knot, &r.rkey, &r.workflow, &r.status); err != nil { + t.Fatalf("scan: %v", err) + } + got = append(got, r) + } + + if len(got) != 3 { + t.Fatalf("pipeline_statuses rows = %d, want 3: %+v", len(got), got) + } + + wantStatuses := []string{"pending", "running", "success"} + gotStatuses := map[string]bool{} + for _, r := range got { + gotStatuses[r.status] = true + if r.spindle != source.Host { + t.Errorf("spindle = %q, want %q", r.spindle, source.Host) + } + if r.knot != workflowId.Knot { + t.Errorf("pipeline_knot = %q, want %q", r.knot, workflowId.Knot) + } + if r.rkey != workflowId.Rkey { + t.Errorf("pipeline_rkey = %q, want %q", r.rkey, workflowId.Rkey) + } + if r.workflow != workflowId.Name { + t.Errorf("workflow = %q, want %q", r.workflow, workflowId.Name) + } + } + for _, want := range wantStatuses { + if !gotStatuses[want] { + t.Errorf("missing status %q in projection", want) + } + } +} diff --git a/appview/state/streams.go b/appview/state/streams.go new file mode 100644 index 00000000..1205c05b --- /dev/null +++ b/appview/state/streams.go @@ -0,0 +1,56 @@ +package state + +import ( + "context" + + "tangled.org/core/appview/cache" + "tangled.org/core/appview/config" + ec "tangled.org/core/eventconsumer" + "tangled.org/core/eventconsumer/cursor" + "tangled.org/core/log" +) + +func bootstrapStream( + ctx context.Context, + name string, + kind ec.Kind, + hosts []string, + redisAddr string, + streamCfg config.ConsumerConfig, + dev bool, + processFn ec.ProcessFunc, +) *ec.Consumer { + logger := log.SubLogger(log.FromContext(ctx), name) + + redisCache := cache.New(redisAddr) + cursorStore := cursor.NewRedisCursorStore(redisCache) + + srcs := make(map[ec.Source]struct{}, len(hosts)) + for _, h := range hosts { + src := ec.Source{Kind: kind, Host: h} + migrateLegacyCursor(&cursorStore, src) + srcs[src] = struct{}{} + } + + return ec.NewConsumer(ec.ConsumerConfig{ + Sources: srcs, + ProcessFunc: processFn, + RetryInterval: streamCfg.RetryInterval, + MaxRetryInterval: streamCfg.MaxRetryInterval, + ConnectionTimeout: streamCfg.ConnectionTimeout, + WorkerCount: streamCfg.WorkerCount, + QueueSize: streamCfg.QueueSize, + Logger: logger, + URLFunc: ec.DefaultURL(dev), + CursorStore: &cursorStore, + }) +} + +func migrateLegacyCursor(store cursor.Store, src ec.Source) { + if store.Get(src.Key()) != 0 { + return + } + if legacy := store.Get(src.Host); legacy != 0 { + store.Set(src.Key(), legacy) + } +} diff --git a/appview/state/streams_test.go b/appview/state/streams_test.go new file mode 100644 index 00000000..2076f1ca --- /dev/null +++ b/appview/state/streams_test.go @@ -0,0 +1,56 @@ +package state + +import ( + "testing" + + ec "tangled.org/core/eventconsumer" + "tangled.org/core/eventconsumer/cursor" +) + +func TestMigrateLegacyCursor_CopiesBareHostToKindKey(t *testing.T) { + store := &cursor.MemoryStore{} + store.Set("clam.oyster.cafe", 1700000000123456789) + + migrateLegacyCursor(store, ec.NewKnotSource("clam.oyster.cafe")) + + if got := store.Get("knot:clam.oyster.cafe"); got != 1700000000123456789 { + t.Fatalf("new key cursor = %d, want legacy value", got) + } +} + +func TestMigrateLegacyCursor_DoesNotClobberAdvancedCursor(t *testing.T) { + store := &cursor.MemoryStore{} + store.Set("knot:whelk.oyster.cafe", 999) + store.Set("whelk.oyster.cafe", 100) + + migrateLegacyCursor(store, ec.NewKnotSource("whelk.oyster.cafe")) + + if got := store.Get("knot:whelk.oyster.cafe"); got != 999 { + t.Fatalf("new key cursor = %d, want it left untouched at 999", got) + } +} + +func TestMigrateLegacyCursor_NoLegacyIsNoOp(t *testing.T) { + store := &cursor.MemoryStore{} + + migrateLegacyCursor(store, ec.NewKnotSource("limpet.nel.pet")) + + if got := store.Get("knot:limpet.nel.pet"); got != 0 { + t.Fatalf("new key cursor = %d, want 0", got) + } +} + +func TestMigrateLegacyCursor_KindsStayNamespaced(t *testing.T) { + store := &cursor.MemoryStore{} + store.Set("mussel.oyster.cafe", 500) + + migrateLegacyCursor(store, ec.NewKnotSource("mussel.oyster.cafe")) + migrateLegacyCursor(store, ec.NewSpindleSource("mussel.oyster.cafe")) + + if got := store.Get("knot:mussel.oyster.cafe"); got != 500 { + t.Fatalf("knot cursor = %d, want 500", got) + } + if got := store.Get("spindle:mussel.oyster.cafe"); got != 500 { + t.Fatalf("spindle cursor = %d, want 500", got) + } +} -- 2.51.2