This repository has no description
1package appview
2
3import (
4 "context"
5 "database/sql"
6 "encoding/json"
7 "errors"
8 "fmt"
9 "slices"
10
11 "github.com/bluesky-social/indigo/atproto/syntax"
12 jmodels "github.com/bluesky-social/jetstream/pkg/models"
13 "tangled.org/core/api/tangled"
14 "tangled.org/core/appview/db"
15 "tangled.org/core/appview/models"
16 "tangled.org/core/orm"
17)
18
19func (i *Ingester) ingestRepo(ctx context.Context, e *jmodels.Event) error {
20 l := i.Logger.With("handler", "ingestRepo", "did", e.Did, "rkey", e.Commit.RKey)
21
22 switch e.Commit.Operation {
23 case jmodels.CommitOperationCreate:
24 return i.ingestRepoCreate(ctx, e)
25 case jmodels.CommitOperationUpdate:
26 return i.ingestRepoUpdate(ctx, e)
27 case jmodels.CommitOperationDelete:
28 return i.ingestRepoDelete(ctx, e)
29 default:
30 l.Info("unknown repo operation", "op", e.Commit.Operation)
31 return nil
32 }
33}
34
35func (i *Ingester) ingestRepoCreate(ctx context.Context, e *jmodels.Event) error {
36 l := i.Logger.With("handler", "ingestRepoCreate", "did", e.Did, "rkey", e.Commit.RKey)
37
38 record := tangled.Repo{}
39 if err := json.Unmarshal(json.RawMessage(e.Commit.Record), &record); err != nil {
40 l.Error("invalid record", "err", err)
41 return err
42 }
43
44 if record.RepoDid == nil || *record.RepoDid == "" {
45 l.Info("skipping repo create from non-DID-migrated knot")
46 return nil
47 }
48 repoDid := *record.RepoDid
49
50 _, err := db.GetRepo(i.Db,
51 orm.FilterEq("did", e.Did),
52 orm.FilterEq("rkey", e.Commit.RKey),
53 )
54 if err == nil {
55 l.Info("repo row already exists, skipping create", "did", e.Did, "rkey", e.Commit.RKey)
56 return nil
57 }
58 if !errors.Is(err, sql.ErrNoRows) {
59 return fmt.Errorf("failed to check existing repo: %w", err)
60 }
61
62 prev, err := db.GetRepoByDid(i.Db, repoDid)
63 if err != nil && !errors.Is(err, sql.ErrNoRows) {
64 return fmt.Errorf("failed to check existing repoDid: %w", err)
65 }
66
67 if prev != nil {
68 l.Info("repoDid exists under different rkey, renaming",
69 "oldRkey", prev.Rkey, "newRkey", e.Commit.RKey)
70
71 oldRepo := *prev
72
73 tx, txErr := i.Db.Begin()
74 if txErr != nil {
75 return fmt.Errorf("failed to begin rename tx: %w", txErr)
76 }
77 defer tx.Rollback()
78
79 newName := derefString(record.Name)
80 if newName == "" {
81 newName = e.Commit.RKey
82 }
83
84 if err := db.RenameRepo(tx, e.Did, prev.Rkey, e.Commit.RKey, newName); err != nil {
85 return fmt.Errorf("failed to rename repo: %w", err)
86 }
87 if err := db.RecordRepoRename(tx, e.Did, prev.Rkey, repoDid); err != nil {
88 return fmt.Errorf("failed to record rename history: %w", err)
89 }
90
91 renamed := *prev
92 renamed.Rkey = e.Commit.RKey
93 renamed.Name = newName
94 desired := repoFromRecord(&renamed, &record)
95 if repoMetadataChanged(&renamed, &desired) {
96 if err := applyRepoMetadata(tx, &renamed, desired); err != nil {
97 return fmt.Errorf("failed to apply metadata after rename: %w", err)
98 }
99 }
100
101 if err := tx.Commit(); err != nil {
102 return fmt.Errorf("failed to commit rename tx: %w", err)
103 }
104
105 newRepo, err := db.GetRepo(i.Db,
106 orm.FilterEq("did", e.Did),
107 orm.FilterEq("rkey", e.Commit.RKey),
108 )
109 if err != nil {
110 l.Warn("failed to fetch repo after rename for notification", "err", err)
111 return nil
112 }
113 i.Notifier.RenameRepo(ctx, syntax.DID(e.Did), &oldRepo, newRepo)
114 return nil
115 }
116
117 rkey := e.Commit.RKey
118 name := derefString(record.Name)
119 if name == "" {
120 name = rkey
121 }
122
123 repo := &models.Repo{
124 Did: e.Did,
125 Name: name,
126 Knot: record.Knot,
127 Rkey: rkey,
128 Description: derefString(record.Description),
129 Website: derefString(record.Website),
130 Topics: append([]string(nil), record.Topics...),
131 Source: derefString(record.Source),
132 Spindle: derefString(record.Spindle),
133 Labels: append([]string(nil), record.Labels...),
134 RepoDid: repoDid,
135 }
136
137 tx, err := i.Db.Begin()
138 if err != nil {
139 return fmt.Errorf("failed to begin insert tx: %w", err)
140 }
141 defer tx.Rollback()
142
143 if err := db.AddRepo(tx, repo); err != nil {
144 return fmt.Errorf("failed to insert repo: %w", err)
145 }
146 if err := tx.Commit(); err != nil {
147 return fmt.Errorf("failed to commit insert tx: %w", err)
148 }
149
150 i.Notifier.NewRepo(ctx, repo)
151 return nil
152}
153
154func (i *Ingester) ingestRepoUpdate(ctx context.Context, e *jmodels.Event) error {
155 l := i.Logger.With("handler", "ingestRepoUpdate", "did", e.Did, "rkey", e.Commit.RKey)
156
157 record := tangled.Repo{}
158 if err := json.Unmarshal(json.RawMessage(e.Commit.Record), &record); err != nil {
159 l.Error("invalid record", "err", err)
160 return err
161 }
162
163 if record.RepoDid == nil || *record.RepoDid == "" {
164 l.Info("skipping repo update from non-DID-migrated knot")
165 return nil
166 }
167
168 current, err := db.GetRepo(i.Db,
169 orm.FilterEq("did", e.Did),
170 orm.FilterEq("rkey", e.Commit.RKey),
171 )
172 if err != nil {
173 if errors.Is(err, sql.ErrNoRows) {
174 l.Info("skipping repo update for unknown row")
175 return nil
176 }
177 return fmt.Errorf("failed to fetch repo for ingest: %w", err)
178 }
179
180 desired := repoFromRecord(current, &record)
181
182 if current.Source != desired.Source {
183 l.Warn("source field changed but mutation is unsupported, ignoring",
184 "current", current.Source, "desired", desired.Source)
185 }
186
187 if !repoMetadataChanged(current, &desired) {
188 return nil
189 }
190
191 tx, err := i.Db.Begin()
192 if err != nil {
193 return fmt.Errorf("failed to begin tx: %w", err)
194 }
195 defer tx.Rollback()
196
197 if err := applyRepoMetadata(tx, current, desired); err != nil {
198 return fmt.Errorf("failed to apply repo metadata: %w", err)
199 }
200 return tx.Commit()
201}
202
203func (i *Ingester) ingestRepoDelete(ctx context.Context, e *jmodels.Event) error {
204 l := i.Logger.With("handler", "ingestRepoDelete", "did", e.Did, "rkey", e.Commit.RKey)
205
206 repo, err := db.GetRepo(i.Db,
207 orm.FilterEq("did", e.Did),
208 orm.FilterEq("rkey", e.Commit.RKey),
209 )
210 if err != nil {
211 if errors.Is(err, sql.ErrNoRows) {
212 l.Info("skipping repo delete for unknown row")
213 return nil
214 }
215 return fmt.Errorf("failed to fetch repo for delete: %w", err)
216 }
217
218 if err := db.RemoveRepo(i.Db, e.Did, e.Commit.RKey); err != nil {
219 return fmt.Errorf("failed to delete repo: %w", err)
220 }
221
222 i.Notifier.DeleteRepo(ctx, repo)
223 l.Info("deleted repo row")
224 return nil
225}
226
227func applyRepoMetadata(tx *sql.Tx, current *models.Repo, desired models.Repo) error {
228 if err := db.PutRepo(tx, desired); err != nil {
229 return err
230 }
231
232 if current.Spindle != desired.Spindle {
233 var spindlePtr *string
234 if desired.Spindle != "" {
235 spindlePtr = &desired.Spindle
236 }
237 if err := db.UpdateSpindle(tx, desired.RepoDid, spindlePtr); err != nil {
238 return err
239 }
240 }
241
242 if !labelsEqual(current.Labels, desired.Labels) {
243 if err := reconcileLabels(tx, current, desired); err != nil {
244 return err
245 }
246 }
247
248 return nil
249}
250
251func reconcileLabels(tx *sql.Tx, current *models.Repo, desired models.Repo) error {
252 added := filterOut(desired.Labels, current.Labels)
253 removed := filterOut(current.Labels, desired.Labels)
254
255 if err := applyEach(added, func(l string) error {
256 return db.SubscribeLabel(tx, &models.RepoLabel{
257 RepoDid: syntax.DID(desired.RepoDid),
258 LabelAt: syntax.ATURI(l),
259 })
260 }); err != nil {
261 return err
262 }
263
264 return applyEach(removed, func(l string) error {
265 return db.UnsubscribeLabel(tx,
266 orm.FilterEq("repo_did", desired.RepoDid),
267 orm.FilterEq("label_at", l),
268 )
269 })
270}
271
272func filterOut(items, exclude []string) []string {
273 return slices.DeleteFunc(slices.Clone(items), func(s string) bool {
274 return slices.Contains(exclude, s)
275 })
276}
277
278func applyEach(items []string, fn func(string) error) error {
279 for _, item := range items {
280 if err := fn(item); err != nil {
281 return err
282 }
283 }
284 return nil
285}
286
287func labelsEqual(a, b []string) bool {
288 if len(a) != len(b) {
289 return false
290 }
291 aSorted := append([]string(nil), a...)
292 bSorted := append([]string(nil), b...)
293 slices.Sort(aSorted)
294 slices.Sort(bSorted)
295 return slices.Equal(aSorted, bSorted)
296}
297
298func repoFromRecord(current *models.Repo, record *tangled.Repo) models.Repo {
299 out := *current
300 out.Name = derefString(record.Name)
301 if out.Name == "" {
302 out.Name = current.Rkey
303 }
304 out.Knot = record.Knot
305 out.Description = derefString(record.Description)
306 out.Website = derefString(record.Website)
307 out.Topics = append([]string(nil), record.Topics...)
308 out.Spindle = derefString(record.Spindle)
309 out.Source = derefString(record.Source)
310 out.Labels = append([]string(nil), record.Labels...)
311 if record.RepoDid != nil {
312 out.RepoDid = *record.RepoDid
313 }
314 return out
315}
316
317func repoMetadataChanged(current *models.Repo, desired *models.Repo) bool {
318 return current.Name != desired.Name ||
319 current.Knot != desired.Knot ||
320 current.Description != desired.Description ||
321 current.Website != desired.Website ||
322 current.TopicStr() != desired.TopicStr() ||
323 current.Spindle != desired.Spindle ||
324 !labelsEqual(current.Labels, desired.Labels)
325}
326
327func derefString(s *string) string {
328 if s == nil {
329 return ""
330 }
331 return *s
332}