add server2

This commit is contained in:
sfja 2026-04-09 13:34:21 +02:00
parent aaa412952d
commit 68578db711
5 changed files with 240 additions and 36 deletions

View File

@ -1,12 +1,9 @@
#include "event_loop.hpp" #include "server2.hpp"
// #include "event_loop.hpp"
#include "json.hpp" #include "json.hpp"
#include "mqtt.hpp" #include "mqtt.hpp"
#include "server.hpp" // #include "server.hpp"
#include <chrono>
#include <cstdio>
#include <iostream>
#include <print> #include <print>
#include <span>
#include <stdexcept> #include <stdexcept>
#include <string_view> #include <string_view>
#include <sys/select.h> #include <sys/select.h>
@ -36,23 +33,23 @@ int main(void)
mqtt_client.subscribe("/skateboard/update", [&](std::string_view text) { mqtt_client.subscribe("/skateboard/update", [&](std::string_view text) {
// //
std::println("Skateboard: {}", text); // std::println("Skateboard: {}", text);
auto result = mst::json::parse(text); // auto result = mst::json::parse(text);
//
if (!result) { // if (!result) {
std::println(stderr, // std::println(stderr,
"error: {} at {}", // "error: {} at {}",
result.error().message, // result.error().message,
result.error().loc.idx); // result.error().loc.idx);
return; // return;
} // }
try { // try {
auto parsed = std::move(result.value()); // auto parsed = std::move(result.value());
std::println(".rotation = {}", // std::println(".rotation = {}",
parsed->query(".rotation").value()->get_f64()); // parsed->query(".rotation").value()->get_f64());
} catch (std::runtime_error& ex) { // } catch (std::runtime_error& ex) {
std::println(stderr, "exception: {}", ex.what()); // std::println(stderr, "exception: {}", ex.what());
} // }
}); });
auto mqtt_thread = std::thread([&]() { auto mqtt_thread = std::thread([&]() {
@ -64,16 +61,19 @@ int main(void)
} }
}); });
auto mgr = mst::event::Manager::create().value(); // auto mgr = mst::event::Manager::create().value();
auto x = mst::Server::bind(mgr, "0.0.0.0", PORT); // auto x = mst::Server::bind(mgr, "0.0.0.0", PORT);
if (!x) { // if (!x) {
std::println("{}", x.error()); // std::println("{}", x.error());
return 1; // return 1;
} // }
std::println("starting"); // std::println("starting");
{ // {
auto x = mgr.start(); // auto x = mgr.start();
} // }
auto server = mst::server2::Server();
server.listen();
mqtt_thread.join(); mqtt_thread.join();
return 0; return 0;

View File

@ -218,7 +218,7 @@ void Client::cb_publish()
void Client::cb_message(std::string_view topic, const void* data, size_t size) void Client::cb_message(std::string_view topic, const void* data, size_t size)
{ {
std::println("[MQTT] Message received"); // std::println("[MQTT] Message received");
auto text = std::string_view(static_cast<const char*>(data), size); auto text = std::string_view(static_cast<const char*>(data), size);

175
backend/src/server2.cpp Normal file
View File

@ -0,0 +1,175 @@
#include "server2.hpp"
#include <algorithm>
#include <cerrno>
#include <cstdint>
#include <cstring>
#include <format>
#include <iterator>
#include <netdb.h>
#include <poll.h>
#include <print>
#include <sys/poll.h>
#include <sys/socket.h>
#include <sys/types.h>
#include <unistd.h>
#include <vector>
namespace {
using namespace mst::server2;
auto get_listener_socket() -> int
{
int listener; // Listening socket descriptor
int status;
struct addrinfo hints = { };
hints.ai_family = AF_INET;
hints.ai_socktype = SOCK_STREAM;
hints.ai_flags = AI_PASSIVE;
struct addrinfo* addr;
if ((status = ::getaddrinfo(NULL, "8881", &hints, &addr)) != 0)
throw Error(std::format("getaddrinfo ({})", ::gai_strerror(status)));
struct addrinfo* p;
for (p = addr; p != NULL; p = p->ai_next) {
listener = ::socket(p->ai_family, p->ai_socktype, p->ai_protocol);
if (listener < 0) {
continue;
}
int reuseaddr_opt = 1;
::setsockopt(
listener, SOL_SOCKET, SO_REUSEADDR, &reuseaddr_opt, sizeof(int));
if (::bind(listener, p->ai_addr, p->ai_addrlen) < 0) {
::close(listener);
continue;
}
break;
}
if (p == NULL)
throw Error("didn't get bound");
::freeaddrinfo(addr);
if (::listen(listener, 10) == -1)
throw Error(std::format("could not listen ({})", strerror(errno)));
return listener;
}
}
namespace mst::server2 {
struct Server::State {
std::vector<::pollfd> pollfds;
std::vector<::pollfd> queued_insertions;
std::vector<std::size_t> queued_deletions;
};
Server::Server()
: m_state(std::make_unique<State>())
{
}
Server::~Server() = default;
void Server::listen()
{
m_listener_fd = get_listener_socket();
m_state->pollfds.push_back(::pollfd {
.fd = m_listener_fd,
.events = POLLIN,
.revents = { },
});
std::println("[mst::server2] listening for connections");
while (true) {
int poll_count
= ::poll(m_state->pollfds.data(), m_state->pollfds.size(), -1);
if (poll_count == -1)
throw Error(std::format("poll (%s)", strerror(errno)));
for (size_t i = 0; i < m_state->pollfds.size(); ++i) {
auto& fd = m_state->pollfds[i];
if (!(fd.revents & (POLLIN | POLLHUP)))
continue;
if (fd.fd == m_listener_fd) {
create_connection();
} else {
try {
handle_request(i);
} catch (Error& ex) {
std::println(stderr,
"[mst::server2] exception in handler for client {}: {}",
fd.fd,
ex.what());
m_state->queued_deletions.push_back(i);
}
}
}
auto& fds = m_state->pollfds;
auto& deletions = m_state->queued_deletions;
std::reverse(deletions.begin(), deletions.end());
for (auto idx : deletions) {
fds.erase(std::next(fds.begin(), static_cast<long>(idx)));
}
m_state->queued_deletions.clear();
for (auto& fd : m_state->queued_insertions) {
m_state->pollfds.push_back(fd);
}
m_state->queued_insertions.clear();
}
}
void Server::create_connection()
{
struct sockaddr_storage remoteaddr;
socklen_t addrlen = sizeof remoteaddr;
int client_fd
= ::accept(m_listener_fd, (struct sockaddr*)&remoteaddr, &addrlen);
if (client_fd == -1)
throw Error(std::format("format ({})", strerror(errno)));
std::println("[mst::server2] client {} connected", client_fd);
m_state->queued_insertions.push_back(::pollfd {
.fd = client_fd,
.events = POLLIN,
.revents = { },
});
}
void Server::handle_request(size_t i)
{
auto& client_fd = m_state->pollfds[i].fd;
auto buffer = std::vector<char>(512);
ssize_t byte_count = ::recv(client_fd, buffer.data(), buffer.size(), 0);
if (byte_count <= 0)
throw Error(std::format("recv: {}", strerror(errno)));
if (byte_count == 0) {
std::println("[mst::server2] client {} disconnected", client_fd);
m_state->queued_deletions.push_back(i);
return;
}
std::println("[mst::server2] received: {:s}", buffer);
}
}

29
backend/src/server2.hpp Normal file
View File

@ -0,0 +1,29 @@
#pragma once
#include <memory>
#include <stdexcept>
namespace mst::server2 {
struct Error : public std::runtime_error {
using std::runtime_error::runtime_error;
};
class Server {
public:
explicit Server();
~Server();
void listen();
private:
struct State;
void create_connection();
void handle_request(size_t i);
int m_listener_fd { };
std::unique_ptr<State> m_state { nullptr };
};
}

View File

@ -42,8 +42,8 @@ auto TcpListener::bind(const std::string& host, uint16_t port)
struct sockaddr_in address = { struct sockaddr_in address = {
.sin_family = AF_INET, .sin_family = AF_INET,
.sin_port = htons(port), .sin_port = ::htons(port),
.sin_addr = in_addr { .s_addr = inet_addr(host.c_str()) }, .sin_addr = in_addr { .s_addr = ::inet_addr(host.c_str()) },
.sin_zero = { }, .sin_zero = { },
}; };