using Backend.Interface; using Npgsql; namespace Backend.DatabaseHandler { public sealed class PostgresNewsSentimentStore : INewsSentimentStore { private readonly string _connectionString; public PostgresNewsSentimentStore(IConfiguration configuration) { _connectionString = configuration.GetConnectionString("MarketDataDb") ?? throw new InvalidOperationException("Missing connection string: MarketDataDb"); } public async Task> GetByKeywordAsync( string keyword, CancellationToken cancellationToken = default) { var results = new List(); await using var conn = new NpgsqlConnection(_connectionString); await conn.OpenAsync(cancellationToken); const string sql = """ SELECT source, keyword, published_at, title, score, sentiment_label, title_hash, publisher, url FROM news_sentiment_cache WHERE keyword = @keyword ORDER BY published_at DESC; """; await using var cmd = new NpgsqlCommand(sql, conn); cmd.Parameters.AddWithValue("keyword", keyword); await using var reader = await cmd.ExecuteReaderAsync(cancellationToken); while (await reader.ReadAsync(cancellationToken)) { results.Add(new NewsSentimentRecord { Source = reader.GetString(0), Keyword = reader.GetString(1), PublishedAt = reader.GetFieldValue(2), Title = reader.GetString(3), Score = reader.GetDouble(4), SentimentLabel = reader.GetString(5), TitleHash = reader.GetString(6).TrimEnd(), Publisher = reader.GetString(7), Url = reader.GetString(8) }); } return results; } public async Task> GetFeedAsync( IEnumerable keywords, int limit, CancellationToken cancellationToken = default) { var keywordList = keywords.ToList(); if (keywordList.Count == 0) return []; var results = new List(); await using var conn = new NpgsqlConnection(_connectionString); await conn.OpenAsync(cancellationToken); const string sql = """ SELECT source, keyword, published_at, title, score, sentiment_label, title_hash, publisher, url FROM news_sentiment_cache WHERE keyword = ANY(@keywords) ORDER BY published_at DESC LIMIT @limit; """; await using var cmd = new NpgsqlCommand(sql, conn); cmd.Parameters.AddWithValue("keywords", keywordList.ToArray()); cmd.Parameters.AddWithValue("limit", limit); await using var reader = await cmd.ExecuteReaderAsync(cancellationToken); while (await reader.ReadAsync(cancellationToken)) { results.Add(new NewsSentimentRecord { Source = reader.GetString(0), Keyword = reader.GetString(1), PublishedAt = reader.GetFieldValue(2), Title = reader.GetString(3), Score = reader.GetDouble(4), SentimentLabel = reader.GetString(5), TitleHash = reader.GetString(6).TrimEnd(), Publisher = reader.GetString(7), Url = reader.GetString(8) }); } return results; } public async Task InsertAsync( IEnumerable records, CancellationToken cancellationToken = default) { var recordList = records.ToList(); if (recordList.Count == 0) return; await using var conn = new NpgsqlConnection(_connectionString); await conn.OpenAsync(cancellationToken); await using var tx = await conn.BeginTransactionAsync(cancellationToken); const string sql = """ INSERT INTO news_sentiment_cache (source, keyword, published_at, title, score, sentiment_label, title_hash, publisher, url) VALUES (@source, @keyword, @published_at, @title, @score, @sentiment_label, @title_hash, @publisher, @url) ON CONFLICT (title_hash) DO NOTHING; """; foreach (var record in recordList) { await using var cmd = new NpgsqlCommand(sql, conn, tx); cmd.Parameters.AddWithValue("source", record.Source); cmd.Parameters.AddWithValue("keyword", record.Keyword); cmd.Parameters.AddWithValue("published_at", record.PublishedAt); cmd.Parameters.AddWithValue("title", record.Title); cmd.Parameters.AddWithValue("score", record.Score); cmd.Parameters.AddWithValue("sentiment_label", record.SentimentLabel); cmd.Parameters.AddWithValue("title_hash", record.TitleHash); cmd.Parameters.AddWithValue("publisher", record.Publisher); cmd.Parameters.AddWithValue("url", record.Url); await cmd.ExecuteNonQueryAsync(cancellationToken); } await tx.CommitAsync(cancellationToken); } public async Task> GetExistingHashesAsync( string keyword, CancellationToken cancellationToken = default) { var hashes = new HashSet(); await using var conn = new NpgsqlConnection(_connectionString); await conn.OpenAsync(cancellationToken); const string sql = """ SELECT title_hash FROM news_sentiment_cache WHERE keyword = @keyword; """; await using var cmd = new NpgsqlCommand(sql, conn); cmd.Parameters.AddWithValue("keyword", keyword); await using var reader = await cmd.ExecuteReaderAsync(cancellationToken); while (await reader.ReadAsync(cancellationToken)) { hashes.Add(reader.GetString(0).TrimEnd()); } return hashes; } } }