d18f6a7a67
Carries scheduledFor, catch-up policy, and maxCatchUp from the scheduler through the request model and transactional outbox. Evidence: targeted Release tests 5/5 passed; TRX SHA256 C1BF3EF274702305A29673D5B6A1C3A98D08B1716DA3CD8CB0EE710B5E6C12E6. Schedules remain disabled.
108 lines
4.2 KiB
C#
108 lines
4.2 KiB
C#
using System.Text.Json;
|
|
using Dapper;
|
|
using KArtSell.BuildingBlocks.Data;
|
|
using KArtSell.BuildingBlocks.Hashing;
|
|
using KArtSell.BuildingBlocks.Reliability;
|
|
using KArtSell.Modules.ModelOperations.Application;
|
|
|
|
namespace KArtSell.Modules.ModelOperations.Infrastructure;
|
|
|
|
public sealed class DapperModelOperationRequestRepository(
|
|
IDbConnectionFactory connectionFactory,
|
|
IOutboxWriter outboxWriter) : IModelOperationRequestRepository
|
|
{
|
|
private const string InsertSql = """
|
|
insert into evaluation.model_operation_request
|
|
(request_id, schedule_id, operation_code, scope_key, automation_mode, idempotency_key,
|
|
dataset_id, data_hash, model_version, config_version, code_sha, contract_version,
|
|
lifecycle_state, correlation_id, status, requested_at, scheduled_for)
|
|
values
|
|
(@RequestId, @ScheduleId, @OperationCode, @ScopeKey, @AutomationMode, @IdempotencyKey,
|
|
@DatasetId, @DataHash, @ModelVersion, @ConfigVersion, @CodeSha, @ContractVersion,
|
|
@LifecycleState, @CorrelationId, 'REQUESTED', @RequestedAt, @ScheduledFor)
|
|
on conflict (idempotency_key) do nothing;
|
|
""";
|
|
|
|
private const string InsertStatusEventSql = """
|
|
insert into evaluation.model_operation_status_event
|
|
(event_id, request_id, from_status, to_status, reason_code, payload_hash,
|
|
actor_type, correlation_id, occurred_at)
|
|
values
|
|
(@EventId, @RequestId, null, 'REQUESTED', null, @PayloadHash,
|
|
'SYSTEM', @CorrelationId, @OccurredAt);
|
|
""";
|
|
|
|
public async Task<bool> AddAsync(ModelOperationRequest request, CancellationToken cancellationToken)
|
|
{
|
|
await using var connection = await connectionFactory.OpenAsync(cancellationToken);
|
|
await using var transaction = await connection.BeginTransactionAsync(cancellationToken);
|
|
|
|
var affected = await connection.ExecuteAsync(new CommandDefinition(
|
|
InsertSql,
|
|
new
|
|
{
|
|
request.RequestId,
|
|
request.ScheduleId,
|
|
request.OperationCode,
|
|
request.ScopeKey,
|
|
request.AutomationMode,
|
|
request.IdempotencyKey,
|
|
request.ScheduledFor,
|
|
request.Context.VersionSet.DatasetId,
|
|
request.Context.VersionSet.DataHash,
|
|
request.Context.VersionSet.ModelVersion,
|
|
request.Context.VersionSet.ConfigVersion,
|
|
request.Context.VersionSet.CodeSha,
|
|
request.Context.VersionSet.ContractVersion,
|
|
request.Context.LifecycleState,
|
|
request.CorrelationId,
|
|
request.RequestedAt
|
|
},
|
|
transaction,
|
|
cancellationToken: cancellationToken));
|
|
|
|
if (affected == 1)
|
|
{
|
|
var payload = JsonSerializer.Serialize(new
|
|
{
|
|
requestId = request.RequestId,
|
|
request.OperationCode,
|
|
request.ScopeKey,
|
|
request.AutomationMode,
|
|
request.ScheduledFor,
|
|
request.CatchUpPolicy,
|
|
request.MaxCatchUp,
|
|
versionSet = request.Context.VersionSet,
|
|
boundary = "NO_AUTO_MODEL_MUTATION"
|
|
});
|
|
var payloadHash = ContentHasher.Sha256(payload);
|
|
|
|
await connection.ExecuteAsync(new CommandDefinition(
|
|
InsertStatusEventSql,
|
|
new
|
|
{
|
|
EventId = Guid.NewGuid(),
|
|
request.RequestId,
|
|
PayloadHash = payloadHash,
|
|
request.CorrelationId,
|
|
OccurredAt = request.RequestedAt
|
|
},
|
|
transaction,
|
|
cancellationToken: cancellationToken));
|
|
|
|
var message = new OutboxMessage(
|
|
request.RequestId,
|
|
"ModelOperationRequested",
|
|
1,
|
|
payload,
|
|
request.CorrelationId,
|
|
request.RequestedAt,
|
|
payloadHash);
|
|
await outboxWriter.AddAsync(connection, transaction, message, cancellationToken);
|
|
}
|
|
|
|
await transaction.CommitAsync(cancellationToken);
|
|
return affected == 1;
|
|
}
|
|
}
|