Documentation
¶
Index ¶
- Constants
- Variables
- func MarshalPayload(payload any) ([]byte, error)
- func RegisterQueueSpecs(consumer TaskQueueConsumer, specs ...QueueSpec) error
- func ReportTaskEvent(ctx context.Context, reporter TaskReporter, event TaskEvent) error
- func ReportTaskProgress(ctx context.Context, reporter TaskReporter, task *QueueTask, progress int, ...) error
- func UnmarshalPayload(task *QueueTask, payload any) error
- type CronTask
- type CronTaskDispatcher
- type HandlerSpec
- type QueueSpec
- type QueueStats
- type QueueTask
- type QueueTaskStatus
- type Task
- type TaskDispatcher
- type TaskEvent
- type TaskFunc
- type TaskQuery
- type TaskQueueConfig
- type TaskQueueConsumer
- type TaskQueueHandler
- type TaskQueueHandlerFunc
- type TaskQueueInspector
- type TaskQueuePublisher
- type TaskReporter
- type TaskStatus
Constants ¶
const DefaultInitialPoolSize = 1024
const DefaultMaxPoolSize = math.MaxInt32
const Lens scene.InfraName = "asynctask"
Variables ¶
var ( ErrInternal = _eg.CreateError(1, "internal error") ErrTaskNotFound = _eg.CreateError(2, "task not found") ErrInvalidQueueTask = _eg.CreateError(3, "invalid queue task") ErrTaskHandlerNotFound = _eg.CreateError(4, "task handler not found") ErrQueueConfigConflict = _eg.CreateError(5, "task queue config conflict") ErrInvalidQueueName = _eg.CreateError(6, "invalid queue name") ErrQueueNotRegistered = _eg.CreateError(7, "task queue not registered") ErrInvalidTaskType = _eg.CreateError(8, "invalid task type") ErrInvalidTaskHandler = _eg.CreateError(9, "invalid task handler") ErrInvalidQueuePayload = _eg.CreateError(10, "invalid queue payload") ErrInvalidTaskQueueConsumer = _eg.CreateError(11, "invalid task queue consumer") )
Functions ¶
func MarshalPayload ¶ added in v0.3.4
func RegisterQueueSpecs ¶ added in v0.3.7
func RegisterQueueSpecs(consumer TaskQueueConsumer, specs ...QueueSpec) error
RegisterQueueSpecs registers queues first, then all handlers in each queue. It is intended to be called from a module-owned worker app, usually in Run().
func ReportTaskEvent ¶ added in v0.3.7
func ReportTaskEvent(ctx context.Context, reporter TaskReporter, event TaskEvent) error
func ReportTaskProgress ¶ added in v0.3.7
func UnmarshalPayload ¶ added in v0.3.4
Types ¶
type CronTask ¶
type CronTask struct {
Name string // Use Identifier method to access
Description string
Func TaskFunc
Total uint64
ErrCount uint64
}
func (*CronTask) Identifier ¶ added in v0.2.10
Identifier is the unique identifier getter, if this CronTask already set a name. the name will be the identifier
type CronTaskDispatcher ¶
type CronTaskDispatcher interface {
scene.Named
// Add will add task with a generated name common uuid
Add(spec string, cmd TaskFunc) (*CronTask, error)
// AddWithName will add task with specific name. name should be unique
AddWithName(spec string, name string, cmd TaskFunc) (*CronTask, error)
// AddTask is the underlying implementation for Add and AddWithName
AddTask(spec string, task *CronTask) error
// Cancel will cancel task have specific identifier
Cancel(id string) error
// GetTask will return the underlying task, not copy of the task info.
// which means user can modify TaskFunc if they need
GetTask(id string) (*CronTask, error)
}
CronTaskDispatcher is a service which will handle all cron task
type HandlerSpec ¶ added in v0.3.7
type HandlerSpec struct {
// Type is the task type routed by TaskQueueConsumer.
Type string
// Handler processes tasks of Type.
Handler TaskQueueHandler
}
HandlerSpec describes one task type handler inside a QueueSpec.
type QueueSpec ¶ added in v0.3.7
type QueueSpec struct {
// Queue is the logical queue name registered on TaskQueueConsumer.
Queue string
// Config is the queue-level runtime configuration.
Config TaskQueueConfig
// Handlers binds task types to handlers in this queue.
Handlers []HandlerSpec
}
QueueSpec describes one logical task queue and the handlers owned by a worker.
type QueueStats ¶ added in v0.3.7
type QueueStats struct {
Queue string `json:"queue"`
Total int64 `json:"total"`
ByStatus map[QueueTaskStatus]int64 `json:"by_status"`
}
type QueueTask ¶ added in v0.3.4
type QueueTask struct {
ID string `json:"id"`
// Queue is the logical queue name.
// It defines the consumption domain and queue-level runtime config.
Queue string `json:"queue"`
// Type is the task type inside a queue.
// Consumers use it to route the task to a specific handler.
Type string `json:"type"`
// Key is an optional business grouping key.
// It can be used by backends or future schedulers for partitioning,
// deduplication, or serializing tasks for the same business object.
Key string `json:"key,omitempty"`
// Payload is the serialized task body, usually JSON.
Payload []byte `json:"payload,omitempty"`
// Headers stores optional metadata for tracing or backend-specific routing.
Headers map[string]string `json:"headers,omitempty"`
// Priority orders queued tasks when the backend supports it.
// Higher values should be consumed before lower values within the same logical queue.
Priority int `json:"priority,omitempty"`
// Delay requests delayed execution.
// The exact behavior depends on the backend implementation.
Delay time.Duration `json:"delay,omitempty"`
// Timeout is the handler execution timeout for this task.
Timeout time.Duration `json:"timeout,omitempty"`
// MaxRetry overrides the queue-level retry limit when greater than zero.
MaxRetry int `json:"max_retry,omitempty"`
CreatedAt time.Time `json:"created_at"`
AvailableAt time.Time `json:"available_at"`
}
func (*QueueTask) Identifier ¶ added in v0.3.4
type QueueTaskStatus ¶ added in v0.3.4
type QueueTaskStatus string
const ( QueueTaskStatusPending QueueTaskStatus = "pending" QueueTaskStatusRunning QueueTaskStatus = "running" QueueTaskStatusRetrying QueueTaskStatus = "retrying" QueueTaskStatusSucceeded QueueTaskStatus = "succeeded" QueueTaskStatusFailed QueueTaskStatus = "failed" )
type Task ¶
func (*Task) Identifier ¶ added in v0.3.2
func (*Task) SetStatus ¶
func (t *Task) SetStatus(status TaskStatus)
func (*Task) Status ¶
func (t *Task) Status() TaskStatus
type TaskDispatcher ¶
type TaskEvent ¶ added in v0.3.7
type TaskEvent struct {
TaskID string `json:"task_id"`
Queue string `json:"queue"`
TaskType string `json:"task_type"`
Key string `json:"key,omitempty"`
Status QueueTaskStatus `json:"status"`
Attempt int `json:"attempt,omitempty"`
Progress int `json:"progress,omitempty"`
Message string `json:"message,omitempty"`
Error string `json:"error,omitempty"`
UpdatedAt time.Time `json:"updated_at"`
}
func NewTaskEvent ¶ added in v0.3.7
type TaskQuery ¶ added in v0.3.7
type TaskQuery struct {
Offset int64 `json:"offset,omitempty"`
Limit int64 `json:"limit,omitempty"`
TaskID string `json:"task_id,omitempty"`
Queue string `json:"queue,omitempty"`
TaskType string `json:"task_type,omitempty"`
Key string `json:"key,omitempty"`
Status []QueueTaskStatus `json:"status,omitempty"`
Keyword string `json:"keyword,omitempty"`
UpdatedFrom time.Time `json:"updated_from,omitempty"`
UpdatedTo time.Time `json:"updated_to,omitempty"`
}
type TaskQueueConfig ¶ added in v0.3.4
type TaskQueueConfig struct {
// Concurrency is the number of workers consuming the logical queue.
Concurrency int
// MaxRetry is the default retry limit for tasks in the queue.
MaxRetry int
// RetryDelay is the default delay before retrying a failed task.
RetryDelay time.Duration
// BufferSize is the in-memory queue buffer size.
// It is only meaningful for in-process backends such as memoryqueue.
BufferSize int
}
type TaskQueueConsumer ¶ added in v0.3.4
type TaskQueueConsumer interface {
scene.Named
// RegisterQueue defines the config of a logical queue.
// Re-registering the same queue requires an identical config.
RegisterQueue(queue string, config TaskQueueConfig) error
// RegisterHandler binds a task type to a logical queue.
// The queue must be registered before handlers are attached.
RegisterHandler(queue string, taskType string, handler TaskQueueHandler) error
}
TaskQueueConsumer registers queues and task handlers, then starts consuming from queues.
type TaskQueueHandler ¶ added in v0.3.4
TaskQueueHandler processes a single queue task delivered by a consumer.
func PayloadHandler ¶ added in v0.3.7
func PayloadHandler[T any](handler func(ctx context.Context, task *QueueTask, payload T) error) TaskQueueHandler
PayloadHandler unmarshals QueueTask.Payload into T before invoking handler.
type TaskQueueHandlerFunc ¶ added in v0.3.4
func (TaskQueueHandlerFunc) HandleTask ¶ added in v0.3.4
func (f TaskQueueHandlerFunc) HandleTask(ctx context.Context, task *QueueTask) error
type TaskQueueInspector ¶ added in v0.3.7
type TaskQueuePublisher ¶ added in v0.3.4
type TaskQueuePublisher interface {
scene.Named
// Publish enqueues a task to the backend and returns the stored task model.
// The implementation may populate fields such as ID and CreatedAt.
Publish(ctx context.Context, task *QueueTask) (*QueueTask, error)
}
TaskQueuePublisher publishes queue tasks into a backend implementation. Implementations should treat Queue and Type as required routing metadata.
type TaskReporter ¶ added in v0.3.7
type TaskStatus ¶
type TaskStatus int32
const ( TaskStatusQueue TaskStatus = iota TaskStatusRunning TaskStatusFinish )