diff --git a/httplib.h b/httplib.h index e4044c57..a53210d7 100644 --- a/httplib.h +++ b/httplib.h @@ -1546,6 +1546,18 @@ public: (void)usec; } + // Bytes already pulled off the socket and sitting in this stream's own + // buffer. Exposing them lets a line reader scan for a terminator in one + // pass instead of asking for a byte at a time. A stream that does no + // buffering of its own reports none, and readers fall back to read(). + virtual const char *buffered_data(size_t &size) const { + size = 0; + return nullptr; + } + + // Discards `size` bytes previously returned by buffered_data(). + virtual void consume_buffered(size_t size) { (void)size; } + ssize_t write(const char *ptr); ssize_t write(const std::string &s); @@ -3370,6 +3382,7 @@ public: private: void append(char c); + void append(const char *data, size_t size); Stream &strm_; char *fixed_buffer_; @@ -5463,6 +5476,46 @@ inline bool stream_line_reader::getline() { #endif for (size_t i = 0;; i++) { + // Fast path: whatever the stream has already buffered can be scanned for + // the terminator in one pass. Asking for a byte at a time costs a virtual + // call, a bounds check and a one-byte copy per character of the request. + size_t buffered_size = 0; + if (auto buffered = strm_.buffered_data(buffered_size)) { + auto take = buffered_size; + auto terminated = false; + + for (size_t at = 0; at < buffered_size;) { + auto nl = static_cast( + memchr(buffered + at, '\n', buffered_size - at)); + if (!nl) { break; } + auto pos = static_cast(nl - buffered); +#ifdef CPPHTTPLIB_ALLOW_LF_AS_LINE_TERMINATOR + take = pos + 1; + terminated = true; + break; +#else + // A bare LF does not end the line; keep looking for CRLF. The CR may + // be the last byte of an earlier chunk, hence prev_byte. + if ((pos > 0 ? buffered[pos - 1] : prev_byte) == '\r') { + take = pos + 1; + terminated = true; + break; + } + at = pos + 1; +#endif + } + + if (size() + take > CPPHTTPLIB_MAX_LINE_LENGTH) { return false; } +#ifndef CPPHTTPLIB_ALLOW_LF_AS_LINE_TERMINATOR + prev_byte = buffered[take - 1]; +#endif + append(buffered, take); + strm_.consume_buffered(take); + i += take; + if (terminated) { return true; } + continue; + } + if (size() >= CPPHTTPLIB_MAX_LINE_LENGTH) { // Treat exceptionally long lines as an error to // prevent infinite loops/memory exhaustion @@ -5494,16 +5547,26 @@ inline bool stream_line_reader::getline() { return true; } -inline void stream_line_reader::append(char c) { - if (fixed_buffer_used_size_ < fixed_buffer_size_ - 1) { - fixed_buffer_[fixed_buffer_used_size_++] = c; +inline void stream_line_reader::append(char c) { append(&c, 1); } + +inline void stream_line_reader::append(const char *data, size_t size) { + // Once the line has outgrown the fixed buffer everything must keep going to + // the growable one, even if a later chunk would have fit. Without the + // emptiness check a short append after a long one would land in the fixed + // buffer, which ptr() and size() no longer look at, and be lost. + if (growable_buffer_.empty() && + fixed_buffer_used_size_ + size < fixed_buffer_size_) { + memcpy(fixed_buffer_ + fixed_buffer_used_size_, data, size); + fixed_buffer_used_size_ += size; fixed_buffer_[fixed_buffer_used_size_] = '\0'; } else { + // Unlike the per-character overload, this can be the very first append of + // the line, so the fixed buffer may hold nothing and carry no terminator + // yet. assign() takes an explicit length and does not need one. if (growable_buffer_.empty()) { - assert(fixed_buffer_[fixed_buffer_used_size_] == '\0'); growable_buffer_.assign(fixed_buffer_, fixed_buffer_used_size_); } - growable_buffer_ += c; + growable_buffer_.append(data, size); } } @@ -5761,8 +5824,17 @@ public: socket_t socket() const override; time_t duration() const override; void set_read_timeout(time_t sec, time_t usec = 0) override; + const char *buffered_data(size_t &size) const override; + void consume_buffered(size_t size) override; + + // The caller has just seen this socket become readable. Lets the next read + // skip its own readiness wait, which would otherwise ask the kernel a + // question that was answered a moment ago. Consumed by that read. + void set_readable_hint() { readable_hint_ = true; } private: + bool ensure_readable(); + socket_t sock_; time_t read_timeout_sec_; time_t read_timeout_usec_; @@ -5774,6 +5846,7 @@ private: std::vector read_buff_; size_t read_buff_off_ = 0; size_t read_buff_content_size_ = 0; + bool readable_hint_ = false; static const size_t read_buff_size_ = 1024l * 4; }; @@ -5841,6 +5914,9 @@ process_server_socket(const std::atomic &svr_sock, socket_t sock, [&](bool close_connection, bool &connection_closed) { SocketStream strm(sock, read_timeout_sec, read_timeout_usec, write_timeout_sec, write_timeout_usec); + // process_server_socket_core() only gets here once keep_alive() has + // seen the socket go readable. + strm.set_readable_hint(); return callback(strm, close_connection, connection_closed); }); } @@ -9204,7 +9280,12 @@ public: time_t duration() const override; void set_read_timeout(time_t sec, time_t usec = 0) override; + // See SocketStream::set_readable_hint(). + void set_readable_hint() { readable_hint_ = true; } + private: + bool ensure_readable(); + socket_t sock_; tls::session_t session_; time_t read_timeout_sec_; @@ -9213,6 +9294,7 @@ private: time_t write_timeout_usec_; time_t max_timeout_msec_; const std::chrono::time_point start_time_; + bool readable_hint_ = false; }; #ifdef CPPHTTPLIB_OPENSSL_SUPPORT @@ -9367,6 +9449,8 @@ inline bool process_server_socket_ssl( [&](bool close_connection, bool &connection_closed) { SSLSocketStream strm(sock, session, read_timeout_sec, read_timeout_usec, write_timeout_sec, write_timeout_usec); + // See the non-TLS path in process_server_socket(). + strm.set_readable_hint(); return callback(strm, close_connection, connection_closed); }); } @@ -10749,6 +10833,24 @@ inline bool SocketStream::wait_writable() const { return select_write(sock_, write_timeout_sec_, write_timeout_usec_) > 0; } +inline bool SocketStream::ensure_readable() { + if (readable_hint_) { + readable_hint_ = false; + return true; + } + return wait_readable(); +} + +inline const char *SocketStream::buffered_data(size_t &size) const { + size = read_buff_content_size_ - read_buff_off_; + return size ? read_buff_.data() + read_buff_off_ : nullptr; +} + +inline void SocketStream::consume_buffered(size_t size) { + assert(size <= read_buff_content_size_ - read_buff_off_); + read_buff_off_ += size; +} + inline bool SocketStream::is_peer_alive() const { return detail::is_socket_alive(sock_); } @@ -10775,7 +10877,7 @@ inline ssize_t SocketStream::read(char *ptr, size_t size) { } } - if (!wait_readable()) { + if (!ensure_readable()) { error_ = Error::Timeout; return -1; } @@ -11253,6 +11355,14 @@ inline bool SSLSocketStream::wait_writable() const { !tls::is_peer_closed(session_, sock_); } +inline bool SSLSocketStream::ensure_readable() { + if (readable_hint_) { + readable_hint_ = false; + return true; + } + return wait_readable(); +} + inline bool SSLSocketStream::is_peer_alive() const { return !tls::is_peer_closed(session_, sock_); } @@ -11265,7 +11375,7 @@ inline ssize_t SSLSocketStream::read(char *ptr, size_t size) { error_ = Error::ConnectionClosed; } return ret; - } else if (wait_readable()) { + } else if (ensure_readable()) { tls::TlsError err; auto ret = tls::read(session_, ptr, size, err); if (ret < 0) {