refactor(dotnet): centralize runtime audit trail
This commit is contained in:
@@ -0,0 +1,6 @@
|
||||
namespace QuantEngine.Application.Interfaces;
|
||||
|
||||
public interface IRuntimeAuditTrailService
|
||||
{
|
||||
void Append<T>(string category, string key, T payload);
|
||||
}
|
||||
@@ -1,6 +1,7 @@
|
||||
using System.Text.Json;
|
||||
using QuantEngine.Core.Domain;
|
||||
using QuantEngine.Core.Interfaces;
|
||||
using QuantEngine.Application.Interfaces;
|
||||
|
||||
namespace QuantEngine.Application.Services;
|
||||
|
||||
@@ -15,34 +16,12 @@ public sealed record FactorComputationAudit(
|
||||
public sealed class FactorComputationService
|
||||
{
|
||||
private readonly HistoryIngestionService _history;
|
||||
private readonly string _auditRoot;
|
||||
private readonly IRuntimeAuditTrailService _auditTrail;
|
||||
|
||||
public FactorComputationService(HistoryIngestionService history)
|
||||
public FactorComputationService(HistoryIngestionService history, IRuntimeAuditTrailService auditTrail)
|
||||
{
|
||||
_history = history;
|
||||
_auditRoot = FindRepoTempRoot();
|
||||
}
|
||||
|
||||
private static string FindRepoTempRoot()
|
||||
{
|
||||
var current = new DirectoryInfo(AppContext.BaseDirectory);
|
||||
while (current != null)
|
||||
{
|
||||
if (Directory.Exists(Path.Combine(current.FullName, ".git")))
|
||||
{
|
||||
return Path.Combine(current.FullName, "Temp", "factor_audit");
|
||||
}
|
||||
current = current.Parent;
|
||||
}
|
||||
|
||||
return Path.Combine(Directory.GetCurrentDirectory(), "Temp", "factor_audit");
|
||||
}
|
||||
|
||||
private void AppendAudit(FactorComputationAudit audit)
|
||||
{
|
||||
Directory.CreateDirectory(_auditRoot);
|
||||
var path = Path.Combine(_auditRoot, $"{audit.Ticker}.jsonl");
|
||||
File.AppendAllText(path, JsonSerializer.Serialize(audit, new JsonSerializerOptions { WriteIndented = false }) + Environment.NewLine);
|
||||
_auditTrail = auditTrail;
|
||||
}
|
||||
|
||||
public FactorOutputs Compute(
|
||||
@@ -53,7 +32,7 @@ public sealed class FactorComputationService
|
||||
{
|
||||
var computedAt = DateTimeOffset.UtcNow;
|
||||
var outputs = FactorCalculator.CalculateFactors(stockBars, indexBars);
|
||||
AppendAudit(new FactorComputationAudit(ticker, stockBars.Count, indexBars.Count, "SUCCEEDED", computedAt, sourceVersion));
|
||||
_auditTrail.Append("factor_audit", ticker, new FactorComputationAudit(ticker, stockBars.Count, indexBars.Count, "SUCCEEDED", computedAt, sourceVersion));
|
||||
return outputs;
|
||||
}
|
||||
|
||||
@@ -71,6 +50,6 @@ public sealed class FactorComputationService
|
||||
await _history.AppendFactorOutputAsync("stdev_20d", sourceVersion, outputs.StDev20D, "PASS", sourceVersion, when);
|
||||
await _history.AppendFactorOutputAsync("beta_60d", sourceVersion, outputs.Beta60D, "PASS", sourceVersion, when);
|
||||
await _history.AppendFactorOutputAsync("rs_20d", sourceVersion, outputs.Rs20D, "PASS", sourceVersion, when);
|
||||
AppendAudit(new FactorComputationAudit(ticker, 0, 0, "PERSISTED", when, sourceVersion));
|
||||
_auditTrail.Append("factor_audit", ticker, new FactorComputationAudit(ticker, 0, 0, "PERSISTED", when, sourceVersion));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -17,42 +17,22 @@ public class KisDataCollectionOrchestrator : ICollectionOrchestrator
|
||||
private readonly PriceDataNormalizer _normalizer;
|
||||
private readonly SourcePriorityResolver _priorityResolver;
|
||||
private readonly ILogger<KisDataCollectionOrchestrator> _logger;
|
||||
private readonly string _auditRoot;
|
||||
private readonly IRuntimeAuditTrailService _auditTrail;
|
||||
|
||||
public KisDataCollectionOrchestrator(
|
||||
IKisApiClient kisApiClient,
|
||||
ICollectionRepository repository,
|
||||
PriceDataNormalizer normalizer,
|
||||
SourcePriorityResolver priorityResolver,
|
||||
ILogger<KisDataCollectionOrchestrator> logger)
|
||||
ILogger<KisDataCollectionOrchestrator> logger,
|
||||
IRuntimeAuditTrailService auditTrail)
|
||||
{
|
||||
_kisApiClient = kisApiClient;
|
||||
_repository = repository;
|
||||
_normalizer = normalizer;
|
||||
_priorityResolver = priorityResolver;
|
||||
_logger = logger;
|
||||
_auditRoot = FindRepoTempRoot();
|
||||
}
|
||||
|
||||
private static string FindRepoTempRoot()
|
||||
{
|
||||
var current = new DirectoryInfo(AppContext.BaseDirectory);
|
||||
while (current != null)
|
||||
{
|
||||
if (Directory.Exists(Path.Combine(current.FullName, ".git")))
|
||||
{
|
||||
return Path.Combine(current.FullName, "Temp", "collection_audit");
|
||||
}
|
||||
current = current.Parent;
|
||||
}
|
||||
return Path.Combine(Directory.GetCurrentDirectory(), "Temp", "collection_audit");
|
||||
}
|
||||
|
||||
private void AppendAudit(CollectionExecutionAudit audit)
|
||||
{
|
||||
Directory.CreateDirectory(_auditRoot);
|
||||
var path = Path.Combine(_auditRoot, $"{audit.RunId}.jsonl");
|
||||
File.AppendAllText(path, JsonSerializer.Serialize(audit, new JsonSerializerOptions { WriteIndented = false }) + Environment.NewLine);
|
||||
_auditTrail = auditTrail;
|
||||
}
|
||||
|
||||
public async Task<CollectionRunResult> RunCollectionAsync(string runId, string account, List<string> tickers)
|
||||
@@ -70,7 +50,7 @@ public class KisDataCollectionOrchestrator : ICollectionOrchestrator
|
||||
try
|
||||
{
|
||||
_logger.LogInformation("Starting collection run {RunId}", runId);
|
||||
AppendAudit(new CollectionExecutionAudit(runId, "RUNNING", DateTimeOffset.UtcNow, null, 0, 0, "started"));
|
||||
_auditTrail.Append("collection_audit", runId, new CollectionExecutionAudit(runId, "RUNNING", DateTimeOffset.UtcNow, null, 0, 0, "started"));
|
||||
|
||||
var kisSource = new KisApiPriceSource(_kisApiClient);
|
||||
var rows = new List<Dictionary<string, object>>();
|
||||
@@ -183,7 +163,7 @@ public class KisDataCollectionOrchestrator : ICollectionOrchestrator
|
||||
result.SourceCounts = sourceCounts;
|
||||
result.Rows = rows;
|
||||
result.Errors = errors;
|
||||
AppendAudit(new CollectionExecutionAudit(runId, result.Status, DateTimeOffset.Parse(startedAt), DateTimeOffset.Parse(finishedAt), result.SuccessCount, result.ErrorCount, "finished"));
|
||||
_auditTrail.Append("collection_audit", runId, new CollectionExecutionAudit(runId, result.Status, DateTimeOffset.Parse(startedAt), DateTimeOffset.Parse(finishedAt), result.SuccessCount, result.ErrorCount, "finished"));
|
||||
|
||||
// Save run record
|
||||
await _repository.SaveRunAsync(new CollectionRunRecord(
|
||||
@@ -231,7 +211,7 @@ public class KisDataCollectionOrchestrator : ICollectionOrchestrator
|
||||
result.Status = "FAILED";
|
||||
result.FinishedAt = DataNormalizationHelper.KstNowIso();
|
||||
result.ErrorMessage = ex.Message;
|
||||
AppendAudit(new CollectionExecutionAudit(runId, result.Status, DateTimeOffset.Parse(startedAt), DateTimeOffset.Parse(result.FinishedAt), result.SuccessCount, result.ErrorCount, ex.Message));
|
||||
_auditTrail.Append("collection_audit", runId, new CollectionExecutionAudit(runId, result.Status, DateTimeOffset.Parse(startedAt), DateTimeOffset.Parse(result.FinishedAt), result.SuccessCount, result.ErrorCount, ex.Message));
|
||||
return result;
|
||||
}
|
||||
}
|
||||
@@ -385,4 +365,3 @@ public class KisDataCollectionOrchestrator : ICollectionOrchestrator
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -0,0 +1,37 @@
|
||||
using System.Text.Json;
|
||||
using QuantEngine.Application.Interfaces;
|
||||
|
||||
namespace QuantEngine.Application.Services;
|
||||
|
||||
public sealed class RuntimeAuditTrailService : IRuntimeAuditTrailService
|
||||
{
|
||||
private readonly string _auditRoot;
|
||||
|
||||
public RuntimeAuditTrailService()
|
||||
{
|
||||
_auditRoot = FindRepoTempRoot();
|
||||
}
|
||||
|
||||
public void Append<T>(string category, string key, T payload)
|
||||
{
|
||||
var root = Path.Combine(_auditRoot, category);
|
||||
Directory.CreateDirectory(root);
|
||||
var path = Path.Combine(root, $"{key}.jsonl");
|
||||
File.AppendAllText(path, JsonSerializer.Serialize(payload, new JsonSerializerOptions { WriteIndented = false }) + Environment.NewLine);
|
||||
}
|
||||
|
||||
private static string FindRepoTempRoot()
|
||||
{
|
||||
var current = new DirectoryInfo(AppContext.BaseDirectory);
|
||||
while (current != null)
|
||||
{
|
||||
if (Directory.Exists(Path.Combine(current.FullName, ".git")))
|
||||
{
|
||||
return Path.Combine(current.FullName, "Temp");
|
||||
}
|
||||
current = current.Parent;
|
||||
}
|
||||
|
||||
return Path.Combine(Directory.GetCurrentDirectory(), "Temp");
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user