From fbfd1f3efaf98742c14b706d96c2de478e7a9ddf Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Mart=C3=ADn=20Lucas=20Golini?= Date: Wed, 2 Sep 2026 02:28:35 -0300 Subject: [PATCH] add cross-platform local IPC and migrate ecode Add a generic, message-oriented EE::System::IPC abstraction using Unix-domain stream sockets on POSIX platforms and named pipes on Windows. - add versioned, binary-safe message framing - provide same-user endpoint isolation and hashed native endpoint names - handle stale POSIX sockets, duplicate listeners, timeouts, and clean shutdown - optimize common messages with inline storage and stack-backed endpoint buffers - document the public IPC API and worker-thread callback contract - add IPC coverage for framing errors, partial messages, binary payloads, sequential clients, duplicate ownership, stale sockets, and relistening - migrate ecode to profile-scoped IPC endpoints and JSON open-file messages - preserve first/latest-instance selection behavior - remove the filesystem watcher-based IPC transport --- include/eepp/system.hpp | 1 + include/eepp/system/ipc.hpp | 110 ++++++ src/eepp/system/ipc.cpp | 599 +++++++++++++++++++++++++++++ src/tests/unit_tests/ipc_tests.cpp | 213 ++++++++++ src/tools/ecode/ecode.cpp | 153 ++++---- src/tools/ecode/ecode.hpp | 5 +- 6 files changed, 1004 insertions(+), 77 deletions(-) create mode 100644 include/eepp/system/ipc.hpp create mode 100644 src/eepp/system/ipc.cpp create mode 100644 src/tests/unit_tests/ipc_tests.cpp diff --git a/include/eepp/system.hpp b/include/eepp/system.hpp index 8d83f3c01..b059bb9a0 100644 --- a/include/eepp/system.hpp +++ b/include/eepp/system.hpp @@ -24,6 +24,7 @@ #include #include #include +#include #include #include #include diff --git a/include/eepp/system/ipc.hpp b/include/eepp/system/ipc.hpp new file mode 100644 index 000000000..2c52a1ebc --- /dev/null +++ b/include/eepp/system/ipc.hpp @@ -0,0 +1,110 @@ +#ifndef EE_SYSTEM_IPC_HPP +#define EE_SYSTEM_IPC_HPP + +#include +#include +#include +#include +#include +#include + +namespace EE { namespace System { + +/** A named, local, same-user, message-oriented inter-process channel. + * + * listen() starts a receiver thread. Message callbacks execute on that thread and callers are + * responsible for dispatching to another thread when required. Each send() creates one connection, + * transfers one complete framed message, and disconnects. + */ +class EE_API IPC : NonCopyable { + public: + /** Result of an IPC operation. Native error codes are mapped to this portable set. */ + enum class Status { + Done, ///< The operation completed successfully. + NotFound, ///< No listener currently owns the destination endpoint. + AlreadyExists, ///< Another listener already owns the endpoint. + Disconnected, ///< The connection closed before the complete message was transferred. + InvalidEndpoint, ///< The logical endpoint or payload pointer is invalid. + MessageTooLarge, ///< The payload exceeds MaxMessageSize. + Timeout, ///< The destination did not become available before the timeout. + Error ///< An otherwise unmapped native transport error occurred. + }; + + /** Callback invoked for each complete message received by a listener. + * + * The callback runs on the listener's worker thread. @p data remains valid only for the + * duration of the callback and may contain embedded NUL bytes. For an empty message, @p data is + * still a valid pointer and @p size is zero. + */ + using MessageFn = std::function; + + /** Maximum accepted payload size, excluding the transport frame header. */ + static constexpr std::size_t MaxMessageSize = 16 * 1024 * 1024; + + /** Constructs a closed IPC listener. */ + IPC(); + + /** Stops the listener, waits for its worker thread, and releases its native endpoint. */ + ~IPC(); + + /** Starts listening on a logical endpoint. + * + * The endpoint is mapped synchronously to a platform-native, same-user local IPC address; the + * borrowed logical name is not retained. Only one IPC instance can listen on a given endpoint + * at a time. A successful call starts one worker thread and invokes @p callback once for every + * complete framed message. Calling listen() on an instance that is already listening closes its + * current endpoint first. + * + * @param endpoint Borrowed logical endpoint name. Native paths or handles are not exposed. + * @param callback Function invoked on the IPC worker thread for each complete message. + * @return Done on success, AlreadyExists when the endpoint is owned by another listener, + * InvalidEndpoint for an empty/invalid endpoint or callback, or Error on a native failure. + */ + Status listen( std::string_view endpoint, MessageFn callback ); + + /** Stops listening and releases the endpoint. + * + * This function interrupts pending native waits, joins the worker thread, and guarantees that + * the callback will not run after it returns. Repeated calls are safe. + */ + void close(); + + /** @return Whether this instance currently owns a listening endpoint. */ + bool isListening() const; + + /** Sends one binary message to a logical endpoint. + * + * A successful call creates a connection, writes exactly one versioned frame containing @p size + * bytes, and disconnects. Embedded NUL bytes and zero-length messages are supported. The + * timeout applies while establishing or waiting for the destination connection. + * + * @param endpoint Borrowed logical destination endpoint name. + * @param data Payload bytes. May be null only when @p size is zero. + * @param size Payload size in bytes, up to MaxMessageSize. + * @param timeout Maximum time to wait for the destination connection. + * @return A portable status describing delivery or connection failure. + */ + static Status send( std::string_view endpoint, const void* data, std::size_t size, + Time timeout = Seconds( 2 ) ); + + /** Sends one string-view payload to a logical endpoint. + * + * The complete view is sent as opaque bytes; embedded NUL bytes are preserved and no terminator + * is appended. + * + * @param endpoint Borrowed logical destination endpoint name. + * @param message Borrowed payload bytes. + * @param timeout Maximum time to wait for the destination connection. + * @return A portable status describing delivery or connection failure. + */ + static Status send( std::string_view endpoint, std::string_view message, + Time timeout = Seconds( 2 ) ); + + private: + class Impl; + std::unique_ptr mImpl; +}; + +}} // namespace EE::System + +#endif diff --git a/src/eepp/system/ipc.cpp b/src/eepp/system/ipc.cpp new file mode 100644 index 000000000..c57ce7f9f --- /dev/null +++ b/src/eepp/system/ipc.cpp @@ -0,0 +1,599 @@ +#include + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#if EE_PLATFORM == EE_PLATFORM_WIN +#include +#include +#include +#elif EE_PLATFORM != EE_PLATFORM_EMSCRIPTEN +#include +#include +#include +#include +#include +#include +#include +#endif + +namespace EE { namespace System { + +namespace { + +constexpr Uint32 IPCMagic = 0x45454950; // "EEIP" +constexpr Uint16 IPCVersion = 1; +constexpr Uint16 IPCFlags = 0; +constexpr std::size_t IPCHeaderSize = 12; +constexpr std::size_t IPCInlinePayloadSize = 512; +constexpr Uint32 IPCPipeBufferSize = 64 * 1024; +using EndpointHash = std::array; +using NativeName = std::array; + +class MessageBuffer { + public: + explicit MessageBuffer( std::size_t size ) : + mHeap( size > IPCInlinePayloadSize ? new Uint8[size] : nullptr ) {} + + Uint8* data() { return mHeap ? mHeap.get() : mInline.data(); } + + private: + std::array mInline; + std::unique_ptr mHeap; +}; + +static_assert( sizeof( MessageBuffer ) <= IPCInlinePayloadSize + sizeof( void* ) ); + +std::array makeHeader( std::size_t payloadSize ) { + std::array header{}; + const Uint32 size = static_cast( payloadSize ); + header[0] = static_cast( IPCMagic >> 24 ); + header[1] = static_cast( IPCMagic >> 16 ); + header[2] = static_cast( IPCMagic >> 8 ); + header[3] = static_cast( IPCMagic ); + header[4] = static_cast( IPCVersion >> 8 ); + header[5] = static_cast( IPCVersion ); + header[6] = static_cast( IPCFlags >> 8 ); + header[7] = static_cast( IPCFlags ); + header[8] = static_cast( size >> 24 ); + header[9] = static_cast( size >> 16 ); + header[10] = static_cast( size >> 8 ); + header[11] = static_cast( size ); + return header; +} + +bool decodeHeader( const std::array& wire, Uint32& payloadSize ) { + const Uint32 magic = static_cast( wire[0] ) << 24 | + static_cast( wire[1] ) << 16 | + static_cast( wire[2] ) << 8 | wire[3]; + const Uint16 version = static_cast( wire[4] ) << 8 | wire[5]; + const Uint16 flags = static_cast( wire[6] ) << 8 | wire[7]; + payloadSize = static_cast( wire[8] ) << 24 | static_cast( wire[9] ) << 16 | + static_cast( wire[10] ) << 8 | wire[11]; + return magic == IPCMagic && version == IPCVersion && flags == IPCFlags && + payloadSize <= IPC::MaxMessageSize; +} + +bool validEndpoint( std::string_view endpoint ) { + return !endpoint.empty() && endpoint.size() <= 4096; +} + +EndpointHash endpointHash( std::string_view endpoint ) { + MD5::Context context; + MD5::init( context ); + MD5::update( context, endpoint.data(), endpoint.size() ); + const MD5::Digest digest = MD5::result( context ).digest; + constexpr char HexDigits[] = "0123456789abcdef"; + EndpointHash hash{}; + for ( std::size_t i = 0; i < digest.size(); ++i ) { + hash[i * 2] = HexDigits[digest[i] >> 4]; + hash[i * 2 + 1] = HexDigits[digest[i] & 0x0F]; + } + return hash; +} + +} // namespace + +class IPC::Impl { + public: + std::atomic running{ false }; + MessageFn callback; + std::thread worker; + std::mutex handlesMutex; + +#if EE_PLATFORM == EE_PLATFORM_WIN + HANDLE listener{ INVALID_HANDLE_VALUE }; + HANDLE ownership{ nullptr }; + NativeName nativeName{}; + PSECURITY_DESCRIPTOR securityDescriptor{ nullptr }; +#elif EE_PLATFORM != EE_PLATFORM_EMSCRIPTEN + int listener{ -1 }; + int client{ -1 }; + sockaddr_un nativeAddress{}; +#endif +}; + +#if EE_PLATFORM == EE_PLATFORM_WIN + +namespace { + +NativeName nativeName( std::wstring_view prefix, const EndpointHash& hash ) { + NativeName name{}; + if ( prefix.size() + hash.size() > name.size() ) + return name; + std::copy( prefix.begin(), prefix.end(), name.begin() ); + for ( std::size_t i = 0; i < hash.size() - 1; ++i ) + name[prefix.size() + i] = static_cast( hash[i] ); + return name; +} + +HANDLE createPipe( const wchar_t* name, PSECURITY_DESCRIPTOR descriptor, bool first ) { + SECURITY_ATTRIBUTES attributes{ sizeof( SECURITY_ATTRIBUTES ), descriptor, FALSE }; + return CreateNamedPipeW( + name, PIPE_ACCESS_INBOUND | ( first ? FILE_FLAG_FIRST_PIPE_INSTANCE : 0 ), + PIPE_TYPE_BYTE | PIPE_READMODE_BYTE | PIPE_WAIT | PIPE_REJECT_REMOTE_CLIENTS, + PIPE_UNLIMITED_INSTANCES, IPCPipeBufferSize, IPCPipeBufferSize, 0, &attributes ); +} + +bool readExact( HANDLE pipe, void* data, std::size_t size ) { + auto* bytes = static_cast( data ); + while ( size > 0 ) { + DWORD read = 0; + const DWORD chunk = static_cast( std::min( size, MAXDWORD ) ); + if ( !ReadFile( pipe, bytes, chunk, &read, nullptr ) || read == 0 ) + return false; + bytes += read; + size -= read; + } + return true; +} + +bool writeExact( HANDLE pipe, const void* data, std::size_t size ) { + const auto* bytes = static_cast( data ); + while ( size > 0 ) { + DWORD written = 0; + const DWORD chunk = static_cast( std::min( size, MAXDWORD ) ); + if ( !WriteFile( pipe, bytes, chunk, &written, nullptr ) || written == 0 ) + return false; + bytes += written; + size -= written; + } + return true; +} + +bool endpointOwned( const wchar_t* ownershipName ) { + HANDLE ownership = OpenMutexW( SYNCHRONIZE, FALSE, ownershipName ); + if ( ownership == nullptr ) + return false; + CloseHandle( ownership ); + return true; +} + +IPC::Status windowsErrorStatus( DWORD error ) { + if ( error == ERROR_FILE_NOT_FOUND || error == ERROR_PATH_NOT_FOUND ) + return IPC::Status::NotFound; + if ( error == ERROR_PIPE_BUSY || error == ERROR_SEM_TIMEOUT ) + return IPC::Status::Timeout; + if ( error == ERROR_BROKEN_PIPE || error == ERROR_NO_DATA ) + return IPC::Status::Disconnected; + return IPC::Status::Error; +} + +} // namespace + +#elif EE_PLATFORM != EE_PLATFORM_EMSCRIPTEN + +namespace { + +const std::string& runtimeDirectory() { + static const std::string directory = [] { + std::string path = Sys::getTempPath(); + if ( path.empty() ) + return std::string{}; + if ( path.back() != '/' ) + path += '/'; + path += "eepp-ipc-" + std::to_string( static_cast( getuid() ) ); + if ( mkdir( path.c_str(), 0700 ) != 0 && errno != EEXIST ) + return std::string{}; + struct stat info; + if ( lstat( path.c_str(), &info ) != 0 || !S_ISDIR( info.st_mode ) || + info.st_uid != getuid() || ( info.st_mode & 077 ) != 0 ) + return std::string{}; + return path; + }(); + return directory; +} + +bool makeAddress( std::string_view endpoint, sockaddr_un& address ) { + const std::string& directory = runtimeDirectory(); + if ( directory.empty() ) + return false; + const EndpointHash hash = endpointHash( endpoint ); + const std::size_t pathSize = directory.size() + 1 + hash.size() - 1; + if ( pathSize >= sizeof( address.sun_path ) ) + return false; + std::memset( &address, 0, sizeof( address ) ); + address.sun_family = AF_UNIX; + std::memcpy( address.sun_path, directory.data(), directory.size() ); + address.sun_path[directory.size()] = '/'; + std::memcpy( address.sun_path + directory.size() + 1, hash.data(), hash.size() - 1 ); + return true; +} + +bool readExact( int fd, void* data, std::size_t size ) { + auto* bytes = static_cast( data ); + while ( size > 0 ) { + const ssize_t readSize = recv( fd, bytes, size, 0 ); + if ( readSize == 0 ) + return false; + if ( readSize < 0 ) { + if ( errno == EINTR ) + continue; + return false; + } + bytes += readSize; + size -= static_cast( readSize ); + } + return true; +} + +bool writeExact( int fd, const void* data, std::size_t size ) { + const auto* bytes = static_cast( data ); + while ( size > 0 ) { +#ifdef MSG_NOSIGNAL + const ssize_t written = send( fd, bytes, size, MSG_NOSIGNAL ); +#else + const ssize_t written = send( fd, bytes, size, 0 ); +#endif + if ( written < 0 ) { + if ( errno == EINTR ) + continue; + return false; + } + if ( written == 0 ) + return false; + bytes += written; + size -= static_cast( written ); + } + return true; +} + +IPC::Status connectSocket( const sockaddr_un& address, Time timeout, int& fd ) { + fd = socket( AF_UNIX, SOCK_STREAM, 0 ); + if ( fd < 0 ) + return IPC::Status::Error; + + const int oldFlags = fcntl( fd, F_GETFL, 0 ); + if ( oldFlags >= 0 ) + fcntl( fd, F_SETFL, oldFlags | O_NONBLOCK ); + if ( connect( fd, reinterpret_cast( &address ), sizeof( address ) ) != 0 ) { + if ( errno != EINPROGRESS ) { + const int error = errno; + ::close( fd ); + fd = -1; + return error == ENOENT || error == ECONNREFUSED ? IPC::Status::NotFound + : IPC::Status::Error; + } + pollfd wait{ fd, POLLOUT, 0 }; + const int timeoutMs = static_cast( eemax( 0, timeout.asMilliseconds() ) ); + int result; + do { + result = poll( &wait, 1, timeoutMs ); + } while ( result < 0 && errno == EINTR ); + if ( result <= 0 ) { + ::close( fd ); + fd = -1; + return result == 0 ? IPC::Status::Timeout : IPC::Status::Error; + } + int error = 0; + socklen_t errorSize = sizeof( error ); + if ( getsockopt( fd, SOL_SOCKET, SO_ERROR, &error, &errorSize ) != 0 || error != 0 ) { + ::close( fd ); + fd = -1; + return error == ENOENT || error == ECONNREFUSED ? IPC::Status::NotFound + : IPC::Status::Error; + } + } + if ( oldFlags >= 0 ) + fcntl( fd, F_SETFL, oldFlags ); + return IPC::Status::Done; +} + +} // namespace + +#endif + +IPC::IPC() : mImpl( std::make_unique() ) {} + +IPC::~IPC() { + close(); +} + +IPC::Status IPC::listen( std::string_view endpoint, MessageFn callback ) { + if ( !validEndpoint( endpoint ) || !callback ) + return Status::InvalidEndpoint; + close(); + +#if EE_PLATFORM == EE_PLATFORM_EMSCRIPTEN + return Status::Error; +#elif EE_PLATFORM == EE_PLATFORM_WIN + const EndpointHash hash = endpointHash( endpoint ); + const NativeName pipe = nativeName( L"\\\\.\\pipe\\eepp-", hash ); + const NativeName ownership = nativeName( L"Local\\eepp-ipc-owner-", hash ); + mImpl->nativeName = pipe; + if ( !ConvertStringSecurityDescriptorToSecurityDescriptorW( + L"D:P(A;;GA;;;OW)", SDDL_REVISION_1, &mImpl->securityDescriptor, nullptr ) ) + return Status::Error; + SECURITY_ATTRIBUTES attributes{ sizeof( SECURITY_ATTRIBUTES ), mImpl->securityDescriptor, + FALSE }; + mImpl->ownership = CreateMutexW( &attributes, FALSE, ownership.data() ); + if ( mImpl->ownership == nullptr || GetLastError() == ERROR_ALREADY_EXISTS ) { + if ( mImpl->ownership != nullptr ) + CloseHandle( mImpl->ownership ); + mImpl->ownership = nullptr; + LocalFree( mImpl->securityDescriptor ); + mImpl->securityDescriptor = nullptr; + return Status::AlreadyExists; + } + mImpl->listener = createPipe( mImpl->nativeName.data(), mImpl->securityDescriptor, true ); + if ( mImpl->listener == INVALID_HANDLE_VALUE ) { + const DWORD error = GetLastError(); + CloseHandle( mImpl->ownership ); + mImpl->ownership = nullptr; + LocalFree( mImpl->securityDescriptor ); + mImpl->securityDescriptor = nullptr; + return error == ERROR_ACCESS_DENIED || error == ERROR_PIPE_BUSY || + error == ERROR_ALREADY_EXISTS + ? Status::AlreadyExists + : Status::Error; + } + mImpl->callback = std::move( callback ); + mImpl->running = true; + mImpl->worker = std::thread( [impl = mImpl.get()] { + while ( impl->running ) { + HANDLE current; + { + std::lock_guard lock( impl->handlesMutex ); + current = impl->listener; + } + const BOOL connected = ConnectNamedPipe( current, nullptr ) + ? TRUE + : GetLastError() == ERROR_PIPE_CONNECTED; + if ( !connected ) { + if ( !impl->running ) + break; + continue; + } + if ( !impl->running ) + break; + std::array header; + Uint32 payloadSize = 0; + if ( readExact( current, header.data(), header.size() ) && + decodeHeader( header, payloadSize ) ) { + MessageBuffer payload( payloadSize ); + if ( readExact( current, payload.data(), payloadSize ) && impl->running ) { + const Uint8 empty = 0; + impl->callback( payloadSize == 0 ? &empty : payload.data(), payloadSize ); + } + } + CloseHandle( current ); + HANDLE replacement = INVALID_HANDLE_VALUE; + if ( impl->running ) + replacement = + createPipe( impl->nativeName.data(), impl->securityDescriptor, false ); + { + std::lock_guard lock( impl->handlesMutex ); + impl->listener = replacement; + } + if ( replacement == INVALID_HANDLE_VALUE ) + break; + } + } ); + return Status::Done; +#else + if ( !makeAddress( endpoint, mImpl->nativeAddress ) ) + return Status::Error; + mImpl->listener = socket( AF_UNIX, SOCK_STREAM, 0 ); + if ( mImpl->listener < 0 ) + return Status::Error; + if ( bind( mImpl->listener, reinterpret_cast( &mImpl->nativeAddress ), + sizeof( mImpl->nativeAddress ) ) != 0 ) { + if ( errno != EADDRINUSE ) { + ::close( mImpl->listener ); + mImpl->listener = -1; + return Status::Error; + } + int probe = -1; + const Status probeStatus = + connectSocket( mImpl->nativeAddress, Milliseconds( 100 ), probe ); + if ( probeStatus == Status::Done ) { + ::close( probe ); + ::close( mImpl->listener ); + mImpl->listener = -1; + return Status::AlreadyExists; + } + if ( probeStatus != Status::NotFound || unlink( mImpl->nativeAddress.sun_path ) != 0 || + bind( mImpl->listener, reinterpret_cast( &mImpl->nativeAddress ), + sizeof( mImpl->nativeAddress ) ) != 0 ) { + ::close( mImpl->listener ); + mImpl->listener = -1; + return Status::Error; + } + } + chmod( mImpl->nativeAddress.sun_path, 0600 ); + if ( ::listen( mImpl->listener, SOMAXCONN ) != 0 ) { + ::close( mImpl->listener ); + mImpl->listener = -1; + unlink( mImpl->nativeAddress.sun_path ); + return Status::Error; + } + mImpl->callback = std::move( callback ); + mImpl->running = true; + mImpl->worker = std::thread( [impl = mImpl.get()] { + while ( impl->running ) { + const int client = accept( impl->listener, nullptr, nullptr ); + if ( client < 0 ) { + if ( !impl->running ) + break; + if ( errno == EINTR ) + continue; + break; + } + { + std::lock_guard lock( impl->handlesMutex ); + impl->client = client; + } + if ( !impl->running ) { + { + std::lock_guard lock( impl->handlesMutex ); + impl->client = -1; + } + ::close( client ); + break; + } + std::array header; + Uint32 payloadSize = 0; + if ( readExact( client, header.data(), header.size() ) && + decodeHeader( header, payloadSize ) ) { + MessageBuffer payload( payloadSize ); + if ( readExact( client, payload.data(), payloadSize ) && impl->running ) { + const Uint8 empty = 0; + impl->callback( payloadSize == 0 ? &empty : payload.data(), payloadSize ); + } + } + { + std::lock_guard lock( impl->handlesMutex ); + if ( impl->client == client ) + impl->client = -1; + } + ::close( client ); + } + } ); + return Status::Done; +#endif +} + +void IPC::close() { + if ( !mImpl->running.exchange( false ) ) { + if ( mImpl->worker.joinable() ) + mImpl->worker.join(); + return; + } +#if EE_PLATFORM == EE_PLATFORM_WIN + { + std::lock_guard lock( mImpl->handlesMutex ); + CancelSynchronousIo( reinterpret_cast( mImpl->worker.native_handle() ) ); + HANDLE wake = CreateFileW( mImpl->nativeName.data(), GENERIC_WRITE, 0, nullptr, + OPEN_EXISTING, FILE_ATTRIBUTE_NORMAL, nullptr ); + if ( wake != INVALID_HANDLE_VALUE ) + CloseHandle( wake ); + } +#elif EE_PLATFORM != EE_PLATFORM_EMSCRIPTEN + { + std::lock_guard lock( mImpl->handlesMutex ); + if ( mImpl->client >= 0 ) + shutdown( mImpl->client, SHUT_RDWR ); + } + if ( mImpl->listener >= 0 ) { + shutdown( mImpl->listener, SHUT_RDWR ); + ::close( mImpl->listener ); + mImpl->listener = -1; + } +#endif + if ( mImpl->worker.joinable() ) + mImpl->worker.join(); +#if EE_PLATFORM == EE_PLATFORM_WIN + { + std::lock_guard lock( mImpl->handlesMutex ); + if ( mImpl->listener != INVALID_HANDLE_VALUE ) { + DisconnectNamedPipe( mImpl->listener ); + CloseHandle( mImpl->listener ); + mImpl->listener = INVALID_HANDLE_VALUE; + } + } + if ( mImpl->securityDescriptor != nullptr ) { + LocalFree( mImpl->securityDescriptor ); + mImpl->securityDescriptor = nullptr; + } + if ( mImpl->ownership != nullptr ) { + CloseHandle( mImpl->ownership ); + mImpl->ownership = nullptr; + } +#endif +#if EE_PLATFORM != EE_PLATFORM_WIN && EE_PLATFORM != EE_PLATFORM_EMSCRIPTEN + if ( mImpl->nativeAddress.sun_path[0] != '\0' ) + unlink( mImpl->nativeAddress.sun_path ); +#endif + mImpl->callback = {}; +} + +bool IPC::isListening() const { + return mImpl->running; +} + +IPC::Status IPC::send( std::string_view endpoint, const void* data, std::size_t size, + Time timeout ) { + if ( !validEndpoint( endpoint ) || ( size > 0 && data == nullptr ) ) + return Status::InvalidEndpoint; + if ( size > MaxMessageSize ) + return Status::MessageTooLarge; + const auto header = makeHeader( size ); + +#if EE_PLATFORM == EE_PLATFORM_EMSCRIPTEN + return Status::Error; +#elif EE_PLATFORM == EE_PLATFORM_WIN + const EndpointHash hash = endpointHash( endpoint ); + const NativeName name = nativeName( L"\\\\.\\pipe\\eepp-", hash ); + const NativeName ownership = nativeName( L"Local\\eepp-ipc-owner-", hash ); + const Uint64 timeoutMs = static_cast( eemax( 0, timeout.asMilliseconds() ) ); + const Uint64 deadline = GetTickCount64() + timeoutMs; + HANDLE pipe; + while ( true ) { + pipe = CreateFileW( name.data(), GENERIC_WRITE, 0, nullptr, OPEN_EXISTING, + FILE_ATTRIBUTE_NORMAL, nullptr ); + if ( pipe != INVALID_HANDLE_VALUE ) + break; + const DWORD error = GetLastError(); + const bool transitioning = + error == ERROR_FILE_NOT_FOUND && endpointOwned( ownership.data() ); + if ( error != ERROR_PIPE_BUSY && error != ERROR_SEM_TIMEOUT && !transitioning ) + return windowsErrorStatus( error ); + if ( GetTickCount64() >= deadline ) + return Status::Timeout; + Sleep( 1 ); + } + const bool written = writeExact( pipe, header.data(), header.size() ) && + ( size == 0 || writeExact( pipe, data, size ) ); + if ( written ) + FlushFileBuffers( pipe ); + CloseHandle( pipe ); + return written ? Status::Done : Status::Disconnected; +#else + sockaddr_un address; + if ( !makeAddress( endpoint, address ) ) + return Status::Error; + int fd = -1; + const Status connected = connectSocket( address, timeout, fd ); + if ( connected != Status::Done ) + return connected; + const bool written = writeExact( fd, header.data(), header.size() ) && + ( size == 0 || writeExact( fd, data, size ) ); + ::close( fd ); + return written ? Status::Done : Status::Disconnected; +#endif +} + +IPC::Status IPC::send( std::string_view endpoint, std::string_view message, Time timeout ) { + return send( endpoint, message.data(), message.size(), timeout ); +} + +}} // namespace EE::System diff --git a/src/tests/unit_tests/ipc_tests.cpp b/src/tests/unit_tests/ipc_tests.cpp new file mode 100644 index 000000000..24abd1117 --- /dev/null +++ b/src/tests/unit_tests/ipc_tests.cpp @@ -0,0 +1,213 @@ +#include "utest.hpp" + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#if EE_PLATFORM != EE_PLATFORM_WIN && EE_PLATFORM != EE_PLATFORM_EMSCRIPTEN +#include +#include +#include +#endif + +using namespace EE; +using namespace EE::System; + +namespace { + +std::string uniqueEndpoint( const char* suffix ) { + static std::atomic counter{ 0 }; + return "eepp-test." + std::to_string( Sys::getProcessID() ) + "." + suffix + "." + + std::to_string( counter++ ); +} + +struct Messages { + void receive( const void* data, std::size_t size ) { + std::lock_guard lock( mutex ); + values.emplace_back( static_cast( data ), size ); + condition.notify_all(); + } + + bool waitFor( std::size_t count ) { + std::unique_lock lock( mutex ); + return condition.wait_for( lock, std::chrono::seconds( 3 ), + [&] { return values.size() >= count; } ); + } + + std::mutex mutex; + std::condition_variable condition; + std::vector values; +}; + +#if EE_PLATFORM != EE_PLATFORM_WIN && EE_PLATFORM != EE_PLATFORM_EMSCRIPTEN + +std::string nativePath( const std::string& endpoint ) { + std::string directory = Sys::getTempPath(); + if ( directory.back() != '/' ) + directory += '/'; + return directory + "eepp-ipc-" + std::to_string( static_cast( getuid() ) ) + + '/' + MD5::fromString( endpoint ).toHexString(); +} + +int rawConnect( const std::string& endpoint ) { + const int fd = socket( AF_UNIX, SOCK_STREAM, 0 ); + if ( fd < 0 ) + return -1; + sockaddr_un address{}; + address.sun_family = AF_UNIX; + const std::string path = nativePath( endpoint ); + std::memcpy( address.sun_path, path.c_str(), path.size() + 1 ); + if ( connect( fd, reinterpret_cast( &address ), sizeof( address ) ) != 0 ) { + close( fd ); + return -1; + } + return fd; +} + +bool sendRaw( const std::string& endpoint, const std::vector& bytes ) { + const int fd = rawConnect( endpoint ); + if ( fd < 0 ) + return false; + const bool sent = + static_cast( bytes.size() ) == write( fd, bytes.data(), bytes.size() ); + close( fd ); + return sent; +} + +#endif + +} // namespace + +UTEST( IPC, BasicBinaryAndEmptyMessages ) { + IPC listener; + Messages messages; + const std::string endpoint = uniqueEndpoint( "basic" ); + ASSERT_EQ( IPC::Status::Done, + listener.listen( endpoint, [&]( const void* data, std::size_t size ) { + messages.receive( data, size ); + } ) ); + + EXPECT_EQ( IPC::Status::Done, IPC::send( endpoint, "hello" ) ); + const std::vector binary{ 0, 1, 2, 0, 255, 3 }; + EXPECT_EQ( IPC::Status::Done, IPC::send( endpoint, binary.data(), binary.size() ) ); + const std::string large( 4096, 'x' ); + EXPECT_EQ( IPC::Status::Done, IPC::send( endpoint, large ) ); + EXPECT_EQ( IPC::Status::Done, IPC::send( endpoint, nullptr, 0 ) ); + ASSERT_TRUE( messages.waitFor( 4 ) ); + ASSERT_EQ( static_cast( 4 ), messages.values.size() ); + EXPECT_TRUE( std::string( "hello" ) == messages.values[0] ); + EXPECT_TRUE( std::string( reinterpret_cast( binary.data() ), binary.size() ) == + messages.values[1] ); + EXPECT_TRUE( large == messages.values[2] ); + EXPECT_TRUE( messages.values[3].empty() ); +} + +UTEST( IPC, SequentialClientsPreserveOrder ) { + IPC listener; + Messages messages; + const std::string endpoint = uniqueEndpoint( "sequence" ); + ASSERT_EQ( IPC::Status::Done, + listener.listen( endpoint, [&]( const void* data, std::size_t size ) { + messages.receive( data, size ); + } ) ); + for ( int i = 0; i < 32; ++i ) + ASSERT_EQ( IPC::Status::Done, IPC::send( endpoint, std::to_string( i ) ) ); + ASSERT_TRUE( messages.waitFor( 32 ) ); + for ( int i = 0; i < 32; ++i ) + EXPECT_TRUE( std::to_string( i ) == messages.values[i] ); +} + +UTEST( IPC, OwnershipIsolationMissingAndOversized ) { + IPC first; + IPC duplicate; + IPC isolated; + const std::string endpoint = uniqueEndpoint( "ownership" ); + ASSERT_EQ( IPC::Status::Done, first.listen( endpoint, []( const void*, std::size_t ) {} ) ); + EXPECT_EQ( IPC::Status::AlreadyExists, + duplicate.listen( endpoint, []( const void*, std::size_t ) {} ) ); + EXPECT_EQ( IPC::Status::Done, + isolated.listen( endpoint + ".other", []( const void*, std::size_t ) {} ) ); + EXPECT_EQ( IPC::Status::NotFound, IPC::send( uniqueEndpoint( "missing" ), "message" ) ); + const Uint8 byte = 0; + EXPECT_EQ( IPC::Status::MessageTooLarge, + IPC::send( endpoint, &byte, IPC::MaxMessageSize + 1 ) ); +} + +UTEST( IPC, CloseAndRelistenWhileBlocked ) { + IPC listener; + const std::string endpoint = uniqueEndpoint( "relisten" ); + ASSERT_EQ( IPC::Status::Done, listener.listen( endpoint, []( const void*, std::size_t ) {} ) ); + listener.close(); + EXPECT_FALSE( listener.isListening() ); + ASSERT_EQ( IPC::Status::Done, listener.listen( endpoint, []( const void*, std::size_t ) {} ) ); + listener.close(); +} + +#if EE_PLATFORM != EE_PLATFORM_WIN && EE_PLATFORM != EE_PLATFORM_EMSCRIPTEN + +UTEST( IPC, RejectsInvalidFramesAndPartialDisconnects ) { + IPC listener; + Messages messages; + const std::string endpoint = uniqueEndpoint( "frames" ); + ASSERT_EQ( IPC::Status::Done, + listener.listen( endpoint, [&]( const void* data, std::size_t size ) { + messages.receive( data, size ); + } ) ); + + ASSERT_TRUE( sendRaw( endpoint, { 0, 0, 0, 0, 0, 1, 0, 0, 0, 0, 0, 0 } ) ); + ASSERT_TRUE( sendRaw( endpoint, { 'E', 'E', 'I', 'P', 0, 2, 0, 0, 0, 0, 0, 0 } ) ); + ASSERT_TRUE( sendRaw( endpoint, { 'E', 'E', 'I', 'P', 0, 1, 0, 1, 0, 0, 0, 0 } ) ); + ASSERT_TRUE( sendRaw( endpoint, { 'E', 'E', 'I', 'P', 0, 1, 0, 0, 1, 0, 0, 1 } ) ); + ASSERT_TRUE( sendRaw( endpoint, { 'E', 'E', 'I', 'P', 0, 1, 0, 0, 0, 0 } ) ); + ASSERT_TRUE( sendRaw( endpoint, { 'E', 'E', 'I', 'P', 0, 1, 0, 0, 0, 0, 0, 8, 1, 2 } ) ); + + EXPECT_EQ( IPC::Status::Done, IPC::send( endpoint, "still-alive" ) ); + ASSERT_TRUE( messages.waitFor( 1 ) ); + EXPECT_TRUE( std::string( "still-alive" ) == messages.values[0] ); +} + +UTEST( IPC, CloseInterruptsPartialClient ) { + IPC listener; + const std::string endpoint = uniqueEndpoint( "partial-close" ); + ASSERT_EQ( IPC::Status::Done, listener.listen( endpoint, []( const void*, std::size_t ) {} ) ); + const int fd = rawConnect( endpoint ); + ASSERT_TRUE( fd >= 0 ); + const Uint8 partial[]{ 'E', 'E' }; + ASSERT_EQ( static_cast( sizeof( partial ) ), write( fd, partial, sizeof( partial ) ) ); + std::this_thread::sleep_for( std::chrono::milliseconds( 20 ) ); + listener.close(); + close( fd ); +} + +UTEST( IPC, RecoversStaleSocketWithoutRemovingLiveSocket ) { + const std::string endpoint = uniqueEndpoint( "stale" ); + IPC directoryCreator; + ASSERT_EQ( IPC::Status::Done, directoryCreator.listen( endpoint + ".prepare", + []( const void*, std::size_t ) {} ) ); + directoryCreator.close(); + + const int stale = socket( AF_UNIX, SOCK_STREAM, 0 ); + ASSERT_TRUE( stale >= 0 ); + sockaddr_un address{}; + address.sun_family = AF_UNIX; + const std::string path = nativePath( endpoint ); + std::memcpy( address.sun_path, path.c_str(), path.size() + 1 ); + ASSERT_EQ( 0, bind( stale, reinterpret_cast( &address ), sizeof( address ) ) ); + close( stale ); + + IPC listener; + ASSERT_EQ( IPC::Status::Done, listener.listen( endpoint, []( const void*, std::size_t ) {} ) ); + IPC duplicate; + EXPECT_EQ( IPC::Status::AlreadyExists, + duplicate.listen( endpoint, []( const void*, std::size_t ) {} ) ); +} + +#endif diff --git a/src/tools/ecode/ecode.cpp b/src/tools/ecode/ecode.cpp index 858f456ac..d9fd15896 100644 --- a/src/tools/ecode/ecode.cpp +++ b/src/tools/ecode/ecode.cpp @@ -18,6 +18,9 @@ #include "version.hpp" #include #include +#include +#include +#include #include #include #include @@ -66,6 +69,33 @@ App* appInstance = nullptr; static const Uint32 APP_LAYOUT_STYLE_MARKER = String::hash( "app_layout_style" ); static const auto NOT_UNIQUE_FILENAME = "not_unique"; +struct IPCEndpoint { + operator std::string_view() const { return { value.data(), size }; } + + std::array value{}; + Uint8 size{ 0 }; +}; + +static_assert( sizeof( IPCEndpoint ) == 81 ); + +static IPCEndpoint ipcEndpoint( const std::string& profileId, Uint64 pid ) { + IPCEndpoint endpoint; + constexpr std::string_view Prefix = "ecode."; + constexpr std::size_t MaxPidDigits = 20; + if ( Prefix.size() + profileId.size() + 1 + MaxPidDigits > endpoint.value.size() ) + return endpoint; + std::memcpy( endpoint.value.data(), Prefix.data(), Prefix.size() ); + endpoint.size = static_cast( Prefix.size() ); + std::memcpy( endpoint.value.data() + endpoint.size, profileId.data(), profileId.size() ); + endpoint.size += profileId.size(); + endpoint.value[endpoint.size++] = '.'; + const auto result = std::to_chars( endpoint.value.data() + endpoint.size, + endpoint.value.data() + endpoint.value.size(), pid ); + endpoint.size = + result.ec == std::errc{} ? static_cast( result.ptr - endpoint.value.data() ) : 0; + return endpoint; +} + void appLoop() { appInstance->mainLoop(); } @@ -802,7 +832,7 @@ bool App::loadConfig( const LogLevel& logLevel, const Sizeu& displaySize, bool s mThemesPath = mConfigPath + "themes"; mScriptsPath = mConfigPath + "scripts"; mPlaygroundPath = mConfigPath + "playground"; - mIpcPath = mConfigPath + "ipc"; + mProfileId = MD5::fromString( FileSystem::getRealPath( mConfigPath ) ).toHexString(); mColorSchemesPath = mConfigPath + "editor" + FileSystem::getOSSlash() + "colorschemes" + FileSystem::getOSSlash(); mTerminalManager = std::make_unique( this ); @@ -831,16 +861,6 @@ bool App::loadConfig( const LogLevel& logLevel, const Sizeu& displaySize, bool s FileSystem::makeDir( mPlaygroundPath ); FileSystem::dirAddSlashAtEnd( mPlaygroundPath ); - if ( !FileSystem::fileExists( mIpcPath ) ) - FileSystem::makeDir( mIpcPath ); - FileSystem::dirAddSlashAtEnd( mIpcPath ); - - Uint64 pid = Sys::getProcessID(); - mPidPath = mIpcPath + String::toString( pid ); - FileSystem::dirAddSlashAtEnd( mPidPath ); - if ( !FileSystem::fileExists( mPidPath ) ) - FileSystem::makeDir( mPidPath ); - mLogsPath = mConfigPath + "ecode.log"; Log::create( mLogsPath, logLevel, stdOutLogs, !disableFileLogs ); @@ -1057,6 +1077,7 @@ App::~App() { mLifetime.invalidate(); appInstance = nullptr; mDestroyingApp = true; + mIPC.close(); if ( mProjectBuildManager ) mProjectBuildManager.reset(); @@ -1075,16 +1096,11 @@ App::~App() { eeSAFE_DELETE( mSplitter ); if ( mFileSystemListener ) { - if ( mIpcListenerId ) - mFileSystemListener->removeListener( mIpcListenerId ); delete mFileSystemListener; mFileSystemListener = nullptr; } mDirTree.reset(); - if ( !FileSystem::dirRemoveAll( mPidPath ) ) - Log::warning( "Failed to remove directory \"%s\"", mPidPath ); - if ( mFirstInstance ) FileSystem::fileRemove( firstInstanceIndicatorPath() ); } @@ -4463,35 +4479,34 @@ bool App::needsRedirectToRunningProcess( std::string file ) { bool useFirstInstance = FileSystem::fileExists( firstInstanceIndicatorPath() ); Uint64 processPid = Sys::getProcessID(); - Uint64 selectedPid = processPid; - Uint64 selCreationTime = useFirstInstance ? std::numeric_limits::max() : 0; - + SmallVector, 4> candidates; + candidates.reserve( pids.size() ); for ( const auto pid : pids ) { if ( pid == processPid ) continue; + candidates.emplace_back( Sys::getProcessCreationTime( pid ), pid ); + } + std::sort( candidates.begin(), candidates.end(), + [useFirstInstance]( const auto& left, const auto& right ) { + return useFirstInstance ? left.first < right.first : left.first > right.first; + } ); - Uint64 creationTime = Sys::getProcessCreationTime( pid ); - - bool shouldUpdate = ( useFirstInstance && creationTime <= selCreationTime ) || - ( !useFirstInstance && creationTime >= selCreationTime ); - - if ( shouldUpdate ) { - selectedPid = pid; - selCreationTime = creationTime; + json message{ { "type", "open" }, { "path", finfo.getFilepath() } }; + if ( position.isValid() ) { + message["line"] = position.line(); + message["column"] = position.column(); + } + const std::string payload = message.dump(); + for ( const auto& candidate : candidates ) { + const auto status = IPC::send( ipcEndpoint( mProfileId, candidate.second ), payload ); + if ( status == IPC::Status::Done ) + return true; + if ( status != IPC::Status::NotFound ) { + Log::warning( "Failed to send IPC message to process %llu", + static_cast( candidate.second ) ); } } - - if ( selectedPid == processPid ) - return false; - - std::string pidPath = mIpcPath + String::toString( selectedPid ); - if ( !FileSystem::isDirectory( pidPath ) ) - return false; - FileSystem::dirAddSlashAtEnd( pidPath ); - FileSystem::fileWrite( pidPath + MD5::fromString( finfo.getFilepath() ).toHexString(), - finfo.getFilepath() + - ( position.isValid() ? position.toPositionString() : "" ) ); - return true; + return false; } void App::tintTitleBar() { @@ -5114,47 +5129,37 @@ void App::init( InitParameters& params ) { mFileWatcher = new efsw::FileWatcher(); mFileSystemListener = new FileSystemListener( mSplitter, mFileSystemModel, { mLogsPath } ); mFileWatcher->addWatch( mPluginsPath, mFileSystemListener ); - mFileWatcher->addWatch( mPidPath, mFileSystemListener ); mFileWatcher->watch(); mPluginManager->setFileSystemListener( mFileSystemListener ); - FileSystemListener::ListenerOptions ipcListenerOptions; - FileSystemListenerFilter ipcListenerFilter; - ipcListenerFilter.eventTypes = - FileSystemListener::eventTypeMask( FileSystemEventType::Add ) | - FileSystemListener::eventTypeMask( FileSystemEventType::Modified ); - ipcListenerFilter.path = mPidPath; - ipcListenerOptions.filters.emplace_back( std::move( ipcListenerFilter ) ); - ipcListenerOptions.affinity = FileSystemListener::ThreadAffinity::Worker; - mIpcListenerId = mFileSystemListener->addListener( - [this]( const FileEvent&, const FileInfo& fi ) { - std::string path; - FileSystem::fileGet( fi.getFilepath(), path ); - String::trimInPlace( path, ' ' ); - String::trimInPlace( path, '\n' ); - - bool hasPosition = pathHasPosition( path ); + const auto ipcStatus = mIPC.listen( + ipcEndpoint( mProfileId, Sys::getProcessID() ), + [this]( const void* data, std::size_t size ) { + json message = json::parse( + std::string_view( static_cast( data ), size ), nullptr, false ); + if ( message.is_discarded() || message.value( "type", "" ) != "open" || + !message.contains( "path" ) || !message["path"].is_string() ) + return; + std::string path = message["path"].get(); TextPosition initialPosition; - if ( hasPosition ) { - auto pathAndPosition = getPathAndPosition( path ); - path = pathAndPosition.first; - initialPosition = pathAndPosition.second; + if ( message.contains( "line" ) && message["line"].is_number_integer() ) { + initialPosition = TextPosition( message["line"].get(), + message.value( "column", 0 ) ); } - if ( FileSystem::fileExists( path ) ) { - mUISceneNode->runOnMainThread( [path, initialPosition, this] { - loadFileFromPathOrFocus( path, true, nullptr, - getForcePositionFn( initialPosition ) ); - - if ( !mWindow->hasFocus() ) { - if ( mWindow->isMinimized() ) - mWindow->restore(); - mWindow->raise(); - } - } ); + mUISceneNode->runOnMainThread( + [this, path = std::move( path ), initialPosition] { + loadFileFromPathOrFocus( path, true, nullptr, + getForcePositionFn( initialPosition ) ); + if ( !mWindow->hasFocus() ) { + if ( mWindow->isMinimized() ) + mWindow->restore(); + mWindow->raise(); + } + } ); } - FileSystem::fileRemove( fi.getFilepath() ); - }, - std::move( ipcListenerOptions ) ); + } ); + if ( ipcStatus != IPC::Status::Done ) + Log::warning( "Failed to initialize IPC listener" ); #endif mNotificationCenter = std::make_unique( diff --git a/src/tools/ecode/ecode.hpp b/src/tools/ecode/ecode.hpp index f2f240c2e..a1e563675 100644 --- a/src/tools/ecode/ecode.hpp +++ b/src/tools/ecode/ecode.hpp @@ -716,8 +716,8 @@ class App : public UICodeEditorSplitter::Client, public PluginContextProvider { std::string mi18nPath; std::string mScriptsPath; std::string mPlaygroundPath; - std::string mIpcPath; - std::string mPidPath; + std::string mProfileId; + IPC mIPC; Float mDisplayDPI{ 96 }; std::shared_ptr mThreadPool; std::shared_ptr mDirTree; @@ -782,7 +782,6 @@ class App : public UICodeEditorSplitter::Client, public PluginContextProvider { UIMenuBar* mMenuBar{ nullptr }; std::unique_ptr mSettingsActions; std::vector mPathsToLoad; - Uint64 mIpcListenerId{ 0 }; std::mutex mAsyncResourcesLoadMutex; std::condition_variable mAsyncResourcesLoadCond; std::vector mColorSchemes;