This repository has no description
0

Configure Feed

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

spindle: ingest pull records only from jetstream

Signed-off-by: noriaki watanabe <nabeyang@gmail.com>

+116 -24
+11 -8
spindle/embedtap.go
··· 50 50 closed atomic.Bool 51 51 } 52 52 53 - func startEmbeddedTap(ctx context.Context, cfg *config.Config, logger *slog.Logger) (*embeddedTap, error) { 54 - if err := assertLoopbackBind(cfg.Server.Tap.Bind); err != nil { 55 - return nil, err 56 - } 57 - 58 - tcfg := tap.Config{ 53 + func newEmbeddedTapConfig(cfg *config.Config) tap.Config { 54 + return tap.Config{ 59 55 DatabaseURL: "sqlite://" + cfg.Server.Tap.DBPath, 60 56 DBMaxConns: 32, 61 57 PLCURL: cfg.Server.PlcUrl, ··· 67 63 RepoFetchTimeout: 5 * time.Minute, 68 64 IdentityCacheSize: 50_000, 69 65 EventCacheSize: 10_000, 70 - SignalCollection: tangled.RepoPullNSID, // HACK: to ingest PRs from any users 71 - CollectionFilters: []string{tangled.RepoNSID, tangled.RepoCollaboratorNSID, tangled.RepoPullNSID}, 66 + CollectionFilters: []string{tangled.RepoNSID, tangled.RepoCollaboratorNSID}, 72 67 AdminPassword: cfg.Server.Tap.AdminPassword, 73 68 RetryTimeout: 60 * time.Second, 74 69 } 70 + } 71 + 72 + func startEmbeddedTap(ctx context.Context, cfg *config.Config, logger *slog.Logger) (*embeddedTap, error) { 73 + if err := assertLoopbackBind(cfg.Server.Tap.Bind); err != nil { 74 + return nil, err 75 + } 76 + 77 + tcfg := newEmbeddedTapConfig(cfg) 75 78 76 79 t, err := tap.New(tcfg) 77 80 if err != nil {
+5 -1
spindle/ingester.go
··· 27 27 switch e.Commit.Collection { 28 28 case tangled.SpindleMemberNSID: 29 29 err = s.ingestMember(ctx, e) 30 - case tangled.RepoNSID, tangled.RepoCollaboratorNSID, tangled.RepoPullNSID: 30 + case tangled.RepoNSID, tangled.RepoCollaboratorNSID: 31 31 if evt, ok := jetstreamToTapEvent(e); ok { 32 32 err = s.tap.processEvent(ctx, evt) 33 + } 34 + case tangled.RepoPullNSID: 35 + if evt, ok := jetstreamToTapEvent(e); ok { 36 + err = s.processPull(ctx, evt.Record) 33 37 } 34 38 } 35 39
+87
spindle/ingester_test.go
··· 1 + package spindle 2 + 3 + import ( 4 + "context" 5 + "encoding/json" 6 + "testing" 7 + 8 + "github.com/bluesky-social/indigo/atproto/syntax" 9 + "github.com/bluesky-social/jetstream/pkg/models" 10 + 11 + "tangled.org/core/api/tangled" 12 + "tangled.org/core/spindle/config" 13 + "tangled.org/core/tapc" 14 + ) 15 + 16 + func TestTapProcessEventIgnoresPullRecords(t *testing.T) { 17 + client := &Tap{} 18 + err := client.processEvent(context.Background(), tapc.Event{ 19 + Type: tapc.EvtRecord, 20 + Record: &tapc.RecordEventData{ 21 + Live: true, 22 + Did: syntax.DID("did:plc:jge3zxi7lgrfnvhzcgrimeo7"), 23 + Collection: syntax.NSID(tangled.RepoPullNSID), 24 + Rkey: syntax.RecordKey("3mrhpypucbsg4"), 25 + Action: tapc.RecordCreateAction, 26 + Record: json.RawMessage(`{`), 27 + }, 28 + }) 29 + if err != nil { 30 + t.Fatalf("Tap.processEvent() returned an error for a pull record: %v", err) 31 + } 32 + } 33 + 34 + func TestJetstreamToTapEventMarksPullRecordsLive(t *testing.T) { 35 + tests := []struct { 36 + name string 37 + operation string 38 + action tapc.RecordAction 39 + }{ 40 + {name: "create", operation: models.CommitOperationCreate, action: tapc.RecordCreateAction}, 41 + {name: "update", operation: models.CommitOperationUpdate, action: tapc.RecordUpdateAction}, 42 + {name: "delete", operation: models.CommitOperationDelete, action: tapc.RecordDeleteAction}, 43 + } 44 + 45 + for _, tt := range tests { 46 + t.Run(tt.name, func(t *testing.T) { 47 + event, ok := jetstreamToTapEvent(&models.Event{ 48 + Did: "did:plc:jge3zxi7lgrfnvhzcgrimeo7", 49 + Kind: models.EventKindCommit, 50 + Commit: &models.Commit{ 51 + Operation: tt.operation, 52 + Collection: tangled.RepoPullNSID, 53 + RKey: "3mrhpypucbsg4", 54 + Record: json.RawMessage(`{"title":"test"}`), 55 + }, 56 + }) 57 + if !ok { 58 + t.Fatal("jetstreamToTapEvent() rejected a valid pull event") 59 + } 60 + if event.Record == nil { 61 + t.Fatal("jetstreamToTapEvent() returned no record") 62 + } 63 + if !event.Record.Live { 64 + t.Error("converted pull event is not live") 65 + } 66 + if event.Record.Collection.String() != tangled.RepoPullNSID { 67 + t.Errorf("collection = %q, want %q", event.Record.Collection, tangled.RepoPullNSID) 68 + } 69 + if event.Record.Action != tt.action { 70 + t.Errorf("action = %q, want %q", event.Record.Action, tt.action) 71 + } 72 + }) 73 + } 74 + } 75 + 76 + func TestEmbeddedTapDoesNotSubscribeToPullRecords(t *testing.T) { 77 + tcfg := newEmbeddedTapConfig(&config.Config{}) 78 + 79 + if tcfg.SignalCollection == tangled.RepoPullNSID { 80 + t.Errorf("SignalCollection = %q, must not ingest pull records", tcfg.SignalCollection) 81 + } 82 + for _, collection := range tcfg.CollectionFilters { 83 + if collection == tangled.RepoPullNSID { 84 + t.Errorf("CollectionFilters includes %q", tangled.RepoPullNSID) 85 + } 86 + } 87 + }
+13 -15
spindle/tapclient.go
··· 82 82 return t.processRepo(ctx, evt.Record) 83 83 case tangled.RepoCollaboratorNSID: 84 84 return t.processCollaborator(ctx, evt.Record) 85 - case tangled.RepoPullNSID: 86 - return t.processPull(ctx, evt.Record) 87 85 } 88 86 return nil 89 87 } ··· 315 313 return nil 316 314 } 317 315 318 - func (t *Tap) processPull(ctx context.Context, evt *tapc.RecordEventData) error { 319 - l := t.logger.With("collection", evt.Collection, "did", evt.Did, "rkey", evt.Rkey) 316 + func (s *Spindle) processPull(ctx context.Context, evt *tapc.RecordEventData) error { 317 + l := s.l.With("component", "ingester", "collection", evt.Collection, "did", evt.Did, "rkey", evt.Rkey) 320 318 321 319 // only listen to live events 322 320 if !evt.Live { ··· 345 343 } 346 344 347 345 // skip if target repo is unknown 348 - repo, err := t.spindle.db.GetRepoByDid(syntax.DID(record.Target.Repo)) 346 + repo, err := s.db.GetRepoByDid(syntax.DID(record.Target.Repo)) 349 347 if err != nil { 350 348 l.Warn("target repo is not ingested yet", "repo", record.Target.Repo, "err", err) 351 349 return fmt.Errorf("target repo is unknown") ··· 357 355 return nil 358 356 } 359 357 360 - latestSubmission, err := t.fetchLatestSubmission(ctx, evt.Did.String(), evt.Rkey.String(), &record) 358 + latestSubmission, err := s.fetchLatestSubmission(ctx, evt.Did.String(), evt.Rkey.String(), &record) 361 359 if err != nil { 362 360 return err 363 361 } 364 362 sourceSha := latestSubmission.SourceRev 365 363 366 364 scheme := "https" 367 - if t.spindle.cfg.Server.Dev { 365 + if s.cfg.Server.Dev { 368 366 scheme = "http" 369 367 } 370 368 client := &indigoxrpc.Client{Host: fmt.Sprintf("%s://%s", scheme, repo.Knot)} ··· 396 394 }, 397 395 } 398 396 399 - repoUri := t.spindle.newRepoCloneUrl(repo.Knot, repo.RepoDid) 400 - repoPath := t.spindle.newRepoPath(repo.RepoDid) 397 + repoUri := s.newRepoCloneUrl(repo.Knot, repo.RepoDid) 398 + repoPath := s.newRepoPath(repo.RepoDid) 401 399 402 400 // load workflow definitions from rev (without spindle context) 403 - rawPipeline, err := t.spindle.loadPipeline(ctx, repoUri, repoPath, sourceSha) 401 + rawPipeline, err := s.loadPipeline(ctx, repoUri, repoPath, sourceSha) 404 402 if err != nil { 405 403 // don't retry 406 404 l.Error("failed loading pipeline", "err", err) ··· 427 425 Knot: tpl.TriggerMetadata.Repo.Knot, 428 426 Rkey: tid.TID(), 429 427 } 430 - if err := t.spindle.db.CreatePipelineEvent(pipelineId.Rkey, tpl, t.spindle.n); err != nil { 428 + if err := s.db.CreatePipelineEvent(pipelineId.Rkey, tpl, s.n); err != nil { 431 429 l.Error("failed to create pipeline event", "err", err) 432 430 return nil 433 431 } 434 - sourceRepo, err := t.spindle.resolvePipelineSourceRepo(ctx, tpl.TriggerMetadata) 432 + sourceRepo, err := s.resolvePipelineSourceRepo(ctx, tpl.TriggerMetadata) 435 433 if err != nil { 436 434 l.Error("failed resolving pipeline source repo", "err", err) 437 435 return nil 438 436 } 439 - err = t.spindle.processPipeline(repo.RepoDid, tpl, pipelineId, sourceRepo) 437 + err = s.processPipeline(repo.RepoDid, tpl, pipelineId, sourceRepo) 440 438 if err != nil { 441 439 // don't retry 442 440 l.Error("failed processing pipeline", "err", err) ··· 516 514 } 517 515 } 518 516 519 - func (t *Tap) fetchLatestSubmission(ctx context.Context, did, rkey string, record *tangled.RepoPull) (*avmodels.PullSubmission, error) { 517 + func (s *Spindle) fetchLatestSubmission(ctx context.Context, did, rkey string, record *tangled.RepoPull) (*avmodels.PullSubmission, error) { 520 518 // resolve the PR owner's identity to fetch the blob from their PDS 521 - prOwnerIdent, err := t.spindle.res.ResolveIdent(ctx, did) 519 + prOwnerIdent, err := s.res.ResolveIdent(ctx, did) 522 520 if err != nil || prOwnerIdent.Handle.IsInvalidHandle() { 523 521 return nil, fmt.Errorf("failed to resolve PR owner handle: %w", err) 524 522 }