Monorepo for Tangled tangled.org
2

Configure Feed

Select the types of activity you want to include in your feed.

spindle: remove `sh.tangled.spindle.member` ingest

Signed-off-by: Seongmin Lee <git@boltless.me>

author
Seongmin Lee
date (Jul 27, 2026, 7:59 PM +0900) commit c47638c2 parent af459bd9 change-id tmtxkynp
-251
-63
spindle/db/member.go
··· 1 - package db 2 - 3 - import ( 4 - "github.com/bluesky-social/indigo/atproto/syntax" 5 - ) 6 - 7 - type SpindleMember struct { 8 - Id int 9 - Did syntax.DID // owner of the record 10 - Rkey string // rkey of the record 11 - Instance string 12 - Subject syntax.DID // the member being added 13 - } 14 - 15 - func AddSpindleMember(q DBTX, member SpindleMember) error { 16 - _, err := q.Exec( 17 - `insert or ignore into spindle_members (did, rkey, instance, subject) values (?, ?, ?, ?)`, 18 - member.Did, 19 - member.Rkey, 20 - member.Instance, 21 - member.Subject, 22 - ) 23 - return err 24 - } 25 - 26 - func RemoveSpindleMember(q DBTX, ownerDid, rkey string) error { 27 - _, err := q.Exec( 28 - "delete from spindle_members where did = ? and rkey = ?", 29 - ownerDid, 30 - rkey, 31 - ) 32 - return err 33 - } 34 - 35 - func CountSpindleMembersBySubject(q DBTX, subject string) (int, error) { 36 - var count int 37 - err := q.QueryRow( 38 - `select count(*) from spindle_members where subject = ?`, 39 - subject, 40 - ).Scan(&count) 41 - return count, err 42 - } 43 - 44 - func GetSpindleMember(q DBTX, did, rkey string) (*SpindleMember, error) { 45 - query := 46 - `select id, did, rkey, instance, subject 47 - from spindle_members 48 - where did = ? and rkey = ?` 49 - 50 - var member SpindleMember 51 - err := q.QueryRow(query, did, rkey).Scan( 52 - &member.Id, 53 - &member.Did, 54 - &member.Rkey, 55 - &member.Instance, 56 - &member.Subject, 57 - ) 58 - if err != nil { 59 - return nil, err 60 - } 61 - 62 - return &member, nil 63 - }
-188
spindle/ingester.go
··· 2 2 3 3 import ( 4 4 "context" 5 - "database/sql" 6 - "encoding/json" 7 - "errors" 8 - "fmt" 9 5 10 6 "tangled.org/core/api/tangled" 11 - "tangled.org/core/spindle/db" 12 7 "tangled.org/core/tapc" 13 8 14 9 "github.com/bluesky-social/indigo/atproto/syntax" ··· 25 20 26 21 var err error 27 22 switch e.Commit.Collection { 28 - case tangled.SpindleMemberNSID: 29 - err = s.ingestMember(ctx, e) 30 23 case tangled.RepoNSID, tangled.RepoCollaboratorNSID: 31 24 if evt, ok := jetstreamToTapEvent(e); ok { 32 25 err = s.tap.processEvent(ctx, evt) ··· 77 70 }, 78 71 }, true 79 72 } 80 - 81 - func (s *Spindle) ingestMember(ctx context.Context, e *models.Event) error { 82 - did := e.Did 83 - rkey := e.Commit.RKey 84 - l := s.l.With("component", "ingester", "record", tangled.SpindleMemberNSID, "did", did, "rkey", rkey) 85 - 86 - switch e.Commit.Operation { 87 - case models.CommitOperationCreate, models.CommitOperationUpdate: 88 - raw := e.Commit.Record 89 - record := tangled.SpindleMember{} 90 - if err := json.Unmarshal(raw, &record); err != nil { 91 - return fmt.Errorf("invalid record: %w", err) 92 - } 93 - 94 - domain := s.cfg.Server.Hostname 95 - recordInstance := record.Instance 96 - 97 - if recordInstance != domain { 98 - return fmt.Errorf("domain mismatch: %s != %s", record.Instance, domain) 99 - } 100 - 101 - subject, err := syntax.ParseDID(record.Subject) 102 - if err != nil { 103 - return fmt.Errorf("invalid subject DID %q: %w", record.Subject, err) 104 - } 105 - 106 - ok, err := s.e.IsSpindleInviteAllowed(did, rbacDomain) 107 - if err != nil { 108 - return fmt.Errorf("failed to enforce permissions: %w", err) 109 - } 110 - if !ok { 111 - return fmt.Errorf("permission denied for %s", did) 112 - } 113 - 114 - sqlTx, err := s.db.BeginTx(ctx, nil) 115 - if err != nil { 116 - return fmt.Errorf("failed to start txn: %w", err) 117 - } 118 - committed := false 119 - defer func() { 120 - if !committed { 121 - sqlTx.Rollback() 122 - } 123 - }() 124 - 125 - existing, err := db.GetSpindleMember(sqlTx, did, rkey) 126 - if err != nil && !errors.Is(err, sql.ErrNoRows) { 127 - return fmt.Errorf("failed to look up existing member: %w", err) 128 - } 129 - 130 - var staleSubject string 131 - if existing != nil && existing.Subject != subject { 132 - staleSubject = existing.Subject.String() 133 - if err := db.RemoveSpindleMember(sqlTx, did, rkey); err != nil { 134 - return fmt.Errorf("failed to remove stale member row: %w", err) 135 - } 136 - } 137 - 138 - if err := db.AddSpindleMember(sqlTx, db.SpindleMember{ 139 - Did: syntax.DID(did), 140 - Rkey: rkey, 141 - Instance: recordInstance, 142 - Subject: subject, 143 - }); err != nil { 144 - return fmt.Errorf("failed to add member: %w", err) 145 - } 146 - 147 - if err := db.AddDid(sqlTx, subject.String()); err != nil { 148 - return fmt.Errorf("failed to add did: %w", err) 149 - } 150 - 151 - dropStaleAcl := false 152 - var staleDidDropped bool 153 - if staleSubject != "" { 154 - remaining, err := db.CountSpindleMembersBySubject(sqlTx, staleSubject) 155 - if err != nil { 156 - return fmt.Errorf("failed to count stale subject rows: %w", err) 157 - } 158 - if remaining == 0 { 159 - dropStaleAcl = true 160 - stillNeeded, err := s.e.WouldHaveAnyPolicyExcludingSpindleMember(staleSubject, rbacDomain) 161 - if err != nil { 162 - return fmt.Errorf("failed to check residual policies for stale subject: %w", err) 163 - } 164 - if !stillNeeded { 165 - if err := db.RemoveDid(sqlTx, staleSubject); err != nil { 166 - return fmt.Errorf("failed to remove stale did: %w", err) 167 - } 168 - staleDidDropped = true 169 - } 170 - } 171 - l.Info("replaced stale spindle member", "old_subject", staleSubject, "new_subject", subject, "stale_did_dropped", staleDidDropped) 172 - } 173 - 174 - if err := sqlTx.Commit(); err != nil { 175 - return fmt.Errorf("failed to commit txn: %w", err) 176 - } 177 - committed = true 178 - 179 - if dropStaleAcl { 180 - if _, err := s.e.TryRemoveSpindleMember(rbacDomain, staleSubject); err != nil { 181 - l.Error("post-commit: failed to remove stale ACL", "subject", staleSubject, "err", err) 182 - } 183 - } 184 - if _, err := s.e.TryAddSpindleMember(rbacDomain, subject.String()); err != nil { 185 - l.Error("post-commit: failed to add member ACL", "subject", subject, "err", err) 186 - } 187 - 188 - if staleDidDropped { 189 - s.jc.RemoveDid(staleSubject) 190 - } 191 - s.jc.AddDid(subject.String()) 192 - l.Info("added member from firehose", "member", subject) 193 - return nil 194 - 195 - case models.CommitOperationDelete: 196 - sqlTx, err := s.db.BeginTx(ctx, nil) 197 - if err != nil { 198 - return fmt.Errorf("failed to start txn: %w", err) 199 - } 200 - committed := false 201 - defer func() { 202 - if !committed { 203 - sqlTx.Rollback() 204 - } 205 - }() 206 - 207 - record, err := db.GetSpindleMember(sqlTx, did, rkey) 208 - if errors.Is(err, sql.ErrNoRows) { 209 - l.Info("spindle member already removed") 210 - return nil 211 - } 212 - if err != nil { 213 - return fmt.Errorf("failed to find member: %w", err) 214 - } 215 - 216 - staleSubject := record.Subject.String() 217 - 218 - if err := db.RemoveSpindleMember(sqlTx, did, rkey); err != nil { 219 - return fmt.Errorf("failed to remove member: %w", err) 220 - } 221 - 222 - remaining, err := db.CountSpindleMembersBySubject(sqlTx, staleSubject) 223 - if err != nil { 224 - return fmt.Errorf("failed to count remaining member rows: %w", err) 225 - } 226 - 227 - dropAcl := false 228 - var staleDidDropped bool 229 - if remaining == 0 { 230 - dropAcl = true 231 - stillNeeded, err := s.e.WouldHaveAnyPolicyExcludingSpindleMember(staleSubject, rbacDomain) 232 - if err != nil { 233 - return fmt.Errorf("failed to check residual policies: %w", err) 234 - } 235 - if !stillNeeded { 236 - if err := db.RemoveDid(sqlTx, staleSubject); err != nil { 237 - return fmt.Errorf("failed to remove did: %w", err) 238 - } 239 - staleDidDropped = true 240 - } 241 - } 242 - 243 - if err := sqlTx.Commit(); err != nil { 244 - return fmt.Errorf("failed to commit txn: %w", err) 245 - } 246 - committed = true 247 - 248 - if dropAcl { 249 - if _, err := s.e.TryRemoveSpindleMember(rbacDomain, staleSubject); err != nil { 250 - l.Error("post-commit: failed to remove member ACL", "subject", staleSubject, "err", err) 251 - } 252 - } 253 - 254 - if staleDidDropped { 255 - s.jc.RemoveDid(staleSubject) 256 - } 257 - l.Info("removed member from firehose", "member", record.Subject, "remaining_rows", remaining, "stale_did_dropped", staleDidDropped) 258 - } 259 - return nil 260 - }