subscribe

This commit is contained in:
kislovdm
2026-03-10 16:24:21 +03:00
parent 4c60efd689
commit e768ee6604
9 changed files with 618 additions and 40 deletions
@@ -0,0 +1,30 @@
using PepperBot.Domain;
namespace PepperBot.Application.Interfaces;
public interface ISubscriptionRepository
{
Task InitializeAsync(CancellationToken cancellationToken);
/// <summary>
/// Получить все активные подписки.
/// </summary>
Task<IReadOnlyList<Subscription>> GetAllActiveAsync(CancellationToken cancellationToken);
/// <summary>
/// Добавить новое правило подписки для указанного чата.
/// Одно правило = один кейворд (но у одного чата может быть много правил).
/// </summary>
Task AddSubscriptionAsync(string chatId, string keyword, CancellationToken cancellationToken);
/// <summary>
/// Получить активные подписки для конкретного чата.
/// </summary>
Task<IReadOnlyList<Subscription>> GetByChatAsync(string chatId, CancellationToken cancellationToken);
/// <summary>
/// Удалить (деактивировать) подписку по идентификатору.
/// </summary>
Task DeleteSubscriptionAsync(long id, CancellationToken cancellationToken);
}
@@ -4,6 +4,19 @@ namespace PepperBot.Application.Interfaces;
public interface ITelegramNotifier
{
/// <summary>
/// Уведомление о новой скидке в общий чат (из переменной окружения TELEGRAM_CHAT_ID).
/// </summary>
Task NotifyNewDealAsync(Deal deal, CancellationToken cancellationToken);
/// <summary>
/// Уведомление о скидке в конкретный Telegram-чат (личная/персональная рассылка).
/// </summary>
Task NotifyDealToChatAsync(Deal deal, string chatId, CancellationToken cancellationToken);
/// <summary>
/// Запустить обработку входящих сообщений Telegram и регистрацию подписок.
/// </summary>
Task StartAsync(CancellationToken cancellationToken);
}
@@ -8,6 +8,7 @@ public class DealService
{
private readonly IPepperClient _pepperClient;
private readonly IDealRepository _repository;
private readonly ISubscriptionRepository _subscriptionRepository;
private readonly ITelegramNotifier _notifier;
private readonly ILogger<DealService> _logger;
private readonly IHealthMonitor _healthMonitor;
@@ -15,12 +16,14 @@ public class DealService
public DealService(
IPepperClient pepperClient,
IDealRepository repository,
ISubscriptionRepository subscriptionRepository,
ITelegramNotifier notifier,
ILogger<DealService> logger,
IHealthMonitor healthMonitor)
{
_pepperClient = pepperClient;
_repository = repository;
_subscriptionRepository = subscriptionRepository;
_notifier = notifier;
_logger = logger;
_healthMonitor = healthMonitor;
@@ -36,6 +39,9 @@ public class DealService
_healthMonitor.ReportSuccess(DateTimeOffset.UtcNow);
}
// Загружаем активные подписки один раз на проход
var subscriptions = await _subscriptionRepository.GetAllActiveAsync(cancellationToken);
foreach (var deal in deals)
{
if (await _repository.ExistsAsync(deal.Id, cancellationToken))
@@ -46,8 +52,63 @@ public class DealService
_logger.LogInformation("Новая скидка: {Title} ({Id})", deal.Title, deal.Id);
await _repository.AddAsync(deal, cancellationToken);
// Общая рассылка в основной чат (как было раньше)
await _notifier.NotifyNewDealAsync(deal, cancellationToken);
if (subscriptions.Count == 0)
{
continue;
}
var searchableText = $"{deal.Title} {deal.StoreName}";
foreach (var subscription in subscriptions)
{
if (!subscription.IsActive)
{
continue;
}
var keywords = ParseKeywords(subscription.Keywords);
if (keywords.Count == 0)
{
continue;
}
if (IsMatch(searchableText, keywords))
{
await _notifier.NotifyDealToChatAsync(deal, subscription.ChatId, cancellationToken);
}
}
}
}
private static List<string> ParseKeywords(string raw)
{
return raw
.Split(new[] { '-',',', ';', '\n', '\r', '\t', ' ' }, StringSplitOptions.RemoveEmptyEntries)
.Select(k => k.Trim())
.Where(k => k.Length > 0)
.ToList();
}
private static bool IsMatch(string text, List<string> keywords)
{
if (keywords.Count == 0)
{
return false;
}
foreach (var keyword in keywords)
{
if (text.Contains(keyword, StringComparison.OrdinalIgnoreCase))
{
return true;
}
}
return false;
}
}
+19
View File
@@ -0,0 +1,19 @@
namespace PepperBot.Domain;
public class Subscription
{
public long Id { get; set; }
/// <summary>
/// Telegram chat id пользователя (может быть user id или id группы).
/// </summary>
public string ChatId { get; set; } = string.Empty;
/// <summary>
/// Строка с ключевыми словами, разделёнными запятыми/пробелами (например: "PLA, ABS, PETG").
/// </summary>
public string Keywords { get; set; } = string.Empty;
public bool IsActive { get; set; } = true;
}
@@ -0,0 +1,144 @@
using Microsoft.Data.Sqlite;
using PepperBot.Application.Interfaces;
using PepperBot.Domain;
namespace PepperBot.Infrastructure.Data;
public class SqliteSubscriptionRepository : ISubscriptionRepository
{
private const string DatabaseFileName = "pepper_bot.db";
private const string DataDirectory = "data";
private string ConnectionString =>
new SqliteConnectionStringBuilder
{
DataSource = Path.Combine(DataDirectory, DatabaseFileName)
}.ToString();
public async Task InitializeAsync(CancellationToken cancellationToken)
{
if (!Directory.Exists(DataDirectory))
{
Directory.CreateDirectory(DataDirectory);
}
await using var connection = new SqliteConnection(ConnectionString);
await connection.OpenAsync(cancellationToken);
var command = connection.CreateCommand();
command.CommandText =
"""
CREATE TABLE IF NOT EXISTS Subscriptions (
Id INTEGER PRIMARY KEY AUTOINCREMENT,
ChatId TEXT NOT NULL,
Keywords TEXT NOT NULL,
IsActive INTEGER NOT NULL DEFAULT 1
);
""";
await command.ExecuteNonQueryAsync(cancellationToken);
}
public async Task<IReadOnlyList<Subscription>> GetAllActiveAsync(CancellationToken cancellationToken)
{
await using var connection = new SqliteConnection(ConnectionString);
await connection.OpenAsync(cancellationToken);
var command = connection.CreateCommand();
command.CommandText =
"""
SELECT Id, ChatId, Keywords, IsActive
FROM Subscriptions
WHERE IsActive = 1;
""";
var result = new List<Subscription>();
await using var reader = await command.ExecuteReaderAsync(cancellationToken);
while (await reader.ReadAsync(cancellationToken))
{
var subscription = new Subscription
{
Id = reader.GetInt64(0),
ChatId = reader.GetString(1),
Keywords = reader.GetString(2),
IsActive = reader.GetInt32(3) == 1
};
result.Add(subscription);
}
return result;
}
public async Task AddSubscriptionAsync(string chatId, string keyword, CancellationToken cancellationToken)
{
await using var connection = new SqliteConnection(ConnectionString);
await connection.OpenAsync(cancellationToken);
var command = connection.CreateCommand();
command.CommandText =
"""
INSERT INTO Subscriptions (ChatId, Keywords, IsActive)
VALUES ($chatId, $keywords, 1);
""";
command.Parameters.AddWithValue("$chatId", chatId);
command.Parameters.AddWithValue("$keywords", keyword);
await command.ExecuteNonQueryAsync(cancellationToken);
}
public async Task<IReadOnlyList<Subscription>> GetByChatAsync(string chatId, CancellationToken cancellationToken)
{
await using var connection = new SqliteConnection(ConnectionString);
await connection.OpenAsync(cancellationToken);
var command = connection.CreateCommand();
command.CommandText =
"""
SELECT Id, ChatId, Keywords, IsActive
FROM Subscriptions
WHERE ChatId = $chatId AND IsActive = 1;
""";
command.Parameters.AddWithValue("$chatId", chatId);
var result = new List<Subscription>();
await using var reader = await command.ExecuteReaderAsync(cancellationToken);
while (await reader.ReadAsync(cancellationToken))
{
var subscription = new Subscription
{
Id = reader.GetInt64(0),
ChatId = reader.GetString(1),
Keywords = reader.GetString(2),
IsActive = reader.GetInt32(3) == 1
};
result.Add(subscription);
}
return result;
}
public async Task DeleteSubscriptionAsync(long id, CancellationToken cancellationToken)
{
await using var connection = new SqliteConnection(ConnectionString);
await connection.OpenAsync(cancellationToken);
var command = connection.CreateCommand();
command.CommandText =
"""
UPDATE Subscriptions
SET IsActive = 0
WHERE Id = $id;
""";
command.Parameters.AddWithValue("$id", id);
await command.ExecuteNonQueryAsync(cancellationToken);
}
}
@@ -1,43 +1,266 @@
using System.Net.Http.Json;
using System.Globalization;
using System.Text;
using Microsoft.Extensions.Logging;
using PepperBot.Application.Interfaces;
using PepperBot.Domain;
using Telegram.Bot;
using Telegram.Bot.Exceptions;
using Telegram.Bot.Polling;
using Telegram.Bot.Types;
using Telegram.Bot.Types.Enums;
using Telegram.Bot.Types.ReplyMarkups;
namespace PepperBot.Infrastructure.Telegram;
public class TelegramNotifier : ITelegramNotifier
{
private readonly HttpClient _httpClient;
private readonly ITelegramBotClient _botClient;
private readonly ISubscriptionRepository _subscriptionRepository;
private readonly ILogger<TelegramNotifier> _logger;
private readonly string _botToken;
private readonly string _chatId;
private readonly string _broadcastChatId;
private bool _started;
public TelegramNotifier(ILogger<TelegramNotifier> logger)
private enum ConversationState
{
None = 0,
AwaitingAddKeyword = 1
}
private readonly Dictionary<string, ConversationState> _chatStates = new();
public TelegramNotifier(
ITelegramBotClient botClient,
ISubscriptionRepository subscriptionRepository,
ILogger<TelegramNotifier> logger)
{
_botClient = botClient;
_subscriptionRepository = subscriptionRepository;
_logger = logger;
_botToken = Environment.GetEnvironmentVariable("TELEGRAM_BOT_TOKEN") ?? string.Empty;
_chatId = Environment.GetEnvironmentVariable("TELEGRAM_CHAT_ID") ?? string.Empty;
_broadcastChatId = Environment.GetEnvironmentVariable("TELEGRAM_CHAT_ID") ?? string.Empty;
if (string.IsNullOrWhiteSpace(_botToken) || string.IsNullOrWhiteSpace(_chatId))
if (string.IsNullOrWhiteSpace(_broadcastChatId))
{
_logger.LogWarning("TELEGRAM_BOT_TOKEN или TELEGRAM_CHAT_ID не заданы. Отправка сообщений отключена.");
_logger.LogInformation("TELEGRAM_CHAT_ID не задан. Общая рассылка по скидкам будет отключена.");
}
}
_httpClient = new HttpClient();
public Task NotifyNewDealAsync(Deal deal, CancellationToken cancellationToken)
{
if (string.IsNullOrWhiteSpace(_broadcastChatId))
{
return Task.CompletedTask;
}
public async Task NotifyNewDealAsync(Deal deal, CancellationToken cancellationToken)
return SendDealAsync(deal, _broadcastChatId, cancellationToken);
}
public Task NotifyDealToChatAsync(Deal deal, string chatId, CancellationToken cancellationToken)
{
if (string.IsNullOrWhiteSpace(_botToken) || string.IsNullOrWhiteSpace(_chatId))
if (string.IsNullOrWhiteSpace(chatId))
{
return Task.CompletedTask;
}
return SendDealAsync(deal, chatId, cancellationToken);
}
public Task StartAsync(CancellationToken cancellationToken)
{
if (_started)
{
return Task.CompletedTask;
}
var receiverOptions = new ReceiverOptions
{
AllowedUpdates = Array.Empty<UpdateType>()
};
_botClient.StartReceiving(
HandleUpdateAsync,
HandleErrorAsync,
receiverOptions,
cancellationToken);
_started = true;
_logger.LogInformation("Запущена обработка входящих сообщений Telegram.");
return Task.CompletedTask;
}
private async Task HandleUpdateAsync(ITelegramBotClient botClient, Update update, CancellationToken cancellationToken)
{
if (update.Type == UpdateType.CallbackQuery)
{
await HandleCallbackQueryAsync(botClient, update.CallbackQuery!, cancellationToken);
{
return;
}
}
if (update.Type != UpdateType.Message)
{
return;
}
var url = $"https://api.telegram.org/bot{_botToken}/sendMessage";
var message = update.Message;
if (message is null)
{
return;
}
if (message.Type != MessageType.Text)
{
return;
}
// Нас интересуют только личные сообщения боту
if (message.Chat.Type != ChatType.Private)
{
return;
}
var chatId = message.Chat.Id.ToString(CultureInfo.InvariantCulture);
var text = message.Text?.Trim();
if (string.IsNullOrWhiteSpace(text))
{
return;
}
// Команда для вызова меню
if (text.Equals("/menu", StringComparison.OrdinalIgnoreCase))
{
await SendMainMenuAsync(message.Chat.Id, cancellationToken);
_chatStates[chatId] = ConversationState.None;
return;
}
// Обработка выбора пунктов меню
if (text.Equals("Добавить подписку", StringComparison.OrdinalIgnoreCase))
{
await botClient.SendMessage(
chatId: message.Chat.Id,
text: "Отправьте одно или несколько ключевых слов для подписки (через пробел, запятую или с новой строки).",
cancellationToken: cancellationToken);
_chatStates[chatId] = ConversationState.AwaitingAddKeyword;
return;
}
if (text.Equals("Удалить подписку", StringComparison.OrdinalIgnoreCase))
{
await ShowDeleteMenuAsync(message.Chat.Id, chatId, cancellationToken);
_chatStates[chatId] = ConversationState.None;
return;
}
if (text.Equals("Список подписок", StringComparison.OrdinalIgnoreCase))
{
await ShowSubscriptionsListAsync(message.Chat.Id, chatId, cancellationToken);
_chatStates[chatId] = ConversationState.None;
return;
}
// Обработка состояний диалога
_chatStates.TryGetValue(chatId, out var state);
if (state == ConversationState.AwaitingAddKeyword)
{
var separators = new[] { ',', ';', '\n', '\r', '\t', ' ' };
var keywords = text
.Split(separators, StringSplitOptions.RemoveEmptyEntries)
.Select(k => k.Trim())
.Where(k => k.Length > 0)
.ToList();
if (keywords.Count == 0)
{
await botClient.SendMessage(
chatId: message.Chat.Id,
text: "Не нашёл ни одного ключевого слова. Попробуйте ещё раз или нажмите /menu.",
cancellationToken: cancellationToken);
return;
}
foreach (var keyword in keywords)
{
await _subscriptionRepository.AddSubscriptionAsync(chatId, keyword, cancellationToken);
}
_chatStates[chatId] = ConversationState.None;
var confirmation = new StringBuilder();
confirmation.AppendLine("Добавлены правила подписки по ключевым словам:");
confirmation.AppendLine(string.Join(", ", keywords));
await botClient.SendMessage(
chatId: message.Chat.Id,
text: confirmation.ToString(),
cancellationToken: cancellationToken);
_logger.LogInformation(
"Для чата {ChatId} добавлены правила подписки по ключевым словам: {Keywords}",
chatId,
string.Join(", ", keywords));
return;
}
// Если состояние не распознано — просто напоминаем про меню
await botClient.SendMessage(
chatId: message.Chat.Id,
text: "Используйте команду /menu для управления подписками.",
cancellationToken: cancellationToken);
}
private async Task HandleCallbackQueryAsync(ITelegramBotClient botClient, CallbackQuery callbackQuery, CancellationToken cancellationToken)
{
if (callbackQuery.Data is null)
{
return;
}
if (callbackQuery.Data.StartsWith("del:", StringComparison.Ordinal))
{
var idPart = callbackQuery.Data["del:".Length..];
if (!long.TryParse(idPart, CultureInfo.InvariantCulture, out var id))
{
return;
}
await _subscriptionRepository.DeleteSubscriptionAsync(id, cancellationToken);
await botClient.AnswerCallbackQuery(
callbackQueryId: callbackQuery.Id,
text: "Подписка удалена.",
cancellationToken: cancellationToken);
if (callbackQuery.Message is not null)
{
var chatId = callbackQuery.Message.Chat.Id;
var chatIdString = chatId.ToString(CultureInfo.InvariantCulture);
await ShowDeleteMenuAsync(chatId, chatIdString, cancellationToken);
}
}
}
private Task HandleErrorAsync(ITelegramBotClient botClient, Exception exception, CancellationToken cancellationToken)
{
var errorMessage = exception switch
{
ApiRequestException apiRequestException =>
$"Ошибка Telegram API: [{apiRequestException.ErrorCode}] {apiRequestException.Message}",
_ => exception.ToString()
};
_logger.LogError("Ошибка Telegram-бота: {Error}", errorMessage);
return Task.CompletedTask;
}
private async Task SendDealAsync(Deal deal, string chatId, CancellationToken cancellationToken)
{
var textBuilder = new StringBuilder();
textBuilder.AppendLine($"🔥 *{Escape(deal.Title)}*");
textBuilder.AppendLine($"🔥 {deal.Title}");
if (deal.CurrentPrice is not null)
{
@@ -55,33 +278,102 @@ public class TelegramNotifier : ITelegramNotifier
}
textBuilder.AppendLine();
textBuilder.AppendLine($"Магазин: {Escape(deal.StoreName)}");
textBuilder.AppendLine($"Магазин: {deal.StoreName}");
textBuilder.AppendLine();
textBuilder.AppendLine($"[Открыть на Pepper.ru]({Escape(deal.DealUrl)})");
textBuilder.AppendLine($"Открыть на Pepper.ru: {deal.DealUrl}");
var payload = new
try
{
chat_id = _chatId,
text = textBuilder.ToString(),
parse_mode = "Markdown"
await _botClient.SendMessage(
chatId: chatId,
text: textBuilder.ToString(),
cancellationToken: cancellationToken);
}
catch (ApiRequestException ex)
{
_logger.LogWarning(
ex,
"Не удалось отправить сообщение в Telegram (чат {ChatId}): [{Code}] {Message}",
chatId,
ex.ErrorCode,
ex.Message);
}
}
private Task SendMainMenuAsync(ChatId chatId, CancellationToken cancellationToken)
{
var keyboard = new ReplyKeyboardMarkup(new[]
{
new[] { new KeyboardButton("Добавить подписку") },
new[] { new KeyboardButton("Удалить подписку") },
new[] { new KeyboardButton("Список подписок") }
})
{
ResizeKeyboard = true,
OneTimeKeyboard = false
};
var response = await _httpClient.PostAsJsonAsync(url, payload, cancellationToken);
return _botClient.SendMessage(
chatId: chatId,
text: "Меню управления подписками:",
replyMarkup: keyboard,
cancellationToken: cancellationToken);
}
if (!response.IsSuccessStatusCode)
private async Task ShowSubscriptionsListAsync(ChatId chatId, string chatIdString, CancellationToken cancellationToken)
{
var body = await response.Content.ReadAsStringAsync(cancellationToken);
_logger.LogWarning("Не удалось отправить сообщение в Telegram: {Status} {Body}", response.StatusCode, body);
}
}
var subs = await _subscriptionRepository.GetByChatAsync(chatIdString, cancellationToken);
private static string Escape(string value)
if (subs.Count == 0)
{
return value
.Replace("_", "\\_")
.Replace("*", "\\*")
.Replace("[", "\\[")
.Replace("`", "\\`");
await _botClient.SendMessage(
chatId: chatId,
text: "У вас пока нет подписок.",
cancellationToken: cancellationToken);
return;
}
var sb = new StringBuilder();
sb.AppendLine("Ваши подписки:");
foreach (var sub in subs)
{
sb.AppendLine($"• [{sub.Id}] {sub.Keywords}");
}
await _botClient.SendMessage(
chatId: chatId,
text: sb.ToString(),
cancellationToken: cancellationToken);
}
private async Task ShowDeleteMenuAsync(ChatId chatId, string chatIdString, CancellationToken cancellationToken)
{
var subs = await _subscriptionRepository.GetByChatAsync(chatIdString, cancellationToken);
if (subs.Count == 0)
{
await _botClient.SendMessage(
chatId: chatId,
text: "Подписок для удаления не найдено.",
cancellationToken: cancellationToken);
return;
}
var buttons = subs
.Select(s => InlineKeyboardButton.WithCallbackData(
text: s.Keywords,
callbackData: $"del:{s.Id}"))
.Chunk(2)
.Select(chunk => chunk.ToArray())
.ToArray();
var keyboard = new InlineKeyboardMarkup(buttons);
await _botClient.SendMessage(
chatId: chatId,
text: "Выберите подписку для удаления:",
replyMarkup: keyboard,
cancellationToken: cancellationToken);
}
}
+1
View File
@@ -15,6 +15,7 @@
<PackageReference Include="Microsoft.Data.Sqlite" Version="10.0.3" />
<PackageReference Include="Microsoft.Extensions.Hosting" Version="10.0.3" />
<PackageReference Include="Microsoft.Extensions.Logging.Console" Version="10.0.3" />
<PackageReference Include="Telegram.Bot" Version="22.9.5.3" />
</ItemGroup>
</Project>
+11 -2
View File
@@ -9,23 +9,33 @@ public class BotWorker : BackgroundService
{
private readonly DealService _dealService;
private readonly IDealRepository _repository;
private readonly ISubscriptionRepository _subscriptionRepository;
private readonly ITelegramNotifier _telegramNotifier;
private readonly ILogger<BotWorker> _logger;
public BotWorker(
DealService dealService,
IDealRepository repository,
ISubscriptionRepository subscriptionRepository,
ITelegramNotifier telegramNotifier,
ILogger<BotWorker> logger)
{
_dealService = dealService;
_repository = repository;
_subscriptionRepository = subscriptionRepository;
_telegramNotifier = telegramNotifier;
_logger = logger;
}
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
await _repository.InitializeAsync(stoppingToken);
await _subscriptionRepository.InitializeAsync(stoppingToken);
_logger.LogInformation("PepperBot запущен. Опрос каждые 30 секунд.");
// Вся работа с Telegram вынесена в TelegramNotifier
await _telegramNotifier.StartAsync(stoppingToken);
_logger.LogInformation("PepperBot запущен. Опрос скидок каждые 30 секунд.");
while (!stoppingToken.IsCancellationRequested)
{
@@ -42,4 +52,3 @@ public class BotWorker : BackgroundService
}
}
}
+9
View File
@@ -9,6 +9,7 @@ using PepperBot.Infrastructure.Data;
using PepperBot.Infrastructure.Health;
using PepperBot.Infrastructure.Pepper;
using PepperBot.Infrastructure.Telegram;
using Telegram.Bot;
var builder = WebApplication.CreateBuilder(args);
@@ -19,11 +20,19 @@ builder.Services.AddLogging(logging =>
});
builder.Services.AddSingleton<IDealRepository, SqliteDealRepository>();
builder.Services.AddSingleton<ISubscriptionRepository, SqliteSubscriptionRepository>();
builder.Services.AddSingleton<IPepperClient, PepperClient>();
builder.Services.AddSingleton<ITelegramNotifier, TelegramNotifier>();
builder.Services.AddSingleton<IHealthMonitor, PepperHealthMonitor>();
builder.Services.AddSingleton<DealService>();
// Клиент Telegram-бота для приёма личных сообщений
builder.Services.AddSingleton<ITelegramBotClient>(_ =>
{
var token = Environment.GetEnvironmentVariable("TELEGRAM_BOT_TOKEN") ?? string.Empty;
return new TelegramBotClient(token);
});
builder.Services.AddHostedService<BotWorker>();
var app = builder.Build();