23 #ifndef INCLUDED_GR_BASIC_BLOCK_H 24 #define INCLUDED_GR_BASIC_BLOCK_H 32 #include <boost/enable_shared_from_this.hpp> 33 #include <boost/function.hpp> 34 #include <boost/foreach.hpp> 35 #include <boost/thread/condition_variable.hpp> 59 public boost::enable_shared_from_this<basic_block>
61 typedef boost::function<void(pmt::pmt_t)> msg_handler_t;
64 typedef std::map<pmt::pmt_t , msg_handler_t, pmt::comparator> d_msg_handlers_t;
65 d_msg_handlers_t d_msg_handlers;
67 typedef std::deque<pmt::pmt_t> msg_queue_t;
68 typedef std::map<pmt::pmt_t, msg_queue_t, pmt::comparator> msg_queue_map_t;
69 typedef std::map<pmt::pmt_t, msg_queue_t, pmt::comparator>::iterator msg_queue_map_itr;
70 std::map<pmt::pmt_t, boost::shared_ptr<boost::condition_variable>,
pmt::comparator> msg_queue_ready;
76 friend class flat_flowgraph;
77 friend class tpb_thread_body;
103 d_input_signature = iosig;
108 d_output_signature = iosig;
121 return (d_msg_handlers.find(which_port) != d_msg_handlers.end());
133 if(has_msg_handler(which_port)) {
134 d_msg_handlers[which_port](msg);
148 std::string
name()
const {
return d_name; }
159 basic_block_sptr to_basic_block();
169 std::string
alias(){
return alias_set()?d_symbol_alias:symbol_name(); }
181 void set_block_alias(std::string name);
184 void message_port_register_in(
pmt::pmt_t port_id);
185 void message_port_register_out(
pmt::pmt_t port_id);
217 if(msg_queue.find(which_port) == msg_queue.end())
218 throw std::runtime_error(
"port does not exist!");
219 return msg_queue[which_port].empty();
223 BOOST_FOREACH(msg_queue_map_t::value_type &i, msg_queue) {
224 rv &= msg_queue[i.first].empty();
231 return (empty_p(which_port) || !has_msg_handler(which_port));
235 BOOST_FOREACH(msg_queue_map_t::value_type &i, msg_queue) {
236 rv &= empty_handled_p(i.first);
243 if(msg_queue.find(which_port) == msg_queue.end())
244 throw std::runtime_error(
"port does not exist!");
245 return msg_queue[which_port].size();
263 return msg_queue[which_port].begin();
267 msg_queue[which_port].erase(it);
271 if(msg_queue.find(which_port) != msg_queue.end()) {
297 void add_rpc_variable(rpcbasic_sptr s)
299 d_rpc_vars.push_back(s);