Monorepo for Tangled
tangled.org
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}