diff --git a/pom.xml b/pom.xml
index 5a2afec..e94bb4c 100644
--- a/pom.xml
+++ b/pom.xml
@@ -107,6 +107,10 @@
org.flywaydb
flyway-core
+
+ org.flywaydb
+ flyway-database-postgresql
+
org.springframework.boot
spring-boot-starter-web
diff --git a/src/main/java/it/cnr/isti/workflow/manager/controllers/BiasAnnotationsController.java b/src/main/java/it/cnr/isti/workflow/manager/controllers/BiasAnnotationsController.java
index 6a890eb..f84b5fc 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/controllers/BiasAnnotationsController.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/controllers/BiasAnnotationsController.java
@@ -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(),
diff --git a/src/main/java/it/cnr/isti/workflow/manager/controllers/BiasCapabilitiesController.java b/src/main/java/it/cnr/isti/workflow/manager/controllers/BiasCapabilitiesController.java
new file mode 100644
index 0000000..e107921
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/controllers/BiasCapabilitiesController.java
@@ -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 blockTypes;
+
+ @Autowired
+ List> blockFactories;
+
+ @Autowired
+ BiasBehaviorAdapterRegistry adapterRegistry;
+
+ public record BiasCapabilityDescriptor(
+ String blockType,
+ boolean supported,
+ boolean isolatedExperimentSupported,
+ boolean fullFlowExperimentSupported,
+ boolean externalSideEffects,
+ List 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 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();
+ }
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/controllers/BiasExperimentsController.java b/src/main/java/it/cnr/isti/workflow/manager/controllers/BiasExperimentsController.java
new file mode 100644
index 0000000..224641d
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/controllers/BiasExperimentsController.java
@@ -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 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();
+ }
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java
index 768dcd6..f8a6c6b 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionContext.java
@@ -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> 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 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());
}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionEventType.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionEventType.java
index 59d7015..ea8b4dd 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionEventType.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionEventType.java
@@ -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
}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionObject.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionObject.java
index 2211a04..799752e 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionObject.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionObject.java
@@ -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 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> 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 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();
}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java
index a789fcd..ee9b817 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java
@@ -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 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> 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> validateBiasActivations(FlowData flow, List activations) {
+ Map nodesById = flow.getNodes().stream()
+ .collect(Collectors.toMap(it.cnr.isti.workflow.manager.flows.model.FlowNode::getId, node -> node));
+ Map> 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 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);
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/api/ExecutionView.java b/src/main/java/it/cnr/isti/workflow/manager/executions/api/ExecutionView.java
index b380488..20c909d 100644
--- a/src/main/java/it/cnr/isti/workflow/manager/executions/api/ExecutionView.java
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/api/ExecutionView.java
@@ -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 missingAuthorizationKeys;
private List 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();
}
}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasActivation.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasActivation.java
new file mode 100644
index 0000000..38b01ae
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasActivation.java
@@ -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();
+ }
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasDownstreamImpact.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasDownstreamImpact.java
new file mode 100644
index 0000000..48e8190
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasDownstreamImpact.java
@@ -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 baselineOutputs,
+ Map biasedOutputs) {
+
+ public BiasDownstreamImpact {
+ baselineOutputs = baselineOutputs == null ? Map.of() : Map.copyOf(baselineOutputs);
+ biasedOutputs = biasedOutputs == null ? Map.of() : Map.copyOf(biasedOutputs);
+ }
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasExecutionContext.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasExecutionContext.java
new file mode 100644
index 0000000..bc215dc
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasExecutionContext.java
@@ -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> activeAnnotationIdsByNode,
+ ExternalSideEffectPolicy externalSideEffectPolicy,
+ boolean externalSideEffectsConfirmed) {
+
+ public BiasExecutionContext {
+ mode = mode == null ? BiasExecutionMode.NORMAL : mode;
+ externalSideEffectPolicy = externalSideEffectPolicy == null ? ExternalSideEffectPolicy.BLOCK : externalSideEffectPolicy;
+ Map> 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 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;
+ }
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasExecutionMode.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasExecutionMode.java
new file mode 100644
index 0000000..473e4ca
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasExecutionMode.java
@@ -0,0 +1,7 @@
+package it.cnr.isti.workflow.manager.executions.bias;
+
+public enum BiasExecutionMode {
+ NORMAL,
+ BIAS_BASELINE,
+ BIAS_VARIANT
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasExperimentKind.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasExperimentKind.java
new file mode 100644
index 0000000..180b3be
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasExperimentKind.java
@@ -0,0 +1,6 @@
+package it.cnr.isti.workflow.manager.executions.bias;
+
+public enum BiasExperimentKind {
+ ISOLATED_STEP,
+ FULL_FLOW
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasImpactExperimentRequest.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasImpactExperimentRequest.java
new file mode 100644
index 0000000..ad25895
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasImpactExperimentRequest.java
@@ -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 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);
+ }
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasImpactReport.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasImpactReport.java
new file mode 100644
index 0000000..5e813c8
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasImpactReport.java
@@ -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 annotationIds,
+ int repetitions,
+ LocalDateTime createdAt,
+ BiasOutputImpact immediateImpact,
+ List downstreamImpact,
+ List routingChanges,
+ String summary,
+ List 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);
+ }
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasImpactService.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasImpactService.java
new file mode 100644
index 0000000..a0f38ac
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasImpactService.java
@@ -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 baselineOutput = outputsOf(baseline, stepId);
+ List