diff --git a/.env.example b/.env.example
index 0df7bc8..aeb2f65 100644
--- a/.env.example
+++ b/.env.example
@@ -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
diff --git a/RUNBOOK.md b/RUNBOOK.md
index 483fe6e..bd98ae6 100644
--- a/RUNBOOK.md
+++ b/RUNBOOK.md
@@ -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:
diff --git a/docker-compose.yml b/docker-compose.yml
index 9e66006..7735b44 100644
--- a/docker-compose.yml
+++ b/docker-compose.yml
@@ -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}
diff --git a/src/eContabil.Api/appsettings.json b/src/eContabil.Api/appsettings.json
index 0a3f940..97704e7 100644
--- a/src/eContabil.Api/appsettings.json
+++ b/src/eContabil.Api/appsettings.json
@@ -85,7 +85,8 @@
"Agendamento": {
"ProcessarTrabalhos": true,
"EmpresasPorCiclo": 500,
- "Workers": 8
+ "Workers": 8,
+ "AtrasoAceitavel": "00:45:00"
},
"AllowedHosts": "*"
}
diff --git a/src/eContabil.Application/Sincronizacao/OpcoesAgendamento.cs b/src/eContabil.Application/Sincronizacao/OpcoesAgendamento.cs
index 03463d6..9edb72f 100644
--- a/src/eContabil.Application/Sincronizacao/OpcoesAgendamento.cs
+++ b/src/eContabil.Application/Sincronizacao/OpcoesAgendamento.cs
@@ -47,4 +47,18 @@ public sealed class OpcoesAgendamento
///
[Range(1, 200)]
public int Workers { get; set; } = 8;
+
+ ///
+ /// Atraso a partir do qual o ciclo passa a avisar no log.
+ ///
+ ///
+ /// 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.
+ ///
+ public TimeSpan AtrasoAceitavel { get; set; } = TimeSpan.FromMinutes(45);
}
diff --git a/src/eContabil.Infrastructure/Agendamento/SincronizacaoGeralJob.cs b/src/eContabil.Infrastructure/Agendamento/SincronizacaoGeralJob.cs
index beffb93..dd1b7a4 100644
--- a/src/eContabil.Infrastructure/Agendamento/SincronizacaoGeralJob.cs
+++ b/src/eContabil.Infrastructure/Agendamento/SincronizacaoGeralJob.cs
@@ -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;
@@ -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 IBackgroundJobClient, 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.
///
public sealed partial class SincronizacaoGeralJob(
IEmpresaRepository empresas,
TimeProvider relogio,
IOptions opcoes,
+ IBackgroundJobClient fila,
+ MetricasDeCaptura metricas,
ILogger log)
{
private readonly OpcoesAgendamento _opcoes =
@@ -38,20 +44,55 @@ public async Task ExecutarAsync(CancellationToken ct)
// workers próprios.
if (empresa.EmCargaInicial)
{
- BackgroundJob.Enqueue(job => job.ExecutarAsync(empresa.Id, default));
+ fila.Enqueue(job => job.ExecutarAsync(empresa.Id, default));
}
else
{
- BackgroundJob.Enqueue(
+ fila.Enqueue(
job => job.ExecutarAsync(empresa.Id, ModoSincronizacao.Incremental, default));
}
}
+ RegistrarAtraso(elegiveis, agora);
+
CicloEnfileirado(log, elegiveis.Count);
}
+ ///
+ /// Publica o atraso do ciclo e avisa quando ele passa do aceitável.
+ ///
+ ///
+ /// 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 EmpresasPorCiclo e Workers. Não
+ /// mede o tempo até uma nota nova aparecer: esse depende também de IntervaloSemNovidade, que
+ /// é com que frequência uma empresa sem novidade volta a ser consultada.
+ ///
+ private void RegistrarAtraso(IReadOnlyList 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);
}
diff --git a/src/eContabil.Infrastructure/InfraestruturaServiceCollectionExtensions.cs b/src/eContabil.Infrastructure/InfraestruturaServiceCollectionExtensions.cs
index 6d691a2..89e0aaf 100644
--- a/src/eContabil.Infrastructure/InfraestruturaServiceCollectionExtensions.cs
+++ b/src/eContabil.Infrastructure/InfraestruturaServiceCollectionExtensions.cs
@@ -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;
@@ -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)
{
@@ -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();
+
services.AddScoped();
services.AddScoped();
services.AddScoped();
diff --git a/src/eContabil.Infrastructure/Observabilidade/MetricasDeCaptura.cs b/src/eContabil.Infrastructure/Observabilidade/MetricasDeCaptura.cs
new file mode 100644
index 0000000..73d2720
--- /dev/null
+++ b/src/eContabil.Infrastructure/Observabilidade/MetricasDeCaptura.cs
@@ -0,0 +1,65 @@
+using System.Diagnostics.Metrics;
+
+namespace eContabil.Infrastructure.Observabilidade;
+
+///
+/// Publica o atraso do ciclo de captura como métrica.
+///
+///
+/// Usa o 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 Polly 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.
+///
+public sealed class MetricasDeCaptura : IDisposable
+{
+ /// Nome do medidor. É por ele que o coletor assina estas métricas.
+ 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.");
+ }
+
+ /// Registra o que o ciclo encontrou.
+ ///
+ /// é do documento mais atrasado da fila, não a média: média esconde
+ /// justamente a empresa que está furando o prazo.
+ ///
+ 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();
+}
diff --git a/tests/eContabil.Infrastructure.Tests/Agendamento/SincronizacaoGeralJobTests.cs b/tests/eContabil.Infrastructure.Tests/Agendamento/SincronizacaoGeralJobTests.cs
new file mode 100644
index 0000000..a9becda
--- /dev/null
+++ b/tests/eContabil.Infrastructure.Tests/Agendamento/SincronizacaoGeralJobTests.cs
@@ -0,0 +1,171 @@
+using System.Diagnostics.Metrics;
+using eContabil.Application.Sincronizacao;
+using eContabil.Domain.Empresas;
+using eContabil.Domain.ValueObjects;
+using eContabil.Infrastructure.Agendamento;
+using eContabil.Infrastructure.Observabilidade;
+using Hangfire;
+using Microsoft.Extensions.Logging;
+using Microsoft.Extensions.Options;
+using Microsoft.Extensions.Time.Testing;
+using NSubstitute;
+using Xunit;
+
+namespace eContabil.Infrastructure.Tests.Agendamento;
+
+///
+/// O ciclo mede há quanto tempo a empresa mais atrasada já deveria ter sido consultada.
+///
+///
+/// Sem esse sinal, saturação é invisível: o excedente escorrega para o ciclo seguinte em silêncio, e
+/// quem descobre é o contador reclamando que a nota não apareceu.
+///
+public class SincronizacaoGeralJobTests
+{
+ private static readonly DateTime AgoraUtc = new(2026, 7, 29, 12, 0, 0, DateTimeKind.Utc);
+
+ private readonly IEmpresaRepository _empresas = Substitute.For();
+
+ [Fact]
+ public async Task ExecutarAsync_PublicaOAtrasoDaEmpresaMaisAntiga()
+ {
+ // A seleção já vem ordenada pela consulta mais antiga, então o atraso sai da primeira da lista —
+ // sem consulta extra ao banco.
+ ComElegiveis(AtrasadaEm(TimeSpan.FromMinutes(20)), AtrasadaEm(TimeSpan.FromMinutes(5)));
+
+ using var metricas = new MetricasDeCaptura();
+ using var coletor = new ColetorDeMetricas();
+
+ await Montar(metricas).ExecutarAsync(TestContext.Current.CancellationToken);
+
+ Assert.Equal(1_200, coletor.Ler("econtabil.captura.atraso_maximo", metricas));
+ Assert.Equal(2, coletor.Ler("econtabil.captura.empresas_elegiveis", metricas));
+ }
+
+ [Fact]
+ public async Task ExecutarAsync_SemEmpresaElegivel_PublicaAtrasoZero()
+ {
+ // Fila vazia é o estado saudável, não ausência de informação: o medidor precisa dizer zero, e
+ // não manter o último valor alto de um ciclo anterior.
+ ComElegiveis();
+
+ using var metricas = new MetricasDeCaptura();
+ using var coletor = new ColetorDeMetricas();
+
+ await Montar(metricas).ExecutarAsync(TestContext.Current.CancellationToken);
+
+ Assert.Equal(0, coletor.Ler("econtabil.captura.atraso_maximo", metricas));
+ }
+
+ [Fact]
+ public async Task ExecutarAsync_ComAtrasoAcimaDoAceitavel_Avisa()
+ {
+ ComElegiveis(AtrasadaEm(TimeSpan.FromMinutes(50)));
+
+ var log = LogQueRegistra();
+ using var metricas = new MetricasDeCaptura();
+
+ await Montar(metricas, log).ExecutarAsync(TestContext.Current.CancellationToken);
+
+ Assert.Equal(1, Avisos(log));
+ }
+
+ [Fact]
+ public async Task ExecutarAsync_ComAtrasoDeUmCicloInteiro_NaoAvisa()
+ {
+ // Empresa que atingiu o teto de lotes volta para a fila na hora e espera até trinta minutos pelo
+ // ciclo seguinte. É operação normal — avisar aqui treinaria o operador a ignorar o alerta.
+ ComElegiveis(AtrasadaEm(TimeSpan.FromMinutes(30)));
+
+ var log = LogQueRegistra();
+ using var metricas = new MetricasDeCaptura();
+
+ await Montar(metricas, log).ExecutarAsync(TestContext.Current.CancellationToken);
+
+ Assert.Equal(0, Avisos(log));
+ }
+
+ private void ComElegiveis(params Empresa[] elegiveis) =>
+ _empresas.ObterElegiveisParaSincronizarAsync(
+ Arg.Any(), Arg.Any(), Arg.Any())
+ .Returns(elegiveis);
+
+ private SincronizacaoGeralJob Montar(
+ MetricasDeCaptura metricas, ILogger? log = null) =>
+ new(_empresas,
+ new FakeTimeProvider(new DateTimeOffset(AgoraUtc, TimeSpan.Zero)),
+ Options.Create(new OpcoesAgendamento()),
+ Substitute.For(),
+ metricas,
+ log ?? Substitute.For>());
+
+ ///
+ /// O código gerado por LoggerMessage consulta IsEnabled antes de montar a mensagem, e
+ /// o dublê responde "não" por padrão — sem isto, nenhum aviso chegaria a ser emitido e o teste
+ /// passaria a verificar o silêncio do dublê, não o comportamento do ciclo.
+ ///
+ private static ILogger LogQueRegistra()
+ {
+ var log = Substitute.For>();
+ log.IsEnabled(Arg.Any()).Returns(true);
+
+ return log;
+ }
+
+ private static int Avisos(ILogger log) =>
+ log.ReceivedCalls()
+ .Count(chamada => chamada.GetMethodInfo().Name == nameof(ILogger.Log)
+ && chamada.GetArguments()[0] is LogLevel.Warning);
+
+ /// Empresa cuja próxima consulta já venceu há .
+ private static Empresa AtrasadaEm(TimeSpan atraso)
+ {
+ var empresa = Empresa.Criar(
+ "NEWNUTRITION BRASIL ALIMENTOS LTDA", null, Cnpj.Criar("11222333000181").Valor, null,
+ Uf.Criar("SP").Valor, AmbienteSefaz.Producao, AgoraUtc).Valor;
+
+ empresa.AgendarProximaConsulta(AgoraUtc - atraso);
+
+ return empresa;
+ }
+
+ /// Lê o valor publicado por um medidor observável.
+ ///
+ /// O medidor só produz valor quando alguém o observa. Sem um ouvinte, o retorno de chamada nunca
+ /// roda — e um teste que apenas chamasse RegistrarCiclo provaria que o campo foi escrito, não
+ /// que a métrica sai.
+ ///
+ private sealed class ColetorDeMetricas : IDisposable
+ {
+ private readonly MeterListener _ouvinte = new();
+ private readonly Dictionary _valores = [];
+
+ public ColetorDeMetricas()
+ {
+ _ouvinte.InstrumentPublished = (instrumento, ouvinte) =>
+ {
+ if (instrumento.Meter.Name == MetricasDeCaptura.Nome)
+ {
+ ouvinte.EnableMeasurementEvents(instrumento);
+ }
+ };
+
+ _ouvinte.SetMeasurementEventCallback(
+ (instrumento, valor, _, _) => _valores[instrumento.Name] = valor);
+
+ _ouvinte.SetMeasurementEventCallback(
+ (instrumento, valor, _, _) => _valores[instrumento.Name] = valor);
+
+ _ouvinte.Start();
+ }
+
+ public double Ler(string instrumento, MetricasDeCaptura _)
+ {
+ _ouvinte.RecordObservableInstruments();
+
+ return _valores.TryGetValue(instrumento, out var valor) ? valor : double.NaN;
+ }
+
+ public void Dispose() => _ouvinte.Dispose();
+ }
+}