diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasImpactService.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasImpactService.java index 3277774..1cfe2d7 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasImpactService.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasImpactService.java @@ -41,15 +41,17 @@ public class BiasImpactService { private final BiasBehaviorAdapterRegistry adapterRegistry; private final BiasIterationComparator iterationComparator; private final BiasImpactJudge judge; + private final BiasJudgementStore judgementStore; public BiasImpactService(ExecutionsService executionsService, BiasImpactReportRepository reportRepository, BiasBehaviorAdapterRegistry adapterRegistry, BiasIterationComparator iterationComparator, - BiasImpactJudge judge) { + BiasImpactJudge judge, BiasJudgementStore judgementStore) { this.executionsService = executionsService; this.reportRepository = reportRepository; this.adapterRegistry = adapterRegistry; this.iterationComparator = iterationComparator; this.judge = judge; + this.judgementStore = judgementStore; } @Transactional @@ -582,15 +584,13 @@ public class BiasImpactService { * would be judging the figures rather than the outputs. Both executions are final, so the * recomputation is the same comparison, and only the verdicts are written onto the stored report. * - *

The evaluation replaces any previous one. The descriptor and timestamp on it say which model - * produced the verdicts now on display, which is what makes a second opinion worth asking for. + *

The evaluation is added to the report's history rather than replacing what is there: asking + * a second model is how one finds out whether the first was reading the change or inventing it, + * and that only works if both answers survive. The descriptor and timestamp on each say which + * model produced it. */ - @Transactional public BiasImpactReport judgeReport(String reportId, LLMDescriptor descriptor, String owner) { - BiasImpactReportEntity entity = reportRepository.findByIdAndOwner(reportId, owner) - .orElseThrow(() -> new BiasApiException(HttpStatus.NOT_FOUND, ValidationErrorCode.BIAS_REPORT_NOT_FOUND, - "biasImpactReport", reportId, "Bias impact report not found: " + reportId)); - BiasImpactReport stored = entity.getReport(); + BiasImpactReport stored = getReport(reportId, owner); BiasImpactReport comparison = comparisonForJudging(stored, owner); ExecutionObject baseline = executionsService.getExecutionByOwner(stored.baselineExecutionId(), owner); @@ -598,11 +598,12 @@ public class BiasImpactService { ? baseline : executionsService.getExecutionByOwner(stored.biasedExecutionId(), owner); + // Deliberately outside any transaction: this is one model call per compared pair, so minutes + // of it, and what comes back is written by judgementStore against the report as it is then - + // not against the copy read above, which by that point may be missing another assessment + // that finished in the meantime. BiasJudgeSummary judgement = judge.evaluate(comparison, descriptor, baseline, biased); - BiasImpactReport judged = BiasJudgeTarget.apply(stored, judgement.verdicts()).withJudgement(judgement); - entity.setReport(judged); - reportRepository.save(entity); - return judged; + return judgementStore.append(reportId, owner, judgement); } /** diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasJudgementStore.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasJudgementStore.java new file mode 100644 index 0000000..1779b2c --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasJudgementStore.java @@ -0,0 +1,89 @@ +package it.cnr.isti.workflow.manager.executions.bias; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.dao.OptimisticLockingFailureException; +import org.springframework.http.HttpStatus; +import org.springframework.stereotype.Component; +import org.springframework.transaction.annotation.Propagation; +import org.springframework.transaction.annotation.Transactional; + +import it.cnr.isti.workflow.manager.executions.bias.persistence.BiasImpactReportEntity; +import it.cnr.isti.workflow.manager.executions.bias.persistence.BiasImpactReportRepository; +import it.cnr.isti.workflow.manager.flows.validation.ValidationErrorCode; + +/** + * Writes an assessment onto a report, against the report as it is now. + * + *

A bean of its own because the transaction has to start here and nowhere earlier. Judging takes + * one model call per compared pair, so minutes; when the reading, the judging and the writing sat in + * one transactional method, two overlapping assessments each began from the report they had read + * before their own model calls - both saw no history - and the one that finished last saved its + * own judgement as the only one on the report. The other was simply gone, which is what a user sees + * as an assessment that vanished. + * + *

It is also why the transaction must not span the model calls: a database transaction held open + * for minutes pins a connection and can only lose to a timeout. + */ +@Component +public class BiasJudgementStore { + + private static final Logger log = LoggerFactory.getLogger(BiasJudgementStore.class); + + /** + * Attempts at the append itself. Each one is a re-read and a save, so a retry is cheap - what + * cost minutes was the judging, and that is already done and held in hand by then. + */ + private static final int MAX_ATTEMPTS = 5; + + private final BiasImpactReportRepository reportRepository; + private final BiasJudgementStore self; + + public BiasJudgementStore(BiasImpactReportRepository reportRepository, + @org.springframework.context.annotation.Lazy BiasJudgementStore self) { + this.reportRepository = reportRepository; + // The retry has to re-enter through the proxy: a plain this.append() would run the next + // attempt inside the transaction that just failed, which can only fail again. + this.self = self; + } + + /** + * Appends an assessment to whatever the report holds now, retrying if someone else got there + * first. + * + *

The retry is the whole point: losing the race means another assessment landed between the + * read and the save, so the next attempt reads that one too and appends after it. Both survive, + * which is what the version column was added for. + */ + public BiasImpactReport append(String reportId, String owner, BiasJudgeSummary judgement) { + for (int attempt = 1; attempt <= MAX_ATTEMPTS; attempt++) { + try { + return self.appendOnce(reportId, owner, judgement); + } catch (OptimisticLockingFailureException conflict) { + log.debug("Report {} changed under an assessment append, attempt {} of {}", reportId, attempt, + MAX_ATTEMPTS); + } + } + throw new BiasApiException(HttpStatus.CONFLICT, ValidationErrorCode.BIAS_JUDGE_FAILED, + "biasImpactReport", reportId, + "The report kept changing while this assessment was being stored. The assessment itself " + + "succeeded; ask for it again to have it recorded."); + } + + /** + * One attempt, in its own transaction: read the report as it is, put this assessment's verdicts + * on those pairs, and add its summary to that history. + */ + @Transactional(propagation = Propagation.REQUIRES_NEW) + public BiasImpactReport appendOnce(String reportId, String owner, BiasJudgeSummary judgement) { + BiasImpactReportEntity entity = reportRepository.findByIdAndOwner(reportId, owner) + .orElseThrow(() -> new BiasApiException(HttpStatus.NOT_FOUND, ValidationErrorCode.BIAS_REPORT_NOT_FOUND, + "biasImpactReport", reportId, "Bias impact report not found: " + reportId)); + + BiasImpactReport current = entity.getReport(); + BiasImpactReport judged = BiasJudgeTarget.apply(current, judgement.verdicts()).withJudgement(judgement); + entity.setReport(judged); + reportRepository.saveAndFlush(entity); + return judged; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/persistence/BiasImpactReportEntity.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/persistence/BiasImpactReportEntity.java index df44629..8c4f64e 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/persistence/BiasImpactReportEntity.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/persistence/BiasImpactReportEntity.java @@ -10,6 +10,7 @@ import jakarta.persistence.Id; import jakarta.persistence.Index; import jakarta.persistence.Table; import jakarta.persistence.UniqueConstraint; +import jakarta.persistence.Version; import jakarta.validation.constraints.NotBlank; import jakarta.validation.constraints.NotNull; import lombok.AllArgsConstructor; @@ -53,4 +54,17 @@ public class BiasImpactReportEntity { @Column(name = "report_data", columnDefinition = "TEXT") @Convert(converter = BiasImpactReportConverter.class) private BiasImpactReport report; + + /** + * Bumped on every write, so a second writer that read the row before the first one saved is + * refused instead of overwriting it. + * + *

Appending an LLM assessment is a read-modify-write over {@link #report}, with minutes of + * model calls in the middle: two overlapping assessments both read a report with no history and + * the one that finished last saved its own judgement as the only one there. The loser now fails + * and is retried against the row as it actually is. + */ + @Version + @Column(name = "version", nullable = false) + private long version; } diff --git a/src/main/resources/db/migration/V10__add_bias_report_version.sql b/src/main/resources/db/migration/V10__add_bias_report_version.sql new file mode 100644 index 0000000..5d627a7 --- /dev/null +++ b/src/main/resources/db/migration/V10__add_bias_report_version.sql @@ -0,0 +1,6 @@ +-- Appending an LLM assessment to a report is a read-modify-write over its JSON, and the model calls +-- in the middle take minutes. Two overlapping assessments both read a report with no history, and +-- whichever saved last kept only its own: the other was silently gone. A version column is what lets +-- the second writer be refused and retried against the row as it actually is. +ALTER TABLE bias_impact_report_entity + ADD COLUMN version BIGINT NOT NULL DEFAULT 0; diff --git a/src/test/java/it/cnr/isti/workflow/manager/executions/bias/BiasExperimentsIntegrationTest.java b/src/test/java/it/cnr/isti/workflow/manager/executions/bias/BiasExperimentsIntegrationTest.java index a0bbd9a..fea0bfd 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/executions/bias/BiasExperimentsIntegrationTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/executions/bias/BiasExperimentsIntegrationTest.java @@ -61,6 +61,12 @@ class BiasExperimentsIntegrationTest { private static final String OWNER = "bias-experiment-user"; + /** + * Two parties, released together: the two overlapping evaluations. Sized for the pair, so each + * one's first model call waits for the other to reach the same point. + */ + private static final java.util.concurrent.Phaser OVERLAP_GATE = new java.util.concurrent.Phaser(2); + private static final LLMDescriptor SIMULATOR = LLMDescriptor.builder() .provider("biasExperimentProvider") .model("simulate-model") @@ -134,6 +140,39 @@ class BiasExperimentsIntegrationTest { }; } + /** + * Blocks inside the model call until a second evaluation has reached it too, so two + * assessments of one report are genuinely in flight at the same time. + */ + @Bean + LLMProvider overlappingJudgeProvider() { + return new LLMProvider() { + @Override + public String getName() { + return "overlappingJudgeProvider"; + } + + @Override + public List getRegisteredModels() { + return List.of("overlapping-judge-model"); + } + + @Override + public String generate(String model, String prompt) { + return "Both assessments were in flight together."; + } + + @Override + public String generateJson(String model, String prompt) { + OVERLAP_GATE.arriveAndAwaitAdvance(); + return """ + {"impact":"SUBSTANTIVE","attribution":"INJECTION","confidence":0.5, + "changedAspects":["ranking"],"rationale":"Assessed while another was running."} + """; + } + }; + } + /** * Answers nothing in JSON mode, like a reasoning model under Ollama's format=json, and the * real thing when asked as text. @@ -704,6 +743,36 @@ class BiasExperimentsIntegrationTest { assertEquals(2, reloaded.judgements().size()); } + @Test + void twoAssessmentsRunningAtOnceBothSurvive() throws Exception { + Block block = annotatedLlmBlock(); + ExecutionObject baseline = completedExecution(block); + String annotationId = block.getBiasAnnotations().getFirst().id(); + ExecutionObject variant = createAndRunVariant(baseline, block, annotationId, BiasInterventionDirection.BIAS); + BiasImpactReport report = biasImpactService.compareFullFlow(baseline.getId(), variant.getId(), true, OWNER); + + // What a second click on "Evaluate impact with LLM" does while the first is still running. + // Both used to start from the report they read before their own model calls, so the one that + // finished last saved its judgement as the only one and the other was lost. + java.util.concurrent.ExecutorService pool = java.util.concurrent.Executors.newFixedThreadPool(2); + try { + var first = pool.submit(() -> biasImpactService.judgeReport(report.id(), + LLMDescriptor.builder().provider("overlappingJudgeProvider").model("model-one").build(), OWNER)); + var second = pool.submit(() -> biasImpactService.judgeReport(report.id(), + LLMDescriptor.builder().provider("overlappingJudgeProvider").model("model-two").build(), OWNER)); + first.get(30, java.util.concurrent.TimeUnit.SECONDS); + second.get(30, java.util.concurrent.TimeUnit.SECONDS); + } finally { + pool.shutdownNow(); + } + + BiasImpactReport reloaded = biasImpactService.getReport(report.id(), OWNER); + assertEquals(2, reloaded.judgements().size()); + assertEquals( + Set.of("model-one", "model-two"), + reloaded.judgements().stream().map(one -> one.judge().model()).collect(java.util.stream.Collectors.toSet())); + } + @Test void theAssessmentHistoryStopsAtItsCapInsteadOfGrowingWithoutBound() { Block block = annotatedLlmBlock();