载入中...
搜索中...
未找到
Network.cpp
浏览该文件的文档.
1#include "network/Network.h"
2#include "data/ByteData.h"
3#include "network/Channel.h"
5#include "network/NetHost.h"
6#include "network/NetRpc.h"
7#include "network/NetStream.h"
8#include "network/NetWorker.h"
9#include "network/Session.h"
10#include "network/TcpSocket.h"
11#include "network/UdpLink.h"
12#include "network/UdpSocket.h"
14
15#include <Poco/Net/StreamSocket.h>
16#include <Poco/Net/ServerSocket.h>
17#include <Poco/Net/DatagramSocket.h>
18#include <Poco/Net/SocketAddress.h>
19#include <Poco/Net/NetException.h>
20#include <Poco/Exception.h>
21
22#include <simplesquirrel/simplesquirrel.hpp>
23#include <algorithm>
24#include <chrono>
25#include <functional>
26#include <thread>
27
28namespace eve::network {
29
30namespace {
31
32void callScript2(const ssq::Object& obj, int64_t a, const std::string& s) {
33 if (obj.isEmpty()) return;
34 ssq::Function f = obj.toFunction();
35 if (f.isEmpty()) return;
36 HSQUIRRELVM raw = f.getHandle();
37 SQInteger top = sq_gettop(raw);
38 sq_pushobject(raw, f.getRaw());
39 sq_pushroottable(raw);
42 sq_call(raw, 3, SQFalse, SQTrue);
43 sq_settop(raw, top);
44}
45
46void callScript3(const ssq::Object& obj, int64_t a, int64_t b, const std::string& s) {
47 if (obj.isEmpty()) return;
48 ssq::Function f = obj.toFunction();
49 if (f.isEmpty()) return;
50 HSQUIRRELVM raw = f.getHandle();
51 SQInteger top = sq_gettop(raw);
52 sq_pushobject(raw, f.getRaw());
53 sq_pushroottable(raw);
57 sq_call(raw, 4, SQFalse, SQTrue);
58 sq_settop(raw, top);
59}
60
61} // namespace
62
64
66 worker_ = std::make_unique<NetWorker>(this);
67 worker_->start();
68 eve::cap::provide<eve::service::INetwork>(this);
69}
70
72 if (worker_) worker_->stop();
73 worker_.reset();
74}
75
76std::unique_ptr<TcpSocket> Network::makeTcp() { return std::make_unique<TcpSocket>(this); }
77TcpSocket* Network::newTcp() { return makeTcp().release(); }
78
79std::unique_ptr<UdpSocket> Network::makeUdp() { return std::make_unique<UdpSocket>(this); }
80UdpSocket* Network::newUdp() { return makeUdp().release(); }
81
82std::unique_ptr<HttpRequest> Network::makeHttp(std::string method, std::string url) {
83 return std::make_unique<HttpRequest>(this, std::move(method), std::move(url));
84}
85HttpRequest* Network::newHttp(std::string method, std::string url) {
86 return makeHttp(std::move(method), std::move(url)).release();
87}
88
89bool Network::httpRequest(const std::string& method, const std::string& url,
90 const std::string& body, int timeoutMs, int& status,
91 std::string& responseBody) {
92 auto ownedRequest = makeHttp(method, url);
93 HttpRequest* req = ownedRequest.get();
94 if (!body.empty()) req->setBodyString(body);
95 req->setTimeout(timeoutMs > 0 ? timeoutMs : 10000);
96 if (!req->submit()) return false;
97
98 const auto deadline = std::chrono::steady_clock::now() +
99 std::chrono::milliseconds(timeoutMs > 0 ? timeoutMs : 1);
100 while (std::chrono::steady_clock::now() < deadline) {
101 std::vector<NetCompletion> out;
102 drainCompletions(out);
103 for (const auto& c : out) {
104 if (c.handle != req) continue;
105 if (c.type == NetEvType::HttpResp) {
106 status = c.status;
107 if (c.bytes) responseBody.assign(c.bytes->begin(), c.bytes->end());
108 return true;
109 }
110 if (c.type == NetEvType::Err) return false;
111 }
112 std::this_thread::sleep_for(std::chrono::milliseconds(10));
113 }
114 return false;
115}
116
117std::unique_ptr<Channel> Network::makeChannel(TcpSocket* socket) {
118 if (!socket) return nullptr;
119 auto ch = std::make_unique<Channel>(socket);
120 bindChannel(socket, ch.get());
121 return ch;
122}
123Channel* Network::newChannel(TcpSocket* socket) { return makeChannel(socket).release(); }
124
125std::unique_ptr<Session> Network::makeSession() { return std::make_unique<Session>(); }
126Session* Network::newSession() { return makeSession().release(); }
127
128std::unique_ptr<NetWriter> Network::makeWriter() { return std::make_unique<NetWriter>(); }
129NetWriter* Network::newWriter() { return makeWriter().release(); }
130
131std::unique_ptr<NetReader> Network::makeReader(std::string bytes) {
132 auto r = std::make_unique<NetReader>();
133 r->setBytes(std::move(bytes));
134 return r;
135}
136NetReader* Network::newReader(std::string bytes) { return makeReader(std::move(bytes)).release(); }
137
138std::unique_ptr<UdpLink> Network::makeUdpLink(UdpSocket* socket) {
139 if (!socket) return nullptr;
140 auto link = std::make_unique<UdpLink>(this, socket);
141 bindUdpLink(socket, link.get());
142 return link;
143}
144UdpLink* Network::newUdpLink(UdpSocket* socket) { return makeUdpLink(socket).release(); }
145
146std::unique_ptr<NetRpc> Network::makeRpc(UdpLink* link) {
147 return link ? std::make_unique<NetRpc>(link) : nullptr;
148}
149NetRpc* Network::newRpc(UdpLink* link) { return makeRpc(link).release(); }
150
151std::unique_ptr<NetHost> Network::makeHost() { return std::make_unique<NetHost>(this); }
152NetHost* Network::newHost() { return makeHost().release(); }
153
154void Network::bindUdpLink(UdpSocket* sock, UdpLink* link) {
155 if (sock) udpLinks_[sock] = link;
156}
157
158void Network::unbindUdpLink(UdpSocket* sock) {
159 if (sock) udpLinks_.erase(sock);
160}
161
162void Network::bindUdpHost(UdpSocket* sock, NetHost* host) {
163 if (sock) udpHosts_[sock] = host;
164}
165
166void Network::unbindUdpHost(UdpSocket* sock) {
167 if (sock) udpHosts_.erase(sock);
168}
169
171 timeoutMs_ = ms;
172}
173
175 return timeoutMs_;
176}
177
178void Network::setVerifySsl(bool verify) {
179 verifySsl_ = verify;
180}
181
183 return verifySsl_;
184}
185
186void Network::post(NetCompletion c) {
187 ++completions_;
188 ++telemetryRevision_;
189 if (c.bytes && (c.type == NetEvType::Data || c.type == NetEvType::HttpResp))
190 receivedBytes_ += c.bytes->size();
191 if (c.type == NetEvType::Err) ++errors_;
192 if (c.type == NetEvType::Conn && c.reason == "ok") ++connections_;
193 if (worker_) worker_->post(std::move(c));
194}
195
196void Network::recordSent(size_t bytes) {
197 sentBytes_ += bytes;
198 ++telemetryRevision_;
199}
200
203 result.revision = telemetryRevision_.load();
204 result.sentBytes = sentBytes_.load();
205 result.receivedBytes = receivedBytes_.load();
206 result.completions = completions_.load();
207 result.errors = errors_.load();
208 result.connections = connections_.load();
209 {
210 std::lock_guard<std::mutex> lock(watchMu_);
211 result.watchedTcp = watchedTcp_.size();
212 result.watchedUdp = watchedUdp_.size();
213 for (const auto* socket : watchedTcp_) if (socket) result.queuedTcpBytes += socket->pendingSendBytes();
214 }
215 {
216 std::lock_guard<std::mutex> lock(channelMu_);
217 result.channels = channels_.size();
218 }
219 return result;
220}
221
223 sentBytes_ = receivedBytes_ = completions_ = errors_ = connections_ = 0;
224 ++telemetryRevision_;
225}
226
227void Network::drainCompletions(std::vector<NetCompletion>& out) {
228 if (worker_) worker_->drain(out);
229}
230
231void Network::watchTcp(TcpSocket* sock) {
232 std::lock_guard<std::mutex> lock(watchMu_);
233 if (std::find(watchedTcp_.begin(), watchedTcp_.end(), sock) == watchedTcp_.end()) {
234 watchedTcp_.push_back(sock);
235 ++telemetryRevision_;
236 }
237}
238
239void Network::unwatchTcp(TcpSocket* sock) {
240 std::lock_guard<std::mutex> lock(watchMu_);
241 const auto before=watchedTcp_.size();
242 watchedTcp_.erase(std::remove(watchedTcp_.begin(), watchedTcp_.end(), sock), watchedTcp_.end());
243 if(before!=watchedTcp_.size())++telemetryRevision_;
244}
245
246void Network::watchUdp(UdpSocket* sock) {
247 std::lock_guard<std::mutex> lock(watchMu_);
248 if (std::find(watchedUdp_.begin(), watchedUdp_.end(), sock) == watchedUdp_.end()) {
249 watchedUdp_.push_back(sock);
250 ++telemetryRevision_;
251 }
252}
253
254void Network::unwatchUdp(UdpSocket* sock) {
255 std::lock_guard<std::mutex> lock(watchMu_);
256 const auto before=watchedUdp_.size();
257 watchedUdp_.erase(std::remove(watchedUdp_.begin(), watchedUdp_.end(), sock), watchedUdp_.end());
258 if(before!=watchedUdp_.size())++telemetryRevision_;
259}
260
261void Network::bindChannel(TcpSocket* sock, Channel* ch) {
262 std::lock_guard<std::mutex> lock(channelMu_);
263 channels_[sock] = ch;
264 ++telemetryRevision_;
265}
266
267void Network::unbindChannel(TcpSocket* sock) {
268 std::lock_guard<std::mutex> lock(channelMu_);
269 if(channels_.erase(sock))++telemetryRevision_;
270}
271
272Channel* Network::channelFor(TcpSocket* sock) const {
273 std::lock_guard<std::mutex> lock(channelMu_);
274 auto it = channels_.find(sock);
275 return it == channels_.end() ? nullptr : it->second;
276}
277
278void Network::pollSockets() {
279 std::vector<TcpSocket*> tcpCopy;
280 std::vector<UdpSocket*> udpCopy;
281 {
282 std::lock_guard<std::mutex> lock(watchMu_);
283 tcpCopy = watchedTcp_;
284 udpCopy = watchedUdp_;
285 }
286 for (TcpSocket* sock : tcpCopy) {
287 if (!sock) continue;
288 if (sock->isListening() && sock->server()) {
289 try {
290 Poco::Net::StreamSocket ss = sock->server()->acceptConnection();
291 ss.setBlocking(false);
292 auto peer = std::make_unique<TcpSocket>(this);
293 auto stream = std::make_unique<Poco::Net::StreamSocket>(ss);
294 peer->setConnectedSocket(std::move(stream));
295 TcpSocket* raw = peer.get();
296 sock->pushAccepted(std::move(peer));
297 watchTcp(raw);
298 NetCompletion c;
299 c.type = NetEvType::Conn;
300 c.kind = NetKind::Tcp;
301 c.handle = raw;
302 c.reason = "ok";
303 c.peer = raw->getPeer();
304 post(std::move(c));
305 } catch (...) {
306 }
307 }
308 if (sock->isConnected() && sock->stream()) {
309 sock->flushSend();
310 try {
311 char buf[64 * 1024];
312 int n = sock->stream()->receiveBytes(buf, sizeof(buf));
313 if (n > 0) {
314 NetCompletion c;
315 c.type = NetEvType::Data;
316 c.kind = NetKind::Tcp;
317 c.handle = sock;
318 c.bytes = std::make_shared<std::vector<char>>(buf, buf + n);
319 c.peer = sock->getPeer();
320 post(std::move(c));
321 } else if (n == 0) {
322 NetCompletion c;
323 c.type = NetEvType::Err;
324 c.kind = NetKind::Tcp;
325 c.handle = sock;
326 c.reason = "closed";
327 post(std::move(c));
328 unwatchTcp(sock);
329 }
330 } catch (const Poco::Exception&) {
331 } catch (...) {
332 }
333 }
334 }
335 for (UdpSocket* sock : udpCopy) {
336 if (!sock || !sock->datagram()) continue;
337 for (int i = 0; i < 64; ++i) {
338 try {
339 char buf[64 * 1024];
340 Poco::Net::SocketAddress sender;
341 int n = sock->datagram()->receiveFrom(buf, sizeof(buf), sender);
342 if (n <= 0) break; // would-block
343 NetCompletion c;
344 c.type = NetEvType::Data;
345 c.kind = NetKind::Udp;
346 c.handle = sock;
347 c.bytes = std::make_shared<std::vector<char>>(buf, buf + n);
348 c.peer = sender.toString();
349 post(std::move(c));
350 } catch (const Poco::Exception&) {
351 break;
352 } catch (...) {
353 break;
354 }
355 }
356 }
357}
358
359void Network::emitCompletion(const NetCompletion& c) {
360 auto* ev = eve::ModuleManager::getInstance<eve::platform_event::PlatformEvent>("PlatformEvent");
361 if (!ev) return;
362
365
366 std::vector<Variant> args;
367 switch (c.type) {
368 case NetEvType::Conn:
369 args.push_back(Variant::makePtr(c.handle));
370 args.push_back(Variant::makeString(c.reason.empty() ? "ok" : c.reason));
371 ev->push(std::make_unique<Message>("netconn", args));
372 break;
373 case NetEvType::Data: {
374 Channel* ch = channelFor(static_cast<TcpSocket*>(c.handle));
375 if (ch && c.bytes) {
376 ch->feed(*c.bytes);
377 return;
378 }
379 eve::data::ByteData* bd = nullptr;
380 if (c.bytes && !c.bytes->empty())
381 bd = new eve::data::ByteData(c.bytes->data(), c.bytes->size());
382 args.push_back(Variant::makePtr(c.handle));
383 args.push_back(Variant::makeOwnedPtr(bd));
384 args.push_back(Variant::makeString(c.peer));
385 ev->push(std::make_unique<Message>("netdata", args));
386 break;
387 }
388 case NetEvType::Err: {
389 Channel* ch = channelFor(static_cast<TcpSocket*>(c.handle));
390 if (ch) {
391 args.push_back(Variant::makePtr(ch));
392 args.push_back(Variant::makeString(c.reason));
393 ev->push(std::make_unique<Message>("chclose", args));
394 }
395 args.clear();
396 args.push_back(Variant::makePtr(c.handle));
397 args.push_back(Variant::makeString(c.reason));
398 ev->push(std::make_unique<Message>("neterr", args));
399 break;
400 }
401 case NetEvType::HttpResp: {
402 eve::data::ByteData* bd = nullptr;
403 if (c.bytes && !c.bytes->empty())
404 bd = new eve::data::ByteData(c.bytes->data(), c.bytes->size());
405 args.push_back(Variant::makePtr(c.handle));
406 args.push_back(Variant::makeInt(c.status));
407 args.push_back(Variant::makeOwnedPtr(bd));
408 ev->push(std::make_unique<Message>("httpresp", args));
409 break;
410 }
411 case NetEvType::ChMsg: {
412 eve::data::ByteData* bd = nullptr;
413 if (c.bytes && !c.bytes->empty())
414 bd = new eve::data::ByteData(c.bytes->data(), c.bytes->size());
415 args.push_back(Variant::makePtr(c.handle));
416 args.push_back(Variant::makeOwnedPtr(bd));
417 ev->push(std::make_unique<Message>("chmsg", args));
418 break;
419 }
421 args.push_back(Variant::makePtr(c.handle));
422 args.push_back(Variant::makeString(c.reason));
423 ev->push(std::make_unique<Message>("chclose", args));
424 break;
425 }
426}
427
429 std::vector<NetCompletion> batch;
430 if (worker_) worker_->drain(batch);
431 const int64_t now = std::chrono::duration_cast<std::chrono::milliseconds>(
432 std::chrono::steady_clock::now().time_since_epoch())
433 .count();
434 for (auto& c : batch) {
435 if (c.type == NetEvType::Data && c.kind == NetKind::Udp) {
436 auto* sock = static_cast<UdpSocket*>(c.handle);
437 auto hit = udpHosts_.find(sock);
438 if (hit != udpHosts_.end()) {
439 if (c.bytes) hit->second->onDatagram(*c.bytes, c.peer);
440 continue;
441 }
442 auto lit = udpLinks_.find(sock);
443 if (lit != udpLinks_.end()) {
444 if (c.bytes) lit->second->onDatagram(*c.bytes, c.peer);
445 continue;
446 }
447 }
448 emitCompletion(c);
449 }
450 for (auto& kv : udpHosts_) kv.second->pump(now);
451 for (auto& kv : udpLinks_) kv.second->pump(now);
452}
453
454void Network::expose(ssq::Table& table) {
455 auto cls = table.addClass(name, Network::create, false);
456 expose(cls);
457
458 auto tcp = table.addClass<TcpSocket>(
459 "TcpSocket", std::function<TcpSocket*()>([]() { return new TcpSocket(Network::create()); }), true);
460 tcp.addFunc("connect", &TcpSocket::connect);
461 tcp.addFunc("listen", &TcpSocket::listen);
462 tcp.addFunc("accept", &TcpSocket::accept);
463 tcp.addFunc("send", &TcpSocket::send);
464 tcp.addFunc("sendString", &TcpSocket::sendString);
465 tcp.addFunc("close", &TcpSocket::close);
466 tcp.addFunc("isConnected", &TcpSocket::isConnected);
467 tcp.addFunc("getPeer", &TcpSocket::getPeer);
468
469 auto udp = table.addClass<UdpSocket>(
470 "UdpSocket", std::function<UdpSocket*()>([]() { return new UdpSocket(Network::create()); }), true);
471 udp.addFunc("bind", &UdpSocket::bind);
472 udp.addFunc("connect", &UdpSocket::connect);
473 udp.addFunc("sendTo", &UdpSocket::sendTo);
474 udp.addFunc("sendToString", &UdpSocket::sendToString);
475 udp.addFunc("send", &UdpSocket::send);
476 udp.addFunc("sendString", &UdpSocket::sendString);
477 udp.addFunc("close", &UdpSocket::close);
478
479 auto http = table.addClass<HttpRequest>(
480 "HttpRequest",
481 std::function<HttpRequest*()>(
482 []() { return new HttpRequest(Network::create(), "GET", "http://127.0.0.1/"); }),
483 true);
484 http.addFunc("setHeader", &HttpRequest::setHeader);
485 http.addFunc("setBody", &HttpRequest::setBody);
486 http.addFunc("setBodyString", &HttpRequest::setBodyString);
487 http.addFunc("setTimeout", &HttpRequest::setTimeout);
488 http.addFunc("setVerifySsl", &HttpRequest::setVerifySsl);
489 http.addFunc("submit", &HttpRequest::submit);
490
491 auto ch = table.addClass<Channel>(
492 "Channel", std::function<Channel*()>([]() { return new Channel(nullptr); }), true);
493 ch.addFunc("sendMsg", &Channel::sendMsg);
494 ch.addFunc("sendMsgString", &Channel::sendMsgString);
495 ch.addFunc("getSocket", &Channel::getSocket);
496
497 auto sess = table.addClass<Session>(
498 "Session", std::function<Session*()>([]() { return new Session(); }), true);
499 sess.addFunc("add", &Session::add);
500 sess.addFunc("get", &Session::get);
501 sess.addFunc("remove", &Session::remove);
502 sess.addFunc("closeAll", &Session::closeAll);
503
504 auto writer = table.addClass<NetWriter>(
505 "NetWriter", std::function<NetWriter*()>([]() { return new NetWriter(); }), true);
506 writer.addFunc("writeU8", [](NetWriter* w, int64_t v) { w->writeU8(static_cast<uint8_t>(v)); });
507 writer.addFunc("writeI8", [](NetWriter* w, int64_t v) { w->writeI8(static_cast<int8_t>(v)); });
508 writer.addFunc("writeU16", [](NetWriter* w, int64_t v) { w->writeU16(static_cast<uint16_t>(v)); });
509 writer.addFunc("writeI16", [](NetWriter* w, int64_t v) { w->writeI16(static_cast<int16_t>(v)); });
510 writer.addFunc("writeU32", [](NetWriter* w, int64_t v) { w->writeU32(static_cast<uint32_t>(v)); });
511 writer.addFunc("writeI32", [](NetWriter* w, int64_t v) { w->writeI32(static_cast<int32_t>(v)); });
512 writer.addFunc("writeU64", [](NetWriter* w, int64_t v) { w->writeU64(static_cast<uint64_t>(v)); });
513 writer.addFunc("writeI64", [](NetWriter* w, int64_t v) { w->writeI64(static_cast<int64_t>(v)); });
514 writer.addFunc("writeF32", &NetWriter::writeF32);
515 writer.addFunc("writeF64", &NetWriter::writeF64);
516 writer.addFunc("writeBool", &NetWriter::writeBool);
517 writer.addFunc("writeString", &NetWriter::writeString);
518 writer.addFunc("writeBytes", [](NetWriter* w, eve::data::ByteData* d) {
519 if (d) w->writeBytes(d->getData(), d->getSize());
520 });
521 writer.addFunc("toString", [](NetWriter* w) { return w->toString(); });
522 writer.addFunc("size", [](NetWriter* w) { return static_cast<int64_t>(w->size()); });
523
524 auto reader = table.addClass<NetReader>(
525 "NetReader", std::function<NetReader*()>([]() { return new NetReader(); }), true);
526 reader.addFunc("init", [](NetReader* r, eve::data::ByteData* d) {
527 return d ? r->init(d->getData(), d->getSize()) : false;
528 });
529 reader.addFunc("initString", [](NetReader* r, const std::string& s) { return r->setBytes(s); });
530 reader.addFunc("u8", [](NetReader* r) { return static_cast<int64_t>(r->u8()); });
531 reader.addFunc("i8", [](NetReader* r) { return static_cast<int64_t>(r->i8()); });
532 reader.addFunc("u16", [](NetReader* r) { return static_cast<int64_t>(r->u16()); });
533 reader.addFunc("i16", [](NetReader* r) { return static_cast<int64_t>(r->i16()); });
534 reader.addFunc("u32", [](NetReader* r) { return static_cast<int64_t>(r->u32()); });
535 reader.addFunc("i32", [](NetReader* r) { return static_cast<int64_t>(r->i32()); });
536 reader.addFunc("u64", [](NetReader* r) { return static_cast<int64_t>(r->u64()); });
537 reader.addFunc("i64", [](NetReader* r) { return static_cast<int64_t>(r->i64()); });
538 reader.addFunc("f32", &NetReader::f32);
539 reader.addFunc("f64", &NetReader::f64);
540 reader.addFunc("bool", &NetReader::b);
541 reader.addFunc("str", [](NetReader* r) { return r->str(); });
542 reader.addFunc("bytes", [](NetReader* r, int64_t n) {
543 auto v = r->bytes(static_cast<size_t>(n));
544 return std::string(v.data(), v.size());
545 });
546 reader.addFunc("remaining", [](NetReader* r) { return static_cast<int64_t>(r->remaining()); });
547 reader.addFunc("pos", [](NetReader* r) { return static_cast<int64_t>(r->pos()); });
548 reader.addFunc("ok", [](NetReader* r) { return r->ok(); });
549
550 auto link = table.addClass<UdpLink>(
551 "UdpLink",
552 std::function<UdpLink*()>([]() { return new UdpLink(nullptr, nullptr); }), true);
553 link.addFunc("setRemote", &UdpLink::setRemote);
554 link.addFunc("setRemoteString", &UdpLink::setRemoteString);
555 link.addFunc("sendReliable", [](UdpLink* l, int64_t ch, const std::string& s) {
556 l->sendString(UdpLink::MsgType::Reliable, static_cast<uint8_t>(ch), s);
557 });
558 link.addFunc("sendUnreliable", [](UdpLink* l, int64_t ch, const std::string& s) {
559 l->sendString(UdpLink::MsgType::Unreliable, static_cast<uint8_t>(ch), s);
560 });
561 link.addFunc("sendOrdered", [](UdpLink* l, int64_t ch, const std::string& s) {
562 l->sendString(UdpLink::MsgType::UnreliableOrdered, static_cast<uint8_t>(ch), s);
563 });
564 link.addFunc("onMessage", [](UdpLink* l, ssq::Object fn) {
565 if (!l) return;
566 l->setMessageHandler(
567 [fn](UdpLink::MsgType, uint8_t ch, const char* d, size_t n) {
568 callScript2(fn, ch, std::string(d, n));
569 });
570 });
571 link.addFunc("setLossRate", &UdpLink::setLossRate);
572 link.addFunc("setTimeoutMs", &UdpLink::setTimeoutMs);
573 link.addFunc("isAlive", [](UdpLink* l) { return l && l->isAlive(); });
574 link.addFunc("peer", &UdpLink::peer);
575 link.addFunc("peerId", [](UdpLink* l) { return static_cast<int64_t>(l->peerId()); });
576 link.addFunc("pendingReliable",
577 [](UdpLink* l) { return static_cast<int64_t>(l->pendingReliable()); });
578 link.addFunc("pendingFragments",
579 [](UdpLink* l) { return static_cast<int64_t>(l->pendingFragments()); });
580
581 auto rpc = table.addClass<NetRpc>(
582 "NetRpc", std::function<NetRpc*()>([]() { return new NetRpc(nullptr); }), true);
583 rpc.addFunc("callRpc", [](NetRpc* r, int64_t msgId, const std::string& payload,
584 bool reliable) {
585 r->callString(static_cast<uint16_t>(msgId), payload, reliable);
586 });
587 rpc.addFunc("registerRpc", [](NetRpc* r, int64_t msgId, ssq::Object fn) {
588 r->registerScript(static_cast<uint16_t>(msgId), fn);
589 });
590
591 auto host = table.addClass<NetHost>(
592 "NetHost", std::function<NetHost*()>([]() { return new NetHost(nullptr); }), true);
593 host.addFunc("start", &NetHost::start);
594 host.addFunc("onMessage", [](NetHost* h, ssq::Object fn) {
595 if (!h) return;
596 h->setMessageHandler(
597 [fn](uint32_t peerId, UdpLink::MsgType, uint8_t ch, const char* d, size_t n) {
598 callScript3(fn, peerId, ch, std::string(d, n));
599 });
600 });
601 host.addFunc("onPeerConnected", [](NetHost* h, ssq::Object fn) {
602 if (!h) return;
603 h->setPeerConnectedHandler([fn](uint32_t id) { callScript2(fn, id, ""); });
604 });
605 host.addFunc("onPeerDisconnected", [](NetHost* h, ssq::Object fn) {
606 if (!h) return;
607 h->setPeerDisconnectedHandler([fn](uint32_t id) { callScript2(fn, id, ""); });
608 });
609 host.addFunc("sendReliable", [](NetHost* h, int64_t peerId, int64_t ch,
610 const std::string& s) {
611 h->sendStringTo(static_cast<uint32_t>(peerId), UdpLink::MsgType::Reliable,
612 static_cast<uint8_t>(ch), s);
613 });
614 host.addFunc("sendUnreliable", [](NetHost* h, int64_t peerId, int64_t ch,
615 const std::string& s) {
616 h->sendStringTo(static_cast<uint32_t>(peerId), UdpLink::MsgType::Unreliable,
617 static_cast<uint8_t>(ch), s);
618 });
619 host.addFunc("sendOrdered", [](NetHost* h, int64_t peerId, int64_t ch,
620 const std::string& s) {
621 h->sendStringTo(static_cast<uint32_t>(peerId), UdpLink::MsgType::UnreliableOrdered,
622 static_cast<uint8_t>(ch), s);
623 });
624 host.addFunc("link", &NetHost::linkByPeerId);
625 host.addFunc("peerCount", [](NetHost* h) { return static_cast<int64_t>(h->peerCount()); });
626 host.addFunc("setLossRate", &NetHost::setLossRate);
627 host.addFunc("setTimeoutMs", &NetHost::setTimeoutMs);
628}
629
630void Network::expose(ssq::Class& cls) {
631 cls.addFunc("getName", &Network::getName);
632 cls.addFunc("newTcp", &Network::newTcp);
633 cls.addFunc("newUdp", &Network::newUdp);
634 cls.addFunc("newHttp", &Network::newHttp);
635 cls.addFunc("newChannel", &Network::newChannel);
636 cls.addFunc("newSession", &Network::newSession);
637 cls.addFunc("newWriter", &Network::newWriter);
638 cls.addFunc("newReader", &Network::newReader);
639 cls.addFunc("newUdpLink", &Network::newUdpLink);
640 cls.addFunc("newRpc", &Network::newRpc);
641 cls.addFunc("newHost", &Network::newHost);
642 cls.addFunc("pump", &Network::pump);
643 cls.addFunc("setTimeout", &Network::setTimeout);
644 cls.addFunc("setVerifySsl", &Network::setVerifySsl);
645}
646
647} // namespace eve::network
Value::Object payload
SQInteger top
float w
Definition AnimClip.cpp:738
const std::string & s
struct SQVM * HSQUIRRELVM
HSQOBJECT cls
Definition ECS.cpp:21
std::uint16_t method
wgpu::PopErrorScopeStatus status
glm::vec3 n
Definition Grass.cpp:63
double r
float v
std::int32_t c
int h
std::uint64_t bytes
std::string name
MeleePoint3 b
Definition MeleeHit.cpp:41
MeleePoint3 a
Definition MeleeHit.cpp:40
#define Module_IMPL(ModuleName, newExpr)
Definition Module.h:26
float f
Duration deadline
bool hit
float d
UIHostHandle host
std::string body
virtual std::string getName() const =0
Returns the name.
In-memory byte buffer implementing eve::Data (ref-counted).
Definition ByteData.h:13
Length-prefixed (big-endian uint32) message framing over a TcpSocket. sendMsg() writes one framed mes...
Definition Channel.h:23
bool sendMsg(eve::data::ByteData *data)
Sends one framed message.
Definition Channel.cpp:38
TcpSocket * getSocket() const
The underlying TCP socket, or nullptr.
Definition Channel.h:39
bool sendMsgString(std::string s)
Sends a string payload as one framed message (empty strings are rejected).
Definition Channel.cpp:59
Asynchronous HTTP request (Poco-based). Configure headers/body, call submit(), then receive the respo...
Definition HttpRequest.h:26
void setBodyString(std::string s)
Sets the request body as a string.
void setTimeout(int ms)
Overrides the module default timeout for this request (ms).
void setBody(eve::data::ByteData *data)
Sets the request body bytes.
void setHeader(std::string k, std::string v)
Sets one request header.
bool submit()
Queues the request on the worker; true when accepted.
void setVerifySsl(bool verify)
Overrides TLS verification for this request.
EVENGINE_API_PLATFORM public API.
Definition NetHost.h:25
void setLossRate(float rate)
Sets the loss rate.
Definition NetHost.cpp:33
bool start(uint16_t port)
Starts .
Definition NetHost.cpp:21
void setTimeoutMs(int ms)
Sets the timeout ms.
Definition NetHost.cpp:38
UdpLink * linkByPeerId(uint32_t peerId) const
Link by peer id.
Definition NetHost.cpp:76
EVENGINE_API_PLATFORM public API.
Definition NetStream.h:71
EVENGINE_API_PLATFORM public API.
Definition NetRpc.h:26
EVENGINE_API_PLATFORM public API.
Definition NetStream.h:20
void writeBool(bool v)
Writes bool.
void writeF32(float v)
Writes f 32.
Definition NetStream.cpp:87
void writeString(const std::string &s)
Writes string.
void writeF64(double v)
Writes f 64.
Definition NetStream.cpp:94
Network module: TCP/UDP/HTTP factories, background worker, and completion event plumbing....
Definition Network.h:35
NetHost * newHost()
Script adapter: creates a caller-owned UDP host.
Definition Network.cpp:152
std::unique_ptr< TcpSocket > makeTcp()
Creates an owning TCP socket handle for C++ callers.
Definition Network.cpp:76
NetReader * newReader(std::string bytes)
Script adapter: creates a caller-owned reader.
Definition Network.cpp:136
std::unique_ptr< NetRpc > makeRpc(UdpLink *link)
Creates an owning RPC handle, or null for a null link.
Definition Network.cpp:146
bool getVerifySsl() const
Returns the verify ssl.
Definition Network.cpp:182
std::unique_ptr< NetReader > makeReader(std::string bytes)
Creates an owning streaming reader handle.
Definition Network.cpp:131
friend class HttpRequest
Definition Network.h:109
void setTimeout(int ms)
Default socket/HTTP timeout in milliseconds.
Definition Network.cpp:170
Session * newSession()
Script adapter: creates a caller-owned session.
Definition Network.cpp:126
void resetTelemetry()
Reset cumulative telemetry counters without affecting live sockets.
Definition Network.cpp:222
std::unique_ptr< Session > makeSession()
Creates an owning session handle.
Definition Network.cpp:125
UdpLink * newUdpLink(UdpSocket *socket)
Script adapter: creates a caller-owned UDP link.
Definition Network.cpp:144
std::unique_ptr< NetWriter > makeWriter()
Creates an owning streaming writer handle.
Definition Network.cpp:128
std::unique_ptr< UdpSocket > makeUdp()
Creates an owning UDP socket handle for C++ callers.
Definition Network.cpp:79
void setVerifySsl(bool verify)
Whether TLS peer certificates are verified (default true).
Definition Network.cpp:178
std::unique_ptr< Channel > makeChannel(TcpSocket *socket)
Creates an owning channel handle, or null for a null socket.
Definition Network.cpp:117
bool httpRequest(const std::string &method, const std::string &url, const std::string &body, int timeoutMs, int &status, std::string &responseBody) override
Synchronous HTTP request (blocks up to timeoutMs on the worker).
Definition Network.cpp:89
HttpRequest * newHttp(std::string method, std::string url)
Script adapter: creates a caller-owned HTTP request.
Definition Network.cpp:85
std::unique_ptr< HttpRequest > makeHttp(std::string method, std::string url)
Creates an owning HTTP request handle for C++ callers.
Definition Network.cpp:82
friend class NetHost
Definition Network.h:110
std::unique_ptr< NetHost > makeHost()
Creates an owning UDP host handle.
Definition Network.cpp:151
int getTimeout() const
Returns the timeout.
Definition Network.cpp:174
~Network() override
Network.
Definition Network.cpp:71
friend class TcpSocket
Definition Network.h:112
friend class Channel
Definition Network.h:108
NetRpc * newRpc(UdpLink *link)
Script adapter: creates a caller-owned RPC client.
Definition Network.cpp:149
void pump()
Drains worker completions and emits them as events; call per frame.
Definition Network.cpp:428
NetWriter * newWriter()
Script adapter: creates a caller-owned writer.
Definition Network.cpp:129
UdpSocket * newUdp()
Script adapter: creates a caller-owned UDP socket.
Definition Network.cpp:80
friend class UdpSocket
Definition Network.h:113
TcpSocket * newTcp()
Script adapter: creates a caller-owned TCP socket.
Definition Network.cpp:77
std::unique_ptr< UdpLink > makeUdpLink(UdpSocket *socket)
Creates an owning UDP link handle, or null for a null socket.
Definition Network.cpp:138
NetTelemetrySnapshot telemetrySnapshot() const
Return a thread-safe copied telemetry snapshot for diagnostics tools.
Definition Network.cpp:201
Channel * newChannel(TcpSocket *socket)
Script adapter: creates a caller-owned channel.
Definition Network.cpp:123
Named collection of Channels (a "session"). Lookup by name, close all. Does not own the channels; the...
Definition Session.h:16
void remove(std::string name)
Removes a channel from the session (does not delete it).
Definition Session.cpp:20
void add(std::string name, Channel *ch)
Registers a channel under a name (replaces an existing entry).
Definition Session.cpp:11
void closeAll()
Closes every channel in the session.
Definition Session.cpp:24
Channel * get(std::string name)
Finds a channel by name, or nullptr.
Definition Session.cpp:15
TCP socket backed by Poco::Net; supports both client (connect) and server (listen/accept) roles....
Definition TcpSocket.h:30
bool listen(uint16_t port)
Starts listening on the given port; true on success.
bool sendString(std::string s)
Sends a string payload.
void close()
Closes the socket.
std::string getPeer() const
Remote peer address string, e.g. "1.2.3.4:5678".
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.
UDP socket backed by Poco::Net; supports connect/bind and datagram send.
Definition UdpSocket.h:26
bool sendTo(eve::data::ByteData *data, std::string host, uint16_t port)
Sends a datagram to an explicit host:port.
Definition UdpSocket.cpp:70
bool sendToString(std::string s, std::string host, uint16_t port)
Sends a string datagram to an explicit host:port.
Definition UdpSocket.cpp:98
void close()
Closes the socket.
bool bind(uint16_t port)
Binds to a local port; true on success.
Definition UdpSocket.cpp:22
bool sendString(std::string s)
Sends a string datagram to the connected peer.
bool connect(std::string host, uint16_t port)
Connects to a remote host:port (restricts send()).
Definition UdpSocket.cpp:47
bool send(eve::data::ByteData *data)
Sends a datagram to the connected peer.
A named event carrying an ordered list of Variant payloads. Pushed messages are heap-allocated; the q...
std::unordered_map< std::string, SkillDefinition > & table()
Definition Skill.cpp:65
void pushValue(HSQUIRRELVM vm, const eve::rx::Value &v)
Definition Rx.cpp:91
SettlementPipeline::Stage fn
One asynchronous network result/event. handle points at the originating TcpSocket/UdpSocket/HttpReque...
Definition NetTypes.h:33
Copied aggregate counters for editor/profiler network inspection.
Definition NetTypes.h:44