1919import org .slf4j .Logger ;
2020import org .slf4j .LoggerFactory ;
2121
22+ import java .lang .invoke .MethodHandle ;
23+ import java .lang .invoke .MethodHandles ;
24+ import java .lang .invoke .MethodType ;
2225import java .lang .ref .WeakReference ;
2326import java .util .ArrayList ;
2427import java .util .Collections ;
6164public class ConcurrentBag <T extends IConcurrentBagEntry > implements AutoCloseable
6265{
6366 private static final Logger LOGGER = LoggerFactory .getLogger (ConcurrentBag .class );
67+ private static final MethodHandle THREAD_IS_VIRTUAL = resolveThreadIsVirtual ();
6468
6569 private final CopyOnWriteArrayList <T > sharedList ;
6670 private final boolean useWeakThreadLocals ;
@@ -72,6 +76,31 @@ public class ConcurrentBag<T extends IConcurrentBagEntry> implements AutoCloseab
7276
7377 private final SynchronousQueue <T > handoffQueue ;
7478
79+ private static MethodHandle resolveThreadIsVirtual ()
80+ {
81+ try {
82+ return MethodHandles .publicLookup ()
83+ .findVirtual (Thread .class , "isVirtual" , MethodType .methodType (boolean .class ));
84+ }
85+ catch (NoSuchMethodException | IllegalAccessException e ) {
86+ return null ;
87+ }
88+ }
89+
90+ static boolean isCurrentThreadVirtual ()
91+ {
92+ if (THREAD_IS_VIRTUAL == null ) {
93+ return false ;
94+ }
95+
96+ try {
97+ return (boolean ) THREAD_IS_VIRTUAL .invokeExact (Thread .currentThread ());
98+ }
99+ catch (Throwable t ) {
100+ return false ;
101+ }
102+ }
103+
75104 /**
76105 * This interface defines the contract for an entry in the ConcurrentBag.
77106 * It provides methods to manage the state of the entry, which can be
@@ -188,11 +217,12 @@ public void requite(final T bagEntry)
188217 {
189218 bagEntry .setState (STATE_NOT_IN_USE );
190219
220+ final var isVirtualThread = isCurrentThreadVirtual ();
191221 for (int i = 1 , waiting = waiters .get (); waiting > 0 ; i ++, waiting = waiters .get ()) {
192222 if (bagEntry .getState () != STATE_NOT_IN_USE || handoffQueue .offer (bagEntry )) {
193223 return ;
194224 }
195- else if ((i & 0xff ) == 0xff || (waiting > 1 && i % waiting == 0 )) {
225+ else if (isVirtualThread || (i & 0xff ) == 0xff || (waiting > 1 && i % waiting == 0 )) {
196226 parkNanos (MICROSECONDS .toNanos (10 ));
197227 }
198228 else {
@@ -319,11 +349,12 @@ public void unreserve(final T bagEntry)
319349 {
320350 if (bagEntry .compareAndSet (STATE_RESERVED , STATE_NOT_IN_USE )) {
321351 // spin until a thread takes it or none are waiting
352+ final var isVirtualThread = isCurrentThreadVirtual ();
322353 for (int i = 1 , waiting = waiters .get (); waiting > 0 ; i ++, waiting = waiters .get ()) {
323354 if (bagEntry .getState () != STATE_NOT_IN_USE || handoffQueue .offer (bagEntry )) {
324355 return ;
325356 }
326- else if ((i & 0xff ) == 0xff || (waiting > 1 && i % waiting == 0 )) {
357+ else if (isVirtualThread || (i & 0xff ) == 0xff || (waiting > 1 && i % waiting == 0 )) {
327358 parkNanos (MICROSECONDS .toNanos (10 ));
328359 }
329360 else {
0 commit comments