Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -68,3 +68,7 @@ AGENDAMENTO_EMPRESAS_POR_CICLO=500
# Trabalhos simultaneos. Cada worker fica bloqueado em rede durante a chamada a SEFAZ, nao em CPU,
# entao este numero pode passar do numero de nucleos.
AGENDAMENTO_WORKERS=8

# Atraso a partir do qual o ciclo avisa no log. Precisa ficar acima de um ciclo inteiro (30 min), que e
# espera normal de quem atingiu o teto de lotes, e abaixo do prazo prometido ao cliente.
AGENDAMENTO_ATRASO_ACEITAVEL=00:45:00
20 changes: 20 additions & 0 deletions RUNBOOK.md
Original file line number Diff line number Diff line change
Expand Up @@ -162,6 +162,26 @@ fica elegível. O padrão de 500 atende **mil empresas** sem atraso. Acima disso
para o ciclo seguinte — nada se perde, porque a ordem é pela consulta mais antiga, mas o intervalo
efetivo cresce em silêncio.

**O ciclo avisa sozinho.** Quando a empresa mais atrasada da fila passa de `AGENDAMENTO_ATRASO_ACEITAVEL`
(45 minutos por padrão), o log recebe:

```
Captura atrasada: a empresa mais antiga da fila esperou 52,3 minuto(s) além do previsto.
500 elegível(is) neste ciclo, teto de 500. Suba EmpresasPorCiclo e Workers juntos se o número se repetir
```

O limiar fica acima de um ciclo inteiro de propósito: empresa que atinge o teto de lotes volta para a
fila na hora e espera até trinta minutos pelo ciclo seguinte, o que é operação normal. Um limiar menor
dispararia sem nada estar errado, e alerta que grita à toa é alerta que ninguém lê.

As mesmas informações saem como métrica no medidor `eContabil.Captura`, para quando houver coletor:

```
econtabil.captura.atraso_maximo (segundos)
econtabil.captura.empresas_elegiveis
econtabil.captura.empresas_enfileiradas
```

Sintoma de saturação: a coluna "Próxima consulta" da tela de empresas fica no passado para muitas
delas ao mesmo tempo. Confirme no banco:

Expand Down
1 change: 1 addition & 0 deletions docker-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -112,6 +112,7 @@ services:
Agendamento__ProcessarTrabalhos: ${AGENDAMENTO_ATIVO:-true}
Agendamento__EmpresasPorCiclo: ${AGENDAMENTO_EMPRESAS_POR_CICLO:-500}
Agendamento__Workers: ${AGENDAMENTO_WORKERS:-8}
Agendamento__AtrasoAceitavel: ${AGENDAMENTO_ATRASO_ACEITAVEL:-00:45:00}
# Origens exatas separadas por ponto e vírgula. Vazio fecha a política — nunca curinga, porque a
# API responde com credenciais.
Cors__OrigensPermitidas: ${CORS_ORIGENS:-http://localhost:4200}
Expand Down
3 changes: 2 additions & 1 deletion src/eContabil.Api/appsettings.json
Original file line number Diff line number Diff line change
Expand Up @@ -85,7 +85,8 @@
"Agendamento": {
"ProcessarTrabalhos": true,
"EmpresasPorCiclo": 500,
"Workers": 8
"Workers": 8,
"AtrasoAceitavel": "00:45:00"
},
"AllowedHosts": "*"
}
14 changes: 14 additions & 0 deletions src/eContabil.Application/Sincronizacao/OpcoesAgendamento.cs
Original file line number Diff line number Diff line change
Expand Up @@ -47,4 +47,18 @@ public sealed class OpcoesAgendamento
/// </remarks>
[Range(1, 200)]
public int Workers { get; set; } = 8;

/// <summary>
/// Atraso a partir do qual o ciclo passa a avisar no log.
/// </summary>
/// <remarks>
/// Mede há quanto tempo a empresa mais atrasada já deveria ter sido consultada. Precisa ficar acima
/// de um ciclo inteiro: empresa que atingiu o teto de lotes volta para a fila imediatamente e espera
/// até trinta minutos pelo ciclo seguinte, o que é operação normal — um limiar de trinta minutos
/// dispararia sem que nada estivesse errado.
///
/// E precisa ficar abaixo do prazo prometido ao cliente, para sobrar tempo de reagir antes de
/// furá-lo. Quarenta e cinco minutos deixam quinze de margem sobre uma promessa de uma hora.
/// </remarks>
public TimeSpan AtrasoAceitavel { get; set; } = TimeSpan.FromMinutes(45);
}
45 changes: 43 additions & 2 deletions src/eContabil.Infrastructure/Agendamento/SincronizacaoGeralJob.cs
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
using eContabil.Application.Sincronizacao;
using eContabil.Domain.Empresas;
using eContabil.Infrastructure.Observabilidade;
using eContabil.Shared;
using Hangfire;
using Microsoft.Extensions.Logging;
Expand All @@ -14,11 +15,16 @@ namespace eContabil.Infrastructure.Agendamento;
/// Este job **não** sincroniza — ele só enfileira. Percorrer as empresas em série aqui faria uma empresa
/// lenta segurar a fila inteira, e uma que falhasse levaria as seguintes junto. Enfileirando, cada uma
/// ganha isolamento de falha e nova tentativa própria.
///
/// Enfileira pelo <c>IBackgroundJobClient</c>, e não pela API estática do Hangfire: é o que a própria
/// biblioteca recomenda, e é o que permite exercitar o ciclo sem subir o armazenamento de trabalhos.
/// </remarks>
public sealed partial class SincronizacaoGeralJob(
IEmpresaRepository empresas,
TimeProvider relogio,
IOptions<OpcoesAgendamento> opcoes,
IBackgroundJobClient fila,
MetricasDeCaptura metricas,
ILogger<SincronizacaoGeralJob> log)
{
private readonly OpcoesAgendamento _opcoes =
Expand All @@ -38,20 +44,55 @@ public async Task ExecutarAsync(CancellationToken ct)
// workers próprios.
if (empresa.EmCargaInicial)
{
BackgroundJob.Enqueue<CargaInicialJob>(job => job.ExecutarAsync(empresa.Id, default));
fila.Enqueue<CargaInicialJob>(job => job.ExecutarAsync(empresa.Id, default));
}
else
{
BackgroundJob.Enqueue<SincronizarEmpresaJob>(
fila.Enqueue<SincronizarEmpresaJob>(
job => job.ExecutarAsync(empresa.Id, ModoSincronizacao.Incremental, default));
}
}

RegistrarAtraso(elegiveis, agora);

CicloEnfileirado(log, elegiveis.Count);
}

/// <summary>
/// Publica o atraso do ciclo e avisa quando ele passa do aceitável.
/// </summary>
/// <remarks>
/// O atraso sai da primeira empresa da lista, sem consulta extra: a seleção já vem ordenada pela
/// consulta mais antiga, então quem está no topo é justamente quem mais esperou.
///
/// Mede o atraso da fila — a parte que se resolve com <c>EmpresasPorCiclo</c> e <c>Workers</c>. Não
/// mede o tempo até uma nota nova aparecer: esse depende também de <c>IntervaloSemNovidade</c>, que
/// é com que frequência uma empresa sem novidade volta a ser consultada.
/// </remarks>
private void RegistrarAtraso(IReadOnlyList<Empresa> elegiveis, DateTime agora)
{
var atraso = elegiveis.Count > 0
? agora - elegiveis[0].ProximaConsultaEm
: TimeSpan.Zero;

metricas.RegistrarCiclo(atraso, elegiveis.Count, elegiveis.Count);

if (atraso > _opcoes.AtrasoAceitavel)
{
CicloAtrasado(log, atraso.TotalMinutes, elegiveis.Count, _opcoes.EmpresasPorCiclo);
}
}

[LoggerMessage(
Level = LogLevel.Information,
Message = "Ciclo de sincronização enfileirou {Quantidade} empresa(s) elegível(is)")]
private static partial void CicloEnfileirado(ILogger log, int quantidade);

[LoggerMessage(
Level = LogLevel.Warning,
Message = "Captura atrasada: a empresa mais antiga da fila esperou {AtrasoEmMinutos:0.0} minuto(s) " +
"além do previsto. {Elegiveis} elegível(is) neste ciclo, teto de {Teto}. " +
"Suba EmpresasPorCiclo e Workers juntos se o número se repetir")]
private static partial void CicloAtrasado(
ILogger log, double atrasoEmMinutos, int elegiveis, int teto);
}
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
using eContabil.Infrastructure.Certificados;
using eContabil.Infrastructure.Diagnostico;
using eContabil.Infrastructure.Fiscal;
using eContabil.Infrastructure.Observabilidade;
using eContabil.Infrastructure.Persistencia;
using eContabil.Infrastructure.Persistencia.Consultas;
using eContabil.Infrastructure.Persistencia.Repositorios;
Expand Down Expand Up @@ -265,7 +266,16 @@ public static IServiceCollection AddAgendamento(
.SetDataCompatibilityLevel(CompatibilityLevel.Version_180)
.UseSimpleAssemblyNameTypeSerializer()
.UseRecommendedSerializerSettings()
.UsePostgreSqlStorage(opcoes => opcoes.UseNpgsqlConnection(conexaoPostgres)));
.UsePostgreSqlStorage(
postgres => postgres.UseNpgsqlConnection(conexaoPostgres),
new PostgreSqlStorageOptions
{
// Declarado, e não herdado do padrão do provedor: é este número que decide quanto
// tempo um trabalho recém-enfileirado espera antes de alguém pegá-lo. O ciclo
// enfileira centenas de uma vez a cada trinta minutos, e uma sondagem longa
// acrescentaria essa espera ao atraso que a métrica acompanha.
QueuePollInterval = TimeSpan.FromSeconds(5)
}));

if (opcoes.ProcessarTrabalhos)
{
Expand All @@ -278,6 +288,10 @@ public static IServiceCollection AddAgendamento(
});
}

// Singleton: o medidor precisa sobreviver aos ciclos para o coletor encontrar o valor entre um
// e outro. Registrado antes dos jobs, que dependem dele.
services.AddSingleton<MetricasDeCaptura>();

services.AddScoped<SincronizacaoGeralJob>();
services.AddScoped<SincronizarEmpresaJob>();
services.AddScoped<CargaInicialJob>();
Expand Down
65 changes: 65 additions & 0 deletions src/eContabil.Infrastructure/Observabilidade/MetricasDeCaptura.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,65 @@
using System.Diagnostics.Metrics;

namespace eContabil.Infrastructure.Observabilidade;

/// <summary>
/// Publica o atraso do ciclo de captura como métrica.
/// </summary>
/// <remarks>
/// Usa o <see cref="Meter"/> da biblioteca padrão, e não OpenTelemetry. A diferença importa: emitir não
/// exige dependência nenhuma, só consumir exige. Dá para instrumentar hoje e ligar um coletor meses
/// depois sem tocar neste código — e quando isso acontecer, o Polly já vem junto, porque ele publica no
/// medidor <c>Polly</c> desde que a nova tentativa do object storage entrou.
///
/// Os medidores são observáveis e leem de campo em memória. Consultar o banco no retorno de chamada
/// faria uma consulta a cada raspagem do coletor — de quinze em quinze segundos, tipicamente — para um
/// valor que só muda quando o ciclo roda, de trinta em trinta minutos.
/// </remarks>
public sealed class MetricasDeCaptura : IDisposable
{
/// <summary>Nome do medidor. É por ele que o coletor assina estas métricas.</summary>
public const string Nome = "eContabil.Captura";

private readonly Meter _medidor;

// Guardado em ticks porque `Interlocked` não opera sobre TimeSpan, e o ciclo escreve enquanto o
// coletor lê.
private long _atrasoEmTicks;
private long _elegiveis;
private long _enfileiradas;

public MetricasDeCaptura()
{
_medidor = new Meter(Nome);

_medidor.CreateObservableGauge(
"econtabil.captura.atraso_maximo",
() => TimeSpan.FromTicks(Interlocked.Read(ref _atrasoEmTicks)).TotalSeconds,
unit: "s",
description: "Há quanto tempo a empresa mais atrasada já deveria ter sido consultada.");

_medidor.CreateObservableGauge(
"econtabil.captura.empresas_elegiveis",
() => Interlocked.Read(ref _elegiveis),
description: "Empresas em condição de sincronizar no início do ciclo.");

_medidor.CreateObservableGauge(
"econtabil.captura.empresas_enfileiradas",
() => Interlocked.Read(ref _enfileiradas),
description: "Empresas que o ciclo conseguiu enfileirar.");
}

/// <summary>Registra o que o ciclo encontrou.</summary>
/// <remarks>
/// <paramref name="atraso"/> é do documento mais atrasado da fila, não a média: média esconde
/// justamente a empresa que está furando o prazo.
/// </remarks>
public void RegistrarCiclo(TimeSpan atraso, int elegiveis, int enfileiradas)
{
Interlocked.Exchange(ref _atrasoEmTicks, Math.Max(atraso.Ticks, 0));
Interlocked.Exchange(ref _elegiveis, elegiveis);
Interlocked.Exchange(ref _enfileiradas, enfileiradas);
}

public void Dispose() => _medidor.Dispose();
}
Loading
Loading