feat: escalate overdue approval tasks via processDue
Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -16,6 +16,7 @@ import com.jetlumen.ordo.api.ProcessEventType;
|
||||
import com.jetlumen.ordo.api.ProcessInstance;
|
||||
import com.jetlumen.ordo.api.ProcessStatus;
|
||||
import com.jetlumen.ordo.api.RoutingCondition;
|
||||
import com.jetlumen.ordo.api.StepDue;
|
||||
import com.jetlumen.ordo.api.StepKind;
|
||||
import com.jetlumen.ordo.api.StepTransition;
|
||||
import com.jetlumen.ordo.api.TaskAction;
|
||||
@@ -195,7 +196,7 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
}
|
||||
Instant now = clock.instant();
|
||||
ApprovalTask updated = new ApprovalTask(task.id(), task.instanceId(), task.stepId(), task.name(),
|
||||
newAssignee, task.status(), task.createdAt(), task.completedAt(), task.action());
|
||||
newAssignee, task.status(), task.createdAt(), task.completedAt(), task.action(), task.dueAt());
|
||||
record(events, task.instanceId(), task.id(), task.stepId(), ProcessEventType.TASK_REASSIGNED, actor,
|
||||
newAssignee, now);
|
||||
return updated;
|
||||
@@ -228,7 +229,8 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
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);
|
||||
pending.name(), pending.assignee(), TaskStatus.SKIPPED, pending.createdAt(), now, action,
|
||||
pending.dueAt());
|
||||
if (taskRepository.completeIfPending(skipped)) {
|
||||
record(events, instanceId, pending.id(), pending.stepId(), ProcessEventType.TASK_SKIPPED, actor,
|
||||
comment, now);
|
||||
@@ -301,12 +303,106 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
return actionExecutionRepository.query(instanceId, pageRequest);
|
||||
}
|
||||
|
||||
@Override
|
||||
public synchronized int processDue(int limit) {
|
||||
if (limit <= 0) {
|
||||
throw new IllegalArgumentException("limit must be positive");
|
||||
}
|
||||
List<PendingAction> queued = new ArrayList<>();
|
||||
List<ProcessEvent> events = new ArrayList<>();
|
||||
Integer processed = transactionExecutor.execute(() -> {
|
||||
Instant now = clock.instant();
|
||||
int count = 0;
|
||||
for (ApprovalTask overdue : taskRepository.findDuePending(now, limit)) {
|
||||
if (escalateDueTask(overdue, now, queued, events)) {
|
||||
count++;
|
||||
}
|
||||
}
|
||||
return count;
|
||||
});
|
||||
finishCommittedWork(queued, events);
|
||||
return processed;
|
||||
}
|
||||
|
||||
private boolean escalateDueTask(ApprovalTask overdue, Instant now, List<PendingAction> queued,
|
||||
List<ProcessEvent> events) {
|
||||
if (!taskRepository.claimIfDue(overdue.id(), overdue.assignee(), now)) {
|
||||
return false;
|
||||
}
|
||||
ProcessInstance instance = requireInstance(overdue.instanceId());
|
||||
if (instance.status() != ProcessStatus.RUNNING) {
|
||||
return false;
|
||||
}
|
||||
ProcessDefinition definition = requireDefinition(instance.definitionId());
|
||||
ApprovalStep step = requireStep(definition, overdue.stepId());
|
||||
StepDue due = step.due();
|
||||
if (due == null) {
|
||||
record(events, overdue.instanceId(), overdue.id(), overdue.stepId(), ProcessEventType.TASK_ESCALATED, null,
|
||||
null, now);
|
||||
return true;
|
||||
}
|
||||
return switch (due.then()) {
|
||||
case REASSIGN -> escalateReassign(overdue, instance, step, due, now, events);
|
||||
case NOTIFY -> escalateNotify(overdue, instance, due, now, queued, events);
|
||||
case GOTO -> escalateGoto(overdue, instance, definition, due, now, queued, events);
|
||||
};
|
||||
}
|
||||
|
||||
private boolean escalateReassign(ApprovalTask overdue, ProcessInstance instance, ApprovalStep step, StepDue due,
|
||||
Instant now, List<ProcessEvent> events) {
|
||||
String newAssignee = assigneeResolver.resolve(due.to(), step, instance.context());
|
||||
requireText(newAssignee, "resolved assignee");
|
||||
boolean duplicatePending = taskRepository.findByInstanceIdAndStepId(overdue.instanceId(), overdue.stepId())
|
||||
.stream()
|
||||
.anyMatch(other -> !other.id().equals(overdue.id())
|
||||
&& other.status() == TaskStatus.PENDING
|
||||
&& other.assignee().equals(newAssignee));
|
||||
if (duplicatePending || overdue.assignee().equals(newAssignee)) {
|
||||
LOG.log(Level.WARNING, "due reassign skipped for task: " + overdue.id());
|
||||
record(events, overdue.instanceId(), overdue.id(), overdue.stepId(), ProcessEventType.TASK_ESCALATED, null,
|
||||
newAssignee, now);
|
||||
return true;
|
||||
}
|
||||
if (!taskRepository.reassignIfPending(overdue.id(), overdue.assignee(), newAssignee)) {
|
||||
return false;
|
||||
}
|
||||
record(events, overdue.instanceId(), overdue.id(), overdue.stepId(), ProcessEventType.TASK_ESCALATED, null,
|
||||
newAssignee, now);
|
||||
return true;
|
||||
}
|
||||
|
||||
private boolean escalateNotify(ApprovalTask overdue, ProcessInstance instance, StepDue due, Instant now,
|
||||
List<PendingAction> queued, List<ProcessEvent> events) {
|
||||
if (due.action() != null) {
|
||||
String executionId = nextId();
|
||||
actionExecutionRepository.insert(new ActionExecution(executionId, instance.id(), overdue.stepId(),
|
||||
due.action(), ActionExecutionStatus.PENDING, null,
|
||||
now.plusMillis(eventSequence.getAndIncrement()), null));
|
||||
queued.add(new PendingAction(executionId, due.action(), instance.id(), overdue.stepId(),
|
||||
instance.context()));
|
||||
}
|
||||
record(events, overdue.instanceId(), overdue.id(), overdue.stepId(), ProcessEventType.TASK_ESCALATED, null,
|
||||
due.action(), now);
|
||||
return true;
|
||||
}
|
||||
|
||||
private boolean escalateGoto(ApprovalTask overdue, ProcessInstance instance, ProcessDefinition definition,
|
||||
StepDue due, Instant now, List<PendingAction> queued, List<ProcessEvent> events) {
|
||||
List<ApprovalTask> siblings = taskRepository.findByInstanceIdAndStepId(instance.id(), overdue.stepId());
|
||||
skipPendingSiblings(siblings, "", null, now, events);
|
||||
enterStep(instance, definition, requireStep(definition, due.to()), now, queued, events);
|
||||
record(events, overdue.instanceId(), overdue.id(), overdue.stepId(), ProcessEventType.TASK_ESCALATED, null,
|
||||
due.to(), now);
|
||||
return true;
|
||||
}
|
||||
|
||||
private void createStepTasks(ProcessInstance instance, ApprovalStep step, Instant now, List<ProcessEvent> events) {
|
||||
for (String candidate : step.candidates()) {
|
||||
String assignee = assigneeResolver.resolve(candidate, step, instance.context());
|
||||
requireText(assignee, "resolved assignee");
|
||||
Instant dueAt = step.due() == null ? null : now.plus(step.due().after());
|
||||
ApprovalTask task = new ApprovalTask(nextId(), instance.id(), step.id(), step.name(), assignee,
|
||||
TaskStatus.PENDING, now, null, null);
|
||||
TaskStatus.PENDING, now, null, null, dueAt);
|
||||
taskRepository.save(task);
|
||||
record(events, instance.id(), task.id(), step.id(), ProcessEventType.TASK_CREATED, assignee, null, now);
|
||||
}
|
||||
@@ -328,7 +424,7 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
private ApprovalTask completeTask(ApprovalTask task, TaskStatus status, TaskAction action,
|
||||
List<ProcessEvent> events) {
|
||||
ApprovalTask completed = new ApprovalTask(task.id(), task.instanceId(), task.stepId(), task.name(),
|
||||
task.assignee(), status, task.createdAt(), action.operatedAt(), action);
|
||||
task.assignee(), status, task.createdAt(), action.operatedAt(), action, task.dueAt());
|
||||
if (!taskRepository.completeIfPending(completed)) {
|
||||
throw new TaskAlreadyCompletedException(task.id());
|
||||
}
|
||||
@@ -480,7 +576,8 @@ public final class DefaultOrdoEngine implements OrdoEngine {
|
||||
continue;
|
||||
}
|
||||
ApprovalTask skipped = new ApprovalTask(sibling.id(), sibling.instanceId(), sibling.stepId(),
|
||||
sibling.name(), sibling.assignee(), TaskStatus.SKIPPED, sibling.createdAt(), now, null);
|
||||
sibling.name(), sibling.assignee(), TaskStatus.SKIPPED, sibling.createdAt(), now, null,
|
||||
sibling.dueAt());
|
||||
// Best-effort: if another concurrent decision already completed this sibling, leave it as-is.
|
||||
if (taskRepository.completeIfPending(skipped)) {
|
||||
record(events, sibling.instanceId(), sibling.id(), sibling.stepId(), ProcessEventType.TASK_SKIPPED,
|
||||
|
||||
@@ -159,4 +159,9 @@ public final class InMemoryOrdoEngine implements OrdoEngine {
|
||||
public Page<ActionExecution> queryActionExecutions(String instanceId, PageRequest pageRequest) {
|
||||
return delegate.queryActionExecutions(instanceId, pageRequest);
|
||||
}
|
||||
|
||||
@Override
|
||||
public int processDue(int limit) {
|
||||
return delegate.processDue(limit);
|
||||
}
|
||||
}
|
||||
|
||||
+33
-1
@@ -9,10 +9,12 @@ import com.jetlumen.ordo.api.query.TaskQuery;
|
||||
import com.jetlumen.ordo.api.repository.ApprovalTaskRepository;
|
||||
import com.jetlumen.ordo.api.repository.ProcessInstanceRepository;
|
||||
|
||||
import java.time.Instant;
|
||||
import java.util.Comparator;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.Optional;
|
||||
|
||||
/** Development-only in-memory implementation of the task storage port. */
|
||||
@@ -92,7 +94,37 @@ public final class InMemoryApprovalTaskRepository implements ApprovalTaskReposit
|
||||
return false;
|
||||
}
|
||||
tasks.put(taskId, new ApprovalTask(current.id(), current.instanceId(), current.stepId(), current.name(),
|
||||
newAssignee, current.status(), current.createdAt(), current.completedAt(), current.action()));
|
||||
newAssignee, current.status(), current.createdAt(), current.completedAt(), current.action(),
|
||||
current.dueAt()));
|
||||
return true;
|
||||
}
|
||||
|
||||
@Override
|
||||
public synchronized List<ApprovalTask> findDuePending(Instant now, int limit) {
|
||||
Objects.requireNonNull(now, "now must not be null");
|
||||
if (limit <= 0) {
|
||||
throw new IllegalArgumentException("limit must be positive");
|
||||
}
|
||||
return tasks.values().stream()
|
||||
.filter(task -> task.status() == TaskStatus.PENDING)
|
||||
.filter(task -> task.dueAt() != null && !task.dueAt().isAfter(now))
|
||||
.sorted(Comparator.comparing(ApprovalTask::dueAt).thenComparing(ApprovalTask::id))
|
||||
.limit(limit)
|
||||
.toList();
|
||||
}
|
||||
|
||||
@Override
|
||||
public synchronized boolean claimIfDue(String taskId, String expectedAssignee, Instant now) {
|
||||
ApprovalTask current = tasks.get(taskId);
|
||||
if (current == null || current.status() != TaskStatus.PENDING
|
||||
|| !current.assignee().equals(expectedAssignee)
|
||||
|| current.dueAt() == null
|
||||
|| current.dueAt().isAfter(now)) {
|
||||
return false;
|
||||
}
|
||||
tasks.put(taskId, new ApprovalTask(current.id(), current.instanceId(), current.stepId(), current.name(),
|
||||
current.assignee(), current.status(), current.createdAt(), current.completedAt(), current.action(),
|
||||
null));
|
||||
return true;
|
||||
}
|
||||
|
||||
|
||||
@@ -13,6 +13,8 @@ import com.jetlumen.ordo.api.ProcessEventType;
|
||||
import com.jetlumen.ordo.api.ProcessInstance;
|
||||
import com.jetlumen.ordo.api.ProcessStatus;
|
||||
import com.jetlumen.ordo.api.RoutingCondition;
|
||||
import com.jetlumen.ordo.api.StepDue;
|
||||
import com.jetlumen.ordo.api.StepKind;
|
||||
import com.jetlumen.ordo.api.StepTransition;
|
||||
import com.jetlumen.ordo.api.TaskStatus;
|
||||
import com.jetlumen.ordo.api.exception.DefinitionAlreadyExistsException;
|
||||
@@ -33,6 +35,9 @@ import org.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.time.Clock;
|
||||
import java.time.Instant;
|
||||
import java.time.ZoneId;
|
||||
import java.time.ZoneOffset;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
|
||||
@@ -677,4 +682,103 @@ class InMemoryOrdoEngineTest {
|
||||
assertTrue(types.contains(ProcessEventType.ACTION_FAILED));
|
||||
assertTrue(types.contains(ProcessEventType.INSTANCE_APPROVED));
|
||||
}
|
||||
|
||||
@Test
|
||||
void processDueReassignsAfterTheStepDueElapses() {
|
||||
MutableClock clock = new MutableClock(Instant.parse("2026-01-15T09:00:00Z"));
|
||||
InMemoryOrdoEngine dueEngine = new InMemoryOrdoEngine(clock);
|
||||
dueEngine.register(ProcessDefinition.linear("leave-due", "Leave request", List.of(
|
||||
new ApprovalStep("manager", "Manager approval", List.of("maria"), ApprovalPolicy.ANY, StepKind.APPROVAL,
|
||||
null, StepDue.reassign(java.time.Duration.ofHours(1), "diana")))));
|
||||
|
||||
var instance = dueEngine.start("leave-due", "alice");
|
||||
ApprovalTask task = dueEngine.findPendingTasksByInstanceId(instance.id()).get(0);
|
||||
assertEquals("maria", task.assignee());
|
||||
assertEquals(0, dueEngine.processDue(10));
|
||||
|
||||
clock.set(Instant.parse("2026-01-15T10:00:00Z"));
|
||||
assertEquals(1, dueEngine.processDue(10));
|
||||
ApprovalTask escalated = dueEngine.findTask(task.id()).orElseThrow();
|
||||
assertEquals("diana", escalated.assignee());
|
||||
assertNull(escalated.dueAt());
|
||||
assertEquals(0, dueEngine.processDue(10));
|
||||
dueEngine.approve(task.id(), "diana");
|
||||
assertEquals(ProcessStatus.APPROVED, dueEngine.findInstance(instance.id()).orElseThrow().status());
|
||||
assertTrue(dueEngine.queryHistory(instance.id(), new PageRequest(0, 20)).content().stream()
|
||||
.anyMatch(event -> event.type() == ProcessEventType.TASK_ESCALATED && "diana".equals(event.detail())));
|
||||
}
|
||||
|
||||
@Test
|
||||
void processDueNotifiesWithoutChangingAssignee() {
|
||||
MutableClock clock = new MutableClock(Instant.parse("2026-01-15T09:00:00Z"));
|
||||
List<String> actions = new java.util.ArrayList<>();
|
||||
InMemoryOrdoEngine dueEngine = new InMemoryOrdoEngine(clock, AssigneeResolver.direct(),
|
||||
RoutingCondition.always(), (key, context) -> actions.add(key));
|
||||
dueEngine.register(ProcessDefinition.linear("leave-notify", "Leave request", List.of(
|
||||
new ApprovalStep("manager", "Manager approval", List.of("maria"), ApprovalPolicy.ANY, StepKind.APPROVAL,
|
||||
null, StepDue.notify(java.time.Duration.ofMinutes(30), "overdue-mail")))));
|
||||
|
||||
var instance = dueEngine.start("leave-notify", "alice");
|
||||
clock.set(Instant.parse("2026-01-15T09:30:00Z"));
|
||||
assertEquals(1, dueEngine.processDue(10));
|
||||
assertEquals("maria", dueEngine.findPendingTasksByInstanceId(instance.id()).get(0).assignee());
|
||||
assertEquals(List.of("overdue-mail"), actions);
|
||||
assertEquals(TaskStatus.PENDING, dueEngine.findPendingTasksByInstanceId(instance.id()).get(0).status());
|
||||
}
|
||||
|
||||
@Test
|
||||
void processDueGotoSkipsTheStepAndEntersTheTarget() {
|
||||
MutableClock clock = new MutableClock(Instant.parse("2026-01-15T09:00:00Z"));
|
||||
InMemoryOrdoEngine dueEngine = new InMemoryOrdoEngine(clock);
|
||||
dueEngine.register(new ProcessDefinition("leave-goto", "Leave request", List.of(
|
||||
new ApprovalStep("manager", "Manager approval", List.of("maria"), ApprovalPolicy.ANY, StepKind.APPROVAL,
|
||||
null, StepDue.gotoStep(java.time.Duration.ofHours(1), "hr")),
|
||||
ApprovalStep.single("hr", "HR approval", "henry")),
|
||||
List.of(
|
||||
StepTransition.always("manager", "hr"),
|
||||
StepTransition.end("hr"))));
|
||||
|
||||
var instance = dueEngine.start("leave-goto", "alice");
|
||||
clock.set(Instant.parse("2026-01-15T10:00:00Z"));
|
||||
assertEquals(1, dueEngine.processDue(10));
|
||||
assertEquals("hr", dueEngine.findPendingTasksByInstanceId(instance.id()).get(0).stepId());
|
||||
assertEquals("henry", dueEngine.findPendingTasksByInstanceId(instance.id()).get(0).assignee());
|
||||
assertEquals(TaskStatus.SKIPPED, dueEngine.findTasks(instance.id()).stream()
|
||||
.filter(task -> task.stepId().equals("manager"))
|
||||
.findFirst()
|
||||
.orElseThrow()
|
||||
.status());
|
||||
}
|
||||
|
||||
@Test
|
||||
void processDueRejectsNonPositiveLimit() {
|
||||
assertThrows(IllegalArgumentException.class, () -> engine.processDue(0));
|
||||
}
|
||||
|
||||
private static final class MutableClock extends Clock {
|
||||
private Instant instant;
|
||||
|
||||
private MutableClock(Instant instant) {
|
||||
this.instant = instant;
|
||||
}
|
||||
|
||||
void set(Instant instant) {
|
||||
this.instant = instant;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ZoneId getZone() {
|
||||
return ZoneOffset.UTC;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Clock withZone(ZoneId zone) {
|
||||
throw new UnsupportedOperationException();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Instant instant() {
|
||||
return instant;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+17
-2
@@ -87,13 +87,28 @@ class InMemoryApprovalTaskRepositoryTest {
|
||||
assertEquals("diana", repository.findById("task-1").orElseThrow().assignee());
|
||||
|
||||
ApprovalTask completed = new ApprovalTask(pending.id(), pending.instanceId(), pending.stepId(), pending.name(),
|
||||
"diana", TaskStatus.APPROVED, pending.createdAt(), CREATED_AT.plusSeconds(1), null);
|
||||
"diana", TaskStatus.APPROVED, pending.createdAt(), CREATED_AT.plusSeconds(1), null, null);
|
||||
assertTrue(repository.completeIfPending(completed));
|
||||
assertFalse(repository.reassignIfPending("task-1", "diana", "henry"));
|
||||
assertFalse(repository.reassignIfPending("missing", "maria", "diana"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void claimIfDueClearsDueAtOnlyForMatchingPendingTasks() {
|
||||
InMemoryApprovalTaskRepository repository = new InMemoryApprovalTaskRepository();
|
||||
Instant dueAt = CREATED_AT.plusSeconds(60);
|
||||
ApprovalTask pending = new ApprovalTask("task-1", "inst-1", "step", "Step", "maria", TaskStatus.PENDING,
|
||||
CREATED_AT, null, null, dueAt);
|
||||
repository.save(pending);
|
||||
|
||||
assertTrue(repository.findDuePending(dueAt, 10).contains(pending));
|
||||
assertTrue(repository.claimIfDue("task-1", "maria", dueAt));
|
||||
assertEquals(null, repository.findById("task-1").orElseThrow().dueAt());
|
||||
assertTrue(repository.findDuePending(dueAt, 10).isEmpty());
|
||||
assertFalse(repository.claimIfDue("task-1", "maria", dueAt));
|
||||
}
|
||||
|
||||
private static ApprovalTask task(String id, String instanceId, String assignee, TaskStatus status, Instant createdAt) {
|
||||
return new ApprovalTask(id, instanceId, "step", "Step", assignee, status, createdAt, null, null);
|
||||
return new ApprovalTask(id, instanceId, "step", "Step", assignee, status, createdAt, null, null, null);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user