Revert "LibIPC: Move message decoding from main thread to I/O thread"

This reverts commit 757795ada4.

Appears to have regressed WPT.
This commit is contained in:
Andreas Kling 2026-01-25 12:19:53 +01:00
parent 81942a84f3
commit a0c389846e
6 changed files with 103 additions and 265 deletions

View file

@ -7,7 +7,6 @@
*/
#include <AK/Vector.h>
#include <LibCore/EventLoop.h>
#include <LibCore/Socket.h>
#include <LibIPC/Connection.h>
#include <LibIPC/Message.h>
@ -20,27 +19,14 @@ ConnectionBase::ConnectionBase(IPC::Stub& local_stub, NonnullOwnPtr<Transport> t
, m_transport(move(transport))
, m_local_endpoint_magic(local_endpoint_magic)
{
}
void ConnectionBase::initialize_messaging()
{
m_event_loop = Core::EventLoop::current_weak();
m_transport->set_message_handler([this](NonnullOwnPtr<IPC::Message> message) {
on_message_received(move(message));
});
m_transport->set_peer_closed_handler([this] {
on_peer_closed();
m_transport->set_up_read_hook([this] {
NonnullRefPtr protect = *this;
drain_messages_from_peer();
handle_messages();
});
}
ConnectionBase::~ConnectionBase()
{
// Close the transport before destroying member variables (especially the
// condition variable). This ensures the I/O thread is stopped and joined
// before we destroy state it may be accessing.
m_transport->close();
}
ConnectionBase::~ConnectionBase() = default;
bool ConnectionBase::is_open() const
{
@ -67,7 +53,6 @@ ErrorOr<void> ConnectionBase::post_message(MessageBuffer buffer)
void ConnectionBase::shutdown()
{
m_transport->close();
m_unprocessed_messages_cv.broadcast();
die();
}
@ -77,45 +62,9 @@ void ConnectionBase::shutdown_with_error(Error const& error)
shutdown();
}
void ConnectionBase::on_message_received(NonnullOwnPtr<IPC::Message> message)
{
// Called from I/O thread - store message and signal waiters
{
Threading::MutexLocker lock(m_unprocessed_messages_mutex);
m_unprocessed_messages.append(move(message));
m_unprocessed_messages_cv.broadcast();
}
// Wake up the main thread's event loop to process messages.
if (m_event_loop) {
NonnullRefPtr<ConnectionBase> strong_this = *this;
m_event_loop->deferred_invoke([strong_this = move(strong_this)] {
strong_this->handle_messages();
});
}
}
void ConnectionBase::on_peer_closed()
{
m_peer_closed.store(true, AK::MemoryOrder::memory_order_release);
m_unprocessed_messages_cv.broadcast();
if (m_event_loop) {
NonnullRefPtr<ConnectionBase> strong_this = *this;
m_event_loop->deferred_invoke([strong_this = move(strong_this)] {
strong_this->shutdown();
});
}
}
void ConnectionBase::handle_messages()
{
Vector<NonnullOwnPtr<Message>> messages;
{
Threading::MutexLocker lock(m_unprocessed_messages_mutex);
messages = move(m_unprocessed_messages);
}
auto messages = move(m_unprocessed_messages);
for (auto& message : messages) {
if (message->endpoint_magic() != m_local_endpoint_magic)
continue;
@ -139,50 +88,80 @@ void ConnectionBase::handle_messages()
}
}
void ConnectionBase::wait_for_transport_to_become_readable()
{
m_transport->wait_until_readable();
}
ConnectionBase::PeerEOF ConnectionBase::drain_messages_from_peer()
{
bool parse_error = false;
auto schedule_shutdown = m_transport->read_as_many_messages_as_possible_without_blocking([&](auto&& raw_message) {
if (auto message = try_parse_message(raw_message.bytes, raw_message.fds)) {
m_unprocessed_messages.append(message.release_nonnull());
} else {
dbgln("Failed to parse IPC message {:hex-dump}", raw_message.bytes);
parse_error = true;
}
});
if (parse_error) {
dbgln("IPC::ConnectionBase ({:p}): Disconnecting misbehaving peer due to malformed message", this);
schedule_shutdown = Transport::ShouldShutdown::Yes;
}
if (!m_unprocessed_messages.is_empty()) {
deferred_invoke([this] {
handle_messages();
});
}
if (schedule_shutdown == Transport::ShouldShutdown::Yes) {
deferred_invoke([this] {
shutdown();
});
return PeerEOF::Yes;
}
return PeerEOF::No;
}
OwnPtr<IPC::Message> ConnectionBase::wait_for_specific_endpoint_message_impl(u32 endpoint_magic, int message_id)
{
{
Threading::MutexLocker lock(m_unprocessed_messages_mutex);
for (;;) {
// Double check we don't already have the message waiting for us.
// Otherwise we might end up blocked for a while for no reason.
for (size_t i = 0; i < m_unprocessed_messages.size(); ++i) {
auto& message = m_unprocessed_messages[i];
if (message->endpoint_magic() != endpoint_magic)
continue;
if (message->message_id() == message_id)
return m_unprocessed_messages.take(i);
}
if (!is_open() || m_peer_closed.load(AK::MemoryOrder::memory_order_acquire))
break;
// Wait for more messages from I/O thread
m_unprocessed_messages_cv.wait();
for (;;) {
// Double check we don't already have the event waiting for us.
// Otherwise we might end up blocked for a while for no reason.
for (size_t i = 0; i < m_unprocessed_messages.size(); ++i) {
auto& message = m_unprocessed_messages[i];
if (message->endpoint_magic() != endpoint_magic)
continue;
if (message->message_id() == message_id)
return m_unprocessed_messages.take(i);
}
if (!is_open())
break;
wait_for_transport_to_become_readable();
if (drain_messages_from_peer() == PeerEOF::Yes)
break;
}
dbgln("Failed to receive message_id: {}", message_id);
bool should_close_transport = false;
{
Threading::MutexLocker lock(m_unprocessed_messages_mutex);
if (!m_unprocessed_messages.is_empty()) {
should_close_transport = true;
dbgln("Transport shutdown with unprocessed messages left: {}", m_unprocessed_messages.size());
for (size_t i = 0; i < m_unprocessed_messages.size(); ++i) {
auto& message = m_unprocessed_messages[i];
dbgln(" Message {:03} is: {:2}-{}", i, message->message_id(), message->message_name());
}
}
}
if (should_close_transport)
if (!m_unprocessed_messages.is_empty()) {
m_transport->close();
dbgln("Handling remaining messages before returning to caller");
handle_messages();
dbgln("Messages handled, returning to caller");
dbgln("Transport shutdown with unprocessed messages left: {}", m_unprocessed_messages.size());
for (size_t i = 0; i < m_unprocessed_messages.size(); ++i) {
auto& message = m_unprocessed_messages[i];
dbgln(" Message {:03} is: {:2}-{}", i, message->message_id(), message->message_name());
}
dbgln("Handling remaining messages before returning to caller");
handle_messages();
dbgln("Messages handled, returning to caller");
}
return {};
}

View file

@ -8,17 +8,13 @@
#pragma once
#include <AK/Atomic.h>
#include <AK/Forward.h>
#include <AK/Queue.h>
#include <LibCore/EventLoop.h>
#include <LibCore/EventReceiver.h>
#include <LibIPC/File.h>
#include <LibIPC/Forward.h>
#include <LibIPC/Message.h>
#include <LibIPC/Transport.h>
#include <LibThreading/ConditionVariable.h>
#include <LibThreading/Mutex.h>
namespace IPC {
@ -40,31 +36,25 @@ public:
protected:
explicit ConnectionBase(IPC::Stub&, NonnullOwnPtr<Transport>, u32 local_endpoint_magic);
// Must be called after setting up the message decoder
void initialize_messaging();
virtual void shutdown_with_error(Error const&);
virtual OwnPtr<Message> try_parse_message(ReadonlyBytes, Queue<File>&) = 0;
OwnPtr<IPC::Message> wait_for_specific_endpoint_message_impl(u32 endpoint_magic, int message_id);
void wait_for_transport_to_become_readable();
enum class PeerEOF {
No,
Yes
};
PeerEOF drain_messages_from_peer();
void handle_messages();
// Called from I/O thread when a message is decoded
void on_message_received(NonnullOwnPtr<IPC::Message> message);
void on_peer_closed();
IPC::Stub& m_local_stub;
NonnullOwnPtr<Transport> m_transport;
Threading::Mutex m_unprocessed_messages_mutex;
Threading::ConditionVariable m_unprocessed_messages_cv { m_unprocessed_messages_mutex };
Vector<NonnullOwnPtr<Message>> m_unprocessed_messages;
RefPtr<Core::WeakEventLoopReference> m_event_loop;
Atomic<bool> m_peer_closed { false };
u32 m_local_endpoint_magic { 0 };
};
@ -74,28 +64,6 @@ public:
Connection(IPC::Stub& local_stub, NonnullOwnPtr<Transport> transport)
: ConnectionBase(local_stub, move(transport), LocalEndpoint::static_magic())
{
// Set up message handler first so we're ready to receive
initialize_messaging();
// Then set up decoder
m_transport->set_message_decoder([](ReadonlyBytes bytes, Queue<File>& fds) -> OwnPtr<IPC::Message> {
auto local = LocalEndpoint::decode_message(bytes, fds);
if (!local.is_error())
return local.release_value();
auto peer = PeerEndpoint::decode_message(bytes, fds);
if (!peer.is_error())
return peer.release_value();
dbgln("Failed to parse IPC message:");
dbgln(" Local endpoint error: {}", local.error());
dbgln(" Peer endpoint error: {}", peer.error());
return nullptr;
});
// Now that handler and decoder are set up, start receiving messages
m_transport->start();
}
template<typename RequestType, typename... Args>
@ -123,6 +91,23 @@ protected:
return message.template release_nonnull<MessageType>();
return {};
}
virtual OwnPtr<Message> try_parse_message(ReadonlyBytes bytes, Queue<File>& fds) override
{
auto local_message = LocalEndpoint::decode_message(bytes, fds);
if (!local_message.is_error())
return local_message.release_value();
auto peer_message = PeerEndpoint::decode_message(bytes, fds);
if (!peer_message.is_error())
return peer_message.release_value();
dbgln("Failed to parse IPC message:");
dbgln(" Local endpoint error: {}", local_message.error());
dbgln(" Peer endpoint error: {}", peer_message.error());
return nullptr;
}
};
}

View file

@ -11,7 +11,6 @@
#include <LibCore/Socket.h>
#include <LibCore/System.h>
#include <LibIPC/Limits.h>
#include <LibIPC/Message.h>
#include <LibIPC/TransportSocket.h>
#include <LibThreading/Thread.h>
@ -68,10 +67,6 @@ TransportSocket::TransportSocket(NonnullOwnPtr<Core::LocalSocket> socket)
}
m_io_thread = Threading::Thread::construct([this] { return io_thread_loop(); });
}
void TransportSocket::start()
{
m_io_thread->start();
}
@ -178,9 +173,6 @@ void TransportSocket::set_up_read_hook(Function<void()> hook)
m_on_read_hook();
};
// Start I/O thread now that the read hook is set up (for raw message path)
start();
{
Threading::MutexLocker locker(m_incoming_mutex);
if (!m_incoming_messages.is_empty()) {
@ -296,21 +288,6 @@ ErrorOr<void> TransportSocket::send_message(Core::LocalSocket& socket, ReadonlyB
return {};
}
void TransportSocket::set_message_decoder(MessageDecoder decoder)
{
m_decoder = move(decoder);
}
void TransportSocket::set_message_handler(MessageHandler handler)
{
m_message_handler = move(handler);
}
void TransportSocket::set_peer_closed_handler(PeerClosedHandler handler)
{
m_peer_closed_handler = move(handler);
}
TransportSocket::TransferState TransportSocket::transfer_data(ReadonlyBytes& bytes, Vector<int>& fds)
{
auto byte_count = bytes.size();
@ -337,7 +314,6 @@ TransportSocket::TransferState TransportSocket::transfer_data(ReadonlyBytes& byt
void TransportSocket::read_incoming_messages()
{
Vector<NonnullOwnPtr<Message>> batch;
Vector<NonnullOwnPtr<IPC::Message>> decoded_batch;
while (m_socket->is_open()) {
u8 buffer[4096];
auto received_fds = Vector<int> {};
@ -408,37 +384,21 @@ void TransportSocket::read_incoming_messages()
break;
if (header.fd_count > m_unprocessed_fds.size())
break;
auto message = make<Message>();
received_fd_count += header.fd_count;
if (received_fd_count.has_overflow()) {
dbgln("TransportSocket: received_fd_count would overflow");
m_peer_eof = true;
break;
}
if (m_decoder) {
// Decode on I/O thread
Queue<File> fds_for_message;
for (size_t i = 0; i < header.fd_count; ++i)
fds_for_message.enqueue(m_unprocessed_fds.dequeue());
ReadonlyBytes payload { m_unprocessed_bytes.data() + index + sizeof(MessageHeader), header.payload_size };
if (auto decoded = m_decoder(payload, fds_for_message)) {
decoded_batch.append(decoded.release_nonnull());
} else {
dbgln("TransportSocket: Failed to decode message");
m_peer_eof = true;
break;
}
} else {
// Legacy path: store raw bytes
auto message = make<Message>();
for (size_t i = 0; i < header.fd_count; ++i)
message->fds.enqueue(m_unprocessed_fds.dequeue());
if (message->bytes.try_append(m_unprocessed_bytes.data() + index + sizeof(MessageHeader), header.payload_size).is_error()) {
dbgln("TransportSocket: Failed to allocate message buffer for payload_size {}", header.payload_size);
m_peer_eof = true;
break;
}
batch.append(move(message));
for (size_t i = 0; i < header.fd_count; ++i)
message->fds.enqueue(m_unprocessed_fds.dequeue());
if (message->bytes.try_append(m_unprocessed_bytes.data() + index + sizeof(MessageHeader), header.payload_size).is_error()) {
dbgln("TransportSocket: Failed to allocate message buffer for payload_size {}", header.payload_size);
m_peer_eof = true;
break;
}
batch.append(move(message));
} else if (header.type == MessageHeader::Type::FileDescriptorAcknowledgement) {
if (header.payload_size != 0) {
dbgln("TransportSocket: FileDescriptorAcknowledgement with non-zero payload_size {}", header.payload_size);
@ -509,17 +469,6 @@ void TransportSocket::read_incoming_messages()
(void)Core::System::write(m_notify_hook_write_fd->value(), bytes);
};
// IPC::Connection path: call message handler for decoded IPC::Message objects.
// The handler is responsible for storing messages and waking the event loop.
if (!decoded_batch.is_empty() && m_message_handler) {
for (auto& msg : decoded_batch) {
m_message_handler(move(msg));
}
}
// Raw message path: store undecoded messages for retrieval via
// read_as_many_messages_as_possible_without_blocking().
// Used by MessagePort which has its own message format (SerializedTransferRecord).
if (!batch.is_empty()) {
Threading::MutexLocker locker(m_incoming_mutex);
m_incoming_messages.extend(move(batch));
@ -528,8 +477,6 @@ void TransportSocket::read_incoming_messages()
}
if (m_peer_eof) {
if (m_peer_closed_handler)
m_peer_closed_handler();
m_incoming_cv.broadcast();
notify_read_available();
}

View file

@ -19,8 +19,6 @@
namespace IPC {
class Message;
class SendQueue : public AtomicRefCounted<SendQueue> {
public:
void enqueue_message(Vector<u8>&& bytes, Vector<int>&& fds);
@ -44,17 +42,6 @@ class TransportSocket {
public:
static constexpr socklen_t SOCKET_BUFFER_SIZE = 128 * KiB;
// IPC::Connection path: Set a decoder to parse raw bytes into IPC::Message objects
// on the I/O thread, and a handler to receive decoded messages.
using MessageDecoder = Function<OwnPtr<IPC::Message>(ReadonlyBytes, Queue<File>&)>;
using MessageHandler = Function<void(NonnullOwnPtr<IPC::Message>)>;
using PeerClosedHandler = Function<void()>;
void set_message_decoder(MessageDecoder decoder);
void set_message_handler(MessageHandler handler);
void set_peer_closed_handler(PeerClosedHandler handler);
void start();
explicit TransportSocket(NonnullOwnPtr<Core::LocalSocket> socket);
~TransportSocket();
@ -72,9 +59,6 @@ public:
No,
Yes,
};
// Raw message path: Used by MessagePort which has its own message format
// (SerializedTransferRecord) rather than IPC::Message.
struct Message {
Vector<u8> bytes;
Queue<File> fds;
@ -130,10 +114,6 @@ private:
RefPtr<AutoCloseFileDescriptor> m_notify_hook_write_fd;
RefPtr<Core::Notifier> m_read_hook_notifier;
Function<void()> m_on_read_hook;
MessageDecoder m_decoder;
MessageHandler m_message_handler;
PeerClosedHandler m_peer_closed_handler;
};
}

View file

@ -10,7 +10,6 @@
#include <AK/Types.h>
#include <LibIPC/HandleType.h>
#include <LibIPC/Limits.h>
#include <LibIPC/Message.h>
#include <LibIPC/TransportSocketWindows.h>
#include <AK/Windows.h>
@ -22,39 +21,6 @@ TransportSocketWindows::TransportSocketWindows(NonnullOwnPtr<Core::LocalSocket>
{
}
void TransportSocketWindows::set_message_decoder(MessageDecoder decoder)
{
m_decoder = move(decoder);
}
void TransportSocketWindows::set_message_handler(MessageHandler handler)
{
m_message_handler = move(handler);
}
void TransportSocketWindows::set_peer_closed_handler(PeerClosedHandler handler)
{
m_peer_closed_handler = move(handler);
}
void TransportSocketWindows::start()
{
// Windows does not use a separate I/O thread.
// Instead, set up a read hook that decodes messages on the main thread.
VERIFY(m_decoder);
VERIFY(m_message_handler);
m_socket->on_ready_to_read = [this] {
auto should_shutdown = read_as_many_messages_as_possible_without_blocking([this](Message&& message) {
Queue<File> fds;
if (auto decoded = m_decoder(message.bytes.span(), fds)) {
m_message_handler(decoded.release_nonnull());
}
});
if (should_shutdown == ShouldShutdown::Yes && m_peer_closed_handler)
m_peer_closed_handler();
};
}
void TransportSocketWindows::set_peer_pid(int pid)
{
m_peer_pid = pid;
@ -62,7 +28,6 @@ void TransportSocketWindows::set_peer_pid(int pid)
void TransportSocketWindows::set_up_read_hook(Function<void()> hook)
{
// Raw message path (used by MessagePort)
VERIFY(m_socket->is_open());
m_socket->on_ready_to_read = move(hook);
}

View file

@ -13,25 +13,11 @@
namespace IPC {
class Message;
class TransportSocketWindows {
AK_MAKE_NONCOPYABLE(TransportSocketWindows);
AK_MAKE_DEFAULT_MOVABLE(TransportSocketWindows);
public:
// IPC::Connection path: Set a decoder to parse raw bytes into IPC::Message objects
// and a handler to receive decoded messages.
// NOTE: Windows does not use a separate I/O thread; decoding happens on the main thread.
using MessageDecoder = Function<OwnPtr<IPC::Message>(ReadonlyBytes, Queue<File>&)>;
using MessageHandler = Function<void(NonnullOwnPtr<IPC::Message>)>;
using PeerClosedHandler = Function<void()>;
void set_message_decoder(MessageDecoder decoder);
void set_message_handler(MessageHandler handler);
void set_peer_closed_handler(PeerClosedHandler handler);
void start();
explicit TransportSocketWindows(NonnullOwnPtr<Core::LocalSocket> socket);
void set_peer_pid(int pid);
@ -67,10 +53,6 @@ private:
NonnullOwnPtr<Core::LocalSocket> m_socket;
ByteBuffer m_unprocessed_bytes;
int m_peer_pid = -1;
MessageDecoder m_decoder;
MessageHandler m_message_handler;
PeerClosedHandler m_peer_closed_handler;
};
}