Files
Finlytic/FinlyticNews/Services/NewsDbService.cs
T

333 lines
12 KiB
C#

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;
/// <summary>
/// Defines database persistence operations for news articles and sources following the FinlyticNews lifecycle pipeline.
/// </summary>
public interface INewsDbService
{
Task<List<ArticleSourceEntity>> GetSourcesAsync();
Task<bool> IsUrlDuplicateAsync(string url);
/// <summary>
/// Phase 1: Discovers and locks a new article URL by setting its status to "Pending".
/// </summary>
Task<NewsArticleEntity> CreatePendingArticleAsync(
string url,
List<string>? discoveredIsins = null,
string? title = null,
string? summary = null,
DateTime? publishedAt = null,
string? language = null);
/// <summary>
/// Transitions the lifecycle state of an article (e.g. Pending -> Processing -> Scraping / Failed / Completed / Analyzed).
/// </summary>
Task UpdateArticleStatusAsync(Guid id, string status);
/// <summary>
/// Updates the target URL of an article if a redirect is resolved during the Processing phase.
/// </summary>
Task UpdateArticleUrlAsync(Guid id, string resolvedUrl);
/// <summary>
/// Phase 4: Saves AI classification from n8n and sets the lifecycle status to "Completed".
/// </summary>
Task<NewsArticleEntity?> SaveArticleClassificationAsync(Guid id, N8nResponsePayload payload, List<MatchedAssetEntity> matchedAssets);
Task<List<NewsArticleEntity>> GetArticlesByStatusAsync(string status);
Task<List<NewsArticleEntity>> GetFilteredNewsAsync(
int limit,
int offset,
string? isin = null,
DateTime? date = null,
string? status = null,
string? searchQuery = null);
Task<NewsArticleEntity?> GetArticleByIdAsync(Guid id);
}
/// <inheritdoc />
public class NewsDbService : INewsDbService
{
private readonly NewsDbContext _context;
private readonly IFinlyticLogger<NewsDbService> _finlyticLogger;
public NewsDbService(NewsDbContext context, IFinlyticLogger<NewsDbService> finlyticLogger)
{
_context = context;
_finlyticLogger = finlyticLogger;
}
/// <inheritdoc />
public async Task<List<ArticleSourceEntity>> GetSourcesAsync()
{
return await _context.ArticleSources.AsNoTracking().ToListAsync();
}
/// <inheritdoc />
public async Task<NewsArticleEntity?> GetArticleByIdAsync(Guid id)
{
return await _context.NewsArticles
.Include(a => a.MatchedAssets)
.AsNoTracking()
.FirstOrDefaultAsync(a => a.Id == id);
}
/// <inheritdoc />
public async Task<bool> IsUrlDuplicateAsync(string url)
{
var trimmedUrl = url.Trim();
return await _context.NewsArticles
.AsNoTracking()
.AnyAsync(a => a.SourceUrl == trimmedUrl);
}
/// <inheritdoc />
public async Task<NewsArticleEntity> CreatePendingArticleAsync(
string url,
List<string>? 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;
}
/// <inheritdoc />
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);
}
}
/// <inheritdoc />
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);
}
}
/// <inheritdoc />
public async Task<NewsArticleEntity?> SaveArticleClassificationAsync(Guid id, N8nResponsePayload payload, List<MatchedAssetEntity> matchedAssets)
{
DateTime? finalPublishedAt = null;
if (DateTime.TryParse(payload.PublishedAt, out var publishedDate))
{
finalPublishedAt = publishedDate.Kind == DateTimeKind.Unspecified
? DateTime.SpecifyKind(publishedDate, DateTimeKind.Utc)
: publishedDate.ToUniversalTime();
}
DateTime? finalScrapedAt = null;
if (DateTime.TryParse(payload.ScrapedAt, out var scrapedDate))
{
finalScrapedAt = scrapedDate.ToUniversalTime();
}
var rowsAffected = await _context.NewsArticles
.Where(a => a.Id == id)
.ExecuteUpdateAsync(s => s
.SetProperty(a => a.Title, payload.Title)
.SetProperty(a => a.Author, payload.Author)
.SetProperty(a => a.Summary, payload.Summary)
.SetProperty(a => a.ContentRaw, payload.ContentRaw)
.SetProperty(a => a.Language, payload.Language)
.SetProperty(a => a.Status, "Completed")
.SetProperty(a => a.PublishedAt, a => finalPublishedAt ?? a.PublishedAt)
.SetProperty(a => a.ScrapedAt, a => finalScrapedAt ?? a.ScrapedAt));
if (rowsAffected == 0)
{
await _finlyticLogger.LogWarningAsync(SettingKeys.NewsChannel, "[NewsChannel] Article with ID {Id} not found for classification update.", id);
return null;
}
// Cleanly delete existing matched assets and insert new ones
await _context.MatchedAssets.Where(m => m.NewsArticleId == id).ExecuteDeleteAsync();
if (matchedAssets != null && matchedAssets.Count > 0)
{
foreach (var asset in matchedAssets)
{
if (asset.Id == Guid.Empty)
{
asset.Id = Guid.NewGuid();
}
asset.NewsArticleId = id;
}
await _context.MatchedAssets.AddRangeAsync(matchedAssets);
await _context.SaveChangesAsync();
}
_context.ChangeTracker.Clear();
var completedArticle = await _context.NewsArticles
.Include(a => a.MatchedAssets)
.AsNoTracking()
.FirstOrDefaultAsync(a => a.Id == id);
await _finlyticLogger.LogInfoAsync(SettingKeys.NewsChannel, "[Lifecycle] Article {Id} successfully classified and marked 'Completed'. Title: '{Title}'", id, payload.Title);
return completedArticle;
}
/// <inheritdoc />
public async Task<List<NewsArticleEntity>> GetArticlesByStatusAsync(string status)
{
return await _context.NewsArticles
.Include(a => a.MatchedAssets)
.Where(a => a.Status == status)
.OrderByDescending(a => a.PublishedAt)
.AsNoTracking()
.ToListAsync();
}
/// <inheritdoc />
public async Task<List<NewsArticleEntity>> 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();
}
}