From 929c6330c26498011e926e1338b12323fd4f40f6 Mon Sep 17 00:00:00 2001 From: Lewis Date: Mon, 15 Jun 2026 10:17:48 +0300 Subject: [PATCH] appview/knotacl: replace cache poll w/ event-driven roster Lewis: May this revision serve well! --- appview/db/db.go | 22 ++ appview/knotacl/cache.go | 134 --------- appview/knotacl/cache_test.go | 226 --------------- appview/knotacl/reader.go | 2 +- appview/knotacl/roster.go | 435 ++++++++++++++++++++++++++++ appview/knotacl/roster_test.go | 486 ++++++++++++++++++++++++++++++++ appview/knotacl/service.go | 21 +- appview/knotacl/service_test.go | 29 +- 8 files changed, 987 insertions(+), 368 deletions(-) delete mode 100644 appview/knotacl/cache_test.go create mode 100644 appview/knotacl/roster.go create mode 100644 appview/knotacl/roster_test.go diff --git a/appview/db/db.go b/appview/db/db.go index c85aa3d4..7a218e20 100644 --- a/appview/db/db.go +++ b/appview/db/db.go @@ -2218,6 +2218,28 @@ func Make(ctx context.Context, dbPath string) (*DB, error) { return err }) + orm.RunMigration(conn, logger, "add-knotacl-sync-table", func(tx *sql.Tx) error { + _, err := tx.Exec(` + create table if not exists knotacl_sync ( + scope_key text primary key, + synced_at text not null + ); + `) + 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 + }) + return &DB{ db, logger, diff --git a/appview/knotacl/cache.go b/appview/knotacl/cache.go index cb764dfd..64a6aa80 100644 --- a/appview/knotacl/cache.go +++ b/appview/knotacl/cache.go @@ -2,18 +2,8 @@ package knotacl import ( "context" - "maps" "net/http" - "slices" "sync" - "time" - - "golang.org/x/sync/singleflight" -) - -const ( - cacheTTL = 15 * time.Second - cacheMaxEntries = 4096 ) type lister interface { @@ -21,130 +11,6 @@ type lister interface { GetRepoCollaborators(ctx context.Context, host, repoDid string) ([]string, error) } -type cacheEntry struct { - subjects []string - storedAt time.Time -} - -type cache struct { - inner lister - ttl time.Duration - now func() time.Time - - mu sync.Mutex - entries map[string]cacheEntry - group singleflight.Group -} - -func newCache(inner lister, ttl time.Duration, now func() time.Time) *cache { - if now == nil { - now = time.Now - } - return &cache{inner: inner, ttl: ttl, now: now, entries: map[string]cacheEntry{}} -} - -func (c *cache) GetKnotMembers(ctx context.Context, host string) ([]string, error) { - return c.fetch(ctx, memberCacheKey(host), func() ([]string, error) { - return c.inner.GetKnotMembers(ctx, host) - }) -} - -func (c *cache) GetRepoCollaborators(ctx context.Context, host, repoDid string) ([]string, error) { - return c.fetch(ctx, collabCacheKey(host, repoDid), func() ([]string, error) { - return c.inner.GetRepoCollaborators(ctx, host, repoDid) - }) -} - -func memberCacheKey(host string) string { return "m\x00" + host } - -func collabCacheKey(host, repoDid string) string { return "c\x00" + host + "\x00" + repoDid } - -func (c *cache) InvalidateMembers(host string) { - c.forget(memberCacheKey(host)) -} - -func (c *cache) InvalidateCollaborators(host, repoDid string) { - c.forget(collabCacheKey(host, repoDid)) -} - -func (c *cache) forget(key string) { - c.mu.Lock() - defer c.mu.Unlock() - delete(c.entries, key) -} - -func (c *cache) fetch(ctx context.Context, key string, load func() ([]string, error)) ([]string, error) { - if memo := memoFrom(ctx); memo != nil { - if v, ok := memo.get(key); ok { - return slices.Clone(v), nil - } - } - v, err := c.load(key, load) - if err != nil { - return nil, err - } - if memo := memoFrom(ctx); memo != nil { - memo.put(key, v) - } - return slices.Clone(v), nil -} - -func (c *cache) load(key string, load func() ([]string, error)) ([]string, error) { - if v, ok := c.lookup(key); ok { - return v, nil - } - v, err, _ := c.group.Do(key, func() (any, error) { - if v, ok := c.lookup(key); ok { - return v, nil - } - fresh, err := load() - if err != nil { - return nil, err - } - c.store(key, fresh) - return fresh, nil - }) - if err != nil { - return nil, err - } - return v.([]string), nil -} - -func (c *cache) lookup(key string) ([]string, bool) { - c.mu.Lock() - defer c.mu.Unlock() - e, ok := c.entries[key] - if !ok || c.now().Sub(e.storedAt) >= c.ttl { - return nil, false - } - return e.subjects, true -} - -func (c *cache) store(key string, subjects []string) { - c.mu.Lock() - defer c.mu.Unlock() - if _, exists := c.entries[key]; !exists && len(c.entries) >= cacheMaxEntries { - maps.DeleteFunc(c.entries, func(_ string, e cacheEntry) bool { - return c.now().Sub(e.storedAt) >= c.ttl - }) - c.evictOldestLocked() - } - c.entries[key] = cacheEntry{subjects: subjects, storedAt: c.now()} -} - -func (c *cache) evictOldestLocked() { - oldestKey := "" - var oldestAt time.Time - for k, e := range c.entries { - if oldestKey == "" || e.storedAt.Before(oldestAt) { - oldestKey, oldestAt = k, e.storedAt - } - } - if len(c.entries) >= cacheMaxEntries && oldestKey != "" { - delete(c.entries, oldestKey) - } -} - type requestMemo struct { mu sync.Mutex entries map[string][]string diff --git a/appview/knotacl/cache_test.go b/appview/knotacl/cache_test.go deleted file mode 100644 index 8e74f35b..00000000 --- a/appview/knotacl/cache_test.go +++ /dev/null @@ -1,226 +0,0 @@ -package knotacl - -import ( - "context" - "errors" - "fmt" - "slices" - "sync" - "testing" - "time" -) - -var cacheTestBase = time.Unix(1700000000, 0) - -type fakeLister struct { - mu sync.Mutex - memberCalls int - members []string - err error - started chan struct{} - block chan struct{} -} - -func (f *fakeLister) GetKnotMembers(ctx context.Context, host string) ([]string, error) { - f.mu.Lock() - f.memberCalls++ - members, err, started, block := f.members, f.err, f.started, f.block - f.mu.Unlock() - if started != nil { - close(started) - } - if block != nil { - <-block - } - if err != nil { - return nil, err - } - return members, nil -} - -func (f *fakeLister) GetRepoCollaborators(ctx context.Context, host, repoDid string) ([]string, error) { - f.mu.Lock() - defer f.mu.Unlock() - return f.members, f.err -} - -func (f *fakeLister) calls() int { - f.mu.Lock() - defer f.mu.Unlock() - return f.memberCalls -} - -func (f *fakeLister) set(members []string, err error) { - f.mu.Lock() - defer f.mu.Unlock() - f.members, f.err = members, err -} - -type fakeClock struct { - mu sync.Mutex - t time.Time -} - -func (c *fakeClock) now() time.Time { - c.mu.Lock() - defer c.mu.Unlock() - return c.t -} - -func (c *fakeClock) advance(d time.Duration) { - c.mu.Lock() - defer c.mu.Unlock() - c.t = c.t.Add(d) -} - -func TestCache_TTLCollapsesThenExpires(t *testing.T) { - clk := &fakeClock{t: cacheTestBase} - f := &fakeLister{members: []string{"did:plc:boltless"}} - c := newCache(f, cacheTTL, clk.now) - ctx := context.Background() - - if _, err := c.GetKnotMembers(ctx, "knot.nel.pet"); err != nil { - t.Fatal(err) - } - if _, err := c.GetKnotMembers(ctx, "knot.nel.pet"); err != nil { - t.Fatal(err) - } - if f.calls() != 1 { - t.Errorf("memberCalls=%d, want 1 within the TTL window", f.calls()) - } - - clk.advance(cacheTTL) - if _, err := c.GetKnotMembers(ctx, "knot.nel.pet"); err != nil { - t.Fatal(err) - } - if f.calls() != 2 { - t.Errorf("memberCalls=%d, want 2 once the entry expired", f.calls()) - } -} - -func TestCache_ErrorsNotCached(t *testing.T) { - clk := &fakeClock{t: cacheTestBase} - f := &fakeLister{err: errors.New("knot unreachable")} - c := newCache(f, cacheTTL, clk.now) - ctx := context.Background() - - if _, err := c.GetKnotMembers(ctx, "knot.nel.pet"); err == nil { - t.Fatal("want error on the first call") - } - f.set([]string{"did:plc:boltless"}, nil) - got, err := c.GetKnotMembers(ctx, "knot.nel.pet") - if err != nil { - t.Fatal(err) - } - if !slices.Equal(got, []string{"did:plc:boltless"}) { - t.Errorf("got %v after recovery, want the live value", got) - } - if f.calls() != 2 { - t.Errorf("memberCalls=%d, want 2; a failed fetch must not be cached", f.calls()) - } -} - -func TestCache_MemoShortCircuitsWithinRequest(t *testing.T) { - clk := &fakeClock{t: cacheTestBase} - f := &fakeLister{members: []string{"did:plc:boltless"}} - c := newCache(f, cacheTTL, clk.now) - ctx := WithMemo(context.Background()) - - first, err := c.GetKnotMembers(ctx, "knot.nel.pet") - if err != nil { - t.Fatal(err) - } - - clk.advance(2 * cacheTTL) - f.set([]string{"did:plc:akshay"}, nil) - - second, err := c.GetKnotMembers(ctx, "knot.nel.pet") - if err != nil { - t.Fatal(err) - } - if !slices.Equal(first, second) { - t.Errorf("memo must hold one snapshot per request: first=%v second=%v", first, second) - } - if f.calls() != 1 { - t.Errorf("memberCalls=%d, want 1; the request memo must not re-query even past the TTL", f.calls()) - } -} - -func TestCache_ReturnedSliceCannotCorruptCache(t *testing.T) { - clk := &fakeClock{t: cacheTestBase} - f := &fakeLister{members: []string{"did:plc:boltless", "did:plc:akshay"}} - c := newCache(f, cacheTTL, clk.now) - ctx := context.Background() - - got, err := c.GetKnotMembers(ctx, "knot.nel.pet") - if err != nil { - t.Fatal(err) - } - for i := range got { - got[i] = "did:plc:squid" - } - - again, err := c.GetKnotMembers(ctx, "knot.nel.pet") - if err != nil { - t.Fatal(err) - } - if slices.Contains(again, "did:plc:squid") { - t.Errorf("a caller mutating its returned slice corrupted the cached entry: %v", again) - } - if f.calls() != 1 { - t.Errorf("memberCalls=%d, want 1; the second read should be served from cache", f.calls()) - } -} - -func TestCache_SingleflightCollapsesConcurrentMisses(t *testing.T) { - clk := &fakeClock{t: cacheTestBase} - started := make(chan struct{}) - release := make(chan struct{}) - f := &fakeLister{members: []string{"did:plc:boltless"}, started: started, block: release} - c := newCache(f, cacheTTL, clk.now) - ctx := context.Background() - - var wg sync.WaitGroup - call := func() { - wg.Add(1) - go func() { - defer wg.Done() - if _, err := c.GetKnotMembers(ctx, "knot.nel.pet"); err != nil { - t.Errorf("GetKnotMembers: %v", err) - } - }() - } - - call() - <-started - for range make([]struct{}, 8) { - call() - } - time.Sleep(20 * time.Millisecond) - close(release) - wg.Wait() - - if f.calls() != 1 { - t.Errorf("memberCalls=%d, want 1; concurrent misses must collapse into a single knot query", f.calls()) - } -} - -func TestCache_CapIsHardUnderFreshFlood(t *testing.T) { - clk := &fakeClock{t: cacheTestBase} - f := &fakeLister{members: []string{"did:plc:limpet"}} - c := newCache(f, cacheTTL, clk.now) - ctx := context.Background() - - for i := 0; i < cacheMaxEntries+100; i++ { - if _, err := c.GetKnotMembers(ctx, fmt.Sprintf("knot-%d.nel.pet", i)); err != nil { - t.Fatal(err) - } - } - - c.mu.Lock() - n := len(c.entries) - c.mu.Unlock() - if n > cacheMaxEntries { - t.Errorf("entries=%d, want <= %d; all-fresh keys must not grow past the cap", n, cacheMaxEntries) - } -} diff --git a/appview/knotacl/reader.go b/appview/knotacl/reader.go index e8aa6c7d..33c10689 100644 --- a/appview/knotacl/reader.go +++ b/appview/knotacl/reader.go @@ -65,7 +65,7 @@ func (r *legacyReader) isKnotMember(ctx context.Context, host, userDid string) b } type nativeReader struct { - client *cache + client *roster execer db.Execer } diff --git a/appview/knotacl/roster.go b/appview/knotacl/roster.go new file mode 100644 index 00000000..e82ca716 --- /dev/null +++ b/appview/knotacl/roster.go @@ -0,0 +1,435 @@ +package knotacl + +import ( + "context" + "database/sql" + "errors" + "fmt" + "log/slog" + "slices" + "sync" + "sync/atomic" + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" + "golang.org/x/sync/singleflight" + + "tangled.org/core/appview/db" + "tangled.org/core/appview/models" + "tangled.org/core/orm" +) + +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 +} + +type roster struct { + store *db.DB + src lister + ttl time.Duration + now func() time.Time + log *slog.Logger + group singleflight.Group + + mu sync.Mutex + scopes map[string]*scopeState +} + +func newRoster(store *db.DB, src lister, ttl time.Duration, now func() time.Time, logger *slog.Logger) *roster { + if now == nil { + now = time.Now + } + if logger == nil { + logger = slog.Default() + } + return &roster{ + store: store, + src: src, + ttl: ttl, + now: now, + log: logger, + scopes: map[string]*scopeState{}, + } +} + +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() ([]string, error) { + rows, err := db.GetKnotMembers(r.store, orm.FilterEq("domain", host)) + if err != nil { + return nil, err + } + return mapSlice(rows, func(km models.KnotMember) string { return km.Subject.String() }), nil + }, + ) +} + +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() ([]string, error) { + rows, err := db.GetCollaborators(r.store, orm.FilterEq("repo_did", repoDid)) + if err != nil { + return nil, err + } + return mapSlice(rows, func(c models.Collaborator) string { return c.SubjectDid.String() }), nil + }, + ) +} + +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)) +} + +func (r *roster) InvalidateCollaborators(host, repoDid string) { + _ = clearSyncedAt(r.store, collabScope(repoDid)) +} + +func (r *roster) serve(ctx context.Context, key string, reconcile func() error, read func() ([]string, error)) ([]string, error) { + if memo := memoFrom(ctx); memo != nil { + if v, ok := memo.get(key); ok { + return slices.Clone(v), nil + } + } + + recErr := r.maybeReconcile(key, reconcile) + + subjects, err := read() + if err != nil { + return nil, err + } + + if len(subjects) == 0 && recErr != nil && !r.everSynced(key) { + return nil, fmt.Errorf("%w: %v", ErrKnotUnreachable, recErr) + } + + subjects = dedup(subjects) + if memo := memoFrom(ctx); memo != nil { + memo.put(key, subjects) + } + return slices.Clone(subjects), nil +} + +func (r *roster) maybeReconcile(key string, reconcile func() error) error { + if r.fresh(key) { + return nil + } + if r.backingOff(key) { + return errReconcileBackoff + } + _, err, _ := r.group.Do(key, func() (any, error) { + if r.fresh(key) { + return nil, nil + } + if r.backingOff(key) { + return nil, errReconcileBackoff + } + return nil, reconcile() + }) + 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) + st := r.acquire(scope) + defer r.release(scope, st) + + genBefore := st.gen.Load() + subjects, err := r.src.GetRepoCollaborators(ctx, host, repoDid) + 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 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) + continue + } + if err := db.AddCollaborator(tx, models.Collaborator{SubjectDid: did, RepoDid: repo}); err != nil { + return err + } + } + return nil + }) +} + +func (r *roster) commitReconcile(st *scopeState, scope string, genBefore uint64, 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 + } + defer tx.Rollback() + if err := replace(tx); err != nil { + return err + } + 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 + } + r.clearFailed(st) + return nil +} + +func (r *roster) acquire(scope string) *scopeState { + r.mu.Lock() + defer r.mu.Unlock() + st := r.scopes[scope] + if st == nil { + st = &scopeState{} + r.scopes[scope] = st + } + st.refs++ + return st +} + +func (r *roster) release(scope string, st *scopeState) { + r.mu.Lock() + defer r.mu.Unlock() + st.refs-- + if st.refs == 0 && st.failedAt.IsZero() { + delete(r.scopes, scope) + } +} + +func (r *roster) markFailed(st *scopeState) { + r.mu.Lock() + defer r.mu.Unlock() + st.failedAt = r.now() +} + +func (r *roster) clearFailed(st *scopeState) { + r.mu.Lock() + defer r.mu.Unlock() + st.failedAt = time.Time{} +} + +func (r *roster) backingOff(scope string) bool { + r.mu.Lock() + defer r.mu.Unlock() + st, ok := r.scopes[scope] + if !ok { + return false + } + return !st.failedAt.IsZero() && r.now().Sub(st.failedAt) < reconcileBackoff +} + +func (r *roster) fresh(key string) bool { + at, ok, err := getSyncedAt(r.store, key) + if err != nil || !ok { + return false + } + return r.now().Sub(at) < r.ttl +} + +func (r *roster) everSynced(key string) bool { + _, ok, _ := getSyncedAt(r.store, key) + return ok +} + +func memberScope(host string) string { return "m\x00" + host } + +func collabScope(repoDid string) string { return "c\x00" + repoDid } + +func getSyncedAt(e db.Execer, key string) (time.Time, bool, error) { + var raw string + err := e.QueryRow(`select synced_at from knotacl_sync where scope_key = ?`, key).Scan(&raw) + if errors.Is(err, sql.ErrNoRows) { + return time.Time{}, false, nil + } + if err != nil { + return time.Time{}, false, err + } + at, err := time.Parse(time.RFC3339, raw) + if err != nil { + return time.Time{}, false, err + } + return at, true, nil +} + +func setSyncedAt(e db.Execer, key string, at time.Time) error { + _, err := e.Exec( + `insert into knotacl_sync (scope_key, synced_at) values (?, ?) + on conflict(scope_key) do update set synced_at = excluded.synced_at`, + key, at.UTC().Format(time.RFC3339), + ) + return err +} + +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 new file mode 100644 index 00000000..c74eabae --- /dev/null +++ b/appview/knotacl/roster_test.go @@ -0,0 +1,486 @@ +package knotacl + +import ( + "context" + "errors" + "path/filepath" + "slices" + "sync" + "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) + +const rosterTTL = 5 * time.Minute + +type fakeLister struct { + mu sync.Mutex + memberCalls int + members []string + err error + started chan struct{} + block chan struct{} +} + +func (f *fakeLister) GetKnotMembers(ctx context.Context, host string) ([]string, error) { + f.mu.Lock() + f.memberCalls++ + members, err, started, block := f.members, f.err, f.started, f.block + f.mu.Unlock() + if started != nil { + close(started) + } + if block != nil { + <-block + } + if err != nil { + return nil, err + } + return members, nil +} + +func (f *fakeLister) GetRepoCollaborators(ctx context.Context, host, repoDid string) ([]string, error) { + f.mu.Lock() + defer f.mu.Unlock() + return f.members, f.err +} + +func (f *fakeLister) calls() int { + f.mu.Lock() + defer f.mu.Unlock() + return f.memberCalls +} + +func (f *fakeLister) set(members []string, err error) { + f.mu.Lock() + defer f.mu.Unlock() + f.members, f.err = members, err +} + +func (f *fakeLister) arm(started, block chan struct{}) { + f.mu.Lock() + defer f.mu.Unlock() + f.started, f.block = started, block +} + +type fakeClock struct { + mu sync.Mutex + t time.Time +} + +func (c *fakeClock) now() time.Time { + c.mu.Lock() + defer c.mu.Unlock() + return c.t +} + +func (c *fakeClock) advance(d time.Duration) { + c.mu.Lock() + defer c.mu.Unlock() + c.t = c.t.Add(d) +} + +func rosterTestDB(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 +} + +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"}} + r := newRoster(rosterTestDB(t), f, rosterTTL, clk.now, nil) + ctx := context.Background() + + got, err := r.GetKnotMembers(ctx, rosterHost) + if err != nil { + t.Fatal(err) + } + if !slices.Equal(got, []string{"did:plc:boltless"}) { + t.Errorf("bootstrap read = %v, want the drained roster", got) + } + + if _, err := r.GetKnotMembers(ctx, rosterHost); err != nil { + t.Fatal(err) + } + if f.calls() != 1 { + t.Errorf("memberCalls=%d, want 1; a read within the reconcile TTL must be served from sqlite", f.calls()) + } + + clk.advance(rosterTTL) + if _, err := r.GetKnotMembers(ctx, rosterHost); err != nil { + t.Fatal(err) + } + if f.calls() != 2 { + t.Errorf("memberCalls=%d, want 2; a stale scope must reconcile from XRPC", f.calls()) + } +} + +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")} + r := newRoster(rosterTestDB(t), f, rosterTTL, clk.now, nil) + + if _, err := r.GetKnotMembers(context.Background(), rosterHost); !errors.Is(err, ErrKnotUnreachable) { + t.Fatalf("err=%v, want ErrKnotUnreachable when cold and the bootstrap drain fails", err) + } +} + +func TestRoster_StaleServedWhenUnreachable(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(2 * rosterTTL) + f.set(nil, errors.New("knot unreachable")) + + got, err := r.GetKnotMembers(ctx, rosterHost) + if err != nil { + t.Fatalf("err=%v, want stale rows served when a once-synced scope goes unreachable", err) + } + if !slices.Equal(got, []string{"did:plc:boltless"}) { + t.Errorf("stale read = %v, want the last good roster", got) + } +} + +func TestRoster_InvalidateForcesReconcile(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) + } + r.InvalidateMembers(rosterHost) + + f.set([]string{"did:plc:akshay"}, nil) + got, err := r.GetKnotMembers(ctx, rosterHost) + if err != nil { + t.Fatal(err) + } + if !slices.Equal(got, []string{"did:plc:akshay"}) { + t.Errorf("post-invalidate read = %v, want a fresh reconcile within the TTL", got) + } + if f.calls() != 2 { + t.Errorf("memberCalls=%d, want 2; invalidation must force the next read to reconcile", f.calls()) + } +} + +func TestRoster_SingleflightCollapsesConcurrentColdReads(t *testing.T) { + clk := &fakeClock{t: rosterTestBase} + started := make(chan struct{}) + release := make(chan struct{}) + f := &fakeLister{members: []string{"did:plc:boltless"}, started: started, block: release} + r := newRoster(rosterTestDB(t), f, rosterTTL, clk.now, nil) + ctx := context.Background() + + var wg sync.WaitGroup + call := func() { + wg.Add(1) + go func() { + defer wg.Done() + if _, err := r.GetKnotMembers(ctx, rosterHost); err != nil { + t.Errorf("GetKnotMembers: %v", err) + } + }() + } + + call() + <-started + for range make([]struct{}, 8) { + call() + } + time.Sleep(20 * time.Millisecond) + close(release) + wg.Wait() + + if f.calls() != 1 { + t.Errorf("memberCalls=%d, want 1; concurrent cold reads must collapse into a single reconcile", f.calls()) + } +} + +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"}} + r := newRoster(rosterTestDB(t), f, rosterTTL, clk.now, nil) + ctx := context.Background() + + if _, err := r.GetKnotMembers(ctx, rosterHost); err != nil { + t.Fatal(err) + } + if f.calls() != 1 { + t.Fatalf("memberCalls=%d, want 1 after bootstrap", f.calls()) + } + + clk.advance(rosterTTL) + f.set(nil, errors.New("knot unreachable")) + + if got, err := r.GetKnotMembers(ctx, rosterHost); err != nil || !slices.Equal(got, []string{"did:plc:boltless"}) { + t.Fatalf("got %v err %v, want the stale roster served", got, err) + } + if f.calls() != 2 { + t.Fatalf("memberCalls=%d, want 2 after the first stale drain", f.calls()) + } + + if got, err := r.GetKnotMembers(ctx, rosterHost); err != nil || !slices.Equal(got, []string{"did:plc:boltless"}) { + t.Fatalf("got %v err %v, want the stale roster served", got, err) + } + if f.calls() != 2 { + t.Errorf("memberCalls=%d, want 2; a down knot must not be re-probed within the backoff window", f.calls()) + } + + clk.advance(reconcileBackoff) + if _, err := r.GetKnotMembers(ctx, rosterHost); err != nil { + t.Fatal(err) + } + if f.calls() != 3 { + t.Errorf("memberCalls=%d, want 3; the probe must resume once the backoff window elapses", f.calls()) + } +} + +func TestRoster_BackoffStillErrorsWhenCold(t *testing.T) { + clk := &fakeClock{t: rosterTestBase} + f := &fakeLister{err: errors.New("knot unreachable")} + r := newRoster(rosterTestDB(t), f, rosterTTL, clk.now, nil) + ctx := context.Background() + + if _, err := r.GetKnotMembers(ctx, rosterHost); !errors.Is(err, ErrKnotUnreachable) { + t.Fatalf("err=%v, want ErrKnotUnreachable on the cold drain failure", err) + } + if f.calls() != 1 { + t.Fatalf("memberCalls=%d, want 1", f.calls()) + } + + if _, err := r.GetKnotMembers(ctx, rosterHost); !errors.Is(err, ErrKnotUnreachable) { + t.Fatalf("err=%v, want ErrKnotUnreachable during backoff for a cold scope", err) + } + if f.calls() != 1 { + t.Errorf("memberCalls=%d, want 1; a cold scope must not be re-probed within the backoff window", f.calls()) + } +} diff --git a/appview/knotacl/service.go b/appview/knotacl/service.go index 8060d756..d91b9454 100644 --- a/appview/knotacl/service.go +++ b/appview/knotacl/service.go @@ -8,6 +8,7 @@ import ( "sync" "time" + "github.com/bluesky-social/indigo/atproto/syntax" "golang.org/x/sync/errgroup" "tangled.org/core/appview/db" @@ -33,12 +34,12 @@ type Service struct { nat *nativeReader } -func NewService(enforcer *rbac.Enforcer, execer db.Execer, dev bool, logger *slog.Logger) *Service { +func NewService(enforcer *rbac.Enforcer, store *db.DB, dev bool, logger *slog.Logger) *Service { return &Service{ dev: dev, log: logger, leg: &legacyReader{enforcer: enforcer}, - nat: &nativeReader{client: newCache(NewClient(dev, logger), cacheTTL, nil), execer: execer}, + nat: &nativeReader{client: newRoster(store, NewClient(dev, logger), reconcileTTL, nil, logger), execer: store}, } } @@ -89,6 +90,22 @@ 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/knotacl/service_test.go b/appview/knotacl/service_test.go index 87a2cff7..e9224515 100644 --- a/appview/knotacl/service_test.go +++ b/appview/knotacl/service_test.go @@ -105,6 +105,20 @@ func testRepo(host string) *models.Repo { return &models.Repo{Did: testOwner, Knot: host, RepoDid: testRepoDid, Name: "anemone"} } +func seedRepoRow(t *testing.T, d *db.DB, repo *models.Repo) { + t.Helper() + tx, err := d.Begin() + if err != nil { + t.Fatalf("begin: %v", err) + } + if err := db.AddRepo(tx, repo); err != nil { + t.Fatalf("AddRepo: %v", err) + } + if err := tx.Commit(); err != nil { + t.Fatalf("commit: %v", err) + } +} + func seedRepoPolicies(t *testing.T, e *rbac.Enforcer, host string) { t.Helper() if err := e.AddRepo(testOwner, host, testRepoDid); err != nil { @@ -141,7 +155,8 @@ func TestService_OldKnotUsesCasbinNoLiveQuery(t *testing.T) { func TestService_ParityOldVsNew(t *testing.T) { ctx := context.Background() oldSvc, _, oldHost := newServiceEnv(t, &fakeKnot{version: "v1.14.0"}, func(e *rbac.Enforcer, h string) { seedRepoPolicies(t, e, h) }) - newSvc, _, newHost := newServiceEnv(t, &fakeKnot{version: "v1.15.0", capabilities: capsKnotACL, collaborators: []string{testCollab}}, nil) + newSvc, newDb, newHost := newServiceEnv(t, &fakeKnot{version: "v1.15.0", capabilities: capsKnotACL, collaborators: []string{testCollab}}, nil) + seedRepoRow(t, newDb, testRepo(newHost)) for _, did := range []string{testOwner, testCollab} { oldRoles := sortedRoles(oldSvc.RolesInRepo(ctx, testRepo(oldHost), did).Roles) @@ -167,7 +182,8 @@ func TestService_NewKnotOwnerFromRecord(t *testing.T) { func TestService_NewKnotCollaboratorFromList(t *testing.T) { ctx := context.Background() - svc, _, host := newServiceEnv(t, &fakeKnot{version: "v1.15.0", capabilities: capsKnotACL, collaborators: []string{testCollab}}, nil) + svc, d, host := newServiceEnv(t, &fakeKnot{version: "v1.15.0", capabilities: capsKnotACL, collaborators: []string{testCollab}}, nil) + seedRepoRow(t, d, testRepo(host)) collab := svc.RolesInRepo(ctx, testRepo(host), testCollab) if !collab.IsCollaborator() || !collab.IsPushAllowed() { @@ -181,7 +197,8 @@ func TestService_NewKnotCollaboratorFromList(t *testing.T) { func TestService_MixedFleet(t *testing.T) { ctx := context.Background() oldSvc, _, oldHost := newServiceEnv(t, &fakeKnot{version: "v1.14.0"}, func(e *rbac.Enforcer, h string) { seedRepoPolicies(t, e, h) }) - newSvc, _, newHost := newServiceEnv(t, &fakeKnot{version: "v1.15.0", capabilities: capsKnotACL, collaborators: []string{testCollab}}, nil) + newSvc, newDb, newHost := newServiceEnv(t, &fakeKnot{version: "v1.15.0", capabilities: capsKnotACL, collaborators: []string{testCollab}}, nil) + seedRepoRow(t, newDb, testRepo(newHost)) if !newSvc.RolesInRepo(ctx, testRepo(newHost), testCollab).IsCollaborator() { t.Error("new-knot collaborator must resolve from the live query with an empty casbin") @@ -307,7 +324,8 @@ func TestService_KnotMembersIncludesOwner(t *testing.T) { func TestService_CollaboratorsNewKnot(t *testing.T) { ctx := context.Background() - svc, _, host := newServiceEnv(t, &fakeKnot{version: "v1.15.0", capabilities: capsKnotACL, collaborators: []string{testCollab}}, nil) + svc, d, host := newServiceEnv(t, &fakeKnot{version: "v1.15.0", capabilities: capsKnotACL, collaborators: []string{testCollab}}, nil) + seedRepoRow(t, d, testRepo(host)) collabs := svc.Collaborators(ctx, testRepo(host)) if len(collabs) != 2 { t.Fatalf("Collaborators = %v, want owner + one collaborator", collabs) @@ -319,7 +337,8 @@ func TestService_CollaboratorsNewKnot(t *testing.T) { t.Errorf("second row = %v, want the collaborator", collabs[1]) } - dupSvc, _, dupHost := newServiceEnv(t, &fakeKnot{version: "v1.15.0", capabilities: capsKnotACL, collaborators: []string{testCollab, testOwner}}, nil) + dupSvc, dupDb, dupHost := newServiceEnv(t, &fakeKnot{version: "v1.15.0", capabilities: capsKnotACL, collaborators: []string{testCollab, testOwner}}, nil) + seedRepoRow(t, dupDb, testRepo(dupHost)) rows := dupSvc.Collaborators(ctx, testRepo(dupHost)) ownerRows := 0 for _, c := range rows { -- 2.51.2