From 133ee5224c272278b9a1398648bc82ac7fe5b3cd Mon Sep 17 00:00:00 2001 From: Tom Martin Date: Wed, 8 Jul 2026 12:47:01 +0100 Subject: [PATCH 01/13] Create a single place for all chunky thread/executors to be registered Closes #1893 - Allows chunky to close them on shutdown without a System.exit() - Additionally provides an api to wait on all chunky threads to be joined. --- .../src/java/se/llbit/chunky/main/Chunky.java | 7 + .../chunky/renderer/DefaultRenderManager.java | 3 +- .../chunky/renderer/RenderWorkerPool.java | 3 +- .../scene/AsynchronousSceneManager.java | 3 +- .../se/llbit/chunky/renderer/scene/Scene.java | 3 +- .../src/java/se/llbit/chunky/ui/ChunkMap.java | 3 +- .../src/java/se/llbit/chunky/ui/ChunkyFx.java | 1 - .../ResourcePackChooserController.java | 3 +- .../ui/controller/SceneChooserController.java | 4 +- .../chunky/world/ChunkTopographyUpdater.java | 4 +- .../se/llbit/chunky/world/SkymapTexture.java | 3 +- .../world/region/RegionChangeWatcher.java | 3 +- .../chunky/world/region/RegionParser.java | 3 +- .../llbit/util/concurrent/ChunkyThread.java | 169 ++++++++++++++++++ 14 files changed, 200 insertions(+), 12 deletions(-) create mode 100644 chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java diff --git a/chunky/src/java/se/llbit/chunky/main/Chunky.java b/chunky/src/java/se/llbit/chunky/main/Chunky.java index 13a6965357..976a5038cd 100644 --- a/chunky/src/java/se/llbit/chunky/main/Chunky.java +++ b/chunky/src/java/se/llbit/chunky/main/Chunky.java @@ -48,6 +48,7 @@ import se.llbit.log.Log; import se.llbit.log.Receiver; import se.llbit.util.TaskTracker; +import se.llbit.util.concurrent.ChunkyThread; import java.io.File; import java.io.FileInputStream; @@ -241,6 +242,12 @@ public static void main(final String[] args) { exitCode = 2; } } + + ChunkyThread.interruptAndJoinAll(); + ForkJoinPool commonThreads = Chunky.getCommonThreads(); + commonThreads.shutdownNow(); // ForkJoinPool doesn't return any tasks that were awaiting execution (all canceled). + // We don't use the ForkJoinPool commonPool in chunky, but if we did there is no shutdown available so nothing changes. + if (exitCode != 0) { System.exit(exitCode); } diff --git a/chunky/src/java/se/llbit/chunky/renderer/DefaultRenderManager.java b/chunky/src/java/se/llbit/chunky/renderer/DefaultRenderManager.java index 9b091c2e43..6074c35775 100644 --- a/chunky/src/java/se/llbit/chunky/renderer/DefaultRenderManager.java +++ b/chunky/src/java/se/llbit/chunky/renderer/DefaultRenderManager.java @@ -28,6 +28,7 @@ import se.llbit.chunky.resources.BitmapImage; import se.llbit.log.Log; import se.llbit.math.ColorUtil; +import se.llbit.util.concurrent.ChunkyThread; import se.llbit.util.TaskTracker; import java.time.Duration; @@ -52,7 +53,7 @@ *

All available final renderers are stored in {@code renderers} and preview renderers * are stored in {@code previewRenderers}. */ -public class DefaultRenderManager extends Thread implements RenderManager { +public class DefaultRenderManager extends ChunkyThread implements RenderManager { /** * Map containing all the final render {@code Renderer}s. The renderer corresponding to * {@code getRendererName()} is used when a render is requested. diff --git a/chunky/src/java/se/llbit/chunky/renderer/RenderWorkerPool.java b/chunky/src/java/se/llbit/chunky/renderer/RenderWorkerPool.java index 51b129d5d8..80cef5ea46 100644 --- a/chunky/src/java/se/llbit/chunky/renderer/RenderWorkerPool.java +++ b/chunky/src/java/se/llbit/chunky/renderer/RenderWorkerPool.java @@ -18,6 +18,7 @@ package se.llbit.chunky.renderer; import se.llbit.log.Log; +import se.llbit.util.concurrent.ChunkyThread; import java.util.ArrayList; import java.util.Random; @@ -39,7 +40,7 @@ public interface Factory { RenderWorkerPool create(int threads, long seed); } - public static class RenderWorker extends Thread { + public static class RenderWorker extends ChunkyThread { private final RenderWorkerPool pool; public final Random random; diff --git a/chunky/src/java/se/llbit/chunky/renderer/scene/AsynchronousSceneManager.java b/chunky/src/java/se/llbit/chunky/renderer/scene/AsynchronousSceneManager.java index 919a4f4c2f..7d08780263 100644 --- a/chunky/src/java/se/llbit/chunky/renderer/scene/AsynchronousSceneManager.java +++ b/chunky/src/java/se/llbit/chunky/renderer/scene/AsynchronousSceneManager.java @@ -26,6 +26,7 @@ import se.llbit.chunky.world.RegionPosition; import se.llbit.chunky.world.World; import se.llbit.log.Log; +import se.llbit.util.concurrent.ChunkyThread; import se.llbit.util.TaskTracker; import java.io.File; @@ -41,7 +42,7 @@ * * @author Jesper Öqvist */ -public class AsynchronousSceneManager extends Thread implements SceneManager { +public class AsynchronousSceneManager extends ChunkyThread implements SceneManager { private final SynchronousSceneManager sceneManager; private final LinkedBlockingQueue taskQueue; diff --git a/chunky/src/java/se/llbit/chunky/renderer/scene/Scene.java b/chunky/src/java/se/llbit/chunky/renderer/scene/Scene.java index bc4bc6bf61..82d0f3f21a 100644 --- a/chunky/src/java/se/llbit/chunky/renderer/scene/Scene.java +++ b/chunky/src/java/se/llbit/chunky/renderer/scene/Scene.java @@ -62,6 +62,7 @@ import se.llbit.nbt.Tag; import se.llbit.util.*; import se.llbit.util.annotation.NotNull; +import se.llbit.util.concurrent.ChunkyThread; import se.llbit.util.io.PositionalInputStream; import se.llbit.util.io.ZipExport; import se.llbit.util.mojangapi.MinecraftProfile; @@ -858,7 +859,7 @@ public synchronized void loadChunks(TaskTracker taskTracker, World world, Map {}; - private ScheduledExecutorService executor = Executors.newSingleThreadScheduledExecutor(); + private final ScheduledExecutorService executor = ChunkyThread.addExecutorService(Executors::newSingleThreadScheduledExecutor); private boolean shouldDrawPlayers = true; diff --git a/chunky/src/java/se/llbit/chunky/ui/ChunkyFx.java b/chunky/src/java/se/llbit/chunky/ui/ChunkyFx.java index 93f6532acc..44a552b5f4 100644 --- a/chunky/src/java/se/llbit/chunky/ui/ChunkyFx.java +++ b/chunky/src/java/se/llbit/chunky/ui/ChunkyFx.java @@ -82,7 +82,6 @@ public void stop() throws Exception { if(mainStage != null) { PersistentSettings.setWindowPosition(new WindowPosition(mainStage)); } - System.exit(0); } public static void startChunkyUI(Chunky chunkyInstance) { diff --git a/chunky/src/java/se/llbit/chunky/ui/controller/ResourcePackChooserController.java b/chunky/src/java/se/llbit/chunky/ui/controller/ResourcePackChooserController.java index dcac327038..251fdff9d5 100644 --- a/chunky/src/java/se/llbit/chunky/ui/controller/ResourcePackChooserController.java +++ b/chunky/src/java/se/llbit/chunky/ui/controller/ResourcePackChooserController.java @@ -49,6 +49,7 @@ import se.llbit.json.JsonParser; import se.llbit.log.Log; import se.llbit.util.MinecraftText; +import se.llbit.util.concurrent.ChunkyThread; import java.awt.*; import java.io.File; @@ -451,7 +452,7 @@ public void populate( private static class PackListItem { - private final static Executor PACK_PARSER_EXECUTOR = Executors.newSingleThreadExecutor(); + private final static Executor PACK_PARSER_EXECUTOR = ChunkyThread.addExecutorService(Executors::newSingleThreadExecutor); private static PackListItem DEFAULT = null; diff --git a/chunky/src/java/se/llbit/chunky/ui/controller/SceneChooserController.java b/chunky/src/java/se/llbit/chunky/ui/controller/SceneChooserController.java index 534926aeac..5c507b0568 100644 --- a/chunky/src/java/se/llbit/chunky/ui/controller/SceneChooserController.java +++ b/chunky/src/java/se/llbit/chunky/ui/controller/SceneChooserController.java @@ -36,6 +36,7 @@ import se.llbit.json.JsonObject; import se.llbit.json.JsonParser; import se.llbit.log.Log; +import se.llbit.util.concurrent.ChunkyThread; import java.io.File; import java.io.FileInputStream; @@ -64,6 +65,8 @@ public class SceneChooserController implements Initializable { private Stage stage; + private static final Executor loadExecutor = ChunkyThread.addExecutorService(Executors::newSingleThreadExecutor); + private ChunkyFxController controller; private static final HashMap sceneListCache = new HashMap<>(); @@ -219,7 +222,6 @@ public void setStage(Stage stage) { private void populateSceneTable(File sceneDir) { this.sceneTbl.setPlaceholder(new Label("Loading scenes…")); - Executor loadExecutor = Executors.newSingleThreadExecutor(); loadExecutor.execute(() -> { List scenes = new ArrayList<>(); diff --git a/chunky/src/java/se/llbit/chunky/world/ChunkTopographyUpdater.java b/chunky/src/java/se/llbit/chunky/world/ChunkTopographyUpdater.java index 8c61d5b591..58ef6be998 100644 --- a/chunky/src/java/se/llbit/chunky/world/ChunkTopographyUpdater.java +++ b/chunky/src/java/se/llbit/chunky/world/ChunkTopographyUpdater.java @@ -16,6 +16,8 @@ */ package se.llbit.chunky.world; +import se.llbit.util.concurrent.ChunkyThread; + import java.util.HashSet; import java.util.Iterator; import java.util.Set; @@ -25,7 +27,7 @@ * * @author Jesper Öqvist (jesper@llbit.se) */ -public class ChunkTopographyUpdater extends Thread { +public class ChunkTopographyUpdater extends ChunkyThread { private final Set queue = new HashSet<>(); diff --git a/chunky/src/java/se/llbit/chunky/world/SkymapTexture.java b/chunky/src/java/se/llbit/chunky/world/SkymapTexture.java index 4b2500a901..d879dc1576 100644 --- a/chunky/src/java/se/llbit/chunky/world/SkymapTexture.java +++ b/chunky/src/java/se/llbit/chunky/world/SkymapTexture.java @@ -25,6 +25,7 @@ import se.llbit.math.QuickMath; import se.llbit.math.Ray; import se.llbit.math.Vector4; +import se.llbit.util.concurrent.ChunkyThread; import se.llbit.util.ImageTools; /** @@ -35,7 +36,7 @@ */ public class SkymapTexture extends Texture { - class TexturePreprocessor extends Thread { + class TexturePreprocessor extends ChunkyThread { private final int x0; private final int x1; private final int y0; diff --git a/chunky/src/java/se/llbit/chunky/world/region/RegionChangeWatcher.java b/chunky/src/java/se/llbit/chunky/world/region/RegionChangeWatcher.java index 01ab9ed1b2..bdc35cd969 100644 --- a/chunky/src/java/se/llbit/chunky/world/region/RegionChangeWatcher.java +++ b/chunky/src/java/se/llbit/chunky/world/region/RegionChangeWatcher.java @@ -20,13 +20,14 @@ import se.llbit.chunky.map.WorldMapLoader; import se.llbit.chunky.renderer.ChunkViewListener; import se.llbit.chunky.world.ChunkView; +import se.llbit.util.concurrent.ChunkyThread; /** * Monitors filesystem for changes to region files. * * @author Jesper Öqvist */ -public abstract class RegionChangeWatcher extends Thread implements ChunkViewListener { +public abstract class RegionChangeWatcher extends ChunkyThread implements ChunkViewListener { protected final WorldMapLoader mapLoader; protected final MapView mapView; protected volatile ChunkView view = ChunkView.EMPTY; diff --git a/chunky/src/java/se/llbit/chunky/world/region/RegionParser.java b/chunky/src/java/se/llbit/chunky/world/region/RegionParser.java index b25f0771c0..841851a51e 100644 --- a/chunky/src/java/se/llbit/chunky/world/region/RegionParser.java +++ b/chunky/src/java/se/llbit/chunky/world/region/RegionParser.java @@ -23,6 +23,7 @@ import se.llbit.chunky.map.WorldMapLoader; import se.llbit.chunky.world.*; import se.llbit.log.Log; +import se.llbit.util.concurrent.ChunkyThread; import se.llbit.util.Mutable; /** @@ -34,7 +35,7 @@ * * @author Jesper Öqvist (jesper@llbit.se) */ -public class RegionParser extends Thread { +public class RegionParser extends ChunkyThread { private final WorldMapLoader mapLoader; private final RegionQueue queue; diff --git a/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java b/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java new file mode 100644 index 0000000000..1ba27c88cb --- /dev/null +++ b/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java @@ -0,0 +1,169 @@ +package se.llbit.util.concurrent; + +import se.llbit.chunky.main.Chunky; + +import java.util.ArrayList; +import java.util.Collection; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.ThreadFactory; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.function.Function; + +/** + * {@link Thread}/{@link ExecutorService} resource management. + *

The goal of this glass is to primarily:

+ * + * + *

Usage

+ *

{@link Thread Threads} in chunky should extend this class, and {@link ExecutorService executor services} should + * be created with {@link #addExecutorService(Function)} to allow chunky to interrupt and join them before chunky closes.

+ * + */ +public class ChunkyThread extends Thread { + /* + * All operations lock. + * When interruptAndJoinAll is called, additional threads/executors can't be added preventing later joins from + * waiting on threads that have not been interrupted. + */ + private static final AtomicBoolean isShutdown = new AtomicBoolean(false); + private static final Collection threads = new ArrayList<>(); + private static final Collection executorServices = new ArrayList<>(); + + /** + * Add a {@link Thread} to be interrupted and joined by chunky on shutdown + * + * @throws IllegalStateException When calling after {@link #interruptAndJoinAll()} has been called. + */ + public synchronized static T addThread(T thread) { + if (isShutdown.get()) { + throw new IllegalStateException("Creating a thread as chunky is stopping."); + } + threads.add(thread); + return thread; + } + + /** + * Add an {@link ExecutorService} to be interrupted and joined by chunky on shutdown + * + @throws IllegalStateException When calling after {@link #interruptAndJoinAll()} has been called. + */ + public synchronized static E addExecutorService(Function executorServiceSupplier) { + if (isShutdown.get()) { + throw new IllegalStateException("Creating an executor service as chunky is stopping."); + } + E e = executorServiceSupplier.apply(ChunkyThread::new); + executorServices.add(e); + return e; + } + + /** + * Await the joining of all threads managed by chunky + * + *

This method is always safe to call.

+ * + *

WARNING: calling this from any thread registered with {@link #addThread(Thread)} may deadlock.

+ */ + public synchronized static void joinAll() { + for (Thread thread : threads) { + try { + thread.join(); + } catch (InterruptedException e) { + // ignored + } + } + for (ExecutorService executorService : executorServices) { + try { + executorService.awaitTermination(1, TimeUnit.MINUTES); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } + } + + /** + * Interrupt and then await joining of all threads managed by chunky. + * + *

This method is always safe to call.

+ * + *

WARNING: calling this from any thread registered with {@link #addThread(Thread)} may deadlock.

+ *

Only to be called by {@link Chunky}

+ */ + public synchronized static void interruptAndJoinAll() { + assert Thread.currentThread().getName().equals("main"); + isShutdown.set(true); + + // shut down executorServices BEFORE threads because they recreate their threads when they are interrupted and stop + // causing an infinite hang. + for (ExecutorService executorService : executorServices) { + executorService.shutdownNow(); + } + + for (Thread thread : ChunkyThread.threads) { + thread.interrupt(); + } + for (Thread thread : ChunkyThread.threads) { + try { + thread.join(); + } catch (InterruptedException e) { + // ignored + } + } + } + + private void setDefaults() { + this.setDaemon(true); + addThread(this); + } + + /* Constructors from super */ + public ChunkyThread() { + super(); + setDefaults(); + } + + public ChunkyThread(Runnable task) { + super(task); + setDefaults(); + } + + public ChunkyThread(ThreadGroup group, Runnable task) { + super(group, task); + setDefaults(); + } + + public ChunkyThread(String name) { + super(name); + setDefaults(); + } + + public ChunkyThread(ThreadGroup group, String name) { + super(group, name); + setDefaults(); + } + + public ChunkyThread(Runnable task, String name) { + super(task, name); + setDefaults(); + } + + public ChunkyThread(ThreadGroup group, Runnable task, String name) { + super(group, task, name); + setDefaults(); + } + + public ChunkyThread(ThreadGroup group, Runnable task, String name, long stackSize) { + super(group, task, name, stackSize); + setDefaults(); + } + + public ChunkyThread(ThreadGroup group, Runnable task, String name, long stackSize, boolean inheritInheritableThreadLocals) { + super(group, task, name, stackSize, inheritInheritableThreadLocals); + setDefaults(); + } +} From 30f3aa033ff9aba03ff7126350ecd8c47f56d17b Mon Sep 17 00:00:00 2001 From: Tom Martin Date: Sat, 11 Jul 2026 09:51:07 +0100 Subject: [PATCH 02/13] Fix race between joinAll and extremely lately added threads/executors --- .../llbit/util/concurrent/ChunkyThread.java | 32 +++++++++++++++---- 1 file changed, 25 insertions(+), 7 deletions(-) diff --git a/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java b/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java index 1ba27c88cb..066efc9aae 100644 --- a/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java +++ b/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java @@ -1,13 +1,14 @@ package se.llbit.util.concurrent; import se.llbit.chunky.main.Chunky; +import se.llbit.log.Log; import java.util.ArrayList; import java.util.Collection; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutorService; import java.util.concurrent.ThreadFactory; import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicBoolean; import java.util.function.Function; /** @@ -31,7 +32,7 @@ public class ChunkyThread extends Thread { * When interruptAndJoinAll is called, additional threads/executors can't be added preventing later joins from * waiting on threads that have not been interrupted. */ - private static final AtomicBoolean isShutdown = new AtomicBoolean(false); + private static final CountDownLatch shutdownLatch = new CountDownLatch(1); private static final Collection threads = new ArrayList<>(); private static final Collection executorServices = new ArrayList<>(); @@ -41,7 +42,7 @@ public class ChunkyThread extends Thread { * @throws IllegalStateException When calling after {@link #interruptAndJoinAll()} has been called. */ public synchronized static T addThread(T thread) { - if (isShutdown.get()) { + if (shutdownLatch.getCount() == 0) { throw new IllegalStateException("Creating a thread as chunky is stopping."); } threads.add(thread); @@ -54,7 +55,7 @@ public synchronized static T addThread(T thread) { @throws IllegalStateException When calling after {@link #interruptAndJoinAll()} has been called. */ public synchronized static E addExecutorService(Function executorServiceSupplier) { - if (isShutdown.get()) { + if (shutdownLatch.getCount() == 0) { throw new IllegalStateException("Creating an executor service as chunky is stopping."); } E e = executorServiceSupplier.apply(ChunkyThread::new); @@ -70,20 +71,37 @@ public synchronized static E addExecutorService(Func *

WARNING: calling this from any thread registered with {@link #addThread(Thread)} may deadlock.

*/ public synchronized static void joinAll() { + boolean interrupted = false; + + while (true) { + try { + // must wait for the latch as hitting the for loop below first causes immediate evaluation of the + // for loop iterator, potentially missing new threads. + shutdownLatch.await(); + break; + } catch (InterruptedException e) { + interrupted = true; + } + } + for (Thread thread : threads) { try { thread.join(); } catch (InterruptedException e) { - // ignored + interrupted = true; } } for (ExecutorService executorService : executorServices) { try { executorService.awaitTermination(1, TimeUnit.MINUTES); } catch (InterruptedException e) { - Thread.currentThread().interrupt(); + interrupted = true; } } + + if (interrupted) { + Thread.currentThread().interrupt(); + } } /** @@ -96,7 +114,7 @@ public synchronized static void joinAll() { */ public synchronized static void interruptAndJoinAll() { assert Thread.currentThread().getName().equals("main"); - isShutdown.set(true); + shutdownLatch.countDown(); // shut down executorServices BEFORE threads because they recreate their threads when they are interrupted and stop // causing an infinite hang. From c675acea49500b77020865b442653b0094471c09 Mon Sep 17 00:00:00 2001 From: Tom Martin Date: Mon, 13 Jul 2026 08:03:18 +0100 Subject: [PATCH 03/13] Make joinAll not synchronized --- .../src/java/se/llbit/util/concurrent/ChunkyThread.java | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java b/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java index 066efc9aae..d10300097a 100644 --- a/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java +++ b/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java @@ -70,7 +70,13 @@ public synchronized static E addExecutorService(Func * *

WARNING: calling this from any thread registered with {@link #addThread(Thread)} may deadlock.

*/ - public synchronized static void joinAll() { + public static void joinAll() { + /* + * This method should not be synchronized because: + * 1. Calls to this method that happen before interruptAndJoinAll will lock the latter interrupting thread, deadlocking. + * 2. shutdownLatch.await is at least acquire memory ordering, and modification is disabled after the latch is zero. + * As such we are guaranteed that no threads can modify the state. + */ boolean interrupted = false; while (true) { From 1d166b1c5ce9e05959ff2d768892f6842fc300c7 Mon Sep 17 00:00:00 2001 From: Tom Martin Date: Mon, 13 Jul 2026 08:49:26 +0100 Subject: [PATCH 04/13] Return whether all threads were joined from joinAll methods --- .../src/java/se/llbit/chunky/main/Chunky.java | 3 +- .../llbit/util/concurrent/ChunkyThread.java | 116 ++++++++++-------- 2 files changed, 68 insertions(+), 51 deletions(-) diff --git a/chunky/src/java/se/llbit/chunky/main/Chunky.java b/chunky/src/java/se/llbit/chunky/main/Chunky.java index 976a5038cd..4888adcb9a 100644 --- a/chunky/src/java/se/llbit/chunky/main/Chunky.java +++ b/chunky/src/java/se/llbit/chunky/main/Chunky.java @@ -57,6 +57,7 @@ import java.nio.file.Path; import java.util.*; import java.util.concurrent.ForkJoinPool; +import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; /** @@ -243,7 +244,7 @@ public static void main(final String[] args) { } } - ChunkyThread.interruptAndJoinAll(); + ChunkyThread.interruptAndJoinAll(5, TimeUnit.SECONDS); ForkJoinPool commonThreads = Chunky.getCommonThreads(); commonThreads.shutdownNow(); // ForkJoinPool doesn't return any tasks that were awaiting execution (all canceled). // We don't use the ForkJoinPool commonPool in chunky, but if we did there is no shutdown available so nothing changes. diff --git a/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java b/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java index d10300097a..3284a7302d 100644 --- a/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java +++ b/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java @@ -1,19 +1,16 @@ package se.llbit.util.concurrent; import se.llbit.chunky.main.Chunky; -import se.llbit.log.Log; +import se.llbit.util.annotation.NotNull; import java.util.ArrayList; import java.util.Collection; -import java.util.concurrent.CountDownLatch; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.ThreadFactory; -import java.util.concurrent.TimeUnit; +import java.util.concurrent.*; import java.util.function.Function; /** * {@link Thread}/{@link ExecutorService} resource management. - *

The goal of this glass is to primarily:

+ *

The goal of this class is to primarily:

*
    *
  • Ensure all chunky threads (even daemons) are interrupted and given a chance clean up before getting killed by * the runtime on termination.
  • @@ -27,11 +24,6 @@ * */ public class ChunkyThread extends Thread { - /* - * All operations lock. - * When interruptAndJoinAll is called, additional threads/executors can't be added preventing later joins from - * waiting on threads that have not been interrupted. - */ private static final CountDownLatch shutdownLatch = new CountDownLatch(1); private static final Collection threads = new ArrayList<>(); private static final Collection executorServices = new ArrayList<>(); @@ -39,7 +31,7 @@ public class ChunkyThread extends Thread { /** * Add a {@link Thread} to be interrupted and joined by chunky on shutdown * - * @throws IllegalStateException When calling after {@link #interruptAndJoinAll()} has been called. + * @throws IllegalStateException When calling after {@link #interruptAndJoinAll(long, TimeUnit)} has been called. */ public synchronized static T addThread(T thread) { if (shutdownLatch.getCount() == 0) { @@ -52,7 +44,7 @@ public synchronized static T addThread(T thread) { /** * Add an {@link ExecutorService} to be interrupted and joined by chunky on shutdown * - @throws IllegalStateException When calling after {@link #interruptAndJoinAll()} has been called. + * @throws IllegalStateException When calling after {@link #interruptAndJoinAll(long, TimeUnit)} has been called. */ public synchronized static E addExecutorService(Function executorServiceSupplier) { if (shutdownLatch.getCount() == 0) { @@ -64,13 +56,18 @@ public synchronized static E addExecutorService(Func } /** - * Await the joining of all threads managed by chunky + * Await the joining of all threads managed by chunky. This method will wait indefinitely until a + * shutdown happens to begin its timeout. * *

    This method is always safe to call.

    * *

    WARNING: calling this from any thread registered with {@link #addThread(Thread)} may deadlock.

    + * + * @param timeout The maximum time to wait AFTER a shutdown is initiated + * @param unit the time unit of the timeout argument + * @return Whether all threads were joined before returning */ - public static void joinAll() { + public static boolean joinAll(long timeout, @NotNull TimeUnit unit) { /* * This method should not be synchronized because: * 1. Calls to this method that happen before interruptAndJoinAll will lock the latter interrupting thread, deadlocking. @@ -90,54 +87,73 @@ public static void joinAll() { } } - for (Thread thread : threads) { - try { - thread.join(); - } catch (InterruptedException e) { - interrupted = true; + long startTime = System.nanoTime(); + long endTime = startTime + unit.toNanos(timeout); + + boolean anyAlive = false; + + try { + for (ExecutorService executorService : executorServices) { + while (System.nanoTime() < endTime) { + try { + long waitTime = endTime - startTime; + if (waitTime > 0) { + executorService.awaitTermination(waitTime, TimeUnit.NANOSECONDS); + } + break; + } catch (InterruptedException e) { + interrupted = true; + } + } + anyAlive |= !executorService.isTerminated(); } - } - for (ExecutorService executorService : executorServices) { - try { - executorService.awaitTermination(1, TimeUnit.MINUTES); - } catch (InterruptedException e) { - interrupted = true; + for (Thread thread : ChunkyThread.threads) { + while (System.nanoTime() < endTime) { + try { + long waitTimeMillis = TimeUnit.NANOSECONDS.toMillis(endTime - startTime); + if (waitTimeMillis > 0) { + thread.join(waitTimeMillis); + } + break; + } catch (InterruptedException e) { + interrupted = true; + } + } + anyAlive |= thread.isAlive(); + } + } finally { + if (interrupted) { + Thread.currentThread().interrupt(); } } - - if (interrupted) { - Thread.currentThread().interrupt(); - } + return !anyAlive; } /** * Interrupt and then await joining of all threads managed by chunky. * - *

    This method is always safe to call.

    - * *

    WARNING: calling this from any thread registered with {@link #addThread(Thread)} may deadlock.

    *

    Only to be called by {@link Chunky}

    + * + * @param timeout The maximum time to wait + * @param unit the time unit of the timeout argument + * @return Whether all threads were joined before the time limit was reached */ - public synchronized static void interruptAndJoinAll() { - assert Thread.currentThread().getName().equals("main"); - shutdownLatch.countDown(); - - // shut down executorServices BEFORE threads because they recreate their threads when they are interrupted and stop - // causing an infinite hang. - for (ExecutorService executorService : executorServices) { - executorService.shutdownNow(); - } - - for (Thread thread : ChunkyThread.threads) { - thread.interrupt(); - } - for (Thread thread : ChunkyThread.threads) { - try { - thread.join(); - } catch (InterruptedException e) { - // ignored + public static boolean interruptAndJoinAll(long timeout, @NotNull TimeUnit unit) { + synchronized (ChunkyThread.class) { + shutdownLatch.countDown(); + + // shut down executorServices BEFORE threads because they recreate their threads when they are interrupted and stop + // causing an infinite hang. + for (ExecutorService executorService : executorServices) { + executorService.shutdownNow(); + } + for (Thread thread : ChunkyThread.threads) { + thread.interrupt(); } } + + return joinAll(timeout, unit); } private void setDefaults() { From 63367bec09dda00d0f8714e41d012b35ad433447 Mon Sep 17 00:00:00 2001 From: Tom Martin Date: Mon, 13 Jul 2026 08:55:51 +0100 Subject: [PATCH 05/13] Add Chunky's ForkJoinPool to ChunkyThread --- .../src/java/se/llbit/chunky/main/Chunky.java | 7 ++----- .../llbit/util/concurrent/ChunkyThread.java | 20 +++++++++++++++++-- 2 files changed, 20 insertions(+), 7 deletions(-) diff --git a/chunky/src/java/se/llbit/chunky/main/Chunky.java b/chunky/src/java/se/llbit/chunky/main/Chunky.java index 4888adcb9a..14e0925915 100644 --- a/chunky/src/java/se/llbit/chunky/main/Chunky.java +++ b/chunky/src/java/se/llbit/chunky/main/Chunky.java @@ -245,9 +245,6 @@ public static void main(final String[] args) { } ChunkyThread.interruptAndJoinAll(5, TimeUnit.SECONDS); - ForkJoinPool commonThreads = Chunky.getCommonThreads(); - commonThreads.shutdownNow(); // ForkJoinPool doesn't return any tasks that were awaiting execution (all canceled). - // We don't use the ForkJoinPool commonPool in chunky, but if we did there is no shutdown available so nothing changes. if (exitCode != 0) { System.exit(exitCode); @@ -350,7 +347,7 @@ public void update() { public static ForkJoinPool getCommonThreads() { if (commonThreads == null) { // use at least two threads to prevent deadlocks in some java versions (see #1631) - commonThreads = new ForkJoinPool(Math.max(PersistentSettings.getNumThreads(), 2)); + commonThreads = ChunkyThread.addForkJoinPool(new ForkJoinPool(Math.max(PersistentSettings.getNumThreads(), 2))); } return commonThreads; } @@ -361,7 +358,7 @@ public static ForkJoinPool getCommonThreads() { */ public static void setCommonThreadsCount(int threads) { ForkJoinPool t = getCommonThreads(); - commonThreads = new ForkJoinPool(threads); + commonThreads = ChunkyThread.addForkJoinPool(new ForkJoinPool(threads)); t.shutdown(); } diff --git a/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java b/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java index 3284a7302d..ea3a544ff7 100644 --- a/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java +++ b/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java @@ -19,8 +19,11 @@ *
* *

Usage

- *

{@link Thread Threads} in chunky should extend this class, and {@link ExecutorService executor services} should - * be created with {@link #addExecutorService(Function)} to allow chunky to interrupt and join them before chunky closes.

+ *
    + *
  • {@link Thread Threads} in chunky should extend this class
  • + *
  • {@link ExecutorService Executor Services} should be created with {@link #addExecutorService(Function)}
  • + *
  • {@link ForkJoinPool Fork Join Pools} should be created with {@link #addForkJoinPool(ForkJoinPool)}
  • + *
* */ public class ChunkyThread extends Thread { @@ -55,6 +58,19 @@ public synchronized static E addExecutorService(Func return e; } + /** + * Add a {@link ForkJoinPool} to be interrupted and joined by chunky on shutdown + * + * @throws IllegalStateException When calling after {@link #interruptAndJoinAll(long, TimeUnit)} has been called. + */ + public synchronized static ForkJoinPool addForkJoinPool(ForkJoinPool pool) { + if (shutdownLatch.getCount() == 0) { + throw new IllegalStateException("Creating a fork join pool as chunky is stopping."); + } + executorServices.add(pool); + return pool; + } + /** * Await the joining of all threads managed by chunky. This method will wait indefinitely until a * shutdown happens to begin its timeout. From f5d4d139357ed04cbbe9c30284ac00754ab128ae Mon Sep 17 00:00:00 2001 From: Tom Martin Date: Mon, 13 Jul 2026 08:57:10 +0100 Subject: [PATCH 06/13] Change ExecutorService thread factory --- .../src/java/se/llbit/util/concurrent/ChunkyThread.java | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java b/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java index ea3a544ff7..e86a080cf8 100644 --- a/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java +++ b/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java @@ -53,7 +53,11 @@ public synchronized static E addExecutorService(Func if (shutdownLatch.getCount() == 0) { throw new IllegalStateException("Creating an executor service as chunky is stopping."); } - E e = executorServiceSupplier.apply(ChunkyThread::new); + E e = executorServiceSupplier.apply(r -> { // executor shutdown interrupts its own threads, so they don't need to be ChunkyThreads + Thread t = new Thread(r); + t.setDaemon(true); + return t; + }); executorServices.add(e); return e; } @@ -159,8 +163,6 @@ public static boolean interruptAndJoinAll(long timeout, @NotNull TimeUnit unit) synchronized (ChunkyThread.class) { shutdownLatch.countDown(); - // shut down executorServices BEFORE threads because they recreate their threads when they are interrupted and stop - // causing an infinite hang. for (ExecutorService executorService : executorServices) { executorService.shutdownNow(); } From 7163545b093ea0fdc33281d10b2353c70e26083f Mon Sep 17 00:00:00 2001 From: Tom Martin Date: Mon, 13 Jul 2026 09:29:45 +0100 Subject: [PATCH 07/13] Add PluginApi to relevant ChunkyThread methods --- chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java b/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java index e86a080cf8..07c9bd8b66 100644 --- a/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java +++ b/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java @@ -1,6 +1,7 @@ package se.llbit.util.concurrent; import se.llbit.chunky.main.Chunky; +import se.llbit.chunky.plugin.PluginApi; import se.llbit.util.annotation.NotNull; import java.util.ArrayList; @@ -36,6 +37,7 @@ public class ChunkyThread extends Thread { * * @throws IllegalStateException When calling after {@link #interruptAndJoinAll(long, TimeUnit)} has been called. */ + @PluginApi public synchronized static T addThread(T thread) { if (shutdownLatch.getCount() == 0) { throw new IllegalStateException("Creating a thread as chunky is stopping."); @@ -49,6 +51,7 @@ public synchronized static T addThread(T thread) { * * @throws IllegalStateException When calling after {@link #interruptAndJoinAll(long, TimeUnit)} has been called. */ + @PluginApi public synchronized static E addExecutorService(Function executorServiceSupplier) { if (shutdownLatch.getCount() == 0) { throw new IllegalStateException("Creating an executor service as chunky is stopping."); @@ -67,6 +70,7 @@ public synchronized static E addExecutorService(Func * * @throws IllegalStateException When calling after {@link #interruptAndJoinAll(long, TimeUnit)} has been called. */ + @PluginApi public synchronized static ForkJoinPool addForkJoinPool(ForkJoinPool pool) { if (shutdownLatch.getCount() == 0) { throw new IllegalStateException("Creating a fork join pool as chunky is stopping."); @@ -87,6 +91,7 @@ public synchronized static ForkJoinPool addForkJoinPool(ForkJoinPool pool) { * @param unit the time unit of the timeout argument * @return Whether all threads were joined before returning */ + @PluginApi public static boolean joinAll(long timeout, @NotNull TimeUnit unit) { /* * This method should not be synchronized because: From 297b535163413d0aabe671c4305d46a7dcacb22a Mon Sep 17 00:00:00 2001 From: Tom Martin Date: Mon, 13 Jul 2026 09:32:28 +0100 Subject: [PATCH 08/13] Add shutdown hook to handle early System.exit() --- chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java b/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java index 07c9bd8b66..47a1161ed1 100644 --- a/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java +++ b/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java @@ -32,6 +32,13 @@ public class ChunkyThread extends Thread { private static final Collection threads = new ArrayList<>(); private static final Collection executorServices = new ArrayList<>(); + static { + // If anyone calls System.exit() we still want to attempt to stop all threads + Runtime.getRuntime().addShutdownHook( + new Thread(() -> ChunkyThread.interruptAndJoinAll(0, TimeUnit.SECONDS)) // intentionally not ChunkyThread + ); + } + /** * Add a {@link Thread} to be interrupted and joined by chunky on shutdown * From 7f0c6615338a470031dd33e34d3a42e912f0def8 Mon Sep 17 00:00:00 2001 From: Tom Martin Date: Mon, 13 Jul 2026 09:39:13 +0100 Subject: [PATCH 09/13] Put interruptAndJoinAll in a finally block --- .../src/java/se/llbit/chunky/main/Chunky.java | 50 ++++++++++--------- 1 file changed, 27 insertions(+), 23 deletions(-) diff --git a/chunky/src/java/se/llbit/chunky/main/Chunky.java b/chunky/src/java/se/llbit/chunky/main/Chunky.java index 14e0925915..0e09e96a18 100644 --- a/chunky/src/java/se/llbit/chunky/main/Chunky.java +++ b/chunky/src/java/se/llbit/chunky/main/Chunky.java @@ -217,35 +217,39 @@ public static void main(final String[] args) { if (cmdline.mode == CommandLineOptions.Mode.CLI_OPERATION) { exitCode = cmdline.exitCode; } else { - // Initialize the common thread pool. - getCommonThreads(); + try { + // Initialize the common thread pool. + getCommonThreads(); - Chunky chunky = new Chunky(cmdline.options); - chunky.headless = cmdline.mode == Mode.HEADLESS_RENDER || cmdline.mode == Mode.CREATE_SNAPSHOT; - chunky.loadPlugins(); + Chunky chunky = new Chunky(cmdline.options); + chunky.headless = cmdline.mode == Mode.HEADLESS_RENDER || cmdline.mode == Mode.CREATE_SNAPSHOT; + chunky.loadPlugins(); - try { - switch (cmdline.mode) { - case HEADLESS_RENDER: - exitCode = chunky.doHeadlessRender(); - break; - case CREATE_SNAPSHOT: - exitCode = chunky.doSnapshot(); - break; - case START_GUI: - ChunkyFx.startChunkyUI(chunky); - break; + try { + switch (cmdline.mode) { + case HEADLESS_RENDER: + exitCode = chunky.doHeadlessRender(); + break; + case CREATE_SNAPSHOT: + exitCode = chunky.doSnapshot(); + break; + case START_GUI: + ChunkyFx.startChunkyUI(chunky); + break; + } + } catch (Throwable t) { + // set receiver in case an exception was thrown before it was set in one of the start modes. + Log.setReceiver(ConsoleReceiver.INSTANCE, Level.INFO, Level.WARNING, Level.ERROR); + Log.error("Unchecked exception caused Chunky to close.", t); + exitCode = 2; + } + } finally { + if (!ChunkyThread.interruptAndJoinAll(5, TimeUnit.SECONDS)) { + Log.warn("Not all Chunky threads stopped before exiting."); } - } catch (Throwable t) { - // set receiver in case an exception was thrown before it was set in one of the start modes. - Log.setReceiver(ConsoleReceiver.INSTANCE, Level.INFO, Level.WARNING, Level.ERROR); - Log.error("Unchecked exception caused Chunky to close.", t); - exitCode = 2; } } - ChunkyThread.interruptAndJoinAll(5, TimeUnit.SECONDS); - if (exitCode != 0) { System.exit(exitCode); } From 22330c69b37437271a140b306f2fdee1421cdbea Mon Sep 17 00:00:00 2001 From: Tom Martin Date: Fri, 24 Jul 2026 20:39:34 +0100 Subject: [PATCH 10/13] Switch back to ConsoleReceiver on UI stop --- chunky/src/java/se/llbit/chunky/ui/ChunkyFx.java | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/chunky/src/java/se/llbit/chunky/ui/ChunkyFx.java b/chunky/src/java/se/llbit/chunky/ui/ChunkyFx.java index 44a552b5f4..6226d4c834 100644 --- a/chunky/src/java/se/llbit/chunky/ui/ChunkyFx.java +++ b/chunky/src/java/se/llbit/chunky/ui/ChunkyFx.java @@ -29,6 +29,8 @@ import se.llbit.chunky.resources.SettingsDirectory; import se.llbit.chunky.ui.controller.ChunkyFxController; import se.llbit.fxutil.WindowPosition; +import se.llbit.log.ConsoleReceiver; +import se.llbit.log.Level; import se.llbit.log.Log; import java.io.File; @@ -79,6 +81,8 @@ public class ChunkyFx extends Application { @Override public void stop() throws Exception { + // UI is closing so we need to replace the receivers + Log.setReceiver(ConsoleReceiver.INSTANCE, Level.INFO, Level.WARNING, Level.ERROR); if(mainStage != null) { PersistentSettings.setWindowPosition(new WindowPosition(mainStage)); } From afb094700438aaf1e89eb8a4e2957377badf2a53 Mon Sep 17 00:00:00 2001 From: Tom Martin Date: Fri, 24 Jul 2026 20:41:34 +0100 Subject: [PATCH 11/13] Move the shutdown hook out of ChunkyThread, now managed by Chunky itself --- .../src/java/se/llbit/chunky/main/Chunky.java | 69 ++++++++++--------- .../llbit/util/concurrent/ChunkyThread.java | 7 -- 2 files changed, 38 insertions(+), 38 deletions(-) diff --git a/chunky/src/java/se/llbit/chunky/main/Chunky.java b/chunky/src/java/se/llbit/chunky/main/Chunky.java index 0e09e96a18..4b69de74d3 100644 --- a/chunky/src/java/se/llbit/chunky/main/Chunky.java +++ b/chunky/src/java/se/llbit/chunky/main/Chunky.java @@ -213,45 +213,52 @@ public static void main(final String[] args) { System.exit(1); } - int exitCode = 0; if (cmdline.mode == CommandLineOptions.Mode.CLI_OPERATION) { - exitCode = cmdline.exitCode; + System.exit(cmdline.exitCode); } else { - try { - // Initialize the common thread pool. - getCommonThreads(); + // Initialize the common thread pool. + getCommonThreads(); - Chunky chunky = new Chunky(cmdline.options); - chunky.headless = cmdline.mode == Mode.HEADLESS_RENDER || cmdline.mode == Mode.CREATE_SNAPSHOT; - chunky.loadPlugins(); + Chunky chunky = new Chunky(cmdline.options); + chunky.headless = cmdline.mode == Mode.HEADLESS_RENDER || cmdline.mode == Mode.CREATE_SNAPSHOT; + chunky.loadPlugins(); - try { - switch (cmdline.mode) { - case HEADLESS_RENDER: - exitCode = chunky.doHeadlessRender(); - break; - case CREATE_SNAPSHOT: - exitCode = chunky.doSnapshot(); - break; - case START_GUI: - ChunkyFx.startChunkyUI(chunky); - break; - } - } catch (Throwable t) { - // set receiver in case an exception was thrown before it was set in one of the start modes. - Log.setReceiver(ConsoleReceiver.INSTANCE, Level.INFO, Level.WARNING, Level.ERROR); - Log.error("Unchecked exception caused Chunky to close.", t); - exitCode = 2; - } - } finally { - if (!ChunkyThread.interruptAndJoinAll(5, TimeUnit.SECONDS)) { - Log.warn("Not all Chunky threads stopped before exiting."); + Runtime.getRuntime().addShutdownHook( + new Thread(() -> { // intentionally not a ChunkyThread + // Within a shutdown hook we need to close quickly, otherwise we risk the user/OS escalating to KILL + chunky.shutdown(1, TimeUnit.SECONDS); + }) + ); + + int exitCode = 0; + try { + switch (cmdline.mode) { + case HEADLESS_RENDER: + exitCode = chunky.doHeadlessRender(); + break; + case CREATE_SNAPSHOT: + exitCode = chunky.doSnapshot(); + break; + case START_GUI: + ChunkyFx.startChunkyUI(chunky); + break; } + } catch (Throwable t) { + // set receiver in case an exception was thrown before it was set in one of the start modes. + Log.setReceiver(ConsoleReceiver.INSTANCE, Level.INFO, Level.WARNING, Level.ERROR); + Log.error("Unchecked exception caused Chunky to close.", t); + exitCode = 2; } + chunky.shutdown(5, TimeUnit.SECONDS); + // Always exit, we've done all shutdown necessary and want to exit whether non-daemon threads exist or not. + // This should prevent hangs if threads aren't cooperating. + System.exit(exitCode); } + } - if (exitCode != 0) { - System.exit(exitCode); + private void shutdown(int timeout, TimeUnit unit) { + if (!ChunkyThread.interruptAndJoinAll(timeout, unit)) { + Log.error("Not all threads were joined before shutting down."); // FIXME: list all alive threads? ThreadGroups are annoying. } } diff --git a/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java b/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java index 47a1161ed1..07c9bd8b66 100644 --- a/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java +++ b/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java @@ -32,13 +32,6 @@ public class ChunkyThread extends Thread { private static final Collection threads = new ArrayList<>(); private static final Collection executorServices = new ArrayList<>(); - static { - // If anyone calls System.exit() we still want to attempt to stop all threads - Runtime.getRuntime().addShutdownHook( - new Thread(() -> ChunkyThread.interruptAndJoinAll(0, TimeUnit.SECONDS)) // intentionally not ChunkyThread - ); - } - /** * Add a {@link Thread} to be interrupted and joined by chunky on shutdown * From f20b005e987f103f19d530a19085c048cf2b63b8 Mon Sep 17 00:00:00 2001 From: Tom Martin Date: Fri, 24 Jul 2026 20:41:47 +0100 Subject: [PATCH 12/13] Remove PluginApi from joinAll --- chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java | 1 - 1 file changed, 1 deletion(-) diff --git a/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java b/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java index 07c9bd8b66..e7aeae8a9a 100644 --- a/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java +++ b/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java @@ -91,7 +91,6 @@ public synchronized static ForkJoinPool addForkJoinPool(ForkJoinPool pool) { * @param unit the time unit of the timeout argument * @return Whether all threads were joined before returning */ - @PluginApi public static boolean joinAll(long timeout, @NotNull TimeUnit unit) { /* * This method should not be synchronized because: From e69129b4fc27a09504edcf0136720300ecb39769 Mon Sep 17 00:00:00 2001 From: Tom Martin Date: Fri, 24 Jul 2026 20:43:50 +0100 Subject: [PATCH 13/13] Refactor ChunkyThread joinAll to give stronger guarantees And hopefully be more readable --- .../llbit/util/concurrent/ChunkyThread.java | 84 +++++++++++-------- 1 file changed, 51 insertions(+), 33 deletions(-) diff --git a/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java b/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java index e7aeae8a9a..51d4631f35 100644 --- a/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java +++ b/chunky/src/java/se/llbit/util/concurrent/ChunkyThread.java @@ -89,6 +89,7 @@ public synchronized static ForkJoinPool addForkJoinPool(ForkJoinPool pool) { * * @param timeout The maximum time to wait AFTER a shutdown is initiated * @param unit the time unit of the timeout argument + * * @return Whether all threads were joined before returning */ public static boolean joinAll(long timeout, @NotNull TimeUnit unit) { @@ -102,8 +103,9 @@ public static boolean joinAll(long timeout, @NotNull TimeUnit unit) { while (true) { try { - // must wait for the latch as hitting the for loop below first causes immediate evaluation of the - // for loop iterator, potentially missing new threads. + // Must wait for the latch as hitting the for loop below first causes immediate evaluation of: + // - The for loop iterator, potentially missing new threads. + // - The end time, meaning waiting starts before shutdown begins. shutdownLatch.await(); break; } catch (InterruptedException e) { @@ -111,46 +113,25 @@ public static boolean joinAll(long timeout, @NotNull TimeUnit unit) { } } - long startTime = System.nanoTime(); - long endTime = startTime + unit.toNanos(timeout); - - boolean anyAlive = false; - + long endTime = System.nanoTime() + unit.toNanos(timeout); try { - for (ExecutorService executorService : executorServices) { - while (System.nanoTime() < endTime) { - try { - long waitTime = endTime - startTime; - if (waitTime > 0) { - executorService.awaitTermination(waitTime, TimeUnit.NANOSECONDS); - } - break; - } catch (InterruptedException e) { - interrupted = true; - } - } - anyAlive |= !executorService.isTerminated(); - } - for (Thread thread : ChunkyThread.threads) { - while (System.nanoTime() < endTime) { - try { - long waitTimeMillis = TimeUnit.NANOSECONDS.toMillis(endTime - startTime); - if (waitTimeMillis > 0) { - thread.join(waitTimeMillis); - } - break; - } catch (InterruptedException e) { - interrupted = true; + while (System.nanoTime() < endTime) { + try { + if (joinAllInterruptable(endTime)) { + // All threads are joined, skip the rest of the wait time. + return true; } + } catch (InterruptedException e) { + interrupted = true; } - anyAlive |= thread.isAlive(); } + // Got to the end of the wait time without joining everything, can give no guarantees + return false; } finally { if (interrupted) { Thread.currentThread().interrupt(); } } - return !anyAlive; } /** @@ -178,6 +159,43 @@ public static boolean interruptAndJoinAll(long timeout, @NotNull TimeUnit unit) return joinAll(timeout, unit); } + /** + * Await the joining of all threads managed by chunky. + * + *

This method is only safe to call if the {@link ChunkyThread#shutdownLatch} has been set.

+ * + *

WARNING: calling this from any thread registered with {@link #addThread(Thread)} may deadlock.

+ * + * @param endTimeNanos The time at which to stop waiting. + * + * @return Whether all threads were joined before returning + * + * @throws InterruptedException Propagates up when interrupted. The caller has no guarantee that shutdown has begun, + * or that any of the inner threads have been joined. + */ + private static boolean joinAllInterruptable(long endTimeNanos) throws InterruptedException { + // The intention here whether we return true or false, is to give the caller the most complete acquire load possible. + // Even if we reach the timeout given by the caller, we still establish a happens-before with every dead thread. + + boolean anyAlive = false; + for (ExecutorService executorService : executorServices) { + long waitTime = endTimeNanos - System.nanoTime(); + anyAlive |= !executorService.awaitTermination(waitTime, TimeUnit.NANOSECONDS); + } + for (Thread thread : threads) { + long waitTime = endTimeNanos - System.nanoTime(); + if (waitTime > 0) { + thread.join(waitTime); // joining with 0 is infinite wait time, very intuitive. + } + // Thread.isAlive() establishes a happens-before with the thread. As such the following are non-issues: + // - Not joining the thread, if waitTime <= 0 + // - The thread stopping between Thread.join() and Thread.isAlive(). + anyAlive |= thread.isAlive(); + } + + return !anyAlive; + } + private void setDefaults() { this.setDaemon(true); addThread(this);