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
13 kB 464 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, 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}