From 1b403b2ab2f61f87885ec67636b72c26692c7957 Mon Sep 17 00:00:00 2001 From: Ralty <78720179+Raltyro@users.noreply.github.com> Date: Sat, 29 Nov 2025 11:24:39 +0700 Subject: [PATCH] oh, theres ThreadUtil --- .../backend/scripting/MultiThreadedScript.hx | 3 +- source/funkin/backend/scripting/Script.hx | 1 + source/funkin/backend/system/Main.hx | 11 +- .../backend/system/updating/UpdateUtil.hx | 40 +++++-- source/funkin/backend/utils/AudioAnalyzer.hx | 2 +- source/funkin/backend/utils/EngineUtil.hx | 35 ------ source/funkin/backend/utils/ThreadUtil.hx | 103 ++++++++++++++---- .../backend/native/NativeAudioSource.hx | 18 +-- 8 files changed, 129 insertions(+), 84 deletions(-) diff --git a/source/funkin/backend/scripting/MultiThreadedScript.hx b/source/funkin/backend/scripting/MultiThreadedScript.hx index 24da0d7a..517dc96e 100644 --- a/source/funkin/backend/scripting/MultiThreadedScript.hx +++ b/source/funkin/backend/scripting/MultiThreadedScript.hx @@ -1,7 +1,6 @@ package funkin.backend.scripting; import hscript.IHScriptCustomBehaviour; -import funkin.backend.utils.EngineUtil; class MultiThreadedScript implements IFlxDestroyable implements IHScriptCustomBehaviour { /** @@ -57,7 +56,7 @@ class MultiThreadedScript implements IFlxDestroyable implements IHScriptCustomBe public function call(func:String, args:Array) { #if ALLOW_MULTITHREADING - EngineUtil.execAsync(() -> { + funkin.backend.utils.ThreadUtil.execAsync(() -> { callEnded = false; returnValue = script.call(func, args); callEnded = true; diff --git a/source/funkin/backend/scripting/Script.hx b/source/funkin/backend/scripting/Script.hx index 8f7ee437..226971cb 100644 --- a/source/funkin/backend/scripting/Script.hx +++ b/source/funkin/backend/scripting/Script.hx @@ -102,6 +102,7 @@ class Script extends FlxBasic implements IFlxDestroyable { #if sys "ZipUtil" => funkin.backend.utils.ZipUtil, #end "MarkdownUtil" => funkin.backend.utils.MarkdownUtil, "EngineUtil" => funkin.backend.utils.EngineUtil, + "ThreadUtil" => funkin.backend.utils.ThreadUtil, "MemoryUtil" => funkin.backend.utils.MemoryUtil, "BitmapUtil" => funkin.backend.utils.BitmapUtil, diff --git a/source/funkin/backend/system/Main.hx b/source/funkin/backend/system/Main.hx index 955aa3f1..31d270d9 100644 --- a/source/funkin/backend/system/Main.hx +++ b/source/funkin/backend/system/Main.hx @@ -13,7 +13,7 @@ import funkin.backend.assets.ModsFolder; import funkin.backend.system.framerate.Framerate; import funkin.backend.system.framerate.SystemInfo; import funkin.backend.system.modules.*; -import funkin.backend.utils.EngineUtil; +import funkin.backend.utils.ThreadUtil; import funkin.editors.SaveWarning; import funkin.options.PlayerSettings; import openfl.Assets; @@ -57,13 +57,6 @@ class Main extends Sprite // You can pretty much ignore everything from here on - your code should go in your states. - #if ALLOW_MULTITHREADING - // DEPRECATED - @:dox(hide) public static var gameThreads(get, set):Array; - static function get_gameThreads() return EngineUtil.gameThreads; - static function set_gameThreads(v) return EngineUtil.gameThreads = v; - #end - public static function preInit() { funkin.backend.utils.NativeAPI.registerAsDPICompatible(); funkin.backend.system.CommandLineHandler.parseCommandLine(Sys.args()); @@ -100,7 +93,7 @@ class Main extends Sprite public static var startedFromSource:Bool = #if TEST_BUILD true #else false #end; // DEPRECATED - @:dox(hide) public static function execAsync(func:Void->Void) EngineUtil.execAsync(func); + @:dox(hide) public static function execAsync(func:Void->Void) ThreadUtil.execAsync(func); private static function getTimer():Int { return time = Lib.getTimer(); diff --git a/source/funkin/backend/system/updating/UpdateUtil.hx b/source/funkin/backend/system/updating/UpdateUtil.hx index c3f49654..6a3bf7b4 100644 --- a/source/funkin/backend/system/updating/UpdateUtil.hx +++ b/source/funkin/backend/system/updating/UpdateUtil.hx @@ -2,22 +2,30 @@ package funkin.backend.system.updating; import funkin.backend.system.github.GitHub; import funkin.backend.system.github.GitHubRelease; +#if ALLOW_MULTITHREADING +import funkin.backend.utils.ThreadUtil; +#end import lime.app.Application; -import sys.thread.Mutex; -import sys.thread.Thread; import sys.FileSystem; import haxe.io.Path; +#if (target.threaded) +import sys.thread.Thread; +import sys.thread.Mutex; +#end + using funkin.backend.system.github.GitHub; class UpdateUtil { public static var lastUpdateCheck:Null; + #if (target.threaded) private static var __waitCallbacks:ArrayVoid>; private static var __mutex:Mutex; + #end public static function init() { // deletes old bak file if it exists @@ -26,26 +34,34 @@ class UpdateUtil { if (FileSystem.exists(bakPath)) FileSystem.deleteFile(bakPath); #end + #if (target.threaded) __waitCallbacks = []; __mutex = new Mutex(); - Thread.create(checkForUpdates.bind(true, false)); + + #if ALLOW_MULTITHREADING ThreadUtil.execAsync #else Thread.create #end(checkForUpdates.bind(true, false)); + #end } public static function waitForUpdates(force = false, callback:UpdateCheckCallback->Void, lazy = false) { + #if (target.threaded) if (__mutex.tryAcquire()) { __mutex.release(); if (__shouldCheck(lazy) || force) { __waitCallbacks.push(callback); - Thread.create(checkForUpdates.bind(force, false)); + #if ALLOW_MULTITHREADING ThreadUtil.execAsync #else Thread.create #end(checkForUpdates.bind(force, false)); } else callback(lastUpdateCheck); } else __waitCallbacks.push(callback); + #else + callback(checkForUpdates(true, false)); + #end } public static function checkForUpdates(force = false, lazy = false):UpdateCheckCallback { + #if (target.threaded) var wasAcquired = !__mutex.tryAcquire(); if (wasAcquired) __mutex.acquire(); @@ -60,8 +76,19 @@ class UpdateUtil { FlxG.signals.preUpdate.addOnce(__callWaitCallbacks); return lastUpdateCheck; + #else + if (!__shouldCheck(lazy)) return lastUpdateCheck; + return lastUpdateCheck = __checkForUpdates(); + #end } + #if (target.threaded) + static function __callWaitCallbacks() { + for (callback in __waitCallbacks) callback(lastUpdateCheck); + __waitCallbacks.resize(0); + } + #end + static function __checkForUpdates():UpdateCheckCallback { var curTag = 'v' + (Flags.VERSION == null ? Application.current.meta.get('version') : Flags.VERSION), error = false; var newUpdates = __doReleaseFiltering(GitHub.getReleases(Flags.REPO_OWNER, Flags.REPO_NAME, (e) -> { @@ -82,11 +109,6 @@ class UpdateUtil { static function __shouldCheck(lazy:Bool):Bool return lastUpdateCheck == null || !lazy && (!lastUpdateCheck.newUpdate || Date.now().getTime() - lastUpdateCheck.date.getTime() > 1800000); - static function __callWaitCallbacks() { - for (callback in __waitCallbacks) callback(lastUpdateCheck); - __waitCallbacks.resize(0); - } - static function __doReleaseFiltering(releases:Array, currentVersionTag:String) { releases = releases.filterReleases(Options.betaUpdates, false); if (releases.length <= 0) diff --git a/source/funkin/backend/utils/AudioAnalyzer.hx b/source/funkin/backend/utils/AudioAnalyzer.hx index 81045098..c6f02cd6 100644 --- a/source/funkin/backend/utils/AudioAnalyzer.hx +++ b/source/funkin/backend/utils/AudioAnalyzer.hx @@ -95,9 +95,9 @@ final class AudioAnalyzer { static var __twiddleImags:Array> = []; static var __freqReals:Array> = []; static var __freqImags:Array> = []; + static var __freqCalculating:Int = 0; #if (target.threaded) static var __mutex:Mutex = new Mutex(); - static var __freqCalculating:Int = 0; #end /** diff --git a/source/funkin/backend/utils/EngineUtil.hx b/source/funkin/backend/utils/EngineUtil.hx index 030e681e..ebfe916c 100644 --- a/source/funkin/backend/utils/EngineUtil.hx +++ b/source/funkin/backend/utils/EngineUtil.hx @@ -1,16 +1,9 @@ package funkin.backend.utils; -#if ALLOW_MULTITHREADING -import sys.thread.Thread; -#end - -#if !macro import funkin.backend.scripting.MultiThreadedScript; import funkin.backend.scripting.Script; -#end final class EngineUtil { - #if !macro /** * Starts a new multithreaded script. * This script will share all the variables with the current one, which means already existing callbacks will be replaced by new ones on conflict. @@ -19,32 +12,4 @@ final class EngineUtil { public static function startMultithreadedScript(path:String) { return new MultiThreadedScript(path, Script.curScript); } - #end - - #if ALLOW_MULTITHREADING - public static var gameThreads:Array = []; - - private static var maxThreads:Int = 4; - private static var threadCycle:Int = 0; - private static var threadsInitialized:Bool = false; - #else - public static var gameThreads:Array = []; - #end - - /** - * Execute a function asynchronously using existing threads when initialized with ALLOW_MULTITHREADING. - * @param func Void -> Void - */ - public static function execAsync(func:Void->Void) { - #if ALLOW_MULTITHREADING - if (!threadsInitialized) { - threadsInitialized = true; - for (i in 0...maxThreads) gameThreads.push(Thread.createWithEventLoop(() -> Thread.current().events.promise())); - } - gameThreads[threadCycle].events.run(func); - if (++threadCycle >= maxThreads) threadCycle = 0; - #else - func(); - #end - } } \ No newline at end of file diff --git a/source/funkin/backend/utils/ThreadUtil.hx b/source/funkin/backend/utils/ThreadUtil.hx index 8c2f654c..9b906ad7 100644 --- a/source/funkin/backend/utils/ThreadUtil.hx +++ b/source/funkin/backend/utils/ThreadUtil.hx @@ -1,32 +1,97 @@ package funkin.backend.utils; -#if ALLOW_MULTITHREADING +#if (target.threaded) +import sys.thread.Deque; +import sys.thread.Thread; +import sys.thread.Mutex; +#else +private typedef Thread = Dynamic; +#end + +#if !macro +import funkin.backend.system.Logs; +#end + final class ThreadUtil { + inline static function error(text:String) { + #if macro + trace(text); + #else + FlxG.signals.preUpdate.addOnce(Logs.error.bind(text)); + #end + } + /** * Creates a new Thread with an error handler. * @param func Function to execute * @param autoRestart Whenever the thread should auto restart itself after crashing. */ - public static function createSafe(func:Void->Void, autoRestart:Bool = false) { - if (autoRestart) { - return sys.thread.Thread.create(function() { - while(true) { - try { - func(); - } catch(e) { - trace(e.details()); - } - } - }); - } else { - return sys.thread.Thread.create(function() { - try { + public static function createSafe(func:Void->Void, autoRestart:Bool = false):Thread { + #if (target.threaded) + try { + return if (autoRestart) Thread.create(() -> { + var restart = true; + while (restart) try { func(); - } catch(e) { - trace(e.details()); + restart = false; } + catch (e) error(e.details()); + }) + else Thread.create(() -> { + try {func();} + catch (e) error(e.details()); }); } + catch (e) error("Failed to safely create a thread: " + e.details()); + #end + return null; } -} -#end \ No newline at end of file + + #if ALLOW_MULTITHREADING + public static var maxThreads:Int = 4; + + static var __threads:Array = []; + static var __pendingExecs:DequeVoid> = new Deque(); + static var __threadMutex:Mutex = new Mutex(); + static var __threadUsed:Int = 0; + + static function __threadExecAsync() { + var callback:Void->Void; + while ((callback = __pendingExecs.pop(true)) != null) { + __threadMutex.acquire(); + __threadUsed++; + __threadMutex.release(); + + callback(); + + __threadMutex.acquire(); + __threadUsed--; + __threadMutex.release(); + } + __threadMutex.acquire(); + __threads.remove(Thread.current()); + __threadMutex.release(); + } + #end + + public static function execAsync(func:Void->Void) { + if (func == null) return; + + #if (ALLOW_MULTITHREADING && !macro) + __pendingExecs.add(func); + if (__threadUsed >= __threads.length) { + if (__threads.length == maxThreads) return; + + __threadMutex.acquire(); + try { + var thread = Thread.create(__threadExecAsync); + __threads.push(thread); + } + catch (e) Logs.warn(e.details()); + __threadMutex.release(); + } + #else + func(); + #end + } +} \ No newline at end of file diff --git a/source/lime/_internal/backend/native/NativeAudioSource.hx b/source/lime/_internal/backend/native/NativeAudioSource.hx index 1c55805d..9f09b244 100644 --- a/source/lime/_internal/backend/native/NativeAudioSource.hx +++ b/source/lime/_internal/backend/native/NativeAudioSource.hx @@ -39,9 +39,9 @@ class NativeAudioSource { public static var STREAM_BUFFER_SAMPLES:Int = 0x2000; // how much buffers will be generating every frequency (doesnt have to be pow of 2?). public static var STREAM_MIN_BUFFERS:Int = 2; // how much buffers can a stream hold on minimum or starting. public static var STREAM_MAX_BUFFERS:Int = 8; // how much limit of a buffers can be used for streamed audios, must be higher than minimum. - public static var STREAM_FLUSH_BUFFERS:Int = 3; // how much buffers can it play. + public static var STREAM_MAX_FLUSH_BUFFERS:Int = 3; // how much buffers can it play. public static var STREAM_PROCESS_BUFFERS:Int = 2; // how much buffers can be processed in a frequency tick. - public static var MAX_POOL_BUFFERS:Int = 32; // how much buffers for the pool to hold. + public static var POOL_MAX_BUFFERS:Int = 32; // how much buffers for the pool to hold. public static var moreFormatsSupported:Null; public static var loopPointsSupported:Null; @@ -191,7 +191,7 @@ class NativeAudioSource { } if (bufferDatas != null) { - for (data in bufferDatas) if (bufferDataPool.length < MAX_POOL_BUFFERS) bufferDataPool.push(data); + for (data in bufferDatas) if (bufferDataPool.length < POOL_MAX_BUFFERS) bufferDataPool.push(data); bufferDatas = null; } @@ -268,7 +268,7 @@ class NativeAudioSource { final length = STREAM_BUFFER_SAMPLES * channels; bufferLength = length * wordSize; - if (buffers == null) buffers = AL.genBuffers(STREAM_FLUSH_BUFFERS); + if (buffers == null) buffers = AL.genBuffers(STREAM_MAX_FLUSH_BUFFERS); if (bufferDatas == null) { bufferDatas = []; bufferTimes = []; @@ -305,7 +305,7 @@ class NativeAudioSource { AL.deleteBuffers(buffers); buffers = null; - for (data in bufferDatas) if (bufferDataPool.length < MAX_POOL_BUFFERS) bufferDataPool.push(data); + for (data in bufferDatas) if (bufferDataPool.length < POOL_MAX_BUFFERS) bufferDataPool.push(data); bufferDatas.resize(0); } @@ -400,7 +400,7 @@ class NativeAudioSource { } } catch (e:haxe.Exception) { - trace('NativeAudioSource readToBufferData Bug! error: ${e.message} | ${e.stack.toString()}, streamEnded: $streamEnded, total: $total, n: $n'); + trace('NativeAudioSource readToBufferData Bug! error: ${e.details()}, streamEnded: $streamEnded, total: $total, n: $n'); return result; } @@ -431,10 +431,10 @@ class NativeAudioSource { inline function flushBuffers() { var i = STREAM_MAX_BUFFERS - (requestBuffers - queuedBuffers); - while (queuedBuffers < STREAM_FLUSH_BUFFERS && queuedBuffers < requestBuffers) { + while (queuedBuffers < STREAM_MAX_FLUSH_BUFFERS && queuedBuffers < requestBuffers) { AL.bufferData(buffers[nextBuffer], format, bufferDatas[i], bufferLengths[i], sampleRate); AL.sourceQueueBuffer(source, buffers[nextBuffer]); - if (++nextBuffer == STREAM_FLUSH_BUFFERS) nextBuffer = 0; + if (++nextBuffer == STREAM_MAX_FLUSH_BUFFERS) nextBuffer = 0; queuedBuffers++; i++; } @@ -521,7 +521,7 @@ class NativeAudioSource { } #if ALLOW_MULTITHREADING - if (streamSources.length != 0) funkin.backend.utils.EngineUtil.execAsync(streamBuffersUpdate); + if (streamSources.length != 0) funkin.backend.utils.ThreadUtil.execAsync(streamBuffersUpdate); #else if (streamSources.length == 0) { if (wasEmpty) {