载入中...
搜索中...
未找到
TcpSocket.cpp
浏览该文件的文档.
1#include "network/TcpSocket.h"
2#include "network/Network.h"
3#include "network/NetWorker.h"
4#include "data/ByteData.h"
5
6#include <Poco/Net/StreamSocket.h>
7#include <Poco/Net/ServerSocket.h>
8#include <Poco/Net/SocketAddress.h>
9#include <Poco/Net/NetException.h>
10#include <Poco/Net/SocketDefs.h>
11#include <Poco/Timespan.h>
12
13#include <cstring>
14
15namespace eve::network {
16
17TcpSocket::TcpSocket(Network* net) : net_(net) {}
18
22
23void TcpSocket::setConnectedSocket(std::unique_ptr<Poco::Net::StreamSocket> sock) {
24 stream_ = std::move(sock);
25 connected_ = stream_ != nullptr;
26 if (stream_) {
27 try {
28 peer_ = stream_->peerAddress().toString();
29 } catch (...) {
30 peer_.clear();
31 }
32 }
33}
34
35Poco::Net::StreamSocket* TcpSocket::stream() {
36 return stream_.get();
37}
38
39Poco::Net::ServerSocket* TcpSocket::server() {
40 return server_.get();
41}
42
43void TcpSocket::pushAccepted(std::unique_ptr<TcpSocket> sock) {
44 std::lock_guard<std::mutex> lock(acceptMu_);
45 accepted_.push_back(std::move(sock));
46}
47
48bool TcpSocket::takeAccepted(std::unique_ptr<TcpSocket>& out) {
49 std::lock_guard<std::mutex> lock(acceptMu_);
50 if (accepted_.empty()) return false;
51 out = std::move(accepted_.front());
52 accepted_.erase(accepted_.begin());
53 return true;
54}
55
56bool TcpSocket::connect(std::string host, uint16_t port) {
57 if (!net_ || !net_->worker()) return false;
58 auto* self = this;
59 net_->worker()->submit([self, host, port]() {
62 c.handle = self;
63 try {
64 auto sock = std::make_unique<Poco::Net::StreamSocket>();
65 int timeoutMs = self->net_->getTimeout();
66 sock->connect(Poco::Net::SocketAddress(host, port),
67 Poco::Timespan(timeoutMs / 1000, (timeoutMs % 1000) * 1000));
68 sock->setBlocking(false);
69 self->setConnectedSocket(std::move(sock));
70 c.type = NetEvType::Conn;
71 c.reason = "ok";
72 c.peer = self->getPeer();
73 self->net_->post(std::move(c));
74 self->net_->watchTcp(self);
75 } catch (const Poco::Net::ConnectionRefusedException&) {
76 c.type = NetEvType::Err;
77 c.reason = "refused";
78 self->net_->post(std::move(c));
79 NetCompletion fail;
80 fail.type = NetEvType::Conn;
81 fail.kind = NetKind::Tcp;
82 fail.handle = self;
83 fail.reason = "fail";
84 self->net_->post(std::move(fail));
85 } catch (const Poco::Exception& ex) {
86 c.type = NetEvType::Err;
87 std::string msg = ex.displayText();
88 if (msg.find("timed") != std::string::npos || msg.find("Timeout") != std::string::npos)
89 c.reason = "timeout";
90 else if (msg.find("host") != std::string::npos || msg.find("DNS") != std::string::npos)
91 c.reason = "dns";
92 else
93 c.reason = "refused";
94 self->net_->post(std::move(c));
95 NetCompletion fail;
96 fail.type = NetEvType::Conn;
97 fail.kind = NetKind::Tcp;
98 fail.handle = self;
99 fail.reason = "fail";
100 self->net_->post(std::move(fail));
101 }
102 });
103 return true;
104}
105
106bool TcpSocket::listen(uint16_t port) {
107 try {
108 server_ = std::make_unique<Poco::Net::ServerSocket>(port);
109 server_->setBlocking(false);
110 listening_ = true;
111 if (net_) net_->watchTcp(this);
112 return true;
113 } catch (...) {
114 listening_ = false;
115 server_.reset();
116 if (net_) {
119 c.kind = NetKind::Tcp;
120 c.handle = this;
121 c.reason = "refused";
122 net_->post(std::move(c));
123 }
124 return false;
125 }
126}
127
129 std::unique_ptr<TcpSocket> sock;
130 if (!takeAccepted(sock)) return nullptr;
131 return sock.release();
132}
133
135 if (!data || !net_) return false;
136 size_t n = data->getSize();
137 return queueSend(data->getData(), n);
138}
139
140bool TcpSocket::sendString(std::string s) {
141 if (s.empty()) return false;
142 eve::data::ByteData data(s.data(), s.size());
143 return send(&data);
144}
145
146bool TcpSocket::queueSend(const void* d, size_t n) {
147 if (!d || n == 0) return false;
148 const char* p = static_cast<const char*>(d);
149 bool overLimit = false;
150 {
151 std::lock_guard<std::mutex> lock(sendMu_);
152 overLimit = pendingSend_.size() + n > kMaxSendBuffer;
153 if (!overLimit) pendingSend_.insert(pendingSend_.end(), p, p + n);
154 }
155 if (overLimit) {
158 c.kind = NetKind::Tcp;
159 c.handle = this;
160 c.reason = "limit";
161 if (net_) net_->post(std::move(c));
162 return false;
163 }
164 return true;
165}
166
168 std::lock_guard<std::mutex> lock(sendMu_);
169 return pendingSend_.size();
170}
171
173 std::lock_guard<std::mutex> lock(sendMu_);
174 pendingSend_.clear();
175}
176
178 if (!stream_ || !connected_) return;
179
180 std::vector<char> local;
181 {
182 std::lock_guard<std::mutex> lock(sendMu_);
183 if (pendingSend_.empty()) return;
184 local.swap(pendingSend_);
185 }
186
187 size_t off = 0;
188 bool failed = false;
189 while (off < local.size()) {
190 try {
191 int n = stream_->sendBytes(local.data() + static_cast<long>(off),
192 static_cast<int>(local.size() - off));
193 if (n <= 0) break; // non-blocking would-block (Poco returns -1) or zero progress
194 off += static_cast<size_t>(n);
195 } catch (const Poco::TimeoutException&) {
196 break; // would block; keep the remainder for the next poll
197 } catch (const Poco::IOException& e) {
198 // Poco maps EAGAIN/WSAEWOULDBLOCK to IOException("Operation would
199 // block") on some platforms instead of returning -1. Non-fatal.
200 if (e.code() == POCO_EWOULDBLOCK || e.code() == POCO_EAGAIN) break;
201 failed = true;
202 break;
203 } catch (...) {
204 failed = true;
205 break;
206 }
207 }
208
209 if (failed && net_) {
212 c.kind = NetKind::Tcp;
213 c.handle = this;
214 c.reason = "closed";
215 net_->post(std::move(c));
216 close();
217 }
218
219 if (!failed && off < local.size()) {
220 std::lock_guard<std::mutex> lock(sendMu_);
221 pendingSend_.insert(pendingSend_.begin(),
222 local.begin() + static_cast<long>(off),
223 local.end());
224 }
225}
226
228 if (net_) net_->unwatchTcp(this);
229 connected_ = false;
230 listening_ = false;
232 try {
233 if (stream_) stream_->close();
234 } catch (...) {
235 }
236 stream_.reset();
237 try {
238 if (server_) server_->close();
239 } catch (...) {
240 }
241 server_.reset();
242}
243
245 return connected_ && stream_ != nullptr;
246}
247
248std::string TcpSocket::getPeer() const {
249 return peer_;
250}
251
252uint16_t TcpSocket::getLocalPort() const {
253 try {
254 if (server_) return static_cast<uint16_t>(server_->address().port());
255 if (stream_) return static_cast<uint16_t>(stream_->address().port());
256 } catch (...) {
257 }
258 return 0;
259}
260
261} // namespace eve::network
glm::vec3 n
Definition Grass.cpp:64
uint32_t c
glm::vec4 p[6]
Light2D::Data * data
int d
uint32_t s
Definition Weather.cpp:28
In-memory byte buffer implementing eve::Data (ref-counted).
Definition ByteData.h:11
void submit(std::function< void()> job)
Queues an arbitrary blocking job to run on the worker thread.
Definition NetWorker.cpp:40
Network module: TCP/UDP/HTTP factories, background worker, and completion event plumbing....
Definition Network.h:30
void watchTcp(TcpSocket *sock)
Internal: register/unregister sockets for pump polling.
Definition Network.cpp:175
NetWorker * worker() const
Background worker thread handle (advanced).
Definition Network.h:73
void unwatchTcp(TcpSocket *sock)
Definition Network.cpp:181
void post(NetCompletion c)
Posts a completion to the worker queue (thread-safe).
Definition Network.cpp:167
TCP socket backed by Poco::Net; supports both client (connect) and server (listen/accept) roles....
Definition TcpSocket.h:28
void pushAccepted(std::unique_ptr< TcpSocket > sock)
Definition TcpSocket.cpp:43
bool listen(uint16_t port)
Starts listening on the given port; true on success.
bool takeAccepted(std::unique_ptr< TcpSocket > &out)
Definition TcpSocket.cpp:48
Poco::Net::StreamSocket * stream()
Definition TcpSocket.cpp:35
uint16_t getLocalPort() const
Local listening port, or 0.
bool queueSend(const void *d, size_t n)
Append bytes to the pending send queue (worker-thread flush via pollSockets).
bool sendString(std::string s)
Sends a string payload.
TcpSocket(Network *net)
Creates an unconnected socket owned by the given module.
Definition TcpSocket.cpp:17
void setConnectedSocket(std::unique_ptr< Poco::Net::StreamSocket > sock)
Definition TcpSocket.cpp:23
void close()
Closes the socket.
void flushSend()
Push pending bytes onto the socket; called from Network::pollSockets.
std::string getPeer() const
Remote peer address string, e.g. "1.2.3.4:5678".
Poco::Net::ServerSocket * server()
Definition TcpSocket.cpp:39
size_t pendingSendBytes() const
Internal: queued-outgoing byte counter for back-pressure.
TcpSocket * accept()
Accepts one pending client, or nullptr if none.
bool isConnected() const
True while a stream is connected.
bool connect(std::string host, uint16_t port)
Connects to host:port; true on success.
Definition TcpSocket.cpp:56
bool send(eve::data::ByteData *data)
Sends a framed byte payload; true when queued/accepted.
constexpr size_t kMaxSendBuffer
Maximum outgoing send buffer size per socket (1 MiB).
Definition NetTypes.h:11
One asynchronous network result/event. handle points at the originating TcpSocket/UdpSocket/HttpReque...
Definition NetTypes.h:33