mirror of
https://github.com/zeromq/libzmq.git
synced 2025-10-19 12:42:32 +02:00
Plug in dbuffer to serve the ZMQ_CONFLATE option
ZMQ_CONFLATE option is passed to pipepair() which creates a usual ypipe_t or ypipe_conflate_t and plugs it into pipe_t under a common abstract base.
This commit is contained in:
@@ -300,7 +300,8 @@ int zmq::session_base_t::zap_connect ()
|
||||
pipe_t *new_pipes [2] = {NULL, NULL};
|
||||
int hwms [2] = {0, 0};
|
||||
bool delays [2] = {false, false};
|
||||
int rc = pipepair (parents, new_pipes, hwms, delays);
|
||||
bool conflates [2] = {false, false};
|
||||
int rc = pipepair (parents, new_pipes, hwms, delays, conflates);
|
||||
errno_assert (rc == 0);
|
||||
|
||||
// Attach local end of the pipe to this socket object.
|
||||
@@ -331,9 +332,19 @@ void zmq::session_base_t::process_attach (i_engine *engine_)
|
||||
if (!pipe && !is_terminating ()) {
|
||||
object_t *parents [2] = {this, socket};
|
||||
pipe_t *pipes [2] = {NULL, NULL};
|
||||
int hwms [2] = {options.rcvhwm, options.sndhwm};
|
||||
|
||||
bool conflate = options.conflate &&
|
||||
(options.type == ZMQ_DEALER ||
|
||||
options.type == ZMQ_PULL ||
|
||||
options.type == ZMQ_PUSH ||
|
||||
options.type == ZMQ_PUB ||
|
||||
options.type == ZMQ_SUB);
|
||||
|
||||
int hwms [2] = {conflate? -1 : options.rcvhwm,
|
||||
conflate? -1 : options.sndhwm};
|
||||
bool delays [2] = {options.delay_on_close, options.delay_on_disconnect};
|
||||
int rc = pipepair (parents, pipes, hwms, delays);
|
||||
bool conflates [2] = {conflate, conflate};
|
||||
int rc = pipepair (parents, pipes, hwms, delays, conflates);
|
||||
errno_assert (rc == 0);
|
||||
|
||||
// Plug the local end of the pipe.
|
||||
|
Reference in New Issue
Block a user