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>
22#include <simplesquirrel/simplesquirrel.hpp>
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;
37 SQInteger
top = sq_gettop(raw);
38 sq_pushobject(raw,
f.getRaw());
39 sq_pushroottable(raw);
42 sq_call(raw, 3, SQFalse, SQTrue);
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;
51 SQInteger
top = sq_gettop(raw);
52 sq_pushobject(raw,
f.getRaw());
53 sq_pushroottable(raw);
57 sq_call(raw, 4, SQFalse, SQTrue);
66 worker_ = std::make_unique<NetWorker>(
this);
68 eve::cap::provide<eve::service::INetwork>(
this);
72 if (worker_) worker_->stop();
76std::unique_ptr<TcpSocket>
Network::makeTcp() {
return std::make_unique<TcpSocket>(
this); }
79std::unique_ptr<UdpSocket>
Network::makeUdp() {
return std::make_unique<UdpSocket>(
this); }
83 return std::make_unique<HttpRequest>(
this, std::move(
method), std::move(url));
90 const std::string&
body,
int timeoutMs,
int&
status,
91 std::string& responseBody) {
95 req->
setTimeout(timeoutMs > 0 ? timeoutMs : 10000);
96 if (!req->
submit())
return false;
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;
107 if (
c.bytes) responseBody.assign(
c.bytes->begin(),
c.bytes->end());
112 std::this_thread::sleep_for(std::chrono::milliseconds(10));
118 if (!socket)
return nullptr;
119 auto ch = std::make_unique<Channel>(socket);
120 bindChannel(socket, ch.get());
132 auto r = std::make_unique<NetReader>();
133 r->setBytes(std::move(
bytes));
139 if (!socket)
return nullptr;
140 auto link = std::make_unique<UdpLink>(
this, socket);
141 bindUdpLink(socket, link.get());
147 return link ? std::make_unique<NetRpc>(link) :
nullptr;
155 if (sock) udpLinks_[sock] = link;
158void Network::unbindUdpLink(UdpSocket* sock) {
159 if (sock) udpLinks_.erase(sock);
162void Network::bindUdpHost(UdpSocket* sock, NetHost*
host) {
163 if (sock) udpHosts_[sock] =
host;
166void Network::unbindUdpHost(UdpSocket* sock) {
167 if (sock) udpHosts_.erase(sock);
188 ++telemetryRevision_;
190 receivedBytes_ +=
c.bytes->size();
193 if (worker_) worker_->post(std::move(
c));
196void Network::recordSent(
size_t bytes) {
198 ++telemetryRevision_;
203 result.
revision = telemetryRevision_.load();
207 result.
errors = errors_.load();
210 std::lock_guard<std::mutex> lock(watchMu_);
213 for (
const auto* socket : watchedTcp_)
if (socket) result.
queuedTcpBytes += socket->pendingSendBytes();
216 std::lock_guard<std::mutex> lock(channelMu_);
223 sentBytes_ = receivedBytes_ = completions_ = errors_ = connections_ = 0;
224 ++telemetryRevision_;
227void Network::drainCompletions(std::vector<NetCompletion>& out) {
228 if (worker_) worker_->drain(out);
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_;
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_;
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_;
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_;
261void Network::bindChannel(TcpSocket* sock, Channel* ch) {
262 std::lock_guard<std::mutex> lock(channelMu_);
263 channels_[sock] = ch;
264 ++telemetryRevision_;
267void Network::unbindChannel(TcpSocket* sock) {
268 std::lock_guard<std::mutex> lock(channelMu_);
269 if(channels_.erase(sock))++telemetryRevision_;
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;
278void Network::pollSockets() {
279 std::vector<TcpSocket*> tcpCopy;
280 std::vector<UdpSocket*> udpCopy;
282 std::lock_guard<std::mutex> lock(watchMu_);
283 tcpCopy = watchedTcp_;
284 udpCopy = watchedUdp_;
288 if (sock->isListening() && sock->server()) {
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));
296 sock->pushAccepted(std::move(peer));
303 c.peer = raw->getPeer();
308 if (sock->isConnected() && sock->stream()) {
312 int n = sock->stream()->receiveBytes(buf,
sizeof(buf));
318 c.bytes = std::make_shared<std::vector<char>>(buf, buf +
n);
319 c.peer = sock->getPeer();
330 }
catch (
const Poco::Exception&) {
336 if (!sock || !sock->datagram())
continue;
337 for (
int i = 0; i < 64; ++i) {
340 Poco::Net::SocketAddress sender;
341 int n = sock->datagram()->receiveFrom(buf,
sizeof(buf), sender);
347 c.bytes = std::make_shared<std::vector<char>>(buf, buf +
n);
348 c.peer = sender.toString();
350 }
catch (
const Poco::Exception&) {
359void Network::emitCompletion(
const NetCompletion&
c) {
360 auto* ev = eve::ModuleManager::getInstance<eve::platform_event::PlatformEvent>(
"PlatformEvent");
366 std::vector<Variant> args;
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));
380 if (
c.bytes && !
c.bytes->empty())
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));
391 args.push_back(Variant::makePtr(ch));
392 args.push_back(Variant::makeString(
c.reason));
393 ev->push(std::make_unique<Message>(
"chclose", args));
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));
403 if (
c.bytes && !
c.bytes->empty())
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));
413 if (
c.bytes && !
c.bytes->empty())
415 args.push_back(Variant::makePtr(
c.handle));
416 args.push_back(Variant::makeOwnedPtr(bd));
417 ev->push(std::make_unique<Message>(
"chmsg", args));
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));
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())
434 for (
auto&
c : batch) {
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);
442 auto lit = udpLinks_.find(sock);
443 if (lit != udpLinks_.end()) {
444 if (
c.bytes) lit->second->onDatagram(*
c.bytes,
c.peer);
450 for (
auto& kv : udpHosts_) kv.second->pump(now);
451 for (
auto& kv : udpLinks_) kv.second->pump(now);
454void Network::expose(ssq::Table& table) {
455 auto cls = table.addClass(
name, Network::create,
false);
459 "TcpSocket", std::function<TcpSocket*()>([]() {
return new TcpSocket(Network::create()); }),
true);
470 "UdpSocket", std::function<UdpSocket*()>([]() {
return new UdpSocket(Network::create()); }),
true);
481 std::function<HttpRequest*()>(
482 []() {
return new HttpRequest(Network::create(),
"GET",
"http://127.0.0.1/"); }),
492 "Channel", std::function<Channel*()>([]() {
return new Channel(
nullptr); }),
true);
497 auto sess =
table.addClass<Session>(
498 "Session", std::function<Session*()>([]() {
return new Session(); }),
true);
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)); });
519 if (
d)
w->writeBytes(
d->getData(),
d->getSize());
521 writer.addFunc(
"toString", [](NetWriter*
w) {
return w->toString(); });
522 writer.addFunc(
"size", [](NetWriter*
w) {
return static_cast<int64_t
>(
w->size()); });
524 auto reader =
table.addClass<NetReader>(
525 "NetReader", std::function<NetReader*()>([]() {
return new NetReader(); }),
true);
527 return d ?
r->init(
d->getData(),
d->getSize()) : false;
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()); });
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());
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(); });
550 auto link =
table.addClass<UdpLink>(
552 std::function<UdpLink*()>([]() {
return new UdpLink(
nullptr,
nullptr); }),
true);
555 link.addFunc(
"sendReliable", [](UdpLink* l, int64_t ch,
const std::string&
s) {
558 link.addFunc(
"sendUnreliable", [](UdpLink* l, int64_t ch,
const std::string&
s) {
561 link.addFunc(
"sendOrdered", [](UdpLink* l, int64_t ch,
const std::string&
s) {
564 link.addFunc(
"onMessage", [](UdpLink* l, ssq::Object
fn) {
566 l->setMessageHandler(
568 callScript2(
fn, ch, std::string(
d,
n));
573 link.addFunc(
"isAlive", [](UdpLink* l) {
return l && l->isAlive(); });
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()); });
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,
585 r->callString(
static_cast<uint16_t
>(msgId),
payload, reliable);
587 rpc.addFunc(
"registerRpc", [](NetRpc*
r, int64_t msgId, ssq::Object
fn) {
588 r->registerScript(
static_cast<uint16_t
>(msgId),
fn);
592 "NetHost", std::function<NetHost*()>([]() {
return new NetHost(
nullptr); }),
true);
596 h->setMessageHandler(
598 callScript3(
fn, peerId, ch, std::string(
d,
n));
603 h->setPeerConnectedHandler([
fn](uint32_t
id) { callScript2(
fn,
id,
""); });
605 host.addFunc(
"onPeerDisconnected", [](
NetHost*
h, ssq::Object
fn) {
607 h->setPeerDisconnectedHandler([
fn](uint32_t
id) { callScript2(
fn,
id,
""); });
609 host.addFunc(
"sendReliable", [](
NetHost*
h, int64_t peerId, int64_t ch,
610 const std::string&
s) {
612 static_cast<uint8_t
>(ch),
s);
614 host.addFunc(
"sendUnreliable", [](
NetHost*
h, int64_t peerId, int64_t ch,
615 const std::string&
s) {
617 static_cast<uint8_t
>(ch),
s);
619 host.addFunc(
"sendOrdered", [](
NetHost*
h, int64_t peerId, int64_t ch,
620 const std::string&
s) {
622 static_cast<uint8_t
>(ch),
s);
625 host.addFunc(
"peerCount", [](
NetHost*
h) {
return static_cast<int64_t
>(
h->peerCount()); });
630void Network::expose(ssq::Class&
cls) {
struct SQVM * HSQUIRRELVM
wgpu::PopErrorScopeStatus status
#define Module_IMPL(ModuleName, newExpr)
virtual std::string getName() const =0
Returns the name.
In-memory byte buffer implementing eve::Data (ref-counted).
Length-prefixed (big-endian uint32) message framing over a TcpSocket. sendMsg() writes one framed mes...
bool sendMsg(eve::data::ByteData *data)
Sends one framed message.
TcpSocket * getSocket() const
The underlying TCP socket, or nullptr.
bool sendMsgString(std::string s)
Sends a string payload as one framed message (empty strings are rejected).
Asynchronous HTTP request (Poco-based). Configure headers/body, call submit(), then receive the respo...
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.
void setLossRate(float rate)
Sets the loss rate.
bool start(uint16_t port)
Starts .
void setTimeoutMs(int ms)
Sets the timeout ms.
UdpLink * linkByPeerId(uint32_t peerId) const
Link by peer id.
EVENGINE_API_PLATFORM public API.
EVENGINE_API_PLATFORM public API.
EVENGINE_API_PLATFORM public API.
void writeBool(bool v)
Writes bool.
void writeF32(float v)
Writes f 32.
void writeString(const std::string &s)
Writes string.
void writeF64(double v)
Writes f 64.
Network module: TCP/UDP/HTTP factories, background worker, and completion event plumbing....
NetHost * newHost()
Script adapter: creates a caller-owned UDP host.
std::unique_ptr< TcpSocket > makeTcp()
Creates an owning TCP socket handle for C++ callers.
NetReader * newReader(std::string bytes)
Script adapter: creates a caller-owned reader.
std::unique_ptr< NetRpc > makeRpc(UdpLink *link)
Creates an owning RPC handle, or null for a null link.
bool getVerifySsl() const
Returns the verify ssl.
std::unique_ptr< NetReader > makeReader(std::string bytes)
Creates an owning streaming reader handle.
void setTimeout(int ms)
Default socket/HTTP timeout in milliseconds.
Session * newSession()
Script adapter: creates a caller-owned session.
void resetTelemetry()
Reset cumulative telemetry counters without affecting live sockets.
std::unique_ptr< Session > makeSession()
Creates an owning session handle.
UdpLink * newUdpLink(UdpSocket *socket)
Script adapter: creates a caller-owned UDP link.
std::unique_ptr< NetWriter > makeWriter()
Creates an owning streaming writer handle.
std::unique_ptr< UdpSocket > makeUdp()
Creates an owning UDP socket handle for C++ callers.
void setVerifySsl(bool verify)
Whether TLS peer certificates are verified (default true).
std::unique_ptr< Channel > makeChannel(TcpSocket *socket)
Creates an owning channel handle, or null for a null socket.
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).
HttpRequest * newHttp(std::string method, std::string url)
Script adapter: creates a caller-owned HTTP request.
std::unique_ptr< HttpRequest > makeHttp(std::string method, std::string url)
Creates an owning HTTP request handle for C++ callers.
std::unique_ptr< NetHost > makeHost()
Creates an owning UDP host handle.
int getTimeout() const
Returns the timeout.
~Network() override
Network.
NetRpc * newRpc(UdpLink *link)
Script adapter: creates a caller-owned RPC client.
void pump()
Drains worker completions and emits them as events; call per frame.
NetWriter * newWriter()
Script adapter: creates a caller-owned writer.
UdpSocket * newUdp()
Script adapter: creates a caller-owned UDP socket.
TcpSocket * newTcp()
Script adapter: creates a caller-owned TCP socket.
std::unique_ptr< UdpLink > makeUdpLink(UdpSocket *socket)
Creates an owning UDP link handle, or null for a null socket.
NetTelemetrySnapshot telemetrySnapshot() const
Return a thread-safe copied telemetry snapshot for diagnostics tools.
Channel * newChannel(TcpSocket *socket)
Script adapter: creates a caller-owned channel.
Named collection of Channels (a "session"). Lookup by name, close all. Does not own the channels; the...
void remove(std::string name)
Removes a channel from the session (does not delete it).
void add(std::string name, Channel *ch)
Registers a channel under a name (replaces an existing entry).
void closeAll()
Closes every channel in the session.
Channel * get(std::string name)
Finds a channel by name, or nullptr.
TCP socket backed by Poco::Net; supports both client (connect) and server (listen/accept) roles....
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.
bool send(eve::data::ByteData *data)
Sends a framed byte payload; true when queued/accepted.
EVENGINE_API_PLATFORM public API.
bool setRemoteString(const std::string &addr)
bool setRemote(std::string host, uint16_t port)
MsgType
MsgType public API.
void setTimeoutMs(int ms)
void setLossRate(float rate)
UDP socket backed by Poco::Net; supports connect/bind and datagram send.
bool sendTo(eve::data::ByteData *data, std::string host, uint16_t port)
Sends a datagram to an explicit host:port.
bool sendToString(std::string s, std::string host, uint16_t port)
Sends a string datagram to an explicit host:port.
void close()
Closes the socket.
bool bind(uint16_t port)
Binds to a local port; true on success.
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()).
bool send(eve::data::ByteData *data)
Sends a datagram to the connected peer.
std::unordered_map< std::string, SkillDefinition > & table()
void pushValue(HSQUIRRELVM vm, const eve::rx::Value &v)
SettlementPipeline::Stage fn
One asynchronous network result/event. handle points at the originating TcpSocket/UdpSocket/HttpReque...
Copied aggregate counters for editor/profiler network inspection.