From e86cb15b28649e64c861b074fbfcaac3dd9074ce Mon Sep 17 00:00:00 2001 From: Varshi Bachu Date: Mon, 20 Jul 2026 16:53:50 -0700 Subject: [PATCH] initial commit --- CHANGELOG.md | 1 + README.md | 24 ++ .../durabletask/ReplaySafeLogger.java | 300 ++++++++++++++ .../durabletask/TaskOrchestrationContext.java | 24 ++ .../durabletask/IntegrationTests.java | 99 +++++ .../durabletask/ReplaySafeLoggerTest.java | 385 ++++++++++++++++++ .../TaskOrchestrationExecutorTest.java | 103 +++++ .../java/com/functions/AzureFunctions.java | 9 +- 8 files changed, 944 insertions(+), 1 deletion(-) create mode 100644 client/src/main/java/com/microsoft/durabletask/ReplaySafeLogger.java create mode 100644 client/src/test/java/com/microsoft/durabletask/ReplaySafeLoggerTest.java diff --git a/CHANGELOG.md b/CHANGELOG.md index cee91be..46f6f9f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,4 +1,5 @@ ## Unreleased +* Add `createReplaySafeLogger` to suppress orchestration log output during replay. * Add `getParentInstance()` API to `TaskOrchestrationContext` for discovering parent orchestration info ([#284](https://github.com/microsoft/durabletask-java/pull/284)) ## v1.9.0 diff --git a/README.md b/README.md index c816022..534fc2d 100644 --- a/README.md +++ b/README.md @@ -16,6 +16,30 @@ result += ctx.callActivity("SayHello", "Seattle", String.class).await(); return result; ``` +### Replay-safe orchestration logging + +Orchestrator code re-executes while rebuilding state from history. Wrap an existing +`java.util.logging.Logger` to suppress log output during those replay segments: + +```java +Logger logger = ctx.createReplaySafeLogger( + Logger.getLogger(MyOrchestration.class.getName())); + +logger.info("Starting orchestration " + ctx.getInstanceId()); +String result = ctx.callActivity("ProcessItem", input, String.class).await(); +logger.info(() -> "Activity returned: " + result); +``` + +In Azure Functions, pass `ExecutionContext.getLogger()` instead of creating a named +logger so the output retains its invocation ID and normal host routing: + +```java +Logger logger = ctx.createReplaySafeLogger(executionContext.getLogger()); +``` + +Replay-safe logging suppresses calls made while replaying; it does not guarantee +exactly-once log delivery across failed or retried live orchestration turns. + ### Reliable fan-out / fan-in orchestration pattern ```java diff --git a/client/src/main/java/com/microsoft/durabletask/ReplaySafeLogger.java b/client/src/main/java/com/microsoft/durabletask/ReplaySafeLogger.java new file mode 100644 index 0000000..c0816aa --- /dev/null +++ b/client/src/main/java/com/microsoft/durabletask/ReplaySafeLogger.java @@ -0,0 +1,300 @@ +// Copyright (c) Microsoft Corporation. All rights reserved. +// Licensed under the MIT License. +package com.microsoft.durabletask; + +import java.util.Objects; +import java.util.ResourceBundle; +import java.util.function.BooleanSupplier; +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; + +final class ReplaySafeLogger extends Logger { + private final Logger delegate; + private final BooleanSupplier isReplaying; + + ReplaySafeLogger(Logger delegate, BooleanSupplier isReplaying) { + super(Objects.requireNonNull(delegate, "delegate").getName(), null); + this.delegate = delegate; + this.isReplaying = Objects.requireNonNull(isReplaying, "isReplaying"); + } + + @Override + public boolean isLoggable(Level level) { + return this.delegate.isLoggable(level); + } + + @Override + public void log(LogRecord record) { + if (!this.isReplaying.getAsBoolean()) { + this.delegate.log(record); + } + } + + @Override + public void log(Level level, String message) { + if (!this.isReplaying.getAsBoolean()) { + this.delegate.log(level, message); + } + } + + @Override + public void log(Level level, Supplier messageSupplier) { + if (!this.isReplaying.getAsBoolean()) { + this.delegate.log(level, messageSupplier); + } + } + + @Override + public void log(Level level, String message, Object parameter) { + if (!this.isReplaying.getAsBoolean()) { + this.delegate.log(level, message, parameter); + } + } + + @Override + public void log(Level level, String message, Object[] parameters) { + if (!this.isReplaying.getAsBoolean()) { + this.delegate.log(level, message, parameters); + } + } + + @Override + public void log(Level level, String message, Throwable thrown) { + if (!this.isReplaying.getAsBoolean()) { + this.delegate.log(level, message, thrown); + } + } + + @Override + public void log(Level level, Throwable thrown, Supplier messageSupplier) { + if (!this.isReplaying.getAsBoolean()) { + this.delegate.log(level, thrown, messageSupplier); + } + } + + @Override + public void logp( + Level level, + String sourceClass, + String sourceMethod, + String message) { + if (!this.isReplaying.getAsBoolean()) { + this.delegate.logp(level, sourceClass, sourceMethod, message); + } + } + + @Override + public void logp( + Level level, + String sourceClass, + String sourceMethod, + Supplier messageSupplier) { + if (!this.isReplaying.getAsBoolean()) { + this.delegate.logp(level, sourceClass, sourceMethod, messageSupplier); + } + } + + @Override + public void logp( + Level level, + String sourceClass, + String sourceMethod, + String message, + Object parameter) { + if (!this.isReplaying.getAsBoolean()) { + this.delegate.logp(level, sourceClass, sourceMethod, message, parameter); + } + } + + @Override + public void logp( + Level level, + String sourceClass, + String sourceMethod, + String message, + Object[] parameters) { + if (!this.isReplaying.getAsBoolean()) { + this.delegate.logp(level, sourceClass, sourceMethod, message, parameters); + } + } + + @Override + public void logp( + Level level, + String sourceClass, + String sourceMethod, + String message, + Throwable thrown) { + if (!this.isReplaying.getAsBoolean()) { + this.delegate.logp(level, sourceClass, sourceMethod, message, thrown); + } + } + + @Override + public void logp( + Level level, + String sourceClass, + String sourceMethod, + Throwable thrown, + Supplier 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); + } +} \ No newline at end of file diff --git a/client/src/main/java/com/microsoft/durabletask/TaskOrchestrationContext.java b/client/src/main/java/com/microsoft/durabletask/TaskOrchestrationContext.java index 1ae6061..71f9c9d 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 59db0b7..6fbc541 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 0000000..d2c0dd1 --- /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++; + } + } +} \ No newline at end of file diff --git a/client/src/test/java/com/microsoft/durabletask/TaskOrchestrationExecutorTest.java b/client/src/test/java/com/microsoft/durabletask/TaskOrchestrationExecutorTest.java index ce21dd9..ef89437 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,107 @@ 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.setUseParentHandlers(false); + delegate.addHandler(new Handler() { + @Override + public void publish(LogRecord record) { + messages.add(record.getMessage()); + } + + @Override + public void flush() { + } + + @Override + public void close() { + } + }); + + 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 755e50a..907238b 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; }