Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions src/sql_jsc/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,8 @@ pub mod shared {
#[path = "SQLDataCell.rs"]
pub mod sql_data_cell;

pub mod socket_teardown;

pub use cached_structure::CachedStructure;
pub(crate) use query_binding_iterator::QueryBindingIterator;
}
5 changes: 4 additions & 1 deletion src/sql_jsc/mysql/MySQLConnection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@ use crate::mysql::js_mysql_connection::JSMySQLConnection;
use crate::mysql::js_mysql_query::JSMySQLQuery;
use crate::mysql::my_sql_request_queue::MySQLRequestQueue;
use crate::mysql::my_sql_statement::{self as mysql_statement, MySQLStatement, Param};
use crate::shared::socket_teardown;
use bun_ptr::RefPtr;

pub use bun_sql::mysql::protocol::error_packet::ErrorPacket;
Expand Down Expand Up @@ -282,8 +283,10 @@ impl MySQLConnection {
}

pub(crate) fn close(&mut self) {
self.socket.close(uws::CloseKind::Normal);
self.write_buffer = OffsetByteList::default();
// A copy: the close dispatches `on_close`, which detaches `self.socket`.
let socket = self.socket;
socket_teardown::close_now(&socket);
}

pub(crate) fn clean_queue_and_close(
Expand Down
7 changes: 5 additions & 2 deletions src/sql_jsc/postgres/PostgresSQLConnection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@ use crate::postgres::postgres_sql_statement::{Error as StatementError, Status as
use crate::postgres::sasl::SASLStatus;
use crate::shared::CachedStructure as PostgresCachedStructure;
use crate::shared::connection_ctor_args::ConnectionCtorArgs;
use crate::shared::socket_teardown;
use bun_sql::postgres::AnyPostgresError;
use bun_sql::postgres::PostgresErrorOptions;
use bun_sql::postgres::PostgresProtocol as protocol;
Expand Down Expand Up @@ -1502,11 +1503,13 @@ impl PostgresSQLConnection {
fn ref_and_close(&self, js_reason: Option<JSValue>) {
// refAndClose is always called when we wanna to disconnect or when we are closed

if !self.socket.get().is_closed() {
// A copy: the close dispatches `on_close`, which detaches `self.socket`.
let socket = *self.socket.get();
if !socket.is_closed() {
// event loop need to be alive to close the socket
self.poll_ref.with_mut(|r| r.ref_(self.vm_ctx()));
// will unref on socket close
self.socket.get().close(uws::CloseKind::Normal);
socket_teardown::close_now(&socket);
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}

// cleanup requests
Expand Down
17 changes: 17 additions & 0 deletions src/sql_jsc/shared/socket_teardown.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
//! Teardown of a database socket the client has given up on.

use bun_uws::{AnySocket, CloseCode};

/// Closes `socket` before this returns, whatever the peer does.
/// `CloseCode::Normal` on a TLS socket would wait for the peer's close_notify.
Comment thread
robobun marked this conversation as resolved.
pub(crate) fn close_now(socket: &AnySocket) {
// shutdown() before the handshake finished makes usockets report the handshake as failed
if matches!(socket, AnySocket::SocketTls(_)) && socket.is_ssl_handshake_finished() {
socket.shutdown();
}
socket.close(CloseCode::FastShutdown);
// usockets parks the first fast shutdown behind unsent ciphertext; the second drops it
if !socket.is_closed() {
socket.close(CloseCode::FastShutdown);
}
}
11 changes: 11 additions & 0 deletions src/uws_sys/socket.rs
Original file line number Diff line number Diff line change
Expand Up @@ -287,6 +287,16 @@ impl<const IS_SSL: bool> NewSocketHandler<IS_SSL> {
)
}

/// Always true for a plain TCP socket. False while connecting or once detached.
pub fn is_ssl_handshake_finished(&self) -> bool {
on_socket!(self.socket;
connected s => s.is_ssl_handshake_finished(),
duplex d => d.is_established(),
pipe p => p.is_established(),
else => false,
)
}

#[inline]
pub fn is_closed_or_has_error(&self) -> bool {
self.is_closed() || self.is_shutdown() || self.get_error() != 0
Expand Down Expand Up @@ -949,6 +959,7 @@ impl AnySocket {
fn is_closed(&self) -> bool;
fn is_shutdown(&self) -> bool;
fn is_established(&self) -> bool;
fn is_ssl_handshake_finished(&self) -> bool;
fn close(&self, code: CloseCode);
fn write(&self, data: &[u8]) -> i32;
fn set_timeout(&self, seconds: c_uint);
Expand Down
6 changes: 6 additions & 0 deletions src/uws_sys/us_socket_t.rs
Original file line number Diff line number Diff line change
Expand Up @@ -461,6 +461,11 @@ impl us_socket_t {
c::us_socket_is_established(self) > 0
}

/// Always true for a plain TCP socket.
pub(crate) fn is_ssl_handshake_finished(&self) -> bool {
c::us_socket_is_ssl_handshake_finished(self) > 0
}

pub(crate) fn queued_input(&self) -> QueuedInput {
match c::us_socket_queued_input(self) {
LIBUS_QUEUED_INPUT_DATA => QueuedInput::Data,
Expand Down Expand Up @@ -568,6 +573,7 @@ mod c {
pub(super) safe fn us_socket_verify_error(s: &us_socket_t) -> us_bun_verify_error_t;
pub(super) safe fn us_socket_get_error(s: &us_socket_t) -> c_int;
pub(super) safe fn us_socket_is_established(s: &us_socket_t) -> i32;
pub(super) safe fn us_socket_is_ssl_handshake_finished(s: &us_socket_t) -> i32;
pub(super) safe fn us_socket_queued_input(s: &us_socket_t) -> c_int;

/// ssl_ctx is required (the whole point); sni may be null.
Expand Down
Loading
Loading