diff --git a/spindle/db/member.go b/spindle/db/member.go deleted file mode 100644 index 31fbfe19..00000000 --- a/spindle/db/member.go +++ /dev/null @@ -1,63 +0,0 @@ -package db - -import ( - "github.com/bluesky-social/indigo/atproto/syntax" -) - -type SpindleMember struct { - Id int - Did syntax.DID // owner of the record - Rkey string // rkey of the record - Instance string - Subject syntax.DID // the member being added -} - -func AddSpindleMember(q DBTX, member SpindleMember) error { - _, err := q.Exec( - `insert or ignore into spindle_members (did, rkey, instance, subject) values (?, ?, ?, ?)`, - member.Did, - member.Rkey, - member.Instance, - member.Subject, - ) - return err -} - -func RemoveSpindleMember(q DBTX, ownerDid, rkey string) error { - _, err := q.Exec( - "delete from spindle_members where did = ? and rkey = ?", - ownerDid, - rkey, - ) - return err -} - -func CountSpindleMembersBySubject(q DBTX, subject string) (int, error) { - var count int - err := q.QueryRow( - `select count(*) from spindle_members where subject = ?`, - subject, - ).Scan(&count) - return count, err -} - -func GetSpindleMember(q DBTX, did, rkey string) (*SpindleMember, error) { - query := - `select id, did, rkey, instance, subject - from spindle_members - where did = ? and rkey = ?` - - var member SpindleMember - err := q.QueryRow(query, did, rkey).Scan( - &member.Id, - &member.Did, - &member.Rkey, - &member.Instance, - &member.Subject, - ) - if err != nil { - return nil, err - } - - return &member, nil -} diff --git a/spindle/ingester.go b/spindle/ingester.go index 4f2909f4..6ecd681d 100644 --- a/spindle/ingester.go +++ b/spindle/ingester.go @@ -2,13 +2,8 @@ package spindle import ( "context" - "database/sql" - "encoding/json" - "errors" - "fmt" "tangled.org/core/api/tangled" - "tangled.org/core/spindle/db" "tangled.org/core/tapc" "github.com/bluesky-social/indigo/atproto/syntax" @@ -25,8 +20,6 @@ func (s *Spindle) ingest() Ingester { var err error switch e.Commit.Collection { - case tangled.SpindleMemberNSID: - err = s.ingestMember(ctx, e) case tangled.RepoNSID, tangled.RepoCollaboratorNSID: if evt, ok := jetstreamToTapEvent(e); ok { err = s.tap.processEvent(ctx, evt) @@ -77,184 +70,3 @@ func jetstreamToTapEvent(e *models.Event) (tapc.Event, bool) { }, }, true } - -func (s *Spindle) ingestMember(ctx context.Context, e *models.Event) error { - did := e.Did - rkey := e.Commit.RKey - l := s.l.With("component", "ingester", "record", tangled.SpindleMemberNSID, "did", did, "rkey", rkey) - - switch e.Commit.Operation { - case models.CommitOperationCreate, models.CommitOperationUpdate: - raw := e.Commit.Record - record := tangled.SpindleMember{} - if err := json.Unmarshal(raw, &record); err != nil { - return fmt.Errorf("invalid record: %w", err) - } - - domain := s.cfg.Server.Hostname - recordInstance := record.Instance - - if recordInstance != domain { - return fmt.Errorf("domain mismatch: %s != %s", record.Instance, domain) - } - - subject, err := syntax.ParseDID(record.Subject) - if err != nil { - return fmt.Errorf("invalid subject DID %q: %w", record.Subject, err) - } - - ok, err := s.e.IsSpindleInviteAllowed(did, rbacDomain) - if err != nil { - return fmt.Errorf("failed to enforce permissions: %w", err) - } - if !ok { - return fmt.Errorf("permission denied for %s", did) - } - - sqlTx, err := s.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.GetSpindleMember(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.RemoveSpindleMember(sqlTx, did, rkey); err != nil { - return fmt.Errorf("failed to remove stale member row: %w", err) - } - } - - if err := db.AddSpindleMember(sqlTx, db.SpindleMember{ - Did: syntax.DID(did), - Rkey: rkey, - Instance: recordInstance, - Subject: subject, - }); err != nil { - return fmt.Errorf("failed to add member: %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.CountSpindleMembersBySubject(sqlTx, staleSubject) - if err != nil { - return fmt.Errorf("failed to count stale subject rows: %w", err) - } - if remaining == 0 { - dropStaleAcl = true - stillNeeded, err := s.e.WouldHaveAnyPolicyExcludingSpindleMember(staleSubject, rbacDomain) - 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 spindle 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 := s.e.TryRemoveSpindleMember(rbacDomain, staleSubject); err != nil { - l.Error("post-commit: failed to remove stale ACL", "subject", staleSubject, "err", err) - } - } - if _, err := s.e.TryAddSpindleMember(rbacDomain, subject.String()); err != nil { - l.Error("post-commit: failed to add member ACL", "subject", subject, "err", err) - } - - if staleDidDropped { - s.jc.RemoveDid(staleSubject) - } - s.jc.AddDid(subject.String()) - l.Info("added member from firehose", "member", subject) - return nil - - case models.CommitOperationDelete: - sqlTx, err := s.db.BeginTx(ctx, nil) - if err != nil { - return fmt.Errorf("failed to start txn: %w", err) - } - committed := false - defer func() { - if !committed { - sqlTx.Rollback() - } - }() - - record, err := db.GetSpindleMember(sqlTx, did, rkey) - if errors.Is(err, sql.ErrNoRows) { - l.Info("spindle member already removed") - return nil - } - if err != nil { - return fmt.Errorf("failed to find member: %w", err) - } - - staleSubject := record.Subject.String() - - if err := db.RemoveSpindleMember(sqlTx, did, rkey); err != nil { - return fmt.Errorf("failed to remove member: %w", err) - } - - remaining, err := db.CountSpindleMembersBySubject(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 := s.e.WouldHaveAnyPolicyExcludingSpindleMember(staleSubject, rbacDomain) - 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 := s.e.TryRemoveSpindleMember(rbacDomain, staleSubject); err != nil { - l.Error("post-commit: failed to remove member ACL", "subject", staleSubject, "err", err) - } - } - - if staleDidDropped { - s.jc.RemoveDid(staleSubject) - } - l.Info("removed member from firehose", "member", record.Subject, "remaining_rows", remaining, "stale_did_dropped", staleDidDropped) - } - return nil -}