summaryrefslogtreecommitdiff
path: root/cpp/src/qpid/client/amqp0_10/SessionImpl.cpp
diff options
context:
space:
mode:
authorAlan Conway <aconway@apache.org>2010-05-13 18:54:39 +0000
committerAlan Conway <aconway@apache.org>2010-05-13 18:54:39 +0000
commit7e13897c9238d0e5a6a64df64eeddceb14c36002 (patch)
tree3e91072bbb28ba7f2403477b2c7f2f945cb3699c /cpp/src/qpid/client/amqp0_10/SessionImpl.cpp
parent211eb170c15e6951f28e900e0ea8284e4e1456eb (diff)
downloadqpid-python-7e13897c9238d0e5a6a64df64eeddceb14c36002.tar.gz
Fix deadlocks & thread safety in new API classes.
git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk/qpid@943973 13f79535-47bb-0310-9956-ffa450edef68
Diffstat (limited to 'cpp/src/qpid/client/amqp0_10/SessionImpl.cpp')
-rw-r--r--cpp/src/qpid/client/amqp0_10/SessionImpl.cpp47
1 files changed, 33 insertions, 14 deletions
diff --git a/cpp/src/qpid/client/amqp0_10/SessionImpl.cpp b/cpp/src/qpid/client/amqp0_10/SessionImpl.cpp
index a55a2737cb..a6067097bb 100644
--- a/cpp/src/qpid/client/amqp0_10/SessionImpl.cpp
+++ b/cpp/src/qpid/client/amqp0_10/SessionImpl.cpp
@@ -53,6 +53,9 @@ namespace qpid {
namespace client {
namespace amqp0_10 {
+typedef qpid::sys::Mutex::ScopedLock ScopedLock;
+typedef qpid::sys::Mutex::ScopedUnlock ScopedUnlock;
+
SessionImpl::SessionImpl(ConnectionImpl& c, bool t) : connection(&c), transactional(t) {}
void SessionImpl::checkError()
@@ -112,23 +115,29 @@ void SessionImpl::release(qpid::messaging::Message& m)
void SessionImpl::close()
{
if (hasError()) {
+ ScopedLock l(lock);
senders.clear();
receivers.clear();
} else {
- //close all the senders and receivers (get copy of names and then
- //make the calls to avoid modifying maps while iterating over
- //them):
- std::vector<std::string> s;
- std::vector<std::string> r;
- {
- qpid::sys::Mutex::ScopedLock l(lock);
- for (Senders::const_iterator i = senders.begin(); i != senders.end(); ++i) s.push_back(i->first);
- for (Receivers::const_iterator i = receivers.begin(); i != receivers.end(); ++i) r.push_back(i->first);
+ while (true) {
+ Sender s;
+ {
+ ScopedLock l(lock);
+ if (senders.empty()) break;
+ s = senders.begin()->second;
+ }
+ s.close(); // outside the lock, will call senderCancelled
+ }
+ while (true) {
+ Receiver r;
+ {
+ ScopedLock l(lock);
+ if (receivers.empty()) break;
+ r = receivers.begin()->second;
+ }
+ r.close(); // outside the lock, will call receiverCancelled
}
- for (std::vector<std::string>::const_iterator i = s.begin(); i != s.end(); ++i) getSender(*i).close();
- for (std::vector<std::string>::const_iterator i = r.begin(); i != r.end(); ++i) getReceiver(*i).close();
}
-
connection->closed(*this);
if (!hasError()) session.close();
}
@@ -151,7 +160,7 @@ template <class T> void getFreeKey(std::string& key, T& map)
void SessionImpl::setSession(qpid::client::Session s)
{
- qpid::sys::Mutex::ScopedLock l(lock);
+ ScopedLock l(lock);
session = s;
incoming.setSession(session);
if (transactional) session.txSelect();
@@ -181,6 +190,7 @@ Receiver SessionImpl::createReceiver(const qpid::messaging::Address& address)
Receiver SessionImpl::createReceiverImpl(const qpid::messaging::Address& address)
{
+ ScopedLock l(lock);
std::string name = address.getName();
getFreeKey(name, receivers);
Receiver receiver(new ReceiverImpl(*this, name, address));
@@ -205,7 +215,8 @@ Sender SessionImpl::createSender(const qpid::messaging::Address& address)
}
Sender SessionImpl::createSenderImpl(const qpid::messaging::Address& address)
-{
+{
+ ScopedLock l(lock);
std::string name = address.getName();
getFreeKey(name, senders);
Sender sender(new SenderImpl(*this, name, address));
@@ -265,6 +276,7 @@ struct IncomingMessageHandler : IncomingMessages::Handler
bool SessionImpl::getNextReceiver(Receiver* receiver, IncomingMessages::MessageTransfer& transfer)
{
+ ScopedLock l(lock);
Receivers::const_iterator i = receivers.find(transfer.getDestination());
if (i == receivers.end()) {
QPID_LOG(error, "Received message for unknown destination " << transfer.getDestination());
@@ -371,6 +383,7 @@ struct SessionImpl::Receivable : Command
uint32_t SessionImpl::getReceivableImpl(const std::string* destination)
{
+ ScopedLock l(lock);
if (destination) {
return incoming.available(*destination);
} else {
@@ -399,6 +412,7 @@ struct SessionImpl::UnsettledAcks : Command
uint32_t SessionImpl::getUnsettledAcksImpl(const std::string* destination)
{
+ ScopedLock l(lock);
if (destination) {
return incoming.pendingAccept(*destination);
} else {
@@ -414,12 +428,14 @@ void SessionImpl::syncImpl(bool block)
void SessionImpl::commitImpl()
{
+ ScopedLock l(lock);
incoming.accept();
session.txCommit();
}
void SessionImpl::rollbackImpl()
{
+ ScopedLock l(lock);
for (Receivers::iterator i = receivers.begin(); i != receivers.end(); ++i) {
getImplPtr<Receiver, ReceiverImpl>(i->second)->stop();
}
@@ -436,6 +452,7 @@ void SessionImpl::rollbackImpl()
void SessionImpl::acknowledgeImpl()
{
+ ScopedLock l(lock);
if (!transactional) incoming.accept();
}
@@ -455,6 +472,7 @@ void SessionImpl::releaseImpl(qpid::messaging::Message& m)
void SessionImpl::receiverCancelled(const std::string& name)
{
+ ScopedLock l(lock);
receivers.erase(name);
session.sync();
incoming.releasePending(name);
@@ -462,6 +480,7 @@ void SessionImpl::receiverCancelled(const std::string& name)
void SessionImpl::senderCancelled(const std::string& name)
{
+ ScopedLock l(lock);
senders.erase(name);
}