diff --git a/knotserver/ingester.go b/knotserver/ingester.go --- a/knotserver/ingester.go +++ b/knotserver/ingester.go @@ -2,9 +2,7 @@ import ( "context" - "database/sql" "encoding/json" - "errors" "fmt" "io" "net/http" @@ -15,7 +13,6 @@ comatproto "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/atproto/syntax" "github.com/bluesky-social/indigo/xrpc" - indigoxrpc "github.com/bluesky-social/indigo/xrpc" jmodels "github.com/bluesky-social/jetstream/pkg/models" "tangled.org/core/api/tangled" "tangled.org/core/appview/models" @@ -24,7 +21,6 @@ "tangled.org/core/knotserver/git" knotxrpc "tangled.org/core/knotserver/xrpc" "tangled.org/core/log" - "tangled.org/core/rbac" "tangled.org/core/tid" "tangled.org/core/workflow" ) @@ -48,189 +44,6 @@ return fmt.Errorf("failed to add public key: %w", err) } l.Info("added public key from firehose", "did", did) - return nil -} - -func (h *Knot) processKnotMember(ctx context.Context, event *jmodels.Event) error { - did := event.Did - rkey := event.Commit.RKey - l := log.FromContext(ctx).With("handler", "processKnotMember", "did", did, "rkey", rkey) - - switch event.Commit.Operation { - case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: - raw := json.RawMessage(event.Commit.Record) - var record tangled.KnotMember - if err := json.Unmarshal(raw, &record); err != nil { - return fmt.Errorf("failed to unmarshal record: %w", err) - } - - if record.Domain != h.c.Server.Hostname { - return fmt.Errorf("domain mismatch: %s != %s", record.Domain, h.c.Server.Hostname) - } - - subject, err := syntax.ParseDID(record.Subject) - if err != nil { - return fmt.Errorf("invalid subject DID %q: %w", record.Subject, err) - } - - ok, err := h.e.E.Enforce(did, rbac.ThisServer, rbac.ThisServer, "server:invite") - if err != nil { - return fmt.Errorf("failed to enforce permissions: %w", err) - } - if !ok { - return fmt.Errorf("permission denied for %s", did) - } - - sqlTx, err := h.db.BeginTx(ctx, nil) - if err != nil { - return fmt.Errorf("failed to start txn: %w", err) - } - committed := false - defer func() { - if !committed { - sqlTx.Rollback() - } - }() - - existing, err := db.GetKnotMember(sqlTx, did, rkey) - if err != nil && !errors.Is(err, sql.ErrNoRows) { - return fmt.Errorf("failed to look up existing member: %w", err) - } - - var staleSubject string - if existing != nil && existing.Subject != subject { - staleSubject = existing.Subject.String() - if err := db.RemoveKnotMember(sqlTx, did, rkey); err != nil { - return fmt.Errorf("failed to remove stale member row: %w", err) - } - } - - if err := db.AddKnotMember(sqlTx, db.KnotMember{ - Did: syntax.DID(did), - Rkey: rkey, - Subject: subject, - }); err != nil { - return fmt.Errorf("failed to persist member row: %w", err) - } - - if err := db.AddDid(sqlTx, subject.String()); err != nil { - return fmt.Errorf("failed to add did: %w", err) - } - - dropStaleAcl := false - var staleDidDropped bool - if staleSubject != "" { - remaining, err := db.CountKnotMembersBySubject(sqlTx, staleSubject) - if err != nil { - return fmt.Errorf("failed to count stale subject rows: %w", err) - } - if remaining == 0 { - dropStaleAcl = true - stillNeeded, err := h.e.WouldHaveAnyPolicyExcludingKnotMember(staleSubject, rbac.ThisServer) - if err != nil { - return fmt.Errorf("failed to check residual policies for stale subject: %w", err) - } - if !stillNeeded { - if err := db.RemoveDid(sqlTx, staleSubject); err != nil { - return fmt.Errorf("failed to remove stale did: %w", err) - } - staleDidDropped = true - } - } - l.Info("replaced stale knot member", "old_subject", staleSubject, "new_subject", subject, "stale_did_dropped", staleDidDropped) - } - - if err := sqlTx.Commit(); err != nil { - return fmt.Errorf("failed to commit txn: %w", err) - } - committed = true - - if dropStaleAcl { - if _, err := h.e.TryRemoveKnotMember(rbac.ThisServer, staleSubject); err != nil { - l.Error("post-commit: failed to remove stale ACL", "subject", staleSubject, "err", err) - } - } - if _, err := h.e.TryAddKnotMember(rbac.ThisServer, subject.String()); err != nil { - l.Error("post-commit: failed to add member ACL", "subject", subject, "err", err) - } - - if staleDidDropped { - h.jc.RemoveDid(staleSubject) - } - h.jc.AddDid(subject.String()) - l.Info("added member from firehose", "member", subject) - - if err := h.fetchAndAddKeys(ctx, subject.String()); err != nil { - return fmt.Errorf("failed to fetch and add keys: %w", err) - } - return nil - - case jmodels.CommitOperationDelete: - sqlTx, err := h.db.BeginTx(ctx, nil) - if err != nil { - return fmt.Errorf("failed to start txn: %w", err) - } - committed := false - defer func() { - if !committed { - sqlTx.Rollback() - } - }() - - member, err := db.GetKnotMember(sqlTx, did, rkey) - if errors.Is(err, sql.ErrNoRows) { - l.Info("knot member already removed") - return nil - } - if err != nil { - return fmt.Errorf("failed to look up knot member: %w", err) - } - - staleSubject := member.Subject.String() - - if err := db.RemoveKnotMember(sqlTx, did, rkey); err != nil { - return fmt.Errorf("failed to remove member row: %w", err) - } - - remaining, err := db.CountKnotMembersBySubject(sqlTx, staleSubject) - if err != nil { - return fmt.Errorf("failed to count remaining member rows: %w", err) - } - - dropAcl := false - var staleDidDropped bool - if remaining == 0 { - dropAcl = true - stillNeeded, err := h.e.WouldHaveAnyPolicyExcludingKnotMember(staleSubject, rbac.ThisServer) - if err != nil { - return fmt.Errorf("failed to check residual policies: %w", err) - } - if !stillNeeded { - if err := db.RemoveDid(sqlTx, staleSubject); err != nil { - return fmt.Errorf("failed to remove did: %w", err) - } - staleDidDropped = true - } - } - - if err := sqlTx.Commit(); err != nil { - return fmt.Errorf("failed to commit txn: %w", err) - } - committed = true - - if dropAcl { - if _, err := h.e.TryRemoveKnotMember(rbac.ThisServer, staleSubject); err != nil { - l.Error("post-commit: failed to remove ACL", "subject", staleSubject, "err", err) - } - } - - if staleDidDropped { - h.jc.RemoveDid(staleSubject) - } - l.Info("removed knot member from firehose", "member", member.Subject, "remaining_rows", remaining, "stale_did_dropped", staleDidDropped) - return nil - } - return nil } @@ -513,128 +326,6 @@ return h.db.InsertEvent(ev, h.n) } -// duplicated from add collaborator -func (h *Knot) processCollaborator(ctx context.Context, event *jmodels.Event) error { - raw := json.RawMessage(event.Commit.Record) - did := event.Did - - var record tangled.RepoCollaborator - if err := json.Unmarshal(raw, &record); err != nil { - return fmt.Errorf("failed to unmarshal record: %w", err) - } - - subjectId, err := h.resolver.ResolveIdent(ctx, record.Subject) - if err != nil || subjectId.Handle.IsInvalidHandle() { - return err - } - - var rbacResource string - switch { - case strings.HasPrefix(record.Repo, "did:"): - ownerDid, _, lookupErr := h.db.GetRepoKeyOwner(record.Repo) - if lookupErr != nil { - return fmt.Errorf("unknown repo DID %s: %w", record.Repo, lookupErr) - } - if ownerDid != did { - return fmt.Errorf("collaborator record author %s does not own repo %s", did, record.Repo) - } - rbacResource = record.Repo - - case strings.Contains(record.Repo, "/"): - // TODO: get rid of this PDS fetch once all repos have DIDs - repoAt, parseErr := syntax.ParseATURI(record.Repo) - if parseErr != nil { - return parseErr - } - - owner, resolveErr := h.resolver.ResolveIdent(ctx, repoAt.Authority().String()) - if resolveErr != nil || owner.Handle.IsInvalidHandle() { - return fmt.Errorf("failed to resolve handle: %w", resolveErr) - } - - xrpcc := xrpc.Client{ - Host: owner.PDSEndpoint(), - } - - resp, getErr := comatproto.RepoGetRecord(ctx, &xrpcc, "", tangled.RepoNSID, repoAt.Authority().String(), repoAt.RecordKey().String()) - if getErr != nil { - return getErr - } - - if _, ok := resp.Value.Val.(*tangled.Repo); !ok { - return fmt.Errorf("record at %s is not a tangled.Repo", repoAt) - } - rkey := repoAt.RecordKey().String() - repoDid, didErr := h.db.GetRepoDid(owner.DID.String(), rkey) - if didErr != nil { - return fmt.Errorf("failed to resolve repo DID for %s/%s: %w", owner.DID.String(), rkey, didErr) - } - rbacResource = repoDid - - default: - return fmt.Errorf("collaborator record has unrecognized repo format: %s", record.Repo) - } - - ok, err := h.e.IsCollaboratorInviteAllowed(did, rbac.ThisServer, rbacResource) - if err != nil { - return fmt.Errorf("failed to check permissions: %w", err) - } - if !ok { - return fmt.Errorf("insufficient permissions: %s, %s, %s", did, "IsCollaboratorInviteAllowed", rbacResource) - } - - if err := db.AddDid(h.db, subjectId.DID.String()); err != nil { - return err - } - h.jc.AddDid(subjectId.DID.String()) - - if err := h.e.AddCollaborator(subjectId.DID.String(), rbac.ThisServer, rbacResource); err != nil { - return err - } - - return h.fetchAndAddKeys(ctx, subjectId.DID.String()) -} - -func (h *Knot) fetchAndAddKeys(ctx context.Context, did string) error { - l := log.FromContext(ctx) - - id, err := h.resolver.Directory().LookupDID(ctx, syntax.DID(did)) - if err != nil { - return fmt.Errorf("lookup did to fetch keys: %w", err) - } - - serviceEndpoint, ok := id.Services["atproto_pds"] - if !ok { - l.Warn("did identity did not contain atproto_pds service while adding their keys", "did", did) - return nil - } - - xrpcc := indigoxrpc.Client{Host: serviceEndpoint.URL} - resp, err := comatproto.RepoListRecords(context.Background(), &xrpcc, tangled.PublicKeyNSID, "", 50, did, false) - if err != nil { - return fmt.Errorf("fetching public keys for did: %w", err) - } - - for _, record := range resp.Records { - if record == nil { - continue - } - key := record.Value.Val.(*tangled.PublicKey) - if key == nil { - continue - } - pk := db.PublicKey{ - Did: did, - PublicKey: *key, - } - err = h.db.AddPublicKey(pk) - if err != nil { - return fmt.Errorf("adding public key to db: %w", err) - } - } - return nil -} - func (h *Knot) processRepo(ctx context.Context, event *jmodels.Event) error { l := log.FromContext(ctx).With("handler", "processRepo", "did", event.Did, "rkey", event.Commit.RKey) @@ -705,14 +396,10 @@ switch event.Commit.Collection { case tangled.PublicKeyNSID: err = h.processPublicKey(ctx, event) - case tangled.KnotMemberNSID: - err = h.processKnotMember(ctx, event) case tangled.RepoNSID: err = h.processRepo(ctx, event) case tangled.RepoPullNSID: err = h.processPull(ctx, event) - case tangled.RepoCollaboratorNSID: - err = h.processCollaborator(ctx, event) } default: return nil diff --git a/knotserver/router.go b/knotserver/router.go --- a/knotserver/router.go +++ b/knotserver/router.go @@ -14,6 +14,7 @@ "tangled.org/core/jetstream" "tangled.org/core/knotserver/config" "tangled.org/core/knotserver/db" + "tangled.org/core/knotserver/keys" "tangled.org/core/knotserver/sandbox" "tangled.org/core/knotserver/xrpc" "tangled.org/core/log" @@ -107,7 +108,9 @@ }) // xrpc apis - r.Mount("/xrpc", h.XrpcRouter()) + x := h.newXrpc() + r.Mount("/xrpc", x.Router()) + r.Mount("/admin", x.AdminRouter()) // Socket that streams git oplogs r.Get("/events", h.Events) @@ -129,12 +132,12 @@ return h.motd } -func (h *Knot) XrpcRouter() http.Handler { +func (h *Knot) newXrpc() *xrpc.Xrpc { serviceAuth := serviceauth.NewServiceAuth(h.l, h.resolver.Directory(), h.c.Server.Did().String()) l := log.SubLogger(h.l, "xrpc") - xrpc := &xrpc.Xrpc{ + return &xrpc.Xrpc{ Config: h.c, Db: h.db, Ingester: h.jc, @@ -145,8 +148,6 @@ ServiceAuth: serviceAuth, Sandbox: h.sandbox, } - - return xrpc.Router() } func (h *Knot) resolveDidRedirect(next http.Handler) http.Handler { @@ -216,7 +217,7 @@ return fmt.Errorf("failed to add owner to RBAC: %w", err) } - err = h.fetchAndAddKeys(ctx, cfgOwner) + err = keys.FetchAndStore(ctx, h.resolver.Directory(), h.db, cfgOwner) if err != nil { h.l.Error("fetching and adding owners public keys", "error", err, "did", cfgOwner) } diff --git a/knotserver/server.go b/knotserver/server.go --- a/knotserver/server.go +++ b/knotserver/server.go @@ -95,10 +95,8 @@ jc, err := jetstream.NewJetstreamClient(c.Server.JetstreamEndpoint, "knotserver", []string{ tangled.PublicKeyNSID, - tangled.KnotMemberNSID, tangled.RepoNSID, tangled.RepoPullNSID, - tangled.RepoCollaboratorNSID, }, nil, log.SubLogger(logger, "jetstream"), db, true, c.Server.LogDids) if err != nil { logger.Error("failed to setup jetstream", "error", err) diff --git a/rbac/util.go b/rbac/util.go --- a/rbac/util.go +++ b/rbac/util.go @@ -95,10 +95,6 @@ return false, nil } -func (e *Enforcer) WouldHaveAnyPolicyExcludingKnotMember(user, domain string) (bool, error) { - return e.wouldHaveAnyPolicyExcludingGrouping(user, "server:member", domain) -} - func (e *Enforcer) WouldHaveAnyPolicyExcludingSpindleMember(user, domain string) (bool, error) { return e.wouldHaveAnyPolicyExcludingGrouping(user, "server:member", intoSpindle(domain)) } diff --git a/knotserver/db/member.go b/knotserver/db/member.go --- a/knotserver/db/member.go +++ b/knotserver/db/member.go @@ -11,7 +11,6 @@ type KnotMember struct { Id int Did syntax.DID - Rkey string Subject syntax.DID Created string } @@ -92,43 +91,4 @@ }, func(m KnotMember) int { return m.Id }, ) -} - -func AddKnotMember(q DBTX, member KnotMember) error { - _, err := q.Exec( - `insert or ignore into knot_members (did, rkey, subject) values (?, ?, ?)`, - member.Did, - member.Rkey, - member.Subject, - ) - return err -} - -func RemoveKnotMember(q DBTX, ownerDid, rkey string) error { - _, err := q.Exec( - "delete from knot_members where did = ? and rkey = ?", - ownerDid, - rkey, - ) - return err -} - -func GetKnotMember(q DBTX, did, rkey string) (*KnotMember, error) { - query := - `select id, did, rkey, subject - from knot_members - where did = ? and rkey = ?` - - var member KnotMember - err := q.QueryRow(query, did, rkey).Scan( - &member.Id, - &member.Did, - &member.Rkey, - &member.Subject, - ) - if err != nil { - return nil, err - } - - return &member, nil }