console

package
v0.14.0 Latest Latest
Warning

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

Go to latest
Published: Aug 23, 2026 License: MIT Imports: 13 Imported by: 0

Documentation

Overview

Package console is the queue's commands.

Each command is a struct built with the collaborators it needs and turned into a console.Command by its Command method, which is how every command in the collection is written: the registry and the compiler read the same slice, so a command missing from it does not exist and one in it with a broken Run does not build.

reg.Add(
	console.NewWorkCommand(worker, manager).Command(),
	console.NewRestartCommand(manager).Command(),
)

Every failed job command takes a tenant

A failed job carries a customer's payload, so `aru queue:failed` is a read like any other and it is scoped: --tenant is required and has no default. A listing that defaulted would print whichever customer happened to sort first.

There is no table generator

A migration for the jobs table, the failed jobs table or the batches table is not written into the application, because the schema belongs to whoever owns the table and travels with them: queue.DatabaseQueue.Migrations is the jobs table and failed.DatabaseFailedJobProvider.Migrations is the failed jobs table, both collected by the module. A generator that copied them into the project would be a second copy that drifts.

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type ClearCommand

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

ClearCommand deletes every job waiting on a queue.

It is `queue:clear`. Parked jobs are not cleared -- a job that gave up is in the dead letter list, not on a queue, and `aru queue:flush` is what empties that.

func NewClearCommand

func NewClearCommand(m *queue.QueueManager) *ClearCommand

NewClearCommand returns the command.

func (*ClearCommand) Command

func (c *ClearCommand) Command() console.Command

Command is the registry entry for queue:clear.

func (*ClearCommand) Handle

func (c *ClearCommand) Handle(ctx context.Context, o *console.IO) error

Handle runs the command.

It asks first, unless --force, and it asks in every environment: a queue with jobs on it is somebody's work whichever one it is.

type FlushFailedCommand

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

FlushFailedCommand deletes the failed jobs.

It is `queue:flush`, with --hours to keep the recent ones.

func NewFlushFailedCommand

func NewFlushFailedCommand(p failed.FailedJobProvider) *FlushFailedCommand

NewFlushFailedCommand returns the command.

func (*FlushFailedCommand) Command

func (c *FlushFailedCommand) Command() console.Command

Command is the registry entry for queue:flush.

func (*FlushFailedCommand) Handle

func (c *FlushFailedCommand) Handle(ctx context.Context, o *console.IO) error

Handle runs the command.

type ForgetFailedCommand

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

ForgetFailedCommand deletes one failed job. It is `queue:forget`.

func NewForgetFailedCommand

func NewForgetFailedCommand(p failed.FailedJobProvider) *ForgetFailedCommand

NewForgetFailedCommand returns the command.

func (*ForgetFailedCommand) Command

func (c *ForgetFailedCommand) Command() console.Command

Command is the registry entry for queue:forget.

func (*ForgetFailedCommand) Handle

func (c *ForgetFailedCommand) Handle(ctx context.Context, o *console.IO) error

Handle runs the command.

type ListFailedCommand

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

ListFailedCommand prints the jobs that gave up. It is `queue:failed`.

func NewListFailedCommand

func NewListFailedCommand(p failed.FailedJobProvider) *ListFailedCommand

NewListFailedCommand returns the command.

func (*ListFailedCommand) Command

func (c *ListFailedCommand) Command() console.Command

Command is the registry entry for queue:failed.

func (*ListFailedCommand) Handle

func (c *ListFailedCommand) Handle(ctx context.Context, o *console.IO) error

Handle runs the command.

type ListenCommand

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

ListenCommand runs a worker in a child process and restarts it when it exits.

It is `queue:listen`, and it is for development, where the point is that a rebuilt binary is picked up without anybody restarting anything -- see queue.Listener for why it is not the way to run a queue in production.

func NewListenCommand

func NewListenCommand(l *queue.Listener, options queue.ListenerOptions) *ListenCommand

NewListenCommand returns the command.

func (*ListenCommand) Command

func (c *ListenCommand) Command() console.Command

Command is the registry entry for queue:listen.

func (*ListenCommand) Handle

func (c *ListenCommand) Handle(ctx context.Context, o *console.IO) error

Handle runs the command.

type MonitorCommand

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

MonitorCommand reports how much work is waiting, and says so loudly when it is too much.

It is `queue:monitor`, meant to run on a schedule: the table is for a person, and the QueueBusy event is for whatever pages one.

func NewMonitorCommand

func NewMonitorCommand(m *queue.QueueManager, d queue.Dispatcher) *MonitorCommand

NewMonitorCommand returns the command.

The dispatcher may be nil, and then the command prints the table and announces nothing.

func (*MonitorCommand) Command

func (c *MonitorCommand) Command() console.Command

Command is the registry entry for queue:monitor.

func (*MonitorCommand) Handle

func (c *MonitorCommand) Handle(ctx context.Context, o *console.IO) error

Handle runs the command.

type PauseCommand

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

PauseCommand stops workers taking new jobs off a queue.

It is `queue:pause`, the switch to reach for when a downstream system is failing and retrying into it is making things worse.

func NewPauseCommand

func NewPauseCommand(m *queue.QueueManager) *PauseCommand

NewPauseCommand returns the command.

func (*PauseCommand) Command

func (c *PauseCommand) Command() console.Command

Command is the registry entry for queue:pause.

func (*PauseCommand) Handle

func (c *PauseCommand) Handle(ctx context.Context, o *console.IO) error

Handle runs the command.

type PruneBatchesCommand

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

PruneBatchesCommand deletes the batch rows that are old enough not to matter.

It is `queue:prune-batches`, with three retentions -- finished, unfinished and cancelled -- because they are three different decisions. A finished batch is history; an unfinished one that is a week old is a bug somebody should have seen.

func NewPruneBatchesCommand

func NewPruneBatchesCommand(b bus.PrunableBatchRepository) *PruneBatchesCommand

NewPruneBatchesCommand returns the command.

It takes the prunable repository rather than the plain one, so a store that cannot prune cannot be wired here at all, instead of being wired and then refusing at three in the morning.

func (*PruneBatchesCommand) Command

func (c *PruneBatchesCommand) Command() console.Command

Command is the registry entry for queue:prune-batches.

func (*PruneBatchesCommand) Handle

func (c *PruneBatchesCommand) Handle(ctx context.Context, o *console.IO) error

Handle runs the command.

type PruneFailedJobsCommand

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

PruneFailedJobsCommand deletes the failed jobs that are old enough not to matter.

It is `queue:prune-failed`, meant to run on a schedule, which is the difference from `queue:flush`: that one is a person deciding, this one is the retention policy.

func NewPruneFailedJobsCommand

func NewPruneFailedJobsCommand(p failed.PrunableFailedJobProvider) *PruneFailedJobsCommand

NewPruneFailedJobsCommand returns the command.

func (*PruneFailedJobsCommand) Command

func (c *PruneFailedJobsCommand) Command() console.Command

Command is the registry entry for queue:prune-failed.

func (*PruneFailedJobsCommand) Handle

Handle runs the command.

type RestartCommand

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

RestartCommand asks every running worker to stop after its current job.

It is `queue:restart`, and it is how a deploy replaces the workers -- the new binary starts, the old processes notice and exit cleanly, and whatever supervises them starts the new image.

func NewRestartCommand

func NewRestartCommand(m *queue.QueueManager) *RestartCommand

NewRestartCommand returns the command.

func (*RestartCommand) Command

func (c *RestartCommand) Command() console.Command

Command is the registry entry for queue:restart.

func (*RestartCommand) Handle

func (c *RestartCommand) Handle(ctx context.Context, o *console.IO) error

Handle runs the command.

type ResumeCommand

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

ResumeCommand lets workers take jobs off a paused queue again. It is `queue:resume`.

func NewResumeCommand

func NewResumeCommand(m *queue.QueueManager) *ResumeCommand

NewResumeCommand returns the command.

func (*ResumeCommand) Command

func (c *ResumeCommand) Command() console.Command

Command is the registry entry for queue:resume.

func (*ResumeCommand) Handle

func (c *ResumeCommand) Handle(ctx context.Context, o *console.IO) error

Handle runs the command.

type RetryBatchCommand

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

RetryBatchCommand puts the failed jobs of a batch back in line.

It is `queue:retry-batch`: `queue:retry` narrowed to one batch, which is the shape the question usually arrives in. The nightly import half-failed, and what has to be retried is the half that belongs to it.

func NewRetryBatchCommand

NewRetryBatchCommand returns the command.

func (*RetryBatchCommand) Command

func (c *RetryBatchCommand) Command() console.Command

Command is the registry entry for queue:retry-batch.

func (*RetryBatchCommand) Handle

func (c *RetryBatchCommand) Handle(ctx context.Context, o *console.IO) error

Handle runs the command.

func (*RetryBatchCommand) IsolatableID

func (c *RetryBatchCommand) IsolatableID(batchID string) string

IsolatableID is the lock this command takes, which is per batch.

It is what makes `--isolated` mean "one at a time per batch" rather than "one at a time". Two people retrying two different batches must not block each other; two people retrying the same one must, or the jobs are pushed twice.

type RetryCommand

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

RetryCommand puts a failed job back in line.

It is `queue:retry`, with the ids as arguments and --queue to take every failure off one queue. Without it the only way out of a dead letter list is SQL by hand, which is how it becomes a table nobody touches.

func NewRetryCommand

NewRetryCommand returns the command.

func (*RetryCommand) Command

func (c *RetryCommand) Command() console.Command

Command is the registry entry for queue:retry.

func (*RetryCommand) Handle

func (c *RetryCommand) Handle(ctx context.Context, o *console.IO) error

Handle runs the command.

The order matters: the job is pushed back before it is forgotten, so a push that fails leaves the failure where it was. Forgetting first and then failing to push loses the job.

type WorkCommand

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

WorkCommand drains a queue.

It is `queue:work`, with the connection as an argument and the worker's options as flags: the process a deployment runs beside the web one, from the same image.

The worker it runs is the one the application built, with its handlers already registered -- nothing resolves a job name to code on its own, so the registry is the thing that has to be passed in.

func NewWorkCommand

func NewWorkCommand(w *queue.Worker, m *queue.QueueManager) *WorkCommand

NewWorkCommand returns the command.

The manager may be nil, and then the worker drains the queue it was built with and the connection argument is refused rather than ignored.

func (*WorkCommand) Command

func (c *WorkCommand) Command() console.Command

Command is the registry entry for queue:work.

func (*WorkCommand) Handle

func (c *WorkCommand) Handle(ctx context.Context, o *console.IO) error

Handle runs the command.

The exit status is the worker's, and it is what a supervisor reads: see queue.WorkerStopReason.

Directories

Path Synopsis
Package concerns is what more than one queue command needs.
Package concerns is what more than one queue command needs.

Jump to

Keyboard shortcuts

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