Documentation
¶
Index ¶
- Constants
- Variables
- func CmdDelayedScheduleNextJob(triggerBy string) loop.Cmd
- func CmdScheduleNextJob(triggerBy string) loop.Cmd
- func NewRequest(ctx context.Context, req *reqctx.RequestDetails, stageIndex int, ...) *pbssinternal.ProcessRangeRequest
- type DelayedMsgScheduleNextJob
- type LaunchQueue
- type MsgJobFailed
- type MsgJobSucceeded
- type MsgPendingShutdown
- type MsgScheduleNextJob
- type RemoteWorker
- type Result
- type RetryableErr
- type SessionWorkerPool
- type SessionWorkerPoolFactory
- type TestWorkerPool
- type Worker
- type WorkerPool
- type WorkerPoolFactory
Constants ¶
const Tier2WorkerServiceName = "t2w"
Variables ¶
var ErrConnectionRefused = errors.New("connection refused")
var ErrorResourceExhausted = errors.New("resource exhausted")
var ErrorResourceExhaustedRampUp = errors.New("resource exhausted during ramp up")
Functions ¶
func CmdDelayedScheduleNextJob ¶ added in v1.13.0
func CmdScheduleNextJob ¶ added in v1.1.9
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 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.
type MsgJobFailed ¶ added in v1.1.9
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 MsgScheduleNextJob ¶ added in v1.1.9
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
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()
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
func (f *SessionWorkerPoolFactory) WorkerPool(ctx context.Context) WorkerPool
type TestWorkerPool ¶ added in v1.13.0
type TestWorkerPool struct {
// contains filtered or unexported fields
}
func NewTestWorkerPool ¶ added in v1.13.0
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()
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