Changes across the board from pollset to pollset_set
This commit is contained in:
parent
279681311f
commit
4afce7e66f
|
|
@ -176,7 +176,7 @@ const grpc_channel_filter grpc_client_census_filter = {
|
|||
grpc_channel_next_op,
|
||||
sizeof(call_data),
|
||||
client_init_call_elem,
|
||||
grpc_call_stack_ignore_set_pollset,
|
||||
grpc_call_stack_ignore_set_pollset_or_pollset_set,
|
||||
client_destroy_call_elem,
|
||||
sizeof(channel_data),
|
||||
init_channel_elem,
|
||||
|
|
@ -189,7 +189,7 @@ const grpc_channel_filter grpc_server_census_filter = {
|
|||
grpc_channel_next_op,
|
||||
sizeof(call_data),
|
||||
server_init_call_elem,
|
||||
grpc_call_stack_ignore_set_pollset,
|
||||
grpc_call_stack_ignore_set_pollset_or_pollset_set,
|
||||
server_destroy_call_elem,
|
||||
sizeof(channel_data),
|
||||
init_channel_elem,
|
||||
|
|
|
|||
|
|
@ -368,7 +368,7 @@ static int cc_pick_subchannel(grpc_exec_ctx *exec_ctx, void *elemp,
|
|||
int r;
|
||||
GRPC_LB_POLICY_REF(lb_policy, "cc_pick_subchannel");
|
||||
gpr_mu_unlock(&chand->mu_config);
|
||||
r = grpc_lb_policy_pick(exec_ctx, lb_policy, calld->pollset,
|
||||
r = grpc_lb_policy_pick(exec_ctx, lb_policy, calld->pollset_set,
|
||||
initial_metadata, initial_metadata_flags,
|
||||
connected_subchannel, on_ready);
|
||||
GRPC_LB_POLICY_UNREF(exec_ctx, lb_policy, "cc_pick_subchannel");
|
||||
|
|
@ -446,10 +446,22 @@ static void destroy_channel_elem(grpc_exec_ctx *exec_ctx,
|
|||
gpr_mu_destroy(&chand->mu_config);
|
||||
}
|
||||
|
||||
static void cc_set_pollset(grpc_exec_ctx *exec_ctx, grpc_call_element *elem,
|
||||
grpc_pollset *pollset) {
|
||||
static void cc_set_pollset_or_pollset_set(grpc_exec_ctx *exec_ctx,
|
||||
grpc_call_element *elem,
|
||||
grpc_pollset *pollset,
|
||||
grpc_pollset_set *or_pollset_set) {
|
||||
GPR_ASSERT(!(pollset != NULL && or_pollset_set != NULL));
|
||||
GPR_ASSERT(pollset != NULL || or_pollset_set != NULL);
|
||||
|
||||
call_data *calld = elem->call_data;
|
||||
calld->pollset = pollset;
|
||||
if (pollset != NULL) {
|
||||
calld->pollset = pollset;
|
||||
grpc_pollset_set_add_pollset(exec_ctx, calld->pollset_set, pollset);
|
||||
} else if (or_pollset_set != NULL) {
|
||||
calld->pollset = NULL;
|
||||
grpc_pollset_set_add_pollset_set(exec_ctx, calld->pollset_set,
|
||||
or_pollset_set);
|
||||
}
|
||||
}
|
||||
|
||||
const grpc_channel_filter grpc_client_channel_filter = {
|
||||
|
|
@ -457,7 +469,7 @@ const grpc_channel_filter grpc_client_channel_filter = {
|
|||
cc_start_transport_op,
|
||||
sizeof(call_data),
|
||||
init_call_elem,
|
||||
cc_set_pollset,
|
||||
cc_set_pollset_or_pollset_set,
|
||||
destroy_call_elem,
|
||||
sizeof(channel_data),
|
||||
init_channel_elem,
|
||||
|
|
|
|||
|
|
@ -99,12 +99,12 @@ void grpc_lb_policy_weak_unref(grpc_exec_ctx *exec_ctx,
|
|||
}
|
||||
|
||||
int grpc_lb_policy_pick(grpc_exec_ctx *exec_ctx, grpc_lb_policy *policy,
|
||||
grpc_pollset *pollset,
|
||||
grpc_pollset_set *pollset_set,
|
||||
grpc_metadata_batch *initial_metadata,
|
||||
uint32_t initial_metadata_flags,
|
||||
grpc_connected_subchannel **target,
|
||||
grpc_closure *on_complete) {
|
||||
return policy->vtable->pick(exec_ctx, policy, pollset, initial_metadata,
|
||||
return policy->vtable->pick(exec_ctx, policy, pollset_set, initial_metadata,
|
||||
initial_metadata_flags, target, on_complete);
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -59,7 +59,8 @@ struct grpc_lb_policy_vtable {
|
|||
|
||||
/** implement grpc_lb_policy_pick */
|
||||
int (*pick)(grpc_exec_ctx *exec_ctx, grpc_lb_policy *policy,
|
||||
grpc_pollset *pollset, grpc_metadata_batch *initial_metadata,
|
||||
grpc_pollset_set *pollset_set,
|
||||
grpc_metadata_batch *initial_metadata,
|
||||
uint32_t initial_metadata_flags,
|
||||
grpc_connected_subchannel **target, grpc_closure *on_complete);
|
||||
void (*cancel_pick)(grpc_exec_ctx *exec_ctx, grpc_lb_policy *policy,
|
||||
|
|
@ -124,7 +125,7 @@ void grpc_lb_policy_init(grpc_lb_policy *policy,
|
|||
\a target.
|
||||
Picking can be asynchronous. Any IO should be done under \a pollset. */
|
||||
int grpc_lb_policy_pick(grpc_exec_ctx *exec_ctx, grpc_lb_policy *policy,
|
||||
grpc_pollset *pollset,
|
||||
grpc_pollset_set *pollset_set,
|
||||
grpc_metadata_batch *initial_metadata,
|
||||
uint32_t initial_metadata_flags,
|
||||
grpc_connected_subchannel **target,
|
||||
|
|
|
|||
|
|
@ -699,7 +699,7 @@ grpc_subchannel_call *grpc_connected_subchannel_create_call(
|
|||
GRPC_CONNECTED_SUBCHANNEL_REF(con, "subchannel_call");
|
||||
grpc_call_stack_init(exec_ctx, chanstk, 1, subchannel_call_destroy, call,
|
||||
NULL, NULL, callstk);
|
||||
grpc_call_stack_set_pollset(exec_ctx, callstk, pollset);
|
||||
grpc_call_stack_set_pollset_or_pollset_set(exec_ctx, callstk, pollset, NULL);
|
||||
return call;
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -68,6 +68,7 @@ void grpc_subchannel_call_holder_init(
|
|||
holder->waiting_ops_capacity = 0;
|
||||
holder->creation_phase = GRPC_SUBCHANNEL_CALL_HOLDER_NOT_CREATING;
|
||||
holder->owning_call = owning_call;
|
||||
holder->pollset_set = grpc_pollset_set_create();
|
||||
}
|
||||
|
||||
void grpc_subchannel_call_holder_destroy(grpc_exec_ctx *exec_ctx,
|
||||
|
|
@ -81,6 +82,7 @@ void grpc_subchannel_call_holder_destroy(grpc_exec_ctx *exec_ctx,
|
|||
gpr_mu_destroy(&holder->mu);
|
||||
GPR_ASSERT(holder->waiting_ops_count == 0);
|
||||
gpr_free(holder->waiting_ops);
|
||||
grpc_pollset_set_destroy(holder->pollset_set);
|
||||
}
|
||||
|
||||
void grpc_subchannel_call_holder_perform_op(grpc_exec_ctx *exec_ctx,
|
||||
|
|
|
|||
|
|
@ -72,6 +72,7 @@ typedef struct grpc_subchannel_call_holder {
|
|||
grpc_subchannel_call_holder_creation_phase creation_phase;
|
||||
grpc_connected_subchannel *connected_subchannel;
|
||||
grpc_pollset *pollset;
|
||||
grpc_pollset_set *pollset_set;
|
||||
|
||||
grpc_transport_stream_op *waiting_ops;
|
||||
size_t waiting_ops_count;
|
||||
|
|
|
|||
|
|
@ -39,7 +39,7 @@
|
|||
|
||||
typedef struct pending_pick {
|
||||
struct pending_pick *next;
|
||||
grpc_pollset *pollset;
|
||||
grpc_pollset_set *pollset_set;
|
||||
uint32_t initial_metadata_flags;
|
||||
grpc_connected_subchannel **target;
|
||||
grpc_closure *on_complete;
|
||||
|
|
@ -118,8 +118,8 @@ static void pf_shutdown(grpc_exec_ctx *exec_ctx, grpc_lb_policy *pol) {
|
|||
while (pp != NULL) {
|
||||
pending_pick *next = pp->next;
|
||||
*pp->target = NULL;
|
||||
grpc_pollset_set_del_pollset(exec_ctx, p->base.interested_parties,
|
||||
pp->pollset);
|
||||
grpc_pollset_set_del_pollset_set(exec_ctx, p->base.interested_parties,
|
||||
pp->pollset_set);
|
||||
grpc_exec_ctx_enqueue(exec_ctx, pp->on_complete, true, NULL);
|
||||
gpr_free(pp);
|
||||
pp = next;
|
||||
|
|
@ -136,8 +136,8 @@ static void pf_cancel_pick(grpc_exec_ctx *exec_ctx, grpc_lb_policy *pol,
|
|||
while (pp != NULL) {
|
||||
pending_pick *next = pp->next;
|
||||
if (pp->target == target) {
|
||||
grpc_pollset_set_del_pollset(exec_ctx, p->base.interested_parties,
|
||||
pp->pollset);
|
||||
grpc_pollset_set_del_pollset_set(exec_ctx, p->base.interested_parties,
|
||||
pp->pollset_set);
|
||||
*target = NULL;
|
||||
grpc_exec_ctx_enqueue(exec_ctx, pp->on_complete, false, NULL);
|
||||
gpr_free(pp);
|
||||
|
|
@ -162,8 +162,8 @@ static void pf_cancel_picks(grpc_exec_ctx *exec_ctx, grpc_lb_policy *pol,
|
|||
pending_pick *next = pp->next;
|
||||
if ((pp->initial_metadata_flags & initial_metadata_flags_mask) ==
|
||||
initial_metadata_flags_eq) {
|
||||
grpc_pollset_set_del_pollset(exec_ctx, p->base.interested_parties,
|
||||
pp->pollset);
|
||||
grpc_pollset_set_del_pollset_set(exec_ctx, p->base.interested_parties,
|
||||
pp->pollset_set);
|
||||
grpc_exec_ctx_enqueue(exec_ctx, pp->on_complete, false, NULL);
|
||||
gpr_free(pp);
|
||||
} else {
|
||||
|
|
@ -196,7 +196,8 @@ static void pf_exit_idle(grpc_exec_ctx *exec_ctx, grpc_lb_policy *pol) {
|
|||
}
|
||||
|
||||
static int pf_pick(grpc_exec_ctx *exec_ctx, grpc_lb_policy *pol,
|
||||
grpc_pollset *pollset, grpc_metadata_batch *initial_metadata,
|
||||
grpc_pollset_set *pollset_set,
|
||||
grpc_metadata_batch *initial_metadata,
|
||||
uint32_t initial_metadata_flags,
|
||||
grpc_connected_subchannel **target,
|
||||
grpc_closure *on_complete) {
|
||||
|
|
@ -221,10 +222,11 @@ static int pf_pick(grpc_exec_ctx *exec_ctx, grpc_lb_policy *pol,
|
|||
if (!p->started_picking) {
|
||||
start_picking(exec_ctx, p);
|
||||
}
|
||||
grpc_pollset_set_add_pollset(exec_ctx, p->base.interested_parties, pollset);
|
||||
grpc_pollset_set_add_pollset_set(exec_ctx, p->base.interested_parties,
|
||||
pollset_set);
|
||||
pp = gpr_malloc(sizeof(*pp));
|
||||
pp->next = p->pending_picks;
|
||||
pp->pollset = pollset;
|
||||
pp->pollset_set = pollset_set;
|
||||
pp->target = target;
|
||||
pp->initial_metadata_flags = initial_metadata_flags;
|
||||
pp->on_complete = on_complete;
|
||||
|
|
@ -304,8 +306,8 @@ static void pf_connectivity_changed(grpc_exec_ctx *exec_ctx, void *arg,
|
|||
while ((pp = p->pending_picks)) {
|
||||
p->pending_picks = pp->next;
|
||||
*pp->target = selected;
|
||||
grpc_pollset_set_del_pollset(exec_ctx, p->base.interested_parties,
|
||||
pp->pollset);
|
||||
grpc_pollset_set_del_pollset_set(exec_ctx, p->base.interested_parties,
|
||||
pp->pollset_set);
|
||||
grpc_exec_ctx_enqueue(exec_ctx, pp->on_complete, true, NULL);
|
||||
gpr_free(pp);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -48,7 +48,7 @@ int grpc_lb_round_robin_trace = 0;
|
|||
* Once a pick is available, \a target is updated and \a on_complete called. */
|
||||
typedef struct pending_pick {
|
||||
struct pending_pick *next;
|
||||
grpc_pollset *pollset;
|
||||
grpc_pollset_set *pollset_set;
|
||||
uint32_t initial_metadata_flags;
|
||||
grpc_connected_subchannel **target;
|
||||
grpc_closure *on_complete;
|
||||
|
|
@ -262,8 +262,8 @@ static void rr_cancel_pick(grpc_exec_ctx *exec_ctx, grpc_lb_policy *pol,
|
|||
while (pp != NULL) {
|
||||
pending_pick *next = pp->next;
|
||||
if (pp->target == target) {
|
||||
grpc_pollset_set_del_pollset(exec_ctx, p->base.interested_parties,
|
||||
pp->pollset);
|
||||
grpc_pollset_set_del_pollset_set(exec_ctx, p->base.interested_parties,
|
||||
pp->pollset_set);
|
||||
*target = NULL;
|
||||
grpc_exec_ctx_enqueue(exec_ctx, pp->on_complete, false, NULL);
|
||||
gpr_free(pp);
|
||||
|
|
@ -288,8 +288,8 @@ static void rr_cancel_picks(grpc_exec_ctx *exec_ctx, grpc_lb_policy *pol,
|
|||
pending_pick *next = pp->next;
|
||||
if ((pp->initial_metadata_flags & initial_metadata_flags_mask) ==
|
||||
initial_metadata_flags_eq) {
|
||||
grpc_pollset_set_del_pollset(exec_ctx, p->base.interested_parties,
|
||||
pp->pollset);
|
||||
grpc_pollset_set_del_pollset_set(exec_ctx, p->base.interested_parties,
|
||||
pp->pollset_set);
|
||||
*pp->target = NULL;
|
||||
grpc_exec_ctx_enqueue(exec_ctx, pp->on_complete, false, NULL);
|
||||
gpr_free(pp);
|
||||
|
|
@ -329,7 +329,8 @@ static void rr_exit_idle(grpc_exec_ctx *exec_ctx, grpc_lb_policy *pol) {
|
|||
}
|
||||
|
||||
static int rr_pick(grpc_exec_ctx *exec_ctx, grpc_lb_policy *pol,
|
||||
grpc_pollset *pollset, grpc_metadata_batch *initial_metadata,
|
||||
grpc_pollset_set *pollset_set,
|
||||
grpc_metadata_batch *initial_metadata,
|
||||
uint32_t initial_metadata_flags,
|
||||
grpc_connected_subchannel **target,
|
||||
grpc_closure *on_complete) {
|
||||
|
|
@ -352,10 +353,11 @@ static int rr_pick(grpc_exec_ctx *exec_ctx, grpc_lb_policy *pol,
|
|||
if (!p->started_picking) {
|
||||
start_picking(exec_ctx, p);
|
||||
}
|
||||
grpc_pollset_set_add_pollset(exec_ctx, p->base.interested_parties, pollset);
|
||||
grpc_pollset_set_add_pollset_set(exec_ctx, p->base.interested_parties,
|
||||
pollset_set);
|
||||
pp = gpr_malloc(sizeof(*pp));
|
||||
pp->next = p->pending_picks;
|
||||
pp->pollset = pollset;
|
||||
pp->pollset_set = pollset_set;
|
||||
pp->target = target;
|
||||
pp->on_complete = on_complete;
|
||||
pp->initial_metadata_flags = initial_metadata_flags;
|
||||
|
|
@ -404,8 +406,8 @@ static void rr_connectivity_changed(grpc_exec_ctx *exec_ctx, void *arg,
|
|||
"[RR CONN CHANGED] TARGET <-- SUBCHANNEL %p (NODE %p)",
|
||||
selected->subchannel, selected);
|
||||
}
|
||||
grpc_pollset_set_del_pollset(exec_ctx, p->base.interested_parties,
|
||||
pp->pollset);
|
||||
grpc_pollset_set_del_pollset_set(exec_ctx, p->base.interested_parties,
|
||||
pp->pollset_set);
|
||||
grpc_exec_ctx_enqueue(exec_ctx, pp->on_complete, true, NULL);
|
||||
gpr_free(pp);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1539,6 +1539,14 @@ static void set_pollset(grpc_exec_ctx *exec_ctx, grpc_transport *gt,
|
|||
unlock(exec_ctx, t);
|
||||
}
|
||||
|
||||
static void set_pollset_set(grpc_exec_ctx *exec_ctx, grpc_transport *gt,
|
||||
grpc_stream *gs, grpc_pollset_set *pollset_set) {
|
||||
grpc_chttp2_transport *t = (grpc_chttp2_transport *)gt;
|
||||
lock(t);
|
||||
add_to_pollset_set_locked(exec_ctx, t, pollset_set);
|
||||
unlock(exec_ctx, t);
|
||||
}
|
||||
|
||||
/*******************************************************************************
|
||||
* BYTE STREAM
|
||||
*/
|
||||
|
|
@ -1792,6 +1800,7 @@ static const grpc_transport_vtable vtable = {sizeof(grpc_chttp2_stream),
|
|||
"chttp2",
|
||||
init_stream,
|
||||
set_pollset,
|
||||
set_pollset_set,
|
||||
perform_stream_op,
|
||||
perform_transport_op,
|
||||
destroy_stream,
|
||||
|
|
|
|||
|
|
@ -189,29 +189,32 @@ void grpc_call_stack_init(grpc_exec_ctx *exec_ctx,
|
|||
}
|
||||
}
|
||||
|
||||
void grpc_call_stack_set_pollset(grpc_exec_ctx *exec_ctx,
|
||||
grpc_call_stack *call_stack,
|
||||
grpc_pollset *pollset) {
|
||||
void grpc_call_stack_set_pollset_or_pollset_set(
|
||||
grpc_exec_ctx *exec_ctx, grpc_call_stack *call_stack, grpc_pollset *pollset,
|
||||
grpc_pollset_set *or_pollset_set) {
|
||||
size_t count = call_stack->count;
|
||||
grpc_call_element *call_elems;
|
||||
char *user_data;
|
||||
size_t i;
|
||||
|
||||
GPR_ASSERT(!(pollset != NULL && or_pollset_set != NULL));
|
||||
GPR_ASSERT(pollset != NULL || or_pollset_set != NULL);
|
||||
call_elems = CALL_ELEMS_FROM_STACK(call_stack);
|
||||
user_data = ((char *)call_elems) +
|
||||
ROUND_UP_TO_ALIGNMENT_SIZE(count * sizeof(grpc_call_element));
|
||||
|
||||
/* init per-filter data */
|
||||
for (i = 0; i < count; i++) {
|
||||
call_elems[i].filter->set_pollset(exec_ctx, &call_elems[i], pollset);
|
||||
call_elems[i].filter->set_pollset_or_pollset_set(exec_ctx, &call_elems[i],
|
||||
pollset, or_pollset_set);
|
||||
user_data +=
|
||||
ROUND_UP_TO_ALIGNMENT_SIZE(call_elems[i].filter->sizeof_call_data);
|
||||
}
|
||||
}
|
||||
|
||||
void grpc_call_stack_ignore_set_pollset(grpc_exec_ctx *exec_ctx,
|
||||
grpc_call_element *elem,
|
||||
grpc_pollset *pollset) {}
|
||||
void grpc_call_stack_ignore_set_pollset_or_pollset_set(
|
||||
grpc_exec_ctx *exec_ctx, grpc_call_element *elem, grpc_pollset *pollset,
|
||||
grpc_pollset_set *or_pollset_set) {}
|
||||
|
||||
void grpc_call_stack_destroy(grpc_exec_ctx *exec_ctx, grpc_call_stack *stack) {
|
||||
grpc_call_element *elems = CALL_ELEMS_FROM_STACK(stack);
|
||||
|
|
|
|||
|
|
@ -101,8 +101,10 @@ typedef struct {
|
|||
argument. */
|
||||
void (*init_call_elem)(grpc_exec_ctx *exec_ctx, grpc_call_element *elem,
|
||||
grpc_call_element_args *args);
|
||||
void (*set_pollset)(grpc_exec_ctx *exec_ctx, grpc_call_element *elem,
|
||||
grpc_pollset *pollset);
|
||||
void (*set_pollset_or_pollset_set)(grpc_exec_ctx *exec_ctx,
|
||||
grpc_call_element *elem,
|
||||
grpc_pollset *pollset,
|
||||
grpc_pollset_set *or_pollset_set);
|
||||
/* Destroy per call data.
|
||||
The filter does not need to do any chaining */
|
||||
void (*destroy_call_elem)(grpc_exec_ctx *exec_ctx, grpc_call_element *elem);
|
||||
|
|
@ -197,10 +199,11 @@ void grpc_call_stack_init(grpc_exec_ctx *exec_ctx,
|
|||
grpc_call_context_element *context,
|
||||
const void *transport_server_data,
|
||||
grpc_call_stack *call_stack);
|
||||
/* Set a pollset for a call stack: must occur before the first op is started */
|
||||
void grpc_call_stack_set_pollset(grpc_exec_ctx *exec_ctx,
|
||||
grpc_call_stack *call_stack,
|
||||
grpc_pollset *pollset);
|
||||
/* Set a pollset or a pollset_set for a call stack: must occur before the first
|
||||
* op is started */
|
||||
void grpc_call_stack_set_pollset_or_pollset_set(
|
||||
grpc_exec_ctx *exec_ctx, grpc_call_stack *call_stack, grpc_pollset *pollset,
|
||||
grpc_pollset_set *or_pollset_set);
|
||||
|
||||
#ifdef GRPC_STREAM_REFCOUNT_DEBUG
|
||||
#define GRPC_CALL_STACK_REF(call_stack, reason) \
|
||||
|
|
@ -225,11 +228,12 @@ void grpc_call_stack_set_pollset(grpc_exec_ctx *exec_ctx,
|
|||
/* Destroy a call stack */
|
||||
void grpc_call_stack_destroy(grpc_exec_ctx *exec_ctx, grpc_call_stack *stack);
|
||||
|
||||
/* Ignore set pollset - used by filters to implement the set_pollset method
|
||||
if they don't care about pollsets at all. Does nothing. */
|
||||
void grpc_call_stack_ignore_set_pollset(grpc_exec_ctx *exec_ctx,
|
||||
grpc_call_element *elem,
|
||||
grpc_pollset *pollset);
|
||||
/* Ignore set pollset{_set} - used by filters to implement the
|
||||
* set_pollset_or_pollset_set method if they don't care about pollsets at all.
|
||||
* Does nothing. */
|
||||
void grpc_call_stack_ignore_set_pollset_or_pollset_set(
|
||||
grpc_exec_ctx *exec_ctx, grpc_call_element *elem, grpc_pollset *pollset,
|
||||
grpc_pollset_set *or_pollset_set);
|
||||
/* Call the next operation in a call stack */
|
||||
void grpc_call_next_op(grpc_exec_ctx *exec_ctx, grpc_call_element *elem,
|
||||
grpc_transport_stream_op *op);
|
||||
|
|
|
|||
|
|
@ -295,7 +295,7 @@ const grpc_channel_filter grpc_compress_filter = {
|
|||
grpc_channel_next_op,
|
||||
sizeof(call_data),
|
||||
init_call_elem,
|
||||
grpc_call_stack_ignore_set_pollset,
|
||||
grpc_call_stack_ignore_set_pollset_or_pollset_set,
|
||||
destroy_call_elem,
|
||||
sizeof(channel_data),
|
||||
init_channel_elem,
|
||||
|
|
|
|||
|
|
@ -93,12 +93,23 @@ static void init_call_elem(grpc_exec_ctx *exec_ctx, grpc_call_element *elem,
|
|||
GPR_ASSERT(r == 0);
|
||||
}
|
||||
|
||||
static void set_pollset(grpc_exec_ctx *exec_ctx, grpc_call_element *elem,
|
||||
grpc_pollset *pollset) {
|
||||
static void set_pollset_or_pollset_set(grpc_exec_ctx *exec_ctx,
|
||||
grpc_call_element *elem,
|
||||
grpc_pollset *pollset,
|
||||
grpc_pollset_set *or_pollset_set) {
|
||||
GPR_ASSERT(!(pollset != NULL && or_pollset_set != NULL));
|
||||
GPR_ASSERT(pollset != NULL || or_pollset_set != NULL);
|
||||
|
||||
call_data *calld = elem->call_data;
|
||||
channel_data *chand = elem->channel_data;
|
||||
grpc_transport_set_pollset(exec_ctx, chand->transport,
|
||||
TRANSPORT_STREAM_FROM_CALL_DATA(calld), pollset);
|
||||
if (pollset != NULL) {
|
||||
grpc_transport_set_pollset(exec_ctx, chand->transport,
|
||||
TRANSPORT_STREAM_FROM_CALL_DATA(calld), pollset);
|
||||
} else if (or_pollset_set != NULL) {
|
||||
grpc_transport_set_pollset_set(exec_ctx, chand->transport,
|
||||
TRANSPORT_STREAM_FROM_CALL_DATA(calld),
|
||||
or_pollset_set);
|
||||
}
|
||||
}
|
||||
|
||||
/* Destructor for call_data */
|
||||
|
|
@ -136,7 +147,7 @@ static const grpc_channel_filter connected_channel_filter = {
|
|||
con_start_transport_op,
|
||||
sizeof(call_data),
|
||||
init_call_elem,
|
||||
set_pollset,
|
||||
set_pollset_or_pollset_set,
|
||||
destroy_call_elem,
|
||||
sizeof(channel_data),
|
||||
init_channel_elem,
|
||||
|
|
|
|||
|
|
@ -250,7 +250,7 @@ const grpc_channel_filter grpc_http_client_filter = {
|
|||
grpc_channel_next_op,
|
||||
sizeof(call_data),
|
||||
init_call_elem,
|
||||
grpc_call_stack_ignore_set_pollset,
|
||||
grpc_call_stack_ignore_set_pollset_or_pollset_set,
|
||||
destroy_call_elem,
|
||||
sizeof(channel_data),
|
||||
init_channel_elem,
|
||||
|
|
|
|||
|
|
@ -239,7 +239,7 @@ const grpc_channel_filter grpc_http_server_filter = {
|
|||
grpc_channel_next_op,
|
||||
sizeof(call_data),
|
||||
init_call_elem,
|
||||
grpc_call_stack_ignore_set_pollset,
|
||||
grpc_call_stack_ignore_set_pollset_or_pollset_set,
|
||||
destroy_call_elem,
|
||||
sizeof(channel_data),
|
||||
init_channel_elem,
|
||||
|
|
|
|||
|
|
@ -62,7 +62,7 @@ typedef struct {
|
|||
grpc_httpcli_response_cb on_response;
|
||||
void *user_data;
|
||||
grpc_httpcli_context *context;
|
||||
grpc_pollset *pollset;
|
||||
grpc_pollset_set *pollset_set;
|
||||
grpc_iomgr_object iomgr_obj;
|
||||
gpr_slice_buffer incoming;
|
||||
gpr_slice_buffer outgoing;
|
||||
|
|
@ -97,8 +97,8 @@ static void next_address(grpc_exec_ctx *exec_ctx, internal_request *req);
|
|||
|
||||
static void finish(grpc_exec_ctx *exec_ctx, internal_request *req,
|
||||
int success) {
|
||||
grpc_pollset_set_del_pollset(exec_ctx, req->context->pollset_set,
|
||||
req->pollset);
|
||||
grpc_pollset_set_del_pollset_set(exec_ctx, req->context->pollset_set,
|
||||
req->pollset_set);
|
||||
req->on_response(exec_ctx, req->user_data,
|
||||
success ? &req->parser.http.response : NULL);
|
||||
grpc_http_parser_destroy(&req->parser);
|
||||
|
|
@ -222,7 +222,7 @@ static void on_resolved(grpc_exec_ctx *exec_ctx, void *arg,
|
|||
|
||||
static void internal_request_begin(
|
||||
grpc_exec_ctx *exec_ctx, grpc_httpcli_context *context,
|
||||
grpc_pollset *pollset, const grpc_httpcli_request *request,
|
||||
grpc_pollset_set *pollset_set, const grpc_httpcli_request *request,
|
||||
gpr_timespec deadline, grpc_httpcli_response_cb on_response,
|
||||
void *user_data, const char *name, gpr_slice request_text) {
|
||||
internal_request *req = gpr_malloc(sizeof(internal_request));
|
||||
|
|
@ -235,7 +235,7 @@ static void internal_request_begin(
|
|||
req->handshaker =
|
||||
request->handshaker ? request->handshaker : &grpc_httpcli_plaintext;
|
||||
req->context = context;
|
||||
req->pollset = pollset;
|
||||
req->pollset_set = pollset_set;
|
||||
grpc_closure_init(&req->on_read, on_read, req);
|
||||
grpc_closure_init(&req->done_write, done_write, req);
|
||||
gpr_slice_buffer_init(&req->incoming);
|
||||
|
|
@ -244,14 +244,14 @@ static void internal_request_begin(
|
|||
req->host = gpr_strdup(request->host);
|
||||
req->ssl_host_override = gpr_strdup(request->ssl_host_override);
|
||||
|
||||
grpc_pollset_set_add_pollset(exec_ctx, req->context->pollset_set,
|
||||
req->pollset);
|
||||
grpc_pollset_set_add_pollset_set(exec_ctx, req->context->pollset_set,
|
||||
req->pollset_set);
|
||||
grpc_resolve_address(request->host, req->handshaker->default_port,
|
||||
on_resolved, req);
|
||||
}
|
||||
|
||||
void grpc_httpcli_get(grpc_exec_ctx *exec_ctx, grpc_httpcli_context *context,
|
||||
grpc_pollset *pollset,
|
||||
grpc_pollset_set *pollset_set,
|
||||
const grpc_httpcli_request *request,
|
||||
gpr_timespec deadline,
|
||||
grpc_httpcli_response_cb on_response, void *user_data) {
|
||||
|
|
@ -261,14 +261,14 @@ void grpc_httpcli_get(grpc_exec_ctx *exec_ctx, grpc_httpcli_context *context,
|
|||
return;
|
||||
}
|
||||
gpr_asprintf(&name, "HTTP:GET:%s:%s", request->host, request->http.path);
|
||||
internal_request_begin(exec_ctx, context, pollset, request, deadline,
|
||||
internal_request_begin(exec_ctx, context, pollset_set, request, deadline,
|
||||
on_response, user_data, name,
|
||||
grpc_httpcli_format_get_request(request));
|
||||
gpr_free(name);
|
||||
}
|
||||
|
||||
void grpc_httpcli_post(grpc_exec_ctx *exec_ctx, grpc_httpcli_context *context,
|
||||
grpc_pollset *pollset,
|
||||
grpc_pollset_set *pollset_set,
|
||||
const grpc_httpcli_request *request,
|
||||
const char *body_bytes, size_t body_size,
|
||||
gpr_timespec deadline,
|
||||
|
|
@ -281,7 +281,7 @@ void grpc_httpcli_post(grpc_exec_ctx *exec_ctx, grpc_httpcli_context *context,
|
|||
}
|
||||
gpr_asprintf(&name, "HTTP:POST:%s:%s", request->host, request->http.path);
|
||||
internal_request_begin(
|
||||
exec_ctx, context, pollset, request, deadline, on_response, user_data,
|
||||
exec_ctx, context, pollset_set, request, deadline, on_response, user_data,
|
||||
name, grpc_httpcli_format_post_request(request, body_bytes, body_size));
|
||||
gpr_free(name);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -100,7 +100,7 @@ void grpc_httpcli_context_destroy(grpc_httpcli_context *context);
|
|||
'on_response' is a callback to report results to (and 'user_data' is a user
|
||||
supplied pointer to pass to said call) */
|
||||
void grpc_httpcli_get(grpc_exec_ctx *exec_ctx, grpc_httpcli_context *context,
|
||||
grpc_pollset *pollset,
|
||||
grpc_pollset_set *pollset_set,
|
||||
const grpc_httpcli_request *request,
|
||||
gpr_timespec deadline,
|
||||
grpc_httpcli_response_cb on_response, void *user_data);
|
||||
|
|
@ -121,7 +121,7 @@ void grpc_httpcli_get(grpc_exec_ctx *exec_ctx, grpc_httpcli_context *context,
|
|||
supplied pointer to pass to said call)
|
||||
Does not support ?var1=val1&var2=val2 in the path. */
|
||||
void grpc_httpcli_post(grpc_exec_ctx *exec_ctx, grpc_httpcli_context *context,
|
||||
grpc_pollset *pollset,
|
||||
grpc_pollset_set *pollset_set,
|
||||
const grpc_httpcli_request *request,
|
||||
const char *body_bytes, size_t body_size,
|
||||
gpr_timespec deadline,
|
||||
|
|
|
|||
|
|
@ -54,11 +54,11 @@ typedef struct {
|
|||
grpc_call_credentials *creds;
|
||||
grpc_mdstr *host;
|
||||
grpc_mdstr *method;
|
||||
/* pollset bound to this call; if we need to make external
|
||||
network requests, they should be done under this pollset
|
||||
so that work can progress when this call wants work to
|
||||
progress */
|
||||
grpc_pollset *pollset;
|
||||
/* pollset_set bound to this call; if we need to make external
|
||||
network requests, they should be done under a pollset added to this
|
||||
pollset_set so that work can progress when this call wants work to progress
|
||||
*/
|
||||
grpc_pollset_set *pollset_set;
|
||||
grpc_transport_stream_op op;
|
||||
uint8_t security_context_set;
|
||||
grpc_linked_mdelem md_links[MAX_CREDENTIALS_METADATA_COUNT];
|
||||
|
|
@ -184,9 +184,9 @@ static void send_security_metadata(grpc_exec_ctx *exec_ctx,
|
|||
build_auth_metadata_context(&chand->security_connector->base,
|
||||
chand->auth_context, calld);
|
||||
calld->op = *op; /* Copy op (originates from the caller's stack). */
|
||||
GPR_ASSERT(calld->pollset);
|
||||
GPR_ASSERT(calld->pollset_set);
|
||||
grpc_call_credentials_get_request_metadata(
|
||||
exec_ctx, calld->creds, calld->pollset, calld->auth_md_context,
|
||||
exec_ctx, calld->creds, calld->pollset_set, calld->auth_md_context,
|
||||
on_credentials_metadata, elem);
|
||||
}
|
||||
|
||||
|
|
@ -268,12 +268,23 @@ static void init_call_elem(grpc_exec_ctx *exec_ctx, grpc_call_element *elem,
|
|||
grpc_call_element_args *args) {
|
||||
call_data *calld = elem->call_data;
|
||||
memset(calld, 0, sizeof(*calld));
|
||||
calld->pollset_set = grpc_pollset_set_create();
|
||||
}
|
||||
|
||||
static void set_pollset(grpc_exec_ctx *exec_ctx, grpc_call_element *elem,
|
||||
grpc_pollset *pollset) {
|
||||
static void set_pollset_or_pollset_set(grpc_exec_ctx *exec_ctx,
|
||||
grpc_call_element *elem,
|
||||
grpc_pollset *pollset,
|
||||
grpc_pollset_set *or_pollset_set) {
|
||||
GPR_ASSERT(!(pollset != NULL && or_pollset_set != NULL));
|
||||
GPR_ASSERT(pollset != NULL || or_pollset_set != NULL);
|
||||
|
||||
call_data *calld = elem->call_data;
|
||||
calld->pollset = pollset;
|
||||
if (pollset != NULL) {
|
||||
grpc_pollset_set_add_pollset(exec_ctx, calld->pollset_set, pollset);
|
||||
} else if (or_pollset_set != NULL) {
|
||||
grpc_pollset_set_add_pollset_set(exec_ctx, calld->pollset_set,
|
||||
or_pollset_set);
|
||||
}
|
||||
}
|
||||
|
||||
/* Destructor for call_data */
|
||||
|
|
@ -288,6 +299,7 @@ static void destroy_call_elem(grpc_exec_ctx *exec_ctx,
|
|||
GRPC_MDSTR_UNREF(calld->method);
|
||||
}
|
||||
reset_auth_metadata_context(&calld->auth_md_context);
|
||||
grpc_pollset_set_destroy(calld->pollset_set);
|
||||
}
|
||||
|
||||
/* Constructor for channel_data */
|
||||
|
|
@ -329,8 +341,14 @@ static void destroy_channel_elem(grpc_exec_ctx *exec_ctx,
|
|||
GRPC_AUTH_CONTEXT_UNREF(chand->auth_context, "client_auth_filter");
|
||||
}
|
||||
|
||||
const grpc_channel_filter grpc_client_auth_filter = {
|
||||
auth_start_transport_op, grpc_channel_next_op, sizeof(call_data),
|
||||
init_call_elem, set_pollset, destroy_call_elem,
|
||||
sizeof(channel_data), init_channel_elem, destroy_channel_elem,
|
||||
grpc_call_next_get_peer, "client-auth"};
|
||||
const grpc_channel_filter grpc_client_auth_filter = {auth_start_transport_op,
|
||||
grpc_channel_next_op,
|
||||
sizeof(call_data),
|
||||
init_call_elem,
|
||||
set_pollset_or_pollset_set,
|
||||
destroy_call_elem,
|
||||
sizeof(channel_data),
|
||||
init_channel_elem,
|
||||
destroy_channel_elem,
|
||||
grpc_call_next_get_peer,
|
||||
"client-auth"};
|
||||
|
|
|
|||
|
|
@ -118,7 +118,7 @@ void grpc_call_credentials_release(grpc_call_credentials *creds) {
|
|||
|
||||
void grpc_call_credentials_get_request_metadata(
|
||||
grpc_exec_ctx *exec_ctx, grpc_call_credentials *creds,
|
||||
grpc_pollset *pollset, grpc_auth_metadata_context context,
|
||||
grpc_pollset_set *pollset_set, grpc_auth_metadata_context context,
|
||||
grpc_credentials_metadata_cb cb, void *user_data) {
|
||||
if (creds == NULL || creds->vtable->get_request_metadata == NULL) {
|
||||
if (cb != NULL) {
|
||||
|
|
@ -126,7 +126,7 @@ void grpc_call_credentials_get_request_metadata(
|
|||
}
|
||||
return;
|
||||
}
|
||||
creds->vtable->get_request_metadata(exec_ctx, creds, pollset, context, cb,
|
||||
creds->vtable->get_request_metadata(exec_ctx, creds, pollset_set, context, cb,
|
||||
user_data);
|
||||
}
|
||||
|
||||
|
|
@ -433,7 +433,7 @@ static void jwt_destruct(grpc_call_credentials *creds) {
|
|||
|
||||
static void jwt_get_request_metadata(grpc_exec_ctx *exec_ctx,
|
||||
grpc_call_credentials *creds,
|
||||
grpc_pollset *pollset,
|
||||
grpc_pollset_set *pollset_set,
|
||||
grpc_auth_metadata_context context,
|
||||
grpc_credentials_metadata_cb cb,
|
||||
void *user_data) {
|
||||
|
|
@ -656,7 +656,7 @@ static void on_oauth2_token_fetcher_http_response(
|
|||
|
||||
static void oauth2_token_fetcher_get_request_metadata(
|
||||
grpc_exec_ctx *exec_ctx, grpc_call_credentials *creds,
|
||||
grpc_pollset *pollset, grpc_auth_metadata_context context,
|
||||
grpc_pollset_set *pollset_set, grpc_auth_metadata_context context,
|
||||
grpc_credentials_metadata_cb cb, void *user_data) {
|
||||
grpc_oauth2_token_fetcher_credentials *c =
|
||||
(grpc_oauth2_token_fetcher_credentials *)creds;
|
||||
|
|
@ -682,7 +682,7 @@ static void oauth2_token_fetcher_get_request_metadata(
|
|||
c->fetch_func(
|
||||
exec_ctx,
|
||||
grpc_credentials_metadata_request_create(creds, cb, user_data),
|
||||
&c->httpcli_context, pollset, on_oauth2_token_fetcher_http_response,
|
||||
&c->httpcli_context, pollset_set, on_oauth2_token_fetcher_http_response,
|
||||
gpr_time_add(gpr_now(GPR_CLOCK_REALTIME), refresh_threshold));
|
||||
}
|
||||
}
|
||||
|
|
@ -705,7 +705,7 @@ static grpc_call_credentials_vtable compute_engine_vtable = {
|
|||
|
||||
static void compute_engine_fetch_oauth2(
|
||||
grpc_exec_ctx *exec_ctx, grpc_credentials_metadata_request *metadata_req,
|
||||
grpc_httpcli_context *httpcli_context, grpc_pollset *pollset,
|
||||
grpc_httpcli_context *httpcli_context, grpc_pollset_set *pollset_set,
|
||||
grpc_httpcli_response_cb response_cb, gpr_timespec deadline) {
|
||||
grpc_http_header header = {"Metadata-Flavor", "Google"};
|
||||
grpc_httpcli_request request;
|
||||
|
|
@ -714,7 +714,7 @@ static void compute_engine_fetch_oauth2(
|
|||
request.http.path = GRPC_COMPUTE_ENGINE_METADATA_TOKEN_PATH;
|
||||
request.http.hdr_count = 1;
|
||||
request.http.hdrs = &header;
|
||||
grpc_httpcli_get(exec_ctx, httpcli_context, pollset, &request, deadline,
|
||||
grpc_httpcli_get(exec_ctx, httpcli_context, pollset_set, &request, deadline,
|
||||
response_cb, metadata_req);
|
||||
}
|
||||
|
||||
|
|
@ -744,7 +744,7 @@ static grpc_call_credentials_vtable refresh_token_vtable = {
|
|||
|
||||
static void refresh_token_fetch_oauth2(
|
||||
grpc_exec_ctx *exec_ctx, grpc_credentials_metadata_request *metadata_req,
|
||||
grpc_httpcli_context *httpcli_context, grpc_pollset *pollset,
|
||||
grpc_httpcli_context *httpcli_context, grpc_pollset_set *pollset_set,
|
||||
grpc_httpcli_response_cb response_cb, gpr_timespec deadline) {
|
||||
grpc_google_refresh_token_credentials *c =
|
||||
(grpc_google_refresh_token_credentials *)metadata_req->creds;
|
||||
|
|
@ -761,7 +761,7 @@ static void refresh_token_fetch_oauth2(
|
|||
request.http.hdr_count = 1;
|
||||
request.http.hdrs = &header;
|
||||
request.handshaker = &grpc_httpcli_ssl;
|
||||
grpc_httpcli_post(exec_ctx, httpcli_context, pollset, &request, body,
|
||||
grpc_httpcli_post(exec_ctx, httpcli_context, pollset_set, &request, body,
|
||||
strlen(body), deadline, response_cb, metadata_req);
|
||||
gpr_free(body);
|
||||
}
|
||||
|
|
@ -812,7 +812,7 @@ static void on_simulated_token_fetch_done(grpc_exec_ctx *exec_ctx,
|
|||
|
||||
static void md_only_test_get_request_metadata(
|
||||
grpc_exec_ctx *exec_ctx, grpc_call_credentials *creds,
|
||||
grpc_pollset *pollset, grpc_auth_metadata_context context,
|
||||
grpc_pollset_set *pollset_set, grpc_auth_metadata_context context,
|
||||
grpc_credentials_metadata_cb cb, void *user_data) {
|
||||
grpc_md_only_test_credentials *c = (grpc_md_only_test_credentials *)creds;
|
||||
|
||||
|
|
@ -852,7 +852,7 @@ static void access_token_destruct(grpc_call_credentials *creds) {
|
|||
|
||||
static void access_token_get_request_metadata(
|
||||
grpc_exec_ctx *exec_ctx, grpc_call_credentials *creds,
|
||||
grpc_pollset *pollset, grpc_auth_metadata_context context,
|
||||
grpc_pollset_set *pollset_set, grpc_auth_metadata_context context,
|
||||
grpc_credentials_metadata_cb cb, void *user_data) {
|
||||
grpc_access_token_credentials *c = (grpc_access_token_credentials *)creds;
|
||||
cb(exec_ctx, user_data, c->access_token_md->entries, 1, GRPC_CREDENTIALS_OK);
|
||||
|
|
@ -936,7 +936,7 @@ typedef struct {
|
|||
grpc_credentials_md_store *md_elems;
|
||||
grpc_auth_metadata_context auth_md_context;
|
||||
void *user_data;
|
||||
grpc_pollset *pollset;
|
||||
grpc_pollset_set *pollset_set;
|
||||
grpc_credentials_metadata_cb cb;
|
||||
} grpc_composite_call_credentials_metadata_context;
|
||||
|
||||
|
|
@ -980,7 +980,7 @@ static void composite_call_metadata_cb(grpc_exec_ctx *exec_ctx, void *user_data,
|
|||
grpc_call_credentials *inner_creds =
|
||||
ctx->composite_creds->inner.creds_array[ctx->creds_index++];
|
||||
grpc_call_credentials_get_request_metadata(
|
||||
exec_ctx, inner_creds, ctx->pollset, ctx->auth_md_context,
|
||||
exec_ctx, inner_creds, ctx->pollset_set, ctx->auth_md_context,
|
||||
composite_call_metadata_cb, ctx);
|
||||
return;
|
||||
}
|
||||
|
|
@ -993,7 +993,7 @@ static void composite_call_metadata_cb(grpc_exec_ctx *exec_ctx, void *user_data,
|
|||
|
||||
static void composite_call_get_request_metadata(
|
||||
grpc_exec_ctx *exec_ctx, grpc_call_credentials *creds,
|
||||
grpc_pollset *pollset, grpc_auth_metadata_context auth_md_context,
|
||||
grpc_pollset_set *pollset_set, grpc_auth_metadata_context auth_md_context,
|
||||
grpc_credentials_metadata_cb cb, void *user_data) {
|
||||
grpc_composite_call_credentials *c = (grpc_composite_call_credentials *)creds;
|
||||
grpc_composite_call_credentials_metadata_context *ctx;
|
||||
|
|
@ -1004,10 +1004,10 @@ static void composite_call_get_request_metadata(
|
|||
ctx->user_data = user_data;
|
||||
ctx->cb = cb;
|
||||
ctx->composite_creds = c;
|
||||
ctx->pollset = pollset;
|
||||
ctx->pollset_set = pollset_set;
|
||||
ctx->md_elems = grpc_credentials_md_store_create(c->inner.num_creds);
|
||||
grpc_call_credentials_get_request_metadata(
|
||||
exec_ctx, c->inner.creds_array[ctx->creds_index++], pollset,
|
||||
exec_ctx, c->inner.creds_array[ctx->creds_index++], pollset_set,
|
||||
auth_md_context, composite_call_metadata_cb, ctx);
|
||||
}
|
||||
|
||||
|
|
@ -1101,7 +1101,7 @@ static void iam_destruct(grpc_call_credentials *creds) {
|
|||
|
||||
static void iam_get_request_metadata(grpc_exec_ctx *exec_ctx,
|
||||
grpc_call_credentials *creds,
|
||||
grpc_pollset *pollset,
|
||||
grpc_pollset_set *pollset_set,
|
||||
grpc_auth_metadata_context context,
|
||||
grpc_credentials_metadata_cb cb,
|
||||
void *user_data) {
|
||||
|
|
@ -1190,7 +1190,7 @@ static void plugin_md_request_metadata_ready(void *request,
|
|||
|
||||
static void plugin_get_request_metadata(grpc_exec_ctx *exec_ctx,
|
||||
grpc_call_credentials *creds,
|
||||
grpc_pollset *pollset,
|
||||
grpc_pollset_set *pollset_set,
|
||||
grpc_auth_metadata_context context,
|
||||
grpc_credentials_metadata_cb cb,
|
||||
void *user_data) {
|
||||
|
|
|
|||
|
|
@ -169,7 +169,8 @@ typedef void (*grpc_credentials_metadata_cb)(grpc_exec_ctx *exec_ctx,
|
|||
typedef struct {
|
||||
void (*destruct)(grpc_call_credentials *c);
|
||||
void (*get_request_metadata)(grpc_exec_ctx *exec_ctx,
|
||||
grpc_call_credentials *c, grpc_pollset *pollset,
|
||||
grpc_call_credentials *c,
|
||||
grpc_pollset_set *pollset_set,
|
||||
grpc_auth_metadata_context context,
|
||||
grpc_credentials_metadata_cb cb,
|
||||
void *user_data);
|
||||
|
|
@ -185,7 +186,7 @@ grpc_call_credentials *grpc_call_credentials_ref(grpc_call_credentials *creds);
|
|||
void grpc_call_credentials_unref(grpc_call_credentials *creds);
|
||||
void grpc_call_credentials_get_request_metadata(
|
||||
grpc_exec_ctx *exec_ctx, grpc_call_credentials *creds,
|
||||
grpc_pollset *pollset, grpc_auth_metadata_context context,
|
||||
grpc_pollset_set *pollset_set, grpc_auth_metadata_context context,
|
||||
grpc_credentials_metadata_cb cb, void *user_data);
|
||||
|
||||
typedef struct {
|
||||
|
|
@ -317,7 +318,7 @@ typedef struct grpc_credentials_metadata_request
|
|||
typedef void (*grpc_fetch_oauth2_func)(grpc_exec_ctx *exec_ctx,
|
||||
grpc_credentials_metadata_request *req,
|
||||
grpc_httpcli_context *http_context,
|
||||
grpc_pollset *pollset,
|
||||
grpc_pollset_set *pollset_set,
|
||||
grpc_httpcli_response_cb response_cb,
|
||||
gpr_timespec deadline);
|
||||
|
||||
|
|
|
|||
|
|
@ -61,6 +61,7 @@ static void init_default_credentials(void) { gpr_mu_init(&g_state_mu); }
|
|||
|
||||
typedef struct {
|
||||
grpc_pollset *pollset;
|
||||
grpc_pollset_set *pollset_set;
|
||||
int is_done;
|
||||
int success;
|
||||
} compute_engine_detector;
|
||||
|
|
@ -105,6 +106,9 @@ static int is_stack_running_on_compute_engine(void) {
|
|||
|
||||
detector.pollset = gpr_malloc(grpc_pollset_size());
|
||||
grpc_pollset_init(detector.pollset, &g_polling_mu);
|
||||
detector.pollset_set = grpc_pollset_set_create();
|
||||
grpc_pollset_set_add_pollset(&exec_ctx, detector.pollset_set,
|
||||
detector.pollset);
|
||||
detector.is_done = 0;
|
||||
detector.success = 0;
|
||||
|
||||
|
|
@ -115,7 +119,7 @@ static int is_stack_running_on_compute_engine(void) {
|
|||
grpc_httpcli_context_init(&context);
|
||||
|
||||
grpc_httpcli_get(
|
||||
&exec_ctx, &context, detector.pollset, &request,
|
||||
&exec_ctx, &context, detector.pollset_set, &request,
|
||||
gpr_time_add(gpr_now(GPR_CLOCK_REALTIME), max_detection_delay),
|
||||
on_compute_engine_detection_http_response, &detector);
|
||||
|
||||
|
|
@ -135,6 +139,7 @@ static int is_stack_running_on_compute_engine(void) {
|
|||
grpc_httpcli_context_destroy(&context);
|
||||
grpc_closure_init(&destroy_closure, destroy_pollset, detector.pollset);
|
||||
grpc_pollset_shutdown(&exec_ctx, detector.pollset, &destroy_closure);
|
||||
grpc_pollset_set_destroy(detector.pollset_set);
|
||||
grpc_exec_ctx_finish(&exec_ctx);
|
||||
g_polling_mu = NULL;
|
||||
|
||||
|
|
|
|||
|
|
@ -321,7 +321,7 @@ grpc_jwt_verifier_status grpc_jwt_claims_check(const grpc_jwt_claims *claims,
|
|||
|
||||
typedef struct {
|
||||
grpc_jwt_verifier *verifier;
|
||||
grpc_pollset *pollset;
|
||||
grpc_pollset_set *pollset_set;
|
||||
jose_header *header;
|
||||
grpc_jwt_claims *claims;
|
||||
char *audience;
|
||||
|
|
@ -337,10 +337,12 @@ static verifier_cb_ctx *verifier_cb_ctx_create(
|
|||
grpc_jwt_claims *claims, const char *audience, gpr_slice signature,
|
||||
const char *signed_jwt, size_t signed_jwt_len, void *user_data,
|
||||
grpc_jwt_verification_done_cb cb) {
|
||||
grpc_exec_ctx exec_ctx = GRPC_EXEC_CTX_INIT;
|
||||
verifier_cb_ctx *ctx = gpr_malloc(sizeof(verifier_cb_ctx));
|
||||
memset(ctx, 0, sizeof(verifier_cb_ctx));
|
||||
ctx->verifier = verifier;
|
||||
ctx->pollset = pollset;
|
||||
ctx->pollset_set = grpc_pollset_set_create();
|
||||
grpc_pollset_set_add_pollset(&exec_ctx, ctx->pollset_set, pollset);
|
||||
ctx->header = header;
|
||||
ctx->audience = gpr_strdup(audience);
|
||||
ctx->claims = claims;
|
||||
|
|
@ -348,6 +350,7 @@ static verifier_cb_ctx *verifier_cb_ctx_create(
|
|||
ctx->signed_data = gpr_slice_from_copied_buffer(signed_jwt, signed_jwt_len);
|
||||
ctx->user_data = user_data;
|
||||
ctx->user_cb = cb;
|
||||
grpc_exec_ctx_finish(&exec_ctx);
|
||||
return ctx;
|
||||
}
|
||||
|
||||
|
|
@ -357,6 +360,7 @@ void verifier_cb_ctx_destroy(verifier_cb_ctx *ctx) {
|
|||
gpr_slice_unref(ctx->signature);
|
||||
gpr_slice_unref(ctx->signed_data);
|
||||
jose_header_destroy(ctx->header);
|
||||
grpc_pollset_set_destroy(ctx->pollset_set);
|
||||
/* TODO: see what to do with claims... */
|
||||
gpr_free(ctx);
|
||||
}
|
||||
|
|
@ -642,7 +646,7 @@ static void on_openid_config_retrieved(grpc_exec_ctx *exec_ctx, void *user_data,
|
|||
*(req.host + (req.http.path - jwks_uri)) = '\0';
|
||||
}
|
||||
grpc_httpcli_get(
|
||||
exec_ctx, &ctx->verifier->http_ctx, ctx->pollset, &req,
|
||||
exec_ctx, &ctx->verifier->http_ctx, ctx->pollset_set, &req,
|
||||
gpr_time_add(gpr_now(GPR_CLOCK_REALTIME), grpc_jwt_verifier_max_delay),
|
||||
on_keys_retrieved, ctx);
|
||||
grpc_json_destroy(json);
|
||||
|
|
@ -745,7 +749,7 @@ static void retrieve_key_and_verify(grpc_exec_ctx *exec_ctx,
|
|||
}
|
||||
|
||||
grpc_httpcli_get(
|
||||
exec_ctx, &ctx->verifier->http_ctx, ctx->pollset, &req,
|
||||
exec_ctx, &ctx->verifier->http_ctx, ctx->pollset_set, &req,
|
||||
gpr_time_add(gpr_now(GPR_CLOCK_REALTIME), grpc_jwt_verifier_max_delay),
|
||||
http_cb, ctx);
|
||||
gpr_free(req.host);
|
||||
|
|
|
|||
|
|
@ -220,9 +220,6 @@ static void init_call_elem(grpc_exec_ctx *exec_ctx, grpc_call_element *elem,
|
|||
grpc_server_security_context_destroy;
|
||||
}
|
||||
|
||||
static void set_pollset(grpc_exec_ctx *exec_ctx, grpc_call_element *elem,
|
||||
grpc_pollset *pollset) {}
|
||||
|
||||
/* Destructor for call_data */
|
||||
static void destroy_call_elem(grpc_exec_ctx *exec_ctx,
|
||||
grpc_call_element *elem) {}
|
||||
|
|
@ -258,7 +255,14 @@ static void destroy_channel_elem(grpc_exec_ctx *exec_ctx,
|
|||
}
|
||||
|
||||
const grpc_channel_filter grpc_server_auth_filter = {
|
||||
auth_start_transport_op, grpc_channel_next_op, sizeof(call_data),
|
||||
init_call_elem, set_pollset, destroy_call_elem,
|
||||
sizeof(channel_data), init_channel_elem, destroy_channel_elem,
|
||||
grpc_call_next_get_peer, "server-auth"};
|
||||
auth_start_transport_op,
|
||||
grpc_channel_next_op,
|
||||
sizeof(call_data),
|
||||
init_call_elem,
|
||||
grpc_call_stack_ignore_set_pollset_or_pollset_set,
|
||||
destroy_call_elem,
|
||||
sizeof(channel_data),
|
||||
init_channel_elem,
|
||||
destroy_channel_elem,
|
||||
grpc_call_next_get_peer,
|
||||
"server-auth"};
|
||||
|
|
|
|||
|
|
@ -135,6 +135,7 @@ typedef struct batch_control {
|
|||
|
||||
struct grpc_call {
|
||||
grpc_completion_queue *cq;
|
||||
grpc_pollset_set *pollset_set;
|
||||
grpc_channel *channel;
|
||||
grpc_call *parent;
|
||||
grpc_call *first_child;
|
||||
|
|
@ -245,13 +246,11 @@ static void destroy_call(grpc_exec_ctx *exec_ctx, void *call_stack,
|
|||
static void receiving_slice_ready(grpc_exec_ctx *exec_ctx, void *bctlp,
|
||||
bool success);
|
||||
|
||||
grpc_call *grpc_call_create(grpc_channel *channel, grpc_call *parent_call,
|
||||
uint32_t propagation_mask,
|
||||
grpc_completion_queue *cq,
|
||||
const void *server_transport_data,
|
||||
grpc_mdelem **add_initial_metadata,
|
||||
size_t add_initial_metadata_count,
|
||||
gpr_timespec send_deadline) {
|
||||
grpc_call *grpc_call_create(
|
||||
grpc_channel *channel, grpc_call *parent_call, uint32_t propagation_mask,
|
||||
grpc_completion_queue *cq, grpc_pollset_set *or_pollset_set,
|
||||
const void *server_transport_data, grpc_mdelem **add_initial_metadata,
|
||||
size_t add_initial_metadata_count, gpr_timespec send_deadline) {
|
||||
size_t i, j;
|
||||
grpc_channel_stack *channel_stack = grpc_channel_get_channel_stack(channel);
|
||||
grpc_exec_ctx exec_ctx = GRPC_EXEC_CTX_INIT;
|
||||
|
|
@ -262,6 +261,12 @@ grpc_call *grpc_call_create(grpc_channel *channel, grpc_call *parent_call,
|
|||
gpr_mu_init(&call->mu);
|
||||
call->channel = channel;
|
||||
call->cq = cq;
|
||||
if (cq != NULL && or_pollset_set != NULL) {
|
||||
gpr_log(GPR_ERROR,
|
||||
"Only one of 'cq' and 'or_pollset_set' should be non-NULL.");
|
||||
abort();
|
||||
}
|
||||
call->pollset_set = or_pollset_set;
|
||||
call->parent = parent_call;
|
||||
call->is_client = server_transport_data == NULL;
|
||||
if (call->is_client) {
|
||||
|
|
@ -287,8 +292,13 @@ grpc_call *grpc_call_create(grpc_channel *channel, grpc_call *parent_call,
|
|||
CALL_STACK_FROM_CALL(call));
|
||||
if (cq != NULL) {
|
||||
GRPC_CQ_INTERNAL_REF(cq, "bind");
|
||||
grpc_call_stack_set_pollset(&exec_ctx, CALL_STACK_FROM_CALL(call),
|
||||
grpc_cq_pollset(cq));
|
||||
grpc_call_stack_set_pollset_or_pollset_set(
|
||||
&exec_ctx, CALL_STACK_FROM_CALL(call), grpc_cq_pollset(cq), NULL);
|
||||
}
|
||||
if (or_pollset_set != NULL) {
|
||||
GPR_ASSERT(cq == NULL);
|
||||
grpc_call_stack_set_pollset_or_pollset_set(
|
||||
&exec_ctx, CALL_STACK_FROM_CALL(call), NULL, or_pollset_set);
|
||||
}
|
||||
if (parent_call != NULL) {
|
||||
GRPC_CALL_INTERNAL_REF(parent_call, "child");
|
||||
|
|
@ -343,10 +353,11 @@ grpc_call *grpc_call_create(grpc_channel *channel, grpc_call *parent_call,
|
|||
void grpc_call_set_completion_queue(grpc_exec_ctx *exec_ctx, grpc_call *call,
|
||||
grpc_completion_queue *cq) {
|
||||
GPR_ASSERT(cq);
|
||||
GPR_ASSERT(call->pollset_set == NULL);
|
||||
call->cq = cq;
|
||||
GRPC_CQ_INTERNAL_REF(cq, "bind");
|
||||
grpc_call_stack_set_pollset(exec_ctx, CALL_STACK_FROM_CALL(call),
|
||||
grpc_cq_pollset(cq));
|
||||
grpc_call_stack_set_pollset_or_pollset_set(
|
||||
exec_ctx, CALL_STACK_FROM_CALL(call), grpc_cq_pollset(cq), NULL);
|
||||
}
|
||||
|
||||
#ifdef GRPC_STREAM_REFCOUNT_DEBUG
|
||||
|
|
|
|||
|
|
@ -53,6 +53,8 @@ typedef void (*grpc_ioreq_completion_func)(grpc_exec_ctx *exec_ctx,
|
|||
grpc_call *grpc_call_create(grpc_channel *channel, grpc_call *parent_call,
|
||||
uint32_t propagation_mask,
|
||||
grpc_completion_queue *cq,
|
||||
/* if not NULL, it'll be used in lieu of \a cq */
|
||||
grpc_pollset_set *or_pollset_set,
|
||||
const void *server_transport_data,
|
||||
grpc_mdelem **add_initial_metadata,
|
||||
size_t add_initial_metadata_count,
|
||||
|
|
|
|||
|
|
@ -166,12 +166,14 @@ char *grpc_channel_get_target(grpc_channel *channel) {
|
|||
|
||||
static grpc_call *grpc_channel_create_call_internal(
|
||||
grpc_channel *channel, grpc_call *parent_call, uint32_t propagation_mask,
|
||||
grpc_completion_queue *cq, grpc_mdelem *path_mdelem,
|
||||
grpc_mdelem *authority_mdelem, gpr_timespec deadline) {
|
||||
grpc_completion_queue *cq, grpc_pollset_set *or_pollset_set,
|
||||
grpc_mdelem *path_mdelem, grpc_mdelem *authority_mdelem,
|
||||
gpr_timespec deadline) {
|
||||
grpc_mdelem *send_metadata[2];
|
||||
size_t num_metadata = 0;
|
||||
|
||||
GPR_ASSERT(channel->is_client);
|
||||
GPR_ASSERT(!(cq != NULL && or_pollset_set != NULL));
|
||||
|
||||
send_metadata[num_metadata++] = path_mdelem;
|
||||
if (authority_mdelem != NULL) {
|
||||
|
|
@ -180,8 +182,9 @@ static grpc_call *grpc_channel_create_call_internal(
|
|||
send_metadata[num_metadata++] = GRPC_MDELEM_REF(channel->default_authority);
|
||||
}
|
||||
|
||||
return grpc_call_create(channel, parent_call, propagation_mask, cq, NULL,
|
||||
send_metadata, num_metadata, deadline);
|
||||
return grpc_call_create(channel, parent_call, propagation_mask, cq,
|
||||
or_pollset_set, NULL, send_metadata, num_metadata,
|
||||
deadline);
|
||||
}
|
||||
|
||||
grpc_call *grpc_channel_create_call(grpc_channel *channel,
|
||||
|
|
@ -201,7 +204,22 @@ grpc_call *grpc_channel_create_call(grpc_channel *channel,
|
|||
(int)deadline.clock_type, reserved));
|
||||
GPR_ASSERT(!reserved);
|
||||
return grpc_channel_create_call_internal(
|
||||
channel, parent_call, propagation_mask, cq,
|
||||
channel, parent_call, propagation_mask, cq, NULL,
|
||||
grpc_mdelem_from_metadata_strings(GRPC_MDSTR_PATH,
|
||||
grpc_mdstr_from_string(method)),
|
||||
host ? grpc_mdelem_from_metadata_strings(GRPC_MDSTR_AUTHORITY,
|
||||
grpc_mdstr_from_string(host))
|
||||
: NULL,
|
||||
deadline);
|
||||
}
|
||||
|
||||
grpc_call *grpc_channel_create_pollset_set_call(
|
||||
grpc_channel *channel, grpc_call *parent_call, uint32_t propagation_mask,
|
||||
grpc_pollset_set *pollset_set, const char *method, const char *host,
|
||||
gpr_timespec deadline, void *reserved) {
|
||||
GPR_ASSERT(!reserved);
|
||||
return grpc_channel_create_call_internal(
|
||||
channel, parent_call, propagation_mask, NULL, pollset_set,
|
||||
grpc_mdelem_from_metadata_strings(GRPC_MDSTR_PATH,
|
||||
grpc_mdstr_from_string(method)),
|
||||
host ? grpc_mdelem_from_metadata_strings(GRPC_MDSTR_AUTHORITY,
|
||||
|
|
@ -245,7 +263,7 @@ grpc_call *grpc_channel_create_registered_call(
|
|||
(int)deadline.tv_nsec, (int)deadline.clock_type, reserved));
|
||||
GPR_ASSERT(!reserved);
|
||||
return grpc_channel_create_call_internal(
|
||||
channel, parent_call, propagation_mask, completion_queue,
|
||||
channel, parent_call, propagation_mask, completion_queue, NULL,
|
||||
GRPC_MDELEM_REF(rc->path),
|
||||
rc->authority ? GRPC_MDELEM_REF(rc->authority) : NULL, deadline);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -42,6 +42,11 @@ grpc_channel *grpc_channel_create(grpc_exec_ctx *exec_ctx, const char *target,
|
|||
grpc_channel_stack_type channel_stack_type,
|
||||
grpc_transport *optional_transport);
|
||||
|
||||
grpc_call *grpc_channel_create_pollset_set_call(
|
||||
grpc_channel *channel, grpc_call *parent_call, uint32_t propagation_mask,
|
||||
grpc_pollset_set *pollset_set, const char *method, const char *host,
|
||||
gpr_timespec deadline, void *reserved);
|
||||
|
||||
/** Get a (borrowed) pointer to this channels underlying channel stack */
|
||||
grpc_channel_stack *grpc_channel_get_channel_stack(grpc_channel *channel);
|
||||
|
||||
|
|
|
|||
|
|
@ -122,7 +122,7 @@ const grpc_channel_filter grpc_lame_filter = {
|
|||
lame_start_transport_op,
|
||||
sizeof(call_data),
|
||||
init_call_elem,
|
||||
grpc_call_stack_ignore_set_pollset,
|
||||
grpc_call_stack_ignore_set_pollset_or_pollset_set,
|
||||
destroy_call_elem,
|
||||
sizeof(channel_data),
|
||||
init_channel_elem,
|
||||
|
|
|
|||
|
|
@ -769,9 +769,9 @@ static void accept_stream(grpc_exec_ctx *exec_ctx, void *cd,
|
|||
const void *transport_server_data) {
|
||||
channel_data *chand = cd;
|
||||
/* create a call */
|
||||
grpc_call *call =
|
||||
grpc_call_create(chand->channel, NULL, 0, NULL, transport_server_data,
|
||||
NULL, 0, gpr_inf_future(GPR_CLOCK_MONOTONIC));
|
||||
grpc_call *call = grpc_call_create(chand->channel, NULL, 0, NULL, NULL,
|
||||
transport_server_data, NULL, 0,
|
||||
gpr_inf_future(GPR_CLOCK_MONOTONIC));
|
||||
grpc_call_element *elem =
|
||||
grpc_call_stack_element(grpc_call_get_call_stack(call), 0);
|
||||
call_data *calld = elem->call_data;
|
||||
|
|
@ -886,7 +886,7 @@ const grpc_channel_filter grpc_server_top_filter = {
|
|||
grpc_channel_next_op,
|
||||
sizeof(call_data),
|
||||
init_call_elem,
|
||||
grpc_call_stack_ignore_set_pollset,
|
||||
grpc_call_stack_ignore_set_pollset_or_pollset_set,
|
||||
destroy_call_elem,
|
||||
sizeof(channel_data),
|
||||
init_channel_elem,
|
||||
|
|
|
|||
|
|
@ -131,6 +131,13 @@ void grpc_transport_set_pollset(grpc_exec_ctx *exec_ctx,
|
|||
transport->vtable->set_pollset(exec_ctx, transport, stream, pollset);
|
||||
}
|
||||
|
||||
void grpc_transport_set_pollset_set(grpc_exec_ctx *exec_ctx,
|
||||
grpc_transport *transport,
|
||||
grpc_stream *stream,
|
||||
grpc_pollset_set *pollset_set) {
|
||||
transport->vtable->set_pollset_set(exec_ctx, transport, stream, pollset_set);
|
||||
}
|
||||
|
||||
void grpc_transport_destroy_stream(grpc_exec_ctx *exec_ctx,
|
||||
grpc_transport *transport,
|
||||
grpc_stream *stream) {
|
||||
|
|
|
|||
|
|
@ -201,6 +201,11 @@ void grpc_transport_set_pollset(grpc_exec_ctx *exec_ctx,
|
|||
grpc_transport *transport, grpc_stream *stream,
|
||||
grpc_pollset *pollset);
|
||||
|
||||
void grpc_transport_set_pollset_set(grpc_exec_ctx *exec_ctx,
|
||||
grpc_transport *transport,
|
||||
grpc_stream *stream,
|
||||
grpc_pollset_set *pollset_set);
|
||||
|
||||
/* Destroy transport data for a stream.
|
||||
|
||||
Requires: a recv_batch with final_state == GRPC_STREAM_CLOSED has been
|
||||
|
|
|
|||
|
|
@ -53,6 +53,10 @@ typedef struct grpc_transport_vtable {
|
|||
void (*set_pollset)(grpc_exec_ctx *exec_ctx, grpc_transport *self,
|
||||
grpc_stream *stream, grpc_pollset *pollset);
|
||||
|
||||
/* implementation of grpc_transport_set_pollset */
|
||||
void (*set_pollset_set)(grpc_exec_ctx *exec_ctx, grpc_transport *self,
|
||||
grpc_stream *stream, grpc_pollset_set *pollset_set);
|
||||
|
||||
/* implementation of grpc_transport_perform_stream_op */
|
||||
void (*perform_stream_op)(grpc_exec_ctx *exec_ctx, grpc_transport *self,
|
||||
grpc_stream *stream, grpc_transport_stream_op *op);
|
||||
|
|
|
|||
|
|
@ -92,17 +92,18 @@ static void free_call(grpc_exec_ctx *exec_ctx, void *arg, bool success) {
|
|||
}
|
||||
|
||||
static void test_create_channel_stack(void) {
|
||||
const grpc_channel_filter filter = {call_func,
|
||||
channel_func,
|
||||
sizeof(int),
|
||||
call_init_func,
|
||||
grpc_call_stack_ignore_set_pollset,
|
||||
call_destroy_func,
|
||||
sizeof(int),
|
||||
channel_init_func,
|
||||
channel_destroy_func,
|
||||
get_peer,
|
||||
"some_test_filter"};
|
||||
const grpc_channel_filter filter = {
|
||||
call_func,
|
||||
channel_func,
|
||||
sizeof(int),
|
||||
call_init_func,
|
||||
grpc_call_stack_ignore_set_pollset_or_pollset_set,
|
||||
call_destroy_func,
|
||||
sizeof(int),
|
||||
channel_init_func,
|
||||
channel_destroy_func,
|
||||
get_peer,
|
||||
"some_test_filter"};
|
||||
const grpc_channel_filter *filters = &filter;
|
||||
grpc_channel_stack *channel_stack;
|
||||
grpc_call_stack *call_stack;
|
||||
|
|
|
|||
|
|
@ -250,7 +250,7 @@ static const grpc_channel_filter test_filter = {
|
|||
grpc_channel_next_op,
|
||||
sizeof(call_data),
|
||||
init_call_elem,
|
||||
grpc_call_stack_ignore_set_pollset,
|
||||
grpc_call_stack_ignore_set_pollset_or_pollset_set,
|
||||
destroy_call_elem,
|
||||
sizeof(channel_data),
|
||||
init_channel_elem,
|
||||
|
|
|
|||
|
|
@ -49,6 +49,7 @@ static int g_done = 0;
|
|||
static grpc_httpcli_context g_context;
|
||||
static gpr_mu *g_mu;
|
||||
static grpc_pollset *g_pollset;
|
||||
static grpc_pollset_set *g_pollset_set;
|
||||
|
||||
static gpr_timespec n_seconds_time(int seconds) {
|
||||
return GRPC_TIMEOUT_SECONDS_TO_DEADLINE(seconds);
|
||||
|
|
@ -86,8 +87,8 @@ static void test_get(int port) {
|
|||
req.http.path = "/get";
|
||||
req.handshaker = &grpc_httpcli_plaintext;
|
||||
|
||||
grpc_httpcli_get(&exec_ctx, &g_context, g_pollset, &req, n_seconds_time(15),
|
||||
on_finish, (void *)42);
|
||||
grpc_httpcli_get(&exec_ctx, &g_context, g_pollset_set, &req,
|
||||
n_seconds_time(15), on_finish, (void *)42);
|
||||
gpr_mu_lock(g_mu);
|
||||
while (!g_done) {
|
||||
grpc_pollset_worker *worker = NULL;
|
||||
|
|
@ -117,7 +118,7 @@ static void test_post(int port) {
|
|||
req.http.path = "/post";
|
||||
req.handshaker = &grpc_httpcli_plaintext;
|
||||
|
||||
grpc_httpcli_post(&exec_ctx, &g_context, g_pollset, &req, "hello", 5,
|
||||
grpc_httpcli_post(&exec_ctx, &g_context, g_pollset_set, &req, "hello", 5,
|
||||
n_seconds_time(15), on_finish, (void *)42);
|
||||
gpr_mu_lock(g_mu);
|
||||
while (!g_done) {
|
||||
|
|
@ -182,6 +183,8 @@ int main(int argc, char **argv) {
|
|||
grpc_httpcli_context_init(&g_context);
|
||||
g_pollset = gpr_malloc(grpc_pollset_size());
|
||||
grpc_pollset_init(g_pollset, &g_mu);
|
||||
g_pollset_set = grpc_pollset_set_create();
|
||||
grpc_pollset_set_add_pollset(&exec_ctx, g_pollset_set, g_pollset);
|
||||
|
||||
test_get(port);
|
||||
test_post(port);
|
||||
|
|
|
|||
|
|
@ -49,6 +49,7 @@ static int g_done = 0;
|
|||
static grpc_httpcli_context g_context;
|
||||
static gpr_mu *g_mu;
|
||||
static grpc_pollset *g_pollset;
|
||||
static grpc_pollset_set *g_pollset_set;
|
||||
|
||||
static gpr_timespec n_seconds_time(int seconds) {
|
||||
return GRPC_TIMEOUT_SECONDS_TO_DEADLINE(seconds);
|
||||
|
|
@ -87,8 +88,8 @@ static void test_get(int port) {
|
|||
req.http.path = "/get";
|
||||
req.handshaker = &grpc_httpcli_ssl;
|
||||
|
||||
grpc_httpcli_get(&exec_ctx, &g_context, g_pollset, &req, n_seconds_time(15),
|
||||
on_finish, (void *)42);
|
||||
grpc_httpcli_get(&exec_ctx, &g_context, g_pollset_set, &req,
|
||||
n_seconds_time(15), on_finish, (void *)42);
|
||||
gpr_mu_lock(g_mu);
|
||||
while (!g_done) {
|
||||
grpc_pollset_worker *worker = NULL;
|
||||
|
|
@ -119,7 +120,7 @@ static void test_post(int port) {
|
|||
req.http.path = "/post";
|
||||
req.handshaker = &grpc_httpcli_ssl;
|
||||
|
||||
grpc_httpcli_post(&exec_ctx, &g_context, g_pollset, &req, "hello", 5,
|
||||
grpc_httpcli_post(&exec_ctx, &g_context, g_pollset_set, &req, "hello", 5,
|
||||
n_seconds_time(15), on_finish, (void *)42);
|
||||
gpr_mu_lock(g_mu);
|
||||
while (!g_done) {
|
||||
|
|
@ -185,6 +186,8 @@ int main(int argc, char **argv) {
|
|||
grpc_httpcli_context_init(&g_context);
|
||||
g_pollset = gpr_malloc(grpc_pollset_size());
|
||||
grpc_pollset_init(g_pollset, &g_mu);
|
||||
g_pollset_set = grpc_pollset_set_create();
|
||||
grpc_pollset_set_add_pollset(&exec_ctx, g_pollset_set, g_pollset);
|
||||
|
||||
test_get(port);
|
||||
test_post(port);
|
||||
|
|
@ -193,6 +196,7 @@ int main(int argc, char **argv) {
|
|||
grpc_closure_init(&destroyed, destroy_pollset, g_pollset);
|
||||
grpc_pollset_shutdown(&exec_ctx, g_pollset, &destroyed);
|
||||
grpc_exec_ctx_finish(&exec_ctx);
|
||||
grpc_pollset_set_destroy(g_pollset_set);
|
||||
grpc_shutdown();
|
||||
|
||||
gpr_free(g_pollset);
|
||||
|
|
|
|||
|
|
@ -47,6 +47,7 @@
|
|||
typedef struct {
|
||||
gpr_mu *mu;
|
||||
grpc_pollset *pollset;
|
||||
grpc_pollset_set *pollset_set;
|
||||
int is_done;
|
||||
char *token;
|
||||
} oauth2_request;
|
||||
|
|
@ -85,13 +86,15 @@ char *grpc_test_fetch_oauth2_token_with_credentials(
|
|||
|
||||
request.pollset = gpr_malloc(grpc_pollset_size());
|
||||
grpc_pollset_init(request.pollset, &request.mu);
|
||||
request.pollset_set = grpc_pollset_set_create();
|
||||
grpc_pollset_set_add_pollset(&exec_ctx, request.pollset_set, request.pollset);
|
||||
request.is_done = 0;
|
||||
|
||||
grpc_closure_init(&do_nothing_closure, do_nothing, NULL);
|
||||
|
||||
grpc_call_credentials_get_request_metadata(&exec_ctx, creds, request.pollset,
|
||||
null_ctx, on_oauth2_response,
|
||||
&request);
|
||||
grpc_call_credentials_get_request_metadata(&exec_ctx, creds,
|
||||
request.pollset_set, null_ctx,
|
||||
on_oauth2_response, &request);
|
||||
|
||||
grpc_exec_ctx_finish(&exec_ctx);
|
||||
|
||||
|
|
@ -107,6 +110,7 @@ char *grpc_test_fetch_oauth2_token_with_credentials(
|
|||
grpc_pollset_shutdown(&exec_ctx, request.pollset, &do_nothing_closure);
|
||||
grpc_exec_ctx_finish(&exec_ctx);
|
||||
grpc_pollset_destroy(request.pollset);
|
||||
grpc_pollset_set_destroy(request.pollset_set);
|
||||
gpr_free(request.pollset);
|
||||
return request.token;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -48,6 +48,7 @@
|
|||
typedef struct {
|
||||
gpr_mu *mu;
|
||||
grpc_pollset *pollset;
|
||||
grpc_pollset_set *pollset_set;
|
||||
int is_done;
|
||||
} synchronizer;
|
||||
|
||||
|
|
@ -95,11 +96,13 @@ int main(int argc, char **argv) {
|
|||
|
||||
sync.pollset = gpr_malloc(grpc_pollset_size());
|
||||
grpc_pollset_init(sync.pollset, &sync.mu);
|
||||
sync.pollset_set = grpc_pollset_set_create();
|
||||
grpc_pollset_set_add_pollset(&exec_ctx, sync.pollset_set, sync.pollset);
|
||||
sync.is_done = 0;
|
||||
|
||||
grpc_call_credentials_get_request_metadata(
|
||||
&exec_ctx, ((grpc_composite_channel_credentials *)creds)->call_creds,
|
||||
sync.pollset, context, on_metadata_response, &sync);
|
||||
sync.pollset_set, context, on_metadata_response, &sync);
|
||||
|
||||
gpr_mu_lock(sync.mu);
|
||||
while (!sync.is_done) {
|
||||
|
|
@ -116,6 +119,7 @@ int main(int argc, char **argv) {
|
|||
grpc_exec_ctx_finish(&exec_ctx);
|
||||
|
||||
grpc_channel_credentials_release(creds);
|
||||
grpc_pollset_set_destroy(sync.pollset_set);
|
||||
gpr_free(sync.pollset);
|
||||
|
||||
end:
|
||||
|
|
|
|||
|
|
@ -52,6 +52,7 @@
|
|||
typedef struct freereq {
|
||||
gpr_mu *mu;
|
||||
grpc_pollset *pollset;
|
||||
grpc_pollset_set *pollset_set;
|
||||
int done;
|
||||
} freereq;
|
||||
|
||||
|
|
@ -86,6 +87,8 @@ void grpc_free_port_using_server(char *server, int port) {
|
|||
|
||||
pr.pollset = gpr_malloc(grpc_pollset_size());
|
||||
grpc_pollset_init(pr.pollset, &pr.mu);
|
||||
pr.pollset_set = grpc_pollset_set_create();
|
||||
grpc_pollset_set_add_pollset(&exec_ctx, pr.pollset_set, pr.pollset);
|
||||
shutdown_closure =
|
||||
grpc_closure_create(destroy_pollset_and_shutdown, pr.pollset);
|
||||
|
||||
|
|
@ -94,7 +97,7 @@ void grpc_free_port_using_server(char *server, int port) {
|
|||
req.http.path = path;
|
||||
|
||||
grpc_httpcli_context_init(&context);
|
||||
grpc_httpcli_get(&exec_ctx, &context, pr.pollset, &req,
|
||||
grpc_httpcli_get(&exec_ctx, &context, pr.pollset_set, &req,
|
||||
GRPC_TIMEOUT_SECONDS_TO_DEADLINE(10), freed_port_from_server,
|
||||
&pr);
|
||||
gpr_mu_lock(pr.mu);
|
||||
|
|
@ -110,12 +113,14 @@ void grpc_free_port_using_server(char *server, int port) {
|
|||
grpc_exec_ctx_finish(&exec_ctx);
|
||||
grpc_pollset_shutdown(&exec_ctx, pr.pollset, shutdown_closure);
|
||||
grpc_exec_ctx_finish(&exec_ctx);
|
||||
grpc_pollset_set_destroy(pr.pollset_set);
|
||||
gpr_free(path);
|
||||
}
|
||||
|
||||
typedef struct portreq {
|
||||
gpr_mu *mu;
|
||||
grpc_pollset *pollset;
|
||||
grpc_pollset_set *pollset_set;
|
||||
int port;
|
||||
int retries;
|
||||
char *server;
|
||||
|
|
@ -151,7 +156,7 @@ static void got_port_from_server(grpc_exec_ctx *exec_ctx, void *arg,
|
|||
pr->retries++;
|
||||
req.host = pr->server;
|
||||
req.http.path = "/get";
|
||||
grpc_httpcli_get(exec_ctx, pr->ctx, pr->pollset, &req,
|
||||
grpc_httpcli_get(exec_ctx, pr->ctx, pr->pollset_set, &req,
|
||||
GRPC_TIMEOUT_SECONDS_TO_DEADLINE(10), got_port_from_server,
|
||||
pr);
|
||||
return;
|
||||
|
|
@ -182,6 +187,8 @@ int grpc_pick_port_using_server(char *server) {
|
|||
memset(&req, 0, sizeof(req));
|
||||
pr.pollset = gpr_malloc(grpc_pollset_size());
|
||||
grpc_pollset_init(pr.pollset, &pr.mu);
|
||||
pr.pollset_set = grpc_pollset_set_create();
|
||||
grpc_pollset_set_add_pollset(&exec_ctx, pr.pollset_set, pr.pollset);
|
||||
shutdown_closure =
|
||||
grpc_closure_create(destroy_pollset_and_shutdown, pr.pollset);
|
||||
pr.port = -1;
|
||||
|
|
@ -192,7 +199,7 @@ int grpc_pick_port_using_server(char *server) {
|
|||
req.http.path = "/get";
|
||||
|
||||
grpc_httpcli_context_init(&context);
|
||||
grpc_httpcli_get(&exec_ctx, &context, pr.pollset, &req,
|
||||
grpc_httpcli_get(&exec_ctx, &context, pr.pollset_set, &req,
|
||||
GRPC_TIMEOUT_SECONDS_TO_DEADLINE(10), got_port_from_server,
|
||||
&pr);
|
||||
grpc_exec_ctx_finish(&exec_ctx);
|
||||
|
|
@ -208,6 +215,7 @@ int grpc_pick_port_using_server(char *server) {
|
|||
grpc_httpcli_context_destroy(&context);
|
||||
grpc_pollset_shutdown(&exec_ctx, pr.pollset, shutdown_closure);
|
||||
grpc_exec_ctx_finish(&exec_ctx);
|
||||
grpc_pollset_set_destroy(pr.pollset_set);
|
||||
|
||||
return pr.port;
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in New Issue