diff --git a/OBSERVABILITY.md b/OBSERVABILITY.md index a47c4e5..a7431dc 100644 --- a/OBSERVABILITY.md +++ b/OBSERVABILITY.md @@ -2,8 +2,19 @@ ## Архитектура сбора метрик -Метрики приложений построены на **OpenTelemetry** и доставляются в Prometheus -по **push-модели через Prometheus Pushgateway**. +Метрики приложений построены на **OpenTelemetry** и отдаются Prometheus +по **pull-модели**: каждый сервис публикует OpenTelemetry-метрики в формате Prometheus +на своём эндпоинте `/metrics`, а Prometheus скрейпит их напрямую. + +- Scrapper: `http://scrapper:8081/metrics` +- Bot: `http://bot:8011/metrics` (отдельный Kestrel-эндпоинт, наружу не публикуется) + +Эндпоинт `/metrics` исключён из общего rate limiter — иначе скрейп раз в 15 секунд +конкурировал бы за лимит с прикладным трафиком и метрики выглядели бы «пропавшими». + +Имена серий, которые видит Prometheus, зафиксированы тестом +`MetricsEndpointTests` — экспортёр переименовывает инструменты, и панели Grafana +завязаны именно на итоговые имена. ## Метрики приложений @@ -43,19 +54,16 @@ > Histogram-метрики дают в Prometheus три серии: `_bucket`, `_sum`, `_count`. -## Конфигурация push +## Конфигурация сбора -Настройки секции `Telemetry:Pushgateway` (переопределяются переменными окружения -с разделителем `__`): +Настраивается на стороне Prometheus в `monitoring/prometheus.yml` — приложениям +никакой конфигурации телеметрии не требуется. Лейбл `job` берётся из имени scrape-job +(`scrapper` / `bot`), `instance` — из адреса цели. -| Ключ | Env-переменная | По умолчанию | +| Job | Цель | Интервал | |---|---|---| -| `Endpoint` | `Telemetry__Pushgateway__Endpoint` | `http://pushgateway:9091/metrics` | -| `Job` | `Telemetry__Pushgateway__Job` | `scrapper` / `bot` | -| `Enabled` | `Telemetry__Pushgateway__Enabled` | `true` | -| `IntervalMilliseconds` | `Telemetry__Pushgateway__IntervalMilliseconds` | `5000` | - -`Instance` берётся из переменной `HOSTNAME` (в контейнере — id контейнера) либо из имени хоста. +| `scrapper` | `scrapper:8081` | 15s (`global.scrape_interval`) | +| `bot` | `bot:8011` | 15s (`global.scrape_interval`) | ## Grafana @@ -75,7 +83,7 @@ ## Запуск мониторинга ```bash -# Поднять инфраструктуру (Kafka, Postgres, Valkey, Pushgateway, Prometheus, Grafana) +# Поднять инфраструктуру (Kafka, Postgres, Valkey, Prometheus, Grafana) docker compose -f docker-compose.yml up -d # Поднять приложения @@ -84,5 +92,4 @@ docker compose -f docker-compose.yml -f docker-compose.apps.yml up -d После запуска: - Prometheus: http://localhost:9090 -- Pushgateway: http://localhost:9091 - Grafana: http://localhost:3000 (admin / admin) \ No newline at end of file diff --git a/docker-compose.apps.yml b/docker-compose.apps.yml index cc75541..80c845c 100644 --- a/docker-compose.apps.yml +++ b/docker-compose.apps.yml @@ -15,8 +15,6 @@ services: condition: service_started valkey-cluster-init: condition: service_completed_successfully - pushgateway: - condition: service_started ports: - "8081:8081" - "8082:8082" @@ -24,36 +22,30 @@ services: - src/LinkTracker.Scrapper.Api/.env environment: ASPNETCORE_ENVIRONMENT: "Docker" - Telemetry__Pushgateway__Endpoint: "http://pushgateway:9091/metrics" - Telemetry__Pushgateway__Job: "scrapper" - bot: - build: - context: . - dockerfile: src/LinkTracker.Bot.Api/Dockerfile - network: host - container_name: linktracker-bot - restart: unless-stopped - depends_on: - scrapper: - condition: service_started - kafka-init: - condition: service_completed_successfully - schema-registry: - condition: service_started - pushgateway: - condition: service_started - ports: - - "8091:8091" - - "8092:8092" - - "8011:8011" - env_file: - - src/LinkTracker.Bot.Api/.env - environment: - ASPNETCORE_ENVIRONMENT: "Docker" - Telemetry__Pushgateway__Endpoint: "http://pushgateway:9091/metrics" - Telemetry__Pushgateway__Job: "bot" - + bot: + build: + context: . + dockerfile: src/LinkTracker.Bot.Api/Dockerfile + network: host + container_name: linktracker-bot + restart: unless-stopped + depends_on: + scrapper: + condition: service_started + kafka-init: + condition: service_completed_successfully + schema-registry: + condition: service_started + ports: + - "8091:8091" + - "8092:8092" + - "8011:8011" + env_file: + - src/LinkTracker.Bot.Api/.env + environment: + ASPNETCORE_ENVIRONMENT: "Docker" + aiagent: build: context: . diff --git a/docker-compose.yml b/docker-compose.yml index a5ccac9..8801d60 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -438,13 +438,6 @@ services: depends_on: - prometheus - pushgateway: - image: prom/pushgateway:latest - container_name: linktracker-pushgateway - restart: unless-stopped - ports: - - "9091:9091" - prometheus: image: prom/prometheus:latest container_name: linktracker-prometheus @@ -457,8 +450,6 @@ services: command: - "--config.file=/etc/prometheus/prometheus.yml" - "--storage.tsdb.path=/prometheus" - depends_on: - - pushgateway volumes: postgres_data: diff --git a/migrations/004_outbox_lease.sql b/migrations/004_outbox_lease.sql new file mode 100644 index 0000000..f704744 --- /dev/null +++ b/migrations/004_outbox_lease.sql @@ -0,0 +1,2 @@ +ALTER TABLE outbox_messages + ADD COLUMN IF NOT EXISTS locked_until TIMESTAMPTZ NULL; diff --git a/migrations/005_drop_filters.sql b/migrations/005_drop_filters.sql new file mode 100644 index 0000000..a6bf276 --- /dev/null +++ b/migrations/005_drop_filters.sql @@ -0,0 +1,2 @@ +DROP TABLE IF EXISTS subscription_filters; +DROP TABLE IF EXISTS filters; \ No newline at end of file diff --git a/monitoring/prometheus.yml b/monitoring/prometheus.yml index d20b951..d40a5be 100644 --- a/monitoring/prometheus.yml +++ b/monitoring/prometheus.yml @@ -9,8 +9,3 @@ scrape_configs: - job_name: bot static_configs: - targets: ["bot:8011"] - - - job_name: pushgateway - honor_labels: true - static_configs: - - targets: [ "pushgateway:9091" ] \ 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 9fa0371..1ff60c7 100644 --- a/src/LinkTracker.AiAgent.Infrastructure/Services/GroupingFlushJob.cs +++ b/src/LinkTracker.AiAgent.Infrastructure/Services/GroupingFlushJob.cs @@ -26,26 +26,26 @@ protected override async Task ExecuteAsync(CancellationToken stoppingToken) private async Task FlushAsync(CancellationToken ct) { - var flushed = buffer.Flush(); + var pending = buffer.Flush() + .SelectMany(entry => grouper + .Group(entry.Updates) + .Select(update => (entry.ChatId, Update: update))) + .OrderByDescending(x => x.Update.Priority) + .ToArray(); - foreach (var (chatId, updates) in flushed) + foreach (var (chatId, update) in pending) { - var grouped = grouper.Group(updates); - - foreach (var update in grouped) + try { - try - { - await publisher.PublishAsync(update, ct); + await publisher.PublishAsync(update, ct); - logger.LogInformation( - "Группа опубликована. ChatId={ChatId}, Count={Count}, Priority={Priority}", - chatId, updates.Count, update.Priority); - } - catch (Exception ex) when (ex is not OperationCanceledException) - { - logger.LogError(ex, "Ошибка публикации сгруппированного обновления. ChatId={ChatId}", chatId); - } + logger.LogInformation( + "Группа опубликована. ChatId={ChatId}, UpdateId={UpdateId}, Priority={Priority}", + chatId, update.Id, update.Priority); + } + catch (Exception ex) when (ex is not OperationCanceledException) + { + logger.LogError(ex, "Ошибка публикации сгруппированного обновления. ChatId={ChatId}", chatId); } } } diff --git a/src/LinkTracker.Bot.Api/Program.cs b/src/LinkTracker.Bot.Api/Program.cs index 98115ff..badfb56 100644 --- a/src/LinkTracker.Bot.Api/Program.cs +++ b/src/LinkTracker.Bot.Api/Program.cs @@ -12,7 +12,7 @@ using LinkTracker.EnvReader; using LinkTracker.Shared.Infrastructure.RateLimiting; using LinkTracker.Shared.Infrastructure.Resilience; -using Prometheus; +using LinkTracker.Shared.Infrastructure.Telemetry; using Serilog; var builder = WebApplication.CreateBuilder(args); @@ -43,7 +43,7 @@ try { - app.MapMetrics().RequireHost("*:8011"); + app.MapMetricsEndpoint().RequireHost("*:8011"); await app.RunAsync(); } finally diff --git a/src/LinkTracker.Bot.Api/appsettings.Docker.json b/src/LinkTracker.Bot.Api/appsettings.Docker.json index ec94803..c7c7c70 100644 --- a/src/LinkTracker.Bot.Api/appsettings.Docker.json +++ b/src/LinkTracker.Bot.Api/appsettings.Docker.json @@ -50,13 +50,5 @@ "WindowSeconds": 60, "SegmentsPerWindow": 6, "QueueLimit": 0 - }, - "Telemetry": { - "Pushgateway": { - "Enabled": true, - "Endpoint": "http://pushgateway:9091/metrics", - "Job": "bot", - "IntervalMilliseconds": 5000 - } } -} \ No newline at end of file +} diff --git a/src/LinkTracker.Bot.Api/appsettings.json b/src/LinkTracker.Bot.Api/appsettings.json index 3eb1ccb..6973521 100644 --- a/src/LinkTracker.Bot.Api/appsettings.json +++ b/src/LinkTracker.Bot.Api/appsettings.json @@ -65,13 +65,5 @@ "WindowSeconds": 60, "SegmentsPerWindow": 6, "QueueLimit": 0 - }, - "Telemetry": { - "Pushgateway": { - "Enabled": true, - "Endpoint": "http://localhost:9091/metrics", - "Job": "bot", - "IntervalMilliseconds": 5000 - } } } diff --git a/src/LinkTracker.Bot.Application/Clients/Scrapper/Contracts/Requests/AddLinkRequest.cs b/src/LinkTracker.Bot.Application/Clients/Scrapper/Contracts/Requests/AddLinkRequest.cs index 4006640..1212ae4 100644 --- a/src/LinkTracker.Bot.Application/Clients/Scrapper/Contracts/Requests/AddLinkRequest.cs +++ b/src/LinkTracker.Bot.Application/Clients/Scrapper/Contracts/Requests/AddLinkRequest.cs @@ -5,6 +5,4 @@ public sealed class AddLinkRequest public Uri? Link { get; init; } public IReadOnlyList Tags { get; init; } = []; - - public IReadOnlyList Filters { get; init; } = []; } \ No newline at end of file diff --git a/src/LinkTracker.Bot.Application/Clients/Scrapper/Contracts/Responses/LinkResponse.cs b/src/LinkTracker.Bot.Application/Clients/Scrapper/Contracts/Responses/LinkResponse.cs index 260da5b..a96a98c 100644 --- a/src/LinkTracker.Bot.Application/Clients/Scrapper/Contracts/Responses/LinkResponse.cs +++ b/src/LinkTracker.Bot.Application/Clients/Scrapper/Contracts/Responses/LinkResponse.cs @@ -7,6 +7,4 @@ public sealed class LinkResponse public Uri Url { get; init; } = default!; public IReadOnlyList Tags { get; init; } = []; - - public IReadOnlyList Filters { get; init; } = []; } \ No newline at end of file diff --git a/src/LinkTracker.Bot.Application/Clients/Scrapper/IScrapperClient.cs b/src/LinkTracker.Bot.Application/Clients/Scrapper/IScrapperClient.cs index c0d5ef4..8d21902 100644 --- a/src/LinkTracker.Bot.Application/Clients/Scrapper/IScrapperClient.cs +++ b/src/LinkTracker.Bot.Application/Clients/Scrapper/IScrapperClient.cs @@ -14,7 +14,6 @@ Task AddLinkAsync( long chatId, Uri link, IReadOnlyList tags, - IReadOnlyList filters, CancellationToken ct = default); Task RemoveLinkAsync(long chatId, Uri link, CancellationToken ct = default); diff --git a/src/LinkTracker.Bot.Application/Dialogs/Implementations/Track/Nodes/TrackConfirmNode.cs b/src/LinkTracker.Bot.Application/Dialogs/Implementations/Track/Nodes/TrackConfirmNode.cs index 92c35cc..8895b41 100644 --- a/src/LinkTracker.Bot.Application/Dialogs/Implementations/Track/Nodes/TrackConfirmNode.cs +++ b/src/LinkTracker.Bot.Application/Dialogs/Implementations/Track/Nodes/TrackConfirmNode.cs @@ -57,7 +57,7 @@ private async Task ConfirmAsync(DialogContext ctx, Cancellatio try { - await scrapperClient.AddLinkAsync(ctx.ChatId, uri, tags, [], ct); + await scrapperClient.AddLinkAsync(ctx.ChatId, uri, tags, ct); return new DialogNodeResult( $"Начал отслеживать:\n{pendingUrl}\nТеги: {tagsText}", diff --git a/src/LinkTracker.Bot.Infrastructure/Clients/Scrapper/ScrapperGrpcClient.cs b/src/LinkTracker.Bot.Infrastructure/Clients/Scrapper/ScrapperGrpcClient.cs index a58766b..f85a4aa 100644 --- a/src/LinkTracker.Bot.Infrastructure/Clients/Scrapper/ScrapperGrpcClient.cs +++ b/src/LinkTracker.Bot.Infrastructure/Clients/Scrapper/ScrapperGrpcClient.cs @@ -104,7 +104,6 @@ public async Task AddLinkAsync( long chatId, Uri link, IReadOnlyList tags, - IReadOnlyList filters, CancellationToken ct = default) { var sw = Stopwatch.StartNew(); @@ -112,11 +111,10 @@ public async Task AddLinkAsync( try { logger.LogInformation( - "Клиент gRPC Scrapper: вызов AddLink. ChatId={ChatId}, Ссылка={Link}, КоличествоТегов={TagsCount}, КоличествоФильтров={FiltersCount}", + "Клиент gRPC Scrapper: вызов AddLink. ChatId={ChatId}, Ссылка={Link}, КоличествоТегов={TagsCount}", chatId, link, - tags.Count, - filters.Count); + tags.Count); try { @@ -180,7 +178,7 @@ public async Task RemoveLinkAsync(long chatId, Uri link, Cancellat private static LinkResponse ToModel(LinkGrpcResponse response) { - return new LinkResponse { Id = response.Id, Url = new Uri(response.Url), Tags = response.Tags.ToArray(), Filters = [] }; + return new LinkResponse { Id = response.Id, Url = new Uri(response.Url), Tags = response.Tags.ToArray() }; } private static ScrapperClientException ToClientException(RpcException ex) diff --git a/src/LinkTracker.Bot.Infrastructure/Clients/Scrapper/ScrapperHttpClient.cs b/src/LinkTracker.Bot.Infrastructure/Clients/Scrapper/ScrapperHttpClient.cs index d34a41e..e558b37 100644 --- a/src/LinkTracker.Bot.Infrastructure/Clients/Scrapper/ScrapperHttpClient.cs +++ b/src/LinkTracker.Bot.Infrastructure/Clients/Scrapper/ScrapperHttpClient.cs @@ -78,7 +78,6 @@ public Task AddLinkAsync( long chatId, Uri link, IReadOnlyList tags, - IReadOnlyList filters, CancellationToken ct = default) { var sw = Stopwatch.StartNew(); @@ -90,7 +89,7 @@ public Task AddLinkAsync( { using var request = new HttpRequestMessage(HttpMethod.Post, "/links"); request.Headers.Add("Tg-Chat-Id", chatId.ToString()); - request.Content = JsonContent.Create(new AddLinkRequest { Link = link, Tags = tags, Filters = [] }); + request.Content = JsonContent.Create(new AddLinkRequest { Link = link, Tags = tags }); using var response = await httpClient.SendAsync(request, token); await EnsureSuccessStatusCodeAsync(response, token); diff --git a/src/LinkTracker.Bot.Infrastructure/Telemetry/Registration/TelemetryModule.cs b/src/LinkTracker.Bot.Infrastructure/Telemetry/Registration/TelemetryModule.cs index 19ec7d0..31193f4 100644 --- a/src/LinkTracker.Bot.Infrastructure/Telemetry/Registration/TelemetryModule.cs +++ b/src/LinkTracker.Bot.Infrastructure/Telemetry/Registration/TelemetryModule.cs @@ -13,9 +13,7 @@ public static IServiceCollection AddTelemetry( { services.AddSingleton(); - services.AddOpenTelemetryMetricsWithPushgateway( - configuration, - "bot", + services.AddOpenTelemetryMetrics( "bot", BotMetrics.MeterName); diff --git a/src/LinkTracker.Bot.Presentation/Telegram/Notifications/LinkUpdateNotifier.cs b/src/LinkTracker.Bot.Presentation/Telegram/Notifications/LinkUpdateNotifier.cs index 216a3bd..76b0e7c 100644 --- a/src/LinkTracker.Bot.Presentation/Telegram/Notifications/LinkUpdateNotifier.cs +++ b/src/LinkTracker.Bot.Presentation/Telegram/Notifications/LinkUpdateNotifier.cs @@ -1,6 +1,7 @@ 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; @@ -26,6 +27,16 @@ private static string BuildMessage(LinkUpdate update) return update.Description[SystemMessageMarkers.FailedLinkReport.Length..]; } - return $"Обновление по ссылке:\n{update.Url}\n\n{update.Description}"; + return $"{BuildHeader(update.Priority)}\n{update.Url}\n\n{update.Description}"; } -} \ No newline at end of file + + private static string BuildHeader(LinkUpdatePriority priority) + { + return priority switch + { + LinkUpdatePriority.High => "‼️ Важное обновление по ссылке:", + LinkUpdatePriority.Low => "Незначительное обновление по ссылке:", + _ => "Обновление по ссылке:" + }; + } +} diff --git a/src/LinkTracker.Scrapper.Api/Dockerfile b/src/LinkTracker.Scrapper.Api/Dockerfile index 5e35185..0396a81 100644 --- a/src/LinkTracker.Scrapper.Api/Dockerfile +++ b/src/LinkTracker.Scrapper.Api/Dockerfile @@ -9,7 +9,6 @@ COPY ["src/LinkTracker.Scrapper.Contracts/LinkTracker.Scrapper.Contracts.csproj" COPY ["src/LinkTracker.Scrapper.Infrastructure/LinkTracker.Scrapper.Infrastructure.csproj", "src/LinkTracker.Scrapper.Infrastructure/"] COPY ["src/LinkTracker.Scrapper.Presentation/LinkTracker.Scrapper.Presentation.csproj", "src/LinkTracker.Scrapper.Presentation/"] COPY ["src/LinkTracker.Scrapper.Storage.Abstractions/LinkTracker.Scrapper.Storage.Abstractions.csproj", "src/LinkTracker.Scrapper.Storage.Abstractions/"] -COPY ["src/LinkTracker.Scrapper.Storage.InMemory/LinkTracker.Scrapper.Storage.InMemory.csproj", "src/LinkTracker.Scrapper.Storage.InMemory/"] COPY ["src/LinkTracker.Scrapper.Storage.Orm/LinkTracker.Scrapper.Storage.Orm.csproj", "src/LinkTracker.Scrapper.Storage.Orm/"] COPY ["src/LinkTracker.Scrapper.Storage.Sql/LinkTracker.Scrapper.Storage.Sql.csproj", "src/LinkTracker.Scrapper.Storage.Sql/"] COPY ["src/LinkTracker.EnvReader/LinkTracker.EnvReader.csproj", "src/LinkTracker.EnvReader/"] diff --git a/src/LinkTracker.Scrapper.Api/Program.cs b/src/LinkTracker.Scrapper.Api/Program.cs index 5537635..511fbe2 100644 --- a/src/LinkTracker.Scrapper.Api/Program.cs +++ b/src/LinkTracker.Scrapper.Api/Program.cs @@ -15,7 +15,7 @@ using LinkTracker.Scrapper.Presentation.Grpc; using LinkTracker.Shared.Infrastructure.RateLimiting; using LinkTracker.Shared.Infrastructure.Resilience; -using Prometheus; +using LinkTracker.Shared.Infrastructure.Telemetry; var builder = WebApplication.CreateBuilder(args); @@ -56,5 +56,5 @@ app.MapScrapperEndpoints(); app.MapGrpcService(); -app.MapMetrics(); +app.MapMetricsEndpoint(); await app.RunAsync(); \ No newline at end of file diff --git a/src/LinkTracker.Scrapper.Api/appsettings.Docker.json b/src/LinkTracker.Scrapper.Api/appsettings.Docker.json index 5e5f3fe..f918e8a 100644 --- a/src/LinkTracker.Scrapper.Api/appsettings.Docker.json +++ b/src/LinkTracker.Scrapper.Api/appsettings.Docker.json @@ -57,13 +57,5 @@ "WindowSeconds": 60, "SegmentsPerWindow": 6, "QueueLimit": 0 - }, - "Telemetry": { - "Pushgateway": { - "Enabled": true, - "Endpoint": "http://pushgateway:9091/metrics", - "Job": "scrapper", - "IntervalMilliseconds": 5000 - } } } diff --git a/src/LinkTracker.Scrapper.Api/appsettings.json b/src/LinkTracker.Scrapper.Api/appsettings.json index 86ae557..4c4dec8 100644 --- a/src/LinkTracker.Scrapper.Api/appsettings.json +++ b/src/LinkTracker.Scrapper.Api/appsettings.json @@ -53,7 +53,8 @@ "Enabled": true, "DispatchIntervalSeconds": 10, "BatchSize": 100, - "MaxRetryCount": 3 + "MaxRetryCount": 3, + "LockSeconds": 60 }, "Valkey": { "Enabled": true, @@ -93,13 +94,5 @@ "WindowSeconds": 60, "SegmentsPerWindow": 6, "QueueLimit": 0 - }, - "Telemetry": { - "Pushgateway": { - "Enabled": true, - "Endpoint": "http://localhost:9091/metrics", - "Job": "scrapper", - "IntervalMilliseconds": 5000 - } } } diff --git a/src/LinkTracker.Scrapper.Application/Abstractions/Tracking/ILinkTrackingService.cs b/src/LinkTracker.Scrapper.Application/Abstractions/Tracking/ILinkTrackingService.cs index a215b11..0b1780f 100644 --- a/src/LinkTracker.Scrapper.Application/Abstractions/Tracking/ILinkTrackingService.cs +++ b/src/LinkTracker.Scrapper.Application/Abstractions/Tracking/ILinkTrackingService.cs @@ -14,7 +14,6 @@ Task AddLinkAsync( long chatId, Uri link, IReadOnlyList tags, - IReadOnlyList filters, CancellationToken ct = default); Task RemoveLinkAsync(long chatId, Uri link, CancellationToken ct = default); diff --git a/src/LinkTracker.Scrapper.Application/Services/Tracking/LinkTrackingService.cs b/src/LinkTracker.Scrapper.Application/Services/Tracking/LinkTrackingService.cs index 1871de8..d8028e8 100644 --- a/src/LinkTracker.Scrapper.Application/Services/Tracking/LinkTrackingService.cs +++ b/src/LinkTracker.Scrapper.Application/Services/Tracking/LinkTrackingService.cs @@ -45,7 +45,6 @@ public async Task AddLinkAsync( long chatId, Uri link, IReadOnlyList tags, - IReadOnlyList filters, CancellationToken ct = default) { ValidateChatId(chatId); @@ -57,7 +56,7 @@ public async Task AddLinkAsync( throw ScrapperErrors.ChatNotFound(chatId); } - var record = await store.TryAddAsync(chatId, link, tags, [], ct); + var record = await store.TryAddAsync(chatId, link, tags, ct); return record ?? throw ScrapperErrors.LinkAlreadyExists(link); } diff --git a/src/LinkTracker.Scrapper.Contracts/Requests/AddLinkRequest.cs b/src/LinkTracker.Scrapper.Contracts/Requests/AddLinkRequest.cs index 5fc457f..0ad1866 100644 --- a/src/LinkTracker.Scrapper.Contracts/Requests/AddLinkRequest.cs +++ b/src/LinkTracker.Scrapper.Contracts/Requests/AddLinkRequest.cs @@ -5,6 +5,4 @@ public sealed class AddLinkRequest public Uri? Link { get; init; } public IReadOnlyList Tags { get; init; } = []; - - public IReadOnlyList Filters { get; init; } = []; } \ No newline at end of file diff --git a/src/LinkTracker.Scrapper.Contracts/Responses/LinkResponse.cs b/src/LinkTracker.Scrapper.Contracts/Responses/LinkResponse.cs index dba3f38..62d2691 100644 --- a/src/LinkTracker.Scrapper.Contracts/Responses/LinkResponse.cs +++ b/src/LinkTracker.Scrapper.Contracts/Responses/LinkResponse.cs @@ -7,6 +7,4 @@ public sealed class LinkResponse public Uri Url { get; init; } = default!; public IReadOnlyList Tags { get; init; } = []; - - public IReadOnlyList Filters { get; init; } = []; } \ No newline at end of file diff --git a/src/LinkTracker.Scrapper.Infrastructure/Database/Registration/DatabaseModule.cs b/src/LinkTracker.Scrapper.Infrastructure/Database/Registration/DatabaseModule.cs index ffc77cf..60131f2 100644 --- a/src/LinkTracker.Scrapper.Infrastructure/Database/Registration/DatabaseModule.cs +++ b/src/LinkTracker.Scrapper.Infrastructure/Database/Registration/DatabaseModule.cs @@ -23,7 +23,7 @@ public static IServiceCollection AddDatabase( return new NpgsqlDataSourceBuilder(options.BuildConnectionString()).Build(); }); - services.AddDbContext((sp, options) => + services.AddDbContextFactory((sp, options) => { var databaseOptions = sp.GetRequiredService>().Value; options.UseNpgsql(databaseOptions.BuildConnectionString()); diff --git a/src/LinkTracker.Scrapper.Infrastructure/Outbox/Abstractions/IOutboxStore.cs b/src/LinkTracker.Scrapper.Infrastructure/Outbox/Abstractions/IOutboxStore.cs index d437f74..7b3af62 100644 --- a/src/LinkTracker.Scrapper.Infrastructure/Outbox/Abstractions/IOutboxStore.cs +++ b/src/LinkTracker.Scrapper.Infrastructure/Outbox/Abstractions/IOutboxStore.cs @@ -14,9 +14,10 @@ Task AddRangeAndSetCursorAsync( IReadOnlyCollection updates, CancellationToken ct); - Task> GetUnprocessedBatchAsync( + Task> ClaimUnprocessedBatchAsync( int batchSize, int maxRetryCount, + TimeSpan lockDuration, CancellationToken ct); Task MarkProcessedAsync(long id, CancellationToken ct); diff --git a/src/LinkTracker.Scrapper.Infrastructure/Outbox/Configuration/OutboxOptions.cs b/src/LinkTracker.Scrapper.Infrastructure/Outbox/Configuration/OutboxOptions.cs index a2b6b1b..6ba74de 100644 --- a/src/LinkTracker.Scrapper.Infrastructure/Outbox/Configuration/OutboxOptions.cs +++ b/src/LinkTracker.Scrapper.Infrastructure/Outbox/Configuration/OutboxOptions.cs @@ -9,4 +9,6 @@ public sealed class OutboxOptions public int BatchSize { get; set; } = 100; public int MaxRetryCount { get; set; } = 3; + + public int LockSeconds { get; set; } = 60; } \ No newline at end of file diff --git a/src/LinkTracker.Scrapper.Infrastructure/Outbox/Jobs/OutboxDispatchJob.cs b/src/LinkTracker.Scrapper.Infrastructure/Outbox/Jobs/OutboxDispatchJob.cs index 46bc10d..654ea4a 100644 --- a/src/LinkTracker.Scrapper.Infrastructure/Outbox/Jobs/OutboxDispatchJob.cs +++ b/src/LinkTracker.Scrapper.Infrastructure/Outbox/Jobs/OutboxDispatchJob.cs @@ -19,9 +19,10 @@ public async Task Execute(IJobExecutionContext context) var ct = context.CancellationToken; var options = outboxOptions.Value; - var messages = await outboxStore.GetUnprocessedBatchAsync( + var messages = await outboxStore.ClaimUnprocessedBatchAsync( options.BatchSize, options.MaxRetryCount, + TimeSpan.FromSeconds(options.LockSeconds), ct); if (messages.Count == 0) diff --git a/src/LinkTracker.Scrapper.Infrastructure/Outbox/PostgresOutboxStore.cs b/src/LinkTracker.Scrapper.Infrastructure/Outbox/PostgresOutboxStore.cs index 4032498..8e8d736 100644 --- a/src/LinkTracker.Scrapper.Infrastructure/Outbox/PostgresOutboxStore.cs +++ b/src/LinkTracker.Scrapper.Infrastructure/Outbox/PostgresOutboxStore.cs @@ -77,16 +77,18 @@ public Task AddRangeAndSetCursorAsync( }); } - public Task> GetUnprocessedBatchAsync( + public Task> ClaimUnprocessedBatchAsync( int batchSize, int maxRetryCount, + TimeSpan lockDuration, CancellationToken ct) { - return MeasureAsync("get_unprocessed_batch", async () => + return MeasureAsync("claim_unprocessed_batch", async () => { - await using var command = dataSource.CreateCommand(OutboxMessageCommands.GetUnprocessedBatch); + await using var command = dataSource.CreateCommand(OutboxMessageCommands.ClaimUnprocessedBatch); command.Parameters.AddWithValue("batchSize", batchSize); command.Parameters.AddWithValue("maxRetryCount", maxRetryCount); + command.Parameters.AddWithValue("lockSeconds", lockDuration.TotalSeconds); await using var reader = await command.ExecuteReaderAsync(ct); @@ -116,7 +118,10 @@ public Task> GetUnprocessedBatchAsync( }); } - return (IReadOnlyList)result; + return (IReadOnlyList)result + .OrderBy(x => x.CreatedAt) + .ThenBy(x => x.Id) + .ToArray(); }); } diff --git a/src/LinkTracker.Scrapper.Infrastructure/Outbox/Registration/OutboxModule.cs b/src/LinkTracker.Scrapper.Infrastructure/Outbox/Registration/OutboxModule.cs index 2977766..82aa4e7 100644 --- a/src/LinkTracker.Scrapper.Infrastructure/Outbox/Registration/OutboxModule.cs +++ b/src/LinkTracker.Scrapper.Infrastructure/Outbox/Registration/OutboxModule.cs @@ -24,6 +24,10 @@ public static IServiceCollection AddOutbox( .Validate(o => o.DispatchIntervalSeconds > 0, "Outbox:DispatchIntervalSeconds must be greater than zero.") .Validate(o => o.BatchSize > 0, "Outbox:BatchSize must be greater than zero.") .Validate(o => o.MaxRetryCount > 0, "Outbox:MaxRetryCount must be greater than zero.") + .Validate(o => o.LockSeconds > 0, "Outbox:LockSeconds must be greater than zero.") + .Validate( + o => o.LockSeconds > o.DispatchIntervalSeconds, + "Outbox:LockSeconds must be greater than Outbox:DispatchIntervalSeconds, otherwise a batch can be picked up twice.") .ValidateOnStart(); services.AddSingleton(); diff --git a/src/LinkTracker.Scrapper.Infrastructure/Outbox/Sql/OutboxMessageCommands.cs b/src/LinkTracker.Scrapper.Infrastructure/Outbox/Sql/OutboxMessageCommands.cs index 6a7da5b..04017e5 100644 --- a/src/LinkTracker.Scrapper.Infrastructure/Outbox/Sql/OutboxMessageCommands.cs +++ b/src/LinkTracker.Scrapper.Infrastructure/Outbox/Sql/OutboxMessageCommands.cs @@ -17,27 +17,35 @@ UPDATE links WHERE id = @linkId; """; - public const string GetUnprocessedBatch = + public const string ClaimUnprocessedBatch = """ - SELECT id, - payload::text, - created_at, - processed_at, - error, - retry_count - FROM outbox_messages - WHERE processed_at IS NULL - AND retry_count < @maxRetryCount - ORDER BY created_at, id - LIMIT @batchSize - FOR UPDATE SKIP LOCKED; + UPDATE outbox_messages AS o + SET locked_until = now() + make_interval(secs => @lockSeconds) + FROM ( + SELECT id + FROM outbox_messages + WHERE processed_at IS NULL + AND retry_count < @maxRetryCount + AND (locked_until IS NULL OR locked_until <= now()) + ORDER BY created_at, id + LIMIT @batchSize + FOR UPDATE SKIP LOCKED + ) AS claimed + WHERE o.id = claimed.id + RETURNING o.id, + o.payload::text, + o.created_at, + o.processed_at, + o.error, + o.retry_count; """; public const string MarkProcessed = """ UPDATE outbox_messages SET processed_at = now(), - error = NULL + error = NULL, + locked_until = NULL WHERE id = @id; """; @@ -45,7 +53,8 @@ UPDATE outbox_messages """ UPDATE outbox_messages SET error = @error, - retry_count = retry_count + 1 + retry_count = retry_count + 1, + locked_until = NULL WHERE id = @id; """; } \ No newline at end of file diff --git a/src/LinkTracker.Scrapper.Infrastructure/Storage/Registration/StorageModule.cs b/src/LinkTracker.Scrapper.Infrastructure/Storage/Registration/StorageModule.cs index 698e92e..11e24be 100644 --- a/src/LinkTracker.Scrapper.Infrastructure/Storage/Registration/StorageModule.cs +++ b/src/LinkTracker.Scrapper.Infrastructure/Storage/Registration/StorageModule.cs @@ -31,7 +31,7 @@ public static IServiceCollection AddStorage( services.AddSingleton(); break; case DatabaseAccessType.Orm: - services.AddScoped(); + services.AddSingleton(); break; default: throw new InvalidOperationException( diff --git a/src/LinkTracker.Scrapper.Infrastructure/Telemetry/Registration/TelemetryModule.cs b/src/LinkTracker.Scrapper.Infrastructure/Telemetry/Registration/TelemetryModule.cs index 1573cd8..0d822ab 100644 --- a/src/LinkTracker.Scrapper.Infrastructure/Telemetry/Registration/TelemetryModule.cs +++ b/src/LinkTracker.Scrapper.Infrastructure/Telemetry/Registration/TelemetryModule.cs @@ -12,9 +12,7 @@ public static IServiceCollection AddTelemetry( { services.AddSingleton(); - services.AddOpenTelemetryMetricsWithPushgateway( - configuration, - "scrapper", + services.AddOpenTelemetryMetrics( "scrapper", ScrapperMetrics.MeterName, "Npgsql"); diff --git a/src/LinkTracker.Scrapper.Presentation/Endpoints/ScrapperEndpoints.cs b/src/LinkTracker.Scrapper.Presentation/Endpoints/ScrapperEndpoints.cs index 1d3c8f3..45b8fc0 100644 --- a/src/LinkTracker.Scrapper.Presentation/Endpoints/ScrapperEndpoints.cs +++ b/src/LinkTracker.Scrapper.Presentation/Endpoints/ScrapperEndpoints.cs @@ -118,7 +118,6 @@ private static async Task AddLink( resolvedChatId, data.Link, data.Tags, - data.Filters, ct); await cache.InvalidateAsync(resolvedChatId, ct); diff --git a/src/LinkTracker.Scrapper.Presentation/Grpc/ScrapperGrpcService.cs b/src/LinkTracker.Scrapper.Presentation/Grpc/ScrapperGrpcService.cs index dfe511c..84f4a22 100644 --- a/src/LinkTracker.Scrapper.Presentation/Grpc/ScrapperGrpcService.cs +++ b/src/LinkTracker.Scrapper.Presentation/Grpc/ScrapperGrpcService.cs @@ -49,7 +49,6 @@ public override async Task AddLink(AddLinkGrpcRequest request, request.ChatId, new Uri(request.Link), request.Tags.ToArray(), - request.Filters.ToArray(), context.CancellationToken); await cache.InvalidateAsync(request.ChatId, context.CancellationToken); @@ -83,7 +82,6 @@ private static LinkGrpcResponse ToGrpc(LinkResponse response) var grpcResponse = new LinkGrpcResponse { Id = response.Id, Url = response.Url.ToString() }; grpcResponse.Tags.AddRange(response.Tags); - grpcResponse.Filters.AddRange(response.Filters); return grpcResponse; } diff --git a/src/LinkTracker.Scrapper.Presentation/Mappers/LinkRequestMappings.cs b/src/LinkTracker.Scrapper.Presentation/Mappers/LinkRequestMappings.cs index ea3d983..0ddbf90 100644 --- a/src/LinkTracker.Scrapper.Presentation/Mappers/LinkRequestMappings.cs +++ b/src/LinkTracker.Scrapper.Presentation/Mappers/LinkRequestMappings.cs @@ -5,10 +5,10 @@ namespace LinkTracker.Scrapper.Presentation.Mappers; public static class LinkRequestMappings { - public static (Uri Link, IReadOnlyList Tags, IReadOnlyList Filters) ToAddLinkData( + public static (Uri Link, IReadOnlyList Tags) ToAddLinkData( this AddLinkRequest? request) { - return request?.Link is null ? throw ScrapperErrors.RequestLinkIsRequired() : (request.Link, request.Tags, request.Filters); + return request?.Link is null ? throw ScrapperErrors.RequestLinkIsRequired() : (request.Link, request.Tags); } public static Uri ToRemoveLinkData(this RemoveLinkRequest? request) diff --git a/src/LinkTracker.Scrapper.Presentation/Mappers/TrackedLinkMappings.cs b/src/LinkTracker.Scrapper.Presentation/Mappers/TrackedLinkMappings.cs index 496b8f7..935ceb6 100644 --- a/src/LinkTracker.Scrapper.Presentation/Mappers/TrackedLinkMappings.cs +++ b/src/LinkTracker.Scrapper.Presentation/Mappers/TrackedLinkMappings.cs @@ -7,6 +7,6 @@ public static class TrackedLinkMappings { public static LinkResponse ToResponse(this TrackedLinkRecord record) { - return new LinkResponse { Id = record.Id, Url = record.Url, Tags = record.Tags, Filters = [] }; + return new LinkResponse { Id = record.Id, Url = record.Url, Tags = record.Tags }; } } \ No newline at end of file diff --git a/src/LinkTracker.Scrapper.Storage.Abstractions/Models/ILinkTrackingStore.cs b/src/LinkTracker.Scrapper.Storage.Abstractions/Models/ILinkTrackingStore.cs index 926fec9..ed9a28b 100644 --- a/src/LinkTracker.Scrapper.Storage.Abstractions/Models/ILinkTrackingStore.cs +++ b/src/LinkTracker.Scrapper.Storage.Abstractions/Models/ILinkTrackingStore.cs @@ -14,7 +14,6 @@ public interface ILinkTrackingStore long chatId, Uri link, IReadOnlyList tags, - IReadOnlyList filters, CancellationToken ct = default); Task TryRemoveAsync(long chatId, Uri link, CancellationToken ct = default); diff --git a/src/LinkTracker.Scrapper.Storage.Abstractions/Models/TrackedLinkRecord.cs b/src/LinkTracker.Scrapper.Storage.Abstractions/Models/TrackedLinkRecord.cs index 416e787..c476fda 100644 --- a/src/LinkTracker.Scrapper.Storage.Abstractions/Models/TrackedLinkRecord.cs +++ b/src/LinkTracker.Scrapper.Storage.Abstractions/Models/TrackedLinkRecord.cs @@ -8,8 +8,7 @@ public sealed class TrackedLinkRecord public IReadOnlyList Tags { get; init; } = []; - public IReadOnlyList Filters { get; init; } = []; - public DateTimeOffset? LastUpdatedAt { get; set; } + public string? LastEventKey { get; set; } } \ No newline at end of file diff --git a/src/LinkTracker.Scrapper.Storage.Orm/AppDbContext.cs b/src/LinkTracker.Scrapper.Storage.Orm/AppDbContext.cs index 05c332c..1b48927 100644 --- a/src/LinkTracker.Scrapper.Storage.Orm/AppDbContext.cs +++ b/src/LinkTracker.Scrapper.Storage.Orm/AppDbContext.cs @@ -9,9 +9,7 @@ public sealed class AppDbContext(DbContextOptions options) : DbCon public DbSet Links => Set(); public DbSet Subscriptions => Set(); public DbSet Tags => Set(); - public DbSet Filters => Set(); public DbSet SubscriptionTags => Set(); - public DbSet SubscriptionFilters => Set(); protected override void OnModelCreating(ModelBuilder modelBuilder) { diff --git a/src/LinkTracker.Scrapper.Storage.Orm/Configurations/FilterEntityConfiguration.cs b/src/LinkTracker.Scrapper.Storage.Orm/Configurations/FilterEntityConfiguration.cs deleted file mode 100644 index 9b8b362..0000000 --- a/src/LinkTracker.Scrapper.Storage.Orm/Configurations/FilterEntityConfiguration.cs +++ /dev/null @@ -1,25 +0,0 @@ -using LinkTracker.Scrapper.Storage.Orm.Entities; -using Microsoft.EntityFrameworkCore; -using Microsoft.EntityFrameworkCore.Metadata.Builders; - -namespace LinkTracker.Scrapper.Storage.Orm.Configurations; - -public sealed class FilterEntityConfiguration : IEntityTypeConfiguration -{ - public void Configure(EntityTypeBuilder builder) - { - builder.ToTable("filters"); - - builder.HasKey(x => x.Id); - - builder.Property(x => x.Id) - .HasColumnName("id"); - - builder.Property(x => x.Value) - .HasColumnName("value") - .IsRequired(); - - builder.HasIndex(x => x.Value) - .IsUnique(); - } -} \ No newline at end of file diff --git a/src/LinkTracker.Scrapper.Storage.Orm/Configurations/SubscriptionFilterEntityConfiguration.cs b/src/LinkTracker.Scrapper.Storage.Orm/Configurations/SubscriptionFilterEntityConfiguration.cs deleted file mode 100644 index eb425a0..0000000 --- a/src/LinkTracker.Scrapper.Storage.Orm/Configurations/SubscriptionFilterEntityConfiguration.cs +++ /dev/null @@ -1,31 +0,0 @@ -using LinkTracker.Scrapper.Storage.Orm.Entities; -using Microsoft.EntityFrameworkCore; -using Microsoft.EntityFrameworkCore.Metadata.Builders; - -namespace LinkTracker.Scrapper.Storage.Orm.Configurations; - -public sealed class SubscriptionFilterEntityConfiguration : IEntityTypeConfiguration -{ - public void Configure(EntityTypeBuilder builder) - { - builder.ToTable("subscription_filters"); - - builder.HasKey(x => new { x.SubscriptionId, x.FilterId }); - - builder.Property(x => x.SubscriptionId) - .HasColumnName("subscription_id"); - - builder.Property(x => x.FilterId) - .HasColumnName("filter_id"); - - builder.HasOne(x => x.Subscription) - .WithMany(x => x.SubscriptionFilters) - .HasForeignKey(x => x.SubscriptionId) - .OnDelete(DeleteBehavior.Cascade); - - builder.HasOne(x => x.Filter) - .WithMany(x => x.SubscriptionFilters) - .HasForeignKey(x => x.FilterId) - .OnDelete(DeleteBehavior.Cascade); - } -} \ No newline at end of file diff --git a/src/LinkTracker.Scrapper.Storage.Orm/Entities/FilterEntity.cs b/src/LinkTracker.Scrapper.Storage.Orm/Entities/FilterEntity.cs deleted file mode 100644 index a5e4553..0000000 --- a/src/LinkTracker.Scrapper.Storage.Orm/Entities/FilterEntity.cs +++ /dev/null @@ -1,10 +0,0 @@ -namespace LinkTracker.Scrapper.Storage.Orm.Entities; - -public sealed class FilterEntity -{ - public long Id { get; set; } - - public string Value { get; set; } = string.Empty; - - public ICollection SubscriptionFilters { get; set; } = []; -} \ No newline at end of file diff --git a/src/LinkTracker.Scrapper.Storage.Orm/Entities/SubscriptionEntity.cs b/src/LinkTracker.Scrapper.Storage.Orm/Entities/SubscriptionEntity.cs index ba88b3e..f58b8b2 100644 --- a/src/LinkTracker.Scrapper.Storage.Orm/Entities/SubscriptionEntity.cs +++ b/src/LinkTracker.Scrapper.Storage.Orm/Entities/SubscriptionEntity.cs @@ -14,5 +14,4 @@ public sealed class SubscriptionEntity public ICollection SubscriptionTags { get; set; } = []; - public ICollection SubscriptionFilters { get; set; } = []; } \ No newline at end of file diff --git a/src/LinkTracker.Scrapper.Storage.Orm/Entities/SubscriptionFilterEntity.cs b/src/LinkTracker.Scrapper.Storage.Orm/Entities/SubscriptionFilterEntity.cs deleted file mode 100644 index 67029c1..0000000 --- a/src/LinkTracker.Scrapper.Storage.Orm/Entities/SubscriptionFilterEntity.cs +++ /dev/null @@ -1,12 +0,0 @@ -namespace LinkTracker.Scrapper.Storage.Orm.Entities; - -public sealed class SubscriptionFilterEntity -{ - public long SubscriptionId { get; set; } - - public SubscriptionEntity Subscription { get; set; } = default!; - - public long FilterId { get; set; } - - public FilterEntity Filter { get; set; } = default!; -} \ No newline at end of file diff --git a/src/LinkTracker.Scrapper.Storage.Orm/OrmLinkTrackingStore.cs b/src/LinkTracker.Scrapper.Storage.Orm/OrmLinkTrackingStore.cs index 80f3a3e..da39241 100644 --- a/src/LinkTracker.Scrapper.Storage.Orm/OrmLinkTrackingStore.cs +++ b/src/LinkTracker.Scrapper.Storage.Orm/OrmLinkTrackingStore.cs @@ -6,10 +6,12 @@ namespace LinkTracker.Scrapper.Storage.Orm; -public sealed class OrmLinkTrackingStore(AppDbContext dbContext) : ILinkTrackingStore +public sealed class OrmLinkTrackingStore(IDbContextFactory dbContextFactory) : ILinkTrackingStore { public async Task TryRegisterChatAsync(long chatId, CancellationToken ct = default) { + await using var dbContext = await dbContextFactory.CreateDbContextAsync(ct); + var exists = await dbContext.Chats .AnyAsync(x => x.Id == chatId, ct); @@ -26,6 +28,8 @@ public async Task TryRegisterChatAsync(long chatId, CancellationToken ct = public async Task TryDeleteChatAsync(long chatId, CancellationToken ct = default) { + await using var dbContext = await dbContextFactory.CreateDbContextAsync(ct); + var chat = await dbContext.Chats .FirstOrDefaultAsync(x => x.Id == chatId, ct); @@ -40,13 +44,17 @@ public async Task TryDeleteChatAsync(long chatId, CancellationToken ct = d return true; } - public Task ChatExistsAsync(long chatId, CancellationToken ct = default) + public async Task ChatExistsAsync(long chatId, CancellationToken ct = default) { - return dbContext.Chats.AnyAsync(x => x.Id == chatId, ct); + await using var dbContext = await dbContextFactory.CreateDbContextAsync(ct); + + return await dbContext.Chats.AnyAsync(x => x.Id == chatId, ct); } public async Task> GetAllTrackedLinkRecordsAsync(long chatId, CancellationToken ct = default) { + await using var dbContext = await dbContextFactory.CreateDbContextAsync(ct); + return await dbContext.Subscriptions .AsNoTracking() .Where(x => x.ChatId == chatId) @@ -61,7 +69,6 @@ public async Task> GetAllTrackedLinkRecordsAsyn .Select(t => t.Tag.Name) .Distinct() .ToArray(), - Filters = Array.Empty() }) .ToArrayAsync(ct); } @@ -70,7 +77,6 @@ public async Task> GetAllTrackedLinkRecordsAsyn long chatId, Uri link, IReadOnlyList tags, - IReadOnlyList filters, CancellationToken ct = default) { var normalizedUrl = TrackedLinkUrl.Normalize(link); @@ -80,6 +86,7 @@ public async Task> GetAllTrackedLinkRecordsAsyn .Distinct(StringComparer.Ordinal) .ToArray(); + await using var dbContext = await dbContextFactory.CreateDbContextAsync(ct); await using var transaction = await dbContext.Database.BeginTransactionAsync(ct); try @@ -127,7 +134,6 @@ public async Task> GetAllTrackedLinkRecordsAsyn Id = linkEntity.Id, Url = link, Tags = cleanedTags, - Filters = [], LastUpdatedAt = linkEntity.LastUpdatedAt, LastEventKey = linkEntity.LastEventKey }; @@ -143,6 +149,7 @@ public async Task> GetAllTrackedLinkRecordsAsyn { var normalizedUrl = TrackedLinkUrl.Normalize(link); + await using var dbContext = await dbContextFactory.CreateDbContextAsync(ct); await using var transaction = await dbContext.Database.BeginTransactionAsync(ct); try @@ -169,7 +176,6 @@ public async Task> GetAllTrackedLinkRecordsAsyn .Select(x => x.Tag.Name) .Distinct() .ToArray(), - Filters = [] }; var linkId = subscription.LinkId; @@ -218,6 +224,8 @@ public async Task TryCreateTagAsync(string tag, CancellationToken ct = def { var normalizedTag = tag.Trim(); + await using var dbContext = await dbContextFactory.CreateDbContextAsync(ct); + var exists = await dbContext.Tags .AnyAsync(x => x.Name == normalizedTag, ct); @@ -234,6 +242,8 @@ public async Task TryCreateTagAsync(string tag, CancellationToken ct = def public async Task> GetTagsAsync(long chatId, CancellationToken ct = default) { + await using var dbContext = await dbContextFactory.CreateDbContextAsync(ct); + return await dbContext.Subscriptions .AsNoTracking() .Where(x => x.ChatId == chatId) @@ -248,6 +258,7 @@ public async Task> GetTagsAsync(long chatId, CancellationT var normalizedUrl = TrackedLinkUrl.Normalize(link); var cleanedTag = tag.Trim(); + await using var dbContext = await dbContextFactory.CreateDbContextAsync(ct); await using var transaction = await dbContext.Database.BeginTransactionAsync(ct); try @@ -300,6 +311,7 @@ public async Task TryRenameTagAsync(long chatId, string tag, string newTag var cleanedOldTag = tag.Trim(); var cleanedNewTag = newTag.Trim(); + await using var dbContext = await dbContextFactory.CreateDbContextAsync(ct); await using var transaction = await dbContext.Database.BeginTransactionAsync(ct); try @@ -354,6 +366,7 @@ public async Task TryDeleteTagAsync(long chatId, string tag, CancellationT { var cleanedTag = tag.Trim(); + await using var dbContext = await dbContextFactory.CreateDbContextAsync(ct); await using var transaction = await dbContext.Database.BeginTransactionAsync(ct); try @@ -387,6 +400,8 @@ public async Task TryDeleteTagAsync(long chatId, string tag, CancellationT public async Task> GetAllSubscriptionsAsync(CancellationToken ct = default) { + await using var dbContext = await dbContextFactory.CreateDbContextAsync(ct); + return await dbContext.Subscriptions .AsNoTracking() .GroupBy(x => new { x.LinkId, x.Link.Url, x.Link.LastUpdatedAt }) @@ -406,6 +421,8 @@ public async Task> GetSubscriptionsBatchA int batchSize, CancellationToken ct = default) { + await using var dbContext = await dbContextFactory.CreateDbContextAsync(ct); + return await dbContext.Links .AsNoTracking() .Where(x => x.Subscriptions.Any()) @@ -428,6 +445,8 @@ public async Task> GetSubscriptionsBatchA public async Task SetCursorAsync(long linkId, DateTimeOffset lastUpdatedAt, string? lastEventKey, CancellationToken ct = default) { + await using var dbContext = await dbContextFactory.CreateDbContextAsync(ct); + var entity = await dbContext.Links .FirstOrDefaultAsync(x => x.Id == linkId, ct); @@ -454,7 +473,6 @@ private static TrackedLinkRecord ToRecord(SubscriptionEntity subscription) .Distinct() .OrderBy(x => x) .ToArray(), - Filters = [] }; } } \ No newline at end of file diff --git a/src/LinkTracker.Scrapper.Storage.Sql/SqlLinkTrackingStore.cs b/src/LinkTracker.Scrapper.Storage.Sql/SqlLinkTrackingStore.cs index e42afe0..9aa8caf 100644 --- a/src/LinkTracker.Scrapper.Storage.Sql/SqlLinkTrackingStore.cs +++ b/src/LinkTracker.Scrapper.Storage.Sql/SqlLinkTrackingStore.cs @@ -67,7 +67,6 @@ public async Task> GetAllTrackedLinkRecordsAsyn long chatId, Uri link, IReadOnlyList tags, - IReadOnlyList filters, CancellationToken ct = default) { var normalizedUrl = TrackedLinkUrl.Normalize(link); @@ -126,7 +125,6 @@ await connection.ExecuteAsync( Id = linkId, Url = link, Tags = cleanedTags, - Filters = [], LastUpdatedAt = null, LastEventKey = null }; @@ -198,7 +196,6 @@ await connection.ExecuteAsync( LastUpdatedAt = row.LastUpdatedAt, LastEventKey = row.LastEventKey, Tags = row.Tags, - Filters = [] }; } catch @@ -449,7 +446,6 @@ await connection.ExecuteAsync( LastUpdatedAt = row.LastUpdatedAt, LastEventKey = row.LastEventKey, Tags = row.Tags, - Filters = [] }; } @@ -462,7 +458,6 @@ private static TrackedLinkRecord MapTrackedLinkRecord(TrackedLinkRow row) LastUpdatedAt = row.LastUpdatedAt, LastEventKey = row.LastEventKey, Tags = row.Tags, - Filters = [] }; } diff --git a/src/LinkTracker.Shared/Contracts/Bot/LinkUpdate.cs b/src/LinkTracker.Shared/Contracts/Bot/LinkUpdate.cs index 374d156..e91091c 100644 --- a/src/LinkTracker.Shared/Contracts/Bot/LinkUpdate.cs +++ b/src/LinkTracker.Shared/Contracts/Bot/LinkUpdate.cs @@ -1,3 +1,5 @@ +using LinkTracker.Shared.Contracts.AiAgent; + namespace LinkTracker.Shared.Contracts.Bot; public sealed class LinkUpdate @@ -11,4 +13,10 @@ public sealed class LinkUpdate public string Author { get; init; } = string.Empty; public IReadOnlyList TgChatIds { get; init; } = []; + + /// + /// Проставляется AI-агентом. На сыром пути Scrapper -> Bot остаётся Medium: + /// приоритет там просто ещё не вычислен. + /// + public LinkUpdatePriority Priority { get; init; } = LinkUpdatePriority.Medium; } \ No newline at end of file diff --git a/src/LinkTracker.Shared/Infrastructure/Telemetry/OpenTelemetryRegistration.cs b/src/LinkTracker.Shared/Infrastructure/Telemetry/OpenTelemetryRegistration.cs index 2ec07c5..77a936a 100644 --- a/src/LinkTracker.Shared/Infrastructure/Telemetry/OpenTelemetryRegistration.cs +++ b/src/LinkTracker.Shared/Infrastructure/Telemetry/OpenTelemetryRegistration.cs @@ -1,4 +1,6 @@ -using Microsoft.Extensions.Configuration; +using Microsoft.AspNetCore.Builder; +using Microsoft.AspNetCore.RateLimiting; +using Microsoft.AspNetCore.Routing; using Microsoft.Extensions.DependencyInjection; using OpenTelemetry.Metrics; using OpenTelemetry.Resources; @@ -7,24 +9,11 @@ namespace LinkTracker.Shared.Infrastructure.Telemetry; public static class OpenTelemetryRegistration { - public static IServiceCollection AddOpenTelemetryMetricsWithPushgateway( + public static IServiceCollection AddOpenTelemetryMetrics( this IServiceCollection services, - IConfiguration configuration, string serviceName, - string job, params string[] meterNames) { - services - .AddOptions() - .Bind(configuration.GetSection(PushgatewayOptions.SectionName)) - .PostConfigure(opt => - { - if (string.IsNullOrWhiteSpace(opt.Job) || opt.Job == "app") - { - opt.Job = job; - } - }); - services .AddOpenTelemetry() .ConfigureResource(resource => resource.AddService(serviceName)) @@ -39,10 +28,17 @@ public static IServiceCollection AddOpenTelemetryMetricsWithPushgateway( { metrics.AddMeter(meterName); } - }); - services.AddHostedService(); + metrics.AddPrometheusExporter(); + }); return services; } -} \ No newline at end of file + + public static IEndpointConventionBuilder MapMetricsEndpoint(this IEndpointRouteBuilder endpoints) + { + return endpoints + .MapPrometheusScrapingEndpoint() + .DisableRateLimiting(); + } +} diff --git a/src/LinkTracker.Shared/Infrastructure/Telemetry/PushgatewayMetricPusherHostedService.cs b/src/LinkTracker.Shared/Infrastructure/Telemetry/PushgatewayMetricPusherHostedService.cs deleted file mode 100644 index aad3d36..0000000 --- a/src/LinkTracker.Shared/Infrastructure/Telemetry/PushgatewayMetricPusherHostedService.cs +++ /dev/null @@ -1,58 +0,0 @@ -using System.Net; -using Microsoft.Extensions.Hosting; -using Microsoft.Extensions.Logging; -using Microsoft.Extensions.Options; -using Prometheus; - -namespace LinkTracker.Shared.Infrastructure.Telemetry; - -public sealed class PushgatewayMetricPusherHostedService( - IOptions options, - ILogger logger) : IHostedService -{ - private MetricPusher? _pusher; - - public Task StartAsync(CancellationToken cancellationToken) - { - var value = options.Value; - - if (!value.Enabled) - { - logger.LogInformation("Pushgateway push отключён (Telemetry:Pushgateway:Enabled = false)."); - return Task.CompletedTask; - } - - var instance = Environment.GetEnvironmentVariable("HOSTNAME") - ?? Dns.GetHostName(); - - _pusher = new MetricPusher(new MetricPusherOptions - { - Endpoint = value.Endpoint, - Job = value.Job, - Instance = instance, - IntervalMilliseconds = value.IntervalMilliseconds, - OnError = ex => logger.LogWarning( - ex, - "Не удалось отправить метрики в Pushgateway. Endpoint={Endpoint}, Job={Job}", - value.Endpoint, - value.Job) - }); - - _pusher.Start(); - - logger.LogInformation( - "Запущен push метрик в Pushgateway. Endpoint={Endpoint}, Job={Job}, Instance={Instance}, IntervalMs={Interval}", - value.Endpoint, - value.Job, - instance, - value.IntervalMilliseconds); - - return Task.CompletedTask; - } - - public Task StopAsync(CancellationToken cancellationToken) - { - _pusher?.Stop(); - return Task.CompletedTask; - } -} \ No newline at end of file diff --git a/src/LinkTracker.Shared/Infrastructure/Telemetry/PushgatewayOptions.cs b/src/LinkTracker.Shared/Infrastructure/Telemetry/PushgatewayOptions.cs deleted file mode 100644 index 67e74ae..0000000 --- a/src/LinkTracker.Shared/Infrastructure/Telemetry/PushgatewayOptions.cs +++ /dev/null @@ -1,14 +0,0 @@ -namespace LinkTracker.Shared.Infrastructure.Telemetry; - -public sealed class PushgatewayOptions -{ - public const string SectionName = "Telemetry:Pushgateway"; - - public bool Enabled { get; set; } = true; - - public string Endpoint { get; set; } = "http://pushgateway:9091/metrics"; - - public string Job { get; set; } = "app"; - - public int IntervalMilliseconds { get; set; } = 5000; -} \ No newline at end of file diff --git a/src/LinkTracker.Shared/LinkTracker.Shared.csproj b/src/LinkTracker.Shared/LinkTracker.Shared.csproj index db916b4..8cb9e5e 100644 --- a/src/LinkTracker.Shared/LinkTracker.Shared.csproj +++ b/src/LinkTracker.Shared/LinkTracker.Shared.csproj @@ -21,7 +21,7 @@ - + @@ -29,10 +29,4 @@ - - - ..\..\..\..\..\..\.nuget\packages\prometheus-net\8.2.1\lib\net7.0\Prometheus.NetStandard.dll - - - \ No newline at end of file diff --git a/src/LinkTracker.Shared/Protos/scrapper.proto b/src/LinkTracker.Shared/Protos/scrapper.proto index 5e2a485..4b803fe 100644 --- a/src/LinkTracker.Shared/Protos/scrapper.proto +++ b/src/LinkTracker.Shared/Protos/scrapper.proto @@ -18,7 +18,7 @@ message AddLinkGrpcRequest { int64 chat_id = 1; string link = 2; repeated string tags = 3; - repeated string filters = 4; + reserved 4; } message RemoveLinkGrpcRequest { @@ -30,7 +30,7 @@ message LinkGrpcResponse { int64 id = 1; string url = 2; repeated string tags = 3; - repeated string filters = 4; + reserved 4; } message ListLinksGrpcResponse { diff --git a/src/LinkTracker.Tests/Bot/Unit/Application/Dialogs/Implementations/Track/Nodes/TrackConfirmNodeTest.cs b/src/LinkTracker.Tests/Bot/Unit/Application/Dialogs/Implementations/Track/Nodes/TrackConfirmNodeTest.cs index f3dad59..4b84edb 100644 --- a/src/LinkTracker.Tests/Bot/Unit/Application/Dialogs/Implementations/Track/Nodes/TrackConfirmNodeTest.cs +++ b/src/LinkTracker.Tests/Bot/Unit/Application/Dialogs/Implementations/Track/Nodes/TrackConfirmNodeTest.cs @@ -23,7 +23,6 @@ public async Task Handle_WhenUserAnswersYesAndLinkAdded_ReturnsSuccessAndEndsDia 123L, uri, Arg.Is>(x => x.SequenceEqual(new[] { "dotnet", "runtime" })), - Arg.Is>(x => x.Count == 0), Arg.Any()) .Returns(new LinkResponse { Id = 1, Url = uri, Tags = ["dotnet", "runtime"] }); @@ -42,7 +41,6 @@ await scrapperClient.Received(1).AddLinkAsync( 123L, uri, Arg.Is>(x => x.SequenceEqual(new[] { "dotnet", "runtime" })), - Arg.Is>(x => x.Count == 0), Arg.Any()); } @@ -68,7 +66,6 @@ await scrapperClient.DidNotReceive() Arg.Any(), Arg.Any(), Arg.Any>(), - Arg.Any>(), Arg.Any()); } @@ -82,7 +79,6 @@ public async Task Handle_WhenUnsupportedLink_ReturnsRetryMessage() 123L, Arg.Any(), Arg.Any>(), - Arg.Any>(), Arg.Any()) .Returns(Task.FromException( CreateException( @@ -112,7 +108,6 @@ public async Task Handle_WhenLinkAlreadyExists_ReturnsFriendlyMessageAndEndsDial 123L, Arg.Any(), Arg.Any>(), - Arg.Any>(), Arg.Any()) .Returns(Task.FromException( CreateException( @@ -140,7 +135,6 @@ public async Task Handle_WhenUnknownScrapperError_ReturnsFallbackMessage() 123L, Arg.Any(), Arg.Any>(), - Arg.Any>(), Arg.Any()) .Returns(Task.FromException( new ScrapperClientException( diff --git a/src/LinkTracker.Tests/Bot/Unit/Presentation/Telegram/Notifications/LinkUpdateNotifierTests.cs b/src/LinkTracker.Tests/Bot/Unit/Presentation/Telegram/Notifications/LinkUpdateNotifierTests.cs new file mode 100644 index 0000000..c94eb44 --- /dev/null +++ b/src/LinkTracker.Tests/Bot/Unit/Presentation/Telegram/Notifications/LinkUpdateNotifierTests.cs @@ -0,0 +1,91 @@ +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; +using Telegram.Bot; +using Telegram.Bot.Requests; +using Telegram.Bot.Types; + +namespace LinkTracker.Tests.Bot.Unit.Presentation.Telegram.Notifications; + +[Trait("Module", "Bot")] +[Trait("Category", "Unit")] +public sealed class LinkUpdateNotifierTests +{ + private static readonly Uri Url = new("https://github.com/dotnet/runtime"); + + [Theory] + [InlineData(LinkUpdatePriority.High, "‼️ Важное обновление по ссылке:")] + [InlineData(LinkUpdatePriority.Medium, "Обновление по ссылке:")] + [InlineData(LinkUpdatePriority.Low, "Незначительное обновление по ссылке:")] + public async Task NotifyAsync_RendersHeaderMatchingPriority( + LinkUpdatePriority priority, + string expectedHeader) + { + var (botClient, sentTexts) = CreateBotClient(); + var sut = new LinkUpdateNotifier(botClient, Substitute.For()); + + await sut.NotifyAsync( + new LinkUpdate { Id = 1, Url = Url, Description = "Новый issue", TgChatIds = [42], Priority = priority }, + CancellationToken.None); + + var text = Assert.Single(sentTexts); + + Assert.StartsWith(expectedHeader, text, StringComparison.Ordinal); + Assert.Contains("Новый issue", text, StringComparison.Ordinal); + } + + [Fact] + public async Task NotifyAsync_WhenPriorityIsNotSet_UsesNeutralHeader() + { + var (botClient, sentTexts) = CreateBotClient(); + var sut = new LinkUpdateNotifier(botClient, Substitute.For()); + + // Сырой путь Scrapper -> Bot приоритет не проставляет. + await sut.NotifyAsync( + new LinkUpdate { Id = 1, Url = Url, Description = "Новый issue", TgChatIds = [42] }, + CancellationToken.None); + + Assert.StartsWith("Обновление по ссылке:", Assert.Single(sentTexts), StringComparison.Ordinal); + } + + [Fact] + public async Task NotifyAsync_WhenSystemReport_StripsMarkerAndSkipsPriorityHeader() + { + var (botClient, sentTexts) = CreateBotClient(); + var sut = new LinkUpdateNotifier(botClient, Substitute.For()); + + await sut.NotifyAsync( + new LinkUpdate + { + Id = 0, + Url = Url, + Description = $"{SystemMessageMarkers.FailedLinkReport}Не удалось проверить часть ссылок", + TgChatIds = [42], + Priority = LinkUpdatePriority.High + }, + CancellationToken.None); + + var text = Assert.Single(sentTexts); + + Assert.Equal("Не удалось проверить часть ссылок", text); + } + + private static (ITelegramBotClient BotClient, List SentTexts) CreateBotClient() + { + var botClient = Substitute.For(); + var sentTexts = new List(); + + botClient + .SendRequest(Arg.Any(), Arg.Any()) + .Returns(call => + { + sentTexts.Add(call.Arg().Text); + return Task.FromResult(new Message()); + }); + + return (botClient, sentTexts); + } +} diff --git a/src/LinkTracker.Tests/Scrapper/Integration/Cache/CacheModuleValkeyTests.cs b/src/LinkTracker.Tests/Scrapper/Integration/Cache/CacheModuleValkeyTests.cs index acbd24b..e75f52d 100644 --- a/src/LinkTracker.Tests/Scrapper/Integration/Cache/CacheModuleValkeyTests.cs +++ b/src/LinkTracker.Tests/Scrapper/Integration/Cache/CacheModuleValkeyTests.cs @@ -50,7 +50,7 @@ public async Task AddCache_WhenValkeyEnabled_ResolvesValkeyImplementationsAndCac Size = 1, Links = [ - new LinkResponse { Id = 42, Url = new Uri("https://github.com/user/repo"), Tags = ["backend"], Filters = [] } + new LinkResponse { Id = 42, Url = new Uri("https://github.com/user/repo"), Tags = ["backend"] } ] }; diff --git a/src/LinkTracker.Tests/Scrapper/Integration/Cache/ScrapperLinksCacheApiTests.cs b/src/LinkTracker.Tests/Scrapper/Integration/Cache/ScrapperLinksCacheApiTests.cs index 6040332..b5a7540 100644 --- a/src/LinkTracker.Tests/Scrapper/Integration/Cache/ScrapperLinksCacheApiTests.cs +++ b/src/LinkTracker.Tests/Scrapper/Integration/Cache/ScrapperLinksCacheApiTests.cs @@ -157,15 +157,12 @@ public async Task PostLinks_WhenUserLinksChanged_InvalidatesCachedList() chatId, Arg.Any(), Arg.Any>(), - Arg.Any>(), Arg.Any()) .Returns(call => { var link = call.ArgAt(1); var tags = call.ArgAt>(2); - var filters = call.ArgAt>(3); - - var record = CreateRecord(2, link.ToString(), tags, filters); + var record = CreateRecord(2, link.ToString(), tags); currentLinks = [.. currentLinks, record]; return Task.FromResult(record); @@ -181,7 +178,7 @@ public async Task PostLinks_WhenUserLinksChanged_InvalidatesCachedList() using var postRequest = new HttpRequestMessage(HttpMethod.Post, "/links"); postRequest.Headers.Add("Tg-Chat-Id", chatId.ToString()); - postRequest.Content = JsonContent.Create(new AddLinkRequest { Link = addedUrl, Tags = ["new"], Filters = [] }); + postRequest.Content = JsonContent.Create(new AddLinkRequest { Link = addedUrl, Tags = ["new"] }); var postResponse = await client.SendAsync(postRequest); @@ -198,7 +195,6 @@ await service.Received(1).AddLinkAsync( chatId, addedUrl, Arg.Any>(), - Arg.Any>(), Arg.Any()); } @@ -315,10 +311,9 @@ private static async Task GetLinksAsync(HttpClient client, lo private static TrackedLinkRecord CreateRecord( long id, string url, - IReadOnlyList? tags = null, - IReadOnlyList? filters = null) + IReadOnlyList? tags = null) { - return new TrackedLinkRecord { Id = id, Url = new Uri(url), Tags = tags ?? [], Filters = filters ?? [] }; + return new TrackedLinkRecord { Id = id, Url = new Uri(url), Tags = tags ?? [] }; } private static string BuildLinksCacheKey(long chatId) diff --git a/src/LinkTracker.Tests/Scrapper/Integration/Cache/ValkeyLinksResponseCacheTests.cs b/src/LinkTracker.Tests/Scrapper/Integration/Cache/ValkeyLinksResponseCacheTests.cs index 9ac4f9e..3f6c6bf 100644 --- a/src/LinkTracker.Tests/Scrapper/Integration/Cache/ValkeyLinksResponseCacheTests.cs +++ b/src/LinkTracker.Tests/Scrapper/Integration/Cache/ValkeyLinksResponseCacheTests.cs @@ -27,8 +27,8 @@ public async Task SetAndGetAsync_WhenResponseCached_RoundTripsListLinksResponse( Size = 2, Links = [ - new LinkResponse { Id = 42, Url = new Uri("https://github.com/user/repo"), Tags = ["backend", "dotnet"], Filters = ["user:alice"] }, - new LinkResponse { Id = 43, Url = new Uri("https://stackoverflow.com/questions/123"), Tags = ["qa"], Filters = [] } + new LinkResponse { Id = 42, Url = new Uri("https://github.com/user/repo"), Tags = ["backend", "dotnet"] }, + new LinkResponse { Id = 43, Url = new Uri("https://stackoverflow.com/questions/123"), Tags = ["qa"] } ] }; @@ -46,14 +46,12 @@ public async Task SetAndGetAsync_WhenResponseCached_RoundTripsListLinksResponse( Assert.Equal(42, first.Id); Assert.Equal(new Uri("https://github.com/user/repo"), first.Url); Assert.Equal(["backend", "dotnet"], first.Tags); - Assert.Equal(["user:alice"], first.Filters); }, second => { Assert.Equal(43, second.Id); Assert.Equal(new Uri("https://stackoverflow.com/questions/123"), second.Url); Assert.Equal(["qa"], second.Tags); - Assert.Empty(second.Filters); }); } @@ -72,7 +70,7 @@ await sut.SetAsync( Size = 1, Links = [ - new LinkResponse { Id = 1, Url = new Uri("https://github.com/user/repo"), Tags = [], Filters = [] } + new LinkResponse { Id = 1, Url = new Uri("https://github.com/user/repo"), Tags = [] } ] }, CancellationToken.None); @@ -128,7 +126,7 @@ await sut.SetAsync( Size = 1, Links = [ - new LinkResponse { Id = 42, Url = new Uri("https://github.com/user/repo"), Tags = ["backend"], Filters = [] } + new LinkResponse { Id = 42, Url = new Uri("https://github.com/user/repo"), Tags = ["backend"] } ] }, CancellationToken.None); diff --git a/src/LinkTracker.Tests/Scrapper/Integration/Outbox/PostgresOutboxStoreTests.cs b/src/LinkTracker.Tests/Scrapper/Integration/Outbox/PostgresOutboxStoreTests.cs index 1eaceb6..bd61f7a 100644 --- a/src/LinkTracker.Tests/Scrapper/Integration/Outbox/PostgresOutboxStoreTests.cs +++ b/src/LinkTracker.Tests/Scrapper/Integration/Outbox/PostgresOutboxStoreTests.cs @@ -12,6 +12,8 @@ namespace LinkTracker.Tests.Scrapper.Integration.Outbox; [Collection("Postgres collection")] public sealed class PostgresOutboxStoreTests(PostgresSqlStorageFixture fixture) { + private static readonly TimeSpan DefaultLock = TimeSpan.FromSeconds(60); + [Fact] public async Task AddRangeAndSetCursorAsync_WhenLinkExists_AddsMessagesAndUpdatesCursor() { @@ -26,7 +28,7 @@ public async Task AddRangeAndSetCursorAsync_WhenLinkExists_AddsMessagesAndUpdate const string eventKey = "issue:123"; await trackingStore.TryRegisterChatAsync(chatId); - var tracked = await trackingStore.TryAddAsync(chatId, url, ["backend"], []); + var tracked = await trackingStore.TryAddAsync(chatId, url, ["backend"]); Assert.NotNull(tracked); @@ -39,7 +41,7 @@ await sut.AddRangeAndSetCursorAsync( ], CancellationToken.None); - var messages = await sut.GetUnprocessedBatchAsync(10, 3, CancellationToken.None); + var messages = await sut.ClaimUnprocessedBatchAsync(10, 3, DefaultLock, CancellationToken.None); var message = Assert.Single(messages); Assert.Equal(tracked.Id, message.Payload.Id); @@ -73,7 +75,7 @@ await Assert.ThrowsAsync(() => ], CancellationToken.None)); - var messages = await sut.GetUnprocessedBatchAsync(10, 3, CancellationToken.None); + var messages = await sut.ClaimUnprocessedBatchAsync(10, 3, DefaultLock, CancellationToken.None); Assert.Empty(messages); } @@ -89,12 +91,12 @@ await sut.AddAsync( new LinkUpdate { Id = 1, Url = new Uri("https://github.com/user/repo"), Description = "Repository updated", TgChatIds = [1001] }, CancellationToken.None); - var before = await sut.GetUnprocessedBatchAsync(10, 3, CancellationToken.None); + var before = await sut.ClaimUnprocessedBatchAsync(10, 3, DefaultLock, CancellationToken.None); var message = Assert.Single(before); await sut.MarkProcessedAsync(message.Id, CancellationToken.None); - var after = await sut.GetUnprocessedBatchAsync(10, 3, CancellationToken.None); + var after = await sut.ClaimUnprocessedBatchAsync(10, 3, DefaultLock, CancellationToken.None); Assert.Empty(after); } @@ -110,12 +112,12 @@ await sut.AddAsync( new LinkUpdate { Id = 1, Url = new Uri("https://github.com/user/repo"), Description = "Repository updated", TgChatIds = [1001] }, CancellationToken.None); - var before = await sut.GetUnprocessedBatchAsync(10, 3, CancellationToken.None); + var before = await sut.ClaimUnprocessedBatchAsync(10, 3, DefaultLock, CancellationToken.None); var message = Assert.Single(before); await sut.MarkFailedAsync(message.Id, "Bot is unavailable", CancellationToken.None); - var after = await sut.GetUnprocessedBatchAsync(10, 3, CancellationToken.None); + var after = await sut.ClaimUnprocessedBatchAsync(10, 3, DefaultLock, CancellationToken.None); var failed = Assert.Single(after); Assert.Equal(1, failed.RetryCount); @@ -123,7 +125,7 @@ await sut.AddAsync( } [Fact] - public async Task GetUnprocessedBatchAsync_WhenRetryCountReachedLimit_DoesNotReturnMessage() + public async Task ClaimUnprocessedBatchAsync_WhenRetryCountReachedLimit_DoesNotReturnMessage() { await fixture.ResetAsync(); @@ -133,21 +135,63 @@ await sut.AddAsync( new LinkUpdate { Id = 1, Url = new Uri("https://github.com/user/repo"), Description = "Repository updated", TgChatIds = [1001] }, CancellationToken.None); - var before = await sut.GetUnprocessedBatchAsync(10, 3, CancellationToken.None); + var before = await sut.ClaimUnprocessedBatchAsync(10, 3, DefaultLock, CancellationToken.None); var message = Assert.Single(before); await sut.MarkFailedAsync(message.Id, "first", CancellationToken.None); await sut.MarkFailedAsync(message.Id, "second", CancellationToken.None); await sut.MarkFailedAsync(message.Id, "third", CancellationToken.None); - var after = await sut.GetUnprocessedBatchAsync( + var after = await sut.ClaimUnprocessedBatchAsync( 10, 3, + DefaultLock, CancellationToken.None); Assert.Empty(after); } + [Fact] + public async Task ClaimUnprocessedBatchAsync_WhenSecondDispatcherClaims_DoesNotReturnAlreadyClaimedMessages() + { + await fixture.ResetAsync(); + + var first = CreateSut(); + var second = CreateSut(); + + foreach (var id in new[] { 1L, 2L, 3L }) + { + await first.AddAsync( + new LinkUpdate { Id = id, Url = new Uri($"https://github.com/user/repo-{id}"), Description = "Repository updated", TgChatIds = [1001] }, + CancellationToken.None); + } + + var firstBatch = await first.ClaimUnprocessedBatchAsync(10, 3, DefaultLock, CancellationToken.None); + var secondBatch = await second.ClaimUnprocessedBatchAsync(10, 3, DefaultLock, CancellationToken.None); + + Assert.Equal(3, firstBatch.Count); + Assert.Empty(secondBatch); + } + + [Fact] + public async Task ClaimUnprocessedBatchAsync_WhenLeaseExpired_ReturnsMessageAgain() + { + await fixture.ResetAsync(); + + var sut = CreateSut(); + + await sut.AddAsync( + new LinkUpdate { Id = 1, Url = new Uri("https://github.com/user/repo"), Description = "Repository updated", TgChatIds = [1001] }, + CancellationToken.None); + + var claimed = await sut.ClaimUnprocessedBatchAsync(10, 3, TimeSpan.Zero, CancellationToken.None); + var reclaimed = await sut.ClaimUnprocessedBatchAsync(10, 3, DefaultLock, CancellationToken.None); + + Assert.Single(claimed); + Assert.Single(reclaimed); + Assert.Equal(claimed[0].Id, reclaimed[0].Id); + } + private PostgresOutboxStore CreateSut() { return new PostgresOutboxStore( diff --git a/src/LinkTracker.Tests/Scrapper/Integration/Storage/DatabaseMigrationTests.cs b/src/LinkTracker.Tests/Scrapper/Integration/Storage/DatabaseMigrationTests.cs index 35dcd2d..f4cf0fb 100644 --- a/src/LinkTracker.Tests/Scrapper/Integration/Storage/DatabaseMigrationTests.cs +++ b/src/LinkTracker.Tests/Scrapper/Integration/Storage/DatabaseMigrationTests.cs @@ -19,12 +19,14 @@ public async Task Migrations_CreateExpectedSchema() Assert.Contains("links", tables); Assert.Contains("subscriptions", tables); Assert.Contains("tags", tables); - Assert.Contains("filters", tables); Assert.Contains("subscription_tags", tables); - Assert.Contains("subscription_filters", tables); Assert.Contains("dbup_schema_versions", tables); Assert.Contains("outbox_messages", tables); + // 005_drop_filters.sql: фича фильтров удалена, таблицы не должны возвращаться. + Assert.DoesNotContain("filters", tables); + Assert.DoesNotContain("subscription_filters", tables); + Assert.Contains("ix_subscriptions_chat_id", indexes); Assert.Contains("ix_subscriptions_link_id", indexes); Assert.Contains("ix_links_normalized_url", indexes); @@ -42,9 +44,7 @@ public async Task Migrations_CreateExpectedConstraints() Assert.Contains("links_normalized_url_key", uniqueConstraints); Assert.Contains("subscriptions_chat_id_link_id_key", uniqueConstraints); Assert.Contains("tags_name_key", uniqueConstraints); - Assert.Contains("filters_value_key", uniqueConstraints); Assert.Contains("subscription_tags_pkey", uniqueConstraints); - Assert.Contains("subscription_filters_pkey", uniqueConstraints); } private static async Task> GetTableNamesAsync(NpgsqlConnection connection) diff --git a/src/LinkTracker.Tests/Scrapper/Integration/Storage/LinkTrackingStoreContractTests.cs b/src/LinkTracker.Tests/Scrapper/Integration/Storage/LinkTrackingStoreContractTests.cs index 58cba2e..8b0a9a8 100644 --- a/src/LinkTracker.Tests/Scrapper/Integration/Storage/LinkTrackingStoreContractTests.cs +++ b/src/LinkTracker.Tests/Scrapper/Integration/Storage/LinkTrackingStoreContractTests.cs @@ -40,20 +40,18 @@ await ExecuteWithSut(async sut => await sut.TryRegisterChatAsync(chatId); - var added = await sut.TryAddAsync(chatId, url, tags, []); + var added = await sut.TryAddAsync(chatId, url, tags); var links = await sut.GetAllTrackedLinkRecordsAsync(chatId); Assert.NotNull(added); Assert.Equal(url, added!.Url); Assert.Equal(["backend", "dotnet"], added.Tags.OrderBy(x => x).ToArray()); - Assert.Empty(added.Filters); Assert.Null(added.LastUpdatedAt); Assert.Null(added.LastEventKey); var only = Assert.Single(links); Assert.Equal(url, only.Url); Assert.Equal(["backend", "dotnet"], only.Tags.OrderBy(x => x).ToArray()); - Assert.Empty(only.Filters); Assert.Null(only.LastUpdatedAt); Assert.Null(only.LastEventKey); }); @@ -69,8 +67,8 @@ await ExecuteWithSut(async sut => await sut.TryRegisterChatAsync(chatId); - var first = await sut.TryAddAsync(chatId, url, ["tag1"], []); - var second = await sut.TryAddAsync(chatId, url, ["tag1"], []); + var first = await sut.TryAddAsync(chatId, url, ["tag1"]); + var second = await sut.TryAddAsync(chatId, url, ["tag1"]); Assert.NotNull(first); Assert.Null(second); @@ -89,8 +87,8 @@ await ExecuteWithSut(async sut => await sut.TryRegisterChatAsync(firstChatId); await sut.TryRegisterChatAsync(secondChatId); - await sut.TryAddAsync(firstChatId, url, ["alpha"], []); - await sut.TryAddAsync(secondChatId, url, ["beta"], []); + await sut.TryAddAsync(firstChatId, url, ["alpha"]); + await sut.TryAddAsync(secondChatId, url, ["beta"]); var removed = await sut.TryRemoveAsync(firstChatId, url); var firstChatLinks = await sut.GetAllTrackedLinkRecordsAsync(firstChatId); @@ -139,7 +137,7 @@ await ExecuteWithSut(async sut => const string eventKey = "issue:123"; await sut.TryRegisterChatAsync(chatId); - var added = await sut.TryAddAsync(chatId, url, ["sync"], []); + var added = await sut.TryAddAsync(chatId, url, ["sync"]); Assert.NotNull(added); await sut.SetCursorAsync(added!.Id, updatedAt, eventKey); @@ -172,8 +170,8 @@ await ExecuteWithSut(async sut => await sut.TryRegisterChatAsync(firstChatId); await sut.TryRegisterChatAsync(secondChatId); - var firstAdded = await sut.TryAddAsync(firstChatId, url, ["alpha"], []); - var secondAdded = await sut.TryAddAsync(secondChatId, url, ["beta"], []); + var firstAdded = await sut.TryAddAsync(firstChatId, url, ["alpha"]); + var secondAdded = await sut.TryAddAsync(secondChatId, url, ["beta"]); Assert.NotNull(firstAdded); Assert.NotNull(secondAdded); @@ -213,7 +211,7 @@ await ExecuteWithSut(async sut => var updatedAt = new DateTimeOffset(2026, 3, 22, 12, 0, 0, TimeSpan.Zero); await sut.TryRegisterChatAsync(chatId); - var added = await sut.TryAddAsync(chatId, url, ["baseline"], []); + var added = await sut.TryAddAsync(chatId, url, ["baseline"]); Assert.NotNull(added); await sut.SetCursorAsync(added!.Id, updatedAt, null); @@ -241,7 +239,7 @@ await ExecuteWithSut(async sut => var url = new Uri("https://github.com/user/repo-tags"); await sut.TryRegisterChatAsync(chatId); - await sut.TryAddAsync(chatId, url, ["alpha"], []); + await sut.TryAddAsync(chatId, url, ["alpha"]); var afterAddTag = await sut.TryAddTagAsync(chatId, url, "beta"); var tagsAfterAdd = await sut.GetTagsAsync(chatId); @@ -283,7 +281,7 @@ await ExecuteWithSut(async sut => if (scenario is not "add-tag-missing-subscription") { - await sut.TryAddAsync(chatId, url, ["alpha"], []); + await sut.TryAddAsync(chatId, url, ["alpha"]); } switch (scenario) @@ -324,7 +322,7 @@ await ExecuteWithSut(async sut => var url = new Uri("https://github.com/user/repo"); await sut.TryRegisterChatAsync(chatId); - await sut.TryAddAsync(chatId, url, ["alpha", "beta"], []); + await sut.TryAddAsync(chatId, url, ["alpha", "beta"]); var renamed = await sut.TryRenameTagAsync(chatId, "alpha", "beta"); var links = await sut.GetAllTrackedLinkRecordsAsync(chatId); @@ -354,19 +352,16 @@ await ExecuteWithSut(async sut => var first = await sut.TryAddAsync( firstChatId, new Uri("https://github.com/user/repo-1"), - [], []); var second = await sut.TryAddAsync( secondChatId, new Uri("https://github.com/user/repo-2"), - [], []); var third = await sut.TryAddAsync( thirdChatId, new Uri("https://github.com/user/repo-3"), - [], []); Assert.NotNull(first); @@ -400,19 +395,16 @@ await ExecuteWithSut(async sut => var first = await sut.TryAddAsync( firstChatId, new Uri("https://github.com/user/repo-1"), - [], []); var second = await sut.TryAddAsync( secondChatId, new Uri("https://github.com/user/repo-2"), - [], []); var third = await sut.TryAddAsync( thirdChatId, new Uri("https://github.com/user/repo-3"), - [], []); Assert.NotNull(first); @@ -445,10 +437,10 @@ await ExecuteWithSut(async sut => await sut.TryRegisterChatAsync(thirdChatId); await sut.TryRegisterChatAsync(fourthChatId); - await sut.TryAddAsync(firstChatId, new Uri("https://github.com/user/repo-1"), [], []); - await sut.TryAddAsync(secondChatId, new Uri("https://github.com/user/repo-2"), [], []); - await sut.TryAddAsync(thirdChatId, new Uri("https://github.com/user/repo-3"), [], []); - await sut.TryAddAsync(fourthChatId, new Uri("https://github.com/user/repo-4"), [], []); + await sut.TryAddAsync(firstChatId, new Uri("https://github.com/user/repo-1"), []); + await sut.TryAddAsync(secondChatId, new Uri("https://github.com/user/repo-2"), []); + await sut.TryAddAsync(thirdChatId, new Uri("https://github.com/user/repo-3"), []); + await sut.TryAddAsync(fourthChatId, new Uri("https://github.com/user/repo-4"), []); var batch = await sut.GetSubscriptionsBatchAsync(null, batchSize); @@ -468,8 +460,8 @@ await ExecuteWithSut(async sut => await sut.TryRegisterChatAsync(firstChatId); await sut.TryRegisterChatAsync(secondChatId); - var first = await sut.TryAddAsync(firstChatId, url, ["alpha"], []); - var second = await sut.TryAddAsync(secondChatId, url, ["beta"], []); + var first = await sut.TryAddAsync(firstChatId, url, ["alpha"]); + var second = await sut.TryAddAsync(secondChatId, url, ["beta"]); Assert.NotNull(first); Assert.NotNull(second); @@ -499,7 +491,7 @@ await ExecuteWithSut(async sut => await sut.TryRegisterChatAsync(chatId); - var added = await sut.TryAddAsync(chatId, url, [], []); + var added = await sut.TryAddAsync(chatId, url, []); Assert.NotNull(added); await sut.SetCursorAsync(added!.Id, updatedAt, eventKey); diff --git a/src/LinkTracker.Tests/Scrapper/Integration/Storage/Orm/OrmLinkTrackingStoreTests.cs b/src/LinkTracker.Tests/Scrapper/Integration/Storage/Orm/OrmLinkTrackingStoreTests.cs index 19850e3..e992f57 100644 --- a/src/LinkTracker.Tests/Scrapper/Integration/Storage/Orm/OrmLinkTrackingStoreTests.cs +++ b/src/LinkTracker.Tests/Scrapper/Integration/Storage/Orm/OrmLinkTrackingStoreTests.cs @@ -13,8 +13,7 @@ protected override async Task ExecuteWithSut(Func test { await fixture.ResetAsync(); - await using var dbContext = fixture.CreateDbContext(); - ILinkTrackingStore sut = new OrmLinkTrackingStore(dbContext); + ILinkTrackingStore sut = new OrmLinkTrackingStore(fixture.CreateDbContextFactory()); await test(sut); } diff --git a/src/LinkTracker.Tests/Scrapper/Integration/Storage/PostgresSqlStorageFixture.cs b/src/LinkTracker.Tests/Scrapper/Integration/Storage/PostgresSqlStorageFixture.cs index a14a249..f4ce7eb 100644 --- a/src/LinkTracker.Tests/Scrapper/Integration/Storage/PostgresSqlStorageFixture.cs +++ b/src/LinkTracker.Tests/Scrapper/Integration/Storage/PostgresSqlStorageFixture.cs @@ -85,11 +85,19 @@ public async Task DisposeAsync() public AppDbContext CreateDbContext() { - var options = new DbContextOptionsBuilder() + return new AppDbContext(BuildDbContextOptions()); + } + + public IDbContextFactory CreateDbContextFactory() + { + return new TestDbContextFactory(BuildDbContextOptions()); + } + + private DbContextOptions BuildDbContextOptions() + { + return new DbContextOptionsBuilder() .UseNpgsql(DataSource) .Options; - - return new AppDbContext(options); } public async Task ResetAsync() @@ -99,10 +107,8 @@ public async Task ResetAsync() truncate table outbox_messages, subscription_tags, - subscription_filters, subscriptions, tags, - filters, links, chats restart identity cascade; @@ -141,6 +147,15 @@ private static async Task WaitUntilDatabaseReadyAsync(string connectionString) throw new InvalidOperationException("PostgreSQL did not become ready in time."); } + private sealed class TestDbContextFactory(DbContextOptions options) + : IDbContextFactory + { + public AppDbContext CreateDbContext() + { + return new AppDbContext(options); + } + } + private sealed class TestWebHostEnvironment : IWebHostEnvironment { public string ApplicationName { get; set; } = "LinkTracker.Tests"; diff --git a/src/LinkTracker.Tests/Scrapper/Integration/Storage/StorageModuleTests.cs b/src/LinkTracker.Tests/Scrapper/Integration/Storage/StorageModuleTests.cs index ed38797..8e930f3 100644 --- a/src/LinkTracker.Tests/Scrapper/Integration/Storage/StorageModuleTests.cs +++ b/src/LinkTracker.Tests/Scrapper/Integration/Storage/StorageModuleTests.cs @@ -32,7 +32,7 @@ public void AddStorage_ResolvesExpectedImplementation( services.AddSingleton(_ => new NpgsqlDataSourceBuilder(fixture.ConnectionString).Build()); - services.AddDbContext(options => + services.AddDbContextFactory(options => { options.UseNpgsql(fixture.ConnectionString); }); @@ -65,7 +65,7 @@ public async Task AddStorage_ResolvedStore_IsUsable(string accessType) services.AddSingleton(_ => new NpgsqlDataSourceBuilder(fixture.ConnectionString).Build()); - services.AddDbContext(options => + services.AddDbContextFactory(options => { options.UseNpgsql(fixture.ConnectionString); }); diff --git a/src/LinkTracker.Tests/Scrapper/Unit/Application/Services/Tracking/LinkTrackingServiceTests.cs b/src/LinkTracker.Tests/Scrapper/Unit/Application/Services/Tracking/LinkTrackingServiceTests.cs index f41312c..feb3d56 100644 --- a/src/LinkTracker.Tests/Scrapper/Unit/Application/Services/Tracking/LinkTrackingServiceTests.cs +++ b/src/LinkTracker.Tests/Scrapper/Unit/Application/Services/Tracking/LinkTrackingServiceTests.cs @@ -23,7 +23,7 @@ public async Task AddLinkAsync_WhenChatIdIsInvalid_ThrowsInvalidChatId(long chat Uri link = new("https://github.com/test/repo"); var exception = await Assert.ThrowsAsync(() => - sut.AddLinkAsync(chatId, link, [], [])); + sut.AddLinkAsync(chatId, link, [])); Assert.Equal(HttpStatusCode.BadRequest, exception.StatusCode); Assert.Equal("invalid_chat_id", exception.Code); @@ -43,7 +43,7 @@ public async Task AddLinkAsync_WhenLinkIsInvalid_ThrowsValidationError(string ra : new Uri(rawLink); var exception = await Assert.ThrowsAsync(() => - sut.AddLinkAsync(1, link, [], [])); + sut.AddLinkAsync(1, link, [])); Assert.Equal(HttpStatusCode.BadRequest, exception.StatusCode); Assert.Equal(expectedCode, exception.Code); @@ -60,7 +60,7 @@ public async Task AddLinkAsync_WhenLinkIsUnsupported_ThrowsUnsupportedLink() _handler.CanHandle(link).Returns(false); var exception = await Assert.ThrowsAsync(() => - sut.AddLinkAsync(1, link, [], [])); + sut.AddLinkAsync(1, link, [])); Assert.Equal(HttpStatusCode.BadRequest, exception.StatusCode); Assert.Equal("unsupported_link", exception.Code); @@ -79,7 +79,7 @@ public async Task AddLinkAsync_WhenChatDoesNotExist_ThrowsChatNotFound() _store.ChatExistsAsync(1, Arg.Any()).Returns(false); var exception = await Assert.ThrowsAsync(() => - sut.AddLinkAsync(1, link, ["tag"], [])); + sut.AddLinkAsync(1, link, ["tag"])); Assert.Equal(HttpStatusCode.NotFound, exception.StatusCode); Assert.Equal("chat_not_found", exception.Code); @@ -89,7 +89,6 @@ await _store.DidNotReceive() Arg.Any(), Arg.Any(), Arg.Any>(), - Arg.Any>(), Arg.Any()); } @@ -105,12 +104,11 @@ public async Task AddLinkAsync_WhenLinkAlreadyExists_ThrowsLinkAlreadyExists() 1, link, Arg.Any>(), - Arg.Any>(), Arg.Any()) .Returns((TrackedLinkRecord?)null); var exception = await Assert.ThrowsAsync(() => - sut.AddLinkAsync(1, link, ["tag"], ["legacy-filter"])); + sut.AddLinkAsync(1, link, ["tag"])); Assert.Equal(HttpStatusCode.Conflict, exception.StatusCode); Assert.Equal("link_already_exists", exception.Code); @@ -128,7 +126,6 @@ public async Task AddLinkAsync_WhenRequestIsValid_ReturnsTrackedLinkRecord() Id = 10, Url = link, Tags = tags, - Filters = [], LastUpdatedAt = DateTimeOffset.UtcNow }; @@ -138,21 +135,16 @@ public async Task AddLinkAsync_WhenRequestIsValid_ReturnsTrackedLinkRecord() 1, link, Arg.Is>(value => value.SequenceEqual(tags)), - Arg.Is>(value => value.SequenceEqual(new List().AsReadOnly())), Arg.Any()) .Returns(expected); - var result = await sut.AddLinkAsync( - 1, - link, - tags, - ["legacy-filter-1", "legacy-filter-2"]); + var result = await sut.AddLinkAsync(1, link, tags); Assert.Same(expected, result); } [Fact] - public async Task AddLinkAsync_WhenRequestIsValid_PassesTagsAndEmptyFiltersToStore() + public async Task AddLinkAsync_WhenRequestIsValid_PassesTagsToStore() { var sut = CreateSut(); Uri link = new("https://github.com/test/repo"); @@ -164,28 +156,21 @@ public async Task AddLinkAsync_WhenRequestIsValid_PassesTagsAndEmptyFiltersToSto Arg.Any(), Arg.Any(), Arg.Any>(), - Arg.Any>(), Arg.Any()) .Returns(new TrackedLinkRecord { Id = 1, Url = link, Tags = tags, - Filters = new List().AsReadOnly(), LastUpdatedAt = DateTimeOffset.UtcNow }); - await sut.AddLinkAsync( - 1, - link, - tags, - new List { "legacy-filter-1", "legacy-filter-2" }.AsReadOnly()); + await sut.AddLinkAsync(1, link, tags); await _store.Received(1).TryAddAsync( 1, link, Arg.Is>(value => value.SequenceEqual(tags)), - Arg.Is>(value => value.SequenceEqual(new List().AsReadOnly())), Arg.Any()); } diff --git a/src/LinkTracker.Tests/Scrapper/Unit/Infrastructure/Outbox/Jobs/OutboxDispatchJobTests.cs b/src/LinkTracker.Tests/Scrapper/Unit/Infrastructure/Outbox/Jobs/OutboxDispatchJobTests.cs index 25da5c2..1db746e 100644 --- a/src/LinkTracker.Tests/Scrapper/Unit/Infrastructure/Outbox/Jobs/OutboxDispatchJobTests.cs +++ b/src/LinkTracker.Tests/Scrapper/Unit/Infrastructure/Outbox/Jobs/OutboxDispatchJobTests.cs @@ -23,7 +23,7 @@ public async Task Execute_WhenMessagesDoNotExist_DoesNothing() { var outboxStore = Substitute.For(); - outboxStore.GetUnprocessedBatchAsync(100, 3, Arg.Any()) + outboxStore.ClaimUnprocessedBatchAsync(100, 3, Arg.Any(), Arg.Any()) .Returns([]); var sut = CreateSut( @@ -47,7 +47,7 @@ public async Task Execute_WhenTransportSucceeds_MarksMessageProcessed() var outboxStore = Substitute.For(); var message = CreateOutboxMessage(); - outboxStore.GetUnprocessedBatchAsync(100, 3, Arg.Any()) + outboxStore.ClaimUnprocessedBatchAsync(100, 3, Arg.Any(), Arg.Any()) .Returns([message]); var sut = CreateSut( @@ -71,7 +71,7 @@ public async Task Execute_WhenHttpTransportUnavailable_FallsBackToKafkaAndMarksM var outboxStore = Substitute.For(); var message = CreateOutboxMessage(); - outboxStore.GetUnprocessedBatchAsync(100, 3, Arg.Any()) + outboxStore.ClaimUnprocessedBatchAsync(100, 3, Arg.Any(), Arg.Any()) .Returns([message]); var kafkaWasCalled = false; @@ -108,7 +108,7 @@ public async Task Execute_WhenNonRetriableHttpErrorOccurs_MarksMessageFailed() var outboxStore = Substitute.For(); var message = CreateOutboxMessage(); - outboxStore.GetUnprocessedBatchAsync(100, 3, Arg.Any()) + outboxStore.ClaimUnprocessedBatchAsync(100, 3, Arg.Any(), Arg.Any()) .Returns([message]); var sut = CreateSut( @@ -138,7 +138,7 @@ public async Task Execute_WhenCancellationRequested_ThrowsAndDoesNotMarkFailed() var outboxStore = Substitute.For(); var message = CreateOutboxMessage(); - outboxStore.GetUnprocessedBatchAsync(100, 3, Arg.Any()) + outboxStore.ClaimUnprocessedBatchAsync(100, 3, Arg.Any(), Arg.Any()) .Returns([message]); using var cts = new CancellationTokenSource(); @@ -169,7 +169,7 @@ private static OutboxDispatchJob CreateSut( { var botOptions = Options.Create(new BotOptions { Transport = TransportKind.Http }); - var outboxOptions = Options.Create(new OutboxOptions { Enabled = true, DispatchIntervalSeconds = 10, BatchSize = 100, MaxRetryCount = 3 }); + var outboxOptions = Options.Create(new OutboxOptions { Enabled = true, DispatchIntervalSeconds = 10, BatchSize = 100, MaxRetryCount = 3, LockSeconds = 60 }); var botClient = new FallbackBotClient( transportClients, diff --git a/src/LinkTracker.Tests/Shared/Integration/Infrastructure/Telemetry/MetricsEndpointTests.cs b/src/LinkTracker.Tests/Shared/Integration/Infrastructure/Telemetry/MetricsEndpointTests.cs new file mode 100644 index 0000000..ee5b04f --- /dev/null +++ b/src/LinkTracker.Tests/Shared/Integration/Infrastructure/Telemetry/MetricsEndpointTests.cs @@ -0,0 +1,157 @@ +using LinkTracker.Bot.Infrastructure.Telemetry; +using LinkTracker.Scrapper.Infrastructure.Telemetry; +using LinkTracker.Shared.Infrastructure.RateLimiting; +using LinkTracker.Shared.Infrastructure.Telemetry; +using Microsoft.AspNetCore.Builder; +using Microsoft.AspNetCore.Hosting; +using Microsoft.AspNetCore.TestHost; +using Microsoft.Extensions.Configuration; +using Microsoft.Extensions.DependencyInjection; + +namespace LinkTracker.Tests.Shared.Integration.Infrastructure.Telemetry; + +[Trait("Module", "Shared")] +[Trait("Category", "Integration")] +public sealed class MetricsEndpointTests +{ + [Fact] + public async Task GetMetrics_ExposesScrapperInstrumentNamesUsedByDashboard() + { + using var metrics = new ScrapperMetrics(); + + using var server = CreateServer(ScrapperMetrics.MeterName); + using var client = server.CreateClient(); + + RecordScrapperSamples(metrics); + + var families = await GetMetricFamilyNamesAsync(client); + + Assert.Contains("api_requests_total", families); + Assert.Contains("sent_updates_total", families); + Assert.Contains("errors_total", families); + Assert.Contains("db_queries_total", families); + Assert.Contains("db_errors_total", families); + Assert.Contains("kafka_produced_total", families); + Assert.Contains("kafka_produce_errors_total", families); + Assert.Contains("request_duration_ms_total", families); + Assert.Contains("db_query_duration_ms_total", families); + Assert.Contains("kafka_produce_duration_ms_total", families); + Assert.Contains("links_on_track_total", families); + Assert.Contains("process_memory_working_set_bytes", families); + Assert.Contains("process_memory_managed_bytes", families); + } + + [Fact] + public async Task GetMetrics_ExposesBotInstrumentNamesUsedByDashboard() + { + using var metrics = new BotMetrics(); + + using var server = CreateServer(BotMetrics.MeterName); + using var client = server.CreateClient(); + + RecordBotSamples(metrics); + + var families = await GetMetricFamilyNamesAsync(client); + + Assert.Contains("bot_requests_total", families); + Assert.Contains("command_requests_total", families); + Assert.Contains("sent_notification_total", families); + Assert.Contains("errors_total", families); + Assert.Contains("kafka_consumed_total", families); + Assert.Contains("kafka_consume_errors_total", families); + Assert.Contains("command_duration_ms_total", families); + Assert.Contains("scrapper_call_duration_ms_total", families); + Assert.Contains("kafka_consume_duration_ms_total", families); + Assert.Contains("process_memory_working_set_bytes", families); + } + + [Fact] + public async Task GetMetrics_WhenRateLimiterIsEnabled_IsNotThrottled() + { + using var metrics = new ScrapperMetrics(); + + using var server = CreateServer(ScrapperMetrics.MeterName); + using var client = server.CreateClient(); + + RecordScrapperSamples(metrics); + + using var first = await client.GetAsync("/metrics"); + using var second = await client.GetAsync("/metrics"); + + first.EnsureSuccessStatusCode(); + second.EnsureSuccessStatusCode(); + } + + private static void RecordScrapperSamples(ScrapperMetrics metrics) + { + metrics.ApiRequests.Add(1, new KeyValuePair("source", "/links")); + metrics.SentUpdates.Add(1); + metrics.Errors.Add( + 1, + new KeyValuePair("scope", "http_api"), + new KeyValuePair("scope_type", "/links"), + new KeyValuePair("reason", "5xx")); + metrics.DbQueries.Add(1, new KeyValuePair("operation", "add")); + metrics.DbErrors.Add(1, new KeyValuePair("operation", "add")); + metrics.KafkaProduced.Add(1, new KeyValuePair("topic", "link.raw-updates")); + metrics.KafkaProduceErrors.Add(1, new KeyValuePair("topic", "link.raw-updates")); + + metrics.RequestDuration.Record( + 1, + new KeyValuePair("scope", "http_api"), + new KeyValuePair("scope_type", "/links")); + metrics.DbQueryDuration.Record(1, new KeyValuePair("operation", "add")); + metrics.KafkaProduceDuration.Record(1, new KeyValuePair("topic", "link.raw-updates")); + + metrics.SetLinksOnTrack("github.com", 1); + } + + private static void RecordBotSamples(BotMetrics metrics) + { + metrics.IncrementRequest("Command"); + metrics.IncrementCommand("start"); + metrics.IncrementSentNotifications(); + metrics.IncrementError("bot_command", "start", "exception"); + metrics.IncrementKafkaConsumed("link.processed-updates"); + metrics.IncrementKafkaConsumeError("link.processed-updates"); + + metrics.ObserveCommandDuration("bot_command", "start", 1); + metrics.ObserveScrapperCallDuration("scrapper_sync_api", "GetLinksAsync", 1); + metrics.ObserveKafkaConsumeDuration("link.processed-updates", 1); + } + + private static async Task> GetMetricFamilyNamesAsync(HttpClient client) + { + using var response = await client.GetAsync("/metrics"); + response.EnsureSuccessStatusCode(); + + var payload = await response.Content.ReadAsStringAsync(); + + return payload + .Split('\n') + .Where(line => line.StartsWith("# TYPE ", StringComparison.Ordinal)) + .Select(line => line.Split(' ')[2]) + .ToHashSet(StringComparer.Ordinal); + } + + private static TestServer CreateServer(string meterName) + { + var configuration = new ConfigurationBuilder() + .AddInMemoryCollection(new Dictionary { ["RateLimiting:PermitLimit"] = "1", ["RateLimiting:WindowSeconds"] = "60", ["RateLimiting:SegmentsPerWindow"] = "1", ["RateLimiting:QueueLimit"] = "0" }) + .Build(); + + return new TestServer(new WebHostBuilder() + .ConfigureServices(services => + { + services.AddRouting(); + services.AddIpRateLimiting(configuration); + services.AddOpenTelemetryMetrics("test", meterName); + }) + .Configure(app => + { + app.UseRouting(); + app.UseRateLimiter(); + app.UseEndpoints(endpoints => endpoints.MapMetricsEndpoint()); + })); + } +} \ No newline at end of file diff --git a/src/LinkTracker.Tests/Shared/Unit/Contracts/ProcessedUpdateWireCompatibilityTests.cs b/src/LinkTracker.Tests/Shared/Unit/Contracts/ProcessedUpdateWireCompatibilityTests.cs new file mode 100644 index 0000000..9d5074d --- /dev/null +++ b/src/LinkTracker.Tests/Shared/Unit/Contracts/ProcessedUpdateWireCompatibilityTests.cs @@ -0,0 +1,45 @@ +using LinkTracker.AiAgent.Infrastructure.Kafka.Serialization; +using LinkTracker.Bot.Infrastructure.Kafka.Deserialization; +using LinkTracker.Shared.Contracts.AiAgent; + +namespace LinkTracker.Tests.Shared.Unit.Contracts; + +/// +/// AI-агент публикует ProcessedLinkUpdate, а Bot читает то же сообщение как LinkUpdate. +/// Контракты разные, связь между ними — только формат на проводе, поэтому она проверяется явно. +/// +[Trait("Module", "Shared")] +[Trait("Category", "Unit")] +public sealed class ProcessedUpdateWireCompatibilityTests +{ + private const string Topic = "link.processed-updates"; + + [Theory] + [InlineData(LinkUpdatePriority.High)] + [InlineData(LinkUpdatePriority.Medium)] + [InlineData(LinkUpdatePriority.Low)] + public async Task BotDeserializer_ReadsPriorityPublishedByAiAgent(LinkUpdatePriority priority) + { + var published = new ProcessedLinkUpdate + { + Id = 42, + Url = new Uri("https://github.com/dotnet/runtime"), + Description = "Новый issue", + TgChatIds = [1001], + Priority = priority + }; + + 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(priority, received!.Priority); + Assert.Equal(published.Id, received.Id); + Assert.Equal(published.Url, received.Url); + Assert.Equal(published.Description, received.Description); + Assert.Equal(published.TgChatIds, received.TgChatIds); + } +}