From a1f2b888f669e6ce0d3f0590b07bfe86b96de49e Mon Sep 17 00:00:00 2001 From: Kleidukos Date: Sat, 15 Aug 2026 21:29:50 +0200 Subject: [PATCH] feat(fundamentals): dynamic settings, IFinlyticLogger, live log streaming, and EF migration --- .../Database/FundamentalsDbContext.cs | 16 +- ...60815183935_AddDynamicSettings.Designer.cs | 430 ++++++++++++++++++ .../20260815183935_AddDynamicSettings.cs | 37 ++ .../FundamentalsDbContextModelSnapshot.cs | 3 +- FinlyticFundamentals/Program.cs | 32 +- .../Services/FundamentalsDbService.cs | 7 +- .../Util/FundamentalsMqttClient.cs | 103 ++++- 7 files changed, 600 insertions(+), 28 deletions(-) create mode 100644 FinlyticFundamentals/Migrations/20260815183935_AddDynamicSettings.Designer.cs create mode 100644 FinlyticFundamentals/Migrations/20260815183935_AddDynamicSettings.cs diff --git a/FinlyticFundamentals/Database/FundamentalsDbContext.cs b/FinlyticFundamentals/Database/FundamentalsDbContext.cs index 157425b..710439b 100644 --- a/FinlyticFundamentals/Database/FundamentalsDbContext.cs +++ b/FinlyticFundamentals/Database/FundamentalsDbContext.cs @@ -1,10 +1,12 @@ +using FinlyticCore.Database; using FinlyticCore.Entities.Settings; using FinlyticFundamentals.Entities; using Microsoft.EntityFrameworkCore; +using Microsoft.EntityFrameworkCore.Design; namespace FinlyticFundamentals.Database; -public class FundamentalsDbContext : DbContext +public class FundamentalsDbContext : DbContext, ISettingsDbContext { public FundamentalsDbContext(DbContextOptions options) : base(options) { @@ -23,7 +25,7 @@ public class FundamentalsDbContext : DbContext modelBuilder.Entity(entity => { entity.HasKey(e => e.Id); - entity.HasIndex(e => e.Key); + entity.HasIndex(e => e.Key).IsUnique(); }); modelBuilder.Entity(entity => @@ -93,3 +95,13 @@ public class FundamentalsDbContext : DbContext }); } } + +public class FundamentalsDbContextFactory : IDesignTimeDbContextFactory +{ + public FundamentalsDbContext CreateDbContext(string[] args) + { + var optionsBuilder = new DbContextOptionsBuilder(); + optionsBuilder.UseNpgsql("Host=localhost;Database=fundamentals;Username=postgres;Password=postgres"); + return new FundamentalsDbContext(optionsBuilder.Options); + } +} diff --git a/FinlyticFundamentals/Migrations/20260815183935_AddDynamicSettings.Designer.cs b/FinlyticFundamentals/Migrations/20260815183935_AddDynamicSettings.Designer.cs new file mode 100644 index 0000000..f443feb --- /dev/null +++ b/FinlyticFundamentals/Migrations/20260815183935_AddDynamicSettings.Designer.cs @@ -0,0 +1,430 @@ +// +using System; +using FinlyticFundamentals.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 FinlyticFundamentals.Migrations +{ + [DbContext(typeof(FundamentalsDbContext))] + [Migration("20260815183935_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("FinlyticFundamentals.Entities.AssetDataEntity", b => + { + b.Property("Isin") + .HasColumnType("text"); + + b.Property("Description") + .IsRequired() + .HasColumnType("text"); + + b.Property("Name") + .IsRequired() + .HasColumnType("text"); + + b.HasKey("Isin"); + + b.ToTable("AssetData"); + }); + + modelBuilder.Entity("FinlyticFundamentals.Entities.AssetEventEntity", b => + { + b.Property("Id") + .ValueGeneratedOnAdd() + .HasColumnType("uuid"); + + b.Property("AssetDataIsin") + .IsRequired() + .HasColumnType("text"); + + b.Property("Date") + .HasColumnType("timestamp with time zone"); + + b.Property("Type") + .IsRequired() + .HasColumnType("text"); + + b.HasKey("Id"); + + b.HasIndex("AssetDataIsin"); + + b.ToTable("AssetEvents"); + }); + + modelBuilder.Entity("FinlyticFundamentals.Entities.FundamentalDataEntity", b => + { + b.Property("Isin") + .HasColumnType("text"); + + b.Property("AssetDataIsin") + .IsRequired() + .HasColumnType("text"); + + b.Property("ConsensusRating") + .HasColumnType("text"); + + b.Property("CurrentRatio") + .HasColumnType("numeric"); + + b.Property("DebtToEquity") + .HasColumnType("numeric"); + + b.Property("DilutedEps") + .HasColumnType("numeric"); + + b.Property("Ebitda") + .HasColumnType("numeric"); + + b.Property("EnterpriseValue") + .HasColumnType("numeric"); + + b.Property("EvToEbitda") + .HasColumnType("numeric"); + + b.Property("FiftyTwoWeekHigh") + .HasColumnType("numeric"); + + b.Property("FiftyTwoWeekLow") + .HasColumnType("numeric"); + + b.Property("ForwardDividendYield") + .HasColumnType("numeric"); + + b.Property("ForwardPe") + .HasColumnType("numeric"); + + b.Property("FreeCashFlow") + .HasColumnType("numeric"); + + b.Property("GrossProfit") + .HasColumnType("numeric"); + + b.Property("LastUpdatedUtc") + .HasColumnType("timestamp with time zone"); + + b.Property("MarketCap") + .HasColumnType("numeric"); + + b.Property("NetIncome") + .HasColumnType("numeric"); + + b.Property("OperatingCashFlow") + .HasColumnType("numeric"); + + b.Property("OperatingIncome") + .HasColumnType("numeric"); + + b.Property("PayoutRatio") + .HasColumnType("numeric"); + + b.Property("PegRatio") + .HasColumnType("numeric"); + + b.Property("PercentHeldByInsiders") + .HasColumnType("numeric"); + + b.Property("PercentHeldByInstitutions") + .HasColumnType("numeric"); + + b.Property("PriceTargetHigh") + .HasColumnType("numeric"); + + b.Property("PriceTargetLow") + .HasColumnType("numeric"); + + b.Property("PriceTargetMean") + .HasColumnType("numeric"); + + b.Property("PriceToBook") + .HasColumnType("numeric"); + + b.Property("PriceToSales") + .HasColumnType("numeric"); + + b.Property("ReturnOnAssets") + .HasColumnType("numeric"); + + b.Property("ReturnOnEquity") + .HasColumnType("numeric"); + + b.Property("RevenueGrowthYoY") + .HasColumnType("numeric"); + + b.Property("ShortPercentOfFloat") + .HasColumnType("numeric"); + + b.Property("ShortRatio") + .HasColumnType("numeric"); + + b.Property("TotalCash") + .HasColumnType("numeric"); + + b.Property("TotalDebt") + .HasColumnType("numeric"); + + b.Property("TotalRevenue") + .HasColumnType("numeric"); + + b.Property("TrailingPe") + .HasColumnType("numeric"); + + b.HasKey("Isin"); + + b.HasIndex("AssetDataIsin"); + + b.ToTable("FundamentalData"); + }); + + modelBuilder.Entity("FinlyticFundamentals.Entities.KeyExecutiveEntity", b => + { + b.Property("Id") + .ValueGeneratedOnAdd() + .HasColumnType("uuid"); + + b.Property("AssetDataIsin") + .IsRequired() + .HasColumnType("text"); + + b.Property("Name") + .IsRequired() + .HasColumnType("text"); + + b.Property("Payment") + .IsRequired() + .HasColumnType("text"); + + b.Property("SortOrder") + .ValueGeneratedOnAdd() + .HasColumnType("integer") + .HasDefaultValue(0); + + b.Property("Title") + .IsRequired() + .HasColumnType("text"); + + b.HasKey("Id"); + + b.HasIndex("AssetDataIsin", "SortOrder"); + + b.ToTable("KeyExecutives"); + }); + + modelBuilder.Entity("FinlyticFundamentals.Entities.AssetDataEntity", b => + { + b.OwnsOne("FinlyticFundamentals.Entities.TickerEntity", "PrimaryTicker", b1 => + { + b1.Property("AssetDataEntityIsin") + .HasColumnType("text"); + + b1.Property("Exchange") + .IsRequired() + .ValueGeneratedOnAdd() + .HasColumnType("text") + .HasDefaultValue("") + .HasColumnName("PrimaryTickerExchange"); + + b1.Property("Ticker") + .IsRequired() + .ValueGeneratedOnAdd() + .HasColumnType("text") + .HasDefaultValue("") + .HasColumnName("PrimaryTicker"); + + b1.HasKey("AssetDataEntityIsin"); + + b1.ToTable("AssetData"); + + b1.WithOwner() + .HasForeignKey("AssetDataEntityIsin"); + }); + + b.OwnsMany("FinlyticFundamentals.Entities.TickerEntity", "AvailableTickers", b1 => + { + b1.Property("Id") + .ValueGeneratedOnAdd() + .HasColumnType("uuid"); + + b1.Property("AssetDataIsin") + .IsRequired() + .HasColumnType("text"); + + b1.Property("Exchange") + .IsRequired() + .HasColumnType("text") + .HasColumnName("Exchange"); + + b1.Property("Ticker") + .IsRequired() + .HasColumnType("text") + .HasColumnName("Ticker"); + + b1.HasKey("Id"); + + b1.HasIndex("AssetDataIsin"); + + b1.HasIndex("Ticker"); + + b1.ToTable("Tickers", (string)null); + + b1.WithOwner() + .HasForeignKey("AssetDataIsin"); + }); + + b.Navigation("AvailableTickers"); + + b.Navigation("PrimaryTicker") + .IsRequired(); + }); + + modelBuilder.Entity("FinlyticFundamentals.Entities.AssetEventEntity", b => + { + b.HasOne("FinlyticFundamentals.Entities.AssetDataEntity", "AssetData") + .WithMany("AssetEvents") + .HasForeignKey("AssetDataIsin") + .OnDelete(DeleteBehavior.Cascade) + .IsRequired(); + + b.OwnsOne("FinlyticFundamentals.Entities.TickerEntity", "Ticker", b1 => + { + b1.Property("AssetEventEntityId") + .HasColumnType("uuid"); + + b1.Property("Exchange") + .IsRequired() + .ValueGeneratedOnAdd() + .HasColumnType("text") + .HasDefaultValue("") + .HasColumnName("TickerExchange"); + + b1.Property("Ticker") + .IsRequired() + .ValueGeneratedOnAdd() + .HasColumnType("text") + .HasDefaultValue("") + .HasColumnName("Ticker"); + + b1.HasKey("AssetEventEntityId"); + + b1.ToTable("AssetEvents"); + + b1.WithOwner() + .HasForeignKey("AssetEventEntityId"); + }); + + b.Navigation("AssetData"); + + b.Navigation("Ticker") + .IsRequired(); + }); + + modelBuilder.Entity("FinlyticFundamentals.Entities.FundamentalDataEntity", b => + { + b.HasOne("FinlyticFundamentals.Entities.AssetDataEntity", "AssetData") + .WithMany("FundamentalData") + .HasForeignKey("AssetDataIsin") + .OnDelete(DeleteBehavior.Cascade) + .IsRequired(); + + b.OwnsOne("FinlyticFundamentals.Entities.TickerEntity", "Ticker", b1 => + { + b1.Property("FundamentalDataEntityIsin") + .HasColumnType("text"); + + b1.Property("Exchange") + .IsRequired() + .ValueGeneratedOnAdd() + .HasColumnType("text") + .HasDefaultValue("") + .HasColumnName("TickerExchange"); + + b1.Property("Ticker") + .IsRequired() + .ValueGeneratedOnAdd() + .HasColumnType("text") + .HasDefaultValue("") + .HasColumnName("Ticker"); + + b1.HasKey("FundamentalDataEntityIsin"); + + b1.ToTable("FundamentalData"); + + b1.WithOwner() + .HasForeignKey("FundamentalDataEntityIsin"); + }); + + b.Navigation("AssetData"); + + b.Navigation("Ticker") + .IsRequired(); + }); + + modelBuilder.Entity("FinlyticFundamentals.Entities.KeyExecutiveEntity", b => + { + b.HasOne("FinlyticFundamentals.Entities.AssetDataEntity", "AssetData") + .WithMany("KeyExecutives") + .HasForeignKey("AssetDataIsin") + .OnDelete(DeleteBehavior.Cascade) + .IsRequired(); + + b.Navigation("AssetData"); + }); + + modelBuilder.Entity("FinlyticFundamentals.Entities.AssetDataEntity", b => + { + b.Navigation("AssetEvents"); + + b.Navigation("FundamentalData"); + + b.Navigation("KeyExecutives"); + }); +#pragma warning restore 612, 618 + } + } +} diff --git a/FinlyticFundamentals/Migrations/20260815183935_AddDynamicSettings.cs b/FinlyticFundamentals/Migrations/20260815183935_AddDynamicSettings.cs new file mode 100644 index 0000000..9331f04 --- /dev/null +++ b/FinlyticFundamentals/Migrations/20260815183935_AddDynamicSettings.cs @@ -0,0 +1,37 @@ +using Microsoft.EntityFrameworkCore.Migrations; + +#nullable disable + +namespace FinlyticFundamentals.Migrations +{ + /// + public partial class AddDynamicSettings : Migration + { + /// + protected override void Up(MigrationBuilder migrationBuilder) + { + migrationBuilder.DropIndex( + name: "IX_DynamicSettings_Key", + table: "DynamicSettings"); + + migrationBuilder.CreateIndex( + name: "IX_DynamicSettings_Key", + table: "DynamicSettings", + column: "Key", + unique: true); + } + + /// + protected override void Down(MigrationBuilder migrationBuilder) + { + migrationBuilder.DropIndex( + name: "IX_DynamicSettings_Key", + table: "DynamicSettings"); + + migrationBuilder.CreateIndex( + name: "IX_DynamicSettings_Key", + table: "DynamicSettings", + column: "Key"); + } + } +} diff --git a/FinlyticFundamentals/Migrations/FundamentalsDbContextModelSnapshot.cs b/FinlyticFundamentals/Migrations/FundamentalsDbContextModelSnapshot.cs index 81e534f..0aa9310 100644 --- a/FinlyticFundamentals/Migrations/FundamentalsDbContextModelSnapshot.cs +++ b/FinlyticFundamentals/Migrations/FundamentalsDbContextModelSnapshot.cs @@ -47,7 +47,8 @@ namespace FinlyticFundamentals.Migrations b.HasKey("Id"); - b.HasIndex("Key"); + b.HasIndex("Key") + .IsUnique(); b.ToTable("DynamicSettings"); }); diff --git a/FinlyticFundamentals/Program.cs b/FinlyticFundamentals/Program.cs index 4d4ecb0..3dbacaa 100644 --- a/FinlyticFundamentals/Program.cs +++ b/FinlyticFundamentals/Program.cs @@ -1,20 +1,26 @@ using FinlyticCore.Clients; +using FinlyticCore.Database; using FinlyticCore.Services; using FinlyticCore.Services.PlaywrightScrapper; using FinlyticCore.Services.TradeRepublic; +using FinlyticCore.Services.Yahoo; using Microsoft.EntityFrameworkCore; using FinlyticFundamentals.Database; using FinlyticFundamentals.Services; using FinlyticFundamentals.Util; - using Microsoft.EntityFrameworkCore.Diagnostics; +using Microsoft.Extensions.Configuration; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Hosting; +using Microsoft.Extensions.Logging; var builder = Host.CreateApplicationBuilder(args); -// Register DB Context +// Register DB Context & ISettingsDbContext builder.Services.AddDbContext(options => options.UseNpgsql(builder.Configuration.GetConnectionString("DefaultConnection")) .ConfigureWarnings(w => w.Ignore(RelationalEventId.PendingModelChangesWarning))); +builder.Services.AddScoped(sp => sp.GetRequiredService()); // Register HTTP Clients builder.Services.AddHttpClient() @@ -25,22 +31,23 @@ builder.Services.AddHttpClient() AllowAutoRedirect = true }); -builder.Services.AddSingleton(); -builder.Services.AddSingleton(); -builder.Services.AddTransient>(); +builder.Services.AddScoped(); +builder.Services.AddScoped(); +builder.Services.AddTransient(); // Register Application Services -builder.Services.AddSingleton(); -builder.Services.AddSingleton(); -builder.Services.AddSingleton(); -builder.Services.AddTransient(); +builder.Services.AddScoped(); +builder.Services.AddScoped(); +builder.Services.AddSingleton(); +builder.Services.AddScoped(); builder.Services.AddScoped(); -builder.Services.AddScoped(typeof(ISettingsService<>), typeof(SettingsService<>)); -builder.Services.AddScoped(typeof(IFinlyticLogger<,>), typeof(FinlyticLogger<,>)); +builder.Services.AddSingleton(); +builder.Services.AddSingleton(typeof(IFinlyticLogger<>), typeof(FinlyticLogger<>)); // Register MQTT Client (as a Hosted Service) -builder.Services.AddHostedService(); +builder.Services.AddSingleton(); +builder.Services.AddHostedService(sp => sp.GetRequiredService()); var host = builder.Build(); @@ -52,7 +59,6 @@ using (var scope = host.Services.CreateScope()) var context = scope.ServiceProvider.GetRequiredService(); await context.Database.MigrateAsync(); Console.WriteLine("Database migrations successfully executed for FinlyticFundamentals."); - } catch (Exception ex) { diff --git a/FinlyticFundamentals/Services/FundamentalsDbService.cs b/FinlyticFundamentals/Services/FundamentalsDbService.cs index 593c95f..0feccb4 100644 --- a/FinlyticFundamentals/Services/FundamentalsDbService.cs +++ b/FinlyticFundamentals/Services/FundamentalsDbService.cs @@ -10,6 +10,7 @@ using FinlyticCore.Dtos.Yahoo; using FinlyticCore.Models.Settings; using FinlyticCore.Services; using FinlyticCore.Services.TradeRepublic; +using FinlyticCore.Services.Yahoo; using FinlyticFundamentals.Database; using FinlyticFundamentals.Entities; using FinlyticFundamentals.Util; @@ -39,13 +40,13 @@ public class FundamentalsDbService : IFundamentalsDbService private readonly IServiceScopeFactory _scopeFactory; private readonly IYahooFinanceScraper _scraper; private readonly ITradeRepublicService _tradeRepublicService; - private readonly IFinlyticLogger _finlyticLogger; + private readonly IFinlyticLogger _finlyticLogger; public FundamentalsDbService( IServiceScopeFactory scopeFactory, IYahooFinanceScraper scraper, ITradeRepublicService tradeRepublicService, - IFinlyticLogger finlyticLogger) + IFinlyticLogger finlyticLogger) { _scopeFactory = scopeFactory; _scraper = scraper; @@ -71,7 +72,7 @@ public class FundamentalsDbService : IFundamentalsDbService { using var scope = _scopeFactory.CreateScope(); var context = scope.ServiceProvider.GetRequiredService(); - var settingsService = scope.ServiceProvider.GetRequiredService>(); + var settingsService = scope.ServiceProvider.GetRequiredService(); // 1. Dynamic Settings lesen bool allowForceRefresh = diff --git a/FinlyticFundamentals/Util/FundamentalsMqttClient.cs b/FinlyticFundamentals/Util/FundamentalsMqttClient.cs index 3bbbd10..5de7336 100644 --- a/FinlyticFundamentals/Util/FundamentalsMqttClient.cs +++ b/FinlyticFundamentals/Util/FundamentalsMqttClient.cs @@ -1,12 +1,13 @@ using System; +using System.Collections.Generic; using System.Text.Json; using System.Threading; using System.Threading.Tasks; using FinlyticCore.Dtos; +using FinlyticCore.Dtos.Settings; using FinlyticCore.Models; using FinlyticCore.Services; using FinlyticCore.Util; -using FinlyticFundamentals.Database; using FinlyticFundamentals.Services; using Microsoft.Extensions.Configuration; using Microsoft.Extensions.DependencyInjection; @@ -60,7 +61,17 @@ public class FundamentalsMqttClient : ManagedMqttClient, IHostedService await SubscribeAsync("services/request/fundamentals_Get/#"); await SubscribeAsync("services/request/events_GetAll/#"); await SubscribeAsync("services/request/events_GetByMonth/#"); + await SubscribeAsync("services/request/fundamentals_settings_GetAll/#"); + await SubscribeAsync("services/request/fundamentals_settings_Update/#"); await SubscribeAsync("services/request/health_Ping/#"); + + FinlyticLogBroadcaster.OnLogPublished = async (logDto) => + { + if (IsConnected && string.Equals(logDto.ServiceName, "FinlyticFundamentals", StringComparison.OrdinalIgnoreCase)) + { + await PublishAsync("finlytic/logs/FinlyticFundamentals", logDto); + } + }; } /// @@ -85,6 +96,14 @@ public class FundamentalsMqttClient : ManagedMqttClient, IHostedService { await OnEventsGetByMonthAsync(payload, correlationId); } + else if (topic.StartsWith("services/request/fundamentals_settings_GetAll", StringComparison.OrdinalIgnoreCase)) + { + await OnSettingsGetAllAsync(correlationId); + } + else if (topic.StartsWith("services/request/fundamentals_settings_Update", StringComparison.OrdinalIgnoreCase)) + { + await OnSettingsUpdateAsync(payload, correlationId); + } else if (topic.StartsWith("services/request/health_Ping", StringComparison.OrdinalIgnoreCase)) { await OnHealthPingAsync(topic, correlationId); @@ -93,8 +112,8 @@ public class FundamentalsMqttClient : ManagedMqttClient, IHostedService private async Task OnFundamentalsGetAsync(string payload, string correlationId) { - using var scope = _scopeFactory.CreateScope(); - var finlyticLogger = scope.ServiceProvider.GetRequiredService>(); + await using var scope = _scopeFactory.CreateAsyncScope(); + var finlyticLogger = scope.ServiceProvider.GetRequiredService>(); var dbService = scope.ServiceProvider.GetRequiredService(); if (string.IsNullOrWhiteSpace(payload)) @@ -129,8 +148,8 @@ public class FundamentalsMqttClient : ManagedMqttClient, IHostedService private async Task OnEventsGetAllAsync(string correlationId) { - using var scope = _scopeFactory.CreateScope(); - var finlyticLogger = scope.ServiceProvider.GetRequiredService>(); + await using var scope = _scopeFactory.CreateAsyncScope(); + var finlyticLogger = scope.ServiceProvider.GetRequiredService>(); var dbService = scope.ServiceProvider.GetRequiredService(); await finlyticLogger.LogInfoAsync(SettingKeys.MqttChannel, "[FinlyticFundamentals] [MQTT_Client] Processing RPC events_GetAll request [CorrelationId: {CorrelationId}]", correlationId); @@ -152,8 +171,8 @@ public class FundamentalsMqttClient : ManagedMqttClient, IHostedService { if (string.IsNullOrWhiteSpace(payload)) return; - using var scope = _scopeFactory.CreateScope(); - var finlyticLogger = scope.ServiceProvider.GetRequiredService>(); + await using var scope = _scopeFactory.CreateAsyncScope(); + var finlyticLogger = scope.ServiceProvider.GetRequiredService>(); var dbService = scope.ServiceProvider.GetRequiredService(); try @@ -175,12 +194,78 @@ public class FundamentalsMqttClient : ManagedMqttClient, IHostedService } } + private async Task OnSettingsGetAllAsync(string correlationId) + { + await using var scope = _scopeFactory.CreateAsyncScope(); + var finlyticLogger = scope.ServiceProvider.GetRequiredService>(); + var settingsService = scope.ServiceProvider.GetRequiredService(); + + await finlyticLogger.LogInfoAsync(SettingKeys.MqttChannel, "[FinlyticFundamentals] [Settings_GetAll] Retrieving all dynamic settings via reflection [CorrelationId: {CorrelationId}]", correlationId); + try + { + var settings = await settingsService.GetAllRegisteredSettingsAsync(new[] { typeof(SettingKeys) }); + var responseTopic = $"services/response/fundamentals_settings_GetAll/{correlationId}"; + + await PublishAsync(responseTopic, settings); + await finlyticLogger.LogInfoAsync(SettingKeys.MqttChannel, "[FinlyticFundamentals] [Settings_GetAll] Published {Count} settings to '{ResponseTopic}'", settings.Count, responseTopic); + } + catch (Exception ex) + { + await finlyticLogger.LogErrorAsync(SettingKeys.MqttChannel, ex, "[FinlyticFundamentals] [Settings_GetAll] Failed to retrieve settings."); + } + } + + private async Task OnSettingsUpdateAsync(string payload, string correlationId) + { + if (string.IsNullOrWhiteSpace(payload)) return; + + await using var scope = _scopeFactory.CreateAsyncScope(); + var finlyticLogger = scope.ServiceProvider.GetRequiredService>(); + var settingsService = scope.ServiceProvider.GetRequiredService(); + + await finlyticLogger.LogInfoAsync(SettingKeys.MqttChannel, "[FinlyticFundamentals] [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, "[FinlyticFundamentals] [Settings_Update] Successfully updated {Count} settings in database and cache.", updates.Count); + } + + var currentSettings = await settingsService.GetAllRegisteredSettingsAsync(new[] { typeof(SettingKeys) }); + var responseTopic = $"services/response/fundamentals_settings_Update/{correlationId}"; + await PublishAsync(responseTopic, currentSettings); + } + catch (Exception ex) + { + await finlyticLogger.LogErrorAsync(SettingKeys.MqttChannel, ex, "[FinlyticFundamentals] [Settings_Update] Failed to update settings."); + } + } + private async Task OnHealthPingAsync(string topic, string correlationId) { if (topic.Contains("FinlyticFundamentals", StringComparison.OrdinalIgnoreCase) || !topic.Contains("/", StringComparison.OrdinalIgnoreCase)) { - using var scope = _scopeFactory.CreateScope(); - var finlyticLogger = scope.ServiceProvider.GetRequiredService>(); + await using var scope = _scopeFactory.CreateAsyncScope(); + var finlyticLogger = scope.ServiceProvider.GetRequiredService>(); var respTopic = $"services/response/health_Ping/{correlationId}"; await PublishAsync(respTopic, new ServiceHealthResponse("FinlyticFundamentals", "Online", DateTime.UtcNow, "Connected"));