diff --git a/backfill/backfill.go b/backfill/backfill.go index 64079a83..6e0f63fe 100644 --- a/backfill/backfill.go +++ b/backfill/backfill.go @@ -35,7 +35,7 @@ type Job interface { // Once done it clears the buffer and marks the job as "complete" // Allowing the Job interface to abstract away the details of how buffered // operations are stored and/or locked - FlushBufferedOps(ctx context.Context, cb func(kind, rev, path string, rec *[]byte, cid *cid.Cid) error) error + FlushBufferedOps(ctx context.Context, cb func(kind repomgr.EventKind, rev, path string, rec *[]byte, cid *cid.Cid) error) error ClearBufferedOps(ctx context.Context) error } @@ -236,9 +236,9 @@ func (b *Backfiller) FlushBuffer(ctx context.Context, job Job) int { repo := job.Repo() // Flush buffered operations, clear the buffer, and mark the job as "complete" - // Clearning and marking are handled by the job interface - err := job.FlushBufferedOps(ctx, func(kind, rev, path string, rec *[]byte, cid *cid.Cid) error { - switch repomgr.EventKind(kind) { + // Clearing and marking are handled by the job interface + err := job.FlushBufferedOps(ctx, func(kind repomgr.EventKind, rev, path string, rec *[]byte, cid *cid.Cid) error { + switch kind { case repomgr.EvtKindCreateRecord: err := b.HandleCreateRecord(ctx, repo, rev, path, rec, cid) if err != nil { @@ -467,22 +467,23 @@ func (bf *Backfiller) HandleEvent(ctx context.Context, evt *atproto.SyncSubscrib var ops []*bufferedOp for _, op := range evt.Ops { - switch op.Action { - case "create", "update": + kind := repomgr.EventKind(op.Action) + switch kind { + case repomgr.EvtKindCreateRecord, repomgr.EvtKindUpdateRecord: cc, rec, err := bf.getRecord(ctx, r, op) if err != nil { return fmt.Errorf("getting record failed (%s,%s): %w", op.Action, op.Path, err) } ops = append(ops, &bufferedOp{ - kind: op.Action, + kind: kind, path: op.Path, rec: rec, cid: &cc, }) - case "delete": + case repomgr.EvtKindDeleteRecord: ops = append(ops, &bufferedOp{ - kind: op.Action, + kind: kind, path: op.Path, }) default: @@ -511,15 +512,15 @@ func (bf *Backfiller) HandleEvent(ctx context.Context, evt *atproto.SyncSubscrib for _, op := range ops { switch op.kind { - case "create": + case repomgr.EvtKindCreateRecord: if err := bf.HandleCreateRecord(ctx, evt.Repo, evt.Rev, op.path, op.rec, op.cid); err != nil { return fmt.Errorf("create record failed: %w", err) } - case "update": + case repomgr.EvtKindUpdateRecord: if err := bf.HandleUpdateRecord(ctx, evt.Repo, evt.Rev, op.path, op.rec, op.cid); err != nil { return fmt.Errorf("update record failed: %w", err) } - case "delete": + case repomgr.EvtKindDeleteRecord: if err := bf.HandleDeleteRecord(ctx, evt.Repo, evt.Rev, op.path); err != nil { return fmt.Errorf("delete record failed: %w", err) } @@ -533,7 +534,7 @@ func (bf *Backfiller) HandleEvent(ctx context.Context, evt *atproto.SyncSubscrib return nil } -func (bf *Backfiller) BufferOp(ctx context.Context, repo string, since *string, rev, kind, path string, rec *[]byte, cid *cid.Cid) (bool, error) { +func (bf *Backfiller) BufferOp(ctx context.Context, repo string, since *string, rev string, kind repomgr.EventKind, path string, rec *[]byte, cid *cid.Cid) (bool, error) { return bf.BufferOps(ctx, repo, since, rev, []*bufferedOp{{ path: path, kind: kind, diff --git a/backfill/backfill_test.go b/backfill/backfill_test.go index 576fbdb1..9e6c25c7 100644 --- a/backfill/backfill_test.go +++ b/backfill/backfill_test.go @@ -63,15 +63,15 @@ func TestBackfill(t *testing.T) { t.Fatal(err) } if s.State() == backfill.StateInProgress { - bf.BufferOp(ctx, testRepos[0], "delete", "app.bsky.feed.follow/1", nil, &cid.Undef) - bf.BufferOp(ctx, testRepos[0], "delete", "app.bsky.feed.follow/2", nil, &cid.Undef) - bf.BufferOp(ctx, testRepos[0], "delete", "app.bsky.feed.follow/3", nil, &cid.Undef) - bf.BufferOp(ctx, testRepos[0], "delete", "app.bsky.feed.follow/4", nil, &cid.Undef) - bf.BufferOp(ctx, testRepos[0], "delete", "app.bsky.feed.follow/5", nil, &cid.Undef) + bf.BufferOp(ctx, testRepos[0], repomgr.EvtKindDeleteRecord, "app.bsky.feed.follow/1", nil, &cid.Undef) + bf.BufferOp(ctx, testRepos[0], repomgr.EvtKindDeleteRecord, "app.bsky.feed.follow/2", nil, &cid.Undef) + bf.BufferOp(ctx, testRepos[0], repomgr.EvtKindDeleteRecord, "app.bsky.feed.follow/3", nil, &cid.Undef) + bf.BufferOp(ctx, testRepos[0], repomgr.EvtKindDeleteRecord, "app.bsky.feed.follow/4", nil, &cid.Undef) + bf.BufferOp(ctx, testRepos[0], repomgr.EvtKindDeleteRecord, "app.bsky.feed.follow/5", nil, &cid.Undef) - bf.BufferOp(ctx, testRepos[0], "create", "app.bsky.feed.follow/1", nil, &cid.Undef) + bf.BufferOp(ctx, testRepos[0], repomgr.EvtKindCreateRecord, "app.bsky.feed.follow/1", nil, &cid.Undef) - bf.BufferOp(ctx, testRepos[0], "update", "app.bsky.feed.follow/1", nil, &cid.Undef) + bf.BufferOp(ctx, testRepos[0], repomgr.EvtKindUpdateRecord, "app.bsky.feed.follow/1", nil, &cid.Undef) break } diff --git a/backfill/gormstore.go b/backfill/gormstore.go index 57cdd604..9913719a 100644 --- a/backfill/gormstore.go +++ b/backfill/gormstore.go @@ -8,6 +8,7 @@ import ( "sync" "time" + "github.com/bluesky-social/indigo/repomgr" "github.com/ipfs/go-cid" "gorm.io/gorm" ) @@ -328,7 +329,7 @@ func (j *Gormjob) SetState(ctx context.Context, state string) error { return j.db.Save(j.dbj).Error } -func (j *Gormjob) FlushBufferedOps(ctx context.Context, fn func(kind, rev, path string, rec *[]byte, cid *cid.Cid) error) error { +func (j *Gormjob) FlushBufferedOps(ctx context.Context, fn func(kind repomgr.EventKind, rev, path string, rec *[]byte, cid *cid.Cid) error) error { // TODO: this will block any events for this repo while this flush is ongoing, is that okay? j.lk.Lock() defer j.lk.Unlock() diff --git a/backfill/memstore.go b/backfill/memstore.go index 516824e9..b5593811 100644 --- a/backfill/memstore.go +++ b/backfill/memstore.go @@ -6,11 +6,12 @@ import ( "sync" "time" + "github.com/bluesky-social/indigo/repomgr" "github.com/ipfs/go-cid" ) type bufferedOp struct { - kind string + kind repomgr.EventKind path string rec *[]byte cid *cid.Cid @@ -81,7 +82,7 @@ func (s *Memstore) EnqueueJobWithState(repo, state string) error { return nil } -func (s *Memstore) BufferOp(ctx context.Context, repo string, since *string, rev, kind, path string, rec *[]byte, cid *cid.Cid) (bool, error) { +func (s *Memstore) BufferOp(ctx context.Context, repo string, since *string, rev string, kind repomgr.EventKind, path string, rec *[]byte, cid *cid.Cid) (bool, error) { s.lk.Lock() // If the job doesn't exist, we can't buffer an op for it @@ -199,7 +200,7 @@ func (j *Memjob) SetRev(ctx context.Context, rev string) error { return nil } -func (j *Memjob) FlushBufferedOps(ctx context.Context, fn func(kind, rev, path string, rec *[]byte, cid *cid.Cid) error) error { +func (j *Memjob) FlushBufferedOps(ctx context.Context, fn func(kind repomgr.EventKind, rev, path string, rec *[]byte, cid *cid.Cid) error) error { panic("TODO: copy what we end up doing from the gormstore") /* j.lk.Lock()