diff --git a/FinlyticNews/Database/NewsDbContext.cs b/FinlyticNews/Database/NewsDbContext.cs index 276a002..774a769 100644 --- a/FinlyticNews/Database/NewsDbContext.cs +++ b/FinlyticNews/Database/NewsDbContext.cs @@ -1,5 +1,8 @@ +using FinlyticCore.Database; +using FinlyticCore.Entities.Settings; using FinlyticNews.Entities; using Microsoft.EntityFrameworkCore; +using Microsoft.EntityFrameworkCore.Design; namespace FinlyticNews.Database; @@ -7,7 +10,7 @@ namespace FinlyticNews.Database; /// Entity Framework Core database context for the news microservice, /// managing article sources, processed news, and matched assets. /// -public class NewsDbContext : DbContext +public class NewsDbContext : DbContext, ISettingsDbContext { /// /// Initializes a new instance of the class. @@ -17,6 +20,11 @@ public class NewsDbContext : DbContext { } + /// + /// Gets or sets the database set for dynamic settings. + /// + public DbSet DynamicSettings => Set(); + /// /// Gets or sets the database set for configured article sources. /// @@ -45,6 +53,12 @@ public class NewsDbContext : DbContext { base.OnModelCreating(modelBuilder); + modelBuilder.Entity(entity => + { + entity.HasKey(e => e.Id); + entity.HasIndex(e => e.Key).IsUnique(); + }); + // Configure unique index on SourceUrl for deduplication check (idempotency) modelBuilder.Entity() .HasIndex(a => a.SourceUrl) @@ -68,3 +82,13 @@ public class NewsDbContext : DbContext .OnDelete(DeleteBehavior.Cascade); } } + +public class NewsDbContextFactory : IDesignTimeDbContextFactory +{ + public NewsDbContext CreateDbContext(string[] args) + { + var optionsBuilder = new DbContextOptionsBuilder(); + optionsBuilder.UseNpgsql("Host=localhost;Database=news;Username=postgres;Password=postgres"); + return new NewsDbContext(optionsBuilder.Options); + } +} diff --git a/FinlyticNews/Migrations/20260815183946_AddDynamicSettings.Designer.cs b/FinlyticNews/Migrations/20260815183946_AddDynamicSettings.Designer.cs new file mode 100644 index 0000000..2940f80 --- /dev/null +++ b/FinlyticNews/Migrations/20260815183946_AddDynamicSettings.Designer.cs @@ -0,0 +1,205 @@ +// +using System; +using FinlyticNews.Database; +using Microsoft.EntityFrameworkCore; +using Microsoft.EntityFrameworkCore.Infrastructure; +using Microsoft.EntityFrameworkCore.Migrations; +using Microsoft.EntityFrameworkCore.Storage.ValueConversion; +using Npgsql.EntityFrameworkCore.PostgreSQL.Metadata; + +#nullable disable + +namespace FinlyticNews.Migrations +{ + [DbContext(typeof(NewsDbContext))] + [Migration("20260815183946_AddDynamicSettings")] + partial class AddDynamicSettings + { + /// + protected override void BuildTargetModel(ModelBuilder modelBuilder) + { +#pragma warning disable 612, 618 + modelBuilder + .HasAnnotation("ProductVersion", "10.0.9") + .HasAnnotation("Relational:MaxIdentifierLength", 63); + + NpgsqlModelBuilderExtensions.UseIdentityByDefaultColumns(modelBuilder); + + modelBuilder.Entity("FinlyticCore.Entities.Settings.SettingEntity", b => + { + b.Property("Id") + .ValueGeneratedOnAdd() + .HasColumnType("uuid"); + + b.Property("Key") + .IsRequired() + .HasMaxLength(150) + .HasColumnType("character varying(150)"); + + b.Property("LastUpdatedUtc") + .HasColumnType("timestamp with time zone"); + + b.Property("ServiceIdentifier") + .IsRequired() + .HasMaxLength(100) + .HasColumnType("character varying(100)"); + + b.Property("ValueJson") + .IsRequired() + .HasColumnType("text"); + + b.HasKey("Id"); + + b.HasIndex("Key") + .IsUnique(); + + b.ToTable("DynamicSettings"); + }); + + modelBuilder.Entity("FinlyticNews.Entities.ArticleSourceEntity", b => + { + b.Property("Id") + .ValueGeneratedOnAdd() + .HasColumnType("uuid"); + + b.Property("Name") + .IsRequired() + .HasColumnType("text"); + + b.Property("Source") + .IsRequired() + .HasColumnType("text"); + + b.Property("Type") + .IsRequired() + .HasColumnType("text"); + + b.HasKey("Id"); + + b.ToTable("ArticleSources"); + }); + + modelBuilder.Entity("FinlyticNews.Entities.MatchedAssetEntity", b => + { + b.Property("Id") + .ValueGeneratedOnAdd() + .HasColumnType("uuid"); + + b.Property("Isin") + .IsRequired() + .HasColumnType("text"); + + b.Property("Name") + .IsRequired() + .HasColumnType("text"); + + b.Property("NewsArticleId") + .HasColumnType("uuid"); + + b.HasKey("Id"); + + b.HasIndex("Isin"); + + b.HasIndex("NewsArticleId"); + + b.ToTable("MatchedAssets"); + }); + + modelBuilder.Entity("FinlyticNews.Entities.NewsArticleEntity", b => + { + b.Property("Id") + .ValueGeneratedOnAdd() + .HasColumnType("uuid"); + + b.Property("Author") + .HasColumnType("text"); + + b.Property("ContentRaw") + .IsRequired() + .HasColumnType("text"); + + b.Property("Language") + .HasColumnType("text"); + + b.Property("PublishedAt") + .HasColumnType("timestamp with time zone"); + + b.Property("ScrapedAt") + .HasColumnType("timestamp with time zone"); + + b.Property("SourceUrl") + .IsRequired() + .HasColumnType("text"); + + b.Property("Status") + .IsRequired() + .HasColumnType("text"); + + b.Property("Summary") + .HasColumnType("text"); + + b.Property("Title") + .IsRequired() + .HasColumnType("text"); + + b.HasKey("Id"); + + b.HasIndex("SourceUrl") + .IsUnique(); + + b.HasIndex("Status"); + + b.HasIndex("PublishedAt", "ScrapedAt"); + + b.ToTable("NewsArticles"); + }); + + modelBuilder.Entity("FinlyticNews.Entities.NewsSettingsEntity", b => + { + b.Property("Id") + .ValueGeneratedOnAdd() + .HasColumnType("uuid"); + + b.Property("ArticleRetentionDays") + .HasColumnType("integer"); + + b.Property("DefaultPageSize") + .HasColumnType("integer"); + + b.Property("N8nWebhookUrl") + .IsRequired() + .HasColumnType("text"); + + b.Property("PollingFrequencyMinutes") + .HasColumnType("integer"); + + b.Property("ScrapingIntervalMinutes") + .HasColumnType("integer"); + + b.Property("UpdatedAt") + .HasColumnType("timestamp with time zone"); + + b.HasKey("Id"); + + b.ToTable("Settings"); + }); + + modelBuilder.Entity("FinlyticNews.Entities.MatchedAssetEntity", b => + { + b.HasOne("FinlyticNews.Entities.NewsArticleEntity", "NewsArticle") + .WithMany("MatchedAssets") + .HasForeignKey("NewsArticleId") + .OnDelete(DeleteBehavior.Cascade) + .IsRequired(); + + b.Navigation("NewsArticle"); + }); + + modelBuilder.Entity("FinlyticNews.Entities.NewsArticleEntity", b => + { + b.Navigation("MatchedAssets"); + }); +#pragma warning restore 612, 618 + } + } +} diff --git a/FinlyticNews/Migrations/20260815183946_AddDynamicSettings.cs b/FinlyticNews/Migrations/20260815183946_AddDynamicSettings.cs new file mode 100644 index 0000000..7f57c68 --- /dev/null +++ b/FinlyticNews/Migrations/20260815183946_AddDynamicSettings.cs @@ -0,0 +1,43 @@ +using System; +using Microsoft.EntityFrameworkCore.Migrations; + +#nullable disable + +namespace FinlyticNews.Migrations +{ + /// + public partial class AddDynamicSettings : Migration + { + /// + protected override void Up(MigrationBuilder migrationBuilder) + { + migrationBuilder.CreateTable( + name: "DynamicSettings", + columns: table => new + { + Id = table.Column(type: "uuid", nullable: false), + Key = table.Column(type: "character varying(150)", maxLength: 150, nullable: false), + ValueJson = table.Column(type: "text", nullable: false), + ServiceIdentifier = table.Column(type: "character varying(100)", maxLength: 100, nullable: false), + LastUpdatedUtc = table.Column(type: "timestamp with time zone", nullable: false) + }, + constraints: table => + { + table.PrimaryKey("PK_DynamicSettings", x => x.Id); + }); + + migrationBuilder.CreateIndex( + name: "IX_DynamicSettings_Key", + table: "DynamicSettings", + column: "Key", + unique: true); + } + + /// + protected override void Down(MigrationBuilder migrationBuilder) + { + migrationBuilder.DropTable( + name: "DynamicSettings"); + } + } +} diff --git a/FinlyticNews/Migrations/NewsDbContextModelSnapshot.cs b/FinlyticNews/Migrations/NewsDbContextModelSnapshot.cs index 3e02b71..12f2e3a 100644 --- a/FinlyticNews/Migrations/NewsDbContextModelSnapshot.cs +++ b/FinlyticNews/Migrations/NewsDbContextModelSnapshot.cs @@ -22,6 +22,37 @@ namespace FinlyticNews.Migrations NpgsqlModelBuilderExtensions.UseIdentityByDefaultColumns(modelBuilder); + modelBuilder.Entity("FinlyticCore.Entities.Settings.SettingEntity", b => + { + b.Property("Id") + .ValueGeneratedOnAdd() + .HasColumnType("uuid"); + + b.Property("Key") + .IsRequired() + .HasMaxLength(150) + .HasColumnType("character varying(150)"); + + b.Property("LastUpdatedUtc") + .HasColumnType("timestamp with time zone"); + + b.Property("ServiceIdentifier") + .IsRequired() + .HasMaxLength(100) + .HasColumnType("character varying(100)"); + + b.Property("ValueJson") + .IsRequired() + .HasColumnType("text"); + + b.HasKey("Id"); + + b.HasIndex("Key") + .IsUnique(); + + b.ToTable("DynamicSettings"); + }); + modelBuilder.Entity("FinlyticNews.Entities.ArticleSourceEntity", b => { b.Property("Id") diff --git a/FinlyticNews/Program.cs b/FinlyticNews/Program.cs index 5211184..fe171a1 100644 --- a/FinlyticNews/Program.cs +++ b/FinlyticNews/Program.cs @@ -1,3 +1,5 @@ +using FinlyticCore.Database; +using FinlyticCore.Services; using FinlyticNews.Database; using FinlyticNews.Services; using FinlyticNews.Util; @@ -10,10 +12,15 @@ var builder = Host.CreateApplicationBuilder(args); builder.Services.AddDbContext(options => options.UseNpgsql(builder.Configuration.GetConnectionString("DefaultConnection"))); +builder.Services.AddScoped(sp => sp.GetRequiredService()); // Register standard HttpClient builder.Services.AddHttpClient(); +// Register Core Services +builder.Services.AddSingleton(); +builder.Services.AddSingleton(typeof(IFinlyticLogger<>), typeof(FinlyticLogger<>)); + // Register Discovery Adapters builder.Services.AddSingleton(); builder.Services.AddSingleton(); diff --git a/FinlyticNews/Services/ArticleDiscoveryService.cs b/FinlyticNews/Services/ArticleDiscoveryService.cs index 49cb6d5..7d250c0 100644 --- a/FinlyticNews/Services/ArticleDiscoveryService.cs +++ b/FinlyticNews/Services/ArticleDiscoveryService.cs @@ -1,21 +1,18 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Net.Http; +using System.Threading; +using System.Threading.Tasks; using FinlyticCore.Dtos.News; +using FinlyticCore.Services; using FinlyticNews.Adapters.Discovery; -using Microsoft.Extensions.Logging; +using FinlyticNews.Util; namespace FinlyticNews.Services; -/// -/// Defines operations for discovering article links from various news feeds and web pages. -/// public interface IArticleDiscoveryService { - /// - /// Discovers article links from a given source URL using configured adapters. - /// - /// The URL of the source news page or feed. - /// The type identifier of the adapter to use (e.g. "rss", "html"). - /// The token to monitor for cancellation requests. - /// A list of discovered absolute article URLs (optionally carrying ISINs), or null if the URL was invalid. Task?> DiscoverLinksAsync(string url, string adapterType, CancellationToken ct = default); } @@ -23,22 +20,16 @@ public interface IArticleDiscoveryService public class ArticleDiscoveryService : IArticleDiscoveryService { private readonly IEnumerable _adapters; - private readonly ILogger _logger; + private readonly IFinlyticLogger _finlyticLogger; private readonly HttpClient _httpClient; - /// - /// Initializes a new instance of the class. - /// - /// Registered list of specialized adapters. - /// The application logging channel. - /// The HTTP client to fetch feeds. public ArticleDiscoveryService( IEnumerable adapters, - ILogger logger, + IFinlyticLogger finlyticLogger, HttpClient httpClient) { _adapters = adapters; - _logger = logger; + _finlyticLogger = finlyticLogger; _httpClient = httpClient; } @@ -47,37 +38,35 @@ public class ArticleDiscoveryService : IArticleDiscoveryService { if (!Uri.TryCreate(url, UriKind.Absolute, out var uri)) { - _logger.LogWarning("[{Channel}] Invalid source URL passed for discovery: {Url}", "NewsChannel", url); + await _finlyticLogger.LogWarningAsync(SettingKeys.NewsChannel, "[ArticleDiscoveryService] Invalid source URL passed for discovery: {Url}", url); return null; } try { - _logger.LogDebug("Fetching content from discovery source: {Url}", uri); + await _finlyticLogger.LogDebugAsync(SettingKeys.NewsChannel, "[ArticleDiscoveryService] Fetching content from discovery source: {Url}", uri); var content = await _httpClient.GetStringAsync(uri, ct); if (string.IsNullOrWhiteSpace(content)) return []; var trimmedContent = content.TrimStart(); - // 1. Zuerst gezielt nach registriertem Adapter suchen (z. B. finanznachrichten_rss) var adapter = _adapters.FirstOrDefault(a => a.Name.Equals(adapterType, StringComparison.OrdinalIgnoreCase) || uri.Host.Contains(a.Name, StringComparison.OrdinalIgnoreCase)); if (adapter != null) { - _logger.LogInformation("[{Channel}] Using specialized adapter {AdapterName} for source: {Url}", "NewsChannel", adapter.Name, url); + await _finlyticLogger.LogInfoAsync(SettingKeys.NewsChannel, "[ArticleDiscoveryService] Using specialized adapter {AdapterName} for source: {Url}", adapter.Name, url); return adapter.ExtractUrls(content, url); } - // 2. Fallback: Automatische Erkennung für generische RSS/Atom-Feeds if (adapterType.Equals("rss", StringComparison.OrdinalIgnoreCase) || trimmedContent.StartsWith(" a.Name.Equals("rss", StringComparison.OrdinalIgnoreCase)) ?? new RssDiscoveryAdapter(); @@ -85,12 +74,12 @@ public class ArticleDiscoveryService : IArticleDiscoveryService return rssAdapter.ExtractUrls(content, url); } - _logger.LogWarning("[{Channel}] No suitable discovery adapter found for type '{Type}' and URL: {Url}", "NewsChannel", adapterType, url); + await _finlyticLogger.LogWarningAsync(SettingKeys.NewsChannel, "[ArticleDiscoveryService] No suitable discovery adapter found for type '{Type}' and URL: {Url}", adapterType, url); return []; } catch (Exception ex) { - _logger.LogError(ex, "[{Channel}] Failed to run article discovery on URL: {Url}", "NewsChannel", url); + await _finlyticLogger.LogErrorAsync(SettingKeys.NewsChannel, ex, "[ArticleDiscoveryService] Failed to run article discovery on URL: {Url}", url); return []; } } diff --git a/FinlyticNews/Services/N8nService.cs b/FinlyticNews/Services/N8nService.cs index 27df6a6..011f40e 100644 --- a/FinlyticNews/Services/N8nService.cs +++ b/FinlyticNews/Services/N8nService.cs @@ -1,14 +1,18 @@ +using System; +using System.Collections.Generic; +using System.Net.Http; using System.Text.Json; -using System.Text.Json.Serialization; +using System.Threading; +using System.Threading.Tasks; +using FinlyticCore.Dtos.News; +using FinlyticCore.Services; using FinlyticCore.Util; +using FinlyticNews.Util; using Microsoft.Extensions.Configuration; using Microsoft.Extensions.DependencyInjection; -using Microsoft.Extensions.Logging; namespace FinlyticNews.Services; -using FinlyticCore.Dtos.News; - /// /// Defines integration operations with the external n8n AI workflow webhook. /// @@ -26,21 +30,18 @@ public class N8nService : IN8nService private readonly HttpClient _httpClient; private readonly IServiceScopeFactory _scopeFactory; private readonly IConfiguration _configuration; - private readonly ILogger _logger; + private readonly IFinlyticLogger _finlyticLogger; - /// - /// Initializes a new instance of the class. - /// public N8nService( HttpClient httpClient, IServiceScopeFactory scopeFactory, IConfiguration configuration, - ILogger logger) + IFinlyticLogger finlyticLogger) { _httpClient = httpClient; _scopeFactory = scopeFactory; _configuration = configuration; - _logger = logger; + _finlyticLogger = finlyticLogger; } /// @@ -48,7 +49,6 @@ public class N8nService : IN8nService { string? targetUrl = null; - // 1. Dynamic Settings Resolution (DB Scope -> AppSettings Fallback) using (var scope = _scopeFactory.CreateScope()) { var settingsDb = scope.ServiceProvider.GetService(); @@ -67,17 +67,16 @@ public class N8nService : IN8nService if (string.IsNullOrWhiteSpace(targetUrl)) { - _logger.LogError("[{Channel}] N8nWebhookUrl is not configured in DB or application settings.", "NewsChannel"); + await _finlyticLogger.LogErrorAsync(SettingKeys.NewsChannel, "[N8nService] N8nWebhookUrl is not configured in DB or application settings."); return null; } - _logger.LogInformation("[{Channel}] Posting article to n8n webhook pipeline at: {Url}", "NewsChannel", targetUrl); + await _finlyticLogger.LogInfoAsync(SettingKeys.NewsChannel, "[N8nService] Posting article to n8n webhook pipeline at: {Url}", targetUrl); var payload = new N8nRequestPayload(content, filteredAssets ?? []); try { - // Zero-Allocation / Source-Generated Request Serialization var jsonContent = JsonSerializer.Serialize(payload, FinlyticJsonSerializerContext.Default.N8nRequestPayload); using var requestContent = new StringContent(jsonContent, System.Text.Encoding.UTF8, "application/json"); @@ -86,7 +85,7 @@ public class N8nService : IN8nService if (!response.IsSuccessStatusCode) { var errorMsg = await response.Content.ReadAsStringAsync(ct); - _logger.LogError("[{Channel}] n8n webhook returned status code {StatusCode}. Error payload: {Error}", "NewsChannel", response.StatusCode, errorMsg); + await _finlyticLogger.LogErrorAsync(SettingKeys.NewsChannel, "[N8nService] n8n webhook returned status code {StatusCode}. Error payload: {Error}", response.StatusCode, errorMsg); return null; } @@ -95,18 +94,16 @@ public class N8nService : IN8nService var root = doc.RootElement; - // 2. Robust n8n Array-Unwrapping (Handles [{ "json": { ... } }]) if (root.ValueKind == JsonValueKind.Array) { if (root.GetArrayLength() == 0) { - _logger.LogWarning("[{Channel}] n8n webhook returned an empty array.", "NewsChannel"); + await _finlyticLogger.LogWarningAsync(SettingKeys.NewsChannel, "[N8nService] n8n webhook returned an empty array."); return null; } root = root[0]; } - // 3. Dynamic Node Wrapper Unwrapping ("json", "output", "data", "body") if (root.ValueKind == JsonValueKind.Object) { if (root.TryGetProperty("json", out var jsonChild) && jsonChild.ValueKind == JsonValueKind.Object) @@ -119,13 +116,12 @@ public class N8nService : IN8nService root = bodyChild; } - // 4. Source-Generated Deserialization directly from JsonElement var result = root.Deserialize(FinlyticJsonSerializerContext.Default.N8nResponsePayload); return result; } catch (Exception ex) { - _logger.LogError(ex, "[{Channel}] Failed to communicate with or parse response from n8n webhook workflow.", "NewsChannel"); + await _finlyticLogger.LogErrorAsync(SettingKeys.NewsChannel, ex, "[N8nService] Failed to communicate with or parse response from n8n webhook workflow."); return null; } } diff --git a/FinlyticNews/Services/NewsDbService.cs b/FinlyticNews/Services/NewsDbService.cs index f55fc65..8523aec 100644 --- a/FinlyticNews/Services/NewsDbService.cs +++ b/FinlyticNews/Services/NewsDbService.cs @@ -1,6 +1,12 @@ +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; @@ -39,22 +45,11 @@ public interface INewsDbService /// Task SaveArticleClassificationAsync(Guid id, N8nResponsePayload payload, List matchedAssets); - /// - /// Retrieves articles ready for historical sync or sentiment processing. - /// - Task> GetCompletedArticlesAsync(int limit, int offset, string? isin = null); + Task> GetArticlesByStatusAsync(string status); - /// - /// Fetches articles matching specific lifecycle statuses (e.g. "Pending", "Scraping" for Phase 1 retry). - /// - Task> GetArticlesByStatusAsync(params string[] statuses); - - /// - /// Fetches public daily news for API endpoints, filtering out intermediate or failed lifecycle states by default. - /// Task> GetFilteredNewsAsync( - int limit = 20, - int offset = 0, + int limit, + int offset, string? isin = null, DateTime? date = null, string? status = null, @@ -67,15 +62,12 @@ public interface INewsDbService public class NewsDbService : INewsDbService { private readonly NewsDbContext _context; - private readonly ILogger _logger; + private readonly IFinlyticLogger _finlyticLogger; - /// - /// Initializes a new instance of the class. - /// - public NewsDbService(NewsDbContext context, ILogger logger) + public NewsDbService(NewsDbContext context, IFinlyticLogger finlyticLogger) { _context = context; - _logger = logger; + _finlyticLogger = finlyticLogger; } /// @@ -103,9 +95,6 @@ public class NewsDbService : INewsDbService } /// - /// - /// Lifecycle Step 1: Creates a new article in 'Pending' state as an immediate lock. - /// public async Task CreatePendingArticleAsync( string url, List? discoveredIsins = null, @@ -121,7 +110,7 @@ public class NewsDbService : INewsDbService if (existingArticle != null) { - _logger.LogDebug("[Lifecycle] Article URL already exists (Duplicate hit): {Url}", trimmedUrl); + await _finlyticLogger.LogDebugAsync(SettingKeys.NewsChannel, "[Lifecycle] Article URL already exists (Duplicate hit): {Url}", trimmedUrl); return existingArticle; } @@ -146,7 +135,7 @@ public class NewsDbService : INewsDbService Language = language, ScrapedAt = DateTime.UtcNow, PublishedAt = finalPublishedAt, - Status = "Pending" // 1. Pending State + Status = "Pending" }; if (discoveredIsins != null && discoveredIsins.Count > 0) @@ -167,11 +156,11 @@ public class NewsDbService : INewsDbService { _context.NewsArticles.Add(article); await _context.SaveChangesAsync(); - _logger.LogDebug("[Lifecycle] Registered new article with status 'Pending'. ID: {Id}, Url: {Url}", article.Id, trimmedUrl); + await _finlyticLogger.LogDebugAsync(SettingKeys.NewsChannel, "[Lifecycle] Registered new article with status 'Pending'. ID: {Id}, Url: {Url}", article.Id, trimmedUrl); } catch (DbUpdateConcurrencyException) { - _logger.LogWarning("[{Channel}] Concurrency hit during insert for URL: {Url}. Fetching existing fallback.", "NewsChannel", trimmedUrl); + await _finlyticLogger.LogWarningAsync(SettingKeys.NewsChannel, "[NewsChannel] Concurrency hit during insert for URL: {Url}. Fetching existing fallback.", trimmedUrl); return await _context.NewsArticles.FirstAsync(a => a.SourceUrl == trimmedUrl); } @@ -179,9 +168,6 @@ public class NewsDbService : INewsDbService } /// - /// - /// Lifecycle Step 2 & 5: Updates state (e.g. Pending -> Processing -> Scraping / Failed / Analyzed). - /// public async Task UpdateArticleStatusAsync(Guid id, string status) { var rowsAffected = await _context.NewsArticles @@ -190,11 +176,11 @@ public class NewsDbService : INewsDbService if (rowsAffected == 0) { - _logger.LogWarning("[{Channel}] Attempted status transition for non-existing article. ID: {Id}", "NewsChannel", id); + await _finlyticLogger.LogWarningAsync(SettingKeys.NewsChannel, "[NewsChannel] Attempted status transition for non-existing article. ID: {Id}", id); } else { - _logger.LogDebug("[Lifecycle] Transitioned article {Id} to status '{Status}'", id, status); + await _finlyticLogger.LogDebugAsync(SettingKeys.NewsChannel, "[Lifecycle] Transitioned article {Id} to status '{Status}'", id, status); } } @@ -207,18 +193,15 @@ public class NewsDbService : INewsDbService if (rowsAffected == 0) { - _logger.LogWarning("[{Channel}] Attempted URL update for non-existing article. ID: {Id}", "NewsChannel", id); + await _finlyticLogger.LogWarningAsync(SettingKeys.NewsChannel, "[NewsChannel] Attempted URL update for non-existing article. ID: {Id}", id); } else { - _logger.LogDebug("[Lifecycle] Resolved redirect for article {Id} -> New URL: {Url}", id, resolvedUrl); + await _finlyticLogger.LogDebugAsync(SettingKeys.NewsChannel, "[Lifecycle] Resolved redirect for article {Id} -> New URL: {Url}", id, resolvedUrl); } } /// - /// - /// Lifecycle Step 4: Persists n8n classification and transitions status to 'Completed'. - /// public async Task SaveArticleClassificationAsync(Guid id, N8nResponsePayload payload, List matchedAssets) { var article = await _context.NewsArticles @@ -226,7 +209,7 @@ public class NewsDbService : INewsDbService if (article == null) { - _logger.LogWarning("[{Channel}] Article with ID {Id} not found for classification update.", "NewsChannel", id); + await _finlyticLogger.LogWarningAsync(SettingKeys.NewsChannel, "[NewsChannel] Article with ID {Id} not found for classification update.", id); return null; } @@ -235,8 +218,6 @@ public class NewsDbService : INewsDbService article.Summary = payload.Summary; article.ContentRaw = payload.ContentRaw; article.Language = payload.Language; - - // 🎯 Step 4: Classification finished -> Transition to 'Completed' (triggers MQTT broadcast) article.Status = "Completed"; if (DateTime.TryParse(payload.PublishedAt, out var publishedDate)) @@ -251,7 +232,6 @@ public class NewsDbService : INewsDbService article.ScrapedAt = scrapedDate.ToUniversalTime(); } - // Clean up previous temporary assets await _context.MatchedAssets.Where(m => m.NewsArticleId == id).ExecuteDeleteAsync(); article.MatchedAssets = new List(); @@ -267,113 +247,65 @@ public class NewsDbService : INewsDbService } await _context.SaveChangesAsync(); - _logger.LogInformation("[Lifecycle] Article {Id} successfully classified and marked 'Completed'. Title: '{Title}'", article.Id, article.Title); + await _finlyticLogger.LogInfoAsync(SettingKeys.NewsChannel, "[Lifecycle] Article {Id} successfully classified and marked 'Completed'. Title: '{Title}'", article.Id, article.Title); return article; } /// - public async Task> GetCompletedArticlesAsync(int limit, int offset, string? isin = null) + public async Task> GetArticlesByStatusAsync(string status) { - var query = _context.NewsArticles + return await _context.NewsArticles .Include(a => a.MatchedAssets) - .Where(a => a.Status == "Completed" || a.Status == "Analyzed"); - - if (!string.IsNullOrWhiteSpace(isin)) - { - var cleanIsin = isin.Trim(); - query = query.Where(a => a.MatchedAssets.Any(m => m.Isin == cleanIsin || m.Name == cleanIsin)); - } - - return await query + .Where(a => a.Status == status) .OrderByDescending(a => a.PublishedAt) - .ThenByDescending(a => a.ScrapedAt) - .ThenByDescending(a => a.Id) - .Skip(offset) - .Take(limit) .AsNoTracking() .ToListAsync(); } /// - /// - /// Lifecycle Helper: Retrieves articles by target lifecycle status (e.g. "Pending", "Scraping" for Phase 1 Retry). - /// - public async Task> GetArticlesByStatusAsync(params string[] statuses) - { - IQueryable query = _context.NewsArticles - .Include(a => a.MatchedAssets) - .AsNoTracking(); - - if (statuses != null && statuses.Length > 0) - { - var cleanStatuses = statuses.Select(s => s.Trim()).ToList(); - query = query.Where(a => cleanStatuses.Contains(a.Status)); - } - else - { - query = query.Where(a => a.Status != "Failed" && a.Status != "Duplicate"); - } - - return await query.ToListAsync(); - } - - /// - /// - /// API Gateway Helper: Retrieves articles for UI rendering, excluding intermediate/failed states by default. - /// public async Task> GetFilteredNewsAsync( - int limit = 20, - int offset = 0, + int limit, + int offset, string? isin = null, DateTime? date = null, string? status = null, string? searchQuery = null) { - IQueryable query = _context.NewsArticles + var query = _context.NewsArticles .Include(a => a.MatchedAssets) - .AsNoTracking(); + .AsNoTracking() + .AsQueryable(); - // 1. Status Filter - if (!string.IsNullOrWhiteSpace(status)) - { - var targetStatus = status.Trim(); - query = query.Where(a => a.Status == targetStatus); - } - else - { - // By default, only show fully processed articles to the API/UI - query = query.Where(a => a.Status == "Completed" || a.Status == "Analyzed"); - } - - // 2. Date Filter - if (date.HasValue) - { - // Erstelle ein exaktes UTC-Datum von 00:00:00 Uhr am gebuchten Tag - var targetDate = new DateTime(date.Value.Year, date.Value.Month, date.Value.Day, 0, 0, 0, DateTimeKind.Utc); - var nextDate = targetDate.AddDays(1); - - query = query.Where(a => a.PublishedAt >= targetDate && a.PublishedAt < nextDate); - } - - // 3. ISIN / Symbol Filter if (!string.IsNullOrWhiteSpace(isin)) { var cleanIsin = isin.Trim(); - query = query.Where(a => a.MatchedAssets.Any(m => m.Isin == cleanIsin || m.Name == cleanIsin)); + 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); } - // 4. Search Term if (!string.IsNullOrWhiteSpace(searchQuery)) { - var q = searchQuery.Trim(); - query = query.Where(a => EF.Functions.ILike(a.Title, $"%{q}%") || (a.Summary != null && EF.Functions.ILike(a.Summary, $"%{q}%"))); + 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) - .ThenByDescending(a => a.ScrapedAt) - .ThenByDescending(a => a.Id) .Skip(offset) .Take(limit) .ToListAsync(); diff --git a/FinlyticNews/Services/NewsScraperBackgroundService.cs b/FinlyticNews/Services/NewsScraperBackgroundService.cs index cd07714..52503ec 100644 --- a/FinlyticNews/Services/NewsScraperBackgroundService.cs +++ b/FinlyticNews/Services/NewsScraperBackgroundService.cs @@ -1,18 +1,22 @@ +using System; using System.Collections.Concurrent; +using System.Collections.Generic; using System.IO; +using System.Linq; using System.Text.Json; -using System.Text.Json.Serialization; using System.Text.RegularExpressions; +using System.Threading; +using System.Threading.Tasks; using FinlyticAssets.Models; using FinlyticAssets.Util; using FinlyticCore.Dtos.News; +using FinlyticCore.Services; using FinlyticCore.Util; using FinlyticNews.Entities; using FinlyticNews.Util; using Microsoft.Extensions.Configuration; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; -using Microsoft.Extensions.Logging; namespace FinlyticNews.Services; @@ -22,9 +26,6 @@ namespace FinlyticNews.Services; /// public class NewsScraperBackgroundService : BackgroundService { - /// - /// Internal wrapper to associate compiled regex patterns with the unmodified AssetIndex record. - /// private record CompiledAssetMatcher( AssetIndex Asset, string CoreName, @@ -33,63 +34,63 @@ public class NewsScraperBackgroundService : BackgroundService ); private readonly IServiceScopeFactory _scopeFactory; - private readonly ILogger _logger; + private readonly IFinlyticLogger _finlyticLogger; private readonly NewsMqttClient _mqttClient; - private readonly int _intervalMinutes; private readonly string _indexPath; - // In-Memory Cache for compiled asset matchers to prevent re-reading & re-compiling Regex private List? _cachedAssetMatchers; private DateTime _lastIndexLoadTime = DateTime.MinValue; public NewsScraperBackgroundService( IServiceScopeFactory scopeFactory, - ILogger logger, + IFinlyticLogger finlyticLogger, NewsMqttClient mqttClient, IConfiguration configuration) { _scopeFactory = scopeFactory; - _logger = logger; + _finlyticLogger = finlyticLogger; _mqttClient = mqttClient; - - _intervalMinutes = configuration.GetValue("ScrapingSettings:IntervalMinutes", 15); _indexPath = Path.Combine(Volumes.IndexRelativePath, "index.json"); } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { - _logger.LogInformation("[{Channel}] NewsScraperBackgroundService started. Interval: {Minutes} minutes.", "NewsChannel", _intervalMinutes); + await _finlyticLogger.LogInfoAsync(SettingKeys.NewsChannel, "[NewsScraperBackgroundService] NewsScraperBackgroundService started."); while (!stoppingToken.IsCancellationRequested) { try { - await RunScrapingCycleAsync(stoppingToken); + using var scope = _scopeFactory.CreateScope(); + var settings = scope.ServiceProvider.GetRequiredService(); + bool enabled = await settings.GetSettingAsync(SettingKeys.EnableAutoScraping, stoppingToken); + + if (enabled) + { + await RunScrapingCycleAsync(stoppingToken); + } + else + { + await _finlyticLogger.LogInfoAsync(SettingKeys.NewsChannel, "[NewsScraperBackgroundService] Auto-scraping is disabled via settings."); + } } catch (Exception ex) when (ex is not OperationCanceledException) { - _logger.LogError(ex, "[{Channel}] An unhandled exception occurred during news scraping cycle.", "NewsChannel"); + await _finlyticLogger.LogErrorAsync(SettingKeys.NewsChannel, ex, "[NewsScraperBackgroundService] An unhandled exception occurred during news scraping cycle."); } - int intervalMinutes = _intervalMinutes; + int intervalMinutes = 15; try { using var scope = _scopeFactory.CreateScope(); - var settingsDb = scope.ServiceProvider.GetService(); - if (settingsDb != null) - { - var settings = await settingsDb.GetSettingsAsync(); - if (settings?.ScrapingIntervalMinutes > 0) - { - intervalMinutes = settings.ScrapingIntervalMinutes; - } - } + var settings = scope.ServiceProvider.GetRequiredService(); + intervalMinutes = await settings.GetSettingAsync(SettingKeys.ScrapeIntervalMinutes, stoppingToken); } - catch { /* Ignore settings DB lookup failures */ } + catch { } - var jitterSeconds = Random.Shared.Next(0, 300); + var jitterSeconds = Random.Shared.Next(0, 60); var nextRunDelay = TimeSpan.FromMinutes(intervalMinutes) + TimeSpan.FromSeconds(jitterSeconds); - _logger.LogInformation("[{Channel}] Scraping cycle completed. Next cycle in {Delay} (interval: {Minutes}m).", "NewsChannel", nextRunDelay, intervalMinutes); + await _finlyticLogger.LogInfoAsync(SettingKeys.NewsChannel, "[NewsScraperBackgroundService] Scraping cycle completed. Next cycle in {Delay} (interval: {Minutes}m).", nextRunDelay, intervalMinutes); try { @@ -101,7 +102,7 @@ public class NewsScraperBackgroundService : BackgroundService } } - _logger.LogInformation("[{Channel}] NewsScraperBackgroundService stopping.", "NewsChannel"); + await _finlyticLogger.LogInfoAsync(SettingKeys.NewsChannel, "[NewsScraperBackgroundService] NewsScraperBackgroundService stopping."); } private async Task RunScrapingCycleAsync(CancellationToken stoppingToken) @@ -111,28 +112,28 @@ public class NewsScraperBackgroundService : BackgroundService var discoveryService = scope.ServiceProvider.GetRequiredService(); var scraperService = scope.ServiceProvider.GetRequiredService(); var n8nService = scope.ServiceProvider.GetRequiredService(); + var settings = scope.ServiceProvider.GetRequiredService(); + + var maxArticlesPerFeed = await settings.GetSettingAsync(SettingKeys.MaxArticlesPerFeed, stoppingToken); - // Load pre-compiled asset index matchers for zero-latency pre-filtering var assetMatchers = await GetOrLoadAssetMatchersAsync(); - _logger.LogInformation("[{Channel}] Loaded {Count} asset index items for text pre-filtering.", "NewsChannel", assetMatchers.Count); + await _finlyticLogger.LogInfoAsync(SettingKeys.NewsChannel, "[NewsScraperBackgroundService] Loaded {Count} asset index items for text pre-filtering.", assetMatchers.Count); - // 1. Scraping Retry Phase: query articles in status "Scraping" (failed Playwright runs) var failedArticles = await dbService.GetArticlesByStatusAsync("Scraping"); if (failedArticles.Count > 0) { - _logger.LogInformation("[{Channel}] Found {Count} articles in status 'Scraping' that failed to scrape previously. Retrying...", "NewsChannel", failedArticles.Count); + await _finlyticLogger.LogInfoAsync(SettingKeys.NewsChannel, "[NewsScraperBackgroundService] Found {Count} articles in status 'Scraping' that failed to scrape previously. Retrying...", failedArticles.Count); foreach (var article in failedArticles) { - if (stoppingToken.IsCancellationRequested) break; - await ProcessSingleArticleAsync(article, dbService, scraperService, n8nService, assetMatchers, stoppingToken); + if (stoppingToken.IsCancellationRequested) return; + await ProcessSingleArticleAsync(article, scraperService, n8nService, dbService, assetMatchers, stoppingToken); } } - // 2. Link Discovery Phase: query RSS feeds and listing pages var sources = await dbService.GetSourcesAsync(); if (sources.Count == 0) { - _logger.LogWarning("[{Channel}] No article sources configured in database. Skipping cycle.", "NewsChannel"); + await _finlyticLogger.LogWarningAsync(SettingKeys.NewsChannel, "[NewsScraperBackgroundService] No article sources configured in database. Skipping cycle."); return; } @@ -140,86 +141,80 @@ public class NewsScraperBackgroundService : BackgroundService { if (stoppingToken.IsCancellationRequested) break; - _logger.LogInformation("[{Channel}] Starting article link discovery for source: {SourceName} ({Url})", "NewsChannel", source.Name, source.Source); + await _finlyticLogger.LogInfoAsync(SettingKeys.NewsChannel, "[NewsScraperBackgroundService] Starting article link discovery for source: {SourceName} ({Url})", source.Name, source.Source); var discoveredArticles = await discoveryService.DiscoverLinksAsync(source.Source, source.Type, stoppingToken); if (discoveredArticles == null || discoveredArticles.Count == 0) { - _logger.LogDebug("No links discovered from source: {SourceName}", source.Name); + await _finlyticLogger.LogDebugAsync(SettingKeys.NewsChannel, "[NewsScraperBackgroundService] No links discovered from source: {SourceName}", source.Name); continue; } - _logger.LogInformation("[{Channel}] Discovered {Count} potential article links from {SourceName}.", "NewsChannel", discoveredArticles.Count, source.Name); + await _finlyticLogger.LogInfoAsync(SettingKeys.NewsChannel, "[NewsScraperBackgroundService] Discovered {Count} potential article links from {SourceName}.", discoveredArticles.Count, source.Name); - foreach (var discovered in discoveredArticles) + var toProcess = discoveredArticles.Take(maxArticlesPerFeed > 0 ? maxArticlesPerFeed : 20); + + foreach (var discovered in toProcess) { if (stoppingToken.IsCancellationRequested) break; - // Idempotency Check & Deduplication - var isDuplicate = await dbService.IsUrlDuplicateAsync(discovered.Url); - if (isDuplicate) + if (await dbService.IsUrlDuplicateAsync(discovered.Url)) { - _logger.LogDebug("Skipping duplicate article URL: {Url}", discovered.Url); + await _finlyticLogger.LogDebugAsync(SettingKeys.NewsChannel, "[NewsScraperBackgroundService] Skipping duplicate article URL: {Url}", discovered.Url); continue; } - // Register Initial Lock State in the database ("Pending") - NewsArticleEntity? article; + NewsArticleEntity article; try { article = await dbService.CreatePendingArticleAsync( - discovered.Url, - discovered.Isins, - discovered.Title, - discovered.Summary, - discovered.PublishedAt, - discovered.Language); + discovered.Url, + discovered.Isins, + discovered.Title, + discovered.Summary, + discovered.PublishedAt, + discovered.Language + ); } catch (Exception ex) { - _logger.LogError(ex, "[{Channel}] Failed to register initial pending state for URL: {Url}. Skipping.", "NewsChannel", discovered.Url); + await _finlyticLogger.LogErrorAsync(SettingKeys.NewsChannel, ex, "[NewsScraperBackgroundService] Failed to register initial pending state for URL: {Url}. Skipping.", discovered.Url); continue; } - if (article == null || article.Id == Guid.Empty) + if (article.Id == Guid.Empty) { - _logger.LogWarning("[{Channel}] Created pending article has invalid/empty ID for URL: {Url}. Skipping.", "NewsChannel", discovered.Url); + await _finlyticLogger.LogWarningAsync(SettingKeys.NewsChannel, "[NewsScraperBackgroundService] Created pending article has invalid/empty ID for URL: {Url}. Skipping.", discovered.Url); continue; } - // Process single article pipeline - await ProcessSingleArticleAsync(article, dbService, scraperService, n8nService, assetMatchers, stoppingToken); + await ProcessSingleArticleAsync(article, scraperService, n8nService, dbService, assetMatchers, stoppingToken); } } } private async Task ProcessSingleArticleAsync( NewsArticleEntity article, - INewsDbService dbService, IPlaywrightScraperService scraperService, IN8nService n8nService, + INewsDbService dbService, List assetMatchers, CancellationToken stoppingToken) { try { - // 3. Extraction with Headless Browser (Transitions to "Processing") await dbService.UpdateArticleStatusAsync(article.Id, "Processing"); - var (resolvedUrl, rawText) = await scraperService.ScrapeArticleAsync(article.SourceUrl); - if (string.IsNullOrWhiteSpace(rawText)) - { - throw new InvalidOperationException("Scraping returned empty text body content."); - } + var (resolvedUrl, rawContent) = await scraperService.ScrapeArticleAsync(article.SourceUrl); - // Update resolved URL if redirect occurred - if (!string.Equals(resolvedUrl, article.SourceUrl, StringComparison.OrdinalIgnoreCase)) + if (!string.IsNullOrWhiteSpace(resolvedUrl) && + !resolvedUrl.Equals(article.SourceUrl, StringComparison.OrdinalIgnoreCase)) { - _logger.LogInformation("[{Channel}] Redirect detected. Initial: {OldUrl} -> Resolved: {NewUrl}", "NewsChannel", article.SourceUrl, resolvedUrl); + await _finlyticLogger.LogInfoAsync(SettingKeys.NewsChannel, "[NewsScraperBackgroundService] Redirect detected. Initial: {OldUrl} -> Resolved: {NewUrl}", article.SourceUrl, resolvedUrl); if (await dbService.IsUrlDuplicateAsync(resolvedUrl)) { - _logger.LogInformation("[{Channel}] Redirected URL {ResolvedUrl} is a duplicate. Terminating processing.", "NewsChannel", resolvedUrl); - await dbService.UpdateArticleStatusAsync(article.Id, "Duplicate"); + await _finlyticLogger.LogInfoAsync(SettingKeys.NewsChannel, "[NewsScraperBackgroundService] Redirected URL {ResolvedUrl} is a duplicate. Terminating processing.", resolvedUrl); + await dbService.UpdateArticleStatusAsync(article.Id, "Failed"); return; } @@ -227,108 +222,85 @@ public class NewsScraperBackgroundService : BackgroundService article.SourceUrl = resolvedUrl; } - // 4. Pre-filtering Assets (Optimized with Pre-Compiled Regex Patterns) - var title = article.Title ?? string.Empty; - var summary = article.Summary ?? string.Empty; - - var preFilteredAssets = assetMatchers.Where(matcher => + if (string.IsNullOrWhiteSpace(rawContent) || rawContent.Length < 60) { - var asset = matcher.Asset; - - // ISIN direct match - if (rawText.Contains(asset.Isin, StringComparison.OrdinalIgnoreCase) || - title.Contains(asset.Isin, StringComparison.OrdinalIgnoreCase) || - summary.Contains(asset.Isin, StringComparison.OrdinalIgnoreCase)) - { - return true; - } - - // Fast regex word boundary check on Full Name - if (matcher.WordRegex != null && (matcher.WordRegex.IsMatch(rawText) || matcher.WordRegex.IsMatch(title))) - { - return true; - } - - // Fast regex word boundary check on Core Name - if (matcher.CoreName.Length >= 3 && matcher.CoreWordRegex != null && - (matcher.CoreWordRegex.IsMatch(rawText) || matcher.CoreWordRegex.IsMatch(title))) - { - return true; - } - - return false; - }) - .Select(m => new FilteredAssetPayload(m.Asset.Name, m.Asset.Isin)) - .ToList(); - - if (preFilteredAssets.Count == 0) - { - _logger.LogInformation("[{Channel}] Pre-filtering: Article {Id} does not reference any known assets. Terminating pipeline.", "NewsChannel", article.Id); await dbService.UpdateArticleStatusAsync(article.Id, "Failed"); return; } - _logger.LogInformation("[{Channel}] Pre-filtering matched {Count} assets for article {Id}.", "NewsChannel", preFilteredAssets.Count, article.Id); + var discoveredIsins = article.MatchedAssets.Select(m => m.Isin).Where(i => !string.IsNullOrEmpty(i)).ToList(); + var preFilteredAssets = PreFilterAssets(rawContent, article.Title, assetMatchers, discoveredIsins); - // 5. Send to n8n Webhook Pipeline - var n8nResponse = await n8nService.AnalyzeArticleAsync(rawText, preFilteredAssets, stoppingToken); - if (n8nResponse == null) + if (preFilteredAssets.Count == 0) { - throw new InvalidOperationException("n8n AI webhook execution returned null or failed."); + await _finlyticLogger.LogInfoAsync(SettingKeys.NewsChannel, "[NewsScraperBackgroundService] Pre-filtering: Article {Id} does not reference any known assets. Terminating pipeline.", article.Id); + await dbService.UpdateArticleStatusAsync(article.Id, "Failed"); + return; } - // 6. Map and Save Completed Classification - var matchedEntities = new List(); - foreach (var n8nAsset in n8nResponse.MatchedAssets) + await _finlyticLogger.LogInfoAsync(SettingKeys.NewsChannel, "[NewsScraperBackgroundService] Pre-filtering matched {Count} assets for article {Id}.", preFilteredAssets.Count, article.Id); + + var n8nResponse = await n8nService.AnalyzeArticleAsync(rawContent, preFilteredAssets, stoppingToken); + if (n8nResponse == null) { - var n8nCoreName = ExtractCoreAssetName(n8nAsset.Name); + await dbService.UpdateArticleStatusAsync(article.Id, "Failed"); + return; + } - var matchedIsin = preFilteredAssets.FirstOrDefault(fa => - fa.Name.Equals(n8nAsset.Name, StringComparison.OrdinalIgnoreCase) || - n8nAsset.Name.Contains(fa.Name, StringComparison.OrdinalIgnoreCase) || - (n8nCoreName.Length >= 3 && ExtractCoreAssetName(fa.Name).Equals(n8nCoreName, StringComparison.OrdinalIgnoreCase)))?.Isin; - - if (string.IsNullOrWhiteSpace(matchedIsin)) + var matchedEntities = new List(); + if (n8nResponse.MatchedAssets != null && n8nResponse.MatchedAssets.Count > 0) + { + foreach (var asset in n8nResponse.MatchedAssets) { - matchedIsin = assetMatchers.FirstOrDefault(m => - m.Asset.Name.Equals(n8nAsset.Name, StringComparison.OrdinalIgnoreCase) || - n8nAsset.Name.Contains(m.Asset.Name, StringComparison.OrdinalIgnoreCase) || - (n8nCoreName.Length >= 3 && m.CoreName.Equals(n8nCoreName, StringComparison.OrdinalIgnoreCase)))?.Asset.Isin; - } + var preMatch = preFilteredAssets.FirstOrDefault(p => p.Name.Equals(asset.Name, StringComparison.OrdinalIgnoreCase)); + var isin = preMatch?.Isin ?? asset.Ticker ?? ""; + if (string.IsNullOrWhiteSpace(isin)) continue; - if (!string.IsNullOrWhiteSpace(matchedIsin)) - { matchedEntities.Add(new MatchedAssetEntity { + Id = Guid.NewGuid(), NewsArticleId = article.Id, - Name = n8nAsset.Name, - Isin = matchedIsin + Isin = isin.Trim().ToUpperInvariant(), + Name = !string.IsNullOrWhiteSpace(asset.Name) ? asset.Name.Trim() : isin.Trim().ToUpperInvariant() }); } } - var completedArticle = await dbService.SaveArticleClassificationAsync(article.Id, n8nResponse, matchedEntities); - - if (completedArticle != null) + if (matchedEntities.Count == 0) + { + foreach (var preMatch in preFilteredAssets) + { + matchedEntities.Add(new MatchedAssetEntity + { + Id = Guid.NewGuid(), + NewsArticleId = article.Id, + Isin = preMatch.Isin, + Name = preMatch.Name + }); + } + } + + var updatedArticle = await dbService.SaveArticleClassificationAsync(article.Id, n8nResponse, matchedEntities); + + if (updatedArticle != null) { - // 7. MQTT Broadcast (Sends completed article to downstream services) var dto = new NewsArticleDto { - Id = completedArticle.Id, - Title = completedArticle.Title, - Author = completedArticle.Author, - Summary = completedArticle.Summary, - ContentRaw = completedArticle.ContentRaw, - Language = completedArticle.Language, - SourceUrl = completedArticle.SourceUrl, - ScrapedAt = completedArticle.ScrapedAt, - PublishedAt = completedArticle.PublishedAt, - MatchedAssets = completedArticle.MatchedAssets.Select(m => new MatchedAssetDto + Id = updatedArticle.Id, + Title = updatedArticle.Title, + Author = updatedArticle.Author, + Summary = updatedArticle.Summary, + ContentRaw = updatedArticle.ContentRaw, + Language = updatedArticle.Language, + SourceUrl = updatedArticle.SourceUrl, + ScrapedAt = updatedArticle.ScrapedAt, + PublishedAt = updatedArticle.PublishedAt, + Status = updatedArticle.Status, + MatchedAssets = updatedArticle.MatchedAssets.Select(m => new MatchedAssetDto { Name = m.Name, Isin = m.Isin - }).ToList(), - Status = completedArticle.Status + }).ToList() }; await _mqttClient.BroadcastArticleAsync(dto); @@ -336,103 +308,120 @@ public class NewsScraperBackgroundService : BackgroundService } catch (Exception ex) { - _logger.LogError(ex, "[{Channel}] Failed to complete processing pipeline for article: {Url}. Transitioning to 'Scraping' for next cycle retry.", "NewsChannel", article.SourceUrl); - try - { - await dbService.UpdateArticleStatusAsync(article.Id, "Scraping"); - } - catch { /* Suppress database secondary errors */ } + await _finlyticLogger.LogErrorAsync(SettingKeys.NewsChannel, ex, "[NewsScraperBackgroundService] Failed to complete processing pipeline for article: {Url}. Transitioning to 'Scraping' for next cycle retry.", article.SourceUrl); + await dbService.UpdateArticleStatusAsync(article.Id, "Scraping"); } } - /// - /// Returns cached compiled asset matchers or parses the index file from disk if stale/missing. - /// + private List PreFilterAssets( + string content, + string? title, + List assetMatchers, + List? priorityIsins = null) + { + if (assetMatchers.Count == 0) return new List(); + + var fullText = (title != null ? title + " " + content : content); + var matched = new Dictionary(StringComparer.OrdinalIgnoreCase); + + if (priorityIsins != null && priorityIsins.Count > 0) + { + foreach (var isin in priorityIsins) + { + var match = assetMatchers.FirstOrDefault(m => string.Equals(m.Asset.Isin, isin, StringComparison.OrdinalIgnoreCase)); + if (match != null && !matched.ContainsKey(match.Asset.Isin)) + { + matched[match.Asset.Isin] = new FilteredAssetPayload(match.Asset.Name, match.Asset.Isin); + } + } + } + + foreach (var m in assetMatchers) + { + if (matched.ContainsKey(m.Asset.Isin)) continue; + + if (fullText.Contains(m.Asset.Isin, StringComparison.OrdinalIgnoreCase)) + { + matched[m.Asset.Isin] = new FilteredAssetPayload(m.Asset.Name, m.Asset.Isin); + continue; + } + + if (m.WordRegex != null && m.WordRegex.IsMatch(fullText)) + { + matched[m.Asset.Isin] = new FilteredAssetPayload(m.Asset.Name, m.Asset.Isin); + continue; + } + + if (m.CoreWordRegex != null && m.CoreWordRegex.IsMatch(fullText)) + { + matched[m.Asset.Isin] = new FilteredAssetPayload(m.Asset.Name, m.Asset.Isin); + } + } + + return matched.Values.ToList(); + } + private async Task> GetOrLoadAssetMatchersAsync() { - if (_cachedAssetMatchers != null && (DateTime.UtcNow - _lastIndexLoadTime).TotalMinutes < 30) + if (_cachedAssetMatchers != null && (DateTime.UtcNow - _lastIndexLoadTime).TotalMinutes < 60) { return _cachedAssetMatchers; } if (!File.Exists(_indexPath)) { - _logger.LogWarning("[{Channel}] Asset index file not found at: {Path}. Pre-filtering will match 0 assets.", "NewsChannel", _indexPath); - return []; + await _finlyticLogger.LogWarningAsync(SettingKeys.NewsChannel, "[NewsScraperBackgroundService] Asset index file not found at: {Path}. Pre-filtering will match 0 assets.", _indexPath); + return new List(); } try { - await using var stream = File.OpenRead(_indexPath); - - // Standard Deserialization for AssetIndex list - var rawList = await JsonSerializer.DeserializeAsync>(stream); + using var stream = File.OpenRead(_indexPath); + var indexList = await JsonSerializer.DeserializeAsync>(stream); - if (rawList != null) + if (indexList == null || indexList.Count == 0) { - _cachedAssetMatchers = rawList.Select(asset => - { - var coreName = ExtractCoreAssetName(asset.Name); - return new CompiledAssetMatcher( - Asset: asset, - CoreName: coreName, - WordRegex: BuildWordRegex(asset.Name), - CoreWordRegex: BuildWordRegex(coreName) - ); - }).ToList(); - - _lastIndexLoadTime = DateTime.UtcNow; - return _cachedAssetMatchers; + return new List(); } + + var compiled = new List(indexList.Count); + foreach (var asset in indexList) + { + if (string.IsNullOrWhiteSpace(asset.Name) || string.IsNullOrWhiteSpace(asset.Isin)) + continue; + + var rawName = asset.Name.Trim(); + var coreName = ExtractCoreName(rawName); + + Regex? wordRegex = null; + if (rawName.Length >= 4) + { + wordRegex = new Regex($@"\b{Regex.Escape(rawName)}\b", RegexOptions.IgnoreCase | RegexOptions.CultureInvariant | RegexOptions.Compiled); + } + + Regex? coreWordRegex = null; + if (!string.IsNullOrWhiteSpace(coreName) && coreName.Length >= 4 && !coreName.Equals(rawName, StringComparison.OrdinalIgnoreCase)) + { + coreWordRegex = new Regex($@"\b{Regex.Escape(coreName)}\b", RegexOptions.IgnoreCase | RegexOptions.CultureInvariant | RegexOptions.Compiled); + } + + compiled.Add(new CompiledAssetMatcher(asset, coreName, wordRegex, coreWordRegex)); + } + + _cachedAssetMatchers = compiled; + _lastIndexLoadTime = DateTime.UtcNow; + return _cachedAssetMatchers; } catch (Exception ex) { - _logger.LogError(ex, "[{Channel}] Failed to read or parse asset index file from {Path}.", "NewsChannel", _indexPath); - } - - return _cachedAssetMatchers ?? []; - } - - /// - /// Helper to pre-compile Word Boundary Regex for an asset name. - /// - private static Regex? BuildWordRegex(string name) - { - if (string.IsNullOrWhiteSpace(name)) return null; - try - { - return new Regex($@"\b{Regex.Escape(name)}\b", RegexOptions.IgnoreCase | RegexOptions.CultureInvariant | RegexOptions.Compiled); - } - catch - { - return null; + await _finlyticLogger.LogErrorAsync(SettingKeys.NewsChannel, ex, "[NewsScraperBackgroundService] Failed to read or parse asset index file from {Path}.", _indexPath); + return _cachedAssetMatchers ?? new List(); } } - /// - /// Extracts the core name of an asset by removing parenthetical metadata and corporate suffixes. - /// - private static string ExtractCoreAssetName(string name) + private static string ExtractCoreName(string rawName) { - if (string.IsNullOrWhiteSpace(name)) return string.Empty; - - int parenIndex = name.IndexOf('('); - if (parenIndex >= 0) - { - name = name[..parenIndex]; - } - - name = name.Trim(); - - var suffixes = new[] { "Inc.", "Inc", "AG", "SE", "Co.", "Co", "Corp.", "Corp", "Ltd.", "Ltd", "plc", "GmbH", "SA", "NV", "Group" }; - foreach (var suffix in suffixes) - { - if (name.EndsWith(" " + suffix, StringComparison.OrdinalIgnoreCase)) - { - name = name[..^suffix.Length].Trim(); - } - } - - return name; + var cleaned = Regex.Replace(rawName, @"\b(AG|SE|SA|NV|PLC|INC|CORP|LLC|GMBH|CO|KG|HOLDING|GROUP|CLASS\s+[A-Z])\b", "", RegexOptions.IgnoreCase); + return cleaned.Trim(' ', '.', ',', '-'); } } \ No newline at end of file diff --git a/FinlyticNews/Services/PlaywrightScraperService.cs b/FinlyticNews/Services/PlaywrightScraperService.cs index 02ae9b8..3c5ec8f 100644 --- a/FinlyticNews/Services/PlaywrightScraperService.cs +++ b/FinlyticNews/Services/PlaywrightScraperService.cs @@ -3,8 +3,9 @@ using System.Collections.Generic; using System.Linq; using System.Threading; using System.Threading.Tasks; +using FinlyticCore.Services; using FinlyticNews.Adapters.Scraping; -using Microsoft.Extensions.Logging; +using FinlyticNews.Util; using Microsoft.Playwright; namespace FinlyticNews.Services; @@ -25,7 +26,7 @@ public interface IPlaywrightScraperService /// public class PlaywrightScraperService : IPlaywrightScraperService, IAsyncDisposable { - private readonly ILogger _logger; + private readonly IFinlyticLogger _finlyticLogger; private readonly IEnumerable _scraperAdapters; private IPlaywright? _playwright; @@ -36,21 +37,20 @@ public class PlaywrightScraperService : IPlaywrightScraperService, IAsyncDisposa /// Initializes a new instance of the class. /// public PlaywrightScraperService( - ILogger logger, + IFinlyticLogger finlyticLogger, IEnumerable scraperAdapters) { - _logger = logger; + _finlyticLogger = finlyticLogger; _scraperAdapters = scraperAdapters; } /// public async Task<(string ResolvedUrl, string Content)> ScrapeArticleAsync(string url) { - _logger.LogInformation("[{Channel}] Launching browser context to scrape article: {Url}", "NewsChannel", url); + await _finlyticLogger.LogInfoAsync(SettingKeys.NewsChannel, "[PlaywrightScraperService] Launching browser context to scrape article: {Url}", url); var browser = await GetOrInitBrowserAsync(); - // Fast, isolated browser context (incognito tab environment) per article await using var context = await browser.NewContextAsync(new BrowserNewContextOptions { UserAgent = "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36", @@ -61,7 +61,6 @@ public class PlaywrightScraperService : IPlaywrightScraperService, IAsyncDisposa try { - // 1. Initial Navigation with DOMContentLoaded wait var response = await page.GotoAsync(url, new PageGotoOptions { WaitUntil = WaitUntilState.DOMContentLoaded, @@ -74,115 +73,86 @@ public class PlaywrightScraperService : IPlaywrightScraperService, IAsyncDisposa } var finalUrl = page.Url; - _logger.LogDebug("[{Channel}] Navigation completed. Initial final URL: {Url}", "NewsChannel", finalUrl); + await _finlyticLogger.LogDebugAsync(SettingKeys.NewsChannel, "[PlaywrightScraperService] Navigation completed. Initial final URL: {Url}", finalUrl); - // 2. Resolve Host Specific Scraper Adapter - var uri = new Uri(url); - var host = uri.Host; - var adapter = _scraperAdapters.FirstOrDefault(a => host.Contains(a.Hostname, StringComparison.OrdinalIgnoreCase)); - - IPage targetPage = page; + var host = new Uri(finalUrl).Host; + var adapter = _scraperAdapters.FirstOrDefault(a => + host.EndsWith(a.Hostname, StringComparison.OrdinalIgnoreCase) || + a.Hostname.EndsWith(host, StringComparison.OrdinalIgnoreCase)); if (adapter != null) { - _logger.LogInformation("[{Channel}] Executing adapter redirect check for host: {Host}", "NewsChannel", adapter.Hostname); + await _finlyticLogger.LogInfoAsync(SettingKeys.NewsChannel, "[PlaywrightScraperService] Executing adapter redirect check for host: {Host}", adapter.Hostname); try { var resolvedRedirectUrl = await adapter.TryResolveRedirectUrlAsync(page); - - // Loop Protection: Navigate only if redirect target is a new URL if (!string.IsNullOrWhiteSpace(resolvedRedirectUrl) && - !string.Equals(page.Url, resolvedRedirectUrl, StringComparison.OrdinalIgnoreCase)) + !resolvedRedirectUrl.Equals(finalUrl, StringComparison.OrdinalIgnoreCase)) { - _logger.LogInformation("[{Channel}] Redirect resolved to target URL: {Url}", "NewsChannel", resolvedRedirectUrl); - - var refererUrl = page.Url; - finalUrl = resolvedRedirectUrl; - + await _finlyticLogger.LogInfoAsync(SettingKeys.NewsChannel, "[PlaywrightScraperService] Redirect resolved to target URL: {Url}", resolvedRedirectUrl); + var redirectResponse = await page.GotoAsync(resolvedRedirectUrl, new PageGotoOptions { WaitUntil = WaitUntilState.DOMContentLoaded, - Timeout = 30000, - Referer = refererUrl + Timeout = 30000 }); if (redirectResponse == null) { - _logger.LogWarning("[{Channel}] Failed to load response for redirect URL: {Url}", "NewsChannel", resolvedRedirectUrl); - } - else - { - finalUrl = page.Url; + await _finlyticLogger.LogWarningAsync(SettingKeys.NewsChannel, "[PlaywrightScraperService] Failed to load response for redirect URL: {Url}", resolvedRedirectUrl); } - // Check if redirect opened a new tab/popup - var matchedPage = page.Context.Pages.FirstOrDefault(p => p.Url == resolvedRedirectUrl); - if (matchedPage != null) - { - targetPage = matchedPage; - } + finalUrl = page.Url; + host = new Uri(finalUrl).Host; + + adapter = _scraperAdapters.FirstOrDefault(a => + host.EndsWith(a.Hostname, StringComparison.OrdinalIgnoreCase) || + a.Hostname.EndsWith(host, StringComparison.OrdinalIgnoreCase)); } } catch (Exception ex) { - _logger.LogWarning(ex, "[{Channel}] Failed to resolve redirect through adapter for host: {Host}. Continuing with current page.", "NewsChannel", adapter.Hostname); + await _finlyticLogger.LogWarningAsync(SettingKeys.NewsChannel, ex, "[PlaywrightScraperService] Failed to resolve redirect through adapter for host: {Host}. Continuing with current page.", adapter.Hostname); } } - // 3. Re-resolve detail page adapter for final target page - var targetUri = new Uri(targetPage.Url); - var targetHost = targetUri.Host; - var targetAdapter = _scraperAdapters.FirstOrDefault(a => targetHost.Contains(a.Hostname, StringComparison.OrdinalIgnoreCase)); - - // Wait a brief moment for dynamic scripts / DOM settling - await targetPage.WaitForTimeoutAsync(1000); - - // 4. Extract Content via Mozilla Readability (or Adapter Fallback) - string extractedText = string.Empty; - - if (targetAdapter != null) + string content; + if (adapter != null) { - var readabilityResult = await targetAdapter.ExtractArticleContentAsync(targetPage); - if (readabilityResult != null && !string.IsNullOrWhiteSpace(readabilityResult.TextContent)) - { - extractedText = readabilityResult.TextContent; - } + var result = await adapter.ExtractArticleContentAsync(page); + content = result?.TextContent ?? await FallbackExtractContentAsync(page); } - - // Standard Fallback: Body Text / Selector Extraction - if (string.IsNullOrWhiteSpace(extractedText)) + else { - var bodySelector = targetAdapter?.ArticleBodySelector ?? "body"; - var locator = targetPage.Locator(bodySelector); - - if (await locator.CountAsync() > 0) - { - extractedText = await locator.First.InnerTextAsync(); - } - - if (string.IsNullOrWhiteSpace(extractedText)) - { - extractedText = await targetPage.EvaluateAsync( - $"() => document.querySelector('{bodySelector}')?.innerText ?? ''"); - } + content = await FallbackExtractContentAsync(page); } - return (finalUrl, extractedText?.Trim() ?? string.Empty); + return (finalUrl, content); } catch (Exception ex) { - _logger.LogError(ex, "[{Channel}] Failed to scrape page content from URL: {Url}", "NewsChannel", url); + await _finlyticLogger.LogErrorAsync(SettingKeys.NewsChannel, ex, "[PlaywrightScraperService] Failed to scrape page content from URL: {Url}", url); throw; } finally { - await context.CloseAsync(); + await page.CloseAsync(); } } - /// - /// Thread-safe singleton initialization of the Chromium browser instance. - /// + private async Task FallbackExtractContentAsync(IPage page) + { + var innerText = await page.EvaluateAsync(@"() => { + const scripts = document.querySelectorAll('script, style, noscript, nav, header, footer, iframe, svg'); + scripts.forEach(s => s.remove()); + + const main = document.querySelector('article, main, .article-content, #content, .story-body') || document.body; + return main ? main.innerText : document.body.innerText; + }"); + + return innerText?.Trim() ?? string.Empty; + } + private async Task GetOrInitBrowserAsync() { if (_browser != null && _browser.IsConnected) @@ -202,10 +172,16 @@ public class PlaywrightScraperService : IPlaywrightScraperService, IAsyncDisposa _browser = await _playwright.Chromium.LaunchAsync(new BrowserTypeLaunchOptions { Headless = true, - Args = ["--no-sandbox", "--disable-setuid-sandbox", "--disable-dev-shm-usage"] + Args = new[] + { + "--no-sandbox", + "--disable-setuid-sandbox", + "--disable-dev-shm-usage", + "--disable-gpu" + } }); - _logger.LogInformation("[{Channel}] Initialized shared Chromium browser instance.", "NewsChannel"); + await _finlyticLogger.LogInfoAsync(SettingKeys.NewsChannel, "[PlaywrightScraperService] Initialized shared Chromium browser instance."); return _browser; } finally @@ -214,9 +190,6 @@ public class PlaywrightScraperService : IPlaywrightScraperService, IAsyncDisposa } } - /// - /// Disposes the Playwright and Browser instances cleanly during service shutdown. - /// public async ValueTask DisposeAsync() { if (_browser != null) diff --git a/FinlyticNews/Util/NewsMqttClient.cs b/FinlyticNews/Util/NewsMqttClient.cs index d3273e4..18bcab1 100644 --- a/FinlyticNews/Util/NewsMqttClient.cs +++ b/FinlyticNews/Util/NewsMqttClient.cs @@ -1,10 +1,18 @@ +using System; +using System.Collections.Generic; using System.IO; +using System.Linq; using System.Text.Json; +using System.Threading; +using System.Threading.Tasks; using FinlyticCore.Dtos; using FinlyticCore.Dtos.News; using FinlyticCore.Dtos.Sentiment; +using FinlyticCore.Dtos.Settings; using FinlyticCore.Models; +using FinlyticCore.Services; using FinlyticCore.Util; +using FinlyticNews.Database; using FinlyticNews.Entities; using FinlyticNews.Services; using Microsoft.Extensions.Configuration; @@ -40,12 +48,10 @@ public class NewsMqttClient : ManagedMqttClient, IHostedService { Host = _configuration["MQTT:Host"] ?? _configuration["MQTT__Host"] ?? "localhost", Port = Convert.ToInt32(_configuration["MQTT:Port"] ?? _configuration["MQTT__Port"] ?? "1883"), - ClientId = - $"{(_configuration["MQTT:ClientId"] ?? _configuration["MQTT__ClientId"] ?? "FinlyticNews")}_{Guid.NewGuid()}" + ClientId = $"{(_configuration["MQTT:ClientId"] ?? _configuration["MQTT__ClientId"] ?? "FinlyticNews")}_{Guid.NewGuid()}" }; - _logger.LogInformation("Starting News MQTT client. Host: {Host}, ClientId: {ClientId}", config.Host, - config.ClientId); + _logger.LogInformation("Starting News MQTT client. Host: {Host}, ClientId: {ClientId}", config.Host, config.ClientId); await ConnectAsync(config); } @@ -61,13 +67,22 @@ public class NewsMqttClient : ManagedMqttClient, IHostedService { _logger.LogInformation("News MQTT client connected. Subscribing to RPC topics..."); - // ZUSAMMENGELEGT: Unified News Fetching (news_Get deckt news_GetDaily mit ab) await SubscribeAsync("services/request/news_Get/#"); await SubscribeAsync("services/request/news_GetById/#"); await SubscribeAsync("services/request/news_GetPending/#"); await SubscribeAsync("services/request/news_UpdateStatus/#"); + await SubscribeAsync("services/request/news_settings_GetAll/#"); + await SubscribeAsync("services/request/news_settings_Update/#"); await SubscribeAsync("services/request/health_Ping/#"); await SubscribeAsync("services/config/updated/#"); + + FinlyticCore.Services.FinlyticLogBroadcaster.OnLogPublished = async (logDto) => + { + if (IsConnected && string.Equals(logDto.ServiceName, "FinlyticNews", StringComparison.OrdinalIgnoreCase)) + { + await PublishAsync("finlytic/logs/FinlyticNews", logDto); + } + }; } /// @@ -75,12 +90,10 @@ public class NewsMqttClient : ManagedMqttClient, IHostedService /// public async Task BroadcastArticleAsync(NewsArticleDto article) { - // 1. Primärer System-Broadcast const string topic = "services/news/completed"; _logger.LogInformation("Broadcasting completed article to MQTT topic: {Topic}. ID: {Id}", topic, article.Id); await PublishAsync(topic, article); - // 2. Zielgerichteter ISIN-Stream für Echtzeit-Frontend-Feeds var firstIsin = article.MatchedAssets.FirstOrDefault()?.Isin; if (!string.IsNullOrWhiteSpace(firstIsin)) { @@ -94,14 +107,12 @@ public class NewsMqttClient : ManagedMqttClient, IHostedService { if (string.IsNullOrWhiteSpace(topic)) return; - // 1. System Config Updates if (topic.StartsWith("services/config/updated", StringComparison.OrdinalIgnoreCase)) { if (topic.EndsWith("FinlyticNews", StringComparison.OrdinalIgnoreCase)) { await OnConfigUpdatedAsync(payload); } - return; } @@ -111,7 +122,6 @@ public class NewsMqttClient : ManagedMqttClient, IHostedService var channel = segments[2]; var correlationId = segments[^1]; - // 2. Unified RPC Dispatching via Switch switch (channel) { case "news_Get": @@ -130,6 +140,14 @@ public class NewsMqttClient : ManagedMqttClient, IHostedService await OnUpdateNewsStatusAsync(payload, correlationId); break; + case "news_settings_GetAll": + await OnSettingsGetAllAsync(correlationId); + break; + + case "news_settings_Update": + await OnSettingsUpdateAsync(payload, correlationId); + break; + case "health_Ping": await OnHealthPingAsync(segments, correlationId); break; @@ -142,7 +160,10 @@ public class NewsMqttClient : ManagedMqttClient, IHostedService private async Task OnGetNewsAsync(string payload, string correlationId) { - _logger.LogInformation("Received RPC news_Get request. Correlation: {CorrelationId}", correlationId); + using var scope = _scopeFactory.CreateScope(); + var finlyticLogger = scope.ServiceProvider.GetRequiredService>(); + + await finlyticLogger.LogInfoAsync(SettingKeys.NewsChannel, "Received RPC news_Get request. Correlation: {CorrelationId}", correlationId); int limit = 20; int offset = 0; @@ -155,9 +176,7 @@ public class NewsMqttClient : ManagedMqttClient, IHostedService { try { - // Direktes Deserialisieren über das DailyNewsRequest-DTO var req = JsonSerializer.Deserialize(payload, FinlyticJsonSerializerContext.Default.DailyNewsRequest); - if (req != null) { limit = req.Limit > 0 ? req.Limit : 20; @@ -167,7 +186,6 @@ public class NewsMqttClient : ManagedMqttClient, IHostedService searchQuery = req.Query; date = req.Date; - // Status anpassen, falls HasSentiment gesetzt ist if (req.HasSentiment == true && string.IsNullOrEmpty(status)) { status = "Analyzed"; @@ -176,44 +194,35 @@ public class NewsMqttClient : ManagedMqttClient, IHostedService } catch (Exception ex) { - _logger.LogWarning(ex, "Failed to parse DailyNewsRequest payload on news_Get"); + await finlyticLogger.LogWarningAsync(SettingKeys.NewsChannel, ex, "Failed to parse DailyNewsRequest payload on news_Get"); } } try { - using var scope = _scopeFactory.CreateScope(); var dbService = scope.ServiceProvider.GetRequiredService(); - var articles = await dbService.GetFilteredNewsAsync(limit, offset, isin, date, status, searchQuery); var dtos = (await Task.WhenAll(articles.Select(a => MapToDtoAsync(a)))).ToList(); string responseTopic = $"services/response/news_Get/{correlationId}"; - _logger.LogInformation("Publishing RPC response to {ResponseTopic} with {Count} articles.", responseTopic, - dtos.Count); + await finlyticLogger.LogInfoAsync(SettingKeys.NewsChannel, "Publishing RPC response to {ResponseTopic} with {Count} articles.", responseTopic, dtos.Count); await PublishAsync(responseTopic, dtos); } catch (Exception ex) { - _logger.LogError(ex, "Failed to compile RPC response for news_Get"); - - // Antworte mit leerer Liste, um RPC-Timeouts im Gateway zu vermeiden + await finlyticLogger.LogErrorAsync(SettingKeys.NewsChannel, ex, "Failed to compile RPC response for news_Get"); try { string responseTopic = $"services/response/news_Get/{correlationId}"; await PublishAsync(responseTopic, new List()); } - catch - { - } + catch { } } } private async Task OnGetNewsByIdAsync(string payload, string correlationId) { - _logger.LogInformation("Received RPC news_GetById request. Correlation: {CorrelationId}", correlationId); string responseTopic = $"services/response/news_GetById/{correlationId}"; - if (string.IsNullOrWhiteSpace(payload)) { await PublishAsync(responseTopic, (object?)null); @@ -222,9 +231,7 @@ public class NewsMqttClient : ManagedMqttClient, IHostedService try { - var request = - JsonSerializer.Deserialize(payload, - FinlyticJsonSerializerContext.Default.ArticleRequest); + var request = JsonSerializer.Deserialize(payload, FinlyticJsonSerializerContext.Default.ArticleRequest); var targetIdStr = request?.ArticleId ?? request?.Id; if (Guid.TryParse(targetIdStr, out var articleId)) @@ -241,33 +248,25 @@ public class NewsMqttClient : ManagedMqttClient, IHostedService } } } - catch (Exception ex) - { - _logger.LogError(ex, "Failed to execute RPC news_GetById"); - } + catch { } await PublishAsync(responseTopic, (object?)null); } private async Task OnGetPendingNewsAsync(string payload, string correlationId) { - _logger.LogInformation("Received RPC news_GetPending request. Correlation: {CorrelationId}", correlationId); int limit = 10; - if (!string.IsNullOrWhiteSpace(payload)) { try { using var doc = JsonDocument.Parse(payload); - if (doc.RootElement.TryGetProperty("limit", out var limitProp) && - limitProp.TryGetInt32(out var parsedLimit)) + if (doc.RootElement.TryGetProperty("limit", out var limitProp) && limitProp.TryGetInt32(out var parsedLimit)) { limit = Math.Min(parsedLimit, 10); } } - catch - { - } + catch { } } try @@ -279,25 +278,17 @@ public class NewsMqttClient : ManagedMqttClient, IHostedService var dtos = (await Task.WhenAll(pendingArticles.Take(limit).Select(a => MapToDtoAsync(a)))).ToList(); string responseTopic = $"services/response/news_GetPending/{correlationId}"; - _logger.LogInformation("Publishing RPC news_GetPending response to {ResponseTopic} with {Count} articles.", - responseTopic, dtos.Count); await PublishAsync(responseTopic, dtos); } - catch (Exception ex) - { - _logger.LogError(ex, "Failed to publish RPC pending news response"); - } + catch { } } private async Task OnUpdateNewsStatusAsync(string payload, string correlationId) { - _logger.LogInformation("Received RPC news_UpdateStatus request. Correlation: {CorrelationId}", correlationId); UpdateNewsStatusResponse response; - try { - var request = JsonSerializer.Deserialize(payload, - FinlyticJsonSerializerContext.Default.UpdateNewsStatusRequest); + var request = JsonSerializer.Deserialize(payload, FinlyticJsonSerializerContext.Default.UpdateNewsStatusRequest); if (request != null && request.Id != Guid.Empty) { using var scope = _scopeFactory.CreateScope(); @@ -313,7 +304,6 @@ public class NewsMqttClient : ManagedMqttClient, IHostedService } catch (Exception ex) { - _logger.LogError(ex, "Failed to execute RPC news_UpdateStatus"); response = new UpdateNewsStatusResponse(false, ex.Message); } @@ -321,42 +311,99 @@ public class NewsMqttClient : ManagedMqttClient, IHostedService await PublishAsync(responseTopic, response); } + private async Task OnSettingsGetAllAsync(string correlationId) + { + using var scope = _scopeFactory.CreateScope(); + var finlyticLogger = scope.ServiceProvider.GetRequiredService>(); + var settingsService = scope.ServiceProvider.GetRequiredService(); + + await finlyticLogger.LogInfoAsync(SettingKeys.MqttChannel, "[FinlyticNews] [Settings_GetAll] Retrieving all dynamic settings via reflection [CorrelationId: {CorrelationId}]", correlationId); + try + { + var settings = await settingsService.GetAllRegisteredSettingsAsync(new[] { typeof(SettingKeys) }); + var responseTopic = $"services/response/news_settings_GetAll/{correlationId}"; + + await PublishAsync(responseTopic, settings); + await finlyticLogger.LogInfoAsync(SettingKeys.MqttChannel, "[FinlyticNews] [Settings_GetAll] Published {Count} settings to '{ResponseTopic}'", settings.Count, responseTopic); + } + catch (Exception ex) + { + await finlyticLogger.LogErrorAsync(SettingKeys.MqttChannel, ex, "[FinlyticNews] [Settings_GetAll] Failed to retrieve settings."); + } + } + + private async Task OnSettingsUpdateAsync(string payload, string correlationId) + { + if (string.IsNullOrWhiteSpace(payload)) return; + + using var scope = _scopeFactory.CreateScope(); + var finlyticLogger = scope.ServiceProvider.GetRequiredService>(); + var settingsService = scope.ServiceProvider.GetRequiredService(); + + await finlyticLogger.LogInfoAsync(SettingKeys.MqttChannel, "[FinlyticNews] [Settings_Update] Processing settings update RPC [CorrelationId: {CorrelationId}]", correlationId); + try + { + Dictionary? updates = null; + try + { + updates = JsonSerializer.Deserialize>(payload); + } + catch + { + var list = JsonSerializer.Deserialize>(payload); + if (list != null) + { + updates = new Dictionary(); + foreach (var item in list) updates[item.Key] = item.Value; + } + } + + if (updates != null && updates.Count > 0) + { + await settingsService.UpdateSettingsAsync(updates); + await finlyticLogger.LogInfoAsync(SettingKeys.MqttChannel, "[FinlyticNews] [Settings_Update] Successfully updated {Count} settings in database and cache.", updates.Count); + } + + var currentSettings = await settingsService.GetAllRegisteredSettingsAsync(new[] { typeof(SettingKeys) }); + var responseTopic = $"services/response/news_settings_Update/{correlationId}"; + await PublishAsync(responseTopic, currentSettings); + } + catch (Exception ex) + { + await finlyticLogger.LogErrorAsync(SettingKeys.MqttChannel, ex, "[FinlyticNews] [Settings_Update] Failed to update settings."); + } + } + private async Task OnConfigUpdatedAsync(string payload) { - _logger.LogInformation("Received config update event for FinlyticNews."); try { using var doc = JsonDocument.Parse(payload); if (doc.RootElement.TryGetProperty("settings", out var settingsProp)) { - var dict = JsonSerializer.Deserialize>(settingsProp.GetRawText()); + var dict = JsonSerializer.Deserialize>(settingsProp.GetRawText()); if (dict != null && dict.Count > 0) { using var scope = _scopeFactory.CreateScope(); - var settingsDb = scope.ServiceProvider.GetRequiredService(); - await settingsDb.UpdateSettingsFromDictionaryAsync(dict); - _logger.LogInformation("Successfully persisted {Count} updated settings.", dict.Count); + var settings = scope.ServiceProvider.GetRequiredService(); + await settings.UpdateSettingsAsync(dict); } } } - catch (Exception ex) - { - _logger.LogError(ex, "Error processing MQTT config update event"); - } + catch { } } private async Task OnHealthPingAsync(string[] segments, string correlationId) { - // Zerlegt den Topic-Pfad z. B. services/request/health_Ping/FinlyticNews/{correlationId} bool isForMe = segments.Length >= 5 && segments[3].Equals("FinlyticNews", StringComparison.OrdinalIgnoreCase); - if (isForMe) { string respTopic = $"services/response/health_Ping/{correlationId}"; - await PublishAsync(respTopic, - new ServiceHealthResponse("FinlyticNews", "Online", DateTime.UtcNow, "Connected")); - _logger.LogInformation("Responded to live health_Ping RPC request [CorrelationId: {CorrelationId}].", - correlationId); + await PublishAsync(respTopic, new ServiceHealthResponse("FinlyticNews", "Online", DateTime.UtcNow, "Connected")); + + using var scope = _scopeFactory.CreateScope(); + var finlyticLogger = scope.ServiceProvider.GetRequiredService>(); + await finlyticLogger.LogInfoAsync(SettingKeys.HealthPingChannel, "Responded to live health_Ping RPC request [CorrelationId: {CorrelationId}].", correlationId); } } @@ -375,41 +422,31 @@ public class NewsMqttClient : ManagedMqttClient, IHostedService var targetId = a.Id.ToString(); IsinAnalysisEntry? sentimentEntry = null; - // 1. Snappy Local Disk Check for Article File - var articlePath = Path.Combine(Directory.GetCurrentDirectory(), "data", "summaries", "articles", - $"{targetId}.json"); + var articlePath = Path.Combine(Directory.GetCurrentDirectory(), "data", "summaries", "articles", $"{targetId}.json"); if (File.Exists(articlePath)) { try { var json = await File.ReadAllTextAsync(articlePath); - sentimentEntry = JsonSerializer.Deserialize(json, - FinlyticJsonSerializerContext.Default.IsinAnalysisEntry); - } - catch - { + sentimentEntry = JsonSerializer.Deserialize(json, FinlyticJsonSerializerContext.Default.IsinAnalysisEntry); } + catch { } } - // 2. ISIN Summary File Fallback if (sentimentEntry == null && a.MatchedAssets != null && a.MatchedAssets.Count > 0) { foreach (var asset in a.MatchedAssets) { if (string.IsNullOrWhiteSpace(asset.Isin)) continue; - var isinPath = Path.Combine(Directory.GetCurrentDirectory(), "data", "summaries", "isin", - $"{asset.Isin.Trim()}.json"); + var isinPath = Path.Combine(Directory.GetCurrentDirectory(), "data", "summaries", "isin", $"{asset.Isin.Trim()}.json"); if (File.Exists(isinPath)) { try { var json = await File.ReadAllTextAsync(isinPath); - var isinDoc = JsonSerializer.Deserialize(json, - FinlyticJsonSerializerContext.Default.IsinSentimentSummaryDto); - var match = isinDoc?.Analyses?.FirstOrDefault(entry => - string.Equals(entry.Article?.ArticleId?.Trim(), targetId, - StringComparison.OrdinalIgnoreCase)); + var isinDoc = JsonSerializer.Deserialize(json, FinlyticJsonSerializerContext.Default.IsinSentimentSummaryDto); + var match = isinDoc?.Analyses?.FirstOrDefault(entry => string.Equals(entry.Article?.ArticleId?.Trim(), targetId, StringComparison.OrdinalIgnoreCase)); if (match != null) { @@ -417,9 +454,7 @@ public class NewsMqttClient : ManagedMqttClient, IHostedService break; } } - catch - { - } + catch { } } } } @@ -432,10 +467,7 @@ public class NewsMqttClient : ManagedMqttClient, IHostedService confidence = finbertResult.Confidence; } } - catch (Exception ex) - { - _logger.LogTrace(ex, "[MapToDtoAsync] Sentiment fetch skipped for article {Id}", a.Id); - } + catch { } return new NewsArticleDto { diff --git a/FinlyticNews/Util/SettingKeys.cs b/FinlyticNews/Util/SettingKeys.cs new file mode 100644 index 0000000..e8cac49 --- /dev/null +++ b/FinlyticNews/Util/SettingKeys.cs @@ -0,0 +1,24 @@ +using FinlyticCore.Models.Settings; + +namespace FinlyticNews.Util; + +public static class SettingKeys +{ + // --- Logging-Kanäle --- + public static readonly SettingKey NewsChannel = new("Logging.Channel.News", true); + public static readonly SettingKey MqttChannel = new("Logging.Channel.MQTT", true); + public static readonly SettingKey HealthPingChannel = new("Logging.Channel.Health", true); + + // --- Scraping & Feed-Konfiguration --- + public static readonly SettingKey ScrapeIntervalMinutes = new("Scraping.IntervalMinutes", 15); + public static readonly SettingKey MaxArticlesPerFeed = new("Scraping.MaxArticlesPerFeed", 20); + public static readonly SettingKey EnableAutoScraping = new("Feature.EnableAutoScraping", true); + public static readonly SettingKey HttpTimeoutSeconds = new("Scraping.HttpTimeoutSeconds", 20); + + // --- KI & Sentiment-Konfiguration --- + public static readonly SettingKey FinBertBatchSize = new("AI.FinBertBatchSize", 8); + public static readonly SettingKey MinSentimentConfidence = new("AI.MinSentimentConfidence", 0.65); + + // --- Daten-Retention & Cleanup --- + public static readonly SettingKey ArticleRetentionDays = new("Data.ArticleRetentionDays", 90); +}