diff --git a/backfill/backfill.go b/backfill/backfill.go index 6e0f63fe..78f2673c 100644 --- a/backfill/backfill.go +++ b/backfill/backfill.go @@ -30,7 +30,10 @@ type Job interface { SetRev(ctx context.Context, rev string) error RetryCount() int - BufferOps(ctx context.Context, since *string, rev string, ops []*bufferedOp) (bool, error) + // BufferOps buffers the given operations and returns true if the operations + // were buffered. + // The given operations move the repo from since to rev. + BufferOps(ctx context.Context, since *string, rev string, ops []*BufferedOp) (bool, error) // FlushBufferedOps calls the given callback for each buffered operation // Once done it clears the buffer and marks the job as "complete" // Allowing the Job interface to abstract away the details of how buffered @@ -465,7 +468,7 @@ func (bf *Backfiller) HandleEvent(ctx context.Context, evt *atproto.SyncSubscrib return fmt.Errorf("failed to read event repo: %w", err) } - var ops []*bufferedOp + var ops []*BufferedOp for _, op := range evt.Ops { kind := repomgr.EventKind(op.Action) switch kind { @@ -474,17 +477,16 @@ func (bf *Backfiller) HandleEvent(ctx context.Context, evt *atproto.SyncSubscrib if err != nil { return fmt.Errorf("getting record failed (%s,%s): %w", op.Action, op.Path, err) } - - ops = append(ops, &bufferedOp{ - kind: kind, - path: op.Path, - rec: rec, - cid: &cc, + ops = append(ops, &BufferedOp{ + Kind: kind, + Path: op.Path, + Record: rec, + Cid: &cc, }) case repomgr.EvtKindDeleteRecord: - ops = append(ops, &bufferedOp{ - kind: kind, - path: op.Path, + ops = append(ops, &BufferedOp{ + Kind: kind, + Path: op.Path, }) default: return fmt.Errorf("invalid op action: %q", op.Action) @@ -511,17 +513,17 @@ func (bf *Backfiller) HandleEvent(ctx context.Context, evt *atproto.SyncSubscrib } for _, op := range ops { - switch op.kind { + switch op.Kind { case repomgr.EvtKindCreateRecord: - if err := bf.HandleCreateRecord(ctx, evt.Repo, evt.Rev, op.path, op.rec, op.cid); err != nil { + if err := bf.HandleCreateRecord(ctx, evt.Repo, evt.Rev, op.Path, op.Record, op.Cid); err != nil { return fmt.Errorf("create record failed: %w", err) } case repomgr.EvtKindUpdateRecord: - if err := bf.HandleUpdateRecord(ctx, evt.Repo, evt.Rev, op.path, op.rec, op.cid); err != nil { + if err := bf.HandleUpdateRecord(ctx, evt.Repo, evt.Rev, op.Path, op.Record, op.Cid); err != nil { return fmt.Errorf("update record failed: %w", err) } case repomgr.EvtKindDeleteRecord: - if err := bf.HandleDeleteRecord(ctx, evt.Repo, evt.Rev, op.path); err != nil { + if err := bf.HandleDeleteRecord(ctx, evt.Repo, evt.Rev, op.Path); err != nil { return fmt.Errorf("delete record failed: %w", err) } } @@ -535,15 +537,15 @@ func (bf *Backfiller) HandleEvent(ctx context.Context, evt *atproto.SyncSubscrib } 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, - rec: rec, - cid: cid, + return bf.BufferOps(ctx, repo, since, rev, []*BufferedOp{{ + Path: path, + Kind: kind, + Record: rec, + Cid: cid, }}) } -func (bf *Backfiller) BufferOps(ctx context.Context, repo string, since *string, rev string, ops []*bufferedOp) (bool, error) { +func (bf *Backfiller) BufferOps(ctx context.Context, repo string, since *string, rev string, ops []*BufferedOp) (bool, error) { j, err := bf.Store.GetJob(ctx, repo) if err != nil { if !errors.Is(err, ErrJobNotFound) { diff --git a/backfill/gormstore.go b/backfill/gormstore.go index 4f8cd535..10a4024a 100644 --- a/backfill/gormstore.go +++ b/backfill/gormstore.go @@ -168,7 +168,7 @@ func (s *Gormstore) createJobForRepo(repo, state string) error { return nil } -func (j *Gormjob) BufferOps(ctx context.Context, since *string, rev string, ops []*bufferedOp) (bool, error) { +func (j *Gormjob) BufferOps(ctx context.Context, since *string, rev string, ops []*BufferedOp) (bool, error) { j.lk.Lock() defer j.lk.Unlock() @@ -373,7 +373,7 @@ func (j *Gormjob) FlushBufferedOps(ctx context.Context, fn func(kind repomgr.Eve } for _, op := range opset.ops { - if err := fn(op.kind, opset.rev, op.path, op.rec, op.cid); err != nil { + if err := fn(op.Kind, opset.rev, op.Path, op.Record, op.Cid); err != nil { return err } } diff --git a/backfill/memstore.go b/backfill/memstore.go index b5593811..8e0aff78 100644 --- a/backfill/memstore.go +++ b/backfill/memstore.go @@ -10,17 +10,22 @@ import ( "github.com/ipfs/go-cid" ) -type bufferedOp struct { - kind repomgr.EventKind - path string - rec *[]byte - cid *cid.Cid +// A BufferedOp is an operation buffered while a repo is being backfilled. +type BufferedOp struct { + // Kind describes the type of operation. + Kind repomgr.EventKind + // Path contains the path the operation applies to. + Path string + // Record contains the serialized record for create and update operations. + Record *[]byte + // Cid is the CID of the record. + Cid *cid.Cid } type opSet struct { since *string rev string - ops []*bufferedOp + ops []*BufferedOp } type Memjob struct { @@ -107,18 +112,18 @@ func (s *Memstore) BufferOp(ctx context.Context, repo string, since *string, rev j.bufferedOps = append(j.bufferedOps, &opSet{ since: since, rev: rev, - ops: []*bufferedOp{&bufferedOp{ - path: path, - kind: kind, - rec: rec, - cid: cid, + ops: []*BufferedOp{&BufferedOp{ + Path: path, + Kind: kind, + Record: rec, + Cid: cid, }}, }) j.updatedAt = time.Now() return true, nil } -func (j *Memjob) BufferOps(ctx context.Context, since *string, rev string, ops []*bufferedOp) (bool, error) { +func (j *Memjob) BufferOps(ctx context.Context, since *string, rev string, ops []*BufferedOp) (bool, error) { j.lk.Lock() defer j.lk.Unlock() @@ -208,13 +213,13 @@ func (j *Memjob) FlushBufferedOps(ctx context.Context, fn func(kind repomgr.Even for _, opset := range j.bufferedOps { for _, op := range opset.ops { - if err := fn(op.kind, op.path, op.rec, op.cid); err != nil { + if err := fn(op.Kind, op.Path, op.Record, op.Cid); err != nil { return err } } } - j.bufferedOps = map[string][]*bufferedOp{} + j.bufferedOps = map[string][]*BufferedOp{} j.state = StateComplete return nil