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
12 kB 455 lines
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}