using System; using System.Collections.Generic; using System.Linq; using System.Threading.Tasks; using FinlyticCore.Dtos.News; using FinlyticCore.Services; using FinlyticNews.Database; using FinlyticNews.Entities; using FinlyticNews.Util; using Microsoft.EntityFrameworkCore; namespace FinlyticNews.Services; /// /// Defines database persistence operations for news articles and sources following the FinlyticNews lifecycle pipeline. /// public interface INewsDbService { Task> GetSourcesAsync(); Task IsUrlDuplicateAsync(string url); /// /// Phase 1: Discovers and locks a new article URL by setting its status to "Pending". /// Task CreatePendingArticleAsync( string url, List? discoveredIsins = null, string? title = null, string? summary = null, DateTime? publishedAt = null, string? language = null); /// /// Transitions the lifecycle state of an article (e.g. Pending -> Processing -> Scraping / Failed / Completed / Analyzed). /// Task UpdateArticleStatusAsync(Guid id, string status); /// /// Updates the target URL of an article if a redirect is resolved during the Processing phase. /// Task UpdateArticleUrlAsync(Guid id, string resolvedUrl); /// /// Phase 4: Saves AI classification from n8n and sets the lifecycle status to "Completed". /// Task SaveArticleClassificationAsync(Guid id, N8nResponsePayload payload, List matchedAssets); Task> GetArticlesByStatusAsync(string status); Task> GetFilteredNewsAsync( int limit, int offset, string? isin = null, DateTime? date = null, string? status = null, string? searchQuery = null); Task GetArticleByIdAsync(Guid id); } /// public class NewsDbService : INewsDbService { private readonly NewsDbContext _context; private readonly IFinlyticLogger _finlyticLogger; public NewsDbService(NewsDbContext context, IFinlyticLogger finlyticLogger) { _context = context; _finlyticLogger = finlyticLogger; } /// public async Task> GetSourcesAsync() { return await _context.ArticleSources.AsNoTracking().ToListAsync(); } /// public async Task GetArticleByIdAsync(Guid id) { return await _context.NewsArticles .Include(a => a.MatchedAssets) .AsNoTracking() .FirstOrDefaultAsync(a => a.Id == id); } /// public async Task IsUrlDuplicateAsync(string url) { var trimmedUrl = url.Trim(); return await _context.NewsArticles .AsNoTracking() .AnyAsync(a => a.SourceUrl == trimmedUrl); } /// public async Task CreatePendingArticleAsync( string url, List? discoveredIsins = null, string? title = null, string? summary = null, DateTime? publishedAt = null, string? language = null) { var trimmedUrl = url.Trim(); var cleanUrl = trimmedUrl.TrimEnd('.', '/'); var existingArticle = await _context.NewsArticles .FirstOrDefaultAsync(a => a.SourceUrl == trimmedUrl || a.SourceUrl == cleanUrl || a.SourceUrl.StartsWith(cleanUrl)); if (existingArticle != null) { await _finlyticLogger.LogDebugAsync(SettingKeys.NewsChannel, "[Lifecycle] Article URL already exists (Duplicate hit): {Url}", trimmedUrl); return existingArticle; } var finalPublishedAt = publishedAt.HasValue ? (publishedAt.Value.Kind == DateTimeKind.Unspecified ? DateTime.SpecifyKind(publishedAt.Value, DateTimeKind.Utc) : publishedAt.Value.ToUniversalTime()) : DateTime.UtcNow; if (finalPublishedAt > DateTime.UtcNow) { finalPublishedAt = DateTime.UtcNow; } var article = new NewsArticleEntity { Id = Guid.NewGuid(), Title = !string.IsNullOrWhiteSpace(title) ? title : "Pending Discovery", Summary = !string.IsNullOrWhiteSpace(summary) ? summary : "Extraction in progress...", ContentRaw = "Extraction in progress...", SourceUrl = trimmedUrl, Language = language, ScrapedAt = DateTime.UtcNow, PublishedAt = finalPublishedAt, Status = "Pending" }; if (discoveredIsins != null && discoveredIsins.Count > 0) { foreach (var isin in discoveredIsins) { article.MatchedAssets.Add(new MatchedAssetEntity { Id = Guid.NewGuid(), NewsArticleId = article.Id, Isin = isin, Name = isin }); } } try { _context.NewsArticles.Add(article); await _context.SaveChangesAsync(); await _finlyticLogger.LogDebugAsync(SettingKeys.NewsChannel, "[Lifecycle] Registered new article with status 'Pending'. ID: {Id}, Url: {Url}", article.Id, trimmedUrl); } catch (DbUpdateException ex) { _context.ChangeTracker.Clear(); await _finlyticLogger.LogWarningAsync(SettingKeys.NewsChannel, "[NewsChannel] Unique constraint or concurrency hit during insert for URL: {Url}. Fetching existing fallback. ({Message})", trimmedUrl, ex.InnerException?.Message ?? ex.Message); var existing = await _context.NewsArticles.AsNoTracking().FirstOrDefaultAsync(a => a.SourceUrl == trimmedUrl || a.SourceUrl == cleanUrl || a.SourceUrl.StartsWith(cleanUrl)); if (existing != null) { return existing; } return article; } return article; } /// public async Task UpdateArticleStatusAsync(Guid id, string status) { var rowsAffected = await _context.NewsArticles .Where(a => a.Id == id) .ExecuteUpdateAsync(s => s.SetProperty(a => a.Status, status)); if (rowsAffected == 0) { await _finlyticLogger.LogWarningAsync(SettingKeys.NewsChannel, "[NewsChannel] Attempted status transition for non-existing article. ID: {Id}", id); } else { await _finlyticLogger.LogDebugAsync(SettingKeys.NewsChannel, "[Lifecycle] Transitioned article {Id} to status '{Status}'", id, status); } } /// public async Task UpdateArticleUrlAsync(Guid id, string resolvedUrl) { var rowsAffected = await _context.NewsArticles .Where(a => a.Id == id) .ExecuteUpdateAsync(s => s.SetProperty(a => a.SourceUrl, resolvedUrl)); if (rowsAffected == 0) { await _finlyticLogger.LogWarningAsync(SettingKeys.NewsChannel, "[NewsChannel] Attempted URL update for non-existing article. ID: {Id}", id); } else { await _finlyticLogger.LogDebugAsync(SettingKeys.NewsChannel, "[Lifecycle] Resolved redirect for article {Id} -> New URL: {Url}", id, resolvedUrl); } } /// public async Task SaveArticleClassificationAsync(Guid id, N8nResponsePayload payload, List matchedAssets) { var article = await _context.NewsArticles .FirstOrDefaultAsync(a => a.Id == id); if (article == null) { await _finlyticLogger.LogWarningAsync(SettingKeys.NewsChannel, "[NewsChannel] Article with ID {Id} not found for classification update.", id); return null; } article.Title = payload.Title; article.Author = payload.Author; article.Summary = payload.Summary; article.ContentRaw = payload.ContentRaw; article.Language = payload.Language; article.Status = "Completed"; if (DateTime.TryParse(payload.PublishedAt, out var publishedDate)) { article.PublishedAt = publishedDate.Kind == DateTimeKind.Unspecified ? DateTime.SpecifyKind(publishedDate, DateTimeKind.Utc) : publishedDate.ToUniversalTime(); } if (DateTime.TryParse(payload.ScrapedAt, out var scrapedDate)) { article.ScrapedAt = scrapedDate.ToUniversalTime(); } await _context.MatchedAssets.Where(m => m.NewsArticleId == id).ExecuteDeleteAsync(); article.MatchedAssets = new List(); foreach (var asset in matchedAssets) { if (asset.Id == Guid.Empty) { asset.Id = Guid.NewGuid(); } _context.MatchedAssets.Add(asset); article.MatchedAssets.Add(asset); } await _context.SaveChangesAsync(); await _finlyticLogger.LogInfoAsync(SettingKeys.NewsChannel, "[Lifecycle] Article {Id} successfully classified and marked 'Completed'. Title: '{Title}'", article.Id, article.Title); return article; } /// public async Task> GetArticlesByStatusAsync(string status) { return await _context.NewsArticles .Include(a => a.MatchedAssets) .Where(a => a.Status == status) .OrderByDescending(a => a.PublishedAt) .AsNoTracking() .ToListAsync(); } /// public async Task> GetFilteredNewsAsync( int limit, int offset, string? isin = null, DateTime? date = null, string? status = null, string? searchQuery = null) { var query = _context.NewsArticles .Include(a => a.MatchedAssets) .AsNoTracking() .AsQueryable(); if (!string.IsNullOrWhiteSpace(isin)) { var cleanIsin = isin.Trim(); query = query.Where(a => a.MatchedAssets.Any(m => m.Isin == cleanIsin)); } if (date.HasValue) { var startUtc = date.Value.Date.ToUniversalTime(); var endUtc = startUtc.AddDays(1); query = query.Where(a => a.PublishedAt >= startUtc && a.PublishedAt < endUtc); } if (!string.IsNullOrWhiteSpace(status)) { query = query.Where(a => a.Status == status); } if (!string.IsNullOrWhiteSpace(searchQuery)) { var cleanSearch = searchQuery.Trim().ToLower(); query = query.Where(a => a.Title.ToLower().Contains(cleanSearch) || a.Summary.ToLower().Contains(cleanSearch) || a.MatchedAssets.Any(m => m.Name.ToLower().Contains(cleanSearch))); } return await query .OrderByDescending(a => a.PublishedAt) .Skip(offset) .Take(limit) .ToListAsync(); } }