diff --git a/CHANGELOG.md b/CHANGELOG.md index f1deaf3b..809b8f92 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,4 +1,7 @@ ## Unreleased +* Recreate a deleted large-payload container and retry the upload once, without allowing stale concurrent failures to invalidate a newly recreated container. +* Preserve backend timestamp precision in orchestration history and client metadata, including export blob names, while retaining existing orchestration replay timestamp behavior. +* Export entity operations and locks using the .NET-compatible `EventSent`/`EventRaised` message representation instead of Java-native entity event objects, preserving missing stack traces in failure details. * Add the `exporthistory` module for durable, checkpointed export of terminal orchestration history to Azure Blob Storage ([#293](https://github.com/microsoft/durabletask-java/pull/293)) * Add client APIs to list terminal instance IDs by completion time (`listInstanceIds`) and read orchestration history (`getOrchestrationHistory`) ([#292](https://github.com/microsoft/durabletask-java/pull/292)) * Add `createReplaySafeLogger` to suppress orchestration log output during replay ([#295](https://github.com/microsoft/durabletask-java/pull/295)). diff --git a/azure-blob-payloads/src/main/java/com/microsoft/durabletask/azureblobpayloads/BlobPayloadStore.java b/azure-blob-payloads/src/main/java/com/microsoft/durabletask/azureblobpayloads/BlobPayloadStore.java index 26f078f9..2e969e83 100644 --- a/azure-blob-payloads/src/main/java/com/microsoft/durabletask/azureblobpayloads/BlobPayloadStore.java +++ b/azure-blob-payloads/src/main/java/com/microsoft/durabletask/azureblobpayloads/BlobPayloadStore.java @@ -8,6 +8,7 @@ import com.azure.storage.blob.BlobServiceClient; import com.azure.storage.blob.BlobServiceClientBuilder; import com.azure.storage.blob.models.BlobDownloadResponse; +import com.azure.storage.blob.models.BlobErrorCode; import com.azure.storage.blob.models.BlobHttpHeaders; import com.azure.storage.blob.models.BlobRequestConditions; import com.azure.storage.blob.models.BlobStorageException; @@ -20,9 +21,8 @@ import java.io.InputStream; import java.nio.charset.StandardCharsets; import java.util.UUID; -import java.util.concurrent.CountDownLatch; -import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; +import java.util.concurrent.locks.ReentrantLock; import java.util.regex.Pattern; import java.util.zip.GZIPInputStream; import java.util.zip.GZIPOutputStream; @@ -32,6 +32,7 @@ *

* Stores payloads as blobs and returns opaque tokens in the form {@code blob:v1::}. * Supports optional gzip compression. The blob container is created automatically on first upload. + * If the container is subsequently deleted, a new upload recreates it and retries once. */ public final class BlobPayloadStore extends PayloadStore { @@ -48,12 +49,9 @@ public final class BlobPayloadStore extends PayloadStore { private final BlobContainerClient containerClient; private final LargePayloadStorageOptions options; - // Container-creation guard. The first thread to call ensureContainerExists() creates the - // latch and performs the RPC. Concurrent callers await the latch so they don't race ahead - // and upload to a container that hasn't been created yet. On failure the reference is - // reset to null so a subsequent call can retry. - private final AtomicReference containerLatch = new AtomicReference<>(); - private volatile boolean containerVerified; + private final ReentrantLock containerInitializationLock = new ReentrantLock(); + // A stale upload failure must not invalidate a container another upload has already recreated. + private final AtomicReference containerGeneration = new AtomicReference<>(); /** * Creates a new {@code BlobPayloadStore} from the given options. @@ -123,118 +121,87 @@ public String upload(String payload) { byte[] payloadBytes = payload.getBytes(StandardCharsets.UTF_8); - // Ensure container exists before uploading. Thread-safe: the first caller creates - // the container while concurrent callers wait for it to complete. - ensureContainerExists(); - - try { - // Defense-in-depth: require the blob to not already exist (If-None-Match: *). - // Blob names are random UUIDs so collisions are astronomically unlikely, but this - // guards against future regressions (e.g. a caller-supplied PayloadStore that - // generates deterministic names or a refactor that reuses names) by failing loudly - // instead of silently overwriting someone else's payload. - BlobRequestConditions conditions = new BlobRequestConditions().setIfNoneMatch("*"); - if (this.options.isCompressionEnabled()) { - ByteArrayOutputStream compressedBuffer = new ByteArrayOutputStream(); - try (GZIPOutputStream gzip = new GZIPOutputStream(compressedBuffer)) { - gzip.write(payloadBytes); - } - byte[] compressedBytes = compressedBuffer.toByteArray(); - BlobHttpHeaders headers = new BlobHttpHeaders().setContentEncoding(CONTENT_ENCODING_GZIP); - try (InputStream stream = new ByteArrayInputStream(compressedBytes)) { - blob.uploadWithResponse( - stream, - compressedBytes.length, - null, // parallelTransferOptions - headers, - null, // metadata - null, // tier - conditions, // requestConditions - null, // timeout - Context.NONE); + boolean retryAfterContainerNotFound = true; + while (true) { + Object generation = ensureContainerExists(); + try { + uploadBlob(blob, payloadBytes); + return encodeToken(this.containerClient.getBlobContainerName(), blobName); + } catch (BlobStorageException e) { + if (retryAfterContainerNotFound + && e.getStatusCode() == 404 + && BlobErrorCode.CONTAINER_NOT_FOUND.equals(e.getErrorCode())) { + this.containerGeneration.compareAndSet(generation, null); + retryAfterContainerNotFound = false; + continue; } - } else { - try (InputStream stream = new ByteArrayInputStream(payloadBytes)) { - blob.uploadWithResponse( - stream, - payloadBytes.length, - null, // parallelTransferOptions - null, // headers - null, // metadata - null, // tier - conditions, // requestConditions - null, // timeout - Context.NONE); + + if (e.getStatusCode() == 409 || e.getStatusCode() == 412) { + throw new PayloadStorageException( + "Payload blob '" + blobName + "' already exists in container '" + + this.containerClient.getBlobContainerName() + + "'. Refusing to overwrite. This should not happen with random UUID blob names " + + "and likely indicates a bug in a custom PayloadStore implementation.", e); } + throw new PayloadStorageException("Failed to upload payload blob '" + blobName + "'.", e); + } catch (IOException e) { + throw new PayloadStorageException("Failed to upload payload blob '" + blobName + "'.", e); } - } catch (BlobStorageException e) { - // 409 BlobAlreadyExists and 412 ConditionNotMet (from If-None-Match: *) both indicate - // a name collision on upload — treat as a hard failure rather than silently overwriting. - if (e.getStatusCode() == 409 || e.getStatusCode() == 412) { - throw new PayloadStorageException( - "Payload blob '" + blobName + "' already exists in container '" - + this.containerClient.getBlobContainerName() - + "'. Refusing to overwrite. This should not happen with random UUID blob names " - + "and likely indicates a bug in a custom PayloadStore implementation.", e); - } - throw new PayloadStorageException("Failed to upload payload blob '" + blobName + "'.", e); - } catch (IOException e) { - throw new PayloadStorageException("Failed to upload payload blob '" + blobName + "'.", e); } + } - return encodeToken(this.containerClient.getBlobContainerName(), blobName); + private void uploadBlob(BlobClient blob, byte[] payloadBytes) throws IOException { + BlobRequestConditions conditions = new BlobRequestConditions().setIfNoneMatch("*"); + if (this.options.isCompressionEnabled()) { + ByteArrayOutputStream compressedBuffer = new ByteArrayOutputStream(); + try (GZIPOutputStream gzip = new GZIPOutputStream(compressedBuffer)) { + gzip.write(payloadBytes); + } + byte[] compressedBytes = compressedBuffer.toByteArray(); + BlobHttpHeaders headers = new BlobHttpHeaders().setContentEncoding(CONTENT_ENCODING_GZIP); + try (InputStream stream = new ByteArrayInputStream(compressedBytes)) { + blob.uploadWithResponse( + stream, compressedBytes.length, null, headers, null, null, conditions, null, Context.NONE); + } + } else { + try (InputStream stream = new ByteArrayInputStream(payloadBytes)) { + blob.uploadWithResponse( + stream, payloadBytes.length, null, null, null, null, conditions, null, Context.NONE); + } + } } - /** - * Ensures the blob container exists, creating it if necessary. Thread-safe: the first - * caller performs the RPC while concurrent callers wait for it to complete. On success - * the check is skipped on all future calls. On failure the guard is reset so a later - * call can retry. - */ - private void ensureContainerExists() { - if (this.containerVerified) { - return; + private Object ensureContainerExists() { + Object generation = this.containerGeneration.get(); + if (generation != null) { + return generation; } - CountDownLatch latch = new CountDownLatch(1); - if (!this.containerLatch.compareAndSet(null, latch)) { - CountDownLatch existing = this.containerLatch.get(); - if (existing == null) { - // Rare race: the creating thread already reset the latch (failure path). - // Retry on the next upload call rather than proceeding without a container. - return; - } - // Another thread is already creating the container — wait for it. - try { - existing.await(); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - throw new PayloadStorageException("Interrupted while waiting for container creation.", e); - } - // If the creating thread failed, containerVerified is still false; the next - // upload attempt will retry. For now, return and let the upload proceed - // (it will fail fast with a clear error if the container doesn't exist). - return; + try { + this.containerInitializationLock.lockInterruptibly(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new PayloadStorageException("Interrupted while waiting for container creation.", e); } - // This thread is responsible for creating the container. try { - this.containerClient.createIfNotExists(); - this.containerVerified = true; - } catch (BlobStorageException e) { - if (e.getStatusCode() == 409) { - // 409 Conflict means it already exists — safe to ignore. - this.containerVerified = true; - } else { - this.containerLatch.set(null); // allow a future call to retry - throw new PayloadStorageException( - "Failed to create blob container '" + this.containerClient.getBlobContainerName() + "'.", e); + generation = this.containerGeneration.get(); + if (generation == null) { + try { + this.containerClient.createIfNotExists(); + } catch (BlobStorageException e) { + if (e.getStatusCode() != 409 + || !BlobErrorCode.CONTAINER_ALREADY_EXISTS.equals(e.getErrorCode())) { + throw new PayloadStorageException( + "Failed to create blob container '" + this.containerClient.getBlobContainerName() + "'.", e); + } + } + generation = new Object(); + this.containerGeneration.set(generation); } - } catch (RuntimeException e) { - this.containerLatch.set(null); // allow a future call to retry - throw e; + return generation; } finally { - latch.countDown(); // unblock waiting threads + this.containerInitializationLock.unlock(); } } diff --git a/azure-blob-payloads/src/test/java/com/microsoft/durabletask/azureblobpayloads/BlobPayloadStoreTest.java b/azure-blob-payloads/src/test/java/com/microsoft/durabletask/azureblobpayloads/BlobPayloadStoreTest.java index bfa7c140..2154d5d5 100644 --- a/azure-blob-payloads/src/test/java/com/microsoft/durabletask/azureblobpayloads/BlobPayloadStoreTest.java +++ b/azure-blob-payloads/src/test/java/com/microsoft/durabletask/azureblobpayloads/BlobPayloadStoreTest.java @@ -6,12 +6,16 @@ import com.azure.storage.blob.BlobContainerClient; import com.azure.storage.blob.models.BlobDownloadHeaders; import com.azure.storage.blob.models.BlobDownloadResponse; +import com.azure.storage.blob.models.BlobErrorCode; import com.azure.storage.blob.models.BlobHttpHeaders; import com.azure.storage.blob.models.BlobRequestConditions; import com.azure.storage.blob.models.BlobStorageException; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.CsvSource; +import org.junit.jupiter.params.provider.ValueSource; import org.mockito.ArgumentCaptor; import java.io.ByteArrayInputStream; @@ -19,6 +23,12 @@ import java.io.InputStream; import java.io.OutputStream; import java.nio.charset.StandardCharsets; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; import java.util.zip.GZIPOutputStream; import static org.junit.jupiter.api.Assertions.*; @@ -134,9 +144,9 @@ void upload_tokenContainsGuidBlobName() { @Test void upload_containerAlreadyExists_succeeds() { - // createIfNotExists throwing 409 should be silently ignored BlobStorageException conflict = mock(BlobStorageException.class); when(conflict.getStatusCode()).thenReturn(409); + when(conflict.getErrorCode()).thenReturn(BlobErrorCode.CONTAINER_ALREADY_EXISTS); doThrow(conflict).when(mockContainerClient).createIfNotExists(); BlobPayloadStore store = new BlobPayloadStore(mockContainerClient, options); @@ -157,6 +167,140 @@ void upload_containerCreationFails_throwsPayloadStorageException() { assertThrows(PayloadStorageException.class, () -> store.upload("test")); } + @ParameterizedTest + @ValueSource(booleans = {true, false}) + void upload_deletedContainer_recreatesAndRetriesSameBlob(boolean compressed) { + options.setCompressionEnabled(compressed); + BlobPayloadStore store = new BlobPayloadStore(mockContainerClient, options); + store.upload("first payload"); + + BlobStorageException missing = storageFailure(404, BlobErrorCode.CONTAINER_NOT_FOUND); + doThrow(missing).doNothing().when(mockBlobClient).uploadWithResponse( + any(InputStream.class), anyLong(), isNull(), nullable(BlobHttpHeaders.class), + isNull(), isNull(), any(BlobRequestConditions.class), isNull(), any()); + + String token = store.upload("payload after deletion"); + + assertTrue(token.startsWith("blob:v1:durabletask-payloads:")); + verify(mockContainerClient, times(2)).createIfNotExists(); + verify(mockContainerClient, times(2)).getBlobClient(anyString()); + verify(mockBlobClient, times(3)).uploadWithResponse( + any(InputStream.class), anyLong(), isNull(), nullable(BlobHttpHeaders.class), + isNull(), isNull(), any(BlobRequestConditions.class), isNull(), any()); + } + + @Test + void upload_containerStillMissing_retriesOnlyOnce() { + BlobStorageException missing = storageFailure(404, BlobErrorCode.CONTAINER_NOT_FOUND); + doThrow(missing).when(mockBlobClient).uploadWithResponse( + any(InputStream.class), anyLong(), isNull(), any(BlobHttpHeaders.class), + isNull(), isNull(), any(BlobRequestConditions.class), isNull(), any()); + BlobPayloadStore store = new BlobPayloadStore(mockContainerClient, options); + + PayloadStorageException failure = assertThrows(PayloadStorageException.class, () -> store.upload("payload")); + + assertSame(missing, failure.getCause()); + verify(mockContainerClient, times(2)).createIfNotExists(); + verify(mockContainerClient).getBlobClient(anyString()); + verify(mockBlobClient, times(2)).uploadWithResponse( + any(InputStream.class), anyLong(), isNull(), any(BlobHttpHeaders.class), + isNull(), isNull(), any(BlobRequestConditions.class), isNull(), any()); + } + + @ParameterizedTest + @CsvSource({"404,BlobNotFound", "404,", "403,AuthorizationFailure", "500,InternalError"}) + void upload_otherStorageFailures_doNotRecreateContainer(int status, String errorCode) { + BlobStorageException storageFailure = storageFailure( + status, errorCode == null ? null : BlobErrorCode.fromString(errorCode)); + doThrow(storageFailure).when(mockBlobClient).uploadWithResponse( + any(InputStream.class), anyLong(), isNull(), any(BlobHttpHeaders.class), + isNull(), isNull(), any(BlobRequestConditions.class), isNull(), any()); + BlobPayloadStore store = new BlobPayloadStore(mockContainerClient, options); + + PayloadStorageException failure = assertThrows(PayloadStorageException.class, () -> store.upload("payload")); + + assertSame(storageFailure, failure.getCause()); + verify(mockContainerClient).createIfNotExists(); + verify(mockBlobClient).uploadWithResponse( + any(InputStream.class), anyLong(), isNull(), any(BlobHttpHeaders.class), + isNull(), isNull(), any(BlobRequestConditions.class), isNull(), any()); + } + + @Test + void upload_containerBeingDeleted_doesNotCacheFailedInitialization() { + BlobStorageException deleting = storageFailure(409, BlobErrorCode.CONTAINER_BEING_DELETED); + doThrow(deleting).doReturn(true).when(mockContainerClient).createIfNotExists(); + BlobPayloadStore store = new BlobPayloadStore(mockContainerClient, options); + + PayloadStorageException failure = assertThrows(PayloadStorageException.class, () -> store.upload("payload")); + assertSame(deleting, failure.getCause()); + store.upload("later payload"); + + verify(mockContainerClient, times(2)).createIfNotExists(); + verify(mockBlobClient).uploadWithResponse( + any(InputStream.class), anyLong(), isNull(), any(BlobHttpHeaders.class), + isNull(), isNull(), any(BlobRequestConditions.class), isNull(), any()); + } + + @Test + void upload_staleContainerNotFound_doesNotInvalidateNewGeneration() throws Exception { + options.setCompressionEnabled(false); + BlobPayloadStore store = new BlobPayloadStore(mockContainerClient, options); + store.upload("initial payload"); + + CountDownLatch slowUploadStarted = new CountDownLatch(1); + CountDownLatch releaseSlowFailure = new CountDownLatch(1); + AtomicInteger slowAttempts = new AtomicInteger(); + AtomicInteger fastAttempts = new AtomicInteger(); + BlobStorageException missing = storageFailure(404, BlobErrorCode.CONTAINER_NOT_FOUND); + doAnswer(invocation -> { + InputStream stream = invocation.getArgument(0); + String payload = new String(stream.readAllBytes(), StandardCharsets.UTF_8); + if (payload.equals("slow") && slowAttempts.incrementAndGet() == 1) { + slowUploadStarted.countDown(); + assertTrue(releaseSlowFailure.await(5, TimeUnit.SECONDS)); + throw missing; + } + if (payload.equals("fast") && fastAttempts.incrementAndGet() == 1) { + throw missing; + } + return null; + }).when(mockBlobClient).uploadWithResponse( + any(InputStream.class), anyLong(), isNull(), isNull(), + isNull(), isNull(), any(BlobRequestConditions.class), isNull(), any()); + + ExecutorService executor = Executors.newFixedThreadPool(2); + try { + Future slow = executor.submit(() -> store.upload("slow")); + assertTrue(slowUploadStarted.await(5, TimeUnit.SECONDS)); + Future fast = executor.submit(() -> store.upload("fast")); + assertNotNull(fast.get(5, TimeUnit.SECONDS)); + releaseSlowFailure.countDown(); + assertNotNull(slow.get(5, TimeUnit.SECONDS)); + + assertEquals(2, slowAttempts.get()); + assertEquals(2, fastAttempts.get()); + verify(mockContainerClient, times(2)).createIfNotExists(); + } finally { + releaseSlowFailure.countDown(); + executor.shutdownNow(); + assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS)); + } + } + + @Test + void upload_interruptedInitialization_preservesInterrupt() { + BlobPayloadStore store = new BlobPayloadStore(mockContainerClient, options); + Thread.currentThread().interrupt(); + try { + assertThrows(PayloadStorageException.class, () -> store.upload("payload")); + assertTrue(Thread.currentThread().isInterrupted()); + verify(mockContainerClient, never()).createIfNotExists(); + } finally { + Thread.interrupted(); + } + } + // ==================== Download tests ==================== @Test @@ -389,6 +533,13 @@ void upload_compressionProducesGzipHeaders() { // ==================== Test helpers ==================== + private static BlobStorageException storageFailure(int status, BlobErrorCode errorCode) { + BlobStorageException failure = mock(BlobStorageException.class); + when(failure.getStatusCode()).thenReturn(status); + when(failure.getErrorCode()).thenReturn(errorCode); + return failure; + } + /** * Sets up the mockBlobClient to return a downloadStreamWithResponse that writes * the given bytes and returns headers with the specified content-encoding. diff --git a/client/src/main/java/com/microsoft/durabletask/FailureDetails.java b/client/src/main/java/com/microsoft/durabletask/FailureDetails.java index f7c634a2..e9f8fd75 100644 --- a/client/src/main/java/com/microsoft/durabletask/FailureDetails.java +++ b/client/src/main/java/com/microsoft/durabletask/FailureDetails.java @@ -82,7 +82,7 @@ public static FailureDetails fromException( FailureDetails(TaskFailureDetails proto) { this(proto.getErrorType(), proto.getErrorMessage(), - proto.getStackTrace().getValue(), + proto.hasStackTrace() ? proto.getStackTrace().getValue() : null, proto.getIsNonRetriable(), proto.hasInnerFailure() ? new FailureDetails(proto.getInnerFailure()) : null, convertProtoProperties(proto.getPropertiesMap())); @@ -111,10 +111,9 @@ public String getErrorMessage() { } /** - * Gets the stack trace of the exception that caused this failure, or {@code null} if the failure was caused by - * a non-exception error. + * Gets the stack trace of the exception that caused this failure, or {@code null} if no stack trace was provided. * - * @return the stack trace of the failure exception or {@code null} if the failure was not caused by an exception + * @return the stack trace of the failure exception or {@code null} if no stack trace was provided */ @Nullable public String getStackTrace() { @@ -210,6 +209,7 @@ static String getFullStackTrace(Throwable e) { /** * Converts this failure to its protocol representation. + * Missing stack traces remain absent; explicitly empty stack traces remain present. * * @return the protocol representation of this failure */ @@ -218,9 +218,12 @@ public TaskFailureDetails toProto() { TaskFailureDetails.Builder builder = TaskFailureDetails.newBuilder() .setErrorType(this.getErrorType()) .setErrorMessage(this.getErrorMessage()) - .setStackTrace(StringValue.of(this.getStackTrace() != null ? this.getStackTrace() : "")) .setIsNonRetriable(this.isNonRetriable); + if (this.stackTrace != null) { + builder.setStackTrace(StringValue.of(this.stackTrace)); + } + if (this.innerFailure != null) { builder.setInnerFailure(this.innerFailure.toProto()); } diff --git a/client/src/main/java/com/microsoft/durabletask/Helpers.java b/client/src/main/java/com/microsoft/durabletask/Helpers.java index deed6e02..ee72b96a 100644 --- a/client/src/main/java/com/microsoft/durabletask/Helpers.java +++ b/client/src/main/java/com/microsoft/durabletask/Helpers.java @@ -2,9 +2,12 @@ // Licensed under the MIT License. package com.microsoft.durabletask; +import com.google.protobuf.Timestamp; + import javax.annotation.Nonnull; import javax.annotation.Nullable; import java.time.Duration; +import java.time.Instant; final class Helpers { final static Duration maxDuration = Duration.ofSeconds(Long.MAX_VALUE, 999999999L); @@ -60,6 +63,11 @@ static boolean isNullOrEmpty(String s) { return s == null || s.isEmpty(); } + // History and metadata retain precision without changing DataConverter's replay-time conversion. + static @Nullable Instant getPreciseInstantFromTimestamp(@Nullable Timestamp timestamp) { + return timestamp == null ? null : Instant.ofEpochSecond(timestamp.getSeconds(), timestamp.getNanos()); + } + // Cannot be instantiated private Helpers() { } diff --git a/client/src/main/java/com/microsoft/durabletask/HistoryEventConverter.java b/client/src/main/java/com/microsoft/durabletask/HistoryEventConverter.java index 53e3a8d6..b1f94ea7 100644 --- a/client/src/main/java/com/microsoft/durabletask/HistoryEventConverter.java +++ b/client/src/main/java/com/microsoft/durabletask/HistoryEventConverter.java @@ -25,7 +25,7 @@ private HistoryEventConverter() { */ static HistoryEvent fromProto(OrchestratorService.HistoryEvent proto) { int id = proto.getEventId(); - Instant ts = DataConverter.getInstantFromTimestamp(proto.getTimestamp()); + Instant ts = Helpers.getPreciseInstantFromTimestamp(proto.getTimestamp()); switch (proto.getEventTypeCase()) { case EXECUTIONSTARTED: { OrchestratorService.ExecutionStartedEvent p = proto.getExecutionStarted(); @@ -36,7 +36,7 @@ static HistoryEvent fromProto(OrchestratorService.HistoryEvent proto) { p.hasOrchestrationInstance() ? toInstance(p.getOrchestrationInstance()) : null, p.hasParentInstance() ? toParentInfo(p.getParentInstance()) : null, p.hasScheduledStartTimestamp() - ? DataConverter.getInstantFromTimestamp(p.getScheduledStartTimestamp()) : null, + ? Helpers.getPreciseInstantFromTimestamp(p.getScheduledStartTimestamp()) : null, p.hasParentTraceContext() ? toTrace(p.getParentTraceContext()) : null, stringOrNull(p.hasOrchestrationSpanID(), p.getOrchestrationSpanID()), p.getTagsMap()); @@ -93,11 +93,11 @@ static HistoryEvent fromProto(OrchestratorService.HistoryEvent proto) { } case TIMERCREATED: { OrchestratorService.TimerCreatedEvent p = proto.getTimerCreated(); - return new TimerCreatedEvent(id, ts, DataConverter.getInstantFromTimestamp(p.getFireAt())); + return new TimerCreatedEvent(id, ts, Helpers.getPreciseInstantFromTimestamp(p.getFireAt())); } case TIMERFIRED: { OrchestratorService.TimerFiredEvent p = proto.getTimerFired(); - return new TimerFiredEvent(id, ts, DataConverter.getInstantFromTimestamp(p.getFireAt()), p.getTimerId()); + return new TimerFiredEvent(id, ts, Helpers.getPreciseInstantFromTimestamp(p.getFireAt()), p.getTimerId()); } case ORCHESTRATORSTARTED: return new OrchestratorStartedEvent(id, ts); @@ -138,7 +138,7 @@ static HistoryEvent fromProto(OrchestratorService.HistoryEvent proto) { return new EntityOperationSignaledEvent(id, ts, p.getRequestId(), p.getOperation(), - p.hasScheduledTime() ? DataConverter.getInstantFromTimestamp(p.getScheduledTime()) : null, + p.hasScheduledTime() ? Helpers.getPreciseInstantFromTimestamp(p.getScheduledTime()) : null, stringOrNull(p.hasInput(), p.getInput()), stringOrNull(p.hasTargetInstanceId(), p.getTargetInstanceId())); } @@ -147,7 +147,7 @@ static HistoryEvent fromProto(OrchestratorService.HistoryEvent proto) { return new EntityOperationCalledEvent(id, ts, p.getRequestId(), p.getOperation(), - p.hasScheduledTime() ? DataConverter.getInstantFromTimestamp(p.getScheduledTime()) : null, + p.hasScheduledTime() ? Helpers.getPreciseInstantFromTimestamp(p.getScheduledTime()) : null, stringOrNull(p.hasInput(), p.getInput()), stringOrNull(p.hasParentInstanceId(), p.getParentInstanceId()), stringOrNull(p.hasParentExecutionId(), p.getParentExecutionId()), @@ -231,10 +231,10 @@ private static OrchestrationState toOrchestrationState(OrchestratorService.Orche stringOrNull(p.hasVersion(), p.getVersion()), OrchestrationRuntimeStatus.fromProtobuf(p.getOrchestrationStatus()), p.hasScheduledStartTimestamp() - ? DataConverter.getInstantFromTimestamp(p.getScheduledStartTimestamp()) : null, - p.hasCreatedTimestamp() ? DataConverter.getInstantFromTimestamp(p.getCreatedTimestamp()) : null, - p.hasLastUpdatedTimestamp() ? DataConverter.getInstantFromTimestamp(p.getLastUpdatedTimestamp()) : null, - p.hasCompletedTimestamp() ? DataConverter.getInstantFromTimestamp(p.getCompletedTimestamp()) : null, + ? Helpers.getPreciseInstantFromTimestamp(p.getScheduledStartTimestamp()) : null, + p.hasCreatedTimestamp() ? Helpers.getPreciseInstantFromTimestamp(p.getCreatedTimestamp()) : null, + p.hasLastUpdatedTimestamp() ? Helpers.getPreciseInstantFromTimestamp(p.getLastUpdatedTimestamp()) : null, + p.hasCompletedTimestamp() ? Helpers.getPreciseInstantFromTimestamp(p.getCompletedTimestamp()) : null, stringOrNull(p.hasInput(), p.getInput()), stringOrNull(p.hasOutput(), p.getOutput()), stringOrNull(p.hasCustomStatus(), p.getCustomStatus()), diff --git a/client/src/main/java/com/microsoft/durabletask/OrchestrationMetadata.java b/client/src/main/java/com/microsoft/durabletask/OrchestrationMetadata.java index 471ee1f6..3c2a9938 100644 --- a/client/src/main/java/com/microsoft/durabletask/OrchestrationMetadata.java +++ b/client/src/main/java/com/microsoft/durabletask/OrchestrationMetadata.java @@ -50,8 +50,8 @@ public final class OrchestrationMetadata { this.name = state.getName(); this.instanceId = state.getInstanceId(); this.runtimeStatus = OrchestrationRuntimeStatus.fromProtobuf(state.getOrchestrationStatus()); - this.createdAt = DataConverter.getInstantFromTimestamp(state.getCreatedTimestamp()); - this.lastUpdatedAt = DataConverter.getInstantFromTimestamp(state.getLastUpdatedTimestamp()); + this.createdAt = Helpers.getPreciseInstantFromTimestamp(state.getCreatedTimestamp()); + this.lastUpdatedAt = Helpers.getPreciseInstantFromTimestamp(state.getLastUpdatedTimestamp()); this.serializedInput = state.getInput().getValue(); this.serializedOutput = state.getOutput().getValue(); this.serializedCustomStatus = state.getCustomStatus().getValue(); @@ -84,7 +84,7 @@ public OrchestrationRuntimeStatus getRuntimeStatus() { } /** - * Gets the orchestration instance's creation time in UTC. + * Gets the orchestration instance's creation time in UTC, preserving the backend's timestamp precision. * @return the orchestration instance's creation time in UTC */ public Instant getCreatedAt() { @@ -92,7 +92,7 @@ public Instant getCreatedAt() { } /** - * Gets the orchestration instance's last updated time in UTC. + * Gets the orchestration instance's last updated time in UTC, preserving the backend's timestamp precision. * @return the orchestration instance's last updated time in UTC */ public Instant getLastUpdatedAt() { diff --git a/client/src/test/java/com/microsoft/durabletask/FailureDetailsTest.java b/client/src/test/java/com/microsoft/durabletask/FailureDetailsTest.java index 280c45ed..54fdda6d 100644 --- a/client/src/test/java/com/microsoft/durabletask/FailureDetailsTest.java +++ b/client/src/test/java/com/microsoft/durabletask/FailureDetailsTest.java @@ -8,6 +8,9 @@ import com.google.protobuf.Value; import com.microsoft.durabletask.implementation.protobuf.OrchestratorService.TaskFailureDetails; import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.NullAndEmptySource; +import org.junit.jupiter.params.provider.ValueSource; import java.io.IOException; import java.util.HashMap; @@ -21,6 +24,33 @@ */ public class FailureDetailsTest { + @ParameterizedTest + @NullAndEmptySource + @ValueSource(strings = {"at App.run(App.java:5)"}) + void protoRoundTrip_preservesStackTracePresence(String stackTrace) { + TaskFailureDetails.Builder inner = TaskFailureDetails.newBuilder() + .setErrorType("Inner") + .setErrorMessage("inner failure"); + TaskFailureDetails.Builder outer = TaskFailureDetails.newBuilder() + .setErrorType("Outer") + .setErrorMessage("outer failure"); + if (stackTrace != null) { + inner.setStackTrace(StringValue.of(stackTrace)); + outer.setStackTrace(StringValue.of(stackTrace)); + } + TaskFailureDetails proto = outer.setInnerFailure(inner).build(); + + FailureDetails details = new FailureDetails(proto); + + assertEquals(stackTrace, details.getStackTrace()); + assertNotNull(details.getInnerFailure()); + assertEquals(stackTrace, details.getInnerFailure().getStackTrace()); + TaskFailureDetails roundTripped = details.toProto(); + assertEquals(stackTrace != null, roundTripped.hasStackTrace()); + assertEquals(stackTrace != null, roundTripped.getInnerFailure().hasStackTrace()); + assertEquals(proto, roundTripped); + } + @Test void constructFromProto_withInnerFailureAndProperties() { TaskFailureDetails innerProto = TaskFailureDetails.newBuilder() diff --git a/client/src/test/java/com/microsoft/durabletask/HistoryEventConverterTest.java b/client/src/test/java/com/microsoft/durabletask/HistoryEventConverterTest.java index e1a86a88..6078123f 100644 --- a/client/src/test/java/com/microsoft/durabletask/HistoryEventConverterTest.java +++ b/client/src/test/java/com/microsoft/durabletask/HistoryEventConverterTest.java @@ -55,12 +55,54 @@ public class HistoryEventConverterTest { private static final long EPOCH_SECONDS = 1_700_000_000L; - private static final Instant EXPECTED_TIMESTAMP = Instant.ofEpochSecond(EPOCH_SECONDS); + private static final int NANOSECONDS = 123_456_700; + private static final Instant EXPECTED_TIMESTAMP = Instant.ofEpochSecond(EPOCH_SECONDS, NANOSECONDS); private static OrchestratorService.HistoryEvent.Builder baseEvent(int eventId) { return OrchestratorService.HistoryEvent.newBuilder() .setEventId(eventId) - .setTimestamp(Timestamp.newBuilder().setSeconds(EPOCH_SECONDS).build()); + .setTimestamp(Timestamp.newBuilder().setSeconds(EPOCH_SECONDS).setNanos(NANOSECONDS).build()); + } + + @Test + void preservesPrecisionInNestedHistoryTimestamps() { + Timestamp precise = DataConverter.getTimestampFromInstant(EXPECTED_TIMESTAMP); + ExecutionStartedEvent started = (ExecutionStartedEvent) HistoryEventConverter.fromProto(baseEvent(1) + .setExecutionStarted(OrchestratorService.ExecutionStartedEvent.newBuilder() + .setScheduledStartTimestamp(precise)) + .build()); + TimerCreatedEvent created = (TimerCreatedEvent) HistoryEventConverter.fromProto(baseEvent(2) + .setTimerCreated(OrchestratorService.TimerCreatedEvent.newBuilder().setFireAt(precise)) + .build()); + TimerFiredEvent fired = (TimerFiredEvent) HistoryEventConverter.fromProto(baseEvent(3) + .setTimerFired(OrchestratorService.TimerFiredEvent.newBuilder().setFireAt(precise)) + .build()); + EntityOperationCalledEvent called = (EntityOperationCalledEvent) HistoryEventConverter.fromProto(baseEvent(4) + .setEntityOperationCalled(OrchestratorService.EntityOperationCalledEvent.newBuilder() + .setScheduledTime(precise)) + .build()); + EntityOperationSignaledEvent signaled = (EntityOperationSignaledEvent) HistoryEventConverter.fromProto(baseEvent(5) + .setEntityOperationSignaled(OrchestratorService.EntityOperationSignaledEvent.newBuilder() + .setScheduledTime(precise)) + .build()); + HistoryStateEvent historyState = (HistoryStateEvent) HistoryEventConverter.fromProto(baseEvent(6) + .setHistoryState(OrchestratorService.HistoryStateEvent.newBuilder() + .setOrchestrationState(OrchestratorService.OrchestrationState.newBuilder() + .setScheduledStartTimestamp(precise) + .setCreatedTimestamp(precise) + .setLastUpdatedTimestamp(precise) + .setCompletedTimestamp(precise))) + .build()); + + assertEquals(EXPECTED_TIMESTAMP, started.getScheduledStartTimestamp()); + assertEquals(EXPECTED_TIMESTAMP, created.getFireAt()); + assertEquals(EXPECTED_TIMESTAMP, fired.getFireAt()); + assertEquals(EXPECTED_TIMESTAMP, called.getScheduledTime()); + assertEquals(EXPECTED_TIMESTAMP, signaled.getScheduledTime()); + assertEquals(EXPECTED_TIMESTAMP, historyState.getState().getScheduledStartTime()); + assertEquals(EXPECTED_TIMESTAMP, historyState.getState().getCreatedTime()); + assertEquals(EXPECTED_TIMESTAMP, historyState.getState().getLastUpdatedTime()); + assertEquals(EXPECTED_TIMESTAMP, historyState.getState().getCompletedTime()); } @Test diff --git a/client/src/test/java/com/microsoft/durabletask/OrchestrationMetadataTest.java b/client/src/test/java/com/microsoft/durabletask/OrchestrationMetadataTest.java new file mode 100644 index 00000000..402e723c --- /dev/null +++ b/client/src/test/java/com/microsoft/durabletask/OrchestrationMetadataTest.java @@ -0,0 +1,40 @@ +// Copyright (c) Microsoft Corporation. All rights reserved. +// Licensed under the MIT License. +package com.microsoft.durabletask; + +import com.google.protobuf.Timestamp; +import com.microsoft.durabletask.implementation.protobuf.OrchestratorService; +import org.junit.jupiter.api.Test; + +import java.time.Instant; +import java.time.temporal.ChronoUnit; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +class OrchestrationMetadataTest { + + @Test + void preservesBackendTimestampPrecision() { + Instant created = Instant.parse("2026-06-30T12:00:00.123456700Z"); + Instant updated = Instant.parse("2026-06-30T12:05:00.765432100Z"); + OrchestratorService.OrchestrationState state = OrchestratorService.OrchestrationState.newBuilder() + .setInstanceId("instance-1") + .setName("Orchestrator") + .setCreatedTimestamp(DataConverter.getTimestampFromInstant(created)) + .setLastUpdatedTimestamp(DataConverter.getTimestampFromInstant(updated)) + .build(); + + OrchestrationMetadata metadata = new OrchestrationMetadata(state, new JacksonDataConverter(), false); + + assertEquals(created, metadata.getCreatedAt()); + assertEquals(updated, metadata.getLastUpdatedAt()); + } + + @Test + void replayTimestampConversionStillUsesMilliseconds() { + Instant precise = Instant.parse("2026-06-30T12:00:00.123456700Z"); + Timestamp timestamp = DataConverter.getTimestampFromInstant(precise); + + assertEquals(precise.truncatedTo(ChronoUnit.MILLIS), DataConverter.getInstantFromTimestamp(timestamp)); + } +} diff --git a/exporthistory/README.md b/exporthistory/README.md index 83b044af..200a7abd 100644 --- a/exporthistory/README.md +++ b/exporthistory/README.md @@ -3,9 +3,7 @@ Durable, resumable export of **terminal orchestration history** to Azure Blob Storage for the Durable Task Java SDK — for compliance, audit, and offline analysis before instances age out of the task hub. -This module is at parity with the .NET `Microsoft.DurableTask.ExportHistory` (preview) feature: a checkpointed -entity + orchestrator that pages terminal instances by completion window, fans out per-instance export activities, -and uploads serialized history (gzipped JSONL by default) to a customer-owned blob container. +This module uses the same checkpointed entity-and-orchestrator workflow as the .NET `Microsoft.DurableTask.ExportHistory` (preview) feature: it pages terminal instances by completion window, fans out per-instance export activities, and uploads serialized history (gzipped JSONL by default) to a customer-owned blob container. > **Status:** preview (`0.1.0`). @@ -78,17 +76,16 @@ supplied, all three are exported. ## Export format -Each blob holds the instance's full history, and the **blob body is byte-for-byte identical to the .NET -`Microsoft.DurableTask.ExportHistory` output** (pinned by a test against golden output captured from -`Microsoft.Azure.DurableTask.Core`): +Each blob holds the instance's full history using .NET-style history-event JSON. Golden fixtures cover core event serialization and the legacy entity-message encoding produced by the .NET SDK and `Microsoft.Azure.DurableTask.Core`: - **JSONL** (default, gzipped) is one JSON object per line; **JSON** is a single array. - Each event is `{"eventType": "...", , "eventId": N, "isPlayed": false, "timestamp": "..."}`. -- camelCase field names, null fields omitted, empty maps as `{}`, enum values in PascalCase (e.g. `"Completed"`), - timestamps as trimmed ISO-8601 ending in `Z`, and the same HTML-safe string escaping (`"` → `\u0022`, - `& < > ' +` and all non-ASCII → `\uXXXX`). +- camelCase field names, null event fields omitted, empty maps as `{}`, enum values in PascalCase (e.g. `"Completed"`), timestamps as trimmed ISO-8601 ending in `Z`, and HTML-safe string escaping (`"` → `\u0022`, `& < > ' +` and all non-ASCII → `\uXXXX`). +- Sub-millisecond timestamps are preserved when reading history and metadata. Exported timestamps and blob-name timestamps use .NET's 100-nanosecond precision; orchestration replay timestamp behavior is unchanged. +- Entity operations and locks are exported as `EventSent`/`EventRaised` events. Their `input` contains the reference SDK's JSON-encoded entity request or response, including parent orchestration context and nested failure details. This replaces the earlier Java-native `Entity...` event format; consumers of that preview format must update their event handling. The public history API continues returning typed Java entity events. +- JSONL uses LF line endings. Identical compressed bytes or identical fields across all SDK versions are not guaranteed. -Blob **names**: a lowercase-hex SHA-256 of `"|"` plus the format extension. +Blob **names**: a lowercase-hex SHA-256 of `"|"` plus the format extension. Re-exporting an instance previously named using millisecond-truncated metadata can produce a new blob name; existing export blobs are not renamed or removed. ## Backend requirement diff --git a/exporthistory/build.gradle b/exporthistory/build.gradle index e24fa0c6..60c122f9 100644 --- a/exporthistory/build.gradle +++ b/exporthistory/build.gradle @@ -47,6 +47,7 @@ dependencies { testImplementation 'org.mockito:mockito-core:5.21.0' testImplementation 'org.mockito:mockito-junit-jupiter:5.21.0' + testImplementation "io.grpc:grpc-inprocess:${grpcVersion}" testImplementation project(':azuremanaged') } diff --git a/exporthistory/src/main/java/com/microsoft/durabletask/exporthistory/ExportBlobNaming.java b/exporthistory/src/main/java/com/microsoft/durabletask/exporthistory/ExportBlobNaming.java index c0c1efcd..0c011998 100644 --- a/exporthistory/src/main/java/com/microsoft/durabletask/exporthistory/ExportBlobNaming.java +++ b/exporthistory/src/main/java/com/microsoft/durabletask/exporthistory/ExportBlobNaming.java @@ -43,8 +43,7 @@ static String blobFileName(Instant completedTimestamp, String instanceId, Export * Formats an instant as {@code yyyy-MM-ddTHH:mm:ss.fffffff+00:00} (seven fractional digits, explicit UTC * offset). The instant is treated as UTC. *

- * Note: instance timestamps are truncated to milliseconds upstream, so the sub-millisecond fractional digits - * are always zero here. + * Preserves sub-millisecond precision up to 100-nanosecond ticks, matching the reference export format. * * @param instant the timestamp (treated as UTC) * @return the formatted timestamp string diff --git a/exporthistory/src/main/java/com/microsoft/durabletask/exporthistory/HistoryEventSerializer.java b/exporthistory/src/main/java/com/microsoft/durabletask/exporthistory/HistoryEventSerializer.java index 8b0537b4..9beaf3ac 100644 --- a/exporthistory/src/main/java/com/microsoft/durabletask/exporthistory/HistoryEventSerializer.java +++ b/exporthistory/src/main/java/com/microsoft/durabletask/exporthistory/HistoryEventSerializer.java @@ -2,15 +2,11 @@ // Licensed under the MIT License. package com.microsoft.durabletask.exporthistory; -import com.fasterxml.jackson.annotation.JsonInclude; import com.fasterxml.jackson.core.JsonFactory; import com.fasterxml.jackson.core.JsonGenerator; -import com.fasterxml.jackson.core.JsonProcessingException; -import com.fasterxml.jackson.databind.ObjectMapper; -import com.fasterxml.jackson.databind.PropertyNamingStrategies; -import com.fasterxml.jackson.databind.SerializationFeature; -import com.fasterxml.jackson.databind.json.JsonMapper; -import com.fasterxml.jackson.databind.node.ObjectNode; +import com.fasterxml.jackson.core.SerializableString; +import com.fasterxml.jackson.core.io.CharacterEscapes; +import com.fasterxml.jackson.core.io.SerializedString; import com.microsoft.durabletask.FailureDetails; import com.microsoft.durabletask.OrchestrationRuntimeStatus; import com.microsoft.durabletask.history.ContinueAsNewEvent; @@ -52,6 +48,8 @@ import java.time.OffsetDateTime; import java.time.ZoneOffset; import java.time.format.DateTimeFormatter; +import java.util.ArrayList; +import java.util.Collections; import java.util.LinkedHashMap; import java.util.List; import java.util.Locale; @@ -63,13 +61,13 @@ * Each event is written as a single JSON object with a leading {@code eventType} discriminator, the type-specific * fields, and a trailing {@code eventId}/{@code isPlayed}/{@code timestamp}: camelCase field names, null fields * omitted, empty maps rendered as {@code {}}, enum values in PascalCase, timestamps as trimmed ISO-8601 with a - * {@code Z} suffix, and (for non-entity events) strings escaped by {@link HtmlSafeJsonEscapes}. + * {@code Z} suffix, and strings escaped by {@link HtmlSafeJsonEscapes}. *

* {@link ExportFormatKind#JSONL} emits one object per line (gzip applied by the blob writer); * {@link ExportFormatKind#JSON} emits a single JSON array. *

- * Entity events have no dedicated representation in this wire format, so they fall back to a Java-native shape: a - * reflective projection of the event with an added {@code eventType} discriminator (serialized with Jackson defaults). + * Entity messages use the reference SDK's legacy {@code EventSent}/{@code EventRaised} representation, including + * the JSON-encoded request or response in the event's {@code input}. */ final class HistoryEventSerializer { @@ -78,13 +76,29 @@ final class HistoryEventSerializer { private static final JsonFactory FACTORY = new JsonFactory(); - // Fallback for entity events, which have no dedicated wire-format representation. - private static final ObjectMapper LEGACY_MAPPER = JsonMapper.builder() - .findAndAddModules() - .propertyNamingStrategy(PropertyNamingStrategies.LOWER_CAMEL_CASE) - .serializationInclusion(JsonInclude.Include.NON_NULL) - .disable(SerializationFeature.WRITE_DATES_AS_TIMESTAMPS) - .build(); + // Entity message JSON uses Newtonsoft's escaping, before the outer event applies HTML-safe escaping. + private static final CharacterEscapes ENTITY_MESSAGE_ESCAPES = new CharacterEscapes() { + private static final long serialVersionUID = 1L; + + @Override + public int[] getEscapeCodesForAscii() { + int[] escapes = CharacterEscapes.standardAsciiEscapesForJSON(); + // Short escapes have positive b/t/n/f/r codes; ESCAPE_STANDARD only marks Unicode escapes. + for (int i = 0; i < 0x20; i++) { + if (escapes[i] == CharacterEscapes.ESCAPE_STANDARD) { + escapes[i] = CharacterEscapes.ESCAPE_CUSTOM; + } + } + return escapes; + } + + @Override + public SerializableString getEscapeSequence(int ch) { + return ch < 0x20 || ch == 0x85 || ch == 0x2028 || ch == 0x2029 + ? new SerializedString(String.format(Locale.ROOT, "\\u%04x", ch)) + : null; + } + }; private HistoryEventSerializer() { } @@ -100,20 +114,22 @@ private HistoryEventSerializer() { * @param historyEvents the ordered history events * @param format the export format * @return the serialized content (JSONL text or JSON array text) - * @throws JsonProcessingException if serialization of an entity event fails + * @throws IllegalArgumentException if an entity message is missing required history context or fields */ - static String serialize(List historyEvents, ExportFormat format) - throws JsonProcessingException { + static String serialize(List historyEvents, ExportFormat format) { StringBuilder sb = new StringBuilder(); boolean json = format.getKind() == ExportFormatKind.JSON; + OrchestrationInstance currentInstance = null; if (json) { sb.append('['); } for (int i = 0; i < historyEvents.size(); i++) { HistoryEvent event = historyEvents.get(i); - String line = isEntityEvent(event) - ? writeEntity(event) - : writeObject(coreMap(event)); + if (event instanceof ExecutionStartedEvent) { + currentInstance = ((ExecutionStartedEvent) event).getOrchestrationInstance(); + } + HistoryEvent exportEvent = isEntityEvent(event) ? toExportEvent(event, currentInstance) : event; + String line = writeObject(coreMap(exportEvent)); if (json) { if (i > 0) { sb.append(','); @@ -382,22 +398,149 @@ private static String formatInstant(Instant t) { return sb.toString(); } - // ---- ordered map -> JSON with the parity escaper ------------------------------------------- + private static HistoryEvent toExportEvent(HistoryEvent event, OrchestrationInstance currentInstance) { + Map message = new LinkedHashMap<>(); + if (event instanceof EntityOperationCalledEvent) { + EntityOperationCalledEvent e = (EntityOperationCalledEvent) event; + message = operationMessage(e.getOperation(), false, e.getInput(), e.getRequestId(), + e.getScheduledTime(), requireCurrentInstance(currentInstance)); + return entityEventSent(e, e.getTargetInstanceId(), operationEventName(e.getScheduledTime()), message); + } else if (event instanceof EntityOperationSignaledEvent) { + EntityOperationSignaledEvent e = (EntityOperationSignaledEvent) event; + message = operationMessage(e.getOperation(), true, e.getInput(), e.getRequestId(), + e.getScheduledTime(), null); + return entityEventSent(e, e.getTargetInstanceId(), operationEventName(e.getScheduledTime()), message); + } else if (event instanceof EntityOperationCompletedEvent) { + EntityOperationCompletedEvent e = (EntityOperationCompletedEvent) event; + message.put("result", e.getOutput()); + return entityEventRaised(e, e.getRequestId(), message); + } else if (event instanceof EntityOperationFailedEvent) { + EntityOperationFailedEvent e = (EntityOperationFailedEvent) event; + FailureDetails failure = e.getFailureDetails(); + if (failure == null) { + throw new IllegalArgumentException("EntityOperationFailed history is missing failure details."); + } + message.put("result", failure.getErrorMessage()); + putIfNotNull(message, "exceptionType", failure.getErrorType()); + message.put("failureDetails", entityFailureMap(failure)); + return entityEventRaised(e, e.getRequestId(), message); + } else if (event instanceof EntityLockRequestedEvent) { + EntityLockRequestedEvent e = (EntityLockRequestedEvent) event; + if (e.getPosition() < 0 || e.getPosition() >= e.getLockSet().size()) { + throw new IllegalArgumentException("Entity lock position must identify a member of the lock set."); + } + message.put("op", null); + message.put("id", e.getCriticalSectionId()); + message.put("parent", requireCurrentInstance(currentInstance).getInstanceId()); + List> lockSet = new ArrayList<>(); + for (String entityId : e.getLockSet()) { + lockSet.add(entityIdMap(entityId)); + } + message.put("lockset", lockSet); + if (e.getPosition() != 0) { + message.put("pos", e.getPosition()); + } + return entityEventSent(e, e.getLockSet().get(e.getPosition()), "op", message); + } else if (event instanceof EntityLockGrantedEvent) { + EntityLockGrantedEvent e = (EntityLockGrantedEvent) event; + message.put("result", "Lock Acquisition Completed"); + return entityEventRaised(e, e.getCriticalSectionId(), message); + } else if (event instanceof EntityUnlockSentEvent) { + EntityUnlockSentEvent e = (EntityUnlockSentEvent) event; + message.put("parent", requireCurrentInstance(currentInstance).getInstanceId()); + message.put("id", e.getCriticalSectionId()); + return entityEventSent(e, e.getTargetInstanceId(), "release", message); + } + throw new IllegalArgumentException("Unsupported entity history event: " + event.getClass().getSimpleName()); + } + + private static Map operationMessage( + String operation, boolean signal, String input, String requestId, + Instant scheduledTime, OrchestrationInstance parent) { + Map message = new LinkedHashMap<>(); + message.put("op", operation); + if (signal) { + message.put("signal", true); + } + putIfNotNull(message, "input", input); + message.put("id", requestId); + if (parent != null) { + message.put("parent", parent.getInstanceId()); + putIfNotNull(message, "parentExecution", parent.getExecutionId()); + } + putIfNotNull(message, "due", formatInstantOrNull(scheduledTime)); + return message; + } + + private static OrchestrationInstance requireCurrentInstance(OrchestrationInstance instance) { + if (instance == null || instance.getInstanceId() == null) { + throw new IllegalArgumentException( + "Entity history export requires an ExecutionStarted event with an orchestration instance."); + } + return instance; + } + + private static Map entityIdMap(String entityId) { + int separator = entityId == null ? -1 : entityId.indexOf('@', 1); + if (entityId == null || !entityId.startsWith("@") || separator < 2) { + throw new IllegalArgumentException("Invalid entity ID in exported lock set: " + entityId); + } + // DT Core preserves name casing and permits empty keys, unlike the SDK's EntityInstanceId. + Map result = new LinkedHashMap<>(); + result.put("name", entityId.substring(1, separator)); + result.put("key", entityId.substring(separator + 1)); + return result; + } - private static String writeEntity(HistoryEvent event) throws JsonProcessingException { - // Entity events have no wire-format equivalent; project the event reflectively and prepend an eventType. - ObjectNode node = LEGACY_MAPPER.valueToTree(event); - ObjectNode withType = LEGACY_MAPPER.createObjectNode(); - withType.put("eventType", eventType(event)); - withType.setAll(node); - return LEGACY_MAPPER.writeValueAsString(withType); + private static Map entityFailureMap(FailureDetails failure) { + if (failure == null) { + return null; + } + Map result = new LinkedHashMap<>(); + result.put("ErrorType", failure.getErrorType()); + result.put("ErrorMessage", failure.getErrorMessage()); + result.put("StackTrace", failure.getStackTrace()); + result.put("InnerFailure", entityFailureMap(failure.getInnerFailure())); + result.put("IsNonRetriable", failure.isNonRetriable()); + result.put("Properties", failure.getProperties() == null ? Collections.emptyMap() : failure.getProperties()); + return result; + } + + private static String operationEventName(Instant scheduledTime) { + if (scheduledTime == null) { + return "op"; + } + OffsetDateTime utc = scheduledTime.atOffset(ZoneOffset.UTC); + return "op@" + DATE_TIME.format(utc) + + String.format(Locale.ROOT, ".%07dZ", utc.getNano() / 100); } + private static HistoryEvent entityEventSent( + HistoryEvent event, String target, String name, Map message) { + if (target == null) { + throw new IllegalArgumentException("Entity history export requires a target instance ID."); + } + return new EventSentEvent(event.getEventId(), event.getTimestamp(), target, name, + writeObject(message, ENTITY_MESSAGE_ESCAPES, 0)); + } + + private static HistoryEvent entityEventRaised(HistoryEvent event, String name, Map message) { + return new EventRaisedEvent(event.getEventId(), event.getTimestamp(), name, + writeObject(message, ENTITY_MESSAGE_ESCAPES, 0)); + } + + // ---- ordered map -> JSON with the parity escaper ------------------------------------------- + private static String writeObject(Map map) { + return writeObject(map, HtmlSafeJsonEscapes.INSTANCE, 0x7F); + } + + private static String writeObject( + Map map, CharacterEscapes escapes, int highestNonEscapedChar) { StringWriter sw = new StringWriter(); try (JsonGenerator g = FACTORY.createGenerator(sw)) { - g.setCharacterEscapes(HtmlSafeJsonEscapes.INSTANCE); - g.setHighestNonEscapedChar(0x7F); + g.setCharacterEscapes(escapes); + g.setHighestNonEscapedChar(highestNonEscapedChar); writeMap(g, map); } catch (IOException ex) { throw new UncheckedIOException(ex); diff --git a/exporthistory/src/test/java/com/microsoft/durabletask/exporthistory/ExportBlobNamingTest.java b/exporthistory/src/test/java/com/microsoft/durabletask/exporthistory/ExportBlobNamingTest.java index 1c3c805d..9eaf5026 100644 --- a/exporthistory/src/test/java/com/microsoft/durabletask/exporthistory/ExportBlobNamingTest.java +++ b/exporthistory/src/test/java/com/microsoft/durabletask/exporthistory/ExportBlobNamingTest.java @@ -62,6 +62,15 @@ void formatTimestamp_hasSevenFractionalDigitsAndUtcOffset() { assertEquals("2026-06-30T12:00:00.0000000+00:00", ExportBlobNaming.formatTimestamp(TS)); assertEquals("2026-06-30T12:00:00.1230000+00:00", ExportBlobNaming.formatTimestamp(Instant.parse("2026-06-30T12:00:00.123Z"))); + assertEquals("2026-06-30T12:00:00.1234567+00:00", + ExportBlobNaming.formatTimestamp(Instant.parse("2026-06-30T12:00:00.123456789Z"))); + } + + @Test + void blobFileName_matchesFullPrecisionReferenceHash() { + assertEquals("8d8ce6e13a2dbef356275361521d0c3da44b809a81474169cd41732e82ecd2e4.jsonl.gz", + ExportBlobNaming.blobFileName( + Instant.parse("2026-09-15T12:34:56.1234567Z"), "instance-1", JSONL)); } @Test diff --git a/exporthistory/src/test/java/com/microsoft/durabletask/exporthistory/ExportInstanceHistoryActivityTest.java b/exporthistory/src/test/java/com/microsoft/durabletask/exporthistory/ExportInstanceHistoryActivityTest.java new file mode 100644 index 00000000..1f7fe599 --- /dev/null +++ b/exporthistory/src/test/java/com/microsoft/durabletask/exporthistory/ExportInstanceHistoryActivityTest.java @@ -0,0 +1,49 @@ +// Copyright (c) Microsoft Corporation. All rights reserved. +// Licensed under the MIT License. +package com.microsoft.durabletask.exporthistory; + +import com.microsoft.durabletask.DurableTaskClient; +import com.microsoft.durabletask.OrchestrationMetadata; +import com.microsoft.durabletask.TaskActivityContext; +import com.microsoft.durabletask.history.GenericEvent; +import org.junit.jupiter.api.Test; + +import java.time.Instant; +import java.util.Collections; + +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Mockito.*; + +class ExportInstanceHistoryActivityTest { + + @Test + void exportPreservesPrecisionInBodyAndBlobName() { + Instant timestamp = Instant.parse("2026-09-15T12:34:56.1234567Z"); + DurableTaskClient client = mock(DurableTaskClient.class); + BlobExportWriter writer = mock(BlobExportWriter.class); + OrchestrationMetadata metadata = mock(OrchestrationMetadata.class); + when(metadata.isInstanceFound()).thenReturn(true); + when(metadata.isCompleted()).thenReturn(true); + when(metadata.getLastUpdatedAt()).thenReturn(timestamp); + when(client.getInstanceMetadata("instance-1", false)).thenReturn(metadata); + when(client.getOrchestrationHistory("instance-1")).thenReturn( + Collections.singletonList(new GenericEvent(1, timestamp, "payload"))); + ExportDestination destination = new ExportDestination("history"); + destination.setPrefix("exports/"); + ExportFormat format = new ExportFormat(ExportFormatKind.JSONL, "1.0"); + TaskActivityContext context = mock(TaskActivityContext.class); + when(context.getInput(ExportRequest.class)).thenReturn( + new ExportRequest("instance-1", destination, format)); + + ExportResult result = (ExportResult) new ExportInstanceHistoryActivity(client, writer).run(context); + + assertTrue(result.isSuccess()); + verify(writer).upload( + "history", + "exports/8d8ce6e13a2dbef356275361521d0c3da44b809a81474169cd41732e82ecd2e4.jsonl.gz", + "{\"eventType\":\"GenericEvent\",\"data\":\"payload\",\"eventId\":1," + + "\"isPlayed\":false,\"timestamp\":\"2026-09-15T12:34:56.1234567Z\"}\n", + format, + "instance-1"); + } +} diff --git a/exporthistory/src/test/java/com/microsoft/durabletask/exporthistory/HistoryEventSerializerParityTest.java b/exporthistory/src/test/java/com/microsoft/durabletask/exporthistory/HistoryEventSerializerParityTest.java index 65dbf18a..778eba9d 100644 --- a/exporthistory/src/test/java/com/microsoft/durabletask/exporthistory/HistoryEventSerializerParityTest.java +++ b/exporthistory/src/test/java/com/microsoft/durabletask/exporthistory/HistoryEventSerializerParityTest.java @@ -2,10 +2,21 @@ // Licensed under the MIT License. package com.microsoft.durabletask.exporthistory; +import com.google.protobuf.StringValue; +import com.google.protobuf.Value; +import com.microsoft.durabletask.DataConverter; +import com.microsoft.durabletask.DurableTaskClient; +import com.microsoft.durabletask.DurableTaskGrpcClientBuilder; import com.microsoft.durabletask.FailureDetails; import com.microsoft.durabletask.OrchestrationRuntimeStatus; import com.microsoft.durabletask.history.ContinueAsNewEvent; import com.microsoft.durabletask.history.EntityLockGrantedEvent; +import com.microsoft.durabletask.history.EntityLockRequestedEvent; +import com.microsoft.durabletask.history.EntityOperationCalledEvent; +import com.microsoft.durabletask.history.EntityOperationCompletedEvent; +import com.microsoft.durabletask.history.EntityOperationFailedEvent; +import com.microsoft.durabletask.history.EntityOperationSignaledEvent; +import com.microsoft.durabletask.history.EntityUnlockSentEvent; import com.microsoft.durabletask.history.EventRaisedEvent; import com.microsoft.durabletask.history.EventSentEvent; import com.microsoft.durabletask.history.ExecutionCompletedEvent; @@ -30,6 +41,13 @@ import com.microsoft.durabletask.history.TaskScheduledEvent; import com.microsoft.durabletask.history.TimerCreatedEvent; import com.microsoft.durabletask.history.TimerFiredEvent; +import com.microsoft.durabletask.implementation.protobuf.OrchestratorService; +import com.microsoft.durabletask.implementation.protobuf.TaskHubSidecarServiceGrpc; +import io.grpc.ManagedChannel; +import io.grpc.Server; +import io.grpc.inprocess.InProcessChannelBuilder; +import io.grpc.inprocess.InProcessServerBuilder; +import io.grpc.stub.StreamObserver; import org.junit.jupiter.api.Test; import java.io.BufferedReader; @@ -43,6 +61,7 @@ import java.util.Collections; import java.util.List; import java.util.Map; +import java.util.concurrent.TimeUnit; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -61,7 +80,7 @@ class HistoryEventSerializerParityTest { @Test void serializesEachEventByteForByteAgainstReference() throws Exception { - List golden = readGolden(); + List golden = readGolden("/golden/reference-history-events.jsonl"); List events = buildEvents(); assertEquals(golden.size(), events.size(), "golden line count vs event count"); @@ -75,7 +94,7 @@ void serializesEachEventByteForByteAgainstReference() throws Exception { @Test void serializesJsonArrayFormat() throws Exception { - List golden = readGolden(); + List golden = readGolden("/golden/reference-history-events.jsonl"); List two = Arrays.asList( new OrchestratorStartedEvent(0, TS), new GenericEvent(15, TS, "some-data")); @@ -93,13 +112,142 @@ void escapesStringsLikeReferenceEncoder() throws Exception { } @Test - void entityEventGetsJavaNativeEventTypeDiscriminator() throws Exception { + void entityInputControlCharactersMatchReferenceEncoding() { + StringBuilder controls = new StringBuilder(); + for (char ch = 0; ch < 0x20; ch++) { + controls.append(ch); + } + HistoryEvent event = new EntityOperationSignaledEvent( + 1, TS, "req-control", "echo", null, controls.toString(), "@counter@one"); + + // Newtonsoft inner-message JSON, wrapped with the reference export's System.Text.Json encoder. + String expected = "{\"eventType\":\"EventSent\",\"instanceId\":\"@counter@one\",\"name\":\"op\"," + + "\"input\":\"{\\u0022op\\u0022:\\u0022echo\\u0022,\\u0022signal\\u0022:true," + + "\\u0022input\\u0022:\\u0022" + + "\\\\u0000\\\\u0001\\\\u0002\\\\u0003\\\\u0004\\\\u0005\\\\u0006\\\\u0007" + + "\\\\b\\\\t\\\\n\\\\u000b\\\\f\\\\r\\\\u000e\\\\u000f" + + "\\\\u0010\\\\u0011\\\\u0012\\\\u0013\\\\u0014\\\\u0015\\\\u0016\\\\u0017" + + "\\\\u0018\\\\u0019\\\\u001a\\\\u001b\\\\u001c\\\\u001d\\\\u001e\\\\u001f" + + "\\u0022,\\u0022id\\u0022:\\u0022req-control\\u0022}\"," + + "\"eventId\":1,\"isPlayed\":false,\"timestamp\":\"2026-06-30T12:00:00Z\"}"; + + assertEquals(expected + "\n", HistoryEventSerializer.serialize(Collections.singletonList(event), JSONL)); + assertEquals("[" + expected + "]", HistoryEventSerializer.serialize(Collections.singletonList(event), JSON)); + } + + @Test + void entityLockGrantUsesReferenceEventRaisedRepresentation() throws Exception { HistoryEvent event = new EntityLockGrantedEvent(3, TS, "cs-1"); String actual = HistoryEventSerializer.serialize(Collections.singletonList(event), JSONL).trim(); - assertTrue(actual.startsWith("{\"eventType\":\"EntityLockGranted\""), actual); + assertTrue(actual.startsWith("{\"eventType\":\"EventRaised\""), actual); + assertTrue(actual.contains("\"name\":\"cs-1\""), actual); + assertTrue(actual.contains("\"isPlayed\":false"), actual); assertTrue(actual.contains("\"eventId\":3"), actual); } + @Test + void entityMessagesMatchReferenceConversionByteForByte() throws Exception { + // Captured using .NET EntityConversionState and DT Core 3.9.0, including Newtonsoft inner-message JSON. + List golden = readGolden("/golden/reference-entity-history-events.jsonl"); + List events = buildEntityEvents(); + String[] actual = HistoryEventSerializer.serialize(events, JSONL).split("\n"); + assertEquals(events.size(), actual.length); + int expectedIndex = 0; + for (int i = 0; i < events.size(); i++) { + if (events.get(i) instanceof ExecutionStartedEvent) { + continue; + } + assertEquals(golden.get(expectedIndex++), actual[i], + "entity wire-format mismatch at event " + events.get(i).getEventId()); + } + assertEquals(golden.size(), expectedIndex); + assertEquals("[" + String.join(",", actual) + "]", HistoryEventSerializer.serialize(events, JSON)); + } + + @Test + void entityFailureFromStreamedHistoryMatchesReference() throws Exception { + OrchestratorService.TaskFailureDetails failure = OrchestratorService.TaskFailureDetails.newBuilder() + .setErrorType("Outer") + .setErrorMessage("boom") + .setStackTrace(StringValue.of("at Foo()")) + .setInnerFailure(OrchestratorService.TaskFailureDetails.newBuilder() + .setErrorType("Inner") + .setErrorMessage("inner") + .setIsNonRetriable(true)) + .putProperties("code", Value.newBuilder().setNumberValue(42).build()) + .build(); + OrchestratorService.HistoryEvent event = OrchestratorService.HistoryEvent.newBuilder() + .setEventId(6) + .setTimestamp(DataConverter.getTimestampFromInstant( + Instant.parse("2026-06-30T12:00:00.1234567Z"))) + .setEntityOperationFailed(OrchestratorService.EntityOperationFailedEvent.newBuilder() + .setRequestId("req-failed") + .setFailureDetails(failure)) + .build(); + String serverName = InProcessServerBuilder.generateName(); + Server server = InProcessServerBuilder.forName(serverName) + .directExecutor() + .addService(new TaskHubSidecarServiceGrpc.TaskHubSidecarServiceImplBase() { + @Override + public void streamInstanceHistory( + OrchestratorService.StreamInstanceHistoryRequest request, + StreamObserver responseObserver) { + responseObserver.onNext(OrchestratorService.HistoryChunk.newBuilder() + .addEvents(event) + .build()); + responseObserver.onCompleted(); + } + }) + .build(); + ManagedChannel channel = InProcessChannelBuilder.forName(serverName).directExecutor().build(); + try { + server.start(); + try (DurableTaskClient client = new DurableTaskGrpcClientBuilder().grpcChannel(channel).build()) { + List history = client.getOrchestrationHistory("instance-1"); + String expected = readGolden("/golden/reference-entity-history-events.jsonl").get(5); + + assertEquals(1, history.size()); + assertEquals(expected + "\n", HistoryEventSerializer.serialize(history, JSONL)); + assertEquals("[" + expected + "]", HistoryEventSerializer.serialize(history, JSON)); + } + } finally { + channel.shutdownNow(); + server.shutdownNow(); + assertTrue(channel.awaitTermination(5, TimeUnit.SECONDS)); + assertTrue(server.awaitTermination(5, TimeUnit.SECONDS)); + } + } + + private static List buildEntityEvents() throws Exception { + Instant timestamp = Instant.parse("2026-06-30T12:00:00.1234567Z"); + Instant due = Instant.parse("2026-06-30T12:05:00.7654321Z"); + FailureDetails inner = failure("Inner", "inner", null, true, null); + FailureDetails outer = failure( + "Outer", "boom", "at Foo()", false, inner, Collections.singletonMap("code", 42.0)); + List lockSet = Arrays.asList("@Counter@", "@counter@a@b"); + return Arrays.asList( + new ExecutionStartedEvent(0, timestamp, "EntityWorkflow", null, null, + new OrchestrationInstance("order-42", "e1"), + null, null, null, null, Collections.emptyMap()), + new EntityOperationCalledEvent(1, timestamp, "req-call", "Add", null, + "caf\u00e9 &<>'+\u001f\u0085\u2028\u2029", + "ignored-parent", "ignored-execution", "@counter@one"), + new EntityOperationSignaledEvent(2, timestamp, "req-signal", "Increment", due, "1", "@counter@two"), + new EntityOperationSignaledEvent(3, timestamp, "req-empty", "Reset", null, null, "@counter@two"), + new EntityOperationCompletedEvent(4, timestamp, "req-null", null), + new EntityOperationCompletedEvent(5, timestamp, "req-complete", "{\"value\":42}"), + new EntityOperationFailedEvent(6, timestamp, "req-failed", outer), + new EntityLockRequestedEvent(7, timestamp, "lock-1", lockSet, 1, "ignored-parent"), + new EntityLockRequestedEvent(8, timestamp, "lock-0", lockSet, 0, "ignored-parent"), + new EntityLockGrantedEvent(9, timestamp, "lock-1"), + new EntityUnlockSentEvent(10, timestamp, "lock-1", "ignored-parent", "@counter@one"), + new ExecutionStartedEvent(11, timestamp, "EntityWorkflow", null, null, + new OrchestrationInstance("order-43", null), + null, null, null, null, Collections.emptyMap()), + new EntityOperationCalledEvent(12, timestamp, "req-next", "Get", + Instant.parse("2026-06-30T12:05:00Z"), null, null, null, "@counter@one")); + } + private static List buildEvents() throws Exception { FailureDetails inner = failure("System.NullReferenceException", "npe", " at Bar()", true, null); FailureDetails outer = failure("System.InvalidOperationException", "boom", " at Foo()", false, inner); @@ -137,10 +285,9 @@ private static List buildEvents() throws Exception { Collections.emptyMap()))); } - private static List readGolden() throws Exception { + private static List readGolden(String resource) throws Exception { List lines = new ArrayList<>(); - try (InputStream in = HistoryEventSerializerParityTest.class - .getResourceAsStream("/golden/reference-history-events.jsonl"); + try (InputStream in = HistoryEventSerializerParityTest.class.getResourceAsStream(resource); BufferedReader reader = new BufferedReader(new InputStreamReader(in, StandardCharsets.UTF_8))) { String line; while ((line = reader.readLine()) != null) { @@ -153,9 +300,15 @@ private static List readGolden() throws Exception { private static FailureDetails failure( String errorType, String message, String stackTrace, boolean nonRetriable, FailureDetails inner) throws Exception { + return failure(errorType, message, stackTrace, nonRetriable, inner, null); + } + + private static FailureDetails failure( + String errorType, String message, String stackTrace, boolean nonRetriable, + FailureDetails inner, Map properties) throws Exception { Constructor ctor = FailureDetails.class.getDeclaredConstructor( String.class, String.class, String.class, boolean.class, FailureDetails.class, Map.class); ctor.setAccessible(true); - return ctor.newInstance(errorType, message, stackTrace, nonRetriable, inner, null); + return ctor.newInstance(errorType, message, stackTrace, nonRetriable, inner, properties); } } diff --git a/exporthistory/src/test/java/com/microsoft/durabletask/exporthistory/HistoryEventSerializerTest.java b/exporthistory/src/test/java/com/microsoft/durabletask/exporthistory/HistoryEventSerializerTest.java index f28a6f02..8659664c 100644 --- a/exporthistory/src/test/java/com/microsoft/durabletask/exporthistory/HistoryEventSerializerTest.java +++ b/exporthistory/src/test/java/com/microsoft/durabletask/exporthistory/HistoryEventSerializerTest.java @@ -3,15 +3,15 @@ package com.microsoft.durabletask.exporthistory; import com.fasterxml.jackson.core.JsonProcessingException; -import com.microsoft.durabletask.history.EntityLockGrantedEvent; import com.microsoft.durabletask.history.EntityLockRequestedEvent; import com.microsoft.durabletask.history.EntityOperationCalledEvent; -import com.microsoft.durabletask.history.EntityOperationCompletedEvent; import com.microsoft.durabletask.history.EntityOperationFailedEvent; import com.microsoft.durabletask.history.EntityOperationSignaledEvent; import com.microsoft.durabletask.history.EntityUnlockSentEvent; +import com.microsoft.durabletask.history.ExecutionStartedEvent; import com.microsoft.durabletask.history.GenericEvent; import com.microsoft.durabletask.history.HistoryEvent; +import com.microsoft.durabletask.history.OrchestrationInstance; import com.microsoft.durabletask.history.TaskCompletedEvent; import org.junit.jupiter.api.Test; @@ -23,6 +23,7 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; /** @@ -105,44 +106,49 @@ void timestampSerialization_isLocaleInvariant() throws JsonProcessingException { } @Test - void entityEvents_serializeReflectivelyWithEventTypeDiscriminator() throws JsonProcessingException { - // The reflective writeEntity path (non-parity by design) covers all 7 entity event types; pin the - // eventType discriminator plus a representative field for each so a regression in the reflective - // projection (or the jsr310 Instant module) is caught. + void entityEvents_requireOrchestrationHistoryContext() { ExportFormat format = new ExportFormat(ExportFormatKind.JSONL, "1.0"); - - assertEntityEvent(format, + List events = Arrays.asList( new EntityOperationCalledEvent(1, TS, "req-1", "Add", null, "\"5\"", "@parent@p", "pe1", "@counter@c1"), - "EntityOperationCalled", "\"requestId\":\"req-1\"", "\"operation\":\"Add\""); - assertEntityEvent(format, - new EntityOperationSignaledEvent(2, TS, "req-2", "Increment", null, "\"1\"", "@counter@c2"), - "EntityOperationSignaled", "\"requestId\":\"req-2\"", "\"operation\":\"Increment\""); - assertEntityEvent(format, - new EntityOperationCompletedEvent(3, TS, "req-3", "\"result\""), - "EntityOperationCompleted", "\"requestId\":\"req-3\""); - assertEntityEvent(format, - new EntityOperationFailedEvent(4, TS, "req-4", null), - "EntityOperationFailed", "\"requestId\":\"req-4\""); - assertEntityEvent(format, new EntityLockRequestedEvent(5, TS, "cs-1", Arrays.asList("@e@a", "@e@b"), 0, "@parent@p"), - "EntityLockRequested", "\"criticalSectionId\":\"cs-1\""); - assertEntityEvent(format, - new EntityLockGrantedEvent(6, TS, "cs-2"), - "EntityLockGranted", "\"criticalSectionId\":\"cs-2\""); - assertEntityEvent(format, - new EntityUnlockSentEvent(7, TS, "cs-3", "@parent@p", "@e@t"), - "EntityUnlockSent", "\"criticalSectionId\":\"cs-3\""); + new EntityUnlockSentEvent(7, TS, "cs-3", "@parent@p", "@e@t")); + for (HistoryEvent event : events) { + assertThrows(IllegalArgumentException.class, + () -> HistoryEventSerializer.serialize(Collections.singletonList(event), format)); + } } - private static void assertEntityEvent( - ExportFormat format, HistoryEvent event, String eventType, String... expectedFragments) - throws JsonProcessingException { - String line = HistoryEventSerializer.serialize(Collections.singletonList(event), format).trim(); - assertTrue(line.startsWith("{\"eventType\":\"" + eventType + "\""), - eventType + " discriminator missing in: " + line); - for (String fragment : expectedFragments) { - assertTrue(line.contains(fragment), eventType + " missing " + fragment + " in: " + line); + @Test + void entityEvents_rejectMissingOrInvalidFields() { + ExportFormat format = new ExportFormat(ExportFormatKind.JSONL, "1.0"); + List invalidEvents = Arrays.asList( + new EntityOperationSignaledEvent(1, TS, "req-1", "Add", null, null, null), + new EntityOperationFailedEvent(2, TS, "req-2", null), + new EntityLockRequestedEvent(3, TS, "cs-1", Collections.singletonList("@counter@one"), 1, null), + new EntityLockRequestedEvent(4, TS, "cs-2", Collections.singletonList("@@invalid"), 0, null)); + for (HistoryEvent event : invalidEvents) { + assertThrows(IllegalArgumentException.class, + () -> HistoryEventSerializer.serialize(withExecutionStarted(event), format)); } } + + @Test + void entityHistoryContext_doesNotLeakAcrossExports() { + ExportFormat format = new ExportFormat(ExportFormatKind.JSONL, "1.0"); + HistoryEvent call = new EntityOperationCalledEvent( + 1, TS, "req-1", "Add", null, null, null, null, "@counter@one"); + HistoryEventSerializer.serialize(withExecutionStarted(call), format); + + assertThrows(IllegalArgumentException.class, + () -> HistoryEventSerializer.serialize(Collections.singletonList(call), format)); + } + + private static List withExecutionStarted(HistoryEvent event) { + return Arrays.asList( + new ExecutionStartedEvent(0, TS, "Orchestrator", null, null, + new OrchestrationInstance("parent", "execution"), + null, null, null, null, Collections.emptyMap()), + event); + } } diff --git a/exporthistory/src/test/resources/golden/reference-entity-history-events.jsonl b/exporthistory/src/test/resources/golden/reference-entity-history-events.jsonl new file mode 100644 index 00000000..78349680 --- /dev/null +++ b/exporthistory/src/test/resources/golden/reference-entity-history-events.jsonl @@ -0,0 +1,11 @@ +{"eventType":"EventSent","instanceId":"@counter@one","name":"op","input":"{\u0022op\u0022:\u0022Add\u0022,\u0022input\u0022:\u0022caf\u00E9 \u0026\u003C\u003E\u0027\u002B\\u001f\\u0085\\u2028\\u2029\u0022,\u0022id\u0022:\u0022req-call\u0022,\u0022parent\u0022:\u0022order-42\u0022,\u0022parentExecution\u0022:\u0022e1\u0022}","eventId":1,"isPlayed":false,"timestamp":"2026-06-30T12:00:00.1234567Z"} +{"eventType":"EventSent","instanceId":"@counter@two","name":"op@2026-06-30T12:05:00.7654321Z","input":"{\u0022op\u0022:\u0022Increment\u0022,\u0022signal\u0022:true,\u0022input\u0022:\u00221\u0022,\u0022id\u0022:\u0022req-signal\u0022,\u0022due\u0022:\u00222026-06-30T12:05:00.7654321Z\u0022}","eventId":2,"isPlayed":false,"timestamp":"2026-06-30T12:00:00.1234567Z"} +{"eventType":"EventSent","instanceId":"@counter@two","name":"op","input":"{\u0022op\u0022:\u0022Reset\u0022,\u0022signal\u0022:true,\u0022id\u0022:\u0022req-empty\u0022}","eventId":3,"isPlayed":false,"timestamp":"2026-06-30T12:00:00.1234567Z"} +{"eventType":"EventRaised","name":"req-null","input":"{\u0022result\u0022:null}","eventId":4,"isPlayed":false,"timestamp":"2026-06-30T12:00:00.1234567Z"} +{"eventType":"EventRaised","name":"req-complete","input":"{\u0022result\u0022:\u0022{\\\u0022value\\\u0022:42}\u0022}","eventId":5,"isPlayed":false,"timestamp":"2026-06-30T12:00:00.1234567Z"} +{"eventType":"EventRaised","name":"req-failed","input":"{\u0022result\u0022:\u0022boom\u0022,\u0022exceptionType\u0022:\u0022Outer\u0022,\u0022failureDetails\u0022:{\u0022ErrorType\u0022:\u0022Outer\u0022,\u0022ErrorMessage\u0022:\u0022boom\u0022,\u0022StackTrace\u0022:\u0022at Foo()\u0022,\u0022InnerFailure\u0022:{\u0022ErrorType\u0022:\u0022Inner\u0022,\u0022ErrorMessage\u0022:\u0022inner\u0022,\u0022StackTrace\u0022:null,\u0022InnerFailure\u0022:null,\u0022IsNonRetriable\u0022:true,\u0022Properties\u0022:{}},\u0022IsNonRetriable\u0022:false,\u0022Properties\u0022:{\u0022code\u0022:42.0}}}","eventId":6,"isPlayed":false,"timestamp":"2026-06-30T12:00:00.1234567Z"} +{"eventType":"EventSent","instanceId":"@counter@a@b","name":"op","input":"{\u0022op\u0022:null,\u0022id\u0022:\u0022lock-1\u0022,\u0022parent\u0022:\u0022order-42\u0022,\u0022lockset\u0022:[{\u0022name\u0022:\u0022Counter\u0022,\u0022key\u0022:\u0022\u0022},{\u0022name\u0022:\u0022counter\u0022,\u0022key\u0022:\u0022a@b\u0022}],\u0022pos\u0022:1}","eventId":7,"isPlayed":false,"timestamp":"2026-06-30T12:00:00.1234567Z"} +{"eventType":"EventSent","instanceId":"@Counter@","name":"op","input":"{\u0022op\u0022:null,\u0022id\u0022:\u0022lock-0\u0022,\u0022parent\u0022:\u0022order-42\u0022,\u0022lockset\u0022:[{\u0022name\u0022:\u0022Counter\u0022,\u0022key\u0022:\u0022\u0022},{\u0022name\u0022:\u0022counter\u0022,\u0022key\u0022:\u0022a@b\u0022}]}","eventId":8,"isPlayed":false,"timestamp":"2026-06-30T12:00:00.1234567Z"} +{"eventType":"EventRaised","name":"lock-1","input":"{\u0022result\u0022:\u0022Lock Acquisition Completed\u0022}","eventId":9,"isPlayed":false,"timestamp":"2026-06-30T12:00:00.1234567Z"} +{"eventType":"EventSent","instanceId":"@counter@one","name":"release","input":"{\u0022parent\u0022:\u0022order-42\u0022,\u0022id\u0022:\u0022lock-1\u0022}","eventId":10,"isPlayed":false,"timestamp":"2026-06-30T12:00:00.1234567Z"} +{"eventType":"EventSent","instanceId":"@counter@one","name":"op@2026-06-30T12:05:00.0000000Z","input":"{\u0022op\u0022:\u0022Get\u0022,\u0022id\u0022:\u0022req-next\u0022,\u0022parent\u0022:\u0022order-43\u0022,\u0022due\u0022:\u00222026-06-30T12:05:00Z\u0022}","eventId":12,"isPlayed":false,"timestamp":"2026-06-30T12:00:00.1234567Z"}