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(); + } +}