using System; using System.Collections.Concurrent; using System.Collections.Generic; using System.Threading; using System.Threading.Tasks; using CodeHollow.FeedReader; using Kruzya.TelegramBot.Core.Extensions; using Kruzya.TelegramBot.RichSiteSummary.Data; using Microsoft.Extensions.Configuration; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; using Telegram.Bot.Types; using Telegram.Bot.Types.Enums; using Feed = Kruzya.TelegramBot.RichSiteSummary.Data.Feed; namespace Kruzya.TelegramBot.RichSiteSummary.Service { /// /// Performs a RSS feed parsing. /// Grabs the all feeds from database. /// /// TODO: move entity grabbing to another service. /// public class RssFetch : IHostedService, IDisposable { private Timer _timer; private readonly ILogger _logger; private readonly IServiceScopeFactory _scopeFactory; private readonly ConcurrentQueue _queue; private TimeSpan timerPeriod => TimeSpan.FromSeconds(45); public RssFetch(ILogger logger, ConcurrentQueue queue, IServiceScopeFactory scopeFactory) { _scopeFactory = scopeFactory; _logger = logger; _queue = queue; } #region IHostedService /// /// Initializes the RSS fetcher timer. /// /// /// public Task StartAsync(CancellationToken cancellationToken) { _logger.LogInformation("RSS fetcher service is starting."); _timer = new Timer(DoFetch, null, TimeSpan.Zero, timerPeriod); return Task.CompletedTask; } /// /// Stops the RSS fetcher timer. /// /// /// public Task StopAsync(CancellationToken cancellationToken) { _logger.LogInformation("RSS fetcher service is stopping."); _timer?.Change(Timeout.Infinite, 0); return Task.CompletedTask; } #endregion #region IDisposable /// /// Disposes the timer. /// public void Dispose() { _timer?.Dispose(); } #endregion /// /// Performs the job of fetching RSS data. /// /// private async void DoFetch(object state) { _timer.Change(Timeout.Infinite, 0); _logger.LogDebug("RSS fetcher service is triggered."); try { using var scope = _scopeFactory.CreateScope(); var dbContext = scope.ServiceProvider.GetRequiredService(); var period = scope.ServiceProvider.GetRequiredService() .GetValue("rssFetchPeriod"); var feeds = await dbContext.Feeds.ForFetching(period); _logger.LogDebug("Received {count} feeds for fetching", new {count = feeds.Length}); foreach (var feed in feeds) { await ProcessFeed(feed, dbContext); feed.UpdatedAt = DateTime.Now; dbContext.MarkAsModified(feed); } await dbContext.SaveChangesAsync(); } finally { _timer.Change(timerPeriod, timerPeriod); } } #region Feeds private async Task ProcessFeed(Feed feed, RichSiteSummaryContext dbContext) { var parsedFeed = await FetchFeed(feed); if (parsedFeed == null) { return; } Post feedPost; var newPosts = new List(); foreach (var post in parsedFeed.Items) { feedPost = await dbContext.Posts.ByFeedAndUrl(post.Link, feed); if (feedPost != null) { // skip. This post already exists. continue; } feedPost = dbContext.Posts.Create(); feedPost.Feed = feed; feedPost.Title = post.Title; feedPost.Url = post.Link; feedPost.PostedAt = post.PublishingDate.GetValueOrDefault(DateTime.Now); newPosts.Add(feedPost); } if (newPosts.Count > 0) { var subscribers = await dbContext.Subscriptions.ByFeed(feed); foreach (var post in newPosts) { var text = post.MessageText; foreach (var subscriber in subscribers) { var message = new UserMessage() { ChatId = new ChatId(subscriber.SubscriberId), DisableWebPagePreview = true, ParseMode = ParseMode.Html, Text = text }; _queue.Enqueue(message); _logger.LogDebug($"Enqueued message for {subscriber.SubscriberId}"); } } } } private async Task FetchFeed(Feed feed) { try { return await FeedReader.ReadAsync(feed.Url.ToString()); } catch (Exception e) { _logger.LogError($"Feed {feed} can't be fetched: {e.Message}"); return null; } } #endregion } }