postgres

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: 16 Imported by: 0

Documentation

Overview

Package postgres e o adaptador de persistencia.

O plano (secao 22) define o banco como fonte da verdade operacional. Aqui vive so o encanamento: pool de conexoes, migrations e o health check. As queries de dominio entram nas fases que as usam.

Index

Constants

This section is empty.

Variables

View Source
var ErrJaExiste = errors.New("run com esta chave de idempotencia ja existe")

ErrJaExiste sinaliza colisao de chave de idempotencia. Tipado para que o chamador distinga "ja criei isso" de erro real — a diferenca entre um retry benigno do scheduler e uma falha de banco.

Functions

func Migrate

func Migrate(ctx context.Context, url, direcao string) error

Migrate aplica as migrations embutidas. Roda pelo subcomando `brevis migrate`, nunca no `serve`: subir a aplicacao e migrar o schema tem blast radius diferente, e juntar as duas faz um restart casual virar um DDL.

Types

type AgendaResumo

type AgendaResumo struct {
	WorkflowSlug string
	Cron         string
	Timezone     string
	Ativo        bool
}

AgendaResumo e o minimo para calcular o proximo disparo.

type Balde

type Balde struct {
	Inicio       time.Time
	Sucesso      int
	Falha        int
	Executando   int
	Fila         int
	DuracaoMedia time.Duration
}

Balde e uma coluna do grafico de execucoes.

func (Balde) Total

func (b Balde) Total() int

Total soma o balde inteiro — a altura da coluna.

type EstadoNo

type EstadoNo struct {
	NodeID    string `json:"node_id"`
	Status    string `json:"status"`
	Tentativa int    `json:"attempt"`
	ExitCode  *int   `json:"exit_code,omitempty"`
	Erro      string `json:"erro,omitempty"`
	DuracaoMs int64  `json:"duracao_ms"`

	// Etapas sao as fases anunciadas por um passo do SDK. Vazio para um passo
	// que nao e do SDK -- e a tela desse passo continua sendo a de sempre.
	Etapas []Etapa `json:"etapas,omitempty"`

	// SdkVersao e a versao que o passo anunciou, vazia quando nao e do SDK.
	SdkVersao string `json:"sdk_versao,omitempty"`
}

EstadoNo e o estado de um passo, para a UI.

type Etapa

type Etapa struct {
	Nome    string         `json:"nome"`
	Estado  string         `json:"estado"`
	Ms      *int64         `json:"ms,omitempty"`
	Em      string         `json:"em"`
	Numeros map[string]any `json:"numeros,omitempty"`
}

Etapa e uma fase de um passo do SDK, para a tela.

type FiltroRuns

type FiltroRuns struct {
	Estado   string
	Workflow string
	De       *time.Time
	Ate      *time.Time
	Limite   int
	Offset   int
}

FiltroRuns e a consulta da tela de execucoes. Campos vazios nao filtram.

type Indicadores

type Indicadores struct {
	Total        int
	Sucesso      int
	Falha        int
	EmExecucao   int
	Pendentes    int
	DuracaoMedia time.Duration
}

Indicadores e o cabecalho do Overview.

func (Indicadores) Razao

func (i Indicadores) Razao(parte int) float64

Razao devolve o percentual de `parte` sobre o total ja concluido.

O denominador exclui o que ainda esta correndo: contar uma run em andamento como "nao-sucesso" faz a taxa despencar durante um pico de trabalho e subir sozinha depois, sem que nada tenha mudado.

type LeituraRepo

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

LeituraRepo serve a UI.

func NewLeituraRepo

func NewLeituraRepo(p *Pool) *LeituraRepo

func (*LeituraRepo) Agendas

func (r *LeituraRepo) Agendas(ctx context.Context) ([]AgendaResumo, error)

Agendas devolve todas as agendas, ativas ou nao. A lista de DAGs mostra as pausadas tambem — some-las da tela seria esconder o motivo de nada rodar.

func (*LeituraRepo) ContagemPorStatus

func (r *LeituraRepo) ContagemPorStatus(ctx context.Context) (map[string]int, error)

ContagemPorStatus alimenta os cartoes do dashboard.

func (*LeituraRepo) ContarRuns

func (r *LeituraRepo) ContarRuns(ctx context.Context, f FiltroRuns) (int, error)

ContarRuns devolve o total do MESMO filtro, para a paginacao saber quantas paginas existem.

func (*LeituraRepo) EmAndamento

func (r *LeituraRepo) EmAndamento(ctx context.Context, limite int) ([]ResumoRun, error)

EmAndamento lista o que esta correndo ou esperando vez, mais antigo primeiro.

A ordem e crescente de proposito: quem esta ha mais tempo na fila e o que merece atencao, e ordenar pelo mais recente esconderia exatamente isso.

func (*LeituraRepo) ExecucoesPorHora

func (r *LeituraRepo) ExecucoesPorHora(ctx context.Context, horas int) ([]Balde, error)

ExecucoesPorHora devolve uma coluna por hora, INCLUSIVE as vazias.

O `generate_series` a esquerda e o ponto: sem ele, uma hora sem execucao simplesmente nao apareceria e o grafico comprimiria o tempo, dando a impressao de atividade continua onde houve um buraco.

func (*LeituraRepo) Indicadores

func (r *LeituraRepo) Indicadores(ctx context.Context, janela time.Duration) (Indicadores, error)

Indicadores agrega a janela recente para os quatro cartoes do topo.

Uma consulta so, com FILTER, em vez de quatro: sao quatro varreduras da mesma tabela sobre o mesmo predicado de tempo.

func (*LeituraRepo) ProfundidadeDaFila

func (r *LeituraRepo) ProfundidadeDaFila(ctx context.Context) (pendentes, reivindicados int, err error)

ProfundidadeDaFila mostra a fila no dashboard.

func (*LeituraRepo) Projetos

func (r *LeituraRepo) Projetos(ctx context.Context) ([]ResumoProjeto, error)

Projetos lista os projetos com seus totais.

func (*LeituraRepo) Runs

func (r *LeituraRepo) Runs(ctx context.Context, f FiltroRuns) ([]ResumoRun, error)

Runs lista execucoes com filtro e paginacao.

func (*LeituraRepo) RunsDoWorkflow

func (r *LeituraRepo) RunsDoWorkflow(ctx context.Context, slug string, limite int) ([]ResumoRun, error)

RunsDoWorkflow lista as execucoes de um workflow so, para a tela dele.

func (*LeituraRepo) UltimasRuns

func (r *LeituraRepo) UltimasRuns(ctx context.Context, limite int) ([]ResumoRun, error)

UltimasRuns lista as execucoes mais recentes.

func (*LeituraRepo) Workflows

func (r *LeituraRepo) Workflows(ctx context.Context) ([]ResumoWorkflow, error)

Workflows lista os workflows publicados com sua agenda e ultimo estado.

type LogDoPasso

type LogDoPasso struct {
	NodeID    string
	Tentativa int
	Status    string
	ExitCode  *int
	Erro      string
	Log       string
	DuracaoMs int64
}

LogDoPasso e a saida de uma tentativa, para a tela da execucao.

type Pool

type Pool struct {
	*pgxpool.Pool
}

Pool envolve o pgxpool. O tipo existe para que o resto do sistema dependa de algo nosso, e nao do driver diretamente.

func New

func New(ctx context.Context, url string) (*Pool, error)

New abre o pool e verifica a conexao antes de devolver. Um pool que so falha no primeiro uso transforma erro de configuracao em erro de request.

func (*Pool) Check

func (p *Pool) Check(ctx context.Context) error

Check e o contrato de health: um ping com prazo. Sem timeout, um banco lento faria o readiness pendurar em vez de reprovar.

type ResumoProjeto

type ResumoProjeto struct {
	Slug      string
	Nome      string
	Workflows int
	Runs      int
	CriadoEm  time.Time
}

ResumoProjeto conta o que existe sob um projeto.

type ResumoRun

type ResumoRun struct {
	ID           string
	WorkflowSlug string
	Status       string
	TriggerType  string
	Tentativa    int
	LogicalDate  *time.Time
	CriadoEm     time.Time
	IniciadoEm   *time.Time
	Duracao      *time.Duration
	Erro         string
}

ResumoRun e uma linha da lista de execucoes.

type ResumoWorkflow

type ResumoWorkflow struct {
	Slug         string
	Nome         string
	Projeto      string
	Cron         string
	Timezone     string
	Catchup      bool
	Ativo        bool
	TemAgenda    bool
	UltimoSlot   *time.Time
	UltimoStatus string
	TotalRuns    int

	// Da ultima execucao — a coluna "Latest Run" da lista.
	UltimaRunID *string
	UltimaRunEm *time.Time

	Tags []string

	// ProximaRun nao vem do banco: e calculada a partir do cron, no consumidor.
	// Guardar no banco exigiria recalcular a cada mudanca de agenda e conviver
	// com o valor obsoleto entre uma e outra.
	ProximaRun *time.Time
}

ResumoWorkflow junta o workflow, sua agenda e o estado da ultima execucao.

type RunRepo

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

RunRepo persiste Runs e TaskRuns.

func NewRunRepo

func NewRunRepo(p *Pool) *RunRepo

func (*RunRepo) Buscar

func (r *RunRepo) Buscar(ctx context.Context, id uuid.UUID) (dom.Run, error)

Buscar le um Run.

func (*RunRepo) ContarPorStatus

func (r *RunRepo) ContarPorStatus(ctx context.Context) (map[dom.Status]int, error)

ContarPorStatus e o que o criterio de aceite da PHASE 2 mede.

func (*RunRepo) ContarPorTrigger

func (r *RunRepo) ContarPorTrigger(ctx context.Context) (map[string]int, error)

ContarPorTrigger mostra a origem dos runs — distinguir backfill de agendado e o que a secao 12 pede ao investigar um incidente.

func (*RunRepo) Criar

func (r *RunRepo) Criar(ctx context.Context, run dom.Run) (dom.Run, error)

Criar insere o Run em CREATED.

A colisao na unique de idempotency_key vira ErrJaExiste, nao erro generico: e o caso da secao 29 — o scheduler caiu depois de criar e tenta de novo ao subir.

func (*RunRepo) EstadoDosNos

func (r *RunRepo) EstadoDosNos(ctx context.Context, runID uuid.UUID) (map[string]EstadoNo, error)

EstadoDosNos devolve o estado de cada no na ULTIMA tentativa de cada um.

`DISTINCT ON` em vez de max(attempt) num subselect: a tentativa mais recente e a que interessa na tela, e uma tentativa antiga que falhou nao deve pintar o no de vermelho depois de o retry ter dado certo.

func (*RunRepo) IncrementarTentativa

func (r *RunRepo) IncrementarTentativa(ctx context.Context, id uuid.UUID) (int, error)

IncrementarTentativa sobe o contador ao reenfileirar por retry.

func (*RunRepo) IniciarTask

func (r *RunRepo) IniciarTask(ctx context.Context, runID uuid.UUID, nodeID string, tentativa int) error

IniciarTask registra o inicio de um passo.

`ON CONFLICT DO UPDATE` na chave (run, node, tentativa): reexecutar o mesmo passo na mesma tentativa e idempotente, o que importa quando o dispatcher recupera um item de worker morto e o refaz.

func (*RunRepo) LogsDaRun

func (r *RunRepo) LogsDaRun(ctx context.Context, runID uuid.UUID) ([]LogDoPasso, error)

LogsDaRun devolve a saida de cada tentativa de cada passo, em ordem de execucao.

TODAS as tentativas, nao so a ultima: quando um passo passa na segunda, o que explica a primeira falha esta justamente na tentativa que a tela descartaria.

func (*RunRepo) PassoJaTeveSucesso

func (r *RunRepo) PassoJaTeveSucesso(ctx context.Context, workflowSlug, nodeID string, exceto uuid.UUID) (bool, error)

PassoJaTeveSucesso responde se este passo, neste workflow, ja terminou bem antes — em qualquer run anterior.

E o que decide se a execucao atual e a PRIMEIRA daquele passo, informacao que vai para o ambiente do passo e que o SDK usa para criar a tabela de destino. A alternativa seria o SDK inferir de "a tabela nao existe", e aí alguem apaga a tabela por engano e a proxima execucao se acha a primeira.

Por (workflow, passo), nao por workflow: um workflow com tres fetchers escrevendo em tres tabelas criaria apenas a do primeiro passo se a resposta fosse do workflow inteiro.

`exceto` e o run corrente, excluido para que a propria tentativa em curso nao conte como sucesso anterior.

func (*RunRepo) PassoQueFalhou

func (r *RunRepo) PassoQueFalhou(ctx context.Context, runID uuid.UUID) (string, string, error)

PassoQueFalhou devolve o node e a saida da ultima tentativa que falhou.

`ORDER BY iniciado_em DESC` e nao `attempt DESC`: num grafo com varios passos, a maior tentativa pode ser de um passo que ja tinha falhado e sido superado — o que interessa e o que falhou POR ULTIMO, que e onde a execucao parou.

Ausencia nao e erro: um run que morreu antes de qualquer passo comecar (imagem inexistente, fila cancelada) nao tem task_run nenhuma, e o alerta sai sem esta parte em vez de nao sair.

func (*RunRepo) RegistrarErro

func (r *RunRepo) RegistrarErro(ctx context.Context, id uuid.UUID, msg string) error

RegistrarErro guarda a causa da falha.

func (*RunRepo) RegistrarEtapas

func (r *RunRepo) RegistrarEtapas(ctx context.Context, runID uuid.UUID, nodeID string,
	tentativa int, sdkVersao string, etapas json.RawMessage) error

RegistrarEtapas grava o avanco das etapas de um passo do SDK.

Sobrescreve o array inteiro em vez de acrescentar: o coletor do runner ja guarda UMA entrada por etapa, com o estado atual, e a tela quer quatro blocos e nao um diario.

func (*RunRepo) TerminarTask

func (r *RunRepo) TerminarTask(ctx context.Context, runID uuid.UUID, nodeID string,
	tentativa int, status dom.Status, exit *int, erro string, log string) error

TerminarTask registra o desfecho.

func (*RunRepo) Transicionar

func (r *RunRepo) Transicionar(ctx context.Context, id uuid.UUID, para dom.Status) error

Transicionar aplica a mudanca de estado, validando ANTES de escrever.

A validacao acontece contra o estado lido dentro da transacao, com FOR UPDATE: ler fora dela permitiria que dois dispatchers lessem "queued" e ambos escrevessem "running".

type ScheduleRepo

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

ScheduleRepo le e atualiza agendas.

func NewScheduleRepo

func NewScheduleRepo(p *Pool) *ScheduleRepo

func (*ScheduleRepo) Alternar

func (r *ScheduleRepo) Alternar(ctx context.Context, slug string) (bool, error)

Alternar inverte o estado atual numa unica ida ao banco.

func (*ScheduleRepo) Ativas

func (r *ScheduleRepo) Ativas(ctx context.Context) ([]sch.Schedule, error)

Ativas lista as agendas que o scheduler deve avaliar.

func (*ScheduleRepo) AvancarSlot

func (r *ScheduleRepo) AvancarSlot(ctx context.Context, slug string, slot time.Time) error

func (*ScheduleRepo) DefinirAtivo

func (r *ScheduleRepo) DefinirAtivo(ctx context.Context, slug string, ativo bool) (bool, error)

AvancarSlot marca ate onde a agenda ja foi materializada.

A condicao `ultimo_slot IS NULL OR ultimo_slot < $2` torna a operacao idempotente e segura sob concorrencia: dois schedulers avaliando a mesma agenda nunca fazem o marcador retroceder. DefinirAtivo pausa ou retoma uma agenda e devolve o estado resultante.

Devolve em vez de so gravar porque a UI alterna sem saber o valor atual: sem o retorno, a tela precisaria de uma segunda consulta e ficaria sujeita a corrida entre dois operadores clicando ao mesmo tempo.

Pausar NAO cancela o que ja esta na fila: os runs materializados sao trabalho aceito, e descarta-los ao pausar surpreenderia quem so queria parar de criar novos.

type WorkflowRepo

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

WorkflowRepo persiste a definicao publicada de um workflow.

func NewWorkflowRepo

func NewWorkflowRepo(p *Pool) *WorkflowRepo

func (*WorkflowRepo) Definicao

func (r *WorkflowRepo) Definicao(ctx context.Context, slug string) (wf.Workflow, error)

Definicao le o grafo publicado.

func (*WorkflowRepo) Podar

func (r *WorkflowRepo) Podar(ctx context.Context, projeto uuid.UUID, manter []string) ([]string, error)

Podar remove do projeto os workflows que NAO estao na lista, junto com suas agendas. Devolve os slugs removidos.

Existe porque publicar so adicionava: tirar um arquivo da pasta nao tirava nada do banco, e o scheduler continuava materializando runs de um workflow que ninguem enxergava mais. Com agendas de 15 minutos, isso e trabalho invisivel rodando para sempre.

O historico (`runs`) NAO e apagado: ele referencia o slug como texto, nao por chave estrangeira, justamente para sobreviver a remocao da definicao. Apagar a execucao junto seria apagar a evidencia do que aconteceu.

func (*WorkflowRepo) Publicar

func (r *WorkflowRepo) Publicar(ctx context.Context, w wf.Workflow, projeto uuid.UUID) error

Publicar grava o workflow e sua agenda numa transacao.

As duas coisas juntas, e nao em chamadas separadas: publicar o grafo sem a agenda deixaria um workflow que nunca dispara, e a agenda sem o grafo faria o scheduler criar runs de algo que nao existe.

Jump to

Keyboard shortcuts

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