293 lines
9.7 KiB
C#
293 lines
9.7 KiB
C#
using System.Data;
|
|
using System.Text.Json;
|
|
using Dapper;
|
|
using KArtSell.BuildingBlocks.Data;
|
|
using KArtSell.BuildingBlocks.Hashing;
|
|
using KArtSell.BuildingBlocks.Reliability;
|
|
using KArtSell.BuildingBlocks.Time;
|
|
using KArtSell.Modules.SignalEngine.Domain;
|
|
|
|
namespace KArtSell.Modules.SignalEngine.Application;
|
|
|
|
public sealed class SellDecisionService(
|
|
IDbConnectionFactory connectionFactory,
|
|
ISellDecisionContextReader contextReader,
|
|
SellPolicyChain policyChain,
|
|
IOutboxWriter outboxWriter,
|
|
IClock clock) : ISellDecisionService
|
|
{
|
|
private static readonly JsonSerializerOptions JsonOptions = new(JsonSerializerDefaults.Web);
|
|
|
|
private const string SelectExistingSql = """
|
|
select
|
|
decision_id as DecisionId,
|
|
action as Action,
|
|
sell_ratio_of_lot as SellRatioOfLot,
|
|
target_security_portfolio_weight_after as TargetSecurityPortfolioWeightAfter,
|
|
policy_id as PolicyId,
|
|
reason_code as ReasonCode,
|
|
evidence_id as EvidenceId,
|
|
dataset_id as DatasetId,
|
|
model_version as ModelVersion,
|
|
config_version as ConfigVersion,
|
|
code_sha as CodeSha,
|
|
decision_contract_version as DecisionContractVersion,
|
|
policy_trace_schema_version as PolicyTraceSchemaVersion,
|
|
reentry_eligible as ReentryEligible,
|
|
policy_trace_json::text as PolicyTraceJson,
|
|
created_at as CreatedAt
|
|
from signal_engine.signal_decision
|
|
where idempotency_key = @IdempotencyKey;
|
|
""";
|
|
|
|
private const string InsertSql = """
|
|
insert into signal_engine.signal_decision
|
|
(
|
|
decision_id,
|
|
position_lot_id,
|
|
cycle_id,
|
|
action,
|
|
sell_ratio_of_lot,
|
|
target_security_portfolio_weight_after,
|
|
policy_id,
|
|
priority,
|
|
reason_code,
|
|
evidence_id,
|
|
dataset_id,
|
|
model_version,
|
|
config_version,
|
|
code_sha,
|
|
decision_contract_version,
|
|
policy_trace_schema_version,
|
|
reentry_eligible,
|
|
policy_trace_json,
|
|
idempotency_key,
|
|
correlation_id,
|
|
created_at,
|
|
decision_hash
|
|
)
|
|
values
|
|
(
|
|
@DecisionId,
|
|
@PositionLotId,
|
|
@CycleId,
|
|
@Action,
|
|
@SellRatioOfLot,
|
|
@TargetSecurityPortfolioWeightAfter,
|
|
@PolicyId,
|
|
@Priority,
|
|
@ReasonCode,
|
|
@EvidenceId,
|
|
@DatasetId,
|
|
@ModelVersion,
|
|
@ConfigVersion,
|
|
@CodeSha,
|
|
@DecisionContractVersion,
|
|
@PolicyTraceSchemaVersion,
|
|
@ReentryEligible,
|
|
cast(@PolicyTraceJson as jsonb),
|
|
@IdempotencyKey,
|
|
@CorrelationId,
|
|
@CreatedAt,
|
|
@DecisionHash
|
|
)
|
|
on conflict (idempotency_key) do nothing;
|
|
""";
|
|
|
|
public async Task<GeneratedSellDecision?> GenerateAsync(
|
|
GenerateSellDecisionCommand command,
|
|
CancellationToken cancellationToken)
|
|
{
|
|
var context = await contextReader.GetApprovedAsync(
|
|
command.PositionLotId,
|
|
command.AsOf,
|
|
cancellationToken);
|
|
|
|
if (context is null)
|
|
{
|
|
return null;
|
|
}
|
|
|
|
await using var connection = await connectionFactory.OpenAsync(cancellationToken);
|
|
await using var transaction = await connection.BeginTransactionAsync(
|
|
IsolationLevel.Serializable,
|
|
cancellationToken);
|
|
|
|
var existing = await connection.QuerySingleOrDefaultAsync<PersistedDecision>(
|
|
new CommandDefinition(
|
|
SelectExistingSql,
|
|
new { command.IdempotencyKey },
|
|
transaction,
|
|
cancellationToken: cancellationToken));
|
|
|
|
if (existing is not null)
|
|
{
|
|
await transaction.CommitAsync(cancellationToken);
|
|
return existing.ToResult(replayed: true);
|
|
}
|
|
|
|
var decision = policyChain.Evaluate(context.ToDomainInput());
|
|
var now = clock.UtcNow;
|
|
var decisionId = Guid.NewGuid();
|
|
var policyTraceJson = JsonSerializer.Serialize(decision.PolicyTrace, JsonOptions);
|
|
|
|
var decisionCanonicalJson = JsonSerializer.Serialize(new
|
|
{
|
|
DecisionId = decisionId,
|
|
context.PositionLotId,
|
|
context.CycleId,
|
|
decision.Action,
|
|
decision.SellRatioOfLot,
|
|
decision.TargetSecurityPortfolioWeightAfter,
|
|
decision.PolicyId,
|
|
decision.Priority,
|
|
decision.ReasonCode,
|
|
decision.EvidenceId,
|
|
decision.DatasetId,
|
|
decision.ModelVersion,
|
|
decision.ConfigVersion,
|
|
decision.CodeSha,
|
|
DecisionContractVersion = SellPolicyContract.DecisionContractVersion,
|
|
PolicyTraceSchemaVersion = SellPolicyContract.PolicyTraceSchemaVersion,
|
|
decision.ReentryEligible,
|
|
PolicyTrace = decision.PolicyTrace,
|
|
CreatedAt = now
|
|
}, JsonOptions);
|
|
var decisionHash = ContentHasher.Sha256(decisionCanonicalJson);
|
|
|
|
var insert = new
|
|
{
|
|
DecisionId = decisionId,
|
|
context.PositionLotId,
|
|
context.CycleId,
|
|
Action = decision.Action.ToString().ToUpperInvariant(),
|
|
decision.SellRatioOfLot,
|
|
decision.TargetSecurityPortfolioWeightAfter,
|
|
decision.PolicyId,
|
|
decision.Priority,
|
|
decision.ReasonCode,
|
|
decision.EvidenceId,
|
|
decision.DatasetId,
|
|
decision.ModelVersion,
|
|
decision.ConfigVersion,
|
|
decision.CodeSha,
|
|
DecisionContractVersion = SellPolicyContract.DecisionContractVersion,
|
|
PolicyTraceSchemaVersion = SellPolicyContract.PolicyTraceSchemaVersion,
|
|
decision.ReentryEligible,
|
|
PolicyTraceJson = policyTraceJson,
|
|
command.IdempotencyKey,
|
|
command.CorrelationId,
|
|
CreatedAt = now,
|
|
DecisionHash = decisionHash
|
|
};
|
|
|
|
var affected = await connection.ExecuteAsync(
|
|
new CommandDefinition(
|
|
InsertSql,
|
|
insert,
|
|
transaction,
|
|
cancellationToken: cancellationToken));
|
|
|
|
if (affected == 0)
|
|
{
|
|
var replay = await connection.QuerySingleAsync<PersistedDecision>(
|
|
new CommandDefinition(
|
|
SelectExistingSql,
|
|
new { command.IdempotencyKey },
|
|
transaction,
|
|
cancellationToken: cancellationToken));
|
|
await transaction.CommitAsync(cancellationToken);
|
|
return replay.ToResult(replayed: true);
|
|
}
|
|
|
|
var payload = JsonSerializer.Serialize(new
|
|
{
|
|
DecisionId = decisionId,
|
|
context.PositionLotId,
|
|
decision.PolicyId,
|
|
Action = decision.Action.ToString(),
|
|
decision.SellRatioOfLot,
|
|
decision.TargetSecurityPortfolioWeightAfter,
|
|
decision.EvidenceId,
|
|
decision.DatasetId,
|
|
DecisionContractVersion = SellPolicyContract.DecisionContractVersion,
|
|
PolicyTraceSchemaVersion = SellPolicyContract.PolicyTraceSchemaVersion,
|
|
CreatedAt = now
|
|
}, JsonOptions);
|
|
|
|
await outboxWriter.AddAsync(
|
|
connection,
|
|
transaction,
|
|
new OutboxMessage(
|
|
Guid.NewGuid(),
|
|
"SignalDecisionCreated",
|
|
2,
|
|
payload,
|
|
command.CorrelationId,
|
|
now,
|
|
ContentHasher.Sha256(payload)),
|
|
cancellationToken);
|
|
|
|
await transaction.CommitAsync(cancellationToken);
|
|
|
|
return new GeneratedSellDecision(
|
|
decisionId,
|
|
decision.Action.ToString(),
|
|
decision.SellRatioOfLot,
|
|
decision.TargetSecurityPortfolioWeightAfter,
|
|
decision.PolicyId,
|
|
decision.ReasonCode,
|
|
decision.EvidenceId,
|
|
decision.DatasetId,
|
|
decision.ModelVersion,
|
|
decision.ConfigVersion,
|
|
decision.CodeSha,
|
|
SellPolicyContract.DecisionContractVersion,
|
|
SellPolicyContract.PolicyTraceSchemaVersion,
|
|
decision.ReentryEligible,
|
|
decision.PolicyTrace,
|
|
now,
|
|
false);
|
|
}
|
|
|
|
private sealed record PersistedDecision(
|
|
Guid DecisionId,
|
|
string Action,
|
|
decimal SellRatioOfLot,
|
|
decimal TargetSecurityPortfolioWeightAfter,
|
|
string PolicyId,
|
|
string ReasonCode,
|
|
string EvidenceId,
|
|
string DatasetId,
|
|
string ModelVersion,
|
|
string ConfigVersion,
|
|
string CodeSha,
|
|
string DecisionContractVersion,
|
|
int PolicyTraceSchemaVersion,
|
|
bool ReentryEligible,
|
|
string PolicyTraceJson,
|
|
DateTimeOffset CreatedAt)
|
|
{
|
|
public GeneratedSellDecision ToResult(bool replayed)
|
|
=> new(
|
|
DecisionId,
|
|
Action,
|
|
SellRatioOfLot,
|
|
TargetSecurityPortfolioWeightAfter,
|
|
PolicyId,
|
|
ReasonCode,
|
|
EvidenceId,
|
|
DatasetId,
|
|
ModelVersion,
|
|
ConfigVersion,
|
|
CodeSha,
|
|
DecisionContractVersion,
|
|
PolicyTraceSchemaVersion,
|
|
ReentryEligible,
|
|
JsonSerializer.Deserialize<List<PolicyTraceEntry>>(PolicyTraceJson, JsonOptions)
|
|
?? new List<PolicyTraceEntry>(),
|
|
CreatedAt,
|
|
replayed);
|
|
}
|
|
}
|