Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -254,6 +254,23 @@ public void testDeleteCheckpointsThroughDbos() throws Exception {
assertEquals(0, DBUtils.getTxStepRows(dataSource, "wf-batch-2").size());
}

@Test
public void testRewindRerunsTheTransaction() throws Exception {
var wfid = "wf-rewind";
var user = "rewindUser";
try (var _o = new WorkflowOptions(wfid).setContext()) {
assertEquals(1, proxy.insertWorkflow(user).greetCount());
}

// The rewind deletes the step's checkpoint, so the transaction runs again rather than
// replaying its recorded result.
WorkflowHandle<FactoryTestService.TestResult, RuntimeException> handle =
dbos.rewindWorkflow(wfid, 0);
assertEquals(new FactoryTestService.TestResult(user, 2), handle.getResult());
assertEquals(2, getGreetCount(user));
assertEquals(1, DBUtils.getTxStepRows(dataSource, wfid).size());
}

@Test
public void testInsert() throws Exception {
var wfid = "wf1";
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -244,6 +244,23 @@ public void testDeleteCheckpointsThroughDbos() throws Exception {
assertEquals(0, DBUtils.getTxStepRows(dataSource, "wf-batch-2").size());
}

@Test
public void testRewindRerunsTheTransaction() throws Exception {
var wfid = "wf-rewind";
var user = "rewindUser";
try (var _o = new WorkflowOptions(wfid).setContext()) {
assertEquals(1, proxy.insertWorkflow(user).greetCount());
}

// The rewind deletes the step's checkpoint, so the transaction runs again rather than
// replaying its recorded result.
WorkflowHandle<FactoryTestService.TestResult, RuntimeException> handle =
dbos.rewindWorkflow(wfid, 0);
assertEquals(new FactoryTestService.TestResult(user, 2), handle.getResult());
assertEquals(2, getGreetCount(user));
assertEquals(1, DBUtils.getTxStepRows(dataSource, wfid).size());
}

@Test
public void testInsert() throws Exception {
var wfid = "wf1";
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -229,6 +229,33 @@ void deleteCheckpointsThroughDbos() throws SQLException {
}
}

@Test
void rewindRerunsTheTransactionalStep() throws SQLException {
try (var db = new TransactionalStepTest.TestDatabase()) {
runner(db)
.run(
ctx -> {
assertThat(ctx).hasNotFailed();
var workflow = ctx.getBean(OrderWorkflowService.class);
var dbos = ctx.getBean(DBOS.class);
var wfid = "wf-jdbc-int-rewind";

try (var _o = new WorkflowOptions(wfid).setContext()) {
workflow.processOrder("ord-r", "Widget", 1);
}
assertThat(TransactionalStepTest.getTxRows(db.dataSource, wfid)).hasSize(1);

// Remove the order the first run placed: a step that runs again places it again,
// while one replayed from a leftover checkpoint does not.
new JdbcTemplate(db.dataSource).update("DELETE FROM orders WHERE id = ?", "ord-r");
dbos.rewindWorkflow(wfid, 0).getResult();

assertThat(orderCount(db.dataSource, "ord-r")).isEqualTo(1);
assertThat(TransactionalStepTest.getTxRows(db.dataSource, wfid)).hasSize(1);
});
}
}

@Test
void isolationLevel() {
try (var db = new TransactionalStepTest.TestDatabase()) {
Expand Down
56 changes: 56 additions & 0 deletions transact/src/main/java/dev/dbos/transact/DBOS.java
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
import dev.dbos.transact.workflow.Queue;
import dev.dbos.transact.workflow.QueueConflictResolution;
import dev.dbos.transact.workflow.QueueOptions;
import dev.dbos.transact.workflow.RewindOptions;
import dev.dbos.transact.workflow.ScheduleStatus;
import dev.dbos.transact.workflow.SendMessage;
import dev.dbos.transact.workflow.SerializationStrategy;
Expand Down Expand Up @@ -1044,6 +1045,61 @@ public void updateWorkflowAttributes(
return forkWorkflow(workflowId, startStep, new ForkOptions());
}

/**
* Rewind a workflow: re-run it in place, under the same ID, from the step provided. Steps before
* {@code startStep} are replayed from their checkpoints; everything from {@code startStep} on is
* discarded and runs again. Only a workflow in a terminal state can be rewound, so cancel a
* running one first.
*
* <p>From {@code startStep} on, the rewind deletes the workflow's step checkpoints, including
* those of registered transactional step factories; rolls back the events it published to their
* last value from before the cut; deletes the messages it consumed, and any it has not consumed
* yet; and removes its streams' close markers so the replay can write to them again. Stream
* entries are kept. The workflow is then re-enqueued.
*
* <p>A message consumed before the system database reached migration 121, or by a DBOS version
* that does not record which step consumed it, cannot be matched to a step. The rewind leaves it
* consumed, so a replayed {@code recv} past the cut waits for a new message instead.
*
* <p>Called from a workflow, the rewind is a step, so a recovered caller does not rewind its
* target again.
*
* @param <T> Return type of the workflow function
* @param <E> Checked exception thrown by the workflow function, if any
* @param workflowId ID of the workflow to rewind
* @param startStep the first step to discard and run again; 0 re-runs the whole workflow
* @param options {@link RewindOptions} containing the queue, partition key and application
* version to re-enqueue the workflow with
* @return handle to the rewound workflow
* @throws dev.dbos.transact.exceptions.DBOSNonExistentWorkflowException if the workflow does not
* exist
* @throws IllegalArgumentException if {@code startStep} is negative, the queue does not exist, or
* the partition key does not match whether the queue is partitioned
* @throws IllegalStateException if the workflow is not in a terminal state, or its status changed
* while it was being rewound; in the second case, retry the rewind
* @throws RuntimeException if a transactional step factory's checkpoints could not be deleted;
* the workflow is not rewound, and the rewind can be retried
*/
public <T, E extends Exception> @NonNull WorkflowHandle<T, E> rewindWorkflow(
@NonNull String workflowId, int startStep, @NonNull RewindOptions options) {
return ensureLaunched("rewindWorkflow").rewindWorkflow(workflowId, startStep, options);
}

/**
* Rewind a workflow: re-run it in place, under the same ID, from the step provided. See {@link
* #rewindWorkflow(String, int, RewindOptions)}.
*
* @param <T> Return type of the workflow function
* @param <E> Checked exception thrown by the workflow function, if any
* @param workflowId ID of the workflow to rewind
* @param startStep the first step to discard and run again; 0 re-runs the whole workflow
* @return handle to the rewound workflow
*/
public <T, E extends Exception> @NonNull WorkflowHandle<T, E> rewindWorkflow(
@NonNull String workflowId, int startStep) {
return rewindWorkflow(workflowId, startStep, new RewindOptions());
}

/**
* List all registered application versions, ordered by timestamp descending.
*
Expand Down
44 changes: 44 additions & 0 deletions transact/src/main/java/dev/dbos/transact/DBOSClient.java
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
import dev.dbos.transact.workflow.Queue;
import dev.dbos.transact.workflow.QueueConflictResolution;
import dev.dbos.transact.workflow.QueueOptions;
import dev.dbos.transact.workflow.RewindOptions;
import dev.dbos.transact.workflow.ScheduleStatus;
import dev.dbos.transact.workflow.SendMessage;
import dev.dbos.transact.workflow.SerializationStrategy;
Expand Down Expand Up @@ -1619,6 +1620,49 @@ public void deleteWorkflows(@NonNull List<String> workflowIds, boolean deleteChi
return retrieveWorkflow(forkedWorkflowId);
}

/**
* Rewind a workflow: re-run it in place, under the same ID, from the step provided. Only a
* workflow in a terminal state can be rewound. See {@link DBOS#rewindWorkflow(String, int,
* RewindOptions)} for what a rewind discards.
*
* <p>A client cannot reach the application's databases, so unlike {@link DBOS#rewindWorkflow}, it
* does not delete the checkpoints transactional step factories keep there. A transactional step
* past the cut that still has one replays its recorded result instead of running again. To re-run
* those transactions, rewind from within the application.
*
* @param <T> Type of the workflow's return value
* @param <E> Type of any checked exception thrown by the workflow
* @param workflowId ID of the workflow to rewind
* @param startStep the first step to discard and run again; 0 re-runs the whole workflow
* @param options Options for the rewind
* @return `WorkflowHandle` for the rewound workflow
* @throws dev.dbos.transact.exceptions.DBOSNonExistentWorkflowException if the workflow does not
* exist
* @throws IllegalArgumentException if {@code startStep} is negative
* @throws IllegalStateException if the workflow is not in a terminal state, or its status changed
* while it was being rewound; in the second case, retry the rewind
*/
public <T, E extends Exception> @NonNull WorkflowHandle<T, E> rewindWorkflow(
@NonNull String workflowId, int startStep, @NonNull RewindOptions options) {
systemDatabase.rewindWorkflow(workflowId, startStep, options);
return retrieveWorkflow(workflowId);
}

/**
* Rewind a workflow: re-run it in place, under the same ID, from the step provided. See {@link
* #rewindWorkflow(String, int, RewindOptions)}.
*
* @param <T> Type of the workflow's return value
* @param <E> Type of any checked exception thrown by the workflow
* @param workflowId ID of the workflow to rewind
* @param startStep the first step to discard and run again; 0 re-runs the whole workflow
* @return `WorkflowHandle` for the rewound workflow
*/
public <T, E extends Exception> @NonNull WorkflowHandle<T, E> rewindWorkflow(
@NonNull String workflowId, int startStep) {
return rewindWorkflow(workflowId, startStep, new RewindOptions());
}

/**
* Get the status of a workflow
*
Expand Down
20 changes: 20 additions & 0 deletions transact/src/main/java/dev/dbos/transact/conductor/Conductor.java
Original file line number Diff line number Diff line change
Expand Up @@ -805,6 +805,7 @@ CompletableFuture<BaseResponse> getResponseAsync(BaseMessage message, WebSocket
case RESUME -> handleResume(this, (ResumeRequest) message);
case RESUME_SCHEDULE -> handleResumeSchedule(this, (ResumeScheduleRequest) message);
case RETENTION -> handleRetention(this, (RetentionRequest) message);
case REWIND_WORKFLOW -> handleRewind(this, (RewindWorkflowRequest) message);
case SET_LATEST_APPLICATION_VERSION ->
handleSetLatestApplicationVersion(this, (SetLatestApplicationVersionRequest) message);
case TRIGGER_SCHEDULE -> handleTriggerSchedule(this, (TriggerScheduleRequest) message);
Expand Down Expand Up @@ -928,6 +929,25 @@ static CompletableFuture<BaseResponse> handleFork(
});
}

static CompletableFuture<BaseResponse> handleRewind(
Conductor conductor, RewindWorkflowRequest request) {
return CompletableFuture.supplyAsync(
() -> {
if (request.body == null || request.body.workflow_id == null) {
return new SuccessResponse(
request, new IllegalArgumentException("Invalid Rewind Workflow Request"));
}
try {
conductor.dbosExecutor.rewindWorkflow(
request.body.workflow_id, request.startStep(), request.toOptions());
return new SuccessResponse(request, true);
} catch (Exception e) {
logger.error("Exception encountered when rewinding workflow {}", request, e);
return new SuccessResponse(request, e);
}
});
}

static CompletableFuture<BaseResponse> handleForkFromFailure(
Conductor conductor, ForkFromFailureRequest request) {
return CompletableFuture.supplyAsync(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,7 @@
@JsonSubTypes.Type(value = ResumeRequest.class, name = "resume"),
@JsonSubTypes.Type(value = ResumeScheduleRequest.class, name = "resume_schedule"),
@JsonSubTypes.Type(value = RetentionRequest.class, name = "retention"),
@JsonSubTypes.Type(value = RewindWorkflowRequest.class, name = "rewind_workflow"),
@JsonSubTypes.Type(
value = SetLatestApplicationVersionRequest.class,
name = "set_latest_application_version"),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@ public enum MessageType {
RESUME("resume"),
RESUME_SCHEDULE("resume_schedule"),
RETENTION("retention"),
REWIND_WORKFLOW("rewind_workflow"),
SET_LATEST_APPLICATION_VERSION("set_latest_application_version"),
TRIGGER_SCHEDULE("trigger_schedule");

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
package dev.dbos.transact.conductor.protocol;

import dev.dbos.transact.workflow.RewindOptions;

import com.fasterxml.jackson.annotation.JsonIgnoreProperties;

public class RewindWorkflowRequest extends BaseMessage {
public RewindWorkflowBody body;

public RewindWorkflowRequest() {}

public RewindWorkflowRequest(String requestId, String workflowId, Integer startStep) {
this.type = MessageType.REWIND_WORKFLOW.getValue();
this.request_id = requestId;
this.body = new RewindWorkflowBody();
this.body.workflow_id = workflowId;
this.body.start_step = startStep;
}

@JsonIgnoreProperties(ignoreUnknown = true)
public static class RewindWorkflowBody {
public String workflow_id;
public Integer start_step; // optional: omitted rewinds the whole history
public String application_version; // optional
public String queue_name; // optional
public String queue_partition_key; // optional
}

public int startStep() {
return body.start_step == null ? 0 : body.start_step;
}

public RewindOptions toOptions() {
return new RewindOptions(body.application_version, body.queue_name, body.queue_partition_key);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@
import dev.dbos.transact.workflow.NotificationInfo;
import dev.dbos.transact.workflow.Queue;
import dev.dbos.transact.workflow.QueueOptions;
import dev.dbos.transact.workflow.RewindOptions;
import dev.dbos.transact.workflow.ScheduleStatus;
import dev.dbos.transact.workflow.SendMessage;
import dev.dbos.transact.workflow.StepAggregateRow;
Expand Down Expand Up @@ -1245,6 +1246,15 @@ public String forkWorkflow(String originalWorkflowId, int startStep, ForkOptions
() -> WorkflowDAO.forkWorkflow(ctx, originalWorkflowId, startStep, options));
}

public void rewindWorkflow(String workflowId, int startStep, RewindOptions options) {
dbRetryIncludingSerializationError(
"rewindWorkflow", () -> WorkflowDAO.rewindWorkflow(ctx, workflowId, startStep, options));
}

public void checkRewindable(String workflowId, int startStep) {
dbRetry(() -> WorkflowDAO.checkRewindable(ctx, workflowId, startStep));
}

public List<String> forkFromFailure(List<String> workflowIds, ForkFromFailureOptions options) {
return dbRetryIncludingSerializationError(
"forkFromFailure", () -> WorkflowDAO.forkFromFailure(ctx, workflowIds, options));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -291,7 +291,7 @@ public static Object recv(
var updateSql =
"""
UPDATE "%1$s".notifications
SET consumed = TRUE
SET consumed = TRUE, consumed_by_function_id = ?
WHERE destination_uuid = ?
AND topic = ?
AND consumed = FALSE
Expand Down Expand Up @@ -321,10 +321,13 @@ public static Object recv(
String serializedMessage = null;
String serialization = null;
try (PreparedStatement stmt = conn.prepareStatement(updateSql)) {
stmt.setString(1, workflowId);
stmt.setString(2, recvTopic);
stmt.setString(3, workflowId);
stmt.setString(4, recvTopic);
// The consuming step is recorded so a rewind can delete exactly the
// messages consumed by the steps it discards.
stmt.setInt(1, stepId);
stmt.setString(2, workflowId);
stmt.setString(3, recvTopic);
stmt.setString(4, workflowId);
stmt.setString(5, recvTopic);

// Note, if there are two executors running the same workflow waiting on the
// same recv, only the first one will return a row here. The second one gets
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -102,6 +102,29 @@ private static void insertStream(
}
}

/**
* Deletes the close sentinels a workflow wrote from {@code fromStepId} on, so a rewound workflow
* can append to its streams again. A leftover sentinel would end every reader before the replay's
* entries.
*/
static void deleteCloseSentinels(
Connection conn, String schema, String workflowId, int fromStepId) throws SQLException {
var closed = SerializationUtil.serializeValue(STREAM_CLOSED_SENTINEL, "portable_json", null);
var sql =
"""
DELETE FROM "%s".streams
WHERE workflow_uuid = ? AND function_id >= ? AND value = ? AND serialization = ?
"""
.formatted(schema);
try (var stmt = conn.prepareStatement(sql)) {
stmt.setString(1, workflowId);
stmt.setInt(2, fromStepId);
stmt.setString(3, closed.serializedValue());
stmt.setString(4, closed.serialization());
stmt.executeUpdate();
}
}

private static int getNextOffsetTx(Connection conn, String schema, String workflowId, String key)
throws SQLException {
String sql =
Expand Down
Loading
Loading