queue

package
v0.5.0 Latest Latest
Warning

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

Go to latest
Published: Sep 6, 2026 License: MIT Imports: 5 Imported by: 0

Documentation

Overview

Package queue e a fila persistente da secao 8 do plano.

A fila mora no Postgres, nao em canal de memoria: "Nunca depender exclusivamente de in-memory channel para jobs criticos". Um processo que morre com itens em canal perde trabalho; um que morre com itens em tabela nao.

O claim usa `FOR UPDATE SKIP LOCKED`, que e o padrao para fila em Postgres: varios dispatchers competem pela mesma tabela sem bloquear uns aos outros e sem entregar o mesmo item duas vezes. A alternativa — SELECT seguido de UPDATE — tem corrida entre as duas instrucoes.

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type Item

type Item struct {
	ID           int64
	RunID        uuid.UUID
	Prioridade   int
	DisponivelEm time.Time
}

Item e uma entrada da fila.

type Queue

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

Queue opera sobre queue_items.

func New

func New(pool *pgxpool.Pool) *Queue

func (*Queue) Claim

func (q *Queue) Claim(ctx context.Context, worker string, limite int) ([]Item, error)

Claim reivindica ate `limite` itens para este worker.

O limite e como a concorrencia e imposta: o dispatcher pede apenas as vagas que tem livres. Nao existe caminho em que mais itens saiam da fila do que a concorrencia permite, porque quem conta as vagas e quem pede.

func (*Queue) Done

func (q *Queue) Done(ctx context.Context, id int64) error

Done remove o item: o trabalho terminou e nao volta.

func (*Queue) Enqueue

func (q *Queue) Enqueue(ctx context.Context, runID uuid.UUID, prioridade int, disponivelEm time.Time) error

Enqueue coloca um run na fila.

`ON CONFLICT DO NOTHING` na unique de run_id: enfileirar duas vezes o mesmo run e no-op, nao erro. E o comportamento que a secao 29 pede — a operacao tolera repeticao. Enqueue poe o run na fila. `disponivelEm` zero significa AGORA, medido pelo relogio do BANCO.

A diferenca importa: o relogio do processo pode estar alguns milissegundos a frente do relogio do Postgres, e um item gravado com `time.Now()` do aplicativo fica invisivel ate o banco alcanca-lo. Nao e perda — o proximo ciclo pega —, mas e latencia inexplicavel, e foi o que fez um teste de concorrencia entregar 4 itens onde 5 estavam prontos.

func (*Queue) Recuperar

func (q *Queue) Recuperar(ctx context.Context, limite time.Duration) ([]Item, error)

Recuperar devolve a fila os itens reivindicados ha mais tempo que `limite`.

E a rede de seguranca contra worker morto: sem isso, um item reivindicado por um processo que caiu ficaria preso para sempre. Era exatamente o modo de falha das execucoes zumbis que travaram pipelines por 33 dias no sistema anterior.

func (*Queue) Release

func (q *Queue) Release(ctx context.Context, id int64, atraso time.Duration) error

Release devolve o item a fila, disponivel apos `atraso`.

Usado no retry e quando um dispatcher e interrompido antes de concluir: o item volta a ficar livre em vez de ficar preso a um worker que morreu.

func (*Queue) Tamanho

func (q *Queue) Tamanho(ctx context.Context) (pendentes, reivindicados int, err error)

Tamanho conta os itens pendentes e os reivindicados, para observabilidade.

Jump to

Keyboard shortcuts

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