Skip to content
Open
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 @@ -815,15 +815,16 @@ private void softEvictPool() {
/**
* The database's clock in epoch milliseconds, for inlining into SQL. The workflow_status and
* queues times that executors compare with each other are stamped and compared on this clock:
* created_at, updated_at, completed_at, started_at_epoch_ms, the rate-limit window, and a
* workflow's deadline. Executors sharing a system database would otherwise write those rows on as
* many clocks as there are hosts, and FIFO order, rate limits, timeouts and retention would all
* inherit the skew between them. Where an executor acts on a deadline itself, it measures the
* time left from a {@link DatabaseTime} reading with its monotonic clock.
* created_at, updated_at, completed_at, started_at_epoch_ms, the rate-limit window, a workflow's
* deadline, and the release of a delayed workflow. Executors sharing a system database would
* otherwise write those rows on as many clocks as there are hosts, and FIFO order, rate limits,
* timeouts, delays and retention would all inherit the skew between them. Where an executor acts
* on a deadline itself, it measures the time left from a {@link DatabaseTime} reading with its
* monotonic clock.
*
* <p>Times that only the writing executor reads stay on the JVM's clock: step timings and the
* durable sleep, recv and getEvent timeouts. So, for now, do delays and their promotion to
* ENQUEUED, and debounce deadlines.
* durable sleep, recv and getEvent timeouts. So, for now, do the end of a relative delay, which
* the enqueuing JVM resolves, and debounce deadlines.
*
* <p>now() is the transaction's start time on Postgres and CockroachDB alike, so every statement
* in one transaction reads the same value. It matches the column defaults, which is what keeps a
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1198,22 +1198,23 @@ public static void updateWorkflowAttributes(
public static void transitionDelayedWorkflows(DbContext ctx) throws SQLException {
var sql =
"""
UPDATE "%s".workflow_status
UPDATE "%1$s".workflow_status
SET status = ?,
deduplication_id = CASE WHEN is_debounced THEN NULL ELSE deduplication_id END
deduplication_id = CASE WHEN is_debounced THEN NULL ELSE deduplication_id END,
updated_at = %2$s
WHERE status = ?
AND delay_until_epoch_ms <= ?
AND delay_until_epoch_ms <= %2$s
"""
.formatted(ctx.schema())
.formatted(ctx.schema(), SystemDatabase.NOW_EPOCH_MS)
+ ctx.andAppScope();

// Released against the database's clock, which every executor shares, rather than this JVM's:
// whichever executor runs the release, a delay ends at the same instant.
try (var conn = ctx.getConnection();
var stmt = conn.prepareStatement(sql)) {
stmt.setString(1, WorkflowState.ENQUEUED.name());
stmt.setString(2, WorkflowState.DELAYED.name());
// This JVM's clock, the one delays are counted from, as in Python and TypeScript.
stmt.setLong(3, System.currentTimeMillis());
ctx.bindAppScope(stmt, 4);
ctx.bindAppScope(stmt, 3);

stmt.executeUpdate();
}
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,154 @@
package dev.dbos.transact.workflow;

import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.junit.jupiter.api.Assumptions.assumeFalse;

import dev.dbos.transact.DBOS;
import dev.dbos.transact.StartWorkflowOptions;
import dev.dbos.transact.utils.DBUtils;
import dev.dbos.transact.utils.PgContainer;

import java.sql.SQLException;
import java.time.Duration;
import java.time.Instant;
import java.util.UUID;

import com.zaxxer.hikari.HikariDataSource;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.AutoClose;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.ValueSource;

interface DelayClockService {
String quick();
}

class DelayClockServiceImpl implements DelayClockService {
@Override
@Workflow
public String quick() {
return "done";
}
}

/**
* A delayed workflow is released against the database's clock, so a JVM whose clock disagrees with
* the database's by far more than the delay still releases it on time.
*
* <p>The skew is the database's: a schema ahead of pg_catalog on the connections' search path
* shadows now() and clock_timestamp() with copies shifted by an hour, which every unqualified call
* in the SDK's SQL resolves to. Postgres only: CockroachDB does not let a function shadow a
* builtin.
*/
public class DelayReleaseClockTest {

private static final String CLOCK_SCHEMA = "dbos_test_skewed_release_clock";
// The workflow's queue, which this executor does not dequeue, so a released row stays ENQUEUED
// with the release's own updated_at.
private static final String QUEUE = "release_clock_queue";
private static final String LISTENED_QUEUE = "release_clock_other_queue";

@AutoClose final PgContainer pgContainer = new PgContainer();
@AutoClose HikariDataSource dataSource;
@AutoClose DBOS dbos;
DelayClockService proxy;

@BeforeEach
void beforeEach() {
assumeFalse(PgContainer.USE_COCKROACH_DB, "CockroachDB cannot shadow now()");
dataSource = pgContainer.dataSource();
}

@AfterEach
void afterEach() throws SQLException {
if (dataSource != null) {
try (var conn = dataSource.getConnection();
var stmt = conn.createStatement()) {
stmt.execute("DROP SCHEMA IF EXISTS %s CASCADE".formatted(CLOCK_SCHEMA));
}
}
}

/** Shifts the database's clock {@code skewHours} from the JVM's and launches DBOS on it. */
private void launchSkewed(int skewHours) throws SQLException {
var interval = "interval '%d hours'".formatted(skewHours);
try (var conn = dataSource.getConnection();
var stmt = conn.createStatement()) {
stmt.execute("CREATE SCHEMA IF NOT EXISTS %s".formatted(CLOCK_SCHEMA));
stmt.execute(
"""
CREATE OR REPLACE FUNCTION %s.now() RETURNS timestamptz LANGUAGE sql STABLE
AS $$ SELECT pg_catalog.now() + %s $$
"""
.formatted(CLOCK_SCHEMA, interval));
stmt.execute(
"""
CREATE OR REPLACE FUNCTION %s.clock_timestamp() RETURNS timestamptz LANGUAGE sql VOLATILE
AS $$ SELECT pg_catalog.clock_timestamp() + %s $$
"""
.formatted(CLOCK_SCHEMA, interval));
}
var url = pgContainer.jdbcUrl() + "?currentSchema=%s,pg_catalog,public".formatted(CLOCK_SCHEMA);
dbos = new DBOS(pgContainer.dbosConfig().withDatabaseUrl(url).withListenQueues(LISTENED_QUEUE));
proxy = dbos.registerProxy(DelayClockService.class, new DelayClockServiceImpl());
dbos.launch();
dbos.registerQueue(QUEUE, new QueueOptions());
dbos.registerQueue(LISTENED_QUEUE, new QueueOptions());

// The harness really did skew the clock the SDK reads.
long skewMs = databaseNowMs() - System.currentTimeMillis();
assertTrue(
Math.abs(skewMs - Duration.ofHours(skewHours).toMillis()) < 60_000,
"database clock is %d ms off the JVM's".formatted(skewMs));
}

/** The skewed clock the SDK reads, in epoch milliseconds. */
private long databaseNowMs() throws SQLException {
try (var conn = dataSource.getConnection();
var stmt = conn.createStatement();
var rs =
stmt.executeQuery(
"SELECT (EXTRACT(epoch FROM %s.now()) * 1000.0)::bigint".formatted(CLOCK_SCHEMA))) {
rs.next();
return rs.getLong(1);
}
}

private static void assertBetween(long low, long high, long actual, String what) {
assertTrue(
low <= actual && actual <= high,
"%s %d not in [%d, %d]".formatted(what, actual, low, high));
}

@ParameterizedTest
@ValueSource(ints = {1, -1})
void aDelayIsReleasedOnTheDatabaseClock(int skewHours) throws Exception {
launchSkewed(skewHours);
var id = UUID.randomUUID().toString();

// Delayed far beyond the skew either way, then moved to end a second from the database's now.
dbos.startWorkflow(
() -> proxy.quick(),
new StartWorkflowOptions(id).withQueue(QUEUE).withDelay(Duration.ofDays(1)));
assertEquals(WorkflowState.DELAYED.name(), DBUtils.getWorkflowRow(dataSource, id).status());
long releaseAt = databaseNowMs() + 1_000;
long start = System.nanoTime();
dbos.setWorkflowDelay(id, Instant.ofEpochMilli(releaseAt));

var row = DBUtils.getWorkflowRow(dataSource, id);
while (WorkflowState.DELAYED.name().equals(row.status())
&& System.nanoTime() - start < Duration.ofSeconds(15).toNanos()) {
Thread.sleep(50);
row = DBUtils.getWorkflowRow(dataSource, id);
}
long elapsedMs = Duration.ofNanos(System.nanoTime() - start).toMillis();
long after = databaseNowMs();

// Released a second later, not an hour early or late, and stamped on the database's clock.
assertEquals(WorkflowState.ENQUEUED.name(), row.status());
assertBetween(900, 10_000, elapsedMs, "release took ms");
assertBetween(releaseAt, after, row.updatedAt(), "updated_at");
}
}
Loading