Skip to content

Commit d8a79bf

Browse files
worker: emit worker exit notifications on BroadcastChannel
Expose worker termination notifications through BroadcastChannel so consumers can observe when a worker exits andinspect its thread ID and exit code. Fixes: #59053 Signed-off-by: SudhansuBandha <bandhasudhansu@gmail.com>
1 parent c66ae4f commit d8a79bf

5 files changed

Lines changed: 156 additions & 1 deletion

File tree

‎doc/api/worker_threads.md‎

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -987,6 +987,25 @@ added: v15.4.0
987987
* Type: {Function} Invoked with a received message cannot be
988988
deserialized.
989989
990+
### `broadcastChannel.onworkerexited`
991+
992+
* Type: {Function} Invoked when worker associated with the
993+
`BroadcastChannel` terminates.
994+
995+
The callback receives an object with the following properties:
996+
997+
* `threadId` {number} The ID of the worker thread that terminated.
998+
* `exitCode` {number} The exit code with which the worker terminated.
999+
1000+
The `exitCode` is the value passed to `process.exit()` when the worker
1001+
explicitly exits. If the worker terminates without explicitly specifying
1002+
an exit code, the corresponding exit code is reported.
1003+
1004+
The `workerexited` event is emitted only when the worker's execution
1005+
environment is stopping. Closing a `BroadcastChannel` or its underlying
1006+
`MessagePort` does not by itself indicate that a worker has exited and
1007+
does not emit this event.
1008+
9901009
### `broadcastChannel.postMessage(message)`
9911010
9921011
<!-- YAML

‎lib/internal/worker/io.js‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -74,6 +74,7 @@ const kIncrementsPortRef = Symbol('kIncrementsPortRef');
7474
const kName = Symbol('kName');
7575
const kOnMessage = Symbol('kOnMessage');
7676
const kOnMessageError = Symbol('kOnMessageError');
77+
const kOnWorkerExited = Symbol('kOnWorkerExited');
7778
const kPort = Symbol('kPort');
7879
const kWaitingStreams = Symbol('kWaitingStreams');
7980
const kWritableCallback = Symbol('kWritableCallback');
@@ -367,8 +368,10 @@ class BroadcastChannel extends EventTarget {
367368
this[kOnMessage] = FunctionPrototypeBind(onMessageEvent, this, 'message');
368369
this[kOnMessageError] =
369370
FunctionPrototypeBind(onMessageEvent, this, 'messageerror');
371+
this[kOnWorkerExited] = FunctionPrototypeBind(onMessageEvent, this, 'workerexited');
370372
this[kHandle].on('message', this[kOnMessage]);
371373
this[kHandle].on('messageerror', this[kOnMessageError]);
374+
this[kHandle].on('workerexited', this[kOnWorkerExited]);
372375
}
373376

374377
[inspect.custom](depth, options) {
@@ -407,8 +410,10 @@ class BroadcastChannel extends EventTarget {
407410
return;
408411
this[kHandle].off('message', this[kOnMessage]);
409412
this[kHandle].off('messageerror', this[kOnMessageError]);
413+
this[kHandle].off('workerexited', this[kOnWorkerExited]);
410414
this[kOnMessage] = undefined;
411415
this[kOnMessageError] = undefined;
416+
this[kOnWorkerExited] = undefined;
412417
this[kHandle].close();
413418
this[kHandle] = undefined;
414419
}
@@ -468,6 +473,7 @@ ObjectDefineProperties(BroadcastChannel.prototype, {
468473

469474
defineEventHandler(BroadcastChannel.prototype, 'message');
470475
defineEventHandler(BroadcastChannel.prototype, 'messageerror');
476+
defineEventHandler(BroadcastChannel.prototype, 'workerexited');
471477

472478
function markAsUncloneable(obj) {
473479
if ((typeof obj !== 'object' && typeof obj !== 'function') || obj === null) {

‎src/node_messaging.cc‎

Lines changed: 88 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -648,6 +648,32 @@ void MessagePortData::AddToIncomingQueue(std::shared_ptr<Message> message) {
648648
}
649649
}
650650

651+
void MessagePortData::AddWorkerExitNotification(uint64_t thread_id,
652+
ExitCode exit_code) {
653+
Mutex::ScopedLock lock(mutex_);
654+
worker_exit_notifications_.emplace_back(WorkerExitNotification{
655+
thread_id,
656+
exit_code,
657+
});
658+
659+
if (owner_ != nullptr) {
660+
Debug(owner_, "Adding worker-exit notification");
661+
owner_->TriggerAsync();
662+
}
663+
}
664+
665+
bool MessagePortData::GetWorkerExitNotification(
666+
WorkerExitNotification* notification) {
667+
Mutex::ScopedLock lock(mutex_);
668+
669+
if (worker_exit_notifications_.empty()) return false;
670+
671+
*notification = worker_exit_notifications_.front();
672+
worker_exit_notifications_.pop_front();
673+
674+
return true;
675+
}
676+
651677
void MessagePortData::Entangle(MessagePortData* a, MessagePortData* b) {
652678
auto group = std::make_shared<SiblingGroup>();
653679
group->Entangle({a, b});
@@ -823,6 +849,7 @@ void MessagePort::OnMessage(MessageProcessingMode mode) {
823849
HandleScope handle_scope(env()->isolate());
824850
Local<Context> context =
825851
object(env()->isolate())->GetCreationContextChecked();
852+
Local<Function> emit_message = PersistentToLocal::Strong(emit_message_fn_);
826853

827854
size_t processing_limit;
828855
if (mode == MessageProcessingMode::kNormalOperation) {
@@ -850,9 +877,43 @@ void MessagePort::OnMessage(MessageProcessingMode mode) {
850877
return;
851878
}
852879

880+
MessagePortData::WorkerExitNotification worker_exit;
881+
882+
if (data_->GetWorkerExitNotification(&worker_exit)) {
883+
Debug(this,
884+
"Worker exited: thread_id=%" PRIu64 ", exit_code=%d",
885+
worker_exit.thread_id,
886+
static_cast<int>(worker_exit.exit_code));
887+
888+
Local<Object> exit_info = Object::New(env()->isolate());
889+
890+
exit_info
891+
->Set(context,
892+
FIXED_ONE_BYTE_STRING(env()->isolate(), "threadId"),
893+
v8::Number::New(env()->isolate(), worker_exit.thread_id))
894+
.Check();
895+
896+
exit_info
897+
->Set(context,
898+
FIXED_ONE_BYTE_STRING(env()->isolate(), "exitCode"),
899+
v8::Integer::New(env()->isolate(),
900+
static_cast<int>(worker_exit.exit_code)))
901+
.Check();
902+
903+
Local<Value> argv[3];
904+
argv[0] = exit_info;
905+
argv[1] = Undefined(env()->isolate());
906+
argv[2] = FIXED_ONE_BYTE_STRING(env()->isolate(), "workerexited");
907+
908+
if (MakeCallback(emit_message, arraysize(argv), argv).IsEmpty()) {
909+
if (data_) TriggerAsync();
910+
return;
911+
}
912+
continue;
913+
}
914+
853915
HandleScope handle_scope(env()->isolate());
854916
Context::Scope context_scope(context);
855-
Local<Function> emit_message = PersistentToLocal::Strong(emit_message_fn_);
856917

857918
Local<Value> payload;
858919
Local<Value> port_list = Undefined(env()->isolate());
@@ -901,6 +962,20 @@ void MessagePort::OnMessage(MessageProcessingMode mode) {
901962
void MessagePort::OnClose() {
902963
Debug(this, "MessagePort::OnClose()");
903964
if (data_) {
965+
Environment* environment = env();
966+
if (environment->is_stopping()) {
967+
const uint64_t thread_id = environment->thread_id();
968+
const ExitCode exit_code = environment->exit_code(ExitCode::kNoFailure);
969+
970+
Debug(this,
971+
"Worker exiting: thread_id=%" PRIu64 ", exit_code=%d",
972+
thread_id,
973+
static_cast<int>(exit_code));
974+
975+
if (data_->group_) {
976+
data_->group_->NotifyWorkerExit(data_.get(), thread_id, exit_code);
977+
}
978+
}
904979
// Detach() returns move(data_).
905980
Detach()->Disentangle();
906981
}
@@ -1587,6 +1662,18 @@ void SiblingGroup::Disentangle(MessagePortData* data) {
15871662
(*(ports_.begin()))->AddToIncomingQueue(std::make_shared<Message>());
15881663
}
15891664

1665+
void SiblingGroup::NotifyWorkerExit(MessagePortData* exiting_port,
1666+
uint64_t thread_id,
1667+
ExitCode exit_code) {
1668+
RwLock::ScopedReadLock lock(group_mutex_);
1669+
1670+
for (MessagePortData* port : ports_) {
1671+
if (port == exiting_port) continue;
1672+
1673+
port->AddWorkerExitNotification(thread_id, exit_code);
1674+
}
1675+
}
1676+
15901677
SiblingGroup::Map SiblingGroup::groups_;
15911678
Mutex SiblingGroup::groups_mutex_;
15921679

‎src/node_messaging.h‎

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -149,6 +149,10 @@ class SiblingGroup final : public std::enable_shared_from_this<SiblingGroup> {
149149
void Entangle(std::initializer_list<MessagePortData*> data);
150150
void Disentangle(MessagePortData* data);
151151

152+
void NotifyWorkerExit(MessagePortData* exiting_port,
153+
uint64_t thread_id,
154+
ExitCode exit_code);
155+
152156
const std::string& name() const { return name_; }
153157

154158
size_t size() const { return ports_.size(); }
@@ -185,6 +189,9 @@ class MessagePortData : public TransferData {
185189
v8::Maybe<bool> Dispatch(
186190
std::shared_ptr<Message> message,
187191
std::string* error = nullptr);
192+
193+
// Internal worker-exit notification.
194+
void AddWorkerExitNotification(uint64_t thread_id, ExitCode exit_code);
188195

189196
// Turns `a` and `b` into siblings, i.e. connects the sending side of one
190197
// to the receiving side of the other. This is not thread-safe.
@@ -213,6 +220,15 @@ class MessagePortData : public TransferData {
213220
// once that is available with C++17, because std::shared_ptr comes with
214221
// overhead that is only necessary for BroadcastChannel.
215222
std::deque<std::shared_ptr<Message>> incoming_messages_;
223+
struct WorkerExitNotification {
224+
uint64_t thread_id;
225+
ExitCode exit_code;
226+
};
227+
228+
bool GetWorkerExitNotification(WorkerExitNotification* notification);
229+
230+
std::deque<WorkerExitNotification> worker_exit_notifications_;
231+
216232
MessagePort* owner_ = nullptr;
217233
std::shared_ptr<SiblingGroup> group_;
218234
friend class MessagePort;

‎test/parallel/test-worker-broadcastchannel.js‎

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -183,3 +183,30 @@ assert.throws(() => new BroadcastChannel(), {
183183
"BroadcastChannel { name: 'channel5', active: false }"
184184
);
185185
}
186+
187+
{
188+
const bc = new BroadcastChannel('channel6');
189+
190+
const worker = new Worker(`
191+
const { BroadcastChannel } = require('worker_threads');
192+
193+
const bc = new BroadcastChannel('channel6');
194+
195+
// Keep the BroadcastChannel alive long enough for the exit
196+
// notification to be observed by the parent.
197+
setImmediate(() => {
198+
process.exit(42);
199+
});
200+
`, { eval: true });
201+
202+
bc.onworkerexited = common.mustCall((event) => {
203+
assert.strictEqual(event.data.threadId, worker.threadId);
204+
assert.strictEqual(event.data.exitCode, 42);
205+
206+
bc.close();
207+
});
208+
209+
worker.on('exit', common.mustCall((exitCode) => {
210+
assert.strictEqual(exitCode, 42);
211+
}));
212+
}

0 commit comments

Comments
 (0)