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
8 changes: 8 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -60,3 +60,11 @@ SEFAZ_AMBIENTE=Producao

# Desligue (false) para subir a aplicacao sem worker: nenhum job roda e nada consulta a SEFAZ.
AGENDAMENTO_ATIVO=true

# Empresas enfileiradas por ciclo de trinta minutos. Empresa sem novidade volta em uma hora, entao
# metade da base fica elegivel a cada ciclo: 500 comporta mil empresas sem atraso.
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
27 changes: 27 additions & 0 deletions RUNBOOK.md
Original file line number Diff line number Diff line change
Expand Up @@ -148,6 +148,33 @@ AGENDAMENTO_ATIVO=false docker compose up -d --build api
A API atende normalmente (grade, relatórios, download), só não roda job nenhum. Para religar, suba de
novo sem a variável: `docker compose up -d api`.

### Capacidade do ciclo — quantas empresas o servidor comporta

Dois números governam isso, e ambos são configuração:

```
AGENDAMENTO_EMPRESAS_POR_CICLO=500 # quantas entram na fila a cada trinta minutos
AGENDAMENTO_WORKERS=8 # quantas saem dela ao mesmo tempo
```

Empresa sem novidade é reagendada para daqui a uma hora, então a cada ciclo cerca de metade da base
fica elegível. O padrão de 500 atende **mil empresas** sem atraso. Acima disso o excedente escorrega
para o ciclo seguinte — nada se perde, porque a ordem é pela consulta mais antiga, mas o intervalo
efetivo cresce em silêncio.

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:

```sql
SELECT count(*) FROM empresas
WHERE ativa AND proxima_consulta_em <= now() AND (bloqueada_ate IS NULL OR bloqueada_ate <= now());
```

Se o número passar de `AGENDAMENTO_EMPRESAS_POR_CICLO`, suba os dois valores juntos. Enfileirar mais do
que os workers drenam só transfere a espera de lugar. Cada worker fica bloqueado em rede durante a
chamada à SEFAZ — cerca de nove segundos —, não em processamento, então o número pode passar bem do
total de núcleos.

### Falha ao arquivar no object storage

O arquivamento repete sozinho até três vezes, com espera crescente e jitter, quando a falha é de rede ou
Expand Down
2 changes: 2 additions & 0 deletions docker-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -110,6 +110,8 @@ services:
# Desligado, a aplicação sobe sem worker: nenhum job roda e nada consulta a SEFAZ. É o que
# permite subir uma versão nova do código sem gastar cota de consulta do CNPJ.
Agendamento__ProcessarTrabalhos: ${AGENDAMENTO_ATIVO:-true}
Agendamento__EmpresasPorCiclo: ${AGENDAMENTO_EMPRESAS_POR_CICLO:-500}
Agendamento__Workers: ${AGENDAMENTO_WORKERS:-8}
# 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
2 changes: 1 addition & 1 deletion src/eContabil.Api/Program.cs
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,7 @@
// Nem toda instância processa fila: desligado, a aplicação segue enfileirando, mas o worker, o
// dashboard e o registro recorrente ficam de fora — nenhum deles é acessível sem tocar o storage.
var processaTrabalhos = builder.Configuration.GetValue("Agendamento:ProcessarTrabalhos", true);
builder.Services.AddAgendamento(conexaoPostgres, processaTrabalhos);
builder.Services.AddAgendamento(builder.Configuration, conexaoPostgres);

// O Kestrel escreve o próprio cabeçalho Server ao montar a resposta, depois de qualquer middleware:
// removê-lo no pipeline não surte efeito, e o nome do servidor só ajuda quem procura exploit conhecido.
Expand Down
4 changes: 3 additions & 1 deletion src/eContabil.Api/appsettings.json
Original file line number Diff line number Diff line change
Expand Up @@ -83,7 +83,9 @@
"TetoDeLotes": 5
},
"Agendamento": {
"ProcessarTrabalhos": true
"ProcessarTrabalhos": true,
"EmpresasPorCiclo": 500,
"Workers": 8
},
"AllowedHosts": "*"
}
50 changes: 50 additions & 0 deletions src/eContabil.Application/Sincronizacao/OpcoesAgendamento.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
using System.ComponentModel.DataAnnotations;

namespace eContabil.Application.Sincronizacao;

/// <summary>
/// Capacidade do agendador.
/// </summary>
/// <remarks>
/// Estes dois números decidem quantas empresas o escritório comporta, e por isso não podem viver como
/// constante no código. Uma base que cresce de 200 para 2.000 clientes precisa deles ajustáveis sem
/// recompilar — e quem opera precisa poder subi-los durante um pico sem esperar um deploy.
///
/// A relação entre eles importa. <see cref="EmpresasPorCiclo"/> limita quantas entram na fila; os
/// <see cref="Workers"/> limitam quantas saem dela dentro do ciclo. Enfileirar muito além do que os
/// workers drenam só transfere a espera de um lugar para o outro.
/// </remarks>
public sealed class OpcoesAgendamento
{
public const string Secao = "Agendamento";

/// <summary>Liga o processamento de trabalhos neste processo.</summary>
/// <remarks>
/// Desligado, a API atende normalmente e não roda job nenhum — é como se sobe uma versão nova sem
/// que o worker retome os ciclos e gaste consulta à SEFAZ numa janela de manutenção.
/// </remarks>
public bool ProcessarTrabalhos { get; set; } = true;

/// <summary>
/// Empresas enfileiradas por ciclo de captura.
/// </summary>
/// <remarks>
/// Empresa sem novidade é reagendada para daqui a uma hora, então a cada ciclo de trinta minutos
/// cerca de metade da base fica elegível. O padrão de 500 comporta mil empresas sem atraso; acima
/// disso o excedente escorrega para o ciclo seguinte, em silêncio e sem perda — a ordem é pela
/// consulta mais antiga, então ninguém fica para trás indefinidamente.
/// </remarks>
[Range(1, 20_000)]
public int EmpresasPorCiclo { get; set; } = 500;

/// <summary>
/// Trabalhos simultâneos por processo.
/// </summary>
/// <remarks>
/// A biblioteca fiscal só expõe chamada síncrona, e cada consulta à SEFAZ leva cerca de nove
/// segundos. O worker fica bloqueado em entrada e saída, não em processamento — por isso este
/// número pode passar bem do número de núcleos sem saturar a máquina.
/// </remarks>
[Range(1, 200)]
public int Workers { get; set; } = 8;
}
Original file line number Diff line number Diff line change
Expand Up @@ -199,6 +199,12 @@ private async Task<Result> ProcessarLoteAsync(
// estouravam o índice único no commit, derrubando o lote inteiro e deixando os XMLs órfãos.
var chavesDoLote = new HashSet<string>(StringComparer.Ordinal);

// Os já gravados saem numa consulta só, antes do laço. Perguntar por um de cada vez custava
// cinquenta consultas por lote e mil numa execução que encadeia vinte, para responder algo que
// um único `IN` resolve.
var jaGravados = await documentos.ObterPorEmpresaEChavesAsync(
empresa.Id, ChavesValidasDe(retorno), ct);

foreach (var recebido in retorno.Documentos)
{
var chave = ChaveAcesso.Criar(recebido.ChaveAcesso);
Expand Down Expand Up @@ -233,8 +239,7 @@ private async Task<Result> ProcessarLoteAsync(
continue;
}

var existente = await documentos.ObterPorEmpresaEChaveAsync(empresa.Id, chave.Valor, ct);
if (existente is not null)
if (jaGravados.TryGetValue(chave.Valor.Valor, out var existente))
{
if (arquivamento.Valor is { } objeto && existente.OrigemConteudo is OrigemConteudo.Resumo)
{
Expand Down Expand Up @@ -276,6 +281,30 @@ private async Task<Result> ProcessarLoteAsync(
return Result.Ok();
}

/// <summary>
/// As chaves aproveitáveis do lote, sem repetição.
/// </summary>
/// <remarks>
/// Item com chave inválida fica de fora — ele seria descartado no laço de qualquer forma, e levá-lo
/// à consulta só alargaria o <c>IN</c> sem chance de casar com nada.
/// </remarks>
private static List<ChaveAcesso> ChavesValidasDe(RetornoDistribuicaoDto retorno)
{
var chaves = new HashSet<ChaveAcesso>();

foreach (var recebido in retorno.Documentos)
{
var chave = ChaveAcesso.Criar(recebido.ChaveAcesso);

if (chave.Sucesso)
{
chaves.Add(chave.Valor);
}
}

return [.. chaves];
}

/// <summary>
/// Arquiva o XML quando houver. Devolve sucesso com valor nulo para resumo, que não tem arquivo.
/// </summary>
Expand Down
11 changes: 11 additions & 0 deletions src/eContabil.Domain/Documentos/IDocumentoFiscalRepository.cs
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,17 @@ public interface IDocumentoFiscalRepository : IRepository<DocumentoFiscal>
{
Task<DocumentoFiscal?> ObterPorEmpresaEChaveAsync(Guid empresaId, ChaveAcesso chave, CancellationToken ct);

/// <summary>
/// Os documentos da empresa entre as chaves informadas, indexados por chave.
/// </summary>
/// <remarks>
/// Existe para o processamento de lote. A SEFAZ entrega até cinquenta documentos por chamada, e
/// perguntar por um de cada vez custava cinquenta consultas por lote — mil numa execução que
/// encadeia vinte. Uma consulta responde o mesmo.
/// </remarks>
Task<IReadOnlyDictionary<string, DocumentoFiscal>> ObterPorEmpresaEChavesAsync(
Guid empresaId, IReadOnlyCollection<ChaveAcesso> chaves, CancellationToken ct);

/// <summary>
/// Indica se o XML daquela chave já foi arquivado por qualquer empresa.
/// </summary>
Expand Down
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
using eContabil.Application.Sincronizacao;
using eContabil.Domain.Documentos;
using Hangfire;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;

namespace eContabil.Infrastructure.Agendamento;

Expand All @@ -14,16 +16,17 @@ namespace eContabil.Infrastructure.Agendamento;
/// </remarks>
public sealed partial class ManifestacaoGeralJob(
IDocumentoFiscalRepository documentos,
IOptions<OpcoesAgendamento> opcoes,
ILogger<ManifestacaoGeralJob> log)
{
/// <summary>Teto de empresas por ciclo, para não inundar a fila de uma vez.</summary>
private const int LimitePorCiclo = 500;
private readonly OpcoesAgendamento _opcoes =
opcoes?.Value ?? throw new ArgumentNullException(nameof(opcoes));

[Queue(Filas.Padrao)]
[DisableConcurrentExecution(timeoutInSeconds: 300)]
public async Task ExecutarAsync(CancellationToken ct)
{
var empresas = await documentos.ObterEmpresasComManifestacaoPendenteAsync(LimitePorCiclo, ct);
var empresas = await documentos.ObterEmpresasComManifestacaoPendenteAsync(_opcoes.EmpresasPorCiclo, ct);

foreach (var empresaId in empresas)
{
Expand Down
Original file line number Diff line number Diff line change
@@ -1,7 +1,9 @@
using eContabil.Application.Sincronizacao;
using eContabil.Domain.Empresas;
using eContabil.Shared;
using Hangfire;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;

namespace eContabil.Infrastructure.Agendamento;

Expand All @@ -16,18 +18,19 @@ namespace eContabil.Infrastructure.Agendamento;
public sealed partial class SincronizacaoGeralJob(
IEmpresaRepository empresas,
TimeProvider relogio,
IOptions<OpcoesAgendamento> opcoes,
ILogger<SincronizacaoGeralJob> log)
{
/// <summary>Teto de empresas enfileiradas por ciclo, para não inundar a fila de uma vez.</summary>
private const int LimitePorCiclo = 500;
private readonly OpcoesAgendamento _opcoes =
opcoes?.Value ?? throw new ArgumentNullException(nameof(opcoes));

[Queue(Filas.Padrao)]
[DisableConcurrentExecution(timeoutInSeconds: 300)]
public async Task ExecutarAsync(CancellationToken ct)
{
var agora = relogio.AgoraUtc();

var elegiveis = await empresas.ObterElegiveisParaSincronizarAsync(agora, LimitePorCiclo, ct);
var elegiveis = await empresas.ObterElegiveisParaSincronizarAsync(agora, _opcoes.EmpresasPorCiclo, ct);

foreach (var empresa in elegiveis)
{
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,79 @@
using System.Security.Cryptography;
using eContabil.Application.Certificados;
using eContabil.Shared;

namespace eContabil.Infrastructure.Certificados;

/// <summary>
/// Evita rebaixar e decifrar o mesmo certificado várias vezes dentro de uma execução.
/// </summary>
/// <remarks>
/// Drenar uma empresa atrasada encadeia até vinte consultas à SEFAZ — duzentas na carga inicial — e cada
/// uma pedia o certificado ao cofre de novo. Era o mesmo arquivo baixado do object storage, decifrado
/// com AES-GCM e derivado com PBKDF2 duzentas vezes seguidas, para produzir exatamente o mesmo material.
///
/// O escopo do cache é a execução, não o processo. Como o registro é <c>Scoped</c> e cada trabalho do
/// Hangfire abre o próprio escopo, a instância nasce e morre com a sincronização de uma empresa. Isso
/// preserva a regra que importa — chave privada não persiste entre execuções, não vai para disco e não
/// atravessa empresas — e elimina só a repetição dentro da janela em que o material já estaria na
/// memória de qualquer forma.
///
/// O descarte zera os bytes decifrados. Não é garantia absoluta contra despejo de memória, mas encurta
/// a janela em que a chave fica legível no heap depois de deixar de ser necessária.
/// </remarks>
public sealed class CofreComCacheDeExecucao(ICertificadoCofre cofre) : ICertificadoCofre, IDisposable
{
private readonly Dictionary<Guid, CertificadoCarregadoDto> _carregados = [];

public Task<Result<CertificadoArmazenadoDto>> GuardarAsync(
Guid certificadoId,
byte[] pfx,
string senha,
string cnpjEsperado,
DateTime agoraUtc,
CancellationToken ct) =>
cofre.GuardarAsync(certificadoId, pfx, senha, cnpjEsperado, agoraUtc, ct);

public async Task<Result<CertificadoCarregadoDto>> CarregarAsync(
CertificadoLocalizacao localizacao, CancellationToken ct)
{
ArgumentNullException.ThrowIfNull(localizacao);

if (_carregados.TryGetValue(localizacao.CertificadoId, out var jaCarregado))
{
return Result<CertificadoCarregadoDto>.Ok(jaCarregado);
}

var carregado = await cofre.CarregarAsync(localizacao, ct);

if (carregado.Sucesso)
{
// A chave é o identificador do certificado, não o da empresa: a substituição de um
// certificado gera um identificador novo, e o material antigo nunca é servido no lugar dele.
_carregados[localizacao.CertificadoId] = carregado.Valor;
}

return carregado;
}

public Task<Result> RemoverAsync(string bucket, string objectName, CancellationToken ct)
{
// Remover invalida tudo: o arquivo deixou de existir, e servir o que estava em memória entregaria
// uma chave privada que o operador acabou de mandar apagar.
Descartar();

return cofre.RemoverAsync(bucket, objectName, ct);
}

public void Dispose() => Descartar();

private void Descartar()
{
foreach (var certificado in _carregados.Values)
{
CryptographicOperations.ZeroMemory(certificado.Pfx);
}

_carregados.Clear();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -155,7 +155,14 @@ argumentos.Outcome.Exception is { } falha
services.AddSingleton<ILeitorDeCertificado, LeitorDeCertificado>();
services.AddSingleton<IExtratorDeDadosFiscais, ExtratorDeDadosFiscais>();
services.AddScoped<IXmlStorage, MinioXmlStorage>();
services.AddScoped<ICertificadoCofre, CertificadoCofre>();

// O cofre real fica atrás do cache de execução. Drenar uma empresa atrasada encadeia até
// duzentas consultas à SEFAZ, e cada uma pedia o mesmo certificado de volta ao object storage
// para decifrá-lo outra vez. O escopo do cache é o trabalho do Hangfire: a chave privada
// continua sem persistir entre execuções.
services.AddScoped<CertificadoCofre>();
services.AddScoped<ICertificadoCofre>(provedor =>
new CofreComCacheDeExecucao(provedor.GetRequiredService<CertificadoCofre>()));

return services;
}
Expand Down Expand Up @@ -239,24 +246,35 @@ public static IServiceCollection AddAutenticacao(
/// enfileirando, mas não executa — é o que permite escalar API e processamento em separado.
/// </remarks>
public static IServiceCollection AddAgendamento(
this IServiceCollection services, string conexaoPostgres, bool processarTrabalhos = true)
this IServiceCollection services, IConfiguration configuracao, string conexaoPostgres)
{
ArgumentNullException.ThrowIfNull(configuracao);
ArgumentException.ThrowIfNullOrWhiteSpace(conexaoPostgres);

services.AddOptions<OpcoesAgendamento>()
.Bind(configuracao.GetSection(OpcoesAgendamento.Secao))
.ValidateDataAnnotations()
.ValidateOnStart();

// Lido aqui, e não por `IOptions`, porque a quantidade de workers é decidida na construção do
// servidor do Hangfire — antes de o provedor de serviços existir.
var opcoes = configuracao.GetSection(OpcoesAgendamento.Secao).Get<OpcoesAgendamento>()
?? new OpcoesAgendamento();

services.AddHangfire(configuracao => configuracao
.SetDataCompatibilityLevel(CompatibilityLevel.Version_180)
.UseSimpleAssemblyNameTypeSerializer()
.UseRecommendedSerializerSettings()
.UsePostgreSqlStorage(opcoes => opcoes.UseNpgsqlConnection(conexaoPostgres)));

if (processarTrabalhos)
if (opcoes.ProcessarTrabalhos)
{
services.AddHangfireServer(opcoes =>
services.AddHangfireServer(servidor =>
{
// A ordem das filas é a prioridade de atendimento. A carga inicial vem por último: é
// longa e não pode atrasar o ciclo recorrente das empresas já em dia.
opcoes.Queues = [Filas.Sincronizacao, Filas.Padrao, Filas.CargaInicial];
opcoes.WorkerCount = 8;
servidor.Queues = [Filas.Sincronizacao, Filas.Padrao, Filas.CargaInicial];
servidor.WorkerCount = opcoes.Workers;
});
}

Expand Down
Loading
Loading