parallel

package
v3.8.0 Latest Latest
Warning

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

Go to latest
Published: Oct 9, 2026 License: Apache-2.0 Imports: 17 Imported by: 0

Documentation

Index

Constants

View Source
const (
	CtxBackgroundTaskIDKey = "background_task_id"

	// CtxTaskStartOrderKey holds the 0-based position of the task in the
	// order tasks actually started, which is also the order their output is
	// printed in.
	CtxTaskStartOrderKey = "task_start_order"
)

Variables

This section is empty.

Functions

func DoTasks

func DoTasks(ctx context.Context, numberOfTasks int, options DoTasksOptions, taskFunc TaskFunc) error

DoTasks executes a specified number of tasks in parallel using a configurable number of workers. Each worker runs a subset of the total tasks, and progress is logged for each task.

Parameters:

  • ctx: The context used to control the operation and provide cancellation support.
  • numberOfTasks: The total number of tasks to be executed.
  • options: A DoTasksOptions struct containing configuration parameters for task execution.
  • taskFunc: A function that performs a single task. It takes a context and a task ID as input and returns an error if one occurs.

func DoTasksDynamic

func DoTasksDynamic(ctx context.Context, options DoTasksOptions, next NextTaskFunc, taskFunc TaskFunc) error

DoTasksDynamic runs workers that each repeatedly pull the next task to run from `next` (instead of a fixed, statically-partitioned task range like DoTasks) until `next` reports there's nothing left. This allows the caller to drive a dynamic dependency-graph scheduler where the set of runnable tasks isn't known upfront and grows as earlier tasks complete.

options.MaxNumberOfWorkers <= 0 means a single worker (the task count is unknown upfront), unlike DoTasks where it means one worker per task.

`next` must be worker-agnostic: any worker must be able to take any runnable task. Workers claim their first task in ID order, so a `next` that reserves a task for a specific worker can block the worker that is holding the rest of them up.

func TaskStartOrder

func TaskStartOrder(ctx context.Context) (int, bool)

TaskStartOrder returns the task's start-order position from a context passed to a TaskFunc, and whether ctx belongs to a parallel task at all. Use it, not the task ID, for "N/Total" progress labels: tasks are printed in start order, so any other numbering comes out scattered in the log.

Types

type DoTasksOptions

type DoTasksOptions struct {
	InitDockerCLIForEachWorker bool
	MaxNumberOfWorkers         int
	OnTaskEnqueued             func(taskID, startOrder int)
}

type NextTaskFunc

type NextTaskFunc func(ctx context.Context) (taskId int, ok bool, err error)

NextTaskFunc returns the next task to run. It may block until a task becomes runnable. ok=false is terminal: the calling worker returns and is never asked again.

type Printer

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

Printer renders task output as coherent, uninterrupted per-task blocks in the order the tasks started: it streams the oldest unfinished task live and fully drains it before moving to the next, rather than interleaving concurrent tasks' output line by line. Every task is enqueued the moment it starts, so its position doubles as the task's start-order index.

Because the head of the queue is always the oldest running task, the live log never goes quiet while something is still building — as soon as the head finishes, the next task in start order takes over, already partially buffered. A task's temp file is removed right after its block is printed, so the disk holds only output that is still waiting to be printed.

func NewPrinter

func NewPrinter() *Printer

func (*Printer) Close

func (p *Printer) Close()

Close tells the Printer no more tasks will be enqueued, so Print returns once the queue is drained.

func (*Printer) Enqueue

func (p *Printer) Enqueue(out *TaskOutput) int

Enqueue appends a started task to the printing queue and returns its start-order index.

func (*Printer) EnqueueWithCallback

func (p *Printer) EnqueueWithCallback(out *TaskOutput, onEnqueued func(startOrder int)) int

EnqueueWithCallback appends a started task to the printing queue. The callback runs under the queue lock, in printing order.

func (*Printer) FailFast

func (p *Printer) FailFast(failed *TaskOutput)

FailFast reorders the queue after a task failed: a failed task that is not being printed yet moves to the end, so its error is the last thing the user sees; a failed task that is being printed right now truncates the queue behind it — the tasks queued after it were canceled and their partial output is noise. A nil output means no task failed, and the queue is left as is.

func (*Printer) Print

func (p *Printer) Print(ctx context.Context) error

Print streams the queue in order and returns once the Printer is closed and drained, or ctx is done. Calling it again resumes from where the previous call stopped.

type TaskFunc

type TaskFunc func(ctx context.Context, taskId int) error

type TaskOutput

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

TaskOutput buffers the log of a single task in a temp file so it can be written by the task and read by the Printer concurrently.

The file is created by the first write, so a task that logs nothing never touches the disk: most tasks of a cleanup run are silent, and a buffer per task allocated up front would put an inode per task in the tmp dir for as long as the printing queue holds them back.

The writer appends only and the reader reads already appended data (or nothing), so they never race on file content. The writer descriptor is closed at HalfClose and the reader descriptor is opened on first Read and closed once everything is drained: a finished task waiting to be printed holds no descriptor, which keeps the count of open files bounded by the number of workers rather than the number of tasks.

Any number of goroutines may write; only one may read, and Close may not run while a Read is in flight (the Printer is the sole reader and is stopped before the outputs are closed).

func NewTaskOutput

func NewTaskOutput(workerID, taskSeq int) *TaskOutput

func (*TaskOutput) Cleanup

func (o *TaskOutput) Cleanup() error

Cleanup removes the tmp file. The Printer calls it as soon as the output is drained, so the disk holds only what is not printed yet; the final sweep in runWorkers calls it again for whatever the Printer never reached. A call after a successful removal is a no-op, while a removal that failed is retried by the sweep: the Printer only warns about it, so giving up after the first failure would leak the file for the rest of the run.

func (*TaskOutput) Close

func (o *TaskOutput) Close() error

Close implements io.Closer: it half-closes the output and releases the reader descriptor if a drain was interrupted.

func (*TaskOutput) HalfClose

func (o *TaskOutput) HalfClose() error

HalfClose stops accepting writes and releases the writer descriptor; later writes are silently dropped. Calling it again is a no-op. It reports a buffer that could not be created, and with it the loss of the task's log.

If the last byte written is not a newline, one is appended first, under the same lock that stops further writes: the block always ends on a line boundary, however late the last write came in, so the printer can never glue the next block onto it.

func (*TaskOutput) Read

func (o *TaskOutput) Read(p []byte) (int, error)

Read implements io.Reader. It resumes reading from "total read offset" and reads until EOF, where EOF is handled with os.File.

A trailing incomplete UTF-8 sequence is held back and returned on the next call as long as more bytes are still to come (either the task is still writing, or half-closed but not yet fully drained). Without this, a fixed-size read can land its boundary in the middle of a multi-byte rune (e.g. a box-drawing character used for log prefixes), and the downstream logger converts each half independently into a replacement character, producing visible mojibake in the terminal. A read that finds nothing but such a fragment returns io.EOF, so the reader comes back for it once the rest has been written.

func (*TaskOutput) Readable

func (o *TaskOutput) Readable() bool

Readable returns true while there is (or may still come) something to read.

func (*TaskOutput) Write

func (o *TaskOutput) Write(p []byte) (int, error)

Write implements io.Writer. It appends to file and accumulates total write offset.

Failing to create the buffer costs the task's log, not the task: logboek discards whatever a log write returns, so the error is kept for HalfClose to report instead.

type Worker

type Worker struct {
	ID int
	// contains filtered or unexported fields
}

Worker owns the TaskOutputs of the tasks it ran and relays the output of the worker-level docker cli to the logger of the task running right now. The cli is created once per worker (a client per task would leak a connection pool per task) and writes only synchronously from inside the task that invoked it, so relaying to the current task's logger streams is exact and its output gets the same indentation and block boundaries as the task's own log lines. Between tasks relayed writes are dropped; a write from a goroutine that outlived its task lands in the block of whatever task the worker runs at that moment.

func NewWorker

func NewWorker(id int) *Worker

func (*Worker) Cleanup

func (w *Worker) Cleanup() error

Cleanup removes the tmp files of every task the worker ran.

func (*Worker) Close

func (w *Worker) Close() error

Close closes the tmp files of every task the worker ran.

func (*Worker) ErrStream

func (w *Worker) ErrStream() io.Writer

func (*Worker) OutStream

func (w *Worker) OutStream() io.Writer

OutStream and ErrStream are what the worker's docker cli is bound to.

type WorkerError

type WorkerError struct {
	ID int
	// contains filtered or unexported fields
}

func NewWorkerError

func NewWorkerError(id int, err error) *WorkerError

func (WorkerError) Error

func (e WorkerError) Error() string

func (WorkerError) Unwrap

func (e WorkerError) Unwrap() error

Jump to

Keyboard shortcuts

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