Documentation
¶
Index ¶
- Variables
- func New(baseDir string) models.QueueStore
- type DualQueue
- func (q *DualQueue) Dequeue(ctx context.Context) (models.QueuedItemData, error)
- func (q *DualQueue) DequeueByDAGRunID(ctx context.Context, dagRunID string) ([]models.QueuedItemData, error)
- func (q *DualQueue) Enqueue(ctx context.Context, priority models.QueuePriority, dagRun digraph.DAGRunRef) error
- func (q *DualQueue) FindByDAGRunID(ctx context.Context, dagRunID string) (models.QueuedItemData, error)
- func (q *DualQueue) Len(ctx context.Context) (int, error)
- func (q *DualQueue) List(ctx context.Context) ([]models.QueuedItemData, error)
- type ItemData
- type Job
- type QueueFile
- func (q *QueueFile) FindByDAGRunID(ctx context.Context, dagRunID string) (*Job, error)
- func (q *QueueFile) Len(ctx context.Context) (int, error)
- func (q *QueueFile) List(ctx context.Context) ([]*Job, error)
- func (q *QueueFile) Pop(ctx context.Context) (*Job, error)
- func (q *QueueFile) PopByDAGRunID(ctx context.Context, dagRunID string) ([]*Job, error)
- func (q *QueueFile) Push(_ context.Context, dagRun digraph.DAGRunRef) error
- type Store
- func (s *Store) All(ctx context.Context) ([]models.QueuedItemData, error)
- func (s *Store) BaseDir() string
- func (s *Store) DequeueByDAGRunID(ctx context.Context, name, dagRunID string) ([]models.QueuedItemData, error)
- func (s *Store) DequeueByName(ctx context.Context, name string) (models.QueuedItemData, error)
- func (s *Store) Enqueue(ctx context.Context, name string, p models.QueuePriority, ...) error
- func (s *Store) Len(ctx context.Context, name string) (int, error)
- func (s *Store) List(ctx context.Context, name string) ([]models.QueuedItemData, error)
- func (s *Store) ListByDAGName(ctx context.Context, name, dagName string) ([]models.QueuedItemData, error)
- func (s *Store) Reader(_ context.Context) models.QueueReader
Constants ¶
This section is empty.
Variables ¶
var ( ErrQueueEmpty = errors.New("queue is empty") ErrQueueItemNotFound = errors.New("queue item not found") )
Errors for the queue
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 ¶
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 ¶
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.
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 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 ¶
NewQueueFile creates a new queue file with the specified base directory and priority
func (*QueueFile) FindByDAGRunID ¶
FindByDAGRunID finds a job by its dag-run ID without removing it from the queue.
func (*QueueFile) PopByDAGRunID ¶
PopByDAGRunID removes jobs from the queue by dag-run ID
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) DequeueByDAGRunID ¶
func (s *Store) DequeueByDAGRunID(ctx context.Context, name, dagRunID string) ([]models.QueuedItemData, error)
DequeueByDAGRunID implements models.QueueStore.
func (*Store) DequeueByName ¶
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.