Commit a24a7c15 authored by Martin Sustrik's avatar Martin Sustrik

Session termination induced by socket fixed

Signed-off-by: 's avatarMartin Sustrik <sustrik@250bpm.com>
parent 0b59866a
...@@ -301,8 +301,11 @@ void zmq::pipe_t::process_pipe_term_ack () ...@@ -301,8 +301,11 @@ void zmq::pipe_t::process_pipe_term_ack ()
delete this; delete this;
} }
void zmq::pipe_t::terminate () void zmq::pipe_t::terminate (bool delay_)
{ {
// Overload the value specified at pipe creation.
delay = delay_;
// If terminate was already called, we can ignore the duplicit invocation. // If terminate was already called, we can ignore the duplicit invocation.
if (state == terminated || state == double_terminated) if (state == terminated || state == double_terminated)
return; return;
...@@ -321,11 +324,15 @@ void zmq::pipe_t::terminate () ...@@ -321,11 +324,15 @@ void zmq::pipe_t::terminate ()
// There are still pending messages available, but the user calls // There are still pending messages available, but the user calls
// 'terminate'. We can act as if all the pending messages were read. // 'terminate'. We can act as if all the pending messages were read.
else if (state == pending) { else if (state == pending && !delay) {
send_pipe_term_ack (peer); send_pipe_term_ack (peer);
state = terminating; state = terminating;
} }
// If there are pending messages still availabe, do nothing.
else if (state == pending && delay) {
}
// We've already got delimiter, but not term command yet. We can ignore // We've already got delimiter, but not term command yet. We can ignore
// the delimiter and ack synchronously terminate as if we were in // the delimiter and ack synchronously terminate as if we were in
// active state. // active state.
...@@ -338,8 +345,7 @@ void zmq::pipe_t::terminate () ...@@ -338,8 +345,7 @@ void zmq::pipe_t::terminate ()
else else
zmq_assert (false); zmq_assert (false);
// Stop inbound and outbound flow of messages. // Stop outbound flow of messages.
in_active = false;
out_active = false; out_active = false;
// Rollback any unfinished outbound messages. // Rollback any unfinished outbound messages.
......
...@@ -94,8 +94,9 @@ namespace zmq ...@@ -94,8 +94,9 @@ namespace zmq
// Ask pipe to terminate. The termination will happen asynchronously // Ask pipe to terminate. The termination will happen asynchronously
// and user will be notified about actual deallocation by 'terminated' // and user will be notified about actual deallocation by 'terminated'
// event. // event. If delay is true, the pending messages will be processed
void terminate (); // before actual shutdown.
void terminate (bool delay_);
private: private:
......
...@@ -228,13 +228,6 @@ void zmq::session_t::process_term (int linger_) ...@@ -228,13 +228,6 @@ void zmq::session_t::process_term (int linger_)
return; return;
} }
// If linger is set to zero, we can ask pipe to terminate without
// waiting for pending messages to be read.
if (linger_ == 0) {
proceed_with_term ();
return;
}
pending = true; pending = true;
// If there's finite linger value, delay the termination. // If there's finite linger value, delay the termination.
...@@ -246,6 +239,11 @@ void zmq::session_t::process_term (int linger_) ...@@ -246,6 +239,11 @@ void zmq::session_t::process_term (int linger_)
has_linger_timer = true; has_linger_timer = true;
} }
// Start pipe termination process. Delay the termination till all messages
// are processed in case the linger time is non-zero.
pipe->terminate (linger_ != 0);
// TODO: Should this go into pipe_t::terminate ?
// In case there's no engine and there's only delimiter in the // In case there's no engine and there's only delimiter in the
// pipe it wouldn't be ever read. Thus we check for it explicitly. // pipe it wouldn't be ever read. Thus we check for it explicitly.
pipe->check_read (); pipe->check_read ();
...@@ -256,13 +254,6 @@ void zmq::session_t::proceed_with_term () ...@@ -256,13 +254,6 @@ void zmq::session_t::proceed_with_term ()
// The pending phase have just ended. // The pending phase have just ended.
pending = false; pending = false;
// If there's pipe attached to the session, we have to wait till it
// terminates.
if (pipe) {
register_term_acks (1);
pipe->terminate ();
}
// Continue with standard termination. // Continue with standard termination.
own_t::process_term (0); own_t::process_term (0);
} }
...@@ -276,7 +267,7 @@ void zmq::session_t::timer_event (int id_) ...@@ -276,7 +267,7 @@ void zmq::session_t::timer_event (int id_)
// Ask pipe to terminate even though there may be pending messages in it. // Ask pipe to terminate even though there may be pending messages in it.
zmq_assert (pipe); zmq_assert (pipe);
proceed_with_term (); pipe->terminate (false);
} }
bool zmq::session_t::has_engine () bool zmq::session_t::has_engine ()
......
...@@ -233,7 +233,7 @@ void zmq::socket_base_t::attach_pipe (pipe_t *pipe_, ...@@ -233,7 +233,7 @@ void zmq::socket_base_t::attach_pipe (pipe_t *pipe_,
// straight away. // straight away.
if (is_terminating ()) { if (is_terminating ()) {
register_term_acks (1); register_term_acks (1);
pipe_->terminate (); pipe_->terminate (false);
} }
} }
...@@ -740,7 +740,7 @@ void zmq::socket_base_t::process_term (int linger_) ...@@ -740,7 +740,7 @@ void zmq::socket_base_t::process_term (int linger_)
// Ask all attached pipes to terminate. // Ask all attached pipes to terminate.
for (pipes_t::size_type i = 0; i != pipes.size (); ++i) for (pipes_t::size_type i = 0; i != pipes.size (); ++i)
pipes [i]->terminate (); pipes [i]->terminate (false);
register_term_acks (pipes.size ()); register_term_acks (pipes.size ());
// Continue the termination process immediately. // Continue the termination process immediately.
......
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment