Stop a second LLM assessment from erasing the one still running
Two assessments of the same report, started minutes apart, left only one: observed in a real run as two REPORT_JUDGE jobs both COMPLETED (14:33-14:37 and 14:36-14:37) against a report that ended up holding a single judgement - the one that finished last. judgeReport was one @Transactional method that read the report, spent minutes in one model call per compared pair, then wrote. The second assessment read the report while the first was still calling models, saw no history, and saved its own judgement as the only one there. Nothing warned, because from each writer's side the write succeeded. The model calls now happen outside any transaction - holding one open across minutes also pins a connection for no reason - and the write moved to BiasJudgementStore, a bean of its own so the transaction starts there and a retry re-enters through the proxy rather than inside the transaction that just failed. It re-reads the report as it is at that moment, puts the new verdicts on those pairs and the summary on that history, and the report entity carries a @Version so a writer working from a stale read is refused instead of overwriting: the next attempt reads the assessment that landed meanwhile and appends after it. Optimistic rather than SELECT ... FOR UPDATE because the tests said so: Hibernate renders PESSIMISTIC_WRITE as "for no key update", which H2 - what the suite runs on - cannot parse. A fix that only holds on one dialect is not one. The regression test runs the two assessments on two threads, with a stub provider that blocks inside the model call until both have reached it, and asserts the report keeps both. Reinstating the old shape fails it. This prevents further losses; it does not recover an assessment already lost. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
parent
02b284d368
commit
e71ba9c290
|
|
@ -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.
|
||||
*
|
||||
* <p>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.
|
||||
* <p>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);
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
*
|
||||
* <p>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.
|
||||
*
|
||||
* <p>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.
|
||||
*
|
||||
* <p>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;
|
||||
}
|
||||
}
|
||||
|
|
@ -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.
|
||||
*
|
||||
* <p>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;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
@ -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<String> 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<LLMBlockType> 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<LLMBlockType> block = annotatedLlmBlock();
|
||||
|
|
|
|||
Loading…
Reference in New Issue