From 7d1937626a6d7f2a2a5961ada1db623f19c60ac6 Mon Sep 17 00:00:00 2001 From: David Rose Date: Tue, 24 Mar 2009 22:24:34 +0000 Subject: [PATCH] more aggressively close threads --- panda/src/net/connection.cxx | 1 + panda/src/net/connectionManager.cxx | 29 +++++++++++++++++++++ panda/src/net/connectionManager.h | 1 + panda/src/net/connectionReader.cxx | 40 +++++++++++++++++++++++++++++ panda/src/net/connectionReader.h | 1 + panda/src/net/connectionWriter.cxx | 30 +++++++++++++++------- panda/src/net/connectionWriter.h | 1 + 7 files changed, 94 insertions(+), 9 deletions(-) diff --git a/panda/src/net/connection.cxx b/panda/src/net/connection.cxx index 1f785d876d..027a45c071 100644 --- a/panda/src/net/connection.cxx +++ b/panda/src/net/connection.cxx @@ -542,6 +542,7 @@ check_send_error(bool okflag) { // Assume any error means the connection has been reset; tell // our manager about it and ignore it. if (_manager != (ConnectionManager *)NULL) { + _manager->flush_read_connection(this); _manager->connection_reset(this, okflag); } return false; diff --git a/panda/src/net/connectionManager.cxx b/panda/src/net/connectionManager.cxx index 6a8257d245..e2d3ec9e4c 100644 --- a/panda/src/net/connectionManager.cxx +++ b/panda/src/net/connectionManager.cxx @@ -301,6 +301,35 @@ new_connection(const PT(Connection) &connection) { _connections.insert(connection); } +//////////////////////////////////////////////////////////////////// +// Function: ConnectionManager::flush_read_connection +// Access: Protected, Virtual +// Description: An internal function called by ConnectionWriter only +// when a write failure has occurred. This method +// ensures that all of the read data has been flushed +// from the pipe before the connection is fully removed. +//////////////////////////////////////////////////////////////////// +void ConnectionManager:: +flush_read_connection(Connection *connection) { + { + LightMutexHolder holder(_set_mutex); + Connections::iterator ci = _connections.find(connection); + if (ci == _connections.end()) { + // Already closed, or not part of this ConnectionManager. + return; + } + _connections.erase(ci); + + Readers::iterator ri; + for (ri = _readers.begin(); ri != _readers.end(); ++ri) { + (*ri)->flush_read_connection(connection); + } + } + + Socket_IP *socket = connection->get_socket(); + socket->Close(); +} + //////////////////////////////////////////////////////////////////// // Function: ConnectionManager::connection_reset // Access: Protected, Virtual diff --git a/panda/src/net/connectionManager.h b/panda/src/net/connectionManager.h index 02fa040d8a..f6e939aa16 100644 --- a/panda/src/net/connectionManager.h +++ b/panda/src/net/connectionManager.h @@ -62,6 +62,7 @@ PUBLISHED: protected: void new_connection(const PT(Connection) &connection); + virtual void flush_read_connection(Connection *connection); virtual void connection_reset(const PT(Connection) &connection, bool okflag); diff --git a/panda/src/net/connectionReader.cxx b/panda/src/net/connectionReader.cxx index 143de566f2..589e311028 100644 --- a/panda/src/net/connectionReader.cxx +++ b/panda/src/net/connectionReader.cxx @@ -386,6 +386,46 @@ get_tcp_header_size() const { return _tcp_header_size; } +//////////////////////////////////////////////////////////////////// +// Function: ConnectionReader::flush_read_connection +// Access: Protected, Virtual +// Description: Attempts to read all the possible data from the +// indicated connection, which has just delivered a +// write error (and has therefore already been closed). +// If the connection is not monitered by this reader, +// does nothing. +//////////////////////////////////////////////////////////////////// +void ConnectionReader:: +flush_read_connection(Connection *connection) { + // Ensure it doesn't get deleted. + SocketInfo sinfo(connection); + + if (!remove_connection(connection)) { + // Not already in the reader. + return; + } + + // The connection was previously in the reader, but has now been + // removed. Now we can flush it completely. We check if there is + // any read data available on just this one socket; we can do this + // right here in this thread, since we've already removed this + // connection from the reader. + + Socket_fdset fdset; + fdset.clear(); + fdset.setForSocket(*(sinfo.get_socket())); + int num_results = fdset.WaitForRead(true, 0); + cerr << "num_results = " << num_results << "\n"; + while (num_results != 0) { + sinfo._busy = true; + process_incoming_data(&sinfo); + fdset.setForSocket(*(sinfo.get_socket())); + num_results = fdset.WaitForRead(true, 0); + cerr << "b num_results = " << num_results << "\n"; + } + cerr << "done\n"; +} + //////////////////////////////////////////////////////////////////// // Function: ConnectionReader::shutdown // Access: Protected diff --git a/panda/src/net/connectionReader.h b/panda/src/net/connectionReader.h index b22978917a..1f2770ee27 100644 --- a/panda/src/net/connectionReader.h +++ b/panda/src/net/connectionReader.h @@ -87,6 +87,7 @@ PUBLISHED: int get_tcp_header_size() const; protected: + virtual void flush_read_connection(Connection *connection); virtual void receive_datagram(const NetDatagram &datagram)=0; class SocketInfo { diff --git a/panda/src/net/connectionWriter.cxx b/panda/src/net/connectionWriter.cxx index 97a2093be7..b429110d82 100644 --- a/panda/src/net/connectionWriter.cxx +++ b/panda/src/net/connectionWriter.cxx @@ -96,15 +96,7 @@ ConnectionWriter:: _manager->remove_writer(this); } - // First, shutdown the queue. This will tell our threads they're - // done. - _queue.shutdown(); - - // Now wait for all threads to terminate. - Threads::iterator ti; - for (ti = _threads.begin(); ti != _threads.end(); ++ti) { - (*ti)->join(); - } + shutdown(); } //////////////////////////////////////////////////////////////////// @@ -341,6 +333,7 @@ get_tcp_header_size() const { void ConnectionWriter:: clear_manager() { _manager = (ConnectionManager *)NULL; + shutdown(); } //////////////////////////////////////////////////////////////////// @@ -362,3 +355,22 @@ thread_run(int thread_index) { } } } + +//////////////////////////////////////////////////////////////////// +// Function: ConnectionWriter::shutdown +// Access: Private +// Description: Stops all the threads and cleans them up. +//////////////////////////////////////////////////////////////////// +void ConnectionWriter:: +shutdown() { + // First, shutdown the queue. This will tell our threads they're + // done. + _queue.shutdown(); + + // Now wait for all threads to terminate. + Threads::iterator ti; + for (ti = _threads.begin(); ti != _threads.end(); ++ti) { + (*ti)->join(); + } + _threads.clear(); +} diff --git a/panda/src/net/connectionWriter.h b/panda/src/net/connectionWriter.h index 9ab4c875e1..30a1ac9322 100644 --- a/panda/src/net/connectionWriter.h +++ b/panda/src/net/connectionWriter.h @@ -71,6 +71,7 @@ protected: private: void thread_run(int thread_index); bool send_datagram(const NetDatagram &datagram); + void shutdown(); protected: ConnectionManager *_manager;