Skip to content

Commit

Permalink
Merge pull request #113 from rikvdh/master
Browse files Browse the repository at this point in the history
Backport issue #1382 & fix for issue #88
  • Loading branch information
hintjens committed May 1, 2015
2 parents 399b71c + 4bba965 commit 95ba372
Show file tree
Hide file tree
Showing 6 changed files with 126 additions and 4 deletions.
1 change: 1 addition & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,7 @@ tests/test_srcfd
tests/test_stream_disconnect
tests/test_proxy_chain
tests/test_bind_src_address
tests/test_proxy_terminate
tests/test*.log
tests/test*.trs
src/platform.hpp*
Expand Down
1 change: 1 addition & 0 deletions AUTHORS
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,7 @@ Philip Kovacs <phil@philkovacs.com>
Pieter Hintjens <ph@imatix.com>
Piotr Trojanek <piotr.trojanek@gmail.com>
Richard Newton <richard_newton@waters.com>
Rik van der Heijden <mail@rikvanderheijden.com>
Robert G. Jakabosky <bobby@sharedrealm.com>
Sebastian Otaegui <feniix@gmail.com>
Stefan Radomski <radomski@tk.informatik.tu-darmstadt.de>
Expand Down
1 change: 1 addition & 0 deletions CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -614,6 +614,7 @@ set(tests
test_inproc_connect
test_issue_566
test_many_sockets
test_proxy_terminate
)
if(NOT WIN32)
list(APPEND tests
Expand Down
12 changes: 8 additions & 4 deletions src/proxy.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -76,7 +76,7 @@ int zmq::proxy (
{ control_, 0, ZMQ_POLLIN, 0 }
};
int qt_poll_items = (control_ ? 3 : 2);

// Proxy can be in these three states
enum {
active,
Expand All @@ -89,7 +89,7 @@ int zmq::proxy (
rc = zmq_poll (&items [0], qt_poll_items, -1);
if (unlikely (rc < 0))
return -1;

// Process a control command if any
if (control_ && items [2].revents & ZMQ_POLLIN) {
rc = control_->recv (&msg, 0);
Expand Down Expand Up @@ -130,7 +130,9 @@ int zmq::proxy (
}
}
// Process a request
if (items [0].revents & ZMQ_POLLIN) {
if (state == active
&& items [0].revents & ZMQ_POLLIN
&& items [1].revents & ZMQ_POLLOUT) {
while (true) {
rc = frontend_->recv (&msg, 0);
if (unlikely (rc < 0))
Expand Down Expand Up @@ -162,7 +164,9 @@ int zmq::proxy (
}
}
// Process a reply
if (items [1].revents & ZMQ_POLLIN) {
if (state == active
&& items [1].revents & ZMQ_POLLIN
&& items [0].revents & ZMQ_POLLOUT) {
while (true) {
rc = backend_->recv (&msg, 0);
if (unlikely (rc < 0))
Expand Down
2 changes: 2 additions & 0 deletions tests/Makefile.am
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@ noinst_PROGRAMS = test_system \
test_inproc_connect \
test_issue_566 \
test_abstract_ipc \
test_proxy_terminate \
test_many_sockets

if !ON_MINGW
Expand Down Expand Up @@ -88,6 +89,7 @@ test_inproc_connect_SOURCES = test_inproc_connect.cpp
test_issue_566_SOURCES = test_issue_566.cpp
test_abstract_ipc_SOURCES = test_abstract_ipc.cpp
test_many_sockets_SOURCES = test_many_sockets.cpp
test_proxy_terminate_SOURCES = test_proxy_terminate.cpp
if !ON_MINGW
test_shutdown_stress_SOURCES = test_shutdown_stress.cpp
test_pair_ipc_SOURCES = test_pair_ipc.cpp testutil.hpp
Expand Down
113 changes: 113 additions & 0 deletions tests/test_proxy_terminate.cpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,113 @@
/*
Copyright (c) 2007-2015 Contributors as noted in the AUTHORS file
This file is part of 0MQ.
0MQ is free software; you can redistribute it and/or modify it under
the terms of the GNU Lesser General Public License as published by
the Free Software Foundation; either version 3 of the License, or
(at your option) any later version.
0MQ is distributed in the hope that it will be useful,
but WITHOUT ANY WARRANTY; without even the implied warranty of
MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
GNU Lesser General Public License for more details.
You should have received a copy of the GNU Lesser General Public License
along with this program. If not, see <http://www.gnu.org/licenses/>.
*/

#include "testutil.hpp"
#include "../include/zmq_utils.h"

// This is a test for issue #1382. The server thread creates a SUB-PUSH
// steerable proxy. The main process then sends messages to the SUB
// but there is no pull on the other side, previously the proxy blocks
// in writing to the backend, preventing the proxy from terminating

void
server_task (void *ctx)
{
// Frontend socket talks to main process
void *frontend = zmq_socket (ctx, ZMQ_SUB);
assert (frontend);
int rc = zmq_setsockopt (frontend, ZMQ_SUBSCRIBE, "", 0);
assert (rc == 0);
rc = zmq_bind (frontend, "tcp://127.0.0.1:15564");
assert (rc == 0);

// Nice socket which is never read
void *backend = zmq_socket (ctx, ZMQ_PUSH);
assert (backend);
rc = zmq_bind (backend, "tcp://127.0.0.1:15563");
assert (rc == 0);

// Control socket receives terminate command from main over inproc
void *control = zmq_socket (ctx, ZMQ_SUB);
assert (control);
rc = zmq_setsockopt (control, ZMQ_SUBSCRIBE, "", 0);
assert (rc == 0);
rc = zmq_connect (control, "inproc://control");
assert (rc == 0);

// Connect backend to frontend via a proxy
zmq_proxy_steerable (frontend, backend, NULL, control);

rc = zmq_close (frontend);
assert (rc == 0);
rc = zmq_close (backend);
assert (rc == 0);
rc = zmq_close (control);
assert (rc == 0);
}


// The main thread simply starts a basic steerable proxy server, publishes some messages, and then
// waits for the server to terminate.

int main (void)
{
setup_test_environment ();

void *ctx = zmq_ctx_new ();
assert (ctx);
// Control socket receives terminate command from main over inproc
void *control = zmq_socket (ctx, ZMQ_PUB);
assert (control);
int rc = zmq_bind (control, "inproc://control");
assert (rc == 0);

void *thread = zmq_threadstart(&server_task, ctx);
msleep (500); // Run for 500 ms

// Start a secondary publisher which writes data to the SUB-PUSH server socket
void *publisher = zmq_socket (ctx, ZMQ_PUB);
assert (publisher);
rc = zmq_connect (publisher, "tcp://127.0.0.1:15564");
assert (rc == 0);

msleep (50);
rc = zmq_send (publisher, "This is a test", 14, 0);
assert (rc == 14);

msleep (50);
rc = zmq_send (publisher, "This is a test", 14, 0);
assert (rc == 14);

msleep (50);
rc = zmq_send (publisher, "This is a test", 14, 0);
assert (rc == 14);
rc = zmq_send (control, "TERMINATE", 9, 0);
assert (rc == 9);

rc = zmq_close (publisher);
assert (rc == 0);
rc = zmq_close (control);
assert (rc == 0);

zmq_threadclose (thread);

rc = zmq_ctx_term (ctx);
assert (rc == 0);
return 0;
}

0 comments on commit 95ba372

Please sign in to comment.