Monorepo for Tangled tangled.org
4

Configure Feed

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

core / spindle / ingester.go
7.2 kB 260 lines
1package spindle 2 3import ( 4 "context" 5 "database/sql" 6 "encoding/json" 7 "errors" 8 "fmt" 9 10 "tangled.org/core/api/tangled" 11 "tangled.org/core/spindle/db" 12 "tangled.org/core/tapc" 13 14 "github.com/bluesky-social/indigo/atproto/syntax" 15 "github.com/bluesky-social/jetstream/pkg/models" 16) 17 18type Ingester func(ctx context.Context, e *models.Event) error 19 20func (s *Spindle) ingest() Ingester { 21 return func(ctx context.Context, e *models.Event) error { 22 if e.Kind != models.EventKindCommit { 23 return nil 24 } 25 26 var err error 27 switch e.Commit.Collection { 28 case tangled.SpindleMemberNSID: 29 err = s.ingestMember(ctx, e) 30 case tangled.RepoNSID, tangled.RepoCollaboratorNSID: 31 if evt, ok := jetstreamToTapEvent(e); ok { 32 err = s.tap.processEvent(ctx, evt) 33 } 34 case tangled.RepoPullNSID: 35 if evt, ok := jetstreamToTapEvent(e); ok { 36 err = s.processPull(ctx, evt.Record) 37 } 38 } 39 40 if err != nil { 41 s.l.Warn("failed to process message, skipping", "nsid", e.Commit.Collection, "did", e.Did, "rkey", e.Commit.RKey, "err", err) 42 } 43 44 return nil 45 } 46} 47 48func jetstreamToTapEvent(e *models.Event) (tapc.Event, bool) { 49 if e.Commit == nil { 50 return tapc.Event{}, false 51 } 52 did, err := syntax.ParseDID(e.Did) 53 if err != nil { 54 return tapc.Event{}, false 55 } 56 var action tapc.RecordAction 57 switch e.Commit.Operation { 58 case models.CommitOperationCreate: 59 action = tapc.RecordCreateAction 60 case models.CommitOperationUpdate: 61 action = tapc.RecordUpdateAction 62 case models.CommitOperationDelete: 63 action = tapc.RecordDeleteAction 64 default: 65 return tapc.Event{}, false 66 } 67 return tapc.Event{ 68 Type: tapc.EvtRecord, 69 Record: &tapc.RecordEventData{ 70 Did: did, 71 Rkey: syntax.RecordKey(e.Commit.RKey), 72 Collection: syntax.NSID(e.Commit.Collection), 73 Action: action, 74 Record: e.Commit.Record, 75 // jetstream is only used for live 76 Live: true, 77 }, 78 }, true 79} 80 81func (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}