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
This commit is contained in:
Martín Lucas Golini
2026-09-02 02:28:35 -03:00
parent 51739165d0
commit fbfd1f3efa
6 changed files with 1004 additions and 77 deletions
+599
View File
@@ -0,0 +1,599 @@
#include <eepp/system/ipc.hpp>
#include <algorithm>
#include <array>
#include <atomic>
#include <cstring>
#include <eepp/system/md5.hpp>
#include <eepp/system/sys.hpp>
#include <memory>
#include <mutex>
#include <string>
#include <string_view>
#include <thread>
#if EE_PLATFORM == EE_PLATFORM_WIN
#include <aclapi.h>
#include <sddl.h>
#include <windows.h>
#elif EE_PLATFORM != EE_PLATFORM_EMSCRIPTEN
#include <cerrno>
#include <fcntl.h>
#include <poll.h>
#include <sys/socket.h>
#include <sys/stat.h>
#include <sys/un.h>
#include <unistd.h>
#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<char, 33>;
using NativeName = std::array<wchar_t, 64>;
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<Uint8, IPCInlinePayloadSize> mInline;
std::unique_ptr<Uint8[]> mHeap;
};
static_assert( sizeof( MessageBuffer ) <= IPCInlinePayloadSize + sizeof( void* ) );
std::array<Uint8, IPCHeaderSize> makeHeader( std::size_t payloadSize ) {
std::array<Uint8, IPCHeaderSize> header{};
const Uint32 size = static_cast<Uint32>( payloadSize );
header[0] = static_cast<Uint8>( IPCMagic >> 24 );
header[1] = static_cast<Uint8>( IPCMagic >> 16 );
header[2] = static_cast<Uint8>( IPCMagic >> 8 );
header[3] = static_cast<Uint8>( IPCMagic );
header[4] = static_cast<Uint8>( IPCVersion >> 8 );
header[5] = static_cast<Uint8>( IPCVersion );
header[6] = static_cast<Uint8>( IPCFlags >> 8 );
header[7] = static_cast<Uint8>( IPCFlags );
header[8] = static_cast<Uint8>( size >> 24 );
header[9] = static_cast<Uint8>( size >> 16 );
header[10] = static_cast<Uint8>( size >> 8 );
header[11] = static_cast<Uint8>( size );
return header;
}
bool decodeHeader( const std::array<Uint8, IPCHeaderSize>& wire, Uint32& payloadSize ) {
const Uint32 magic = static_cast<Uint32>( wire[0] ) << 24 |
static_cast<Uint32>( wire[1] ) << 16 |
static_cast<Uint32>( wire[2] ) << 8 | wire[3];
const Uint16 version = static_cast<Uint16>( wire[4] ) << 8 | wire[5];
const Uint16 flags = static_cast<Uint16>( wire[6] ) << 8 | wire[7];
payloadSize = static_cast<Uint32>( wire[8] ) << 24 | static_cast<Uint32>( wire[9] ) << 16 |
static_cast<Uint32>( 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<bool> 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<wchar_t>( 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<Uint8*>( data );
while ( size > 0 ) {
DWORD read = 0;
const DWORD chunk = static_cast<DWORD>( std::min<std::size_t>( 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<const Uint8*>( data );
while ( size > 0 ) {
DWORD written = 0;
const DWORD chunk = static_cast<DWORD>( std::min<std::size_t>( 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<unsigned long long>( 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<Uint8*>( 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<std::size_t>( readSize );
}
return true;
}
bool writeExact( int fd, const void* data, std::size_t size ) {
const auto* bytes = static_cast<const Uint8*>( 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<std::size_t>( 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<const sockaddr*>( &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<int>( eemax<double>( 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<Impl>() ) {}
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<std::mutex> 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<Uint8, IPCHeaderSize> 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<std::mutex> 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<sockaddr*>( &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<sockaddr*>( &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<std::mutex> lock( impl->handlesMutex );
impl->client = client;
}
if ( !impl->running ) {
{
std::lock_guard<std::mutex> lock( impl->handlesMutex );
impl->client = -1;
}
::close( client );
break;
}
std::array<Uint8, IPCHeaderSize> 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<std::mutex> 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<std::mutex> lock( mImpl->handlesMutex );
CancelSynchronousIo( reinterpret_cast<HANDLE>( 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<std::mutex> 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<std::mutex> 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<Uint64>( eemax<double>( 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
+213
View File
@@ -0,0 +1,213 @@
#include "utest.hpp"
#include <atomic>
#include <chrono>
#include <condition_variable>
#include <cstring>
#include <eepp/system/ipc.hpp>
#include <eepp/system/md5.hpp>
#include <eepp/system/sys.hpp>
#include <mutex>
#include <string>
#include <thread>
#include <vector>
#if EE_PLATFORM != EE_PLATFORM_WIN && EE_PLATFORM != EE_PLATFORM_EMSCRIPTEN
#include <sys/socket.h>
#include <sys/un.h>
#include <unistd.h>
#endif
using namespace EE;
using namespace EE::System;
namespace {
std::string uniqueEndpoint( const char* suffix ) {
static std::atomic<unsigned int> 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<std::mutex> lock( mutex );
values.emplace_back( static_cast<const char*>( data ), size );
condition.notify_all();
}
bool waitFor( std::size_t count ) {
std::unique_lock<std::mutex> lock( mutex );
return condition.wait_for( lock, std::chrono::seconds( 3 ),
[&] { return values.size() >= count; } );
}
std::mutex mutex;
std::condition_variable condition;
std::vector<std::string> 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<unsigned long long>( 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<sockaddr*>( &address ), sizeof( address ) ) != 0 ) {
close( fd );
return -1;
}
return fd;
}
bool sendRaw( const std::string& endpoint, const std::vector<Uint8>& bytes ) {
const int fd = rawConnect( endpoint );
if ( fd < 0 )
return false;
const bool sent =
static_cast<ssize_t>( 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<Uint8> 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<std::size_t>( 4 ), messages.values.size() );
EXPECT_TRUE( std::string( "hello" ) == messages.values[0] );
EXPECT_TRUE( std::string( reinterpret_cast<const char*>( 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<ssize_t>( 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<sockaddr*>( &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
+79 -74
View File
@@ -18,6 +18,9 @@
#include "version.hpp"
#include <algorithm>
#include <args/args.hxx>
#include <array>
#include <charconv>
#include <cstring>
#include <eepp/graphics/fontfamily.hpp>
#include <eepp/system/iostreammemory.hpp>
#include <eepp/ui/doc/languagessyntaxhighlighting.hpp>
@@ -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<char, 80> 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<Uint8>( 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<Uint8>( 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<TerminalManager>( 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<Uint64>::max() : 0;
SmallVector<std::pair<Uint64, Uint64>, 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<unsigned long long>( 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<const char*>( 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<std::string>();
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<Int64>(),
message.value<Int64>( "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<NotificationCenter>(
+2 -3
View File
@@ -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<ThreadPool> mThreadPool;
std::shared_ptr<ProjectDirectoryTree> mDirTree;
@@ -782,7 +782,6 @@ class App : public UICodeEditorSplitter::Client, public PluginContextProvider {
UIMenuBar* mMenuBar{ nullptr };
std::unique_ptr<SettingsActions> mSettingsActions;
std::vector<std::string> mPathsToLoad;
Uint64 mIpcListenerId{ 0 };
std::mutex mAsyncResourcesLoadMutex;
std::condition_variable mAsyncResourcesLoadCond;
std::vector<SyntaxColorScheme> mColorSchemes;