feat: replace process definition graphs when no instance is running
Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -6,6 +6,12 @@ import java.util.Optional;
|
|||||||
/** Public entry point for definition registration and approval operations. */
|
/** Public entry point for definition registration and approval operations. */
|
||||||
public interface OrdoEngine {
|
public interface OrdoEngine {
|
||||||
void register(ProcessDefinition definition);
|
void register(ProcessDefinition definition);
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Inserts the definition, or replaces its entire graph when the id already exists.
|
||||||
|
* Refuses replacement while any instance of this definition is {@code RUNNING}.
|
||||||
|
*/
|
||||||
|
void replace(ProcessDefinition definition);
|
||||||
default ProcessInstance start(String definitionId, String initiator) {
|
default ProcessInstance start(String definitionId, String initiator) {
|
||||||
return start(definitionId, initiator, ProcessContext.empty());
|
return start(definitionId, initiator, ProcessContext.empty());
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,7 @@
|
|||||||
|
package com.jetlumen.ordo.api.exception;
|
||||||
|
|
||||||
|
public final class DefinitionInUseException extends OrdoException {
|
||||||
|
public DefinitionInUseException(String definitionId) {
|
||||||
|
super("definition in use: " + definitionId);
|
||||||
|
}
|
||||||
|
}
|
||||||
+7
-1
@@ -4,7 +4,7 @@ import com.jetlumen.ordo.api.ProcessDefinition;
|
|||||||
|
|
||||||
import java.util.Optional;
|
import java.util.Optional;
|
||||||
|
|
||||||
/** Storage port for immutable process definitions. */
|
/** Storage port for process definitions. */
|
||||||
public interface ProcessDefinitionRepository {
|
public interface ProcessDefinitionRepository {
|
||||||
/**
|
/**
|
||||||
* Inserts the definition if no definition with the same id exists.
|
* Inserts the definition if no definition with the same id exists.
|
||||||
@@ -13,5 +13,11 @@ public interface ProcessDefinitionRepository {
|
|||||||
*/
|
*/
|
||||||
boolean insertIfAbsent(ProcessDefinition definition);
|
boolean insertIfAbsent(ProcessDefinition definition);
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Inserts the definition, or replaces its name, steps, candidates and transitions
|
||||||
|
* when the id already exists.
|
||||||
|
*/
|
||||||
|
void upsert(ProcessDefinition definition);
|
||||||
|
|
||||||
Optional<ProcessDefinition> findById(String definitionId);
|
Optional<ProcessDefinition> findById(String definitionId);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -13,4 +13,6 @@ public interface ProcessInstanceRepository {
|
|||||||
void update(ProcessInstance instance);
|
void update(ProcessInstance instance);
|
||||||
|
|
||||||
Optional<ProcessInstance> findById(String instanceId);
|
Optional<ProcessInstance> findById(String instanceId);
|
||||||
|
|
||||||
|
boolean existsRunning(String definitionId);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -15,6 +15,7 @@ import com.jetlumen.ordo.api.TaskAction;
|
|||||||
import com.jetlumen.ordo.api.TaskStatus;
|
import com.jetlumen.ordo.api.TaskStatus;
|
||||||
import com.jetlumen.ordo.api.TransactionExecutor;
|
import com.jetlumen.ordo.api.TransactionExecutor;
|
||||||
import com.jetlumen.ordo.api.exception.DefinitionAlreadyExistsException;
|
import com.jetlumen.ordo.api.exception.DefinitionAlreadyExistsException;
|
||||||
|
import com.jetlumen.ordo.api.exception.DefinitionInUseException;
|
||||||
import com.jetlumen.ordo.api.exception.DefinitionNotFoundException;
|
import com.jetlumen.ordo.api.exception.DefinitionNotFoundException;
|
||||||
import com.jetlumen.ordo.api.exception.NoRouteFoundException;
|
import com.jetlumen.ordo.api.exception.NoRouteFoundException;
|
||||||
import com.jetlumen.ordo.api.exception.TaskAlreadyCompletedException;
|
import com.jetlumen.ordo.api.exception.TaskAlreadyCompletedException;
|
||||||
@@ -73,6 +74,18 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public synchronized void replace(ProcessDefinition definition) {
|
||||||
|
Objects.requireNonNull(definition, "definition must not be null");
|
||||||
|
transactionExecutor.execute(() -> {
|
||||||
|
if (instanceRepository.existsRunning(definition.id())) {
|
||||||
|
throw new DefinitionInUseException(definition.id());
|
||||||
|
}
|
||||||
|
definitionRepository.upsert(definition);
|
||||||
|
return null;
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public synchronized ProcessInstance start(String definitionId, String initiator, ProcessContext context) {
|
public synchronized ProcessInstance start(String definitionId, String initiator, ProcessContext context) {
|
||||||
requireText(initiator, "initiator");
|
requireText(initiator, "initiator");
|
||||||
|
|||||||
@@ -55,6 +55,11 @@ public final class InMemoryOrdoEngine implements OrdoEngine {
|
|||||||
delegate.register(definition);
|
delegate.register(definition);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void replace(ProcessDefinition definition) {
|
||||||
|
delegate.replace(definition);
|
||||||
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public ProcessInstance start(String definitionId, String initiator, ProcessContext context) {
|
public ProcessInstance start(String definitionId, String initiator, ProcessContext context) {
|
||||||
return delegate.start(definitionId, initiator, context);
|
return delegate.start(definitionId, initiator, context);
|
||||||
|
|||||||
+5
@@ -16,6 +16,11 @@ public final class InMemoryProcessDefinitionRepository implements ProcessDefinit
|
|||||||
return definitions.putIfAbsent(definition.id(), definition) == null;
|
return definitions.putIfAbsent(definition.id(), definition) == null;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public synchronized void upsert(ProcessDefinition definition) {
|
||||||
|
definitions.put(definition.id(), definition);
|
||||||
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public synchronized Optional<ProcessDefinition> findById(String definitionId) {
|
public synchronized Optional<ProcessDefinition> findById(String definitionId) {
|
||||||
return Optional.ofNullable(definitions.get(definitionId));
|
return Optional.ofNullable(definitions.get(definitionId));
|
||||||
|
|||||||
+8
@@ -1,6 +1,7 @@
|
|||||||
package com.jetlumen.ordo.core.repository;
|
package com.jetlumen.ordo.core.repository;
|
||||||
|
|
||||||
import com.jetlumen.ordo.api.ProcessInstance;
|
import com.jetlumen.ordo.api.ProcessInstance;
|
||||||
|
import com.jetlumen.ordo.api.ProcessStatus;
|
||||||
import com.jetlumen.ordo.api.repository.ProcessInstanceRepository;
|
import com.jetlumen.ordo.api.repository.ProcessInstanceRepository;
|
||||||
|
|
||||||
import java.util.HashMap;
|
import java.util.HashMap;
|
||||||
@@ -25,4 +26,11 @@ public final class InMemoryProcessInstanceRepository implements ProcessInstanceR
|
|||||||
public synchronized Optional<ProcessInstance> findById(String instanceId) {
|
public synchronized Optional<ProcessInstance> findById(String instanceId) {
|
||||||
return Optional.ofNullable(instances.get(instanceId));
|
return Optional.ofNullable(instances.get(instanceId));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public synchronized boolean existsRunning(String definitionId) {
|
||||||
|
return instances.values().stream()
|
||||||
|
.anyMatch(instance -> instance.definitionId().equals(definitionId)
|
||||||
|
&& instance.status() == ProcessStatus.RUNNING);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ import com.jetlumen.ordo.api.RoutingCondition;
|
|||||||
import com.jetlumen.ordo.api.StepTransition;
|
import com.jetlumen.ordo.api.StepTransition;
|
||||||
import com.jetlumen.ordo.api.TaskStatus;
|
import com.jetlumen.ordo.api.TaskStatus;
|
||||||
import com.jetlumen.ordo.api.exception.DefinitionAlreadyExistsException;
|
import com.jetlumen.ordo.api.exception.DefinitionAlreadyExistsException;
|
||||||
|
import com.jetlumen.ordo.api.exception.DefinitionInUseException;
|
||||||
import com.jetlumen.ordo.api.exception.DefinitionNotFoundException;
|
import com.jetlumen.ordo.api.exception.DefinitionNotFoundException;
|
||||||
import com.jetlumen.ordo.api.exception.NoRouteFoundException;
|
import com.jetlumen.ordo.api.exception.NoRouteFoundException;
|
||||||
import com.jetlumen.ordo.api.exception.TaskAlreadyCompletedException;
|
import com.jetlumen.ordo.api.exception.TaskAlreadyCompletedException;
|
||||||
@@ -86,6 +87,46 @@ class InMemoryOrdoEngineTest {
|
|||||||
"leave", "Another leave request", List.of(ApprovalStep.single("lead", "Lead approval", "lee")))));
|
"leave", "Another leave request", List.of(ApprovalStep.single("lead", "Lead approval", "lee")))));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void replaceInsertsWhenTheDefinitionIsMissing() {
|
||||||
|
InMemoryOrdoEngine empty = new InMemoryOrdoEngine();
|
||||||
|
empty.replace(ProcessDefinition.linear("expense", "Expense request", List.of(
|
||||||
|
ApprovalStep.single("director", "Director approval", "diana"))));
|
||||||
|
|
||||||
|
var instance = empty.start("expense", "alice");
|
||||||
|
assertEquals("diana", empty.findTasks(instance.id()).getFirst().assignee());
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void replaceSwapsTheGraphWhenNoInstanceIsRunning() {
|
||||||
|
engine.replace(ProcessDefinition.linear("leave", "Leave request v2", List.of(
|
||||||
|
ApprovalStep.single("director", "Director approval", "diana"))));
|
||||||
|
|
||||||
|
var instance = engine.start("leave", "alice");
|
||||||
|
ApprovalTask task = engine.findTasks(instance.id()).getFirst();
|
||||||
|
assertEquals("director", task.stepId());
|
||||||
|
assertEquals("diana", task.assignee());
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void replaceIsRejectedWhileAnInstanceIsRunning() {
|
||||||
|
engine.start("leave", "alice");
|
||||||
|
assertThrows(DefinitionInUseException.class, () -> engine.replace(ProcessDefinition.linear(
|
||||||
|
"leave", "Leave request v2", List.of(ApprovalStep.single("director", "Director approval", "diana")))));
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void replaceSucceedsAfterInstancesReachATerminalStatus() {
|
||||||
|
var instance = engine.start("leave", "alice");
|
||||||
|
engine.reject(engine.findTasks(instance.id()).getFirst().id(), "maria");
|
||||||
|
|
||||||
|
engine.replace(ProcessDefinition.linear("leave", "Leave request v2", List.of(
|
||||||
|
ApprovalStep.single("director", "Director approval", "diana"))));
|
||||||
|
|
||||||
|
var next = engine.start("leave", "bob");
|
||||||
|
assertEquals("diana", engine.findTasks(next.id()).getFirst().assignee());
|
||||||
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
void rejectsBlankRuntimeArguments() {
|
void rejectsBlankRuntimeArguments() {
|
||||||
assertThrows(IllegalArgumentException.class, () -> engine.start("leave", " "));
|
assertThrows(IllegalArgumentException.class, () -> engine.start("leave", " "));
|
||||||
|
|||||||
+2
-7
@@ -3,7 +3,6 @@ package com.jetlumen.ordo.spring;
|
|||||||
import com.jetlumen.ordo.api.OrdoEngine;
|
import com.jetlumen.ordo.api.OrdoEngine;
|
||||||
import com.jetlumen.ordo.api.ProcessDefinition;
|
import com.jetlumen.ordo.api.ProcessDefinition;
|
||||||
import com.jetlumen.ordo.api.ProcessDefinitionParser;
|
import com.jetlumen.ordo.api.ProcessDefinitionParser;
|
||||||
import com.jetlumen.ordo.api.exception.DefinitionAlreadyExistsException;
|
|
||||||
import org.springframework.boot.ApplicationArguments;
|
import org.springframework.boot.ApplicationArguments;
|
||||||
import org.springframework.boot.ApplicationRunner;
|
import org.springframework.boot.ApplicationRunner;
|
||||||
import org.springframework.core.io.Resource;
|
import org.springframework.core.io.Resource;
|
||||||
@@ -13,7 +12,7 @@ import org.springframework.core.io.support.ResourcePatternResolver;
|
|||||||
import java.io.InputStream;
|
import java.io.InputStream;
|
||||||
import java.util.Objects;
|
import java.util.Objects;
|
||||||
|
|
||||||
/** Loads structural JSON process definitions and registers them if absent. */
|
/** Loads structural JSON process definitions and upserts them when no instance is running. */
|
||||||
public final class OrdoDefinitionLoader implements ApplicationRunner {
|
public final class OrdoDefinitionLoader implements ApplicationRunner {
|
||||||
private final OrdoEngine engine;
|
private final OrdoEngine engine;
|
||||||
private final String location;
|
private final String location;
|
||||||
@@ -49,11 +48,7 @@ public final class OrdoDefinitionLoader implements ApplicationRunner {
|
|||||||
} catch (RuntimeException e) {
|
} catch (RuntimeException e) {
|
||||||
throw new IllegalStateException("failed to parse ordo definition from " + describe(resource), e);
|
throw new IllegalStateException("failed to parse ordo definition from " + describe(resource), e);
|
||||||
}
|
}
|
||||||
try {
|
engine.replace(definition);
|
||||||
engine.register(definition);
|
|
||||||
} catch (DefinitionAlreadyExistsException ignored) {
|
|
||||||
// classpath JSON is an import source; existing ids stay as stored
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+17
@@ -94,6 +94,23 @@ class OrdoJdbcAutoConfigurationTest {
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void reloadsClasspathJsonDefinitionsWhenNoInstanceIsRunning() {
|
||||||
|
withDataSourceRunner.run(context -> {
|
||||||
|
OrdoEngine engine = context.getBean(OrdoEngine.class);
|
||||||
|
engine.replace(ProcessDefinition.linear("leave-request-routed", "stale", List.of(
|
||||||
|
ApprovalStep.single("lead", "Lead approval", "lee"))));
|
||||||
|
|
||||||
|
context.getBean(OrdoDefinitionLoader.class).load();
|
||||||
|
|
||||||
|
ProcessDefinition definition = context.getBean(ProcessDefinitionRepository.class)
|
||||||
|
.findById("leave-request-routed")
|
||||||
|
.orElseThrow();
|
||||||
|
assertThat(definition.name()).isEqualTo("Leave request");
|
||||||
|
assertThat(definition.steps()).hasSize(2);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
void honoursUserDefinedRoutingCondition() {
|
void honoursUserDefinedRoutingCondition() {
|
||||||
withDataSourceRunner.withUserConfiguration(CustomRoutingConditionConfig.class)
|
withDataSourceRunner.withUserConfiguration(CustomRoutingConditionConfig.class)
|
||||||
|
|||||||
+98
-39
@@ -30,6 +30,14 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR
|
|||||||
"INSERT INTO ordo_step_candidate (definition_id, step_id, candidate, candidate_order) VALUES (?, ?, ?, ?)";
|
"INSERT INTO ordo_step_candidate (definition_id, step_id, candidate, candidate_order) VALUES (?, ?, ?, ?)";
|
||||||
private static final String INSERT_TRANSITION =
|
private static final String INSERT_TRANSITION =
|
||||||
"INSERT INTO ordo_step_transition (definition_id, from_step_id, to_step_id, condition_key, priority) VALUES (?, ?, ?, ?, ?)";
|
"INSERT INTO ordo_step_transition (definition_id, from_step_id, to_step_id, condition_key, priority) VALUES (?, ?, ?, ?, ?)";
|
||||||
|
private static final String UPDATE_DEFINITION_NAME =
|
||||||
|
"UPDATE ordo_process_definition SET name = ? WHERE id = ?";
|
||||||
|
private static final String DELETE_TRANSITIONS =
|
||||||
|
"DELETE FROM ordo_step_transition WHERE definition_id = ?";
|
||||||
|
private static final String DELETE_CANDIDATES =
|
||||||
|
"DELETE FROM ordo_step_candidate WHERE definition_id = ?";
|
||||||
|
private static final String DELETE_STEPS =
|
||||||
|
"DELETE FROM ordo_approval_step WHERE definition_id = ?";
|
||||||
private static final String SELECT_DEFINITION =
|
private static final String SELECT_DEFINITION =
|
||||||
"SELECT id, name FROM ordo_process_definition WHERE id = ?";
|
"SELECT id, name FROM ordo_process_definition WHERE id = ?";
|
||||||
private static final String SELECT_STEPS =
|
private static final String SELECT_STEPS =
|
||||||
@@ -51,47 +59,10 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR
|
|||||||
Objects.requireNonNull(definition, "definition must not be null");
|
Objects.requireNonNull(definition, "definition must not be null");
|
||||||
Connection connection = connectionProvider.getConnection();
|
Connection connection = connectionProvider.getConnection();
|
||||||
try {
|
try {
|
||||||
try (PreparedStatement insertDefinition = connection.prepareStatement(INSERT_DEFINITION)) {
|
if (!insertDefinitionRow(connection, definition)) {
|
||||||
insertDefinition.setString(1, definition.id());
|
|
||||||
insertDefinition.setString(2, definition.name());
|
|
||||||
insertDefinition.executeUpdate();
|
|
||||||
} catch (SQLException e) {
|
|
||||||
if (isDuplicateKey(e)) {
|
|
||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
throw new JdbcStorageException("failed to insert definition: " + definition.id(), e);
|
insertGraph(connection, definition);
|
||||||
}
|
|
||||||
int stepOrder = 0;
|
|
||||||
for (ApprovalStep step : definition.steps()) {
|
|
||||||
try (PreparedStatement insertStep = connection.prepareStatement(INSERT_STEP)) {
|
|
||||||
insertStep.setString(1, definition.id());
|
|
||||||
insertStep.setString(2, step.id());
|
|
||||||
insertStep.setString(3, step.name());
|
|
||||||
insertStep.setString(4, step.policy().name());
|
|
||||||
insertStep.setInt(5, stepOrder++);
|
|
||||||
insertStep.executeUpdate();
|
|
||||||
}
|
|
||||||
int candidateOrder = 0;
|
|
||||||
for (String candidate : step.candidates()) {
|
|
||||||
try (PreparedStatement insertCandidate = connection.prepareStatement(INSERT_CANDIDATE)) {
|
|
||||||
insertCandidate.setString(1, definition.id());
|
|
||||||
insertCandidate.setString(2, step.id());
|
|
||||||
insertCandidate.setString(3, candidate);
|
|
||||||
insertCandidate.setInt(4, candidateOrder++);
|
|
||||||
insertCandidate.executeUpdate();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
for (StepTransition transition : definition.transitions()) {
|
|
||||||
try (PreparedStatement insertTransition = connection.prepareStatement(INSERT_TRANSITION)) {
|
|
||||||
insertTransition.setString(1, definition.id());
|
|
||||||
insertTransition.setString(2, transition.fromStepId());
|
|
||||||
insertTransition.setString(3, transition.toStepId());
|
|
||||||
insertTransition.setString(4, transition.conditionKey());
|
|
||||||
insertTransition.setInt(5, transition.priority());
|
|
||||||
insertTransition.executeUpdate();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return true;
|
return true;
|
||||||
} catch (SQLException e) {
|
} catch (SQLException e) {
|
||||||
throw new JdbcStorageException("failed to insert definition: " + definition.id(), e);
|
throw new JdbcStorageException("failed to insert definition: " + definition.id(), e);
|
||||||
@@ -100,6 +71,23 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void upsert(ProcessDefinition definition) {
|
||||||
|
Objects.requireNonNull(definition, "definition must not be null");
|
||||||
|
Connection connection = connectionProvider.getConnection();
|
||||||
|
try {
|
||||||
|
if (!insertDefinitionRow(connection, definition)) {
|
||||||
|
updateDefinitionName(connection, definition);
|
||||||
|
deleteGraph(connection, definition.id());
|
||||||
|
}
|
||||||
|
insertGraph(connection, definition);
|
||||||
|
} catch (SQLException e) {
|
||||||
|
throw new JdbcStorageException("failed to upsert definition: " + definition.id(), e);
|
||||||
|
} finally {
|
||||||
|
connectionProvider.close(connection);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public Optional<ProcessDefinition> findById(String definitionId) {
|
public Optional<ProcessDefinition> findById(String definitionId) {
|
||||||
Connection connection = connectionProvider.getConnection();
|
Connection connection = connectionProvider.getConnection();
|
||||||
@@ -157,6 +145,77 @@ public final class JdbcProcessDefinitionRepository implements ProcessDefinitionR
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private static boolean insertDefinitionRow(Connection connection, ProcessDefinition definition) throws SQLException {
|
||||||
|
try (PreparedStatement insertDefinition = connection.prepareStatement(INSERT_DEFINITION)) {
|
||||||
|
insertDefinition.setString(1, definition.id());
|
||||||
|
insertDefinition.setString(2, definition.name());
|
||||||
|
insertDefinition.executeUpdate();
|
||||||
|
return true;
|
||||||
|
} catch (SQLException e) {
|
||||||
|
if (isDuplicateKey(e)) {
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
throw e;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private static void updateDefinitionName(Connection connection, ProcessDefinition definition) throws SQLException {
|
||||||
|
try (PreparedStatement update = connection.prepareStatement(UPDATE_DEFINITION_NAME)) {
|
||||||
|
update.setString(1, definition.name());
|
||||||
|
update.setString(2, definition.id());
|
||||||
|
update.executeUpdate();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private static void deleteGraph(Connection connection, String definitionId) throws SQLException {
|
||||||
|
try (PreparedStatement deleteTransitions = connection.prepareStatement(DELETE_TRANSITIONS)) {
|
||||||
|
deleteTransitions.setString(1, definitionId);
|
||||||
|
deleteTransitions.executeUpdate();
|
||||||
|
}
|
||||||
|
try (PreparedStatement deleteCandidates = connection.prepareStatement(DELETE_CANDIDATES)) {
|
||||||
|
deleteCandidates.setString(1, definitionId);
|
||||||
|
deleteCandidates.executeUpdate();
|
||||||
|
}
|
||||||
|
try (PreparedStatement deleteSteps = connection.prepareStatement(DELETE_STEPS)) {
|
||||||
|
deleteSteps.setString(1, definitionId);
|
||||||
|
deleteSteps.executeUpdate();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private static void insertGraph(Connection connection, ProcessDefinition definition) throws SQLException {
|
||||||
|
int stepOrder = 0;
|
||||||
|
for (ApprovalStep step : definition.steps()) {
|
||||||
|
try (PreparedStatement insertStep = connection.prepareStatement(INSERT_STEP)) {
|
||||||
|
insertStep.setString(1, definition.id());
|
||||||
|
insertStep.setString(2, step.id());
|
||||||
|
insertStep.setString(3, step.name());
|
||||||
|
insertStep.setString(4, step.policy().name());
|
||||||
|
insertStep.setInt(5, stepOrder++);
|
||||||
|
insertStep.executeUpdate();
|
||||||
|
}
|
||||||
|
int candidateOrder = 0;
|
||||||
|
for (String candidate : step.candidates()) {
|
||||||
|
try (PreparedStatement insertCandidate = connection.prepareStatement(INSERT_CANDIDATE)) {
|
||||||
|
insertCandidate.setString(1, definition.id());
|
||||||
|
insertCandidate.setString(2, step.id());
|
||||||
|
insertCandidate.setString(3, candidate);
|
||||||
|
insertCandidate.setInt(4, candidateOrder++);
|
||||||
|
insertCandidate.executeUpdate();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
for (StepTransition transition : definition.transitions()) {
|
||||||
|
try (PreparedStatement insertTransition = connection.prepareStatement(INSERT_TRANSITION)) {
|
||||||
|
insertTransition.setString(1, definition.id());
|
||||||
|
insertTransition.setString(2, transition.fromStepId());
|
||||||
|
insertTransition.setString(3, transition.toStepId());
|
||||||
|
insertTransition.setString(4, transition.conditionKey());
|
||||||
|
insertTransition.setInt(5, transition.priority());
|
||||||
|
insertTransition.executeUpdate();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
private static boolean isDuplicateKey(SQLException e) {
|
private static boolean isDuplicateKey(SQLException e) {
|
||||||
return "23505".equals(e.getSQLState());
|
return "23505".equals(e.getSQLState());
|
||||||
}
|
}
|
||||||
|
|||||||
+19
@@ -1,6 +1,7 @@
|
|||||||
package com.jetlumen.ordo.storage.jdbc;
|
package com.jetlumen.ordo.storage.jdbc;
|
||||||
|
|
||||||
import com.jetlumen.ordo.api.ProcessInstance;
|
import com.jetlumen.ordo.api.ProcessInstance;
|
||||||
|
import com.jetlumen.ordo.api.ProcessStatus;
|
||||||
import com.jetlumen.ordo.api.repository.ProcessInstanceRepository;
|
import com.jetlumen.ordo.api.repository.ProcessInstanceRepository;
|
||||||
import com.jetlumen.ordo.storage.jdbc.mapper.ProcessInstanceMapper;
|
import com.jetlumen.ordo.storage.jdbc.mapper.ProcessInstanceMapper;
|
||||||
|
|
||||||
@@ -21,6 +22,8 @@ public final class JdbcProcessInstanceRepository implements ProcessInstanceRepos
|
|||||||
private static final String SELECT_INSTANCE =
|
private static final String SELECT_INSTANCE =
|
||||||
"SELECT id, definition_id, initiator, status, context_json, started_at, finished_at"
|
"SELECT id, definition_id, initiator, status, context_json, started_at, finished_at"
|
||||||
+ " FROM ordo_process_instance WHERE id = ?";
|
+ " FROM ordo_process_instance WHERE id = ?";
|
||||||
|
private static final String EXISTS_RUNNING =
|
||||||
|
"SELECT 1 FROM ordo_process_instance WHERE definition_id = ? AND status = ? LIMIT 1";
|
||||||
|
|
||||||
private final JdbcConnectionProvider connectionProvider;
|
private final JdbcConnectionProvider connectionProvider;
|
||||||
|
|
||||||
@@ -72,4 +75,20 @@ public final class JdbcProcessInstanceRepository implements ProcessInstanceRepos
|
|||||||
connectionProvider.close(connection);
|
connectionProvider.close(connection);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean existsRunning(String definitionId) {
|
||||||
|
Connection connection = connectionProvider.getConnection();
|
||||||
|
try (PreparedStatement select = connection.prepareStatement(EXISTS_RUNNING)) {
|
||||||
|
select.setString(1, definitionId);
|
||||||
|
select.setString(2, ProcessStatus.RUNNING.name());
|
||||||
|
try (ResultSet resultSet = select.executeQuery()) {
|
||||||
|
return resultSet.next();
|
||||||
|
}
|
||||||
|
} catch (SQLException e) {
|
||||||
|
throw new JdbcStorageException("failed to check running instances for definition: " + definitionId, e);
|
||||||
|
} finally {
|
||||||
|
connectionProvider.close(connection);
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+17
@@ -11,6 +11,7 @@ import com.jetlumen.ordo.api.ProcessStatus;
|
|||||||
import com.jetlumen.ordo.api.RoutingCondition;
|
import com.jetlumen.ordo.api.RoutingCondition;
|
||||||
import com.jetlumen.ordo.api.StepTransition;
|
import com.jetlumen.ordo.api.StepTransition;
|
||||||
import com.jetlumen.ordo.api.TaskStatus;
|
import com.jetlumen.ordo.api.TaskStatus;
|
||||||
|
import com.jetlumen.ordo.api.exception.DefinitionInUseException;
|
||||||
import com.jetlumen.ordo.api.exception.NoRouteFoundException;
|
import com.jetlumen.ordo.api.exception.NoRouteFoundException;
|
||||||
import com.jetlumen.ordo.api.exception.TaskAlreadyCompletedException;
|
import com.jetlumen.ordo.api.exception.TaskAlreadyCompletedException;
|
||||||
import com.jetlumen.ordo.core.DefaultOrdoEngine;
|
import com.jetlumen.ordo.core.DefaultOrdoEngine;
|
||||||
@@ -148,6 +149,22 @@ class JdbcOrdoEngineIntegrationTest {
|
|||||||
assertEquals("henry", engine.findPendingTasksByInstanceId(instance.id()).getFirst().assignee());
|
assertEquals("henry", engine.findPendingTasksByInstanceId(instance.id()).getFirst().assignee());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void replaceSwapsTheGraphWhenNoInstanceIsRunning() {
|
||||||
|
engine.replace(ProcessDefinition.linear("leave", "Leave request v2", List.of(
|
||||||
|
ApprovalStep.single("director", "Director approval", "diana"))));
|
||||||
|
|
||||||
|
ProcessInstance instance = engine.start("leave", "alice");
|
||||||
|
assertEquals("diana", engine.findPendingTasksByInstanceId(instance.id()).getFirst().assignee());
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void replaceIsRejectedWhileAnInstanceIsRunning() {
|
||||||
|
engine.start("leave", "alice");
|
||||||
|
assertThrows(DefinitionInUseException.class, () -> engine.replace(ProcessDefinition.linear(
|
||||||
|
"leave", "Leave request v2", List.of(ApprovalStep.single("director", "Director approval", "diana")))));
|
||||||
|
}
|
||||||
|
|
||||||
private OrdoEngine newEngine(AssigneeResolver assigneeResolver) {
|
private OrdoEngine newEngine(AssigneeResolver assigneeResolver) {
|
||||||
return new DefaultOrdoEngine(Clock.fixed(NOW, ZoneOffset.UTC), assigneeResolver, RoutingCondition.always(),
|
return new DefaultOrdoEngine(Clock.fixed(NOW, ZoneOffset.UTC), assigneeResolver, RoutingCondition.always(),
|
||||||
new JdbcTransactionExecutor(connectionProvider),
|
new JdbcTransactionExecutor(connectionProvider),
|
||||||
|
|||||||
+14
@@ -39,6 +39,20 @@ class JdbcProcessDefinitionRepositoryTest {
|
|||||||
assertEquals("Leave request v1", repository.findById("leave").orElseThrow().name());
|
assertEquals("Leave request v1", repository.findById("leave").orElseThrow().name());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void upsertReplacesNameStepsCandidatesAndTransitions() {
|
||||||
|
repository.insertIfAbsent(ProcessDefinition.linear("leave", "Leave request v1", List.of(
|
||||||
|
ApprovalStep.single("manager", "Manager approval", "maria"),
|
||||||
|
ApprovalStep.single("hr", "HR approval", "henry"))));
|
||||||
|
|
||||||
|
ProcessDefinition replacement = new ProcessDefinition("leave", "Leave request v2", List.of(
|
||||||
|
ApprovalStep.single("director", "Director approval", "diana")),
|
||||||
|
List.of(StepTransition.end("director")));
|
||||||
|
repository.upsert(replacement);
|
||||||
|
|
||||||
|
assertEquals(replacement, repository.findById("leave").orElseThrow());
|
||||||
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
void returnsEmptyForAnUnknownDefinition() {
|
void returnsEmptyForAnUnknownDefinition() {
|
||||||
assertTrue(repository.findById("missing").isEmpty());
|
assertTrue(repository.findById("missing").isEmpty());
|
||||||
|
|||||||
+15
@@ -13,6 +13,7 @@ import java.util.List;
|
|||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
|
|
||||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||||
|
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||||
|
|
||||||
@@ -66,6 +67,20 @@ class JdbcProcessInstanceRepositoryTest {
|
|||||||
assertTrue(repository.findById("missing").isEmpty());
|
assertTrue(repository.findById("missing").isEmpty());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void existsRunningOnlyCountsRunningInstancesOfThatDefinition() {
|
||||||
|
assertFalse(repository.existsRunning("leave"));
|
||||||
|
|
||||||
|
repository.insert(new ProcessInstance("inst-1", "leave", "alice", ProcessStatus.RUNNING,
|
||||||
|
STARTED_AT, null, ProcessContext.empty()));
|
||||||
|
assertTrue(repository.existsRunning("leave"));
|
||||||
|
assertFalse(repository.existsRunning("other"));
|
||||||
|
|
||||||
|
repository.update(new ProcessInstance("inst-1", "leave", "alice", ProcessStatus.APPROVED,
|
||||||
|
STARTED_AT, STARTED_AT.plusSeconds(60), ProcessContext.empty()));
|
||||||
|
assertFalse(repository.existsRunning("leave"));
|
||||||
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
void updateOfAnUnknownInstanceFails() {
|
void updateOfAnUnknownInstanceFails() {
|
||||||
assertThrows(IllegalStateException.class, () -> repository.update(new ProcessInstance("missing", "leave",
|
assertThrows(IllegalStateException.class, () -> repository.update(new ProcessInstance("missing", "leave",
|
||||||
|
|||||||
Reference in New Issue
Block a user