Fix possible crashes on Threading
This commit is contained in:
@@ -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<Dynamic>) {
|
||||
#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
|
||||
}
|
||||
}
|
||||
@@ -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<Thread> = [];
|
||||
// DEPRECATED
|
||||
@:dox(hide) public static var gameThreads(get, set):Array<sys.thread.Thread>;
|
||||
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();
|
||||
|
||||
|
||||
@@ -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<Thread> = [];
|
||||
|
||||
private static var maxThreads:Int = 4;
|
||||
private static var threadCycle:Int = 0;
|
||||
private static var threadsInitialized:Bool = false;
|
||||
#else
|
||||
public static var gameThreads:Array<Dynamic> = [];
|
||||
#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
|
||||
}
|
||||
}
|
||||
@@ -122,14 +122,18 @@ class NativeAudioSource {
|
||||
var arrayType:TypedArrayType;
|
||||
var loopPoints:Array<Int>; // In Samples
|
||||
|
||||
static var threadRunning:Bool = false;
|
||||
static var streamSources:Array<NativeAudioSource> = [];
|
||||
static var queuedStreamSources:Array<NativeAudioSource> = [];
|
||||
|
||||
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() {
|
||||
|
||||
Reference in New Issue
Block a user