Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,65 @@
package com.idea2strategy.backend.batch;

import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.json.JsonMapper;
import com.idea2strategy.backend.application.competition.FinalRoomResultService;
import com.idea2strategy.backend.application.competition.RoomFinalizationService;
import com.idea2strategy.backend.application.competition.ScoringEvidenceService;
import com.idea2strategy.backend.application.competition.ScoringTemplateCatalogService;
import com.idea2strategy.backend.application.competition.VirtualLiquidationService;
import com.idea2strategy.backend.persistence.competition.CanonicalVirtualLiquidationQuoteAdapter;
import com.idea2strategy.backend.persistence.competition.FinalRoomResultJooqAdapter;
import com.idea2strategy.backend.persistence.competition.RoomFinalizationWorkJooqAdapter;
import com.idea2strategy.backend.persistence.competition.ScoringEvidenceJooqAdapter;
import com.idea2strategy.backend.persistence.competition.ScoringTemplateCatalogJooqQueryAdapter;
import com.idea2strategy.backend.persistence.competition.VirtualLiquidationJooqAdapter;
import java.time.Clock;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Import;
import org.springframework.scheduling.annotation.EnableScheduling;

@Configuration(proxyBeanMethods = false)
@EnableScheduling
@ConditionalOnProperty(
name = "idea2strategy.batch.room-finalization.enabled",
havingValue = "true",
matchIfMissing = true)
@Import({
RoomFinalizationWorkJooqAdapter.class,
VirtualLiquidationJooqAdapter.class,
CanonicalVirtualLiquidationQuoteAdapter.class,
ScoringEvidenceJooqAdapter.class,
FinalRoomResultJooqAdapter.class,
ScoringTemplateCatalogJooqQueryAdapter.class
})
class RoomFinalizationBatchConfiguration {
@Bean
RoomFinalizationService roomFinalizationService(
RoomFinalizationWorkJooqAdapter work,
VirtualLiquidationJooqAdapter liquidationStore,
CanonicalVirtualLiquidationQuoteAdapter quote,
ScoringEvidenceJooqAdapter evidence,
FinalRoomResultJooqAdapter results,
ScoringTemplateCatalogJooqQueryAdapter templates) {
Clock clock = Clock.systemUTC();
ObjectMapper mapper = JsonMapper.builder().build();
return new RoomFinalizationService(
work,
new VirtualLiquidationService(liquidationStore, quote, liquidationStore),
new ScoringEvidenceService(evidence),
new FinalRoomResultService(results, clock),
new ScoringTemplateCatalogService(templates, clock, mapper),
clock,
mapper);
}

@Bean
RoomFinalizationBatchRunner roomFinalizationBatchRunner(
RoomFinalizationService service,
@Value("${idea2strategy.batch.room-finalization.batch-size:100}") int batchSize) {
return new RoomFinalizationBatchRunner(service, batchSize);
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
package com.idea2strategy.backend.batch;

import com.idea2strategy.backend.application.competition.RoomFinalizationService;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.scheduling.annotation.Scheduled;

class RoomFinalizationBatchRunner {
private static final Logger log = LoggerFactory.getLogger(RoomFinalizationBatchRunner.class);
private final RoomFinalizationService service;
private final int batchSize;

RoomFinalizationBatchRunner(RoomFinalizationService service, int batchSize) {
this.service = service;
this.batchSize = batchSize;
}

@Scheduled(fixedDelayString = "${idea2strategy.batch.room-finalization.fixed-delay:PT10S}")
void run() {
var report = service.run(batchSize);
log.info(
"Room finalization batch completed: roomsAttempted={}, roomsFinalized={}, "
+ "participationsFinalized={}, failures={}, observedAt={}",
report.roomsAttempted(), report.roomsFinalized(), report.participationsFinalized(),
report.failures().size(), report.observedAt());
report.failures().forEach(failure -> log.warn(
"Room finalization remains retryable: roomId={}, reason={}",
failure.roomId(), failure.reason()));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
"idea2strategy.batch.pending-registration-cleanup.enabled=false",
"idea2strategy.batch.room-schedule-transition.enabled=false",
"idea2strategy.batch.room-evaluation-start.enabled=false",
"idea2strategy.batch.room-finalization.enabled=false",
"idea2strategy.batch.private-continuation-transition.enabled=false",
"idea2strategy.batch.post-evaluation-stop-transition.enabled=false"
})
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,8 @@ public final class DatabaseAccessPolicy {
private static final Set<QualifiedTable> BATCH_UPDATED_TABLES = Set.of(
new QualifiedTable("competition", "rooms"),
new QualifiedTable("competition", "participations"),
new QualifiedTable("competition", "live_evaluation_segments"),
new QualifiedTable("competition", "leaderboard_snapshots"),
new QualifiedTable("bot", "bots"),
new QualifiedTable("bot", "continuation_deadlines"),
new QualifiedTable("identity", "accounts"),
Expand All @@ -93,6 +95,9 @@ public final class DatabaseAccessPolicy {
new QualifiedTable("competition", "participation_events"),
new QualifiedTable("competition", "backtest_period_runs"),
new QualifiedTable("competition", "live_evaluation_segments"),
new QualifiedTable("competition", "leaderboard_snapshots"),
new QualifiedTable("competition", "leaderboard_entries"),
new QualifiedTable("competition", "room_final_access_grants"),
new QualifiedTable("bot", "continuation_deadlines"),
new QualifiedTable("backtest", "runs"),
new QualifiedTable("identity", "account_closure_runs"),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -276,6 +276,8 @@ void grantsTheBatchRoleTheWritesItsScheduledJobsPerform() {
for (var target : List.of(
new DatabaseAccessPolicy.QualifiedTable("competition", "rooms"),
new DatabaseAccessPolicy.QualifiedTable("competition", "participations"),
new DatabaseAccessPolicy.QualifiedTable("competition", "live_evaluation_segments"),
new DatabaseAccessPolicy.QualifiedTable("competition", "leaderboard_snapshots"),
new DatabaseAccessPolicy.QualifiedTable("bot", "bots"),
new DatabaseAccessPolicy.QualifiedTable("bot", "continuation_deadlines"))) {
assertTrue(
Expand All @@ -291,6 +293,9 @@ void grantsTheBatchRoleTheWritesItsScheduledJobsPerform() {
new DatabaseAccessPolicy.QualifiedTable("competition", "participation_events"),
new DatabaseAccessPolicy.QualifiedTable("competition", "backtest_period_runs"),
new DatabaseAccessPolicy.QualifiedTable("competition", "live_evaluation_segments"),
new DatabaseAccessPolicy.QualifiedTable("competition", "leaderboard_snapshots"),
new DatabaseAccessPolicy.QualifiedTable("competition", "leaderboard_entries"),
new DatabaseAccessPolicy.QualifiedTable("competition", "room_final_access_grants"),
new DatabaseAccessPolicy.QualifiedTable("bot", "continuation_deadlines"),
new DatabaseAccessPolicy.QualifiedTable("backtest", "runs"))) {
assertTrue(
Expand Down Expand Up @@ -333,6 +338,17 @@ void grantsTheBatchRoleTheWritesItsScheduledJobsPerform() {
target.table()),
"batch has no write path into " + target.schema() + "." + target.table());
}

// The batch can only EXPIRE an existing sanction. Manual APPLY commands are rejected by
// DeadlineBatchConfiguration, so granting INSERT here would widen the runtime role beyond
// the only sanction mutation the scheduled job can perform.
assertFalse(
DatabaseAccessPolicy.allows(
DatabaseAccessPolicy.ApplicationRole.BATCH,
DatabaseAccessPolicy.Access.INSERT,
"identity",
"account_sanctions"),
"sanction expiry updates an existing row and must not create a sanction");
}
@Test
void grantsTheBacktestRoleTheBotReadsItsExecutorPerforms() {
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
package com.idea2strategy.backend.application.competition;

import java.util.Objects;
import java.util.UUID;

public record RoomFinalizationCandidateSource(
UUID participationId,
UUID evaluationSegmentId,
UUID performanceSnapshotId,
long scheduledEvaluationSeconds,
long actualOperationSeconds,
int actualFillCount,
long baseRequiredOperationSeconds,
int baseRequiredFillCount) {
public RoomFinalizationCandidateSource {
Objects.requireNonNull(participationId, "participationId");
Objects.requireNonNull(evaluationSegmentId, "evaluationSegmentId");
Objects.requireNonNull(performanceSnapshotId, "performanceSnapshotId");
if (scheduledEvaluationSeconds <= 0
|| actualOperationSeconds < 0
|| actualFillCount < 0
|| baseRequiredOperationSeconds < 0
|| baseRequiredFillCount < 0) {
throw new IllegalArgumentException("room finalization eligibility evidence is invalid");
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
package com.idea2strategy.backend.application.competition;

import java.util.Objects;
import java.util.UUID;

public record RoomFinalizationFailure(UUID roomId, String reason) {
public RoomFinalizationFailure {
Objects.requireNonNull(roomId, "roomId");
Objects.requireNonNull(reason, "reason");
if (reason.isBlank()) {
throw new IllegalArgumentException("reason must not be blank");
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
package com.idea2strategy.backend.application.competition;

import java.time.Instant;
import java.util.List;
import java.util.Objects;

public record RoomFinalizationReport(
Instant observedAt,
int roomsAttempted,
int roomsFinalized,
int participationsFinalized,
List<RoomFinalizationFailure> failures) {
public RoomFinalizationReport {
Objects.requireNonNull(observedAt, "observedAt");
failures = List.copyOf(Objects.requireNonNull(failures, "failures"));
if (roomsAttempted < 0 || roomsFinalized < 0 || participationsFinalized < 0
|| roomsFinalized + failures.size() > roomsAttempted) {
throw new IllegalArgumentException("room finalization report counts are invalid");
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,133 @@
package com.idea2strategy.backend.application.competition;

import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import java.time.Clock;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.UUID;

/** Completes due live-paper evidence and publishes one immutable FINAL leaderboard per room. */
public final class RoomFinalizationService {
private final RoomFinalizationWorkPort workPort;
private final VirtualLiquidationService liquidationService;
private final ScoringEvidenceService evidenceService;
private final FinalRoomResultService resultService;
private final ScoringTemplateCatalogService templateCatalog;
private final OfficialScoringCalculator scoring = new OfficialScoringCalculator();
private final Clock clock;
private final ObjectMapper mapper;

public RoomFinalizationService(
RoomFinalizationWorkPort workPort,
VirtualLiquidationService liquidationService,
ScoringEvidenceService evidenceService,
FinalRoomResultService resultService,
ScoringTemplateCatalogService templateCatalog,
Clock clock,
ObjectMapper mapper) {
this.workPort = Objects.requireNonNull(workPort, "workPort");
this.liquidationService = Objects.requireNonNull(liquidationService, "liquidationService");
this.evidenceService = Objects.requireNonNull(evidenceService, "evidenceService");
this.resultService = Objects.requireNonNull(resultService, "resultService");
this.templateCatalog = Objects.requireNonNull(templateCatalog, "templateCatalog");
this.clock = Objects.requireNonNull(clock, "clock");
this.mapper = Objects.requireNonNull(mapper, "mapper");
}

public RoomFinalizationReport run(int limit) {
if (limit <= 0) {
throw new IllegalArgumentException("limit must be positive");
}
var observedAt = clock.instant();
List<UUID> roomIds = workPort.findDueRoomIds(observedAt, limit);
List<RoomFinalizationFailure> failures = new ArrayList<>();
int roomsFinalized = 0;
int participationsFinalized = 0;
for (UUID roomId : roomIds) {
try {
for (VirtualLiquidationRequest request : workPort.findPendingLiquidations(roomId)) {
liquidationService.finalizeEvaluation(request);
workPort.markEvaluationCompleted(request, observedAt);
participationsFinalized++;
}
var ready = workPort.loadReadyResult(roomId);
if (ready.isEmpty()) {
continue;
}
resultService.finalize(command(ready.orElseThrow()));
roomsFinalized++;
} catch (RuntimeException exception) {
String message = exception.getMessage();
failures.add(new RoomFinalizationFailure(
roomId,
exception.getClass().getSimpleName() + (message == null ? "" : ": " + message)));
}
}
return new RoomFinalizationReport(
observedAt, roomIds.size(), roomsFinalized, participationsFinalized, failures);
}

private FinalRoomResultCommand command(RoomFinalizationSource source) {
var template = templateCatalog.parseLocked(source.scoringTemplate());
List<FinalRoomResultCandidate> candidates = source.candidates().stream()
.map(candidate -> candidate(source, candidate))
.toList();
return new FinalRoomResultCommand(
source.roomId(), template.id(), source.cutoffAt(), template, candidates);
}

private FinalRoomResultCandidate candidate(
RoomFinalizationSource room, RoomFinalizationCandidateSource candidate) {
var evidence = evidenceService.prepare(new ScoringEvidenceRequest(
candidate.participationId(),
candidate.evaluationSegmentId(),
candidate.performanceSnapshotId(),
room.scoringTemplate().id()));
var source = evidence.source();
var eligibility = scoring.eligibility(
candidate.scheduledEvaluationSeconds(),
candidate.scheduledEvaluationSeconds(),
candidate.actualOperationSeconds(),
candidate.baseRequiredOperationSeconds(),
candidate.actualFillCount(),
candidate.baseRequiredFillCount(),
true);
return new FinalRoomResultCandidate(
candidate.participationId(),
candidate.performanceSnapshotId(),
new OfficialScoringMetrics(
source.totalReturnPct(), source.maxDrawdownPct(), source.sharpeRatio()),
eligibility,
evidence.provenanceHash(),
calculationDocument(evidence, candidate, eligibility));
}

private String calculationDocument(
ScoringEvidenceBundle evidence,
RoomFinalizationCandidateSource candidate,
OfficialScoringEligibility eligibility) {
Map<String, Object> document = new LinkedHashMap<>();
document.put("schemaVersion", "live-room-finalization.v1");
document.put("provenanceVersion", evidence.provenanceVersion());
document.put("performanceSnapshotHash", evidence.source().performanceSnapshotHash());
document.put("roomRulesHash", evidence.source().roomRulesHash());
document.put("scoringTemplateRulesHash", evidence.source().lockedScoringTemplateRulesHash());
document.put("scheduledEvaluationSeconds", candidate.scheduledEvaluationSeconds());
document.put("normalEvaluationSeconds", candidate.scheduledEvaluationSeconds());
document.put("actualOperationSeconds", candidate.actualOperationSeconds());
document.put("actualFillCount", candidate.actualFillCount());
document.put("requiredOperationSeconds", eligibility.requiredOperationSeconds());
document.put("requiredFillCount", eligibility.requiredFillCount());
document.put("coverage", eligibility.coverage().toPlainString());
document.put("eligibilityReasons", eligibility.reasons().stream().map(Enum::name).toList());
try {
return mapper.writeValueAsString(document);
} catch (JsonProcessingException exception) {
throw new IllegalArgumentException("room finalization evidence is not JSON serializable", exception);
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,22 @@
package com.idea2strategy.backend.application.competition;

import java.time.Instant;
import java.util.List;
import java.util.Objects;
import java.util.UUID;

public record RoomFinalizationSource(
UUID roomId,
Instant cutoffAt,
ScoringTemplateCatalogRecord scoringTemplate,
List<RoomFinalizationCandidateSource> candidates) {
public RoomFinalizationSource {
Objects.requireNonNull(roomId, "roomId");
Objects.requireNonNull(cutoffAt, "cutoffAt");
Objects.requireNonNull(scoringTemplate, "scoringTemplate");
candidates = List.copyOf(Objects.requireNonNull(candidates, "candidates"));
if (candidates.isEmpty()) {
throw new IllegalArgumentException("room finalization candidates must not be empty");
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,16 @@
package com.idea2strategy.backend.application.competition;

import java.time.Instant;
import java.util.List;
import java.util.Optional;
import java.util.UUID;

public interface RoomFinalizationWorkPort {
List<UUID> findDueRoomIds(Instant observedAt, int limit);

List<VirtualLiquidationRequest> findPendingLiquidations(UUID roomId);

void markEvaluationCompleted(VirtualLiquidationRequest request, Instant completedAt);

Optional<RoomFinalizationSource> loadReadyResult(UUID roomId);
}
Loading