diff --git a/nflow-engine/src/main/java/io/nflow/engine/internal/dao/WorkflowInstanceDao.java b/nflow-engine/src/main/java/io/nflow/engine/internal/dao/WorkflowInstanceDao.java index 491333772..db82142d8 100644 --- a/nflow-engine/src/main/java/io/nflow/engine/internal/dao/WorkflowInstanceDao.java +++ b/nflow-engine/src/main/java/io/nflow/engine/internal/dao/WorkflowInstanceDao.java @@ -284,13 +284,12 @@ void insertVariables(final long id, final long actionId, Map cha } private void insertVariablesWithMultipleUpdates(final long id, final long actionId, Map changedStateVariables) { - for (Entry entry : changedStateVariables.entrySet()) { - int updated = jdbc.update(insertWorkflowInstanceStateSql() + " values (?,?,?,?)", id, actionId, entry.getKey(), - entry.getValue()); + changedStateVariables.forEach((key, value) -> { + int updated = jdbc.update(insertWorkflowInstanceStateSql() + " values (?,?,?,?)", id, actionId, key, value); if (updated != 1) { - throw new IllegalStateException("Failed to insert state variable " + entry.getKey()); + throw new IllegalStateException("Failed to insert state variable " + key); } - } + }); } private void insertVariablesWithBatchUpdate(final long id, final long actionId, Map changedStateVariables) { @@ -310,8 +309,8 @@ protected boolean setValuesIfAvailable(PreparedStatement ps, int i) throws SQLEx return true; } }); - int updatedRows = 0; boolean unknownResults = false; + AtomicInteger updatedRows = new AtomicInteger(0); for (int i = 0; i < updateStatus.length; ++i) { if (updateStatus[i] == Statement.SUCCESS_NO_INFO) { unknownResults = true; @@ -320,12 +319,13 @@ protected boolean setValuesIfAvailable(PreparedStatement ps, int i) throws SQLEx if (updateStatus[i] == Statement.EXECUTE_FAILED) { throw new IllegalStateException("Failed to insert/update state variable at index " + i + " (" + updateStatus[i] + ")"); } - updatedRows += updateStatus[i]; + updatedRows.addAndGet(updateStatus[i]); } + int updatedRowsCount = updatedRows.get(); int changedVariables = changedStateVariables.size(); - if (!unknownResults && updatedRows != changedVariables) { + if (!unknownResults && updatedRowsCount != changedVariables) { throw new IllegalStateException( - "Failed to insert/update state variables, expected update count " + changedVariables + ", actual " + updatedRows); + "Failed to insert/update state variables, expected update count " + changedVariables + ", actual " + updatedRowsCount); } } @@ -385,14 +385,12 @@ protected void doInTransactionWithoutResult(TransactionStatus status) { } long parentActionId = insertWorkflowInstanceAction(action); insertVariables(action.workflowInstanceId, parentActionId, changedStateVariables); - for (WorkflowInstance childTemplate : childWorkflows) { + childWorkflows.forEach(childTemplate -> { WorkflowInstance childWorkflow = new WorkflowInstance.Builder(childTemplate).setParentWorkflowId(instance.id) .setParentActionId(parentActionId).build(); insertWorkflowInstance(childWorkflow); - } - for (WorkflowInstance workflow : workflows) { - insertWorkflowInstance(workflow); - } + }); + workflows.forEach(workflow -> insertWorkflowInstance(workflow)); } }); } @@ -655,18 +653,18 @@ private List updateNextWorkflowInstancesWithBatchUpdate(List batchArgs = instances.stream() .map(instance -> new Object[] { instance.id, sqlVariants.tuneTimestampForDb(instance.modified) }).collect(toList()); int[] updateStatuses = jdbc.batchUpdate(sql, batchArgs); - List ids = new ArrayList<>(instances.size()); - for (int i = 0; i < updateStatuses.length; ++i) { - int status = updateStatuses[i]; - if (status == 1) { - ids.add(instances.get(i).id); - } else if (status != 0) { - disableBatchUpdates.set(true); - throw new PollingBatchException( - "Database was unable to provide information about affected rows in a batch update. Disabling batch updates."); - } - } - return ids; + return Stream.iterate(0, i -> i < updateStatuses.length, i -> i + 1) + .peek(i -> { + int status = updateStatuses[i]; + if (status != 1 && status != 0) { + disableBatchUpdates.set(true); + throw new PollingBatchException( + "Database was unable to provide information about affected rows in a batch update. Disabling batch updates."); + } + }) + .filter(i -> updateStatuses[i] == 1) + .map(i -> instances.get(i).id) + .collect(toList()); } private static class OptimisticLockKey extends ModelObject implements Comparable { @@ -799,10 +797,9 @@ private void fillChildWorkflowIds(final WorkflowInstance instance, boolean query } private long getMaxResults(Long maxResults) { - if (maxResults == null) { - return workflowInstanceQueryMaxResultsDefault; - } - return min(maxResults, workflowInstanceQueryMaxResults); + return Optional.ofNullable(maxResults) + .map(m -> min(m, workflowInstanceQueryMaxResults)) + .orElse(workflowInstanceQueryMaxResultsDefault); } private void fillActions(WorkflowInstance instance, boolean includeStateVariables, Long requestedMaxActions) { @@ -815,20 +812,17 @@ private void fillActions(WorkflowInstance instance, boolean includeStateVariable if (includeStateVariables) { Map> actionStates = fetchActionStateVariables(instance, actionBuilders.size(), maxActions); actionBuilders.forEach(builder -> { - Map actionState = actionStates.get(builder.getId()); - if (actionState != null) { - builder.setUpdatedStateVariables(actionState); - } + Map actionState = actionStates.get(builder.getId()); + Optional.ofNullable(actionState).ifPresent(builder::setUpdatedStateVariables); }); } actionBuilders.stream().map(WorkflowInstanceAction.Builder::build).forEach(instance.actions::add); } private long getMaxActions(Long maxActions) { - if (maxActions == null) { - return workflowInstanceQueryMaxActionsDefault; - } - return min(maxActions, workflowInstanceQueryMaxActions); + return Optional.ofNullable(maxActions) + .map(m -> min(m, workflowInstanceQueryMaxActions)) + .orElse(workflowInstanceQueryMaxActionsDefault); } private Map> fetchActionStateVariables(WorkflowInstance instance, long actions, long maxActions) { diff --git a/nflow-engine/src/main/java/io/nflow/engine/internal/executor/WorkflowDispatcher.java b/nflow-engine/src/main/java/io/nflow/engine/internal/executor/WorkflowDispatcher.java index 0fe65db46..27f857f81 100644 --- a/nflow-engine/src/main/java/io/nflow/engine/internal/executor/WorkflowDispatcher.java +++ b/nflow-engine/src/main/java/io/nflow/engine/internal/executor/WorkflowDispatcher.java @@ -183,9 +183,7 @@ private void dispatch(List nextInstanceIds) { return; } logger.debug("Found {} workflow instances, dispatching executors.", nextInstanceIds.size()); - for (Long instanceId : nextInstanceIds) { - executor.execute(stateProcessorFactory.createProcessor(instanceId, shutdownRequested::get)); - } + nextInstanceIds.forEach(instanceId -> executor.execute(stateProcessorFactory.createProcessor(instanceId, shutdownRequested::get))); } private List getNextInstanceIds() { diff --git a/nflow-engine/src/main/java/io/nflow/engine/internal/executor/WorkflowStateProcessor.java b/nflow-engine/src/main/java/io/nflow/engine/internal/executor/WorkflowStateProcessor.java index 4c9f27f54..d42826970 100644 --- a/nflow-engine/src/main/java/io/nflow/engine/internal/executor/WorkflowStateProcessor.java +++ b/nflow-engine/src/main/java/io/nflow/engine/internal/executor/WorkflowStateProcessor.java @@ -9,6 +9,7 @@ import static io.nflow.engine.workflow.instance.WorkflowInstanceAction.WorkflowActionType.stateExecutionFailed; import static java.lang.Thread.currentThread; import static java.util.Arrays.asList; +import static java.util.Arrays.stream; import static java.util.Collections.emptyList; import static java.util.concurrent.TimeUnit.SECONDS; import static org.apache.commons.lang3.exception.ExceptionUtils.getStackTrace; @@ -24,7 +25,9 @@ import java.util.Iterator; import java.util.List; import java.util.Map; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.function.Supplier; +import java.util.stream.Collectors; import org.joda.time.DateTime; import org.joda.time.Duration; @@ -563,33 +566,33 @@ public NextAction processState() { } private void processBeforeListeners() { - for (WorkflowExecutorListener listener : executorListeners) { + executorListeners.forEach(listener -> { try { listener.beforeProcessing(listenerContext); } catch (Throwable t) { logger.error("Error in {}.beforeProcessing ({})", listener.getClass().getName(), t.getMessage(), t); } - } + }); } private void processAfterListeners() { - for (WorkflowExecutorListener listener : executorListeners) { + executorListeners.forEach(listener -> { try { listener.afterProcessing(listenerContext); } catch (Throwable t) { logger.error("Error in {}.afterProcessing ({})", listener.getClass().getName(), t.getMessage(), t); } - } + }); } private void processAfterFailureListeners(Throwable ex) { - for (WorkflowExecutorListener listener : executorListeners) { + executorListeners.forEach(listener -> { try { listener.afterFailure(listenerContext, ex); } catch (Throwable t) { logger.error("Error in {}.afterFailure ({})", listener.getClass().getName(), t.getMessage(), t); } - } + }); } public DateTime getStartTime() { @@ -602,25 +605,26 @@ public void logPotentiallyStuck(long processingTimeSeconds) { } private StringBuilder getStackTraceAsString() { - StringBuilder sb = new StringBuilder(2000); - for (StackTraceElement element : thread.getStackTrace()) { - sb.append(element).append('\n'); + String stack = stream(thread.getStackTrace()).map(Object::toString).collect(Collectors.joining("\n")); + StringBuilder sb = new StringBuilder(stack.length() + 2); + if (!stack.isEmpty()) { + sb.append(stack).append('\n'); } return sb; } public void handlePotentiallyStuck(Duration processingTime) { - boolean interrupt = false; - for (WorkflowExecutorListener listener : executorListeners) { + AtomicBoolean interrupt = new AtomicBoolean(false); + executorListeners.forEach(listener -> { try { if (listener.handlePotentiallyStuck(listenerContext, processingTime)) { - interrupt = true; + interrupt.set(true); } } catch (Throwable t) { logger.error("Error in " + listener.getClass().getName() + ".handleStuck (" + t.getMessage() + ")", t); } - } - if (interrupt) { + }); + if (interrupt.get()) { thread.interrupt(); } } diff --git a/nflow-engine/src/main/java/io/nflow/engine/internal/executor/WorkflowStateProcessorFactory.java b/nflow-engine/src/main/java/io/nflow/engine/internal/executor/WorkflowStateProcessorFactory.java index ace6cd31a..eb4760d25 100644 --- a/nflow-engine/src/main/java/io/nflow/engine/internal/executor/WorkflowStateProcessorFactory.java +++ b/nflow-engine/src/main/java/io/nflow/engine/internal/executor/WorkflowStateProcessorFactory.java @@ -4,6 +4,7 @@ import java.util.Map; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.atomic.AtomicInteger; import java.util.function.Supplier; import jakarta.inject.Inject; @@ -65,16 +66,16 @@ public WorkflowStateProcessor createProcessor(long instanceId, Supplier public int getPotentiallyStuckProcessors() { DateTime currentTime = now(); - int potentiallyStuck = 0; - for (WorkflowStateProcessor processor : processingInstances.values()) { + AtomicInteger potentiallyStuck = new AtomicInteger(0); + processingInstances.values().forEach(processor -> { Duration processingTime = new Duration(processor.getStartTime(), currentTime); long processingTimeSeconds = processingTime.getStandardSeconds(); if (processingTimeSeconds > stuckThreadThresholdSeconds) { - potentiallyStuck++; + potentiallyStuck.incrementAndGet(); processor.logPotentiallyStuck(processingTimeSeconds); processor.handlePotentiallyStuck(processingTime); } - } - return potentiallyStuck; + }); + return potentiallyStuck.get(); } } diff --git a/nflow-engine/src/main/java/io/nflow/engine/internal/workflow/ObjectStringMapper.java b/nflow-engine/src/main/java/io/nflow/engine/internal/workflow/ObjectStringMapper.java index 62bb47f7e..7360982e8 100644 --- a/nflow-engine/src/main/java/io/nflow/engine/internal/workflow/ObjectStringMapper.java +++ b/nflow-engine/src/main/java/io/nflow/engine/internal/workflow/ObjectStringMapper.java @@ -4,6 +4,7 @@ import java.lang.reflect.Constructor; import java.lang.reflect.InvocationTargetException; import java.lang.reflect.Type; +import java.util.Optional; import io.nflow.engine.config.EngineConfiguration.EngineObjectMapperSupplier; import jakarta.inject.Inject; @@ -81,22 +82,22 @@ public void storeArguments(StateExecution execution, continue; } Object value = args[i + 1]; - if (value == null) { - continue; - } - String sVal; - if (param.mutable) { - value = ((Mutable) value).val; - if (value == null) { - continue; + Optional.ofNullable(value).ifPresent(v -> { + Object actual = v; + if (param.mutable) { + actual = ((Mutable) actual).val; + if (actual == null) { + return; + } } - } - if (String.class.equals(param.type)) { - sVal = (String) value; - } else { - sVal = convertFromObject(param.key, value); - } - execution.setVariable(param.key, sVal); + String sVal; + if (String.class.equals(param.type)) { + sVal = (String) actual; + } else { + sVal = convertFromObject(param.key, actual); + } + execution.setVariable(param.key, sVal); + }); } } diff --git a/nflow-engine/src/main/java/io/nflow/engine/internal/workflow/WorkflowInstancePreProcessor.java b/nflow-engine/src/main/java/io/nflow/engine/internal/workflow/WorkflowInstancePreProcessor.java index 753137aa3..0d1d0d3b2 100644 --- a/nflow-engine/src/main/java/io/nflow/engine/internal/workflow/WorkflowInstancePreProcessor.java +++ b/nflow-engine/src/main/java/io/nflow/engine/internal/workflow/WorkflowInstancePreProcessor.java @@ -13,6 +13,8 @@ import io.nflow.engine.workflow.definition.WorkflowDefinition; import io.nflow.engine.workflow.instance.WorkflowInstance; +import java.util.Optional; + @Component public class WorkflowInstancePreProcessor { @@ -28,10 +30,8 @@ public WorkflowInstancePreProcessor(WorkflowDefinitionService workflowDefinition } public WorkflowInstance process(WorkflowInstance instance) { - WorkflowDefinition def = workflowDefinitionService.getWorkflowDefinition(instance.type); - if (def == null) { - throw new IllegalArgumentException("No workflow definition found for type [" + instance.type + "]"); - } + WorkflowDefinition def = Optional.ofNullable(workflowDefinitionService.getWorkflowDefinition(instance.type)) + .orElseThrow(() -> new IllegalArgumentException("No workflow definition found for type [" + instance.type + "]")); WorkflowInstance.Builder builder = new WorkflowInstance.Builder(instance); if (instance.state == null) { builder.setState(def.getInitialState()); diff --git a/nflow-rest-api-common/src/main/java/io/nflow/rest/v1/converter/CreateWorkflowConverter.java b/nflow-rest-api-common/src/main/java/io/nflow/rest/v1/converter/CreateWorkflowConverter.java index 2c373dbd4..cf2865d15 100644 --- a/nflow-rest-api-common/src/main/java/io/nflow/rest/v1/converter/CreateWorkflowConverter.java +++ b/nflow-rest-api-common/src/main/java/io/nflow/rest/v1/converter/CreateWorkflowConverter.java @@ -1,6 +1,7 @@ package io.nflow.rest.v1.converter; import static java.lang.Boolean.FALSE; +import static java.util.Optional.ofNullable; import static org.apache.commons.lang3.StringUtils.isNotEmpty; import java.util.Map.Entry; @@ -27,9 +28,7 @@ public WorkflowInstance convert(CreateWorkflowInstanceRequest req) { WorkflowInstance.Builder builder = factory.newWorkflowInstanceBuilder().setType(req.type).setBusinessKey(req.businessKey) .setExternalId(req.externalId); if (!FALSE.equals(req.activate)) { - if (req.activationTime != null) { - builder.setNextActivation(req.activationTime); - } + ofNullable(req.activationTime).ifPresent(builder::setNextActivation); } else { builder.setNextActivation(null); } @@ -37,14 +36,14 @@ public WorkflowInstance convert(CreateWorkflowInstanceRequest req) { if (isNotEmpty(req.startState)) { builder.setState(req.startState); } - for (Entry entry : req.stateVariables.entrySet()) { + req.stateVariables.entrySet().forEach(entry -> { Object value = entry.getValue(); if (value instanceof String) { builder.putStateVariable(entry.getKey(), (String) value); } else { builder.putStateVariable(entry.getKey(), value); } - } + }); return builder.build(); } diff --git a/nflow-rest-api-common/src/main/java/io/nflow/rest/v1/converter/ListWorkflowDefinitionConverter.java b/nflow-rest-api-common/src/main/java/io/nflow/rest/v1/converter/ListWorkflowDefinitionConverter.java index de9808812..7173bf61f 100644 --- a/nflow-rest-api-common/src/main/java/io/nflow/rest/v1/converter/ListWorkflowDefinitionConverter.java +++ b/nflow-rest-api-common/src/main/java/io/nflow/rest/v1/converter/ListWorkflowDefinitionConverter.java @@ -2,11 +2,8 @@ import static java.util.stream.Collectors.toMap; -import java.util.ArrayList; -import java.util.Collection; -import java.util.List; import java.util.Map; -import java.util.Map.Entry; +import java.util.Optional; import org.springframework.stereotype.Component; @@ -30,16 +27,11 @@ public ListWorkflowDefinitionResponse convert(WorkflowDefinition definition) { resp.description = definition.getDescription(); resp.onError = definition.getErrorState().name(); Map states = definition.getStates().stream().collect(toMap(WorkflowState::name, this::toState)); - for (Entry> entry : definition.getAllowedTransitions().entrySet()) { - State state = states.get(entry.getKey()); - state.transitions.addAll(entry.getValue()); - } - for (Entry entry : definition.getFailureTransitions().entrySet()) { - State state = states.get(entry.getKey()); - state.onFailure = entry.getValue().name(); - } - Collection values = states.values(); - resp.states = values.toArray(new State[values.size()]); + definition.getAllowedTransitions().forEach((key, targets) -> + Optional.ofNullable(states.get(key)).ifPresent(s -> s.transitions.addAll(targets))); + definition.getFailureTransitions().forEach((key, wfState) -> + Optional.ofNullable(states.get(key)).ifPresent(s -> s.onFailure = wfState.name())); + resp.states = states.values().toArray(new State[0]); WorkflowSettings workflowSettings = definition.getSettings(); TransitionDelays transitionDelays = new TransitionDelays(); @@ -72,14 +64,12 @@ public ListWorkflowDefinitionResponse convert(StoredWorkflowDefinition storedDef resp.type = storedDefinition.type; resp.description = storedDefinition.description; resp.onError = storedDefinition.onError; - List states = new ArrayList<>(storedDefinition.states.size()); - for (StoredWorkflowDefinition.State state : storedDefinition.states) { + resp.states = storedDefinition.states.stream().map(state -> { State tmp = new State(state.id, state.type, state.description); tmp.transitions.addAll(state.transitions); tmp.onFailure = state.onFailure; - states.add(tmp); - } - resp.states = states.toArray(new State[states.size()]); + return tmp; + }).toArray(size -> new State[size]); resp.supportedSignals = storedDefinition.supportedSignals.stream().map(s -> { Signal signal = new Signal(); signal.value = s.value; diff --git a/nflow-rest-api-common/src/main/java/io/nflow/rest/v1/converter/ListWorkflowInstanceConverter.java b/nflow-rest-api-common/src/main/java/io/nflow/rest/v1/converter/ListWorkflowInstanceConverter.java index 867a8287e..c7c36602f 100644 --- a/nflow-rest-api-common/src/main/java/io/nflow/rest/v1/converter/ListWorkflowInstanceConverter.java +++ b/nflow-rest-api-common/src/main/java/io/nflow/rest/v1/converter/ListWorkflowInstanceConverter.java @@ -4,6 +4,7 @@ import static io.nflow.rest.v1.ApiWorkflowInstanceInclude.actions; import static io.nflow.rest.v1.ApiWorkflowInstanceInclude.childWorkflows; import static io.nflow.rest.v1.ApiWorkflowInstanceInclude.currentStateVariables; +import static java.util.stream.Collectors.toList; import static java.util.stream.Collectors.toMap; import static org.springframework.util.CollectionUtils.isEmpty; @@ -57,16 +58,13 @@ public ListWorkflowInstanceResponse convert(WorkflowInstance instance, Set(); - for (WorkflowInstanceAction action : instance.actions) { - if (includes.contains(actionStateVariables)) { - resp.actions.add(new Action(action.id, action.type.name(), action.state, action.stateText, action.retryNo, - action.executionStart, action.executionEnd, action.executorId, stateVariablesToJson(action.updatedStateVariables))); - } else { - resp.actions.add(new Action(action.id, action.type.name(), action.state, action.stateText, action.retryNo, - action.executionStart, action.executionEnd, action.executorId)); - } - } + resp.actions = instance.actions.stream() + .map(action -> includes.contains(actionStateVariables) + ? new Action(action.id, action.type.name(), action.state, action.stateText, action.retryNo, + action.executionStart, action.executionEnd, action.executorId, stateVariablesToJson(action.updatedStateVariables)) + : new Action(action.id, action.type.name(), action.state, action.stateText, action.retryNo, + action.executionStart, action.executionEnd, action.executorId)) + .collect(toList()); } if (includes.contains(currentStateVariables)) { resp.stateVariables = stateVariablesToJson(instance.stateVariables); diff --git a/nflow-rest-api-common/src/main/java/io/nflow/rest/v1/converter/StatisticsConverter.java b/nflow-rest-api-common/src/main/java/io/nflow/rest/v1/converter/StatisticsConverter.java index c13768094..39c79d77d 100644 --- a/nflow-rest-api-common/src/main/java/io/nflow/rest/v1/converter/StatisticsConverter.java +++ b/nflow-rest-api-common/src/main/java/io/nflow/rest/v1/converter/StatisticsConverter.java @@ -28,10 +28,10 @@ public StatisticsResponse convert(Statistics stats) { public WorkflowDefinitionStatisticsResponse convert(Map> stats) { WorkflowDefinitionStatisticsResponse resp = new WorkflowDefinitionStatisticsResponse(); - for (Entry> entry : stats.entrySet()) { + stats.entrySet().forEach(entry -> { StateStatistics stateStats = new StateStatistics(); resp.stateStatistics.put(entry.getKey(), stateStats); - for (Entry statusEntry : entry.getValue().entrySet()) { + entry.getValue().entrySet().forEach(statusEntry -> { WorkflowDefinitionStatistics value = statusEntry.getValue(); switch (statusEntry.getKey()) { case "created": @@ -54,8 +54,8 @@ public WorkflowDefinitionStatisticsResponse convert(Map