using Dapper; using IBKRTrader.Core.Database; using IBKRTrader.Core.Logging; using IBKRTrader.Modules.CongressTrading.Models; namespace IBKRTrader.Modules.CongressTrading.Database; /// /// Datenbankzugriff für das CongressTrading-Modul. /// Kapselt alle INSERT/SELECT auf ct_congressMember und ct_trade. /// public class CongressRepository { private readonly DatabaseService _db; private readonly LoggingService _logger; public CongressRepository(DatabaseService db, LoggingService logger) { _db = db; _logger = logger; } // ─── CongressMember ─────────────────────────────────────────────────────── public async Task MemberExistsAsync(string bioId) { var count = await _db.ExecuteScalarAsync( "SELECT COUNT(*) FROM `ct_congressMember` WHERE bio_id = @bioId", new { bioId }); return count > 0; } public async Task UpsertMemberAsync(CongressMember member) { await _db.ExecuteAsync(@" INSERT INTO `ct_congressMember` (bio_id, name, party, state, chamber, profile_url) VALUES (@BioId, @Name, @Party, @State, @Chamber, @ProfileUrl) ON DUPLICATE KEY UPDATE name = VALUES(name), party = VALUES(party), state = VALUES(state), chamber = VALUES(chamber), profile_url = VALUES(profile_url), last_updated = CURRENT_TIMESTAMP", member); } public async Task> GetAllMemberBioIdsAsync() { var ids = await _db.QueryAsync( "SELECT bio_id FROM `ct_congressMember`"); return ids.ToHashSet(); } // ─── CongressTrade ──────────────────────────────────────────────────────── public async Task TradeExistsAsync(string tradeId) { var count = await _db.ExecuteScalarAsync( "SELECT COUNT(*) FROM `ct_trade` WHERE trade_id = @tradeId", new { tradeId }); return count > 0; } public async Task> GetAllTradeIdsAsync() { var ids = await _db.QueryAsync( "SELECT trade_id FROM `ct_trade`"); return ids.ToHashSet(); } /// Fügt einen Trade ein. Ignoriert Duplikate (INSERT IGNORE). public async Task InsertTradeAsync(CongressTrade trade, bool detailsFetched = true) { await _db.ExecuteAsync(@" INSERT IGNORE INTO `ct_trade` (trade_id, member_bio_id, issuer_name, issuer_id, ticker, trade_type, chamber, owner, value, trade_date, published_date, detail_url, details_fetched) VALUES (@TradeId, @MemberBioId, @IssuerName, @IssuerId, @Ticker, @TradeType, @Chamber, @Owner, @Value, @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 }); } /// Aktualisiert die Detail-Felder eines bereits vorhandenen Trades. public async Task UpdateTradeDetailsAsync(CongressTrade trade) { await _db.ExecuteAsync(@" UPDATE `ct_trade` SET ticker = @Ticker, trade_type = @TradeType, trade_date = @TradeDate, published_date = @PublishedDate, details_fetched = 1 WHERE trade_id = @TradeId AND details_fetched = 0", new { trade.TradeId, trade.Ticker, trade.TradeType, TradeDate = trade.TradeDate?.ToString("yyyy-MM-dd"), PublishedDate = trade.PublishedDate?.ToString("yyyy-MM-dd") }); } /// Gibt Trade-IDs zurück für die noch keine Details geholt wurden. public async Task> GetTradesWithoutDetailsAsync(int limit = 500) => await _db.QueryAsync( "SELECT trade_id FROM `ct_trade` WHERE details_fetched = 0 LIMIT @limit", new { limit }); public async Task GetTradeCountAsync() => await _db.ExecuteScalarAsync("SELECT COUNT(*) FROM `ct_trade`"); public async Task GetMemberCountAsync() => await _db.ExecuteScalarAsync("SELECT COUNT(*) FROM `ct_congressMember`"); // ─── core_settings Flags ────────────────────────────────────────────────── public async Task GetSettingAsync(string key) { return await _db.QueryFirstOrDefaultAsync( "SELECT `value` FROM `core_settings` WHERE `key` = @key", new { key }); } public async Task SetSettingAsync(string key, string value) { 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) { await _db.ExecuteAsync( "DELETE FROM `core_settings` WHERE `key` = @key", new { key }); } // ─── Reset ──────────────────────────────────────────────────────────────── /// /// 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. /// public async Task ResetHistoryImportAsync() { await DeleteSettingAsync("ct.history_import_done"); await DeleteSettingAsync("ct.history_import_page"); await DeleteSettingAsync("ct.history_import_total"); await _db.ExecuteAsync("DELETE FROM `ct_trade`"); await _db.ExecuteAsync("DELETE FROM `ct_congressMember`"); _logger.Warn("CT", "History-Import-Reset durchgeführt – alle CT-Daten gelöscht."); } }