This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / appview / ingester_repo.go
8.8 kB 332 lines
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}