!14 thread synchronization problem

Merge pull request !14 from yaojian16/master
This commit is contained in:
openharmony_ci
2021-10-08 01:27:49 +00:00
committed by Gitee
2 changed files with 525 additions and 81 deletions
+485 -75
View File
@@ -73,7 +73,7 @@ void CallWorkCallback(napi_env env, napi_value recv, size_t argc, const napi_val
}
}
void Worker::PrepareForWorkerInstance(const Worker* worker)
bool Worker::PrepareForWorkerInstance(const Worker* worker)
{
napi_env env = worker->GetWorkerEnv();
// 1. init worker environment
@@ -84,7 +84,7 @@ void Worker::PrepareForWorkerInstance(const Worker* worker)
if (OHOS::CCRuntime::Worker::WorkerCore::getAssertFunc == NULL) {
HILOG_ERROR("worker::getAssertFunc is null");
napi_throw_error(env, nullptr, "worker::getAssertFunc is null");
return;
return false;
}
std::vector<uint8_t> scriptContent;
OHOS::CCRuntime::Worker::WorkerCore::getAssertFunc(worker->GetScript(), scriptContent);
@@ -95,7 +95,7 @@ void Worker::PrepareForWorkerInstance(const Worker* worker)
// An exception occurred when running the script.
HILOG_ERROR("worker:: run script exception occurs, will handle exception");
(const_cast<Worker*>(worker))->HandleException();
return;
return false;
}
// 3. register postMessage in DedicatedWorkerGlobalScope
@@ -115,6 +115,7 @@ void Worker::PrepareForWorkerInstance(const Worker* worker)
napi_create_string_utf8(env, workerName.c_str(), workerName.length(), &nameValue);
NapiValueHelp::SetNamePropertyInGlobal(env, "name", nameValue);
}
return true;
}
bool Worker::UpdateWorkerState(RunnerState state)
@@ -131,50 +132,77 @@ bool Worker::UpdateWorkerState(RunnerState state)
return true;
}
bool Worker::UpdateMainState(MainState state)
{
bool done = false;
do {
MainState oldState = mainState_.load(std::memory_order_acquire);
if (oldState >= state) {
// make sure state sequence is ACTIVE, INACTIVE
return false;
}
done = mainState_.compare_exchange_strong(oldState, state);
} while (!done);
return true;
}
void Worker::PublishWorkerOverSignal()
{
// post NULL tell main worker is not running
mainMessageQueue_.EnQueue(NULL);
uv_async_send(&mainOnMessageSignal_);
TriggerPostTask();
if (!MainIsStop()) {
mainMessageQueue_.EnQueue(NULL);
uv_async_send(&mainOnMessageSignal_);
TriggerPostTask();
}
}
void Worker::ExecuteInThread(const void* data)
{
auto worker = reinterpret_cast<Worker*>(const_cast<void*>(data));
// 1. create a runtime, nativeengine
napi_env env = worker->GetMainEnv();
napi_env newEnv = nullptr;
napi_create_runtime(env, &newEnv);
if (newEnv == nullptr) {
napi_throw_error(env, nullptr, "Worker create runtime error");
return;
napi_env workerEnv = nullptr;
{
std::lock_guard<std::recursive_mutex> lock(worker->liveStatusLock_);
if (worker->MainIsStop()) {
CloseHelp::DeletePointer(worker, false);
return;
}
napi_env env = worker->GetMainEnv();
napi_create_runtime(env, &workerEnv);
if (workerEnv == nullptr) {
napi_throw_error(env, nullptr, "Worker create runtime error");
return;
}
// mark worker env is subThread
reinterpret_cast<NativeEngine*>(workerEnv)->MarkSubThread();
worker->SetWorkerEnv(workerEnv);
}
// mark worker env is subThread
reinterpret_cast<NativeEngine*>(newEnv)->MarkSubThread();
worker->SetWorkerEnv(newEnv);
uv_loop_t* loop = worker->GetWorkerLoop();
if (loop == nullptr) {
napi_throw_error(env, nullptr, "Worker loop is nullptr");
HILOG_ERROR("worker:: Worker loop is nullptr");
return;
}
uv_async_init(loop, &worker->workerOnMessageSignal_, reinterpret_cast<uv_async_cb>(Worker::WorkerOnMessage));
if (worker->UpdateWorkerState(RUNNING)) {
// 2. add some preparation for the worker
PrepareForWorkerInstance(worker);
// 2. add some preparation for the worker
if (PrepareForWorkerInstance(worker)) {
uv_async_init(loop, &worker->workerOnMessageSignal_, reinterpret_cast<uv_async_cb>(Worker::WorkerOnMessage));
worker->UpdateWorkerState(RUNNING);
// in order to invoke worker send before subThread start
uv_async_send(&worker->workerOnMessageSignal_);
// 3. start worker loop
if (worker->GetWorkerEnv() == nullptr) {
HILOG_ERROR("worker::worker engine is null");
} else {
uv_async_send(&worker->workerOnMessageSignal_);
worker->Loop();
}
worker->Loop();
} else {
worker->CloseInner();
HILOG_ERROR("worker:: worker PrepareForWorkerInstance failure");
worker->UpdateWorkerState(TERMINATED);
}
worker->ReleaseWorkerThreadContent();
std::lock_guard<std::recursive_mutex> lock(worker->liveStatusLock_);
if (worker->MainIsStop()) {
CloseHelp::DeletePointer(worker, false);
} else {
worker->PublishWorkerOverSignal();
}
worker->PublishWorkerOverSignal();
}
void Worker::MainOnMessage(const uv_async_t* req)
@@ -189,6 +217,10 @@ void Worker::MainOnMessage(const uv_async_t* req)
void Worker::MainOnErrorInner()
{
if (mainEnv_ == nullptr || MainIsStop()) {
HILOG_ERROR("worker:: main thread maybe is over");
return;
}
napi_value callback = nullptr;
napi_value obj = nullptr;
napi_get_reference_value(mainEnv_, workerWrapper_, &obj);
@@ -239,12 +271,6 @@ void Worker::CloseMainCallback() const
napi_create_int32(mainEnv_, 1, &exitValue);
napi_value argv[1] = { exitValue };
CallMainFunction(1, argv, "onexit");
std::lock_guard<std::mutex> lock(g_workersMutex);
std::list<Worker*>::iterator it = std::find(g_workers.begin(), g_workers.end(), this);
if (it != g_workers.end()) {
g_workers.erase(it);
}
CloseHelp::DeletePointer(this, false);
}
@@ -274,6 +300,10 @@ void Worker::HandleEventListeners(napi_env env, napi_value recv, size_t argc, co
void Worker::MainOnMessageInner()
{
if (mainEnv_ == nullptr || MainIsStop()) {
HILOG_ERROR("worker:: main thread maybe is over");
return;
}
napi_value callback = nullptr;
napi_value obj = nullptr;
napi_get_reference_value(mainEnv_, workerWrapper_, &obj);
@@ -315,7 +345,6 @@ void Worker::MainOnMessageInner()
void Worker::TerminateWorker()
{
// when there is no active handle, worker loop will stop automatic.
std::lock_guard<std::mutex> lock(workerAsyncMutex_);
uv_close((uv_handle_t*)&workerOnMessageSignal_, nullptr);
CloseWorkerCallback();
uv_loop_t* loop = GetWorkerLoop();
@@ -344,9 +373,14 @@ void Worker::HandleException()
if (mainEnv_ != nullptr) {
napi_value data = nullptr;
napi_serialize(workerEnv_, exception, NapiValueHelp::GetUndefinedValue(workerEnv_), &data);
errorQueue_.EnQueue(data);
uv_async_send(&mainOnErrorSignal_);
TriggerPostTask();
{
std::lock_guard<std::recursive_mutex> lock(liveStatusLock_);
if (!MainIsStop()) {
errorQueue_.EnQueue(data);
uv_async_send(&mainOnErrorSignal_);
TriggerPostTask();
}
}
} else {
HILOG_ERROR("worker:: main engine is nullptr.");
}
@@ -385,6 +419,10 @@ void Worker::WorkerOnMessageInner()
void Worker::MainOnMessageErrorInner()
{
if (mainEnv_ == nullptr || MainIsStop()) {
HILOG_ERROR("worker:: main thread maybe is over");
return;
}
napi_value obj = nullptr;
napi_get_reference_value(mainEnv_, workerWrapper_, &obj);
CallMainFunction(0, nullptr, "onmessageerror");
@@ -487,7 +525,8 @@ napi_value Worker::PostMessageToMain(napi_env env, napi_callback_info cbinfo)
void Worker::PostMessageToMainInner(MessageDataType data)
{
if (mainEnv_ != nullptr) {
std::lock_guard<std::recursive_mutex> lock(liveStatusLock_);
if (mainEnv_ != nullptr && !MainIsStop()) {
mainMessageQueue_.EnQueue(data);
uv_async_send(&mainOnMessageSignal_);
TriggerPostTask();
@@ -506,7 +545,6 @@ void Worker::PostMessageInner(MessageDataType data)
HILOG_INFO("worker:: worker has been terminated.");
return;
}
std::lock_guard<std::mutex> lock(workerAsyncMutex_);
workerMessageQueue_.EnQueue(data);
if (IsRunning()) {
uv_async_send(&workerOnMessageSignal_);
@@ -533,6 +571,10 @@ napi_value Worker::Terminate(napi_env env, napi_callback_info cbinfo)
void Worker::TerminateInner()
{
if (IsTerminated() || IsTerminating()) {
HILOG_INFO("worker:: worker is not in running");
return;
}
// 1. send null signal
PostMessageInner(NULL);
UpdateWorkerState(TERMINATEING);
@@ -540,27 +582,54 @@ void Worker::TerminateInner()
Worker::~Worker()
{
workerMessageQueue_.Clear(mainEnv_);
mainMessageQueue_.Clear(workerEnv_);
// set thisVar's nativepointer is null
napi_value thisVar = nullptr;
napi_get_reference_value(mainEnv_, workerWrapper_, &thisVar);
Worker* worker = nullptr;
napi_remove_wrap(mainEnv_, thisVar, (void**)&worker);
napi_delete_reference(mainEnv_, workerWrapper_);
workerWrapper_ = nullptr;
napi_delete_reference(mainEnv_, parentPort_);
parentPort_ = nullptr;
CloseHelp::DeletePointer(reinterpret_cast<NativeEngine*>(workerEnv_), false);
workerEnv_ = nullptr;
mainEnv_ = nullptr;
if (!MainIsStop()) {
ReleaseMainThreadContent();
}
RemoveAllListenerInner();
}
napi_value Worker::CancelTask(napi_env env, napi_callback_info cbinfo)
{
napi_value thisVar = nullptr;
napi_get_cb_info(env, cbinfo, nullptr, nullptr, &thisVar, nullptr);
Worker* worker = nullptr;
napi_unwrap(env, thisVar, (void**)&worker);
if (worker == nullptr) {
HILOG_ERROR("worker:: worker is nullptr when CancelTask, maybe worker is terminated");
return nullptr;
}
if (worker->IsTerminated() || worker->IsTerminating()) {
HILOG_INFO("worker:: worker is not in running");
return nullptr;
}
if (!worker->ClearWorkerTasks()) {
HILOG_ERROR("worker:: clear worker task error");
}
return NapiValueHelp::GetUndefinedValue(env);
}
napi_value Worker::ParentPortCancelTask(napi_env env, napi_callback_info cbinfo)
{
Worker* worker = nullptr;
napi_get_cb_info(env, cbinfo, nullptr, nullptr, nullptr, (void**)&worker);
if (worker == nullptr) {
HILOG_ERROR("worker:: worker is nullptr when CancelTask, maybe worker is terminated");
return nullptr;
}
if (worker->IsTerminated() || worker->IsTerminating()) {
HILOG_INFO("worker:: worker is not in running");
return nullptr;
}
if (!worker->ClearWorkerTasks()) {
HILOG_ERROR("worker:: clear worker task error");
}
return NapiValueHelp::GetUndefinedValue(env);
}
napi_value Worker::WorkerConstructor(napi_env env, napi_callback_info cbinfo)
{
// check argv count
@@ -580,17 +649,23 @@ napi_value Worker::WorkerConstructor(napi_env env, napi_callback_info cbinfo)
napi_throw_error(env, nullptr, "Worker 1st param must be string with new");
return nullptr;
}
Worker* worker = nullptr;
{
std::lock_guard<std::mutex> lock(g_workersMutex);
if (g_workers.size() >= MAXWORKERS) {
napi_throw_error(env, nullptr, "Too many workers, the number of workers exceeds the maximum.");
return nullptr;
}
std::lock_guard<std::mutex> lock(g_workersMutex);
if (g_workers.size() >= MAXWORKERS) {
napi_throw_error(env, nullptr, "Too many workers, the number of workers exceeds the maximum.");
return nullptr;
// 2. new worker instance
worker = new Worker(env, nullptr);
if (worker == nullptr) {
napi_throw_error(env, nullptr, "create worker error");
return nullptr;
}
g_workers.push_back(worker);
}
// 2. new worker instance
Worker* worker = new Worker(env, nullptr);
g_workers.push_back(worker);
if (argc > 1 && NapiValueHelp::IsObject(args[1])) {
napi_value nameValue = nullptr;
napi_get_named_property(env, args[1], "name", &nameValue);
@@ -637,21 +712,27 @@ napi_value Worker::WorkerConstructor(napi_env env, napi_callback_info cbinfo)
napi_throw_error(env, nullptr, "worker script create error, please check.");
return nullptr;
}
HILOG_INFO("worker:: script is %{public}s", script);
worker->StartExecuteInThread(env, script);
napi_wrap(
env, thisVar, worker,
[](napi_env env, void* data, void* hint) {
Worker* worker = (Worker*)data;
auto iter = std::find(g_workers.begin(), g_workers.end(), worker);
if (iter == g_workers.end()) {
return;
{
std::lock_guard<std::recursive_mutex> lock(worker->liveStatusLock_);
if (worker->UpdateMainState(INACTIVE)) {
uv_unref((uv_handle_t*)&worker->mainOnMessageSignal_);
uv_close((uv_handle_t*)&worker->mainOnMessageSignal_, nullptr);
uv_unref((uv_handle_t*)&worker->mainOnErrorSignal_);
uv_close((uv_handle_t*)&worker->mainOnErrorSignal_, nullptr);
worker->ReleaseMainThreadContent();
}
if (!worker->IsRunning()) {
HILOG_INFO("worker:: worker is not in running");
return;
}
worker->TerminateInner();
}
if (worker->IsTerminated() || worker->IsTerminating()) {
HILOG_INFO("worker:: worker is not in running");
return;
}
worker->TerminateInner();
},
nullptr, nullptr);
napi_create_reference(env, thisVar, 1, &worker->workerWrapper_);
@@ -935,6 +1016,7 @@ napi_value Worker::InitWorker(napi_env env, napi_value exports)
DECLARE_NAPI_FUNCTION("dispatchEvent", DispatchEvent),
DECLARE_NAPI_FUNCTION("removeEventListener", RemoveEventListener),
DECLARE_NAPI_FUNCTION("removeAllListener", RemoveAllListener),
DECLARE_NAPI_FUNCTION("cancelTasks", CancelTask),
};
napi_value workerClazz = nullptr;
napi_define_class(env, className, sizeof(className), Worker::WorkerConstructor, nullptr,
@@ -955,10 +1037,23 @@ napi_value Worker::InitWorker(napi_env env, napi_value exports)
napi_property_descriptor properties[] = {
DECLARE_NAPI_FUNCTION_WITH_DATA("postMessage", PostMessageToMain, worker),
DECLARE_NAPI_FUNCTION_WITH_DATA("close", CloseWorker, worker),
DECLARE_NAPI_FUNCTION_WITH_DATA("cancelTasks", ParentPortCancelTask, worker),
DECLARE_NAPI_FUNCTION_WITH_DATA("addEventListener", ParentPortAddEventListener, worker),
DECLARE_NAPI_FUNCTION_WITH_DATA("dispatchEvent", ParentPortDispatchEvent, worker),
DECLARE_NAPI_FUNCTION_WITH_DATA("removeEventListener", ParentPortRemoveEventListener, worker),
DECLARE_NAPI_FUNCTION_WITH_DATA("removeAllListener", ParentPortRemoveAllListener, worker),
};
napi_value parentPortObj = nullptr;
napi_create_object(env, &parentPortObj);
napi_define_properties(env, parentPortObj, sizeof(properties) / sizeof(properties[0]), properties);
// 5. register worker name in DedicatedWorkerGlobalScope
std::string workerName = worker->GetName();
if (!workerName.empty()) {
napi_value nameValue = nullptr;
napi_create_string_utf8(env, workerName.c_str(), workerName.length(), &nameValue);
napi_set_named_property(env, parentPortObj, "name", nameValue);
}
napi_set_named_property(env, exports, "parentPort", parentPortObj);
// register worker parentPort.
@@ -975,6 +1070,9 @@ void Worker::WorkerOnErrorInner(napi_value error)
bool Worker::CallWorkerFunction(int argc, const napi_value* argv, const char* methodName, bool tryCatch)
{
if (workerEnv_ == nullptr) {
return false;
}
napi_value callback = NapiValueHelp::GetNamePropertyInParentPort(workerEnv_, parentPort_, methodName);
bool isCallable = NapiValueHelp::IsCallable(workerEnv_, callback);
if (!isCallable) {
@@ -1002,6 +1100,10 @@ void Worker::CloseWorkerCallback()
void Worker::CallMainFunction(int argc, const napi_value* argv, const char* methodName) const
{
if (mainEnv_ == nullptr || MainIsStop()) {
HILOG_ERROR("worker:: main thread maybe is over");
return;
}
napi_value callback = nullptr;
napi_value obj = nullptr;
napi_get_reference_value(mainEnv_, workerWrapper_, &obj);
@@ -1014,4 +1116,312 @@ void Worker::CallMainFunction(int argc, const napi_value* argv, const char* meth
napi_value callbackResult = nullptr;
napi_call_function(mainEnv_, obj, callback, argc, argv, &callbackResult);
}
void Worker::ReleaseWorkerThreadContent()
{
// 1. remove worker instance count
{
std::lock_guard<std::mutex> lock(g_workersMutex);
std::list<Worker*>::iterator it = std::find(g_workers.begin(), g_workers.end(), this);
if (it != g_workers.end()) {
g_workers.erase(it);
}
}
ParentPortRemoveAllListenerInner();
// 2. delete worker's parentPort
napi_delete_reference(workerEnv_, parentPort_);
parentPort_ = nullptr;
// 3. clear message send to worker thread
workerMessageQueue_.Clear(workerEnv_);
// 4. delete NativeEngine created in worker thread
CloseHelp::DeletePointer(reinterpret_cast<NativeEngine*>(workerEnv_), false);
workerEnv_ = nullptr;
}
void Worker::ReleaseMainThreadContent()
{
// 1. clear message send to main thread
mainMessageQueue_.Clear(mainEnv_);
// 2. clear error queue send to main thread
errorQueue_.Clear(mainEnv_);
if (!MainIsStop()) {
// 3. set thisVar's nativepointer be null
napi_value thisVar = nullptr;
napi_get_reference_value(mainEnv_, workerWrapper_, &thisVar);
Worker* worker = nullptr;
napi_remove_wrap(mainEnv_, thisVar, (void**)&worker);
// 4. set workerWrapper_ be null
napi_delete_reference(mainEnv_, workerWrapper_);
}
mainEnv_ = nullptr;
workerWrapper_ = nullptr;
}
napi_value Worker::ParentPortAddEventListener(napi_env env, napi_callback_info cbinfo)
{
size_t argc = NapiValueHelp::GetCallbackInfoArgc(env, cbinfo);
if (argc < WORKERPARAMNUM) {
napi_throw_error(env, nullptr, "Worker param count must be more than WORKPARAMNUM with on");
return nullptr;
}
napi_value* args = new napi_value[argc];
[[maybe_unused]] ObjectScope<napi_value> scope(args, true);
Worker* worker = nullptr;
napi_get_cb_info(env, cbinfo, &argc, args, nullptr, (void**)&worker);
if (!NapiValueHelp::IsString(args[0])) {
napi_throw_error(env, nullptr, "Worker 1st param must be string with on");
return nullptr;
}
if (!NapiValueHelp::IsCallable(env, args[1])) {
napi_throw_error(env, nullptr, "Worker 2st param must be callable with on");
return nullptr;
}
if (worker == nullptr) {
HILOG_ERROR("worker:: when post message to main occur worker is nullptr");
return nullptr;
}
if (!worker->IsRunning()) {
// if worker is not running, don't send any message to main thread
HILOG_INFO("worker:: when post message to main occur worker is not in running.");
return nullptr;
}
auto listener = new WorkerListener(worker, PERMANENT);
if (argc > WORKERPARAMNUM && NapiValueHelp::IsObject(args[WORKERPARAMNUM])) {
napi_value onceValue = nullptr;
napi_get_named_property(env, args[WORKERPARAMNUM], "once", &onceValue);
bool isOnce = false;
napi_get_value_bool(env, onceValue, &isOnce);
if (isOnce) {
listener->SetMode(ONCE);
}
}
listener->SetCallable(env, args[1]);
char* typeStr = NapiValueHelp::GetString(env, args[0]);
if (typeStr == nullptr) {
CloseHelp::DeletePointer(listener, false);
napi_throw_error(env, nullptr, "worker listener type create error, please check.");
return nullptr;
}
worker->ParentPortAddListenerInner(env, typeStr, listener);
CloseHelp::DeletePointer(typeStr, true);
return NapiValueHelp::GetUndefinedValue(env);
}
napi_value Worker::ParentPortRemoveAllListener(napi_env env, napi_callback_info cbinfo)
{
Worker* worker = nullptr;
napi_get_cb_info(env, cbinfo, nullptr, nullptr, nullptr, (void**)&worker);
if (worker == nullptr) {
HILOG_ERROR("worker:: when post message to main occur worker is nullptr");
return nullptr;
}
if (!worker->IsRunning()) {
// if worker is not running, don't send any message to main thread
HILOG_INFO("worker:: when post message to main occur worker is not in running.");
return nullptr;
}
worker->ParentPortRemoveAllListenerInner();
return NapiValueHelp::GetUndefinedValue(env);
}
napi_value Worker::ParentPortDispatchEvent(napi_env env, napi_callback_info cbinfo)
{
size_t argc = NapiValueHelp::GetCallbackInfoArgc(env, cbinfo);
if (argc < 1) {
napi_throw_error(env, nullptr, "worker:: DispatchEvent param count must be more than 1");
return NapiValueHelp::GetBooleanValue(env, false);
}
napi_value* args = new napi_value[argc];
[[maybe_unused]] ObjectScope<napi_value> scope(args, true);
Worker* worker = nullptr;
napi_get_cb_info(env, cbinfo, &argc, args, nullptr, (void**)&worker);
if (!NapiValueHelp::IsObject(args[0])) {
napi_throw_error(env, nullptr, "worker DispatchEvent 1st param must be Event");
return NapiValueHelp::GetBooleanValue(env, false);
}
napi_value typeValue = nullptr;
napi_get_named_property(env, args[0], "type", &typeValue);
if (!NapiValueHelp::IsString(typeValue)) {
napi_throw_error(env, nullptr, "worker event type must be string");
return NapiValueHelp::GetBooleanValue(env, false);
}
if (worker == nullptr) {
HILOG_ERROR("worker:: when post message to main occur worker is nullptr");
return NapiValueHelp::GetBooleanValue(env, false);
}
if (!worker->IsRunning()) {
// if worker is not running, don't send any message to main thread
HILOG_INFO("worker:: when post message to main occur worker is not in running.");
return NapiValueHelp::GetBooleanValue(env, false);
}
napi_value argv[1] = { args[0] };
char* typeStr = NapiValueHelp::GetString(env, typeValue);
if (typeStr == nullptr) {
napi_throw_error(env, nullptr, "worker listener type create error, please check.");
return NapiValueHelp::GetBooleanValue(env, false);
}
napi_value obj = nullptr;
napi_get_reference_value(env, worker->parentPort_, &obj);
if (strcmp(typeStr, "error") == 0) {
CallWorkCallback(env, obj, 1, argv, "onerror");
} else if (strcmp(typeStr, "messageerror") == 0) {
CallWorkCallback(env, obj, 1, argv, "onmessageerror");
} else if (strcmp(typeStr, "message") == 0) {
CallWorkCallback(env, obj, 1, argv, "onmessage");
}
worker->ParentPortHandleEventListeners(env, obj, 1, argv, typeStr);
CloseHelp::DeletePointer(typeStr, true);
return NapiValueHelp::GetBooleanValue(env, true);
}
napi_value Worker::ParentPortRemoveEventListener(napi_env env, napi_callback_info cbinfo)
{
size_t argc = NapiValueHelp::GetCallbackInfoArgc(env, cbinfo);
if (argc < 1) {
napi_throw_error(env, nullptr, "Worker param count must be more than 2 with on");
return nullptr;
}
napi_value* args = new napi_value[argc];
[[maybe_unused]] ObjectScope<napi_value> scope(args, true);
Worker* worker = nullptr;
napi_get_cb_info(env, cbinfo, &argc, args, nullptr, (void**)&worker);
if (!NapiValueHelp::IsString(args[0])) {
napi_throw_error(env, nullptr, "Worker 1st param must be string with on");
return nullptr;
}
if (argc > 1 && !NapiValueHelp::IsCallable(env, args[1])) {
napi_throw_error(env, nullptr, "Worker 2st param must be callable with on");
return nullptr;
}
if (worker == nullptr) {
HILOG_ERROR("worker:: when post message to main occur worker is nullptr");
return nullptr;
}
if (!worker->IsRunning()) {
// if worker is not running, don't send any message to main thread
HILOG_INFO("worker:: when post message to main occur worker is not in running.");
return nullptr;
}
napi_ref callback = nullptr;
if (argc > 1 && NapiValueHelp::IsCallable(env, args[1])) {
napi_create_reference(env, args[1], 1, &callback);
}
char* typeStr = NapiValueHelp::GetString(env, args[0]);
if (typeStr == nullptr) {
napi_throw_error(env, nullptr, "worker listener type create error, please check.");
return nullptr;
}
worker->ParentPortRemoveListenerInner(env, typeStr, callback);
CloseHelp::DeletePointer(typeStr, true);
napi_delete_reference(env, callback);
return NapiValueHelp::GetUndefinedValue(env);
}
void Worker::ParentPortAddListenerInner(napi_env env, const char* type, const WorkerListener* listener)
{
std::string typestr(type);
auto iter = parentPortEventListeners_.find(typestr);
if (iter == parentPortEventListeners_.end()) {
std::list<WorkerListener*> listeners;
listeners.emplace_back(const_cast<WorkerListener*>(listener));
parentPortEventListeners_[typestr] = listeners;
} else {
std::list<WorkerListener*>& listenerList = iter->second;
std::list<WorkerListener*>::iterator it = std::find_if(
listenerList.begin(), listenerList.end(), Worker::FindWorkerListener(env, listener->callback_));
if (it != listenerList.end()) {
return;
}
listenerList.emplace_back(const_cast<WorkerListener*>(listener));
}
}
void Worker::ParentPortRemoveAllListenerInner()
{
for (auto iter = parentPortEventListeners_.begin(); iter != parentPortEventListeners_.end(); iter++) {
std::list<WorkerListener*>& listeners = iter->second;
for (auto item = listeners.begin(); item != listeners.end(); item++) {
WorkerListener* listener = *item;
CloseHelp::DeletePointer(listener, false);
}
}
parentPortEventListeners_.clear();
}
void Worker::ParentPortRemoveListenerInner(napi_env env, const char* type, napi_ref callback)
{
std::string typestr(type);
auto iter = parentPortEventListeners_.find(typestr);
if (iter == parentPortEventListeners_.end()) {
return;
}
std::list<WorkerListener*>& listenerList = iter->second;
if (callback != nullptr) {
std::list<WorkerListener*>::iterator it =
std::find_if(listenerList.begin(), listenerList.end(), Worker::FindWorkerListener(env, callback));
if (it != listenerList.end()) {
CloseHelp::DeletePointer(*it, false);
listenerList.erase(it);
}
} else {
for (auto it = listenerList.begin(); it != listenerList.end(); it++) {
CloseHelp::DeletePointer(*it, false);
}
parentPortEventListeners_.erase(typestr);
}
}
void Worker::ParentPortHandleEventListeners(napi_env env, napi_value recv,
size_t argc, const napi_value* argv, const char* type)
{
std::string listener(type);
auto iter = parentPortEventListeners_.find(listener);
if (iter == parentPortEventListeners_.end()) {
HILOG_INFO("worker:: there is no listener for type %{public}s", type);
return;
}
std::list<WorkerListener*>& listeners = iter->second;
std::list<WorkerListener*>::iterator it = listeners.begin();
while (it != listeners.end()) {
WorkerListener* data = *it++;
napi_value callbackObj = nullptr;
napi_get_reference_value(env, data->callback_, &callbackObj);
napi_value callbackResult = nullptr;
napi_call_function(env, recv, callbackObj, argc, argv, &callbackResult);
if (!data->NextIsAvailable()) {
listeners.remove(data);
CloseHelp::DeletePointer(data, false);
}
}
}
} // namespace OHOS::CCRuntime::Worker
+40 -6
View File
@@ -34,7 +34,7 @@ public:
static const int8_t WORKERPARAMNUM = 2;
enum RunnerState { STARTING, RUNNING, TERMINATEING, TERMINATED };
enum MainState { ACTIVE, INACTIVE };
enum ListenerMode { ONCE, PERMANENT };
enum ScriptMode { CLASSIC, MODULE };
@@ -100,7 +100,7 @@ public:
static void MainOnError(const uv_async_t* req);
static void WorkerOnMessage(const uv_async_t* req);
static void ExecuteInThread(const void* data);
static void PrepareForWorkerInstance(const Worker* worker);
static bool PrepareForWorkerInstance(const Worker* worker);
static napi_value PostMessage(napi_env env, napi_callback_info cbinfo);
static napi_value PostMessageToMain(napi_env env, napi_callback_info cbinfo);
@@ -120,9 +120,18 @@ public:
static napi_value WorkerConstructor(napi_env env, napi_callback_info cbinfo);
static napi_value InitWorker(napi_env env, napi_value exports);
static napi_value CancelTask(napi_env env, napi_callback_info cbinfo);
static napi_value ParentPortCancelTask(napi_env env, napi_callback_info cbinfo);
static napi_value ParentPortAddEventListener(napi_env env, napi_callback_info cbinfo);
static napi_value ParentPortRemoveAllListener(napi_env env, napi_callback_info cbinfo);
static napi_value ParentPortDispatchEvent(napi_env env, napi_callback_info cbinfo);
static napi_value ParentPortRemoveEventListener(napi_env env, napi_callback_info cbinfo);
void StartExecuteInThread(napi_env env, const char* script);
bool UpdateWorkerState(RunnerState state);
bool UpdateMainState(MainState state);
bool IsRunning() const
{
@@ -179,9 +188,13 @@ public:
return nullptr;
}
bool IsSameWorkerEnv(napi_env env) const
bool ClearWorkerTasks()
{
return workerEnv_ == env;
if (mainEnv_ != nullptr) {
workerMessageQueue_.Clear(mainEnv_);
return true;
}
return false;
}
void TriggerPostTask()
@@ -191,10 +204,20 @@ public:
}
}
bool MainIsStop() const
{
return mainState_.load(std::memory_order_acquire) == INACTIVE;
}
bool IsSameWorkerEnv(napi_env env) const
{
return workerEnv_ == env;
}
void Loop()
{
if (workerEnv_ != nullptr) {
return reinterpret_cast<NativeEngine*>(workerEnv_)->Loop(LOOP_DEFAULT);
reinterpret_cast<NativeEngine*>(workerEnv_)->Loop(LOOP_DEFAULT);
}
}
@@ -211,6 +234,8 @@ private:
void CallMainFunction(int argc, const napi_value* argv, const char* methodName) const;
void HandleEventListeners(napi_env env, napi_value recv, size_t argc, const napi_value* argv, const char* type);
void ParentPortHandleEventListeners(napi_env env, napi_value recv,
size_t argc, const napi_value* argv, const char* type);
void TerminateInner();
void PostMessageInner(MessageDataType data);
@@ -223,6 +248,13 @@ private:
void CloseWorkerCallback();
void CloseMainCallback() const;
void ReleaseWorkerThreadContent();
void ReleaseMainThreadContent();
void ParentPortAddListenerInner(napi_env env, const char* type, const WorkerListener* listener);
void ParentPortRemoveAllListenerInner();
void ParentPortRemoveListenerInner(napi_env env, const char* type, napi_ref callback);
napi_env GetMainEnv() const
{
return mainEnv_;
@@ -246,6 +278,7 @@ private:
uv_async_t mainOnErrorSignal_ {};
std::atomic<RunnerState> runnerState_ {STARTING};
std::atomic<MainState> mainState_ {ACTIVE};
std::unique_ptr<WorkerRunner> runner_ {};
napi_env mainEnv_ {nullptr};
@@ -255,8 +288,9 @@ private:
napi_ref parentPort_ {nullptr};
std::map<std::string, std::list<WorkerListener*>> eventListeners_ {};
std::map<std::string, std::list<WorkerListener*>> parentPortEventListeners_ {};
std::mutex workerAsyncMutex_ {};
std::recursive_mutex liveStatusLock_ {};
};
} // namespace OHOS::CCRuntime::Worker
#endif // FOUNDATION_CCRUNTIME_JSAPI_WORKER_H