This commit is contained in:
Bart Beumer 2025-11-06 19:50:29 +01:00
parent 7c2ac8456e
commit aeb0d3c445
6 changed files with 61 additions and 81 deletions

View File

@ -14,6 +14,7 @@ RUN apk update && \
libstdc++ \
libtool \
linux-headers \
ninja \
m4 \
perl \
python3 \

View File

@ -126,7 +126,7 @@ void fractal(
void handle_request_http(const address_type& addr, const request_type& req, bmrshared::web::response_promise promise) override
{
std::string_view target = req.target();
std::string target = req.target();
std::cout << "HTTP GET: " << target << std::endl;
if (req.method() != boost::beast::http::verb::get)
{
@ -150,8 +150,9 @@ void fractal(
double offset_x = std::stoi(base_match[2]);
double z = std::stoi(base_match[3]);
auto renderfn = [offset_y, offset_x, z, req, promise]() mutable
auto renderfn = [offset_y, offset_x, z, req, target, promise]() mutable
{
std::cout << "HTTP inside render function " << target << std::endl;
constexpr int pixel_width = 256;
constexpr int pixel_height = 256;
constexpr int max_iterations = 100;
@ -195,6 +196,8 @@ void fractal(
ok.keep_alive(req.keep_alive());
ok.prepare_payload();
promise.SendResponse(std::move(ok));
std::cout << " DONE HTTP inside render function " << target << std::endl;
};
m_ioc.post(renderfn);

View File

@ -27,11 +27,12 @@ class response_promise final
template<typename Body, typename Fields>
void SendResponse(boost::beast::http::response<Body, Fields> response)
{
response_sender responder = [&response](boost::beast::tcp_stream& stream)
response_sender responder = [r = std::move(response)](boost::beast::tcp_stream& stream)
{
boost::beast::http::write(stream, response);
boost::beast::http::write(stream, r);
};
m_call_on_response(responder);
m_call_on_response(std::move(responder));
}
private:

View File

@ -152,7 +152,7 @@ void directory_request_handler::handle_request_http(const address_type& address,
head.set(http::field::cache_control, "public, max-age=360");
head.set(http::field::content_length, std::to_string(file_size));
head.keep_alive(req.keep_alive());
promise.SendResponse(std::move(head));
//promise.SendResponse(std::move(head));
return;
}
}

View File

@ -73,6 +73,24 @@ void internal_server::do_accept(boost::system::error_code ec, boost::asio::ip::t
void internal_server::session_process_queued()
{
// Send out all processed responses.
for (auto& session : m_sessions)
{
if (session)
{
http_session& s = *session;
auto& our = s.m_outstanding_requests;
while (!our.empty() && our.front() && our.front()->m_response_sender)
{
auto& resp = our.front();
(*resp->m_response_sender)(s.stream);
our.pop();
}
}
}
// Determine how many outstanding requests we have.
std::size_t count_responding = 0;
for(const auto& session : m_sessions)
@ -85,78 +103,14 @@ void internal_server::session_process_queued()
{
for(auto& session : m_sessions)
{
if (count_responding > m_max_simultaneous_requests)
if (!session || (count_responding > m_max_simultaneous_requests))
{
break;
// do nothing.
}
else if (session && (session->m_state == session_state::queue_reading))
else if (session->m_state == session_state::queue_reading)
{
http_session& s = *session;
s.state = session_state::reading;
s.parser.emplace(); // Create empty parser object.
s.parser->body_limit(m_request_body_limit);
s.stream.expires_after(m_incoming_request_timeout);
auto handler = [weak_server = weak_from_this(), session](boost::beast::error_code ec, std::size_t size)
{
auto shared_server = weak_server.lock();
if (shared_server)
{
shared_server->session_on_read(session, ec, size);
}
};
boost::beast::http::async_read(s.stream, s.buffer, *s.parser, handler);
}
}
}
// First empty out all sessions queued for getting a response as much as we are allowed.
auto iter = m_sessions.begin();
while ((count_responding < m_max_simultaneous_requests) && (iter != m_sessions.end()))
{
iter = std::find_if(
iter,
m_sessions.end(),
[](std::shared_ptr<http_session> session)
{
return session && (session->state == session_state::queue_responding);
});
if (iter != m_sessions.end())
{
auto& session = *iter;
session->state = session_state::responding;
++count_responding;
// Regular http request. Pass to request handler.
auto finalize_response = [session, weak_server = weak_from_this()](const response_promise::response_sender& rs)
{
auto shared_server = weak_server.lock();
if (shared_server)
{
rs(session->stream);
// When we receive the response and have sent it, we proceed with reading
shared_server->session_finalize_response(session);
}
};
const auto& address = session->stream.socket().remote_endpoint().address();
m_handler.handle_request_http(address, session->parser->get(), response_promise(finalize_response));
}
}
if (count_responding < m_max_simultaneous_requests)
{
// Now we can dispatch reading.
for (const auto& session : m_sessions)
{
if (session && (session->state == session_state::queue_reading))
{
http_session& s = *session;
s.state = session_state::reading;
s.m_state = session_state::reading;
s.parser.emplace(); // Create empty parser object.
s.parser->body_limit(m_request_body_limit);
s.stream.expires_after(m_incoming_request_timeout);
@ -202,8 +156,27 @@ void internal_server::session_on_read(std::shared_ptr<http_session> session,
}
else
{
// Just a regular http request. Queue for handling later.
s.state = session_state::queue_responding;
std::shared_ptr<http_request> req = s.m_outstanding_requests.emplace(std::make_shared<http_request>());
// A regular HTTP request.
auto finalize_response = [session, weak_server = weak_from_this(), req](const response_promise::response_sender& rs)
{
// We have a response, add it to our spot in the queue.
auto shared_server = weak_server.lock();
if (shared_server)
{
// We have a response, add it to our spot in the queue.
req->m_response_sender = rs;
// Queue processing of the queue so we can continue.
shared_server->m_io_context.post([shared_server]{shared_server->session_process_queued();});
}
};
const auto& address = session->stream.socket().remote_endpoint().address();
m_handler.handle_request_http(address, session->parser->get(), response_promise(finalize_response));
s.m_state = session_state::queue_reading;
session_process_queued();
}
}
@ -212,7 +185,7 @@ void internal_server::session_finalize_response(std::shared_ptr<http_session> se
{
if(session)
{
session->state = session_state::queue_reading;
session->m_state = session_state::queue_reading;
session_process_queued();
}
}

View File

@ -21,6 +21,11 @@ namespace detail {
reading
};
struct http_request
{
std::optional<response_promise::response_sender> m_response_sender;
};
struct http_session {
explicit http_session(boost::asio::ip::tcp::socket s)
: stream(std::move(s))
@ -29,12 +34,9 @@ namespace detail {
boost::beast::flat_buffer buffer;
boost::optional<boost::beast::http::request_parser<boost::beast::http::string_body>> parser;
session_state m_state;
std::queue<> m_outstanding_requests;
std::queue<std::shared_ptr<http_request>> m_outstanding_requests;
};
class internal_server final : public std::enable_shared_from_this<internal_server> {
public:
using address_type = bmrshared::web::server::address_type;