R3 Slice 3: CongressTrading-Modul + Worker-Log + core_settings auf EF Core

- Modul: CongressTradingDbContext (ct_congressMember/ct_trade) + Design-Time-Factory;
  CongressRepository von Dapper auf EF; CongressTrade.DetailsFetched ergaenzt; Dapper entfernt
- Core: CoreSettingsService (core_settings via EF); WorkerBase-Log auf EF (core_worker_log);
  7 Worker-Ctors DatabaseService -> IDbContextFactory<CoreDbContext>
- CoreMigrations + CongressMigrations (Dapper) entfernt; core_-Schema nun rein EF
- EF-Migration InitialCongressTrading; AddCorePersistence-Fallback fuer leeren Connection-String
  (App startet ohne DB); Connection aus appsettings.Local.json (gitignored)
- Tests: +4 CongressRepository (EF-InMemory) -> 39/39 gruen; Build + smoke-ui + App-Start ok

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
@
This commit is contained in:
Richard
2026-07-28 11:34:42 +02:00
parent 7e3791b82d
commit 6c8e36dc3c
27 changed files with 832 additions and 323 deletions
+2 -2
View File
@@ -131,15 +131,15 @@ public sealed class LauncherForm : Form
if (Enum.TryParse<AppLogLevel>(levelStr, true, out var level)) if (Enum.TryParse<AppLogLevel>(levelStr, true, out var level))
_logger.SetMinLevel(level); _logger.SetMinLevel(level);
// core_-Schema läuft über EF-Migrationen (extern angewendet). Nur ibkr_ (Dapper) noch hier.
SetStatus("Migrationen..."); SetStatus("Migrationen...");
try try
{ {
await _services.GetRequiredService<CoreMigrations>().RunAsync();
await _services.GetRequiredService<IBKRMigrations>().RunAsync(); await _services.GetRequiredService<IBKRMigrations>().RunAsync();
} }
catch (Exception ex) catch (Exception ex)
{ {
_logger.Error("Core", "Core-Datenbankfehler beim Start.", ex); _logger.Error("Core", "IBKR-Datenbankfehler beim Start.", ex);
} }
foreach (var module in _modules) foreach (var module in _modules)
+3 -3
View File
@@ -94,9 +94,9 @@ internal static class Program
services.AddSingleton<LoggingService>(); services.AddSingleton<LoggingService>();
services.AddSingleton<DatabaseService>(); services.AddSingleton<DatabaseService>(); // noch für IBKR-Marktdaten (ibkr_, Dapper)
services.AddSingleton<CoreMigrations>(); services.AddSingleton<IBKRMigrations>(); // ibkr_-Tabellen (Dapper); core_ läuft über EF
services.AddSingleton<IBKRMigrations>(); services.AddSingleton<CoreSettingsService>(); // core_settings via EF
services.AddSingleton<IBKRGatewayService>(); services.AddSingleton<IBKRGatewayService>();
services.AddSingleton<IBKRMarketDataRepository>(); services.AddSingleton<IBKRMarketDataRepository>();
+2 -2
View File
@@ -94,8 +94,8 @@ Pin `new MariaDbServerVersion(new Version(11, 8, 6))`. Verbindung aus `appsettin
- [x] `--db-version`-Diagnose (Serverversion für den EF-Pin) - [x] `--db-version`-Diagnose (Serverversion für den EF-Pin)
- [x] **Slice 1:** `Configuration/DatabaseOptions` + `DatabaseServerVersion`-Pin (MariaDB 11.8.6); `AddCorePersistence` (`AddDbContextFactory`) + `CoreDbContext` + Entities (core_position, core_trade_history, core_budget, core_worker_log, core_settings) + Design-Time-Factory; **EF-Migration `InitialCore` erzeugt** - [x] **Slice 1:** `Configuration/DatabaseOptions` + `DatabaseServerVersion`-Pin (MariaDB 11.8.6); `AddCorePersistence` (`AddDbContextFactory`) + `CoreDbContext` + Entities (core_position, core_trade_history, core_budget, core_worker_log, core_settings) + Design-Time-Factory; **EF-Migration `InitialCore` erzeugt**
- [x] **Slice 2:** `PortfolioService`, `BudgetService`, `TradeHistoryService` auf EF (`IDbContextFactory<CoreDbContext>`) umgestellt; **5 EF-InMemory-Unit-Tests** (Buchführung real verifiziert). WorkerBase-Log folgt in Slice 4. - [x] **Slice 2:** `PortfolioService`, `BudgetService`, `TradeHistoryService` auf EF (`IDbContextFactory<CoreDbContext>`) umgestellt; **5 EF-InMemory-Unit-Tests** (Buchführung real verifiziert). WorkerBase-Log folgt in Slice 4.
- [ ] **Slice 3:** Modul-DbContext (ct_) im CongressTrading-Projekt + Umstellung - [x] **Slice 3:** Modul auf EF (`CongressTradingDbContext` ct_ + `CongressRepository`); `WorkerBase`-Log auf EF (core_worker_log); `CoreSettingsService` (core_settings); `CoreMigrations`/`CongressMigrations` (Dapper) entfernt; EF-Migration `InitialCongressTrading`; `appsettings.Local.json` (gitignored) als Connection-Quelle; Fallback für leeren Connection-String. **+4 EF-InMemory-Tests** (CongressRepository)
- [ ] **Slice 4:** Dapper + `DatabaseService` + manuelle Migrationen (CoreMigrations/IBKRMigrations/CongressMigrations) entfernen; IBKR-Marktdaten auf EF - [ ] **Slice 4:** IBKR-Marktdaten (`ibkr_`) + `IBKRMigrations` auf EF; dann restliches Dapper + `DatabaseService` entfernen
- Hinweis: nur build-verifizierbar (Unit-Tests ohne DB); Schema-Anwendung extern via `dotnet ef database update` (env `IBKRTRADER_MYSQL`) - Hinweis: nur build-verifizierbar (Unit-Tests ohne DB); Schema-Anwendung extern via `dotnet ef database update` (env `IBKRTRADER_MYSQL`)
### R4 Trading-Kern einфügen ### R4 Trading-Kern einфügen
@@ -1,91 +0,0 @@
using IBKRTrader.Core.Logging;
namespace IBKRTrader.Core.Database.Migrations;
/// <summary>
/// Erstellt alle core_xxx-Tabellen idempotent (IF NOT EXISTS).
/// Wird einmalig beim App-Start ausgeführt.
/// </summary>
public class CoreMigrations
{
private readonly DatabaseService _db;
private readonly LoggingService _logger;
public CoreMigrations(DatabaseService db, LoggingService logger)
{
_db = db;
_logger = logger;
}
public async Task RunAsync()
{
_logger.Info("Core", "Starte Core-Datenbankmigrationen...");
await CreateCoreSettingsAsync();
await CreateCoreWorkerLogAsync();
await CreateCoreTradeHistoryAsync();
await CreateCoreBudgetAsync();
await CreateCorePositionAsync();
_logger.Info("Core", "Core-Migrationen abgeschlossen.");
}
// ─── Tabellen ─────────────────────────────────────────────────────────────
private Task CreateCoreSettingsAsync() => _db.ExecuteAsync(@"
CREATE TABLE IF NOT EXISTS `core_settings` (
`key` VARCHAR(100) NOT NULL PRIMARY KEY,
`value` TEXT,
`updated_at` DATETIME DEFAULT CURRENT_TIMESTAMP
ON UPDATE CURRENT_TIMESTAMP
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;");
private Task CreateCoreWorkerLogAsync() => _db.ExecuteAsync(@"
CREATE TABLE IF NOT EXISTS `core_worker_log` (
`id` BIGINT NOT NULL AUTO_INCREMENT PRIMARY KEY,
`worker_name` VARCHAR(100) NOT NULL,
`module` VARCHAR(50) NOT NULL DEFAULT 'Core',
`started_at` DATETIME,
`finished_at` DATETIME,
`status` ENUM('Running','Success','Error') DEFAULT 'Running',
`message` TEXT,
INDEX `idx_worker` (`worker_name`, `started_at`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;");
private Task CreateCoreTradeHistoryAsync() => _db.ExecuteAsync(@"
CREATE TABLE IF NOT EXISTS `core_trade_history` (
`id` BIGINT NOT NULL AUTO_INCREMENT PRIMARY KEY,
`module` VARCHAR(50) NOT NULL,
`symbol` VARCHAR(20) NOT NULL,
`action` ENUM('BUY','SELL') NOT NULL,
`quantity` DECIMAL(18,4),
`price` DECIMAL(18,4),
`total_value` DECIMAL(18,4),
`traded_at` DATETIME,
`ibkr_order_id` VARCHAR(100),
`status` VARCHAR(50),
`notes` TEXT,
`created_at` DATETIME DEFAULT CURRENT_TIMESTAMP,
INDEX `idx_symbol` (`symbol`),
INDEX `idx_module` (`module`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;");
private Task CreateCoreBudgetAsync() => _db.ExecuteAsync(@"
CREATE TABLE IF NOT EXISTS `core_budget` (
`module` VARCHAR(50) NOT NULL PRIMARY KEY,
`total_budget` DECIMAL(18,2) DEFAULT 0,
`used_budget` DECIMAL(18,2) DEFAULT 0,
`max_per_trade` DECIMAL(18,2) DEFAULT 0,
`updated_at` DATETIME DEFAULT CURRENT_TIMESTAMP
ON UPDATE CURRENT_TIMESTAMP
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;");
private Task CreateCorePositionAsync() => _db.ExecuteAsync(@"
CREATE TABLE IF NOT EXISTS `core_position` (
`module` VARCHAR(50) NOT NULL,
`symbol` VARCHAR(20) NOT NULL,
`quantity` INT NOT NULL DEFAULT 0,
`avg_price` DECIMAL(18,4) NOT NULL DEFAULT 0,
`updated_at` DATETIME DEFAULT CURRENT_TIMESTAMP
ON UPDATE CURRENT_TIMESTAMP,
PRIMARY KEY (`module`, `symbol`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;");
}
@@ -14,8 +14,18 @@ public static class ServiceCollectionExtensions
public static IServiceCollection AddCorePersistence(this IServiceCollection services, DatabaseOptions options) public static IServiceCollection AddCorePersistence(this IServiceCollection services, DatabaseOptions options)
{ {
services.AddDbContextFactory<CoreDbContext>(o => services.AddDbContextFactory<CoreDbContext>(o =>
o.UseMySql(options.MySqlConnectionString, DatabaseServerVersion.Value)); o.UseMySql(EffectiveConnectionString(options.MySqlConnectionString), DatabaseServerVersion.Value));
return services; return services;
} }
/// <summary>
/// Liefert den Connection-String oder bei leerem Wert (z. B. fehlende appsettings.Local.json)
/// einen unschädlichen Platzhalter, damit der DbContext-Optionsbau nicht wirft. Die App startet so
/// unabhängig von der DB; echte Queries schlagen erst bei tatsächlichem Zugriff fehl.
/// </summary>
public static string EffectiveConnectionString(string? connectionString)
=> string.IsNullOrWhiteSpace(connectionString)
? "Server=localhost;Port=3306;Database=ibkrtrader;User ID=root;Password=;"
: connectionString;
} }
@@ -0,0 +1,48 @@
using IBKRTrader.Core.Persistence.Ef;
using IBKRTrader.Core.Persistence.Entities;
using Microsoft.EntityFrameworkCore;
namespace IBKRTrader.Core.Settings;
/// <summary>
/// Key-Value-Zugriff auf core_settings (EF Core). Von Core und Modulen für Fortschritts-/Flag-Werte
/// genutzt (z. B. "ct.history_import_page").
/// </summary>
public class CoreSettingsService
{
private readonly IDbContextFactory<CoreDbContext> _dbf;
public CoreSettingsService(IDbContextFactory<CoreDbContext> dbf) => _dbf = dbf;
public async Task<string?> GetAsync(string key)
{
await using var db = await _dbf.CreateDbContextAsync();
var s = await db.Settings.FindAsync(key);
return s?.Value;
}
public async Task SetAsync(string key, string value)
{
await using var db = await _dbf.CreateDbContextAsync();
var s = await db.Settings.FindAsync(key);
if (s is null)
db.Settings.Add(new CoreSetting { Key = key, Value = value, UpdatedAt = DateTime.UtcNow });
else
{
s.Value = value;
s.UpdatedAt = DateTime.UtcNow;
}
await db.SaveChangesAsync();
}
public async Task DeleteAsync(string key)
{
await using var db = await _dbf.CreateDbContextAsync();
var s = await db.Settings.FindAsync(key);
if (s is not null)
{
db.Settings.Remove(s);
await db.SaveChangesAsync();
}
}
}
@@ -1,3 +1,5 @@
using IBKRTrader.Core.Persistence.Ef;
using Microsoft.EntityFrameworkCore;
using IBKRTrader.Core.Database; using IBKRTrader.Core.Database;
using IBKRTrader.Core.Logging; using IBKRTrader.Core.Logging;
using IBKRTrader.Core.Settings; using IBKRTrader.Core.Settings;
@@ -21,7 +23,7 @@ public class BackupWorker : WorkerBase
protected override TimeSpan? Interval => protected override TimeSpan? Interval =>
TimeSpan.FromMinutes(_settings.Settings.WorkerSettings.BackupWorker.IntervalMinutes); TimeSpan.FromMinutes(_settings.Settings.WorkerSettings.BackupWorker.IntervalMinutes);
public BackupWorker(LoggingService logger, DatabaseService db, SettingsService settings) public BackupWorker(LoggingService logger, IDbContextFactory<CoreDbContext> db, SettingsService settings)
: base(logger, db) : base(logger, db)
{ {
_settings = settings; _settings = settings;
@@ -1,3 +1,5 @@
using IBKRTrader.Core.Persistence.Ef;
using Microsoft.EntityFrameworkCore;
using IBKRTrader.Core.Database; using IBKRTrader.Core.Database;
using IBKRTrader.Core.IBKR; using IBKRTrader.Core.IBKR;
using IBKRTrader.Core.Logging; using IBKRTrader.Core.Logging;
@@ -22,7 +24,7 @@ public class IBKRInstrumentSyncWorker : WorkerBase
private readonly IBKRMarketDataRepository _repo; private readonly IBKRMarketDataRepository _repo;
public IBKRInstrumentSyncWorker( public IBKRInstrumentSyncWorker(
LoggingService logger, DatabaseService db, LoggingService logger, IDbContextFactory<CoreDbContext> db,
SettingsService settings, IBKRGatewayService gateway, SettingsService settings, IBKRGatewayService gateway,
IBKRMarketDataRepository repo) IBKRMarketDataRepository repo)
: base(logger, db) : base(logger, db)
@@ -1,3 +1,5 @@
using IBKRTrader.Core.Persistence.Ef;
using Microsoft.EntityFrameworkCore;
using IBKRTrader.Core.Database; using IBKRTrader.Core.Database;
using IBKRTrader.Core.IBKR; using IBKRTrader.Core.IBKR;
using IBKRTrader.Core.Logging; using IBKRTrader.Core.Logging;
@@ -29,7 +31,7 @@ public class IBKRPriceHistoryWorker : WorkerBase
private readonly IBKRMarketDataRepository _repo; private readonly IBKRMarketDataRepository _repo;
public IBKRPriceHistoryWorker( public IBKRPriceHistoryWorker(
LoggingService logger, DatabaseService db, LoggingService logger, IDbContextFactory<CoreDbContext> db,
SettingsService settings, IBKRGatewayService gateway, SettingsService settings, IBKRGatewayService gateway,
IBKRMarketDataRepository repo) IBKRMarketDataRepository repo)
: base(logger, db) : base(logger, db)
@@ -1,3 +1,5 @@
using IBKRTrader.Core.Persistence.Ef;
using Microsoft.EntityFrameworkCore;
using System.Net; using System.Net;
using System.Text; using System.Text;
using System.Text.Json; using System.Text.Json;
@@ -27,7 +29,7 @@ public class WebApiService : WorkerBase
protected override TimeSpan? Interval => null; protected override TimeSpan? Interval => null;
public WebApiService(LoggingService logger, DatabaseService db, SettingsService settings) public WebApiService(LoggingService logger, IDbContextFactory<CoreDbContext> db, SettingsService settings)
: base(logger, db) : base(logger, db)
{ {
_settings = settings; _settings = settings;
@@ -1,3 +1,5 @@
using IBKRTrader.Core.Persistence.Ef;
using Microsoft.EntityFrameworkCore;
using System.Net; using System.Net;
using System.Text; using System.Text;
using IBKRTrader.Core.Database; using IBKRTrader.Core.Database;
@@ -22,7 +24,7 @@ public class WebserverService : WorkerBase
protected override TimeSpan? Interval => null; // Service = permanent protected override TimeSpan? Interval => null; // Service = permanent
public WebserverService(LoggingService logger, DatabaseService db, SettingsService settings) public WebserverService(LoggingService logger, IDbContextFactory<CoreDbContext> db, SettingsService settings)
: base(logger, db) : base(logger, db)
{ {
_settings = settings; _settings = settings;
+32 -10
View File
@@ -1,5 +1,7 @@
using IBKRTrader.Core.Database;
using IBKRTrader.Core.Logging; using IBKRTrader.Core.Logging;
using IBKRTrader.Core.Persistence.Ef;
using IBKRTrader.Core.Persistence.Entities;
using Microsoft.EntityFrameworkCore;
namespace IBKRTrader.Core.Workers; namespace IBKRTrader.Core.Workers;
@@ -26,8 +28,8 @@ public abstract class WorkerBase : IWorker
// ─── Infrastruktur ──────────────────────────────────────────────────────── // ─── Infrastruktur ────────────────────────────────────────────────────────
protected readonly LoggingService Logger; protected readonly LoggingService Logger;
protected readonly DatabaseService Db; private readonly IDbContextFactory<CoreDbContext> _dbf;
public WorkerInfo Info { get; } = new(); public WorkerInfo Info { get; } = new();
@@ -35,10 +37,10 @@ public abstract class WorkerBase : IWorker
private Task? _runLoop; private Task? _runLoop;
private readonly SemaphoreSlim _triggerSemaphore = new(0, 1); private readonly SemaphoreSlim _triggerSemaphore = new(0, 1);
protected WorkerBase(LoggingService logger, DatabaseService db) protected WorkerBase(LoggingService logger, IDbContextFactory<CoreDbContext> dbf)
{ {
Logger = logger; Logger = logger;
Db = db; _dbf = dbf;
Info.WorkerName = Name; Info.WorkerName = Name;
Info.Module = Module; Info.Module = Module;
@@ -153,13 +155,33 @@ public abstract class WorkerBase : IWorker
// ─── DB-Log-Seam (überschreibbar für Unit-Tests) ────────────────────────── // ─── DB-Log-Seam (überschreibbar für Unit-Tests) ──────────────────────────
/// <summary>Legt den Worker-Log-Eintrag an. Kapselt den DB-Zugriff (testbar).</summary> /// <summary>Legt den Worker-Log-Eintrag an (core_worker_log via EF). Kapselt den DB-Zugriff (testbar).</summary>
protected virtual Task<long> BeginRunLogAsync() protected virtual async Task<long> BeginRunLogAsync()
=> Db.BeginWorkerLogAsync(Name, Module); {
await using var db = await _dbf.CreateDbContextAsync();
var log = new CoreWorkerLog
{
WorkerName = Name,
Module = Module,
StartedAt = DateTime.UtcNow,
Status = "Running"
};
db.WorkerLog.Add(log);
await db.SaveChangesAsync();
return log.Id;
}
/// <summary>Schließt den Worker-Log-Eintrag ab. Kapselt den DB-Zugriff (testbar).</summary> /// <summary>Schließt den Worker-Log-Eintrag ab. Kapselt den DB-Zugriff (testbar).</summary>
protected virtual Task EndRunLogAsync(long logId, bool success, string? message = null) protected virtual async Task EndRunLogAsync(long logId, bool success, string? message = null)
=> Db.EndWorkerLogAsync(logId, success, message); {
await using var db = await _dbf.CreateDbContextAsync();
var log = await db.WorkerLog.FindAsync(logId);
if (log is null) return;
log.FinishedAt = DateTime.UtcNow;
log.Status = success ? "Success" : "Error";
log.Message = message;
await db.SaveChangesAsync();
}
// ─── Hilfsmethoden ──────────────────────────────────────────────────────── // ─── Hilfsmethoden ────────────────────────────────────────────────────────
@@ -1,10 +1,14 @@
using IBKRTrader.Core.Configuration;
using IBKRTrader.Core.DependencyInjection;
using IBKRTrader.Core.Logging; using IBKRTrader.Core.Logging;
using IBKRTrader.Core.Modularity; using IBKRTrader.Core.Modularity;
using IBKRTrader.Core.Workers; using IBKRTrader.Core.Workers;
using IBKRTrader.Modules.CongressTrading.Database; using IBKRTrader.Modules.CongressTrading.Database;
using IBKRTrader.Modules.CongressTrading.Persistence.Ef;
using IBKRTrader.Modules.CongressTrading.Scraper; using IBKRTrader.Modules.CongressTrading.Scraper;
using IBKRTrader.Modules.CongressTrading.UI; using IBKRTrader.Modules.CongressTrading.UI;
using IBKRTrader.Modules.CongressTrading.Workers; using IBKRTrader.Modules.CongressTrading.Workers;
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.Configuration; using Microsoft.Extensions.Configuration;
using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.DependencyInjection;
@@ -22,12 +26,12 @@ public sealed class CongressTradingModule : IModule
public string Name => "CongressTrading"; public string Name => "CongressTrading";
public string DbPrefix => "ct_"; public string DbPrefix => "ct_";
// In RegisterUi gesetzt (hat den Provider); für Modul-Migrationen in StartAsync genutzt.
private IServiceProvider? _services;
public void RegisterServices(IServiceCollection services, IConfiguration configuration) public void RegisterServices(IServiceCollection services, IConfiguration configuration)
{ {
services.AddSingleton<CongressMigrations>(); // Modul-Persistenz: eigener DbContext (ct_-Tabellen) in derselben MariaDB.
var conn = ServiceCollectionExtensions.EffectiveConnectionString(configuration["Database:MySqlConnectionString"]);
services.AddDbContextFactory<CongressTradingDbContext>(o => o.UseMySql(conn, DatabaseServerVersion.Value));
services.AddSingleton<CongressRepository>(); services.AddSingleton<CongressRepository>();
services.AddSingleton<CapitolTradesScraper>(); services.AddSingleton<CapitolTradesScraper>();
@@ -41,8 +45,6 @@ public sealed class CongressTradingModule : IModule
public void RegisterUi(IModuleUiHost host, IServiceProvider services) public void RegisterUi(IModuleUiHost host, IServiceProvider services)
{ {
_services = services;
host.RegisterView(new ModuleView host.RegisterView(new ModuleView
{ {
Id = "congresstrading.main", Id = "congresstrading.main",
@@ -56,13 +58,8 @@ public sealed class CongressTradingModule : IModule
}); });
} }
public async Task StartAsync(CancellationToken cancellationToken) // DB-Schema wird extern per `dotnet ef database update` angewendet (keine Laufzeit-Migration).
{ public Task StartAsync(CancellationToken cancellationToken) => Task.CompletedTask;
// Modul-Migrationen laufen nach der Core-Initialisierung.
if (_services is null) return;
var migrations = _services.GetRequiredService<CongressMigrations>();
await migrations.RunAsync();
}
public Task StopAsync(CancellationToken cancellationToken) => Task.CompletedTask; public Task StopAsync(CancellationToken cancellationToken) => Task.CompletedTask;
} }
@@ -1,80 +0,0 @@
using IBKRTrader.Core.Database;
using IBKRTrader.Core.Logging;
namespace IBKRTrader.Modules.CongressTrading.Database;
/// <summary>
/// Legt die ct_xxx-Tabellen idempotent an.
/// Namensschema: ct_{tabellenname} gemäß Architekturregeln.
/// </summary>
public class CongressMigrations
{
private readonly DatabaseService _db;
private readonly LoggingService _logger;
public CongressMigrations(DatabaseService db, LoggingService logger)
{
_db = db;
_logger = logger;
}
public async Task RunAsync()
{
_logger.Info("CT", "Starte CongressTrading-Datenbankmigrationen...");
await CreateCongressMemberTableAsync();
await CreateCongressTradeTableAsync();
_logger.Info("CT", "CongressTrading-Migrationen abgeschlossen.");
}
private Task CreateCongressMemberTableAsync() => _db.ExecuteAsync(@"
CREATE TABLE IF NOT EXISTS `ct_congressMember` (
`id` BIGINT NOT NULL AUTO_INCREMENT PRIMARY KEY,
`bio_id` VARCHAR(20) NOT NULL UNIQUE,
`name` VARCHAR(200),
`party` VARCHAR(50),
`state` VARCHAR(50),
`chamber` VARCHAR(10),
`profile_url` VARCHAR(500),
`first_seen` DATETIME DEFAULT CURRENT_TIMESTAMP,
`last_updated` DATETIME DEFAULT CURRENT_TIMESTAMP
ON UPDATE CURRENT_TIMESTAMP,
INDEX `idx_bio_id` (`bio_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;");
private async Task CreateCongressTradeTableAsync()
{
// Tabelle anlegen
await _db.ExecuteAsync(@"
CREATE TABLE IF NOT EXISTS `ct_trade` (
`id` BIGINT NOT NULL AUTO_INCREMENT PRIMARY KEY,
`trade_id` VARCHAR(20) NOT NULL UNIQUE,
`member_bio_id` VARCHAR(20) NOT NULL,
`issuer_name` VARCHAR(500),
`issuer_id` VARCHAR(20),
`ticker` VARCHAR(20),
`trade_type` VARCHAR(50),
`chamber` VARCHAR(20),
`owner` VARCHAR(50),
`value` DECIMAL(18,2),
`size_range_low` BIGINT,
`size_range_high` BIGINT,
`trade_date` DATE,
`published_date` DATE,
`detail_url` VARCHAR(500),
`details_fetched` TINYINT(1) NOT NULL DEFAULT 1,
`scraped_at` DATETIME DEFAULT CURRENT_TIMESTAMP,
INDEX `idx_member` (`member_bio_id`),
INDEX `idx_trade_date` (`trade_date`),
INDEX `idx_ticker` (`ticker`),
INDEX `idx_details` (`details_fetched`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;");
// Idempotente Spalten-Migrationen (für bestehende DBs)
await _db.ExecuteAsync("ALTER TABLE `ct_trade` ADD COLUMN IF NOT EXISTS `details_fetched` TINYINT(1) NOT NULL DEFAULT 1;");
await _db.ExecuteAsync("ALTER TABLE `ct_trade` ADD COLUMN IF NOT EXISTS `chamber` VARCHAR(20);");
await _db.ExecuteAsync("ALTER TABLE `ct_trade` ADD COLUMN IF NOT EXISTS `owner` VARCHAR(50);");
await _db.ExecuteAsync("ALTER TABLE `ct_trade` ADD COLUMN IF NOT EXISTS `value` DECIMAL(18,2);");
await _db.ExecuteAsync("ALTER TABLE `ct_trade` ADD INDEX IF NOT EXISTS `idx_details` (`details_fetched`);");
}
}
@@ -1,53 +1,65 @@
using Dapper;
using IBKRTrader.Core.Database;
using IBKRTrader.Core.Logging; using IBKRTrader.Core.Logging;
using IBKRTrader.Core.Settings;
using IBKRTrader.Modules.CongressTrading.Models; using IBKRTrader.Modules.CongressTrading.Models;
using IBKRTrader.Modules.CongressTrading.Persistence.Ef;
using Microsoft.EntityFrameworkCore;
namespace IBKRTrader.Modules.CongressTrading.Database; namespace IBKRTrader.Modules.CongressTrading.Database;
/// <summary> /// <summary>
/// Datenbankzugriff für das CongressTrading-Modul. /// EF-gestützter Datenbankzugriff für das CongressTrading-Modul (ct_congressMember, ct_trade).
/// Kapselt alle INSERT/SELECT auf ct_congressMember und ct_trade. /// Fortschritts-Flags liegen in core_settings (über <see cref="CoreSettingsService"/>).
/// </summary> /// </summary>
public class CongressRepository public class CongressRepository
{ {
private readonly DatabaseService _db; private readonly IDbContextFactory<CongressTradingDbContext> _dbf;
private readonly LoggingService _logger; private readonly CoreSettingsService _settings;
private readonly LoggingService _logger;
public CongressRepository(DatabaseService db, LoggingService logger) public CongressRepository(
IDbContextFactory<CongressTradingDbContext> dbf,
CoreSettingsService settings,
LoggingService logger)
{ {
_db = db; _dbf = dbf;
_logger = logger; _settings = settings;
_logger = logger;
} }
// ─── CongressMember ─────────────────────────────────────────────────────── // ─── CongressMember ───────────────────────────────────────────────────────
public async Task<bool> MemberExistsAsync(string bioId) public async Task<bool> MemberExistsAsync(string bioId)
{ {
var count = await _db.ExecuteScalarAsync<int>( await using var db = await _dbf.CreateDbContextAsync();
"SELECT COUNT(*) FROM `ct_congressMember` WHERE bio_id = @bioId", return await db.Members.AnyAsync(m => m.BioId == bioId);
new { bioId });
return count > 0;
} }
public async Task UpsertMemberAsync(CongressMember member) public async Task UpsertMemberAsync(CongressMember member)
{ {
await _db.ExecuteAsync(@" await using var db = await _dbf.CreateDbContextAsync();
INSERT INTO `ct_congressMember` (bio_id, name, party, state, chamber, profile_url) var existing = await db.Members.FirstOrDefaultAsync(m => m.BioId == member.BioId);
VALUES (@BioId, @Name, @Party, @State, @Chamber, @ProfileUrl) if (existing is null)
ON DUPLICATE KEY UPDATE {
name = VALUES(name), member.FirstSeen = DateTime.UtcNow;
party = VALUES(party), member.LastUpdated = DateTime.UtcNow;
state = VALUES(state), db.Members.Add(member);
chamber = VALUES(chamber), }
profile_url = VALUES(profile_url), else
last_updated = CURRENT_TIMESTAMP", member); {
existing.Name = member.Name;
existing.Party = member.Party;
existing.State = member.State;
existing.Chamber = member.Chamber;
existing.ProfileUrl = member.ProfileUrl;
existing.LastUpdated = DateTime.UtcNow;
}
await db.SaveChangesAsync();
} }
public async Task<HashSet<string>> GetAllMemberBioIdsAsync() public async Task<HashSet<string>> GetAllMemberBioIdsAsync()
{ {
var ids = await _db.QueryAsync<string>( await using var db = await _dbf.CreateDbContextAsync();
"SELECT bio_id FROM `ct_congressMember`"); var ids = await db.Members.Select(m => m.BioId).ToListAsync();
return ids.ToHashSet(); return ids.ToHashSet();
} }
@@ -55,120 +67,85 @@ public class CongressRepository
public async Task<bool> TradeExistsAsync(string tradeId) public async Task<bool> TradeExistsAsync(string tradeId)
{ {
var count = await _db.ExecuteScalarAsync<int>( await using var db = await _dbf.CreateDbContextAsync();
"SELECT COUNT(*) FROM `ct_trade` WHERE trade_id = @tradeId", return await db.Trades.AnyAsync(t => t.TradeId == tradeId);
new { tradeId });
return count > 0;
} }
public async Task<HashSet<string>> GetAllTradeIdsAsync() public async Task<HashSet<string>> GetAllTradeIdsAsync()
{ {
var ids = await _db.QueryAsync<string>( await using var db = await _dbf.CreateDbContextAsync();
"SELECT trade_id FROM `ct_trade`"); var ids = await db.Trades.Select(t => t.TradeId).ToListAsync();
return ids.ToHashSet(); return ids.ToHashSet();
} }
/// <summary>Fügt einen Trade ein. Ignoriert Duplikate (INSERT IGNORE).</summary> /// <summary>Fügt einen Trade ein, sofern die TradeId noch nicht existiert.</summary>
public async Task InsertTradeAsync(CongressTrade trade, bool detailsFetched = true) public async Task InsertTradeAsync(CongressTrade trade, bool detailsFetched = true)
{ {
await _db.ExecuteAsync(@" await using var db = await _dbf.CreateDbContextAsync();
INSERT IGNORE INTO `ct_trade` if (await db.Trades.AnyAsync(t => t.TradeId == trade.TradeId)) return;
(trade_id, member_bio_id, issuer_name, issuer_id, ticker,
trade_type, chamber, owner, value, trade.DetailsFetched = detailsFetched;
trade_date, published_date, detail_url, details_fetched) trade.ScrapedAt = DateTime.UtcNow;
VALUES trade.MemberSnapshot = null; // transient
(@TradeId, @MemberBioId, @IssuerName, @IssuerId, @Ticker, db.Trades.Add(trade);
@TradeType, @Chamber, @Owner, @Value, await db.SaveChangesAsync();
@TradeDate, @PublishedDate, @DetailUrl, @DetailsFetched)",
new
{
trade.TradeId,
trade.MemberBioId,
trade.IssuerName,
trade.IssuerId,
trade.Ticker,
trade.TradeType,
trade.Chamber,
trade.Owner,
trade.Value,
TradeDate = trade.TradeDate?.ToString("yyyy-MM-dd"),
PublishedDate = trade.PublishedDate?.ToString("yyyy-MM-dd"),
trade.DetailUrl,
DetailsFetched = detailsFetched ? 1 : 0
});
} }
/// <summary>Aktualisiert die Detail-Felder eines bereits vorhandenen Trades.</summary> /// <summary>Aktualisiert die Detail-Felder eines noch unvollständigen Trades.</summary>
public async Task UpdateTradeDetailsAsync(CongressTrade trade) public async Task UpdateTradeDetailsAsync(CongressTrade trade)
{ {
await _db.ExecuteAsync(@" await using var db = await _dbf.CreateDbContextAsync();
UPDATE `ct_trade` var existing = await db.Trades.FirstOrDefaultAsync(t => t.TradeId == trade.TradeId && !t.DetailsFetched);
SET ticker = @Ticker, if (existing is null) return;
trade_type = @TradeType,
trade_date = @TradeDate, existing.Ticker = trade.Ticker;
published_date = @PublishedDate, existing.TradeType = trade.TradeType;
details_fetched = 1 existing.TradeDate = trade.TradeDate;
WHERE trade_id = @TradeId existing.PublishedDate = trade.PublishedDate;
AND details_fetched = 0", existing.DetailsFetched = true;
new await db.SaveChangesAsync();
{
trade.TradeId,
trade.Ticker,
trade.TradeType,
TradeDate = trade.TradeDate?.ToString("yyyy-MM-dd"),
PublishedDate = trade.PublishedDate?.ToString("yyyy-MM-dd")
});
} }
/// <summary>Gibt Trade-IDs zurück für die noch keine Details geholt wurden.</summary>
public async Task<IEnumerable<string>> GetTradesWithoutDetailsAsync(int limit = 500) public async Task<IEnumerable<string>> GetTradesWithoutDetailsAsync(int limit = 500)
=> await _db.QueryAsync<string>( {
"SELECT trade_id FROM `ct_trade` WHERE details_fetched = 0 LIMIT @limit", await using var db = await _dbf.CreateDbContextAsync();
new { limit }); return await db.Trades
.Where(t => !t.DetailsFetched)
.Select(t => t.TradeId)
.Take(limit)
.ToListAsync();
}
public async Task<int> GetTradeCountAsync() public async Task<int> GetTradeCountAsync()
=> await _db.ExecuteScalarAsync<int>("SELECT COUNT(*) FROM `ct_trade`"); {
await using var db = await _dbf.CreateDbContextAsync();
return await db.Trades.CountAsync();
}
public async Task<int> GetMemberCountAsync() public async Task<int> GetMemberCountAsync()
=> await _db.ExecuteScalarAsync<int>("SELECT COUNT(*) FROM `ct_congressMember`");
// ─── core_settings Flags ──────────────────────────────────────────────────
public async Task<string?> GetSettingAsync(string key)
{ {
return await _db.QueryFirstOrDefaultAsync<string>( await using var db = await _dbf.CreateDbContextAsync();
"SELECT `value` FROM `core_settings` WHERE `key` = @key", return await db.Members.CountAsync();
new { key });
} }
public async Task SetSettingAsync(string key, string value) // ─── Fortschritts-Flags (core_settings via Core) ──────────────────────────
{
await _db.ExecuteAsync(@"
INSERT INTO `core_settings` (`key`, `value`)
VALUES (@key, @value)
ON DUPLICATE KEY UPDATE `value` = @value",
new { key, value });
}
public async Task DeleteSettingAsync(string key) public Task<string?> GetSettingAsync(string key) => _settings.GetAsync(key);
{ public Task SetSettingAsync(string key, string v) => _settings.SetAsync(key, v);
await _db.ExecuteAsync( public Task DeleteSettingAsync(string key) => _settings.DeleteAsync(key);
"DELETE FROM `core_settings` WHERE `key` = @key", new { key });
}
// ─── Reset ──────────────────────────────────────────────────────────────── // ─── Reset ────────────────────────────────────────────────────────────────
/// <summary> /// <summary>Löscht alle CT-Fortschritts-Flags und alle CT-Daten aus der DB.</summary>
/// Löscht alle CT-Fortschritts-Flags und alle CT-Daten aus der DB.
/// Danach startet der History-Import beim nächsten Worker-Run von Seite 1.
/// </summary>
public async Task ResetHistoryImportAsync() public async Task ResetHistoryImportAsync()
{ {
await DeleteSettingAsync("ct.history_import_done"); await DeleteSettingAsync("ct.history_import_done");
await DeleteSettingAsync("ct.history_import_page"); await DeleteSettingAsync("ct.history_import_page");
await DeleteSettingAsync("ct.history_import_total"); await DeleteSettingAsync("ct.history_import_total");
await _db.ExecuteAsync("DELETE FROM `ct_trade`");
await _db.ExecuteAsync("DELETE FROM `ct_congressMember`"); await using var db = await _dbf.CreateDbContextAsync();
await db.Trades.ExecuteDeleteAsync();
await db.Members.ExecuteDeleteAsync();
_logger.Warn("CT", "History-Import-Reset durchgeführt alle CT-Daten gelöscht."); _logger.Warn("CT", "History-Import-Reset durchgeführt alle CT-Daten gelöscht.");
} }
} }
@@ -10,7 +10,11 @@
</PropertyGroup> </PropertyGroup>
<ItemGroup> <ItemGroup>
<PackageReference Include="Dapper" Version="2.1.35" /> <PackageReference Include="Pomelo.EntityFrameworkCore.MySql" Version="8.0.3" />
<PackageReference Include="Microsoft.EntityFrameworkCore.Design" Version="8.0.11">
<PrivateAssets>all</PrivateAssets>
<IncludeAssets>runtime; build; native; contentfiles; analyzers; buildtransitive</IncludeAssets>
</PackageReference>
</ItemGroup> </ItemGroup>
<ItemGroup> <ItemGroup>
@@ -16,8 +16,9 @@ public class CongressTrade
public DateOnly? TradeDate { get; set; } public DateOnly? TradeDate { get; set; }
public DateOnly? PublishedDate { get; set; } public DateOnly? PublishedDate { get; set; }
public string DetailUrl { get; set; } = ""; public string DetailUrl { get; set; } = "";
public bool DetailsFetched { get; set; } = true;
public DateTime ScrapedAt { get; set; } = DateTime.UtcNow; public DateTime ScrapedAt { get; set; } = DateTime.UtcNow;
/// <summary>Transient wird direkt in ct_congressMember gespeichert.</summary> /// <summary>Transient wird direkt in ct_congressMember gespeichert (nicht persistiert).</summary>
public CongressMember? MemberSnapshot { get; set; } public CongressMember? MemberSnapshot { get; set; }
} }
@@ -0,0 +1,52 @@
using IBKRTrader.Modules.CongressTrading.Models;
using Microsoft.EntityFrameworkCore;
namespace IBKRTrader.Modules.CongressTrading.Persistence.Ef;
/// <summary>
/// EF-Core-Kontext des CongressTrading-Moduls (ct_-Tabellen in derselben MariaDB wie der Core).
/// </summary>
public class CongressTradingDbContext : DbContext
{
public CongressTradingDbContext(DbContextOptions<CongressTradingDbContext> options) : base(options) { }
public DbSet<CongressMember> Members => Set<CongressMember>();
public DbSet<CongressTrade> Trades => Set<CongressTrade>();
protected override void OnModelCreating(ModelBuilder b)
{
b.Entity<CongressMember>(e =>
{
e.ToTable("ct_congressMember");
e.HasKey(x => x.Id);
e.HasIndex(x => x.BioId).IsUnique();
e.Property(x => x.BioId).HasMaxLength(20);
e.Property(x => x.Name).HasMaxLength(200);
e.Property(x => x.Party).HasMaxLength(50);
e.Property(x => x.State).HasMaxLength(50);
e.Property(x => x.Chamber).HasMaxLength(10);
e.Property(x => x.ProfileUrl).HasMaxLength(500);
});
b.Entity<CongressTrade>(e =>
{
e.ToTable("ct_trade");
e.HasKey(x => x.Id);
e.Ignore(x => x.MemberSnapshot); // transient nicht persistieren
e.HasIndex(x => x.TradeId).IsUnique();
e.HasIndex(x => x.MemberBioId);
e.HasIndex(x => x.Ticker);
e.HasIndex(x => x.DetailsFetched);
e.Property(x => x.TradeId).HasMaxLength(20);
e.Property(x => x.MemberBioId).HasMaxLength(20);
e.Property(x => x.IssuerName).HasMaxLength(500);
e.Property(x => x.IssuerId).HasMaxLength(20);
e.Property(x => x.Ticker).HasMaxLength(20);
e.Property(x => x.TradeType).HasMaxLength(50);
e.Property(x => x.Chamber).HasMaxLength(20);
e.Property(x => x.Owner).HasMaxLength(50);
e.Property(x => x.Value).HasPrecision(18, 2);
e.Property(x => x.DetailUrl).HasMaxLength(500);
});
}
}
@@ -0,0 +1,21 @@
using IBKRTrader.Core.Configuration;
using Microsoft.EntityFrameworkCore;
using Microsoft.EntityFrameworkCore.Design;
namespace IBKRTrader.Modules.CongressTrading.Persistence.Ef;
/// <summary>Design-Time-Factory für EF-Tooling (dotnet ef). Connection aus env IBKRTRADER_MYSQL.</summary>
public class CongressTradingDbContextFactory : IDesignTimeDbContextFactory<CongressTradingDbContext>
{
public CongressTradingDbContext CreateDbContext(string[] args)
{
var conn = Environment.GetEnvironmentVariable("IBKRTRADER_MYSQL")
?? "Server=localhost;Port=3306;Database=ibkrtrader;User ID=root;Password=;";
var options = new DbContextOptionsBuilder<CongressTradingDbContext>()
.UseMySql(conn, DatabaseServerVersion.Value)
.Options;
return new CongressTradingDbContext(options);
}
}
@@ -0,0 +1,165 @@
// <auto-generated />
using System;
using IBKRTrader.Modules.CongressTrading.Persistence.Ef;
using Microsoft.EntityFrameworkCore;
using Microsoft.EntityFrameworkCore.Infrastructure;
using Microsoft.EntityFrameworkCore.Metadata;
using Microsoft.EntityFrameworkCore.Migrations;
using Microsoft.EntityFrameworkCore.Storage.ValueConversion;
#nullable disable
namespace IBKRTrader.Modules.CongressTrading.Persistence.Ef.Migrations
{
[DbContext(typeof(CongressTradingDbContext))]
[Migration("20260728092729_InitialCongressTrading")]
partial class InitialCongressTrading
{
/// <inheritdoc />
protected override void BuildTargetModel(ModelBuilder modelBuilder)
{
#pragma warning disable 612, 618
modelBuilder
.HasAnnotation("ProductVersion", "8.0.13")
.HasAnnotation("Relational:MaxIdentifierLength", 64);
MySqlModelBuilderExtensions.AutoIncrementColumns(modelBuilder);
modelBuilder.Entity("IBKRTrader.Modules.CongressTrading.Models.CongressMember", b =>
{
b.Property<long>("Id")
.ValueGeneratedOnAdd()
.HasColumnType("bigint");
MySqlPropertyBuilderExtensions.UseMySqlIdentityColumn(b.Property<long>("Id"));
b.Property<string>("BioId")
.IsRequired()
.HasMaxLength(20)
.HasColumnType("varchar(20)");
b.Property<string>("Chamber")
.IsRequired()
.HasMaxLength(10)
.HasColumnType("varchar(10)");
b.Property<DateTime>("FirstSeen")
.HasColumnType("datetime(6)");
b.Property<DateTime>("LastUpdated")
.HasColumnType("datetime(6)");
b.Property<string>("Name")
.IsRequired()
.HasMaxLength(200)
.HasColumnType("varchar(200)");
b.Property<string>("Party")
.IsRequired()
.HasMaxLength(50)
.HasColumnType("varchar(50)");
b.Property<string>("ProfileUrl")
.IsRequired()
.HasMaxLength(500)
.HasColumnType("varchar(500)");
b.Property<string>("State")
.IsRequired()
.HasMaxLength(50)
.HasColumnType("varchar(50)");
b.HasKey("Id");
b.HasIndex("BioId")
.IsUnique();
b.ToTable("ct_congressMember", (string)null);
});
modelBuilder.Entity("IBKRTrader.Modules.CongressTrading.Models.CongressTrade", b =>
{
b.Property<long>("Id")
.ValueGeneratedOnAdd()
.HasColumnType("bigint");
MySqlPropertyBuilderExtensions.UseMySqlIdentityColumn(b.Property<long>("Id"));
b.Property<string>("Chamber")
.IsRequired()
.HasMaxLength(20)
.HasColumnType("varchar(20)");
b.Property<string>("DetailUrl")
.IsRequired()
.HasMaxLength(500)
.HasColumnType("varchar(500)");
b.Property<bool>("DetailsFetched")
.HasColumnType("tinyint(1)");
b.Property<string>("IssuerId")
.IsRequired()
.HasMaxLength(20)
.HasColumnType("varchar(20)");
b.Property<string>("IssuerName")
.IsRequired()
.HasMaxLength(500)
.HasColumnType("varchar(500)");
b.Property<string>("MemberBioId")
.IsRequired()
.HasMaxLength(20)
.HasColumnType("varchar(20)");
b.Property<string>("Owner")
.IsRequired()
.HasMaxLength(50)
.HasColumnType("varchar(50)");
b.Property<DateOnly?>("PublishedDate")
.HasColumnType("date");
b.Property<DateTime>("ScrapedAt")
.HasColumnType("datetime(6)");
b.Property<string>("Ticker")
.IsRequired()
.HasMaxLength(20)
.HasColumnType("varchar(20)");
b.Property<DateOnly?>("TradeDate")
.HasColumnType("date");
b.Property<string>("TradeId")
.IsRequired()
.HasMaxLength(20)
.HasColumnType("varchar(20)");
b.Property<string>("TradeType")
.IsRequired()
.HasMaxLength(50)
.HasColumnType("varchar(50)");
b.Property<decimal?>("Value")
.HasPrecision(18, 2)
.HasColumnType("decimal(18,2)");
b.HasKey("Id");
b.HasIndex("DetailsFetched");
b.HasIndex("MemberBioId");
b.HasIndex("Ticker");
b.HasIndex("TradeId")
.IsUnique();
b.ToTable("ct_trade", (string)null);
});
#pragma warning restore 612, 618
}
}
}
@@ -0,0 +1,119 @@
using System;
using Microsoft.EntityFrameworkCore.Metadata;
using Microsoft.EntityFrameworkCore.Migrations;
#nullable disable
namespace IBKRTrader.Modules.CongressTrading.Persistence.Ef.Migrations
{
/// <inheritdoc />
public partial class InitialCongressTrading : Migration
{
/// <inheritdoc />
protected override void Up(MigrationBuilder migrationBuilder)
{
migrationBuilder.AlterDatabase()
.Annotation("MySql:CharSet", "utf8mb4");
migrationBuilder.CreateTable(
name: "ct_congressMember",
columns: table => new
{
Id = table.Column<long>(type: "bigint", nullable: false)
.Annotation("MySql:ValueGenerationStrategy", MySqlValueGenerationStrategy.IdentityColumn),
BioId = table.Column<string>(type: "varchar(20)", maxLength: 20, nullable: false)
.Annotation("MySql:CharSet", "utf8mb4"),
Name = table.Column<string>(type: "varchar(200)", maxLength: 200, nullable: false)
.Annotation("MySql:CharSet", "utf8mb4"),
Party = table.Column<string>(type: "varchar(50)", maxLength: 50, nullable: false)
.Annotation("MySql:CharSet", "utf8mb4"),
State = table.Column<string>(type: "varchar(50)", maxLength: 50, nullable: false)
.Annotation("MySql:CharSet", "utf8mb4"),
Chamber = table.Column<string>(type: "varchar(10)", maxLength: 10, nullable: false)
.Annotation("MySql:CharSet", "utf8mb4"),
ProfileUrl = table.Column<string>(type: "varchar(500)", maxLength: 500, nullable: false)
.Annotation("MySql:CharSet", "utf8mb4"),
FirstSeen = table.Column<DateTime>(type: "datetime(6)", nullable: false),
LastUpdated = table.Column<DateTime>(type: "datetime(6)", nullable: false)
},
constraints: table =>
{
table.PrimaryKey("PK_ct_congressMember", x => x.Id);
})
.Annotation("MySql:CharSet", "utf8mb4");
migrationBuilder.CreateTable(
name: "ct_trade",
columns: table => new
{
Id = table.Column<long>(type: "bigint", nullable: false)
.Annotation("MySql:ValueGenerationStrategy", MySqlValueGenerationStrategy.IdentityColumn),
TradeId = table.Column<string>(type: "varchar(20)", maxLength: 20, nullable: false)
.Annotation("MySql:CharSet", "utf8mb4"),
MemberBioId = table.Column<string>(type: "varchar(20)", maxLength: 20, nullable: false)
.Annotation("MySql:CharSet", "utf8mb4"),
IssuerName = table.Column<string>(type: "varchar(500)", maxLength: 500, nullable: false)
.Annotation("MySql:CharSet", "utf8mb4"),
IssuerId = table.Column<string>(type: "varchar(20)", maxLength: 20, nullable: false)
.Annotation("MySql:CharSet", "utf8mb4"),
Ticker = table.Column<string>(type: "varchar(20)", maxLength: 20, nullable: false)
.Annotation("MySql:CharSet", "utf8mb4"),
TradeType = table.Column<string>(type: "varchar(50)", maxLength: 50, nullable: false)
.Annotation("MySql:CharSet", "utf8mb4"),
Chamber = table.Column<string>(type: "varchar(20)", maxLength: 20, nullable: false)
.Annotation("MySql:CharSet", "utf8mb4"),
Owner = table.Column<string>(type: "varchar(50)", maxLength: 50, nullable: false)
.Annotation("MySql:CharSet", "utf8mb4"),
Value = table.Column<decimal>(type: "decimal(18,2)", precision: 18, scale: 2, nullable: true),
TradeDate = table.Column<DateOnly>(type: "date", nullable: true),
PublishedDate = table.Column<DateOnly>(type: "date", nullable: true),
DetailUrl = table.Column<string>(type: "varchar(500)", maxLength: 500, nullable: false)
.Annotation("MySql:CharSet", "utf8mb4"),
DetailsFetched = table.Column<bool>(type: "tinyint(1)", nullable: false),
ScrapedAt = table.Column<DateTime>(type: "datetime(6)", nullable: false)
},
constraints: table =>
{
table.PrimaryKey("PK_ct_trade", x => x.Id);
})
.Annotation("MySql:CharSet", "utf8mb4");
migrationBuilder.CreateIndex(
name: "IX_ct_congressMember_BioId",
table: "ct_congressMember",
column: "BioId",
unique: true);
migrationBuilder.CreateIndex(
name: "IX_ct_trade_DetailsFetched",
table: "ct_trade",
column: "DetailsFetched");
migrationBuilder.CreateIndex(
name: "IX_ct_trade_MemberBioId",
table: "ct_trade",
column: "MemberBioId");
migrationBuilder.CreateIndex(
name: "IX_ct_trade_Ticker",
table: "ct_trade",
column: "Ticker");
migrationBuilder.CreateIndex(
name: "IX_ct_trade_TradeId",
table: "ct_trade",
column: "TradeId",
unique: true);
}
/// <inheritdoc />
protected override void Down(MigrationBuilder migrationBuilder)
{
migrationBuilder.DropTable(
name: "ct_congressMember");
migrationBuilder.DropTable(
name: "ct_trade");
}
}
}
@@ -0,0 +1,162 @@
// <auto-generated />
using System;
using IBKRTrader.Modules.CongressTrading.Persistence.Ef;
using Microsoft.EntityFrameworkCore;
using Microsoft.EntityFrameworkCore.Infrastructure;
using Microsoft.EntityFrameworkCore.Metadata;
using Microsoft.EntityFrameworkCore.Storage.ValueConversion;
#nullable disable
namespace IBKRTrader.Modules.CongressTrading.Persistence.Ef.Migrations
{
[DbContext(typeof(CongressTradingDbContext))]
partial class CongressTradingDbContextModelSnapshot : ModelSnapshot
{
protected override void BuildModel(ModelBuilder modelBuilder)
{
#pragma warning disable 612, 618
modelBuilder
.HasAnnotation("ProductVersion", "8.0.13")
.HasAnnotation("Relational:MaxIdentifierLength", 64);
MySqlModelBuilderExtensions.AutoIncrementColumns(modelBuilder);
modelBuilder.Entity("IBKRTrader.Modules.CongressTrading.Models.CongressMember", b =>
{
b.Property<long>("Id")
.ValueGeneratedOnAdd()
.HasColumnType("bigint");
MySqlPropertyBuilderExtensions.UseMySqlIdentityColumn(b.Property<long>("Id"));
b.Property<string>("BioId")
.IsRequired()
.HasMaxLength(20)
.HasColumnType("varchar(20)");
b.Property<string>("Chamber")
.IsRequired()
.HasMaxLength(10)
.HasColumnType("varchar(10)");
b.Property<DateTime>("FirstSeen")
.HasColumnType("datetime(6)");
b.Property<DateTime>("LastUpdated")
.HasColumnType("datetime(6)");
b.Property<string>("Name")
.IsRequired()
.HasMaxLength(200)
.HasColumnType("varchar(200)");
b.Property<string>("Party")
.IsRequired()
.HasMaxLength(50)
.HasColumnType("varchar(50)");
b.Property<string>("ProfileUrl")
.IsRequired()
.HasMaxLength(500)
.HasColumnType("varchar(500)");
b.Property<string>("State")
.IsRequired()
.HasMaxLength(50)
.HasColumnType("varchar(50)");
b.HasKey("Id");
b.HasIndex("BioId")
.IsUnique();
b.ToTable("ct_congressMember", (string)null);
});
modelBuilder.Entity("IBKRTrader.Modules.CongressTrading.Models.CongressTrade", b =>
{
b.Property<long>("Id")
.ValueGeneratedOnAdd()
.HasColumnType("bigint");
MySqlPropertyBuilderExtensions.UseMySqlIdentityColumn(b.Property<long>("Id"));
b.Property<string>("Chamber")
.IsRequired()
.HasMaxLength(20)
.HasColumnType("varchar(20)");
b.Property<string>("DetailUrl")
.IsRequired()
.HasMaxLength(500)
.HasColumnType("varchar(500)");
b.Property<bool>("DetailsFetched")
.HasColumnType("tinyint(1)");
b.Property<string>("IssuerId")
.IsRequired()
.HasMaxLength(20)
.HasColumnType("varchar(20)");
b.Property<string>("IssuerName")
.IsRequired()
.HasMaxLength(500)
.HasColumnType("varchar(500)");
b.Property<string>("MemberBioId")
.IsRequired()
.HasMaxLength(20)
.HasColumnType("varchar(20)");
b.Property<string>("Owner")
.IsRequired()
.HasMaxLength(50)
.HasColumnType("varchar(50)");
b.Property<DateOnly?>("PublishedDate")
.HasColumnType("date");
b.Property<DateTime>("ScrapedAt")
.HasColumnType("datetime(6)");
b.Property<string>("Ticker")
.IsRequired()
.HasMaxLength(20)
.HasColumnType("varchar(20)");
b.Property<DateOnly?>("TradeDate")
.HasColumnType("date");
b.Property<string>("TradeId")
.IsRequired()
.HasMaxLength(20)
.HasColumnType("varchar(20)");
b.Property<string>("TradeType")
.IsRequired()
.HasMaxLength(50)
.HasColumnType("varchar(50)");
b.Property<decimal?>("Value")
.HasPrecision(18, 2)
.HasColumnType("decimal(18,2)");
b.HasKey("Id");
b.HasIndex("DetailsFetched");
b.HasIndex("MemberBioId");
b.HasIndex("Ticker");
b.HasIndex("TradeId")
.IsUnique();
b.ToTable("ct_trade", (string)null);
});
#pragma warning restore 612, 618
}
}
}
@@ -1,3 +1,5 @@
using IBKRTrader.Core.Persistence.Ef;
using Microsoft.EntityFrameworkCore;
using IBKRTrader.Core.Database; using IBKRTrader.Core.Database;
using IBKRTrader.Core.Logging; using IBKRTrader.Core.Logging;
using IBKRTrader.Modules.CongressTrading.Database; using IBKRTrader.Modules.CongressTrading.Database;
@@ -34,7 +36,7 @@ public class CongressHistoryImportWorker : WorkerBase
public CongressHistoryImportWorker( public CongressHistoryImportWorker(
LoggingService logger, LoggingService logger,
DatabaseService db, IDbContextFactory<CoreDbContext> db,
CapitolTradesScraper scraper, CapitolTradesScraper scraper,
CongressRepository repo) CongressRepository repo)
: base(logger, db) : base(logger, db)
@@ -1,3 +1,5 @@
using IBKRTrader.Core.Persistence.Ef;
using Microsoft.EntityFrameworkCore;
using IBKRTrader.Core.Database; using IBKRTrader.Core.Database;
using IBKRTrader.Core.Logging; using IBKRTrader.Core.Logging;
using IBKRTrader.Core.Settings; using IBKRTrader.Core.Settings;
@@ -25,7 +27,7 @@ public class CongressScrapeWorker : WorkerBase
public CongressScrapeWorker( public CongressScrapeWorker(
LoggingService logger, LoggingService logger,
DatabaseService db, IDbContextFactory<CoreDbContext> db,
CapitolTradesScraper scraper, CapitolTradesScraper scraper,
CongressRepository repo, CongressRepository repo,
SettingsService settings) SettingsService settings)
@@ -0,0 +1,81 @@
using FluentAssertions;
using IBKRTrader.Core.Logging;
using IBKRTrader.Core.Persistence.Ef;
using IBKRTrader.Core.Settings;
using IBKRTrader.Modules.CongressTrading.Database;
using IBKRTrader.Modules.CongressTrading.Models;
using IBKRTrader.Modules.CongressTrading.Persistence.Ef;
using Microsoft.EntityFrameworkCore;
namespace IBKRTrader.Tests.Modules;
/// <summary>EF-Repo des Moduls gegen EF-InMemory (deterministisch, kein externer DB-Zugriff).</summary>
[Trait("cat", "unit")]
public class CongressRepositoryTests
{
private sealed class Factory<T>(DbContextOptions<T> options) : IDbContextFactory<T> where T : DbContext
{
public T CreateDbContext() => (T)Activator.CreateInstance(typeof(T), options)!;
}
private static CongressRepository CreateSut()
{
var ctOpts = new DbContextOptionsBuilder<CongressTradingDbContext>()
.UseInMemoryDatabase(Guid.NewGuid().ToString()).Options;
var coreOpts = new DbContextOptionsBuilder<CoreDbContext>()
.UseInMemoryDatabase(Guid.NewGuid().ToString()).Options;
var settings = new CoreSettingsService(new Factory<CoreDbContext>(coreOpts));
return new CongressRepository(new Factory<CongressTradingDbContext>(ctOpts), settings, new LoggingService());
}
[Fact]
public async Task UpsertMember_IsIdempotent_ByBioId()
{
var repo = CreateSut();
await repo.UpsertMemberAsync(new CongressMember { BioId = "W1", Name = "Alice" });
await repo.UpsertMemberAsync(new CongressMember { BioId = "W1", Name = "Alice B." });
(await repo.MemberExistsAsync("W1")).Should().BeTrue();
(await repo.GetMemberCountAsync()).Should().Be(1);
(await repo.GetAllMemberBioIdsAsync()).Should().Contain("W1");
}
[Fact]
public async Task InsertTrade_IgnoresDuplicateTradeId()
{
var repo = CreateSut();
await repo.InsertTradeAsync(new CongressTrade { TradeId = "T1", MemberBioId = "W1", Ticker = "AAPL", TradeType = "buy" });
await repo.InsertTradeAsync(new CongressTrade { TradeId = "T1", MemberBioId = "W1", Ticker = "AAPL" });
(await repo.TradeExistsAsync("T1")).Should().BeTrue();
(await repo.GetTradeCountAsync()).Should().Be(1);
}
[Fact]
public async Task DetailsFlow_MarksTradeComplete()
{
var repo = CreateSut();
await repo.InsertTradeAsync(new CongressTrade { TradeId = "T2" }, detailsFetched: false);
(await repo.GetTradesWithoutDetailsAsync()).Should().Contain("T2");
await repo.UpdateTradeDetailsAsync(new CongressTrade { TradeId = "T2", Ticker = "MSFT", TradeType = "sell" });
(await repo.GetTradesWithoutDetailsAsync()).Should().NotContain("T2");
}
[Fact]
public async Task Settings_RoundTripThroughCoreSettings()
{
var repo = CreateSut();
await repo.SetSettingAsync("ct.history_import_page", "7");
(await repo.GetSettingAsync("ct.history_import_page")).Should().Be("7");
await repo.DeleteSettingAsync("ct.history_import_page");
(await repo.GetSettingAsync("ct.history_import_page")).Should().BeNull();
}
}
@@ -45,7 +45,6 @@ public class CongressTradingModuleTests
var types = services.Select(d => d.ServiceType).ToList(); var types = services.Select(d => d.ServiceType).ToList();
types.Should().Contain(new[] types.Should().Contain(new[]
{ {
typeof(CongressMigrations),
typeof(CongressRepository), typeof(CongressRepository),
typeof(CapitolTradesScraper), typeof(CapitolTradesScraper),
typeof(CongressHistoryImportWorker), typeof(CongressHistoryImportWorker),
@@ -1,9 +1,9 @@
using System.Diagnostics; using System.Diagnostics;
using FluentAssertions; using FluentAssertions;
using IBKRTrader.Core.Database;
using IBKRTrader.Core.Logging; using IBKRTrader.Core.Logging;
using IBKRTrader.Core.Settings; using IBKRTrader.Core.Persistence.Ef;
using IBKRTrader.Core.Workers; using IBKRTrader.Core.Workers;
using Microsoft.EntityFrameworkCore;
namespace IBKRTrader.Tests.Workers; namespace IBKRTrader.Tests.Workers;
@@ -26,9 +26,17 @@ public class WorkerBaseTests
public override string Module => "TEST"; public override string Module => "TEST";
protected override TimeSpan? Interval => _interval; protected override TimeSpan? Interval => _interval;
private sealed class InMemoryFactory(DbContextOptions<CoreDbContext> o) : IDbContextFactory<CoreDbContext>
{
public CoreDbContext CreateDbContext() => new(o);
}
private static IDbContextFactory<CoreDbContext> Dbf() =>
new InMemoryFactory(new DbContextOptionsBuilder<CoreDbContext>()
.UseInMemoryDatabase(Guid.NewGuid().ToString()).Options);
public TestWorker(TimeSpan? interval, Func<CancellationToken, Task> body) public TestWorker(TimeSpan? interval, Func<CancellationToken, Task> body)
: base(new LoggingService(), : base(new LoggingService(), Dbf())
new DatabaseService(new SettingsService(), new LoggingService()))
{ {
_interval = interval; _interval = interval;
_body = body; _body = body;