SoftTraderBackend/DatabaseHandler/PostgresNewsSentimentStore.cs
2026-03-20 10:42:48 -04:00

124 lines
3.9 KiB
C#

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<List<NewsSentimentRecord>> GetByKeywordAsync(
string keyword,
CancellationToken cancellationToken = default)
{
var results = new List<NewsSentimentRecord>();
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, 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<DateTimeOffset>(2),
Title = reader.GetString(3),
Score = reader.GetDouble(4),
SentimentLabel = reader.GetString(5),
TitleHash = reader.GetString(6).TrimEnd(),
Url = reader.GetString(7)
});
}
return results;
}
public async Task InsertAsync(
IEnumerable<NewsSentimentRecord> 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, url)
VALUES
(@source, @keyword, @published_at, @title, @score,
@sentiment_label, @title_hash, @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("url", record.Url);
await cmd.ExecuteNonQueryAsync(cancellationToken);
}
await tx.CommitAsync(cancellationToken);
}
public async Task<HashSet<string>> GetExistingHashesAsync(
string keyword,
CancellationToken cancellationToken = default)
{
var hashes = new HashSet<string>();
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;
}
}
}