This repository has no description
1package appview
2
3import (
4 "context"
5 "database/sql"
6 "encoding/json"
7 "errors"
8 "fmt"
9 "log/slog"
10 "slices"
11 "strings"
12
13 "github.com/bluesky-social/indigo/atproto/syntax"
14 jmodels "github.com/bluesky-social/jetstream/pkg/models"
15 "tangled.org/core/api/tangled"
16 "tangled.org/core/appview/db"
17 "tangled.org/core/appview/models"
18 "tangled.org/core/orm"
19 "tangled.org/core/repoident"
20)
21
22func (i *Ingester) ingestRepo(ctx context.Context, e *jmodels.Event, l *slog.Logger) error {
23 l = l.With("handler", "ingestRepo")
24
25 switch e.Commit.Operation {
26 case jmodels.CommitOperationCreate:
27 return i.ingestRepoCreate(ctx, e, l)
28 case jmodels.CommitOperationUpdate:
29 return i.ingestRepoUpdate(ctx, e, l)
30 case jmodels.CommitOperationDelete:
31 return i.ingestRepoDelete(ctx, e, l)
32 default:
33 l.Info("unknown repo operation")
34 return nil
35 }
36}
37
38func (i *Ingester) ingestRepoCreate(ctx context.Context, e *jmodels.Event, l *slog.Logger) error {
39 l = l.With("handler", "ingestRepoCreate")
40
41 record := tangled.Repo{}
42 if err := json.Unmarshal(json.RawMessage(e.Commit.Record), &record); err != nil {
43 l.Error("invalid record", "err", err)
44 return err
45 }
46
47 if record.RepoDid == nil || *record.RepoDid == "" {
48 l.Info("skipping repo create from non-DID-migrated knot")
49 return nil
50 }
51 repoDid := *record.RepoDid
52
53 proceed, err := i.verifyOwnership(ctx, l, repoDid, e.Did, record.Knot)
54 if err != nil {
55 return err
56 }
57 if !proceed {
58 return nil
59 }
60
61 existing, err := db.GetRepo(i.Db,
62 orm.FilterEq("did", e.Did),
63 orm.FilterEq("rkey", e.Commit.RKey),
64 )
65 if err == nil {
66 l.Info("repo row already exists, skipping create", "did", e.Did, "rkey", e.Commit.RKey)
67 if err := i.ensureRepoOwnerPermissions(e.Did, existing.Knot, existing.RepoIdentifier()); err != nil {
68 return fmt.Errorf("failed to ensure repo owner permissions: %w", err)
69 }
70 return nil
71 }
72 if !errors.Is(err, sql.ErrNoRows) {
73 return fmt.Errorf("failed to check existing repo: %w", err)
74 }
75
76 prev, err := db.GetRepoByDid(i.Db, repoDid)
77 if err != nil && !errors.Is(err, sql.ErrNoRows) {
78 return fmt.Errorf("failed to check existing repoDid: %w", err)
79 }
80
81 if prev != nil {
82 l.Info("repoDid exists under different rkey, renaming",
83 "oldRkey", prev.Rkey, "newRkey", e.Commit.RKey)
84
85 oldRepo := *prev
86
87 tx, txErr := i.Db.Begin()
88 if txErr != nil {
89 return fmt.Errorf("failed to begin rename tx: %w", txErr)
90 }
91 defer tx.Rollback()
92
93 newName := derefString(record.Name)
94 if newName == "" {
95 newName = e.Commit.RKey
96 }
97
98 if err := db.RenameRepo(tx, e.Did, prev.Rkey, e.Commit.RKey, newName); err != nil {
99 return fmt.Errorf("failed to rename repo: %w", err)
100 }
101 if err := db.RecordRepoRename(tx, e.Did, prev.Rkey, repoDid); err != nil {
102 return fmt.Errorf("failed to record rename history: %w", err)
103 }
104 if err := db.DeleteRepoRename(tx, e.Did, strings.ToLower(newName)); err != nil {
105 return fmt.Errorf("failed to clear colliding rename alias: %w", err)
106 }
107
108 renamed := *prev
109 renamed.Rkey = e.Commit.RKey
110 renamed.Name = newName
111 desired := repoFromRecord(&renamed, &record)
112 if repoMetadataChanged(&renamed, &desired) {
113 if err := applyRepoMetadata(tx, &renamed, desired); err != nil {
114 return fmt.Errorf("failed to apply metadata after rename: %w", err)
115 }
116 }
117
118 if err := tx.Commit(); err != nil {
119 return fmt.Errorf("failed to commit rename tx: %w", err)
120 }
121
122 newRepo, err := db.GetRepo(i.Db,
123 orm.FilterEq("did", e.Did),
124 orm.FilterEq("rkey", e.Commit.RKey),
125 )
126 if err != nil {
127 l.Warn("failed to fetch repo after rename for notification", "err", err)
128 return nil
129 }
130 if err := i.ensureRepoOwnerPermissions(e.Did, newRepo.Knot, newRepo.RepoIdentifier()); err != nil {
131 return fmt.Errorf("failed to ensure repo owner permissions: %w", err)
132 }
133 i.Notifier.RenameRepo(ctx, syntax.DID(e.Did), &oldRepo, newRepo)
134 return nil
135 }
136
137 rkey := e.Commit.RKey
138 name := derefString(record.Name)
139 if name == "" {
140 name = rkey
141 }
142
143 repo := &models.Repo{
144 Did: e.Did,
145 Name: name,
146 Knot: record.Knot,
147 Rkey: rkey,
148 Description: derefString(record.Description),
149 Website: derefString(record.Website),
150 Topics: append([]string(nil), record.Topics...),
151 Source: derefString(record.Source),
152 Spindle: derefString(record.Spindle),
153 Labels: append([]string(nil), record.Labels...),
154 RepoDid: repoDid,
155 }
156
157 tx, err := i.Db.Begin()
158 if err != nil {
159 return fmt.Errorf("failed to begin insert tx: %w", err)
160 }
161 defer tx.Rollback()
162
163 if err := db.AddRepo(tx, repo); err != nil {
164 return fmt.Errorf("failed to insert repo: %w", err)
165 }
166 if err := db.DeleteRepoRename(tx, e.Did, strings.ToLower(repo.Slug())); err != nil {
167 return fmt.Errorf("failed to clear colliding rename alias: %w", err)
168 }
169 if err := tx.Commit(); err != nil {
170 return fmt.Errorf("failed to commit insert tx: %w", err)
171 }
172
173 if err := i.ensureRepoOwnerPermissions(e.Did, repo.Knot, repo.RepoIdentifier()); err != nil {
174 return fmt.Errorf("failed to ensure repo owner permissions: %w", err)
175 }
176
177 i.Notifier.NewRepo(ctx, repo)
178 return nil
179}
180
181func (i *Ingester) ensureRepoOwnerPermissions(ownerDid, knot, repo string) error {
182 if i.Enforcer == nil {
183 return fmt.Errorf("ingester has no RBAC enforcer configured")
184 }
185 if err := i.Enforcer.AddRepo(ownerDid, knot, repo); err != nil {
186 return err
187 }
188 return i.Enforcer.E.SavePolicy()
189}
190
191func (i *Ingester) ingestRepoUpdate(ctx context.Context, e *jmodels.Event, l *slog.Logger) error {
192 l = l.With("handler", "ingestRepoUpdate")
193
194 record := tangled.Repo{}
195 if err := json.Unmarshal(json.RawMessage(e.Commit.Record), &record); err != nil {
196 l.Error("invalid record", "err", err)
197 return err
198 }
199
200 if record.RepoDid == nil || *record.RepoDid == "" {
201 l.Info("skipping repo update from non-DID-migrated knot")
202 return nil
203 }
204
205 proceed, err := i.verifyOwnership(ctx, l, *record.RepoDid, e.Did, record.Knot)
206 if err != nil {
207 return err
208 }
209 if !proceed {
210 return nil
211 }
212
213 current, err := db.GetRepo(i.Db,
214 orm.FilterEq("did", e.Did),
215 orm.FilterEq("rkey", e.Commit.RKey),
216 )
217 if err != nil {
218 if errors.Is(err, sql.ErrNoRows) {
219 l.Info("skipping repo update for unknown row")
220 return nil
221 }
222 return fmt.Errorf("failed to fetch repo for ingest: %w", err)
223 }
224
225 if current.RepoDid != "" && current.RepoDid != *record.RepoDid {
226 l.Warn("rejecting repo update: repoDid is immutable",
227 "currentRepoDid", current.RepoDid,
228 "recordRepoDid", *record.RepoDid,
229 )
230 return nil
231 }
232
233 desired := repoFromRecord(current, &record)
234
235 if current.Source != desired.Source {
236 l.Warn("source field changed but mutation is unsupported, ignoring",
237 "current", current.Source, "desired", desired.Source)
238 }
239
240 if !repoMetadataChanged(current, &desired) {
241 return nil
242 }
243
244 tx, err := i.Db.Begin()
245 if err != nil {
246 return fmt.Errorf("failed to begin tx: %w", err)
247 }
248 defer tx.Rollback()
249
250 if err := applyRepoMetadata(tx, current, desired); err != nil {
251 return fmt.Errorf("failed to apply repo metadata: %w", err)
252 }
253 if err := tx.Commit(); err != nil {
254 return err
255 }
256
257 return nil
258}
259
260func (i *Ingester) ingestRepoDelete(ctx context.Context, e *jmodels.Event, l *slog.Logger) error {
261 l = l.With("handler", "ingestRepoDelete")
262
263 repo, err := db.GetRepo(i.Db,
264 orm.FilterEq("did", e.Did),
265 orm.FilterEq("rkey", e.Commit.RKey),
266 )
267 if err != nil {
268 if errors.Is(err, sql.ErrNoRows) {
269 l.Info("skipping repo delete for unknown row")
270 return nil
271 }
272 return fmt.Errorf("failed to fetch repo for delete: %w", err)
273 }
274
275 if i.Enforcer == nil {
276 return fmt.Errorf("ingester has no RBAC enforcer configured")
277 }
278
279 tx, err := i.Db.Begin()
280 if err != nil {
281 return fmt.Errorf("failed to start txn: %w", err)
282 }
283 committed := false
284 defer func() {
285 if committed {
286 return
287 }
288 tx.Rollback()
289 i.Enforcer.E.LoadPolicy()
290 }()
291
292 if err := db.RemoveRepo(tx, e.Did, e.Commit.RKey); err != nil {
293 return fmt.Errorf("failed to delete repo: %w", err)
294 }
295
296 if err := i.Enforcer.WipeRepoPolicies(repo.Knot, repo.RepoIdentifier()); err != nil {
297 return fmt.Errorf("failed to wipe repo permissions: %w", err)
298 }
299
300 if err := tx.Commit(); err != nil {
301 return fmt.Errorf("failed to commit txn: %w", err)
302 }
303
304 if err := i.Enforcer.E.SavePolicy(); err != nil {
305 return fmt.Errorf("failed to save ACLs: %w", err)
306 }
307 committed = true
308
309 i.Notifier.DeleteRepo(ctx, repo)
310 l.Info("deleted repo row")
311 return nil
312}
313
314func applyRepoMetadata(tx *sql.Tx, current *models.Repo, desired models.Repo) error {
315 if err := db.PutRepo(tx, desired); err != nil {
316 return err
317 }
318
319 if current.Spindle != desired.Spindle {
320 var spindlePtr *string
321 if desired.Spindle != "" {
322 spindlePtr = &desired.Spindle
323 }
324 if err := db.UpdateSpindle(tx, desired.RepoDid, spindlePtr); err != nil {
325 return err
326 }
327 }
328
329 if !labelsEqual(current.Labels, desired.Labels) {
330 if err := reconcileLabels(tx, current, desired); err != nil {
331 return err
332 }
333 }
334
335 return nil
336}
337
338func reconcileLabels(tx *sql.Tx, current *models.Repo, desired models.Repo) error {
339 added := filterOut(desired.Labels, current.Labels)
340 removed := filterOut(current.Labels, desired.Labels)
341
342 if err := applyEach(added, func(l string) error {
343 return db.SubscribeLabel(tx, &models.RepoLabel{
344 RepoDid: syntax.DID(desired.RepoDid),
345 LabelAt: syntax.ATURI(l),
346 })
347 }); err != nil {
348 return err
349 }
350
351 return applyEach(removed, func(l string) error {
352 return db.UnsubscribeLabel(tx,
353 orm.FilterEq("repo_did", desired.RepoDid),
354 orm.FilterEq("label_at", l),
355 )
356 })
357}
358
359func filterOut(items, exclude []string) []string {
360 return slices.DeleteFunc(slices.Clone(items), func(s string) bool {
361 return slices.Contains(exclude, s)
362 })
363}
364
365func applyEach(items []string, fn func(string) error) error {
366 for _, item := range items {
367 if err := fn(item); err != nil {
368 return err
369 }
370 }
371 return nil
372}
373
374func labelsEqual(a, b []string) bool {
375 if len(a) != len(b) {
376 return false
377 }
378 aSorted := append([]string(nil), a...)
379 bSorted := append([]string(nil), b...)
380 slices.Sort(aSorted)
381 slices.Sort(bSorted)
382 return slices.Equal(aSorted, bSorted)
383}
384
385func repoFromRecord(current *models.Repo, record *tangled.Repo) models.Repo {
386 out := *current
387 out.Name = derefString(record.Name)
388 if out.Name == "" {
389 out.Name = current.Rkey
390 }
391 out.Knot = record.Knot
392 out.Description = derefString(record.Description)
393 out.Website = derefString(record.Website)
394 out.Topics = append([]string(nil), record.Topics...)
395 out.Spindle = derefString(record.Spindle)
396 out.Source = derefString(record.Source)
397 out.Labels = append([]string(nil), record.Labels...)
398 if record.RepoDid != nil {
399 out.RepoDid = *record.RepoDid
400 }
401 return out
402}
403
404func repoMetadataChanged(current *models.Repo, desired *models.Repo) bool {
405 return current.Name != desired.Name ||
406 current.Knot != desired.Knot ||
407 current.Description != desired.Description ||
408 current.Website != desired.Website ||
409 current.TopicStr() != desired.TopicStr() ||
410 current.Spindle != desired.Spindle ||
411 !labelsEqual(current.Labels, desired.Labels)
412}
413
414func derefString(s *string) string {
415 if s == nil {
416 return ""
417 }
418 return *s
419}
420
421func (i *Ingester) verifyOwnership(ctx context.Context, l *slog.Logger, repoDid, eventDid, recordKnot string) (bool, error) {
422 if i.Verifier == nil {
423 return false, fmt.Errorf("ingester has no repo ownership verifier configured")
424 }
425 rd, err := repoident.NewRepoDid(repoDid)
426 if err != nil {
427 l.Warn("rejecting repo event: invalid repoDid on record", "repoDid", repoDid, "err", err)
428 return false, nil
429 }
430 result, err := i.Verifier(ctx, rd)
431 if err != nil {
432 return false, fmt.Errorf("verify repo ownership: %w", err)
433 }
434 if result.OwnerDid == "" {
435 l.Warn("knot lacks RepoDescribeRepo, skipping owner check; upgrade knot to 1.14+",
436 "repoDid", repoDid, "knot", result.KnotURL.String())
437 } else if result.OwnerDid.String() != eventDid {
438 l.Warn("rejecting repo event: owner mismatch",
439 "repoDid", repoDid,
440 "claimedOwner", eventDid,
441 "knotOwner", result.OwnerDid.String(),
442 "knot", result.KnotURL.String(),
443 )
444 return false, nil
445 }
446 if !strings.EqualFold(recordKnot, result.KnotURL.Host) {
447 l.Warn("rejecting repo event: record knot does not match DID-doc endpoint",
448 "repoDid", repoDid,
449 "recordKnot", recordKnot,
450 "canonicalKnot", result.KnotURL.Host,
451 )
452 return false, nil
453 }
454 return true, nil
455}