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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -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)).
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand All @@ -32,6 +32,7 @@
* <p>
* Stores payloads as blobs and returns opaque tokens in the form {@code blob:v1:<container>:<blobName>}.
* 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 {

Expand All @@ -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<CountDownLatch> 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<Object> containerGeneration = new AtomicReference<>();

/**
* Creates a new {@code BlobPayloadStore} from the given options.
Expand Down Expand Up @@ -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();
}
}

Expand Down
Loading
Loading