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
1 change: 1 addition & 0 deletions Directory.Packages.props
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@
<PackageVersion Include="Minio" Version="7.0.0" />
<PackageVersion Include="Npgsql.EntityFrameworkCore.PostgreSQL" Version="10.0.3" />
<PackageVersion Include="NSubstitute" Version="6.0.0" />
<PackageVersion Include="Polly.Extensions" Version="8.7.0" />
<PackageVersion Include="Serilog.AspNetCore" Version="10.0.0" />
<PackageVersion Include="Serilog.Enrichers.Environment" Version="3.0.1" />
<PackageVersion Include="Serilog.Enrichers.Thread" Version="4.0.0" />
Expand Down
1 change: 1 addition & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ auditável de tudo que foi emitido contra ou pelos seus clientes.
| Banco | PostgreSQL 17 · EF Core 10 (Npgsql) |
| Storage | MinIO (S3-compatible) |
| Agendamento | Hangfire |
| Resiliência | Polly 8 (nova tentativa no object storage) |
| Fiscal | ZeusFiscal (`Hercules.NET.NFe.NFCe`) |
| Frontend | Angular · Signals · NgRx SignalStore · Angular Material |
| Infra | Docker + Docker Compose |
Expand Down
13 changes: 13 additions & 0 deletions RUNBOOK.md
Original file line number Diff line number Diff line change
Expand Up @@ -148,6 +148,19 @@ 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`.

### 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
o MinIO responde 5xx. Erro permanente — arquivo inexistente, bucket errado, credencial inválida — falha
na primeira tentativa, porque insistir só atrasaria o lote para chegar ao mesmo lugar.

A repetição existe por causa da SEFAZ, não do MinIO: falha ao arquivar aborta o lote, o ponteiro de NSU
não avança, e a rodada seguinte busca os mesmos documentos na SEFAZ de novo. A consulta é o recurso
racionado.

Se `storage.indisponivel` aparecer no log da sincronização, as quatro tentativas já se esgotaram — o
MinIO está fora, não oscilando. Confira `/health/ready` antes de mexer no job.

### "A nota não tem itens, CFOP nem impostos na tela"

O documento foi capturado antes de a extração existir: o XML está arquivado, mas nada foi lido dele.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
using eContabil.Infrastructure.Persistencia;
using eContabil.Infrastructure.Persistencia.Consultas;
using eContabil.Infrastructure.Persistencia.Repositorios;
using eContabil.Infrastructure.Resiliencia;
using eContabil.Infrastructure.Sefaz;
using eContabil.Infrastructure.Seguranca;
using eContabil.Infrastructure.Storage;
Expand All @@ -31,9 +32,12 @@
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.Configuration;
using Microsoft.Extensions.Diagnostics.HealthChecks;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
using Microsoft.IdentityModel.Tokens;
using Minio;
using Polly;
using Polly.Retry;

namespace Microsoft.Extensions.DependencyInjection;

Expand Down Expand Up @@ -112,6 +116,41 @@ public static IServiceCollection AddArmazenamentoESeguranca(
.Build();
});

// Nova tentativa curta no object storage.
//
// Não é resiliência genérica: é economia de cota da SEFAZ. Uma falha ao arquivar aborta o lote
// inteiro, o ponteiro de NSU não avança, e a rodada seguinte busca **os mesmos documentos na
// SEFAZ de novo** — a consulta é o recurso racionado, a gravação no MinIO não custa nada. Um
// soluço de rede de meio segundo não pode custar uma ida à SEFAZ.
//
// Sem disjuntor de propósito. Ele é por processo e em memória, enquanto o que precisa ser
// contido aqui já é contido no lugar certo: o bloqueio por empresa, que vive no banco, sobrevive
// a reinício e vale para todas as instâncias. Um disjuntor sobre "o MinIO" ainda pararia a
// captura de todas as empresas por causa de uma, sem nada em troca — se o storage caiu de vez, o
// job já falha rápido e o ciclo seguinte retoma.
services.AddResiliencePipeline(FalhaTransitoriaDeStorage.Pipeline, (construtor, contexto) =>
{
construtor
.ConfigureTelemetry(contexto.ServiceProvider.GetRequiredService<ILoggerFactory>())
.AddRetry(new RetryStrategyOptions
{
MaxRetryAttempts = 3,
BackoffType = DelayBackoffType.Exponential,

// Jitter porque a sincronização roda em paralelo por empresa: sem ele, as que
// falharem juntas voltariam juntas, e a segunda tentativa recriaria a rajada.
UseJitter = true,
Delay = TimeSpan.FromMilliseconds(200),
ShouldHandle = argumentos => ValueTask.FromResult(
argumentos.Outcome.Exception is { } falha
&& FalhaTransitoriaDeStorage.Ehtransitoria(falha))
})

// Mais interno: vale por tentativa. Sem ele, uma conexão pendurada seguraria o worker
// até o tempo limite do próprio job, e as demais empresas ficariam esperando atrás.
.AddTimeout(TimeSpan.FromSeconds(30));
});

services.AddSingleton<ICifrador, CifradorAesGcm>();
services.AddSingleton<ILeitorDeCertificado, LeitorDeCertificado>();
services.AddSingleton<IExtratorDeDadosFiscais, ExtratorDeDadosFiscais>();
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,42 @@
using System.Net.Sockets;
using Minio.Exceptions;
using Polly.Timeout;

namespace eContabil.Infrastructure.Resiliencia;

/// <summary>
/// Decide se uma falha do object storage merece nova tentativa.
/// </summary>
/// <remarks>
/// A regra é de <b>permissão</b>, não de exclusão: repete só o que se sabe transitório, e trata todo o
/// resto como definitivo. O SDK do MinIO faz quase toda exceção herdar de <see cref="MinioException"/> —
/// inclusive <c>ObjectNotFoundException</c>, <c>AccessDeniedException</c> e <c>InvalidBucketNameException</c>.
/// Um <c>Handle&lt;MinioException&gt;</c> genérico repetiria arquivo inexistente e credencial errada,
/// gastando tempo do job para chegar ao mesmo lugar.
///
/// A lista de permissão também envelhece melhor: uma exceção nova do SDK entra como definitiva, que é o
/// comportamento seguro. Numa lista de exclusão, ela entraria como repetível sem ninguém decidir isso.
/// </remarks>
public static class FalhaTransitoriaDeStorage
{
/// <summary>Nome do pipeline registrado na injeção de dependência.</summary>
public const string Pipeline = "storage-xml";

public static bool Ehtransitoria(Exception excecao) => excecao switch
{
// Cancelamento é decisão de quem chamou, nunca falha do storage.
OperationCanceledException => false,

// O MinIO não respondeu, respondeu 5xx ou cortou a resposta no meio.
ConnectionException or InternalServerException or UnexpectedMinioException
or UnexpectedShortReadException => true,

// Rede abaixo do SDK: conexão recusada, reset, DNS momentaneamente fora.
HttpRequestException or SocketException or IOException => true,

// Estouro do limite por tentativa imposto pelo próprio pipeline.
TimeoutRejectedException => true,

_ => false
};
}
112 changes: 78 additions & 34 deletions src/eContabil.Infrastructure/Storage/MinioXmlStorage.cs
Original file line number Diff line number Diff line change
Expand Up @@ -3,11 +3,15 @@
using eContabil.Application.Armazenamento;
using eContabil.Domain.ValueObjects;
using eContabil.Infrastructure.Compactacao;
using eContabil.Infrastructure.Resiliencia;
using eContabil.Shared;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Options;
using Minio;
using Minio.DataModel.Args;
using Minio.Exceptions;
using Polly;
using Polly.Registry;

namespace eContabil.Infrastructure.Storage;

Expand All @@ -17,10 +21,29 @@ namespace eContabil.Infrastructure.Storage;
/// <remarks>
/// O nome vem da chave, e não do identificador da empresa, porque é isso que permite **um único objeto**
/// quando a mesma NF-e interessa a duas empresas do escritório. As duas apontam para ele.
///
/// Toda chamada ao SDK passa por um pipeline de nova tentativa. As operações usadas aqui são todas
/// idempotentes — gravar o mesmo objeto com o mesmo conteúdo, ler, conferir existência —, então repetir
/// não duplica nem corrompe nada.
/// </remarks>
public sealed class MinioXmlStorage(IMinioClient cliente, IOptions<MinioOptions> opcoes) : IXmlStorage
public sealed class MinioXmlStorage : IXmlStorage
{
private readonly MinioOptions _opcoes = opcoes?.Value ?? throw new ArgumentNullException(nameof(opcoes));
private readonly IMinioClient _cliente;
private readonly MinioOptions _opcoes;
private readonly ResiliencePipeline _tentativas;

public MinioXmlStorage(
IMinioClient cliente,
IOptions<MinioOptions> opcoes,
ResiliencePipelineProvider<string> pipelines)
{
ArgumentNullException.ThrowIfNull(opcoes);
ArgumentNullException.ThrowIfNull(pipelines);

_cliente = cliente;
_opcoes = opcoes.Value;
_tentativas = pipelines.GetPipeline(FalhaTransitoriaDeStorage.Pipeline);
}

public Task<Result<ObjetoArmazenadoDto>> SalvarAsync(
string chaveAcesso, Stream xml, CancellationToken ct) =>
Expand Down Expand Up @@ -64,16 +87,7 @@ private async Task<Result<ObjetoArmazenadoDto>> GravarAsync(
_opcoes.BucketXmls, objectName, hash, conteudo.Length, JaExistia: true));
}

using var origem = new MemoryStream(conteudo);

await cliente.PutObjectAsync(
new PutObjectArgs()
.WithBucket(_opcoes.BucketXmls)
.WithObject(objectName)
.WithStreamData(origem)
.WithObjectSize(conteudo.Length)
.WithContentType("application/xml"),
ct);
await EnviarAsync(_opcoes.BucketXmls, objectName, conteudo, ct);

return Result<ObjetoArmazenadoDto>.Ok(new ObjetoArmazenadoDto(
_opcoes.BucketXmls, objectName, hash, conteudo.Length, JaExistia: false));
Expand All @@ -98,18 +112,27 @@ public async Task<Result<Stream>> ObterAsync(string bucket, string objectName, C
{
try
{
// O conteúdo é copiado para um buffer porque o SDK entrega o corpo por callback e fecha a
// conexão ao final da chamada. XML de NF-e tem dezenas de KB — o custo é irrelevante.
var destino = new MemoryStream();

await cliente.GetObjectAsync(
new GetObjectArgs()
.WithBucket(bucket)
.WithObject(objectName)
.WithCallbackStream((fluxo, cancelamento) => fluxo.CopyToAsync(destino, cancelamento)),
// O buffer nasce dentro do pipeline: uma tentativa que falhou no meio da cópia deixa bytes
// parciais nele, e reaproveitá-lo na tentativa seguinte concatenaria os dois pedaços.
var lido = await _tentativas.ExecuteAsync(
async cancelamento =>
{
// O conteúdo é copiado para um buffer porque o SDK entrega o corpo por callback e
// fecha a conexão ao final da chamada. XML de NF-e tem dezenas de KB.
var destino = new MemoryStream();

await _cliente.GetObjectAsync(
new GetObjectArgs()
.WithBucket(bucket)
.WithObject(objectName)
.WithCallbackStream((fluxo, token) => fluxo.CopyToAsync(destino, token)),
cancelamento);

return destino.ToArray();
},
ct);

var conteudo = ConteudoGzip.Descompactar(destino.ToArray());
var conteudo = ConteudoGzip.Descompactar(lido);

return Result<Stream>.Ok(new MemoryStream(conteudo, writable: false));
}
Expand Down Expand Up @@ -142,16 +165,7 @@ public async Task<Result<ObjetoArmazenadoDto>> SubstituirAsync(

try
{
using var origem = new MemoryStream(bytes);

await cliente.PutObjectAsync(
new PutObjectArgs()
.WithBucket(bucket)
.WithObject(objectName)
.WithStreamData(origem)
.WithObjectSize(bytes.Length)
.WithContentType("application/xml"),
ct);
await EnviarAsync(bucket, objectName, bytes, ct);

return Result<ObjetoArmazenadoDto>.Ok(new ObjetoArmazenadoDto(
bucket,
Expand All @@ -171,8 +185,10 @@ public async Task<bool> ExisteAsync(string bucket, string objectName, Cancellati
{
try
{
await cliente.StatObjectAsync(
new StatObjectArgs().WithBucket(bucket).WithObject(objectName), ct);
await _tentativas.ExecuteAsync(
async cancelamento => await _cliente.StatObjectAsync(
new StatObjectArgs().WithBucket(bucket).WithObject(objectName), cancelamento),
ct);

return true;
}
Expand All @@ -182,6 +198,34 @@ await cliente.StatObjectAsync(
}
}

/// <summary>
/// Envia o conteúdo para o objeto, repetindo enquanto a falha for transitória.
/// </summary>
/// <remarks>
/// O <see cref="MemoryStream"/> é criado <b>dentro</b> do pipeline, uma vez por tentativa. Um stream
/// só é lido do início ao fim: reaproveitá-lo faria a segunda tentativa encontrá-lo no fim e gravar
/// um objeto de zero byte — sem erro nenhum, porque do ponto de vista do SDK o envio deu certo.
///
/// Repetir é seguro porque a operação é idempotente: o mesmo nome com o mesmo conteúdo produz
/// exatamente o mesmo objeto, e o XML autorizado nunca muda.
/// </remarks>
private async Task EnviarAsync(string bucket, string objectName, byte[] conteudo, CancellationToken ct) =>
await _tentativas.ExecuteAsync(
async cancelamento =>
{
using var origem = new MemoryStream(conteudo, writable: false);

await _cliente.PutObjectAsync(
new PutObjectArgs()
.WithBucket(bucket)
.WithObject(objectName)
.WithStreamData(origem)
.WithObjectSize(conteudo.Length)
.WithContentType("application/xml"),
cancelamento);
},
ct);

/// <summary>
/// Particiona por ano e mês da competência, que a própria chave carrega.
/// </summary>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
<PackageReference Include="Microsoft.IdentityModel.JsonWebTokens" />
<PackageReference Include="Minio" />
<PackageReference Include="Npgsql.EntityFrameworkCore.PostgreSQL" />
<PackageReference Include="Polly.Extensions" />
</ItemGroup>

</Project>
Loading
Loading