diff --git a/src/main/java/it/cnr/isti/workflow/manager/auth/services/AuthService.java b/src/main/java/it/cnr/isti/workflow/manager/auth/services/AuthService.java index 00046cf..1354b3c 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/auth/services/AuthService.java +++ b/src/main/java/it/cnr/isti/workflow/manager/auth/services/AuthService.java @@ -16,6 +16,7 @@ import it.cnr.isti.workflow.manager.auth.repo.AuthRepository; import it.cnr.isti.workflow.manager.auth.repo.LoginEntity; import it.cnr.isti.workflow.manager.executions.repo.ExecutionRepository; import it.cnr.isti.workflow.manager.flows.repo.FlowRepository; +import it.cnr.isti.workflow.manager.projects.repo.ProjectRepository; import static it.cnr.isti.workflow.manager.auth.services.PasswordHasher.*; @@ -34,6 +35,9 @@ public class AuthService { @Autowired ExecutionRepository executionRepository; + @Autowired + ProjectRepository projectRepository; + public boolean validateUser(String username, String password) throws UsernameNotFoundException { LoginEntity le = authRepository.findByUsernameAndActiveTrue(username).orElse(null); if (le == null) { @@ -106,6 +110,11 @@ public class AuthService { } executionRepository.deleteByOwner(username); + // Detach before dropping the projects: the flows deliberately kept below are the finalized + // ones, and leaving them pointing at deleted projects is what a project foreign key would + // otherwise cascade away - destroying exactly the rows this code preserves. + flowRepository.clearProjectAssignmentForOwner(username); + projectRepository.deleteByOwner(username); flowRepository.deleteAll(flowRepository.findByOwnerAndFinalizedFalse(username)); user.setActive(false); authRepository.save(user); diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/PlaceholderInputs.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/PlaceholderInputs.java index e542ea6..dab6cb5 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/PlaceholderInputs.java +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/PlaceholderInputs.java @@ -18,13 +18,15 @@ final class PlaceholderInputs { return placeholder != null && (placeholder.startsWith("vars.") || placeholder.startsWith("global.") - || placeholder.startsWith("context.")); + || placeholder.startsWith("context.") + || placeholder.startsWith("project.")); } static boolean isRuntimeSpelVariable(String variableName) { return "vars".equals(variableName) || "global".equals(variableName) || "context".equals(variableName) + || "project".equals(variableName) || "iteration".equals(variableName); } diff --git a/src/main/java/it/cnr/isti/workflow/manager/controllers/FlowController.java b/src/main/java/it/cnr/isti/workflow/manager/controllers/FlowController.java index 025c81e..d72258e 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/controllers/FlowController.java +++ b/src/main/java/it/cnr/isti/workflow/manager/controllers/FlowController.java @@ -24,6 +24,8 @@ import it.cnr.isti.workflow.manager.flows.FlowNotFoundException; import it.cnr.isti.workflow.manager.flows.FlowService; import it.cnr.isti.workflow.manager.flows.model.FlowCreateRequest; import it.cnr.isti.workflow.manager.flows.model.FlowFlagUpdateRequest; +import it.cnr.isti.workflow.manager.flows.model.FlowProjectUpdateRequest; +import it.cnr.isti.workflow.manager.flows.model.FlowSummaryView; import it.cnr.isti.workflow.manager.flows.model.FlowView; import it.cnr.isti.workflow.manager.flows.validation.GroupedFlowValidation; import it.cnr.isti.workflow.manager.flows.validation.ValidationError; @@ -44,6 +46,14 @@ public class FlowController { return flowService.getAllFlows(userDetails.getUsername()); } + @GetMapping("/summaries") + @Operation(summary = "Get all flow summaries", + description = "Returns the same listing as GET /flows but without the flow graphs, plus each flow's " + + "project. Project membership is disclosed only to the flow's owner.") + public List getAllFlowSummaries(@AuthenticationPrincipal LoginEntity userDetails) { + return flowService.getAllFlowSummaries(userDetails.getUsername()); + } + @GetMapping("/{id}") @Operation(summary = "Get flow", description = "Returns a single flow visible to the authenticated user.") public ResponseEntity getFlow(@PathVariable String id, @AuthenticationPrincipal LoginEntity userDetails) { @@ -93,7 +103,7 @@ public class FlowController { @PostMapping @Operation(summary = "Create flow", description = "Creates a new flow owned by the authenticated user.") public ResponseEntity createFlow(@RequestBody @Valid FlowCreateRequest flow, - @AuthenticationPrincipal + @AuthenticationPrincipal LoginEntity userDetails) { try { logger.debug("Creating flow for user {}", userDetails.getUsername()); @@ -168,6 +178,25 @@ public class FlowController { } } + @PutMapping("/{id}/project") + @Operation(summary = "Move flow to a project", + description = "Assigns a flow owned by the authenticated user to one of their projects, or detaches it " + + "when projectId is null. A finalized flow can still be moved: project membership is workspace " + + "metadata, not flow content.") + public ResponseEntity updateProject(@PathVariable String id, + @RequestBody @Valid FlowProjectUpdateRequest request, + @AuthenticationPrincipal LoginEntity userDetails) { + try { + return ResponseEntity.ok(flowService.updateProject(id, userDetails.getUsername(), request)); + } catch (FlowNotFoundException e) { + logger.warn("Update project for flow {} failed for user {}: not found", id, userDetails.getUsername()); + return ResponseEntity.status(HttpStatus.NOT_FOUND).build(); + } catch (FlowAccessDeniedException e) { + logger.warn("Update project for flow {} forbidden for user {}", id, userDetails.getUsername()); + return ResponseEntity.status(HttpStatus.FORBIDDEN).build(); + } + } + @PutMapping("/{id}/finalized") @Operation(summary = "Update finalized flag", description = "Updates the finalized flag of a flow owned by the authenticated user.") public ResponseEntity updateFinalized(@PathVariable String id, diff --git a/src/main/java/it/cnr/isti/workflow/manager/controllers/ProjectController.java b/src/main/java/it/cnr/isti/workflow/manager/controllers/ProjectController.java new file mode 100644 index 0000000..391252d --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/controllers/ProjectController.java @@ -0,0 +1,152 @@ +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.DeleteMapping; +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.PutMapping; +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.ResponseStatus; +import org.springframework.web.bind.annotation.RestController; + +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.projects.ProjectExecutionService; +import it.cnr.isti.workflow.manager.projects.ProjectService; +import it.cnr.isti.workflow.manager.projects.model.ProjectContext; +import it.cnr.isti.workflow.manager.projects.model.ProjectCreateRequest; +import it.cnr.isti.workflow.manager.projects.model.ProjectExecuteRequest; +import it.cnr.isti.workflow.manager.projects.model.ProjectFlowOrderRequest; +import it.cnr.isti.workflow.manager.projects.model.ProjectExecutionPlanView; +import it.cnr.isti.workflow.manager.projects.model.ProjectRunView; +import it.cnr.isti.workflow.manager.projects.model.ProjectSummaryView; +import it.cnr.isti.workflow.manager.projects.model.ProjectUpdateRequest; +import it.cnr.isti.workflow.manager.projects.model.ProjectView; +import jakarta.validation.Valid; + +/** + * Projects are strictly owner-scoped: another user's project is reported as 404 rather than 403, + * because a project is private workspace structure whose existence should not leak - unlike a flow, + * which lives in a shared, publishable space. + */ +@RestController +@RequestMapping("/projects") +@SecurityRequirement(name = "bearerAuth") +public class ProjectController { + + private final ProjectService projectService; + private final ProjectExecutionService projectExecutionService; + + public ProjectController(ProjectService projectService, ProjectExecutionService projectExecutionService) { + this.projectService = projectService; + this.projectExecutionService = projectExecutionService; + } + + @GetMapping + @Operation(summary = "List projects", + description = "Returns the authenticated user's projects with their flow counts. Carries no flows and no " + + "flow graphs.") + public List list(@AuthenticationPrincipal LoginEntity userDetails) { + return projectService.list(userDetails.getUsername()); + } + + @GetMapping("/{projectId}") + @Operation(summary = "Get project", + description = "Returns a project owned by the authenticated user, its shared context, and summaries of " + + "the flows it contains.") + public ProjectView get(@PathVariable String projectId, @AuthenticationPrincipal LoginEntity userDetails) { + return projectService.get(userDetails.getUsername(), projectId); + } + + @PostMapping + @ResponseStatus(HttpStatus.CREATED) + @Operation(summary = "Create project", description = "Creates a project owned by the authenticated user.") + public ProjectView create(@RequestBody @Valid ProjectCreateRequest request, + @AuthenticationPrincipal LoginEntity userDetails) { + return projectService.create(userDetails.getUsername(), request); + } + + @PutMapping("/{projectId}") + @Operation(summary = "Update project", + description = "Updates name and description. A null field is left unchanged. The shared context is not " + + "editable here - it has its own endpoint, so clearing it is always explicit.") + public ProjectView update(@PathVariable String projectId, @RequestBody @Valid ProjectUpdateRequest request, + @AuthenticationPrincipal LoginEntity userDetails) { + return projectService.update(userDetails.getUsername(), projectId, request); + } + + @PutMapping("/{projectId}/context") + @Operation(summary = "Replace shared context", + description = "Replaces the values the project shares with its flows. Each entry name is referenced from a " + + "flow as ${{project.}}, so names must be valid placeholder names and must not start with a " + + "reserved prefix (project., global., context., vars.).") + public ProjectView updateContext(@PathVariable String projectId, @RequestBody @Valid ProjectContext context, + @AuthenticationPrincipal LoginEntity userDetails) { + return projectService.updateContext(userDetails.getUsername(), projectId, context); + } + + @PostMapping("/{projectId}/execute") + @ResponseStatus(HttpStatus.CREATED) + @Operation(summary = "Run a project", + description = "Creates one execution per flow in the project, all sharing a projectRunId. It does NOT " + + "start them: as with POST /executions, the client supplies global inputs and credentials and " + + "then calls PUT /executions/{id}/start. If any flow is not executable the whole call is refused " + + "with 409 and nothing is created - pass skipNonExecutable=true to run the rest and have the " + + "skipped flows reported back.") + public ProjectExecutionPlanView execute(@PathVariable String projectId, + @RequestBody(required = false) ProjectExecuteRequest request, + @RequestParam(defaultValue = "false") boolean skipNonExecutable, + @AuthenticationPrincipal LoginEntity userDetails) { + return projectExecutionService.execute(userDetails.getUsername(), projectId, + request == null ? null : request.name(), skipNonExecutable); + } + + @PostMapping("/{projectId}/runs/{projectRunId}/start") + @Operation(summary = "Start or resume a project run", + description = "Runs the project's flows one at a time, in the project's order: each step starts only when " + + "the previous one succeeded. Idempotent - calling it while a step is running starts nothing. A " + + "failed step stops the run, and a step still missing inputs or credentials blocks it; fix either " + + "and call this again to resume from where it stopped.") + public ProjectRunView startRun(@PathVariable String projectId, @PathVariable String projectRunId, + @AuthenticationPrincipal LoginEntity userDetails) { + return projectExecutionService.startRun(userDetails.getUsername(), projectId, projectRunId); + } + + @PutMapping("/{projectId}/flow-order") + @Operation(summary = "Reorder the project's flows", + description = "Sets the order a project run executes the flows in. Flows omitted from the list keep their " + + "relative order after the listed ones.") + public ProjectView updateFlowOrder(@PathVariable String projectId, + @RequestBody @Valid ProjectFlowOrderRequest request, + @AuthenticationPrincipal LoginEntity userDetails) { + return projectService.updateFlowOrder(userDetails.getUsername(), projectId, request.flowIds()); + } + + @GetMapping("/{projectId}/runs") + @Operation(summary = "List project runs", + description = "Every run of the project, most recent first. Distinct from GET /executions/groups, which " + + "groups the rerun history of a single flow and is unaffected by project runs.") + public List listRuns(@PathVariable String projectId, + @AuthenticationPrincipal LoginEntity userDetails) { + return projectExecutionService.listRuns(userDetails.getUsername(), projectId); + } + + @DeleteMapping("/{projectId}") + @ResponseStatus(HttpStatus.NO_CONTENT) + @Operation(summary = "Delete project", + description = "Deletes the project and every flow it contains, finalized flows included. Without " + + "confirm=true a non-empty project is refused with 409, whose errors[] lists the flows that would " + + "be destroyed. Executions are kept: each holds its own snapshot of the flow graph.") + public void delete(@PathVariable String projectId, + @RequestParam(defaultValue = "false") boolean confirm, + @AuthenticationPrincipal LoginEntity userDetails) { + projectService.delete(userDetails.getUsername(), projectId, confirm); + } +} 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 350a99c..a7c502e 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 @@ -36,6 +36,12 @@ public class ExecutionContext implements ExecutionListener { Map globalInputs = new HashMap<>(); Map globalInputDescriptors = new HashMap<>(); + /** + * Values inherited from the flow's project, captured once when the execution is created and + * never re-read: editing the project afterwards must not rewrite a run that already happened. + */ + Map projectContext = new HashMap<>(); + @Getter(AccessLevel.NONE) @JsonIgnore Map runtimeContextVariables = new HashMap<>(); @@ -533,6 +539,7 @@ public class ExecutionContext implements ExecutionListener { this.executionVariableDescriptors.clear(); this.globalInputs.clear(); this.globalInputDescriptors.clear(); + this.projectContext.clear(); this.runtimeContextVariables.clear(); this.runtimeExecutionVariables.clear(); this.errors.clear(); @@ -555,7 +562,7 @@ public class ExecutionContext implements ExecutionListener { public ExecutionSnapshot snapshot(Map providedAuthorizations) { return ExecutionSnapshot.builder() - .snapshotVersion(3) + .snapshotVersion(4) .status(this.status) .startTime(this.startTime) .endTime(this.endTime) @@ -566,6 +573,7 @@ public class ExecutionContext implements ExecutionListener { .executionVariableDescriptors(new HashMap<>(this.executionVariableDescriptors)) .globalInputs(new HashMap<>(this.globalInputs)) .globalInputDescriptors(new HashMap<>(this.globalInputDescriptors)) + .projectContext(new HashMap<>(this.projectContext)) .providedAuthorizations(providedAuthorizations == null ? Map.of() : new HashMap<>(providedAuthorizations)) .inputs(new HashMap<>(this.inputs)) .result(new HashMap<>(this.result)) @@ -663,6 +671,10 @@ public class ExecutionContext implements ExecutionListener { if (snapshot.getWaitingSteps() != null) { this.waitingSteps.addAll(snapshot.getWaitingSteps()); } + this.projectContext.clear(); + if (snapshot.getProjectContext() != null) { + this.projectContext.putAll(snapshot.getProjectContext()); + } this.startTime = snapshot.getStartTime(); this.endTime = snapshot.getEndTime(); this.setInteractionSimulationEnabled(snapshot.isInteractionSimulationEnabled()); @@ -723,6 +735,14 @@ public class ExecutionContext implements ExecutionListener { } } + protected void setProjectContext(Map projectContext) { + this.projectContext.clear(); + if (projectContext != null) { + this.projectContext.putAll(projectContext); + } + refreshRuntimeExecutionVariables(); + } + protected void setRuntimeContextValues(Map runtimeContextValues) { this.runtimeContextVariables.clear(); if (runtimeContextValues != null) { @@ -737,6 +757,7 @@ public class ExecutionContext implements ExecutionListener { this.runtimeExecutionVariables.putAll(ExecutionRuntimeContextSupport.runtimeValues( this.executionVariables, this.globalInputs, + this.projectContext, this.runtimeContextVariables)); } 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 f2de03c..a7b4c37 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 @@ -92,11 +92,21 @@ public class ExecutionObject { BiasExecutionContext biasExecutionContext = BiasExecutionContext.normal(); + /** The project the source flow belonged to when this run was created, if any. */ + String projectId; + + /** Groups the executions started together by one project run. Null for a single-flow run. */ + String projectRunId; + + /** Position in the project run; the run starts one execution at a time in this order. */ + Integer projectRunOrder; + @Builder public ExecutionObject(String executionName, FlowData flow, List requiredAuthorizations, String owner, String runGroupId, String sourceFlowId, String rerunOfExecutionId, Integer runNumber, BiasExecutionContext biasExecutionContext, ExecutionKind executionKind, String parentExecutionId, - String parentStepId, Integer parentIterationIndex, ContainerSubflowRole subflowRole) { + String parentStepId, Integer parentIterationIndex, ContainerSubflowRole subflowRole, + String projectId, String projectRunId, Integer projectRunOrder, Map projectContext) { this.name = executionName; this.owner = owner; this.executionKind = executionKind == null ? ExecutionKind.TOP_LEVEL : executionKind; @@ -116,6 +126,9 @@ public class ExecutionObject { 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; + this.projectId = projectId; + this.projectRunId = projectRunId; + this.projectRunOrder = projectRunOrder; List> steps = getStepsFromFlow(flow); @@ -124,6 +137,7 @@ public class ExecutionObject { this.context = new ExecutionContext(steps.stream().collect(Collectors.toMap(Step::getId, Function.identity()))); this.context.getSteps().values().forEach(step -> step.setParentExecutionId(this.id)); this.context.setBiasExecutionContext(this.biasExecutionContext); + this.context.setProjectContext(projectContext); ensureRuntimeContextVariables(); registerGlobalInputs(); this.requiredAuthorizations.stream() diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionRuntimeContextSupport.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionRuntimeContextSupport.java index aeabd54..34d1389 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionRuntimeContextSupport.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionRuntimeContextSupport.java @@ -9,6 +9,7 @@ public final class ExecutionRuntimeContextSupport { public static final String CONTEXT_PREFIX = "context."; public static final String GLOBAL_PREFIX = "global."; + public static final String PROJECT_PREFIX = "project."; public static final String EXECUTION_ID = CONTEXT_PREFIX + "executionId"; public static final String EXECUTION_NAME = CONTEXT_PREFIX + "executionName"; public static final String EXECUTION_OWNER = CONTEXT_PREFIX + "executionOwner"; @@ -29,6 +30,22 @@ public final class ExecutionRuntimeContextSupport { public static Map runtimeValues(Map executionVariables, Map globalInputs, Map contextValues) { + return runtimeValues(executionVariables, globalInputs, null, contextValues); + } + + /** + * Merges every source of runtime values into the single flat map the {@code ${{...}}} resolver + * reads. + * + *

Each source keeps its own prefix, so {@code project.x}, {@code global.x}, {@code context.x} + * and a bare execution variable {@code x} are four distinct keys: there is no precedence + * question to settle. The only way to manufacture a collision is an entry name that already + * carries a reserved prefix, which is refused when the project context is saved. + * + *

Context values go in last so the system-reserved namespace stays authoritative. + */ + public static Map runtimeValues(Map executionVariables, Map globalInputs, + Map projectValues, Map contextValues) { LinkedHashMap runtime = new LinkedHashMap<>(); if (executionVariables != null) { runtime.putAll(executionVariables); @@ -36,6 +53,9 @@ public final class ExecutionRuntimeContextSupport { if (globalInputs != null) { globalInputs.forEach((key, value) -> runtime.put(GLOBAL_PREFIX + key, value)); } + if (projectValues != null) { + projectValues.forEach((key, value) -> runtime.put(PROJECT_PREFIX + key, value)); + } if (contextValues != null) { runtime.putAll(contextValues); } @@ -55,16 +75,24 @@ public final class ExecutionRuntimeContextSupport { return context; } + public static Map projectView(Map runtimeVariables) { + return viewOfPrefix(runtimeVariables, PROJECT_PREFIX); + } + public static Map globalView(Map runtimeVariables) { - LinkedHashMap global = new LinkedHashMap<>(); + return viewOfPrefix(runtimeVariables, GLOBAL_PREFIX); + } + + private static Map viewOfPrefix(Map runtimeVariables, String prefix) { + LinkedHashMap view = new LinkedHashMap<>(); if (runtimeVariables == null) { - return global; + return view; } runtimeVariables.forEach((key, value) -> { - if (StringUtils.hasText(key) && key.startsWith(GLOBAL_PREFIX) && key.length() > GLOBAL_PREFIX.length()) { - global.put(key.substring(GLOBAL_PREFIX.length()), value); + if (StringUtils.hasText(key) && key.startsWith(prefix) && key.length() > prefix.length()) { + view.put(key.substring(prefix.length()), value); } }); - return global; + return view; } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionTemplateResolver.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionTemplateResolver.java index de0d92d..de1bbb4 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionTemplateResolver.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionTemplateResolver.java @@ -50,7 +50,8 @@ public final class ExecutionTemplateResolver { } String formatted = formatValue(entry.getValue()); if (entry.getKey().startsWith(ExecutionRuntimeContextSupport.GLOBAL_PREFIX) - || entry.getKey().startsWith(ExecutionRuntimeContextSupport.CONTEXT_PREFIX)) { + || entry.getKey().startsWith(ExecutionRuntimeContextSupport.CONTEXT_PREFIX) + || entry.getKey().startsWith(ExecutionRuntimeContextSupport.PROJECT_PREFIX)) { substitutions.put("${{" + entry.getKey() + "}}", formatted); // Also accept the "[]" (array) marker on a global/context reference, so a // ${{global.name[]}} placeholder (used to declare the global as multi-valued) 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 13d97b9..0aef4a5 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 @@ -13,6 +13,7 @@ import java.util.stream.Collectors; import jakarta.annotation.PostConstruct; +import org.springframework.beans.factory.ObjectProvider; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.context.event.ApplicationReadyEvent; @@ -56,6 +57,8 @@ 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; import it.cnr.isti.workflow.manager.flows.model.FlowData; +import it.cnr.isti.workflow.manager.projects.ProjectContextResolver; +import it.cnr.isti.workflow.manager.projects.ProjectRunOrchestrator; import it.cnr.isti.workflow.manager.flows.validation.FlowExecutionValidator; import it.cnr.isti.workflow.manager.llms.LLMDescriptor; import it.cnr.isti.workflow.manager.llms.providers.LLMProvider; @@ -80,6 +83,16 @@ public class ExecutionsService { @Autowired FlowExecutionValidator flowExecutionValidator; + @Autowired + ProjectContextResolver projectContextResolver; + + /** + * Looked up lazily: the orchestrator needs this service to start executions, so a direct + * injection either way would be a cycle. + */ + @Autowired + ObjectProvider projectRunOrchestrator; + @Autowired Map llmProviders; @@ -164,8 +177,90 @@ public class ExecutionsService { @Transactional public ExecutionObject createExecutionForFlow(String sourceFlowId, String executionName, FlowData flow, String owner) { + return createExecutionForFlow(sourceFlowId, executionName, flow, owner, null); + } + + /** + * Creates a top-level execution for a flow, inheriting its project's shared context. + * + *

The context is resolved and frozen here, at creation: a later edit to the project must not + * retroactively change a run that already exists. {@code projectRunId} ties together the + * executions started by one project run and is null for an ordinary single-flow run. + */ + @Transactional + public ExecutionObject createExecutionForFlow(String sourceFlowId, String executionName, FlowData flow, String owner, + String projectRunId) { + return createExecutionForFlow(sourceFlowId, executionName, flow, owner, projectRunId, null); + } + + @Transactional + public ExecutionObject createExecutionForFlow(String sourceFlowId, String executionName, FlowData flow, String owner, + String projectRunId, Integer projectRunOrder) { int runNumber = nextRunNumberForFlow(sourceFlowId, owner); - return createExecution(executionName, flow, owner, sourceFlowId, sourceFlowId, null, runNumber); + ProjectContextResolver.ResolvedProjectContext projectContext = + projectContextResolver.resolveForFlow(sourceFlowId, owner); + ExecutionObject execution = createExecution(executionName, flow, owner, sourceFlowId, sourceFlowId, null, + runNumber, BiasExecutionContext.normal(), ExecutionKind.TOP_LEVEL, null, null, null, null, + // The execution is tagged with the flow's project even when the context resolved + // empty (a stranger running a published flow), so grouping still works. + projectContextResolver.projectIdOfFlow(sourceFlowId), projectRunId, projectRunOrder, + projectContext.values()); + seedGlobalInputsFromProjectContext(execution, projectContext.values()); + return execution; + } + + /** + * Pre-fills a flow's global inputs from the project's shared values. + * + *

This is what makes a one-click project run possible: an execution whose inputs and + * authorizations are all satisfied becomes READY by itself, so the run can start it without a + * per-flow round trip. + * + *

Matching is deliberately conservative - same name and same declared type, and only inputs + * that have no value yet. A project value never overwrites something the user supplied, and a + * name that happens to collide with a differently-typed input is left alone rather than being + * coerced. + */ + private void seedGlobalInputsFromProjectContext(ExecutionObject execution, Map projectValues) { + if (projectValues == null || projectValues.isEmpty()) { + return; + } + Map seed = new LinkedHashMap<>(); + for (IODescriptor required : execution.getRequiredGlobalInputs()) { + Object value = projectValues.get(required.getName()); + if (value == null) { + continue; + } + ExecutionVariableDescriptor existing = execution.getContext().getGlobalInputDescriptor(required.getName()); + if (existing != null && existing.getValue() != null) { + continue; + } + if (matchesDeclaredType(required, value)) { + seed.put(required.getName(), value); + } else { + logger.debug("Project value {} not seeded into global input: declared type {} does not match", + required.getName(), required.getType()); + } + } + if (!seed.isEmpty()) { + execution.setGlobalInputs(seed); + persist(execution); + } + } + + private boolean matchesDeclaredType(IODescriptor descriptor, Object value) { + if (descriptor.isMultiple() != (value instanceof java.util.Collection)) { + return false; + } + if (descriptor.getType() == null) { + return false; + } + return switch (descriptor.getType()) { + case TEXT, CSV, JSON, ANY -> true; + case BOOLEAN -> value instanceof Boolean || value instanceof String; + // A file is an uploaded artifact, never a literal carried in a project's shared values. + case FILE -> false; + }; } @Transactional @@ -186,6 +281,16 @@ public class ExecutionsService { String sourceFlowId, String rerunOfExecutionId, Integer runNumber, BiasExecutionContext biasExecutionContext, ExecutionKind executionKind, String parentExecutionId, String parentStepId, Integer parentIterationIndex, ContainerSubflowRole subflowRole) { + return createExecution(executionName, flow, owner, runGroupId, sourceFlowId, rerunOfExecutionId, runNumber, + biasExecutionContext, executionKind, parentExecutionId, parentStepId, parentIterationIndex, + subflowRole, null, null, null, null); + } + + private ExecutionObject createExecution(String executionName, FlowData flow, String owner, String runGroupId, + String sourceFlowId, String rerunOfExecutionId, Integer runNumber, + BiasExecutionContext biasExecutionContext, ExecutionKind executionKind, String parentExecutionId, + String parentStepId, Integer parentIterationIndex, ContainerSubflowRole subflowRole, + String projectId, String projectRunId, Integer projectRunOrder, Map projectContext) { flowExecutionValidator.validate(flow); List requiredAuthorizations = AuthorizationRequirementResolver.resolveRequiredAuthorizations(flow, llmProviders); ExecutionObject execObject = ExecutionObject.builder() @@ -203,6 +308,10 @@ public class ExecutionsService { .parentStepId(parentStepId) .parentIterationIndex(parentIterationIndex) .subflowRole(subflowRole) + .projectId(projectId) + .projectRunId(projectRunId) + .projectRunOrder(projectRunOrder) + .projectContext(projectContext) .build(); attachPersistence(execObject); executions.put(execObject.getId(), execObject); @@ -1232,10 +1341,35 @@ public class ExecutionsService { executionObject.setStateChangeListener(() -> { persist(executionObject); cleanupManagedResourcesIfFinal(executionObject); + advanceProjectRunIfFinished(executionObject); }); executionObject.setErrorStateListener(executionObject::abortOnError); } + /** + * Moves a project run to its next flow once this one is done. Runs on the finishing execution's + * own thread, but only ever starts a different execution, and the orchestrator + * serializes per run - and never throws into the callback, so a broken run can never take the + * completing execution down with it. + */ + private void advanceProjectRunIfFinished(ExecutionObject executionObject) { + String projectRunId = executionObject.getProjectRunId(); + if (projectRunId == null || projectRunId.isBlank()) { + return; + } + ExecutionStatus status = executionObject.getContext().getStatus(); + if (!status.isFinalState()) { + return; + } + try { + projectRunOrchestrator.getObject() + .onExecutionFinished(executionObject.getOwner(), projectRunId, status); + } catch (RuntimeException e) { + logger.error("Advancing project run {} after execution {} failed", + projectRunId, executionObject.getId(), e); + } + } + private void persist(ExecutionObject executionObject) { executionRepository.save(ExecutionEntity.builder() .id(executionObject.getId()) @@ -1252,6 +1386,9 @@ public class ExecutionsService { .parentStepId(executionObject.getParentStepId()) .parentIterationIndex(executionObject.getParentIterationIndex()) .subflowRole(executionObject.getSubflowRole()) + .projectId(executionObject.getProjectId()) + .projectRunId(executionObject.getProjectRunId()) + .projectRunOrder(executionObject.getProjectRunOrder()) .flow(executionObject.getFlow()) .snapshot(executionObject.snapshot()) .build()); @@ -1276,6 +1413,12 @@ public class ExecutionsService { .parentStepId(entity.getParentStepId()) .parentIterationIndex(entity.getParentIterationIndex()) .subflowRole(entity.getSubflowRole()) + .projectId(entity.getProjectId()) + .projectRunId(entity.getProjectRunId()) + .projectRunOrder(entity.getProjectRunOrder()) + // The snapshot is the only source here: the project may have been edited, or + // deleted, since the run was created. + .projectContext(snapshot == null ? Map.of() : snapshot.getProjectContext()) .build(); executionObject.restore(entity.getId(), entity.getCreationTime(), snapshot == null ? Map.of() : snapshot.getProvidedAuthorizations(), diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/api/ExecutionContextView.java b/src/main/java/it/cnr/isti/workflow/manager/executions/api/ExecutionContextView.java index 3c56e6f..a5b7e81 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/api/ExecutionContextView.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/api/ExecutionContextView.java @@ -21,6 +21,8 @@ public class ExecutionContextView { private Map executionVariableDescriptors; private Map globalInputs; private Map globalInputDescriptors; + /** Values inherited from the flow's project, frozen when the run was created. */ + private Map projectContext; private Map result; private Map partialResult; private Long startTime; @@ -41,6 +43,7 @@ public class ExecutionContextView { .executionVariableDescriptors(context.getExecutionVariableDescriptors()) .globalInputs(context.getGlobalInputs()) .globalInputDescriptors(context.getGlobalInputDescriptors()) + .projectContext(context.getProjectContext()) .result(context.getResult()) .partialResult(context.getPartialResult()) .startTime(context.getStartTime()) 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 08160c9..8031c44 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 @@ -25,6 +25,8 @@ public class ExecutionView { private long creationTime; private String name; private String runGroupId; + private String projectId; + private String projectRunId; private String sourceFlowId; private String rerunOfExecutionId; private int runNumber; @@ -51,6 +53,8 @@ public class ExecutionView { .creationTime(execution.getCreationTime()) .name(execution.getName()) .runGroupId(execution.getRunGroupId()) + .projectId(execution.getProjectId()) + .projectRunId(execution.getProjectRunId()) .sourceFlowId(execution.getSourceFlowId()) .rerunOfExecutionId(execution.getRerunOfExecutionId()) .runNumber(execution.getRunNumber()) diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/ExpressionEvaluationSupport.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/ExpressionEvaluationSupport.java index e17f082..b2ad671 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/ExpressionEvaluationSupport.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/ExpressionEvaluationSupport.java @@ -44,6 +44,7 @@ final class ExpressionEvaluationSupport { context.setVariable("vars", executionVariables == null ? Map.of() : executionVariables); context.setVariable("global", ExecutionRuntimeContextSupport.globalView(executionVariables)); context.setVariable("context", ExecutionRuntimeContextSupport.contextView(executionVariables)); + context.setVariable("project", ExecutionRuntimeContextSupport.projectView(executionVariables)); return context; } } 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 c1da5f9..be5863b 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 @@ -31,6 +31,11 @@ public class ExecutionSnapshot { private Map executionVariableDescriptors; private Map globalInputs; private Map globalInputDescriptors; + /** + * The project values this run inherited, frozen at creation. Snapshot version 4 added it; it + * is a JSON blob, so older rows simply deserialize it as null. + */ + private Map projectContext; private Map providedAuthorizations; private Map inputs; private Map result; diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/repo/ExecutionEntity.java b/src/main/java/it/cnr/isti/workflow/manager/executions/repo/ExecutionEntity.java index e9b1ad2..e2b51f2 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/repo/ExecutionEntity.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/repo/ExecutionEntity.java @@ -39,6 +39,18 @@ public class ExecutionEntity { private String runGroupId; + /** Project the source flow belonged to when the run was created; null for an unassigned flow. */ + @Column(name = "project_id") + private String projectId; + + /** Ties together the executions started by one project run; null for a single-flow run. */ + @Column(name = "project_run_id") + private String projectRunId; + + /** Position in the project run: the sequence starts one flow at a time in this order. */ + @Column(name = "project_run_order") + private Integer projectRunOrder; + private String sourceFlowId; private String rerunOfExecutionId; diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/repo/ExecutionRepository.java b/src/main/java/it/cnr/isti/workflow/manager/executions/repo/ExecutionRepository.java index 5ecf4d2..1ef8a91 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/repo/ExecutionRepository.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/repo/ExecutionRepository.java @@ -20,6 +20,11 @@ public interface ExecutionRepository extends JpaRepository findBySourceFlowIdAndOwnerAndExecutionKindOrderByRunNumberAscCreationTimeAsc( String sourceFlowId, String owner, ExecutionKind executionKind); + List findByProjectRunIdAndOwnerOrderByCreationTimeAsc(String projectRunId, String owner); + + List findByProjectIdAndOwnerAndExecutionKindOrderByCreationTimeAsc( + String projectId, String owner, ExecutionKind executionKind); + List findByParentExecutionIdOrderByCreationTimeAsc(String parentExecutionId); List findByParentExecutionIdAndParentStepIdOrderByParentIterationIndexAscCreationTimeAsc( String parentExecutionId, String parentStepId); diff --git a/src/main/java/it/cnr/isti/workflow/manager/flows/FlowMapper.java b/src/main/java/it/cnr/isti/workflow/manager/flows/FlowMapper.java index df19a7d..d8b4e66 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/flows/FlowMapper.java +++ b/src/main/java/it/cnr/isti/workflow/manager/flows/FlowMapper.java @@ -3,6 +3,7 @@ package it.cnr.isti.workflow.manager.flows; import java.time.LocalDateTime; import it.cnr.isti.workflow.manager.flows.model.FlowCreateRequest; +import it.cnr.isti.workflow.manager.flows.model.FlowSummaryView; import it.cnr.isti.workflow.manager.flows.model.FlowView; import it.cnr.isti.workflow.manager.flows.model.FlowViewStatus; import it.cnr.isti.workflow.manager.flows.repo.FlowEntity; @@ -24,6 +25,11 @@ public final class FlowMapper { return entity; } + /** + * Deliberately does not touch {@code projectId}: this is the body of the full-replace PUT the + * editor issues on every save, and no client sends a project. Assignment goes through + * {@code PUT /flows/{id}/project} instead, so a save can never silently detach a flow. + */ public static void updateEntity(FlowEntity entity, FlowCreateRequest request) { entity.setName(request.name()); entity.setDescription(request.description()); @@ -31,7 +37,12 @@ public final class FlowMapper { entity.setLastUpdateAt(LocalDateTime.now()); } - public static FlowView toView(FlowEntity entity, FlowViewStatus status) { + /** + * {@code projectId} and {@code projectName} are passed in already privacy-filtered rather than + * read off the entity: project membership is disclosed only to the flow's owner, and reading the + * id straight from the entity here would leak it even when the name is withheld. + */ + public static FlowView toView(FlowEntity entity, FlowViewStatus status, String projectId, String projectName) { return new FlowView( entity.getId(), entity.getName(), @@ -42,6 +53,24 @@ public final class FlowMapper { status, entity.isPublished(), entity.isFinalized(), + projectId, + projectName, entity.getFlow()); } + + public static FlowSummaryView toSummaryView(FlowEntity entity, FlowViewStatus status, String projectId, + String projectName) { + return new FlowSummaryView( + entity.getId(), + entity.getName(), + entity.getDescription(), + entity.getCreatedAt(), + entity.getLastUpdateAt(), + entity.getOwner(), + status, + entity.isPublished(), + entity.isFinalized(), + projectId, + projectName); + } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/flows/FlowService.java b/src/main/java/it/cnr/isti/workflow/manager/flows/FlowService.java index 25d28d6..c8df6ce 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/flows/FlowService.java +++ b/src/main/java/it/cnr/isti/workflow/manager/flows/FlowService.java @@ -9,10 +9,14 @@ import org.springframework.web.server.ResponseStatusException; import it.cnr.isti.workflow.manager.flows.model.FlowCreateRequest; import it.cnr.isti.workflow.manager.flows.model.FlowFlagUpdateRequest; +import it.cnr.isti.workflow.manager.flows.model.FlowProjectUpdateRequest; +import it.cnr.isti.workflow.manager.flows.model.FlowSummaryView; import it.cnr.isti.workflow.manager.flows.model.FlowView; import it.cnr.isti.workflow.manager.flows.model.FlowViewStatus; import it.cnr.isti.workflow.manager.flows.repo.FlowEntity; import it.cnr.isti.workflow.manager.flows.repo.FlowRepository; +import it.cnr.isti.workflow.manager.projects.repo.ProjectEntity; +import it.cnr.isti.workflow.manager.projects.repo.ProjectRepository; import it.cnr.isti.workflow.manager.flows.validation.FlowExecutionValidator; import it.cnr.isti.workflow.manager.flows.validation.GroupedFlowValidation; import it.cnr.isti.workflow.manager.flows.validation.ValidationError; @@ -23,6 +27,7 @@ import jakarta.validation.Validator; import java.util.ArrayList; import java.util.EnumSet; +import java.util.Map; import java.util.Set; @Service @@ -39,11 +44,14 @@ public class FlowService { @Autowired FlowExecutionValidator flowExecutionValidator; + @Autowired + ProjectRepository projectRepository; + public FlowView createFlow(String owner, FlowCreateRequest request) { validateFlow(request); FlowEntity flowEntity = FlowMapper.toNewEntity(owner, request); FlowEntity savedEntity = flowRepository.save(flowEntity); - return toView(savedEntity); + return toView(savedEntity, owner); } public FlowView updateFlow(String id, String owner, FlowCreateRequest request) { @@ -59,7 +67,7 @@ public class FlowService { validateFlow(request); FlowMapper.updateEntity(entity, request); FlowEntity savedEntity = flowRepository.save(entity); - return toView(savedEntity); + return toView(savedEntity, owner); } public void deleteFlow(String id, String owner) { @@ -86,7 +94,37 @@ public class FlowService { entity.setPublished(Boolean.TRUE.equals(request.value())); entity.setLastUpdateAt(java.time.LocalDateTime.now()); - return toView(flowRepository.save(entity)); + return toView(flowRepository.save(entity), owner); + } + + /** + * Moves a flow into a project, or detaches it when {@code projectId} is null. + * + *

Deliberately does not call {@link #ensureMutable}: project membership is workspace + * metadata, not flow content - the same reason {@link #updatePublished} does not. Finalizing a + * flow must not freeze it out of reorganization forever, especially as finalizing is + * irreversible. + */ + public FlowView updateProject(String id, String owner, FlowProjectUpdateRequest request) { + FlowEntity entity = flowRepository.findById(id) + .orElseThrow(() -> new FlowNotFoundException(id)); + + if (!entity.getOwner().equals(owner)) { + logger.warn("User {} attempted to change the project of flow {} owned by {}", owner, id, entity.getOwner()); + throw new FlowAccessDeniedException(); + } + + String requestedProjectId = request == null || request.projectId() == null + || request.projectId().isBlank() + ? null + : request.projectId(); + // Resolving through findByIdAndOwner is what prevents attaching a flow to someone + // else's project, and keeps the project.owner == flow.owner invariant total. + entity.setProjectId(requestedProjectId == null + ? null + : requireOwnedProject(requestedProjectId, owner).getId()); + entity.setLastUpdateAt(java.time.LocalDateTime.now()); + return toView(flowRepository.save(entity), owner); } public FlowView updateFinalized(String id, String owner, FlowFlagUpdateRequest request) { @@ -99,18 +137,30 @@ public class FlowService { } if (entity.isFinalized()) { if (Boolean.TRUE.equals(request.value())) { - return toView(entity); + return toView(entity, owner); } throw new ResponseStatusException(HttpStatus.CONFLICT, "Flow is finalized"); } entity.setFinalized(Boolean.TRUE.equals(request.value())); entity.setLastUpdateAt(java.time.LocalDateTime.now()); - return toView(flowRepository.save(entity)); + return toView(flowRepository.save(entity), owner); } public List getAllFlows(String owner) { - return flowRepository.findFlowsByOwnerOrPublic(owner).stream().map(this::toView).toList(); + List entities = flowRepository.findFlowsByOwnerOrPublic(owner); + Map projectNames = projectNamesFor(owner); + return entities.stream().map(entity -> toView(entity, owner, projectNames)).toList(); + } + + /** The same listing without the flow graphs - see {@link FlowSummaryView}. */ + public List getAllFlowSummaries(String owner) { + List entities = flowRepository.findFlowsByOwnerOrPublic(owner); + Map projectNames = projectNamesFor(owner); + return entities.stream() + .map(entity -> FlowMapper.toSummaryView(entity, statusOf(entity), + visibleProjectId(entity, owner), projectNameFor(entity, owner, projectNames))) + .toList(); } public FlowView getFlow(String id, String owner) { @@ -122,7 +172,7 @@ public class FlowService { throw new FlowAccessDeniedException(); } - return toView(entity); + return toView(entity, owner); } public List getFlowValidation(String id, String owner) { @@ -228,10 +278,61 @@ public class FlowService { } } - private FlowView toView(FlowEntity entity) { - FlowViewStatus status = flowExecutionValidator.isExecutable(entity.getFlow()) + private ProjectEntity requireOwnedProject(String projectId, String owner) { + return projectRepository.findByIdAndOwner(projectId, owner) + .orElseThrow(() -> new ResponseStatusException(HttpStatus.NOT_FOUND, + "Project with id " + projectId + " not found")); + } + + private FlowViewStatus statusOf(FlowEntity entity) { + return flowExecutionValidator.isExecutable(entity.getFlow()) ? FlowViewStatus.EXECUTABLE : FlowViewStatus.DRAFT; - return FlowMapper.toView(entity, status); + } + + private FlowView toView(FlowEntity entity, String viewer) { + return FlowMapper.toView(entity, statusOf(entity), visibleProjectId(entity, viewer), + resolveProjectName(entity, viewer)); + } + + private FlowView toView(FlowEntity entity, String viewer, Map projectNames) { + return FlowMapper.toView(entity, statusOf(entity), visibleProjectId(entity, viewer), + projectNameFor(entity, viewer, projectNames)); + } + + /** + * Project membership is disclosed only to the flow's owner. A published flow in someone's + * project stays readable by everyone - project membership is organizational structure, not an + * access-control mechanism - but the project name can itself be private information, so a + * non-owner sees the flow as unassigned. + */ + private boolean maySeeProject(FlowEntity entity, String viewer) { + return entity.getProjectId() != null && viewer != null && viewer.equals(entity.getOwner()); + } + + private String visibleProjectId(FlowEntity entity, String viewer) { + return maySeeProject(entity, viewer) ? entity.getProjectId() : null; + } + + private String resolveProjectName(FlowEntity entity, String viewer) { + if (!maySeeProject(entity, viewer)) { + return null; + } + return projectRepository.findByIdAndOwner(entity.getProjectId(), entity.getOwner()) + .map(ProjectEntity::getName) + .orElse(null); + } + + private String projectNameFor(FlowEntity entity, String viewer, Map projectNames) { + if (!maySeeProject(entity, viewer)) { + return null; + } + return projectNames.get(entity.getProjectId()); + } + + /** Loaded once per listing: never resolve a project name with a query per flow. */ + private Map projectNamesFor(String owner) { + return projectRepository.findByOwnerOrderByNameAsc(owner).stream() + .collect(java.util.stream.Collectors.toMap(ProjectEntity::getId, ProjectEntity::getName)); } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/flows/model/FlowProjectUpdateRequest.java b/src/main/java/it/cnr/isti/workflow/manager/flows/model/FlowProjectUpdateRequest.java new file mode 100644 index 0000000..a0fdf09 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/flows/model/FlowProjectUpdateRequest.java @@ -0,0 +1,12 @@ +package it.cnr.isti.workflow.manager.flows.model; + +/** + * Moves a flow into a project. A null or blank {@code projectId} detaches it. + * + *

Its own endpoint rather than a field on {@link FlowCreateRequest}, which is the body of the + * full-replace PUT the editor issues on every save: a client that omitted the field there would + * silently detach the flow. Mirrors the existing {@code /published} and {@code /finalized} flag + * endpoints, giving idempotent move and detach in one call. + */ +public record FlowProjectUpdateRequest(String projectId) { +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/flows/model/FlowSummaryView.java b/src/main/java/it/cnr/isti/workflow/manager/flows/model/FlowSummaryView.java new file mode 100644 index 0000000..f9598a2 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/flows/model/FlowSummaryView.java @@ -0,0 +1,23 @@ +package it.cnr.isti.workflow.manager.flows.model; + +import java.time.LocalDateTime; + +/** + * A flow without its graph. + * + *

{@link FlowView} embeds the whole {@link FlowData} even in list responses, which dominates the + * payload of any listing. This is the shape a project-grouped list should consume. + */ +public record FlowSummaryView( + String id, + String name, + String description, + LocalDateTime createdAt, + LocalDateTime lastUpdateAt, + String owner, + FlowViewStatus status, + boolean published, + boolean finalized, + String projectId, + String projectName) { +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/flows/model/FlowView.java b/src/main/java/it/cnr/isti/workflow/manager/flows/model/FlowView.java index 9a3ed69..7e4e037 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/flows/model/FlowView.java +++ b/src/main/java/it/cnr/isti/workflow/manager/flows/model/FlowView.java @@ -12,5 +12,11 @@ public record FlowView( FlowViewStatus status, boolean published, boolean finalized, + /** + * Project membership, disclosed only to the flow's owner: a project name can itself be + * private information, and a non-owner reading a published flow has no use for it. + */ + String projectId, + String projectName, FlowData flow) { } diff --git a/src/main/java/it/cnr/isti/workflow/manager/flows/repo/FlowEntity.java b/src/main/java/it/cnr/isti/workflow/manager/flows/repo/FlowEntity.java index 72c3435..5ae4597 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/flows/repo/FlowEntity.java +++ b/src/main/java/it/cnr/isti/workflow/manager/flows/repo/FlowEntity.java @@ -50,6 +50,17 @@ public class FlowEntity { @Builder.Default private boolean finalized = false; + /** + * Owning project, or null for an unassigned flow. A plain id column, not a {@code @ManyToOne}: + * every association in this codebase is modelled this way (see {@code ExecutionEntity}). + */ + @Column(name = "project_id") + private String projectId; + + /** Stable display order within the project. Reserved for project execution; unused for now. */ + @Column(name = "project_order") + private Integer projectOrder; + @Column(name = "flow_data", columnDefinition = "TEXT") @Convert(converter = FlowConverter.class) private FlowData flow; diff --git a/src/main/java/it/cnr/isti/workflow/manager/flows/repo/FlowRepository.java b/src/main/java/it/cnr/isti/workflow/manager/flows/repo/FlowRepository.java index 0d6ab5d..82e0fb5 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/flows/repo/FlowRepository.java +++ b/src/main/java/it/cnr/isti/workflow/manager/flows/repo/FlowRepository.java @@ -3,6 +3,7 @@ package it.cnr.isti.workflow.manager.flows.repo; import java.util.List; import org.springframework.data.jpa.repository.JpaRepository; +import org.springframework.data.jpa.repository.Modifying; import org.springframework.data.jpa.repository.Query; import org.springframework.data.repository.query.Param; import org.springframework.stereotype.Repository; @@ -19,4 +20,27 @@ public interface FlowRepository extends JpaRepository { List findByOwner(String owner); List findByOwnerAndFinalizedFalse(String owner); + + List findByProjectIdOrderByProjectOrderAscNameAsc(String projectId); + + List findByOwnerAndProjectId(String owner, String projectId); + + long countByProjectId(String projectId); + + /** + * Flow counts for every project of an owner in one query, so a project listing never issues a + * count per project. + */ + @Query("SELECT f.projectId AS projectId, COUNT(f) AS flowCount FROM FlowEntity f " + + "WHERE f.owner = :owner AND f.projectId IS NOT NULL GROUP BY f.projectId") + List countFlowsByProjectForOwner(@Param("owner") String owner); + + /** + * Detaches every flow of an owner from its project. Used before deleting that owner's projects + * on account deletion: {@code AuthService} deliberately preserves finalized flows, and the + * {@code project} foreign key would otherwise cascade them away. + */ + @Modifying + @Query("UPDATE FlowEntity f SET f.projectId = null WHERE f.owner = :owner") + int clearProjectAssignmentForOwner(@Param("owner") String owner); } diff --git a/src/main/java/it/cnr/isti/workflow/manager/flows/repo/ProjectFlowCount.java b/src/main/java/it/cnr/isti/workflow/manager/flows/repo/ProjectFlowCount.java new file mode 100644 index 0000000..b4f8715 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/flows/repo/ProjectFlowCount.java @@ -0,0 +1,9 @@ +package it.cnr.isti.workflow.manager.flows.repo; + +/** Projection for {@link FlowRepository#countFlowsByProjectForOwner(String)}. */ +public interface ProjectFlowCount { + + String getProjectId(); + + long getFlowCount(); +} 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 cec47cc..f7f3876 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 @@ -86,5 +86,10 @@ public enum ValidationErrorCode { GLOBAL_INPUT_NOT_DECLARED, SHARED_SESSION_NAME_NOT_UNIQUE, SHARED_SESSION_PRODUCER_UNREACHABLE, - EXECUTION_DEADLOCK + EXECUTION_DEADLOCK, + PROJECT_DELETE_REQUIRES_CONFIRMATION, + PROJECT_CONTEXT_NAME_REQUIRED, + PROJECT_CONTEXT_NAME_INVALID, + PROJECT_CONTEXT_NAME_RESERVED, + PROJECT_CONTEXT_NAME_DUPLICATE } diff --git a/src/main/java/it/cnr/isti/workflow/manager/projects/ProjectContextConverter.java b/src/main/java/it/cnr/isti/workflow/manager/projects/ProjectContextConverter.java new file mode 100644 index 0000000..d32af27 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/projects/ProjectContextConverter.java @@ -0,0 +1,42 @@ +package it.cnr.isti.workflow.manager.projects; + +import com.fasterxml.jackson.annotation.JsonAutoDetect; + +import it.cnr.isti.workflow.manager.projects.model.ProjectContext; +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 ProjectContextConverter implements AttributeConverter { + + private static final ObjectMapper objectMapper = JsonMapper.builder() + .changeDefaultVisibility(vc -> vc.withFieldVisibility(JsonAutoDetect.Visibility.ANY)) + .build(); + + @Override + public String convertToDatabaseColumn(ProjectContext context) { + if (context == null) { + return null; + } + try { + return objectMapper.writeValueAsString(context); + } catch (JacksonException e) { + throw new IllegalArgumentException("Errore nella serializzazione di ProjectContext in JSON", e); + } + } + + @Override + public ProjectContext convertToEntityAttribute(String dbData) { + if (dbData == null || dbData.isBlank()) { + return null; + } + try { + return objectMapper.readValue(dbData, ProjectContext.class); + } catch (Exception e) { + throw new IllegalArgumentException("Errore nella deserializzazione di JSON in ProjectContext", e); + } + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/projects/ProjectContextResolver.java b/src/main/java/it/cnr/isti/workflow/manager/projects/ProjectContextResolver.java new file mode 100644 index 0000000..4fec9b3 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/projects/ProjectContextResolver.java @@ -0,0 +1,87 @@ +package it.cnr.isti.workflow.manager.projects; + +import java.util.Map; + +import org.springframework.stereotype.Component; +import org.springframework.transaction.annotation.Transactional; + +import it.cnr.isti.workflow.manager.flows.repo.FlowEntity; +import it.cnr.isti.workflow.manager.flows.repo.FlowRepository; +import it.cnr.isti.workflow.manager.projects.repo.ProjectEntity; +import it.cnr.isti.workflow.manager.projects.repo.ProjectRepository; + +/** + * Resolves the shared values a flow inherits from its project at execution time. + * + *

An execution is always created from a flow id, so the project is reached through the flow. + * This lives in its own component rather than inside {@code ExecutionsService} so every path that + * creates an execution - including container subflows and any future caller - gets it for free. + */ +@Component +public class ProjectContextResolver { + + private static final org.slf4j.Logger logger = org.slf4j.LoggerFactory.getLogger(ProjectContextResolver.class); + + public record ResolvedProjectContext(String projectId, Map values) { + + public static ResolvedProjectContext empty() { + return new ResolvedProjectContext(null, Map.of()); + } + + public boolean isEmpty() { + return values == null || values.isEmpty(); + } + } + + private final FlowRepository flowRepository; + private final ProjectRepository projectRepository; + + public ProjectContextResolver(FlowRepository flowRepository, ProjectRepository projectRepository) { + this.flowRepository = flowRepository; + this.projectRepository = projectRepository; + } + + /** + * The values to publish under the {@code project.} namespace for a run of {@code flowId} + * started by {@code executingOwner}. + * + *

The context is applied only when the executing user owns the project. A + * published flow can be run by anyone, and its project's values may be prompts, URLs or + * identifiers its owner considers private; they must never leak to a stranger's run. This also + * covers the anonymous and internal execution paths, where the owner is null. + */ + @Transactional(readOnly = true) + public ResolvedProjectContext resolveForFlow(String flowId, String executingOwner) { + if (flowId == null || flowId.isBlank() || executingOwner == null || executingOwner.isBlank()) { + return ResolvedProjectContext.empty(); + } + + FlowEntity flow = flowRepository.findById(flowId).orElse(null); + if (flow == null || flow.getProjectId() == null) { + return ResolvedProjectContext.empty(); + } + + ProjectEntity project = projectRepository.findById(flow.getProjectId()).orElse(null); + if (project == null) { + return ResolvedProjectContext.empty(); + } + if (!executingOwner.equals(project.getOwner())) { + logger.debug("User {} runs flow {} of project {} owned by {}: shared context not applied", + executingOwner, flowId, project.getId(), project.getOwner()); + return ResolvedProjectContext.empty(); + } + + return new ResolvedProjectContext( + project.getId(), + project.getSharedContext() == null ? Map.of() : project.getSharedContext().toValueMap()); + } + + /** The project id alone, without the values - used to tag an execution for later grouping. */ + @Transactional(readOnly = true) + public String projectIdOfFlow(String flowId) { + if (flowId == null || flowId.isBlank()) { + return null; + } + return flowRepository.findById(flowId).map(FlowEntity::getProjectId).orElse(null); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/projects/ProjectExecutionService.java b/src/main/java/it/cnr/isti/workflow/manager/projects/ProjectExecutionService.java new file mode 100644 index 0000000..588bc56 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/projects/ProjectExecutionService.java @@ -0,0 +1,191 @@ +package it.cnr.isti.workflow.manager.projects; + +import java.util.ArrayList; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +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.executions.ExecutionObject; +import it.cnr.isti.workflow.manager.executions.ExecutionsService; +import it.cnr.isti.workflow.manager.executions.api.ExecutionView; +import it.cnr.isti.workflow.manager.executions.repo.ExecutionEntity; +import it.cnr.isti.workflow.manager.executions.repo.ExecutionRepository; +import it.cnr.isti.workflow.manager.executions.ExecutionKind; +import it.cnr.isti.workflow.manager.flows.repo.FlowEntity; +import it.cnr.isti.workflow.manager.flows.repo.FlowRepository; +import it.cnr.isti.workflow.manager.flows.validation.FlowExecutionValidator; +import it.cnr.isti.workflow.manager.flows.validation.ValidationError; +import it.cnr.isti.workflow.manager.flows.validation.ValidationErrorCodec; +import it.cnr.isti.workflow.manager.projects.model.ProjectExecutionPlanView; +import it.cnr.isti.workflow.manager.projects.model.ProjectExecutionPlanView.SkippedFlow; +import it.cnr.isti.workflow.manager.projects.model.ProjectRunView; +import it.cnr.isti.workflow.manager.projects.repo.ProjectEntity; +import it.cnr.isti.workflow.manager.projects.repo.ProjectRepository; + +/** + * Runs a whole project: one execution per flow, tied together by a {@code projectRunId}. + * + *

It creates the executions, it does not start them. That is the existing + * contract, not timidity: {@code POST /executions} only creates, and the client then supplies + * global inputs and credentials before calling start. Any flow in a project may declare either, and + * there is no way to provide per-flow values through a single project-level call - so this delivers + * N executions, bound together, ready to be filled in and started. + */ +@Service +public class ProjectExecutionService { + + private static final org.slf4j.Logger logger = org.slf4j.LoggerFactory.getLogger(ProjectExecutionService.class); + + private final ProjectRepository projectRepository; + private final FlowRepository flowRepository; + private final ExecutionRepository executionRepository; + private final ExecutionsService executionsService; + private final FlowExecutionValidator flowExecutionValidator; + private final ProjectRunOrchestrator projectRunOrchestrator; + + public ProjectExecutionService(ProjectRepository projectRepository, FlowRepository flowRepository, + ExecutionRepository executionRepository, ExecutionsService executionsService, + FlowExecutionValidator flowExecutionValidator, ProjectRunOrchestrator projectRunOrchestrator) { + this.projectRepository = projectRepository; + this.flowRepository = flowRepository; + this.executionRepository = executionRepository; + this.executionsService = executionsService; + this.flowExecutionValidator = flowExecutionValidator; + this.projectRunOrchestrator = projectRunOrchestrator; + } + + /** + * Creates one execution per flow in the project. + * + *

Non-executable flows are refused up front, all-or-nothing: {@code createExecution} would + * otherwise throw partway through the loop and leave orphan executions behind, and a half-created + * run group is worse than a clear refusal. Pass {@code skipNonExecutable} to run the rest and + * have the skipped flows reported back instead. + */ + @Transactional + public ProjectExecutionPlanView execute(String owner, String projectId, String runName, + boolean skipNonExecutable) { + ProjectEntity project = requireOwnedProject(owner, projectId); + List flows = flowRepository.findByProjectIdOrderByProjectOrderAscNameAsc(projectId); + + if (flows.isEmpty()) { + throw new ResponseStatusException(HttpStatus.BAD_REQUEST, "Project contains no flows to run"); + } + + Map> notExecutable = new LinkedHashMap<>(); + for (FlowEntity flow : flows) { + if (!flowExecutionValidator.isExecutable(flow.getFlow())) { + notExecutable.put(flow, flowExecutionValidator.collectErrors(flow.getFlow())); + } + } + + if (!notExecutable.isEmpty() && !skipNonExecutable) { + throw new ResponseStatusException(HttpStatus.CONFLICT, encodeNotExecutable(notExecutable)); + } + + List runnable = flows.stream().filter(flow -> !notExecutable.containsKey(flow)).toList(); + if (runnable.isEmpty()) { + throw new ResponseStatusException(HttpStatus.CONFLICT, encodeNotExecutable(notExecutable)); + } + + String projectRunId = UUID.randomUUID().toString(); + String name = runName == null || runName.isBlank() ? project.getName() : runName.trim(); + + List created = new ArrayList<>(); + for (int order = 0; order < runnable.size(); order++) { + FlowEntity flow = runnable.get(order); + ExecutionObject execution = executionsService.createExecutionForFlow( + flow.getId(), name + " - " + flow.getName(), flow.getFlow(), owner, projectRunId, order); + created.add(ExecutionView.fromExecution(execution)); + } + logger.info("Project run {} for project {} created {} execution(s) for user {}", + projectRunId, projectId, created.size(), owner); + + List skipped = notExecutable.entrySet().stream() + .map(entry -> new SkippedFlow(entry.getKey().getId(), entry.getKey().getName(), + "Flow is not executable")) + .toList(); + + return new ProjectExecutionPlanView( + runView(projectId, projectRunId, name, + created.isEmpty() ? System.currentTimeMillis() : created.get(0).getCreationTime(), + created, owner), + skipped); + } + + /** Every run of a project, most recent first. */ + @Transactional(readOnly = true) + public List listRuns(String owner, String projectId) { + requireOwnedProject(owner, projectId); + + Map> byRun = new LinkedHashMap<>(); + for (ExecutionEntity entity : executionRepository + .findByProjectIdAndOwnerAndExecutionKindOrderByCreationTimeAsc(projectId, owner, + ExecutionKind.TOP_LEVEL)) { + if (entity.getProjectRunId() != null) { + byRun.computeIfAbsent(entity.getProjectRunId(), key -> new ArrayList<>()).add(entity); + } + } + + return byRun.entrySet().stream() + .map(entry -> toRunView(owner, projectId, entry.getKey(), entry.getValue())) + .sorted((left, right) -> Long.compare(right.createdAt(), left.createdAt())) + .toList(); + } + + private ProjectRunView toRunView(String owner, String projectId, String projectRunId, + List entities) { + List executions = entities.stream() + .map(entity -> ExecutionView.fromExecution(executionsService.getExecution(entity.getId()))) + .toList(); + ExecutionEntity first = entities.get(0); + return runView(projectId, projectRunId, first.getName(), first.getCreationTime(), executions, owner); + } + + private ProjectRunView runView(String projectId, String projectRunId, String name, long createdAt, + List executions, String owner) { + ProjectRunOrchestrator.ProjectRunProgress progress = projectRunOrchestrator.progress(owner, projectRunId); + return new ProjectRunView(projectRunId, projectId, name, createdAt, executions.size(), + progress.status(), progress.currentExecutionId(), progress.completed(), progress.blockedReason(), + executions); + } + + /** + * Starts the run, or resumes it after it stopped. The run advances one flow at a time in the + * project's order; each step starts only when the previous one succeeded. + */ + @Transactional + public ProjectRunView startRun(String owner, String projectId, String projectRunId) { + requireOwnedProject(owner, projectId); + projectRunOrchestrator.start(owner, projectRunId); + return listRuns(owner, projectId).stream() + .filter(run -> run.projectRunId().equals(projectRunId)) + .findFirst() + .orElseThrow(() -> new ResponseStatusException(HttpStatus.NOT_FOUND, "Project run not found")); + } + + private String encodeNotExecutable(Map> notExecutable) { + List errors = new ArrayList<>(); + notExecutable.forEach((flow, flowErrors) -> { + if (flowErrors.isEmpty()) { + errors.add(new ValidationError("flow", flow.getId(), null, + "Flow \"" + flow.getName() + "\" is not executable")); + return; + } + flowErrors.forEach(error -> errors.add(new ValidationError(error.code(), "flow", flow.getId(), + error.field(), "\"" + flow.getName() + "\": " + error.message()))); + }); + return ValidationErrorCodec.encode(errors); + } + + private ProjectEntity requireOwnedProject(String owner, String projectId) { + return projectRepository.findByIdAndOwner(projectId, owner) + .orElseThrow(() -> new ResponseStatusException(HttpStatus.NOT_FOUND, "Project not found")); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/projects/ProjectMapper.java b/src/main/java/it/cnr/isti/workflow/manager/projects/ProjectMapper.java new file mode 100644 index 0000000..073d642 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/projects/ProjectMapper.java @@ -0,0 +1,63 @@ +package it.cnr.isti.workflow.manager.projects; + +import java.util.List; + +import org.springframework.stereotype.Component; + +import it.cnr.isti.workflow.manager.flows.FlowMapper; +import it.cnr.isti.workflow.manager.flows.model.FlowSummaryView; +import it.cnr.isti.workflow.manager.flows.model.FlowViewStatus; +import it.cnr.isti.workflow.manager.flows.repo.FlowEntity; +import it.cnr.isti.workflow.manager.flows.validation.FlowExecutionValidator; +import it.cnr.isti.workflow.manager.projects.model.ProjectSummaryView; +import it.cnr.isti.workflow.manager.projects.model.ProjectView; +import it.cnr.isti.workflow.manager.projects.repo.ProjectEntity; + +/** + * A bean rather than a static helper like {@code FlowMapper}, because rendering a project's flows + * needs {@link FlowExecutionValidator} to derive each flow's status - which is computed at read + * time, never stored. + */ +@Component +public class ProjectMapper { + + private final FlowExecutionValidator flowExecutionValidator; + + public ProjectMapper(FlowExecutionValidator flowExecutionValidator) { + this.flowExecutionValidator = flowExecutionValidator; + } + + public static ProjectSummaryView toSummaryView(ProjectEntity project, long flowCount) { + return new ProjectSummaryView( + project.getId(), + project.getName(), + project.getDescription(), + project.getOwner(), + project.getCreatedAt(), + project.getLastUpdateAt(), + flowCount, + project.getSharedContext() == null ? 0 : project.getSharedContext().entryCount()); + } + + public ProjectView toView(ProjectEntity project, List flows) { + List flowViews = flows.stream() + // Reached only through the owner's own project, so membership is always visible here. + .map(flow -> FlowMapper.toSummaryView(flow, statusOf(flow), project.getId(), project.getName())) + .toList(); + return new ProjectView( + project.getId(), + project.getName(), + project.getDescription(), + project.getOwner(), + project.getCreatedAt(), + project.getLastUpdateAt(), + project.getSharedContext(), + flowViews); + } + + private FlowViewStatus statusOf(FlowEntity entity) { + return flowExecutionValidator.isExecutable(entity.getFlow()) + ? FlowViewStatus.EXECUTABLE + : FlowViewStatus.DRAFT; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/projects/ProjectRunOrchestrator.java b/src/main/java/it/cnr/isti/workflow/manager/projects/ProjectRunOrchestrator.java new file mode 100644 index 0000000..2bfcb1b --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/projects/ProjectRunOrchestrator.java @@ -0,0 +1,212 @@ +package it.cnr.isti.workflow.manager.projects; + +import java.util.ArrayList; +import java.util.Comparator; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; + +import org.springframework.http.HttpStatus; +import org.springframework.stereotype.Component; +import org.springframework.web.server.ResponseStatusException; + +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.repo.ExecutionEntity; +import it.cnr.isti.workflow.manager.executions.repo.ExecutionRepository; + +/** + * Drives a project run: one flow at a time, in the project's order. + * + *

The run has no state of its own. Its position is derived from the executions + * it created - they already carry {@code projectRunId} and {@code projectRunOrder} - so there is no + * second status machine that could drift out of sync with the executions, and a restart loses + * nothing. "Advance" simply means: find the first execution that has not finished, and start it if + * it is ready. + * + *

A failed or cancelled step stops the run. The remaining executions stay CREATED, so nothing is + * lost and the user can fix the problem and call start again to resume from where it stopped. + */ +@Component +public class ProjectRunOrchestrator { + + private static final org.slf4j.Logger logger = org.slf4j.LoggerFactory.getLogger(ProjectRunOrchestrator.class); + + /** The state of a project run, derived from its executions. */ + public enum ProjectRunStatus { + /** Nothing has been started yet. */ + PENDING, + /** One execution is running or waiting for an interaction. */ + RUNNING, + /** Stopped: the next execution needs inputs or credentials before it can start. */ + BLOCKED, + /** Stopped: an execution failed or was cancelled. */ + STOPPED, + /** Every execution finished successfully. */ + COMPLETED + } + + public record ProjectRunProgress(ProjectRunStatus status, String currentExecutionId, int completed, int total, + String blockedReason) { + } + + private final ExecutionRepository executionRepository; + private final ExecutionsService executionsService; + + /** + * One lock per run, so "find the next execution and start it" is atomic. Without it the + * completion callbacks of two runs - or a manual start racing a completion - could start the + * same step twice. + */ + private final Map runLocks = new ConcurrentHashMap<>(); + + public ProjectRunOrchestrator(ExecutionRepository executionRepository, ExecutionsService executionsService) { + this.executionRepository = executionRepository; + this.executionsService = executionsService; + } + + /** + * Starts the run, or resumes it after it stopped. Idempotent: if a step is already running this + * reports the current progress and starts nothing. + */ + public ProjectRunProgress start(String owner, String projectRunId) { + List ordered = requireRun(owner, projectRunId); + synchronized (lockFor(projectRunId)) { + return advanceLocked(owner, projectRunId, ordered, true); + } + } + + /** + * Called when an execution of a project run reaches a final state. Starts the next step, unless + * this one failed. + */ + public void onExecutionFinished(String owner, String projectRunId, ExecutionStatus finishedStatus) { + if (projectRunId == null || projectRunId.isBlank()) { + return; + } + if (finishedStatus != ExecutionStatus.SUCCESS) { + logger.info("Project run {} stops: an execution ended as {}", projectRunId, finishedStatus); + return; + } + synchronized (lockFor(projectRunId)) { + List ordered = orderedExecutions(owner, projectRunId); + if (!ordered.isEmpty()) { + advanceLocked(owner, projectRunId, ordered, false); + } + } + } + + public ProjectRunProgress progress(String owner, String projectRunId) { + return describe(orderedExecutions(owner, projectRunId)); + } + + private ProjectRunProgress advanceLocked(String owner, String projectRunId, List ordered, + boolean explicitStart) { + int completed = 0; + for (ExecutionEntity entity : ordered) { + ExecutionObject execution = executionsService.getExecution(entity.getId()); + ExecutionStatus status = execution.getContext().getStatus(); + + if (status == ExecutionStatus.SUCCESS) { + completed++; + continue; + } + if (status == ExecutionStatus.ERROR || status == ExecutionStatus.CANCELLED) { + return new ProjectRunProgress(ProjectRunStatus.STOPPED, execution.getId(), completed, ordered.size(), + "Execution " + execution.getId() + " ended as " + status); + } + if (!status.isInitState()) { + // Already running, waiting for an interaction, or suspended: nothing to start. + return new ProjectRunProgress(ProjectRunStatus.RUNNING, execution.getId(), completed, ordered.size(), + null); + } + + String blocked = notReadyReason(execution); + if (blocked != null) { + logger.debug("Project run {} blocked at execution {}: {}", projectRunId, execution.getId(), blocked); + return new ProjectRunProgress(ProjectRunStatus.BLOCKED, execution.getId(), completed, ordered.size(), + blocked); + } + + logger.info("Project run {} starts execution {} ({}/{})", + projectRunId, execution.getId(), completed + 1, ordered.size()); + executionsService.startExecution(execution.getId()); + return new ProjectRunProgress(ProjectRunStatus.RUNNING, execution.getId(), completed, ordered.size(), null); + } + + if (explicitStart) { + logger.debug("Project run {} has nothing left to start", projectRunId); + } + return new ProjectRunProgress(ProjectRunStatus.COMPLETED, null, completed, ordered.size(), null); + } + + /** Null when the execution can start; otherwise what it is still waiting for. */ + private String notReadyReason(ExecutionObject execution) { + List missingInputs = execution.getMissingGlobalInputKeys(); + List missingAuthorizations = execution.getMissingAuthorizationKeys(); + List reasons = new ArrayList<>(); + if (!missingInputs.isEmpty()) { + reasons.add("missing global inputs: " + String.join(", ", missingInputs)); + } + if (!missingAuthorizations.isEmpty()) { + reasons.add("missing credentials: " + String.join(", ", missingAuthorizations)); + } + if (!reasons.isEmpty()) { + return String.join("; ", reasons); + } + return execution.getContext().getStatus() == ExecutionStatus.READY + ? null + : "execution is not ready to start"; + } + + private ProjectRunProgress describe(List ordered) { + if (ordered.isEmpty()) { + return new ProjectRunProgress(ProjectRunStatus.PENDING, null, 0, 0, null); + } + int completed = 0; + for (ExecutionEntity entity : ordered) { + ExecutionObject execution = executionsService.getExecution(entity.getId()); + ExecutionStatus status = execution.getContext().getStatus(); + if (status == ExecutionStatus.SUCCESS) { + completed++; + continue; + } + if (status == ExecutionStatus.ERROR || status == ExecutionStatus.CANCELLED) { + return new ProjectRunProgress(ProjectRunStatus.STOPPED, execution.getId(), completed, ordered.size(), + "Execution " + execution.getId() + " ended as " + status); + } + if (!status.isInitState()) { + return new ProjectRunProgress(ProjectRunStatus.RUNNING, execution.getId(), completed, ordered.size(), + null); + } + String blocked = notReadyReason(execution); + return blocked == null + ? new ProjectRunProgress(ProjectRunStatus.PENDING, execution.getId(), completed, ordered.size(), + null) + : new ProjectRunProgress(ProjectRunStatus.BLOCKED, execution.getId(), completed, ordered.size(), + blocked); + } + return new ProjectRunProgress(ProjectRunStatus.COMPLETED, null, completed, ordered.size(), null); + } + + private List requireRun(String owner, String projectRunId) { + List ordered = orderedExecutions(owner, projectRunId); + if (ordered.isEmpty()) { + throw new ResponseStatusException(HttpStatus.NOT_FOUND, "Project run not found"); + } + return ordered; + } + + private List orderedExecutions(String owner, String projectRunId) { + return executionRepository.findByProjectRunIdAndOwnerOrderByCreationTimeAsc(projectRunId, owner).stream() + .sorted(Comparator.comparing(entity -> entity.getProjectRunOrder() == null + ? Integer.MAX_VALUE + : entity.getProjectRunOrder())) + .toList(); + } + + private Object lockFor(String projectRunId) { + return runLocks.computeIfAbsent(projectRunId, key -> new Object()); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/projects/ProjectService.java b/src/main/java/it/cnr/isti/workflow/manager/projects/ProjectService.java new file mode 100644 index 0000000..078b13a --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/projects/ProjectService.java @@ -0,0 +1,260 @@ +package it.cnr.isti.workflow.manager.projects; + +import java.time.LocalDateTime; +import java.util.ArrayList; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.regex.Pattern; +import java.util.stream.Collectors; + +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.flows.repo.FlowEntity; +import it.cnr.isti.workflow.manager.flows.repo.FlowRepository; +import it.cnr.isti.workflow.manager.flows.repo.ProjectFlowCount; +import it.cnr.isti.workflow.manager.flows.validation.ValidationError; +import it.cnr.isti.workflow.manager.flows.validation.ValidationErrorCode; +import it.cnr.isti.workflow.manager.flows.validation.ValidationErrorCodec; +import it.cnr.isti.workflow.manager.projects.model.ProjectContext; +import it.cnr.isti.workflow.manager.projects.model.ProjectContextEntry; +import it.cnr.isti.workflow.manager.projects.model.ProjectCreateRequest; +import it.cnr.isti.workflow.manager.projects.model.ProjectSummaryView; +import it.cnr.isti.workflow.manager.projects.model.ProjectUpdateRequest; +import it.cnr.isti.workflow.manager.projects.model.ProjectView; +import it.cnr.isti.workflow.manager.projects.repo.ProjectEntity; +import it.cnr.isti.workflow.manager.projects.repo.ProjectRepository; + +@Service +public class ProjectService { + + private static final org.slf4j.Logger logger = org.slf4j.LoggerFactory.getLogger(ProjectService.class); + + /** + * A context entry name is written by flow authors as {@code ${{project.}}}, so it must be + * a legal placeholder name - same shape the template machinery accepts. + */ + private static final Pattern VALID_ENTRY_NAME = Pattern.compile("^[A-Za-z][A-Za-z0-9_.-]*$"); + + /** + * The runtime publishes project values under the {@code project.} prefix alongside + * {@code global.}, {@code context.} and {@code vars.}. Namespacing means those four can never + * collide - unless an entry name smuggles a prefix in, which is refused here. + */ + private static final List RESERVED_PREFIXES = List.of("project.", "global.", "context.", "vars."); + + private final ProjectRepository projectRepository; + private final FlowRepository flowRepository; + private final ProjectMapper projectMapper; + + public ProjectService(ProjectRepository projectRepository, FlowRepository flowRepository, + ProjectMapper projectMapper) { + this.projectRepository = projectRepository; + this.flowRepository = flowRepository; + this.projectMapper = projectMapper; + } + + @Transactional(readOnly = true) + public List list(String owner) { + Map flowCounts = flowRepository.countFlowsByProjectForOwner(owner).stream() + .collect(Collectors.toMap(ProjectFlowCount::getProjectId, ProjectFlowCount::getFlowCount)); + return projectRepository.findByOwnerOrderByNameAsc(owner).stream() + .map(project -> ProjectMapper.toSummaryView(project, flowCounts.getOrDefault(project.getId(), 0L))) + .toList(); + } + + @Transactional(readOnly = true) + public ProjectView get(String owner, String projectId) { + return toView(requireOwnedProject(owner, projectId)); + } + + @Transactional + public ProjectView create(String owner, ProjectCreateRequest request) { + String name = requiredTrimmed(request.name(), "name"); + if (projectRepository.existsByOwnerAndName(owner, name)) { + throw new ResponseStatusException(HttpStatus.CONFLICT, "A project already uses this name"); + } + LocalDateTime now = LocalDateTime.now(); + ProjectEntity project = projectRepository.save(ProjectEntity.builder() + .owner(owner) + .name(name) + .description(trimToNull(request.description())) + .createdAt(now) + .lastUpdateAt(now) + .build()); + return toView(project); + } + + @Transactional + public ProjectView update(String owner, String projectId, ProjectUpdateRequest request) { + ProjectEntity project = requireOwnedProject(owner, projectId); + if (request.name() != null) { + String name = requiredTrimmed(request.name(), "name"); + if (projectRepository.existsByOwnerAndNameAndIdNot(owner, name, projectId)) { + throw new ResponseStatusException(HttpStatus.CONFLICT, "A project already uses this name"); + } + project.setName(name); + } + if (request.description() != null) { + project.setDescription(trimToNull(request.description())); + } + project.setLastUpdateAt(LocalDateTime.now()); + return toView(projectRepository.save(project)); + } + + @Transactional + public ProjectView updateContext(String owner, String projectId, ProjectContext context) { + ProjectEntity project = requireOwnedProject(owner, projectId); + validateContext(context); + project.setSharedContext(context); + project.setLastUpdateAt(LocalDateTime.now()); + return toView(projectRepository.save(project)); + } + + /** + * Sets the order a project run executes the project's flows in. + * + *

Ids that do not belong to the project are refused rather than ignored: silently dropping + * one would leave the caller believing an order it does not have. Flows omitted from the list + * keep their relative order after the listed ones. + */ + @Transactional + public ProjectView updateFlowOrder(String owner, String projectId, List flowIds) { + requireOwnedProject(owner, projectId); + List flows = flowRepository.findByProjectIdOrderByProjectOrderAscNameAsc(projectId); + Map byId = flows.stream() + .collect(Collectors.toMap(FlowEntity::getId, flow -> flow)); + + List requested = flowIds == null ? List.of() : flowIds; + Set unknown = requested.stream().filter(id -> !byId.containsKey(id)).collect(Collectors.toSet()); + if (!unknown.isEmpty()) { + throw new ResponseStatusException(HttpStatus.BAD_REQUEST, + "These flows do not belong to the project: " + String.join(", ", unknown)); + } + if (requested.size() != new HashSet<>(requested).size()) { + throw new ResponseStatusException(HttpStatus.BAD_REQUEST, "The flow order contains duplicates"); + } + + int position = 0; + Set placed = new HashSet<>(); + for (String flowId : requested) { + byId.get(flowId).setProjectOrder(position++); + placed.add(flowId); + } + for (FlowEntity flow : flows) { + if (!placed.contains(flow.getId())) { + flow.setProjectOrder(position++); + } + } + flowRepository.saveAll(flows); + return toView(requireOwnedProject(owner, projectId)); + } + + /** + * Deletes a project and every flow in it. + * + *

The cascade deliberately bypasses {@code FlowService.ensureMutable}, which makes a + * finalized flow undeletable, and so deletes finalized flows too. Refusing instead would be + * unrecoverable: finalizing is irreversible ({@code updateFinalized} rejects un-finalizing with + * 409), so a project holding one finalized flow could never be deleted by any API. The + * protection {@code finalized} provides is restored at the right granularity by requiring + * {@code confirm}. + * + *

Executions are left alone: each snapshots its own copy of the flow graph, so history stays + * renderable - exactly as it already does after a plain {@code DELETE /flows/{id}}. + */ + @Transactional + public void delete(String owner, String projectId, boolean confirm) { + ProjectEntity project = requireOwnedProject(owner, projectId); + List flows = flowRepository.findByProjectIdOrderByProjectOrderAscNameAsc(projectId); + + if (!flows.isEmpty() && !confirm) { + throw new ResponseStatusException(HttpStatus.CONFLICT, encodeConfirmationRequired(project, flows)); + } + if (!flows.isEmpty()) { + logger.info("Cascade-deleting project {} for user {} with {} flows ({} finalized)", + projectId, owner, flows.size(), flows.stream().filter(FlowEntity::isFinalized).count()); + flowRepository.deleteAll(flows); + } + projectRepository.delete(project); + } + + private String encodeConfirmationRequired(ProjectEntity project, List flows) { + long finalizedCount = flows.stream().filter(FlowEntity::isFinalized).count(); + List errors = new ArrayList<>(); + errors.add(new ValidationError(ValidationErrorCode.PROJECT_DELETE_REQUIRES_CONFIRMATION, "project", + project.getId(), null, + "Deleting project \"" + project.getName() + "\" will also delete " + flows.size() + + " flow(s), " + finalizedCount + " of them finalized. " + + "Repeat the request with confirm=true to proceed.")); + for (FlowEntity flow : flows) { + errors.add(new ValidationError(ValidationErrorCode.PROJECT_DELETE_REQUIRES_CONFIRMATION, "flow", + flow.getId(), flow.isFinalized() ? "finalized" : null, flow.getName())); + } + return ValidationErrorCodec.encode(errors); + } + + private void validateContext(ProjectContext context) { + if (context == null || context.getEntries() == null) { + return; + } + List errors = new ArrayList<>(); + Set seen = new HashSet<>(); + for (ProjectContextEntry entry : context.getEntries()) { + String name = entry == null ? null : trimToNull(entry.getName()); + if (name == null) { + errors.add(contextError(ValidationErrorCode.PROJECT_CONTEXT_NAME_REQUIRED, null, + "Shared context entry name is required")); + continue; + } + if (RESERVED_PREFIXES.stream().anyMatch(prefix -> name.toLowerCase().startsWith(prefix))) { + errors.add(contextError(ValidationErrorCode.PROJECT_CONTEXT_NAME_RESERVED, name, + "Shared context entry name must not start with a reserved prefix: " + name)); + } else if (!VALID_ENTRY_NAME.matcher(name).matches()) { + errors.add(contextError(ValidationErrorCode.PROJECT_CONTEXT_NAME_INVALID, name, + "Shared context entry name is not a valid placeholder name: " + name)); + } + if (!seen.add(name.toLowerCase())) { + errors.add(contextError(ValidationErrorCode.PROJECT_CONTEXT_NAME_DUPLICATE, name, + "Shared context entry name must be unique: " + name)); + } + } + if (!errors.isEmpty()) { + throw new ResponseStatusException(HttpStatus.BAD_REQUEST, ValidationErrorCodec.encode(errors)); + } + } + + private ValidationError contextError(ValidationErrorCode code, String name, String message) { + return new ValidationError(code, "projectContext", null, name, message); + } + + private ProjectView toView(ProjectEntity project) { + List flows = flowRepository.findByProjectIdOrderByProjectOrderAscNameAsc(project.getId()); + return projectMapper.toView(project, flows); + } + + private ProjectEntity requireOwnedProject(String owner, String projectId) { + return projectRepository.findByIdAndOwner(projectId, owner) + .orElseThrow(() -> new ResponseStatusException(HttpStatus.NOT_FOUND, "Project not found")); + } + + private String requiredTrimmed(String value, String field) { + String trimmed = trimToNull(value); + if (trimmed == null) { + throw new ResponseStatusException(HttpStatus.BAD_REQUEST, "Project " + field + " is required"); + } + return trimmed; + } + + private String trimToNull(String value) { + if (value == null) { + return null; + } + String trimmed = value.trim(); + return trimmed.isEmpty() ? null : trimmed; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectContext.java b/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectContext.java new file mode 100644 index 0000000..3ebe073 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectContext.java @@ -0,0 +1,51 @@ +package it.cnr.isti.workflow.manager.projects.model; + +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; + +import com.fasterxml.jackson.annotation.JsonIgnore; + +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.NoArgsConstructor; +import lombok.Singular; + +/** + * The values a project shares with its flows. + * + *

Wrapped in an object rather than persisted as a bare map because this shape is expected to + * grow (a project-level default LLM descriptor and vault credential references are planned): it is + * stored as a single TEXT JSON column via {@code ProjectContextConverter}, so additive fields cost + * no migration - the same reason {@code ExecutionSnapshot} gains fields freely. + */ +@Data +@Builder +@NoArgsConstructor(access = lombok.AccessLevel.PROTECTED) +@AllArgsConstructor +public class ProjectContext { + + @Singular + List entries; + + /** Flattens the entries to the {@code name -> value} map the execution runtime publishes. */ + @JsonIgnore + public Map toValueMap() { + Map values = new LinkedHashMap<>(); + if (entries == null) { + return values; + } + for (ProjectContextEntry entry : entries) { + if (entry != null && entry.getName() != null && !entry.getName().isBlank()) { + values.put(entry.getName(), entry.getValue()); + } + } + return values; + } + + @JsonIgnore + public int entryCount() { + return entries == null ? 0 : entries.size(); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectContextEntry.java b/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectContextEntry.java new file mode 100644 index 0000000..9173ffc --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectContextEntry.java @@ -0,0 +1,32 @@ +package it.cnr.isti.workflow.manager.projects.model; + +import it.cnr.isti.workflow.manager.ios.IOType; +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.NoArgsConstructor; + +/** + * One shared value a project hands to the flows it contains. + * + *

The {@code name} is what a flow author writes as {@code ${{project.}}}, so it is + * validated against the same rules a template placeholder must satisfy - see + * {@code ProjectService.validateContext}. {@link IOType} is deliberately the same enum global + * inputs use, so the editor can offer the identical value widgets for both. + */ +@Data +@Builder +@NoArgsConstructor(access = lombok.AccessLevel.PROTECTED) +@AllArgsConstructor +public class ProjectContextEntry { + + String name; + + IOType type; + + boolean multiple; + + Object value; + + String description; +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectCreateRequest.java b/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectCreateRequest.java new file mode 100644 index 0000000..043a39a --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectCreateRequest.java @@ -0,0 +1,9 @@ +package it.cnr.isti.workflow.manager.projects.model; + +import jakarta.validation.constraints.NotBlank; +import jakarta.validation.constraints.Size; + +public record ProjectCreateRequest( + @NotBlank String name, + @Size(max = 1000) String description) { +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectExecuteRequest.java b/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectExecuteRequest.java new file mode 100644 index 0000000..126a903 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectExecuteRequest.java @@ -0,0 +1,8 @@ +package it.cnr.isti.workflow.manager.projects.model; + +/** + * Optional name for the run. When absent the project name is used, so every execution of one run + * is recognisable in the executions list. + */ +public record ProjectExecuteRequest(String name) { +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectExecutionPlanView.java b/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectExecutionPlanView.java new file mode 100644 index 0000000..3368f16 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectExecutionPlanView.java @@ -0,0 +1,13 @@ +package it.cnr.isti.workflow.manager.projects.model; + +import java.util.List; + +/** What a project run created, and what it deliberately left out. */ +public record ProjectExecutionPlanView( + ProjectRunView run, + /** Flows not started because they are not executable, when skipNonExecutable was set. */ + List skipped) { + + public record SkippedFlow(String flowId, String flowName, String reason) { + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectFlowOrderRequest.java b/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectFlowOrderRequest.java new file mode 100644 index 0000000..351bd37 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectFlowOrderRequest.java @@ -0,0 +1,12 @@ +package it.cnr.isti.workflow.manager.projects.model; + +import java.util.List; + +import jakarta.validation.constraints.NotNull; + +/** + * The flows of a project in the order a project run should execute them. Ids not belonging to the + * project are refused, and any flow left out keeps its place after the listed ones. + */ +public record ProjectFlowOrderRequest(@NotNull List flowIds) { +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectRunView.java b/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectRunView.java new file mode 100644 index 0000000..8b20848 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectRunView.java @@ -0,0 +1,28 @@ +package it.cnr.isti.workflow.manager.projects.model; + +import java.util.List; + +import it.cnr.isti.workflow.manager.executions.api.ExecutionView; +import it.cnr.isti.workflow.manager.projects.ProjectRunOrchestrator.ProjectRunStatus; + +/** + * One run of a project: the executions started together, tied by {@code projectRunId}. + * + *

Deliberately separate from the execution groups of {@code GET /executions/groups}, + * which mean "the rerun history of a single flow" and assume every execution in them shares a + * source flow. + */ +public record ProjectRunView( + String projectRunId, + String projectId, + String name, + long createdAt, + int executionCount, + /** Derived from the executions themselves; the run keeps no separate state. */ + ProjectRunStatus status, + String currentExecutionId, + int completedCount, + /** What the run is waiting for when BLOCKED, otherwise null. */ + String blockedReason, + List executions) { +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectSummaryView.java b/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectSummaryView.java new file mode 100644 index 0000000..8f3ae2a --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectSummaryView.java @@ -0,0 +1,18 @@ +package it.cnr.isti.workflow.manager.projects.model; + +import java.time.LocalDateTime; + +/** + * Listing shape: carries no flows and no flow graphs. {@code flowCount} comes from a single + * group-by query, never from loading the flows. + */ +public record ProjectSummaryView( + String id, + String name, + String description, + String owner, + LocalDateTime createdAt, + LocalDateTime lastUpdateAt, + long flowCount, + int sharedContextEntryCount) { +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectUpdateRequest.java b/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectUpdateRequest.java new file mode 100644 index 0000000..5bfc38c --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectUpdateRequest.java @@ -0,0 +1,14 @@ +package it.cnr.isti.workflow.manager.projects.model; + +import jakarta.validation.constraints.Size; + +/** + * Partial update: a null field means "leave unchanged", matching {@code VaultSecretUpdateRequest}. + * + *

The shared context is deliberately absent - it has its own {@code PUT /projects/{id}/context}, + * so clearing it is always an explicit act and can never happen by omission. + */ +public record ProjectUpdateRequest( + String name, + @Size(max = 1000) String description) { +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectView.java b/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectView.java new file mode 100644 index 0000000..93796af --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/projects/model/ProjectView.java @@ -0,0 +1,17 @@ +package it.cnr.isti.workflow.manager.projects.model; + +import java.time.LocalDateTime; +import java.util.List; + +import it.cnr.isti.workflow.manager.flows.model.FlowSummaryView; + +public record ProjectView( + String id, + String name, + String description, + String owner, + LocalDateTime createdAt, + LocalDateTime lastUpdateAt, + ProjectContext sharedContext, + List flows) { +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/projects/repo/ProjectEntity.java b/src/main/java/it/cnr/isti/workflow/manager/projects/repo/ProjectEntity.java new file mode 100644 index 0000000..afbd0c8 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/projects/repo/ProjectEntity.java @@ -0,0 +1,51 @@ +package it.cnr.isti.workflow.manager.projects.repo; + +import java.time.LocalDateTime; + +import it.cnr.isti.workflow.manager.projects.ProjectContextConverter; +import it.cnr.isti.workflow.manager.projects.model.ProjectContext; +import jakarta.persistence.Column; +import jakarta.persistence.Convert; +import jakarta.persistence.Entity; +import jakarta.persistence.GeneratedValue; +import jakarta.persistence.GenerationType; +import jakarta.persistence.Id; +import jakarta.persistence.Table; +import jakarta.persistence.UniqueConstraint; +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.NoArgsConstructor; + +@Entity +@Table(name = "project", uniqueConstraints = @UniqueConstraint(name = "uk_project_owner_name", + columnNames = { "owner", "name" })) +@Data +@Builder +@NoArgsConstructor +@AllArgsConstructor +public class ProjectEntity { + + @Id + @GeneratedValue(strategy = GenerationType.UUID) + private String id; + + @Column(nullable = false) + private String name; + + @Column(length = 1000) + private String description; + + @Column(nullable = false) + private String owner; + + @Column(nullable = false) + private LocalDateTime createdAt; + + @Column(nullable = false) + private LocalDateTime lastUpdateAt; + + @Column(name = "shared_context", columnDefinition = "TEXT") + @Convert(converter = ProjectContextConverter.class) + private ProjectContext sharedContext; +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/projects/repo/ProjectRepository.java b/src/main/java/it/cnr/isti/workflow/manager/projects/repo/ProjectRepository.java new file mode 100644 index 0000000..02486c8 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/projects/repo/ProjectRepository.java @@ -0,0 +1,26 @@ +package it.cnr.isti.workflow.manager.projects.repo; + +import java.util.List; +import java.util.Optional; + +import org.springframework.data.jpa.repository.JpaRepository; +import org.springframework.stereotype.Repository; + +@Repository +public interface ProjectRepository extends JpaRepository { + + List findByOwnerOrderByNameAsc(String owner); + + /** + * Every read and write goes through this, so another user's project is indistinguishable from a + * missing one and surfaces as 404 - projects are private workspace structure, unlike flows, + * whose existence is not secret. + */ + Optional findByIdAndOwner(String id, String owner); + + boolean existsByOwnerAndName(String owner, String name); + + boolean existsByOwnerAndNameAndIdNot(String owner, String name, String id); + + void deleteByOwner(String owner); +} diff --git a/src/main/resources/db/migration/V6__create_project.sql b/src/main/resources/db/migration/V6__create_project.sql new file mode 100644 index 0000000..8d1e2a8 --- /dev/null +++ b/src/main/resources/db/migration/V6__create_project.sql @@ -0,0 +1,25 @@ +CREATE TABLE project ( + id VARCHAR(255) PRIMARY KEY, + name VARCHAR(255) NOT NULL, + description VARCHAR(1000), + owner VARCHAR(255) NOT NULL, + created_at TIMESTAMP(6) NOT NULL, + last_update_at TIMESTAMP(6) NOT NULL, + shared_context TEXT, + CONSTRAINT uk_project_owner_name UNIQUE (owner, name) +); + +CREATE INDEX idx_project_owner ON project (owner); + +-- No foreign key on project_id, on purpose. The tests run on H2 with ddl-auto=create-drop and +-- spring.flyway.enabled=false, and Hibernate's validate does not check foreign keys, so a +-- constraint added here would exist only in production and never be exercised anywhere. The +-- cascade is enforced in ProjectService, which deletes a project's flows before the project, and a +-- dangling project_id degrades safely: the project name simply resolves to null. + +-- project_order is unused for now: it is the stable display order a project run will need, and a +-- nullable column costs nothing today while sparing an ALTER on flow_entity later. +ALTER TABLE flow_entity ADD COLUMN project_id VARCHAR(255); +ALTER TABLE flow_entity ADD COLUMN project_order INTEGER; + +CREATE INDEX idx_flow_owner_project ON flow_entity (owner, project_id); diff --git a/src/main/resources/db/migration/V7__add_project_execution_columns.sql b/src/main/resources/db/migration/V7__add_project_execution_columns.sql new file mode 100644 index 0000000..cad3dcc --- /dev/null +++ b/src/main/resources/db/migration/V7__add_project_execution_columns.sql @@ -0,0 +1,11 @@ +-- Executions carry the project their source flow belonged to, plus the run that started them. +-- +-- project_run_id is a new column rather than a reuse of run_group_id: that one means "the rerun +-- history of one flow" (resolveHistoryGroupId keys off source_flow_id first, and run numbering +-- depends on it), so overloading it would both fail to group a project run and corrupt per-flow +-- rerun numbering. +ALTER TABLE execution_entity ADD COLUMN project_id VARCHAR(255); +ALTER TABLE execution_entity ADD COLUMN project_run_id VARCHAR(255); + +CREATE INDEX idx_execution_project_run ON execution_entity (owner, project_run_id); +CREATE INDEX idx_execution_owner_project ON execution_entity (owner, project_id); diff --git a/src/main/resources/db/migration/V8__add_project_run_order.sql b/src/main/resources/db/migration/V8__add_project_run_order.sql new file mode 100644 index 0000000..8bdf2ac --- /dev/null +++ b/src/main/resources/db/migration/V8__add_project_run_order.sql @@ -0,0 +1,4 @@ +-- Position of an execution inside its project run. A project run starts one flow at a time in this +-- order, so the sequence has to survive a restart: it is the run's only state, deliberately, rather +-- than a separate run table with a status machine to keep in sync with the executions themselves. +ALTER TABLE execution_entity ADD COLUMN project_run_order INTEGER; diff --git a/src/test/java/it/cnr/isti/workflow/manager/controllers/ProjectControllerTest.java b/src/test/java/it/cnr/isti/workflow/manager/controllers/ProjectControllerTest.java new file mode 100644 index 0000000..e9f114c --- /dev/null +++ b/src/test/java/it/cnr/isti/workflow/manager/controllers/ProjectControllerTest.java @@ -0,0 +1,298 @@ +package it.cnr.isti.workflow.manager.controllers; + +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.assertNull; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.util.List; + +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.boot.webmvc.test.autoconfigure.AutoConfigureMockMvc; +import org.springframework.http.HttpStatus; +import org.springframework.http.ResponseEntity; +import org.springframework.test.context.TestPropertySource; +import org.springframework.web.server.ResponseStatusException; + +import it.cnr.isti.workflow.manager.auth.repo.LoginEntity; +import it.cnr.isti.workflow.manager.flows.model.FlowCreateRequest; +import it.cnr.isti.workflow.manager.flows.model.FlowData; +import it.cnr.isti.workflow.manager.flows.model.FlowFlagUpdateRequest; +import it.cnr.isti.workflow.manager.flows.model.FlowProjectUpdateRequest; +import it.cnr.isti.workflow.manager.flows.model.FlowSummaryView; +import it.cnr.isti.workflow.manager.flows.model.FlowView; +import it.cnr.isti.workflow.manager.flows.repo.FlowRepository; +import it.cnr.isti.workflow.manager.projects.model.ProjectContext; +import it.cnr.isti.workflow.manager.projects.model.ProjectContextEntry; +import it.cnr.isti.workflow.manager.projects.model.ProjectCreateRequest; +import it.cnr.isti.workflow.manager.projects.model.ProjectSummaryView; +import it.cnr.isti.workflow.manager.projects.model.ProjectUpdateRequest; +import it.cnr.isti.workflow.manager.projects.model.ProjectView; +import it.cnr.isti.workflow.manager.projects.repo.ProjectRepository; +import it.cnr.isti.workflow.manager.ios.IOType; + +@SpringBootTest +@AutoConfigureMockMvc +@TestPropertySource(locations = "classpath:test.properties") +public class ProjectControllerTest { + + @Autowired + private ProjectController projectController; + + @Autowired + private FlowController flowController; + + @Autowired + private FlowRepository flowRepository; + + @Autowired + private ProjectRepository projectRepository; + + private static LoginEntity user(String username) { + return new LoginEntity(username, "password"); + } + + private ProjectView createProject(String owner, String name) { + return projectController.create(new ProjectCreateRequest(name, "desc"), user(owner)); + } + + /** An empty flow is intentionally fine: FlowService saves non-executable flows as DRAFT. */ + private FlowView createFlow(String owner, String name) { + ResponseEntity response = flowController.createFlow( + new FlowCreateRequest(name, "d", FlowData.builder().build()), user(owner)); + return response.getBody(); + } + + private FlowView assign(String owner, String flowId, String projectId) { + return flowController.updateProject(flowId, new FlowProjectUpdateRequest(projectId), user(owner)).getBody(); + } + + @Test + public void createUpdateAndListProject() { + ProjectView created = createProject("proj-user-1", "Recruiting"); + assertNotNull(created.id()); + assertEquals("Recruiting", created.name()); + assertEquals("proj-user-1", created.owner()); + assertTrue(created.flows().isEmpty()); + + ProjectView renamed = projectController.update(created.id(), + new ProjectUpdateRequest("Hiring", null), user("proj-user-1")); + assertEquals("Hiring", renamed.name()); + // null means "leave unchanged", so the description survives a rename. + assertEquals("desc", renamed.description()); + + List listed = projectController.list(user("proj-user-1")); + assertEquals(1, listed.size()); + assertEquals(0, listed.get(0).flowCount()); + } + + @Test + public void duplicateNamePerOwnerIsRejectedButNotAcrossOwners() { + createProject("proj-user-2", "Shared name"); + + ResponseStatusException conflict = assertThrows(ResponseStatusException.class, + () -> createProject("proj-user-2", "Shared name")); + assertEquals(HttpStatus.CONFLICT, conflict.getStatusCode()); + + // The uniqueness is per owner, so another user may reuse the name. + assertNotNull(createProject("proj-user-3", "Shared name").id()); + } + + @Test + public void anotherUsersProjectIsReportedAsNotFound() { + ProjectView owned = createProject("proj-owner", "Private"); + + // 404 rather than 403 on purpose: a project is private workspace structure, so its + // existence must not leak - unlike a flow, which lives in a publishable space. + for (Runnable call : List.of( + () -> projectController.get(owned.id(), user("proj-intruder")), + () -> projectController.update(owned.id(), new ProjectUpdateRequest("x", null), user("proj-intruder")), + () -> projectController.delete(owned.id(), true, user("proj-intruder")))) { + ResponseStatusException e = assertThrows(ResponseStatusException.class, call::run); + assertEquals(HttpStatus.NOT_FOUND, e.getStatusCode()); + } + } + + @Test + public void assignMoveAndDetachFlow() { + ProjectView first = createProject("assign-user", "First"); + ProjectView second = createProject("assign-user", "Second"); + FlowView flow = createFlow("assign-user", "Assignable flow"); + assertNull(flow.projectId()); + + FlowView assigned = assign("assign-user", flow.id(), first.id()); + assertEquals(first.id(), assigned.projectId()); + assertEquals("First", assigned.projectName()); + + FlowView moved = assign("assign-user", flow.id(), second.id()); + assertEquals(second.id(), moved.projectId()); + assertEquals("Second", moved.projectName()); + + FlowView detached = assign("assign-user", flow.id(), null); + assertNull(detached.projectId()); + assertNull(detached.projectName()); + } + + @Test + public void assigningToAnotherUsersProjectIsNotFound() { + ProjectView foreign = createProject("other-owner", "Foreign"); + FlowView flow = createFlow("assign-user-2", "My flow"); + + ResponseStatusException e = assertThrows(ResponseStatusException.class, + () -> assign("assign-user-2", flow.id(), foreign.id())); + assertEquals(HttpStatus.NOT_FOUND, e.getStatusCode()); + } + + @Test + public void assigningAnotherUsersFlowIsForbidden() { + ProjectView project = createProject("assign-user-3", "Mine"); + FlowView foreignFlow = createFlow("flow-owner", "Not mine"); + + ResponseEntity response = flowController.updateProject(foreignFlow.id(), + new FlowProjectUpdateRequest(project.id()), user("assign-user-3")); + assertEquals(HttpStatus.FORBIDDEN, response.getStatusCode()); + } + + @Test + public void finalizedFlowCanStillBeReassigned() { + ProjectView project = createProject("final-user", "Archive"); + FlowView flow = createFlow("final-user", "To finalize"); + flowController.updateFinalized(flow.id(), new FlowFlagUpdateRequest(true), user("final-user")); + + // Finalizing is irreversible, so freezing a flow out of reorganization would be permanent. + // Project membership is workspace metadata, not flow content. + FlowView assigned = assign("final-user", flow.id(), project.id()); + assertEquals(project.id(), assigned.projectId()); + } + + @Test + public void deleteWithoutConfirmIsRefusedAndListsTheFlows() { + ProjectView project = createProject("cascade-user", "Doomed"); + FlowView keep = createFlow("cascade-user", "Flow A"); + assign("cascade-user", keep.id(), project.id()); + + ResponseStatusException conflict = assertThrows(ResponseStatusException.class, + () -> projectController.delete(project.id(), false, user("cascade-user"))); + assertEquals(HttpStatus.CONFLICT, conflict.getStatusCode()); + assertTrue(conflict.getReason().contains(keep.id())); + + // Nothing was destroyed by the refused call. + assertTrue(flowRepository.findById(keep.id()).isPresent()); + assertTrue(projectRepository.findById(project.id()).isPresent()); + } + + @Test + public void confirmedDeleteCascadesIncludingFinalizedFlows() { + ProjectView project = createProject("cascade-user-2", "Doomed too"); + FlowView plain = createFlow("cascade-user-2", "Plain"); + FlowView finalized = createFlow("cascade-user-2", "Finalized"); + assign("cascade-user-2", plain.id(), project.id()); + assign("cascade-user-2", finalized.id(), project.id()); + flowController.updateFinalized(finalized.id(), new FlowFlagUpdateRequest(true), user("cascade-user-2")); + + projectController.delete(project.id(), true, user("cascade-user-2")); + + // The cascade deliberately bypasses the finalized-is-undeletable rule: refusing would be + // unrecoverable, since a flow can never be un-finalized. + assertFalse(flowRepository.findById(plain.id()).isPresent()); + assertFalse(flowRepository.findById(finalized.id()).isPresent()); + assertFalse(projectRepository.findById(project.id()).isPresent()); + } + + @Test + public void emptyProjectDeletesWithoutConfirmation() { + ProjectView project = createProject("cascade-user-3", "Empty"); + projectController.delete(project.id(), false, user("cascade-user-3")); + assertFalse(projectRepository.findById(project.id()).isPresent()); + } + + @Test + public void listingCarriesFlowCountAndNoFlowGraphs() { + ProjectView project = createProject("count-user", "Counted"); + for (int i = 0; i < 3; i++) { + assign("count-user", createFlow("count-user", "Flow " + i).id(), project.id()); + } + + ProjectSummaryView summary = projectController.list(user("count-user")).stream() + .filter(p -> p.id().equals(project.id())) + .findFirst() + .orElseThrow(); + assertEquals(3L, summary.flowCount()); + + // The detail view exposes flows as summaries, which carry no FlowData by construction. + ProjectView detail = projectController.get(project.id(), user("count-user")); + assertEquals(3, detail.flows().size()); + for (FlowSummaryView flow : detail.flows()) { + assertEquals("Counted", flow.projectName()); + } + } + + @Test + public void publishedFlowHidesItsProjectFromOtherUsers() { + ProjectView project = createProject("privacy-owner", "Client X confidential"); + FlowView flow = createFlow("privacy-owner", "Published flow"); + assign("privacy-owner", flow.id(), project.id()); + flowController.updatePublished(flow.id(), new FlowFlagUpdateRequest(true), user("privacy-owner")); + + // The flow stays readable by everyone - project membership is organizational structure, + // not access control - but the project name can itself be private information, so a + // non-owner must see the flow as unassigned. + FlowView asStranger = flowController.getFlow(flow.id(), user("privacy-stranger")).getBody(); + assertNotNull(asStranger); + assertNull(asStranger.projectId()); + assertNull(asStranger.projectName()); + + FlowView asOwner = flowController.getFlow(flow.id(), user("privacy-owner")).getBody(); + assertEquals(project.id(), asOwner.projectId()); + assertEquals("Client X confidential", asOwner.projectName()); + + // The same must hold for the listings, which resolve project names in bulk. + FlowView listedForStranger = flowController.getAllFlows(user("privacy-stranger")).stream() + .filter(f -> f.id().equals(flow.id())) + .findFirst() + .orElseThrow(); + assertNull(listedForStranger.projectId()); + + FlowSummaryView summaryForStranger = flowController.getAllFlowSummaries(user("privacy-stranger")).stream() + .filter(f -> f.id().equals(flow.id())) + .findFirst() + .orElseThrow(); + assertNull(summaryForStranger.projectId()); + assertNull(summaryForStranger.projectName()); + } + + @Test + public void sharedContextRoundTripsAndRejectsReservedNames() { + ProjectView project = createProject("ctx-user", "Context"); + + ProjectView saved = projectController.updateContext(project.id(), + ProjectContext.builder() + .entry(ProjectContextEntry.builder().name("tone").type(IOType.TEXT).value("formal").build()) + .build(), + user("ctx-user")); + assertEquals(1, saved.sharedContext().entryCount()); + assertEquals("formal", saved.sharedContext().toValueMap().get("tone")); + + // A name carrying a reserved prefix is the only way to manufacture a runtime key collision. + ResponseStatusException reserved = assertThrows(ResponseStatusException.class, + () -> projectController.updateContext(project.id(), + ProjectContext.builder() + .entry(ProjectContextEntry.builder().name("global.tone").type(IOType.TEXT).build()) + .build(), + user("ctx-user"))); + assertEquals(HttpStatus.BAD_REQUEST, reserved.getStatusCode()); + + ResponseStatusException duplicate = assertThrows(ResponseStatusException.class, + () -> projectController.updateContext(project.id(), + ProjectContext.builder() + .entry(ProjectContextEntry.builder().name("tone").type(IOType.TEXT).build()) + .entry(ProjectContextEntry.builder().name("Tone").type(IOType.TEXT).build()) + .build(), + user("ctx-user"))); + assertEquals(HttpStatus.BAD_REQUEST, duplicate.getStatusCode()); + } +} diff --git a/src/test/java/it/cnr/isti/workflow/manager/controllers/StatsControllerTest.java b/src/test/java/it/cnr/isti/workflow/manager/controllers/StatsControllerTest.java index 92d1494..860b712 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/controllers/StatsControllerTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/controllers/StatsControllerTest.java @@ -8,8 +8,11 @@ import java.time.LocalDateTime; import java.util.List; import java.util.Map; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.TestInstance; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.test.context.SpringBootTest; import org.springframework.http.ResponseEntity; @@ -30,6 +33,7 @@ import it.cnr.isti.workflow.manager.stats.model.UserUsageStatsView; @SpringBootTest @TestPropertySource(locations = "classpath:test.properties") +@TestInstance(TestInstance.Lifecycle.PER_CLASS) public class StatsControllerTest { @Autowired @@ -44,6 +48,24 @@ public class StatsControllerTest { @Autowired private ExecutionRepository executionRepository; + private List preservedUsers = List.of(); + + /** + * These tests need a known dataset, so they wipe the users table - but the Spring context, and + * with it the H2 database, is shared by the whole suite. Losing the seeded "testuser" makes + * every later class that authenticates with a JWT for it fail with 401, depending purely on + * class ordering. Capture the shared users up front and put them back afterwards. + */ + @BeforeAll + void captureSharedUsers() { + preservedUsers = List.copyOf(authRepository.findAll()); + } + + @AfterAll + void restoreSharedUsers() { + authRepository.saveAll(preservedUsers); + } + @BeforeEach void cleanRepositories() { executionRepository.deleteAll(); diff --git a/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionTemplateResolverTest.java b/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionTemplateResolverTest.java index 4111519..d894992 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionTemplateResolverTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionTemplateResolverTest.java @@ -153,4 +153,39 @@ public class ExecutionTemplateResolverTest { null); assertEquals("Dossier: {\"hiringNeed\":\"Backend engineer\"}", result); } + + @Test + public void resolvesProjectContextWithProjectPrefix() { + // Guards the one-line prefix check in ExecutionTemplateResolver: without it a project. + // key is published as ${{vars.project.x}} and ${{project.x}} silently never resolves. + Map runtime = new LinkedHashMap<>(); + runtime.put("project.tone", "formal"); + + assertEquals("Write in a formal tone", + ExecutionTemplateResolver.resolve("Write in a ${{project.tone}} tone", Map.of(), runtime)); + } + + @Test + public void doesNotExposeProjectContextUnderTheVarsPrefix() { + Map runtime = new LinkedHashMap<>(); + runtime.put("project.tone", "formal"); + + assertEquals("still ${{vars.project.tone}}", + ExecutionTemplateResolver.resolve("still ${{vars.project.tone}}", Map.of(), runtime)); + } + + @Test + public void keepsProjectGlobalAndContextNamespacesDistinct() { + // Namespacing is why there is no precedence question to settle between the sources. + Map runtime = new LinkedHashMap<>(); + runtime.put("project.tone", "from-project"); + runtime.put("global.tone", "from-global"); + runtime.put("context.tone", "from-context"); + runtime.put("tone", "bare"); + + assertEquals("from-project from-global from-context bare", + ExecutionTemplateResolver.resolve( + "${{project.tone}} ${{global.tone}} ${{context.tone}} ${{vars.tone}}", + Map.of(), runtime)); + } } diff --git a/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionTest.java b/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionTest.java index bd003a7..d439087 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/executions/ExecutionTest.java @@ -950,7 +950,8 @@ public class ExecutionTest { execution.getContext().getSteps().get(rejected.getId()).getInputs().getFirst().getResolutionState()); assertEquals(InputResolutionState.NOT_SELECTED, execution.getContext().getSteps().get(rejectedDownstream.getId()).getInputs().getFirst().getResolutionState()); - assertEquals(3, execution.snapshot().getSnapshotVersion()); + // Version 4 added the project shared context the run inherited. + assertEquals(4, execution.snapshot().getSnapshotVersion()); } @Test diff --git a/src/test/java/it/cnr/isti/workflow/manager/projects/ProjectContextExecutionTest.java b/src/test/java/it/cnr/isti/workflow/manager/projects/ProjectContextExecutionTest.java new file mode 100644 index 0000000..3caf66e --- /dev/null +++ b/src/test/java/it/cnr/isti/workflow/manager/projects/ProjectContextExecutionTest.java @@ -0,0 +1,178 @@ +package it.cnr.isti.workflow.manager.projects; + +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.UUID; + +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.test.context.TestPropertySource; + +import it.cnr.isti.workflow.manager.auth.repo.LoginEntity; +import it.cnr.isti.workflow.manager.blocks.Block; +import it.cnr.isti.workflow.manager.blocks.configurations.LLMBlockConfiguration; +import it.cnr.isti.workflow.manager.blocks.types.LLMBlockType; +import it.cnr.isti.workflow.manager.controllers.BlocksController; +import it.cnr.isti.workflow.manager.controllers.FlowController; +import it.cnr.isti.workflow.manager.controllers.ProjectController; +import it.cnr.isti.workflow.manager.executions.ExecutionObject; +import it.cnr.isti.workflow.manager.executions.ExecutionsService; +import it.cnr.isti.workflow.manager.flows.model.FlowCreateRequest; +import it.cnr.isti.workflow.manager.flows.model.FlowData; +import it.cnr.isti.workflow.manager.flows.model.FlowProjectUpdateRequest; +import it.cnr.isti.workflow.manager.flows.model.FlowView; +import it.cnr.isti.workflow.manager.flows.model.FlowViewStatus; +import it.cnr.isti.workflow.manager.ios.IOType; +import it.cnr.isti.workflow.manager.llms.LLMDescriptor; +import it.cnr.isti.workflow.manager.projects.model.ProjectContext; +import it.cnr.isti.workflow.manager.projects.model.ProjectContextEntry; +import it.cnr.isti.workflow.manager.projects.model.ProjectCreateRequest; +import it.cnr.isti.workflow.manager.projects.model.ProjectView; + +/** End-to-end cover for the project shared context reaching a running flow. */ +@SpringBootTest +@TestPropertySource(locations = "classpath:test.properties") +public class ProjectContextExecutionTest { + + private static final String OWNER = "ctx-exec-user"; + + @Autowired + private ProjectController projectController; + + @Autowired + private FlowController flowController; + + @Autowired + private BlocksController blocksController; + + @Autowired + private ExecutionsService executionsService; + + private static LoginEntity user(String username) { + return new LoginEntity(username, "password"); + } + + private ProjectView projectWithTone(String owner, String tone) { + ProjectView project = projectController.create( + new ProjectCreateRequest("Ctx " + UUID.randomUUID(), null), user(owner)); + return projectController.updateContext(project.id(), + ProjectContext.builder() + .entry(ProjectContextEntry.builder() + .name("tone").type(IOType.TEXT).value(tone).build()) + .build(), + user(owner)); + } + + /** A one-block flow whose prompt reads the project value. */ + private FlowView flowUsingProjectTone(String owner, String projectId) { + Block block = blocksController.create(LLMBlockConfiguration.builder() + .name("Writer") + .llmDescriptor(LLMDescriptor.builder().provider("testProvider").model("testModel").build()) + .prompt("Write in a ${{project.tone}} tone") + .build()); + + FlowView flow = flowController.createFlow( + new FlowCreateRequest("Project context flow " + UUID.randomUUID(), null, + FlowData.builder().block(block).build()), + user(owner)).getBody(); + return flowController.updateProject(flow.id(), new FlowProjectUpdateRequest(projectId), user(owner)).getBody(); + } + + @Test + public void projectPlaceholderDoesNotBecomeABlockInput() { + // Guards the PlaceholderInputs whitelist: without project. in it, ${{project.tone}} would be + // harvested into a dangling block input and the flow would never be executable. + Block block = blocksController.create(LLMBlockConfiguration.builder() + .name("Writer") + .llmDescriptor(LLMDescriptor.builder().provider("testProvider").model("testModel").build()) + .prompt("Write in a ${{project.tone}} tone") + .build()); + + assertTrue(block.getInputs().stream().noneMatch(input -> input.getName().startsWith("project")), + "a project placeholder must not be turned into a block input"); + assertTrue(block.getInputs().isEmpty()); + } + + @Test + public void aFlowUsingOnlyProjectPlaceholdersIsExecutable() { + ProjectView project = projectWithTone(OWNER, "formal"); + FlowView flow = flowUsingProjectTone(OWNER, project.id()); + + assertEquals(FlowViewStatus.EXECUTABLE, flow.status()); + } + + @Test + public void executionInheritsTheProjectContext() { + ProjectView project = projectWithTone(OWNER, "formal"); + FlowView flow = flowUsingProjectTone(OWNER, project.id()); + + ExecutionObject execution = executionsService.createExecutionForFlow( + flow.id(), flow.name(), flow.flow(), OWNER); + + assertEquals(project.id(), execution.getProjectId()); + assertEquals("formal", execution.getContext().getProjectContext().get("tone")); + } + + @Test + public void aStrangerRunningAPublishedFlowInheritsNothing() { + ProjectView project = projectWithTone(OWNER, "formal"); + FlowView flow = flowUsingProjectTone(OWNER, project.id()); + + ExecutionObject execution = executionsService.createExecutionForFlow( + flow.id(), flow.name(), flow.flow(), "ctx-exec-stranger"); + + assertTrue(execution.getContext().getProjectContext().isEmpty()); + // The execution is still tagged with the project, so grouping keeps working. + assertEquals(project.id(), execution.getProjectId()); + } + + @Test + public void editingTheProjectAfterwardsDoesNotRewriteAnExistingExecution() { + ProjectView project = projectWithTone(OWNER, "formal"); + FlowView flow = flowUsingProjectTone(OWNER, project.id()); + + ExecutionObject execution = executionsService.createExecutionForFlow( + flow.id(), flow.name(), flow.flow(), OWNER); + String executionId = execution.getId(); + + // Change the project after the run exists. + projectController.updateContext(project.id(), + ProjectContext.builder() + .entry(ProjectContextEntry.builder() + .name("tone").type(IOType.TEXT).value("casual").build()) + .build(), + user(OWNER)); + + // Force the execution to be rebuilt from its persisted snapshot. + executionsService.clearInMemoryExecutions(); + ExecutionObject reloaded = executionsService.getExecution(executionId); + + assertNotNull(reloaded); + assertEquals("formal", reloaded.getContext().getProjectContext().get("tone"), + "a run must keep the context it was created with"); + assertFalse(reloaded.getContext().getProjectContext().containsValue("casual")); + } + + @Test + public void aFlowWithNoProjectInheritsNothing() { + Block block = blocksController.create(LLMBlockConfiguration.builder() + .name("Plain") + .llmDescriptor(LLMDescriptor.builder().provider("testProvider").model("testModel").build()) + .prompt("Just write something") + .build()); + FlowView flow = flowController.createFlow( + new FlowCreateRequest("No project flow " + UUID.randomUUID(), null, + FlowData.builder().block(block).build()), + user(OWNER)).getBody(); + + ExecutionObject execution = executionsService.createExecutionForFlow( + flow.id(), flow.name(), flow.flow(), OWNER); + + assertTrue(execution.getContext().getProjectContext().isEmpty()); + assertEquals(null, execution.getProjectId()); + } +} diff --git a/src/test/java/it/cnr/isti/workflow/manager/projects/ProjectContextResolverTest.java b/src/test/java/it/cnr/isti/workflow/manager/projects/ProjectContextResolverTest.java new file mode 100644 index 0000000..43193ce --- /dev/null +++ b/src/test/java/it/cnr/isti/workflow/manager/projects/ProjectContextResolverTest.java @@ -0,0 +1,137 @@ +package it.cnr.isti.workflow.manager.projects; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.time.LocalDateTime; +import java.util.UUID; + +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.test.context.TestPropertySource; + +import it.cnr.isti.workflow.manager.flows.model.FlowData; +import it.cnr.isti.workflow.manager.flows.repo.FlowEntity; +import it.cnr.isti.workflow.manager.flows.repo.FlowRepository; +import it.cnr.isti.workflow.manager.ios.IOType; +import it.cnr.isti.workflow.manager.projects.ProjectContextResolver.ResolvedProjectContext; +import it.cnr.isti.workflow.manager.projects.model.ProjectContext; +import it.cnr.isti.workflow.manager.projects.model.ProjectContextEntry; +import it.cnr.isti.workflow.manager.projects.repo.ProjectEntity; +import it.cnr.isti.workflow.manager.projects.repo.ProjectRepository; + +@SpringBootTest +@TestPropertySource(locations = "classpath:test.properties") +public class ProjectContextResolverTest { + + @Autowired + private ProjectContextResolver resolver; + + @Autowired + private FlowRepository flowRepository; + + @Autowired + private ProjectRepository projectRepository; + + private ProjectEntity saveProject(String owner, String entryName, Object value) { + LocalDateTime now = LocalDateTime.now(); + return projectRepository.save(ProjectEntity.builder() + .owner(owner) + .name("Ctx " + UUID.randomUUID()) + .createdAt(now) + .lastUpdateAt(now) + .sharedContext(entryName == null + ? null + : ProjectContext.builder() + .entry(ProjectContextEntry.builder() + .name(entryName) + .type(IOType.TEXT) + .value(value) + .build()) + .build()) + .build()); + } + + private FlowEntity saveFlow(String owner, String projectId) { + LocalDateTime now = LocalDateTime.now(); + return flowRepository.save(FlowEntity.builder() + .name("Flow " + UUID.randomUUID()) + .owner(owner) + .createdAt(now) + .lastUpdateAt(now) + .projectId(projectId) + .flow(FlowData.builder().build()) + .build()); + } + + @Test + public void resolvesTheProjectValuesForTheOwner() { + ProjectEntity project = saveProject("ctx-owner", "tone", "formal"); + FlowEntity flow = saveFlow("ctx-owner", project.getId()); + + ResolvedProjectContext resolved = resolver.resolveForFlow(flow.getId(), "ctx-owner"); + + assertEquals(project.getId(), resolved.projectId()); + assertEquals("formal", resolved.values().get("tone")); + } + + @Test + public void resolvesEmptyForAFlowWithNoProject() { + FlowEntity flow = saveFlow("ctx-owner", null); + + ResolvedProjectContext resolved = resolver.resolveForFlow(flow.getId(), "ctx-owner"); + assertTrue(resolved.isEmpty()); + assertNull(resolved.projectId()); + } + + @Test + public void resolvesEmptyWhenTheProjectIsGone() { + FlowEntity flow = saveFlow("ctx-owner", "a-project-that-no-longer-exists"); + + assertTrue(resolver.resolveForFlow(flow.getId(), "ctx-owner").isEmpty()); + } + + @Test + public void doesNotApplyTheContextToSomeoneElsesRun() { + // The security decision: a published flow can be run by anyone, and the project's values + // may be prompts, URLs or identifiers its owner considers private. + ProjectEntity project = saveProject("ctx-owner", "secretUrl", "https://internal"); + FlowEntity flow = saveFlow("ctx-owner", project.getId()); + + ResolvedProjectContext resolved = resolver.resolveForFlow(flow.getId(), "ctx-stranger"); + + assertTrue(resolved.isEmpty()); + assertNull(resolved.projectId()); + } + + @Test + public void resolvesEmptyForAnAnonymousRun() { + ProjectEntity project = saveProject("ctx-owner", "tone", "formal"); + FlowEntity flow = saveFlow("ctx-owner", project.getId()); + + assertTrue(resolver.resolveForFlow(flow.getId(), null).isEmpty()); + } + + @Test + public void keepsTheProjectIdWhenTheContextIsUnset() { + ProjectEntity project = saveProject("ctx-owner", null, null); + FlowEntity flow = saveFlow("ctx-owner", project.getId()); + + ResolvedProjectContext resolved = resolver.resolveForFlow(flow.getId(), "ctx-owner"); + + assertEquals(project.getId(), resolved.projectId()); + assertTrue(resolved.isEmpty()); + } + + @Test + public void reportsTheProjectIdOfAFlowRegardlessOfOwnership() { + // Used to tag an execution for grouping, so it must not depend on who is running it. + ProjectEntity project = saveProject("ctx-owner", "tone", "formal"); + FlowEntity flow = saveFlow("ctx-owner", project.getId()); + + assertEquals(project.getId(), resolver.projectIdOfFlow(flow.getId())); + assertNull(resolver.projectIdOfFlow(null)); + } +} diff --git a/src/test/java/it/cnr/isti/workflow/manager/projects/ProjectExecutionTest.java b/src/test/java/it/cnr/isti/workflow/manager/projects/ProjectExecutionTest.java new file mode 100644 index 0000000..80063ad --- /dev/null +++ b/src/test/java/it/cnr/isti/workflow/manager/projects/ProjectExecutionTest.java @@ -0,0 +1,218 @@ +package it.cnr.isti.workflow.manager.projects; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.util.List; +import java.util.UUID; + +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.http.HttpStatus; +import org.springframework.test.context.TestPropertySource; +import org.springframework.web.server.ResponseStatusException; + +import it.cnr.isti.workflow.manager.auth.repo.LoginEntity; +import it.cnr.isti.workflow.manager.blocks.Block; +import it.cnr.isti.workflow.manager.blocks.configurations.LLMBlockConfiguration; +import it.cnr.isti.workflow.manager.blocks.types.LLMBlockType; +import it.cnr.isti.workflow.manager.controllers.BlocksController; +import it.cnr.isti.workflow.manager.controllers.ExecutionsController; +import it.cnr.isti.workflow.manager.controllers.FlowController; +import it.cnr.isti.workflow.manager.controllers.ProjectController; +import it.cnr.isti.workflow.manager.executions.api.ExecutionGroupView; +import it.cnr.isti.workflow.manager.executions.api.ExecutionView; +import it.cnr.isti.workflow.manager.executions.repo.ExecutionRepository; +import it.cnr.isti.workflow.manager.flows.model.FlowCreateRequest; +import it.cnr.isti.workflow.manager.flows.model.FlowData; +import it.cnr.isti.workflow.manager.flows.model.FlowProjectUpdateRequest; +import it.cnr.isti.workflow.manager.flows.model.FlowView; +import it.cnr.isti.workflow.manager.llms.LLMDescriptor; +import it.cnr.isti.workflow.manager.projects.model.ProjectCreateRequest; +import it.cnr.isti.workflow.manager.projects.model.ProjectExecuteRequest; +import it.cnr.isti.workflow.manager.projects.model.ProjectExecutionPlanView; +import it.cnr.isti.workflow.manager.projects.model.ProjectRunView; +import it.cnr.isti.workflow.manager.projects.model.ProjectView; + +@SpringBootTest +@TestPropertySource(locations = "classpath:test.properties") +public class ProjectExecutionTest { + + private static final String OWNER = "run-user"; + + @Autowired + private ProjectController projectController; + + @Autowired + private FlowController flowController; + + @Autowired + private BlocksController blocksController; + + @Autowired + private ExecutionsController executionsController; + + @Autowired + private ExecutionRepository executionRepository; + + private static LoginEntity user(String username) { + return new LoginEntity(username, "password"); + } + + private ProjectView newProject(String owner) { + return projectController.create(new ProjectCreateRequest("Run " + UUID.randomUUID(), null), user(owner)); + } + + /** A single-block flow with no inputs, so it is executable straight away. */ + private FlowView executableFlow(String owner, String projectId, String name) { + Block block = blocksController.create(LLMBlockConfiguration.builder() + .name(name) + .llmDescriptor(LLMDescriptor.builder().provider("testProvider").model("testModel").build()) + .prompt("Do the thing") + .build()); + FlowView flow = flowController.createFlow( + new FlowCreateRequest(name + " " + UUID.randomUUID(), null, + FlowData.builder().block(block).build()), + user(owner)).getBody(); + return projectId == null + ? flow + : flowController.updateProject(flow.id(), new FlowProjectUpdateRequest(projectId), user(owner)).getBody(); + } + + /** An empty flow is saved as DRAFT and is never executable. */ + private FlowView draftFlow(String owner, String projectId) { + FlowView flow = flowController.createFlow( + new FlowCreateRequest("Draft " + UUID.randomUUID(), null, FlowData.builder().build()), + user(owner)).getBody(); + return flowController.updateProject(flow.id(), new FlowProjectUpdateRequest(projectId), user(owner)).getBody(); + } + + @Test + public void createsOneExecutionPerFlowSharingOneProjectRunId() { + ProjectView project = newProject(OWNER); + executableFlow(OWNER, project.id(), "First"); + executableFlow(OWNER, project.id(), "Second"); + + ProjectExecutionPlanView plan = projectController.execute( + project.id(), new ProjectExecuteRequest("Nightly"), false, user(OWNER)); + + assertEquals(2, plan.run().executionCount()); + assertTrue(plan.skipped().isEmpty()); + assertNotNull(plan.run().projectRunId()); + // One run id across executions of different flows is the whole point. + assertEquals(1, plan.run().executions().stream() + .map(ExecutionView::getProjectRunId).distinct().count()); + assertEquals(project.id(), plan.run().projectId()); + assertTrue(plan.run().executions().stream().allMatch(e -> e.getName().startsWith("Nightly - "))); + } + + @Test + public void createsNothingWhenAnyFlowIsNotExecutable() { + ProjectView project = newProject(OWNER); + executableFlow(OWNER, project.id(), "Fine"); + draftFlow(OWNER, project.id()); + + ResponseStatusException conflict = assertThrows(ResponseStatusException.class, + () -> projectController.execute(project.id(), null, false, user(OWNER))); + assertEquals(HttpStatus.CONFLICT, conflict.getStatusCode()); + + // All-or-nothing: a half-created run group is worse than a clear refusal. + assertTrue(executionRepository.findByProjectIdAndOwnerAndExecutionKindOrderByCreationTimeAsc( + project.id(), OWNER, it.cnr.isti.workflow.manager.executions.ExecutionKind.TOP_LEVEL).isEmpty()); + } + + @Test + public void skipsNonExecutableFlowsWhenAsked() { + ProjectView project = newProject(OWNER); + executableFlow(OWNER, project.id(), "Fine"); + FlowView draft = draftFlow(OWNER, project.id()); + + ProjectExecutionPlanView plan = projectController.execute(project.id(), null, true, user(OWNER)); + + assertEquals(1, plan.run().executionCount()); + assertEquals(1, plan.skipped().size()); + assertEquals(draft.id(), plan.skipped().get(0).flowId()); + } + + @Test + public void refusesAnEmptyProject() { + ProjectView project = newProject(OWNER); + + ResponseStatusException error = assertThrows(ResponseStatusException.class, + () -> projectController.execute(project.id(), null, false, user(OWNER))); + assertEquals(HttpStatus.BAD_REQUEST, error.getStatusCode()); + } + + @Test + public void refusesWhenEverySkippableFlowIsNotExecutable() { + ProjectView project = newProject(OWNER); + draftFlow(OWNER, project.id()); + + ResponseStatusException error = assertThrows(ResponseStatusException.class, + () -> projectController.execute(project.id(), null, true, user(OWNER))); + assertEquals(HttpStatus.CONFLICT, error.getStatusCode()); + } + + @Test + public void anotherUsersProjectCannotBeRun() { + ProjectView project = newProject("run-owner"); + executableFlow("run-owner", project.id(), "Theirs"); + + ResponseStatusException error = assertThrows(ResponseStatusException.class, + () -> projectController.execute(project.id(), null, false, user("run-intruder"))); + assertEquals(HttpStatus.NOT_FOUND, error.getStatusCode()); + } + + @Test + public void listsRunsMostRecentFirst() { + ProjectView project = newProject(OWNER); + executableFlow(OWNER, project.id(), "Only"); + + String first = projectController.execute(project.id(), null, false, user(OWNER)).run().projectRunId(); + String second = projectController.execute(project.id(), null, false, user(OWNER)).run().projectRunId(); + + List runs = projectController.listRuns(project.id(), user(OWNER)); + + assertEquals(2, runs.size()); + assertTrue(runs.stream().anyMatch(run -> run.projectRunId().equals(first))); + assertTrue(runs.stream().anyMatch(run -> run.projectRunId().equals(second))); + assertTrue(runs.get(0).createdAt() >= runs.get(1).createdAt()); + } + + @Test + public void executionsAreCreatedButNotStarted() { + ProjectView project = newProject(OWNER); + executableFlow(OWNER, project.id(), "Only"); + + ProjectExecutionPlanView plan = projectController.execute(project.id(), null, false, user(OWNER)); + + // Creation only, matching POST /executions: the client fills inputs and credentials, then starts. + assertTrue(plan.run().executions().stream() + .allMatch(execution -> !execution.getContext().getStatus().isFinalState() + && execution.getContext().getStatus().isInitState())); + } + + @Test + public void aProjectRunDoesNotDisturbTheSingleFlowExecutionGroups() { + // Regression guard for the decision not to reuse runGroupId: an execution group means "the + // rerun history of one flow" and takes its sourceFlowId from the first execution, so a + // project run sharing a runGroupId would have polluted it. + ProjectView project = newProject("group-user"); + FlowView flowA = executableFlow("group-user", project.id(), "A"); + FlowView flowB = executableFlow("group-user", project.id(), "B"); + + projectController.execute(project.id(), null, false, user("group-user")); + + List groups = executionsController.getGroups(user("group-user")); + List forFlows = groups.stream() + .filter(group -> flowA.id().equals(group.getSourceFlowId()) + || flowB.id().equals(group.getSourceFlowId())) + .toList(); + + assertEquals(2, forFlows.size(), "each flow keeps its own single-flow history group"); + assertTrue(forFlows.stream().allMatch(group -> group.getExecutionCount() == 1)); + } +} diff --git a/src/test/java/it/cnr/isti/workflow/manager/projects/ProjectRunOrchestrationTest.java b/src/test/java/it/cnr/isti/workflow/manager/projects/ProjectRunOrchestrationTest.java new file mode 100644 index 0000000..ba5bf5f --- /dev/null +++ b/src/test/java/it/cnr/isti/workflow/manager/projects/ProjectRunOrchestrationTest.java @@ -0,0 +1,313 @@ +package it.cnr.isti.workflow.manager.projects; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.util.List; +import java.util.UUID; + +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.http.HttpStatus; +import org.springframework.test.context.TestPropertySource; +import org.springframework.web.server.ResponseStatusException; + +import it.cnr.isti.workflow.manager.auth.repo.LoginEntity; +import it.cnr.isti.workflow.manager.blocks.Block; +import it.cnr.isti.workflow.manager.blocks.configurations.LLMBlockConfiguration; +import it.cnr.isti.workflow.manager.blocks.types.LLMBlockType; +import it.cnr.isti.workflow.manager.controllers.BlocksController; +import it.cnr.isti.workflow.manager.controllers.FlowController; +import it.cnr.isti.workflow.manager.controllers.ProjectController; +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.api.ExecutionView; +import it.cnr.isti.workflow.manager.flows.model.FlowCreateRequest; +import it.cnr.isti.workflow.manager.flows.model.FlowData; +import it.cnr.isti.workflow.manager.flows.model.FlowProjectUpdateRequest; +import it.cnr.isti.workflow.manager.flows.model.FlowView; +import it.cnr.isti.workflow.manager.ios.IODescriptor; +import it.cnr.isti.workflow.manager.ios.IOType; +import it.cnr.isti.workflow.manager.llms.LLMDescriptor; +import it.cnr.isti.workflow.manager.llms.providers.LLMProvider; +import it.cnr.isti.workflow.manager.projects.ProjectRunOrchestrator.ProjectRunStatus; +import it.cnr.isti.workflow.manager.projects.model.ProjectContext; +import it.cnr.isti.workflow.manager.projects.model.ProjectContextEntry; +import it.cnr.isti.workflow.manager.projects.model.ProjectCreateRequest; +import it.cnr.isti.workflow.manager.projects.model.ProjectExecutionPlanView; +import it.cnr.isti.workflow.manager.projects.model.ProjectFlowOrderRequest; +import it.cnr.isti.workflow.manager.projects.model.ProjectRunView; +import it.cnr.isti.workflow.manager.projects.model.ProjectView; + +@SpringBootTest +@TestPropertySource(locations = "classpath:test.properties") +public class ProjectRunOrchestrationTest { + + private static final String OWNER = "orch-user"; + private static final String FAILING_PROMPT = "PLEASE FAIL"; + + /** A provider that actually answers, so a step can complete and the run can advance. */ + @TestConfiguration + static class TestConfig { + + @Bean + public LLMProvider orchestrationTestProvider() { + return new LLMProvider() { + + @Override + public String getName() { + return "testProvider"; + } + + @Override + public List getRegisteredModels() { + return List.of("testModel"); + } + + @Override + public String generate(String model, String prompt) { + if (prompt.contains(FAILING_PROMPT)) { + throw new IllegalStateException("deliberate failure"); + } + return "done"; + } + }; + } + } + + @Autowired + private ProjectController projectController; + + @Autowired + private FlowController flowController; + + @Autowired + private BlocksController blocksController; + + @Autowired + private ExecutionsService executionsService; + + private static LoginEntity user(String username) { + return new LoginEntity(username, "password"); + } + + private ProjectView newProject(String owner) { + return projectController.create(new ProjectCreateRequest("Orch " + UUID.randomUUID(), null), user(owner)); + } + + private FlowView flowInProject(String owner, String projectId, String name, String prompt, + String globalInputName) { + Block block = blocksController.create(LLMBlockConfiguration.builder() + .name(name) + .llmDescriptor(LLMDescriptor.builder().provider("testProvider").model("testModel").build()) + .prompt(prompt) + .build()); + + FlowData.FlowDataBuilder data = FlowData.builder().block(block); + if (globalInputName != null) { + data.globalInput(IODescriptor.of(globalInputName, IOType.TEXT)); + } + + FlowView flow = flowController.createFlow( + new FlowCreateRequest(name + " " + UUID.randomUUID(), null, data.build()), user(owner)).getBody(); + return flowController.updateProject(flow.id(), new FlowProjectUpdateRequest(projectId), user(owner)).getBody(); + } + + private FlowView plainFlow(String owner, String projectId, String name) { + return flowInProject(owner, projectId, name, "Do the thing", null); + } + + /** + * Polls until the run reaches a state it will not leave on its own. RUNNING and PENDING are both + * transient: PENDING is the brief gap between one step finishing and the next being started. + */ + private ProjectRunView awaitSettled(String projectId, String projectRunId) { + long deadline = System.currentTimeMillis() + 20_000; + ProjectRunView run = null; + while (System.currentTimeMillis() < deadline) { + run = projectController.listRuns(projectId, user(OWNER)).stream() + .filter(candidate -> candidate.projectRunId().equals(projectRunId)) + .findFirst() + .orElseThrow(); + if (run.status() != ProjectRunStatus.RUNNING && run.status() != ProjectRunStatus.PENDING) { + return run; + } + try { + Thread.sleep(50); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + break; + } + } + return run; + } + + @Test + public void runsTheFlowsOneAtATimeInTheProjectOrder() { + ProjectView project = newProject(OWNER); + FlowView first = plainFlow(OWNER, project.id(), "AAA"); + FlowView second = plainFlow(OWNER, project.id(), "BBB"); + // Deliberately the reverse of the alphabetical order, so only the project order can explain it. + projectController.updateFlowOrder(project.id(), + new ProjectFlowOrderRequest(List.of(second.id(), first.id())), user(OWNER)); + + ProjectExecutionPlanView plan = projectController.execute(project.id(), null, false, user(OWNER)); + assertEquals(List.of(second.id(), first.id()), + plan.run().executions().stream().map(ExecutionView::getSourceFlowId).toList()); + assertEquals(ProjectRunStatus.PENDING, plan.run().status()); + + projectController.startRun(project.id(), plan.run().projectRunId(), user(OWNER)); + ProjectRunView settled = awaitSettled(project.id(), plan.run().projectRunId()); + + assertEquals(ProjectRunStatus.COMPLETED, settled.status()); + assertEquals(2, settled.completedCount()); + + ExecutionObject step1 = executionsService.getExecution(plan.run().executions().get(0).getId()); + ExecutionObject step2 = executionsService.getExecution(plan.run().executions().get(1).getId()); + assertEquals(ExecutionStatus.SUCCESS, step1.getContext().getStatus()); + assertEquals(ExecutionStatus.SUCCESS, step2.getContext().getStatus()); + // Sequential, not parallel: the second step only began after the first had finished. + assertTrue(step2.getContext().getStartTime() >= step1.getContext().getEndTime(), + "the second step must start after the first one ends"); + } + + @Test + public void aFailedStepStopsTheRunAndLeavesTheRestUntouched() { + ProjectView project = newProject(OWNER); + FlowView failing = flowInProject(OWNER, project.id(), "Boom", FAILING_PROMPT, null); + FlowView later = plainFlow(OWNER, project.id(), "Later"); + projectController.updateFlowOrder(project.id(), + new ProjectFlowOrderRequest(List.of(failing.id(), later.id())), user(OWNER)); + + ProjectExecutionPlanView plan = projectController.execute(project.id(), null, false, user(OWNER)); + projectController.startRun(project.id(), plan.run().projectRunId(), user(OWNER)); + ProjectRunView settled = awaitSettled(project.id(), plan.run().projectRunId()); + + assertEquals(ProjectRunStatus.STOPPED, settled.status()); + // Nothing is lost: the remaining step stays untouched so it can be resumed later. + ExecutionObject next = executionsService.getExecution(plan.run().executions().get(1).getId()); + assertTrue(next.getContext().getStatus().isInitState()); + } + + @Test + public void aRunIsBlockedWhenTheNextStepStillNeedsAGlobalInput() { + ProjectView project = newProject(OWNER); + flowInProject(OWNER, project.id(), "NeedsInput", "Do it with ${{global.requirements}}", "requirements"); + + ProjectExecutionPlanView plan = projectController.execute(project.id(), null, false, user(OWNER)); + ProjectRunView run = projectController.startRun(project.id(), plan.run().projectRunId(), user(OWNER)); + + assertEquals(ProjectRunStatus.BLOCKED, run.status()); + assertNotNull(run.blockedReason()); + assertTrue(run.blockedReason().contains("requirements")); + } + + @Test + public void aBlockedRunResumesOnceTheInputIsSupplied() { + ProjectView project = newProject(OWNER); + flowInProject(OWNER, project.id(), "NeedsInput", "Do it with ${{global.requirements}}", "requirements"); + + ProjectExecutionPlanView plan = projectController.execute(project.id(), null, false, user(OWNER)); + projectController.startRun(project.id(), plan.run().projectRunId(), user(OWNER)); + + executionsService.setGlobalInput(plan.run().executions().get(0).getId(), "requirements", "Ship it"); + projectController.startRun(project.id(), plan.run().projectRunId(), user(OWNER)); + + assertEquals(ProjectRunStatus.COMPLETED, awaitSettled(project.id(), plan.run().projectRunId()).status()); + } + + @Test + public void theProjectContextPreFillsAMatchingGlobalInputSoTheRunNeedsNoRoundTrip() { + // This is what makes a one-click project run possible: the execution reaches READY on its + // own, with no per-flow round trip to supply the value. + ProjectView project = newProject(OWNER); + projectController.updateContext(project.id(), + ProjectContext.builder() + .entry(ProjectContextEntry.builder() + .name("requirements").type(IOType.TEXT).value("Ship it").build()) + .build(), + user(OWNER)); + flowInProject(OWNER, project.id(), "NeedsInput", "Do it with ${{global.requirements}}", "requirements"); + + ProjectExecutionPlanView plan = projectController.execute(project.id(), null, false, user(OWNER)); + + assertTrue(executionsService.getExecution(plan.run().executions().get(0).getId()) + .getMissingGlobalInputKeys().isEmpty()); + + projectController.startRun(project.id(), plan.run().projectRunId(), user(OWNER)); + assertEquals(ProjectRunStatus.COMPLETED, awaitSettled(project.id(), plan.run().projectRunId()).status()); + } + + @Test + public void aProjectValueNeverOverwritesOneTheUserAlreadySupplied() { + ProjectView project = newProject(OWNER); + projectController.updateContext(project.id(), + ProjectContext.builder() + .entry(ProjectContextEntry.builder() + .name("requirements").type(IOType.TEXT).value("from-project").build()) + .build(), + user(OWNER)); + flowInProject(OWNER, project.id(), "NeedsInput", "Do it with ${{global.requirements}}", "requirements"); + + ProjectExecutionPlanView plan = projectController.execute(project.id(), null, false, user(OWNER)); + String executionId = plan.run().executions().get(0).getId(); + executionsService.setGlobalInput(executionId, "requirements", "from-user"); + + assertEquals("from-user", + executionsService.getExecution(executionId).getContext().getGlobalInputs().get("requirements")); + } + + @Test + public void reorderingRefusesAFlowThatIsNotInTheProject() { + ProjectView project = newProject(OWNER); + plainFlow(OWNER, project.id(), "Mine"); + FlowView outsider = plainFlow(OWNER, newProject(OWNER).id(), "Other"); + + ResponseStatusException error = assertThrows(ResponseStatusException.class, + () -> projectController.updateFlowOrder(project.id(), + new ProjectFlowOrderRequest(List.of(outsider.id())), user(OWNER))); + assertEquals(HttpStatus.BAD_REQUEST, error.getStatusCode()); + } + + @Test + public void reorderingKeepsOmittedFlowsAfterTheListedOnes() { + ProjectView project = newProject(OWNER); + FlowView a = plainFlow(OWNER, project.id(), "AAA"); + FlowView b = plainFlow(OWNER, project.id(), "BBB"); + FlowView c = plainFlow(OWNER, project.id(), "CCC"); + + ProjectView reordered = projectController.updateFlowOrder(project.id(), + new ProjectFlowOrderRequest(List.of(c.id())), user(OWNER)); + + List order = reordered.flows().stream().map(flow -> flow.id()).toList(); + assertEquals(c.id(), order.get(0)); + assertTrue(order.containsAll(List.of(a.id(), b.id()))); + } + + @Test + public void startingAnUnknownRunIsNotFound() { + ProjectView project = newProject(OWNER); + plainFlow(OWNER, project.id(), "One"); + + ResponseStatusException error = assertThrows(ResponseStatusException.class, + () -> projectController.startRun(project.id(), "no-such-run", user(OWNER))); + assertEquals(HttpStatus.NOT_FOUND, error.getStatusCode()); + } + + @Test + public void anotherUsersRunCannotBeStarted() { + ProjectView project = newProject("orch-owner"); + plainFlow("orch-owner", project.id(), "Theirs"); + ProjectExecutionPlanView plan = projectController.execute(project.id(), null, false, user("orch-owner")); + + ResponseStatusException error = assertThrows(ResponseStatusException.class, + () -> projectController.startRun(project.id(), plan.run().projectRunId(), user("orch-intruder"))); + assertEquals(HttpStatus.NOT_FOUND, error.getStatusCode()); + } +}