From 48559343e207ff0ae2799d7c1722a758198cef48 Mon Sep 17 00:00:00 2001 From: matevip Date: Fri, 8 May 2026 15:04:16 +0800 Subject: [PATCH] feat(workflow): add await_approval pause / resume to the runtime --- .../vip/mate/workflow/runtime/StepResult.java | 24 +- .../workflow/runtime/WorkflowResumer.java | 210 ++++++++++++++++++ .../mate/workflow/runtime/WorkflowRunner.java | 44 +++- .../mode/AwaitApprovalStepAdapter.java | 81 +++++++ 4 files changed, 346 insertions(+), 13 deletions(-) create mode 100644 mateclaw-server/src/main/java/vip/mate/workflow/runtime/WorkflowResumer.java create mode 100644 mateclaw-server/src/main/java/vip/mate/workflow/runtime/mode/AwaitApprovalStepAdapter.java diff --git a/mateclaw-server/src/main/java/vip/mate/workflow/runtime/StepResult.java b/mateclaw-server/src/main/java/vip/mate/workflow/runtime/StepResult.java index acfa3362..b49ee83a 100644 --- a/mateclaw-server/src/main/java/vip/mate/workflow/runtime/StepResult.java +++ b/mateclaw-server/src/main/java/vip/mate/workflow/runtime/StepResult.java @@ -3,10 +3,11 @@ package vip.mate.workflow.runtime; /** * Outcome reported by a {@link StepAdapter#execute}. Records: * */ public record StepResult( @@ -24,20 +27,25 @@ public record StepResult( String outputContentType, Object outputValue, String outputSummary, - String errorMessage + String errorMessage, + String pauseToken ) { - public enum State { SUCCEEDED, SKIPPED, FAILED } + public enum State { SUCCEEDED, SKIPPED, FAILED, PAUSED } public static StepResult succeeded(String payloadUri, String contentType, Object value, String summary) { - return new StepResult(State.SUCCEEDED, payloadUri, contentType, value, summary, null); + return new StepResult(State.SUCCEEDED, payloadUri, contentType, value, summary, null, null); } public static StepResult skipped(String reason) { - return new StepResult(State.SKIPPED, null, null, null, reason, null); + return new StepResult(State.SKIPPED, null, null, null, reason, null, null); } public static StepResult failed(String errorMessage) { - return new StepResult(State.FAILED, null, null, null, null, errorMessage); + return new StepResult(State.FAILED, null, null, null, null, errorMessage, null); + } + + public static StepResult paused(String pauseToken, String summary) { + return new StepResult(State.PAUSED, null, null, null, summary, null, pauseToken); } } diff --git a/mateclaw-server/src/main/java/vip/mate/workflow/runtime/WorkflowResumer.java b/mateclaw-server/src/main/java/vip/mate/workflow/runtime/WorkflowResumer.java new file mode 100644 index 00000000..836a7bbb --- /dev/null +++ b/mateclaw-server/src/main/java/vip/mate/workflow/runtime/WorkflowResumer.java @@ -0,0 +1,210 @@ +package vip.mate.workflow.runtime; + +import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; +import com.fasterxml.jackson.databind.ObjectMapper; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Service; +import vip.mate.workflow.compiler.ir.WorkflowGraph; +import vip.mate.workflow.model.WorkflowRunEntity; +import vip.mate.workflow.model.WorkflowRunPauseEntity; +import vip.mate.workflow.model.WorkflowRunStepEntity; +import vip.mate.workflow.repository.WorkflowRunMapper; +import vip.mate.workflow.repository.WorkflowRunPauseMapper; +import vip.mate.workflow.repository.WorkflowRunStepMapper; + +import java.time.LocalDateTime; +import java.util.List; +import java.util.Map; + +/** + * Settles a paused workflow run. Callers (approval callbacks, timeout sweeper, + * REST endpoints) hand in a {@code pauseToken} and an outcome; the resumer + * marks the pause and the await_approval step row, hydrates a fresh + * {@link WorkflowRunContext} from the persisted step rows, and delegates back + * to {@link WorkflowRunner#continueFromIndex} for the post-pause tail. + * + *

Idempotent: a pause that has already been resumed yields + * {@link Outcome#alreadyResolved(long)} without touching DB or memory. The + * graph is loaded by the caller (typically via a revision-id lookup) since the + * resumer has no opinion on storage. + */ +@Slf4j +@Service +public class WorkflowResumer { + + private static final String STATE_SUCCEEDED = "succeeded"; + private static final String STATE_FAILED = "failed"; + + private final WorkflowRunMapper runMapper; + private final WorkflowRunStepMapper stepMapper; + private final WorkflowRunPauseMapper pauseMapper; + private final WorkflowRunner runner; + private final PayloadStore payloadStore; + private final ObjectMapper objectMapper; + + public WorkflowResumer(WorkflowRunMapper runMapper, + WorkflowRunStepMapper stepMapper, + WorkflowRunPauseMapper pauseMapper, + WorkflowRunner runner, + PayloadStore payloadStore, + ObjectMapper objectMapper) { + this.runMapper = runMapper; + this.stepMapper = stepMapper; + this.pauseMapper = pauseMapper; + this.runner = runner; + this.payloadStore = payloadStore; + this.objectMapper = objectMapper; + } + + public Outcome resume(WorkflowGraph graph, String pauseToken, + ResumeOutcome outcome, byte[] resumePayloadBody) { + WorkflowRunPauseEntity pause = pauseMapper.selectOne(new LambdaQueryWrapper() + .eq(WorkflowRunPauseEntity::getPauseToken, pauseToken)); + if (pause == null) { + return Outcome.notFound(pauseToken); + } + if (pause.getResumedAt() != null) { + return Outcome.alreadyResolved(pause.getRunId()); + } + + WorkflowRunEntity runRow = runMapper.selectById(pause.getRunId()); + if (runRow == null) { + return Outcome.notFound(pauseToken); + } + + WorkflowRunStepEntity stepRow = stepMapper.selectById(pause.getStepId()); + if (stepRow == null) { + return Outcome.notFound(pauseToken); + } + + // Persist the pause row before doing any further work so a crash mid-resume + // leaves a clear audit trail (the pause is settled even if the post-resume + // execution never started). + String resumePayloadRef = null; + if (resumePayloadBody != null && resumePayloadBody.length > 0) { + resumePayloadRef = payloadStore.storeBytes(runRow.getWorkspaceId(), + resumePayloadBody, "application/octet-stream"); + } + pause.setResumedAt(LocalDateTime.now()); + pause.setResumeOutcome(outcome.token()); + pause.setResumePayloadRef(resumePayloadRef); + pauseMapper.updateById(pause); + + // Settle the await_approval step row first. + stepRow.setState(outcome == ResumeOutcome.APPROVED ? STATE_SUCCEEDED : STATE_FAILED); + stepRow.setOutputSummary("resumed: " + outcome.token()); + stepRow.setCompletedAt(LocalDateTime.now()); + if (outcome != ResumeOutcome.APPROVED) { + stepRow.setErrorMessage("approval " + outcome.token()); + } + stepMapper.updateById(stepRow); + + if (outcome != ResumeOutcome.APPROVED) { + // Failed approval ends the run — no further steps. + runRow.setState(STATE_FAILED); + runRow.setErrorMessage("paused step '" + stepRow.getStepName() + "' " + outcome.token()); + runRow.setCompletedAt(LocalDateTime.now()); + runMapper.updateById(runRow); + return Outcome.failed(runRow.getId(), runRow.getErrorMessage()); + } + + // Hydrate the run context from prior step rows so post-resume steps can + // reference {{ outputs.xxx }} from steps that completed before the pause. + WorkflowRunContext ctx = hydrateContext(runRow, graph, stepRow.getStepIndex()); + String priorOutputRef = lastSucceededOutputRef(runRow.getId(), stepRow.getStepIndex()); + + WorkflowRunResult result = runner.continueFromIndex( + graph, ctx, runRow, stepRow.getStepIndex() + 1, priorOutputRef); + return Outcome.continued(result); + } + + private WorkflowRunContext hydrateContext(WorkflowRunEntity runRow, WorkflowGraph graph, + int pausedStepIndex) { + Map inputs = (runRow.getInitialInputRef() == null) + ? Map.of() + : payloadStore.readJson(runRow.getInitialInputRef(), Map.class); + WorkflowRunContext ctx = new WorkflowRunContext( + runRow.getId(), + runRow.getWorkspaceId(), + runRow.getWorkflowId(), + runRow.getRevisionId(), + inputs); + + // Replay the rolling outputs map: walk completed succeeded step rows + // up to the pause and put their parsed payloads back into the context + // under their declared outputVar. + List rows = stepMapper.selectList(new LambdaQueryWrapper() + .eq(WorkflowRunStepEntity::getRunId, runRow.getId()) + .lt(WorkflowRunStepEntity::getStepIndex, pausedStepIndex) + .orderByAsc(WorkflowRunStepEntity::getStepIndex) + .orderByAsc(WorkflowRunStepEntity::getIterationIndex)); + for (WorkflowRunStepEntity row : rows) { + if (!STATE_SUCCEEDED.equals(row.getState()) || row.getOutputRef() == null) continue; + int idx = row.getStepIndex(); + if (idx < 0 || idx >= graph.steps().size()) continue; + var step = graph.steps().get(idx); + if (step.outputVar() == null || step.outputVar().isBlank()) continue; + Object value = decodeOutput(row); + if (value != null) ctx.putOutput(step.outputVar(), value); + } + return ctx; + } + + private Object decodeOutput(WorkflowRunStepEntity row) { + try { + byte[] body = payloadStore.readBytes(row.getOutputRef()); + if ("json".equals(row.getOutputContentType())) { + return objectMapper.readValue(body, Object.class); + } + return new String(body, java.nio.charset.StandardCharsets.UTF_8); + } catch (Exception e) { + log.warn("Workflow resume: failed to decode prior step output ref={}: {}", + row.getOutputRef(), e.getMessage()); + return null; + } + } + + private String lastSucceededOutputRef(long runId, int beforeStepIndex) { + WorkflowRunStepEntity row = stepMapper.selectOne(new LambdaQueryWrapper() + .eq(WorkflowRunStepEntity::getRunId, runId) + .eq(WorkflowRunStepEntity::getState, STATE_SUCCEEDED) + .lt(WorkflowRunStepEntity::getStepIndex, beforeStepIndex) + .isNotNull(WorkflowRunStepEntity::getOutputRef) + .orderByDesc(WorkflowRunStepEntity::getStepIndex) + .orderByDesc(WorkflowRunStepEntity::getIterationIndex) + .last("LIMIT 1")); + return row == null ? null : row.getOutputRef(); + } + + /** Outcome label written to {@code mate_workflow_run_pause.resume_outcome}. */ + public enum ResumeOutcome { + APPROVED("approved"), + REJECTED("rejected"), + TIMEOUT("timeout"), + CANCELLED("cancelled"); + + private final String token; + + ResumeOutcome(String token) { this.token = token; } + + public String token() { return token; } + } + + /** Result of attempting a resume — exposes the final run state when completed inline. */ + public record Outcome(Kind kind, Long runId, WorkflowRunResult finalResult, String errorMessage) { + public enum Kind { CONTINUED, FAILED, ALREADY_RESOLVED, NOT_FOUND } + + public static Outcome continued(WorkflowRunResult r) { + return new Outcome(Kind.CONTINUED, r.runId(), r, null); + } + public static Outcome failed(long runId, String err) { + return new Outcome(Kind.FAILED, runId, null, err); + } + public static Outcome alreadyResolved(long runId) { + return new Outcome(Kind.ALREADY_RESOLVED, runId, null, null); + } + public static Outcome notFound(String token) { + return new Outcome(Kind.NOT_FOUND, null, null, "pause token not found: " + token); + } + } +} diff --git a/mateclaw-server/src/main/java/vip/mate/workflow/runtime/WorkflowRunner.java b/mateclaw-server/src/main/java/vip/mate/workflow/runtime/WorkflowRunner.java index 6c9feca5..1a7b8c54 100644 --- a/mateclaw-server/src/main/java/vip/mate/workflow/runtime/WorkflowRunner.java +++ b/mateclaw-server/src/main/java/vip/mate/workflow/runtime/WorkflowRunner.java @@ -40,6 +40,7 @@ public class WorkflowRunner { private static final String STATE_SUCCEEDED = "succeeded"; private static final String STATE_FAILED = "failed"; private static final String STATE_SKIPPED = "skipped"; + private static final String STATE_PAUSED = "paused"; private static final ExecutorService FAN_OUT_EXECUTOR = Executors.newVirtualThreadPerTaskExecutor(); @@ -61,21 +62,43 @@ public class WorkflowRunner { public WorkflowRunResult run(WorkflowGraph graph, WorkflowRunRequest request) { WorkflowRunEntity runRow = openRun(request); - long runId = runRow.getId(); String inputsRef = payloadStore.storeJson(request.workspaceId(), request.inputs()); runRow.setInitialInputRef(inputsRef); runMapper.updateById(runRow); WorkflowRunContext ctx = new WorkflowRunContext( - runId, + runRow.getId(), request.workspaceId(), request.workflowId(), request.revisionId(), request.inputs()); - String lastSucceededOutputRef = null; + return executeFromIndex(graph, ctx, runRow, /*fromIndex*/ 0, /*priorOutputRef*/ null); + } + + /** + * Continue an already-open run from {@code fromIndex}. Used by the resumer + * after a pause settles. {@code priorOutputRef} is the last successful + * step's output URI from before the pause — propagated so the + * {@code final_output_ref} on success still points at meaningful data when + * the post-resume tail of the run produces no further output. + */ + public WorkflowRunResult continueFromIndex(WorkflowGraph graph, WorkflowRunContext ctx, + WorkflowRunEntity runRow, int fromIndex, + String priorOutputRef) { + // Move the run row back to running so step-completion timestamps make + // sense and the GC sweeper does not see a stale paused row. + runRow.setState(STATE_RUNNING); + runMapper.updateById(runRow); + return executeFromIndex(graph, ctx, runRow, fromIndex, priorOutputRef); + } + + private WorkflowRunResult executeFromIndex(WorkflowGraph graph, WorkflowRunContext ctx, + WorkflowRunEntity runRow, int fromIndex, + String priorOutputRef) { + String lastSucceededOutputRef = priorOutputRef; try { - int i = 0; + int i = fromIndex; while (i < graph.steps().size()) { WorkflowStep step = graph.steps().get(i); int groupEnd = scanFanOutGroup(graph.steps(), i); @@ -91,6 +114,9 @@ public class WorkflowRunner { if (result.state() == StepResult.State.FAILED) { return finishFailed(runRow, result.errorMessage()); } + if (result.state() == StepResult.State.PAUSED) { + return finishPaused(runRow, result.pauseToken()); + } if (result.outputPayloadUri() != null) { lastSucceededOutputRef = result.outputPayloadUri(); } @@ -99,7 +125,7 @@ public class WorkflowRunner { } return finishSucceeded(runRow, lastSucceededOutputRef); } catch (RuntimeException e) { - log.error("Workflow run {} aborted by unexpected exception", runId, e); + log.error("Workflow run {} aborted by unexpected exception", ctx.runId(), e); return finishFailed(runRow, "runtime error: " + e.getMessage()); } } @@ -226,6 +252,13 @@ public class WorkflowRunner { return new WorkflowRunResult(runRow.getId(), STATE_FAILED, null, errorMessage); } + private WorkflowRunResult finishPaused(WorkflowRunEntity runRow, String pauseToken) { + runRow.setState(STATE_PAUSED); + // Pause leaves the run open — completedAt stays null until resume settles it. + runMapper.updateById(runRow); + return new WorkflowRunResult(runRow.getId(), STATE_PAUSED, null, "pauseToken=" + pauseToken); + } + private WorkflowRunStepEntity openStep(long runId, int stepIndex, Integer iterationIndex, WorkflowStep step) { WorkflowRunStepEntity row = new WorkflowRunStepEntity(); @@ -246,6 +279,7 @@ public class WorkflowRunner { case SUCCEEDED -> row.setState(STATE_SUCCEEDED); case SKIPPED -> row.setState(STATE_SKIPPED); case FAILED -> row.setState(STATE_FAILED); + case PAUSED -> row.setState(STATE_PAUSED); } row.setOutputRef(result.outputPayloadUri()); if (result.outputContentType() != null) { diff --git a/mateclaw-server/src/main/java/vip/mate/workflow/runtime/mode/AwaitApprovalStepAdapter.java b/mateclaw-server/src/main/java/vip/mate/workflow/runtime/mode/AwaitApprovalStepAdapter.java new file mode 100644 index 00000000..8e236855 --- /dev/null +++ b/mateclaw-server/src/main/java/vip/mate/workflow/runtime/mode/AwaitApprovalStepAdapter.java @@ -0,0 +1,81 @@ +package vip.mate.workflow.runtime.mode; + +import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; +import org.springframework.stereotype.Component; +import vip.mate.workflow.compiler.ir.StepMode; +import vip.mate.workflow.compiler.ir.WorkflowStep; +import vip.mate.workflow.model.WorkflowRunPauseEntity; +import vip.mate.workflow.model.WorkflowRunStepEntity; +import vip.mate.workflow.repository.WorkflowRunPauseMapper; +import vip.mate.workflow.repository.WorkflowRunStepMapper; +import vip.mate.workflow.runtime.StepAdapter; +import vip.mate.workflow.runtime.StepResult; +import vip.mate.workflow.runtime.WorkflowRunContext; + +import java.time.LocalDateTime; +import java.util.UUID; + +/** + * {@code await_approval} — pauses the run pending an external approval + * decision. Inserts a {@code mate_workflow_run_pause} row keyed by a fresh + * {@code pauseToken}, then returns {@link StepResult.State#PAUSED} so the + * runner can short-circuit and mark the run row {@code paused}. Resume is + * orchestrated by {@code WorkflowResumer} once the external approval (or + * timeout) lands. + * + *

v0 stores no extra approval metadata — the workflow's {@code mate_tool_approval} + * link will be wired in once the approval-driven resume callback exists. + * The pause row's {@code resume_deadline} is honoured when the step + * declares a {@code timeoutSecs}; otherwise it stays {@code null} and the + * resumer treats the pause as open-ended. + */ +@Component +public class AwaitApprovalStepAdapter implements StepAdapter { + + private final WorkflowRunPauseMapper pauseMapper; + private final WorkflowRunStepMapper stepMapper; + + public AwaitApprovalStepAdapter(WorkflowRunPauseMapper pauseMapper, + WorkflowRunStepMapper stepMapper) { + this.pauseMapper = pauseMapper; + this.stepMapper = stepMapper; + } + + @Override + public String typeName() { return "await_approval"; } + + @Override + public StepResult execute(WorkflowStep step, WorkflowRunContext context) { + if (!(step.mode() instanceof StepMode.AwaitApproval cfg)) { + return StepResult.failed("await_approval adapter received non-await mode: " + + step.mode().typeName()); + } + + // Look up the freshly opened step row so we can link the pause to it. + WorkflowRunStepEntity stepRow = stepMapper.selectOne(new LambdaQueryWrapper() + .eq(WorkflowRunStepEntity::getRunId, context.runId()) + .eq(WorkflowRunStepEntity::getStepName, step.name()) + .orderByDesc(WorkflowRunStepEntity::getId) + .last("LIMIT 1")); + if (stepRow == null) { + return StepResult.failed("await_approval could not locate its run-step row"); + } + + String pauseToken = UUID.randomUUID().toString(); + LocalDateTime now = LocalDateTime.now(); + + WorkflowRunPauseEntity pause = new WorkflowRunPauseEntity(); + pause.setRunId(context.runId()); + pause.setStepId(stepRow.getId()); + pause.setPauseKind("await_approval"); + pause.setPauseToken(pauseToken); + pause.setPausedAt(now); + if (cfg.timeoutSecs() != null && cfg.timeoutSecs() > 0) { + pause.setResumeDeadline(now.plusSeconds(cfg.timeoutSecs())); + } + pauseMapper.insert(pause); + + return StepResult.paused(pauseToken, + "awaiting " + (cfg.approvalKind() == null ? "approval" : cfg.approvalKind())); + } +}