diff --git a/include/fluent-bit/flb_io.h b/include/fluent-bit/flb_io.h index 9a6afd1e0cd..ec1a0130eae 100644 --- a/include/fluent-bit/flb_io.h +++ b/include/fluent-bit/flb_io.h @@ -39,6 +39,15 @@ /* Other features */ #define FLB_IO_IPV6 32 /* network I/O uses IPv6 */ +/* Retryable network I/O results */ +#define FLB_IO_WANT_READ -0x7e4 +#define FLB_IO_WANT_WRITE -0x7e6 + +static inline int flb_io_net_is_retry(ssize_t result) +{ + return result == FLB_IO_WANT_READ || result == FLB_IO_WANT_WRITE; +} + struct flb_connection; int flb_io_net_accept(struct flb_connection *connection, diff --git a/include/fluent-bit/tls/flb_tls.h b/include/fluent-bit/tls/flb_tls.h index 018231218e2..adefb119870 100644 --- a/include/fluent-bit/tls/flb_tls.h +++ b/include/fluent-bit/tls/flb_tls.h @@ -25,6 +25,7 @@ #include #include #include +#include #include #define FLB_TLS_ALPN_MAX_LENGTH 16 @@ -32,8 +33,8 @@ #define FLB_TLS_CLIENT "Fluent Bit" /* TLS backend return status on read/write */ -#define FLB_TLS_WANT_READ -0x7e4 -#define FLB_TLS_WANT_WRITE -0x7e6 +#define FLB_TLS_WANT_READ FLB_IO_WANT_READ +#define FLB_TLS_WANT_WRITE FLB_IO_WANT_WRITE /* Cert Flags */ #define FLB_TLS_CA_ROOT 1 diff --git a/plugins/in_forward/fw_conn.c b/plugins/in_forward/fw_conn.c index aca3c80eb7c..ee32f5ccf8d 100644 --- a/plugins/in_forward/fw_conn.c +++ b/plugins/in_forward/fw_conn.c @@ -62,6 +62,10 @@ static int fw_conn_event_internal(struct flb_connection *connection) flb_plg_trace(ctx->ins, "handshake status = %d", conn->handshake_status); ret = fw_prot_secure_forward_handshake(ctx->ins, conn); + if (flb_io_net_is_retry(ret)) { + return 0; + } + if (ret == -1) { flb_plg_trace(ctx->ins, "fd=%i closed connection", event->fd); fw_conn_del(conn); @@ -108,6 +112,10 @@ static int fw_conn_event_internal(struct flb_connection *connection) (void *) &conn->buf[conn->buf_len], available); + if (flb_io_net_is_retry(bytes)) { + return 0; + } + if (bytes > 0) { read_bytes = (size_t) bytes; flb_plg_trace(ctx->ins, "read()=%zd pre_len=%zu now_len=%zu", diff --git a/plugins/in_forward/fw_prot.c b/plugins/in_forward/fw_prot.c index ef8464bad46..32fde63bf7b 100644 --- a/plugins/in_forward/fw_prot.c +++ b/plugins/in_forward/fw_prot.c @@ -140,38 +140,46 @@ static inline void print_msgpack_error_code(struct flb_input_instance *in, /* Read a secure forward msgpack message for handshake */ static int secure_forward_read(struct flb_input_instance *in, - struct flb_connection *connection, - char *buf, size_t size, size_t *out_len) + struct fw_conn *conn, size_t size, + size_t *out_len) { int ret; size_t off; size_t avail; - size_t buf_off = 0; msgpack_unpacked result; msgpack_unpacked_init(&result); while (1) { - avail = size - buf_off; - if (avail < 1) { + if (conn->buf_len >= size) { goto error; } + avail = size - conn->buf_len; + /* Read the message */ - ret = flb_io_net_read(connection, buf + buf_off, size - buf_off); + ret = flb_io_net_read(conn->connection, + conn->buf + conn->buf_len, avail); + if (flb_io_net_is_retry(ret)) { + msgpack_unpacked_destroy(&result); + return ret; + } + if (ret <= 0) { flb_plg_debug(in, "read %d byte(s)", ret); goto error; } - buf_off += ret; + conn->buf_len += ret; /* Validate */ off = 0; - ret = msgpack_unpack_next(&result, buf, buf_off, &off); + ret = msgpack_unpack_next(&result, conn->buf, conn->buf_len, &off); switch (ret) { case MSGPACK_UNPACK_SUCCESS: msgpack_unpacked_destroy(&result); - *out_len = buf_off; + *out_len = conn->buf_len; return 0; + case MSGPACK_UNPACK_CONTINUE: + break; default: print_msgpack_error_code(in, ret, "handshake"); goto error; @@ -179,6 +187,7 @@ static int secure_forward_read(struct flb_input_instance *in, } error: + conn->buf_len = 0; msgpack_unpacked_destroy(&result); return -1; } @@ -506,7 +515,6 @@ static int check_ping(struct flb_input_instance *ins, flb_sds_t *out_shared_key_salt) { int ret; - char buf[1024]; size_t out_len; size_t off; msgpack_unpacked result; @@ -531,7 +539,12 @@ static int check_ping(struct flb_input_instance *ins, } /* Wait for client PING */ - ret = secure_forward_read(ins, conn->connection, buf, sizeof(buf) - 1, &out_len); + ret = secure_forward_read(ins, conn, 1023, &out_len); + if (flb_io_net_is_retry(ret)) { + flb_free(serverside); + return ret; + } + if (ret == -1) { flb_free(serverside); flb_plg_error(ins, "handshake error expecting PING"); @@ -541,7 +554,8 @@ static int check_ping(struct flb_input_instance *ins, /* Unpack message and validate */ off = 0; msgpack_unpacked_init(&result); - ret = msgpack_unpack_next(&result, buf, out_len, &off); + ret = msgpack_unpack_next(&result, conn->buf, out_len, &off); + conn->buf_len = 0; if (ret != MSGPACK_UNPACK_SUCCESS) { flb_free(serverside); print_msgpack_error_code(ins, ret, "PING"); @@ -1206,6 +1220,11 @@ int fw_prot_secure_forward_handshake(struct flb_input_instance *ins, reason = flb_sds_create_size(32); flb_plg_debug(ins, "protocol: checking PING"); ping_ret = check_ping(ins, conn, &shared_key_salt); + if (flb_io_net_is_retry(ping_ret)) { + flb_sds_destroy(reason); + return ping_ret; + } + if (ping_ret == -1) { flb_plg_error(ins, "handshake error checking PING"); diff --git a/plugins/in_http/http_conn.c b/plugins/in_http/http_conn.c index 73f1ed53117..21801965e65 100644 --- a/plugins/in_http/http_conn.c +++ b/plugins/in_http/http_conn.c @@ -107,6 +107,10 @@ static int http_conn_event(void *data) (void *) &conn->buf_data[conn->buf_len], available); + if (flb_io_net_is_retry(bytes)) { + return 0; + } + if (bytes <= 0) { flb_plg_trace(ctx->ins, "fd=%i closed connection", event->fd); http_conn_del(conn); diff --git a/plugins/in_mqtt/mqtt_conn.c b/plugins/in_mqtt/mqtt_conn.c index bceecde89bd..f35163a184b 100644 --- a/plugins/in_mqtt/mqtt_conn.c +++ b/plugins/in_mqtt/mqtt_conn.c @@ -54,6 +54,10 @@ int mqtt_conn_event(void *data) (void *) &conn->buf[conn->buf_len], available); + if (flb_io_net_is_retry(bytes)) { + return 0; + } + if (bytes > 0) { conn->buf_len += bytes; flb_plg_trace(ctx->ins, "[fd=%i] read()=%i bytes", diff --git a/plugins/in_syslog/syslog_conn.c b/plugins/in_syslog/syslog_conn.c index 0aecbebb3e2..08559512bda 100644 --- a/plugins/in_syslog/syslog_conn.c +++ b/plugins/in_syslog/syslog_conn.c @@ -97,6 +97,10 @@ int syslog_stream_conn_event(void *data) (void *) &conn->buf_data[conn->buf_len], available); + if (flb_io_net_is_retry(bytes)) { + return 0; + } + if (bytes > 0) { flb_plg_trace(ctx->ins, "read()=%i pre_len=%zu now_len=%zu", bytes, conn->buf_len, conn->buf_len + bytes); diff --git a/plugins/in_tcp/tcp_conn.c b/plugins/in_tcp/tcp_conn.c index ddd15e57ab1..79e32c8a05f 100644 --- a/plugins/in_tcp/tcp_conn.c +++ b/plugins/in_tcp/tcp_conn.c @@ -452,6 +452,11 @@ int tcp_conn_event(void *data) (void *) &conn->buf_data[conn->buf_len], available); + if (flb_io_net_is_retry(bytes)) { + conn->busy = FLB_FALSE; + return 0; + } + if (bytes <= 0) { flb_plg_trace(ctx->ins, "fd=%i closed connection", event->fd); conn->busy = FLB_FALSE; diff --git a/plugins/in_unix_socket/unix_socket_conn.c b/plugins/in_unix_socket/unix_socket_conn.c index f65ced31e65..ef287cda87e 100644 --- a/plugins/in_unix_socket/unix_socket_conn.c +++ b/plugins/in_unix_socket/unix_socket_conn.c @@ -269,6 +269,10 @@ int unix_socket_conn_event(void *data) (void *) &conn->buf_data[conn->buf_len], available); + if (flb_io_net_is_retry(bytes)) { + return 0; + } + if (bytes <= 0) { if (!ctx->dgram_mode_flag) { flb_plg_trace(ctx->ins, "fd=%i closed connection", event->fd); diff --git a/src/http_server/flb_http_server.c b/src/http_server/flb_http_server.c index aa6c5c6e1a0..773a47746fd 100644 --- a/src/http_server/flb_http_server.c +++ b/src/http_server/flb_http_server.c @@ -247,6 +247,10 @@ static int flb_http_server_session_read(struct flb_http_server_session *session) (void *) session->read_buffer, session->read_buffer_size); + if (flb_io_net_is_retry(result)) { + return 0; + } + if (result <= 0) { return -1; } diff --git a/src/tls/flb_tls.c b/src/tls/flb_tls.c index 40961b6a12a..a7c936ce02e 100644 --- a/src/tls/flb_tls.c +++ b/src/tls/flb_tls.c @@ -162,7 +162,7 @@ static inline int io_tls_event_switch(struct flb_tls_session *session, event = &session->connection->event; event_loop = session->connection->evl; - if ((event->mask & mask) == 0) { + if ((event->mask & (MK_EVENT_READ | MK_EVENT_WRITE)) != mask) { ret = mk_event_add(event_loop, event->fd, FLB_ENGINE_EV_THREAD, @@ -345,12 +345,19 @@ int flb_tls_set_client_thumbprints(struct flb_tls *tls, const char *thumbprints) int flb_tls_net_read(struct flb_tls_session *session, void *buf, size_t len) { + int event_driven; time_t timeout_timestamp; time_t current_timestamp; struct flb_tls *tls; int ret; tls = session->tls; + event_driven = FLB_FALSE; + + if (session->connection->type == FLB_DOWNSTREAM_CONNECTION && + MK_EVENT_IS_REGISTERED((&session->connection->event))) { + event_driven = FLB_TRUE; + } if (session->connection->net->io_timeout > 0) { timeout_timestamp = time(NULL) + session->connection->net->io_timeout; @@ -365,6 +372,14 @@ int flb_tls_net_read(struct flb_tls_session *session, void *buf, size_t len) current_timestamp = time(NULL); if (ret == FLB_TLS_WANT_READ) { + if (event_driven) { + if (io_tls_event_switch(session, MK_EVENT_READ) == -1) { + return -1; + } + + return FLB_TLS_WANT_READ; + } + if (timeout_timestamp > 0 && timeout_timestamp <= current_timestamp) { return ret; @@ -373,6 +388,14 @@ int flb_tls_net_read(struct flb_tls_session *session, void *buf, size_t len) goto retry_read; } else if (ret == FLB_TLS_WANT_WRITE) { + if (event_driven) { + if (io_tls_event_switch(session, MK_EVENT_READ | MK_EVENT_WRITE) == -1) { + return -1; + } + + return FLB_TLS_WANT_WRITE; + } + goto retry_read; } else if (ret < 0) { @@ -382,6 +405,10 @@ int flb_tls_net_read(struct flb_tls_session *session, void *buf, size_t len) return -1; } + if (event_driven && io_tls_event_switch(session, MK_EVENT_READ) == -1) { + return -1; + } + return ret; }