diff --git a/appview/config/config.go b/appview/config/config.go index 88194e2be..651294ab6 100644 --- a/appview/config/config.go +++ b/appview/config/config.go @@ -68,14 +68,6 @@ type JetstreamConfig struct { Endpoint string `env:"ENDPOINT, default=wss://jetstream1.us-east.bsky.network/subscribe"` } -type ConsumerConfig struct { - RetryInterval time.Duration `env:"RETRY_INTERVAL, default=60s"` - MaxRetryInterval time.Duration `env:"MAX_RETRY_INTERVAL, default=120m"` - ConnectionTimeout time.Duration `env:"CONNECTION_TIMEOUT, default=5s"` - WorkerCount int `env:"WORKER_COUNT, default=64"` - QueueSize int `env:"QUEUE_SIZE, default=100"` -} - type ResendConfig struct { ApiKey string `env:"API_KEY"` SentFrom string `env:"SENT_FROM, default=noreply@notifs.tangled.sh"` @@ -194,27 +186,25 @@ func (cfg RedisConfig) ToURL() string { } type Config struct { - Core CoreConfig `env:",prefix=TANGLED_"` - Jetstream JetstreamConfig `env:",prefix=TANGLED_JETSTREAM_"` - Knotstream ConsumerConfig `env:",prefix=TANGLED_KNOTSTREAM_"` - Spindlestream ConsumerConfig `env:",prefix=TANGLED_SPINDLESTREAM_"` - Resend ResendConfig `env:",prefix=TANGLED_RESEND_"` - Posthog PosthogConfig `env:",prefix=TANGLED_POSTHOG_"` - Camo CamoConfig `env:",prefix=TANGLED_CAMO_"` - Avatar AvatarConfig `env:",prefix=TANGLED_AVATAR_"` - OAuth OAuthConfig `env:",prefix=TANGLED_OAUTH_"` - Redis RedisConfig `env:",prefix=TANGLED_REDIS_"` - Plc PlcConfig `env:",prefix=TANGLED_PLC_"` - Pds PdsConfig `env:",prefix=TANGLED_PDS_"` - Knot KnotConfig `env:",prefix=TANGLED_KNOT_"` - Cloudflare Cloudflare `env:",prefix=TANGLED_CLOUDFLARE_"` - Label LabelConfig `env:",prefix=TANGLED_LABEL_"` - Bluesky BlueskyConfig `env:",prefix=TANGLED_BLUESKY_"` - Sites SitesConfig `env:",prefix=TANGLED_SITES_"` - KnotMirror KnotMirrorConfig `env:",prefix=TANGLED_KNOTMIRROR_"` - Ogre OgreConfig `env:",prefix=TANGLED_OGRE_"` - SSH SSHConfig `env:",prefix=TANGLED_SSH_"` - CodeSearch CodeSearchConfig `env:",prefix=TANGLED_CODESEARCH_"` + Core CoreConfig `env:",prefix=TANGLED_"` + Jetstream JetstreamConfig `env:",prefix=TANGLED_JETSTREAM_"` + Resend ResendConfig `env:",prefix=TANGLED_RESEND_"` + Posthog PosthogConfig `env:",prefix=TANGLED_POSTHOG_"` + Camo CamoConfig `env:",prefix=TANGLED_CAMO_"` + Avatar AvatarConfig `env:",prefix=TANGLED_AVATAR_"` + OAuth OAuthConfig `env:",prefix=TANGLED_OAUTH_"` + Redis RedisConfig `env:",prefix=TANGLED_REDIS_"` + Plc PlcConfig `env:",prefix=TANGLED_PLC_"` + Pds PdsConfig `env:",prefix=TANGLED_PDS_"` + Knot KnotConfig `env:",prefix=TANGLED_KNOT_"` + Cloudflare Cloudflare `env:",prefix=TANGLED_CLOUDFLARE_"` + Label LabelConfig `env:",prefix=TANGLED_LABEL_"` + Bluesky BlueskyConfig `env:",prefix=TANGLED_BLUESKY_"` + Sites SitesConfig `env:",prefix=TANGLED_SITES_"` + KnotMirror KnotMirrorConfig `env:",prefix=TANGLED_KNOTMIRROR_"` + Ogre OgreConfig `env:",prefix=TANGLED_OGRE_"` + SSH SSHConfig `env:",prefix=TANGLED_SSH_"` + CodeSearch CodeSearchConfig `env:",prefix=TANGLED_CODESEARCH_"` } func LoadConfig(ctx context.Context) (*Config, error) { diff --git a/appview/db/db.go b/appview/db/db.go index b344e09f8..36a1b76e8 100644 --- a/appview/db/db.go +++ b/appview/db/db.go @@ -2252,18 +2252,6 @@ func Make(ctx context.Context, dbPath string) (*DB, error) { return err }) - orm.RunMigration(conn, logger, "add-knotacl-delta-cursor-table", func(tx *sql.Tx) error { - _, err := tx.Exec(` - create table if not exists knotacl_delta_cursor ( - scope_key text not null, - subject text not null, - cursor integer not null, - primary key (scope_key, subject) - ); - `) - return err - }) - orm.RunMigration(conn, logger, "migrate-knots-to-knot-owned-acl", func(tx *sql.Tx) error { _, err := tx.Exec(` update registrations set needs_upgrade = 1; diff --git a/appview/knotacl/client.go b/appview/knotacl/client.go index db26a4346..cdd6d2eda 100644 --- a/appview/knotacl/client.go +++ b/appview/knotacl/client.go @@ -10,6 +10,7 @@ import ( indigoxrpc "github.com/bluesky-social/indigo/xrpc" "tangled.org/core/api/tangled" + "tangled.org/core/xrpc/serviceauth" ) const ( @@ -41,85 +42,81 @@ func (c *Client) xrpcClient(host string) *indigoxrpc.Client { } func (c *Client) GetKnotMembers(ctx context.Context, host string) ([]string, error) { - ctx, cancel := context.WithTimeout(ctx, listDrainBudget) - defer cancel() + return getRoster(c, ctx, host, serviceauth.DidWeb(host).String(), + func(ctx context.Context, xc *indigoxrpc.Client, cursor, subject string) ([]string, *string, error) { + out, err := tangled.KnotListMembers(ctx, xc, cursor, listPageLimit, 0, "", subject) + if err != nil { + return nil, nil, err + } + return mapSlice(out.Items, func(i *tangled.KnotListMembers_ListItem) string { return i.Subject }), out.Cursor, nil + }, + func() { + c.logger.Warn("knot member list truncated before draining all pages", "host", host, "limit", maxListPages) + }) +} - xc := c.xrpcClient(host) - subjects, truncated, err := drainList( - "", - make(map[string]struct{}), - func(cursor string) ([]*tangled.KnotListMembers_ListItem, *string, error) { - out, err := tangled.KnotListMembers(ctx, xc, cursor, listPageLimit, 0, "", host) +func (c *Client) GetRepoCollaborators(ctx context.Context, host, repoDid string) ([]string, error) { + return getRoster(c, ctx, host, repoDid, + func(ctx context.Context, xc *indigoxrpc.Client, cursor, subject string) ([]string, *string, error) { + out, err := tangled.RepoListCollaborators(ctx, xc, cursor, listPageLimit, 0, "", subject) if err != nil { return nil, nil, err } - return out.Items, out.Cursor, nil + return mapSlice(out.Items, func(i *tangled.RepoListCollaborators_ListItem) string { return i.Subject }), out.Cursor, nil }, - func(i *tangled.KnotListMembers_ListItem) string { return i.Subject }, - ) - if err != nil { - return nil, err - } - if truncated { - c.logger.Warn("knot member list truncated before draining all pages", "host", host, "limit", maxListPages) + func() { + c.logger.Warn("repo collaborator list truncated before draining all pages", "host", host, "repoDid", repoDid, "limit", maxListPages) + }) +} + +func drainList( + page func(cursor string) ([]string, *string, error), +) (subjects []string, truncated bool, err error) { + seen := make(map[string]struct{}) + cursor := "" + for { + if len(seen) >= maxListPages { + return subjects, true, nil + } + if _, repeated := seen[cursor]; repeated { + return subjects, true, nil + } + seen[cursor] = struct{}{} + items, next, err := page(cursor) + if err != nil { + return nil, false, err + } + subjects = append(subjects, items...) + if len(items) == 0 || next == nil || *next == "" { + return subjects, false, nil + } + cursor = *next } - return dedup(subjects), nil } -func (c *Client) GetRepoCollaborators(ctx context.Context, host, repoDid string) ([]string, error) { +func getRoster( + c *Client, + ctx context.Context, + host, subject string, + list func(ctx context.Context, xc *indigoxrpc.Client, cursor, subject string) ([]string, *string, error), + warn func(), +) ([]string, error) { ctx, cancel := context.WithTimeout(ctx, listDrainBudget) defer cancel() xc := c.xrpcClient(host) subjects, truncated, err := drainList( - "", - make(map[string]struct{}), - func(cursor string) ([]*tangled.RepoListCollaborators_ListItem, *string, error) { - out, err := tangled.RepoListCollaborators(ctx, xc, cursor, listPageLimit, 0, "", repoDid) - if err != nil { - return nil, nil, err - } - return out.Items, out.Cursor, nil - }, - func(i *tangled.RepoListCollaborators_ListItem) string { return i.Subject }, + func(cursor string) ([]string, *string, error) { return list(ctx, xc, cursor, subject) }, ) if err != nil { return nil, err } if truncated { - c.logger.Warn("repo collaborator list truncated before draining all pages", "host", host, "repoDid", repoDid, "limit", maxListPages) + warn() } return dedup(subjects), nil } -func drainList[T any]( - cursor string, - seen map[string]struct{}, - page func(cursor string) ([]*T, *string, error), - subject func(*T) string, -) (subjects []string, truncated bool, err error) { - if len(seen) >= maxListPages { - return nil, true, nil - } - if _, repeated := seen[cursor]; repeated { - return nil, true, nil - } - seen[cursor] = struct{}{} - items, next, err := page(cursor) - if err != nil { - return nil, false, err - } - subjects = mapSlice(items, subject) - if len(items) == 0 || next == nil || *next == "" { - return subjects, false, nil - } - rest, truncated, err := drainList(*next, seen, page, subject) - if err != nil { - return nil, false, err - } - return append(subjects, rest...), truncated, nil -} - func dedup(subjects []string) []string { slices.Sort(subjects) return slices.Compact(subjects) diff --git a/appview/knotacl/client_test.go b/appview/knotacl/client_test.go index edd56654c..1517c364b 100644 --- a/appview/knotacl/client_test.go +++ b/appview/knotacl/client_test.go @@ -7,12 +7,14 @@ import ( "log/slog" "net/http" "net/http/httptest" + "net/url" "slices" "strings" "sync" "testing" "tangled.org/core/api/tangled" + "tangled.org/core/xrpc/serviceauth" ) func testLogger() *slog.Logger { @@ -145,6 +147,19 @@ func TestGetRepoCollaborators_PassesRepoDidAsSubject(t *testing.T) { t.Errorf("subject param missing repo DID: %v", calls) } } +func TestGetKnotMembers_PassesKnotDidAsSubject(t *testing.T) { + c, knot, host := devClientFor(t, func(w http.ResponseWriter, r *http.Request) { + json.NewEncoder(w).Encode(memberPage([]string{testCollab}, "")) + }) + + if _, err := c.GetKnotMembers(context.Background(), host); err != nil { + t.Fatalf("GetKnotMembers: %v", err) + } + want := url.Values{"subject": []string{serviceauth.DidWeb(host).String()}}.Encode() + if calls := knot.calls(); len(calls) != 1 || !strings.Contains(calls[0], want) { + t.Errorf("subject param missing the knot DID: %v", calls) + } +} func TestDrainStopsOnRepeatedCursor(t *testing.T) { c, knot, host := devClientFor(t, func(w http.ResponseWriter, r *http.Request) { diff --git a/appview/knotacl/roster.go b/appview/knotacl/roster.go index e82ca7168..8ec27fb55 100644 --- a/appview/knotacl/roster.go +++ b/appview/knotacl/roster.go @@ -8,7 +8,6 @@ import ( "log/slog" "slices" "sync" - "sync/atomic" "time" "github.com/bluesky-social/indigo/atproto/syntax" @@ -22,16 +21,12 @@ import ( const ( reconcileTTL = 5 * time.Minute reconcileBackoff = 15 * time.Second - cursorRetention = 1 * time.Hour ) var errReconcileBackoff = errors.New("reconcile suppressed during backoff") -type Cursor int64 - type scopeState struct { mu sync.Mutex - gen atomic.Uint64 failedAt time.Time refs int } @@ -67,7 +62,18 @@ func newRoster(store *db.DB, src lister, ttl time.Duration, now func() time.Time func (r *roster) GetKnotMembers(ctx context.Context, host string) ([]string, error) { return r.serve(ctx, memberScope(host), - func() error { return r.reconcileMembers(ctx, host) }, + func() error { + return r.reconcile(memberScope(host), + func() ([]string, error) { return r.src.GetKnotMembers(ctx, host) }, + func(tx *sql.Tx) error { return db.RemoveKnotMember(tx, orm.FilterEq("domain", host)) }, + func(tx *sql.Tx, did syntax.DID) error { + return db.AddKnotMember(tx, models.KnotMember{Domain: host, Subject: did}) + }, + func(subject string, perr error) { + r.log.Warn("dropping malformed member DID from reconcile", "host", host, "subject", subject, "error", perr) + }, + ) + }, func() ([]string, error) { rows, err := db.GetKnotMembers(r.store, orm.FilterEq("domain", host)) if err != nil { @@ -80,7 +86,22 @@ func (r *roster) GetKnotMembers(ctx context.Context, host string) ([]string, err func (r *roster) GetRepoCollaborators(ctx context.Context, host, repoDid string) ([]string, error) { return r.serve(ctx, collabScope(repoDid), - func() error { return r.reconcileCollaborators(ctx, host, repoDid) }, + func() error { + repo, perr := syntax.ParseDID(repoDid) + if perr != nil { + return perr + } + return r.reconcile(collabScope(repoDid), + func() ([]string, error) { return r.src.GetRepoCollaborators(ctx, host, repoDid) }, + func(tx *sql.Tx) error { return db.DeleteCollaborator(tx, orm.FilterEq("repo_did", repoDid)) }, + func(tx *sql.Tx, did syntax.DID) error { + return db.AddCollaborator(tx, models.Collaborator{SubjectDid: did, RepoDid: repo}) + }, + func(subject string, perr error) { + r.log.Warn("dropping malformed collaborator DID from reconcile", "repo_did", repoDid, "subject", subject, "error", perr) + }, + ) + }, func() ([]string, error) { rows, err := db.GetCollaborators(r.store, orm.FilterEq("repo_did", repoDid)) if err != nil { @@ -91,78 +112,6 @@ func (r *roster) GetRepoCollaborators(ctx context.Context, host, repoDid string) ) } -func (r *roster) AddKnotMember(host string, subject syntax.DID, cursor Cursor) error { - return r.applyDelta(memberScope(host), subject, cursor, func(tx *sql.Tx) error { - if err := db.RemoveKnotMember(tx, - orm.FilterEq("domain", host), - orm.FilterEq("subject", subject.String()), - ); err != nil { - return err - } - return db.AddKnotMember(tx, models.KnotMember{Domain: host, Subject: subject}) - }) -} - -func (r *roster) RemoveKnotMember(host string, subject syntax.DID, cursor Cursor) error { - return r.applyDelta(memberScope(host), subject, cursor, func(tx *sql.Tx) error { - return db.RemoveKnotMember(tx, - orm.FilterEq("domain", host), - orm.FilterEq("subject", subject.String()), - ) - }) -} - -// NOTE: maybe TODO or not, but no Did of adder means no ability to suggest a vouch for the person they just added as collaborator. -func (r *roster) AddCollaborator(repoDid, subject syntax.DID, cursor Cursor) error { - return r.applyDelta(collabScope(repoDid.String()), subject, cursor, func(tx *sql.Tx) error { - return db.AddCollaborator(tx, models.Collaborator{SubjectDid: subject, RepoDid: repoDid}) - }) -} - -func (r *roster) RemoveCollaborator(repoDid, subject syntax.DID, cursor Cursor) error { - return r.applyDelta(collabScope(repoDid.String()), subject, cursor, func(tx *sql.Tx) error { - return db.DeleteCollaborator(tx, - orm.FilterEq("repo_did", repoDid.String()), - orm.FilterEq("subject_did", subject.String()), - ) - }) -} - -func (r *roster) applyDelta(scope string, subject syntax.DID, cursor Cursor, mutate func(*sql.Tx) error) error { - st := r.acquire(scope) - defer r.release(scope, st) - - st.mu.Lock() - defer st.mu.Unlock() - - tx, err := r.store.Begin() - if err != nil { - return err - } - defer tx.Rollback() - - seen, ok, err := seenCursor(tx, scope, subject) - if err != nil { - return err - } - if ok && cursor <= seen { - return nil - } - - if err := mutate(tx); err != nil { - return err - } - if err := recordCursor(tx, scope, subject, cursor); err != nil { - return err - } - if err := tx.Commit(); err != nil { - return err - } - - st.gen.Add(1) - return nil -} - func (r *roster) InvalidateMembers(host string) { _ = clearSyncedAt(r.store, memberScope(host)) } @@ -215,61 +164,32 @@ func (r *roster) maybeReconcile(key string, reconcile func() error) error { return err } -func (r *roster) reconcileMembers(ctx context.Context, host string) error { - scope := memberScope(host) - st := r.acquire(scope) - defer r.release(scope, st) - - genBefore := st.gen.Load() - subjects, err := r.src.GetKnotMembers(ctx, host) - if err != nil { - r.markFailed(st) - return err - } - return r.commitReconcile(st, scope, genBefore, func(tx *sql.Tx) error { - if err := db.RemoveKnotMember(tx, orm.FilterEq("domain", host)); err != nil { - return err - } - for _, s := range subjects { - did, perr := syntax.ParseDID(s) - if perr != nil { - r.log.Warn("dropping malformed member DID from reconcile", "host", host, "subject", s, "error", perr) - continue - } - if err := db.AddKnotMember(tx, models.KnotMember{Domain: host, Subject: did}); err != nil { - return err - } - } - return nil - }) -} - -func (r *roster) reconcileCollaborators(ctx context.Context, host, repoDid string) error { - repo, perr := syntax.ParseDID(repoDid) - if perr != nil { - return perr - } - scope := collabScope(repoDid) +func (r *roster) reconcile( + scope string, + fetch func() ([]string, error), + drop func(*sql.Tx) error, + add func(*sql.Tx, syntax.DID) error, + warn func(subject string, perr error), +) error { st := r.acquire(scope) defer r.release(scope, st) - genBefore := st.gen.Load() - subjects, err := r.src.GetRepoCollaborators(ctx, host, repoDid) + subjects, err := fetch() if err != nil { r.markFailed(st) return err } - return r.commitReconcile(st, scope, genBefore, func(tx *sql.Tx) error { - if err := db.DeleteCollaborator(tx, orm.FilterEq("repo_did", repoDid)); err != nil { + return r.commitReconcile(st, scope, func(tx *sql.Tx) error { + if err := drop(tx); err != nil { return err } for _, s := range subjects { did, perr := syntax.ParseDID(s) if perr != nil { - r.log.Warn("dropping malformed collaborator DID from reconcile", "repo_did", repoDid, "subject", s, "error", perr) + warn(s, perr) continue } - if err := db.AddCollaborator(tx, models.Collaborator{SubjectDid: did, RepoDid: repo}); err != nil { + if err := add(tx, did); err != nil { return err } } @@ -277,18 +197,10 @@ func (r *roster) reconcileCollaborators(ctx context.Context, host, repoDid strin }) } -func (r *roster) commitReconcile(st *scopeState, scope string, genBefore uint64, replace func(*sql.Tx) error) error { +func (r *roster) commitReconcile(st *scopeState, scope string, replace func(*sql.Tx) error) error { st.mu.Lock() defer st.mu.Unlock() - if st.gen.Load() != genBefore { - if err := setSyncedAt(r.store, scope, r.now()); err != nil { - return err - } - r.clearFailed(st) - return nil - } - tx, err := r.store.Begin() if err != nil { return err @@ -300,9 +212,6 @@ func (r *roster) commitReconcile(st *scopeState, scope string, genBefore uint64, if err := setSyncedAt(tx, scope, r.now()); err != nil { return err } - if err := pruneCursors(tx, scope, int64(cursorRetention)); err != nil { - return err - } if err := tx.Commit(); err != nil { return err } @@ -399,37 +308,3 @@ func clearSyncedAt(e db.Execer, key string) error { _, err := e.Exec(`delete from knotacl_sync where scope_key = ?`, key) return err } - -func seenCursor(e db.Execer, scope string, subject syntax.DID) (Cursor, bool, error) { - var raw int64 - err := e.QueryRow( - `select cursor from knotacl_delta_cursor where scope_key = ? and subject = ?`, - scope, subject.String(), - ).Scan(&raw) - if errors.Is(err, sql.ErrNoRows) { - return 0, false, nil - } - if err != nil { - return 0, false, err - } - return Cursor(raw), true, nil -} - -func recordCursor(e db.Execer, scope string, subject syntax.DID, cursor Cursor) error { - _, err := e.Exec( - `insert into knotacl_delta_cursor (scope_key, subject, cursor) values (?, ?, ?) - on conflict(scope_key, subject) do update set cursor = excluded.cursor`, - scope, subject.String(), int64(cursor), - ) - return err -} - -func pruneCursors(e db.Execer, scope string, retentionNanos int64) error { - _, err := e.Exec( - `delete from knotacl_delta_cursor - where scope_key = ? - and cursor < (select max(cursor) from knotacl_delta_cursor where scope_key = ?) - ?`, - scope, scope, retentionNanos, - ) - return err -} diff --git a/appview/knotacl/roster_test.go b/appview/knotacl/roster_test.go index c74eabae5..7d6aa3ac4 100644 --- a/appview/knotacl/roster_test.go +++ b/appview/knotacl/roster_test.go @@ -9,11 +9,7 @@ import ( "testing" "time" - "github.com/bluesky-social/indigo/atproto/syntax" - "tangled.org/core/appview/db" - "tangled.org/core/appview/models" - "tangled.org/core/orm" ) var rosterTestBase = time.Unix(1700000000, 0) @@ -99,15 +95,6 @@ func rosterTestDB(t *testing.T) *db.DB { const rosterHost = "knot.nel.pet" -func membersForHost(t *testing.T, d *db.DB, host string) []models.KnotMember { - t.Helper() - rows, err := db.GetKnotMembers(d, orm.FilterEq("domain", host)) - if err != nil { - t.Fatalf("GetKnotMembers: %v", err) - } - return rows -} - func TestRoster_BootstrapThenServesFromSqlite(t *testing.T) { clk := &fakeClock{t: rosterTestBase} f := &fakeLister{members: []string{"did:plc:boltless"}} @@ -138,33 +125,6 @@ func TestRoster_BootstrapThenServesFromSqlite(t *testing.T) { } } -func TestRoster_EventDeltaVisibleWithoutXRPC(t *testing.T) { - clk := &fakeClock{t: rosterTestBase} - f := &fakeLister{members: []string{"did:plc:boltless"}} - d := rosterTestDB(t) - r := newRoster(d, f, rosterTTL, clk.now, nil) - ctx := context.Background() - - if _, err := r.GetKnotMembers(ctx, rosterHost); err != nil { - t.Fatal(err) - } - - if err := r.AddKnotMember(rosterHost, syntax.DID("did:plc:akshay"), 1); err != nil { - t.Fatalf("AddKnotMember: %v", err) - } - - got, err := r.GetKnotMembers(ctx, rosterHost) - if err != nil { - t.Fatal(err) - } - if !slices.Equal(got, []string{"did:plc:akshay", "did:plc:boltless"}) { - t.Errorf("post-delta read = %v, want the pushed member reflected", got) - } - if f.calls() != 1 { - t.Errorf("memberCalls=%d, want 1; a pushed delta must not trigger an XRPC reconcile", f.calls()) - } -} - func TestRoster_ColdUnreachableErrors(t *testing.T) { clk := &fakeClock{t: rosterTestBase} f := &fakeLister{err: errors.New("knot unreachable")} @@ -254,177 +214,6 @@ func TestRoster_SingleflightCollapsesConcurrentColdReads(t *testing.T) { } } -func TestRoster_MemberDeltaIdempotentAndScoped(t *testing.T) { - clk := &fakeClock{t: rosterTestBase} - d := rosterTestDB(t) - r := newRoster(d, &fakeLister{}, rosterTTL, clk.now, nil) - - if err := r.AddKnotMember(rosterHost, syntax.DID("did:plc:boltless"), 1); err != nil { - t.Fatalf("AddKnotMember: %v", err) - } - if err := r.AddKnotMember(rosterHost, syntax.DID("did:plc:boltless"), 2); err != nil { - t.Fatalf("AddKnotMember: %v", err) - } - if err := r.RemoveKnotMember("other.nel.pet", syntax.DID("did:plc:boltless"), 1); err != nil { - t.Fatalf("RemoveKnotMember other host: %v", err) - } - - if rows := membersForHost(t, d, rosterHost); len(rows) != 1 { - t.Fatalf("members = %v, want a single row after a duplicate add and an unrelated-host remove", rows) - } - - if err := r.RemoveKnotMember(rosterHost, syntax.DID("did:plc:boltless"), 3); err != nil { - t.Fatalf("RemoveKnotMember: %v", err) - } - if rows := membersForHost(t, d, rosterHost); len(rows) != 0 { - t.Fatalf("members = %v, want empty after remove", rows) - } -} - -func TestRoster_NativeAddReplacesLegacyRow(t *testing.T) { - clk := &fakeClock{t: rosterTestBase} - d := rosterTestDB(t) - r := newRoster(d, &fakeLister{}, rosterTTL, clk.now, nil) - - if err := db.AddKnotMember(d, models.KnotMember{ - Did: syntax.DID("did:plc:akshay"), - Rkey: "legacy-rkey", - Domain: rosterHost, - Subject: syntax.DID("did:plc:boltless"), - }); err != nil { - t.Fatalf("seed legacy row: %v", err) - } - - if err := r.AddKnotMember(rosterHost, syntax.DID("did:plc:boltless"), 1); err != nil { - t.Fatalf("AddKnotMember: %v", err) - } - - rows := membersForHost(t, d, rosterHost) - if len(rows) != 1 { - t.Fatalf("members = %v, want a single canonical row after a native add over a legacy row", rows) - } - if rows[0].Did != "" { - t.Errorf("row did = %q, want empty; the native write must own the (domain, subject) identity", rows[0].Did) - } -} - -func TestRoster_StaleDeltaIgnored(t *testing.T) { - clk := &fakeClock{t: rosterTestBase} - d := rosterTestDB(t) - r := newRoster(d, &fakeLister{}, rosterTTL, clk.now, nil) - - if err := r.AddKnotMember(rosterHost, syntax.DID("did:plc:boltless"), 100); err != nil { - t.Fatalf("AddKnotMember: %v", err) - } - if err := r.RemoveKnotMember(rosterHost, syntax.DID("did:plc:boltless"), 50); err != nil { - t.Fatalf("RemoveKnotMember: %v", err) - } - - if rows := membersForHost(t, d, rosterHost); len(rows) != 1 { - t.Fatalf("members = %v, want the add preserved; a lower-cursor remove arriving late must be ignored", rows) - } -} - -func TestRoster_NewerRemoveWins(t *testing.T) { - clk := &fakeClock{t: rosterTestBase} - d := rosterTestDB(t) - r := newRoster(d, &fakeLister{}, rosterTTL, clk.now, nil) - - if err := r.AddKnotMember(rosterHost, syntax.DID("did:plc:boltless"), 100); err != nil { - t.Fatalf("AddKnotMember: %v", err) - } - if err := r.RemoveKnotMember(rosterHost, syntax.DID("did:plc:boltless"), 101); err != nil { - t.Fatalf("RemoveKnotMember: %v", err) - } - - if rows := membersForHost(t, d, rosterHost); len(rows) != 0 { - t.Fatalf("members = %v, want empty; a higher-cursor remove must win", rows) - } -} - -func TestRoster_LateAddAfterRemoveIgnored(t *testing.T) { - clk := &fakeClock{t: rosterTestBase} - d := rosterTestDB(t) - r := newRoster(d, &fakeLister{}, rosterTTL, clk.now, nil) - - if err := r.RemoveKnotMember(rosterHost, syntax.DID("did:plc:boltless"), 101); err != nil { - t.Fatalf("RemoveKnotMember: %v", err) - } - if err := r.AddKnotMember(rosterHost, syntax.DID("did:plc:boltless"), 100); err != nil { - t.Fatalf("AddKnotMember: %v", err) - } - - if rows := membersForHost(t, d, rosterHost); len(rows) != 0 { - t.Fatalf("members = %v, want empty; an add older than the applied remove must be ignored", rows) - } -} - -func TestRoster_ReconcilePrunesStaleCursors(t *testing.T) { - clk := &fakeClock{t: rosterTestBase} - d := rosterTestDB(t) - r := newRoster(d, &fakeLister{}, rosterTTL, clk.now, nil) - ctx := context.Background() - - stale := Cursor(rosterTestBase.Add(-2 * cursorRetention).UnixNano()) - fresh := Cursor(rosterTestBase.UnixNano()) - - if err := r.AddKnotMember(rosterHost, syntax.DID("did:plc:clam"), stale); err != nil { - t.Fatalf("AddKnotMember stale: %v", err) - } - if err := r.AddKnotMember(rosterHost, syntax.DID("did:plc:whelk"), fresh); err != nil { - t.Fatalf("AddKnotMember fresh: %v", err) - } - - if _, err := r.GetKnotMembers(ctx, rosterHost); err != nil { - t.Fatalf("reconcile: %v", err) - } - - scope := memberScope(rosterHost) - if _, ok, err := seenCursor(d, scope, syntax.DID("did:plc:clam")); err != nil || ok { - t.Errorf("stale cursor present (ok=%v err=%v); reconcile must prune cursors older than the retention window", ok, err) - } - if _, ok, err := seenCursor(d, scope, syntax.DID("did:plc:whelk")); err != nil || !ok { - t.Errorf("fresh cursor missing (ok=%v err=%v); reconcile must keep cursors within the retention window", ok, err) - } -} - -func TestRoster_ReconcileDoesNotClobberConcurrentDelta(t *testing.T) { - clk := &fakeClock{t: rosterTestBase} - f := &fakeLister{members: []string{"did:plc:boltless"}} - r := newRoster(rosterTestDB(t), f, rosterTTL, clk.now, nil) - ctx := context.Background() - - if _, err := r.GetKnotMembers(ctx, rosterHost); err != nil { - t.Fatal(err) - } - - clk.advance(rosterTTL) - - started := make(chan struct{}) - release := make(chan struct{}) - f.arm(started, release) - - done := make(chan []string, 1) - go func() { - got, err := r.GetKnotMembers(ctx, rosterHost) - if err != nil { - t.Errorf("reconcile read: %v", err) - } - done <- got - }() - - <-started - if err := r.AddKnotMember(rosterHost, syntax.DID("did:plc:akshay"), 1); err != nil { - t.Fatalf("AddKnotMember: %v", err) - } - close(release) - - got := <-done - if !slices.Contains(got, "did:plc:akshay") { - t.Errorf("post-reconcile read = %v, want the delta that landed during the drain preserved", got) - } -} - func TestRoster_BackoffSuppressesRepeatedDrains(t *testing.T) { clk := &fakeClock{t: rosterTestBase} f := &fakeLister{members: []string{"did:plc:boltless"}} diff --git a/appview/knotacl/service.go b/appview/knotacl/service.go index d91b9454d..abdb39d84 100644 --- a/appview/knotacl/service.go +++ b/appview/knotacl/service.go @@ -8,7 +8,6 @@ import ( "sync" "time" - "github.com/bluesky-social/indigo/atproto/syntax" "golang.org/x/sync/errgroup" "tangled.org/core/appview/db" @@ -90,22 +89,6 @@ func (s *Service) InvalidateCollaborators(host, repoDid string) { s.nat.client.InvalidateCollaborators(host, repoDid) } -func (s *Service) AddKnotMember(host string, subject syntax.DID, cursor Cursor) error { - return s.nat.client.AddKnotMember(host, subject, cursor) -} - -func (s *Service) RemoveKnotMember(host string, subject syntax.DID, cursor Cursor) error { - return s.nat.client.RemoveKnotMember(host, subject, cursor) -} - -func (s *Service) AddCollaborator(repoDid, subject syntax.DID, cursor Cursor) error { - return s.nat.client.AddCollaborator(repoDid, subject, cursor) -} - -func (s *Service) RemoveCollaborator(repoDid, subject syntax.DID, cursor Cursor) error { - return s.nat.client.RemoveCollaborator(repoDid, subject, cursor) -} - func (s *Service) KnotsForUser(ctx context.Context, userDid string) []string { legacyKnots, err := s.leg.enforcer.GetKnotsForUser(userDid) if err != nil { diff --git a/appview/knots/knots.go b/appview/knots/knots.go index bfdfbbd96..d2ec12255 100644 --- a/appview/knots/knots.go +++ b/appview/knots/knots.go @@ -20,8 +20,8 @@ import ( "tangled.org/core/appview/oauth" "tangled.org/core/appview/pages" "tangled.org/core/appview/serververify" + "tangled.org/core/appview/sitefeed" "tangled.org/core/consts" - "tangled.org/core/eventconsumer" "tangled.org/core/idresolver" "tangled.org/core/orm" "tangled.org/core/rbac" @@ -41,8 +41,8 @@ type Knots struct { Enforcer *rbac.Enforcer Acl *knotacl.Service IdResolver *idresolver.Resolver + SiteFeed *sitefeed.Feed Logger *slog.Logger - Knotstream *eventconsumer.Consumer } func (k *Knots) Router() http.Handler { @@ -224,7 +224,7 @@ func (k *Knots) register(w http.ResponseWriter, r *http.Request) { return } - go k.Knotstream.AddSource(r.Context(), eventconsumer.NewKnotSource(domain)) + k.SiteFeed.Subscribe(context.Background(), domain) k.Pages.HxRefresh(w) } @@ -343,7 +343,7 @@ func (k *Knots) delete(w http.ResponseWriter, r *http.Request) { 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)) + k.SiteFeed.Unsubscribe(domain) } } @@ -459,11 +459,7 @@ func (k *Knots) retry(w http.ResponseWriter, r *http.Request) { } } - // add this knot to knotstream - go k.Knotstream.AddSource( - r.Context(), - eventconsumer.NewKnotSource(domain), - ) + k.SiteFeed.Subscribe(context.Background(), domain) shouldRefresh := r.Header.Get("shouldRefresh") if shouldRefresh == "true" { diff --git a/appview/sitefeed/feed.go b/appview/sitefeed/feed.go new file mode 100644 index 000000000..9b563fd13 --- /dev/null +++ b/appview/sitefeed/feed.go @@ -0,0 +1,325 @@ +package sitefeed + +import ( + "context" + "database/sql" + "errors" + "fmt" + "log/slog" + "strconv" + "sync" + "time" + + "tangled.org/core/appview/cache" + "tangled.org/core/appview/cloudflare" + "tangled.org/core/appview/config" + "tangled.org/core/appview/db" + "tangled.org/core/appview/models" + "tangled.org/core/appview/sites" + "tangled.org/core/hostutil" + "tangled.org/core/knotfeed" + "tangled.org/core/log" + "tangled.org/core/orm" + + "github.com/bluesky-social/indigo/atproto/syntax" + "github.com/go-git/go-git/v5/plumbing" + "github.com/redis/go-redis/v9" +) + +const ( + registrationRetry = 30 * time.Second + maxConcurrentDeploys = 4 +) + +type feedSource struct { + cancel context.CancelFunc +} + +type Feed struct { + d *db.DB + rdb *cache.Cache + cfg *config.Config + cf *cloudflare.Client + logger *slog.Logger + + deploy func(context.Context, *cloudflare.Client, *config.Config, *models.Repo, string, string) error + + deploySlots chan struct{} + + deployLocks sync.Map + pendingSha sync.Map + + mu sync.Mutex + sources map[string]*feedSource +} + +func New(d *db.DB, rdb *cache.Cache, cfg *config.Config, cf *cloudflare.Client, logger *slog.Logger) *Feed { + return &Feed{ + d: d, + rdb: rdb, + cfg: cfg, + cf: cf, + logger: log.SubLogger(logger, "sitefeed"), + deploy: sites.Deploy, + deploySlots: make(chan struct{}, maxConcurrentDeploys), + sources: make(map[string]*feedSource), + } +} + +func (f *Feed) Start(ctx context.Context) { + for { + knots, err := db.GetRegistrations(f.d, orm.FilterIsNot("registered", "null")) + if err == nil { + for _, k := range knots { + f.Subscribe(ctx, k.Domain) + } + return + } + f.logger.Error("failed to list registered knots, retrying", "err", err) + select { + case <-ctx.Done(): + return + case <-time.After(registrationRetry): + } + } +} + +func (f *Feed) Subscribe(ctx context.Context, domain string) { + host, noTLS, err := hostutil.ParseHostname(domain) + if err != nil { + f.logger.Warn("unparseable knot domain, not subscribing", "domain", domain, "err", err) + return + } + + f.mu.Lock() + if _, ok := f.sources[host]; ok { + f.mu.Unlock() + return + } + srcCtx, cancel := context.WithCancel(ctx) + src := &feedSource{cancel: cancel} + f.sources[host] = src + f.mu.Unlock() + + go func() { + defer f.forget(host, src) + f.consumer(host, noTLS).Run(srcCtx) + }() +} + +func (f *Feed) Unsubscribe(domain string) { + host, _, err := hostutil.ParseHostname(domain) + if err != nil { + f.logger.Warn("unparseable knot domain, not unsubscribing", "domain", domain, "err", err) + return + } + + f.mu.Lock() + src, ok := f.sources[host] + delete(f.sources, host) + f.mu.Unlock() + if ok { + f.logger.Info("unsubscribed from knot firehose", "host", host) + src.cancel() + } +} + +func (f *Feed) forget(host string, src *feedSource) { + f.mu.Lock() + defer f.mu.Unlock() + if f.sources[host] == src { + delete(f.sources, host) + } +} + +func (f *Feed) consumer(host string, noTLS bool) *knotfeed.Consumer { + return &knotfeed.Consumer{ + Host: host, + NoTLS: noTLS, + Logger: f.logger, + + LoadCursor: f.loadCursor(host), + StoreCursor: f.storeCursor(host), + Handle: f.handle(host), + OutdatedReplay: func(context.Context) int64 { + f.logger.Warn("site feed cursor is behind the knot, resuming live", "host", host) + return 0 + }, + } +} + +func (f *Feed) loadCursor(host string) func(context.Context) (int64, error) { + return func(ctx context.Context) (int64, error) { + if f.rdb == nil { + return 0, nil + } + raw, err := f.rdb.Get(ctx, f.cursorKey(host)).Result() + if errors.Is(err, redis.Nil) { + return 0, nil + } + if err != nil { + return 0, fmt.Errorf("loading the site feed cursor: %w", err) + } + seq, err := strconv.ParseInt(raw, 10, 64) + if err != nil { + f.logger.Warn("unreadable site feed cursor, resuming live", "host", host, "err", err) + if err := f.rdb.Del(ctx, f.cursorKey(host)).Err(); err != nil { + return 0, fmt.Errorf("clearing an unreadable site feed cursor: %w", err) + } + return 0, nil + } + return seq, nil + } +} + +func (f *Feed) storeCursor(host string) func(context.Context, int64) error { + return func(ctx context.Context, seq int64) error { + if f.rdb == nil { + return nil + } + return f.rdb.Set(ctx, f.cursorKey(host), seq, 0).Err() + } +} + +func (f *Feed) cursorKey(host string) string { + return "sitefeed:cursor:" + host +} + +func (f *Feed) handle(host string) func(context.Context, knotfeed.Message) error { + return func(ctx context.Context, msg knotfeed.Message) error { + if msg.Type != knotfeed.TypeCommit || msg.Commit == nil { + return nil + } + var errs []error + for _, op := range msg.Commit.Records { + if op.Collection != knotfeed.GitRefCollection { + continue + } + errs = append(errs, f.refOp(ctx, host, msg.Commit.Repo, op)) + } + return errors.Join(errs...) + } +} + +func (f *Feed) refOp(ctx context.Context, host, repoDid string, op knotfeed.RecordOp) error { + logger := f.logger.With("knot", host, "repo_did", repoDid) + + refname, ok := knotfeed.UnescapeRkey(op.Rkey) + if !ok { + logger.Warn("undecodable git ref rkey", "rkey", op.Rkey) + return nil + } + + if op.Deleted() { + return nil + } + + record, err := knotfeed.DecodeRefRecord(op.Bytes) + if err != nil { + logger.Warn("undecodable git ref record", "ref", refname, "err", err) + return nil + } + + repo, err := db.GetRepoByDid(f.d, repoDid) + if errors.Is(err, sql.ErrNoRows) { + return nil + } + if err != nil { + return fmt.Errorf("looking up repo %s: %w", repoDid, err) + } + repoKnot, _, err := hostutil.ParseHostname(repo.Knot) + if err != nil { + logger.Warn("repo names an unparseable knot, dropping the record", "repo_knot", repo.Knot) + return nil + } + if repoKnot != host { + logger.Info("dropping a ref record from a knot that doesn't host the repo", "repo_knot", repo.Knot) + return nil + } + + ref := plumbing.ReferenceName(refname) + if ref.IsBranch() { + return f.maybeDeploy(ctx, repo, ref.Short(), record.Sha) + } + return nil +} + +func (f *Feed) maybeDeploy(ctx context.Context, repo *models.Repo, branch, sha string) error { + if f.cf == nil || !f.cf.Enabled() { + return nil + } + siteConfig, err := db.GetRepoSiteConfig(f.d, repo.RepoDid) + if err != nil { + return fmt.Errorf("reading the site config for %s: %w", repo.RepoDid, err) + } + if siteConfig == nil { + return nil + } + if siteConfig.Branch != branch { + return nil + } + f.pendingSha.Store(repo.RepoDid, sha) + select { + case f.deploySlots <- struct{}{}: + default: + f.logger.Warn("deploy queue saturated, dropping the deploy", "repo", repo.RepoDid, "branch", branch) + f.recordDroppedDeploy(repo, siteConfig, sha, "deploy queue saturated") + return nil + } + go func() { + defer func() { <-f.deploySlots }() + mu, _ := f.deployLocks.LoadOrStore(repo.RepoDid, &sync.Mutex{}) + mu.(*sync.Mutex).Lock() + defer mu.(*sync.Mutex).Unlock() + if latest, ok := f.pendingSha.Load(repo.RepoDid); !ok || latest.(string) != sha { + f.logger.Info("superseded push deploy skipped", "repo", repo.RepoDid, "branch", branch, "sha", sha) + return + } + f.triggerDeploy(context.WithoutCancel(ctx), repo, siteConfig, sha) + }() + return nil +} + +func (f *Feed) recordDroppedDeploy(repo *models.Repo, siteConfig *models.RepoSite, sha, reason string) { + deploy := &models.SiteDeploy{ + RepoDid: syntax.DID(repo.RepoDid), + Branch: siteConfig.Branch, + Dir: siteConfig.Dir, + CommitSHA: sha, + Trigger: models.SiteDeployTriggerPush, + Status: models.SiteDeployStatusFailure, + Error: reason, + } + if err := db.AddSiteDeploy(f.d, deploy); err != nil { + f.logger.Error("failed to record the dropped deploy", "repo", repo.RepoDid, "err", err) + } +} + +func (f *Feed) triggerDeploy(ctx context.Context, repo *models.Repo, siteConfig *models.RepoSite, sha string) { + logger := f.logger.With("repo", repo.RepoIdentifier()) + + deploy := &models.SiteDeploy{ + RepoDid: syntax.DID(repo.RepoDid), + Branch: siteConfig.Branch, + Dir: siteConfig.Dir, + CommitSHA: sha, + Trigger: models.SiteDeployTriggerPush, + } + + deployErr := f.deploy(ctx, f.cf, f.cfg, repo, siteConfig.Branch, siteConfig.Dir) + if deployErr != nil { + logger.Error("sites: R2 sync failed on push", "err", deployErr) + deploy.Status = models.SiteDeployStatusFailure + deploy.Error = deployErr.Error() + } else { + deploy.Status = models.SiteDeployStatusSuccess + } + + if err := db.AddSiteDeploy(f.d, deploy); err != nil { + logger.Error("sites: failed to record deploy", "err", err) + } + + if deployErr == nil { + logger.Info("site deployed to r2") + } +} diff --git a/appview/sitefeed/feed_test.go b/appview/sitefeed/feed_test.go new file mode 100644 index 000000000..158ec5d6d --- /dev/null +++ b/appview/sitefeed/feed_test.go @@ -0,0 +1,399 @@ +package sitefeed + +import ( + "bytes" + "context" + "errors" + "path/filepath" + "sync" + "testing" + "time" + + "tangled.org/core/appview/cloudflare" + "tangled.org/core/appview/config" + "tangled.org/core/appview/db" + "tangled.org/core/appview/models" + "tangled.org/core/knotfeed" + "tangled.org/core/log" + + cbg "github.com/whyrusleeping/cbor-gen" +) + +const ( + feedTestHost = "knot.invalid" + feedTestOtherHost = "other.knot.invalid" + feedTestRepoDid = "did:plc:limpet" + feedTestOwner = "did:plc:akshay" +) + +var errBoom = errors.New("boom") + +func feedTestDB(t *testing.T) *db.DB { + t.Helper() + d, err := db.Make(context.Background(), filepath.Join(t.TempDir(), "appview.db")) + if err != nil { + t.Fatalf("db.Make: %v", err) + } + t.Cleanup(func() { d.Close() }) + return d +} + +func seedFeedRepo(t *testing.T, d *db.DB) { + t.Helper() + tx, err := d.Begin() + if err != nil { + t.Fatalf("begin: %v", err) + } + if err := db.AddRepo(tx, &models.Repo{ + Did: feedTestOwner, + Knot: feedTestHost, + RepoDid: feedTestRepoDid, + Name: "anemone", + }); err != nil { + t.Fatalf("AddRepo: %v", err) + } + if err := tx.Commit(); err != nil { + t.Fatalf("commit: %v", err) + } +} + +func setSiteConfig(t *testing.T, d *db.DB) { + t.Helper() + if err := db.SetRepoSiteConfig(d, feedTestRepoDid, "main", "/", false); err != nil { + t.Fatalf("SetRepoSiteConfig: %v", err) + } +} + +func enabledCf(t *testing.T) *cloudflare.Client { + t.Helper() + cf, err := cloudflare.New(&config.Config{Cloudflare: config.Cloudflare{ + AccountId: "acct", + KV: config.KVConfig{NamespaceId: "ns"}, + R2: config.R2Config{Bucket: "bucket"}, + }}) + if err != nil { + t.Fatalf("cloudflare.New: %v", err) + } + return cf +} + +type deployCall struct { + repoDid string + branch string + dir string +} + +type deployRecorder struct { + mu sync.Mutex + calls []deployCall +} + +func (r *deployRecorder) hook(err error) func(context.Context, *cloudflare.Client, *config.Config, *models.Repo, string, string) error { + return func(_ context.Context, _ *cloudflare.Client, _ *config.Config, repo *models.Repo, branch, dir string) error { + r.mu.Lock() + r.calls = append(r.calls, deployCall{repoDid: repo.RepoDid, branch: branch, dir: dir}) + r.mu.Unlock() + return err + } +} + +func (r *deployRecorder) got() []deployCall { + r.mu.Lock() + defer r.mu.Unlock() + return append([]deployCall(nil), r.calls...) +} + +func writeText(w *bytes.Buffer, s string) { + if err := cbg.CborWriteHeader(w, cbg.MajTextString, uint64(len(s))); err != nil { + w.Reset() + } + w.WriteString(s) +} + +func encodeRefRecord(t *testing.T, sha string) []byte { + t.Helper() + var out bytes.Buffer + if err := cbg.CborWriteHeader(&out, cbg.MajMap, 1); err != nil { + t.Fatalf("map header: %v", err) + } + writeText(&out, "sha") + writeText(&out, sha) + return out.Bytes() +} + +func refOpFor(t *testing.T, refname, action, sha string) knotfeed.RecordOp { + t.Helper() + rkey, _ := knotfeed.EscapeRefname(refname) + op := knotfeed.RecordOp{Action: action, Collection: knotfeed.GitRefCollection, Rkey: rkey} + if sha != "" { + op.Bytes = encodeRefRecord(t, sha) + } + return op +} + +func pollFor(t *testing.T, what string, ready func() bool) { + t.Helper() + deadline := time.Now().Add(2 * time.Second) + for time.Now().Before(deadline) { + if ready() { + return + } + time.Sleep(5 * time.Millisecond) + } + t.Fatalf("%s never appeared", what) +} + +func feedFor(t *testing.T, d *db.DB, hook func(context.Context, *cloudflare.Client, *config.Config, *models.Repo, string, string) error) *Feed { + t.Helper() + f := New(d, nil, &config.Config{}, enabledCf(t), log.New("test")) + f.deploy = hook + t.Cleanup(func() { + f.Unsubscribe(feedTestHost) + f.Unsubscribe(feedTestOtherHost) + }) + return f +} + +func TestFeedRefOp_BranchMatchTriggersDeploy(t *testing.T) { + d := feedTestDB(t) + seedFeedRepo(t, d) + setSiteConfig(t, d) + r := &deployRecorder{} + f := feedFor(t, d, r.hook(nil)) + + f.refOp(context.Background(), feedTestHost, feedTestRepoDid, refOpFor(t, "refs/heads/main", "create", "def456")) + + pollFor(t, "the deploy call", func() bool { return len(r.got()) == 1 }) + calls := r.got() + if calls[0] != (deployCall{repoDid: feedTestRepoDid, branch: "main", dir: "/"}) { + t.Fatalf("deploys = %+v, want one push on main", calls) + } + + var deploy *models.SiteDeploy + pollFor(t, "the deploy row", func() bool { + rows, err := db.GetSiteDeploys(d, feedTestRepoDid, 10) + if err == nil && len(rows) > 0 { + deploy = &rows[0] + return true + } + return false + }) + if deploy.Status != models.SiteDeployStatusSuccess { + t.Errorf("status = %q, want success", deploy.Status) + } + if deploy.Trigger != models.SiteDeployTriggerPush { + t.Errorf("trigger = %q, want push", deploy.Trigger) + } + if deploy.CommitSHA != "def456" || deploy.Branch != "main" || deploy.Dir != "/" { + t.Errorf("deploy row = %+v, want the pushed sha, branch and dir", deploy) + } +} + +func TestFeedSkipsNonDeployingOps(t *testing.T) { + seedSiteRepo := func(t *testing.T, d *db.DB) { + seedFeedRepo(t, d) + setSiteConfig(t, d) + } + for _, tc := range []struct { + name string + seed func(*testing.T, *db.DB) + op func(*testing.T, *Feed) + }{ + { + name: "push to a branch other than the site branch", + seed: seedSiteRepo, + op: func(_ *testing.T, f *Feed) { + f.maybeDeploy(context.Background(), &models.Repo{RepoDid: feedTestRepoDid}, "other", "abc123") + }, + }, + { + name: "push without a site config", + seed: seedFeedRepo, + op: func(_ *testing.T, f *Feed) { + f.maybeDeploy(context.Background(), &models.Repo{RepoDid: feedTestRepoDid}, "main", "abc123") + }, + }, + { + name: "tag push", + seed: seedSiteRepo, + op: func(t *testing.T, f *Feed) { + f.refOp(context.Background(), feedTestHost, feedTestRepoDid, refOpFor(t, "refs/tags/v1", "create", "def456")) + }, + }, + { + name: "deleted ref", + seed: seedSiteRepo, + op: func(t *testing.T, f *Feed) { + f.refOp(context.Background(), feedTestHost, feedTestRepoDid, refOpFor(t, "refs/heads/main", "delete", "")) + }, + }, + { + name: "record from a knot that doesn't host the repo", + seed: seedSiteRepo, + op: func(t *testing.T, f *Feed) { + f.refOp(context.Background(), feedTestOtherHost, feedTestRepoDid, refOpFor(t, "refs/heads/main", "create", "def456")) + }, + }, + { + name: "repo missing from the index", + op: func(t *testing.T, f *Feed) { + f.refOp(context.Background(), feedTestHost, "did:plc:unknown", refOpFor(t, "refs/heads/main", "create", "def456")) + }, + }, + { + name: "garbage rkey", + seed: seedSiteRepo, + op: func(_ *testing.T, f *Feed) { + f.refOp(context.Background(), feedTestHost, feedTestRepoDid, knotfeed.RecordOp{Action: "create", Collection: knotfeed.GitRefCollection, Rkey: "~zz", Bytes: []byte{0xff}}) + }, + }, + { + name: "garbage record bytes", + seed: seedSiteRepo, + op: func(_ *testing.T, f *Feed) { + f.refOp(context.Background(), feedTestHost, feedTestRepoDid, knotfeed.RecordOp{Action: "create", Collection: knotfeed.GitRefCollection, Rkey: "refs~2fheads~2fmain", Bytes: []byte{0xff}}) + }, + }, + { + name: "commit with only foreign collections", + seed: seedFeedRepo, + op: func(t *testing.T, f *Feed) { + msg := knotfeed.Message{Type: knotfeed.TypeCommit, Commit: &knotfeed.Commit{ + Repo: feedTestRepoDid, + Seq: 7, + Records: []knotfeed.RecordOp{ + {Action: "create", Collection: "sh.tangled.repo", Rkey: "abc", Bytes: []byte("x")}, + }, + }} + if err := f.handle(feedTestHost)(context.Background(), msg); err != nil { + t.Errorf("handle: %v", err) + } + }, + }, + } { + t.Run(tc.name, func(t *testing.T) { + d := feedTestDB(t) + if tc.seed != nil { + tc.seed(t, d) + } + r := &deployRecorder{} + f := feedFor(t, d, r.hook(nil)) + tc.op(t, f) + if calls := r.got(); len(calls) != 0 { + t.Errorf("deploys = %v, want none", calls) + } + }) + } +} + +func TestFeedTrigger_RecordsDeployFailure(t *testing.T) { + d := feedTestDB(t) + seedFeedRepo(t, d) + setSiteConfig(t, d) + r := &deployRecorder{} + f := feedFor(t, d, r.hook(errBoom)) + + f.maybeDeploy(context.Background(), &models.Repo{RepoDid: feedTestRepoDid, Knot: feedTestHost}, "main", "abc123") + + pollFor(t, "the deploy call", func() bool { return len(r.got()) == 1 }) + var deploy *models.SiteDeploy + pollFor(t, "the deploy row", func() bool { + rows, err := db.GetSiteDeploys(d, feedTestRepoDid, 10) + if err == nil && len(rows) > 0 { + deploy = &rows[0] + return true + } + return false + }) + if deploy.Status != models.SiteDeployStatusFailure { + t.Errorf("status = %q, want failure", deploy.Status) + } + if deploy.Error != "boom" { + t.Errorf("error = %q, want the sync failure text", deploy.Error) + } +} + +func TestFeedHandleFailsTheCommitWhenTheDatabaseIsDown(t *testing.T) { + d := feedTestDB(t) + seedFeedRepo(t, d) + setSiteConfig(t, d) + r := &deployRecorder{} + f := feedFor(t, d, r.hook(nil)) + + if err := d.Close(); err != nil { + t.Fatalf("Close: %v", err) + } + + msg := knotfeed.Message{Type: knotfeed.TypeCommit, Commit: &knotfeed.Commit{ + Repo: feedTestRepoDid, + Seq: 9, + Records: []knotfeed.RecordOp{ + refOpFor(t, "refs/heads/main", "create", "def456"), + }, + }} + if err := f.handle(feedTestHost)(context.Background(), msg); err == nil { + t.Fatal("handle accepted a commit while the database was down; the cursor would advance past a dropped deploy") + } + if calls := r.got(); len(calls) != 0 { + t.Errorf("deploys = %v, want none while the database is down", calls) + } +} + +func TestFeedSubscribeDeduplicatesAndUnsubscribes(t *testing.T) { + d := feedTestDB(t) + f := feedFor(t, d, nil) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + f.Subscribe(ctx, feedTestHost) + f.Subscribe(ctx, feedTestHost) + + f.mu.Lock() + count := len(f.sources) + f.mu.Unlock() + if count != 1 { + t.Fatalf("sources = %d, want 1 after a duplicate subscribe", count) + } + + f.Unsubscribe(feedTestHost) + f.mu.Lock() + count = len(f.sources) + f.mu.Unlock() + if count != 0 { + t.Fatalf("sources = %d, want 0 after unsubscribe", count) + } + + f.Unsubscribe(feedTestOtherHost) + f.mu.Lock() + count = len(f.sources) + f.mu.Unlock() + if count != 0 { + t.Fatalf("sources = %d, want 0 after unsubscribing a host with no source", count) + } +} + +func TestFeedTrigger_LastPushWinsOverARacingDeploy(t *testing.T) { + d := feedTestDB(t) + seedFeedRepo(t, d) + setSiteConfig(t, d) + r := &deployRecorder{} + f := feedFor(t, d, r.hook(nil)) + repo := &models.Repo{RepoDid: feedTestRepoDid, Knot: feedTestHost} + + f.maybeDeploy(context.Background(), repo, "main", "shaA") + f.maybeDeploy(context.Background(), repo, "main", "shaB") + + var winner models.SiteDeploy + pollFor(t, "a successful shaB deploy row", func() bool { + rows, err := db.GetSiteDeploys(d, feedTestRepoDid, 10) + if err != nil { + t.Fatalf("GetSiteDeploys: %v", err) + } + for _, row := range rows { + if row.Id > winner.Id { + winner = row + } + } + return winner.CommitSHA == "shaB" && winner.Status == models.SiteDeployStatusSuccess + }) +} diff --git a/appview/state/knotstream.go b/appview/state/knotstream.go deleted file mode 100644 index c2ad968c2..000000000 --- a/appview/state/knotstream.go +++ /dev/null @@ -1,477 +0,0 @@ -package state - -import ( - "context" - "database/sql" - "encoding/json" - "errors" - "fmt" - "slices" - "strings" - "time" - - "tangled.org/core/appview/cloudflare" - "tangled.org/core/appview/notify" - - "tangled.org/core/api/tangled" - "tangled.org/core/appview/config" - "tangled.org/core/appview/db" - "tangled.org/core/appview/knotacl" - "tangled.org/core/appview/knotcompat" - "tangled.org/core/appview/models" - "tangled.org/core/appview/sites" - "tangled.org/core/consts" - ec "tangled.org/core/eventconsumer" - "tangled.org/core/eventstream" - knotdb "tangled.org/core/knotserver/db" - "tangled.org/core/log" - "tangled.org/core/orm" - "tangled.org/core/rbac" - - "github.com/bluesky-social/indigo/atproto/syntax" - "github.com/go-git/go-git/v5/plumbing" - "github.com/posthog/posthog-go" -) - -type aclRoster interface { - AddKnotMember(host string, subject syntax.DID, cursor knotacl.Cursor) error - RemoveKnotMember(host string, subject syntax.DID, cursor knotacl.Cursor) error - AddCollaborator(repoDid, subject syntax.DID, cursor knotacl.Cursor) error - RemoveCollaborator(repoDid, subject syntax.DID, cursor knotacl.Cursor) error - InvalidateMembers(host string) - InvalidateCollaborators(host, repoDid string) -} - -func Knotstream(ctx context.Context, c *config.Config, d *db.DB, acl *knotacl.Service, enforcer *rbac.Enforcer, posthog posthog.Client, notifier notify.Notifier, cfClient *cloudflare.Client) (*ec.Consumer, error) { - knots, err := db.GetRegistrations(d, orm.FilterIsNot("registered", "null")) - if err != nil { - return nil, err - } - - hosts := make([]string, len(knots)) - for i, k := range knots { - hosts[i] = k.Domain - } - - return bootstrapStream( - ctx, "knotstream", ec.KindKnot, hosts, c.Redis.Addr, - c.Knotstream, - knotIngester(d, acl, enforcer, posthog, notifier, c.Core.Dev, c, cfClient), - ), nil -} - -func resolveRepo(d *db.DB, repoDid *string, ownerDid, repoName string) (*models.Repo, error) { - if repoDid != nil && *repoDid != "" { - return db.GetRepoByDid(d, *repoDid) - } - repos, err := db.GetRepos(d, orm.FilterEq("did", ownerDid), orm.FilterEq("rkey", strings.ToLower(repoName))) - if err != nil { - return nil, err - } - if len(repos) == 0 { - return nil, sql.ErrNoRows - } - return &repos[0], nil -} - -func knotIngester(d *db.DB, acl aclRoster, 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 eventstream.Event) error { - switch msg.Nsid { - case tangled.GitRefUpdateNSID: - return ingestRefUpdate(ctx, d, enforcer, posthog, notifier, dev, c, cfClient, source, msg) - case knotdb.RepoDIDAssignNSID: - return ingestDIDAssign(d, enforcer, source, msg, ctx) - case knotdb.KnotMemberUpdateNSID: - return ingestKnotMemberUpdate(acl, source, msg) - case knotdb.RepoCollaboratorUpdateNSID: - return ingestCollaboratorUpdate(ctx, d, acl, source, msg) - } - - return nil - } -} - -const ( - aclIngestAttempts = 3 - aclIngestBackoff = 50 * time.Millisecond -) - -func withAclRetry(attempts int, backoff time.Duration, op func() error) error { - if err := op(); err == nil || attempts <= 1 { - return err - } - time.Sleep(backoff) - return withAclRetry(attempts-1, backoff, op) -} - -func ingestKnotMemberUpdate(acl aclRoster, source ec.Source, msg eventstream.Event) error { - var rec knotdb.KnotMemberUpdate - if err := json.Unmarshal(msg.EventJson, &rec); err != nil { - return fmt.Errorf("unmarshal memberUpdate: %w", err) - } - - subject, err := syntax.ParseDID(rec.Subject) - if err != nil { - return fmt.Errorf("memberUpdate bad subject %q: %w", rec.Subject, err) - } - - cursor := knotacl.Cursor(msg.Created) - switch rec.Op { - case knotdb.AclOpAdd: - err = withAclRetry(aclIngestAttempts, aclIngestBackoff, func() error { - return acl.AddKnotMember(source.Host, subject, cursor) - }) - case knotdb.AclOpRemove: - err = withAclRetry(aclIngestAttempts, aclIngestBackoff, func() error { - return acl.RemoveKnotMember(source.Host, subject, cursor) - }) - default: - return fmt.Errorf("memberUpdate unknown op %q", rec.Op) - } - - if err != nil { - acl.InvalidateMembers(source.Host) - } - return err -} - -func ingestCollaboratorUpdate(ctx context.Context, d *db.DB, acl aclRoster, source ec.Source, msg eventstream.Event) error { - var rec knotdb.RepoCollaboratorUpdate - if err := json.Unmarshal(msg.EventJson, &rec); err != nil { - return fmt.Errorf("unmarshal collaboratorUpdate: %w", err) - } - - subject, err := syntax.ParseDID(rec.Subject) - if err != nil { - return fmt.Errorf("collaboratorUpdate bad subject %q: %w", rec.Subject, err) - } - repoDid, err := syntax.ParseDID(rec.Repo) - if err != nil { - return fmt.Errorf("collaboratorUpdate bad repo %q: %w", rec.Repo, err) - } - - cursor := knotacl.Cursor(msg.Created) - switch rec.Op { - case knotdb.AclOpAdd: - err = withAclRetry(aclIngestAttempts, aclIngestBackoff, func() error { - owned, err := repoOwnedBySource(ctx, d, source, repoDid, subject) - if err != nil || !owned { - return err - } - return acl.AddCollaborator(repoDid, subject, cursor) - }) - case knotdb.AclOpRemove: - err = withAclRetry(aclIngestAttempts, aclIngestBackoff, func() error { - owned, err := repoOwnedBySource(ctx, d, source, repoDid, subject) - if err != nil || !owned { - return err - } - return acl.RemoveCollaborator(repoDid, subject, cursor) - }) - default: - return fmt.Errorf("collaboratorUpdate unknown op %q", rec.Op) - } - - if err != nil { - acl.InvalidateCollaborators(source.Host, repoDid.String()) - } - return err -} - -func repoOwnedBySource(ctx context.Context, d *db.DB, source ec.Source, repoDid, subject syntax.DID) (bool, error) { - repo, err := db.GetRepoByDid(d, repoDid.String()) - if errors.Is(err, sql.ErrNoRows) { - log.FromContext(ctx).Warn("collaboratorUpdate for unindexed repo, skipping until reconcile", - "repo_did", repoDid, "subject", subject) - return false, nil - } - if err != nil { - return false, err - } - if repo.Knot != source.Host { - log.FromContext(ctx).Warn("collaboratorUpdate for a repo this knot does not host, dropping", - "repo_did", repoDid, "subject", subject, "claimed_by", source.Host, "owner", repo.Knot) - return false, nil - } - return true, nil -} - -// 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 eventstream.Event) error { - logger := log.FromContext(ctx) - - var record tangled.GitRefUpdate - err := json.Unmarshal(msg.EventJson, &record) - if err != nil { - return err - } - - if !knotcompat.KnotHasCapability(ctx, source.Host, dev, consts.CapKnotACL) { - knownKnots, err := enforcer.GetKnotsForUser(record.CommitterDid) - switch { - case err != nil: - logger.Warn("gitRefUpdate membership lookup failed, ingesting without the sanity check", "committer", record.CommitterDid, "knot", source.Host, "err", err) - case !slices.Contains(knownKnots, source.Host): - logger.Warn("gitRefUpdate committer is not a known member of the knot, ingesting anyway", "committer", record.CommitterDid, "knot", source.Host) - } - } - - if record.Repo == "" { - return fmt.Errorf("gitRefUpdate from %s missing repo", source.Host) - } - - repo, lookupErr := db.GetRepoByDid(d, record.Repo) - if lookupErr != nil { - return fmt.Errorf("failed to look up repo: %w", lookupErr) - } - - logger.Info("processing gitRefUpdate event", - "repo", repo.RepoIdentifier(), - "ref", record.Ref, - "old_sha", record.OldSha, - "new_sha", record.NewSha) - - notifier.Push(ctx, repo, record.Ref, record.OldSha, record.NewSha, record.CommitterDid) - - errPunchcard := populatePunchcard(d, record) - errLanguages := updateRepoLanguages(d, record) - - var errPosthog error - if !dev && record.CommitterDid != "" { - errPosthog = pc.Enqueue(posthog.Capture{ - DistinctId: record.CommitterDid, - Event: "git_ref_update", - }) - } - - // Trigger a sites redeploy if this push is to the configured sites branch. - if cfClient.Enabled() { - go triggerSitesDeployIfNeeded(ctx, d, cfClient, c, record, source) - } - - return errors.Join(errPunchcard, errLanguages, errPosthog) -} - -// triggerSitesDeployIfNeeded checks whether the pushed ref matches the sites -// branch configured for this repo and, if so, syncs the site to R2 -func triggerSitesDeployIfNeeded(ctx context.Context, d *db.DB, cfClient *cloudflare.Client, cfg *config.Config, record tangled.GitRefUpdate, source ec.Source) { - logger := log.FromContext(ctx) - - ref := plumbing.ReferenceName(record.Ref) - if !ref.IsBranch() { - return - } - pushedBranch := ref.Short() - - repo, err := db.GetRepoByDid(d, record.Repo) - if err != nil { - return - } - - siteConfig, err := db.GetRepoSiteConfig(d, repo.RepoDid) - if err != nil || siteConfig == nil { - return - } - if siteConfig.Branch != pushedBranch { - return - } - - deploy := &models.SiteDeploy{ - RepoDid: syntax.DID(repo.RepoDid), - Branch: siteConfig.Branch, - Dir: siteConfig.Dir, - CommitSHA: record.NewSha, - Trigger: models.SiteDeployTriggerPush, - } - - deployErr := sites.Deploy(ctx, cfClient, cfg, repo, siteConfig.Branch, siteConfig.Dir) - if deployErr != nil { - logger.Error("sites: R2 sync failed on push", "repo", repo.RepoIdentifier(), "err", deployErr) - deploy.Status = models.SiteDeployStatusFailure - deploy.Error = deployErr.Error() - } else { - deploy.Status = models.SiteDeployStatusSuccess - } - - if err := db.AddSiteDeploy(d, deploy); err != nil { - logger.Error("sites: failed to record deploy", "repo", repo.RepoIdentifier(), "err", err) - } - - if deployErr == nil { - logger.Info("site deployed to r2", "repo", repo.RepoIdentifier()) - } -} - -func populatePunchcard(d *db.DB, record tangled.GitRefUpdate) error { - if record.CommitterDid == "" { - return nil - } - - knownEmails, err := db.GetAllEmails(d, record.CommitterDid) - if err != nil { - return err - } - - count := 0 - for _, ke := range knownEmails { - if record.Meta == nil { - continue - } - if record.Meta.CommitCount == nil { - continue - } - for _, ce := range record.Meta.CommitCount.ByEmail { - if ce == nil { - continue - } - if ce.Email == ke.Address || ce.Email == record.CommitterDid { - count += int(ce.Count) - } - } - } - - punch := models.Punch{ - Did: record.CommitterDid, - Date: time.Now(), - Count: count, - } - return db.AddPunch(d, punch) -} - -func updateRepoLanguages(d *db.DB, record tangled.GitRefUpdate) error { - if record.Meta == nil || record.Meta.LangBreakdown == nil || record.Meta.LangBreakdown.Inputs == nil { - return fmt.Errorf("empty language data for repo: %s", record.Repo) - } - - r, lookupErr := db.GetRepoByDid(d, record.Repo) - if lookupErr != nil { - return fmt.Errorf("failed to look up repo: %w", lookupErr) - } - repo := *r - - ref := plumbing.ReferenceName(record.Ref) - if !ref.IsBranch() { - return fmt.Errorf("%s is not a valid reference name", ref) - } - - var langs []models.RepoLanguage - for _, l := range record.Meta.LangBreakdown.Inputs { - if l == nil { - continue - } - - langs = append(langs, models.RepoLanguage{ - RepoDid: syntax.DID(repo.RepoDid), - Ref: ref.Short(), - IsDefaultRef: record.Meta.IsDefaultRef, - Language: l.Lang, - Bytes: l.Size, - }) - } - - tx, err := d.Begin() - if err != nil { - return err - } - defer tx.Rollback() - - // update appview's cache - err = db.UpdateRepoLanguages(tx, syntax.DID(repo.RepoDid), ref.Short(), langs) - if err != nil { - fmt.Printf("failed; %s\n", err) - // non-fatal - } - - return tx.Commit() -} - -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 - if err := json.Unmarshal(msg.EventJson, &record); err != nil { - return fmt.Errorf("unmarshal didAssign: %w", err) - } - - if record.RepoDid == "" || record.OwnerDid == "" || record.RepoName == "" { - return fmt.Errorf("didAssign missing required fields: repoDid=%q ownerDid=%q repoName=%q", - record.RepoDid, record.OwnerDid, record.RepoName) - } - - logger.Info("processing didAssign event", - "repo_did", record.RepoDid, - "owner_did", record.OwnerDid, - "repo_name", record.RepoName) - - repos, err := db.GetRepos(d, - orm.FilterEq("did", record.OwnerDid), - orm.FilterEq("name", record.RepoName), - ) - if err != nil || len(repos) == 0 { - logger.Warn("didAssign for unknown repo, skipping", - "owner_did", record.OwnerDid, - "repo_name", record.RepoName) - return nil - } - repo := repos[0] - knot := source.Host - - if repo.Knot != knot { - return fmt.Errorf("didAssign from %s for repo hosted on %s, rejecting", knot, repo.Knot) - } - - repoAtUri := repo.RepoAt().String() - legacyResource := record.OwnerDid + "/" + record.RepoName - - if repo.RepoDid != record.RepoDid { - tx, err := d.Begin() - if err != nil { - return fmt.Errorf("begin didAssign txn: %w", err) - } - defer tx.Rollback() - - if err := db.CascadeRepoDid(tx, repoAtUri, record.RepoDid); err != nil { - return fmt.Errorf("cascade repo_did: %w", err) - } - - if err := db.EnqueuePdsRewritesForRepo(tx, record.RepoDid, repoAtUri); err != nil { - return fmt.Errorf("enqueue pds rewrites: %w", err) - } - - if err := tx.Commit(); err != nil { - return fmt.Errorf("commit didAssign txn: %w", err) - } - } - - if err := enforcer.RemoveRepo(record.OwnerDid, knot, legacyResource); err != nil { - return fmt.Errorf("remove legacy RBAC policies for %s: %w", legacyResource, err) - } - if err := enforcer.AddRepo(record.OwnerDid, knot, record.RepoDid); err != nil { - return fmt.Errorf("add RBAC policies for %s: %w", record.RepoDid, err) - } - - collabs, collabErr := db.GetCollaborators(d, orm.FilterEq("repo_did", record.RepoDid)) - if collabErr != nil { - return fmt.Errorf("get collaborators for RBAC update: %w", collabErr) - } - for _, c := range collabs { - collabDid := c.SubjectDid.String() - if err := enforcer.RemoveCollaborator(collabDid, knot, legacyResource); err != nil { - return fmt.Errorf("remove collaborator RBAC for %s: %w", collabDid, err) - } - if err := enforcer.AddCollaborator(collabDid, knot, record.RepoDid); err != nil { - return fmt.Errorf("add collaborator RBAC for %s: %w", collabDid, err) - } - } - - if err := enforcer.E.SavePolicy(); err != nil { - return fmt.Errorf("save RBAC policies after didAssign: %w", err) - } - - logger.Info("didAssign processed successfully", - "repo_did", record.RepoDid, - "owner_did", record.OwnerDid, - "repo_name", record.RepoName) - - return nil -} diff --git a/appview/state/knotstream_test.go b/appview/state/knotstream_test.go deleted file mode 100644 index f00938dd2..000000000 --- a/appview/state/knotstream_test.go +++ /dev/null @@ -1,333 +0,0 @@ -package state - -import ( - "context" - "encoding/json" - "errors" - "path/filepath" - "testing" - - "github.com/bluesky-social/indigo/atproto/syntax" - - "tangled.org/core/appview/db" - "tangled.org/core/appview/knotacl" - "tangled.org/core/appview/models" - ec "tangled.org/core/eventconsumer" - "tangled.org/core/eventstream" - knotdb "tangled.org/core/knotserver/db" -) - -const ( - aclTestHost = "knot.nel.pet" - aclTestRepoDid = "did:plc:limpet" - aclTestOwner = "did:plc:akshay" - aclTestSubject = "did:plc:boltless" -) - -type memberCall struct { - host string - subject string -} - -type collabCall struct { - repoDid string - subject string -} - -type recordingAcl struct { - memberAdd []memberCall - memberRemove []memberCall - collabAdd []collabCall - collabRemove []collabCall - membersInvalid []string - collabsInvalid []collabCall -} - -func (r *recordingAcl) AddKnotMember(host string, subject syntax.DID, cursor knotacl.Cursor) error { - r.memberAdd = append(r.memberAdd, memberCall{host, subject.String()}) - return nil -} - -func (r *recordingAcl) RemoveKnotMember(host string, subject syntax.DID, cursor knotacl.Cursor) error { - r.memberRemove = append(r.memberRemove, memberCall{host, subject.String()}) - return nil -} - -func (r *recordingAcl) AddCollaborator(repoDid, subject syntax.DID, cursor knotacl.Cursor) error { - r.collabAdd = append(r.collabAdd, collabCall{repoDid.String(), subject.String()}) - return nil -} - -func (r *recordingAcl) RemoveCollaborator(repoDid, subject syntax.DID, cursor knotacl.Cursor) error { - r.collabRemove = append(r.collabRemove, collabCall{repoDid.String(), subject.String()}) - return nil -} - -func (r *recordingAcl) InvalidateMembers(host string) { - r.membersInvalid = append(r.membersInvalid, host) -} - -func (r *recordingAcl) InvalidateCollaborators(host, repoDid string) { - r.collabsInvalid = append(r.collabsInvalid, collabCall{repoDid, host}) -} - -type flakyAcl struct { - failsLeft int - calls int - membersInvalid int - collabsInvalid int -} - -func (a *flakyAcl) try() error { - a.calls++ - if a.failsLeft > 0 { - a.failsLeft-- - return errors.New("transient store error") - } - return nil -} - -func (a *flakyAcl) AddKnotMember(host string, subject syntax.DID, cursor knotacl.Cursor) error { - return a.try() -} -func (a *flakyAcl) RemoveKnotMember(host string, subject syntax.DID, cursor knotacl.Cursor) error { - return a.try() -} -func (a *flakyAcl) AddCollaborator(repoDid, subject syntax.DID, cursor knotacl.Cursor) error { - return a.try() -} -func (a *flakyAcl) RemoveCollaborator(repoDid, subject syntax.DID, cursor knotacl.Cursor) error { - return a.try() -} -func (a *flakyAcl) InvalidateMembers(host string) { a.membersInvalid++ } -func (a *flakyAcl) InvalidateCollaborators(host, repo string) { a.collabsInvalid++ } - -func aclTestDB(t *testing.T) *db.DB { - t.Helper() - d, err := db.Make(context.Background(), filepath.Join(t.TempDir(), "appview.db")) - if err != nil { - t.Fatalf("db.Make: %v", err) - } - t.Cleanup(func() { d.Close() }) - return d -} - -func seedAclRepo(t *testing.T, d *db.DB) { - t.Helper() - tx, err := d.Begin() - if err != nil { - t.Fatalf("begin: %v", err) - } - if err := db.AddRepo(tx, &models.Repo{ - Did: aclTestOwner, - Knot: aclTestHost, - RepoDid: aclTestRepoDid, - Name: "anemone", - }); err != nil { - t.Fatalf("AddRepo: %v", err) - } - if err := tx.Commit(); err != nil { - t.Fatalf("commit: %v", err) - } -} - -func memberEvent(t *testing.T, op knotdb.AclOp, subject string) eventstream.Event { - t.Helper() - payload, err := json.Marshal(knotdb.KnotMemberUpdate{Op: op, Subject: subject}) - if err != nil { - t.Fatalf("marshal memberUpdate: %v", err) - } - return eventstream.Event{Rkey: "evt", Nsid: knotdb.KnotMemberUpdateNSID, EventJson: payload} -} - -func collabEvent(t *testing.T, op knotdb.AclOp, subject, repoDid string) eventstream.Event { - t.Helper() - payload, err := json.Marshal(knotdb.RepoCollaboratorUpdate{Op: op, Subject: subject, Repo: repoDid}) - if err != nil { - t.Fatalf("marshal collaboratorUpdate: %v", err) - } - return eventstream.Event{Rkey: "evt", Nsid: knotdb.RepoCollaboratorUpdateNSID, EventJson: payload} -} - -func TestIngestKnotMemberUpdate_DispatchesAddThenRemove(t *testing.T) { - acl := &recordingAcl{} - source := ec.Source{Kind: ec.KindKnot, Host: aclTestHost} - - if err := ingestKnotMemberUpdate(acl, source, memberEvent(t, knotdb.AclOpAdd, aclTestSubject)); err != nil { - t.Fatalf("add: %v", err) - } - if err := ingestKnotMemberUpdate(acl, source, memberEvent(t, knotdb.AclOpRemove, aclTestSubject)); err != nil { - t.Fatalf("remove: %v", err) - } - - if len(acl.memberAdd) != 1 || acl.memberAdd[0] != (memberCall{aclTestHost, aclTestSubject}) { - t.Errorf("memberAdd = %v, want one add scoped to the source host", acl.memberAdd) - } - if len(acl.memberRemove) != 1 || acl.memberRemove[0] != (memberCall{aclTestHost, aclTestSubject}) { - t.Errorf("memberRemove = %v, want one remove scoped to the source host", acl.memberRemove) - } -} - -func TestIngestKnotMemberUpdate_UnknownOpErrors(t *testing.T) { - acl := &recordingAcl{} - source := ec.Source{Kind: ec.KindKnot, Host: aclTestHost} - if err := ingestKnotMemberUpdate(acl, source, memberEvent(t, knotdb.AclOp("bogus"), aclTestSubject)); err == nil { - t.Fatal("an unknown op must be rejected") - } - if len(acl.memberAdd)+len(acl.memberRemove) != 0 { - t.Errorf("an unknown op must not reach the roster: %+v", acl) - } -} - -func TestIngestKnotMemberUpdate_BadSubjectErrors(t *testing.T) { - acl := &recordingAcl{} - source := ec.Source{Kind: ec.KindKnot, Host: aclTestHost} - if err := ingestKnotMemberUpdate(acl, source, memberEvent(t, knotdb.AclOpAdd, "not-a-did")); err == nil { - t.Fatal("a malformed subject DID must be rejected") - } - if len(acl.memberAdd) != 0 { - t.Errorf("a malformed subject must not reach the roster: %v", acl.memberAdd) - } -} - -func TestIngestCollaboratorUpdate_DispatchesAddThenRemove(t *testing.T) { - ctx := context.Background() - d := aclTestDB(t) - seedAclRepo(t, d) - acl := &recordingAcl{} - source := ec.Source{Kind: ec.KindKnot, Host: aclTestHost} - - if err := ingestCollaboratorUpdate(ctx, d, acl, source, collabEvent(t, knotdb.AclOpAdd, aclTestSubject, aclTestRepoDid)); err != nil { - t.Fatalf("add: %v", err) - } - if err := ingestCollaboratorUpdate(ctx, d, acl, source, collabEvent(t, knotdb.AclOpRemove, aclTestSubject, aclTestRepoDid)); err != nil { - t.Fatalf("remove: %v", err) - } - - if len(acl.collabAdd) != 1 || acl.collabAdd[0] != (collabCall{aclTestRepoDid, aclTestSubject}) { - t.Errorf("collabAdd = %v, want one add for the repo", acl.collabAdd) - } - if len(acl.collabRemove) != 1 || acl.collabRemove[0] != (collabCall{aclTestRepoDid, aclTestSubject}) { - t.Errorf("collabRemove = %v, want one remove for the repo", acl.collabRemove) - } -} - -func TestIngestCollaboratorUpdate_UnindexedRepoSkips(t *testing.T) { - ctx := context.Background() - d := aclTestDB(t) - acl := &recordingAcl{} - source := ec.Source{Kind: ec.KindKnot, Host: aclTestHost} - - if err := ingestCollaboratorUpdate(ctx, d, acl, source, collabEvent(t, knotdb.AclOpAdd, aclTestSubject, aclTestRepoDid)); err != nil { - t.Fatalf("add for unindexed repo must not error, got: %v", err) - } - if len(acl.collabAdd) != 0 { - t.Errorf("an add for an unindexed repo must not reach the roster: %v", acl.collabAdd) - } -} - -func TestIngestCollaboratorUpdate_ForeignKnotDropped(t *testing.T) { - ctx := context.Background() - d := aclTestDB(t) - seedAclRepo(t, d) - acl := &recordingAcl{} - source := ec.Source{Kind: ec.KindKnot, Host: "barnacle.nel.pet"} - - if err := ingestCollaboratorUpdate(ctx, d, acl, source, collabEvent(t, knotdb.AclOpAdd, aclTestSubject, aclTestRepoDid)); err != nil { - t.Fatalf("a foreign-knot collaboratorUpdate must be dropped, not error: %v", err) - } - if len(acl.collabAdd) != 0 { - t.Errorf("a knot that does not host the repo must not mutate its collaborators: %v", acl.collabAdd) - } - - if err := ingestCollaboratorUpdate(ctx, d, acl, source, collabEvent(t, knotdb.AclOpRemove, aclTestSubject, aclTestRepoDid)); err != nil { - t.Fatalf("a foreign-knot remove must be dropped, not error: %v", err) - } - if len(acl.collabRemove) != 0 { - t.Errorf("a knot that does not host the repo must not remove its collaborators: %v", acl.collabRemove) - } -} - -func TestIngestCollaboratorUpdate_BadDidErrors(t *testing.T) { - ctx := context.Background() - d := aclTestDB(t) - acl := &recordingAcl{} - source := ec.Source{Kind: ec.KindKnot, Host: aclTestHost} - if err := ingestCollaboratorUpdate(ctx, d, acl, source, collabEvent(t, knotdb.AclOpAdd, "not-a-did", aclTestRepoDid)); err == nil { - t.Fatal("a malformed subject DID must be rejected") - } - if len(acl.collabAdd) != 0 { - t.Errorf("a malformed subject must not reach the roster: %v", acl.collabAdd) - } -} - -func TestIngestKnotMemberUpdate_RetriesTransientThenSucceeds(t *testing.T) { - acl := &flakyAcl{failsLeft: aclIngestAttempts - 1} - source := ec.Source{Kind: ec.KindKnot, Host: aclTestHost} - - if err := ingestKnotMemberUpdate(acl, source, memberEvent(t, knotdb.AclOpAdd, aclTestSubject)); err != nil { - t.Fatalf("a transient store error within the retry budget must recover, got: %v", err) - } - if acl.calls != aclIngestAttempts { - t.Errorf("calls = %d, want %d; the write must retry until it lands", acl.calls, aclIngestAttempts) - } -} - -func TestIngestKnotMemberUpdate_GivesUpAfterAttempts(t *testing.T) { - acl := &flakyAcl{failsLeft: aclIngestAttempts + 5} - source := ec.Source{Kind: ec.KindKnot, Host: aclTestHost} - - if err := ingestKnotMemberUpdate(acl, source, memberEvent(t, knotdb.AclOpAdd, aclTestSubject)); err == nil { - t.Fatal("a persistent store error must surface so the failure is logged") - } - if acl.calls != aclIngestAttempts { - t.Errorf("calls = %d, want %d; the retry must be bounded", acl.calls, aclIngestAttempts) - } - if acl.membersInvalid != 1 { - t.Errorf("membersInvalid = %d, want 1; a dropped delta must invalidate the scope so the next read reconciles instead of waiting out the TTL", acl.membersInvalid) - } -} - -func TestIngestCollaboratorUpdate_InvalidatesScopeOnGiveUp(t *testing.T) { - ctx := context.Background() - d := aclTestDB(t) - seedAclRepo(t, d) - acl := &flakyAcl{failsLeft: aclIngestAttempts + 5} - source := ec.Source{Kind: ec.KindKnot, Host: aclTestHost} - - if err := ingestCollaboratorUpdate(ctx, d, acl, source, collabEvent(t, knotdb.AclOpAdd, aclTestSubject, aclTestRepoDid)); err == nil { - t.Fatal("a persistent store error must surface so the failure is logged") - } - if acl.collabsInvalid != 1 { - t.Errorf("collabsInvalid = %d, want 1; a dropped delta must invalidate the scope", acl.collabsInvalid) - } -} - -func TestIngestKnotMemberUpdate_NoInvalidateOnSuccess(t *testing.T) { - acl := &flakyAcl{failsLeft: aclIngestAttempts - 1} - source := ec.Source{Kind: ec.KindKnot, Host: aclTestHost} - - if err := ingestKnotMemberUpdate(acl, source, memberEvent(t, knotdb.AclOpAdd, aclTestSubject)); err != nil { - t.Fatalf("a recoverable delta must not error: %v", err) - } - if acl.membersInvalid != 0 { - t.Errorf("membersInvalid = %d, want 0; a delta that lands must not force a reconcile", acl.membersInvalid) - } -} - -func TestIngestCollaboratorUpdate_StoreErrorPropagates(t *testing.T) { - ctx := context.Background() - d := aclTestDB(t) - if err := d.Close(); err != nil { - t.Fatalf("close: %v", err) - } - acl := &recordingAcl{} - source := ec.Source{Kind: ec.KindKnot, Host: aclTestHost} - - if err := ingestCollaboratorUpdate(ctx, d, acl, source, collabEvent(t, knotdb.AclOpAdd, aclTestSubject, aclTestRepoDid)); err == nil { - t.Fatal("a store error on the repo lookup must surface, not be swallowed as an unindexed-repo skip") - } - if len(acl.collabAdd) != 0 { - t.Errorf("a failed repo lookup must not reach the roster: %v", acl.collabAdd) - } -} diff --git a/appview/state/router.go b/appview/state/router.go index 08745d67d..041066ec0 100644 --- a/appview/state/router.go +++ b/appview/state/router.go @@ -365,7 +365,7 @@ func (s *State) KnotsRouter() http.Handler { Enforcer: s.enforcer, Acl: s.aclService, IdResolver: s.idResolver, - Knotstream: s.knotstream, + SiteFeed: s.siteFeed, Logger: logger, } diff --git a/appview/state/state.go b/appview/state/state.go index 9d4e041cc..d99523a05 100644 --- a/appview/state/state.go +++ b/appview/state/state.go @@ -35,8 +35,8 @@ import ( "tangled.org/core/appview/pages" pipelinessh "tangled.org/core/appview/pipelines/ssh" "tangled.org/core/appview/reporesolver" + "tangled.org/core/appview/sitefeed" "tangled.org/core/consts" - "tangled.org/core/eventconsumer" "tangled.org/core/idresolver" "tangled.org/core/jetstream" "tangled.org/core/log" @@ -71,7 +71,7 @@ type State struct { config *config.Config repoResolver *reporesolver.RepoResolver aclService *knotacl.Service - knotstream *eventconsumer.Consumer + siteFeed *sitefeed.Feed logger *slog.Logger cfClient *cloudflare.Client codesearch *codesearch.CodeSearch @@ -218,11 +218,8 @@ func Make(ctx context.Context, config *config.Config) (*State, error) { } } - knotstream, err := Knotstream(ctx, config, d, aclService, enforcer, posthog, notifier, cfClient) - if err != nil { - return nil, fmt.Errorf("failed to start knotstream consumer: %w", err) - } - knotstream.Start(ctx) + siteFeed := sitefeed.New(d, rdb, config, cfClient, logger) + go siteFeed.Start(ctx) state := &State{ db: d, @@ -239,7 +236,7 @@ func Make(ctx context.Context, config *config.Config) (*State, error) { config: config, repoResolver: repoResolver, aclService: aclService, - knotstream: knotstream, + siteFeed: siteFeed, logger: logger, cfClient: cfClient, codesearch: &codesearch.CodeSearch{Host: config.CodeSearch.ZoektUrl}, diff --git a/appview/state/streams.go b/appview/state/streams.go deleted file mode 100644 index 4e7c748a6..000000000 --- a/appview/state/streams.go +++ /dev/null @@ -1,45 +0,0 @@ -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, - 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} - ec.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, - CursorStore: &cursorStore, - }) -} diff --git a/eventconsumer/consumer.go b/eventconsumer/consumer.go deleted file mode 100644 index 7d5d4c4cf..000000000 --- a/eventconsumer/consumer.go +++ /dev/null @@ -1,366 +0,0 @@ -package eventconsumer - -import ( - "context" - "encoding/json" - "log/slog" - "net/http" - "sync" - "time" - - "tangled.org/core/eventconsumer/cursor" - "tangled.org/core/eventstream" - "tangled.org/core/log" - - "github.com/avast/retry-go/v4" - "github.com/gorilla/websocket" -) - -type ProcessFunc func(ctx context.Context, source Source, event eventstream.Event) error - -// server sends a ping every 30s, so any silence longer than this means the -// connection is half-open, dropped without a close frame. -// so without a read deadline, ReadMessage only notices when the kernel's -// tcp keepalive gives up. -const livenessTimeout = 90 * time.Second - -type ConsumerConfig struct { - Sources map[Source]struct{} - ProcessFunc ProcessFunc - RetryInterval time.Duration - MaxRetryInterval time.Duration - ConnectionTimeout time.Duration - WorkerCount int - QueueSize int - Logger *slog.Logger - CursorStore cursor.Store - - Dialer *websocket.Dialer - RequestHeader http.Header - MaxRetryAttempts uint - OnConnectExceeded func(Source, error) -} - -func NewConsumerConfig() *ConsumerConfig { - return &ConsumerConfig{ - Sources: make(map[Source]struct{}), - } -} - -type Consumer struct { - sourceWg sync.WaitGroup - workerWg sync.WaitGroup - dialer *websocket.Dialer - jobQueue chan job - logger *slog.Logger - - // sourcesMu guards sources. It must only be held for short, non-blocking - // map operations; never across a blocking call (dial, read, close). - sourcesMu sync.Mutex - sources map[Source]*sourceState - - cfg ConsumerConfig -} - -type sourceState struct { - cancel context.CancelFunc - conn *websocket.Conn - - cursorMu sync.Mutex - cursorMax int64 -} - -type job struct { - source Source - message []byte -} - -func NewConsumer(cfg ConsumerConfig) *Consumer { - if cfg.RetryInterval == 0 { - cfg.RetryInterval = 15 * time.Minute - } - if cfg.ConnectionTimeout == 0 { - cfg.ConnectionTimeout = 10 * time.Second - } - if cfg.WorkerCount <= 0 { - cfg.WorkerCount = 5 - } - if cfg.MaxRetryInterval == 0 { - cfg.MaxRetryInterval = 1 * time.Hour - } - if cfg.Logger == nil { - cfg.Logger = log.New("consumer") - } - if cfg.QueueSize == 0 { - cfg.QueueSize = 100 - } - if cfg.CursorStore == nil { - cfg.CursorStore = &cursor.MemoryStore{} - } - dialer := cfg.Dialer - if dialer == nil { - dialer = websocket.DefaultDialer - } - return &Consumer{ - cfg: cfg, - dialer: dialer, - jobQueue: make(chan job, cfg.QueueSize), - logger: cfg.Logger, - sources: make(map[Source]*sourceState), - } -} - -func (c *Consumer) Start(ctx context.Context) { - c.cfg.Logger.Info("starting consumer", "config", c.cfg) - - for range c.cfg.WorkerCount { - c.workerWg.Add(1) - go c.worker(ctx) - } - - for source := range c.cfg.Sources { - c.AddSource(ctx, source) - } -} - -func (c *Consumer) Stop() { - // snapshot cancels and conns under lock so we don't hold sourcesMu across Close - c.sourcesMu.Lock() - cancels := make([]context.CancelFunc, 0, len(c.sources)) - conns := make([]*websocket.Conn, 0, len(c.sources)) - for _, st := range c.sources { - if st.cancel != nil { - cancels = append(cancels, st.cancel) - } - if st.conn != nil { - conns = append(conns, st.conn) - } - } - c.sourcesMu.Unlock() - - for _, cancel := range cancels { - cancel() - } - for _, conn := range conns { - conn.Close() - } - - c.sourceWg.Wait() - close(c.jobQueue) - c.workerWg.Wait() -} - -func (c *Consumer) AddSource(ctx context.Context, s Source) { - c.sourcesMu.Lock() - if _, ok := c.sources[s]; ok { - c.sourcesMu.Unlock() - c.logger.Info("source already present", "source", s) - return - } - srcCtx, cancel := context.WithCancel(ctx) - c.sources[s] = &sourceState{cancel: cancel} - c.sourcesMu.Unlock() - - c.sourceWg.Add(1) - go c.startConnectionLoop(srcCtx, s) -} - -func (c *Consumer) RemoveSource(s Source) { - c.sourcesMu.Lock() - st, ok := c.sources[s] - if !ok { - c.sourcesMu.Unlock() - c.logger.Info("source not present", "source", s) - return - } - delete(c.sources, s) - cancel := st.cancel - conn := st.conn - c.sourcesMu.Unlock() - - // release lock before any potentially blocking call - if cancel != nil { - cancel() - } - if conn != nil { - conn.Close() - } -} - -func (c *Consumer) worker(ctx context.Context) { - defer c.workerWg.Done() - for { - select { - case <-ctx.Done(): - return - case j, ok := <-c.jobQueue: - if !ok { - return - } - - var ev eventstream.Event - err := json.Unmarshal(j.message, &ev) - if err != nil { - c.logger.Error("error deserializing message", "source", j.source.Key(), "err", err) - continue - } - - if err := c.cfg.ProcessFunc(ctx, j.source, ev); err != nil { - c.logger.Error("error processing message", "source", j.source, "err", err) - } - - c.advanceCursor(j.source, ev.Created) - } - } -} - -func (c *Consumer) advanceCursor(s Source, newCursor int64) { - if newCursor == 0 { - return - } - c.sourcesMu.Lock() - st, ok := c.sources[s] - c.sourcesMu.Unlock() - if !ok { - return - } - - st.cursorMu.Lock() - defer st.cursorMu.Unlock() - if newCursor <= st.cursorMax { - return - } - st.cursorMax = newCursor - c.cfg.CursorStore.Set(s.Key(), newCursor) -} - -func (c *Consumer) startConnectionLoop(ctx context.Context, source Source) { - defer c.sourceWg.Done() - - // attempt connection initially - err := c.runConnection(ctx, source) - if err != nil { - c.logger.Error("failed to run connection", "err", err) - } - - timer := time.NewTimer(1 * time.Minute) - defer timer.Stop() - - // every subsequent attempt is delayed by 1 minute - for { - select { - case <-ctx.Done(): - return - case <-timer.C: - err := c.runConnection(ctx, source) - if err != nil { - c.logger.Error("failed to run connection", "err", err) - } - timer.Reset(1 * time.Minute) - } - } -} - -func (c *Consumer) runConnection(ctx context.Context, source Source) error { - cursor := c.cfg.CursorStore.Get(source.Key()) - - u, err := source.URL(cursor) - if err != nil { - return err - } - - c.logger.Info("connecting", "url", u.String()) - - retryOpts := []retry.Option{ - retry.Attempts(c.cfg.MaxRetryAttempts), - retry.DelayType(retry.BackOffDelay), - retry.Delay(c.cfg.RetryInterval), - retry.MaxDelay(c.cfg.MaxRetryInterval), - retry.MaxJitter(c.cfg.RetryInterval / 5), - retry.OnRetry(func(n uint, err error) { - c.logger.Info("retrying connection", - "source", source, - "url", u.String(), - "attempt", n+1, - "err", err, - ) - }), - retry.Context(ctx), - } - - var conn *websocket.Conn - - err = retry.Do(func() error { - connCtx, cancel := context.WithTimeout(ctx, c.cfg.ConnectionTimeout) - defer cancel() - conn, _, err = c.dialer.DialContext(connCtx, u.String(), c.cfg.RequestHeader) - return err - }, retryOpts...) - if err != nil { - if c.cfg.OnConnectExceeded != nil { - c.cfg.OnConnectExceeded(source, err) - } - return err - } - - // Register the conn. If the source was removed (or our ctx cancelled) - // while we were dialing, drop this conn instead of installing it. - c.sourcesMu.Lock() - st, ok := c.sources[source] - if !ok || ctx.Err() != nil { - c.sourcesMu.Unlock() - conn.Close() - if ctx.Err() != nil { - return ctx.Err() - } - return nil - } - st.conn = conn - c.sourcesMu.Unlock() - - defer func() { - // Clear the conn from state, but only if it's still our conn (a - // concurrent RemoveSource may have already done it). - c.sourcesMu.Lock() - if st, ok := c.sources[source]; ok && st.conn == conn { - st.conn = nil - } - c.sourcesMu.Unlock() - conn.Close() - }() - - c.logger.Info("connected", "source", source) - - conn.SetReadDeadline(time.Now().Add(livenessTimeout)) - conn.SetPongHandler(func(string) error { - return conn.SetReadDeadline(time.Now().Add(livenessTimeout)) - }) - conn.SetPingHandler(func(appData string) error { - err := conn.WriteControl(websocket.PongMessage, []byte(appData), time.Now().Add(10*time.Second)) - if err != nil { - return err - } - return conn.SetReadDeadline(time.Now().Add(livenessTimeout)) - }) - - for { - select { - case <-ctx.Done(): - return nil - default: - msgType, msg, err := conn.ReadMessage() - if err != nil { - return err - } - if msgType != websocket.TextMessage { - continue - } - conn.SetReadDeadline(time.Now().Add(livenessTimeout)) - select { - case c.jobQueue <- job{source: source, message: msg}: - case <-ctx.Done(): - return nil - } - } - } -} diff --git a/eventconsumer/consumer_test.go b/eventconsumer/consumer_test.go deleted file mode 100644 index 2cd19082b..000000000 --- a/eventconsumer/consumer_test.go +++ /dev/null @@ -1,278 +0,0 @@ -package eventconsumer - -import ( - "context" - "encoding/json" - "fmt" - "io" - "log/slog" - "net/http" - "net/http/httptest" - "strings" - "sync" - "testing" - "time" - - "tangled.org/core/eventconsumer/cursor" - "tangled.org/core/eventstream" - "tangled.org/core/notifier" -) - -type memSrc struct { - mu sync.Mutex - events []eventstream.Event -} - -func (s *memSrc) add(ev eventstream.Event) { - s.mu.Lock() - defer s.mu.Unlock() - s.events = append(s.events, ev) -} - -func (s *memSrc) GetEvents(cursor int64, limit int) ([]eventstream.Event, error) { - s.mu.Lock() - defer s.mu.Unlock() - out := []eventstream.Event{} - for _, ev := range s.events { - if ev.Created > cursor { - out = append(out, ev) - if len(out) == limit { - break - } - } - } - return out, nil -} - -func mkEv(i int) eventstream.Event { - return eventstream.Event{ - Rkey: fmt.Sprintf("rk-%04d", i), - Nsid: "sh.tangled.test", - EventJson: json.RawMessage(fmt.Sprintf(`{"i":%d}`, i)), - Created: int64(i + 1), - } -} - -func startEventServer(t *testing.T, src *memSrc) (Source, *notifier.Notifier) { - t.Helper() - n := notifier.New() - mux := http.NewServeMux() - mux.HandleFunc("/events", func(w http.ResponseWriter, r *http.Request) { - _ = eventstream.Stream(w, r, eventstream.StreamConfig{ - Backend: src, - Notifier: &n, - Logger: slog.New(slog.NewTextHandler(io.Discard, nil)), - BatchSize: 5, - MaxBatchesPerDrain: 100, - }) - }) - srv := httptest.NewServer(mux) - t.Cleanup(srv.Close) - addr := strings.TrimPrefix(srv.URL, "http://") - return Source{Kind: "test", Host: addr, NoTLS: true}, &n -} - -func TestConsumer_DrainAdvancesCursor(t *testing.T) { - src := &memSrc{} - for i := range 8 { - src.add(mkEv(i)) - } - - source, _ := startEventServer(t, src) - - store := &cursor.MemoryStore{} - seenMu := sync.Mutex{} - seen := []int64{} - - cfg := ConsumerConfig{ - ProcessFunc: func(ctx context.Context, _ Source, msg eventstream.Event) error { - seenMu.Lock() - seen = append(seen, msg.Created) - seenMu.Unlock() - return nil - }, - WorkerCount: 1, - QueueSize: 16, - ConnectionTimeout: 2 * time.Second, - CursorStore: store, - Logger: slog.New(slog.NewTextHandler(io.Discard, nil)), - } - c := NewConsumer(cfg) - - ctx, cancel := context.WithCancel(context.Background()) - defer cancel() - - c.Start(ctx) - c.AddSource(ctx, source) - - deadline := time.Now().Add(3 * time.Second) - for time.Now().Before(deadline) { - seenMu.Lock() - n := len(seen) - seenMu.Unlock() - if n >= 8 { - break - } - time.Sleep(20 * time.Millisecond) - } - - seenMu.Lock() - defer seenMu.Unlock() - if len(seen) != 8 { - t.Fatalf("processed %d events, want 8: %v", len(seen), seen) - } - for i, got := range seen { - if got != int64(i+1) { - t.Fatalf("event %d: got created=%d want %d", i, got, i+1) - } - } - - if final := store.Get(source.Key()); final != 8 { - t.Fatalf("cursor = %d, want 8", final) - } -} - -func TestConsumer_CursorMonotonic_OutOfOrderWorkers(t *testing.T) { - src := &memSrc{} - for i := range 4 { - src.add(mkEv(i)) - } - - source, _ := startEventServer(t, src) - - store := &cursor.MemoryStore{} - - releaseFirst := make(chan struct{}) - processed := make(chan int64, 4) - - cfg := ConsumerConfig{ - ProcessFunc: func(ctx context.Context, _ Source, msg eventstream.Event) error { - if msg.Created == 1 { - <-releaseFirst - } - processed <- msg.Created - return nil - }, - WorkerCount: 4, - QueueSize: 16, - ConnectionTimeout: 2 * time.Second, - CursorStore: store, - Logger: slog.New(slog.NewTextHandler(io.Discard, nil)), - } - c := NewConsumer(cfg) - - ctx, cancel := context.WithCancel(context.Background()) - defer cancel() - - c.Start(ctx) - c.AddSource(ctx, source) - - for range 3 { - select { - case <-processed: - case <-time.After(3 * time.Second): - t.Fatal("timed out waiting for events 2-4 to be processed") - } - } - - if cur := store.Get(source.Key()); cur != 4 { - t.Fatalf("cursor before slow worker finished = %d, want 4", cur) - } - - close(releaseFirst) - select { - case <-processed: - case <-time.After(3 * time.Second): - t.Fatal("timed out waiting for slow worker") - } - - if cur := store.Get(source.Key()); cur != 4 { - t.Fatalf("cursor regressed after slow worker: %d, want 4", cur) - } -} - -func TestConsumer_StopTerminatesWithoutCtxCancel(t *testing.T) { - src := &memSrc{} - source, _ := startEventServer(t, src) - - cfg := ConsumerConfig{ - ProcessFunc: func(ctx context.Context, _ Source, _ eventstream.Event) error { return nil }, - WorkerCount: 2, - QueueSize: 8, - ConnectionTimeout: 2 * time.Second, - CursorStore: &cursor.MemoryStore{}, - Logger: slog.New(slog.NewTextHandler(io.Discard, nil)), - } - c := NewConsumer(cfg) - - c.Start(context.Background()) - c.AddSource(context.Background(), source) - - done := make(chan struct{}) - go func() { - c.Stop() - close(done) - }() - - select { - case <-done: - case <-time.After(5 * time.Second): - t.Fatal("Stop did not return within 5s") - } -} - -func TestConsumer_ResumesFromStoredCursor(t *testing.T) { - src := &memSrc{} - for i := range 5 { - src.add(mkEv(i)) - } - - source, _ := startEventServer(t, src) - - store := &cursor.MemoryStore{} - store.Set(source.Key(), 3) - - seenMu := sync.Mutex{} - seen := []int64{} - - cfg := ConsumerConfig{ - ProcessFunc: func(ctx context.Context, _ Source, msg eventstream.Event) error { - seenMu.Lock() - seen = append(seen, msg.Created) - seenMu.Unlock() - return nil - }, - WorkerCount: 1, - QueueSize: 16, - ConnectionTimeout: 2 * time.Second, - CursorStore: store, - Logger: slog.New(slog.NewTextHandler(io.Discard, nil)), - } - c := NewConsumer(cfg) - - ctx, cancel := context.WithCancel(context.Background()) - defer cancel() - - c.Start(ctx) - c.AddSource(ctx, source) - - deadline := time.Now().Add(3 * time.Second) - for time.Now().Before(deadline) { - seenMu.Lock() - n := len(seen) - seenMu.Unlock() - if n >= 2 { - break - } - time.Sleep(20 * time.Millisecond) - } - - seenMu.Lock() - defer seenMu.Unlock() - if len(seen) < 2 { - t.Fatalf("processed %d events, want 2: %v", len(seen), seen) - } - if seen[0] != 4 || seen[1] != 5 { - t.Fatalf("resumed events = %v, want [4 5]", seen) - } -} diff --git a/eventconsumer/cursor/memory.go b/eventconsumer/cursor/memory.go deleted file mode 100644 index 722348ca5..000000000 --- a/eventconsumer/cursor/memory.go +++ /dev/null @@ -1,22 +0,0 @@ -package cursor - -import ( - "sync" -) - -type MemoryStore struct { - store sync.Map -} - -func (m *MemoryStore) Set(key string, cursor int64) { - m.store.Store(key, cursor) -} - -func (m *MemoryStore) Get(key string) (cursor int64) { - if result, ok := m.store.Load(key); ok { - if val, ok := result.(int64); ok { - return val - } - } - return 0 -} diff --git a/eventconsumer/cursor/redis.go b/eventconsumer/cursor/redis.go deleted file mode 100644 index 3e67104db..000000000 --- a/eventconsumer/cursor/redis.go +++ /dev/null @@ -1,41 +0,0 @@ -package cursor - -import ( - "context" - "fmt" - "strconv" - - "tangled.org/core/appview/cache" -) - -const ( - cursorKey = "cursor:%s" -) - -type RedisStore struct { - rdb *cache.Cache -} - -func NewRedisCursorStore(cache *cache.Cache) RedisStore { - return RedisStore{ - rdb: cache, - } -} - -func (r *RedisStore) Set(key string, cursor int64) { - k := fmt.Sprintf(cursorKey, key) - r.rdb.Set(context.Background(), k, cursor, 0) -} - -func (r *RedisStore) Get(key string) (cursor int64) { - k := fmt.Sprintf(cursorKey, key) - val, err := r.rdb.Get(context.Background(), k).Result() - if err != nil { - return 0 - } - parsed, err := strconv.ParseInt(val, 10, 64) - if err != nil { - return 0 - } - return parsed -} diff --git a/eventconsumer/cursor/sqlite.go b/eventconsumer/cursor/sqlite.go deleted file mode 100644 index 4d1b74288..000000000 --- a/eventconsumer/cursor/sqlite.go +++ /dev/null @@ -1,83 +0,0 @@ -package cursor - -import ( - "database/sql" - "errors" - "fmt" - "log/slog" - - _ "github.com/mattn/go-sqlite3" -) - -type SqliteStore struct { - db *sql.DB - tableName string -} - -type SqliteStoreOpt func(*SqliteStore) - -func WithTableName(name string) SqliteStoreOpt { - return func(s *SqliteStore) { - s.tableName = name - } -} - -func NewSQLiteStore(dbPath string, opts ...SqliteStoreOpt) (*SqliteStore, error) { - db, err := sql.Open("sqlite3", dbPath+"?_foreign_keys=1") - if err != nil { - return nil, fmt.Errorf("failed to open sqlite database: %w", err) - } - - store := &SqliteStore{ - db: db, - tableName: "cursors", - } - - for _, o := range opts { - o(store) - } - - if err := store.init(); err != nil { - return nil, err - } - - return store, nil -} - -func (s *SqliteStore) init() error { - createTable := fmt.Sprintf(` - create table if not exists %s ( - knot text primary key, - cursor integer - );`, s.tableName) - _, err := s.db.Exec(createTable) - return err -} - -func (s *SqliteStore) Set(key string, cursor int64) { - query := fmt.Sprintf(` - insert into %s (knot, cursor) - values (?, ?) - on conflict(knot) do update set cursor=excluded.cursor; - `, s.tableName) - - if _, err := s.db.Exec(query, key, cursor); err != nil { - slog.Default().Error("cursor sqlite set failed", "key", key, "cursor", cursor, "err", err) - } -} - -func (s *SqliteStore) Get(key string) (cursor int64) { - query := fmt.Sprintf(` - select cursor from %s where knot = ?; - `, s.tableName) - err := s.db.QueryRow(query, key).Scan(&cursor) - - if err != nil { - if !errors.Is(err, sql.ErrNoRows) { - slog.Default().Error("cursor sqlite get failed", "key", key, "err", err) - } - return 0 - } - - return cursor -} diff --git a/eventconsumer/cursor/store.go b/eventconsumer/cursor/store.go deleted file mode 100644 index c6c8db556..000000000 --- a/eventconsumer/cursor/store.go +++ /dev/null @@ -1,6 +0,0 @@ -package cursor - -type Store interface { - Set(key string, cursor int64) - Get(key string) (cursor int64) -} diff --git a/eventconsumer/migrate_test.go b/eventconsumer/migrate_test.go deleted file mode 100644 index 5187f12ac..000000000 --- a/eventconsumer/migrate_test.go +++ /dev/null @@ -1,67 +0,0 @@ -package eventconsumer - -import ( - "testing" - - "tangled.org/core/eventconsumer/cursor" -) - -func TestMigrateLegacyCursor(t *testing.T) { - const host = "whelk.knot.tld" - src := NewKnotSource(host) - - t.Run("copies legacy bare-host cursor to namespaced key", func(t *testing.T) { - store := &cursor.MemoryStore{} - store.Set(host, 42) - - MigrateLegacyCursor(store, src) - - if got := store.Get(src.Key()); got != 42 { - t.Fatalf("namespaced key = %d, want 42", got) - } - if got := store.Get(host); got != 42 { - t.Fatalf("legacy key = %d, want it left at 42", got) - } - }) - - t.Run("does not pave over an already-migrated cursor", func(t *testing.T) { - store := &cursor.MemoryStore{} - store.Set(src.Key(), 100) - store.Set(host, 42) - - MigrateLegacyCursor(store, src) - - if got := store.Get(src.Key()); got != 100 { - t.Fatalf("namespaced key = %d, want 100", got) - } - }) - - t.Run("noop when neither key is set", func(t *testing.T) { - store := &cursor.MemoryStore{} - - MigrateLegacyCursor(store, src) - - if got := store.Get(src.Key()); got != 0 { - t.Fatalf("namespaced key = %d, want 0", got) - } - }) - - t.Run("namespaces the same host by kind", func(t *testing.T) { - const shared = "mussel.knot.tld" - knot := NewKnotSource(shared) - spindle := NewSpindleSource(shared) - - store := &cursor.MemoryStore{} - store.Set(shared, 500) - - MigrateLegacyCursor(store, knot) - MigrateLegacyCursor(store, spindle) - - if got := store.Get(knot.Key()); got != 500 { - t.Fatalf("knot key = %d, want 500", got) - } - if got := store.Get(spindle.Key()); got != 500 { - t.Fatalf("spindle key = %d, want 500", got) - } - }) -} diff --git a/eventconsumer/source.go b/eventconsumer/source.go deleted file mode 100644 index 010d42678..000000000 --- a/eventconsumer/source.go +++ /dev/null @@ -1,59 +0,0 @@ -package eventconsumer - -import ( - "net/url" - "strconv" - - "tangled.org/core/eventconsumer/cursor" - "tangled.org/core/hostutil" -) - -type Kind string - -const ( - KindKnot Kind = "knot" - KindSpindle Kind = "spindle" -) - -type Source struct { - Kind Kind - Host string - NoTLS bool // use TLS by default -} - -func NewKnotSource(host string) Source { - host, noTLS, _ := hostutil.ParseHostname(host) - return Source{Kind: KindKnot, Host: host, NoTLS: noTLS} -} -func NewSpindleSource(host string) Source { - host, noTLS, _ := hostutil.ParseHostname(host) - return Source{Kind: KindSpindle, Host: host, NoTLS: noTLS} -} - -func (s Source) Key() string { return string(s.Kind) + ":" + s.Host } - -func MigrateLegacyCursor(store cursor.Store, s Source) { - if store.Get(s.Key()) != 0 { - return - } - if legacy := store.Get(s.Host); legacy != 0 { - store.Set(s.Key(), legacy) - } -} - -func (s Source) URL(cursor int64) (*url.URL, error) { - scheme := "wss" - if s.NoTLS { - scheme = "ws" - } - u, err := url.Parse(scheme + "://" + s.Host + "/events") - if err != nil { - return nil, err - } - if cursor != 0 { - q := url.Values{} - q.Add("cursor", strconv.FormatInt(cursor, 10)) - u.RawQuery = q.Encode() - } - return u, nil -} diff --git a/eventconsumer/upgrade_test.go b/eventconsumer/upgrade_test.go deleted file mode 100644 index 6376645b2..000000000 --- a/eventconsumer/upgrade_test.go +++ /dev/null @@ -1,105 +0,0 @@ -package eventconsumer - -import ( - "context" - "io" - "log/slog" - "path/filepath" - "sync" - "testing" - "time" - - "tangled.org/core/eventconsumer/cursor" - "tangled.org/core/eventstream" -) - -func sqliteCursorStore(t *testing.T) cursor.Store { - t.Helper() - store, err := cursor.NewSQLiteStore(filepath.Join(t.TempDir(), "spindle.db")) - if err != nil { - t.Fatalf("new sqlite cursor store: %v", err) - } - return store -} - -func drainProcessed(t *testing.T, store cursor.Store, source Source, expected int) []int64 { - t.Helper() - - var mu sync.Mutex - var seen []int64 - - c := NewConsumer(ConsumerConfig{ - ProcessFunc: func(_ context.Context, _ Source, msg eventstream.Event) error { - mu.Lock() - seen = append(seen, msg.Created) - mu.Unlock() - return nil - }, - WorkerCount: 1, - QueueSize: 16, - ConnectionTimeout: 2 * time.Second, - CursorStore: store, - Logger: slog.New(slog.NewTextHandler(io.Discard, nil)), - }) - - ctx, cancel := context.WithCancel(context.Background()) - defer cancel() - c.Start(ctx) - c.AddSource(ctx, source) - - deadline := time.Now().Add(3 * time.Second) - for time.Now().Before(deadline) { - mu.Lock() - n := len(seen) - mu.Unlock() - if n >= expected { - break - } - time.Sleep(20 * time.Millisecond) - } - - c.Stop() - - mu.Lock() - defer mu.Unlock() - return append([]int64(nil), seen...) -} - -func TestSpindleUpgrade_OrphanedCursorReplaysFromZero(t *testing.T) { - src := &memSrc{} - for i := range 8 { - src.add(mkEv(i)) - } - source, _ := startEventServer(t, src) - - store := sqliteCursorStore(t) - store.Set(source.Host, 5) - - seen := drainProcessed(t, store, source, 8) - - if len(seen) != 8 { - t.Fatalf("orphaned bare-host cursor processed %d events, want a full replay of 8: %v", len(seen), seen) - } -} - -func TestSpindleUpgrade_MigratedCursorResumesNoReplay(t *testing.T) { - src := &memSrc{} - for i := range 8 { - src.add(mkEv(i)) - } - source, _ := startEventServer(t, src) - - store := sqliteCursorStore(t) - store.Set(source.Host, 5) - - MigrateLegacyCursor(store, source) - - seen := drainProcessed(t, store, source, 3) - - if len(seen) != 3 { - t.Fatalf("migrated cursor processed %d events, want a resume of 3: %v", len(seen), seen) - } - if seen[0] != 6 || seen[2] != 8 { - t.Fatalf("resumed events = %v, want [6 7 8]", seen) - } -} diff --git a/knotmirror/db/hosts.go b/knotmirror/db/hosts.go index 5fbb19cb5..2faec40e2 100644 --- a/knotmirror/db/hosts.go +++ b/knotmirror/db/hosts.go @@ -56,7 +56,7 @@ func StoreCursors(ctx context.Context, e *sql.DB, cursors []models.HostCursor) e } defer tx.Rollback() for _, cur := range cursors { - if cur.LastSeq <= 0 { + if cur.LastSeq < 0 { continue } if _, err := tx.ExecContext(ctx, diff --git a/knotmirror/db/repos.go b/knotmirror/db/repos.go index 106282135..906120b87 100644 --- a/knotmirror/db/repos.go +++ b/knotmirror/db/repos.go @@ -61,17 +61,6 @@ func UpdateRepoState(ctx context.Context, e DBTX, repoDid syntax.DID, state mode return nil } -func DeleteRepo(ctx context.Context, e DBTX, did syntax.DID, rkey syntax.RecordKey) error { - if _, err := e.ExecContext(ctx, - `delete from repos where did = $1 and rkey = $2`, - did, - rkey, - ); err != nil { - return fmt.Errorf("deleting repo: %w", err) - } - return nil -} - const repoColumns = ` did, rkey, @@ -124,6 +113,26 @@ func GetRepoByRepoDid(ctx context.Context, e DBTX, repoDid syntax.DID) (*models. return repo, nil } +func DesynchronizeReposByKnotHost(ctx context.Context, e DBTX, host *models.Host) (int64, error) { + res, err := e.ExecContext(ctx, + `update repos + set state = $1 + where state = $2 and knot_domain in ($3, $4)`, + models.RepoStateDesynchronized, + models.RepoStateActive, + host.URL(), + host.Hostname, + ) + if err != nil { + return 0, fmt.Errorf("desynchronizing repos by knot host: %w", err) + } + n, err := res.RowsAffected() + if err != nil { + return 0, fmt.Errorf("counting desynchronized repos: %w", err) + } + return n, nil +} + func GetRepoByAtUri(ctx context.Context, e DBTX, aturi syntax.ATURI) (*models.Repo, error) { row := e.QueryRowContext(ctx, `select`+repoColumns+` diff --git a/knotmirror/knotmirror.go b/knotmirror/knotmirror.go index 692deebe0..b2f042db4 100644 --- a/knotmirror/knotmirror.go +++ b/knotmirror/knotmirror.go @@ -68,7 +68,7 @@ func Run(ctx context.Context, cfg *config.Config) error { adminpage := NewAdminServer(logger, db, resyncer, xrpc, resolver) // maintain repository list with tap - // NOTE: this can be removed once we introduce did-for-repo because then we can just listen to KnotStream for #identity events. + // NOTE: this can be removed once we introduce did-for-repo because then we can just listen to the knot firehose for #identity frames. tap := NewTapClient(logger, cfg, db, gitm, knotstream) // start http server diff --git a/knotmirror/knotstream/metrics.go b/knotmirror/knotstream/metrics.go index c5845d708..e5c0996d9 100644 --- a/knotmirror/knotstream/metrics.go +++ b/knotmirror/knotstream/metrics.go @@ -5,24 +5,26 @@ import ( "github.com/prometheus/client_golang/prometheus/promauto" ) -// KnotStream metrics var ( - knotstreamEventsReceived = promauto.NewCounter(prometheus.CounterOpts{ - Name: "knotmirror_knotstream_events_received_total", - Help: "Total number of events received from knotstream", + firehoseRefOpsReceived = promauto.NewCounter(prometheus.CounterOpts{ + Name: "knotmirror_firehose_ref_ops_received_total", + Help: "Total number of git.ref ops received from knot firehoses", }) - knotstreamEventsProcessed = promauto.NewCounter(prometheus.CounterOpts{ - Name: "knotmirror_knotstream_events_processed_total", - Help: "Total number of events successfully processed", + firehoseRefOpsProcessed = promauto.NewCounter(prometheus.CounterOpts{ + Name: "knotmirror_firehose_ref_ops_processed_total", + Help: "Total number of git.ref ops that marked a repo desynchronized", }) - knotstreamEventsSkipped = promauto.NewCounter(prometheus.CounterOpts{ - Name: "knotmirror_knotstream_events_skipped_total", - Help: "Total number of events skipped (not tracked)", + firehoseRefOpsFailed = promauto.NewCounter(prometheus.CounterOpts{ + Name: "knotmirror_firehose_ref_ops_failed_total", + Help: "Total number of git.ref ops whose desynchronization mark failed", + }) + firehoseRefOpsSkipped = promauto.NewCounter(prometheus.CounterOpts{ + Name: "knotmirror_firehose_ref_ops_skipped_total", + Help: "Total number of git.ref ops skipped (unknown or non-active repo)", }) ) -// slurper metrics var connectedInbound = promauto.NewGauge(prometheus.GaugeOpts{ Name: "knotmirror_connected_inbound", - Help: "Number of inbound knotstream we are consuming", + Help: "Number of inbound knot firehoses we are consuming", }) diff --git a/knotmirror/knotstream/scheduler.go b/knotmirror/knotstream/scheduler.go index b3fd537b7..69721e3aa 100644 --- a/knotmirror/knotstream/scheduler.go +++ b/knotmirror/knotstream/scheduler.go @@ -4,9 +4,9 @@ import ( "context" "log/slog" "sync" - "sync/atomic" "time" + "tangled.org/core/knotfeed" "tangled.org/core/log" ) @@ -18,17 +18,19 @@ type ParallelScheduler struct { feeder chan *Task lk sync.Mutex scheduled map[string][]*Task - lastSeq atomic.Int64 logger *slog.Logger } type Task struct { - Key string - message []byte + Key string + + repoDid string + seq int64 + op knotfeed.RecordOp } -func NewParallelScheduler(maxC int, ident string, do func(context.Context, *Task) error) *ParallelScheduler { +func NewParallelScheduler(maxC int, do func(context.Context, *Task) error) *ParallelScheduler { return &ParallelScheduler{ concurrency: maxC, do: do, @@ -44,6 +46,25 @@ func (s *ParallelScheduler) Start(ctx context.Context) { } } +func (s *ParallelScheduler) Drain(ctx context.Context) { + for { + s.lk.Lock() + remaining := 0 + for _, queue := range s.scheduled { + remaining += len(queue) + } + s.lk.Unlock() + if remaining == 0 { + return + } + select { + case <-ctx.Done(): + return + case <-time.After(10 * time.Millisecond): + } + } +} + func (s *ParallelScheduler) AddTask(ctx context.Context, task *Task) { s.lk.Lock() if st, ok := s.scheduled[task.Key]; ok { @@ -72,6 +93,7 @@ func (s *ParallelScheduler) ForEach(ctx context.Context, fn func(context.Context default: } if err := fn(ctx, task); err != nil { + firehoseRefOpsFailed.Inc() s.logger.Error("event handler failed", "err", err) } @@ -88,15 +110,8 @@ func (s *ParallelScheduler) ForEach(ctx context.Context, fn func(context.Context task = rem[0] s.scheduled[task.Key] = rem[1:] } - - // TODO: update seq from received message - s.lastSeq.Store(time.Now().UnixNano()) }() s.lk.Unlock() } } } - -func (s *ParallelScheduler) LastSeq() int64 { - return s.lastSeq.Load() -} diff --git a/knotmirror/knotstream/slurper.go b/knotmirror/knotstream/slurper.go index 36f4318d5..f6dca55fe 100644 --- a/knotmirror/knotstream/slurper.go +++ b/knotmirror/knotstream/slurper.go @@ -3,24 +3,23 @@ package knotstream import ( "context" "database/sql" - "encoding/json" "fmt" "log/slog" - "math/rand" - "net/http" "sync" "time" "github.com/bluesky-social/indigo/atproto/syntax" - "github.com/bluesky-social/indigo/util/ssrf" - "github.com/carlmjohnson/versioninfo" "github.com/gorilla/websocket" + "tangled.org/core/knotfeed" "tangled.org/core/knotmirror/config" "tangled.org/core/knotmirror/db" "tangled.org/core/knotmirror/models" "tangled.org/core/log" + "tangled.org/core/netutil" ) +const maxConnectFailures = 30 + type KnotSlurper struct { logger *slog.Logger db *sql.DB @@ -31,6 +30,11 @@ type KnotSlurper struct { subs map[string]*subscription } +const ( + markAttempts = 3 + markRetryDelay = 100 * time.Millisecond +) + func NewKnotSlurper(l *slog.Logger, db *sql.DB, cfg *config.Config) *KnotSlurper { return &KnotSlurper{ logger: log.SubLogger(l, "slurper"), @@ -63,6 +67,17 @@ func (s *KnotSlurper) CheckIfSubscribed(hostname string) bool { } func (s *KnotSlurper) Shutdown(ctx context.Context) error { + s.subsLk.Lock() + subs := make([]*subscription, 0, len(s.subs)) + for _, sub := range s.subs { + subs = append(subs, sub) + } + s.subsLk.Unlock() + + for _, sub := range subs { + sub.scheduler.Drain(ctx) + } + s.logger.Info("starting shutdown host cursor flush") err := s.persistCursors(ctx) if err != nil { @@ -73,9 +88,6 @@ func (s *KnotSlurper) Shutdown(ctx context.Context) error { } func (s *KnotSlurper) persistCursors(ctx context.Context) error { - // // gather cursor list from subscriptions and store them to DB - // start := time.Now() - s.subsLk.Lock() cursors := make([]models.HostCursor, len(s.subs)) i := 0 @@ -85,9 +97,10 @@ func (s *KnotSlurper) persistCursors(ctx context.Context) error { } s.subsLk.Unlock() - err := db.StoreCursors(ctx, s.db, cursors) - // s.logger.Info("finished persisting cursors", "count", len(cursors), "duration", time.Since(start).String(), "err", err) - return err + if len(cursors) == 0 { + return nil + } + return db.StoreCursors(ctx, s.db, cursors) } func (s *KnotSlurper) Subscribe(host models.Host) error { @@ -99,27 +112,41 @@ func (s *KnotSlurper) Subscribe(host models.Host) error { return fmt.Errorf("already subscribed: %s", host.Hostname) } - // TODO: include `cancel` function to kill subscription by hostname sub := &subscription{ hostname: host.Hostname, - scheduler: NewParallelScheduler( - s.cfg.ConcurrencyPerHost, - host.Hostname, - s.ProcessEvent, - ), } + do := func(ctx context.Context, task *Task) error { + err := s.ProcessEvent(ctx, task) + if err == nil { + sub.MarkApplied(task.seq) + } + return err + } + sub.scheduler = NewParallelScheduler( + s.cfg.ConcurrencyPerHost, + do, + ) + sub.lastSeq.Store(host.LastSeq) + sub.appliedSeq.Store(host.LastSeq) s.subs[host.Hostname] = sub - // TODO: use service level context, not the top-most one. - // Using top-most context should be avoided to do graceful shutdown. ctx := context.TODO() sub.scheduler.Start(ctx) - go s.subscribeWithRedialer(ctx, host, sub) + go s.runConsumer(ctx, host, sub) return nil } -func (s *KnotSlurper) subscribeWithRedialer(ctx context.Context, host models.Host, sub *subscription) { +func (s *KnotSlurper) dialer(host models.Host) *websocket.Dialer { + if !host.NoSSL || s.ssrf { + d := netutil.SSRFWebsocketDialer(false) + d.HandshakeTimeout = 5 * time.Second + return d + } + return &websocket.Dialer{HandshakeTimeout: 5 * time.Second} +} + +func (s *KnotSlurper) runConsumer(ctx context.Context, host models.Host, sub *subscription) { l := s.logger.With("host", host.Hostname) defer func() { s.subsLk.Lock() @@ -129,237 +156,147 @@ func (s *KnotSlurper) subscribeWithRedialer(ctx context.Context, host models.Hos delete(s.subs, host.Hostname) }() - dialer := websocket.Dialer{ - HandshakeTimeout: time.Second * 5, - } - - // if this isn't a localhost / private connection, then we should enable SSRF protections - if !host.NoSSL || s.ssrf { - netDialer := ssrf.PublicOnlyDialer() - dialer.NetDialContext = netDialer.DialContext - } - - cursor := host.LastSeq + ctx, cancel := context.WithCancel(ctx) + defer cancel() connectedInbound.Inc() defer connectedInbound.Dec() - var backoff int - for { - select { - case <-ctx.Done(): - return - default: - } - u := host.LegacyEventsURL(cursor) - l.Debug("made url with cursor", "cursor", cursor, "url", u) - - // NOTE: manual backoff retry implementation to explicitly handle fails - hdr := make(http.Header) - hdr.Add("User-Agent", userAgent()) - conn, resp, err := dialer.DialContext(ctx, u, hdr) - if err != nil { - l.Warn("dialing failed", "err", err, "backoff", backoff) - time.Sleep(sleepForBackoff(backoff)) - backoff++ - if backoff > 30 { - l.Warn("host does not appear to be online, disabling for now") - host.Status = models.HostStatusOffline - if err := db.UpsertHost(ctx, s.db, &host); err != nil { - l.Error("failed to update host status", "err", err) + highWater := sub.LastSeq() + connectFailures := 0 + + consumer := &knotfeed.Consumer{ + Host: host.Hostname, + NoTLS: host.NoSSL, + Dialer: s.dialer(host), + Logger: l, + LoadCursor: func(context.Context) (int64, error) { + return sub.AppliedSeq(), nil + }, + StoreCursor: func(_ context.Context, seq int64) error { + sub.lastSeq.Store(seq) + return nil + }, + Handle: func(ctx context.Context, msg knotfeed.Message) error { + if msg.Type != knotfeed.TypeCommit || msg.Commit == nil { + return nil + } + ops := 0 + var first knotfeed.RecordOp + for _, op := range msg.Commit.Records { + if op.Collection != knotfeed.GitRefCollection { + continue + } + ops++ + if ops == 1 { + first = op } - return } - continue - } - - l.Debug("knot event subscription response", "code", resp.StatusCode, "url", u) - - if err := s.handleConnection(ctx, conn, sub); err != nil { - // TODO: measure the last N connection error times and if they're coming too fast reconnect slower or don't reconnect and wait for requestCrawl - l.Warn("host connection failed", "err", err, "backoff", backoff) - } - - updatedCursor := sub.LastSeq() - didProgress := updatedCursor > cursor - l.Debug("cursor compare", "cursor", cursor, "updatedCursor", updatedCursor, "didProgress", didProgress) - if cursor == 0 || didProgress { - cursor = updatedCursor - backoff = 0 - - batch := []models.HostCursor{sub.HostCursor()} - if err := db.StoreCursors(ctx, s.db, batch); err != nil { - l.Error("failed to store cursors", "err", err) + if ops == 0 { + return nil } - } - } -} - -// handleConnection handles websocket connection. -// Schedules task from received event and return when connection is closed -func (s *KnotSlurper) handleConnection(ctx context.Context, conn *websocket.Conn, sub *subscription) error { - // ping on every 30s - ctx, cancel := context.WithCancel(ctx) - defer cancel() // close the background ping job on connection close - go func() { - t := time.NewTicker(30 * time.Second) - defer t.Stop() - failcount := 0 - - for { - select { - case <-t.C: - if err := conn.WriteControl(websocket.PingMessage, []byte{}, time.Now().Add(time.Second*10)); err != nil { - s.logger.Warn("failed to ping", "err", err) - failcount++ - if failcount >= 4 { - s.logger.Error("too many ping fails", "count", failcount) - _ = conn.Close() - return - } - } else { - failcount = 0 // ok ping - } - case <-ctx.Done(): - _ = conn.Close() + firehoseRefOpsReceived.Add(float64(ops)) + sub.scheduler.AddTask(ctx, &Task{ + Key: msg.Commit.Repo, + repoDid: msg.Commit.Repo, + seq: msg.Commit.Seq, + op: first, + }) + return nil + }, + OutdatedReplay: func(ctx context.Context) int64 { + if _, err := s.desynchronizeHostRepos(ctx, host); err != nil { + l.Warn("couldn't mark repos desynchronized, retrying on the next outdated notice") + return sub.LastSeq() + } + return 0 + }, + OnConnectError: func(connectErr error) { + if seq := sub.LastSeq(); seq > highWater { + highWater = seq + connectFailures = 0 + } + connectFailures++ + l.Warn("dialing failed", "err", connectErr, "failures", connectFailures) + if connectFailures <= maxConnectFailures { return } - } - }() - - conn.SetPingHandler(func(message string) error { - err := conn.WriteControl(websocket.PongMessage, []byte(message), time.Now().Add(time.Minute)) - if err == websocket.ErrCloseSent { - return nil - } - return err - }) - conn.SetPongHandler(func(_ string) error { - if err := conn.SetReadDeadline(time.Now().Add(time.Minute)); err != nil { - s.logger.Error("failed to set read deadline", "err", err) - } - return nil - }) - - for { - select { - case <-ctx.Done(): - return ctx.Err() - default: - } - msgType, msg, err := conn.ReadMessage() - if err != nil { - return err - } - - if msgType != websocket.TextMessage { - continue - } - - sub.scheduler.AddTask(ctx, &Task{ - Key: sub.hostname, // TODO: replace to repository AT-URI for better concurrency - message: msg, - }) + l.Warn("host doesn't appear to be online, disabling for now") + host.Status = models.HostStatusOffline + host.LastSeq = sub.AppliedSeq() + if err := db.UpsertHost(ctx, s.db, &host); err != nil { + l.Error("failed to update host status", "err", err) + } + cancel() + }, + } + if err := consumer.Run(ctx); err != nil { + l.Error("firehose consumer stopped", "err", err) } } -type legacyGitRefUpdate struct { - OwnerDid *string `json:"ownerDid,omitempty"` - Repo *string `json:"repo,omitempty"` - LegacyRepoDid *string `json:"repoDid,omitempty"` -} - -type LegacyGitEvent struct { - Rkey string - Nsid string - Event legacyGitRefUpdate +func (s *KnotSlurper) desynchronizeHostRepos(ctx context.Context, host models.Host) (int64, error) { + n, err := db.DesynchronizeReposByKnotHost(ctx, s.db, &host) + if err != nil { + s.logger.With("host", host.Hostname).Error("failed to mark repos desynchronized", "err", err) + return 0, err + } + s.logger.With("host", host.Hostname).Warn("firehose cursor outdated, marked repos desynchronized", "repos", n) + return n, nil } func (s *KnotSlurper) ProcessEvent(ctx context.Context, task *Task) error { - var legacyMessage LegacyGitEvent - if err := json.Unmarshal(task.message, &legacyMessage); err != nil { - return fmt.Errorf("unmarshaling message: %w", err) - } - - if err := s.ProcessLegacyGitRefUpdate(ctx, task.Key, &legacyMessage); err != nil { - return fmt.Errorf("processing gitRefUpdate: %w", err) - } - return nil -} + l := s.logger.With("src", task.Key, "repo", task.repoDid) -// lookupRepoForRefUpdate resolves the local repo row for an incoming refUpdate -// via the stable RepoDid join. Returns (nil, "", nil) when the event has no -// repoDid (unjoinable) and (nil, key, nil) on a clean miss. -func (s *KnotSlurper) lookupRepoForRefUpdate(ctx context.Context, evt *LegacyGitEvent) (*models.Repo, string, error) { - raw := evt.Event.Repo - if raw == nil || *raw == "" { - raw = evt.Event.LegacyRepoDid + if task.repoDid == "" { + l.Warn("skipping ref op: commit frame has no repo did") + firehoseRefOpsSkipped.Inc() + return nil } - if raw == nil || *raw == "" { - return nil, "", nil + repoDid, err := syntax.ParseDID(task.repoDid) + if err != nil { + l.Warn("skipping ref op: commit frame has invalid repo did", "err", err) + firehoseRefOpsSkipped.Inc() + return nil } - repoDid := syntax.DID(*raw) - curr, err := db.GetRepoByRepoDid(ctx, s.db, repoDid) - return curr, repoDid.String(), err -} - -func (s *KnotSlurper) ProcessLegacyGitRefUpdate(ctx context.Context, source string, evt *LegacyGitEvent) error { - knotstreamEventsReceived.Inc() - - l := s.logger.With("src", source) - curr, lookupKey, err := s.lookupRepoForRefUpdate(ctx, evt) + curr, err := db.GetRepoByRepoDid(ctx, s.db, repoDid) if err != nil { - return fmt.Errorf("failed to get repo '%s': %w", lookupKey, err) + return fmt.Errorf("failed to get repo '%s': %w", task.repoDid, err) } if curr == nil { - if lookupKey == "" { - l.Warn("skipping gitRefUpdate: event has no fields to join on", - "repo", evt.Event.Repo, "legacy_repo_did", evt.Event.LegacyRepoDid) - } else { - // if repo doesn't exist in DB, just ignore the event. That repo is unknown. - // Hopefully crawler/tap will sync it later. - l.Warn("skipping event from unknown repo", "key", lookupKey) - } - knotstreamEventsSkipped.Inc() + l.Warn("skipping event from unknown repo", "key", task.repoDid) + firehoseRefOpsSkipped.Inc() return nil } l = l.With("repoAt", curr.AtUri()) - // TODO: should plan resync to resyncBuffer on RepoStateResyncing if curr.State != models.RepoStateActive { l.Debug("skipping non-active repo") - knotstreamEventsSkipped.Inc() + firehoseRefOpsSkipped.Inc() return nil } - if curr.GitRev != "" && evt.Rkey <= curr.GitRev.String() { - l.Debug("skipping replayed event", "event.Rkey", evt.Rkey, "currentRev", curr.GitRev) - knotstreamEventsSkipped.Inc() - return nil + var markErr error + for attempt := range markAttempts { + if attempt > 0 { + select { + case <-ctx.Done(): + return ctx.Err() + case <-time.After(markRetryDelay << attempt): + } + } + markErr = db.UpdateRepoState(ctx, s.db, curr.RepoDid, models.RepoStateDesynchronized) + if markErr == nil { + break + } } - - // can't skip anything, update repo state - if err := db.UpdateRepoState(ctx, s.db, curr.RepoDid, models.RepoStateDesynchronized); err != nil { - return err + if markErr != nil { + return fmt.Errorf("marking %s desynchronized: %w", curr.RepoDid, markErr) } - l.Info("event processed", "eventRev", evt.Rkey) + l.Info("event processed", "action", task.op.Action, "rkey", task.op.Rkey) - knotstreamEventsProcessed.Inc() + firehoseRefOpsProcessed.Inc() return nil } - -func userAgent() string { - return fmt.Sprintf("knotmirror/%s", versioninfo.Short()) -} - -func sleepForBackoff(b int) time.Duration { - if b == 0 { - return 0 - } - if b < 10 { - return time.Millisecond * time.Duration((50*b)+rand.Intn(500)) - } - return time.Second * 30 -} diff --git a/knotmirror/knotstream/subscription.go b/knotmirror/knotstream/subscription.go index 974e9e858..1231aef68 100644 --- a/knotmirror/knotstream/subscription.go +++ b/knotmirror/knotstream/subscription.go @@ -1,22 +1,40 @@ package knotstream -import "tangled.org/core/knotmirror/models" +import ( + "sync/atomic" + + "tangled.org/core/knotmirror/models" +) -// subscription represents websocket connection with that host type subscription struct { hostname string - // embedded parallel job scheduler + lastSeq atomic.Int64 + appliedSeq atomic.Int64 + scheduler *ParallelScheduler } func (s *subscription) LastSeq() int64 { - return s.scheduler.LastSeq() + return s.lastSeq.Load() +} + +func (s *subscription) AppliedSeq() int64 { + return s.appliedSeq.Load() +} + +func (s *subscription) MarkApplied(seq int64) { + for { + current := s.appliedSeq.Load() + if seq <= current || s.appliedSeq.CompareAndSwap(current, seq) { + return + } + } } func (s *subscription) HostCursor() models.HostCursor { return models.HostCursor{ Hostname: s.hostname, - LastSeq: s.LastSeq(), + LastSeq: s.AppliedSeq(), } } diff --git a/knotmirror/models/models.go b/knotmirror/models/models.go index 4833f527a..da21d292b 100644 --- a/knotmirror/models/models.go +++ b/knotmirror/models/models.go @@ -100,31 +100,3 @@ func (h *Host) URL() string { return fmt.Sprintf("https://%s", h.Hostname) } } - -func (h *Host) WsURL() string { - if h.NoSSL { - return fmt.Sprintf("ws://%s", h.Hostname) - } else { - return fmt.Sprintf("wss://%s", h.Hostname) - } -} - -// func (h *Host) SubscribeGitRefsURL(cursor int64) string { -// scheme := "wss" -// if h.NoSSL { -// scheme = "ws" -// } -// u := fmt.Sprintf("%s://%s/xrpc/%s", scheme, h.Hostname, tangled.SubscribeGitRefsNSID) -// if cursor > 0 { -// u = fmt.Sprintf("%s?cursor=%d", u, h.LastSeq) -// } -// return u -// } - -func (h *Host) LegacyEventsURL(cursor int64) string { - u := fmt.Sprintf("%s/events", h.WsURL()) - if cursor > 0 { - u = fmt.Sprintf("%s?cursor=%d", u, cursor) - } - return u -} diff --git a/knotmirror/readme.md b/knotmirror/readme.md index 723520d08..66f071280 100644 --- a/knotmirror/readme.md +++ b/knotmirror/readme.md @@ -2,7 +2,7 @@ KnotMirror is a git mirror service for all known repos. Heavily inspired by [indigo/relay] and [indigo/tap]. -KnotMirror syncs repo list using tap and subscribe to all known knots as KnotStream. +KnotMirror syncs repo list using tap and subscribes to every known knot's atproto firehose. [indigo/relay]: https://github.com/bluesky-social/indigo/tree/main/cmd/relay [indigo/tap]: https://github.com/bluesky-social/indigo/tree/main/cmd/tap diff --git a/knotmirror/repoindexer/language.go b/knotmirror/repoindexer/language.go index e5331399d..dbac64b6f 100644 --- a/knotmirror/repoindexer/language.go +++ b/knotmirror/repoindexer/language.go @@ -61,7 +61,6 @@ func NewIndexer(l *slog.Logger, cfg *config.Config, rdb *redis.Client) *Indexer func NewBackgroundIndexScheduler(l *slog.Logger, cfg *config.Config, e *sql.DB, indexer *Indexer) *knotstream.ParallelScheduler { return knotstream.NewParallelScheduler( 4, - "repo_stats_update", // NOTE: this is unused func(ctx context.Context, t *knotstream.Task) error { start := time.Now() repoId := syntax.DID(t.Key) diff --git a/knotmirror/resyncer.go b/knotmirror/resyncer.go index 966b04adc..449a6faf6 100644 --- a/knotmirror/resyncer.go +++ b/knotmirror/resyncer.go @@ -240,7 +240,7 @@ func (r *Resyncer) doResync(ctx context.Context, repoDid syntax.DID) (bool, erro r.knotBackoffMu.Unlock() return false, nil } - // TODO: suspend repo on 404. KnotStream updates will change the repo state back online + // TODO: suspend repo on 404. Firehose updates will change the repo state back online return false, fmt.Errorf("knot unreachable: %w", err) } @@ -260,7 +260,7 @@ func (r *Resyncer) doResync(ctx context.Context, repoDid syntax.DID) (bool, erro } // request index to zoekt server - // NOTE: indexing after full git resync is bad design. We are doing _after_ the sync because knotstream event doesn't include repository refs. + // NOTE: indexing after full git resync is bad design. We are doing _after_ the sync because the firehose's git.ref ops don't carry the repository's full ref set. // NOTE: and zoekt indexer should directly subscribe to the knot. remove this when we have knotrelay. if r.cfg.Search.ZoektUrl != "" { go func() { diff --git a/nix/modules/appview.nix b/nix/modules/appview.nix index 01f87bcc1..15124bf4d 100644 --- a/nix/modules/appview.nix +++ b/nix/modules/appview.nix @@ -93,73 +93,6 @@ in description = "Jetstream WebSocket endpoint"; }; }; - - # knotstream consumer configuration - knotstream = { - retryInterval = mkOption { - type = types.str; - default = "60s"; - description = "Initial retry interval for knotstream consumer"; - }; - - maxRetryInterval = mkOption { - type = types.str; - default = "120m"; - description = "Maximum retry interval for knotstream consumer"; - }; - - connectionTimeout = mkOption { - type = types.str; - default = "5s"; - description = "Connection timeout for knotstream consumer"; - }; - - workerCount = mkOption { - type = types.int; - default = 64; - description = "Number of workers for knotstream consumer"; - }; - - queueSize = mkOption { - type = types.int; - default = 100; - description = "Queue size for knotstream consumer"; - }; - }; - - # spindlestream consumer configuration - spindlestream = { - retryInterval = mkOption { - type = types.str; - default = "60s"; - description = "Initial retry interval for spindlestream consumer"; - }; - - maxRetryInterval = mkOption { - type = types.str; - default = "120m"; - description = "Maximum retry interval for spindlestream consumer"; - }; - - connectionTimeout = mkOption { - type = types.str; - default = "5s"; - description = "Connection timeout for spindlestream consumer"; - }; - - workerCount = mkOption { - type = types.int; - default = 64; - description = "Number of workers for spindlestream consumer"; - }; - - queueSize = mkOption { - type = types.int; - default = 100; - description = "Queue size for spindlestream consumer"; - }; - }; - # resend configuration resend = { sentFrom = mkOption { @@ -339,18 +272,6 @@ in TANGLED_JETSTREAM_ENDPOINT = cfg.jetstream.endpoint; - TANGLED_KNOTSTREAM_RETRY_INTERVAL = cfg.knotstream.retryInterval; - TANGLED_KNOTSTREAM_MAX_RETRY_INTERVAL = cfg.knotstream.maxRetryInterval; - TANGLED_KNOTSTREAM_CONNECTION_TIMEOUT = cfg.knotstream.connectionTimeout; - TANGLED_KNOTSTREAM_WORKER_COUNT = toString cfg.knotstream.workerCount; - TANGLED_KNOTSTREAM_QUEUE_SIZE = toString cfg.knotstream.queueSize; - - TANGLED_SPINDLESTREAM_RETRY_INTERVAL = cfg.spindlestream.retryInterval; - TANGLED_SPINDLESTREAM_MAX_RETRY_INTERVAL = cfg.spindlestream.maxRetryInterval; - TANGLED_SPINDLESTREAM_CONNECTION_TIMEOUT = cfg.spindlestream.connectionTimeout; - TANGLED_SPINDLESTREAM_WORKER_COUNT = toString cfg.spindlestream.workerCount; - TANGLED_SPINDLESTREAM_QUEUE_SIZE = toString cfg.spindlestream.queueSize; - TANGLED_RESEND_SENT_FROM = cfg.resend.sentFrom; TANGLED_POSTHOG_ENDPOINT = cfg.posthog.endpoint;