Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
95 changes: 95 additions & 0 deletions job.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,95 @@
package feedx

import "context"

// Job is a regular job.
type Job struct {
versionCheck VersionCheck
writerOpt *WriterOptions
beforeHooks []BeforeHook
}

// NewJob creates a new job.
func NewJob(check VersionCheck) *Job {
return &Job{}
}

// BeforeSync adds custom before hooks.
func (j *Job) BeforeSync(hooks ...BeforeHook) *Job {
j.beforeHooks = append(j.beforeHooks, hooks...)
return j
}

// WithWriterOptions sets custom writer options for producers.
func (j *Job) WithWriterOptions(opt *WriterOptions) *Job {
j.writerOpt = opt
return j
}

// WithVersionCheck sets a custom version check for producers.
func (j *Job) WithVersionCheck(fn VersionCheck) *Job {
j.versionCheck = fn
return j
}

// Produce starts a producer job.
func (j *Job) Produce(ctx context.Context, remoteURL string, pfn ProduceFunc) (*Status, error) {
pcr, err := NewProducer(ctx, remoteURL)
if err != nil {
return nil, err
}
defer pcr.Close()

return j.ProduceWith(ctx, pcr, pfn)
}

// Produce starts an incremental producer job.
func (j *Job) ProduceIncrementally(ctx context.Context, remoteURL string, pfn IncrementalProduceFunc) (*Status, error) {
pcr, err := NewIncrementalProducer(ctx, remoteURL)
if err != nil {
return nil, err
}
defer pcr.Close()

return j.ProduceIncrementallyWith(ctx, pcr, pfn)
}

// ProduceWith starts a producer job with an existing producer.
func (j *Job) ProduceWith(ctx context.Context, pcr *Producer, pfn ProduceFunc) (*Status, error) {
return j.produce(ctx, func(ctx context.Context, version int64) (*Status, error) {
return pcr.Produce(ctx, version, j.writerOpt, pfn)
})
}

// ProduceIncrementallyFrom starts an incremental producer job with an existing producer.
func (j *Job) ProduceIncrementallyWith(ctx context.Context, pcr *IncrementalProducer, pfn IncrementalProduceFunc) (*Status, error) {
return j.produce(ctx, func(ctx context.Context, version int64) (*Status, error) {
return pcr.Produce(ctx, version, j.writerOpt, pfn)
})
}

func (j *Job) produce(ctx context.Context, fn func(context.Context, int64) (*Status, error)) (*Status, error) {
var version int64
if j.versionCheck != nil {
latest, err := j.versionCheck(ctx)
if err != nil {
return nil, err
}
version = latest
}

if !j.runBeforeHooks(version) {
return &Status{Skipped: true, LocalVersion: version}, nil
}

return fn(ctx, version)
}

func (j *Job) runBeforeHooks(version int64) bool {
for _, hook := range j.beforeHooks {
if !hook(version) {
return false
}
}
return true
}
131 changes: 131 additions & 0 deletions job_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,131 @@
package feedx_test

import (
"context"
"errors"
"fmt"
"sync/atomic"
"testing"

"github.com/bsm/bfs"
"github.com/bsm/feedx"
)

func TestJob(t *testing.T) {
beforeCallbacks := new(atomic.Int32)
numCycles := new(atomic.Int32)

resetCounters := func() {
beforeCallbacks.Store(0)
numCycles.Store(0)
}

obj := bfs.NewInMemObject("file.json")
defer obj.Close()

t.Run("produce", func(t *testing.T) {
resetCounters()

pcr := feedx.NewProducerForRemote(obj)
defer pcr.Close()

status, err := feedx.NewJob(nil).
BeforeSync(func(_ int64) bool {
beforeCallbacks.Add(1)
return true
}).
WithVersionCheck(func(_ context.Context) (int64, error) {
return 101, nil
}).
ProduceWith(context.Background(), pcr, func(w *feedx.Writer) error {
numCycles.Add(1)
return nil
})
if err != nil {
t.Fatal("unexpected error", err)
}
if status == nil {
t.Fatal("expected status, got nil")
}
if exp, got := int32(1), numCycles.Load(); exp != got {
t.Errorf("expected %d, got %d", exp, got)
}
if exp, got := int32(1), beforeCallbacks.Load(); exp != got {
t.Errorf("expected %d, got %d", exp, got)
}
if exp, got := int64(101), status.LocalVersion; exp != got {
t.Errorf("expected %d, got %d", exp, got)
}
})

t.Run("produce skipped by before hook", func(t *testing.T) {
resetCounters()

pcr := feedx.NewProducerForRemote(obj)
defer pcr.Close()

status, err := feedx.NewJob(nil).
BeforeSync(func(_ int64) bool {
beforeCallbacks.Add(1)
return false
}).
ProduceWith(context.Background(), pcr, func(w *feedx.Writer) error {
numCycles.Add(1)
return nil
})
if err != nil {
t.Fatal("unexpected error", err)
}
if status == nil {
t.Fatal("expected status, got nil")
}
if !status.Skipped {
t.Error("expected status to be skipped")
}
if exp, got := int32(0), numCycles.Load(); exp != got {
t.Errorf("expected %d, got %d", exp, got)
}
if exp, got := int32(1), beforeCallbacks.Load(); exp != got {
t.Errorf("expected %d, got %d", exp, got)
}
})

t.Run("produce may fail", func(t *testing.T) {
resetCounters()

pcr := feedx.NewProducerForRemote(obj)
defer pcr.Close()

exp := fmt.Errorf("failed!")
_, err := feedx.NewJob(nil).
ProduceWith(context.Background(), pcr, func(w *feedx.Writer) error {
return exp
})
if !errors.Is(err, exp) {
t.Errorf("expected %v, got %v", exp, err)
}
})

t.Run("produce version check may fail", func(t *testing.T) {
resetCounters()

pcr := feedx.NewProducerForRemote(obj)
defer pcr.Close()

exp := fmt.Errorf("version check failed!")
_, err := feedx.NewJob(nil).
WithVersionCheck(func(_ context.Context) (int64, error) {
return 0, exp
}).
ProduceWith(context.Background(), pcr, func(w *feedx.Writer) error {
numCycles.Add(1)
return nil
})
if !errors.Is(err, exp) {
t.Errorf("expected %v, got %v", exp, err)
}
if exp, got := int32(0), numCycles.Load(); exp != got {
t.Errorf("expected %d, got %d", exp, got)
}
})
}