diff --git a/transact/src/main/java/dev/dbos/transact/database/SystemDatabase.java b/transact/src/main/java/dev/dbos/transact/database/SystemDatabase.java index 131dcd480..567776c6a 100644 --- a/transact/src/main/java/dev/dbos/transact/database/SystemDatabase.java +++ b/transact/src/main/java/dev/dbos/transact/database/SystemDatabase.java @@ -437,6 +437,16 @@ public static Duration toDuration(Long ms) { return ms != null ? Duration.ofMillis(ms) : null; } + /** + * The error column of a workflow or a step, or null when it records no error. + * + *

The Go SDK stores the error as a non-nullable string, so a workflow that succeeded leaves "" + * behind where this one writes NULL. Empty only: no SDK writes a blank non-empty value. + */ + public static String errorOrNull(String error) { + return error != null && error.isEmpty() ? null : error; + } + /** * Initializes the status of a workflow. * diff --git a/transact/src/main/java/dev/dbos/transact/database/dao/StepsDAO.java b/transact/src/main/java/dev/dbos/transact/database/dao/StepsDAO.java index 2791a4ace..582117fff 100644 --- a/transact/src/main/java/dev/dbos/transact/database/dao/StepsDAO.java +++ b/transact/src/main/java/dev/dbos/transact/database/dao/StepsDAO.java @@ -1,6 +1,7 @@ package dev.dbos.transact.database.dao; import dev.dbos.transact.database.DbContext; +import dev.dbos.transact.database.SystemDatabase; import dev.dbos.transact.exceptions.*; import dev.dbos.transact.internal.DebugTriggers; import dev.dbos.transact.json.DBOSSerializer; @@ -144,7 +145,7 @@ public static StepResult checkStepResult( try (ResultSet rs = pstmt.executeQuery()) { if (rs.next()) { // Check if any operation output row exists String output = rs.getString("output"); - String error = rs.getString("error"); + String error = SystemDatabase.errorOrNull(rs.getString("error")); String _stepName = rs.getString("function_name"); String serialization = rs.getString("serialization"); result = @@ -222,7 +223,7 @@ static List listWorkflowSteps( int functionId = rs.getInt("function_id"); String functionName = rs.getString("function_name"); String outputData = rs.getString("output"); - String errorData = rs.getString("error"); + String errorData = SystemDatabase.errorOrNull(rs.getString("error")); String childWorkflowId = rs.getString("child_workflow_id"); Long startedAt = rs.getObject("started_at_epoch_ms", Long.class); Long completedAt = rs.getObject("completed_at_epoch_ms", Long.class); @@ -231,7 +232,10 @@ static List listWorkflowSteps( Object outputVal = null; ErrorResult stepError = null; - if (Objects.requireNonNullElse(loadOutput, true)) { + // As for a workflow's status: the steps are what is wanted, the payloads one field of + // each. See SerializationUtil.canDeserialize. + if (Objects.requireNonNullElse(loadOutput, true) + && SerializationUtil.canDeserialize(serialization, serializer)) { if (outputData != null) { try { outputVal = diff --git a/transact/src/main/java/dev/dbos/transact/database/dao/WorkflowDAO.java b/transact/src/main/java/dev/dbos/transact/database/dao/WorkflowDAO.java index 894d4a5c3..752069f79 100644 --- a/transact/src/main/java/dev/dbos/transact/database/dao/WorkflowDAO.java +++ b/transact/src/main/java/dev/dbos/transact/database/dao/WorkflowDAO.java @@ -25,6 +25,7 @@ import dev.dbos.transact.workflow.GetWorkflowAggregatesInput; import dev.dbos.transact.workflow.ListWorkflowsInput; import dev.dbos.transact.workflow.StepAggregateRow; +import dev.dbos.transact.workflow.StepInfo; import dev.dbos.transact.workflow.WorkflowAggregateRow; import dev.dbos.transact.workflow.WorkflowDelay; import dev.dbos.transact.workflow.WorkflowEvent; @@ -47,6 +48,7 @@ import java.util.HashMap; import java.util.HashSet; import java.util.LinkedHashMap; +import java.util.LinkedHashSet; import java.util.List; import java.util.Map; import java.util.Objects; @@ -1233,8 +1235,11 @@ private static WorkflowStatus resultsToWorkflowStatus( String attributesJson = rs.getString("attributes"); String serializedInput = loadInput ? rs.getString("inputs") : null; String serializedOutput = loadOutput ? rs.getString("output") : null; - String serializedError = loadOutput ? rs.getString("error") : null; + String serializedError = loadOutput ? SystemDatabase.errorOrNull(rs.getString("error")) : null; String serialization = loadInput || loadOutput ? rs.getString("serialization") : null; + // A status read reaches other applications' rows on purpose, and wants their metadata; an + // unreadable payload comes back null rather than failing the read. + boolean readable = SerializationUtil.canDeserialize(serialization, serializer); WorkflowStatus info = new WorkflowStatus( rs.getString("workflow_uuid"), @@ -1247,14 +1252,16 @@ private static WorkflowStatus resultsToWorkflowStatus( (authenticatedRolesJson != null) ? JsonUtility.fromJson(authenticatedRolesJson, new TypeReference>() {}) : null, - loadInput + loadInput && readable ? SerializationUtil.deserializePositionalArgs( serializedInput, serialization, serializer) : null, - loadOutput + loadOutput && readable ? SerializationUtil.deserializeValue(serializedOutput, serialization, serializer) : null, - loadOutput ? ErrorResult.deserialize(serializedError, serialization, serializer) : null, + loadOutput && readable + ? ErrorResult.deserialize(serializedError, serialization, serializer) + : null, rs.getString("executor_id"), SystemDatabase.toInstant(rs.getObject("created_at", Long.class)), SystemDatabase.toInstant(rs.getObject("updated_at", Long.class)), @@ -1327,7 +1334,7 @@ public static Result awaitWorkflowResult( } case ERROR -> { - String error = rs.getString("error"); + String error = SystemDatabase.errorOrNull(rs.getString("error")); Throwable t = SerializationUtil.deserializeError(error, serialization, serializer); return Result.failure(t); } @@ -2246,6 +2253,35 @@ static List listWorkflowStreams(Connection conn, String schema, return streams; } + /** + * Refuse to move a workflow whose payloads this runtime cannot handle. + * + *

Export and import round-trip the payloads through this runtime's serializers, so neither can + * settle for the null a status read reports: a payload dropped on the way through restores a + * workflow that never had it. + */ + private static void requireSerializerFor( + String action, + String workflowId, + String workflowSerialization, + List steps, + DBOSSerializer serializer) { + var formats = new LinkedHashSet(); + if (!SerializationUtil.canDeserialize(workflowSerialization, serializer)) { + formats.add(workflowSerialization); + } + for (var step : steps) { + if (!SerializationUtil.canDeserialize(step.serialization(), serializer)) { + formats.add(step.serialization()); + } + } + if (!formats.isEmpty()) { + throw new IllegalStateException( + "Cannot %s workflow %s: it is serialized as %s, which this application has no serializer for" + .formatted(action, workflowId, String.join(", ", formats))); + } + } + public static List exportWorkflow( DbContext ctx, String workflowId, boolean exportChildren) throws SQLException { @@ -2263,6 +2299,9 @@ public static List exportWorkflow( var steps = StepsDAO.listWorkflowSteps( conn, ctx.schema(), ctx.serializer(), wfid, true, null, null); + if (status != null) { + requireSerializerFor("export", wfid, status.serialization(), steps, ctx.serializer()); + } var events = listWorkflowEvents(conn, ctx.schema(), wfid); var eventHistory = listWorkflowEventHistory(conn, ctx.schema(), wfid); var streams = listWorkflowStreams(conn, ctx.schema(), wfid); @@ -2276,6 +2315,14 @@ public static void importWorkflow(DbContext ctx, List workflow throws SQLException { DBOSSerializer serializer = ctx.serializer(); + // The whole batch, before anything is written: export and import are a single-SDK affair but + // not a single-configuration one, and a payload we cannot re-serialize imports empty. + for (var workflow : workflows) { + var s = workflow.status(); + requireSerializerFor( + "import", s.workflowId(), s.serialization(), workflow.steps(), serializer); + } + var wfSQL = """ INSERT INTO "%s".workflow_status ( diff --git a/transact/src/main/java/dev/dbos/transact/execution/DBOSExecutor.java b/transact/src/main/java/dev/dbos/transact/execution/DBOSExecutor.java index 714d26ae9..035822144 100644 --- a/transact/src/main/java/dev/dbos/transact/execution/DBOSExecutor.java +++ b/transact/src/main/java/dev/dbos/transact/execution/DBOSExecutor.java @@ -1696,6 +1696,22 @@ public WorkflowHandle executeWorkflowById( throw new DBOSNonExistentWorkflowException(workflowId); } + // Reading the row reports an unreadable payload as null, so running the workflow refuses for + // itself: a workflow invoked with arguments it never had is worse than one marked ERROR. + if (!SerializationUtil.canDeserialize(status.serialization(), serializer)) { + var e = + new IllegalStateException( + "Cannot run workflow %s: its arguments are serialized as %s, which this application has no deserializer for" + .formatted(workflowId, status.serialization())); + logger.error("Unreadable serialization for workflow {}", workflowId, e); + // Recorded in this runtime's own format, not the row's: a format we cannot read is one + // we cannot write either, and serializing the error into it would throw and leave the + // workflow PENDING forever — the hang this refusal exists to prevent. The status is the + // part that has to land. + persistWorkflowError(workflowId, e, null); + throw e; + } + Object[] inputs = status.input(); var wfName = RegisteredWorkflow.fullyQualifiedName( diff --git a/transact/src/main/java/dev/dbos/transact/json/SerializationUtil.java b/transact/src/main/java/dev/dbos/transact/json/SerializationUtil.java index e58cbcd2a..890daa810 100644 --- a/transact/src/main/java/dev/dbos/transact/json/SerializationUtil.java +++ b/transact/src/main/java/dev/dbos/transact/json/SerializationUtil.java @@ -26,6 +26,19 @@ public final class SerializationUtil { private SerializationUtil() {} + /** + * Whether this runtime can read the given serialization format. + * + *

The two built-in formats always, and a custom serializer only for the format it names. A + * null format is the native one, written before the column existed. + */ + public static boolean canDeserialize(String serialization, DBOSSerializer customSerializer) { + if (serialization == null || PORTABLE.equals(serialization) || NATIVE.equals(serialization)) { + return true; + } + return customSerializer != null && customSerializer.name().equals(serialization); + } + // ============ Value Serialization ============ /** diff --git a/transact/src/main/java/dev/dbos/transact/txstep/TxStepSchema.java b/transact/src/main/java/dev/dbos/transact/txstep/TxStepSchema.java index 0abf5cc7f..4d7947de8 100644 --- a/transact/src/main/java/dev/dbos/transact/txstep/TxStepSchema.java +++ b/transact/src/main/java/dev/dbos/transact/txstep/TxStepSchema.java @@ -1,5 +1,6 @@ package dev.dbos.transact.txstep; +import dev.dbos.transact.database.SystemDatabase; import dev.dbos.transact.workflow.internal.StepResult; import java.sql.Connection; @@ -66,7 +67,7 @@ public static Optional readResult( stepId, stepName, rs.getString("output"), - rs.getString("error"), + SystemDatabase.errorOrNull(rs.getString("error")), null, rs.getString("serialization"))); } diff --git a/transact/src/test/java/dev/dbos/transact/json/InteropTest.java b/transact/src/test/java/dev/dbos/transact/json/InteropTest.java index 09565c3e8..121d55781 100644 --- a/transact/src/test/java/dev/dbos/transact/json/InteropTest.java +++ b/transact/src/test/java/dev/dbos/transact/json/InteropTest.java @@ -4,15 +4,20 @@ import dev.dbos.transact.DBOS; import dev.dbos.transact.DBOSClient; +import dev.dbos.transact.DBOSTestAccess; import dev.dbos.transact.config.DBOSConfig; import dev.dbos.transact.utils.DBUtils; import dev.dbos.transact.utils.PgContainer; +import dev.dbos.transact.workflow.ExportedWorkflow; +import dev.dbos.transact.workflow.ListWorkflowsInput; import dev.dbos.transact.workflow.Queue; import dev.dbos.transact.workflow.SerializationStrategy; +import dev.dbos.transact.workflow.StepInfo; import dev.dbos.transact.workflow.Workflow; import dev.dbos.transact.workflow.WorkflowClassName; import dev.dbos.transact.workflow.WorkflowHandle; import dev.dbos.transact.workflow.WorkflowState; +import dev.dbos.transact.workflow.internal.StepResult; import java.sql.Connection; import java.sql.PreparedStatement; @@ -229,6 +234,66 @@ private void insertPortableWorkflowRow( } } + /** + * Insert a completed workflow row as another SDK would leave it. + * + *

{@code serialization} and {@code error} are what vary between them: Go stores the error as a + * non-nullable string and so writes "" for a workflow that succeeded, and every SDK writes its + * own native format unless the workflow was declared portable. + */ + private void insertPeerWorkflowRow( + String workflowId, String serialization, String inputs, String output, String error) + throws Exception { + try (Connection conn = dataSource.getConnection()) { + String sql = + """ + INSERT INTO dbos.workflow_status( + workflow_uuid, name, class_name, config_name, + status, inputs, output, error, created_at, serialization, application_name + ) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + """; + try (PreparedStatement stmt = conn.prepareStatement(sql)) { + stmt.setString(1, workflowId); + stmt.setString(2, "echoWorkflow"); + stmt.setString(3, "interop"); + stmt.setString(4, null); + stmt.setString(5, "SUCCESS"); + stmt.setString(6, inputs); + stmt.setString(7, output); + stmt.setString(8, error); + stmt.setLong(9, System.currentTimeMillis()); + stmt.setString(10, serialization); + stmt.setString(11, "interop-peer"); + stmt.executeUpdate(); + } + } + } + + /** Insert a completed step row as another SDK would leave it. */ + private void insertPeerStepRow( + String workflowId, String serialization, String output, String error) throws Exception { + try (Connection conn = dataSource.getConnection()) { + String sql = + """ + INSERT INTO dbos.operation_outputs( + workflow_uuid, function_id, function_name, output, error, serialization, application_name + ) + VALUES (?, ?, ?, ?, ?, ?, ?) + """; + try (PreparedStatement stmt = conn.prepareStatement(sql)) { + stmt.setString(1, workflowId); + stmt.setInt(2, 0); + stmt.setString(3, "echoStep"); + stmt.setString(4, output); + stmt.setString(5, error); + stmt.setString(6, serialization); + stmt.setString(7, "interop-peer"); + stmt.executeUpdate(); + } + } + } + private void insertPortableNotification(String destinationUuid, String topic, String messageJson) throws Exception { try (Connection conn = dataSource.getConnection()) { @@ -422,4 +487,189 @@ public void testInteropNamedArgs() throws Exception { assertEquals(Arrays.asList("a", "b"), storedNamedArgs.get("tags")); } } + + // ============================================================================ + // Test: reading a workflow another application owns + // ============================================================================ + + /** + * A workflow that succeeded, whose error column holds "" rather than NULL. + * + *

That is what the Go SDK leaves behind — it stores the error as a non-nullable string — and + * on a shared system database this runtime is routinely asked about such a row. An empty column + * has to mean the same thing as an absent one; parsing it as JSON does not. + */ + @Test + public void testStatusReadTreatsAnEmptyErrorColumnAsNoError() throws Exception { + String workflowId = "peer-empty-error"; + insertPeerWorkflowRow( + workflowId, "portable_json", "{\"positionalArgs\":[]}", "{\"ok\":true}", ""); + + dbos.launch(); + var status = dbos.getWorkflowStatus(workflowId).orElseThrow(); + + assertEquals(WorkflowState.SUCCESS, status.status()); + assertNull(status.error(), "an empty error column is not an error"); + } + + /** + * A workflow another application ran, in a serialization format this runtime cannot read. + * + *

Workflow IDs address the whole system database, so a status read reaches a peer's rows on + * purpose — and what is wanted there is the metadata: who owns it, what it is, whether it + * finished. A payload in the peer's own format must not take that away, so the fields that cannot + * be deserialized come back null and the read succeeds. Python behaves the same way, in + * safe_deserialize. + */ + @Test + public void testStatusReadOfAPeerRowInAnUnreadableFormat() throws Exception { + String workflowId = "peer-foreign-format"; + insertPeerWorkflowRow(workflowId, "py_pickle", "gASVCgAAAA==", "gASVBAAAAA==", null); + + dbos.launch(); + var status = dbos.getWorkflowStatus(workflowId).orElseThrow(); + + // The metadata is the point, and it is all there. + assertEquals(workflowId, status.workflowId()); + assertEquals("echoWorkflow", status.workflowName()); + assertEquals(WorkflowState.SUCCESS, status.status()); + assertEquals("interop-peer", status.applicationName()); + assertEquals("py_pickle", status.serialization()); + + // The payloads are not, and saying so beats failing the whole read. + assertNull(status.input()); + assertNull(status.output()); + assertNull(status.error()); + } + + /** The same row, through the listing rather than by ID. */ + @Test + public void testListingAPeerRowInAnUnreadableFormat() throws Exception { + String workflowId = "peer-foreign-format-listed"; + insertPeerWorkflowRow(workflowId, "py_pickle", "gASVCgAAAA==", "gASVBAAAAA==", ""); + + dbos.launch(); + var listed = dbos.listWorkflows(new ListWorkflowsInput().withWorkflowIds(workflowId)).get(0); + + assertEquals(workflowId, listed.workflowId()); + assertEquals("interop-peer", listed.applicationName()); + assertNull(listed.input()); + assertNull(listed.output()); + assertNull(listed.error()); + } + + /** + * A step of a peer's workflow, recorded with an empty error column. + * + *

Steps carry the same column with the same quirk, so they get the same reading. + */ + @Test + public void testStepReadTreatsAnEmptyErrorColumnAsNoError() throws Exception { + String workflowId = "peer-step-empty-error"; + insertPeerWorkflowRow(workflowId, "portable_json", "{\"positionalArgs\":[]}", "{}", ""); + insertPeerStepRow(workflowId, "portable_json", "{\"ok\":true}", ""); + + dbos.launch(); + var steps = dbos.listWorkflowSteps(workflowId); + + assertEquals(1, steps.size()); + assertNull(steps.get(0).error(), "an empty error column is not an error"); + } + + /** + * The same gap without another language in it: two Java applications on one system database, one + * configured with a custom serializer and one not. + * + *

`custom_base64` is a format this runtime has no deserializer for, exactly as `py_pickle` is, + * and canDeserialize says so without having to try and fail. + */ + @Test + public void testStatusReadOfARowWrittenByACustomSerializerWeLack() throws Exception { + String workflowId = "peer-custom-serializer"; + insertPeerWorkflowRow(workflowId, "custom_base64", "cG9zaXRpb25hbA==", "b3V0cHV0", null); + + dbos.launch(); + var status = dbos.getWorkflowStatus(workflowId).orElseThrow(); + + assertEquals(WorkflowState.SUCCESS, status.status()); + assertEquals("custom_base64", status.serialization()); + assertNull(status.input()); + assertNull(status.output()); + } + + /** + * Exporting a workflow this runtime cannot read is refused rather than silently emptied. + * + *

A status read reports an unreadable payload as null, which is right when the metadata is + * what was asked for. An export is imported back, so the same null would restore a workflow that + * never had an input. + */ + @Test + public void testExportRefusesAWorkflowItCannotRead() throws Exception { + String workflowId = "peer-export"; + insertPeerWorkflowRow(workflowId, "py_pickle", "gASVCgAAAA==", "gASVBAAAAA==", null); + + dbos.launch(); + var systemDatabase = DBOSTestAccess.getSystemDatabase(dbos); + + var thrown = + assertThrows( + IllegalStateException.class, () -> systemDatabase.exportWorkflow(workflowId, false)); + assertTrue( + thrown.getMessage().contains("py_pickle"), + "The refusal should name the format it could not read, got: " + thrown.getMessage()); + } + + /** + * Importing one is refused too, before anything is written. + * + *

Export and import are a single-SDK affair, but not a single-configuration one: the source + * application may have had a serializer this one does not. Import re-serializes the payloads, so + * it needs that serializer as much as export did — and a payload the export already carried as + * null would otherwise be written as NULL and committed. + */ + @Test + public void testImportRefusesAWorkflowItCannotRead() throws Exception { + String workflowId = "peer-import"; + insertPeerWorkflowRow(workflowId, "portable_json", "{\"positionalArgs\":[]}", "{}", null); + + dbos.launch(); + var systemDatabase = DBOSTestAccess.getSystemDatabase(dbos); + + // Exported readably, then given a step in a format this application has no serializer for — + // which is what a batch arriving from an application configured differently looks like. + var exported = systemDatabase.exportWorkflow(workflowId, false).get(0); + var foreignStep = + new StepInfo( + 0, "echoStep", "cG9zaXRpb25hbA==", null, null, null, null, "custom_base64", null); + var batch = + List.of( + new ExportedWorkflow( + exported.status(), + List.of(foreignStep), + exported.events(), + exported.eventHistory(), + exported.streams())); + + var thrown = + assertThrows(IllegalStateException.class, () -> systemDatabase.importWorkflow(batch)); + assertTrue( + thrown.getMessage().contains("custom_base64"), + "The refusal should name the format, got: " + thrown.getMessage()); + } + + /** + * Replaying a step whose result this application cannot read fails; it does not replay as null. + * + *

The tolerant reads are the ones that assemble a record. A recorded step result is consumed + * to continue a workflow — a step that "returned null" because its output could not be read would + * corrupt the run it is replaying, silently and durably. + */ + @Test + public void testStepResultRefusesToReplayWhatItCannotRead() { + var recorded = + new StepResult("wf", 0, "echoStep", "cG9zaXRpb25hbA==", null, null, "custom_base64"); + + assertThrows(IllegalArgumentException.class, () -> recorded.toResult(null)); + } } diff --git a/transact/src/test/java/dev/dbos/transact/json/PortableSerializationTest.java b/transact/src/test/java/dev/dbos/transact/json/PortableSerializationTest.java index 701b74961..e39788e7d 100644 --- a/transact/src/test/java/dev/dbos/transact/json/PortableSerializationTest.java +++ b/transact/src/test/java/dev/dbos/transact/json/PortableSerializationTest.java @@ -1036,6 +1036,60 @@ private WorkflowStatusRow waitForWorkflowTerminal(String workflowId, Duration ti throw new AssertionError("Workflow " + workflowId + " did not reach terminal state in time"); } + /** Insert an enqueued workflow row carrying an arbitrary serialization format. */ + private void insertEnqueuedRowWithSerialization( + String workflowId, String queueName, String inputsJson, String serialization) + throws Exception { + try (Connection conn = dataSource.getConnection()) { + String sql = + """ + INSERT INTO dbos.workflow_status( + workflow_uuid, name, class_name, config_name, + queue_name, status, inputs, created_at, serialization + ) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) + """; + try (PreparedStatement stmt = conn.prepareStatement(sql)) { + stmt.setString(1, workflowId); + stmt.setString(2, "recvWorkflow"); + stmt.setString(3, "PortableTestService"); + stmt.setString(4, null); + stmt.setString(5, queueName); + stmt.setString(6, "ENQUEUED"); + stmt.setString(7, inputsJson); + stmt.setLong(8, System.currentTimeMillis()); + stmt.setString(9, serialization); + stmt.executeUpdate(); + } + } + } + + /** + * A workflow whose arguments are in a format this application cannot read is marked ERROR, not + * left in PENDING. + * + *

Reading the row does not fail on such a format any more — a status read reports the + * arguments as null and hands back the metadata — so running the workflow has to refuse for + * itself, and say which format it would have taken. + */ + @Test + public void testUnreadableSerializationMarksTheWorkflowErrored() throws Exception { + Queue testQueue = new Queue("testq"); + dbos.registerQueue(testQueue); + dbos.registerProxy(PortableTestService.class, new PortableTestServiceImpl(dbos)); + dbos.launch(); + + String workflowId = UUID.randomUUID().toString(); + insertEnqueuedRowWithSerialization(workflowId, "testq", "cGlja2xlZA==", "py_pickle"); + + var row = waitForWorkflowTerminal(workflowId, Duration.ofSeconds(30)); + assertEquals(WorkflowState.ERROR.name(), row.status()); + assertNotNull(row.error()); + assertTrue( + row.error().contains("py_pickle"), + "The error should name the format it could not read, got: " + row.error()); + } + /** * Tests that completely invalid (unparseable) JSON in the inputs column results in the workflow * being marked as ERROR rather than being stuck in PENDING forever.