Versions in this module Expand all Collapse all v0 v0.1.1 Sep 22, 2026 v0.1.0 Sep 21, 2026 Changes in this version + var ErrAlreadyCompleted = errors.New("batchx: job instance already completed") + var ErrAlreadyRunning = errors.New("batchx: job instance already running") + var ErrCheckpointBeyondEnd = errors.New("batchx: checkpoint beyond end of input") + var ErrFilter = errors.New("batchx: item filtered") + var ErrInvalidConfig = errors.New("batchx: invalid configuration") + var ErrNotFound = errors.New("batchx: job record not found") + var ErrSkipLimitExceeded = errors.New("batchx: skip limit exceeded") + func ConstantBackoff(d time.Duration) func(int) time.Duration + func ExponentialBackoff(base, max time.Duration) func(int) time.Duration + func IsSkippable(err error) bool + func Skippable(err error) error + type CSVReader struct + func NewCSVReader(r io.Reader, hasHeader bool, configure ...func(*csv.Reader)) *CSVReader + func (c *CSVReader) Header() []string + func (c *CSVReader) Read(ctx context.Context) ([]string, error) + type CSVWriter struct + func NewCSVWriter(w io.Writer, configure ...func(*csv.Writer)) *CSVWriter + func (c *CSVWriter) Write(_ context.Context, recs [][]string) error + type ChanReader struct + func NewChanReader[T any](ch <-chan T) *ChanReader[T] + func (c *ChanReader[T]) Read(ctx context.Context) (T, error) + type FileRepository struct + func NewFileRepository(dir string) (*FileRepository, error) + func (f *FileRepository) Load(_ context.Context, key string) (JobRecord, error) + func (f *FileRepository) Save(_ context.Context, rec JobRecord) error + type JSONLReader struct + func NewJSONLReader[T any](r io.Reader) *JSONLReader[T] + func (j *JSONLReader[T]) Read(ctx context.Context) (T, error) + type JSONLWriter struct + func NewJSONLWriter[T any](w io.Writer) *JSONLWriter[T] + func (j *JSONLWriter[T]) Write(_ context.Context, items []T) error + type Job struct + func NewJob(name string, repo Repository, opts ...JobOption) *Job + func (j *Job) Name() string + func (j *Job) Run(ctx context.Context, params Params) (JobRecord, error) + func (j *Job) Then(s Step) *Job + func (j *Job) ThenParallel(steps ...Step) *Job + type JobOption func(*Job) + func WithClock(now func() time.Time) JobOption + func WithListener(l Listener) JobOption + type JobRecord struct + Attempts int + EndedAt time.Time + Err string + Key string + Name string + Params Params + StartedAt time.Time + Status Status + Steps []StepRecord + UpdatedAt time.Time + Version int + func (r JobRecord) Clone() JobRecord + func (r JobRecord) Step(name string) (StepRecord, bool) + type LinesReader struct + func NewLinesReader(r io.Reader) *LinesReader + func (l *LinesReader) Read(ctx context.Context) (string, error) + type Listener struct + AfterChunk func(ctx context.Context, job string, rec StepRecord) + AfterJob func(ctx context.Context, rec JobRecord) + AfterStep func(ctx context.Context, job string, rec StepRecord) + BeforeJob func(ctx context.Context, rec JobRecord) + BeforeStep func(ctx context.Context, job, step string) + OnRetry func(ctx context.Context, step string, phase Phase, attempt int, err error) + OnSkip func(ctx context.Context, step string, phase Phase, item any, err error) + type MemoryRepository struct + func (m *MemoryRepository) Load(_ context.Context, key string) (JobRecord, error) + func (m *MemoryRepository) Save(_ context.Context, rec JobRecord) error + type Params map[string]string + type Phase string + const PhaseProcess + const PhaseRead + const PhaseWrite + type Processor interface + Process func(ctx context.Context, in I) (O, error) + type ProcessorFunc func(ctx context.Context, in I) (O, error) + func (f ProcessorFunc[I, O]) Process(ctx context.Context, in I) (O, error) + type Reader interface + Read func(ctx context.Context) (I, error) + type ReaderFunc func(ctx context.Context) (I, error) + func (f ReaderFunc[I]) Read(ctx context.Context) (I, error) + type Repository interface + Load func(ctx context.Context, key string) (JobRecord, error) + Save func(ctx context.Context, rec JobRecord) error + type Retry struct + Backoff func(attempt int) time.Duration + If func(error) bool + MaxAttempts int + type Seeker interface + Seek func(ctx context.Context, n int64) error + type SeqReader struct + func NewSeqReader[T any](seq iter.Seq[T]) *SeqReader[T] + func (s *SeqReader[T]) Close() error + func (s *SeqReader[T]) Read(context.Context) (T, error) + type SkippableError struct + Err error + func (e *SkippableError) Error() string + func (e *SkippableError) Unwrap() error + type SliceReader struct + func NewSliceReader[T any](items []T) *SliceReader[T] + func (s *SliceReader[T]) Read(context.Context) (T, error) + func (s *SliceReader[T]) Seek(_ context.Context, n int64) error + type SliceWriter struct + func (s *SliceWriter[T]) Items() []T + func (s *SliceWriter[T]) Write(_ context.Context, items []T) error + type Status string + const StatusCompleted + const StatusFailed + const StatusRunning + const StatusStopped + type Step interface + Execute func(ctx context.Context, sc *StepContext) error + Name func() string + func NewCopyStep[T any](name string, r Reader[T], w Writer[T], opts ...StepOption) (Step, error) + func NewStep[I, O any](name string, r Reader[I], p Processor[I, O], w Writer[O], opts ...StepOption) (Step, error) + func NewTaskletStep(name string, fn func(context.Context) error) Step + type StepContext struct + Job string + Record StepRecord + func (sc *StepContext) Commit(ctx context.Context) error + type StepError struct + Err error + Job string + Step string + func (e *StepError) Error() string + func (e *StepError) Unwrap() error + type StepOption func(*stepConfig) + func WithChunkSize(n int) StepOption + func WithRetry(r Retry) StepOption + func WithSkipIf(f func(error) bool) StepOption + func WithSkipLimit(n int) StepOption + func WithSleep(f func(context.Context, time.Duration) error) StepOption + func WithStepListener(l Listener) StepOption + type StepRecord struct + Attempts int + Commits int64 + EndedAt time.Time + Err string + Filtered int64 + Name string + Position int64 + Read int64 + Retries int64 + Skipped int64 + StartedAt time.Time + Status Status + Written int64 + type Writer interface + Write func(ctx context.Context, items []O) error + type WriterFunc func(ctx context.Context, items []O) error + func (f WriterFunc[O]) Write(ctx context.Context, items []O) error