using Dapper; using IBKRTrader.Core.Database; using IBKRTrader.Core.Logging; namespace IBKRTrader.Core.IBKR; /// /// Datenzugriffsschicht für die core_ibkr_xxx-Tabellen. /// Alle IBKR-Marktdaten-Operationen laufen über diese Klasse. /// public class IBKRMarketDataRepository { private readonly DatabaseService _db; private readonly LoggingService _logger; public IBKRMarketDataRepository(DatabaseService db, LoggingService logger) { _db = db; _logger = logger; } // ─── Instruments ───────────────────────────────────────────────────────── /// /// Legt ein neues Instrument an oder aktualisiert ein bestehendes (UPSERT via conid). /// Gibt die Instrument-ID zurück. /// public async Task UpsertInstrumentAsync(IBKRInstrument instr) { const string sql = @" INSERT INTO `core_ibkr_instruments` (`ibkr_conid`, `symbol`, `sec_type`, `exchange`, `primary_exchange`, `currency`, `company_name`, `isin`, `sector`, `industry`, `description`, `active`, `last_fetched`) VALUES (@IbkrConid, @Symbol, @SecType, @Exchange, @PrimaryExchange, @Currency, @CompanyName, @Isin, @Sector, @Industry, @Description, @Active, @LastFetched) ON DUPLICATE KEY UPDATE `symbol` = VALUES(`symbol`), `sec_type` = VALUES(`sec_type`), `exchange` = VALUES(`exchange`), `primary_exchange` = VALUES(`primary_exchange`), `currency` = VALUES(`currency`), `company_name` = VALUES(`company_name`), `isin` = VALUES(`isin`), `sector` = VALUES(`sector`), `industry` = VALUES(`industry`), `description` = VALUES(`description`), `last_fetched` = VALUES(`last_fetched`); SELECT `id` FROM `core_ibkr_instruments` WHERE `ibkr_conid` = @IbkrConid;"; await using var conn = _db.CreateConnection(); return await conn.ExecuteScalarAsync(sql, instr); } /// Gibt alle aktiven Instrumente zurück. public Task> GetAllActiveInstrumentsAsync() => _db.QueryAsync( "SELECT * FROM `core_ibkr_instruments` WHERE `active` = 1 ORDER BY `symbol`"); /// Findet ein Instrument anhand seiner IBKR ConID. public Task GetInstrumentByConidAsync(long conid) => _db.QueryFirstOrDefaultAsync( "SELECT * FROM `core_ibkr_instruments` WHERE `ibkr_conid` = @conid", new { conid }); /// Gibt die Anzahl aktiver Instrumente zurück. public Task GetActiveInstrumentCountAsync() => _db.ExecuteScalarAsync( "SELECT COUNT(*) FROM `core_ibkr_instruments` WHERE `active` = 1"); // ─── Market Data ───────────────────────────────────────────────────────── /// /// Fügt Marktdaten-Balken via UPSERT ein (ON DUPLICATE KEY UPDATE). /// public async Task UpsertMarketDataBatchAsync(IEnumerable bars) { const string sql = @" INSERT INTO `core_ibkr_market_data` (`instrument_id`, `bar_size`, `timestamp`, `open`, `high`, `low`, `close`, `volume`, `wap`, `bar_count`) VALUES (@InstrumentId, @BarSize, @Timestamp, @Open, @High, @Low, @Close, @Volume, @Wap, @BarCount) ON DUPLICATE KEY UPDATE `open` = VALUES(`open`), `high` = VALUES(`high`), `low` = VALUES(`low`), `close` = VALUES(`close`), `volume` = VALUES(`volume`), `wap` = VALUES(`wap`), `bar_count` = VALUES(`bar_count`)"; await using var conn = _db.CreateConnection(); await conn.OpenAsync(); await conn.ExecuteAsync(sql, bars); } /// Gibt den neuesten Timestamp für ein Instrument zurück. public Task GetLatestBarTimestampAsync(long instrumentId, string barSize = "daily") => _db.QueryFirstOrDefaultAsync( @"SELECT MAX(`timestamp`) FROM `core_ibkr_market_data` WHERE `instrument_id` = @instrumentId AND `bar_size` = @barSize", new { instrumentId, barSize }); /// Gibt die Anzahl Bars für ein Instrument zurück. public Task GetBarCountAsync(long instrumentId, string barSize = "daily") => _db.ExecuteScalarAsync( @"SELECT COUNT(*) FROM `core_ibkr_market_data` WHERE `instrument_id` = @instrumentId AND `bar_size` = @barSize", new { instrumentId, barSize }); // ─── External Identifiers ──────────────────────────────────────────────── /// /// Erstellt oder ignoriert ein External-Identifier-Mapping (IGNORE bei Duplikat). /// public Task UpsertExternalIdentifierAsync(long instrumentId, string source, string ticker) => _db.ExecuteAsync(@" INSERT IGNORE INTO `core_ibkr_external_identifiers` (`instrument_id`, `source`, `ticker`) VALUES (@instrumentId, @source, @ticker)", new { instrumentId, source, ticker }); /// /// Findet ein Instrument anhand eines externen Tickers (z.B. aus ct_trade). /// public Task FindInstrumentByExternalTickerAsync(string source, string ticker) => _db.QueryFirstOrDefaultAsync(@" SELECT i.* FROM `core_ibkr_instruments` i INNER JOIN `core_ibkr_external_identifiers` e ON e.`instrument_id` = i.`id` WHERE e.`source` = @source AND e.`ticker` = @ticker LIMIT 1", new { source, ticker }); /// /// Findet alle einzigartigen Ticker aus ct_trade, die noch kein IBKR-Mapping haben. /// public Task> GetUnmappedTickersFromCongressTradesAsync() => _db.QueryAsync(@" SELECT DISTINCT t.`ticker` FROM `ct_trade` t WHERE t.`ticker` IS NOT NULL AND t.`ticker` != '' AND NOT EXISTS ( SELECT 1 FROM `core_ibkr_external_identifiers` e WHERE e.`source` = 'capitoltrades' AND e.`ticker` = t.`ticker` ) ORDER BY t.`ticker`"); }