Skip to content
Projects
Groups
Snippets
Help
Loading...
Sign in / Register
Toggle navigation
L
libzmq
Project
Project
Details
Activity
Cycle Analytics
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Charts
Issues
0
Issues
0
List
Board
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Charts
Packages
Packages
Wiki
Wiki
Snippets
Snippets
Members
Members
Collapse sidebar
Close sidebar
Activity
Graph
Charts
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
submodule
libzmq
Commits
c590873f
Commit
c590873f
authored
Aug 15, 2018
by
Simon Giesecke
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
Problem: complexity of start_connecting
Solution: extract functions for each protocol
parent
31b0a1df
Show whitespace changes
Inline
Side-by-side
Showing
2 changed files
with
160 additions
and
73 deletions
+160
-73
session_base.cpp
src/session_base.cpp
+133
-73
session_base.hpp
src/session_base.hpp
+27
-0
No files found.
src/session_base.cpp
View file @
c590873f
...
@@ -550,6 +550,51 @@ void zmq::session_base_t::reconnect ()
...
@@ -550,6 +550,51 @@ void zmq::session_base_t::reconnect ()
_pipe
->
hiccup
();
_pipe
->
hiccup
();
}
}
zmq
::
session_base_t
::
connecter_factory_entry_t
zmq
::
session_base_t
::
_connecter_factories
[]
=
{
std
::
make_pair
(
protocol_name
::
tcp
,
&
zmq
::
session_base_t
::
create_connecter_tcp
),
#if !defined ZMQ_HAVE_WINDOWS && !defined ZMQ_HAVE_OPENVMS \
&& !defined ZMQ_HAVE_VXWORKS
std
::
make_pair
(
protocol_name
::
ipc
,
&
zmq
::
session_base_t
::
create_connecter_ipc
),
#endif
#if defined ZMQ_HAVE_TIPC
std
::
make_pair
(
protocol_name
::
tipc
,
&
zmq
::
session_base_t
::
create_connecter_tipc
),
#endif
#if defined ZMQ_HAVE_VMCI
std
::
make_pair
(
protocol_name
::
vmci
,
&
zmq
::
session_base_t
::
create_connecter_vmci
),
#endif
};
zmq
::
session_base_t
::
connecter_factory_map_t
zmq
::
session_base_t
::
_connecter_factories_map
(
_connecter_factories
,
_connecter_factories
+
sizeof
(
_connecter_factories
)
/
sizeof
(
_connecter_factories
[
0
]));
zmq
::
session_base_t
::
start_connecting_entry_t
zmq
::
session_base_t
::
_start_connecting_entries
[]
=
{
std
::
make_pair
(
protocol_name
::
udp
,
&
zmq
::
session_base_t
::
start_connecting_udp
),
#if defined ZMQ_HAVE_OPENPGM
std
::
make_pair
(
"pgm"
,
&
zmq
::
session_base_t
::
start_connecting_pgm
),
std
::
make_pair
(
"epgm"
,
&
zmq
::
session_base_t
::
start_connecting_pgm
),
#endif
#if defined ZMQ_HAVE_NORM
std
::
make_pair
(
"norm"
,
&
zmq
::
session_base_t
::
start_connecting_norm
),
#endif
};
zmq
::
session_base_t
::
start_connecting_map_t
zmq
::
session_base_t
::
_start_connecting_map
(
_start_connecting_entries
,
_start_connecting_entries
+
sizeof
(
_start_connecting_entries
)
/
sizeof
(
_start_connecting_entries
[
0
]));
void
zmq
::
session_base_t
::
start_connecting
(
bool
wait_
)
void
zmq
::
session_base_t
::
start_connecting
(
bool
wait_
)
{
{
zmq_assert
(
_active
);
zmq_assert
(
_active
);
...
@@ -560,78 +605,71 @@ void zmq::session_base_t::start_connecting (bool wait_)
...
@@ -560,78 +605,71 @@ void zmq::session_base_t::start_connecting (bool wait_)
zmq_assert
(
io_thread
);
zmq_assert
(
io_thread
);
// Create the connecter object.
// Create the connecter object.
own_t
*
connecter
=
NULL
;
const
connecter_factory_map_t
::
const_iterator
connecter_factories_it
=
if
(
_addr
->
protocol
==
protocol_name
::
tcp
)
{
_connecter_factories_map
.
find
(
_addr
->
protocol
);
if
(
!
options
.
socks_proxy_address
.
empty
())
{
if
(
connecter_factories_it
!=
_connecter_factories_map
.
end
())
{
address_t
*
proxy_address
=
new
(
std
::
nothrow
)
own_t
*
connecter
=
address_t
(
protocol_name
::
tcp
,
options
.
socks_proxy_address
,
(
this
->*
connecter_factories_it
->
second
)
(
io_thread
,
wait_
);
this
->
get_ctx
());
alloc_assert
(
proxy_address
);
connecter
=
new
(
std
::
nothrow
)
socks_connecter_t
(
io_thread
,
this
,
options
,
_addr
,
proxy_address
,
wait_
);
}
else
{
connecter
=
new
(
std
::
nothrow
)
tcp_connecter_t
(
io_thread
,
this
,
options
,
_addr
,
wait_
);
}
}
#if !defined ZMQ_HAVE_WINDOWS && !defined ZMQ_HAVE_OPENVMS \
&& !defined ZMQ_HAVE_VXWORKS
else
if
(
_addr
->
protocol
==
protocol_name
::
ipc
)
{
connecter
=
new
(
std
::
nothrow
)
ipc_connecter_t
(
io_thread
,
this
,
options
,
_addr
,
wait_
);
}
#endif
#if defined ZMQ_HAVE_TIPC
else
if
(
_addr
->
protocol
==
protocol_name
::
tipc
)
{
connecter
=
new
(
std
::
nothrow
)
tipc_connecter_t
(
io_thread
,
this
,
options
,
_addr
,
wait_
);
}
#endif
#if defined ZMQ_HAVE_VMCI
else
if
(
_addr
->
protocol
==
protocol_name
::
vmci
)
{
connecter
=
new
(
std
::
nothrow
)
vmci_connecter_t
(
io_thread
,
this
,
options
,
_addr
,
wait_
);
}
#endif
if
(
connecter
!=
NULL
)
{
alloc_assert
(
connecter
);
alloc_assert
(
connecter
);
launch_child
(
connecter
);
launch_child
(
connecter
);
return
;
return
;
}
}
const
start_connecting_map_t
::
const_iterator
start_connecting_it
=
_start_connecting_map
.
find
(
_addr
->
protocol
);
if
(
start_connecting_it
!=
_start_connecting_map
.
end
())
{
(
this
->*
start_connecting_it
->
second
)
(
io_thread
);
return
;
}
if
(
_addr
->
protocol
==
protocol_name
::
udp
)
{
zmq_assert
(
false
);
zmq_assert
(
options
.
type
==
ZMQ_DISH
||
options
.
type
==
ZMQ_RADIO
}
||
options
.
type
==
ZMQ_DGRAM
);
udp_engine_t
*
engine
=
new
(
std
::
nothrow
)
udp_engine_t
(
options
);
alloc_assert
(
engine
);
bool
recv
=
false
;
bool
send
=
false
;
if
(
options
.
type
==
ZMQ_RADIO
)
{
#if defined ZMQ_HAVE_VMCI
send
=
true
;
zmq
::
own_t
*
zmq
::
session_base_t
::
create_connecter_vmci
(
io_thread_t
*
io_thread_
,
recv
=
false
;
bool
wait_
)
}
else
if
(
options
.
type
==
ZMQ_DISH
)
{
{
send
=
false
;
return
new
(
std
::
nothrow
)
recv
=
true
;
vmci_connecter_t
(
io_thread_
,
this
,
options
,
_addr
,
wait_
);
}
else
if
(
options
.
type
==
ZMQ_DGRAM
)
{
}
send
=
true
;
#endif
recv
=
true
;
}
int
rc
=
engine
->
init
(
_addr
,
send
,
recv
);
#if defined ZMQ_HAVE_TIPC
errno_assert
(
rc
==
0
);
zmq
::
own_t
*
zmq
::
session_base_t
::
create_connecter_tipc
(
io_thread_t
*
io_thread_
,
bool
wait_
)
{
return
new
(
std
::
nothrow
)
tipc_connecter_t
(
io_thread_
,
this
,
options
,
_addr
,
wait_
);
}
#endif
send_attach
(
this
,
engine
);
#if !defined ZMQ_HAVE_WINDOWS && !defined ZMQ_HAVE_OPENVMS \
&& !defined ZMQ_HAVE_VXWORKS
zmq
::
own_t
*
zmq
::
session_base_t
::
create_connecter_ipc
(
io_thread_t
*
io_thread_
,
bool
wait_
)
{
return
new
(
std
::
nothrow
)
ipc_connecter_t
(
io_thread_
,
this
,
options
,
_addr
,
wait_
);
}
#endif
return
;
zmq
::
own_t
*
zmq
::
session_base_t
::
create_connecter_tcp
(
io_thread_t
*
io_thread_
,
bool
wait_
)
{
if
(
!
options
.
socks_proxy_address
.
empty
())
{
address_t
*
proxy_address
=
new
(
std
::
nothrow
)
address_t
(
protocol_name
::
tcp
,
options
.
socks_proxy_address
,
this
->
get_ctx
());
alloc_assert
(
proxy_address
);
return
new
(
std
::
nothrow
)
socks_connecter_t
(
io_thread_
,
this
,
options
,
_addr
,
proxy_address
,
wait_
);
}
}
return
new
(
std
::
nothrow
)
tcp_connecter_t
(
io_thread_
,
this
,
options
,
_addr
,
wait_
);
}
#ifdef ZMQ_HAVE_OPENPGM
#ifdef ZMQ_HAVE_OPENPGM
void
zmq
::
session_base_t
::
start_connecting_pgm
(
io_thread_t
*
io_thread_
)
// Both PGM and EPGM transports are using the same infrastructure.
{
if
(
_addr
->
protocol
==
"pgm"
||
_addr
->
protocol
==
"epgm"
)
{
zmq_assert
(
options
.
type
==
ZMQ_PUB
||
options
.
type
==
ZMQ_XPUB
zmq_assert
(
options
.
type
==
ZMQ_PUB
||
options
.
type
==
ZMQ_XPUB
||
options
.
type
==
ZMQ_SUB
||
options
.
type
==
ZMQ_XSUB
);
||
options
.
type
==
ZMQ_SUB
||
options
.
type
==
ZMQ_XSUB
);
...
@@ -644,18 +682,17 @@ void zmq::session_base_t::start_connecting (bool wait_)
...
@@ -644,18 +682,17 @@ void zmq::session_base_t::start_connecting (bool wait_)
if
(
options
.
type
==
ZMQ_PUB
||
options
.
type
==
ZMQ_XPUB
)
{
if
(
options
.
type
==
ZMQ_PUB
||
options
.
type
==
ZMQ_XPUB
)
{
// PGM sender.
// PGM sender.
pgm_sender_t
*
pgm_sender
=
pgm_sender_t
*
pgm_sender
=
new
(
std
::
nothrow
)
pgm_sender_t
(
io_thread
,
options
);
new
(
std
::
nothrow
)
pgm_sender_t
(
io_thread_
,
options
);
alloc_assert
(
pgm_sender
);
alloc_assert
(
pgm_sender
);
int
rc
=
int
rc
=
pgm_sender
->
init
(
udp_encapsulation
,
_addr
->
address
.
c_str
());
pgm_sender
->
init
(
udp_encapsulation
,
_addr
->
address
.
c_str
());
errno_assert
(
rc
==
0
);
errno_assert
(
rc
==
0
);
send_attach
(
this
,
pgm_sender
);
send_attach
(
this
,
pgm_sender
);
}
else
{
}
else
{
// PGM receiver.
// PGM receiver.
pgm_receiver_t
*
pgm_receiver
=
pgm_receiver_t
*
pgm_receiver
=
new
(
std
::
nothrow
)
pgm_receiver_t
(
io_thread
,
options
);
new
(
std
::
nothrow
)
pgm_receiver_t
(
io_thread_
,
options
);
alloc_assert
(
pgm_receiver
);
alloc_assert
(
pgm_receiver
);
int
rc
=
int
rc
=
...
@@ -664,20 +701,19 @@ void zmq::session_base_t::start_connecting (bool wait_)
...
@@ -664,20 +701,19 @@ void zmq::session_base_t::start_connecting (bool wait_)
send_attach
(
this
,
pgm_receiver
);
send_attach
(
this
,
pgm_receiver
);
}
}
}
return
;
}
#endif
#endif
#ifdef ZMQ_HAVE_NORM
#ifdef ZMQ_HAVE_NORM
if
(
_addr
->
protocol
==
"norm"
)
{
void
zmq
::
session_base_t
::
start_connecting_norm
(
io_thread_t
*
io_thread_
)
{
// At this point we'll create message pipes to the session straight
// At this point we'll create message pipes to the session straight
// away. There's no point in delaying it as no concept of 'connect'
// away. There's no point in delaying it as no concept of 'connect'
// exists with NORM anyway.
// exists with NORM anyway.
if
(
options
.
type
==
ZMQ_PUB
||
options
.
type
==
ZMQ_XPUB
)
{
if
(
options
.
type
==
ZMQ_PUB
||
options
.
type
==
ZMQ_XPUB
)
{
// NORM sender.
// NORM sender.
norm_engine_t
*
norm_sender
=
norm_engine_t
*
norm_sender
=
new
(
std
::
nothrow
)
norm_engine_t
(
io_thread
,
options
);
new
(
std
::
nothrow
)
norm_engine_t
(
io_thread_
,
options
);
alloc_assert
(
norm_sender
);
alloc_assert
(
norm_sender
);
int
rc
=
norm_sender
->
init
(
_addr
->
address
.
c_str
(),
true
,
false
);
int
rc
=
norm_sender
->
init
(
_addr
->
address
.
c_str
(),
true
,
false
);
...
@@ -688,7 +724,7 @@ void zmq::session_base_t::start_connecting (bool wait_)
...
@@ -688,7 +724,7 @@ void zmq::session_base_t::start_connecting (bool wait_)
// NORM receiver.
// NORM receiver.
norm_engine_t
*
norm_receiver
=
norm_engine_t
*
norm_receiver
=
new
(
std
::
nothrow
)
norm_engine_t
(
io_thread
,
options
);
new
(
std
::
nothrow
)
norm_engine_t
(
io_thread_
,
options
);
alloc_assert
(
norm_receiver
);
alloc_assert
(
norm_receiver
);
int
rc
=
norm_receiver
->
init
(
_addr
->
address
.
c_str
(),
false
,
true
);
int
rc
=
norm_receiver
->
init
(
_addr
->
address
.
c_str
(),
false
,
true
);
...
@@ -696,9 +732,33 @@ void zmq::session_base_t::start_connecting (bool wait_)
...
@@ -696,9 +732,33 @@ void zmq::session_base_t::start_connecting (bool wait_)
send_attach
(
this
,
norm_receiver
);
send_attach
(
this
,
norm_receiver
);
}
}
return
;
}
#endif
void
zmq
::
session_base_t
::
start_connecting_udp
(
io_thread_t
*
/*io_thread_*/
)
{
zmq_assert
(
options
.
type
==
ZMQ_DISH
||
options
.
type
==
ZMQ_RADIO
||
options
.
type
==
ZMQ_DGRAM
);
udp_engine_t
*
engine
=
new
(
std
::
nothrow
)
udp_engine_t
(
options
);
alloc_assert
(
engine
);
bool
recv
=
false
;
bool
send
=
false
;
if
(
options
.
type
==
ZMQ_RADIO
)
{
send
=
true
;
recv
=
false
;
}
else
if
(
options
.
type
==
ZMQ_DISH
)
{
send
=
false
;
recv
=
true
;
}
else
if
(
options
.
type
==
ZMQ_DGRAM
)
{
send
=
true
;
recv
=
true
;
}
}
#endif // ZMQ_HAVE_NORM
zmq_assert
(
false
);
const
int
rc
=
engine
->
init
(
_addr
,
send
,
recv
);
errno_assert
(
rc
==
0
);
send_attach
(
this
,
engine
);
}
}
src/session_base.hpp
View file @
c590873f
...
@@ -105,6 +105,33 @@ class session_base_t : public own_t, public io_object_t, public i_pipe_events
...
@@ -105,6 +105,33 @@ class session_base_t : public own_t, public io_object_t, public i_pipe_events
private
:
private
:
void
start_connecting
(
bool
wait_
);
void
start_connecting
(
bool
wait_
);
typedef
own_t
*
(
session_base_t
::*
connecter_factory_fun_t
)
(
io_thread_t
*
io_thread
,
bool
wait_
);
typedef
std
::
pair
<
std
::
string
,
connecter_factory_fun_t
>
connecter_factory_entry_t
;
static
connecter_factory_entry_t
_connecter_factories
[];
typedef
std
::
map
<
std
::
string
,
connecter_factory_fun_t
>
connecter_factory_map_t
;
static
connecter_factory_map_t
_connecter_factories_map
;
own_t
*
create_connecter_vmci
(
io_thread_t
*
io_thread_
,
bool
wait_
);
own_t
*
create_connecter_tipc
(
io_thread_t
*
io_thread_
,
bool
wait_
);
own_t
*
create_connecter_ipc
(
io_thread_t
*
io_thread_
,
bool
wait_
);
own_t
*
create_connecter_tcp
(
io_thread_t
*
io_thread_
,
bool
wait_
);
typedef
void
(
session_base_t
::*
start_connecting_fun_t
)
(
io_thread_t
*
io_thread
);
typedef
std
::
pair
<
std
::
string
,
start_connecting_fun_t
>
start_connecting_entry_t
;
static
start_connecting_entry_t
_start_connecting_entries
[];
typedef
std
::
map
<
std
::
string
,
start_connecting_fun_t
>
start_connecting_map_t
;
static
start_connecting_map_t
_start_connecting_map
;
void
start_connecting_pgm
(
io_thread_t
*
io_thread_
);
void
start_connecting_norm
(
io_thread_t
*
io_thread_
);
void
start_connecting_udp
(
io_thread_t
*
io_thread_
);
void
reconnect
();
void
reconnect
();
// Handlers for incoming commands.
// Handlers for incoming commands.
...
...
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment