ラベル Boost の投稿を表示しています。 すべての投稿を表示
ラベル Boost の投稿を表示しています。 すべての投稿を表示

2011年11月14日月曜日

BOOST::ASIOを使ってみる(5)

非同期処理であるASIOをマルチスレッド上で動かしてみます。
その前に前回までのソースコードを少し変更しています。簡単に関数オブジェクトを実装できるboost::format<T>ですが、前回まではSessionのメンバ変数の値として持っていましたが、ServerModuleが先に解放された場合、クリアできなくなるので参照として持つようにしました。その代わりServerModuleのメンバ変数として存在するようにしました。値型と参照やスコープなどを考えないといけないところが私にとってC++が難しく感じるところかもしれません。
boost::format<T>などの使い方も勉強していきたいと思います。Lamdaもあって便利だと思いますので。

今回、boost::thread_groupを使いたいと思います。Threadをグループとして扱うことができます。boost::io_service::ioはハンドラーを効率よくスレッドに振り分けるようで、何も考えずにスレッドの中でio.run()を実行します。スレッドを開始するクラスをtemplateを使って少し汎用的にして別ファイルにしました。

前回少し説明したio_service::workクラスを使っています。今回、マルチスレッドの各スレッドでio.run()を実行してもキューが空だと終了してしまうからです。io_service::workを代入して空ではない状態にします。逆に終了するときはio_service::workを破棄します。

以下、ソースコードです。

ModuleThread.hpp
#include <iostream>
#include <boost/asio.hpp>
#include <boost/bind.hpp>
#include <boost/thread.hpp>
#include <boost/shared_ptr.hpp>

namespace thorny_road{

// T is necessary 'stop', 'constructor(io,short port)'
template<typename T>
class ModuleThread : private boost::noncopyable
{
private:
    boost::thread_group tg ;
    boost::shared_ptr<T> _module ; // Target Module
    boost::shared_ptr<boost::asio::io_service> io_ptr ; // for creating service io
    boost::asio::io_service& io ; // service IO
    boost::shared_ptr<boost::asio::io_service::work> work;
    short port ;
public:
    ModuleThread(short Port, int NoThreads)
        : port(Port),
        io_ptr(new boost::asio::io_service()),
        io(*io_ptr),
        work(new boost::asio::io_service::work(io)) // add waiting work in io loop
    {
        io.reset() ;
        // Create some threads. and then run io.
        for (int i = 0 ; i < NoThreads; i++)
            tg.create_thread( boost::bind(&boost::asio::io_service::run, &io) ) ;

        // Create module
        _module.reset(new T(io,Port)) ;
    }

    // Destructor. reset IO work and waiting by stopping threads
    ~ModuleThread()
    {
        stop() ;
    }

    void stop()
    {
        work.reset() ; // destroy work object which is the waiting object.
        tg.join_all() ; // waiting by stopping threads.

        _module.reset() ;
    }

    T& getModule() { return *_module ; }
    T& getModule() const { return static_cast<const T&>(getModule()); }
} ;

} // thorny_road

#endif  //__ASIO_MODULE_THREAD__

ServerModule.hpp

#pragma once

#ifndef __ASIO_SERVER_MODULE__
#define __ASIO_SERVER_MODULE__

#include <iostream>
#include <boost/asio.hpp>
#include <boost/bind.hpp>
#include <boost/function.hpp>
#include <boost/shared_ptr.hpp>

namespace thorny_road{

using namespace boost::asio ;

// For calling event function handler easily.
typedef boost::function<void ()> SimpleProcedure ;

class ServerSession : private boost::noncopyable
{
private:
    io_service& io ;
    ip::tcp::socket socket ;
    boost::asio::streambuf buf ;
    SimpleProcedure& acc_stop ;
public:
    ServerSession(io_service& io,SimpleProcedure& acc_stop)
        : io(io),socket(io),acc_stop(acc_stop)    {}

    ~ServerSession(){}

    void start()
    {
        // Enqueue Read Handler
        async_read_until(
            socket,
            buf,
            '\n',
            boost::bind(&ServerSession::read_ok,this,_1) ) ;
    }

private:
    void read_ok(const boost::system::error_code e)
    {
        if (!e)
        {
            std::iostream ios(&buf) ;

            std::string tmp ;
            ios >> tmp ; // get input stream
            std::cout << tmp ;
            if (tmp == "end") { delete this; return; }
            else if (tmp == "bye")
            {
                acc_stop() ;
                delete this;
                return;
            }

            ios << tmp <<std::endl ; // retrun as it is

            // Enqueue Write Handler
            async_write(
                socket,
                buf,
                boost::bind(&ServerSession::write_ok,this,_1) ) ;
        }
        else if (e == boost::asio::error::operation_aborted)
        {
            // abort
            std::cout << "connection abort" ;
            delete this;
            return ;
        }
        else
        {
            std::cout << "error" ;
            delete this;
            return ;
        }
    }

    void write_ok(const boost::system::error_code e)
    {
        if (!e)
        {
            start() ; // Restart
        }
        else if (e == boost::asio::error::operation_aborted)
        {
            // abort
            std::cout << "connection abort" ;
            delete this;
            return ;
        }
        else
        {
            std::cout << "error" ;
            delete this;
            return ;
        }
    }

public:
    ip::tcp::socket& getSocket()
    {
        return (socket) ;
    }

} ;

class ServerModule : private boost::noncopyable
{
private:
    io_service& io ;   
    ip::tcp::acceptor accept ;
    ServerSession* session ;
    SimpleProcedure stopEvent ;
public:
    ServerModule(io_service& io, const short port)
        :io(io),accept(io,ip::tcp::endpoint(ip::tcp::v4(),port))
    {
        // Initialize Function Object
        stopEvent = boost::bind(&ServerModule::stop,this) ;
        start_accept() ;
    }

    ~ServerModule() {}

    void start_accept()
    {
        session = new ServerSession(io,stopEvent) ;
        // Enqueue Accept Handler
        accept.async_accept(
            session->getSocket(),
            boost::bind(&ServerModule::accept_ok,this,_1) ) ;
    }

    // Accept service will be closed. Then IOServce throw operation_aborted.
    // After queue get be empty, the io.run() will stop.
    void stop()
    {
        accept.close() ;
    }
private:

    void accept_ok(const boost::system::error_code e)
    {
        if (!e)
        {
            session->start() ;
            start_accept() ;
        }
        else if (e == boost::asio::error::operation_aborted)
        {
            // abort
            std::cout << "connection abort" ;
            delete session ;
            session = NULL ;
            return ;
        }
        else
        {
            // error
            std::cout << "error" ;
            delete session ;
            session = NULL ;
            return ;
        }
    }

} ;

} // namespace thorny_road
#endif //__ASIO_SERVER_MODULE__
Main

#include "stdafx.h"
#include "ServerModule.hpp"
#include "ModuleThread.hpp" // for threading

using namespace thorny_road ;

int _tmain(int argc, _TCHAR* argv[])
{
    ModuleThread<ServerModule> module(2085,2) ; // create module.

    return 0;
}
ModuleThreadの引数でポート番号とスレッドの数を指定します。"bye"と入力があると、ハンドラーが破棄され、アプリケーションが終了します。

次回は送信側を勉強してみます。

ソースコードは自由にご使用ください。ただし問題が起きても責任はとれません。また、ソースコードに対する著作権は放棄していません。

2011年11月12日土曜日

BOOST::ASIOを使ってみる(4)

今回はio_service.run()の挙動を簡単に確認してみたいと思います。run()はキューに入っているハンドラーが無くなるまでループしています。前回のソースコードでは接続ハンドラーaccept_ok(...)の処理中に次の接続ハンドラーを登録していますので、常にキューにハンドラーが存在している状態になっています。

ちょっとした実験を兼ねて、"bye"と入力すると、接続ハンドラーが破棄されるようなものを作ってみます。

デリゲートのような関数オブジェクトを簡単に実装できるものとしてboost::functionがあります。使い方は簡単でboost::function<戻り値型 (引数型)> というように型を宣言します。

今回は単純な戻り値・引数が無い関数オブジェクトを使うので

typedef boost::function<void ()> SimpleProcedure ;


という風にしています。

以下、ソースコードです。

ServerModule.hpp



#pragma once

#ifndef __ASIO_SERVER_MODULE__
#define __ASIO_SERVER_MODULE__

#include <iostream>
#include <boost/asio.hpp>
#include <boost/bind.hpp>
#include <boost/function.hpp>

namespace thorny_road{

using namespace boost::asio ;

// For calling event function handler easily.
typedef boost::function<void ()> SimpleProcedure ;

class ServerSession : private boost::noncopyable
{
private:
    io_service& io ;
    ip::tcp::socket socket ;
    boost::asio::streambuf buf ;
    SimpleProcedure acc_stop ;
public:
    ServerSession(io_service& io,SimpleProcedure acc_stop)
        : io(io),socket(io),acc_stop(acc_stop) {}

    ~ServerSession() {}

    void start()
    {
        // Enqueue Read Handler
        async_read_until(
            socket,
            buf,
            '\n',
            boost::bind(&ServerSession::read_ok,this,_1) ) ;
    }

private:
    void read_ok(const boost::system::error_code e)
    {
        if (!e)
        {
            std::iostream ios(&buf) ;

            std::string tmp ;
            ios >> tmp ; // get input stream
            std::cout << tmp ;
            if (tmp == "end") { delete this; return; }
            else if (tmp == "bye")
            {
                acc_stop() ;
                delete this;
                return;
            }

            ios << tmp <<std::endl ; // retrun as it is

            // Enqueue Write Handler
            async_write(
                socket,
                buf,
                boost::bind(&ServerSession::write_ok,this,_1) ) ;
        }
        else if (e == boost::asio::error::operation_aborted)
        {
            // abort
            std::cout << "connection abort" ;
            delete this;
            return ;
        }
        else
        {
            std::cout << "error" ;
            delete this;
            return ;
        }
    }

    void write_ok(const boost::system::error_code e)
    {
        if (!e)
        {
            start() ; // Restart
        }
        else if (e == boost::asio::error::operation_aborted)
        {
            // abort
            std::cout << "connection abort" ;
            delete this;
            return ;
        }
        else
        {
            std::cout << "error" ;
            delete this;
            return ;
        }
    }

public:
    ip::tcp::socket& getSocket()
    {
        return (socket) ;
    }

} ;

class ServerModule : private boost::noncopyable
{
private:
    io_service& io ;   
    ip::tcp::acceptor accept ;
    ServerSession* session ;
public:
   
    ServerModule(io_service& io, const short port)
        :io(io),accept(io,ip::tcp::endpoint(ip::tcp::v4(),port))
    {
        start_accept() ;
    }

    void start_accept()
    {
       
        session = new ServerSession(io,boost::bind(&ServerModule::stop_accept,this)) ;
        // Enqueue Accept Handler
        accept.async_accept(
            session->getSocket(),
            boost::bind(&ServerModule::accept_ok,this,_1) ) ;
    }

    // Accept service will be closed. Then IOServce throw operation_aborted.
    // After queue get be empty, the io.run() will stop.
    void stop_accept()
    {
        accept.close() ;
    }

    void accept_ok(const boost::system::error_code e)
    {
        if (!e)
        {
            session->start() ;
            start_accept() ;
        }
        else if (e == boost::asio::error::operation_aborted)
        {
            // abort
            std::cout << "connection abort" ;
            delete session ;
            session = NULL ;
            return ;
        }
        else
        {
            // error
            std::cout << "error" ;
            delete session ;
            session = NULL ;
            return ;
        }
    }

} ;

} // namespace thorny_road
#endif //__ASIO_SERVER_MODULE__

ServerSessionクラスのコンストラクタでacc_stopとして関数オブジェクトを受け取るようにしています。"bye"と入力されると、ServerModule::stop_accept()が呼ばれます。ServerModule::stop_accept()ではaccept.close()が実行され、キューに入っている接続用のハンドラーが破棄されます。

ハンドラーが破棄されると、破棄されたハンドラー関数にboost::asio::error::operation_abortedが入ってきます。今回、分りやすいように標準出力に"abort"と出力するようにしています。

前回同様、telnetでアクセスし、byeと入力するとアプリケーションが終了すると思います。これは接続ハンドラーが破棄され、Sessionの中でも新たにハンドラーがキューに代入されないからです。

たとえば、telnetを2つ起動し、両方から接続していると、片方でbyeと入力してもアプリケーションは終了しません。これはもう片方のtelnetのハンドラーが残っているからです。

キューにハンドラーが無くなってもio.run()が終了しないようにio_service::workクラスが用意されています。
今後、試してみたいと思っています。

単体のスレッドで、複数のクライアントからアクセスできるのは素晴らしいことですが、やはり逐次処理のため、どれかが処理している間は他の処理がブロックされてしまいます。この問題を回避するためにはマルチスレッドを使う必要がありそうです。


ソースコードは自由にご使用ください。ただし問題が起きても責任はとれません。また、ソースコードに対する著作権は放棄していません。

2011年11月11日金曜日

BOOST::ASIOを使ってみる(6)

今回は非同期送信処理を作ってみます。基本的な構造は受信側と同じでModuleThreadクラスを使ってマルチスレッド化しています。送信処理の流れは名前解決→接続(コネクション)→データの送信→返信データの受信となります。

以下送信側のソースコードです。

#pragma once

#ifndef __ASIO_SENDER_MODULE__
#define __ASIO_SENDER_MODULE__

#include <iostream>
#include <boost/asio.hpp>
#include <boost/bind.hpp>
#include <boost/function.hpp>
#include <boost/thread.hpp>
#include <boost/shared_ptr.hpp>
#include <boost/lexical_cast.hpp>

namespace thorny_road{

using namespace boost::asio ;

// Simple Message data socket
class SimpleMsg
{
private :
    std::string msg ;
    boost::asio::streambuf buf ;
public:
    SimpleMsg(const std::string msg) : msg(msg)
    {
        std::ostream os(&buf) ;
        os.write((char *)(msg.c_str()),msg.size()) ;

        char ch = '\n' ;
        os.write(&ch,sizeof(char)) ; // Append Return Code
    }

    ~SimpleMsg() {}

    boost::asio::streambuf& rdbuf()
    {
        return buf ;
    }
} ;


template<typename SOK_PTYPE>
class SenderSession : private boost::noncopyable
{
private:
    io_service& io ;
    ip::tcp::socket socket ;
    boost::asio::streambuf buf ;
    std::string address ;
    short port ;
    boost::shared_ptr<ip::tcp::resolver> resolver ;
public:
    SenderSession(std::string& address,short port, SOK_PTYPE pdata,io_service& io)
        : io(io),socket(io),address(address),port(port)
    {
        // make resolver
        resolver.reset(new ip::tcp::resolver(io) );
        // make query
        ip::tcp::resolver::query query(address,boost::lexical_cast<std::string>(port));

        resolver->async_resolve(
            query,
            boost::bind(&SenderSession::resolve_ok, this,
                boost::asio::placeholders::error,
                pdata,
                boost::asio::placeholders::iterator));
    }

    ~SenderSession(){}

private:
    void resolve_ok(const boost::system::error_code e,SOK_PTYPE pdata,
            boost::asio::ip::tcp::resolver::iterator endpoint_iterator)
    {
        if (!e)
        {
            // Enqueue Connect Handler
            boost::asio::async_connect(
                socket,
                endpoint_iterator,
                boost::bind(&SenderSession::connect_ok, this,
                    boost::asio::placeholders::error,
                    pdata,
                    endpoint_iterator));

        }
        else if (e == boost::asio::error::operation_aborted)
        {
            // abort
            std::cout << "connection abort" ;
            std::cout << e.message() ;
            delete this;
            return ;
        }
        else
        {
            std::cout << "error" ;
            std::cout << e.message() ;
            delete this;
            return ;
        }
    }

    void connect_ok(const boost::system::error_code e,SOK_PTYPE pdata,
            boost::asio::ip::tcp::resolver::iterator endpoint_iterator)
    {
        if (!e)
        {
            // Enqueue Write Handler
            boost::asio::async_write(
                socket,
                pdata->rdbuf(),
                boost::bind(&SenderSession::receive_read, this,
                    boost::asio::placeholders::error,
                    pdata));
           
        }
        else if (e == boost::asio::error::operation_aborted)
        {
            // abort
            std::cout << "connection abort" ;
            std::cout << e.message() ;
            delete this;
            return ;
        }
        else if (endpoint_iterator != boost::asio::ip::tcp::resolver::iterator())
        {
            // failed to resolve endpoint.
            boost::asio::async_connect(
                socket,
                endpoint_iterator,
                boost::bind(&SenderSession::connect_ok, this,
                    boost::asio::placeholders::error,
                    pdata,
                    ++endpoint_iterator)); // Try next resolved endpoint name.
        }
        else
        {
            std::cout << "error" ;
            std::cout << e.message() ;
            delete this;
            return ;
        }
    }

    void receive_read(const boost::system::error_code e,SOK_PTYPE pdata)
    {
        if (!e)
        {
            // Call Self until end of receive data.
            async_read(
                socket,
                pdata->rdbuf(),
                boost::asio::transfer_at_least(1),
                boost::bind(&SenderSession::receive_read,this,
                    boost::asio::placeholders::error,
                    pdata)) ;
        }
        else if (e == boost::asio::error::operation_aborted)
        {
            // abort
            std::cout << "connection abort" ;
            std::cout << e.message() ;
            delete this;
            return ;
        }
        else if (e != boost::asio::error::eof)
        {
            std::cout << "error" ;
            std::cout << e.message() ;
            delete this;
            return ;
        }
        else
        {
            delete this ; // end of session.
        }
    }


public:
    ip::tcp::socket& getSocket()
    {
        return (socket) ;
    }

} ;

class SenderModule : private boost::noncopyable
{
private:
    io_service& io ;   
    short port ;
    boost::mutex _lock ; // Lock Object
public:
    SenderModule(io_service& io, const short port)
        :io(io),port(port)
    {}

    ~SenderModule() {}

    template<typename U>
    void send(std::string address, U pdata)
    {
        boost::mutex::scoped_lock cs(_lock) ; // Lock until escaping this scope

        SenderSession *session = new SenderSession(address,port,pdata,io) ;
    }

} ;

} // namespace thorny_road
#endif //__ASIO_SENDER_MODULE__
メイン
#include "stdafx.h"
#include "SenderModule.hpp"
#include "ModuleThread.hpp" // for threading
#include <boost/date_time/posix_time/posix_time.hpp>
#include <boost/bind.hpp>


using namespace thorny_road ;

int _tmain(int argc, _TCHAR* argv[])
{
    ModuleThread<SenderModule> sender(2085,2) ; // create module.

    boost::shared_ptr<SimpleMsg> msg(new SimpleMsg("[message]")) ;
    sender.getModule()->send("localhost",msg) ;

    boost::this_thread::sleep(boost::posix_time::milliseconds(1000)) ;

    return 0;
}
今回は単純に文字列を送っています。上の例では"[message]"を送信しています。送信先はsend(...)の第一引数にアドレスを入れます。今回は"localhost"に送るようにしています。第二引数にはパケットデータクラスのポインタを入れます。パケットデータクラスは以下のようにバッファを参照するためのメソッド( rdbuf() )が用意されていればコンパイルが通ります(今のところは・・・)。
class SimpleMsg
{
private :
    std::string msg ;
    boost::asio::streambuf buf ;
public:
    SimpleMsg(const std::string msg) : msg(msg)
    {
        std::ostream os(&buf) ;
        os.write((char *)(msg.c_str()),msg.size()) ;

        char ch = '\n' ;
        os.write(&ch,sizeof(char)) ; // Append Return Code
    }

    ~SimpleMsg() {}

    boost::asio::streambuf& rdbuf()
    {
        return buf ;
    }
} ;
今回の例では単純に文字列をバッファストリームに書き込んでいます。受信側では最初の受信では'\n'まで読み込むので最後に'\n'を足しています。
また、ModuleThreadgetModule()は前回まではモジュールクラスのポインタを返していましたが、スマートポインタで返すように変更しています。私はポインターは面白くて好きですが、苦い経験も多々あるので最近はスマートポインタを多用しています。DELPHIでは参照カウンターを持っているInterfaceをよく使います。

受信側のプログラムが実行している状態で、送信側のプログラムを実行すると送受信が行われますが、送信セッションが受信待ち状態のままになります。これは受信側のセッションが終了していないためです。受信側で返信した後セッション閉じるようにすれば送信側もセッションが切れます。
    void write_ok(const boost::system::error_code e)
    {
        if (!e)
        {
            // start() ; // restart
            delete this; // finish this packet
        }
DELPHIにもC++のTemplateのような機能が欲しいと思うことが多々ありますが、ビルドが高速というのがDELPHIの良い点のひとつなので難しいところです。

ソースコードは自由にご使用ください。ただし問題が起きても責任はとれません。また、ソースコードに対する著作権は放棄していません。

2011年11月10日木曜日

BOOST::ASIOを使ってみる(3)

前回までのソースコードでboost::noncopyableの使い方が間違えていました。boost::noncopyableは簡単にクラスのコピーコンストラクタとコピー代入演算子を隠蔽するものですが、privateで継承しないと意味が無かったです。

class ServerModule : private boost::noncopyable

今回は、接続後のやり取りを分けてみたいと思います。行うことはシンプルで、非同期処理関数をクラスにまとめて分けます。

ソースコードは以下です。

ServerModule.hpp

#pragma once

#ifndef __ASIO_SERVER_MODULE__
#define __ASIO_SERVER_MODULE__

#include <iostream>
#include <boost/asio.hpp>
#include <boost/bind.hpp>

namespace thorny_road{

using namespace boost::asio ;

class ServerSession : private boost::noncopyable
{
private:
    io_service& io ;
    ip::tcp::socket socket ;
    boost::asio::streambuf buf ;

public:
    ServerSession(io_service& io) : io(io),socket(io) {}

    ~ServerSession() {}

    void start()
    {
        // Enqueue Read Handler
        async_read_until(
            socket,
            buf,
            '\n',
            boost::bind(&ServerSession::read_ok,this,_1) ) ;
    }

private:
    void read_ok(const boost::system::error_code e)
    {
        if (!e)
        {
            std::iostream ios(&buf) ;

            std::string tmp ;
            ios >> tmp ; // get input stream
            if (tmp == "end") { delete this; return; }
            ios << tmp <<std::endl ; // retrun as it is

            // Enqueue Write Handler
            async_write(
                socket,
                buf,
                boost::bind(&ServerSession::write_ok,this,_1) ) ;
        }
    }

    void write_ok(const boost::system::error_code e)
    {
        if (!e)
        {
            start() ; // Restart
        }
    }

public:
    ip::tcp::socket& getSocket()
    {
        return (socket) ;
    }

} ;

class ServerModule : private boost::noncopyable
{
private:
    io_service& io ;    
    ip::tcp::acceptor accept ;
    ServerSession* session ;
public:
    
    ServerModule(io_service& io, const short port)
        :io(io),accept(io,ip::tcp::endpoint(ip::tcp::v4(),port))
    {
        start_accept() ;
    }

    void start_accept()
    {
        session = new ServerSession(io) ;
        // Enqueue Accept Handler
        accept.async_accept(
            session->getSocket(),
            boost::bind(&ServerModule::accept_ok,this,_1) ) ;
    }

    void accept_ok(const boost::system::error_code e)
    {
        if (!e)
        {
            session->start() ;
            start_accept() ;
        }
        else
        {
            // error
            delete session ;
            session = NULL ;
            return ;
        }
    }

} ;

} // namespace thorny_road
#endif //__ASIO_SERVER_MODULE__


メインは変更はありません。接続後の処理はServerSessionクラスで行うようになりました。このため、SocketはServerSessionが管理するようになります。
ServerModuleでは接続完了後(accept_ok(...)が呼ばれると)、start_accept()を呼び、新たにSessionを作成し接続待ち状態になります。
ServerSessionでは前回同様入力をそのままエコーバックします。ただし、入力が"end"の場合、ServerSessionオブジェクトを破棄され、接続を切断します。

if (tmp == "end") { delete this; return; }

また、今回は簡単なboost::system::error_codeの処理をいれています。ハンドラーの引数であるboost::system::error_codeは正常の場合、0になります。

ようやく「らしく」なってきました。次はどうしようか考え中です。

ソースコードは自由にご使用ください。ただし問題が起きても責任はとれません。また、ソースコードに対する著作権は放棄していません。

BOOST::ASIOを使ってみる(2)

ネットワークでパケットのやりとりを行う際に同期処理(接続があるまで待っている)だけで行うと、単体スレッドでは実用的ではありません。その場合、マルチスレッド処理を使って行いますが、パフォーマンスが落ちる場合があります。
Boost::Asioでは、非同期処理を推奨しており、必要なときにイベントを発生させて処理させることで、単体スレッドでもパケットのやり取りができるようになるようです。
Boost::AsioではProactorデザインパターンを使っているようです。

http://www.boost.org/doc/libs/1_47_0/doc/html/boost_asio/overview/core/async.html

どのようなデザインパターンかは勉強不足のため分りません(汗)が、キューに(イベント)ハンドラーを入れていき、イベントが完了したら完了(イベント)ハンドラーを呼ぶらしいです。

ソースコードは以下です。

ServerModule.hpp

#pragma once

#ifndef __ASIO_SERVER_MODULE__
#define __ASIO_SERVER_MODULE__

#include <iostream>
#include <boost/asio.hpp>
#include <boost/bind.hpp>

namespace thorny_road{

using namespace boost::asio ;

class ServerModule : public boost::noncopyable
{
private:
    io_service& io ;
    ip::tcp::acceptor accept ;
    boost::asio::streambuf buf ;
    ip::tcp::socket socket ;
public:
   
    ServerModule(io_service& io, const short port)
        :io(io),accept(io,ip::tcp::endpoint(ip::tcp::v4(),port)),socket(io)
    {
        start_accept() ;
    }

    void start_accept()
    {
        // Enqueue Accept Handler
        accept.async_accept(
            socket,
            boost::bind(&ServerModule::accept_ok,this,_1) ) ;
    }

    void accept_ok(const boost::system::error_code e)
    {
        // Enqueue Read Handler
        async_read_until(
            socket,
            buf,
            '\n',
            boost::bind(&ServerModule::read_ok,this,_1) ) ;
    }

    void read_ok(const boost::system::error_code e)
    {
        std::iostream ios(&buf) ;

        std::string tmp ;
        ios >> tmp ; // get input stream
        ios << tmp <<std::endl ; // retrun as it is

        // Enqueue Read Handler
        async_write(
            socket,
            buf,
            boost::bind(&ServerModule::write_ok,this,_1) ) ;
    }

    void write_ok(const boost::system::error_code e)
    {
        // Finish
    }


} ;

} // namespace snb
#endif //__ASIO_SERVER_MODULE__

Main.cpp
#include "stdafx.h"
#include "ServerModule.hpp"

#include <boost/asio.hpp>

using namespace thorny_road ;

int _tmain(int argc, _TCHAR* argv[])
{
    boost::asio::io_service io ;

    ServerModule svr(io,2085) ; // Start Server Module

    io.run() ;

    return 0;
}


非同期処理ではSocketを使って相互処理を行います。同期処理ではaccept.accept(...)を使いましたが、非同期処理ではaccept.async_accept(...)を使います。第一引数がSocketで、第二引数が完了ハンドラーとなります。ですので、この場合、接続が完了するとaccept_ok(...)が呼ばれます。boost::bind(...)は関数オブジェクトを簡単に使うためにBoostで用意されているものです。DELPHIの「.... of object」みたいなものかな?
async_read_until(...)は指定した完了条件になるまで受信する非同期処理関数です。この場合、第三引数に改行コードを設定しているので、改行を受信と完了ハンドラーが呼ばれます。
async_write(...)は指定したストリームに入っている情報を送信します。
一般にBoost::Asioではboost::asio::streambufを使ってバッファ処理をするのが良いようです。勉強不足(汗)で利点はよく分りません・・・。
メイン処理(_tmain)ではio.run()を追加しています。これは非同期処理ではハンドラーをキューに入れると、待機せずに処理が戻ってくるためです。io.run()で「ハンドラーが無くなる」までループします。
前回同様telnetから接続すると、入力したものがそのままエコーバックされると思います。

現状のままだと、一度やり取りをするとプログラムが終わってしまいます。そこで、write_okの中でstart_accept()を呼ぶようにします。こうすることで再び接続待ちになります。

    void write_ok(const boost::system::error_code e)
    {
        // Finish
        start_accept() ;
    }


ただし、この場合、io.run()が無限にループするようになります。

次はパケットのやり取り部分を分けたいと思います。

ソースコードは自由に使用ください。ただし、問題が起きても責任はとりません。

2011年11月8日火曜日

BOOST::ASIOを使ってみる(1)

C++とBoostの勉強のため、Boost::Asioを使ってみたいと思います。おかしな部分も多いので、ご指摘頂けると嬉しいです。
Boostのインストール等は省略します。環境は以下です。

-WindowsXP
-Boost1.47.0

BoostのAsioはネットワークやローレベルのIOを扱うための同期・非同期モデルです。これを使えば、同期・非同期でのネットワーク処理などが行えます。
(# http://www.boost.org/doc/libs/1_47_0/doc/html/boost_asio.html

まずは簡単なものから作りたいので、文字を打ったらエコーバックするだけのものを作りたいと思います。

以下ソースコードです。

ServerModule.hpp

#pragma once

#ifndef __ASIO_SERVER_MODULE__
#define __ASIO_SERVER_MODULE__

#include <iostream>
#include <boost/asio.hpp>
#include <boost/bind.hpp>

namespace thorny_road{

using namespace boost::asio ;

class ServerModule : public boost::noncopyable
{
private:
    io_service& io ;
    ip::tcp::acceptor accept ;
    ip::tcp::iostream buf ;
public:
    
    ServerModule(io_service& io, const short port)
        :io(io),accept(io,ip::tcp::endpoint(ip::tcp::v4(),port))
    {
        start_accept() ;
    }

    void start_accept()
    {
        accept.accept( *buf.rdbuf() ) ;
        std::string tmp ;
        buf >> tmp ; // get input stream
        buf << tmp <<std::endl ; // retrun as it is
    }


} ;

} // namespace snb
#endif //__ASIO_SERVER_MODULE__


Main.cpp
#include "stdafx.h"
#include "ServerModule.hpp"

#include <boost/asio.hpp>

using namespace thorny_road ;

int _tmain(int argc, _TCHAR* argv[])
{
    boost::asio::io_service io ;

    ServerModule svr(io,2085) ; // Start Server Module

    return 0;
}

勉強を兼ねてServerはクラスで構築しました。また、ヘッダーファイルにコードを書いています。その方が何かと載せやすいというのと、DELPHIもよく使うため慣れているからです。

まず最初にio_serviceを作成します。それを引数にServerModuleのインスタンスを作成しています。ServerModuleのコンストラクタが呼ばれると、start_accept()関数が呼ばれて接続待ちになります。accept(...)は同期処理になるので、接続が来るまで待つことになります。

これを実行すると、インプット待ち状態になります。
以下のようにTelnetでアクセスすると、接続できるようになります。
telnet localhost 2085

簡単にネットワークのプログラムが書けました。Boostは便利ですね。

次は非同期処理に挑戦したいですね。

ソースコードはご自由に使用ください。ただし、問題が起きても責任は取りません。