diff --git a/jetstream.go b/jetstream.go index c898c02..54f996e 100644 --- a/jetstream.go +++ b/jetstream.go @@ -78,12 +78,13 @@ func startJetstream(ctx context.Context, cfg config, st *store, knots KnotConsum clientCfg.WebsocketURL = cfg.JetstreamURL clientCfg.WantedCollections = collections - // The handler closes over `st`, `knots` and the spindle hostname so - // the scheduler signature stays plain `func(ctx, *Event) error` and - // applyCommit can hand the knot consumer new sources as soon as - // matching repo records arrive. + // The handler closes over `st`, `knots`, the spindle hostname and + // the owner DID so the scheduler signature stays plain + // `func(ctx, *Event) error` and applyCommit can hand the knot + // consumer new sources as soon as matching repo records arrive + // (gated on the publisher being an authorized actor). handler := func(ctx context.Context, evt *jsmodels.Event) error { - return handleJetstreamEvent(ctx, st, knots, cfg.Hostname, evt) + return handleJetstreamEvent(ctx, st, knots, cfg.Hostname, cfg.OwnerDID, evt) } // Re-attach the component-scoped logger so handler — which the @@ -177,7 +178,7 @@ func startJetstream(ctx context.Context, cfg config, st *store, knots KnotConsum // lose a row that the rest of tack relies on. Leaving the cursor // in place gives the next reconnect (which rewinds by // jetstreamRewind) a chance to re-deliver and re-apply it. -func handleJetstreamEvent(ctx context.Context, st *store, knots KnotConsumer, hostname string, evt *jsmodels.Event) error { +func handleJetstreamEvent(ctx context.Context, st *store, knots KnotConsumer, hostname, ownerDID string, evt *jsmodels.Event) error { // We only care about commits, which are the actual record CRUD // operations on a user's PDS. Account/identity events are ignored // for now; if we ever care about handle changes we can add them. @@ -189,7 +190,7 @@ func handleJetstreamEvent(ctx context.Context, st *store, knots KnotConsumer, ho // Dispatch on collection. Unknown collections shouldn't happen given // our wantedCollections filter, but be defensive — jetstream may // send schema changes ahead of us updating the filter. - applyErr := applyCommit(ctx, st, knots, hostname, evt) + applyErr := applyCommit(ctx, st, knots, hostname, ownerDID, evt) if applyErr != nil { logger.Error("apply commit", "err", applyErr, @@ -261,13 +262,13 @@ func isBadRecord(err error) bool { // applyCommit routes a commit to the right store mutation based on its // collection NSID and operation. -func applyCommit(ctx context.Context, st *store, knots KnotConsumer, hostname string, evt *jsmodels.Event) error { +func applyCommit(ctx context.Context, st *store, knots KnotConsumer, hostname, ownerDID string, evt *jsmodels.Event) error { c := evt.Commit switch c.Collection { case tangled.SpindleMemberNSID: - return applySpindleMember(ctx, st, evt.Did, c) + return applySpindleMember(ctx, st, knots, hostname, ownerDID, evt.Did, c) case tangled.RepoNSID: - return applyRepo(ctx, st, knots, hostname, evt.Did, c) + return applyRepo(ctx, st, knots, hostname, ownerDID, evt.Did, c) case tangled.RepoCollaboratorNSID: return applyRepoCollaborator(ctx, st, evt.Did, c) default: @@ -279,23 +280,61 @@ func applyCommit(ctx context.Context, st *store, knots KnotConsumer, hostname st } } -func applySpindleMember(ctx context.Context, st *store, did string, c *jsmodels.Commit) error { +func applySpindleMember(ctx context.Context, st *store, knots KnotConsumer, hostname, ownerDID, did string, c *jsmodels.Commit) error { switch c.Operation { case jsOpCreate, jsOpUpdate: + // Capture the previous subject. Necessary because a same-rkey + // update can move the grant to a different DID, in which case + // the *old* subject's knots may need to be released even as + // the new subject's are picked up. + oldSubject, err := st.GetSpindleMember(ctx, did, c.RKey) + if err != nil { + return err + } + var rec tangled.SpindleMember if err := json.Unmarshal(c.Record, &rec); err != nil { // Decode failures are a permanent property of the record's // bytes; mark as bad so the cursor can advance past it. return badRecord(fmt.Errorf("decode spindle.member: %w", err)) } - return st.UpsertSpindleMember(ctx, did, c.RKey, rec.Instance, rec.Subject, rec.CreatedAt) + if err := st.UpsertSpindleMember(ctx, did, c.RKey, rec.Instance, rec.Subject, rec.CreatedAt); err != nil { + return err + } + + // Only grants published by the spindle owner actually change + // authorization (see IsAuthorizedActor). Forged grants are + // stored but don't move the needle, so don't bother + // reconciling on them: it would be a no-op at best and + // add log noise. + if did != ownerDID { + return nil + } + if oldSubject != "" && oldSubject != rec.Subject { + if err := reconcileMember(ctx, st, knots, hostname, ownerDID, oldSubject); err != nil { + return err + } + } + return reconcileMember(ctx, st, knots, hostname, ownerDID, rec.Subject) + case jsOpDelete: - return st.DeleteSpindleMember(ctx, did, c.RKey) + oldSubject, err := st.GetSpindleMember(ctx, did, c.RKey) + if err != nil { + return err + } + if err := st.DeleteSpindleMember(ctx, did, c.RKey); err != nil { + return err + } + // As above: only the owner's grants matter for authorization. + if did != ownerDID || oldSubject == "" { + return nil + } + return reconcileMember(ctx, st, knots, hostname, ownerDID, oldSubject) } return nil } -func applyRepo(ctx context.Context, st *store, knots KnotConsumer, hostname string, did string, c *jsmodels.Commit) error { +func applyRepo(ctx context.Context, st *store, knots KnotConsumer, hostname, ownerDID, did string, c *jsmodels.Commit) error { switch c.Operation { case jsOpCreate, jsOpUpdate: var rec tangled.Repo @@ -322,7 +361,7 @@ func applyRepo(ctx context.Context, st *store, knots KnotConsumer, hostname stri } newSpindle := deref(rec.Spindle) - return reconcileKnot(ctx, st, knots, hostname, + return reconcileKnot(ctx, st, knots, hostname, ownerDID, did, oldKnot, oldSpindle, rec.Knot, newSpindle, ) @@ -339,7 +378,7 @@ func applyRepo(ctx context.Context, st *store, knots KnotConsumer, hostname stri if err := st.DeleteRepo(ctx, did, c.RKey); err != nil { return err } - return reconcileKnot(ctx, st, knots, hostname, + return reconcileKnot(ctx, st, knots, hostname, ownerDID, did, oldKnot, oldSpindle, "", "", ) @@ -353,19 +392,21 @@ func applyRepo(ctx context.Context, st *store, knots KnotConsumer, hostname stri // 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 new record names us as its spindle AND the publisher is +// an authorized actor (spindle owner or owner-vouched member), +// ensure we're subscribed to its knot. Without the membership +// check, any firehose publisher could pin us to an attacker-chosen +// knot just by publishing a sh.tangled.repo record naming us. // - 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. +// other authorized 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, + hostname, ownerDID, publisherDID string, oldKnot, oldSpindle string, newKnot, newSpindle string, ) error { @@ -377,17 +418,29 @@ func reconcileKnot( } if newSpindle == hostname && newKnot != "" { - knots.AddKnot(ctx, newKnot) + ok, err := st.IsAuthorizedActor(ctx, ownerDID, publisherDID) + if err != nil { + return err + } + if ok { + knots.AddKnot(ctx, newKnot) + } else { + loggerFrom(ctx).Warn("ignoring repo from unauthorized publisher", + "publisher_did", publisherDID, + "knot", 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. + // from oldKnot, but only if no *other* authorized repo still has + // us on it. releasedOld := oldSpindle == hostname && oldKnot != "" && (newSpindle != hostname || newKnot != oldKnot) if releasedOld { - stillWanted, err := st.IsKnotWanted(ctx, hostname, oldKnot) + stillWanted, err := st.IsKnotWanted(ctx, hostname, ownerDID, oldKnot) if err != nil { return err } @@ -398,6 +451,56 @@ func reconcileKnot( return nil } +// reconcileMember adjusts knot subscriptions after a membership grant +// or revocation may have changed `subject`'s authorization status. +// For each knot named by subject's repos that point at us: +// +// - if subject is now authorized, AddKnot (idempotent: already- +// subscribed knots are no-ops in the consumer); +// - if subject is now unauthorized, ask IsKnotWanted whether any +// *other* authorized repo still holds the knot; if not, +// RemoveKnot. +// +// Without this, a member's repos picked up over the firehose before +// the grant arrived would never get subscribed (the grant doesn't +// re-deliver the older repo events), and a revocation would leave +// the now-unauthorized publisher's knot subscribed until restart. +func reconcileMember( + ctx context.Context, + st *store, + knots KnotConsumer, + hostname, ownerDID, subject string, +) error { + if knots == nil || subject == "" { + return nil + } + knotsForSubject, err := st.KnotsForOwner(ctx, hostname, subject) + if err != nil { + return err + } + if len(knotsForSubject) == 0 { + return nil + } + authorized, err := st.IsAuthorizedActor(ctx, ownerDID, subject) + if err != nil { + return err + } + for _, k := range knotsForSubject { + if authorized { + knots.AddKnot(ctx, k) + continue + } + stillWanted, err := st.IsKnotWanted(ctx, hostname, ownerDID, k) + if err != nil { + return err + } + if !stillWanted { + knots.RemoveKnot(ctx, k) + } + } + return nil +} + func applyRepoCollaborator(ctx context.Context, st *store, did string, c *jsmodels.Commit) error { switch c.Operation { case jsOpCreate, jsOpUpdate: diff --git a/jetstream_test.go b/jetstream_test.go index 3c88a18..7f35a9b 100644 --- a/jetstream_test.go +++ b/jetstream_test.go @@ -67,7 +67,7 @@ func TestHandleNonCommitEvent(t *testing.T) { TimeUS: 100, Kind: jsmodels.EventKindAccount, } - if err := handleJetstreamEvent(ctx, s, nil, "", evt); err != nil { + if err := handleJetstreamEvent(ctx, s, nil, "", "", evt); err != nil { t.Fatalf("handle: %v", err) } got, err := s.LoadCursor(ctx) @@ -91,7 +91,7 @@ func TestHandleSpindleMemberCreate(t *testing.T) { CreatedAt: "2026-01-01T00:00:00Z", } evt := commitEvent(12345, "did:plc:owner", tangled.SpindleMemberNSID, jsOpCreate, "rk1", rec) - if err := handleJetstreamEvent(ctx, s, nil, "", evt); err != nil { + if err := handleJetstreamEvent(ctx, s, nil, "", "", evt); err != nil { t.Fatalf("handle: %v", err) } @@ -119,7 +119,7 @@ func TestHandleSpindleMemberDelete(t *testing.T) { t.Fatalf("seed: %v", err) } evt := commitEvent(99, "did:plc:owner", tangled.SpindleMemberNSID, jsOpDelete, "rk1", nil) - if err := handleJetstreamEvent(ctx, s, nil, "", evt); err != nil { + if err := handleJetstreamEvent(ctx, s, nil, "", "", evt); err != nil { t.Fatalf("handle: %v", err) } if n := countRows(t, s, "spindle_members"); n != 0 { @@ -145,7 +145,7 @@ func TestHandleRepoCreateOptionals(t *testing.T) { CreatedAt: "2026-01-01T00:00:00Z", } evt := commitEvent(7, "did:plc:owner", tangled.RepoNSID, jsOpCreate, "repo1", rec) - if err := handleJetstreamEvent(ctx, s, nil, "", evt); err != nil { + if err := handleJetstreamEvent(ctx, s, nil, "", "", evt); err != nil { t.Fatalf("handle: %v", err) } @@ -169,7 +169,7 @@ func TestHandleRepoCreateOptionals(t *testing.T) { CreatedAt: "2026-01-01T00:00:00Z", } evt2 := commitEvent(8, "did:plc:owner", tangled.RepoNSID, jsOpCreate, "repo2", rec2) - if err := handleJetstreamEvent(ctx, s, nil, "", evt2); err != nil { + if err := handleJetstreamEvent(ctx, s, nil, "", "", evt2); err != nil { t.Fatalf("handle nil-optionals: %v", err) } err = s.db.QueryRowContext(ctx, @@ -198,7 +198,7 @@ func TestHandleRepoCollaboratorCreate(t *testing.T) { CreatedAt: "2026-01-01T00:00:00Z", } evt := commitEvent(55, "did:plc:owner", tangled.RepoCollaboratorNSID, jsOpCreate, "c1", rec) - if err := handleJetstreamEvent(ctx, s, nil, "", evt); err != nil { + if err := handleJetstreamEvent(ctx, s, nil, "", "", evt); err != nil { t.Fatalf("handle: %v", err) } @@ -224,7 +224,7 @@ func TestHandleUnknownCollection(t *testing.T) { ctx := context.Background() evt := commitEvent(42, "did:plc:owner", "app.bsky.feed.post", jsOpCreate, "rk", map[string]string{"text": "hi"}) - if err := handleJetstreamEvent(ctx, s, nil, "", evt); err != nil { + if err := handleJetstreamEvent(ctx, s, nil, "", "", evt); err != nil { t.Fatalf("handle: %v", err) } requireCursor(t, s, 42) @@ -249,7 +249,7 @@ func TestHandleBadRecordAdvancesCursor(t *testing.T) { Record: json.RawMessage(`{not valid json`), }, } - if err := handleJetstreamEvent(ctx, s, nil, "", evt); err != nil { + if err := handleJetstreamEvent(ctx, s, nil, "", "", evt); err != nil { t.Fatalf("handle should swallow decode error, got: %v", err) } if n := countRows(t, s, "spindle_members"); n != 0 { @@ -295,7 +295,7 @@ func TestHandleTransientStoreErrorDoesNotAdvanceCursor(t *testing.T) { evt := commitEvent(1000, "did:plc:owner", tangled.SpindleMemberNSID, jsOpCreate, "rk", rec) // Expect a non-nil error: handler propagates transient failures. - if err := handleJetstreamEvent(ctx, s, nil, "", evt); err == nil { + if err := handleJetstreamEvent(ctx, s, nil, "", "", evt); err == nil { t.Fatalf("handle should return transient error, got nil") } @@ -334,7 +334,7 @@ func TestRepoEventSubscribesKnotForOurSpindle(t *testing.T) { evt := commitEvent(1, "did:plc:owner", tangled.RepoNSID, jsOpCreate, "rk", rec) fake := &fakeKnotConsumer{} - if err := handleJetstreamEvent(ctx, s, fake, ours, evt); err != nil { + if err := handleJetstreamEvent(ctx, s, fake, ours, "did:plc:owner", evt); err != nil { t.Fatalf("handle: %v", err) } added := fake.Added() @@ -343,6 +343,159 @@ func TestRepoEventSubscribesKnotForOurSpindle(t *testing.T) { } } +// TestRepoEventIgnoresUnauthorizedPublisher pins the high-severity +// security gate: a sh.tangled.repo record naming us as its spindle +// must NOT cause an outbound knot subscription unless its publisher +// is the spindle owner or an owner-vouched member. Without this gate +// any firehose publisher could force tack to dial an attacker-chosen +// websocket simply by minting a matching repo record. +func TestRepoEventIgnoresUnauthorizedPublisher(t *testing.T) { + s := newTestStore(t) + ctx := context.Background() + + const ours = "tack.example" + const owner = "did:plc:owner" + + // did:plc:rando is neither the owner nor a member. + spindle := ours + rec := tangled.Repo{ + Knot: "evil.example", + Name: "myrepo", + Spindle: &spindle, + CreatedAt: "2026-01-01T00:00:00Z", + } + evt := commitEvent(1, "did:plc:rando", tangled.RepoNSID, jsOpCreate, "rk", rec) + + fake := &fakeKnotConsumer{} + if err := handleJetstreamEvent(ctx, s, fake, ours, owner, evt); err != nil { + t.Fatalf("handle: %v", err) + } + if added := fake.Added(); len(added) != 0 { + t.Fatalf("AddKnot calls = %v, want none (publisher is unauthorized)", added) + } +} + +// TestSpindleMemberGrantSubscribesPendingKnots covers the case where +// a member publishes their sh.tangled.repo *before* the owner's +// sh.tangled.spindle.member grant arrives over the firehose. Until +// the grant lands the publisher is unauthorized, so the repo doesn't +// pull a subscription. Once the grant arrives, the same-rkey reconcile +// must catch up by AddKnot-ing every knot that member's pending repos +// already named. +func TestSpindleMemberGrantSubscribesPendingKnots(t *testing.T) { + s := newTestStore(t) + ctx := context.Background() + + const ours = "tack.example" + const owner = "did:plc:owner" + const alice = "did:plc:alice" + + // Step 1: alice publishes a repo claiming us. She isn't a member + // yet, so no subscription should happen. + spindle := ours + repoRec := tangled.Repo{ + Knot: "knot.example", + Name: "myrepo", + Spindle: &spindle, + CreatedAt: "2026-01-01T00:00:00Z", + } + repoEvt := commitEvent(1, alice, tangled.RepoNSID, jsOpCreate, "rk", repoRec) + + fake := &fakeKnotConsumer{} + if err := handleJetstreamEvent(ctx, s, fake, ours, owner, repoEvt); err != nil { + t.Fatalf("handle repo: %v", err) + } + if added := fake.Added(); len(added) != 0 { + t.Fatalf("AddKnot before grant = %v, want none", added) + } + + // Step 2: owner publishes a membership grant naming alice. The + // reconcile triggered by that grant must subscribe to alice's + // already-stored repo's knot. + grant := tangled.SpindleMember{ + Instance: ours, + Subject: alice, + CreatedAt: "2026-01-02T00:00:00Z", + } + grantEvt := commitEvent(2, owner, tangled.SpindleMemberNSID, jsOpCreate, "mk", grant) + if err := handleJetstreamEvent(ctx, s, fake, ours, owner, grantEvt); err != nil { + t.Fatalf("handle grant: %v", err) + } + if added := fake.Added(); len(added) != 1 || added[0] != "knot.example" { + t.Fatalf("AddKnot after grant = %v, want [knot.example]", added) + } +} + +// TestSpindleMemberRevokeUnsubscribesKnot verifies that revoking a +// previously-granted membership tears down the only-held-by-that-member +// knot subscription. Without this, an attacker who briefly held +// membership could leave us dialing their chosen knot until restart, +// the exact lingering-subscription concern called out in KNOWN_ISSUES. +func TestSpindleMemberRevokeUnsubscribesKnot(t *testing.T) { + s := newTestStore(t) + ctx := context.Background() + + const ours = "tack.example" + const owner = "did:plc:owner" + const alice = "did:plc:alice" + + // Seed: owner has previously vouched for alice, and alice has + // already published a repo on knot.example pointing at us. + if err := s.UpsertSpindleMember(ctx, owner, "mk", ours, alice, "t"); err != nil { + t.Fatal(err) + } + if err := s.UpsertRepo(ctx, alice, "rk", "knot.example", "myrepo", ours, "", "t"); err != nil { + t.Fatal(err) + } + + // Owner publishes a delete of the membership grant. + revoke := commitEvent(10, owner, tangled.SpindleMemberNSID, jsOpDelete, "mk", nil) + + fake := &fakeKnotConsumer{} + if err := handleJetstreamEvent(ctx, s, fake, ours, owner, revoke); err != nil { + t.Fatalf("handle revoke: %v", err) + } + if removed := fake.Removed(); len(removed) != 1 || removed[0] != "knot.example" { + t.Fatalf("RemoveKnot calls = %v, want [knot.example]", removed) + } +} + +// TestSpindleMemberGrantFromNonOwnerIgnored confirms that a forged +// membership record published by anyone other than the spindle owner +// does NOT trigger a knot subscription, even if a matching repo for +// the named subject exists. This mirrors AuthorizePipelineActor's +// rule that the membership grant's publisher must equal the spindle +// owner. +func TestSpindleMemberGrantFromNonOwnerIgnored(t *testing.T) { + s := newTestStore(t) + ctx := context.Background() + + const ours = "tack.example" + const owner = "did:plc:owner" + const alice = "did:plc:alice" + + if err := s.UpsertRepo(ctx, alice, "rk", "knot.example", "myrepo", ours, "", "t"); err != nil { + t.Fatal(err) + } + + // Forged grant: alice "vouches" for herself. Stored, but must + // not flip her into authorized status. + forgery := tangled.SpindleMember{ + Instance: ours, + Subject: alice, + CreatedAt: "2026-01-02T00:00:00Z", + } + evt := commitEvent(2, alice, tangled.SpindleMemberNSID, jsOpCreate, "mk", forgery) + + fake := &fakeKnotConsumer{} + if err := handleJetstreamEvent(ctx, s, fake, ours, owner, evt); err != nil { + t.Fatalf("handle forged grant: %v", err) + } + if added := fake.Added(); len(added) != 0 { + t.Fatalf("AddKnot calls = %v, want none (grant publisher != owner)", added) + } +} + // TestRepoEventIgnoresKnotForOtherSpindle confirms repos pointing at a // *different* spindle do not pull us into watching their knot. Without // this guard, tack would dial every knot named in any sh.tangled.repo @@ -361,7 +514,7 @@ func TestRepoEventIgnoresKnotForOtherSpindle(t *testing.T) { evt := commitEvent(1, "did:plc:owner", tangled.RepoNSID, jsOpCreate, "rk", rec) fake := &fakeKnotConsumer{} - if err := handleJetstreamEvent(ctx, s, fake, "tack.example", evt); err != nil { + if err := handleJetstreamEvent(ctx, s, fake, "tack.example", "did:plc:owner", evt); err != nil { t.Fatalf("handle: %v", err) } if added := fake.Added(); len(added) != 0 { @@ -396,7 +549,7 @@ func TestRepoUpdateSpindleAwayFromUsRemovesKnot(t *testing.T) { evt := commitEvent(1, "did:plc:a", tangled.RepoNSID, jsOpUpdate, "rk", rec) fake := &fakeKnotConsumer{} - if err := handleJetstreamEvent(ctx, s, fake, ours, evt); err != nil { + if err := handleJetstreamEvent(ctx, s, fake, ours, "did:plc:owner", evt); err != nil { t.Fatalf("handle: %v", err) } if added := fake.Added(); len(added) != 0 { @@ -418,6 +571,17 @@ func TestRepoUpdateSpindleAwayFromUsKeepsKnotIfShared(t *testing.T) { const ours = "tack.example" const knot = "knot.example" + // Both publishers must be vouched for; otherwise IsKnotWanted's + // membership filter would zero them out and we'd unsubscribe + // regardless. The reconciliation we want to test only matters + // when there is a real authorized hold to keep. + if err := s.UpsertSpindleMember(ctx, "did:plc:owner", "mk1", ours, "did:plc:a", "t"); err != nil { + t.Fatal(err) + } + if err := s.UpsertSpindleMember(ctx, "did:plc:owner", "mk2", ours, "did:plc:b", "t"); err != nil { + t.Fatal(err) + } + // 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) @@ -437,7 +601,7 @@ func TestRepoUpdateSpindleAwayFromUsKeepsKnotIfShared(t *testing.T) { evt := commitEvent(1, "did:plc:a", tangled.RepoNSID, jsOpUpdate, "rk1", rec) fake := &fakeKnotConsumer{} - if err := handleJetstreamEvent(ctx, s, fake, ours, evt); err != nil { + if err := handleJetstreamEvent(ctx, s, fake, ours, "did:plc:owner", evt); err != nil { t.Fatalf("handle: %v", err) } if removed := fake.Removed(); len(removed) != 0 { @@ -456,6 +620,13 @@ func TestRepoUpdateChangingKnotSwapsSubscription(t *testing.T) { const oldKnot = "old.example" const newKnot = "new.example" + // Vouch for did:plc:a so reconcileKnot will actually call + // AddKnot for the new host. Without the grant the publisher + // is unauthorized and the AddKnot half is (correctly) skipped. + if err := s.UpsertSpindleMember(ctx, "did:plc:owner", "mk1", ours, "did:plc:a", "t"); err != nil { + t.Fatal(err) + } + if err := s.UpsertRepo(ctx, "did:plc:a", "rk", oldKnot, "a", ours, "", "t"); err != nil { t.Fatal(err) } @@ -470,7 +641,7 @@ func TestRepoUpdateChangingKnotSwapsSubscription(t *testing.T) { evt := commitEvent(1, "did:plc:a", tangled.RepoNSID, jsOpUpdate, "rk", rec) fake := &fakeKnotConsumer{} - if err := handleJetstreamEvent(ctx, s, fake, ours, evt); err != nil { + if err := handleJetstreamEvent(ctx, s, fake, ours, "did:plc:owner", evt); err != nil { t.Fatalf("handle: %v", err) } if added := fake.Added(); len(added) != 1 || added[0] != newKnot { @@ -496,7 +667,7 @@ func TestRepoDeleteRemovesKnotWhenLast(t *testing.T) { evt := commitEvent(1, "did:plc:a", tangled.RepoNSID, jsOpDelete, "rk", nil) fake := &fakeKnotConsumer{} - if err := handleJetstreamEvent(ctx, s, fake, ours, evt); err != nil { + if err := handleJetstreamEvent(ctx, s, fake, ours, "did:plc:owner", evt); err != nil { t.Fatalf("handle: %v", err) } if removed := fake.Removed(); len(removed) != 1 || removed[0] != knot { @@ -513,6 +684,15 @@ func TestRepoDeleteKeepsKnotIfShared(t *testing.T) { const ours = "tack.example" const knot = "knot.example" + // Both publishers vouched for; otherwise IsKnotWanted would + // (correctly) report neither holds the knot and we'd unsubscribe. + if err := s.UpsertSpindleMember(ctx, "did:plc:owner", "mk1", ours, "did:plc:a", "t"); err != nil { + t.Fatal(err) + } + if err := s.UpsertSpindleMember(ctx, "did:plc:owner", "mk2", ours, "did:plc:b", "t"); err != nil { + t.Fatal(err) + } + if err := s.UpsertRepo(ctx, "did:plc:a", "rk1", knot, "a", ours, "", "t"); err != nil { t.Fatal(err) } @@ -522,7 +702,7 @@ func TestRepoDeleteKeepsKnotIfShared(t *testing.T) { evt := commitEvent(1, "did:plc:a", tangled.RepoNSID, jsOpDelete, "rk1", nil) fake := &fakeKnotConsumer{} - if err := handleJetstreamEvent(ctx, s, fake, ours, evt); err != nil { + if err := handleJetstreamEvent(ctx, s, fake, ours, "did:plc:owner", evt); err != nil { t.Fatalf("handle: %v", err) } if removed := fake.Removed(); len(removed) != 0 { diff --git a/knot.go b/knot.go index 4db1994..cd44c9f 100644 --- a/knot.go +++ b/knot.go @@ -120,7 +120,7 @@ var _ KnotConsumer = (*knotConsumer)(nil) func startKnotConsumer(ctx context.Context, cfg config, st *store, provider Provider) (*knotConsumer, error) { logger := loggerFrom(ctx).With("component", "knotconsumer") - knots, err := st.KnotsForSpindle(ctx, cfg.Hostname) + knots, err := st.KnotsForSpindle(ctx, cfg.Hostname, cfg.OwnerDID) if err != nil { return nil, fmt.Errorf("load known knots: %w", err) } diff --git a/store.go b/store.go index 04e7303..bcfe560 100644 --- a/store.go +++ b/store.go @@ -1,22 +1,18 @@ package main -// SQLite-backed persistence for tack. +// SQLite-backed persistence for tack. Holds: // -// Two responsibilities live here: +// - the jetstream cursor, so restarts resume the firehose where the +// previous run left off +// - mirrored Tangled state (spindle members, repos, collaborators) +// and the authorization helpers built on top of it, used both to +// gate inbound pipeline triggers and to decide which knots we +// should be subscribed to +// - the outbound event log fed to /events websocket subscribers +// - the Buildkite build → pipeline mapping the webhook and /logs +// handlers look up // -// 1. The jetstream cursor — a microsecond unix timestamp the AT Proto -// firehose uses to resume from a specific point. Without persistence -// every restart begins at "now" and we'd silently miss any record -// published while we were down. -// -// 2. Tangled membership state derived from jetstream commits: -// sh.tangled.spindle.member, sh.tangled.repo, and -// sh.tangled.repo.collaborator. We need this to later answer -// "is DID X allowed to trigger a build on this spindle for repo Y?" -// -// We use mattn/go-sqlite3, which is the most battle-tested SQLite driver -// for Go. It requires CGo, which is fine for tack — the project already -// builds under Nix where a C toolchain is readily available. +// Per-method docs go into more detail. import ( "context" @@ -294,15 +290,30 @@ func (s *store) AuthorizePipelineActor( return true, "", 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) { +// IsKnotWanted reports whether any *authorized* 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. +// +// "Authorized" means the repo's publisher (the row's `did` column) is +// either the spindle owner or has been vouched for by the spindle owner +// via a sh.tangled.spindle.member record. Without this filter, a non- +// member could pin us to an arbitrary attacker-chosen knot just by +// publishing a sh.tangled.repo record naming us as its spindle. See +// the matching gate in IsAuthorizedActor and AuthorizePipelineActor. +func (s *store) IsKnotWanted(ctx context.Context, hostname, ownerDID, knot string) (bool, error) { var n int err := s.db.QueryRowContext(ctx, - `SELECT COUNT(*) FROM repos WHERE spindle = ? AND knot = ?`, - hostname, knot, + `SELECT COUNT(*) FROM repos r + WHERE r.spindle = ? AND r.knot = ? + AND ( + r.did = ? + OR EXISTS ( + SELECT 1 FROM spindle_members m + WHERE m.did = ? AND m.subject = r.did + ) + )`, + hostname, knot, ownerDID, ownerDID, ).Scan(&n) if err != nil { return false, fmt.Errorf("count repos for knot: %w", err) @@ -310,16 +321,110 @@ func (s *store) IsKnotWanted(ctx context.Context, hostname, knot string) (bool, return n > 0, nil } +// IsAuthorizedActor reports whether did is the spindle owner or has +// been authorized by the spindle owner via a sh.tangled.spindle.member +// record. The membership record's publisher (its `did` column) must +// equal ownerDID; anyone can publish a membership record naming +// anyone, so trusting unsigned grants would let any DID grant itself +// access. This is the same trust rule AuthorizePipelineActor enforces; +// it's pulled out here so knot-subscription decisions in jetstream.go +// gate on the same check as pipeline-spawning decisions in knot.go. +func (s *store) IsAuthorizedActor(ctx context.Context, ownerDID, did string) (bool, error) { + if did == "" { + return false, nil + } + if did == ownerDID { + return true, nil + } + var n int + err := s.db.QueryRowContext(ctx, + `SELECT COUNT(*) FROM spindle_members + WHERE did = ? AND subject = ?`, + ownerDID, did, + ).Scan(&n) + if err != nil { + return false, fmt.Errorf("count membership: %w", err) + } + return n > 0, nil +} + +// GetSpindleMember returns the subject DID currently stored for a +// (did, rkey) spindle.member row, or "" when no such row exists. +// +// Used by the jetstream handler to learn whose authorization a +// delete/update of a membership record affects, so it can reconcile +// that subject's knot subscriptions in the same step as the mutation. +func (s *store) GetSpindleMember(ctx context.Context, did, rkey string) (string, error) { + var subject string + err := s.db.QueryRowContext(ctx, + `SELECT subject FROM spindle_members WHERE did = ? AND rkey = ?`, + did, rkey, + ).Scan(&subject) + if errors.Is(err, sql.ErrNoRows) { + return "", nil + } + if err != nil { + return "", fmt.Errorf("get spindle_member: %w", err) + } + return subject, nil +} + +// KnotsForOwner returns the distinct knot hostnames of repos published +// by `did` whose spindle field equals hostname. Used after a membership +// change to find which knots that DID's repos want, so we can subscribe +// (newly granted) or potentially unsubscribe (revoked) in lockstep with +// the grant. +// +// Returns an empty slice (not nil) on no matches so callers can range +// over the result without a nil check. +func (s *store) KnotsForOwner(ctx context.Context, hostname, did string) ([]string, error) { + rows, err := s.db.QueryContext(ctx, + `SELECT DISTINCT knot FROM repos + WHERE did = ? AND spindle = ? AND knot <> ''`, + did, hostname, + ) + if err != nil { + return nil, fmt.Errorf("query knots for owner: %w", err) + } + defer rows.Close() + out := []string{} + for rows.Next() { + var k string + if err := rows.Scan(&k); err != nil { + return nil, fmt.Errorf("scan knot: %w", err) + } + out = append(out, k) + } + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("iterate knots: %w", err) + } + return out, 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. +// (a) have declared the given spindle hostname as their CI spindle, and +// (b) were published by an *authorized* DID: either the spindle owner +// itself or a member they vouched for via sh.tangled.spindle.member. +// +// The membership filter is critical: without it, any DID that publishes +// a sh.tangled.repo record naming us as its spindle could force us to +// dial an attacker-chosen knot at startup. See IsAuthorizedActor for +// the matching trust check applied at firehose-event time. // // Returns an empty slice (not nil) when nothing matches, so callers can // range over the result without a nil check. -func (s *store) KnotsForSpindle(ctx context.Context, hostname string) ([]string, error) { +func (s *store) KnotsForSpindle(ctx context.Context, hostname, ownerDID string) ([]string, error) { rows, err := s.db.QueryContext(ctx, - `SELECT DISTINCT knot FROM repos WHERE spindle = ? AND knot <> ''`, - hostname, + `SELECT DISTINCT r.knot FROM repos r + WHERE r.spindle = ? AND r.knot <> '' + AND ( + r.did = ? + OR EXISTS ( + SELECT 1 FROM spindle_members m + WHERE m.did = ? AND m.subject = r.did + ) + )`, + hostname, ownerDID, ownerDID, ) if err != nil { return nil, fmt.Errorf("query knots: %w", err) diff --git a/store_test.go b/store_test.go index ba18f46..fdb0b67 100644 --- a/store_test.go +++ b/store_test.go @@ -241,36 +241,58 @@ func TestRepoCollaboratorLifecycle(t *testing.T) { } // TestKnotsForSpindle verifies the query returns only knots from repos -// whose .spindle field matches the given hostname, and that duplicate -// knots collapse to a single entry. +// whose .spindle field matches the given hostname *and* whose publisher +// is an authorized actor (owner or owner-vouched member). Duplicate +// knots collapse to a single entry. Repos from non-members are excluded: +// that's the gate that prevents an attacker-published repo record from +// forcing us to dial an attacker-chosen knot. func TestKnotsForSpindle(t *testing.T) { s := newTestStore(t) ctx := context.Background() const ours = "tack.example" const other = "other.example" + const owner = "did:plc:owner" - // Two repos on the same knot pointing at us — should collapse to 1. + // Vouch for did:plc:a and did:plc:b; did:plc:c is the owner so + // no membership grant is needed; did:plc:zzz is unvouched. + if err := s.UpsertSpindleMember(ctx, owner, "mk1", ours, "did:plc:a", "t"); err != nil { + t.Fatal(err) + } + if err := s.UpsertSpindleMember(ctx, owner, "mk2", ours, "did:plc:b", "t"); err != nil { + t.Fatal(err) + } + + // Two member-published repos on the same knot pointing at us; + // should collapse to one entry. if err := s.UpsertRepo(ctx, "did:plc:a", "rk1", "knot1.example", "repo-a", ours, "", "t"); err != nil { t.Fatal(err) } if err := s.UpsertRepo(ctx, "did:plc:b", "rk2", "knot1.example", "repo-b", ours, "", "t"); err != nil { t.Fatal(err) } - // A second knot pointing at us. - if err := s.UpsertRepo(ctx, "did:plc:c", "rk3", "knot2.example", "repo-c", ours, "", "t"); err != nil { + // Owner-published repo on a second knot. Owner is implicitly + // authorized and must be included. + if err := s.UpsertRepo(ctx, owner, "rk3", "knot2.example", "repo-c", ours, "", "t"); err != nil { + t.Fatal(err) + } + // Repo pointing at a different spindle: must be excluded. + if err := s.UpsertRepo(ctx, "did:plc:a", "rk4", "knot3.example", "repo-d", other, "", "t"); err != nil { t.Fatal(err) } - // A repo pointing at a different spindle — must be excluded. - if err := s.UpsertRepo(ctx, "did:plc:d", "rk4", "knot3.example", "repo-d", other, "", "t"); err != nil { + // Repo with no spindle declared: must be excluded. + if err := s.UpsertRepo(ctx, "did:plc:a", "rk5", "knot4.example", "repo-e", "", "", "t"); err != nil { t.Fatal(err) } - // A repo with no spindle declared — must be excluded. - if err := s.UpsertRepo(ctx, "did:plc:e", "rk5", "knot4.example", "repo-e", "", "", "t"); err != nil { + // Repo published by an unvouched DID: must be excluded even + // though it points at us. This is the security-relevant case: + // without the membership filter, a stranger could pin us to + // "evil.example" just by publishing a sh.tangled.repo record. + if err := s.UpsertRepo(ctx, "did:plc:zzz", "rk6", "evil.example", "repo-z", ours, "", "t"); err != nil { t.Fatal(err) } - got, err := s.KnotsForSpindle(ctx, ours) + got, err := s.KnotsForSpindle(ctx, ours, owner) if err != nil { t.Fatalf("KnotsForSpindle: %v", err) }