Commit eda39d4d authored by jan.koester's avatar jan.koester
Browse files

test

parent 02c448e9
Loading
Loading
Loading
Loading
+93 −48
Original line number Diff line number Diff line
@@ -25,6 +25,61 @@ ON ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT
SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
*******************************************************************************/

#include <iostream>
#include <algorithm>
#include <chrono>
#include <mutex>
#include <cstring>

#include <signal.h>
#include <stdlib.h>
#include <stdint.h>

#include <winsock2.h>
#include <ws2ipdef.h>
#include <mswsock.h>
#include <strsafe.h>

#include <errno.h>

#include "socket.h"
#include "exception.h"
#include "eventapi.h"
#include "connection.h"
#include <assert.h>

#define READEVENT 0
#define SENDEVENT 1

#define BLOCKSIZE 16384

/*******************************************************************************
Copyright (c) 2014, Jan Koester jan.koester@gmx.net
All rights reserved.

Redistribution and use in source and binary forms, with or without
modification, are permitted provided that the following conditions are met:
    * Redistributions of source code must retain the above copyright
      notice, this list of conditions and the following disclaimer.
    * Redistributions in binary form must reproduce the above copyright
      notice, this list of conditions and the following disclaimer in the
      documentation and/or other materials provided with the distribution.
    * Neither the name of the <organization> nor the
      names of its contributors may be used to endorse or promote products
      derived from this software without specific prior written permission.

THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS IS" AND
ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED
WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE ARE
DISCLAIMED. IN NO EVENT SHALL <COPYRIGHT HOLDER> BE LIABLE FOR ANY
DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES
(INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES;
LOSS OF USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND
ON ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT
(INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE OF THIS
SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
*******************************************************************************/

#include <iostream>
#include <algorithm>
#include <chrono>
@@ -97,6 +152,8 @@ namespace netplus {
        void       *args;
    };

    class EventWorker {
    public:
class EventWorker {
public:
    /**
@@ -167,52 +224,40 @@ namespace netplus {
                if (pIoCtx->operation == OP_READ) {
                    if (c.csock->_Type == sockettype::SSL) {
                        ssl* sslSocket = static_cast<ssl*>(c.csock.get());

                        // Bridge: Feed ciphertext from network to SSL engine for decryption
                        buffer plain(BLOCKSIZE);
                        size_t decrypted = sslSocket->recvData(plain, 0);

                        if (decrypted > 0) {
                            c.RecvData.append(plain.data.buf, decrypted);
                                eargs->event->RequestEvent(c, tid, args);
                            }
                            // Prepare fresh buffer for next overlapped read
                            buffer readBuf(pClientContext->readCtx.buffer, BLOCKSIZE);
                            sslSocket->recvDataWSA(readBuf, 0);
                            eargs->event->RequestEvent(c, tid, (ULONG_PTR)eargs->args);
                        }
                        else {
                    } else {
                        // Plaintext TCP/UDP
                        c.RecvData.append(pIoCtx->buffer, dwBytesTransfered);
                            eargs->event->RequestEvent(c, tid, args);
                            buffer readBuf(pClientContext->readCtx.buffer, BLOCKSIZE);
                            static_cast<tcp*>(c.csock.get())->recvDataWSA(readBuf, 0);
                        eargs->event->RequestEvent(c, tid, (ULONG_PTR)eargs->args);
                    }
                    }
                    else if (pIoCtx->operation == OP_WRITE) {
                        // Correct iterator-based erasure
                        if (dwBytesTransfered > 0)
                            c.SendData.erase(c.SendData.begin(), c.SendData.begin() + dwBytesTransfered);

                        if (c.SendData.empty()) {
                            eargs->event->ResponseEvent(c, tid, args);
                        }
                    // Continue Read/Write cycle
                    if (!c.SendData.empty()) start_write(pClientContext);
                    else start_read(pClientContext);

                        if (!c.SendData.empty()) {
                            size_t toSend = (std::min)((size_t)BLOCKSIZE, c.SendData.size());
                            buffer out(c.SendData.data(), toSend);
                } else if (pIoCtx->operation == OP_WRITE) {
                    // Update buffer: for SSL, dwBytesTransfered is the encrypted wire size
                    c.SendData.erase(0, dwBytesTransfered);

                            if (c.csock->_Type == sockettype::TCP) {
                                static_cast<tcp*>(c.csock.get())->sendDataWSA(out, 0);
                            }
                            else if (c.csock->_Type == sockettype::UDP) {
                                static_cast<udp*>(c.csock.get())->sendDataWSA(out, 0);
                            }
                            else if (c.csock->_Type == sockettype::SSL) {
                                // Now correctly calling sendDataWSA
                                static_cast<ssl*>(c.csock.get())->sendDataWSA(out, 0);
                            }
                    if (c.SendData.empty()) {
                        eargs->event->ResponseEvent(c, tid, (ULONG_PTR)eargs->args);
                    }

                    if (!c.SendData.empty()) start_write(pClientContext);
                    else start_read(pClientContext);
                }
            } catch (NetException& e) {
                    std::cerr << "IOCP Worker Error: " << e.what() << std::endl;
                    eargs->event->DisconnectEvent(c, tid, args);
                // Ignore non-critical notes, disconnect on actual errors
                if (e.getErrorType() != NetException::Note) {
                    eargs->event->DisconnectEvent(c, tid, (ULONG_PTR)eargs->args);
                    delete pClientContext;
                }
            }
@@ -289,4 +334,4 @@ namespace netplus {
        for (auto& t : threadpool) t.join();
        CloseHandle(iocp);
    }
}
};