diff --git a/OSC/Source/OSC/Private/Receive/OscDispatcher.cpp b/OSC/Source/OSC/Private/Receive/OscDispatcher.cpp index 7550cb4..855e07d 100644 --- a/OSC/Source/OSC/Private/Receive/OscDispatcher.cpp +++ b/OSC/Source/OSC/Private/Receive/OscDispatcher.cpp @@ -165,8 +165,6 @@ static void SendBundle(TCircularQueueGetData(), data->Num()); if(packet.State() != osc::SUCCESS) { @@ -184,9 +182,9 @@ void UOscDispatcher::Callback(const FArrayReaderPtr& data, const FIPv4Endpoint&) } // Set a single callback in the main thread per frame. - if(wasEmpty && !_pendingMessages.IsEmpty()) + if(!_pendingMessages.IsEmpty() && !_runPendingMessagesTask) { - FSimpleDelegateGraphTask::CreateAndDispatchWhenReady( + _runPendingMessagesTask = FSimpleDelegateGraphTask::CreateAndDispatchWhenReady( FSimpleDelegateGraphTask::FDelegate::CreateUObject(this, &UOscDispatcher::CallbackMainThread), TStatId(), nullptr, @@ -205,6 +203,11 @@ void UOscDispatcher::Callback(const FArrayReaderPtr& data, const FIPv4Endpoint&) void UOscDispatcher::CallbackMainThread() { + // Release before dequeue. + // If it was released after dequeue, when a message arrives after the while + // loop and before the release, it would not be processed. + _runPendingMessagesTask.SafeRelease(); + FScopeLock ScopeLock(&_receiversMutex); std::pair> message; diff --git a/OSC/Source/OSC/Private/Receive/OscDispatcher.h b/OSC/Source/OSC/Private/Receive/OscDispatcher.h index e2ae7b1..441433a 100644 --- a/OSC/Source/OSC/Private/Receive/OscDispatcher.h +++ b/OSC/Source/OSC/Private/Receive/OscDispatcher.h @@ -43,13 +43,14 @@ private: void CallbackMainThread(); void BeginDestroy() override; - + private: TArray _receivers; std::pair _listening; FSocket * _socket; FUdpSocketReceiver * _socketReceiver; TCircularQueue>> _pendingMessages; + FGraphEventRef _runPendingMessagesTask; /// Protects _receivers FCriticalSection _receiversMutex;