Skip to content
Merged
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
33 changes: 24 additions & 9 deletions protocols/turnloop-mysql/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,10 @@ caching_sha2 RSA authentication obtains fresh OAEP entropy from the TLS provider
Construct Config and Connection, connect transport in the host, then feed
plaintext into `receive`. Pull `next_event` until None; `Progress` means a control
packet was consumed and polling should continue. Write `output`, acknowledging
only successfully written bytes via `consume_output`. Borrowed outputs/events
only successfully written bytes via `consume_output`. When it returns `true`, an
event is ready without further input: call `next_event` immediately. COM_STMT_CLOSE
and COM_QUIT get no server reply, so that acknowledgement is their only wakeup;
a completion-driven host needs no timer for them. Borrowed outputs/events
remain valid until the next mutable call. Preserve the output borrow through
write completion or copy into a reusable host transport buffer.

Expand All @@ -58,14 +61,18 @@ Provide absolute deadlines; schedule `next_timeout` and call
deadline; Config does not read a clock. On EOF/TLS failure or a parsing error,
call `abort(error)` and drain events. A parsing error is terminal; never continue
the byte stream after it. Each accepted command yields exactly one Completed;
server Error and result Ok events are informational. Close emits one Closed.
server Error, Ok and Eof events are informational. Close emits one Closed.
Command rejection accepts no token. There are no callbacks from this core.

MySQL permits one active command. Busy calls return backpressure; the adapter
MySQL permits one active command. Busy calls return backpressure;
`can_accept()` reports that exact admission decision up front (the command
methods call it themselves), so a host queue never copies the rule. The adapter
queues commands in JS submission order. This prevents unsynchronized packet
sequence resets. COM_QUERY supports multiple result sets (multiple statements
are opt-in); EOF negotiation deliberately selects legacy EOF, which MySQL 9.6
supports. Result Ok includes affected_rows, last_insert_id, warnings and status.
supports. `Ok` is a real OK packet (affected_rows, last_insert_id, warnings,
status); `Eof` ends a result set's rows and carries only warnings and status,
so a host need not track whether a result set is open to read either one.

Prepare emits parameter/column metadata and a Statement ID. Execute accepts
mysql_common Value parameters and emits binary rows. Reset/close statements,
Expand Down Expand Up @@ -94,18 +101,26 @@ Zstd and MariaDB extensions are out of scope.
## mysql2 conversions and hooks

Row iterators return borrowed bytes or mysql_common numeric/calendar scalars.
Each value is validated and decoded once, as the iterator reaches it (the last
one also rejects trailing bytes); the first error ends the row. Such an error
is a parsing error: abort the connection.
`types::decode` implements a default conversion policy. `Column` retains names,
original names/table/schema, flags, charset, length, decimals and wire type.
original names/table/schema, flags, charset, length, decimals and wire type;
every `Row` also carries its columns' `ColumnTypeInfo` (type, flags, charset,
decimals) via `Row::columns` and `Row::typed`, so a host can implement its own
policy per MySQL column type without reimplementing value decoding.

| MySQL type | Default policy / host action |
|---|---|
| integer, FLOAT, DOUBLE | Number; TINYINT(1) remains numeric |
| BIGINT | Number by default; support_big_numbers preserves unsafe-range values as strings; big_number_strings forces strings when enabled |
| DECIMAL/NEWDECIMAL | exact string; decimal_numbers opts into f64 |
| DATE/DATETIME/TIMESTAMP | explicit Date request with raw text/calendar components; date_strings formats strings |
| TIME | string, including negative and >24-hour values |
| DATE/DATETIME/TIMESTAMP | explicit Date request with raw text/calendar components; date_strings formats strings; truncate_fraction_to_decimals cuts binary-protocol fractions to the column's decimals like mysql2 |
| TIME | string, including negative and >24-hour values (by column type, even though its charset is 63) |
| JSON | explicit host JSON parse request; json_strings returns raw text |
| BLOB/binary charset 63, BIT, geometry | Buffer bytes |
| BIT | Buffer bytes |
| GEOMETRY | explicit Geometry request with the SRID+WKB bytes; the host builds mysql2's objects |
| BLOB/other string types with binary charset 63 | Buffer bytes |
| text | UTF-8 string; other character sets require host decoding |
| NULL | Null |

Expand All @@ -114,7 +129,7 @@ behavior), parses JSON and materializes row objects or rowsAsArray tuples.
The protocol preserves microseconds; JS Date loses sub-millisecond precision.
`typeCast` can inspect Column and RawValue in the event-dispatch layer and invoke
`types::decode` for next(). Callback invocation, field.string/buffer single-use
semantics and field.geometry parsing are adapter work, not core callbacks.
semantics and building field.geometry objects are adapter work, not core callbacks.
Per-type dateStrings arrays are not implemented (only the boolean option).
BIGINT inside JSON is still host policy. This is documented surface support,
not a drop-in mysql2 API.
Expand Down
4 changes: 3 additions & 1 deletion protocols/turnloop-mysql/src/asynchronous.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,9 @@ impl Output for crate::Connection {
self.output()
}
fn consume_output(&mut self, n: usize) -> io::Result<()> {
self.consume_output(n).map_err(io::Error::other)
// The driver polls `event` after every drain, so the "call next_event
// now" signal is already honoured here.
self.consume_output(n).map(drop).map_err(io::Error::other)
}
}
impl SansIo for crate::Connection {
Expand Down
75 changes: 62 additions & 13 deletions protocols/turnloop-mysql/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -140,10 +140,21 @@ pub enum Event<'a> {
token: Token,
row: Row<'a>,
},
/// A genuine OK packet: a command or one statement of a multi-statement
/// query finished without a result set. `affected_rows`, `last_insert_id`
/// and `info` are meaningful.
Ok {
token: Token,
packet: OkPacket<'a>,
},
/// The EOF packet ending one result set's rows. It carries no row counts,
/// only the warning count and status flags (`SERVER_MORE_RESULTS_EXISTS`
/// announces another result).
Eof {
token: Token,
warnings: u16,
status: StatusFlags,
},
Prepared {
token: Token,
statement: Statement,
Expand Down Expand Up @@ -193,6 +204,8 @@ enum CommandKind {
Execute,
Reset,
ChangeUser,
/// COM_STMT_CLOSE for this statement ID; forgotten once its bytes are sent.
CloseStatement(u32),
Other,
}
struct Pending {
Expand Down Expand Up @@ -270,16 +283,36 @@ impl Connection {
pub fn output(&self) -> &[u8] {
&self.output[self.output_at..]
}
pub fn consume_output(&mut self, n: usize) -> Result<()> {
/// Acknowledge `n` written bytes of `output()`.
///
/// Returns `true` when an event is now deliverable without further input,
/// in which case the host must call `next_event()` right away. This matters
/// for COM_STMT_CLOSE and COM_QUIT: the server never answers them, so their
/// `Completed` / `Closed` is produced by this acknowledgement itself and no
/// read completion or timer would otherwise wake a completion-driven host.
pub fn consume_output(&mut self, n: usize) -> Result<bool> {
if n > self.output().len() {
return Err(Error::State("invalid output acknowledgement"));
}
self.output_at += n;
if self.output_at == self.output.len() {
self.output.clear();
self.output_at = 0;
if self.state == State::NoResponse {
// The server holds the statement until these bytes reach it, so
// it is dropped from the bookkeeping only now; a failed flush
// aborts instead and leaves it listed.
if let Some(Pending {
kind: CommandKind::CloseStatement(id),
..
}) = self.pending
{
self.statements.retain(|s| s.id != id);
}
self.complete(Outcome::Success);
}
}
Ok(())
Ok(self.completion.is_some() || (self.state == State::Closing && self.output.is_empty()))
}
pub fn receive(&mut self, b: &[u8]) -> Result<()> {
if matches!(
Expand Down Expand Up @@ -403,8 +436,16 @@ impl Connection {
self.state = State::Auth;
Ok(())
}
/// Whether a command submitted now would be admitted. This is the exact rule
/// every command method applies: the session is authenticated and idle, the
/// previous command's `Completed` has been delivered and all output has been
/// acknowledged. A host queue can gate submission on it instead of copying
/// the preconditions. Unlike `is_ready`, it also requires flushed output.
pub fn can_accept(&self) -> bool {
self.state == State::Ready && self.pending.is_none() && self.output().is_empty()
}
fn accept(&self) -> Result<()> {
if self.state != State::Ready || self.pending.is_some() || !self.output().is_empty() {
if !self.can_accept() {
Err(Error::State("connection busy or closed"))
} else {
Ok(())
Expand Down Expand Up @@ -537,9 +578,18 @@ impl Connection {
self.simple(token, 0x1a, Some(id), CommandKind::Other, State::Header)
}
pub fn close_statement(&mut self, token: Token, id: u32) -> Result<()> {
self.simple(token, 0x19, Some(id), CommandKind::Other, State::NoResponse)?;
self.statements.retain(|s| s.id != id);
Ok(())
self.simple(
token,
0x19,
Some(id),
CommandKind::CloseStatement(id),
State::NoResponse,
)
}
/// Prepared statements the server still holds for this session. A statement
/// being closed stays listed until its COM_STMT_CLOSE bytes are acknowledged.
pub fn statements(&self) -> &[Statement] {
&self.statements
}
pub fn change_user(
&mut self,
Expand Down Expand Up @@ -642,9 +692,6 @@ impl Connection {
self.state = State::Ready;
}
pub fn next_event(&mut self) -> Result<Option<Event<'_>>> {
if self.state == State::NoResponse && self.output().is_empty() {
self.complete(Outcome::Success);
}
if let Some(outcome) = self.completion.take() {
let p = self
.pending
Expand Down Expand Up @@ -902,8 +949,9 @@ impl Connection {
}
State::Rows => {
if self.packet[0] == 0xfe && self.packet.len() < 9 {
let ok = parse_eof(&self.packet, self.caps)?;
self.status = ok.status_flags();
let eof = parse_eof(&self.packet, self.caps)?;
let warnings = eof.warnings();
self.status = eof.status_flags();
if self
.status
.contains(StatusFlags::SERVER_MORE_RESULTS_EXISTS)
Expand All @@ -912,9 +960,10 @@ impl Connection {
} else {
self.complete(Outcome::Success);
}
return Ok(Some(Event::Ok {
return Ok(Some(Event::Eof {
token: self.token()?,
packet: parse_eof(&self.packet, self.caps)?,
warnings,
status: self.status,
}));
}
return Ok(Some(Event::Row {
Expand Down
104 changes: 100 additions & 4 deletions protocols/turnloop-mysql/src/types.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,12 @@ pub struct Options {
pub big_number_strings: bool,
pub date_strings: bool,
pub json_strings: bool,
/// mysql2's binary-protocol DATETIME/TIMESTAMP string policy: cut the
/// fractional seconds of a formatted `date_strings` value to the column's
/// declared `decimals` (DATETIME(3) gives `.123`, not `.123000`). Without
/// it all six digits are kept. Text-protocol values already arrive with the
/// declared digits and are returned verbatim either way.
pub truncate_fraction_to_decimals: bool,
}
#[derive(Debug, Clone, PartialEq)]
pub enum Date<'a> {
Expand All @@ -31,6 +37,9 @@ pub enum JsValue<'a> {
Buffer(Cow<'a, [u8]>),
Json(&'a str),
Date(Date<'a>),
/// A GEOMETRY value (4-byte little-endian SRID, then WKB). mysql2 turns it
/// into point/array objects; like `Json` and `Date` that is host work.
Geometry(&'a [u8]),
}
fn utf8(b: &[u8]) -> Result<&str> {
std::str::from_utf8(b)
Expand All @@ -52,7 +61,13 @@ fn bigint(value: i128, options: Options) -> JsValue<'static> {
}
}
/// Invoke this default conversion after a host typeCast hook chooses `next()`;
/// raw fields and Column metadata are also available without conversion.
/// raw fields and Column metadata are also available without conversion
/// (`Row::typed` pairs every value with its `ColumnTypeInfo`).
///
/// The policy is chosen by MySQL column type first. `character_set == 63`
/// (binary) is reported for every non-string column, so it only selects Buffer
/// for the string/BLOB family; TIME stays a string, BIT a Buffer and GEOMETRY a
/// `Geometry` request, matching mysql2.
pub fn decode<'a>(
info: ColumnTypeInfo,
value: RawValue<'a>,
Expand Down Expand Up @@ -98,7 +113,9 @@ pub fn decode<'a>(
JsValue::Json(utf8(bytes)?)
}
}
16 | 255 => JsValue::Buffer(bytes.into()),
11 => JsValue::String(utf8(bytes)?.into()),
16 => JsValue::Buffer(bytes.into()),
255 => JsValue::Geometry(bytes),
_ if info.character_set == 63 => JsValue::Buffer(bytes.into()),
_ => JsValue::String(utf8(bytes)?.into()),
},
Expand Down Expand Up @@ -129,8 +146,15 @@ pub fn decode<'a>(
} else {
let mut s =
format!("{year:04}-{month:02}-{day:02} {hour:02}:{minute:02}:{second:02}");
if microsecond != 0 {
s.push_str(&format!(".{microsecond:06}"));
let digits = if options.truncate_fraction_to_decimals {
usize::from(info.decimals.min(6))
} else {
6
};
if microsecond != 0 && digits > 0 {
let fraction = format!("{microsecond:06}");
s.push('.');
s.push_str(&fraction[..digits]);
}
s
};
Expand Down Expand Up @@ -169,8 +193,80 @@ mod tests {
column_type: t,
flags: ColumnFlags::empty(),
character_set: 45,
decimals: 0,
}
}
/// MySQL reports charset 63 for every non-string column; the column type,
/// not the charset, must pick TIME/BIT/GEOMETRY's policy.
#[test]
fn binary_charset_does_not_collapse_time_bit_and_geometry() {
let binary = |t| ColumnTypeInfo {
character_set: 63,
..info(t)
};
let decode = |t, bytes| decode(binary(t), RawValue::Bytes(bytes), Options::default());
assert_eq!(
decode(ColumnType::MYSQL_TYPE_TIME, b"-838:59:59.5"),
Ok(JsValue::String("-838:59:59.5".into()))
);
assert_eq!(
decode(ColumnType::MYSQL_TYPE_BIT, &[0b101]),
Ok(JsValue::Buffer(Cow::Borrowed(&[0b101])))
);
let point = [0, 0, 0, 0, 1, 1, 0, 0, 0];
assert_eq!(
decode(ColumnType::MYSQL_TYPE_GEOMETRY, &point),
Ok(JsValue::Geometry(&point))
);
assert_eq!(
decode(ColumnType::MYSQL_TYPE_VAR_STRING, b"\xff\x00"),
Ok(JsValue::Buffer(Cow::Borrowed(b"\xff\x00"))),
"VARBINARY still follows the binary charset"
);
}
#[test]
fn datetime_fraction_follows_the_column_decimals_on_request() {
let value = || RawValue::Scalar(Value::Date(2026, 9, 22, 10, 11, 12, 123_456));
let datetime = |decimals| ColumnTypeInfo {
decimals,
..info(ColumnType::MYSQL_TYPE_DATETIME)
};
let strings = Options {
date_strings: true,
..Options::default()
};
let truncated = Options {
truncate_fraction_to_decimals: true,
..strings
};
assert_eq!(
decode(datetime(3), value(), strings),
Ok(JsValue::String("2026-09-22 10:11:12.123456".into()))
);
assert_eq!(
decode(datetime(3), value(), truncated),
Ok(JsValue::String("2026-09-22 10:11:12.123".into()))
);
assert_eq!(
decode(datetime(6), value(), truncated),
Ok(JsValue::String("2026-09-22 10:11:12.123456".into()))
);
assert_eq!(
decode(datetime(0), value(), truncated),
Ok(JsValue::String("2026-09-22 10:11:12".into()))
);
assert_eq!(
decode(
ColumnTypeInfo {
decimals: 2,
..info(ColumnType::MYSQL_TYPE_TIMESTAMP)
},
value(),
truncated
),
Ok(JsValue::String("2026-09-22 10:11:12.12".into()))
);
}
#[test]
fn node_conversion_options() {
let decimal = info(ColumnType::MYSQL_TYPE_NEWDECIMAL);
Expand Down
Loading