feat: let initiators withdraw running process instances

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
0264408
2026-09-11 16:39:17 +08:00
co-authored by Cursor
parent 946c9d3c52
commit 87b449e45c
15 changed files with 213 additions and 4 deletions
@@ -24,6 +24,10 @@ public interface OrdoEngine {
return reject(taskId, actor, null); return reject(taskId, actor, null);
} }
ApprovalTask reject(String taskId, String actor, String comment); ApprovalTask reject(String taskId, String actor, String comment);
default ProcessInstance withdraw(String instanceId, String actor) {
return withdraw(instanceId, actor, null);
}
ProcessInstance withdraw(String instanceId, String actor, String comment);
Optional<ProcessInstance> findInstance(String instanceId); Optional<ProcessInstance> findInstance(String instanceId);
Optional<ApprovalTask> findTask(String taskId); Optional<ApprovalTask> findTask(String taskId);
List<ApprovalTask> findTasks(String instanceId); List<ApprovalTask> findTasks(String instanceId);
@@ -1,5 +1,5 @@
package com.jetlumen.ordo.api; package com.jetlumen.ordo.api;
public enum ProcessStatus { public enum ProcessStatus {
RUNNING, APPROVED, REJECTED RUNNING, APPROVED, REJECTED, WITHDRAWN
} }
@@ -0,0 +1,7 @@
package com.jetlumen.ordo.api.exception;
public final class InstanceAlreadyCompletedException extends OrdoException {
public InstanceAlreadyCompletedException(String instanceId) {
super("instance is already completed: " + instanceId);
}
}
@@ -0,0 +1,7 @@
package com.jetlumen.ordo.api.exception;
public final class InstanceNotFoundException extends OrdoException {
public InstanceNotFoundException(String instanceId) {
super("instance not found: " + instanceId);
}
}
@@ -0,0 +1,7 @@
package com.jetlumen.ordo.api.exception;
public final class UnauthorizedInstanceOperationException extends OrdoException {
public UnauthorizedInstanceOperationException(String instanceId, String actor) {
super("actor '" + actor + "' is not the initiator for instance: " + instanceId);
}
}
@@ -15,4 +15,11 @@ public interface ProcessInstanceRepository {
Optional<ProcessInstance> findById(String instanceId); Optional<ProcessInstance> findById(String instanceId);
boolean existsRunning(String definitionId); boolean existsRunning(String definitionId);
/**
* Completes the instance only if it is still running.
*
* @return true if the update was applied, false if the instance was missing or already terminal
*/
boolean completeIfRunning(ProcessInstance completed);
} }
@@ -17,9 +17,12 @@ 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.DefinitionInUseException;
import com.jetlumen.ordo.api.exception.DefinitionNotFoundException; import com.jetlumen.ordo.api.exception.DefinitionNotFoundException;
import com.jetlumen.ordo.api.exception.InstanceAlreadyCompletedException;
import com.jetlumen.ordo.api.exception.InstanceNotFoundException;
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.api.exception.TaskNotFoundException; import com.jetlumen.ordo.api.exception.TaskNotFoundException;
import com.jetlumen.ordo.api.exception.UnauthorizedInstanceOperationException;
import com.jetlumen.ordo.api.exception.UnauthorizedTaskOperationException; import com.jetlumen.ordo.api.exception.UnauthorizedTaskOperationException;
import com.jetlumen.ordo.api.repository.ApprovalTaskRepository; import com.jetlumen.ordo.api.repository.ApprovalTaskRepository;
import com.jetlumen.ordo.api.repository.ProcessDefinitionRepository; import com.jetlumen.ordo.api.repository.ProcessDefinitionRepository;
@@ -123,6 +126,35 @@ public final class DefaultOrdoEngine implements OrdoEngine {
}); });
} }
@Override
public synchronized ProcessInstance withdraw(String instanceId, String actor, String comment) {
requireText(instanceId, "instance id");
requireText(actor, "actor");
return transactionExecutor.execute(() -> {
ProcessInstance instance = instanceRepository.findById(instanceId)
.orElseThrow(() -> new InstanceNotFoundException(instanceId));
if (!instance.initiator().equals(actor)) {
throw new UnauthorizedInstanceOperationException(instanceId, actor);
}
if (instance.status() != ProcessStatus.RUNNING) {
throw new InstanceAlreadyCompletedException(instanceId);
}
Instant now = clock.instant();
ProcessInstance withdrawn = new ProcessInstance(instance.id(), instance.definitionId(), instance.initiator(),
ProcessStatus.WITHDRAWN, instance.startedAt(), now, instance.context());
if (!instanceRepository.completeIfRunning(withdrawn)) {
throw new InstanceAlreadyCompletedException(instanceId);
}
TaskAction action = new TaskAction(actor, comment, now);
for (ApprovalTask pending : taskRepository.findPendingByInstanceId(instanceId)) {
ApprovalTask skipped = new ApprovalTask(pending.id(), pending.instanceId(), pending.stepId(),
pending.name(), pending.assignee(), TaskStatus.SKIPPED, pending.createdAt(), now, action);
taskRepository.completeIfPending(skipped);
}
return withdrawn;
});
}
@Override @Override
public synchronized Optional<ProcessInstance> findInstance(String instanceId) { public synchronized Optional<ProcessInstance> findInstance(String instanceId) {
return instanceRepository.findById(instanceId); return instanceRepository.findById(instanceId);
@@ -259,8 +291,11 @@ public final class DefaultOrdoEngine implements OrdoEngine {
} }
private void completeInstance(ProcessInstance instance, ProcessStatus status, Instant now) { private void completeInstance(ProcessInstance instance, ProcessStatus status, Instant now) {
instanceRepository.update(new ProcessInstance(instance.id(), instance.definitionId(), instance.initiator(), ProcessInstance completed = new ProcessInstance(instance.id(), instance.definitionId(), instance.initiator(),
status, instance.startedAt(), now, instance.context())); status, instance.startedAt(), now, instance.context());
if (!instanceRepository.completeIfRunning(completed)) {
throw new InstanceAlreadyCompletedException(instance.id());
}
} }
private ProcessDefinition requireDefinition(String definitionId) { private ProcessDefinition requireDefinition(String definitionId) {
@@ -75,6 +75,11 @@ public final class InMemoryOrdoEngine implements OrdoEngine {
return delegate.reject(taskId, actor, comment); return delegate.reject(taskId, actor, comment);
} }
@Override
public ProcessInstance withdraw(String instanceId, String actor, String comment) {
return delegate.withdraw(instanceId, actor, comment);
}
@Override @Override
public Optional<ProcessInstance> findInstance(String instanceId) { public Optional<ProcessInstance> findInstance(String instanceId) {
return delegate.findInstance(instanceId); return delegate.findInstance(instanceId);
@@ -33,4 +33,14 @@ public final class InMemoryProcessInstanceRepository implements ProcessInstanceR
.anyMatch(instance -> instance.definitionId().equals(definitionId) .anyMatch(instance -> instance.definitionId().equals(definitionId)
&& instance.status() == ProcessStatus.RUNNING); && instance.status() == ProcessStatus.RUNNING);
} }
@Override
public synchronized boolean completeIfRunning(ProcessInstance completed) {
ProcessInstance current = instances.get(completed.id());
if (current == null || current.status() != ProcessStatus.RUNNING) {
return false;
}
instances.put(completed.id(), completed);
return true;
}
} }
@@ -5,6 +5,7 @@ import com.jetlumen.ordo.api.ApprovalStep;
import com.jetlumen.ordo.api.ApprovalTask; import com.jetlumen.ordo.api.ApprovalTask;
import com.jetlumen.ordo.api.ProcessContext; import com.jetlumen.ordo.api.ProcessContext;
import com.jetlumen.ordo.api.ProcessDefinition; import com.jetlumen.ordo.api.ProcessDefinition;
import com.jetlumen.ordo.api.ProcessInstance;
import com.jetlumen.ordo.api.ProcessStatus; 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;
@@ -12,9 +13,12 @@ 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.DefinitionInUseException;
import com.jetlumen.ordo.api.exception.DefinitionNotFoundException; import com.jetlumen.ordo.api.exception.DefinitionNotFoundException;
import com.jetlumen.ordo.api.exception.InstanceAlreadyCompletedException;
import com.jetlumen.ordo.api.exception.InstanceNotFoundException;
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.api.exception.TaskNotFoundException; import com.jetlumen.ordo.api.exception.TaskNotFoundException;
import com.jetlumen.ordo.api.exception.UnauthorizedInstanceOperationException;
import com.jetlumen.ordo.api.exception.UnauthorizedTaskOperationException; import com.jetlumen.ordo.api.exception.UnauthorizedTaskOperationException;
import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
@@ -127,6 +131,41 @@ class InMemoryOrdoEngineTest {
assertEquals("diana", engine.findTasks(next.id()).getFirst().assignee()); assertEquals("diana", engine.findTasks(next.id()).getFirst().assignee());
} }
@Test
void initiatorCanWithdrawWhileALaterStepIsPending() {
var instance = engine.start("leave", "alice");
engine.approve(engine.findTasks(instance.id()).getFirst().id(), "maria");
ProcessInstance withdrawn = engine.withdraw(instance.id(), "alice", "changed plans");
assertEquals(ProcessStatus.WITHDRAWN, withdrawn.status());
assertEquals(ProcessStatus.WITHDRAWN, engine.findInstance(instance.id()).orElseThrow().status());
assertTrue(engine.findPendingTasksByInstanceId(instance.id()).isEmpty());
ApprovalTask hrTask = engine.findTasks(instance.id()).stream()
.filter(task -> task.stepId().equals("hr"))
.findFirst()
.orElseThrow();
assertEquals(TaskStatus.SKIPPED, hrTask.status());
assertEquals("alice", hrTask.action().actor());
assertEquals("changed plans", hrTask.action().comment());
assertEquals(TaskStatus.APPROVED, engine.findTasks(instance.id()).stream()
.filter(task -> task.stepId().equals("manager"))
.findFirst()
.orElseThrow()
.status());
}
@Test
void withdrawIsRejectedForNonInitiatorMissingAndCompletedInstances() {
var instance = engine.start("leave", "alice");
assertThrows(UnauthorizedInstanceOperationException.class,
() -> engine.withdraw(instance.id(), "mallory"));
assertThrows(InstanceNotFoundException.class, () -> engine.withdraw("missing", "alice"));
engine.reject(engine.findTasks(instance.id()).getFirst().id(), "maria");
assertThrows(InstanceAlreadyCompletedException.class, () -> engine.withdraw(instance.id(), "alice"));
}
@Test @Test
void rejectsBlankRuntimeArguments() { void rejectsBlankRuntimeArguments() {
assertThrows(IllegalArgumentException.class, () -> engine.start("leave", " ")); assertThrows(IllegalArgumentException.class, () -> engine.start("leave", " "));
@@ -3,6 +3,7 @@ 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.DefinitionInUseException;
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;
@@ -48,7 +49,11 @@ 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.replace(definition);
} catch (DefinitionInUseException ignored) {
// keep the stored graph while RUNNING instances still reference it
}
} }
} }
@@ -111,6 +111,24 @@ class OrdoJdbcAutoConfigurationTest {
}); });
} }
@Test
void keepsStoredDefinitionWhenReloadFindsRunningInstances() {
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"))));
engine.start("leave-request-routed", "alice");
context.getBean(OrdoDefinitionLoader.class).load();
ProcessDefinition definition = context.getBean(ProcessDefinitionRepository.class)
.findById("leave-request-routed")
.orElseThrow();
assertThat(definition.name()).isEqualTo("stale");
assertThat(definition.steps()).hasSize(1);
});
}
@Test @Test
void honoursUserDefinedRoutingCondition() { void honoursUserDefinedRoutingCondition() {
withDataSourceRunner.withUserConfiguration(CustomRoutingConditionConfig.class) withDataSourceRunner.withUserConfiguration(CustomRoutingConditionConfig.class)
@@ -19,6 +19,8 @@ public final class JdbcProcessInstanceRepository implements ProcessInstanceRepos
+ " VALUES (?, ?, ?, ?, ?, ?, ?)"; + " VALUES (?, ?, ?, ?, ?, ?, ?)";
private static final String UPDATE_INSTANCE = private static final String UPDATE_INSTANCE =
"UPDATE ordo_process_instance SET status = ?, finished_at = ? WHERE id = ?"; "UPDATE ordo_process_instance SET status = ?, finished_at = ? WHERE id = ?";
private static final String COMPLETE_IF_RUNNING =
"UPDATE ordo_process_instance SET status = ?, finished_at = ? WHERE id = ? AND status = 'RUNNING'";
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 = ?";
@@ -91,4 +93,18 @@ public final class JdbcProcessInstanceRepository implements ProcessInstanceRepos
connectionProvider.close(connection); connectionProvider.close(connection);
} }
} }
@Override
public boolean completeIfRunning(ProcessInstance completed) {
Objects.requireNonNull(completed, "completed must not be null");
Connection connection = connectionProvider.getConnection();
try (PreparedStatement update = connection.prepareStatement(COMPLETE_IF_RUNNING)) {
ProcessInstanceMapper.bindUpdate(update, completed);
return update.executeUpdate() == 1;
} catch (SQLException e) {
throw new JdbcStorageException("failed to complete instance: " + completed.id(), e);
} finally {
connectionProvider.close(connection);
}
}
} }
@@ -12,6 +12,8 @@ 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.DefinitionInUseException;
import com.jetlumen.ordo.api.exception.InstanceAlreadyCompletedException;
import com.jetlumen.ordo.api.exception.UnauthorizedInstanceOperationException;
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;
@@ -165,6 +167,34 @@ class JdbcOrdoEngineIntegrationTest {
"leave", "Leave request v2", List.of(ApprovalStep.single("director", "Director approval", "diana"))))); "leave", "Leave request v2", List.of(ApprovalStep.single("director", "Director approval", "diana")))));
} }
@Test
void initiatorCanWithdrawWhileALaterStepIsPending() {
ProcessInstance instance = engine.start("leave", "alice");
engine.approve(engine.findPendingTasksByInstanceId(instance.id()).getFirst().id(), "maria");
ProcessInstance withdrawn = engine.withdraw(instance.id(), "alice", "changed plans");
assertEquals(ProcessStatus.WITHDRAWN, withdrawn.status());
assertEquals(NOW, withdrawn.finishedAt());
assertTrue(engine.findPendingTasksByInstanceId(instance.id()).isEmpty());
ApprovalTask hrTask = engine.findTasks(instance.id()).stream()
.filter(task -> task.stepId().equals("hr"))
.findFirst()
.orElseThrow();
assertEquals(TaskStatus.SKIPPED, hrTask.status());
assertEquals("alice", hrTask.action().actor());
assertEquals("changed plans", hrTask.action().comment());
}
@Test
void withdrawIsRejectedForNonInitiatorAndCompletedInstances() {
ProcessInstance instance = engine.start("leave", "alice");
assertThrows(UnauthorizedInstanceOperationException.class,
() -> engine.withdraw(instance.id(), "mallory"));
engine.reject(engine.findPendingTasksByInstanceId(instance.id()).getFirst().id(), "maria");
assertThrows(InstanceAlreadyCompletedException.class, () -> engine.withdraw(instance.id(), "alice"));
}
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),
@@ -81,6 +81,25 @@ class JdbcProcessInstanceRepositoryTest {
assertFalse(repository.existsRunning("leave")); assertFalse(repository.existsRunning("leave"));
} }
@Test
void completeIfRunningOnlyUpdatesARunningInstance() {
ProcessInstance running = new ProcessInstance("inst-1", "leave", "alice", ProcessStatus.RUNNING,
STARTED_AT, null, ProcessContext.empty());
repository.insert(running);
Instant finishedAt = STARTED_AT.plusSeconds(30);
ProcessInstance withdrawn = new ProcessInstance("inst-1", "leave", "alice", ProcessStatus.WITHDRAWN,
STARTED_AT, finishedAt, ProcessContext.empty());
assertTrue(repository.completeIfRunning(withdrawn));
assertEquals(ProcessStatus.WITHDRAWN, repository.findById("inst-1").orElseThrow().status());
assertEquals(finishedAt, repository.findById("inst-1").orElseThrow().finishedAt());
assertFalse(repository.completeIfRunning(new ProcessInstance("inst-1", "leave", "alice",
ProcessStatus.APPROVED, STARTED_AT, finishedAt, ProcessContext.empty())));
assertFalse(repository.completeIfRunning(new ProcessInstance("missing", "leave", "alice",
ProcessStatus.WITHDRAWN, STARTED_AT, finishedAt, ProcessContext.empty())));
}
@Test @Test
void updateOfAnUnknownInstanceFails() { void updateOfAnUnknownInstanceFails() {
assertThrows(IllegalStateException.class, () -> repository.update(new ProcessInstance("missing", "leave", assertThrows(IllegalStateException.class, () -> repository.update(new ProcessInstance("missing", "leave",