filequeue

package
v1.22.10 Latest Latest
Warning

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

Go to latest
Published: Oct 6, 2025 License: GPL-3.0 Imports: 19 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

View Source
var (
	ErrQueueEmpty        = errors.New("queue is empty")
	ErrQueueItemNotFound = errors.New("queue item not found")
)

Errors for the queue

View Source
var (
	ErrQueueFileEmpty        = fmt.Errorf("queue file is empty")
	ErrQueueFileItemNotFound = fmt.Errorf("queue file item not found")
)

Errors for the queue file

Functions

func New

func New(baseDir string) models.QueueStore

Types

type DualQueue

type DualQueue struct {
	dirlock.DirLock // Embed DirLock to ensure thread-safe access to the queue files
	// contains filtered or unexported fields
}

DualQueue represents a queue for storing dag-runs with two priorities: high and low. It uses two queue files to store the items.

func NewDualQueue

func NewDualQueue(baseDir, name string) *DualQueue

NewDualQueue creates a new queue with the specified base directory and name It initializes the queue files for high and low priority

func (*DualQueue) Dequeue

func (q *DualQueue) Dequeue(ctx context.Context) (models.QueuedItemData, error)

Dequeue retrieves a dag-run from the queue and removes it. It checks the high-priority queue first, then the low-priority queue

func (*DualQueue) DequeueByDAGRunID

func (q *DualQueue) DequeueByDAGRunID(ctx context.Context, dagRunID string) ([]models.QueuedItemData, error)

DequeueByDAGRunID retrieves a dag-run from the queue by its dag-run ID

func (*DualQueue) Enqueue

func (q *DualQueue) Enqueue(ctx context.Context, priority models.QueuePriority, dagRun digraph.DAGRunRef) error

Enqueue adds a dag-run to the queue with the specified priority

func (*DualQueue) FindByDAGRunID

func (q *DualQueue) FindByDAGRunID(ctx context.Context, dagRunID string) (models.QueuedItemData, error)

FindByDAGRunID retrieves a dag-run from the queue by its dag-run ID without removing it. It returns the first found item in the queue files. If the item is not found in any of the queue files, it returns ErrQueueItemNotFound.

func (*DualQueue) Len

func (q *DualQueue) Len(ctx context.Context) (int, error)

Len returns the total number of items in the queue

func (*DualQueue) List

func (q *DualQueue) List(ctx context.Context) ([]models.QueuedItemData, error)

List returns all items in the queue

type ItemData

type ItemData struct {
	FileName string            `json:"fileName"`
	DAGRun   digraph.DAGRunRef `json:"dagRun"`
	QueuedAt time.Time         `json:"queuedAt"`
}

ItemData represents the data stored in the queue file

type Job

type Job struct {
	ItemData
	// contains filtered or unexported fields
}

func NewJob

func NewJob(data ItemData) *Job

func (*Job) Data

func (j *Job) Data() digraph.DAGRunRef

Data implements models.QueuedItem.

func (*Job) ID

func (j *Job) ID() string

ID implements models.QueuedJob.

type QueueFile

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

QueueFile is a simple queue implementation using files It stores the queued items in JSON files in a specified directory The timestamp is in UTC and the milliseconds are added to the filename Since this relies on the file system, it is not thread-safe and should be accessed by a single process at a time.

func NewQueueFile

func NewQueueFile(baseDir, prefix string) *QueueFile

NewQueueFile creates a new queue file with the specified base directory and priority

func (*QueueFile) FindByDAGRunID

func (q *QueueFile) FindByDAGRunID(ctx context.Context, dagRunID string) (*Job, error)

FindByDAGRunID finds a job by its dag-run ID without removing it from the queue.

func (*QueueFile) Len

func (q *QueueFile) Len(ctx context.Context) (int, error)

Len returns the number of items in the queue

func (*QueueFile) List

func (q *QueueFile) List(ctx context.Context) ([]*Job, error)

func (*QueueFile) Pop

func (q *QueueFile) Pop(ctx context.Context) (*Job, error)

func (*QueueFile) PopByDAGRunID

func (q *QueueFile) PopByDAGRunID(ctx context.Context, dagRunID string) ([]*Job, error)

PopByDAGRunID removes jobs from the queue by dag-run ID

func (*QueueFile) Push

func (q *QueueFile) Push(_ context.Context, dagRun digraph.DAGRunRef) error

Push adds a job to the queue Since it's a prototype, it just create a json file with the job ID and dag-run reference

type Store

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

Store implements models.QueueStore. It provides a dead-simple queue implementation using files. Since implementing a queue is not trivial, this implementation provides as a prototype for a more complex queue implementation.

func (*Store) All

func (s *Store) All(ctx context.Context) ([]models.QueuedItemData, error)

All implements models.QueueStore.

func (*Store) BaseDir

func (s *Store) BaseDir() string

BaseDir returns the base directory of the queue store

func (*Store) DequeueByDAGRunID

func (s *Store) DequeueByDAGRunID(ctx context.Context, name, dagRunID string) ([]models.QueuedItemData, error)

DequeueByDAGRunID implements models.QueueStore.

func (*Store) DequeueByName

func (s *Store) DequeueByName(ctx context.Context, name string) (models.QueuedItemData, error)

DequeueByName implements models.QueueStore.

func (*Store) Enqueue

func (s *Store) Enqueue(ctx context.Context, name string, p models.QueuePriority, dagRun digraph.DAGRunRef) error

Enqueue implements models.QueueStore.

func (*Store) Len

func (s *Store) Len(ctx context.Context, name string) (int, error)

Len implements models.QueueStore.

func (*Store) List

func (s *Store) List(ctx context.Context, name string) ([]models.QueuedItemData, error)

List implements models.QueueStore.

func (*Store) ListByDAGName added in v1.22.0

func (s *Store) ListByDAGName(ctx context.Context, name, dagName string) ([]models.QueuedItemData, error)

func (*Store) Reader

func (s *Store) Reader(_ context.Context) models.QueueReader

Reader implements models.QueueStore.

Jump to

Keyboard shortcuts

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