diff --git a/jsapi/worker/worker.cpp b/jsapi/worker/worker.cpp index 9b013d9..8993a0d 100644 --- a/jsapi/worker/worker.cpp +++ b/jsapi/worker/worker.cpp @@ -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 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))->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(const_cast(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 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(workerEnv)->MarkSubThread(); + worker->SetWorkerEnv(workerEnv); } - // mark worker env is subThread - reinterpret_cast(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(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(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 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 lock(g_workersMutex); - std::list::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 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 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 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 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(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 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 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 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 lock(g_workersMutex); + std::list::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(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 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 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 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 listeners; + listeners.emplace_back(const_cast(listener)); + parentPortEventListeners_[typestr] = listeners; + } else { + std::list& listenerList = iter->second; + std::list::iterator it = std::find_if( + listenerList.begin(), listenerList.end(), Worker::FindWorkerListener(env, listener->callback_)); + if (it != listenerList.end()) { + return; + } + listenerList.emplace_back(const_cast(listener)); + } +} + +void Worker::ParentPortRemoveAllListenerInner() +{ + for (auto iter = parentPortEventListeners_.begin(); iter != parentPortEventListeners_.end(); iter++) { + std::list& 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& listenerList = iter->second; + if (callback != nullptr) { + std::list::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& listeners = iter->second; + std::list::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 diff --git a/jsapi/worker/worker.h b/jsapi/worker/worker.h index 212bb93..9a60b60 100644 --- a/jsapi/worker/worker.h +++ b/jsapi/worker/worker.h @@ -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(workerEnv_)->Loop(LOOP_DEFAULT); + reinterpret_cast(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_ {STARTING}; + std::atomic mainState_ {ACTIVE}; std::unique_ptr runner_ {}; napi_env mainEnv_ {nullptr}; @@ -255,8 +288,9 @@ private: napi_ref parentPort_ {nullptr}; std::map> eventListeners_ {}; + std::map> parentPortEventListeners_ {}; - std::mutex workerAsyncMutex_ {}; + std::recursive_mutex liveStatusLock_ {}; }; } // namespace OHOS::CCRuntime::Worker #endif // FOUNDATION_CCRUNTIME_JSAPI_WORKER_H \ No newline at end of file