diff --git a/source/funkin/backend/scripting/MultiThreadedScript.hx b/source/funkin/backend/scripting/MultiThreadedScript.hx index 24754c7c..24da0d7a 100644 --- a/source/funkin/backend/scripting/MultiThreadedScript.hx +++ b/source/funkin/backend/scripting/MultiThreadedScript.hx @@ -1,14 +1,9 @@ package funkin.backend.scripting; -#if ALLOW_MULTITHREADING -import sys.thread.Thread; -#end - import hscript.IHScriptCustomBehaviour; +import funkin.backend.utils.EngineUtil; class MultiThreadedScript implements IFlxDestroyable implements IHScriptCustomBehaviour { - var thread:#if ALLOW_MULTITHREADING Thread #else Dynamic #end; - /** * Script being ran. */ @@ -46,13 +41,6 @@ class MultiThreadedScript implements IFlxDestroyable implements IHScriptCustomBe script.load(); - #if ALLOW_MULTITHREADING - thread = Thread.createWithEventLoop(function() { - // Prevent the thread from being auto deleted - Thread.current().events.promise(); - }); - #end - __variables = Type.getInstanceFields(Type.getClass(this)); } @@ -69,7 +57,7 @@ class MultiThreadedScript implements IFlxDestroyable implements IHScriptCustomBe public function call(func:String, args:Array) { #if ALLOW_MULTITHREADING - thread.events.run(function() { + EngineUtil.execAsync(() -> { callEnded = false; returnValue = script.call(func, args); callEnded = true; @@ -85,13 +73,5 @@ class MultiThreadedScript implements IFlxDestroyable implements IHScriptCustomBe script.call("destroy"); script.destroy(); } - - #if ALLOW_MULTITHREADING - if (thread != null) { - thread.events.runPromised(function() { - // close the thing - }); - } - #end } } \ No newline at end of file diff --git a/source/funkin/backend/system/Main.hx b/source/funkin/backend/system/Main.hx index ebffd227..955aa3f1 100644 --- a/source/funkin/backend/system/Main.hx +++ b/source/funkin/backend/system/Main.hx @@ -13,6 +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.editors.SaveWarning; import funkin.options.PlayerSettings; import openfl.Assets; @@ -22,10 +23,6 @@ import openfl.text.TextFormat; import openfl.utils.AssetLibrary; import sys.FileSystem; import sys.io.File; - -#if ALLOW_MULTITHREADING -import sys.thread.Thread; -#end #if android import android.content.Context; import android.os.Build; @@ -61,7 +58,10 @@ class Main extends Sprite // You can pretty much ignore everything from here on - your code should go in your states. #if ALLOW_MULTITHREADING - public static var gameThreads:Array = []; + // 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() { @@ -99,16 +99,8 @@ class Main extends Sprite #end; public static var startedFromSource:Bool = #if TEST_BUILD true #else false #end; - - private static var __threadCycle:Int = 0; - public static function execAsync(func:Void->Void) { - #if ALLOW_MULTITHREADING - var thread = gameThreads[(__threadCycle++) % gameThreads.length]; - thread.events.run(func); - #else - func(); - #end - } + // DEPRECATED + @:dox(hide) public static function execAsync(func:Void->Void) EngineUtil.execAsync(func); private static function getTimer():Int { return time = Lib.getTimer(); @@ -120,10 +112,6 @@ class Main extends Sprite MemoryUtil.init(); @:privateAccess FlxG.game.getTimer = getTimer; - #if ALLOW_MULTITHREADING - for(i in 0...4) - gameThreads.push(Thread.createWithEventLoop(function() {Thread.current().events.promise();})); - #end FunkinCache.init(); Paths.assetsTree = new AssetsLibraryList(); diff --git a/source/funkin/backend/utils/EngineUtil.hx b/source/funkin/backend/utils/EngineUtil.hx index ebfe916c..030e681e 100644 --- a/source/funkin/backend/utils/EngineUtil.hx +++ b/source/funkin/backend/utils/EngineUtil.hx @@ -1,9 +1,16 @@ 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. @@ -12,4 +19,32 @@ 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/lime/_internal/backend/native/NativeAudioSource.hx b/source/lime/_internal/backend/native/NativeAudioSource.hx index 9a7b46b7..1c55805d 100644 --- a/source/lime/_internal/backend/native/NativeAudioSource.hx +++ b/source/lime/_internal/backend/native/NativeAudioSource.hx @@ -122,14 +122,18 @@ class NativeAudioSource { var arrayType:TypedArrayType; var loopPoints:Array; // In Samples - static var threadRunning:Bool = false; static var streamSources:Array = []; static var queuedStreamSources:Array = []; static var streamMutex:Mutex = new Mutex(); - static var streamThread:Thread; static var streamTimer:Timer; + #if !ALLOW_MULTITHREADING + static var wasEmpty:Bool = false; + static var threadRunning:Bool = false; + static var streamThread:Thread; + #end + var streamRemove:Bool; var bufferLength:Int; // Size in bytes for current streamed audio buffers. @@ -469,27 +473,31 @@ class NativeAudioSource { streamMutex.release(); } - static function streamThreadRun() { - var i:Int, source:NativeAudioSource, process:Int, v:Int; + static function streamBuffersUpdate() { + streamMutex.acquire(); - while ((i = Thread.readMessage(true)) != 0) { - streamMutex.acquire(); - while (i-- > 0) { - if ((source = streamSources[i]).streamRemove) continue; - else if (source.parent.buffer == null) { - source.stopStream(); - continue; - } - - process = source.requestBuffers < STREAM_MIN_BUFFERS ? STREAM_MIN_BUFFERS - source.requestBuffers : 0; - process = STREAM_PROCESS_BUFFERS > process ? STREAM_PROCESS_BUFFERS : process; - if ((process = (v = STREAM_MAX_BUFFERS - source.requestBuffers) > process ? process : v) > 0) source.fillBuffers(process); + var i:Int = streamSources.length, source:NativeAudioSource, process:Int, v:Int; + while (i-- > 0) { + if ((source = streamSources[i]).streamRemove) continue; + else if (source.parent.buffer == null) { + source.stopStream(); + continue; } - streamMutex.release(); + + process = source.requestBuffers < STREAM_MIN_BUFFERS ? STREAM_MIN_BUFFERS - source.requestBuffers : 0; + process = STREAM_PROCESS_BUFFERS > process ? STREAM_PROCESS_BUFFERS : process; + if ((process = (v = STREAM_MAX_BUFFERS - source.requestBuffers) > process ? process : v) > 0) source.fillBuffers(process); } + streamMutex.release(); + } + + #if !ALLOW_MULTITHREADING + static function streamThreadRun() { + while (Thread.readMessage(true)) streamBuffersUpdate(); threadRunning = false; } + #end static function streamUpdate() { if (!streamMutex.tryAcquire()) return; @@ -512,13 +520,28 @@ class NativeAudioSource { } } - streamMutex.release(); + #if ALLOW_MULTITHREADING + if (streamSources.length != 0) funkin.backend.utils.EngineUtil.execAsync(streamBuffersUpdate); + #else if (streamSources.length == 0) { - streamTimer.stop(); - if (threadRunning) streamThread.sendMessage(0); + if (wasEmpty) { + wasEmpty = false; + streamTimer.stop(); + if (threadRunning) streamThread.sendMessage(1); + } + else { + wasEmpty = true; + streamTimer = resetTimer(streamTimer, 1000, streamUpdate); + } } - else if (threadRunning || (threadRunning = (streamThread = Thread.create(streamThreadRun)) != null)) - streamThread.sendMessage(streamSources.length); + else { + wasEmpty = false; + if (threadRunning || (threadRunning = (streamThread = Thread.create(streamThreadRun)) != null)) + streamThread.sendMessage(1); + } + #end + + streamMutex.release(); } function removeStream() {