165 lines
5.2 KiB
C#
165 lines
5.2 KiB
C#
using Npgsql;
|
|
using Backend.Interface;
|
|
|
|
namespace Backend.DatabaseHandler;
|
|
|
|
public sealed class PostgresAlertStore : IAlertStore
|
|
{
|
|
private readonly string _connectionString;
|
|
|
|
public PostgresAlertStore(IConfiguration configuration)
|
|
{
|
|
_connectionString = configuration.GetConnectionString("MarketDataDb")
|
|
?? throw new InvalidOperationException("Missing connection string: MarketDataDb");
|
|
|
|
EnsureTableAsync().GetAwaiter().GetResult();
|
|
}
|
|
|
|
private async Task EnsureTableAsync()
|
|
{
|
|
const string sql = """
|
|
CREATE TABLE IF NOT EXISTS user_alerts (
|
|
id BIGSERIAL PRIMARY KEY,
|
|
username TEXT NOT NULL,
|
|
symbol TEXT NOT NULL,
|
|
alert_type TEXT NOT NULL,
|
|
parameters JSONB NOT NULL DEFAULT '{}'::jsonb,
|
|
last_state JSONB NOT NULL DEFAULT '{}'::jsonb,
|
|
is_triggered BOOLEAN NOT NULL DEFAULT FALSE,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
|
);
|
|
|
|
CREATE INDEX IF NOT EXISTS idx_user_alerts_username
|
|
ON user_alerts (username);
|
|
|
|
CREATE INDEX IF NOT EXISTS idx_user_alerts_active
|
|
ON user_alerts (is_triggered) WHERE is_triggered = FALSE;
|
|
""";
|
|
|
|
await using var conn = new NpgsqlConnection(_connectionString);
|
|
await conn.OpenAsync();
|
|
await using var cmd = new NpgsqlCommand(sql, conn);
|
|
await cmd.ExecuteNonQueryAsync();
|
|
}
|
|
|
|
public async Task<List<AlertRecord>> GetActiveAlertsAsync(CancellationToken cancellationToken = default)
|
|
{
|
|
const string sql = """
|
|
SELECT id, username, symbol, alert_type, parameters, last_state, is_triggered, created_at
|
|
FROM user_alerts
|
|
WHERE is_triggered = FALSE
|
|
ORDER BY id;
|
|
""";
|
|
|
|
return await QueryAlertsAsync(sql, null, cancellationToken);
|
|
}
|
|
|
|
public async Task<List<AlertRecord>> GetAlertsByUserAsync(string username, CancellationToken cancellationToken = default)
|
|
{
|
|
const string sql = """
|
|
SELECT id, username, symbol, alert_type, parameters, last_state, is_triggered, created_at
|
|
FROM user_alerts
|
|
WHERE username = @username
|
|
ORDER BY created_at DESC;
|
|
""";
|
|
|
|
return await QueryAlertsAsync(sql, cmd => cmd.Parameters.AddWithValue("username", username.Trim()), cancellationToken);
|
|
}
|
|
|
|
public async Task<AlertRecord> CreateAlertAsync(AlertRecord alert, CancellationToken cancellationToken = default)
|
|
{
|
|
const string sql = """
|
|
INSERT INTO user_alerts (username, symbol, alert_type, parameters, last_state)
|
|
VALUES (@username, @symbol, @alert_type, @parameters::jsonb, @last_state::jsonb)
|
|
RETURNING id, created_at;
|
|
""";
|
|
|
|
await using var conn = new NpgsqlConnection(_connectionString);
|
|
await conn.OpenAsync(cancellationToken);
|
|
|
|
await using var cmd = new NpgsqlCommand(sql, conn);
|
|
cmd.Parameters.AddWithValue("username", alert.Username.Trim());
|
|
cmd.Parameters.AddWithValue("symbol", alert.Symbol.Trim().ToUpperInvariant());
|
|
cmd.Parameters.AddWithValue("alert_type", alert.AlertType);
|
|
cmd.Parameters.AddWithValue("parameters", alert.Parameters);
|
|
cmd.Parameters.AddWithValue("last_state", alert.LastState);
|
|
|
|
await using var reader = await cmd.ExecuteReaderAsync(cancellationToken);
|
|
await reader.ReadAsync(cancellationToken);
|
|
|
|
alert.Id = reader.GetInt64(0);
|
|
alert.CreatedAt = reader.GetDateTime(1);
|
|
return alert;
|
|
}
|
|
|
|
public async Task<bool> DeleteAlertAsync(string username, long alertId, CancellationToken cancellationToken = default)
|
|
{
|
|
const string sql = """
|
|
DELETE FROM user_alerts
|
|
WHERE id = @id AND username = @username;
|
|
""";
|
|
|
|
await using var conn = new NpgsqlConnection(_connectionString);
|
|
await conn.OpenAsync(cancellationToken);
|
|
|
|
await using var cmd = new NpgsqlCommand(sql, conn);
|
|
cmd.Parameters.AddWithValue("id", alertId);
|
|
cmd.Parameters.AddWithValue("username", username.Trim());
|
|
|
|
int rows = await cmd.ExecuteNonQueryAsync(cancellationToken);
|
|
return rows > 0;
|
|
}
|
|
|
|
public async Task UpdateAlertAsync(AlertRecord alert, CancellationToken cancellationToken = default)
|
|
{
|
|
const string sql = """
|
|
UPDATE user_alerts
|
|
SET is_triggered = @is_triggered,
|
|
last_state = @last_state::jsonb
|
|
WHERE id = @id;
|
|
""";
|
|
|
|
await using var conn = new NpgsqlConnection(_connectionString);
|
|
await conn.OpenAsync(cancellationToken);
|
|
|
|
await using var cmd = new NpgsqlCommand(sql, conn);
|
|
cmd.Parameters.AddWithValue("id", alert.Id);
|
|
cmd.Parameters.AddWithValue("is_triggered", alert.IsTriggered);
|
|
cmd.Parameters.AddWithValue("last_state", alert.LastState);
|
|
|
|
await cmd.ExecuteNonQueryAsync(cancellationToken);
|
|
}
|
|
|
|
private async Task<List<AlertRecord>> QueryAlertsAsync(
|
|
string sql,
|
|
Action<NpgsqlCommand>? configureParams,
|
|
CancellationToken cancellationToken)
|
|
{
|
|
var results = new List<AlertRecord>();
|
|
|
|
await using var conn = new NpgsqlConnection(_connectionString);
|
|
await conn.OpenAsync(cancellationToken);
|
|
|
|
await using var cmd = new NpgsqlCommand(sql, conn);
|
|
configureParams?.Invoke(cmd);
|
|
|
|
await using var reader = await cmd.ExecuteReaderAsync(cancellationToken);
|
|
while (await reader.ReadAsync(cancellationToken))
|
|
{
|
|
results.Add(new AlertRecord
|
|
{
|
|
Id = reader.GetInt64(0),
|
|
Username = reader.GetString(1),
|
|
Symbol = reader.GetString(2),
|
|
AlertType = reader.GetString(3),
|
|
Parameters = reader.GetString(4),
|
|
LastState = reader.GetString(5),
|
|
IsTriggered = reader.GetBoolean(6),
|
|
CreatedAt = reader.GetDateTime(7)
|
|
});
|
|
}
|
|
|
|
return results;
|
|
}
|
|
}
|