Add projects: grouping, shared context, and sequential project runs

A project groups 1..N flows. A flow's project is optional, so flows without one
stay fully valid, and the association is a plain String column - this codebase
has no JPA relations, and the hottest entity is not the place to introduce the
first one.

Ownership and visibility
- Projects are strictly owner-scoped and report another user's project as 404,
  not 403: unlike a flow, whose existence is not secret, a project is private
  workspace structure.
- projectId and projectName are disclosed only to the flow's owner. A published
  flow inside a project stays readable by everyone - membership is organizational
  structure, not access control - but a project name can itself be sensitive, so
  a non-owner sees the flow as unassigned.
- Assignment gets its own PUT /flows/{id}/project rather than a field on the flow
  body: that body is the full-replace PUT the editor issues on every save, so a
  project carried there would be silently dropped on each save. A finalized flow
  stays reassignable, since finalizing is irreversible and must not freeze a flow
  out of reorganization forever.

Deleting a project deletes its flows, finalized ones included, and requires
confirm=true. Refusing the cascade was not an option: finalizing cannot be
undone, so a project holding one finalized flow could never be deleted by any
API. Executions are kept - each snapshots its own copy of the flow graph.

Shared context
Values a project shares with its flows, readable as ${{project.x}} and as
#project['x'] in conditional expressions. Three changes make that work, each
silent if missed: the project. prefix is preserved by ExecutionTemplateResolver
(otherwise keys publish as ${{vars.project.x}} and never resolve), admitted by
PlaceholderInputs (otherwise the placeholder becomes a dangling block input), and
bound in the SpEL context. Values are frozen into the execution snapshot at
creation, so editing a project never rewrites a run that already happened, and
they are applied only when the person running owns the project.

Project runs
POST /projects/{id}/execute creates one execution per flow, sharing a
projectRunId, and refuses the lot if any flow is not executable - a half-created
run group is worse than a clear refusal. POST .../runs/{runId}/start then runs
them one at a time in the project's order, each step starting only when the
previous succeeded. The run keeps no state of its own: its position is derived
from the executions, so there is no second status machine to drift out of sync.
A failed step stops the run and leaves the rest untouched, so it can be resumed.
Project values pre-fill matching global inputs - same name and type, never
overwriting a user's own value - which is what lets a run start without a
per-flow round trip.

Also fixes an ordering hazard in AuthService: account deletion deliberately
preserves finalized flows, so their project assignment is now cleared before the
projects are removed.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
Lucio Lelii 2026-09-03 07:53:25 +02:00
parent 61bbbe50d4
commit f9b94d823e
53 changed files with 3049 additions and 24 deletions

View File

@ -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);

View File

@ -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);
}

View File

@ -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<FlowSummaryView> 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<FlowView> 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<FlowView> 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<FlowView> 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<FlowView> updateFinalized(@PathVariable String id,

View File

@ -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<ProjectSummaryView> 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.<name>}}, 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<ProjectRunView> 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);
}
}

View File

@ -36,6 +36,12 @@ public class ExecutionContext implements ExecutionListener {
Map<String, Object> globalInputs = new HashMap<>();
Map<String, ExecutionVariableDescriptor> 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<String, Object> projectContext = new HashMap<>();
@Getter(AccessLevel.NONE)
@JsonIgnore
Map<String, Object> 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<String, Object> 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<String, Object> projectContext) {
this.projectContext.clear();
if (projectContext != null) {
this.projectContext.putAll(projectContext);
}
refreshRuntimeExecutionVariables();
}
protected void setRuntimeContextValues(Map<String, Object> 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));
}

View File

@ -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<ExecutionAuthorizationRequirement> 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<String, Object> 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<Step<?>> 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()

View File

@ -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<String, Object> runtimeValues(Map<String, Object> executionVariables, Map<String, Object> globalInputs,
Map<String, Object> contextValues) {
return runtimeValues(executionVariables, globalInputs, null, contextValues);
}
/**
* Merges every source of runtime values into the single flat map the {@code ${{...}}} resolver
* reads.
*
* <p>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.
*
* <p>Context values go in last so the system-reserved namespace stays authoritative.
*/
public static Map<String, Object> runtimeValues(Map<String, Object> executionVariables, Map<String, Object> globalInputs,
Map<String, Object> projectValues, Map<String, Object> contextValues) {
LinkedHashMap<String, Object> 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<String, Object> projectView(Map<String, Object> runtimeVariables) {
return viewOfPrefix(runtimeVariables, PROJECT_PREFIX);
}
public static Map<String, Object> globalView(Map<String, Object> runtimeVariables) {
LinkedHashMap<String, Object> global = new LinkedHashMap<>();
return viewOfPrefix(runtimeVariables, GLOBAL_PREFIX);
}
private static Map<String, Object> viewOfPrefix(Map<String, Object> runtimeVariables, String prefix) {
LinkedHashMap<String, Object> 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;
}
}

View File

@ -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)

View File

@ -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> projectRunOrchestrator;
@Autowired
Map<String, LLMProvider> 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.
*
* <p>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.
*
* <p>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.
*
* <p>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<String, Object> projectValues) {
if (projectValues == null || projectValues.isEmpty()) {
return;
}
Map<String, Object> 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<String, Object> projectContext) {
flowExecutionValidator.validate(flow);
List<ExecutionAuthorizationRequirement> 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 <em>different</em> 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(),

View File

@ -21,6 +21,8 @@ public class ExecutionContextView {
private Map<String, ExecutionVariableDescriptor> executionVariableDescriptors;
private Map<String, Object> globalInputs;
private Map<String, ExecutionVariableDescriptor> globalInputDescriptors;
/** Values inherited from the flow's project, frozen when the run was created. */
private Map<String, Object> projectContext;
private Map<FieldKey, Object> result;
private Map<FieldKey, Object> 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())

View File

@ -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())

View File

@ -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;
}
}

View File

@ -31,6 +31,11 @@ public class ExecutionSnapshot {
private Map<String, ExecutionVariableDescriptor> executionVariableDescriptors;
private Map<String, Object> globalInputs;
private Map<String, ExecutionVariableDescriptor> 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<String, Object> projectContext;
private Map<String, Object> providedAuthorizations;
private Map<FieldKey, Object> inputs;
private Map<FieldKey, Object> result;

View File

@ -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;

View File

@ -20,6 +20,11 @@ public interface ExecutionRepository extends JpaRepository<ExecutionEntity, Stri
String runGroupId, String owner, ExecutionKind executionKind);
List<ExecutionEntity> findBySourceFlowIdAndOwnerAndExecutionKindOrderByRunNumberAscCreationTimeAsc(
String sourceFlowId, String owner, ExecutionKind executionKind);
List<ExecutionEntity> findByProjectRunIdAndOwnerOrderByCreationTimeAsc(String projectRunId, String owner);
List<ExecutionEntity> findByProjectIdAndOwnerAndExecutionKindOrderByCreationTimeAsc(
String projectId, String owner, ExecutionKind executionKind);
List<ExecutionEntity> findByParentExecutionIdOrderByCreationTimeAsc(String parentExecutionId);
List<ExecutionEntity> findByParentExecutionIdAndParentStepIdOrderByParentIterationIndexAscCreationTimeAsc(
String parentExecutionId, String parentStepId);

View File

@ -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);
}
}

View File

@ -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.
*
* <p>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<FlowView> getAllFlows(String owner) {
return flowRepository.findFlowsByOwnerOrPublic(owner).stream().map(this::toView).toList();
List<FlowEntity> entities = flowRepository.findFlowsByOwnerOrPublic(owner);
Map<String, String> 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<FlowSummaryView> getAllFlowSummaries(String owner) {
List<FlowEntity> entities = flowRepository.findFlowsByOwnerOrPublic(owner);
Map<String, String> 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<ValidationError> 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<String, String> 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<String, String> 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<String, String> projectNamesFor(String owner) {
return projectRepository.findByOwnerOrderByNameAsc(owner).stream()
.collect(java.util.stream.Collectors.toMap(ProjectEntity::getId, ProjectEntity::getName));
}
}

View File

@ -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.
*
* <p>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) {
}

View File

@ -0,0 +1,23 @@
package it.cnr.isti.workflow.manager.flows.model;
import java.time.LocalDateTime;
/**
* A flow without its graph.
*
* <p>{@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) {
}

View File

@ -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) {
}

View File

@ -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;

View File

@ -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<FlowEntity, String> {
List<FlowEntity> findByOwner(String owner);
List<FlowEntity> findByOwnerAndFinalizedFalse(String owner);
List<FlowEntity> findByProjectIdOrderByProjectOrderAscNameAsc(String projectId);
List<FlowEntity> 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<ProjectFlowCount> 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);
}

View File

@ -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();
}

View File

@ -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
}

View File

@ -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<ProjectContext, String> {
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);
}
}
}

View File

@ -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.
*
* <p>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<String, Object> 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}.
*
* <p>The context is applied <strong>only when the executing user owns the project</strong>. 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);
}
}

View File

@ -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}.
*
* <p><strong>It creates the executions, it does not start them.</strong> 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.
*
* <p>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<FlowEntity> flows = flowRepository.findByProjectIdOrderByProjectOrderAscNameAsc(projectId);
if (flows.isEmpty()) {
throw new ResponseStatusException(HttpStatus.BAD_REQUEST, "Project contains no flows to run");
}
Map<FlowEntity, List<ValidationError>> 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<FlowEntity> 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<ExecutionView> 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<SkippedFlow> 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<ProjectRunView> listRuns(String owner, String projectId) {
requireOwnedProject(owner, projectId);
Map<String, List<ExecutionEntity>> 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<ExecutionEntity> entities) {
List<ExecutionView> 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<ExecutionView> 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<FlowEntity, List<ValidationError>> notExecutable) {
List<ValidationError> 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"));
}
}

View File

@ -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<FlowEntity> flows) {
List<FlowSummaryView> 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;
}
}

View File

@ -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.
*
* <p>The run has <strong>no state of its own</strong>. 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.
*
* <p>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<String, Object> 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<ExecutionEntity> 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<ExecutionEntity> 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<ExecutionEntity> 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<String> missingInputs = execution.getMissingGlobalInputKeys();
List<String> missingAuthorizations = execution.getMissingAuthorizationKeys();
List<String> 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<ExecutionEntity> 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<ExecutionEntity> requireRun(String owner, String projectRunId) {
List<ExecutionEntity> ordered = orderedExecutions(owner, projectRunId);
if (ordered.isEmpty()) {
throw new ResponseStatusException(HttpStatus.NOT_FOUND, "Project run not found");
}
return ordered;
}
private List<ExecutionEntity> 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());
}
}

View File

@ -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.<name>}}}, 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<String> 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<ProjectSummaryView> list(String owner) {
Map<String, Long> 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.
*
* <p>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<String> flowIds) {
requireOwnedProject(owner, projectId);
List<FlowEntity> flows = flowRepository.findByProjectIdOrderByProjectOrderAscNameAsc(projectId);
Map<String, FlowEntity> byId = flows.stream()
.collect(Collectors.toMap(FlowEntity::getId, flow -> flow));
List<String> requested = flowIds == null ? List.of() : flowIds;
Set<String> 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<String> 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.
*
* <p>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}.
*
* <p>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<FlowEntity> 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<FlowEntity> flows) {
long finalizedCount = flows.stream().filter(FlowEntity::isFinalized).count();
List<ValidationError> 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<ValidationError> errors = new ArrayList<>();
Set<String> 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<FlowEntity> 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;
}
}

View File

@ -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.
*
* <p>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<ProjectContextEntry> entries;
/** Flattens the entries to the {@code name -> value} map the execution runtime publishes. */
@JsonIgnore
public Map<String, Object> toValueMap() {
Map<String, Object> 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();
}
}

View File

@ -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.
*
* <p>The {@code name} is what a flow author writes as {@code ${{project.<name>}}}, 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;
}

View File

@ -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) {
}

View File

@ -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) {
}

View File

@ -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<SkippedFlow> skipped) {
public record SkippedFlow(String flowId, String flowName, String reason) {
}
}

View File

@ -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<String> flowIds) {
}

View File

@ -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}.
*
* <p>Deliberately separate from the execution <em>groups</em> 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<ExecutionView> executions) {
}

View File

@ -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) {
}

View File

@ -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}.
*
* <p>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) {
}

View File

@ -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<FlowSummaryView> flows) {
}

View File

@ -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;
}

View File

@ -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<ProjectEntity, String> {
List<ProjectEntity> 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<ProjectEntity> findByIdAndOwner(String id, String owner);
boolean existsByOwnerAndName(String owner, String name);
boolean existsByOwnerAndNameAndIdNot(String owner, String name, String id);
void deleteByOwner(String owner);
}

View File

@ -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);

View File

@ -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);

View File

@ -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;

View File

@ -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<FlowView> 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<ProjectSummaryView> 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.<Runnable>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<FlowView> 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());
}
}

View File

@ -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<LoginEntity> 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();

View File

@ -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<String, Object> 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<String, Object> 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<String, Object> 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));
}
}

View File

@ -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

View File

@ -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<LLMBlockType> 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<LLMBlockType> 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<LLMBlockType> 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());
}
}

View File

@ -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));
}
}

View File

@ -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<LLMBlockType> 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<ProjectRunView> 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<ExecutionGroupView> groups = executionsController.getGroups(user("group-user"));
List<ExecutionGroupView> 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));
}
}

View File

@ -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<String> 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<LLMBlockType> 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<String> 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());
}
}