diff --git a/jetstream.go b/jetstream.go index 526f1e2..bb2e538 100644 --- a/jetstream.go +++ b/jetstream.go @@ -215,6 +215,16 @@ func applyRepo(ctx context.Context, st *store, knots KnotConsumer, hostname stri if err := json.Unmarshal(c.Record, &rec); err != nil { return fmt.Errorf("decode repo: %w", err) } + + // Capture the prior (knot, spindle) before the upsert so the + // post-mutation reconcile below can detect transitions like + // "repo used to point at us, no longer does" — which would + // otherwise leave a knot subscription dangling. + oldKnot, oldSpindle, err := st.GetRepo(ctx, did, c.RKey) + if err != nil { + return err + } + if err := st.UpsertRepo(ctx, did, c.RKey, rec.Knot, rec.Name, deref(rec.Spindle), deref(rec.RepoDid), @@ -223,21 +233,79 @@ func applyRepo(ctx context.Context, st *store, knots KnotConsumer, hostname stri return err } - // If this repo just declared us as its spindle, start (or - // continue) listening to its knot for pipeline triggers. The - // knot consumer dedupes on its own so this is safe to call - // even on update events that don't change the spindle field. - if knots != nil && rec.Spindle != nil && *rec.Spindle == hostname && rec.Knot != "" { - knots.AddKnot(ctx, rec.Knot) + newSpindle := deref(rec.Spindle) + return reconcileKnot(ctx, st, knots, hostname, + oldKnot, oldSpindle, + rec.Knot, newSpindle, + ) + + case jsOpDelete: + // Same shape as the update path, just with no "new" side: we + // have to read the row out before deleting so we can decide + // whether deletion freed the last hold on a knot we'd been + // subscribed to. + oldKnot, oldSpindle, err := st.GetRepo(ctx, did, c.RKey) + if err != nil { + return err + } + if err := st.DeleteRepo(ctx, did, c.RKey); err != nil { + return err } + return reconcileKnot(ctx, st, knots, hostname, + oldKnot, oldSpindle, + "", "", + ) + } + return nil +} +// reconcileKnot brings the knot consumer's subscriptions in line with +// the latest store state after a single repo mutation. It is called +// after the mutation has been applied so IsKnotWanted reflects the +// post-mutation truth. +// +// Logic: +// - If the new record names us as its spindle, ensure we're +// subscribed to its knot (AddKnot is idempotent, so calling it on +// an already-watched knot is cheap). +// - If the old record named us as its spindle, check whether any +// other repo still references that knot through us; if not, +// RemoveKnot it. Skip this when the knot didn't change AND the +// spindle didn't move away from us, because then nothing actually +// released our hold. +func reconcileKnot( + ctx context.Context, + st *store, + knots KnotConsumer, + hostname string, + oldKnot, oldSpindle string, + newKnot, newSpindle string, +) error { + // Tests pass nil for the consumer when they only care about the + // store mutation half of the handler; tolerate that here so + // callers don't have to special-case it. + if knots == nil { return nil - case jsOpDelete: - // We don't unsubscribe from the knot here: other repos may - // still want us to watch it. A periodic reconciliation pass - // (not yet implemented) is the right place to drop unused - // subscriptions. - return st.DeleteRepo(ctx, did, c.RKey) + } + + if newSpindle == hostname && newKnot != "" { + knots.AddKnot(ctx, newKnot) + } + + // Did we just lose our claim on oldKnot? Two ways that can happen: + // the spindle field moved off of us, or the knot field moved to a + // different host. Either is a reason to consider unsubscribing + // from oldKnot — but only if no *other* repo still has us on it. + releasedOld := oldSpindle == hostname && oldKnot != "" && + (newSpindle != hostname || newKnot != oldKnot) + if releasedOld { + stillWanted, err := st.IsKnotWanted(ctx, hostname, oldKnot) + if err != nil { + return err + } + if !stillWanted { + knots.RemoveKnot(ctx, oldKnot) + } } return nil } diff --git a/jetstream_test.go b/jetstream_test.go index ae9d93b..0e819fc 100644 --- a/jetstream_test.go +++ b/jetstream_test.go @@ -310,3 +310,164 @@ func TestRepoEventIgnoresKnotForOtherSpindle(t *testing.T) { t.Fatalf("AddKnot calls = %v, want none", added) } } + +// TestRepoUpdateSpindleAwayFromUsRemovesKnot covers the case where a +// repo we'd previously been watching gets its .spindle field flipped to +// some other spindle. Once that's the only repo we had on that knot, +// the reconciliation must call RemoveKnot. +func TestRepoUpdateSpindleAwayFromUsRemovesKnot(t *testing.T) { + s := newTestStore(t) + ctx := context.Background() + + const ours = "tack.example" + const knot = "knot.example" + + // Seed: a repo that names us as its spindle on `knot`. + if err := s.UpsertRepo(ctx, "did:plc:a", "rk", knot, "repo-a", ours, "", "t"); err != nil { + t.Fatal(err) + } + + // Update: same record, now points at a different spindle. + other := "other.example" + rec := tangled.Repo{ + Knot: knot, + Name: "repo-a", + Spindle: &other, + CreatedAt: "2026-01-01T00:00:00Z", + } + evt := commitEvent(1, "did:plc:a", tangled.RepoNSID, jsOpUpdate, "rk", rec) + + fake := &fakeKnotConsumer{} + if err := handleJetstreamEvent(ctx, s, fake, ours, evt); err != nil { + t.Fatalf("handle: %v", err) + } + if added := fake.Added(); len(added) != 0 { + t.Fatalf("AddKnot calls = %v, want none", added) + } + if removed := fake.Removed(); len(removed) != 1 || removed[0] != knot { + t.Fatalf("RemoveKnot calls = %v, want [%s]", removed, knot) + } +} + +// TestRepoUpdateSpindleAwayFromUsKeepsKnotIfShared ensures we don't +// over-eagerly unsubscribe from a knot when one of multiple repos on +// that knot moves away from us. The other repo's subscription must keep +// the knot in our wanted set. +func TestRepoUpdateSpindleAwayFromUsKeepsKnotIfShared(t *testing.T) { + s := newTestStore(t) + ctx := context.Background() + + const ours = "tack.example" + const knot = "knot.example" + + // Two repos sharing one knot, both pointed at us. + if err := s.UpsertRepo(ctx, "did:plc:a", "rk1", knot, "a", ours, "", "t"); err != nil { + t.Fatal(err) + } + if err := s.UpsertRepo(ctx, "did:plc:b", "rk2", knot, "b", ours, "", "t"); err != nil { + t.Fatal(err) + } + + // Repo A flips to a different spindle. B still wants us. + other := "other.example" + rec := tangled.Repo{ + Knot: knot, + Name: "a", + Spindle: &other, + CreatedAt: "2026-01-01T00:00:00Z", + } + evt := commitEvent(1, "did:plc:a", tangled.RepoNSID, jsOpUpdate, "rk1", rec) + + fake := &fakeKnotConsumer{} + if err := handleJetstreamEvent(ctx, s, fake, ours, evt); err != nil { + t.Fatalf("handle: %v", err) + } + if removed := fake.Removed(); len(removed) != 0 { + t.Fatalf("RemoveKnot calls = %v, want none (B still wants us on %s)", removed, knot) + } +} + +// TestRepoUpdateChangingKnotSwapsSubscription verifies that a repo +// staying with us but changing its .knot field unsubscribes the old +// knot (if no other repo holds it) and subscribes the new one. +func TestRepoUpdateChangingKnotSwapsSubscription(t *testing.T) { + s := newTestStore(t) + ctx := context.Background() + + const ours = "tack.example" + const oldKnot = "old.example" + const newKnot = "new.example" + + if err := s.UpsertRepo(ctx, "did:plc:a", "rk", oldKnot, "a", ours, "", "t"); err != nil { + t.Fatal(err) + } + + spindle := ours + rec := tangled.Repo{ + Knot: newKnot, + Name: "a", + Spindle: &spindle, + CreatedAt: "2026-01-01T00:00:00Z", + } + evt := commitEvent(1, "did:plc:a", tangled.RepoNSID, jsOpUpdate, "rk", rec) + + fake := &fakeKnotConsumer{} + if err := handleJetstreamEvent(ctx, s, fake, ours, evt); err != nil { + t.Fatalf("handle: %v", err) + } + if added := fake.Added(); len(added) != 1 || added[0] != newKnot { + t.Fatalf("AddKnot calls = %v, want [%s]", added, newKnot) + } + if removed := fake.Removed(); len(removed) != 1 || removed[0] != oldKnot { + t.Fatalf("RemoveKnot calls = %v, want [%s]", removed, oldKnot) + } +} + +// TestRepoDeleteRemovesKnotWhenLast confirms deleting the last repo on +// a knot we cared about triggers RemoveKnot. +func TestRepoDeleteRemovesKnotWhenLast(t *testing.T) { + s := newTestStore(t) + ctx := context.Background() + + const ours = "tack.example" + const knot = "knot.example" + + if err := s.UpsertRepo(ctx, "did:plc:a", "rk", knot, "a", ours, "", "t"); err != nil { + t.Fatal(err) + } + evt := commitEvent(1, "did:plc:a", tangled.RepoNSID, jsOpDelete, "rk", nil) + + fake := &fakeKnotConsumer{} + if err := handleJetstreamEvent(ctx, s, fake, ours, evt); err != nil { + t.Fatalf("handle: %v", err) + } + if removed := fake.Removed(); len(removed) != 1 || removed[0] != knot { + t.Fatalf("RemoveKnot calls = %v, want [%s]", removed, knot) + } +} + +// TestRepoDeleteKeepsKnotIfShared ensures deleting one of multiple +// repos on a knot does not unsubscribe — the survivors still want it. +func TestRepoDeleteKeepsKnotIfShared(t *testing.T) { + s := newTestStore(t) + ctx := context.Background() + + const ours = "tack.example" + const knot = "knot.example" + + if err := s.UpsertRepo(ctx, "did:plc:a", "rk1", knot, "a", ours, "", "t"); err != nil { + t.Fatal(err) + } + if err := s.UpsertRepo(ctx, "did:plc:b", "rk2", knot, "b", ours, "", "t"); err != nil { + t.Fatal(err) + } + evt := commitEvent(1, "did:plc:a", tangled.RepoNSID, jsOpDelete, "rk1", nil) + + fake := &fakeKnotConsumer{} + if err := handleJetstreamEvent(ctx, s, fake, ours, evt); err != nil { + t.Fatalf("handle: %v", err) + } + if removed := fake.Removed(); len(removed) != 0 { + t.Fatalf("RemoveKnot calls = %v, want none", removed) + } +} diff --git a/knot.go b/knot.go index 5023cde..aacb21d 100644 --- a/knot.go +++ b/knot.go @@ -19,10 +19,10 @@ package main // Once the build pipeline is wired up this is where pipeline // triggers will be translated into Buildkite builds. // -// The jetstream consumer also gets a back-reference (via the knotAdder -// interface) so it can dynamically subscribe to a new knot the moment a -// matching sh.tangled.repo record arrives, without waiting for a tack -// restart. +// The jetstream consumer also gets a back-reference (via the +// KnotConsumer interface) so it can dynamically subscribe a new knot — +// or, conceptually, unsubscribe one — the moment a matching +// sh.tangled.repo record arrives, without waiting for a tack restart. import ( "context" @@ -50,6 +50,16 @@ type KnotConsumer interface { // a no-op. An empty knot string is ignored. The supplied context // scopes the dial; cancelling it tears the subscription down. AddKnot(ctx context.Context, knot string) + + // RemoveKnot stops processing events from the given knot. It is the + // inverse of AddKnot and must tolerate being called for a knot that + // was never added (no-op). An empty knot string is ignored. + // + // The production implementation is currently a no-op: tangled-core's + // eventconsumer does not expose a way to drop an individual source's + // websocket. Tracked upstream as + // https://tangled.org/did:plc:j5hmlfdrwkvtxm7cjmu7j2is/issues/510 + RemoveKnot(ctx context.Context, knot string) } // knotConsumer is the production KnotConsumer. It wraps @@ -118,6 +128,19 @@ func (k *knotConsumer) AddKnot(ctx context.Context, knot string) { k.c.AddSource(ctx, eventconsumer.NewKnotSource(knot)) } +// RemoveKnot is currently a no-op — see the interface comment for the +// upstream blocker. We still log at info so operators can see when the +// reconciliation logic *would* have unsubscribed; once eventconsumer +// gains a RemoveSource primitive, swap the body for a real call. +func (k *knotConsumer) RemoveKnot(_ context.Context, knot string) { + if knot == "" { + return + } + k.log.Info("remove knot source requested (no-op until upstream supports it)", + "knot", knot, + ) +} + // Stop tears down all knot websocket connections and waits for the // consumer's goroutines to exit. It must be called exactly once. func (k *knotConsumer) Stop() { diff --git a/knot_fake.go b/knot_fake.go index d3147d5..05894d5 100644 --- a/knot_fake.go +++ b/knot_fake.go @@ -6,8 +6,8 @@ package main // any future test files (and, if we ever split tack into subpackages, can // be promoted to an exported helper without moving code around). // -// It does no I/O: AddKnot just records the knot it was handed so tests -// can assert on the side effect. +// It does no I/O: AddKnot and RemoveKnot just record the knot they were +// handed so tests can assert on the side effect. import ( "context" @@ -15,10 +15,11 @@ import ( ) // fakeKnotConsumer is an in-memory KnotConsumer suitable for tests. The -// zero value is ready to use; concurrent calls to AddKnot are safe. +// zero value is ready to use; concurrent calls are safe. type fakeKnotConsumer struct { - mu sync.Mutex - added []string + mu sync.Mutex + added []string + removed []string } // Compile-time interface conformance check — keeps the fake honest if @@ -32,6 +33,13 @@ func (f *fakeKnotConsumer) AddKnot(_ context.Context, knot string) { f.added = append(f.added, knot) } +// RemoveKnot records the knot for later inspection via Removed(). +func (f *fakeKnotConsumer) RemoveKnot(_ context.Context, knot string) { + f.mu.Lock() + defer f.mu.Unlock() + f.removed = append(f.removed, knot) +} + // Added returns a copy of the knots passed to AddKnot, in call order. // A copy is returned so callers can't accidentally mutate the fake's // internal slice while comparing. @@ -42,3 +50,13 @@ func (f *fakeKnotConsumer) Added() []string { copy(out, f.added) return out } + +// Removed returns a copy of the knots passed to RemoveKnot, in call +// order. See Added() for the rationale behind copying. +func (f *fakeKnotConsumer) Removed() []string { + f.mu.Lock() + defer f.mu.Unlock() + out := make([]string, len(f.removed)) + copy(out, f.removed) + return out +} diff --git a/store.go b/store.go index eb838a7..51d5d75 100644 --- a/store.go +++ b/store.go @@ -194,6 +194,45 @@ func (s *store) UpsertRepoCollaborator(ctx context.Context, did, rkey, repo, rep return nil } +// GetRepo returns the (knot, spindle) currently stored for a (did, rkey) +// pair. Both are returned as empty strings when no row exists; callers +// that need to distinguish "absent" from "stored but empty" should +// pre-check existence themselves. +// +// This exists so applyRepo can read the *previous* spindle/knot of a +// record before applying a mutation, which is what makes it possible to +// detect transitions like "this repo used to be ours, now it isn't" and +// trigger a knot unsubscribe. +func (s *store) GetRepo(ctx context.Context, did, rkey string) (knot, spindle string, err error) { + err = s.db.QueryRowContext(ctx, + `SELECT knot, spindle FROM repos WHERE did = ? AND rkey = ?`, + did, rkey, + ).Scan(&knot, &spindle) + if errors.Is(err, sql.ErrNoRows) { + return "", "", nil + } + if err != nil { + return "", "", fmt.Errorf("get repo: %w", err) + } + return knot, spindle, nil +} + +// IsKnotWanted reports whether any repo currently stored still names the +// given hostname as its spindle and the given knot as its host. After a +// repo update or delete this is the question we ask to decide whether +// to keep watching that knot or unsubscribe from it. +func (s *store) IsKnotWanted(ctx context.Context, hostname, knot string) (bool, error) { + var n int + err := s.db.QueryRowContext(ctx, + `SELECT COUNT(*) FROM repos WHERE spindle = ? AND knot = ?`, + hostname, knot, + ).Scan(&n) + if err != nil { + return false, fmt.Errorf("count repos for knot: %w", err) + } + return n > 0, nil +} + // KnotsForSpindle returns the distinct knot hostnames of all repos that // have declared the given spindle hostname as their CI spindle. The knot // event-stream subscriber uses this to decide which knots to dial.