레이블이 socket인 게시물을 표시합니다. 모든 게시물 표시
레이블이 socket인 게시물을 표시합니다. 모든 게시물 표시

2021년 4월 30일 금요일

C# 멀티캐스트 소켓 생성

 protected Socket socket;   
 ...  

 void CreateMulticastSocket(string multicastIP, int port, bool exclusive)     
 {  
   sock = new Socket(AddressFamily.InterNetwork, SocketType.Dgram, ProtocolType.Udp);  
   
   socket.Blocking = false;  
   socket.ExclusiveAddressUse = exclusive;  
      
   if (!exclusive) socket.SetSocketOption(SocketOptionLevel.Socket, SocketOptionName.ReuseAddress, 1);  
      
   IPEndPoint endpoint = new IPEndPoint(IPAddress.Any, port);  
   socket.Bind(endpoint);  
    
   IPAddress ip = IPAddress.Parse(multicastIP);  
   socket.SetSocketOption(SocketOptionLevel.IP, SocketOptionName.AddMembership, new MulticastOption(ip, IPAddress.Any));          
 }  

2017년 1월 23일 월요일

C# 유용한 소켓 Select 함수 핸들러 소스코드 - useful socket select method handler source code

소켓 수신 이벤트 처리에 유용한 코드이다.
소켓 이벤트 핸들러를 등록하면 수신 이벤트 발생시 콜백으로 호출해준다.
TCP/UDP-서버/클라이언트 구분없이 모두 사용할 수 있다.

 using System;  
 using System.Collections.Generic;  
 using System.Text;  
 using System.Diagnostics;  
 using System.Net;  
 using System.Net.Sockets;  
 using System.Collections;  
 using System.Threading;  
   
 namespace MySocketLib  
 {  
   public delegate void SocketReadHandlerCallback(Object data);  
   
   public class SocketHandler  
   {  
     public Socket sock;  
     public SocketReadHandlerCallback handler;  
     public object data;  
   }  
   
   public class SocketTaskScheduler  
   {  
     protected Dictionary<Socket, SocketHandler> sockHandlerTable = new Dictionary<Socket, SocketHandler>();  
     protected bool isRunning = false;  
     protected Thread thread;  
     protected int TIMEOUT = 1000000;  
   
     public SocketTaskScheduler()  
     {  
     }  
   
     public void RegisterSocketHandler(Socket sock, SocketReadHandlerCallback handler, object data)  
     {  
       lock (sockHandlerTable)  
       {  
         if (sock != null)  
         {  
           SocketHandler sockHandler = new SocketHandler();  
           sockHandler.sock = sock;  
           sockHandler.handler = handler;  
           sockHandler.data = data;  
           sockHandlerTable.Add(sockHandler.sock, sockHandler);  
         }  
       }  
     }  
   
     public void UnregisterSocketHandler(Socket sock)  
     {  
       lock (sockHandlerTable)  
       {  
         if (sock != null && sockHandlerTable.ContainsKey(sock))  
         {  
           sockHandlerTable.Remove(sock);  
         }  
       }  
     }  
   
     protected void SingleStep()  
     {  
       lock (sockHandlerTable)  
       {  
         try  
         {  
           ArrayList selectList = new ArrayList();  
           foreach (var handler in sockHandlerTable)  
           {  
             selectList.Add(handler.Key);  
           }  
   
           if (selectList.Count == 0)  
           {  
             Thread.Sleep(10);  
             return;  
           }  
   
           Socket.Select(selectList, null, null, TIMEOUT);  
   
           foreach (Socket sock in selectList)  
           {  
             var handler = sockHandlerTable[sock];  
             if (handler != null && handler.handler != null) handler.handler(handler.data);  
           }            
         }  
         catch (Exception ex)  
         {  
           Trace.WriteLine(ex.ToString());  
         }  
       }  
     }  
   
     protected void DoEventLoop()  
     {  
       while (isRunning == true)  
       {  
         SingleStep();  
       }  
     }  
   
     public void StartEventLoop()  
     {  
       if (isRunning == true) return;  
   
       thread = new Thread(new ThreadStart(DoEventLoop));  
       thread.IsBackground = true;  
       isRunning = true;  
       thread.Start();  
     }  
   
     public void StopEventLoop()  
     {  
       isRunning = false;  
       if (thread != null)  
       {  
         if (Thread.CurrentThread != thread)  
         {  
           thread.Join();  
           thread = null;  
         }  
       }  
       sockHandlerTable.Clear();  
     }  
   }  
 }  
   


< 사용법 >

Socket clientSock;
SocketTaskScheduler task = new SocketTaskScheduler();
...
clientSock.Blocking = false;
clientSock.Connect(serverIP, serverPort);

task.RegisterSocketHandler(clientSock, IncomingPacketHandler, this);
task.StartEventLoop();
...

private void IncomingPacketHandler(object data)
{
            int ret = clientSock.Receive(recvBuffer, 0, recvBuffer.Length, SocketFlags.None);
            if (ret <= 0)
            {
                task.UnregisterSocketHandler(clientSock);
                CloseSocket();

                task.StopEventLoop();
            }
            else
            {
                 // process recvBuffer
            }
}




2016년 7월 13일 수요일

소켓 프로그래밍에 대한 조언

언어와 플랫폼을 불문하고 소켓 프로그래밍은 형태가 조금씩 다를뿐이지 내부구현은 완전히 동일하다.
기본적으로 POSIX 소켓 API 로 제공되는 low level API 들이 있고 각 플랫폼마다 이것들을 래핑한 상위 레이어들이 존재하는 형태이다.

응용 프로그램에서는 이러한 low level API를 가지고 개발할 것인지 좀더 편리한 상위 라이브러리를
이용할 것인지 결정해서 개발하면된다.

소켓 프로그래밍을 처음 입문하는 개발자들은 개발에 용이한 상위 라이브러리가 아무래도 많이 끌릴것이다. 심지어 경력 개발자들도 이러한 라이브러리를 이용해서 개발하는것을 선호하기도한다.

개인적으로 소켓 프로그래밍은 반드시 low level API로 개발해야한다고 생각한다. 사실 말이 low level이지 그냥 소켓 API 들이다. C#/Java 등의 언어에서 제공되는 상위 클래스들은 이러한 소켓 API 들을 래핑해서 사용하기 편리하게 만들어주는 역할밖에 하지않는다.

문제는 이러한 라이브러리를 사용해서 개발하다보면 내부구현을 볼 수가 없어서 다양한 상황들이 발생하는 네트워크 프로그래밍에서 올바른 대처를 하기가 어렵다는 점이다. 물론 잘 사용하면 대부분 문제는 없지만 이러한 방법에 길들여지면 소켓 프로그래밍이나 TCP/IP의 기본원리를 이해하기 힘들다는 점도있다.

사실 소켓 프로그래밍을 시작하기 전에 기본적인 TCP/IP 에 대한 이해가 필요하다. 개발을 해가면서 이해하든 미리 공부를 하든 기본적인 TCP/UDP의 동작정도는 알고 시작하는것이 좋다.

여기서 중요한 점은 응용 프로그램과 TCP/IP간의 관계에 대해서 잘 알아야한다는것이다. 
사실 실제적인 데이터 전송과 수신은 TCP/IP가 하는것이지 응용 프로그램이 하는것이 아니다.
응용 프로그램과 TCP/IP는 서로 다른 계층에 존재하고 소켓 API를 통해 서로 데이터를 주고받는데 응용 프로그램의 역할과 TCP/IP의 역할을 잘 알아야 특정한 상황이 발생했을때 이것을 분석해서 원인을 유추하는것이 가능하다.
또한 응용 프로그램간의 프로토콜을 정의하고 구현하는데도 이러한 지식이 필요하다.

마지막으로 한가지 더하자면 통신 프로그램은 아주 정교하게 작성해야한다. TCP 서버/클라이언트 개발자는 잘 알겠지만 통신 상황에서는 너무나도 다양한 케이스들이 발생한다. 대상이 일반 유저라면 생각할 수 있는 거의 모든 경우의 수가 다 발생한다고 보면된다. 
적당히 데이터 주고받는 테스트 정도만 수행하고 릴리즈를 했다간 큰 낭패를 볼 수 있다. 네트워크 상황은 시시각각 변하기때문에 특정 상황에 대한 스냅샷을 잡기가 굉장히 힘들고 이것을 사후에 분석하는것도 어렵기때문에 애초에 robust한 프로그래밍과 테스트가 필요하다.

2016년 7월 4일 월요일

C# 소켓 Select 함수 사용방법 - C# Socket Select

C# 소켓 프로그래밍에서도 기존 POSIX select 함수를 그대로 사용할 수 있다. 방식은 제공되는 Socket 클래스의 Select 정적 멤버함수를 사용하면된다.

< 일반적인 select 사용 >
         try  
         {  
           ArrayList selectList = new ArrayList();  
             
           selectList.Add(mySock);             
   
           if (selectList.Count == 0)  
           {  
             Thread.Sleep(10);  
             return;  
           }  
   
           Socket.Select(selectList, null, null, 1000000);  
   
           foreach (Socket sock in selectList)  
           {                    
             if (sock == mySock)
             {
                 // do something with mySock...
             }
           }            
         }  
         catch (Exception ex)  
         {  
           Trace.WriteLine(ex.ToString());  
         }  

< TCP Connection 타임아웃 >
      
Socket clientSock = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp);
clientSock.Blocking = false;
...
public bool Connect(IPAddress serverIP, int serverPort, int timeout) 
{
       try  
       {  
         clientSock.Connect(serverIP, serverPort);        
         return true;  
       }  
       catch (SocketException ex)  
       {          
         if (ex.SocketErrorCode == SocketError.WouldBlock)  
         {  
           ArrayList selectArray = new ArrayList();  
           selectArray.Add(clientSock);  
   
           Socket.Select(null, selectArray, null, timeout * 1000000);  
   
           if (selectArray.Count == 0)  
           {
             Trace.WriteLine(ex.ToString());  
             return false;  
           }
   
           return true;  
         }    
       }  
       catch (Exception ex)  
       {          
         Trace.WriteLine(ex.ToString());  
         return false;  
       }  
}

2015년 12월 10일 목요일

java/안드로이드 넌블러킹 소켓을 사용하는 TCP 클라이언트/서버 소스 - java/android nonblocking socket tcp server/client source code

자바/안드로이드 환경에서 TCP 서버는 ServerSocketChannel 클래스를 사용해서 넌블러킹 TCP 서버소켓을 생성후 Selector 를 통해 소켓 select 함수 기능을 수행한다.(소켓 다중화)
TCP 클라이언트는 SocketChannel 클래스를 사용해서 넌블러킹 소켓을 구현한다.

* 서버 종료시 현재 연결된 클라이언트들에 대해 3-way handshake 를 통해 자연스럽게 연결종료 시키는 것에 주목(stopServer 메소드)

< TCPServer.java >
 package com.kimdh.dxmediaplayer;  
   
 import java.net.InetSocketAddress;  
 import java.nio.ByteBuffer;  
 import java.nio.channels.SelectionKey;  
 import java.nio.channels.Selector;  
 import java.nio.channels.ServerSocketChannel;  
 import java.nio.channels.SocketChannel;  
 import java.util.ArrayList;  
 import java.util.Iterator;  
 import java.util.List;  
 import java.util.Set;  
   
 import android.os.StrictMode;  
   
 public class TcpServer {  
      protected ServerSocketChannel mChannel;  
      protected ServerThread mThread;  
        
      private static final int BUFFER_SIZE = 1024 * 1024;  
        
      public boolean isOpened() {  
           if (mChannel == null) return false;  
           return mChannel.isOpen();  
      }       
   
      protected void setNetworkThreadPolicy() {  
           StrictMode.ThreadPolicy policy = new StrictMode.ThreadPolicy.Builder().permitAll().build();  
           StrictMode.setThreadPolicy(policy);            
      }  
        
      public boolean startServer(int port, ReceiveEventHandler handler) {  
           setNetworkThreadPolicy();  
             
           try {  
                if (mChannel != null) return false;  
                  
                mChannel = ServerSocketChannel.open();  
                mChannel.configureBlocking(false);  
                mChannel.socket().bind(new InetSocketAddress(port));  
                  
                Selector selector = Selector.open();  
                mChannel.register(selector, SelectionKey.OP_ACCEPT);  
                  
                mThread = new ServerThread(mChannel, selector, handler);                                
                mThread.start();                 
                  
                return true;  
           } catch (Exception ex) {  
                ex.printStackTrace();  
                return false;  
           }            
      }  
        
      public void stopServer() {  
           setNetworkThreadPolicy();  
           try {  
                if (mChannel != null) {       
                     mThread.closeAllClient();  
                       
                     while (mThread.getClientCount() > 0)  
                          Thread.sleep(100);  
                       
                     if (mThread != null) {  
                          mThread.mIsRunning = false;  
                          mThread.join();  
                     }            
                     mChannel.close();  
                     mChannel = null;  
                     System.out.println("tcp server channel closed");  
                }  
           } catch (Exception ex) {  
                ex.printStackTrace();  
           }            
      }  
        
      public void closeAllClient() {  
           setNetworkThreadPolicy();  
           if (mChannel != null)       
                mThread.closeAllClient();            
      }  
        
      public boolean send(SocketChannel client, ByteBuffer buffer) {  
           try {  
                if (client.write(buffer) == buffer.limit())   
                     return true;  
           } catch (Exception ex) {  
                ex.printStackTrace();  
           }  
           return false;  
      }       
        
      public interface ReceiveEventHandler {  
           public void onClientConnected(SocketChannel client);  
           public void onReceived(SocketChannel client, ByteBuffer buffer, int len);  
           public void onClientDisconnected(SocketChannel client);  
      }       
        
      protected class ServerThread extends Thread {  
           private ServerSocketChannel mChannel;  
           private Selector mSelector;  
           private List<SocketChannel> mClientList = new ArrayList<SocketChannel>();  
             
           public boolean mIsRunning = false;  
             
           private ReceiveEventHandler mHandler;  
             
           public ServerThread(ServerSocketChannel channel, Selector selector, ReceiveEventHandler handler) {  
                mChannel = channel;  
                mSelector = selector;  
                mIsRunning = true;  
                mHandler = handler;  
           }  
             
           public void closeAllClient() {  
                synchronized (mClientList) {  
                     for (int i=0; i<mClientList.size(); i++) {  
                          try {  
                               mClientList.get(i).socket().shutdownOutput();  
                          } catch (Exception ex) {  
                               ex.printStackTrace();  
                          }  
                     }                      
                }  
           }  
             
           public int getClientCount() {   
                synchronized (mClientList) {  
                     return mClientList.size();                      
                }                 
           }  
             
           @Override  
           public void run() {  
                System.out.println("server thread start");  
                  
                try {  
                     while (mIsRunning) {  
                          mSelector.select(2*1000);  
                            
                          Set keys = mSelector.selectedKeys();  
                          Iterator i = keys.iterator();  
                            
                          while (i.hasNext()) {  
                               SelectionKey key = (SelectionKey)i.next();  
                               i.remove();  
                                 
                               if (key.isAcceptable()) {  
                                    SocketChannel client = mChannel.accept();  
                                    client.configureBlocking(false);  
                                    client.register(mSelector, SelectionKey.OP_READ);  
                                    synchronized (mClientList) {  
                                         mClientList.add(client);                                          
                                    }                                     
                                    System.out.println("client connected : " + client.socket().getRemoteSocketAddress().toString());  
                                    if (mHandler != null) mHandler.onClientConnected(client);  
                                    continue;  
                               }  
                                 
                               if (key.isReadable()) {  
                                    SocketChannel client = (SocketChannel)key.channel();  
                                    ByteBuffer buffer = ByteBuffer.allocate(BUFFER_SIZE);  
                                    buffer.clear();  
                                      
                                    int ret;  
                                    try {  
                                         ret = client.read(buffer);  
                                    } catch (Exception ex) {  
                                         ex.printStackTrace();  
                                         key.cancel();  
                                         synchronized (mClientList) {  
                                              mClientList.remove(client);                                               
                                         }                                          
                                         System.out.println("client disconnected : " + client.socket().getRemoteSocketAddress().toString());  
                                         if (mHandler != null) mHandler.onClientDisconnected(client);  
                                         continue;  
                                    }  
                                      
                                    if (ret <= 0) {  
                                         key.cancel();  
                                         synchronized (mClientList) {  
                                              mClientList.remove(client);                                               
                                         }                                                                                       
                                         System.out.println("socket read : " + ret + ", client disconnected : " + client.socket().getRemoteSocketAddress().toString());  
                                         if (mHandler != null) mHandler.onClientDisconnected(client);  
                                         continue;  
                                    }  
                                      
                                    buffer.rewind();  
                                    if (mHandler != null) mHandler.onReceived(client, buffer, ret);  
                               }  
                          }  
                     }  
                } catch (Exception ex) {  
                     ex.printStackTrace();  
                }  
                  
                System.out.println("server thread end");  
           }            
      }  
 }  
   

< TCPClient.java >
 package com.kimdh.dxmediaplayer;  
   
 import java.io.IOException;  
 import java.net.InetSocketAddress;  
 import java.nio.ByteBuffer;  
 import java.nio.channels.SelectionKey;  
 import java.nio.channels.Selector;  
 import java.nio.channels.SocketChannel;  
 import java.util.Iterator;  
 import java.util.Set;  
   
 import android.os.StrictMode;  
   
 public class TcpClient {  
      protected SocketChannel mChannel;  
      protected ReceiveThread mThread;  
        
      public int RECV_BUFFER_SIZE = 1024 * 1024;  
        
      public boolean isConnected() {  
           if (mChannel == null) return false;  
           return mChannel.isConnected();  
      }  
        
      protected void setNetworkThreadPolicy() {  
           StrictMode.ThreadPolicy policy = new StrictMode.ThreadPolicy.Builder().permitAll().build();  
           StrictMode.setThreadPolicy(policy);            
      }  
             
      public boolean connect(String ipAddress, short port, int timeout, ReceiveEventHandler handler) {  
           setNetworkThreadPolicy();  
                                 
           try     {  
                if (mChannel != null && mChannel.isConnected() == true)  
                     return false;  
                  
                mChannel = SocketChannel.open();  
                mChannel.configureBlocking(false);  
                mChannel.socket().setReceiveBufferSize(RECV_BUFFER_SIZE);  
                  
                mChannel.connect(new InetSocketAddress(ipAddress, port));  
                  
                Selector selector = Selector.open();                 
                SelectionKey clientKey = mChannel.register(selector, SelectionKey.OP_CONNECT);  
                                 
                if (selector.select(timeout*1000) > 0) {  
                     if (clientKey.isConnectable()) {  
                          if (mChannel.finishConnect()) {  
                               mThread = new ReceiveThread(mChannel, handler);  
                               mThread.start();  
                               return true;  
                          }  
                     }                      
                     mChannel.close();  
                     mChannel = null;  
                     return false;                      
                } else {  
                     return false;  
                }  
           } catch (Exception ex) {  
                ex.printStackTrace();  
                return false;  
           }       
      }  
        
      public void close() {  
           setNetworkThreadPolicy();  
           try {  
                if (mChannel != null) {                      
                     if (mThread != null) {  
                          mThread.mIsRunning = false;  
                          mThread.join();  
                     }  
                     mChannel.close();  
                     mChannel = null;  
                     System.out.println("tcp client channel closed");  
                }  
           } catch (Exception ex) {  
                ex.printStackTrace();  
           }  
      }  
             
      public int send(ByteBuffer buffer) {  
           setNetworkThreadPolicy();  
           if (mChannel == null || !mChannel.isConnected()) return -1;            
           try {  
                return mChannel.write(buffer);  
           } catch (IOException ex) {  
                ex.printStackTrace();  
                return -1;  
           }            
      }  
                  
      public interface ReceiveEventHandler {  
           public void onReceived(ByteBuffer buffer, int len);  
           public void onClosed();  
           public void onThreadEvent();  
      }  
             
      protected class ReceiveThread extends Thread {  
           private SocketChannel mChannel;  
           public boolean mIsRunning = false;  
             
           private ReceiveEventHandler mHandler;  
             
           private static final int BUFFER_SIZE = 1024 * 4;  
             
           public ReceiveThread(SocketChannel channel, ReceiveEventHandler handler) {  
                mChannel = channel;                 
                mHandler = handler;  
                mIsRunning = true;  
           }  
             
           @Override  
           public void run() {  
                System.out.println("receive thread start");  
                  
                try {  
                     Selector selector = Selector.open();  
                     mChannel.register(selector, SelectionKey.OP_READ);  
                  
                     while (mIsRunning) {                 
                          if (selector.select(2*1000) > 0) {  
                               Set<SelectionKey> selectedKey = selector.selectedKeys();  
                               Iterator<SelectionKey> iterator = selectedKey.iterator();  
                                 
                               while (iterator.hasNext()) {  
                                    SelectionKey key = iterator.next();  
                                    iterator.remove();  
                                      
                                    if (key.isReadable()) {  
                                         SocketChannel channel = (SocketChannel)key.channel();  
                                         ByteBuffer buffer = ByteBuffer.allocate(BUFFER_SIZE);  
                                         buffer.clear();  
                                           
                                         int ret;  
                                         try {  
                                              ret = channel.read(buffer);  
                                         } catch (IOException ex) {  
                                              ex.printStackTrace();  
                                              if (mHandler != null) mHandler.onClosed();  
                                              key.cancel();  
                                              continue;  
                                         }  
                                           
                                         if (ret <= 0) {  
                                              System.out.println("SocketChannel.read returned " + ret);  
                                              if (mHandler != null) mHandler.onClosed();  
                                              key.cancel();  
                                              continue;  
                                         }  
                                           
                                         buffer.rewind();                                                                             
                                                                                   
                                         if (mHandler != null) mHandler.onReceived(buffer, ret);  
                                    }                                     
                               }  
                          }  
                          if (mHandler != null) mHandler.onThreadEvent();  
                     }  
                } catch (Exception ex) {  
                     ex.printStackTrace();  
                }  
                  
                System.out.println("receive thread end");  
           }            
      }  
 }  
   


< TCP 서버 사용 >

public class MainActivity extends Activity implements TcpServer.ReceiveEventHandler {
private TcpServer mServer = new TcpServer();

private byte[] mReceiveBuffer = new byte[1024*4];
private int mReceiveBufferIndex = 0;
        ...
public void onClick(View v) {
        ...
mServer.startServer(8112, this);
    }

protected void onDestroy() {
        ...
mServer.stopServer();
    }

    @Override
public void onClientConnected(SocketChannel client) {
            // do something when client connected

    @Override
public void onReceived(SocketChannel client, ByteBuffer buffer, int len) {
System.out.println("onReceived : " + len);
buffer.get(mReceiveBuffer, mReceiveBufferIndex, len);
mReceiveBufferIndex += len;
// process receive buffer
         ...
         ByteBuffer response = ByteBuffer.allocate(len+4);
         buffer.putInt(len);
         buffer.put(4, payload);
         buffer.rewind();
         mServer.send(client, response);
}

    @Override
public void onClientDisconnected(SocketChannel client) {
// do something when client disconnected
}
}

2015년 12월 2일 수요일

윈도우/리눅스/안드로이드 공통 소켓 라이브러리 소스 - common socket library source code for window/linux/android

live555 의 소켓관련 소스를 기본으로 윈도우/리눅스/안드로이드 환경에서 공통으로 사용가능한 소켓 라이브러리 소스이다. 자신의 환경에 맞게 약간의 수정만 하면 각 플랫폼에서 쉽게 사용가능하다.

< NetCommon.h >
 #ifndef __NETCOMMON_H__  
 #define __NETCOMMON_H__  
   
 #ifdef WIN32  
   
 #include <WinSock2.h>  
 #include <ws2tcpip.h>  
   
 #define closeSocket     closesocket  
 #define EWOULDBLOCK     WSAEWOULDBLOCK  
 #define EINPROGRESS WSAEWOULDBLOCK  
 #define EINTR          WSAEINTR  
   
 #define _strcasecmp     _strnicmp  
 #define snprintf     _snprintf  
   
 #else  
   
 #include <sys/socket.h>  
 #include <netinet/in.h>  
 #include <netinet/tcp.h>  
 #include <arpa/inet.h>  
 #include <unistd.h>  
 #include <fcntl.h>  
 #include <errno.h>  
   
 #define closeSocket               close  
 #define WSAGetLastError()     errno  
   
 #include <ctype.h>  
 #include <stdlib.h>  
 #define _strcasecmp strncasecmp  
 #endif  
   
 #endif  
   

< SockCommon.h >
 #ifndef __SOCK_COMMON_H__  
 #define __SOCK_COMMON_H__  
   
 #include "NetCommon.h"  
   
 int setupStreamSock(short port, int makeNonBlocking);  
 int setupDatagramSock(short port, int makeNonBlocking);  
 int setupServerSock(short port, int makeNonBlocking);  
 int setupClientSock(int serverSock, int makeNonBlocking, struct sockaddr_in& clientAddr);  
 int makeSocketNonBlocking(int sock);  
   
 int makeTCP_NoDelay(int sock);  
   
 unsigned setSendBufferTo(int sock, unsigned requestedSize);  
 unsigned setReceiveBufferTo(int sock, unsigned requestedSize);  
 unsigned getSendBufferSize(int sock);  
 unsigned getReceiveBufferSize(int sock);  
   
 unsigned getBufferSize(int bufOptName, int sock);  
 unsigned setBufferSizeTo(int bufOptName, int sock, int requestedSize);  
   
 int blockUntilReadable(int sock, struct timeval* timeout);  
   
 int readSocket1(int sock, char *buffer, unsigned bufferSize, struct sockaddr_in &fromAddress);  
 int readSocket(int sock, char *buffer, unsigned bufferSize, struct sockaddr_in &fromAddress, struct timeval *timeout = NULL);  
 int readSocketExact(int sock, char *buffer, unsigned bufferSize, struct sockaddr_in &fromAddress, struct timeval *timeout = NULL);  
   
 int writeSocket(int sock, char *buffer, unsigned bufferSize);  
 int writeSocket(int sock, char *buffer, unsigned bufferSize, struct sockaddr_in& toAddress);  
   
 int sendRTPOverTCP(int sock, char *buffer, int len, unsigned char streamChannelId);  
   
 void shutdown(int sock);  
   
 bool isMulticastAddress(unsigned int address);  
 bool socketJoinGroupSSM(int sock, unsigned int groupAddress, unsigned int sourceFilterAddr);  
 bool socketLeaveGroupSSM(int sock, unsigned int groupAddress, unsigned int sourceFilterAddr);  
 bool socketJoinGroup(int sock, unsigned int groupAddress);  
 bool socketLeaveGroup(int sock, unsigned int groupAddress);  
   
 unsigned int ourIPAddress();  
   
 extern unsigned int ReceivingInterfaceAddr;  
 #endif  
   

< SockCommon.cpp >
 #include "SockCommon.h"  
 #include "RTSPCommonEnv.h"  
 #include <stdio.h>  
   
 #ifdef WIN32  
 #pragma comment(lib, "ws2_32.lib")  
 #elif defined(LINUX)  
 #include <string.h>  
 #endif  
   
 #define MAKE_SOCKADDR_IN(var,adr,prt) /*adr,prt must be in network order*/\  
   struct sockaddr_in var;\  
   var.sin_family = AF_INET;\  
   var.sin_addr.s_addr = (adr);\  
   var.sin_port = htons(prt);\  
   
 #ifdef WIN32  
 #define WS_VERSION_CHOICE1 0x202/*MAKEWORD(2,2)*/  
 #define WS_VERSION_CHOICE2 0x101/*MAKEWORD(1,1)*/  
 int initializeWinsockIfNecessary(void) {  
      /* We need to call an initialization routine before  
      * we can do anything with winsock. (How fucking lame!):  
      */  
      static int _haveInitializedWinsock = 0;  
      WSADATA     wsadata;  
   
      if (!_haveInitializedWinsock) {  
           if ((WSAStartup(WS_VERSION_CHOICE1, &wsadata) != 0)  
                && ((WSAStartup(WS_VERSION_CHOICE2, &wsadata)) != 0)) {  
                     return 0; /* error in initialization */  
           }  
           if ((wsadata.wVersion != WS_VERSION_CHOICE1)  
                && (wsadata.wVersion != WS_VERSION_CHOICE2)) {  
                     WSACleanup();  
                     return 0; /* desired Winsock version was not available */  
           }  
           _haveInitializedWinsock = 1;  
      }  
   
      return 1;  
 }  
 #else  
 #define initializeWinsockIfNecessary()     1  
 #endif  
   
 void socketErr(char *lpszFormat,...)  
 {  
      va_list args;  
      int len;  
      char *buffer;  
   
      va_start(args, lpszFormat);  
   
      len = _vscprintf(lpszFormat, args) + 32;  
      buffer = (char *)malloc(len * sizeof(char));  
   
      vsprintf(buffer, lpszFormat, args);  
   
 #ifdef WIN32  
      if (RTSPCommonEnv::nDebugPrint == 0) {  
           fprintf(stdout, buffer);  
           fprintf(stdout, "%d\n", WSAGetLastError());  
      } else if (RTSPCommonEnv::nDebugPrint == 1) {  
           OutputDebugString(buffer);  
           char tmp[16] = {0};  
           sprintf(tmp, "%d\n", WSAGetLastError());  
           OutputDebugString(tmp);  
      }  
 #elif defined(ANDROID)  
      DPRINTF0(buffer);  
      char tmp[16] = {0};  
      sprintf(tmp, "%d\n", WSAGetLastError());  
      DPRINTF0(tmp);  
 #else  
      fprintf(stdout, buffer);  
      fprintf(stdout, "%d\n", WSAGetLastError());  
 #endif  
   
      free(buffer);  
 }  
   
 static int reuseFlag = 1;  
   
 int setupStreamSock(short port, int makeNonBlocking)  
 {  
      if (!initializeWinsockIfNecessary()) {  
           socketErr("[%s] Failed to initialize 'winsock': ", __FUNCTION__);  
           return -1;  
      }  
   
      int newSocket = socket(AF_INET, SOCK_STREAM, 0);  
      if (newSocket < 0) {  
           DPRINTF("%s:%d\n",__FUNCTION__,__LINE__);  
           socketErr("[%s] unable to create stream socket: ", __FUNCTION__);  
           return newSocket;  
      }  
 #if 0  
      if (setsockopt(newSocket, SOL_SOCKET, SO_REUSEADDR,  
           (const char*)&reuseFlag, sizeof reuseFlag) != 0) {  
                socketErr("[%s] setsockopt(SO_REUSEADDR) error: ", __FUNCTION__);  
                closeSocket(newSocket);  
                return -1;  
      }  
 #endif  
      struct sockaddr_in c_addr;  
      memset(&c_addr, 0, sizeof(c_addr));  
      c_addr.sin_addr.s_addr = INADDR_ANY;  
      c_addr.sin_family = AF_INET;  
      c_addr.sin_port = htons(port);  
   
      if (bind(newSocket, (struct sockaddr*)&c_addr, sizeof c_addr) != 0) {  
           socketErr("[%s] bind() error (port number: %d): ", __FUNCTION__, port);  
           closeSocket(newSocket);  
           return -1;  
      }  
   
      if (makeNonBlocking) {  
           if (!makeSocketNonBlocking(newSocket)) {  
                socketErr("[%s] failed to make non-blocking: ", __FUNCTION__);  
                closeSocket(newSocket);  
                return -1;  
           }  
      }  
   
      return newSocket;  
 }  
   
 int setupDatagramSock(short port, int makeNonBlocking)  
 {  
      if (!initializeWinsockIfNecessary()) {  
           socketErr("[%s] Failed to initialize 'winsock': ", __FUNCTION__);  
           return -1;  
      }  
   
      int newSocket = socket(AF_INET, SOCK_DGRAM, 0);  
      if (newSocket < 0) {  
           socketErr("[%s] unable to create datagram socket: ", __FUNCTION__);  
           return newSocket;  
      }  
 #if 0  
      if (setsockopt(newSocket, SOL_SOCKET, SO_REUSEADDR,  
           (const char*)&reuseFlag, sizeof reuseFlag) < 0) {  
                socketErr("setsockopt(SO_REUSEADDR) error: ", __FUNCTION__);  
                closeSocket(newSocket);  
                return -1;  
      }  
 #endif  
      struct sockaddr_in c_addr;  
      memset(&c_addr, 0, sizeof(c_addr));  
      c_addr.sin_addr.s_addr = INADDR_ANY;  
      c_addr.sin_family = AF_INET;  
      c_addr.sin_port = htons(port);  
   
      if (bind(newSocket, (struct sockaddr*)&c_addr, sizeof c_addr) != 0) {  
           socketErr("[%s] bind() error (port number: %d): ", __FUNCTION__, port);  
           closeSocket(newSocket);  
           return -1;  
      }  
   
      if (makeNonBlocking) {  
           if (!makeSocketNonBlocking(newSocket)) {  
                socketErr("[%s] failed to make non-blocking: ", __FUNCTION__);  
                closeSocket(newSocket);  
                return -1;  
           }  
      }  
   
      return newSocket;  
 }  
   
 #define LISTEN_BACKLOG_SIZE 20  
   
 int setupServerSock(short port, int makeNonBlocking)  
 {  
      int sock = setupStreamSock(port, makeNonBlocking);  
      if (sock < 0) return sock;  
   
      if (listen(sock, LISTEN_BACKLOG_SIZE) != 0) {  
           socketErr("[%s] failed to listen sock: ", __FUNCTION__);  
           closeSocket(sock);  
           return -1;  
      }  
   
      return sock;  
 }  
   
 int setupClientSock(int serverSock, int makeNonBlocking, struct sockaddr_in& clientAddr)  
 {  
      socklen_t clientAddrLen = sizeof clientAddr;  
      int clientSock = accept(serverSock, (struct sockaddr*)&clientAddr, &clientAddrLen);  
      if (clientSock < 0) {  
           int err = WSAGetLastError();  
           if (err != EWOULDBLOCK) {  
                socketErr("[%s] accept() failed: ", __FUNCTION__);  
                closeSocket(clientSock);  
                return -1;  
           }  
           return 0;  
      }  
      makeSocketNonBlocking(clientSock);  
   
      return clientSock;  
 }  
   
 int makeSocketNonBlocking(int sock)  
 {  
 #ifdef WIN32  
      unsigned long arg = 1;  
      return ioctlsocket(sock, FIONBIO, &arg) == 0;  
 #else  
      int curFlags = fcntl(sock, F_GETFL, 0);  
      return fcntl(sock, F_SETFL, curFlags|O_NONBLOCK) >= 0;  
 #endif  
 }  
   
 int makeTCP_NoDelay(int sock)  
 {  
      int flag = 1;  
      int err = setsockopt(sock, IPPROTO_TCP, TCP_NODELAY, (char *)&flag, sizeof(flag));  
   
      if(err != 0)  
           socketErr("[%s] setsocket TCPNODELAY error: ", __FUNCTION__);  
   
      return 0;  
 }  
   
 unsigned setSendBufferTo(int sock, unsigned requestedSize)  
 {  
      return setBufferSizeTo(SO_SNDBUF, sock, requestedSize);  
 }  
   
 unsigned setReceiveBufferTo(int sock, unsigned requestedSize)  
 {  
      return setBufferSizeTo(SO_RCVBUF, sock, requestedSize);  
 }  
   
 unsigned getSendBufferSize(int sock)  
 {  
      return getBufferSize(SO_SNDBUF, sock);  
 }  
   
 unsigned getReceiveBufferSize(int sock)  
 {  
      return getBufferSize(SO_RCVBUF, sock);  
 }  
   
 unsigned getBufferSize(int bufOptName, int sock)  
 {  
      unsigned curSize;  
      socklen_t sizeSize = sizeof curSize;  
      if (getsockopt(sock, SOL_SOCKET, bufOptName,  
           (char*)&curSize, &sizeSize) < 0) {  
                socketErr("getBufferSize() error: ", __FUNCTION__);  
                return 0;  
      }  
   
      return curSize;  
 }  
   
 unsigned setBufferSizeTo(int bufOptName, int sock, int requestedSize)  
 {  
      socklen_t sizeSize = sizeof requestedSize;  
      if (setsockopt(sock, SOL_SOCKET, bufOptName, (char*)&requestedSize, sizeSize) != 0)  
           socketErr("setBufferSizeTo() error: ", __FUNCTION__);  
      return getBufferSize(bufOptName, sock);  
 }  
   
 int blockUntilReadable(int sock, timeval *timeout)  
 {  
      int result = -1;  
   
      do {  
           fd_set rd_set;  
           FD_ZERO(&rd_set);  
           if (sock < 0) break;  
   
           FD_SET((unsigned) sock, &rd_set);  
           const unsigned numFds = sock+1;  
   
           result = select(numFds, &rd_set, NULL, NULL, timeout);  
           if (timeout != NULL && result == 0) {  
                break; // this is OK - timeout occurred  
           } else if (result <= 0) {  
                int err = WSAGetLastError();  
                if (err == EINTR || err == EWOULDBLOCK) continue;  
                socketErr("[%s] select() error: ", __FUNCTION__);  
                break;  
           }  
   
           if (!FD_ISSET(sock, &rd_set)) {  
                socketErr("[%s] select() error - !FD_ISSET", __FUNCTION__);  
                break;  
           }  
      } while (0);  
   
      return result;  
 }  
   
 int readSocket1(int sock, char *buffer, unsigned bufferSize, struct sockaddr_in &fromAddress)  
 {  
      int bytesRead;  
      socklen_t addressSize = sizeof fromAddress;  
   
      bytesRead = recvfrom(sock, buffer, bufferSize, 0, (struct sockaddr*)&fromAddress, &addressSize);  
   
      return bytesRead;  
 }  
   
 int readSocket(int sock, char *buffer, unsigned int bufferSize, sockaddr_in &fromAddress, timeval *timeout)  
 {  
      int bytesRead = -1;  
   
      do {  
           int result = blockUntilReadable(sock, timeout);  
           if (timeout != NULL && result == 0) {  
                bytesRead = 0;  
                break;  
           } else if (result <= 0) {  
                break;  
           }  
   
           socklen_t addressSize = sizeof fromAddress;  
           bytesRead = recvfrom(sock, buffer, bufferSize, 0,  
                (struct sockaddr*)&fromAddress,  
                &addressSize);  
           if (bytesRead < 0) {  
                int err = WSAGetLastError();  
                if (err == 111 /*ECONNREFUSED (Linux)*/  
                     // What a piece of crap Windows is. Sometimes  
                     // recvfrom() returns -1, but with an 'errno' of 0.  
                     // This appears not to be a real error; just treat  
                     // it as if it were a read of zero bytes, and hope  
                     // we don't have to do anything else to 'reset'  
                     // this alleged error:  
                     || err == 0 || err == EWOULDBLOCK  
                     || err == 113 /*EHOSTUNREACH (Linux)*/) {  
                     //Why does Linux return this for datagram sock?  
                     fromAddress.sin_addr.s_addr = 0;  
                     return 0;  
                }  
                socketErr("[%s] recvfrom() error: ", __FUNCTION__);  
                break;  
           }  
      } while (0);  
   
      return bytesRead;  
 }  
   
 int readSocketExact(int sock, char *buffer, unsigned bufferSize, struct sockaddr_in& fromAddress, struct timeval* timeout)  
 {  
      int bsize = bufferSize;  
      int bytesRead = 0;  
      int totBytesRead = 0;  
      do   
      {  
           bytesRead = readSocket (sock, buffer + totBytesRead, bsize, fromAddress, timeout);  
           if (bytesRead <= 0) break;  
           totBytesRead += bytesRead;  
           bsize -= bytesRead;  
      } while (bsize != 0);  
   
      return totBytesRead;  
 }  
   
 int writeSocket(int sock, char *buffer, unsigned bufferSize)  
 {  
      return send(sock, buffer, bufferSize, 0);  
 }  
   
 int writeSocket(int sock, char *buffer, unsigned int bufferSize, sockaddr_in &toAddress)  
 {  
      return sendto(sock, buffer, bufferSize, 0, (struct sockaddr *)&toAddress, sizeof(struct sockaddr_in));  
 }  
   
 bool writeSocket(int socket, struct in_addr address, unsigned short port,  
                      unsigned char* buffer, unsigned bufferSize)   
 {  
  do {  
       MAKE_SOCKADDR_IN(dest, address.s_addr, port);  
       int bytesSent = sendto(socket, (char*)buffer, bufferSize, 0, (struct sockaddr*)&dest, sizeof dest);  
       if (bytesSent != (int)bufferSize) {  
            char tmpBuf[100];  
            sprintf(tmpBuf, "writeSocket(%d), sendTo() error: wrote %d bytes instead of %u: ", socket, bytesSent, bufferSize);  
            socketErr(tmpBuf);  
            break;  
       }  
   
       return true;  
  } while (0);  
   
  return false;  
 }  
   
 bool writeSocket(int socket, struct in_addr address, unsigned short port,  
                      unsigned char ttlArg,  
                      unsigned char* buffer, unsigned bufferSize)   
 {  
      // Before sending, set the socket's TTL:  
 #if defined(__WIN32__) || defined(_WIN32)  
 #define TTL_TYPE int  
 #else  
 #define TTL_TYPE u_int8_t  
 #endif  
      TTL_TYPE ttl = (TTL_TYPE)ttlArg;  
      if (setsockopt(socket, IPPROTO_IP, IP_MULTICAST_TTL, (const char*)&ttl, sizeof ttl) < 0) {  
           socketErr("setsockopt(IP_MULTICAST_TTL) error: ");  
           return false;  
      }  
   
      return writeSocket(socket, address, port, buffer, bufferSize);  
 }  
   
 int sendRTPOverTCP(int sock, char *buffer, int len, unsigned char streamChannelId)  
 {  
      char const dollar = '$';  
   
      if (send(sock, &dollar, 1, 0) != 1) return -1;  
      if (send(sock, (char*)&streamChannelId, 1, 0) != 1) return -1;  
   
      char sz[2];  
      sz[0] = (char)((len&0xFF00)>>8);  
      sz[1] = (char)(len&0x00FF);  
      if (send(sock, sz, 2, 0) != 2) return -1;  
   
      if (send(sock, buffer, len, 0) != len) return -1;  
   
      return 0;  
 }  
   
 void shutdown(int sock)  
 {  
 #ifdef WIN32  
      if (shutdown(sock, SD_SEND) != 0)  
           socketErr("shutdown error: ");  
 #else  
      if (shutdown(sock, SHUT_RD) != 0)  
           socketErr("shutdown error: ");  
 #endif  
 }  
   
 bool isMulticastAddress(unsigned int address)  
 {  
      // Note: We return False for addresses in the range 224.0.0.0  
      // through 224.0.0.255, because these are non-routable  
      // Note: IPv4-specific #####  
      unsigned int addressInHostOrder = ntohl(address);  
      return addressInHostOrder > 0xE00000FF &&  
           addressInHostOrder <= 0xEFFFFFFF;  
 }  
   
 bool socketJoinGroupSSM(int sock, unsigned int groupAddress, unsigned int sourceFilterAddr)  
 {  
      if (!isMulticastAddress(groupAddress)) return true; // ignore this case  
   
      struct ip_mreq_source imr;  
 #ifdef ANDROID  
   imr.imr_multiaddr = groupAddress;  
   imr.imr_sourceaddr = sourceFilterAddr;  
   imr.imr_interface = INADDR_ANY;  
 #else  
      imr.imr_multiaddr.s_addr = groupAddress;  
      imr.imr_sourceaddr.s_addr = sourceFilterAddr;  
      imr.imr_interface.s_addr = INADDR_ANY;  
 #endif  
      if (setsockopt(sock, IPPROTO_IP, IP_ADD_SOURCE_MEMBERSHIP, (const char*)&imr, sizeof (struct ip_mreq_source)) < 0) {  
           socketErr("setsockopt(IP_ADD_SOURCE_MEMBERSHIP) error: ", __FUNCTION__);  
           return false;  
      }  
   
      return true;  
 }  
   
 bool socketLeaveGroupSSM(int sock, unsigned int groupAddress, unsigned int sourceFilterAddr)  
 {  
      if (!isMulticastAddress(groupAddress)) return true; // ignore this case  
   
      struct ip_mreq_source imr;  
 #ifdef ANDROID  
   imr.imr_multiaddr = groupAddress;  
   imr.imr_sourceaddr = sourceFilterAddr;  
   imr.imr_interface = INADDR_ANY;  
 #else  
      imr.imr_multiaddr.s_addr = groupAddress;  
      imr.imr_sourceaddr.s_addr = sourceFilterAddr;  
      imr.imr_interface.s_addr = INADDR_ANY;  
 #endif  
      if (setsockopt(sock, IPPROTO_IP, IP_DROP_SOURCE_MEMBERSHIP, (const char*)&imr, sizeof (struct ip_mreq_source)) < 0) {  
           return false;  
      }  
   
      return true;  
 }  
   
 bool socketJoinGroup(int sock, unsigned int groupAddress)  
 {  
      if (!isMulticastAddress(groupAddress)) return true; // ignore this case  
   
      struct ip_mreq imr;  
      imr.imr_multiaddr.s_addr = groupAddress;  
      imr.imr_interface.s_addr = INADDR_ANY;  
      if (setsockopt(sock, IPPROTO_IP, IP_ADD_MEMBERSHIP, (const char*)&imr, sizeof (struct ip_mreq)) < 0) {  
 #if defined(__WIN32__) || defined(_WIN32)  
           if (WSAGetLastError() != 0) {  
                // That piece-of-shit toy operating system (Windows) sometimes lies  
                // about setsockopt() failing!  
 #endif  
                socketErr("setsockopt(IP_ADD_MEMBERSHIP) error: ", __FUNCTION__);  
                return false;  
 #if defined(__WIN32__) || defined(_WIN32)  
           }  
 #endif  
      }  
   
      return true;  
 }  
   
 bool socketLeaveGroup(int sock, unsigned int groupAddress)  
 {  
      if (!isMulticastAddress(groupAddress)) return true; // ignore this case  
   
      struct ip_mreq imr;  
      imr.imr_multiaddr.s_addr = groupAddress;  
      imr.imr_interface.s_addr = INADDR_ANY;  
      if (setsockopt(sock, IPPROTO_IP, IP_DROP_MEMBERSHIP, (const char*)&imr, sizeof (struct ip_mreq)) < 0) {  
           return false;  
      }  
   
      return true;  
 }  
   
 typedef unsigned int     u_int32_t;  
 u_int32_t ReceivingInterfaceAddr = INADDR_ANY;  
 static bool loopbackWorks = 1;  
   
 static bool badAddressForUs(u_int32_t addr)   
 {  
      // Check for some possible erroneous addresses:  
      u_int32_t nAddr = htonl(addr);  
      return (nAddr == 0x7F000001 /* 127.0.0.1 */  
           || nAddr == 0  
           || nAddr == (u_int32_t)(~0));  
 }  
   
 u_int32_t ourIPAddress()   
 {  
      static u_int32_t ourAddress = 0;  
      int sock = -1;  
      struct in_addr testAddr;  
   
      if (ReceivingInterfaceAddr != INADDR_ANY) {  
           // Hack: If we were told to receive on a specific interface address, then   
           // define this to be our ip address:  
           ourAddress = ReceivingInterfaceAddr;  
      }  
   
      if (ourAddress == 0) {  
           // We need to find our source address  
           struct sockaddr_in fromAddr;  
           fromAddr.sin_addr.s_addr = 0;  
   
           // Get our address by sending a (0-TTL) multicast packet,  
           // receiving it, and looking at the source address used.  
           // (This is kinda bogus, but it provides the best guarantee  
           // that other nodes will think our address is the same as we do.)  
           do {  
                loopbackWorks = 0; // until we learn otherwise  
   
                testAddr.s_addr = inet_addr("228.67.43.91"); // arbitrary  
                unsigned short testPort = 15947; // ditto  
   
                sock = setupDatagramSock(testPort, false);  
                if (sock < 0) break;  
   
                if (!socketJoinGroup(sock, testAddr.s_addr)) break;  
   
                unsigned char testString[] = "hostIdTest";  
                unsigned testStringLength = sizeof testString;  
   
                if (!writeSocket(sock, testAddr, testPort, 0, testString, testStringLength))   
                     break;  
   
                // Block until the socket is readable (with a 5-second timeout):  
                fd_set rd_set;  
                FD_ZERO(&rd_set);  
                FD_SET((unsigned)sock, &rd_set);  
                const unsigned numFds = sock+1;  
                struct timeval timeout;  
                timeout.tv_sec = 5;  
                timeout.tv_usec = 0;  
                int result = select(numFds, &rd_set, NULL, NULL, &timeout);  
                if (result <= 0) break;  
   
                unsigned char readBuffer[20];  
                int bytesRead = readSocket1(sock, (char *)readBuffer, sizeof readBuffer, fromAddr);  
                if (bytesRead != (int)testStringLength  
                     || strncmp((char*)readBuffer, (char*)testString, testStringLength) != 0) {  
                          break;  
                }  
   
                // We use this packet's source address, if it's good:  
                loopbackWorks = !badAddressForUs(fromAddr.sin_addr.s_addr);  
           } while (0);  
   
           if (sock >= 0) {  
                socketLeaveGroup(sock, testAddr.s_addr);  
                closeSocket(sock);  
           }  
   
           // Make sure we have a good address:  
           u_int32_t from = fromAddr.sin_addr.s_addr;  
           if (badAddressForUs(from)) {  
                DPRINTF("This computer has an invalid IP address\n");  
                from = 0;  
           }  
   
           ourAddress = from;  
      }  
   
      return ourAddress;  
 }  
   

2015년 11월 26일 목요일

live555 TaskScheduler 수정소스

live555 의 TaskScheduler 를 사용하기 쉽게 수정한 소스이다. 윈도우/리눅스 둘 다 적용가능하며 소켓통신 프로그램 개발시 간편하게 사용할 수 있다.
Mutex/Thread 관련 소스는 http://greenday96.blogspot.com/2015/08/blog-post.html 포스트 참조

< NetCommon.h >
 #ifndef __NETCOMMON_H__  
 #define __NETCOMMON_H__  
   
 #ifdef WIN32  
   
 #include <WinSock2.h>  
 #include <ws2tcpip.h>  
   
 #define closeSocket     closesocket  
 #define EWOULDBLOCK     WSAEWOULDBLOCK  
 #define EINPROGRESS WSAEWOULDBLOCK  
 #define EINTR          WSAEINTR  
   
 #define _strcasecmp     _strnicmp  
 #define snprintf     _snprintf  
   
 #else  
   
 #include <sys/socket.h>  
 #include <netinet/in.h>  
 #include <netinet/tcp.h>  
 #include <arpa/inet.h>  
 #include <unistd.h>  
 #include <fcntl.h>  
 #include <errno.h>  
   
 #define closeSocket               close  
 #define WSAGetLastError()     errno  
   
 #include <ctype.h>  
 #include <stdlib.h>  
 #define _strcasecmp strncasecmp  
 #endif  
   
 #endif  
   


< TaskScheduler.h >
 #ifndef __TASK_SCHEDULER_H__  
 #define __TASK_SCHEDULER_H__  
   
 #include "NetCommon.h"  
 #include "Mutex.h"  
 #include "Thread.h"  
   
 #define SOCKET_READABLE  (1<<1)  
 #define SOCKET_WRITABLE  (1<<2)  
 #define SOCKET_EXCEPTION  (1<<3)  
   
 class HandlerSet;  
   
 class TaskScheduler   
 {  
 public:       
      TaskScheduler();  
      virtual ~TaskScheduler();  
   
      typedef void BackgroundHandlerProc(void* clientData, int mask);  
   
      void turnOnBackgroundReadHandling(int socketNum, BackgroundHandlerProc* handlerProc, void *clientData);  
      void turnOffBackgroundReadHandling(int socketNum);       
   
      int startEventLoop();  
      void stopEventLoop();  
      void doEventLoop();  
   
      int isRunning() { return fTaskLoop; }  
   
 protected:            
      virtual void SingleStep();  
      void taskLock();  
      void taskUnlock();  
   
 protected:  
      int                         fTaskLoop;  
      MUTEX                    fMutex;  
      THREAD                    fThread;  
   
      HandlerSet     *fReadHandlers;  
      int               fLastHandledSocketNum;  
   
      int          fMaxNumSockets;  
      fd_set     fReadSet;  
 };  
   
 class HandlerDescriptor {  
      HandlerDescriptor(HandlerDescriptor* nextHandler);  
      virtual ~HandlerDescriptor();  
   
 public:  
      int socketNum;  
      TaskScheduler::BackgroundHandlerProc* handlerProc;  
      void* clientData;  
   
 private:  
      // Descriptors are linked together in a doubly-linked list:  
      friend class HandlerSet;  
      friend class HandlerIterator;  
      HandlerDescriptor* fNextHandler;  
      HandlerDescriptor* fPrevHandler;  
 };  
   
 class HandlerSet {  
 public:  
      HandlerSet();  
      virtual ~HandlerSet();  
   
      void assignHandler(int socketNum, TaskScheduler::BackgroundHandlerProc* handlerProc, void* clientData);  
      void removeHandler(int socketNum);  
      void moveHandler(int oldSocketNum, int newSocketNum);  
   
 private:  
      HandlerDescriptor* lookupHandler(int socketNum);  
   
 private:  
      friend class HandlerIterator;  
      HandlerDescriptor fHandlers;  
 };  
   
 class HandlerIterator {  
 public:  
      HandlerIterator(HandlerSet& handlerSet);  
      virtual ~HandlerIterator();  
   
      HandlerDescriptor* next(); // returns NULL if none  
      void reset();  
   
 private:  
      HandlerSet& fOurSet;  
      HandlerDescriptor* fNextPtr;  
 };  
   
 #endif  
   


< TaskScheduler.cpp >
 #include "TaskScheduler.h"  
 #include "RTSPCommonEnv.h"  
 #include <stdio.h>  
   
 THREAD_FUNC DoEventThread(void* lpParam)  
 {  
      TaskScheduler *scheduler = (TaskScheduler *)lpParam;  
      scheduler->doEventLoop();  
      return 0;  
 }  
   
 TaskScheduler::TaskScheduler()  
 {  
      fTaskLoop = 0;  
      MUTEX_INIT(&fMutex);  
      FD_ZERO(&fReadSet);  
      fMaxNumSockets = 0;  
      fThread = NULL;  
      fReadHandlers = new HandlerSet();  
 }  
   
 TaskScheduler::~TaskScheduler()  
 {  
      stopEventLoop();  
   
      delete fReadHandlers;  
   
      THREAD_DESTROY(&fThread);  
   
      MUTEX_DESTROY(&fMutex);  
 }  
   
 void TaskScheduler::taskLock()  
 {  
      MUTEX_LOCK(&fMutex);  
 }  
   
 void TaskScheduler::taskUnlock()  
 {  
      MUTEX_UNLOCK(&fMutex);  
 }  
   
 void TaskScheduler::turnOnBackgroundReadHandling(int socketNum, BackgroundHandlerProc* handlerProc, void *clientData)   
 {  
	taskLock();

	if (socketNum < 0) goto exit;

	FD_SET((unsigned)socketNum, &fReadSet);
	fReadHandlers->assignHandler(socketNum, handlerProc, clientData);

	if (socketNum+1 > fMaxNumSockets) {
		fMaxNumSockets = socketNum+1;
	}

exit:
	taskUnlock();
 }  
   
 void TaskScheduler::turnOffBackgroundReadHandling(int socketNum)   
 {  
	taskLock();

	if (socketNum < 0) goto exit;

	FD_CLR((unsigned)socketNum, &fReadSet);
	fReadHandlers->removeHandler(socketNum);

    HandlerIterator iter(*fReadHandlers);
	HandlerDescriptor* handler;

	int maxSocketNum = 0;
	while ((handler = iter.next()) != NULL) {
		if (handler->socketNum+1 > maxSocketNum) {
			maxSocketNum = handler->socketNum+1;
		}
	}
	fMaxNumSockets = maxSocketNum;	

exit:
	taskUnlock();
 }  
   
 int TaskScheduler::startEventLoop()  
 {  
      if (fTaskLoop != 0)  
           return -1;  
   
      fTaskLoop = 1;  
      THREAD_CREATE(&fThread, DoEventThread, this);  
      if (!fThread) {  
           DPRINTF("failed to create event loop thread\n");  
           fTaskLoop = 0;  
           return -1;  
      }  
   
      return 0;  
 }  
   
 void TaskScheduler::stopEventLoop()  
 {  
      fTaskLoop = 0;  
   
      THREAD_JOIN(&fThread);  
      THREAD_DESTROY(&fThread);  
 }  
   
 void TaskScheduler::doEventLoop()   
 {  
      while (fTaskLoop)  
      {  
           SingleStep();  
      }  
 }  
   
 void TaskScheduler::SingleStep()  
 {  
	taskLock();

	fd_set readSet = fReadSet;

	struct timeval timeout;
	timeout.tv_sec = 1;
	timeout.tv_usec = 0;

	int selectResult = select(fMaxNumSockets, &readSet, NULL, NULL, &timeout);
	if (selectResult < 0) {
		int err = WSAGetLastError();
		DPRINTF("TaskScheduler::SingleStep(): select() fails : %d\n", err);
		taskUnlock();
		return;
	}

	HandlerIterator iter(*fReadHandlers);
	HandlerDescriptor* handler;

	while ((handler = iter.next()) != NULL) {
		if (FD_ISSET(handler->socketNum, &readSet) && handler->handlerProc != NULL) {
			(*handler->handlerProc)(handler->clientData, SOCKET_READABLE);
		}
	}

	taskUnlock();
 }  
   
   
 HandlerDescriptor::HandlerDescriptor(HandlerDescriptor* nextHandler)  
 : handlerProc(NULL) {  
      // Link this descriptor into a doubly-linked list:  
      if (nextHandler == this) { // initialization  
           fNextHandler = fPrevHandler = this;  
      } else {  
           fNextHandler = nextHandler;  
           fPrevHandler = nextHandler->fPrevHandler;  
           nextHandler->fPrevHandler = this;  
           fPrevHandler->fNextHandler = this;  
      }  
 }  
   
 HandlerDescriptor::~HandlerDescriptor() {  
      // Unlink this descriptor from a doubly-linked list:  
      fNextHandler->fPrevHandler = fPrevHandler;  
      fPrevHandler->fNextHandler = fNextHandler;  
 }  
   
 HandlerSet::HandlerSet()  
 : fHandlers(&fHandlers) {  
      fHandlers.socketNum = -1; // shouldn't ever get looked at, but in case...  
 }  
   
 HandlerSet::~HandlerSet() {  
      // Delete each handler descriptor:  
      while (fHandlers.fNextHandler != &fHandlers) {  
           delete fHandlers.fNextHandler; // changes fHandlers->fNextHandler  
      }  
 }  
   
 void HandlerSet  
 ::assignHandler(int socketNum, TaskScheduler::BackgroundHandlerProc* handlerProc, void* clientData) {  
      // First, see if there's already a handler for this socket:  
      HandlerDescriptor* handler = lookupHandler(socketNum);  
      if (handler == NULL) { // No existing handler, so create a new descr:  
           handler = new HandlerDescriptor(fHandlers.fNextHandler);  
           handler->socketNum = socketNum;  
      }  
   
      handler->handlerProc = handlerProc;  
      handler->clientData = clientData;  
 }  
   
 void HandlerSet::removeHandler(int socketNum) {  
      HandlerDescriptor* handler = lookupHandler(socketNum);  
      delete handler;  
 }  
   
 void HandlerSet::moveHandler(int oldSocketNum, int newSocketNum) {  
      HandlerDescriptor* handler = lookupHandler(oldSocketNum);  
      if (handler != NULL) {  
           handler->socketNum = newSocketNum;  
      }  
 }  
   
 HandlerDescriptor* HandlerSet::lookupHandler(int socketNum) {  
      HandlerDescriptor* handler;  
      HandlerIterator iter(*this);  
      while ((handler = iter.next()) != NULL) {  
           if (handler->socketNum == socketNum) break;  
      }  
      return handler;  
 }  
   
 HandlerIterator::HandlerIterator(HandlerSet& handlerSet)  
 : fOurSet(handlerSet) {  
      reset();  
 }  
   
 HandlerIterator::~HandlerIterator() {  
 }  
   
 void HandlerIterator::reset() {  
      fNextPtr = fOurSet.fHandlers.fNextHandler;  
 }  
   
 HandlerDescriptor* HandlerIterator::next() {  
      HandlerDescriptor* result = fNextPtr;  
      if (result == &fOurSet.fHandlers) { // no more  
           result = NULL;  
      } else {  
           fNextPtr = fNextPtr->fNextHandler;  
      }  
   
      return result;  
 }  
   


< TaskScheduler 사용 >

TaskScheduler* fTask = new TaskScheduler();
...
fTask->turnOnBackgroundReadHandling(sock, &incomingHandler, this); // 소켓 핸들러 등록
...
fTask->turnOffBackgroundReadHandling(sock); // 소켓 핸들러 해제
delete fTask;

static void incomingHandler(void *data, int)
{
    recvfrom();
    ...
}

2012년 6월 5일 화요일

TCP 소켓 Connection 타임아웃 주기(윈도우/리눅스 공통) - TCP Socket Connection Timeout


TCP 소켓통신에서 서버로 연결할때 상대서버가 동작하지 않거나 연결에 문제가 생겼을때
connect 함수에서 블러킹이 걸려 한참 동안 빠져나오지 못할때가 있다.
이런 현상을 방지하려면 소켓을 non-blocking 으로 만든다음 타임아웃을 주고 connect 후
select 함수에서 연결결과를 감지하면 된다.

 #ifdef WIN32  
   
 #include <WinSock2.h>  
 #include <ws2tcpip.h>  
   
 #define closeSocket     closesocket  
 #define EWOULDBLOCK     WSAEWOULDBLOCK  
 #define EINPROGRESS WSAEWOULDBLOCK  
 #define EINTR          WSAEINTR  
     
 #else  
   
 #include <sys/socket.h>  
 #include <netinet/in.h>  
 #include <netinet/tcp.h>  
 #include <arpa/inet.h>  
 #include <unistd.h>  
 #include <fcntl.h>  
 #include <errno.h>  
   
 #define closeSocket               close  
 #define WSAGetLastError()     errno  
   
 #include <ctype.h>  
 #include <stdlib.h>  
 #endif  
   
 #ifdef WIN32  
 #define WS_VERSION_CHOICE1 0x202/*MAKEWORD(2,2)*/  
 #define WS_VERSION_CHOICE2 0x101/*MAKEWORD(1,1)*/  
 int initializeWinsockIfNecessary(void) {  
      /* We need to call an initialization routine before  
      * we can do anything with winsock. (How fucking lame!):  
      */  
      static int _haveInitializedWinsock = 0;  
      WSADATA     wsadata;  
   
      if (!_haveInitializedWinsock) {  
           if ((WSAStartup(WS_VERSION_CHOICE1, &wsadata) != 0)  
                && ((WSAStartup(WS_VERSION_CHOICE2, &wsadata)) != 0)) {  
                     return 0; /* error in initialization */  
           }  
           if ((wsadata.wVersion != WS_VERSION_CHOICE1)  
                && (wsadata.wVersion != WS_VERSION_CHOICE2)) {  
                     WSACleanup();  
                     return 0; /* desired Winsock version was not available */  
           }  
           _haveInitializedWinsock = 1;  
      }  
   
      return 1;  
 }  
 #else  
 #define initializeWinsockIfNecessary()     1  
 #endif  
    
 int makeSocketNonBlocking(int sock)  
 {  
 #ifdef WIN32  
      unsigned long arg = 1;  
      return ioctlsocket(sock, FIONBIO, &arg) == 0;  
 #else  
      int curFlags = fcntl(sock, F_GETFL, 0);  
      return fcntl(sock, F_SETFL, curFlags|O_NONBLOCK) >= 0;  
 #endif  
 }  
   
 static int reuseFlag = 1;  
   
 int setupStreamSock(short port, int makeNonBlocking)  
 {  
      if (!initializeWinsockIfNecessary()) {  
           socketErr("[%s] Failed to initialize 'winsock': ", __FUNCTION__);  
           return -1;  
      }  
   
      int newSocket = socket(AF_INET, SOCK_STREAM, 0);  
      if (newSocket < 0) {  
           DPRINTF("%s:%d\n",__FUNCTION__,__LINE__);  
           socketErr("[%s] unable to create stream socket: ", __FUNCTION__);  
           return newSocket;  
      }  
  
      if (setsockopt(newSocket, SOL_SOCKET, SO_REUSEADDR,  
           (const char*)&reuseFlag, sizeof reuseFlag) != 0) {  
                socketErr("[%s] setsockopt(SO_REUSEADDR) error: ", __FUNCTION__);  
                closeSocket(newSocket);  
                return -1;  
      }  
 
      struct sockaddr_in c_addr;  
      memset(&c_addr, 0, sizeof(c_addr));  
      c_addr.sin_addr.s_addr = INADDR_ANY;  
      c_addr.sin_family = AF_INET;  
      c_addr.sin_port = htons(port);  
   
      if (bind(newSocket, (struct sockaddr*)&c_addr, sizeof c_addr) != 0) {  
           socketErr("[%s] bind() error (port number: %d): ", __FUNCTION__, port);  
           closeSocket(newSocket);  
           return -1;  
      }  
   
      if (makeNonBlocking) {  
           if (!makeSocketNonBlocking(newSocket)) {  
                socketErr("[%s] failed to make non-blocking: ", __FUNCTION__);  
                closeSocket(newSocket);  
                return -1;  
           }  
      }  
   
      return newSocket;  
 }  
   
 int connectToServer(char *svrIP, int svrPort)  
 {  
      int sock = setupStreamSocket(0, 1);  
      int ret, err;  
   
      struct sockaddr_in svr_addr;  
      memset(&svr_addr, 0, sizeof(svr_addr));  
      svr_addr.sin_addr.s_addr = inet_addr(svrIP);  
      svr_addr.sin_family = AF_INET;  
      svr_addr.sin_port = htons(svrPort);  
   
      fd_set set;  
      FD_ZERO(&set);  
      timeval tvout = {2, 0};  // 2 seconds timeout  
   
      FD_SET(sock, &set);  
   
      if ((ret=connect(sock, (struct sockaddr *)&svr_addr, sizeof(svr_addr))) != 0) {  
           err = WSAGetLastError();  
           if (err != EINPROGRESS && err != EWOULDBLOCK) {  
                printf("connect() failed : %d\n", err);  
                return -1;  
           }  
           if (select(sock+1, NULL, &set, NULL, &tvout) <= 0) {  
                printf("select/connect() failed : %d\n", WSAGetLastError());  
                return -1;  
           } 
           err = 0;
	   socklen_t len = sizeof(err);
	   if (getsockopt(sock, SOL_SOCKET, SO_ERROR, (char*)&err, &len) < 0 || err != 0 ) {
		printf("getsockopt() error: %d\n", err);
		return -1;
	   }
      }  
      return 0;  
 }  

c# 소스코드

2012년 3월 2일 금요일

자기 ip 알아내기 - obtaining host ip address

struct hostent *fHost;
char fMyIpAddr[64];

void getMyIP()
{
        char buffer[1024];

if (gethostname(buffer, sizeof(buffer)) == SOCKET_ERROR) {
printf("%s gethostname error !!!\r\n", __FUNCTION__);
return;
}

fHost = gethostbyname(buffer);
if (fHost == NULL) {
printf("%s gethostbyname error !!!\r\n", __FUNCTION__);
return;
}

sprintf(fMyIpAddr, "%d.%d.%d.%d", 
((struct in_addr *)(fHost->h_addr))->S_un.S_un_b.s_b1,
((struct in_addr *)(fHost->h_addr))->S_un.S_un_b.s_b2,
((struct in_addr *)(fHost->h_addr))->S_un.S_un_b.s_b3,
((struct in_addr *)(fHost->h_addr))->S_un.S_un_b.s_b4
);

}

2012년 2월 16일 목요일

win32 shutdown 함수 사용법 - how to use win32 shutdown function

shutdown 함수는 tcp 소켓통신에서 송신/수신 채널만을 닫을때 쓰는 함수이다.
select 함수를 쓰는 구조에서 소켓통신을 종료할때 closesocket(close) 을 바로 호출하면
select 함수에서 오류가 발생하고 소켓 전체적으로 문제가 발생한다. 이때 shutdown 함수를
호출해서 통신하는 상대방에게 FIN 패킷을 전달하여 상대방 recv 함수에서 0를 리턴하게
하여 상대방쪽에서 closesocket(close) 시켜 정상적으로 tcp 소켓을 닫을때 주로 사용한다.
이때 MS 의 페이크가 있는데 리눅스와 윈도우의 shutdown 함수의 인수는 서로 반대의 의미
를 가진다. 즉,
linux                                   windows
shutdown(SHUT_RD)  => shutdown(SD_SEND)
shutdown(SHUT_WR) => shutdown(SD_RECEIVE)

tcp 통신중 상대방에게 FIN 메시지를 전달해서 상대방으로 하여금 closesocket(close) 함수를
호출하도록 유도하려면 shutdown(SD_SEND) 함수를 호출

2011년 10월 18일 화요일

윈도우에서 IOCP 를 사용해야하는 이유

네트웍 프로그램에서 대량의 트래픽이나 클라이언트/서버를 상대로 데이터를 송수신할때 보통 단일쓰레드로 select 함수를 사용해서 처리한다.
리눅스에서는 이것이 보통 합당한 방식인데 왜냐하면 멀티코어 cpu 상에서 cpu 사용률을 보면 각 cpu가 골고루 사용되는것을 볼수 있다. 이것은 이더넷 카드의 인터럽트가 각 cpu 에 골고루 신호를 보내기때문이다.
하지만 윈도우에서 같은 방식으로 싱글쓰레드 select 함수를 사용하면 cpu 1개가 100% 로 치솟고 나머지 cpu 들은 거의 0% 로 놀고있는 현상을 볼수있다. 이것은 이더넷 카드의 인터럽트가 쿼드코어기준으로 4개의 cpu 에 골고루 인터럽트를 보내지않고 1개의 cpu 에게만 집중적으로 보내기 때문이다. 이것을 막기위해 레지스트리 수정등 여러방법을 써봤지만 다 효과가 없었고 오직 IOCP 기술을 사용해서 부하분산을 할 수 있었다.
IO Completion Port 는 일종의 IO 큐로써 등록된 각 소켓 디스크립터들의 이벤트들을 담고있고 초기에 사용자가 만들었던 쓰레드풀에서 각 쓰레드들이 랜덤하게 큐의 이벤트들을 pop 해서 가져가서 처리하는 방식이다.
IO Completion Port 가 일종의 select 함수 역할을 하는 것이고 이벤트에 대한 처리는 각 쓰레드들이 알아서 하는 방식이다. 이렇게 하면 쓰레드들이 골고루 일을 하기때문에 cpu 부하분산을 시킬수 있다.

2011년 10월 14일 금요일

IOCP로 UDP 데이터 수신 - how to use IOCP WSARecvFrom


< 변수들 - variables >

// IOCP 컨텍스트 정의
class IOContext {
WSAOVERLAPPED wsaOverlapped;
int socket;
WSABUF wsaBuf;
char buf[10240];
HANDLE hCompletionPort;
    HANDLE hCloseEvent;
    BOOL bClosed;

    IOContext(int sock);
    ~IOContext();

};


IOContext::IOContext(int sock)
{
socket = sock;
wsaBuf.len = sizeof(buf);
wsaBuf.buf = buf;
SecureZeroMemory((PVOID)&wsaOverlapped, sizeof(WSAOVERLAPPED));

hCloseEvent = CreateEvent(NULL, FALSE, FALSE, NULL);
bClosed = FALSE;
}

IOContext::~IOContext()
{
CloseHandle(hCloseEvent);
}


List<IOContext> contextList;    // IOCP 컨텍스트들 링크드 리스트로 관리
HANDLE hCompletionPort;
SYSTEM_INFO systemInfo;
DWORD dwThreadCount;
HANDLE hThreads[16];


< IOCP 초기화 및 쓰레드 생성 - initialize IOCP Handle and create worker threads >

void init()
{

    hCompletionPort = CreateIoCompletionPort(INVALID_HANDLE_VALUE, NULL, 0, 0);
    GetSystemInfo(&systemInfo);
    dwThreadCount = systemInfo.dwNumberOfProcessors*2;
 
    start();
}

void start()
{
    for (int i=0; i < dwThreadCount; i++)
    {
  hThreads[i] = (HANDLE)_beginthreadex(NULL, 0, WorkerThread, this, 0, NULL);
    }
}

< IOCP 에 소켓등록 - register socket to IOCP Handle >

void registerSocket(int socketNum)
{
     IOContext *context = new IOContext(socketNum);

     contextList.insert(context);
     context->hCompletionPort = CreateIoCompletionPort((HANDLE)socketNum, hCompletionPort, (ULONG_PTR)context, 0);

int err;
if (!PostQueuedCompletionStatus(hCompletionPort, 0, (ULONG_PTR)context, &context->wsaOverlapped)) {
  if ((err = WSAGetLastError()) != WSA_IO_PENDING)
    printf("[%s] PostQueuedCompletionStatus error: %d\r\n", err);
}
}


< IOCP 에 소켓등록해제 - unregister socket from IOCP Handle >

void unregisterSocket(int socketNum)
{
     IOContext *context = (IOContext *)contextList.search(socketNum);
     closesocket(socketNum);
     if (WaitForSingleObject(context->hCloseEvent, 1000*5) == WAIT_TIMEOUT) {
         printf("CloseEvent Wait Timeout !!!\n");
     }
     contextList.remove(context);
}


< Worker Thread >

unsigned int __stdcall WorkerThread(LPVOID lpParam)
{
    BOOL bSuccess = FALSE;
    DWORD dwIoSize;
    LPOVERLAPPED lpOverlapped = NULL;
    IOContext *context;

    while (1)
    {
        bSuccess = GetQueuedCompletionStatus(hCompletionPort, &dwIoSize, (PULONG_PTR)&context, &lpOverlapped, INFINITE);

if (context->bClosed) continue;

  if (!bSuccess) {
    err = GetLastError();
    if (err == WSA_OPERATION_ABORTED) {
                context->bClosed = TRUE;
SetEvent(context->hCloseEvent);         
            } else if (err == WSAENOTSOCK) {
                context->bClosed = TRUE;
                SetEvent(context->hCloseEvent);
            }
            continue;
        }

        // DO SOMETHING
        ......
  struct sockaddr_in fromAddress;
  int len = sizeof(context->buf);
  int addressSize = sizeof(fromAddress);
  int bytesRead;
  DWORD flag = 0;
  int err;

  context->wsaBuf.len = sizeof(buf);
  context->wsaBuf.buf = context->buf;
  if (WSARecvFrom(context->socket, &context->wsaBuf, 1, (LPDWORD)&bytesRead, (LPDWORD)&flag, (struct sockaddr*)&fromAddress, (socklen_t *)&addressSize, &context->wsaOverlapped, NULL) == SOCKET_ERROR) {
    if ((err=WSAGetLastError()) != WSA_IO_PENDING) {
      printf("[%s] WSARecvFrom error:%d, sock:%d, bytesRead:%d\r\n", __FUNCTION__, err, context->socket, bytesRead);
    }

            if (err == WSAENOTSOCK) { // invalid socket (Socket operation on nonsocket.)
                context->bClosed = TRUE;
                SetEvent(context->hCloseEvent);
            }
  }
    }
    return 0;
}


2011년 9월 30일 금요일

윈도우에서 writev 함수구현 - windows writev function implementation


#include <winsock2.h>

int writev(int sock, struct iovec *iov, int nvecs)
{
DWORD ret;
if (WSASend(sock, (LPWSABUF)iov, nvecs, &ret, 0, NULL, NULL) == 0) {
return ret;
}
return -1;
}

윈도우/리눅스 non-blocking socket 만들기 - making windows/linux non-blocking socket


Boolean makeSocketNonBlocking(int sock) {
#if defined(WIN32) || defined(_WIN32) || defined(IMN_PIM)
  unsigned long arg = 1;
  return ioctlsocket(sock, FIONBIO, &arg) == 0;
#elif defined(VXWORKS)
  int arg = 1;
  return ioctl(sock, FIONBIO, (int)&arg) == 0;
#else
  int curFlags = fcntl(sock, F_GETFL, 0);
  return fcntl(sock, F_SETFL, curFlags|O_NONBLOCK) >= 0;
#endif
}

윈도우는 ioctlsocket, 리눅스는 fcntl 사용

2011년 9월 29일 목요일

윈도우 소켓버퍼 사이즈 늘리기 - increasing windows socket buffer size

< 현재 소켓 버퍼 사이즈 읽어오기 >

unsigned getBufferSize(int bufOptName,
     int socket) {
  unsigned curSize;
  SOCKLEN_T sizeSize = sizeof curSize;
  if (getsockopt(socket, SOL_SOCKET, bufOptName,
(char*)&curSize, &sizeSize) < 0) {
    socketErr("getBufferSize() error: ");
    return 0;
  }

  return curSize;
}

< 소켓 버퍼 사이즈 늘리기 설정 >

#ifndef SOCKLEN_T
#define SOCKLEN_T socklen_t
#endif

unsigned increaseBufferTo( int bufOptName,
int socket, unsigned requestedSize) {
  // First, get the current buffer size.  If it's already at least
  // as big as what we're requesting, do nothing.
  unsigned curSize = getBufferSize(bufOptName, socket);

  // Next, try to increase the buffer to the requested size,
  // or to some smaller size, if that's not possible:
  while (requestedSize > curSize) {
    SOCKLEN_T sizeSize = sizeof requestedSize;
    if (setsockopt(socket, SOL_SOCKET, bufOptName,
  (char*)&requestedSize, sizeSize) >= 0) {
      // success
      return requestedSize;
    }
    requestedSize = (requestedSize+curSize)/2;
  }

  return getBufferSize( bufOptName, socket);
}

< 함수 사용 >

// 송신버퍼 사이즈 늘리기
unsigned increaseSendBufferTo(int socket, unsigned requestedSize) {
  return increaseBufferTo( SO_SNDBUF, socket, requestedSize);
}

// 수신버퍼 사이즈 늘리기
unsigned increaseReceiveBufferTo( int socket, unsigned requestedSize) {
  return increaseBufferTo( SO_RCVBUF, socket, requestedSize);
}

< 함수 사용 >

int ret, size = 1024*1024;

if ((ret=increaseSendBufferTo(fsock, size)) != size)
printf("%s failed to increase send buffer size (size:%d, ret:%d)\r\n", __FUNCTION__, size, ret);


*** 버퍼 사이즈를 원하는 만큼 늘릴 수 없으면 윈도우 레지스트리 설정을 추가/수정해야한다.

[HKEY_LOCAL_MACHINE \SYSTEM \CurrentControlSet \Services \Afd \Parameters]
DefaultReceiveWindow = 16384
DefaultSendWindow = 16384

* 리부팅할 필요없고 위 레지스트리 값만 넣어주면된다.






winsock socket 연결 64개 이상으로 늘리는 법 - increasing winsock 64 sockets limitation

윈속 소켓 연결시 한 쓰레드에 64개까지만 연결을 제한한다. (서버측 accept 혹은 select 할 소켓 디스크립터 갯수)
이것은 winsock2.h 에 FD_SET 이 64 로 정의되어있기 때문인데 FD_SET 을 아래 코드와 같이 정의해서 연결을 늘려줄 수 있다.


#define FD_SETSIZE 1024
#include <winsock2.h>

< 출처 : http://tangentsoft.net/wskfaq/advanced.html >

2011년 9월 26일 월요일

vc++에서 winsock2 헤더파일 include 및 link(링크)

< winsock2 include >

#ifdef WIN32
#include <winsock2.h>
#include <ws2tcpip.h>
#endif

빌드에러가 발생하면 프로젝트 속성에서 C/C++ -> Preprocessor -> Preprocessor Definitions 에
_WINSOCKAPI_ 추가


< winsock2 library 링크 >

프로젝트 속성 -> Librarian -> General -> Additional Dependencies 에
ws2_32.lib 추가

WSAStartup() 함수로 winsock 초기화


윈도우에서 winsock 사용시 WSAStartup 함수로 winsock 을 초기화해야한다.


#if defined(__WIN32__) || defined(_WIN32)
#ifndef IMN_PIM
#define WS_VERSION_CHOICE1 0x202/*MAKEWORD(2,2)*/
#define WS_VERSION_CHOICE2 0x101/*MAKEWORD(1,1)*/
int initializeWinsockIfNecessary(void) {
/* We need to call an initialization routine before
* we can do anything with winsock.  (How fucking lame!):
*/
static int _haveInitializedWinsock = 0;
WSADATA wsadata;

if (!_haveInitializedWinsock) {
if ((WSAStartup(WS_VERSION_CHOICE1, &wsadata) != 0)
   && ((WSAStartup(WS_VERSION_CHOICE2, &wsadata)) != 0)) {
return 0; /* error in initialization */
}
    if ((wsadata.wVersion != WS_VERSION_CHOICE1)
       && (wsadata.wVersion != WS_VERSION_CHOICE2)) {
        WSACleanup();
return 0; /* desired Winsock version was not available */
}
_haveInitializedWinsock = 1;
}

return 1;
}
#else
int initializeWinsockIfNecessary(void) { return 1; }
#endif
#else
#define initializeWinsockIfNecessary() 1
#endif

...

// 윈속 초기화

  if (!initializeWinsockIfNecessary()) {
    socketErr(env, "Failed to initialize 'winsock': ");
    return -1;
  }


windows 에서 winsock select 함수사용


winsock select 함수 사용시 fd_set 이 비어있으면 select 함수에서 에러발생
이때 WSAGetLastError() 함수로 에러값을 보고 dummy 소켓을 생성해서 fd_set에
넣어줘야함. 아래 소스 참조


  fd_set readSet = fReadSet; // make a copy for this select() call

  struct timeval tv_timeToDelay;
  tv_timeToDelay.tv_sec = 1;
  tv_timeToDelay.tv_usec = 0;


  int selectResult = select(fMaxNumSockets, &readSet, NULL, NULL,
   &tv_timeToDelay);
  if (selectResult < 0) {
#if defined(__WIN32__) || defined(_WIN32)
    int err = WSAGetLastError();
    // For some unknown reason, select() in Windoze sometimes fails with WSAEINVAL if
    // it was called with no entries set in "readSet".  If this happens, ignore it:
    if (err == WSAEINVAL && readSet.fd_count == 0) {
      err = 0;
      // To stop this from happening again, create a dummy readable socket:
      int dummySocketNum = socket(AF_INET, SOCK_DGRAM, 0);
      FD_SET((unsigned)dummySocketNum, &fReadSet);
    }
    if (err != 0) {
#else
    if (errno != EINTR && errno != EAGAIN) {
#endif
// Unexpected error - treat this as fatal:
#if !defined(_WIN32_WCE)
perror("BasicTaskScheduler::SingleStep(): select() fails");
#endif
// exit(0);
      }
  }