Опытный
 
Профиль
Группа: Участник
Сообщений: 943
Регистрация: 17.6.2009
Репутация: нет Всего: 13
|
Привет! Помогите, пожалуйста, разобраться с асинхрнными таймерами. Не понимаю, почему они перестают тикать...похоже проблема не в таймерах... | Код | #include <iostream> #include <string> #include <vector>
#include <boost/date_time/posix_time/posix_time.hpp> #include <boost/asio.hpp> #include <boost/shared_ptr.hpp> #include <boost/shared_array.hpp> #include <boost/bind.hpp> #include <boost/thread.hpp>
#include "boost/simple_atomic.hpp"
boost::atomic::simple_atomic<int> termination(0);
class UdpServer { public: UdpServer(boost::asio::io_service & ioService, unsigned short port);
private: /// @brief handle receiving void handleReceive_(const boost::system::error_code& error, size_t bytesReceived);
/// @brief handle termination void checkProcessController_(const boost::system::error_code& error);
boost::asio::ip::udp::socket socket_; boost::asio::deadline_timer timer_;
boost::shared_ptr<std::string> packet_; enum {MAX_LEN_ = 0xFFFF}; char data_[MAX_LEN_]; };
UdpServer::UdpServer(boost::asio::io_service & ioService, unsigned short port): socket_(ioService), timer_(ioService) { using namespace boost::asio; using namespace boost::posix_time;
socket_.open(boost::asio::ip::udp::v4()); socket_.set_option(boost::asio::ip::udp::socket::reuse_address(true)); socket_.bind(boost::asio::ip::udp::endpoint(boost::asio::ip::udp::v4(), port));
socket_.async_receive(buffer(data_, MAX_LEN_), boost::bind(&UdpServer::handleReceive_, this, placeholders::error, placeholders::bytes_transferred));
timer_.expires_from_now(millisec(500)); timer_.async_wait(boost::bind(&UdpServer::checkProcessController_, this, boost::asio::placeholders::error)); }
void UdpServer::handleReceive_(const boost::system::error_code& error, size_t bytesReceived) { using namespace boost::asio; socket_.async_receive(buffer(data_, MAX_LEN_), boost::bind(&UdpServer::handleReceive_, this, placeholders::error, placeholders::bytes_transferred)); }
void UdpServer::checkProcessController_(const boost::system::error_code& error) { using namespace boost::posix_time; if(termination.get() == 1) { socket_.close(); socket_.get_io_service().stop(); } else { timer_.expires_from_now(millisec(500)); timer_.async_wait(boost::bind(&UdpServer::checkProcessController_, this, boost::asio::placeholders::error)); }
if(error) { std::cout << "timeout error: " << error.message() << std::endl; } }
void runServer(unsigned short port) { boost::asio::io_service ioService; UdpServer server(ioService, port); ioService.run(); }
void generate(unsigned short port) { using boost::asio::ip::udp;
boost::asio::io_service ioService; boost::asio::ip::udp::socket socket(ioService, udp::endpoint(udp::v4(), 0)); boost::asio::ip::address_v4 addr(0x7F000001UL); //127.0.0.1 udp::endpoint endpoint(addr, port);
std::string message("message"); while(termination.get() == 0) { socket.send_to(boost::asio::buffer(message), endpoint); } }
int main() { using namespace boost::posix_time;
boost::thread s0(boost::bind(runServer, 10000)); boost::thread s1(boost::bind(runServer, 10001)); boost::thread s2(boost::bind(runServer, 10002)); boost::thread g0(boost::bind(generate, 10000)); boost::thread g1(boost::bind(generate, 10001)); boost::thread g2(boost::bind(generate, 10002));
boost::this_thread::sleep(seconds(3)); termination.exchange(1); g0.join(); g1.join(); g2.join(); s0.join(); s1.join(); s2.join(); std::cout << "terminated" << std::endl; }
|
| Код | /// /// @file /// @brief simple_atomic variable definition /// /// Copyright (c) 2011 Andrey Lizunov /// /// Use, modification, and distribution are subject to the /// Boost Software License, Version 1.0. (See at http://www.boost.org/LICENSE_1_0.txt) ///
#ifndef BOOST_ATOMIC_SIMPLE_ATOMIC_HPP_INCLUDE #define BOOST_ATOMIC_SIMPLE_ATOMIC_HPP_INCLUDE
#include <utility> #include <boost/noncopyable.hpp> #include <boost/static_assert.hpp> #include <boost/type_traits.hpp> #include <boost/function.hpp>
#if defined(_MSC_VER) #include <intrin.h> #elif (__GNUC__ * 10000 + __GNUC_MINOR__ * 100 + __GNUC_PATCHLEVEL__) > 40100 #else # error no implementation #endif
namespace boost { namespace atomic { /// @brief facility to work with simple_atomic variables /// @tparam any POD type fit to sizeof(void *) template<typename T> class simple_atomic: private boost::noncopyable { public:
/// @brief constructor /// @param value - initial value explicit simple_atomic(T value = T()): instance_((void volatile*)value) { BOOST_STATIC_ASSERT(boost::is_pod<T>::type::value && sizeof(T) <= sizeof(void *) && sizeof(std::size_t) == sizeof(void *)); }
/// @brief get value with full memory barrier /// @return current value T get() const;
/// @brief Compare And Swap facility /// @param [in, out] oldValue - value to compare with /// if CAS fails contains current value /// @param newValue - value to load into memory /// @return result of the CAS /// @retval false if fails and true otherwise bool strong_cas(T & oldValue, T const& newValue);
/// @brief exchange values /// @param newValue - value to load into memory /// @return the previous value T exchange(T const& newValue);
/// @brief aplly user defined operation atomically /// @param operation - user defined operation /// @return pair of values /// @retval first - oldvalue, second - new value std::pair<T, T> apply(boost::function<T(T)> const& operation);
private: void volatile* instance_;
/// @brief full memory barrier static void inline memory_barrier_(); };
// member function definitions template<typename T> T simple_atomic<T>::get() const { this->memory_barrier_(); // intermediate cast to std::size_t is a workaround for atomic enums T res = (T)(std::size_t)(this->instance_); this->memory_barrier_(); return res; }
template<typename T> bool simple_atomic<T>::strong_cas(T & oldValue, T const& newValue) { void *volatile* dest = (void *volatile*) &this->instance_; void * res;
#if defined(_MSC_VER) res = InterlockedCompareExchangePointer(dest, (void *)newValue, (void *)oldValue); #elif (__GNUC__ * 10000 + __GNUC_MINOR__ * 100 + __GNUC_PATCHLEVEL__) > 40100 res = __sync_val_compare_and_swap(dest, (void *)oldValue, (void *)newValue); #else # error no implementation #endif
if(res != (void *)oldValue) { // intermediate cast to std::size_t is a workaround for atomic enums oldValue = (T)(std::size_t)res; return false; } return true; }
template<typename T> T simple_atomic<T>::exchange(T const& newValue) { T oldValue = this->get(); while(!this->strong_cas(oldValue, newValue)); return oldValue; }
template<typename T> std::pair<T, T> simple_atomic<T>::apply(boost::function<T(T)> const& operation) { T oldValue = this->get(); T newValue; do { newValue = operation(oldValue); } while(!this->strong_cas(oldValue, newValue));
return std::make_pair(oldValue, newValue); }
template<typename T> void simple_atomic<T>::memory_barrier_() { #if defined(_MSC_VER) _ReadWriteBarrier(); #elif (__GNUC__ * 10000 + __GNUC_MINOR__ * 100 + __GNUC_PATCHLEVEL__) > 40100 __sync_synchronize(); #else # error No implementation #endif }
}// namespace atomic }// namespace boost
#endif
|
Это сообщение отредактировал(а) Леопольд - 11.8.2011, 11:17
--------------------
вопросов больше чем ответов
|