107 lines
5.7 KiB
C#
107 lines
5.7 KiB
C#
using System.Data;
|
|
using System.Text.Json;
|
|
using Dapper;
|
|
using QuantEngine.Infrastructure.Data;
|
|
using QuantEngine.Core.Interfaces;
|
|
|
|
namespace QuantEngine.Infrastructure.Repositories
|
|
{
|
|
/// <summary>
|
|
/// PostgreSQL Dapper Repository implementation for History-First Operating Model.
|
|
/// Manages 3NF Data Integrity, Waterfall Auditing, and Shadow Ledger persistence.
|
|
/// </summary>
|
|
public class PostgresqlHistoryStore : IPostgresqlHistoryStore
|
|
{
|
|
private readonly IDbConnectionFactory _connectionFactory;
|
|
|
|
private static readonly IReadOnlyDictionary<string, string[]> DomainColumns = new Dictionary<string, string[]>
|
|
{
|
|
["market_raw_history"] = new[] { "source_id", "observed_at", "source_name", "instrument_id", "field_name", "field_value", "unit" },
|
|
["factor_version_history"] = new[] { "factor_id", "factor_version", "effective_from", "effective_to", "formula_id", "source_version" },
|
|
["factor_output_history"] = new[] { "factor_output_id", "observed_at", "factor_id", "factor_version", "output_value", "output_gate", "source_version" },
|
|
["decision_result_history"] = new[] { "decision_id", "decided_at", "instrument_id", "action", "gate", "score", "source_version" },
|
|
["market_vs_engine_gap_history"] = new[] { "gap_id", "observed_at", "instrument_id", "metric_name", "market_value", "engine_value", "gap_value", "gap_pct", "source_version" }
|
|
};
|
|
|
|
public PostgresqlHistoryStore(IDbConnectionFactory connectionFactory)
|
|
{
|
|
_connectionFactory = connectionFactory;
|
|
}
|
|
|
|
internal static IReadOnlyDictionary<string, string[]> GetDomainColumns() => DomainColumns;
|
|
|
|
internal static string BuildInsertSql(string domain, IReadOnlyList<string> insertColumns)
|
|
=> $@"INSERT INTO engine_history.{domain} ({string.Join(", ", insertColumns)}) VALUES ({string.Join(", ", insertColumns.Select(column => column == "provenance" ? "CAST(@provenance AS jsonb)" : $"@{column}"))})";
|
|
|
|
internal static string BuildSnapshotSql(string domain, int limit)
|
|
=> $@"SELECT * FROM engine_history.{domain} ORDER BY created_at DESC LIMIT @Limit";
|
|
|
|
public async Task<int> AppendAsync(string domain, IDictionary<string, object?> payload)
|
|
{
|
|
if (!DomainColumns.TryGetValue(domain, out var columns))
|
|
throw new ArgumentException($"Unsupported history domain: {domain}", nameof(domain));
|
|
|
|
using var conn = _connectionFactory.CreateConnection();
|
|
conn.Open();
|
|
|
|
var values = new DynamicParameters();
|
|
var insertColumns = new List<string>(columns.Length + 1);
|
|
foreach (var column in columns)
|
|
{
|
|
insertColumns.Add(column);
|
|
values.Add(column, payload.TryGetValue(column, out var value) ? value : null);
|
|
}
|
|
|
|
insertColumns.Add("provenance");
|
|
var provenance = payload.TryGetValue("provenance", out var provenanceValue) ? provenanceValue : new Dictionary<string, object?>();
|
|
values.Add("provenance", provenance is string s ? s : JsonSerializer.Serialize(provenance));
|
|
|
|
var sql = BuildInsertSql(domain, insertColumns);
|
|
return await conn.ExecuteAsync(sql, values);
|
|
}
|
|
|
|
public async Task<IReadOnlyList<IDictionary<string, object?>>> SnapshotAsync(string domain, int limit = 500)
|
|
{
|
|
if (!DomainColumns.ContainsKey(domain))
|
|
throw new ArgumentException($"Unsupported history domain: {domain}", nameof(domain));
|
|
|
|
using var conn = _connectionFactory.CreateConnection();
|
|
conn.Open();
|
|
|
|
var sql = BuildSnapshotSql(domain, limit);
|
|
var rows = await conn.QueryAsync(sql, new { Limit = limit });
|
|
return rows.Select(row => (IDictionary<string, object?>)row).ToList();
|
|
}
|
|
|
|
public async Task<long> RecordWaterfallExecutionAsync(string runId, string ticker, int rank, string stage, string action, int targetQty, decimal? targetPrice, decimal? bidAskSpreadBps, decimal? slippageBps, string status, string rationale)
|
|
{
|
|
using var conn = _connectionFactory.CreateConnection();
|
|
conn.Open();
|
|
|
|
const string sql = @"
|
|
INSERT INTO quantengine.order_waterfall_execution_history
|
|
(run_id, ticker, sell_priority_rank, waterfall_stage, action, target_qty, target_price, bid_ask_spread_bps, slippage_bps, status, rationale)
|
|
VALUES
|
|
(@runId, @ticker, @rank, @stage, @action, @targetQty, @targetPrice, @bidAskSpreadBps, @slippageBps, @status, @rationale)
|
|
RETURNING id;";
|
|
|
|
return await conn.ExecuteScalarAsync<long>(sql, new { runId, ticker, rank, stage, action, targetQty, targetPrice, bidAskSpreadBps, slippageBps, status, rationale });
|
|
}
|
|
|
|
public async Task<long> RecordShadowLedgerAsync(string runId, string ticker, string blockedGate, string blockedReason, decimal shadowPrice, int shadowQty, decimal? shadowTpPrice, decimal? shadowSlPrice)
|
|
{
|
|
using var conn = _connectionFactory.CreateConnection();
|
|
conn.Open();
|
|
|
|
const string sql = @"
|
|
INSERT INTO quantengine.shadow_ledger_history
|
|
(run_id, ticker, blocked_gate, blocked_reason, shadow_price, shadow_qty, shadow_tp_price, shadow_sl_price)
|
|
VALUES
|
|
(@runId, @ticker, @blockedGate, @blockedReason, @shadowPrice, @shadowQty, @shadowTpPrice, @shadowSlPrice)
|
|
RETURNING id;";
|
|
|
|
return await conn.ExecuteScalarAsync<long>(sql, new { runId, ticker, blockedGate, blockedReason, shadowPrice, shadowQty, shadowTpPrice, shadowSlPrice });
|
|
}
|
|
}
|
|
}
|