Boost.Asio Server dan RAII

Oct 08 2020

Saya mencoba menerapkan aplikasi server jaringan di C ++ menggunakan Boost.Asio.

Berikut adalah persyaratan yang saya coba penuhi:

  • Aplikasi hanya membuat satu contoh boost::io_context.
  • Single io_contextsedang run()oleh Thread Pool bersama. Jumlah utas tidak ditentukan.
  • Aplikasi dapat membuat instance beberapa objek Server. Server baru dapat muncul dan dibunuh kapan saja.
  • Setiap Server dapat menangani koneksi dari banyak klien.

Saya mencoba menerapkan pola RAII untuk kelas Server. Yang ingin saya jamin adalah bahwa ketika Server dibatalkan alokasinya, semua koneksinya akan ditutup sepenuhnya. Ada 3 cara untuk menutup setiap koneksi:

  1. Klien merespons dan tidak ada lagi pekerjaan yang harus dilakukan dalam sebuah koneksi.
  2. Server sedang dibatalkan alokasinya dan menyebabkan semua koneksi yang hidup ditutup.
  3. Koneksi dihentikan secara manual dengan stop()metode pemanggilan .

Saya telah sampai pada solusi yang tampaknya memenuhi semua kriteria di atas tetapi karena Boost.Asio masih cukup baru bagi saya, saya ingin memverifikasi bahwa apa yang saya lakukan benar. Juga ada beberapa hal yang secara khusus saya tidak 100% yakin:

  • Saya mencoba untuk menghapus mutexdari kelas Server dan sebagai gantinya menggunakan a stranduntuk semua sinkronisasi tetapi saya tidak dapat menemukan cara yang jelas untuk melakukannya.
  • Karena Thread Pool hanya dapat terdiri dari 1 utas dan utas ini mungkin adalah apa yang memanggil penghancur Server yang harus saya panggil io_context::poll_one()dari destruktor untuk memberikan kesempatan bagi semua koneksi yang tertunda untuk menyelesaikan pematian dan mencegah kemungkinan kebuntuan.
  • Saya akan menerima saran lain untuk perbaikan yang dapat Anda pikirkan.

Bagaimanapun, inilah kode dengan beberapa tes unit (versi langsung di Coliru: http://coliru.stacked-crooked.com/a/1afb0dc34dd09008 ):

#include <boost/asio/io_context.hpp>
#include <boost/asio/io_context_strand.hpp>
#include <boost/asio/executor.hpp>
#include <boost/asio/deadline_timer.hpp>
#include <boost/asio/dispatch.hpp>
#include <iostream>
#include <string>
#include <vector>
#include <memory>
#include <list>
using namespace std;
using namespace boost::asio;
using namespace std::placeholders;


class Connection;


class ConnectionDelegate
{
public:
    virtual ~ConnectionDelegate() { }
    
    virtual class executor executor() const = 0;
    virtual void didReceiveResponse(shared_ptr<Connection> connection) = 0;
};


class Connection: public enable_shared_from_this<Connection>
{
public:
    Connection(string name, io_context& ioContext)
    : _name(name)
    , _ioContext(ioContext)
    , _timer(ioContext)
    {
    }
    
    const string& name() const
    {
        return _name;
    }
    void setDelegate(ConnectionDelegate *delegate)
    {
        _delegate = delegate;
    }
    
    void start()
    {
        // Simulate a network request
        _timer.expires_from_now(boost::posix_time::seconds(3));
        _timer.async_wait(bind(&Connection::handleResponse, shared_from_this(), _1));
    }
    void stop()
    {
        _timer.cancel();
    }
    
private:
    string _name;
    io_context& _ioContext;
    boost::asio::deadline_timer _timer;
    ConnectionDelegate *_delegate;
    
    void handleResponse(const boost::system::error_code& errorCode)
    {
        if (errorCode == error::operation_aborted)
        {
            return;
        }
        dispatch(_delegate->executor(),
                 bind(&ConnectionDelegate::didReceiveResponse, _delegate, shared_from_this()));
    }
};


class Server: public ConnectionDelegate
{
public:
    Server(string name, io_context& ioContext)
    : _name(name)
    , _ioContext(ioContext)
    , _strand(_ioContext)
    {
    }
    ~Server()
    {
        stop();
        assert(_connections.empty());
        assert(_connectionIterators.empty());
    }
    weak_ptr<Connection> addConnection(string name)
    {
        auto connection = shared_ptr<Connection>(new Connection(name, _ioContext), bind(&Server::deleteConnection, this, _1));
        {
            lock_guard<mutex> lock(_mutex);
            _connectionIterators[connection.get()] = _connections.insert(_connections.end(), connection);
        }
        connection->setDelegate(this);
        connection->start();
        return connection;
    }
    
    vector<shared_ptr<Connection>> connections()
    {
        lock_guard<mutex> lock(_mutex);
        
        vector<shared_ptr<Connection>> connections;
        for (auto weakConnection: _connections)
        {
            if (auto connection = weakConnection.lock())
            {
                connections.push_back(connection);
            }
        }
        return connections;
    }
    void stop()
    {
        auto connectionsCount = 0;
        for (auto connection: connections())
        {
            ++connectionsCount;
            connection->stop();
        }
        
        while (connectionsCount != 0)
        {
            _ioContext.poll_one();
            connectionsCount = connections().size();
        }
    }
    
    // MARK: - ConnectionDelegate
    class executor executor() const override
    {
        return _strand;
    }
    void didReceiveResponse(shared_ptr<Connection> connection) override
    {
        // Strand to protect shared resourcess to be accessed by this method.
        assert(_strand.running_in_this_thread());
        
        // Here I plan to execute some business logic and I need both Server & Connection to be alive.
        std::cout << "didReceiveResponse - server: " << _name << ", connection: " << connection->name() << endl;
    }
    
private:
    typedef list<weak_ptr<Connection>> ConnectionsList;
    typedef unordered_map<Connection*, ConnectionsList::iterator> ConnectionsIteratorMap;
    
    string _name;
    io_context& _ioContext;
    io_context::strand _strand;
    ConnectionsList _connections;
    ConnectionsIteratorMap _connectionIterators;
    mutex _mutex;
    
    void deleteConnection(Connection *connection)
    {
        {
            lock_guard<mutex> lock(_mutex);
            auto iterator = _connectionIterators[connection];
            _connections.erase(iterator);
            _connectionIterators.erase(connection);
        }
        default_delete<Connection>()(connection);
    }
};


void testConnectionClosedByTheServer()
{
    io_context ioContext;
    auto server = make_unique<Server>("server1", ioContext);
    
    auto weakConnection = server->addConnection("connection1");
    assert(weakConnection.expired() == false);
    assert(server->connections().size() == 1);
    
    server.reset();
    assert(weakConnection.expired() == true);
}

void testConnectionClosedAfterResponse()
{
    io_context ioContext;
    auto server = make_unique<Server>("server1", ioContext);
    
    auto weakConnection = server->addConnection("connection1");
    assert(weakConnection.expired() == false);
    assert(server->connections().size() == 1);
    
    while (!weakConnection.expired())
    {
        ioContext.poll_one();
    }
    assert(server->connections().size() == 0);
}

void testConnectionClosedManually()
{
    io_context ioContext;
    auto server = make_unique<Server>("server1", ioContext);
    
    auto weakConnection = server->addConnection("connection1");
    assert(weakConnection.expired() == false);
    assert(server->connections().size() == 1);
    
    weakConnection.lock()->stop();
    ioContext.run();
    
    assert(weakConnection.expired() == true);
    assert(server->connections().size() == 0);
}

void testMultipleServers()
{
    io_context ioContext;
    auto server1 = make_unique<Server>("server1", ioContext);
    auto server2 = make_unique<Server>("server2", ioContext);

    auto weakConnection1 = server1->addConnection("connection1");
    auto weakConnection2 = server2->addConnection("connection2");

    server1.reset();
    assert(weakConnection1.expired() == true);
    assert(weakConnection2.expired() == false);
}

void testDeadLock()
{
    io_context ioContext;
    auto server = make_unique<Server>("server1", ioContext);
    
    auto weakConnection = server->addConnection("connection1");
    assert(weakConnection.expired() == false);
    assert(server->connections().size() == 1);
    
    auto connection = weakConnection.lock();
    server.reset(); // <-- deadlock, but that's OK, i will try to prevent it by other means
}


int main()
{
    testConnectionClosedByTheServer();
    testConnectionClosedAfterResponse();
    testConnectionClosedManually();
    // testDeadLock();
}

Hormat kami, Marek

Jawaban

2 Quuxplusone Oct 10 2020 at 05:16

Saya tidak cukup tahu tentang Asio untuk memberikan umpan balik yang saya tahu Anda inginkan, tetapi berikut adalah beberapa pembersihan kecil yang dapat Anda lakukan:

  • Jangan using namespace std. Anda mungkin juga harus menghindari using namespacehal lain, hanya untuk kejelasan.

  • virtual ~ConnectionDelegate() { }bisa jadi virtual ~ConnectionDelegate() = default;sebagai gantinya. Ini menunjukkan niat Anda sedikit lebih baik.

  • ~Server()seharusnya ~Server() override, untuk menunjukkan bahwa itu menimpa fungsi anggota virtual. Secara umum, Anda harus menggunakan di overridemana pun itu secara fisik diizinkan oleh bahasa. (Saya pikir Anda melakukannya dengan benar di mana saja kecuali di destruktor.)

  • Connection(string name,dan Server(string name,keduanya tidak perlu membuat salinan string name.

  • Semua konstruktor Anda seharusnya explicit, untuk memberi tahu compiler bahwa misalnya pasangan braced {"hello world", myIOContext}tidak boleh secara implisit diperlakukan sebagai (atau secara implisit diubah menjadi) Serverobjek, bahkan tidak secara kebetulan.

  • Secara pribadi, saya menemukan penggunaan typedefs untuk ConnectionsListdan ConnectionsIteratorMapmenjadi lapisan tipuan yang tidak perlu. Saya lebih suka melihat std::list<std::weak_ptr<Connection>> _connections;di sana dalam antrean. Kalau saya butuh nama untuk tipe itu, saya bisa ucapkan saja decltype(_connections).

  • default_delete<Connection>()(connection)adalah cara mengucapkan yang bertele-tele delete connection. Bersikaplah langsung.

  • class executor executor()sangat membingungkan. Fakta bahwa Anda harus mengatakan classseharusnya ada bendera merah yang executorbukan nama yang tepat untuk kelas, atau executor()bukan nama yang tepat untuk metode ini. Pertimbangkan untuk mengubah nama metode menjadi get_executor(), misalnya. Saya berasumsi bahwa Anda tidak dapat mengubah nama class executorkarena tidak dideklarasikan dalam file ini; itu pasti berasal dari beberapa namespace Boost yang Anda usingbuat, bukan? (Jangan usingruang nama!)


Anda melewatkan banyak kesempatan untuk menghindari salinan melalui referensi dan / atau memindahkan semantik. Misalnya, dalam Server::connections(), saya akan menulis:

std::vector<std::shared_ptr<Connection>> connections() {
    std::lock_guard<std::mutex> lock(_mutex);
    std::vector<std::shared_ptr<Connection>> result;
    for (const auto& weakConnection : _connections) {
        if (auto sptr = weakConnection.lock()) {
            result.push_back(std::move(sptr));
        }
    }
    return result;
}

Ini menghindari menabrak refcount lemah dengan membuat weakConnectionreferensi alih-alih salinan, dan kemudian menghindari menabrak refcount kuat dengan menggunakan pindah alih-alih salin masuk push_back. Empat operasi atom diselamatkan! (Bukan berarti ini penting dalam kehidupan nyata, mungkin, tapi hei, selamat datang di Code Review.)


dispatch(_delegate->executor(),
         bind(&ConnectionDelegate::didReceiveResponse, _delegate, shared_from_this()));

Saya menemukan penggunaan bindmembingungkan, tetapi saya tidak tahu pasti (dan sebenarnya berharap seseorang akan berkomentar dan mencerahkan saya) - bind diperlukan di sini? Ini pasti akan lebih jelas untuk dibaca, lebih cepat untuk dikompilasi, dan tidak lebih lambat pada waktu runtime untuk menulis.

dispatch(
    _delegate->executor(),
    [self = shared_from_this(), d = _delegate]() {
        d->didReceiveResponse(self);
    }
);

Ini akan membuatnya sedikit lebih jelas tentang apa yang sebenarnya sedang disalin (satu shared_ptrtetap *thishidup, dan satu penunjuk mentah). Sebenarnya, saya bertanya-tanya apakah kita perlu menyembunyikan salinan pointer mentah; bisakah kita lolos dengan ini saja?

dispatch(
    _delegate->executor(),
    [self = shared_from_this()]() {
        self->_delegate->didReceiveResponse(self);
    }
);

Atau apakah Anda berharap terkadang Anda akan masuk ke tubuh lambda itu dengan d != self->_delegatedan itulah mengapa Anda memerlukan penunjuk tambahan?


Saya juga bertanya-tanya apakah itu mungkin untuk digunakan std::chrono::secondssebagai pengganti boost::posix_time::seconds. Bisakah Boost.Asio beroperasi dengan C ++ 11 saat std::chronoini?


_connectionIterators[connection.get()] = _connections.insert(_connections.end(), connection);

Saya merasa "kecerdasan" di sini berada di sisi yang salah dari tanda sama dengan. _connections.insert(_connections.end(), connection)sepertinya cara penulisan yang panjang _connections.push_back(connection). Begitu pula sebaliknya, saya terbiasa melihat orang mengganti map[k] = vdengan map.emplace(k, v)untuk kinerja dan kejelasan. Ingat bahwa map[k] = vpertama default-konstruksi map[k] , dan kemudian memberikan nilai baru ke dalamnya.

Ah, saya mengerti, Anda perlu menggunakan insertkarena insertmengembalikan iterator dan push_backtidak.

Tapi itu hanya menimbulkan pertanyaan: Mengapa Anda mencoba menyatukan dua efek samping menjadi satu baris? Jika kami diizinkan dua baris, kami hanya melakukan push_backdan kemudian mengatur map.emplace(connection.get(), std::prev(_connections.end())). Atau, heck, pada saat itu saya tidak akan benar-benar mengeluh

auto it = _connections.insert(_connections.end(), connection);
_connectionIterators.emplace(connection.get(), it);

Setelah melihat bendera merah, gali lebih dalam: apa perbedaan antara satu baris dan dua baris yang lebih jelas? Aha! Perbedaannya adalah yang terjadi jika _connections.insert(...)kehabisan memori dan lemparan. Dengan two-liner, _connectionIteratorstetap tak tersentuh. Dengan one-liner, Anda terlebih dahulu membuat secara default beberapa sampah berbahaya di _connectionIterators[connection.get()]dan kemudian menyebarkan pengecualian.

Jadi saya pikir ada argumen yang masuk akal untuk dibuat mendukung dua garis, hanya pada prinsip-prinsip umum.


Sekali lagi, jawaban ini tidak benar-benar menjawab perhatian utama Anda tentang RAII, tetapi saya berharap jawaban ini bisa memberi sedikit pemikiran.