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