diff --git a/src/Models/ClefEvent.cs b/src/Models/ClefEvent.cs index 1f50fc5..9eb4455 100644 --- a/src/Models/ClefEvent.cs +++ b/src/Models/ClefEvent.cs @@ -12,6 +12,11 @@ public class ClefEvent public string? MessageTemplate { get; set; } public string? Exception { get; set; } public string? SourceFile { get; set; } - public Dictionary? Properties { get; set; } + /// + /// Propriedades estruturadas do evento. A interface (e não Dictionary) + /// permite ao parser publicar a forma compacta de ; + /// testes e integrações continuam atribuindo um Dictionary comum. + /// + public IReadOnlyDictionary? Properties { get; set; } } } diff --git a/src/Models/PropriedadesEvento.cs b/src/Models/PropriedadesEvento.cs new file mode 100644 index 0000000..89f0901 --- /dev/null +++ b/src/Models/PropriedadesEvento.cs @@ -0,0 +1,141 @@ +using System.Collections; +using Serilog.Events; +using Serilog.Parsing; + +namespace ClefExplorer.Models +{ + /// + /// As propriedades de um evento, em dois arrays paralelos com o tamanho exato. + /// + /// O que ficava aqui custava caro no + /// agregado: buckets + entries + overhead de objeto por evento, multiplicados por + /// centenas de milhares de eventos — para um conjunto que tem ~18 pares, nunca muda + /// depois de construído e quase sempre é lido por enumeração, não por chave. O lookup + /// linear com primeiro (as chaves vêm do + /// pool da sessão, então a MESMA instância aparece em todo evento) empata com o hash + /// para n desse tamanho. + /// + public sealed class PropriedadesEvento : IReadOnlyDictionary + { + public static PropriedadesEvento Vazio { get; } = new(null); + + private readonly string[] _chaves; + private readonly LogEventPropertyValue[] _valores; + + /// + /// Constrói a partir da lista do parser, com a MESMA semântica do dicionário que + /// substitui: chave repetida (OrdinalIgnoreCase) faz o último valor vencer. + /// + public PropriedadesEvento(IReadOnlyList? itens) + { + if (itens is null || itens.Count == 0) + { + _chaves = Array.Empty(); + _valores = Array.Empty(); + return; + } + + var chaves = new string[itens.Count]; + var valores = new LogEventPropertyValue[itens.Count]; + var quantidade = 0; + + for (var i = 0; i < itens.Count; i++) + { + var nome = itens[i].Name; + + // Um único loop com o teste combinado: duplicata é raríssima, e rodar o + // fallback OrdinalIgnoreCase completo por inserção (como o lookup de + // leitura faz) custava caro multiplicado por 18 propriedades por evento. + var indice = -1; + for (var j = 0; j < quantidade; j++) + { + if (ReferenceEquals(chaves[j], nome) + || string.Equals(chaves[j], nome, StringComparison.OrdinalIgnoreCase)) + { + indice = j; + break; + } + } + + if (indice >= 0) + { + valores[indice] = itens[i].Value; + continue; + } + + chaves[quantidade] = nome; + valores[quantidade] = itens[i].Value; + quantidade++; + } + + if (quantidade == itens.Count) + { + _chaves = chaves; + _valores = valores; + } + else + { + // Havia duplicata: os arrays exatos evitam carregar as sobras para sempre. + _chaves = chaves.AsSpan(0, quantidade).ToArray(); + _valores = valores.AsSpan(0, quantidade).ToArray(); + } + } + + /// Acesso direto para os caminhos quentes (filtro), sem enumerator. + internal LogEventPropertyValue[] ValoresInternos => _valores; + + internal string[] ChavesInternas => _chaves; + + public int Count => _chaves.Length; + + public IEnumerable Keys => _chaves; + + public IEnumerable Values => _valores; + + public LogEventPropertyValue this[string key] => + TryGetValue(key, out var valor) ? valor : throw new KeyNotFoundException(key); + + public bool ContainsKey(string key) => TryGetValue(key, out _); + + public bool TryGetValue(string key, out LogEventPropertyValue value) + { + var indice = IndiceDe(_chaves, _chaves.Length, key); + if (indice >= 0) + { + value = _valores[indice]; + return true; + } + + value = null!; + return false; + } + + private static int IndiceDe(string[] chaves, int quantidade, string chave) + { + // Referência primeiro: com o pool, "SourceContext" é UMA instância no processo + // inteiro e a comparação é um ponteiro. O fallback mantém o contrato do + // dicionário antigo (OrdinalIgnoreCase) para chaves vindas de fora do pool. + for (var i = 0; i < quantidade; i++) + { + if (ReferenceEquals(chaves[i], chave)) return i; + } + + for (var i = 0; i < quantidade; i++) + { + if (string.Equals(chaves[i], chave, StringComparison.OrdinalIgnoreCase)) return i; + } + + return -1; + } + + public IEnumerator> GetEnumerator() + { + for (var i = 0; i < _chaves.Length; i++) + { + yield return new KeyValuePair(_chaves[i], _valores[i]); + } + } + + IEnumerator IEnumerable.GetEnumerator() => GetEnumerator(); + } +} diff --git a/src/Services/LeitorArquivoLog.cs b/src/Services/LeitorArquivoLog.cs index fcec829..946d747 100644 --- a/src/Services/LeitorArquivoLog.cs +++ b/src/Services/LeitorArquivoLog.cs @@ -61,6 +61,32 @@ public sealed class LeitorArquivoLog : ILeitorArquivoLog // grande em @x é comum). private const int TamanhoBuffer = 64 * 1024; + private readonly long _limiarParalelo; + private readonly long _segmentoMinimo; + + public LeitorArquivoLog() + : this(limiarParalelo: 32L * 1024 * 1024, segmentoMinimo: 8L * 1024 * 1024) + { + } + + /// + /// Ajusta quando a leitura de UM arquivo passa a ser dividida entre núcleos. + /// Existe como construtor (e não como constantes) para os testes exercitarem o + /// caminho paralelo com arquivos de quilobytes em vez de gigabytes. + /// + /// Tamanho a partir do qual o arquivo é segmentado. + /// + /// Menor fatia que vale um worker: abaixo disso o custo de abrir o stream e alinhar + /// a fronteira supera o ganho de paralelismo. + /// + public LeitorArquivoLog(long limiarParalelo, long segmentoMinimo) + { + ArgumentOutOfRangeException.ThrowIfLessThan(limiarParalelo, 1); + ArgumentOutOfRangeException.ThrowIfLessThan(segmentoMinimo, 1); + _limiarParalelo = limiarParalelo; + _segmentoMinimo = segmentoMinimo; + } + public async Task LerAsync( string arquivo, IReadOnlyList textosIgnorados, @@ -80,12 +106,28 @@ public async Task LerAsync( if (arquivo.EndsWith(".gz", StringComparison.OrdinalIgnoreCase)) { + // gz é um fluxo: não há como posicionar num byte do meio sem descomprimir + // tudo antes dele, então o paralelismo por segmento não se aplica. await using var compactado = new GZipStream(streamArquivo, CompressionMode.Decompress); return await LerStreamAsync( compactado, arquivo, textosIgnorados, offsetFinal: null, + CacheDeTemplates.Para(pool), + verificarBom: true, + cancellationToken).ConfigureAwait(false); + } + + // O Parallel.ForEachAsync do LogStore distribui por ARQUIVO: um único .clef + // de 1 GB ocupava um núcleo só enquanto os outros 19 esperavam. Acima do + // limiar, o próprio arquivo é dividido em segmentos alinhados a '\n'. + if (streamArquivo.Length >= _limiarParalelo) + { + return await LerParaleloAsync( + streamArquivo, + arquivo, + textosIgnorados, pool, cancellationToken).ConfigureAwait(false); } @@ -95,22 +137,200 @@ public async Task LerAsync( arquivo, textosIgnorados, offsetFinal: () => streamArquivo.Position, - pool, + CacheDeTemplates.Para(pool), + verificarBom: true, cancellationToken).ConfigureAwait(false); } + /// + /// Divide o arquivo em segmentos que começam exatamente em início de linha e lê + /// cada um num worker próprio, preservando a ordem do arquivo na concatenação. + /// + /// As fronteiras são resolvidas ANTES dos workers: um seek por fronteira, + /// avançando até o primeiro \n. Assim cada worker roda o laço sequencial + /// já existente sobre um recorte fechado, sem coordenação entre eles — a linha + /// que cruza uma fronteira bruta pertence, por construção, ao segmento anterior. + /// + private async Task LerParaleloAsync( + FileStream stream, + string arquivo, + IReadOnlyList textosIgnorados, + PoolDeTextos? pool, + CancellationToken cancellationToken) + { + // Comprimento capturado uma vez: o arquivo pode continuar crescendo (tail), + // e ler além de L produziria segmentos com fronteiras móveis. O que entrar + // depois fica para o acompanhamento ao vivo — OffsetFinal diz até onde fomos. + var comprimento = stream.Length; + var alvo = (long)Math.Clamp(Environment.ProcessorCount, 2, 32); + var tamanhoSegmento = Math.Max(_segmentoMinimo, comprimento / alvo); + + var fronteiras = new List { 0 }; + for (var bruta = tamanhoSegmento; bruta < comprimento; bruta += tamanhoSegmento) + { + var alinhada = await AcharInicioDeLinhaAsync(stream, bruta, comprimento, cancellationToken) + .ConfigureAwait(false); + // Linha gigante pode empurrar a fronteira além da próxima bruta; só + // fronteiras estritamente crescentes viram segmento. + if (alinhada > fronteiras[^1] && alinhada < comprimento) fronteiras.Add(alinhada); + } + fronteiras.Add(comprimento); + + var cache = CacheDeTemplates.Para(pool); + var tarefas = new Task[fronteiras.Count - 1]; + for (var i = 0; i < tarefas.Length; i++) + { + var inicio = fronteiras[i]; + var fim = fronteiras[i + 1]; + var primeiro = i == 0; + tarefas[i] = Task.Run( + () => LerSegmentoAsync( + arquivo, inicio, fim - inicio, primeiro, textosIgnorados, cache, cancellationToken), + cancellationToken); + } + + var partes = await Task.WhenAll(tarefas).ConfigureAwait(false); + + var eventos = new List(partes.Sum(p => p.Eventos.Count)); + var invalidas = 0; + string? primeiroErro = null; + foreach (var parte in partes) + { + eventos.AddRange(parte.Eventos); + invalidas += parte.LinhasInvalidas; + primeiroErro ??= parte.PrimeiroErro; + } + + return new ResultadoLeituraArquivoLog(eventos, comprimento, invalidas, primeiroErro); + } + + private static async Task LerSegmentoAsync( + string arquivo, + long inicio, + long comprimento, + bool primeiroSegmento, + IReadOnlyList textosIgnorados, + CacheDeTemplates cache, + CancellationToken cancellationToken) + { + await using var stream = new FileStream( + arquivo, + FileMode.Open, + FileAccess.Read, + FileShare.ReadWrite, + bufferSize: TamanhoBuffer, + useAsync: true); + stream.Position = inicio; + + return await LerStreamAsync( + new RecorteDeLeitura(stream, comprimento), + arquivo, + textosIgnorados, + offsetFinal: null, + cache, + verificarBom: primeiroSegmento, + cancellationToken).ConfigureAwait(false); + } + + /// + /// Limita a leitura a um trecho do stream subjacente. Permite que o laço de + /// leitura sequencial rode inalterado sobre UM segmento do arquivo: para ele, o + /// fim do recorte é indistinguível do fim do arquivo. + /// + private sealed class RecorteDeLeitura : Stream + { + private readonly Stream _origem; + private long _restante; + + public RecorteDeLeitura(Stream origem, long comprimento) + { + _origem = origem; + _restante = comprimento; + } + + public override bool CanRead => true; + public override bool CanSeek => false; + public override bool CanWrite => false; + public override long Length => throw new NotSupportedException(); + public override long Position + { + get => throw new NotSupportedException(); + set => throw new NotSupportedException(); + } + + public override async ValueTask ReadAsync( + Memory destino, CancellationToken cancellationToken = default) + { + if (_restante <= 0) return 0; + var maximo = (int)Math.Min(destino.Length, _restante); + var lidos = await _origem.ReadAsync(destino[..maximo], cancellationToken) + .ConfigureAwait(false); + _restante -= lidos; + return lidos; + } + + public override int Read(byte[] buffer, int offset, int count) + { + if (_restante <= 0) return 0; + var maximo = (int)Math.Min(count, _restante); + var lidos = _origem.Read(buffer, offset, maximo); + _restante -= lidos; + return lidos; + } + + public override void Flush() { } + public override long Seek(long offset, SeekOrigin origin) => throw new NotSupportedException(); + public override void SetLength(long value) => throw new NotSupportedException(); + public override void Write(byte[] buffer, int offset, int count) => throw new NotSupportedException(); + } + + /// Avança a partir de um ponto bruto até logo após o primeiro \n. + private static async Task AcharInicioDeLinhaAsync( + FileStream stream, + long posicaoBruta, + long comprimento, + CancellationToken cancellationToken) + { + stream.Position = posicaoBruta; + var buffer = ArrayPool.Shared.Rent(TamanhoBuffer); + try + { + var posicao = posicaoBruta; + while (posicao < comprimento) + { + var maximo = (int)Math.Min(buffer.Length, comprimento - posicao); + var lidos = await stream.ReadAsync(buffer.AsMemory(0, maximo), cancellationToken) + .ConfigureAwait(false); + if (lidos == 0) break; + + var quebra = buffer.AsSpan(0, lidos).IndexOf((byte)'\n'); + if (quebra >= 0) return posicao + quebra + 1; + posicao += lidos; + } + + return comprimento; + } + finally + { + ArrayPool.Shared.Return(buffer); + } + } + private static async Task LerStreamAsync( Stream stream, string arquivo, IReadOnlyList textosIgnorados, Func? offsetFinal, - PoolDeTextos? pool, + CacheDeTemplates cache, + bool verificarBom, CancellationToken cancellationToken) { - var acumulador = new Acumulador(arquivo, textosIgnorados, CacheDeTemplates.Para(pool)); + var acumulador = new Acumulador(arquivo, textosIgnorados, cache); var buffer = ArrayPool.Shared.Rent(TamanhoBuffer); var preenchido = 0; - var bomVerificado = false; + // Segmentos que não são o primeiro começam no meio do arquivo: três bytes + // EF BB BF ali seriam conteúdo de linha, nunca BOM. + var bomVerificado = !verificarBom; try { diff --git a/src/Services/LeitorClef.cs b/src/Services/LeitorClef.cs index 08c9fd3..df19521 100644 --- a/src/Services/LeitorClef.cs +++ b/src/Services/LeitorClef.cs @@ -99,6 +99,73 @@ internal string Compartilhar(ReadOnlySpan texto) return Compartilhar(new string(texto)); } + + // ── Pool de valores escalares ──────────────────────────────────────────── + // + // A premissa "valor é único por evento" caiu na medição: nos logs reais, os + // valores das propriedades têm 9,9% de cardinalidade — "VAREJO", "PDV OMNI", + // o nome da máquina e o CNPJ se repetem em TODA linha, cada uma criando seu + // próprio ScalarValue com sua própria string. Compartilhar a INSTÂNCIA é seguro + // (ScalarValue é imutável no Serilog) e elimina as duas alocações de uma vez. + + /// true/false/null são três valores no mundo — três objetos no processo. + public static readonly ScalarValue EscalarVerdadeiro = new(true); + public static readonly ScalarValue EscalarFalso = new(false); + public static readonly ScalarValue EscalarNulo = new(null); + + // Contagens pequenas (ProcessorCount, códigos de loja, quantidades) dominam os + // inteiros dos logs reais; o cache evita o boxing E o ScalarValue por linha. + private static readonly ScalarValue[] LongsPequenos = CriarLongsPequenos(); + + private static ScalarValue[] CriarLongsPequenos() + { + var valores = new ScalarValue[1024]; + for (var i = 0; i < valores.Length; i++) valores[i] = new ScalarValue((long)i); + return valores; + } + + // Caps: um log despeja um GUID novo por linha (SpanId) e, sem teto, o pool + // guardaria 315 mil strings que nunca repetem. Os valores QUENTES aparecem nas + // primeiras linhas e entram antes de o teto ser atingido. + private const int MaximoDeEscalares = 64 * 1024; + private const int MaximoDeCandidatos = 64 * 1024; + private const int MaiorTextoCompartilhavel = 128; + + private readonly ConcurrentDictionary _escalares = new(StringComparer.Ordinal); + private readonly ConcurrentDictionary _candidatos = new(StringComparer.Ordinal); + + // Contadores próprios: ConcurrentDictionary.Count ADQUIRE TODOS OS LOCKS da + // tabela — chamado uma vez por valor de cada linha, com 20 workers, serializava + // a leitura paralela inteira (medido: a carga triplicou por causa disso). + private int _totalEscalares; + private int _totalCandidatos; + + public ScalarValue EscalarDe(string? texto) + { + if (texto is null) return EscalarNulo; + if (texto.Length is 0 or > MaiorTextoCompartilhavel) return new ScalarValue(texto); + if (_escalares.TryGetValue(texto, out var existente)) return existente; + + // Promoção na SEGUNDA vista: um log despeja um GUID único por linha (SpanId), + // e inseri-lo direto no pool o encheria de texto que nunca repete. O que não + // reaparece morre na lista de candidatos; só o que repete é promovido. + if (Volatile.Read(ref _totalCandidatos) < MaximoDeCandidatos && _candidatos.TryAdd(texto, true)) + { + Interlocked.Increment(ref _totalCandidatos); + return new ScalarValue(texto); + } + + if (_candidatos.TryRemove(texto, out _) && Volatile.Read(ref _totalEscalares) < MaximoDeEscalares) + { + Interlocked.Increment(ref _totalEscalares); + return _escalares.GetOrAdd(texto, static t => new ScalarValue(t)); + } + + return new ScalarValue(texto); + } + + public static ScalarValue EscalarDeNumero(object bruto) => + bruto is long inteiro and >= 0 and < 1024 ? LongsPequenos[(int)inteiro] : new ScalarValue(bruto); } /// @@ -335,16 +402,14 @@ private static ClefEvent Ler(ReadOnlySpan linha, string arquivo, CacheDeTe : template.Template.Render(new PropriedadesOrdinais(propriedades), CultureInfo.InvariantCulture), Exception = excecao, SourceFile = arquivo, - Properties = new Dictionary( - propriedades?.Count ?? 0, - StringComparer.OrdinalIgnoreCase), + // A forma compacta em vez de Dictionary: 18 pares imutáveis não precisam + // de buckets, e multiplicado por centenas de milhares de eventos o + // dicionário era uma das maiores fatias da memória retida. + Properties = propriedades is null + ? PropriedadesEvento.Vazio + : new PropriedadesEvento(propriedades), }; - if (propriedades is not null) - { - foreach (var propriedade in propriedades) evento.Properties[propriedade.Name] = propriedade.Value; - } - return evento; } @@ -581,17 +646,20 @@ private static LogEventPropertyValue LerValor(ref Utf8JsonReader reader, scoped switch (reader.TokenType) { case JsonTokenType.String: - // O texto do valor é único por evento — compartilhar aqui só encheria o pool. - return new ScalarValue(reader.GetString()); + // Compartilhado: a premissa "valor é único por evento" caiu na medição — + // nos logs reais os valores têm ~10% de cardinalidade ("VAREJO", máquina, + // CNPJ repetem em toda linha). O EscalarDe tem teto para o que realmente + // é único (SpanId) não inchar o pool. + return cache.EscalarDe(reader.GetString()); case JsonTokenType.Number: - return new ScalarValue(LerNumero(ref reader)); + return CacheDeTemplates.EscalarDeNumero(LerNumero(ref reader)); case JsonTokenType.True: - return new ScalarValue(true); + return CacheDeTemplates.EscalarVerdadeiro; case JsonTokenType.False: - return new ScalarValue(false); + return CacheDeTemplates.EscalarFalso; case JsonTokenType.Null: // null NÃO é o mesmo que ausente: o app distingue os dois na descoberta de colunas. - return new ScalarValue(null); + return CacheDeTemplates.EscalarNulo; case JsonTokenType.StartObject: return LerObjeto(ref reader, rascunho, cache); case JsonTokenType.StartArray: diff --git a/src/Services/LogStore.cs b/src/Services/LogStore.cs index f21158c..2038b34 100644 --- a/src/Services/LogStore.cs +++ b/src/Services/LogStore.cs @@ -68,6 +68,13 @@ public LogStore(SettingsService settingsService) public bool IsLoading => Volatile.Read(ref _isLoading) == 1; + /// + /// Intervalo mínimo entre publicações parciais durante a carga. 400 ms equilibra + /// "ver os primeiros eventos cedo" com o custo de cada mesclagem (que copia o + /// conjunto publicado). Os testes zeram para publicar a cada arquivo concluído. + /// + public int PublicacaoParcialMs { get; set; } = 400; + public int Count => Volatile.Read(ref _events).Length; /// @@ -153,12 +160,66 @@ public async Task LoadFromPathsAsync(IEnumerable paths) var intervaloProgresso = Math.Max(1, arquivosCarregar.Length / 20); ReportProgress(operacao, $"Lendo 0/{arquivosCarregar.Length:N0} arquivos…"); + // Publicação incremental: cada arquivo concluído entra numa fila e, no + // máximo a cada PublicacaoParcialMs, o que acumulou é mesclado no conjunto + // publicado. Quem abre uma pasta grande passa a LER os primeiros eventos em + // segundos, com o resto chegando por baixo — antes a tela ficava num + // spinner até o último arquivo terminar. + var lotesPendentes = new ConcurrentQueue>(); + var ultimaPublicacao = Environment.TickCount64; + var publicando = 0; + var metadadosPublicados = false; + + void PublicarParcial() + { + if (Environment.TickCount64 - Volatile.Read(ref ultimaPublicacao) < PublicacaoParcialMs) return; + // Um publicador por vez; quem perder a corrida deixa o lote na fila + // para a próxima rodada — nada se perde, só atrasa um ciclo. + if (Interlocked.Exchange(ref publicando, 1) == 1) return; + try + { + var lote = new List(); + while (lotesPendentes.TryDequeue(out var parte)) lote.AddRange(parte); + if (lote.Count == 0) return; + lote.Sort(PorTimestampDecrescente); + + lock (_stateGate) + { + if (!IsCurrent(operacao)) return; + + // O primeiro lote leva junto os metadados da carga: eventos + // novos ao lado da lista de arquivos da sessão ANTERIOR seria + // um estado que nenhuma outra parte do app sabe interpretar. + if (!metadadosPublicados) + { + metadadosPublicados = true; + _fileName = pathList.Length == 1 ? pathList[0] : "Múltiplos locais"; + Volatile.Write(ref _openedPaths, pathList); + Volatile.Write(ref _availableFiles, descoberta.Arquivos.ToArray()); + Volatile.Write(ref _loadedFiles, arquivosCarregar); + Volatile.Write(ref _events, Array.Empty()); + } + + Volatile.Write(ref _events, MesclarPorRegiao(Volatile.Read(ref _events), lote)); + Interlocked.Increment(ref _stateVersion); + } + + Volatile.Write(ref ultimaPublicacao, Environment.TickCount64); + Changed?.Invoke(); + } + finally + { + Volatile.Write(ref publicando, 0); + } + } + await Parallel.ForEachAsync(arquivosCarregar, opcoes, async (arquivo, token) => { try { var leitura = await _leitor.LerAsync(arquivo, textosIgnorados, pool, token).ConfigureAwait(false); foreach (var evento in leitura.Eventos) eventos.Add(evento); + if (leitura.Eventos.Count > 0) lotesPendentes.Enqueue(leitura.Eventos); if (leitura.OffsetFinal is { } offset) offsets[arquivo] = offset; if (leitura.LinhasInvalidas > 0) @@ -178,6 +239,7 @@ await Parallel.ForEachAsync(arquivosCarregar, opcoes, async (arquivo, token) => } finally { + PublicarParcial(); var processados = Interlocked.Increment(ref arquivosLidos); if (processados == arquivosCarregar.Length || processados % intervaloProgresso == 0) { @@ -324,6 +386,7 @@ await Parallel.ForEachAsync(adicionar, opcoes, async (arquivo, token) => { _fileOffsets.Remove(arquivo); _oversizedTailScanOffsets.Remove(arquivo); + _ultimaAtividadeTail.Remove(arquivo); } foreach (var offset in offsets) _fileOffsets[offset.Key] = offset.Value; @@ -449,6 +512,13 @@ private static string[] Concatenar(string[] atuais, IReadOnlyList novos) private readonly Dictionary _fileOffsets = new(StringComparer.OrdinalIgnoreCase); private readonly Dictionary _oversizedTailScanOffsets = new(StringComparer.OrdinalIgnoreCase); + + // Tick (Environment.TickCount64) da última vez que cada arquivo moveu o offset no + // poll. Enquanto está dentro da janela, o arquivo é aberto TODO tick, ignorando o + // tamanho enumerado — que é stale para arquivos mantidos abertos pelo logger. + private readonly Dictionary _ultimaAtividadeTail = new(StringComparer.OrdinalIgnoreCase); + private int _rodizioTail; + private CancellationTokenSource? _tailCts; private Task? _tailTask; private int _tailEnabled; @@ -456,6 +526,15 @@ private static string[] Concatenar(string[] atuais, IReadOnlyList novos) private static readonly TimeSpan TailInterval = TimeSpan.FromSeconds(1); + // Janela em que um arquivo que acabou de entregar continua sendo aberto todo tick + // (latência de 1 s durante rajadas), e quantos arquivos "frios" cada tick abre no + // rodízio mesmo com a enumeração dizendo "nada mudou". As duas constantes trabalham + // juntas contra o índice stale: o Auditing do PDV escreve a cada ~30 s, então 60 s o + // mantém sempre quente; um arquivo quieto há mais tempo volta a ser visto em até + // total/fatia ticks (104 arquivos ÷ 8 = 13 s) e é promovido a quente na hora. + private const long JanelaAtividadeTailMs = 60_000; + private const int FatiaDoRodizioTail = 8; + // A procura por arquivos novos tem cadência PRÓPRIA porque ela custa uma varredura // recursiva da árvore que o usuário abriu (medida em 29 s numa pasta com 738 mil // entradas), enquanto o poll precisa rodar a cada segundo. Sem ela o acompanhamento @@ -669,8 +748,11 @@ private async Task PollTailAsync(CancellationToken cancellationToken) var novos = new List(); var novosOffsets = new Dictionary(StringComparer.OrdinalIgnoreCase); - foreach (var arquivo in arquivos) + var agora = Environment.TickCount64; + var inicioRodizio = _rodizioTail; + for (var i = 0; i < arquivos.Length; i++) { + var arquivo = arquivos[i]; cancellationToken.ThrowIfCancellationRequested(); if (arquivo.EndsWith(".gz", StringComparison.OrdinalIgnoreCase)) continue; @@ -679,7 +761,23 @@ private async Task PollTailAsync(CancellationToken cancellationToken) // avança sem o arquivo crescer, e um portão por variação deixaria o resto do // arquivo preso para sempre. Tamanho menor que o offset também abre — é // truncamento, e a releitura precisa reposicionar no início. - if (tamanhos.TryGetValue(arquivo, out var tamanho) && tamanho == OffsetLido(arquivo)) continue; + // + // ATENÇÃO: o tamanho enumerado vem do ÍNDICE DO DIRETÓRIO, que o NTFS só + // atualiza de vez em quando enquanto o logger mantém o arquivo aberto — no + // PDV real a enumeração devolvia o tamanho da carga indefinidamente e o modo + // ao vivo morria em silêncio (o harness sintético não pegava: o escritor de + // teste fecha o arquivo a cada linha, o que atualiza o índice). Só o handle + // (FileStream.Length) diz o tamanho verdadeiro. Por isso o portão nunca é a + // única palavra: arquivo com atividade recente é aberto todo tick, e os + // demais passam por um rodízio que abre uma fatia por tick — nenhum arquivo + // fica invisível para sempre atrás do índice stale. + if (tamanhos.TryGetValue(arquivo, out var tamanho) + && tamanho == OffsetLido(arquivo) + && !TeveAtividadeRecente(arquivo, agora) + && !NaFatiaDoRodizio(i, inicioRodizio, arquivos.Length)) + { + continue; + } try { @@ -689,7 +787,11 @@ private async Task PollTailAsync(CancellationToken cancellationToken) pool, cancellationToken).ConfigureAwait(false); novos.AddRange(trecho.Eventos); - if (trecho.NovoOffset is { } offset) novosOffsets[arquivo] = offset; + if (trecho.NovoOffset is { } offset) + { + novosOffsets[arquivo] = offset; + MarcarAtividadeTail(arquivo, agora); + } } catch (OperationCanceledException) { @@ -701,6 +803,11 @@ private async Task PollTailAsync(CancellationToken cancellationToken) } } + if (arquivos.Length > 0) + { + _rodizioTail = (inicioRodizio + FatiaDoRodizioTail) % arquivos.Length; + } + cancellationToken.ThrowIfCancellationRequested(); if (IsLoading || versao != Volatile.Read(ref _stateVersion)) return; novos.Sort(PorTimestampDecrescente); @@ -771,6 +878,31 @@ private long OffsetLido(string arquivo) } } + private bool TeveAtividadeRecente(string arquivo, long agora) + { + lock (_fileOffsets) + { + return _ultimaAtividadeTail.TryGetValue(arquivo, out var ultimo) + && agora - ultimo < JanelaAtividadeTailMs; + } + } + + private void MarcarAtividadeTail(string arquivo, long agora) + { + lock (_fileOffsets) { _ultimaAtividadeTail[arquivo] = agora; } + } + + // Fatia circular [inicio, inicio+fatia) sobre a lista de arquivos; com poucos + // arquivos todos entram todo tick. + private static bool NaFatiaDoRodizio(int indice, int inicio, int total) + { + if (total <= FatiaDoRodizioTail) return true; + var fim = (inicio + FatiaDoRodizioTail) % total; + return inicio <= fim + ? indice >= inicio && indice < fim + : indice >= inicio || indice < fim; + } + private sealed record ResultadoTail(IReadOnlyList Eventos, long? NovoOffset); private async Task ReadNewEventsAsync( @@ -795,6 +927,22 @@ private async Task ReadNewEventsAsync( offset = 0; lock (_fileOffsets) { _oversizedTailScanOffsets.Remove(arquivo); } } + else if (offset > 0 && stream.Length > offset) + { + // Truncado e reescrito MAIOR que o offset antigo — a única rotação que a + // comparação de tamanho não enxerga (pego por prova de ponta a ponta: o + // conteúdo novo era maior que o antigo e o tail lia do meio de uma linha, + // descartando-a como inválida). O offset avança sempre até logo depois de + // um '\n'; se o byte anterior não é '\n', o conteúdo sob o offset não é o + // que foi lido antes. + stream.Seek(offset - 1, SeekOrigin.Begin); + if (stream.ReadByte() != '\n') + { + offset = 0; + stream.Seek(0, SeekOrigin.Begin); + lock (_fileOffsets) { _oversizedTailScanOffsets.Remove(arquivo); } + } + } if (stream.Length == offset) return new ResultadoTail(Array.Empty(), null); var pendente = Math.Min(stream.Length - offset, MaxTailReadBytes); @@ -898,6 +1046,7 @@ private void PublicarOffsets( { _fileOffsets.Clear(); _oversizedTailScanOffsets.Clear(); + _ultimaAtividadeTail.Clear(); foreach (var arquivo in arquivosCarregados) { if (offsets.TryGetValue(arquivo, out var offset)) _fileOffsets[arquivo] = offset; diff --git a/test/ClefExplorer.Tests/LeitorArquivoLogTests.cs b/test/ClefExplorer.Tests/LeitorArquivoLogTests.cs index 80063de..ecb4df5 100644 --- a/test/ClefExplorer.Tests/LeitorArquivoLogTests.cs +++ b/test/ClefExplorer.Tests/LeitorArquivoLogTests.cs @@ -268,4 +268,90 @@ await Parallel.ForEachAsync( Assert.All(totais, t => Assert.Equal(200, t)); } + + // ── Leitura paralela de um arquivo grande ─────────────────────────────────── + // + // Acima do limiar, o arquivo é dividido em segmentos alinhados à quebra de linha e + // cada um é lido num worker. O contrato é EQUIVALÊNCIA TOTAL com o caminho sequencial: + // mesmos eventos, na mesma ordem, mesmas linhas inválidas, mesmo offset. Os testes + // usam limiar/segmento minúsculos para exercitar o caminho com arquivos de KB. + + private string GerarArquivoParaleloTeste(int linhas, bool comBom, bool crlf, bool semQuebraFinal) + { + var caminho = Path.Combine(_pasta, Guid.NewGuid().ToString("N") + ".clef"); + var sb = new StringBuilder(); + for (var i = 0; i < linhas; i++) + { + if (i % 97 == 3) + { + sb.Append("{ isto nao e um evento CLEF }"); // inválida no meio + } + else if (i % 53 == 7) + { + // linha longa: um @x de vários KB atravessa fronteiras de segmento + sb.Append("{\"@t\":\"2026-07-01T00:00:") + .Append((i % 60).ToString("00")) + .Append("Z\",\"@mt\":\"grande {N}\",\"N\":").Append(i) + .Append(",\"@x\":\"").Append(new string('x', 4000)).Append("\"}"); + } + else + { + sb.Append("{\"@t\":\"2026-07-01T00:00:") + .Append((i % 60).ToString("00")) + .Append("Z\",\"@mt\":\"evento {N}\",\"N\":").Append(i).Append('}'); + } + + if (!semQuebraFinal || i < linhas - 1) sb.Append(crlf ? "\r\n" : "\n"); + } + + var corpo = Encoding.UTF8.GetBytes(sb.ToString()); + using var fs = File.Create(caminho); + if (comBom) fs.Write(new byte[] { 0xEF, 0xBB, 0xBF }); + fs.Write(corpo); + return caminho; + } + + private static async Task CompararComSequencialAsync(string caminho) + { + // Limiar alto = nunca paralelo; limiar 1 = sempre paralelo, com segmentos de 2 KB. + var sequencial = await new LeitorArquivoLog(limiarParalelo: long.MaxValue, segmentoMinimo: 1) + .LerAsync(caminho, Array.Empty()); + var paralelo = await new LeitorArquivoLog(limiarParalelo: 1, segmentoMinimo: 2048) + .LerAsync(caminho, Array.Empty()); + + Assert.Equal(sequencial.Eventos.Count, paralelo.Eventos.Count); + for (var i = 0; i < sequencial.Eventos.Count; i++) + { + Assert.Equal(sequencial.Eventos[i].Message, paralelo.Eventos[i].Message); + Assert.Equal(sequencial.Eventos[i].Timestamp, paralelo.Eventos[i].Timestamp); + } + Assert.Equal(sequencial.LinhasInvalidas, paralelo.LinhasInvalidas); + Assert.Equal(sequencial.PrimeiroErro, paralelo.PrimeiroErro); + Assert.Equal(sequencial.OffsetFinal, paralelo.OffsetFinal); + } + + [Fact] + public async Task Parallel_reading_is_indistinguishable_from_sequential() + { + var caminho = GerarArquivoParaleloTeste(2_000, comBom: false, crlf: false, semQuebraFinal: false); + await CompararComSequencialAsync(caminho); + } + + [Fact] + public async Task Parallel_reading_handles_bom_crlf_and_a_missing_final_newline() + { + // BOM só existe no segmento 0; CRLF e a última linha sem quebra são dos workers + // das pontas — os três juntos cobrem as costuras entre segmentos. + var caminho = GerarArquivoParaleloTeste(1_500, comBom: true, crlf: true, semQuebraFinal: true); + await CompararComSequencialAsync(caminho); + } + + [Fact] + public async Task A_line_longer_than_a_whole_segment_still_comes_out_once() + { + // Segmento de 2 KB com linhas de 4 KB: toda fronteira bruta cai DENTRO de alguma + // linha, o que força o alinhamento a atravessar segmentos inteiros. + var caminho = GerarArquivoParaleloTeste(300, comBom: false, crlf: false, semQuebraFinal: false); + await CompararComSequencialAsync(caminho); + } } diff --git a/test/ClefExplorer.Tests/LogColumnDiscoveryTests.cs b/test/ClefExplorer.Tests/LogColumnDiscoveryTests.cs index 21b33d3..830f000 100644 --- a/test/ClefExplorer.Tests/LogColumnDiscoveryTests.cs +++ b/test/ClefExplorer.Tests/LogColumnDiscoveryTests.cs @@ -13,19 +13,18 @@ public class LogColumnDiscoveryTests { private static ClefEvent Event(params (string Key, object Value)[] props) { - var ev = new ClefEvent - { - Level = "Information", - Timestamp = DateTimeOffset.UtcNow, - Properties = new Dictionary(StringComparer.OrdinalIgnoreCase), - }; - + var propriedades = new Dictionary(StringComparer.OrdinalIgnoreCase); foreach (var (key, value) in props) { - ev.Properties![key] = new ScalarValue(value); + propriedades[key] = new ScalarValue(value); } - return ev; + return new ClefEvent + { + Level = "Information", + Timestamp = DateTimeOffset.UtcNow, + Properties = propriedades, + }; } // --- Descoberta -------------------------------------------------------------- diff --git a/test/ClefExplorer.Tests/LogStoreTests.cs b/test/ClefExplorer.Tests/LogStoreTests.cs index 79a80e9..0ae75ca 100644 --- a/test/ClefExplorer.Tests/LogStoreTests.cs +++ b/test/ClefExplorer.Tests/LogStoreTests.cs @@ -634,10 +634,11 @@ public void CancelLoad_without_a_load_in_progress_is_harmless() } [Fact] - public async Task Cancelling_a_load_preserves_the_previous_content() + public async Task Cancelling_a_load_never_mixes_the_two_sessions() { - // O estado só é trocado ao final de uma leitura completa, então um - // carregamento abortado não pode deixar a lista pela metade. + // Com a publicação incremental, cancelar no meio PODE deixar um parcial da carga + // nova — é o que o usuário estava lendo quando cancelou. O que continua proibido + // é o meio-termo entre SESSÕES: eventos antigos e novos no mesmo conjunto. WriteClef("pequeno.clef", ClefLine("conteúdo anterior")); var store = NewStore(); await store.LoadFromFolderAsync(_root); @@ -650,11 +651,75 @@ public async Task Cancelling_a_load_preserves_the_previous_content() store.CancelLoad(); await carregando; - // Ou o conteúdo anterior (cancelou a tempo) ou o novo completo — nunca um meio-termo. - Assert.True(store.Count == 1 || store.Count == 40_000, $"estado inconsistente: {store.Count}"); + var snapshot = store.Snapshot(); + if (snapshot.Length == 1 && snapshot[0].Message == "conteúdo anterior") + { + // Cancelou antes do primeiro lote parcial: sessão anterior intacta. + } + else + { + // Publicou parcial (ou completou): tudo tem que ser da carga NOVA. + Assert.DoesNotContain(snapshot, e => e.Message == "conteúdo anterior"); + } + Assert.False(store.IsLoading); + } + + [Fact] + public async Task A_load_publishes_events_before_finishing() + { + // O ponto da carga incremental: com um arquivo entregue e outro ainda em leitura, + // o que já chegou aparece — antes, a tela ficava vazia até o último arquivo. + var rapido = WriteClef("rapido.clef", ClefLine("do arquivo rápido")); + var lento = WriteClef("lento.clef", ClefLine("do arquivo lento")); + + // No fake, o "segundo" arquivo fica preso até LiberarSegundo; o rápido entrega na + // hora (e vira um evento cujo Message é o NOME do arquivo). + var leitor = new LeitorControlado( + Path.Combine(_root, "nunca-usado.clef"), + lento); + var store = NewStore(leitor); + store.PublicacaoParcialMs = 0; // publica a cada arquivo concluído + + var carga = store.LoadFromFolderAsync(_root); + await leitor.SegundoIniciado.Task.WaitAsync(TimeSpan.FromSeconds(10)); + + // Espera o lote parcial do arquivo rápido ser publicado. + var limite = DateTime.UtcNow.AddSeconds(10); + while (store.Count == 0 && DateTime.UtcNow < limite) + { + await Task.Delay(25); + } + + Assert.True(store.IsLoading, "a carga já tinha terminado — o parcial não foi observável"); + Assert.Equal(1, store.Count); + Assert.Equal("rapido.clef", store.Snapshot()[0].Message); + // Os metadados acompanham o primeiro lote: a lista de arquivos é a da carga NOVA. + Assert.Contains(store.LoadedFiles, f => f.EndsWith("rapido.clef", StringComparison.OrdinalIgnoreCase)); + GC.KeepAlive(rapido); + + leitor.LiberarSegundo.SetResult(); + await carga; + + Assert.Equal(2, store.Count); Assert.False(store.IsLoading); } + [Fact] + public async Task Partial_batches_do_not_duplicate_events_in_the_final_set() + { + // A publicação final substitui o conjunto pelo total ordenado; um evento que + // apareceu num lote parcial não pode aparecer duas vezes ao terminar. + WriteClef("a.clef", ClefLine("um"), ClefLine("dois")); + WriteClef("b.clef", ClefLine("três")); + var store = NewStore(); + store.PublicacaoParcialMs = 0; + + await store.LoadFromFolderAsync(_root); + + Assert.Equal(3, store.Count); + Assert.Equal(3, store.Snapshot().Select(e => e.Message).Distinct().Count()); + } + [Fact] public async Task A_new_load_supersedes_the_one_in_progress() { @@ -750,6 +815,45 @@ public async Task Tail_picks_up_lines_appended_after_the_load() } } + [Fact] + public async Task Tail_picks_up_lines_from_a_writer_that_keeps_the_file_open() + { + // O cenário do PDV real: o logger abre o arquivo UMA vez e nunca fecha. Enquanto o + // handle está aberto, o tamanho no índice do diretório (o que a enumeração do poll + // enxerga) pode ficar congelado no valor antigo — só o handle diz o tamanho real. + // O gate do poll não pode confiar na enumeração como prova de "nada mudou". + var file = WriteClef("app.clef", ClefLine("inicial")); + var store = NewStore(); + await store.LoadFromFolderAsync(_root); + Assert.Equal(1, store.Count); + + await using var escritor = new FileStream( + file, FileMode.Append, FileAccess.Write, FileShare.ReadWrite); + + store.SetTailEnabled(true); + try + { + var linha = Encoding.UTF8.GetBytes(ClefLine("do handle aberto") + "\n"); + await escritor.WriteAsync(linha); + await escritor.FlushAsync(); + + Assert.True(await WaitUntil(() => store.Count == 2), "o evento do escritor de handle aberto não foi captado"); + Assert.Contains(store.Snapshot(), e => e.Message == "do handle aberto"); + + // Segunda rajada com o MESMO handle: o arquivo agora está "quente" e precisa + // continuar sendo lido tick a tick. + var linha2 = Encoding.UTF8.GetBytes(ClefLine("segunda rajada") + "\n"); + await escritor.WriteAsync(linha2); + await escritor.FlushAsync(); + + Assert.True(await WaitUntil(() => store.Count == 3), "a segunda rajada não foi captada"); + } + finally + { + store.SetTailEnabled(false); + } + } + [Fact] public async Task Tail_does_not_duplicate_lines_written_during_the_initial_load() { @@ -927,6 +1031,40 @@ public async Task Tail_restarts_from_the_beginning_when_the_file_is_truncated() } } + [Fact] + public async Task Tail_detects_a_file_rewritten_LARGER_than_the_old_offset() + { + // O caso que a comparação de tamanho não enxerga — pego por prova de ponta a + // ponta, com os testes acima verdes: se a rotação reescreve o arquivo com MAIS + // bytes que o offset antigo, `Length < offset` é falso e o tail lia do meio de + // uma linha, descartando-a como inválida. A detecção agora usa o invariante de + // que o offset sempre para logo depois de um '\n'. + var file = WriteClef("app.clef", ClefLine("curto")); + var store = NewStore(); + await store.LoadFromFolderAsync(_root); + Assert.Equal(1, store.Count); + + store.SetTailEnabled(true); + try + { + // Reescrito de uma vez, com conteúdo maior que o arquivo original inteiro. + File.WriteAllText( + file, + ClefLine("primeira linha bem mais longa que o conteúdo antigo do arquivo") + Environment.NewLine + + ClefLine("segunda linha depois da rotação") + Environment.NewLine); + + Assert.True( + await WaitUntil(() => store.Snapshot().Any( + e => e.Message == "primeira linha bem mais longa que o conteúdo antigo do arquivo")), + "a primeira linha do arquivo reescrito foi perdida"); + Assert.Contains(store.Snapshot(), e => e.Message == "segunda linha depois da rotação"); + } + finally + { + store.SetTailEnabled(false); + } + } + [Fact] public async Task Tail_descarta_o_BOM_ao_reler_do_inicio_o_arquivo_truncado() { diff --git a/test/ClefExplorer.Tests/PropriedadesEventoTests.cs b/test/ClefExplorer.Tests/PropriedadesEventoTests.cs new file mode 100644 index 0000000..e3b2e07 --- /dev/null +++ b/test/ClefExplorer.Tests/PropriedadesEventoTests.cs @@ -0,0 +1,114 @@ +using ClefExplorer.Models; +using ClefExplorer.Services; +using Serilog.Events; + +namespace ClefExplorer.Tests; + +/// +/// Contrato do — a forma compacta que substituiu o +/// Dictionary por evento. O que importa: comportamento IDÊNTICO ao dicionário +/// OrdinalIgnoreCase que ficava ali, com dois arrays no lugar de buckets. +/// +public class PropriedadesEventoTests +{ + private static LogEventProperty Prop(string nome, object? valor) => + new(nome, new ScalarValue(valor)); + + [Fact] + public void Lookup_ignores_case_like_the_dictionary_it_replaced() + { + var props = new PropriedadesEvento(new[] { Prop("SourceContext", "Api") }); + + Assert.True(props.TryGetValue("sourcecontext", out var valor)); + Assert.Equal("Api", Assert.IsType(valor).Value); + Assert.True(props.ContainsKey("SOURCECONTEXT")); + } + + [Fact] + public void A_repeated_key_keeps_the_last_value() + { + // Semântica do Dictionary[k] = v em laço: o último vence — inclusive quando a + // repetição só existe ignorando maiúsculas. + var props = new PropriedadesEvento(new[] + { + Prop("Chave", 1), + Prop("chave", 2), + }); + + Assert.Single(props); + Assert.Equal(2, Assert.IsType(props["Chave"]).Value); + } + + [Fact] + public void Enumeration_preserves_the_file_order() + { + var props = new PropriedadesEvento(new[] { Prop("B", 1), Prop("A", 2), Prop("C", 3) }); + + Assert.Equal(new[] { "B", "A", "C" }, props.Keys); + Assert.Equal(3, props.Count); + } + + [Fact] + public void A_missing_key_behaves_like_the_dictionary() + { + var props = new PropriedadesEvento(new[] { Prop("A", 1) }); + + Assert.False(props.TryGetValue("Z", out _)); + Assert.Throws(() => props["Z"]); + } + + [Fact] + public void Vazio_is_a_single_shared_instance() + { + Assert.Same(PropriedadesEvento.Vazio, PropriedadesEvento.Vazio); + Assert.Empty(PropriedadesEvento.Vazio); + } + + // ── Pool de escalares (CacheDeTemplates) ──────────────────────────────────── + + [Fact] + public void True_false_and_null_are_process_wide_singletons() + { + Assert.Same(CacheDeTemplates.EscalarVerdadeiro, CacheDeTemplates.EscalarVerdadeiro); + Assert.Equal(true, CacheDeTemplates.EscalarVerdadeiro.Value); + Assert.Equal(false, CacheDeTemplates.EscalarFalso.Value); + Assert.Null(CacheDeTemplates.EscalarNulo.Value); + } + + [Fact] + public void Small_longs_share_one_instance_and_big_ones_do_not() + { + Assert.Same(CacheDeTemplates.EscalarDeNumero(20L), CacheDeTemplates.EscalarDeNumero(20L)); + Assert.NotSame(CacheDeTemplates.EscalarDeNumero(1_000_000L), CacheDeTemplates.EscalarDeNumero(1_000_000L)); + // O VALOR continua exato nos dois casos. + Assert.Equal(1_000_000L, CacheDeTemplates.EscalarDeNumero(1_000_000L).Value); + } + + [Fact] + public void Repeated_string_values_share_one_instance_from_the_second_sighting_on() + { + // "VAREJO" aparece em toda linha dos logs reais. A promoção é na SEGUNDA vista: + // a primeira ocorrência fica avulsa de propósito — inserir tudo no pool encheria + // ele de GUIDs que nunca repetem e criava contenção entre os workers da carga. + var cache = new CacheDeTemplates(new PoolDeTextos()); + + var primeira = cache.EscalarDe("VAREJO"); + var segunda = cache.EscalarDe("VAREJO"); + var terceira = cache.EscalarDe("VAREJO"); + + Assert.NotSame(primeira, segunda); + Assert.Same(segunda, terceira); + Assert.Equal("VAREJO", terceira.Value); + } + + [Fact] + public void A_huge_string_value_is_not_pooled() + { + // Stack traces e payloads não entram: o teto de tamanho protege o pool do que + // quase nunca repete. + var cache = new CacheDeTemplates(new PoolDeTextos()); + var grande = new string('x', 4_000); + + Assert.NotSame(cache.EscalarDe(grande), cache.EscalarDe(grande)); + } +}