From b4ccbabc5669394be03292f6aed6b31510c5d985 Mon Sep 17 00:00:00 2001 From: Dongnyoung Date: Fri, 4 Sep 2026 23:41:09 +0900 Subject: [PATCH 1/3] test: reproduce startup idle tracking race --- .../EventLoopSchedulerStartupRaceTest.java | 97 +++++++++++++++++++ 1 file changed, 97 insertions(+) create mode 100644 bootstrap/src/test/java/io/netty/loom/scheduler/EventLoopSchedulerStartupRaceTest.java diff --git a/bootstrap/src/test/java/io/netty/loom/scheduler/EventLoopSchedulerStartupRaceTest.java b/bootstrap/src/test/java/io/netty/loom/scheduler/EventLoopSchedulerStartupRaceTest.java new file mode 100644 index 0000000..06b1886 --- /dev/null +++ b/bootstrap/src/test/java/io/netty/loom/scheduler/EventLoopSchedulerStartupRaceTest.java @@ -0,0 +1,97 @@ +/* + * Copyright 2026 The Netty VirtualThread Scheduler Project + * + * The Netty VirtualThread Scheduler Project 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: + * + * https://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 io.netty.loom.scheduler; + +import static org.junit.jupiter.api.Assertions.*; + +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.locks.LockSupport; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; + +class EventLoopSchedulerStartupRaceTest { + + @Test + @Timeout(10) + void parkedCarrierStartedBeforeClusterStateMustBeDiscoverableAfterClusterStateIsConnected() throws Exception { + var carrierAtStartupHook = new CountDownLatch(1); + var releaseCarrier = new CountDownLatch(1); + var carrierLeftStartupHook = new CountDownLatch(1); + var victimAtStartupHook = new CountDownLatch(1); + var releaseVictim = new CountDownLatch(1); + + var scheduler = new EventLoopScheduler(0, Thread.ofPlatform().daemon(true).factory(), 16, null, () -> { + carrierAtStartupHook.countDown(); + await(releaseCarrier); + carrierLeftStartupHook.countDown(); + }); + var victim = new EventLoopScheduler(1, Thread.ofPlatform().daemon(true).factory(), 16, null, () -> { + victimAtStartupHook.countDown(); + await(releaseVictim); + }); + + ClusterState clusterState = null; + boolean searchStarted = false; + try { + assertTrue(carrierAtStartupHook.await(5, TimeUnit.SECONDS), "carrier did not reach startup hook"); + assertTrue(victimAtStartupHook.await(5, TimeUnit.SECONDS), "victim did not reach startup hook"); + assertNull(scheduler.clusterState, "test must release the carrier before ClusterState is connected"); + + releaseCarrier.countDown(); + assertTrue(carrierLeftStartupHook.await(5, TimeUnit.SECONDS), "carrier did not leave startup hook"); + awaitCarrierWaiting(scheduler); + + clusterState = new ClusterState(new EventLoopScheduler[]{scheduler, victim}); + scheduler.clusterState = clusterState; + victim.clusterState = clusterState; + + searchStarted = clusterState.tryStartSearcher(); + assertTrue(searchStarted, "test setup should be able to start one searcher"); + assertTrue(clusterState.wakeFirstIdle(victim), + "carrier is actually PARKED and connected to ClusterState, so it must be discoverable as idle"); + } finally { + if (searchStarted && clusterState != null && clusterState.nSearching() > 0) { + clusterState.stoppedSearching(); + } + releaseCarrier.countDown(); + releaseVictim.countDown(); + scheduler.wakeup(); + victim.wakeup(); + } + } + + private static void await(CountDownLatch latch) { + try { + latch.await(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new AssertionError(e); + } + } + + private static void awaitCarrierWaiting(EventLoopScheduler scheduler) { + var carrier = scheduler.carrierThread(); + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(5); + while (System.nanoTime() < deadline) { + if (carrier.getState() == Thread.State.WAITING && LockSupport.getBlocker(carrier) == null) { + return; + } + Thread.onSpinWait(); + } + fail("carrier did not enter WAITING state after startup; actual state=" + carrier.getState()); + } +} From c8f03adbdc7e568b4950aac2645a081e070ff670 Mon Sep 17 00:00:00 2001 From: Dongnyoung Date: Fri, 4 Sep 2026 23:50:55 +0900 Subject: [PATCH 2/3] Fix startup idle tracking race in EventLoopScheduler --- .../loom/scheduler/EventLoopScheduler.java | 3 ++ .../scheduler/EventLoopSchedulerGroup.java | 3 ++ .../EventLoopSchedulerStartupRaceTest.java | 32 ++++++++----------- 3 files changed, 20 insertions(+), 18 deletions(-) diff --git a/bootstrap/src/main/java/io/netty/loom/scheduler/EventLoopScheduler.java b/bootstrap/src/main/java/io/netty/loom/scheduler/EventLoopScheduler.java index d959b76..97aa92e 100644 --- a/bootstrap/src/main/java/io/netty/loom/scheduler/EventLoopScheduler.java +++ b/bootstrap/src/main/java/io/netty/loom/scheduler/EventLoopScheduler.java @@ -233,6 +233,9 @@ static SchedulingContext currentThreadSchedulerContext() { carrierThread = threadFactory.newThread(this::virtualThreadSchedulerLoop); carrierThread.setName("carrier-" + id); carrierThread.setDaemon(true); + } + + void startCarrier() { carrierThread.start(); } diff --git a/bootstrap/src/main/java/io/netty/loom/scheduler/EventLoopSchedulerGroup.java b/bootstrap/src/main/java/io/netty/loom/scheduler/EventLoopSchedulerGroup.java index d84317c..3263a00 100644 --- a/bootstrap/src/main/java/io/netty/loom/scheduler/EventLoopSchedulerGroup.java +++ b/bootstrap/src/main/java/io/netty/loom/scheduler/EventLoopSchedulerGroup.java @@ -105,6 +105,9 @@ private EventLoopSchedulerGroup(int size, NettyScheduler scheduler) { } } } + for (int i = 0; i < size; i++) { + schedulers[i].startCarrier(); + } } private static boolean isAllowedPeer(int self, int other, CarrierTopology.StealScope scope, diff --git a/bootstrap/src/test/java/io/netty/loom/scheduler/EventLoopSchedulerStartupRaceTest.java b/bootstrap/src/test/java/io/netty/loom/scheduler/EventLoopSchedulerStartupRaceTest.java index 06b1886..4303076 100644 --- a/bootstrap/src/test/java/io/netty/loom/scheduler/EventLoopSchedulerStartupRaceTest.java +++ b/bootstrap/src/test/java/io/netty/loom/scheduler/EventLoopSchedulerStartupRaceTest.java @@ -27,50 +27,46 @@ class EventLoopSchedulerStartupRaceTest { @Test @Timeout(10) - void parkedCarrierStartedBeforeClusterStateMustBeDiscoverableAfterClusterStateIsConnected() throws Exception { + void carrierStartedAfterClusterStateIsConnectedMustBeDiscoverableAsIdle() throws Exception { + var supportGroup = new EventLoopSchedulerGroup(null); var carrierAtStartupHook = new CountDownLatch(1); var releaseCarrier = new CountDownLatch(1); var carrierLeftStartupHook = new CountDownLatch(1); - var victimAtStartupHook = new CountDownLatch(1); - var releaseVictim = new CountDownLatch(1); var scheduler = new EventLoopScheduler(0, Thread.ofPlatform().daemon(true).factory(), 16, null, () -> { carrierAtStartupHook.countDown(); await(releaseCarrier); carrierLeftStartupHook.countDown(); }); - var victim = new EventLoopScheduler(1, Thread.ofPlatform().daemon(true).factory(), 16, null, () -> { - victimAtStartupHook.countDown(); - await(releaseVictim); - }); ClusterState clusterState = null; boolean searchStarted = false; + boolean searchWakeIssued = false; try { - assertTrue(carrierAtStartupHook.await(5, TimeUnit.SECONDS), "carrier did not reach startup hook"); - assertTrue(victimAtStartupHook.await(5, TimeUnit.SECONDS), "victim did not reach startup hook"); - assertNull(scheduler.clusterState, "test must release the carrier before ClusterState is connected"); + assertEquals(Thread.State.NEW, scheduler.carrierThread().getState(), + "carrier must not start before ClusterState is connected"); + + clusterState = new ClusterState(new EventLoopScheduler[]{scheduler}); + scheduler.group = supportGroup; + scheduler.clusterState = clusterState; + scheduler.startCarrier(); + assertTrue(carrierAtStartupHook.await(5, TimeUnit.SECONDS), "carrier did not reach startup hook"); releaseCarrier.countDown(); assertTrue(carrierLeftStartupHook.await(5, TimeUnit.SECONDS), "carrier did not leave startup hook"); awaitCarrierWaiting(scheduler); - clusterState = new ClusterState(new EventLoopScheduler[]{scheduler, victim}); - scheduler.clusterState = clusterState; - victim.clusterState = clusterState; - searchStarted = clusterState.tryStartSearcher(); assertTrue(searchStarted, "test setup should be able to start one searcher"); - assertTrue(clusterState.wakeFirstIdle(victim), + searchWakeIssued = clusterState.wakeFirstIdle(supportGroup.scheduler(0)); + assertTrue(searchWakeIssued, "carrier is actually PARKED and connected to ClusterState, so it must be discoverable as idle"); } finally { - if (searchStarted && clusterState != null && clusterState.nSearching() > 0) { + if (searchStarted && !searchWakeIssued && clusterState != null && clusterState.nSearching() > 0) { clusterState.stoppedSearching(); } releaseCarrier.countDown(); - releaseVictim.countDown(); scheduler.wakeup(); - victim.wakeup(); } } From cdcaa4dc931a1c9f80091e0a02ff29167ff4be8a Mon Sep 17 00:00:00 2001 From: Dongnyoung Date: Sat, 5 Sep 2026 00:13:41 +0900 Subject: [PATCH 3/3] Refine startup idle tracking regression test --- .../EventLoopSchedulerStartupRaceTest.java | 14 ++++---------- 1 file changed, 4 insertions(+), 10 deletions(-) diff --git a/bootstrap/src/test/java/io/netty/loom/scheduler/EventLoopSchedulerStartupRaceTest.java b/bootstrap/src/test/java/io/netty/loom/scheduler/EventLoopSchedulerStartupRaceTest.java index 4303076..17a6971 100644 --- a/bootstrap/src/test/java/io/netty/loom/scheduler/EventLoopSchedulerStartupRaceTest.java +++ b/bootstrap/src/test/java/io/netty/loom/scheduler/EventLoopSchedulerStartupRaceTest.java @@ -27,25 +27,20 @@ class EventLoopSchedulerStartupRaceTest { @Test @Timeout(10) - void carrierStartedAfterClusterStateIsConnectedMustBeDiscoverableAsIdle() throws Exception { + void parkedCarrierConnectedToClusterStateMustBeDiscoverableAsIdle() throws Exception { var supportGroup = new EventLoopSchedulerGroup(null); var carrierAtStartupHook = new CountDownLatch(1); var releaseCarrier = new CountDownLatch(1); - var carrierLeftStartupHook = new CountDownLatch(1); var scheduler = new EventLoopScheduler(0, Thread.ofPlatform().daemon(true).factory(), 16, null, () -> { carrierAtStartupHook.countDown(); await(releaseCarrier); - carrierLeftStartupHook.countDown(); }); ClusterState clusterState = null; boolean searchStarted = false; boolean searchWakeIssued = false; try { - assertEquals(Thread.State.NEW, scheduler.carrierThread().getState(), - "carrier must not start before ClusterState is connected"); - clusterState = new ClusterState(new EventLoopScheduler[]{scheduler}); scheduler.group = supportGroup; scheduler.clusterState = clusterState; @@ -53,8 +48,7 @@ void carrierStartedAfterClusterStateIsConnectedMustBeDiscoverableAsIdle() throws scheduler.startCarrier(); assertTrue(carrierAtStartupHook.await(5, TimeUnit.SECONDS), "carrier did not reach startup hook"); releaseCarrier.countDown(); - assertTrue(carrierLeftStartupHook.await(5, TimeUnit.SECONDS), "carrier did not leave startup hook"); - awaitCarrierWaiting(scheduler); + awaitCarrierParked(scheduler); searchStarted = clusterState.tryStartSearcher(); assertTrue(searchStarted, "test setup should be able to start one searcher"); @@ -79,7 +73,7 @@ private static void await(CountDownLatch latch) { } } - private static void awaitCarrierWaiting(EventLoopScheduler scheduler) { + private static void awaitCarrierParked(EventLoopScheduler scheduler) { var carrier = scheduler.carrierThread(); long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(5); while (System.nanoTime() < deadline) { @@ -88,6 +82,6 @@ private static void awaitCarrierWaiting(EventLoopScheduler scheduler) { } Thread.onSpinWait(); } - fail("carrier did not enter WAITING state after startup; actual state=" + carrier.getState()); + fail("carrier did not park after startup; actual state=" + carrier.getState()); } }