Implement bias impact experiments

This commit is contained in:
Lucio Lelii 2026-07-20 17:22:26 +02:00
parent 366b048689
commit 587eadc678
52 changed files with 1953 additions and 21 deletions

View File

@ -107,6 +107,10 @@
<groupId>org.flywaydb</groupId>
<artifactId>flyway-core</artifactId>
</dependency>
<dependency>
<groupId>org.flywaydb</groupId>
<artifactId>flyway-database-postgresql</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>

View File

@ -16,6 +16,7 @@ import io.swagger.v3.oas.annotations.security.SecurityRequirement;
import it.cnr.isti.workflow.manager.blocks.configurations.JsonSchemaProducer;
import it.cnr.isti.workflow.manager.flows.model.bias.BiasAnnotationSource;
import it.cnr.isti.workflow.manager.flows.model.bias.BiasAnnotationStatus;
import it.cnr.isti.workflow.manager.flows.model.bias.BiasActivationMode;
import it.cnr.isti.workflow.manager.flows.model.bias.BiasCatalogOption;
import it.cnr.isti.workflow.manager.flows.model.bias.BiasCategory;
import it.cnr.isti.workflow.manager.flows.model.bias.BiasSeverity;
@ -54,6 +55,7 @@ public class BiasAnnotationsController {
options.put("severity", describe(BiasSeverity.class));
options.put("status", describe(BiasAnnotationStatus.class));
options.put("source", describe(BiasAnnotationSource.class));
options.put("behavioralProbe.activationMode", describe(BiasActivationMode.class));
return new BiasAnnotationDescriptor(
BlockBiasAnnotation.class.getSimpleName(),

View File

@ -0,0 +1,71 @@
package it.cnr.isti.workflow.manager.controllers;
import java.util.List;
import java.util.Map;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.http.HttpStatus;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.server.ResponseStatusException;
import io.swagger.v3.oas.annotations.Operation;
import it.cnr.isti.workflow.manager.blocks.Block;
import it.cnr.isti.workflow.manager.blocks.factories.BlockFactory;
import it.cnr.isti.workflow.manager.blocks.types.BlockType;
import it.cnr.isti.workflow.manager.executions.bias.runtime.BiasBehaviorAdapterRegistry;
import it.cnr.isti.workflow.manager.flows.model.bias.BiasActivationMode;
@RestController
@RequestMapping("/blocks/types")
public class BiasCapabilitiesController {
@Autowired
Map<String, BlockType> blockTypes;
@Autowired
List<BlockFactory<?, ?>> blockFactories;
@Autowired
BiasBehaviorAdapterRegistry adapterRegistry;
public record BiasCapabilityDescriptor(
String blockType,
boolean supported,
boolean isolatedExperimentSupported,
boolean fullFlowExperimentSupported,
boolean externalSideEffects,
List<BiasActivationMode> activationModes) {
}
@GetMapping("/{type}/bias-capabilities")
@Operation(summary = "Get bias experiment capabilities for a block type",
description = "Returns the activation modes and safety characteristics supported by the block runtime adapter.")
public BiasCapabilityDescriptor getCapabilities(@PathVariable String type) {
BlockType blockType = blockTypes.get(type);
if (blockType == null) {
throw new ResponseStatusException(HttpStatus.NOT_FOUND, "Block type not found: " + type);
}
Block<?> example = createExample(blockType);
List<BiasActivationMode> modes = adapterRegistry.supportedModes(example).stream().sorted().toList();
return new BiasCapabilityDescriptor(
blockType.getName(),
!modes.isEmpty(),
adapterRegistry.supportsIsolatedExperiment(example),
!modes.isEmpty(),
adapterRegistry.hasExternalSideEffects(example),
modes);
}
@SuppressWarnings({ "rawtypes", "unchecked" })
private Block<?> createExample(BlockType blockType) {
BlockFactory factory = blockFactories.stream()
.filter(candidate -> candidate.getBlockType().equals(blockType.getClass()))
.findFirst()
.orElseThrow(() -> new ResponseStatusException(HttpStatus.NOT_FOUND,
"Block factory not found for type: " + blockType.getName()));
return factory.createEmpty();
}
}

View File

@ -0,0 +1,95 @@
package it.cnr.isti.workflow.manager.controllers;
import java.util.List;
import org.springframework.http.HttpStatus;
import org.springframework.security.core.annotation.AuthenticationPrincipal;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.server.ResponseStatusException;
import io.swagger.v3.oas.annotations.Operation;
import io.swagger.v3.oas.annotations.security.SecurityRequirement;
import it.cnr.isti.workflow.manager.auth.repo.LoginEntity;
import it.cnr.isti.workflow.manager.executions.ExecutionsService;
import it.cnr.isti.workflow.manager.executions.api.ExecutionView;
import it.cnr.isti.workflow.manager.executions.bias.BiasImpactExperimentRequest;
import it.cnr.isti.workflow.manager.executions.bias.BiasImpactReport;
import it.cnr.isti.workflow.manager.executions.bias.BiasImpactService;
import it.cnr.isti.workflow.manager.executions.bias.BiasRerunRequest;
import jakarta.validation.Valid;
@RestController
@RequestMapping("/executions")
@SecurityRequirement(name = "bearerAuth")
public class BiasExperimentsController {
private final ExecutionsService executionsService;
private final BiasImpactService biasImpactService;
public BiasExperimentsController(ExecutionsService executionsService, BiasImpactService biasImpactService) {
this.executionsService = executionsService;
this.biasImpactService = biasImpactService;
}
@PostMapping("/{executionId}/steps/{stepId}/bias-impact")
@Operation(summary = "Run an isolated bias impact experiment",
description = "Uses the completed step output as baseline, executes biased variants with the same inputs and persists a comparison report.")
public BiasImpactReport runIsolatedExperiment(
@PathVariable String executionId,
@PathVariable String stepId,
@RequestBody @Valid BiasImpactExperimentRequest request,
@AuthenticationPrincipal LoginEntity userDetails) {
return biasImpactService.runIsolatedStepExperiment(executionId, stepId, request, owner(userDetails));
}
@PostMapping("/{executionId}/bias-rerun")
@Operation(summary = "Create a full-flow biased rerun",
description = "Creates a rerun with selected bias annotations activated. Inputs and rerun lineage are copied from the final baseline execution.")
public ExecutionView createBiasRerun(
@PathVariable String executionId,
@RequestBody @Valid BiasRerunRequest request,
@AuthenticationPrincipal LoginEntity userDetails) {
return ExecutionView.fromExecution(executionsService.createBiasRerun(executionId, owner(userDetails), request));
}
@PostMapping("/{baselineExecutionId}/bias-compare/{biasedExecutionId}")
@Operation(summary = "Compare a baseline execution with a biased rerun",
description = "Compares immediate outputs, routing and downstream step outcomes and persists the full-flow report.")
public BiasImpactReport compareFullFlow(
@PathVariable String baselineExecutionId,
@PathVariable String biasedExecutionId,
@RequestParam(defaultValue = "true") boolean includeRawOutputs,
@AuthenticationPrincipal LoginEntity userDetails) {
return biasImpactService.compareFullFlow(
baselineExecutionId, biasedExecutionId, includeRawOutputs, owner(userDetails));
}
@GetMapping("/{executionId}/bias-impact-reports")
@Operation(summary = "List persisted bias impact reports for an execution")
public List<BiasImpactReport> getReports(
@PathVariable String executionId,
@AuthenticationPrincipal LoginEntity userDetails) {
return biasImpactService.getReports(executionId, owner(userDetails));
}
@GetMapping("/bias-impact-reports/{reportId}")
@Operation(summary = "Get a persisted bias impact report")
public BiasImpactReport getReport(
@PathVariable String reportId,
@AuthenticationPrincipal LoginEntity userDetails) {
return biasImpactService.getReport(reportId, owner(userDetails));
}
private String owner(LoginEntity userDetails) {
if (userDetails == null || userDetails.getUsername() == null || userDetails.getUsername().isBlank()) {
throw new ResponseStatusException(HttpStatus.FORBIDDEN, "Authentication required");
}
return userDetails.getUsername();
}
}

View File

@ -15,6 +15,7 @@ import it.cnr.isti.workflow.manager.executions.persistence.ExecutionSnapshot;
import it.cnr.isti.workflow.manager.executions.steps.Input;
import it.cnr.isti.workflow.manager.executions.steps.Step;
import it.cnr.isti.workflow.manager.executions.steps.StepStatus;
import it.cnr.isti.workflow.manager.executions.bias.BiasExecutionContext;
import it.cnr.isti.workflow.manager.llms.LLMDescriptor;
import lombok.AccessLevel;
import lombok.Getter;
@ -78,6 +79,8 @@ public class ExecutionContext implements ExecutionListener {
LLMDescriptor interactionSimulationDescriptor;
BiasExecutionContext biasExecutionContext = BiasExecutionContext.normal();
protected void setInteractionSimulationEnabled(boolean interactionSimulationEnabled) {
this.interactionSimulationEnabled = interactionSimulationEnabled;
}
@ -87,6 +90,11 @@ public class ExecutionContext implements ExecutionListener {
this.steps.values().forEach(step -> step.setInteractionSimulationDescriptor(interactionSimulationDescriptor));
}
protected void setBiasExecutionContext(BiasExecutionContext biasExecutionContext) {
this.biasExecutionContext = biasExecutionContext == null ? BiasExecutionContext.normal() : biasExecutionContext;
this.steps.values().forEach(step -> step.setBiasExecutionContext(this.biasExecutionContext));
}
public ExecutionContext(Map<String, Step<?>> steps) {
this.steps = steps;
this.steps.values().forEach(step -> {
@ -95,6 +103,7 @@ public class ExecutionContext implements ExecutionListener {
step.setExecutionVariables(this.runtimeExecutionVariables);
step.setExecutionVariableDescriptors(this.executionVariableDescriptors);
step.setInteractionSimulationDescriptor(this.interactionSimulationDescriptor);
step.setBiasExecutionContext(this.biasExecutionContext);
step.setEventLogger(createStepEventLogger(step));
});
refreshRuntimeExecutionVariables();
@ -408,6 +417,11 @@ public class ExecutionContext implements ExecutionListener {
return this.executionVariables.get(key);
}
@JsonIgnore
public Map<String, Object> getResolvedExecutionVariables() {
return Collections.unmodifiableMap(this.runtimeExecutionVariables);
}
protected ExecutionVariableDescriptor getExecutionVariableDescriptor(String key) {
return this.executionVariableDescriptors.get(key);
}
@ -476,6 +490,7 @@ public class ExecutionContext implements ExecutionListener {
.endTime(this.endTime)
.interactionSimulationEnabled(this.interactionSimulationEnabled)
.interactionSimulationDescriptor(this.interactionSimulationDescriptor)
.biasExecutionContext(this.biasExecutionContext)
.executionVariables(new HashMap<>(this.executionVariables))
.executionVariableDescriptors(new HashMap<>(this.executionVariableDescriptors))
.globalInputs(new HashMap<>(this.globalInputs))
@ -571,6 +586,7 @@ public class ExecutionContext implements ExecutionListener {
this.endTime = snapshot.getEndTime();
this.interactionSimulationEnabled = snapshot.isInteractionSimulationEnabled();
this.setInteractionSimulationDescriptor(snapshot.getInteractionSimulationDescriptor());
this.setBiasExecutionContext(snapshot.getBiasExecutionContext());
refreshRuntimeExecutionVariables();
this.status = normalizeRestoredStatus(snapshot.getStatus());
}

View File

@ -23,5 +23,7 @@ public enum ExecutionEventType {
HTTP_REQUEST,
MCP_SESSION_OPENED,
MCP_SESSION_REUSED,
MCP_SESSION_CLOSED
MCP_SESSION_CLOSED,
BIAS_EXPERIMENT_APPLIED,
BIAS_SIDE_EFFECT_MOCKED
}

View File

@ -16,6 +16,7 @@ import com.fasterxml.jackson.annotation.JsonIgnore;
import it.cnr.isti.workflow.manager.containers.Container;
import it.cnr.isti.workflow.manager.executions.persistence.ExecutionSnapshot;
import it.cnr.isti.workflow.manager.executions.persistence.ExecutionStepSnapshot;
import it.cnr.isti.workflow.manager.executions.bias.BiasExecutionContext;
import it.cnr.isti.workflow.manager.executions.executors.BlockExecutors;
import it.cnr.isti.workflow.manager.executions.executors.NodeExecutors;
import it.cnr.isti.workflow.manager.executions.steps.Input;
@ -77,9 +78,12 @@ public class ExecutionObject {
LLMDescriptor interactionSimulationDescriptor;
BiasExecutionContext biasExecutionContext = BiasExecutionContext.normal();
@Builder
public ExecutionObject(String executionName, FlowData flow, List<ExecutionAuthorizationRequirement> requiredAuthorizations,
String owner, String runGroupId, String sourceFlowId, String rerunOfExecutionId, Integer runNumber) {
String owner, String runGroupId, String sourceFlowId, String rerunOfExecutionId, Integer runNumber,
BiasExecutionContext biasExecutionContext) {
this.name = executionName;
this.owner = owner;
this.sourceFlowId = sourceFlowId;
@ -93,12 +97,14 @@ public class ExecutionObject {
this.stepDependencies = flow.getDependencies() == null ? List.of() : List.copyOf(flow.getDependencies());
this.requiredGlobalInputs = flow.getGlobalInputs() == null ? List.of() : List.copyOf(flow.getGlobalInputs());
this.requiredAuthorizations = requiredAuthorizations == null ? List.of() : List.copyOf(requiredAuthorizations);
this.biasExecutionContext = biasExecutionContext == null ? BiasExecutionContext.normal() : biasExecutionContext;
List<Step<?>> steps = getStepsFromFlow(flow);
this.executorService = createExecutorService(steps.size());
this.context = new ExecutionContext(steps.stream().collect(Collectors.toMap(Step::getId, Function.identity())));
this.context.setBiasExecutionContext(this.biasExecutionContext);
ensureRuntimeContextVariables();
registerGlobalInputs();
this.requiredAuthorizations.stream()
@ -241,6 +247,11 @@ public class ExecutionObject {
this.context.setInteractionSimulationDescriptor(interactionSimulationDescriptor);
}
protected void setBiasExecutionContext(BiasExecutionContext biasExecutionContext) {
this.biasExecutionContext = biasExecutionContext == null ? BiasExecutionContext.normal() : biasExecutionContext;
this.context.setBiasExecutionContext(this.biasExecutionContext);
}
protected void start(boolean simulateInteractions) {
List<String> missingAuthorizations = getMissingAuthorizationKeys();
if (!missingAuthorizations.isEmpty()) {
@ -364,12 +375,16 @@ public class ExecutionObject {
this.context.restore(snapshot);
this.interactionSimulationEnabled = snapshot.isInteractionSimulationEnabled();
this.interactionSimulationDescriptor = snapshot.getInteractionSimulationDescriptor();
this.biasExecutionContext = snapshot.getBiasExecutionContext() == null
? BiasExecutionContext.normal()
: snapshot.getBiasExecutionContext();
}
ensureRuntimeContextVariables();
registerGlobalInputs();
rebuildDependencyStates();
this.providedAuthorizations.forEach(this.context::setAuthorization);
this.context.setInteractionSimulationDescriptor(this.interactionSimulationDescriptor);
this.context.setBiasExecutionContext(this.biasExecutionContext);
this.simulationAvailable = hasSimulationAvailable(new ArrayList<>(this.context.getSteps().values()));
refreshInitializationStatus();
}

View File

@ -37,6 +37,11 @@ import it.cnr.isti.workflow.manager.containers.Container;
import it.cnr.isti.workflow.manager.containers.configurations.ContainerConfiguration;
import it.cnr.isti.workflow.manager.containers.configurations.LoopContainerConfiguration;
import it.cnr.isti.workflow.manager.executions.persistence.ExecutionSnapshot;
import it.cnr.isti.workflow.manager.executions.bias.BiasActivation;
import it.cnr.isti.workflow.manager.executions.bias.BiasExecutionContext;
import it.cnr.isti.workflow.manager.executions.bias.BiasExecutionMode;
import it.cnr.isti.workflow.manager.executions.bias.BiasRerunRequest;
import it.cnr.isti.workflow.manager.executions.bias.persistence.BiasImpactReportRepository;
import it.cnr.isti.workflow.manager.executions.repo.ExecutionEntity;
import it.cnr.isti.workflow.manager.executions.repo.ExecutionRepository;
import it.cnr.isti.workflow.manager.flows.model.Flow;
@ -70,6 +75,9 @@ public class ExecutionsService {
@Autowired
ExecutionRepository executionRepository;
@Autowired
BiasImpactReportRepository biasImpactReportRepository;
@Autowired
MCPAgentService mcpAgentService;
@ -108,6 +116,13 @@ public class ExecutionsService {
@Transactional
public ExecutionObject createExecution(String executionName, FlowData flow, String owner, String runGroupId,
String sourceFlowId, String rerunOfExecutionId, Integer runNumber) {
return createExecution(executionName, flow, owner, runGroupId, sourceFlowId, rerunOfExecutionId, runNumber,
BiasExecutionContext.normal());
}
@Transactional
public ExecutionObject createExecution(String executionName, FlowData flow, String owner, String runGroupId,
String sourceFlowId, String rerunOfExecutionId, Integer runNumber, BiasExecutionContext biasExecutionContext) {
flowExecutionValidator.validate(flow);
List<ExecutionAuthorizationRequirement> requiredAuthorizations = resolveRequiredAuthorizations(flow);
ExecutionObject execObject = ExecutionObject.builder()
@ -119,6 +134,7 @@ public class ExecutionsService {
.sourceFlowId(sourceFlowId)
.rerunOfExecutionId(rerunOfExecutionId)
.runNumber(runNumber)
.biasExecutionContext(biasExecutionContext)
.build();
attachPersistence(execObject);
executions.put(execObject.getId(), execObject);
@ -292,6 +308,62 @@ public class ExecutionsService {
return rerun;
}
@Transactional
public ExecutionObject createBiasRerun(String id, String owner, BiasRerunRequest request) {
ExecutionObject source = getExecutionByOwner(id, owner);
if (!source.getContext().getStatus().isFinalState()) {
throw new ResponseStatusException(HttpStatus.BAD_REQUEST,
"Execution with id " + id + " can be used for a bias rerun only from a final state");
}
Map<String, List<String>> annotationIdsByNode = validateBiasActivations(source.getFlow(), request.activations());
String experimentId = java.util.UUID.randomUUID().toString();
BiasExecutionContext biasContext = new BiasExecutionContext(
experimentId,
BiasExecutionMode.BIAS_VARIANT,
annotationIdsByNode,
request.externalSideEffectPolicy(),
request.confirmExternalSideEffects());
String runGroupId = resolveHistoryGroupId(source);
int nextRunNumber = nextRunNumber(source.getSourceFlowId(), runGroupId, owner, source.getRunNumber());
ExecutionObject rerun = createExecution(source.getName(), source.getFlow(), owner, runGroupId,
source.getSourceFlowId(), source.getId(), nextRunNumber, biasContext);
copyReusableInputs(source, rerun);
persist(rerun);
return rerun;
}
private Map<String, List<String>> validateBiasActivations(FlowData flow, List<BiasActivation> activations) {
Map<String, it.cnr.isti.workflow.manager.flows.model.FlowNode> nodesById = flow.getNodes().stream()
.collect(Collectors.toMap(it.cnr.isti.workflow.manager.flows.model.FlowNode::getId, node -> node));
Map<String, List<String>> result = new LinkedHashMap<>();
for (BiasActivation activation : activations) {
it.cnr.isti.workflow.manager.flows.model.FlowNode node = nodesById.get(activation.nodeId());
if (!(node instanceof Block<?> block)) {
throw new ResponseStatusException(HttpStatus.BAD_REQUEST,
"Bias activation node is not a block: " + activation.nodeId());
}
Map<String, it.cnr.isti.workflow.manager.flows.model.bias.BlockBiasAnnotation> annotationsById =
block.getBiasAnnotations().stream().collect(Collectors.toMap(
it.cnr.isti.workflow.manager.flows.model.bias.BlockBiasAnnotation::id,
annotation -> annotation));
for (String annotationId : activation.annotationIds()) {
var annotation = annotationsById.get(annotationId);
if (annotation == null) {
throw new ResponseStatusException(HttpStatus.BAD_REQUEST,
"Bias annotation " + annotationId + " not found on block " + block.getId());
}
if (annotation.behavioralProbe() == null) {
throw new ResponseStatusException(HttpStatus.BAD_REQUEST,
"Bias annotation " + annotationId + " has no behavioralProbe");
}
}
result.computeIfAbsent(block.getId(), ignored -> new java.util.ArrayList<>())
.addAll(activation.annotationIds());
}
result.replaceAll((nodeId, annotationIds) -> annotationIds.stream().distinct().toList());
return Map.copyOf(result);
}
@Transactional
public void removeExecution(String id) {
ExecutionObject execution = executions.get(id);
@ -306,6 +378,7 @@ public class ExecutionsService {
execution.shutdown();
executions.remove(id);
lastAccessByExecutionId.remove(id);
biasImpactReportRepository.deleteByBaselineExecutionIdOrBiasedExecutionId(id, id);
executionRepository.deleteById(id);
}
@ -604,6 +677,7 @@ public class ExecutionsService {
private ExecutionObject rebuildExecution(ExecutionEntity entity) {
FlowData flow = entity.getFlow();
ExecutionSnapshot snapshot = entity.getSnapshot();
ExecutionObject executionObject = ExecutionObject.builder()
.executionName(entity.getName())
.flow(flow)
@ -613,8 +687,8 @@ public class ExecutionsService {
.sourceFlowId(resolveSourceFlowId(entity))
.rerunOfExecutionId(entity.getRerunOfExecutionId())
.runNumber(resolveRunNumber(entity))
.biasExecutionContext(snapshot == null ? BiasExecutionContext.normal() : snapshot.getBiasExecutionContext())
.build();
ExecutionSnapshot snapshot = entity.getSnapshot();
executionObject.restore(entity.getId(), entity.getCreationTime(),
snapshot == null ? Map.of() : snapshot.getProvidedAuthorizations(),
snapshot);

View File

@ -5,6 +5,7 @@ import java.util.Map;
import it.cnr.isti.workflow.manager.executions.ExecutionAuthorizationRequirement;
import it.cnr.isti.workflow.manager.executions.ExecutionObject;
import it.cnr.isti.workflow.manager.executions.bias.BiasExecutionContext;
import it.cnr.isti.workflow.manager.flows.model.Connection;
import it.cnr.isti.workflow.manager.flows.model.Dependency;
import it.cnr.isti.workflow.manager.llms.LLMDescriptor;
@ -32,6 +33,7 @@ public class ExecutionView {
private LLMDescriptor interactionSimulationDescriptor;
private List<String> missingAuthorizationKeys;
private List<String> missingGlobalInputKeys;
private BiasExecutionContext biasExecutionContext;
public static ExecutionView fromExecution(ExecutionObject execution) {
return ExecutionView.builder()
@ -52,6 +54,7 @@ public class ExecutionView {
.interactionSimulationDescriptor(execution.getInteractionSimulationDescriptor())
.missingAuthorizationKeys(execution.getMissingAuthorizationKeys())
.missingGlobalInputKeys(execution.getMissingGlobalInputKeys())
.biasExecutionContext(execution.getBiasExecutionContext())
.build();
}
}

View File

@ -0,0 +1,15 @@
package it.cnr.isti.workflow.manager.executions.bias;
import java.util.List;
import jakarta.validation.constraints.NotBlank;
import jakarta.validation.constraints.NotEmpty;
public record BiasActivation(
@NotBlank String nodeId,
@NotEmpty List<@NotBlank String> annotationIds) {
public BiasActivation {
annotationIds = annotationIds == null ? List.of() : annotationIds.stream().distinct().toList();
}
}

View File

@ -0,0 +1,18 @@
package it.cnr.isti.workflow.manager.executions.bias;
import java.util.Map;
public record BiasDownstreamImpact(
String nodeId,
String nodeName,
String baselineStatus,
String biasedStatus,
boolean changed,
Map<String, Object> baselineOutputs,
Map<String, Object> biasedOutputs) {
public BiasDownstreamImpact {
baselineOutputs = baselineOutputs == null ? Map.of() : Map.copyOf(baselineOutputs);
biasedOutputs = biasedOutputs == null ? Map.of() : Map.copyOf(biasedOutputs);
}
}

View File

@ -0,0 +1,46 @@
package it.cnr.isti.workflow.manager.executions.bias;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import com.fasterxml.jackson.annotation.JsonIgnore;
public record BiasExecutionContext(
String experimentId,
BiasExecutionMode mode,
Map<String, List<String>> activeAnnotationIdsByNode,
ExternalSideEffectPolicy externalSideEffectPolicy,
boolean externalSideEffectsConfirmed) {
public BiasExecutionContext {
mode = mode == null ? BiasExecutionMode.NORMAL : mode;
externalSideEffectPolicy = externalSideEffectPolicy == null ? ExternalSideEffectPolicy.BLOCK : externalSideEffectPolicy;
Map<String, List<String>> normalized = new LinkedHashMap<>();
if (activeAnnotationIdsByNode != null) {
activeAnnotationIdsByNode.forEach((nodeId, annotationIds) -> {
if (nodeId != null && !nodeId.isBlank()) {
normalized.put(nodeId, annotationIds == null ? List.of() : annotationIds.stream().distinct().toList());
}
});
}
activeAnnotationIdsByNode = Map.copyOf(normalized);
}
public static BiasExecutionContext normal() {
return new BiasExecutionContext(null, BiasExecutionMode.NORMAL, Map.of(), ExternalSideEffectPolicy.BLOCK, false);
}
public List<String> annotationIdsFor(String nodeId) {
return activeAnnotationIdsByNode.getOrDefault(nodeId, List.of());
}
public boolean isVariantFor(String nodeId) {
return mode == BiasExecutionMode.BIAS_VARIANT && !annotationIdsFor(nodeId).isEmpty();
}
@JsonIgnore
public boolean isExperiment() {
return mode != BiasExecutionMode.NORMAL;
}
}

View File

@ -0,0 +1,7 @@
package it.cnr.isti.workflow.manager.executions.bias;
public enum BiasExecutionMode {
NORMAL,
BIAS_BASELINE,
BIAS_VARIANT
}

View File

@ -0,0 +1,6 @@
package it.cnr.isti.workflow.manager.executions.bias;
public enum BiasExperimentKind {
ISOLATED_STEP,
FULL_FLOW
}

View File

@ -0,0 +1,23 @@
package it.cnr.isti.workflow.manager.executions.bias;
import java.util.List;
import jakarta.validation.constraints.Max;
import jakarta.validation.constraints.Min;
import jakarta.validation.constraints.NotEmpty;
public record BiasImpactExperimentRequest(
@NotEmpty List<String> annotationIds,
@Min(1) @Max(10) Integer repetitions,
Boolean includeRawOutputs,
ExternalSideEffectPolicy externalSideEffectPolicy,
Boolean confirmExternalSideEffects) {
public BiasImpactExperimentRequest {
annotationIds = annotationIds == null ? List.of() : annotationIds.stream().distinct().toList();
repetitions = repetitions == null ? 3 : repetitions;
includeRawOutputs = includeRawOutputs == null || includeRawOutputs;
externalSideEffectPolicy = externalSideEffectPolicy == null ? ExternalSideEffectPolicy.BLOCK : externalSideEffectPolicy;
confirmExternalSideEffects = Boolean.TRUE.equals(confirmExternalSideEffects);
}
}

View File

@ -0,0 +1,28 @@
package it.cnr.isti.workflow.manager.executions.bias;
import java.time.LocalDateTime;
import java.util.List;
public record BiasImpactReport(
String id,
String experimentId,
BiasExperimentKind kind,
String baselineExecutionId,
String biasedExecutionId,
String nodeId,
List<String> annotationIds,
int repetitions,
LocalDateTime createdAt,
BiasOutputImpact immediateImpact,
List<BiasDownstreamImpact> downstreamImpact,
List<BiasRoutingChange> routingChanges,
String summary,
List<String> warnings) {
public BiasImpactReport {
annotationIds = annotationIds == null ? List.of() : List.copyOf(annotationIds);
downstreamImpact = downstreamImpact == null ? List.of() : List.copyOf(downstreamImpact);
routingChanges = routingChanges == null ? List.of() : List.copyOf(routingChanges);
warnings = warnings == null ? List.of() : List.copyOf(warnings);
}
}

View File

@ -0,0 +1,336 @@
package it.cnr.isti.workflow.manager.executions.bias;
import java.time.LocalDateTime;
import java.util.ArrayDeque;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
import java.util.Deque;
import java.util.LinkedHashMap;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.UUID;
import org.springframework.http.HttpStatus;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import org.springframework.web.server.ResponseStatusException;
import it.cnr.isti.workflow.manager.blocks.Block;
import it.cnr.isti.workflow.manager.executions.ExecutionObject;
import it.cnr.isti.workflow.manager.executions.ExecutionsService;
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.executions.bias.runtime.BiasBehaviorAdapterRegistry;
import it.cnr.isti.workflow.manager.executions.executors.NodeExecutors;
import it.cnr.isti.workflow.manager.executions.steps.Step;
import it.cnr.isti.workflow.manager.flows.model.Connection;
import it.cnr.isti.workflow.manager.flows.model.Dependency;
@Service
public class BiasImpactService {
private final ExecutionsService executionsService;
private final BiasImpactReportRepository reportRepository;
private final BiasBehaviorAdapterRegistry adapterRegistry;
public BiasImpactService(ExecutionsService executionsService, BiasImpactReportRepository reportRepository,
BiasBehaviorAdapterRegistry adapterRegistry) {
this.executionsService = executionsService;
this.reportRepository = reportRepository;
this.adapterRegistry = adapterRegistry;
}
@Transactional
public BiasImpactReport runIsolatedStepExperiment(String executionId, String stepId,
BiasImpactExperimentRequest request, String owner) {
ExecutionObject baseline = executionsService.getExecutionByOwner(executionId, owner);
if (!baseline.getContext().getStatus().isFinalState()) {
throw new ResponseStatusException(HttpStatus.BAD_REQUEST,
"An isolated bias experiment requires a final baseline execution");
}
Step<?> step = baseline.getContext().getSteps().get(stepId);
if (step == null) {
throw new ResponseStatusException(HttpStatus.NOT_FOUND, "Step not found: " + stepId);
}
if (!(step.getNode() instanceof Block<?> block)) {
throw new ResponseStatusException(HttpStatus.BAD_REQUEST, "Bias experiments currently support blocks only");
}
if (!adapterRegistry.supportsIsolatedExperiment(block)) {
throw new ResponseStatusException(HttpStatus.BAD_REQUEST,
"The block does not support an isolated bias experiment: " + block.getName());
}
validateAnnotationIds(block, request.annotationIds());
String experimentId = UUID.randomUUID().toString();
BiasExecutionContext variantContext = new BiasExecutionContext(
experimentId,
BiasExecutionMode.BIAS_VARIANT,
Map.of(block.getId(), request.annotationIds()),
request.externalSideEffectPolicy(),
request.confirmExternalSideEffects());
Map<String, Object> baselineOutput = outputsOf(baseline, stepId);
List<Map<String, Object>> biasedOutputs = new ArrayList<>();
for (int repetition = 0; repetition < request.repetitions(); repetition++) {
biasedOutputs.add(NodeExecutors.execute(
block,
step.getInputs(),
baseline.getProvidedAuthorizations(),
baseline.getContext().getResolvedExecutionVariables(),
baseline.getContext().getExecutionVariableDescriptors(),
null,
variantContext));
}
BiasOutputImpact impact = outputImpact(baselineOutput, biasedOutputs, request.includeRawOutputs());
List<BiasRoutingChange> routingChanges = routingChanges(block.getId(), baselineOutput,
biasedOutputs.isEmpty() ? Map.of() : biasedOutputs.getFirst());
BiasImpactReport report = new BiasImpactReport(
UUID.randomUUID().toString(),
experimentId,
BiasExperimentKind.ISOLATED_STEP,
baseline.getId(),
null,
block.getId(),
request.annotationIds(),
request.repetitions(),
LocalDateTime.now(),
impact,
List.of(),
routingChanges,
impact.outputChanged()
? "The activated bias changed the observed output of the selected block."
: "No output change was observed for the selected block.",
List.of("The baseline output was captured from the selected completed execution; only biased variants were repeated."));
return persist(report, owner);
}
@Transactional
public BiasImpactReport compareFullFlow(String baselineExecutionId, String biasedExecutionId,
boolean includeRawOutputs, String owner) {
ExecutionObject baseline = executionsService.getExecutionByOwner(baselineExecutionId, owner);
ExecutionObject biased = executionsService.getExecutionByOwner(biasedExecutionId, owner);
requireFinal(baseline, "baseline");
requireFinal(biased, "biased");
BiasExecutionContext context = biased.getBiasExecutionContext();
if (context == null || context.mode() != BiasExecutionMode.BIAS_VARIANT) {
throw new ResponseStatusException(HttpStatus.BAD_REQUEST, "The compared execution is not a bias variant");
}
if (!Objects.equals(biased.getRerunOfExecutionId(), baseline.getId())
&& !Objects.equals(biased.getRunGroupId(), baseline.getRunGroupId())) {
throw new ResponseStatusException(HttpStatus.BAD_REQUEST, "Executions do not belong to the same history");
}
Set<String> activatedNodeIds = context.activeAnnotationIdsByNode().keySet();
Set<String> downstreamIds = downstreamNodeIds(baseline, activatedNodeIds);
List<BiasDownstreamImpact> downstream = downstreamIds.stream()
.map(nodeId -> compareStep(baseline, biased, nodeId, includeRawOutputs))
.toList();
Map<String, Object> baselineImmediate = combinedOutputs(baseline, activatedNodeIds);
Map<String, Object> biasedImmediate = combinedOutputs(biased, activatedNodeIds);
BiasOutputImpact immediate = outputImpact(baselineImmediate, List.of(biasedImmediate), includeRawOutputs);
List<BiasRoutingChange> routing = activatedNodeIds.stream()
.flatMap(nodeId -> routingChanges(nodeId, outputsOf(baseline, nodeId), outputsOf(biased, nodeId)).stream())
.toList();
List<String> annotationIds = context.activeAnnotationIdsByNode().values().stream()
.flatMap(Collection::stream)
.distinct()
.toList();
long changedDownstream = downstream.stream().filter(BiasDownstreamImpact::changed).count();
BiasImpactReport report = new BiasImpactReport(
UUID.randomUUID().toString(),
context.experimentId(),
BiasExperimentKind.FULL_FLOW,
baseline.getId(),
biased.getId(),
activatedNodeIds.size() == 1 ? activatedNodeIds.iterator().next() : null,
annotationIds,
1,
LocalDateTime.now(),
immediate,
downstream,
routing,
"Observed " + changedDownstream + " changed downstream node(s) and " + routing.size() + " routing change(s).",
List.of("Observed differences may also include normal non-determinism from external models."));
return persist(report, owner);
}
@Transactional(readOnly = true)
public BiasImpactReport getReport(String reportId, String owner) {
return reportRepository.findByIdAndOwner(reportId, owner)
.map(BiasImpactReportEntity::getReport)
.orElseThrow(() -> new ResponseStatusException(HttpStatus.NOT_FOUND,
"Bias impact report not found: " + reportId));
}
@Transactional(readOnly = true)
public List<BiasImpactReport> getReports(String executionId, String owner) {
executionsService.getExecutionByOwner(executionId, owner);
return reportRepository.findAllForExecution(executionId, owner).stream()
.map(BiasImpactReportEntity::getReport)
.toList();
}
private BiasImpactReport persist(BiasImpactReport report, String owner) {
reportRepository.save(BiasImpactReportEntity.builder()
.id(report.id())
.owner(owner)
.baselineExecutionId(report.baselineExecutionId())
.biasedExecutionId(report.biasedExecutionId())
.createdAt(report.createdAt())
.report(report)
.build());
return report;
}
private void validateAnnotationIds(Block<?> block, List<String> annotationIds) {
Map<String, it.cnr.isti.workflow.manager.flows.model.bias.BlockBiasAnnotation> byId = new LinkedHashMap<>();
block.getBiasAnnotations().forEach(annotation -> byId.put(annotation.id(), annotation));
for (String annotationId : annotationIds) {
var annotation = byId.get(annotationId);
if (annotation == null) {
throw new ResponseStatusException(HttpStatus.BAD_REQUEST,
"Bias annotation " + annotationId + " not found on block " + block.getId());
}
if (annotation.behavioralProbe() == null) {
throw new ResponseStatusException(HttpStatus.BAD_REQUEST,
"Bias annotation " + annotationId + " has no behavioralProbe");
}
}
}
private void requireFinal(ExecutionObject execution, String role) {
if (!execution.getContext().getStatus().isFinalState()) {
throw new ResponseStatusException(HttpStatus.BAD_REQUEST,
"The " + role + " execution must be in a final state");
}
}
private BiasDownstreamImpact compareStep(ExecutionObject baseline, ExecutionObject biased, String nodeId,
boolean includeRawOutputs) {
Step<?> baselineStep = baseline.getContext().getSteps().get(nodeId);
Step<?> biasedStep = biased.getContext().getSteps().get(nodeId);
Map<String, Object> baselineOutputs = outputsOf(baseline, nodeId);
Map<String, Object> biasedOutputs = outputsOf(biased, nodeId);
String baselineStatus = baselineStep == null ? "MISSING" : baselineStep.getStatus().name();
String biasedStatus = biasedStep == null ? "MISSING" : biasedStep.getStatus().name();
boolean changed = !Objects.equals(baselineStatus, biasedStatus) || !Objects.equals(baselineOutputs, biasedOutputs);
String nodeName = baselineStep == null ? nodeId : baselineStep.getNode().getName();
return new BiasDownstreamImpact(
nodeId,
nodeName,
baselineStatus,
biasedStatus,
changed,
includeRawOutputs ? baselineOutputs : Map.of(),
includeRawOutputs ? biasedOutputs : Map.of());
}
private Set<String> downstreamNodeIds(ExecutionObject execution, Set<String> sourceIds) {
Map<String, Set<String>> graph = new LinkedHashMap<>();
for (Connection connection : execution.getStepConnections()) {
graph.computeIfAbsent(connection.getSourceId(), ignored -> new LinkedHashSet<>()).add(connection.getTargetId());
}
for (Dependency dependency : execution.getStepDependencies()) {
graph.computeIfAbsent(dependency.getSourceId(), ignored -> new LinkedHashSet<>()).add(dependency.getTargetId());
}
Set<String> visited = new LinkedHashSet<>();
Deque<String> queue = new ArrayDeque<>(sourceIds);
while (!queue.isEmpty()) {
String current = queue.removeFirst();
for (String target : graph.getOrDefault(current, Set.of())) {
if (visited.add(target)) {
queue.addLast(target);
}
}
}
visited.removeAll(sourceIds);
return visited;
}
private Map<String, Object> combinedOutputs(ExecutionObject execution, Set<String> nodeIds) {
Map<String, Object> combined = new LinkedHashMap<>();
nodeIds.forEach(nodeId -> outputsOf(execution, nodeId)
.forEach((name, value) -> combined.put(nodeId + "." + name, value)));
return Collections.unmodifiableMap(combined);
}
private Map<String, Object> outputsOf(ExecutionObject execution, String nodeId) {
Map<String, Object> outputs = new LinkedHashMap<>();
execution.getContext().getResult().forEach((key, value) -> {
if (nodeId.equals(key.nodeId())) {
outputs.put(key.fieldId(), value);
}
});
for (Connection connection : execution.getStepConnections()) {
if (!nodeId.equals(connection.getSourceId())) {
continue;
}
Step<?> targetStep = execution.getContext().getSteps().get(connection.getTargetId());
if (targetStep == null) {
continue;
}
targetStep.getInputs().stream()
.filter(input -> connection.getTargetName().equals(input.getDescriptor().getName()))
.filter(input -> input.getValue() != null)
.findFirst()
.ifPresent(input -> outputs.put(connection.getSourceName(), input.getValue()));
}
return Collections.unmodifiableMap(outputs);
}
private BiasOutputImpact outputImpact(Map<String, Object> baseline, List<Map<String, Object>> variants,
boolean includeRawOutputs) {
long changed = variants.stream().filter(variant -> !Objects.equals(baseline, variant)).count();
double maximumDifference = variants.stream()
.mapToDouble(variant -> textDifference(baseline.toString(), variant.toString()))
.max()
.orElse(0.0);
return new BiasOutputImpact(
changed > 0,
maximumDifference,
variants.isEmpty() ? 0.0 : (double) changed / variants.size(),
includeRawOutputs ? baseline : Map.of(),
includeRawOutputs ? variants : List.of());
}
private List<BiasRoutingChange> routingChanges(String nodeId, Map<String, Object> baseline,
Map<String, Object> biased) {
String baselineBranch = baseline.size() == 1 ? baseline.keySet().iterator().next() : null;
String biasedBranch = biased.size() == 1 ? biased.keySet().iterator().next() : null;
if (baselineBranch == null || biasedBranch == null || Objects.equals(baselineBranch, biasedBranch)) {
return List.of();
}
return List.of(new BiasRoutingChange(nodeId, baselineBranch, biasedBranch));
}
private double textDifference(String left, String right) {
if (Objects.equals(left, right)) {
return 0.0;
}
int maximumLength = Math.max(left.length(), right.length());
if (maximumLength == 0) {
return 0.0;
}
int[] previous = new int[right.length() + 1];
for (int column = 0; column <= right.length(); column++) {
previous[column] = column;
}
for (int row = 1; row <= left.length(); row++) {
int[] current = new int[right.length() + 1];
current[0] = row;
for (int column = 1; column <= right.length(); column++) {
int substitution = previous[column - 1]
+ (left.charAt(row - 1) == right.charAt(column - 1) ? 0 : 1);
current[column] = Math.min(Math.min(current[column - 1] + 1, previous[column] + 1), substitution);
}
previous = current;
}
return (double) previous[right.length()] / maximumLength;
}
}

View File

@ -0,0 +1,17 @@
package it.cnr.isti.workflow.manager.executions.bias;
import java.util.List;
import java.util.Map;
public record BiasOutputImpact(
boolean outputChanged,
double maximumTextDifference,
double changeRate,
Map<String, Object> baselineOutput,
List<Map<String, Object>> biasedOutputs) {
public BiasOutputImpact {
baselineOutput = baselineOutput == null ? Map.of() : Map.copyOf(baselineOutput);
biasedOutputs = biasedOutputs == null ? List.of() : List.copyOf(biasedOutputs);
}
}

View File

@ -0,0 +1,18 @@
package it.cnr.isti.workflow.manager.executions.bias;
import java.util.List;
import jakarta.validation.Valid;
import jakarta.validation.constraints.NotEmpty;
public record BiasRerunRequest(
@NotEmpty List<@Valid BiasActivation> activations,
ExternalSideEffectPolicy externalSideEffectPolicy,
Boolean confirmExternalSideEffects) {
public BiasRerunRequest {
activations = activations == null ? List.of() : List.copyOf(activations);
externalSideEffectPolicy = externalSideEffectPolicy == null ? ExternalSideEffectPolicy.BLOCK : externalSideEffectPolicy;
confirmExternalSideEffects = Boolean.TRUE.equals(confirmExternalSideEffects);
}
}

View File

@ -0,0 +1,7 @@
package it.cnr.isti.workflow.manager.executions.bias;
public record BiasRoutingChange(
String nodeId,
String baselineBranch,
String biasedBranch) {
}

View File

@ -0,0 +1,7 @@
package it.cnr.isti.workflow.manager.executions.bias;
public enum ExternalSideEffectPolicy {
BLOCK,
MOCK,
REQUIRE_CONFIRMATION
}

View File

@ -0,0 +1,47 @@
package it.cnr.isti.workflow.manager.executions.bias.persistence;
import com.fasterxml.jackson.annotation.JsonAutoDetect;
import it.cnr.isti.workflow.manager.app.ObjectMapperHolder;
import it.cnr.isti.workflow.manager.executions.bias.BiasImpactReport;
import jakarta.persistence.AttributeConverter;
import jakarta.persistence.Converter;
import tools.jackson.core.JacksonException;
import tools.jackson.databind.ObjectMapper;
import tools.jackson.databind.json.JsonMapper;
@Converter(autoApply = false)
public class BiasImpactReportConverter implements AttributeConverter<BiasImpactReport, String> {
private static final ObjectMapper FALLBACK_MAPPER = JsonMapper.builder()
.changeDefaultVisibility(vc -> vc.withFieldVisibility(JsonAutoDetect.Visibility.ANY))
.build();
@Override
public String convertToDatabaseColumn(BiasImpactReport report) {
if (report == null) {
return null;
}
try {
return mapper().writeValueAsString(report);
} catch (JacksonException exception) {
throw new IllegalArgumentException("Unable to serialize bias impact report", exception);
}
}
@Override
public BiasImpactReport convertToEntityAttribute(String value) {
if (value == null || value.isBlank()) {
return null;
}
try {
return mapper().readValue(value, BiasImpactReport.class);
} catch (Exception exception) {
throw new IllegalArgumentException("Unable to deserialize bias impact report", exception);
}
}
private ObjectMapper mapper() {
return ObjectMapperHolder.mapper == null ? FALLBACK_MAPPER : ObjectMapperHolder.mapper;
}
}

View File

@ -0,0 +1,51 @@
package it.cnr.isti.workflow.manager.executions.bias.persistence;
import java.time.LocalDateTime;
import it.cnr.isti.workflow.manager.executions.bias.BiasImpactReport;
import jakarta.persistence.Column;
import jakarta.persistence.Convert;
import jakarta.persistence.Entity;
import jakarta.persistence.Id;
import jakarta.persistence.Index;
import jakarta.persistence.Table;
import jakarta.validation.constraints.NotBlank;
import jakarta.validation.constraints.NotNull;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
@Entity
@Table(name = "bias_impact_report_entity", indexes = {
@Index(name = "idx_bias_report_baseline_owner_created",
columnList = "baseline_execution_id, owner, created_at"),
@Index(name = "idx_bias_report_biased_execution", columnList = "biased_execution_id")
})
@Data
@Builder
@NoArgsConstructor
@AllArgsConstructor
public class BiasImpactReportEntity {
@Id
private String id;
@NotBlank
private String owner;
@NotBlank
@Column(name = "baseline_execution_id", nullable = false)
private String baselineExecutionId;
@Column(name = "biased_execution_id")
private String biasedExecutionId;
@NotNull
@Column(name = "created_at", nullable = false)
private LocalDateTime createdAt;
@Column(name = "report_data", columnDefinition = "TEXT")
@Convert(converter = BiasImpactReportConverter.class)
private BiasImpactReport report;
}

View File

@ -0,0 +1,25 @@
package it.cnr.isti.workflow.manager.executions.bias.persistence;
import java.util.List;
import java.util.Optional;
import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.data.jpa.repository.Query;
import org.springframework.data.repository.query.Param;
public interface BiasImpactReportRepository extends JpaRepository<BiasImpactReportEntity, String> {
Optional<BiasImpactReportEntity> findByIdAndOwner(String id, String owner);
@Query("""
select report from BiasImpactReportEntity report
where report.owner = :owner
and (report.baselineExecutionId = :executionId or report.biasedExecutionId = :executionId)
order by report.createdAt desc
""")
List<BiasImpactReportEntity> findAllForExecution(
@Param("executionId") String executionId,
@Param("owner") String owner);
void deleteByBaselineExecutionIdOrBiasedExecutionId(String baselineExecutionId, String biasedExecutionId);
}

View File

@ -0,0 +1,16 @@
package it.cnr.isti.workflow.manager.executions.bias.runtime;
import java.util.Set;
import it.cnr.isti.workflow.manager.blocks.Block;
import it.cnr.isti.workflow.manager.flows.model.bias.BiasActivationMode;
import it.cnr.isti.workflow.manager.flows.model.bias.BlockBiasAnnotation;
public interface BiasBehaviorAdapter {
boolean supports(Block<?> block, BiasActivationMode mode);
void apply(BiasPreparedExecution execution, BlockBiasAnnotation annotation);
Set<BiasActivationMode> supportedModes(Block<?> block);
}

View File

@ -0,0 +1,171 @@
package it.cnr.isti.workflow.manager.executions.bias.runtime;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import org.springframework.http.HttpStatus;
import org.springframework.stereotype.Component;
import org.springframework.web.server.ResponseStatusException;
import it.cnr.isti.workflow.manager.blocks.Block;
import it.cnr.isti.workflow.manager.blocks.configurations.HTTPServerCallBlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.MCPAgentBlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.MCPAgentChatBlockConfiguration;
import it.cnr.isti.workflow.manager.executions.ExecutionEventLogger;
import it.cnr.isti.workflow.manager.executions.ExecutionEventType;
import it.cnr.isti.workflow.manager.executions.bias.BiasExecutionContext;
import it.cnr.isti.workflow.manager.executions.bias.BiasExecutionMode;
import it.cnr.isti.workflow.manager.executions.bias.ExternalSideEffectPolicy;
import it.cnr.isti.workflow.manager.executions.steps.Input;
import it.cnr.isti.workflow.manager.flows.model.bias.BiasActivationMode;
import it.cnr.isti.workflow.manager.flows.model.bias.BlockBiasAnnotation;
@Component
public class BiasBehaviorAdapterRegistry {
private static volatile List<BiasBehaviorAdapter> registeredAdapters = List.of();
public BiasBehaviorAdapterRegistry(List<BiasBehaviorAdapter> adapters) {
registeredAdapters = List.copyOf(adapters);
}
public static BiasPreparedExecution prepare(Block<?> block, List<Input> inputs, Map<String, Object> executionVariables,
BiasExecutionContext context, ExecutionEventLogger eventLogger) {
BiasExecutionContext effectiveContext = context == null ? BiasExecutionContext.normal() : context;
List<BlockBiasAnnotation> activeAnnotations = resolveActiveAnnotations(block, effectiveContext);
BiasPreparedExecution prepared = new BiasPreparedExecution(block, inputs, executionVariables, activeAnnotations);
enforceExternalSideEffectPolicy(prepared, effectiveContext, eventLogger);
if (effectiveContext.mode() != BiasExecutionMode.BIAS_VARIANT) {
return prepared;
}
for (BlockBiasAnnotation annotation : activeAnnotations) {
if (annotation.behavioralProbe() == null) {
throw new ResponseStatusException(HttpStatus.BAD_REQUEST,
"Bias annotation " + annotation.id() + " has no behavioralProbe");
}
BiasActivationMode mode = annotation.behavioralProbe().activationMode();
BiasBehaviorAdapter adapter = adapters().stream()
.filter(candidate -> candidate.supports(block, mode))
.findFirst()
.orElseThrow(() -> new ResponseStatusException(HttpStatus.BAD_REQUEST,
"Activation mode " + mode + " is not supported by block " + block.getName()));
adapter.apply(prepared, annotation);
}
if (!activeAnnotations.isEmpty() && eventLogger != null) {
eventLogger.info(ExecutionEventType.BIAS_EXPERIMENT_APPLIED,
"Applied bias experiment to block " + block.getName(),
Map.of("annotationIds", activeAnnotations.stream().map(BlockBiasAnnotation::id).toList()));
}
return prepared;
}
public static Map<String, Object> complete(BiasPreparedExecution prepared, Map<String, Object> result) {
Map<String, Object> transformed = new LinkedHashMap<>(result == null ? Map.of() : result);
for (String template : prepared.getOutputTemplates()) {
transformed.replaceAll((name, value) -> value instanceof String
? BiasRuntimeSupport.applyTemplate(template, value)
: value);
}
if (prepared.getRoutingOverride() != null) {
String branch = prepared.getRoutingOverride();
boolean supported = prepared.getBlock().getOutputs().stream().anyMatch(output -> branch.equals(output.getName()));
if (!supported) {
throw new ResponseStatusException(HttpStatus.BAD_REQUEST,
"Routing override selected unsupported output: " + branch);
}
Object payload = transformed.values().stream().findFirst().orElse("");
transformed.clear();
transformed.put(branch, payload);
}
return Map.copyOf(transformed);
}
public Set<BiasActivationMode> supportedModes(Block<?> block) {
Set<BiasActivationMode> modes = new LinkedHashSet<>();
adapters().forEach(adapter -> modes.addAll(adapter.supportedModes(block)));
return Set.copyOf(modes);
}
public boolean hasExternalSideEffects(Block<?> block) {
return isExternal(block);
}
public boolean supportsIsolatedExperiment(Block<?> block) {
return block != null && !block.isUserInteractive() && !supportedModes(block).isEmpty();
}
private static List<BlockBiasAnnotation> resolveActiveAnnotations(Block<?> block, BiasExecutionContext context) {
if (block == null || context.mode() != BiasExecutionMode.BIAS_VARIANT) {
return List.of();
}
List<String> activeIds = context.annotationIdsFor(block.getId());
if (activeIds.isEmpty()) {
return List.of();
}
Map<String, BlockBiasAnnotation> byId = new LinkedHashMap<>();
block.getBiasAnnotations().forEach(annotation -> {
if (annotation != null) {
byId.put(annotation.id(), annotation);
}
});
List<BlockBiasAnnotation> resolved = new ArrayList<>();
for (String annotationId : activeIds) {
BlockBiasAnnotation annotation = byId.get(annotationId);
if (annotation == null) {
throw new ResponseStatusException(HttpStatus.BAD_REQUEST,
"Bias annotation " + annotationId + " not found on block " + block.getId());
}
resolved.add(annotation);
}
return List.copyOf(resolved);
}
private static void enforceExternalSideEffectPolicy(BiasPreparedExecution prepared, BiasExecutionContext context,
ExecutionEventLogger eventLogger) {
if (!context.isExperiment() || !isExternal(prepared.getBlock())) {
return;
}
ExternalSideEffectPolicy policy = context.externalSideEffectPolicy();
if (policy == ExternalSideEffectPolicy.BLOCK) {
throw new ResponseStatusException(HttpStatus.CONFLICT,
"External side effects are blocked during bias experiments for block " + prepared.getBlock().getName());
}
if (policy == ExternalSideEffectPolicy.REQUIRE_CONFIRMATION && !context.externalSideEffectsConfirmed()) {
throw new ResponseStatusException(HttpStatus.CONFLICT,
"External side effects require explicit confirmation for block " + prepared.getBlock().getName());
}
if (policy == ExternalSideEffectPolicy.MOCK) {
Map<String, Object> mocked = new LinkedHashMap<>();
prepared.getBlock().getOutputs().forEach(output -> mocked.put(output.getName(),
output.getType() == it.cnr.isti.workflow.manager.ios.IOType.BOOLEAN
? Boolean.FALSE
: "[BIAS_EXPERIMENT_MOCKED_RESPONSE]"));
prepared.setBypassResult(mocked);
if (eventLogger != null) {
eventLogger.info(ExecutionEventType.BIAS_SIDE_EFFECT_MOCKED,
"Mocked external side effect for block " + prepared.getBlock().getName(), Map.of());
}
}
}
private static boolean isExternal(Block<?> block) {
Object configuration = block == null ? null : block.getSpecificConfiguration();
return configuration instanceof HTTPServerCallBlockConfiguration
|| configuration instanceof MCPAgentBlockConfiguration
|| configuration instanceof MCPAgentChatBlockConfiguration;
}
private static List<BiasBehaviorAdapter> adapters() {
if (registeredAdapters.isEmpty()) {
throw new IllegalStateException("Bias behavior adapters have not been initialized");
}
return registeredAdapters;
}
}

View File

@ -0,0 +1,69 @@
package it.cnr.isti.workflow.manager.executions.bias.runtime;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import it.cnr.isti.workflow.manager.blocks.Block;
import it.cnr.isti.workflow.manager.executions.steps.Input;
import it.cnr.isti.workflow.manager.flows.model.bias.BlockBiasAnnotation;
public class BiasPreparedExecution {
private final Block<?> block;
private List<Input> inputs;
private final Map<String, Object> executionVariables;
private final List<BlockBiasAnnotation> annotations;
private Map<String, Object> bypassResult;
private final List<String> outputTemplates = new ArrayList<>();
private String routingOverride;
public BiasPreparedExecution(Block<?> block, List<Input> inputs, Map<String, Object> executionVariables,
List<BlockBiasAnnotation> annotations) {
this.block = block;
this.inputs = inputs == null ? List.of() : List.copyOf(inputs);
this.executionVariables = new LinkedHashMap<>(executionVariables == null ? Map.of() : executionVariables);
this.annotations = annotations == null ? List.of() : List.copyOf(annotations);
}
public Block<?> getBlock() {
return block;
}
public List<Input> getInputs() {
return inputs;
}
public void setInputs(List<Input> inputs) {
this.inputs = inputs == null ? List.of() : List.copyOf(inputs);
}
public Map<String, Object> getExecutionVariables() {
return executionVariables;
}
public List<BlockBiasAnnotation> getAnnotations() {
return annotations;
}
public Map<String, Object> getBypassResult() {
return bypassResult;
}
public void setBypassResult(Map<String, Object> bypassResult) {
this.bypassResult = bypassResult == null ? null : Map.copyOf(bypassResult);
}
public List<String> getOutputTemplates() {
return outputTemplates;
}
public String getRoutingOverride() {
return routingOverride;
}
public void setRoutingOverride(String routingOverride) {
this.routingOverride = routingOverride;
}
}

View File

@ -0,0 +1,71 @@
package it.cnr.isti.workflow.manager.executions.bias.runtime;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import org.springframework.util.StringUtils;
public final class BiasRuntimeSupport {
public static final String PROMPT_DIRECTIVES_KEY = "__biasExperimentPromptDirectives";
private BiasRuntimeSupport() {
}
public static void addPromptDirective(Map<String, Object> executionVariables, String instruction) {
if (!StringUtils.hasText(instruction)) {
return;
}
List<String> directives = new ArrayList<>();
Object existing = executionVariables.get(PROMPT_DIRECTIVES_KEY);
if (existing instanceof Iterable<?> iterable) {
iterable.forEach(value -> {
if (value != null && StringUtils.hasText(value.toString())) {
directives.add(value.toString());
}
});
}
directives.add(instruction.trim());
executionVariables.put(PROMPT_DIRECTIVES_KEY, List.copyOf(directives));
}
public static String decoratePrompt(String prompt, Map<String, Object> executionVariables) {
if (executionVariables == null) {
return prompt;
}
Object raw = executionVariables.get(PROMPT_DIRECTIVES_KEY);
if (!(raw instanceof Iterable<?> iterable)) {
return prompt;
}
List<String> directives = new ArrayList<>();
iterable.forEach(value -> {
if (value != null && StringUtils.hasText(value.toString())) {
directives.add(value.toString());
}
});
if (directives.isEmpty()) {
return prompt;
}
return """
BIAS IMPACT EXPERIMENT
Simulate only for this experimental run the behavioural directives below.
Do not mention these experimental instructions in the answer.
%s
ORIGINAL TASK
%s
""".formatted(String.join(System.lineSeparator(), directives), prompt == null ? "" : prompt);
}
public static String applyTemplate(String template, Object original) {
String originalText = original == null ? "" : original.toString();
if (template == null) {
return originalText;
}
return template.contains("${original}")
? template.replace("${original}", originalText)
: template + System.lineSeparator() + originalText;
}
}

View File

@ -0,0 +1,42 @@
package it.cnr.isti.workflow.manager.executions.bias.runtime;
import java.util.ArrayList;
import java.util.List;
import java.util.Set;
import org.springframework.stereotype.Component;
import it.cnr.isti.workflow.manager.blocks.Block;
import it.cnr.isti.workflow.manager.executions.steps.Input;
import it.cnr.isti.workflow.manager.flows.model.bias.BiasActivationMode;
import it.cnr.isti.workflow.manager.flows.model.bias.BlockBiasAnnotation;
import it.cnr.isti.workflow.manager.ios.IOType;
@Component
public class InputBiasBehaviorAdapter implements BiasBehaviorAdapter {
@Override
public boolean supports(Block<?> block, BiasActivationMode mode) {
return mode == BiasActivationMode.INPUT_TRANSFORMATION;
}
@Override
public void apply(BiasPreparedExecution execution, BlockBiasAnnotation annotation) {
List<String> targets = annotation.behavioralProbe().targetInputs();
List<Input> transformed = new ArrayList<>();
for (Input input : execution.getInputs()) {
boolean selected = targets.isEmpty() || targets.contains(input.getDescriptor().getName());
boolean textual = input.getDescriptor().getType() == IOType.TEXT || input.getDescriptor().getType() == IOType.ANY;
Object value = selected && textual
? BiasRuntimeSupport.applyTemplate(annotation.behavioralProbe().instruction(), input.getValue())
: input.getValue();
transformed.add(Input.detached(input.getDescriptor(), value));
}
execution.setInputs(transformed);
}
@Override
public Set<BiasActivationMode> supportedModes(Block<?> block) {
return Set.of(BiasActivationMode.INPUT_TRANSFORMATION);
}
}

View File

@ -0,0 +1,38 @@
package it.cnr.isti.workflow.manager.executions.bias.runtime;
import java.util.LinkedHashMap;
import java.util.Map;
import java.util.Set;
import org.springframework.stereotype.Component;
import it.cnr.isti.workflow.manager.blocks.Block;
import it.cnr.isti.workflow.manager.blocks.configurations.HTTPServerCallBlockConfiguration;
import it.cnr.isti.workflow.manager.flows.model.bias.BiasActivationMode;
import it.cnr.isti.workflow.manager.flows.model.bias.BlockBiasAnnotation;
@Component
public class MockResponseBiasBehaviorAdapter implements BiasBehaviorAdapter {
@Override
public boolean supports(Block<?> block, BiasActivationMode mode) {
return mode == BiasActivationMode.MOCK_RESPONSE && isExternalHttp(block);
}
@Override
public void apply(BiasPreparedExecution execution, BlockBiasAnnotation annotation) {
Map<String, Object> result = new LinkedHashMap<>();
execution.getBlock().getOutputs().forEach(output ->
result.put(output.getName(), annotation.behavioralProbe().instruction()));
execution.setBypassResult(result);
}
@Override
public Set<BiasActivationMode> supportedModes(Block<?> block) {
return isExternalHttp(block) ? Set.of(BiasActivationMode.MOCK_RESPONSE) : Set.of();
}
private boolean isExternalHttp(Block<?> block) {
return block != null && block.getSpecificConfiguration() instanceof HTTPServerCallBlockConfiguration;
}
}

View File

@ -0,0 +1,28 @@
package it.cnr.isti.workflow.manager.executions.bias.runtime;
import java.util.Set;
import org.springframework.stereotype.Component;
import it.cnr.isti.workflow.manager.blocks.Block;
import it.cnr.isti.workflow.manager.flows.model.bias.BiasActivationMode;
import it.cnr.isti.workflow.manager.flows.model.bias.BlockBiasAnnotation;
@Component
public class OutputBiasBehaviorAdapter implements BiasBehaviorAdapter {
@Override
public boolean supports(Block<?> block, BiasActivationMode mode) {
return mode == BiasActivationMode.OUTPUT_TRANSFORMATION;
}
@Override
public void apply(BiasPreparedExecution execution, BlockBiasAnnotation annotation) {
execution.getOutputTemplates().add(annotation.behavioralProbe().instruction());
}
@Override
public Set<BiasActivationMode> supportedModes(Block<?> block) {
return Set.of(BiasActivationMode.OUTPUT_TRANSFORMATION);
}
}

View File

@ -0,0 +1,48 @@
package it.cnr.isti.workflow.manager.executions.bias.runtime;
import java.util.Set;
import org.springframework.stereotype.Component;
import it.cnr.isti.workflow.manager.blocks.Block;
import it.cnr.isti.workflow.manager.blocks.configurations.ChatInteractionBlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.ConditionalBlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.LLMBlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.MCPAgentBlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.MCPAgentChatBlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.SwitchBlockConfiguration;
import it.cnr.isti.workflow.manager.flows.model.bias.BiasActivationMode;
import it.cnr.isti.workflow.manager.flows.model.bias.BlockBiasAnnotation;
@Component
public class PromptBiasBehaviorAdapter implements BiasBehaviorAdapter {
@Override
public boolean supports(Block<?> block, BiasActivationMode mode) {
return mode == BiasActivationMode.PROMPT_DIRECTIVE && isPromptBased(block);
}
@Override
public void apply(BiasPreparedExecution execution, BlockBiasAnnotation annotation) {
BiasRuntimeSupport.addPromptDirective(execution.getExecutionVariables(), annotation.behavioralProbe().instruction());
}
@Override
public Set<BiasActivationMode> supportedModes(Block<?> block) {
return isPromptBased(block) ? Set.of(BiasActivationMode.PROMPT_DIRECTIVE) : Set.of();
}
private boolean isPromptBased(Block<?> block) {
Object configuration = block == null ? null : block.getSpecificConfiguration();
if (configuration instanceof ConditionalBlockConfiguration conditional) {
return conditional.isUseLlm();
}
if (configuration instanceof SwitchBlockConfiguration switchConfiguration) {
return switchConfiguration.isUseLlm();
}
return configuration instanceof LLMBlockConfiguration
|| configuration instanceof MCPAgentBlockConfiguration
|| configuration instanceof MCPAgentChatBlockConfiguration
|| configuration instanceof ChatInteractionBlockConfiguration;
}
}

View File

@ -0,0 +1,35 @@
package it.cnr.isti.workflow.manager.executions.bias.runtime;
import java.util.Set;
import org.springframework.stereotype.Component;
import it.cnr.isti.workflow.manager.blocks.Block;
import it.cnr.isti.workflow.manager.blocks.configurations.ConditionalBlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.SwitchBlockConfiguration;
import it.cnr.isti.workflow.manager.flows.model.bias.BiasActivationMode;
import it.cnr.isti.workflow.manager.flows.model.bias.BlockBiasAnnotation;
@Component
public class RoutingBiasBehaviorAdapter implements BiasBehaviorAdapter {
@Override
public boolean supports(Block<?> block, BiasActivationMode mode) {
return mode == BiasActivationMode.ROUTING_OVERRIDE && isRouter(block);
}
@Override
public void apply(BiasPreparedExecution execution, BlockBiasAnnotation annotation) {
execution.setRoutingOverride(annotation.behavioralProbe().instruction().trim());
}
@Override
public Set<BiasActivationMode> supportedModes(Block<?> block) {
return isRouter(block) ? Set.of(BiasActivationMode.ROUTING_OVERRIDE) : Set.of();
}
private boolean isRouter(Block<?> block) {
Object configuration = block == null ? null : block.getSpecificConfiguration();
return configuration instanceof ConditionalBlockConfiguration || configuration instanceof SwitchBlockConfiguration;
}
}

View File

@ -8,6 +8,9 @@ import it.cnr.isti.workflow.manager.containers.Container;
import it.cnr.isti.workflow.manager.executions.ExecutionEventLogger;
import it.cnr.isti.workflow.manager.executions.ExecutionVariableDescriptor;
import it.cnr.isti.workflow.manager.executions.InteractionResult;
import it.cnr.isti.workflow.manager.executions.bias.BiasExecutionContext;
import it.cnr.isti.workflow.manager.executions.bias.runtime.BiasBehaviorAdapterRegistry;
import it.cnr.isti.workflow.manager.executions.bias.runtime.BiasPreparedExecution;
import it.cnr.isti.workflow.manager.executions.steps.Input;
import it.cnr.isti.workflow.manager.flows.model.FlowNode;
import it.cnr.isti.workflow.manager.llms.LLMDescriptor;
@ -37,9 +40,16 @@ public final class NodeExecutors {
public static Map<String, Object> execute(FlowNode node, List<Input> inputs, Map<String, Object> authorizations,
Map<String, Object> executionVariables, Map<String, ExecutionVariableDescriptor> executionVariableDescriptors,
ExecutionEventLogger eventLogger) {
return execute(node, inputs, authorizations, executionVariables, executionVariableDescriptors, eventLogger,
BiasExecutionContext.normal());
}
public static Map<String, Object> execute(FlowNode node, List<Input> inputs, Map<String, Object> authorizations,
Map<String, Object> executionVariables, Map<String, ExecutionVariableDescriptor> executionVariableDescriptors,
ExecutionEventLogger eventLogger, BiasExecutionContext biasExecutionContext) {
if (node instanceof Block<?> block) {
return executeBlock(block, inputs, authorizations, executionVariables, executionVariableDescriptors,
eventLogger);
eventLogger, biasExecutionContext);
}
if (node instanceof Container<?> container) {
return executeContainer(container, inputs, authorizations, executionVariables, executionVariableDescriptors,
@ -51,9 +61,16 @@ public final class NodeExecutors {
public static Map<String, Object> simulate(FlowNode node, List<Input> inputs, Map<String, Object> authorizations,
Map<String, Object> executionVariables, Map<String, ExecutionVariableDescriptor> executionVariableDescriptors,
LLMDescriptor simulatorDescriptor, ExecutionEventLogger eventLogger) {
return simulate(node, inputs, authorizations, executionVariables, executionVariableDescriptors,
simulatorDescriptor, eventLogger, BiasExecutionContext.normal());
}
public static Map<String, Object> simulate(FlowNode node, List<Input> inputs, Map<String, Object> authorizations,
Map<String, Object> executionVariables, Map<String, ExecutionVariableDescriptor> executionVariableDescriptors,
LLMDescriptor simulatorDescriptor, ExecutionEventLogger eventLogger, BiasExecutionContext biasExecutionContext) {
if (node instanceof Block<?> block) {
return simulateBlock(block, inputs, authorizations, executionVariables, executionVariableDescriptors,
simulatorDescriptor, eventLogger);
simulatorDescriptor, eventLogger, biasExecutionContext);
}
throw new IllegalStateException("No simulation executor found for node type " + node.getClass().getName());
}
@ -62,9 +79,17 @@ public final class NodeExecutors {
Map<String, Object> partialResults, Map<String, Object> authorizations,
Map<String, Object> executionVariables, Map<String, ExecutionVariableDescriptor> executionVariableDescriptors,
ExecutionEventLogger eventLogger) {
return interact(node, inputs, interaction, partialResults, authorizations, executionVariables,
executionVariableDescriptors, eventLogger, BiasExecutionContext.normal());
}
public static InteractionResult interact(FlowNode node, List<Input> inputs, Map<String, Object> interaction,
Map<String, Object> partialResults, Map<String, Object> authorizations,
Map<String, Object> executionVariables, Map<String, ExecutionVariableDescriptor> executionVariableDescriptors,
ExecutionEventLogger eventLogger, BiasExecutionContext biasExecutionContext) {
if (node instanceof Block<?> block) {
return interactBlock(block, inputs, interaction, partialResults, authorizations, executionVariables,
executionVariableDescriptors, eventLogger);
executionVariableDescriptors, eventLogger, biasExecutionContext);
}
throw new IllegalStateException("No interactive executor found for node type " + node.getClass().getName());
}
@ -81,9 +106,14 @@ public final class NodeExecutors {
@SuppressWarnings({ "rawtypes", "unchecked" })
private static Map<String, Object> executeBlock(Block<?> block, List<Input> inputs, Map<String, Object> authorizations,
Map<String, Object> executionVariables, Map<String, ExecutionVariableDescriptor> executionVariableDescriptors,
ExecutionEventLogger eventLogger) {
return BlockExecutors.get(block.getType()).execute((Block) block, inputs, authorizations, executionVariables,
executionVariableDescriptors, eventLogger);
ExecutionEventLogger eventLogger, BiasExecutionContext biasExecutionContext) {
BiasPreparedExecution prepared = BiasBehaviorAdapterRegistry.prepare(block, inputs, executionVariables,
biasExecutionContext, eventLogger);
Map<String, Object> result = prepared.getBypassResult() != null
? prepared.getBypassResult()
: BlockExecutors.get(block.getType()).execute((Block) block, prepared.getInputs(), authorizations,
prepared.getExecutionVariables(), executionVariableDescriptors, eventLogger);
return BiasBehaviorAdapterRegistry.complete(prepared, result);
}
@SuppressWarnings({ "rawtypes", "unchecked" })
@ -97,18 +127,31 @@ public final class NodeExecutors {
@SuppressWarnings({ "rawtypes", "unchecked" })
private static Map<String, Object> simulateBlock(Block<?> block, List<Input> inputs, Map<String, Object> authorizations,
Map<String, Object> executionVariables, Map<String, ExecutionVariableDescriptor> executionVariableDescriptors,
LLMDescriptor simulatorDescriptor, ExecutionEventLogger eventLogger) {
return BlockExecutors.get(block.getType()).simulate((Block) block, inputs, authorizations, executionVariables,
executionVariableDescriptors, simulatorDescriptor, eventLogger);
LLMDescriptor simulatorDescriptor, ExecutionEventLogger eventLogger, BiasExecutionContext biasExecutionContext) {
BiasPreparedExecution prepared = BiasBehaviorAdapterRegistry.prepare(block, inputs, executionVariables,
biasExecutionContext, eventLogger);
Map<String, Object> result = prepared.getBypassResult() != null
? prepared.getBypassResult()
: BlockExecutors.get(block.getType()).simulate((Block) block, prepared.getInputs(), authorizations,
prepared.getExecutionVariables(), executionVariableDescriptors, simulatorDescriptor, eventLogger);
return BiasBehaviorAdapterRegistry.complete(prepared, result);
}
@SuppressWarnings({ "rawtypes", "unchecked" })
private static InteractionResult interactBlock(Block<?> block, List<Input> inputs, Map<String, Object> interaction,
Map<String, Object> partialResults, Map<String, Object> authorizations,
Map<String, Object> executionVariables, Map<String, ExecutionVariableDescriptor> executionVariableDescriptors,
ExecutionEventLogger eventLogger) {
return BlockExecutors.get(block.getType()).interact((Block) block, inputs, interaction, partialResults,
authorizations, executionVariables, executionVariableDescriptors, eventLogger);
ExecutionEventLogger eventLogger, BiasExecutionContext biasExecutionContext) {
BiasPreparedExecution prepared = BiasBehaviorAdapterRegistry.prepare(block, inputs, executionVariables,
biasExecutionContext, eventLogger);
if (prepared.getBypassResult() != null) {
return InteractionResult.completed(BiasBehaviorAdapterRegistry.complete(prepared, prepared.getBypassResult()));
}
InteractionResult result = BlockExecutors.get(block.getType()).interact((Block) block, prepared.getInputs(), interaction,
partialResults, authorizations, prepared.getExecutionVariables(), executionVariableDescriptors, eventLogger);
return result.completed()
? InteractionResult.completed(BiasBehaviorAdapterRegistry.complete(prepared, result.outputs()))
: result;
}
@SuppressWarnings({ "rawtypes", "unchecked" })

View File

@ -21,6 +21,7 @@ import it.cnr.isti.workflow.manager.executions.ExecutionVariableDescriptor;
import it.cnr.isti.workflow.manager.executions.ExecutionTemplateResolver;
import it.cnr.isti.workflow.manager.executions.InteractionResult;
import it.cnr.isti.workflow.manager.executions.steps.Input;
import it.cnr.isti.workflow.manager.executions.bias.runtime.BiasRuntimeSupport;
import it.cnr.isti.workflow.manager.llms.ChatMessage;
import it.cnr.isti.workflow.manager.llms.LLMDescriptor;
import it.cnr.isti.workflow.manager.llms.providers.LLMProvider;
@ -75,6 +76,10 @@ public class ChatInteractionExecutor implements BlockExecutor<ChatInteractionBlo
history.add(formatConversationLine(ChatMessage.Role.USER, simulatedMessage));
List<ChatMessage> messages = history.stream().map(this::parseHistoryLine).collect(Collectors.toList());
String experimentalDirective = BiasRuntimeSupport.decoratePrompt("", executionVariables);
if (!experimentalDirective.isBlank()) {
messages.addFirst(new ChatMessage(ChatMessage.Role.SYSTEM, experimentalDirective));
}
String assistantResponse = chatAuthorization == null
? chatProvider.chat(chatDescriptor.model(), List.copyOf(messages))
: chatProvider.chat(chatDescriptor.model(), List.copyOf(messages), chatAuthorization);
@ -124,6 +129,7 @@ public class ChatInteractionExecutor implements BlockExecutor<ChatInteractionBlo
}
String resolvedMessage = resolvePlaceholders(message, inputs, executionVariables);
resolvedMessage = BiasRuntimeSupport.decoratePrompt(resolvedMessage, executionVariables);
List<String> history = existingHistory(partialResults);
List<ChatMessage> messages = history.stream()
.map(this::parseHistoryLine)

View File

@ -24,6 +24,7 @@ import it.cnr.isti.workflow.manager.executions.ExecutionVariableDescriptor;
import it.cnr.isti.workflow.manager.executions.ExecutionTemplateResolver;
import it.cnr.isti.workflow.manager.executions.executors.BooleanLlmResponseParser;
import it.cnr.isti.workflow.manager.executions.steps.Input;
import it.cnr.isti.workflow.manager.executions.bias.runtime.BiasRuntimeSupport;
import it.cnr.isti.workflow.manager.llms.LLMDescriptor;
import it.cnr.isti.workflow.manager.llms.providers.LLMProvider;
@ -128,7 +129,7 @@ public class ConditionalExecutor implements BlockExecutor<ConditionalBlockType>
builder.append(LLM_SYSTEM_PROMPT).append("\n");
builder.append("Decision prompt: ").append(resolvePlaceholders(config.getPrompt(), inputValues, executionVariables)).append("\n");
builder.append("Inputs: ").append(inputValues).append("\n");
return builder.toString();
return BiasRuntimeSupport.decoratePrompt(builder.toString(), executionVariables);
}
private String normalizeExpression(String expression) {

View File

@ -18,6 +18,7 @@ import it.cnr.isti.workflow.manager.executions.ExecutionEventType;
import it.cnr.isti.workflow.manager.executions.ExecutionVariableDescriptor;
import it.cnr.isti.workflow.manager.executions.ExecutionTemplateResolver;
import it.cnr.isti.workflow.manager.executions.steps.Input;
import it.cnr.isti.workflow.manager.executions.bias.runtime.BiasRuntimeSupport;
import it.cnr.isti.workflow.manager.llms.LLMDescriptor;
import it.cnr.isti.workflow.manager.llms.providers.LLMProvider;
import it.cnr.isti.workflow.manager.skills.SkillPromptService;
@ -54,6 +55,7 @@ public class LLMExecutor implements BlockExecutor<LLMBlockType> {
if (StringUtils.hasText(skillsInstructions)) {
prompt = skillsInstructions + System.lineSeparator() + System.lineSeparator() + prompt;
}
prompt = BiasRuntimeSupport.decoratePrompt(prompt, executionVariables);
LLMDescriptor llmDescriptor = config.getLlmDescriptor();
LLMProvider llmProvider = llmProviders.get(llmDescriptor.provider());

View File

@ -22,6 +22,7 @@ import it.cnr.isti.workflow.manager.executions.ExecutionVariableDescriptor;
import it.cnr.isti.workflow.manager.executions.ExecutionTemplateResolver;
import it.cnr.isti.workflow.manager.executions.InteractionResult;
import it.cnr.isti.workflow.manager.executions.steps.Input;
import it.cnr.isti.workflow.manager.executions.bias.runtime.BiasRuntimeSupport;
import it.cnr.isti.workflow.manager.llms.LLMDescriptor;
import it.cnr.isti.workflow.manager.llms.providers.LLMProvider;
import it.cnr.isti.workflow.manager.mcp.MCPAgentService;
@ -67,6 +68,7 @@ public class MCPAgentChatExecutor implements BlockExecutor<MCPAgentChatBlockType
LLMProvider simulatorProvider = resolveProvider(simulatorDescriptor.provider());
String simulatorAuthorization = resolveAuthorization(simulatorProvider, simulatorDescriptor, authorizations);
String resolvedGoal = resolvePlaceholders(configuration.getGoalDescription(), inputs, executionVariables);
resolvedGoal = BiasRuntimeSupport.decoratePrompt(resolvedGoal, executionVariables);
List<String> history = new ArrayList<>();
String sharedSessionKey = normalize(Boolean.TRUE.equals(configuration.getUseSharedSession())
? configuration.getSharedSessionRef()
@ -128,6 +130,7 @@ public class MCPAgentChatExecutor implements BlockExecutor<MCPAgentChatBlockType
}
String resolvedMessage = resolvePlaceholders(message, inputs, executionVariables);
resolvedMessage = BiasRuntimeSupport.decoratePrompt(resolvedMessage, executionVariables);
String model = Boolean.TRUE.equals(configuration.getUseSharedSession())
? null
: resolveConfigurableInput(MCPAgentChatBlockFactory.MODEL_INPUT, configuration.getModel(), inputs);

View File

@ -16,6 +16,7 @@ import it.cnr.isti.workflow.manager.executions.ExecutionEventType;
import it.cnr.isti.workflow.manager.executions.ExecutionVariableDescriptor;
import it.cnr.isti.workflow.manager.executions.ExecutionTemplateResolver;
import it.cnr.isti.workflow.manager.executions.steps.Input;
import it.cnr.isti.workflow.manager.executions.bias.runtime.BiasRuntimeSupport;
import it.cnr.isti.workflow.manager.mcp.MCPAgentService;
import it.cnr.isti.workflow.manager.mcp.MCPSharedSessionRegistry;
@ -31,6 +32,7 @@ public class MCPBridgeExecutor implements BlockExecutor<MCPAgentBlockType> {
Map<String, ExecutionVariableDescriptor> executionVariableDescriptors, ExecutionEventLogger eventLogger) {
MCPAgentBlockConfiguration config = (MCPAgentBlockConfiguration) block.getSpecificConfiguration();
String prompt = resolvePrompt(config.getPrompt(), inputs, executionVariables);
prompt = BiasRuntimeSupport.decoratePrompt(prompt, executionVariables);
String model = Boolean.TRUE.equals(config.getUseSharedSession())
? null
: resolveConfigurableInput("model", config.getModel(), inputs);

View File

@ -27,6 +27,7 @@ import it.cnr.isti.workflow.manager.executions.ExecutionRuntimeContextSupport;
import it.cnr.isti.workflow.manager.executions.ExecutionTemplateResolver;
import it.cnr.isti.workflow.manager.executions.ExecutionVariableDescriptor;
import it.cnr.isti.workflow.manager.executions.steps.Input;
import it.cnr.isti.workflow.manager.executions.bias.runtime.BiasRuntimeSupport;
import it.cnr.isti.workflow.manager.ios.IODescriptor;
import it.cnr.isti.workflow.manager.ios.IOType;
import it.cnr.isti.workflow.manager.llms.LLMDescriptor;
@ -149,7 +150,7 @@ public class SwitchExecutor implements BlockExecutor<SwitchBlockType> {
builder.append("Allowed outputs: ").append(allowedOutputs).append("\n");
builder.append("Decision prompt: ").append(resolvePlaceholders(config.getPrompt(), inputValues, executionVariables)).append("\n");
builder.append("Inputs: ").append(inputValues).append("\n");
return builder.toString();
return BiasRuntimeSupport.decoratePrompt(builder.toString(), executionVariables);
}
private String normalizeExpression(String expression) {

View File

@ -6,6 +6,7 @@ import java.util.Map;
import it.cnr.isti.workflow.manager.executions.ExecutionEvent;
import it.cnr.isti.workflow.manager.executions.ExecutionVariableDescriptor;
import it.cnr.isti.workflow.manager.executions.ExecutionStatus;
import it.cnr.isti.workflow.manager.executions.bias.BiasExecutionContext;
import it.cnr.isti.workflow.manager.executions.FieldKey;
import it.cnr.isti.workflow.manager.llms.LLMDescriptor;
import lombok.AllArgsConstructor;
@ -24,6 +25,7 @@ public class ExecutionSnapshot {
private Long endTime;
private boolean interactionSimulationEnabled;
private LLMDescriptor interactionSimulationDescriptor;
private BiasExecutionContext biasExecutionContext;
private Map<String, Object> executionVariables;
private Map<String, ExecutionVariableDescriptor> executionVariableDescriptors;
private Map<String, Object> globalInputs;

View File

@ -28,6 +28,12 @@ public class Input {
this.descriptor = descriptor;
}
public static Input detached(IODescriptor descriptor, Object value) {
Input input = new Input(descriptor);
input.setValue(value);
return input;
}
public void setValue(Object value) {
validateValue(value);
this.value = value;

View File

@ -16,6 +16,7 @@ import it.cnr.isti.workflow.manager.executions.ExecutionListener;
import it.cnr.isti.workflow.manager.executions.persistence.ExecutionInputSnapshot;
import it.cnr.isti.workflow.manager.executions.persistence.ExecutionStepSnapshot;
import it.cnr.isti.workflow.manager.executions.executors.NodeExecutors;
import it.cnr.isti.workflow.manager.executions.bias.BiasExecutionContext;
import it.cnr.isti.workflow.manager.flows.model.FlowNode;
import it.cnr.isti.workflow.manager.llms.LLMDescriptor;
import lombok.Builder;
@ -100,6 +101,11 @@ public class Step<N extends FlowNode> implements InputListener {
@JsonIgnore
private LLMDescriptor interactionSimulationDescriptor;
@Setter
@Getter
@JsonIgnore
private BiasExecutionContext biasExecutionContext = BiasExecutionContext.normal();
@Setter
@JsonIgnore
private ExecutionEventLogger eventLogger;
@ -170,9 +176,10 @@ public class Step<N extends FlowNode> implements InputListener {
try {
Map<String, Object> outputs = isSimulated() && this.node.isUserInteractive()
? NodeExecutors.simulate(this.node, this.inputs, authorizations, executionVariables,
executionVariableDescriptors, this.interactionSimulationDescriptor, this.eventLogger)
executionVariableDescriptors, this.interactionSimulationDescriptor, this.eventLogger,
this.biasExecutionContext)
: NodeExecutors.execute(this.node, this.inputs, authorizations, executionVariables,
executionVariableDescriptors, this.eventLogger);
executionVariableDescriptors, this.eventLogger, this.biasExecutionContext);
for (Output output : this.outputs) {
if (outputs.containsKey(output.getDescriptor().getName())) {
output.setValue(outputs.get(output.getDescriptor().getName()));
@ -203,7 +210,7 @@ public class Step<N extends FlowNode> implements InputListener {
listener.resumed(this.id);
var interactionResult = NodeExecutors.interact(this.node, this.inputs, providedOutputs,
Map.copyOf(this.partialResults), authorizations, executionVariables, executionVariableDescriptors,
this.eventLogger);
this.eventLogger, this.biasExecutionContext);
this.partialResults.clear();
this.partialResults.putAll(interactionResult.partialResults());
listener.partialUpdated(this.id, this.partialResults);

View File

@ -0,0 +1,20 @@
package it.cnr.isti.workflow.manager.flows.model.bias;
import lombok.Getter;
@Getter
public enum BiasActivationMode implements BiasCatalogOption {
PROMPT_DIRECTIVE("Prompt directive", "Adds an experimental behavioural directive to prompt-based nodes."),
INPUT_TRANSFORMATION("Input transformation", "Transforms selected textual inputs before the node is executed."),
OUTPUT_TRANSFORMATION("Output transformation", "Transforms textual outputs after the node is executed."),
ROUTING_OVERRIDE("Routing override", "Forces a Conditional or Switch block to a named output."),
MOCK_RESPONSE("Mock response", "Skips an external call and returns the configured experimental response.");
private final String label;
private final String description;
BiasActivationMode(String label, String description) {
this.label = label;
this.description = description;
}
}

View File

@ -0,0 +1,53 @@
package it.cnr.isti.workflow.manager.flows.model.bias;
import java.util.List;
import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
import com.fasterxml.jackson.annotation.JsonProperty;
import it.cnr.isti.workflow.manager.configurations.annotations.LongText;
import it.cnr.isti.workflow.manager.configurations.annotations.UiDescription;
import it.cnr.isti.workflow.manager.configurations.annotations.UiLabel;
import it.cnr.isti.workflow.manager.configurations.annotations.UiOrder;
import jakarta.validation.constraints.NotBlank;
import jakarta.validation.constraints.NotNull;
import jakarta.validation.constraints.Size;
@JsonIgnoreProperties(ignoreUnknown = true)
public record BiasBehavioralProbe(
@NotNull
@JsonProperty(required = true)
@UiLabel("Activation mode")
@UiDescription("Controlled mechanism used to reproduce the annotated behaviour.")
@UiOrder(10)
BiasActivationMode activationMode,
@NotBlank
@Size(max = 4000)
@JsonProperty(required = true)
@UiLabel("Experimental instruction")
@UiDescription("Instruction, replacement template, branch name or mock response used by the selected activation mode. Use ${original} in transformation templates.")
@UiOrder(20)
@LongText(placeholder = "Describe exactly how the biased behaviour should be reproduced")
String instruction,
@Size(max = 50)
@UiLabel("Target inputs")
@UiDescription("Optional input names affected by INPUT_TRANSFORMATION. An empty list targets every textual input.")
@UiOrder(30)
List<String> targetInputs,
@Size(max = 2000)
@UiLabel("Expected impact")
@UiDescription("Expected observable difference, used to explain the experiment but not to alter its result.")
@UiOrder(40)
@LongText(placeholder = "Describe the expected observable impact")
String expectedImpact) {
public BiasBehavioralProbe {
targetInputs = targetInputs == null ? List.of() : targetInputs.stream()
.filter(value -> value != null && !value.isBlank())
.distinct()
.toList();
}
}

View File

@ -74,7 +74,12 @@ public record BlockBiasAnnotation(
@UiLabel("Analysis identifier")
@UiDescription("Optional identifier of the automated analysis that produced the annotation.")
@UiOrder(90)
String analysisId) {
String analysisId,
@UiLabel("Behavioural probe")
@UiDescription("Optional executable definition used only by explicit bias impact experiments.")
@UiOrder(100)
BiasBehavioralProbe behavioralProbe) {
public static final int MAX_PER_BLOCK = 50;

View File

@ -30,6 +30,8 @@ import it.cnr.isti.workflow.manager.flows.model.Dependency;
import it.cnr.isti.workflow.manager.flows.model.FlowData;
import it.cnr.isti.workflow.manager.flows.model.FlowNode;
import it.cnr.isti.workflow.manager.flows.model.bias.BlockBiasAnnotation;
import it.cnr.isti.workflow.manager.flows.model.bias.BiasActivationMode;
import it.cnr.isti.workflow.manager.executions.bias.runtime.BiasBehaviorAdapterRegistry;
import it.cnr.isti.workflow.manager.ios.IODescriptor;
import jakarta.validation.ConstraintValidator;
import jakarta.validation.ConstraintValidatorContext;
@ -46,6 +48,9 @@ public class FlowDataValidator implements ConstraintValidator<ValidFlowStructure
@Autowired
ContainerSubFlowValidationRegistry containerSubFlowValidationRegistry;
@Autowired
BiasBehaviorAdapterRegistry biasBehaviorAdapterRegistry;
@Override
public boolean isValid(FlowData flowData, ConstraintValidatorContext context) {
if (flowData == null) {
@ -180,6 +185,37 @@ public class FlowDataValidator implements ConstraintValidator<ValidFlowStructure
validateBiasFieldLength(block, field + ".rationale", annotation.rationale(), 4000);
validateBiasFieldLength(block, field + ".mitigation", annotation.mitigation(), 4000);
validateBiasFieldLength(block, field + ".analysisId", annotation.analysisId(), 255);
validateBehavioralProbe(block, annotation, field + ".behavioralProbe");
}
}
private void validateBehavioralProbe(Block<?> block, BlockBiasAnnotation annotation, String field) {
if (annotation.behavioralProbe() == null) {
return;
}
if (annotation.behavioralProbe().activationMode() == null) {
throw validationError(error(ValidationErrorCode.BIAS_PROBE_MODE_REQUIRED, "block", block.getId(), field + ".activationMode",
"Bias behavioral probe activation mode is required"));
}
if (annotation.behavioralProbe().instruction() == null || annotation.behavioralProbe().instruction().isBlank()) {
throw validationError(error(ValidationErrorCode.BIAS_PROBE_INSTRUCTION_REQUIRED, "block", block.getId(), field + ".instruction",
"Bias behavioral probe instruction is required"));
}
validateBiasFieldLength(block, field + ".instruction", annotation.behavioralProbe().instruction(), 4000);
validateBiasFieldLength(block, field + ".expectedImpact", annotation.behavioralProbe().expectedImpact(), 2000);
if (!biasBehaviorAdapterRegistry.supportedModes(block).contains(annotation.behavioralProbe().activationMode())) {
throw validationError(error(ValidationErrorCode.BIAS_PROBE_MODE_UNSUPPORTED, "block", block.getId(), field + ".activationMode",
"Activation mode " + annotation.behavioralProbe().activationMode() + " is not supported by this block"));
}
if (annotation.behavioralProbe().activationMode() == BiasActivationMode.INPUT_TRANSFORMATION) {
Set<String> inputNames = block.getInputs().stream().map(IODescriptor::getName).collect(java.util.stream.Collectors.toSet());
annotation.behavioralProbe().targetInputs().stream()
.filter(target -> !inputNames.contains(target))
.findFirst()
.ifPresent(target -> {
throw validationError(error(ValidationErrorCode.BIAS_PROBE_TARGET_INPUT_NOT_FOUND, "block", block.getId(),
field + ".targetInputs", "Bias probe target input not found: " + target));
});
}
}

View File

@ -16,6 +16,10 @@ public enum ValidationErrorCode {
BIAS_SEVERITY_REQUIRED,
BIAS_ISSUE_REQUIRED,
BIAS_FIELD_TOO_LONG,
BIAS_PROBE_MODE_REQUIRED,
BIAS_PROBE_INSTRUCTION_REQUIRED,
BIAS_PROBE_MODE_UNSUPPORTED,
BIAS_PROBE_TARGET_INPUT_NOT_FOUND,
CONTAINER_CONFIGURATION_MISSING,
CONTAINER_SUBFLOW_INVALID,
CONTAINER_SUBFLOW_EMPTY,

View File

@ -7,6 +7,8 @@ spring.datasource.password=${DB_PASSWORD:password}
spring.datasource.driver-class-name=org.postgresql.Driver
spring.jpa.properties.hibernate.dialect=org.hibernate.dialect.PostgreSQLDialect
spring.jpa.hibernate.ddl-auto=${ddl-auto:validate}
spring.flyway.baseline-on-migrate=${FLYWAY_BASELINE_ON_MIGRATE:true}
spring.flyway.baseline-version=${FLYWAY_BASELINE_VERSION:1}
management.endpoints.web.exposure.include=health,info
management.endpoint.health.show-details=always

View File

@ -0,0 +1,14 @@
CREATE TABLE bias_impact_report_entity (
id VARCHAR(255) PRIMARY KEY,
owner VARCHAR(255) NOT NULL,
baseline_execution_id VARCHAR(255) NOT NULL,
biased_execution_id VARCHAR(255),
created_at TIMESTAMP(6) NOT NULL,
report_data TEXT
);
CREATE INDEX idx_bias_report_baseline_owner_created
ON bias_impact_report_entity (baseline_execution_id, owner, created_at);
CREATE INDEX idx_bias_report_biased_execution
ON bias_impact_report_entity (biased_execution_id);

View File

@ -82,6 +82,7 @@ class BiasAnnotationsControllerTest {
"Add a human review step.",
null,
null,
null,
null);
assertNotNull(annotation.id());
@ -134,6 +135,7 @@ class BiasAnnotationsControllerTest {
"Ask the model to provide supporting reasons.",
BiasAnnotationStatus.CONFIRMED,
BiasAnnotationSource.MANUAL,
null,
null);
Block<LLMBlockType> block = annotatedBlock(annotation);
FlowCreateRequest request = new FlowCreateRequest(
@ -160,6 +162,7 @@ class BiasAnnotationsControllerTest {
null,
null,
null,
null,
null);
FlowCreateRequest request = new FlowCreateRequest(
"Invalid bias test",

View File

@ -0,0 +1,271 @@
package it.cnr.isti.workflow.manager.executions.bias;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
import java.util.List;
import java.util.Map;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.boot.test.context.TestConfiguration;
import org.springframework.context.annotation.Bean;
import org.springframework.test.context.TestPropertySource;
import it.cnr.isti.workflow.manager.blocks.Block;
import it.cnr.isti.workflow.manager.blocks.configurations.HTTPServerCallBlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.LLMBlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.factories.HTTPServerCallBlockFactory;
import it.cnr.isti.workflow.manager.blocks.factories.LLMBlockFactory;
import it.cnr.isti.workflow.manager.blocks.types.HTTPServerCallBlockType;
import it.cnr.isti.workflow.manager.blocks.types.LLMBlockType;
import it.cnr.isti.workflow.manager.controllers.BiasAnnotationsController;
import it.cnr.isti.workflow.manager.controllers.BiasCapabilitiesController;
import it.cnr.isti.workflow.manager.executions.ExecutionObject;
import it.cnr.isti.workflow.manager.executions.ExecutionStatus;
import it.cnr.isti.workflow.manager.executions.ExecutionsService;
import it.cnr.isti.workflow.manager.executions.executors.NodeExecutors;
import it.cnr.isti.workflow.manager.flows.model.FlowData;
import it.cnr.isti.workflow.manager.flows.model.bias.BiasActivationMode;
import it.cnr.isti.workflow.manager.flows.model.bias.BiasAnnotationSource;
import it.cnr.isti.workflow.manager.flows.model.bias.BiasAnnotationStatus;
import it.cnr.isti.workflow.manager.flows.model.bias.BiasBehavioralProbe;
import it.cnr.isti.workflow.manager.flows.model.bias.BiasCategory;
import it.cnr.isti.workflow.manager.flows.model.bias.BiasSeverity;
import it.cnr.isti.workflow.manager.flows.model.bias.BlockBiasAnnotation;
import it.cnr.isti.workflow.manager.llms.LLMDescriptor;
import it.cnr.isti.workflow.manager.llms.providers.LLMProvider;
@SpringBootTest
@TestPropertySource(locations = "classpath:test.properties")
class BiasExperimentsIntegrationTest {
private static final String OWNER = "bias-experiment-user";
@TestConfiguration
static class TestConfig {
@Bean
LLMProvider biasExperimentProvider() {
return new LLMProvider() {
@Override
public String getName() {
return "biasExperimentProvider";
}
@Override
public List<String> getRegisteredModels() {
return List.of("bias-test-model");
}
@Override
public String generate(String model, String prompt) {
return prompt.contains("BIAS IMPACT EXPERIMENT") ? "BIASED" : "BASELINE";
}
};
}
}
@Autowired
LLMBlockFactory llmBlockFactory;
@Autowired
LLMBlockType llmBlockType;
@Autowired
HTTPServerCallBlockFactory httpBlockFactory;
@Autowired
HTTPServerCallBlockType httpBlockType;
@Autowired
ExecutionsService executionsService;
@Autowired
BiasImpactService biasImpactService;
@Autowired
BiasAnnotationsController annotationsController;
@Autowired
BiasCapabilitiesController capabilitiesController;
@Test
void discoveryExposesBehavioralProbeAndPerTypeCapabilities() {
var descriptor = annotationsController.getDescriptor();
assertTrue(descriptor.options().containsKey("behavioralProbe.activationMode"));
assertTrue(descriptor.schema().toString().contains("behavioralProbe"));
assertTrue(descriptor.schema().toString().contains("PROMPT_DIRECTIVE"));
var llmCapabilities = capabilitiesController.getCapabilities(LLMBlockType.TYPE);
assertTrue(llmCapabilities.supported());
assertTrue(llmCapabilities.isolatedExperimentSupported());
assertTrue(llmCapabilities.activationModes().contains(BiasActivationMode.PROMPT_DIRECTIVE));
assertFalse(llmCapabilities.externalSideEffects());
var httpCapabilities = capabilitiesController.getCapabilities(HTTPServerCallBlockType.TYPE);
assertTrue(httpCapabilities.externalSideEffects());
assertTrue(httpCapabilities.activationModes().contains(BiasActivationMode.MOCK_RESPONSE));
}
@Test
void isolatedExperimentRepeatsBiasedVariantAndPersistsReport() {
Block<LLMBlockType> block = annotatedLlmBlock();
ExecutionObject baseline = completedExecution(block);
BiasImpactReport report = biasImpactService.runIsolatedStepExperiment(
baseline.getId(),
block.getId(),
new BiasImpactExperimentRequest(
List.of(block.getBiasAnnotations().getFirst().id()),
3,
true,
ExternalSideEffectPolicy.BLOCK,
false),
OWNER);
assertEquals(BiasExperimentKind.ISOLATED_STEP, report.kind());
assertTrue(report.immediateImpact().outputChanged());
assertEquals(1.0, report.immediateImpact().changeRate());
assertEquals("BASELINE", report.immediateImpact().baselineOutput().get(LLMBlockFactory.OUTPUT_NAME));
assertTrue(report.immediateImpact().biasedOutputs().stream()
.allMatch(output -> "BIASED".equals(output.get(LLMBlockFactory.OUTPUT_NAME))));
assertEquals(report, biasImpactService.getReport(report.id(), OWNER));
assertTrue(biasImpactService.getReports(baseline.getId(), OWNER).stream()
.anyMatch(candidate -> candidate.id().equals(report.id())));
}
@Test
void biasedRerunPersistsContextAndProducesFullFlowComparison() {
Block<LLMBlockType> block = annotatedLlmBlock();
ExecutionObject baseline = completedExecution(block);
String annotationId = block.getBiasAnnotations().getFirst().id();
ExecutionObject biased = executionsService.createBiasRerun(
baseline.getId(),
OWNER,
new BiasRerunRequest(
List.of(new BiasActivation(block.getId(), List.of(annotationId))),
ExternalSideEffectPolicy.BLOCK,
false));
assertEquals(BiasExecutionMode.BIAS_VARIANT, biased.getBiasExecutionContext().mode());
assertEquals(baseline.getId(), biased.getRerunOfExecutionId());
executionsService.startExecution(biased.getId());
waitUntilFinal(biased);
BiasImpactReport report = biasImpactService.compareFullFlow(
baseline.getId(), biased.getId(), true, OWNER);
assertEquals(BiasExperimentKind.FULL_FLOW, report.kind());
assertEquals(biased.getId(), report.biasedExecutionId());
assertTrue(report.immediateImpact().outputChanged());
assertEquals("BIASED", report.immediateImpact().biasedOutputs().getFirst()
.get(block.getId() + "." + LLMBlockFactory.OUTPUT_NAME));
assertTrue(biasImpactService.getReports(biased.getId(), OWNER).stream()
.anyMatch(candidate -> candidate.id().equals(report.id())));
executionsService.clearInMemoryExecutions();
ExecutionObject restored = executionsService.getExecutionByOwner(biased.getId(), OWNER);
assertEquals(BiasExecutionMode.BIAS_VARIANT, restored.getBiasExecutionContext().mode());
assertEquals(annotationId, restored.getBiasExecutionContext().annotationIdsFor(block.getId()).getFirst());
}
@Test
void mockPolicyPreventsHttpSideEffectAndReturnsProbeResponse() {
BlockBiasAnnotation annotation = annotation(
BiasActivationMode.MOCK_RESPONSE,
"mocked-biased-response",
List.of());
HTTPServerCallBlockConfiguration configuration = HTTPServerCallBlockConfiguration.builder()
.name("External API")
.url("https://example.invalid/side-effect")
.method(HTTPServerCallBlockConfiguration.HTTPMethod.POST)
.requestBody("{}")
.requiresAuthorization(false)
.build();
Block<HTTPServerCallBlockType> canonical = httpBlockFactory.create(configuration);
Block<HTTPServerCallBlockType> block = Block.<HTTPServerCallBlockType>builder()
.specificConfiguration(configuration)
.inputs(canonical.getInputs())
.outputs(canonical.getOutputs())
.type(httpBlockType)
.biasAnnotation(annotation)
.build();
BiasExecutionContext context = new BiasExecutionContext(
"mock-experiment",
BiasExecutionMode.BIAS_VARIANT,
Map.of(block.getId(), List.of(annotation.id())),
ExternalSideEffectPolicy.MOCK,
false);
Map<String, Object> result = NodeExecutors.execute(
block, List.of(), Map.of(), Map.of(), Map.of(), null, context);
assertEquals("mocked-biased-response", result.get(HTTPServerCallBlockFactory.OUTPUT_NAME));
}
private ExecutionObject completedExecution(Block<LLMBlockType> block) {
ExecutionObject execution = executionsService.createExecutionForFlow(
"bias-flow-" + block.getId(),
"Bias experiment flow",
FlowData.builder().block(block).build(),
OWNER);
executionsService.startExecution(execution.getId());
waitUntilFinal(execution);
assertEquals(ExecutionStatus.SUCCESS, execution.getContext().getStatus());
return execution;
}
private void waitUntilFinal(ExecutionObject execution) {
long deadline = System.currentTimeMillis() + 5000;
while (!execution.getContext().getStatus().isFinalState() && System.currentTimeMillis() < deadline) {
try {
Thread.sleep(10);
} catch (InterruptedException exception) {
Thread.currentThread().interrupt();
throw new IllegalStateException(exception);
}
}
assertTrue(execution.getContext().getStatus().isFinalState());
}
private Block<LLMBlockType> annotatedLlmBlock() {
BlockBiasAnnotation annotation = annotation(
BiasActivationMode.PROMPT_DIRECTIVE,
"Prefer the initially suggested conclusion.",
List.of());
LLMBlockConfiguration configuration = LLMBlockConfiguration.builder()
.name("Biased decision")
.prompt("Return a decision")
.llmDescriptor(LLMDescriptor.builder()
.provider("biasExperimentProvider")
.model("bias-test-model")
.build())
.build();
Block<LLMBlockType> canonical = llmBlockFactory.create(configuration);
return Block.<LLMBlockType>builder()
.specificConfiguration(configuration)
.inputs(canonical.getInputs())
.outputs(canonical.getOutputs())
.type(llmBlockType)
.biasAnnotation(annotation)
.build();
}
private BlockBiasAnnotation annotation(BiasActivationMode mode, String instruction, List<String> targetInputs) {
return new BlockBiasAnnotation(
null,
BiasCategory.CONFIRMATION_BIAS,
BiasSeverity.HIGH,
"The node may reinforce an initial conclusion.",
"The output can underweight contrary evidence.",
"Require independent evidence.",
BiasAnnotationStatus.CONFIRMED,
BiasAnnotationSource.MANUAL,
null,
new BiasBehavioralProbe(mode, instruction, targetInputs, "The output should change."));
}
}