mirror of
https://github.com/FairRootGroup/FairMQ.git
synced 2025-10-13 16:46:47 +00:00
611 lines
90 KiB
HTML
611 lines
90 KiB
HTML
<!DOCTYPE html PUBLIC "-//W3C//DTD XHTML 1.0 Transitional//EN" "https://www.w3.org/TR/xhtml1/DTD/xhtml1-transitional.dtd">
|
|
<html xmlns="http://www.w3.org/1999/xhtml">
|
|
<head>
|
|
<meta http-equiv="Content-Type" content="text/xhtml;charset=UTF-8"/>
|
|
<meta http-equiv="X-UA-Compatible" content="IE=9"/>
|
|
<meta name="generator" content="Doxygen 1.8.18"/>
|
|
<meta name="viewport" content="width=device-width, initial-scale=1"/>
|
|
<title>FairMQ: fairmq/shmem/Socket.h Source File</title>
|
|
<link href="tabs.css" rel="stylesheet" type="text/css"/>
|
|
<script type="text/javascript" src="jquery.js"></script>
|
|
<script type="text/javascript" src="dynsections.js"></script>
|
|
<link href="search/search.css" rel="stylesheet" type="text/css"/>
|
|
<script type="text/javascript" src="search/searchdata.js"></script>
|
|
<script type="text/javascript" src="search/search.js"></script>
|
|
<link href="doxygen.css" rel="stylesheet" type="text/css" />
|
|
</head>
|
|
<body>
|
|
<div id="top"><!-- do not remove this div, it is closed by doxygen! -->
|
|
<div id="titlearea">
|
|
<table cellspacing="0" cellpadding="0">
|
|
<tbody>
|
|
<tr style="height: 56px;">
|
|
<td id="projectalign" style="padding-left: 0.5em;">
|
|
<div id="projectname">FairMQ
|
|
 <span id="projectnumber">1.4.33</span>
|
|
</div>
|
|
<div id="projectbrief">C++ Message Queuing Library and Framework</div>
|
|
</td>
|
|
</tr>
|
|
</tbody>
|
|
</table>
|
|
</div>
|
|
<!-- end header part -->
|
|
<!-- Generated by Doxygen 1.8.18 -->
|
|
<script type="text/javascript">
|
|
/* @license magnet:?xt=urn:btih:cf05388f2679ee054f2beb29a391d25f4e673ac3&dn=gpl-2.0.txt GPL-v2 */
|
|
var searchBox = new SearchBox("searchBox", "search",false,'Search');
|
|
/* @license-end */
|
|
</script>
|
|
<script type="text/javascript" src="menudata.js"></script>
|
|
<script type="text/javascript" src="menu.js"></script>
|
|
<script type="text/javascript">
|
|
/* @license magnet:?xt=urn:btih:cf05388f2679ee054f2beb29a391d25f4e673ac3&dn=gpl-2.0.txt GPL-v2 */
|
|
$(function() {
|
|
initMenu('',true,false,'search.php','Search');
|
|
$(document).ready(function() { init_search(); });
|
|
});
|
|
/* @license-end */</script>
|
|
<div id="main-nav"></div>
|
|
<!-- window showing the filter options -->
|
|
<div id="MSearchSelectWindow"
|
|
onmouseover="return searchBox.OnSearchSelectShow()"
|
|
onmouseout="return searchBox.OnSearchSelectHide()"
|
|
onkeydown="return searchBox.OnSearchSelectKey(event)">
|
|
</div>
|
|
|
|
<!-- iframe showing the search results (closed by default) -->
|
|
<div id="MSearchResultsWindow">
|
|
<iframe src="javascript:void(0)" frameborder="0"
|
|
name="MSearchResults" id="MSearchResults">
|
|
</iframe>
|
|
</div>
|
|
|
|
<div id="nav-path" class="navpath">
|
|
<ul>
|
|
<li class="navelem"><a class="el" href="dir_d6b28f7731906a8cbc4171450df4b180.html">fairmq</a></li><li class="navelem"><a class="el" href="dir_6475741fe3587c0a949798307da6131d.html">shmem</a></li> </ul>
|
|
</div>
|
|
</div><!-- top -->
|
|
<div class="header">
|
|
<div class="headertitle">
|
|
<div class="title">Socket.h</div> </div>
|
|
</div><!--header-->
|
|
<div class="contents">
|
|
<div class="fragment"><div class="line"><a name="l00001"></a><span class="lineno"> 1</span> <span class="comment">/********************************************************************************</span></div>
|
|
<div class="line"><a name="l00002"></a><span class="lineno"> 2</span> <span class="comment"> * Copyright (C) 2014 GSI Helmholtzzentrum fuer Schwerionenforschung GmbH *</span></div>
|
|
<div class="line"><a name="l00003"></a><span class="lineno"> 3</span> <span class="comment"> * *</span></div>
|
|
<div class="line"><a name="l00004"></a><span class="lineno"> 4</span> <span class="comment"> * This software is distributed under the terms of the *</span></div>
|
|
<div class="line"><a name="l00005"></a><span class="lineno"> 5</span> <span class="comment"> * GNU Lesser General Public Licence (LGPL) version 3, *</span></div>
|
|
<div class="line"><a name="l00006"></a><span class="lineno"> 6</span> <span class="comment"> * copied verbatim in the file "LICENSE" *</span></div>
|
|
<div class="line"><a name="l00007"></a><span class="lineno"> 7</span> <span class="comment"> ********************************************************************************/</span></div>
|
|
<div class="line"><a name="l00008"></a><span class="lineno"> 8</span> <span class="preprocessor">#ifndef FAIR_MQ_SHMEM_SOCKET_H_</span></div>
|
|
<div class="line"><a name="l00009"></a><span class="lineno"> 9</span> <span class="preprocessor">#define FAIR_MQ_SHMEM_SOCKET_H_</span></div>
|
|
<div class="line"><a name="l00010"></a><span class="lineno"> 10</span>  </div>
|
|
<div class="line"><a name="l00011"></a><span class="lineno"> 11</span> <span class="preprocessor">#include "Common.h"</span></div>
|
|
<div class="line"><a name="l00012"></a><span class="lineno"> 12</span> <span class="preprocessor">#include "Manager.h"</span></div>
|
|
<div class="line"><a name="l00013"></a><span class="lineno"> 13</span> <span class="preprocessor">#include "Message.h"</span></div>
|
|
<div class="line"><a name="l00014"></a><span class="lineno"> 14</span>  </div>
|
|
<div class="line"><a name="l00015"></a><span class="lineno"> 15</span> <span class="preprocessor">#include <FairMQSocket.h></span></div>
|
|
<div class="line"><a name="l00016"></a><span class="lineno"> 16</span> <span class="preprocessor">#include <FairMQMessage.h></span></div>
|
|
<div class="line"><a name="l00017"></a><span class="lineno"> 17</span> <span class="preprocessor">#include <FairMQLogger.h></span></div>
|
|
<div class="line"><a name="l00018"></a><span class="lineno"> 18</span> <span class="preprocessor">#include <fairmq/tools/Strings.h></span></div>
|
|
<div class="line"><a name="l00019"></a><span class="lineno"> 19</span>  </div>
|
|
<div class="line"><a name="l00020"></a><span class="lineno"> 20</span> <span class="preprocessor">#include <zmq.h></span></div>
|
|
<div class="line"><a name="l00021"></a><span class="lineno"> 21</span>  </div>
|
|
<div class="line"><a name="l00022"></a><span class="lineno"> 22</span> <span class="preprocessor">#include <atomic></span></div>
|
|
<div class="line"><a name="l00023"></a><span class="lineno"> 23</span> <span class="preprocessor">#include <memory></span> <span class="comment">// make_unique</span></div>
|
|
<div class="line"><a name="l00024"></a><span class="lineno"> 24</span>  </div>
|
|
<div class="line"><a name="l00025"></a><span class="lineno"> 25</span> <span class="keyword">class </span><a class="code" href="classFairMQTransportFactory.html">FairMQTransportFactory</a>;</div>
|
|
<div class="line"><a name="l00026"></a><span class="lineno"> 26</span>  </div>
|
|
<div class="line"><a name="l00027"></a><span class="lineno"> 27</span> <span class="keyword">namespace </span><a class="code" href="namespacefair_1_1mq_1_1shmem.html">fair::mq::shmem</a></div>
|
|
<div class="line"><a name="l00028"></a><span class="lineno"> 28</span> {</div>
|
|
<div class="line"><a name="l00029"></a><span class="lineno"> 29</span>  </div>
|
|
<div class="line"><a name="l00030"></a><span class="lineno"><a class="line" href="structfair_1_1mq_1_1shmem_1_1ZMsg.html"> 30</a></span> <span class="keyword">struct </span><a class="code" href="structfair_1_1mq_1_1shmem_1_1ZMsg.html">ZMsg</a></div>
|
|
<div class="line"><a name="l00031"></a><span class="lineno"> 31</span> {</div>
|
|
<div class="line"><a name="l00032"></a><span class="lineno"> 32</span>  <a class="code" href="structfair_1_1mq_1_1shmem_1_1ZMsg.html">ZMsg</a>() { <span class="keywordtype">int</span> rc __attribute__((unused)) = zmq_msg_init(&fMsg); assert(rc == 0); }</div>
|
|
<div class="line"><a name="l00033"></a><span class="lineno"> 33</span>  <span class="keyword">explicit</span> <a class="code" href="structfair_1_1mq_1_1shmem_1_1ZMsg.html">ZMsg</a>(<span class="keywordtype">size_t</span> size) { <span class="keywordtype">int</span> rc __attribute__((unused)) = zmq_msg_init_size(&fMsg, size); assert(rc == 0); }</div>
|
|
<div class="line"><a name="l00034"></a><span class="lineno"> 34</span>  ~<a class="code" href="structfair_1_1mq_1_1shmem_1_1ZMsg.html">ZMsg</a>() { <span class="keywordtype">int</span> rc __attribute__((unused)) = zmq_msg_close(&fMsg); assert(rc == 0); }</div>
|
|
<div class="line"><a name="l00035"></a><span class="lineno"> 35</span>  </div>
|
|
<div class="line"><a name="l00036"></a><span class="lineno"> 36</span>  <span class="keywordtype">void</span>* Data() { <span class="keywordflow">return</span> zmq_msg_data(&fMsg); }</div>
|
|
<div class="line"><a name="l00037"></a><span class="lineno"> 37</span>  <span class="keywordtype">size_t</span> Size() { <span class="keywordflow">return</span> zmq_msg_size(&fMsg); }</div>
|
|
<div class="line"><a name="l00038"></a><span class="lineno"> 38</span>  zmq_msg_t* Msg() { <span class="keywordflow">return</span> &fMsg; }</div>
|
|
<div class="line"><a name="l00039"></a><span class="lineno"> 39</span>  </div>
|
|
<div class="line"><a name="l00040"></a><span class="lineno"> 40</span>  zmq_msg_t fMsg;</div>
|
|
<div class="line"><a name="l00041"></a><span class="lineno"> 41</span> };</div>
|
|
<div class="line"><a name="l00042"></a><span class="lineno"> 42</span>  </div>
|
|
<div class="line"><a name="l00043"></a><span class="lineno"><a class="line" href="classfair_1_1mq_1_1shmem_1_1Socket.html"> 43</a></span> <span class="keyword">class </span><a class="code" href="classfair_1_1mq_1_1shmem_1_1Socket.html">Socket</a> final : <span class="keyword">public</span> <a class="code" href="classFairMQSocket.html">fair::mq::Socket</a></div>
|
|
<div class="line"><a name="l00044"></a><span class="lineno"> 44</span> {</div>
|
|
<div class="line"><a name="l00045"></a><span class="lineno"> 45</span>  <span class="keyword">public</span>:</div>
|
|
<div class="line"><a name="l00046"></a><span class="lineno"> 46</span>  <a class="code" href="classfair_1_1mq_1_1shmem_1_1Socket.html">Socket</a>(<a class="code" href="classfair_1_1mq_1_1shmem_1_1Manager.html">Manager</a>& manager, <span class="keyword">const</span> std::string& type, <span class="keyword">const</span> std::string& name, <span class="keyword">const</span> std::string& <span class="keywordtype">id</span>, <span class="keywordtype">void</span>* context, <a class="code" href="classFairMQTransportFactory.html">FairMQTransportFactory</a>* fac = <span class="keyword">nullptr</span>)</div>
|
|
<div class="line"><a name="l00047"></a><span class="lineno"> 47</span>  : <a class="code" href="classFairMQSocket.html">fair::mq::Socket</a>(fac)</div>
|
|
<div class="line"><a name="l00048"></a><span class="lineno"> 48</span>  , fSocket(<span class="keyword">nullptr</span>)</div>
|
|
<div class="line"><a name="l00049"></a><span class="lineno"> 49</span>  , fManager(manager)</div>
|
|
<div class="line"><a name="l00050"></a><span class="lineno"> 50</span>  , fId(<span class="keywordtype">id</span> + <span class="stringliteral">"."</span> + name + <span class="stringliteral">"."</span> + type)</div>
|
|
<div class="line"><a name="l00051"></a><span class="lineno"> 51</span>  , fBytesTx(0)</div>
|
|
<div class="line"><a name="l00052"></a><span class="lineno"> 52</span>  , fBytesRx(0)</div>
|
|
<div class="line"><a name="l00053"></a><span class="lineno"> 53</span>  , fMessagesTx(0)</div>
|
|
<div class="line"><a name="l00054"></a><span class="lineno"> 54</span>  , fMessagesRx(0)</div>
|
|
<div class="line"><a name="l00055"></a><span class="lineno"> 55</span>  , fTimeout(100)</div>
|
|
<div class="line"><a name="l00056"></a><span class="lineno"> 56</span>  {</div>
|
|
<div class="line"><a name="l00057"></a><span class="lineno"> 57</span>  assert(context);</div>
|
|
<div class="line"><a name="l00058"></a><span class="lineno"> 58</span>  </div>
|
|
<div class="line"><a name="l00059"></a><span class="lineno"> 59</span>  <span class="keywordflow">if</span> (type == <span class="stringliteral">"sub"</span> || type == <span class="stringliteral">"pub"</span>) {</div>
|
|
<div class="line"><a name="l00060"></a><span class="lineno"> 60</span>  LOG(error) << <span class="stringliteral">"PUB/SUB socket type is not supported for shared memory transport"</span>;</div>
|
|
<div class="line"><a name="l00061"></a><span class="lineno"> 61</span>  <span class="keywordflow">throw</span> <a class="code" href="structfair_1_1mq_1_1SocketError.html">SocketError</a>(<span class="stringliteral">"PUB/SUB socket type is not supported for shared memory transport"</span>);</div>
|
|
<div class="line"><a name="l00062"></a><span class="lineno"> 62</span>  }</div>
|
|
<div class="line"><a name="l00063"></a><span class="lineno"> 63</span>  </div>
|
|
<div class="line"><a name="l00064"></a><span class="lineno"> 64</span>  fSocket = zmq_socket(context, GetConstant(type));</div>
|
|
<div class="line"><a name="l00065"></a><span class="lineno"> 65</span>  </div>
|
|
<div class="line"><a name="l00066"></a><span class="lineno"> 66</span>  <span class="keywordflow">if</span> (fSocket == <span class="keyword">nullptr</span>) {</div>
|
|
<div class="line"><a name="l00067"></a><span class="lineno"> 67</span>  LOG(error) << <span class="stringliteral">"Failed creating socket "</span> << fId << <span class="stringliteral">", reason: "</span> << zmq_strerror(errno);</div>
|
|
<div class="line"><a name="l00068"></a><span class="lineno"> 68</span>  <span class="keywordflow">throw</span> <a class="code" href="structfair_1_1mq_1_1SocketError.html">SocketError</a>(tools::ToString(<span class="stringliteral">"Failed creating socket "</span>, fId, <span class="stringliteral">", reason: "</span>, zmq_strerror(errno)));</div>
|
|
<div class="line"><a name="l00069"></a><span class="lineno"> 69</span>  }</div>
|
|
<div class="line"><a name="l00070"></a><span class="lineno"> 70</span>  </div>
|
|
<div class="line"><a name="l00071"></a><span class="lineno"> 71</span>  <span class="keywordflow">if</span> (zmq_setsockopt(fSocket, ZMQ_IDENTITY, fId.c_str(), fId.length()) != 0) {</div>
|
|
<div class="line"><a name="l00072"></a><span class="lineno"> 72</span>  LOG(error) << <span class="stringliteral">"Failed setting ZMQ_IDENTITY socket option, reason: "</span> << zmq_strerror(errno);</div>
|
|
<div class="line"><a name="l00073"></a><span class="lineno"> 73</span>  }</div>
|
|
<div class="line"><a name="l00074"></a><span class="lineno"> 74</span>  </div>
|
|
<div class="line"><a name="l00075"></a><span class="lineno"> 75</span>  <span class="comment">// Tell socket to try and send/receive outstanding messages for <linger> milliseconds before terminating.</span></div>
|
|
<div class="line"><a name="l00076"></a><span class="lineno"> 76</span>  <span class="comment">// Default value for ZeroMQ is -1, which is to wait forever.</span></div>
|
|
<div class="line"><a name="l00077"></a><span class="lineno"> 77</span>  <span class="keywordtype">int</span> linger = 1000;</div>
|
|
<div class="line"><a name="l00078"></a><span class="lineno"> 78</span>  <span class="keywordflow">if</span> (zmq_setsockopt(fSocket, ZMQ_LINGER, &linger, <span class="keyword">sizeof</span>(linger)) != 0) {</div>
|
|
<div class="line"><a name="l00079"></a><span class="lineno"> 79</span>  LOG(error) << <span class="stringliteral">"Failed setting ZMQ_LINGER socket option, reason: "</span> << zmq_strerror(errno);</div>
|
|
<div class="line"><a name="l00080"></a><span class="lineno"> 80</span>  }</div>
|
|
<div class="line"><a name="l00081"></a><span class="lineno"> 81</span>  </div>
|
|
<div class="line"><a name="l00082"></a><span class="lineno"> 82</span>  <span class="keywordflow">if</span> (zmq_setsockopt(fSocket, ZMQ_SNDTIMEO, &fTimeout, <span class="keyword">sizeof</span>(fTimeout)) != 0) {</div>
|
|
<div class="line"><a name="l00083"></a><span class="lineno"> 83</span>  LOG(error) << <span class="stringliteral">"Failed setting ZMQ_SNDTIMEO socket option, reason: "</span> << zmq_strerror(errno);</div>
|
|
<div class="line"><a name="l00084"></a><span class="lineno"> 84</span>  }</div>
|
|
<div class="line"><a name="l00085"></a><span class="lineno"> 85</span>  </div>
|
|
<div class="line"><a name="l00086"></a><span class="lineno"> 86</span>  <span class="keywordflow">if</span> (zmq_setsockopt(fSocket, ZMQ_RCVTIMEO, &fTimeout, <span class="keyword">sizeof</span>(fTimeout)) != 0) {</div>
|
|
<div class="line"><a name="l00087"></a><span class="lineno"> 87</span>  LOG(error) << <span class="stringliteral">"Failed setting ZMQ_RCVTIMEO socket option, reason: "</span> << zmq_strerror(errno);</div>
|
|
<div class="line"><a name="l00088"></a><span class="lineno"> 88</span>  }</div>
|
|
<div class="line"><a name="l00089"></a><span class="lineno"> 89</span>  </div>
|
|
<div class="line"><a name="l00090"></a><span class="lineno"> 90</span>  <span class="comment">// if (type == "sub")</span></div>
|
|
<div class="line"><a name="l00091"></a><span class="lineno"> 91</span>  <span class="comment">// {</span></div>
|
|
<div class="line"><a name="l00092"></a><span class="lineno"> 92</span>  <span class="comment">// if (zmq_setsockopt(fSocket, ZMQ_SUBSCRIBE, nullptr, 0) != 0)</span></div>
|
|
<div class="line"><a name="l00093"></a><span class="lineno"> 93</span>  <span class="comment">// {</span></div>
|
|
<div class="line"><a name="l00094"></a><span class="lineno"> 94</span>  <span class="comment">// LOG(error) << "Failed setting ZMQ_SUBSCRIBE socket option, reason: " << zmq_strerror(errno);</span></div>
|
|
<div class="line"><a name="l00095"></a><span class="lineno"> 95</span>  <span class="comment">// }</span></div>
|
|
<div class="line"><a name="l00096"></a><span class="lineno"> 96</span>  <span class="comment">// }</span></div>
|
|
<div class="line"><a name="l00097"></a><span class="lineno"> 97</span>  LOG(debug) << <span class="stringliteral">"Created socket "</span> << GetId();</div>
|
|
<div class="line"><a name="l00098"></a><span class="lineno"> 98</span>  }</div>
|
|
<div class="line"><a name="l00099"></a><span class="lineno"> 99</span>  </div>
|
|
<div class="line"><a name="l00100"></a><span class="lineno"> 100</span>  <a class="code" href="classfair_1_1mq_1_1shmem_1_1Socket.html">Socket</a>(<span class="keyword">const</span> <a class="code" href="classfair_1_1mq_1_1shmem_1_1Socket.html">Socket</a>&) = <span class="keyword">delete</span>;</div>
|
|
<div class="line"><a name="l00101"></a><span class="lineno"> 101</span>  <a class="code" href="classfair_1_1mq_1_1shmem_1_1Socket.html">Socket</a> operator=(<span class="keyword">const</span> <a class="code" href="classfair_1_1mq_1_1shmem_1_1Socket.html">Socket</a>&) = <span class="keyword">delete</span>;</div>
|
|
<div class="line"><a name="l00102"></a><span class="lineno"> 102</span>  </div>
|
|
<div class="line"><a name="l00103"></a><span class="lineno"> 103</span>  std::string GetId()<span class="keyword"> const override </span>{ <span class="keywordflow">return</span> fId; }</div>
|
|
<div class="line"><a name="l00104"></a><span class="lineno"> 104</span>  </div>
|
|
<div class="line"><a name="l00105"></a><span class="lineno"> 105</span>  <span class="keywordtype">bool</span> Bind(<span class="keyword">const</span> std::string& address)<span class="keyword"> override</span></div>
|
|
<div class="line"><a name="l00106"></a><span class="lineno"> 106</span> <span class="keyword"> </span>{</div>
|
|
<div class="line"><a name="l00107"></a><span class="lineno"> 107</span>  <span class="comment">// LOG(info) << "binding socket " << fId << " on " << address;</span></div>
|
|
<div class="line"><a name="l00108"></a><span class="lineno"> 108</span>  <span class="keywordflow">if</span> (zmq_bind(fSocket, address.c_str()) != 0) {</div>
|
|
<div class="line"><a name="l00109"></a><span class="lineno"> 109</span>  <span class="keywordflow">if</span> (errno == EADDRINUSE) {</div>
|
|
<div class="line"><a name="l00110"></a><span class="lineno"> 110</span>  <span class="comment">// do not print error in this case, this is handled by FairMQDevice in case no connection could be established after trying a number of random ports from a range.</span></div>
|
|
<div class="line"><a name="l00111"></a><span class="lineno"> 111</span>  <span class="keywordflow">return</span> <span class="keyword">false</span>;</div>
|
|
<div class="line"><a name="l00112"></a><span class="lineno"> 112</span>  }</div>
|
|
<div class="line"><a name="l00113"></a><span class="lineno"> 113</span>  LOG(error) << <span class="stringliteral">"Failed binding socket "</span> << fId << <span class="stringliteral">", reason: "</span> << zmq_strerror(errno);</div>
|
|
<div class="line"><a name="l00114"></a><span class="lineno"> 114</span>  <span class="keywordflow">return</span> <span class="keyword">false</span>;</div>
|
|
<div class="line"><a name="l00115"></a><span class="lineno"> 115</span>  }</div>
|
|
<div class="line"><a name="l00116"></a><span class="lineno"> 116</span>  <span class="keywordflow">return</span> <span class="keyword">true</span>;</div>
|
|
<div class="line"><a name="l00117"></a><span class="lineno"> 117</span>  }</div>
|
|
<div class="line"><a name="l00118"></a><span class="lineno"> 118</span>  </div>
|
|
<div class="line"><a name="l00119"></a><span class="lineno"> 119</span>  <span class="keywordtype">bool</span> Connect(<span class="keyword">const</span> std::string& address)<span class="keyword"> override</span></div>
|
|
<div class="line"><a name="l00120"></a><span class="lineno"> 120</span> <span class="keyword"> </span>{</div>
|
|
<div class="line"><a name="l00121"></a><span class="lineno"> 121</span>  <span class="comment">// LOG(info) << "connecting socket " << fId << " on " << address;</span></div>
|
|
<div class="line"><a name="l00122"></a><span class="lineno"> 122</span>  <span class="keywordflow">if</span> (zmq_connect(fSocket, address.c_str()) != 0) {</div>
|
|
<div class="line"><a name="l00123"></a><span class="lineno"> 123</span>  LOG(error) << <span class="stringliteral">"Failed connecting socket "</span> << fId << <span class="stringliteral">", reason: "</span> << zmq_strerror(errno);</div>
|
|
<div class="line"><a name="l00124"></a><span class="lineno"> 124</span>  <span class="keywordflow">return</span> <span class="keyword">false</span>;</div>
|
|
<div class="line"><a name="l00125"></a><span class="lineno"> 125</span>  }</div>
|
|
<div class="line"><a name="l00126"></a><span class="lineno"> 126</span>  <span class="keywordflow">return</span> <span class="keyword">true</span>;</div>
|
|
<div class="line"><a name="l00127"></a><span class="lineno"> 127</span>  }</div>
|
|
<div class="line"><a name="l00128"></a><span class="lineno"> 128</span>  </div>
|
|
<div class="line"><a name="l00129"></a><span class="lineno"> 129</span>  <span class="keywordtype">bool</span> ShouldRetry(<span class="keywordtype">int</span> flags, <span class="keywordtype">int</span> timeout, <span class="keywordtype">int</span>& elapsed)<span class="keyword"> const</span></div>
|
|
<div class="line"><a name="l00130"></a><span class="lineno"> 130</span> <span class="keyword"> </span>{</div>
|
|
<div class="line"><a name="l00131"></a><span class="lineno"> 131</span>  <span class="keywordflow">if</span> ((flags & ZMQ_DONTWAIT) == 0) {</div>
|
|
<div class="line"><a name="l00132"></a><span class="lineno"> 132</span>  <span class="keywordflow">if</span> (timeout > 0) {</div>
|
|
<div class="line"><a name="l00133"></a><span class="lineno"> 133</span>  elapsed += fTimeout;</div>
|
|
<div class="line"><a name="l00134"></a><span class="lineno"> 134</span>  <span class="keywordflow">if</span> (elapsed >= timeout) {</div>
|
|
<div class="line"><a name="l00135"></a><span class="lineno"> 135</span>  <span class="keywordflow">return</span> <span class="keyword">false</span>;</div>
|
|
<div class="line"><a name="l00136"></a><span class="lineno"> 136</span>  }</div>
|
|
<div class="line"><a name="l00137"></a><span class="lineno"> 137</span>  }</div>
|
|
<div class="line"><a name="l00138"></a><span class="lineno"> 138</span>  <span class="keywordflow">return</span> <span class="keyword">true</span>;</div>
|
|
<div class="line"><a name="l00139"></a><span class="lineno"> 139</span>  } <span class="keywordflow">else</span> {</div>
|
|
<div class="line"><a name="l00140"></a><span class="lineno"> 140</span>  <span class="keywordflow">return</span> <span class="keyword">false</span>;</div>
|
|
<div class="line"><a name="l00141"></a><span class="lineno"> 141</span>  }</div>
|
|
<div class="line"><a name="l00142"></a><span class="lineno"> 142</span>  }</div>
|
|
<div class="line"><a name="l00143"></a><span class="lineno"> 143</span>  </div>
|
|
<div class="line"><a name="l00144"></a><span class="lineno"> 144</span>  <span class="keywordtype">int</span> HandleErrors()<span class="keyword"> const</span></div>
|
|
<div class="line"><a name="l00145"></a><span class="lineno"> 145</span> <span class="keyword"> </span>{</div>
|
|
<div class="line"><a name="l00146"></a><span class="lineno"> 146</span>  <span class="keywordflow">if</span> (zmq_errno() == ETERM) {</div>
|
|
<div class="line"><a name="l00147"></a><span class="lineno"> 147</span>  LOG(debug) << <span class="stringliteral">"Terminating socket "</span> << fId;</div>
|
|
<div class="line"><a name="l00148"></a><span class="lineno"> 148</span>  <span class="keywordflow">return</span> <span class="keyword">static_cast<</span><span class="keywordtype">int</span><span class="keyword">></span>(TransferCode::error);</div>
|
|
<div class="line"><a name="l00149"></a><span class="lineno"> 149</span>  } <span class="keywordflow">else</span> {</div>
|
|
<div class="line"><a name="l00150"></a><span class="lineno"> 150</span>  LOG(error) << <span class="stringliteral">"Failed transfer on socket "</span> << fId << <span class="stringliteral">", reason: "</span> << zmq_strerror(errno);</div>
|
|
<div class="line"><a name="l00151"></a><span class="lineno"> 151</span>  <span class="keywordflow">return</span> <span class="keyword">static_cast<</span><span class="keywordtype">int</span><span class="keyword">></span>(TransferCode::error);</div>
|
|
<div class="line"><a name="l00152"></a><span class="lineno"> 152</span>  }</div>
|
|
<div class="line"><a name="l00153"></a><span class="lineno"> 153</span>  }</div>
|
|
<div class="line"><a name="l00154"></a><span class="lineno"> 154</span>  </div>
|
|
<div class="line"><a name="l00155"></a><span class="lineno"> 155</span>  int64_t Send(MessagePtr& msg, <span class="keyword">const</span> <span class="keywordtype">int</span> timeout = -1)<span class="keyword"> override</span></div>
|
|
<div class="line"><a name="l00156"></a><span class="lineno"> 156</span> <span class="keyword"> </span>{</div>
|
|
<div class="line"><a name="l00157"></a><span class="lineno"> 157</span>  <span class="keywordtype">int</span> flags = 0;</div>
|
|
<div class="line"><a name="l00158"></a><span class="lineno"> 158</span>  <span class="keywordflow">if</span> (timeout == 0) {</div>
|
|
<div class="line"><a name="l00159"></a><span class="lineno"> 159</span>  flags = ZMQ_DONTWAIT;</div>
|
|
<div class="line"><a name="l00160"></a><span class="lineno"> 160</span>  }</div>
|
|
<div class="line"><a name="l00161"></a><span class="lineno"> 161</span>  <span class="keywordtype">int</span> elapsed = 0;</div>
|
|
<div class="line"><a name="l00162"></a><span class="lineno"> 162</span>  </div>
|
|
<div class="line"><a name="l00163"></a><span class="lineno"> 163</span>  <a class="code" href="classfair_1_1mq_1_1shmem_1_1Message.html">Message</a>* shmMsg = <span class="keyword">static_cast<</span><a class="code" href="classfair_1_1mq_1_1shmem_1_1Message.html">Message</a>*<span class="keyword">></span>(msg.get());</div>
|
|
<div class="line"><a name="l00164"></a><span class="lineno"> 164</span>  <a class="code" href="structfair_1_1mq_1_1shmem_1_1ZMsg.html">ZMsg</a> zmqMsg(<span class="keyword">sizeof</span>(<a class="code" href="structfair_1_1mq_1_1shmem_1_1MetaHeader.html">MetaHeader</a>));</div>
|
|
<div class="line"><a name="l00165"></a><span class="lineno"> 165</span>  std::memcpy(zmqMsg.Data(), &(shmMsg->fMeta), <span class="keyword">sizeof</span>(<a class="code" href="structfair_1_1mq_1_1shmem_1_1MetaHeader.html">MetaHeader</a>));</div>
|
|
<div class="line"><a name="l00166"></a><span class="lineno"> 166</span>  </div>
|
|
<div class="line"><a name="l00167"></a><span class="lineno"> 167</span>  <span class="keywordflow">while</span> (<span class="keyword">true</span>) {</div>
|
|
<div class="line"><a name="l00168"></a><span class="lineno"> 168</span>  <span class="keywordtype">int</span> nbytes = zmq_msg_send(zmqMsg.Msg(), fSocket, flags);</div>
|
|
<div class="line"><a name="l00169"></a><span class="lineno"> 169</span>  <span class="keywordflow">if</span> (nbytes > 0) {</div>
|
|
<div class="line"><a name="l00170"></a><span class="lineno"> 170</span>  shmMsg->fQueued = <span class="keyword">true</span>;</div>
|
|
<div class="line"><a name="l00171"></a><span class="lineno"> 171</span>  ++fMessagesTx;</div>
|
|
<div class="line"><a name="l00172"></a><span class="lineno"> 172</span>  <span class="keywordtype">size_t</span> size = msg->GetSize();</div>
|
|
<div class="line"><a name="l00173"></a><span class="lineno"> 173</span>  fBytesTx += size;</div>
|
|
<div class="line"><a name="l00174"></a><span class="lineno"> 174</span>  <span class="keywordflow">return</span> size;</div>
|
|
<div class="line"><a name="l00175"></a><span class="lineno"> 175</span>  } <span class="keywordflow">else</span> <span class="keywordflow">if</span> (zmq_errno() == EAGAIN || zmq_errno() == EINTR) {</div>
|
|
<div class="line"><a name="l00176"></a><span class="lineno"> 176</span>  <span class="keywordflow">if</span> (fManager.Interrupted()) {</div>
|
|
<div class="line"><a name="l00177"></a><span class="lineno"> 177</span>  <span class="keywordflow">return</span> <span class="keyword">static_cast<</span><span class="keywordtype">int</span><span class="keyword">></span>(TransferCode::interrupted);</div>
|
|
<div class="line"><a name="l00178"></a><span class="lineno"> 178</span>  } <span class="keywordflow">else</span> <span class="keywordflow">if</span> (ShouldRetry(flags, timeout, elapsed)) {</div>
|
|
<div class="line"><a name="l00179"></a><span class="lineno"> 179</span>  <span class="keywordflow">continue</span>;</div>
|
|
<div class="line"><a name="l00180"></a><span class="lineno"> 180</span>  } <span class="keywordflow">else</span> {</div>
|
|
<div class="line"><a name="l00181"></a><span class="lineno"> 181</span>  <span class="keywordflow">return</span> <span class="keyword">static_cast<</span><span class="keywordtype">int</span><span class="keyword">></span>(TransferCode::timeout);</div>
|
|
<div class="line"><a name="l00182"></a><span class="lineno"> 182</span>  }</div>
|
|
<div class="line"><a name="l00183"></a><span class="lineno"> 183</span>  } <span class="keywordflow">else</span> {</div>
|
|
<div class="line"><a name="l00184"></a><span class="lineno"> 184</span>  <span class="keywordflow">return</span> HandleErrors();</div>
|
|
<div class="line"><a name="l00185"></a><span class="lineno"> 185</span>  }</div>
|
|
<div class="line"><a name="l00186"></a><span class="lineno"> 186</span>  }</div>
|
|
<div class="line"><a name="l00187"></a><span class="lineno"> 187</span>  </div>
|
|
<div class="line"><a name="l00188"></a><span class="lineno"> 188</span>  <span class="keywordflow">return</span> <span class="keyword">static_cast<</span><span class="keywordtype">int</span><span class="keyword">></span>(TransferCode::error);</div>
|
|
<div class="line"><a name="l00189"></a><span class="lineno"> 189</span>  }</div>
|
|
<div class="line"><a name="l00190"></a><span class="lineno"> 190</span>  </div>
|
|
<div class="line"><a name="l00191"></a><span class="lineno"> 191</span>  int64_t Receive(MessagePtr& msg, <span class="keyword">const</span> <span class="keywordtype">int</span> timeout = -1)<span class="keyword"> override</span></div>
|
|
<div class="line"><a name="l00192"></a><span class="lineno"> 192</span> <span class="keyword"> </span>{</div>
|
|
<div class="line"><a name="l00193"></a><span class="lineno"> 193</span>  <span class="keywordtype">int</span> flags = 0;</div>
|
|
<div class="line"><a name="l00194"></a><span class="lineno"> 194</span>  <span class="keywordflow">if</span> (timeout == 0) {</div>
|
|
<div class="line"><a name="l00195"></a><span class="lineno"> 195</span>  flags = ZMQ_DONTWAIT;</div>
|
|
<div class="line"><a name="l00196"></a><span class="lineno"> 196</span>  }</div>
|
|
<div class="line"><a name="l00197"></a><span class="lineno"> 197</span>  <span class="keywordtype">int</span> elapsed = 0;</div>
|
|
<div class="line"><a name="l00198"></a><span class="lineno"> 198</span>  </div>
|
|
<div class="line"><a name="l00199"></a><span class="lineno"> 199</span>  <a class="code" href="structfair_1_1mq_1_1shmem_1_1ZMsg.html">ZMsg</a> zmqMsg;</div>
|
|
<div class="line"><a name="l00200"></a><span class="lineno"> 200</span>  </div>
|
|
<div class="line"><a name="l00201"></a><span class="lineno"> 201</span>  <span class="keywordflow">while</span> (<span class="keyword">true</span>) {</div>
|
|
<div class="line"><a name="l00202"></a><span class="lineno"> 202</span>  <a class="code" href="classfair_1_1mq_1_1shmem_1_1Message.html">Message</a>* shmMsg = <span class="keyword">static_cast<</span><a class="code" href="classfair_1_1mq_1_1shmem_1_1Message.html">Message</a>*<span class="keyword">></span>(msg.get());</div>
|
|
<div class="line"><a name="l00203"></a><span class="lineno"> 203</span>  <span class="keywordtype">int</span> nbytes = zmq_msg_recv(zmqMsg.Msg(), fSocket, flags);</div>
|
|
<div class="line"><a name="l00204"></a><span class="lineno"> 204</span>  <span class="keywordflow">if</span> (nbytes > 0) {</div>
|
|
<div class="line"><a name="l00205"></a><span class="lineno"> 205</span>  <span class="comment">// check for number of received messages. must be 1</span></div>
|
|
<div class="line"><a name="l00206"></a><span class="lineno"> 206</span>  <span class="keywordflow">if</span> (nbytes != <span class="keyword">sizeof</span>(<a class="code" href="structfair_1_1mq_1_1shmem_1_1MetaHeader.html">MetaHeader</a>)) {</div>
|
|
<div class="line"><a name="l00207"></a><span class="lineno"> 207</span>  <span class="keywordflow">throw</span> <a class="code" href="structfair_1_1mq_1_1SocketError.html">SocketError</a>(</div>
|
|
<div class="line"><a name="l00208"></a><span class="lineno"> 208</span>  tools::ToString(<span class="stringliteral">"Received message is not a valid FairMQ shared memory message. "</span>,</div>
|
|
<div class="line"><a name="l00209"></a><span class="lineno"> 209</span>  <span class="stringliteral">"Possibly due to a misconfigured transport on the sender side. "</span>,</div>
|
|
<div class="line"><a name="l00210"></a><span class="lineno"> 210</span>  <span class="stringliteral">"Expected size of "</span>, <span class="keyword">sizeof</span>(<a class="code" href="structfair_1_1mq_1_1shmem_1_1MetaHeader.html">MetaHeader</a>), <span class="stringliteral">" bytes, received "</span>, nbytes));</div>
|
|
<div class="line"><a name="l00211"></a><span class="lineno"> 211</span>  }</div>
|
|
<div class="line"><a name="l00212"></a><span class="lineno"> 212</span>  </div>
|
|
<div class="line"><a name="l00213"></a><span class="lineno"> 213</span>  <a class="code" href="structfair_1_1mq_1_1shmem_1_1MetaHeader.html">MetaHeader</a>* hdr = <span class="keyword">static_cast<</span><a class="code" href="structfair_1_1mq_1_1shmem_1_1MetaHeader.html">MetaHeader</a>*<span class="keyword">></span>(zmqMsg.Data());</div>
|
|
<div class="line"><a name="l00214"></a><span class="lineno"> 214</span>  <span class="keywordtype">size_t</span> size = hdr->fSize;</div>
|
|
<div class="line"><a name="l00215"></a><span class="lineno"> 215</span>  shmMsg->fMeta = *hdr;</div>
|
|
<div class="line"><a name="l00216"></a><span class="lineno"> 216</span>  </div>
|
|
<div class="line"><a name="l00217"></a><span class="lineno"> 217</span>  fBytesRx += size;</div>
|
|
<div class="line"><a name="l00218"></a><span class="lineno"> 218</span>  ++fMessagesRx;</div>
|
|
<div class="line"><a name="l00219"></a><span class="lineno"> 219</span>  <span class="keywordflow">return</span> size;</div>
|
|
<div class="line"><a name="l00220"></a><span class="lineno"> 220</span>  } <span class="keywordflow">else</span> <span class="keywordflow">if</span> (zmq_errno() == EAGAIN || zmq_errno() == EINTR) {</div>
|
|
<div class="line"><a name="l00221"></a><span class="lineno"> 221</span>  <span class="keywordflow">if</span> (fManager.Interrupted()) {</div>
|
|
<div class="line"><a name="l00222"></a><span class="lineno"> 222</span>  <span class="keywordflow">return</span> <span class="keyword">static_cast<</span><span class="keywordtype">int</span><span class="keyword">></span>(TransferCode::interrupted);</div>
|
|
<div class="line"><a name="l00223"></a><span class="lineno"> 223</span>  } <span class="keywordflow">else</span> <span class="keywordflow">if</span> (ShouldRetry(flags, timeout, elapsed)) {</div>
|
|
<div class="line"><a name="l00224"></a><span class="lineno"> 224</span>  <span class="keywordflow">continue</span>;</div>
|
|
<div class="line"><a name="l00225"></a><span class="lineno"> 225</span>  } <span class="keywordflow">else</span> {</div>
|
|
<div class="line"><a name="l00226"></a><span class="lineno"> 226</span>  <span class="keywordflow">return</span> <span class="keyword">static_cast<</span><span class="keywordtype">int</span><span class="keyword">></span>(TransferCode::timeout);</div>
|
|
<div class="line"><a name="l00227"></a><span class="lineno"> 227</span>  }</div>
|
|
<div class="line"><a name="l00228"></a><span class="lineno"> 228</span>  } <span class="keywordflow">else</span> {</div>
|
|
<div class="line"><a name="l00229"></a><span class="lineno"> 229</span>  <span class="keywordflow">return</span> HandleErrors();</div>
|
|
<div class="line"><a name="l00230"></a><span class="lineno"> 230</span>  }</div>
|
|
<div class="line"><a name="l00231"></a><span class="lineno"> 231</span>  }</div>
|
|
<div class="line"><a name="l00232"></a><span class="lineno"> 232</span>  }</div>
|
|
<div class="line"><a name="l00233"></a><span class="lineno"> 233</span>  </div>
|
|
<div class="line"><a name="l00234"></a><span class="lineno"> 234</span>  int64_t Send(std::vector<MessagePtr>& msgVec, <span class="keyword">const</span> <span class="keywordtype">int</span> timeout = -1)<span class="keyword"> override</span></div>
|
|
<div class="line"><a name="l00235"></a><span class="lineno"> 235</span> <span class="keyword"> </span>{</div>
|
|
<div class="line"><a name="l00236"></a><span class="lineno"> 236</span>  <span class="keywordtype">int</span> flags = 0;</div>
|
|
<div class="line"><a name="l00237"></a><span class="lineno"> 237</span>  <span class="keywordflow">if</span> (timeout == 0) {</div>
|
|
<div class="line"><a name="l00238"></a><span class="lineno"> 238</span>  flags = ZMQ_DONTWAIT;</div>
|
|
<div class="line"><a name="l00239"></a><span class="lineno"> 239</span>  }</div>
|
|
<div class="line"><a name="l00240"></a><span class="lineno"> 240</span>  <span class="keywordtype">int</span> elapsed = 0;</div>
|
|
<div class="line"><a name="l00241"></a><span class="lineno"> 241</span>  </div>
|
|
<div class="line"><a name="l00242"></a><span class="lineno"> 242</span>  <span class="comment">// put it into zmq message</span></div>
|
|
<div class="line"><a name="l00243"></a><span class="lineno"> 243</span>  <span class="keyword">const</span> <span class="keywordtype">unsigned</span> <span class="keywordtype">int</span> vecSize = msgVec.size();</div>
|
|
<div class="line"><a name="l00244"></a><span class="lineno"> 244</span>  <a class="code" href="structfair_1_1mq_1_1shmem_1_1ZMsg.html">ZMsg</a> zmqMsg(vecSize * <span class="keyword">sizeof</span>(<a class="code" href="structfair_1_1mq_1_1shmem_1_1MetaHeader.html">MetaHeader</a>));</div>
|
|
<div class="line"><a name="l00245"></a><span class="lineno"> 245</span>  </div>
|
|
<div class="line"><a name="l00246"></a><span class="lineno"> 246</span>  <span class="comment">// prepare the message with shm metas</span></div>
|
|
<div class="line"><a name="l00247"></a><span class="lineno"> 247</span>  <a class="code" href="structfair_1_1mq_1_1shmem_1_1MetaHeader.html">MetaHeader</a>* metas = <span class="keyword">static_cast<</span><a class="code" href="structfair_1_1mq_1_1shmem_1_1MetaHeader.html">MetaHeader</a>*<span class="keyword">></span>(zmqMsg.Data());</div>
|
|
<div class="line"><a name="l00248"></a><span class="lineno"> 248</span>  </div>
|
|
<div class="line"><a name="l00249"></a><span class="lineno"> 249</span>  <span class="keywordflow">for</span> (<span class="keyword">auto</span>& msg : msgVec) {</div>
|
|
<div class="line"><a name="l00250"></a><span class="lineno"> 250</span>  <a class="code" href="classfair_1_1mq_1_1shmem_1_1Message.html">Message</a>* shmMsg = <span class="keyword">static_cast<</span><a class="code" href="classfair_1_1mq_1_1shmem_1_1Message.html">Message</a>*<span class="keyword">></span>(msg.get());</div>
|
|
<div class="line"><a name="l00251"></a><span class="lineno"> 251</span>  std::memcpy(metas++, &(shmMsg->fMeta), <span class="keyword">sizeof</span>(<a class="code" href="structfair_1_1mq_1_1shmem_1_1MetaHeader.html">MetaHeader</a>));</div>
|
|
<div class="line"><a name="l00252"></a><span class="lineno"> 252</span>  }</div>
|
|
<div class="line"><a name="l00253"></a><span class="lineno"> 253</span>  </div>
|
|
<div class="line"><a name="l00254"></a><span class="lineno"> 254</span>  <span class="keywordflow">while</span> (<span class="keyword">true</span>) {</div>
|
|
<div class="line"><a name="l00255"></a><span class="lineno"> 255</span>  int64_t totalSize = 0;</div>
|
|
<div class="line"><a name="l00256"></a><span class="lineno"> 256</span>  <span class="keywordtype">int</span> nbytes = zmq_msg_send(zmqMsg.Msg(), fSocket, flags);</div>
|
|
<div class="line"><a name="l00257"></a><span class="lineno"> 257</span>  <span class="keywordflow">if</span> (nbytes > 0) {</div>
|
|
<div class="line"><a name="l00258"></a><span class="lineno"> 258</span>  assert(<span class="keyword">static_cast<</span><span class="keywordtype">unsigned</span> <span class="keywordtype">int</span><span class="keyword">></span>(nbytes) == (vecSize * <span class="keyword">sizeof</span>(<a class="code" href="structfair_1_1mq_1_1shmem_1_1MetaHeader.html">MetaHeader</a>))); <span class="comment">// all or nothing</span></div>
|
|
<div class="line"><a name="l00259"></a><span class="lineno"> 259</span>  </div>
|
|
<div class="line"><a name="l00260"></a><span class="lineno"> 260</span>  <span class="keywordflow">for</span> (<span class="keyword">auto</span>& msg : msgVec) {</div>
|
|
<div class="line"><a name="l00261"></a><span class="lineno"> 261</span>  <a class="code" href="classfair_1_1mq_1_1shmem_1_1Message.html">Message</a>* shmMsg = <span class="keyword">static_cast<</span><a class="code" href="classfair_1_1mq_1_1shmem_1_1Message.html">Message</a>*<span class="keyword">></span>(msg.get());</div>
|
|
<div class="line"><a name="l00262"></a><span class="lineno"> 262</span>  shmMsg->fQueued = <span class="keyword">true</span>;</div>
|
|
<div class="line"><a name="l00263"></a><span class="lineno"> 263</span>  totalSize += shmMsg->fMeta.fSize;</div>
|
|
<div class="line"><a name="l00264"></a><span class="lineno"> 264</span>  }</div>
|
|
<div class="line"><a name="l00265"></a><span class="lineno"> 265</span>  </div>
|
|
<div class="line"><a name="l00266"></a><span class="lineno"> 266</span>  <span class="comment">// store statistics on how many messages have been sent</span></div>
|
|
<div class="line"><a name="l00267"></a><span class="lineno"> 267</span>  fMessagesTx++;</div>
|
|
<div class="line"><a name="l00268"></a><span class="lineno"> 268</span>  fBytesTx += totalSize;</div>
|
|
<div class="line"><a name="l00269"></a><span class="lineno"> 269</span>  </div>
|
|
<div class="line"><a name="l00270"></a><span class="lineno"> 270</span>  <span class="keywordflow">return</span> totalSize;</div>
|
|
<div class="line"><a name="l00271"></a><span class="lineno"> 271</span>  } <span class="keywordflow">else</span> <span class="keywordflow">if</span> (zmq_errno() == EAGAIN || zmq_errno() == EINTR) {</div>
|
|
<div class="line"><a name="l00272"></a><span class="lineno"> 272</span>  <span class="keywordflow">if</span> (fManager.Interrupted()) {</div>
|
|
<div class="line"><a name="l00273"></a><span class="lineno"> 273</span>  <span class="keywordflow">return</span> <span class="keyword">static_cast<</span><span class="keywordtype">int</span><span class="keyword">></span>(TransferCode::interrupted);</div>
|
|
<div class="line"><a name="l00274"></a><span class="lineno"> 274</span>  } <span class="keywordflow">else</span> <span class="keywordflow">if</span> (ShouldRetry(flags, timeout, elapsed)) {</div>
|
|
<div class="line"><a name="l00275"></a><span class="lineno"> 275</span>  <span class="keywordflow">continue</span>;</div>
|
|
<div class="line"><a name="l00276"></a><span class="lineno"> 276</span>  } <span class="keywordflow">else</span> {</div>
|
|
<div class="line"><a name="l00277"></a><span class="lineno"> 277</span>  <span class="keywordflow">return</span> <span class="keyword">static_cast<</span><span class="keywordtype">int</span><span class="keyword">></span>(TransferCode::timeout);</div>
|
|
<div class="line"><a name="l00278"></a><span class="lineno"> 278</span>  }</div>
|
|
<div class="line"><a name="l00279"></a><span class="lineno"> 279</span>  } <span class="keywordflow">else</span> {</div>
|
|
<div class="line"><a name="l00280"></a><span class="lineno"> 280</span>  <span class="keywordflow">return</span> HandleErrors();</div>
|
|
<div class="line"><a name="l00281"></a><span class="lineno"> 281</span>  }</div>
|
|
<div class="line"><a name="l00282"></a><span class="lineno"> 282</span>  }</div>
|
|
<div class="line"><a name="l00283"></a><span class="lineno"> 283</span>  </div>
|
|
<div class="line"><a name="l00284"></a><span class="lineno"> 284</span>  <span class="keywordflow">return</span> <span class="keyword">static_cast<</span><span class="keywordtype">int</span><span class="keyword">></span>(TransferCode::error);</div>
|
|
<div class="line"><a name="l00285"></a><span class="lineno"> 285</span>  }</div>
|
|
<div class="line"><a name="l00286"></a><span class="lineno"> 286</span>  </div>
|
|
<div class="line"><a name="l00287"></a><span class="lineno"> 287</span>  int64_t Receive(std::vector<MessagePtr>& msgVec, <span class="keyword">const</span> <span class="keywordtype">int</span> timeout = -1)<span class="keyword"> override</span></div>
|
|
<div class="line"><a name="l00288"></a><span class="lineno"> 288</span> <span class="keyword"> </span>{</div>
|
|
<div class="line"><a name="l00289"></a><span class="lineno"> 289</span>  <span class="keywordtype">int</span> flags = 0;</div>
|
|
<div class="line"><a name="l00290"></a><span class="lineno"> 290</span>  <span class="keywordflow">if</span> (timeout == 0) {</div>
|
|
<div class="line"><a name="l00291"></a><span class="lineno"> 291</span>  flags = ZMQ_DONTWAIT;</div>
|
|
<div class="line"><a name="l00292"></a><span class="lineno"> 292</span>  }</div>
|
|
<div class="line"><a name="l00293"></a><span class="lineno"> 293</span>  <span class="keywordtype">int</span> elapsed = 0;</div>
|
|
<div class="line"><a name="l00294"></a><span class="lineno"> 294</span>  </div>
|
|
<div class="line"><a name="l00295"></a><span class="lineno"> 295</span>  <a class="code" href="structfair_1_1mq_1_1shmem_1_1ZMsg.html">ZMsg</a> zmqMsg;</div>
|
|
<div class="line"><a name="l00296"></a><span class="lineno"> 296</span>  </div>
|
|
<div class="line"><a name="l00297"></a><span class="lineno"> 297</span>  <span class="keywordflow">while</span> (<span class="keyword">true</span>) {</div>
|
|
<div class="line"><a name="l00298"></a><span class="lineno"> 298</span>  int64_t totalSize = 0;</div>
|
|
<div class="line"><a name="l00299"></a><span class="lineno"> 299</span>  <span class="keywordtype">int</span> nbytes = zmq_msg_recv(zmqMsg.Msg(), fSocket, flags);</div>
|
|
<div class="line"><a name="l00300"></a><span class="lineno"> 300</span>  <span class="keywordflow">if</span> (nbytes > 0) {</div>
|
|
<div class="line"><a name="l00301"></a><span class="lineno"> 301</span>  <a class="code" href="structfair_1_1mq_1_1shmem_1_1MetaHeader.html">MetaHeader</a>* hdrVec = <span class="keyword">static_cast<</span><a class="code" href="structfair_1_1mq_1_1shmem_1_1MetaHeader.html">MetaHeader</a>*<span class="keyword">></span>(zmqMsg.Data());</div>
|
|
<div class="line"><a name="l00302"></a><span class="lineno"> 302</span>  <span class="keyword">const</span> <span class="keyword">auto</span> hdrVecSize = zmqMsg.Size();</div>
|
|
<div class="line"><a name="l00303"></a><span class="lineno"> 303</span>  </div>
|
|
<div class="line"><a name="l00304"></a><span class="lineno"> 304</span>  assert(hdrVecSize > 0);</div>
|
|
<div class="line"><a name="l00305"></a><span class="lineno"> 305</span>  <span class="keywordflow">if</span> (hdrVecSize % <span class="keyword">sizeof</span>(<a class="code" href="structfair_1_1mq_1_1shmem_1_1MetaHeader.html">MetaHeader</a>) != 0) {</div>
|
|
<div class="line"><a name="l00306"></a><span class="lineno"> 306</span>  <span class="keywordflow">throw</span> <a class="code" href="structfair_1_1mq_1_1SocketError.html">SocketError</a>(</div>
|
|
<div class="line"><a name="l00307"></a><span class="lineno"> 307</span>  tools::ToString(<span class="stringliteral">"Received message is not a valid FairMQ shared memory message. "</span>,</div>
|
|
<div class="line"><a name="l00308"></a><span class="lineno"> 308</span>  <span class="stringliteral">"Possibly due to a misconfigured transport on the sender side. "</span>,</div>
|
|
<div class="line"><a name="l00309"></a><span class="lineno"> 309</span>  <span class="stringliteral">"Expected size of "</span>, <span class="keyword">sizeof</span>(<a class="code" href="structfair_1_1mq_1_1shmem_1_1MetaHeader.html">MetaHeader</a>), <span class="stringliteral">" bytes, received "</span>, nbytes));</div>
|
|
<div class="line"><a name="l00310"></a><span class="lineno"> 310</span>  }</div>
|
|
<div class="line"><a name="l00311"></a><span class="lineno"> 311</span>  </div>
|
|
<div class="line"><a name="l00312"></a><span class="lineno"> 312</span>  <span class="keyword">const</span> <span class="keyword">auto</span> numMessages = hdrVecSize / <span class="keyword">sizeof</span>(<a class="code" href="structfair_1_1mq_1_1shmem_1_1MetaHeader.html">MetaHeader</a>);</div>
|
|
<div class="line"><a name="l00313"></a><span class="lineno"> 313</span>  msgVec.reserve(numMessages);</div>
|
|
<div class="line"><a name="l00314"></a><span class="lineno"> 314</span>  </div>
|
|
<div class="line"><a name="l00315"></a><span class="lineno"> 315</span>  <span class="keywordflow">for</span> (<span class="keywordtype">size_t</span> m = 0; m < numMessages; m++) {</div>
|
|
<div class="line"><a name="l00316"></a><span class="lineno"> 316</span>  <span class="comment">// create new message (part)</span></div>
|
|
<div class="line"><a name="l00317"></a><span class="lineno"> 317</span>  msgVec.emplace_back(std::make_unique<Message>(fManager, hdrVec[m], GetTransport()));</div>
|
|
<div class="line"><a name="l00318"></a><span class="lineno"> 318</span>  <a class="code" href="classfair_1_1mq_1_1shmem_1_1Message.html">Message</a>* shmMsg = <span class="keyword">static_cast<</span><a class="code" href="classfair_1_1mq_1_1shmem_1_1Message.html">Message</a>*<span class="keyword">></span>(msgVec.back().get());</div>
|
|
<div class="line"><a name="l00319"></a><span class="lineno"> 319</span>  totalSize += shmMsg->GetSize();</div>
|
|
<div class="line"><a name="l00320"></a><span class="lineno"> 320</span>  }</div>
|
|
<div class="line"><a name="l00321"></a><span class="lineno"> 321</span>  </div>
|
|
<div class="line"><a name="l00322"></a><span class="lineno"> 322</span>  <span class="comment">// store statistics on how many messages have been received (handle all parts as a single message)</span></div>
|
|
<div class="line"><a name="l00323"></a><span class="lineno"> 323</span>  fMessagesRx++;</div>
|
|
<div class="line"><a name="l00324"></a><span class="lineno"> 324</span>  fBytesRx += totalSize;</div>
|
|
<div class="line"><a name="l00325"></a><span class="lineno"> 325</span>  </div>
|
|
<div class="line"><a name="l00326"></a><span class="lineno"> 326</span>  <span class="keywordflow">return</span> totalSize;</div>
|
|
<div class="line"><a name="l00327"></a><span class="lineno"> 327</span>  } <span class="keywordflow">else</span> <span class="keywordflow">if</span> (zmq_errno() == EAGAIN || zmq_errno() == EINTR) {</div>
|
|
<div class="line"><a name="l00328"></a><span class="lineno"> 328</span>  <span class="keywordflow">if</span> (fManager.Interrupted()) {</div>
|
|
<div class="line"><a name="l00329"></a><span class="lineno"> 329</span>  <span class="keywordflow">return</span> <span class="keyword">static_cast<</span><span class="keywordtype">int</span><span class="keyword">></span>(TransferCode::interrupted);</div>
|
|
<div class="line"><a name="l00330"></a><span class="lineno"> 330</span>  } <span class="keywordflow">else</span> <span class="keywordflow">if</span> (ShouldRetry(flags, timeout, elapsed)) {</div>
|
|
<div class="line"><a name="l00331"></a><span class="lineno"> 331</span>  <span class="keywordflow">continue</span>;</div>
|
|
<div class="line"><a name="l00332"></a><span class="lineno"> 332</span>  } <span class="keywordflow">else</span> {</div>
|
|
<div class="line"><a name="l00333"></a><span class="lineno"> 333</span>  <span class="keywordflow">return</span> <span class="keyword">static_cast<</span><span class="keywordtype">int</span><span class="keyword">></span>(TransferCode::timeout);</div>
|
|
<div class="line"><a name="l00334"></a><span class="lineno"> 334</span>  }</div>
|
|
<div class="line"><a name="l00335"></a><span class="lineno"> 335</span>  } <span class="keywordflow">else</span> {</div>
|
|
<div class="line"><a name="l00336"></a><span class="lineno"> 336</span>  <span class="keywordflow">return</span> HandleErrors();</div>
|
|
<div class="line"><a name="l00337"></a><span class="lineno"> 337</span>  }</div>
|
|
<div class="line"><a name="l00338"></a><span class="lineno"> 338</span>  }</div>
|
|
<div class="line"><a name="l00339"></a><span class="lineno"> 339</span>  </div>
|
|
<div class="line"><a name="l00340"></a><span class="lineno"> 340</span>  <span class="keywordflow">return</span> <span class="keyword">static_cast<</span><span class="keywordtype">int</span><span class="keyword">></span>(TransferCode::error);</div>
|
|
<div class="line"><a name="l00341"></a><span class="lineno"> 341</span>  }</div>
|
|
<div class="line"><a name="l00342"></a><span class="lineno"> 342</span>  </div>
|
|
<div class="line"><a name="l00343"></a><span class="lineno"> 343</span>  <span class="keywordtype">void</span>* GetSocket()<span class="keyword"> const </span>{ <span class="keywordflow">return</span> fSocket; }</div>
|
|
<div class="line"><a name="l00344"></a><span class="lineno"> 344</span>  </div>
|
|
<div class="line"><a name="l00345"></a><span class="lineno"> 345</span>  <span class="keywordtype">void</span> Close()<span class="keyword"> override</span></div>
|
|
<div class="line"><a name="l00346"></a><span class="lineno"> 346</span> <span class="keyword"> </span>{</div>
|
|
<div class="line"><a name="l00347"></a><span class="lineno"> 347</span>  <span class="comment">// LOG(debug) << "Closing socket " << fId;</span></div>
|
|
<div class="line"><a name="l00348"></a><span class="lineno"> 348</span>  </div>
|
|
<div class="line"><a name="l00349"></a><span class="lineno"> 349</span>  <span class="keywordflow">if</span> (fSocket == <span class="keyword">nullptr</span>) {</div>
|
|
<div class="line"><a name="l00350"></a><span class="lineno"> 350</span>  <span class="keywordflow">return</span>;</div>
|
|
<div class="line"><a name="l00351"></a><span class="lineno"> 351</span>  }</div>
|
|
<div class="line"><a name="l00352"></a><span class="lineno"> 352</span>  </div>
|
|
<div class="line"><a name="l00353"></a><span class="lineno"> 353</span>  <span class="keywordflow">if</span> (zmq_close(fSocket) != 0) {</div>
|
|
<div class="line"><a name="l00354"></a><span class="lineno"> 354</span>  LOG(error) << <span class="stringliteral">"Failed closing socket "</span> << fId << <span class="stringliteral">", reason: "</span> << zmq_strerror(errno);</div>
|
|
<div class="line"><a name="l00355"></a><span class="lineno"> 355</span>  }</div>
|
|
<div class="line"><a name="l00356"></a><span class="lineno"> 356</span>  </div>
|
|
<div class="line"><a name="l00357"></a><span class="lineno"> 357</span>  fSocket = <span class="keyword">nullptr</span>;</div>
|
|
<div class="line"><a name="l00358"></a><span class="lineno"> 358</span>  }</div>
|
|
<div class="line"><a name="l00359"></a><span class="lineno"> 359</span>  </div>
|
|
<div class="line"><a name="l00360"></a><span class="lineno"> 360</span>  <span class="keywordtype">void</span> SetOption(<span class="keyword">const</span> std::string& option, <span class="keyword">const</span> <span class="keywordtype">void</span>* value, <span class="keywordtype">size_t</span> valueSize)<span class="keyword"> override</span></div>
|
|
<div class="line"><a name="l00361"></a><span class="lineno"> 361</span> <span class="keyword"> </span>{</div>
|
|
<div class="line"><a name="l00362"></a><span class="lineno"> 362</span>  <span class="keywordflow">if</span> (zmq_setsockopt(fSocket, GetConstant(option), value, valueSize) < 0) {</div>
|
|
<div class="line"><a name="l00363"></a><span class="lineno"> 363</span>  LOG(error) << <span class="stringliteral">"Failed setting socket option, reason: "</span> << zmq_strerror(errno);</div>
|
|
<div class="line"><a name="l00364"></a><span class="lineno"> 364</span>  }</div>
|
|
<div class="line"><a name="l00365"></a><span class="lineno"> 365</span>  }</div>
|
|
<div class="line"><a name="l00366"></a><span class="lineno"> 366</span>  </div>
|
|
<div class="line"><a name="l00367"></a><span class="lineno"> 367</span>  <span class="keywordtype">void</span> GetOption(<span class="keyword">const</span> std::string& option, <span class="keywordtype">void</span>* value, <span class="keywordtype">size_t</span>* valueSize)<span class="keyword"> override</span></div>
|
|
<div class="line"><a name="l00368"></a><span class="lineno"> 368</span> <span class="keyword"> </span>{</div>
|
|
<div class="line"><a name="l00369"></a><span class="lineno"> 369</span>  <span class="keywordflow">if</span> (zmq_getsockopt(fSocket, GetConstant(option), value, valueSize) < 0) {</div>
|
|
<div class="line"><a name="l00370"></a><span class="lineno"> 370</span>  LOG(error) << <span class="stringliteral">"Failed getting socket option, reason: "</span> << zmq_strerror(errno);</div>
|
|
<div class="line"><a name="l00371"></a><span class="lineno"> 371</span>  }</div>
|
|
<div class="line"><a name="l00372"></a><span class="lineno"> 372</span>  }</div>
|
|
<div class="line"><a name="l00373"></a><span class="lineno"> 373</span>  </div>
|
|
<div class="line"><a name="l00374"></a><span class="lineno"> 374</span>  <span class="keywordtype">void</span> SetLinger(<span class="keyword">const</span> <span class="keywordtype">int</span> value)<span class="keyword"> override</span></div>
|
|
<div class="line"><a name="l00375"></a><span class="lineno"> 375</span> <span class="keyword"> </span>{</div>
|
|
<div class="line"><a name="l00376"></a><span class="lineno"> 376</span>  <span class="keywordflow">if</span> (zmq_setsockopt(fSocket, ZMQ_LINGER, &value, <span class="keyword">sizeof</span>(value)) < 0) {</div>
|
|
<div class="line"><a name="l00377"></a><span class="lineno"> 377</span>  <span class="keywordflow">throw</span> <a class="code" href="structfair_1_1mq_1_1SocketError.html">SocketError</a>(tools::ToString(<span class="stringliteral">"failed setting ZMQ_LINGER, reason: "</span>, zmq_strerror(errno)));</div>
|
|
<div class="line"><a name="l00378"></a><span class="lineno"> 378</span>  }</div>
|
|
<div class="line"><a name="l00379"></a><span class="lineno"> 379</span>  }</div>
|
|
<div class="line"><a name="l00380"></a><span class="lineno"> 380</span>  </div>
|
|
<div class="line"><a name="l00381"></a><span class="lineno"><a class="line" href="classfair_1_1mq_1_1shmem_1_1Socket.html#acd1bc3ce745e748eefad50ac19c175dd"> 381</a></span>  <span class="keywordtype">void</span> <a class="code" href="classfair_1_1mq_1_1shmem_1_1Socket.html#acd1bc3ce745e748eefad50ac19c175dd">Events</a>(uint32_t* events)<span class="keyword"> override</span></div>
|
|
<div class="line"><a name="l00382"></a><span class="lineno"> 382</span> <span class="keyword"> </span>{</div>
|
|
<div class="line"><a name="l00383"></a><span class="lineno"> 383</span>  <span class="keywordtype">size_t</span> eventsSize = <span class="keyword">sizeof</span>(uint32_t);</div>
|
|
<div class="line"><a name="l00384"></a><span class="lineno"> 384</span>  <span class="keywordflow">if</span> (zmq_getsockopt(fSocket, ZMQ_EVENTS, events, &eventsSize) < 0) {</div>
|
|
<div class="line"><a name="l00385"></a><span class="lineno"> 385</span>  <span class="keywordflow">throw</span> <a class="code" href="structfair_1_1mq_1_1SocketError.html">SocketError</a>(tools::ToString(<span class="stringliteral">"failed setting ZMQ_EVENTS, reason: "</span>, zmq_strerror(errno)));</div>
|
|
<div class="line"><a name="l00386"></a><span class="lineno"> 386</span>  }</div>
|
|
<div class="line"><a name="l00387"></a><span class="lineno"> 387</span>  }</div>
|
|
<div class="line"><a name="l00388"></a><span class="lineno"> 388</span>  </div>
|
|
<div class="line"><a name="l00389"></a><span class="lineno"> 389</span>  <span class="keywordtype">int</span> GetLinger()<span class="keyword"> const override</span></div>
|
|
<div class="line"><a name="l00390"></a><span class="lineno"> 390</span> <span class="keyword"> </span>{</div>
|
|
<div class="line"><a name="l00391"></a><span class="lineno"> 391</span>  <span class="keywordtype">int</span> value = 0;</div>
|
|
<div class="line"><a name="l00392"></a><span class="lineno"> 392</span>  <span class="keywordtype">size_t</span> valueSize = <span class="keyword">sizeof</span>(value);</div>
|
|
<div class="line"><a name="l00393"></a><span class="lineno"> 393</span>  <span class="keywordflow">if</span> (zmq_getsockopt(fSocket, ZMQ_LINGER, &value, &valueSize) < 0) {</div>
|
|
<div class="line"><a name="l00394"></a><span class="lineno"> 394</span>  <span class="keywordflow">throw</span> <a class="code" href="structfair_1_1mq_1_1SocketError.html">SocketError</a>(tools::ToString(<span class="stringliteral">"failed getting ZMQ_LINGER, reason: "</span>, zmq_strerror(errno)));</div>
|
|
<div class="line"><a name="l00395"></a><span class="lineno"> 395</span>  }</div>
|
|
<div class="line"><a name="l00396"></a><span class="lineno"> 396</span>  <span class="keywordflow">return</span> value;</div>
|
|
<div class="line"><a name="l00397"></a><span class="lineno"> 397</span>  }</div>
|
|
<div class="line"><a name="l00398"></a><span class="lineno"> 398</span>  </div>
|
|
<div class="line"><a name="l00399"></a><span class="lineno"> 399</span>  <span class="keywordtype">void</span> SetSndBufSize(<span class="keyword">const</span> <span class="keywordtype">int</span> value)<span class="keyword"> override</span></div>
|
|
<div class="line"><a name="l00400"></a><span class="lineno"> 400</span> <span class="keyword"> </span>{</div>
|
|
<div class="line"><a name="l00401"></a><span class="lineno"> 401</span>  <span class="keywordflow">if</span> (zmq_setsockopt(fSocket, ZMQ_SNDHWM, &value, <span class="keyword">sizeof</span>(value)) < 0) {</div>
|
|
<div class="line"><a name="l00402"></a><span class="lineno"> 402</span>  <span class="keywordflow">throw</span> SocketError(tools::ToString(<span class="stringliteral">"failed setting ZMQ_SNDHWM, reason: "</span>, zmq_strerror(errno)));</div>
|
|
<div class="line"><a name="l00403"></a><span class="lineno"> 403</span>  }</div>
|
|
<div class="line"><a name="l00404"></a><span class="lineno"> 404</span>  }</div>
|
|
<div class="line"><a name="l00405"></a><span class="lineno"> 405</span>  </div>
|
|
<div class="line"><a name="l00406"></a><span class="lineno"> 406</span>  <span class="keywordtype">int</span> GetSndBufSize()<span class="keyword"> const override</span></div>
|
|
<div class="line"><a name="l00407"></a><span class="lineno"> 407</span> <span class="keyword"> </span>{</div>
|
|
<div class="line"><a name="l00408"></a><span class="lineno"> 408</span>  <span class="keywordtype">int</span> value = 0;</div>
|
|
<div class="line"><a name="l00409"></a><span class="lineno"> 409</span>  <span class="keywordtype">size_t</span> valueSize = <span class="keyword">sizeof</span>(value);</div>
|
|
<div class="line"><a name="l00410"></a><span class="lineno"> 410</span>  <span class="keywordflow">if</span> (zmq_getsockopt(fSocket, ZMQ_SNDHWM, &value, &valueSize) < 0) {</div>
|
|
<div class="line"><a name="l00411"></a><span class="lineno"> 411</span>  <span class="keywordflow">throw</span> SocketError(tools::ToString(<span class="stringliteral">"failed getting ZMQ_SNDHWM, reason: "</span>, zmq_strerror(errno)));</div>
|
|
<div class="line"><a name="l00412"></a><span class="lineno"> 412</span>  }</div>
|
|
<div class="line"><a name="l00413"></a><span class="lineno"> 413</span>  <span class="keywordflow">return</span> value;</div>
|
|
<div class="line"><a name="l00414"></a><span class="lineno"> 414</span>  }</div>
|
|
<div class="line"><a name="l00415"></a><span class="lineno"> 415</span>  </div>
|
|
<div class="line"><a name="l00416"></a><span class="lineno"> 416</span>  <span class="keywordtype">void</span> SetRcvBufSize(<span class="keyword">const</span> <span class="keywordtype">int</span> value)<span class="keyword"> override</span></div>
|
|
<div class="line"><a name="l00417"></a><span class="lineno"> 417</span> <span class="keyword"> </span>{</div>
|
|
<div class="line"><a name="l00418"></a><span class="lineno"> 418</span>  <span class="keywordflow">if</span> (zmq_setsockopt(fSocket, ZMQ_RCVHWM, &value, <span class="keyword">sizeof</span>(value)) < 0) {</div>
|
|
<div class="line"><a name="l00419"></a><span class="lineno"> 419</span>  <span class="keywordflow">throw</span> SocketError(tools::ToString(<span class="stringliteral">"failed setting ZMQ_RCVHWM, reason: "</span>, zmq_strerror(errno)));</div>
|
|
<div class="line"><a name="l00420"></a><span class="lineno"> 420</span>  }</div>
|
|
<div class="line"><a name="l00421"></a><span class="lineno"> 421</span>  }</div>
|
|
<div class="line"><a name="l00422"></a><span class="lineno"> 422</span>  </div>
|
|
<div class="line"><a name="l00423"></a><span class="lineno"> 423</span>  <span class="keywordtype">int</span> GetRcvBufSize()<span class="keyword"> const override</span></div>
|
|
<div class="line"><a name="l00424"></a><span class="lineno"> 424</span> <span class="keyword"> </span>{</div>
|
|
<div class="line"><a name="l00425"></a><span class="lineno"> 425</span>  <span class="keywordtype">int</span> value = 0;</div>
|
|
<div class="line"><a name="l00426"></a><span class="lineno"> 426</span>  <span class="keywordtype">size_t</span> valueSize = <span class="keyword">sizeof</span>(value);</div>
|
|
<div class="line"><a name="l00427"></a><span class="lineno"> 427</span>  <span class="keywordflow">if</span> (zmq_getsockopt(fSocket, ZMQ_RCVHWM, &value, &valueSize) < 0) {</div>
|
|
<div class="line"><a name="l00428"></a><span class="lineno"> 428</span>  <span class="keywordflow">throw</span> SocketError(tools::ToString(<span class="stringliteral">"failed getting ZMQ_RCVHWM, reason: "</span>, zmq_strerror(errno)));</div>
|
|
<div class="line"><a name="l00429"></a><span class="lineno"> 429</span>  }</div>
|
|
<div class="line"><a name="l00430"></a><span class="lineno"> 430</span>  <span class="keywordflow">return</span> value;</div>
|
|
<div class="line"><a name="l00431"></a><span class="lineno"> 431</span>  }</div>
|
|
<div class="line"><a name="l00432"></a><span class="lineno"> 432</span>  </div>
|
|
<div class="line"><a name="l00433"></a><span class="lineno"> 433</span>  <span class="keywordtype">void</span> SetSndKernelSize(<span class="keyword">const</span> <span class="keywordtype">int</span> value)<span class="keyword"> override</span></div>
|
|
<div class="line"><a name="l00434"></a><span class="lineno"> 434</span> <span class="keyword"> </span>{</div>
|
|
<div class="line"><a name="l00435"></a><span class="lineno"> 435</span>  <span class="keywordflow">if</span> (zmq_setsockopt(fSocket, ZMQ_SNDBUF, &value, <span class="keyword">sizeof</span>(value)) < 0) {</div>
|
|
<div class="line"><a name="l00436"></a><span class="lineno"> 436</span>  <span class="keywordflow">throw</span> SocketError(tools::ToString(<span class="stringliteral">"failed getting ZMQ_SNDBUF, reason: "</span>, zmq_strerror(errno)));</div>
|
|
<div class="line"><a name="l00437"></a><span class="lineno"> 437</span>  }</div>
|
|
<div class="line"><a name="l00438"></a><span class="lineno"> 438</span>  }</div>
|
|
<div class="line"><a name="l00439"></a><span class="lineno"> 439</span>  </div>
|
|
<div class="line"><a name="l00440"></a><span class="lineno"> 440</span>  <span class="keywordtype">int</span> GetSndKernelSize()<span class="keyword"> const override</span></div>
|
|
<div class="line"><a name="l00441"></a><span class="lineno"> 441</span> <span class="keyword"> </span>{</div>
|
|
<div class="line"><a name="l00442"></a><span class="lineno"> 442</span>  <span class="keywordtype">int</span> value = 0;</div>
|
|
<div class="line"><a name="l00443"></a><span class="lineno"> 443</span>  <span class="keywordtype">size_t</span> valueSize = <span class="keyword">sizeof</span>(value);</div>
|
|
<div class="line"><a name="l00444"></a><span class="lineno"> 444</span>  <span class="keywordflow">if</span> (zmq_getsockopt(fSocket, ZMQ_SNDBUF, &value, &valueSize) < 0) {</div>
|
|
<div class="line"><a name="l00445"></a><span class="lineno"> 445</span>  <span class="keywordflow">throw</span> SocketError(tools::ToString(<span class="stringliteral">"failed getting ZMQ_SNDBUF, reason: "</span>, zmq_strerror(errno)));</div>
|
|
<div class="line"><a name="l00446"></a><span class="lineno"> 446</span>  }</div>
|
|
<div class="line"><a name="l00447"></a><span class="lineno"> 447</span>  <span class="keywordflow">return</span> value;</div>
|
|
<div class="line"><a name="l00448"></a><span class="lineno"> 448</span>  }</div>
|
|
<div class="line"><a name="l00449"></a><span class="lineno"> 449</span>  </div>
|
|
<div class="line"><a name="l00450"></a><span class="lineno"> 450</span>  <span class="keywordtype">void</span> SetRcvKernelSize(<span class="keyword">const</span> <span class="keywordtype">int</span> value)<span class="keyword"> override</span></div>
|
|
<div class="line"><a name="l00451"></a><span class="lineno"> 451</span> <span class="keyword"> </span>{</div>
|
|
<div class="line"><a name="l00452"></a><span class="lineno"> 452</span>  <span class="keywordflow">if</span> (zmq_setsockopt(fSocket, ZMQ_RCVBUF, &value, <span class="keyword">sizeof</span>(value)) < 0) {</div>
|
|
<div class="line"><a name="l00453"></a><span class="lineno"> 453</span>  <span class="keywordflow">throw</span> SocketError(tools::ToString(<span class="stringliteral">"failed getting ZMQ_RCVBUF, reason: "</span>, zmq_strerror(errno)));</div>
|
|
<div class="line"><a name="l00454"></a><span class="lineno"> 454</span>  }</div>
|
|
<div class="line"><a name="l00455"></a><span class="lineno"> 455</span>  }</div>
|
|
<div class="line"><a name="l00456"></a><span class="lineno"> 456</span>  </div>
|
|
<div class="line"><a name="l00457"></a><span class="lineno"> 457</span>  <span class="keywordtype">int</span> GetRcvKernelSize()<span class="keyword"> const override</span></div>
|
|
<div class="line"><a name="l00458"></a><span class="lineno"> 458</span> <span class="keyword"> </span>{</div>
|
|
<div class="line"><a name="l00459"></a><span class="lineno"> 459</span>  <span class="keywordtype">int</span> value = 0;</div>
|
|
<div class="line"><a name="l00460"></a><span class="lineno"> 460</span>  <span class="keywordtype">size_t</span> valueSize = <span class="keyword">sizeof</span>(value);</div>
|
|
<div class="line"><a name="l00461"></a><span class="lineno"> 461</span>  <span class="keywordflow">if</span> (zmq_getsockopt(fSocket, ZMQ_RCVBUF, &value, &valueSize) < 0) {</div>
|
|
<div class="line"><a name="l00462"></a><span class="lineno"> 462</span>  <span class="keywordflow">throw</span> SocketError(tools::ToString(<span class="stringliteral">"failed getting ZMQ_RCVBUF, reason: "</span>, zmq_strerror(errno)));</div>
|
|
<div class="line"><a name="l00463"></a><span class="lineno"> 463</span>  }</div>
|
|
<div class="line"><a name="l00464"></a><span class="lineno"> 464</span>  <span class="keywordflow">return</span> value;</div>
|
|
<div class="line"><a name="l00465"></a><span class="lineno"> 465</span>  }</div>
|
|
<div class="line"><a name="l00466"></a><span class="lineno"> 466</span>  </div>
|
|
<div class="line"><a name="l00467"></a><span class="lineno"> 467</span>  <span class="keywordtype">unsigned</span> <span class="keywordtype">long</span> GetBytesTx()<span class="keyword"> const override </span>{ <span class="keywordflow">return</span> fBytesTx; }</div>
|
|
<div class="line"><a name="l00468"></a><span class="lineno"> 468</span>  <span class="keywordtype">unsigned</span> <span class="keywordtype">long</span> GetBytesRx()<span class="keyword"> const override </span>{ <span class="keywordflow">return</span> fBytesRx; }</div>
|
|
<div class="line"><a name="l00469"></a><span class="lineno"> 469</span>  <span class="keywordtype">unsigned</span> <span class="keywordtype">long</span> GetMessagesTx()<span class="keyword"> const override </span>{ <span class="keywordflow">return</span> fMessagesTx; }</div>
|
|
<div class="line"><a name="l00470"></a><span class="lineno"> 470</span>  <span class="keywordtype">unsigned</span> <span class="keywordtype">long</span> GetMessagesRx()<span class="keyword"> const override </span>{ <span class="keywordflow">return</span> fMessagesRx; }</div>
|
|
<div class="line"><a name="l00471"></a><span class="lineno"> 471</span>  </div>
|
|
<div class="line"><a name="l00472"></a><span class="lineno"> 472</span>  <span class="keyword">static</span> <span class="keywordtype">int</span> GetConstant(<span class="keyword">const</span> std::string& constant)</div>
|
|
<div class="line"><a name="l00473"></a><span class="lineno"> 473</span>  {</div>
|
|
<div class="line"><a name="l00474"></a><span class="lineno"> 474</span>  <span class="keywordflow">if</span> (constant == <span class="stringliteral">""</span>) <span class="keywordflow">return</span> 0;</div>
|
|
<div class="line"><a name="l00475"></a><span class="lineno"> 475</span>  <span class="keywordflow">if</span> (constant == <span class="stringliteral">"sub"</span>) <span class="keywordflow">return</span> ZMQ_SUB;</div>
|
|
<div class="line"><a name="l00476"></a><span class="lineno"> 476</span>  <span class="keywordflow">if</span> (constant == <span class="stringliteral">"pub"</span>) <span class="keywordflow">return</span> ZMQ_PUB;</div>
|
|
<div class="line"><a name="l00477"></a><span class="lineno"> 477</span>  <span class="keywordflow">if</span> (constant == <span class="stringliteral">"xsub"</span>) <span class="keywordflow">return</span> ZMQ_XSUB;</div>
|
|
<div class="line"><a name="l00478"></a><span class="lineno"> 478</span>  <span class="keywordflow">if</span> (constant == <span class="stringliteral">"xpub"</span>) <span class="keywordflow">return</span> ZMQ_XPUB;</div>
|
|
<div class="line"><a name="l00479"></a><span class="lineno"> 479</span>  <span class="keywordflow">if</span> (constant == <span class="stringliteral">"push"</span>) <span class="keywordflow">return</span> ZMQ_PUSH;</div>
|
|
<div class="line"><a name="l00480"></a><span class="lineno"> 480</span>  <span class="keywordflow">if</span> (constant == <span class="stringliteral">"pull"</span>) <span class="keywordflow">return</span> ZMQ_PULL;</div>
|
|
<div class="line"><a name="l00481"></a><span class="lineno"> 481</span>  <span class="keywordflow">if</span> (constant == <span class="stringliteral">"req"</span>) <span class="keywordflow">return</span> ZMQ_REQ;</div>
|
|
<div class="line"><a name="l00482"></a><span class="lineno"> 482</span>  <span class="keywordflow">if</span> (constant == <span class="stringliteral">"rep"</span>) <span class="keywordflow">return</span> ZMQ_REP;</div>
|
|
<div class="line"><a name="l00483"></a><span class="lineno"> 483</span>  <span class="keywordflow">if</span> (constant == <span class="stringliteral">"dealer"</span>) <span class="keywordflow">return</span> ZMQ_DEALER;</div>
|
|
<div class="line"><a name="l00484"></a><span class="lineno"> 484</span>  <span class="keywordflow">if</span> (constant == <span class="stringliteral">"router"</span>) <span class="keywordflow">return</span> ZMQ_ROUTER;</div>
|
|
<div class="line"><a name="l00485"></a><span class="lineno"> 485</span>  <span class="keywordflow">if</span> (constant == <span class="stringliteral">"pair"</span>) <span class="keywordflow">return</span> ZMQ_PAIR;</div>
|
|
<div class="line"><a name="l00486"></a><span class="lineno"> 486</span>  </div>
|
|
<div class="line"><a name="l00487"></a><span class="lineno"> 487</span>  <span class="keywordflow">if</span> (constant == <span class="stringliteral">"snd-hwm"</span>) <span class="keywordflow">return</span> ZMQ_SNDHWM;</div>
|
|
<div class="line"><a name="l00488"></a><span class="lineno"> 488</span>  <span class="keywordflow">if</span> (constant == <span class="stringliteral">"rcv-hwm"</span>) <span class="keywordflow">return</span> ZMQ_RCVHWM;</div>
|
|
<div class="line"><a name="l00489"></a><span class="lineno"> 489</span>  <span class="keywordflow">if</span> (constant == <span class="stringliteral">"snd-size"</span>) <span class="keywordflow">return</span> ZMQ_SNDBUF;</div>
|
|
<div class="line"><a name="l00490"></a><span class="lineno"> 490</span>  <span class="keywordflow">if</span> (constant == <span class="stringliteral">"rcv-size"</span>) <span class="keywordflow">return</span> ZMQ_RCVBUF;</div>
|
|
<div class="line"><a name="l00491"></a><span class="lineno"> 491</span>  <span class="keywordflow">if</span> (constant == <span class="stringliteral">"snd-more"</span>) <span class="keywordflow">return</span> ZMQ_SNDMORE;</div>
|
|
<div class="line"><a name="l00492"></a><span class="lineno"> 492</span>  <span class="keywordflow">if</span> (constant == <span class="stringliteral">"rcv-more"</span>) <span class="keywordflow">return</span> ZMQ_RCVMORE;</div>
|
|
<div class="line"><a name="l00493"></a><span class="lineno"> 493</span>  </div>
|
|
<div class="line"><a name="l00494"></a><span class="lineno"> 494</span>  <span class="keywordflow">if</span> (constant == <span class="stringliteral">"linger"</span>) <span class="keywordflow">return</span> ZMQ_LINGER;</div>
|
|
<div class="line"><a name="l00495"></a><span class="lineno"> 495</span>  <span class="keywordflow">if</span> (constant == <span class="stringliteral">"no-block"</span>) <span class="keywordflow">return</span> ZMQ_DONTWAIT;</div>
|
|
<div class="line"><a name="l00496"></a><span class="lineno"> 496</span>  <span class="keywordflow">if</span> (constant == <span class="stringliteral">"snd-more no-block"</span>) <span class="keywordflow">return</span> ZMQ_DONTWAIT|ZMQ_SNDMORE;</div>
|
|
<div class="line"><a name="l00497"></a><span class="lineno"> 497</span>  </div>
|
|
<div class="line"><a name="l00498"></a><span class="lineno"> 498</span>  <span class="keywordflow">if</span> (constant == <span class="stringliteral">"fd"</span>) <span class="keywordflow">return</span> ZMQ_FD;</div>
|
|
<div class="line"><a name="l00499"></a><span class="lineno"> 499</span>  <span class="keywordflow">if</span> (constant == <span class="stringliteral">"events"</span>)</div>
|
|
<div class="line"><a name="l00500"></a><span class="lineno"> 500</span>  <span class="keywordflow">return</span> ZMQ_EVENTS;</div>
|
|
<div class="line"><a name="l00501"></a><span class="lineno"> 501</span>  <span class="keywordflow">if</span> (constant == <span class="stringliteral">"pollin"</span>)</div>
|
|
<div class="line"><a name="l00502"></a><span class="lineno"> 502</span>  <span class="keywordflow">return</span> ZMQ_POLLIN;</div>
|
|
<div class="line"><a name="l00503"></a><span class="lineno"> 503</span>  <span class="keywordflow">if</span> (constant == <span class="stringliteral">"pollout"</span>)</div>
|
|
<div class="line"><a name="l00504"></a><span class="lineno"> 504</span>  <span class="keywordflow">return</span> ZMQ_POLLOUT;</div>
|
|
<div class="line"><a name="l00505"></a><span class="lineno"> 505</span>  </div>
|
|
<div class="line"><a name="l00506"></a><span class="lineno"> 506</span>  <span class="keywordflow">throw</span> SocketError(tools::ToString(<span class="stringliteral">"GetConstant called with an invalid argument: "</span>, constant));</div>
|
|
<div class="line"><a name="l00507"></a><span class="lineno"> 507</span>  }</div>
|
|
<div class="line"><a name="l00508"></a><span class="lineno"> 508</span>  </div>
|
|
<div class="line"><a name="l00509"></a><span class="lineno"> 509</span>  ~Socket()<span class="keyword"> override </span>{ Close(); }</div>
|
|
<div class="line"><a name="l00510"></a><span class="lineno"> 510</span>  </div>
|
|
<div class="line"><a name="l00511"></a><span class="lineno"> 511</span>  <span class="keyword">private</span>:</div>
|
|
<div class="line"><a name="l00512"></a><span class="lineno"> 512</span>  <span class="keywordtype">void</span>* fSocket;</div>
|
|
<div class="line"><a name="l00513"></a><span class="lineno"> 513</span>  Manager& fManager;</div>
|
|
<div class="line"><a name="l00514"></a><span class="lineno"> 514</span>  std::string fId;</div>
|
|
<div class="line"><a name="l00515"></a><span class="lineno"> 515</span>  std::atomic<unsigned long> fBytesTx;</div>
|
|
<div class="line"><a name="l00516"></a><span class="lineno"> 516</span>  std::atomic<unsigned long> fBytesRx;</div>
|
|
<div class="line"><a name="l00517"></a><span class="lineno"> 517</span>  std::atomic<unsigned long> fMessagesTx;</div>
|
|
<div class="line"><a name="l00518"></a><span class="lineno"> 518</span>  std::atomic<unsigned long> fMessagesRx;</div>
|
|
<div class="line"><a name="l00519"></a><span class="lineno"> 519</span>  </div>
|
|
<div class="line"><a name="l00520"></a><span class="lineno"> 520</span>  <span class="keywordtype">int</span> fTimeout;</div>
|
|
<div class="line"><a name="l00521"></a><span class="lineno"> 521</span> };</div>
|
|
<div class="line"><a name="l00522"></a><span class="lineno"> 522</span>  </div>
|
|
<div class="line"><a name="l00523"></a><span class="lineno"> 523</span> } <span class="comment">// namespace fair::mq::shmem</span></div>
|
|
<div class="line"><a name="l00524"></a><span class="lineno"> 524</span>  </div>
|
|
<div class="line"><a name="l00525"></a><span class="lineno"> 525</span> <span class="preprocessor">#endif </span><span class="comment">/* FAIR_MQ_SHMEM_SOCKET_H_ */</span><span class="preprocessor"></span></div>
|
|
</div><!-- fragment --></div><!-- contents -->
|
|
<div class="ttc" id="aclassFairMQSocket_html"><div class="ttname"><a href="classFairMQSocket.html">FairMQSocket</a></div><div class="ttdef"><b>Definition:</b> FairMQSocket.h:36</div></div>
|
|
<div class="ttc" id="aclassfair_1_1mq_1_1shmem_1_1Socket_html"><div class="ttname"><a href="classfair_1_1mq_1_1shmem_1_1Socket.html">fair::mq::shmem::Socket</a></div><div class="ttdef"><b>Definition:</b> Socket.h:44</div></div>
|
|
<div class="ttc" id="aclassfair_1_1mq_1_1shmem_1_1Socket_html_acd1bc3ce745e748eefad50ac19c175dd"><div class="ttname"><a href="classfair_1_1mq_1_1shmem_1_1Socket.html#acd1bc3ce745e748eefad50ac19c175dd">fair::mq::shmem::Socket::Events</a></div><div class="ttdeci">void Events(uint32_t *events) override</div><div class="ttdef"><b>Definition:</b> Socket.h:381</div></div>
|
|
<div class="ttc" id="aclassfair_1_1mq_1_1shmem_1_1Manager_html"><div class="ttname"><a href="classfair_1_1mq_1_1shmem_1_1Manager.html">fair::mq::shmem::Manager</a></div><div class="ttdef"><b>Definition:</b> Manager.h:61</div></div>
|
|
<div class="ttc" id="astructfair_1_1mq_1_1SocketError_html"><div class="ttname"><a href="structfair_1_1mq_1_1SocketError.html">fair::mq::SocketError</a></div><div class="ttdef"><b>Definition:</b> FairMQSocket.h:92</div></div>
|
|
<div class="ttc" id="astructfair_1_1mq_1_1shmem_1_1ZMsg_html"><div class="ttname"><a href="structfair_1_1mq_1_1shmem_1_1ZMsg.html">fair::mq::shmem::ZMsg</a></div><div class="ttdef"><b>Definition:</b> Socket.h:31</div></div>
|
|
<div class="ttc" id="astructfair_1_1mq_1_1shmem_1_1MetaHeader_html"><div class="ttname"><a href="structfair_1_1mq_1_1shmem_1_1MetaHeader.html">fair::mq::shmem::MetaHeader</a></div><div class="ttdef"><b>Definition:</b> Common.h:132</div></div>
|
|
<div class="ttc" id="aclassfair_1_1mq_1_1shmem_1_1Message_html"><div class="ttname"><a href="classfair_1_1mq_1_1shmem_1_1Message.html">fair::mq::shmem::Message</a></div><div class="ttdef"><b>Definition:</b> Message.h:39</div></div>
|
|
<div class="ttc" id="anamespacefair_1_1mq_1_1shmem_html"><div class="ttname"><a href="namespacefair_1_1mq_1_1shmem.html">fair::mq::shmem</a></div><div class="ttdef"><b>Definition:</b> Common.h:33</div></div>
|
|
<div class="ttc" id="aclassFairMQTransportFactory_html"><div class="ttname"><a href="classFairMQTransportFactory.html">FairMQTransportFactory</a></div><div class="ttdef"><b>Definition:</b> FairMQTransportFactory.h:30</div></div>
|
|
<p style="margin: 0 12px 10px 12px;"><a href="https://help.github.com/articles/github-privacy-statement/">privacy</a></p>
|