From eee31a3ff2c4bf78786cb386461c8cc618b3003b Mon Sep 17 00:00:00 2001 From: Chris Constable Date: Tue, 18 Aug 2026 16:24:22 -0400 Subject: [PATCH 1/2] refactor(extstore): general extstore refactoring. rename MessageTransformer to ExternalStorage, create a lazy extstore resolving data converter. --- .../PayloadAndFailureDataConverter.java | 6 + .../payload/storage/ExternalStorage.java | 162 +++++++++ .../ExternalStorageMessageTransformer.java | 75 ---- ...ExternalStorageNotConfiguredException.java | 16 + .../storage/ExternalStorageReferences.java | 11 +- .../payload/visitor/MessageVisitor.java | 2 +- .../visitor/PayloadVisitorOptions.java | 2 +- .../internal/worker/SingleWorkerOptions.java | 23 +- .../storage/ExternalStorageOptions.java | 30 +- .../ExternalStorageReferenceGuardTest.java | 51 +++ ...ExternalStorageMessageTransformerTest.java | 161 --------- .../payload/storage/ExternalStorageTest.java | 323 ++++++++++++++++++ .../storage/ExternalStorageOptionsTest.java | 18 + 13 files changed, 635 insertions(+), 245 deletions(-) create mode 100644 temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStorage.java delete mode 100644 temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStorageMessageTransformer.java create mode 100644 temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStorageNotConfiguredException.java create mode 100644 temporal-sdk/src/test/java/io/temporal/common/converter/ExternalStorageReferenceGuardTest.java delete mode 100644 temporal-sdk/src/test/java/io/temporal/internal/payload/storage/ExternalStorageMessageTransformerTest.java create mode 100644 temporal-sdk/src/test/java/io/temporal/internal/payload/storage/ExternalStorageTest.java diff --git a/temporal-sdk/src/main/java/io/temporal/common/converter/PayloadAndFailureDataConverter.java b/temporal-sdk/src/main/java/io/temporal/common/converter/PayloadAndFailureDataConverter.java index 935fd8462c..a628b4118b 100644 --- a/temporal-sdk/src/main/java/io/temporal/common/converter/PayloadAndFailureDataConverter.java +++ b/temporal-sdk/src/main/java/io/temporal/common/converter/PayloadAndFailureDataConverter.java @@ -8,6 +8,8 @@ import io.temporal.api.common.v1.Payloads; import io.temporal.api.failure.v1.Failure; import io.temporal.failure.DefaultFailureConverter; +import io.temporal.internal.payload.storage.ExternalStorageNotConfiguredException; +import io.temporal.internal.payload.storage.ExternalStorageReferences; import io.temporal.payload.context.SerializationContext; import java.lang.reflect.Type; import java.util.*; @@ -71,6 +73,10 @@ public T fromPayload(Payload payload, Class valueClass, Type valueType) return (T) new RawValue(payload); } + if (ExternalStorageReferences.isReference(payload)) { + throw new ExternalStorageNotConfiguredException(); + } + try { String encoding = payload.getMetadataOrThrow(EncodingKeys.METADATA_ENCODING_KEY).toString(UTF_8); diff --git a/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStorage.java b/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStorage.java new file mode 100644 index 0000000000..ec225fc837 --- /dev/null +++ b/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStorage.java @@ -0,0 +1,162 @@ +package io.temporal.internal.payload.storage; + +import com.google.common.base.Throwables; +import com.google.protobuf.Message; +import io.temporal.api.common.v1.Payload; +import io.temporal.api.sdk.v1.ExternalStorageReference; +import io.temporal.common.CancellationToken; +import io.temporal.internal.payload.visitor.MessageVisitor; +import io.temporal.internal.payload.visitor.PayloadVisitorOptions; +import io.temporal.internal.payload.visitor.PayloadVisitors; +import io.temporal.payload.storage.ExternalStorageOptions; +import io.temporal.payload.storage.StorageDriver; +import io.temporal.payload.storage.StorageDriverTargetInfo; +import java.util.List; +import java.util.concurrent.CancellationException; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionException; +import java.util.concurrent.ExecutionException; +import javax.annotation.Nullable; + +/** + * External storage offloads large payloads via {@link StorageDriver}s. It walks messages using + * {@link PayloadVisitors} transforming payloads to and from {@link ExternalStorageReference} using + * {@link ExternalStoragePayloadTransformer}. Use {@link ExternalStorageOptions} via {@link#create} + * to configure external storage. + */ +public final class ExternalStorage { + private final ExternalStoragePayloadTransformer payloadTransformer; + private final int payloadVisitConcurrency; + + public static ExternalStorage create(ExternalStorageOptions options) { + return new ExternalStorage( + ExternalStoragePayloadTransformer.fromOptions(options), + options.getMaxConcurrentPayloadVisits()); + } + + ExternalStorage( + ExternalStoragePayloadTransformer payloadTransformer, int payloadVisitConcurrency) { + this.payloadTransformer = payloadTransformer; + this.payloadVisitConcurrency = payloadVisitConcurrency; + } + + public T storeBlocking(T message, @Nullable StorageDriverTargetInfo target) { + return storeBlocking(message, target, CancellationToken.none()); + } + + public T storeBlocking( + T message, + @Nullable StorageDriverTargetInfo target, + CancellationToken cancellationToken) { + return getOrThrowIfCancelled(store(message, target, cancellationToken), cancellationToken); + } + + public T storeBlocking( + T message, + @Nullable StorageDriverTargetInfo target, + @Nullable MessageVisitor targetVisitor) { + CancellationToken cancellationToken = CancellationToken.none(); + return getOrThrowIfCancelled( + PayloadVisitors.visit(message, storeOptions(target, targetVisitor, cancellationToken)), + cancellationToken); + } + + public T retrieveBlocking(T message) { + CancellationToken cancellationToken = CancellationToken.none(); + return getOrThrowIfCancelled(retrieve(message, cancellationToken), cancellationToken); + } + + public CompletableFuture retrieveAsync(T message) { + return retrieve(message, CancellationToken.none()); + } + + /** + * Throws {@link ExternalStorageNotConfiguredException} if {@code message} contains any reference + * payload. Used at inbound task boundaries when external storage is not configured. + */ + public static void throwIfContainsReference(Message message) { + PayloadVisitorOptions options = + PayloadVisitorOptions.newBuilder( + (context, payloads) -> { + for (Payload payload : payloads) { + if (ExternalStorageReferences.isReference(payload)) { + CompletableFuture> found = new CompletableFuture<>(); + found.completeExceptionally(new ExternalStorageNotConfiguredException()); + return found; + } + } + return CompletableFuture.completedFuture(payloads); + }) + .setSkipSearchAttributes(true) + .build(); + try { + PayloadVisitors.visit(message.toBuilder(), options).join(); + } catch (CompletionException e) { + Throwable cause = e.getCause() != null ? e.getCause() : e; + Throwables.throwIfUnchecked(cause); + throw e; + } + } + + private static T getOrThrowIfCancelled( + CompletableFuture future, CancellationToken cancellationToken) { + try { + CompletableFuture.anyOf(future, cancellationToken.getCancellationFuture()).get(); + return future.get(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new CancellationException("External storage operation interrupted"); + } catch (ExecutionException e) { + Throwable cause = e.getCause() != null ? e.getCause() : e; + Throwables.throwIfUnchecked(cause); + throw new CompletionException(cause); + } + } + + CompletableFuture store( + T message, + @Nullable StorageDriverTargetInfo target, + CancellationToken cancellationToken) { + return PayloadVisitors.visit(message, storeOptions(target, null, cancellationToken)); + } + + CompletableFuture store( + Message.Builder builder, + @Nullable StorageDriverTargetInfo target, + CancellationToken cancellationToken) { + return PayloadVisitors.visit(builder, storeOptions(target, null, cancellationToken)); + } + + CompletableFuture retrieve( + T message, CancellationToken cancellationToken) { + return PayloadVisitors.visit(message, retrieveOptions(cancellationToken)); + } + + CompletableFuture retrieve( + Message.Builder builder, CancellationToken cancellationToken) { + return PayloadVisitors.visit(builder, retrieveOptions(cancellationToken)); + } + + private PayloadVisitorOptions storeOptions( + @Nullable StorageDriverTargetInfo target, + @Nullable MessageVisitor targetVisitor, + CancellationToken cancellationToken) { + return PayloadVisitorOptions.newBuilder( + (visitedTarget, payloads) -> + payloadTransformer.store(payloads, visitedTarget, cancellationToken)) + .setInitialContext(target) + .setMessageVisitor(targetVisitor) + .setConcurrency(payloadVisitConcurrency) + .setSkipSearchAttributes(true) + .build(); + } + + private PayloadVisitorOptions retrieveOptions( + CancellationToken cancellationToken) { + return PayloadVisitorOptions.newBuilder( + (context, payloads) -> payloadTransformer.retrieve(payloads, cancellationToken)) + .setConcurrency(payloadVisitConcurrency) + .setSkipSearchAttributes(true) + .build(); + } +} diff --git a/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStorageMessageTransformer.java b/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStorageMessageTransformer.java deleted file mode 100644 index 7385f99009..0000000000 --- a/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStorageMessageTransformer.java +++ /dev/null @@ -1,75 +0,0 @@ -package io.temporal.internal.payload.storage; - -import com.google.protobuf.Message; -import io.temporal.common.CancellationToken; -import io.temporal.internal.payload.visitor.PayloadVisitorOptions; -import io.temporal.internal.payload.visitor.PayloadVisitors; -import io.temporal.payload.storage.StorageDriverTargetInfo; -import java.util.concurrent.CancellationException; -import java.util.concurrent.CompletableFuture; -import javax.annotation.Nullable; - -/** - * Transforms payload lists reachable from a proto message by delegating each visited list to {@link - * ExternalStoragePayloadTransformer}. - * - *

Search attributes stay inline because the server indexes and validates their payload values. - * - *

The {@link Message.Builder} overloads transform in place; the {@link Message} overloads copy - * through a builder and complete with the copy. - */ -final class ExternalStorageMessageTransformer { - private final ExternalStoragePayloadTransformer payloadTransformer; - private final int payloadVisitConcurrency; - - ExternalStorageMessageTransformer( - ExternalStoragePayloadTransformer payloadTransformer, int payloadVisitConcurrency) { - this.payloadTransformer = payloadTransformer; - this.payloadVisitConcurrency = payloadVisitConcurrency; - } - - CompletableFuture store( - T message, - @Nullable StorageDriverTargetInfo target, - CancellationToken cancellationToken) { - return PayloadVisitors.visit(message, storeOptions(target, cancellationToken)); - } - - CompletableFuture store( - Message.Builder builder, - @Nullable StorageDriverTargetInfo target, - CancellationToken cancellationToken) { - return PayloadVisitors.visit(builder, storeOptions(target, cancellationToken)); - } - - CompletableFuture retrieve( - T message, CancellationToken cancellationToken) { - return PayloadVisitors.visit(message, retrieveOptions(cancellationToken)); - } - - CompletableFuture retrieve( - Message.Builder builder, CancellationToken cancellationToken) { - return PayloadVisitors.visit(builder, retrieveOptions(cancellationToken)); - } - - private PayloadVisitorOptions storeOptions( - @Nullable StorageDriverTargetInfo target, - CancellationToken cancellationToken) { - return PayloadVisitorOptions.newBuilder( - (visitedTarget, payloads) -> - payloadTransformer.store(payloads, visitedTarget, cancellationToken)) - .setInitialContext(target) - .setConcurrency(payloadVisitConcurrency) - .setSkipSearchAttributes(true) - .build(); - } - - private PayloadVisitorOptions retrieveOptions( - CancellationToken cancellationToken) { - return PayloadVisitorOptions.newBuilder( - (context, payloads) -> payloadTransformer.retrieve(payloads, cancellationToken)) - .setConcurrency(payloadVisitConcurrency) - .setSkipSearchAttributes(true) - .build(); - } -} diff --git a/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStorageNotConfiguredException.java b/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStorageNotConfiguredException.java new file mode 100644 index 0000000000..441e41851a --- /dev/null +++ b/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStorageNotConfiguredException.java @@ -0,0 +1,16 @@ +package io.temporal.internal.payload.storage; + +import io.temporal.common.converter.DataConverterException; + +/** + * Signals that an external storage reference reached a data converter without storage configured. + * Logs a TMPRL1105 error. + */ +public final class ExternalStorageNotConfiguredException extends DataConverterException { + public ExternalStorageNotConfiguredException() { + super( + "[TMPRL1105] Encountered an external-storage reference payload but external storage is not " + + "configured. Configure WorkflowClientOptions.Builder.setExternalStorage(...) with a " + + "driver able to retrieve it."); + } +} diff --git a/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStorageReferences.java b/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStorageReferences.java index 3a68c6bb66..1a81e3b676 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStorageReferences.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStorageReferences.java @@ -9,7 +9,7 @@ import javax.annotation.Nonnull; import javax.annotation.Nullable; -final class ExternalStorageReferences { +public final class ExternalStorageReferences { private static final String ENCODING_PROTOBUF_JSON = "json/protobuf"; private static final String REFERENCE_MESSAGE_TYPE = ExternalStorageReference.getDescriptor().getFullName(); @@ -64,8 +64,7 @@ static Payload toReferencePayload( * producer that omits it still yields a readable reference. */ static @Nullable ParsedReference tryParseReference(@Nonnull Payload payload) { - if (!hasMetadata(payload, EncodingKeys.METADATA_ENCODING_KEY, ENCODING_PROTOBUF_JSON) - || !hasMetadata(payload, EncodingKeys.METADATA_MESSAGE_TYPE_KEY, REFERENCE_MESSAGE_TYPE)) { + if (!isReference(payload)) { return null; } ExternalStorageReference.Builder builder = ExternalStorageReference.newBuilder(); @@ -79,6 +78,12 @@ static Payload toReferencePayload( reference.getDriverName(), new StorageDriverClaim(reference.getClaimDataMap())); } + /** True if {@code payload} has an external storage reference encoding and message type. */ + public static boolean isReference(Payload payload) { + return hasMetadata(payload, EncodingKeys.METADATA_ENCODING_KEY, ENCODING_PROTOBUF_JSON) + && hasMetadata(payload, EncodingKeys.METADATA_MESSAGE_TYPE_KEY, REFERENCE_MESSAGE_TYPE); + } + private static boolean hasMetadata(Payload payload, String key, String expected) { ByteString value = payload.getMetadataMap().get(key); return value != null && expected.equals(value.toStringUtf8()); diff --git a/temporal-sdk/src/main/java/io/temporal/internal/payload/visitor/MessageVisitor.java b/temporal-sdk/src/main/java/io/temporal/internal/payload/visitor/MessageVisitor.java index 21268e41d7..4bb6083e3e 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/payload/visitor/MessageVisitor.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/payload/visitor/MessageVisitor.java @@ -11,7 +11,7 @@ * @param type of the contextual value */ @FunctionalInterface -interface MessageVisitor { +public interface MessageVisitor { /** * Handles a message being entered and returns the contextual value for it and its contents. * diff --git a/temporal-sdk/src/main/java/io/temporal/internal/payload/visitor/PayloadVisitorOptions.java b/temporal-sdk/src/main/java/io/temporal/internal/payload/visitor/PayloadVisitorOptions.java index 4eac39be46..e4d6c89e47 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/payload/visitor/PayloadVisitorOptions.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/payload/visitor/PayloadVisitorOptions.java @@ -69,7 +69,7 @@ private Builder(@Nonnull PayloadVisitor payloadVisitor) { this.payloadVisitor = Objects.requireNonNull(payloadVisitor, "payloadVisitor"); } - Builder setMessageVisitor(@Nullable MessageVisitor messageVisitor) { + public Builder setMessageVisitor(@Nullable MessageVisitor messageVisitor) { this.messageVisitor = messageVisitor; return this; } diff --git a/temporal-sdk/src/main/java/io/temporal/internal/worker/SingleWorkerOptions.java b/temporal-sdk/src/main/java/io/temporal/internal/worker/SingleWorkerOptions.java index 8e0288566e..9b1d055110 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/worker/SingleWorkerOptions.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/worker/SingleWorkerOptions.java @@ -7,10 +7,12 @@ import io.temporal.common.converter.DataConverter; import io.temporal.common.converter.GlobalDataConverter; import io.temporal.common.interceptors.WorkerInterceptor; +import io.temporal.internal.payload.storage.ExternalStorage; import io.temporal.worker.PreferredVersionProvider; import io.temporal.worker.WorkerDeploymentOptions; import java.time.Duration; import java.util.List; +import javax.annotation.Nullable; public final class SingleWorkerOptions { @@ -45,6 +47,7 @@ public static final class Builder { private boolean allowActivityHeartbeatDuringShutdown; private String workerControlTaskQueue; private PreferredVersionProvider preferredVersionProvider; + private @Nullable ExternalStorage externalStorage; private Builder() {} @@ -73,6 +76,7 @@ private Builder(SingleWorkerOptions options) { this.allowActivityHeartbeatDuringShutdown = options.getAllowActivityHeartbeatDuringShutdown(); this.workerControlTaskQueue = options.getWorkerControlTaskQueue(); this.preferredVersionProvider = options.getPreferredVersionProvider(); + this.externalStorage = options.getExternalStorage(); } public Builder setIdentity(String identity) { @@ -185,6 +189,11 @@ public Builder setPreferredVersionProvider(PreferredVersionProvider preferredVer return this; } + public Builder setExternalStorage(@Nullable ExternalStorage externalStorage) { + this.externalStorage = externalStorage; + return this; + } + public SingleWorkerOptions build() { PollerOptions pollerOptions = this.pollerOptions; if (pollerOptions == null) { @@ -227,7 +236,8 @@ public SingleWorkerOptions build() { this.workerInstanceKey, this.allowActivityHeartbeatDuringShutdown, this.workerControlTaskQueue, - this.preferredVersionProvider); + this.preferredVersionProvider, + this.externalStorage); } } @@ -252,6 +262,7 @@ public SingleWorkerOptions build() { private final boolean allowActivityHeartbeatDuringShutdown; private final String workerControlTaskQueue; private final PreferredVersionProvider preferredVersionProvider; + private final @Nullable ExternalStorage externalStorage; private SingleWorkerOptions( String identity, @@ -274,7 +285,8 @@ private SingleWorkerOptions( String workerInstanceKey, boolean allowActivityHeartbeatDuringShutdown, String workerControlTaskQueue, - PreferredVersionProvider preferredVersionProvider) { + PreferredVersionProvider preferredVersionProvider, + @Nullable ExternalStorage externalStorage) { this.identity = identity; this.binaryChecksum = binaryChecksum; this.buildId = buildId; @@ -296,6 +308,7 @@ private SingleWorkerOptions( this.allowActivityHeartbeatDuringShutdown = allowActivityHeartbeatDuringShutdown; this.workerControlTaskQueue = workerControlTaskQueue; this.preferredVersionProvider = preferredVersionProvider; + this.externalStorage = externalStorage; } public String getIdentity() { @@ -393,6 +406,12 @@ public PreferredVersionProvider getPreferredVersionProvider() { return preferredVersionProvider; } + /** The external-storage message transformer for this worker, or null when disabled. */ + @Nullable + public ExternalStorage getExternalStorage() { + return externalStorage; + } + public WorkerVersioningOptions getWorkerVersioningOptions() { return new WorkerVersioningOptions( this.getBuildId(), this.isUsingBuildIdForVersioning(), this.getDeploymentOptions()); diff --git a/temporal-sdk/src/main/java/io/temporal/payload/storage/ExternalStorageOptions.java b/temporal-sdk/src/main/java/io/temporal/payload/storage/ExternalStorageOptions.java index 1486fb76b0..3abc969291 100644 --- a/temporal-sdk/src/main/java/io/temporal/payload/storage/ExternalStorageOptions.java +++ b/temporal-sdk/src/main/java/io/temporal/payload/storage/ExternalStorageOptions.java @@ -15,6 +15,7 @@ @Experimental public final class ExternalStorageOptions { static final int DEFAULT_PAYLOAD_SIZE_THRESHOLD = 256 * 1024; + static final int DEFAULT_MAX_CONCURRENT_PAYLOAD_VISITS = 3; public static Builder newBuilder() { return new Builder(); @@ -23,14 +24,17 @@ public static Builder newBuilder() { private final @Nonnull List drivers; private final @Nonnull StorageDriverSelector driverSelector; private final int payloadSizeThreshold; + private final int maxConcurrentPayloadVisits; private ExternalStorageOptions( @Nonnull List drivers, @Nonnull StorageDriverSelector driverSelector, - int payloadSizeThreshold) { + int payloadSizeThreshold, + int maxConcurrentPayloadVisits) { this.drivers = Collections.unmodifiableList(new ArrayList<>(drivers)); this.driverSelector = driverSelector; this.payloadSizeThreshold = payloadSizeThreshold; + this.maxConcurrentPayloadVisits = maxConcurrentPayloadVisits; } @Nonnull @@ -51,10 +55,20 @@ public int getPayloadSizeThreshold() { return payloadSizeThreshold; } + /** + * Maximum number of payload lists visited concurrently while offloading or restoring the payloads + * of a single message. Defaults to 3. + */ + public int getMaxConcurrentPayloadVisits() { + return maxConcurrentPayloadVisits; + } + public static final class Builder { private List drivers = Collections.emptyList(); private StorageDriverSelector driverSelector; private int payloadSizeThreshold = ExternalStorageOptions.DEFAULT_PAYLOAD_SIZE_THRESHOLD; + private int maxConcurrentPayloadVisits = + ExternalStorageOptions.DEFAULT_MAX_CONCURRENT_PAYLOAD_VISITS; private Builder() {} @@ -84,10 +98,21 @@ public Builder setPayloadSizeThreshold(int payloadSizeThreshold) { return this; } + /** + * Maximum number of payload lists visited concurrently while offloading or restoring the + * payloads of a single message. Must be at least 1. Defaults to 3. + */ + public Builder setMaxConcurrentPayloadVisits(int maxConcurrentPayloadVisits) { + this.maxConcurrentPayloadVisits = maxConcurrentPayloadVisits; + return this; + } + public ExternalStorageOptions build() { Preconditions.checkState(!drivers.isEmpty(), "At least one driver must be provided"); Preconditions.checkState( payloadSizeThreshold >= 0, "payloadSizeThreshold must be greater than or equal to zero"); + Preconditions.checkState( + maxConcurrentPayloadVisits >= 1, "maxConcurrentPayloadVisits must be at least 1"); Set names = new HashSet<>(); for (StorageDriver driver : drivers) { String name = driver.getName(); @@ -102,7 +127,8 @@ public ExternalStorageOptions build() { StorageDriver driver = drivers.get(0); selector = (context, payload) -> driver; } - return new ExternalStorageOptions(drivers, selector, payloadSizeThreshold); + return new ExternalStorageOptions( + drivers, selector, payloadSizeThreshold, maxConcurrentPayloadVisits); } } } diff --git a/temporal-sdk/src/test/java/io/temporal/common/converter/ExternalStorageReferenceGuardTest.java b/temporal-sdk/src/test/java/io/temporal/common/converter/ExternalStorageReferenceGuardTest.java new file mode 100644 index 0000000000..31c9698692 --- /dev/null +++ b/temporal-sdk/src/test/java/io/temporal/common/converter/ExternalStorageReferenceGuardTest.java @@ -0,0 +1,51 @@ +package io.temporal.common.converter; + +import static org.junit.Assert.assertThrows; +import static org.junit.Assert.assertTrue; + +import com.google.protobuf.ByteString; +import io.temporal.api.common.v1.Payload; +import io.temporal.api.sdk.v1.ExternalStorageReference; +import io.temporal.internal.payload.storage.ExternalStorageNotConfiguredException; +import org.junit.Test; + +/** + * When external storage is not configured, an inbound reference payload reaching value + * deserialization must fail with the clear {@code [TMPRL1105]} error instead of an opaque decoding + * failure. + */ +public class ExternalStorageReferenceGuardTest { + + private final DataConverter dataConverter = DefaultDataConverter.newDefaultInstance(); + + @Test + public void referencePayloadWithoutConfiguredStorageThrows() { + Payload reference = + Payload.newBuilder() + .putMetadata( + EncodingKeys.METADATA_ENCODING_KEY, ByteString.copyFromUtf8("json/protobuf")) + .putMetadata( + EncodingKeys.METADATA_MESSAGE_TYPE_KEY, + ByteString.copyFromUtf8(ExternalStorageReference.getDescriptor().getFullName())) + .setData(ByteString.copyFromUtf8("{}")) + .build(); + + ExternalStorageNotConfiguredException e = + assertThrows( + ExternalStorageNotConfiguredException.class, + () -> dataConverter.fromPayload(reference, String.class, String.class)); + assertTrue(e.getMessage(), e.getMessage().contains("[TMPRL1105]")); + } + + @Test + public void rawValueBypassesTheGuard() { + Payload reference = + Payload.newBuilder() + .addExternalPayloads( + Payload.ExternalPayloadDetails.newBuilder().setSizeBytes(1024).build()) + .build(); + + RawValue raw = dataConverter.fromPayload(reference, RawValue.class, RawValue.class); + assertTrue(raw.getPayload().getExternalPayloadsCount() > 0); + } +} diff --git a/temporal-sdk/src/test/java/io/temporal/internal/payload/storage/ExternalStorageMessageTransformerTest.java b/temporal-sdk/src/test/java/io/temporal/internal/payload/storage/ExternalStorageMessageTransformerTest.java deleted file mode 100644 index f17bcff47a..0000000000 --- a/temporal-sdk/src/test/java/io/temporal/internal/payload/storage/ExternalStorageMessageTransformerTest.java +++ /dev/null @@ -1,161 +0,0 @@ -package io.temporal.internal.payload.storage; - -import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertNotNull; -import static org.junit.Assert.assertNull; -import static org.junit.Assert.assertTrue; - -import com.google.protobuf.ByteString; -import io.temporal.api.command.v1.Command; -import io.temporal.api.command.v1.ScheduleActivityTaskCommandAttributes; -import io.temporal.api.command.v1.StartChildWorkflowExecutionCommandAttributes; -import io.temporal.api.common.v1.Payload; -import io.temporal.api.common.v1.Payloads; -import io.temporal.api.common.v1.SearchAttributes; -import io.temporal.common.CancellationToken; -import io.temporal.payload.storage.ExternalStorageOptions; -import io.temporal.payload.storage.StorageDriver; -import io.temporal.payload.storage.StorageDriverClaim; -import io.temporal.payload.storage.StorageDriverRetrieveContext; -import io.temporal.payload.storage.StorageDriverStoreContext; -import java.util.ArrayList; -import java.util.Collections; -import java.util.HashMap; -import java.util.List; -import java.util.Map; -import java.util.concurrent.CompletableFuture; -import org.junit.Test; - -/** Tests external storage message conversion. */ -public class ExternalStorageMessageTransformerTest { - - @Test - public void storeAndRetrieveRoundTripsOverAMessage() throws Exception { - InMemoryDriver driver = new InMemoryDriver("d1"); - ExternalStorageMessageTransformer transformer = transformer(driver, 0); - Payloads message = - Payloads.newBuilder().addPayloads(payload("a")).addPayloads(payload("b")).build(); - - Payloads stored = transformer.store(message, null, CancellationToken.none()).get(); - - assertNotNull(ExternalStorageReferences.tryParseReference(stored.getPayloads(0))); - assertNotNull(ExternalStorageReferences.tryParseReference(stored.getPayloads(1))); - - Payloads retrieved = transformer.retrieve(stored, CancellationToken.none()).get(); - assertEquals(message, retrieved); - } - - @Test - public void walksNestedPayloads() throws Exception { - InMemoryDriver driver = new InMemoryDriver("d1"); - ExternalStorageMessageTransformer transformer = transformer(driver, 0); - Command command = - Command.newBuilder() - .setScheduleActivityTaskCommandAttributes( - ScheduleActivityTaskCommandAttributes.newBuilder() - .setInput(Payloads.newBuilder().addPayloads(payload("deep")))) - .build(); - - Command stored = transformer.store(command, null, CancellationToken.none()).get(); - - Payload nested = stored.getScheduleActivityTaskCommandAttributes().getInput().getPayloads(0); - assertNotNull(ExternalStorageReferences.tryParseReference(nested)); - assertEquals(command, transformer.retrieve(stored, CancellationToken.none()).get()); - } - - @Test - public void payloadBelowThresholdLeavesMessageUnchanged() throws Exception { - InMemoryDriver driver = new InMemoryDriver("d1"); - ExternalStorageMessageTransformer transformer = transformer(driver, 1024); - Payloads message = Payloads.newBuilder().addPayloads(payload("small")).build(); - - Payloads stored = transformer.store(message, null, CancellationToken.none()).get(); - - assertNull(ExternalStorageReferences.tryParseReference(stored.getPayloads(0))); - assertEquals(message, stored); - assertTrue(driver.storeBatchSizes.isEmpty()); - } - - @Test - public void searchAttributesAreNotOffloaded() throws Exception { - InMemoryDriver driver = new InMemoryDriver("d1"); - ExternalStorageMessageTransformer transformer = transformer(driver, 0); - Command command = - Command.newBuilder() - .setStartChildWorkflowExecutionCommandAttributes( - StartChildWorkflowExecutionCommandAttributes.newBuilder() - .setInput(Payloads.newBuilder().addPayloads(payload("input"))) - .setSearchAttributes( - SearchAttributes.newBuilder() - .putIndexedFields("k", payload("indexed-value")))) - .build(); - - Command stored = transformer.store(command, null, CancellationToken.none()).get(); - - StartChildWorkflowExecutionCommandAttributes attrs = - stored.getStartChildWorkflowExecutionCommandAttributes(); - assertNotNull(ExternalStorageReferences.tryParseReference(attrs.getInput().getPayloads(0))); - Payload indexed = attrs.getSearchAttributes().getIndexedFieldsOrThrow("k"); - assertNull(ExternalStorageReferences.tryParseReference(indexed)); - assertEquals(payload("indexed-value"), indexed); - } - - private static ExternalStorageMessageTransformer transformer( - StorageDriver driver, int threshold) { - ExternalStoragePayloadTransformer payloadTransformer = - ExternalStoragePayloadTransformer.fromOptions( - ExternalStorageOptions.newBuilder() - .setDriver(driver) - .setPayloadSizeThreshold(threshold) - .build()); - return new ExternalStorageMessageTransformer(payloadTransformer, 4); - } - - private static Payload payload(String data) { - return Payload.newBuilder().setData(ByteString.copyFromUtf8(data)).build(); - } - - private static final class InMemoryDriver implements StorageDriver { - private final String name; - private final Map objects = new HashMap<>(); - final List storeBatchSizes = new ArrayList<>(); - private int counter = 0; - - InMemoryDriver(String name) { - this.name = name; - } - - @Override - public String getName() { - return name; - } - - @Override - public String getType() { - return "test.inmemory"; - } - - @Override - public synchronized CompletableFuture> store( - StorageDriverStoreContext context, List payloads) { - storeBatchSizes.add(payloads.size()); - List claims = new ArrayList<>(); - for (Payload payload : payloads) { - String key = name + "-" + (counter++); - objects.put(key, payload); - claims.add(new StorageDriverClaim(Collections.singletonMap("key", key))); - } - return CompletableFuture.completedFuture(claims); - } - - @Override - public synchronized CompletableFuture> retrieve( - StorageDriverRetrieveContext context, List claims) { - List payloads = new ArrayList<>(); - for (StorageDriverClaim claim : claims) { - payloads.add(objects.get(claim.getClaimData().get("key"))); - } - return CompletableFuture.completedFuture(payloads); - } - } -} diff --git a/temporal-sdk/src/test/java/io/temporal/internal/payload/storage/ExternalStorageTest.java b/temporal-sdk/src/test/java/io/temporal/internal/payload/storage/ExternalStorageTest.java new file mode 100644 index 0000000000..44576a87f8 --- /dev/null +++ b/temporal-sdk/src/test/java/io/temporal/internal/payload/storage/ExternalStorageTest.java @@ -0,0 +1,323 @@ +package io.temporal.internal.payload.storage; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertThrows; +import static org.junit.Assert.assertTrue; + +import com.google.protobuf.ByteString; +import io.temporal.api.command.v1.Command; +import io.temporal.api.command.v1.CompleteWorkflowExecutionCommandAttributes; +import io.temporal.api.command.v1.ScheduleActivityTaskCommandAttributes; +import io.temporal.api.command.v1.ScheduleActivityTaskCommandAttributesOrBuilder; +import io.temporal.api.command.v1.StartChildWorkflowExecutionCommandAttributes; +import io.temporal.api.common.v1.ActivityType; +import io.temporal.api.common.v1.Payload; +import io.temporal.api.common.v1.Payloads; +import io.temporal.api.common.v1.SearchAttributes; +import io.temporal.api.workflowservice.v1.RespondWorkflowTaskCompletedRequest; +import io.temporal.common.CancellationToken; +import io.temporal.internal.concurrent.structured.CancelSource; +import io.temporal.internal.payload.visitor.MessageVisitor; +import io.temporal.payload.storage.ExternalStorageOptions; +import io.temporal.payload.storage.StorageDriver; +import io.temporal.payload.storage.StorageDriverActivityInfo; +import io.temporal.payload.storage.StorageDriverClaim; +import io.temporal.payload.storage.StorageDriverRetrieveContext; +import io.temporal.payload.storage.StorageDriverStoreContext; +import io.temporal.payload.storage.StorageDriverTargetInfo; +import io.temporal.payload.storage.StorageDriverWorkflowInfo; +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CancellationException; +import java.util.concurrent.CompletableFuture; +import org.junit.Test; + +/** Tests external storage message conversion. */ +public class ExternalStorageTest { + + @Test + public void storeAndRetrieveRoundTripsOverAMessage() throws Exception { + InMemoryDriver driver = new InMemoryDriver("d1"); + ExternalStorage transformer = transformer(driver, 0); + Payloads message = + Payloads.newBuilder().addPayloads(payload("a")).addPayloads(payload("b")).build(); + + Payloads stored = transformer.store(message, null, CancellationToken.none()).get(); + + assertNotNull(ExternalStorageReferences.tryParseReference(stored.getPayloads(0))); + assertNotNull(ExternalStorageReferences.tryParseReference(stored.getPayloads(1))); + + Payloads retrieved = transformer.retrieve(stored, CancellationToken.none()).get(); + assertEquals(message, retrieved); + } + + @Test + public void walksNestedPayloads() throws Exception { + InMemoryDriver driver = new InMemoryDriver("d1"); + ExternalStorage transformer = transformer(driver, 0); + Command command = + Command.newBuilder() + .setScheduleActivityTaskCommandAttributes( + ScheduleActivityTaskCommandAttributes.newBuilder() + .setInput(Payloads.newBuilder().addPayloads(payload("deep")))) + .build(); + + Command stored = transformer.store(command, null, CancellationToken.none()).get(); + + Payload nested = stored.getScheduleActivityTaskCommandAttributes().getInput().getPayloads(0); + assertNotNull(ExternalStorageReferences.tryParseReference(nested)); + assertEquals(command, transformer.retrieve(stored, CancellationToken.none()).get()); + } + + @Test + public void payloadBelowThresholdLeavesMessageUnchanged() throws Exception { + InMemoryDriver driver = new InMemoryDriver("d1"); + ExternalStorage transformer = transformer(driver, 1024); + Payloads message = Payloads.newBuilder().addPayloads(payload("small")).build(); + + Payloads stored = transformer.store(message, null, CancellationToken.none()).get(); + + assertNull(ExternalStorageReferences.tryParseReference(stored.getPayloads(0))); + assertEquals(message, stored); + assertTrue(driver.storeBatchSizes.isEmpty()); + } + + @Test + public void searchAttributesAreNotOffloaded() throws Exception { + InMemoryDriver driver = new InMemoryDriver("d1"); + ExternalStorage transformer = transformer(driver, 0); + Command command = + Command.newBuilder() + .setStartChildWorkflowExecutionCommandAttributes( + StartChildWorkflowExecutionCommandAttributes.newBuilder() + .setInput(Payloads.newBuilder().addPayloads(payload("input"))) + .setSearchAttributes( + SearchAttributes.newBuilder() + .putIndexedFields("k", payload("indexed-value")))) + .build(); + + Command stored = transformer.store(command, null, CancellationToken.none()).get(); + + StartChildWorkflowExecutionCommandAttributes attrs = + stored.getStartChildWorkflowExecutionCommandAttributes(); + assertNotNull(ExternalStorageReferences.tryParseReference(attrs.getInput().getPayloads(0))); + Payload indexed = attrs.getSearchAttributes().getIndexedFieldsOrThrow("k"); + assertNull(ExternalStorageReferences.tryParseReference(indexed)); + assertEquals(payload("indexed-value"), indexed); + } + + @Test + public void throwIfContainsReferenceThrowsOnReference() throws Exception { + InMemoryDriver driver = new InMemoryDriver("d1"); + ExternalStorage transformer = transformer(driver, 0); + Payloads stored = + transformer + .store( + Payloads.newBuilder().addPayloads(payload("a")).build(), + null, + CancellationToken.none()) + .get(); + + ExternalStorageNotConfiguredException e = + assertThrows( + ExternalStorageNotConfiguredException.class, + () -> ExternalStorage.throwIfContainsReference(stored)); + assertTrue(e.getMessage(), e.getMessage().contains("[TMPRL1105]")); + } + + @Test + public void throwIfContainsReferenceAllowsInlinePayloads() { + Payloads inline = Payloads.newBuilder().addPayloads(payload("a")).build(); + ExternalStorage.throwIfContainsReference(inline); + } + + @Test + public void storeAppliesPerCommandTargetFromMessageVisitor() { + TargetCapturingDriver driver = new TargetCapturingDriver("d1"); + ExternalStorage storage = transformer(driver, 0); + + RespondWorkflowTaskCompletedRequest request = + RespondWorkflowTaskCompletedRequest.newBuilder() + .addCommands( + Command.newBuilder() + .setScheduleActivityTaskCommandAttributes( + ScheduleActivityTaskCommandAttributes.newBuilder() + .setActivityId("act-1") + .setActivityType(ActivityType.newBuilder().setName("MyActivity")) + .setInput( + Payloads.newBuilder().addPayloads(payload("activity-input"))))) + .addCommands( + Command.newBuilder() + .setCompleteWorkflowExecutionCommandAttributes( + CompleteWorkflowExecutionCommandAttributes.newBuilder() + .setResult(Payloads.newBuilder().addPayloads(payload("wf-result"))))) + .build(); + + StorageDriverTargetInfo workflowTarget = + new StorageDriverWorkflowInfo("ns", "wf-1", "run-1", "MyWorkflow"); + MessageVisitor visitor = + (current, message) -> { + if (message instanceof ScheduleActivityTaskCommandAttributesOrBuilder) { + ScheduleActivityTaskCommandAttributesOrBuilder attrs = + (ScheduleActivityTaskCommandAttributesOrBuilder) message; + return new StorageDriverActivityInfo( + "ns", attrs.getActivityId(), null, attrs.getActivityType().getName()); + } + return current; + }; + + storage.storeBlocking(request, workflowTarget, visitor); + + assertEquals( + new StorageDriverActivityInfo("ns", "act-1", null, "MyActivity"), + driver.targetFor("activity-input")); + assertEquals(workflowTarget, driver.targetFor("wf-result")); + } + + @Test + public void callerCancellationAbortsStore() { + ExternalStorage storage = transformer(new HangingDriver("d1"), 0); + CancelSource caller = new CancelSource<>(CancellationException::new); + caller.cancel(); + Payloads message = Payloads.newBuilder().addPayloads(payload("big")).build(); + + assertThrows( + CancellationException.class, () -> storage.storeBlocking(message, null, caller.token())); + } + + private static ExternalStorage transformer(StorageDriver driver, int threshold) { + ExternalStoragePayloadTransformer payloadTransformer = + ExternalStoragePayloadTransformer.fromOptions( + ExternalStorageOptions.newBuilder() + .setDriver(driver) + .setPayloadSizeThreshold(threshold) + .build()); + return new ExternalStorage(payloadTransformer, 4); + } + + private static Payload payload(String data) { + return Payload.newBuilder().setData(ByteString.copyFromUtf8(data)).build(); + } + + private static final class InMemoryDriver implements StorageDriver { + private final String name; + private final Map objects = new HashMap<>(); + final List storeBatchSizes = new ArrayList<>(); + private int counter = 0; + + InMemoryDriver(String name) { + this.name = name; + } + + @Override + public String getName() { + return name; + } + + @Override + public String getType() { + return "test.inmemory"; + } + + @Override + public synchronized CompletableFuture> store( + StorageDriverStoreContext context, List payloads) { + storeBatchSizes.add(payloads.size()); + List claims = new ArrayList<>(); + for (Payload payload : payloads) { + String key = name + "-" + (counter++); + objects.put(key, payload); + claims.add(new StorageDriverClaim(Collections.singletonMap("key", key))); + } + return CompletableFuture.completedFuture(claims); + } + + @Override + public synchronized CompletableFuture> retrieve( + StorageDriverRetrieveContext context, List claims) { + List payloads = new ArrayList<>(); + for (StorageDriverClaim claim : claims) { + payloads.add(objects.get(claim.getClaimData().get("key"))); + } + return CompletableFuture.completedFuture(payloads); + } + } + + private static final class TargetCapturingDriver implements StorageDriver { + private final String name; + private final Map targetByData = new HashMap<>(); + private int counter = 0; + + TargetCapturingDriver(String name) { + this.name = name; + } + + @Override + public String getName() { + return name; + } + + @Override + public String getType() { + return "test.capture"; + } + + @Override + public synchronized CompletableFuture> store( + StorageDriverStoreContext context, List payloads) { + List claims = new ArrayList<>(); + for (Payload payload : payloads) { + targetByData.put(payload.getData().toStringUtf8(), context.getTarget()); + claims.add( + new StorageDriverClaim(Collections.singletonMap("key", name + "-" + (counter++)))); + } + return CompletableFuture.completedFuture(claims); + } + + synchronized StorageDriverTargetInfo targetFor(String data) { + return targetByData.get(data); + } + + @Override + public CompletableFuture> retrieve( + StorageDriverRetrieveContext context, List claims) { + throw new UnsupportedOperationException(); + } + } + + /** Driver whose operations never settle, so only cancellation can end a blocking call. */ + private static final class HangingDriver implements StorageDriver { + private final String name; + + HangingDriver(String name) { + this.name = name; + } + + @Override + public String getName() { + return name; + } + + @Override + public String getType() { + return "test.hanging"; + } + + @Override + public CompletableFuture> store( + StorageDriverStoreContext context, List payloads) { + return new CompletableFuture<>(); + } + + @Override + public CompletableFuture> retrieve( + StorageDriverRetrieveContext context, List claims) { + return new CompletableFuture<>(); + } + } +} diff --git a/temporal-sdk/src/test/java/io/temporal/payload/storage/ExternalStorageOptionsTest.java b/temporal-sdk/src/test/java/io/temporal/payload/storage/ExternalStorageOptionsTest.java index 2c7ffc782f..b932986ad6 100644 --- a/temporal-sdk/src/test/java/io/temporal/payload/storage/ExternalStorageOptionsTest.java +++ b/temporal-sdk/src/test/java/io/temporal/payload/storage/ExternalStorageOptionsTest.java @@ -118,4 +118,22 @@ public void negativeThresholdRejected() { .setPayloadSizeThreshold(-1) .build(); } + + @Test + public void maxConcurrentPayloadVisitsDefaultsToThree() { + assertEquals( + 3, + ExternalStorageOptions.newBuilder() + .setDriver(driver("a")) + .build() + .getMaxConcurrentPayloadVisits()); + } + + @Test(expected = IllegalStateException.class) + public void zeroMaxConcurrentPayloadVisitsRejected() { + ExternalStorageOptions.newBuilder() + .setDriver(driver("a")) + .setMaxConcurrentPayloadVisits(0) + .build(); + } } From b3804da7dad7331916b8d2956ca4a5dece49b627 Mon Sep 17 00:00:00 2001 From: Chris Constable Date: Fri, 21 Aug 2026 16:34:50 -0400 Subject: [PATCH 2/2] refactor(extstore): move external storage to data converter and rename to match other sdks. --- .../common/converter/CodecDataConverter.java | 7 ++++ .../common/converter/DataConverter.java | 12 +++++++ .../converter/DefaultDataConverter.java | 11 +++++++ .../PayloadAndFailureDataConverter.java | 13 +++++++- .../ExternalStoragePayloadTransformer.java | 4 +-- ...torage.java => ExternalStorageRunner.java} | 14 ++++---- .../internal/worker/SingleWorkerOptions.java | 12 +++---- ...orageOptions.java => ExternalStorage.java} | 13 ++++---- .../payload/storage/StorageDriver.java | 4 +-- .../storage/StorageDriverSelector.java | 2 +- ...ExternalStoragePayloadTransformerTest.java | 17 ++++------ ...st.java => ExternalStorageRunnerTest.java} | 28 ++++++++-------- ...ionsTest.java => ExternalStorageTest.java} | 33 ++++++++----------- 13 files changed, 101 insertions(+), 69 deletions(-) rename temporal-sdk/src/main/java/io/temporal/internal/payload/storage/{ExternalStorage.java => ExternalStorageRunner.java} (95%) rename temporal-sdk/src/main/java/io/temporal/payload/storage/{ExternalStorageOptions.java => ExternalStorage.java} (92%) rename temporal-sdk/src/test/java/io/temporal/internal/payload/storage/{ExternalStorageTest.java => ExternalStorageRunnerTest.java} (93%) rename temporal-sdk/src/test/java/io/temporal/payload/storage/{ExternalStorageOptionsTest.java => ExternalStorageTest.java} (79%) diff --git a/temporal-sdk/src/main/java/io/temporal/common/converter/CodecDataConverter.java b/temporal-sdk/src/main/java/io/temporal/common/converter/CodecDataConverter.java index a82172348b..3fe83a3d0b 100644 --- a/temporal-sdk/src/main/java/io/temporal/common/converter/CodecDataConverter.java +++ b/temporal-sdk/src/main/java/io/temporal/common/converter/CodecDataConverter.java @@ -12,6 +12,7 @@ import io.temporal.payload.codec.ChainCodec; import io.temporal.payload.codec.PayloadCodec; import io.temporal.payload.context.SerializationContext; +import io.temporal.payload.storage.ExternalStorage; import java.lang.reflect.Type; import java.util.Collection; import java.util.Collections; @@ -190,6 +191,12 @@ public CodecDataConverter withContext(@Nonnull SerializationContext context) { return new CodecDataConverter(dataConverter, chainCodec, encodeFailureAttributes, context); } + @Override + @Nullable + public ExternalStorage getExternalStorage() { + return dataConverter.getExternalStorage(); + } + @Nonnull @Override public List encode(@Nonnull List payloads) { diff --git a/temporal-sdk/src/main/java/io/temporal/common/converter/DataConverter.java b/temporal-sdk/src/main/java/io/temporal/common/converter/DataConverter.java index decf6181b0..e299a3c2b7 100644 --- a/temporal-sdk/src/main/java/io/temporal/common/converter/DataConverter.java +++ b/temporal-sdk/src/main/java/io/temporal/common/converter/DataConverter.java @@ -11,10 +11,12 @@ import io.temporal.failure.DefaultFailureConverter; import io.temporal.payload.codec.PayloadCodec; import io.temporal.payload.context.SerializationContext; +import io.temporal.payload.storage.ExternalStorage; import java.lang.reflect.Type; import java.util.Arrays; import java.util.Optional; import javax.annotation.Nonnull; +import javax.annotation.Nullable; /** * Used by the framework to serialize/deserialize method parameters that need to be sent over the @@ -202,6 +204,16 @@ default DataConverter withContext(@Nonnull SerializationContext context) { return this; } + /** + * External storage offloads large payloads. This should not be used from inside workflow code as + * it performs nondeterministic operations. + */ + @Experimental + @Nullable + default ExternalStorage getExternalStorage() { + return null; + } + /** * @deprecated use {@link DataConverter#fromPayloads(int, Optional, Class, Type)}. This is an SDK * implementation detail and never was expected to be exposed to users. diff --git a/temporal-sdk/src/main/java/io/temporal/common/converter/DefaultDataConverter.java b/temporal-sdk/src/main/java/io/temporal/common/converter/DefaultDataConverter.java index f05c3f95a9..2df6393f76 100644 --- a/temporal-sdk/src/main/java/io/temporal/common/converter/DefaultDataConverter.java +++ b/temporal-sdk/src/main/java/io/temporal/common/converter/DefaultDataConverter.java @@ -1,8 +1,10 @@ package io.temporal.common.converter; import com.google.common.base.Preconditions; +import io.temporal.payload.storage.ExternalStorage; import java.util.*; import javax.annotation.Nonnull; +import javax.annotation.Nullable; /** * A {@link DataConverter} that delegates payload conversion to type specific {@link @@ -101,4 +103,13 @@ public DefaultDataConverter withFailureConverter(@Nonnull FailureConverter failu this.failureConverter = Preconditions.checkNotNull(failureConverter, "failureConverter"); return this; } + + /** + * Modifies this {@code DefaultDataConverter} by attaching external storage, used to offload and + * restore large payloads at task and RPC boundaries. + */ + public DefaultDataConverter withExternalStorage(@Nullable ExternalStorage externalStorage) { + this.externalStorage = externalStorage; + return this; + } } diff --git a/temporal-sdk/src/main/java/io/temporal/common/converter/PayloadAndFailureDataConverter.java b/temporal-sdk/src/main/java/io/temporal/common/converter/PayloadAndFailureDataConverter.java index a628b4118b..e848e5d9db 100644 --- a/temporal-sdk/src/main/java/io/temporal/common/converter/PayloadAndFailureDataConverter.java +++ b/temporal-sdk/src/main/java/io/temporal/common/converter/PayloadAndFailureDataConverter.java @@ -11,6 +11,7 @@ import io.temporal.internal.payload.storage.ExternalStorageNotConfiguredException; import io.temporal.internal.payload.storage.ExternalStorageReferences; import io.temporal.payload.context.SerializationContext; +import io.temporal.payload.storage.ExternalStorage; import java.lang.reflect.Type; import java.util.*; import javax.annotation.Nonnull; @@ -24,6 +25,7 @@ class PayloadAndFailureDataConverter implements DataConverter { volatile List converters; volatile Map convertersMap; volatile FailureConverter failureConverter; + volatile @Nullable ExternalStorage externalStorage; private final @Nullable SerializationContext serializationContext; public PayloadAndFailureDataConverter(@Nonnull List converters) { @@ -149,9 +151,18 @@ public Failure exceptionToFailure(@Nonnull Throwable throwable) { .exceptionToFailure(throwable, this); } + @Override + @Nullable + public ExternalStorage getExternalStorage() { + return externalStorage; + } + @Override public @Nonnull DataConverter withContext(@Nonnull SerializationContext context) { - return new PayloadAndFailureDataConverter(converters, convertersMap, failureConverter, context); + PayloadAndFailureDataConverter copy = + new PayloadAndFailureDataConverter(converters, convertersMap, failureConverter, context); + copy.externalStorage = this.externalStorage; + return copy; } static Map createConvertersMap(List converters) { diff --git a/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStoragePayloadTransformer.java b/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStoragePayloadTransformer.java index 6e0d4d770c..ef13075f9e 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStoragePayloadTransformer.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStoragePayloadTransformer.java @@ -4,7 +4,7 @@ import io.temporal.common.CancellationToken; import io.temporal.internal.common.ListUtils; import io.temporal.internal.concurrent.structured.TaskScope; -import io.temporal.payload.storage.ExternalStorageOptions; +import io.temporal.payload.storage.ExternalStorage; import io.temporal.payload.storage.StorageDriver; import io.temporal.payload.storage.StorageDriverClaim; import io.temporal.payload.storage.StorageDriverRetrieveContext; @@ -30,7 +30,7 @@ final class ExternalStoragePayloadTransformer { private final StorageDriverSelector selector; private final int payloadSizeThreshold; - static ExternalStoragePayloadTransformer fromOptions(ExternalStorageOptions options) { + static ExternalStoragePayloadTransformer fromOptions(ExternalStorage options) { Map driversByName = new LinkedHashMap<>(); for (StorageDriver driver : options.getDrivers()) { driversByName.put(driver.getName(), driver); diff --git a/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStorage.java b/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStorageRunner.java similarity index 95% rename from temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStorage.java rename to temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStorageRunner.java index ec225fc837..4eb516be81 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStorage.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/payload/storage/ExternalStorageRunner.java @@ -8,7 +8,7 @@ import io.temporal.internal.payload.visitor.MessageVisitor; import io.temporal.internal.payload.visitor.PayloadVisitorOptions; import io.temporal.internal.payload.visitor.PayloadVisitors; -import io.temporal.payload.storage.ExternalStorageOptions; +import io.temporal.payload.storage.ExternalStorage; import io.temporal.payload.storage.StorageDriver; import io.temporal.payload.storage.StorageDriverTargetInfo; import java.util.List; @@ -21,20 +21,20 @@ /** * External storage offloads large payloads via {@link StorageDriver}s. It walks messages using * {@link PayloadVisitors} transforming payloads to and from {@link ExternalStorageReference} using - * {@link ExternalStoragePayloadTransformer}. Use {@link ExternalStorageOptions} via {@link#create} - * to configure external storage. + * {@link ExternalStoragePayloadTransformer}. Use {@link ExternalStorage} via {@link#create} to + * configure external storage. */ -public final class ExternalStorage { +public final class ExternalStorageRunner { private final ExternalStoragePayloadTransformer payloadTransformer; private final int payloadVisitConcurrency; - public static ExternalStorage create(ExternalStorageOptions options) { - return new ExternalStorage( + public static ExternalStorageRunner create(ExternalStorage options) { + return new ExternalStorageRunner( ExternalStoragePayloadTransformer.fromOptions(options), options.getMaxConcurrentPayloadVisits()); } - ExternalStorage( + ExternalStorageRunner( ExternalStoragePayloadTransformer payloadTransformer, int payloadVisitConcurrency) { this.payloadTransformer = payloadTransformer; this.payloadVisitConcurrency = payloadVisitConcurrency; diff --git a/temporal-sdk/src/main/java/io/temporal/internal/worker/SingleWorkerOptions.java b/temporal-sdk/src/main/java/io/temporal/internal/worker/SingleWorkerOptions.java index 9b1d055110..23eb93625a 100644 --- a/temporal-sdk/src/main/java/io/temporal/internal/worker/SingleWorkerOptions.java +++ b/temporal-sdk/src/main/java/io/temporal/internal/worker/SingleWorkerOptions.java @@ -7,7 +7,7 @@ import io.temporal.common.converter.DataConverter; import io.temporal.common.converter.GlobalDataConverter; import io.temporal.common.interceptors.WorkerInterceptor; -import io.temporal.internal.payload.storage.ExternalStorage; +import io.temporal.internal.payload.storage.ExternalStorageRunner; import io.temporal.worker.PreferredVersionProvider; import io.temporal.worker.WorkerDeploymentOptions; import java.time.Duration; @@ -47,7 +47,7 @@ public static final class Builder { private boolean allowActivityHeartbeatDuringShutdown; private String workerControlTaskQueue; private PreferredVersionProvider preferredVersionProvider; - private @Nullable ExternalStorage externalStorage; + private @Nullable ExternalStorageRunner externalStorage; private Builder() {} @@ -189,7 +189,7 @@ public Builder setPreferredVersionProvider(PreferredVersionProvider preferredVer return this; } - public Builder setExternalStorage(@Nullable ExternalStorage externalStorage) { + public Builder setExternalStorage(@Nullable ExternalStorageRunner externalStorage) { this.externalStorage = externalStorage; return this; } @@ -262,7 +262,7 @@ public SingleWorkerOptions build() { private final boolean allowActivityHeartbeatDuringShutdown; private final String workerControlTaskQueue; private final PreferredVersionProvider preferredVersionProvider; - private final @Nullable ExternalStorage externalStorage; + private final @Nullable ExternalStorageRunner externalStorage; private SingleWorkerOptions( String identity, @@ -286,7 +286,7 @@ private SingleWorkerOptions( boolean allowActivityHeartbeatDuringShutdown, String workerControlTaskQueue, PreferredVersionProvider preferredVersionProvider, - @Nullable ExternalStorage externalStorage) { + @Nullable ExternalStorageRunner externalStorage) { this.identity = identity; this.binaryChecksum = binaryChecksum; this.buildId = buildId; @@ -408,7 +408,7 @@ public PreferredVersionProvider getPreferredVersionProvider() { /** The external-storage message transformer for this worker, or null when disabled. */ @Nullable - public ExternalStorage getExternalStorage() { + public ExternalStorageRunner getExternalStorage() { return externalStorage; } diff --git a/temporal-sdk/src/main/java/io/temporal/payload/storage/ExternalStorageOptions.java b/temporal-sdk/src/main/java/io/temporal/payload/storage/ExternalStorage.java similarity index 92% rename from temporal-sdk/src/main/java/io/temporal/payload/storage/ExternalStorageOptions.java rename to temporal-sdk/src/main/java/io/temporal/payload/storage/ExternalStorage.java index 3abc969291..854254ef04 100644 --- a/temporal-sdk/src/main/java/io/temporal/payload/storage/ExternalStorageOptions.java +++ b/temporal-sdk/src/main/java/io/temporal/payload/storage/ExternalStorage.java @@ -13,7 +13,7 @@ /** Configuration for offloading large payloads to external storage. */ @Experimental -public final class ExternalStorageOptions { +public final class ExternalStorage { static final int DEFAULT_PAYLOAD_SIZE_THRESHOLD = 256 * 1024; static final int DEFAULT_MAX_CONCURRENT_PAYLOAD_VISITS = 3; @@ -26,7 +26,7 @@ public static Builder newBuilder() { private final int payloadSizeThreshold; private final int maxConcurrentPayloadVisits; - private ExternalStorageOptions( + private ExternalStorage( @Nonnull List drivers, @Nonnull StorageDriverSelector driverSelector, int payloadSizeThreshold, @@ -66,9 +66,8 @@ public int getMaxConcurrentPayloadVisits() { public static final class Builder { private List drivers = Collections.emptyList(); private StorageDriverSelector driverSelector; - private int payloadSizeThreshold = ExternalStorageOptions.DEFAULT_PAYLOAD_SIZE_THRESHOLD; - private int maxConcurrentPayloadVisits = - ExternalStorageOptions.DEFAULT_MAX_CONCURRENT_PAYLOAD_VISITS; + private int payloadSizeThreshold = ExternalStorage.DEFAULT_PAYLOAD_SIZE_THRESHOLD; + private int maxConcurrentPayloadVisits = ExternalStorage.DEFAULT_MAX_CONCURRENT_PAYLOAD_VISITS; private Builder() {} @@ -107,7 +106,7 @@ public Builder setMaxConcurrentPayloadVisits(int maxConcurrentPayloadVisits) { return this; } - public ExternalStorageOptions build() { + public ExternalStorage build() { Preconditions.checkState(!drivers.isEmpty(), "At least one driver must be provided"); Preconditions.checkState( payloadSizeThreshold >= 0, "payloadSizeThreshold must be greater than or equal to zero"); @@ -127,7 +126,7 @@ public ExternalStorageOptions build() { StorageDriver driver = drivers.get(0); selector = (context, payload) -> driver; } - return new ExternalStorageOptions( + return new ExternalStorage( drivers, selector, payloadSizeThreshold, maxConcurrentPayloadVisits); } } diff --git a/temporal-sdk/src/main/java/io/temporal/payload/storage/StorageDriver.java b/temporal-sdk/src/main/java/io/temporal/payload/storage/StorageDriver.java index 01d851fbe6..bf332c38d7 100644 --- a/temporal-sdk/src/main/java/io/temporal/payload/storage/StorageDriver.java +++ b/temporal-sdk/src/main/java/io/temporal/payload/storage/StorageDriver.java @@ -11,8 +11,8 @@ public interface StorageDriver { /** * Name of this driver instance, unique among the drivers registered in a single {@link - * ExternalStorageOptions}. Used as the routing key recorded in a stored payload's reference and - * resolved back to this driver on retrieval. + * ExternalStorage}. Used as the routing key recorded in a stored payload's reference and resolved + * back to this driver on retrieval. */ @Nonnull String getName(); diff --git a/temporal-sdk/src/main/java/io/temporal/payload/storage/StorageDriverSelector.java b/temporal-sdk/src/main/java/io/temporal/payload/storage/StorageDriverSelector.java index 431622e2fa..966e52e68d 100644 --- a/temporal-sdk/src/main/java/io/temporal/payload/storage/StorageDriverSelector.java +++ b/temporal-sdk/src/main/java/io/temporal/payload/storage/StorageDriverSelector.java @@ -11,7 +11,7 @@ public interface StorageDriverSelector { /** * Returns the driver to store {@code payload}, which must be one of the drivers registered in the - * {@link ExternalStorageOptions}, or {@code null} to leave the payload stored inline. + * {@link ExternalStorage}, or {@code null} to leave the payload stored inline. */ @Nullable StorageDriver selectDriver(@Nonnull StorageDriverStoreContext context, @Nonnull Payload payload); diff --git a/temporal-sdk/src/test/java/io/temporal/internal/payload/storage/ExternalStoragePayloadTransformerTest.java b/temporal-sdk/src/test/java/io/temporal/internal/payload/storage/ExternalStoragePayloadTransformerTest.java index f1632ca81e..2dcb58f384 100644 --- a/temporal-sdk/src/test/java/io/temporal/internal/payload/storage/ExternalStoragePayloadTransformerTest.java +++ b/temporal-sdk/src/test/java/io/temporal/internal/payload/storage/ExternalStoragePayloadTransformerTest.java @@ -12,7 +12,7 @@ import io.temporal.api.common.v1.Payload; import io.temporal.common.CancellationToken; import io.temporal.internal.concurrent.structured.CancelSource; -import io.temporal.payload.storage.ExternalStorageOptions; +import io.temporal.payload.storage.ExternalStorage; import io.temporal.payload.storage.StorageDriver; import io.temporal.payload.storage.StorageDriverClaim; import io.temporal.payload.storage.StorageDriverRetrieveContext; @@ -75,7 +75,7 @@ public void selectorReturningNullKeepsInline() throws Exception { InMemoryDriver driver = new InMemoryDriver("d1"); ExternalStoragePayloadTransformer transformer = ExternalStoragePayloadTransformer.fromOptions( - ExternalStorageOptions.newBuilder() + ExternalStorage.newBuilder() .setDriver(driver) .setDriverSelector((context, payload) -> null) .setPayloadSizeThreshold(0) @@ -101,7 +101,7 @@ public void multipleDriversBatchPerDriverAndPreserveOrder() throws Exception { (context, payload) -> byPrefix.get(payload.getData().toStringUtf8().substring(0, 1)); ExternalStoragePayloadTransformer transformer = ExternalStoragePayloadTransformer.fromOptions( - ExternalStorageOptions.newBuilder() + ExternalStorage.newBuilder() .setDrivers(Arrays.asList(d1, d2)) .setDriverSelector(selector) .setPayloadSizeThreshold(0) @@ -198,7 +198,7 @@ public void selectorReturningUnregisteredDriverFails() { InMemoryDriver stranger = new InMemoryDriver("d2"); ExternalStoragePayloadTransformer transformer = ExternalStoragePayloadTransformer.fromOptions( - ExternalStorageOptions.newBuilder() + ExternalStorage.newBuilder() .setDriver(registered) .setDriverSelector((context, payload) -> stranger) .setPayloadSizeThreshold(0) @@ -232,7 +232,7 @@ public CompletableFuture> store( byPrefix.put("2", doomed); ExternalStoragePayloadTransformer transformer = ExternalStoragePayloadTransformer.fromOptions( - ExternalStorageOptions.newBuilder() + ExternalStorage.newBuilder() .setDrivers(Arrays.asList(slow, doomed)) .setDriverSelector( (context, payload) -> @@ -314,7 +314,7 @@ public void selectorObservesCallerCancellationToken() { AtomicReference> observed = new AtomicReference<>(); ExternalStoragePayloadTransformer transformer = ExternalStoragePayloadTransformer.fromOptions( - ExternalStorageOptions.newBuilder() + ExternalStorage.newBuilder() .setDriver(driver) .setDriverSelector( (context, payload) -> { @@ -332,10 +332,7 @@ public void selectorObservesCallerCancellationToken() { private static ExternalStoragePayloadTransformer transformer( StorageDriver driver, int threshold) { return ExternalStoragePayloadTransformer.fromOptions( - ExternalStorageOptions.newBuilder() - .setDriver(driver) - .setPayloadSizeThreshold(threshold) - .build()); + ExternalStorage.newBuilder().setDriver(driver).setPayloadSizeThreshold(threshold).build()); } private static Payload payload(String data) { diff --git a/temporal-sdk/src/test/java/io/temporal/internal/payload/storage/ExternalStorageTest.java b/temporal-sdk/src/test/java/io/temporal/internal/payload/storage/ExternalStorageRunnerTest.java similarity index 93% rename from temporal-sdk/src/test/java/io/temporal/internal/payload/storage/ExternalStorageTest.java rename to temporal-sdk/src/test/java/io/temporal/internal/payload/storage/ExternalStorageRunnerTest.java index 44576a87f8..be0f773d37 100644 --- a/temporal-sdk/src/test/java/io/temporal/internal/payload/storage/ExternalStorageTest.java +++ b/temporal-sdk/src/test/java/io/temporal/internal/payload/storage/ExternalStorageRunnerTest.java @@ -20,7 +20,7 @@ import io.temporal.common.CancellationToken; import io.temporal.internal.concurrent.structured.CancelSource; import io.temporal.internal.payload.visitor.MessageVisitor; -import io.temporal.payload.storage.ExternalStorageOptions; +import io.temporal.payload.storage.ExternalStorage; import io.temporal.payload.storage.StorageDriver; import io.temporal.payload.storage.StorageDriverActivityInfo; import io.temporal.payload.storage.StorageDriverClaim; @@ -38,12 +38,12 @@ import org.junit.Test; /** Tests external storage message conversion. */ -public class ExternalStorageTest { +public class ExternalStorageRunnerTest { @Test public void storeAndRetrieveRoundTripsOverAMessage() throws Exception { InMemoryDriver driver = new InMemoryDriver("d1"); - ExternalStorage transformer = transformer(driver, 0); + ExternalStorageRunner transformer = transformer(driver, 0); Payloads message = Payloads.newBuilder().addPayloads(payload("a")).addPayloads(payload("b")).build(); @@ -59,7 +59,7 @@ public void storeAndRetrieveRoundTripsOverAMessage() throws Exception { @Test public void walksNestedPayloads() throws Exception { InMemoryDriver driver = new InMemoryDriver("d1"); - ExternalStorage transformer = transformer(driver, 0); + ExternalStorageRunner transformer = transformer(driver, 0); Command command = Command.newBuilder() .setScheduleActivityTaskCommandAttributes( @@ -77,7 +77,7 @@ public void walksNestedPayloads() throws Exception { @Test public void payloadBelowThresholdLeavesMessageUnchanged() throws Exception { InMemoryDriver driver = new InMemoryDriver("d1"); - ExternalStorage transformer = transformer(driver, 1024); + ExternalStorageRunner transformer = transformer(driver, 1024); Payloads message = Payloads.newBuilder().addPayloads(payload("small")).build(); Payloads stored = transformer.store(message, null, CancellationToken.none()).get(); @@ -90,7 +90,7 @@ public void payloadBelowThresholdLeavesMessageUnchanged() throws Exception { @Test public void searchAttributesAreNotOffloaded() throws Exception { InMemoryDriver driver = new InMemoryDriver("d1"); - ExternalStorage transformer = transformer(driver, 0); + ExternalStorageRunner transformer = transformer(driver, 0); Command command = Command.newBuilder() .setStartChildWorkflowExecutionCommandAttributes( @@ -114,7 +114,7 @@ public void searchAttributesAreNotOffloaded() throws Exception { @Test public void throwIfContainsReferenceThrowsOnReference() throws Exception { InMemoryDriver driver = new InMemoryDriver("d1"); - ExternalStorage transformer = transformer(driver, 0); + ExternalStorageRunner transformer = transformer(driver, 0); Payloads stored = transformer .store( @@ -126,20 +126,20 @@ public void throwIfContainsReferenceThrowsOnReference() throws Exception { ExternalStorageNotConfiguredException e = assertThrows( ExternalStorageNotConfiguredException.class, - () -> ExternalStorage.throwIfContainsReference(stored)); + () -> ExternalStorageRunner.throwIfContainsReference(stored)); assertTrue(e.getMessage(), e.getMessage().contains("[TMPRL1105]")); } @Test public void throwIfContainsReferenceAllowsInlinePayloads() { Payloads inline = Payloads.newBuilder().addPayloads(payload("a")).build(); - ExternalStorage.throwIfContainsReference(inline); + ExternalStorageRunner.throwIfContainsReference(inline); } @Test public void storeAppliesPerCommandTargetFromMessageVisitor() { TargetCapturingDriver driver = new TargetCapturingDriver("d1"); - ExternalStorage storage = transformer(driver, 0); + ExternalStorageRunner storage = transformer(driver, 0); RespondWorkflowTaskCompletedRequest request = RespondWorkflowTaskCompletedRequest.newBuilder() @@ -181,7 +181,7 @@ public void storeAppliesPerCommandTargetFromMessageVisitor() { @Test public void callerCancellationAbortsStore() { - ExternalStorage storage = transformer(new HangingDriver("d1"), 0); + ExternalStorageRunner storage = transformer(new HangingDriver("d1"), 0); CancelSource caller = new CancelSource<>(CancellationException::new); caller.cancel(); Payloads message = Payloads.newBuilder().addPayloads(payload("big")).build(); @@ -190,14 +190,14 @@ public void callerCancellationAbortsStore() { CancellationException.class, () -> storage.storeBlocking(message, null, caller.token())); } - private static ExternalStorage transformer(StorageDriver driver, int threshold) { + private static ExternalStorageRunner transformer(StorageDriver driver, int threshold) { ExternalStoragePayloadTransformer payloadTransformer = ExternalStoragePayloadTransformer.fromOptions( - ExternalStorageOptions.newBuilder() + ExternalStorage.newBuilder() .setDriver(driver) .setPayloadSizeThreshold(threshold) .build()); - return new ExternalStorage(payloadTransformer, 4); + return new ExternalStorageRunner(payloadTransformer, 4); } private static Payload payload(String data) { diff --git a/temporal-sdk/src/test/java/io/temporal/payload/storage/ExternalStorageOptionsTest.java b/temporal-sdk/src/test/java/io/temporal/payload/storage/ExternalStorageTest.java similarity index 79% rename from temporal-sdk/src/test/java/io/temporal/payload/storage/ExternalStorageOptionsTest.java rename to temporal-sdk/src/test/java/io/temporal/payload/storage/ExternalStorageTest.java index b932986ad6..e68b13b3a0 100644 --- a/temporal-sdk/src/test/java/io/temporal/payload/storage/ExternalStorageOptionsTest.java +++ b/temporal-sdk/src/test/java/io/temporal/payload/storage/ExternalStorageTest.java @@ -12,7 +12,7 @@ import org.junit.Test; /** Tests external storage option validation and defaults. */ -public class ExternalStorageOptionsTest { +public class ExternalStorageTest { private static StorageDriverStoreContext storeContext(StorageDriverTargetInfo target) { return new StorageDriverStoreContext() { @@ -52,7 +52,7 @@ public CompletableFuture> retrieve( @Test public void singleDriverNoSelectorSynthesizesSelector() { StorageDriver a = driver("a"); - ExternalStorageOptions storage = ExternalStorageOptions.newBuilder().setDriver(a).build(); + ExternalStorage storage = ExternalStorage.newBuilder().setDriver(a).build(); assertEquals(1, storage.getDrivers().size()); StorageDriverSelector selector = storage.getDriverSelector(); assertNotNull(selector); @@ -62,8 +62,8 @@ public void singleDriverNoSelectorSynthesizesSelector() { @Test public void multipleDriversWithSelectorIsValid() { StorageDriver a = driver("a"); - ExternalStorageOptions storage = - ExternalStorageOptions.newBuilder() + ExternalStorage storage = + ExternalStorage.newBuilder() .setDrivers(Arrays.asList(a, driver("b"))) .setDriverSelector((context, payload) -> a) .build(); @@ -76,8 +76,8 @@ public void lastSetDriversWins() { StorageDriver a = driver("a"); StorageDriver b = driver("b"); StorageDriver c = driver("c"); - ExternalStorageOptions storage = - ExternalStorageOptions.newBuilder() + ExternalStorage storage = + ExternalStorage.newBuilder() .setDrivers(Arrays.asList(a, b)) .setDrivers(Collections.singletonList(c)) .build(); @@ -86,8 +86,8 @@ public void lastSetDriversWins() { @Test public void zeroThresholdStoresAll() { - ExternalStorageOptions storage = - ExternalStorageOptions.newBuilder() + ExternalStorage storage = + ExternalStorage.newBuilder() .setDrivers(Collections.singletonList(driver("a"))) .setPayloadSizeThreshold(0) .build(); @@ -96,24 +96,22 @@ public void zeroThresholdStoresAll() { @Test(expected = IllegalStateException.class) public void noDriversRejected() { - ExternalStorageOptions.newBuilder().build(); + ExternalStorage.newBuilder().build(); } @Test(expected = IllegalStateException.class) public void duplicateDriverNamesRejected() { - ExternalStorageOptions.newBuilder() - .setDrivers(Arrays.asList(driver("dup"), driver("dup"))) - .build(); + ExternalStorage.newBuilder().setDrivers(Arrays.asList(driver("dup"), driver("dup"))).build(); } @Test(expected = IllegalStateException.class) public void multipleDriversRequireSelector() { - ExternalStorageOptions.newBuilder().setDrivers(Arrays.asList(driver("a"), driver("b"))).build(); + ExternalStorage.newBuilder().setDrivers(Arrays.asList(driver("a"), driver("b"))).build(); } @Test(expected = IllegalStateException.class) public void negativeThresholdRejected() { - ExternalStorageOptions.newBuilder() + ExternalStorage.newBuilder() .setDrivers(Collections.singletonList(driver("a"))) .setPayloadSizeThreshold(-1) .build(); @@ -123,7 +121,7 @@ public void negativeThresholdRejected() { public void maxConcurrentPayloadVisitsDefaultsToThree() { assertEquals( 3, - ExternalStorageOptions.newBuilder() + ExternalStorage.newBuilder() .setDriver(driver("a")) .build() .getMaxConcurrentPayloadVisits()); @@ -131,9 +129,6 @@ public void maxConcurrentPayloadVisitsDefaultsToThree() { @Test(expected = IllegalStateException.class) public void zeroMaxConcurrentPayloadVisitsRejected() { - ExternalStorageOptions.newBuilder() - .setDriver(driver("a")) - .setMaxConcurrentPayloadVisits(0) - .build(); + ExternalStorage.newBuilder().setDriver(driver("a")).setMaxConcurrentPayloadVisits(0).build(); } }