Skip to content

Commit 477092e

Browse files
committed
Clarify synchronous handshake blocking and error handling
1 parent b3cb1e0 commit 477092e

1 file changed

Lines changed: 13 additions & 2 deletions

File tree

‎src/tcp_server.cpp‎

Lines changed: 13 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -108,6 +108,9 @@ class sync_write_handler {
108108
apply_send_timeout(*sock);
109109
// An inlet may return from open_stream() as soon as these bytes arrive.
110110
// Keep sample writes locked out until the socket is registered below.
111+
// This synchronous-mode handshake runs on the server IO thread while
112+
// holding the sample-write lock. It uses the configured sync send timeout;
113+
// a timeout of zero permits unbounded blocking of both IO and sample writes.
111114
asio::write(*sock, header);
112115
if (reverse_byte_order) {
113116
sockets_swapped_.push_back(std::move(sock));
@@ -714,15 +717,23 @@ void client_session::handle_read_feedparams(
714717
// The synchronous writer sends the header under its sample-write lock,
715718
// then registers the socket before an immediate push can acquire it.
716719
auto protocol = sock_.local_endpoint().protocol();
717-
serv->sync_handler_->add_socket(
718-
sock_.release(), protocol, reverse_byte_order_, feedbuf_.data());
720+
try {
721+
serv->sync_handler_->add_socket(
722+
sock_.release(), protocol, reverse_byte_order_, feedbuf_.data());
723+
} catch (const asio::system_error &e) {
724+
// add_socket owns and closes the released handle on failure. Peer
725+
// resets here are connection failures, not serialization errors.
726+
LOG_F(INFO, "Synchronous stream handshake failed: %s", e.what());
727+
}
719728
serv->unregister_inflight_session(this);
720729
return;
721730
}
722731

723732
// Subscribe before sending the handshake: receiving its test patterns is
724733
// what lets the inlet return from open_stream(). Queue samples until the
725734
// header write completes so data cannot overtake the handshake.
735+
// have_consumers()/wait_for_consumers() can now observe this subscription
736+
// before the handshake completes; open_stream() remains the inlet's gate.
726737
std::shared_ptr<consumer_queue> queue;
727738
if (max_buffered_ > 0) queue = serv->send_buffer_->new_consumer(max_buffered_);
728739
async_write(sock_, feedbuf_.data(),

0 commit comments

Comments
 (0)