Skip to content

Add workflow rewind, with system-database migration 121 - #608

Merged
devhawk merged 5 commits into
mainfrom
workflow-rewind
Oct 7, 2026
Merged

devhawk merged 5 commits into
mainfrom
workflow-rewind

Conversation

@devhawk

@devhawk devhawk commented Oct 7, 2026 •

Copy link
Copy Markdown
Collaborator

Closes #537.

Rewind re-runs a terminal workflow in place, under the same ID, from a start step. Steps before the cut replay from their checkpoints; everything from the cut on is discarded and runs again. A workflow that is still PENDING, ENQUEUED or DELAYED is refused ("cancel it first"). Python, TypeScript and Go already have it; this follows Python and TypeScript, including clearing the legacy output/error columns.

What a rewind does

In one system-database transaction, from the start step on:

  • rolls each event back to its last value published before the cut, from workflow_events_history, then deletes those history rows;
  • deletes the step checkpoints in operation_outputs;
  • deletes the streams' close sentinels so the replay can append again. Stream entries are kept, since their offsets are what readers resume from;
  • deletes the messages consumed by discarded steps, and any not yet consumed;
  • deletes the workflow_output row and clears the legacy columns;
  • re-enqueues the workflow (internal queue, or the given queue and partition), resetting recovery attempts, deadline, dedup ID and timestamps and clearing owner_xid; the application version is restamped only when given;
  • re-checks the status it read, failing with "retry the rewind" if it changed.

Before that, DBOS.rewindWorkflow deletes the transactional step factories' tx_step_outputs rows from the cut on, through the checkpoint-store registry from #607. If a delete fails, the workflow keeps its terminal status and is not re-enqueued, and the rewind can be retried. Called from a workflow, the rewind is a step (DBOS.rewindWorkflow), so a recovered caller does not rewind again once that step is recorded. Like the other built-in steps that write, it runs again on recovery if the process dies after the rewind commits and before its checkpoint is written (#591).

API

  • DBOS.rewindWorkflow(workflowId, startStep[, RewindOptions]) and the same on DBOSClient. RewindOptions carries the queue, partition key and application version. startStep 0 re-runs the whole workflow.
  • A client cannot reach the application databases, so DBOSClient.rewindWorkflow leaves the step factories' checkpoints in place: a transactional step past the cut replays its recorded result. The javadoc says so and points to rewinding from the application.

Migration 121

Adds notifications.consumed_by_function_id as INT4 (INTEGER would be 8 bytes on CockroachDB), and recv now records the consuming step there, so a rewind can delete exactly the messages consumed past the cut. A message consumed before 121 (or by a version that doesn't record the step) can't be matched and stays consumed; the javadoc says so.

Conductor

Answers rewind_workflow (workflow_id, start_step, application_version, queue_name, queue_partition_key; a missing start_step means the whole history).

Known limitation

The step-factory deletes run after an unlocked status check and before the system-database transaction, as in the other SDKs. A concurrent rewind or resume landing in between can delete checkpoints a replay has just written. The window is narrow and the fix is being settled across SDKs.

Tests

RewindTest covers events, streams, messages (including one consumed without a recorded step, as before migration 121), legacy outcome columns, step-factory checkpoints (including a failed delete and an active workflow), the client, in-workflow and in-step calls with recovery, and refusal of active and missing workflows. Jdbi, jOOQ and Spring tests check that a rewind re-runs the transaction. ConductorTest covers the message with and without a start step; MigrationManagerTest covers 121, including a database stopped at 120.

🤖 Generated with Claude Code

Rewind re-runs a terminal workflow in place, under the same ID, from a
start step. Steps before the cut replay from their checkpoints;
everything from the cut on is discarded and runs again. A non-terminal
workflow is refused ("cancel it first").

In one system-database transaction, from the start step on, a rewind:

- rolls each event back to its last value published before the cut,
  from `workflow_events_history`, then deletes the history rows;
- deletes the step checkpoints in `operation_outputs`;
- deletes the streams' close sentinels, so the replay can write to them
  again (stream entries are kept);
- deletes the messages consumed by discarded steps, and any not yet
  consumed;
- deletes the `workflow_output` row and clears the legacy output and
  error columns;
- re-enqueues the workflow (internal queue, or the given queue and
  partition) with recovery attempts, deadline, dedup ID and timestamps
  reset and `owner_xid` cleared, restamping the application version
  only when one is given;
- re-checks the status it read, failing with "retry the rewind" if it
  changed.

Before the system database, `DBOS.rewindWorkflow` deletes the
checkpoints transactional step factories keep in `tx_step_outputs` from
the cut on. If that fails, the workflow keeps its terminal status and is
not re-enqueued, and the rewind can be retried. Called from a workflow,
the rewind is a step.

- **API:** `DBOS.rewindWorkflow` and `DBOSClient.rewindWorkflow`, with
  `RewindOptions` (queue, partition key, application version). A client
  cannot reach the application databases, so it leaves the step
  factories' checkpoints in place.
- **Migration 121** adds `notifications.consumed_by_function_id`, and
  `recv` records the consuming step there.
- **Conductor** answers `rewind_workflow`.
List every exception DBOS.rewindWorkflow and DBOSClient.rewindWorkflow
can throw, including the status race that asks the caller to retry and
a failed step-factory checkpoint delete.

A message consumed before migration 121, or by a DBOS version that does
not record the consuming step, has no consumed_by_function_id, so a
rewind leaves it consumed and a replayed recv past the cut waits for a
new message. Say so in the javadoc.
The rewind's status check ran its own copy of getWorkflowState's query
and kept its own set of active statuses. Use getWorkflowState and
WorkflowState.isActive instead.

The copy compared statuses as strings, so a status this SDK does not
know counted as terminal and the rewind went ahead. getWorkflowState
rejects an unknown status, as every other status read in the SDK
already does, so that test is removed.

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Changes recommended

A crash between committing the rewind and checkpointing its caller can cause recovery to fail or rewind the target twice.

2 open findings
What changed in this PR

Adds in-place workflow rewind support, migration 121, Conductor integration, and transactional-step checkpoint cleanup.

Changes:

  • Adds public rewind APIs and transactional database logic.
  • Tracks message-consuming steps and restores events, streams, and notifications.
  • Adds Conductor support and broad integration tests.
File Description
RewindTestService.java Provides rewind test workflows.
RewindTest.java Tests rewind behavior and validation.
MigrationManagerTest.java Tests migration 121.
ConductorTest.java Tests rewind protocol handling.
RewindOptions.java Defines rewind configuration.
MigrationManager.java Adds migration 121.
DBOSExecutor.java Coordinates rewind and checkpoint deletion.
DBOSClient.java Exposes client rewind APIs.
DBOS.java Exposes application rewind APIs.
SystemDatabase.java Adds rewind database entry points.
WorkflowDAO.java Implements transactional rewind.
StreamsDAO.java Removes discarded close sentinels.
NotificationsDAO.java Records the consuming function ID.
RewindWorkflowRequest.java Defines the Conductor request.
MessageType.java Registers the rewind message type.
BaseMessage.java Adds rewind deserialization.
Conductor.java Handles rewind requests.
TransactionalStepJdbcIntegrationTest.java Tests Spring transactional replay.
JooqStepFactoryTest.java Tests jOOQ transactional replay.
JdbiStepFactoryTest.java Tests Jdbi transactional replay.

🧠 Review effort: Balanced


💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

Comment thread transact/src/test/java/dev/dbos/transact/workflow/RewindTest.java
INTEGER is INT4 on Postgres but INT8 on CockroachDB, so migration 121
gave CockroachDB an 8-byte column while every other step-ID column is
INT4. JDBC then refuses to read it as an Integer, which failed
rewindDeletesNotifications on CockroachDB.
A message consumed before migration 121 has no consumed_by_function_id,
so the rewind's delete cannot tie it to a discarded step and must leave
it consumed. Clear the marker on a consumed message, rewind past its
recv, and check that it stays consumed and the replayed recv takes the
next message.
@devhawk
devhawk merged commit 723ecc8 into main Oct 7, 2026
12 checks passed
@devhawk
devhawk deleted the workflow-rewind branch October 7, 2026 21:42
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Rewind

3 participants