blob: e6c253a3f06a062fa850eb4748856ed79ca56442 [file] [log] [blame]
murgatroid999030c812016-09-16 13:25:08 -07001/*
2 *
3 * Copyright 2016, Google Inc.
4 * All rights reserved.
5 *
6 * Redistribution and use in source and binary forms, with or without
7 * modification, are permitted provided that the following conditions are
8 * met:
9 *
10 * * Redistributions of source code must retain the above copyright
11 * notice, this list of conditions and the following disclaimer.
12 * * Redistributions in binary form must reproduce the above
13 * copyright notice, this list of conditions and the following disclaimer
14 * in the documentation and/or other materials provided with the
15 * distribution.
16 * * Neither the name of Google Inc. nor the names of its
17 * contributors may be used to endorse or promote products derived from
18 * this software without specific prior written permission.
19 *
20 * THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS
21 * "AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT
22 * LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR
23 * A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT
24 * OWNER OR CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL,
25 * SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT
26 * LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE,
27 * DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY
28 * THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT
29 * (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE
30 * OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
31 *
32 */
33
34#include "src/core/lib/iomgr/port.h"
35
36#ifdef GRPC_UV
37
38#include <limits.h>
39#include <string.h>
40
murgatroid99e2672c92016-11-09 15:12:22 -080041#include <grpc/slice_buffer.h>
42
murgatroid999030c812016-09-16 13:25:08 -070043#include <grpc/support/alloc.h>
44#include <grpc/support/log.h>
murgatroid999030c812016-09-16 13:25:08 -070045#include <grpc/support/string_util.h>
46
47#include "src/core/lib/iomgr/error.h"
48#include "src/core/lib/iomgr/network_status_tracker.h"
murgatroid99e2672c92016-11-09 15:12:22 -080049#include "src/core/lib/iomgr/resource_quota.h"
murgatroid999030c812016-09-16 13:25:08 -070050#include "src/core/lib/iomgr/tcp_uv.h"
Craig Tiller7c70b6c2017-01-23 07:48:42 -080051#include "src/core/lib/slice/slice_internal.h"
murgatroid99e2672c92016-11-09 15:12:22 -080052#include "src/core/lib/slice/slice_string_helpers.h"
murgatroid999030c812016-09-16 13:25:08 -070053#include "src/core/lib/support/string.h"
54
Craig Tiller020176d2017-05-09 08:36:44 -070055grpc_tracer_flag grpc_tcp_trace = GRPC_TRACER_INITIALIZER(false);
murgatroid999030c812016-09-16 13:25:08 -070056
57typedef struct {
58 grpc_endpoint base;
59 gpr_refcount refcount;
60
murgatroid9969259d42016-10-31 14:34:10 -070061 uv_write_t write_req;
62 uv_shutdown_t shutdown_req;
63
murgatroid999030c812016-09-16 13:25:08 -070064 uv_tcp_t *handle;
65
66 grpc_closure *read_cb;
67 grpc_closure *write_cb;
68
murgatroid99e2672c92016-11-09 15:12:22 -080069 grpc_slice read_slice;
70 grpc_slice_buffer *read_slices;
71 grpc_slice_buffer *write_slices;
murgatroid999030c812016-09-16 13:25:08 -070072 uv_buf_t *write_buffers;
73
murgatroid99e2672c92016-11-09 15:12:22 -080074 grpc_resource_user *resource_user;
murgatroid9969259d42016-10-31 14:34:10 -070075
murgatroid992c287ca2016-10-07 09:55:35 -070076 bool shutting_down;
murgatroid9969259d42016-10-31 14:34:10 -070077
murgatroid999030c812016-09-16 13:25:08 -070078 char *peer_string;
79 grpc_pollset *pollset;
80} grpc_tcp;
81
murgatroid99e2672c92016-11-09 15:12:22 -080082static void tcp_free(grpc_exec_ctx *exec_ctx, grpc_tcp *tcp) {
murgatroid99bb8b1c92017-06-23 15:10:18 -070083 grpc_slice_unref(tcp->read_slice);
murgatroid99e2672c92016-11-09 15:12:22 -080084 grpc_resource_user_unref(exec_ctx, tcp->resource_user);
murgatroid9969259d42016-10-31 14:34:10 -070085 gpr_free(tcp);
murgatroid9969259d42016-10-31 14:34:10 -070086}
murgatroid999030c812016-09-16 13:25:08 -070087
88/*#define GRPC_TCP_REFCOUNT_DEBUG*/
89#ifdef GRPC_TCP_REFCOUNT_DEBUG
murgatroid992e012342016-11-10 18:24:08 -080090#define TCP_UNREF(exec_ctx, tcp, reason) \
91 tcp_unref((exec_ctx), (tcp), (reason), __FILE__, __LINE__)
murgatroid99a1137ce2017-05-09 14:21:12 -070092#define TCP_REF(tcp, reason) tcp_ref((tcp), (reason), __FILE__, __LINE__)
murgatroid992e012342016-11-10 18:24:08 -080093static void tcp_unref(grpc_exec_ctx *exec_ctx, grpc_tcp *tcp,
94 const char *reason, const char *file, int line) {
murgatroid99a1137ce2017-05-09 14:21:12 -070095 gpr_log(file, line, GPR_LOG_SEVERITY_DEBUG,
96 "TCP unref %p : %s %" PRIiPTR " -> %" PRIiPTR, tcp, reason,
97 tcp->refcount.count, tcp->refcount.count - 1);
murgatroid999030c812016-09-16 13:25:08 -070098 if (gpr_unref(&tcp->refcount)) {
murgatroid99e2672c92016-11-09 15:12:22 -080099 tcp_free(exec_ctx, tcp);
murgatroid999030c812016-09-16 13:25:08 -0700100 }
101}
102
103static void tcp_ref(grpc_tcp *tcp, const char *reason, const char *file,
104 int line) {
murgatroid99a1137ce2017-05-09 14:21:12 -0700105 gpr_log(file, line, GPR_LOG_SEVERITY_DEBUG,
106 "TCP ref %p : %s %" PRIiPTR " -> %" PRIiPTR, tcp, reason,
107 tcp->refcount.count, tcp->refcount.count + 1);
murgatroid999030c812016-09-16 13:25:08 -0700108 gpr_ref(&tcp->refcount);
109}
110#else
murgatroid99e2672c92016-11-09 15:12:22 -0800111#define TCP_UNREF(exec_ctx, tcp, reason) tcp_unref((exec_ctx), (tcp))
murgatroid999030c812016-09-16 13:25:08 -0700112#define TCP_REF(tcp, reason) tcp_ref((tcp))
murgatroid99e2672c92016-11-09 15:12:22 -0800113static void tcp_unref(grpc_exec_ctx *exec_ctx, grpc_tcp *tcp) {
murgatroid999030c812016-09-16 13:25:08 -0700114 if (gpr_unref(&tcp->refcount)) {
murgatroid99e2672c92016-11-09 15:12:22 -0800115 tcp_free(exec_ctx, tcp);
murgatroid999030c812016-09-16 13:25:08 -0700116 }
117}
118
119static void tcp_ref(grpc_tcp *tcp) { gpr_ref(&tcp->refcount); }
120#endif
121
murgatroid991191b722017-02-08 11:56:52 -0800122static void uv_close_callback(uv_handle_t *handle) {
123 grpc_exec_ctx exec_ctx = GRPC_EXEC_CTX_INIT;
124 grpc_tcp *tcp = handle->data;
125 TCP_UNREF(&exec_ctx, tcp, "destroy");
126 grpc_exec_ctx_finish(&exec_ctx);
127}
128
murgatroid99bb8b1c92017-06-23 15:10:18 -0700129static grpc_slice alloc_read_slice(grpc_exec_ctx *exec_ctx,
130 grpc_resource_user *resource_user) {
131 return grpc_resource_user_slice_malloc(exec_ctx, resource_user,
132 GRPC_TCP_DEFAULT_READ_SLICE_SIZE);
133}
134
murgatroid99dedb9232016-09-26 13:54:04 -0700135static void alloc_uv_buf(uv_handle_t *handle, size_t suggested_size,
136 uv_buf_t *buf) {
murgatroid9969259d42016-10-31 14:34:10 -0700137 grpc_exec_ctx exec_ctx = GRPC_EXEC_CTX_INIT;
murgatroid999030c812016-09-16 13:25:08 -0700138 grpc_tcp *tcp = handle->data;
139 (void)suggested_size;
Craig Tiller618e67d2016-10-26 21:08:10 -0700140 buf->base = (char *)GRPC_SLICE_START_PTR(tcp->read_slice);
141 buf->len = GRPC_SLICE_LENGTH(tcp->read_slice);
murgatroid9969259d42016-10-31 14:34:10 -0700142 grpc_exec_ctx_finish(&exec_ctx);
murgatroid999030c812016-09-16 13:25:08 -0700143}
144
murgatroid99dedb9232016-09-26 13:54:04 -0700145static void read_callback(uv_stream_t *stream, ssize_t nread,
146 const uv_buf_t *buf) {
murgatroid99e2672c92016-11-09 15:12:22 -0800147 grpc_slice sub;
murgatroid999030c812016-09-16 13:25:08 -0700148 grpc_error *error;
149 grpc_exec_ctx exec_ctx = GRPC_EXEC_CTX_INIT;
150 grpc_tcp *tcp = stream->data;
151 grpc_closure *cb = tcp->read_cb;
152 if (nread == 0) {
153 // Nothing happened. Wait for the next callback
154 return;
155 }
murgatroid99e2672c92016-11-09 15:12:22 -0800156 TCP_UNREF(&exec_ctx, tcp, "read");
murgatroid999030c812016-09-16 13:25:08 -0700157 tcp->read_cb = NULL;
158 // TODO(murgatroid99): figure out what the return value here means
159 uv_read_stop(stream);
160 if (nread == UV_EOF) {
ncteisen4b36a3d2017-03-13 19:08:06 -0700161 error = GRPC_ERROR_CREATE_FROM_STATIC_STRING("EOF");
murgatroid999030c812016-09-16 13:25:08 -0700162 } else if (nread > 0) {
163 // Successful read
murgatroid99e2672c92016-11-09 15:12:22 -0800164 sub = grpc_slice_sub_no_ref(tcp->read_slice, 0, (size_t)nread);
165 grpc_slice_buffer_add(tcp->read_slices, sub);
murgatroid99bb8b1c92017-06-23 15:10:18 -0700166 tcp->read_slice = alloc_read_slice(&exec_ctx, tcp->resource_user);
murgatroid999030c812016-09-16 13:25:08 -0700167 error = GRPC_ERROR_NONE;
Craig Tiller341bcc52017-05-09 08:37:44 -0700168 if (GRPC_TRACER_ON(grpc_tcp_trace)) {
murgatroid999030c812016-09-16 13:25:08 -0700169 size_t i;
170 const char *str = grpc_error_string(error);
171 gpr_log(GPR_DEBUG, "read: error=%s", str);
Craig Tiller7c70b6c2017-01-23 07:48:42 -0800172
murgatroid999030c812016-09-16 13:25:08 -0700173 for (i = 0; i < tcp->read_slices->count; i++) {
murgatroid99e2672c92016-11-09 15:12:22 -0800174 char *dump = grpc_dump_slice(tcp->read_slices->slices[i],
murgatroid992e012342016-11-10 18:24:08 -0800175 GPR_DUMP_HEX | GPR_DUMP_ASCII);
murgatroid99dedb9232016-09-26 13:54:04 -0700176 gpr_log(GPR_DEBUG, "READ %p (peer=%s): %s", tcp, tcp->peer_string,
177 dump);
murgatroid999030c812016-09-16 13:25:08 -0700178 gpr_free(dump);
179 }
180 }
181 } else {
182 // nread < 0: Error
ncteisen4b36a3d2017-03-13 19:08:06 -0700183 error = GRPC_ERROR_CREATE_FROM_STATIC_STRING("TCP Read failed");
murgatroid999030c812016-09-16 13:25:08 -0700184 }
Craig Tiller91031da2016-12-28 15:44:25 -0800185 grpc_closure_sched(&exec_ctx, cb, error);
murgatroid999030c812016-09-16 13:25:08 -0700186 grpc_exec_ctx_finish(&exec_ctx);
187}
188
189static void uv_endpoint_read(grpc_exec_ctx *exec_ctx, grpc_endpoint *ep,
murgatroid99e2672c92016-11-09 15:12:22 -0800190 grpc_slice_buffer *read_slices, grpc_closure *cb) {
murgatroid999030c812016-09-16 13:25:08 -0700191 grpc_tcp *tcp = (grpc_tcp *)ep;
192 int status;
193 grpc_error *error = GRPC_ERROR_NONE;
194 GPR_ASSERT(tcp->read_cb == NULL);
195 tcp->read_cb = cb;
196 tcp->read_slices = read_slices;
Craig Tillerab7b2d82016-12-07 07:26:34 -0800197 grpc_slice_buffer_reset_and_unref_internal(exec_ctx, read_slices);
murgatroid999030c812016-09-16 13:25:08 -0700198 TCP_REF(tcp, "read");
199 // TODO(murgatroid99): figure out what the return value here means
murgatroid99dedb9232016-09-26 13:54:04 -0700200 status =
201 uv_read_start((uv_stream_t *)tcp->handle, alloc_uv_buf, read_callback);
murgatroid999030c812016-09-16 13:25:08 -0700202 if (status != 0) {
ncteisen4b36a3d2017-03-13 19:08:06 -0700203 error = GRPC_ERROR_CREATE_FROM_STATIC_STRING("TCP Read failed at start");
murgatroid99dedb9232016-09-26 13:54:04 -0700204 error =
ncteisen4b36a3d2017-03-13 19:08:06 -0700205 grpc_error_set_str(error, GRPC_ERROR_STR_OS_ERROR,
206 grpc_slice_from_static_string(uv_strerror(status)));
Craig Tiller91031da2016-12-28 15:44:25 -0800207 grpc_closure_sched(exec_ctx, cb, error);
murgatroid999030c812016-09-16 13:25:08 -0700208 }
Craig Tiller341bcc52017-05-09 08:37:44 -0700209 if (GRPC_TRACER_ON(grpc_tcp_trace)) {
murgatroid999030c812016-09-16 13:25:08 -0700210 const char *str = grpc_error_string(error);
211 gpr_log(GPR_DEBUG, "Initiating read on %p: error=%s", tcp, str);
212 }
213}
214
215static void write_callback(uv_write_t *req, int status) {
216 grpc_tcp *tcp = req->data;
217 grpc_error *error;
218 grpc_exec_ctx exec_ctx = GRPC_EXEC_CTX_INIT;
219 grpc_closure *cb = tcp->write_cb;
220 tcp->write_cb = NULL;
murgatroid99e2672c92016-11-09 15:12:22 -0800221 TCP_UNREF(&exec_ctx, tcp, "write");
murgatroid999030c812016-09-16 13:25:08 -0700222 if (status == 0) {
223 error = GRPC_ERROR_NONE;
224 } else {
ncteisen4b36a3d2017-03-13 19:08:06 -0700225 error = GRPC_ERROR_CREATE_FROM_STATIC_STRING("TCP Write failed");
murgatroid999030c812016-09-16 13:25:08 -0700226 }
Craig Tiller341bcc52017-05-09 08:37:44 -0700227 if (GRPC_TRACER_ON(grpc_tcp_trace)) {
murgatroid999030c812016-09-16 13:25:08 -0700228 const char *str = grpc_error_string(error);
229 gpr_log(GPR_DEBUG, "write complete on %p: error=%s", tcp, str);
230 }
231 gpr_free(tcp->write_buffers);
murgatroid99e2672c92016-11-09 15:12:22 -0800232 grpc_resource_user_free(&exec_ctx, tcp->resource_user,
murgatroid9969259d42016-10-31 14:34:10 -0700233 sizeof(uv_buf_t) * tcp->write_slices->count);
Craig Tiller91031da2016-12-28 15:44:25 -0800234 grpc_closure_sched(&exec_ctx, cb, error);
murgatroid999030c812016-09-16 13:25:08 -0700235 grpc_exec_ctx_finish(&exec_ctx);
236}
237
238static void uv_endpoint_write(grpc_exec_ctx *exec_ctx, grpc_endpoint *ep,
murgatroid99e2672c92016-11-09 15:12:22 -0800239 grpc_slice_buffer *write_slices,
murgatroid999030c812016-09-16 13:25:08 -0700240 grpc_closure *cb) {
241 grpc_tcp *tcp = (grpc_tcp *)ep;
242 uv_buf_t *buffers;
243 unsigned int buffer_count;
244 unsigned int i;
murgatroid99e2672c92016-11-09 15:12:22 -0800245 grpc_slice *slice;
murgatroid999030c812016-09-16 13:25:08 -0700246 uv_write_t *write_req;
247
Craig Tiller341bcc52017-05-09 08:37:44 -0700248 if (GRPC_TRACER_ON(grpc_tcp_trace)) {
murgatroid99c36f6ea2016-10-03 09:24:09 -0700249 size_t j;
murgatroid999030c812016-09-16 13:25:08 -0700250
murgatroid99c36f6ea2016-10-03 09:24:09 -0700251 for (j = 0; j < write_slices->count; j++) {
murgatroid99e2672c92016-11-09 15:12:22 -0800252 char *data = grpc_dump_slice(write_slices->slices[j],
murgatroid992e012342016-11-10 18:24:08 -0800253 GPR_DUMP_HEX | GPR_DUMP_ASCII);
murgatroid999030c812016-09-16 13:25:08 -0700254 gpr_log(GPR_DEBUG, "WRITE %p (peer=%s): %s", tcp, tcp->peer_string, data);
255 gpr_free(data);
256 }
257 }
258
259 if (tcp->shutting_down) {
ncteisen4b36a3d2017-03-13 19:08:06 -0700260 grpc_closure_sched(exec_ctx, cb, GRPC_ERROR_CREATE_FROM_STATIC_STRING(
261 "TCP socket is shutting down"));
murgatroid999030c812016-09-16 13:25:08 -0700262 return;
263 }
264
265 GPR_ASSERT(tcp->write_cb == NULL);
266 tcp->write_slices = write_slices;
267 GPR_ASSERT(tcp->write_slices->count <= UINT_MAX);
268 if (tcp->write_slices->count == 0) {
269 // No slices means we don't have to do anything,
270 // and libuv doesn't like empty writes
Craig Tiller91031da2016-12-28 15:44:25 -0800271 grpc_closure_sched(exec_ctx, cb, GRPC_ERROR_NONE);
murgatroid999030c812016-09-16 13:25:08 -0700272 return;
273 }
274
275 tcp->write_cb = cb;
276 buffer_count = (unsigned int)tcp->write_slices->count;
277 buffers = gpr_malloc(sizeof(uv_buf_t) * buffer_count);
murgatroid99e2672c92016-11-09 15:12:22 -0800278 grpc_resource_user_alloc(exec_ctx, tcp->resource_user,
murgatroid9969259d42016-10-31 14:34:10 -0700279 sizeof(uv_buf_t) * buffer_count, NULL);
murgatroid999030c812016-09-16 13:25:08 -0700280 for (i = 0; i < buffer_count; i++) {
281 slice = &tcp->write_slices->slices[i];
Craig Tiller618e67d2016-10-26 21:08:10 -0700282 buffers[i].base = (char *)GRPC_SLICE_START_PTR(*slice);
283 buffers[i].len = GRPC_SLICE_LENGTH(*slice);
murgatroid999030c812016-09-16 13:25:08 -0700284 }
murgatroid9969259d42016-10-31 14:34:10 -0700285 tcp->write_buffers = buffers;
286 write_req = &tcp->write_req;
murgatroid999030c812016-09-16 13:25:08 -0700287 write_req->data = tcp;
288 TCP_REF(tcp, "write");
289 // TODO(murgatroid99): figure out what the return value here means
290 uv_write(write_req, (uv_stream_t *)tcp->handle, buffers, buffer_count,
291 write_callback);
292}
293
294static void uv_add_to_pollset(grpc_exec_ctx *exec_ctx, grpc_endpoint *ep,
295 grpc_pollset *pollset) {
296 // No-op. We're ignoring pollsets currently
murgatroid99dedb9232016-09-26 13:54:04 -0700297 (void)exec_ctx;
298 (void)ep;
299 (void)pollset;
300 grpc_tcp *tcp = (grpc_tcp *)ep;
murgatroid999030c812016-09-16 13:25:08 -0700301 tcp->pollset = pollset;
302}
303
304static void uv_add_to_pollset_set(grpc_exec_ctx *exec_ctx, grpc_endpoint *ep,
305 grpc_pollset_set *pollset) {
306 // No-op. We're ignoring pollsets currently
murgatroid99dedb9232016-09-26 13:54:04 -0700307 (void)exec_ctx;
308 (void)ep;
309 (void)pollset;
murgatroid999030c812016-09-16 13:25:08 -0700310}
311
murgatroid9969259d42016-10-31 14:34:10 -0700312static void shutdown_callback(uv_shutdown_t *req, int status) {}
313
Craig Tiller22f13fb2017-01-27 11:43:25 -0800314static void uv_endpoint_shutdown(grpc_exec_ctx *exec_ctx, grpc_endpoint *ep,
315 grpc_error *why) {
murgatroid999030c812016-09-16 13:25:08 -0700316 grpc_tcp *tcp = (grpc_tcp *)ep;
murgatroid992c287ca2016-10-07 09:55:35 -0700317 if (!tcp->shutting_down) {
318 tcp->shutting_down = true;
murgatroid9969259d42016-10-31 14:34:10 -0700319 uv_shutdown_t *req = &tcp->shutdown_req;
murgatroid992c287ca2016-10-07 09:55:35 -0700320 uv_shutdown(req, (uv_stream_t *)tcp->handle, shutdown_callback);
murgatroid99a1137ce2017-05-09 14:21:12 -0700321 grpc_resource_user_shutdown(exec_ctx, tcp->resource_user);
murgatroid992c287ca2016-10-07 09:55:35 -0700322 }
Craig Tiller22f13fb2017-01-27 11:43:25 -0800323 GRPC_ERROR_UNREF(why);
murgatroid999030c812016-09-16 13:25:08 -0700324}
325
326static void uv_destroy(grpc_exec_ctx *exec_ctx, grpc_endpoint *ep) {
327 grpc_network_status_unregister_endpoint(ep);
328 grpc_tcp *tcp = (grpc_tcp *)ep;
murgatroid999030c812016-09-16 13:25:08 -0700329 uv_close((uv_handle_t *)tcp->handle, uv_close_callback);
murgatroid999030c812016-09-16 13:25:08 -0700330}
331
332static char *uv_get_peer(grpc_endpoint *ep) {
333 grpc_tcp *tcp = (grpc_tcp *)ep;
334 return gpr_strdup(tcp->peer_string);
335}
336
murgatroid9969259d42016-10-31 14:34:10 -0700337static grpc_resource_user *uv_get_resource_user(grpc_endpoint *ep) {
338 grpc_tcp *tcp = (grpc_tcp *)ep;
murgatroid99e2672c92016-11-09 15:12:22 -0800339 return tcp->resource_user;
murgatroid9969259d42016-10-31 14:34:10 -0700340}
341
murgatroid99dedb9232016-09-26 13:54:04 -0700342static grpc_workqueue *uv_get_workqueue(grpc_endpoint *ep) { return NULL; }
murgatroid999030c812016-09-16 13:25:08 -0700343
murgatroid992e012342016-11-10 18:24:08 -0800344static int uv_get_fd(grpc_endpoint *ep) { return -1; }
345
murgatroid9969259d42016-10-31 14:34:10 -0700346static grpc_endpoint_vtable vtable = {
347 uv_endpoint_read, uv_endpoint_write, uv_get_workqueue,
348 uv_add_to_pollset, uv_add_to_pollset_set, uv_endpoint_shutdown,
murgatroid992e012342016-11-10 18:24:08 -0800349 uv_destroy, uv_get_resource_user, uv_get_peer,
350 uv_get_fd};
murgatroid999030c812016-09-16 13:25:08 -0700351
murgatroid9969259d42016-10-31 14:34:10 -0700352grpc_endpoint *grpc_tcp_create(uv_tcp_t *handle,
353 grpc_resource_quota *resource_quota,
354 char *peer_string) {
murgatroid999030c812016-09-16 13:25:08 -0700355 grpc_tcp *tcp = (grpc_tcp *)gpr_malloc(sizeof(grpc_tcp));
murgatroid99bb8b1c92017-06-23 15:10:18 -0700356 grpc_exec_ctx exec_ctx = GRPC_EXEC_CTX_INIT;
murgatroid999030c812016-09-16 13:25:08 -0700357
Craig Tiller341bcc52017-05-09 08:37:44 -0700358 if (GRPC_TRACER_ON(grpc_tcp_trace)) {
murgatroid999030c812016-09-16 13:25:08 -0700359 gpr_log(GPR_DEBUG, "Creating TCP endpoint %p", tcp);
360 }
361
murgatroid9904d28292016-10-27 14:36:57 -0700362 /* Disable Nagle's Algorithm */
murgatroid9912e57752016-10-27 14:30:41 -0700363 uv_tcp_nodelay(handle, 1);
364
murgatroid999030c812016-09-16 13:25:08 -0700365 memset(tcp, 0, sizeof(grpc_tcp));
366 tcp->base.vtable = &vtable;
367 tcp->handle = handle;
368 handle->data = tcp;
369 gpr_ref_init(&tcp->refcount, 1);
370 tcp->peer_string = gpr_strdup(peer_string);
murgatroid992c287ca2016-10-07 09:55:35 -0700371 tcp->shutting_down = false;
murgatroid99e2672c92016-11-09 15:12:22 -0800372 tcp->resource_user = grpc_resource_user_create(resource_quota, peer_string);
murgatroid99bb8b1c92017-06-23 15:10:18 -0700373 tcp->read_slice = alloc_read_slice(&exec_ctx, tcp->resource_user);
murgatroid999030c812016-09-16 13:25:08 -0700374 /* Tell network status tracking code about the new endpoint */
375 grpc_network_status_register_endpoint(&tcp->base);
376
377#ifndef GRPC_UV_TCP_HOLD_LOOP
378 uv_unref((uv_handle_t *)handle);
379#endif
380
murgatroid99bb8b1c92017-06-23 15:10:18 -0700381 grpc_exec_ctx_finish(&exec_ctx);
382
murgatroid999030c812016-09-16 13:25:08 -0700383 return &tcp->base;
384}
385
386#endif /* GRPC_UV */