aboutsummaryrefslogtreecommitdiff
path: root/binding/unreal
diff options
context:
space:
mode:
author魏曹先生 <1992414357@qq.com>2026-06-22 19:54:48 +0800
committer魏曹先生 <1992414357@qq.com>2026-06-22 19:54:48 +0800
commit80120a325597fcbdc8a987258caceaf9cc9c448f (patch)
treeec4a8ceadc4abd9548622a85ea037e1fe45eb4d1 /binding/unreal
parentc553a53f25ddc72ebbc6df13bc61bfd49980044c (diff)
refactor(unreal): migrate socket I/O to FRunnable worker thread
Diffstat (limited to 'binding/unreal')
-rw-r--r--binding/unreal/DMVOPBridge/Source/DMVOPBridge/DMVOPBridgeClient.cpp269
-rw-r--r--binding/unreal/DMVOPBridge/Source/DMVOPBridge/DMVOPBridgeClient.h34
2 files changed, 237 insertions, 66 deletions
diff --git a/binding/unreal/DMVOPBridge/Source/DMVOPBridge/DMVOPBridgeClient.cpp b/binding/unreal/DMVOPBridge/Source/DMVOPBridge/DMVOPBridgeClient.cpp
index 15a11e0..8a436e0 100644
--- a/binding/unreal/DMVOPBridge/Source/DMVOPBridge/DMVOPBridgeClient.cpp
+++ b/binding/unreal/DMVOPBridge/Source/DMVOPBridge/DMVOPBridgeClient.cpp
@@ -1,75 +1,170 @@
#include "DMVOPBridgeClient.h"
-#include "Common/TcpSocketBuilder.h"
#include "IPAddress.h"
+#include "Interfaces/IPv4/IPv4Address.h"
+
+#ifdef _WIN32
+#include <winsock2.h>
+#include <ws2tcpip.h>
+// Undef Windows macros that conflict with UE types
+#ifdef SetPort
+#undef SetPort
+#endif
+#endif
+
+// ═════════════════════════════════════════════════════════════════════════════
+// FDMVOPWorker
+// ═════════════════════════════════════════════════════════════════════════════
+
+FDMVOPWorker::FDMVOPWorker(TWeakObjectPtr<UDMVOPClient> InOwner, FString InHost,
+ int32 InPort)
+ : bRun(false), Socket(nullptr), Owner(InOwner), Host(MoveTemp(InHost)),
+ Port(InPort), Thread(nullptr) {}
+
+FDMVOPWorker::~FDMVOPWorker() {
+ Stop();
+ if (Thread) {
+ Thread->WaitForCompletion();
+ delete Thread;
+ Thread = nullptr;
+ }
+}
-UDMVOPClient::UDMVOPClient() : Socket(nullptr) {}
+void FDMVOPWorker::Start() {
+ Thread = FRunnableThread::Create(
+ this, *FString::Printf(TEXT("DMVOP %s:%d"), *Host, Port), 128 * 1024,
+ TPri_Normal);
+}
-UDMVOPClient::~UDMVOPClient() { cleanup(); }
+bool FDMVOPWorker::Init() {
+ bRun = true;
+ return true;
+}
-void UDMVOPClient::cleanup() {
- // SAFETY: signal the reader thread to stop, then close the socket so
- // Recv() unblocks immediately. After the thread exits, destroy the socket.
- // The thread checks bRunning on every iteration and after each Recv().
- bRunning.store(false);
- if (Socket) {
- Socket->Close(); // unblocks Recv() in the reader thread
- }
- if (ReaderThread.joinable())
- ReaderThread.join(); // thread exits after Recv() fails + bRunning check
- if (Socket) {
- ISocketSubsystem::Get()->DestroySocket(Socket);
- Socket = nullptr;
+uint32 FDMVOPWorker::Run() {
+ // ── Create socket ──
+ Socket = ISocketSubsystem::Get(PLATFORM_SOCKETSUBSYSTEM)
+ ->CreateSocket(NAME_Stream, TEXT("DMVOP"), false);
+ if (!Socket) {
+ return 0;
}
-}
-void UDMVOPClient::Connect(const FString &Host, int32 Port) {
- if (Socket)
- return;
+ int32 RecvSize = 0, SendSize = 0;
+ Socket->SetReceiveBufferSize(16384, RecvSize);
+ Socket->SetSendBufferSize(16384, SendSize);
- FIPv4Address IP;
- if (!FIPv4Address::Parse(Host, IP)) {
- UE_LOG(LogTemp, Error, TEXT("DMVOP: Invalid IP %s"), *Host);
- return;
+ // ── Resolve address ──
+ FIPv4Address AddrIP;
+ if (!FIPv4Address::Parse(Host, AddrIP)) {
+ Socket->Close();
+ delete Socket;
+ Socket = nullptr;
+ return 0;
}
- TSharedRef<FInternetAddr> Addr =
+ TSharedRef<FInternetAddr> DstAddr =
ISocketSubsystem::Get(PLATFORM_SOCKETSUBSYSTEM)->CreateInternetAddr();
- Addr->SetIp(IP.Value);
- Addr->SetPort(Port);
-
- Socket = FTcpSocketBuilder(TEXT("DMVOPSocket")).AsNonBlocking().Build();
-
- if (!Socket->Connect(*Addr)) {
- UE_LOG(LogTemp, Error, TEXT("DMVOP: Failed to connect"));
+ DstAddr->SetIp(AddrIP.Value);
+ DstAddr->SetPort(Port);
+
+ // ── Connect ──
+ if (!Socket->Connect(*DstAddr)) {
+ TWeakObjectPtr<UDMVOPClient> WeakOwner = Owner;
+ AsyncTask(ENamedThreads::GameThread, [WeakOwner]() {
+ if (auto *Self = WeakOwner.Get())
+ UE_LOG(LogTemp, Error, TEXT("DMVOP: Failed to connect"));
+ });
Socket->Close();
- ISocketSubsystem::Get()->DestroySocket(Socket);
+ delete Socket;
Socket = nullptr;
- return;
+ return 0;
+ }
+
+ TWeakObjectPtr<UDMVOPClient> WeakOwner = Owner;
+ AsyncTask(ENamedThreads::GameThread, [WeakOwner]() {
+ if (auto *Self = WeakOwner.Get())
+ UE_LOG(LogTemp, Log, TEXT("DMVOP: Connected"));
+ });
+
+ // ── Main loop ──
+ TArray<uint8> Buf;
+ Buf.SetNumUninitialized(4096);
+ FString Partial;
+
+#ifdef _WIN32
+ // Raw socket for comparison (same destination, bypasses UE layer)
+ SOCKET RawSock = INVALID_SOCKET;
+ bool bRawOK = false;
+ RawSock = socket(AF_INET, SOCK_STREAM, IPPROTO_TCP);
+ if (RawSock != INVALID_SOCKET) {
+ struct sockaddr_in RawAddr;
+ FMemory::Memzero(&RawAddr, sizeof(RawAddr));
+ RawAddr.sin_family = AF_INET;
+ RawAddr.sin_port = htons(Port);
+ inet_pton(AF_INET, TCHAR_TO_UTF8(*Host), &RawAddr.sin_addr);
+ if (connect(RawSock, (struct sockaddr *)&RawAddr, sizeof(RawAddr)) == 0) {
+ bRawOK = true;
+ u_long Mode = 1;
+ ioctlsocket(RawSock, FIONBIO, &Mode);
+ UE_LOG(LogTemp, Log, TEXT("DMVOP: Raw socket connected for comparison"));
+ } else {
+ closesocket(RawSock);
+ RawSock = INVALID_SOCKET;
+ }
}
+#endif
- UE_LOG(LogTemp, Log, TEXT("DMVOP: Connected to %s:%d"), *Host, Port);
-
- // Start reader thread
- bRunning.store(true);
- TWeakObjectPtr<UDMVOPClient> WeakThis(this);
- FSocket *Sock = Socket;
-
- ReaderThread = std::thread([WeakThis, Sock, this]() {
- TArray<uint8> Buf;
- Buf.SetNumUninitialized(4096);
- FString Partial;
-
- while (this->bRunning.load()) {
- int32 Read = 0;
- if (!Sock->Recv(Buf.GetData(), Buf.Num(), Read,
- ESocketReceiveFlags::None) ||
- Read == 0) {
- std::this_thread::sleep_for(std::chrono::milliseconds(10));
- continue;
+ while (bRun) {
+ FDateTime TickStart = FDateTime::UtcNow();
+
+ // Peek check
+ Socket->SetNonBlocking(true);
+ int32 DummyRead = 0;
+ uint8 Dummy;
+ bool bPeekOK =
+ Socket->Recv(&Dummy, 1, DummyRead, ESocketReceiveFlags::Peek);
+ Socket->SetNonBlocking(false);
+
+ if (!bPeekOK) {
+ break;
+ }
+
+ // Receive data
+ while (bRun) {
+ uint32 PendingSize = 0;
+ bool bHasData = Socket->HasPendingData(PendingSize) && PendingSize > 0;
+
+#ifdef _WIN32
+ // Compare: does raw socket have data when UE socket doesn't?
+ bool bRawHasData = false;
+ if (bRawOK) {
+ char RawBuf[1];
+ int RawRead = recv(RawSock, RawBuf, 1, MSG_PEEK);
+ bRawHasData = (RawRead > 0);
+ if (bHasData != bRawHasData) {
+ static int LogCount = 0;
+ if (++LogCount <= 5) {
+ UE_LOG(LogTemp, Log, TEXT("DMVOP: COMPARE UE=%d Raw=%d (Peek=%d)"),
+ bHasData, bRawHasData, RawRead);
+ }
+ }
}
+#endif
- Partial += FString(
- Read, UTF8_TO_TCHAR(reinterpret_cast<const char *>(Buf.GetData())));
+ if (!bHasData) {
+ break;
+ }
+
+ Buf.SetNumUninitialized(FMath::Max(PendingSize, 4096u));
+ int32 ReadNow = 0;
+ if (!Socket->Recv(Buf.GetData(), (int32)PendingSize, ReadNow,
+ ESocketReceiveFlags::None) ||
+ ReadNow <= 0) {
+ break;
+ }
+
+ Partial +=
+ FString(ReadNow,
+ UTF8_TO_TCHAR(reinterpret_cast<const char *>(Buf.GetData())));
int32 Idx;
while (Partial.FindChar('\n', Idx)) {
@@ -86,13 +181,67 @@ void UDMVOPClient::Connect(const FString &Host, int32 Port) {
Text = Line.Mid(Comma + 1);
}
- AsyncTask(ENamedThreads::GameThread, [WeakThis, Vol, Text]() {
- if (auto *Self = WeakThis.Get())
- Self->OnVoiceInput.Broadcast(Vol, Text);
+ TWeakObjectPtr<UDMVOPClient> InnerOwner = Owner;
+ AsyncTask(ENamedThreads::GameThread, [InnerOwner, Vol, Text]() {
+ if (auto *Self = InnerOwner.Get())
+ Self->DispatchVoiceInput(Vol, Text);
});
}
}
- });
+
+ // Sleep
+ FTimespan Elapsed = FDateTime::UtcNow() - TickStart;
+ float SleepSec = 0.008f - (float)Elapsed.GetTotalSeconds();
+ if (SleepSec > 0.0f) {
+ FPlatformProcess::Sleep(SleepSec);
+ }
+ }
+
+ // ── Cleanup ──
+#ifdef _WIN32
+ if (bRawOK) {
+ closesocket(RawSock);
+ }
+#endif
+
+ if (Socket) {
+ Socket->Close();
+ delete Socket;
+ Socket = nullptr;
+ }
+ return 0;
+}
+
+void FDMVOPWorker::Stop() { bRun = false; }
+
+// ═════════════════════════════════════════════════════════════════════════════
+// UDMVOPClient
+// ═════════════════════════════════════════════════════════════════════════════
+
+UDMVOPClient::UDMVOPClient() {}
+
+UDMVOPClient::~UDMVOPClient() {
+ if (Worker.IsValid()) {
+ Worker.Reset();
+ }
+}
+
+void UDMVOPClient::Connect(const FString &Host, int32 Port) {
+ if (Worker.IsValid()) {
+ UE_LOG(LogTemp, Warning, TEXT("DMVOP: Already connected"));
+ return;
+ }
+ UE_LOG(LogTemp, Log, TEXT("DMVOP: Connecting to %s:%d..."), *Host, Port);
+ Worker = MakeShareable(new FDMVOPWorker(this, Host, Port));
+ Worker->Start();
}
-void UDMVOPClient::Disconnect() { cleanup(); }
+void UDMVOPClient::Disconnect() {
+ if (Worker.IsValid()) {
+ Worker.Reset();
+ }
+}
+
+void UDMVOPClient::DispatchVoiceInput(float Vol, const FString &Text) {
+ OnVoiceInput.Broadcast(Vol, Text);
+}
diff --git a/binding/unreal/DMVOPBridge/Source/DMVOPBridge/DMVOPBridgeClient.h b/binding/unreal/DMVOPBridge/Source/DMVOPBridge/DMVOPBridgeClient.h
index b82c2dc..63542c7 100644
--- a/binding/unreal/DMVOPBridge/Source/DMVOPBridge/DMVOPBridgeClient.h
+++ b/binding/unreal/DMVOPBridge/Source/DMVOPBridge/DMVOPBridgeClient.h
@@ -2,18 +2,42 @@
#include "Async/Async.h"
#include "CoreMinimal.h"
+#include "HAL/Runnable.h"
+#include "HAL/RunnableThread.h"
+#include "HAL/ThreadSafeBool.h"
#include "SocketSubsystem.h"
#include "Sockets.h"
#include "UObject/NoExportTypes.h"
#include "UObject/WeakObjectPtr.h"
#include <atomic>
-#include <thread>
#include "DMVOPBridgeClient.generated.h"
DECLARE_DYNAMIC_MULTICAST_DELEGATE_TwoParams(FOnDMVOPVoiceInput, float, Volume,
const FString &, Sentence);
+class FDMVOPWorker : public FRunnable {
+public:
+ FDMVOPWorker(TWeakObjectPtr<class UDMVOPClient> InOwner, FString InHost,
+ int32 InPort);
+ virtual ~FDMVOPWorker();
+
+ void Start();
+
+ // FRunnable
+ virtual bool Init() override;
+ virtual uint32 Run() override;
+ virtual void Stop() override;
+
+private:
+ FThreadSafeBool bRun;
+ FSocket *Socket;
+ TWeakObjectPtr<UDMVOPClient> Owner;
+ FString Host;
+ int32 Port;
+ FRunnableThread *Thread;
+};
+
UCLASS(BlueprintType)
class DMVOPBRIDGE_API UDMVOPClient : public UObject {
GENERATED_BODY()
@@ -31,11 +55,9 @@ public:
UPROPERTY(BlueprintAssignable, Category = "DMVOP")
FOnDMVOPVoiceInput OnVoiceInput;
-private:
- void cleanup();
+ void DispatchVoiceInput(float Vol, const FString &Text);
- FSocket *Socket;
+private:
+ TSharedPtr<FDMVOPWorker> Worker;
FString PartialLine;
- std::thread ReaderThread;
- std::atomic<bool> bRunning{false};
};