|
|
@ -35,42 +35,36 @@ |
|
|
|
#define EXPECTED_CONTENT_TYPE_LENGTH sizeof(EXPECTED_CONTENT_TYPE) - 1 |
|
|
|
#define EXPECTED_CONTENT_TYPE_LENGTH sizeof(EXPECTED_CONTENT_TYPE) - 1 |
|
|
|
|
|
|
|
|
|
|
|
namespace { |
|
|
|
namespace { |
|
|
|
|
|
|
|
|
|
|
|
struct call_data { |
|
|
|
struct call_data { |
|
|
|
grpc_call_combiner* call_combiner; |
|
|
|
grpc_call_combiner* call_combiner; |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
// Outgoing headers to add to send_initial_metadata.
|
|
|
|
grpc_linked_mdelem status; |
|
|
|
grpc_linked_mdelem status; |
|
|
|
grpc_linked_mdelem content_type; |
|
|
|
grpc_linked_mdelem content_type; |
|
|
|
|
|
|
|
|
|
|
|
/* did this request come with path query containing request payload */ |
|
|
|
// If we see the recv_message contents in the GET query string, we
|
|
|
|
bool seen_path_with_query; |
|
|
|
// store it here.
|
|
|
|
/* flag to ensure payload_bin is delivered only once */ |
|
|
|
grpc_core::ManualConstructor<grpc_core::SliceBufferByteStream> read_stream; |
|
|
|
bool payload_bin_delivered; |
|
|
|
bool have_read_stream; |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
// State for intercepting recv_initial_metadata.
|
|
|
|
|
|
|
|
grpc_closure recv_initial_metadata_ready; |
|
|
|
|
|
|
|
grpc_closure* original_recv_initial_metadata_ready; |
|
|
|
grpc_metadata_batch* recv_initial_metadata; |
|
|
|
grpc_metadata_batch* recv_initial_metadata; |
|
|
|
uint32_t* recv_initial_metadata_flags; |
|
|
|
uint32_t* recv_initial_metadata_flags; |
|
|
|
/** Closure to call when finished with the hs_on_recv hook */ |
|
|
|
bool seen_recv_initial_metadata_ready; |
|
|
|
grpc_closure* on_done_recv; |
|
|
|
|
|
|
|
/** Closure to call when we retrieve read message from the path URI
|
|
|
|
|
|
|
|
*/ |
|
|
|
|
|
|
|
grpc_closure* recv_message_ready; |
|
|
|
|
|
|
|
grpc_closure* on_complete; |
|
|
|
|
|
|
|
grpc_core::OrphanablePtr<grpc_core::ByteStream>* pp_recv_message; |
|
|
|
|
|
|
|
grpc_core::ManualConstructor<grpc_core::SliceBufferByteStream> read_stream; |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
/** Receive closures are chained: we inject this closure as the on_done_recv
|
|
|
|
// State for intercepting recv_message.
|
|
|
|
up-call on transport_op, and remember to call our on_done_recv member |
|
|
|
grpc_closure* original_recv_message_ready; |
|
|
|
after handling it. */ |
|
|
|
grpc_closure recv_message_ready; |
|
|
|
grpc_closure hs_on_recv; |
|
|
|
grpc_core::OrphanablePtr<grpc_core::ByteStream>* recv_message; |
|
|
|
grpc_closure hs_on_complete; |
|
|
|
bool seen_recv_message_ready; |
|
|
|
grpc_closure hs_recv_message_ready; |
|
|
|
|
|
|
|
}; |
|
|
|
}; |
|
|
|
|
|
|
|
|
|
|
|
struct channel_data { |
|
|
|
|
|
|
|
uint8_t unused; |
|
|
|
|
|
|
|
}; |
|
|
|
|
|
|
|
} // namespace
|
|
|
|
} // namespace
|
|
|
|
|
|
|
|
|
|
|
|
static grpc_error* server_filter_outgoing_metadata(grpc_call_element* elem, |
|
|
|
static grpc_error* hs_filter_outgoing_metadata(grpc_call_element* elem, |
|
|
|
grpc_metadata_batch* b) { |
|
|
|
grpc_metadata_batch* b) { |
|
|
|
if (b->idx.named.grpc_message != nullptr) { |
|
|
|
if (b->idx.named.grpc_message != nullptr) { |
|
|
|
grpc_slice pct_encoded_msg = grpc_percent_encode_slice( |
|
|
|
grpc_slice pct_encoded_msg = grpc_percent_encode_slice( |
|
|
@ -86,7 +80,7 @@ static grpc_error* server_filter_outgoing_metadata(grpc_call_element* elem, |
|
|
|
return GRPC_ERROR_NONE; |
|
|
|
return GRPC_ERROR_NONE; |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
static void add_error(const char* error_name, grpc_error** cumulative, |
|
|
|
static void hs_add_error(const char* error_name, grpc_error** cumulative, |
|
|
|
grpc_error* new_err) { |
|
|
|
grpc_error* new_err) { |
|
|
|
if (new_err == GRPC_ERROR_NONE) return; |
|
|
|
if (new_err == GRPC_ERROR_NONE) return; |
|
|
|
if (*cumulative == GRPC_ERROR_NONE) { |
|
|
|
if (*cumulative == GRPC_ERROR_NONE) { |
|
|
@ -95,7 +89,7 @@ static void add_error(const char* error_name, grpc_error** cumulative, |
|
|
|
*cumulative = grpc_error_add_child(*cumulative, new_err); |
|
|
|
*cumulative = grpc_error_add_child(*cumulative, new_err); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
static grpc_error* server_filter_incoming_metadata(grpc_call_element* elem, |
|
|
|
static grpc_error* hs_filter_incoming_metadata(grpc_call_element* elem, |
|
|
|
grpc_metadata_batch* b) { |
|
|
|
grpc_metadata_batch* b) { |
|
|
|
call_data* calld = static_cast<call_data*>(elem->call_data); |
|
|
|
call_data* calld = static_cast<call_data*>(elem->call_data); |
|
|
|
grpc_error* error = GRPC_ERROR_NONE; |
|
|
|
grpc_error* error = GRPC_ERROR_NONE; |
|
|
@ -119,14 +113,14 @@ static grpc_error* server_filter_incoming_metadata(grpc_call_element* elem, |
|
|
|
*calld->recv_initial_metadata_flags &= |
|
|
|
*calld->recv_initial_metadata_flags &= |
|
|
|
~GRPC_INITIAL_METADATA_IDEMPOTENT_REQUEST; |
|
|
|
~GRPC_INITIAL_METADATA_IDEMPOTENT_REQUEST; |
|
|
|
} else { |
|
|
|
} else { |
|
|
|
add_error(error_name, &error, |
|
|
|
hs_add_error(error_name, &error, |
|
|
|
grpc_attach_md_to_error( |
|
|
|
grpc_attach_md_to_error( |
|
|
|
GRPC_ERROR_CREATE_FROM_STATIC_STRING("Bad header"), |
|
|
|
GRPC_ERROR_CREATE_FROM_STATIC_STRING("Bad header"), |
|
|
|
b->idx.named.method->md)); |
|
|
|
b->idx.named.method->md)); |
|
|
|
} |
|
|
|
} |
|
|
|
grpc_metadata_batch_remove(b, b->idx.named.method); |
|
|
|
grpc_metadata_batch_remove(b, b->idx.named.method); |
|
|
|
} else { |
|
|
|
} else { |
|
|
|
add_error( |
|
|
|
hs_add_error( |
|
|
|
error_name, &error, |
|
|
|
error_name, &error, |
|
|
|
grpc_error_set_str( |
|
|
|
grpc_error_set_str( |
|
|
|
GRPC_ERROR_CREATE_FROM_STATIC_STRING("Missing header"), |
|
|
|
GRPC_ERROR_CREATE_FROM_STATIC_STRING("Missing header"), |
|
|
@ -135,14 +129,14 @@ static grpc_error* server_filter_incoming_metadata(grpc_call_element* elem, |
|
|
|
|
|
|
|
|
|
|
|
if (b->idx.named.te != nullptr) { |
|
|
|
if (b->idx.named.te != nullptr) { |
|
|
|
if (!grpc_mdelem_eq(b->idx.named.te->md, GRPC_MDELEM_TE_TRAILERS)) { |
|
|
|
if (!grpc_mdelem_eq(b->idx.named.te->md, GRPC_MDELEM_TE_TRAILERS)) { |
|
|
|
add_error(error_name, &error, |
|
|
|
hs_add_error(error_name, &error, |
|
|
|
grpc_attach_md_to_error( |
|
|
|
grpc_attach_md_to_error( |
|
|
|
GRPC_ERROR_CREATE_FROM_STATIC_STRING("Bad header"), |
|
|
|
GRPC_ERROR_CREATE_FROM_STATIC_STRING("Bad header"), |
|
|
|
b->idx.named.te->md)); |
|
|
|
b->idx.named.te->md)); |
|
|
|
} |
|
|
|
} |
|
|
|
grpc_metadata_batch_remove(b, b->idx.named.te); |
|
|
|
grpc_metadata_batch_remove(b, b->idx.named.te); |
|
|
|
} else { |
|
|
|
} else { |
|
|
|
add_error(error_name, &error, |
|
|
|
hs_add_error(error_name, &error, |
|
|
|
grpc_error_set_str( |
|
|
|
grpc_error_set_str( |
|
|
|
GRPC_ERROR_CREATE_FROM_STATIC_STRING("Missing header"), |
|
|
|
GRPC_ERROR_CREATE_FROM_STATIC_STRING("Missing header"), |
|
|
|
GRPC_ERROR_STR_KEY, grpc_slice_from_static_string("te"))); |
|
|
|
GRPC_ERROR_STR_KEY, grpc_slice_from_static_string("te"))); |
|
|
@ -152,14 +146,14 @@ static grpc_error* server_filter_incoming_metadata(grpc_call_element* elem, |
|
|
|
if (!grpc_mdelem_eq(b->idx.named.scheme->md, GRPC_MDELEM_SCHEME_HTTP) && |
|
|
|
if (!grpc_mdelem_eq(b->idx.named.scheme->md, GRPC_MDELEM_SCHEME_HTTP) && |
|
|
|
!grpc_mdelem_eq(b->idx.named.scheme->md, GRPC_MDELEM_SCHEME_HTTPS) && |
|
|
|
!grpc_mdelem_eq(b->idx.named.scheme->md, GRPC_MDELEM_SCHEME_HTTPS) && |
|
|
|
!grpc_mdelem_eq(b->idx.named.scheme->md, GRPC_MDELEM_SCHEME_GRPC)) { |
|
|
|
!grpc_mdelem_eq(b->idx.named.scheme->md, GRPC_MDELEM_SCHEME_GRPC)) { |
|
|
|
add_error(error_name, &error, |
|
|
|
hs_add_error(error_name, &error, |
|
|
|
grpc_attach_md_to_error( |
|
|
|
grpc_attach_md_to_error( |
|
|
|
GRPC_ERROR_CREATE_FROM_STATIC_STRING("Bad header"), |
|
|
|
GRPC_ERROR_CREATE_FROM_STATIC_STRING("Bad header"), |
|
|
|
b->idx.named.scheme->md)); |
|
|
|
b->idx.named.scheme->md)); |
|
|
|
} |
|
|
|
} |
|
|
|
grpc_metadata_batch_remove(b, b->idx.named.scheme); |
|
|
|
grpc_metadata_batch_remove(b, b->idx.named.scheme); |
|
|
|
} else { |
|
|
|
} else { |
|
|
|
add_error( |
|
|
|
hs_add_error( |
|
|
|
error_name, &error, |
|
|
|
error_name, &error, |
|
|
|
grpc_error_set_str( |
|
|
|
grpc_error_set_str( |
|
|
|
GRPC_ERROR_CREATE_FROM_STATIC_STRING("Missing header"), |
|
|
|
GRPC_ERROR_CREATE_FROM_STATIC_STRING("Missing header"), |
|
|
@ -196,7 +190,8 @@ static grpc_error* server_filter_incoming_metadata(grpc_call_element* elem, |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
if (b->idx.named.path == nullptr) { |
|
|
|
if (b->idx.named.path == nullptr) { |
|
|
|
add_error(error_name, &error, |
|
|
|
hs_add_error( |
|
|
|
|
|
|
|
error_name, &error, |
|
|
|
grpc_error_set_str( |
|
|
|
grpc_error_set_str( |
|
|
|
GRPC_ERROR_CREATE_FROM_STATIC_STRING("Missing header"), |
|
|
|
GRPC_ERROR_CREATE_FROM_STATIC_STRING("Missing header"), |
|
|
|
GRPC_ERROR_STR_KEY, grpc_slice_from_static_string(":path"))); |
|
|
|
GRPC_ERROR_STR_KEY, grpc_slice_from_static_string(":path"))); |
|
|
@ -235,7 +230,7 @@ static grpc_error* server_filter_incoming_metadata(grpc_call_element* elem, |
|
|
|
GRPC_SLICE_LENGTH(query_slice), k_url_safe)); |
|
|
|
GRPC_SLICE_LENGTH(query_slice), k_url_safe)); |
|
|
|
calld->read_stream.Init(&read_slice_buffer, 0); |
|
|
|
calld->read_stream.Init(&read_slice_buffer, 0); |
|
|
|
grpc_slice_buffer_destroy_internal(&read_slice_buffer); |
|
|
|
grpc_slice_buffer_destroy_internal(&read_slice_buffer); |
|
|
|
calld->seen_path_with_query = true; |
|
|
|
calld->have_read_stream = true; |
|
|
|
grpc_slice_unref_internal(query_slice); |
|
|
|
grpc_slice_unref_internal(query_slice); |
|
|
|
} else { |
|
|
|
} else { |
|
|
|
gpr_log(GPR_ERROR, "GET request without QUERY"); |
|
|
|
gpr_log(GPR_ERROR, "GET request without QUERY"); |
|
|
@ -246,7 +241,7 @@ static grpc_error* server_filter_incoming_metadata(grpc_call_element* elem, |
|
|
|
grpc_linked_mdelem* el = b->idx.named.host; |
|
|
|
grpc_linked_mdelem* el = b->idx.named.host; |
|
|
|
grpc_mdelem md = GRPC_MDELEM_REF(el->md); |
|
|
|
grpc_mdelem md = GRPC_MDELEM_REF(el->md); |
|
|
|
grpc_metadata_batch_remove(b, el); |
|
|
|
grpc_metadata_batch_remove(b, el); |
|
|
|
add_error(error_name, &error, |
|
|
|
hs_add_error(error_name, &error, |
|
|
|
grpc_metadata_batch_add_head( |
|
|
|
grpc_metadata_batch_add_head( |
|
|
|
b, el, |
|
|
|
b, el, |
|
|
|
grpc_mdelem_from_slices( |
|
|
|
grpc_mdelem_from_slices( |
|
|
@ -256,7 +251,7 @@ static grpc_error* server_filter_incoming_metadata(grpc_call_element* elem, |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
if (b->idx.named.authority == nullptr) { |
|
|
|
if (b->idx.named.authority == nullptr) { |
|
|
|
add_error( |
|
|
|
hs_add_error( |
|
|
|
error_name, &error, |
|
|
|
error_name, &error, |
|
|
|
grpc_error_set_str( |
|
|
|
grpc_error_set_str( |
|
|
|
GRPC_ERROR_CREATE_FROM_STATIC_STRING("Missing header"), |
|
|
|
GRPC_ERROR_CREATE_FROM_STATIC_STRING("Missing header"), |
|
|
@ -266,49 +261,55 @@ static grpc_error* server_filter_incoming_metadata(grpc_call_element* elem, |
|
|
|
return error; |
|
|
|
return error; |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
static void hs_on_recv(void* user_data, grpc_error* err) { |
|
|
|
static void hs_recv_initial_metadata_ready(void* user_data, grpc_error* err) { |
|
|
|
grpc_call_element* elem = static_cast<grpc_call_element*>(user_data); |
|
|
|
grpc_call_element* elem = static_cast<grpc_call_element*>(user_data); |
|
|
|
call_data* calld = static_cast<call_data*>(elem->call_data); |
|
|
|
call_data* calld = static_cast<call_data*>(elem->call_data); |
|
|
|
|
|
|
|
calld->seen_recv_initial_metadata_ready = true; |
|
|
|
if (err == GRPC_ERROR_NONE) { |
|
|
|
if (err == GRPC_ERROR_NONE) { |
|
|
|
err = server_filter_incoming_metadata(elem, calld->recv_initial_metadata); |
|
|
|
err = hs_filter_incoming_metadata(elem, calld->recv_initial_metadata); |
|
|
|
|
|
|
|
if (calld->seen_recv_message_ready) { |
|
|
|
|
|
|
|
// We've already seen the recv_message callback, but we previously
|
|
|
|
|
|
|
|
// deferred it, so we need to return it here.
|
|
|
|
|
|
|
|
// Replace the recv_message byte stream if needed.
|
|
|
|
|
|
|
|
if (calld->have_read_stream) { |
|
|
|
|
|
|
|
calld->recv_message->reset(calld->read_stream.get()); |
|
|
|
|
|
|
|
calld->have_read_stream = false; |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
// Re-enter call combiner for original_recv_message_ready, since the
|
|
|
|
|
|
|
|
// surface code will release the call combiner for each callback it
|
|
|
|
|
|
|
|
// receives.
|
|
|
|
|
|
|
|
GRPC_CALL_COMBINER_START( |
|
|
|
|
|
|
|
calld->call_combiner, calld->original_recv_message_ready, |
|
|
|
|
|
|
|
GRPC_ERROR_REF(err), |
|
|
|
|
|
|
|
"resuming recv_message_ready from recv_initial_metadata_ready"); |
|
|
|
|
|
|
|
} |
|
|
|
} else { |
|
|
|
} else { |
|
|
|
GRPC_ERROR_REF(err); |
|
|
|
GRPC_ERROR_REF(err); |
|
|
|
} |
|
|
|
} |
|
|
|
GRPC_CLOSURE_RUN(calld->on_done_recv, err); |
|
|
|
GRPC_CLOSURE_RUN(calld->original_recv_initial_metadata_ready, err); |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
static void hs_on_complete(void* user_data, grpc_error* err) { |
|
|
|
|
|
|
|
grpc_call_element* elem = static_cast<grpc_call_element*>(user_data); |
|
|
|
|
|
|
|
call_data* calld = static_cast<call_data*>(elem->call_data); |
|
|
|
|
|
|
|
/* Call recv_message_ready if we got the payload via the path field */ |
|
|
|
|
|
|
|
if (calld->seen_path_with_query && calld->recv_message_ready != nullptr) { |
|
|
|
|
|
|
|
calld->pp_recv_message->reset( |
|
|
|
|
|
|
|
calld->payload_bin_delivered ? nullptr |
|
|
|
|
|
|
|
: reinterpret_cast<grpc_core::ByteStream*>( |
|
|
|
|
|
|
|
calld->read_stream.get())); |
|
|
|
|
|
|
|
// Re-enter call combiner for recv_message_ready, since the surface
|
|
|
|
|
|
|
|
// code will release the call combiner for each callback it receives.
|
|
|
|
|
|
|
|
GRPC_CALL_COMBINER_START(calld->call_combiner, calld->recv_message_ready, |
|
|
|
|
|
|
|
GRPC_ERROR_REF(err), |
|
|
|
|
|
|
|
"resuming recv_message_ready from on_complete"); |
|
|
|
|
|
|
|
calld->recv_message_ready = nullptr; |
|
|
|
|
|
|
|
calld->payload_bin_delivered = true; |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
GRPC_CLOSURE_RUN(calld->on_complete, GRPC_ERROR_REF(err)); |
|
|
|
|
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
static void hs_recv_message_ready(void* user_data, grpc_error* err) { |
|
|
|
static void hs_recv_message_ready(void* user_data, grpc_error* err) { |
|
|
|
grpc_call_element* elem = static_cast<grpc_call_element*>(user_data); |
|
|
|
grpc_call_element* elem = static_cast<grpc_call_element*>(user_data); |
|
|
|
call_data* calld = static_cast<call_data*>(elem->call_data); |
|
|
|
call_data* calld = static_cast<call_data*>(elem->call_data); |
|
|
|
if (calld->seen_path_with_query) { |
|
|
|
calld->seen_recv_message_ready = true; |
|
|
|
// Do nothing. This is probably a GET request, and payload will be
|
|
|
|
if (calld->seen_recv_initial_metadata_ready) { |
|
|
|
// returned in hs_on_complete callback.
|
|
|
|
// We've already seen the recv_initial_metadata callback, so
|
|
|
|
|
|
|
|
// replace the recv_message byte stream if needed and invoke the
|
|
|
|
|
|
|
|
// original recv_message callback immediately.
|
|
|
|
|
|
|
|
if (calld->have_read_stream) { |
|
|
|
|
|
|
|
calld->recv_message->reset(calld->read_stream.get()); |
|
|
|
|
|
|
|
calld->have_read_stream = false; |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
GRPC_CLOSURE_RUN(calld->original_recv_message_ready, GRPC_ERROR_REF(err)); |
|
|
|
|
|
|
|
} else { |
|
|
|
|
|
|
|
// We have not yet seen the recv_initial_metadata callback, so we
|
|
|
|
|
|
|
|
// need to wait to see if this is a GET request.
|
|
|
|
// Note that we release the call combiner here, so that other
|
|
|
|
// Note that we release the call combiner here, so that other
|
|
|
|
// callbacks can run.
|
|
|
|
// callbacks can run.
|
|
|
|
GRPC_CALL_COMBINER_STOP(calld->call_combiner, |
|
|
|
GRPC_CALL_COMBINER_STOP( |
|
|
|
"pausing recv_message_ready until on_complete"); |
|
|
|
calld->call_combiner, |
|
|
|
} else { |
|
|
|
"pausing recv_message_ready until recv_initial_metadata_ready"); |
|
|
|
GRPC_CLOSURE_RUN(calld->recv_message_ready, GRPC_ERROR_REF(err)); |
|
|
|
|
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
@ -320,18 +321,18 @@ static grpc_error* hs_mutate_op(grpc_call_element* elem, |
|
|
|
if (op->send_initial_metadata) { |
|
|
|
if (op->send_initial_metadata) { |
|
|
|
grpc_error* error = GRPC_ERROR_NONE; |
|
|
|
grpc_error* error = GRPC_ERROR_NONE; |
|
|
|
static const char* error_name = "Failed sending initial metadata"; |
|
|
|
static const char* error_name = "Failed sending initial metadata"; |
|
|
|
add_error(error_name, &error, |
|
|
|
hs_add_error(error_name, &error, |
|
|
|
grpc_metadata_batch_add_head( |
|
|
|
grpc_metadata_batch_add_head( |
|
|
|
op->payload->send_initial_metadata.send_initial_metadata, |
|
|
|
op->payload->send_initial_metadata.send_initial_metadata, |
|
|
|
&calld->status, GRPC_MDELEM_STATUS_200)); |
|
|
|
&calld->status, GRPC_MDELEM_STATUS_200)); |
|
|
|
add_error(error_name, &error, |
|
|
|
hs_add_error(error_name, &error, |
|
|
|
grpc_metadata_batch_add_tail( |
|
|
|
grpc_metadata_batch_add_tail( |
|
|
|
op->payload->send_initial_metadata.send_initial_metadata, |
|
|
|
op->payload->send_initial_metadata.send_initial_metadata, |
|
|
|
&calld->content_type, |
|
|
|
&calld->content_type, |
|
|
|
GRPC_MDELEM_CONTENT_TYPE_APPLICATION_SLASH_GRPC)); |
|
|
|
GRPC_MDELEM_CONTENT_TYPE_APPLICATION_SLASH_GRPC)); |
|
|
|
add_error( |
|
|
|
hs_add_error( |
|
|
|
error_name, &error, |
|
|
|
error_name, &error, |
|
|
|
server_filter_outgoing_metadata( |
|
|
|
hs_filter_outgoing_metadata( |
|
|
|
elem, op->payload->send_initial_metadata.send_initial_metadata)); |
|
|
|
elem, op->payload->send_initial_metadata.send_initial_metadata)); |
|
|
|
if (error != GRPC_ERROR_NONE) return error; |
|
|
|
if (error != GRPC_ERROR_NONE) return error; |
|
|
|
} |
|
|
|
} |
|
|
@ -343,27 +344,21 @@ static grpc_error* hs_mutate_op(grpc_call_element* elem, |
|
|
|
op->payload->recv_initial_metadata.recv_initial_metadata; |
|
|
|
op->payload->recv_initial_metadata.recv_initial_metadata; |
|
|
|
calld->recv_initial_metadata_flags = |
|
|
|
calld->recv_initial_metadata_flags = |
|
|
|
op->payload->recv_initial_metadata.recv_flags; |
|
|
|
op->payload->recv_initial_metadata.recv_flags; |
|
|
|
calld->on_done_recv = |
|
|
|
calld->original_recv_initial_metadata_ready = |
|
|
|
op->payload->recv_initial_metadata.recv_initial_metadata_ready; |
|
|
|
op->payload->recv_initial_metadata.recv_initial_metadata_ready; |
|
|
|
op->payload->recv_initial_metadata.recv_initial_metadata_ready = |
|
|
|
op->payload->recv_initial_metadata.recv_initial_metadata_ready = |
|
|
|
&calld->hs_on_recv; |
|
|
|
&calld->recv_initial_metadata_ready; |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
if (op->recv_message) { |
|
|
|
if (op->recv_message) { |
|
|
|
calld->recv_message_ready = op->payload->recv_message.recv_message_ready; |
|
|
|
calld->recv_message = op->payload->recv_message.recv_message; |
|
|
|
calld->pp_recv_message = op->payload->recv_message.recv_message; |
|
|
|
calld->original_recv_message_ready = |
|
|
|
if (op->payload->recv_message.recv_message_ready) { |
|
|
|
op->payload->recv_message.recv_message_ready; |
|
|
|
op->payload->recv_message.recv_message_ready = |
|
|
|
op->payload->recv_message.recv_message_ready = &calld->recv_message_ready; |
|
|
|
&calld->hs_recv_message_ready; |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
if (op->on_complete) { |
|
|
|
|
|
|
|
calld->on_complete = op->on_complete; |
|
|
|
|
|
|
|
op->on_complete = &calld->hs_on_complete; |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
if (op->send_trailing_metadata) { |
|
|
|
if (op->send_trailing_metadata) { |
|
|
|
grpc_error* error = server_filter_outgoing_metadata( |
|
|
|
grpc_error* error = hs_filter_outgoing_metadata( |
|
|
|
elem, op->payload->send_trailing_metadata.send_trailing_metadata); |
|
|
|
elem, op->payload->send_trailing_metadata.send_trailing_metadata); |
|
|
|
if (error != GRPC_ERROR_NONE) return error; |
|
|
|
if (error != GRPC_ERROR_NONE) return error; |
|
|
|
} |
|
|
|
} |
|
|
@ -385,50 +380,47 @@ static void hs_start_transport_stream_op_batch( |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
/* Constructor for call_data */ |
|
|
|
/* Constructor for call_data */ |
|
|
|
static grpc_error* init_call_elem(grpc_call_element* elem, |
|
|
|
static grpc_error* hs_init_call_elem(grpc_call_element* elem, |
|
|
|
const grpc_call_element_args* args) { |
|
|
|
const grpc_call_element_args* args) { |
|
|
|
/* grab pointers to our data from the call element */ |
|
|
|
|
|
|
|
call_data* calld = static_cast<call_data*>(elem->call_data); |
|
|
|
call_data* calld = static_cast<call_data*>(elem->call_data); |
|
|
|
/* initialize members */ |
|
|
|
|
|
|
|
calld->call_combiner = args->call_combiner; |
|
|
|
calld->call_combiner = args->call_combiner; |
|
|
|
GRPC_CLOSURE_INIT(&calld->hs_on_recv, hs_on_recv, elem, |
|
|
|
GRPC_CLOSURE_INIT(&calld->recv_initial_metadata_ready, |
|
|
|
grpc_schedule_on_exec_ctx); |
|
|
|
hs_recv_initial_metadata_ready, elem, |
|
|
|
GRPC_CLOSURE_INIT(&calld->hs_on_complete, hs_on_complete, elem, |
|
|
|
|
|
|
|
grpc_schedule_on_exec_ctx); |
|
|
|
grpc_schedule_on_exec_ctx); |
|
|
|
GRPC_CLOSURE_INIT(&calld->hs_recv_message_ready, hs_recv_message_ready, elem, |
|
|
|
GRPC_CLOSURE_INIT(&calld->recv_message_ready, hs_recv_message_ready, elem, |
|
|
|
grpc_schedule_on_exec_ctx); |
|
|
|
grpc_schedule_on_exec_ctx); |
|
|
|
return GRPC_ERROR_NONE; |
|
|
|
return GRPC_ERROR_NONE; |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
/* Destructor for call_data */ |
|
|
|
/* Destructor for call_data */ |
|
|
|
static void destroy_call_elem(grpc_call_element* elem, |
|
|
|
static void hs_destroy_call_elem(grpc_call_element* elem, |
|
|
|
const grpc_call_final_info* final_info, |
|
|
|
const grpc_call_final_info* final_info, |
|
|
|
grpc_closure* ignored) { |
|
|
|
grpc_closure* ignored) { |
|
|
|
call_data* calld = static_cast<call_data*>(elem->call_data); |
|
|
|
call_data* calld = static_cast<call_data*>(elem->call_data); |
|
|
|
if (calld->seen_path_with_query && !calld->payload_bin_delivered) { |
|
|
|
if (calld->have_read_stream) { |
|
|
|
calld->read_stream->Orphan(); |
|
|
|
calld->read_stream->Orphan(); |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
/* Constructor for channel_data */ |
|
|
|
/* Constructor for channel_data */ |
|
|
|
static grpc_error* init_channel_elem(grpc_channel_element* elem, |
|
|
|
static grpc_error* hs_init_channel_elem(grpc_channel_element* elem, |
|
|
|
grpc_channel_element_args* args) { |
|
|
|
grpc_channel_element_args* args) { |
|
|
|
GPR_ASSERT(!args->is_last); |
|
|
|
GPR_ASSERT(!args->is_last); |
|
|
|
return GRPC_ERROR_NONE; |
|
|
|
return GRPC_ERROR_NONE; |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
/* Destructor for channel data */ |
|
|
|
/* Destructor for channel data */ |
|
|
|
static void destroy_channel_elem(grpc_channel_element* elem) {} |
|
|
|
static void hs_destroy_channel_elem(grpc_channel_element* elem) {} |
|
|
|
|
|
|
|
|
|
|
|
const grpc_channel_filter grpc_http_server_filter = { |
|
|
|
const grpc_channel_filter grpc_http_server_filter = { |
|
|
|
hs_start_transport_stream_op_batch, |
|
|
|
hs_start_transport_stream_op_batch, |
|
|
|
grpc_channel_next_op, |
|
|
|
grpc_channel_next_op, |
|
|
|
sizeof(call_data), |
|
|
|
sizeof(call_data), |
|
|
|
init_call_elem, |
|
|
|
hs_init_call_elem, |
|
|
|
grpc_call_stack_ignore_set_pollset_or_pollset_set, |
|
|
|
grpc_call_stack_ignore_set_pollset_or_pollset_set, |
|
|
|
destroy_call_elem, |
|
|
|
hs_destroy_call_elem, |
|
|
|
sizeof(channel_data), |
|
|
|
0, |
|
|
|
init_channel_elem, |
|
|
|
hs_init_channel_elem, |
|
|
|
destroy_channel_elem, |
|
|
|
hs_destroy_channel_elem, |
|
|
|
grpc_channel_next_get_info, |
|
|
|
grpc_channel_next_get_info, |
|
|
|
"http-server"}; |
|
|
|
"http-server"}; |
|
|
|