Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
277 changes: 43 additions & 234 deletions src/coreclr/src/debug/debug-pal/unix/diagnosticsipc.cpp
Original file line numberDiff line numberDiff line change
Expand Up@@ -12,20 +12,13 @@
#include "diagnosticsipc.h"
#include "processdescriptor.h"

#if __GNUC__
#include <poll.h>
#else
#include <sys/poll.h>
#endif // __GNUC__

IpcStream::DiagnosticsIpc::DiagnosticsIpc(const int serverSocket, sockaddr_un *const pServerAddress, ConnectionMode mode) :
mode(mode),
IpcStream::DiagnosticsIpc::DiagnosticsIpc(const int serverSocket, sockaddr_un *const pServerAddress) :
_serverSocket(serverSocket),
_pServerAddress(new sockaddr_un),
_isClosed(false),
_isListening(false)
_isClosed(false)
{
_ASSERTE(_pServerAddress != nullptr);
_ASSERTE(_serverSocket != -1);
_ASSERTE(pServerAddress != nullptr);

if (_pServerAddress == nullptr || pServerAddress == nullptr)
Expand All@@ -39,8 +32,24 @@ IpcStream::DiagnosticsIpc::~DiagnosticsIpc()
delete _pServerAddress;
}

IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const pIpcName, ConnectionMode mode, ErrorCallback callback)
IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const pIpcName, ErrorCallback callback)
{
#ifdef __APPLE__
mode_t prev_mask = umask(~(S_IRUSR | S_IWUSR)); // This will set the default permission bit to 600
#endif // __APPLE__

const int serverSocket = ::socket(AF_UNIX, SOCK_STREAM, 0);
if (serverSocket == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__
_ASSERTE(!"Failed to create diagnostics IPC socket.");
return nullptr;
}

sockaddr_un serverAddress{};
serverAddress.sun_family = AF_UNIX;

Expand All@@ -62,24 +71,6 @@ IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const p
"socket");
}

if (mode == ConnectionMode::CLIENT)
return new IpcStream::DiagnosticsIpc(-1, &serverAddress, ConnectionMode::CLIENT);

#ifdef __APPLE__
mode_t prev_mask = umask(~(S_IRUSR | S_IWUSR)); // This will set the default permission bit to 600
#endif // __APPLE__

const int serverSocket = ::socket(AF_UNIX, SOCK_STREAM, 0);
if (serverSocket == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__
_ASSERTE(!"Failed to create diagnostics IPC socket.");
return nullptr;
}

#ifndef __APPLE__
if (fchmod(serverSocket, S_IRUSR | S_IWUSR) == -1)
Expand DownExpand Up@@ -108,52 +99,33 @@ IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const p
return nullptr;
}

#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__

return new IpcStream::DiagnosticsIpc(serverSocket, &serverAddress, mode);
}

bool IpcStream::DiagnosticsIpc::Listen(ErrorCallback callback)
{
_ASSERTE(mode == ConnectionMode::SERVER);
if (mode != ConnectionMode::SERVER)
{
if (callback != nullptr)
callback("Cannot call Listen on a client connection", -1);
return false;
}

if (_isListening)
return true;

const int fSuccessfulListen = ::listen(_serverSocket, /* backlog */ 255);
const int fSuccessfulListen = ::listen(serverSocket, /* backlog */ 255);
if (fSuccessfulListen == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
_ASSERTE(fSuccessfulListen != -1);

const int fSuccessUnlink = ::unlink(_pServerAddress->sun_path);
const int fSuccessUnlink = ::unlink(serverAddress.sun_path);
_ASSERTE(fSuccessUnlink != -1);

const int fSuccessClose = ::close(_serverSocket);
const int fSuccessClose = ::close(serverSocket);
_ASSERTE(fSuccessClose != -1);
return false;
}
else
{
_isListening = true;
return true;
#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__
return nullptr;
}

#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__

return new IpcStream::DiagnosticsIpc(serverSocket, &serverAddress);
}

IpcStream *IpcStream::DiagnosticsIpc::Accept(ErrorCallback callback)
IpcStream *IpcStream::DiagnosticsIpc::Accept(ErrorCallback callback) const
{
_ASSERTE(mode == ConnectionMode::SERVER);
_ASSERTE(_isListening);

sockaddr_un from;
socklen_t fromlen = sizeof(from);
const int clientSocket = ::accept(_serverSocket, (sockaddr *)&from, &fromlen);
Expand All@@ -164,114 +136,7 @@ IpcStream *IpcStream::DiagnosticsIpc::Accept(ErrorCallback callback)
return nullptr;
}

return new IpcStream(clientSocket, mode);
}

IpcStream *IpcStream::DiagnosticsIpc::Connect(ErrorCallback callback)
{
_ASSERTE(mode == ConnectionMode::CLIENT);

sockaddr_un clientAddress{};
clientAddress.sun_family = AF_UNIX;
const int clientSocket = ::socket(AF_UNIX, SOCK_STREAM, 0);
if (clientSocket == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
return nullptr;
}

// We don't expect this to block since this is a Unix Domain Socket. `connect` may block until the
// TCP handshake is complete for TCP/IP sockets, but UDS don't use TCP. `connect` will return even if
// the server hasn't called `accept`.
if (::connect(clientSocket, (struct sockaddr *)_pServerAddress, sizeof(*_pServerAddress)) < 0)
{
if (callback != nullptr)
callback(strerror(errno), errno);
return nullptr;
}

return new IpcStream(clientSocket, ConnectionMode::CLIENT);
}

int32_t IpcStream::DiagnosticsIpc::Poll(IpcPollHandle *rgIpcPollHandles, uint32_t nHandles, int32_t timeoutMs, ErrorCallback callback)
{
// prepare the pollfd structs
pollfd *pollfds = new pollfd[nHandles];
for (uint32_t i = 0; i < nHandles; i++)
{
rgIpcPollHandles[i].revents = 0; // ignore any values in revents
int fd = -1;
if (rgIpcPollHandles[i].pIpc != nullptr)
{
// SERVER
_ASSERTE(rgIpcPollHandles[i].pIpc->mode == ConnectionMode::SERVER);
fd = rgIpcPollHandles[i].pIpc->_serverSocket;
}
else
{
// CLIENT
_ASSERTE(rgIpcPollHandles[i].pStream != nullptr);
fd = rgIpcPollHandles[i].pStream->_clientSocket;
}

pollfds[i].fd = fd;
pollfds[i].events = POLLIN;
}

int retval = poll(pollfds, nHandles, timeoutMs);

// Check results
if (retval < 0)
{
for (uint32_t i = 0; i < nHandles; i++)
{
if ((pollfds[i].revents & POLLERR) && callback != nullptr)
callback(strerror(errno), errno);
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::ERR;
}
delete[] pollfds;
return -1;
}
else if (retval == 0)
{
// we timed out
delete[] pollfds;
return 0;
}

for (uint32_t i = 0; i < nHandles; i++)
{
if (pollfds[i].revents != 0)
{
// error check FIRST
if (pollfds[i].revents & POLLHUP)
{
// check for hangup first because a closed socket
// will technically meet the requirements for POLLIN
// i.e., a call to recv/read won't block
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::HANGUP;
delete[] pollfds;
return -1;
}
else if ((pollfds[i].revents & (POLLERR|POLLNVAL)))
{
if (callback != nullptr)
callback("Poll error", (uint32_t)pollfds[i].revents);
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::ERR;
delete[] pollfds;
return -1;
}
else if (pollfds[i].revents & POLLIN)
{
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::SIGNALED;
break;
}
}
}

delete[] pollfds;
return 1;
return new IpcStream(clientSocket);
}

void IpcStream::DiagnosticsIpc::Close(ErrorCallback callback)
Expand DownExpand Up@@ -307,101 +172,45 @@ void IpcStream::DiagnosticsIpc::Unlink(ErrorCallback callback)
}

IpcStream::~IpcStream()
{
Close();
}

void IpcStream::Close(ErrorCallback)
{
if (_clientSocket != -1)
{
Flush();

const int fSuccessClose = ::close(_clientSocket);
_ASSERTE(fSuccessClose != -1);
_clientSocket = -1;
}
}

bool IpcStream::Read(void *lpBuffer, const uint32_t nBytesToRead, uint32_t &nBytesRead, const int32_t timeoutMs)
bool IpcStream::Read(void *lpBuffer, const uint32_t nBytesToRead, uint32_t &nBytesRead) const
{
_ASSERTE(lpBuffer != nullptr);

if (timeoutMs != InfiniteTimeout)
{
pollfd pfd;
pfd.fd = _clientSocket;
pfd.events = POLLIN;
int retval = poll(&pfd, 1, timeoutMs);
if (retval <= 0 || pfd.revents != POLLIN)
{
// timeout or error
return false;
}
// else fallthrough
}

uint8_t *lpBufferCursor = (uint8_t*)lpBuffer;
ssize_t currentBytesRead = 0;
ssize_t totalBytesRead = 0;
bool fSuccess = true;
while (fSuccess && nBytesToRead - totalBytesRead > 0)
{
currentBytesRead = ::recv(_clientSocket, lpBufferCursor, nBytesToRead - totalBytesRead, 0);
fSuccess = currentBytesRead != 0;
if (!fSuccess)
break;
totalBytesRead += currentBytesRead;
lpBufferCursor += currentBytesRead;
}
const ssize_t ssize = ::recv(_clientSocket, lpBuffer, nBytesToRead, 0);
const bool fSuccess = ssize != -1;

if (!fSuccess)
{
// TODO: Add error handling.
}

nBytesRead = static_cast<uint32_t>(totalBytesRead);
nBytesRead = static_cast<uint32_t>(ssize);
return fSuccess;
}

bool IpcStream::Write(const void *lpBuffer, const uint32_t nBytesToWrite, uint32_t &nBytesWritten, const int32_t timeoutMs)
bool IpcStream::Write(const void *lpBuffer, const uint32_t nBytesToWrite, uint32_t &nBytesWritten) const
{
_ASSERTE(lpBuffer != nullptr);

if (timeoutMs != InfiniteTimeout)
{
pollfd pfd;
pfd.fd = _clientSocket;
pfd.events = POLLOUT;
int retval = poll(&pfd, 1, timeoutMs);
if (retval <= 0 || pfd.revents != POLLOUT)
{
// timeout or error
return false;
}
// else fallthrough
}

uint8_t *lpBufferCursor = (uint8_t*)lpBuffer;
ssize_t currentBytesWritten = 0;
ssize_t totalBytesWritten = 0;
bool fSuccess = true;
while (fSuccess && nBytesToWrite - totalBytesWritten > 0)
{
currentBytesWritten = ::send(_clientSocket, lpBufferCursor, nBytesToWrite - totalBytesWritten, 0);
fSuccess = currentBytesWritten != -1;
if (!fSuccess)
break;
lpBufferCursor += currentBytesWritten;
totalBytesWritten += currentBytesWritten;
}
const ssize_t ssize = ::send(_clientSocket, lpBuffer, nBytesToWrite, 0);
const bool fSuccess = ssize != -1;

if (!fSuccess)
{
// TODO: Add error handling.
}

nBytesWritten = static_cast<uint32_t>(totalBytesWritten);
nBytesWritten = static_cast<uint32_t>(ssize);
return fSuccess;
}

Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
277 changes: 43 additions & 234 deletions src/coreclr/src/debug/debug-pal/unix/diagnosticsipc.cpp
Original file line numberDiff line numberDiff line change
Expand Up@@ -12,20 +12,13 @@
#include "diagnosticsipc.h"
#include "processdescriptor.h"

#if __GNUC__
#include <poll.h>
#else
#include <sys/poll.h>
#endif // __GNUC__

IpcStream::DiagnosticsIpc::DiagnosticsIpc(const int serverSocket, sockaddr_un *const pServerAddress, ConnectionMode mode) :
mode(mode),
IpcStream::DiagnosticsIpc::DiagnosticsIpc(const int serverSocket, sockaddr_un *const pServerAddress) :
_serverSocket(serverSocket),
_pServerAddress(new sockaddr_un),
_isClosed(false),
_isListening(false)
_isClosed(false)
{
_ASSERTE(_pServerAddress != nullptr);
_ASSERTE(_serverSocket != -1);
_ASSERTE(pServerAddress != nullptr);

if (_pServerAddress == nullptr || pServerAddress == nullptr)
Expand All@@ -39,8 +32,24 @@ IpcStream::DiagnosticsIpc::~DiagnosticsIpc()
delete _pServerAddress;
}

IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const pIpcName, ConnectionMode mode, ErrorCallback callback)
IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const pIpcName, ErrorCallback callback)
{
#ifdef __APPLE__
mode_t prev_mask = umask(~(S_IRUSR | S_IWUSR)); // This will set the default permission bit to 600
#endif // __APPLE__

const int serverSocket = ::socket(AF_UNIX, SOCK_STREAM, 0);
if (serverSocket == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__
_ASSERTE(!"Failed to create diagnostics IPC socket.");
return nullptr;
}

sockaddr_un serverAddress{};
serverAddress.sun_family = AF_UNIX;

Expand All@@ -62,24 +71,6 @@ IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const p
"socket");
}

if (mode == ConnectionMode::CLIENT)
return new IpcStream::DiagnosticsIpc(-1, &serverAddress, ConnectionMode::CLIENT);

#ifdef __APPLE__
mode_t prev_mask = umask(~(S_IRUSR | S_IWUSR)); // This will set the default permission bit to 600
#endif // __APPLE__

const int serverSocket = ::socket(AF_UNIX, SOCK_STREAM, 0);
if (serverSocket == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__
_ASSERTE(!"Failed to create diagnostics IPC socket.");
return nullptr;
}

#ifndef __APPLE__
if (fchmod(serverSocket, S_IRUSR | S_IWUSR) == -1)
Expand DownExpand Up@@ -108,52 +99,33 @@ IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const p
return nullptr;
}

#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__

return new IpcStream::DiagnosticsIpc(serverSocket, &serverAddress, mode);
}

bool IpcStream::DiagnosticsIpc::Listen(ErrorCallback callback)
{
_ASSERTE(mode == ConnectionMode::SERVER);
if (mode != ConnectionMode::SERVER)
{
if (callback != nullptr)
callback("Cannot call Listen on a client connection", -1);
return false;
}

if (_isListening)
return true;

const int fSuccessfulListen = ::listen(_serverSocket, /* backlog */ 255);
const int fSuccessfulListen = ::listen(serverSocket, /* backlog */ 255);
if (fSuccessfulListen == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
_ASSERTE(fSuccessfulListen != -1);

const int fSuccessUnlink = ::unlink(_pServerAddress->sun_path);
const int fSuccessUnlink = ::unlink(serverAddress.sun_path);
_ASSERTE(fSuccessUnlink != -1);

const int fSuccessClose = ::close(_serverSocket);
const int fSuccessClose = ::close(serverSocket);
_ASSERTE(fSuccessClose != -1);
return false;
}
else
{
_isListening = true;
return true;
#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__
return nullptr;
}

#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__

return new IpcStream::DiagnosticsIpc(serverSocket, &serverAddress);
}

IpcStream *IpcStream::DiagnosticsIpc::Accept(ErrorCallback callback)
IpcStream *IpcStream::DiagnosticsIpc::Accept(ErrorCallback callback) const
{
_ASSERTE(mode == ConnectionMode::SERVER);
_ASSERTE(_isListening);

sockaddr_un from;
socklen_t fromlen = sizeof(from);
const int clientSocket = ::accept(_serverSocket, (sockaddr *)&from, &fromlen);
Expand All@@ -164,114 +136,7 @@ IpcStream *IpcStream::DiagnosticsIpc::Accept(ErrorCallback callback)
return nullptr;
}

return new IpcStream(clientSocket, mode);
}

IpcStream *IpcStream::DiagnosticsIpc::Connect(ErrorCallback callback)
{
_ASSERTE(mode == ConnectionMode::CLIENT);

sockaddr_un clientAddress{};
clientAddress.sun_family = AF_UNIX;
const int clientSocket = ::socket(AF_UNIX, SOCK_STREAM, 0);
if (clientSocket == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
return nullptr;
}

// We don't expect this to block since this is a Unix Domain Socket. `connect` may block until the
// TCP handshake is complete for TCP/IP sockets, but UDS don't use TCP. `connect` will return even if
// the server hasn't called `accept`.
if (::connect(clientSocket, (struct sockaddr *)_pServerAddress, sizeof(*_pServerAddress)) < 0)
{
if (callback != nullptr)
callback(strerror(errno), errno);
return nullptr;
}

return new IpcStream(clientSocket, ConnectionMode::CLIENT);
}

int32_t IpcStream::DiagnosticsIpc::Poll(IpcPollHandle *rgIpcPollHandles, uint32_t nHandles, int32_t timeoutMs, ErrorCallback callback)
{
// prepare the pollfd structs
pollfd *pollfds = new pollfd[nHandles];
for (uint32_t i = 0; i < nHandles; i++)
{
rgIpcPollHandles[i].revents = 0; // ignore any values in revents
int fd = -1;
if (rgIpcPollHandles[i].pIpc != nullptr)
{
// SERVER
_ASSERTE(rgIpcPollHandles[i].pIpc->mode == ConnectionMode::SERVER);
fd = rgIpcPollHandles[i].pIpc->_serverSocket;
}
else
{
// CLIENT
_ASSERTE(rgIpcPollHandles[i].pStream != nullptr);
fd = rgIpcPollHandles[i].pStream->_clientSocket;
}

pollfds[i].fd = fd;
pollfds[i].events = POLLIN;
}

int retval = poll(pollfds, nHandles, timeoutMs);

// Check results
if (retval < 0)
{
for (uint32_t i = 0; i < nHandles; i++)
{
if ((pollfds[i].revents & POLLERR) && callback != nullptr)
callback(strerror(errno), errno);
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::ERR;
}
delete[] pollfds;
return -1;
}
else if (retval == 0)
{
// we timed out
delete[] pollfds;
return 0;
}

for (uint32_t i = 0; i < nHandles; i++)
{
if (pollfds[i].revents != 0)
{
// error check FIRST
if (pollfds[i].revents & POLLHUP)
{
// check for hangup first because a closed socket
// will technically meet the requirements for POLLIN
// i.e., a call to recv/read won't block
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::HANGUP;
delete[] pollfds;
return -1;
}
else if ((pollfds[i].revents & (POLLERR|POLLNVAL)))
{
if (callback != nullptr)
callback("Poll error", (uint32_t)pollfds[i].revents);
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::ERR;
delete[] pollfds;
return -1;
}
else if (pollfds[i].revents & POLLIN)
{
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::SIGNALED;
break;
}
}
}

delete[] pollfds;
return 1;
return new IpcStream(clientSocket);
}

void IpcStream::DiagnosticsIpc::Close(ErrorCallback callback)
Expand DownExpand Up@@ -307,101 +172,45 @@ void IpcStream::DiagnosticsIpc::Unlink(ErrorCallback callback)
}

IpcStream::~IpcStream()
{
Close();
}

void IpcStream::Close(ErrorCallback)
{
if (_clientSocket != -1)
{
Flush();

const int fSuccessClose = ::close(_clientSocket);
_ASSERTE(fSuccessClose != -1);
_clientSocket = -1;
}
}

bool IpcStream::Read(void *lpBuffer, const uint32_t nBytesToRead, uint32_t &nBytesRead, const int32_t timeoutMs)
bool IpcStream::Read(void *lpBuffer, const uint32_t nBytesToRead, uint32_t &nBytesRead) const
{
_ASSERTE(lpBuffer != nullptr);

if (timeoutMs != InfiniteTimeout)
{
pollfd pfd;
pfd.fd = _clientSocket;
pfd.events = POLLIN;
int retval = poll(&pfd, 1, timeoutMs);
if (retval <= 0 || pfd.revents != POLLIN)
{
// timeout or error
return false;
}
// else fallthrough
}

uint8_t *lpBufferCursor = (uint8_t*)lpBuffer;
ssize_t currentBytesRead = 0;
ssize_t totalBytesRead = 0;
bool fSuccess = true;
while (fSuccess && nBytesToRead - totalBytesRead > 0)
{
currentBytesRead = ::recv(_clientSocket, lpBufferCursor, nBytesToRead - totalBytesRead, 0);
fSuccess = currentBytesRead != 0;
if (!fSuccess)
break;
totalBytesRead += currentBytesRead;
lpBufferCursor += currentBytesRead;
}
const ssize_t ssize = ::recv(_clientSocket, lpBuffer, nBytesToRead, 0);
const bool fSuccess = ssize != -1;

if (!fSuccess)
{
// TODO: Add error handling.
}

nBytesRead = static_cast<uint32_t>(totalBytesRead);
nBytesRead = static_cast<uint32_t>(ssize);
return fSuccess;
}

bool IpcStream::Write(const void *lpBuffer, const uint32_t nBytesToWrite, uint32_t &nBytesWritten, const int32_t timeoutMs)
bool IpcStream::Write(const void *lpBuffer, const uint32_t nBytesToWrite, uint32_t &nBytesWritten) const
{
_ASSERTE(lpBuffer != nullptr);

if (timeoutMs != InfiniteTimeout)
{
pollfd pfd;
pfd.fd = _clientSocket;
pfd.events = POLLOUT;
int retval = poll(&pfd, 1, timeoutMs);
if (retval <= 0 || pfd.revents != POLLOUT)
{
// timeout or error
return false;
}
// else fallthrough
}

uint8_t *lpBufferCursor = (uint8_t*)lpBuffer;
ssize_t currentBytesWritten = 0;
ssize_t totalBytesWritten = 0;
bool fSuccess = true;
while (fSuccess && nBytesToWrite - totalBytesWritten > 0)
{
currentBytesWritten = ::send(_clientSocket, lpBufferCursor, nBytesToWrite - totalBytesWritten, 0);
fSuccess = currentBytesWritten != -1;
if (!fSuccess)
break;
lpBufferCursor += currentBytesWritten;
totalBytesWritten += currentBytesWritten;
}
const ssize_t ssize = ::send(_clientSocket, lpBuffer, nBytesToWrite, 0);
const bool fSuccess = ssize != -1;

if (!fSuccess)
{
// TODO: Add error handling.
}

nBytesWritten = static_cast<uint32_t>(totalBytesWritten);
nBytesWritten = static_cast<uint32_t>(ssize);
return fSuccess;
}

Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
277 changes: 43 additions & 234 deletions src/coreclr/src/debug/debug-pal/unix/diagnosticsipc.cpp
Original file line numberDiff line numberDiff line change
Expand Up@@ -12,20 +12,13 @@
#include "diagnosticsipc.h"
#include "processdescriptor.h"

#if __GNUC__
#include <poll.h>
#else
#include <sys/poll.h>
#endif // __GNUC__

IpcStream::DiagnosticsIpc::DiagnosticsIpc(const int serverSocket, sockaddr_un *const pServerAddress, ConnectionMode mode) :
mode(mode),
IpcStream::DiagnosticsIpc::DiagnosticsIpc(const int serverSocket, sockaddr_un *const pServerAddress) :
_serverSocket(serverSocket),
_pServerAddress(new sockaddr_un),
_isClosed(false),
_isListening(false)
_isClosed(false)
{
_ASSERTE(_pServerAddress != nullptr);
_ASSERTE(_serverSocket != -1);
_ASSERTE(pServerAddress != nullptr);

if (_pServerAddress == nullptr || pServerAddress == nullptr)
Expand All@@ -39,8 +32,24 @@ IpcStream::DiagnosticsIpc::~DiagnosticsIpc()
delete _pServerAddress;
}

IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const pIpcName, ConnectionMode mode, ErrorCallback callback)
IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const pIpcName, ErrorCallback callback)
{
#ifdef __APPLE__
mode_t prev_mask = umask(~(S_IRUSR | S_IWUSR)); // This will set the default permission bit to 600
#endif // __APPLE__

const int serverSocket = ::socket(AF_UNIX, SOCK_STREAM, 0);
if (serverSocket == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__
_ASSERTE(!"Failed to create diagnostics IPC socket.");
return nullptr;
}

sockaddr_un serverAddress{};
serverAddress.sun_family = AF_UNIX;

Expand All@@ -62,24 +71,6 @@ IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const p
"socket");
}

if (mode == ConnectionMode::CLIENT)
return new IpcStream::DiagnosticsIpc(-1, &serverAddress, ConnectionMode::CLIENT);

#ifdef __APPLE__
mode_t prev_mask = umask(~(S_IRUSR | S_IWUSR)); // This will set the default permission bit to 600
#endif // __APPLE__

const int serverSocket = ::socket(AF_UNIX, SOCK_STREAM, 0);
if (serverSocket == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__
_ASSERTE(!"Failed to create diagnostics IPC socket.");
return nullptr;
}

#ifndef __APPLE__
if (fchmod(serverSocket, S_IRUSR | S_IWUSR) == -1)
Expand DownExpand Up@@ -108,52 +99,33 @@ IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const p
return nullptr;
}

#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__

return new IpcStream::DiagnosticsIpc(serverSocket, &serverAddress, mode);
}

bool IpcStream::DiagnosticsIpc::Listen(ErrorCallback callback)
{
_ASSERTE(mode == ConnectionMode::SERVER);
if (mode != ConnectionMode::SERVER)
{
if (callback != nullptr)
callback("Cannot call Listen on a client connection", -1);
return false;
}

if (_isListening)
return true;

const int fSuccessfulListen = ::listen(_serverSocket, /* backlog */ 255);
const int fSuccessfulListen = ::listen(serverSocket, /* backlog */ 255);
if (fSuccessfulListen == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
_ASSERTE(fSuccessfulListen != -1);

const int fSuccessUnlink = ::unlink(_pServerAddress->sun_path);
const int fSuccessUnlink = ::unlink(serverAddress.sun_path);
_ASSERTE(fSuccessUnlink != -1);

const int fSuccessClose = ::close(_serverSocket);
const int fSuccessClose = ::close(serverSocket);
_ASSERTE(fSuccessClose != -1);
return false;
}
else
{
_isListening = true;
return true;
#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__
return nullptr;
}

#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__

return new IpcStream::DiagnosticsIpc(serverSocket, &serverAddress);
}

IpcStream *IpcStream::DiagnosticsIpc::Accept(ErrorCallback callback)
IpcStream *IpcStream::DiagnosticsIpc::Accept(ErrorCallback callback) const
{
_ASSERTE(mode == ConnectionMode::SERVER);
_ASSERTE(_isListening);

sockaddr_un from;
socklen_t fromlen = sizeof(from);
const int clientSocket = ::accept(_serverSocket, (sockaddr *)&from, &fromlen);
Expand All@@ -164,114 +136,7 @@ IpcStream *IpcStream::DiagnosticsIpc::Accept(ErrorCallback callback)
return nullptr;
}

return new IpcStream(clientSocket, mode);
}

IpcStream *IpcStream::DiagnosticsIpc::Connect(ErrorCallback callback)
{
_ASSERTE(mode == ConnectionMode::CLIENT);

sockaddr_un clientAddress{};
clientAddress.sun_family = AF_UNIX;
const int clientSocket = ::socket(AF_UNIX, SOCK_STREAM, 0);
if (clientSocket == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
return nullptr;
}

// We don't expect this to block since this is a Unix Domain Socket. `connect` may block until the
// TCP handshake is complete for TCP/IP sockets, but UDS don't use TCP. `connect` will return even if
// the server hasn't called `accept`.
if (::connect(clientSocket, (struct sockaddr *)_pServerAddress, sizeof(*_pServerAddress)) < 0)
{
if (callback != nullptr)
callback(strerror(errno), errno);
return nullptr;
}

return new IpcStream(clientSocket, ConnectionMode::CLIENT);
}

int32_t IpcStream::DiagnosticsIpc::Poll(IpcPollHandle *rgIpcPollHandles, uint32_t nHandles, int32_t timeoutMs, ErrorCallback callback)
{
// prepare the pollfd structs
pollfd *pollfds = new pollfd[nHandles];
for (uint32_t i = 0; i < nHandles; i++)
{
rgIpcPollHandles[i].revents = 0; // ignore any values in revents
int fd = -1;
if (rgIpcPollHandles[i].pIpc != nullptr)
{
// SERVER
_ASSERTE(rgIpcPollHandles[i].pIpc->mode == ConnectionMode::SERVER);
fd = rgIpcPollHandles[i].pIpc->_serverSocket;
}
else
{
// CLIENT
_ASSERTE(rgIpcPollHandles[i].pStream != nullptr);
fd = rgIpcPollHandles[i].pStream->_clientSocket;
}

pollfds[i].fd = fd;
pollfds[i].events = POLLIN;
}

int retval = poll(pollfds, nHandles, timeoutMs);

// Check results
if (retval < 0)
{
for (uint32_t i = 0; i < nHandles; i++)
{
if ((pollfds[i].revents & POLLERR) && callback != nullptr)
callback(strerror(errno), errno);
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::ERR;
}
delete[] pollfds;
return -1;
}
else if (retval == 0)
{
// we timed out
delete[] pollfds;
return 0;
}

for (uint32_t i = 0; i < nHandles; i++)
{
if (pollfds[i].revents != 0)
{
// error check FIRST
if (pollfds[i].revents & POLLHUP)
{
// check for hangup first because a closed socket
// will technically meet the requirements for POLLIN
// i.e., a call to recv/read won't block
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::HANGUP;
delete[] pollfds;
return -1;
}
else if ((pollfds[i].revents & (POLLERR|POLLNVAL)))
{
if (callback != nullptr)
callback("Poll error", (uint32_t)pollfds[i].revents);
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::ERR;
delete[] pollfds;
return -1;
}
else if (pollfds[i].revents & POLLIN)
{
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::SIGNALED;
break;
}
}
}

delete[] pollfds;
return 1;
return new IpcStream(clientSocket);
}

void IpcStream::DiagnosticsIpc::Close(ErrorCallback callback)
Expand DownExpand Up@@ -307,101 +172,45 @@ void IpcStream::DiagnosticsIpc::Unlink(ErrorCallback callback)
}

IpcStream::~IpcStream()
{
Close();
}

void IpcStream::Close(ErrorCallback)
{
if (_clientSocket != -1)
{
Flush();

const int fSuccessClose = ::close(_clientSocket);
_ASSERTE(fSuccessClose != -1);
_clientSocket = -1;
}
}

bool IpcStream::Read(void *lpBuffer, const uint32_t nBytesToRead, uint32_t &nBytesRead, const int32_t timeoutMs)
bool IpcStream::Read(void *lpBuffer, const uint32_t nBytesToRead, uint32_t &nBytesRead) const
{
_ASSERTE(lpBuffer != nullptr);

if (timeoutMs != InfiniteTimeout)
{
pollfd pfd;
pfd.fd = _clientSocket;
pfd.events = POLLIN;
int retval = poll(&pfd, 1, timeoutMs);
if (retval <= 0 || pfd.revents != POLLIN)
{
// timeout or error
return false;
}
// else fallthrough
}

uint8_t *lpBufferCursor = (uint8_t*)lpBuffer;
ssize_t currentBytesRead = 0;
ssize_t totalBytesRead = 0;
bool fSuccess = true;
while (fSuccess && nBytesToRead - totalBytesRead > 0)
{
currentBytesRead = ::recv(_clientSocket, lpBufferCursor, nBytesToRead - totalBytesRead, 0);
fSuccess = currentBytesRead != 0;
if (!fSuccess)
break;
totalBytesRead += currentBytesRead;
lpBufferCursor += currentBytesRead;
}
const ssize_t ssize = ::recv(_clientSocket, lpBuffer, nBytesToRead, 0);
const bool fSuccess = ssize != -1;

if (!fSuccess)
{
// TODO: Add error handling.
}

nBytesRead = static_cast<uint32_t>(totalBytesRead);
nBytesRead = static_cast<uint32_t>(ssize);
return fSuccess;
}

bool IpcStream::Write(const void *lpBuffer, const uint32_t nBytesToWrite, uint32_t &nBytesWritten, const int32_t timeoutMs)
bool IpcStream::Write(const void *lpBuffer, const uint32_t nBytesToWrite, uint32_t &nBytesWritten) const
{
_ASSERTE(lpBuffer != nullptr);

if (timeoutMs != InfiniteTimeout)
{
pollfd pfd;
pfd.fd = _clientSocket;
pfd.events = POLLOUT;
int retval = poll(&pfd, 1, timeoutMs);
if (retval <= 0 || pfd.revents != POLLOUT)
{
// timeout or error
return false;
}
// else fallthrough
}

uint8_t *lpBufferCursor = (uint8_t*)lpBuffer;
ssize_t currentBytesWritten = 0;
ssize_t totalBytesWritten = 0;
bool fSuccess = true;
while (fSuccess && nBytesToWrite - totalBytesWritten > 0)
{
currentBytesWritten = ::send(_clientSocket, lpBufferCursor, nBytesToWrite - totalBytesWritten, 0);
fSuccess = currentBytesWritten != -1;
if (!fSuccess)
break;
lpBufferCursor += currentBytesWritten;
totalBytesWritten += currentBytesWritten;
}
const ssize_t ssize = ::send(_clientSocket, lpBuffer, nBytesToWrite, 0);
const bool fSuccess = ssize != -1;

if (!fSuccess)
{
// TODO: Add error handling.
}

nBytesWritten = static_cast<uint32_t>(totalBytesWritten);
nBytesWritten = static_cast<uint32_t>(ssize);
return fSuccess;
}

Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
277 changes: 43 additions & 234 deletions src/coreclr/src/debug/debug-pal/unix/diagnosticsipc.cpp
Original file line numberDiff line numberDiff line change
Expand Up@@ -12,20 +12,13 @@
#include "diagnosticsipc.h"
#include "processdescriptor.h"

#if __GNUC__
#include <poll.h>
#else
#include <sys/poll.h>
#endif // __GNUC__

IpcStream::DiagnosticsIpc::DiagnosticsIpc(const int serverSocket, sockaddr_un *const pServerAddress, ConnectionMode mode) :
mode(mode),
IpcStream::DiagnosticsIpc::DiagnosticsIpc(const int serverSocket, sockaddr_un *const pServerAddress) :
_serverSocket(serverSocket),
_pServerAddress(new sockaddr_un),
_isClosed(false),
_isListening(false)
_isClosed(false)
{
_ASSERTE(_pServerAddress != nullptr);
_ASSERTE(_serverSocket != -1);
_ASSERTE(pServerAddress != nullptr);

if (_pServerAddress == nullptr || pServerAddress == nullptr)
Expand All@@ -39,8 +32,24 @@ IpcStream::DiagnosticsIpc::~DiagnosticsIpc()
delete _pServerAddress;
}

IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const pIpcName, ConnectionMode mode, ErrorCallback callback)
IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const pIpcName, ErrorCallback callback)
{
#ifdef __APPLE__
mode_t prev_mask = umask(~(S_IRUSR | S_IWUSR)); // This will set the default permission bit to 600
#endif // __APPLE__

const int serverSocket = ::socket(AF_UNIX, SOCK_STREAM, 0);
if (serverSocket == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__
_ASSERTE(!"Failed to create diagnostics IPC socket.");
return nullptr;
}

sockaddr_un serverAddress{};
serverAddress.sun_family = AF_UNIX;

Expand All@@ -62,24 +71,6 @@ IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const p
"socket");
}

if (mode == ConnectionMode::CLIENT)
return new IpcStream::DiagnosticsIpc(-1, &serverAddress, ConnectionMode::CLIENT);

#ifdef __APPLE__
mode_t prev_mask = umask(~(S_IRUSR | S_IWUSR)); // This will set the default permission bit to 600
#endif // __APPLE__

const int serverSocket = ::socket(AF_UNIX, SOCK_STREAM, 0);
if (serverSocket == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__
_ASSERTE(!"Failed to create diagnostics IPC socket.");
return nullptr;
}

#ifndef __APPLE__
if (fchmod(serverSocket, S_IRUSR | S_IWUSR) == -1)
Expand DownExpand Up@@ -108,52 +99,33 @@ IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const p
return nullptr;
}

#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__

return new IpcStream::DiagnosticsIpc(serverSocket, &serverAddress, mode);
}

bool IpcStream::DiagnosticsIpc::Listen(ErrorCallback callback)
{
_ASSERTE(mode == ConnectionMode::SERVER);
if (mode != ConnectionMode::SERVER)
{
if (callback != nullptr)
callback("Cannot call Listen on a client connection", -1);
return false;
}

if (_isListening)
return true;

const int fSuccessfulListen = ::listen(_serverSocket, /* backlog */ 255);
const int fSuccessfulListen = ::listen(serverSocket, /* backlog */ 255);
if (fSuccessfulListen == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
_ASSERTE(fSuccessfulListen != -1);

const int fSuccessUnlink = ::unlink(_pServerAddress->sun_path);
const int fSuccessUnlink = ::unlink(serverAddress.sun_path);
_ASSERTE(fSuccessUnlink != -1);

const int fSuccessClose = ::close(_serverSocket);
const int fSuccessClose = ::close(serverSocket);
_ASSERTE(fSuccessClose != -1);
return false;
}
else
{
_isListening = true;
return true;
#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__
return nullptr;
}

#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__

return new IpcStream::DiagnosticsIpc(serverSocket, &serverAddress);
}

IpcStream *IpcStream::DiagnosticsIpc::Accept(ErrorCallback callback)
IpcStream *IpcStream::DiagnosticsIpc::Accept(ErrorCallback callback) const
{
_ASSERTE(mode == ConnectionMode::SERVER);
_ASSERTE(_isListening);

sockaddr_un from;
socklen_t fromlen = sizeof(from);
const int clientSocket = ::accept(_serverSocket, (sockaddr *)&from, &fromlen);
Expand All@@ -164,114 +136,7 @@ IpcStream *IpcStream::DiagnosticsIpc::Accept(ErrorCallback callback)
return nullptr;
}

return new IpcStream(clientSocket, mode);
}

IpcStream *IpcStream::DiagnosticsIpc::Connect(ErrorCallback callback)
{
_ASSERTE(mode == ConnectionMode::CLIENT);

sockaddr_un clientAddress{};
clientAddress.sun_family = AF_UNIX;
const int clientSocket = ::socket(AF_UNIX, SOCK_STREAM, 0);
if (clientSocket == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
return nullptr;
}

// We don't expect this to block since this is a Unix Domain Socket. `connect` may block until the
// TCP handshake is complete for TCP/IP sockets, but UDS don't use TCP. `connect` will return even if
// the server hasn't called `accept`.
if (::connect(clientSocket, (struct sockaddr *)_pServerAddress, sizeof(*_pServerAddress)) < 0)
{
if (callback != nullptr)
callback(strerror(errno), errno);
return nullptr;
}

return new IpcStream(clientSocket, ConnectionMode::CLIENT);
}

int32_t IpcStream::DiagnosticsIpc::Poll(IpcPollHandle *rgIpcPollHandles, uint32_t nHandles, int32_t timeoutMs, ErrorCallback callback)
{
// prepare the pollfd structs
pollfd *pollfds = new pollfd[nHandles];
for (uint32_t i = 0; i < nHandles; i++)
{
rgIpcPollHandles[i].revents = 0; // ignore any values in revents
int fd = -1;
if (rgIpcPollHandles[i].pIpc != nullptr)
{
// SERVER
_ASSERTE(rgIpcPollHandles[i].pIpc->mode == ConnectionMode::SERVER);
fd = rgIpcPollHandles[i].pIpc->_serverSocket;
}
else
{
// CLIENT
_ASSERTE(rgIpcPollHandles[i].pStream != nullptr);
fd = rgIpcPollHandles[i].pStream->_clientSocket;
}

pollfds[i].fd = fd;
pollfds[i].events = POLLIN;
}

int retval = poll(pollfds, nHandles, timeoutMs);

// Check results
if (retval < 0)
{
for (uint32_t i = 0; i < nHandles; i++)
{
if ((pollfds[i].revents & POLLERR) && callback != nullptr)
callback(strerror(errno), errno);
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::ERR;
}
delete[] pollfds;
return -1;
}
else if (retval == 0)
{
// we timed out
delete[] pollfds;
return 0;
}

for (uint32_t i = 0; i < nHandles; i++)
{
if (pollfds[i].revents != 0)
{
// error check FIRST
if (pollfds[i].revents & POLLHUP)
{
// check for hangup first because a closed socket
// will technically meet the requirements for POLLIN
// i.e., a call to recv/read won't block
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::HANGUP;
delete[] pollfds;
return -1;
}
else if ((pollfds[i].revents & (POLLERR|POLLNVAL)))
{
if (callback != nullptr)
callback("Poll error", (uint32_t)pollfds[i].revents);
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::ERR;
delete[] pollfds;
return -1;
}
else if (pollfds[i].revents & POLLIN)
{
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::SIGNALED;
break;
}
}
}

delete[] pollfds;
return 1;
return new IpcStream(clientSocket);
}

void IpcStream::DiagnosticsIpc::Close(ErrorCallback callback)
Expand DownExpand Up@@ -307,101 +172,45 @@ void IpcStream::DiagnosticsIpc::Unlink(ErrorCallback callback)
}

IpcStream::~IpcStream()
{
Close();
}

void IpcStream::Close(ErrorCallback)
{
if (_clientSocket != -1)
{
Flush();

const int fSuccessClose = ::close(_clientSocket);
_ASSERTE(fSuccessClose != -1);
_clientSocket = -1;
}
}

bool IpcStream::Read(void *lpBuffer, const uint32_t nBytesToRead, uint32_t &nBytesRead, const int32_t timeoutMs)
bool IpcStream::Read(void *lpBuffer, const uint32_t nBytesToRead, uint32_t &nBytesRead) const
{
_ASSERTE(lpBuffer != nullptr);

if (timeoutMs != InfiniteTimeout)
{
pollfd pfd;
pfd.fd = _clientSocket;
pfd.events = POLLIN;
int retval = poll(&pfd, 1, timeoutMs);
if (retval <= 0 || pfd.revents != POLLIN)
{
// timeout or error
return false;
}
// else fallthrough
}

uint8_t *lpBufferCursor = (uint8_t*)lpBuffer;
ssize_t currentBytesRead = 0;
ssize_t totalBytesRead = 0;
bool fSuccess = true;
while (fSuccess && nBytesToRead - totalBytesRead > 0)
{
currentBytesRead = ::recv(_clientSocket, lpBufferCursor, nBytesToRead - totalBytesRead, 0);
fSuccess = currentBytesRead != 0;
if (!fSuccess)
break;
totalBytesRead += currentBytesRead;
lpBufferCursor += currentBytesRead;
}
const ssize_t ssize = ::recv(_clientSocket, lpBuffer, nBytesToRead, 0);
const bool fSuccess = ssize != -1;

if (!fSuccess)
{
// TODO: Add error handling.
}

nBytesRead = static_cast<uint32_t>(totalBytesRead);
nBytesRead = static_cast<uint32_t>(ssize);
return fSuccess;
}

bool IpcStream::Write(const void *lpBuffer, const uint32_t nBytesToWrite, uint32_t &nBytesWritten, const int32_t timeoutMs)
bool IpcStream::Write(const void *lpBuffer, const uint32_t nBytesToWrite, uint32_t &nBytesWritten) const
{
_ASSERTE(lpBuffer != nullptr);

if (timeoutMs != InfiniteTimeout)
{
pollfd pfd;
pfd.fd = _clientSocket;
pfd.events = POLLOUT;
int retval = poll(&pfd, 1, timeoutMs);
if (retval <= 0 || pfd.revents != POLLOUT)
{
// timeout or error
return false;
}
// else fallthrough
}

uint8_t *lpBufferCursor = (uint8_t*)lpBuffer;
ssize_t currentBytesWritten = 0;
ssize_t totalBytesWritten = 0;
bool fSuccess = true;
while (fSuccess && nBytesToWrite - totalBytesWritten > 0)
{
currentBytesWritten = ::send(_clientSocket, lpBufferCursor, nBytesToWrite - totalBytesWritten, 0);
fSuccess = currentBytesWritten != -1;
if (!fSuccess)
break;
lpBufferCursor += currentBytesWritten;
totalBytesWritten += currentBytesWritten;
}
const ssize_t ssize = ::send(_clientSocket, lpBuffer, nBytesToWrite, 0);
const bool fSuccess = ssize != -1;

if (!fSuccess)
{
// TODO: Add error handling.
}

nBytesWritten = static_cast<uint32_t>(totalBytesWritten);
nBytesWritten = static_cast<uint32_t>(ssize);
return fSuccess;
}

Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
277 changes: 43 additions & 234 deletions src/coreclr/src/debug/debug-pal/unix/diagnosticsipc.cpp
Original file line numberDiff line numberDiff line change
Expand Up@@ -12,20 +12,13 @@
#include "diagnosticsipc.h"
#include "processdescriptor.h"

#if __GNUC__
#include <poll.h>
#else
#include <sys/poll.h>
#endif // __GNUC__

IpcStream::DiagnosticsIpc::DiagnosticsIpc(const int serverSocket, sockaddr_un *const pServerAddress, ConnectionMode mode) :
mode(mode),
IpcStream::DiagnosticsIpc::DiagnosticsIpc(const int serverSocket, sockaddr_un *const pServerAddress) :
_serverSocket(serverSocket),
_pServerAddress(new sockaddr_un),
_isClosed(false),
_isListening(false)
_isClosed(false)
{
_ASSERTE(_pServerAddress != nullptr);
_ASSERTE(_serverSocket != -1);
_ASSERTE(pServerAddress != nullptr);

if (_pServerAddress == nullptr || pServerAddress == nullptr)
Expand All@@ -39,8 +32,24 @@ IpcStream::DiagnosticsIpc::~DiagnosticsIpc()
delete _pServerAddress;
}

IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const pIpcName, ConnectionMode mode, ErrorCallback callback)
IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const pIpcName, ErrorCallback callback)
{
#ifdef __APPLE__
mode_t prev_mask = umask(~(S_IRUSR | S_IWUSR)); // This will set the default permission bit to 600
#endif // __APPLE__

const int serverSocket = ::socket(AF_UNIX, SOCK_STREAM, 0);
if (serverSocket == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__
_ASSERTE(!"Failed to create diagnostics IPC socket.");
return nullptr;
}

sockaddr_un serverAddress{};
serverAddress.sun_family = AF_UNIX;

Expand All@@ -62,24 +71,6 @@ IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const p
"socket");
}

if (mode == ConnectionMode::CLIENT)
return new IpcStream::DiagnosticsIpc(-1, &serverAddress, ConnectionMode::CLIENT);

#ifdef __APPLE__
mode_t prev_mask = umask(~(S_IRUSR | S_IWUSR)); // This will set the default permission bit to 600
#endif // __APPLE__

const int serverSocket = ::socket(AF_UNIX, SOCK_STREAM, 0);
if (serverSocket == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__
_ASSERTE(!"Failed to create diagnostics IPC socket.");
return nullptr;
}

#ifndef __APPLE__
if (fchmod(serverSocket, S_IRUSR | S_IWUSR) == -1)
Expand DownExpand Up@@ -108,52 +99,33 @@ IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const p
return nullptr;
}

#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__

return new IpcStream::DiagnosticsIpc(serverSocket, &serverAddress, mode);
}

bool IpcStream::DiagnosticsIpc::Listen(ErrorCallback callback)
{
_ASSERTE(mode == ConnectionMode::SERVER);
if (mode != ConnectionMode::SERVER)
{
if (callback != nullptr)
callback("Cannot call Listen on a client connection", -1);
return false;
}

if (_isListening)
return true;

const int fSuccessfulListen = ::listen(_serverSocket, /* backlog */ 255);
const int fSuccessfulListen = ::listen(serverSocket, /* backlog */ 255);
if (fSuccessfulListen == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
_ASSERTE(fSuccessfulListen != -1);

const int fSuccessUnlink = ::unlink(_pServerAddress->sun_path);
const int fSuccessUnlink = ::unlink(serverAddress.sun_path);
_ASSERTE(fSuccessUnlink != -1);

const int fSuccessClose = ::close(_serverSocket);
const int fSuccessClose = ::close(serverSocket);
_ASSERTE(fSuccessClose != -1);
return false;
}
else
{
_isListening = true;
return true;
#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__
return nullptr;
}

#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__

return new IpcStream::DiagnosticsIpc(serverSocket, &serverAddress);
}

IpcStream *IpcStream::DiagnosticsIpc::Accept(ErrorCallback callback)
IpcStream *IpcStream::DiagnosticsIpc::Accept(ErrorCallback callback) const
{
_ASSERTE(mode == ConnectionMode::SERVER);
_ASSERTE(_isListening);

sockaddr_un from;
socklen_t fromlen = sizeof(from);
const int clientSocket = ::accept(_serverSocket, (sockaddr *)&from, &fromlen);
Expand All@@ -164,114 +136,7 @@ IpcStream *IpcStream::DiagnosticsIpc::Accept(ErrorCallback callback)
return nullptr;
}

return new IpcStream(clientSocket, mode);
}

IpcStream *IpcStream::DiagnosticsIpc::Connect(ErrorCallback callback)
{
_ASSERTE(mode == ConnectionMode::CLIENT);

sockaddr_un clientAddress{};
clientAddress.sun_family = AF_UNIX;
const int clientSocket = ::socket(AF_UNIX, SOCK_STREAM, 0);
if (clientSocket == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
return nullptr;
}

// We don't expect this to block since this is a Unix Domain Socket. `connect` may block until the
// TCP handshake is complete for TCP/IP sockets, but UDS don't use TCP. `connect` will return even if
// the server hasn't called `accept`.
if (::connect(clientSocket, (struct sockaddr *)_pServerAddress, sizeof(*_pServerAddress)) < 0)
{
if (callback != nullptr)
callback(strerror(errno), errno);
return nullptr;
}

return new IpcStream(clientSocket, ConnectionMode::CLIENT);
}

int32_t IpcStream::DiagnosticsIpc::Poll(IpcPollHandle *rgIpcPollHandles, uint32_t nHandles, int32_t timeoutMs, ErrorCallback callback)
{
// prepare the pollfd structs
pollfd *pollfds = new pollfd[nHandles];
for (uint32_t i = 0; i < nHandles; i++)
{
rgIpcPollHandles[i].revents = 0; // ignore any values in revents
int fd = -1;
if (rgIpcPollHandles[i].pIpc != nullptr)
{
// SERVER
_ASSERTE(rgIpcPollHandles[i].pIpc->mode == ConnectionMode::SERVER);
fd = rgIpcPollHandles[i].pIpc->_serverSocket;
}
else
{
// CLIENT
_ASSERTE(rgIpcPollHandles[i].pStream != nullptr);
fd = rgIpcPollHandles[i].pStream->_clientSocket;
}

pollfds[i].fd = fd;
pollfds[i].events = POLLIN;
}

int retval = poll(pollfds, nHandles, timeoutMs);

// Check results
if (retval < 0)
{
for (uint32_t i = 0; i < nHandles; i++)
{
if ((pollfds[i].revents & POLLERR) && callback != nullptr)
callback(strerror(errno), errno);
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::ERR;
}
delete[] pollfds;
return -1;
}
else if (retval == 0)
{
// we timed out
delete[] pollfds;
return 0;
}

for (uint32_t i = 0; i < nHandles; i++)
{
if (pollfds[i].revents != 0)
{
// error check FIRST
if (pollfds[i].revents & POLLHUP)
{
// check for hangup first because a closed socket
// will technically meet the requirements for POLLIN
// i.e., a call to recv/read won't block
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::HANGUP;
delete[] pollfds;
return -1;
}
else if ((pollfds[i].revents & (POLLERR|POLLNVAL)))
{
if (callback != nullptr)
callback("Poll error", (uint32_t)pollfds[i].revents);
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::ERR;
delete[] pollfds;
return -1;
}
else if (pollfds[i].revents & POLLIN)
{
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::SIGNALED;
break;
}
}
}

delete[] pollfds;
return 1;
return new IpcStream(clientSocket);
}

void IpcStream::DiagnosticsIpc::Close(ErrorCallback callback)
Expand DownExpand Up@@ -307,101 +172,45 @@ void IpcStream::DiagnosticsIpc::Unlink(ErrorCallback callback)
}

IpcStream::~IpcStream()
{
Close();
}

void IpcStream::Close(ErrorCallback)
{
if (_clientSocket != -1)
{
Flush();

const int fSuccessClose = ::close(_clientSocket);
_ASSERTE(fSuccessClose != -1);
_clientSocket = -1;
}
}

bool IpcStream::Read(void *lpBuffer, const uint32_t nBytesToRead, uint32_t &nBytesRead, const int32_t timeoutMs)
bool IpcStream::Read(void *lpBuffer, const uint32_t nBytesToRead, uint32_t &nBytesRead) const
{
_ASSERTE(lpBuffer != nullptr);

if (timeoutMs != InfiniteTimeout)
{
pollfd pfd;
pfd.fd = _clientSocket;
pfd.events = POLLIN;
int retval = poll(&pfd, 1, timeoutMs);
if (retval <= 0 || pfd.revents != POLLIN)
{
// timeout or error
return false;
}
// else fallthrough
}

uint8_t *lpBufferCursor = (uint8_t*)lpBuffer;
ssize_t currentBytesRead = 0;
ssize_t totalBytesRead = 0;
bool fSuccess = true;
while (fSuccess && nBytesToRead - totalBytesRead > 0)
{
currentBytesRead = ::recv(_clientSocket, lpBufferCursor, nBytesToRead - totalBytesRead, 0);
fSuccess = currentBytesRead != 0;
if (!fSuccess)
break;
totalBytesRead += currentBytesRead;
lpBufferCursor += currentBytesRead;
}
const ssize_t ssize = ::recv(_clientSocket, lpBuffer, nBytesToRead, 0);
const bool fSuccess = ssize != -1;

if (!fSuccess)
{
// TODO: Add error handling.
}

nBytesRead = static_cast<uint32_t>(totalBytesRead);
nBytesRead = static_cast<uint32_t>(ssize);
return fSuccess;
}

bool IpcStream::Write(const void *lpBuffer, const uint32_t nBytesToWrite, uint32_t &nBytesWritten, const int32_t timeoutMs)
bool IpcStream::Write(const void *lpBuffer, const uint32_t nBytesToWrite, uint32_t &nBytesWritten) const
{
_ASSERTE(lpBuffer != nullptr);

if (timeoutMs != InfiniteTimeout)
{
pollfd pfd;
pfd.fd = _clientSocket;
pfd.events = POLLOUT;
int retval = poll(&pfd, 1, timeoutMs);
if (retval <= 0 || pfd.revents != POLLOUT)
{
// timeout or error
return false;
}
// else fallthrough
}

uint8_t *lpBufferCursor = (uint8_t*)lpBuffer;
ssize_t currentBytesWritten = 0;
ssize_t totalBytesWritten = 0;
bool fSuccess = true;
while (fSuccess && nBytesToWrite - totalBytesWritten > 0)
{
currentBytesWritten = ::send(_clientSocket, lpBufferCursor, nBytesToWrite - totalBytesWritten, 0);
fSuccess = currentBytesWritten != -1;
if (!fSuccess)
break;
lpBufferCursor += currentBytesWritten;
totalBytesWritten += currentBytesWritten;
}
const ssize_t ssize = ::send(_clientSocket, lpBuffer, nBytesToWrite, 0);
const bool fSuccess = ssize != -1;

if (!fSuccess)
{
// TODO: Add error handling.
}

nBytesWritten = static_cast<uint32_t>(totalBytesWritten);
nBytesWritten = static_cast<uint32_t>(ssize);
return fSuccess;
}

Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
277 changes: 43 additions & 234 deletions src/coreclr/src/debug/debug-pal/unix/diagnosticsipc.cpp
Original file line numberDiff line numberDiff line change
Expand Up@@ -12,20 +12,13 @@
#include "diagnosticsipc.h"
#include "processdescriptor.h"

#if __GNUC__
#include <poll.h>
#else
#include <sys/poll.h>
#endif // __GNUC__

IpcStream::DiagnosticsIpc::DiagnosticsIpc(const int serverSocket, sockaddr_un *const pServerAddress, ConnectionMode mode) :
mode(mode),
IpcStream::DiagnosticsIpc::DiagnosticsIpc(const int serverSocket, sockaddr_un *const pServerAddress) :
_serverSocket(serverSocket),
_pServerAddress(new sockaddr_un),
_isClosed(false),
_isListening(false)
_isClosed(false)
{
_ASSERTE(_pServerAddress != nullptr);
_ASSERTE(_serverSocket != -1);
_ASSERTE(pServerAddress != nullptr);

if (_pServerAddress == nullptr || pServerAddress == nullptr)
Expand All@@ -39,8 +32,24 @@ IpcStream::DiagnosticsIpc::~DiagnosticsIpc()
delete _pServerAddress;
}

IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const pIpcName, ConnectionMode mode, ErrorCallback callback)
IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const pIpcName, ErrorCallback callback)
{
#ifdef __APPLE__
mode_t prev_mask = umask(~(S_IRUSR | S_IWUSR)); // This will set the default permission bit to 600
#endif // __APPLE__

const int serverSocket = ::socket(AF_UNIX, SOCK_STREAM, 0);
if (serverSocket == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__
_ASSERTE(!"Failed to create diagnostics IPC socket.");
return nullptr;
}

sockaddr_un serverAddress{};
serverAddress.sun_family = AF_UNIX;

Expand All@@ -62,24 +71,6 @@ IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const p
"socket");
}

if (mode == ConnectionMode::CLIENT)
return new IpcStream::DiagnosticsIpc(-1, &serverAddress, ConnectionMode::CLIENT);

#ifdef __APPLE__
mode_t prev_mask = umask(~(S_IRUSR | S_IWUSR)); // This will set the default permission bit to 600
#endif // __APPLE__

const int serverSocket = ::socket(AF_UNIX, SOCK_STREAM, 0);
if (serverSocket == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__
_ASSERTE(!"Failed to create diagnostics IPC socket.");
return nullptr;
}

#ifndef __APPLE__
if (fchmod(serverSocket, S_IRUSR | S_IWUSR) == -1)
Expand DownExpand Up@@ -108,52 +99,33 @@ IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const p
return nullptr;
}

#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__

return new IpcStream::DiagnosticsIpc(serverSocket, &serverAddress, mode);
}

bool IpcStream::DiagnosticsIpc::Listen(ErrorCallback callback)
{
_ASSERTE(mode == ConnectionMode::SERVER);
if (mode != ConnectionMode::SERVER)
{
if (callback != nullptr)
callback("Cannot call Listen on a client connection", -1);
return false;
}

if (_isListening)
return true;

const int fSuccessfulListen = ::listen(_serverSocket, /* backlog */ 255);
const int fSuccessfulListen = ::listen(serverSocket, /* backlog */ 255);
if (fSuccessfulListen == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
_ASSERTE(fSuccessfulListen != -1);

const int fSuccessUnlink = ::unlink(_pServerAddress->sun_path);
const int fSuccessUnlink = ::unlink(serverAddress.sun_path);
_ASSERTE(fSuccessUnlink != -1);

const int fSuccessClose = ::close(_serverSocket);
const int fSuccessClose = ::close(serverSocket);
_ASSERTE(fSuccessClose != -1);
return false;
}
else
{
_isListening = true;
return true;
#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__
return nullptr;
}

#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__

return new IpcStream::DiagnosticsIpc(serverSocket, &serverAddress);
}

IpcStream *IpcStream::DiagnosticsIpc::Accept(ErrorCallback callback)
IpcStream *IpcStream::DiagnosticsIpc::Accept(ErrorCallback callback) const
{
_ASSERTE(mode == ConnectionMode::SERVER);
_ASSERTE(_isListening);

sockaddr_un from;
socklen_t fromlen = sizeof(from);
const int clientSocket = ::accept(_serverSocket, (sockaddr *)&from, &fromlen);
Expand All@@ -164,114 +136,7 @@ IpcStream *IpcStream::DiagnosticsIpc::Accept(ErrorCallback callback)
return nullptr;
}

return new IpcStream(clientSocket, mode);
}

IpcStream *IpcStream::DiagnosticsIpc::Connect(ErrorCallback callback)
{
_ASSERTE(mode == ConnectionMode::CLIENT);

sockaddr_un clientAddress{};
clientAddress.sun_family = AF_UNIX;
const int clientSocket = ::socket(AF_UNIX, SOCK_STREAM, 0);
if (clientSocket == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
return nullptr;
}

// We don't expect this to block since this is a Unix Domain Socket. `connect` may block until the
// TCP handshake is complete for TCP/IP sockets, but UDS don't use TCP. `connect` will return even if
// the server hasn't called `accept`.
if (::connect(clientSocket, (struct sockaddr *)_pServerAddress, sizeof(*_pServerAddress)) < 0)
{
if (callback != nullptr)
callback(strerror(errno), errno);
return nullptr;
}

return new IpcStream(clientSocket, ConnectionMode::CLIENT);
}

int32_t IpcStream::DiagnosticsIpc::Poll(IpcPollHandle *rgIpcPollHandles, uint32_t nHandles, int32_t timeoutMs, ErrorCallback callback)
{
// prepare the pollfd structs
pollfd *pollfds = new pollfd[nHandles];
for (uint32_t i = 0; i < nHandles; i++)
{
rgIpcPollHandles[i].revents = 0; // ignore any values in revents
int fd = -1;
if (rgIpcPollHandles[i].pIpc != nullptr)
{
// SERVER
_ASSERTE(rgIpcPollHandles[i].pIpc->mode == ConnectionMode::SERVER);
fd = rgIpcPollHandles[i].pIpc->_serverSocket;
}
else
{
// CLIENT
_ASSERTE(rgIpcPollHandles[i].pStream != nullptr);
fd = rgIpcPollHandles[i].pStream->_clientSocket;
}

pollfds[i].fd = fd;
pollfds[i].events = POLLIN;
}

int retval = poll(pollfds, nHandles, timeoutMs);

// Check results
if (retval < 0)
{
for (uint32_t i = 0; i < nHandles; i++)
{
if ((pollfds[i].revents & POLLERR) && callback != nullptr)
callback(strerror(errno), errno);
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::ERR;
}
delete[] pollfds;
return -1;
}
else if (retval == 0)
{
// we timed out
delete[] pollfds;
return 0;
}

for (uint32_t i = 0; i < nHandles; i++)
{
if (pollfds[i].revents != 0)
{
// error check FIRST
if (pollfds[i].revents & POLLHUP)
{
// check for hangup first because a closed socket
// will technically meet the requirements for POLLIN
// i.e., a call to recv/read won't block
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::HANGUP;
delete[] pollfds;
return -1;
}
else if ((pollfds[i].revents & (POLLERR|POLLNVAL)))
{
if (callback != nullptr)
callback("Poll error", (uint32_t)pollfds[i].revents);
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::ERR;
delete[] pollfds;
return -1;
}
else if (pollfds[i].revents & POLLIN)
{
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::SIGNALED;
break;
}
}
}

delete[] pollfds;
return 1;
return new IpcStream(clientSocket);
}

void IpcStream::DiagnosticsIpc::Close(ErrorCallback callback)
Expand DownExpand Up@@ -307,101 +172,45 @@ void IpcStream::DiagnosticsIpc::Unlink(ErrorCallback callback)
}

IpcStream::~IpcStream()
{
Close();
}

void IpcStream::Close(ErrorCallback)
{
if (_clientSocket != -1)
{
Flush();

const int fSuccessClose = ::close(_clientSocket);
_ASSERTE(fSuccessClose != -1);
_clientSocket = -1;
}
}

bool IpcStream::Read(void *lpBuffer, const uint32_t nBytesToRead, uint32_t &nBytesRead, const int32_t timeoutMs)
bool IpcStream::Read(void *lpBuffer, const uint32_t nBytesToRead, uint32_t &nBytesRead) const
{
_ASSERTE(lpBuffer != nullptr);

if (timeoutMs != InfiniteTimeout)
{
pollfd pfd;
pfd.fd = _clientSocket;
pfd.events = POLLIN;
int retval = poll(&pfd, 1, timeoutMs);
if (retval <= 0 || pfd.revents != POLLIN)
{
// timeout or error
return false;
}
// else fallthrough
}

uint8_t *lpBufferCursor = (uint8_t*)lpBuffer;
ssize_t currentBytesRead = 0;
ssize_t totalBytesRead = 0;
bool fSuccess = true;
while (fSuccess && nBytesToRead - totalBytesRead > 0)
{
currentBytesRead = ::recv(_clientSocket, lpBufferCursor, nBytesToRead - totalBytesRead, 0);
fSuccess = currentBytesRead != 0;
if (!fSuccess)
break;
totalBytesRead += currentBytesRead;
lpBufferCursor += currentBytesRead;
}
const ssize_t ssize = ::recv(_clientSocket, lpBuffer, nBytesToRead, 0);
const bool fSuccess = ssize != -1;

if (!fSuccess)
{
// TODO: Add error handling.
}

nBytesRead = static_cast<uint32_t>(totalBytesRead);
nBytesRead = static_cast<uint32_t>(ssize);
return fSuccess;
}

bool IpcStream::Write(const void *lpBuffer, const uint32_t nBytesToWrite, uint32_t &nBytesWritten, const int32_t timeoutMs)
bool IpcStream::Write(const void *lpBuffer, const uint32_t nBytesToWrite, uint32_t &nBytesWritten) const
{
_ASSERTE(lpBuffer != nullptr);

if (timeoutMs != InfiniteTimeout)
{
pollfd pfd;
pfd.fd = _clientSocket;
pfd.events = POLLOUT;
int retval = poll(&pfd, 1, timeoutMs);
if (retval <= 0 || pfd.revents != POLLOUT)
{
// timeout or error
return false;
}
// else fallthrough
}

uint8_t *lpBufferCursor = (uint8_t*)lpBuffer;
ssize_t currentBytesWritten = 0;
ssize_t totalBytesWritten = 0;
bool fSuccess = true;
while (fSuccess && nBytesToWrite - totalBytesWritten > 0)
{
currentBytesWritten = ::send(_clientSocket, lpBufferCursor, nBytesToWrite - totalBytesWritten, 0);
fSuccess = currentBytesWritten != -1;
if (!fSuccess)
break;
lpBufferCursor += currentBytesWritten;
totalBytesWritten += currentBytesWritten;
}
const ssize_t ssize = ::send(_clientSocket, lpBuffer, nBytesToWrite, 0);
const bool fSuccess = ssize != -1;

if (!fSuccess)
{
// TODO: Add error handling.
}

nBytesWritten = static_cast<uint32_t>(totalBytesWritten);
nBytesWritten = static_cast<uint32_t>(ssize);
return fSuccess;
}

Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
277 changes: 43 additions & 234 deletions src/coreclr/src/debug/debug-pal/unix/diagnosticsipc.cpp
Original file line numberDiff line numberDiff line change
Expand Up@@ -12,20 +12,13 @@
#include "diagnosticsipc.h"
#include "processdescriptor.h"

#if __GNUC__
#include <poll.h>
#else
#include <sys/poll.h>
#endif // __GNUC__

IpcStream::DiagnosticsIpc::DiagnosticsIpc(const int serverSocket, sockaddr_un *const pServerAddress, ConnectionMode mode) :
mode(mode),
IpcStream::DiagnosticsIpc::DiagnosticsIpc(const int serverSocket, sockaddr_un *const pServerAddress) :
_serverSocket(serverSocket),
_pServerAddress(new sockaddr_un),
_isClosed(false),
_isListening(false)
_isClosed(false)
{
_ASSERTE(_pServerAddress != nullptr);
_ASSERTE(_serverSocket != -1);
_ASSERTE(pServerAddress != nullptr);

if (_pServerAddress == nullptr || pServerAddress == nullptr)
Expand All@@ -39,8 +32,24 @@ IpcStream::DiagnosticsIpc::~DiagnosticsIpc()
delete _pServerAddress;
}

IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const pIpcName, ConnectionMode mode, ErrorCallback callback)
IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const pIpcName, ErrorCallback callback)
{
#ifdef __APPLE__
mode_t prev_mask = umask(~(S_IRUSR | S_IWUSR)); // This will set the default permission bit to 600
#endif // __APPLE__

const int serverSocket = ::socket(AF_UNIX, SOCK_STREAM, 0);
if (serverSocket == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__
_ASSERTE(!"Failed to create diagnostics IPC socket.");
return nullptr;
}

sockaddr_un serverAddress{};
serverAddress.sun_family = AF_UNIX;

Expand All@@ -62,24 +71,6 @@ IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const p
"socket");
}

if (mode == ConnectionMode::CLIENT)
return new IpcStream::DiagnosticsIpc(-1, &serverAddress, ConnectionMode::CLIENT);

#ifdef __APPLE__
mode_t prev_mask = umask(~(S_IRUSR | S_IWUSR)); // This will set the default permission bit to 600
#endif // __APPLE__

const int serverSocket = ::socket(AF_UNIX, SOCK_STREAM, 0);
if (serverSocket == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__
_ASSERTE(!"Failed to create diagnostics IPC socket.");
return nullptr;
}

#ifndef __APPLE__
if (fchmod(serverSocket, S_IRUSR | S_IWUSR) == -1)
Expand DownExpand Up@@ -108,52 +99,33 @@ IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const p
return nullptr;
}

#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__

return new IpcStream::DiagnosticsIpc(serverSocket, &serverAddress, mode);
}

bool IpcStream::DiagnosticsIpc::Listen(ErrorCallback callback)
{
_ASSERTE(mode == ConnectionMode::SERVER);
if (mode != ConnectionMode::SERVER)
{
if (callback != nullptr)
callback("Cannot call Listen on a client connection", -1);
return false;
}

if (_isListening)
return true;

const int fSuccessfulListen = ::listen(_serverSocket, /* backlog */ 255);
const int fSuccessfulListen = ::listen(serverSocket, /* backlog */ 255);
if (fSuccessfulListen == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
_ASSERTE(fSuccessfulListen != -1);

const int fSuccessUnlink = ::unlink(_pServerAddress->sun_path);
const int fSuccessUnlink = ::unlink(serverAddress.sun_path);
_ASSERTE(fSuccessUnlink != -1);

const int fSuccessClose = ::close(_serverSocket);
const int fSuccessClose = ::close(serverSocket);
_ASSERTE(fSuccessClose != -1);
return false;
}
else
{
_isListening = true;
return true;
#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__
return nullptr;
}

#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__

return new IpcStream::DiagnosticsIpc(serverSocket, &serverAddress);
}

IpcStream *IpcStream::DiagnosticsIpc::Accept(ErrorCallback callback)
IpcStream *IpcStream::DiagnosticsIpc::Accept(ErrorCallback callback) const
{
_ASSERTE(mode == ConnectionMode::SERVER);
_ASSERTE(_isListening);

sockaddr_un from;
socklen_t fromlen = sizeof(from);
const int clientSocket = ::accept(_serverSocket, (sockaddr *)&from, &fromlen);
Expand All@@ -164,114 +136,7 @@ IpcStream *IpcStream::DiagnosticsIpc::Accept(ErrorCallback callback)
return nullptr;
}

return new IpcStream(clientSocket, mode);
}

IpcStream *IpcStream::DiagnosticsIpc::Connect(ErrorCallback callback)
{
_ASSERTE(mode == ConnectionMode::CLIENT);

sockaddr_un clientAddress{};
clientAddress.sun_family = AF_UNIX;
const int clientSocket = ::socket(AF_UNIX, SOCK_STREAM, 0);
if (clientSocket == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
return nullptr;
}

// We don't expect this to block since this is a Unix Domain Socket. `connect` may block until the
// TCP handshake is complete for TCP/IP sockets, but UDS don't use TCP. `connect` will return even if
// the server hasn't called `accept`.
if (::connect(clientSocket, (struct sockaddr *)_pServerAddress, sizeof(*_pServerAddress)) < 0)
{
if (callback != nullptr)
callback(strerror(errno), errno);
return nullptr;
}

return new IpcStream(clientSocket, ConnectionMode::CLIENT);
}

int32_t IpcStream::DiagnosticsIpc::Poll(IpcPollHandle *rgIpcPollHandles, uint32_t nHandles, int32_t timeoutMs, ErrorCallback callback)
{
// prepare the pollfd structs
pollfd *pollfds = new pollfd[nHandles];
for (uint32_t i = 0; i < nHandles; i++)
{
rgIpcPollHandles[i].revents = 0; // ignore any values in revents
int fd = -1;
if (rgIpcPollHandles[i].pIpc != nullptr)
{
// SERVER
_ASSERTE(rgIpcPollHandles[i].pIpc->mode == ConnectionMode::SERVER);
fd = rgIpcPollHandles[i].pIpc->_serverSocket;
}
else
{
// CLIENT
_ASSERTE(rgIpcPollHandles[i].pStream != nullptr);
fd = rgIpcPollHandles[i].pStream->_clientSocket;
}

pollfds[i].fd = fd;
pollfds[i].events = POLLIN;
}

int retval = poll(pollfds, nHandles, timeoutMs);

// Check results
if (retval < 0)
{
for (uint32_t i = 0; i < nHandles; i++)
{
if ((pollfds[i].revents & POLLERR) && callback != nullptr)
callback(strerror(errno), errno);
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::ERR;
}
delete[] pollfds;
return -1;
}
else if (retval == 0)
{
// we timed out
delete[] pollfds;
return 0;
}

for (uint32_t i = 0; i < nHandles; i++)
{
if (pollfds[i].revents != 0)
{
// error check FIRST
if (pollfds[i].revents & POLLHUP)
{
// check for hangup first because a closed socket
// will technically meet the requirements for POLLIN
// i.e., a call to recv/read won't block
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::HANGUP;
delete[] pollfds;
return -1;
}
else if ((pollfds[i].revents & (POLLERR|POLLNVAL)))
{
if (callback != nullptr)
callback("Poll error", (uint32_t)pollfds[i].revents);
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::ERR;
delete[] pollfds;
return -1;
}
else if (pollfds[i].revents & POLLIN)
{
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::SIGNALED;
break;
}
}
}

delete[] pollfds;
return 1;
return new IpcStream(clientSocket);
}

void IpcStream::DiagnosticsIpc::Close(ErrorCallback callback)
Expand DownExpand Up@@ -307,101 +172,45 @@ void IpcStream::DiagnosticsIpc::Unlink(ErrorCallback callback)
}

IpcStream::~IpcStream()
{
Close();
}

void IpcStream::Close(ErrorCallback)
{
if (_clientSocket != -1)
{
Flush();

const int fSuccessClose = ::close(_clientSocket);
_ASSERTE(fSuccessClose != -1);
_clientSocket = -1;
}
}

bool IpcStream::Read(void *lpBuffer, const uint32_t nBytesToRead, uint32_t &nBytesRead, const int32_t timeoutMs)
bool IpcStream::Read(void *lpBuffer, const uint32_t nBytesToRead, uint32_t &nBytesRead) const
{
_ASSERTE(lpBuffer != nullptr);

if (timeoutMs != InfiniteTimeout)
{
pollfd pfd;
pfd.fd = _clientSocket;
pfd.events = POLLIN;
int retval = poll(&pfd, 1, timeoutMs);
if (retval <= 0 || pfd.revents != POLLIN)
{
// timeout or error
return false;
}
// else fallthrough
}

uint8_t *lpBufferCursor = (uint8_t*)lpBuffer;
ssize_t currentBytesRead = 0;
ssize_t totalBytesRead = 0;
bool fSuccess = true;
while (fSuccess && nBytesToRead - totalBytesRead > 0)
{
currentBytesRead = ::recv(_clientSocket, lpBufferCursor, nBytesToRead - totalBytesRead, 0);
fSuccess = currentBytesRead != 0;
if (!fSuccess)
break;
totalBytesRead += currentBytesRead;
lpBufferCursor += currentBytesRead;
}
const ssize_t ssize = ::recv(_clientSocket, lpBuffer, nBytesToRead, 0);
const bool fSuccess = ssize != -1;

if (!fSuccess)
{
// TODO: Add error handling.
}

nBytesRead = static_cast<uint32_t>(totalBytesRead);
nBytesRead = static_cast<uint32_t>(ssize);
return fSuccess;
}

bool IpcStream::Write(const void *lpBuffer, const uint32_t nBytesToWrite, uint32_t &nBytesWritten, const int32_t timeoutMs)
bool IpcStream::Write(const void *lpBuffer, const uint32_t nBytesToWrite, uint32_t &nBytesWritten) const
{
_ASSERTE(lpBuffer != nullptr);

if (timeoutMs != InfiniteTimeout)
{
pollfd pfd;
pfd.fd = _clientSocket;
pfd.events = POLLOUT;
int retval = poll(&pfd, 1, timeoutMs);
if (retval <= 0 || pfd.revents != POLLOUT)
{
// timeout or error
return false;
}
// else fallthrough
}

uint8_t *lpBufferCursor = (uint8_t*)lpBuffer;
ssize_t currentBytesWritten = 0;
ssize_t totalBytesWritten = 0;
bool fSuccess = true;
while (fSuccess && nBytesToWrite - totalBytesWritten > 0)
{
currentBytesWritten = ::send(_clientSocket, lpBufferCursor, nBytesToWrite - totalBytesWritten, 0);
fSuccess = currentBytesWritten != -1;
if (!fSuccess)
break;
lpBufferCursor += currentBytesWritten;
totalBytesWritten += currentBytesWritten;
}
const ssize_t ssize = ::send(_clientSocket, lpBuffer, nBytesToWrite, 0);
const bool fSuccess = ssize != -1;

if (!fSuccess)
{
// TODO: Add error handling.
}

nBytesWritten = static_cast<uint32_t>(totalBytesWritten);
nBytesWritten = static_cast<uint32_t>(ssize);
return fSuccess;
}

Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
277 changes: 43 additions & 234 deletions src/coreclr/src/debug/debug-pal/unix/diagnosticsipc.cpp
Original file line numberDiff line numberDiff line change
Expand Up@@ -12,20 +12,13 @@
#include "diagnosticsipc.h"
#include "processdescriptor.h"

#if __GNUC__
#include <poll.h>
#else
#include <sys/poll.h>
#endif // __GNUC__

IpcStream::DiagnosticsIpc::DiagnosticsIpc(const int serverSocket, sockaddr_un *const pServerAddress, ConnectionMode mode) :
mode(mode),
IpcStream::DiagnosticsIpc::DiagnosticsIpc(const int serverSocket, sockaddr_un *const pServerAddress) :
_serverSocket(serverSocket),
_pServerAddress(new sockaddr_un),
_isClosed(false),
_isListening(false)
_isClosed(false)
{
_ASSERTE(_pServerAddress != nullptr);
_ASSERTE(_serverSocket != -1);
_ASSERTE(pServerAddress != nullptr);

if (_pServerAddress == nullptr || pServerAddress == nullptr)
Expand All@@ -39,8 +32,24 @@ IpcStream::DiagnosticsIpc::~DiagnosticsIpc()
delete _pServerAddress;
}

IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const pIpcName, ConnectionMode mode, ErrorCallback callback)
IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const pIpcName, ErrorCallback callback)
{
#ifdef __APPLE__
mode_t prev_mask = umask(~(S_IRUSR | S_IWUSR)); // This will set the default permission bit to 600
#endif // __APPLE__

const int serverSocket = ::socket(AF_UNIX, SOCK_STREAM, 0);
if (serverSocket == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__
_ASSERTE(!"Failed to create diagnostics IPC socket.");
return nullptr;
}

sockaddr_un serverAddress{};
serverAddress.sun_family = AF_UNIX;

Expand All@@ -62,24 +71,6 @@ IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const p
"socket");
}

if (mode == ConnectionMode::CLIENT)
return new IpcStream::DiagnosticsIpc(-1, &serverAddress, ConnectionMode::CLIENT);

#ifdef __APPLE__
mode_t prev_mask = umask(~(S_IRUSR | S_IWUSR)); // This will set the default permission bit to 600
#endif // __APPLE__

const int serverSocket = ::socket(AF_UNIX, SOCK_STREAM, 0);
if (serverSocket == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__
_ASSERTE(!"Failed to create diagnostics IPC socket.");
return nullptr;
}

#ifndef __APPLE__
if (fchmod(serverSocket, S_IRUSR | S_IWUSR) == -1)
Expand DownExpand Up@@ -108,52 +99,33 @@ IpcStream::DiagnosticsIpc *IpcStream::DiagnosticsIpc::Create(const char *const p
return nullptr;
}

#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__

return new IpcStream::DiagnosticsIpc(serverSocket, &serverAddress, mode);
}

bool IpcStream::DiagnosticsIpc::Listen(ErrorCallback callback)
{
_ASSERTE(mode == ConnectionMode::SERVER);
if (mode != ConnectionMode::SERVER)
{
if (callback != nullptr)
callback("Cannot call Listen on a client connection", -1);
return false;
}

if (_isListening)
return true;

const int fSuccessfulListen = ::listen(_serverSocket, /* backlog */ 255);
const int fSuccessfulListen = ::listen(serverSocket, /* backlog */ 255);
if (fSuccessfulListen == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
_ASSERTE(fSuccessfulListen != -1);

const int fSuccessUnlink = ::unlink(_pServerAddress->sun_path);
const int fSuccessUnlink = ::unlink(serverAddress.sun_path);
_ASSERTE(fSuccessUnlink != -1);

const int fSuccessClose = ::close(_serverSocket);
const int fSuccessClose = ::close(serverSocket);
_ASSERTE(fSuccessClose != -1);
return false;
}
else
{
_isListening = true;
return true;
#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__
return nullptr;
}

#ifdef __APPLE__
umask(prev_mask);
#endif // __APPLE__

return new IpcStream::DiagnosticsIpc(serverSocket, &serverAddress);
}

IpcStream *IpcStream::DiagnosticsIpc::Accept(ErrorCallback callback)
IpcStream *IpcStream::DiagnosticsIpc::Accept(ErrorCallback callback) const
{
_ASSERTE(mode == ConnectionMode::SERVER);
_ASSERTE(_isListening);

sockaddr_un from;
socklen_t fromlen = sizeof(from);
const int clientSocket = ::accept(_serverSocket, (sockaddr *)&from, &fromlen);
Expand All@@ -164,114 +136,7 @@ IpcStream *IpcStream::DiagnosticsIpc::Accept(ErrorCallback callback)
return nullptr;
}

return new IpcStream(clientSocket, mode);
}

IpcStream *IpcStream::DiagnosticsIpc::Connect(ErrorCallback callback)
{
_ASSERTE(mode == ConnectionMode::CLIENT);

sockaddr_un clientAddress{};
clientAddress.sun_family = AF_UNIX;
const int clientSocket = ::socket(AF_UNIX, SOCK_STREAM, 0);
if (clientSocket == -1)
{
if (callback != nullptr)
callback(strerror(errno), errno);
return nullptr;
}

// We don't expect this to block since this is a Unix Domain Socket. `connect` may block until the
// TCP handshake is complete for TCP/IP sockets, but UDS don't use TCP. `connect` will return even if
// the server hasn't called `accept`.
if (::connect(clientSocket, (struct sockaddr *)_pServerAddress, sizeof(*_pServerAddress)) < 0)
{
if (callback != nullptr)
callback(strerror(errno), errno);
return nullptr;
}

return new IpcStream(clientSocket, ConnectionMode::CLIENT);
}

int32_t IpcStream::DiagnosticsIpc::Poll(IpcPollHandle *rgIpcPollHandles, uint32_t nHandles, int32_t timeoutMs, ErrorCallback callback)
{
// prepare the pollfd structs
pollfd *pollfds = new pollfd[nHandles];
for (uint32_t i = 0; i < nHandles; i++)
{
rgIpcPollHandles[i].revents = 0; // ignore any values in revents
int fd = -1;
if (rgIpcPollHandles[i].pIpc != nullptr)
{
// SERVER
_ASSERTE(rgIpcPollHandles[i].pIpc->mode == ConnectionMode::SERVER);
fd = rgIpcPollHandles[i].pIpc->_serverSocket;
}
else
{
// CLIENT
_ASSERTE(rgIpcPollHandles[i].pStream != nullptr);
fd = rgIpcPollHandles[i].pStream->_clientSocket;
}

pollfds[i].fd = fd;
pollfds[i].events = POLLIN;
}

int retval = poll(pollfds, nHandles, timeoutMs);

// Check results
if (retval < 0)
{
for (uint32_t i = 0; i < nHandles; i++)
{
if ((pollfds[i].revents & POLLERR) && callback != nullptr)
callback(strerror(errno), errno);
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::ERR;
}
delete[] pollfds;
return -1;
}
else if (retval == 0)
{
// we timed out
delete[] pollfds;
return 0;
}

for (uint32_t i = 0; i < nHandles; i++)
{
if (pollfds[i].revents != 0)
{
// error check FIRST
if (pollfds[i].revents & POLLHUP)
{
// check for hangup first because a closed socket
// will technically meet the requirements for POLLIN
// i.e., a call to recv/read won't block
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::HANGUP;
delete[] pollfds;
return -1;
}
else if ((pollfds[i].revents & (POLLERR|POLLNVAL)))
{
if (callback != nullptr)
callback("Poll error", (uint32_t)pollfds[i].revents);
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::ERR;
delete[] pollfds;
return -1;
}
else if (pollfds[i].revents & POLLIN)
{
rgIpcPollHandles[i].revents = (uint8_t)PollEvents::SIGNALED;
break;
}
}
}

delete[] pollfds;
return 1;
return new IpcStream(clientSocket);
}

void IpcStream::DiagnosticsIpc::Close(ErrorCallback callback)
Expand DownExpand Up@@ -307,101 +172,45 @@ void IpcStream::DiagnosticsIpc::Unlink(ErrorCallback callback)
}

IpcStream::~IpcStream()
{
Close();
}

void IpcStream::Close(ErrorCallback)
{
if (_clientSocket != -1)
{
Flush();

const int fSuccessClose = ::close(_clientSocket);
_ASSERTE(fSuccessClose != -1);
_clientSocket = -1;
}
}

bool IpcStream::Read(void *lpBuffer, const uint32_t nBytesToRead, uint32_t &nBytesRead, const int32_t timeoutMs)
bool IpcStream::Read(void *lpBuffer, const uint32_t nBytesToRead, uint32_t &nBytesRead) const
{
_ASSERTE(lpBuffer != nullptr);

if (timeoutMs != InfiniteTimeout)
{
pollfd pfd;
pfd.fd = _clientSocket;
pfd.events = POLLIN;
int retval = poll(&pfd, 1, timeoutMs);
if (retval <= 0 || pfd.revents != POLLIN)
{
// timeout or error
return false;
}
// else fallthrough
}

uint8_t *lpBufferCursor = (uint8_t*)lpBuffer;
ssize_t currentBytesRead = 0;
ssize_t totalBytesRead = 0;
bool fSuccess = true;
while (fSuccess && nBytesToRead - totalBytesRead > 0)
{
currentBytesRead = ::recv(_clientSocket, lpBufferCursor, nBytesToRead - totalBytesRead, 0);
fSuccess = currentBytesRead != 0;
if (!fSuccess)
break;
totalBytesRead += currentBytesRead;
lpBufferCursor += currentBytesRead;
}
const ssize_t ssize = ::recv(_clientSocket, lpBuffer, nBytesToRead, 0);
const bool fSuccess = ssize != -1;

if (!fSuccess)
{
// TODO: Add error handling.
}

nBytesRead = static_cast<uint32_t>(totalBytesRead);
nBytesRead = static_cast<uint32_t>(ssize);
return fSuccess;
}

bool IpcStream::Write(const void *lpBuffer, const uint32_t nBytesToWrite, uint32_t &nBytesWritten, const int32_t timeoutMs)
bool IpcStream::Write(const void *lpBuffer, const uint32_t nBytesToWrite, uint32_t &nBytesWritten) const
{
_ASSERTE(lpBuffer != nullptr);

if (timeoutMs != InfiniteTimeout)
{
pollfd pfd;
pfd.fd = _clientSocket;
pfd.events = POLLOUT;
int retval = poll(&pfd, 1, timeoutMs);
if (retval <= 0 || pfd.revents != POLLOUT)
{
// timeout or error
return false;
}
// else fallthrough
}

uint8_t *lpBufferCursor = (uint8_t*)lpBuffer;
ssize_t currentBytesWritten = 0;
ssize_t totalBytesWritten = 0;
bool fSuccess = true;
while (fSuccess && nBytesToWrite - totalBytesWritten > 0)
{
currentBytesWritten = ::send(_clientSocket, lpBufferCursor, nBytesToWrite - totalBytesWritten, 0);
fSuccess = currentBytesWritten != -1;
if (!fSuccess)
break;
lpBufferCursor += currentBytesWritten;
totalBytesWritten += currentBytesWritten;
}
const ssize_t ssize = ::send(_clientSocket, lpBuffer, nBytesToWrite, 0);
const bool fSuccess = ssize != -1;

if (!fSuccess)
{
// TODO: Add error handling.
}

nBytesWritten = static_cast<uint32_t>(totalBytesWritten);
nBytesWritten = static_cast<uint32_t>(ssize);
return fSuccess;
}

Expand Down
Loading