work

package
v1.25.0 Latest Latest
Warning

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

Go to latest
Published: Oct 8, 2026 License: Apache-2.0 Imports: 31 Imported by: 1

Documentation

Index

Constants

View Source
const Tier2WorkerServiceName = "t2w"

Variables

View Source
var ErrConnectionRefused = errors.New("connection refused")
View Source
var ErrorResourceExhausted = errors.New("resource exhausted")
View Source
var ErrorResourceExhaustedRampUp = errors.New("resource exhausted during ramp up")

Functions

func CmdDelayedScheduleNextJob added in v1.13.0

func CmdDelayedScheduleNextJob(triggerBy string) loop.Cmd

func CmdScheduleNextJob added in v1.1.9

func CmdScheduleNextJob(triggerBy string) loop.Cmd

func NewRequest added in v1.1.9

func NewRequest(ctx context.Context, req *reqctx.RequestDetails, stageIndex int, startBlock uint64, streamOutput bool) *pbssinternal.ProcessRangeRequest

Types

type DelayedMsgScheduleNextJob added in v1.13.0

type DelayedMsgScheduleNextJob struct {
	loop.IsMsg
	TriggerBy string
}

type LaunchQueue added in v1.23.0

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

LaunchQueue decides which jobs of one tier1 request may send a request to a tier2, and when. A tier2 at its concurrent-request limit refuses the job without doing any work, so the job dials again and can land on an instance that has room; when every job of the request dials at the same rate, the segment the client reads first is no more likely to land than the one it reads last, and a whole request can idle behind one unlucky low segment.

The queue holds the jobs waiting to get into a tier2 — turned away for capacity, waiting out a failure, or held back behind those — ordered by what the client reads first: lowest segment, then highest stage, the job that also produces the partials of the stages under it. Only the first windowSize jobs of that queue may dial at all; the ones behind them send nothing and wait. A job leaves the queue the moment a tier2 takes it, which lets exactly one more job start dialing.

A job with nothing queued ahead of it dials immediately and never joins, so a fleet with room is paced no differently than without the queue.

func NewLaunchQueue added in v1.23.0

func NewLaunchQueue(maxWorkers int) *LaunchQueue

func (*LaunchQueue) Leave added in v1.23.0

func (q *LaunchQueue) Leave(unit stage.Unit)

Leave takes the job out of the queue, moving up every job behind it. It is called once the job is through to a tier2, and again when the job ends.

func (*LaunchQueue) Retry added in v1.23.0

func (q *LaunchQueue) Retry(unit stage.Unit, after time.Duration)

Retry holds the job back for at least the given delay before it dials again. It keeps its place in the queue, ahead of the jobs the client reads after it.

func (*LaunchQueue) WaitTurn added in v1.23.0

func (q *LaunchQueue) WaitTurn(ctx context.Context, unit stage.Unit) error

WaitTurn blocks until the job may send its request to a tier2.

type MsgJobFailed added in v1.1.9

type MsgJobFailed struct {
	loop.IsMsg
	Unit   stage.Unit
	Worker Worker
	Error  error
}

type MsgJobSucceeded added in v1.1.9

type MsgJobSucceeded struct {
	loop.IsMsg
	Unit     stage.Unit
	Worker   Worker
	Streamed bool

	// ProcessedBlocks is what the worker reported on completion: one count per block and per
	// stage it actually executed, so blocks skipped by a block index and blocks served from
	// the cache are not in it.
	ProcessedBlocks uint64
}

type MsgPendingShutdown added in v1.14.3

type MsgPendingShutdown struct {
	loop.IsMsg
}

type MsgScheduleNextJob added in v1.1.9

type MsgScheduleNextJob struct {
	loop.IsMsg
	TriggerBy string
}

type RemoteWorker

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

func NewRemoteWorker

func NewRemoteWorker(clientFactory client.InternalClientFactory, id string, logger *zap.Logger, launchQueue *LaunchQueue) *RemoteWorker

func (*RemoteWorker) ID added in v0.2.0

func (w *RemoteWorker) ID() string

func (*RemoteWorker) Work

func (w *RemoteWorker) Work(ctx context.Context, unit stage.Unit, startBlock uint64, moduleNames []string, upstream *response.Stream, streamOutput bool) loop.Cmd

type Result

type Result struct {
	Error error

	// ProcessedBlocks is what the worker reported on completion: one count per block and per
	// stage it actually executed, so blocks skipped by a block index are not in it.
	ProcessedBlocks uint64
}

type RetryableErr

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

func NewRetryableErr added in v1.1.1

func NewRetryableErr(cause error) *RetryableErr

func (*RetryableErr) Error

func (r *RetryableErr) Error() string

func (*RetryableErr) Unwrap added in v1.14.3

func (r *RetryableErr) Unwrap() error

type SessionWorkerPool added in v1.16.5

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

func NewSessionWorkerPool added in v1.16.5

func NewSessionWorkerPool(
	ctx context.Context,
	sessionKey string,
	sessionPool dsession.SessionPool,
	clientFactory client.InternalClientFactory,
) *SessionWorkerPool

func (*SessionWorkerPool) Borrow added in v1.16.5

func (p *SessionWorkerPool) Borrow(ctx context.Context) (Worker, error)

func (*SessionWorkerPool) ReleaseAll added in v1.23.0

func (p *SessionWorkerPool) ReleaseAll()

func (*SessionWorkerPool) Return added in v1.16.5

func (p *SessionWorkerPool) Return(ctx context.Context, worker Worker)

type SessionWorkerPoolFactory added in v1.16.5

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

func NewSessionWorkerPoolFactory added in v1.16.5

func NewSessionWorkerPoolFactory(sessionPool dsession.SessionPool, clientFactory client.InternalClientFactory) *SessionWorkerPoolFactory

func (*SessionWorkerPoolFactory) WorkerPool added in v1.16.5

type TestWorkerPool added in v1.13.0

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

func NewTestWorkerPool added in v1.13.0

func NewTestWorkerPool(t *testing.T, workerFactory func(ctx context.Context) Worker) *TestWorkerPool

func (*TestWorkerPool) Borrow added in v1.13.0

func (t *TestWorkerPool) Borrow(ctx context.Context) (Worker, error)

func (*TestWorkerPool) RampingUp added in v1.13.0

func (t *TestWorkerPool) RampingUp() bool

func (*TestWorkerPool) ReleaseAll added in v1.23.0

func (t *TestWorkerPool) ReleaseAll()

func (*TestWorkerPool) Return added in v1.13.0

func (t *TestWorkerPool) Return(ctx context.Context, worker Worker)

type Worker

type Worker interface {
	ID() string
	Work(ctx context.Context, unit stage.Unit, startBlock uint64, moduleNames []string, upstream *response.Stream, streamOutput bool) loop.Cmd // *Result
}

type WorkerPool

type WorkerPool interface {
	Borrow(ctx context.Context) (Worker, error)
	Return(ctx context.Context, worker Worker)
	// ReleaseAll returns every worker still borrowed. Called when the parallel processing
	// ends, before the request releases its session: the session pool cannot account for
	// workers returned once their session is gone.
	ReleaseAll()
}

type WorkerPoolFactory added in v1.13.0

type WorkerPoolFactory func(ctx context.Context) WorkerPool

Jump to

Keyboard shortcuts

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