using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
using TeleWave.Application.Common.Interfaces;
using TeleWave.Application.Notifications;
using TeleWave.Domain.Notifications;
namespace TeleWave.Infrastructure.Notifications;
///
/// Опрос обновлений (getUpdates) — режим для установки без публичного адреса.
///
/// Служба крутится всегда, а не заводится по настройке: бота включают и выключают в рантайме,
/// и перезапускать контейнер ради этого никто не должен. Пока бот выключен или выбран вебхук,
/// цикл просто дремлет.
///
public sealed class TelegramPollingBackgroundService(
IServiceScopeFactory scopeFactory,
ILogger logger
) : BackgroundService
{
/// Пауза, когда работать нечего: бот выключен, не настроен или выбран вебхук.
private static readonly TimeSpan IdleDelay = TimeSpan.FromSeconds(30);
/// Пауза после ошибки — чтобы недоступный прокси не превращался в поток запросов.
private static readonly TimeSpan ErrorDelay = TimeSpan.FromSeconds(15);
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
while (!stoppingToken.IsCancellationRequested)
{
TimeSpan delay;
try
{
delay = await TickAsync(stoppingToken);
}
catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
{
return;
}
catch (Exception exception)
{
logger.LogError(exception, "Ошибка опроса Telegram");
delay = ErrorDelay;
}
if (delay > TimeSpan.Zero)
{
try
{
await Task.Delay(delay, stoppingToken);
}
catch (OperationCanceledException)
{
return;
}
}
}
}
/// Один заход опроса. Возвращает, сколько ждать перед следующим.
private async Task TickAsync(CancellationToken cancellationToken)
{
await using var scope = scopeFactory.CreateAsyncScope();
var db = scope.ServiceProvider.GetRequiredService();
var settings = await db.TelegramSettings.FirstOrDefaultAsync(cancellationToken);
if (
settings is null
|| !settings.IsOperational
|| settings.Transport != TelegramTransport.Polling
)
return IdleDelay;
var api = scope.ServiceProvider.GetRequiredService();
var bot = scope.ServiceProvider.GetRequiredService();
IReadOnlyList updates;
try
{
// Смещение — «следующее за последним разобранным»: так Telegram сам подтверждает приём
// и не отдаёт то же обновление второй раз.
updates = await api.GetUpdatesAsync(
settings,
settings.LastUpdateId + 1,
cancellationToken
);
}
catch (Exception exception) when (exception is not OperationCanceledException)
{
settings.MarkError(exception.Message);
await db.SaveChangesAsync(cancellationToken);
logger.LogWarning(exception, "Telegram недоступен");
return ErrorDelay;
}
settings.MarkContact(DateTimeOffset.UtcNow);
foreach (var update in updates)
{
try
{
await bot.HandleAsync(settings, update, DateTimeOffset.UtcNow, cancellationToken);
}
catch (Exception exception) when (exception is not OperationCanceledException)
{
// Одно кривое обновление не должно останавливать разбор остальных — и, главное,
// не должно повторяться вечно: смещение двигаем в любом случае.
logger.LogError(exception, "Не разобрано обновление {UpdateId}", update.UpdateId);
}
settings.AdvanceUpdateId(update.UpdateId);
}
await db.SaveChangesAsync(cancellationToken);
// Длинный опрос сам держит соединение — своей паузы между заходами не нужно.
return TimeSpan.Zero;
}
}