Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions test-app/app/src/main/assets/app/mainpage.js
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand Down
Original file line number Diff line number Diff line change
@@ -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);
})();
});
});
Original file line number Diff line number Diff line change
@@ -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);
}
};
12 changes: 7 additions & 5 deletions test-app/runtime/src/main/cpp/WorkerInspectorClient.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down Expand Up @@ -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) {
Expand All @@ -225,7 +227,7 @@ void WorkerInspectorClient::runMessageLoopOnPause(int contextGroupId) {
->GetEventLoop(isolate_)
->RunNestableV8Tasks();

if (shouldWait && !pauseTerminated_ && !dying_) {
if (shouldWait && !pauseTerminated_ && !dying_ && !workerTerminating_) {
std::unique_lock<std::mutex> lock(messageArrivedMutex_);
messageArrived_.wait_for(lock, std::chrono::milliseconds(1));
}
Expand Down
8 changes: 7 additions & 1 deletion test-app/runtime/src/main/cpp/WorkerInspectorClient.h
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -126,6 +128,10 @@ class WorkerInspectorClient final : public v8_inspector::V8InspectorClient,

std::atomic<bool> dying_{false};
std::atomic<bool> 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<bool> runningPauseLoop_{false}; // written on the worker thread only
bool pendingReset_ = false; // worker thread only
};
Expand Down
39 changes: 25 additions & 14 deletions test-app/runtime/src/main/cpp/WorkerWrapper.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -142,8 +142,7 @@ void WorkerWrapper::PostMessageToParent(std::shared_ptr<worker::Message> message
}

void WorkerWrapper::Terminate() {
if (isClosing_ || isDisposed_) {
// The worker is already shutting down on its own; nothing to do.
if (isDisposed_) {
return;
}

Expand All @@ -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<std::mutex> 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();
}
}
}

Expand Down Expand Up @@ -656,7 +662,11 @@ void WorkerWrapper::BackgroundLooper(std::shared_ptr<WorkerWrapper> 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<std::mutex> lock(workerIsolateMutex_);
workerIsolate_.store(nullptr);
}
Isolate* isolate = runtime_->GetIsolate();
{
v8::Locker locker(isolate);
Expand Down Expand Up @@ -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<std::mutex> lock(inspectorMutex_);
inspector_ = client;
Expand Down
4 changes: 4 additions & 0 deletions test-app/runtime/src/main/cpp/WorkerWrapper.h
Original file line number Diff line number Diff line change
Expand Up @@ -178,7 +178,11 @@ class WorkerWrapper : public std::enable_shared_from_this<WorkerWrapper> {
// 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<EventLoop> 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<v8::Isolate*> workerIsolate_;
std::mutex workerIsolateMutex_;
Runtime* runtime_;

const int workerId_;
Expand Down
Loading