From 7e6fd441ebb70166137828c3fbceecc8edabcbb8 Mon Sep 17 00:00:00 2001 From: Echways Date: Sun, 16 Aug 2026 22:52:17 +0300 Subject: [PATCH 1/2] fix dlq logic --- docker-compose.yml | 43 ++-- .../grafana/provisioning/alerting/rules.yml | 92 +++++++- monitoring/prometheus.yml | 4 + src/LinkTracker.AiAgent.Api/Program.cs | 7 +- src/LinkTracker.AiAgent.Api/appsettings.json | 12 ++ .../Telemetry/Abstractions/IAiAgentMetrics.cs | 14 ++ .../Clients/Kafka/RawUpdatesKafkaConsumer.cs | 15 ++ .../RawUpdatesKafkaDeadLetterPublisher.cs | 48 +++++ .../Kafka/RawUpdatesKafkaMessageHandler.cs | 124 ++++++++++- .../Clients/Registration/ClientsModule.cs | 6 +- .../Kafka/RawUpdatesKafkaOptions.cs | 5 +- .../IRawUpdateDeadLetterPublisher.cs | 12 ++ .../Kafka/RawUpdatesDeadLetterKafkaMessage.cs | 18 ++ .../Telemetry/AiAgentMetrics.cs | 100 +++++++++ .../Telemetry/Registration/TelemetryModule.cs | 22 ++ .../appsettings.Docker.json | 2 +- src/LinkTracker.Bot.Api/appsettings.json | 6 +- .../Telemetry/Abstractions/IBotMetrics.cs | 4 + .../Kafka/LinkUpdatesKafkaMessageHandler.cs | 13 +- .../Kafka/LinkUpdatesKafkaOptions.cs | 2 +- .../Telemetry/BotMetrics.cs | 24 +++ src/LinkTracker.Scrapper.Api/appsettings.json | 6 +- ...RawUpdatesKafkaConsumerIntegrationTests.cs | 93 ++++++++- .../RawUpdatesKafkaMessageHandlerTests.cs | 197 ++++++++++++++++++ ...inkUpdatesKafkaConsumerIntegrationTests.cs | 2 + .../LinkUpdatesKafkaMessageHandlerTests.cs | 62 +++++- 26 files changed, 876 insertions(+), 57 deletions(-) create mode 100644 src/LinkTracker.AiAgent.Application/Telemetry/Abstractions/IAiAgentMetrics.cs create mode 100644 src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/RawUpdatesKafkaDeadLetterPublisher.cs create mode 100644 src/LinkTracker.AiAgent.Infrastructure/Kafka/Abstractions/IRawUpdateDeadLetterPublisher.cs create mode 100644 src/LinkTracker.AiAgent.Infrastructure/Models/Kafka/RawUpdatesDeadLetterKafkaMessage.cs create mode 100644 src/LinkTracker.AiAgent.Infrastructure/Telemetry/AiAgentMetrics.cs create mode 100644 src/LinkTracker.AiAgent.Infrastructure/Telemetry/Registration/TelemetryModule.cs create mode 100644 src/LinkTracker.Tests/AiAgent/Unit/Infrastructure/Clients/Kafka/RawUpdatesKafkaMessageHandlerTests.cs diff --git a/docker-compose.yml b/docker-compose.yml index 8801d60..b21a14c 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -103,33 +103,24 @@ services: - /bin/sh - -c - | - /opt/kafka/bin/kafka-topics.sh \ - --bootstrap-server kafka-1:9092 \ - --create \ - --if-not-exists \ - --topic link.processed-updates \ - --partitions 6 \ - --replication-factor 3 \ - --config min.insync.replicas=2 - - /opt/kafka/bin/kafka-topics.sh \ - --bootstrap-server kafka-1:9092 \ - --create \ - --if-not-exists \ - --topic link.raw-updates \ - --partitions 6 \ - --replication-factor 3 \ - --config min.insync.replicas=2 - - /opt/kafka/bin/kafka-topics.sh \ - --bootstrap-server kafka-1:9092 \ - --create \ - --if-not-exists \ - --topic link.raw-updates-dlq \ - --partitions 6 \ - --replication-factor 3 \ - --config min.insync.replicas=2 + set -e + for topic in \ + link.raw-updates \ + link.raw-updates-dlq \ + link.processed-updates \ + link.processed-updates-dlq + do + /opt/kafka/bin/kafka-topics.sh \ + --bootstrap-server kafka-1:9092 \ + --create \ + --if-not-exists \ + --topic "$$topic" \ + --partitions 6 \ + --replication-factor 3 \ + --config min.insync.replicas=2 + done + schema-registry: image: confluentinc/cp-schema-registry:7.6.1 container_name: linktracker-schema-registry diff --git a/monitoring/grafana/provisioning/alerting/rules.yml b/monitoring/grafana/provisioning/alerting/rules.yml index 2d4fe19..e7123d2 100644 --- a/monitoring/grafana/provisioning/alerting/rules.yml +++ b/monitoring/grafana/provisioning/alerting/rules.yml @@ -26,7 +26,7 @@ groups: model: refId: A editorMode: code - expr: max by (job, instance) (process_memory_working_set_bytes{job=~"scrapper|bot"}) + expr: max by (job, instance) (process_memory_working_set_bytes{job=~"scrapper|bot|aiagent"}) instant: true intervalMs: 1000 maxDataPoints: 43200 @@ -44,4 +44,92 @@ groups: evaluator: type: gt params: - - 524288000 \ No newline at end of file + - 524288000 + + - orgId: 1 + name: kafka + folder: LinkTracker + interval: 1m + rules: + - uid: kafka_dlq_publish_failed + title: Kafka DLQ is unavailable + condition: C + for: 2m + noDataState: OK + execErrState: Error + labels: + severity: critical + annotations: + summary: "{{ $labels.job }} не может писать в DLQ (topic {{ $labels.topic }})" + description: >- + kafka_dead_letter_errors_total растет: отправка в DLQ падает, offset не коммитится, + и «ядовитое» сообщение переигрывается бесконечно. Проверь, что DLQ-топик существует + (kafka-init) и что DeadLetterTopic в конфиге сервиса совпадает с ним. + data: + - refId: A + relativeTimeRange: + from: 600 + to: 0 + datasourceUid: prometheus + model: + refId: A + editorMode: code + expr: sum by (job, topic) (increase(kafka_dead_letter_errors_total{job=~"bot|aiagent"}[5m])) + instant: true + intervalMs: 1000 + maxDataPoints: 43200 + - refId: C + relativeTimeRange: + from: 600 + to: 0 + datasourceUid: __expr__ + model: + refId: C + type: threshold + expression: A + conditions: + - type: query + evaluator: + type: gt + params: + - 0 + + - uid: kafka_dlq_growth + title: Kafka DLQ is filling up + condition: C + for: 5m + noDataState: OK + execErrState: Error + labels: + severity: warning + annotations: + summary: "Сообщения уходят в DLQ у {{ $labels.job }} (topic {{ $labels.topic }})" + description: "kafka_dead_letter_total растет: сообщения не проходят десериализацию, валидацию или обработку." + data: + - refId: A + relativeTimeRange: + from: 600 + to: 0 + datasourceUid: prometheus + model: + refId: A + editorMode: code + expr: sum by (job, topic) (increase(kafka_dead_letter_total{job=~"bot|aiagent"}[5m])) + instant: true + intervalMs: 1000 + maxDataPoints: 43200 + - refId: C + relativeTimeRange: + from: 600 + to: 0 + datasourceUid: __expr__ + model: + refId: C + type: threshold + expression: A + conditions: + - type: query + evaluator: + type: gt + params: + - 0 \ No newline at end of file diff --git a/monitoring/prometheus.yml b/monitoring/prometheus.yml index d40a5be..ffe67cb 100644 --- a/monitoring/prometheus.yml +++ b/monitoring/prometheus.yml @@ -9,3 +9,7 @@ scrape_configs: - job_name: bot static_configs: - targets: ["bot:8011"] + + - job_name: aiagent + static_configs: + - targets: ["aiagent:8102"] diff --git a/src/LinkTracker.AiAgent.Api/Program.cs b/src/LinkTracker.AiAgent.Api/Program.cs index 2f2943d..4368cf6 100644 --- a/src/LinkTracker.AiAgent.Api/Program.cs +++ b/src/LinkTracker.AiAgent.Api/Program.cs @@ -1,6 +1,8 @@ using LinkTracker.AiAgent.Application.Registration; using LinkTracker.AiAgent.Infrastructure.Clients.Registration; +using LinkTracker.AiAgent.Infrastructure.Telemetry.Registration; using LinkTracker.EnvReader; +using LinkTracker.Shared.Infrastructure.Telemetry; var builder = WebApplication.CreateBuilder(args); @@ -8,9 +10,12 @@ builder.Services.AddAiAgentApplication(); builder.Services.AddAiAgentInfrastructure(builder.Configuration); +builder.Services.AddTelemetry(builder.Configuration); var app = builder.Build(); app.MapGet("/health", () => Results.Ok()); -await app.RunAsync(); \ No newline at end of file +app.MapMetricsEndpoint().RequireHost("*:8102"); + +await app.RunAsync(); diff --git a/src/LinkTracker.AiAgent.Api/appsettings.json b/src/LinkTracker.AiAgent.Api/appsettings.json index 104c726..6bd95d6 100644 --- a/src/LinkTracker.AiAgent.Api/appsettings.json +++ b/src/LinkTracker.AiAgent.Api/appsettings.json @@ -6,6 +6,18 @@ } }, "AllowedHosts": "*", + "Kestrel": { + "Endpoints": { + "Rest": { + "Url": "http://0.0.0.0:8101", + "Protocols": "Http1" + }, + "Metrics": { + "Url": "http://0.0.0.0:8102", + "Protocols": "Http1" + } + } + }, "Kafka": { "Consumer": { "BootstrapServers": "localhost:9094,localhost:9095,localhost:9096", diff --git a/src/LinkTracker.AiAgent.Application/Telemetry/Abstractions/IAiAgentMetrics.cs b/src/LinkTracker.AiAgent.Application/Telemetry/Abstractions/IAiAgentMetrics.cs new file mode 100644 index 0000000..a458c45 --- /dev/null +++ b/src/LinkTracker.AiAgent.Application/Telemetry/Abstractions/IAiAgentMetrics.cs @@ -0,0 +1,14 @@ +namespace LinkTracker.AiAgent.Application.Telemetry.Abstractions; + +public interface IAiAgentMetrics +{ + void IncrementKafkaConsumed(string topic); + + void IncrementKafkaConsumeError(string topic); + + void ObserveKafkaConsumeDuration(string topic, double milliseconds); + + void IncrementKafkaDeadLetter(string topic); + + void IncrementKafkaDeadLetterError(string topic); +} diff --git a/src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/RawUpdatesKafkaConsumer.cs b/src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/RawUpdatesKafkaConsumer.cs index e8e43e4..fdde3df 100644 --- a/src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/RawUpdatesKafkaConsumer.cs +++ b/src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/RawUpdatesKafkaConsumer.cs @@ -1,4 +1,6 @@ +using System.Diagnostics; using Confluent.Kafka; +using LinkTracker.AiAgent.Application.Telemetry.Abstractions; using LinkTracker.AiAgent.Infrastructure.Configuration.Kafka; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; @@ -10,6 +12,7 @@ internal sealed class RawUpdatesKafkaConsumer( IConsumer consumer, IRawUpdatesKafkaMessageHandler messageHandler, IOptions kafkaOptions, + IAiAgentMetrics metrics, ILogger logger) : BackgroundService { protected override Task ExecuteAsync(CancellationToken stoppingToken) @@ -63,10 +66,17 @@ private async Task ConsumeLoopAsync(CancellationToken stoppingToken) private async Task ProcessMessageAsync(ConsumeResult result, CancellationToken ct) { + var sw = Stopwatch.StartNew(); + try { var shouldCommit = await messageHandler.HandleAsync(result, ct); + sw.Stop(); + + metrics.IncrementKafkaConsumed(result.Topic); + metrics.ObserveKafkaConsumeDuration(result.Topic, sw.Elapsed.TotalMilliseconds); + if (shouldCommit) { TryCommit(result); @@ -78,6 +88,11 @@ private async Task ProcessMessageAsync(ConsumeResult result, Can } catch (Exception ex) { + sw.Stop(); + + metrics.IncrementKafkaConsumeError(result.Topic); + metrics.ObserveKafkaConsumeDuration(result.Topic, sw.Elapsed.TotalMilliseconds); + logger.LogError( ex, "Ошибка обработки Kafka сообщения. Offset не будет подтвержден. Topic={Topic}, Partition={Partition}, Offset={Offset}", diff --git a/src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/RawUpdatesKafkaDeadLetterPublisher.cs b/src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/RawUpdatesKafkaDeadLetterPublisher.cs new file mode 100644 index 0000000..63ffce2 --- /dev/null +++ b/src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/RawUpdatesKafkaDeadLetterPublisher.cs @@ -0,0 +1,48 @@ +using System.Text.Json; +using Confluent.Kafka; +using LinkTracker.AiAgent.Infrastructure.Configuration.Kafka; +using LinkTracker.AiAgent.Infrastructure.Kafka.Abstractions; +using LinkTracker.AiAgent.Infrastructure.Models.Kafka; +using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Options; + +namespace LinkTracker.AiAgent.Infrastructure.Clients.Kafka; + +internal sealed class RawUpdatesKafkaDeadLetterPublisher( + IProducer producer, + IOptions options, + ILogger logger) : IRawUpdateDeadLetterPublisher +{ + private static readonly JsonSerializerOptions JsonSerializerOptions = new(JsonSerializerDefaults.Web); + + public async Task PublishAsync( + ConsumeResult sourceMessage, + string reason, + Exception? exception, + CancellationToken ct) + { + var deadLetterMessage = new RawUpdatesDeadLetterKafkaMessage + { + Payload = Convert.ToBase64String(sourceMessage.Message.Value), + Reason = reason, + ExceptionType = exception?.GetType().FullName, + SourceTopic = sourceMessage.Topic, + SourcePartition = sourceMessage.Partition.Value, + SourceOffset = sourceMessage.Offset.Value + }; + + var payload = JsonSerializer.SerializeToUtf8Bytes(deadLetterMessage, JsonSerializerOptions); + + var result = await producer.ProduceAsync( + options.Value.DeadLetterTopic, + new Message { Key = sourceMessage.Message.Key, Value = payload }, + ct); + + logger.LogWarning( + "Kafka сообщение отправлено в DLQ. Topic={Topic}, Partition={Partition}, Offset={Offset}, Reason={Reason}", + result.Topic, + result.Partition.Value, + result.Offset.Value, + reason); + } +} diff --git a/src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/RawUpdatesKafkaMessageHandler.cs b/src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/RawUpdatesKafkaMessageHandler.cs index 7307004..27a469e 100644 --- a/src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/RawUpdatesKafkaMessageHandler.cs +++ b/src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/RawUpdatesKafkaMessageHandler.cs @@ -1,31 +1,130 @@ using Confluent.Kafka; using LinkTracker.AiAgent.Application.Abstractions; +using LinkTracker.AiAgent.Application.Telemetry.Abstractions; +using LinkTracker.AiAgent.Infrastructure.Configuration.Kafka; using LinkTracker.AiAgent.Infrastructure.Kafka.Abstractions; +using LinkTracker.Shared.Contracts.Bot; using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Options; namespace LinkTracker.AiAgent.Infrastructure.Clients.Kafka; internal sealed class RawUpdatesKafkaMessageHandler( IRawLinkUpdateKafkaDeserializer deserializer, ILinkUpdateProcessingService processingService, + IRawUpdateDeadLetterPublisher deadLetterPublisher, + IOptions kafkaOptions, + IAiAgentMetrics metrics, ILogger logger) : IRawUpdatesKafkaMessageHandler { public async Task HandleAsync(ConsumeResult result, CancellationToken ct) { + LinkUpdate? update; + try { - var update = await deserializer.DeserializeAsync(result.Message.Value, result.Topic, ct); + update = await deserializer.DeserializeAsync(result.Message.Value, result.Topic, ct); + } + catch (OperationCanceledException) when (ct.IsCancellationRequested) + { + throw; + } + catch (Exception ex) + { + return await TryPublishToDeadLetterAsync( + result, + $"Kafka сообщение не удалось десериализовать: {ex.Message}", + ex, + ct); + } - if (update is null) + if (update is null) + { + return await TryPublishToDeadLetterAsync( + result, + "Kafka сообщение десериализовалось в null.", + null, + ct); + } + + var processingError = await TryProcessWithRetriesAsync(update, ct); + + if (processingError is not null) + { + return await TryPublishToDeadLetterAsync( + result, + "Исчерпаны попытки обработки Kafka сообщения.", + processingError, + ct); + } + + logger.LogInformation( + "Kafka сообщение обработано. Topic={Topic}, Partition={Partition}, Offset={Offset}, UpdateId={UpdateId}", + result.Topic, + result.Partition.Value, + result.Offset.Value, + update.Id); + + return true; + } + + private async Task TryProcessWithRetriesAsync(LinkUpdate update, CancellationToken ct) + { + var attempts = Math.Max(1, kafkaOptions.Value.RetryAttempts); + var backoff = TimeSpan.FromMilliseconds(Math.Max(0, kafkaOptions.Value.RetryBackoffMilliseconds)); + + for (var attempt = 1; attempt <= attempts; attempt++) + { + try + { + await processingService.ProcessAsync(update, ct); + return null; + } + catch (OperationCanceledException) when (ct.IsCancellationRequested) + { + throw; + } + catch (Exception ex) when (attempt < attempts) + { + logger.LogWarning( + ex, + "Ошибка обработки Kafka сообщения. Будет повторная попытка. Attempt={Attempt}, MaxAttempts={MaxAttempts}, UpdateId={UpdateId}", + attempt, + attempts, + update.Id); + + if (backoff > TimeSpan.Zero) + { + await Task.Delay(backoff, ct); + } + } + catch (Exception ex) { logger.LogWarning( - "Kafka сообщение десериализовалось в null. Topic={Topic}, Offset={Offset}", - result.Topic, result.Offset.Value); + ex, + "Ошибка обработки Kafka сообщения. Повторные попытки закончились. Attempts={Attempts}, UpdateId={UpdateId}", + attempts, + update.Id); - return true; + return ex; } + } + + return null; + } + + private async Task TryPublishToDeadLetterAsync( + ConsumeResult result, + string reason, + Exception? exception, + CancellationToken ct) + { + try + { + await deadLetterPublisher.PublishAsync(result, reason, exception, ct); + + metrics.IncrementKafkaDeadLetter(result.Topic); - await processingService.ProcessAsync(update, ct); return true; } catch (OperationCanceledException) when (ct.IsCancellationRequested) @@ -34,12 +133,17 @@ public async Task HandleAsync(ConsumeResult result, Cancel } catch (Exception ex) { + metrics.IncrementKafkaDeadLetterError(result.Topic); + logger.LogError( ex, - "Ошибка обработки Kafka сообщения. Topic={Topic}, Partition={Partition}, Offset={Offset}", - result.Topic, result.Partition.Value, result.Offset.Value); + "Не удалось отправить Kafka сообщение в DLQ. Offset не будет подтвержден, сообщение будет переигрываться. Topic={Topic}, Partition={Partition}, Offset={Offset}, DeadLetterTopic={DeadLetterTopic}", + result.Topic, + result.Partition.Value, + result.Offset.Value, + kafkaOptions.Value.DeadLetterTopic); - return true; + return false; } } -} \ No newline at end of file +} diff --git a/src/LinkTracker.AiAgent.Infrastructure/Clients/Registration/ClientsModule.cs b/src/LinkTracker.AiAgent.Infrastructure/Clients/Registration/ClientsModule.cs index 3f937db..bc6a3d4 100644 --- a/src/LinkTracker.AiAgent.Infrastructure/Clients/Registration/ClientsModule.cs +++ b/src/LinkTracker.AiAgent.Infrastructure/Clients/Registration/ClientsModule.cs @@ -37,6 +37,9 @@ public static IServiceCollection AddAiAgentInfrastructure( .Validate(o => !string.IsNullOrWhiteSpace(o.BootstrapServers), "Kafka:Consumer:BootstrapServers must be set") .Validate(o => !string.IsNullOrWhiteSpace(o.Topic), "Kafka:Consumer:Topic must be set") .Validate(o => !string.IsNullOrWhiteSpace(o.GroupId), "Kafka:Consumer:GroupId must be set") + .Validate(o => !string.IsNullOrWhiteSpace(o.DeadLetterTopic), "Kafka:Consumer:DeadLetterTopic must be set") + .Validate(o => o.RetryAttempts > 0, "Kafka:Consumer:RetryAttempts must be positive") + .Validate(o => o.RetryBackoffMilliseconds >= 0, "Kafka:Consumer:RetryBackoffMilliseconds must not be negative") .ValidateOnStart(); services @@ -67,7 +70,7 @@ public static IServiceCollection AddAiAgentInfrastructure( services.AddSingleton>(sp => { var opts = sp.GetRequiredService>().Value; - return new ProducerBuilder(new ProducerConfig { BootstrapServers = opts.BootstrapServers, Acks = Acks.All, AllowAutoCreateTopics = false }).Build(); + return new ProducerBuilder(new ProducerConfig { BootstrapServers = opts.BootstrapServers, Acks = Acks.All, EnableIdempotence = true, AllowAutoCreateTopics = false }).Build(); }); services.AddHttpClient(nameof(YandexAiHttpClient), (sp, client) => @@ -85,6 +88,7 @@ public static IServiceCollection AddAiAgentInfrastructure( services.AddSingleton(); services.AddSingleton(); + services.AddSingleton(); services.AddSingleton(); services.AddHostedService(); diff --git a/src/LinkTracker.AiAgent.Infrastructure/Configuration/Kafka/RawUpdatesKafkaOptions.cs b/src/LinkTracker.AiAgent.Infrastructure/Configuration/Kafka/RawUpdatesKafkaOptions.cs index 1578cd0..d5c798e 100644 --- a/src/LinkTracker.AiAgent.Infrastructure/Configuration/Kafka/RawUpdatesKafkaOptions.cs +++ b/src/LinkTracker.AiAgent.Infrastructure/Configuration/Kafka/RawUpdatesKafkaOptions.cs @@ -5,4 +5,7 @@ public sealed class RawUpdatesKafkaOptions public string BootstrapServers { get; set; } = "localhost:9094,localhost:9095,localhost:9096"; public string Topic { get; set; } = "link.raw-updates"; public string GroupId { get; set; } = "linktracker-ai-agent"; -} \ No newline at end of file + public string DeadLetterTopic { get; set; } = "link.raw-updates-dlq"; + public int RetryAttempts { get; set; } = 3; + public int RetryBackoffMilliseconds { get; set; } = 500; +} diff --git a/src/LinkTracker.AiAgent.Infrastructure/Kafka/Abstractions/IRawUpdateDeadLetterPublisher.cs b/src/LinkTracker.AiAgent.Infrastructure/Kafka/Abstractions/IRawUpdateDeadLetterPublisher.cs new file mode 100644 index 0000000..9e63068 --- /dev/null +++ b/src/LinkTracker.AiAgent.Infrastructure/Kafka/Abstractions/IRawUpdateDeadLetterPublisher.cs @@ -0,0 +1,12 @@ +using Confluent.Kafka; + +namespace LinkTracker.AiAgent.Infrastructure.Kafka.Abstractions; + +internal interface IRawUpdateDeadLetterPublisher +{ + Task PublishAsync( + ConsumeResult sourceMessage, + string reason, + Exception? exception, + CancellationToken ct); +} diff --git a/src/LinkTracker.AiAgent.Infrastructure/Models/Kafka/RawUpdatesDeadLetterKafkaMessage.cs b/src/LinkTracker.AiAgent.Infrastructure/Models/Kafka/RawUpdatesDeadLetterKafkaMessage.cs new file mode 100644 index 0000000..782c88e --- /dev/null +++ b/src/LinkTracker.AiAgent.Infrastructure/Models/Kafka/RawUpdatesDeadLetterKafkaMessage.cs @@ -0,0 +1,18 @@ +namespace LinkTracker.AiAgent.Infrastructure.Models.Kafka; + +internal sealed class RawUpdatesDeadLetterKafkaMessage +{ + public string Payload { get; init; } = string.Empty; + + public string Reason { get; init; } = string.Empty; + + public string? ExceptionType { get; init; } + + public string SourceTopic { get; init; } = string.Empty; + + public int SourcePartition { get; init; } + + public long SourceOffset { get; init; } + + public DateTimeOffset CreatedAt { get; init; } = DateTimeOffset.UtcNow; +} diff --git a/src/LinkTracker.AiAgent.Infrastructure/Telemetry/AiAgentMetrics.cs b/src/LinkTracker.AiAgent.Infrastructure/Telemetry/AiAgentMetrics.cs new file mode 100644 index 0000000..5db8dcb --- /dev/null +++ b/src/LinkTracker.AiAgent.Infrastructure/Telemetry/AiAgentMetrics.cs @@ -0,0 +1,100 @@ +using System.Diagnostics; +using System.Diagnostics.Metrics; +using LinkTracker.AiAgent.Application.Telemetry.Abstractions; + +namespace LinkTracker.AiAgent.Infrastructure.Telemetry; + +public sealed class AiAgentMetrics : IAiAgentMetrics, IDisposable +{ + public const string MeterName = "LinkTracker.AiAgent"; + + private static readonly double[] DurationBuckets = + [5, 10, 25, 50, 100, 250, 500, 1000, 2500]; + + private readonly Counter _kafkaConsumed; + private readonly Histogram _kafkaConsumeDuration; + private readonly Counter _kafkaConsumeErrors; + private readonly Counter _kafkaDeadLetterErrors; + private readonly Counter _kafkaDeadLetters; + + private readonly Meter _meter; + + public AiAgentMetrics() + { + _meter = new Meter(MeterName); + + _kafkaConsumed = _meter.CreateCounter( + "kafka_consumed_total", + description: "Количество обработанных сообщений из Kafka"); + + _kafkaConsumeErrors = _meter.CreateCounter( + "kafka_consume_errors_total", + description: "Количество ошибок обработки сообщений из Kafka"); + + _kafkaDeadLetters = _meter.CreateCounter( + "kafka_dead_letter_total", + description: "Количество сообщений, отправленных в DLQ"); + + _kafkaDeadLetterErrors = _meter.CreateCounter( + "kafka_dead_letter_errors_total", + description: "Количество неудачных отправок в DLQ (offset не подтверждается, сообщение переигрывается)"); + + _kafkaConsumeDuration = _meter.CreateHistogram( + "kafka_consume_duration_ms_total", + null, + "Длительность обработки сообщения из Kafka в миллисекундах", + advice: new InstrumentAdvice { HistogramBucketBoundaries = DurationBuckets }); + + _meter.CreateObservableGauge( + "process_memory_working_set_bytes", + static () => Process.GetCurrentProcess().WorkingSet64, + null, + "Резидентная память процесса (working set) в байтах"); + + _meter.CreateObservableGauge( + "process_memory_managed_bytes", + static () => GC.GetTotalMemory(false), + null, + "Управляемая память (GC) в байтах"); + } + + public void IncrementKafkaConsumed(string topic) + { + _kafkaConsumed.Add( + 1, + new KeyValuePair("topic", topic)); + } + + public void IncrementKafkaConsumeError(string topic) + { + _kafkaConsumeErrors.Add( + 1, + new KeyValuePair("topic", topic)); + } + + public void ObserveKafkaConsumeDuration(string topic, double milliseconds) + { + _kafkaConsumeDuration.Record( + milliseconds, + new KeyValuePair("topic", topic)); + } + + public void IncrementKafkaDeadLetter(string topic) + { + _kafkaDeadLetters.Add( + 1, + new KeyValuePair("topic", topic)); + } + + public void IncrementKafkaDeadLetterError(string topic) + { + _kafkaDeadLetterErrors.Add( + 1, + new KeyValuePair("topic", topic)); + } + + public void Dispose() + { + _meter.Dispose(); + } +} diff --git a/src/LinkTracker.AiAgent.Infrastructure/Telemetry/Registration/TelemetryModule.cs b/src/LinkTracker.AiAgent.Infrastructure/Telemetry/Registration/TelemetryModule.cs new file mode 100644 index 0000000..788bc20 --- /dev/null +++ b/src/LinkTracker.AiAgent.Infrastructure/Telemetry/Registration/TelemetryModule.cs @@ -0,0 +1,22 @@ +using LinkTracker.AiAgent.Application.Telemetry.Abstractions; +using LinkTracker.Shared.Infrastructure.Telemetry; +using Microsoft.Extensions.Configuration; +using Microsoft.Extensions.DependencyInjection; + +namespace LinkTracker.AiAgent.Infrastructure.Telemetry.Registration; + +public static class TelemetryModule +{ + public static IServiceCollection AddTelemetry( + this IServiceCollection services, + IConfiguration configuration) + { + services.AddSingleton(); + + services.AddOpenTelemetryMetrics( + "aiagent", + AiAgentMetrics.MeterName); + + return services; + } +} diff --git a/src/LinkTracker.Bot.Api/appsettings.Docker.json b/src/LinkTracker.Bot.Api/appsettings.Docker.json index c7c7c70..05f63a5 100644 --- a/src/LinkTracker.Bot.Api/appsettings.Docker.json +++ b/src/LinkTracker.Bot.Api/appsettings.Docker.json @@ -15,7 +15,7 @@ "BootstrapServers": "kafka-1:9092,kafka-2:9092,kafka-3:9092", "Topic": "link.processed-updates", "GroupId": "linktracker-bot", - "DeadLetterTopic": "link-updates-dlq", + "DeadLetterTopic": "link.processed-updates-dlq", "RetryAttempts": 3, "RetryBackoffMilliseconds": 500, "Serialization": "Json", diff --git a/src/LinkTracker.Bot.Api/appsettings.json b/src/LinkTracker.Bot.Api/appsettings.json index 6973521..b7ff7a2 100644 --- a/src/LinkTracker.Bot.Api/appsettings.json +++ b/src/LinkTracker.Bot.Api/appsettings.json @@ -28,12 +28,12 @@ }, "Kafka": { "BootstrapServers": "localhost:9094,localhost:9095,localhost:9096", - "Topic": "link-updates", + "Topic": "link.processed-updates", "GroupId": "linktracker-bot", - "DeadLetterTopic": "link-updates-dlq", + "DeadLetterTopic": "link.processed-updates-dlq", "RetryAttempts": 3, "RetryBackoffMilliseconds": 500, - "Serialization": "Avro", + "Serialization": "Json", "SchemaRegistryUrl": "http://localhost:8071" }, "Resilience": { diff --git a/src/LinkTracker.Bot.Application/Telemetry/Abstractions/IBotMetrics.cs b/src/LinkTracker.Bot.Application/Telemetry/Abstractions/IBotMetrics.cs index bd9fc96..a74f4de 100644 --- a/src/LinkTracker.Bot.Application/Telemetry/Abstractions/IBotMetrics.cs +++ b/src/LinkTracker.Bot.Application/Telemetry/Abstractions/IBotMetrics.cs @@ -19,4 +19,8 @@ public interface IBotMetrics void IncrementKafkaConsumeError(string topic); void ObserveKafkaConsumeDuration(string topic, double milliseconds); + + void IncrementKafkaDeadLetter(string topic); + + void IncrementKafkaDeadLetterError(string topic); } \ No newline at end of file diff --git a/src/LinkTracker.Bot.Infrastructure/Clients/Kafka/LinkUpdatesKafkaMessageHandler.cs b/src/LinkTracker.Bot.Infrastructure/Clients/Kafka/LinkUpdatesKafkaMessageHandler.cs index d4b8077..f25d997 100644 --- a/src/LinkTracker.Bot.Infrastructure/Clients/Kafka/LinkUpdatesKafkaMessageHandler.cs +++ b/src/LinkTracker.Bot.Infrastructure/Clients/Kafka/LinkUpdatesKafkaMessageHandler.cs @@ -1,4 +1,5 @@ using Confluent.Kafka; +using LinkTracker.Bot.Application.Telemetry.Abstractions; using LinkTracker.Bot.Application.Updates.Abstractions; using LinkTracker.Bot.Infrastructure.Abstractions.Kafka; using LinkTracker.Bot.Infrastructure.Configuration.Kafka; @@ -15,6 +16,7 @@ internal sealed class LinkUpdatesKafkaMessageHandler( ILinkUpdateDeadLetterPublisher deadLetterPublisher, ILinkUpdateNotifier notifier, IOptions kafkaOptions, + IBotMetrics metrics, ILogger logger) : ILinkUpdatesKafkaMessageHandler { public async Task HandleAsync(ConsumeResult result, CancellationToken ct) @@ -122,6 +124,9 @@ private async Task TryPublishToDeadLetterAsync( try { await deadLetterPublisher.PublishAsync(result, reason, exception, ct); + + metrics.IncrementKafkaDeadLetter(result.Topic); + return true; } catch (OperationCanceledException) when (ct.IsCancellationRequested) @@ -130,12 +135,16 @@ private async Task TryPublishToDeadLetterAsync( } catch (Exception ex) { + metrics.IncrementKafkaDeadLetterError(result.Topic); + metrics.IncrementError("kafka_dead_letter", result.Topic, "publish_failed"); + logger.LogError( ex, - "Не удалось отправить Kafka сообщение в DLQ. Offset не будет подтвержден. Topic={Topic}, Partition={Partition}, Offset={Offset}", + "Не удалось отправить Kafka сообщение в DLQ. Offset не будет подтвержден, сообщение будет переигрываться. Topic={Topic}, Partition={Partition}, Offset={Offset}, DeadLetterTopic={DeadLetterTopic}", result.Topic, result.Partition.Value, - result.Offset.Value); + result.Offset.Value, + kafkaOptions.Value.DeadLetterTopic); return false; } diff --git a/src/LinkTracker.Bot.Infrastructure/Configuration/Kafka/LinkUpdatesKafkaOptions.cs b/src/LinkTracker.Bot.Infrastructure/Configuration/Kafka/LinkUpdatesKafkaOptions.cs index 712b025..74ef28b 100644 --- a/src/LinkTracker.Bot.Infrastructure/Configuration/Kafka/LinkUpdatesKafkaOptions.cs +++ b/src/LinkTracker.Bot.Infrastructure/Configuration/Kafka/LinkUpdatesKafkaOptions.cs @@ -7,7 +7,7 @@ public sealed class LinkUpdatesKafkaOptions public string BootstrapServers { get; set; } = "localhost:9094,localhost:9095,localhost:9096"; public string Topic { get; set; } = "link.processed-updates"; public string GroupId { get; set; } = "linktracker-bot"; - public string DeadLetterTopic { get; set; } = "link-updates-dlq"; + public string DeadLetterTopic { get; set; } = "link.processed-updates-dlq"; public int RetryAttempts { get; set; } = 3; public int RetryBackoffMilliseconds { get; set; } = 500; public KafkaSerializationKind Serialization { get; set; } = KafkaSerializationKind.Json; diff --git a/src/LinkTracker.Bot.Infrastructure/Telemetry/BotMetrics.cs b/src/LinkTracker.Bot.Infrastructure/Telemetry/BotMetrics.cs index c079134..20e4ab1 100644 --- a/src/LinkTracker.Bot.Infrastructure/Telemetry/BotMetrics.cs +++ b/src/LinkTracker.Bot.Infrastructure/Telemetry/BotMetrics.cs @@ -18,6 +18,8 @@ public sealed class BotMetrics : IBotMetrics, IDisposable private readonly Counter _kafkaConsumed; private readonly Histogram _kafkaConsumeDuration; private readonly Counter _kafkaConsumeErrors; + private readonly Counter _kafkaDeadLetterErrors; + private readonly Counter _kafkaDeadLetters; private readonly Meter _meter; private readonly Histogram _scrapperCallDuration; @@ -64,6 +66,14 @@ public BotMetrics() "kafka_consume_errors_total", description: "Количество ошибок обработки сообщений из Kafka"); + _kafkaDeadLetters = _meter.CreateCounter( + "kafka_dead_letter_total", + description: "Количество сообщений, отправленных в DLQ"); + + _kafkaDeadLetterErrors = _meter.CreateCounter( + "kafka_dead_letter_errors_total", + description: "Количество неудачных отправок в DLQ (offset не подтверждается, сообщение переигрывается)"); + _kafkaConsumeDuration = _meter.CreateHistogram( "kafka_consume_duration_ms_total", null, @@ -148,6 +158,20 @@ public void ObserveKafkaConsumeDuration(string topic, double milliseconds) new KeyValuePair("topic", topic)); } + public void IncrementKafkaDeadLetter(string topic) + { + _kafkaDeadLetters.Add( + 1, + new KeyValuePair("topic", topic)); + } + + public void IncrementKafkaDeadLetterError(string topic) + { + _kafkaDeadLetterErrors.Add( + 1, + new KeyValuePair("topic", topic)); + } + public void Dispose() { _meter.Dispose(); diff --git a/src/LinkTracker.Scrapper.Api/appsettings.json b/src/LinkTracker.Scrapper.Api/appsettings.json index 4c4dec8..108076c 100644 --- a/src/LinkTracker.Scrapper.Api/appsettings.json +++ b/src/LinkTracker.Scrapper.Api/appsettings.json @@ -20,7 +20,7 @@ }, "Bot": { "BaseUrl": "http://localhost:8091", - "Transport": "Http" + "Transport": "Kafka" }, "Scheduling": { "IntervalSeconds": 30, @@ -45,8 +45,8 @@ }, "Kafka": { "BootstrapServers": "localhost:9094,localhost:9095,localhost:9096", - "Topic": "link-updates", - "Serialization": "Avro", + "Topic": "link.raw-updates", + "Serialization": "Json", "SchemaRegistryUrl": "http://localhost:8071" }, "Outbox": { diff --git a/src/LinkTracker.Tests/AiAgent/Integration/Infrastructure/Clients/Kafka/RawUpdatesKafkaConsumerIntegrationTests.cs b/src/LinkTracker.Tests/AiAgent/Integration/Infrastructure/Clients/Kafka/RawUpdatesKafkaConsumerIntegrationTests.cs index 19b2cb7..ba3cf4e 100644 --- a/src/LinkTracker.Tests/AiAgent/Integration/Infrastructure/Clients/Kafka/RawUpdatesKafkaConsumerIntegrationTests.cs +++ b/src/LinkTracker.Tests/AiAgent/Integration/Infrastructure/Clients/Kafka/RawUpdatesKafkaConsumerIntegrationTests.cs @@ -3,9 +3,11 @@ using Confluent.Kafka; using LinkTracker.AiAgent.Application.Abstractions; using LinkTracker.AiAgent.Application.Services; +using LinkTracker.AiAgent.Application.Telemetry.Abstractions; using LinkTracker.AiAgent.Infrastructure.Clients.Kafka; using LinkTracker.AiAgent.Infrastructure.Configuration.AiAgent; using LinkTracker.AiAgent.Infrastructure.Configuration.Kafka; +using LinkTracker.AiAgent.Infrastructure.Kafka.Abstractions; using LinkTracker.AiAgent.Infrastructure.Kafka.Deserialization; using LinkTracker.AiAgent.Infrastructure.Services; using LinkTracker.Shared.Contracts.AiAgent; @@ -125,7 +127,82 @@ public async Task Consumer_WhenMalformedMessagePublished_DoesNotCrash() } } - private RawUpdatesKafkaConsumer BuildConsumer(string topic, IGroupingBuffer groupingBuffer) + [Fact] + public async Task Consumer_WhenMalformedMessagePublished_PublishesMessageToDeadLetterTopic() + { + var topic = $"link-raw-dlq-{Guid.NewGuid():N}"; + var deadLetterTopic = $"link-raw-dlq-{Guid.NewGuid():N}-dlq"; + + await kafkaFixture.CreateTopicAsync(topic); + await kafkaFixture.CreateTopicAsync(deadLetterTopic); + + using var deadLetterProducer = new ProducerBuilder( + new ProducerConfig { BootstrapServers = kafkaFixture.BootstrapServers, Acks = Acks.All, EnableIdempotence = true }).Build(); + + var deadLetterPublisher = new RawUpdatesKafkaDeadLetterPublisher( + deadLetterProducer, + Options.Create(new RawUpdatesKafkaOptions { DeadLetterTopic = deadLetterTopic }), + NullLogger.Instance); + + var groupingBuffer = Substitute.For(); + using var consumer = BuildConsumer(topic, groupingBuffer, deadLetterPublisher, deadLetterTopic); + + await consumer.StartAsync(CancellationToken.None); + + try + { + await ProduceRawAsync(topic, "{ this is not valid json !!!"); + + var deadLettered = ReadFirstMessage(deadLetterTopic); + + Assert.NotNull(deadLettered); + + using var document = JsonDocument.Parse(deadLettered); + var root = document.RootElement; + + Assert.Equal(topic, root.GetProperty("sourceTopic").GetString()); + Assert.Contains("десериализовать", root.GetProperty("reason").GetString()); + Assert.Equal( + "{ this is not valid json !!!", + Encoding.UTF8.GetString(Convert.FromBase64String(root.GetProperty("payload").GetString()!))); + + groupingBuffer.DidNotReceive().Add(Arg.Any(), Arg.Any()); + } + finally + { + await consumer.StopAsync(CancellationToken.None); + } + } + + private string? ReadFirstMessage(string topic) + { + using var consumer = new ConsumerBuilder(new ConsumerConfig + { + BootstrapServers = kafkaFixture.BootstrapServers, + GroupId = $"dlq-reader-{Guid.NewGuid():N}", + AutoOffsetReset = AutoOffsetReset.Earliest, + EnableAutoCommit = false + }).Build(); + + consumer.Subscribe(topic); + + try + { + var result = consumer.Consume(TimeSpan.FromSeconds(30)); + + return result is null ? null : Encoding.UTF8.GetString(result.Message.Value); + } + finally + { + consumer.Close(); + } + } + + private RawUpdatesKafkaConsumer BuildConsumer( + string topic, + IGroupingBuffer groupingBuffer, + IRawUpdateDeadLetterPublisher? deadLetterPublisher = null, + string deadLetterTopic = "unused-dlq") { var summarizer = Substitute.For(); summarizer.SummarizeAsync(Arg.Any(), Arg.Any()) @@ -134,7 +211,15 @@ private RawUpdatesKafkaConsumer BuildConsumer(string topic, IGroupingBuffer grou var prioritizer = Substitute.For(); prioritizer.Prioritize(Arg.Any()).Returns(LinkUpdatePriority.Medium); - var consumerOpts = Options.Create(new RawUpdatesKafkaOptions { BootstrapServers = kafkaFixture.BootstrapServers, Topic = topic, GroupId = $"test-group-{Guid.NewGuid():N}" }); + var consumerOpts = Options.Create(new RawUpdatesKafkaOptions + { + BootstrapServers = kafkaFixture.BootstrapServers, + Topic = topic, + GroupId = $"test-group-{Guid.NewGuid():N}", + DeadLetterTopic = deadLetterTopic, + RetryAttempts = 1, + RetryBackoffMilliseconds = 0 + }); var aiAgentOpts = Options.Create(new AiAgentOptions { Filtering = new FilteringOptions { MinLength = 10, StopWords = [], ExcludedAuthors = [] }, Summarization = new SummarizationOptions { Threshold = 1000 } }); @@ -148,6 +233,9 @@ private RawUpdatesKafkaConsumer BuildConsumer(string topic, IGroupingBuffer grou var messageHandler = new RawUpdatesKafkaMessageHandler( new JsonRawLinkUpdateKafkaDeserializer(), processingService, + deadLetterPublisher ?? Substitute.For(), + consumerOpts, + Substitute.For(), NullLogger.Instance); var kafkaConsumer = new ConsumerBuilder(new ConsumerConfig @@ -161,6 +249,7 @@ private RawUpdatesKafkaConsumer BuildConsumer(string topic, IGroupingBuffer grou return new RawUpdatesKafkaConsumer( kafkaConsumer, messageHandler, consumerOpts, + Substitute.For(), NullLogger.Instance); } diff --git a/src/LinkTracker.Tests/AiAgent/Unit/Infrastructure/Clients/Kafka/RawUpdatesKafkaMessageHandlerTests.cs b/src/LinkTracker.Tests/AiAgent/Unit/Infrastructure/Clients/Kafka/RawUpdatesKafkaMessageHandlerTests.cs new file mode 100644 index 0000000..49c6b16 --- /dev/null +++ b/src/LinkTracker.Tests/AiAgent/Unit/Infrastructure/Clients/Kafka/RawUpdatesKafkaMessageHandlerTests.cs @@ -0,0 +1,197 @@ +using System.Text; +using Confluent.Kafka; +using LinkTracker.AiAgent.Application.Abstractions; +using LinkTracker.AiAgent.Application.Telemetry.Abstractions; +using LinkTracker.AiAgent.Infrastructure.Clients.Kafka; +using LinkTracker.AiAgent.Infrastructure.Configuration.Kafka; +using LinkTracker.AiAgent.Infrastructure.Kafka.Abstractions; +using LinkTracker.AiAgent.Infrastructure.Kafka.Deserialization; +using LinkTracker.Shared.Contracts.Bot; +using Microsoft.Extensions.Logging.Abstractions; +using Microsoft.Extensions.Options; +using NSubstitute; + +namespace LinkTracker.Tests.AiAgent.Unit.Infrastructure.Clients.Kafka; + +[Trait("Module", "AiAgent")] +[Trait("Category", "Unit")] +public sealed class RawUpdatesKafkaMessageHandlerTests +{ + [Fact] + public async Task HandleAsync_WhenMessageIsValid_ProcessesAndReturnsTrue() + { + var processingService = Substitute.For(); + var deadLetterPublisher = Substitute.For(); + + var sut = CreateSut(processingService, deadLetterPublisher, 3); + + var result = await sut.HandleAsync(CreateValidConsumeResult(), CancellationToken.None); + + Assert.True(result); + + await processingService.Received(1).ProcessAsync( + Arg.Is(u => u.Id == 42), + Arg.Any()); + + await deadLetterPublisher.DidNotReceive().PublishAsync( + Arg.Any>(), + Arg.Any(), + Arg.Any(), + Arg.Any()); + } + + [Fact] + public async Task HandleAsync_WhenMessageIsMalformed_PublishesToDeadLetterAndReturnsTrue() + { + var processingService = Substitute.For(); + var deadLetterPublisher = Substitute.For(); + + var sut = CreateSut(processingService, deadLetterPublisher, 3); + + var message = CreateConsumeResult("{ invalid json"); + + var result = await sut.HandleAsync(message, CancellationToken.None); + + Assert.True(result); + + await processingService.DidNotReceive().ProcessAsync( + Arg.Any(), + Arg.Any()); + + await deadLetterPublisher.Received(1).PublishAsync( + message, + Arg.Any(), + Arg.Any(), + Arg.Any()); + } + + [Fact] + public async Task HandleAsync_WhenProcessingFails_RetriesThenPublishesToDeadLetter() + { + var processingService = Substitute.For(); + var deadLetterPublisher = Substitute.For(); + + processingService + .ProcessAsync(Arg.Any(), Arg.Any()) + .Returns(Task.FromException(new InvalidOperationException("YandexAi failed"))); + + var sut = CreateSut(processingService, deadLetterPublisher, 3); + + var message = CreateValidConsumeResult(); + + var result = await sut.HandleAsync(message, CancellationToken.None); + + Assert.True(result); + + await processingService.Received(3).ProcessAsync( + Arg.Any(), + Arg.Any()); + + await deadLetterPublisher.Received(1).PublishAsync( + message, + "Исчерпаны попытки обработки Kafka сообщения.", + Arg.Is(ex => ex.Message == "YandexAi failed"), + Arg.Any()); + } + + [Fact] + public async Task HandleAsync_WhenDeadLetterPublishingFails_ReturnsFalseAndIncrementsMetric() + { + var processingService = Substitute.For(); + var deadLetterPublisher = Substitute.For(); + var metrics = Substitute.For(); + + deadLetterPublisher + .PublishAsync( + Arg.Any>(), + Arg.Any(), + Arg.Any(), + Arg.Any()) + .Returns(Task.FromException(new InvalidOperationException("Kafka DLQ failed"))); + + var sut = CreateSut(processingService, deadLetterPublisher, 3, metrics: metrics); + + var result = await sut.HandleAsync(CreateConsumeResult("{ invalid json"), CancellationToken.None); + + Assert.False(result); + + metrics.Received(1).IncrementKafkaDeadLetterError("link.raw-updates"); + metrics.DidNotReceive().IncrementKafkaDeadLetter(Arg.Any()); + } + + [Fact] + public async Task HandleAsync_WhenCancellationRequested_ThrowsOperationCanceledException() + { + var processingService = Substitute.For(); + var deadLetterPublisher = Substitute.For(); + + using var cts = new CancellationTokenSource(); + + processingService + .ProcessAsync(Arg.Any(), Arg.Any()) + .Returns(_ => Task.FromCanceled(cts.Token)); + + var sut = CreateSut(processingService, deadLetterPublisher, 3); + + await cts.CancelAsync(); + + await Assert.ThrowsAnyAsync(() => + sut.HandleAsync(CreateValidConsumeResult(), cts.Token)); + + await deadLetterPublisher.DidNotReceive().PublishAsync( + Arg.Any>(), + Arg.Any(), + Arg.Any(), + Arg.Any()); + } + + private static RawUpdatesKafkaMessageHandler CreateSut( + ILinkUpdateProcessingService processingService, + IRawUpdateDeadLetterPublisher deadLetterPublisher, + int retryAttempts, + int retryBackoffMilliseconds = 0, + IAiAgentMetrics? metrics = null) + { + var options = Options.Create(new RawUpdatesKafkaOptions + { + BootstrapServers = "localhost:9092", + Topic = "link.raw-updates", + GroupId = "linktracker-ai-agent", + DeadLetterTopic = "link.raw-updates-dlq", + RetryAttempts = retryAttempts, + RetryBackoffMilliseconds = retryBackoffMilliseconds + }); + + return new RawUpdatesKafkaMessageHandler( + new JsonRawLinkUpdateKafkaDeserializer(), + processingService, + deadLetterPublisher, + options, + metrics ?? Substitute.For(), + NullLogger.Instance); + } + + private static ConsumeResult CreateValidConsumeResult() + { + return CreateConsumeResult( + """ + { + "id": 42, + "url": "https://github.com/user/repo", + "description": "Repository updated", + "tgChatIds": [123] + } + """); + } + + private static ConsumeResult CreateConsumeResult(string payload) + { + return new ConsumeResult + { + Topic = "link.raw-updates", + Partition = new Partition(0), + Offset = new Offset(1), + Message = new Message { Key = "key", Value = Encoding.UTF8.GetBytes(payload) } + }; + } +} diff --git a/src/LinkTracker.Tests/Bot/Integration/Infrastructure/Clients/Kafka/LinkUpdatesKafkaConsumerIntegrationTests.cs b/src/LinkTracker.Tests/Bot/Integration/Infrastructure/Clients/Kafka/LinkUpdatesKafkaConsumerIntegrationTests.cs index 0c16d80..f527c6d 100644 --- a/src/LinkTracker.Tests/Bot/Integration/Infrastructure/Clients/Kafka/LinkUpdatesKafkaConsumerIntegrationTests.cs +++ b/src/LinkTracker.Tests/Bot/Integration/Infrastructure/Clients/Kafka/LinkUpdatesKafkaConsumerIntegrationTests.cs @@ -323,6 +323,7 @@ private LinkUpdatesKafkaConsumer CreateConsumer( deadLetterPublisher, notifier, options, + Substitute.For(), NullLogger.Instance); return new LinkUpdatesKafkaConsumer( @@ -377,6 +378,7 @@ private LinkUpdatesKafkaConsumer CreateAvroConsumer( deadLetterPublisher, notifier, options, + Substitute.For(), NullLogger.Instance); return new LinkUpdatesKafkaConsumer( diff --git a/src/LinkTracker.Tests/Bot/Unit/Infrastructure/Clients/Kafka/LinkUpdatesKafkaMessageHandlerTests.cs b/src/LinkTracker.Tests/Bot/Unit/Infrastructure/Clients/Kafka/LinkUpdatesKafkaMessageHandlerTests.cs index eea16d2..a0c0298 100644 --- a/src/LinkTracker.Tests/Bot/Unit/Infrastructure/Clients/Kafka/LinkUpdatesKafkaMessageHandlerTests.cs +++ b/src/LinkTracker.Tests/Bot/Unit/Infrastructure/Clients/Kafka/LinkUpdatesKafkaMessageHandlerTests.cs @@ -1,6 +1,7 @@ using System.Text; using System.Text.Json; using Confluent.Kafka; +using LinkTracker.Bot.Application.Telemetry.Abstractions; using LinkTracker.Bot.Application.Updates.Abstractions; using LinkTracker.Bot.Infrastructure.Abstractions.Kafka; using LinkTracker.Bot.Infrastructure.Clients.Kafka; @@ -194,6 +195,57 @@ await deadLetterPublisher.Received(1).PublishAsync( Arg.Any()); } + [Fact] + public async Task HandleAsync_WhenDeadLetterPublishingFails_IncrementsDeadLetterErrorMetric() + { + var notifier = Substitute.For(); + var deadLetterPublisher = Substitute.For(); + var metrics = Substitute.For(); + + deadLetterPublisher + .PublishAsync( + Arg.Any>(), + Arg.Any(), + Arg.Any(), + Arg.Any()) + .Returns(Task.FromException(new InvalidOperationException("Kafka DLQ failed"))); + + var sut = CreateSut( + notifier, + deadLetterPublisher, + 3, + metrics: metrics); + + var result = await sut.HandleAsync(CreateConsumeResult("{ invalid json"), CancellationToken.None); + + Assert.False(result); + + metrics.Received(1).IncrementKafkaDeadLetterError("link.processed-updates"); + metrics.Received(1).IncrementError("kafka_dead_letter", "link.processed-updates", "publish_failed"); + metrics.DidNotReceive().IncrementKafkaDeadLetter(Arg.Any()); + } + + [Fact] + public async Task HandleAsync_WhenDeadLetterPublishingSucceeds_IncrementsDeadLetterMetric() + { + var notifier = Substitute.For(); + var deadLetterPublisher = Substitute.For(); + var metrics = Substitute.For(); + + var sut = CreateSut( + notifier, + deadLetterPublisher, + 3, + metrics: metrics); + + var result = await sut.HandleAsync(CreateConsumeResult("{ invalid json"), CancellationToken.None); + + Assert.True(result); + + metrics.Received(1).IncrementKafkaDeadLetter("link.processed-updates"); + metrics.DidNotReceive().IncrementKafkaDeadLetterError(Arg.Any()); + } + [Fact] public async Task HandleAsync_WhenCancellationRequested_ThrowsOperationCanceledException() { @@ -229,12 +281,13 @@ private static LinkUpdatesKafkaMessageHandler CreateSut( ILinkUpdateNotifier notifier, ILinkUpdateDeadLetterPublisher deadLetterPublisher, int retryAttempts, - int retryBackoffMilliseconds = 0) + int retryBackoffMilliseconds = 0, + IBotMetrics? metrics = null) { var options = Options.Create(new LinkUpdatesKafkaOptions { - Topic = "link-updates", - DeadLetterTopic = "link-updates-dlq", + Topic = "link.processed-updates", + DeadLetterTopic = "link.processed-updates-dlq", GroupId = "linktracker-bot", BootstrapServers = "localhost:9092", RetryAttempts = retryAttempts, @@ -247,6 +300,7 @@ private static LinkUpdatesKafkaMessageHandler CreateSut( deadLetterPublisher, notifier, options, + metrics ?? Substitute.For(), NullLogger.Instance); } @@ -265,6 +319,6 @@ private static ConsumeResult CreateValidConsumeResult() private static ConsumeResult CreateConsumeResult(string payload) { - return new ConsumeResult { Topic = "link-updates", Partition = new Partition(0), Offset = new Offset(1), Message = new Message { Key = "key", Value = Encoding.UTF8.GetBytes(payload) } }; + return new ConsumeResult { Topic = "link.processed-updates", Partition = new Partition(0), Offset = new Offset(1), Message = new Message { Key = "key", Value = Encoding.UTF8.GetBytes(payload) } }; } } \ No newline at end of file From 64450f51be58da25e51c199b39240f34561e73cb Mon Sep 17 00:00:00 2001 From: Echways Date: Sun, 16 Aug 2026 23:14:07 +0300 Subject: [PATCH 2/2] fix sys header logic && fix buffer && fix avro --- .../Abstractions/IGroupingBuffer.cs | 11 +- .../ILinkUpdateProcessingService.cs | 2 +- .../Abstractions/IMessageAck.cs | 8 ++ .../Services/LinkUpdateProcessingService.cs | 50 +++++-- .../Kafka/IRawUpdatesKafkaMessageHandler.cs | 3 +- .../Clients/Kafka/KafkaOffsetTracker.cs | 111 +++++++++++++++ .../Clients/Kafka/RawUpdatesKafkaConsumer.cs | 44 +++--- .../Kafka/RawUpdatesKafkaMessageHandler.cs | 14 +- .../Clients/Registration/ClientsModule.cs | 10 +- .../Clients/YandexAi/YandexAiHttpClient.cs | 38 +----- .../Services/GroupingFlushJob.cs | 85 +++++++++--- .../Services/TimeWindowGroupingBuffer.cs | 41 +++--- .../AvroLinkUpdateKafkaDeserializer.cs | 24 +++- .../Notifications/LinkUpdateNotifier.cs | 5 +- .../AvroLinkUpdateKafkaSerializer.cs | 6 + .../Quartz/Jobs/LinkUpdatesJob.cs | 11 +- .../Constants/SystemMessageMarkers.cs | 6 - .../Contracts/AiAgent/ProcessedLinkUpdate.cs | 4 + .../Contracts/Bot/LinkUpdate.cs | 2 + .../Contracts/Bot/LinkUpdateAvroSchema.cs | 19 +++ .../Contracts/Bot/LinkUpdateKind.cs | 8 ++ ...RawUpdatesKafkaConsumerIntegrationTests.cs | 12 +- .../LinkUpdateProcessingServiceTests.cs | 84 ++++++++++-- .../Clients/Kafka/KafkaOffsetTrackerTests.cs | 116 ++++++++++++++++ .../RawUpdatesKafkaMessageHandlerTests.cs | 17 ++- .../Services/GroupingFlushJobTests.cs | 108 +++++++++++++++ .../Services/TimeWindowGroupingBufferTests.cs | 128 ++++++++++++++++++ ...inkUpdatesKafkaConsumerIntegrationTests.cs | 15 +- .../Notifications/LinkUpdateNotifierTests.cs | 8 +- .../Quartz/Jobs/LinkUpdatesJobTests.cs | 13 +- .../Contracts/LinkUpdateAvroSchemaTests.cs | 35 +++++ .../ProcessedUpdateWireCompatibilityTests.cs | 25 ++++ 32 files changed, 908 insertions(+), 155 deletions(-) create mode 100644 src/LinkTracker.AiAgent.Application/Abstractions/IMessageAck.cs create mode 100644 src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/KafkaOffsetTracker.cs delete mode 100644 src/LinkTracker.Shared/Constants/SystemMessageMarkers.cs create mode 100644 src/LinkTracker.Shared/Contracts/Bot/LinkUpdateKind.cs create mode 100644 src/LinkTracker.Tests/AiAgent/Unit/Infrastructure/Clients/Kafka/KafkaOffsetTrackerTests.cs create mode 100644 src/LinkTracker.Tests/AiAgent/Unit/Infrastructure/Services/GroupingFlushJobTests.cs create mode 100644 src/LinkTracker.Tests/AiAgent/Unit/Infrastructure/Services/TimeWindowGroupingBufferTests.cs create mode 100644 src/LinkTracker.Tests/Shared/Unit/Contracts/LinkUpdateAvroSchemaTests.cs diff --git a/src/LinkTracker.AiAgent.Application/Abstractions/IGroupingBuffer.cs b/src/LinkTracker.AiAgent.Application/Abstractions/IGroupingBuffer.cs index e07e881..d1bc9a6 100644 --- a/src/LinkTracker.AiAgent.Application/Abstractions/IGroupingBuffer.cs +++ b/src/LinkTracker.AiAgent.Application/Abstractions/IGroupingBuffer.cs @@ -2,8 +2,15 @@ namespace LinkTracker.AiAgent.Application.Abstractions; +public sealed record BufferedLinkUpdate(ProcessedLinkUpdate Update, IMessageAck Ack); + +public sealed record GroupingBucket(long ChatId, IReadOnlyList Updates); + public interface IGroupingBuffer { - void Add(long tgChatId, ProcessedLinkUpdate update); - IReadOnlyList<(long ChatId, IReadOnlyList Updates)> Flush(); + void Add(long tgChatId, ProcessedLinkUpdate update, IMessageAck ack); + + IReadOnlyList Flush(bool force = false); + + void Requeue(GroupingBucket bucket); } \ No newline at end of file diff --git a/src/LinkTracker.AiAgent.Application/Abstractions/ILinkUpdateProcessingService.cs b/src/LinkTracker.AiAgent.Application/Abstractions/ILinkUpdateProcessingService.cs index aa961c9..ab7dc40 100644 --- a/src/LinkTracker.AiAgent.Application/Abstractions/ILinkUpdateProcessingService.cs +++ b/src/LinkTracker.AiAgent.Application/Abstractions/ILinkUpdateProcessingService.cs @@ -4,5 +4,5 @@ namespace LinkTracker.AiAgent.Application.Abstractions; public interface ILinkUpdateProcessingService { - Task ProcessAsync(LinkUpdate update, CancellationToken ct = default); + Task ProcessAsync(LinkUpdate update, IMessageAck ack, CancellationToken ct = default); } \ No newline at end of file diff --git a/src/LinkTracker.AiAgent.Application/Abstractions/IMessageAck.cs b/src/LinkTracker.AiAgent.Application/Abstractions/IMessageAck.cs new file mode 100644 index 0000000..4f0a8ef --- /dev/null +++ b/src/LinkTracker.AiAgent.Application/Abstractions/IMessageAck.cs @@ -0,0 +1,8 @@ +namespace LinkTracker.AiAgent.Application.Abstractions; + +public interface IMessageAck +{ + void Retain(); + + void Release(); +} diff --git a/src/LinkTracker.AiAgent.Application/Services/LinkUpdateProcessingService.cs b/src/LinkTracker.AiAgent.Application/Services/LinkUpdateProcessingService.cs index 06f1e7c..2804f5a 100644 --- a/src/LinkTracker.AiAgent.Application/Services/LinkUpdateProcessingService.cs +++ b/src/LinkTracker.AiAgent.Application/Services/LinkUpdateProcessingService.cs @@ -10,10 +10,17 @@ public sealed class LinkUpdateProcessingService( ILinkUpdateSummarizer summarizer, ILinkUpdatePrioritizer prioritizer, IGroupingBuffer groupingBuffer, + IProcessedUpdatePublisher publisher, ILogger logger) : ILinkUpdateProcessingService { - public async Task ProcessAsync(LinkUpdate update, CancellationToken ct) + public async Task ProcessAsync(LinkUpdate update, IMessageAck ack, CancellationToken ct = default) { + if (update.Kind == LinkUpdateKind.SystemReport) + { + await PublishSystemReportAsync(update, ct); + return; + } + if (filter.ShouldFilter(update)) { logger.LogDebug( @@ -27,18 +34,41 @@ public async Task ProcessAsync(LinkUpdate update, CancellationToken ct) foreach (var chatId in update.TgChatIds) { - groupingBuffer.Add(chatId, new ProcessedLinkUpdate - { - Id = update.Id, - Url = update.Url, - Description = description, - TgChatIds = [chatId], - Priority = priority - }); + groupingBuffer.Add( + chatId, + new ProcessedLinkUpdate + { + Id = update.Id, + Url = update.Url, + Description = description, + TgChatIds = [chatId], + Priority = priority + }, + ack); } logger.LogDebug( "Обновление добавлено в буфер. UpdateId={UpdateId}, Priority={Priority}", update.Id, priority); } -} \ No newline at end of file + + private async Task PublishSystemReportAsync(LinkUpdate update, CancellationToken ct) + { + foreach (var chatId in update.TgChatIds) + { + await publisher.PublishAsync( + new ProcessedLinkUpdate + { + Id = update.Id, + Url = update.Url, + Description = update.Description, + TgChatIds = [chatId], + Priority = update.Priority, + Kind = LinkUpdateKind.SystemReport + }, + ct); + } + + logger.LogDebug("Служебный отчёт опубликован без обработки. UpdateId={UpdateId}", update.Id); + } +} diff --git a/src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/IRawUpdatesKafkaMessageHandler.cs b/src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/IRawUpdatesKafkaMessageHandler.cs index d0e3ee1..7a6fe19 100644 --- a/src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/IRawUpdatesKafkaMessageHandler.cs +++ b/src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/IRawUpdatesKafkaMessageHandler.cs @@ -1,8 +1,9 @@ using Confluent.Kafka; +using LinkTracker.AiAgent.Application.Abstractions; namespace LinkTracker.AiAgent.Infrastructure.Clients.Kafka; internal interface IRawUpdatesKafkaMessageHandler { - Task HandleAsync(ConsumeResult result, CancellationToken ct); + Task HandleAsync(ConsumeResult result, IMessageAck ack, CancellationToken ct); } \ No newline at end of file diff --git a/src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/KafkaOffsetTracker.cs b/src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/KafkaOffsetTracker.cs new file mode 100644 index 0000000..c85e018 --- /dev/null +++ b/src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/KafkaOffsetTracker.cs @@ -0,0 +1,111 @@ +using Confluent.Kafka; +using LinkTracker.AiAgent.Application.Abstractions; + +namespace LinkTracker.AiAgent.Infrastructure.Clients.Kafka; + +internal sealed class KafkaOffsetTracker +{ + private readonly Dictionary _partitions = []; + private readonly Lock _lock = new(); + private long _generation; + + public IMessageAck Track(ConsumeResult result) + { + var partition = result.TopicPartition; + var offset = result.Offset.Value; + + lock (_lock) + { + if (!_partitions.TryGetValue(partition, out var state)) + { + state = new PartitionState(_generation) { LastCommitted = offset - 1 }; + _partitions[partition] = state; + } + + state.InFlight.Add(offset); + state.MaxSeen = Math.Max(state.MaxSeen, offset); + + return new MessageAck(this, partition, state.Generation, offset); + } + } + + public IReadOnlyList TakeCommittableOffsets() + { + lock (_lock) + { + var result = new List(); + + foreach (var (partition, state) in _partitions) + { + var watermark = state.InFlight.Count == 0 ? state.MaxSeen : state.InFlight.Min - 1; + + if (watermark < 0 || watermark <= state.LastCommitted) + { + continue; + } + + state.LastCommitted = watermark; + result.Add(new TopicPartitionOffset(partition, new Offset(watermark + 1))); + } + + return result; + } + } + + public void Forget(IEnumerable partitions) + { + lock (_lock) + { + _generation++; + + foreach (var partition in partitions) + { + _partitions.Remove(partition); + } + } + } + + private void Complete(TopicPartition partition, long generation, long offset) + { + lock (_lock) + { + if (_partitions.TryGetValue(partition, out var state) && state.Generation == generation) + { + state.InFlight.Remove(offset); + } + } + } + + private sealed class PartitionState(long generation) + { + public SortedSet InFlight { get; } = []; + + public long MaxSeen { get; set; } = -1; + + public long LastCommitted { get; set; } = -1; + + public long Generation { get; } = generation; + } + + private sealed class MessageAck( + KafkaOffsetTracker tracker, + TopicPartition partition, + long generation, + long offset) : IMessageAck + { + private int _references = 1; + + public void Retain() + { + Interlocked.Increment(ref _references); + } + + public void Release() + { + if (Interlocked.Decrement(ref _references) == 0) + { + tracker.Complete(partition, generation, offset); + } + } + } +} diff --git a/src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/RawUpdatesKafkaConsumer.cs b/src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/RawUpdatesKafkaConsumer.cs index fdde3df..511f77b 100644 --- a/src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/RawUpdatesKafkaConsumer.cs +++ b/src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/RawUpdatesKafkaConsumer.cs @@ -11,13 +11,16 @@ namespace LinkTracker.AiAgent.Infrastructure.Clients.Kafka; internal sealed class RawUpdatesKafkaConsumer( IConsumer consumer, IRawUpdatesKafkaMessageHandler messageHandler, + KafkaOffsetTracker offsetTracker, IOptions kafkaOptions, IAiAgentMetrics metrics, ILogger logger) : BackgroundService { + private static readonly TimeSpan PollTimeout = TimeSpan.FromMilliseconds(500); + protected override Task ExecuteAsync(CancellationToken stoppingToken) { - return ConsumeLoopAsync(stoppingToken); + return Task.Run(() => ConsumeLoopAsync(stoppingToken), stoppingToken); } private async Task ConsumeLoopAsync(CancellationToken stoppingToken) @@ -38,7 +41,7 @@ private async Task ConsumeLoopAsync(CancellationToken stoppingToken) try { - result = consumer.Consume(stoppingToken); + result = consumer.Consume(PollTimeout); } catch (ConsumeException ex) { @@ -46,12 +49,12 @@ private async Task ConsumeLoopAsync(CancellationToken stoppingToken) continue; } - if (result is null) + if (result is not null && !result.IsPartitionEOF) { - continue; + await ProcessMessageAsync(result, stoppingToken); } - await ProcessMessageAsync(result, stoppingToken); + CommitCompleted(); } } catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) @@ -60,6 +63,7 @@ private async Task ConsumeLoopAsync(CancellationToken stoppingToken) } finally { + CommitCompleted(); consumer.Close(); } } @@ -68,18 +72,20 @@ private async Task ProcessMessageAsync(ConsumeResult result, Can { var sw = Stopwatch.StartNew(); + var ack = offsetTracker.Track(result); + try { - var shouldCommit = await messageHandler.HandleAsync(result, ct); + var handled = await messageHandler.HandleAsync(result, ack, ct); sw.Stop(); metrics.IncrementKafkaConsumed(result.Topic); metrics.ObserveKafkaConsumeDuration(result.Topic, sw.Elapsed.TotalMilliseconds); - if (shouldCommit) + if (handled) { - TryCommit(result); + ack.Release(); } } catch (OperationCanceledException) when (ct.IsCancellationRequested) @@ -102,23 +108,25 @@ private async Task ProcessMessageAsync(ConsumeResult result, Can } } - private bool TryCommit(ConsumeResult result) + private void CommitCompleted() { + var offsets = offsetTracker.TakeCommittableOffsets(); + + if (offsets.Count == 0) + { + return; + } + try { - consumer.Commit(result); - return true; + consumer.Commit(offsets); } catch (KafkaException ex) { logger.LogError( ex, - "Не удалось подтвердить Kafka offset. Topic={Topic}, Partition={Partition}, Offset={Offset}", - result.Topic, - result.Partition.Value, - result.Offset.Value); - - return false; + "Не удалось подтвердить Kafka offsets. Offsets={Offsets}", + string.Join(", ", offsets)); } } -} \ No newline at end of file +} diff --git a/src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/RawUpdatesKafkaMessageHandler.cs b/src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/RawUpdatesKafkaMessageHandler.cs index 27a469e..3e29f68 100644 --- a/src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/RawUpdatesKafkaMessageHandler.cs +++ b/src/LinkTracker.AiAgent.Infrastructure/Clients/Kafka/RawUpdatesKafkaMessageHandler.cs @@ -17,7 +17,10 @@ internal sealed class RawUpdatesKafkaMessageHandler( IAiAgentMetrics metrics, ILogger logger) : IRawUpdatesKafkaMessageHandler { - public async Task HandleAsync(ConsumeResult result, CancellationToken ct) + public async Task HandleAsync( + ConsumeResult result, + IMessageAck ack, + CancellationToken ct) { LinkUpdate? update; @@ -47,7 +50,7 @@ public async Task HandleAsync(ConsumeResult result, Cancel ct); } - var processingError = await TryProcessWithRetriesAsync(update, ct); + var processingError = await TryProcessWithRetriesAsync(update, ack, ct); if (processingError is not null) { @@ -68,7 +71,10 @@ public async Task HandleAsync(ConsumeResult result, Cancel return true; } - private async Task TryProcessWithRetriesAsync(LinkUpdate update, CancellationToken ct) + private async Task TryProcessWithRetriesAsync( + LinkUpdate update, + IMessageAck ack, + CancellationToken ct) { var attempts = Math.Max(1, kafkaOptions.Value.RetryAttempts); var backoff = TimeSpan.FromMilliseconds(Math.Max(0, kafkaOptions.Value.RetryBackoffMilliseconds)); @@ -77,7 +83,7 @@ public async Task HandleAsync(ConsumeResult result, Cancel { try { - await processingService.ProcessAsync(update, ct); + await processingService.ProcessAsync(update, ack, ct); return null; } catch (OperationCanceledException) when (ct.IsCancellationRequested) diff --git a/src/LinkTracker.AiAgent.Infrastructure/Clients/Registration/ClientsModule.cs b/src/LinkTracker.AiAgent.Infrastructure/Clients/Registration/ClientsModule.cs index bc6a3d4..d652995 100644 --- a/src/LinkTracker.AiAgent.Infrastructure/Clients/Registration/ClientsModule.cs +++ b/src/LinkTracker.AiAgent.Infrastructure/Clients/Registration/ClientsModule.cs @@ -54,9 +54,13 @@ public static IServiceCollection AddAiAgentInfrastructure( .Bind(configuration.GetSection("YandexAi")) .ValidateOnStart(); + services.AddSingleton(); + services.AddSingleton>(sp => { var opts = sp.GetRequiredService>().Value; + var offsetTracker = sp.GetRequiredService(); + return new ConsumerBuilder(new ConsumerConfig { BootstrapServers = opts.BootstrapServers, @@ -64,7 +68,10 @@ public static IServiceCollection AddAiAgentInfrastructure( EnableAutoCommit = false, AutoOffsetReset = AutoOffsetReset.Earliest, AllowAutoCreateTopics = false - }).Build(); + }) + .SetPartitionsRevokedHandler((_, partitions) => + offsetTracker.Forget(partitions.Select(x => x.TopicPartition))) + .Build(); }); services.AddSingleton>(sp => @@ -95,6 +102,7 @@ public static IServiceCollection AddAiAgentInfrastructure( services.AddSingleton(); services.AddSingleton(); services.AddSingleton(); + services.AddHostedService(); return services; diff --git a/src/LinkTracker.AiAgent.Infrastructure/Clients/YandexAi/YandexAiHttpClient.cs b/src/LinkTracker.AiAgent.Infrastructure/Clients/YandexAi/YandexAiHttpClient.cs index a4c64b1..f352aa9 100644 --- a/src/LinkTracker.AiAgent.Infrastructure/Clients/YandexAi/YandexAiHttpClient.cs +++ b/src/LinkTracker.AiAgent.Infrastructure/Clients/YandexAi/YandexAiHttpClient.cs @@ -4,7 +4,6 @@ using LinkTracker.AiAgent.Infrastructure.Clients.YandexAi.Contracts; using LinkTracker.AiAgent.Infrastructure.Configuration.AiAgent; using LinkTracker.AiAgent.Infrastructure.Configuration.YandexAi; -using LinkTracker.Shared.Constants; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; @@ -71,43 +70,18 @@ private async Task CallApiAsync(string text, CancellationToken ct) private static string FallbackTruncate(string text, int threshold) { - var markerIndex = text.IndexOf(SystemMessageMarkers.FailedLinkReport, StringComparison.Ordinal); + var cutAt = threshold; - if (markerIndex < 0) + while (cutAt > 0 && text[cutAt - 1] != '\n') { - var cutAt = threshold; - while (cutAt > 0 && text[cutAt - 1] != '\n') - { - cutAt--; - } - - if (cutAt == 0) - { - cutAt = threshold; - } - - return string.Concat(text.AsSpan(0, cutAt).TrimEnd(), "\n..."); + cutAt--; } - var mainContent = text[..markerIndex].TrimEnd(); - var systemSuffix = text[markerIndex..]; - - if (mainContent.Length <= threshold) - { - return text; - } - - var mainCutAt = threshold; - while (mainCutAt > 0 && mainContent[mainCutAt - 1] != '\n') - { - mainCutAt--; - } - - if (mainCutAt == 0) + if (cutAt == 0) { - mainCutAt = threshold; + cutAt = threshold; } - return string.Concat(mainContent.AsSpan(0, mainCutAt).TrimEnd(), "\n...\n", systemSuffix); + return string.Concat(text.AsSpan(0, cutAt).TrimEnd(), "\n..."); } } \ No newline at end of file diff --git a/src/LinkTracker.AiAgent.Infrastructure/Services/GroupingFlushJob.cs b/src/LinkTracker.AiAgent.Infrastructure/Services/GroupingFlushJob.cs index 1ff60c7..c3516a1 100644 --- a/src/LinkTracker.AiAgent.Infrastructure/Services/GroupingFlushJob.cs +++ b/src/LinkTracker.AiAgent.Infrastructure/Services/GroupingFlushJob.cs @@ -1,5 +1,6 @@ using LinkTracker.AiAgent.Application.Abstractions; using LinkTracker.AiAgent.Infrastructure.Configuration.AiAgent; +using LinkTracker.Shared.Contracts.AiAgent; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; @@ -19,34 +20,84 @@ protected override async Task ExecuteAsync(CancellationToken stoppingToken) while (!stoppingToken.IsCancellationRequested) { - await Task.Delay(interval, stoppingToken); - await FlushAsync(stoppingToken); + try + { + await Task.Delay(interval, stoppingToken); + } + catch (OperationCanceledException) + { + break; + } + + await FlushAsync(false, stoppingToken); } } - private async Task FlushAsync(CancellationToken ct) + // Окна, не успевшие закрыться, публикуются на остановке: иначе они умрут вместе + // с процессом, а сообщения будут переигрываться с последнего подтвержденного оффсета. + public override async Task StopAsync(CancellationToken cancellationToken) + { + await base.StopAsync(cancellationToken); + await FlushAsync(true, cancellationToken); + } + + private async Task FlushAsync(bool force, CancellationToken ct) { - var pending = buffer.Flush() - .SelectMany(entry => grouper - .Group(entry.Updates) - .Select(update => (entry.ChatId, Update: update))) - .OrderByDescending(x => x.Update.Priority) + var pending = buffer.Flush(force) + .Select(bucket => new + { + Bucket = bucket, + Groups = grouper.Group(bucket.Updates.Select(x => x.Update).ToArray()) + }) + .OrderByDescending(x => x.Groups.Max(group => group.Priority)) .ToArray(); - foreach (var (chatId, update) in pending) + foreach (var entry in pending) { - try + if (!await TryPublishAsync(entry.Bucket, entry.Groups, ct)) { - await publisher.PublishAsync(update, ct); + buffer.Requeue(entry.Bucket); - logger.LogInformation( - "Группа опубликована. ChatId={ChatId}, UpdateId={UpdateId}, Priority={Priority}", - chatId, update.Id, update.Priority); + if (ct.IsCancellationRequested) + { + break; + } } - catch (Exception ex) when (ex is not OperationCanceledException) + } + } + + private async Task TryPublishAsync( + GroupingBucket bucket, + IReadOnlyList groups, + CancellationToken ct) + { + try + { + foreach (var group in groups) { - logger.LogError(ex, "Ошибка публикации сгруппированного обновления. ChatId={ChatId}", chatId); + await publisher.PublishAsync(group, ct); + + logger.LogInformation( + "Группа опубликована. ChatId={ChatId}, UpdateId={UpdateId}, Priority={Priority}", + bucket.ChatId, group.Id, group.Priority); } } + catch (Exception ex) + { + logger.LogError( + ex, + "Ошибка публикации сгруппированного обновления, окно возвращено в буфер. ChatId={ChatId}", + bucket.ChatId); + + return false; + } + + // Оффсеты исходных сообщений подтверждаются только после успешной публикации. + foreach (var buffered in bucket.Updates) + { + buffered.Ack.Release(); + } + + return true; } -} \ No newline at end of file +} diff --git a/src/LinkTracker.AiAgent.Infrastructure/Services/TimeWindowGroupingBuffer.cs b/src/LinkTracker.AiAgent.Infrastructure/Services/TimeWindowGroupingBuffer.cs index 3aa0f35..4de6aec 100644 --- a/src/LinkTracker.AiAgent.Infrastructure/Services/TimeWindowGroupingBuffer.cs +++ b/src/LinkTracker.AiAgent.Infrastructure/Services/TimeWindowGroupingBuffer.cs @@ -7,15 +7,19 @@ namespace LinkTracker.AiAgent.Infrastructure.Services; internal sealed class TimeWindowGroupingBuffer(IOptions options) : IGroupingBuffer { - private readonly Dictionary Updates, DateTimeOffset WindowStart)> _buckets = []; - private readonly List<(long ChatId, List Updates)> _evicted = []; + private readonly Dictionary Updates, DateTimeOffset WindowStart)> _buckets = []; + private readonly List _evicted = []; private readonly Lock _lock = new(); - public void Add(long tgChatId, ProcessedLinkUpdate update) + public void Add(long tgChatId, ProcessedLinkUpdate update, IMessageAck ack) { var windowMs = options.Value.Grouping.WindowMs; var now = DateTimeOffset.UtcNow; + ack.Retain(); + + var buffered = new BufferedLinkUpdate(update, ack); + lock (_lock) { if (_buckets.TryGetValue(tgChatId, out var bucket)) @@ -24,44 +28,49 @@ public void Add(long tgChatId, ProcessedLinkUpdate update) if (elapsed <= windowMs) { - bucket.Updates.Add(update); + bucket.Updates.Add(buffered); return; } - _evicted.Add((tgChatId, bucket.Updates)); + _evicted.Add(new GroupingBucket(tgChatId, bucket.Updates)); } - _buckets[tgChatId] = ([update], now); + _buckets[tgChatId] = ([buffered], now); } } - public IReadOnlyList<(long ChatId, IReadOnlyList Updates)> Flush() + public IReadOnlyList Flush(bool force = false) { var windowMs = options.Value.Grouping.WindowMs; var now = DateTimeOffset.UtcNow; - var result = new List<(long, IReadOnlyList)>(); + var result = new List(); lock (_lock) { - foreach (var (chatId, updates) in _evicted) - { - result.Add((chatId, updates.AsReadOnly())); - } - + result.AddRange(_evicted); _evicted.Clear(); var expiredKeys = _buckets - .Where(kv => (now - kv.Value.WindowStart).TotalMilliseconds > windowMs) + .Where(kv => force || (now - kv.Value.WindowStart).TotalMilliseconds > windowMs) .Select(kv => kv.Key) .ToList(); foreach (var key in expiredKeys) { - result.Add((key, _buckets[key].Updates.AsReadOnly())); + result.Add(new GroupingBucket(key, _buckets[key].Updates)); _buckets.Remove(key); } } return result; } -} \ No newline at end of file + + public void Requeue(GroupingBucket bucket) + { + foreach (var buffered in bucket.Updates) + { + Add(bucket.ChatId, buffered.Update, buffered.Ack); + buffered.Ack.Release(); + } + } +} diff --git a/src/LinkTracker.Bot.Infrastructure/Kafka/Deserialization/AvroLinkUpdateKafkaDeserializer.cs b/src/LinkTracker.Bot.Infrastructure/Kafka/Deserialization/AvroLinkUpdateKafkaDeserializer.cs index 58b33ad..ac0ad83 100644 --- a/src/LinkTracker.Bot.Infrastructure/Kafka/Deserialization/AvroLinkUpdateKafkaDeserializer.cs +++ b/src/LinkTracker.Bot.Infrastructure/Kafka/Deserialization/AvroLinkUpdateKafkaDeserializer.cs @@ -3,6 +3,7 @@ using Confluent.Kafka; using Confluent.SchemaRegistry.Serdes; using LinkTracker.Bot.Infrastructure.Kafka.Abstractions; +using LinkTracker.Shared.Contracts.AiAgent; using LinkTracker.Shared.Contracts.Bot; namespace LinkTracker.Bot.Infrastructure.Kafka.Deserialization; @@ -22,10 +23,29 @@ internal sealed class AvroLinkUpdateKafkaDeserializer( Id = Convert.ToInt64(record["id"]), Url = new Uri((string)record["url"]), Description = (string)record["description"], + Author = ReadOptional(record, "author") as string ?? string.Empty, TgChatIds = ((IEnumerable)record["tgChatIds"]) .Cast() .Select(Convert.ToInt64) - .ToArray() + .ToArray(), + Priority = ReadEnum(record, "priority", LinkUpdatePriority.Medium), + Kind = ReadEnum(record, "kind", LinkUpdateKind.Content) }; } -} \ No newline at end of file + + private static object? ReadOptional(GenericRecord record, string field) + { + return record.TryGetValue(field, out var value) ? value : null; + } + + private static TEnum ReadEnum(GenericRecord record, string field, TEnum fallback) + where TEnum : struct, Enum + { + if (ReadOptional(record, field) is not GenericEnum value) + { + return fallback; + } + + return Enum.TryParse(value.Value, out var parsed) ? parsed : fallback; + } +} diff --git a/src/LinkTracker.Bot.Presentation/Telegram/Notifications/LinkUpdateNotifier.cs b/src/LinkTracker.Bot.Presentation/Telegram/Notifications/LinkUpdateNotifier.cs index 76b0e7c..bf6f07a 100644 --- a/src/LinkTracker.Bot.Presentation/Telegram/Notifications/LinkUpdateNotifier.cs +++ b/src/LinkTracker.Bot.Presentation/Telegram/Notifications/LinkUpdateNotifier.cs @@ -1,6 +1,5 @@ using LinkTracker.Bot.Application.Telemetry.Abstractions; using LinkTracker.Bot.Application.Updates.Abstractions; -using LinkTracker.Shared.Constants; using LinkTracker.Shared.Contracts.AiAgent; using LinkTracker.Shared.Contracts.Bot; using Telegram.Bot; @@ -22,9 +21,9 @@ public async Task NotifyAsync(LinkUpdate update, CancellationToken ct = default) private static string BuildMessage(LinkUpdate update) { - if (update.Description.StartsWith(SystemMessageMarkers.FailedLinkReport, StringComparison.Ordinal)) + if (update.Kind == LinkUpdateKind.SystemReport) { - return update.Description[SystemMessageMarkers.FailedLinkReport.Length..]; + return update.Description; } return $"{BuildHeader(update.Priority)}\n{update.Url}\n\n{update.Description}"; diff --git a/src/LinkTracker.Scrapper.Infrastructure/Kafka/Serialization/AvroLinkUpdateKafkaSerializer.cs b/src/LinkTracker.Scrapper.Infrastructure/Kafka/Serialization/AvroLinkUpdateKafkaSerializer.cs index bc3b9c3..dc130ec 100644 --- a/src/LinkTracker.Scrapper.Infrastructure/Kafka/Serialization/AvroLinkUpdateKafkaSerializer.cs +++ b/src/LinkTracker.Scrapper.Infrastructure/Kafka/Serialization/AvroLinkUpdateKafkaSerializer.cs @@ -13,6 +13,9 @@ internal sealed class AvroLinkUpdateKafkaSerializer( private static readonly RecordSchema Schema = (RecordSchema)Avro.Schema.Parse(LinkUpdateAvroSchema.Value); + private static readonly EnumSchema PrioritySchema = (EnumSchema)Schema["priority"].Schema; + private static readonly EnumSchema KindSchema = (EnumSchema)Schema["kind"].Schema; + public Task SerializeAsync(LinkUpdate update, string topic, CancellationToken ct) { var record = new GenericRecord(Schema); @@ -20,6 +23,9 @@ public Task SerializeAsync(LinkUpdate update, string topic, Cancellation record.Add("url", update.Url.ToString()); record.Add("description", update.Description); record.Add("tgChatIds", update.TgChatIds.ToArray()); + record.Add("author", update.Author); + record.Add("priority", new GenericEnum(PrioritySchema, update.Priority.ToString())); + record.Add("kind", new GenericEnum(KindSchema, update.Kind.ToString())); return serializer.SerializeAsync( record, diff --git a/src/LinkTracker.Scrapper.Infrastructure/Quartz/Jobs/LinkUpdatesJob.cs b/src/LinkTracker.Scrapper.Infrastructure/Quartz/Jobs/LinkUpdatesJob.cs index fc8965c..5d2ff06 100644 --- a/src/LinkTracker.Scrapper.Infrastructure/Quartz/Jobs/LinkUpdatesJob.cs +++ b/src/LinkTracker.Scrapper.Infrastructure/Quartz/Jobs/LinkUpdatesJob.cs @@ -8,7 +8,6 @@ using LinkTracker.Scrapper.Infrastructure.Quartz.Configuration; using LinkTracker.Scrapper.Infrastructure.Telemetry; using LinkTracker.Scrapper.Storage.Abstractions.Models; -using LinkTracker.Shared.Constants; using LinkTracker.Shared.Contracts.Bot; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; @@ -203,7 +202,14 @@ private async Task SendFailedReportsAsync( try { await botClient.SendUpdateAsync( - new LinkUpdate { Id = 0, Url = report.Urls[0], TgChatIds = [report.ChatId], Description = BuildFailedReportDescription(report.Urls) }, + new LinkUpdate + { + Id = 0, + Url = report.Urls[0], + TgChatIds = [report.ChatId], + Description = BuildFailedReportDescription(report.Urls), + Kind = LinkUpdateKind.SystemReport + }, ct); } catch (Exception ex) @@ -223,7 +229,6 @@ private static string BuildFailedReportDescription(IReadOnlyCollection urls var urlLines = string.Join(newLine, urls.Select(url => $"- {url}")); return - $"{SystemMessageMarkers.FailedLinkReport}" + $"Не удалось проверить часть ссылок в текущем цикле:{newLine}" + $"{urlLines}{newLine}{newLine}" + "Остальные ссылки были обработаны. Повторим попытку в следующем запуске."; diff --git a/src/LinkTracker.Shared/Constants/SystemMessageMarkers.cs b/src/LinkTracker.Shared/Constants/SystemMessageMarkers.cs deleted file mode 100644 index d7baaae..0000000 --- a/src/LinkTracker.Shared/Constants/SystemMessageMarkers.cs +++ /dev/null @@ -1,6 +0,0 @@ -namespace LinkTracker.Shared.Constants; - -public static class SystemMessageMarkers -{ - public const string FailedLinkReport = "[FAILED_LINK_REPORT]"; -} \ No newline at end of file diff --git a/src/LinkTracker.Shared/Contracts/AiAgent/ProcessedLinkUpdate.cs b/src/LinkTracker.Shared/Contracts/AiAgent/ProcessedLinkUpdate.cs index 7322dd8..edc85c9 100644 --- a/src/LinkTracker.Shared/Contracts/AiAgent/ProcessedLinkUpdate.cs +++ b/src/LinkTracker.Shared/Contracts/AiAgent/ProcessedLinkUpdate.cs @@ -1,3 +1,5 @@ +using LinkTracker.Shared.Contracts.Bot; + namespace LinkTracker.Shared.Contracts.AiAgent; public sealed class ProcessedLinkUpdate @@ -11,4 +13,6 @@ public sealed class ProcessedLinkUpdate public IReadOnlyList TgChatIds { get; init; } = []; public LinkUpdatePriority Priority { get; init; } = LinkUpdatePriority.Medium; + + public LinkUpdateKind Kind { get; init; } = LinkUpdateKind.Content; } \ No newline at end of file diff --git a/src/LinkTracker.Shared/Contracts/Bot/LinkUpdate.cs b/src/LinkTracker.Shared/Contracts/Bot/LinkUpdate.cs index e91091c..82ea1a8 100644 --- a/src/LinkTracker.Shared/Contracts/Bot/LinkUpdate.cs +++ b/src/LinkTracker.Shared/Contracts/Bot/LinkUpdate.cs @@ -19,4 +19,6 @@ public sealed class LinkUpdate /// приоритет там просто ещё не вычислен. /// public LinkUpdatePriority Priority { get; init; } = LinkUpdatePriority.Medium; + + public LinkUpdateKind Kind { get; init; } = LinkUpdateKind.Content; } \ No newline at end of file diff --git a/src/LinkTracker.Shared/Contracts/Bot/LinkUpdateAvroSchema.cs b/src/LinkTracker.Shared/Contracts/Bot/LinkUpdateAvroSchema.cs index ff85f71..52d2fb1 100644 --- a/src/LinkTracker.Shared/Contracts/Bot/LinkUpdateAvroSchema.cs +++ b/src/LinkTracker.Shared/Contracts/Bot/LinkUpdateAvroSchema.cs @@ -18,6 +18,25 @@ public static class LinkUpdateAvroSchema "type": "array", "items": "long" } + }, + { "name": "author", "type": "string", "default": "" }, + { + "name": "priority", + "type": { + "type": "enum", + "name": "LinkUpdatePriority", + "symbols": ["Low", "Medium", "High"] + }, + "default": "Medium" + }, + { + "name": "kind", + "type": { + "type": "enum", + "name": "LinkUpdateKind", + "symbols": ["Content", "SystemReport"] + }, + "default": "Content" } ] } diff --git a/src/LinkTracker.Shared/Contracts/Bot/LinkUpdateKind.cs b/src/LinkTracker.Shared/Contracts/Bot/LinkUpdateKind.cs new file mode 100644 index 0000000..21dba86 --- /dev/null +++ b/src/LinkTracker.Shared/Contracts/Bot/LinkUpdateKind.cs @@ -0,0 +1,8 @@ +namespace LinkTracker.Shared.Contracts.Bot; + +public enum LinkUpdateKind +{ + Content = 0, + + SystemReport = 1 +} \ No newline at end of file diff --git a/src/LinkTracker.Tests/AiAgent/Integration/Infrastructure/Clients/Kafka/RawUpdatesKafkaConsumerIntegrationTests.cs b/src/LinkTracker.Tests/AiAgent/Integration/Infrastructure/Clients/Kafka/RawUpdatesKafkaConsumerIntegrationTests.cs index ba3cf4e..64cbeba 100644 --- a/src/LinkTracker.Tests/AiAgent/Integration/Infrastructure/Clients/Kafka/RawUpdatesKafkaConsumerIntegrationTests.cs +++ b/src/LinkTracker.Tests/AiAgent/Integration/Infrastructure/Clients/Kafka/RawUpdatesKafkaConsumerIntegrationTests.cs @@ -54,7 +54,8 @@ await WaitUntilAsync(() => { groupingBuffer.Received(1).Add( 42L, - Arg.Is(u => u.Id == 100)); + Arg.Is(u => u.Id == 100), + Arg.Any()); return Task.FromResult(true); } @@ -94,7 +95,7 @@ public async Task Consumer_WhenFilteredMessagePublished_DoesNotAddToGroupingBuff await Task.Delay(TimeSpan.FromSeconds(3)); - groupingBuffer.DidNotReceive().Add(Arg.Any(), Arg.Any()); + groupingBuffer.DidNotReceive().Add(Arg.Any(), Arg.Any(), Arg.Any()); } finally { @@ -119,7 +120,7 @@ public async Task Consumer_WhenMalformedMessagePublished_DoesNotCrash() await Task.Delay(TimeSpan.FromSeconds(3)); - groupingBuffer.DidNotReceive().Add(Arg.Any(), Arg.Any()); + groupingBuffer.DidNotReceive().Add(Arg.Any(), Arg.Any(), Arg.Any()); } finally { @@ -166,7 +167,7 @@ public async Task Consumer_WhenMalformedMessagePublished_PublishesMessageToDeadL "{ this is not valid json !!!", Encoding.UTF8.GetString(Convert.FromBase64String(root.GetProperty("payload").GetString()!))); - groupingBuffer.DidNotReceive().Add(Arg.Any(), Arg.Any()); + groupingBuffer.DidNotReceive().Add(Arg.Any(), Arg.Any(), Arg.Any()); } finally { @@ -228,6 +229,7 @@ private RawUpdatesKafkaConsumer BuildConsumer( summarizer, prioritizer, groupingBuffer, + Substitute.For(), NullLogger.Instance); var messageHandler = new RawUpdatesKafkaMessageHandler( @@ -248,7 +250,7 @@ private RawUpdatesKafkaConsumer BuildConsumer( }).Build(); return new RawUpdatesKafkaConsumer( - kafkaConsumer, messageHandler, consumerOpts, + kafkaConsumer, messageHandler, new KafkaOffsetTracker(), consumerOpts, Substitute.For(), NullLogger.Instance); } diff --git a/src/LinkTracker.Tests/AiAgent/Unit/Application/Services/LinkUpdateProcessingServiceTests.cs b/src/LinkTracker.Tests/AiAgent/Unit/Application/Services/LinkUpdateProcessingServiceTests.cs index 5e95cdc..c58b227 100644 --- a/src/LinkTracker.Tests/AiAgent/Unit/Application/Services/LinkUpdateProcessingServiceTests.cs +++ b/src/LinkTracker.Tests/AiAgent/Unit/Application/Services/LinkUpdateProcessingServiceTests.cs @@ -14,11 +14,13 @@ public sealed class LinkUpdateProcessingServiceTests private readonly ILinkUpdateFilter _filter = Substitute.For(); private readonly IGroupingBuffer _groupingBuffer = Substitute.For(); private readonly ILinkUpdatePrioritizer _prioritizer = Substitute.For(); + private readonly IProcessedUpdatePublisher _publisher = Substitute.For(); private readonly ILinkUpdateSummarizer _summarizer = Substitute.For(); + private readonly IMessageAck _ack = Substitute.For(); private LinkUpdateProcessingService CreateService() { - return new LinkUpdateProcessingService(_filter, _summarizer, _prioritizer, _groupingBuffer, + return new LinkUpdateProcessingService(_filter, _summarizer, _prioritizer, _groupingBuffer, _publisher, NullLogger.Instance); } @@ -39,9 +41,9 @@ public async Task ProcessAsync_WhenFilterRejectsUpdate_DoesNotAddToBuffer() { _filter.ShouldFilter(Arg.Any()).Returns(true); - await CreateService().ProcessAsync(BuildUpdate(), CancellationToken.None); + await CreateService().ProcessAsync(BuildUpdate(), _ack, CancellationToken.None); - _groupingBuffer.DidNotReceive().Add(Arg.Any(), Arg.Any()); + _groupingBuffer.DidNotReceive().Add(Arg.Any(), Arg.Any(), Arg.Any()); } [Fact] @@ -49,7 +51,7 @@ public async Task ProcessAsync_WhenFilterRejectsUpdate_DoesNotCallSummarizer() { _filter.ShouldFilter(Arg.Any()).Returns(true); - await CreateService().ProcessAsync(BuildUpdate(), CancellationToken.None); + await CreateService().ProcessAsync(BuildUpdate(), _ack, CancellationToken.None); await _summarizer.DidNotReceive().SummarizeAsync( Arg.Any(), @@ -61,7 +63,7 @@ public async Task ProcessAsync_WhenFilterRejectsUpdate_DoesNotCallPrioritizer() { _filter.ShouldFilter(Arg.Any()).Returns(true); - await CreateService().ProcessAsync(BuildUpdate(), CancellationToken.None); + await CreateService().ProcessAsync(BuildUpdate(), _ack, CancellationToken.None); _prioritizer.DidNotReceive().Prioritize(Arg.Any()); } @@ -74,7 +76,7 @@ public async Task ProcessAsync_WhenUpdatePasses_CallsSummarizerWithOriginalDescr _summarizer.SummarizeAsync("original description", Arg.Any()) .Returns("summarized"); - await CreateService().ProcessAsync(update, CancellationToken.None); + await CreateService().ProcessAsync(update, _ack, CancellationToken.None); await _summarizer.Received(1).SummarizeAsync( "original description", @@ -89,7 +91,7 @@ public async Task ProcessAsync_WhenUpdatePasses_CallsPrioritizerWithSummarizedDe .Returns("summarized text"); _prioritizer.Prioritize(Arg.Any()).Returns(LinkUpdatePriority.Medium); - await CreateService().ProcessAsync(BuildUpdate(), CancellationToken.None); + await CreateService().ProcessAsync(BuildUpdate(), _ack, CancellationToken.None); _prioritizer.Received(1).Prioritize("summarized text"); } @@ -111,10 +113,10 @@ public async Task ProcessAsync_WhenUpdatePasses_AddsToBufferForEachChatId() .Returns("summarized"); _prioritizer.Prioritize(Arg.Any()).Returns(LinkUpdatePriority.High); - await CreateService().ProcessAsync(update, CancellationToken.None); + await CreateService().ProcessAsync(update, _ack, CancellationToken.None); - _groupingBuffer.Received(1).Add(42L, Arg.Any()); - _groupingBuffer.Received(1).Add(99L, Arg.Any()); + _groupingBuffer.Received(1).Add(42L, Arg.Any(), _ack); + _groupingBuffer.Received(1).Add(99L, Arg.Any(), _ack); } [Fact] @@ -125,11 +127,12 @@ public async Task ProcessAsync_WhenUpdatePasses_BufferedUpdateHasCorrectPriority .Returns("text"); _prioritizer.Prioritize(Arg.Any()).Returns(LinkUpdatePriority.High); - await CreateService().ProcessAsync(BuildUpdate(), CancellationToken.None); + await CreateService().ProcessAsync(BuildUpdate(), _ack, CancellationToken.None); _groupingBuffer.Received(1).Add( Arg.Any(), - Arg.Is(p => p.Priority == LinkUpdatePriority.High)); + Arg.Is(p => p.Priority == LinkUpdatePriority.High), + Arg.Any()); } [Fact] @@ -140,10 +143,61 @@ public async Task ProcessAsync_WhenUpdatePasses_BufferedUpdateHasSummarizedDescr .Returns("summarized text"); _prioritizer.Prioritize(Arg.Any()).Returns(LinkUpdatePriority.Medium); - await CreateService().ProcessAsync(BuildUpdate(), CancellationToken.None); + await CreateService().ProcessAsync(BuildUpdate(), _ack, CancellationToken.None); _groupingBuffer.Received(1).Add( Arg.Any(), - Arg.Is(p => p.Description == "summarized text")); + Arg.Is(p => p.Description == "summarized text"), + Arg.Any()); } -} \ No newline at end of file + + [Fact] + public async Task ProcessAsync_WhenSystemReport_PublishesDirectlyWithoutFilterSummarizationAndGrouping() + { + var update = new LinkUpdate + { + Id = 0, + Url = new Uri("https://github.com/user/repo"), + Description = "Не удалось проверить часть ссылок: spam", + TgChatIds = [42], + Kind = LinkUpdateKind.SystemReport + }; + + // Стоп-слово в тексте отчета не должно приводить к его отбрасыванию. + _filter.ShouldFilter(Arg.Any()).Returns(true); + + await CreateService().ProcessAsync(update, _ack, CancellationToken.None); + + await _publisher.Received(1).PublishAsync( + Arg.Is(p => + p.Kind == LinkUpdateKind.SystemReport && + p.Description == "Не удалось проверить часть ссылок: spam" && + p.TgChatIds.Count == 1 && + p.TgChatIds[0] == 42L), + Arg.Any()); + + _filter.DidNotReceive().ShouldFilter(Arg.Any()); + await _summarizer.DidNotReceive().SummarizeAsync(Arg.Any(), Arg.Any()); + _groupingBuffer.DidNotReceive().Add(Arg.Any(), Arg.Any(), Arg.Any()); + } + + [Fact] + public async Task ProcessAsync_WhenSystemReportPublishingFails_PropagatesException() + { + var update = new LinkUpdate + { + Id = 0, + Url = new Uri("https://github.com/user/repo"), + Description = "Не удалось проверить часть ссылок", + TgChatIds = [42], + Kind = LinkUpdateKind.SystemReport + }; + + _publisher + .PublishAsync(Arg.Any(), Arg.Any()) + .Returns(Task.FromException(new InvalidOperationException("Kafka down"))); + + await Assert.ThrowsAsync( + () => CreateService().ProcessAsync(update, _ack, CancellationToken.None)); + } +} diff --git a/src/LinkTracker.Tests/AiAgent/Unit/Infrastructure/Clients/Kafka/KafkaOffsetTrackerTests.cs b/src/LinkTracker.Tests/AiAgent/Unit/Infrastructure/Clients/Kafka/KafkaOffsetTrackerTests.cs new file mode 100644 index 0000000..891a4ed --- /dev/null +++ b/src/LinkTracker.Tests/AiAgent/Unit/Infrastructure/Clients/Kafka/KafkaOffsetTrackerTests.cs @@ -0,0 +1,116 @@ +using Confluent.Kafka; +using LinkTracker.AiAgent.Infrastructure.Clients.Kafka; + +namespace LinkTracker.Tests.AiAgent.Unit.Infrastructure.Clients.Kafka; + +[Trait("Module", "AiAgent")] +[Trait("Category", "Unit")] +public sealed class KafkaOffsetTrackerTests +{ + private const string Topic = "link.raw-updates"; + + [Fact] + public void TakeCommittableOffsets_WhenMessageNotReleased_ReturnsNothing() + { + var tracker = new KafkaOffsetTracker(); + + tracker.Track(CreateResult(0)); + + Assert.Empty(tracker.TakeCommittableOffsets()); + } + + [Fact] + public void TakeCommittableOffsets_WhenMessageReleased_ReturnsNextOffset() + { + var tracker = new KafkaOffsetTracker(); + + tracker.Track(CreateResult(7)).Release(); + + var offset = Assert.Single(tracker.TakeCommittableOffsets()); + + Assert.Equal(8, offset.Offset.Value); + Assert.Equal(Topic, offset.Topic); + } + + [Fact] + public void TakeCommittableOffsets_WhenBufferStillHoldsMessage_ReturnsNothing() + { + var tracker = new KafkaOffsetTracker(); + var ack = tracker.Track(CreateResult(0)); + + // Копия обновления лежит в буфере группировки. + ack.Retain(); + ack.Release(); + + Assert.Empty(tracker.TakeCommittableOffsets()); + + ack.Release(); + + Assert.Equal(1, Assert.Single(tracker.TakeCommittableOffsets()).Offset.Value); + } + + [Fact] + public void TakeCommittableOffsets_WhenEarlierMessageStillInFlight_StopsBeforeIt() + { + var tracker = new KafkaOffsetTracker(); + + var first = tracker.Track(CreateResult(0)); + tracker.Track(CreateResult(1)).Release(); + tracker.Track(CreateResult(2)).Release(); + + Assert.Empty(tracker.TakeCommittableOffsets()); + + first.Release(); + + Assert.Equal(3, Assert.Single(tracker.TakeCommittableOffsets()).Offset.Value); + } + + [Fact] + public void TakeCommittableOffsets_WhenNoNewProgress_ReturnsNothingOnSecondCall() + { + var tracker = new KafkaOffsetTracker(); + + tracker.Track(CreateResult(0)).Release(); + + Assert.Single(tracker.TakeCommittableOffsets()); + Assert.Empty(tracker.TakeCommittableOffsets()); + } + + [Fact] + public void TakeCommittableOffsets_TracksPartitionsIndependently() + { + var tracker = new KafkaOffsetTracker(); + + tracker.Track(CreateResult(5, partition: 0)); + tracker.Track(CreateResult(9, partition: 1)).Release(); + + var offset = Assert.Single(tracker.TakeCommittableOffsets()); + + Assert.Equal(1, offset.Partition.Value); + Assert.Equal(10, offset.Offset.Value); + } + + [Fact] + public void TakeCommittableOffsets_AfterRevoke_IgnoresStaleAcks() + { + var tracker = new KafkaOffsetTracker(); + var ack = tracker.Track(CreateResult(3)); + + tracker.Forget([new TopicPartition(Topic, new Partition(0))]); + + ack.Release(); + + Assert.Empty(tracker.TakeCommittableOffsets()); + } + + private static ConsumeResult CreateResult(long offset, int partition = 0) + { + return new ConsumeResult + { + Topic = Topic, + Partition = new Partition(partition), + Offset = new Offset(offset), + Message = new Message { Key = "key", Value = [] } + }; + } +} diff --git a/src/LinkTracker.Tests/AiAgent/Unit/Infrastructure/Clients/Kafka/RawUpdatesKafkaMessageHandlerTests.cs b/src/LinkTracker.Tests/AiAgent/Unit/Infrastructure/Clients/Kafka/RawUpdatesKafkaMessageHandlerTests.cs index 49c6b16..7ee6915 100644 --- a/src/LinkTracker.Tests/AiAgent/Unit/Infrastructure/Clients/Kafka/RawUpdatesKafkaMessageHandlerTests.cs +++ b/src/LinkTracker.Tests/AiAgent/Unit/Infrastructure/Clients/Kafka/RawUpdatesKafkaMessageHandlerTests.cs @@ -25,12 +25,13 @@ public async Task HandleAsync_WhenMessageIsValid_ProcessesAndReturnsTrue() var sut = CreateSut(processingService, deadLetterPublisher, 3); - var result = await sut.HandleAsync(CreateValidConsumeResult(), CancellationToken.None); + var result = await sut.HandleAsync(CreateValidConsumeResult(), Substitute.For(), CancellationToken.None); Assert.True(result); await processingService.Received(1).ProcessAsync( Arg.Is(u => u.Id == 42), + Arg.Any(), Arg.Any()); await deadLetterPublisher.DidNotReceive().PublishAsync( @@ -50,12 +51,13 @@ public async Task HandleAsync_WhenMessageIsMalformed_PublishesToDeadLetterAndRet var message = CreateConsumeResult("{ invalid json"); - var result = await sut.HandleAsync(message, CancellationToken.None); + var result = await sut.HandleAsync(message, Substitute.For(), CancellationToken.None); Assert.True(result); await processingService.DidNotReceive().ProcessAsync( Arg.Any(), + Arg.Any(), Arg.Any()); await deadLetterPublisher.Received(1).PublishAsync( @@ -72,19 +74,20 @@ public async Task HandleAsync_WhenProcessingFails_RetriesThenPublishesToDeadLett var deadLetterPublisher = Substitute.For(); processingService - .ProcessAsync(Arg.Any(), Arg.Any()) + .ProcessAsync(Arg.Any(), Arg.Any(), Arg.Any()) .Returns(Task.FromException(new InvalidOperationException("YandexAi failed"))); var sut = CreateSut(processingService, deadLetterPublisher, 3); var message = CreateValidConsumeResult(); - var result = await sut.HandleAsync(message, CancellationToken.None); + var result = await sut.HandleAsync(message, Substitute.For(), CancellationToken.None); Assert.True(result); await processingService.Received(3).ProcessAsync( Arg.Any(), + Arg.Any(), Arg.Any()); await deadLetterPublisher.Received(1).PublishAsync( @@ -111,7 +114,7 @@ public async Task HandleAsync_WhenDeadLetterPublishingFails_ReturnsFalseAndIncre var sut = CreateSut(processingService, deadLetterPublisher, 3, metrics: metrics); - var result = await sut.HandleAsync(CreateConsumeResult("{ invalid json"), CancellationToken.None); + var result = await sut.HandleAsync(CreateConsumeResult("{ invalid json"), Substitute.For(), CancellationToken.None); Assert.False(result); @@ -128,7 +131,7 @@ public async Task HandleAsync_WhenCancellationRequested_ThrowsOperationCanceledE using var cts = new CancellationTokenSource(); processingService - .ProcessAsync(Arg.Any(), Arg.Any()) + .ProcessAsync(Arg.Any(), Arg.Any(), Arg.Any()) .Returns(_ => Task.FromCanceled(cts.Token)); var sut = CreateSut(processingService, deadLetterPublisher, 3); @@ -136,7 +139,7 @@ public async Task HandleAsync_WhenCancellationRequested_ThrowsOperationCanceledE await cts.CancelAsync(); await Assert.ThrowsAnyAsync(() => - sut.HandleAsync(CreateValidConsumeResult(), cts.Token)); + sut.HandleAsync(CreateValidConsumeResult(), Substitute.For(), cts.Token)); await deadLetterPublisher.DidNotReceive().PublishAsync( Arg.Any>(), diff --git a/src/LinkTracker.Tests/AiAgent/Unit/Infrastructure/Services/GroupingFlushJobTests.cs b/src/LinkTracker.Tests/AiAgent/Unit/Infrastructure/Services/GroupingFlushJobTests.cs new file mode 100644 index 0000000..52e3879 --- /dev/null +++ b/src/LinkTracker.Tests/AiAgent/Unit/Infrastructure/Services/GroupingFlushJobTests.cs @@ -0,0 +1,108 @@ +using LinkTracker.AiAgent.Application.Abstractions; +using LinkTracker.AiAgent.Infrastructure.Configuration.AiAgent; +using LinkTracker.AiAgent.Infrastructure.Services; +using LinkTracker.Shared.Contracts.AiAgent; +using Microsoft.Extensions.Logging.Abstractions; +using Microsoft.Extensions.Options; +using NSubstitute; + +namespace LinkTracker.Tests.AiAgent.Unit.Infrastructure.Services; + +[Trait("Module", "AiAgent")] +[Trait("Category", "Unit")] +public sealed class GroupingFlushJobTests +{ + private readonly IGroupingBuffer _buffer = Substitute.For(); + private readonly IProcessedUpdatePublisher _publisher = Substitute.For(); + + [Fact] + public async Task StopAsync_PublishesWindowsThatDidNotCloseYet() + { + var ack = Substitute.For(); + _buffer.Flush(true).Returns([new GroupingBucket(42, [new BufferedLinkUpdate(BuildUpdate(), ack)])]); + + await CreateJob().StopAsync(CancellationToken.None); + + await _publisher.Received(1).PublishAsync( + Arg.Is(u => u.Id == 1), + Arg.Any()); + } + + [Fact] + public async Task StopAsync_WhenPublished_ReleasesSourceMessages() + { + var ack = Substitute.For(); + _buffer.Flush(true).Returns([new GroupingBucket(42, [new BufferedLinkUpdate(BuildUpdate(), ack)])]); + + await CreateJob().StopAsync(CancellationToken.None); + + ack.Received(1).Release(); + _buffer.DidNotReceive().Requeue(Arg.Any()); + } + + [Fact] + public async Task StopAsync_WhenPublishingFails_KeepsMessageUnacknowledgedAndRequeuesWindow() + { + var ack = Substitute.For(); + var bucket = new GroupingBucket(42, [new BufferedLinkUpdate(BuildUpdate(), ack)]); + + _buffer.Flush(true).Returns([bucket]); + + _publisher + .PublishAsync(Arg.Any(), Arg.Any()) + .Returns(Task.FromException(new InvalidOperationException("Kafka down"))); + + await CreateJob().StopAsync(CancellationToken.None); + + ack.DidNotReceive().Release(); + _buffer.Received(1).Requeue(bucket); + } + + [Fact] + public async Task StopAsync_PublishesHighPriorityWindowsFirst() + { + var published = new List(); + + _publisher + .PublishAsync(Arg.Any(), Arg.Any()) + .Returns(call => + { + published.Add(call.Arg().Id); + return Task.CompletedTask; + }); + + _buffer.Flush(true).Returns( + [ + new GroupingBucket(1, [new BufferedLinkUpdate(BuildUpdate(10, LinkUpdatePriority.Low), Substitute.For())]), + new GroupingBucket(2, [new BufferedLinkUpdate(BuildUpdate(20, LinkUpdatePriority.High), Substitute.For())]) + ]); + + await CreateJob().StopAsync(CancellationToken.None); + + Assert.Equal([20L, 10L], published); + } + + private GroupingFlushJob CreateJob() + { + return new GroupingFlushJob( + _buffer, + new WindowLinkUpdateGrouper(), + _publisher, + Options.Create(new AiAgentOptions()), + NullLogger.Instance); + } + + private static ProcessedLinkUpdate BuildUpdate( + long id = 1, + LinkUpdatePriority priority = LinkUpdatePriority.Medium) + { + return new ProcessedLinkUpdate + { + Id = id, + Url = new Uri("https://github.com/user/repo"), + Description = "Обновление", + TgChatIds = [42], + Priority = priority + }; + } +} diff --git a/src/LinkTracker.Tests/AiAgent/Unit/Infrastructure/Services/TimeWindowGroupingBufferTests.cs b/src/LinkTracker.Tests/AiAgent/Unit/Infrastructure/Services/TimeWindowGroupingBufferTests.cs new file mode 100644 index 0000000..5ae4b39 --- /dev/null +++ b/src/LinkTracker.Tests/AiAgent/Unit/Infrastructure/Services/TimeWindowGroupingBufferTests.cs @@ -0,0 +1,128 @@ +using LinkTracker.AiAgent.Application.Abstractions; +using LinkTracker.AiAgent.Infrastructure.Configuration.AiAgent; +using LinkTracker.AiAgent.Infrastructure.Services; +using LinkTracker.Shared.Contracts.AiAgent; +using Microsoft.Extensions.Options; + +namespace LinkTracker.Tests.AiAgent.Unit.Infrastructure.Services; + +[Trait("Module", "AiAgent")] +[Trait("Category", "Unit")] +public sealed class TimeWindowGroupingBufferTests +{ + [Fact] + public void Add_HoldsMessageUntilBufferIsFlushed() + { + var buffer = CreateBuffer(); + var ack = new CountingAck(); + + buffer.Add(42, BuildUpdate(), ack); + + Assert.False(ack.IsCompleted); + } + + [Fact] + public void Flush_WhenWindowIsStillOpen_ReturnsNothing() + { + var buffer = CreateBuffer(); + + buffer.Add(42, BuildUpdate(), new CountingAck()); + + Assert.Empty(buffer.Flush()); + } + + [Fact] + public void Flush_WhenForced_ReturnsOpenWindows() + { + var buffer = CreateBuffer(); + var ack = new CountingAck(); + + buffer.Add(42, BuildUpdate(1), ack); + buffer.Add(42, BuildUpdate(2), ack); + + var bucket = Assert.Single(buffer.Flush(true)); + + Assert.Equal(42, bucket.ChatId); + Assert.Equal([1L, 2L], bucket.Updates.Select(x => x.Update.Id)); + Assert.All(bucket.Updates, x => Assert.Same(ack, x.Ack)); + } + + [Fact] + public void Flush_WhenWindowExpired_ReturnsBucket() + { + var buffer = CreateBuffer(1); + + buffer.Add(42, BuildUpdate(), new CountingAck()); + + Thread.Sleep(20); + + Assert.Single(buffer.Flush()); + } + + [Fact] + public void Requeue_ReturnsUpdatesToBufferWithoutCompletingMessage() + { + var buffer = CreateBuffer(); + var ack = new CountingAck(); + + buffer.Add(42, BuildUpdate(), ack); + + var bucket = Assert.Single(buffer.Flush(true)); + + buffer.Requeue(bucket); + + Assert.False(ack.IsCompleted); + + var requeued = Assert.Single(buffer.Flush(true)); + + Assert.Equal(bucket.Updates[0].Update.Id, Assert.Single(requeued.Updates).Update.Id); + + foreach (var buffered in requeued.Updates) + { + buffered.Ack.Release(); + } + + Assert.False(ack.IsCompleted); + + ack.Release(); + + Assert.True(ack.IsCompleted); + } + + private static TimeWindowGroupingBuffer CreateBuffer(int windowMs = 60000) + { + return new TimeWindowGroupingBuffer( + Options.Create(new AiAgentOptions { Grouping = new GroupingOptions { WindowMs = windowMs } })); + } + + private static ProcessedLinkUpdate BuildUpdate(long id = 1) + { + return new ProcessedLinkUpdate + { + Id = id, + Url = new Uri("https://github.com/user/repo"), + Description = "Обновление", + TgChatIds = [42] + }; + } + + private sealed class CountingAck : IMessageAck + { + private int _references = 1; + + public bool IsCompleted { get; private set; } + + public void Retain() + { + Interlocked.Increment(ref _references); + } + + public void Release() + { + if (Interlocked.Decrement(ref _references) == 0) + { + IsCompleted = true; + } + } + } +} diff --git a/src/LinkTracker.Tests/Bot/Integration/Infrastructure/Clients/Kafka/LinkUpdatesKafkaConsumerIntegrationTests.cs b/src/LinkTracker.Tests/Bot/Integration/Infrastructure/Clients/Kafka/LinkUpdatesKafkaConsumerIntegrationTests.cs index f527c6d..2de4d65 100644 --- a/src/LinkTracker.Tests/Bot/Integration/Infrastructure/Clients/Kafka/LinkUpdatesKafkaConsumerIntegrationTests.cs +++ b/src/LinkTracker.Tests/Bot/Integration/Infrastructure/Clients/Kafka/LinkUpdatesKafkaConsumerIntegrationTests.cs @@ -15,6 +15,7 @@ using LinkTracker.Scrapper.Infrastructure.Configuration.Kafka; using LinkTracker.Scrapper.Infrastructure.Kafka.Serialization; using LinkTracker.Scrapper.Infrastructure.Telemetry; +using LinkTracker.Shared.Contracts.AiAgent; using LinkTracker.Shared.Contracts.Bot; using LinkTracker.Shared.Infrastructure; using LinkTracker.Tests.Bot.Integration.Kafka; @@ -252,7 +253,16 @@ public async Task ScrapperKafkaClient_WhenUsesAvro_BotConsumerNotifiesUser() try { await scrapperKafkaClient.Client.SendUpdateAsync( - new LinkUpdate { Id = 42, Url = new Uri("https://github.com/user/repo"), Description = "Repository updated with Avro", TgChatIds = [123] }); + new LinkUpdate + { + Id = 42, + Url = new Uri("https://github.com/user/repo"), + Description = "Repository updated with Avro", + Author = "octocat", + TgChatIds = [123], + Priority = LinkUpdatePriority.High, + Kind = LinkUpdateKind.SystemReport + }); await WaitUntilAsync(async () => { @@ -263,6 +273,9 @@ await notifier.Received(1).NotifyAsync( update.Id == 42 && update.Url == new Uri("https://github.com/user/repo") && update.Description == "Repository updated with Avro" + && update.Author == "octocat" + && update.Priority == LinkUpdatePriority.High + && update.Kind == LinkUpdateKind.SystemReport && update.TgChatIds.SequenceEqual(new[] { 123L })), Arg.Any()); diff --git a/src/LinkTracker.Tests/Bot/Unit/Presentation/Telegram/Notifications/LinkUpdateNotifierTests.cs b/src/LinkTracker.Tests/Bot/Unit/Presentation/Telegram/Notifications/LinkUpdateNotifierTests.cs index c94eb44..88be82a 100644 --- a/src/LinkTracker.Tests/Bot/Unit/Presentation/Telegram/Notifications/LinkUpdateNotifierTests.cs +++ b/src/LinkTracker.Tests/Bot/Unit/Presentation/Telegram/Notifications/LinkUpdateNotifierTests.cs @@ -1,6 +1,5 @@ using LinkTracker.Bot.Application.Telemetry.Abstractions; using LinkTracker.Bot.Presentation.Telegram.Notifications; -using LinkTracker.Shared.Constants; using LinkTracker.Shared.Contracts.AiAgent; using LinkTracker.Shared.Contracts.Bot; using NSubstitute; @@ -52,7 +51,7 @@ await sut.NotifyAsync( } [Fact] - public async Task NotifyAsync_WhenSystemReport_StripsMarkerAndSkipsPriorityHeader() + public async Task NotifyAsync_WhenSystemReport_SendsDescriptionAsIsWithoutHeader() { var (botClient, sentTexts) = CreateBotClient(); var sut = new LinkUpdateNotifier(botClient, Substitute.For()); @@ -62,9 +61,10 @@ await sut.NotifyAsync( { Id = 0, Url = Url, - Description = $"{SystemMessageMarkers.FailedLinkReport}Не удалось проверить часть ссылок", + Description = "Не удалось проверить часть ссылок", TgChatIds = [42], - Priority = LinkUpdatePriority.High + Priority = LinkUpdatePriority.High, + Kind = LinkUpdateKind.SystemReport }, CancellationToken.None); diff --git a/src/LinkTracker.Tests/Scrapper/Unit/Infrastructure/Quartz/Jobs/LinkUpdatesJobTests.cs b/src/LinkTracker.Tests/Scrapper/Unit/Infrastructure/Quartz/Jobs/LinkUpdatesJobTests.cs index 99d9dfc..690bd90 100644 --- a/src/LinkTracker.Tests/Scrapper/Unit/Infrastructure/Quartz/Jobs/LinkUpdatesJobTests.cs +++ b/src/LinkTracker.Tests/Scrapper/Unit/Infrastructure/Quartz/Jobs/LinkUpdatesJobTests.cs @@ -7,7 +7,6 @@ using LinkTracker.Scrapper.Infrastructure.Quartz.Jobs; using LinkTracker.Scrapper.Infrastructure.Telemetry; using LinkTracker.Scrapper.Storage.Abstractions.Models; -using LinkTracker.Shared.Constants; using LinkTracker.Shared.Contracts.Bot; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; @@ -161,7 +160,7 @@ await botClient.Received(1) Arg.Is(x => x.TgChatIds.Count == 1 && x.TgChatIds.Contains(1001L) && - x.Description.Contains(SystemMessageMarkers.FailedLinkReport) && + x.Kind == LinkUpdateKind.SystemReport && x.Description.Contains(subscription.Url.ToString())), Arg.Any()); } @@ -509,7 +508,7 @@ await botClient.Received(1) x.Url == successfulSubscription.Url && x.TgChatIds.Count == 1 && x.TgChatIds.Contains(1002L) && - !x.Description.Contains(SystemMessageMarkers.FailedLinkReport) && + x.Kind == LinkUpdateKind.Content && x.Description.Contains("Заголовок: Issue")), Arg.Any()); @@ -518,7 +517,7 @@ await botClient.Received(1) Arg.Is(x => x.TgChatIds.Count == 1 && x.TgChatIds.Contains(1001L) && - x.Description.Contains(SystemMessageMarkers.FailedLinkReport) && + x.Kind == LinkUpdateKind.SystemReport && x.Description.Contains(failedSubscription.Url.ToString())), Arg.Any()); @@ -760,7 +759,7 @@ await botClient.Received(1) Arg.Is(x => x.TgChatIds.Count == 1 && x.TgChatIds.Contains(1001L) && - x.Description.Contains(SystemMessageMarkers.FailedLinkReport) && + x.Kind == LinkUpdateKind.SystemReport && x.Description.Contains(firstFailedSubscription.Url.ToString()) && !x.Description.Contains(secondFailedSubscription.Url.ToString())), Arg.Any()); @@ -770,7 +769,7 @@ await botClient.Received(1) Arg.Is(x => x.TgChatIds.Count == 1 && x.TgChatIds.Contains(1002L) && - x.Description.Contains(SystemMessageMarkers.FailedLinkReport) && + x.Kind == LinkUpdateKind.SystemReport && x.Description.Contains(firstFailedSubscription.Url.ToString()) && x.Description.Contains(secondFailedSubscription.Url.ToString())), Arg.Any()); @@ -981,7 +980,7 @@ await botClient.Received(1).SendUpdateAsync( Arg.Is(x => x.Id == 0 && x.TgChatIds.SequenceEqual(new[] { 1001L }) && - x.Description.Contains(SystemMessageMarkers.FailedLinkReport) && + x.Kind == LinkUpdateKind.SystemReport && x.Description.Contains(subscription.Url.ToString())), Arg.Any()); } diff --git a/src/LinkTracker.Tests/Shared/Unit/Contracts/LinkUpdateAvroSchemaTests.cs b/src/LinkTracker.Tests/Shared/Unit/Contracts/LinkUpdateAvroSchemaTests.cs new file mode 100644 index 0000000..7904f5f --- /dev/null +++ b/src/LinkTracker.Tests/Shared/Unit/Contracts/LinkUpdateAvroSchemaTests.cs @@ -0,0 +1,35 @@ +using Avro; +using LinkTracker.Shared.Contracts.Bot; + +namespace LinkTracker.Tests.Shared.Unit.Contracts; + +[Trait("Module", "Shared")] +[Trait("Category", "Unit")] +public sealed class LinkUpdateAvroSchemaTests +{ + // Avro-ветка обязана переносить тот же контракт, что и JSON, иначе выбор + // сериализации молча меняет данные (так терялся author). + [Fact] + public void Schema_CoversEveryLinkUpdateProperty() + { + var schema = (RecordSchema)Avro.Schema.Parse(LinkUpdateAvroSchema.Value); + + var contractFields = typeof(LinkUpdate) + .GetProperties() + .Select(property => char.ToLowerInvariant(property.Name[0]) + property.Name[1..]) + .Order(); + + Assert.Equal(contractFields, schema.Fields.Select(field => field.Name).Order()); + } + + [Theory] + [InlineData("author")] + [InlineData("priority")] + [InlineData("kind")] + public void Schema_FieldsAddedAfterFirstVersion_HaveDefaults(string field) + { + var schema = (RecordSchema)Avro.Schema.Parse(LinkUpdateAvroSchema.Value); + + Assert.NotNull(schema[field].DefaultValue); + } +} diff --git a/src/LinkTracker.Tests/Shared/Unit/Contracts/ProcessedUpdateWireCompatibilityTests.cs b/src/LinkTracker.Tests/Shared/Unit/Contracts/ProcessedUpdateWireCompatibilityTests.cs index 9d5074d..2b1e36d 100644 --- a/src/LinkTracker.Tests/Shared/Unit/Contracts/ProcessedUpdateWireCompatibilityTests.cs +++ b/src/LinkTracker.Tests/Shared/Unit/Contracts/ProcessedUpdateWireCompatibilityTests.cs @@ -1,6 +1,7 @@ using LinkTracker.AiAgent.Infrastructure.Kafka.Serialization; using LinkTracker.Bot.Infrastructure.Kafka.Deserialization; using LinkTracker.Shared.Contracts.AiAgent; +using LinkTracker.Shared.Contracts.Bot; namespace LinkTracker.Tests.Shared.Unit.Contracts; @@ -42,4 +43,28 @@ public async Task BotDeserializer_ReadsPriorityPublishedByAiAgent(LinkUpdatePrio Assert.Equal(published.Description, received.Description); Assert.Equal(published.TgChatIds, received.TgChatIds); } + + [Theory] + [InlineData(LinkUpdateKind.Content)] + [InlineData(LinkUpdateKind.SystemReport)] + public async Task BotDeserializer_ReadsKindPublishedByAiAgent(LinkUpdateKind kind) + { + var published = new ProcessedLinkUpdate + { + Id = 0, + Url = new Uri("https://github.com/dotnet/runtime"), + Description = "Не удалось проверить часть ссылок", + TgChatIds = [1001], + Kind = kind + }; + + var payload = await new JsonProcessedLinkUpdateKafkaSerializer() + .SerializeAsync(published, Topic, CancellationToken.None); + + var received = await new JsonLinkUpdateKafkaDeserializer() + .DeserializeAsync(payload, Topic, CancellationToken.None); + + Assert.NotNull(received); + Assert.Equal(kind, received!.Kind); + } }