This is an automated email from the ASF dual-hosted git repository.

merlimat pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/bookkeeper.git


The following commit(s) were added to refs/heads/master by this push:
     new 9ea1eea418 Add SingleThreadExecutor.executeOrRun and issue reads 
inline from the ledger's own thread (#4883)
9ea1eea418 is described below

commit 9ea1eea418d106c9815622a7bd08bddb0312c075
Author: Matteo Merli <[email protected]>
AuthorDate: Fri Sep 11 09:46:59 2026 -0700

    Add SingleThreadExecutor.executeOrRun and issue reads inline from the 
ledger's own thread (#4883)
    
    * Run reads inline when issued from the ledger's own thread
    
    Add SingleThreadExecutor.isCurrentThread() and executeOrRun(Runnable): the
    task runs inline when called from the executor's own thread and is queued
    otherwise. The shortcut is opt-in because it bypasses the queue: a task
    submitted this way from the executor thread runs before the tasks already
    queued, nested inside the submitting task. execute() keeps its FIFO 
contract.
    Inline failures are logged and counted like those of queued tasks.
    
    LedgerHandle uses it for reads on a writable handle, which were queued onto
    the ledger's thread even when the caller was already on it. Adds already
    initiate inline on the caller's thread and reads on read-only handles 
already
    bypass the executor, so this was the remaining same-thread hop. Applications
    that run their own work on a thread from OrderedExecutor.chooseThread(...) 
can
    use the same check for their re-dispatches.
    
    * Type LedgerHandle.executor as SingleThreadExecutor and share the safe-run 
path
    
    Review follow-ups: the inline path of SingleThreadExecutor.executeOrRun
    calls safeRunTask directly, so inline and queued tasks share the same
    failure isolation, logging and counters; the queued path only adds the
    pending-count decrement, which must not apply to inline runs.
    
    LedgerHandle keeps its worker thread as a SingleThreadExecutor instead of
    checking with instanceof: the client's main worker pool is a plain
    OrderedExecutor, whose threads are SingleThreadExecutor instances. The two
    test fixtures that handed the client an OrderedScheduler as its worker pool
    (MockClientContext, MockBookKeeperTestCase) now use an OrderedExecutor, with
    the mock bookie client dispatching on that same pool.
    
    * Fix checkstyle line length in MockClientContext and reject executeOrRun 
after shutdown
    
    The scheduler line in MockClientContext.create was one character over the
    limit; the PR validation job runs checkstyle over all modules and failed on
    it.
    
    executeOrRun now rejects the task once the executor is shut down, like
    execute() does, instead of still running it inline when called from the
    executor thread.
    
    * Expose the thread identity through the OrderedExecutor decorators
    
    OrderedExecutor wraps each worker thread in a forwarding decorator when
    task tracing or MDC preservation is enabled, so chooseThread() returned the
    wrapper rather than the SingleThreadExecutor and the cast in the
    LedgerHandle constructor failed for such clients. The auditor and the
    replication worker run with task tracing enabled in the tests, so every
    ledger open there failed with UnexpectedConditionException and
    re-replication never completed.
    
    Add the ThreadBoundExecutor interface (isCurrentThread, executeOrRun),
    implemented by SingleThreadExecutor and by the decorator, which applies the
    timing and MDC wrappers to inline runs as well. LedgerHandle keeps its
    executor as a ThreadBoundExecutor.
---
 .../bookkeeper/common/util/OrderedExecutor.java    | 111 ++++++++++------
 .../common/util/SingleThreadExecutor.java          |  46 ++++++-
 .../common/util/ThreadBoundExecutor.java           |  46 +++++++
 .../common/util/TestOrderedExecutorDecorators.java |  26 ++++
 .../common/util/TestSingleThreadExecutor.java      | 105 +++++++++++++++
 .../org/apache/bookkeeper/client/LedgerHandle.java |  20 +--
 .../client/LedgerHandleInlineReadTest.java         | 148 +++++++++++++++++++++
 .../bookkeeper/client/MockBookKeeperTestCase.java  |   2 +-
 .../bookkeeper/client/MockClientContext.java       |   8 +-
 9 files changed, 450 insertions(+), 62 deletions(-)

diff --git 
a/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/OrderedExecutor.java
 
b/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/OrderedExecutor.java
index f32d25c3ce..699c1b1395 100644
--- 
a/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/OrderedExecutor.java
+++ 
b/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/OrderedExecutor.java
@@ -306,58 +306,81 @@ public class OrderedExecutor implements ExecutorService {
     }
 
     protected ExecutorService addExecutorDecorators(ExecutorService executor) {
-        return new ForwardingExecutorService() {
-            @Override
-            protected ExecutorService delegate() {
-                return executor;
-            }
+        checkArgument(executor instanceof ThreadBoundExecutor, "Expected a 
ThreadBoundExecutor, got %s", executor);
+        return new DecoratedThread((ThreadBoundExecutor) executor);
+    }
 
-            @Override
-            public <T> List<Future<T>> invokeAll(Collection<? extends 
Callable<T>> tasks)
-                    throws InterruptedException {
-                return super.invokeAll(timedCallables(tasks));
-            }
+    /**
+     * One of the pool's threads with the task tracing and MDC decorators 
applied, which keeps exposing the thread
+     * identity of the underlying executor.
+     */
+    private class DecoratedThread extends ForwardingExecutorService implements 
ThreadBoundExecutor {
+        private final ThreadBoundExecutor executor;
 
-            @Override
-            public <T> List<Future<T>> invokeAll(Collection<? extends 
Callable<T>> tasks,
-                                                 long timeout, TimeUnit unit)
-                    throws InterruptedException {
-                return super.invokeAll(timedCallables(tasks), timeout, unit);
-            }
+        DecoratedThread(ThreadBoundExecutor executor) {
+            this.executor = executor;
+        }
 
-            @Override
-            public <T> T invokeAny(Collection<? extends Callable<T>> tasks)
-                    throws InterruptedException, ExecutionException {
-                return super.invokeAny(timedCallables(tasks));
-            }
+        @Override
+        protected ExecutorService delegate() {
+            return executor;
+        }
 
-            @Override
-            public <T> T invokeAny(Collection<? extends Callable<T>> tasks,
-                                   long timeout, TimeUnit unit)
-                    throws InterruptedException, ExecutionException, 
TimeoutException {
-                return super.invokeAny(timedCallables(tasks), timeout, unit);
-            }
+        @Override
+        public boolean isCurrentThread() {
+            return executor.isCurrentThread();
+        }
 
-            @Override
-            public void execute(Runnable command) {
-                super.execute(timedRunnable(command));
-            }
+        @Override
+        public void executeOrRun(Runnable r) {
+            executor.executeOrRun(timedRunnable(r));
+        }
 
-            @Override
-            public <T> Future<T> submit(Callable<T> task) {
-                return super.submit(timedCallable(task));
-            }
+        @Override
+        public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> 
tasks)
+                throws InterruptedException {
+            return super.invokeAll(timedCallables(tasks));
+        }
 
-            @Override
-            public Future<?> submit(Runnable task) {
-                return super.submit(timedRunnable(task));
-            }
+        @Override
+        public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> 
tasks,
+                                             long timeout, TimeUnit unit)
+                throws InterruptedException {
+            return super.invokeAll(timedCallables(tasks), timeout, unit);
+        }
 
-            @Override
-            public <T> Future<T> submit(Runnable task, T result) {
-                return super.submit(timedRunnable(task), result);
-            }
-        };
+        @Override
+        public <T> T invokeAny(Collection<? extends Callable<T>> tasks)
+                throws InterruptedException, ExecutionException {
+            return super.invokeAny(timedCallables(tasks));
+        }
+
+        @Override
+        public <T> T invokeAny(Collection<? extends Callable<T>> tasks,
+                               long timeout, TimeUnit unit)
+                throws InterruptedException, ExecutionException, 
TimeoutException {
+            return super.invokeAny(timedCallables(tasks), timeout, unit);
+        }
+
+        @Override
+        public void execute(Runnable command) {
+            super.execute(timedRunnable(command));
+        }
+
+        @Override
+        public <T> Future<T> submit(Callable<T> task) {
+            return super.submit(timedCallable(task));
+        }
+
+        @Override
+        public Future<?> submit(Runnable task) {
+            return super.submit(timedRunnable(task));
+        }
+
+        @Override
+        public <T> Future<T> submit(Runnable task, T result) {
+            return super.submit(timedRunnable(task), result);
+        }
     }
 
     /**
diff --git 
a/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/SingleThreadExecutor.java
 
b/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/SingleThreadExecutor.java
index 840a4c6dec..e26a3da5bc 100644
--- 
a/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/SingleThreadExecutor.java
+++ 
b/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/SingleThreadExecutor.java
@@ -24,7 +24,6 @@ import java.util.ArrayList;
 import java.util.List;
 import java.util.concurrent.AbstractExecutorService;
 import java.util.concurrent.CountDownLatch;
-import java.util.concurrent.ExecutorService;
 import java.util.concurrent.RejectedExecutionException;
 import java.util.concurrent.ThreadFactory;
 import java.util.concurrent.TimeUnit;
@@ -45,7 +44,7 @@ import org.apache.bookkeeper.stats.StatsLogger;
  * proceed with the next tasks.
  */
 @CustomLog
-public class SingleThreadExecutor extends AbstractExecutorService implements 
ExecutorService, Runnable {
+public class SingleThreadExecutor extends AbstractExecutorService implements 
ThreadBoundExecutor, Runnable {
 
     private static final int MAX_DRAIN_BATCH_SIZE = 1024;
 
@@ -120,7 +119,7 @@ public class SingleThreadExecutor extends 
AbstractExecutorService implements Exe
                 for (int i = 0; i < n; i++) {
                     Runnable task = localTasks[i];
                     localTasks[i] = null;
-                    if (!safeRunTask(task)) {
+                    if (!runQueuedTask(task)) {
                         return;
                     }
                 }
@@ -129,7 +128,7 @@ public class SingleThreadExecutor extends 
AbstractExecutorService implements Exe
             // Clear the queue in orderly shutdown
             Runnable task;
             while ((task = queue.poll()) != null) {
-                safeRunTask(task);
+                runQueuedTask(task);
             }
         } catch (InterruptedException ie) {
             // Exit loop when interrupted
@@ -142,6 +141,19 @@ public class SingleThreadExecutor extends 
AbstractExecutorService implements Exe
         }
     }
 
+    private boolean runQueuedTask(Runnable r) {
+        try {
+            return safeRunTask(r);
+        } finally {
+            decrementPendingTaskCount(1);
+        }
+    }
+
+    /**
+     * Runs a task, logging and counting a failure instead of propagating it.
+     *
+     * @return false when the task was interrupted
+     */
     private boolean safeRunTask(Runnable r) {
         try {
             r.run();
@@ -154,8 +166,6 @@ public class SingleThreadExecutor extends 
AbstractExecutorService implements Exe
                 tasksFailed.increment();
                 log.error().exception(t).log("Error while running task");
             }
-        } finally {
-            decrementPendingTaskCount(1);
         }
 
         return true;
@@ -220,6 +230,30 @@ public class SingleThreadExecutor extends 
AbstractExecutorService implements Exe
         executeRunnableOrList(r, null);
     }
 
+    @Override
+    public boolean isCurrentThread() {
+        return Thread.currentThread() == runner;
+    }
+
+    /**
+     * {@inheritDoc}
+     *
+     * <p>Failures of an inline run are logged and counted like those of 
queued tasks.
+     */
+    @Override
+    public void executeOrRun(Runnable r) {
+        if (state != State.Running) {
+            throw new RejectedExecutionException("Executor is shutting down");
+        }
+
+        if (isCurrentThread()) {
+            tasksCount.increment();
+            safeRunTask(r);
+        } else {
+            execute(r);
+        }
+    }
+
     @VisibleForTesting
     void executeRunnableOrList(Runnable runnable, List<Runnable> runnableList) 
{
         if (state != State.Running) {
diff --git 
a/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/ThreadBoundExecutor.java
 
b/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/ThreadBoundExecutor.java
new file mode 100644
index 0000000000..7d9098a3a2
--- /dev/null
+++ 
b/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/ThreadBoundExecutor.java
@@ -0,0 +1,46 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.bookkeeper.common.util;
+
+import java.util.concurrent.ExecutorService;
+
+/**
+ * An {@link ExecutorService} backed by a single thread, which can tell 
whether the caller is already on that
+ * thread and then run a task inline instead of queueing it.
+ *
+ * <p>Implemented by {@link SingleThreadExecutor} and by the threads an {@link 
OrderedExecutor} hands out from
+ * {@link OrderedExecutor#chooseThread(long)}, which are decorated when task 
tracing or MDC preservation is enabled.
+ */
+public interface ThreadBoundExecutor extends ExecutorService {
+
+    /**
+     * Whether the calling thread is the thread of this executor.
+     */
+    boolean isCurrentThread();
+
+    /**
+     * Runs the task inline when called from this executor's own thread, 
otherwise submits it like
+     * {@link #execute(Runnable)}. Like {@code execute}, it rejects the task 
once the executor is shut down.
+     *
+     * <p>The inline run bypasses the queue: a task submitted this way from 
the executor thread runs before the
+     * tasks already queued, nested inside the task that submitted it. Use it 
only where that reordering is
+     * acceptable.
+     */
+    void executeOrRun(Runnable r);
+}
diff --git 
a/bookkeeper-common/src/test/java/org/apache/bookkeeper/common/util/TestOrderedExecutorDecorators.java
 
b/bookkeeper-common/src/test/java/org/apache/bookkeeper/common/util/TestOrderedExecutorDecorators.java
index 69b0570899..475215313d 100644
--- 
a/bookkeeper-common/src/test/java/org/apache/bookkeeper/common/util/TestOrderedExecutorDecorators.java
+++ 
b/bookkeeper-common/src/test/java/org/apache/bookkeeper/common/util/TestOrderedExecutorDecorators.java
@@ -22,7 +22,9 @@
 package org.apache.bookkeeper.common.util;
 
 import static org.hamcrest.Matchers.hasItem;
+import static org.junit.Assert.assertFalse;
 import static org.junit.Assert.assertThat;
+import static org.junit.Assert.assertTrue;
 import static org.mockito.AdditionalAnswers.answerVoid;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.Mockito.doAnswer;
@@ -30,6 +32,7 @@ import static org.mockito.Mockito.spy;
 
 import java.util.Queue;
 import java.util.UUID;
+import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.ConcurrentLinkedQueue;
 import java.util.concurrent.TimeUnit;
 import lombok.CustomLog;
@@ -82,6 +85,29 @@ public class TestOrderedExecutorDecorators {
         ThreadContext.clearMap();
     }
 
+    @Test
+    public void testDecoratedThreadsAreThreadBound() throws Exception {
+        for (OrderedExecutor executor : new OrderedExecutor[] {
+                
OrderedExecutor.newBuilder().name("traced").numThreads(2).traceTaskExecution(true).build(),
+                
OrderedExecutor.newBuilder().name("mdc").numThreads(2).preserveMdcForTaskExecution(true).build()})
 {
+            try {
+                ThreadBoundExecutor thread = (ThreadBoundExecutor) 
executor.chooseThread(10);
+                assertFalse(thread.isCurrentThread());
+
+                // From the thread itself, executeOrRun runs the task before 
returning
+                CompletableFuture<Boolean> ranInline = new 
CompletableFuture<>();
+                thread.execute(() -> {
+                    boolean[] ran = new boolean[1];
+                    thread.executeOrRun(() -> ran[0] = 
thread.isCurrentThread());
+                    ranInline.complete(ran[0]);
+                });
+                assertTrue(ranInline.get(10, TimeUnit.SECONDS));
+            } finally {
+                executor.shutdown();
+            }
+        }
+    }
+
     @Test
     public void testMDCInvokeOrdered() throws Exception {
         OrderedExecutor executor = OrderedExecutor.newBuilder()
diff --git 
a/bookkeeper-common/src/test/java/org/apache/bookkeeper/common/util/TestSingleThreadExecutor.java
 
b/bookkeeper-common/src/test/java/org/apache/bookkeeper/common/util/TestSingleThreadExecutor.java
index ed72704e97..0f6cdadf76 100644
--- 
a/bookkeeper-common/src/test/java/org/apache/bookkeeper/common/util/TestSingleThreadExecutor.java
+++ 
b/bookkeeper-common/src/test/java/org/apache/bookkeeper/common/util/TestSingleThreadExecutor.java
@@ -20,6 +20,8 @@ package org.apache.bookkeeper.common.util;
 
 import static org.junit.Assert.assertEquals;
 import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertNotSame;
+import static org.junit.Assert.assertNull;
 import static org.junit.Assert.assertThrows;
 import static org.junit.Assert.assertTrue;
 import static org.junit.Assert.fail;
@@ -30,6 +32,8 @@ import java.util.ArrayList;
 import java.util.Collections;
 import java.util.List;
 import java.util.concurrent.BrokenBarrierException;
+import java.util.concurrent.Callable;
+import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.CountDownLatch;
 import java.util.concurrent.CyclicBarrier;
 import java.util.concurrent.Future;
@@ -38,6 +42,7 @@ import java.util.concurrent.ThreadFactory;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.TimeoutException;
 import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicReference;
 import lombok.Cleanup;
 import org.awaitility.Awaitility;
 import org.junit.Test;
@@ -369,4 +374,104 @@ public class TestSingleThreadExecutor {
 
         future.get();
     }
+
+    @Test
+    public void testIsCurrentThread() throws Exception {
+        @Cleanup("shutdown")
+        SingleThreadExecutor ste = new SingleThreadExecutor(THREAD_FACTORY);
+
+        assertFalse(ste.isCurrentThread());
+        Callable<Boolean> onExecutorThread = ste::isCurrentThread;
+        assertTrue(ste.submit(onExecutorThread).get(10, TimeUnit.SECONDS));
+    }
+
+    @Test
+    public void testExecuteOrRunInlineOnOwnThread() throws Exception {
+        @Cleanup("shutdown")
+        SingleThreadExecutor ste = new SingleThreadExecutor(THREAD_FACTORY);
+
+        // From the executor thread the task runs before executeOrRun returns, 
ahead of the queued task
+        List<String> events = Collections.synchronizedList(new ArrayList<>());
+        CompletableFuture<List<String>> seenBeforeReturn = new 
CompletableFuture<>();
+        ste.execute(() -> {
+            ste.execute(() -> events.add("queued"));
+            ste.executeOrRun(() -> events.add("inline"));
+            seenBeforeReturn.complete(new ArrayList<>(events));
+        });
+
+        assertEquals(Lists.newArrayList("inline"), seenBeforeReturn.get(10, 
TimeUnit.SECONDS));
+        Awaitility.await().until(() -> events.size() == 2);
+        assertEquals(Lists.newArrayList("inline", "queued"), events);
+        Awaitility.await().until(() -> ste.getCompletedTasksCount() == 3);
+        assertEquals(3, ste.getSubmittedTasksCount());
+        assertEquals(0, ste.getQueuedTasksCount());
+    }
+
+    @Test
+    public void testExecuteOrRunFromOtherThreadIsQueued() throws Exception {
+        @Cleanup("shutdown")
+        SingleThreadExecutor ste = new SingleThreadExecutor(THREAD_FACTORY);
+
+        // Hold the executor thread so that a queued task cannot run yet
+        CountDownLatch release = new CountDownLatch(1);
+        ste.execute(() -> {
+            try {
+                release.await();
+            } catch (InterruptedException e) {
+                Thread.currentThread().interrupt();
+            }
+        });
+
+        AtomicReference<Thread> ranOn = new AtomicReference<>();
+        ste.executeOrRun(() -> ranOn.set(Thread.currentThread()));
+        assertNull(ranOn.get());
+
+        release.countDown();
+        Awaitility.await().until(() -> ranOn.get() != null);
+        assertNotSame(Thread.currentThread(), ranOn.get());
+        Callable<Boolean> ranOnExecutorThread = () -> ranOn.get() == 
Thread.currentThread();
+        assertTrue(ste.submit(ranOnExecutorThread).get(10, TimeUnit.SECONDS));
+    }
+
+    @Test
+    public void testExecuteOrRunRejectedAfterShutdown() throws Exception {
+        SingleThreadExecutor ste = new SingleThreadExecutor(THREAD_FACTORY);
+
+        // Shut down from the executor thread itself: the inline path must 
reject like execute() does
+        CompletableFuture<Boolean> rejectedInline = new CompletableFuture<>();
+        ste.execute(() -> {
+            ste.shutdown();
+            try {
+                ste.executeOrRun(() -> {
+                });
+                rejectedInline.complete(false);
+            } catch (RejectedExecutionException e) {
+                rejectedInline.complete(true);
+            }
+        });
+        assertTrue(rejectedInline.get(10, TimeUnit.SECONDS));
+
+        ste.awaitTermination(10, TimeUnit.SECONDS);
+        assertThrows(RejectedExecutionException.class, () -> 
ste.executeOrRun(() -> {
+        }));
+    }
+
+    @Test
+    public void testExecuteOrRunInlineFailureIsIsolated() throws Exception {
+        @Cleanup("shutdown")
+        SingleThreadExecutor ste = new SingleThreadExecutor(THREAD_FACTORY);
+
+        CompletableFuture<Boolean> submitterCompleted = new 
CompletableFuture<>();
+        ste.execute(() -> {
+            ste.executeOrRun(() -> {
+                throw new RuntimeException("test");
+            });
+            submitterCompleted.complete(true);
+        });
+
+        assertTrue(submitterCompleted.get(10, TimeUnit.SECONDS));
+        Awaitility.await().until(() -> ste.getCompletedTasksCount() == 1);
+        assertEquals(1, ste.getFailedTasksCount());
+        assertEquals(2, ste.getSubmittedTasksCount());
+    }
 }
diff --git 
a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/LedgerHandle.java
 
b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/LedgerHandle.java
index a4699a15b7..71116dd3d1 100644
--- 
a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/LedgerHandle.java
+++ 
b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/LedgerHandle.java
@@ -51,7 +51,6 @@ import java.util.Set;
 import java.util.concurrent.Callable;
 import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.ConcurrentLinkedQueue;
-import java.util.concurrent.ExecutorService;
 import java.util.concurrent.RejectedExecutionException;
 import java.util.concurrent.ScheduledFuture;
 import java.util.concurrent.TimeUnit;
@@ -82,6 +81,7 @@ import org.apache.bookkeeper.client.impl.LedgerEntryImpl;
 import org.apache.bookkeeper.common.concurrent.FutureEventListener;
 import org.apache.bookkeeper.common.concurrent.FutureUtils;
 import org.apache.bookkeeper.common.util.MathUtils;
+import org.apache.bookkeeper.common.util.ThreadBoundExecutor;
 import org.apache.bookkeeper.net.BookieId;
 import org.apache.bookkeeper.proto.BookieProtocol;
 import org.apache.bookkeeper.proto.checksum.DigestManager;
@@ -106,7 +106,7 @@ public class LedgerHandle implements WriteHandle {
     final byte[] ledgerKey;
     private Versioned<LedgerMetadata> versionedMetadata;
     final long ledgerId;
-    final ExecutorService executor;
+    final ThreadBoundExecutor executor;
     long lastAddPushed;
     boolean notSupportBatch;
 
@@ -233,9 +233,11 @@ public class LedgerHandle implements WriteHandle {
         // Two calls on purpose: chooseThread(long) hashes the raw id while 
chooseThread(Object) goes through
         // hashCode(), and Long.hashCode folds the high bits. Boxing the id 
would move ledgers with ids >= 2^31
         // to a different thread than the other ledger-id keyed dispatches 
(e.g. OrderedGenericCallback).
-        this.executor = orderingKey == null
+        // The main worker pool is an OrderedExecutor, whose threads implement 
ThreadBoundExecutor whether or not
+        // they are decorated for task tracing or MDC preservation.
+        this.executor = (ThreadBoundExecutor) (orderingKey == null
                 ? clientCtx.getMainWorkerPool().chooseThread(ledgerId)
-                : clientCtx.getMainWorkerPool().chooseThread(orderingKey);
+                : clientCtx.getMainWorkerPool().chooseThread(orderingKey));
 
         if (clientCtx.getConf().enableStickyReads
                 && getLedgerMetadata().getEnsembleSize() == 
getLedgerMetadata().getWriteQuorumSize()) {
@@ -1119,8 +1121,9 @@ public class LedgerHandle implements WriteHandle {
             }
 
             if (isHandleWritable()) {
-                // Ledger handle in read/write mode: submit to OSE for ordered 
execution.
-                executeOrdered(op);
+                // Ledger handle in read/write mode: submit to OSE for ordered 
execution, unless the
+                // caller is already on the ledger's thread.
+                executor.executeOrRun(op);
             } else {
                 // Read-only ledger handle: bypass OSE and execute read 
directly in client thread.
                 // This avoids a context-switch to OSE thread and thus reduces 
latency.
@@ -1299,8 +1302,9 @@ public class LedgerHandle implements WriteHandle {
             }
 
             if (isHandleWritable()) {
-                // Ledger handle in read/write mode: submit to OSE for ordered 
execution.
-                executeOrdered(op);
+                // Ledger handle in read/write mode: submit to OSE for ordered 
execution, unless the
+                // caller is already on the ledger's thread.
+                executor.executeOrRun(op);
             } else {
                 // Read-only ledger handle: bypass OSE and execute read 
directly in client thread.
                 // This avoids a context-switch to OSE thread and thus reduces 
latency.
diff --git 
a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/LedgerHandleInlineReadTest.java
 
b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/LedgerHandleInlineReadTest.java
new file mode 100644
index 0000000000..a6b76b0136
--- /dev/null
+++ 
b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/LedgerHandleInlineReadTest.java
@@ -0,0 +1,148 @@
+/*
+ *
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ *
+ */
+package org.apache.bookkeeper.client;
+
+import static java.nio.charset.StandardCharsets.UTF_8;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertSame;
+import static org.junit.Assert.assertTrue;
+
+import com.google.common.collect.Lists;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
+import java.util.function.Consumer;
+import org.apache.bookkeeper.client.api.LedgerEntries;
+import org.apache.bookkeeper.client.api.LedgerMetadata;
+import org.apache.bookkeeper.client.api.WriteFlag;
+import org.apache.bookkeeper.common.concurrent.FutureUtils;
+import org.apache.bookkeeper.common.util.OrderedExecutor;
+import org.apache.bookkeeper.net.BookieId;
+import org.apache.bookkeeper.net.BookieSocketAddress;
+import org.apache.bookkeeper.proto.MockBookieClient;
+import org.apache.bookkeeper.proto.MockBookies;
+import org.apache.bookkeeper.versioning.Versioned;
+import org.junit.Test;
+
+/**
+ * Reads issued from the ledger's own worker thread are initiated inline 
instead of being queued on it.
+ */
+public class LedgerHandleInlineReadTest {
+
+    private static final BookieId b1 = new BookieSocketAddress("b1", 
3181).toBookieId();
+    private static final BookieId b2 = new BookieSocketAddress("b2", 
3181).toBookieId();
+    private static final BookieId b3 = new BookieSocketAddress("b3", 
3181).toBookieId();
+
+    private MockClientContext clientCtx;
+    private LedgerHandle lh;
+    private final AtomicReference<Thread> readIssuedOn = new 
AtomicReference<>();
+
+    private void setup(boolean decoratedThreads) throws Exception {
+        if (decoratedThreads) {
+            // Task tracing wraps the pool's threads in decorators; the handle 
must see through them
+            OrderedExecutor pool = 
OrderedExecutor.newBuilder().name("inline-read-test").numThreads(1)
+                    .traceTaskExecution(true).build();
+            MockBookies mockBookies = new MockBookies();
+            clientCtx = MockClientContext.create(mockBookies)
+                    .setMainWorkerPool(pool)
+                    .setBookieClient(new MockBookieClient(pool, mockBookies));
+        } else {
+            clientCtx = MockClientContext.create();
+        }
+        Versioned<LedgerMetadata> md = ClientUtil.setupLedger(clientCtx, 10L,
+                LedgerMetadataBuilder.create().newEnsembleEntry(0L, 
Lists.newArrayList(b1, b2, b3)));
+        lh = new LedgerHandle(clientCtx, 10L, md, 
BookKeeper.DigestType.CRC32C, ClientUtil.PASSWD,
+                WriteFlag.NONE);
+        lh.append("entry".getBytes(UTF_8));
+
+        // The hook runs synchronously inside the bookie client's readEntry, 
so it records the thread
+        // that initiates the read request.
+        clientCtx.getMockBookieClient().setPreReadHook((bookie, ledgerId, 
entryId) -> {
+            readIssuedOn.compareAndSet(null, Thread.currentThread());
+            return FutureUtils.value(null);
+        });
+    }
+
+    @Test(timeout = 30000)
+    public void testReadFromLedgerThreadIsInitiatedInline() throws Exception {
+        setup(false);
+        assertReadsFromLedgerThreadAreInitiatedInline();
+    }
+
+    @Test(timeout = 30000)
+    public void 
testReadFromLedgerThreadIsInitiatedInlineWithDecoratedThreads() throws 
Exception {
+        setup(true);
+        assertReadsFromLedgerThreadAreInitiatedInline();
+    }
+
+    private void assertReadsFromLedgerThreadAreInitiatedInline() throws 
Exception {
+        assertTrue(issuedBeforeReturnOnLedgerThread(result -> {
+            lh.readAsync(0, 0).whenComplete((entries, ex) -> complete(result, 
entries, ex));
+        }));
+
+        readIssuedOn.set(null);
+        assertTrue(issuedBeforeReturnOnLedgerThread(result -> {
+            lh.asyncReadEntries(0, 0, (rc, handle, entries, ctx) -> {
+                if (rc == BKException.Code.OK) {
+                    result.complete(null);
+                } else {
+                    result.completeExceptionally(BKException.create(rc));
+                }
+            }, null);
+        }));
+    }
+
+    @Test(timeout = 30000)
+    public void testReadFromOtherThreadIsQueuedOnLedgerThread() throws 
Exception {
+        setup(false);
+        try (LedgerEntries entries = lh.readAsync(0, 0).get(10, 
TimeUnit.SECONDS)) {
+            assertFalse(Thread.currentThread() == readIssuedOn.get());
+        }
+        CompletableFuture<Thread> ledgerThread = new CompletableFuture<>();
+        lh.executor.execute(() -> 
ledgerThread.complete(Thread.currentThread()));
+        assertSame(ledgerThread.get(10, TimeUnit.SECONDS), readIssuedOn.get());
+    }
+
+    /**
+     * Runs {@code read} on the ledger thread and returns whether the read 
request had reached the bookie
+     * client, on that same thread, by the time the call returned. Also waits 
for the read to complete.
+     */
+    private boolean 
issuedBeforeReturnOnLedgerThread(Consumer<CompletableFuture<Void>> read) throws 
Exception {
+        CompletableFuture<Void> completed = new CompletableFuture<>();
+        CompletableFuture<Boolean> issuedBeforeReturn = new 
CompletableFuture<>();
+        lh.executor.execute(() -> {
+            read.accept(completed);
+            issuedBeforeReturn.complete(readIssuedOn.get() == 
Thread.currentThread());
+        });
+        boolean inline = issuedBeforeReturn.get(10, TimeUnit.SECONDS);
+        completed.get(10, TimeUnit.SECONDS);
+        return inline;
+    }
+
+    private static void complete(CompletableFuture<Void> result, LedgerEntries 
entries, Throwable ex) {
+        if (ex != null) {
+            result.completeExceptionally(ex);
+        } else {
+            entries.close();
+            result.complete(null);
+        }
+    }
+}
diff --git 
a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/MockBookKeeperTestCase.java
 
b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/MockBookKeeperTestCase.java
index ffa59835ba..7a3ff584ba 100644
--- 
a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/MockBookKeeperTestCase.java
+++ 
b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/MockBookKeeperTestCase.java
@@ -210,7 +210,7 @@ public abstract class MockBookKeeperTestCase {
 
                 @Override
                 public OrderedExecutor getMainWorkerPool() {
-                    return scheduler;
+                    return executor;
                 }
 
                 @Override
diff --git 
a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/MockClientContext.java
 
b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/MockClientContext.java
index 93078a0512..8f0e9365d8 100644
--- 
a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/MockClientContext.java
+++ 
b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/MockClientContext.java
@@ -56,7 +56,9 @@ public class MockClientContext implements ClientContext {
 
     static MockClientContext create(MockBookies mockBookies) throws Exception {
         ClientConfiguration conf = new ClientConfiguration();
-        OrderedScheduler scheduler = 
OrderedScheduler.newSchedulerBuilder().name("mock-executor").numThreads(1).build();
+        OrderedExecutor executor = 
OrderedExecutor.newBuilder().name("mock-executor").numThreads(1).build();
+        OrderedScheduler scheduler = OrderedScheduler.newSchedulerBuilder()
+                .name("mock-scheduler").numThreads(1).build();
         MockRegistrationClient regClient = new MockRegistrationClient();
         EnsemblePlacementPolicy placementPolicy = new 
DefaultEnsemblePlacementPolicy();
         BookieWatcherImpl bookieWatcherImpl = new BookieWatcherImpl(conf, 
placementPolicy,
@@ -71,9 +73,9 @@ public class MockClientContext implements ClientContext {
                 .setBookieWatcher(bookieWatcherImpl)
                 .setPlacementPolicy(placementPolicy)
                 .setRegistrationClient(regClient)
-                .setBookieClient(new MockBookieClient(scheduler, mockBookies))
+                .setBookieClient(new MockBookieClient(executor, mockBookies))
                 .setByteBufAllocator(UnpooledByteBufAllocator.DEFAULT)
-                .setMainWorkerPool(scheduler)
+                .setMainWorkerPool(executor)
                 .setScheduler(scheduler)
                 
.setClientStats(BookKeeperClientStats.newInstance(NullStatsLogger.INSTANCE))
                 .setIsClientClosed(() -> false);

Reply via email to