#include "mungal4.h" #pragma hdrstop #include "l4nettransport.h" #include "..\munga\appmgr.h" #include //######################################################################## // Waiting for a peer without hanging the window //######################################################################## namespace { Logical gInWait = False; // re-entrancy guard Logical gWaitCancelled = False; char gOriginalTitle[256] = ""; Logical gTitleSaved = False; } //~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ // // How long Connect keeps redialling. // // It was two minutes, which was right for what it was written for: an // arcade console redialling a pod that is still booting and WILL answer, // on a LAN with nothing else to go wrong. Over Steam the same two minutes // is a host staring at a dead window, because a peer that has not answered // in twenty seconds is not coming - it never launched, or it is behind // something the relay cannot cross. // int NetTransport_ConnectWaitSeconds() { static int cached = -1; if (cached < 0) { const char *setting = getenv("RP412CONNECTWAIT"); cached = (setting != NULL) ? atoi(setting) : 20; if (cached < 2) { cached = 2; } if (cached > 300) { cached = 300; } } return cached; } //~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ // void NetTransport_SetWaitProgress(const char *text) { if (ghWnd == 0) { return; } if (!gTitleSaved) { GetWindowTextA(ghWnd, gOriginalTitle, sizeof(gOriginalTitle) - 1); gTitleSaved = True; } // // The title bar because it is the one surface guaranteed to exist // here: this runs before the engine block, so there is no renderer to // draw a progress screen with yet. // SetWindowTextA(ghWnd, (text != NULL) ? text : gOriginalTitle); } //~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ // NetWaitResult NetTransport_PumpAndSleep(int milliseconds) { // // Dispatching a message can run application code, and that code can // reach a connect of its own. Nested pumping would then deliver the // same messages twice and let an inner wait consume the escape meant // for the outer one, so an inner call just sleeps. // if (gInWait) { Sleep(milliseconds); return NetWaitContinue; } gInWait = True; NetWaitResult result = NetWaitContinue; MSG message; while (PeekMessage(&message, NULL, 0, 0, PM_REMOVE)) { if (message.message == WM_QUIT) { // // Put it back: the message loop that owns the shutdown has // to see this, not us. // PostQuitMessage((int) message.wParam); result = NetWaitQuit; break; } if (message.message == WM_KEYDOWN && message.wParam == VK_ESCAPE) { gWaitCancelled = True; } TranslateMessage(&message); DispatchMessage(&message); } if (result == NetWaitContinue && gWaitCancelled) { result = NetWaitCancelled; } if (result == NetWaitContinue) { Sleep(milliseconds); } gInWait = False; return result; } //~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ // // Cleared when a connect sequence starts, so an escape pressed during one // race does not cancel the next one. // void NetTransport_ClearWaitCancel() { gWaitCancelled = False; } //######################################################################## // Outbound send queue - shared machinery for both transports. // See the block comment in l4nettransport.h for why this exists. //######################################################################## Logical NetTransport_SendQueueEnabled() { static int cached = -1; if (cached < 0) { const char *setting = getenv("RP412NETSENDQ"); cached = (setting != NULL && atoi(setting) == 0) ? 0 : 1; } return cached ? True : False; } int NetTransport_SendQueueCapBytes() { static int cached = -1; if (cached < 0) { const char *setting = getenv("RP412NETQUEUEKB"); int kilobytes = (setting != NULL) ? atoi(setting) : 256; if (kilobytes < 16) { kilobytes = 16; } if (kilobytes > 4096) { kilobytes = 4096; } cached = kilobytes * 1024; } return cached; } Logical NetTransport_StatsEnabled() { static int cached = -1; if (cached < 0) { const char *setting = getenv("RP412NETSTATS"); cached = (setting != NULL && atoi(setting) != 0) ? 1 : 0; } return cached ? True : False; } NetSendQueueTable::NetSendQueueTable() { for (int i = 0; i < MaxQueues; ++i) { queues[i].connection = NetTransport::InvalidConnection; queues[i].buffer = NULL; queues[i].head = 0; queues[i].tail = 0; queues[i].dead = False; queues[i].sentMessages = 0; queues[i].sentBytes = 0; queues[i].wireWrites = 0; queues[i].partialWrites = 0; queues[i].wouldBlocks = 0; queues[i].queuedHighWater = 0; } } NetSendQueueTable::~NetSendQueueTable() { ReleaseAll(); } NetSendQueue * NetSendQueueTable::Find(NetTransport::Connection connection) { for (int i = 0; i < MaxQueues; ++i) { if (queues[i].connection == connection) { return &queues[i]; } } return NULL; } NetSendQueue * NetSendQueueTable::FindOrAdopt(NetTransport::Connection connection) { NetSendQueue *queue = Find(connection); if (queue != NULL) { return queue; } for (int i = 0; i < MaxQueues; ++i) { if (queues[i].connection == NetTransport::InvalidConnection) { queue = &queues[i]; queue->connection = connection; queue->head = 0; queue->tail = 0; queue->dead = False; queue->sentMessages = 0; queue->sentBytes = 0; queue->wireWrites = 0; queue->partialWrites = 0; queue->wouldBlocks = 0; queue->queuedHighWater = 0; return queue; } } DEBUG_STREAM << "NetSendQueue: table full - connection " << (unsigned long) connection << " gets no queue\n" << std::flush; return NULL; } void NetSendQueueTable::Release(NetTransport::Connection connection) { NetSendQueue *queue = Find(connection); if (queue == NULL) { return; } if (queue->buffer != NULL) { free(queue->buffer); queue->buffer = NULL; } queue->connection = NetTransport::InvalidConnection; queue->head = 0; queue->tail = 0; queue->dead = False; } void NetSendQueueTable::ReleaseAll() { for (int i = 0; i < MaxQueues; ++i) { if (queues[i].connection != NetTransport::InvalidConnection) { Release(queues[i].connection); } } } Logical NetSendQueueTable::Append( NetSendQueue *queue, const void *data, int size) { int capacity = NetTransport_SendQueueCapBytes(); int pending = queue->PendingBytes(); if (pending + size > capacity) { // // Overflow. These are reliable ordered messages - dropping any of // them desyncs the stream - so the only honest answer is to give // up on the connection and let the disconnect path run. // DEBUG_STREAM << "NetSendQueue: connection " << (unsigned long) queue->connection << " overflowed (" << pending << " pending + " << size << " > " << capacity << ") - declaring it dead\n" << std::flush; if (queue->buffer != NULL) { free(queue->buffer); queue->buffer = NULL; } queue->head = 0; queue->tail = 0; queue->dead = True; return False; } if (queue->buffer == NULL) { queue->buffer = (unsigned char *) malloc(capacity); if (queue->buffer == NULL) { queue->dead = True; return False; } } // // Rejustify before appending when the tail would run off the end - // the bytes ahead of head are already on the wire. // if (queue->tail + size > capacity && queue->head > 0) { memmove(queue->buffer, &queue->buffer[queue->head], pending); queue->head = 0; queue->tail = pending; } memcpy(&queue->buffer[queue->tail], data, size); queue->tail += size; if ((unsigned long) queue->PendingBytes() > queue->queuedHighWater) { queue->queuedHighWater = (unsigned long) queue->PendingBytes(); } return True; } //######################################################################## // WinsockNetTransport - the TCP wire the arcade always used, moved // verbatim out of L4NET.CPP behind the NetTransport seam. Behavior // notes preserved from the original: // - Connect retries while the remote refuses (WSAECONNREFUSED): the // console feeds eggs in [pilots] order with ACKs, so earlier pods // are listening before later pods open - the retry covers the race. // - Sockets go nonblocking once connected; listeners are nonblocking // from creation. //######################################################################## namespace { class WinsockNetTransport: public NetTransport { public: WinsockNetTransport(): started(False), statsNextTick(0) { } Logical Startup() { if (started) { return True; } int result = WSAStartup(MAKEWORD(2, 2), &winsockData); if (result != NO_ERROR) { DEBUG_STREAM << "ERROR: WSAStartup() failed with " << result << "!\n" << std::flush; return False; } started = True; return True; } void Cleanup() { if (started) { FlushSends(); sendQueues.ReleaseAll(); WSACleanup(); started = False; } } int GetLocalAddresses( unsigned long *addresses, int max_count ) { char name[255]; PHOSTENT hostinfo; if (gethostname(name, sizeof(name)) != 0) { DEBUG_STREAM << "ERROR: gethostname() failed!" << std::endl << std::flush; return 0; } if ((hostinfo = gethostbyname(name)) == NULL) { DEBUG_STREAM << "ERROR: gethostbyname() failed!" << std::endl << std::flush; return 0; } int count = 0; for (int i = 0; hostinfo->h_addr_list[i] != NULL && count < max_count; ++i) { addresses[count++] = *((unsigned long *) hostinfo->h_addr_list[i]); } // loopback rounds out the list (single-machine testing) if (count < max_count) { addresses[count++] = 0x0100007F; } return count; } Logical Resolve( const char *host_name, SOCKADDR_IN *address ) { // numeric ip[:port] only - the egg carries addresses, not // names; port stays 0 when absent (caller applies default) int buffer_size = sizeof(SOCKADDR_IN); return WSAStringToAddressA( (LPSTR) host_name, AF_INET, NULL, (LPSOCKADDR) address, &buffer_size) == 0; } Connection Connect( const SOCKADDR_IN *remote, int local_port ) { DEBUG_STREAM << "Opening connection to " << inet_ntoa(remote->sin_addr) << ":" << ntohs(remote->sin_port) << "...\n" << std::flush; // // Retry-while-refused, bounded: the egg-ACK ordering means // the peer may not be listening yet. A refused TCP socket // is dead - every attempt needs a fresh one. // // // NOTE the connect() below is still BLOCKING - the socket only // goes nonblocking once it has succeeded - so an unreachable // host (filtered rather than refused) sits in the OS SYN retry // for around twenty seconds with nothing we can do about it. // Pumping between redials fixes the refused case, which is the // common one; the unreachable case needs the socket made // nonblocking before connect() and a select() on our own // timeout. Not done here because this is the LAN/direct path - // a Steam host goes through SteamNetTransport::Connect. // DWORD deadline = GetTickCount() + (DWORD) NetTransport_ConnectWaitSeconds() * 1000; char progress[256]; sprintf( progress, "Red Planet - connecting to %s ...", inet_ntoa(remote->sin_addr) ); NetTransport_SetWaitProgress(progress); for (;;) { SOCKET sock = socket(AF_INET, SOCK_STREAM, IPPROTO_TCP); if (sock == INVALID_SOCKET) { DEBUG_STREAM << "ERROR: socket() failed with " << WSAGetLastError() << "!\n"; return InvalidConnection; } bool reuse_address = true; if (setsockopt(sock, SOL_SOCKET, SO_REUSEADDR, (char *) &reuse_address, sizeof(bool))) { DEBUG_STREAM << "ERROR: Could not set SO_REUSEADDR on socket; setsockopt() failed with " << WSAGetLastError() << "!\n" << std::flush; closesocket(sock); return InvalidConnection; } // bind the game port: the mesh identifies peers by // address AND port, so the source port must be ours // (port 0 = ephemeral, for console channels) sockaddr_in local_endpoint; memset(&local_endpoint, 0, sizeof(local_endpoint)); local_endpoint.sin_family = AF_INET; local_endpoint.sin_port = htons((unsigned short) local_port); local_endpoint.sin_addr.S_un.S_addr = INADDR_ANY; if (bind(sock, (sockaddr *) &local_endpoint, sizeof(local_endpoint))) { DEBUG_STREAM << "ERROR: Could not bind local socket; bind() failed with " << WSAGetLastError() << "!\n" << std::flush; closesocket(sock); return InvalidConnection; } if (connect(sock, (sockaddr *) remote, sizeof(SOCKADDR_IN)) == 0) { unsigned long enable = 1; if (ioctlsocket(sock, FIONBIO, &enable)) { DEBUG_STREAM << "ERROR: Could not set actively opened socket to nonblocking; ioctlsocket() failed with " << WSAGetLastError() << "!\n" << std::flush; closesocket(sock); return InvalidConnection; } NetTransport_SetWaitProgress(NULL); return (Connection) sock; } int wsa_error = WSAGetLastError(); closesocket(sock); if (wsa_error != WSAECONNREFUSED || (LONG)(GetTickCount() - deadline) >= 0) { DEBUG_STREAM << "ERROR: connect() failed with " << wsa_error << "!\n" << std::flush; NetTransport_SetWaitProgress(NULL); return InvalidConnection; } // // Peer not up yet - redial, answering the window meanwhile. // DWORD left = (DWORD)((LONG)(deadline - GetTickCount()) / 1000); sprintf( progress, "Red Planet - connecting to %s ... %us left, ESC to cancel", inet_ntoa(remote->sin_addr), (unsigned) left ); NetTransport_SetWaitProgress(progress); for (int slept = 0; slept < 250; slept += 50) { NetWaitResult wait = NetTransport_PumpAndSleep(50); if (wait != NetWaitContinue) { DEBUG_STREAM << "Connect " << ((wait == NetWaitCancelled) ? "cancelled" : "abandoned, quitting") << " by the user\n" << std::flush; NetTransport_SetWaitProgress(NULL); return InvalidConnection; } } } } Connection Listen( int local_port, int backlog ) { DEBUG_STREAM << "Starting to listen on port " << local_port << "...\n" << std::flush; SOCKET listener = socket(AF_INET, SOCK_STREAM, IPPROTO_TCP); if (listener == INVALID_SOCKET) { DEBUG_STREAM << "ERROR: Could not create listener socket; socket() failed with " << WSAGetLastError() << "!\n" << std::flush; return InvalidConnection; } bool reuse_address = true; if (setsockopt(listener, SOL_SOCKET, SO_REUSEADDR, (char *) &reuse_address, sizeof(bool))) { DEBUG_STREAM << "ERROR: Could not set SO_REUSEADDR on listener socket; setsockopt() failed with " << WSAGetLastError() << "!\n" << std::flush; closesocket(listener); return InvalidConnection; } sockaddr_in local_endpoint; memset(&local_endpoint, 0, sizeof(local_endpoint)); local_endpoint.sin_family = AF_INET; local_endpoint.sin_port = htons((unsigned short) local_port); local_endpoint.sin_addr.S_un.S_addr = INADDR_ANY; if (bind(listener, (sockaddr *) &local_endpoint, sizeof(local_endpoint))) { DEBUG_STREAM << "ERROR: Could not bind listener socket; bind() failed with " << WSAGetLastError() << "!\n" << std::flush; closesocket(listener); return InvalidConnection; } if (listen(listener, backlog)) { DEBUG_STREAM << "ERROR: Could not listen on listener socket; listen() failed with " << WSAGetLastError() << "!\n" << std::flush; closesocket(listener); return InvalidConnection; } unsigned long enable = 1; if (ioctlsocket(listener, FIONBIO, &enable)) { DEBUG_STREAM << "ERROR: Could not set listener socket to nonblocking; ioctlsocket() failed with " << WSAGetLastError() << "!\n" << std::flush; closesocket(listener); return InvalidConnection; } return (Connection) listener; } Connection Accept(Connection listener) { SOCKET accepted = accept((SOCKET) listener, NULL, 0); if (accepted == INVALID_SOCKET) { return InvalidConnection; } return (Connection) accepted; } void Close(Connection connection) { // // Best-effort last drain, so an orderly close does not strand // bytes a slow connection still owes the wire. // NetSendQueue *queue = sendQueues.Find(connection); if (queue != NULL) { if (!queue->dead && queue->PendingBytes() > 0) { DrainQueue(queue); } sendQueues.Release(connection); } shutdown((SOCKET) connection, SD_BOTH); closesocket((SOCKET) connection); } int Send( Connection connection, const void *data, int size ) { NetSendQueue *queue = sendQueues.Find(connection); if (queue != NULL && queue->dead) { return -1; } if (!NetTransport_SendQueueEnabled()) { // // Legacy direct send (RP412NETSENDQ=0) - but count and // name the losses the queue would have prevented. // int sent = send((SOCKET) connection, (const char *) data, size, 0); if (sent != size) { static int said = 0; if (said < 8) { ++said; DEBUG_STREAM << "WinsockNetTransport::Send: legacy mode " << ((sent < 0) ? "lost a message of " : "sheared the stream at ") << size << " bytes (send returned " << sent << ", error " << WSAGetLastError() << ")\n" << std::flush; } } return sent; } // // Every live connection gets a table slot (the buffer itself // stays unallocated until something actually queues), so the // RP412NETSTATS counters exist for healthy connections too. // if (queue == NULL) { queue = sendQueues.FindOrAdopt(connection); if (queue == NULL) { // table full: behave exactly like the legacy send return send((SOCKET) connection, (const char *) data, size, 0); } } queue->sentMessages++; queue->sentBytes += size; // // FIFO discipline: bytes already queued must go first, so a // queued connection only ever appends. // if (queue->PendingBytes() > 0) { if (!sendQueues.Append(queue, data, size)) { return -1; } DrainQueue(queue); return size; } int sent = send((SOCKET) connection, (const char *) data, size, 0); if (sent == size) { queue->wireWrites++; return size; } // // The wire would not take all of it. Keep the byte-exact // remainder - a partial send left unhealed is a sheared // stream, the failure this queue exists to close. // if (sent < 0) { DWORD error = WSAGetLastError(); if (error != WSAEWOULDBLOCK) { DEBUG_STREAM << "WinsockNetTransport::Send: hard error " << error << " - declaring the connection dead\n" << std::flush; queue->dead = True; return -1; } queue->wouldBlocks++; sent = 0; } else { queue->wireWrites++; queue->partialWrites++; } if (!sendQueues.Append(queue, (const char *) data + sent, size - sent)) { return -1; } return size; } void FlushSends() { for (int i = 0; i < NetSendQueueTable::MaxQueues; ++i) { NetSendQueue *queue = &sendQueues.queues[i]; if (queue->connection != InvalidConnection && !queue->dead && queue->PendingBytes() > 0) { DrainQueue(queue); } } // // RP412NETSTATS: 5 s per-connection rollup, from here because // FlushSends is the one transport entry the engine calls on a // steady cadence. // if (NetTransport_StatsEnabled()) { DWORD now = GetTickCount(); if ((LONG) (now - statsNextTick) >= 0) { statsNextTick = now + 5000; for (int j = 0; j < NetSendQueueTable::MaxQueues; ++j) { NetSendQueue *queue = &sendQueues.queues[j]; if (queue->connection == InvalidConnection || queue->sentMessages == 0) { continue; } DEBUG_STREAM << "NetStats: tcp conn " << (unsigned long) queue->connection << ": " << queue->sentMessages << " msgs " << queue->sentBytes << " B, " << queue->wireWrites << " writes (" << queue->partialWrites << " partial, " << queue->wouldBlocks << " would-block), queued high " << queue->queuedHighWater << " B, pending " << queue->PendingBytes() << (queue->dead ? " B, DEAD" : " B") << "\n" << std::flush; } } } } int Receive( Connection connection, void *buffer, int size ) { // // A connection the send side declared dead reports as a // disconnect here, so the one disconnect path upstairs runs. // NetSendQueue *queue = sendQueues.Find(connection); if (queue != NULL && queue->dead) { return ReceiveDisconnected; } int received = recv((SOCKET) connection, (char *) buffer, size, 0); if (received > 0) { return received; } if (received == 0) { return ReceiveDisconnected; } DWORD error = WSAGetLastError(); switch (error) { case WSAECONNRESET: // hard drop reads the same as an orderly close upstairs return ReceiveDisconnected; case WSAEWOULDBLOCK: return ReceiveNoData; default: DEBUG_STREAM << "WinsockNetTransport::Receive: recv returned an unexpected error: WSAGetLastError = " << error << std::endl << std::flush; return ReceiveNoData; } } Logical GetRemoteAddress( Connection connection, SOCKADDR_IN *address ) { int size = sizeof(SOCKADDR_IN); memset(address, 0, size); return getpeername((SOCKET) connection, (sockaddr *) address, &size) == 0; } private: // // Push a queue's pending bytes at the wire until it is empty or // the wire refuses. Partial progress advances head byte-exactly. // void DrainQueue(NetSendQueue *queue) { while (queue->PendingBytes() > 0 && !queue->dead) { int sent = send( (SOCKET) queue->connection, (const char *) &queue->buffer[queue->head], queue->PendingBytes(), 0); if (sent > 0) { queue->wireWrites++; queue->head += sent; if (queue->PendingBytes() == 0) { queue->head = 0; queue->tail = 0; } else { queue->partialWrites++; break; } continue; } DWORD error = WSAGetLastError(); if (sent < 0 && error == WSAEWOULDBLOCK) { queue->wouldBlocks++; break; } DEBUG_STREAM << "WinsockNetTransport::DrainQueue: hard error " << error << " with " << queue->PendingBytes() << " bytes pending - declaring the connection dead\n" << std::flush; queue->dead = True; break; } } Logical started; WSADATA winsockData; NetSendQueueTable sendQueues; DWORD statsNextTick; }; WinsockNetTransport gWinsockTransport; NetTransport *gTransport = &gWinsockTransport; } NetTransport * NetTransport_Get() { return gTransport; } void NetTransport_Set(NetTransport *transport) { gTransport = (transport != NULL) ? transport : &gWinsockTransport; }