xsi_message_queue & xsi_semaphore: Use condition variables to wait.

This removes a lot of custom logic for managing waiting threads,
which was not even correct in all cases (and the code actually
acknowledged this with a big TODO about it, which weinhold
added all the way back in 2008!)
This commit is contained in:
Augustin Cavalier
2023-04-26 17:16:07 -04:00
parent 9747721a43
commit 6acd708e97
2 changed files with 76 additions and 183 deletions
+36 -90
View File
@@ -1,5 +1,5 @@
/* /*
* Copyright 2008-2011, Haiku, Inc. All rights reserved. * Copyright 2008-2023, Haiku, Inc. All rights reserved.
* Distributed under the terms of the MIT License. * Distributed under the terms of the MIT License.
* *
* Authors: * Authors:
@@ -36,23 +36,6 @@
namespace { namespace {
// Queue for holding blocked threads
struct queued_thread : DoublyLinkedListLinkImpl<queued_thread> {
queued_thread(Thread *_thread, int32 _message_length)
:
thread(_thread),
message_length(_message_length),
queued(false)
{
}
Thread *thread;
int32 message_length;
bool queued;
};
typedef DoublyLinkedList<queued_thread> ThreadQueue;
struct queued_message : DoublyLinkedListLinkImpl<queued_message> { struct queued_message : DoublyLinkedListLinkImpl<queued_message> {
queued_message(const void *_message, ssize_t _length) queued_message(const void *_message, ssize_t _length)
@@ -106,11 +89,12 @@ class XsiMessageQueue {
public: public:
XsiMessageQueue(int flags) XsiMessageQueue(int flags)
: :
fBytesInQueue(0), fBytesInQueue(0)
fThreadsWaitingToReceive(0),
fThreadsWaitingToSend(0)
{ {
mutex_init(&fLock, "XsiMessageQueue private mutex"); mutex_init(&fLock, "XsiMessageQueue private mutex");
fWaitingToReceive.Init(this, "XsiMessageQueue");
fWaitingToSend.Init(this, "XsiMessageQueue");
SetIpcKey((key_t)-1); SetIpcKey((key_t)-1);
SetPermissions(flags); SetPermissions(flags);
// Initialize all fields to zero // Initialize all fields to zero
@@ -122,18 +106,11 @@ public:
// Implemented after sXsiMessageCount is declared // Implemented after sXsiMessageCount is declared
~XsiMessageQueue(); ~XsiMessageQueue();
status_t BlockAndUnlock(Thread *thread, MutexLocker *queueLocker) status_t BlockAndUnlock(ConditionVariableEntry *queueEntry, MutexLocker *queueLocker)
{ {
thread_prepare_to_block(thread, B_CAN_INTERRUPT,
THREAD_BLOCK_TYPE_OTHER, (void*)"xsi message queue");
// Unlock the queue before blocking // Unlock the queue before blocking
queueLocker->Unlock(); queueLocker->Unlock();
return queueEntry->Wait(B_CAN_INTERRUPT);
// TODO: We've got a serious race condition: If BlockAndUnlock() returned due to
// interruption, we will still be queued. A WakeUpThread() at this point will
// call thread_unblock() and might thus screw with our trying to re-lock the
// mutex.
return thread_block();
} }
void DoIpcSet(struct msqid_ds *result) void DoIpcSet(struct msqid_ds *result)
@@ -146,29 +123,18 @@ public:
fMessageQueue.msg_ctime = (time_t)real_time_clock(); fMessageQueue.msg_ctime = (time_t)real_time_clock();
} }
void Deque(queued_thread *queueEntry, bool waitForMessage) void Dequeue(ConditionVariableEntry *queueEntry, bool waitForMessage)
{ {
if (queueEntry->queued) { queueEntry->Wait(B_RELATIVE_TIMEOUT, 0);
if (waitForMessage) {
fWaitingToReceive.Remove(queueEntry);
fThreadsWaitingToReceive--;
} else {
fWaitingToSend.Remove(queueEntry);
fThreadsWaitingToSend--;
}
}
} }
void Enqueue(queued_thread *queueEntry, bool waitForMessage) void Enqueue(ConditionVariableEntry *queueEntry, bool waitForMessage)
{ {
if (waitForMessage) { if (waitForMessage) {
fWaitingToReceive.Add(queueEntry); fWaitingToReceive.Add(queueEntry);
fThreadsWaitingToReceive++;
} else { } else {
fWaitingToSend.Add(queueEntry); fWaitingToSend.Add(queueEntry);
fThreadsWaitingToSend++;
} }
queueEntry->queued = true;
} }
struct msqid_ds &GetMessageQueue() struct msqid_ds &GetMessageQueue()
@@ -252,18 +218,10 @@ public:
// Wake up all waiting thread for a message // Wake up all waiting thread for a message
// TODO: this can cause starvation for any // TODO: this can cause starvation for any
// very-unlucky-and-slow thread // very-unlucky-and-slow thread
while (queued_thread *entry = fWaitingToReceive.RemoveHead()) { fWaitingToReceive.NotifyAll();
entry->queued = false;
fThreadsWaitingToReceive--;
thread_unblock(entry->thread, 0);
}
} else { } else {
// Wake up only one thread waiting to send // Wake up only one thread waiting to send
if (queued_thread *entry = fWaitingToSend.RemoveHead()) { fWaitingToSend.NotifyOne();
entry->queued = false;
fThreadsWaitingToSend--;
thread_unblock(entry->thread, 0);
}
} }
} }
@@ -279,11 +237,9 @@ private:
MessageQueue fMessage; MessageQueue fMessage;
struct msqid_ds fMessageQueue; struct msqid_ds fMessageQueue;
uint32 fSequenceNumber; uint32 fSequenceNumber;
uint32 fThreadsWaitingToReceive;
uint32 fThreadsWaitingToSend;
ThreadQueue fWaitingToReceive; ConditionVariable fWaitingToReceive;
ThreadQueue fWaitingToSend; ConditionVariable fWaitingToSend;
XsiMessageQueue* fLink; XsiMessageQueue* fLink;
}; };
@@ -402,16 +358,8 @@ XsiMessageQueue::~XsiMessageQueue()
mutex_destroy(&fLock); mutex_destroy(&fLock);
// Wake up any threads still waiting // Wake up any threads still waiting
if (fThreadsWaitingToSend || fThreadsWaitingToReceive) { fWaitingToReceive.NotifyAll(EIDRM);
while (queued_thread *entry = fWaitingToReceive.RemoveHead()) { fWaitingToSend.NotifyAll(EIDRM);
entry->queued = false;
thread_unblock(entry->thread, EIDRM);
}
while (queued_thread *entry = fWaitingToSend.RemoveHead()) {
entry->queued = false;
thread_unblock(entry->thread, EIDRM);
}
}
// Free up any remaining messages // Free up any remaining messages
if (fMessageQueue.msg_qnum) { if (fMessageQueue.msg_qnum) {
@@ -447,8 +395,8 @@ XsiMessageQueue::Insert(queued_message *message)
fMessageQueue.msg_lspid = getpid(); fMessageQueue.msg_lspid = getpid();
fMessageQueue.msg_stime = real_time_clock(); fMessageQueue.msg_stime = real_time_clock();
fBytesInQueue += message->length; fBytesInQueue += message->length;
if (fThreadsWaitingToReceive)
WakeUpThread(true /* WaitForMessage */); WakeUpThread(true /* WaitForMessage */);
return false; return false;
} }
@@ -492,8 +440,8 @@ XsiMessageQueue::Remove(long typeRequested)
fMessageQueue.msg_rtime = real_time_clock(); fMessageQueue.msg_rtime = real_time_clock();
fBytesInQueue -= message->length; fBytesInQueue -= message->length;
atomic_add(&sXsiMessageCount, -1); atomic_add(&sXsiMessageCount, -1);
if (fThreadsWaitingToSend)
WakeUpThread(false /* WaitForMessage */); WakeUpThread(false /* WaitForMessage */);
return message; return message;
} }
@@ -763,30 +711,29 @@ _user_xsi_msgrcv(int messageQueueID, void *messagePointer,
if (message == NULL && !(messageFlags & IPC_NOWAIT)) { if (message == NULL && !(messageFlags & IPC_NOWAIT)) {
// We are going to sleep // We are going to sleep
Thread *thread = thread_get_current_thread(); ConditionVariableEntry queueEntry;
queued_thread queueEntry(thread, messageSize);
messageQueue->Enqueue(&queueEntry, /* waitForMessage */ true); messageQueue->Enqueue(&queueEntry, /* waitForMessage */ true);
uint32 sequenceNumber = messageQueue->SequenceNumber(); uint32 sequenceNumber = messageQueue->SequenceNumber();
TRACE(("xsi_msgrcv: thread %d going to sleep\n", (int)thread->id)); TRACE(("xsi_msgrcv: thread %d going to sleep\n", (int)thread_get_current_thread_id()));
status_t result status_t result
= messageQueue->BlockAndUnlock(thread, &messageQueueLocker); = messageQueue->BlockAndUnlock(&queueEntry, &messageQueueLocker);
TRACE(("xsi_msgrcv: thread %d back to life\n", (int)thread->id)); TRACE(("xsi_msgrcv: thread %d back to life\n", (int)thread_get_current_thread_id()));
messageQueueHashLocker.Lock(); messageQueueHashLocker.Lock();
messageQueue = sMessageQueueHashTable.Lookup(messageQueueID); messageQueue = sMessageQueueHashTable.Lookup(messageQueueID);
if (result == EIDRM || messageQueue == NULL || (messageQueue != NULL if (result == EIDRM || messageQueue == NULL || (messageQueue != NULL
&& sequenceNumber != messageQueue->SequenceNumber())) { && sequenceNumber != messageQueue->SequenceNumber())) {
TRACE_ERROR(("xsi_msgrcv: message queue id %d (sequence = " TRACE(("xsi_msgrcv: message queue id %d (sequence = "
"%" B_PRIu32 ") got destroyed\n", messageQueueID, "%" B_PRIu32 ") got destroyed\n", messageQueueID,
sequenceNumber)); sequenceNumber));
return EIDRM; return EIDRM;
} else if (result == B_INTERRUPTED) { } else if (result == B_INTERRUPTED) {
TRACE_ERROR(("xsi_msgrcv: thread %d got interrupted while " TRACE(("xsi_msgrcv: thread %d got interrupted while "
"waiting on message queue %d\n",(int)thread->id, "waiting on message queue %d\n", (int)thread_get_current_thread_id(),
messageQueueID)); messageQueueID));
messageQueue->Deque(&queueEntry, /* waitForMessage */ true); messageQueue->Dequeue(&queueEntry, /* waitForMessage */ true);
return EINTR; return EINTR;
} else { } else {
messageQueueLocker.Lock(); messageQueueLocker.Lock();
@@ -871,31 +818,30 @@ _user_xsi_msgsnd(int messageQueueID, const void *messagePointer,
if (goToSleep && !(messageFlags & IPC_NOWAIT)) { if (goToSleep && !(messageFlags & IPC_NOWAIT)) {
// We are going to sleep // We are going to sleep
Thread *thread = thread_get_current_thread(); ConditionVariableEntry queueEntry;
queued_thread queueEntry(thread, messageSize);
messageQueue->Enqueue(&queueEntry, /* waitForMessage */ false); messageQueue->Enqueue(&queueEntry, /* waitForMessage */ false);
uint32 sequenceNumber = messageQueue->SequenceNumber(); uint32 sequenceNumber = messageQueue->SequenceNumber();
TRACE(("xsi_msgsnd: thread %d going to sleep\n", (int)thread->id)); TRACE(("xsi_msgsnd: thread %d going to sleep\n", (int)thread_get_current_thread_id()));
result = messageQueue->BlockAndUnlock(thread, &messageQueueLocker); result = messageQueue->BlockAndUnlock(&queueEntry, &messageQueueLocker);
TRACE(("xsi_msgsnd: thread %d back to life\n", (int)thread->id)); TRACE(("xsi_msgsnd: thread %d back to life\n", (int)thread_get_current_thread_id()));
messageQueueHashLocker.Lock(); messageQueueHashLocker.Lock();
messageQueue = sMessageQueueHashTable.Lookup(messageQueueID); messageQueue = sMessageQueueHashTable.Lookup(messageQueueID);
if (result == EIDRM || messageQueue == NULL || (messageQueue != NULL if (result == EIDRM || messageQueue == NULL || (messageQueue != NULL
&& sequenceNumber != messageQueue->SequenceNumber())) { && sequenceNumber != messageQueue->SequenceNumber())) {
TRACE_ERROR(("xsi_msgsnd: message queue id %d (sequence = " TRACE(("xsi_msgsnd: message queue id %d (sequence = "
"%" B_PRIu32 ") got destroyed\n", messageQueueID, "%" B_PRIu32 ") got destroyed\n", messageQueueID,
sequenceNumber)); sequenceNumber));
delete message; delete message;
notSent = false; notSent = false;
result = EIDRM; result = EIDRM;
} else if (result == B_INTERRUPTED) { } else if (result == B_INTERRUPTED) {
TRACE_ERROR(("xsi_msgsnd: thread %d got interrupted while " TRACE(("xsi_msgsnd: thread %d got interrupted while "
"waiting on message queue %d\n",(int)thread->id, "waiting on message queue %d\n", (int)thread_get_current_thread_id(),
messageQueueID)); messageQueueID));
messageQueue->Deque(&queueEntry, /* waitForMessage */ false); messageQueue->Dequeue(&queueEntry, /* waitForMessage */ false);
delete message; delete message;
notSent = false; notSent = false;
result = EINTR; result = EINTR;
+40 -93
View File
@@ -1,5 +1,5 @@
/* /*
* Copyright 2008-2011, Haiku, Inc. All rights reserved. * Copyright 2008-2023, Haiku, Inc. All rights reserved.
* Distributed under the terms of the MIT License. * Distributed under the terms of the MIT License.
* *
* Authors: * Authors:
@@ -37,23 +37,6 @@
namespace { namespace {
// Queue for holding blocked threads
struct queued_thread : DoublyLinkedListLinkImpl<queued_thread> {
queued_thread(Thread *thread, int32 count)
:
thread(thread),
count(count),
queued(false)
{
}
Thread *thread;
int32 count;
bool queued;
};
typedef DoublyLinkedList<queued_thread> ThreadQueue;
class XsiSemaphoreSet; class XsiSemaphoreSet;
struct sem_undo : DoublyLinkedListLinkImpl<sem_undo> { struct sem_undo : DoublyLinkedListLinkImpl<sem_undo> {
@@ -101,25 +84,21 @@ namespace {
class XsiSemaphore { class XsiSemaphore {
public: public:
XsiSemaphore() XsiSemaphore()
: fLastPidOperation(0), :
fThreadsWaitingToIncrease(0), fLastPidOperation(0),
fThreadsWaitingToBeZero(0),
fValue(0) fValue(0)
{ {
fWaitingToIncrease.Init(this, "XsiSemaphore");
fWaitingToBeZero.Init(this, "XsiSemaphore");
} }
~XsiSemaphore() ~XsiSemaphore()
{ {
// For some reason the semaphore is getting destroyed. // For some reason the semaphore is getting destroyed.
// Wake up any remaing awaiting threads // Wake up any remaing awaiting threads
while (queued_thread *entry = fWaitingToIncreaseQueue.RemoveHead()) { fWaitingToIncrease.NotifyAll(EIDRM);
entry->queued = false; fWaitingToBeZero.NotifyAll(EIDRM);
thread_unblock(entry->thread, EIDRM);
}
while (queued_thread *entry = fWaitingToBeZeroQueue.RemoveHead()) {
entry->queued = false;
thread_unblock(entry->thread, EIDRM);
}
// No need to remove any sem_undo request still // No need to remove any sem_undo request still
// hanging. When the process exit and doesn't found // hanging. When the process exit and doesn't found
// the semaphore set, it'll just ignore the sem_undo // the semaphore set, it'll just ignore the sem_undo
@@ -138,51 +117,33 @@ public:
return true; return true;
} else { } else {
fValue += value; fValue += value;
if (fValue == 0 && fThreadsWaitingToBeZero > 0) if (fValue == 0)
WakeUpThread(true); WakeUpThreads(true);
else if (fValue > 0 && fThreadsWaitingToIncrease > 0) else if (fValue > 0)
WakeUpThread(false); WakeUpThreads(false);
return false; return false;
} }
} }
status_t BlockAndUnlock(Thread *thread, MutexLocker *setLocker) status_t BlockAndUnlock(ConditionVariableEntry *queueEntry, MutexLocker *setLocker)
{ {
thread_prepare_to_block(thread, B_CAN_INTERRUPT,
THREAD_BLOCK_TYPE_OTHER, (void*)"xsi semaphore");
// Unlock the set before blocking // Unlock the set before blocking
setLocker->Unlock(); setLocker->Unlock();
return queueEntry->Wait(B_CAN_INTERRUPT);
// TODO: We've got a serious race condition: If BlockAndUnlock() returned due to
// interruption, we will still be queued. A WakeUpThread() at this point will
// call thread_unblock() and might thus screw with our trying to re-lock the
// mutex.
return thread_block();
} }
void Deque(queued_thread *queueEntry, bool waitForZero) void Dequeue(ConditionVariableEntry *queueEntry, bool waitForZero)
{ {
if (queueEntry->queued) { queueEntry->Wait(B_RELATIVE_TIMEOUT, 0);
if (waitForZero) {
fWaitingToBeZeroQueue.Remove(queueEntry);
fThreadsWaitingToBeZero--;
} else {
fWaitingToIncreaseQueue.Remove(queueEntry);
fThreadsWaitingToIncrease--;
}
}
} }
void Enqueue(queued_thread *queueEntry, bool waitForZero) void Enqueue(ConditionVariableEntry *queueEntry, bool waitForZero)
{ {
if (waitForZero) { if (waitForZero) {
fWaitingToBeZeroQueue.Add(queueEntry); fWaitingToBeZero.Add(queueEntry);
fThreadsWaitingToBeZero++;
} else { } else {
fWaitingToIncreaseQueue.Add(queueEntry); fWaitingToIncrease.Add(queueEntry);
fThreadsWaitingToIncrease++;
} }
queueEntry->queued = true;
} }
pid_t LastPid() const pid_t LastPid() const
@@ -193,10 +154,10 @@ public:
void Revert(short value) void Revert(short value)
{ {
fValue -= value; fValue -= value;
if (fValue == 0 && fThreadsWaitingToBeZero > 0) if (fValue == 0)
WakeUpThread(true); WakeUpThreads(true);
else if (fValue > 0 && fThreadsWaitingToIncrease > 0) else if (fValue > 0)
WakeUpThread(false); WakeUpThreads(false);
} }
void SetPid(pid_t pid) void SetPid(pid_t pid)
@@ -209,14 +170,14 @@ public:
fValue = value; fValue = value;
} }
ushort ThreadsWaitingToIncrease() const ushort ThreadsWaitingToIncrease()
{ {
return fThreadsWaitingToIncrease; return fWaitingToIncrease.EntriesCount();
} }
ushort ThreadsWaitingToBeZero() const ushort ThreadsWaitingToBeZero()
{ {
return fThreadsWaitingToBeZero; return fWaitingToBeZero.EntriesCount();
} }
ushort Value() const ushort Value() const
@@ -224,33 +185,21 @@ public:
return fValue; return fValue;
} }
void WakeUpThread(bool waitingForZero) void WakeUpThreads(bool waitingForZero)
{ {
if (waitingForZero) { if (waitingForZero) {
// Wake up all threads waiting on zero fWaitingToBeZero.NotifyAll();
while (queued_thread *entry = fWaitingToBeZeroQueue.RemoveHead()) {
entry->queued = false;
fThreadsWaitingToBeZero--;
thread_unblock(entry->thread, 0);
}
} else { } else {
// Wake up all threads even though they might go back to sleep fWaitingToIncrease.NotifyAll();
while (queued_thread *entry = fWaitingToIncreaseQueue.RemoveHead()) {
entry->queued = false;
fThreadsWaitingToIncrease--;
thread_unblock(entry->thread, 0);
}
} }
} }
private: private:
pid_t fLastPidOperation; // sempid pid_t fLastPidOperation; // sempid
ushort fThreadsWaitingToIncrease; // semncnt
ushort fThreadsWaitingToBeZero; // semzcnt
ushort fValue; // semval ushort fValue; // semval
ThreadQueue fWaitingToIncreaseQueue; ConditionVariable fWaitingToIncrease;
ThreadQueue fWaitingToBeZeroQueue; ConditionVariable fWaitingToBeZero;
}; };
#define MAX_XSI_SEMS_PER_TEAM 128 #define MAX_XSI_SEMS_PER_TEAM 128
@@ -1191,33 +1140,31 @@ _user_xsi_semop(int semaphoreID, struct sembuf *ops, size_t numOps)
if (operations[i].sem_op != 0) if (operations[i].sem_op != 0)
waitOnZero = false; waitOnZero = false;
Thread *thread = thread_get_current_thread(); ConditionVariableEntry queueEntry;
queued_thread queueEntry(thread, (int32)operations[i].sem_op);
semaphore->Enqueue(&queueEntry, waitOnZero); semaphore->Enqueue(&queueEntry, waitOnZero);
uint32 sequenceNumber = semaphoreSet->SequenceNumber(); uint32 sequenceNumber = semaphoreSet->SequenceNumber();
TRACE(("xsi_semop: thread %d going to sleep\n", (int)thread->id)); TRACE(("xsi_semop: thread %d going to sleep\n", (int)thread->id));
result = semaphore->BlockAndUnlock(thread, &setLocker); result = semaphore->BlockAndUnlock(&queueEntry, &setLocker);
TRACE(("xsi_semop: thread %d back to life\n", (int)thread->id)); TRACE(("xsi_semop: thread %d back to life\n", (int)thread->id));
// We are back to life. Find out why! // We are back to life. Find out why!
// Make sure the set hasn't been deleted or worst yet // Make sure the set hasn't been deleted or worst yet replaced.
// replaced.
setHashLocker.Lock(); setHashLocker.Lock();
semaphoreSet = sSemaphoreHashTable.Lookup(semaphoreID); semaphoreSet = sSemaphoreHashTable.Lookup(semaphoreID);
if (result == EIDRM || semaphoreSet == NULL || (semaphoreSet != NULL if (result == EIDRM || semaphoreSet == NULL || (semaphoreSet != NULL
&& sequenceNumber != semaphoreSet->SequenceNumber())) { && sequenceNumber != semaphoreSet->SequenceNumber())) {
TRACE_ERROR(("xsi_semop: semaphore set id %d (sequence = " TRACE(("xsi_semop: semaphore set id %d (sequence = "
"%" B_PRIu32 ") got destroyed\n", semaphoreID, "%" B_PRIu32 ") got destroyed\n", semaphoreID,
sequenceNumber)); sequenceNumber));
notDone = false; notDone = false;
result = EIDRM; result = EIDRM;
} else if (result == B_INTERRUPTED) { } else if (result == B_INTERRUPTED) {
TRACE_ERROR(("xsi_semop: thread %d got interrupted while " TRACE(("xsi_semop: thread %d got interrupted while "
"waiting on semaphore set id %d\n",(int)thread->id, "waiting on semaphore set id %d\n", (int)thread_get_current_thread_id(),
semaphoreID)); semaphoreID));
semaphore->Deque(&queueEntry, waitOnZero); semaphore->Dequeue(&queueEntry, waitOnZero);
result = EINTR; result = EINTR;
notDone = false; notDone = false;
} else { } else {