messageSupplier) {
+ if (!this.isReplaying.getAsBoolean()) {
+ this.delegate.logp(level, sourceClass, sourceMethod, thrown, messageSupplier);
+ }
+ }
+
+ @Override
+ public void logrb(
+ Level level,
+ String sourceClass,
+ String sourceMethod,
+ String bundleName,
+ String message) {
+ if (!this.isReplaying.getAsBoolean()) {
+ this.delegate.logrb(level, sourceClass, sourceMethod, bundleName, message);
+ }
+ }
+
+ @Override
+ public void logrb(
+ Level level,
+ String sourceClass,
+ String sourceMethod,
+ String bundleName,
+ String message,
+ Object parameter) {
+ if (!this.isReplaying.getAsBoolean()) {
+ this.delegate.logrb(level, sourceClass, sourceMethod, bundleName, message, parameter);
+ }
+ }
+
+ @Override
+ public void logrb(
+ Level level,
+ String sourceClass,
+ String sourceMethod,
+ String bundleName,
+ String message,
+ Object[] parameters) {
+ if (!this.isReplaying.getAsBoolean()) {
+ this.delegate.logrb(level, sourceClass, sourceMethod, bundleName, message, parameters);
+ }
+ }
+
+ @Override
+ public void logrb(
+ Level level,
+ String sourceClass,
+ String sourceMethod,
+ ResourceBundle bundle,
+ String message,
+ Object... parameters) {
+ if (!this.isReplaying.getAsBoolean()) {
+ this.delegate.logrb(level, sourceClass, sourceMethod, bundle, message, parameters);
+ }
+ }
+
+ @Override
+ public void logrb(
+ Level level,
+ String sourceClass,
+ String sourceMethod,
+ String bundleName,
+ String message,
+ Throwable thrown) {
+ if (!this.isReplaying.getAsBoolean()) {
+ this.delegate.logrb(level, sourceClass, sourceMethod, bundleName, message, thrown);
+ }
+ }
+
+ @Override
+ public void logrb(
+ Level level,
+ String sourceClass,
+ String sourceMethod,
+ ResourceBundle bundle,
+ String message,
+ Throwable thrown) {
+ if (!this.isReplaying.getAsBoolean()) {
+ this.delegate.logrb(level, sourceClass, sourceMethod, bundle, message, thrown);
+ }
+ }
+
+ @Override
+ public String getName() {
+ return this.delegate.getName();
+ }
+
+ @Override
+ public ResourceBundle getResourceBundle() {
+ return this.delegate.getResourceBundle();
+ }
+
+ @Override
+ public String getResourceBundleName() {
+ return this.delegate.getResourceBundleName();
+ }
+
+ @Override
+ public void setResourceBundle(ResourceBundle bundle) {
+ this.delegate.setResourceBundle(bundle);
+ }
+
+ @Override
+ public Filter getFilter() {
+ return this.delegate.getFilter();
+ }
+
+ @Override
+ public void setFilter(Filter filter) {
+ this.delegate.setFilter(filter);
+ }
+
+ @Override
+ public Level getLevel() {
+ return this.delegate.getLevel();
+ }
+
+ @Override
+ public void setLevel(Level level) {
+ this.delegate.setLevel(level);
+ }
+
+ @Override
+ public Handler[] getHandlers() {
+ return this.delegate.getHandlers();
+ }
+
+ @Override
+ public void addHandler(Handler handler) {
+ this.delegate.addHandler(handler);
+ }
+
+ @Override
+ public void removeHandler(Handler handler) {
+ this.delegate.removeHandler(handler);
+ }
+
+ @Override
+ public Logger getParent() {
+ return this.delegate.getParent();
+ }
+
+ @Override
+ public void setParent(Logger parent) {
+ this.delegate.setParent(parent);
+ }
+
+ @Override
+ public boolean getUseParentHandlers() {
+ return this.delegate.getUseParentHandlers();
+ }
+
+ @Override
+ public void setUseParentHandlers(boolean useParentHandlers) {
+ this.delegate.setUseParentHandlers(useParentHandlers);
+ }
+}
diff --git a/client/src/main/java/com/microsoft/durabletask/TaskOrchestrationContext.java b/client/src/main/java/com/microsoft/durabletask/TaskOrchestrationContext.java
index 1ae60612..71f9c9d2 100644
--- a/client/src/main/java/com/microsoft/durabletask/TaskOrchestrationContext.java
+++ b/client/src/main/java/com/microsoft/durabletask/TaskOrchestrationContext.java
@@ -10,7 +10,9 @@
import java.util.Arrays;
import java.util.List;
import java.util.Map;
+import java.util.Objects;
import java.util.UUID;
+import java.util.logging.Logger;
import javax.annotation.Nonnull;
/**
@@ -62,6 +64,28 @@ public interface TaskOrchestrationContext {
*/
boolean getIsReplaying();
+ /**
+ * Creates a logger that suppresses output while this orchestration is replaying.
+ *
+ * The returned logger does not own {@code logger} or any of its handlers. Replay-safe logging suppresses replay
+ * output but does not guarantee exactly-once delivery across failed or retried live orchestration turns.
+ *
+ * {@link Logger#isLoggable} continues to report the supplied logger's level state during replay. Use
+ * supplier-based logging methods to avoid expensive message construction while replaying.
+ *
+ * Automatic source-class and source-method inference is not preserved by all JUL convenience methods. Use
+ * {@link Logger#logp} when explicit source metadata is required.
+ *
+ * @param logger the configured logger to wrap
+ * @return a logger that emits through {@code logger} only when not replaying
+ * @throws NullPointerException if {@code logger} is {@code null}
+ */
+ default Logger createReplaySafeLogger(Logger logger) {
+ return new ReplaySafeLogger(
+ Objects.requireNonNull(logger, "logger"),
+ this::getIsReplaying);
+ }
+
/**
* Gets the version of the orchestration that this context represents.
*
diff --git a/client/src/test/java/com/microsoft/durabletask/IntegrationTests.java b/client/src/test/java/com/microsoft/durabletask/IntegrationTests.java
index 59db0b7f..6fbc541d 100644
--- a/client/src/test/java/com/microsoft/durabletask/IntegrationTests.java
+++ b/client/src/test/java/com/microsoft/durabletask/IntegrationTests.java
@@ -5,11 +5,16 @@
import java.io.IOException;
import java.time.*;
import java.util.*;
+import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicReferenceArray;
+import java.util.logging.Handler;
+import java.util.logging.Level;
+import java.util.logging.LogRecord;
+import java.util.logging.Logger;
import java.util.stream.Collectors;
import java.util.stream.IntStream;
import java.util.stream.Stream;
@@ -413,6 +418,100 @@ void subOrchestration() throws TimeoutException {
}
}
+ @Test
+ void replaySafeLogger_parentSubOrchestrationAndActivity_logEachLiveSegmentOnce() throws TimeoutException {
+ final String parentOrchestratorName = "ReplaySafeLoggerParent";
+ final String childOrchestratorName = "ReplaySafeLoggerChild";
+ final String activityName = "ReplaySafeLoggerActivity";
+ final String parentBefore = "parent before child";
+ final String parentAfter = "parent after child";
+ final String childBefore = "child before activity";
+ final String childAfter = "child after activity";
+ final String activityMessage = "activity executed";
+ final List expectedMessages = Arrays.asList(
+ parentBefore,
+ parentAfter,
+ childBefore,
+ childAfter,
+ activityMessage);
+ final Map logCounts = new ConcurrentHashMap<>();
+ final AtomicInteger parentExecutions = new AtomicInteger();
+ final AtomicInteger childExecutions = new AtomicInteger();
+ final AtomicInteger activityExecutions = new AtomicInteger();
+ final Logger delegate = Logger.getAnonymousLogger();
+ delegate.setLevel(Level.ALL);
+ delegate.setUseParentHandlers(false);
+ Handler handler = new Handler() {
+ @Override
+ public void publish(LogRecord record) {
+ logCounts.computeIfAbsent(record.getMessage(), ignored -> new AtomicInteger()).incrementAndGet();
+ }
+
+ @Override
+ public void flush() {
+ }
+
+ @Override
+ public void close() {
+ }
+ };
+ delegate.addHandler(handler);
+
+ DurableTaskGrpcWorker worker = this.createWorkerBuilder()
+ .addOrchestrator(parentOrchestratorName, ctx -> {
+ parentExecutions.incrementAndGet();
+ Logger logger = ctx.createReplaySafeLogger(delegate);
+ logger.info(parentBefore);
+ String result = ctx.callSubOrchestrator(
+ childOrchestratorName,
+ null,
+ String.class).await();
+ logger.info(parentAfter);
+ ctx.complete(result);
+ })
+ .addOrchestrator(childOrchestratorName, ctx -> {
+ childExecutions.incrementAndGet();
+ Logger logger = ctx.createReplaySafeLogger(delegate);
+ logger.info(childBefore);
+ String result = ctx.callActivity(activityName, null, String.class).await();
+ logger.info(childAfter);
+ ctx.complete(result);
+ })
+ .addActivity(activityName, ctx -> {
+ activityExecutions.incrementAndGet();
+ delegate.info(activityMessage);
+ return "done";
+ })
+ .buildAndStart();
+
+ DurableTaskClient client = this.createClientBuilder().build();
+ try (worker; client) {
+ String instanceId = client.scheduleNewOrchestrationInstance(parentOrchestratorName);
+ OrchestrationMetadata instance = client.waitForInstanceCompletion(
+ instanceId,
+ defaultTimeout,
+ true);
+
+ assertNotNull(instance);
+ assertEquals(OrchestrationRuntimeStatus.COMPLETED, instance.getRuntimeStatus());
+ assertEquals("done", instance.readOutputAs(String.class));
+ assertTrue(parentExecutions.get() >= 2, "Parent orchestrator should replay after the child completes.");
+ assertTrue(childExecutions.get() >= 2, "Child orchestrator should replay after the activity completes.");
+ assertEquals(1, activityExecutions.get());
+ for (String expectedMessage : expectedMessages) {
+ assertEquals(
+ 1,
+ logCounts.getOrDefault(expectedMessage, new AtomicInteger()).get(),
+ expectedMessage);
+ }
+ assertEquals(
+ expectedMessages.size(),
+ logCounts.values().stream().mapToInt(AtomicInteger::get).sum());
+ } finally {
+ delegate.removeHandler(handler);
+ }
+ }
+
@Test
void continueAsNew() throws TimeoutException {
final String orchestratorName = "continueAsNew";
diff --git a/client/src/test/java/com/microsoft/durabletask/ReplaySafeLoggerTest.java b/client/src/test/java/com/microsoft/durabletask/ReplaySafeLoggerTest.java
new file mode 100644
index 00000000..d01b25c2
--- /dev/null
+++ b/client/src/test/java/com/microsoft/durabletask/ReplaySafeLoggerTest.java
@@ -0,0 +1,385 @@
+// Copyright (c) Microsoft Corporation. All rights reserved.
+// Licensed under the MIT License.
+package com.microsoft.durabletask;
+
+import org.junit.jupiter.api.Test;
+
+import java.lang.reflect.Method;
+import java.lang.reflect.Modifier;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.HashSet;
+import java.util.List;
+import java.util.ListResourceBundle;
+import java.util.ResourceBundle;
+import java.util.Set;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.function.Consumer;
+import java.util.function.Supplier;
+import java.util.logging.Filter;
+import java.util.logging.Handler;
+import java.util.logging.Level;
+import java.util.logging.LogRecord;
+import java.util.logging.Logger;
+import java.util.stream.Collectors;
+
+import static org.junit.jupiter.api.Assertions.*;
+
+class ReplaySafeLoggerTest {
+ private static final ResourceBundle TEST_BUNDLE = new ListResourceBundle() {
+ @Override
+ protected Object[][] getContents() {
+ return new Object[][]{{"message", "localized message"}};
+ }
+
+ @Override
+ public String getBaseBundleName() {
+ return "test.bundle";
+ }
+ };
+
+ @Test
+ void eagerMethodFamilies_areSuppressedDuringReplayAndEmittedWhenLive() {
+ AtomicBoolean replaying = new AtomicBoolean(true);
+ RecordingHandler handler = new RecordingHandler();
+ Logger logger = newReplaySafeLogger(replaying, handler);
+ IllegalStateException failure = new IllegalStateException("failure");
+ List> logCalls = Arrays.asList(
+ value -> value.severe("severe"),
+ value -> value.warning("warning"),
+ value -> value.info("info"),
+ value -> value.config("config"),
+ value -> value.fine("fine"),
+ value -> value.finer("finer"),
+ value -> value.finest("finest"),
+ value -> value.log(Level.INFO, "log"),
+ value -> value.log(Level.INFO, "log {0}", "argument"),
+ value -> value.log(Level.INFO, "log {0}", new Object[]{"argument"}),
+ value -> value.log(Level.INFO, "log", failure),
+ value -> value.logp(Level.INFO, "Source", "method", "logp"),
+ value -> value.logp(Level.INFO, "Source", "method", "logp {0}", "argument"),
+ value -> value.logp(Level.INFO, "Source", "method", "logp {0}", new Object[]{"argument"}),
+ value -> value.logp(Level.INFO, "Source", "method", "logp", failure),
+ value -> value.logrb(Level.INFO, "Source", "method", (String) null, "logrb"),
+ value -> value.logrb(Level.INFO, "Source", "method", (String) null, "logrb {0}", "argument"),
+ value -> value.logrb(
+ Level.INFO,
+ "Source",
+ "method",
+ (String) null,
+ "logrb {0}",
+ new Object[]{"argument"}),
+ value -> value.logrb(Level.INFO, "Source", "method", (String) null, "logrb", failure),
+ value -> value.logrb(Level.INFO, "Source", "method", TEST_BUNDLE, "message", "argument"),
+ value -> value.logrb(Level.INFO, TEST_BUNDLE, "message", "argument"),
+ value -> value.logrb(Level.INFO, "Source", "method", TEST_BUNDLE, "message", failure),
+ value -> value.logrb(Level.INFO, TEST_BUNDLE, "message", failure),
+ value -> value.entering("Source", "method"),
+ value -> value.entering("Source", "method", "argument"),
+ value -> value.entering("Source", "method", new Object[]{"argument"}),
+ value -> value.exiting("Source", "method"),
+ value -> value.exiting("Source", "method", "result"),
+ value -> value.throwing("Source", "method", failure));
+
+ logCalls.forEach(call -> call.accept(logger));
+ assertTrue(handler.records.isEmpty());
+
+ replaying.set(false);
+ logCalls.forEach(call -> call.accept(logger));
+ assertEquals(logCalls.size(), handler.records.size());
+ }
+
+ @Test
+ void supplierMethods_doNotEvaluateDuringReplayAndEvaluateOnceWhenLive() {
+ AtomicBoolean replaying = new AtomicBoolean(true);
+ AtomicInteger evaluations = new AtomicInteger();
+ RecordingHandler handler = new RecordingHandler();
+ Logger logger = newReplaySafeLogger(replaying, handler);
+ IllegalStateException failure = new IllegalStateException("failure");
+ Supplier supplier = () -> {
+ evaluations.incrementAndGet();
+ return "message";
+ };
+ List> logCalls = Arrays.asList(
+ value -> value.log(Level.INFO, supplier),
+ value -> value.log(Level.INFO, failure, supplier),
+ value -> value.logp(Level.INFO, "Source", "method", supplier),
+ value -> value.logp(Level.INFO, "Source", "method", failure, supplier),
+ value -> value.info(supplier));
+
+ logCalls.forEach(call -> call.accept(logger));
+ assertEquals(0, evaluations.get());
+ assertTrue(handler.records.isEmpty());
+
+ replaying.set(false);
+ logCalls.forEach(call -> call.accept(logger));
+ assertEquals(logCalls.size(), evaluations.get());
+ assertEquals(logCalls.size(), handler.records.size());
+ }
+
+ @Test
+ void loggerTracksReplayStateTransitions() {
+ AtomicBoolean replaying = new AtomicBoolean(true);
+ RecordingHandler handler = new RecordingHandler();
+ Logger logger = newReplaySafeLogger(replaying, handler);
+
+ logger.info("replay");
+ replaying.set(false);
+ logger.info("live");
+ replaying.set(true);
+ logger.info("replay again");
+
+ assertEquals(1, handler.records.size());
+ assertEquals("live", handler.records.get(0).getMessage());
+ }
+
+ @Test
+ void isLoggableAlwaysDelegates() {
+ AtomicBoolean replaying = new AtomicBoolean(true);
+ TestLogger delegate = new TestLogger("delegate");
+ delegate.setLevel(Level.WARNING);
+ Logger logger = new ReplaySafeLogger(delegate, replaying::get);
+
+ assertFalse(logger.isLoggable(Level.INFO));
+ assertTrue(logger.isLoggable(Level.WARNING));
+
+ replaying.set(false);
+ assertFalse(logger.isLoggable(Level.INFO));
+ assertTrue(logger.isLoggable(Level.WARNING));
+ }
+
+ @Test
+ void liveRecordsUseDelegateFilterAndHandler() {
+ AtomicBoolean replaying = new AtomicBoolean(true);
+ AtomicBoolean accepted = new AtomicBoolean(false);
+ AtomicInteger filterCalls = new AtomicInteger();
+ RecordingHandler handler = new RecordingHandler();
+ TestLogger delegate = configuredLogger(handler);
+ delegate.setFilter(record -> {
+ filterCalls.incrementAndGet();
+ return accepted.get();
+ });
+ Logger logger = new ReplaySafeLogger(delegate, replaying::get);
+
+ logger.info("replay");
+ assertEquals(0, filterCalls.get());
+
+ replaying.set(false);
+ logger.info("filtered");
+ assertEquals(1, filterCalls.get());
+ assertTrue(handler.records.isEmpty());
+
+ accepted.set(true);
+ logger.info("accepted");
+ assertEquals(2, filterCalls.get());
+ assertEquals(1, handler.records.size());
+ assertEquals("accepted", handler.records.get(0).getMessage());
+ }
+
+ @Test
+ void explicitLogRecordMetadataIsPreserved() {
+ RecordingHandler handler = new RecordingHandler();
+ Logger logger = new ReplaySafeLogger(configuredLogger(handler), () -> false);
+ IllegalStateException failure = new IllegalStateException("failure");
+ LogRecord record = new LogRecord(Level.WARNING, "message");
+ record.setLoggerName("category");
+ record.setParameters(new Object[]{"argument"});
+ record.setThrown(failure);
+ record.setSourceClassName("CustomerOrchestrator");
+ record.setSourceMethodName("run");
+ record.setResourceBundle(TEST_BUNDLE);
+ record.setResourceBundleName(TEST_BUNDLE.getBaseBundleName());
+
+ logger.log(record);
+
+ assertEquals(1, handler.records.size());
+ LogRecord actual = handler.records.get(0);
+ assertSame(record, actual);
+ assertEquals(Level.WARNING, actual.getLevel());
+ assertEquals("message", actual.getMessage());
+ assertEquals("category", actual.getLoggerName());
+ assertArrayEquals(new Object[]{"argument"}, actual.getParameters());
+ assertSame(failure, actual.getThrown());
+ assertEquals("CustomerOrchestrator", actual.getSourceClassName());
+ assertEquals("run", actual.getSourceMethodName());
+ assertSame(TEST_BUNDLE, actual.getResourceBundle());
+ assertEquals(TEST_BUNDLE.getBaseBundleName(), actual.getResourceBundleName());
+ }
+
+ @Test
+ void explicitSourceAndDirectResourceBundleArePreserved() {
+ RecordingHandler handler = new RecordingHandler();
+ TestLogger delegate = configuredLogger(handler);
+ delegate.setResourceBundle(TEST_BUNDLE);
+ Logger logger = new ReplaySafeLogger(delegate, () -> false);
+
+ logger.logp(Level.INFO, "CustomerOrchestrator", "run", "message");
+
+ assertEquals(1, handler.records.size());
+ LogRecord record = handler.records.get(0);
+ assertEquals("CustomerOrchestrator", record.getSourceClassName());
+ assertEquals("run", record.getSourceMethodName());
+ assertSame(TEST_BUNDLE, record.getResourceBundle());
+ assertEquals(TEST_BUNDLE.getBaseBundleName(), record.getResourceBundleName());
+ }
+
+ @Test
+ void parentInheritedResourceBundleIsPreserved() {
+ RecordingHandler handler = new RecordingHandler();
+ TestLogger parent = new TestLogger("parent");
+ parent.setResourceBundle(TEST_BUNDLE);
+ TestLogger delegate = configuredLogger(handler);
+ delegate.setParent(parent);
+ Logger logger = new ReplaySafeLogger(delegate, () -> false);
+
+ logger.info("message");
+
+ assertEquals(1, handler.records.size());
+ LogRecord record = handler.records.get(0);
+ assertSame(TEST_BUNDLE, record.getResourceBundle());
+ assertEquals(TEST_BUNDLE.getBaseBundleName(), record.getResourceBundleName());
+ }
+
+ @Test
+ void configurationMethodsOperateOnDelegate() {
+ TestLogger delegate = new TestLogger("delegate");
+ TestLogger parent = new TestLogger("parent");
+ Logger logger = new ReplaySafeLogger(delegate, () -> false);
+ Filter filter = record -> true;
+ RecordingHandler handler = new RecordingHandler();
+
+ logger.setResourceBundle(TEST_BUNDLE);
+ logger.setFilter(filter);
+ logger.setLevel(Level.FINE);
+ logger.addHandler(handler);
+ logger.setParent(parent);
+ logger.setUseParentHandlers(false);
+
+ assertEquals("delegate", logger.getName());
+ assertSame(TEST_BUNDLE, delegate.getResourceBundle());
+ assertSame(TEST_BUNDLE, logger.getResourceBundle());
+ assertEquals(TEST_BUNDLE.getBaseBundleName(), logger.getResourceBundleName());
+ assertSame(filter, delegate.getFilter());
+ assertSame(filter, logger.getFilter());
+ assertEquals(Level.FINE, delegate.getLevel());
+ assertEquals(Level.FINE, logger.getLevel());
+ assertArrayEquals(new Handler[]{handler}, delegate.getHandlers());
+ assertArrayEquals(new Handler[]{handler}, logger.getHandlers());
+ assertNotSame(logger.getHandlers(), logger.getHandlers());
+ assertSame(parent, delegate.getParent());
+ assertSame(parent, logger.getParent());
+ assertFalse(delegate.getUseParentHandlers());
+ assertFalse(logger.getUseParentHandlers());
+
+ logger.removeHandler(handler);
+ assertEquals(0, delegate.getHandlers().length);
+ assertEquals(0, handler.closeCalls);
+ }
+
+ @Test
+ void constructorRejectsNullDependencies() {
+ NullPointerException delegateException = assertThrows(
+ NullPointerException.class,
+ () -> new ReplaySafeLogger(null, () -> false));
+ NullPointerException replayException = assertThrows(
+ NullPointerException.class,
+ () -> new ReplaySafeLogger(new TestLogger("delegate"), null));
+
+ assertEquals("delegate", delegateException.getMessage());
+ assertEquals("isReplaying", replayException.getMessage());
+ }
+
+ @Test
+ void allPublicLoggerMethodsAreClassified() {
+ Set expectedInheritedEmissionMethods = new HashSet<>(Arrays.asList(
+ signature("logrb", Level.class, ResourceBundle.class, String.class, Object[].class),
+ signature("logrb", Level.class, ResourceBundle.class, String.class, Throwable.class),
+ signature("entering", String.class, String.class),
+ signature("entering", String.class, String.class, Object.class),
+ signature("entering", String.class, String.class, Object[].class),
+ signature("exiting", String.class, String.class),
+ signature("exiting", String.class, String.class, Object.class),
+ signature("throwing", String.class, String.class, Throwable.class),
+ signature("severe", String.class),
+ signature("warning", String.class),
+ signature("info", String.class),
+ signature("config", String.class),
+ signature("fine", String.class),
+ signature("finer", String.class),
+ signature("finest", String.class),
+ signature("severe", Supplier.class),
+ signature("warning", Supplier.class),
+ signature("info", Supplier.class),
+ signature("config", Supplier.class),
+ signature("fine", Supplier.class),
+ signature("finer", Supplier.class),
+ signature("finest", Supplier.class)));
+
+ Set actualInheritedMethods = Arrays.stream(Logger.class.getDeclaredMethods())
+ .filter(method -> Modifier.isPublic(method.getModifiers()))
+ .filter(method -> !Modifier.isStatic(method.getModifiers()))
+ .filter(method -> !Modifier.isFinal(method.getModifiers()))
+ .filter(method -> !isOverridden(method))
+ .map(ReplaySafeLoggerTest::signature)
+ .collect(Collectors.toSet());
+
+ assertEquals(expectedInheritedEmissionMethods, actualInheritedMethods);
+ }
+
+ private static Logger newReplaySafeLogger(AtomicBoolean replaying, RecordingHandler handler) {
+ return new ReplaySafeLogger(configuredLogger(handler), replaying::get);
+ }
+
+ private static TestLogger configuredLogger(RecordingHandler handler) {
+ TestLogger logger = new TestLogger("delegate");
+ logger.setLevel(Level.ALL);
+ logger.setUseParentHandlers(false);
+ logger.addHandler(handler);
+ return logger;
+ }
+
+ private static boolean isOverridden(Method method) {
+ try {
+ ReplaySafeLogger.class.getDeclaredMethod(method.getName(), method.getParameterTypes());
+ return true;
+ } catch (NoSuchMethodException ignored) {
+ return false;
+ }
+ }
+
+ private static String signature(Method method) {
+ return signature(method.getName(), method.getParameterTypes());
+ }
+
+ private static String signature(String name, Class>... parameterTypes) {
+ return name + Arrays.stream(parameterTypes)
+ .map(Class::getName)
+ .collect(Collectors.joining(",", "(", ")"));
+ }
+
+ private static final class TestLogger extends Logger {
+ TestLogger(String name) {
+ super(name, null);
+ }
+ }
+
+ private static final class RecordingHandler extends Handler {
+ private final List records = new ArrayList<>();
+ private int closeCalls;
+
+ @Override
+ public void publish(LogRecord record) {
+ this.records.add(record);
+ }
+
+ @Override
+ public void flush() {
+ }
+
+ @Override
+ public void close() {
+ this.closeCalls++;
+ }
+ }
+}
diff --git a/client/src/test/java/com/microsoft/durabletask/TaskOrchestrationExecutorTest.java b/client/src/test/java/com/microsoft/durabletask/TaskOrchestrationExecutorTest.java
index ce21dd9d..47e71529 100644
--- a/client/src/test/java/com/microsoft/durabletask/TaskOrchestrationExecutorTest.java
+++ b/client/src/test/java/com/microsoft/durabletask/TaskOrchestrationExecutorTest.java
@@ -12,6 +12,8 @@
import java.time.Instant;
import java.time.ZonedDateTime;
import java.util.*;
+import java.util.logging.Handler;
+import java.util.logging.LogRecord;
import java.util.logging.Logger;
import java.util.stream.Collectors;
@@ -673,6 +675,110 @@ public TaskOrchestration create() {
assertEquals("parent-456", captured[0].getInstanceId());
}
+ @Test
+ void execute_replaySafeLogger_suppressesReplayAndLogsLiveSegments() {
+ String orchestrationName = "LoggingOrchestration";
+ String activityName = "LoggingActivity";
+ List messages = new ArrayList<>();
+ Logger delegate = Logger.getAnonymousLogger();
+ delegate.setLevel(java.util.logging.Level.ALL);
+ delegate.setUseParentHandlers(false);
+ Handler handler = new Handler() {
+ @Override
+ public void publish(LogRecord record) {
+ messages.add(record.getMessage());
+ }
+
+ @Override
+ public void flush() {
+ }
+
+ @Override
+ public void close() {
+ }
+ };
+ handler.setLevel(java.util.logging.Level.ALL);
+ delegate.addHandler(handler);
+
+ HashMap factories = new HashMap<>();
+ factories.put(orchestrationName, new TaskOrchestrationFactory() {
+ @Override
+ public String getName() {
+ return orchestrationName;
+ }
+
+ @Override
+ public TaskOrchestration create() {
+ return ctx -> {
+ NullPointerException exception = assertThrows(
+ NullPointerException.class,
+ () -> ctx.createReplaySafeLogger(null));
+ assertEquals("logger", exception.getMessage());
+
+ Logger logger = ctx.createReplaySafeLogger(delegate);
+ logger.info("before activity");
+ String result = ctx.callActivity(activityName, null, String.class).await();
+ logger.info("after activity: " + result);
+ ctx.complete(result);
+ };
+ }
+ });
+
+ TaskOrchestrationExecutor executor = new TaskOrchestrationExecutor(
+ factories, new JacksonDataConverter(), Duration.ofDays(3), logger, null);
+ HistoryEvent orchestratorStarted = HistoryEvent.newBuilder()
+ .setEventId(-1)
+ .setTimestamp(Timestamp.getDefaultInstance())
+ .setOrchestratorStarted(OrchestratorStartedEvent.getDefaultInstance())
+ .build();
+ HistoryEvent executionStarted = HistoryEvent.newBuilder()
+ .setEventId(-1)
+ .setTimestamp(Timestamp.getDefaultInstance())
+ .setExecutionStarted(ExecutionStartedEvent.newBuilder()
+ .setName(orchestrationName)
+ .setVersion(StringValue.of(""))
+ .setInput(StringValue.of(""))
+ .setOrchestrationInstance(OrchestrationInstance.newBuilder()
+ .setInstanceId("logging-instance")
+ .build())
+ .build())
+ .build();
+ HistoryEvent orchestratorCompleted = HistoryEvent.newBuilder()
+ .setEventId(-1)
+ .setTimestamp(Timestamp.getDefaultInstance())
+ .setOrchestratorCompleted(OrchestratorCompletedEvent.getDefaultInstance())
+ .build();
+
+ executor.execute(
+ Collections.emptyList(),
+ Arrays.asList(orchestratorStarted, executionStarted, orchestratorCompleted),
+ null);
+ assertEquals(Collections.singletonList("before activity"), messages);
+
+ HistoryEvent taskScheduled = HistoryEvent.newBuilder()
+ .setEventId(0)
+ .setTimestamp(Timestamp.getDefaultInstance())
+ .setTaskScheduled(TaskScheduledEvent.newBuilder()
+ .setName(activityName)
+ .build())
+ .build();
+ HistoryEvent taskCompleted = HistoryEvent.newBuilder()
+ .setEventId(1)
+ .setTimestamp(Timestamp.getDefaultInstance())
+ .setTaskCompleted(TaskCompletedEvent.newBuilder()
+ .setTaskScheduledId(0)
+ .setResult(StringValue.of("\"done\""))
+ .build())
+ .build();
+
+ executor.execute(
+ Arrays.asList(orchestratorStarted, executionStarted, taskScheduled, orchestratorCompleted),
+ Arrays.asList(orchestratorStarted, taskCompleted, orchestratorCompleted),
+ null);
+
+ assertEquals(Arrays.asList("before activity", "after activity: done"), messages);
+ }
+
@Test
void parentOrchestrationInstance_equalsAndHashCode() {
ParentOrchestrationInstance a = new ParentOrchestrationInstance("Orch", "id-1");
diff --git a/samples-azure-functions/src/main/java/com/functions/AzureFunctions.java b/samples-azure-functions/src/main/java/com/functions/AzureFunctions.java
index 755e50ac..907238bf 100644
--- a/samples-azure-functions/src/main/java/com/functions/AzureFunctions.java
+++ b/samples-azure-functions/src/main/java/com/functions/AzureFunctions.java
@@ -4,6 +4,7 @@
import com.microsoft.azure.functions.annotation.*;
import com.microsoft.azure.functions.*;
import java.util.*;
+import java.util.logging.Logger;
import com.microsoft.durabletask.*;
import com.microsoft.durabletask.azurefunctions.DurableActivityTrigger;
@@ -37,12 +38,18 @@ public HttpResponseMessage startOrchestration(
*/
@FunctionName("Cities")
public String citiesOrchestrator(
- @DurableOrchestrationTrigger(name = "ctx") TaskOrchestrationContext ctx) {
+ @DurableOrchestrationTrigger(name = "ctx") TaskOrchestrationContext ctx,
+ final ExecutionContext context) {
+ Logger logger = ctx.createReplaySafeLogger(context.getLogger());
+ logger.info("Starting Cities orchestration.");
+
String result = "";
result += ctx.callActivity("Capitalize", "Tokyo", String.class).await() + ", ";
result += ctx.callActivity("Capitalize", "London", String.class).await() + ", ";
result += ctx.callActivity("Capitalize", "Seattle", String.class).await() + ", ";
result += ctx.callActivity("Capitalize", "Austin", String.class).await();
+
+ logger.info("Cities orchestration completed.");
return result;
}