diff --git a/test-app/app/src/main/assets/app/mainpage.js b/test-app/app/src/main/assets/app/mainpage.js index 62aeb52eb..7b56d7ba8 100644 --- a/test-app/app/src/main/assets/app/mainpage.js +++ b/test-app/app/src/main/assets/app/mainpage.js @@ -27,6 +27,7 @@ require("./tests/testWebAssembly"); require("./tests/testEventLoop"); require("./tests/testMultithreadedJavascript"); require("./tests/testWorkerTerminateDuringLoad"); +require("./tests/testWorkerTerminateAfterClose"); require("./tests/testWorkerOptions"); require("./tests/testWorkerResourceLimits"); require("./tests/testInterfaceDefaultMethods"); diff --git a/test-app/app/src/main/assets/app/tests/testWorkerTerminateAfterClose.js b/test-app/app/src/main/assets/app/tests/testWorkerTerminateAfterClose.js new file mode 100644 index 000000000..4eff24878 --- /dev/null +++ b/test-app/app/src/main/assets/app/tests/testWorkerTerminateAfterClose.js @@ -0,0 +1,56 @@ +describe("Worker terminate after close", function () { + var START_DEADLINE = 3000; + var SETTLE_AFTER = 300; + // [0] counts the worker's loop iterations; [1] stops the loop, so a worker + // that terminate() failed to stop does not outlive the spec. + var shared; + var timers = []; + var originalTimeout; + + beforeEach(function () { + originalTimeout = jasmine.DEFAULT_TIMEOUT_INTERVAL; + jasmine.DEFAULT_TIMEOUT_INTERVAL = 10000; + }); + + afterEach(function () { + timers.forEach(clearTimeout); + timers = []; + if (shared) { + Atomics.store(shared, 1, 1); + shared = null; + } + jasmine.DEFAULT_TIMEOUT_INTERVAL = originalTimeout; + }); + + function later(fn, ms) { + timers.push(setTimeout(fn, ms)); + } + + it("stops a worker that keeps running after it called close()", function (done) { + shared = new Int32Array(new SharedArrayBuffer(8)); + var counter = shared; + var worker = new Worker("./workerCloseThenSpinWorker.js"); + worker.postMessage(counter.buffer); + var started = Date.now(); + + (function waitForSpin() { + if (Atomics.load(counter, 0) === 0) { + if (Date.now() - started > START_DEADLINE) { + expect("the worker never started running").toBeNull(); + done(); + return; + } + later(waitForSpin, 20); + return; + } + worker.terminate(); + later(function () { + var afterTerminate = Atomics.load(counter, 0); + later(function () { + expect(Atomics.load(counter, 0)).toBe(afterTerminate); + done(); + }, SETTLE_AFTER); + }, SETTLE_AFTER); + })(); + }); +}); diff --git a/test-app/app/src/main/assets/app/tests/workerCloseThenSpinWorker.js b/test-app/app/src/main/assets/app/tests/workerCloseThenSpinWorker.js new file mode 100644 index 000000000..00cd2905a --- /dev/null +++ b/test-app/app/src/main/assets/app/tests/workerCloseThenSpinWorker.js @@ -0,0 +1,9 @@ +// close() lets the running callback finish, so this one never returns unless +// terminate() interrupts it, or the spec raises the stop flag to clean up. +onmessage = function (event) { + var shared = new Int32Array(event.data); + close(); + while (Atomics.load(shared, 1) === 0) { + Atomics.add(shared, 0, 1); + } +}; diff --git a/test-app/runtime/src/main/cpp/WorkerInspectorClient.cpp b/test-app/runtime/src/main/cpp/WorkerInspectorClient.cpp index fbe7a7192..09544e425 100644 --- a/test-app/runtime/src/main/cpp/WorkerInspectorClient.cpp +++ b/test-app/runtime/src/main/cpp/WorkerInspectorClient.cpp @@ -34,13 +34,15 @@ std::string ToUtf8String(const StringView& view) { } // namespace WorkerInspectorClient::WorkerInspectorClient(int workerId, Isolate* isolate, ALooper* workerLooper, - const std::string& url) + const std::string& url, + const std::atomic_bool& workerTerminating) : workerId_(workerId), sessionId_("NS_WORKER_" + std::to_string(workerId)), targetId_("ns-worker-" + std::to_string(workerId)), url_(url), isolate_(isolate), - workerLooper_(workerLooper) { + workerLooper_(workerLooper), + workerTerminating_(workerTerminating) { // Wakes the worker looper when CDP messages arrive on the socket thread; // same mechanism as the worker's message inbox (ConcurrentQueue). eventFd_ = eventfd(0, EFD_NONBLOCK | EFD_CLOEXEC); @@ -206,13 +208,13 @@ void WorkerInspectorClient::MaybeResetSession() { } void WorkerInspectorClient::runMessageLoopOnPause(int contextGroupId) { - if (runningPauseLoop_.load(std::memory_order_acquire) || dying_) { + if (runningPauseLoop_.load(std::memory_order_acquire) || dying_ || workerTerminating_) { return; } runningPauseLoop_.store(true, std::memory_order_release); pauseTerminated_ = false; - while (!pauseTerminated_ && !dying_) { + while (!pauseTerminated_ && !dying_ && !workerTerminating_) { std::string message = this->PopMessage(); bool shouldWait = message.empty(); if (!shouldWait) { @@ -225,7 +227,7 @@ void WorkerInspectorClient::runMessageLoopOnPause(int contextGroupId) { ->GetEventLoop(isolate_) ->RunNestableV8Tasks(); - if (shouldWait && !pauseTerminated_ && !dying_) { + if (shouldWait && !pauseTerminated_ && !dying_ && !workerTerminating_) { std::unique_lock lock(messageArrivedMutex_); messageArrived_.wait_for(lock, std::chrono::milliseconds(1)); } diff --git a/test-app/runtime/src/main/cpp/WorkerInspectorClient.h b/test-app/runtime/src/main/cpp/WorkerInspectorClient.h index a3deb572a..3ca4acb3b 100644 --- a/test-app/runtime/src/main/cpp/WorkerInspectorClient.h +++ b/test-app/runtime/src/main/cpp/WorkerInspectorClient.h @@ -35,8 +35,10 @@ class WorkerInspectorClient final : public v8_inspector::V8InspectorClient, public v8_inspector::V8Inspector::Channel { public: // Worker thread, with the worker isolate locked and its context created. + // `workerTerminating` is the worker's own termination flag; the + // worker outlives this client. WorkerInspectorClient(int workerId, v8::Isolate* isolate, ALooper* workerLooper, - const std::string& url); + const std::string& url, const std::atomic_bool& workerTerminating); ~WorkerInspectorClient() override; int WorkerId() const { @@ -126,6 +128,10 @@ class WorkerInspectorClient final : public v8_inspector::V8InspectorClient, std::atomic dying_{false}; std::atomic pauseTerminated_{false}; + // Read by the pause loop, which leaves on it without NotifyTerminating: + // a pause entered while the worker holds its inspector mutex, inside + // console.log, would otherwise keep Terminate() waiting for that mutex. + const std::atomic_bool& workerTerminating_; std::atomic runningPauseLoop_{false}; // written on the worker thread only bool pendingReset_ = false; // worker thread only }; diff --git a/test-app/runtime/src/main/cpp/WorkerWrapper.cpp b/test-app/runtime/src/main/cpp/WorkerWrapper.cpp index 04abc68fd..a097c871e 100644 --- a/test-app/runtime/src/main/cpp/WorkerWrapper.cpp +++ b/test-app/runtime/src/main/cpp/WorkerWrapper.cpp @@ -142,8 +142,7 @@ void WorkerWrapper::PostMessageToParent(std::shared_ptr message } void WorkerWrapper::Terminate() { - if (isClosing_ || isDisposed_) { - // The worker is already shutting down on its own; nothing to do. + if (isDisposed_) { return; } @@ -152,16 +151,23 @@ void WorkerWrapper::Terminate() { return; } - Isolate* isolate = workerIsolate_.load(); - if (isolate != nullptr) { - // The only v8 call that is legal from another thread - interrupts any - // JS currently running on the worker (e.g. a busy loop). - isolate->TerminateExecution(); - // A pump parked with nothing queued runs no JS, so the interrupt - // above never materializes for it - the loop's own flag ends it. - auto loop = NativeScriptPlatform::Instance()->LookupEventLoop(isolate); - if (loop != nullptr) { - loop->NoteTerminationRequested(); + { + // Held across the use, not just the read: the worker thread withdraws + // the isolate under the same mutex before disposing it, so a terminate + // that already read it finishes with it first, and a later one finds + // null. + std::lock_guard lock(workerIsolateMutex_); + Isolate* isolate = workerIsolate_.load(); + if (isolate != nullptr) { + // Legal from any thread: interrupts any JS currently running on the + // worker (e.g. a busy loop, including one that follows close()). + isolate->TerminateExecution(); + // A pump parked with nothing queued runs no JS, so the interrupt + // above never materializes for it - the loop's own flag ends it. + auto loop = NativeScriptPlatform::Instance()->LookupEventLoop(isolate); + if (loop != nullptr) { + loop->NoteTerminationRequested(); + } } } @@ -656,7 +662,11 @@ void WorkerWrapper::BackgroundLooper(std::shared_ptr self) { // bootstrap failed between initWorkerRuntime and the workerIsolate_ // publish (e.g. a JNI error while resolving the looper), the atomic is // still null while the isolate very much needs disposing. - workerIsolate_.store(nullptr); + { + // Waits out a Terminate() on another thread that is still using it. + std::lock_guard lock(workerIsolateMutex_); + workerIsolate_.store(nullptr); + } Isolate* isolate = runtime_->GetIsolate(); { v8::Locker locker(isolate); @@ -808,7 +818,8 @@ void WorkerWrapper::CreateInspector(Isolate* isolate) { ? workerPath_ : "file://" + workerPath_; - auto* client = new WorkerInspectorClient(workerId_, isolate, ALooper_forThread(), url); + auto* client = new WorkerInspectorClient(workerId_, isolate, ALooper_forThread(), url, + isTerminating_); { std::lock_guard lock(inspectorMutex_); inspector_ = client; diff --git a/test-app/runtime/src/main/cpp/WorkerWrapper.h b/test-app/runtime/src/main/cpp/WorkerWrapper.h index 5dd3b2fa8..b79941d87 100644 --- a/test-app/runtime/src/main/cpp/WorkerWrapper.h +++ b/test-app/runtime/src/main/cpp/WorkerWrapper.h @@ -178,7 +178,11 @@ class WorkerWrapper : public std::enable_shared_from_this { // The parent runtime's task queue; weak so a child outliving its parent // just drops its posts instead of touching a dead runtime. std::weak_ptr parentTasks_; + // Written by the worker thread only: published once the runtime is up, + // withdrawn before the isolate is disposed. Any other thread reads and uses + // it under workerIsolateMutex_. std::atomic workerIsolate_; + std::mutex workerIsolateMutex_; Runtime* runtime_; const int workerId_;