executors

package
v0.0.5 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 2, 2026 License: MIT Imports: 7 Imported by: 0

Documentation

Overview

Package executors provides batching executors that coalesce sporadic tasks into fewer, larger executions: periodical (time-driven), bulk (count-driven), delay (debounced), chunk (byte-size-driven) and less (rate-limited).

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type BulkExecutor

type BulkExecutor struct {
	// contains filtered or unexported fields
}

BulkExecutor executes tasks when either the batch size is reached or the flush interval elapses.

func NewBulkExecutor

func NewBulkExecutor(execute Execute, opts ...BulkOption) *BulkExecutor

NewBulkExecutor returns a BulkExecutor that executes batches with execute.

func (*BulkExecutor) Add

func (be *BulkExecutor) Add(task any) error

Add adds a task to the batch.

func (*BulkExecutor) Flush

func (be *BulkExecutor) Flush()

Flush forces a flush.

func (*BulkExecutor) Wait

func (be *BulkExecutor) Wait()

Wait flushes and waits for execution to finish.

type BulkOption

type BulkOption func(*bulkOptions)

BulkOption customizes a BulkExecutor.

func WithBulkInterval

func WithBulkInterval(duration time.Duration) BulkOption

WithBulkInterval sets the flush interval.

func WithBulkTasks

func WithBulkTasks(tasks int) BulkOption

WithBulkTasks sets the batch size that triggers a flush.

type ChunkExecutor

type ChunkExecutor struct {
	// contains filtered or unexported fields
}

ChunkExecutor executes tasks when either the accumulated chunk size is reached or the flush interval elapses.

func NewChunkExecutor

func NewChunkExecutor(execute Execute, opts ...ChunkOption) *ChunkExecutor

NewChunkExecutor returns a ChunkExecutor that executes chunks with execute.

func (*ChunkExecutor) Add

func (ce *ChunkExecutor) Add(task any, size int) error

Add adds a task with the given byte size to the chunk.

func (*ChunkExecutor) Flush

func (ce *ChunkExecutor) Flush()

Flush forces a flush.

func (*ChunkExecutor) Wait

func (ce *ChunkExecutor) Wait()

Wait flushes and waits for execution to finish.

type ChunkOption

type ChunkOption func(*chunkOptions)

ChunkOption customizes a ChunkExecutor.

func WithChunkBytes

func WithChunkBytes(size int) ChunkOption

WithChunkBytes sets the chunk size that triggers a flush.

func WithFlushInterval

func WithFlushInterval(duration time.Duration) ChunkOption

WithFlushInterval sets the flush interval.

type DelayExecutor

type DelayExecutor struct {
	// contains filtered or unexported fields
}

DelayExecutor runs a task after a debounce delay, coalescing repeated triggers within the delay window.

func NewDelayExecutor

func NewDelayExecutor(fn func(), delay time.Duration) *DelayExecutor

NewDelayExecutor returns a DelayExecutor running fn after delay.

func (*DelayExecutor) Trigger

func (de *DelayExecutor) Trigger()

Trigger schedules fn to run after the delay. It is safe to call repeatedly; subsequent triggers within the delay window are coalesced.

type Execute

type Execute func(tasks []any)

Execute handles a batch of tasks.

type LessExecutor

type LessExecutor struct {
	// contains filtered or unexported fields
}

LessExecutor runs a task at most once per threshold interval.

func NewLessExecutor

func NewLessExecutor(threshold time.Duration) *LessExecutor

NewLessExecutor returns a LessExecutor with the given threshold interval.

func (*LessExecutor) DoOrDiscard

func (le *LessExecutor) DoOrDiscard(fn func()) bool

DoOrDiscard runs fn if threshold has elapsed since the last execution, and discards it otherwise. It reports whether fn ran. The last-execution time is updated with a compare-and-swap so concurrent callers within a threshold window still execute at most once.

type PeriodicalExecutor

type PeriodicalExecutor struct {
	// contains filtered or unexported fields
}

PeriodicalExecutor coalesces tasks and executes them periodically, or immediately when the container signals fullness.

func NewPeriodicalExecutor

func NewPeriodicalExecutor(interval time.Duration, container TaskContainer) *PeriodicalExecutor

NewPeriodicalExecutor returns a PeriodicalExecutor with the given interval and container. Call Flush (or Wait) before shutdown to drain pending tasks.

func (*PeriodicalExecutor) Add

func (pe *PeriodicalExecutor) Add(task any)

Add adds a task to be executed.

func (*PeriodicalExecutor) Flush

func (pe *PeriodicalExecutor) Flush() bool

Flush forces pe to execute all pending tasks.

func (*PeriodicalExecutor) Sync

func (pe *PeriodicalExecutor) Sync(fn func())

Sync runs fn while holding the executor's container lock.

func (*PeriodicalExecutor) Wait

func (pe *PeriodicalExecutor) Wait()

Wait flushes and waits for all executions to finish.

type TaskContainer

type TaskContainer interface {
	// AddTask adds task, returning true if the container needs flushing now.
	AddTask(task any) bool
	// Execute handles the collected tasks.
	Execute(tasks any)
	// RemoveAll removes and returns all collected tasks.
	RemoveAll() any
}

TaskContainer accumulates tasks for a PeriodicalExecutor.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL