using System.Collections.Generic; using System.Threading.Tasks; using QuantEngine.Core.Domain; using QuantEngine.Core.Interfaces; namespace QuantEngine.Application.Services { public class HistoryIngestionService { private readonly IPostgresqlHistoryStore _store; public HistoryIngestionService(IPostgresqlHistoryStore store) { _store = store; } public Task AppendDecisionAsync(IDictionary payload) => _store.AppendAsync("decision_result_history", RequirePayload(payload)); public Task AppendFactorOutputAsync(IDictionary payload) => _store.AppendAsync("factor_output_history", RequirePayload(payload)); public Task AppendMarketRawAsync(IDictionary payload) => _store.AppendAsync("market_raw_history", RequirePayload(payload)); public Task AppendGapAsync(IDictionary payload) => _store.AppendAsync("market_vs_engine_gap_history", RequirePayload(payload)); public Task AppendDecisionAsync( FinalDecisionResult decision, SellDecisionResult? sellDecision = null, TimingDecisionResult? timingDecision = null, string? instrumentId = null, string? sourceVersion = null, string? gate = null) { ArgumentNullException.ThrowIfNull(decision); var normalizedInstrumentId = NormalizeOptional(instrumentId); var normalizedSourceVersion = NormalizeOptional(sourceVersion) ?? RequireValue(decision.DecisionSource, nameof(decision.DecisionSource)); var normalizedGate = NormalizeOptional(gate) ?? (string.IsNullOrWhiteSpace(sellDecision?.Validation) ? "PASS" : sellDecision.Validation!.Trim()); var normalizedAction = RequireValue(decision.FinalAction, nameof(decision.FinalAction)); var payload = new Dictionary { ["decision_id"] = Guid.NewGuid().ToString("N"), ["decided_at"] = DateTimeOffset.UtcNow, ["instrument_id"] = normalizedInstrumentId ?? string.Empty, ["action"] = normalizedAction, ["gate"] = normalizedGate, ["score"] = decision.PriorityScore, ["source_version"] = normalizedSourceVersion, ["provenance"] = new Dictionary { ["final_action"] = normalizedAction, ["action_priority"] = decision.ActionPriority, ["priority_score"] = decision.PriorityScore, ["decision_source"] = decision.DecisionSource, ["sell_action"] = NormalizeOptional(sellDecision?.Action), ["sell_validation"] = NormalizeOptional(sellDecision?.Validation), ["timing_action"] = NormalizeOptional(timingDecision?.Action), ["timing_reason"] = NormalizeOptional(timingDecision?.Reason) } }; return _store.AppendAsync("decision_result_history", payload); } public Task AppendFactorOutputAsync( string factorId, string factorVersion, double outputValue, string outputGate, string? sourceVersion = null, DateTimeOffset? observedAt = null) { factorId = RequireValue(factorId, nameof(factorId)); factorVersion = RequireValue(factorVersion, nameof(factorVersion)); outputGate = RequireValue(outputGate, nameof(outputGate)); sourceVersion = NormalizeOptional(sourceVersion) ?? factorVersion; var payload = new Dictionary { ["factor_output_id"] = Guid.NewGuid().ToString("N"), ["observed_at"] = observedAt ?? DateTimeOffset.UtcNow, ["factor_id"] = factorId, ["factor_version"] = factorVersion, ["output_value"] = outputValue, ["output_gate"] = outputGate, ["source_version"] = sourceVersion, ["provenance"] = new Dictionary { ["factor_id"] = factorId, ["factor_version"] = factorVersion, ["output_value"] = outputValue, ["output_gate"] = outputGate, ["source_version"] = sourceVersion ?? factorVersion } }; return _store.AppendAsync("factor_output_history", payload); } private static string RequireValue(string value, string parameterName) { if (string.IsNullOrWhiteSpace(value)) { throw new ArgumentException("Value is required.", parameterName); } return value.Trim(); } private static string? NormalizeOptional(string? value) => string.IsNullOrWhiteSpace(value) ? null : value.Trim(); private static IDictionary RequirePayload(IDictionary payload) { ArgumentNullException.ThrowIfNull(payload); return payload; } } }