From 587eadc6781021fa26afee17ea3f4e0d53117c78 Mon Sep 17 00:00:00 2001 From: Lucio Lelii Date: Mon, 20 Jul 2026 17:22:26 +0200 Subject: [PATCH] Implement bias impact experiments --- pom.xml | 4 + .../BiasAnnotationsController.java | 2 + .../BiasCapabilitiesController.java | 71 ++++ .../BiasExperimentsController.java | 95 +++++ .../manager/executions/ExecutionContext.java | 16 + .../executions/ExecutionEventType.java | 4 +- .../manager/executions/ExecutionObject.java | 17 +- .../manager/executions/ExecutionsService.java | 76 +++- .../manager/executions/api/ExecutionView.java | 3 + .../executions/bias/BiasActivation.java | 15 + .../executions/bias/BiasDownstreamImpact.java | 18 + .../executions/bias/BiasExecutionContext.java | 46 +++ .../executions/bias/BiasExecutionMode.java | 7 + .../executions/bias/BiasExperimentKind.java | 6 + .../bias/BiasImpactExperimentRequest.java | 23 ++ .../executions/bias/BiasImpactReport.java | 28 ++ .../executions/bias/BiasImpactService.java | 336 ++++++++++++++++++ .../executions/bias/BiasOutputImpact.java | 17 + .../executions/bias/BiasRerunRequest.java | 18 + .../executions/bias/BiasRoutingChange.java | 7 + .../bias/ExternalSideEffectPolicy.java | 7 + .../BiasImpactReportConverter.java | 47 +++ .../persistence/BiasImpactReportEntity.java | 51 +++ .../BiasImpactReportRepository.java | 25 ++ .../bias/runtime/BiasBehaviorAdapter.java | 16 + .../runtime/BiasBehaviorAdapterRegistry.java | 171 +++++++++ .../bias/runtime/BiasPreparedExecution.java | 69 ++++ .../bias/runtime/BiasRuntimeSupport.java | 71 ++++ .../runtime/InputBiasBehaviorAdapter.java | 42 +++ .../MockResponseBiasBehaviorAdapter.java | 38 ++ .../runtime/OutputBiasBehaviorAdapter.java | 28 ++ .../runtime/PromptBiasBehaviorAdapter.java | 48 +++ .../runtime/RoutingBiasBehaviorAdapter.java | 35 ++ .../executions/executors/NodeExecutors.java | 67 +++- .../blocks/ChatInteractionExecutor.java | 6 + .../executors/blocks/ConditionalExecutor.java | 3 +- .../executors/blocks/LLMExecutor.java | 2 + .../blocks/MCPAgentChatExecutor.java | 3 + .../executors/blocks/MCPBridgeExecutor.java | 2 + .../executors/blocks/SwitchExecutor.java | 3 +- .../persistence/ExecutionSnapshot.java | 2 + .../manager/executions/steps/Input.java | 6 + .../manager/executions/steps/Step.java | 13 +- .../flows/model/bias/BiasActivationMode.java | 20 ++ .../flows/model/bias/BiasBehavioralProbe.java | 53 +++ .../flows/model/bias/BlockBiasAnnotation.java | 7 +- .../flows/validation/FlowDataValidator.java | 36 ++ .../flows/validation/ValidationErrorCode.java | 4 + src/main/resources/application.properties | 2 + .../V2__create_bias_impact_reports.sql | 14 + .../BiasAnnotationsControllerTest.java | 3 + .../bias/BiasExperimentsIntegrationTest.java | 271 ++++++++++++++ 52 files changed, 1953 insertions(+), 21 deletions(-) create mode 100644 src/main/java/it/cnr/isti/workflow/manager/controllers/BiasCapabilitiesController.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/controllers/BiasExperimentsController.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasActivation.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasDownstreamImpact.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasExecutionContext.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasExecutionMode.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasExperimentKind.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasImpactExperimentRequest.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasImpactReport.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasImpactService.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasOutputImpact.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasRerunRequest.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasRoutingChange.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/bias/ExternalSideEffectPolicy.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/bias/persistence/BiasImpactReportConverter.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/bias/persistence/BiasImpactReportEntity.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/bias/persistence/BiasImpactReportRepository.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/BiasBehaviorAdapter.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/BiasBehaviorAdapterRegistry.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/BiasPreparedExecution.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/BiasRuntimeSupport.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/InputBiasBehaviorAdapter.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/MockResponseBiasBehaviorAdapter.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/OutputBiasBehaviorAdapter.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/PromptBiasBehaviorAdapter.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/RoutingBiasBehaviorAdapter.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/flows/model/bias/BiasActivationMode.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/flows/model/bias/BiasBehavioralProbe.java create mode 100644 src/main/resources/db/migration/V2__create_bias_impact_reports.sql create mode 100644 src/test/java/it/cnr/isti/workflow/manager/executions/bias/BiasExperimentsIntegrationTest.java 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> 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 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 activatedNodeIds = context.activeAnnotationIdsByNode().keySet(); + Set downstreamIds = downstreamNodeIds(baseline, activatedNodeIds); + List downstream = downstreamIds.stream() + .map(nodeId -> compareStep(baseline, biased, nodeId, includeRawOutputs)) + .toList(); + + Map baselineImmediate = combinedOutputs(baseline, activatedNodeIds); + Map biasedImmediate = combinedOutputs(biased, activatedNodeIds); + BiasOutputImpact immediate = outputImpact(baselineImmediate, List.of(biasedImmediate), includeRawOutputs); + List routing = activatedNodeIds.stream() + .flatMap(nodeId -> routingChanges(nodeId, outputsOf(baseline, nodeId), outputsOf(biased, nodeId)).stream()) + .toList(); + List 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 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 annotationIds) { + Map 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 baselineOutputs = outputsOf(baseline, nodeId); + Map 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 downstreamNodeIds(ExecutionObject execution, Set sourceIds) { + Map> 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 visited = new LinkedHashSet<>(); + Deque 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 combinedOutputs(ExecutionObject execution, Set nodeIds) { + Map combined = new LinkedHashMap<>(); + nodeIds.forEach(nodeId -> outputsOf(execution, nodeId) + .forEach((name, value) -> combined.put(nodeId + "." + name, value))); + return Collections.unmodifiableMap(combined); + } + + private Map outputsOf(ExecutionObject execution, String nodeId) { + Map 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 baseline, List> 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 routingChanges(String nodeId, Map baseline, + Map 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; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasOutputImpact.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasOutputImpact.java new file mode 100644 index 0000000..96e6cc2 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasOutputImpact.java @@ -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 baselineOutput, + List> biasedOutputs) { + + public BiasOutputImpact { + baselineOutput = baselineOutput == null ? Map.of() : Map.copyOf(baselineOutput); + biasedOutputs = biasedOutputs == null ? List.of() : List.copyOf(biasedOutputs); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasRerunRequest.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasRerunRequest.java new file mode 100644 index 0000000..348609a --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasRerunRequest.java @@ -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); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasRoutingChange.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasRoutingChange.java new file mode 100644 index 0000000..be0d15c --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasRoutingChange.java @@ -0,0 +1,7 @@ +package it.cnr.isti.workflow.manager.executions.bias; + +public record BiasRoutingChange( + String nodeId, + String baselineBranch, + String biasedBranch) { +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/ExternalSideEffectPolicy.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/ExternalSideEffectPolicy.java new file mode 100644 index 0000000..3cb13df --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/ExternalSideEffectPolicy.java @@ -0,0 +1,7 @@ +package it.cnr.isti.workflow.manager.executions.bias; + +public enum ExternalSideEffectPolicy { + BLOCK, + MOCK, + REQUIRE_CONFIRMATION +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/persistence/BiasImpactReportConverter.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/persistence/BiasImpactReportConverter.java new file mode 100644 index 0000000..0063957 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/persistence/BiasImpactReportConverter.java @@ -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 { + + 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; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/persistence/BiasImpactReportEntity.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/persistence/BiasImpactReportEntity.java new file mode 100644 index 0000000..9abfaa5 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/persistence/BiasImpactReportEntity.java @@ -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; +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/persistence/BiasImpactReportRepository.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/persistence/BiasImpactReportRepository.java new file mode 100644 index 0000000..0bbec43 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/persistence/BiasImpactReportRepository.java @@ -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 { + + Optional 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 findAllForExecution( + @Param("executionId") String executionId, + @Param("owner") String owner); + + void deleteByBaselineExecutionIdOrBiasedExecutionId(String baselineExecutionId, String biasedExecutionId); +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/BiasBehaviorAdapter.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/BiasBehaviorAdapter.java new file mode 100644 index 0000000..448945c --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/BiasBehaviorAdapter.java @@ -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 supportedModes(Block block); +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/BiasBehaviorAdapterRegistry.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/BiasBehaviorAdapterRegistry.java new file mode 100644 index 0000000..38e7be8 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/BiasBehaviorAdapterRegistry.java @@ -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 registeredAdapters = List.of(); + + public BiasBehaviorAdapterRegistry(List adapters) { + registeredAdapters = List.copyOf(adapters); + } + + public static BiasPreparedExecution prepare(Block block, List inputs, Map executionVariables, + BiasExecutionContext context, ExecutionEventLogger eventLogger) { + BiasExecutionContext effectiveContext = context == null ? BiasExecutionContext.normal() : context; + List 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 complete(BiasPreparedExecution prepared, Map result) { + Map 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 supportedModes(Block block) { + Set 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 resolveActiveAnnotations(Block block, BiasExecutionContext context) { + if (block == null || context.mode() != BiasExecutionMode.BIAS_VARIANT) { + return List.of(); + } + List activeIds = context.annotationIdsFor(block.getId()); + if (activeIds.isEmpty()) { + return List.of(); + } + Map byId = new LinkedHashMap<>(); + block.getBiasAnnotations().forEach(annotation -> { + if (annotation != null) { + byId.put(annotation.id(), annotation); + } + }); + List 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 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 adapters() { + if (registeredAdapters.isEmpty()) { + throw new IllegalStateException("Bias behavior adapters have not been initialized"); + } + return registeredAdapters; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/BiasPreparedExecution.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/BiasPreparedExecution.java new file mode 100644 index 0000000..598e735 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/BiasPreparedExecution.java @@ -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 inputs; + private final Map executionVariables; + private final List annotations; + private Map bypassResult; + private final List outputTemplates = new ArrayList<>(); + private String routingOverride; + + public BiasPreparedExecution(Block block, List inputs, Map executionVariables, + List 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 getInputs() { + return inputs; + } + + public void setInputs(List inputs) { + this.inputs = inputs == null ? List.of() : List.copyOf(inputs); + } + + public Map getExecutionVariables() { + return executionVariables; + } + + public List getAnnotations() { + return annotations; + } + + public Map getBypassResult() { + return bypassResult; + } + + public void setBypassResult(Map bypassResult) { + this.bypassResult = bypassResult == null ? null : Map.copyOf(bypassResult); + } + + public List getOutputTemplates() { + return outputTemplates; + } + + public String getRoutingOverride() { + return routingOverride; + } + + public void setRoutingOverride(String routingOverride) { + this.routingOverride = routingOverride; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/BiasRuntimeSupport.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/BiasRuntimeSupport.java new file mode 100644 index 0000000..f3a7565 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/BiasRuntimeSupport.java @@ -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 executionVariables, String instruction) { + if (!StringUtils.hasText(instruction)) { + return; + } + List 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 executionVariables) { + if (executionVariables == null) { + return prompt; + } + Object raw = executionVariables.get(PROMPT_DIRECTIVES_KEY); + if (!(raw instanceof Iterable iterable)) { + return prompt; + } + List 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; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/InputBiasBehaviorAdapter.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/InputBiasBehaviorAdapter.java new file mode 100644 index 0000000..0ef8921 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/InputBiasBehaviorAdapter.java @@ -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 targets = annotation.behavioralProbe().targetInputs(); + List 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 supportedModes(Block block) { + return Set.of(BiasActivationMode.INPUT_TRANSFORMATION); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/MockResponseBiasBehaviorAdapter.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/MockResponseBiasBehaviorAdapter.java new file mode 100644 index 0000000..626a2ca --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/MockResponseBiasBehaviorAdapter.java @@ -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 result = new LinkedHashMap<>(); + execution.getBlock().getOutputs().forEach(output -> + result.put(output.getName(), annotation.behavioralProbe().instruction())); + execution.setBypassResult(result); + } + + @Override + public Set 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; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/OutputBiasBehaviorAdapter.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/OutputBiasBehaviorAdapter.java new file mode 100644 index 0000000..2113a89 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/OutputBiasBehaviorAdapter.java @@ -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 supportedModes(Block block) { + return Set.of(BiasActivationMode.OUTPUT_TRANSFORMATION); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/PromptBiasBehaviorAdapter.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/PromptBiasBehaviorAdapter.java new file mode 100644 index 0000000..dcc677a --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/PromptBiasBehaviorAdapter.java @@ -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 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; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/RoutingBiasBehaviorAdapter.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/RoutingBiasBehaviorAdapter.java new file mode 100644 index 0000000..51e4e10 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/RoutingBiasBehaviorAdapter.java @@ -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 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; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/NodeExecutors.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/NodeExecutors.java index 9b002d7..2f7a909 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/NodeExecutors.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/NodeExecutors.java @@ -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 execute(FlowNode node, List inputs, Map authorizations, Map executionVariables, Map executionVariableDescriptors, ExecutionEventLogger eventLogger) { + return execute(node, inputs, authorizations, executionVariables, executionVariableDescriptors, eventLogger, + BiasExecutionContext.normal()); + } + + public static Map execute(FlowNode node, List inputs, Map authorizations, + Map executionVariables, Map 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 simulate(FlowNode node, List inputs, Map authorizations, Map executionVariables, Map executionVariableDescriptors, LLMDescriptor simulatorDescriptor, ExecutionEventLogger eventLogger) { + return simulate(node, inputs, authorizations, executionVariables, executionVariableDescriptors, + simulatorDescriptor, eventLogger, BiasExecutionContext.normal()); + } + + public static Map simulate(FlowNode node, List inputs, Map authorizations, + Map executionVariables, Map 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 partialResults, Map authorizations, Map executionVariables, Map executionVariableDescriptors, ExecutionEventLogger eventLogger) { + return interact(node, inputs, interaction, partialResults, authorizations, executionVariables, + executionVariableDescriptors, eventLogger, BiasExecutionContext.normal()); + } + + public static InteractionResult interact(FlowNode node, List inputs, Map interaction, + Map partialResults, Map authorizations, + Map executionVariables, Map 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 executeBlock(Block block, List inputs, Map authorizations, Map executionVariables, Map 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 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 simulateBlock(Block block, List inputs, Map authorizations, Map executionVariables, Map 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 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 inputs, Map interaction, Map partialResults, Map authorizations, Map executionVariables, Map 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" }) diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/ChatInteractionExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/ChatInteractionExecutor.java index 01e40c0..719c357 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/ChatInteractionExecutor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/ChatInteractionExecutor.java @@ -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 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 history = existingHistory(partialResults); List messages = history.stream() .map(this::parseHistoryLine) diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/ConditionalExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/ConditionalExecutor.java index eef5e24..7268950 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/ConditionalExecutor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/ConditionalExecutor.java @@ -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 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) { diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/LLMExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/LLMExecutor.java index 6a0a21d..260eda3 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/LLMExecutor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/LLMExecutor.java @@ -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 { 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()); diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/MCPAgentChatExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/MCPAgentChatExecutor.java index d3340c5..2f143de 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/MCPAgentChatExecutor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/MCPAgentChatExecutor.java @@ -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 history = new ArrayList<>(); String sharedSessionKey = normalize(Boolean.TRUE.equals(configuration.getUseSharedSession()) ? configuration.getSharedSessionRef() @@ -128,6 +130,7 @@ public class MCPAgentChatExecutor implements BlockExecutor { Map 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); diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/SwitchExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/SwitchExecutor.java index c219e35..3b4f40e 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/SwitchExecutor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/SwitchExecutor.java @@ -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 { 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) { diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/persistence/ExecutionSnapshot.java b/src/main/java/it/cnr/isti/workflow/manager/executions/persistence/ExecutionSnapshot.java index b0b3113..dd1aef4 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/persistence/ExecutionSnapshot.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/persistence/ExecutionSnapshot.java @@ -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 executionVariables; private Map executionVariableDescriptors; private Map globalInputs; diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Input.java b/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Input.java index 9b5c160..f8ccfb8 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Input.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Input.java @@ -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; diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Step.java b/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Step.java index d6cbd44..5af01e5 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Step.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Step.java @@ -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 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 implements InputListener { try { Map 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 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); diff --git a/src/main/java/it/cnr/isti/workflow/manager/flows/model/bias/BiasActivationMode.java b/src/main/java/it/cnr/isti/workflow/manager/flows/model/bias/BiasActivationMode.java new file mode 100644 index 0000000..5af25c3 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/flows/model/bias/BiasActivationMode.java @@ -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; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/flows/model/bias/BiasBehavioralProbe.java b/src/main/java/it/cnr/isti/workflow/manager/flows/model/bias/BiasBehavioralProbe.java new file mode 100644 index 0000000..d58209d --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/flows/model/bias/BiasBehavioralProbe.java @@ -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 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(); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/flows/model/bias/BlockBiasAnnotation.java b/src/main/java/it/cnr/isti/workflow/manager/flows/model/bias/BlockBiasAnnotation.java index 08a7f3b..9d0508b 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/flows/model/bias/BlockBiasAnnotation.java +++ b/src/main/java/it/cnr/isti/workflow/manager/flows/model/bias/BlockBiasAnnotation.java @@ -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; diff --git a/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowDataValidator.java b/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowDataValidator.java index c8cd4f7..3248b15 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowDataValidator.java +++ b/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowDataValidator.java @@ -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 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 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)); + }); } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/flows/validation/ValidationErrorCode.java b/src/main/java/it/cnr/isti/workflow/manager/flows/validation/ValidationErrorCode.java index 0482e13..1071a97 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/flows/validation/ValidationErrorCode.java +++ b/src/main/java/it/cnr/isti/workflow/manager/flows/validation/ValidationErrorCode.java @@ -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, diff --git a/src/main/resources/application.properties b/src/main/resources/application.properties index 7179a12..e95f2d5 100644 --- a/src/main/resources/application.properties +++ b/src/main/resources/application.properties @@ -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 diff --git a/src/main/resources/db/migration/V2__create_bias_impact_reports.sql b/src/main/resources/db/migration/V2__create_bias_impact_reports.sql new file mode 100644 index 0000000..01e6a51 --- /dev/null +++ b/src/main/resources/db/migration/V2__create_bias_impact_reports.sql @@ -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); diff --git a/src/test/java/it/cnr/isti/workflow/manager/controllers/BiasAnnotationsControllerTest.java b/src/test/java/it/cnr/isti/workflow/manager/controllers/BiasAnnotationsControllerTest.java index beb8d2c..f177c84 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/controllers/BiasAnnotationsControllerTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/controllers/BiasAnnotationsControllerTest.java @@ -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 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", diff --git a/src/test/java/it/cnr/isti/workflow/manager/executions/bias/BiasExperimentsIntegrationTest.java b/src/test/java/it/cnr/isti/workflow/manager/executions/bias/BiasExperimentsIntegrationTest.java new file mode 100644 index 0000000..37c5fc3 --- /dev/null +++ b/src/test/java/it/cnr/isti/workflow/manager/executions/bias/BiasExperimentsIntegrationTest.java @@ -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 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 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 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 canonical = httpBlockFactory.create(configuration); + Block block = Block.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 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 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 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 canonical = llmBlockFactory.create(configuration); + return Block.builder() + .specificConfiguration(configuration) + .inputs(canonical.getInputs()) + .outputs(canonical.getOutputs()) + .type(llmBlockType) + .biasAnnotation(annotation) + .build(); + } + + private BlockBiasAnnotation annotation(BiasActivationMode mode, String instruction, List 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.")); + } +}