Moved boost::asio::strand yet again.
This commit is contained in:
parent
c6806e6fde
commit
4bc50785ac
2 changed files with 14 additions and 13 deletions
|
|
@ -21,14 +21,14 @@ namespace SimpleWeb {
|
||||||
|
|
||||||
boost::asio::streambuf streambuf;
|
boost::asio::streambuf streambuf;
|
||||||
|
|
||||||
std::shared_ptr<socket_type> socket;
|
socket_type &socket;
|
||||||
|
|
||||||
Response(boost::asio::io_service& io_service, std::shared_ptr<socket_type> socket, boost::asio::yield_context& yield):
|
Response(boost::asio::io_service& io_service, socket_type &socket, boost::asio::yield_context& yield):
|
||||||
yield(yield), socket(socket), stream(&streambuf) {}
|
yield(yield), socket(socket), stream(&streambuf) {}
|
||||||
|
|
||||||
void flush() {
|
void flush() {
|
||||||
boost::system::error_code ec;
|
boost::system::error_code ec;
|
||||||
boost::asio::async_write(*socket, streambuf, yield[ec]);
|
boost::asio::async_write(socket, streambuf, yield[ec]);
|
||||||
|
|
||||||
if(ec)
|
if(ec)
|
||||||
throw std::runtime_error(ec.message());
|
throw std::runtime_error(ec.message());
|
||||||
|
|
@ -77,10 +77,12 @@ namespace SimpleWeb {
|
||||||
unsigned short remote_endpoint_port;
|
unsigned short remote_endpoint_port;
|
||||||
|
|
||||||
private:
|
private:
|
||||||
Request(): content(&streambuf) {}
|
Request(boost::asio::io_service &io_service): content(&streambuf), strand(io_service) {}
|
||||||
|
|
||||||
boost::asio::streambuf streambuf;
|
boost::asio::streambuf streambuf;
|
||||||
|
|
||||||
|
boost::asio::strand strand;
|
||||||
|
|
||||||
void read_remote_endpoint_data(socket_type& socket) {
|
void read_remote_endpoint_data(socket_type& socket) {
|
||||||
try {
|
try {
|
||||||
remote_endpoint_address=socket.lowest_layer().remote_endpoint().address().to_string();
|
remote_endpoint_address=socket.lowest_layer().remote_endpoint().address().to_string();
|
||||||
|
|
@ -149,7 +151,6 @@ namespace SimpleWeb {
|
||||||
boost::asio::io_service io_service;
|
boost::asio::io_service io_service;
|
||||||
boost::asio::ip::tcp::endpoint endpoint;
|
boost::asio::ip::tcp::endpoint endpoint;
|
||||||
boost::asio::ip::tcp::acceptor acceptor;
|
boost::asio::ip::tcp::acceptor acceptor;
|
||||||
boost::asio::strand strand;
|
|
||||||
size_t num_threads;
|
size_t num_threads;
|
||||||
std::vector<std::thread> threads;
|
std::vector<std::thread> threads;
|
||||||
|
|
||||||
|
|
@ -157,7 +158,7 @@ namespace SimpleWeb {
|
||||||
size_t timeout_content;
|
size_t timeout_content;
|
||||||
|
|
||||||
ServerBase(unsigned short port, size_t num_threads, size_t timeout_request, size_t timeout_send_or_receive) :
|
ServerBase(unsigned short port, size_t num_threads, size_t timeout_request, size_t timeout_send_or_receive) :
|
||||||
endpoint(boost::asio::ip::tcp::v4(), port), acceptor(io_service, endpoint), strand(io_service),
|
endpoint(boost::asio::ip::tcp::v4(), port), acceptor(io_service, endpoint),
|
||||||
num_threads(num_threads), timeout_request(timeout_request), timeout_content(timeout_send_or_receive) {}
|
num_threads(num_threads), timeout_request(timeout_request), timeout_content(timeout_send_or_receive) {}
|
||||||
|
|
||||||
virtual void accept()=0;
|
virtual void accept()=0;
|
||||||
|
|
@ -175,10 +176,10 @@ namespace SimpleWeb {
|
||||||
return timer;
|
return timer;
|
||||||
}
|
}
|
||||||
|
|
||||||
std::shared_ptr<boost::asio::deadline_timer> set_timeout_on_socket_strand_wrapped(std::shared_ptr<socket_type> socket, size_t seconds) {
|
std::shared_ptr<boost::asio::deadline_timer> set_timeout_on_socket(std::shared_ptr<socket_type> socket, std::shared_ptr<Request> request, size_t seconds) {
|
||||||
std::shared_ptr<boost::asio::deadline_timer> timer(new boost::asio::deadline_timer(io_service));
|
std::shared_ptr<boost::asio::deadline_timer> timer(new boost::asio::deadline_timer(io_service));
|
||||||
timer->expires_from_now(boost::posix_time::seconds(seconds));
|
timer->expires_from_now(boost::posix_time::seconds(seconds));
|
||||||
timer->async_wait(strand.wrap([socket](const boost::system::error_code& ec){
|
timer->async_wait(request->strand.wrap([socket](const boost::system::error_code& ec){
|
||||||
if(!ec) {
|
if(!ec) {
|
||||||
boost::system::error_code ec;
|
boost::system::error_code ec;
|
||||||
socket->lowest_layer().shutdown(boost::asio::ip::tcp::socket::shutdown_both, ec);
|
socket->lowest_layer().shutdown(boost::asio::ip::tcp::socket::shutdown_both, ec);
|
||||||
|
|
@ -191,7 +192,7 @@ namespace SimpleWeb {
|
||||||
void read_request_and_content(std::shared_ptr<socket_type> socket) {
|
void read_request_and_content(std::shared_ptr<socket_type> socket) {
|
||||||
//Create new streambuf (Request::streambuf) for async_read_until()
|
//Create new streambuf (Request::streambuf) for async_read_until()
|
||||||
//shared_ptr is used to pass temporary objects to the asynchronous functions
|
//shared_ptr is used to pass temporary objects to the asynchronous functions
|
||||||
std::shared_ptr<Request> request(new Request());
|
std::shared_ptr<Request> request(new Request(io_service));
|
||||||
request->read_remote_endpoint_data(*socket);
|
request->read_remote_endpoint_data(*socket);
|
||||||
|
|
||||||
//Set timeout on the following boost::asio::async-read or write function
|
//Set timeout on the following boost::asio::async-read or write function
|
||||||
|
|
@ -294,10 +295,10 @@ namespace SimpleWeb {
|
||||||
//Set timeout on the following boost::asio::async-read or write function
|
//Set timeout on the following boost::asio::async-read or write function
|
||||||
std::shared_ptr<boost::asio::deadline_timer> timer;
|
std::shared_ptr<boost::asio::deadline_timer> timer;
|
||||||
if(timeout_content>0)
|
if(timeout_content>0)
|
||||||
timer=set_timeout_on_socket_strand_wrapped(socket, timeout_content);
|
timer=set_timeout_on_socket(socket, request, timeout_content);
|
||||||
|
|
||||||
boost::asio::spawn(strand, [this, &resource_function, socket, request, timer](boost::asio::yield_context yield) {
|
boost::asio::spawn(request->strand, [this, &resource_function, socket, request, timer](boost::asio::yield_context yield) {
|
||||||
Response response(io_service, socket, yield);
|
Response response(io_service, *socket, yield);
|
||||||
|
|
||||||
try {
|
try {
|
||||||
resource_function(response, request);
|
resource_function(response, request);
|
||||||
|
|
|
||||||
|
|
@ -15,7 +15,7 @@ public:
|
||||||
void accept() {}
|
void accept() {}
|
||||||
|
|
||||||
bool parse_request_test() {
|
bool parse_request_test() {
|
||||||
std::shared_ptr<Request> request(new Request());
|
std::shared_ptr<Request> request(new Request(io_service));
|
||||||
|
|
||||||
stringstream ss;
|
stringstream ss;
|
||||||
ss << "GET /test/ HTTP/1.1\r\n";
|
ss << "GET /test/ HTTP/1.1\r\n";
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue