This repository has no description
0

Configure Feed

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

core / spindle / ingester.go
7.1 kB 256 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, tangled.RepoPullNSID: 31 if evt, ok := jetstreamToTapEvent(e); ok { 32 err = s.tap.processEvent(ctx, evt) 33 } 34 } 35 36 if err != nil { 37 s.l.Warn("failed to process message, skipping", "nsid", e.Commit.Collection, "did", e.Did, "rkey", e.Commit.RKey, "err", err) 38 } 39 40 return nil 41 } 42} 43 44func jetstreamToTapEvent(e *models.Event) (tapc.Event, bool) { 45 if e.Commit == nil { 46 return tapc.Event{}, false 47 } 48 did, err := syntax.ParseDID(e.Did) 49 if err != nil { 50 return tapc.Event{}, false 51 } 52 var action tapc.RecordAction 53 switch e.Commit.Operation { 54 case models.CommitOperationCreate: 55 action = tapc.RecordCreateAction 56 case models.CommitOperationUpdate: 57 action = tapc.RecordUpdateAction 58 case models.CommitOperationDelete: 59 action = tapc.RecordDeleteAction 60 default: 61 return tapc.Event{}, false 62 } 63 return tapc.Event{ 64 Type: tapc.EvtRecord, 65 Record: &tapc.RecordEventData{ 66 Did: did, 67 Rkey: syntax.RecordKey(e.Commit.RKey), 68 Collection: syntax.NSID(e.Commit.Collection), 69 Action: action, 70 Record: e.Commit.Record, 71 // jetstream is only used for live 72 Live: true, 73 }, 74 }, true 75} 76 77func (s *Spindle) ingestMember(ctx context.Context, e *models.Event) error { 78 did := e.Did 79 rkey := e.Commit.RKey 80 l := s.l.With("component", "ingester", "record", tangled.SpindleMemberNSID, "did", did, "rkey", rkey) 81 82 switch e.Commit.Operation { 83 case models.CommitOperationCreate, models.CommitOperationUpdate: 84 raw := e.Commit.Record 85 record := tangled.SpindleMember{} 86 if err := json.Unmarshal(raw, &record); err != nil { 87 return fmt.Errorf("invalid record: %w", err) 88 } 89 90 domain := s.cfg.Server.Hostname 91 recordInstance := record.Instance 92 93 if recordInstance != domain { 94 return fmt.Errorf("domain mismatch: %s != %s", record.Instance, domain) 95 } 96 97 subject, err := syntax.ParseDID(record.Subject) 98 if err != nil { 99 return fmt.Errorf("invalid subject DID %q: %w", record.Subject, err) 100 } 101 102 ok, err := s.e.IsSpindleInviteAllowed(did, rbacDomain) 103 if err != nil { 104 return fmt.Errorf("failed to enforce permissions: %w", err) 105 } 106 if !ok { 107 return fmt.Errorf("permission denied for %s", did) 108 } 109 110 sqlTx, err := s.db.BeginTx(ctx, nil) 111 if err != nil { 112 return fmt.Errorf("failed to start txn: %w", err) 113 } 114 committed := false 115 defer func() { 116 if !committed { 117 sqlTx.Rollback() 118 } 119 }() 120 121 existing, err := db.GetSpindleMember(sqlTx, did, rkey) 122 if err != nil && !errors.Is(err, sql.ErrNoRows) { 123 return fmt.Errorf("failed to look up existing member: %w", err) 124 } 125 126 var staleSubject string 127 if existing != nil && existing.Subject != subject { 128 staleSubject = existing.Subject.String() 129 if err := db.RemoveSpindleMember(sqlTx, did, rkey); err != nil { 130 return fmt.Errorf("failed to remove stale member row: %w", err) 131 } 132 } 133 134 if err := db.AddSpindleMember(sqlTx, db.SpindleMember{ 135 Did: syntax.DID(did), 136 Rkey: rkey, 137 Instance: recordInstance, 138 Subject: subject, 139 }); err != nil { 140 return fmt.Errorf("failed to add member: %w", err) 141 } 142 143 if err := db.AddDid(sqlTx, subject.String()); err != nil { 144 return fmt.Errorf("failed to add did: %w", err) 145 } 146 147 dropStaleAcl := false 148 var staleDidDropped bool 149 if staleSubject != "" { 150 remaining, err := db.CountSpindleMembersBySubject(sqlTx, staleSubject) 151 if err != nil { 152 return fmt.Errorf("failed to count stale subject rows: %w", err) 153 } 154 if remaining == 0 { 155 dropStaleAcl = true 156 stillNeeded, err := s.e.WouldHaveAnyPolicyExcludingSpindleMember(staleSubject, rbacDomain) 157 if err != nil { 158 return fmt.Errorf("failed to check residual policies for stale subject: %w", err) 159 } 160 if !stillNeeded { 161 if err := db.RemoveDid(sqlTx, staleSubject); err != nil { 162 return fmt.Errorf("failed to remove stale did: %w", err) 163 } 164 staleDidDropped = true 165 } 166 } 167 l.Info("replaced stale spindle member", "old_subject", staleSubject, "new_subject", subject, "stale_did_dropped", staleDidDropped) 168 } 169 170 if err := sqlTx.Commit(); err != nil { 171 return fmt.Errorf("failed to commit txn: %w", err) 172 } 173 committed = true 174 175 if dropStaleAcl { 176 if _, err := s.e.TryRemoveSpindleMember(rbacDomain, staleSubject); err != nil { 177 l.Error("post-commit: failed to remove stale ACL", "subject", staleSubject, "err", err) 178 } 179 } 180 if _, err := s.e.TryAddSpindleMember(rbacDomain, subject.String()); err != nil { 181 l.Error("post-commit: failed to add member ACL", "subject", subject, "err", err) 182 } 183 184 if staleDidDropped { 185 s.jc.RemoveDid(staleSubject) 186 } 187 s.jc.AddDid(subject.String()) 188 l.Info("added member from firehose", "member", subject) 189 return nil 190 191 case models.CommitOperationDelete: 192 sqlTx, err := s.db.BeginTx(ctx, nil) 193 if err != nil { 194 return fmt.Errorf("failed to start txn: %w", err) 195 } 196 committed := false 197 defer func() { 198 if !committed { 199 sqlTx.Rollback() 200 } 201 }() 202 203 record, err := db.GetSpindleMember(sqlTx, did, rkey) 204 if errors.Is(err, sql.ErrNoRows) { 205 l.Info("spindle member already removed") 206 return nil 207 } 208 if err != nil { 209 return fmt.Errorf("failed to find member: %w", err) 210 } 211 212 staleSubject := record.Subject.String() 213 214 if err := db.RemoveSpindleMember(sqlTx, did, rkey); err != nil { 215 return fmt.Errorf("failed to remove member: %w", err) 216 } 217 218 remaining, err := db.CountSpindleMembersBySubject(sqlTx, staleSubject) 219 if err != nil { 220 return fmt.Errorf("failed to count remaining member rows: %w", err) 221 } 222 223 dropAcl := false 224 var staleDidDropped bool 225 if remaining == 0 { 226 dropAcl = true 227 stillNeeded, err := s.e.WouldHaveAnyPolicyExcludingSpindleMember(staleSubject, rbacDomain) 228 if err != nil { 229 return fmt.Errorf("failed to check residual policies: %w", err) 230 } 231 if !stillNeeded { 232 if err := db.RemoveDid(sqlTx, staleSubject); err != nil { 233 return fmt.Errorf("failed to remove did: %w", err) 234 } 235 staleDidDropped = true 236 } 237 } 238 239 if err := sqlTx.Commit(); err != nil { 240 return fmt.Errorf("failed to commit txn: %w", err) 241 } 242 committed = true 243 244 if dropAcl { 245 if _, err := s.e.TryRemoveSpindleMember(rbacDomain, staleSubject); err != nil { 246 l.Error("post-commit: failed to remove member ACL", "subject", staleSubject, "err", err) 247 } 248 } 249 250 if staleDidDropped { 251 s.jc.RemoveDid(staleSubject) 252 } 253 l.Info("removed member from firehose", "member", record.Subject, "remaining_rows", remaining, "stale_did_dropped", staleDidDropped) 254 } 255 return nil 256}