This repository has no description
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}