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
dfdaff5e
Commit
dfdaff5e
authored
Mar 20, 2010
by
Martin Sustrik
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
XREP-style prefixing/trimming messages removed
parent
cbaf1097
Hide whitespace changes
Inline
Side-by-side
Showing
16 changed files
with
17 additions
and
145 deletions
+17
-145
i_engine.hpp
src/i_engine.hpp
+2
-8
options.cpp
src/options.cpp
+1
-2
options.hpp
src/options.hpp
+0
-3
pgm_receiver.cpp
src/pgm_receiver.cpp
+0
-12
pgm_receiver.hpp
src/pgm_receiver.hpp
+0
-2
pgm_sender.cpp
src/pgm_sender.cpp
+0
-12
pgm_sender.hpp
src/pgm_sender.hpp
+0
-2
session.cpp
src/session.cpp
+0
-5
xrep.cpp
src/xrep.cpp
+2
-3
zmq_decoder.cpp
src/zmq_decoder.cpp
+7
-41
zmq_decoder.hpp
src/zmq_decoder.hpp
+0
-8
zmq_encoder.cpp
src/zmq_encoder.cpp
+4
-27
zmq_encoder.hpp
src/zmq_encoder.hpp
+0
-6
zmq_engine.cpp
src/zmq_engine.cpp
+0
-10
zmq_engine.hpp
src/zmq_engine.hpp
+0
-2
zmq_init.cpp
src/zmq_init.cpp
+1
-2
No files found.
src/i_engine.hpp
View file @
dfdaff5e
...
@@ -22,8 +22,6 @@
...
@@ -22,8 +22,6 @@
#include <stddef.h>
#include <stddef.h>
#include "blob.hpp"
namespace
zmq
namespace
zmq
{
{
...
@@ -41,13 +39,9 @@ namespace zmq
...
@@ -41,13 +39,9 @@ namespace zmq
// are messages to send available.
// are messages to send available.
virtual
void
revive
()
=
0
;
virtual
void
revive
()
=
0
;
// This method is called by the session to signalise that more
// messages can be written to the pipe.
virtual
void
resume_input
()
=
0
;
virtual
void
resume_input
()
=
0
;
// Engine should add the prefix supplied to all inbound messages.
virtual
void
add_prefix
(
const
blob_t
&
identity_
)
=
0
;
// Engine should trim prefix from all the outbound messages.
virtual
void
trim_prefix
()
=
0
;
};
};
}
}
...
...
src/options.cpp
View file @
dfdaff5e
...
@@ -34,8 +34,7 @@ zmq::options_t::options_t () :
...
@@ -34,8 +34,7 @@ zmq::options_t::options_t () :
rcvbuf
(
0
),
rcvbuf
(
0
),
requires_in
(
false
),
requires_in
(
false
),
requires_out
(
false
),
requires_out
(
false
),
immediate_connect
(
true
),
immediate_connect
(
true
)
traceroute
(
false
)
{
{
}
}
...
...
src/options.hpp
View file @
dfdaff5e
...
@@ -61,9 +61,6 @@ namespace zmq
...
@@ -61,9 +61,6 @@ namespace zmq
// is not aware of the peer's identity, however, it is able to send
// is not aware of the peer's identity, however, it is able to send
// messages straight away.
// messages straight away.
bool
immediate_connect
;
bool
immediate_connect
;
// If true, socket requires tracerouting the messages.
bool
traceroute
;
};
};
}
}
...
...
src/pgm_receiver.cpp
View file @
dfdaff5e
...
@@ -121,18 +121,6 @@ void zmq::pgm_receiver_t::resume_input ()
...
@@ -121,18 +121,6 @@ void zmq::pgm_receiver_t::resume_input ()
in_event
();
in_event
();
}
}
void
zmq
::
pgm_receiver_t
::
add_prefix
(
const
blob_t
&
identity_
)
{
// No need for tracerouting functionality in PGM socket at the moment.
zmq_assert
(
false
);
}
void
zmq
::
pgm_receiver_t
::
trim_prefix
()
{
// No need for tracerouting functionality in PGM socket at the moment.
zmq_assert
(
false
);
}
void
zmq
::
pgm_receiver_t
::
in_event
()
void
zmq
::
pgm_receiver_t
::
in_event
()
{
{
// Read data from the underlying pgm_socket.
// Read data from the underlying pgm_socket.
...
...
src/pgm_receiver.hpp
View file @
dfdaff5e
...
@@ -55,8 +55,6 @@ namespace zmq
...
@@ -55,8 +55,6 @@ namespace zmq
void
unplug
();
void
unplug
();
void
revive
();
void
revive
();
void
resume_input
();
void
resume_input
();
void
add_prefix
(
const
blob_t
&
identity_
);
void
trim_prefix
();
// i_poll_events interface implementation.
// i_poll_events interface implementation.
void
in_event
();
void
in_event
();
...
...
src/pgm_sender.cpp
View file @
dfdaff5e
...
@@ -107,18 +107,6 @@ void zmq::pgm_sender_t::resume_input ()
...
@@ -107,18 +107,6 @@ void zmq::pgm_sender_t::resume_input ()
zmq_assert
(
false
);
zmq_assert
(
false
);
}
}
void
zmq
::
pgm_sender_t
::
add_prefix
(
const
blob_t
&
identity_
)
{
// No need for tracerouting functionality in PGM socket at the moment.
zmq_assert
(
false
);
}
void
zmq
::
pgm_sender_t
::
trim_prefix
()
{
// No need for tracerouting functionality in PGM socket at the moment.
zmq_assert
(
false
);
}
zmq
::
pgm_sender_t
::~
pgm_sender_t
()
zmq
::
pgm_sender_t
::~
pgm_sender_t
()
{
{
if
(
out_buffer
)
{
if
(
out_buffer
)
{
...
...
src/pgm_sender.hpp
View file @
dfdaff5e
...
@@ -53,8 +53,6 @@ namespace zmq
...
@@ -53,8 +53,6 @@ namespace zmq
void
unplug
();
void
unplug
();
void
revive
();
void
revive
();
void
resume_input
();
void
resume_input
();
void
add_prefix
(
const
blob_t
&
identity_
);
void
trim_prefix
();
// i_poll_events interface implementation.
// i_poll_events interface implementation.
void
in_event
();
void
in_event
();
...
...
src/session.cpp
View file @
dfdaff5e
...
@@ -264,9 +264,4 @@ void zmq::session_t::process_attach (i_engine *engine_,
...
@@ -264,9 +264,4 @@ void zmq::session_t::process_attach (i_engine *engine_,
zmq_assert
(
engine_
);
zmq_assert
(
engine_
);
engine
=
engine_
;
engine
=
engine_
;
engine
->
plug
(
this
);
engine
->
plug
(
this
);
// Once the initial handshaking is over tracerouting should trim prefixes
// from outbound messages.
if
(
options
.
traceroute
)
engine
->
trim_prefix
();
}
}
src/xrep.cpp
View file @
dfdaff5e
...
@@ -33,9 +33,8 @@ zmq::xrep_t::xrep_t (class app_thread_t *parent_) :
...
@@ -33,9 +33,8 @@ zmq::xrep_t::xrep_t (class app_thread_t *parent_) :
// That way we are aware of the peer's identity when binding to the pipes.
// That way we are aware of the peer's identity when binding to the pipes.
options
.
immediate_connect
=
false
;
options
.
immediate_connect
=
false
;
// XREP socket adds identity to inbound messages and strips identity
// XREP is unfunctional at the moment. Crash here!
// from the outbound messages.
zmq_assert
(
false
);
options
.
traceroute
=
true
;
}
}
zmq
::
xrep_t
::~
xrep_t
()
zmq
::
xrep_t
::~
xrep_t
()
...
...
src/zmq_decoder.cpp
View file @
dfdaff5e
...
@@ -45,11 +45,6 @@ void zmq::zmq_decoder_t::set_inout (i_inout *destination_)
...
@@ -45,11 +45,6 @@ void zmq::zmq_decoder_t::set_inout (i_inout *destination_)
destination
=
destination_
;
destination
=
destination_
;
}
}
void
zmq
::
zmq_decoder_t
::
add_prefix
(
const
blob_t
&
prefix_
)
{
prefix
=
prefix_
;
}
bool
zmq
::
zmq_decoder_t
::
one_byte_size_ready
()
bool
zmq
::
zmq_decoder_t
::
one_byte_size_ready
()
{
{
// First byte of size is read. If it is 0xff read 8-byte size.
// First byte of size is read. If it is 0xff read 8-byte size.
...
@@ -64,19 +59,8 @@ bool zmq::zmq_decoder_t::one_byte_size_ready ()
...
@@ -64,19 +59,8 @@ bool zmq::zmq_decoder_t::one_byte_size_ready ()
// in_progress is initialised at this point so in theory we should
// in_progress is initialised at this point so in theory we should
// close it before calling zmq_msg_init_size, however, it's a 0-byte
// close it before calling zmq_msg_init_size, however, it's a 0-byte
// message and thus we can treat it as uninitialised...
// message and thus we can treat it as uninitialised...
if
(
prefix
.
empty
())
{
int
rc
=
zmq_msg_init_size
(
&
in_progress
,
*
tmpbuf
-
1
);
int
rc
=
zmq_msg_init_size
(
&
in_progress
,
*
tmpbuf
-
1
);
errno_assert
(
rc
==
0
);
errno_assert
(
rc
==
0
);
}
else
{
int
rc
=
zmq_msg_init_size
(
&
in_progress
,
1
+
prefix
.
size
()
+
*
tmpbuf
-
1
);
errno_assert
(
rc
==
0
);
unsigned
char
*
data
=
(
unsigned
char
*
)
zmq_msg_data
(
&
in_progress
);
*
data
=
(
unsigned
char
)
prefix
.
size
();
memcpy
(
data
+
1
,
prefix
.
data
(),
*
data
);
}
next_step
(
tmpbuf
,
1
,
&
zmq_decoder_t
::
flags_ready
);
next_step
(
tmpbuf
,
1
,
&
zmq_decoder_t
::
flags_ready
);
}
}
return
true
;
return
true
;
...
@@ -93,18 +77,8 @@ bool zmq::zmq_decoder_t::eight_byte_size_ready ()
...
@@ -93,18 +77,8 @@ bool zmq::zmq_decoder_t::eight_byte_size_ready ()
// in_progress is initialised at this point so in theory we should
// in_progress is initialised at this point so in theory we should
// close it before calling zmq_msg_init_size, however, it's a 0-byte
// close it before calling zmq_msg_init_size, however, it's a 0-byte
// message and thus we can treat it as uninitialised...
// message and thus we can treat it as uninitialised...
if
(
prefix
.
empty
())
{
int
rc
=
zmq_msg_init_size
(
&
in_progress
,
size
-
1
);
int
rc
=
zmq_msg_init_size
(
&
in_progress
,
size
-
1
);
errno_assert
(
rc
==
0
);
errno_assert
(
rc
==
0
);
}
else
{
int
rc
=
zmq_msg_init_size
(
&
in_progress
,
1
+
prefix
.
size
()
+
size
-
1
);
errno_assert
(
rc
==
0
);
unsigned
char
*
data
=
(
unsigned
char
*
)
zmq_msg_data
(
&
in_progress
);
*
data
=
(
unsigned
char
)
prefix
.
size
();
memcpy
(
data
+
1
,
prefix
.
data
(),
*
data
);
}
next_step
(
tmpbuf
,
1
,
&
zmq_decoder_t
::
flags_ready
);
next_step
(
tmpbuf
,
1
,
&
zmq_decoder_t
::
flags_ready
);
return
true
;
return
true
;
...
@@ -115,17 +89,9 @@ bool zmq::zmq_decoder_t::flags_ready ()
...
@@ -115,17 +89,9 @@ bool zmq::zmq_decoder_t::flags_ready ()
// Store the flags from the wire into the message structure.
// Store the flags from the wire into the message structure.
in_progress
.
flags
=
tmpbuf
[
0
];
in_progress
.
flags
=
tmpbuf
[
0
];
if
(
prefix
.
empty
())
{
next_step
(
zmq_msg_data
(
&
in_progress
),
zmq_msg_size
(
&
in_progress
),
next_step
(
zmq_msg_data
(
&
in_progress
),
zmq_msg_size
(
&
in_progress
),
&
zmq_decoder_t
::
message_ready
);
&
zmq_decoder_t
::
message_ready
);
}
else
{
next_step
((
unsigned
char
*
)
zmq_msg_data
(
&
in_progress
)
+
prefix
.
size
()
+
1
,
zmq_msg_size
(
&
in_progress
)
-
prefix
.
size
()
-
1
,
&
zmq_decoder_t
::
message_ready
);
}
return
true
;
return
true
;
}
}
...
...
src/zmq_decoder.hpp
View file @
dfdaff5e
...
@@ -33,17 +33,11 @@ namespace zmq
...
@@ -33,17 +33,11 @@ namespace zmq
{
{
public
:
public
:
// If prefix is not NULL, it will be glued to the beginning of every
// decoded message.
zmq_decoder_t
(
size_t
bufsize_
);
zmq_decoder_t
(
size_t
bufsize_
);
~
zmq_decoder_t
();
~
zmq_decoder_t
();
void
set_inout
(
struct
i_inout
*
destination_
);
void
set_inout
(
struct
i_inout
*
destination_
);
// Once called, all decoded messages will be prefixed by the specified
// prefix.
void
add_prefix
(
const
blob_t
&
prefix_
);
private
:
private
:
bool
one_byte_size_ready
();
bool
one_byte_size_ready
();
...
@@ -55,8 +49,6 @@ namespace zmq
...
@@ -55,8 +49,6 @@ namespace zmq
unsigned
char
tmpbuf
[
8
];
unsigned
char
tmpbuf
[
8
];
::
zmq_msg_t
in_progress
;
::
zmq_msg_t
in_progress
;
blob_t
prefix
;
zmq_decoder_t
(
const
zmq_decoder_t
&
);
zmq_decoder_t
(
const
zmq_decoder_t
&
);
void
operator
=
(
const
zmq_decoder_t
&
);
void
operator
=
(
const
zmq_decoder_t
&
);
};
};
...
...
src/zmq_encoder.cpp
View file @
dfdaff5e
...
@@ -23,8 +23,7 @@
...
@@ -23,8 +23,7 @@
zmq
::
zmq_encoder_t
::
zmq_encoder_t
(
size_t
bufsize_
)
:
zmq
::
zmq_encoder_t
::
zmq_encoder_t
(
size_t
bufsize_
)
:
encoder_t
<
zmq_encoder_t
>
(
bufsize_
),
encoder_t
<
zmq_encoder_t
>
(
bufsize_
),
source
(
NULL
),
source
(
NULL
)
trim
(
false
)
{
{
zmq_msg_init
(
&
in_progress
);
zmq_msg_init
(
&
in_progress
);
...
@@ -42,25 +41,11 @@ void zmq::zmq_encoder_t::set_inout (i_inout *source_)
...
@@ -42,25 +41,11 @@ void zmq::zmq_encoder_t::set_inout (i_inout *source_)
source
=
source_
;
source
=
source_
;
}
}
void
zmq
::
zmq_encoder_t
::
trim_prefix
()
{
trim
=
true
;
}
bool
zmq
::
zmq_encoder_t
::
size_ready
()
bool
zmq
::
zmq_encoder_t
::
size_ready
()
{
{
// Write message body into the buffer.
// Write message body into the buffer.
if
(
!
trim
)
{
next_step
(
zmq_msg_data
(
&
in_progress
),
zmq_msg_size
(
&
in_progress
),
next_step
(
zmq_msg_data
(
&
in_progress
),
zmq_msg_size
(
&
in_progress
),
&
zmq_encoder_t
::
message_ready
,
false
);
&
zmq_encoder_t
::
message_ready
,
false
);
}
else
{
size_t
prefix_size
=
*
(
unsigned
char
*
)
zmq_msg_data
(
&
in_progress
);
next_step
(
(
unsigned
char
*
)
zmq_msg_data
(
&
in_progress
)
+
prefix_size
+
1
,
zmq_msg_size
(
&
in_progress
)
-
prefix_size
-
1
,
&
zmq_encoder_t
::
message_ready
,
false
);
}
return
true
;
return
true
;
}
}
...
@@ -78,16 +63,8 @@ bool zmq::zmq_encoder_t::message_ready ()
...
@@ -78,16 +63,8 @@ bool zmq::zmq_encoder_t::message_ready ()
return
false
;
return
false
;
}
}
// Get the message size. If the prefix is not to be sent, adjust the
// Get the message size.
// size accordingly.
size_t
size
=
zmq_msg_size
(
&
in_progress
);
size_t
size
=
zmq_msg_size
(
&
in_progress
);
if
(
trim
)
{
zmq_assert
(
size
);
size_t
prefix_size
=
(
*
(
unsigned
char
*
)
zmq_msg_data
(
&
in_progress
))
+
1
;
zmq_assert
(
prefix_size
<=
size
);
size
-=
prefix_size
;
}
// Account for the 'flags' byte.
// Account for the 'flags' byte.
size
++
;
size
++
;
...
...
src/zmq_encoder.hpp
View file @
dfdaff5e
...
@@ -37,10 +37,6 @@ namespace zmq
...
@@ -37,10 +37,6 @@ namespace zmq
void
set_inout
(
struct
i_inout
*
source_
);
void
set_inout
(
struct
i_inout
*
source_
);
// Once called, encoder will start trimming frefixes from outbound
// messages.
void
trim_prefix
();
private
:
private
:
bool
size_ready
();
bool
size_ready
();
...
@@ -50,8 +46,6 @@ namespace zmq
...
@@ -50,8 +46,6 @@ namespace zmq
::
zmq_msg_t
in_progress
;
::
zmq_msg_t
in_progress
;
unsigned
char
tmpbuf
[
10
];
unsigned
char
tmpbuf
[
10
];
bool
trim
;
zmq_encoder_t
(
const
zmq_encoder_t
&
);
zmq_encoder_t
(
const
zmq_encoder_t
&
);
void
operator
=
(
const
zmq_encoder_t
&
);
void
operator
=
(
const
zmq_encoder_t
&
);
};
};
...
...
src/zmq_engine.cpp
View file @
dfdaff5e
...
@@ -169,16 +169,6 @@ void zmq::zmq_engine_t::resume_input ()
...
@@ -169,16 +169,6 @@ void zmq::zmq_engine_t::resume_input ()
in_event
();
in_event
();
}
}
void
zmq
::
zmq_engine_t
::
add_prefix
(
const
blob_t
&
identity_
)
{
decoder
.
add_prefix
(
identity_
);
}
void
zmq
::
zmq_engine_t
::
trim_prefix
()
{
encoder
.
trim_prefix
();
}
void
zmq
::
zmq_engine_t
::
error
()
void
zmq
::
zmq_engine_t
::
error
()
{
{
zmq_assert
(
inout
);
zmq_assert
(
inout
);
...
...
src/zmq_engine.hpp
View file @
dfdaff5e
...
@@ -48,8 +48,6 @@ namespace zmq
...
@@ -48,8 +48,6 @@ namespace zmq
void
unplug
();
void
unplug
();
void
revive
();
void
revive
();
void
resume_input
();
void
resume_input
();
void
add_prefix
(
const
blob_t
&
identity_
);
void
trim_prefix
();
// i_poll_events interface implementation.
// i_poll_events interface implementation.
void
in_event
();
void
in_event
();
...
...
src/zmq_init.cpp
View file @
dfdaff5e
...
@@ -85,8 +85,7 @@ bool zmq::zmq_init_t::write (::zmq_msg_t *msg_)
...
@@ -85,8 +85,7 @@ bool zmq::zmq_init_t::write (::zmq_msg_t *msg_)
peer_identity
.
assign
((
const
unsigned
char
*
)
zmq_msg_data
(
msg_
),
peer_identity
.
assign
((
const
unsigned
char
*
)
zmq_msg_data
(
msg_
),
zmq_msg_size
(
msg_
));
zmq_msg_size
(
msg_
));
}
}
if
(
options
.
traceroute
)
engine
->
add_prefix
(
peer_identity
);
received
=
true
;
received
=
true
;
return
true
;
return
true
;
...
...
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