diff --git a/src/commit_pipeline.cpp b/src/commit_pipeline.cpp index c551d08..29293c2 100644 --- a/src/commit_pipeline.cpp +++ b/src/commit_pipeline.cpp @@ -335,7 +335,7 @@ void CommitPipeline::run_release_stage(int thread_index) { // Send the JSON response using protocol-agnostic interface // HTTP formatting will happen in on_preprocess_writes() - conn_ref->send_response(commit_entry.protocol_context, + conn_ref->send_response(commit_entry.handle, commit_entry.response_json, std::move(commit_entry.request_arena)); } else if constexpr (std::is_same_v) { @@ -349,7 +349,7 @@ void CommitPipeline::run_release_stage(int thread_index) { // Send the JSON response using protocol-agnostic interface // HTTP formatting will happen in on_preprocess_writes() - conn_ref->send_response(status_entry.protocol_context, + conn_ref->send_response(status_entry.handle, status_entry.response_json, std::move(status_entry.request_arena)); } else if constexpr (std::is_same_v) { @@ -364,8 +364,7 @@ void CommitPipeline::run_release_stage(int thread_index) { // Send the response using protocol-agnostic interface // HTTP formatting will happen in on_preprocess_writes() conn_ref->send_response( - health_check_entry.protocol_context, - health_check_entry.response_json, + health_check_entry.handle, health_check_entry.response_json, std::move(health_check_entry.request_arena)); } else if constexpr (std::is_same_v) { auto &get_version_entry = e; @@ -378,8 +377,7 @@ void CommitPipeline::run_release_stage(int thread_index) { // Send the response using protocol-agnostic interface // HTTP formatting will happen in on_preprocess_writes() conn_ref->send_response( - get_version_entry.protocol_context, - get_version_entry.response_json, + get_version_entry.handle, get_version_entry.response_json, std::move(get_version_entry.request_arena)); } }, diff --git a/src/connection.cpp b/src/connection.cpp index 0652965..fe092e0 100644 --- a/src/connection.cpp +++ b/src/connection.cpp @@ -125,7 +125,7 @@ void Connection::append_bytes(std::span data_parts, } // May be called from a foreign thread! -void Connection::send_response(void *protocol_context, +void Connection::send_response(ProtocolHandle handle, std::string_view response_json, Arena arena) { std::unique_lock lock(mutex_); @@ -136,7 +136,7 @@ void Connection::send_response(void *protocol_context, // Store response in queue for protocol handler processing pending_response_queue_.emplace_back( - PendingResponse{protocol_context, response_json, std::move(arena)}); + PendingResponse{handle, response_json, std::move(arena)}); // Trigger epoll interest if this is the first pending response if (pending_response_queue_.size() == 1) { diff --git a/src/connection.hpp b/src/connection.hpp index 7d458ba..81903ff 100644 --- a/src/connection.hpp +++ b/src/connection.hpp @@ -40,25 +40,24 @@ enum class ConnectionShutdown { */ struct MessageSender { /** - * @brief Send response with protocol-specific context for ordering. + * @brief Send response with protocol-specific handle for correlation. * * Thread-safe method for pipeline threads to send responses back to clients. * Delegates to the connection's protocol handler for ordering logic. * The protocol handler may queue the response or send it immediately. * - * @param protocol_context Arena-allocated protocol-specific context - * @param data Response data parts (may be empty for deferred serialization) + * @param handle Protocol-specific handle for correlating this response + * @param response_json JSON response body (may be empty for deferred + * serialization) * @param arena Arena containing response data and context * * Example usage: * ```cpp - * auto* ctx = arena.allocate(); - * ctx->sequence_id = 42; - * auto response_data = format_response(arena); - * conn.send_response(ctx, response_data, std::move(arena)); + * ProtocolHandle handle = handler.allocate_response_context(arena); + * conn.send_response(handle, response_json, std::move(arena)); * ``` */ - virtual void send_response(void *protocol_context, + virtual void send_response(ProtocolHandle handle, std::string_view response_json, Arena arena) = 0; virtual ~MessageSender() = default; @@ -141,7 +140,7 @@ struct Connection : MessageSender { append_bytes(std::span data_parts, Arena arena, ConnectionShutdown shutdown_mode = ConnectionShutdown::None); - void send_response(void *protocol_context, std::string_view response_json, + void send_response(ProtocolHandle handle, std::string_view response_json, Arena arena) override; /** @@ -166,7 +165,7 @@ struct Connection : MessageSender { * if (auto conn = weak_conn.lock()) { * Arena arena; * auto response = process_request(request_data, arena); - * conn->send_response(response_context, response_json, std::move(arena)); + * conn->send_response(handle, response_json, std::move(arena)); * } * }); * ``` diff --git a/src/connection_handler.hpp b/src/connection_handler.hpp index 9543f7a..fb2a43a 100644 --- a/src/connection_handler.hpp +++ b/src/connection_handler.hpp @@ -8,13 +8,19 @@ struct Connection; // Include Arena header since PendingResponse uses Arena by value #include "arena.hpp" +#include + +// Opaque handle used to correlate responses with protocol-specific state. +// Each ConnectionHandler implementation defines the meaning of a handle value +// and resolves it back to full context in on_preprocess_writes(). +using ProtocolHandle = int64_t; /** * Represents a response queued by pipeline threads for protocol processing. * Contains JSON response data that can be wrapped by any protocol. */ struct PendingResponse { - void *protocol_context; // Arena-allocated protocol-specific context + ProtocolHandle handle; // Protocol-specific handle for correlation std::string_view response_json; // JSON response body (arena-allocated) Arena arena; // Arena containing response data and context }; diff --git a/src/http_handler.cpp b/src/http_handler.cpp index 11bbf32..320ca44 100644 --- a/src/http_handler.cpp +++ b/src/http_handler.cpp @@ -52,6 +52,59 @@ HttpRequestState::HttpRequestState() current_header_field_buf(ArenaStlAllocator(&arena)), current_header_value_buf(ArenaStlAllocator(&arena)) {} +// HttpConnectionState implementation +HttpConnectionState::~HttpConnectionState() = default; + +void HttpConnectionState::register_response_context( + int64_t sequence_id, std::unique_ptr ctx) { + contexts_[sequence_id] = std::move(ctx); +} + +HttpResponseContext * +HttpConnectionState::resolve_response_context(ProtocolHandle handle) { + auto it = contexts_.find(handle); + if (it == contexts_.end()) { + return nullptr; + } + return it->second.get(); +} + +void HttpConnectionState::send_ordered_response( + Connection &conn, HttpResponseContext *ctx, + std::span http_response, Arena arena) { + assert(ctx); + int64_t sequence_id = ctx->sequence_id; + + std::unique_ptr owned_ctx; + auto it = contexts_.find(sequence_id); + if (it != contexts_.end()) { + owned_ctx = std::move(it->second); + } else { + owned_ctx.reset(ctx); + } + + owned_ctx->ready = true; + owned_ctx->data = http_response; + owned_ctx->arena = std::move(arena); + contexts_[sequence_id] = std::move(owned_ctx); + + // Process ready responses in order and send via append_bytes + auto iter = contexts_.begin(); + while (iter != contexts_.end() && iter->first == next_sequence_to_send) { + auto *context = iter->second.get(); + if (!context || !context->ready) { + break; + } + + // Send through append_bytes which handles write interest + conn.append_bytes(context->data, std::move(context->arena), + context->connection_close ? ConnectionShutdown::WriteOnly + : ConnectionShutdown::None); + next_sequence_to_send++; + iter = contexts_.erase(iter); + } +} + // HttpHandler implementation void HttpHandler::on_connection_established(Connection &conn) { // Allocate HTTP state using server-provided arena for connection lifecycle @@ -72,7 +125,11 @@ void HttpHandler::on_preprocess_writes( // Process incoming responses and add to reorder queue { for (auto &pending : pending_responses) { - auto *ctx = static_cast(pending.protocol_context); + auto *ctx = state->resolve_response_context(pending.handle); + // Handle stale or invalid handles defensively + if (!ctx) { + continue; + } // Determine HTTP status code and content type from response content int status_code = 200; @@ -96,9 +153,8 @@ void HttpHandler::on_preprocess_writes( status_code, content_type, pending.response_json, pending.arena, ctx->http_request_id, ctx->connection_close); - state->send_ordered_response(conn, ctx->sequence_id, http_response, - std::move(pending.arena), - ctx->connection_close); + state->send_ordered_response(conn, ctx, http_response, + std::move(pending.arena)); } } } @@ -116,11 +172,18 @@ void HttpHandler::on_batch_complete(std::span batch) { int64_t sequence_id = state->get_next_sequence_id(); req.sequence_id = sequence_id; - // Create HttpResponseContext for this request - auto *ctx = req.arena.allocate(1); + // Create HttpResponseContext for this request; sequence_id doubles as the + // ProtocolHandle for async response correlation. The context is + // registered in HttpConnectionState so async responses can resolve it by + // handle. + auto ctx = std::make_unique(); ctx->sequence_id = sequence_id; ctx->http_request_id = req.http_request_id; ctx->connection_close = req.connection_close; + HttpResponseContext *ctx_ptr = ctx.get(); + req.response_context = ctx_ptr; + state->register_response_context(sequence_id, std::move(ctx)); + ProtocolHandle handle = sequence_id; RouteMatch route_match; auto parse_result = @@ -132,9 +195,9 @@ void HttpHandler::on_batch_complete(std::span batch) { auto json_response = R"({"error":"Malformed URL encoding"})"; auto http_response = format_json_response(400, json_response, req.arena, 0, true); - state->send_ordered_response(*conn, ctx->sequence_id, http_response, - std::move(req.arena), - ctx->connection_close); + ctx_ptr->connection_close = true; + state->send_ordered_response(*conn, ctx_ptr, http_response, + std::move(req.arena)); break; } req.route = route_match.route; @@ -176,25 +239,25 @@ void HttpHandler::on_batch_complete(std::span batch) { // Create CommitEntry for commit requests if (req.route == HttpRoute::PostCommit && req.commit_request && req.parsing_commit && req.basic_validation_passed) { - g_batch_entries.emplace_back(CommitEntry(conn->get_weak_ref(), ctx, + g_batch_entries.emplace_back(CommitEntry(conn->get_weak_ref(), handle, req.commit_request.get(), std::move(req.arena))); } // Create StatusEntry for status requests else if (req.route == HttpRoute::GetStatus) { - g_batch_entries.emplace_back(StatusEntry(conn->get_weak_ref(), ctx, + g_batch_entries.emplace_back(StatusEntry(conn->get_weak_ref(), handle, req.status_request_id, std::move(req.arena))); } // Create HealthCheckEntry for health check requests else if (req.route == HttpRoute::GetOk) { - g_batch_entries.emplace_back( - HealthCheckEntry(conn->get_weak_ref(), ctx, std::move(req.arena))); + g_batch_entries.emplace_back(HealthCheckEntry( + conn->get_weak_ref(), handle, std::move(req.arena))); } // Create GetVersionEntry for version requests else if (req.route == HttpRoute::GetVersion) { g_batch_entries.emplace_back( - GetVersionEntry(conn->get_weak_ref(), ctx, std::move(req.arena), + GetVersionEntry(conn->get_weak_ref(), handle, std::move(req.arena), commit_pipeline_.get_committed_version())); } } @@ -235,13 +298,17 @@ void HttpHandler::on_data_arrived(std::string_view data, Connection &conn) { break; } // Parse error - send response directly since this is before sequence - // assignment + // assignment. Use a temporary response context to carry metadata. auto json_response = R"({"error":"Bad request"})"; + auto ctx = std::make_unique(); + ctx->sequence_id = state->get_next_sequence_id(); + ctx->http_request_id = 0; + ctx->connection_close = true; auto http_response = - format_json_response(400, json_response, state->pending.arena, 0, true); - state->send_ordered_response(conn, state->get_next_sequence_id(), - http_response, std::move(state->pending.arena), - true); + format_json_response(400, json_response, state->pending.arena, + ctx->http_request_id, ctx->connection_close); + state->send_ordered_response(conn, ctx.get(), http_response, + std::move(state->pending.arena)); return; } } @@ -255,17 +322,18 @@ void HttpHandler::handle_get_version(Connection &, HttpRequestState &) { void HttpHandler::handle_post_commit(Connection &conn, HttpRequestState &state) { commit_counter.inc(); + auto *ctx = state.response_context; + assert(ctx); // Check if streaming parse was successful if (!state.commit_request || !state.parsing_commit) { auto json_response = R"({"error":"Parse failed"})"; auto http_response = format_json_response(400, json_response, state.arena, - state.http_request_id, state.connection_close); + ctx->http_request_id, ctx->connection_close); auto *conn_state = static_cast(conn.user_data); - conn_state->send_ordered_response(conn, state.sequence_id, http_response, - std::move(state.arena), - state.connection_close); + conn_state->send_ordered_response(conn, ctx, http_response, + std::move(state.arena)); return; } @@ -310,12 +378,11 @@ void HttpHandler::handle_post_commit(Connection &conn, static_cast(error_msg.size()), error_msg.data()); auto http_response = format_json_response(400, json_response, state.arena, - state.http_request_id, state.connection_close); + ctx->http_request_id, ctx->connection_close); auto *conn_state = static_cast(conn.user_data); - conn_state->send_ordered_response(conn, state.sequence_id, http_response, - std::move(state.arena), - state.connection_close); + conn_state->send_ordered_response(conn, ctx, http_response, + std::move(state.arena)); return; } @@ -326,22 +393,25 @@ void HttpHandler::handle_post_commit(Connection &conn, void HttpHandler::handle_get_subscribe(Connection &conn, HttpRequestState &state) { + auto *ctx = state.response_context; + assert(ctx); // TODO: Implement subscription streaming auto json_response = R"({"message":"Subscription endpoint - streaming not yet implemented"})"; auto http_response = format_json_response(200, json_response, state.arena, - state.http_request_id, state.connection_close); + ctx->http_request_id, ctx->connection_close); auto *conn_state = static_cast(conn.user_data); - conn_state->send_ordered_response(conn, state.sequence_id, http_response, - std::move(state.arena), - state.connection_close); + conn_state->send_ordered_response(conn, ctx, http_response, + std::move(state.arena)); } void HttpHandler::handle_get_status(Connection &conn, HttpRequestState &state, const RouteMatch &route_match) { status_counter.inc(); + auto *ctx = state.response_context; + assert(ctx); // Status requests are processed through the pipeline // Response will be generated in the sequence stage // This handler extracts request_id from query parameters and prepares for @@ -354,13 +424,12 @@ void HttpHandler::handle_get_status(Connection &conn, HttpRequestState &state, R"({"error":"Missing required query parameter: request_id"})"; auto http_response = format_json_response(400, json_response, state.arena, - state.http_request_id, state.connection_close); + ctx->http_request_id, ctx->connection_close); // Add directly to response queue with proper sequencing auto *conn_state = static_cast(conn.user_data); - conn_state->send_ordered_response(conn, state.sequence_id, http_response, - std::move(state.arena), - state.connection_close); + conn_state->send_ordered_response(conn, ctx, http_response, + std::move(state.arena)); return; } @@ -368,13 +437,12 @@ void HttpHandler::handle_get_status(Connection &conn, HttpRequestState &state, auto json_response = R"({"error":"Empty request_id parameter"})"; auto http_response = format_json_response(400, json_response, state.arena, - state.http_request_id, state.connection_close); + ctx->http_request_id, ctx->connection_close); // Add directly to response queue with proper sequencing auto *conn_state = static_cast(conn.user_data); - conn_state->send_ordered_response(conn, state.sequence_id, http_response, - std::move(state.arena), - state.connection_close); + conn_state->send_ordered_response(conn, ctx, http_response, + std::move(state.arena)); return; } @@ -387,53 +455,58 @@ void HttpHandler::handle_get_status(Connection &conn, HttpRequestState &state, void HttpHandler::handle_put_retention(Connection &conn, HttpRequestState &state, const RouteMatch &) { + auto *ctx = state.response_context; + assert(ctx); // TODO: Parse retention policy from body and store auto json_response = R"({"policy_id":"example","status":"created"})"; auto http_response = format_json_response(200, json_response, state.arena, - state.http_request_id, state.connection_close); + ctx->http_request_id, ctx->connection_close); // Send through reorder queue and preprocessing to maintain proper ordering auto *conn_state = static_cast(conn.user_data); - conn_state->send_ordered_response(conn, state.sequence_id, http_response, - std::move(state.arena), - state.connection_close); + conn_state->send_ordered_response(conn, ctx, http_response, + std::move(state.arena)); } void HttpHandler::handle_get_retention(Connection &conn, HttpRequestState &state, const RouteMatch &) { + auto *ctx = state.response_context; + assert(ctx); // TODO: Extract policy_id from URL or return all policies auto json_response = R"({"policies":[]})"; auto http_response = format_json_response(200, json_response, state.arena, - state.http_request_id, state.connection_close); + ctx->http_request_id, ctx->connection_close); // Send through reorder queue and preprocessing to maintain proper ordering auto *conn_state = static_cast(conn.user_data); - conn_state->send_ordered_response(conn, state.sequence_id, http_response, - std::move(state.arena), - state.connection_close); + conn_state->send_ordered_response(conn, ctx, http_response, + std::move(state.arena)); } void HttpHandler::handle_delete_retention(Connection &conn, HttpRequestState &state, const RouteMatch &) { + auto *ctx = state.response_context; + assert(ctx); // TODO: Extract policy_id from URL and delete auto json_response = R"({"policy_id":"example","status":"deleted"})"; auto http_response = format_json_response(200, json_response, state.arena, - state.http_request_id, state.connection_close); + ctx->http_request_id, ctx->connection_close); // Send through reorder queue and preprocessing to maintain proper ordering auto *conn_state = static_cast(conn.user_data); - conn_state->send_ordered_response(conn, state.sequence_id, http_response, - std::move(state.arena), - state.connection_close); + conn_state->send_ordered_response(conn, ctx, http_response, + std::move(state.arena)); } void HttpHandler::handle_get_metrics(Connection &conn, HttpRequestState &state) { + auto *ctx = state.response_context; + assert(ctx); metrics_counter.inc(); auto metrics_span = metric::render(state.arena); @@ -450,19 +523,19 @@ void HttpHandler::handle_get_metrics(Connection &conn, // Build HTTP headers std::string_view headers; - if (state.connection_close) { + if (ctx->connection_close) { headers = static_format( state.arena, "HTTP/1.1 200 OK\r\n", "Content-Type: text/plain; version=0.0.4\r\n", "Content-Length: ", static_cast(total_size), "\r\n", - "X-Response-ID: ", static_cast(state.http_request_id), "\r\n", + "X-Response-ID: ", static_cast(ctx->http_request_id), "\r\n", "Connection: close\r\n", "\r\n"); } else { headers = static_format( state.arena, "HTTP/1.1 200 OK\r\n", "Content-Type: text/plain; version=0.0.4\r\n", "Content-Length: ", static_cast(total_size), "\r\n", - "X-Response-ID: ", static_cast(state.http_request_id), "\r\n", + "X-Response-ID: ", static_cast(ctx->http_request_id), "\r\n", "Connection: keep-alive\r\n", "\r\n"); } @@ -472,9 +545,7 @@ void HttpHandler::handle_get_metrics(Connection &conn, } auto *conn_state = static_cast(conn.user_data); - conn_state->send_ordered_response(conn, state.sequence_id, result, - std::move(state.arena), - state.connection_close); + conn_state->send_ordered_response(conn, ctx, result, std::move(state.arena)); } void HttpHandler::handle_get_ok(Connection &, HttpRequestState &) { @@ -485,41 +556,17 @@ void HttpHandler::handle_get_ok(Connection &, HttpRequestState &) { } void HttpHandler::handle_not_found(Connection &conn, HttpRequestState &state) { + auto *ctx = state.response_context; + assert(ctx); not_found_counter.inc(); auto json_response = R"({"error":"Not found"})"; auto http_response = format_json_response(404, json_response, state.arena, - state.http_request_id, state.connection_close); + ctx->http_request_id, ctx->connection_close); auto *conn_state = static_cast(conn.user_data); - conn_state->send_ordered_response(conn, state.sequence_id, http_response, - std::move(state.arena), - state.connection_close); -} - -void HttpConnectionState::send_ordered_response( - Connection &conn, int64_t sequence_id, - std::span http_response, Arena arena, - bool close_connection) { - - // Add to reorder queue with proper sequencing - ready_responses[sequence_id] = - ResponseData{http_response, std::move(arena), close_connection}; - - // Process ready responses in order and send via append_bytes - auto iter = ready_responses.begin(); - while (iter != ready_responses.end() && - iter->first == next_sequence_to_send) { - auto &[sequence_id, response_data] = *iter; - - // Send through append_bytes which handles write interest - conn.append_bytes(response_data.data, std::move(response_data.arena), - response_data.connection_close - ? ConnectionShutdown::WriteOnly - : ConnectionShutdown::None); - next_sequence_to_send++; - iter = ready_responses.erase(iter); - } + conn_state->send_ordered_response(conn, ctx, http_response, + std::move(state.arena)); } std::span diff --git a/src/http_handler.hpp b/src/http_handler.hpp index 18d1abc..3da8bdb 100644 --- a/src/http_handler.hpp +++ b/src/http_handler.hpp @@ -1,6 +1,7 @@ #pragma once #include +#include #include #include @@ -19,22 +20,18 @@ struct RouteMatch; /** * HTTP-specific response context stored in pipeline entries. - * Arena-allocated and passed through pipeline for response correlation. + * Arena-allocated and identified by a ProtocolHandle (sequence_id). + * Also owns the response data once it becomes ready. */ struct HttpResponseContext { int64_t sequence_id; // For response ordering in pipelining int64_t http_request_id; // For X-Response-ID header bool connection_close; // Whether to close connection after response -}; -/** - * Response data ready to send (sequence_id -> response data). - * Absence from map indicates response not ready yet. - */ -struct ResponseData { + // Response payload; populated when the response is ready. + bool ready = false; std::span data; Arena arena; - bool connection_close; }; /** @@ -69,6 +66,10 @@ struct HttpRequestState { 0; // X-Request-Id header value (for tracing/logging) int64_t sequence_id = 0; // Assigned for response ordering in pipelining + // HTTP response context for this request (set during batch processing). + // Non-owning: the context is owned by HttpConnectionState::contexts_. + HttpResponseContext *response_context = nullptr; + // Streaming parser for POST requests Arena::Ptr commit_parser; Arena::Ptr commit_request; @@ -89,15 +90,27 @@ struct HttpConnectionState { int64_t get_next_sequence_id() { return next_sequence_id++; } HttpConnectionState(); + ~HttpConnectionState(); - void send_ordered_response(Connection &conn, int64_t sequence_id, + // Register an HttpResponseContext so the pipeline can resolve it by handle. + void register_response_context(int64_t sequence_id, + std::unique_ptr ctx); + + // Resolve a ProtocolHandle back to its HttpResponseContext. + HttpResponseContext *resolve_response_context(ProtocolHandle handle); + + // Mark a context's response data as ready and try to send in-order. + // The context may already be registered in contexts_; if so, it is moved out + // and populated. A transient context pointer (e.g., for parse errors) works + // too. + void send_ordered_response(Connection &conn, HttpResponseContext *ctx, std::span http_response, - Arena arena, bool close_connection); + Arena arena); private: - // Response ordering for HTTP pipelining - std::map - ready_responses; // sequence_id -> response data + // All registered response contexts keyed by sequence_id. + // A context stays here until its response has been dispatched. + std::map> contexts_; int64_t next_sequence_to_send = 0; int64_t next_sequence_id = 0; }; diff --git a/src/pipeline_entry.hpp b/src/pipeline_entry.hpp index 74c12cf..e923812 100644 --- a/src/pipeline_entry.hpp +++ b/src/pipeline_entry.hpp @@ -17,8 +17,8 @@ struct CommitEntry { bool resolve_success = false; // Set by resolve stage bool persist_success = false; // Set by persist stage - // Protocol-agnostic context (arena-allocated, protocol-specific) - void *protocol_context = nullptr; + // Protocol-agnostic handle for correlating the response + ProtocolHandle handle = -1; const CommitRequest *commit_request = nullptr; // Points to request_arena data // Request arena contains parsed request data and response data @@ -28,9 +28,9 @@ struct CommitEntry { std::string_view response_json; CommitEntry() = default; // Default constructor for variant - explicit CommitEntry(WeakRef conn, void *ctx, + explicit CommitEntry(WeakRef conn, ProtocolHandle handle, const CommitRequest *req, Arena arena) - : connection(std::move(conn)), protocol_context(ctx), commit_request(req), + : connection(std::move(conn)), handle(handle), commit_request(req), request_arena(std::move(arena)) {} }; @@ -42,8 +42,8 @@ struct StatusEntry { WeakRef connection; int64_t version_upper_bound = 0; // Set by sequence stage - // Protocol-agnostic context (arena-allocated, protocol-specific) - void *protocol_context = nullptr; + // Protocol-agnostic handle for correlating the response + ProtocolHandle handle = -1; std::string_view status_request_id; // Points to request_arena data // Request arena for request data @@ -53,9 +53,9 @@ struct StatusEntry { std::string_view response_json; StatusEntry() = default; // Default constructor for variant - explicit StatusEntry(WeakRef conn, void *ctx, + explicit StatusEntry(WeakRef conn, ProtocolHandle handle, std::string_view request_id, Arena arena) - : connection(std::move(conn)), protocol_context(ctx), + : connection(std::move(conn)), handle(handle), status_request_id(request_id), request_arena(std::move(arena)) {} }; @@ -67,8 +67,8 @@ struct StatusEntry { struct HealthCheckEntry { WeakRef connection; - // Protocol-agnostic context (arena-allocated, protocol-specific) - void *protocol_context = nullptr; + // Protocol-agnostic handle for correlating the response + ProtocolHandle handle = -1; // Request arena for response data Arena request_arena; @@ -77,8 +77,9 @@ struct HealthCheckEntry { std::string_view response_json; HealthCheckEntry() = default; // Default constructor for variant - explicit HealthCheckEntry(WeakRef conn, void *ctx, Arena arena) - : connection(std::move(conn)), protocol_context(ctx), + explicit HealthCheckEntry(WeakRef conn, ProtocolHandle handle, + Arena arena) + : connection(std::move(conn)), handle(handle), request_arena(std::move(arena)) {} }; @@ -89,8 +90,8 @@ struct HealthCheckEntry { struct GetVersionEntry { WeakRef connection; - // Protocol-agnostic context (arena-allocated, protocol-specific) - void *protocol_context = nullptr; + // Protocol-agnostic handle for correlating the response + ProtocolHandle handle = -1; // Request arena for response data Arena request_arena; @@ -102,9 +103,9 @@ struct GetVersionEntry { int64_t version; GetVersionEntry() = default; // Default constructor for variant - explicit GetVersionEntry(WeakRef conn, void *ctx, Arena arena, - int64_t version) - : connection(std::move(conn)), protocol_context(ctx), + explicit GetVersionEntry(WeakRef conn, ProtocolHandle handle, + Arena arena, int64_t version) + : connection(std::move(conn)), handle(handle), request_arena(std::move(arena)), version(version) {} };