Monorepo for Tangled
0

Configure Feed

Select the types of activity you want to include in your feed.

appview/state: switch streams to eventstream, backfill legacy cursors

Lewis: May this revision serve well! <lewis@tangled.org>

Lewis (May 22, 2026, 12:38 PM +0300) 20301f18 90df2c63

+381 -79
+12
appview/knots/knots.go
··· 314 314 return 315 315 } 316 316 317 + if registration.Registered != nil { 318 + remaining, rErr := db.GetRegistrations(k.Db, 319 + orm.FilterEq("domain", domain), 320 + orm.FilterIsNot("registered", "null"), 321 + ) 322 + if rErr != nil { 323 + l.Warn("failed to check remaining registrations after delete", "err", rErr) 324 + } else if len(remaining) == 0 { 325 + go k.Knotstream.RemoveSource(eventconsumer.NewKnotSource(domain)) 326 + } 327 + } 328 + 317 329 shouldRedirect := r.Header.Get("shouldRedirect") 318 330 if shouldRedirect == "true" { 319 331 k.Pages.HxRedirect(w, "/knots")
+10 -1
appview/repo/repo.go
··· 168 168 return 169 169 } 170 170 171 + oldSpindle := f.Spindle 172 + if oldSpindle != "" && oldSpindle != newSpindle { 173 + remaining, qErr := db.GetRepos(rp.db, orm.FilterEq("spindle", oldSpindle)) 174 + if qErr != nil { 175 + l.Warn("failed to count repos using old spindle", "err", qErr) 176 + } else if len(remaining) == 0 { 177 + rp.spindlestream.RemoveSource(eventconsumer.NewSpindleSource(oldSpindle)) 178 + } 179 + } 180 + 171 181 if !removingSpindle { 172 - // add this spindle to spindle stream 173 182 rp.spindlestream.AddSource( 174 183 context.Background(), 175 184 eventconsumer.NewSpindleSource(newSpindle),
+50
appview/state/cursor_migrate.go
··· 1 + package state 2 + 3 + import ( 4 + "context" 5 + "log/slog" 6 + "strconv" 7 + 8 + "tangled.org/core/appview/cache" 9 + "tangled.org/core/tid" 10 + ) 11 + 12 + func legacyCursorTID(oldVal string) (string, bool) { 13 + nanos, err := strconv.ParseInt(oldVal, 10, 64) 14 + if err != nil { 15 + return "", false 16 + } 17 + gen := tid.MonotonicGenerator{} 18 + return gen.FromNanos(nanos), true 19 + } 20 + 21 + func migrateLegacyCursor(ctx context.Context, logger *slog.Logger, c *cache.Cache, host, newPrefix string) { 22 + if c == nil { 23 + return 24 + } 25 + oldKey := "cursor:" + host 26 + newKey := "cursor:" + newPrefix + ":" + host 27 + 28 + if n, _ := c.Exists(ctx, newKey).Result(); n > 0 { 29 + return 30 + } 31 + 32 + oldVal, err := c.Get(ctx, oldKey).Result() 33 + if err != nil || oldVal == "" { 34 + return 35 + } 36 + 37 + newTid, ok := legacyCursorTID(oldVal) 38 + if !ok { 39 + return 40 + } 41 + 42 + if err := c.Set(ctx, newKey, newTid, 0).Err(); err != nil { 43 + logger.Warn("cursor backfill: set new key failed", "host", host, "err", err) 44 + return 45 + } 46 + if err := c.Del(ctx, oldKey).Err(); err != nil { 47 + logger.Warn("cursor backfill: delete old key failed", "host", host, "err", err) 48 + } 49 + logger.Info("cursor backfill: migrated", "host", host, "nanos", oldVal, "tid", newTid, "prefix", newPrefix) 50 + }
+83
appview/state/cursor_migrate_test.go
··· 1 + package state 2 + 3 + import ( 4 + "testing" 5 + 6 + "github.com/bluesky-social/indigo/atproto/syntax" 7 + ) 8 + 9 + func TestLegacyCursorTID_ConvertsNanos(t *testing.T) { 10 + nanos := int64(1_700_000_000_123_456_789) 11 + got, ok := legacyCursorTID("1700000000123456789") 12 + if !ok { 13 + t.Fatal("conversion failed for valid nanos") 14 + } 15 + parsed, err := syntax.ParseTID(got) 16 + if err != nil { 17 + t.Fatalf("not a valid TID: %q err=%v", got, err) 18 + } 19 + wantMicros := nanos / 1000 20 + if parsed.Integer() != uint64(wantMicros)<<10 { 21 + t.Fatalf("TID encodes micros=%d, want %d", parsed.Integer()>>10, wantMicros) 22 + } 23 + } 24 + 25 + func TestLegacyCursorTID_RejectsNonNumeric(t *testing.T) { 26 + _, ok := legacyCursorTID("3kmwd6xk6ww22") 27 + if ok { 28 + t.Fatal("should reject existing TID-format value") 29 + } 30 + _, ok = legacyCursorTID("") 31 + if ok { 32 + t.Fatal("should reject empty string") 33 + } 34 + _, ok = legacyCursorTID("not-a-number") 35 + if ok { 36 + t.Fatal("should reject garbage") 37 + } 38 + } 39 + 40 + func TestLegacyCursorTID_RoundTripsThroughTIDClock(t *testing.T) { 41 + cases := []int64{ 42 + 1_000_000, 43 + 1_700_000_000_000_000_000, 44 + 1_700_000_000_000_001_000, 45 + 1_700_000_000_000_002_000, 46 + } 47 + seen := map[string]struct{}{} 48 + for _, n := range cases { 49 + got, ok := legacyCursorTID(itoa(n)) 50 + if !ok { 51 + t.Fatalf("conversion failed for %d", n) 52 + } 53 + if _, err := syntax.ParseTID(got); err != nil { 54 + t.Fatalf("nanos=%d produced non-TID %q: %v", n, got, err) 55 + } 56 + if _, dup := seen[got]; dup { 57 + t.Fatalf("nanos=%d produced duplicate TID %q", n, got) 58 + } 59 + seen[got] = struct{}{} 60 + } 61 + } 62 + 63 + func itoa(n int64) string { 64 + var b [20]byte 65 + i := len(b) 66 + neg := n < 0 67 + if neg { 68 + n = -n 69 + } 70 + if n == 0 { 71 + return "0" 72 + } 73 + for n > 0 { 74 + i-- 75 + b[i] = byte('0' + n%10) 76 + n /= 10 77 + } 78 + if neg { 79 + i-- 80 + b[i] = '-' 81 + } 82 + return string(b[i:]) 83 + }
+19 -39
appview/state/knotstream.go
··· 14 14 "tangled.org/core/appview/notify" 15 15 16 16 "tangled.org/core/api/tangled" 17 - "tangled.org/core/appview/cache" 18 17 "tangled.org/core/appview/config" 19 18 "tangled.org/core/appview/db" 20 19 "tangled.org/core/appview/models" 21 20 "tangled.org/core/appview/sites" 22 21 ec "tangled.org/core/eventconsumer" 23 - "tangled.org/core/eventconsumer/cursor" 22 + "tangled.org/core/eventstream" 24 23 knotdb "tangled.org/core/knotserver/db" 25 24 "tangled.org/core/log" 26 25 "tangled.org/core/orm" ··· 33 32 ) 34 33 35 34 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) { 36 - logger := log.FromContext(ctx) 37 - logger = log.SubLogger(logger, "knotstream") 38 - 39 - knots, err := db.GetRegistrations( 40 - d, 41 - orm.FilterIsNot("registered", "null"), 42 - ) 35 + knots, err := db.GetRegistrations(d, orm.FilterIsNot("registered", "null")) 43 36 if err != nil { 44 37 return nil, err 45 38 } 46 39 47 - srcs := make(map[ec.Source]struct{}) 48 - for _, k := range knots { 49 - s := ec.NewKnotSource(k.Domain) 50 - srcs[s] = struct{}{} 51 - } 52 - 53 - cache := cache.New(c.Redis.Addr) 54 - cursorStore := cursor.NewRedisCursorStore(cache) 55 - 56 - cfg := ec.ConsumerConfig{ 57 - Sources: srcs, 58 - ProcessFunc: knotIngester(d, enforcer, posthog, notifier, c.Core.Dev, c, cfClient), 59 - RetryInterval: c.Knotstream.RetryInterval, 60 - MaxRetryInterval: c.Knotstream.MaxRetryInterval, 61 - ConnectionTimeout: c.Knotstream.ConnectionTimeout, 62 - WorkerCount: c.Knotstream.WorkerCount, 63 - QueueSize: c.Knotstream.QueueSize, 64 - Logger: logger, 65 - Dev: c.Core.Dev, 66 - CursorStore: &cursorStore, 40 + hosts := make([]string, len(knots)) 41 + for i, k := range knots { 42 + hosts[i] = k.Domain 67 43 } 68 44 69 - return ec.NewConsumer(cfg), nil 45 + return bootstrapStream( 46 + ctx, "knotstream", ec.KindKnot, hosts, c.Redis.Addr, 47 + c.Knotstream, c.Core.Dev, 48 + knotIngester(d, enforcer, posthog, notifier, c.Core.Dev, c, cfClient), 49 + ), nil 70 50 } 71 51 72 52 func resolveRepo(d *db.DB, repoDid *string, ownerDid, repoName string) (*models.Repo, error) { ··· 84 64 } 85 65 86 66 func knotIngester(d *db.DB, enforcer *rbac.Enforcer, posthog posthog.Client, notifier notify.Notifier, dev bool, c *config.Config, cfClient *cloudflare.Client) ec.ProcessFunc { 87 - return func(ctx context.Context, source ec.Source, msg ec.Message) error { 67 + return func(ctx context.Context, source ec.Source, msg eventstream.Event) error { 88 68 switch msg.Nsid { 89 69 case tangled.GitRefUpdateNSID: 90 70 return ingestRefUpdate(ctx, d, enforcer, posthog, notifier, dev, c, cfClient, source, msg) ··· 99 79 } 100 80 101 81 // TODO(boltless): remove this. knotmirror should do all sort of indexing 102 - 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 { 82 + 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 { 103 83 logger := log.FromContext(ctx) 104 84 105 85 var record tangled.GitRefUpdate ··· 112 92 if err != nil { 113 93 return err 114 94 } 115 - if !slices.Contains(knownKnots, source.Key()) { 116 - return fmt.Errorf("%s does not belong to %s, something is fishy", record.CommitterDid, source.Key()) 95 + if !slices.Contains(knownKnots, source.Host) { 96 + return fmt.Errorf("%s does not belong to %s, something is fishy", record.CommitterDid, source.Host) 117 97 } 118 98 119 99 if record.Repo == "" { 120 - return fmt.Errorf("gitRefUpdate from %s missing repo", source.Key()) 100 + return fmt.Errorf("gitRefUpdate from %s missing repo", source.Host) 121 101 } 122 102 123 103 repo, lookupErr := db.GetRepoByDid(d, record.Repo) ··· 285 265 return tx.Commit() 286 266 } 287 267 288 - func ingestPipeline(d *db.DB, source ec.Source, msg ec.Message) error { 268 + func ingestPipeline(d *db.DB, source ec.Source, msg eventstream.Event) error { 289 269 var record tangled.Pipeline 290 270 err := json.Unmarshal(msg.EventJson, &record) 291 271 if err != nil { ··· 343 323 344 324 pipeline := models.Pipeline{ 345 325 Rkey: msg.Rkey, 346 - Knot: source.Key(), 326 + Knot: source.Host, 347 327 RepoOwner: syntax.DID(record.TriggerMetadata.Repo.Did), 348 328 RepoName: repoName, 349 329 RepoDid: repo.RepoDid, ··· 364 344 return nil 365 345 } 366 346 367 - func ingestDIDAssign(d *db.DB, enforcer *rbac.Enforcer, source ec.Source, msg ec.Message, ctx context.Context) error { 347 + func ingestDIDAssign(d *db.DB, enforcer *rbac.Enforcer, source ec.Source, msg eventstream.Event, ctx context.Context) error { 368 348 logger := log.FromContext(ctx) 369 349 370 350 var record knotdb.RepoDIDAssign ··· 393 373 return nil 394 374 } 395 375 repo := repos[0] 396 - knot := source.Key() 376 + knot := source.Host 397 377 398 378 if repo.Knot != knot { 399 379 return fmt.Errorf("didAssign from %s for repo hosted on %s, rejecting", knot, repo.Knot)
+15 -39
appview/state/spindlestream.go
··· 4 4 "context" 5 5 "encoding/json" 6 6 "fmt" 7 - "log/slog" 8 7 "strings" 9 8 "time" 10 9 11 10 "github.com/bluesky-social/indigo/atproto/syntax" 12 11 "tangled.org/core/api/tangled" 13 - "tangled.org/core/appview/cache" 14 12 "tangled.org/core/appview/config" 15 13 "tangled.org/core/appview/db" 16 14 "tangled.org/core/appview/models" 17 15 ec "tangled.org/core/eventconsumer" 18 - "tangled.org/core/eventconsumer/cursor" 19 - "tangled.org/core/log" 16 + "tangled.org/core/eventstream" 20 17 "tangled.org/core/orm" 21 18 "tangled.org/core/rbac" 22 19 spindle "tangled.org/core/spindle/models" 23 20 ) 24 21 25 22 func Spindlestream(ctx context.Context, c *config.Config, d *db.DB, enforcer *rbac.Enforcer) (*ec.Consumer, error) { 26 - logger := log.FromContext(ctx) 27 - logger = log.SubLogger(logger, "spindlestream") 28 - 29 - spindles, err := db.GetSpindles( 30 - ctx, 31 - d, 32 - orm.FilterIsNot("verified", "null"), 33 - ) 23 + spindles, err := db.GetSpindles(ctx, d, orm.FilterIsNot("verified", "null")) 34 24 if err != nil { 35 25 return nil, err 36 26 } 37 27 38 - srcs := make(map[ec.Source]struct{}) 39 - for _, s := range spindles { 40 - src := ec.NewSpindleSource(s.Instance) 41 - srcs[src] = struct{}{} 42 - } 43 - 44 - cache := cache.New(c.Redis.Addr) 45 - cursorStore := cursor.NewRedisCursorStore(cache) 46 - 47 - cfg := ec.ConsumerConfig{ 48 - Sources: srcs, 49 - ProcessFunc: spindleIngester(ctx, logger, d), 50 - RetryInterval: c.Spindlestream.RetryInterval, 51 - MaxRetryInterval: c.Spindlestream.MaxRetryInterval, 52 - ConnectionTimeout: c.Spindlestream.ConnectionTimeout, 53 - WorkerCount: c.Spindlestream.WorkerCount, 54 - QueueSize: c.Spindlestream.QueueSize, 55 - Logger: logger, 56 - Dev: c.Core.Dev, 57 - CursorStore: &cursorStore, 28 + hosts := make([]string, len(spindles)) 29 + for i, s := range spindles { 30 + hosts[i] = s.Instance 58 31 } 59 32 60 - return ec.NewConsumer(cfg), nil 33 + return bootstrapStream( 34 + ctx, "spindlestream", ec.KindSpindle, hosts, c.Redis.Addr, 35 + c.Spindlestream, c.Core.Dev, 36 + spindleIngester(d), 37 + ), nil 61 38 } 62 39 63 - func spindleIngester(ctx context.Context, logger *slog.Logger, d *db.DB) ec.ProcessFunc { 64 - return func(ctx context.Context, source ec.Source, msg ec.Message) error { 40 + func spindleIngester(d *db.DB) ec.ProcessFunc { 41 + return func(ctx context.Context, source ec.Source, msg eventstream.Event) error { 65 42 switch msg.Nsid { 66 43 case tangled.PipelineStatusNSID: 67 - return ingestPipelineStatus(ctx, logger, d, source, msg) 44 + return ingestPipelineStatus(ctx, d, source, msg) 68 45 } 69 - 70 46 return nil 71 47 } 72 48 } 73 49 74 - func ingestPipelineStatus(ctx context.Context, logger *slog.Logger, d *db.DB, source ec.Source, msg ec.Message) error { 50 + func ingestPipelineStatus(ctx context.Context, d *db.DB, source ec.Source, msg eventstream.Event) error { 75 51 var record tangled.PipelineStatus 76 52 err := json.Unmarshal(msg.EventJson, &record) 77 53 if err != nil { ··· 95 71 } 96 72 97 73 status := models.PipelineStatus{ 98 - Spindle: source.Key(), 74 + Spindle: source.Host, 99 75 Rkey: msg.Rkey, 100 76 PipelineKnot: strings.TrimPrefix(pipelineUri.Authority().String(), "did:web:"), 101 77 PipelineRkey: pipelineUri.RecordKey().String(),
+144
appview/state/spindlestream_test.go
··· 1 + package state 2 + 3 + import ( 4 + "context" 5 + "io" 6 + "log/slog" 7 + "net/http" 8 + "net/http/httptest" 9 + "path/filepath" 10 + "strings" 11 + "testing" 12 + "time" 13 + 14 + "tangled.org/core/appview/db" 15 + ec "tangled.org/core/eventconsumer" 16 + "tangled.org/core/eventconsumer/cursor" 17 + "tangled.org/core/eventstream" 18 + "tangled.org/core/notifier" 19 + spindledb "tangled.org/core/spindle/db" 20 + spindlemodels "tangled.org/core/spindle/models" 21 + ) 22 + 23 + func TestColdStart_SpindleEventsRebuildPipelineStatuses(t *testing.T) { 24 + ctx := t.Context() 25 + 26 + spindleDB, err := spindledb.Make(ctx, filepath.Join(t.TempDir(), "spindle.db")) 27 + if err != nil { 28 + t.Fatalf("spindle Make: %v", err) 29 + } 30 + t.Cleanup(func() { spindleDB.Close() }) 31 + 32 + n := notifier.New() 33 + workflowId := spindlemodels.WorkflowId{ 34 + PipelineId: spindlemodels.PipelineId{Knot: "knot.boltless.example", Rkey: "pipeline-rk1"}, 35 + Name: "build", 36 + } 37 + for _, step := range []func() error{ 38 + func() error { return spindleDB.StatusPending(workflowId, &n) }, 39 + func() error { return spindleDB.StatusRunning(workflowId, &n) }, 40 + func() error { return spindleDB.StatusSuccess(workflowId, &n) }, 41 + } { 42 + if err := step(); err != nil { 43 + t.Fatalf("seed spindle event: %v", err) 44 + } 45 + } 46 + 47 + mux := http.NewServeMux() 48 + mux.HandleFunc("/events", func(w http.ResponseWriter, r *http.Request) { 49 + _ = eventstream.Stream(w, r, eventstream.StreamConfig{ 50 + Backend: spindleDB, 51 + Notifier: &n, 52 + Logger: slog.New(slog.NewTextHandler(io.Discard, nil)), 53 + }) 54 + }) 55 + srv := httptest.NewServer(mux) 56 + t.Cleanup(srv.Close) 57 + source := ec.Source{Kind: "test", Host: strings.TrimPrefix(srv.URL, "http://")} 58 + 59 + appviewDB, err := db.Make(ctx, filepath.Join(t.TempDir(), "appview.db")) 60 + if err != nil { 61 + t.Fatalf("appview Make: %v", err) 62 + } 63 + t.Cleanup(func() { appviewDB.Close() }) 64 + 65 + logger := slog.New(slog.NewTextHandler(io.Discard, nil)) 66 + processFunc := spindleIngester(appviewDB) 67 + 68 + cfg := ec.ConsumerConfig{ 69 + ProcessFunc: processFunc, 70 + WorkerCount: 1, 71 + QueueSize: 16, 72 + ConnectionTimeout: 2 * time.Second, 73 + CursorStore: &cursor.MemoryStore{}, 74 + URLFunc: ec.DefaultURL(true), 75 + Logger: logger, 76 + } 77 + c := ec.NewConsumer(cfg) 78 + 79 + consumerCtx, cancel := context.WithCancel(ctx) 80 + defer cancel() 81 + c.Start(consumerCtx) 82 + c.AddSource(consumerCtx, source) 83 + 84 + deadline := time.Now().Add(3 * time.Second) 85 + for time.Now().Before(deadline) { 86 + var n int 87 + if err := appviewDB.QueryRow(`select count(*) from pipeline_statuses`).Scan(&n); err != nil { 88 + t.Fatalf("count: %v", err) 89 + } 90 + if n >= 3 { 91 + break 92 + } 93 + time.Sleep(20 * time.Millisecond) 94 + } 95 + 96 + rows, err := appviewDB.Query(` 97 + select spindle, pipeline_knot, pipeline_rkey, workflow, status 98 + from pipeline_statuses 99 + order by created asc 100 + `) 101 + if err != nil { 102 + t.Fatalf("query: %v", err) 103 + } 104 + defer rows.Close() 105 + 106 + type rec struct { 107 + spindle, knot, rkey, workflow, status string 108 + } 109 + var got []rec 110 + for rows.Next() { 111 + var r rec 112 + if err := rows.Scan(&r.spindle, &r.knot, &r.rkey, &r.workflow, &r.status); err != nil { 113 + t.Fatalf("scan: %v", err) 114 + } 115 + got = append(got, r) 116 + } 117 + 118 + if len(got) != 3 { 119 + t.Fatalf("pipeline_statuses rows = %d, want 3: %+v", len(got), got) 120 + } 121 + 122 + wantStatuses := []string{"pending", "running", "success"} 123 + gotStatuses := map[string]bool{} 124 + for _, r := range got { 125 + gotStatuses[r.status] = true 126 + if r.spindle != source.Host { 127 + t.Errorf("spindle = %q, want %q", r.spindle, source.Host) 128 + } 129 + if r.knot != workflowId.Knot { 130 + t.Errorf("pipeline_knot = %q, want %q", r.knot, workflowId.Knot) 131 + } 132 + if r.rkey != workflowId.Rkey { 133 + t.Errorf("pipeline_rkey = %q, want %q", r.rkey, workflowId.Rkey) 134 + } 135 + if r.workflow != workflowId.Name { 136 + t.Errorf("workflow = %q, want %q", r.workflow, workflowId.Name) 137 + } 138 + } 139 + for _, want := range wantStatuses { 140 + if !gotStatuses[want] { 141 + t.Errorf("missing status %q in projection", want) 142 + } 143 + } 144 + }
+48
appview/state/streams.go
··· 1 + package state 2 + 3 + import ( 4 + "context" 5 + 6 + "tangled.org/core/appview/cache" 7 + "tangled.org/core/appview/config" 8 + ec "tangled.org/core/eventconsumer" 9 + "tangled.org/core/eventconsumer/cursor" 10 + "tangled.org/core/log" 11 + ) 12 + 13 + func bootstrapStream( 14 + ctx context.Context, 15 + name string, 16 + kind ec.Kind, 17 + hosts []string, 18 + redisAddr string, 19 + streamCfg config.ConsumerConfig, 20 + dev bool, 21 + processFn ec.ProcessFunc, 22 + ) *ec.Consumer { 23 + logger := log.SubLogger(log.FromContext(ctx), name) 24 + 25 + srcs := make(map[ec.Source]struct{}, len(hosts)) 26 + for _, h := range hosts { 27 + srcs[ec.Source{Kind: kind, Host: h}] = struct{}{} 28 + } 29 + 30 + redisCache := cache.New(redisAddr) 31 + for _, h := range hosts { 32 + migrateLegacyCursor(ctx, logger, redisCache, h, string(kind)) 33 + } 34 + cursorStore := cursor.NewRedisCursorStore(redisCache) 35 + 36 + return ec.NewConsumer(ec.ConsumerConfig{ 37 + Sources: srcs, 38 + ProcessFunc: processFn, 39 + RetryInterval: streamCfg.RetryInterval, 40 + MaxRetryInterval: streamCfg.MaxRetryInterval, 41 + ConnectionTimeout: streamCfg.ConnectionTimeout, 42 + WorkerCount: streamCfg.WorkerCount, 43 + QueueSize: streamCfg.QueueSize, 44 + Logger: logger, 45 + URLFunc: ec.DefaultURL(dev), 46 + CursorStore: &cursorStore, 47 + }) 48 + }