Monorepo for Tangled
0

Configure Feed

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

knotserver/events: switch to eventstream, TID migration

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

Lewis (May 22, 2026, 12:38 PM +0300) fecc6aab ef9f8e4c

+324 -207
+9 -1
knotserver/db/db.go
··· 10 10 11 11 securejoin "github.com/cyphar/filepath-securejoin" 12 12 _ "github.com/mattn/go-sqlite3" 13 + "tangled.org/core/eventstream" 13 14 "tangled.org/core/log" 14 15 "tangled.org/core/orm" 15 16 ) ··· 68 69 ); 69 70 70 71 create table if not exists events ( 72 + tid text not null, 71 73 rkey text not null, 72 74 nsid text not null, 73 75 event text not null, -- json 74 - created integer not null default (strftime('%s', 'now')), 76 + created integer not null, 75 77 primary key (rkey, nsid) 76 78 ); 77 79 ··· 175 177 on repo_keys(owner_did, repo_name); 176 178 `) 177 179 return mErr 180 + }); err != nil { 181 + return nil, err 182 + } 183 + 184 + if err := orm.RunMigration(conn, logger, "add-tid-to-events", func(tx *sql.Tx) error { 185 + return eventstream.MigrateAddTID(ctx, tx) 178 186 }); err != nil { 179 187 return nil, err 180 188 }
+7 -59
knotserver/db/events.go
··· 3 3 import ( 4 4 "encoding/json" 5 5 "fmt" 6 - "time" 7 6 7 + "tangled.org/core/eventstream" 8 8 "tangled.org/core/notifier" 9 9 "tangled.org/core/tid" 10 10 ) 11 11 12 - type Event struct { 13 - Rkey string `json:"rkey"` 14 - Nsid string `json:"nsid"` 15 - EventJson string `json:"event"` 16 - Created int64 `json:"created"` 17 - } 18 - 19 - func (d *DB) InsertEvent(event Event, notifier *notifier.Notifier) error { 20 - 21 - _, err := d.db.Exec( 22 - `insert into events (rkey, nsid, event, created) values (?, ?, ?, ?)`, 23 - event.Rkey, 24 - event.Nsid, 25 - event.EventJson, 26 - time.Now().UnixNano(), 27 - ) 28 - 29 - notifier.NotifyAll() 30 - 31 - return err 12 + func (d *DB) InsertEvent(event eventstream.Event, n *notifier.Notifier) error { 13 + return eventstream.Insert(d.db, event, n) 32 14 } 33 15 34 16 func (d *DB) EmitDIDAssign(n *notifier.Notifier, ownerDid, repoName, repoDid string) error { ··· 43 25 return fmt.Errorf("marshal didAssign event: %w", err) 44 26 } 45 27 46 - return d.InsertEvent(Event{ 28 + return d.InsertEvent(eventstream.Event{ 47 29 Rkey: tid.TID(), 48 30 Nsid: RepoDIDAssignNSID, 49 - EventJson: string(eventJson), 31 + EventJson: eventJson, 50 32 }, n) 51 33 } 52 34 53 - func (d *DB) GetEvents(cursor int64) ([]Event, error) { 54 - whereClause := "" 55 - args := []any{} 56 - if cursor > 0 { 57 - whereClause = "where created > ?" 58 - args = append(args, cursor) 59 - } 60 - 61 - query := fmt.Sprintf(` 62 - select rkey, nsid, event, created 63 - from events 64 - %s 65 - order by created asc 66 - limit 100 67 - `, whereClause) 68 - 69 - rows, err := d.db.Query(query, args...) 70 - if err != nil { 71 - return nil, err 72 - } 73 - defer rows.Close() 74 - 75 - var evts []Event 76 - for rows.Next() { 77 - var ev Event 78 - if err := rows.Scan(&ev.Rkey, &ev.Nsid, &ev.EventJson, &ev.Created); err != nil { 79 - return nil, err 80 - } 81 - evts = append(evts, ev) 82 - } 83 - 84 - if err := rows.Err(); err != nil { 85 - return nil, err 86 - } 87 - 88 - return evts, nil 35 + func (d *DB) GetEvents(cursorTid string, limit int) ([]eventstream.Event, error) { 36 + return eventstream.List(d.db, cursorTid, limit) 89 37 }
+273
knotserver/db/events_migration_test.go
··· 1 + package db 2 + 3 + import ( 4 + "context" 5 + "database/sql" 6 + "path/filepath" 7 + "sort" 8 + "testing" 9 + 10 + "github.com/bluesky-social/indigo/atproto/syntax" 11 + _ "github.com/mattn/go-sqlite3" 12 + "tangled.org/core/eventstream" 13 + ) 14 + 15 + func openLegacyDB(t *testing.T) *sql.DB { 16 + t.Helper() 17 + path := filepath.Join(t.TempDir(), "knot.db") 18 + d, err := sql.Open("sqlite3", path+"?_foreign_keys=1") 19 + if err != nil { 20 + t.Fatalf("open: %v", err) 21 + } 22 + t.Cleanup(func() { d.Close() }) 23 + 24 + if _, err := d.Exec(` 25 + create table events ( 26 + rkey text not null, 27 + nsid text not null, 28 + event text not null, 29 + created integer not null default (strftime('%s', 'now')), 30 + primary key (rkey, nsid) 31 + ); 32 + create table migrations ( 33 + id integer primary key autoincrement, 34 + name text unique 35 + ); 36 + `); err != nil { 37 + t.Fatalf("legacy schema: %v", err) 38 + } 39 + return d 40 + } 41 + 42 + func runEventsMigration(t *testing.T, d *sql.DB) { 43 + t.Helper() 44 + tx, err := d.Begin() 45 + if err != nil { 46 + t.Fatalf("begin: %v", err) 47 + } 48 + if err := eventstream.MigrateAddTID(context.Background(), tx); err != nil { 49 + tx.Rollback() 50 + t.Fatalf("migrate: %v", err) 51 + } 52 + if err := tx.Commit(); err != nil { 53 + t.Fatalf("commit: %v", err) 54 + } 55 + } 56 + 57 + func TestMigrateEventsAddTID_PreservesData(t *testing.T) { 58 + d := openLegacyDB(t) 59 + 60 + seed := []struct { 61 + rkey, nsid, event string 62 + created int64 63 + }{ 64 + {"r1", "sh.tangled.test", `{"i":1}`, 1_700_000_000_000_000_000}, 65 + {"r2", "sh.tangled.test", `{"i":2}`, 1_700_000_000_000_001_000}, 66 + {"r3", "sh.tangled.test", `{"i":3}`, 1_700_000_000_000_002_000}, 67 + {"r4", "sh.tangled.other", `{"i":4}`, 1_700_000_000_000_002_000}, 68 + } 69 + for _, s := range seed { 70 + if _, err := d.Exec( 71 + `insert into events (rkey, nsid, event, created) values (?, ?, ?, ?)`, 72 + s.rkey, s.nsid, s.event, s.created, 73 + ); err != nil { 74 + t.Fatalf("seed insert: %v", err) 75 + } 76 + } 77 + 78 + runEventsMigration(t, d) 79 + 80 + var count int 81 + if err := d.QueryRow(`select count(*) from events`).Scan(&count); err != nil { 82 + t.Fatalf("count: %v", err) 83 + } 84 + if count != len(seed) { 85 + t.Fatalf("row count = %d, want %d", count, len(seed)) 86 + } 87 + 88 + var hasTid int 89 + if err := d.QueryRow( 90 + `select count(*) from pragma_table_info('events') where name = 'tid'`, 91 + ).Scan(&hasTid); err != nil { 92 + t.Fatalf("pragma: %v", err) 93 + } 94 + if hasTid != 1 { 95 + t.Fatal("tid column missing after migration") 96 + } 97 + 98 + rows, err := d.Query(`select tid, rkey, created from events order by tid asc`) 99 + if err != nil { 100 + t.Fatalf("read back: %v", err) 101 + } 102 + defer rows.Close() 103 + 104 + type row struct { 105 + tid, rkey string 106 + created int64 107 + } 108 + var got []row 109 + for rows.Next() { 110 + var r row 111 + if err := rows.Scan(&r.tid, &r.rkey, &r.created); err != nil { 112 + t.Fatalf("scan: %v", err) 113 + } 114 + got = append(got, r) 115 + } 116 + 117 + for _, r := range got { 118 + if _, err := syntax.ParseTID(r.tid); err != nil { 119 + t.Fatalf("row %s has invalid TID %q: %v", r.rkey, r.tid, err) 120 + } 121 + } 122 + 123 + sorted := append([]row(nil), got...) 124 + sort.Slice(sorted, func(i, j int) bool { return sorted[i].created < sorted[j].created }) 125 + for i := range got { 126 + if got[i].rkey != sorted[i].rkey && got[i].created != sorted[i].created { 127 + t.Errorf("TID order diverges from created order at %d: %+v vs %+v", i, got[i], sorted[i]) 128 + } 129 + } 130 + 131 + for i := 1; i < len(got); i++ { 132 + if got[i].tid <= got[i-1].tid { 133 + t.Errorf("TID not strictly monotonic: %s <= %s", got[i].tid, got[i-1].tid) 134 + } 135 + } 136 + } 137 + 138 + func TestMigrateEventsAddTID_Idempotent(t *testing.T) { 139 + d := openLegacyDB(t) 140 + if _, err := d.Exec( 141 + `insert into events (rkey, nsid, event, created) values ('r1', 'n', '{}', 1)`, 142 + ); err != nil { 143 + t.Fatalf("seed: %v", err) 144 + } 145 + 146 + runEventsMigration(t, d) 147 + 148 + var firstTid string 149 + if err := d.QueryRow(`select tid from events where rkey = 'r1'`).Scan(&firstTid); err != nil { 150 + t.Fatalf("get first: %v", err) 151 + } 152 + 153 + runEventsMigration(t, d) 154 + 155 + var secondTid string 156 + if err := d.QueryRow(`select tid from events where rkey = 'r1'`).Scan(&secondTid); err != nil { 157 + t.Fatalf("get second: %v", err) 158 + } 159 + if firstTid != secondTid { 160 + t.Fatalf("second migration mutated TID: %q != %q", firstTid, secondTid) 161 + } 162 + } 163 + 164 + func TestSetupBootsLegacyDB(t *testing.T) { 165 + path := filepath.Join(t.TempDir(), "knot.db") 166 + d, err := sql.Open("sqlite3", path+"?_foreign_keys=1") 167 + if err != nil { 168 + t.Fatalf("open: %v", err) 169 + } 170 + if _, err := d.Exec(` 171 + create table events ( 172 + rkey text not null, 173 + nsid text not null, 174 + event text not null, 175 + created integer not null default (strftime('%s', 'now')), 176 + primary key (rkey, nsid) 177 + ); 178 + insert into events (rkey, nsid, event, created) values ('r1', 'n', '{}', 1700000000000000000); 179 + `); err != nil { 180 + d.Close() 181 + t.Fatalf("seed legacy schema: %v", err) 182 + } 183 + d.Close() 184 + 185 + db, err := Setup(context.Background(), path) 186 + if err != nil { 187 + t.Fatalf("Setup on legacy DB: %v", err) 188 + } 189 + t.Cleanup(func() { db.db.Close() }) 190 + 191 + var hasTid int 192 + if err := db.db.QueryRow( 193 + `select count(*) from pragma_table_info('events') where name = 'tid'`, 194 + ).Scan(&hasTid); err != nil { 195 + t.Fatalf("pragma: %v", err) 196 + } 197 + if hasTid != 1 { 198 + t.Fatal("tid column missing after Setup") 199 + } 200 + 201 + var hasIndex int 202 + if err := db.db.QueryRow( 203 + `select count(*) from sqlite_master where type='index' and name='idx_events_tid'`, 204 + ).Scan(&hasIndex); err != nil { 205 + t.Fatalf("pragma index: %v", err) 206 + } 207 + if hasIndex != 1 { 208 + t.Fatal("idx_events_tid missing after Setup") 209 + } 210 + } 211 + 212 + func TestSetupIdempotentOnFreshDB(t *testing.T) { 213 + path := filepath.Join(t.TempDir(), "knot.db") 214 + db, err := Setup(context.Background(), path) 215 + if err != nil { 216 + t.Fatalf("first Setup: %v", err) 217 + } 218 + db.db.Close() 219 + 220 + db, err = Setup(context.Background(), path) 221 + if err != nil { 222 + t.Fatalf("second Setup: %v", err) 223 + } 224 + t.Cleanup(func() { db.db.Close() }) 225 + 226 + var hasIndex int 227 + if err := db.db.QueryRow( 228 + `select count(*) from sqlite_master where type='index' and name='idx_events_tid'`, 229 + ).Scan(&hasIndex); err != nil { 230 + t.Fatalf("pragma index: %v", err) 231 + } 232 + if hasIndex != 1 { 233 + t.Fatal("idx_events_tid missing on fresh DB") 234 + } 235 + } 236 + 237 + func TestMigrateEventsAddTID_HandlesTiedCreated(t *testing.T) { 238 + d := openLegacyDB(t) 239 + if _, err := d.Exec(` 240 + insert into events (rkey, nsid, event, created) values 241 + ('a', 'n', '{}', 100), 242 + ('b', 'n', '{}', 100), 243 + ('c', 'n', '{}', 100); 244 + `); err != nil { 245 + t.Fatalf("seed: %v", err) 246 + } 247 + 248 + runEventsMigration(t, d) 249 + 250 + rows, err := d.Query(`select tid from events order by rkey asc`) 251 + if err != nil { 252 + t.Fatalf("read: %v", err) 253 + } 254 + defer rows.Close() 255 + var tids []string 256 + for rows.Next() { 257 + var t string 258 + if err := rows.Scan(&t); err != nil { 259 + break 260 + } 261 + tids = append(tids, t) 262 + } 263 + if len(tids) != 3 { 264 + t.Fatalf("got %d tids, want 3", len(tids)) 265 + } 266 + uniq := map[string]struct{}{} 267 + for _, tt := range tids { 268 + uniq[tt] = struct{}{} 269 + } 270 + if len(uniq) != 3 { 271 + t.Fatalf("TIDs not unique for tied created: %v", tids) 272 + } 273 + }
+13 -123
knotserver/events.go
··· 2 2 3 3 import ( 4 4 "context" 5 - "encoding/json" 5 + "errors" 6 6 "net/http" 7 - "strconv" 8 7 "time" 9 8 10 9 "github.com/bluesky-social/indigo/xrpc" 11 - "github.com/gorilla/websocket" 12 10 "tangled.org/core/api/tangled" 11 + "tangled.org/core/eventstream" 13 12 "tangled.org/core/log" 14 13 ) 15 14 16 - var upgrader = websocket.Upgrader{ 17 - ReadBufferSize: 1024, 18 - WriteBufferSize: 1024, 19 - } 20 - 21 15 func (h *Knot) Events(w http.ResponseWriter, r *http.Request) { 22 16 l := log.SubLogger(h.l, "eventstream") 23 17 l.Debug("received new connection") 24 18 25 - conn, err := upgrader.Upgrade(w, r, nil) 26 - if err != nil { 27 - l.Error("websocket upgrade failed", "err", err) 28 - w.WriteHeader(http.StatusInternalServerError) 29 - return 19 + err := eventstream.Stream(w, r, eventstream.StreamConfig{ 20 + Backend: h.db, 21 + Notifier: h.n, 22 + Logger: l, 23 + }) 24 + if err != nil && !errors.Is(err, eventstream.ErrDrainCap) { 25 + l.Error("event stream ended with error", "err", err) 30 26 } 31 - defer conn.Close() 32 - l.Debug("upgraded http to wss") 33 27 34 - ch := h.n.Subscribe() 35 - defer h.n.Unsubscribe(ch) 36 - 37 - ctx, cancel := context.WithCancel(r.Context()) 38 - defer cancel() 39 28 go func() { 40 - for { 41 - if _, _, err := conn.NextReader(); err != nil { 42 - l.Error("failed to read", "err", err) 43 - cancel() 44 - return 45 - } 29 + retryCtx, retryCancel := context.WithTimeout(context.Background(), 10*time.Second) 30 + defer retryCancel() 31 + if err := h.requestCrawl(retryCtx); err != nil { 32 + l.Error("error requesting crawls", "err", err) 46 33 } 47 34 }() 48 - 49 - var cursor int64 50 - cursorStr := r.URL.Query().Get("cursor") 51 - if cursorStr != "" { 52 - cursor, err = strconv.ParseInt(cursorStr, 10, 64) 53 - if err != nil { 54 - l.Error("invalid cursor, starting from beginning", "invalidCursor", cursorStr) 55 - cursor = 0 56 - } 57 - } 58 - 59 - l.Debug("going through backfill", "cursor", cursor) 60 - if err := h.drainBackfill(conn, &cursor, 10_000); err != nil { 61 - l.Error("failed to backfill", "err", err) 62 - return 63 - } 64 - 65 - // try request crawl when connection closed 66 - defer func() { 67 - go func() { 68 - retryCtx, retryCancel := context.WithTimeout(context.Background(), 10*time.Second) 69 - defer retryCancel() 70 - if err := h.requestCrawl(retryCtx); err != nil { 71 - l.Error("error requesting crawls", "err", err) 72 - } 73 - }() 74 - }() 75 - 76 - for { 77 - // wait for new data or timeout 78 - select { 79 - case <-ctx.Done(): 80 - l.Debug("stopping stream: client closed connection") 81 - return 82 - case <-ch: 83 - l.Debug("going through live data", "cursor", cursor) 84 - if _, err := h.streamOps(conn, &cursor); err != nil { 85 - l.Error("failed to stream", "err", err) 86 - return 87 - } 88 - case <-time.After(30 * time.Second): 89 - // send a keep-alive 90 - if err = conn.WriteControl(websocket.PingMessage, []byte{}, time.Now().Add(time.Second)); err != nil { 91 - l.Error("failed to write control", "err", err) 92 - } 93 - } 94 - } 95 - } 96 - 97 - func (h *Knot) drainBackfill(conn *websocket.Conn, cursor *int64, maxBatches int) error { 98 - for range maxBatches { 99 - n, err := h.streamOps(conn, cursor) 100 - if err != nil { 101 - return err 102 - } 103 - if n < 100 { 104 - return nil 105 - } 106 - } 107 - h.l.Warn("backfill hit batch limit", "maxBatches", maxBatches, "cursor", *cursor) 108 - return nil 109 - } 110 - 111 - func (h *Knot) streamOps(conn *websocket.Conn, cursor *int64) (int, error) { 112 - events, err := h.db.GetEvents(*cursor) 113 - if err != nil { 114 - h.l.Error("failed to fetch events from db", "err", err, "cursor", cursor) 115 - return 0, err 116 - } 117 - 118 - for _, event := range events { 119 - var eventJson map[string]any 120 - err := json.Unmarshal([]byte(event.EventJson), &eventJson) 121 - if err != nil { 122 - h.l.Error("failed to unmarshal event", "err", err) 123 - return 0, err 124 - } 125 - 126 - jsonMsg, err := json.Marshal(map[string]any{ 127 - "rkey": event.Rkey, 128 - "nsid": event.Nsid, 129 - "event": eventJson, 130 - "created": event.Created, 131 - }) 132 - if err != nil { 133 - h.l.Error("failed to marshal record", "err", err) 134 - return 0, err 135 - } 136 - 137 - if err := conn.WriteMessage(websocket.TextMessage, jsonMsg); err != nil { 138 - h.l.Debug("err", "err", err) 139 - return 0, err 140 - } 141 - *cursor = event.Created 142 - } 143 - 144 - return len(events), nil 145 35 } 146 36 147 37 func (h *Knot) requestCrawl(ctx context.Context) error {
+5 -3
knotserver/ingester.go
··· 16 16 jmodels "github.com/bluesky-social/jetstream/pkg/models" 17 17 "tangled.org/core/api/tangled" 18 18 "tangled.org/core/appview/models" 19 + "tangled.org/core/eventstream" 19 20 "tangled.org/core/knotserver/db" 20 21 "tangled.org/core/knotserver/git" 21 22 knotxrpc "tangled.org/core/knotserver/xrpc" 22 23 "tangled.org/core/log" 23 24 "tangled.org/core/rbac" 25 + "tangled.org/core/tid" 24 26 "tangled.org/core/workflow" 25 27 ) 26 28 ··· 355 357 return fmt.Errorf("failed to marshal pipeline event: %w", err) 356 358 } 357 359 358 - ev := db.Event{ 359 - Rkey: TID(), 360 + ev := eventstream.Event{ 361 + Rkey: tid.TID(), 360 362 Nsid: tangled.PipelineNSID, 361 - EventJson: string(eventJson), 363 + EventJson: eventJson, 362 364 } 363 365 364 366 l.Info("inserting pipeline event")
+8 -6
knotserver/internal.go
··· 16 16 "github.com/go-chi/chi/v5/middleware" 17 17 "github.com/go-git/go-git/v5/plumbing" 18 18 "tangled.org/core/api/tangled" 19 + "tangled.org/core/eventstream" 19 20 "tangled.org/core/hook" 20 21 "tangled.org/core/idresolver" 21 22 "tangled.org/core/knotserver/config" ··· 24 25 "tangled.org/core/log" 25 26 "tangled.org/core/notifier" 26 27 "tangled.org/core/rbac" 28 + "tangled.org/core/tid" 27 29 "tangled.org/core/workflow" 28 30 ) 29 31 ··· 310 312 return err 311 313 } 312 314 313 - event := db.Event{ 314 - Rkey: TID(), 315 + event := eventstream.Event{ 316 + Rkey: tid.TID(), 315 317 Nsid: tangled.GitRefUpdateNSID, 316 - EventJson: string(eventJson), 318 + EventJson: eventJson, 317 319 } 318 320 319 321 return h.db.InsertEvent(event, h.n) ··· 414 416 return nil 415 417 } 416 418 417 - event := db.Event{ 418 - Rkey: TID(), 419 + event := eventstream.Event{ 420 + Rkey: tid.TID(), 419 421 Nsid: tangled.PipelineNSID, 420 - EventJson: string(eventJson), 422 + EventJson: eventJson, 421 423 } 422 424 423 425 return h.db.InsertEvent(event, h.n)
-11
knotserver/util.go
··· 1 - package knotserver 2 - 3 - import ( 4 - "github.com/bluesky-social/indigo/atproto/syntax" 5 - ) 6 - 7 - var TIDClock = syntax.NewTIDClock(0) 8 - 9 - func TID() string { 10 - return TIDClock.Next().String() 11 - }
+9 -4
knotserver/xrpc/set_default_branch.go
··· 9 9 "github.com/bluesky-social/indigo/atproto/syntax" 10 10 "github.com/bluesky-social/indigo/xrpc" 11 11 "tangled.org/core/api/tangled" 12 - "tangled.org/core/knotserver/db" 12 + "tangled.org/core/eventstream" 13 13 "tangled.org/core/knotserver/git" 14 14 "tangled.org/core/rbac" 15 15 "tangled.org/core/tid" ··· 101 101 } 102 102 eventJson, err := json.Marshal(refUpdate) 103 103 if err != nil { 104 + fail(xrpcerr.GenericError(err)) 104 105 return 105 106 } 106 107 107 - x.Db.InsertEvent(db.Event{ 108 + if err := x.Db.InsertEvent(eventstream.Event{ 108 109 Rkey: tid.TID(), 109 110 Nsid: tangled.GitRefUpdateNSID, 110 - EventJson: string(eventJson), 111 - }, x.Notifier) 111 + EventJson: eventJson, 112 + }, x.Notifier); err != nil { 113 + l.Error("failed to insert event", "error", err) 114 + writeError(w, xrpcerr.GenericError(err), http.StatusInternalServerError) 115 + return 116 + } 112 117 113 118 w.WriteHeader(http.StatusOK) 114 119 }